Move existing accounting writebacks to SQLite .NET batches
This commit is contained in:
@@ -1,2 +1,2 @@
|
||||
__version__ = "1.25.7"
|
||||
__version__ = "1.26.0"
|
||||
AGENT_NAME = "ERP Local Agent"
|
||||
|
||||
@@ -1157,7 +1157,20 @@ class LocalAccountingStore:
|
||||
run=db.execute("SELECT * FROM it_depreciation_runs WHERE id=?",(int(run_id),)).fetchone()
|
||||
if not run: raise ValueError("Depreciation draft was not found.")
|
||||
lines=db.execute("SELECT * FROM it_depreciation_run_lines WHERE run_id=? ORDER BY line_no",(int(run_id),)).fetchall()
|
||||
result=dict(run); result["lines"]=[dict(x) for x in lines]
|
||||
result=dict(run)
|
||||
decoded_lines=[]
|
||||
for item in lines:
|
||||
row=dict(item)
|
||||
try:
|
||||
payload=json.loads(row.get("payload_json") or "{}")
|
||||
except Exception:
|
||||
payload={}
|
||||
if isinstance(payload,dict):
|
||||
row={**row,**payload}
|
||||
decoded_lines.append(row)
|
||||
result["lines"]=decoded_lines
|
||||
result["total_recorded_depreciation"]=round(sum(float(x.get("recorded_depreciation") or 0) for x in decoded_lines),2)
|
||||
result["total_remaining_unrecorded_depreciation"]=round(sum(float(x.get("remaining_unrecorded_depreciation") if x.get("remaining_unrecorded_depreciation") is not None else x.get("depreciation_amount") or 0) for x in decoded_lines),2)
|
||||
result["writeback_attempts"] = self.get_writeback_attempts(client_id, run_id) if "tally_writeback_attempts" else []
|
||||
result["no_tally_writeback"] = False
|
||||
return result
|
||||
|
||||
@@ -13,6 +13,7 @@ from . import __version__
|
||||
from .accounting_store import LocalAccountingStore
|
||||
from .mirror_tally import MirrorFirstTallyConnector
|
||||
from .native_voucher_engine import NativeVoucherEngine
|
||||
from .tally_batch_queue import TallyBatchQueue
|
||||
|
||||
|
||||
_CASH_TALLY_EXTRACTION_LOCK = threading.Lock()
|
||||
@@ -88,6 +89,7 @@ class AgentCommandProcessor:
|
||||
self.logger = logger
|
||||
self.store = LocalAccountingStore(config.storage_root)
|
||||
self.tally = MirrorFirstTallyConnector(self.store, logger)
|
||||
self.write_queue = TallyBatchQueue(self.store, self.tally, logger)
|
||||
|
||||
def process(self, command: dict[str, Any]) -> dict[str, Any]:
|
||||
command_id = str(command.get("command_id") or "").strip()
|
||||
@@ -148,6 +150,8 @@ class AgentCommandProcessor:
|
||||
result = self._cash_payment_mirror_analyze(payload)
|
||||
elif action == "accounting_cash_payment_compliance":
|
||||
result = self._cash_payment_compliance(payload)
|
||||
elif action == "accounting_post_cash_allocation":
|
||||
result = self._post_cash_allocation(payload)
|
||||
elif action == "accounting_tds_compliance":
|
||||
result = self._tds_compliance(payload)
|
||||
elif action == "accounting_post_tds_liability":
|
||||
@@ -235,6 +239,7 @@ class AgentCommandProcessor:
|
||||
"cash_payment_durable_vps_transfer_capability": True,
|
||||
"opening_balance_balance_sheet_only_capability": True,
|
||||
"tally_writeback_capability": True,
|
||||
"dotnet_sqlite_batch_writeback_capability": True,
|
||||
}
|
||||
|
||||
def _status(self, payload: dict[str, Any]) -> dict[str, Any]:
|
||||
@@ -789,48 +794,104 @@ class AgentCommandProcessor:
|
||||
return abs(float(left or 0)-float(right or 0))<=tolerance
|
||||
|
||||
def _opening_balance_apply(self, payload: dict[str, Any]) -> dict[str, Any]:
|
||||
current=self._opening_balance_company(str(payload.get("current_company_name") or ""))
|
||||
expected_guid=str(payload.get("current_company_guid") or "").strip()
|
||||
if expected_guid and current.guid and expected_guid!=current.guid: raise ValueError("Current Tally company GUID changed since comparison. No opening balance was altered.")
|
||||
masters=self.tally.fetch_accounting_masters(current.name)
|
||||
ledgers={str(r.get("name") or "").strip().casefold():r for r in masters.get("ledgers") or []}
|
||||
stocks={str(r.get("name") or "").strip().casefold():r for r in masters.get("stock_items") or []}
|
||||
lr=[]; sr=[]; le=[]; se=[]
|
||||
current = self._opening_balance_company(str(payload.get("current_company_name") or ""))
|
||||
expected_guid = str(payload.get("current_company_guid") or "").strip()
|
||||
if expected_guid and current.guid and expected_guid != current.guid:
|
||||
raise ValueError("Current Tally company GUID changed since comparison. No opening balance was altered.")
|
||||
|
||||
masters = self.tally.fetch_accounting_masters(current.name)
|
||||
ledgers = {str(r.get("name") or "").strip().casefold(): r for r in masters.get("ledgers") or []}
|
||||
stocks = {str(r.get("name") or "").strip().casefold(): r for r in masters.get("stock_items") or []}
|
||||
accepted_ledgers = []
|
||||
accepted_stocks = []
|
||||
precheck_results_l = []
|
||||
precheck_results_s = []
|
||||
|
||||
for req in list(payload.get("ledgers") or []):
|
||||
name=str(req.get("name") or "").strip(); row=ledgers.get(name.casefold())
|
||||
if not row: lr.append({"item_id":req.get("item_id"),"verified":False,"message":f"Current Ledger '{name}' no longer exists."}); continue
|
||||
expected=float(req.get("expected_opening") or 0); actual=float(row.get("opening_balance") or 0)
|
||||
if not self._opening_near(expected,actual): lr.append({"item_id":req.get("item_id"),"verified":False,"verified_opening":actual,"message":"Ledger opening changed after comparison; not altered."}); continue
|
||||
le.append(req)
|
||||
name = str(req.get("name") or "").strip()
|
||||
row = ledgers.get(name.casefold())
|
||||
if not row:
|
||||
precheck_results_l.append({"item_id": req.get("item_id"), "verified": False, "message": f"Current Ledger '{name}' no longer exists."})
|
||||
continue
|
||||
expected = float(req.get("expected_opening") or 0)
|
||||
actual = float(row.get("opening_balance") or 0)
|
||||
if not self._opening_near(expected, actual):
|
||||
precheck_results_l.append({"item_id": req.get("item_id"), "verified": False, "verified_opening": actual, "message": "Ledger opening changed after comparison; not altered."})
|
||||
continue
|
||||
accepted_ledgers.append(req)
|
||||
|
||||
for req in list(payload.get("stock_items") or []):
|
||||
name=str(req.get("name") or "").strip(); row=stocks.get(name.casefold())
|
||||
if not row: sr.append({"item_id":req.get("item_id"),"verified":False,"message":f"Current Stock Item '{name}' no longer exists."}); continue
|
||||
eq=float(req.get("expected_opening_qty") or 0); ev=float(req.get("expected_opening_value") or 0); aq=float(row.get("opening_balance") or 0); av=float(row.get("opening_value") or 0)
|
||||
if not self._opening_near(eq,aq,0.000001) or not self._opening_near(ev,av): sr.append({"item_id":req.get("item_id"),"verified":False,"verified_opening_qty":aq,"verified_opening_value":av,"message":"Stock opening changed after comparison; not altered."}); continue
|
||||
cu=str(row.get("base_units") or "").strip(); ru=str(req.get("unit") or "").strip()
|
||||
if cu and ru and cu.casefold()!=ru.casefold(): sr.append({"item_id":req.get("item_id"),"verified":False,"message":f"Current Tally unit is '{cu}', not '{ru}'."}); continue
|
||||
se.append(req)
|
||||
for req in le:
|
||||
try:
|
||||
self.tally.alter_ledger_opening_balance(current.name,ledger_name=str(req.get("name") or "").strip(),opening_balance=float(req.get("target_opening") or 0))
|
||||
lr.append({"item_id":req.get("item_id"),"verified":None,"target_opening":float(req.get("target_opening") or 0),"message":"Altered; verification pending."})
|
||||
except Exception as exc: lr.append({"item_id":req.get("item_id"),"verified":False,"message":str(exc)})
|
||||
for req in se:
|
||||
try:
|
||||
self.tally.alter_stock_item_opening(current.name,stock_item_name=str(req.get("name") or "").strip(),opening_qty=float(req.get("target_opening_qty") or 0),unit=str(req.get("unit") or "").strip(),opening_value=float(req.get("target_opening_value") or 0),opening_rate=float(req.get("target_opening_rate") or 0))
|
||||
sr.append({"item_id":req.get("item_id"),"verified":None,"target_opening_qty":float(req.get("target_opening_qty") or 0),"target_opening_value":float(req.get("target_opening_value") or 0),"message":"Altered; verification pending."})
|
||||
except Exception as exc: sr.append({"item_id":req.get("item_id"),"verified":False,"message":str(exc)})
|
||||
name = str(req.get("name") or "").strip()
|
||||
row = stocks.get(name.casefold())
|
||||
if not row:
|
||||
precheck_results_s.append({"item_id": req.get("item_id"), "verified": False, "message": f"Current Stock Item '{name}' no longer exists."})
|
||||
continue
|
||||
eq = float(req.get("expected_opening_qty") or 0)
|
||||
ev = float(req.get("expected_opening_value") or 0)
|
||||
aq = float(row.get("opening_balance") or 0)
|
||||
av = float(row.get("opening_value") or 0)
|
||||
if not self._opening_near(eq, aq, 0.000001) or not self._opening_near(ev, av):
|
||||
precheck_results_s.append({"item_id": req.get("item_id"), "verified": False, "verified_opening_qty": aq, "verified_opening_value": av, "message": "Stock opening changed after comparison; not altered."})
|
||||
continue
|
||||
cu = str(row.get("base_units") or "").strip()
|
||||
ru = str(req.get("unit") or "").strip()
|
||||
if cu and ru and cu.casefold() != ru.casefold():
|
||||
precheck_results_s.append({"item_id": req.get("item_id"), "verified": False, "message": f"Current Tally unit is '{cu}', not '{ru}'."})
|
||||
continue
|
||||
accepted_stocks.append(req)
|
||||
|
||||
if not accepted_ledgers and not accepted_stocks:
|
||||
return {
|
||||
"company_name": current.name, "company_guid": current.guid,
|
||||
"ledgers": precheck_results_l, "stock_items": precheck_results_s,
|
||||
"write_operation": False, "verified_after_write": False, "agent": self._agent_info(),
|
||||
}
|
||||
|
||||
client_id = int(payload.get("client_id") or 0)
|
||||
source_key = str(payload.get("source_key") or "").strip() or (
|
||||
"OPENING:" + str(payload.get("current_company_guid") or current.guid or current.name) + ":"
|
||||
+ ",".join(str(x.get("item_id") or 0) for x in accepted_ledgers + accepted_stocks)
|
||||
)
|
||||
batch = self.write_queue.queue_opening_balances(
|
||||
client_id,
|
||||
source_key=source_key,
|
||||
company_guid=str(current.guid or expected_guid or ""),
|
||||
company_name=current.name,
|
||||
created_by_user_id=int(payload.get("requested_by_user_id") or 0) or None,
|
||||
ledgers=accepted_ledgers,
|
||||
stock_items=accepted_stocks,
|
||||
)
|
||||
executed = self.write_queue.execute(client_id, int(batch["batch"]["id"]))
|
||||
by_key = {str(x.get("idempotency_key") or ""): x for x in (executed.get("entries") or [])}
|
||||
lr = list(precheck_results_l)
|
||||
sr = list(precheck_results_s)
|
||||
for req in accepted_ledgers:
|
||||
key = f"{source_key}:LEDGER:{int(req.get('item_id') or 0)}"
|
||||
row = by_key.get(key) or {}
|
||||
ok = str(row.get("status") or "") == "verified"
|
||||
lr.append({
|
||||
"item_id": req.get("item_id"), "verified": ok,
|
||||
"target_opening": float(req.get("target_opening") or 0),
|
||||
"message": "Ledger opening altered by .NET batch and re-read successfully." if ok else str(row.get("last_error") or "Ledger opening verification failed."),
|
||||
})
|
||||
for req in accepted_stocks:
|
||||
key = f"{source_key}:STOCK:{int(req.get('item_id') or 0)}"
|
||||
row = by_key.get(key) or {}
|
||||
ok = str(row.get("status") or "") == "verified"
|
||||
sr.append({
|
||||
"item_id": req.get("item_id"), "verified": ok,
|
||||
"target_opening_qty": float(req.get("target_opening_qty") or 0),
|
||||
"target_opening_value": float(req.get("target_opening_value") or 0),
|
||||
"message": "Stock opening altered by .NET batch and re-read successfully." if ok else str(row.get("last_error") or "Stock opening verification failed."),
|
||||
})
|
||||
self.tally.refresh_mirror(current.name, str(current.guid or ""))
|
||||
verify=self.tally.fetch_accounting_masters(current.name)
|
||||
vl={str(r.get("name") or "").strip().casefold():r for r in verify.get("ledgers") or []}; vs={str(r.get("name") or "").strip().casefold():r for r in verify.get("stock_items") or []}
|
||||
lem={int(r.get("item_id") or 0):r for r in le}; sem={int(r.get("item_id") or 0):r for r in se}
|
||||
for res in lr:
|
||||
if res.get("verified") is not None: continue
|
||||
req=lem.get(int(res.get("item_id") or 0)); row=vl.get(str(req.get("name") or "").strip().casefold()) if req else None; actual=float((row or {}).get("opening_balance") or 0); target=float((req or {}).get("target_opening") or 0); ok=self._opening_near(actual,target); res.update({"verified":ok,"verified_opening":actual,"message":"Ledger opening re-read and verified." if ok else f"Verification returned {actual:.2f}; expected {target:.2f}."})
|
||||
for res in sr:
|
||||
if res.get("verified") is not None: continue
|
||||
req=sem.get(int(res.get("item_id") or 0)); row=vs.get(str(req.get("name") or "").strip().casefold()) if req else None; aq=float((row or {}).get("opening_balance") or 0); av=float((row or {}).get("opening_value") or 0); tq=float((req or {}).get("target_opening_qty") or 0); tv=float((req or {}).get("target_opening_value") or 0); ok=self._opening_near(aq,tq,0.000001) and self._opening_near(av,tv); res.update({"verified":ok,"verified_opening_qty":aq,"verified_opening_value":av,"message":"Stock opening re-read and verified." if ok else "Stock opening verification did not equal the approved quantity/value."})
|
||||
return {"company_name":current.name,"company_guid":current.guid,"ledgers":lr,"stock_items":sr,"write_operation":True,"verified_after_write":True,"agent":self._agent_info()}
|
||||
return {
|
||||
"company_name": current.name, "company_guid": current.guid,
|
||||
"ledgers": lr, "stock_items": sr, "write_operation": True,
|
||||
"verified_after_write": True, "batch_id": int(executed["batch"]["id"]),
|
||||
"batch": executed, "agent": self._agent_info(),
|
||||
}
|
||||
|
||||
|
||||
def _stock_master_intelligence(self, payload: dict[str, Any]) -> dict[str, Any]:
|
||||
# Read-only against Tally. Existing local .act master storage remains
|
||||
@@ -3327,13 +3388,54 @@ class AgentCommandProcessor:
|
||||
|
||||
def _post_tds_liability(self, payload: dict[str, Any]) -> dict[str, Any]:
|
||||
company, company_name = self._resolve_open_company(payload)
|
||||
voucher_date=str(payload.get("voucher_date") or "").strip(); party=str(payload.get("party_ledger") or "").strip(); tds=str(payload.get("tds_ledger") or "").strip(); amount=round(float(payload.get("amount") or 0),2); ref=str(payload.get("reference") or "").strip()
|
||||
if not voucher_date or not party or not tds or amount<=0 or not ref: raise ValueError("Voucher date, party ledger, TDS ledger, positive amount and ERP reference are required.")
|
||||
for v in self.tally.export_vouchers(company_name,voucher_date,voucher_date):
|
||||
if ref and ref in {str(v.get("reference") or "").strip(), str(v.get("narration") or "").strip()}:
|
||||
raise ValueError(f"TDS entry {ref} already exists in Tally; duplicate posting blocked.")
|
||||
result=self.tally.post_journal_voucher(company_name,voucher_date=voucher_date,debit_ledger=party,credit_ledger=tds,amount=amount,narration=str(payload.get("narration") or f"TDS liability {ref}"),reference=ref)
|
||||
return {"posted":True,"company_name":company_name,"reference":ref,"last_voucher_id":result.get("last_voucher_id"),"result":result,"agent":self._agent_info()}
|
||||
client_id = int(payload.get("client_id") or 0)
|
||||
voucher_date = str(payload.get("voucher_date") or "").strip()
|
||||
party = str(payload.get("party_ledger") or "").strip()
|
||||
tds = str(payload.get("tds_ledger") or "").strip()
|
||||
amount = round(float(payload.get("amount") or 0), 2)
|
||||
ref = str(payload.get("reference") or "").strip()
|
||||
narration = str(payload.get("narration") or f"TDS liability {ref}")
|
||||
posted_by = int(payload.get("posted_by_user_id") or payload.get("requested_by_user_id") or 0) or None
|
||||
if client_id <= 0 or not voucher_date or not party or not tds or amount <= 0 or not ref:
|
||||
raise ValueError("Client, voucher date, party ledger, TDS ledger, positive amount and ERP reference are required.")
|
||||
|
||||
batch = self.write_queue.queue_vouchers(
|
||||
client_id,
|
||||
source_tool="TDS_LIABILITY",
|
||||
source_key=ref,
|
||||
company_guid=str(company.get("guid") or payload.get("tally_guid") or ""),
|
||||
company_name=company_name,
|
||||
created_by_user_id=posted_by,
|
||||
vouchers=[{
|
||||
"idempotency_key": ref,
|
||||
"voucher_type": "Journal",
|
||||
"voucher_date": voucher_date,
|
||||
"reference": ref,
|
||||
"narration": narration,
|
||||
"total_amount": amount,
|
||||
"lines": [
|
||||
{"ledger_name": party, "dr_cr": "DR", "amount": amount},
|
||||
{"ledger_name": tds, "dr_cr": "CR", "amount": amount},
|
||||
],
|
||||
}],
|
||||
)
|
||||
executed = self.write_queue.execute(client_id, int(batch["batch"]["id"]))
|
||||
entries = executed.get("entries") or []
|
||||
first = entries[0] if entries else {}
|
||||
if str(first.get("status") or "") != "verified":
|
||||
raise ValueError(str(first.get("last_error") or "TDS batch posting could not be verified in Tally."))
|
||||
return {
|
||||
"posted": True,
|
||||
"verified": True,
|
||||
"company_name": company_name,
|
||||
"reference": ref,
|
||||
"batch_id": int(executed["batch"]["id"]),
|
||||
"last_voucher_id": str(first.get("tally_voucher_id") or ""),
|
||||
"voucher_number": str(first.get("tally_voucher_number") or ""),
|
||||
"batch": executed,
|
||||
"agent": self._agent_info(),
|
||||
}
|
||||
|
||||
|
||||
def _bank_posting_preflight(self, payload: dict[str, Any]) -> dict[str, Any]:
|
||||
company, company_name = self._resolve_open_company(payload)
|
||||
@@ -3596,11 +3698,7 @@ class AgentCommandProcessor:
|
||||
tally_status = self.tally.status()
|
||||
if not tally_status.get("connected"):
|
||||
raise ValueError(str(tally_status.get("error") or "TallyPrime is not connected."))
|
||||
loaded = next(
|
||||
(row for row in (tally_status.get("companies") or [])
|
||||
if str(row.get("guid") or "").strip() == tally_guid),
|
||||
None,
|
||||
)
|
||||
loaded = next((row for row in (tally_status.get("companies") or []) if str(row.get("guid") or "").strip() == tally_guid), None)
|
||||
if not loaded:
|
||||
raise ValueError("The mapped Tally company is not currently open. Open the mapped company in TallyPrime and retry.")
|
||||
if str(loaded.get("name") or "").strip() != company_name:
|
||||
@@ -3611,36 +3709,113 @@ class AgentCommandProcessor:
|
||||
narration = f"Income-tax depreciation for FY {run.get('fy_start')} to {run.get('fy_end')} · ERP draft #{run_id}"
|
||||
debit = str(run.get("depreciation_expense_ledger") or "").strip()
|
||||
credit = str(run.get("depreciation_reserve_ledger") or "").strip()
|
||||
amount = float(run.get("total_depreciation") or 0)
|
||||
request_xml = self.tally.build_journal_import_xml(
|
||||
company_name, voucher_date=voucher_date, debit_ledger=debit,
|
||||
credit_ledger=credit, amount=amount, narration=narration, reference=reference,
|
||||
)
|
||||
remaining = run.get("total_remaining_unrecorded_depreciation")
|
||||
amount = round(float(remaining if remaining is not None else run.get("total_depreciation") or 0), 2)
|
||||
if amount <= 0:
|
||||
raise ValueError("No unrecorded depreciation remains to be posted for this approved draft.")
|
||||
|
||||
attempt = self.store.begin_writeback_attempt(
|
||||
client_id, run_id, posted_by, voucher_date=voucher_date,
|
||||
reference=reference, request_xml=request_xml,
|
||||
client_id, run_id, posted_by, voucher_date=voucher_date, reference=reference,
|
||||
request_xml=f"Queued to .NET batch writer; amount={amount:.2f}",
|
||||
)
|
||||
attempt_id = int(attempt["attempt_id"])
|
||||
try:
|
||||
result = self.tally.post_journal_voucher(
|
||||
company_name, voucher_date=voucher_date, debit_ledger=debit,
|
||||
credit_ledger=credit, amount=amount, narration=narration, reference=reference,
|
||||
batch = self.write_queue.queue_vouchers(
|
||||
client_id,
|
||||
source_tool="DEPRECIATION_IT",
|
||||
source_key=f"DEPRECIATION:{run_id}",
|
||||
company_guid=tally_guid,
|
||||
company_name=company_name,
|
||||
created_by_user_id=posted_by,
|
||||
vouchers=[{
|
||||
"idempotency_key": reference,
|
||||
"voucher_type": "Journal",
|
||||
"voucher_date": voucher_date,
|
||||
"reference": reference,
|
||||
"narration": narration,
|
||||
"total_amount": amount,
|
||||
"lines": [
|
||||
{"ledger_name": debit, "dr_cr": "DR", "amount": amount},
|
||||
{"ledger_name": credit, "dr_cr": "CR", "amount": amount},
|
||||
],
|
||||
}],
|
||||
)
|
||||
executed = self.write_queue.execute(client_id, int(batch["batch"]["id"]))
|
||||
entries = executed.get("entries") or []
|
||||
first = entries[0] if entries else {}
|
||||
if str(first.get("status") or "") != "verified":
|
||||
raise ValueError(str(first.get("last_error") or "Depreciation batch posting could not be verified in Tally."))
|
||||
result = {
|
||||
"created": 1, "altered": 0, "errors": 0,
|
||||
"last_voucher_id": str(first.get("tally_voucher_id") or ""),
|
||||
"voucher_number": str(first.get("tally_voucher_number") or ""),
|
||||
"batch_id": int(executed["batch"]["id"]),
|
||||
}
|
||||
updated = self.store.finish_writeback_attempt(
|
||||
client_id, attempt_id, posted=True, result=result,
|
||||
response_xml=str(result.get("raw_response") or ""), posted_by_user_id=posted_by,
|
||||
response_xml=json.dumps(executed, ensure_ascii=False), posted_by_user_id=posted_by,
|
||||
)
|
||||
self.logger.warning(
|
||||
"CONTROLLED TALLY WRITEBACK posted client_id=%s company=%s run_id=%s amount=%.2f by_user=%s tally_voucher=%s",
|
||||
client_id, company_name, run_id, amount, posted_by, result.get("last_voucher_id"),
|
||||
"CONTROLLED .NET BATCH TALLY WRITEBACK verified client_id=%s company=%s run_id=%s amount=%.2f by_user=%s batch_id=%s",
|
||||
client_id, company_name, run_id, amount, posted_by, executed["batch"]["id"],
|
||||
)
|
||||
return {"posted": True, "depreciation": updated, "tally_result": {k:v for k,v in result.items() if k not in {"raw_response","request_xml"}}, "agent": self._agent_info()}
|
||||
return {"posted": True, "verified": True, "depreciation": updated, "tally_result": result, "batch": executed, "agent": self._agent_info()}
|
||||
except Exception as exc:
|
||||
self.store.finish_writeback_attempt(
|
||||
client_id, attempt_id, posted=False, result={}, error_message=str(exc), posted_by_user_id=posted_by,
|
||||
)
|
||||
raise
|
||||
|
||||
def _post_cash_allocation(self, payload: dict[str, Any]) -> dict[str, Any]:
|
||||
client_id = int(payload.get("client_id") or 0)
|
||||
if client_id > 0 and not str(payload.get("tally_guid") or "").strip():
|
||||
mappings = self.store.list_active_mappings(client_id)
|
||||
if len(mappings) == 1:
|
||||
payload = {
|
||||
**payload,
|
||||
"tally_guid": str(mappings[0].get("tally_guid") or ""),
|
||||
"company_name": str(mappings[0].get("company_name") or ""),
|
||||
}
|
||||
company, company_name = self._resolve_open_company(payload)
|
||||
cash_ledger = str(payload.get("cash_ledger") or "Cash").strip()
|
||||
debit_ledger = str(payload.get("debit_ledger") or "").strip()
|
||||
party_name = str(payload.get("party_name") or "").strip()
|
||||
entries = list(payload.get("entries") or [])
|
||||
posted_by = int(payload.get("posted_by_user_id") or payload.get("requested_by_user_id") or 0) or None
|
||||
if client_id <= 0 or not debit_ledger or not cash_ledger or not entries:
|
||||
raise ValueError("Client, debit ledger, cash ledger and approved split entries are required.")
|
||||
vouchers = []
|
||||
source_key = str(payload.get("source_key") or "").strip() or f"CASHALLOC:{client_id}:{uuid.uuid4().hex[:12]}"
|
||||
for index, row in enumerate(entries, 1):
|
||||
voucher_date = str(row.get("date") or "").strip()
|
||||
amount = round(float(row.get("amount") or 0), 2)
|
||||
if not voucher_date or amount <= 0:
|
||||
raise ValueError("Every cash allocation row requires an actual date and positive amount.")
|
||||
reference = f"{source_key}-{index:03d}"
|
||||
narration = str(payload.get("narration") or "").strip() or f"Cash payment allocation{(' - ' + party_name) if party_name else ''} · {source_key}"
|
||||
vouchers.append({
|
||||
"idempotency_key": reference,
|
||||
"voucher_type": "Payment",
|
||||
"voucher_date": voucher_date,
|
||||
"reference": reference,
|
||||
"narration": narration,
|
||||
"total_amount": amount,
|
||||
"lines": [
|
||||
{"ledger_name": debit_ledger, "dr_cr": "DR", "amount": amount},
|
||||
{"ledger_name": cash_ledger, "dr_cr": "CR", "amount": amount},
|
||||
],
|
||||
})
|
||||
batch = self.write_queue.queue_vouchers(
|
||||
client_id, source_tool="CASH_PAYMENT_ALLOCATION", source_key=source_key,
|
||||
company_guid=str(company.get("guid") or payload.get("tally_guid") or ""), company_name=company_name,
|
||||
created_by_user_id=posted_by, vouchers=vouchers,
|
||||
)
|
||||
executed = self.write_queue.execute(client_id, int(batch["batch"]["id"]))
|
||||
failed = [x for x in (executed.get("entries") or []) if str(x.get("status") or "") != "verified"]
|
||||
if failed:
|
||||
raise ValueError(f"{len(failed)} cash allocation voucher(s) could not be verified in Tally. Batch #{executed['batch']['id']} is retained in the client .act for safe retry/audit.")
|
||||
return {"posted": True, "verified": True, "batch_id": int(executed["batch"]["id"]), "batch": executed, "agent": self._agent_info()}
|
||||
|
||||
|
||||
|
||||
@staticmethod
|
||||
|
||||
+94
@@ -0,0 +1,94 @@
|
||||
using System;
|
||||
using System.IO;
|
||||
using System.Net;
|
||||
using System.Text;
|
||||
|
||||
class TallyBatchWriterV1260
|
||||
{
|
||||
static int Main(string[] args)
|
||||
{
|
||||
string url = "http://127.0.0.1:9000";
|
||||
string requestPath = "";
|
||||
string responsePath = "";
|
||||
int timeoutSeconds = 180;
|
||||
|
||||
for (int i = 0; i < args.Length; i++)
|
||||
{
|
||||
if (args[i] == "--url" && i + 1 < args.Length) url = args[++i];
|
||||
else if (args[i] == "--request" && i + 1 < args.Length) requestPath = args[++i];
|
||||
else if (args[i] == "--response" && i + 1 < args.Length) responsePath = args[++i];
|
||||
else if (args[i] == "--timeout" && i + 1 < args.Length) timeoutSeconds = Math.Max(10, Int32.Parse(args[++i]));
|
||||
}
|
||||
|
||||
if (String.IsNullOrWhiteSpace(requestPath) || !File.Exists(requestPath))
|
||||
{
|
||||
Console.Error.WriteLine("Request XML file was not found.");
|
||||
return 2;
|
||||
}
|
||||
if (String.IsNullOrWhiteSpace(responsePath))
|
||||
{
|
||||
Console.Error.WriteLine("Response XML path is required.");
|
||||
return 2;
|
||||
}
|
||||
|
||||
try
|
||||
{
|
||||
byte[] payload = File.ReadAllBytes(requestPath);
|
||||
HttpWebRequest request = (HttpWebRequest)WebRequest.Create(url);
|
||||
request.Method = "POST";
|
||||
request.ContentType = "application/xml; charset=utf-8";
|
||||
request.ContentLength = payload.Length;
|
||||
request.Timeout = timeoutSeconds * 1000;
|
||||
request.ReadWriteTimeout = timeoutSeconds * 1000;
|
||||
request.KeepAlive = true;
|
||||
|
||||
using (Stream stream = request.GetRequestStream())
|
||||
{
|
||||
stream.Write(payload, 0, payload.Length);
|
||||
}
|
||||
|
||||
string responseText;
|
||||
using (HttpWebResponse response = (HttpWebResponse)request.GetResponse())
|
||||
using (Stream responseStream = response.GetResponseStream())
|
||||
using (StreamReader reader = new StreamReader(responseStream, Encoding.UTF8, true))
|
||||
{
|
||||
responseText = reader.ReadToEnd();
|
||||
}
|
||||
|
||||
Directory.CreateDirectory(Path.GetDirectoryName(Path.GetFullPath(responsePath)));
|
||||
File.WriteAllText(responsePath, responseText, new UTF8Encoding(false));
|
||||
Console.WriteLine("SUCCESS");
|
||||
Console.WriteLine("REQUEST_BYTES=" + payload.Length);
|
||||
Console.WriteLine("RESPONSE_BYTES=" + Encoding.UTF8.GetByteCount(responseText));
|
||||
return 0;
|
||||
}
|
||||
catch (WebException ex)
|
||||
{
|
||||
string body = "";
|
||||
try
|
||||
{
|
||||
if (ex.Response != null)
|
||||
{
|
||||
using (Stream s = ex.Response.GetResponseStream())
|
||||
using (StreamReader r = new StreamReader(s, Encoding.UTF8, true))
|
||||
body = r.ReadToEnd();
|
||||
}
|
||||
}
|
||||
catch { }
|
||||
try
|
||||
{
|
||||
if (!String.IsNullOrWhiteSpace(responsePath))
|
||||
File.WriteAllText(responsePath, body, new UTF8Encoding(false));
|
||||
}
|
||||
catch { }
|
||||
Console.Error.WriteLine("HTTP ERROR: " + ex.Message);
|
||||
if (!String.IsNullOrWhiteSpace(body)) Console.Error.WriteLine(body);
|
||||
return 3;
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
Console.Error.WriteLine("ERROR: " + ex.GetType().FullName + ": " + ex.Message);
|
||||
return 4;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,633 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime, timezone
|
||||
import json
|
||||
import re
|
||||
import sqlite3
|
||||
import subprocess
|
||||
import xml.etree.ElementTree as ET
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
from .native_voucher_engine import NativeVoucherEngine
|
||||
|
||||
|
||||
def _utc_now_iso() -> str:
|
||||
return datetime.now(timezone.utc).isoformat()
|
||||
|
||||
|
||||
def _clean_xml_response(text: str) -> str:
|
||||
text = str(text or "")
|
||||
text = re.sub(r"&#(0?[0-8]|1[0-9]|2[0-9]|3[01]);", "", text)
|
||||
return re.sub(r"[\x00-\x08\x0B\x0C\x0E-\x1F]", "", text)
|
||||
|
||||
|
||||
def _first_text(root: ET.Element, names: set[str]) -> str:
|
||||
names = {x.upper() for x in names}
|
||||
for node in root.iter():
|
||||
tag = str(node.tag).split("}")[-1].upper()
|
||||
if tag in names and node.text:
|
||||
return node.text.strip()
|
||||
return ""
|
||||
|
||||
|
||||
class TallyBatchQueue:
|
||||
"""Internal client-.act posting queue used by existing ERP accounting tools.
|
||||
|
||||
This class is deliberately not exposed as a separate ERP tool. Existing
|
||||
Depreciation, TDS Liability, Opening Balance and Cash Payment Allocation
|
||||
screens create approved instructions here and immediately execute them through
|
||||
the shared .NET batch writer.
|
||||
"""
|
||||
|
||||
SCHEMA_VERSION = "1"
|
||||
DEFAULT_BATCH_SIZE = 25
|
||||
|
||||
def __init__(self, store, tally, logger):
|
||||
self.store = store
|
||||
self.tally = tally
|
||||
self.logger = logger
|
||||
self.runtime_root = Path(__file__).resolve().parent / "mirror_runtime"
|
||||
self.agent_runtime = Path(__file__).resolve().parents[1] / "data" / "accounting_write_runtime"
|
||||
self.agent_runtime.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
def _ensure_schema(self, client_id: int) -> None:
|
||||
self.store.initialize(client_id)
|
||||
with self.store.connect(client_id) as db:
|
||||
db.executescript(
|
||||
"""
|
||||
CREATE TABLE IF NOT EXISTS tally_write_batches_v2 (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
source_tool TEXT NOT NULL,
|
||||
source_key TEXT NOT NULL,
|
||||
operation_kind TEXT NOT NULL,
|
||||
company_guid TEXT NOT NULL DEFAULT '',
|
||||
company_name TEXT NOT NULL,
|
||||
status TEXT NOT NULL DEFAULT 'queued',
|
||||
created_by_user_id INTEGER,
|
||||
created_at_utc TEXT NOT NULL,
|
||||
started_at_utc TEXT,
|
||||
completed_at_utc TEXT,
|
||||
entry_count INTEGER NOT NULL DEFAULT 0,
|
||||
verified_count INTEGER NOT NULL DEFAULT 0,
|
||||
failed_count INTEGER NOT NULL DEFAULT 0,
|
||||
skipped_count INTEGER NOT NULL DEFAULT 0,
|
||||
last_error TEXT NOT NULL DEFAULT '',
|
||||
UNIQUE(source_tool, source_key)
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS ix_tally_write_batches_v2_status
|
||||
ON tally_write_batches_v2(status, id);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS tally_write_entries_v2 (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
batch_id INTEGER NOT NULL,
|
||||
line_no INTEGER NOT NULL,
|
||||
idempotency_key TEXT NOT NULL,
|
||||
operation_type TEXT NOT NULL,
|
||||
voucher_type TEXT NOT NULL DEFAULT '',
|
||||
voucher_date TEXT NOT NULL DEFAULT '',
|
||||
reference TEXT NOT NULL DEFAULT '',
|
||||
narration TEXT NOT NULL DEFAULT '',
|
||||
payload_json TEXT NOT NULL DEFAULT '{}',
|
||||
status TEXT NOT NULL DEFAULT 'queued',
|
||||
tally_voucher_id TEXT NOT NULL DEFAULT '',
|
||||
tally_voucher_number TEXT NOT NULL DEFAULT '',
|
||||
verified_at_utc TEXT,
|
||||
last_error TEXT NOT NULL DEFAULT '',
|
||||
created_at_utc TEXT NOT NULL,
|
||||
updated_at_utc TEXT NOT NULL,
|
||||
FOREIGN KEY(batch_id) REFERENCES tally_write_batches_v2(id) ON DELETE CASCADE,
|
||||
UNIQUE(idempotency_key)
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS ix_tally_write_entries_v2_batch
|
||||
ON tally_write_entries_v2(batch_id, line_no);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS tally_write_attempts_v2 (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
batch_id INTEGER NOT NULL,
|
||||
attempt_no INTEGER NOT NULL,
|
||||
entry_ids_json TEXT NOT NULL DEFAULT '[]',
|
||||
request_xml TEXT NOT NULL DEFAULT '',
|
||||
response_xml TEXT NOT NULL DEFAULT '',
|
||||
status TEXT NOT NULL,
|
||||
started_at_utc TEXT NOT NULL,
|
||||
completed_at_utc TEXT,
|
||||
error_message TEXT NOT NULL DEFAULT '',
|
||||
FOREIGN KEY(batch_id) REFERENCES tally_write_batches_v2(id) ON DELETE CASCADE
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS ix_tally_write_attempts_v2_batch
|
||||
ON tally_write_attempts_v2(batch_id, attempt_no);
|
||||
"""
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def _find_csc() -> Path:
|
||||
import os
|
||||
windir = Path(os.environ.get("WINDIR") or r"C:\Windows")
|
||||
for candidate in (
|
||||
windir / "Microsoft.NET/Framework64/v4.0.30319/csc.exe",
|
||||
windir / "Microsoft.NET/Framework/v4.0.30319/csc.exe",
|
||||
):
|
||||
if candidate.is_file():
|
||||
return candidate
|
||||
raise RuntimeError("Microsoft .NET Framework 4.x C# compiler was not found.")
|
||||
|
||||
def _ensure_writer(self) -> Path:
|
||||
source = self.runtime_root / "TallyBatchWriterV1260.cs"
|
||||
exe = self.agent_runtime / "TallyBatchWriterV1260.exe"
|
||||
rebuild = not exe.is_file()
|
||||
if not rebuild:
|
||||
try:
|
||||
rebuild = source.stat().st_mtime > exe.stat().st_mtime
|
||||
except OSError:
|
||||
rebuild = True
|
||||
if rebuild:
|
||||
csc = self._find_csc()
|
||||
cp = subprocess.run(
|
||||
[str(csc), "/nologo", "/target:exe", f"/out:{exe}", str(source)],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
timeout=120,
|
||||
)
|
||||
if cp.returncode != 0 or not exe.is_file():
|
||||
raise RuntimeError(
|
||||
"Tally .NET batch writer build failed: "
|
||||
+ ((cp.stdout or "") + "\n" + (cp.stderr or "")).strip()
|
||||
)
|
||||
return exe
|
||||
|
||||
def _call_writer(self, client_id: int, batch_id: int, attempt_no: int, request_xml: str) -> str:
|
||||
exe = self._ensure_writer()
|
||||
work = self.store.db_path(client_id).parent / ".tally_write_queue"
|
||||
work.mkdir(parents=True, exist_ok=True)
|
||||
request_path = work / f"batch_{batch_id:08d}_attempt_{attempt_no:03d}_request.xml"
|
||||
response_path = work / f"batch_{batch_id:08d}_attempt_{attempt_no:03d}_response.xml"
|
||||
request_path.write_text(request_xml, encoding="utf-8")
|
||||
cp = subprocess.run(
|
||||
[
|
||||
str(exe),
|
||||
"--url", self.tally.url,
|
||||
"--request", str(request_path),
|
||||
"--response", str(response_path),
|
||||
"--timeout", "180",
|
||||
],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
timeout=210,
|
||||
)
|
||||
response_xml = response_path.read_text(encoding="utf-8", errors="replace") if response_path.is_file() else ""
|
||||
if cp.returncode != 0:
|
||||
raise RuntimeError(
|
||||
"Tally .NET batch writer failed: "
|
||||
+ ((cp.stdout or "") + "\n" + (cp.stderr or "") + "\n" + response_xml).strip()
|
||||
)
|
||||
return response_xml
|
||||
|
||||
@staticmethod
|
||||
def _xml_import_result(response_xml: str) -> dict[str, Any]:
|
||||
cleaned = _clean_xml_response(response_xml)
|
||||
try:
|
||||
root = ET.fromstring(cleaned.encode("utf-8"))
|
||||
except Exception as exc:
|
||||
raise ValueError(f"Tally returned invalid XML after .NET batch import: {exc}") from exc
|
||||
|
||||
def to_int(name: str) -> int:
|
||||
try:
|
||||
return int(float(_first_text(root, {name}) or 0))
|
||||
except Exception:
|
||||
return 0
|
||||
|
||||
result = {
|
||||
"created": to_int("CREATED"),
|
||||
"altered": to_int("ALTERED"),
|
||||
"errors": to_int("ERRORS"),
|
||||
"last_voucher_id": _first_text(root, {"LASTVCHID", "LASTMID", "LASTVOUCHERID"}),
|
||||
"line_error": _first_text(root, {"LINEERROR"}),
|
||||
}
|
||||
return result
|
||||
|
||||
def _get_or_create_batch(
|
||||
self,
|
||||
client_id: int,
|
||||
*,
|
||||
source_tool: str,
|
||||
source_key: str,
|
||||
operation_kind: str,
|
||||
company_guid: str,
|
||||
company_name: str,
|
||||
created_by_user_id: int | None,
|
||||
entries: list[dict[str, Any]],
|
||||
) -> dict[str, Any]:
|
||||
self._ensure_schema(client_id)
|
||||
now = _utc_now_iso()
|
||||
with self.store.connect(client_id) as db:
|
||||
existing = db.execute(
|
||||
"SELECT * FROM tally_write_batches_v2 WHERE source_tool=? AND source_key=? LIMIT 1",
|
||||
(source_tool, source_key),
|
||||
).fetchone()
|
||||
if existing:
|
||||
batch_id = int(existing["id"])
|
||||
else:
|
||||
cur = db.execute(
|
||||
"""INSERT INTO tally_write_batches_v2(
|
||||
source_tool,source_key,operation_kind,company_guid,company_name,status,
|
||||
created_by_user_id,created_at_utc,entry_count
|
||||
) VALUES(?,?,?,?,?,'queued',?,?,?)""",
|
||||
(
|
||||
source_tool,
|
||||
source_key,
|
||||
operation_kind,
|
||||
str(company_guid or ""),
|
||||
str(company_name or ""),
|
||||
created_by_user_id,
|
||||
now,
|
||||
len(entries),
|
||||
),
|
||||
)
|
||||
batch_id = int(cur.lastrowid)
|
||||
|
||||
existing_keys = {
|
||||
str(row["idempotency_key"])
|
||||
for row in db.execute(
|
||||
"SELECT idempotency_key FROM tally_write_entries_v2 WHERE batch_id=?",
|
||||
(batch_id,),
|
||||
).fetchall()
|
||||
}
|
||||
for index, entry in enumerate(entries, 1):
|
||||
key = str(entry.get("idempotency_key") or "").strip()
|
||||
if not key:
|
||||
raise ValueError("Every Tally write instruction requires an idempotency key.")
|
||||
if key in existing_keys:
|
||||
continue
|
||||
db.execute(
|
||||
"""INSERT INTO tally_write_entries_v2(
|
||||
batch_id,line_no,idempotency_key,operation_type,voucher_type,voucher_date,
|
||||
reference,narration,payload_json,status,created_at_utc,updated_at_utc
|
||||
) VALUES(?,?,?,?,?,?,?,?,?,'queued',?,?)""",
|
||||
(
|
||||
batch_id,
|
||||
index,
|
||||
key,
|
||||
str(entry.get("operation_type") or ""),
|
||||
str(entry.get("voucher_type") or ""),
|
||||
str(entry.get("voucher_date") or ""),
|
||||
str(entry.get("reference") or ""),
|
||||
str(entry.get("narration") or ""),
|
||||
json.dumps(entry, ensure_ascii=False, separators=(",", ":")),
|
||||
now,
|
||||
now,
|
||||
),
|
||||
)
|
||||
db.execute(
|
||||
"UPDATE tally_write_batches_v2 SET entry_count=(SELECT COUNT(*) FROM tally_write_entries_v2 WHERE batch_id=?) WHERE id=?",
|
||||
(batch_id, batch_id),
|
||||
)
|
||||
return self.batch_status(client_id, batch_id)
|
||||
|
||||
def queue_vouchers(
|
||||
self,
|
||||
client_id: int,
|
||||
*,
|
||||
source_tool: str,
|
||||
source_key: str,
|
||||
company_guid: str,
|
||||
company_name: str,
|
||||
created_by_user_id: int | None,
|
||||
vouchers: list[dict[str, Any]],
|
||||
) -> dict[str, Any]:
|
||||
entries = []
|
||||
for row in vouchers:
|
||||
entries.append({
|
||||
**row,
|
||||
"operation_type": "VOUCHER_CREATE",
|
||||
"idempotency_key": str(row.get("idempotency_key") or row.get("reference") or "").strip(),
|
||||
})
|
||||
return self._get_or_create_batch(
|
||||
client_id,
|
||||
source_tool=source_tool,
|
||||
source_key=source_key,
|
||||
operation_kind="VOUCHER_CREATE",
|
||||
company_guid=company_guid,
|
||||
company_name=company_name,
|
||||
created_by_user_id=created_by_user_id,
|
||||
entries=entries,
|
||||
)
|
||||
|
||||
def queue_opening_balances(
|
||||
self,
|
||||
client_id: int,
|
||||
*,
|
||||
source_key: str,
|
||||
company_guid: str,
|
||||
company_name: str,
|
||||
created_by_user_id: int | None,
|
||||
ledgers: list[dict[str, Any]],
|
||||
stock_items: list[dict[str, Any]],
|
||||
) -> dict[str, Any]:
|
||||
entries = []
|
||||
for row in ledgers:
|
||||
entries.append({
|
||||
"operation_type": "MASTER_LEDGER_OPENING",
|
||||
"idempotency_key": f"{source_key}:LEDGER:{int(row.get('item_id') or 0)}",
|
||||
"payload": row,
|
||||
})
|
||||
for row in stock_items:
|
||||
entries.append({
|
||||
"operation_type": "MASTER_STOCK_OPENING",
|
||||
"idempotency_key": f"{source_key}:STOCK:{int(row.get('item_id') or 0)}",
|
||||
"payload": row,
|
||||
})
|
||||
return self._get_or_create_batch(
|
||||
client_id,
|
||||
source_tool="OPENING_BALANCE",
|
||||
source_key=source_key,
|
||||
operation_kind="MASTER_OPENING_BALANCE_UPDATE",
|
||||
company_guid=company_guid,
|
||||
company_name=company_name,
|
||||
created_by_user_id=created_by_user_id,
|
||||
entries=entries,
|
||||
)
|
||||
|
||||
def batch_status(self, client_id: int, batch_id: int) -> dict[str, Any]:
|
||||
self._ensure_schema(client_id)
|
||||
with self.store.connect(client_id) as db:
|
||||
batch = db.execute("SELECT * FROM tally_write_batches_v2 WHERE id=?", (int(batch_id),)).fetchone()
|
||||
if not batch:
|
||||
raise ValueError("Tally write batch was not found.")
|
||||
entries = db.execute(
|
||||
"SELECT * FROM tally_write_entries_v2 WHERE batch_id=? ORDER BY line_no,id",
|
||||
(int(batch_id),),
|
||||
).fetchall()
|
||||
return {"batch": dict(batch), "entries": [dict(x) for x in entries]}
|
||||
|
||||
def _entry_payload(self, row: sqlite3.Row | dict[str, Any]) -> dict[str, Any]:
|
||||
try:
|
||||
return json.loads(str(row["payload_json"] or "{}"))
|
||||
except Exception:
|
||||
return {}
|
||||
|
||||
def _voucher_message(self, company_name: str, payload: dict[str, Any]) -> str:
|
||||
engine = NativeVoucherEngine(self.tally)
|
||||
engine.preflight(
|
||||
company_name,
|
||||
voucher_type=str(payload.get("voucher_type") or "Journal"),
|
||||
voucher_date=str(payload.get("voucher_date") or ""),
|
||||
reference=str(payload.get("reference") or ""),
|
||||
total_amount=float(payload.get("total_amount") or 0),
|
||||
lines=list(payload.get("lines") or []),
|
||||
items=list(payload.get("items") or []),
|
||||
)
|
||||
xml = engine.build_xml(
|
||||
company_name,
|
||||
voucher_type=str(payload.get("voucher_type") or "Journal"),
|
||||
voucher_date=str(payload.get("voucher_date") or ""),
|
||||
reference=str(payload.get("reference") or ""),
|
||||
narration=str(payload.get("narration") or ""),
|
||||
lines=list(payload.get("lines") or []),
|
||||
items=list(payload.get("items") or []),
|
||||
)
|
||||
match = re.search(r"(<TALLYMESSAGE\b.*?</TALLYMESSAGE>)", xml, flags=re.I | re.S)
|
||||
if not match:
|
||||
raise ValueError("Could not build native Tally voucher batch message.")
|
||||
return match.group(1)
|
||||
|
||||
def _master_message(self, payload: dict[str, Any]) -> str:
|
||||
operation = str(payload.get("operation_type") or "")
|
||||
row = payload.get("payload") or {}
|
||||
esc = self.tally._xml_escape
|
||||
if operation == "MASTER_LEDGER_OPENING":
|
||||
name = str(row.get("name") or "").strip()
|
||||
amount = float(row.get("target_opening") or 0)
|
||||
if not name:
|
||||
raise ValueError("Ledger name is required for opening-balance update.")
|
||||
return (
|
||||
'<TALLYMESSAGE xmlns:UDF="TallyUDF">'
|
||||
f'<LEDGER NAME="{esc(name)}" ACTION="Alter"><NAME>{esc(name)}</NAME>'
|
||||
f'<OPENINGBALANCE>{amount:.2f}</OPENINGBALANCE></LEDGER></TALLYMESSAGE>'
|
||||
)
|
||||
if operation == "MASTER_STOCK_OPENING":
|
||||
name = str(row.get("name") or "").strip()
|
||||
unit = str(row.get("unit") or "").strip()
|
||||
qty = float(row.get("target_opening_qty") or 0)
|
||||
value = float(row.get("target_opening_value") or 0)
|
||||
rate = abs(float(row.get("target_opening_rate") or 0))
|
||||
if not name or not unit:
|
||||
raise ValueError("Stock Item name and unit are required for opening-balance update.")
|
||||
qty_text = f"{qty:.6f}".rstrip("0").rstrip(".")
|
||||
rate_text = f"{rate:.6f}".rstrip("0").rstrip(".")
|
||||
return (
|
||||
'<TALLYMESSAGE xmlns:UDF="TallyUDF">'
|
||||
f'<STOCKITEM NAME="{esc(name)}" ACTION="Alter"><NAME>{esc(name)}</NAME>'
|
||||
f'<OPENINGBALANCE>{qty_text} {esc(unit)}</OPENINGBALANCE>'
|
||||
f'<OPENINGVALUE>{value:.2f}</OPENINGVALUE>'
|
||||
f'<OPENINGRATE>{rate_text}/{esc(unit)}</OPENINGRATE>'
|
||||
'</STOCKITEM></TALLYMESSAGE>'
|
||||
)
|
||||
raise ValueError(f"Unsupported master batch operation: {operation}")
|
||||
|
||||
def _batch_envelope(self, company_name: str, messages: list[str], *, masters: bool) -> str:
|
||||
target = "All Masters" if masters else "Vouchers"
|
||||
return (
|
||||
"<ENVELOPE><HEADER><VERSION>1</VERSION><TALLYREQUEST>Import</TALLYREQUEST>"
|
||||
f"<TYPE>Data</TYPE><ID>{target}</ID></HEADER><BODY><DESC><STATICVARIABLES>"
|
||||
f"<SVCURRENTCOMPANY>{self.tally._xml_escape(company_name)}</SVCURRENTCOMPANY>"
|
||||
"</STATICVARIABLES></DESC><DATA>" + "".join(messages) + "</DATA></BODY></ENVELOPE>"
|
||||
)
|
||||
|
||||
def _mark_attempt(self, client_id: int, batch_id: int, entry_ids: list[int], request_xml: str) -> tuple[int, int]:
|
||||
now = _utc_now_iso()
|
||||
with self.store.connect(client_id) as db:
|
||||
current = db.execute(
|
||||
"SELECT COALESCE(MAX(attempt_no),0) FROM tally_write_attempts_v2 WHERE batch_id=?",
|
||||
(batch_id,),
|
||||
).fetchone()[0]
|
||||
attempt_no = int(current or 0) + 1
|
||||
cur = db.execute(
|
||||
"""INSERT INTO tally_write_attempts_v2(
|
||||
batch_id,attempt_no,entry_ids_json,request_xml,status,started_at_utc
|
||||
) VALUES(?,?,?,?, 'running', ?)""",
|
||||
(batch_id, attempt_no, json.dumps(entry_ids), request_xml, now),
|
||||
)
|
||||
db.execute("UPDATE tally_write_batches_v2 SET status='posting',started_at_utc=COALESCE(started_at_utc,?) WHERE id=?", (now, batch_id))
|
||||
for entry_id in entry_ids:
|
||||
db.execute("UPDATE tally_write_entries_v2 SET status='posting',updated_at_utc=? WHERE id=?", (now, entry_id))
|
||||
return int(cur.lastrowid), attempt_no
|
||||
|
||||
def _finish_attempt(self, client_id: int, attempt_id: int, response_xml: str, error: str = "") -> None:
|
||||
with self.store.connect(client_id) as db:
|
||||
db.execute(
|
||||
"UPDATE tally_write_attempts_v2 SET status=?,response_xml=?,error_message=?,completed_at_utc=? WHERE id=?",
|
||||
("failed" if error else "completed", response_xml, error, _utc_now_iso(), attempt_id),
|
||||
)
|
||||
|
||||
def _verify_voucher_entry(self, company_name: str, payload: dict[str, Any]) -> dict[str, Any]:
|
||||
voucher_date = str(payload.get("voucher_date") or "")
|
||||
reference = str(payload.get("reference") or "").strip()
|
||||
voucher_type = str(payload.get("voucher_type") or "").strip().casefold()
|
||||
total = round(abs(float(payload.get("total_amount") or 0)), 2)
|
||||
live = getattr(self.tally, "live", self.tally)
|
||||
for voucher in live.export_vouchers(company_name, voucher_date, voucher_date):
|
||||
if voucher_type and str(voucher.get("voucher_type_name") or "").strip().casefold() != voucher_type:
|
||||
continue
|
||||
ref = str(voucher.get("reference") or "").strip()
|
||||
narration = str(voucher.get("narration") or "").strip()
|
||||
if reference and reference not in {ref, narration} and reference not in narration:
|
||||
continue
|
||||
values = [abs(float(x.get("amount") or 0)) for x in (voucher.get("ledger_entries") or [])]
|
||||
if total and values and all(abs(v - total) > max(1.0, total * 0.002) for v in values):
|
||||
continue
|
||||
return {
|
||||
"verified": True,
|
||||
"voucher_id": str(voucher.get("master_id") or voucher.get("guid") or ""),
|
||||
"voucher_number": str(voucher.get("voucher_number") or ""),
|
||||
}
|
||||
return {"verified": False, "voucher_id": "", "voucher_number": ""}
|
||||
|
||||
def _verify_master_entries(self, company_name: str, entries: list[sqlite3.Row]) -> dict[int, tuple[bool, str]]:
|
||||
live = getattr(self.tally, "live", self.tally)
|
||||
masters = live.fetch_accounting_masters(company_name)
|
||||
ledgers = {str(x.get("name") or "").strip().casefold(): x for x in (masters.get("ledgers") or [])}
|
||||
stocks = {str(x.get("name") or "").strip().casefold(): x for x in (masters.get("stock_items") or [])}
|
||||
results: dict[int, tuple[bool, str]] = {}
|
||||
for row in entries:
|
||||
payload = self._entry_payload(row)
|
||||
op = str(payload.get("operation_type") or "")
|
||||
req = payload.get("payload") or {}
|
||||
if op == "MASTER_LEDGER_OPENING":
|
||||
name = str(req.get("name") or "").strip()
|
||||
current = ledgers.get(name.casefold())
|
||||
actual = float((current or {}).get("opening_balance") or 0)
|
||||
target = float(req.get("target_opening") or 0)
|
||||
ok = current is not None and abs(actual - target) <= 0.01
|
||||
results[int(row["id"])] = (ok, f"Ledger opening {actual:.2f}; expected {target:.2f}.")
|
||||
elif op == "MASTER_STOCK_OPENING":
|
||||
name = str(req.get("name") or "").strip()
|
||||
current = stocks.get(name.casefold())
|
||||
aq = float((current or {}).get("opening_balance") or 0)
|
||||
av = float((current or {}).get("opening_value") or 0)
|
||||
tq = float(req.get("target_opening_qty") or 0)
|
||||
tv = float(req.get("target_opening_value") or 0)
|
||||
ok = current is not None and abs(aq - tq) <= 0.000001 and abs(av - tv) <= 0.01
|
||||
results[int(row["id"])] = (ok, f"Stock opening qty/value {aq}/{av:.2f}; expected {tq}/{tv:.2f}.")
|
||||
return results
|
||||
|
||||
def execute(self, client_id: int, batch_id: int, *, batch_size: int | None = None) -> dict[str, Any]:
|
||||
self._ensure_schema(client_id)
|
||||
size = max(1, min(100, int(batch_size or self.DEFAULT_BATCH_SIZE)))
|
||||
status = self.batch_status(client_id, batch_id)
|
||||
batch = status["batch"]
|
||||
company_name = str(batch.get("company_name") or "").strip()
|
||||
company_guid = str(batch.get("company_guid") or "").strip()
|
||||
if not company_name:
|
||||
raise ValueError("Tally company name is missing from the posting batch.")
|
||||
|
||||
loaded = self.tally.get_loaded_companies()
|
||||
matching = [x for x in loaded if (not company_guid or str(x.guid or "").strip() == company_guid)]
|
||||
if company_guid and not matching:
|
||||
raise ValueError("The mapped Tally company is not currently open. Open it in TallyPrime before posting.")
|
||||
if company_guid and matching and str(matching[0].name or "").strip() != company_name:
|
||||
raise ValueError("The open Tally company GUID matches but the company name differs from the approved batch.")
|
||||
if not company_guid:
|
||||
named = [x for x in loaded if str(x.name or "").strip().casefold() == company_name.casefold()]
|
||||
if len(named) != 1:
|
||||
raise ValueError(f"Tally company '{company_name}' is not uniquely open.")
|
||||
|
||||
with self.store.connect(client_id) as db:
|
||||
rows = db.execute(
|
||||
"""SELECT * FROM tally_write_entries_v2
|
||||
WHERE batch_id=? AND status NOT IN ('verified','skipped')
|
||||
ORDER BY line_no,id""",
|
||||
(batch_id,),
|
||||
).fetchall()
|
||||
|
||||
operation_kind = str(batch.get("operation_kind") or "")
|
||||
for offset in range(0, len(rows), size):
|
||||
chunk = rows[offset:offset + size]
|
||||
messages: list[str] = []
|
||||
send_rows: list[sqlite3.Row] = []
|
||||
|
||||
for row in chunk:
|
||||
payload = self._entry_payload(row)
|
||||
if operation_kind == "VOUCHER_CREATE":
|
||||
# Idempotency against Tally itself before any retry.
|
||||
found = self._verify_voucher_entry(company_name, payload)
|
||||
if found.get("verified"):
|
||||
with self.store.connect(client_id) as db:
|
||||
db.execute(
|
||||
"""UPDATE tally_write_entries_v2
|
||||
SET status='verified',tally_voucher_id=?,tally_voucher_number=?,verified_at_utc=?,updated_at_utc=?,last_error=''
|
||||
WHERE id=?""",
|
||||
(found.get("voucher_id") or "", found.get("voucher_number") or "", _utc_now_iso(), _utc_now_iso(), int(row["id"])),
|
||||
)
|
||||
continue
|
||||
messages.append(self._voucher_message(company_name, payload))
|
||||
send_rows.append(row)
|
||||
else:
|
||||
messages.append(self._master_message(payload))
|
||||
send_rows.append(row)
|
||||
|
||||
if not send_rows:
|
||||
continue
|
||||
|
||||
request_xml = self._batch_envelope(company_name, messages, masters=(operation_kind != "VOUCHER_CREATE"))
|
||||
entry_ids = [int(x["id"]) for x in send_rows]
|
||||
attempt_id, attempt_no = self._mark_attempt(client_id, batch_id, entry_ids, request_xml)
|
||||
response_xml = ""
|
||||
try:
|
||||
response_xml = self._call_writer(client_id, batch_id, attempt_no, request_xml)
|
||||
parsed = self._xml_import_result(response_xml)
|
||||
self._finish_attempt(client_id, attempt_id, response_xml)
|
||||
|
||||
if operation_kind == "VOUCHER_CREATE":
|
||||
for row in send_rows:
|
||||
payload = self._entry_payload(row)
|
||||
verified = self._verify_voucher_entry(company_name, payload)
|
||||
with self.store.connect(client_id) as db:
|
||||
if verified.get("verified"):
|
||||
db.execute(
|
||||
"""UPDATE tally_write_entries_v2 SET status='verified',tally_voucher_id=?,tally_voucher_number=?,verified_at_utc=?,updated_at_utc=?,last_error='' WHERE id=?""",
|
||||
(verified.get("voucher_id") or "", verified.get("voucher_number") or "", _utc_now_iso(), _utc_now_iso(), int(row["id"])),
|
||||
)
|
||||
else:
|
||||
err = parsed.get("line_error") or (
|
||||
f"Tally batch response created={parsed.get('created')} altered={parsed.get('altered')} errors={parsed.get('errors')}; voucher was not found during read-back verification."
|
||||
)
|
||||
db.execute("UPDATE tally_write_entries_v2 SET status='failed',last_error=?,updated_at_utc=? WHERE id=?", (err, _utc_now_iso(), int(row["id"])))
|
||||
else:
|
||||
verify = self._verify_master_entries(company_name, send_rows)
|
||||
with self.store.connect(client_id) as db:
|
||||
for row in send_rows:
|
||||
ok, message = verify.get(int(row["id"]), (False, "Verification result unavailable."))
|
||||
db.execute(
|
||||
"UPDATE tally_write_entries_v2 SET status=?,verified_at_utc=?,last_error=?,updated_at_utc=? WHERE id=?",
|
||||
("verified" if ok else "failed", _utc_now_iso() if ok else None, "" if ok else message, _utc_now_iso(), int(row["id"])),
|
||||
)
|
||||
except Exception as exc:
|
||||
self._finish_attempt(client_id, attempt_id, response_xml, str(exc))
|
||||
with self.store.connect(client_id) as db:
|
||||
for row in send_rows:
|
||||
db.execute("UPDATE tally_write_entries_v2 SET status='failed',last_error=?,updated_at_utc=? WHERE id=?", (str(exc), _utc_now_iso(), int(row["id"])))
|
||||
|
||||
with self.store.connect(client_id) as db:
|
||||
counts = db.execute(
|
||||
"""SELECT
|
||||
SUM(CASE WHEN status='verified' THEN 1 ELSE 0 END) AS verified,
|
||||
SUM(CASE WHEN status='failed' THEN 1 ELSE 0 END) AS failed,
|
||||
SUM(CASE WHEN status='skipped' THEN 1 ELSE 0 END) AS skipped,
|
||||
COUNT(*) AS total
|
||||
FROM tally_write_entries_v2 WHERE batch_id=?""",
|
||||
(batch_id,),
|
||||
).fetchone()
|
||||
verified = int(counts["verified"] or 0)
|
||||
failed = int(counts["failed"] or 0)
|
||||
skipped = int(counts["skipped"] or 0)
|
||||
total = int(counts["total"] or 0)
|
||||
batch_status = "verified" if verified + skipped == total and total > 0 else ("partial" if verified > 0 else "failed")
|
||||
db.execute(
|
||||
"""UPDATE tally_write_batches_v2
|
||||
SET status=?,completed_at_utc=?,verified_count=?,failed_count=?,skipped_count=?,last_error=?
|
||||
WHERE id=?""",
|
||||
(batch_status, _utc_now_iso(), verified, failed, skipped, "" if failed == 0 else f"{failed} instruction(s) failed verification.", batch_id),
|
||||
)
|
||||
|
||||
return self.batch_status(client_id, batch_id)
|
||||
Reference in New Issue
Block a user