From b8fbaf55ecd62e04e86a47774a642aa3bd113665 Mon Sep 17 00:00:00 2001 From: A R R R Associates Date: Mon, 24 Aug 2026 18:58:58 +0530 Subject: [PATCH] Add Phase 20 stored bank reuse and richer reconciliation --- ...red_reuse_richer_reconciliation_phase20.py | 147 +++++++ .../accounting/bank_reconciliation_models.py | 35 +- .../accounting/bank_reconciliation_service.py | 367 ++++++++++++++++-- .../accounting/bank_reconciliation_ui.py | 184 +++++++++ app/modules/accounting/bank_service.py | 3 + .../accounting/bank_reconciliation.html | 203 +++++++--- .../bank_statement_analyzer/client_context.py | 25 ++ app/modules/bank_statement_analyzer/models.py | 5 + .../bank_statement_analyzer/service.py | 95 ++++- .../stored_statement_service.py | 254 ++++++++++++ .../bank_statement_analyzer/index.html | 75 +++- app/modules/bank_statement_analyzer/ui.py | 104 ++++- 12 files changed, 1419 insertions(+), 78 deletions(-) create mode 100644 alembic/versions/20260824_bank_analyzer_stored_reuse_richer_reconciliation_phase20.py create mode 100644 app/modules/bank_statement_analyzer/stored_statement_service.py 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

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.

@@ -47,54 +47,140 @@ {% if selected_client and selected_tally_guid %}
-

Start a reconciliation run

-

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.

+

1. Select the Bank Analyzer job

+

Choose the completed client-bound analysis whose imported bank transactions you want to reconcile. The job can contain any number of bank accounts.

+
+
+ + + +
+
+
+ + {% if selected_job_id %} +
+
+

2. Map client bank accounts to Tally Bank ledgers

+

Mappings are remembered client-wise and company-wise. Once all accounts are mapped, future reconciliation runs do not require selecting the bank ledger again.

-
+
+ {% for account in source_accounts %} + {% set existing = namespace(value=None) %} + {% for mapping in bank_mappings %} + {% if mapping.account_number == account.account_number %}{% set existing.value = mapping %}{% endif %} + {% endfor %} + + + + + + {% for company in companies %}{% if company.guid==selected_tally_guid %}{% endif %}{% endfor %} + + + +
+
Statement account
+
{{ account.bank_name or 'Bank' }} · {{ account.account_number }}
+
{{ account.transaction_count }} imported transaction(s)
+
+ +
+ + {% else %} +
No imported bank accounts were found for this job. Ensure its Accounting import status is Completed.
+ {% endfor %} +
+
+ +
+
+

3. Reconcile mapped bank accounts

+

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.

