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()))