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

665 lines
24 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.
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 {},
}