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("mirror_db_path") or mirror.get("path") or job.get("accounting_db_path") or f"client_{int(client_id):08d}_mirror.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("mirror_db_path") or mirror.get("path") or job.get("accounting_db_path") or row.mirror_file_name or f"client_{int(client_id):08d}_mirror.act")).name mirror_local_path = str(job.get("mirror_db_path") or mirror.get("path") or job.get("accounting_db_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("mirror_db_path") or (mirror or {}).get("path") or (job or {}).get("accounting_db_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