Unify accounting analysis on local SQLite mirror

This commit is contained in:
A R R R Associates
2026-09-06 20:13:51 +05:30
parent 6c2a2081ed
commit 7adbbde220
11 changed files with 848 additions and 129 deletions
+1 -1
View File
@@ -4,7 +4,7 @@ import io
from pathlib import Path
import zipfile
ERP_LOCAL_AGENT_VERSION = "1.24.6"
ERP_LOCAL_AGENT_VERSION = "1.25.0"
ERP_LOCAL_AGENT_NAME = "ERP Local Agent"
RUNTIME_ROOT = Path(__file__).resolve().parent / "local_agent_runtime"
_DETERMINISTIC_ZIP_TIMESTAMP = (2026, 1, 1, 0, 0, 0)
@@ -1,2 +1,2 @@
__version__ = "1.24.6"
__version__ = "1.25.0"
AGENT_NAME = "ERP Local Agent"
@@ -178,6 +178,14 @@ class AgentCommandProcessor:
result = self._accounting_full_export_status(payload)
elif action == "accounting_mirror_query":
result = self._accounting_mirror_query(payload)
elif action == "accounting_analysis_save":
result = self._accounting_analysis_save(payload)
elif action == "accounting_analysis_history":
result = self._accounting_analysis_history(payload)
elif action == "accounting_analysis_get":
result = self._accounting_analysis_get(payload)
elif action == "accounting_opening_balance_mirror_snapshot":
result = self._opening_balance_mirror_snapshot(payload)
else:
raise ValueError(f"Unsupported local-agent command: {action}")
ok = True
@@ -205,6 +213,8 @@ class AgentCommandProcessor:
"full_accounting_export_capability": True,
"accounting_mirror_progress_capability": True,
"mirror_first_accounting_reads": True,
"accounting_analysis_history_capability": True,
"cross_fy_mirror_analysis_capability": True,
"historical_learning_read_capability": True,
"purchase_posting_preflight_capability": True,
"purchase_voucher_write_capability": True,
@@ -471,6 +481,274 @@ class AgentCommandProcessor:
def _ensure_analysis_tables(self, client_id: int) -> None:
db = self.store.connect(client_id)
try:
db.executescript(
"""
CREATE TABLE IF NOT EXISTS accounting_analysis_runs (
id INTEGER PRIMARY KEY AUTOINCREMENT,
analysis_type TEXT NOT NULL,
financial_year TEXT NOT NULL DEFAULT '',
period_from TEXT NOT NULL DEFAULT '',
period_to TEXT NOT NULL DEFAULT '',
company_guid TEXT NOT NULL DEFAULT '',
company_name TEXT NOT NULL DEFAULT '',
status TEXT NOT NULL DEFAULT 'completed',
requested_by_user_id INTEGER,
source_mirror_path TEXT NOT NULL DEFAULT '',
source_mirror_synced_at TEXT NOT NULL DEFAULT '',
parameters_json TEXT NOT NULL DEFAULT '{}',
summary_json TEXT NOT NULL DEFAULT '{}',
result_json TEXT NOT NULL DEFAULT '{}',
created_at_utc TEXT NOT NULL,
completed_at_utc TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS ix_accounting_analysis_runs_type
ON accounting_analysis_runs(analysis_type, created_at_utc DESC);
"""
)
db.commit()
finally:
db.close()
def _save_analysis_run(
self,
client_id: int,
*,
analysis_type: str,
financial_year: str = "",
period_from: str = "",
period_to: str = "",
company_guid: str = "",
company_name: str = "",
requested_by_user_id: int | None = None,
source_mirror_path: str = "",
source_mirror_synced_at: str = "",
parameters: dict[str, Any] | None = None,
summary: dict[str, Any] | None = None,
result: dict[str, Any] | None = None,
) -> int:
self._ensure_analysis_tables(client_id)
now = datetime.now(timezone.utc).isoformat()
db = self.store.connect(client_id)
try:
cur = db.execute(
"""INSERT INTO accounting_analysis_runs(
analysis_type,financial_year,period_from,period_to,company_guid,company_name,
status,requested_by_user_id,source_mirror_path,source_mirror_synced_at,
parameters_json,summary_json,result_json,created_at_utc,completed_at_utc
) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
(
str(analysis_type or "").strip(),
str(financial_year or "").strip(),
str(period_from or "").strip(),
str(period_to or "").strip(),
str(company_guid or "").strip(),
str(company_name or "").strip(),
"completed",
int(requested_by_user_id) if requested_by_user_id not in (None, "") else None,
str(source_mirror_path or "").strip(),
str(source_mirror_synced_at or "").strip(),
json.dumps(parameters or {}, ensure_ascii=False, default=str),
json.dumps(summary or {}, ensure_ascii=False, default=str),
json.dumps(result or {}, ensure_ascii=False, default=str),
now,
now,
),
)
db.commit()
return int(cur.lastrowid)
finally:
db.close()
def _accounting_analysis_save(self, payload: dict[str, Any]) -> dict[str, Any]:
client_id = int(payload.get("client_id") or 0)
if client_id <= 0:
raise ValueError("client_id is required.")
run_id = self._save_analysis_run(
client_id,
analysis_type=str(payload.get("analysis_type") or "ACCOUNTING_ANALYSIS"),
financial_year=str(payload.get("financial_year") or ""),
period_from=str(payload.get("period_from") or ""),
period_to=str(payload.get("period_to") or ""),
company_guid=str(payload.get("company_guid") or ""),
company_name=str(payload.get("company_name") or ""),
requested_by_user_id=payload.get("requested_by_user_id"),
source_mirror_path=str(payload.get("source_mirror_path") or ""),
source_mirror_synced_at=str(payload.get("source_mirror_synced_at") or ""),
parameters=dict(payload.get("parameters") or {}),
summary=dict(payload.get("summary") or {}),
result=dict(payload.get("result") or {}),
)
return {"analysis_run_id": run_id, "agent": self._agent_info()}
def _accounting_analysis_history(self, payload: dict[str, Any]) -> dict[str, Any]:
client_id = int(payload.get("client_id") or 0)
if client_id <= 0:
raise ValueError("client_id is required.")
self._ensure_analysis_tables(client_id)
limit = max(1, min(100, int(payload.get("limit") or 30)))
analysis_type = str(payload.get("analysis_type") or "").strip()
db = self.store.connect(client_id)
try:
if analysis_type:
rows = db.execute(
"""SELECT id,analysis_type,financial_year,period_from,period_to,company_name,
status,summary_json,created_at_utc,completed_at_utc
FROM accounting_analysis_runs
WHERE analysis_type=?
ORDER BY id DESC LIMIT ?""",
(analysis_type, limit),
).fetchall()
else:
rows = db.execute(
"""SELECT id,analysis_type,financial_year,period_from,period_to,company_name,
status,summary_json,created_at_utc,completed_at_utc
FROM accounting_analysis_runs
ORDER BY id DESC LIMIT ?""",
(limit,),
).fetchall()
result = []
for row in rows:
item = dict(row)
try:
item["summary"] = json.loads(item.pop("summary_json") or "{}")
except Exception:
item["summary"] = {}
result.append(item)
return {"runs": result, "agent": self._agent_info()}
finally:
db.close()
def _accounting_analysis_get(self, payload: dict[str, Any]) -> dict[str, Any]:
client_id = int(payload.get("client_id") or 0)
run_id = int(payload.get("run_id") or 0)
if client_id <= 0 or run_id <= 0:
raise ValueError("client_id and run_id are required.")
self._ensure_analysis_tables(client_id)
db = self.store.connect(client_id)
try:
row = db.execute(
"SELECT * FROM accounting_analysis_runs WHERE id=?",
(run_id,),
).fetchone()
if not row:
raise ValueError("Analysis run was not found in the client SQLite database.")
item = dict(row)
for field in ("parameters_json", "summary_json", "result_json"):
target = field[:-5]
try:
item[target] = json.loads(item.pop(field) or "{}")
except Exception:
item[target] = {}
return {"run": item, "agent": self._agent_info()}
finally:
db.close()
def _mirror_db_from_relative_dir(self, client_id: int, relative_dir: str):
from pathlib import Path as _Path
import os as _os
relative = _Path(str(relative_dir or "").replace("\\", "/"))
if relative.is_absolute() or ".." in relative.parts:
raise ValueError("Invalid Accounting Mirror storage path.")
base = self.store.storage_root.resolve()
folder = (base / relative).resolve()
if _os.path.commonpath([str(base), str(folder)]) != str(base):
raise ValueError("Accounting Mirror storage path escapes the Storage Node root.")
return folder / f"client_{int(client_id):08d}_mirror.act"
@staticmethod
def _mirror_master_snapshot_from_file(db_path, *, previous_year: bool = False):
import sqlite3 as _sqlite3
if not db_path.is_file():
raise ValueError(f"Accounting Mirror was not found: {db_path}")
db = _sqlite3.connect(db_path, timeout=60)
db.row_factory = _sqlite3.Row
try:
company = db.execute("SELECT * FROM company_master ORDER BY synced_at DESC LIMIT 1").fetchone()
ledgers = []
for row in db.execute("SELECT * FROM ledger_master ORDER BY ledger_name").fetchall():
ledgers.append({
"guid": row["ledger_guid"] or "",
"name": row["ledger_name"] or "",
"parent": row["parent_group"] or "",
"opening_balance": float(row["opening_balance"] or 0),
"closing_balance": float(row["closing_balance"] or 0),
"is_revenue": row["is_revenue"] or "",
"is_active": "Yes",
})
stocks = []
for row in db.execute("SELECT * FROM stock_item_master ORDER BY product_name").fetchall():
opening_qty = float(row["opening_qty"] or 0)
opening_value = float(row["opening_value"] or 0)
closing_qty = opening_qty
closing_value = opening_value
if previous_year:
movement = db.execute(
"""SELECT
COALESCE(SUM(CASE WHEN UPPER(direction)='INWARD' THEN ABS(quantity)
WHEN UPPER(direction)='OUTWARD' THEN -ABS(quantity) ELSE 0 END),0),
COALESCE(SUM(CASE WHEN UPPER(direction)='INWARD' THEN ABS(value)
WHEN UPPER(direction)='OUTWARD' THEN -ABS(value) ELSE 0 END),0)
FROM inventory_movement WHERE stock_item_guid=?""",
(row["stock_item_guid"],),
).fetchone()
closing_qty = opening_qty + float(movement[0] or 0)
closing_value = opening_value + float(movement[1] or 0)
stocks.append({
"guid": row["stock_item_guid"] or "",
"name": row["product_name"] or "",
"parent": row["parent_group"] or "",
"base_units": row["base_uom"] or "",
"opening_balance": opening_qty,
"opening_value": opening_value,
"closing_balance": closing_qty,
"closing_value": closing_value,
"hsn_code": row["current_hsn"] or "",
"is_active": "Yes",
})
company_dict = dict(company) if company else {}
return {
"company": {
"name": str(company_dict.get("company_name") or ""),
"guid": str(company_dict.get("company_guid") or ""),
"gstin": str(company_dict.get("gstin") or ""),
},
"masters": {"ledgers": ledgers, "stock_items": stocks},
"path": str(db_path),
"sync": company_dict,
}
finally:
db.close()
def _opening_balance_mirror_snapshot(self, payload: dict[str, Any]) -> dict[str, Any]:
client_id = int(payload.get("client_id") or 0)
previous_relative_dir = str(payload.get("previous_accounting_relative_dir") or "").strip()
current_relative_dir = str(payload.get("current_accounting_relative_dir") or payload.get("accounting_relative_dir") or "").strip()
if not previous_relative_dir or not current_relative_dir:
raise ValueError("Previous-year and current-year Accounting Mirror paths are required.")
previous_path = self._mirror_db_from_relative_dir(client_id, previous_relative_dir)
current_path = self._mirror_db_from_relative_dir(client_id, current_relative_dir)
previous = self._mirror_master_snapshot_from_file(previous_path, previous_year=True)
current = self._mirror_master_snapshot_from_file(current_path, previous_year=False)
if previous["company"].get("guid") and current["company"].get("guid") and previous["company"]["guid"] != current["company"]["guid"]:
self.logger.info(
"Opening Balance cross-FY mirror company GUID differs previous=%s current=%s; continuing because FY-split Tally companies are supported.",
previous["company"]["guid"], current["company"]["guid"],
)
return {
"previous_company": previous["company"],
"current_company": current["company"],
"previous_masters": previous["masters"],
"current_masters": current["masters"],
"previous_mirror_path": previous["path"],
"current_mirror_path": current["path"],
"read_only": True,
"agent": self._agent_info(),
}
def _opening_balance_company(self, company_name: str):
wanted=str(company_name or "").strip()
if not wanted: raise ValueError("Tally company name is required.")
@@ -2391,8 +2669,7 @@ class AgentCommandProcessor:
exceptions.sort(key=lambda x: (x["date"], x["party"].casefold()))
split_patterns.sort(key=lambda x: (x["date_from"], x["party"].casefold()))
return {
"cash_payment_review": {
review = {
"company_name": company_name,
"company_guid": company_guid,
"date_from": date_from,
@@ -2418,10 +2695,24 @@ class AgentCommandProcessor:
"tally_retry_count": 0,
"tally_pause_ms": 0,
"read_only": True,
},
"mirror": mirror,
"agent": self._agent_info(),
}
}
analysis_run_id = self._save_analysis_run(
client_id,
analysis_type="CASH_PAYMENT_COMPLIANCE",
financial_year=str(payload.get("financial_year") or ""),
period_from=date_from,
period_to=date_to,
company_guid=company_guid,
company_name=company_name,
requested_by_user_id=payload.get("requested_by_user_id"),
source_mirror_path=str(mirror_status.get("mirror_db_path") or mirror.get("path") or ""),
source_mirror_synced_at=str((mirror.get("sync") or {}).get("last_sync_at") or ""),
parameters={"cash_limit": cash_limit, "split_window_days": split_window_days, "near_limit_percent": near_limit_percent},
summary=review.get("summary") or {},
result=review,
)
review["analysis_run_id"] = analysis_run_id
return {"cash_payment_review": review, "mirror": mirror, "agent": self._agent_info()}
def _cash_payment_cache_analyze(self, payload: dict[str, Any]) -> dict[str, Any]:
from collections import defaultdict
@@ -2832,15 +3123,24 @@ class AgentCommandProcessor:
from collections import defaultdict
from datetime import date as _date
company, company_name = self._resolve_open_company(payload)
client_id = int(payload.get("client_id") or 0)
if client_id <= 0:
raise ValueError("client_id is required.")
mirror_status = self.tally.mirror.status(client_id)
mirror = mirror_status.get("mirror") or {}
if not mirror.get("ready"):
raise ValueError("Accounting Mirror is not ready. Run 'Mirror Tally to SQLite' first.")
company_info = mirror.get("company") or {}
company_name = str(company_info.get("company_name") or "").strip()
company_guid = str(company_info.get("company_guid") or "").strip()
date_from = str(payload.get("date_from") or "").strip()
date_to = str(payload.get("date_to") or "").strip()
rules = list(payload.get("rules") or [])
if not rules:
raise ValueError("No active TDS rules were supplied by ERP.")
vouchers = self.tally.export_vouchers(company_name, date_from, date_to)
# TDS ledgers are identified by both ledger name and parent/group text; this is deliberately broader than one fixed ledger name.
masters = self.tally.export_master_collection(company_name, "ledgers")
vouchers = (self.tally.mirror.transactions(client_id, company_name, date_from, date_to, company_guid).get("vouchers") or [])
# TDS ledgers are identified from the local mirror ledger master only.
masters = (self.tally.mirror.master_snapshot(client_id, company_name, company_guid).get("ledgers") or [])
tds_ledger_names = set()
for led in masters:
name = str(led.get("name") or "").strip()
@@ -2894,7 +3194,24 @@ class AgentCommandProcessor:
liability=round(max(0.0,expected-actual),2)
status="OK" if expected>0 and liability<=0.009 else ("TDS NOT DEDUCTED" if expected>0 and actual<=0.009 else ("SHORT DEDUCTION" if liability>0 else "BELOW / OUTSIDE THRESHOLD"))
transactions.append({**{k:v for k,v in x.items() if k!='rule'},"rule_id":r.get("id"),"rule_code":r.get("rule_code"),"rule_name":r.get("name"),"legacy_section":r.get("legacy_section"),"statutory_reference":r.get("statutory_reference"),"cumulative_amount":cumulative,"expected_tds":expected,"actual_tds":actual,"liability":liability,"status":status})
return {"tds_review":{"company_name":company_name,"company_guid":str(company.get("guid") or ""),"date_from":date_from,"date_to":date_to,"summary":{"vouchers_reviewed":len(vouchers),"candidate_transactions":len(transactions),"tds_not_deducted":sum(1 for x in transactions if x["status"]=="TDS NOT DEDUCTED"),"short_deduction":sum(1 for x in transactions if x["status"]=="SHORT DEDUCTION"),"purchase_transactions":sum(1 for x in transactions if x["rule_code"]=="PURCHASE_GOODS")},"transactions":transactions,"tds_ledgers":sorted(tds_ledger_names),"read_only":True},"agent":self._agent_info()}
review = {"company_name":company_name,"company_guid":company_guid,"date_from":date_from,"date_to":date_to,"summary":{"vouchers_reviewed":len(vouchers),"candidate_transactions":len(transactions),"tds_not_deducted":sum(1 for x in transactions if x["status"]=="TDS NOT DEDUCTED"),"short_deduction":sum(1 for x in transactions if x["status"]=="SHORT DEDUCTION"),"purchase_transactions":sum(1 for x in transactions if x["rule_code"]=="PURCHASE_GOODS")},"transactions":transactions,"tds_ledgers":sorted(tds_ledger_names),"read_only":True,"sqlite_mirror":True}
analysis_run_id = self._save_analysis_run(
client_id,
analysis_type="TDS_COMPLIANCE",
financial_year=str(payload.get("financial_year") or ""),
period_from=date_from,
period_to=date_to,
company_guid=company_guid,
company_name=company_name,
requested_by_user_id=payload.get("requested_by_user_id"),
source_mirror_path=str(mirror_status.get("mirror_db_path") or mirror.get("path") or ""),
source_mirror_synced_at=str((mirror.get("sync") or {}).get("last_sync_at") or ""),
parameters={"rule_count": len(rules)},
summary=review["summary"],
result=review,
)
review["analysis_run_id"] = analysis_run_id
return {"tds_review":review,"mirror":mirror,"agent":self._agent_info()}
def _post_tds_liability(self, payload: dict[str, Any]) -> dict[str, Any]:
company, company_name = self._resolve_open_company(payload)
@@ -3045,29 +3362,68 @@ class AgentCommandProcessor:
"agent": self._agent_info(),
}
def _depreciation_preview(self, payload: dict[str, Any]) -> dict[str, Any]:
client_id = int(payload.get("client_id"))
tally_guid = str(payload.get("tally_guid") or "").strip()
def _depreciation_mirror_context(self, payload: dict[str, Any]):
client_id = int(payload.get("client_id") or 0)
if client_id <= 0:
raise ValueError("client_id is required.")
mirror_status = self.tally.mirror.status(client_id)
mirror = mirror_status.get("mirror") or {}
if not mirror.get("ready"):
raise ValueError("Accounting Mirror is not ready for this financial year. Run 'Mirror Tally to SQLite' first.")
company = mirror.get("company") or {}
company_name = str(company.get("company_name") or "").strip()
company_guid = str(company.get("company_guid") or "").strip()
if not company_name:
raise ValueError("Accounting Mirror does not contain a company identity.")
mapping = {}
if company_guid:
try:
mapping = self.store.get_active_mapping_by_guid(client_id, company_guid)
except Exception:
mapping = {}
if not mapping:
mappings = self.store.list_active_mappings(client_id)
mapping = (mappings[0] if mappings else {
"tally_guid": company_guid,
"company_name": company_name,
"gstin": str(company.get("gstin") or ""),
})
return client_id, mirror_status, mirror, company_name, company_guid, mapping
def _depreciation_refresh_local_snapshot(self, payload: dict[str, Any]):
client_id, mirror_status, mirror, company_name, company_guid, mapping = self._depreciation_mirror_context(payload)
fy_start = str(payload.get("fy_start") or "").strip()
fy_end = str(payload.get("fy_end") or "").strip()
if not tally_guid: raise ValueError("Select a mapped Tally company for depreciation.")
sync_payload = dict(payload)
sync_payload["date_from"] = fy_start
sync_payload["date_to"] = fy_end
self._sync_masters(sync_payload)
self._sync_transactions(sync_payload)
return {"preview": self.store.depreciation_preview(client_id, tally_guid=tally_guid, fy_start=fy_start, fy_end=fy_end), "agent": self._agent_info()}
masters = self.tally.mirror.master_snapshot(client_id, company_name, company_guid)
transactions = self.tally.mirror.transactions(client_id, company_name, fy_start, fy_end, company_guid)
self.store.replace_master_snapshot(
client_id,
mapping={**mapping, "company_name": company_name, "tally_guid": company_guid},
masters=masters,
requested_by_user_id=int(payload.get("requested_by_user_id")) if payload.get("requested_by_user_id") not in (None, "") else None,
)
self.store.replace_transaction_snapshot(
client_id,
mapping={**mapping, "company_name": company_name, "tally_guid": company_guid},
transactions=transactions,
requested_by_user_id=int(payload.get("requested_by_user_id")) if payload.get("requested_by_user_id") not in (None, "") else None,
)
return client_id, mirror_status, mirror, company_name, company_guid
def _depreciation_preview(self, payload: dict[str, Any]) -> dict[str, Any]:
client_id, mirror_status, mirror, company_name, company_guid = self._depreciation_refresh_local_snapshot(payload)
fy_start = str(payload.get("fy_start") or "").strip()
fy_end = str(payload.get("fy_end") or "").strip()
preview = self.store.depreciation_preview(
client_id, tally_guid=company_guid, fy_start=fy_start, fy_end=fy_end
)
return {"preview": preview, "mirror": mirror, "agent": self._agent_info()}
def _calculate_it_depreciation(self, payload: dict[str, Any]) -> dict[str, Any]:
client_id = int(payload.get("client_id"))
sync_payload = dict(payload)
sync_payload["date_from"] = str(payload.get("fy_start") or "").strip()
sync_payload["date_to"] = str(payload.get("fy_end") or "").strip()
self._sync_masters(sync_payload)
self._sync_transactions(sync_payload)
client_id, mirror_status, mirror, company_name, company_guid = self._depreciation_refresh_local_snapshot(payload)
result = self.store.calculate_it_depreciation(
client_id,
tally_guid=str(payload.get("tally_guid") or "").strip(),
tally_guid=company_guid,
fy_start=str(payload.get("fy_start") or "").strip(),
fy_end=str(payload.get("fy_end") or "").strip(),
assignments=list(payload.get("assignments") or []),
@@ -3075,8 +3431,31 @@ class AgentCommandProcessor:
depreciation_reserve_ledger=str(payload.get("depreciation_reserve_ledger") or "").strip(),
requested_by_user_id=int(payload.get("requested_by_user_id")) if payload.get("requested_by_user_id") not in (None, "") else None,
)
self.logger.info("IT depreciation draft calculated client_id=%s company=%s run_id=%s total=%s", client_id, result.get("company_name"), result.get("run_id"), result.get("total_depreciation"))
return {"calculated": True, "depreciation": result, "accounting": self.store.snapshot(client_id), "agent": self._agent_info()}
analysis_run_id = self._save_analysis_run(
client_id,
analysis_type="DEPRECIATION_IT",
financial_year=str(payload.get("financial_year") or ""),
period_from=str(payload.get("fy_start") or ""),
period_to=str(payload.get("fy_end") or ""),
company_guid=company_guid,
company_name=company_name,
requested_by_user_id=payload.get("requested_by_user_id"),
source_mirror_path=str(mirror_status.get("mirror_db_path") or mirror.get("path") or ""),
source_mirror_synced_at=str((mirror.get("sync") or {}).get("last_sync_at") or ""),
parameters={
"depreciation_expense_ledger": str(payload.get("depreciation_expense_ledger") or ""),
"depreciation_reserve_ledger": str(payload.get("depreciation_reserve_ledger") or ""),
},
summary={
"total_depreciation": result.get("total_depreciation"),
"line_count": len(result.get("lines") or []),
"status": result.get("status"),
},
result=result,
)
result["analysis_run_id"] = analysis_run_id
self.logger.info("IT depreciation draft calculated from SQLite mirror client_id=%s company=%s run_id=%s total=%s", client_id, result.get("company_name"), result.get("run_id"), result.get("total_depreciation"))
return {"calculated": True, "depreciation": result, "mirror": mirror, "accounting": self.store.snapshot(client_id), "agent": self._agent_info()}
def _get_it_depreciation_run(self, payload: dict[str, Any]) -> dict[str, Any]:
client_id = int(payload.get("client_id")); run_id = int(payload.get("run_id"))