Add SQLite cached slow Tally extraction for cash payment review
This commit is contained in:
@@ -1,2 +1,2 @@
|
||||
__version__ = "1.22.7"
|
||||
__version__ = "1.22.8"
|
||||
AGENT_NAME = "ERP Local Agent"
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user