Add async SQLite job manager for cash payment extraction

This commit is contained in:
A R R R Associates
2026-09-04 00:05:09 +05:30
parent 4739ad321e
commit 471205053a
5 changed files with 135 additions and 9 deletions
@@ -1,2 +1,2 @@
__version__ = "1.22.8"
__version__ = "1.22.9"
AGENT_NAME = "ERP Local Agent"
@@ -18,6 +18,7 @@ _CASH_TALLY_EXTRACTION_LOCK = threading.Lock()
_CASH_CACHE_BUSY = threading.Event()
_CASH_CACHE_THREADS: dict[str, threading.Thread] = {}
_CASH_CACHE_THREADS_LOCK = threading.Lock()
_CASH_CACHE_CANCEL_REQUESTS: set[str] = set()
def is_tally_cache_busy() -> bool:
return _CASH_CACHE_BUSY.is_set()
@@ -73,6 +74,8 @@ class AgentCommandProcessor:
result = self._cash_payment_cache_start(payload)
elif action == "accounting_cash_payment_cache_status":
result = self._cash_payment_cache_status(payload)
elif action == "accounting_cash_payment_cache_cancel":
result = self._cash_payment_cache_cancel(payload)
elif action == "accounting_cash_payment_cache_analyze":
result = self._cash_payment_cache_analyze(payload)
elif action == "accounting_cash_payment_compliance":
@@ -768,10 +771,15 @@ class AgentCommandProcessor:
return {"job": self._cash_payment_cache_status({"client_id": client_id, "job_id": prior["job_id"]})["job"], "resumed": False, "already_running": True}
raise ValueError("Another Tally extraction is already running on this workstation. Please wait for it to finish.")
company, company_name = self._resolve_open_company(payload)
if str(company.get("guid") or "").strip() != tally_guid:
raise ValueError("The selected Tally company is not the currently resolved company.")
# IMPORTANT: starting a cache job must never call TallyPrime. The UI/ERP
# command must receive a job id immediately, while all Tally I/O happens
# only inside the background worker. Resolve the mapped company name
# from the local .act SQLite store; the worker validates the actually
# open Tally company after the command result has already been returned.
mapping = self.store.get_active_mapping_by_guid(client_id, tally_guid)
company_name = str(mapping.get("company_name") or "").strip()
if not company_name:
raise ValueError("The mapped Tally company name is not available in local accounting storage. Refresh the Tally mapping first.")
pause_seconds = max(2.0, min(float(payload.get("tally_pause_seconds") or 3.0), 15.0))
settings = {
"cash_limit": float(payload.get("cash_limit") or 10000.0),
@@ -810,9 +818,12 @@ class AgentCommandProcessor:
worker_payload = dict(payload)
worker_payload["tally_pause_seconds"] = pause_seconds
worker_payload["company_name"] = company_name
with _CASH_CACHE_THREADS_LOCK:
_CASH_CACHE_CANCEL_REQUESTS.discard(job_id)
thread = threading.Thread(
target=self._run_cash_payment_cache_job,
args=(job_id, worker_payload, dict(company), company_name, dict(mapping)),
args=(job_id, worker_payload, dict(mapping)),
name=f"tally-cash-cache-{job_id[:8]}",
daemon=True,
)
@@ -821,13 +832,25 @@ class AgentCommandProcessor:
thread.start()
return {"job": self._cash_payment_cache_status({"client_id": client_id, "job_id": job_id})["job"], "resumed": resume, "already_running": False}
def _run_cash_payment_cache_job(self, job_id: str, payload: dict[str, Any], company: dict[str, Any], company_name: str, mapping: dict[str, Any]) -> None:
def _run_cash_payment_cache_job(self, job_id: str, payload: dict[str, Any], mapping: dict[str, Any]) -> None:
client_id = int(payload.get("client_id"))
pause_seconds = max(2.0, min(float(payload.get("tally_pause_seconds") or 3.0), 15.0))
requested_by = payload.get("requested_by_user_id")
_CASH_TALLY_EXTRACTION_LOCK.acquire()
_CASH_CACHE_BUSY.set()
try:
self._cash_cache_update(client_id, job_id, status="running", stage="Connecting to TallyPrime in background", current_date="")
if job_id in _CASH_CACHE_CANCEL_REQUESTS:
self._cash_cache_update(client_id, job_id, status="cancelled", stage="Cancelled before Tally connection", completed_at_utc=datetime.now(timezone.utc).isoformat())
return
# This is deliberately inside the worker. It may be slow, but it can
# no longer hold open the ERP start-job command or block its timeout.
company, company_name = self._resolve_open_company(payload)
expected_guid = str(payload.get("tally_guid") or "").strip()
if str(company.get("guid") or "").strip() != expected_guid:
raise ValueError("The selected Tally company is not the currently resolved company.")
self._cash_cache_update(client_id, job_id, status="running", stage="Reading cash ledger masters", current_date="")
# One master request only, then a generous cooling gap before vouchers.
masters = None
@@ -862,6 +885,9 @@ class AgentCommandProcessor:
(job_id,),
).fetchall()
for day_row in pending:
if job_id in _CASH_CACHE_CANCEL_REQUESTS:
self._cash_cache_update(client_id, job_id, status="cancelled", stage="Cancelled safely", current_date="", completed_at_utc=datetime.now(timezone.utc).isoformat())
return
day = str(day_row[0])
self._cash_cache_update(client_id, job_id, status="running", stage="Extracting one day from TallyPrime", current_date=day)
vouchers = None
@@ -917,6 +943,21 @@ class AgentCommandProcessor:
_CASH_TALLY_EXTRACTION_LOCK.release()
with _CASH_CACHE_THREADS_LOCK:
_CASH_CACHE_THREADS.pop(job_id, None)
_CASH_CACHE_CANCEL_REQUESTS.discard(job_id)
def _cash_payment_cache_cancel(self, payload: dict[str, Any]) -> dict[str, Any]:
client_id = int(payload.get("client_id"))
job_id = str(payload.get("job_id") or "").strip()
row = self._cash_cache_job_row(client_id, job_id)
if not row:
raise ValueError("Cash payment cache job was not found on this workstation.")
status = str(row.get("status") or "")
if status in {"completed", "cancelled", "failed", "paused"}:
return {"job": self._cash_payment_cache_status({"client_id": client_id, "job_id": job_id})["job"], "cancel_requested": False}
with _CASH_CACHE_THREADS_LOCK:
_CASH_CACHE_CANCEL_REQUESTS.add(job_id)
self._cash_cache_update(client_id, job_id, stage="Cancellation requested; stopping after the current safe Tally request")
return {"job": self._cash_payment_cache_status({"client_id": client_id, "job_id": job_id})["job"], "cancel_requested": True}
def _cash_payment_cache_status(self, payload: dict[str, Any]) -> dict[str, Any]:
client_id = int(payload.get("client_id"))
@@ -948,7 +989,8 @@ class AgentCommandProcessor:
"completed_at_utc": row.get("completed_at_utc"),
"background": True,
"sqlite_cache": True,
"can_resume": str(row.get("status") or "") in {"paused", "failed", "interrupted"},
"can_resume": str(row.get("status") or "") in {"paused", "failed", "interrupted", "cancelled"},
"can_cancel": str(row.get("status") or "") in {"queued", "running", "cooling"},
},
"agent": self._agent_info(),
}