from __future__ import annotations import hashlib import json from datetime import datetime, timezone from pathlib import Path import pandas as pd from sqlalchemy import func, select from app.modules.accounting.bank_models import AccountingBankTransaction from app.modules.accounting.ai_service import ai_assist_bank as run_ai_assist_bank, mark_review_outcome from app.modules.accounting.internal_model_service import predict_bank, record_prediction_review from app.modules.accounting.ledger_learning_service import ( active_natures, available_tally_guids, rank_suggestions, record_review, ) 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 PREFLIGHT_ACTION = "accounting_bank_posting_preflight" POST_ACTION = "accounting_post_bank_voucher" def _utcnow(): return datetime.now(timezone.utc) def _s(value): if value is None: return "" try: if pd.isna(value): return "" except Exception: pass return str(value).strip() def _f(value): try: if pd.isna(value): return 0.0 return round(float(value or 0), 2) except Exception: return 0.0 def _date(value): dt = pd.to_datetime(value, errors="coerce") return "" if pd.isna(dt) else dt.strftime("%Y-%m-%d") def _fingerprint(client_id, row): raw = "|".join([ str(client_id), _s(row.get("bank_name")), _s(row.get("account_number")), _date(row.get("transaction_date")), "D" if _f(row.get("debit")) > 0 else "C", f"{max(_f(row.get('debit')), _f(row.get('credit'))):.2f}", _s(row.get("reference_no")), _s(row.get("transfer_reference")), _s(row.get("narration")).upper(), ]) return hashlib.sha256(raw.encode("utf-8", "ignore")).hexdigest() def _voucher_type(row): if _s(row.get("contra_pair_id")): return "Contra" return "Payment" if _f(row.get("debit")) > 0 else "Receipt" def import_completed_job(db, *, tenant_id, client_id, job_id, user_id): job = db.get(BankStatementAnalysisJob, job_id) if not job or job.status != "completed" or not job.output_file or not Path(job.output_file).is_file(): raise ValueError("Completed Bank Analyzer workbook was not found.") if job.tenant_id not in (None, tenant_id): raise ValueError("Bank Analyzer job belongs to a different firm.") if getattr(job, "client_id", None) and int(job.client_id) != int(client_id): raise ValueError( "This Bank Analyzer job is already bound to a different ERP client. " "Import it using the client confirmed during statement ownership validation." ) if getattr(job, "client_id", None) and getattr(job, "ownership_status", "") != "confirmed": raise ValueError("Client-bound Bank Analyzer job has not passed ownership validation.") frame = pd.read_excel(job.output_file, sheet_name="Transaction Classification") tally_options = available_tally_guids(db, tenant_id, client_id) tally_guid = tally_options[0][0] if len(tally_options) == 1 else "" inserted = 0 skipped = 0 for _, row in frame.iterrows(): fp = _fingerprint(client_id, row) exists = db.execute(select(AccountingBankTransaction.id).where( AccountingBankTransaction.tenant_id == tenant_id, AccountingBankTransaction.client_id == client_id, AccountingBankTransaction.fingerprint == fp, )).scalar_one_or_none() if exists: skipped += 1 continue contra_pair = _s(row.get("contra_pair_id")) suggestions = [] if not contra_pair: suggestions = rank_suggestions( db, tenant_id=tenant_id, client_id=client_id, tally_guid=tally_guid, supplier_name=_s(row.get("auto_party")), supplier_gstin="", hsn_code="", description=_s(row.get("narration")), amount=max(_f(row.get("debit")), _f(row.get("credit"))), ) top = suggestions[0] if suggestions else None tx = AccountingBankTransaction( tenant_id=tenant_id, client_id=client_id, source_job_id=job.id, fingerprint=fp, statement_id=_s(row.get("statement_id")), bank_name=_s(row.get("bank_name")), account_number=_s(row.get("account_number")), customer_name=_s(row.get("customer_name")), transaction_date=_date(row.get("transaction_date")), value_date=_date(row.get("value_date")), narration=_s(row.get("narration")), reference_no=_s(row.get("reference_no")), transfer_reference=_s(row.get("transfer_reference")), direction=_s(row.get("direction")), debit=_f(row.get("debit")), credit=_f(row.get("credit")), amount=max(_f(row.get("debit")), _f(row.get("credit"))), auto_party=_s(row.get("auto_party")), analyzer_category=_s(row.get("auto_category")), analyzer_nature=_s(row.get("auto_nature")), analyzer_ledger=_s(row.get("suggested_ledger")), contra_pair_id=contra_pair, contra_counter_bank=_s(row.get("contra_counter_bank")), contra_counter_account=_s(row.get("contra_counter_account")), contra_confidence=_f(row.get("contra_confidence")), contra_reason=_s(row.get("contra_reason")) or None, tally_guid=tally_guid, suggested_nature_id=(top["nature"].id if top else None), suggested_ledger_name=(top["suggested_ledger"] if top else ""), suggested_confidence=(int(top["confidence"]) if top else int(_f(row.get("contra_confidence")))), suggestion_reason_json=json.dumps(top["reasons"] if top else ([f"Matched multi-bank contra pair {contra_pair}."] if contra_pair else []), ensure_ascii=False), suggested_voucher_type=_voucher_type(row), final_voucher_type=_voucher_type(row), ) db.add(tx) inserted += 1 db.commit() return inserted, skipped def queue_rows(db, *, tenant_id, client_id, status="", page=1, per_page=25): stmt = select(AccountingBankTransaction).where( AccountingBankTransaction.tenant_id == tenant_id, AccountingBankTransaction.client_id == client_id, ) count_stmt = select(func.count()).select_from(AccountingBankTransaction).where( AccountingBankTransaction.tenant_id == tenant_id, AccountingBankTransaction.client_id == client_id, ) if status: stmt = stmt.where(AccountingBankTransaction.review_status == status) count_stmt = count_stmt.where(AccountingBankTransaction.review_status == status) total = int(db.scalar(count_stmt) or 0) pages = max(1, (total + per_page - 1) // per_page) page = max(1, min(page, pages)) rows = list(db.execute( stmt.order_by(AccountingBankTransaction.transaction_date.desc(), AccountingBankTransaction.id.desc()) .offset((page - 1) * per_page).limit(per_page) ).scalars().all()) return rows, total, page, pages def confirm_review(db, *, tx_id, tenant_id, client_id, nature_id, ledger_name, voucher_type, user_id, note=""): tx = db.get(AccountingBankTransaction, tx_id) if not tx or tx.tenant_id != tenant_id or tx.client_id != client_id: raise ValueError("Bank transaction was not found.") voucher_type = (voucher_type or tx.suggested_voucher_type or "").title() if voucher_type not in {"Payment", "Receipt", "Contra"}: raise ValueError("Voucher type must be Payment, Receipt or Contra.") tx.final_voucher_type = voucher_type tx.review_note = (note or "").strip() or None tx.reviewed_by_user_id = user_id tx.reviewed_at_utc = _utcnow() tx.review_status = "reviewed" if voucher_type == "Contra": if not tx.contra_pair_id: raise ValueError("Contra treatment requires a matched multi-bank contra pair.") tx.final_nature_id = None tx.final_ledger_name = "" else: if not nature_id: raise ValueError("Select the final accounting nature.") tx.final_nature_id = int(nature_id) tx.final_ledger_name = (ledger_name or "").strip() record_review( db, tenant_id=tenant_id, client_id=client_id, tally_guid=tx.tally_guid or "", supplier_name=tx.auto_party, supplier_gstin="", hsn_code="", description=tx.narration, amount=tx.amount, suggested_nature_id=tx.suggested_nature_id, suggested_ledger_name=tx.suggested_ledger_name, suggested_confidence=tx.suggested_confidence, final_nature_id=int(nature_id), final_ledger_name=tx.final_ledger_name, user_id=user_id, explanation=json.loads(tx.suggestion_reason_json or "[]"), ) # record_review commits; refresh before final mutation persistence tx = db.get(AccountingBankTransaction, tx_id) tx.final_voucher_type = voucher_type tx.review_note = (note or "").strip() or None tx.reviewed_by_user_id = user_id tx.reviewed_at_utc = _utcnow() tx.review_status = "reviewed" tx.final_nature_id = int(nature_id) tx.final_ledger_name = (ledger_name or "").strip() db.add(tx) db.commit() mark_review_outcome( db, tenant_id=tenant_id, source_type="bank", source_record_id=tx.id, final_nature_id=tx.final_nature_id, final_ledger_name=tx.final_ledger_name, user_id=user_id, ) record_prediction_review( db, tenant_id=tenant_id, source_type="bank", source_record_id=tx.id, final_nature_id=tx.final_nature_id, user_id=user_id, ) return tx def ai_assist_transaction(db, *, tx_id, tenant_id, client_id, user_id): tx = db.get(AccountingBankTransaction, int(tx_id)) if not tx or int(tx.tenant_id) != int(tenant_id) or int(tx.client_id) != int(client_id): raise ValueError("Bank transaction was not found.") return run_ai_assist_bank(db, tx=tx, user_id=user_id) def internal_model_predict_transaction(db, *, tx_id, tenant_id, client_id, force_shadow=True): tx = db.get(AccountingBankTransaction, int(tx_id)) if not tx or int(tx.tenant_id) != int(tenant_id) or int(tx.client_id) != int(client_id): raise ValueError("Bank transaction was not found.") return predict_bank(db, row=tx, force_shadow=force_shadow) def visible_workstations(db, tenant_id, branch_id=None): stmt = select(ERPWorkstationAgent).where( ERPWorkstationAgent.tenant_id == tenant_id, ERPWorkstationAgent.is_active.is_(True), ) if branch_id is not None: stmt = stmt.where(ERPWorkstationAgent.branch_id == branch_id) return list(db.execute(stmt.order_by(ERPWorkstationAgent.tally_connected.desc(), ERPWorkstationAgent.id)).scalars().all()) def _ensure_reviewed(tx): if tx.review_status != "reviewed": raise ValueError("Review the bank transaction before Tally posting.") if tx.posting_status == "posted": raise ValueError("This bank transaction is already posted.") def queue_preflight(db, *, tx_id, tenant_id, client_id, workstation_id, user_id): tx = db.get(AccountingBankTransaction, tx_id) if not tx or tx.tenant_id != tenant_id or tx.client_id != client_id: raise ValueError("Bank transaction was not found.") _ensure_reviewed(tx) ws = db.get(ERPWorkstationAgent, int(workstation_id)) if not ws or ws.tenant_id != tenant_id or not ws.is_active or not ws.tally_connected: raise ValueError("Selected workstation/Tally connection is unavailable.") if not tx.tally_guid: raise ValueError("Map the client to a Tally company before bank posting.") payload = { "tenant_id": tenant_id, "client_id": client_id, "tally_guid": tx.tally_guid, "bank_name": tx.bank_name, "account_number": tx.account_number, "counter_bank": tx.contra_counter_bank, "counter_account": tx.contra_counter_account, "party_hint": tx.auto_party, "counter_ledger_hint": tx.final_ledger_name, "voucher_type": tx.final_voucher_type, "transaction_date": tx.transaction_date, "amount": tx.amount, "reference": tx.transfer_reference or tx.reference_no, "erp_bank_transaction_id": tx.id, } job = enqueue_agent_job( db, workstation_agent_id=ws.id, action=PREFLIGHT_ACTION, payload=payload, idempotency_key=f"bank:{client_id}:{tx.id}:preflight:{ws.id}", priority=9, max_attempts=2, created_by_user_id=user_id, ) tx.workstation_agent_id = ws.id tx.preflight_job_id = job.id tx.posting_status = "preflight_queued" tx.last_error = None db.add(tx) db.commit() return tx def _loads(value): try: return json.loads(value or "{}") except Exception: return {} def sync_posting(db, tx): changed = False if tx.preflight_job_id and tx.posting_status.startswith("preflight"): job = db.get(ERPAgentJob, tx.preflight_job_id) if job: if job.status == "claimed": tx.posting_status = "preflight_claimed"; changed = True elif job.status == "succeeded": tx.preflight_result_json = job.result_json tx.posting_status = "preflight_ready"; tx.last_error = None; changed = True elif job.status in {"failed", "cancelled"}: tx.posting_status = "preflight_failed"; tx.last_error = job.last_error or job.status; changed = True if tx.posting_job_id and tx.posting_status.startswith("posting"): job = db.get(ERPAgentJob, tx.posting_job_id) if job: if job.status == "claimed": tx.posting_status = "posting_claimed"; changed = True elif job.status == "succeeded": result = _loads(job.result_json) tally = result.get("tally_result") or result tx.posting_result_json = job.result_json tx.tally_voucher_id = str(tally.get("last_voucher_id") or tally.get("voucher_id") or "") tx.tally_voucher_number = str(tally.get("voucher_number") or tally.get("last_voucher_id") or "") tx.posting_status = "posted" tx.posted_at_utc = job.completed_at_utc or _utcnow() tx.last_error = None changed = True elif job.status == "failed": err = job.last_error or "Tally bank posting failed." tx.posting_status = "posting_indeterminate" if "Verify Tally before retrying" in err else "posting_failed" tx.last_error = err; changed = True elif job.status == "cancelled": tx.posting_status = "posting_failed"; tx.last_error = "Posting job cancelled."; changed = True if changed: tx.updated_at_utc = _utcnow() db.add(tx); db.commit() return tx def preflight_choices(tx): return _loads(tx.preflight_result_json) def queue_post(db, *, tx_id, tenant_id, client_id, bank_ledger_name, counter_ledger_name, other_bank_ledger_name, user_id): tx = db.get(AccountingBankTransaction, tx_id) if not tx or tx.tenant_id != tenant_id or tx.client_id != client_id: raise ValueError("Bank transaction was not found.") tx = sync_posting(db, tx) if tx.posting_status != "preflight_ready": raise ValueError("Successful workstation preflight is required.") choices = preflight_choices(tx) banks = {str(x.get("name") or "") for x in choices.get("bank_ledgers") or []} ledgers = {str(x.get("name") or "") for x in choices.get("all_ledgers") or []} bank_ledger_name = (bank_ledger_name or "").strip() counter_ledger_name = (counter_ledger_name or "").strip() other_bank_ledger_name = (other_bank_ledger_name or "").strip() if bank_ledger_name not in banks: raise ValueError("Select a bank ledger returned by the current Tally preflight.") if tx.final_voucher_type == "Contra": if other_bank_ledger_name not in banks or other_bank_ledger_name == bank_ledger_name: raise ValueError("Select the other bank ledger for this contra pair.") else: if counter_ledger_name not in ledgers: raise ValueError("Select a valid Tally counter ledger.") if choices.get("duplicate_candidates"): raise ValueError("Possible duplicate bank voucher exists in Tally. Verify before posting.") payload = { "tenant_id": tenant_id, "client_id": client_id, "tally_guid": tx.tally_guid, "erp_bank_transaction_id": tx.id, "voucher_type": tx.final_voucher_type, "transaction_date": tx.transaction_date, "bank_ledger_name": bank_ledger_name, "counter_ledger_name": counter_ledger_name, "other_bank_ledger_name": other_bank_ledger_name, "direction": tx.direction, "amount": tx.amount, "reference": tx.transfer_reference or tx.reference_no or f"ERP-BANK-{tx.id}", "narration": f"ERP Bank Analyzer #{tx.id}: {tx.narration}"[:1000], } ws = db.get(ERPWorkstationAgent, tx.workstation_agent_id) job = enqueue_agent_job( db, workstation_agent_id=ws.id, action=POST_ACTION, payload=payload, idempotency_key=f"bank:{client_id}:{tx.id}:post:{ws.id}", priority=10, max_attempts=1, created_by_user_id=user_id, ) tx.final_bank_ledger_name = bank_ledger_name tx.final_party_ledger_name = counter_ledger_name tx.final_other_bank_ledger_name = other_bank_ledger_name tx.posting_job_id = job.id tx.posted_by_user_id = user_id tx.posting_status = "posting_queued" db.add(tx); db.commit() return tx