255 lines
8.1 KiB
Python
255 lines
8.1 KiB
Python
from __future__ import annotations
|
|
|
|
import hashlib
|
|
import json
|
|
import shutil
|
|
from pathlib import Path
|
|
|
|
from sqlalchemy import select
|
|
from sqlalchemy.orm import selectinload
|
|
|
|
from app.modules.documents.models import (
|
|
DocumentDownloadRequest,
|
|
EngagementDocument,
|
|
EngagementDocumentVersion,
|
|
)
|
|
from app.modules.documents.services import (
|
|
create_download_request_for_version,
|
|
download_request_cache_path,
|
|
version_absolute_path,
|
|
)
|
|
|
|
from .client_context import analyzer_client_context
|
|
|
|
|
|
def _sha256(path: Path) -> str:
|
|
hasher = hashlib.sha256()
|
|
with path.open("rb") as handle:
|
|
while True:
|
|
chunk = handle.read(1024 * 1024)
|
|
if not chunk:
|
|
break
|
|
hasher.update(chunk)
|
|
return hasher.hexdigest()
|
|
|
|
|
|
def _allowed_engagement_ids(db, *, request, user, roles, client_id: int) -> set[int]:
|
|
context = analyzer_client_context(db, request=request, user=user, roles=roles)
|
|
visible_clients = {int(row.id) for row in context["clients"]}
|
|
if int(client_id) not in visible_clients:
|
|
raise PermissionError("Selected client is not visible for Bank Statement Analyzer.")
|
|
return {
|
|
int(row.id)
|
|
for row in context["engagements"]
|
|
if int(row.client_id) == int(client_id)
|
|
}
|
|
|
|
|
|
def list_stored_bank_statements(
|
|
db,
|
|
*,
|
|
request,
|
|
user,
|
|
roles,
|
|
client_id: int,
|
|
limit: int = 200,
|
|
):
|
|
allowed_engagements = _allowed_engagement_ids(
|
|
db,
|
|
request=request,
|
|
user=user,
|
|
roles=roles,
|
|
client_id=client_id,
|
|
)
|
|
if not allowed_engagements:
|
|
return []
|
|
|
|
docs = list(
|
|
db.execute(
|
|
select(EngagementDocument)
|
|
.options(selectinload(EngagementDocument.versions))
|
|
.where(
|
|
EngagementDocument.client_id == int(client_id),
|
|
EngagementDocument.engagement_id.in_(allowed_engagements),
|
|
EngagementDocument.document_type == "BANK_STATEMENT",
|
|
EngagementDocument.is_deleted.is_(False),
|
|
)
|
|
.order_by(EngagementDocument.updated_at_utc.desc(), EngagementDocument.id.desc())
|
|
.limit(max(1, min(500, int(limit))))
|
|
).scalars().all()
|
|
)
|
|
|
|
result = []
|
|
for doc in docs:
|
|
version = doc.versions[0] if doc.versions else None
|
|
if not version:
|
|
continue
|
|
result.append(
|
|
{
|
|
"document_id": int(doc.id),
|
|
"version_id": int(version.id),
|
|
"engagement_id": int(doc.engagement_id),
|
|
"title": doc.title,
|
|
"filename": version.original_filename,
|
|
"size_bytes": int(version.file_size_bytes or 0),
|
|
"sha256": version.file_hash_sha256,
|
|
"uploaded_at": version.uploaded_at_utc.isoformat() if version.uploaded_at_utc else "",
|
|
"storage_status": version.storage_status,
|
|
}
|
|
)
|
|
return result
|
|
|
|
|
|
def _ready_cached_path(db, *, version_id: int, user_id: int):
|
|
req = db.execute(
|
|
select(DocumentDownloadRequest)
|
|
.where(
|
|
DocumentDownloadRequest.version_id == int(version_id),
|
|
DocumentDownloadRequest.requested_by_user_id == int(user_id),
|
|
DocumentDownloadRequest.request_status == "ready",
|
|
)
|
|
.order_by(DocumentDownloadRequest.fulfilled_at_utc.desc(), DocumentDownloadRequest.id.desc())
|
|
.limit(1)
|
|
).scalar_one_or_none()
|
|
if not req:
|
|
return None
|
|
path = download_request_cache_path(req)
|
|
return path if path and path.exists() else None
|
|
|
|
|
|
def prepare_stored_bank_statements(
|
|
db,
|
|
*,
|
|
request,
|
|
user,
|
|
roles,
|
|
client_id: int,
|
|
version_ids: list[int],
|
|
input_dir: Path,
|
|
):
|
|
selected_ids = list(dict.fromkeys(int(value) for value in version_ids if int(value) > 0))
|
|
if not selected_ids:
|
|
return [], [], []
|
|
|
|
allowed_engagements = _allowed_engagement_ids(
|
|
db,
|
|
request=request,
|
|
user=user,
|
|
roles=roles,
|
|
client_id=client_id,
|
|
)
|
|
|
|
versions = list(
|
|
db.execute(
|
|
select(EngagementDocumentVersion)
|
|
.join(
|
|
EngagementDocument,
|
|
EngagementDocument.id == EngagementDocumentVersion.document_id,
|
|
)
|
|
.where(
|
|
EngagementDocumentVersion.id.in_(selected_ids),
|
|
EngagementDocumentVersion.client_id == int(client_id),
|
|
EngagementDocumentVersion.engagement_id.in_(allowed_engagements),
|
|
EngagementDocument.document_type == "BANK_STATEMENT",
|
|
EngagementDocument.is_deleted.is_(False),
|
|
)
|
|
).scalars().all()
|
|
)
|
|
by_id = {int(row.id): row for row in versions}
|
|
|
|
if set(selected_ids) != set(by_id):
|
|
raise PermissionError("One or more stored bank statements are not available to this user/client.")
|
|
|
|
prepared = []
|
|
provenance = []
|
|
pending = []
|
|
|
|
for order, version_id in enumerate(selected_ids, start=1):
|
|
version = by_id[version_id]
|
|
source = version_absolute_path(version)
|
|
if not source.exists():
|
|
source = _ready_cached_path(
|
|
db,
|
|
version_id=version.id,
|
|
user_id=user.id,
|
|
)
|
|
|
|
if source is None or not source.exists():
|
|
req = create_download_request_for_version(
|
|
db,
|
|
version=version,
|
|
user=user,
|
|
request=request,
|
|
)
|
|
if req is not None and req.id is None:
|
|
db.flush()
|
|
if req is not None and req.request_status == "ready":
|
|
cached = download_request_cache_path(req)
|
|
if cached is None or not cached.exists():
|
|
# A stale ready request must be made retrievable again.
|
|
req.request_status = "retry"
|
|
req.cached_relative_path = None
|
|
req.cached_hash_sha256 = None
|
|
req.last_error = "Cached retrieval copy was missing; queued again for branch Local Agent."
|
|
db.add(req)
|
|
db.flush()
|
|
if not req:
|
|
raise FileNotFoundError(
|
|
f"{version.original_filename}: the stored file is not currently available "
|
|
"on ERP or a registered branch storage node."
|
|
)
|
|
pending.append(
|
|
{
|
|
"version_id": int(version.id),
|
|
"filename": version.original_filename,
|
|
"request_id": int(req.id) if req.id else None,
|
|
"status": req.request_status,
|
|
}
|
|
)
|
|
continue
|
|
|
|
actual_hash = _sha256(source)
|
|
expected = (version.file_hash_sha256 or "").strip().lower()
|
|
if expected and actual_hash.lower() != expected:
|
|
raise ValueError(
|
|
f"{version.original_filename}: stored file hash verification failed."
|
|
)
|
|
|
|
safe_name = Path(version.original_filename or f"statement_{version.id}.pdf").name
|
|
if Path(safe_name).suffix.lower() != ".pdf":
|
|
raise ValueError(f"{safe_name}: only stored PDF bank statements can be reused.")
|
|
|
|
target = input_dir / f"stored_{order:03d}_v{version.id}_{safe_name}"
|
|
shutil.copy2(source, target)
|
|
prepared.append(target)
|
|
provenance.append(
|
|
{
|
|
"version_id": int(version.id),
|
|
"document_id": int(version.document_id),
|
|
"engagement_id": int(version.engagement_id),
|
|
"original_filename": version.original_filename,
|
|
"copied_filename": target.name,
|
|
"sha256": actual_hash,
|
|
}
|
|
)
|
|
|
|
db.flush()
|
|
return prepared, provenance, pending
|
|
|
|
|
|
def deduplicate_source_paths(paths: list[Path]):
|
|
unique = []
|
|
hashes = []
|
|
seen = set()
|
|
|
|
for path in paths:
|
|
digest = _sha256(path)
|
|
if digest in seen:
|
|
path.unlink(missing_ok=True)
|
|
continue
|
|
seen.add(digest)
|
|
unique.append(path)
|
|
hashes.append({"filename": path.name, "sha256": digest})
|
|
|
|
return unique, hashes
|