Add instant ACK for async cash payment cache jobs

This commit is contained in:
A R R R Associates
2026-09-04 00:19:59 +05:30
parent 471205053a
commit 5d13373e36
3 changed files with 158 additions and 70 deletions
@@ -1,2 +1,2 @@
__version__ = "1.22.9"
__version__ = "1.22.10"
AGENT_NAME = "ERP Local Agent"
@@ -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(),
}