diff --git a/alembic/versions/20260822_bank_to_tally_phase11.py b/alembic/versions/20260822_bank_to_tally_phase11.py new file mode 100644 index 0000000..ec8733e --- /dev/null +++ b/alembic/versions/20260822_bank_to_tally_phase11.py @@ -0,0 +1,81 @@ +"""Phase 11 Bank to Tally + multi-bank contra. + +Revision ID: 20260822_bank_to_tally_p11 +Revises: 20260822_purchase_posting_p10 +""" +from alembic import op +import sqlalchemy as sa + +revision = "20260822_bank_to_tally_p11" +down_revision = "20260822_purchase_posting_p10" +branch_labels = None +depends_on = None + + +def upgrade(): + op.create_table( + "accounting_bank_transactions", + 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("source_job_id", sa.String(32), sa.ForeignKey("bank_statement_analysis_jobs.id", ondelete="SET NULL"), nullable=True), + sa.Column("fingerprint", sa.String(80), nullable=False), + sa.Column("statement_id", sa.String(40), nullable=False, server_default=""), + sa.Column("bank_name", sa.String(160), nullable=False, server_default=""), + sa.Column("account_number", sa.String(100), nullable=False, server_default=""), + sa.Column("customer_name", sa.String(255), nullable=False, server_default=""), + sa.Column("transaction_date", sa.String(20), nullable=False, server_default=""), + sa.Column("value_date", sa.String(20), nullable=False, server_default=""), + sa.Column("narration", sa.Text(), nullable=False, server_default=""), + sa.Column("reference_no", sa.String(180), nullable=False, server_default=""), + sa.Column("transfer_reference", sa.String(180), nullable=False, server_default=""), + sa.Column("direction", sa.String(10), nullable=False, server_default=""), + sa.Column("debit", sa.Float(), nullable=False, server_default="0"), + sa.Column("credit", sa.Float(), nullable=False, server_default="0"), + sa.Column("amount", sa.Float(), nullable=False, server_default="0"), + sa.Column("auto_party", sa.String(255), nullable=False, server_default=""), + sa.Column("analyzer_category", sa.String(255), nullable=False, server_default=""), + sa.Column("analyzer_nature", sa.String(160), nullable=False, server_default=""), + sa.Column("analyzer_ledger", sa.String(255), nullable=False, server_default=""), + sa.Column("contra_pair_id", sa.String(60), nullable=False, server_default=""), + sa.Column("contra_counter_bank", sa.String(160), nullable=False, server_default=""), + sa.Column("contra_counter_account", sa.String(100), nullable=False, server_default=""), + sa.Column("contra_confidence", sa.Float(), nullable=False, server_default="0"), + sa.Column("contra_reason", sa.Text(), nullable=True), + sa.Column("tally_guid", sa.String(120), nullable=False, server_default=""), + sa.Column("suggested_nature_id", sa.Integer(), sa.ForeignKey("accounting_natures.id", ondelete="SET NULL"), nullable=True), + sa.Column("suggested_ledger_name", sa.String(255), nullable=False, server_default=""), + sa.Column("suggested_confidence", sa.Integer(), nullable=False, server_default="0"), + sa.Column("suggestion_reason_json", sa.Text(), nullable=True), + sa.Column("suggested_voucher_type", sa.String(20), nullable=False, server_default=""), + sa.Column("review_status", sa.String(20), nullable=False, server_default="pending"), + sa.Column("final_nature_id", sa.Integer(), sa.ForeignKey("accounting_natures.id", ondelete="SET NULL"), nullable=True), + sa.Column("final_ledger_name", sa.String(255), nullable=False, server_default=""), + sa.Column("final_party_ledger_name", sa.String(255), nullable=False, server_default=""), + sa.Column("final_bank_ledger_name", sa.String(255), nullable=False, server_default=""), + sa.Column("final_other_bank_ledger_name", sa.String(255), nullable=False, server_default=""), + sa.Column("final_voucher_type", sa.String(20), nullable=False, server_default=""), + sa.Column("review_note", sa.Text(), nullable=True), + sa.Column("reviewed_by_user_id", sa.Integer(), sa.ForeignKey("users.id", ondelete="SET NULL"), nullable=True), + sa.Column("reviewed_at_utc", sa.DateTime(timezone=True), nullable=True), + sa.Column("workstation_agent_id", sa.Integer(), sa.ForeignKey("erp_workstation_agents.id", ondelete="SET NULL"), nullable=True), + sa.Column("preflight_job_id", sa.Integer(), sa.ForeignKey("erp_agent_jobs.id", ondelete="SET NULL"), nullable=True), + sa.Column("posting_job_id", sa.Integer(), sa.ForeignKey("erp_agent_jobs.id", ondelete="SET NULL"), nullable=True), + sa.Column("preflight_result_json", sa.Text(), nullable=True), + sa.Column("posting_result_json", sa.Text(), nullable=True), + sa.Column("posting_status", sa.String(30), nullable=False, server_default="not_started"), + sa.Column("last_error", sa.Text(), nullable=True), + sa.Column("tally_voucher_id", sa.String(120), nullable=False, server_default=""), + sa.Column("tally_voucher_number", sa.String(160), nullable=False, server_default=""), + sa.Column("posted_by_user_id", sa.Integer(), sa.ForeignKey("users.id", ondelete="SET NULL"), nullable=True), + sa.Column("posted_at_utc", sa.DateTime(timezone=True), 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", "fingerprint", name="uq_accounting_bank_tx_fingerprint"), + ) + for name in ("tenant_id","client_id","source_job_id","fingerprint","transaction_date","contra_pair_id","review_status","posting_status","created_at_utc"): + op.create_index(f"ix_accounting_bank_transactions_{name}", "accounting_bank_transactions", [name]) + + +def downgrade(): + op.drop_table("accounting_bank_transactions") diff --git a/app/modules/accounting/bank_models.py b/app/modules/accounting/bank_models.py new file mode 100644 index 0000000..df71bc8 --- /dev/null +++ b/app/modules/accounting/bank_models.py @@ -0,0 +1,79 @@ +from __future__ import annotations + +from datetime import datetime, timezone + +from sqlalchemy import DateTime, Float, ForeignKey, Integer, String, Text, UniqueConstraint +from sqlalchemy.orm import Mapped, mapped_column + +from app.core.db.common import CommonBase + + +class AccountingBankTransaction(CommonBase): + __tablename__ = "accounting_bank_transactions" + __table_args__ = ( + UniqueConstraint("tenant_id", "client_id", "fingerprint", name="uq_accounting_bank_tx_fingerprint"), + ) + + 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) + source_job_id: Mapped[str | None] = mapped_column(ForeignKey("bank_statement_analysis_jobs.id", ondelete="SET NULL"), nullable=True, index=True) + fingerprint: Mapped[str] = mapped_column(String(80), nullable=False, index=True) + + statement_id: Mapped[str] = mapped_column(String(40), nullable=False, default="") + bank_name: Mapped[str] = mapped_column(String(160), nullable=False, default="") + account_number: Mapped[str] = mapped_column(String(100), nullable=False, default="") + customer_name: Mapped[str] = mapped_column(String(255), nullable=False, default="") + transaction_date: Mapped[str] = mapped_column(String(20), nullable=False, default="", index=True) + value_date: Mapped[str] = mapped_column(String(20), nullable=False, default="") + narration: Mapped[str] = mapped_column(Text, nullable=False, default="") + reference_no: Mapped[str] = mapped_column(String(180), nullable=False, default="") + transfer_reference: Mapped[str] = mapped_column(String(180), nullable=False, default="") + direction: Mapped[str] = mapped_column(String(10), nullable=False, default="") + debit: Mapped[float] = mapped_column(Float, nullable=False, default=0) + credit: Mapped[float] = mapped_column(Float, nullable=False, default=0) + amount: Mapped[float] = mapped_column(Float, nullable=False, default=0) + + auto_party: Mapped[str] = mapped_column(String(255), nullable=False, default="") + analyzer_category: Mapped[str] = mapped_column(String(255), nullable=False, default="") + analyzer_nature: Mapped[str] = mapped_column(String(160), nullable=False, default="") + analyzer_ledger: Mapped[str] = mapped_column(String(255), nullable=False, default="") + + contra_pair_id: Mapped[str] = mapped_column(String(60), nullable=False, default="", index=True) + contra_counter_bank: Mapped[str] = mapped_column(String(160), nullable=False, default="") + contra_counter_account: Mapped[str] = mapped_column(String(100), nullable=False, default="") + contra_confidence: Mapped[float] = mapped_column(Float, nullable=False, default=0) + contra_reason: Mapped[str | None] = mapped_column(Text, nullable=True) + + tally_guid: Mapped[str] = mapped_column(String(120), nullable=False, default="") + suggested_nature_id: Mapped[int | None] = mapped_column(ForeignKey("accounting_natures.id", ondelete="SET NULL"), nullable=True) + suggested_ledger_name: Mapped[str] = mapped_column(String(255), nullable=False, default="") + suggested_confidence: Mapped[int] = mapped_column(Integer, nullable=False, default=0) + suggestion_reason_json: Mapped[str | None] = mapped_column(Text, nullable=True) + suggested_voucher_type: Mapped[str] = mapped_column(String(20), nullable=False, default="") + + review_status: Mapped[str] = mapped_column(String(20), nullable=False, default="pending", index=True) + final_nature_id: Mapped[int | None] = mapped_column(ForeignKey("accounting_natures.id", ondelete="SET NULL"), nullable=True) + final_ledger_name: Mapped[str] = mapped_column(String(255), nullable=False, default="") + final_party_ledger_name: Mapped[str] = mapped_column(String(255), nullable=False, default="") + final_bank_ledger_name: Mapped[str] = mapped_column(String(255), nullable=False, default="") + final_other_bank_ledger_name: Mapped[str] = mapped_column(String(255), nullable=False, default="") + final_voucher_type: Mapped[str] = mapped_column(String(20), nullable=False, default="") + review_note: Mapped[str | None] = mapped_column(Text, nullable=True) + reviewed_by_user_id: Mapped[int | None] = mapped_column(ForeignKey("users.id", ondelete="SET NULL"), nullable=True) + reviewed_at_utc: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + + workstation_agent_id: Mapped[int | None] = mapped_column(ForeignKey("erp_workstation_agents.id", ondelete="SET NULL"), nullable=True) + preflight_job_id: Mapped[int | None] = mapped_column(ForeignKey("erp_agent_jobs.id", ondelete="SET NULL"), nullable=True) + posting_job_id: Mapped[int | None] = mapped_column(ForeignKey("erp_agent_jobs.id", ondelete="SET NULL"), nullable=True) + preflight_result_json: Mapped[str | None] = mapped_column(Text, nullable=True) + posting_result_json: Mapped[str | None] = mapped_column(Text, nullable=True) + posting_status: Mapped[str] = mapped_column(String(30), nullable=False, default="not_started", index=True) + last_error: Mapped[str | None] = mapped_column(Text, nullable=True) + tally_voucher_id: Mapped[str] = mapped_column(String(120), nullable=False, default="") + tally_voucher_number: Mapped[str] = mapped_column(String(160), nullable=False, default="") + posted_by_user_id: Mapped[int | None] = mapped_column(ForeignKey("users.id", ondelete="SET NULL"), nullable=True) + posted_at_utc: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + + created_at_utc: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=lambda: datetime.now(timezone.utc), index=True) + updated_at_utc: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=lambda: datetime.now(timezone.utc), onupdate=lambda: datetime.now(timezone.utc)) diff --git a/app/modules/accounting/bank_service.py b/app/modules/accounting/bank_service.py new file mode 100644 index 0000000..15bba97 --- /dev/null +++ b/app/modules/accounting/bank_service.py @@ -0,0 +1,412 @@ +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.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.") + + 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() + return tx + + +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 diff --git a/app/modules/accounting/bank_ui.py b/app/modules/accounting/bank_ui.py new file mode 100644 index 0000000..bee99d2 --- /dev/null +++ b/app/modules/accounting/bank_ui.py @@ -0,0 +1,123 @@ +from __future__ import annotations + +from urllib.parse import urlencode + +from fastapi import APIRouter, Form, Request +from fastapi.responses import RedirectResponse + +from app.core.db.common import CommonSessionLocal +from app.core.security.csrf import get_or_create_csrf_token, validate_csrf +from app.core.templating import templates +from app.modules.accounting.bank_service import ( + confirm_review, import_completed_job, preflight_choices, queue_post, queue_preflight, + queue_rows, sync_posting, visible_workstations, +) +from app.modules.accounting.ledger_learning_service import active_natures +from app.modules.accounting.ui import _find_visible_client, _require_partner, _visible_clients +from app.modules.bank_statement_analyzer.models import BankStatementAnalysisJob +from app.modules.core.rbac.deps import get_user_permissions, get_user_roles +from sqlalchemy import select + +router = APIRouter(prefix="/tools/accounting/bank-posting", tags=["accounting-bank-posting-ui"]) + + +def _go(client_id=0, message="", error=""): + q = {"client_id": client_id} if client_id else {} + if message: q["message"] = message[:250] + if error: q["error"] = error[:250] + return RedirectResponse("/tools/accounting/bank-posting" + ("?" + urlencode(q) if q else ""), status_code=303) + + +@router.get("") +def page(request: Request, client_id: int | None = None, status: str = "", page: int = 1, per_page: int = 25, message: str = "", error: str = ""): + db = CommonSessionLocal() + try: + user, denied = _require_partner(request, db, "accounting.learning.view") + if denied: return denied + clients, scope = _visible_clients(db, request, user) + selected = next((c for c in clients if client_id and int(c.id) == int(client_id)), None) + rows = []; total = 0; pages = 1 + choices = {}; workstations = [] + if selected: + rows, total, page, pages = queue_rows(db, tenant_id=scope.tenant_id, client_id=selected.id, status=status, page=page, per_page=per_page) + for row in rows: + sync_posting(db, row) + choices = {r.id: preflight_choices(r) for r in rows if r.preflight_result_json} + workstations = visible_workstations(db, scope.tenant_id, getattr(user, "branch_id", None)) + jobs = list(db.execute(select(BankStatementAnalysisJob).where( + BankStatementAnalysisJob.user_id == user.id, + BankStatementAnalysisJob.status == "completed", + ).order_by(BankStatementAnalysisJob.completed_at_utc.desc()).limit(30)).scalars().all()) + return templates.TemplateResponse("modules/accounting/templates/accounting/bank_posting.html", { + "request": request, "current_user": user, + "current_user_roles": get_user_roles(db, user.id), + "current_user_permissions": get_user_permissions(db, user.id), + "csrf_token": get_or_create_csrf_token(request), + "title": "Bank to Tally", "clients": clients, "selected_client": selected, + "rows": rows, "total": total, "page": page, "pages": pages, "per_page": per_page, + "status_filter": status, "jobs": jobs, "natures": active_natures(db, scope.tenant_id), + "choices": choices, "workstations": workstations, "message": message, "error": error, + }) + finally: + db.close() + + +@router.post("/import") +def import_job(request: Request, client_id: int = Form(...), job_id: 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, _, scope = _find_visible_client(db, request, user, client_id) + if not client: return _go(error="Client is not visible.") + added, skipped = import_completed_job(db, tenant_id=scope.tenant_id, client_id=client.id, job_id=job_id, user_id=user.id) + return _go(client.id, message=f"Imported {added} bank transaction(s); {skipped} duplicate fingerprint(s) skipped.") + except Exception as exc: + db.rollback(); return _go(client_id, error=str(exc)) + finally: db.close() + + +@router.post("/{tx_id}/review") +def review(request: Request, tx_id: int, client_id: int = Form(...), nature_id: int | None = Form(None), ledger_name: str = Form(""), voucher_type: str = Form(...), review_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, _, scope = _find_visible_client(db, request, user, client_id) + confirm_review(db, tx_id=tx_id, tenant_id=scope.tenant_id, client_id=client.id, nature_id=nature_id, ledger_name=ledger_name, voucher_type=voucher_type, user_id=user.id, note=review_note) + return _go(client.id, message="Bank transaction review saved.") + except Exception as exc: + db.rollback(); return _go(client_id, error=str(exc)) + finally: db.close() + + +@router.post("/{tx_id}/preflight") +def preflight(request: Request, tx_id: int, client_id: int = Form(...), 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, _, scope = _find_visible_client(db, request, user, client_id) + queue_preflight(db, tx_id=tx_id, tenant_id=scope.tenant_id, client_id=client.id, workstation_id=workstation_id, user_id=user.id) + return _go(client.id, message="Bank posting preflight queued.") + except Exception as exc: + db.rollback(); return _go(client_id, error=str(exc)) + finally: db.close() + + +@router.post("/{tx_id}/post") +def post(request: Request, tx_id: int, client_id: int = Form(...), bank_ledger_name: str = Form(...), counter_ledger_name: str = Form(""), other_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, _, scope = _find_visible_client(db, request, user, client_id) + queue_post(db, tx_id=tx_id, tenant_id=scope.tenant_id, client_id=client.id, bank_ledger_name=bank_ledger_name, counter_ledger_name=counter_ledger_name, other_bank_ledger_name=other_bank_ledger_name, user_id=user.id) + return _go(client.id, message="Controlled bank voucher job queued.") + except Exception as exc: + db.rollback(); return _go(client_id, error=str(exc)) + finally: db.close() diff --git a/app/modules/accounting/templates/accounting/bank_posting.html b/app/modules/accounting/templates/accounting/bank_posting.html new file mode 100644 index 0000000..3cd83a9 --- /dev/null +++ b/app/modules/accounting/templates/accounting/bank_posting.html @@ -0,0 +1,76 @@ +{% extends "ui/templates/base/layout.html" %} +{% block content %} +
Accounting Intelligence · Phase 11
Import a completed Bank Analyzer job, reuse the common ledger-learning engine for non-contra transactions, and post only reviewed Payment / Receipt / Contra vouchers through the durable Local Agent job system.
{{ total }} bank transaction(s). Matched multi-bank contra rows are visibly linked by pair ID.
Upload supported PDF statements. Up to three analyses run globally at one time; additional jobs are queued safely.
Upload one or several PDF statements for the same client. Multiple bank accounts can be analysed together so equal-and-opposite inter-bank transfers can be identified as probable contra. Up to three analyses run globally at one time; additional jobs are queued safely.
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' }}.
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' }}.
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