diff --git a/app/modules/accounting/cash_payment_ui.py b/app/modules/accounting/cash_payment_ui.py index 19b38d2..5a6db0c 100644 --- a/app/modules/accounting/cash_payment_ui.py +++ b/app/modules/accounting/cash_payment_ui.py @@ -3,6 +3,7 @@ from __future__ import annotations from datetime import date, datetime from decimal import Decimal, InvalidOperation from urllib.parse import quote +import uuid from fastapi import APIRouter, Request from fastapi.responses import RedirectResponse, JSONResponse @@ -252,10 +253,11 @@ def cash_payment_review( @router.post("/ledgers") async def cash_payment_ledgers(request: Request): - """Load Cash ledgers from the common client .act SQLite master snapshot. + """Queue Cash-ledger preparation and return a durable VPS-known job id. - Normal loads are local-only. A live Tally master refresh happens only when - the user explicitly clicks Refresh Master Data or when no local snapshot exists. + The Local Agent reads the common client .act master in the background. Normal + loads are SQLite-only. Refresh Master Data is the only path that is allowed to + contact TallyPrime for a master refresh. """ form = await request.form() validate_csrf(request, str(form.get("csrf_token") or "")) @@ -284,37 +286,72 @@ async def cash_payment_ledgers(request: Request): if not node or not _node_online(node): return JSONResponse({"ok": False, "error": "ERP Local Agent is offline."}, status_code=409) + # The VPS generates the job id before sending the command. Even when a + # WebSocket ACK is lost, the browser can poll this known id and retrieve the + # ledger payload once the Local Agent has completed the SQLite read. + job_id = "CASHMASTER-" + uuid.uuid4().hex[:16].upper() try: fy = _financial_year_for_date(date.fromisoformat(date_from)) payload = { **_accounting_storage_payload(selected_client, fy), "tally_guid": tally_guid, "company_name": company_name, + "ledger_scope": ledger_scope, + "force_refresh": force_refresh, + "job_id": job_id, + "requested_by_user_id": int(user.id), } - # Reuse the common accounting master tables already stored in the - # client .act SQLite file. This avoids repeatedly querying TallyPrime - # just to populate a ledger selector. The Local Agent performs a live - # master refresh only when force_refresh is explicitly requested, or - # when no usable local snapshot exists yet. - payload["ledger_scope"] = ledger_scope - payload["force_refresh"] = force_refresh result = request_agent_command( node.node_code, "accounting_cash_payment_cash_ledgers", payload, - timeout_seconds=180 if force_refresh else 45, + timeout_seconds=12, ) - except Exception as exc: - return JSONResponse({"ok": False, "error": str(exc)}, status_code=502) - if not result.get("ok"): - return JSONResponse({"ok": False, "error": str(result.get("error") or "Could not load Tally ledgers.")}, status_code=409) - return JSONResponse({"ok": True, **(result.get("result") or {})}) + if result.get("ok"): + body = result.get("result") or {} + job = body.get("job") or { + "job_id": job_id, + "status": "queued", + "stage": "Local Agent accepted Cash ledger job", + "percent": 5, + } + return JSONResponse({"ok": True, "accepted": True, "job": job, **body}) + return JSONResponse({ + "ok": False, + "error": str(result.get("error") or "Local Agent rejected the Cash ledger request."), + }, status_code=409) + except Exception: + # A transport timeout does not mean the Local Agent failed. The job id + # is durable and known to the VPS; return it immediately and let the + # browser poll. This is specifically designed for intermittent tunnel + # keepalive/ACK loss. + return JSONResponse({ + "ok": True, + "accepted": True, + "transport_pending": True, + "job": { + "job_id": job_id, + "status": "queued", + "stage": "Command sent; waiting for Local Agent acknowledgment", + "percent": 5, + "company_name": company_name, + }, + }) finally: db.close() @router.get("/ledgers/progress") -def cash_payment_ledgers_progress(request: Request, client_id: int, job_id: str, date_from: str): +def cash_payment_ledgers_progress( + request: Request, + client_id: int, + job_id: str, + date_from: str, + tally_guid: str = "", + company_name: str = "", + ledger_scope: str = "cash", + force_refresh: int = 0, +): db = CommonSessionLocal() try: user, response = _require_partner(request, db, "accounting.tally.view") @@ -326,23 +363,79 @@ def cash_payment_ledgers_progress(request: Request, client_id: int, job_id: str, return JSONResponse({"ok": False, "error": "Client not found."}, status_code=404) node = get_active_storage_node_for_branch(db, scope.tenant_id, scope.branch_id) if not node or not _node_online(node): - return JSONResponse({"ok": False, "error": "ERP Local Agent is offline."}, status_code=409) + return JSONResponse({ + "ok": True, + "transport_pending": True, + "job": { + "job_id": job_id, + "status": "running", + "stage": "Waiting for ERP Local Agent tunnel", + "percent": 85, + "company_name": company_name, + }, + }) + fy = _financial_year_for_date(date.fromisoformat(date_from)) + payload = { + **_accounting_storage_payload(selected_client, fy), + "job_id": str(job_id), + "tally_guid": str(tally_guid or ""), + "company_name": str(company_name or ""), + "ledger_scope": str(ledger_scope or "cash"), + "force_refresh": bool(force_refresh), + "requested_by_user_id": int(user.id), + } try: - fy = _financial_year_for_date(date.fromisoformat(date_from)) result = request_agent_command( node.node_code, "accounting_cash_payment_ledgers_status", - { - **_accounting_storage_payload(selected_client, fy), - "job_id": str(job_id), - }, - timeout_seconds=15, + payload, + timeout_seconds=10, ) - except Exception as exc: - return JSONResponse({"ok": False, "error": str(exc)}, status_code=502) - if not result.get("ok"): - return JSONResponse({"ok": False, "error": str(result.get("error") or "Could not read ledger-sync progress.")}, status_code=409) - return JSONResponse({"ok": True, **(result.get("result") or {})}) + except Exception: + return JSONResponse({ + "ok": True, + "transport_pending": True, + "job": { + "job_id": job_id, + "status": "running", + "stage": "Local work complete or running; waiting to transfer result to VPS", + "percent": 90, + "company_name": company_name, + }, + }) + + if result.get("ok"): + return JSONResponse({"ok": True, **(result.get("result") or {})}) + + error = str(result.get("error") or "") + # If the initial start command itself was lost before reaching the Local + # Agent, resend the same deterministic job id. The Local Agent treats this + # idempotently, so this cannot create duplicate work. + if "not found" in error.casefold(): + try: + restarted = request_agent_command( + node.node_code, + "accounting_cash_payment_cash_ledgers", + payload, + timeout_seconds=10, + ) + if restarted.get("ok"): + return JSONResponse({"ok": True, **(restarted.get("result") or {})}) + except Exception: + pass + return JSONResponse({ + "ok": True, + "transport_pending": True, + "job": { + "job_id": job_id, + "status": "queued", + "stage": "Retrying Cash ledger job delivery to Local Agent", + "percent": 10, + "company_name": company_name, + }, + }) + + return JSONResponse({"ok": False, "error": error or "Could not read Cash ledger progress."}, status_code=409) finally: db.close() diff --git a/app/modules/accounting/templates/accounting/cash_payment_review.html b/app/modules/accounting/templates/accounting/cash_payment_review.html index d4890c7..b5a49dc 100644 --- a/app/modules/accounting/templates/accounting/cash_payment_review.html +++ b/app/modules/accounting/templates/accounting/cash_payment_review.html @@ -56,30 +56,35 @@
- + -
- + + + + +
-

Reads the common ledger master already stored in the client .act SQLite database. TallyPrime is contacted only when you click “Refresh Master Data” or when no local master snapshot exists yet.

+ +

Load Cash Ledgers reads the common master from the client .act SQLite database only. Refresh Master Data explicitly refreshes that common master from TallyPrime.

-
+
1. Local Master
.act SQLite
2. Cash Filter
Cash-in-Hand
-
3. Ledger Ready
No Tally query
+
3. Transfer to VPS
Tunnel result
+
4. Ledger Ready
Dropdown
-
Master source: local .act SQLite. Use Refresh Master Data only when Tally masters have changed.
+
Master source: local .act SQLite. Normal loading never queries TallyPrime.
@@ -241,11 +246,12 @@ const masterStep1 = document.getElementById('cash-master-step-1'); const masterStep2 = document.getElementById('cash-master-step-2'); const masterStep3 = document.getElementById('cash-master-step-3'); + const masterStep4 = document.getElementById('cash-master-step-4'); const masterMeta = document.getElementById('cash-master-meta'); let ledgerPollTimer = null; function paintMasterSteps(step) { - [masterStep1, masterStep2, masterStep3].forEach(function(el, idx) { + [masterStep1, masterStep2, masterStep3, masterStep4].forEach(function(el, idx) { if (!el) return; el.className = 'rounded-lg border px-2 py-2 text-center ' + (idx < step ? 'border-emerald-200 bg-emerald-50 text-emerald-900' @@ -286,7 +292,7 @@ from_cache: Boolean(data.from_cache) }; renderLedgerOptions(); - paintMasterSteps(3); + paintMasterSteps(4); if (masterMeta) { const synced = data.master_synced_at_utc ? new Date(data.master_synced_at_utc).toLocaleString() : 'time not available'; const source = data.master_refreshed ? 'TallyPrime → local .act SQLite (refreshed now)' : 'local .act SQLite'; @@ -300,27 +306,51 @@ ledgerProgressBar.style.width=p+'%'; ledgerProgressPct.textContent=Math.round(p)+'%'; ledgerProgressStage.textContent=job.stage||job.status||'Working…'; - if (p >= 90) paintMasterSteps(3); + if (p >= 100) paintMasterSteps(4); + else if (p >= 90) paintMasterSteps(3); else if (p >= 60) paintMasterSteps(2); else if (p >= 15) paintMasterSteps(1); } - async function pollLedgers(jobId, clientId, dateFrom) { + async function pollLedgers(jobId, clientId, dateFrom, tallyGuid, companyName, selectedScope, forceRefresh) { try { const url=new URL('/tools/accounting/cash-payments/ledgers/progress',window.location.origin); - url.searchParams.set('client_id',clientId); url.searchParams.set('job_id',jobId); url.searchParams.set('date_from',dateFrom); + url.searchParams.set('client_id',clientId); + url.searchParams.set('job_id',jobId); + url.searchParams.set('date_from',dateFrom); + url.searchParams.set('tally_guid',tallyGuid || ''); + url.searchParams.set('company_name',companyName || ''); + url.searchParams.set('ledger_scope',selectedScope || 'cash'); + url.searchParams.set('force_refresh',forceRefresh ? '1' : '0'); const response=await fetch(url.toString(),{headers:{'Accept':'application/json'},cache:'no-store'}); const data=await response.json(); - if(!response.ok||!data.ok) throw new Error(data.error||'Could not read Tally ledger progress.'); - const job=data.job||{}; updateLedgerProgress(job); - if(job.status==='completed'){ - populateLedgers(data); ledgerButton.disabled=false; if(ledgerRefreshButton) ledgerRefreshButton.disabled=false; - window.setTimeout(()=>ledgerProgress.classList.add('hidden'),1800); return; + if(!response.ok||!data.ok) throw new Error(data.error||'Could not read Cash ledger transfer progress.'); + const job=data.job||{}; + updateLedgerProgress(job); + if(data.transport_pending){ + ledgerHelp.textContent='Local Agent work is queued/running. Waiting for the tunnel to deliver the ledger result to the VPS…'; } - if(job.status==='failed') throw new Error(job.error||'Tally ledger sync failed.'); - ledgerPollTimer=window.setTimeout(()=>pollLedgers(jobId,clientId,dateFrom),1500); + if(job.status==='completed' && Array.isArray(data.ledgers) && data.ledgers.length >= 0){ + populateLedgers(data); + ledgerButton.disabled=false; + if(ledgerRefreshButton) ledgerRefreshButton.disabled=false; + ledgerProgress.classList.remove('hidden'); + ledgerProgressBar.style.width='100%'; + ledgerProgressPct.textContent='100%'; + ledgerProgressStage.textContent=String(data.cash_ledger_count || 0)+' Cash-in-Hand ledger(s) transmitted to VPS'; + paintMasterSteps(4); + ledgerHelp.textContent=String(data.cash_ledger_count || 0)+' Cash-in-Hand ledger(s) received by VPS and ready for selection.'; + window.setTimeout(()=>ledgerProgress.classList.add('hidden'),1800); + return; + } + if(job.status==='failed') throw new Error(job.error||'Cash ledger preparation failed.'); + ledgerPollTimer=window.setTimeout( + ()=>pollLedgers(jobId,clientId,dateFrom,tallyGuid,companyName,selectedScope,forceRefresh), + 1500 + ); } catch(err) { - ledgerHelp.textContent=err&&err.message?err.message:String(err); ledgerButton.disabled=false; + ledgerHelp.textContent=err&&err.message?err.message:String(err); + ledgerButton.disabled=false; if(ledgerRefreshButton) ledgerRefreshButton.disabled=false; ledgerProgressStage.textContent=ledgerHelp.textContent; } @@ -364,14 +394,23 @@ ledgerProgressStage.textContent = data.cash_only ? String(data.ledger_count || 0)+' Cash-in-Hand ledger(s) ready from SQLite' : String(data.ledger_count || 0)+' ledger(s) ready from SQLite'; - paintMasterSteps(3); + paintMasterSteps(4); window.setTimeout(()=>ledgerProgress.classList.add('hidden'),1200); return; } const job=data.job||{}; if(!job.job_id) throw new Error('Local Agent did not return a ledger-sync job id.'); updateLedgerProgress(job); - pollLedgers(job.job_id,String(fd.get('client_id')||''),String(fd.get('date_from')||'')); + if(data.transport_pending){ ledgerHelp.textContent='Command sent. Waiting for the Local Agent result to cross the tunnel back to VPS…'; } + pollLedgers( + job.job_id, + String(fd.get('client_id')||''), + String(fd.get('date_from')||''), + String(fd.get('tally_guid')||''), + companyName, + selectedScope, + Boolean(forceRefresh) + ); } catch(err) { ledgerHelp.textContent = err && err.message ? err.message : String(err); ledgerButton.disabled=false; if(ledgerRefreshButton) ledgerRefreshButton.disabled=false; diff --git a/app/modules/documents/agent_package.py b/app/modules/documents/agent_package.py index 18f0a77..4737f57 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.22" +ERP_LOCAL_AGENT_VERSION = "1.22.23" 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 aa7ece2..f65a262 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.22" +__version__ = "1.22.23" 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 00ae1bf..034e3c0 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 @@ -181,6 +181,7 @@ class AgentCommandProcessor: "cash_payment_selected_company_resolver_capability": True, "cash_payment_common_act_master_capability": True, "cash_payment_sqlite_first_ledger_capability": True, + "cash_payment_durable_vps_transfer_capability": True, "opening_balance_balance_sheet_only_capability": True, "tally_writeback_capability": True, } @@ -744,248 +745,479 @@ class AgentCommandProcessor: } - def _cash_payment_cash_ledgers(self, payload: dict[str, Any]) -> dict[str, Any]: - """Return Cash ledgers from the common client .act SQLite master snapshot. + def _cash_common_master_snapshot( + self, + client_id: int, + tally_guid: str, + company_name: str, + ) -> tuple[str, str, list[dict[str, Any]], list[dict[str, Any]], str]: + """Read the shared accounting master snapshot from the client .act SQLite DB.""" + with self.store.connect(client_id) as db: + guid = str(tally_guid or "").strip() + name = str(company_name or "").strip() - The normal Cash Payment selector does not contact TallyPrime. It reuses the - same common ``tally_groups`` and ``tally_ledgers`` tables populated by the - accounting master/opening-balance workflow. A live Tally master refresh is - performed only when the user explicitly requests Refresh Master Data, or - when no usable local snapshot exists for the selected company. + if guid: + count_row = db.execute( + "SELECT COUNT(*) FROM tally_ledgers WHERE tally_guid=?", + (guid,), + ).fetchone() + if not count_row or int(count_row[0] or 0) == 0: + guid = "" + + if not guid and name: + guid_row = db.execute( + """SELECT tally_guid, company_name, MAX(synced_at_utc) AS synced_at + FROM tally_ledgers + WHERE lower(trim(company_name))=lower(trim(?)) + GROUP BY tally_guid, company_name + ORDER BY MAX(synced_at_utc) DESC + LIMIT 1""", + (name,), + ).fetchone() + if guid_row: + guid = str(guid_row["tally_guid"] or "") + name = str(guid_row["company_name"] or name) + + if not guid: + return "", name, [], [], "" + + group_rows = [ + dict(row) for row in db.execute( + """SELECT master_guid AS guid, name, parent, reserved_name, synced_at_utc + FROM tally_groups + WHERE tally_guid=? + ORDER BY name COLLATE NOCASE""", + (guid,), + ).fetchall() + ] + ledger_rows = [ + dict(row) for row in db.execute( + """SELECT master_guid AS guid, name, parent, reserved_name, + opening_balance, closing_balance, synced_at_utc + FROM tally_ledgers + WHERE tally_guid=? + ORDER BY name COLLATE NOCASE""", + (guid,), + ).fetchall() + ] + sync_row = db.execute( + "SELECT MAX(synced_at_utc) AS synced_at FROM tally_ledgers WHERE tally_guid=?", + (guid,), + ).fetchone() + synced_at = str((sync_row["synced_at"] if sync_row else "") or "") + company_row = db.execute( + """SELECT company_name + FROM tally_ledgers + WHERE tally_guid=? AND trim(company_name)<>'' + ORDER BY synced_at_utc DESC LIMIT 1""", + (guid,), + ).fetchone() + if company_row: + name = str(company_row["company_name"] or name) + return guid, name, group_rows, ledger_rows, synced_at + + def _cash_common_master_ledgers( + self, + group_rows: list[dict[str, Any]], + ledger_rows: list[dict[str, Any]], + ) -> tuple[list[dict[str, Any]], list[dict[str, Any]]]: + """Return all ledgers plus the recursively resolved Cash-in-Hand subset.""" + cash_group_keys = {"cash-in-hand", "cash in hand"} + changed = True + while changed: + changed = False + normalized_cash = { + x.casefold().replace("-", " ").strip() for x in cash_group_keys + } + for group in group_rows or []: + name = str(group.get("name") or "").strip() + parent = str(group.get("parent") or "").strip() + if not name: + continue + parent_key = parent.casefold().replace("-", " ").strip() + name_key = name.casefold().replace("-", " ").strip() + if parent_key in normalized_cash and name_key not in normalized_cash: + cash_group_keys.add(name_key) + changed = True + + all_ledgers = self._cash_ledger_rows(list(ledger_rows or [])) + normalized_groups = { + x.casefold().replace("-", " ").strip() for x in cash_group_keys + } + cash_ledgers: list[dict[str, Any]] = [] + for row in all_ledgers: + parent_key = str(row.get("parent") or "").casefold().replace("-", " ").strip() + reserved = str(row.get("reserved_name") or "").strip().casefold() + name_key = str(row.get("name") or "").strip().casefold() + is_cash = ( + parent_key in normalized_groups + or reserved == "cash" + or name_key == "cash" + ) + row["is_cash_candidate"] = bool(is_cash) + if is_cash: + cash_ledgers.append(row) + + all_ledgers.sort( + key=lambda x: ( + not bool(x.get("is_cash_candidate")), + str(x.get("name") or "").casefold(), + ) + ) + cash_ledgers.sort(key=lambda x: str(x.get("name") or "").casefold()) + return all_ledgers, cash_ledgers + + def _cash_payment_cash_ledgers(self, payload: dict[str, Any]) -> dict[str, Any]: + """Queue a durable SQLite-first Cash ledger transfer and ACK immediately. + + The VPS supplies a deterministic job id. The Local Agent stores progress in + the client .act SQLite database, so a lost WebSocket response cannot lose the + work. Normal loads never query TallyPrime. Refresh Master Data performs the + existing common master refresh in the background and then exposes the result + through the same status endpoint. """ client_id = int(payload.get("client_id")) expected_company_name = str(payload.get("company_name") or "").strip() requested_guid = str(payload.get("tally_guid") or "").strip() - ledger_scope = str(payload.get("ledger_scope") or "cash").strip().lower() force_refresh = bool(payload.get("force_refresh")) - if ledger_scope not in {"cash", "all"}: - ledger_scope = "cash" + job_id = str(payload.get("job_id") or "").strip() + if not job_id: + job_id = "CASHMASTER-" + uuid.uuid4().hex[:16].upper() + if not job_id.startswith("CASHMASTER-"): + job_id = "CASHMASTER-" + job_id[-16:].upper() + + self._ensure_cash_ledger_schema(client_id) + with self.store.connect(client_id) as db: + existing = db.execute( + """SELECT status FROM cash_payment_ledger_jobs + WHERE job_id=? AND client_id=?""", + (job_id, client_id), + ).fetchone() + if existing: + return self._cash_payment_common_ledgers_status( + {**payload, "job_id": job_id} + ) + + now = datetime.now(timezone.utc).isoformat() + db.execute( + """INSERT INTO cash_payment_ledger_jobs + (job_id,client_id,tally_guid,company_name,status,stage,percent, + ledger_count,error,started_at_utc,finished_at_utc) + VALUES (?,?,?,?,?,?,?,?,?,?,?)""", + ( + job_id, + client_id, + requested_guid, + expected_company_name, + "queued", + "Queued local SQLite master read", + 5, + 0, + "", + now, + "", + ), + ) + db.commit() - progress_id = "CASHMASTER-" + uuid.uuid4().hex[:12].upper() _cash_progress_snapshot({ - "job_id": progress_id, - "job_kind": "cash_ledger_discovery", - "status": "running", - "stage": "Reading common ledger master from local SQLite", - "percent": 20, + "job_id": job_id, + "job_kind": "cash_ledger_transfer", + "status": "queued", + "stage": "Queued local SQLite master read", + "percent": 5, "company_name": expected_company_name, "ledger_count": 0, "source": "Local .act SQLite", }) - def _read_local_snapshot(tally_guid: str, company_name: str) -> tuple[str, str, list[dict[str, Any]], list[dict[str, Any]], str]: - with self.store.connect(client_id) as db: - guid = str(tally_guid or "").strip() - name = str(company_name or "").strip() - - if guid: - count_row = db.execute( - "SELECT COUNT(*) FROM tally_ledgers WHERE tally_guid=?", - (guid,), - ).fetchone() - if not count_row or int(count_row[0] or 0) == 0: - guid = "" - - if not guid and name: - guid_row = db.execute( - """SELECT tally_guid, company_name, MAX(synced_at_utc) AS synced_at - FROM tally_ledgers - WHERE lower(trim(company_name))=lower(trim(?)) - GROUP BY tally_guid, company_name - ORDER BY MAX(synced_at_utc) DESC - LIMIT 1""", - (name,), - ).fetchone() - if guid_row: - guid = str(guid_row["tally_guid"] or "") - name = str(guid_row["company_name"] or name) - - if not guid: - return "", name, [], [], "" - - group_rows = [ - dict(row) for row in db.execute( - """SELECT master_guid AS guid, name, parent, reserved_name, synced_at_utc - FROM tally_groups - WHERE tally_guid=? - ORDER BY name COLLATE NOCASE""", - (guid,), - ).fetchall() - ] - ledger_rows = [ - dict(row) for row in db.execute( - """SELECT master_guid AS guid, name, parent, reserved_name, - opening_balance, closing_balance, synced_at_utc - FROM tally_ledgers - WHERE tally_guid=? - ORDER BY name COLLATE NOCASE""", - (guid,), - ).fetchall() - ] - sync_row = db.execute( - """SELECT MAX(synced_at_utc) AS synced_at - FROM tally_ledgers - WHERE tally_guid=?""", - (guid,), - ).fetchone() - synced_at = str((sync_row["synced_at"] if sync_row else "") or "") - company_row = db.execute( - """SELECT company_name - FROM tally_ledgers - WHERE tally_guid=? AND trim(company_name)<>'' - ORDER BY synced_at_utc DESC LIMIT 1""", - (guid,), - ).fetchone() - if company_row: - name = str(company_row["company_name"] or name) - return guid, name, group_rows, ledger_rows, synced_at - - try: - resolved_guid, resolved_name, group_rows, ledger_rows, synced_at = _read_local_snapshot( - requested_guid, expected_company_name - ) - refreshed = False - - if force_refresh or not ledger_rows: + def worker() -> None: + job_guid = requested_guid + job_company = expected_company_name + try: + with self.store.connect(client_id) as db: + db.execute( + """UPDATE cash_payment_ledger_jobs + SET status='running',stage=?,percent=15 + WHERE job_id=?""", + ("Reading common ledger master from local SQLite", job_id), + ) + db.commit() _cash_progress_snapshot({ - "job_id": progress_id, - "job_kind": "cash_ledger_discovery", + "job_id": job_id, + "job_kind": "cash_ledger_transfer", "status": "running", - "stage": "Refreshing common accounting masters from TallyPrime", - "percent": 45, - "company_name": expected_company_name, - "ledger_count": len(ledger_rows), - "source": "TallyPrime → .act SQLite", + "stage": "Reading common ledger master from local SQLite", + "percent": 15, + "company_name": job_company, + "ledger_count": 0, + "source": "Local .act SQLite", }) - company, actual_guid, guid_refreshed = self._cash_payment_selected_company(payload) - actual_guid = str(actual_guid or requested_guid or "").strip() - if not actual_guid: - raise ValueError("Tally did not return a GUID for the selected company.") - masters = self.tally.fetch_accounting_masters(company.name) - mapping = self.store.get_active_mapping_by_guid(client_id, actual_guid) - self.store.replace_master_snapshot( - client_id, - mapping={**mapping, "tally_guid": actual_guid, "company_name": company.name}, - 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 - ), - ) - refreshed = True - resolved_guid, resolved_name, group_rows, ledger_rows, synced_at = _read_local_snapshot( - actual_guid, company.name + resolved_guid, resolved_name, group_rows, ledger_rows, synced_at = ( + self._cash_common_master_snapshot(client_id, job_guid, job_company) ) + + if force_refresh: + with self.store.connect(client_id) as db: + db.execute( + """UPDATE cash_payment_ledger_jobs + SET stage=?,percent=35 + WHERE job_id=?""", + ("Refreshing common accounting masters from TallyPrime", job_id), + ) + db.commit() + _cash_progress_snapshot({ + "job_id": job_id, + "job_kind": "cash_ledger_transfer", + "status": "running", + "stage": "Refreshing common accounting masters from TallyPrime", + "percent": 35, + "company_name": job_company, + "ledger_count": len(ledger_rows), + "source": "TallyPrime → .act SQLite", + }) + + company, actual_guid, _guid_refreshed = self._cash_payment_selected_company(payload) + actual_guid = str(actual_guid or job_guid or "").strip() + if not actual_guid: + raise ValueError("Tally did not return a GUID for the selected company.") + masters = self.tally.fetch_accounting_masters(company.name) + mapping = self.store.get_active_mapping_by_guid(client_id, actual_guid) + self.store.replace_master_snapshot( + client_id, + mapping={ + **mapping, + "tally_guid": actual_guid, + "company_name": company.name, + }, + 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 + ), + ) + resolved_guid, resolved_name, group_rows, ledger_rows, synced_at = ( + self._cash_common_master_snapshot(client_id, actual_guid, company.name) + ) + job_guid = resolved_guid or actual_guid + job_company = resolved_name or company.name + else: + if not ledger_rows: + raise ValueError( + "No common ledger master is stored locally for this Tally company. " + "Click Refresh Master Data once; normal Load Cash Ledgers never queries TallyPrime." + ) + job_guid = resolved_guid or job_guid + job_company = resolved_name or job_company + if not ledger_rows: - raise ValueError("The refreshed Tally master snapshot did not contain any ledgers.") - else: - guid_refreshed = bool( - requested_guid and resolved_guid and requested_guid != resolved_guid + raise ValueError("The common accounting master snapshot contains no ledgers.") + + with self.store.connect(client_id) as db: + db.execute( + """UPDATE cash_payment_ledger_jobs + SET tally_guid=?,company_name=?,stage=?,percent=60,ledger_count=? + WHERE job_id=?""", + ( + job_guid, + job_company, + "Filtering Cash-in-Hand hierarchy from local SQLite", + len(ledger_rows), + job_id, + ), + ) + db.commit() + _cash_progress_snapshot({ + "job_id": job_id, + "job_kind": "cash_ledger_transfer", + "status": "running", + "stage": "Filtering Cash-in-Hand hierarchy from local SQLite", + "percent": 60, + "company_name": job_company, + "ledger_count": len(ledger_rows), + "source": "Local .act SQLite", + }) + + all_ledgers, cash_ledgers = self._cash_common_master_ledgers( + group_rows, ledger_rows ) - _cash_progress_snapshot({ - "job_id": progress_id, - "job_kind": "cash_ledger_discovery", - "status": "running", - "stage": "Filtering Cash-in-Hand hierarchy from local SQLite", - "percent": 75, - "company_name": resolved_name or expected_company_name, - "ledger_count": len(ledger_rows), - "source": "Local .act SQLite", - }) + with self.store.connect(client_id) as db: + db.execute( + """UPDATE cash_payment_ledger_jobs + SET status='completed',tally_guid=?,company_name=?, + stage=?,percent=90,ledger_count=?,error='', + finished_at_utc=? + WHERE job_id=?""", + ( + job_guid, + job_company, + f"{len(cash_ledgers)} Cash-in-Hand ledger(s) ready for VPS transfer", + len(all_ledgers), + datetime.now(timezone.utc).isoformat(), + job_id, + ), + ) + db.commit() - cash_group_keys = {"cash-in-hand", "cash in hand"} - changed = True - while changed: - changed = False - normalized_cash = { - x.casefold().replace("-", " ").strip() for x in cash_group_keys - } - for group in group_rows: - name = str(group.get("name") or "").strip() - parent = str(group.get("parent") or "").strip() - if not name: - continue - parent_key = parent.casefold().replace("-", " ").strip() - name_key = name.casefold().replace("-", " ").strip() - if parent_key in normalized_cash and name_key not in normalized_cash: - cash_group_keys.add(name_key) - changed = True - - all_ledgers = self._cash_ledger_rows(ledger_rows) - normalized_groups = { - x.casefold().replace("-", " ").strip() for x in cash_group_keys - } - cash_ledgers: list[dict[str, Any]] = [] - for row in all_ledgers: - parent_key = str(row.get("parent") or "").casefold().replace("-", " ").strip() - reserved = str(row.get("reserved_name") or "").strip().casefold() - name_key = str(row.get("name") or "").strip().casefold() - is_cash = ( - parent_key in normalized_groups - or reserved == "cash" - or name_key == "cash" + _cash_progress_snapshot({ + "job_id": job_id, + "job_kind": "cash_ledger_transfer", + "status": "completed", + "stage": f"{len(cash_ledgers)} Cash-in-Hand ledger(s) ready for VPS transfer", + "percent": 90, + "company_name": job_company, + "ledger_count": len(cash_ledgers), + "total_ledgers": len(all_ledgers), + "master_synced_at_utc": synced_at, + "source": "Local .act SQLite", + }) + except Exception as exc: + finished = datetime.now(timezone.utc).isoformat() + try: + with self.store.connect(client_id) as db: + db.execute( + """UPDATE cash_payment_ledger_jobs + SET status='failed',stage='Cash ledger preparation failed', + percent=100,error=?,finished_at_utc=? + WHERE job_id=?""", + (str(exc), finished, job_id), + ) + db.commit() + except Exception: + pass + self.logger.exception( + "Cash ledger SQLite transfer failed job_id=%s: %s", job_id, exc ) - row["is_cash_candidate"] = bool(is_cash) - if is_cash: - cash_ledgers.append(row) + _cash_progress_snapshot({ + "job_id": job_id, + "job_kind": "cash_ledger_transfer", + "status": "failed", + "stage": "Cash ledger preparation failed", + "percent": 100, + "company_name": job_company, + "ledger_count": 0, + "source": "Local .act SQLite", + "error": str(exc), + }) + finally: + with _CASH_LEDGER_THREADS_LOCK: + _CASH_LEDGER_THREADS.pop(job_id, None) - all_ledgers.sort( - key=lambda x: ( - not bool(x.get("is_cash_candidate")), - str(x.get("name") or "").casefold(), + thread = threading.Thread( + target=worker, + name=f"cash-master-{job_id[-6:]}", + daemon=True, + ) + with _CASH_LEDGER_THREADS_LOCK: + _CASH_LEDGER_THREADS[job_id] = thread + thread.start() + + # Return only the durable job record. The actual ledger list is delivered + # through status polling, so a lost ACK cannot discard a completed SQLite read. + return self._cash_payment_common_ledgers_status({**payload, "job_id": job_id}) + + def _cash_payment_common_ledgers_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() + if not job_id: + raise ValueError("Cash ledger transfer job id is required.") + self._ensure_cash_ledger_schema(client_id) + with self.store.connect(client_id) as db: + row = db.execute( + """SELECT job_id,tally_guid,company_name,status,stage,percent, + ledger_count,error,started_at_utc,finished_at_utc + FROM cash_payment_ledger_jobs + WHERE job_id=? AND client_id=?""", + (job_id, client_id), + ).fetchone() + if not row: + raise ValueError("Cash ledger transfer job was not found.") + + job = { + "job_id": str(row[0]), + "tally_guid": str(row[1] or ""), + "company_name": str(row[2] or ""), + "status": str(row[3] or ""), + "stage": str(row[4] or ""), + "percent": int(row[5] or 0), + "ledger_count": int(row[6] or 0), + "error": str(row[7] or ""), + "started_at_utc": str(row[8] or ""), + "finished_at_utc": str(row[9] or ""), + } + result: dict[str, Any] = { + "job": job, + "ledgers": [], + "cash_candidates": [], + "ledger_count": int(job["ledger_count"]), + "cash_ledger_count": 0, + "company_name": job["company_name"], + "company_guid": job["tally_guid"], + "from_cache": True, + "source": "common_act_sqlite_master", + "agent": self._agent_info(), + } + + if job["status"] == "completed": + resolved_guid, resolved_name, group_rows, ledger_rows, synced_at = ( + self._cash_common_master_snapshot( + client_id, + job["tally_guid"], + job["company_name"], ) ) - cash_ledgers.sort(key=lambda x: str(x.get("name") or "").casefold()) - returned = all_ledgers if ledger_scope == "all" else cash_ledgers - - _cash_progress_snapshot({ - "job_id": progress_id, - "job_kind": "cash_ledger_discovery", - "status": "completed", - "stage": f"{len(returned)} ledger(s) ready from local SQLite", - "percent": 100, - "company_name": resolved_name or expected_company_name, - "ledger_count": len(returned), - "source": "Local .act SQLite", - }) - return { - "status": "completed", - "company": { - "guid": resolved_guid, - "name": resolved_name or expected_company_name, - }, - "company_name": resolved_name or expected_company_name, - "company_guid": resolved_guid, - "selected_guid_refreshed": bool(guid_refreshed), - "ledgers": returned, + if not ledger_rows: + raise ValueError( + "The completed Cash ledger job no longer has a common local master snapshot." + ) + all_ledgers, cash_ledgers = self._cash_common_master_ledgers( + group_rows, ledger_rows + ) + result.update({ + "ledgers": all_ledgers, "cash_candidates": cash_ledgers, - "ledger_count": len(returned), - "total_ledgers_examined": len(all_ledgers), + "ledger_count": len(all_ledgers), "cash_ledger_count": len(cash_ledgers), - "cash_only": ledger_scope != "all", - "from_cache": not refreshed, - "master_refreshed": refreshed, + "total_ledgers_examined": len(all_ledgers), + "company_name": resolved_name or job["company_name"], + "company_guid": resolved_guid or job["tally_guid"], "master_synced_at_utc": synced_at, - "source": "common_act_sqlite_master", + "master_refreshed": bool(payload.get("force_refresh")), "read_only": True, - "agent": self._agent_info(), - } - except Exception as exc: - _cash_progress_snapshot({ - "job_id": progress_id, - "job_kind": "cash_ledger_discovery", - "status": "failed", - "stage": "Cash ledger preparation failed", - "percent": 100, - "company_name": expected_company_name, - "ledger_count": 0, - "source": "Local .act SQLite", - "error": str(exc), }) - raise + + # This status call is the actual transfer boundary: the VPS has now + # received the ledger payload. Reflect that truthfully on the dashboard. + transfer_stage = ( + f"{len(cash_ledgers)} Cash-in-Hand ledger(s) transmitted to VPS" + ) + result["job"] = {**job, "stage": transfer_stage, "percent": 100} + _cash_progress_snapshot({ + "job_id": job_id, + "job_kind": "cash_ledger_transfer", + "status": "completed", + "stage": transfer_stage, + "percent": 100, + "company_name": resolved_name or job["company_name"], + "ledger_count": len(cash_ledgers), + "total_ledgers": len(all_ledgers), + "source": "Local .act SQLite → VPS", + }) + try: + with self.store.connect(client_id) as db: + db.execute( + """UPDATE cash_payment_ledger_jobs + SET stage=?,percent=100 + WHERE job_id=?""", + (transfer_stage, job_id), + ) + db.commit() + except Exception: + pass + return result @staticmethod def _cash_ledger_rows(rows: list[dict[str, Any]]) -> list[dict[str, Any]]: @@ -1273,8 +1505,10 @@ class AgentCommandProcessor: return self._cash_payment_ledgers_status({**payload, "job_id": job_id}) def _cash_payment_ledgers_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() + if job_id.startswith("CASHMASTER-"): + return self._cash_payment_common_ledgers_status(payload) + client_id = int(payload.get("client_id")) if not job_id: raise ValueError("Cash ledger master job id is required.") self._ensure_cash_ledger_schema(client_id)