Fix cash ledger SQLite result transport to VPS

This commit is contained in:
A R R R Associates
2026-09-04 23:27:29 +05:30
parent 8536dc0457
commit 9cceec0283
5 changed files with 634 additions and 268 deletions
+1 -1
View File
@@ -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)
@@ -1,2 +1,2 @@
__version__ = "1.22.22"
__version__ = "1.22.23"
AGENT_NAME = "ERP Local Agent"
@@ -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)