from __future__ import annotations import json import re from collections import defaultdict from datetime import datetime, timezone from sqlalchemy import delete, func, select from app.modules.accounting.sales_learning_models import ( AccountingSalesCustomerMapping, AccountingSalesHistoricalEvidence, AccountingSalesHistoricalRun, ) from app.modules.accounting.sales_models import AccountingSalesIncomeTransaction from app.modules.accounting.taxonomy_models import AccountingNature from app.modules.documents.models import ERPAgentJob NOISE_COUNTER_GROUPS = { "sundry debtors", "cash-in-hand", "cash in hand", "bank accounts", "duties & taxes", "duties and taxes", "sales accounts", } def _utcnow(): return datetime.now(timezone.utc) def _s(value): return str(value or "").strip() def normalize_party(value: str) -> str: text = re.sub(r"[^A-Z0-9]+", " ", _s(value).upper()).strip() noise = { "PRIVATE", "PVT", "LIMITED", "LTD", "LLP", "THE", "INDIA", "M", "S", "MS", "MR", "MRS", } return " ".join(token for token in text.split() if token not in noise) def identity_key(customer_name: str, customer_gstin: str) -> str: gstin = re.sub(r"\s+", "", _s(customer_gstin).upper()) if gstin: return "GSTIN:" + gstin return "NAME:" + normalize_party(customer_name) def sales_mappings(db, *, tenant_id: int, client_id: int): return list(db.execute( select(AccountingSalesCustomerMapping).where( AccountingSalesCustomerMapping.tenant_id == int(tenant_id), AccountingSalesCustomerMapping.client_id == int(client_id), ).order_by( AccountingSalesCustomerMapping.confidence_percent.desc(), AccountingSalesCustomerMapping.customer_name, ) ).scalars().all()) def mapping_for_customer(db, *, tenant_id: int, client_id: int, customer_name: str, customer_gstin: str): key = identity_key(customer_name, customer_gstin) if key in {"NAME:", "GSTIN:"}: return None return db.execute( select(AccountingSalesCustomerMapping).where( AccountingSalesCustomerMapping.tenant_id == int(tenant_id), AccountingSalesCustomerMapping.client_id == int(client_id), AccountingSalesCustomerMapping.identity_key == key, ) ).scalar_one_or_none() def save_review_mapping( db, *, tenant_id: int, client_id: int, customer_name: str, customer_gstin: str, nature_id: int, sales_ledger_name: str, tally_guid: str = "", tally_customer_ledger_name: str = "", user_id: int, amount: float = 0, invoice_date: str = "", ): key = identity_key(customer_name, customer_gstin) if key in {"NAME:", "GSTIN:"}: return None row = db.execute( select(AccountingSalesCustomerMapping).where( AccountingSalesCustomerMapping.tenant_id == int(tenant_id), AccountingSalesCustomerMapping.client_id == int(client_id), AccountingSalesCustomerMapping.identity_key == key, ) ).scalar_one_or_none() if not row: row = AccountingSalesCustomerMapping( tenant_id=int(tenant_id), client_id=int(client_id), identity_key=key, customer_name=_s(customer_name), normalized_customer_name=normalize_party(customer_name), customer_gstin=re.sub(r"\s+", "", _s(customer_gstin).upper()), evidence_count=0, ) row.customer_name = _s(customer_name) or row.customer_name row.normalized_customer_name = normalize_party(row.customer_name) row.customer_gstin = re.sub(r"\s+", "", _s(customer_gstin).upper()) or row.customer_gstin row.nature_id = int(nature_id) row.sales_ledger_name = _s(sales_ledger_name) row.tally_guid = _s(tally_guid) or row.tally_guid row.tally_customer_ledger_name = _s(tally_customer_ledger_name) or row.tally_customer_ledger_name row.source = "review" row.evidence_count = int(row.evidence_count or 0) + 1 row.total_amount = float(row.total_amount or 0) + abs(float(amount or 0)) row.last_seen_date = _s(invoice_date) or row.last_seen_date row.confidence_percent = min(100, max(90, 70 + row.evidence_count * 5)) 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 ingest_completed_sales_run(db, run: AccountingSalesHistoricalRun): 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 Sales evidence collection failed." run.completed_at_utc = _utcnow() db.commit() return run try: result = json.loads(job.result_json or "{}") except Exception: result = {} rows = result.get("evidence") or [] now = _utcnow() company_name = _s(result.get("company_name") or run.company_name) # Sales evidence is stored in its own table; Purchase history is untouched. db.execute( delete(AccountingSalesHistoricalEvidence).where( AccountingSalesHistoricalEvidence.tenant_id == run.tenant_id, AccountingSalesHistoricalEvidence.client_id == run.client_id, AccountingSalesHistoricalEvidence.tally_guid == run.tally_guid, ) ) inserted = 0 for item in rows: voucher_type = _s(item.get("voucher_type_name")) if "sales" not in voucher_type.casefold(): continue party = _s(item.get("party_ledger_name")) counter = _s(item.get("counter_ledger_name")) if not party or not counter or party.casefold() == counter.casefold(): continue db.add( AccountingSalesHistoricalEvidence( 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=_s(item.get("party_parent_group")), counter_ledger_name=counter, counter_parent_group=_s(item.get("counter_parent_group")), voucher_type_name=voucher_type or "Sales", voucher_count=int(item.get("voucher_count") or 0), absolute_amount_total=float(item.get("absolute_amount_total") or 0), first_voucher_date=_s(item.get("first_voucher_date")) or None, last_voucher_date=_s(item.get("last_voucher_date")) or None, sample_narration=_s(item.get("sample_narration"))[: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 historical_sales_rows(db, *, tenant_id: int, client_id: int): return list(db.execute( select(AccountingSalesHistoricalEvidence).where( AccountingSalesHistoricalEvidence.tenant_id == int(tenant_id), AccountingSalesHistoricalEvidence.client_id == int(client_id), ).order_by( AccountingSalesHistoricalEvidence.voucher_count.desc(), AccountingSalesHistoricalEvidence.party_ledger_name, AccountingSalesHistoricalEvidence.counter_ledger_name, ) ).scalars().all()) def historical_party_rank(db, *, tenant_id: int, client_id: int, customer_name: str): target = normalize_party(customer_name) if not target: return [] candidates = [] for row in historical_sales_rows(db, tenant_id=tenant_id, client_id=client_id): if normalize_party(row.party_ledger_name) != target: continue parent = _s(row.counter_parent_group).casefold() # In sales vouchers the useful counter ledger is normally the sales/income ledger. if parent in {"sundry debtors", "bank accounts", "cash-in-hand", "cash in hand"}: continue candidates.append(row) return sorted( candidates, key=lambda row: (-int(row.voucher_count or 0), -float(row.absolute_amount_total or 0)), ) def suggestion_for_customer(db, *, tenant_id: int, client_id: int, customer_name: str, customer_gstin: str): mapping = mapping_for_customer( db, tenant_id=tenant_id, client_id=client_id, customer_name=customer_name, customer_gstin=customer_gstin, ) if mapping and mapping.nature_id: nature = db.get(AccountingNature, mapping.nature_id) if nature and nature.is_active: return { "nature": nature, "ledger_name": mapping.sales_ledger_name, "customer_ledger_name": mapping.tally_customer_ledger_name, "confidence": int(mapping.confidence_percent or 0), "reason": ( f"Confirmed customer mapping from {int(mapping.evidence_count or 0)} " f"reviewed transaction(s)." ), "source": "confirmed_customer_mapping", } history = historical_party_rank( db, tenant_id=tenant_id, client_id=client_id, customer_name=customer_name, ) if history: best = history[0] # Historical Tally alone knows the sales ledger but not canonical nature. # Try an existing mapping that uses this exact sales ledger. linked = db.execute( select(AccountingSalesCustomerMapping).where( AccountingSalesCustomerMapping.tenant_id == int(tenant_id), AccountingSalesCustomerMapping.client_id == int(client_id), func.lower(AccountingSalesCustomerMapping.sales_ledger_name) == best.counter_ledger_name.casefold(), AccountingSalesCustomerMapping.nature_id.is_not(None), ).order_by(AccountingSalesCustomerMapping.confidence_percent.desc()).limit(1) ).scalar_one_or_none() nature = db.get(AccountingNature, linked.nature_id) if linked and linked.nature_id else None return { "nature": nature, "ledger_name": best.counter_ledger_name, "customer_ledger_name": best.party_ledger_name, "confidence": min(88, 55 + int(best.voucher_count or 0) * 3), "reason": ( f"Historical Tally Sales evidence: {int(best.voucher_count or 0)} voucher(s) " f"used ledger '{best.counter_ledger_name}'." ), "source": "historical_tally_sales", } return None def learning_summary(db, *, tenant_id: int, client_id: int): mappings = sales_mappings(db, tenant_id=tenant_id, client_id=client_id) history = historical_sales_rows(db, tenant_id=tenant_id, client_id=client_id) return { "mapping_count": len(mappings), "historical_rows": len(history), "review_confirmed": sum(1 for row in mappings if row.source == "review"), "high_confidence": sum(1 for row in mappings if int(row.confidence_percent or 0) >= 90), }