from __future__ import annotations from typing import Any from app.modules.accounting.agent_bridge import request_agent_command DEFAULT_QUERY_TIMEOUT_SECONDS = 45 DEFAULT_SYNC_TIMEOUT_SECONDS = 330 class AccountingMirrorError(RuntimeError): pass def _unwrap(response: dict[str, Any]) -> dict[str, Any]: if not isinstance(response, dict): raise AccountingMirrorError("ERP Local Agent returned an invalid response.") if not response.get("ok"): raise AccountingMirrorError(str(response.get("error") or "ERP Local Agent command failed.")) result = response.get("result") if not isinstance(result, dict): raise AccountingMirrorError("ERP Local Agent returned an invalid mirror result.") return result def mirror_status(*, node_code: str, accounting_payload: dict[str, Any]) -> dict[str, Any]: response = request_agent_command( node_code, "accounting_mirror_status", accounting_payload, timeout_seconds=DEFAULT_QUERY_TIMEOUT_SECONDS, ) return _unwrap(response) def mirror_sync( *, node_code: str, accounting_payload: dict[str, Any], company_name: str = "", company_guid: str = "", dsn: str = "TallyODBC64_9000", timeout_seconds: int = DEFAULT_SYNC_TIMEOUT_SECONDS, ) -> dict[str, Any]: payload = dict(accounting_payload or {}) payload.update( { "company_name": str(company_name or "").strip(), "company_guid": str(company_guid or "").strip(), "dsn": str(dsn or "TallyODBC64_9000").strip(), "timeout_seconds": max(60, int(timeout_seconds)), } ) response = request_agent_command( node_code, "accounting_mirror_sync", payload, timeout_seconds=max(60, int(timeout_seconds)), ) return _unwrap(response) def mirror_query( *, node_code: str, accounting_payload: dict[str, Any], query: str, filters: dict[str, Any] | None = None, timeout_seconds: int = DEFAULT_QUERY_TIMEOUT_SECONDS, ) -> dict[str, Any]: payload = dict(accounting_payload or {}) payload.update({"query": str(query or "").strip(), "filters": filters or {}}) response = request_agent_command( node_code, "accounting_mirror_query", payload, timeout_seconds=max(10, int(timeout_seconds)), ) return _unwrap(response) def daybook(*, node_code: str, accounting_payload: dict[str, Any], from_date="", to_date="", limit=500): return mirror_query( node_code=node_code, accounting_payload=accounting_payload, query="daybook", filters={"from_date": from_date, "to_date": to_date, "limit": limit}, ) def trial_balance(*, node_code: str, accounting_payload: dict[str, Any], limit=5000): return mirror_query( node_code=node_code, accounting_payload=accounting_payload, query="trial_balance", filters={"limit": limit}, ) def ledger_transactions( *, node_code: str, accounting_payload: dict[str, Any], ledger_name: str, from_date="", to_date="", limit=5000, ): return mirror_query( node_code=node_code, accounting_payload=accounting_payload, query="ledger_transactions", filters={ "ledger_name": ledger_name, "from_date": from_date, "to_date": to_date, "limit": limit, }, ) def inventory( *, node_code: str, accounting_payload: dict[str, Any], stock_item_name: str = "", from_date="", to_date="", limit=5000, ): return mirror_query( node_code=node_code, accounting_payload=accounting_payload, query="inventory", filters={ "stock_item_name": stock_item_name, "from_date": from_date, "to_date": to_date, "limit": limit, }, ) def stock_items(*, node_code: str, accounting_payload: dict[str, Any], limit=5000): return mirror_query( node_code=node_code, accounting_payload=accounting_payload, query="stock_items", filters={"limit": limit}, ) def hsn_history(*, node_code: str, accounting_payload: dict[str, Any], stock_item_name: str): return mirror_query( node_code=node_code, accounting_payload=accounting_payload, query="hsn_history", filters={"stock_item_name": stock_item_name}, ) def gst_history(*, node_code: str, accounting_payload: dict[str, Any], stock_item_name: str): return mirror_query( node_code=node_code, accounting_payload=accounting_payload, query="gst_history", filters={"stock_item_name": stock_item_name}, ) def cash_transactions( *, node_code: str, accounting_payload: dict[str, Any], ledger_name: str = "Cash", from_date="", to_date="", limit=5000, ): return mirror_query( node_code=node_code, accounting_payload=accounting_payload, query="cash_transactions", filters={ "ledger_name": ledger_name, "from_date": from_date, "to_date": to_date, "limit": limit, }, ) def voucher(*, node_code: str, accounting_payload: dict[str, Any], voucher_guid: str): return mirror_query( node_code=node_code, accounting_payload=accounting_payload, query="voucher", filters={"voucher_guid": voucher_guid}, ) def mirror_exceptions(*, node_code: str, accounting_payload: dict[str, Any], limit=5000): return mirror_query( node_code=node_code, accounting_payload=accounting_payload, query="exceptions", filters={"limit": limit}, ) def refresh_masters(*, node_code: str, accounting_payload: dict[str, Any], tally_guid: str, requested_by_user_id: int | None = None) -> dict[str, Any]: payload = dict(accounting_payload or {}) payload["tally_guid"] = str(tally_guid or "").strip() if requested_by_user_id is not None: payload["requested_by_user_id"] = int(requested_by_user_id) return _unwrap(request_agent_command( node_code, "accounting_sync_masters", payload, timeout_seconds=330 )) def refresh_transactions(*, node_code: str, accounting_payload: dict[str, Any], tally_guid: str, date_from: str, date_to: str, requested_by_user_id: int | None = None) -> dict[str, Any]: payload = dict(accounting_payload or {}) payload.update({"tally_guid": str(tally_guid or "").strip(), "date_from": str(date_from or "").strip(), "date_to": str(date_to or "").strip()}) if requested_by_user_id is not None: payload["requested_by_user_id"] = int(requested_by_user_id) return _unwrap(request_agent_command( node_code, "accounting_sync_transactions", payload, timeout_seconds=330 )) def _creditor_movements( *, node_code: str, accounting_payload: dict[str, Any], from_date: str, to_date: str, limit: int = 50000, ) -> dict[str, Any]: return mirror_query( node_code=node_code, accounting_payload=accounting_payload, query="sundry_creditor_movements", filters={ "from_date": str(from_date or "").strip(), "to_date": str(to_date or "").strip(), "limit": max(1, min(100000, int(limit or 50000))), }, timeout_seconds=90, ) def sundry_creditors_aging( *, node_code: str, accounting_payload: dict[str, Any], fy_start: str, fy_end: str, follow_up_accounting_payload: dict[str, Any] | None = None, follow_up_start: str = "", follow_up_end: str = "", limit: int = 50000, ) -> dict[str, Any]: """Age closing Sundry Creditors using FIFO from the Accounting Mirror. The selected-FY closing is reconstructed from the mirror's ledger opening balance plus all creditor movements available up to the selected FY end. FIFO means the oldest credit lots are settled first; therefore the closing unpaid balance is reconstructed from the newest surviving credit lots. Voucher dates are normalized before ageing. This is intentionally strict: a dated movement that cannot be parsed is not silently treated as an old opening balance, because doing so would incorrectly push current balances into the >180-day bucket. """ from datetime import date as _date, datetime as _datetime, timedelta as _timedelta # Pull everything available up to the selected FY end. This lets a mirror # containing more than one FY retain the original date of an older surviving # creditor lot. If the mirror starts at the selected FY, any unrepresented # residual is treated as brought forward and is necessarily >180 days at FY end. current = _creditor_movements( node_code=node_code, accounting_payload=accounting_payload, from_date="", to_date=fy_end, limit=limit, ) later = {"ledgers": [], "rows": []} if follow_up_accounting_payload is not None and follow_up_start and follow_up_end: later = _creditor_movements( node_code=node_code, accounting_payload=follow_up_accounting_payload, from_date=follow_up_start, to_date=follow_up_end, limit=limit, ) if current.get("truncated"): raise AccountingMirrorError( "Sundry Creditor movement extraction reached the safety limit. " "Increase the mirror query limit before relying on this analysis." ) if int(current.get("date_parse_error_count") or 0) > 0: sample = (current.get("date_parse_errors") or [{}])[0] raise AccountingMirrorError( "Accounting Mirror contains creditor voucher date(s) that cannot be normalized. " f"First affected entry: {sample.get('ledger_name') or 'Unknown ledger'} / " f"{sample.get('voucher_number') or 'No voucher number'} / " f"{sample.get('voucher_date') or 'blank date'}. Refresh the Accounting Mirror and retry." ) if later.get("truncated"): raise AccountingMirrorError( "Follow-up Sundry Creditor movement extraction reached the safety limit. " "Increase the mirror query limit before relying on the subsequent-payment analysis." ) if int(later.get("date_parse_error_count") or 0) > 0: sample = (later.get("date_parse_errors") or [{}])[0] raise AccountingMirrorError( "Follow-up Accounting Mirror contains creditor voucher date(s) that cannot be normalized. " f"First affected entry: {sample.get('ledger_name') or 'Unknown ledger'} / " f"{sample.get('voucher_number') or 'No voucher number'} / " f"{sample.get('voucher_date') or 'blank date'}. Refresh the follow-up Accounting Mirror and retry." ) fy_start_date = _date.fromisoformat(fy_start) fy_end_date = _date.fromisoformat(fy_end) opening_fallback_date = fy_start_date - _timedelta(days=1) def _money(value: Any) -> float: try: return round(float(value or 0), 2) except (TypeError, ValueError): return 0.0 def _parse_date(value: Any) -> _date | None: text = str(value or "").strip() if not text: return None # Fast path for the current .act ISO format. try: return _date.fromisoformat(text[:10]) except ValueError: pass candidates = ( "%Y%m%d", "%d-%m-%Y", "%d/%m/%Y", "%d-%b-%Y", "%d-%b-%y", "%m/%d/%Y", "%m/%d/%Y %H:%M:%S", "%Y-%m-%d %H:%M:%S", ) for fmt in candidates: try: return _datetime.strptime(text, fmt).date() except ValueError: continue return None def _iso_date(value: Any) -> str: parsed = _parse_date(value) return parsed.isoformat() if parsed else "" ledger_meta = { str(x.get("ledger_name") or "").casefold(): x for x in (current.get("ledgers") or []) if str(x.get("ledger_name") or "").strip() } by_ledger: dict[str, list[dict[str, Any]]] = {} invalid_dates: list[dict[str, str]] = [] for raw_row in current.get("rows") or []: row = dict(raw_row) name = str(row.get("ledger_name") or "").strip() if not name: continue parsed = _parse_date(row.get("voucher_date")) if parsed is None: invalid_dates.append({ "party": name, "voucher_number": str(row.get("voucher_number") or ""), "voucher_date": str(row.get("voucher_date") or ""), }) continue if parsed > fy_end_date: continue row["voucher_date"] = parsed.isoformat() row["_parsed_date"] = parsed by_ledger.setdefault(name.casefold(), []).append(row) if invalid_dates: first = invalid_dates[0] raise AccountingMirrorError( "Accounting Mirror contains creditor voucher date(s) that could not be parsed. " f"First affected entry: {first['party']} / {first['voucher_number']} / " f"{first['voucher_date'] or 'blank date'}. Refresh the Accounting Mirror and retry." ) later_by_ledger: dict[str, list[dict[str, Any]]] = {} invalid_later_dates: list[dict[str, str]] = [] for raw_row in later.get("rows") or []: row = dict(raw_row) name = str(row.get("ledger_name") or "").strip() if not name: continue parsed = _parse_date(row.get("voucher_date")) if parsed is None: invalid_later_dates.append({ "party": name, "voucher_number": str(row.get("voucher_number") or ""), "voucher_date": str(row.get("voucher_date") or ""), }) continue row["voucher_date"] = parsed.isoformat() row["_parsed_date"] = parsed later_by_ledger.setdefault(name.casefold(), []).append(row) if invalid_later_dates: first = invalid_later_dates[0] raise AccountingMirrorError( "Follow-up Accounting Mirror contains creditor voucher date(s) that could not be parsed. " f"First affected entry: {first['party']} / {first['voucher_number']} / " f"{first['voucher_date'] or 'blank date'}. Refresh the follow-up Accounting Mirror and retry." ) parties: list[dict[str, Any]] = [] details: list[dict[str, Any]] = [] party_names = sorted( { str(x.get("ledger_name") or "").strip() for x in (current.get("ledgers") or []) if str(x.get("ledger_name") or "").strip() }, key=str.casefold, ) for party_name in party_names: key = party_name.casefold() meta = ledger_meta.get(key) or {} rows = sorted( by_ledger.get(key, []), key=lambda r: (r.get("_parsed_date"), int(r.get("line_no") or 0)), ) # Preserve the selected-FY closing-balance control already validated in v2: # ledger opening + movements inside the selected FY. Earlier rows, when the # mirror has them, are used only to identify the original date/composition of # brought-forward FIFO lots; they are not added again to the closing balance. opening_balance = _money(meta.get("opening_balance")) fy_rows = [ row for row in rows if fy_start_date <= row.get("_parsed_date") <= fy_end_date ] period_credit = round(sum( _money(row.get("amount")) for row in fy_rows if str(row.get("dr_cr") or "").upper() == "CR" ), 2) period_debit = round(sum( _money(row.get("amount")) for row in fy_rows if str(row.get("dr_cr") or "").upper() == "DR" ), 2) historical_closing = round(opening_balance + period_credit - period_debit, 2) closing_balance = max(0.0, historical_closing) if closing_balance <= 0.009: continue # True forward FIFO reconstruction. The selected-FY opening balance is # the oldest outstanding lot. Every creditor-increasing movement creates # a new dated lot; every creditor-reducing movement settles the oldest lot # first. The lots that survive at FY end are therefore the actual FIFO # composition of the closing balance. closing_lots: list[dict[str, Any]] = [] if opening_balance > 0.009: closing_lots.append({ "party_name": party_name, "source": "Opening / brought forward", "voucher_date": opening_fallback_date.isoformat(), "voucher_type": "Opening", "voucher_number": "", "reference": "", "original_credit": opening_balance, "remaining": opening_balance, "outstanding_at_fy_end": 0.0, "paid_subsequently": 0.0, "final_payment_date": "", }) def settle_fifo(lots: list[dict[str, Any]], amount: float) -> None: remaining_payment = max(0.0, _money(amount)) for lot in lots: if remaining_payment <= 0.009: break available = max(0.0, _money(lot.get("remaining"))) if available <= 0.009: continue applied = round(min(available, remaining_payment), 2) lot["remaining"] = round(available - applied, 2) remaining_payment = round(remaining_payment - applied, 2) for row in fy_rows: side = str(row.get("dr_cr") or "").strip().upper() amount = _money(row.get("amount")) if amount <= 0.009: continue if side == "CR": closing_lots.append({ "party_name": party_name, "source": "Voucher", "voucher_date": str(row.get("voucher_date") or ""), "voucher_type": str(row.get("voucher_type") or ""), "voucher_number": str(row.get("voucher_number") or ""), "reference": str(row.get("reference") or ""), "original_credit": amount, "remaining": amount, "outstanding_at_fy_end": 0.0, "paid_subsequently": 0.0, "final_payment_date": "", }) elif side == "DR": settle_fifo(closing_lots, amount) # Remove fully settled lots and freeze the surviving FY-end outstanding # amounts before any optional follow-up-year settlement is applied. closing_lots = [ lot for lot in closing_lots if _money(lot.get("remaining")) > 0.009 ] for lot in closing_lots: lot["outstanding_at_fy_end"] = _money(lot.get("remaining")) fifo_total = round(sum(_money(lot.get("remaining")) for lot in closing_lots), 2) if abs(fifo_total - closing_balance) > 0.01: raise AccountingMirrorError( f"FIFO movement reconstruction does not reconcile for {party_name}: " f"selected-FY closing {closing_balance:.2f}, FIFO lots {fifo_total:.2f}. " "Refresh the Accounting Mirror for the selected FY and retry." ) closing_lots = sorted( closing_lots, key=lambda lot: ( _parse_date(lot.get("voucher_date")) or opening_fallback_date, str(lot.get("voucher_number") or ""), ), ) def settle_follow_up(amount: float, paid_on: str) -> None: remaining_payment = max(0.0, _money(amount)) for lot in closing_lots: if remaining_payment <= 0.009: break available = max(0.0, _money(lot.get("remaining"))) if available <= 0.009: continue applied = round(min(available, remaining_payment), 2) lot["remaining"] = round(available - applied, 2) lot["paid_subsequently"] = round(_money(lot.get("paid_subsequently")) + applied, 2) remaining_payment = round(remaining_payment - applied, 2) if _money(lot.get("remaining")) <= 0.009: lot["final_payment_date"] = paid_on for row in sorted( later_by_ledger.get(key, []), key=lambda r: (r.get("_parsed_date"), int(r.get("line_no") or 0)), ): if str(row.get("dr_cr") or "").upper() == "DR": settle_follow_up(_money(row.get("amount")), str(row.get("voucher_date") or "")) within_180 = 0.0 over_180 = 0.0 paid_later = 0.0 still_unpaid = 0.0 over_180_lots: list[dict[str, Any]] = [] for lot in closing_lots: lot_date = _parse_date(lot.get("voucher_date")) or opening_fallback_date age_days = max(0, (fy_end_date - lot_date).days) bucket = ">180 Days" if age_days > 180 else "≤180 Days" outstanding = _money(lot.get("outstanding_at_fy_end")) paid = round(min(outstanding, _money(lot.get("paid_subsequently"))), 2) after = round(max(0.0, outstanding - paid), 2) if bucket == ">180 Days": over_180 = round(over_180 + outstanding, 2) paid_later = round(paid_later + paid, 2) still_unpaid = round(still_unpaid + after, 2) over_180_lots.append(lot) else: within_180 = round(within_180 + outstanding, 2) details.append({ "party_name": lot.get("party_name"), "source": lot.get("source"), "voucher_date": lot.get("voucher_date"), "voucher_type": lot.get("voucher_type"), "voucher_number": lot.get("voucher_number"), "reference": lot.get("reference"), "original_credit": lot.get("original_credit"), "outstanding_at_fy_end": outstanding, "age_days": age_days, "age_bucket": bucket, "paid_subsequently": paid, "balance_after_follow_up": after, "final_payment_date": lot.get("final_payment_date"), "allocation_basis": "FIFO from dated Accounting Mirror creditor movements", }) final_payment_date = "" if over_180_lots and all(_money(lot.get("remaining")) <= 0.009 for lot in over_180_lots): dates = [ str(lot.get("final_payment_date") or "") for lot in over_180_lots if lot.get("final_payment_date") ] if dates: final_payment_date = max(dates) reconstructed = round(within_180 + over_180, 2) variance = round(closing_balance - reconstructed, 2) if abs(variance) > 0.01: raise AccountingMirrorError( f"FIFO reconstruction does not reconcile for {party_name}: " f"selected-FY closing {closing_balance:.2f}, reconstructed {reconstructed:.2f}." ) parties.append({ "party_name": party_name, "within_180": within_180, "over_180": over_180, "closing_balance": closing_balance, "over_180_paid_later": paid_later, "over_180_still_unpaid": still_unpaid, "final_payment_date": final_payment_date, }) summary = { "party_count": len(parties), "total_closing": round(sum(x["closing_balance"] for x in parties), 2), "within_180": round(sum(x["within_180"] for x in parties), 2), "over_180": round(sum(x["over_180"] for x in parties), 2), "over_180_paid_later": round(sum(x["over_180_paid_later"] for x in parties), 2), "over_180_still_unpaid": round(sum(x["over_180_still_unpaid"] for x in parties), 2), } return { "summary": summary, "parties": parties, "details": details, "allocation_basis": ( "FIFO from dated Accounting Mirror movements; historical closing is " "reconstructed from ledger opening plus selected-FY movements, " "and debit-side settlements clear the oldest outstanding lots first" ), "current_mirror": current.get("mirror") or {}, "follow_up_mirror": later.get("mirror") or {}, }