Integrate read-only accounting mirror analytics runtime

This commit is contained in:
A R R R Associates
2026-09-06 13:07:50 +05:30
parent 9cceec0283
commit 7ae8a40b49
9 changed files with 1789 additions and 3 deletions
@@ -0,0 +1,80 @@
from __future__ import annotations
from typing import Any
from app.modules.accounting.accounting_mirror_service import (
cash_transactions,
daybook,
gst_history,
hsn_history,
inventory,
ledger_transactions,
mirror_exceptions,
mirror_query,
mirror_status,
stock_items,
trial_balance,
voucher,
)
class AccountingAnalyticsDataSource:
"""Common read-only data source for accounting/audit analytical procedures.
This class is additive. Existing tools are not silently redirected in v1.23.0.
New or migrated tools can use the mirror immediately while current Opening
Balance, Cash Payment, Depreciation and posting workflows continue unchanged.
"""
def __init__(self, *, node_code: str, accounting_payload: dict[str, Any]):
self.node_code = str(node_code or "").strip()
self.accounting_payload = dict(accounting_payload or {})
if not self.node_code:
raise ValueError("node_code is required.")
def status(self):
return mirror_status(node_code=self.node_code, accounting_payload=self.accounting_payload)
def daybook(self, **kwargs):
return daybook(node_code=self.node_code, accounting_payload=self.accounting_payload, **kwargs)
def trial_balance(self, **kwargs):
return trial_balance(node_code=self.node_code, accounting_payload=self.accounting_payload, **kwargs)
def ledger_transactions(self, **kwargs):
return ledger_transactions(
node_code=self.node_code, accounting_payload=self.accounting_payload, **kwargs
)
def inventory(self, **kwargs):
return inventory(node_code=self.node_code, accounting_payload=self.accounting_payload, **kwargs)
def stock_items(self, **kwargs):
return stock_items(node_code=self.node_code, accounting_payload=self.accounting_payload, **kwargs)
def hsn_history(self, **kwargs):
return hsn_history(node_code=self.node_code, accounting_payload=self.accounting_payload, **kwargs)
def gst_history(self, **kwargs):
return gst_history(node_code=self.node_code, accounting_payload=self.accounting_payload, **kwargs)
def cash_transactions(self, **kwargs):
return cash_transactions(
node_code=self.node_code, accounting_payload=self.accounting_payload, **kwargs
)
def voucher(self, **kwargs):
return voucher(node_code=self.node_code, accounting_payload=self.accounting_payload, **kwargs)
def exceptions(self, **kwargs):
return mirror_exceptions(
node_code=self.node_code, accounting_payload=self.accounting_payload, **kwargs
)
def raw(self, query: str, filters: dict[str, Any] | None = None):
return mirror_query(
node_code=self.node_code,
accounting_payload=self.accounting_payload,
query=query,
filters=filters or {},
)
@@ -0,0 +1,209 @@
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},
)