from __future__ import annotations import json import re from datetime import date, datetime, timezone from difflib import SequenceMatcher 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, ) from app.modules.accounting.chart_models import AccountingChartLedger from app.modules.accounting.chart_service import effective_role from app.modules.bank_statement_analyzer.models import BankStatementAnalysisJob from app.modules.documents.agent_jobs import enqueue_agent_job from app.modules.documents.models import ERPAgentJob, ERPWorkstationAgent ACTION = "accounting_bank_reconciliation_extract" def _utcnow(): return datetime.now(timezone.utc) def _s(value): return str(value or "").strip() def _loads(value, default=None): try: return json.loads(value or "") except Exception: return {} if default is None else default def _norm_ref(value): return re.sub(r"[^A-Z0-9]", "", _s(value).upper()) def _norm_text(value): return re.sub(r"[^A-Z0-9]+", " ", _s(value).upper()).strip() def _date_obj(value): try: return date.fromisoformat(_s(value)[:10]) except Exception: return None def visible_bank_ledgers(db, *, tenant_id: int, client_id: int, tally_guid: str): rows = list( db.execute( select(AccountingChartLedger) .where( AccountingChartLedger.tenant_id == int(tenant_id), AccountingChartLedger.client_id == int(client_id), AccountingChartLedger.tally_guid == _s(tally_guid), ) .order_by(AccountingChartLedger.name) ).scalars().all() ) # Prefer the semantic role, but also trust Tally's immediate parent group. # This keeps reconciliation usable even before a previously mis-resolved # Phase 16 root hierarchy is manually refreshed. return [ row for row in rows if effective_role(row) == "BANK" or _norm_text(row.parent_group_name) == "BANK ACCOUNTS" or _norm_text(row.root_group_name) == "BANK ACCOUNTS" ] def visible_workstations(db, *, tenant_id: int, branch_id: int | None = None): stmt = select(ERPWorkstationAgent).where( ERPWorkstationAgent.tenant_id == int(tenant_id), ERPWorkstationAgent.is_active.is_(True), ERPWorkstationAgent.tally_connected.is_(True), ) if branch_id is not None: stmt = stmt.where(ERPWorkstationAgent.branch_id == int(branch_id)) return list( db.execute( stmt.order_by( ERPWorkstationAgent.last_seen_at_utc.desc(), ERPWorkstationAgent.id.desc(), ) ).scalars().all() ) def completed_client_jobs(db, *, tenant_id: int, client_id: int, limit: int = 50): return list( db.execute( select(BankStatementAnalysisJob) .where( BankStatementAnalysisJob.tenant_id == int(tenant_id), BankStatementAnalysisJob.client_id == int(client_id), BankStatementAnalysisJob.status == "completed", BankStatementAnalysisJob.ownership_status == "confirmed", ) .order_by(BankStatementAnalysisJob.completed_at_utc.desc()) .limit(limit) ).scalars().all() ) def source_accounts(db, *, tenant_id: int, client_id: int, job_id: str): rows = list( db.execute( 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, ) .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/account. " "Use a Bank Reconciliation purpose job or import the completed job into Accounting first." ) return values[0][:10], values[-1][:10] def queue_reconciliation( db, *, tenant_id: int, client_id: int, source_job_id: str, tally_guid: str, company_name: str, 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 ( not source_job or source_job.status != "completed" or int(source_job.client_id or 0) != int(client_id) or int(source_job.tenant_id or 0) != int(tenant_id) or source_job.ownership_status != "confirmed" ): raise ValueError("The selected completed Bank Analyzer job is not valid for this client.") if source_job.accounting_import_status != "completed": raise ValueError( "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 = { 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 Bank ledger from the client's synchronized Chart of Accounts.") ws = db.get(ERPWorkstationAgent, int(workstation_id)) if ( not ws or not ws.is_active or not ws.tally_connected or int(ws.tenant_id) != int(tenant_id) ): raise ValueError("Selected workstation is unavailable or Tally is not connected.") run = BankReconciliationRun( tenant_id=int(tenant_id), client_id=int(client_id), source_job_id=source_job_id, 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), ) db.add(run) db.commit() db.refresh(run) job = enqueue_agent_job( db, workstation_agent_id=ws.id, action=ACTION, payload={ "tenant_id": int(tenant_id), "client_id": int(client_id), "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}:{account_number}:{tally_guid}:{bank_ledger_name}:{run.id}", priority=8, max_attempts=2, created_by_user_id=user_id, ) run.agent_job_id = job.id db.add(run) db.commit() 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 []) bank_entries = [ row for row in entries if _s(row.get("ledger_name")).casefold() == bank_key ] if not bank_entries: return None # Tally exports in this project use negative amount for a Debit ledger and # positive amount for Credit. Bank statement direction is opposite from the # bank ledger accounting side: bank Debit statement = bank ledger Credit. amount = round(abs(float(bank_entries[0].get("amount") or 0)), 2) signed = float(bank_entries[0].get("amount") or 0) direction = "DEBIT" if signed > 0 else "CREDIT" return { "guid": _s(voucher.get("guid")), "voucher_number": _s(voucher.get("voucher_number")), "voucher_type": _s(voucher.get("voucher_type_name")), "date": _s(voucher.get("date"))[:10], "reference": _s(voucher.get("reference")), "narration": _s(voucher.get("narration")), "amount": amount, "direction": direction, } 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(): return 0, "" bd = _date_obj(bank.transaction_date) td = _date_obj(tally["date"]) if not bd or not td: return 0, "" gap = abs((bd - td).days) 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) tally_ref = _norm_ref(tally.get("reference")) ref_exact = bool(bank_ref and tally_ref and bank_ref == tally_ref) bank_text = _norm_text(f"{bank.auto_party} {bank.narration}") tally_text = _norm_text(tally.get("narration")) similarity = int(round(SequenceMatcher(None, bank_text, tally_text).ratio() * 100)) if bank_text and tally_text else 0 score = 60 reasons = ["same amount", "same bank direction"] if gap == 0: score += 22 reasons.append("same date") elif gap <= 2: score += 12 reasons.append(f"{gap}-day timing difference") elif gap <= 4: score += 7 reasons.append(f"{gap}-day timing difference") 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 reasons.append("same reference") elif similarity >= 85: score += 10 reasons.append("high narration similarity") elif similarity >= 65: score += 5 reasons.append("narration similarity") return min(100, score), ", ".join(reasons) def build_reconciliation(db, run: BankReconciliationRun, vouchers: list[dict]): db.execute(delete(BankReconciliationItem).where(BankReconciliationItem.run_id == run.id)) db.commit() bank_rows = list( db.execute( select(AccountingBankTransaction) .where( 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, AccountingBankTransaction.id, ) ).scalars().all() ) tally_rows = [] for voucher in vouchers: if _s(voucher.get("is_cancelled")).lower() in {"yes", "true", "1"}: continue row = _tally_side(voucher, run.bank_ledger_name) if row and row["amount"] > 0: tally_rows.append(row) candidates = {} for bank in bank_rows: scored = [] for index, tally in enumerate(tally_rows): 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 = timing = bank_only = duplicates = 0 for bank in bank_rows: options = [ item for item in candidates.get(bank.id, []) if item[1] not in used_tally ] if not options: status = "bank_only" confidence = 0 reason = "No Tally bank-ledger voucher matched amount, direction and permitted date window." tally = None bank_only += 1 else: top_score, tally_index, reason = options[0] tied = [item for item in options if item[0] == top_score] if len(tied) > 1 and top_score < 95: status = "duplicate_candidate" confidence = top_score tally = tally_rows[tally_index] duplicates += 1 else: used_tally.add(tally_index) tally = tally_rows[tally_index] confidence = top_score 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( run_id=run.id, bank_transaction_id=bank.id, match_status=status, confidence=int(confidence), match_reason=reason, bank_date=_s(bank.transaction_date)[:10], bank_direction=_s(bank.direction).upper(), bank_amount=float(bank.amount or 0), bank_reference=_s(bank.transfer_reference or bank.reference_no), bank_narration=_s(bank.narration), ) if tally: item.tally_guid = tally["guid"] item.tally_voucher_number = tally["voucher_number"] item.tally_voucher_type = tally["voucher_type"] item.tally_date = tally["date"] item.tally_reference = tally["reference"] item.tally_narration = tally["narration"] item.tally_amount = tally["amount"] item.tally_direction = tally["direction"] bank.reconciliation_status = status bank.last_reconciliation_run_id = run.id db.add(bank) db.add(item) books_only = 0 for index, tally in enumerate(tally_rows): if index in used_tally: continue books_only += 1 db.add( BankReconciliationItem( run_id=run.id, bank_transaction_id=None, match_status="books_only", confidence=0, match_reason="Tally bank-ledger voucher has no matching transaction in the uploaded bank statement set.", tally_guid=tally["guid"], tally_voucher_number=tally["voucher_number"], tally_voucher_type=tally["voucher_type"], tally_date=tally["date"], tally_reference=tally["reference"], tally_narration=tally["narration"], tally_amount=tally["amount"], tally_direction=tally["direction"], ) ) 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 = "" db.add(run) db.commit() return run def sync_run(db, run: BankReconciliationRun): if not run.agent_job_id or run.status == "completed": return run job = db.get(ERPAgentJob, int(run.agent_job_id)) if not job: return run if job.status in {"queued", "claimed"}: next_status = "extracting" if job.status == "claimed" else "queued" if run.status != next_status: run.status = next_status db.add(run) db.commit() return run if job.status == "succeeded": result = _loads(job.result_json, {}) vouchers = list(result.get("vouchers") or []) return build_reconciliation(db, run, vouchers) if job.status in {"failed", "cancelled"}: run.status = "failed" run.last_error = _s(job.last_error) or f"Local Agent reconciliation extraction {job.status}." run.completed_at_utc = _utcnow() db.add(run) db.commit() 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( select(BankReconciliationRun) .where( BankReconciliationRun.tenant_id == int(tenant_id), BankReconciliationRun.client_id == int(client_id), ) .order_by(BankReconciliationRun.id.desc()) .limit(limit) ).scalars().all() ) for row in rows: sync_run(db, row) return rows def run_items(db, *, run_id: int, status: str = ""): stmt = select(BankReconciliationItem).where( BankReconciliationItem.run_id == int(run_id) ) if _s(status): stmt = stmt.where(BankReconciliationItem.match_status == _s(status)) return list( db.execute( stmt.order_by( BankReconciliationItem.bank_date, BankReconciliationItem.tally_date, BankReconciliationItem.id, ) ).scalars().all() )