Add Phase 11 bank to Tally with multi-bank contra intelligence
This commit is contained in:
@@ -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
|
||||
Reference in New Issue
Block a user