From 5d13373e364b04f6553e413322e392557216210c Mon Sep 17 00:00:00 2001 From: A R R R Associates Date: Fri, 4 Sep 2026 00:19:59 +0530 Subject: [PATCH] Add instant ACK for async cash payment cache jobs --- app/modules/documents/agent_package.py | 2 +- .../erp_local_agent/__init__.py | 2 +- .../erp_local_agent/commands.py | 224 ++++++++++++------ 3 files changed, 158 insertions(+), 70 deletions(-) diff --git a/app/modules/documents/agent_package.py b/app/modules/documents/agent_package.py index 5a62d0a..9214c96 100644 --- a/app/modules/documents/agent_package.py +++ b/app/modules/documents/agent_package.py @@ -4,7 +4,7 @@ import io from pathlib import Path import zipfile -ERP_LOCAL_AGENT_VERSION = "1.22.9" +ERP_LOCAL_AGENT_VERSION = "1.22.10" 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) diff --git a/app/modules/documents/local_agent_runtime/erp_local_agent/__init__.py b/app/modules/documents/local_agent_runtime/erp_local_agent/__init__.py index dc243cd..1e69a72 100644 --- a/app/modules/documents/local_agent_runtime/erp_local_agent/__init__.py +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/__init__.py @@ -1,2 +1,2 @@ -__version__ = "1.22.9" +__version__ = "1.22.10" AGENT_NAME = "ERP Local Agent" diff --git a/app/modules/documents/local_agent_runtime/erp_local_agent/commands.py b/app/modules/documents/local_agent_runtime/erp_local_agent/commands.py index 019ed9c..39965dc 100644 --- a/app/modules/documents/local_agent_runtime/erp_local_agent/commands.py +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/commands.py @@ -19,6 +19,7 @@ _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() +_CASH_CACHE_STARTING: dict[str, dict[str, Any]] = {} def is_tally_cache_busy() -> bool: return _CASH_CACHE_BUSY.is_set() @@ -751,35 +752,66 @@ class AgentCommandProcessor: raise ValueError("Cash payment cache dates must be valid YYYY-MM-DD dates.") from exc if end_day < start_day: raise ValueError("Cash payment cache To date cannot be before From date.") - if not self.store.exists(client_id): - raise ValueError("Accounting local storage is not initialized for this client.") - self._ensure_cash_cache_schema(client_id) - # If a previous identical job was interrupted/paused, resume it instead of - # discarding already cached days. Completed jobs start a fresh refresh. - with self.store.connect(client_id) as db: - prior = db.execute( - """SELECT * FROM cash_payment_cache_jobs - WHERE client_id=? AND tally_guid=? AND date_from=? AND date_to=? - ORDER BY started_at_utc DESC LIMIT 1""", - (client_id, tally_guid, date_from, date_to), - ).fetchone() - prior = dict(prior) if prior else None + # Instant acknowledgement: do not open SQLite and do not contact TallyPrime + # on the command thread. The background worker performs all potentially + # slow local-storage/mapping/schema work after this method returns. + total_days = (end_day - start_day).days + 1 + job_id = str(uuid.uuid4()) + pause_seconds = max(2.0, min(float(payload.get("tally_pause_seconds") or 3.0), 15.0)) + now = datetime.now(timezone.utc).isoformat() + starting = { + "job_id": job_id, + "client_id": client_id, + "tally_guid": tally_guid, + "date_from": date_from, + "date_to": date_to, + "status": "starting", + "stage": "Preparing local SQLite cache in background", + "current_date": "", + "total_days": total_days, + "completed_days": 0, + "percent": 0, + "vouchers_cached": 0, + "tally_requests": 0, + "retry_count": 0, + "pause_seconds": pause_seconds, + "error": "", + "started_at_utc": now, + "updated_at_utc": now, + "completed_at_utc": None, + "background": True, + "sqlite_cache": True, + "can_resume": False, + "can_cancel": True, + } with _CASH_CACHE_THREADS_LOCK: - if _CASH_CACHE_BUSY.is_set(): - if prior and str(prior.get("status") or "") in {"queued", "running", "cooling"}: - return {"job": self._cash_payment_cache_status({"client_id": client_id, "job_id": prior["job_id"]})["job"], "resumed": False, "already_running": True} + if _CASH_CACHE_BUSY.is_set() or any(t.is_alive() for t in _CASH_CACHE_THREADS.values()): raise ValueError("Another Tally extraction is already running on this workstation. Please wait for it to finish.") + _CASH_CACHE_STARTING[job_id] = dict(starting) + _CASH_CACHE_CANCEL_REQUESTS.discard(job_id) - # 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.") + worker_payload = dict(payload) + worker_payload["tally_pause_seconds"] = pause_seconds + thread = threading.Thread( + target=self._initialize_and_run_cash_payment_cache_job, + args=(job_id, worker_payload), + name=f"tally-cash-cache-{job_id[:8]}", + daemon=True, + ) + with _CASH_CACHE_THREADS_LOCK: + _CASH_CACHE_THREADS[job_id] = thread + thread.start() + return {"job": dict(starting), "resumed": False, "already_running": False, "instant_ack": True} + + def _initialize_and_run_cash_payment_cache_job(self, job_id: str, payload: dict[str, Any]) -> None: + client_id = int(payload.get("client_id")) + tally_guid = str(payload.get("tally_guid") or "").strip() + date_from = str(payload.get("date_from") or "").strip() + date_to = str(payload.get("date_to") or "").strip() + start_day = _dt_date.fromisoformat(date_from) + end_day = _dt_date.fromisoformat(date_to) + total_days = (end_day - start_day).days + 1 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), @@ -788,49 +820,94 @@ class AgentCommandProcessor: "requested_by_user_id": payload.get("requested_by_user_id"), } now = datetime.now(timezone.utc).isoformat() - resume = bool(prior and str(prior.get("status") or "") in {"failed", "paused", "interrupted", "running", "cooling"}) - job_id = str(prior.get("job_id")) if resume else str(uuid.uuid4()) - total_days = (end_day - start_day).days + 1 - with self.store.connect(client_id) as db: - if resume: - db.execute( - """UPDATE cash_payment_cache_jobs SET status='queued', stage='Queued for resume', error_message=NULL, - pause_seconds=?, settings_json=?, updated_at_utc=? WHERE job_id=?""", - (pause_seconds, json.dumps(settings, separators=(",", ":")), now, job_id), - ) - else: - db.execute( - """INSERT INTO cash_payment_cache_jobs( - job_id, client_id, tally_guid, company_name, date_from, date_to, status, stage, - total_days, completed_days, vouchers_cached, tally_requests, retry_count, - pause_seconds, settings_json, started_at_utc, updated_at_utc - ) VALUES (?, ?, ?, ?, ?, ?, 'queued', 'Queued', ?, 0, 0, 0, 0, ?, ?, ?, ?)""", - (job_id, client_id, tally_guid, company_name, date_from, date_to, total_days, pause_seconds, - json.dumps(settings, separators=(",", ":")), now, now), - ) - cursor = start_day - while cursor <= end_day: - db.execute( - "INSERT INTO cash_payment_cache_days(job_id, voucher_date, status, updated_at_utc) VALUES (?, ?, 'pending', ?)", - (job_id, cursor.isoformat(), now), - ) - cursor += _timedelta(days=1) + try: + with _CASH_CACHE_THREADS_LOCK: + if job_id in _CASH_CACHE_CANCEL_REQUESTS: + state = _CASH_CACHE_STARTING.get(job_id) + if state is not None: + state.update({"status": "cancelled", "stage": "Cancelled before SQLite initialization", "can_cancel": False, "completed_at_utc": now, "updated_at_utc": now}) + return - 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(mapping)), - name=f"tally-cash-cache-{job_id[:8]}", - daemon=True, - ) - with _CASH_CACHE_THREADS_LOCK: - _CASH_CACHE_THREADS[job_id] = thread - thread.start() - return {"job": self._cash_payment_cache_status({"client_id": client_id, "job_id": job_id})["job"], "resumed": resume, "already_running": False} + if not self.store.exists(client_id): + raise ValueError("Accounting local storage is not initialized for this client.") + self._ensure_cash_cache_schema(client_id) + 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.") + + # Reuse a paused/interrupted identical cache only inside the worker. + with self.store.connect(client_id) as db: + prior = db.execute( + """SELECT * FROM cash_payment_cache_jobs + WHERE client_id=? AND tally_guid=? AND date_from=? AND date_to=? + ORDER BY started_at_utc DESC LIMIT 1""", + (client_id, tally_guid, date_from, date_to), + ).fetchone() + prior = dict(prior) if prior else None + resume = bool(prior and str(prior.get("status") or "") in {"failed", "paused", "interrupted", "running", "cooling", "cancelled"}) + if resume: + db_job_id = str(prior.get("job_id")) + # Keep the instant-ACK id stable in the UI by copying the + # prior checkpoints into this new job id in one transaction. + db.execute( + """INSERT INTO cash_payment_cache_jobs( + job_id, client_id, tally_guid, company_name, date_from, date_to, status, stage, + total_days, completed_days, vouchers_cached, tally_requests, retry_count, + pause_seconds, settings_json, cash_ledgers_json, started_at_utc, updated_at_utc + ) VALUES (?, ?, ?, ?, ?, ?, 'queued', 'Queued for resume', ?, ?, ?, 0, 0, ?, ?, ?, ?, ?)""", + (job_id, client_id, tally_guid, company_name, date_from, date_to, total_days, + int(prior.get("completed_days") or 0), int(prior.get("vouchers_cached") or 0), pause_seconds, + json.dumps(settings, separators=(",", ":")), str(prior.get("cash_ledgers_json") or "[]"), now, now), + ) + rows = db.execute( + "SELECT voucher_date, status, voucher_count, attempts, last_error FROM cash_payment_cache_days WHERE job_id=? ORDER BY voucher_date", + (db_job_id,), + ).fetchall() + db.executemany( + "INSERT INTO cash_payment_cache_days(job_id, voucher_date, status, voucher_count, attempts, last_error, updated_at_utc) VALUES (?, ?, ?, ?, ?, ?, ?)", + [(job_id, r[0], r[1], int(r[2] or 0), int(r[3] or 0), r[4], now) for r in rows], + ) + else: + db.execute( + """INSERT INTO cash_payment_cache_jobs( + job_id, client_id, tally_guid, company_name, date_from, date_to, status, stage, + total_days, completed_days, vouchers_cached, tally_requests, retry_count, + pause_seconds, settings_json, started_at_utc, updated_at_utc + ) VALUES (?, ?, ?, ?, ?, ?, 'queued', 'Queued', ?, 0, 0, 0, 0, ?, ?, ?, ?)""", + (job_id, client_id, tally_guid, company_name, date_from, date_to, total_days, pause_seconds, + json.dumps(settings, separators=(",", ":")), now, now), + ) + dates = [] + cursor = start_day + while cursor <= end_day: + dates.append((job_id, cursor.isoformat(), "pending", now)) + cursor += _timedelta(days=1) + db.executemany( + "INSERT INTO cash_payment_cache_days(job_id, voucher_date, status, updated_at_utc) VALUES (?, ?, ?, ?)", + dates, + ) + + payload["company_name"] = company_name + with _CASH_CACHE_THREADS_LOCK: + _CASH_CACHE_STARTING.pop(job_id, None) + self._run_cash_payment_cache_job(job_id, payload, dict(mapping)) + except Exception as exc: + self.logger.exception("Cash payment cache initialization failed job_id=%s: %s", job_id, exc) + failed_at = datetime.now(timezone.utc).isoformat() + try: + if self.store.exists(client_id): + self._ensure_cash_cache_schema(client_id) + row = self._cash_cache_job_row(client_id, job_id) + if row: + self._cash_cache_update(client_id, job_id, status="paused", stage="Initialization paused safely", error_message=str(exc), completed_at_utc=failed_at) + except Exception: + pass + with _CASH_CACHE_THREADS_LOCK: + state = _CASH_CACHE_STARTING.get(job_id) + if state is not None: + state.update({"status": "paused", "stage": "Initialization paused safely", "error": str(exc), "can_cancel": False, "can_resume": True, "completed_at_utc": failed_at, "updated_at_utc": failed_at}) + _CASH_CACHE_THREADS.pop(job_id, 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")) @@ -948,6 +1025,13 @@ class AgentCommandProcessor: 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() + with _CASH_CACHE_THREADS_LOCK: + mem = _CASH_CACHE_STARTING.get(job_id) + if mem is not None: + _CASH_CACHE_CANCEL_REQUESTS.add(job_id) + mem["stage"] = "Cancellation requested; stopping safely" + mem["updated_at_utc"] = datetime.now(timezone.utc).isoformat() + return {"job": dict(mem), "cancel_requested": True} 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.") @@ -962,6 +1046,10 @@ class AgentCommandProcessor: def _cash_payment_cache_status(self, payload: dict[str, Any]) -> dict[str, Any]: client_id = int(payload.get("client_id")) job_id = str(payload.get("job_id") or "").strip() + with _CASH_CACHE_THREADS_LOCK: + mem = _CASH_CACHE_STARTING.get(job_id) + if mem is not None: + return {"job": dict(mem), "agent": self._agent_info()} 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.") @@ -990,7 +1078,7 @@ class AgentCommandProcessor: "background": True, "sqlite_cache": True, "can_resume": str(row.get("status") or "") in {"paused", "failed", "interrupted", "cancelled"}, - "can_cancel": str(row.get("status") or "") in {"queued", "running", "cooling"}, + "can_cancel": str(row.get("status") or "") in {"queued", "running", "cooling", "starting"}, }, "agent": self._agent_info(), }