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 ( 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 _job_period(db, *, tenant_id: int, client_id: int, job_id: str): values = list( db.execute( select(AccountingBankTransaction.transaction_date) .where( AccountingBankTransaction.tenant_id == int(tenant_id), AccountingBankTransaction.client_id == int(client_id), AccountingBankTransaction.source_job_id == job_id, ) .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. " "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, ): 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." ) date_from, date_to = _job_period( db, tenant_id=tenant_id, client_id=client_id, job_id=source_job_id, ) 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, date_from=date_from, date_to=date_to, 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, "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}", 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 _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): 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) if gap > 7: 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") else: score += 2 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, ) .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) 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 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 status = "matched" if top_score >= 90 else "probable_match" if status == "matched": exact += 1 else: 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"], ) ) 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, ) 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 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() )