{% endif %}
+
+ {% endif %}
+
{% if review %}
- {% if review.safe_paced_mode %}
-
Completed in adaptive safe-paced mode using {{ review.tally_request_batches or 0 }} successful voucher batch(es) across {{ review.tally_request_count or 0 }} Tally request(s), with {{ review.tally_retry_count or 0 }} retry/retries and {{ review.tally_fallback_batches or 0 }} adaptive fallback(s). The agent pauses approximately {{ review.tally_pause_ms or 0 }} ms between Tally calls.
+ {% if review.sqlite_cache %}
+
Analysis completed from the Local Agent SQLite cache. TallyPrime was used only for the slow background extraction and was released before this analysis ran. {{ review.tally_request_count or 0 }} Tally request(s), {{ review.tally_retry_count or 0 }} retry/retries, and approximately {{ review.tally_pause_ms or 0 }} ms cooling between successful day reads.
{% endif %}
{% for label, value in [('Vouchers reviewed', review.summary.vouchers_reviewed), ('Cash payments', review.summary.cash_payment_vouchers), ('Single-voucher exceptions', review.summary.single_voucher_exceptions), ('Same-day exceptions', review.summary.same_day_exceptions), ('Possible split patterns', review.summary.possible_split_patterns)] %}
@@ -119,60 +133,126 @@
{% endif %}
-
+
-
Cash Payment Analysis in Progress
-
Sending the review request to the ERP Local Agent…
+
Slow Tally Extraction to Local SQLite
+
Submitting a background extraction job to the ERP Local Agent…
-
5%
+
0%
-
+
-
- AgentTallyVouchersAnalysis
+
+
Days
0 / 0
+
Vouchers cached
0
+
Tally requests
0
+
Current date
—
-
Tally is deliberately read in 7-day batches with pauses between requests. If a batch times out, the agent automatically retries and reduces it to 3-day and then 1-day batches. Please keep TallyPrime and the ERP Local Agent open until this finishes.
+
Only one day is requested from TallyPrime at a time. The agent waits about 3 seconds after every successful day. If Tally responds slowly, it cools for 20 seconds and retries only once. The tunnel heartbeat stays responsive because progress is read from SQLite, not from Tally.
+
+
diff --git a/app/modules/documents/agent_package.py b/app/modules/documents/agent_package.py
index 21584d2..db3082f 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.7"
+ERP_LOCAL_AGENT_VERSION = "1.22.8"
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 f174659..32e9365 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.7"
+__version__ = "1.22.8"
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 d1901eb..88a81f2 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
@@ -4,6 +4,8 @@ from datetime import datetime, timezone, date as _dt_date, timedelta as _timedel
import re
import time
import threading
+import json
+import uuid
from typing import Any
from . import __version__
@@ -13,6 +15,12 @@ from .native_voucher_engine import NativeVoucherEngine
_CASH_TALLY_EXTRACTION_LOCK = threading.Lock()
+_CASH_CACHE_BUSY = threading.Event()
+_CASH_CACHE_THREADS: dict[str, threading.Thread] = {}
+_CASH_CACHE_THREADS_LOCK = threading.Lock()
+
+def is_tally_cache_busy() -> bool:
+ return _CASH_CACHE_BUSY.is_set()
class AgentCommandProcessor:
@@ -61,6 +69,12 @@ class AgentCommandProcessor:
result = self._bank_posting_preflight(payload)
elif action == "accounting_bank_reconciliation_extract":
result = self._bank_reconciliation_extract(payload)
+ elif action == "accounting_cash_payment_cache_start":
+ 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_analyze":
+ result = self._cash_payment_cache_analyze(payload)
elif action == "accounting_cash_payment_compliance":
result = self._cash_payment_compliance(payload)
elif action == "accounting_tds_compliance":
@@ -109,6 +123,8 @@ class AgentCommandProcessor:
"purchase_voucher_write_capability": True,
"it_depreciation_capability": True,
"cash_payment_compliance_capability": True,
+ "cash_payment_sqlite_cache_capability": True,
+ "cash_payment_background_sync_capability": True,
"tally_writeback_capability": True,
}
@@ -649,6 +665,424 @@ class AgentCommandProcessor:
}
+ def _ensure_cash_cache_schema(self, client_id: int) -> None:
+ with self.store.connect(client_id) as db:
+ db.executescript(
+ """
+ CREATE TABLE IF NOT EXISTS cash_payment_cache_jobs (
+ job_id TEXT PRIMARY KEY,
+ client_id INTEGER NOT NULL,
+ tally_guid TEXT NOT NULL,
+ company_name TEXT NOT NULL,
+ date_from TEXT NOT NULL,
+ date_to TEXT NOT NULL,
+ status TEXT NOT NULL,
+ stage TEXT NOT NULL DEFAULT '',
+ current_date TEXT NOT NULL DEFAULT '',
+ total_days INTEGER NOT NULL DEFAULT 0,
+ completed_days INTEGER NOT NULL DEFAULT 0,
+ vouchers_cached INTEGER NOT NULL DEFAULT 0,
+ tally_requests INTEGER NOT NULL DEFAULT 0,
+ retry_count INTEGER NOT NULL DEFAULT 0,
+ pause_seconds REAL NOT NULL DEFAULT 3,
+ cash_ledgers_json TEXT NOT NULL DEFAULT '[]',
+ settings_json TEXT NOT NULL DEFAULT '{}',
+ started_at_utc TEXT NOT NULL,
+ updated_at_utc TEXT NOT NULL,
+ completed_at_utc TEXT,
+ error_message TEXT
+ );
+ CREATE INDEX IF NOT EXISTS ix_cash_cache_jobs_company
+ ON cash_payment_cache_jobs(client_id, tally_guid, started_at_utc);
+ CREATE TABLE IF NOT EXISTS cash_payment_cache_days (
+ job_id TEXT NOT NULL,
+ voucher_date TEXT NOT NULL,
+ status TEXT NOT NULL,
+ voucher_count INTEGER NOT NULL DEFAULT 0,
+ attempts INTEGER NOT NULL DEFAULT 0,
+ last_error TEXT,
+ updated_at_utc TEXT NOT NULL,
+ PRIMARY KEY(job_id, voucher_date),
+ FOREIGN KEY(job_id) REFERENCES cash_payment_cache_jobs(job_id) ON DELETE CASCADE
+ );
+ CREATE INDEX IF NOT EXISTS ix_cash_cache_days_status
+ ON cash_payment_cache_days(job_id, status, voucher_date);
+ """
+ )
+
+ def _cash_cache_job_row(self, client_id: int, job_id: str = "") -> dict[str, Any] | None:
+ self._ensure_cash_cache_schema(client_id)
+ with self.store.connect(client_id) as db:
+ if job_id:
+ row = db.execute(
+ "SELECT * FROM cash_payment_cache_jobs WHERE job_id=? AND client_id=? LIMIT 1",
+ (job_id, client_id),
+ ).fetchone()
+ else:
+ row = db.execute(
+ "SELECT * FROM cash_payment_cache_jobs WHERE client_id=? ORDER BY started_at_utc DESC LIMIT 1",
+ (client_id,),
+ ).fetchone()
+ return dict(row) if row else None
+
+ def _cash_cache_update(self, client_id: int, job_id: str, **values: Any) -> None:
+ if not values:
+ return
+ values["updated_at_utc"] = datetime.now(timezone.utc).isoformat()
+ cols = list(values)
+ sql = "UPDATE cash_payment_cache_jobs SET " + ", ".join(f"{name}=?" for name in cols) + " WHERE job_id=? AND client_id=?"
+ with self.store.connect(client_id) as db:
+ db.execute(sql, [values[name] for name in cols] + [job_id, client_id])
+
+ def _cash_payment_cache_start(self, payload: dict[str, Any]) -> dict[str, Any]:
+ 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()
+ if not tally_guid:
+ raise ValueError("Select a mapped Tally company before starting extraction.")
+ try:
+ start_day = _dt_date.fromisoformat(date_from)
+ end_day = _dt_date.fromisoformat(date_to)
+ except ValueError as exc:
+ 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
+ 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}
+ 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.")
+ mapping = self.store.get_active_mapping_by_guid(client_id, tally_guid)
+ 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),
+ "split_window_days": int(payload.get("split_window_days") or 3),
+ "near_limit_percent": float(payload.get("near_limit_percent") or 80.0),
+ "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)
+
+ worker_payload = dict(payload)
+ worker_payload["tally_pause_seconds"] = pause_seconds
+ thread = threading.Thread(
+ target=self._run_cash_payment_cache_job,
+ args=(job_id, worker_payload, dict(company), company_name, 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}
+
+ 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:
+ 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="Reading cash ledger masters", current_date="")
+ # One master request only, then a generous cooling gap before vouchers.
+ masters = None
+ master_error = None
+ for attempt in range(1, 3):
+ try:
+ masters = self.tally.export_master_collection(company_name, "ledgers") or []
+ self._cash_cache_update(client_id, job_id, tally_requests=(self._cash_cache_job_row(client_id, job_id) or {}).get("tally_requests", 0) + 1)
+ break
+ except Exception as exc:
+ master_error = exc
+ row = self._cash_cache_job_row(client_id, job_id) or {}
+ self._cash_cache_update(client_id, job_id, retry_count=int(row.get("retry_count") or 0) + 1, stage="Cooling after slow Tally response")
+ time.sleep(15.0)
+ if masters is None:
+ raise ValueError(f"TallyPrime did not return ledger masters after a slow retry. Details: {master_error}")
+ cash_ledgers = set()
+ for ledger in masters:
+ name = str(ledger.get("name") or "").strip()
+ parent = str(ledger.get("parent") or "").strip().casefold().replace("-", " ")
+ reserved = str(ledger.get("reserved_name") or "").strip().casefold()
+ if name and (parent == "cash in hand" or reserved == "cash" or name.casefold() == "cash"):
+ cash_ledgers.add(name.casefold())
+ if not cash_ledgers:
+ cash_ledgers.add("cash")
+ self._cash_cache_update(client_id, job_id, cash_ledgers_json=json.dumps(sorted(cash_ledgers)), stage="Cooling before voucher extraction")
+ time.sleep(max(5.0, pause_seconds))
+
+ with self.store.connect(client_id) as db:
+ pending = db.execute(
+ "SELECT voucher_date FROM cash_payment_cache_days WHERE job_id=? AND status!='completed' ORDER BY voucher_date",
+ (job_id,),
+ ).fetchall()
+ for day_row in pending:
+ 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
+ last_exc = None
+ for attempt in range(1, 3):
+ try:
+ row = self._cash_cache_job_row(client_id, job_id) or {}
+ self._cash_cache_update(client_id, job_id, tally_requests=int(row.get("tally_requests") or 0) + 1)
+ vouchers = self.tally.export_vouchers(company_name, day, day) or []
+ break
+ except Exception as exc:
+ last_exc = exc
+ with self.store.connect(client_id) as db:
+ db.execute(
+ "UPDATE cash_payment_cache_days SET attempts=attempts+1, status='retrying', last_error=?, updated_at_utc=? WHERE job_id=? AND voucher_date=?",
+ (str(exc), datetime.now(timezone.utc).isoformat(), job_id, day),
+ )
+ row = self._cash_cache_job_row(client_id, job_id) or {}
+ self._cash_cache_update(client_id, job_id, status="cooling", stage="Tally is slow; cooling for 20 seconds before one retry", retry_count=int(row.get("retry_count") or 0) + 1)
+ time.sleep(20.0)
+ if vouchers is None:
+ with self.store.connect(client_id) as db:
+ db.execute(
+ "UPDATE cash_payment_cache_days SET status='failed', last_error=?, updated_at_utc=? WHERE job_id=? AND voucher_date=?",
+ (str(last_exc), datetime.now(timezone.utc).isoformat(), job_id, day),
+ )
+ raise ValueError(f"Extraction paused at {day} because TallyPrime did not answer two deliberately slow one-day requests. Restart/verify TallyPrime and run again to resume from this date. Details: {last_exc}")
+
+ self.store.replace_transaction_snapshot(
+ client_id,
+ mapping=mapping,
+ transactions={"date_from": day, "date_to": day, "vouchers": list(vouchers)},
+ requested_by_user_id=int(requested_by) if requested_by not in (None, "") else None,
+ )
+ with self.store.connect(client_id) as db:
+ db.execute(
+ "UPDATE cash_payment_cache_days SET status='completed', voucher_count=?, attempts=attempts+1, last_error=NULL, updated_at_utc=? WHERE job_id=? AND voucher_date=?",
+ (len(vouchers), datetime.now(timezone.utc).isoformat(), job_id, day),
+ )
+ agg = db.execute(
+ "SELECT COUNT(*), COALESCE(SUM(voucher_count),0) FROM cash_payment_cache_days WHERE job_id=? AND status='completed'",
+ (job_id,),
+ ).fetchone()
+ self._cash_cache_update(client_id, job_id, completed_days=int(agg[0]), vouchers_cached=int(agg[1]), stage="Cooling between one-day Tally reads")
+ time.sleep(pause_seconds)
+
+ self._cash_cache_update(client_id, job_id, status="completed", stage="SQLite cache ready for analysis", current_date="", completed_at_utc=datetime.now(timezone.utc).isoformat())
+ except Exception as exc:
+ self.logger.exception("Cash-payment SQLite cache job failed job_id=%s: %s", job_id, exc)
+ self._cash_cache_update(client_id, job_id, status="paused", stage="Paused safely", error_message=str(exc))
+ finally:
+ _CASH_CACHE_BUSY.clear()
+ _CASH_TALLY_EXTRACTION_LOCK.release()
+ with _CASH_CACHE_THREADS_LOCK:
+ _CASH_CACHE_THREADS.pop(job_id, None)
+
+ 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()
+ 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.")
+ total = max(1, int(row.get("total_days") or 1))
+ completed = int(row.get("completed_days") or 0)
+ percent = min(100, int(round(completed * 100 / total)))
+ return {
+ "job": {
+ "job_id": row["job_id"],
+ "status": row.get("status"),
+ "stage": row.get("stage"),
+ "current_date": row.get("current_date"),
+ "date_from": row.get("date_from"),
+ "date_to": row.get("date_to"),
+ "total_days": total,
+ "completed_days": completed,
+ "percent": percent,
+ "vouchers_cached": int(row.get("vouchers_cached") or 0),
+ "tally_requests": int(row.get("tally_requests") or 0),
+ "retry_count": int(row.get("retry_count") or 0),
+ "pause_seconds": float(row.get("pause_seconds") or 0),
+ "error": row.get("error_message") or "",
+ "started_at_utc": row.get("started_at_utc"),
+ "updated_at_utc": row.get("updated_at_utc"),
+ "completed_at_utc": row.get("completed_at_utc"),
+ "background": True,
+ "sqlite_cache": True,
+ "can_resume": str(row.get("status") or "") in {"paused", "failed", "interrupted"},
+ },
+ "agent": self._agent_info(),
+ }
+
+ def _cash_payment_cache_analyze(self, payload: dict[str, Any]) -> dict[str, Any]:
+ from collections import defaultdict
+ from datetime import date as _date
+ 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.")
+ if str(row.get("status") or "") != "completed":
+ raise ValueError("Tally extraction is not complete yet. Wait for the SQLite cache to finish before analysis.")
+ tally_guid = str(row.get("tally_guid") or "")
+ date_from = str(row.get("date_from") or "")
+ date_to = str(row.get("date_to") or "")
+ settings = json.loads(str(row.get("settings_json") or "{}"))
+ cash_limit = round(float(payload.get("cash_limit") or settings.get("cash_limit") or 10000.0), 2)
+ split_window_days = int(payload.get("split_window_days") or settings.get("split_window_days") or 3)
+ near_limit_percent = float(payload.get("near_limit_percent") or settings.get("near_limit_percent") or 80.0)
+ cash_ledgers = set(json.loads(str(row.get("cash_ledgers_json") or "[]"))) or {"cash"}
+
+ with self.store.connect(client_id) as db:
+ voucher_rows = db.execute(
+ """SELECT * FROM tally_vouchers WHERE tally_guid=? AND voucher_date>=? AND voucher_date<=?
+ ORDER BY voucher_date, id""",
+ (tally_guid, date_from, date_to),
+ ).fetchall()
+ vouchers = []
+ for vr in voucher_rows:
+ entries = db.execute(
+ "SELECT ledger_name, amount, is_deemed_positive FROM tally_voucher_ledger_entries WHERE voucher_id=? ORDER BY line_no",
+ (int(vr["id"]),),
+ ).fetchall()
+ vouchers.append({
+ "date": vr["voucher_date"], "effective_date": vr["effective_date"],
+ "voucher_number": vr["voucher_number"], "voucher_type_name": vr["voucher_type_name"],
+ "reference": vr["reference"], "narration": vr["narration"], "party_ledger_name": vr["party_ledger_name"],
+ "is_cancelled": vr["is_cancelled"], "is_optional": vr["is_optional"],
+ "ledger_entries": [dict(e) for e in entries],
+ })
+
+ payments = []
+ for voucher in vouchers:
+ if str(voucher.get("is_cancelled") or "").strip().lower() in {"yes", "true", "1"}:
+ continue
+ if str(voucher.get("is_optional") or "").strip().lower() in {"yes", "true", "1"}:
+ continue
+ entries = list(voucher.get("ledger_entries") or [])
+ cash_credit = []
+ for entry in entries:
+ lname = str(entry.get("ledger_name") or "").strip()
+ if lname.casefold() not in cash_ledgers:
+ continue
+ amount = float(entry.get("amount") or 0)
+ deemed = str(entry.get("is_deemed_positive") or "").strip().lower()
+ if amount > 0 or deemed == "no":
+ cash_credit.append(abs(amount))
+ cash_amount = round(sum(cash_credit), 2)
+ if cash_amount <= 0:
+ continue
+ non_cash = [e for e in entries if str(e.get("ledger_name") or "").strip().casefold() not in cash_ledgers]
+ party = str(voucher.get("party_ledger_name") or "").strip()
+ if not party or party.casefold() in cash_ledgers:
+ candidates = sorted(non_cash, key=lambda e: abs(float(e.get("amount") or 0)), reverse=True)
+ party = str(candidates[0].get("ledger_name") or "").strip() if candidates else "Unidentified counter-ledger"
+ payments.append({
+ "date": str(voucher.get("date") or voucher.get("effective_date") or ""), "party": party or "Unidentified counter-ledger",
+ "amount": cash_amount, "voucher_number": str(voucher.get("voucher_number") or ""),
+ "voucher_type": str(voucher.get("voucher_type_name") or ""), "reference": str(voucher.get("reference") or ""),
+ "narration": str(voucher.get("narration") or ""),
+ })
+
+ by_party_date = defaultdict(list)
+ for item in payments:
+ by_party_date[(item["party"].casefold(), item["date"])].append(item)
+ exceptions, seen = [], set()
+ single_count = same_day_count = 0
+ for item in payments:
+ if item["amount"] > cash_limit + 0.009:
+ single_count += 1
+ key = (item["party"].casefold(), item["date"], "single", item["voucher_number"])
+ if key not in seen:
+ seen.add(key); exceptions.append({"date": item["date"], "party": item["party"], "amount": item["amount"], "reason": "Single cash-payment voucher exceeds configured limit", "voucher_numbers": [item["voucher_number"] or "-"]})
+ for (party_key, paid_on), rows in by_party_date.items():
+ total = round(sum(x["amount"] for x in rows), 2)
+ if total > cash_limit + 0.009 and len(rows) > 1:
+ same_day_count += 1
+ key = (party_key, paid_on, "aggregate")
+ if key not in seen:
+ seen.add(key); exceptions.append({"date": paid_on, "party": rows[0]["party"], "amount": total, "reason": "Same-day aggregate cash payments to the same party exceed configured limit", "voucher_numbers": [x["voucher_number"] or "-" for x in rows]})
+ daily_by_party = defaultdict(lambda: defaultdict(float))
+ for item in payments:
+ daily_by_party[item["party"].casefold()][item["date"]] += item["amount"]
+ near_floor = cash_limit * near_limit_percent / 100.0
+ split_patterns = []
+ for party_key, day_map in daily_by_party.items():
+ day_rows = sorted((_date.fromisoformat(day), round(amount, 2)) for day, amount in day_map.items() if day)
+ for start_idx in range(len(day_rows)):
+ window = []
+ for idx in range(start_idx, len(day_rows)):
+ d, amount = day_rows[idx]
+ if (d - day_rows[start_idx][0]).days >= split_window_days:
+ break
+ window.append((d, amount))
+ if len(window) < 2 or any(amount > cash_limit + 0.009 for _, amount in window):
+ continue
+ total = round(sum(amount for _, amount in window), 2)
+ if total <= cash_limit + 0.009 or sum(1 for _, amount in window if amount >= near_floor) < 2:
+ continue
+ party_name = next((x["party"] for x in payments if x["party"].casefold() == party_key), party_key)
+ signature = (party_key, window[0][0].isoformat(), window[-1][0].isoformat())
+ if any((r["party"].casefold(), r["date_from"], r["date_to"]) == signature for r in split_patterns):
+ continue
+ split_patterns.append({"party": party_name, "date_from": window[0][0].isoformat(), "date_to": window[-1][0].isoformat(), "total_amount": total, "days": [{"date": d.isoformat(), "amount": a} for d, a in window], "review_only": True})
+ exceptions.sort(key=lambda x: (x["date"], x["party"].casefold()))
+ split_patterns.sort(key=lambda x: (x["date_from"], x["party"].casefold()))
+ return {
+ "cash_payment_review": {
+ "company_name": row.get("company_name"), "company_guid": tally_guid,
+ "date_from": date_from, "date_to": date_to, "cash_limit": cash_limit,
+ "cash_ledgers": sorted(cash_ledgers),
+ "summary": {"vouchers_reviewed": len(vouchers), "cash_payment_vouchers": len(payments), "single_voucher_exceptions": single_count, "same_day_exceptions": same_day_count, "possible_split_patterns": len(split_patterns)},
+ "exceptions": exceptions, "possible_split_payments": split_patterns, "cash_payments": payments,
+ "read_only": True, "sqlite_cache": True, "background_extraction": True,
+ "cache_job_id": job_id, "cache_completed_at_utc": row.get("completed_at_utc"),
+ "tally_request_count": int(row.get("tally_requests") or 0), "tally_retry_count": int(row.get("retry_count") or 0),
+ "tally_pause_ms": int(float(row.get("pause_seconds") or 0) * 1000),
+ "review_note": "Analysis ran entirely from the Local Agent SQLite cache after Tally extraction completed. Possible split-payment patterns are review indicators only.",
+ },
+ "agent": self._agent_info(),
+ }
+
def _cash_payment_compliance(self, payload: dict[str, Any]) -> dict[str, Any]:
from collections import defaultdict
from datetime import date as _date
diff --git a/app/modules/documents/local_agent_runtime/erp_local_agent/tunnel.py b/app/modules/documents/local_agent_runtime/erp_local_agent/tunnel.py
index 172a277..362f9df 100644
--- a/app/modules/documents/local_agent_runtime/erp_local_agent/tunnel.py
+++ b/app/modules/documents/local_agent_runtime/erp_local_agent/tunnel.py
@@ -8,7 +8,7 @@ from typing import Any
import websockets
-from .commands import AgentCommandProcessor
+from .commands import AgentCommandProcessor, is_tally_cache_busy
from . import __version__
from .tally import TallyLiveConnector
from .workstation_identity import load_or_create
@@ -199,7 +199,10 @@ class StorageAgentTunnel:
def _workstation_status(self, force_tally: bool = False) -> dict[str, Any]:
now = time.monotonic()
- if force_tally or now - self._last_tally_status_at >= 30:
+ # While the background SQLite cache is deliberately reading Tally one day
+ # at a time, do not inject extra heartbeat/status XML calls into TallyPrime.
+ # The tunnel remains responsive and reuses the last known Tally status.
+ if not is_tally_cache_busy() and (force_tally or now - self._last_tally_status_at >= 30):
try:
self._cached_tally_status = self.tally.status()
except Exception as exc: