from __future__ import annotations import hashlib import json from datetime import datetime, timezone from sqlalchemy import select from app.modules.accounting.gstr2b_models import AccountingGSTR2BPurchase from app.modules.accounting.purchase_posting_models import AccountingPurchasePostingAttempt from app.modules.documents.agent_jobs import enqueue_agent_job from app.modules.documents.models import ERPAgentJob, ERPWorkstationAgent PREFLIGHT_ACTION = "accounting_purchase_posting_preflight" POST_ACTION = "accounting_post_purchase_voucher" TERMINAL = {"succeeded", "failed", "cancelled"} def _utcnow(): return datetime.now(timezone.utc) def _loads(value, default=None): if not value: return default if default is not None else {} try: parsed = json.loads(value) return parsed except Exception: return default if default is not None else {} def _dumps(value): return json.dumps(value, ensure_ascii=False, separators=(",", ":"), default=str) def _base_key(*, tenant_id: int, client_id: int, purchase: AccountingGSTR2BPurchase) -> str: raw = "|".join([ str(tenant_id), str(client_id), str(purchase.id), purchase.supplier_gstin or "", purchase.invoice_number or "", purchase.invoice_date or "", purchase.document_type or "", ]) digest = hashlib.sha256(raw.encode("utf-8", "ignore")).hexdigest()[:32] return f"purchase:{client_id}:{purchase.id}:{digest}" def visible_workstations(db, *, tenant_id: int, branch_id: int | None = None): stmt = select(ERPWorkstationAgent).where( ERPWorkstationAgent.tenant_id == tenant_id, ERPWorkstationAgent.is_active.is_(True), ) if branch_id is not None: stmt = stmt.where(ERPWorkstationAgent.branch_id == branch_id) return list(db.execute( stmt.order_by( ERPWorkstationAgent.tally_connected.desc(), ERPWorkstationAgent.last_seen_at_utc.desc(), ERPWorkstationAgent.id.desc(), ) ).scalars().all()) def _purchase(db, *, tenant_id: int, client_id: int, purchase_id: int): return db.execute(select(AccountingGSTR2BPurchase).where( AccountingGSTR2BPurchase.id == purchase_id, AccountingGSTR2BPurchase.tenant_id == tenant_id, AccountingGSTR2BPurchase.client_id == client_id, )).scalar_one_or_none() def get_attempt(db, *, tenant_id: int, client_id: int, purchase_id: int): return db.execute(select(AccountingPurchasePostingAttempt).where( AccountingPurchasePostingAttempt.tenant_id == tenant_id, AccountingPurchasePostingAttempt.client_id == client_id, AccountingPurchasePostingAttempt.purchase_id == purchase_id, )).scalar_one_or_none() def attempts_for_purchases(db, *, tenant_id: int, client_id: int, purchase_ids): ids = [int(x) for x in purchase_ids if x] if not ids: return {} rows = list(db.execute(select(AccountingPurchasePostingAttempt).where( AccountingPurchasePostingAttempt.tenant_id == tenant_id, AccountingPurchasePostingAttempt.client_id == client_id, AccountingPurchasePostingAttempt.purchase_id.in_(ids), )).scalars().all()) return {r.purchase_id: r for r in rows} def _ensure_reviewed(purchase: AccountingGSTR2BPurchase): if purchase.review_status != "reviewed": raise ValueError("Only accountant-reviewed purchases may proceed to Tally posting.") if not int(purchase.final_nature_id or 0): raise ValueError("The purchase has no final accounting nature.") if not (purchase.final_ledger_name or "").strip(): raise ValueError("The reviewed purchase has no mapped Tally purchase/expense ledger.") if (purchase.posting_status or "").lower() == "posted": raise ValueError("This purchase is already marked as posted to Tally.") def _accounting_relative_dir(client_id: int) -> str: # This matches the existing Local Agent accounting-store binding convention. return f"Accounting/clients/{int(client_id)}" def queue_preflight( db, *, tenant_id: int, client_id: int, purchase_id: int, workstation_id: int, user_id: int, ): purchase = _purchase(db, tenant_id=tenant_id, client_id=client_id, purchase_id=purchase_id) if not purchase: raise ValueError("Purchase record was not found.") _ensure_reviewed(purchase) workstation = db.get(ERPWorkstationAgent, int(workstation_id)) if not workstation or not workstation.is_active or workstation.tenant_id != tenant_id: raise ValueError("Selected workstation is unavailable.") if not workstation.tally_connected: raise ValueError("Selected workstation is not currently reporting a Tally connection.") tally_guid = (purchase.tally_guid or "").strip() if not tally_guid: raise ValueError("The purchase is not associated with a mapped Tally company.") attempt = get_attempt( db, tenant_id=tenant_id, client_id=client_id, purchase_id=purchase.id ) if not attempt: attempt = AccountingPurchasePostingAttempt( tenant_id=tenant_id, client_id=client_id, purchase_id=purchase.id, tally_guid=tally_guid, purchase_ledger_name=(purchase.final_ledger_name or "").strip(), status="not_started", requested_by_user_id=user_id, idempotency_key=_base_key( tenant_id=tenant_id, client_id=client_id, purchase=purchase ), ) db.add(attempt) db.flush() if attempt.status == "posted": raise ValueError("This purchase has already been posted.") if attempt.status in {"posting_queued", "posting_claimed", "posting_indeterminate"}: raise ValueError("A posting job is already active or requires manual verification.") attempt.workstation_agent_id = workstation.id attempt.tally_guid = tally_guid attempt.purchase_ledger_name = (purchase.final_ledger_name or "").strip() attempt.requested_by_user_id = user_id attempt.last_error = None payload = { "client_id": client_id, "tenant_id": tenant_id, "accounting_relative_dir": _accounting_relative_dir(client_id), "tally_guid": tally_guid, "supplier_gstin": purchase.supplier_gstin or "", "supplier_name": purchase.supplier_name or "", "invoice_number": purchase.invoice_number or "", "invoice_date": purchase.invoice_date or "", "invoice_value": float(purchase.invoice_value or 0), "taxable_value": float(purchase.taxable_value or 0), "igst": float(purchase.igst or 0), "cgst": float(purchase.cgst or 0), "sgst": float(purchase.sgst or 0), "cess": float(purchase.cess or 0), "purchase_ledger_name": attempt.purchase_ledger_name, "erp_purchase_id": purchase.id, "requested_by_user_id": user_id, } job = enqueue_agent_job( db, workstation_agent_id=workstation.id, action=PREFLIGHT_ACTION, payload=payload, idempotency_key=f"{attempt.idempotency_key}:preflight:{workstation.id}", priority=9, max_attempts=2, created_by_user_id=user_id, ) attempt.preflight_job_id = job.id attempt.status = "preflight_queued" attempt.updated_at_utc = _utcnow() db.add(attempt) db.commit() db.refresh(attempt) return attempt def _job_result(job: ERPAgentJob | None): if not job: return {} return _loads(job.result_json, {}) def sync_attempt(db, attempt: AccountingPurchasePostingAttempt): changed = False if attempt.preflight_job_id and attempt.status.startswith("preflight"): job = db.get(ERPAgentJob, attempt.preflight_job_id) if job: if job.status == "claimed" and attempt.status != "preflight_claimed": attempt.status = "preflight_claimed" changed = True elif job.status == "succeeded": result = _job_result(job) attempt.preflight_result_json = _dumps(result) attempt.tally_company_name = str( result.get("company_name") or (result.get("company") or {}).get("name") or "" )[:240] attempt.status = "preflight_ready" attempt.preflight_completed_at_utc = job.completed_at_utc or _utcnow() attempt.last_error = None changed = True elif job.status in {"failed", "cancelled"}: attempt.status = "preflight_failed" attempt.last_error = job.last_error or f"Preflight job {job.status}." changed = True if attempt.posting_job_id and attempt.status.startswith("posting"): job = db.get(ERPAgentJob, attempt.posting_job_id) if job: if job.status == "claimed" and attempt.status != "posting_claimed": attempt.status = "posting_claimed" changed = True elif job.status == "succeeded": result = _job_result(job) attempt.posting_result_json = _dumps(result) tally_result = result.get("tally_result") or result attempt.tally_voucher_id = str( tally_result.get("last_voucher_id") or tally_result.get("voucher_id") or "" )[:120] attempt.tally_voucher_number = str( tally_result.get("voucher_number") or tally_result.get("last_voucher_id") or "" )[:160] attempt.tally_reference = str( result.get("reference") or tally_result.get("reference") or "" )[:160] attempt.status = "posted" attempt.posted_at_utc = job.completed_at_utc or _utcnow() attempt.last_error = None purchase = db.get(AccountingGSTR2BPurchase, attempt.purchase_id) if purchase: purchase.posting_status = "posted" db.add(purchase) changed = True elif job.status == "failed": error = job.last_error or "Tally posting failed." # Phase 2 agent uses this phrase when execution began but acknowledgement # was lost. Never automatically requeue a financial write in this case. if "Verify Tally before retrying" in error or "interrupted before acknowledgement" in error.lower(): attempt.status = "posting_indeterminate" else: attempt.status = "posting_failed" attempt.last_error = error changed = True elif job.status == "cancelled": attempt.status = "posting_failed" attempt.last_error = "Posting job was cancelled." changed = True if changed: attempt.updated_at_utc = _utcnow() db.add(attempt) db.commit() db.refresh(attempt) return attempt def sync_attempts(db, attempts): return [sync_attempt(db, a) for a in attempts] def preflight_choices(attempt: AccountingPurchasePostingAttempt): result = _loads(attempt.preflight_result_json, {}) return { "party_candidates": result.get("party_candidates") or [], "igst_candidates": result.get("igst_candidates") or [], "cgst_candidates": result.get("cgst_candidates") or [], "sgst_candidates": result.get("sgst_candidates") or [], "cess_candidates": result.get("cess_candidates") or [], "duplicate_candidates": result.get("duplicate_candidates") or [], "purchase_ledger_verified": bool(result.get("purchase_ledger_verified")), "company_name": result.get("company_name") or "", "company_guid": result.get("company_guid") or "", } def queue_post( db, *, tenant_id: int, client_id: int, purchase_id: int, party_ledger_name: str, input_igst_ledger_name: str, input_cgst_ledger_name: str, input_sgst_ledger_name: str, input_cess_ledger_name: str, user_id: int, ): purchase = _purchase(db, tenant_id=tenant_id, client_id=client_id, purchase_id=purchase_id) if not purchase: raise ValueError("Purchase record was not found.") _ensure_reviewed(purchase) attempt = get_attempt( db, tenant_id=tenant_id, client_id=client_id, purchase_id=purchase.id ) if not attempt: raise ValueError("Run workstation preflight before posting.") attempt = sync_attempt(db, attempt) if attempt.status != "preflight_ready": raise ValueError("Successful workstation preflight is required before posting.") if not attempt.workstation_agent_id: raise ValueError("No workstation is attached to this posting attempt.") choices = preflight_choices(attempt) if not choices["purchase_ledger_verified"]: raise ValueError("The reviewed purchase/expense ledger was not verified in the open Tally company.") if choices["duplicate_candidates"]: raise ValueError( "A possible duplicate Purchase voucher exists in Tally. Verify it before posting." ) def allowed(name, key, required): value = (name or "").strip() candidates = {str(x.get("name") or "").strip() for x in choices[key]} if required and not value: raise ValueError(f"Select {key.replace('_candidates','').upper()} ledger.") if value and value not in candidates: raise ValueError("Selected ledger was not returned by the current workstation preflight.") return value party = allowed(party_ledger_name, "party_candidates", True) igst = allowed(input_igst_ledger_name, "igst_candidates", float(purchase.igst or 0) > 0) cgst = allowed(input_cgst_ledger_name, "cgst_candidates", float(purchase.cgst or 0) > 0) sgst = allowed(input_sgst_ledger_name, "sgst_candidates", float(purchase.sgst or 0) > 0) cess = allowed(input_cess_ledger_name, "cess_candidates", float(purchase.cess or 0) > 0) if not attempt.tally_company_name: raise ValueError("Preflight did not return a valid Tally company name.") reference = (purchase.invoice_number or f"ERP-{purchase.id}")[:160] narration = ( f"Purchase imported from reviewed ERP GSTR-2B document #{purchase.id}; " f"supplier GSTIN {purchase.supplier_gstin or '-'}; invoice {reference}" )[:1000] payload = { "client_id": client_id, "tenant_id": tenant_id, "accounting_relative_dir": _accounting_relative_dir(client_id), "tally_guid": attempt.tally_guid, "company_name": attempt.tally_company_name, "erp_purchase_id": purchase.id, "invoice_number": purchase.invoice_number or "", "invoice_date": purchase.invoice_date or "", "supplier_gstin": purchase.supplier_gstin or "", "supplier_name": purchase.supplier_name or "", "party_ledger_name": party, "purchase_ledger_name": attempt.purchase_ledger_name, "input_igst_ledger_name": igst, "input_cgst_ledger_name": cgst, "input_sgst_ledger_name": sgst, "input_cess_ledger_name": cess, "taxable_value": round(float(purchase.taxable_value or 0), 2), "igst": round(float(purchase.igst or 0), 2), "cgst": round(float(purchase.cgst or 0), 2), "sgst": round(float(purchase.sgst or 0), 2), "cess": round(float(purchase.cess or 0), 2), "invoice_value": round(float(purchase.invoice_value or 0), 2), "reference": reference, "narration": narration, "requested_by_user_id": user_id, "posting_attempt_id": attempt.id, } workstation = db.get(ERPWorkstationAgent, attempt.workstation_agent_id) if not workstation or not workstation.is_active: raise ValueError("Selected workstation is no longer active.") # Deterministic key means clicking Post again returns the same durable job. job = enqueue_agent_job( db, workstation_agent_id=workstation.id, action=POST_ACTION, payload=payload, idempotency_key=f"{attempt.idempotency_key}:post:{workstation.id}", priority=10, max_attempts=1, # Financial writes are never automatically retried. created_by_user_id=user_id, ) attempt.party_ledger_name = party attempt.input_igst_ledger_name = igst attempt.input_cgst_ledger_name = cgst attempt.input_sgst_ledger_name = sgst attempt.input_cess_ledger_name = cess attempt.posting_job_id = job.id attempt.posting_requested_at_utc = _utcnow() attempt.posted_by_user_id = user_id attempt.status = "posting_queued" attempt.last_error = None attempt.updated_at_utc = _utcnow() db.add(attempt) purchase.posting_status = "queued" db.add(purchase) db.commit() db.refresh(attempt) return attempt