Add Phase 5 historical Tally learning foundation
This commit is contained in:
@@ -0,0 +1,168 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from collections import defaultdict
|
||||
from datetime import datetime, timezone
|
||||
|
||||
from sqlalchemy import delete, select
|
||||
|
||||
from app.modules.accounting.historical_learning_models import (
|
||||
AccountingHistoricalLedgerEvidence, AccountingHistoricalLearningRun, AccountingLedgerNatureMapping,
|
||||
)
|
||||
from app.modules.accounting.taxonomy_models import AccountingNature
|
||||
from app.modules.documents.models import ERPAgentJob
|
||||
|
||||
|
||||
NOISE_GROUPS = {
|
||||
"bank accounts", "bank od a/c", "cash-in-hand", "cash in hand", "sundry creditors", "sundry debtors",
|
||||
"duties & taxes", "duties and taxes", "loans (liability)", "secured loans", "unsecured loans",
|
||||
}
|
||||
|
||||
|
||||
def _utcnow():
|
||||
return datetime.now(timezone.utc)
|
||||
|
||||
|
||||
def _json(value):
|
||||
try:
|
||||
return json.loads(value or "{}")
|
||||
except Exception:
|
||||
return {}
|
||||
|
||||
|
||||
def active_natures(db, tenant_id: int):
|
||||
return list(db.execute(select(AccountingNature).where(
|
||||
AccountingNature.tenant_id == tenant_id, AccountingNature.is_active.is_(True), AccountingNature.is_posting_nature.is_(True)
|
||||
).order_by(AccountingNature.sort_order, AccountingNature.name)).scalars().all())
|
||||
|
||||
|
||||
def ledger_mappings(db, tenant_id: int, client_id: int, tally_guid: str = ""):
|
||||
stmt = select(AccountingLedgerNatureMapping).where(
|
||||
AccountingLedgerNatureMapping.tenant_id == tenant_id, AccountingLedgerNatureMapping.client_id == client_id
|
||||
)
|
||||
if tally_guid:
|
||||
stmt = stmt.where(AccountingLedgerNatureMapping.tally_guid == tally_guid)
|
||||
return list(db.execute(stmt.order_by(AccountingLedgerNatureMapping.ledger_name)).scalars().all())
|
||||
|
||||
|
||||
def save_ledger_mapping(db, *, tenant_id: int, client_id: int, tally_guid: str, company_name: str, ledger_name: str, parent_group: str, nature_id: int, user_id: int):
|
||||
nature = db.execute(select(AccountingNature).where(
|
||||
AccountingNature.id == nature_id, AccountingNature.tenant_id == tenant_id, AccountingNature.is_active.is_(True), AccountingNature.is_posting_nature.is_(True)
|
||||
)).scalar_one_or_none()
|
||||
if not nature:
|
||||
raise ValueError("Select an active posting accounting nature.")
|
||||
row = db.execute(select(AccountingLedgerNatureMapping).where(
|
||||
AccountingLedgerNatureMapping.tenant_id == tenant_id, AccountingLedgerNatureMapping.client_id == client_id,
|
||||
AccountingLedgerNatureMapping.tally_guid == tally_guid, AccountingLedgerNatureMapping.ledger_name == ledger_name
|
||||
)).scalar_one_or_none()
|
||||
if not row:
|
||||
row = AccountingLedgerNatureMapping(tenant_id=tenant_id, client_id=client_id, tally_guid=tally_guid, ledger_name=ledger_name, nature_id=nature.id)
|
||||
row.company_name = company_name or row.company_name or ""
|
||||
row.parent_group = parent_group or row.parent_group or ""
|
||||
row.nature_id = nature.id
|
||||
row.source = "manual"
|
||||
row.confidence_percent = 100
|
||||
row.confirmed_by_user_id = user_id
|
||||
row.confirmed_at_utc = _utcnow()
|
||||
row.updated_at_utc = _utcnow()
|
||||
db.add(row); db.commit(); db.refresh(row)
|
||||
return row
|
||||
|
||||
|
||||
def remove_ledger_mapping(db, *, tenant_id: int, client_id: int, mapping_id: int):
|
||||
row = db.execute(select(AccountingLedgerNatureMapping).where(
|
||||
AccountingLedgerNatureMapping.id == mapping_id, AccountingLedgerNatureMapping.tenant_id == tenant_id, AccountingLedgerNatureMapping.client_id == client_id
|
||||
)).scalar_one_or_none()
|
||||
if row:
|
||||
db.delete(row); db.commit()
|
||||
return row
|
||||
|
||||
|
||||
def ingest_completed_run(db, run: AccountingHistoricalLearningRun):
|
||||
if run.status == "completed" or not run.agent_job_id:
|
||||
return run
|
||||
job = db.get(ERPAgentJob, run.agent_job_id)
|
||||
if not job:
|
||||
run.status = "failed"; run.error_message = "Agent job no longer exists."; run.completed_at_utc = _utcnow(); db.commit(); return run
|
||||
if job.status in {"queued", "claimed"}:
|
||||
run.status = job.status; db.commit(); return run
|
||||
if job.status != "succeeded":
|
||||
run.status = "failed"; run.error_message = job.last_error or "Historical evidence collection failed."; run.completed_at_utc = _utcnow(); db.commit(); return run
|
||||
result = _json(job.result_json)
|
||||
rows = result.get("evidence") or []
|
||||
company_name = str(result.get("company_name") or run.company_name or "")
|
||||
now = _utcnow()
|
||||
# Snapshot semantics for this company: remove old aggregated rows then replace with current evidence.
|
||||
db.execute(delete(AccountingHistoricalLedgerEvidence).where(
|
||||
AccountingHistoricalLedgerEvidence.tenant_id == run.tenant_id, AccountingHistoricalLedgerEvidence.client_id == run.client_id,
|
||||
AccountingHistoricalLedgerEvidence.tally_guid == run.tally_guid
|
||||
))
|
||||
inserted = 0
|
||||
for item in rows:
|
||||
party = str(item.get("party_ledger_name") or "").strip()
|
||||
counter = str(item.get("counter_ledger_name") or "").strip()
|
||||
if not party or not counter or party.casefold() == counter.casefold():
|
||||
continue
|
||||
db.add(AccountingHistoricalLedgerEvidence(
|
||||
tenant_id=run.tenant_id, client_id=run.client_id, tally_guid=run.tally_guid, company_name=company_name,
|
||||
party_ledger_name=party, party_parent_group=str(item.get("party_parent_group") or ""),
|
||||
counter_ledger_name=counter, counter_parent_group=str(item.get("counter_parent_group") or ""),
|
||||
voucher_type_name=str(item.get("voucher_type_name") or "Purchase"), voucher_count=int(item.get("voucher_count") or 0),
|
||||
absolute_amount_total=float(item.get("absolute_amount_total") or 0), first_voucher_date=str(item.get("first_voucher_date") or "") or None,
|
||||
last_voucher_date=str(item.get("last_voucher_date") or "") or None, sample_narration=str(item.get("sample_narration") or "")[:2000] or None, refreshed_at_utc=now
|
||||
)); inserted += 1
|
||||
run.company_name = company_name
|
||||
run.status = "completed"
|
||||
run.evidence_rows = inserted
|
||||
run.error_message = None
|
||||
run.completed_at_utc = now
|
||||
db.commit(); db.refresh(run)
|
||||
return run
|
||||
|
||||
|
||||
def evidence_rows(db, tenant_id: int, client_id: int):
|
||||
return list(db.execute(select(AccountingHistoricalLedgerEvidence).where(
|
||||
AccountingHistoricalLedgerEvidence.tenant_id == tenant_id, AccountingHistoricalLedgerEvidence.client_id == client_id
|
||||
).order_by(AccountingHistoricalLedgerEvidence.voucher_count.desc(), AccountingHistoricalLedgerEvidence.party_ledger_name, AccountingHistoricalLedgerEvidence.counter_ledger_name)).scalars().all())
|
||||
|
||||
|
||||
def classifiable_ledgers(db, tenant_id: int, client_id: int):
|
||||
rows = evidence_rows(db, tenant_id, client_id)
|
||||
seen = {}
|
||||
for r in rows:
|
||||
group = (r.counter_parent_group or "").strip().casefold()
|
||||
if group in NOISE_GROUPS:
|
||||
continue
|
||||
key = (r.tally_guid, r.counter_ledger_name.casefold())
|
||||
item = seen.setdefault(key, {"tally_guid": r.tally_guid, "company_name": r.company_name, "ledger_name": r.counter_ledger_name, "parent_group": r.counter_parent_group, "voucher_count": 0, "amount": 0.0})
|
||||
item["voucher_count"] += int(r.voucher_count or 0); item["amount"] += float(r.absolute_amount_total or 0)
|
||||
mappings = {(m.tally_guid, m.ledger_name.casefold()): m for m in ledger_mappings(db, tenant_id, client_id)}
|
||||
for key, item in seen.items(): item["mapping"] = mappings.get(key)
|
||||
return sorted(seen.values(), key=lambda x: (-x["voucher_count"], x["ledger_name"].casefold()))
|
||||
|
||||
|
||||
def party_suggestions(db, tenant_id: int, client_id: int):
|
||||
rows = evidence_rows(db, tenant_id, client_id)
|
||||
mappings = {(m.tally_guid, m.ledger_name.casefold()): m for m in ledger_mappings(db, tenant_id, client_id)}
|
||||
natures = {n.id: n for n in db.execute(select(AccountingNature).where(AccountingNature.tenant_id == tenant_id)).scalars().all()}
|
||||
by_party = defaultdict(lambda: defaultdict(lambda: {"count": 0, "amount": 0.0, "ledgers": set()}))
|
||||
party_meta = {}
|
||||
for r in rows:
|
||||
m = mappings.get((r.tally_guid, r.counter_ledger_name.casefold()))
|
||||
if not m:
|
||||
continue
|
||||
key=(r.tally_guid, r.party_ledger_name)
|
||||
party_meta[key]=(r.company_name, r.party_parent_group)
|
||||
bucket=by_party[key][m.nature_id]
|
||||
bucket["count"] += int(r.voucher_count or 0); bucket["amount"] += float(r.absolute_amount_total or 0); bucket["ledgers"].add(r.counter_ledger_name)
|
||||
out=[]
|
||||
for key, nature_buckets in by_party.items():
|
||||
total=sum(v["count"] for v in nature_buckets.values())
|
||||
if total <= 0: continue
|
||||
nature_id, best=max(nature_buckets.items(), key=lambda kv:(kv[1]["count"], kv[1]["amount"]))
|
||||
nature=natures.get(nature_id)
|
||||
if not nature: continue
|
||||
confidence=round(best["count"]*100/total)
|
||||
company, parent=party_meta[key]
|
||||
out.append({"tally_guid": key[0], "company_name": company, "party_ledger_name": key[1], "party_parent_group": parent, "nature": nature, "voucher_count": best["count"], "mapped_voucher_count": total, "confidence": confidence, "amount": best["amount"], "ledgers": sorted(best["ledgers"])})
|
||||
return sorted(out, key=lambda x:(-x["confidence"], -x["voucher_count"], x["party_ledger_name"].casefold()))
|
||||
Reference in New Issue
Block a user