diff --git a/alembic/versions/20260824_bank_analyzer_stored_reuse_richer_reconciliation_phase20.py b/alembic/versions/20260824_bank_analyzer_stored_reuse_richer_reconciliation_phase20.py new file mode 100644 index 0000000..9796784 --- /dev/null +++ b/alembic/versions/20260824_bank_analyzer_stored_reuse_richer_reconciliation_phase20.py @@ -0,0 +1,147 @@ +"""Phase 20 Bank Analyzer stored-statement reuse + richer reconciliation. + +Revision ID: 20260824_bank_analyzer_reuse_p20 +Revises: 20260824_bank_reconciliation_p19 +""" +from alembic import op +import sqlalchemy as sa + + +revision = "20260824_bank_analyzer_reuse_p20" +down_revision = "20260824_bank_reconciliation_p19" +branch_labels = None +depends_on = None + + +def upgrade(): + op.add_column( + "bank_statement_analysis_jobs", + sa.Column("stored_source_versions_json", sa.Text(), nullable=True), + ) + op.add_column( + "bank_statement_analysis_jobs", + sa.Column("source_hashes_json", sa.Text(), nullable=True), + ) + + op.add_column( + "accounting_bank_reconciliation_runs", + sa.Column("account_number", sa.String(100), nullable=False, server_default=""), + ) + op.add_column( + "accounting_bank_reconciliation_runs", + sa.Column("date_tolerance_days", sa.Integer(), nullable=False, server_default="15"), + ) + op.create_index( + "ix_accounting_bank_reconciliation_runs_account_number", + "accounting_bank_reconciliation_runs", + ["account_number"], + ) + + op.add_column( + "accounting_bank_reconciliation_items", + sa.Column("resolution_status", sa.String(30), nullable=False, server_default="unresolved"), + ) + op.add_column( + "accounting_bank_reconciliation_items", + sa.Column("resolution_note", sa.Text(), nullable=False, server_default=""), + ) + op.add_column( + "accounting_bank_reconciliation_items", + sa.Column( + "resolved_by_user_id", + sa.Integer(), + sa.ForeignKey("users.id", ondelete="SET NULL"), + nullable=True, + ), + ) + op.add_column( + "accounting_bank_reconciliation_items", + sa.Column("resolved_at_utc", sa.DateTime(timezone=True), nullable=True), + ) + op.create_index( + "ix_accounting_bank_reconciliation_items_resolution_status", + "accounting_bank_reconciliation_items", + ["resolution_status"], + ) + + op.create_table( + "accounting_bank_ledger_mappings", + sa.Column("id", sa.Integer(), primary_key=True), + sa.Column( + "tenant_id", + sa.Integer(), + sa.ForeignKey("tenants.id", ondelete="CASCADE"), + nullable=False, + ), + sa.Column( + "client_id", + sa.Integer(), + sa.ForeignKey("clients.id", ondelete="CASCADE"), + nullable=False, + ), + sa.Column("account_number", sa.String(100), nullable=False), + sa.Column("bank_name", sa.String(160), nullable=False, server_default=""), + sa.Column("tally_guid", sa.String(120), nullable=False), + sa.Column("company_name", sa.String(255), nullable=False, server_default=""), + sa.Column("bank_ledger_name", sa.String(255), nullable=False), + sa.Column( + "created_by_user_id", + sa.Integer(), + sa.ForeignKey("users.id", ondelete="SET NULL"), + nullable=True, + ), + sa.Column( + "updated_by_user_id", + sa.Integer(), + sa.ForeignKey("users.id", ondelete="SET NULL"), + nullable=True, + ), + sa.Column( + "created_at_utc", + sa.DateTime(timezone=True), + nullable=False, + server_default=sa.func.now(), + ), + sa.Column( + "updated_at_utc", + sa.DateTime(timezone=True), + nullable=False, + server_default=sa.func.now(), + ), + sa.UniqueConstraint( + "tenant_id", + "client_id", + "account_number", + "tally_guid", + name="uq_accounting_bank_ledger_mapping", + ), + ) + for name in ("tenant_id", "client_id", "account_number", "tally_guid"): + op.create_index( + f"ix_accounting_bank_ledger_mappings_{name}", + "accounting_bank_ledger_mappings", + [name], + ) + + +def downgrade(): + op.drop_table("accounting_bank_ledger_mappings") + + op.drop_index( + "ix_accounting_bank_reconciliation_items_resolution_status", + table_name="accounting_bank_reconciliation_items", + ) + op.drop_column("accounting_bank_reconciliation_items", "resolved_at_utc") + op.drop_column("accounting_bank_reconciliation_items", "resolved_by_user_id") + op.drop_column("accounting_bank_reconciliation_items", "resolution_note") + op.drop_column("accounting_bank_reconciliation_items", "resolution_status") + + op.drop_index( + "ix_accounting_bank_reconciliation_runs_account_number", + table_name="accounting_bank_reconciliation_runs", + ) + op.drop_column("accounting_bank_reconciliation_runs", "date_tolerance_days") + op.drop_column("accounting_bank_reconciliation_runs", "account_number") + + op.drop_column("bank_statement_analysis_jobs", "source_hashes_json") + op.drop_column("bank_statement_analysis_jobs", "stored_source_versions_json") diff --git a/app/modules/accounting/bank_reconciliation_models.py b/app/modules/accounting/bank_reconciliation_models.py index 858276a..c345cec 100644 --- a/app/modules/accounting/bank_reconciliation_models.py +++ b/app/modules/accounting/bank_reconciliation_models.py @@ -2,7 +2,7 @@ from __future__ import annotations from datetime import datetime, timezone -from sqlalchemy import DateTime, Float, ForeignKey, Integer, String, Text +from sqlalchemy import DateTime, Float, ForeignKey, Integer, String, Text, UniqueConstraint from sqlalchemy.orm import Mapped, mapped_column from app.core.db.common import CommonBase @@ -22,8 +22,10 @@ class BankReconciliationRun(CommonBase): tally_guid: Mapped[str] = mapped_column(String(120), nullable=False, index=True) company_name: Mapped[str] = mapped_column(String(255), nullable=False, default="") bank_ledger_name: Mapped[str] = mapped_column(String(255), nullable=False) + account_number: Mapped[str] = mapped_column(String(100), nullable=False, default="", index=True) date_from: Mapped[str] = mapped_column(String(10), nullable=False) date_to: Mapped[str] = mapped_column(String(10), nullable=False) + date_tolerance_days: Mapped[int] = mapped_column(Integer, nullable=False, default=15) workstation_agent_id: Mapped[int] = mapped_column(ForeignKey("erp_workstation_agents.id", ondelete="RESTRICT"), nullable=False, index=True) agent_job_id: Mapped[int | None] = mapped_column(ForeignKey("erp_agent_jobs.id", ondelete="SET NULL"), nullable=True, index=True) @@ -62,3 +64,34 @@ class BankReconciliationItem(CommonBase): tally_narration: Mapped[str] = mapped_column(Text, nullable=False, default="") tally_amount: Mapped[float] = mapped_column(Float, nullable=False, default=0) tally_direction: Mapped[str] = mapped_column(String(10), nullable=False, default="") + + resolution_status: Mapped[str] = mapped_column(String(30), nullable=False, default="unresolved", index=True) + resolution_note: Mapped[str] = mapped_column(Text, nullable=False, default="") + resolved_by_user_id: Mapped[int | None] = mapped_column(ForeignKey("users.id", ondelete="SET NULL"), nullable=True) + resolved_at_utc: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + + +class AccountingBankLedgerMapping(CommonBase): + __tablename__ = "accounting_bank_ledger_mappings" + __table_args__ = ( + UniqueConstraint( + "tenant_id", + "client_id", + "account_number", + "tally_guid", + name="uq_accounting_bank_ledger_mapping", + ), + ) + + id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True) + tenant_id: Mapped[int] = mapped_column(ForeignKey("tenants.id", ondelete="CASCADE"), nullable=False, index=True) + client_id: Mapped[int] = mapped_column(ForeignKey("clients.id", ondelete="CASCADE"), nullable=False, index=True) + account_number: Mapped[str] = mapped_column(String(100), nullable=False, index=True) + bank_name: Mapped[str] = mapped_column(String(160), nullable=False, default="") + tally_guid: Mapped[str] = mapped_column(String(120), nullable=False, index=True) + company_name: Mapped[str] = mapped_column(String(255), nullable=False, default="") + bank_ledger_name: Mapped[str] = mapped_column(String(255), nullable=False) + created_by_user_id: Mapped[int | None] = mapped_column(ForeignKey("users.id", ondelete="SET NULL"), nullable=True) + updated_by_user_id: Mapped[int | None] = mapped_column(ForeignKey("users.id", ondelete="SET NULL"), nullable=True) + created_at_utc: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=_utcnow) + updated_at_utc: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=_utcnow) diff --git a/app/modules/accounting/bank_reconciliation_service.py b/app/modules/accounting/bank_reconciliation_service.py index d0ae477..6b0b837 100644 --- a/app/modules/accounting/bank_reconciliation_service.py +++ b/app/modules/accounting/bank_reconciliation_service.py @@ -9,6 +9,7 @@ from sqlalchemy import delete, func, select from app.modules.accounting.bank_models import AccountingBankTransaction from app.modules.accounting.bank_reconciliation_models import ( + AccountingBankLedgerMapping, BankReconciliationItem, BankReconciliationRun, ) @@ -109,22 +110,131 @@ def completed_client_jobs(db, *, tenant_id: int, client_id: int, limit: int = 50 ) -def _job_period(db, *, tenant_id: int, client_id: int, job_id: str): - values = list( +def source_accounts(db, *, tenant_id: int, client_id: int, job_id: str): + rows = list( db.execute( - select(AccountingBankTransaction.transaction_date) + select( + AccountingBankTransaction.account_number, + AccountingBankTransaction.bank_name, + func.count(AccountingBankTransaction.id), + ) .where( AccountingBankTransaction.tenant_id == int(tenant_id), AccountingBankTransaction.client_id == int(client_id), AccountingBankTransaction.source_job_id == job_id, ) - .order_by(AccountingBankTransaction.transaction_date) + .group_by( + AccountingBankTransaction.account_number, + AccountingBankTransaction.bank_name, + ) + .order_by( + AccountingBankTransaction.bank_name, + AccountingBankTransaction.account_number, + ) + ).all() + ) + return [ + { + "account_number": _s(account), + "bank_name": _s(bank), + "transaction_count": int(count or 0), + } + for account, bank, count in rows + if _s(account) + ] + + +def list_bank_mappings(db, *, tenant_id: int, client_id: int, tally_guid: str = ""): + stmt = select(AccountingBankLedgerMapping).where( + AccountingBankLedgerMapping.tenant_id == int(tenant_id), + AccountingBankLedgerMapping.client_id == int(client_id), + ) + if _s(tally_guid): + stmt = stmt.where(AccountingBankLedgerMapping.tally_guid == _s(tally_guid)) + return list( + db.execute( + stmt.order_by( + AccountingBankLedgerMapping.bank_name, + AccountingBankLedgerMapping.account_number, + ) + ).scalars().all() + ) + + +def save_bank_mapping( + db, + *, + tenant_id: int, + client_id: int, + account_number: str, + bank_name: str, + tally_guid: str, + company_name: str, + bank_ledger_name: str, + user_id: int, +): + account_number = _s(account_number) + if not account_number: + raise ValueError("Bank account number is required.") + + allowed = { + row.name + for row in visible_bank_ledgers( + db, + tenant_id=tenant_id, + client_id=client_id, + tally_guid=tally_guid, + ) + } + if bank_ledger_name not in allowed: + raise ValueError("Select a valid Tally Bank ledger from Chart of Accounts.") + + mapping = db.execute( + select(AccountingBankLedgerMapping).where( + AccountingBankLedgerMapping.tenant_id == int(tenant_id), + AccountingBankLedgerMapping.client_id == int(client_id), + AccountingBankLedgerMapping.account_number == account_number, + AccountingBankLedgerMapping.tally_guid == _s(tally_guid), + ) + ).scalar_one_or_none() + + if mapping is None: + mapping = AccountingBankLedgerMapping( + tenant_id=int(tenant_id), + client_id=int(client_id), + account_number=account_number, + tally_guid=_s(tally_guid), + created_by_user_id=int(user_id), + ) + + mapping.bank_name = _s(bank_name) + mapping.company_name = _s(company_name) + mapping.bank_ledger_name = _s(bank_ledger_name) + mapping.updated_by_user_id = int(user_id) + mapping.updated_at_utc = _utcnow() + db.add(mapping) + db.commit() + db.refresh(mapping) + return mapping + + +def _job_period(db, *, tenant_id: int, client_id: int, job_id: str, account_number: str = ""): + stmt = select(AccountingBankTransaction.transaction_date).where( + AccountingBankTransaction.tenant_id == int(tenant_id), + AccountingBankTransaction.client_id == int(client_id), + AccountingBankTransaction.source_job_id == job_id, + ) + if _s(account_number): + stmt = stmt.where(AccountingBankTransaction.account_number == _s(account_number)) + values = list( + db.execute( + stmt.order_by(AccountingBankTransaction.transaction_date) ).scalars().all() ) values = [v for v in values if _date_obj(v)] if not values: raise ValueError( - "No imported bank transactions were found for this Bank Analyzer job. " + "No imported bank transactions were found for this Bank Analyzer job/account. " "Use a Bank Reconciliation purpose job or import the completed job into Accounting first." ) return values[0][:10], values[-1][:10] @@ -141,6 +251,8 @@ def queue_reconciliation( bank_ledger_name: str, workstation_id: int, user_id: int, + account_number: str = "", + date_tolerance_days: int = 15, ): source_job = db.get(BankStatementAnalysisJob, source_job_id) if ( @@ -157,11 +269,15 @@ def queue_reconciliation( "The selected Bank Analyzer job has not been imported successfully into the Accounting bank queue." ) + account_number = _s(account_number) + date_tolerance_days = max(0, min(90, int(date_tolerance_days or 15))) + date_from, date_to = _job_period( db, tenant_id=tenant_id, client_id=client_id, job_id=source_job_id, + account_number=account_number, ) allowed = { @@ -192,8 +308,10 @@ def queue_reconciliation( tally_guid=_s(tally_guid), company_name=_s(company_name), bank_ledger_name=bank_ledger_name, + account_number=account_number, date_from=date_from, date_to=date_to, + date_tolerance_days=date_tolerance_days, workstation_agent_id=ws.id, status="queued", created_by_user_id=int(user_id), @@ -212,11 +330,12 @@ def queue_reconciliation( "tally_guid": _s(tally_guid), "company_name": _s(company_name), "bank_ledger_name": bank_ledger_name, + "account_number": account_number, "date_from": date_from, "date_to": date_to, "reconciliation_run_id": run.id, }, - idempotency_key=f"bank-recon:{tenant_id}:{client_id}:{source_job_id}:{tally_guid}:{bank_ledger_name}:{run.id}", + idempotency_key=f"bank-recon:{tenant_id}:{client_id}:{source_job_id}:{account_number}:{tally_guid}:{bank_ledger_name}:{run.id}", priority=8, max_attempts=2, created_by_user_id=user_id, @@ -227,6 +346,74 @@ def queue_reconciliation( return run +def queue_all_mapped_reconciliations( + db, + *, + tenant_id: int, + client_id: int, + source_job_id: str, + tally_guid: str, + company_name: str, + workstation_id: int, + user_id: int, + date_tolerance_days: int = 15, +): + accounts = source_accounts( + db, + tenant_id=tenant_id, + client_id=client_id, + job_id=source_job_id, + ) + if not accounts: + raise ValueError("No bank accounts were imported from the selected Bank Analyzer job.") + + mappings = { + row.account_number: row + for row in list_bank_mappings( + db, + tenant_id=tenant_id, + client_id=client_id, + tally_guid=tally_guid, + ) + } + + missing = [ + row + for row in accounts + if row["account_number"] not in mappings + ] + if missing: + labels = ", ".join( + f"{row['bank_name']} {row['account_number']}" + for row in missing + ) + raise ValueError( + "Save Tally bank-ledger mapping for every bank account before " + f"running all accounts together. Missing: {labels}" + ) + + runs = [] + for account in accounts: + mapping = mappings[account["account_number"]] + runs.append( + queue_reconciliation( + db, + tenant_id=tenant_id, + client_id=client_id, + source_job_id=source_job_id, + tally_guid=tally_guid, + company_name=company_name, + bank_ledger_name=mapping.bank_ledger_name, + workstation_id=workstation_id, + user_id=user_id, + account_number=account["account_number"], + date_tolerance_days=date_tolerance_days, + ) + ) + return runs + + + def _tally_side(voucher: dict, bank_ledger_name: str): bank_key = bank_ledger_name.casefold() entries = list(voucher.get("ledger_entries") or []) @@ -256,7 +443,7 @@ def _tally_side(voucher: dict, bank_ledger_name: str): } -def _pair_score(bank: AccountingBankTransaction, tally: dict): +def _pair_score(bank: AccountingBankTransaction, tally: dict, tolerance_days: int = 15): if round(float(bank.amount or 0), 2) != round(float(tally["amount"] or 0), 2): return 0, "" if _s(bank.direction).upper() != _s(tally["direction"]).upper(): @@ -268,7 +455,8 @@ def _pair_score(bank: AccountingBankTransaction, tally: dict): return 0, "" gap = abs((bd - td).days) - if gap > 7: + tolerance_days = max(0, min(90, int(tolerance_days or 15))) + if gap > tolerance_days: return 0, "" bank_ref = _norm_ref(bank.transfer_reference or bank.reference_no) @@ -291,9 +479,14 @@ def _pair_score(bank: AccountingBankTransaction, tally: dict): elif gap <= 4: score += 7 reasons.append(f"{gap}-day timing difference") - else: + elif gap <= 7: score += 2 reasons.append(f"{gap}-day timing difference") + else: + # Long timing differences remain match candidates only because amount + # and bank direction are exact. They are surfaced separately rather + # than silently treated as a normal probable match. + reasons.append(f"{gap}-day timing difference") if ref_exact: score += 18 @@ -319,6 +512,11 @@ def build_reconciliation(db, run: BankReconciliationRun, vouchers: list[dict]): AccountingBankTransaction.tenant_id == run.tenant_id, AccountingBankTransaction.client_id == run.client_id, AccountingBankTransaction.source_job_id == run.source_job_id, + *( + [AccountingBankTransaction.account_number == run.account_number] + if _s(run.account_number) + else [] + ), ) .order_by( AccountingBankTransaction.transaction_date, @@ -339,14 +537,14 @@ def build_reconciliation(db, run: BankReconciliationRun, vouchers: list[dict]): for bank in bank_rows: scored = [] for index, tally in enumerate(tally_rows): - score, reason = _pair_score(bank, tally) + score, reason = _pair_score(bank, tally, run.date_tolerance_days) if score: scored.append((score, index, reason)) scored.sort(key=lambda item: (-item[0], item[1])) candidates[bank.id] = scored used_tally = set() - exact = probable = bank_only = duplicates = 0 + exact = probable = timing = bank_only = duplicates = 0 for bank in bank_rows: options = [ @@ -372,10 +570,17 @@ def build_reconciliation(db, run: BankReconciliationRun, vouchers: list[dict]): used_tally.add(tally_index) tally = tally_rows[tally_index] confidence = top_score - status = "matched" if top_score >= 90 else "probable_match" - if status == "matched": + bank_date = _date_obj(bank.transaction_date) + tally_date = _date_obj(tally["date"]) + gap = abs((bank_date - tally_date).days) if bank_date and tally_date else 999 + if top_score >= 90: + status = "matched" exact += 1 + elif gap > 7: + status = "timing_difference" + timing += 1 else: + status = "probable_match" probable += 1 item = BankReconciliationItem( @@ -428,18 +633,74 @@ def build_reconciliation(db, run: BankReconciliationRun, vouchers: list[dict]): ) ) - run.summary_json = json.dumps( - { - "bank_transactions": len(bank_rows), - "tally_bank_vouchers": len(tally_rows), - "matched": exact, - "probable_match": probable, - "bank_only": bank_only, - "books_only": books_only, - "duplicate_candidate": duplicates, - }, - ensure_ascii=False, + bank_only_rows = [ + row for row in db.execute( + select(BankReconciliationItem).where( + BankReconciliationItem.run_id == run.id, + BankReconciliationItem.match_status == "bank_only", + ) + ).scalars().all() + ] + books_only_rows = [ + row for row in db.execute( + select(BankReconciliationItem).where( + BankReconciliationItem.run_id == run.id, + BankReconciliationItem.match_status == "books_only", + ) + ).scalars().all() + ] + + def _direction_totals(rows, prefix): + result = { + f"{prefix}_debit_amount": 0.0, + f"{prefix}_credit_amount": 0.0, + } + for row in rows: + direction = ( + row.bank_direction if prefix == "bank_only" + else row.tally_direction + ) + amount = ( + row.bank_amount if prefix == "bank_only" + else row.tally_amount + ) + key = ( + f"{prefix}_debit_amount" + if _s(direction).upper() == "DEBIT" + else f"{prefix}_credit_amount" + ) + result[key] = round(result[key] + float(amount or 0), 2) + return result + + summary = { + "bank_transactions": len(bank_rows), + "tally_bank_vouchers": len(tally_rows), + "matched": exact, + "probable_match": probable, + "timing_difference": timing, + "bank_only": bank_only, + "books_only": books_only, + "duplicate_candidate": duplicates, + "account_number": run.account_number, + "bank_ledger_name": run.bank_ledger_name, + "date_tolerance_days": int(run.date_tolerance_days or 15), + } + summary.update(_direction_totals(bank_only_rows, "bank_only")) + summary.update(_direction_totals(books_only_rows, "books_only")) + summary["bank_only_net"] = round( + summary["bank_only_credit_amount"] - summary["bank_only_debit_amount"], + 2, ) + summary["books_only_net"] = round( + summary["books_only_credit_amount"] - summary["books_only_debit_amount"], + 2, + ) + summary["unreconciled_net_difference"] = round( + summary["bank_only_net"] - summary["books_only_net"], + 2, + ) + + run.summary_json = json.dumps(summary, ensure_ascii=False) run.status = "completed" run.completed_at_utc = _utcnow() run.last_error = "" @@ -478,6 +739,64 @@ def sync_run(db, run: BankReconciliationRun): return run +def resolve_reconciliation_item( + db, + *, + run_id: int, + item_id: int, + action: str, + note: str, + user_id: int, +): + item = db.execute( + select(BankReconciliationItem).where( + BankReconciliationItem.id == int(item_id), + BankReconciliationItem.run_id == int(run_id), + ) + ).scalar_one_or_none() + if not item: + raise ValueError("Reconciliation item was not found.") + + action = _s(action) + allowed = { + "confirm_match", + "confirm_timing", + "reject_match_bank_only", + "confirm_bank_only", + "confirm_books_only", + "needs_follow_up", + "reopen", + } + if action not in allowed: + raise ValueError("Unsupported reconciliation resolution.") + + item.resolution_status = action + item.resolution_note = _s(note) + item.resolved_by_user_id = int(user_id) + item.resolved_at_utc = _utcnow() + + if item.bank_transaction_id: + tx = db.get(AccountingBankTransaction, int(item.bank_transaction_id)) + if tx: + if action == "confirm_match": + tx.reconciliation_status = "matched" + elif action == "confirm_timing": + tx.reconciliation_status = "timing_difference_confirmed" + elif action in {"reject_match_bank_only", "confirm_bank_only"}: + tx.reconciliation_status = "bank_only" + elif action == "needs_follow_up": + tx.reconciliation_status = "needs_review" + elif action == "reopen": + tx.reconciliation_status = item.match_status + tx.last_reconciliation_run_id = int(run_id) + db.add(tx) + + db.add(item) + db.commit() + return item + + + def list_runs(db, *, tenant_id: int, client_id: int, limit: int = 30): rows = list( db.execute( diff --git a/app/modules/accounting/bank_reconciliation_ui.py b/app/modules/accounting/bank_reconciliation_ui.py index 286dfe1..e4732aa 100644 --- a/app/modules/accounting/bank_reconciliation_ui.py +++ b/app/modules/accounting/bank_reconciliation_ui.py @@ -11,9 +11,14 @@ from app.core.security.csrf import get_or_create_csrf_token, validate_csrf from app.core.templating import templates from app.modules.accounting.bank_reconciliation_service import ( completed_client_jobs, + list_bank_mappings, list_runs, + queue_all_mapped_reconciliations, queue_reconciliation, + resolve_reconciliation_item, run_items, + save_bank_mapping, + source_accounts, visible_bank_ledgers, visible_workstations, ) @@ -91,6 +96,9 @@ def page( runs = [] items = [] selected_run = None + accounts = [] + mappings = [] + run_summaries = {} if selected: jobs = completed_client_jobs( @@ -99,6 +107,13 @@ def page( client_id=selected.id, limit=50, ) + if job_id: + accounts = source_accounts( + db, + tenant_id=scope.tenant_id, + client_id=selected.id, + job_id=job_id, + ) companies = _company_options(db, scope.tenant_id, selected.id) if not tally_guid and len(companies) == 1: tally_guid = companies[0]["guid"] @@ -109,6 +124,12 @@ def page( client_id=selected.id, tally_guid=tally_guid, ) + mappings = list_bank_mappings( + db, + tenant_id=scope.tenant_id, + client_id=selected.id, + tally_guid=tally_guid, + ) workstations = visible_workstations( db, tenant_id=scope.tenant_id, @@ -136,6 +157,12 @@ def page( status=status, ) + for row in runs: + try: + run_summaries[int(row.id)] = json.loads(row.summary_json or "{}") + except Exception: + run_summaries[int(row.id)] = {} + return templates.TemplateResponse( "modules/accounting/templates/accounting/bank_reconciliation.html", { @@ -157,6 +184,9 @@ def page( "selected_run": selected_run, "items": items, "status_filter": status, + "source_accounts": accounts, + "bank_mappings": mappings, + "run_summaries": run_summaries, "message": message, "error": error, }, @@ -173,6 +203,8 @@ def create_run( tally_guid: str = Form(...), company_name: str = Form(...), bank_ledger_name: str = Form(...), + account_number: str = Form(""), + date_tolerance_days: int = Form(15), workstation_id: int = Form(...), csrf_token: str = Form(...), ): @@ -197,6 +229,8 @@ def create_run( bank_ledger_name=bank_ledger_name, workstation_id=workstation_id, user_id=user.id, + account_number=account_number, + date_tolerance_days=date_tolerance_days, ) return _go( client.id, @@ -216,3 +250,153 @@ def create_run( ) finally: db.close() + + +@router.post("/mapping") +def save_mapping( + request: Request, + client_id: int = Form(...), + source_job_id: str = Form(""), + tally_guid: str = Form(...), + company_name: str = Form(...), + account_number: str = Form(...), + bank_name: str = Form(""), + bank_ledger_name: str = Form(...), + csrf_token: str = Form(...), +): + validate_csrf(request, csrf_token) + db = CommonSessionLocal() + try: + user, denied = _require_partner(request, db, "accounting.learning.manage") + if denied: + return denied + client, _clients, scope = _find_visible_client(db, request, user, client_id) + if not client: + return _go(error="Client is not visible.") + + save_bank_mapping( + db, + tenant_id=scope.tenant_id, + client_id=client.id, + account_number=account_number, + bank_name=bank_name, + tally_guid=tally_guid, + company_name=company_name, + bank_ledger_name=bank_ledger_name, + user_id=user.id, + ) + return _go( + client.id, + job_id=source_job_id, + tally_guid=tally_guid, + message=f"Bank account {account_number} mapped to Tally ledger '{bank_ledger_name}'.", + ) + except Exception as exc: + db.rollback() + return _go(client_id, job_id=source_job_id, tally_guid=tally_guid, error=str(exc)) + finally: + db.close() + + +@router.post("/run-all") +def create_all_runs( + request: Request, + client_id: int = Form(...), + source_job_id: str = Form(...), + tally_guid: str = Form(...), + company_name: str = Form(...), + date_tolerance_days: int = Form(15), + workstation_id: int = Form(...), + csrf_token: str = Form(...), +): + validate_csrf(request, csrf_token) + db = CommonSessionLocal() + try: + user, denied = _require_partner(request, db, "accounting.learning.manage") + if denied: + return denied + client, _clients, scope = _find_visible_client(db, request, user, client_id) + if not client: + return _go(error="Client is not visible.") + + runs = queue_all_mapped_reconciliations( + db, + tenant_id=scope.tenant_id, + client_id=client.id, + source_job_id=source_job_id, + tally_guid=tally_guid, + company_name=company_name, + workstation_id=workstation_id, + user_id=user.id, + date_tolerance_days=date_tolerance_days, + ) + return _go( + client.id, + job_id=source_job_id, + tally_guid=tally_guid, + run_id=(runs[0].id if runs else None), + message=f"Queued {len(runs)} reconciliation run(s), one for each mapped bank account.", + ) + except Exception as exc: + db.rollback() + return _go(client_id, job_id=source_job_id, tally_guid=tally_guid, error=str(exc)) + finally: + db.close() + + +@router.post("/runs/{run_id}/items/{item_id}/resolve") +def resolve_item( + request: Request, + run_id: int, + item_id: int, + client_id: int = Form(...), + tally_guid: str = Form(""), + action: str = Form(...), + note: str = Form(""), + csrf_token: str = Form(...), +): + validate_csrf(request, csrf_token) + db = CommonSessionLocal() + try: + user, denied = _require_partner(request, db, "accounting.learning.manage") + if denied: + return denied + client, _clients, scope = _find_visible_client(db, request, user, client_id) + if not client: + return _go(error="Client is not visible.") + + run = next( + ( + row + for row in list_runs( + db, + tenant_id=scope.tenant_id, + client_id=client.id, + limit=100, + ) + if int(row.id) == int(run_id) + ), + None, + ) + if not run: + raise ValueError("Reconciliation run was not found.") + + resolve_reconciliation_item( + db, + run_id=run.id, + item_id=item_id, + action=action, + note=note, + user_id=user.id, + ) + return _go( + client.id, + tally_guid=(tally_guid or run.tally_guid), + run_id=run.id, + message="Reconciliation review decision saved.", + ) + except Exception as exc: + db.rollback() + return _go(client_id, tally_guid=tally_guid, run_id=run_id, error=str(exc)) + finally: + db.close() diff --git a/app/modules/accounting/bank_service.py b/app/modules/accounting/bank_service.py index 60eb058..328ed88 100644 --- a/app/modules/accounting/bank_service.py +++ b/app/modules/accounting/bank_service.py @@ -299,6 +299,9 @@ def _ensure_reviewed(tx): "matched", "probable_match", "duplicate_candidate", + "timing_difference", + "timing_difference_confirmed", + "needs_review", }: raise ValueError( "Bank Reconciliation indicates that this transaction already has, " diff --git a/app/modules/accounting/templates/accounting/bank_reconciliation.html b/app/modules/accounting/templates/accounting/bank_reconciliation.html index 0741e54..de1b542 100644 --- a/app/modules/accounting/templates/accounting/bank_reconciliation.html +++ b/app/modules/accounting/templates/accounting/bank_reconciliation.html @@ -6,7 +6,7 @@
Accounting · Bank Reconciliation
- Compare the normalized transactions from a completed client-bound Bank Analyzer job with one selected Tally Bank ledger. This screen is read-only against Tally; unmatched bank rows can later be handled through the existing Accounting Bank Queue. + Compare normalized Bank Analyzer transactions with Tally Bank ledgers. Phase 20 remembers each client bank-account ↔ Tally-ledger mapping, can reconcile every mapped account in one action, supports wider cheque/clearing timing windows, and keeps reviewer resolutions. Tally access remains read-only.
The Local Agent only reads vouchers from the selected Tally company/date range and filters them to the selected Bank ledger. No accounting entry is posted by this action.
+Choose the completed client-bound analysis whose imported bank transactions you want to reconcile. The job can contain any number of bank accounts.
+Mappings are remembered client-wise and company-wise. Once all accounts are mapped, future reconciliation runs do not require selecting the bank ledger again.
Use a wider timing tolerance for cheque issue/presentation or deposit/clearing differences. Amount and bank direction must still match exactly before the engine considers a Tally voucher as a candidate.
+| Run | Bank Ledger | Period | Status | Result | ||||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
| Run | Bank / Account | Period | Status | Result | ||||||||||
| #{{ run.id }} | -{{ run.bank_ledger_name }} | +{{ run.bank_ledger_name }} {{ run.account_number or "All imported accounts" }} |
{{ run.date_from }} → {{ run.date_to }} | {{ run.status|replace('_',' ')|title }}{% if run.last_error %} {{ run.last_error }} {% endif %} |
+ {% set s = run_summaries.get(run.id, {}) %}
{% if run.status=='completed' %}
- {% set s = run.summary_json|from_json if false else none %}
- Open run to view matched / probable / bank-only / books-only detail.
+ Matched {{ s.get('matched',0) }} · Probable {{ s.get('probable_match',0) }} · Timing {{ s.get('timing_difference',0) }}
+ Bank only {{ s.get('bank_only',0) }} · Books only {{ s.get('books_only',0) }} · Duplicate {{ s.get('duplicate_candidate',0) }}
+ Unreconciled net difference ₹{{ '%.2f'|format(s.get('unreconciled_net_difference',0)) }}
{% else %}—{% endif %}
|
Open | @@ -144,7 +232,7 @@ @@ -158,6 +246,15 @@ {% elif selected_run.status=='failed' %}
| Result | Bank Statement | Tally | Why | Action | {{ item.match_reason }} | -+ |
+ {% if item.resolution_status and item.resolution_status != 'unresolved' %}
+ Reviewed: {{ item.resolution_status|replace('_',' ')|title }}{% if item.resolution_note %}
+ {% endif %}
+ {{ item.resolution_note }} {% endif %} |
diff --git a/app/modules/bank_statement_analyzer/client_context.py b/app/modules/bank_statement_analyzer/client_context.py
index 32cc872..f21a50c 100644
--- a/app/modules/bank_statement_analyzer/client_context.py
+++ b/app/modules/bank_statement_analyzer/client_context.py
@@ -365,12 +365,37 @@ def archive_analysis_to_engagement(
return {"status": "not_requested", "documents": []}
archived = []
+
+ reused = []
+ try:
+ reused = json.loads(getattr(job, "stored_source_versions_json", None) or "[]")
+ except Exception:
+ reused = []
+ reused_by_copied_name = {
+ _text(row.get("copied_filename")): row
+ for row in reused
+ if isinstance(row, dict) and _text(row.get("copied_filename"))
+ }
+
meta_by_source = {
Path(_text(getattr(meta, "source_file", ""))).name: meta
for meta in metas
}
for source in source_paths:
+ reused_row = reused_by_copied_name.get(source.name)
+ if reused_row:
+ archived.append(
+ {
+ "document_id": int(reused_row.get("document_id") or 0),
+ "version_id": int(reused_row.get("version_id") or 0),
+ "filename": reused_row.get("original_filename") or source.name,
+ "type": "BANK_STATEMENT",
+ "reused_existing_document": True,
+ }
+ )
+ continue
+
upload = _upload_proxy(source, "application/pdf")
try:
meta = meta_by_source.get(source.name)
diff --git a/app/modules/bank_statement_analyzer/models.py b/app/modules/bank_statement_analyzer/models.py
index d1be60f..ef99058 100644
--- a/app/modules/bank_statement_analyzer/models.py
+++ b/app/modules/bank_statement_analyzer/models.py
@@ -39,6 +39,11 @@ class BankStatementAnalysisJob(CommonBase):
accounting_import_status: Mapped[str] = mapped_column(String(30), nullable=False, default="not_requested", index=True)
accounting_import_json: Mapped[str | None] = mapped_column(Text, nullable=True)
+ # Phase 20: stored-statement reuse. These fields only record provenance;
+ # the physical files remain in the existing Engagement Documents storage.
+ stored_source_versions_json: Mapped[str | None] = mapped_column(Text, nullable=True)
+ source_hashes_json: Mapped[str | None] = mapped_column(Text, nullable=True)
+
status: Mapped[str] = mapped_column(String(20), nullable=False, default="queued", index=True)
progress_percent: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
file_count: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
diff --git a/app/modules/bank_statement_analyzer/service.py b/app/modules/bank_statement_analyzer/service.py
index 23cc46b..4929078 100644
--- a/app/modules/bank_statement_analyzer/service.py
+++ b/app/modules/bank_statement_analyzer/service.py
@@ -68,6 +68,80 @@ def _segment(value: object, default: str = "NA") -> str:
+def _statement_period(value):
+ text_value = str(value or "").strip()
+ if not text_value:
+ return None
+ for fmt in (
+ "%Y-%m-%d",
+ "%d/%m/%Y",
+ "%d-%m-%Y",
+ "%d.%m.%Y",
+ "%d-%b-%Y",
+ "%d %b %Y",
+ "%d %B %Y",
+ "%Y/%m/%d",
+ ):
+ try:
+ return datetime.strptime(text_value[:10], fmt).date()
+ except Exception:
+ pass
+ try:
+ return datetime.fromisoformat(text_value.replace("Z", "+00:00")).date()
+ except Exception:
+ return None
+
+
+def _statement_overlap_warnings(metas: list) -> list[dict]:
+ warnings = []
+ groups = {}
+ for meta in metas:
+ account = str(getattr(meta, "account_number", "") or "").strip()
+ bank = str(getattr(meta, "bank_name", "") or "").strip()
+ start = _statement_period(getattr(meta, "period_from", ""))
+ end = _statement_period(getattr(meta, "period_to", ""))
+ if not account or not start or not end:
+ continue
+ if end < start:
+ start, end = end, start
+ groups.setdefault((bank.casefold(), account), []).append(
+ {
+ "bank_name": bank,
+ "account_number": account,
+ "period_from": start,
+ "period_to": end,
+ "source_file": Path(str(getattr(meta, "source_file", "") or "")).name,
+ }
+ )
+
+ for (_bank_key, _account), rows in groups.items():
+ rows.sort(key=lambda row: (row["period_from"], row["period_to"]))
+ for idx, left in enumerate(rows):
+ for right in rows[idx + 1:]:
+ if right["period_from"] > left["period_to"]:
+ break
+ overlap_from = max(left["period_from"], right["period_from"])
+ overlap_to = min(left["period_to"], right["period_to"])
+ if overlap_from <= overlap_to:
+ warnings.append(
+ {
+ "bank_name": left["bank_name"],
+ "account_number": left["account_number"],
+ "first_file": left["source_file"],
+ "second_file": right["source_file"],
+ "overlap_from": overlap_from.isoformat(),
+ "overlap_to": overlap_to.isoformat(),
+ "message": (
+ f"Overlapping statement periods for account "
+ f"{left['account_number']}: {left['source_file']} and "
+ f"{right['source_file']} overlap from "
+ f"{overlap_from.isoformat()} to {overlap_to.isoformat()}."
+ ),
+ }
+ )
+ return warnings
+
+
def _workbook_filename(metas: list) -> str:
"""Build a safe BankName_ClientName.xlsx filename for every bank."""
bank_names = [str(getattr(meta, "bank_name", "") or "").strip() for meta in metas]
@@ -162,7 +236,7 @@ def pending_count_for_user(user_id: int) -> int:
db.close()
-def enqueue_job(*, user, roles: Iterable[str], job_id: str, paths: list[Path], job_dir: Path, bank_selection: str, financial_year: str, customer_override: str, account_override: str, classification_enabled: bool, client_id: int | None = None, engagement_id: int | None = None, ownership_confirmation: bool = False, purpose: str = "analyze_only") -> BankStatementAnalysisJob:
+def enqueue_job(*, user, roles: Iterable[str], job_id: str, paths: list[Path], job_dir: Path, bank_selection: str, financial_year: str, customer_override: str, account_override: str, classification_enabled: bool, client_id: int | None = None, engagement_id: int | None = None, ownership_confirmation: bool = False, purpose: str = "analyze_only", stored_source_versions: list[dict] | None = None, source_hashes: list[dict] | None = None) -> BankStatementAnalysisJob:
if pending_count_for_user(int(user.id)) >= MAX_PENDING_PER_USER:
shutil.rmtree(job_dir, ignore_errors=True)
raise ValueError("You already have three queued or processing analyses. Please wait for one to complete before submitting another.")
@@ -192,6 +266,16 @@ def enqueue_job(*, user, roles: Iterable[str], job_id: str, paths: list[Path], j
if purpose in {"accounting_entries", "bank_reconciliation"}
else "not_requested"
),
+ stored_source_versions_json=(
+ json.dumps(stored_source_versions, ensure_ascii=False)
+ if stored_source_versions
+ else None
+ ),
+ source_hashes_json=(
+ json.dumps(source_hashes, ensure_ascii=False)
+ if source_hashes
+ else None
+ ),
status="queued",
progress_percent=0,
file_count=len(paths),
@@ -328,6 +412,13 @@ def _process_job(job_id: str) -> None:
"engagement_id": job.engagement_id,
"ownership_status": job.ownership_status,
"engagement_archive_status": job.engagement_archive_status,
+ "stored_statement_reuse_count": len(
+ json.loads(job.stored_source_versions_json or "[]")
+ ) if getattr(job, "stored_source_versions_json", None) else 0,
+ "source_hash_count": len(
+ json.loads(job.source_hashes_json or "[]")
+ ) if getattr(job, "source_hashes_json", None) else 0,
+ "statement_overlap_warnings": _statement_overlap_warnings(metas),
}
# Keep the original uploaded statements until the job expiry time.
# This applies equally to completed and failed jobs and allows the
@@ -530,6 +621,8 @@ def job_view(job: BankStatementAnalysisJob) -> dict:
"purpose": getattr(job, "purpose", "analyze_only"),
"accounting_import_status": getattr(job, "accounting_import_status", "not_requested"),
"accounting_import": json.loads(job.accounting_import_json) if getattr(job, "accounting_import_json", None) else {},
+ "stored_source_versions": json.loads(job.stored_source_versions_json) if getattr(job, "stored_source_versions_json", None) else [],
+ "source_hashes": json.loads(job.source_hashes_json) if getattr(job, "source_hashes_json", None) else [],
"submitted_at": job.submitted_at_utc,
"completed_at": job.completed_at_utc,
"expires_at": job.expires_at_utc,
diff --git a/app/modules/bank_statement_analyzer/stored_statement_service.py b/app/modules/bank_statement_analyzer/stored_statement_service.py
new file mode 100644
index 0000000..2a81e21
--- /dev/null
+++ b/app/modules/bank_statement_analyzer/stored_statement_service.py
@@ -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
diff --git a/app/modules/bank_statement_analyzer/templates/bank_statement_analyzer/index.html b/app/modules/bank_statement_analyzer/templates/bank_statement_analyzer/index.html
index a40c421..99d9097 100644
--- a/app/modules/bank_statement_analyzer/templates/bank_statement_analyzer/index.html
+++ b/app/modules/bank_statement_analyzer/templates/bank_statement_analyzer/index.html
@@ -65,7 +65,7 @@
I confirm that all statements uploaded in this analysis belong to the selected client or one of its ERP-recorded legal/trade/business names.
-
|---|