diff --git a/alembic/versions/20260822_controlled_tally_purchase_posting_phase10.py b/alembic/versions/20260822_controlled_tally_purchase_posting_phase10.py new file mode 100644 index 0000000..935c791 --- /dev/null +++ b/alembic/versions/20260822_controlled_tally_purchase_posting_phase10.py @@ -0,0 +1,64 @@ +"""Phase 10 controlled Tally purchase posting. + +Revision ID: 20260822_purchase_posting_p10 +Revises: 20260822_purchase_enrich_p8 + +Phase 9 is migration-free, so Phase 10 correctly follows the Phase 8 migration head. +""" +from alembic import op +import sqlalchemy as sa + +revision = "20260822_purchase_posting_p10" +down_revision = "20260822_purchase_enrich_p8" +branch_labels = None +depends_on = None + + +def upgrade(): + op.create_table( + "accounting_purchase_posting_attempts", + sa.Column("id", sa.Integer(), primary_key=True), + sa.Column("tenant_id", sa.Integer(), sa.ForeignKey("tenants.id", ondelete="CASCADE"), nullable=False), + sa.Column("client_id", sa.Integer(), sa.ForeignKey("clients.id", ondelete="CASCADE"), nullable=False), + sa.Column("purchase_id", sa.Integer(), sa.ForeignKey("accounting_gstr2b_purchases.id", ondelete="CASCADE"), nullable=False), + sa.Column("workstation_agent_id", sa.Integer(), sa.ForeignKey("erp_workstation_agents.id", ondelete="SET NULL"), nullable=True), + sa.Column("tally_guid", sa.String(120), nullable=False, server_default=""), + sa.Column("tally_company_name", sa.String(240), nullable=False, server_default=""), + sa.Column("status", sa.String(40), nullable=False, server_default="not_started"), + sa.Column("preflight_job_id", sa.Integer(), sa.ForeignKey("erp_agent_jobs.id", ondelete="SET NULL"), nullable=True), + sa.Column("posting_job_id", sa.Integer(), sa.ForeignKey("erp_agent_jobs.id", ondelete="SET NULL"), nullable=True), + sa.Column("preflight_result_json", sa.Text(), nullable=True), + sa.Column("posting_result_json", sa.Text(), nullable=True), + sa.Column("last_error", sa.Text(), nullable=True), + sa.Column("party_ledger_name", sa.String(240), nullable=False, server_default=""), + sa.Column("purchase_ledger_name", sa.String(240), nullable=False, server_default=""), + sa.Column("input_igst_ledger_name", sa.String(240), nullable=False, server_default=""), + sa.Column("input_cgst_ledger_name", sa.String(240), nullable=False, server_default=""), + sa.Column("input_sgst_ledger_name", sa.String(240), nullable=False, server_default=""), + sa.Column("input_cess_ledger_name", sa.String(240), nullable=False, server_default=""), + sa.Column("tally_voucher_id", sa.String(120), nullable=False, server_default=""), + sa.Column("tally_voucher_number", sa.String(160), nullable=False, server_default=""), + sa.Column("tally_reference", sa.String(160), nullable=False, server_default=""), + sa.Column("idempotency_key", sa.String(200), nullable=False, server_default=""), + sa.Column("requested_by_user_id", sa.Integer(), sa.ForeignKey("users.id", ondelete="SET NULL"), nullable=True), + sa.Column("posted_by_user_id", sa.Integer(), sa.ForeignKey("users.id", ondelete="SET NULL"), nullable=True), + sa.Column("created_at_utc", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()), + sa.Column("updated_at_utc", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()), + sa.Column("preflight_completed_at_utc", sa.DateTime(timezone=True), nullable=True), + sa.Column("posting_requested_at_utc", sa.DateTime(timezone=True), nullable=True), + sa.Column("posted_at_utc", sa.DateTime(timezone=True), nullable=True), + sa.UniqueConstraint("tenant_id", "client_id", "purchase_id", name="uq_accounting_purchase_posting_purchase"), + ) + for col in ( + "tenant_id", "client_id", "purchase_id", "workstation_agent_id", "tally_guid", + "status", "preflight_job_id", "posting_job_id", "idempotency_key", "created_at_utc", + ): + op.create_index( + f"ix_accounting_purchase_posting_attempts_{col}", + "accounting_purchase_posting_attempts", + [col], + ) + + +def downgrade(): + op.drop_table("accounting_purchase_posting_attempts") diff --git a/app/modules/accounting/purchase_posting_models.py b/app/modules/accounting/purchase_posting_models.py new file mode 100644 index 0000000..b909d51 --- /dev/null +++ b/app/modules/accounting/purchase_posting_models.py @@ -0,0 +1,77 @@ +from __future__ import annotations + +from datetime import datetime, timezone + +from sqlalchemy import DateTime, ForeignKey, Integer, String, Text, UniqueConstraint +from sqlalchemy.orm import Mapped, mapped_column + +from app.core.db.common import CommonBase + + +class AccountingPurchasePostingAttempt(CommonBase): + """Controlled posting state for one reviewed GSTR-2B purchase. + + Phase 10 intentionally keeps posting separate from the purchase review record so + accounting review history is never overwritten by workstation/Tally execution state. + """ + + __tablename__ = "accounting_purchase_posting_attempts" + __table_args__ = ( + UniqueConstraint( + "tenant_id", "client_id", "purchase_id", + name="uq_accounting_purchase_posting_purchase", + ), + ) + + id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True) + tenant_id: Mapped[int] = mapped_column(ForeignKey("tenants.id", ondelete="CASCADE"), nullable=False, index=True) + client_id: Mapped[int] = mapped_column(ForeignKey("clients.id", ondelete="CASCADE"), nullable=False, index=True) + purchase_id: Mapped[int] = mapped_column( + ForeignKey("accounting_gstr2b_purchases.id", ondelete="CASCADE"), nullable=False, index=True + ) + + workstation_agent_id: Mapped[int | None] = mapped_column( + ForeignKey("erp_workstation_agents.id", ondelete="SET NULL"), nullable=True, index=True + ) + tally_guid: Mapped[str] = mapped_column(String(120), nullable=False, default="", index=True) + tally_company_name: Mapped[str] = mapped_column(String(240), nullable=False, default="") + + status: Mapped[str] = mapped_column(String(40), nullable=False, default="not_started", index=True) + preflight_job_id: Mapped[int | None] = mapped_column( + ForeignKey("erp_agent_jobs.id", ondelete="SET NULL"), nullable=True, index=True + ) + posting_job_id: Mapped[int | None] = mapped_column( + ForeignKey("erp_agent_jobs.id", ondelete="SET NULL"), nullable=True, index=True + ) + + preflight_result_json: Mapped[str | None] = mapped_column(Text, nullable=True) + posting_result_json: Mapped[str | None] = mapped_column(Text, nullable=True) + last_error: Mapped[str | None] = mapped_column(Text, nullable=True) + + party_ledger_name: Mapped[str] = mapped_column(String(240), nullable=False, default="") + purchase_ledger_name: Mapped[str] = mapped_column(String(240), nullable=False, default="") + input_igst_ledger_name: Mapped[str] = mapped_column(String(240), nullable=False, default="") + input_cgst_ledger_name: Mapped[str] = mapped_column(String(240), nullable=False, default="") + input_sgst_ledger_name: Mapped[str] = mapped_column(String(240), nullable=False, default="") + input_cess_ledger_name: Mapped[str] = mapped_column(String(240), nullable=False, default="") + + tally_voucher_id: Mapped[str] = mapped_column(String(120), nullable=False, default="") + tally_voucher_number: Mapped[str] = mapped_column(String(160), nullable=False, default="") + tally_reference: Mapped[str] = mapped_column(String(160), nullable=False, default="") + idempotency_key: Mapped[str] = mapped_column(String(200), nullable=False, default="", index=True) + + requested_by_user_id: Mapped[int | None] = mapped_column(ForeignKey("users.id", ondelete="SET NULL"), nullable=True) + posted_by_user_id: Mapped[int | None] = mapped_column(ForeignKey("users.id", ondelete="SET NULL"), nullable=True) + + created_at_utc: Mapped[datetime] = mapped_column( + DateTime(timezone=True), default=lambda: datetime.now(timezone.utc), nullable=False, index=True + ) + updated_at_utc: Mapped[datetime] = mapped_column( + DateTime(timezone=True), + default=lambda: datetime.now(timezone.utc), + onupdate=lambda: datetime.now(timezone.utc), + nullable=False, + ) + preflight_completed_at_utc: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + posting_requested_at_utc: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + posted_at_utc: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) diff --git a/app/modules/accounting/purchase_posting_service.py b/app/modules/accounting/purchase_posting_service.py new file mode 100644 index 0000000..9d45151 --- /dev/null +++ b/app/modules/accounting/purchase_posting_service.py @@ -0,0 +1,436 @@ +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 diff --git a/app/modules/accounting/purchase_posting_ui.py b/app/modules/accounting/purchase_posting_ui.py new file mode 100644 index 0000000..a2414cb --- /dev/null +++ b/app/modules/accounting/purchase_posting_ui.py @@ -0,0 +1,201 @@ +from __future__ import annotations + +from urllib.parse import urlencode + +from fastapi import APIRouter, Form, Request +from fastapi.responses import RedirectResponse + +from app.core.db.common import CommonSessionLocal +from app.core.security.csrf import get_or_create_csrf_token, validate_csrf +from app.core.templating import templates +from app.modules.accounting.purchase_posting_service import ( + attempts_for_purchases, + preflight_choices, + queue_post, + queue_preflight, + sync_attempts, + visible_workstations, +) +from app.modules.accounting.purchase_review_service import review_queue +from app.modules.accounting.ui import _find_visible_client, _require_partner, _visible_clients +from app.modules.core.rbac.deps import get_user_permissions, get_user_roles + +router = APIRouter(prefix="/tools/accounting/purchase-posting", tags=["accounting-purchase-posting-ui"]) + + +def _redirect(client_id: int, *, message: str = "", error: str = "", page: int = 1, per_page: int = 25): + params = {"client_id": client_id, "page": page, "per_page": per_page} + if message: + params["message"] = message[:240] + if error: + params["error"] = error[:240] + return RedirectResponse( + url="/tools/accounting/purchase-posting?" + urlencode(params), + status_code=303, + ) + + +@router.get("") +def page( + request: Request, + client_id: int | None = None, + page: int = 1, + per_page: int = 25, + message: str = "", + error: str = "", +): + db = CommonSessionLocal() + try: + user, response = _require_partner(request, db, "accounting.learning.view") + if response: + return response + clients, scope = _visible_clients(db, request, user) + selected = next((c for c in clients if client_id and int(c.id) == int(client_id)), None) + + result = None + attempts = {} + choices = {} + workstations = [] + + if selected: + # Phase 10 intentionally starts only from reviewed purchases. + result = review_queue( + db, + tenant_id=scope.tenant_id, + client_id=selected.id, + status="reviewed", + page=page, + per_page=per_page, + ) + ids = [row[0].id for row in result.rows] + attempts = attempts_for_purchases( + db, tenant_id=scope.tenant_id, client_id=selected.id, purchase_ids=ids + ) + sync_attempts(db, list(attempts.values())) + attempts = attempts_for_purchases( + db, tenant_id=scope.tenant_id, client_id=selected.id, purchase_ids=ids + ) + choices = { + pid: preflight_choices(attempt) + for pid, attempt in attempts.items() + if attempt.preflight_result_json + } + + branch_id = getattr(user, "branch_id", None) + workstations = visible_workstations( + db, tenant_id=scope.tenant_id, branch_id=branch_id + ) + + return templates.TemplateResponse( + "modules/accounting/templates/accounting/purchase_posting.html", + { + "request": request, + "current_user": user, + "current_user_roles": get_user_roles(db, user.id), + "current_user_permissions": get_user_permissions(db, user.id), + "csrf_token": get_or_create_csrf_token(request), + "title": "Controlled Tally Purchase Posting", + "clients": clients, + "selected_client": selected, + "result": result, + "attempts": attempts, + "choices": choices, + "workstations": workstations, + "message": message, + "error": error, + }, + ) + finally: + db.close() + + +@router.post("/purchase/{purchase_id}/preflight") +def preflight( + request: Request, + purchase_id: int, + client_id: int = Form(...), + workstation_id: int = Form(...), + page: int = Form(1), + per_page: int = Form(25), + csrf_token: str = Form(...), +): + validate_csrf(request, csrf_token) + db = CommonSessionLocal() + try: + user, response = _require_partner(request, db, "accounting.learning.manage") + if response: + return response + client, _, scope = _find_visible_client(db, request, user, client_id) + if not client: + from app.core.http_responses import ui_access_denied + return ui_access_denied() + + queue_preflight( + db, + tenant_id=scope.tenant_id, + client_id=client.id, + purchase_id=purchase_id, + workstation_id=workstation_id, + user_id=user.id, + ) + return _redirect( + client.id, + message="Workstation preflight queued. Refresh this page after the agent processes the job.", + page=page, + per_page=per_page, + ) + except Exception as exc: + db.rollback() + return _redirect(client_id, error=str(exc), page=page, per_page=per_page) + finally: + db.close() + + +@router.post("/purchase/{purchase_id}/post") +def post( + request: Request, + purchase_id: int, + client_id: int = Form(...), + party_ledger_name: str = Form(...), + input_igst_ledger_name: str = Form(""), + input_cgst_ledger_name: str = Form(""), + input_sgst_ledger_name: str = Form(""), + input_cess_ledger_name: str = Form(""), + page: int = Form(1), + per_page: int = Form(25), + csrf_token: str = Form(...), +): + validate_csrf(request, csrf_token) + db = CommonSessionLocal() + try: + user, response = _require_partner(request, db, "accounting.learning.manage") + if response: + return response + client, _, scope = _find_visible_client(db, request, user, client_id) + if not client: + from app.core.http_responses import ui_access_denied + return ui_access_denied() + + queue_post( + db, + tenant_id=scope.tenant_id, + client_id=client.id, + purchase_id=purchase_id, + party_ledger_name=party_ledger_name, + input_igst_ledger_name=input_igst_ledger_name, + input_cgst_ledger_name=input_cgst_ledger_name, + input_sgst_ledger_name=input_sgst_ledger_name, + input_cess_ledger_name=input_cess_ledger_name, + user_id=user.id, + ) + return _redirect( + client.id, + message="Controlled Purchase voucher job queued. Refresh to see the workstation/Tally result.", + page=page, + per_page=per_page, + ) + except Exception as exc: + db.rollback() + return _redirect(client_id, error=str(exc), page=page, per_page=per_page) + finally: + db.close() diff --git a/app/modules/accounting/templates/accounting/purchase_posting.html b/app/modules/accounting/templates/accounting/purchase_posting.html new file mode 100644 index 0000000..fe862e5 --- /dev/null +++ b/app/modules/accounting/templates/accounting/purchase_posting.html @@ -0,0 +1,252 @@ +{% extends "ui/templates/base/layout.html" %} +{% block content %} +
Tools · Accounting Intelligence
+Phase 10 posts only accountant-reviewed purchases. Every voucher must pass workstation preflight, open-company verification, Tally ledger validation, duplicate screening and a durable Phase 2 idempotent job. Financial write jobs are never automatically retried.
+{{ result.total }} reviewed purchase(s) · page {{ result.page }} of {{ result.pages }}
+