Files
arrr-erp/app/modules/accounting/purchase_posting_service.py
T
2026-08-22 16:12:22 +05:30

437 lines
17 KiB
Python

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