Files
arrr-erp/app/modules/accounting/accounting_mirror_service.py
T
2026-09-09 16:17:08 +05:30

540 lines
19 KiB
Python

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.
FIFO means debit-side settlements consume the oldest creditor balance first.
At a reporting date, the unpaid closing balance is therefore represented by
the *latest* credit-side additions, with any residual carried from opening.
We reconstruct that closing composition backwards from the authoritative
ledger closing balance. This is algebraically the same FIFO result as a
complete forward replay, but it is materially safer for an extracted mirror:
it cannot mis-age a closing balance merely because the opening balance is a
brought-forward aggregate rather than individual bills.
Tally bill-allocation references are not yet persisted in the .act mirror, so
this remains a transparent FIFO analysis rather than Agst Ref/New Ref matching.
"""
from datetime import date as _date, timedelta as _timedelta
current = _creditor_movements(
node_code=node_code,
accounting_payload=accounting_payload,
from_date=fy_start,
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 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."
)
fy_start_date = _date.fromisoformat(fy_start)
fy_end_date = _date.fromisoformat(fy_end)
opening_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 _row_date(row: dict[str, Any]) -> str:
return str(row.get("voucher_date") or "").strip()
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]]] = {}
for row in current.get("rows") or []:
name = str(row.get("ledger_name") or "").strip()
if name:
by_ledger.setdefault(name.casefold(), []).append(row)
later_by_ledger: dict[str, list[dict[str, Any]]] = {}
for row in later.get("rows") or []:
name = str(row.get("ledger_name") or "").strip()
if name:
later_by_ledger.setdefault(name.casefold(), []).append(row)
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 {}
# Sundry Creditor liability balances in the canonical mirror are credit
# balances (positive). Debit/zero closing balances are deliberately not
# reported as closing creditors.
closing_balance = max(0.0, _money(meta.get("closing_balance")))
if closing_balance <= 0.009:
continue
rows = by_ledger.get(key, [])
credit_rows = [
row for row in rows
if str(row.get("dr_cr") or "").upper() == "CR" and _money(row.get("amount")) > 0.009
]
# Reverse reconstruction of the closing balance under FIFO:
# oldest balances are settled first, therefore the newest credits remain.
amount_to_allocate = closing_balance
lots_newest_first: list[dict[str, Any]] = []
for row in sorted(
credit_rows,
key=lambda r: (_row_date(r), int(r.get("line_no") or 0)),
reverse=True,
):
if amount_to_allocate <= 0.009:
break
source_amount = _money(row.get("amount"))
allocated = round(min(source_amount, amount_to_allocate), 2)
if allocated <= 0.009:
continue
lots_newest_first.append({
"party_name": party_name,
"source": "Voucher",
"voucher_date": _row_date(row),
"voucher_type": str(row.get("voucher_type") or ""),
"voucher_number": str(row.get("voucher_number") or ""),
"reference": str(row.get("reference") or ""),
"original_credit": source_amount,
"remaining": allocated,
"outstanding_at_fy_end": allocated,
"paid_subsequently": 0.0,
"final_payment_date": "",
})
amount_to_allocate = round(amount_to_allocate - allocated, 2)
# Anything not represented by current-year credits is necessarily a
# brought-forward closing component and is old at this FY close.
if amount_to_allocate > 0.009:
lots_newest_first.append({
"party_name": party_name,
"source": "Opening / brought forward",
"voucher_date": opening_date.isoformat(),
"voucher_type": "Opening",
"voucher_number": "",
"reference": "",
"original_credit": amount_to_allocate,
"remaining": amount_to_allocate,
"outstanding_at_fy_end": amount_to_allocate,
"paid_subsequently": 0.0,
"final_payment_date": "",
})
amount_to_allocate = 0.0
# FIFO settlement in the follow-up period must consume the oldest FY-end
# lots first. Later credits are new liabilities and do not change which
# FY-end lot a FIFO debit settles.
closing_lots = sorted(
lots_newest_first,
key=lambda lot: (str(lot.get("voucher_date") or ""), 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: (_row_date(r), int(r.get("line_no") or 0)),
):
if str(row.get("dr_cr") or "").upper() == "DR":
settle_follow_up(_money(row.get("amount")), _row_date(row))
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:
try:
lot_date = _date.fromisoformat(str(lot.get("voucher_date") or ""))
except ValueError:
# A malformed/missing date must not be classified as current.
lot_date = opening_date
age_days = (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({
**{
k: lot.get(k)
for k in (
"party_name", "source", "voucher_date", "voucher_type",
"voucher_number", "reference", "original_credit", "final_payment_date",
)
},
"outstanding_at_fy_end": outstanding,
"age_days": age_days,
"age_bucket": bucket,
"paid_subsequently": paid,
"balance_after_follow_up": after,
"allocation_basis": "FIFO reconstructed from Accounting Mirror closing balance",
})
# Date on which the party's >180 FY-end component became fully paid.
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)
# Maintain the mirror ledger closing balance as the authoritative control.
variance = round(closing_balance - reconstructed, 2)
if abs(variance) > 0.01:
raise AccountingMirrorError(
f"FIFO reconstruction does not reconcile for {party_name}: "
f"mirror 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 reconstructed from the Accounting Mirror closing balance; "
"newest credit-side additions form the unpaid closing balance and "
"debit-side settlements are deemed to clear the oldest balance first"
),
"current_mirror": current.get("mirror") or {},
"follow_up_mirror": later.get("mirror") or {},
}