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