Files
2026-08-24 18:58:58 +05:30

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