+
+ +
+ + {% for company in companies %}{% if company.guid==selected_tally_guid %}{% endif %}{% endfor %} - - - - - - -
{% endif %} + {% endif %} {% if selected_client %}
@@ -104,19 +190,21 @@
- + {% for run in runs %} {% set summary = run.summary_json|default('{}') %} - + @@ -144,7 +232,7 @@ @@ -158,6 +246,15 @@ {% elif selected_run.status=='failed' %}
{{ selected_run.last_error }}
{% elif selected_run.status=='completed' %} + {% set rs = run_summaries.get(selected_run.id, {}) %} +
+
Matched
{{ rs.get('matched',0) }}
+
Probable
{{ rs.get('probable_match',0) }}
+
Timing differences
{{ rs.get('timing_difference',0) }}
+
Bank only
{{ rs.get('bank_only',0) }}
Dr ₹{{ '%.2f'|format(rs.get('bank_only_debit_amount',0)) }} · Cr ₹{{ '%.2f'|format(rs.get('bank_only_credit_amount',0)) }}
+
Books only
{{ rs.get('books_only',0) }}
Dr ₹{{ '%.2f'|format(rs.get('books_only_debit_amount',0)) }} · Cr ₹{{ '%.2f'|format(rs.get('books_only_credit_amount',0)) }}
+
Unreconciled net
₹{{ '%.2f'|format(rs.get('unreconciled_net_difference',0)) }}
Tolerance {{ rs.get('date_tolerance_days',15) }} days
+
RunBank LedgerPeriodStatusResult
RunBank / AccountPeriodStatusResult
#{{ 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
@@ -168,6 +265,7 @@ @@ -191,13 +289,28 @@ {% else %}—{% 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. -

When a client is selected, the analyzer validates every statement owner against the client name, trade name, business-unit names, registration legal/trade names and branch names. On successful analysis, the source PDFs and generated workbook are archived using the existing Engagement Documents storage pipeline.

+

When a client is selected, the analyzer validates every statement owner against the client name, trade name, business-unit names, registration legal/trade names and branch names. On successful analysis, newly uploaded source PDFs and the generated workbook are archived using the existing Engagement Documents storage pipeline. Previously stored PDFs selected for reuse remain linked to their existing document/version and are not duplicated.

{% endif %}

Keep Auto Detect or select a bank for direct parser validation.

@@ -74,7 +74,36 @@
-
Multi-bank contra detection: when the client has more than one bank account, upload all relevant statements in the same analysis job. There is no fixed statement-count limit unless your administrator configures one. The analyzer pairs only conservative equal-and-opposite transfers across different account numbers and keeps every pair reviewable.
+
Multi-bank contra detection: when the client has more than one bank account, upload all relevant statements in the same analysis job. There is no fixed statement-count limit unless your administrator configures one. The analyzer pairs only conservative equal-and-opposite transfers across different account numbers and keeps every pair reviewable.
+
+ + +

Select any number of new statements. They can be combined with previously stored statements below.

+
+ +
+
+
+
Reuse bank statements already stored for this client
+

Existing BANK_STATEMENT documents are reused from Engagement Documents/local storage and are not archived again as duplicates.

+
+ +
+
+ {% for row in stored_bank_statements or [] %} + + {% else %} +
Select an ERP client to load previously stored bank statements.
+ {% endfor %} +
+
+
Queue limits
Maximum three processing jobs across all users. Each user may have up to three queued or processing jobs. Completed workbooks and original uploaded statements remain available for 24 hours. Failed-job statements are also retained for 24 hours for debugging.
@@ -85,7 +114,7 @@
Files
{{ active_job.file_count }}
Queue position
{{ active_job.queue_position or '—' }}
Estimated wait
{{ active_job.estimated_wait or '—' }}
Progress
{{ active_job.progress_percent }}%
{% if active_job.status == 'completed' %} -

Analysis completed

Statements
{{ active_job.summary.statement_count or 0 }}
Transactions extracted
{{ active_job.summary.rows_extracted or 0 }}
Exact duplicates
{{ active_job.summary.exact_duplicate_rows or 0 }}
Review items
{{ active_job.summary.review_items or 0 }}
Inter-bank contra pairs
{{ active_job.summary.contra_pairs or 0 }}
Download Excel{% for file in active_job.original_files %}Download Statement {{ file.index }}{% endfor %}Analyze Another Bank

The workbook and original statement{{ 's' if active_job.file_count != 1 else '' }} remain available until {{ active_job.expires_at or '24 hours after completion' }}.

+

Analysis completed

Statements
{{ active_job.summary.statement_count or 0 }}
Transactions extracted
{{ active_job.summary.rows_extracted or 0 }}
Exact duplicates
{{ active_job.summary.exact_duplicate_rows or 0 }}
Review items
{{ active_job.summary.review_items or 0 }}
Inter-bank contra pairs
{{ active_job.summary.contra_pairs or 0 }}
{% if active_job.summary.stored_statement_reuse_count %}
Reused {{ active_job.summary.stored_statement_reuse_count }} statement(s) from existing Engagement Documents without creating duplicate source documents.
{% endif %}{% if active_job.summary.statement_overlap_warnings %}
Overlapping statement periods detected
{% for warning in active_job.summary.statement_overlap_warnings %}
{{ warning.message }}
{% endfor %}
Transactions are still deduplicated by the existing analyzer rules; verify overlapping-period coverage before reconciliation.
{% endif %}
Download Excel{% for file in active_job.original_files %}Download Statement {{ file.index }}{% endfor %}Analyze Another Bank

The workbook and original statement{{ 's' if active_job.file_count != 1 else '' }} remain available until {{ active_job.expires_at or '24 hours after completion' }}.

{% elif active_job.status == 'failed' %}
Analysis failed.
{{ active_job.error_message }}
{% if active_job.original_files %}
{% for file in active_job.original_files %}Download Statement {{ file.index }}{% endfor %}

Original statement{{ 's are' if active_job.file_count != 1 else ' is' }} retained until {{ active_job.expires_at or '24 hours after failure' }} for debugging.

{% endif %}Analyze another statement
{% else %}
{% if active_job.status == 'queued' %}Your job is queued. You may safely leave this page and return through My Analysis Jobs.{% else %}Your statements are being processed.{% endif %}
{% endif %} {% if active_job.status == 'completed' and active_job.purpose in ['accounting_entries','bank_reconciliation'] %}
Accounting import: {{ active_job.accounting_import_status|replace('_',' ')|title }}{% if active_job.summary.client_id %} · Open accounting queue{% endif %}
{% endif %} @@ -114,8 +143,46 @@ if (!selectedStillVisible) engagement.value = ""; engagement.required = !!selectedClient; } - client.addEventListener("change", filterEngagements); + async function loadStoredStatements() { + const box = document.getElementById("stored-bank-statements"); + if (!box) return; + const selectedClient = client.value; + if (!selectedClient) { + box.innerHTML = '
Select an ERP client to load previously stored bank statements.
'; + return; + } + box.innerHTML = '
Loading stored statements…
'; + try { + const response = await fetch(`/tools/bank-statement-analyzer/stored-statements?client_id=${encodeURIComponent(selectedClient)}`, {headers: {"Accept":"application/json"}}); + const data = await response.json(); + if (!response.ok) throw new Error(data.detail || "Unable to load stored statements."); + const items = data.items || []; + if (!items.length) { + box.innerHTML = '
No stored BANK_STATEMENT documents are available for this client.
'; + return; + } + box.innerHTML = ""; + items.forEach((row) => { + const label = document.createElement("label"); + label.className = "flex items-start gap-2 rounded-lg border border-indigo-100 bg-white p-3 text-xs"; + label.innerHTML = ``; + label.querySelector(".font-semibold").textContent = row.title || row.filename; + label.querySelector(".text-slate-500").textContent = `${row.filename} · version ${row.version_id} · ${(row.storage_status || "").replaceAll("_"," ")}`; + box.appendChild(label); + }); + } catch (error) { + box.innerHTML = `
${error.message}
`; + } + } + + client.addEventListener("change", () => { + filterEngagements(); + loadStoredStatements(); + }); + const refreshButton = document.getElementById("refresh-stored-bank-statements"); + if (refreshButton) refreshButton.addEventListener("click", loadStoredStatements); filterEngagements(); + if (client.value) loadStoredStatements(); })(); diff --git a/app/modules/bank_statement_analyzer/ui.py b/app/modules/bank_statement_analyzer/ui.py index c532bfc..6441077 100644 --- a/app/modules/bank_statement_analyzer/ui.py +++ b/app/modules/bank_statement_analyzer/ui.py @@ -17,6 +17,11 @@ from app.modules.core.rbac.deps import get_user_permissions, get_user_roles from .parsers.registry import BANK_OPTIONS from .client_context import analyzer_client_context, validate_selected_client_and_engagement +from .stored_statement_service import ( + deduplicate_source_paths, + list_stored_bank_statements, + prepare_stored_bank_statements, +) from .service import ( can_use, create_job_folder, @@ -107,7 +112,7 @@ def _auth(request, db): @router.get("") -def index(request: Request, job: str | None = None): +def index(request: Request, job: str | None = None, client_id: int | None = None): ensure_worker_started() db = CommonSessionLocal() try: @@ -119,6 +124,19 @@ def index(request: Request, job: str | None = None): recent = [_localised_job_view(item, timezone_name) for item in list_user_jobs(user.id, limit=8)] active_job = _localised_job_view(selected_job, timezone_name) if selected_job else None client_context = analyzer_client_context(db, request=request, user=user, roles=roles) + visible_client_ids = {int(row.id) for row in client_context["clients"]} + selected_client_id = int(client_id) if client_id and int(client_id) in visible_client_ids else None + stored_statements = ( + list_stored_bank_statements( + db, + request=request, + user=user, + roles=roles, + client_id=selected_client_id, + ) + if selected_client_id + else [] + ) return templates.TemplateResponse( "modules/bank_statement_analyzer/templates/bank_statement_analyzer/index.html", _ctx( @@ -131,14 +149,38 @@ def index(request: Request, job: str | None = None): display_timezone=timezone_name, analyzer_clients=client_context["clients"], analyzer_engagements=client_context["engagements"], + selected_client_id=selected_client_id, + stored_bank_statements=stored_statements, ), ) finally: db.close() +@router.get("/stored-statements") +def stored_statements(request: Request, client_id: int): + db = CommonSessionLocal() + try: + user, roles, denied = _auth(request, db) + if denied: + return JSONResponse({"detail": "Access denied"}, status_code=403) + try: + rows = list_stored_bank_statements( + db, + request=request, + user=user, + roles=roles, + client_id=int(client_id), + ) + return JSONResponse({"items": rows}) + except Exception as exc: + return JSONResponse({"detail": str(exc)}, status_code=400) + finally: + db.close() + + @router.post("/analyze") -async def analyze(request: Request, csrf_token: str = Form(...), bank_selection: str = Form("auto"), financial_year: str = Form(""), customer_name: str = Form(""), account_number: str = Form(""), client_id: str = Form(""), engagement_id: str = Form(""), confirm_same_client: str | None = Form(None), enable_classification: str | None = Form(None), purpose: str = Form("analyze_only"), statements: list[UploadFile] = File(...)): +async def analyze(request: Request, csrf_token: str = Form(...), bank_selection: str = Form("auto"), financial_year: str = Form(""), customer_name: str = Form(""), account_number: str = Form(""), client_id: str = Form(""), engagement_id: str = Form(""), confirm_same_client: str | None = Form(None), enable_classification: str | None = Form(None), purpose: str = Form("analyze_only"), stored_version_ids: list[int] = Form([]), statements: list[UploadFile] = File(default=[])): db = CommonSessionLocal() selected_bank = bank_selection if bank_selection in dict(BANK_OPTIONS) else "auto" classification_enabled = enable_classification == "1" @@ -169,7 +211,50 @@ async def analyze(request: Request, csrf_token: str = Form(...), bank_selection: job_id, input_dir, _output_dir = create_job_folder(user, roles) job_dir = input_dir.parent - paths = await save_uploads(statements, input_dir) + + uploaded = [] + usable_uploads = [item for item in (statements or []) if item and (item.filename or "").strip()] + if usable_uploads: + uploaded = await save_uploads(usable_uploads, input_dir) + + stored_paths = [] + stored_provenance = [] + pending_retrieval = [] + if stored_version_ids: + if not resolved_client_id: + raise ValueError("Select an ERP client before reusing stored bank statements.") + stored_paths, stored_provenance, pending_retrieval = prepare_stored_bank_statements( + db, + request=request, + user=user, + roles=roles, + client_id=resolved_client_id, + version_ids=stored_version_ids, + input_dir=input_dir, + ) + + if pending_retrieval: + db.commit() + details = ", ".join( + f"{row['filename']} (retrieval #{row.get('request_id') or 'pending'})" + for row in pending_retrieval + ) + raise ValueError( + "Stored statement retrieval has been queued from branch local storage: " + + details + + ". Wait for the Local Agent to return the file, then submit the analysis again." + ) + + paths, source_hashes = deduplicate_source_paths(uploaded + stored_paths) + retained_names = {path.name for path in paths} + stored_provenance = [ + row + for row in stored_provenance + if row.get("copied_filename") in retained_names + ] + if not paths: + raise ValueError("Select at least one new PDF or one stored bank statement.") + enqueue_job( user=user, roles=roles, @@ -185,6 +270,8 @@ async def analyze(request: Request, csrf_token: str = Form(...), bank_selection: engagement_id=resolved_engagement_id, ownership_confirmation=(confirm_same_client == "1"), purpose=purpose, + stored_source_versions=stored_provenance, + source_hashes=source_hashes, ) return RedirectResponse(f"/tools/bank-statement-analyzer?job={job_id}#analysis-status", status_code=303) except Exception as exc: @@ -215,6 +302,17 @@ async def analyze(request: Request, csrf_token: str = Form(...), bank_selection: analyzer_engagements=client_context["engagements"], selected_client_id=(int(client_id) if str(client_id).strip().isdigit() else None), selected_engagement_id=(int(engagement_id) if str(engagement_id).strip().isdigit() else None), + stored_bank_statements=( + list_stored_bank_statements( + db, + request=request, + user=user, + roles=get_user_roles(db, user.id), + client_id=int(client_id), + ) + if str(client_id).strip().isdigit() + else [] + ), ), status_code=400, )
ResultBank StatementTallyWhyAction
{{ item.match_reason }} + + {% if item.resolution_status and item.resolution_status != 'unresolved' %} +
Reviewed: {{ item.resolution_status|replace('_',' ')|title }}{% if item.resolution_note %}
{{ item.resolution_note }}
{% endif %}
+ {% endif %} +
+ + + + + + +
{% if item.match_status=='bank_only' and item.bank_transaction_id %} - Review in Accounting Queue - {% elif item.match_status=='books_only' %} - Investigate timing / statement coverage. - {% else %} - Review if needed. + Open Accounting Queue {% endif %}