420 lines
17 KiB
Python
420 lines
17 KiB
Python
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.")
|
|
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()
|
|
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
|