1112 lines
43 KiB
Python
1112 lines
43 KiB
Python
from __future__ import annotations
|
|
|
|
from typing import Any
|
|
from datetime import date, datetime, timezone
|
|
from pathlib import Path
|
|
import re
|
|
|
|
from sqlalchemy import select, func
|
|
|
|
from app.modules.accounting.agent_bridge import request_agent_command
|
|
from app.modules.accounting.accounting_mirror_models import AccountingMirrorRegistry
|
|
|
|
|
|
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,
|
|
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 {},
|
|
}
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Accounting Mirror metadata/version helpers
|
|
# Kept inside the existing Accounting module/service so every downstream tool
|
|
# resolves Client + FY mirrors through one source of truth.
|
|
# ---------------------------------------------------------------------------
|
|
|
|
def _mirror_utcnow() -> datetime:
|
|
return datetime.now(timezone.utc)
|
|
|
|
|
|
def _mirror_fy_key(value: str) -> int:
|
|
match = re.fullmatch(r"(\d{4})-(\d{2})", str(value or "").strip())
|
|
return int(match.group(1)) if match else -1
|
|
|
|
|
|
def _parse_mirror_date(value: Any) -> date | None:
|
|
text = str(value or "").strip()
|
|
if not text:
|
|
return None
|
|
try:
|
|
return date.fromisoformat(text[:10])
|
|
except Exception:
|
|
return None
|
|
|
|
|
|
def _parse_mirror_datetime(value: Any) -> datetime | None:
|
|
text = str(value or "").strip()
|
|
if not text:
|
|
return None
|
|
try:
|
|
parsed = datetime.fromisoformat(text.replace("Z", "+00:00"))
|
|
return parsed if parsed.tzinfo else parsed.replace(tzinfo=timezone.utc)
|
|
except Exception:
|
|
return None
|
|
|
|
|
|
def get_current_accounting_mirror(db, tenant_id: int, client_id: int, financial_year: str) -> AccountingMirrorRegistry | None:
|
|
return db.execute(
|
|
select(AccountingMirrorRegistry).where(
|
|
AccountingMirrorRegistry.tenant_id == int(tenant_id),
|
|
AccountingMirrorRegistry.client_id == int(client_id),
|
|
AccountingMirrorRegistry.financial_year == str(financial_year or "").strip(),
|
|
AccountingMirrorRegistry.is_current.is_(True),
|
|
AccountingMirrorRegistry.is_active.is_(True),
|
|
AccountingMirrorRegistry.status == "active",
|
|
)
|
|
).scalar_one_or_none()
|
|
|
|
|
|
def get_registered_mirror(db, tenant_id: int, client_id: int, financial_year: str) -> AccountingMirrorRegistry | None:
|
|
"""Backward-compatible name used by existing Accounting tools."""
|
|
return get_current_accounting_mirror(db, tenant_id, client_id, financial_year)
|
|
|
|
|
|
def registered_mirror_payload(
|
|
db,
|
|
*,
|
|
tenant_id: int,
|
|
client_id: int,
|
|
financial_year: str,
|
|
fallback: dict[str, Any] | None = None,
|
|
) -> dict[str, Any]:
|
|
"""Return the standard Local-Agent payload using the current registered mirror.
|
|
|
|
Existing callers may pass the old deterministic payload as ``fallback``. When a
|
|
current Client/FY mirror is registered, its stored directory/file metadata wins.
|
|
This keeps all Accounting tools on one source of truth without breaking clients
|
|
that have not yet been backfilled into the registry.
|
|
"""
|
|
payload = dict(fallback or {})
|
|
payload.setdefault("client_id", int(client_id))
|
|
payload["financial_year"] = str(financial_year or payload.get("financial_year") or "").strip()
|
|
row = get_current_accounting_mirror(db, int(tenant_id), int(client_id), payload["financial_year"])
|
|
if not row:
|
|
payload["mirror_registered"] = False
|
|
return payload
|
|
payload.update({
|
|
"accounting_relative_dir": str(row.accounting_relative_dir or "").strip(),
|
|
"mirror_file_name": str(row.mirror_file_name or "").strip(),
|
|
"mirror_local_path": str(getattr(row, "mirror_local_path", "") or "").strip(),
|
|
"tally_data_relative_path": str(getattr(row, "tally_data_relative_path", "") or "").strip(),
|
|
"tally_data_local_path": str(getattr(row, "tally_data_local_path", "") or "").strip(),
|
|
"mirror_registry_id": int(row.id),
|
|
"mirror_version_no": int(row.version_no or 1),
|
|
"mirror_registered": True,
|
|
})
|
|
return payload
|
|
|
|
|
|
def update_accounting_mirror_paths(
|
|
db,
|
|
*,
|
|
tenant_id: int,
|
|
client_id: int,
|
|
financial_year: str,
|
|
mirror_local_path: str | None = None,
|
|
tally_data_relative_path: str | None = None,
|
|
tally_data_local_path: str | None = None,
|
|
) -> AccountingMirrorRegistry | None:
|
|
"""Idempotently enrich the current Client/FY mirror with workstation paths."""
|
|
row = get_current_accounting_mirror(db, tenant_id, client_id, financial_year)
|
|
if not row:
|
|
return None
|
|
changed = False
|
|
for attr, value in (
|
|
("mirror_local_path", mirror_local_path),
|
|
("tally_data_relative_path", tally_data_relative_path),
|
|
("tally_data_local_path", tally_data_local_path),
|
|
):
|
|
if value is None:
|
|
continue
|
|
value = str(value or "").strip()
|
|
if value and str(getattr(row, attr, "") or "").strip() != value:
|
|
setattr(row, attr, value)
|
|
changed = True
|
|
if changed:
|
|
row.updated_at_utc = _mirror_utcnow()
|
|
db.flush()
|
|
return row
|
|
|
|
|
|
def list_accounting_mirror_versions(db, tenant_id: int, client_id: int, financial_year: str) -> list[AccountingMirrorRegistry]:
|
|
rows = db.execute(
|
|
select(AccountingMirrorRegistry).where(
|
|
AccountingMirrorRegistry.tenant_id == int(tenant_id),
|
|
AccountingMirrorRegistry.client_id == int(client_id),
|
|
AccountingMirrorRegistry.financial_year == str(financial_year or "").strip(),
|
|
).order_by(AccountingMirrorRegistry.version_no.desc(), AccountingMirrorRegistry.id.desc())
|
|
).scalars().all()
|
|
return list(rows)
|
|
|
|
|
|
def list_registered_mirrors(db, tenant_id: int, client_id: int) -> list[AccountingMirrorRegistry]:
|
|
rows = db.execute(
|
|
select(AccountingMirrorRegistry).where(
|
|
AccountingMirrorRegistry.tenant_id == int(tenant_id),
|
|
AccountingMirrorRegistry.client_id == int(client_id),
|
|
AccountingMirrorRegistry.is_current.is_(True),
|
|
AccountingMirrorRegistry.is_active.is_(True),
|
|
AccountingMirrorRegistry.status == "active",
|
|
)
|
|
).scalars().all()
|
|
return sorted(rows, key=lambda row: _mirror_fy_key(row.financial_year), reverse=True)
|
|
|
|
|
|
def next_accounting_mirror_version(db, tenant_id: int, client_id: int, financial_year: str) -> int:
|
|
value = db.execute(
|
|
select(func.max(AccountingMirrorRegistry.version_no)).where(
|
|
AccountingMirrorRegistry.tenant_id == int(tenant_id),
|
|
AccountingMirrorRegistry.client_id == int(client_id),
|
|
AccountingMirrorRegistry.financial_year == str(financial_year or "").strip(),
|
|
)
|
|
).scalar_one_or_none()
|
|
return int(value or 0) + 1
|
|
|
|
|
|
def register_accounting_mirror_version(
|
|
db,
|
|
*,
|
|
tenant_id: int,
|
|
client_id: int,
|
|
financial_year: str,
|
|
accounting_relative_dir: str,
|
|
storage_node_id: int | None,
|
|
mirror: dict[str, Any] | None = None,
|
|
job: dict[str, Any] | None = None,
|
|
requested_by_user_id: int | None = None,
|
|
version_no: int | None = None,
|
|
archived_previous_file_name: str = "",
|
|
) -> AccountingMirrorRegistry:
|
|
mirror = dict(mirror or {})
|
|
job = dict(job or {})
|
|
company = dict(mirror.get("company") or {})
|
|
voucher_period = dict(mirror.get("voucher_period") or {})
|
|
fy = str(financial_year or "").strip()
|
|
|
|
current = get_current_accounting_mirror(db, tenant_id, client_id, fy)
|
|
target_version = int(version_no or next_accounting_mirror_version(db, tenant_id, client_id, fy))
|
|
|
|
# Idempotent status polling: if this version was already registered, update it
|
|
# instead of creating another row.
|
|
row = db.execute(
|
|
select(AccountingMirrorRegistry).where(
|
|
AccountingMirrorRegistry.tenant_id == int(tenant_id),
|
|
AccountingMirrorRegistry.client_id == int(client_id),
|
|
AccountingMirrorRegistry.financial_year == fy,
|
|
AccountingMirrorRegistry.version_no == target_version,
|
|
)
|
|
).scalar_one_or_none()
|
|
|
|
if row is None:
|
|
row = AccountingMirrorRegistry(
|
|
tenant_id=int(tenant_id),
|
|
client_id=int(client_id),
|
|
financial_year=fy,
|
|
version_no=target_version,
|
|
is_current=True,
|
|
is_active=True,
|
|
status="active",
|
|
accounting_relative_dir=str(accounting_relative_dir or "").strip(),
|
|
mirror_file_name=Path(str(job.get("accounting_db_path") or f"client_{int(client_id):08d}.act")).name,
|
|
created_by_user_id=requested_by_user_id,
|
|
supersedes_mirror_id=(int(current.id) if current else None),
|
|
)
|
|
db.add(row)
|
|
db.flush()
|
|
|
|
if current and int(current.id) != int(row.id):
|
|
current.is_current = False
|
|
current.status = "superseded"
|
|
current.updated_at_utc = _mirror_utcnow()
|
|
if archived_previous_file_name:
|
|
current.mirror_file_name = Path(str(archived_previous_file_name)).name
|
|
current.replacement_count = int(current.replacement_count or 0) + 1
|
|
|
|
# Ensure no other row remains current for the same Client/FY.
|
|
others = db.execute(
|
|
select(AccountingMirrorRegistry).where(
|
|
AccountingMirrorRegistry.tenant_id == int(tenant_id),
|
|
AccountingMirrorRegistry.client_id == int(client_id),
|
|
AccountingMirrorRegistry.financial_year == fy,
|
|
AccountingMirrorRegistry.id != int(row.id),
|
|
AccountingMirrorRegistry.is_current.is_(True),
|
|
)
|
|
).scalars().all()
|
|
for other in others:
|
|
other.is_current = False
|
|
if other.status == "active":
|
|
other.status = "superseded"
|
|
other.updated_at_utc = _mirror_utcnow()
|
|
|
|
row.storage_node_id = int(storage_node_id) if storage_node_id else None
|
|
row.accounting_relative_dir = str(accounting_relative_dir or row.accounting_relative_dir or "").strip()
|
|
row.mirror_file_name = Path(str(job.get("accounting_db_path") or row.mirror_file_name or f"client_{int(client_id):08d}.act")).name
|
|
mirror_local_path = str(job.get("accounting_db_path") or mirror.get("path") or "").strip()
|
|
if mirror_local_path:
|
|
row.mirror_local_path = mirror_local_path
|
|
tally_relative = str(job.get("tally_data_relative_path") or mirror.get("tally_data_relative_path") or "").strip()
|
|
tally_local = str(job.get("tally_data_local_path") or mirror.get("tally_data_local_path") or "").strip()
|
|
if tally_relative:
|
|
row.tally_data_relative_path = tally_relative
|
|
if tally_local:
|
|
row.tally_data_local_path = tally_local
|
|
row.company_name = str(job.get("company_name") or company.get("company_name") or row.company_name or "").strip()
|
|
row.company_guid = str(job.get("tally_guid") or company.get("company_guid") or row.company_guid or "").strip()
|
|
row.voucher_from_date = _parse_mirror_date(voucher_period.get("from_date")) or row.voucher_from_date
|
|
row.voucher_to_date = _parse_mirror_date(voucher_period.get("to_date")) or row.voucher_to_date
|
|
try:
|
|
row.file_size_bytes = int(mirror.get("size_bytes") or row.file_size_bytes or 0)
|
|
except Exception:
|
|
pass
|
|
row.status = "active"
|
|
row.is_active = True
|
|
row.is_current = True
|
|
row.last_synced_by_user_id = requested_by_user_id
|
|
row.last_synced_at_utc = _parse_mirror_datetime(job.get("finished_at_utc")) or _parse_mirror_datetime(job.get("updated_at_utc")) or _mirror_utcnow()
|
|
row.updated_at_utc = _mirror_utcnow()
|
|
db.flush()
|
|
return row
|
|
|
|
|
|
def upsert_registered_mirror(
|
|
db,
|
|
*,
|
|
tenant_id: int,
|
|
client_id: int,
|
|
financial_year: str,
|
|
accounting_relative_dir: str,
|
|
storage_node_id: int | None,
|
|
mirror: dict[str, Any] | None = None,
|
|
job: dict[str, Any] | None = None,
|
|
requested_by_user_id: int | None = None,
|
|
replacement: bool = False,
|
|
) -> AccountingMirrorRegistry:
|
|
"""Compatibility/backfill helper.
|
|
|
|
Discovery of a pre-versioning canonical mirror creates/updates v1 only when no
|
|
current row exists. A successful explicit replacement is registered through
|
|
register_accounting_mirror_version() with the job's requested version.
|
|
"""
|
|
fy = str(financial_year or "").strip()
|
|
current = get_current_accounting_mirror(db, tenant_id, client_id, fy)
|
|
if replacement:
|
|
return register_accounting_mirror_version(
|
|
db,
|
|
tenant_id=tenant_id,
|
|
client_id=client_id,
|
|
financial_year=fy,
|
|
accounting_relative_dir=accounting_relative_dir,
|
|
storage_node_id=storage_node_id,
|
|
mirror=mirror,
|
|
job=job,
|
|
requested_by_user_id=requested_by_user_id,
|
|
version_no=int((job or {}).get("mirror_version_no") or 0) or None,
|
|
archived_previous_file_name=str((job or {}).get("archived_previous_file_name") or ""),
|
|
)
|
|
if current:
|
|
# Refresh metadata for the current physical canonical mirror without
|
|
# manufacturing another version during discovery/status polling.
|
|
current.storage_node_id = int(storage_node_id) if storage_node_id else None
|
|
current.accounting_relative_dir = str(accounting_relative_dir or current.accounting_relative_dir or "").strip()
|
|
current.company_name = str(((mirror or {}).get("company") or {}).get("company_name") or current.company_name or "").strip()
|
|
current.company_guid = str(((mirror or {}).get("company") or {}).get("company_guid") or current.company_guid or "").strip()
|
|
current.file_size_bytes = int((mirror or {}).get("size_bytes") or current.file_size_bytes or 0)
|
|
discovered_path = str((job or {}).get("accounting_db_path") or (mirror or {}).get("path") or "").strip()
|
|
if discovered_path:
|
|
current.mirror_local_path = discovered_path
|
|
current.updated_at_utc = _mirror_utcnow()
|
|
db.flush()
|
|
return current
|
|
return register_accounting_mirror_version(
|
|
db,
|
|
tenant_id=tenant_id,
|
|
client_id=client_id,
|
|
financial_year=fy,
|
|
accounting_relative_dir=accounting_relative_dir,
|
|
storage_node_id=storage_node_id,
|
|
mirror=mirror,
|
|
job=job,
|
|
requested_by_user_id=requested_by_user_id,
|
|
version_no=1,
|
|
)
|
|
|
|
|
|
def sync_discovered_mirrors(
|
|
db,
|
|
*,
|
|
tenant_id: int,
|
|
client_id: int,
|
|
storage_node_id: int | None,
|
|
discovered: list[dict[str, Any]],
|
|
requested_by_user_id: int | None = None,
|
|
) -> list[AccountingMirrorRegistry]:
|
|
output: list[AccountingMirrorRegistry] = []
|
|
discovered_fys: set[str] = set()
|
|
for item in discovered or []:
|
|
fy = str(item.get("financial_year") or "").strip()
|
|
relative_dir = str(item.get("accounting_relative_dir") or "").strip()
|
|
mirror = dict(item.get("mirror") or {})
|
|
if not re.fullmatch(r"\d{4}-\d{2}", fy) or not relative_dir or not mirror.get("ready"):
|
|
continue
|
|
discovered_fys.add(fy)
|
|
row = upsert_registered_mirror(
|
|
db,
|
|
tenant_id=tenant_id,
|
|
client_id=client_id,
|
|
financial_year=fy,
|
|
accounting_relative_dir=relative_dir,
|
|
storage_node_id=storage_node_id,
|
|
mirror=mirror,
|
|
job=item,
|
|
requested_by_user_id=requested_by_user_id,
|
|
replacement=False,
|
|
)
|
|
output.append(row)
|
|
|
|
# Only current rows are availability indicators. Historical/superseded
|
|
# versions are intentionally left untouched.
|
|
current_rows = db.execute(
|
|
select(AccountingMirrorRegistry).where(
|
|
AccountingMirrorRegistry.tenant_id == int(tenant_id),
|
|
AccountingMirrorRegistry.client_id == int(client_id),
|
|
AccountingMirrorRegistry.is_current.is_(True),
|
|
AccountingMirrorRegistry.is_active.is_(True),
|
|
)
|
|
).scalars().all()
|
|
for row in current_rows:
|
|
if row.financial_year not in discovered_fys:
|
|
row.status = "missing"
|
|
row.is_active = False
|
|
row.updated_at_utc = _mirror_utcnow()
|
|
db.commit()
|
|
return output
|