Files
arrr-erp/app/modules/accounting/accounting_mirror_service.py
T
2026-09-10 16:34:57 +05:30

731 lines
28 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 in the selected FY. Ageing is then
reconstructed forward in transaction order: genuine Purchase credits create
payable lots, debit-side movements clear the oldest payable lots first, and
Receipt/other credits are never mislabelled as purchase bills.
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
# Strict signed forward FIFO. The opening credit balance is the first
# liability lot and therefore every debit-side settlement clears it before
# any later purchase/credit lot. Once the opening lot is exhausted, the
# same debit continues against subsequent credit lots in chronological
# order. If a debit exceeds all credit lots, the excess becomes a debit
# carry and offsets the next credit before that credit can create a new
# outstanding lot. This mirrors the user's FIFO requirement exactly.
#
# Non-purchase credits are retained only when they genuinely survive in
# the signed party balance (for example an advance/other credit). They are
# never labelled as purchase bills in the drill-down.
credit_lots: list[dict[str, Any]] = []
debit_carry = round(max(0.0, -opening_balance), 2)
if opening_balance > 0.009:
credit_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 _is_purchase_voucher(row: dict[str, Any]) -> bool:
return "purchase" in str(row.get("voucher_type") or "").strip().casefold()
def _settle_oldest(lots: list[dict[str, Any]], amount: float) -> float:
"""Apply a debit against the oldest surviving credit lots first."""
remaining = max(0.0, _money(amount))
for lot in lots:
if remaining <= 0.009:
break
available = max(0.0, _money(lot.get("remaining")))
if available <= 0.009:
continue
applied = round(min(available, remaining), 2)
lot["remaining"] = round(available - applied, 2)
remaining = round(remaining - applied, 2)
return remaining
def _append_credit_lot(row: dict[str, Any], amount: float, source: str) -> None:
amount = _money(amount)
if amount <= 0.009:
return
credit_lots.append({
"party_name": party_name,
"source": source,
"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": "",
})
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 == "DR":
# FIFO rule: opening balance is physically the first entry in
# credit_lots, so payments/debit-side adjustments necessarily
# exhaust opening first and then move to subsequent credits.
remaining = _settle_oldest(credit_lots, amount)
if remaining > 0.009:
debit_carry = round(debit_carry + remaining, 2)
continue
if side != "CR":
continue
# A prior excess debit is settled before this credit can become a
# fresh outstanding lot. Only the residual credit survives.
residual = amount
if debit_carry > 0.009:
offset = round(min(debit_carry, residual), 2)
debit_carry = round(debit_carry - offset, 2)
residual = round(residual - offset, 2)
if residual <= 0.009:
continue
source = "Purchase voucher" if _is_purchase_voucher(row) else "Other credit / advance"
_append_credit_lot(row, residual, source)
# Freeze the year-end signed FIFO position. The queue order has never
# changed, so any surviving lots are precisely the oldest-to-newest
# credit components after opening and all FY settlements have been
# applied.
closing_lots = [
lot for lot in credit_lots
if _money(lot.get("remaining")) > 0.009
]
for lot in closing_lots:
lot["outstanding_at_fy_end"] = _money(lot.get("remaining"))
reconstructed_credit = round(
sum(_money(lot.get("remaining")) for lot in closing_lots) - debit_carry,
2,
)
if abs(reconstructed_credit - closing_balance) > 0.01:
raise AccountingMirrorError(
f"FIFO movement reconstruction does not reconcile for {party_name}: "
f"selected-FY closing {closing_balance:.2f}, surviving credit lots "
f"{sum(_money(lot.get('remaining')) for lot in closing_lots):.2f}, "
f"debit carry {debit_carry:.2f}. Refresh the Accounting Mirror "
"for the selected FY and retry."
)
# Keep all surviving credit-side components in date order for ageing.
# Non-purchase credits remain visibly distinguished in the drill-down;
# they are not represented as purchase bills.
closing_lots = sorted(
closing_lots + non_purchase_credit_lots,
key=lambda lot: (
_parse_date(lot.get("voucher_date")) or opening_fallback_date,
str(lot.get("voucher_number") or ""),
),
)
fifo_total = round(sum(_money(lot.get("outstanding_at_fy_end")) for lot in closing_lots), 2)
if abs(fifo_total - closing_balance) > 0.01:
raise AccountingMirrorError(
f"Purchase-bill FIFO reconstruction does not reconcile for {party_name}: "
f"selected-FY closing {closing_balance:.2f}, allocated 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": (
"Strict forward FIFO from dated Accounting Mirror movements; opening credit is cleared first, "
"then subsequent credit lots are cleared chronologically. Purchase vouchers remain "
"identified separately from other surviving credits/advances."
),
"current_mirror": current.get("mirror") or {},
"follow_up_mirror": later.get("mirror") or {},
}