Add Phase 20 stored bank reuse and richer reconciliation
This commit is contained in:
@@ -0,0 +1,254 @@
|
||||
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
|
||||
Reference in New Issue
Block a user