diff --git a/alembic/versions/20260822_einvoice_ewaybill_enrichment_phase8.py b/alembic/versions/20260822_einvoice_ewaybill_enrichment_phase8.py
new file mode 100644
index 0000000..5c9f175
--- /dev/null
+++ b/alembic/versions/20260822_einvoice_ewaybill_enrichment_phase8.py
@@ -0,0 +1,112 @@
+"""Phase 8 E-Invoice and E-Way Bill purchase enrichment.
+
+Revision ID: 20260822_purchase_enrich_p8
+Revises: 20260822_gstr2b_purchase_p7
+"""
+from alembic import op
+import sqlalchemy as sa
+
+revision = "20260822_purchase_enrich_p8"
+down_revision = "20260822_gstr2b_purchase_p7"
+branch_labels = None
+depends_on = None
+
+
+def upgrade():
+ op.create_table(
+ "accounting_purchase_enrichment_batches",
+ 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("tally_guid", sa.String(120), nullable=False, server_default=""),
+ sa.Column("source_type", sa.String(30), nullable=False),
+ sa.Column("original_filename", sa.String(260), nullable=False, server_default=""),
+ sa.Column("file_sha256", sa.String(64), nullable=False),
+ sa.Column("source_period", sa.String(20), nullable=False, server_default=""),
+ sa.Column("status", sa.String(30), nullable=False, server_default="imported"),
+ sa.Column("records_read", sa.Integer(), nullable=False, server_default="0"),
+ sa.Column("records_imported", sa.Integer(), nullable=False, server_default="0"),
+ sa.Column("records_duplicate", sa.Integer(), nullable=False, server_default="0"),
+ sa.Column("records_linked", sa.Integer(), nullable=False, server_default="0"),
+ sa.Column("records_unmatched", sa.Integer(), nullable=False, server_default="0"),
+ sa.Column("records_ambiguous", sa.Integer(), nullable=False, server_default="0"),
+ sa.Column("imported_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("completed_at_utc", sa.DateTime(timezone=True), nullable=True),
+ )
+ for col in ("tenant_id", "client_id", "tally_guid", "source_type", "file_sha256", "source_period", "status", "created_at_utc"):
+ op.create_index(f"ix_accounting_purchase_enrichment_batches_{col}", "accounting_purchase_enrichment_batches", [col])
+
+ op.create_table(
+ "accounting_purchase_enrichment_records",
+ 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("batch_id", sa.Integer(), sa.ForeignKey("accounting_purchase_enrichment_batches.id", ondelete="CASCADE"), nullable=False),
+ sa.Column("tally_guid", sa.String(120), nullable=False, server_default=""),
+ sa.Column("source_type", sa.String(30), nullable=False),
+ sa.Column("source_document_key", sa.String(500), nullable=False),
+ sa.Column("external_reference", sa.String(180), nullable=False, server_default=""),
+ sa.Column("supplier_gstin", sa.String(20), nullable=False, server_default=""),
+ sa.Column("supplier_name", sa.String(260), nullable=False, server_default=""),
+ sa.Column("recipient_gstin", sa.String(20), nullable=False, server_default=""),
+ sa.Column("document_number", sa.String(160), nullable=False, server_default=""),
+ sa.Column("document_date", sa.String(20), nullable=False, server_default=""),
+ sa.Column("document_type", sa.String(40), nullable=False, server_default="invoice"),
+ sa.Column("taxable_value", sa.Float(), nullable=False, server_default="0"),
+ sa.Column("igst", sa.Float(), nullable=False, server_default="0"),
+ sa.Column("cgst", sa.Float(), nullable=False, server_default="0"),
+ sa.Column("sgst", sa.Float(), nullable=False, server_default="0"),
+ sa.Column("cess", sa.Float(), nullable=False, server_default="0"),
+ sa.Column("invoice_value", sa.Float(), nullable=False, server_default="0"),
+ sa.Column("place_of_supply", sa.String(120), nullable=False, server_default=""),
+ sa.Column("transport_mode", sa.String(80), nullable=False, server_default=""),
+ sa.Column("vehicle_number", sa.String(40), nullable=False, server_default=""),
+ sa.Column("transporter_id", sa.String(40), nullable=False, server_default=""),
+ sa.Column("raw_summary_json", sa.Text(), nullable=True),
+ sa.Column("source_row_hash", sa.String(64), nullable=False, server_default=""),
+ sa.Column("gstr2b_purchase_id", sa.Integer(), sa.ForeignKey("accounting_gstr2b_purchases.id", ondelete="SET NULL"), nullable=True),
+ sa.Column("match_status", sa.String(30), nullable=False, server_default="unmatched"),
+ sa.Column("match_method", sa.String(80), nullable=False, server_default=""),
+ sa.Column("match_confidence", sa.Integer(), nullable=False, server_default="0"),
+ sa.Column("manually_linked", sa.Boolean(), nullable=False, server_default=sa.false()),
+ 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.UniqueConstraint(
+ "tenant_id", "client_id", "source_type", "source_document_key",
+ name="uq_accounting_purchase_enrichment_source_document",
+ ),
+ )
+ for col in (
+ "tenant_id", "client_id", "batch_id", "tally_guid", "source_type", "source_document_key",
+ "external_reference", "supplier_gstin", "supplier_name", "recipient_gstin", "document_number",
+ "document_date", "document_type", "source_row_hash", "gstr2b_purchase_id", "match_status", "created_at_utc",
+ ):
+ op.create_index(f"ix_accounting_purchase_enrichment_records_{col}", "accounting_purchase_enrichment_records", [col])
+
+ op.create_table(
+ "accounting_purchase_enrichment_items",
+ sa.Column("id", sa.Integer(), primary_key=True),
+ sa.Column("record_id", sa.Integer(), sa.ForeignKey("accounting_purchase_enrichment_records.id", ondelete="CASCADE"), nullable=False),
+ sa.Column("line_number", sa.Integer(), nullable=False, server_default="0"),
+ sa.Column("product_name", sa.String(500), nullable=False, server_default=""),
+ sa.Column("description_text", sa.Text(), nullable=True),
+ sa.Column("hsn_code", sa.String(20), nullable=False, server_default=""),
+ sa.Column("quantity", sa.Float(), nullable=False, server_default="0"),
+ sa.Column("unit", sa.String(30), nullable=False, server_default=""),
+ sa.Column("unit_price", sa.Float(), nullable=False, server_default="0"),
+ sa.Column("taxable_value", sa.Float(), nullable=False, server_default="0"),
+ sa.Column("gst_rate", sa.Float(), nullable=False, server_default="0"),
+ sa.Column("igst", sa.Float(), nullable=False, server_default="0"),
+ sa.Column("cgst", sa.Float(), nullable=False, server_default="0"),
+ sa.Column("sgst", sa.Float(), nullable=False, server_default="0"),
+ sa.Column("cess", sa.Float(), nullable=False, server_default="0"),
+ )
+ op.create_index("ix_accounting_purchase_enrichment_items_record_id", "accounting_purchase_enrichment_items", ["record_id"])
+ op.create_index("ix_accounting_purchase_enrichment_items_hsn_code", "accounting_purchase_enrichment_items", ["hsn_code"])
+
+
+def downgrade():
+ op.drop_table("accounting_purchase_enrichment_items")
+ op.drop_table("accounting_purchase_enrichment_records")
+ op.drop_table("accounting_purchase_enrichment_batches")
diff --git a/app/modules/accounting/purchase_enrichment_models.py b/app/modules/accounting/purchase_enrichment_models.py
new file mode 100644
index 0000000..38b8b54
--- /dev/null
+++ b/app/modules/accounting/purchase_enrichment_models.py
@@ -0,0 +1,121 @@
+from __future__ import annotations
+
+from datetime import datetime, timezone
+
+from sqlalchemy import Boolean, DateTime, Float, ForeignKey, Integer, String, Text, UniqueConstraint
+from sqlalchemy.orm import Mapped, mapped_column
+
+from app.core.db.common import CommonBase
+
+
+class AccountingPurchaseEnrichmentBatch(CommonBase):
+ __tablename__ = "accounting_purchase_enrichment_batches"
+
+ 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)
+ tally_guid: Mapped[str] = mapped_column(String(120), nullable=False, default="", index=True)
+
+ source_type: Mapped[str] = mapped_column(String(30), nullable=False, index=True)
+ original_filename: Mapped[str] = mapped_column(String(260), nullable=False, default="")
+ file_sha256: Mapped[str] = mapped_column(String(64), nullable=False, index=True)
+ source_period: Mapped[str] = mapped_column(String(20), nullable=False, default="", index=True)
+ status: Mapped[str] = mapped_column(String(30), nullable=False, default="imported", index=True)
+
+ records_read: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
+ records_imported: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
+ records_duplicate: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
+ records_linked: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
+ records_unmatched: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
+ records_ambiguous: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
+
+ imported_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
+ )
+ completed_at_utc: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
+
+
+class AccountingPurchaseEnrichmentRecord(CommonBase):
+ __tablename__ = "accounting_purchase_enrichment_records"
+ __table_args__ = (
+ UniqueConstraint(
+ "tenant_id", "client_id", "source_type", "source_document_key",
+ name="uq_accounting_purchase_enrichment_source_document",
+ ),
+ )
+
+ 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)
+ batch_id: Mapped[int] = mapped_column(
+ ForeignKey("accounting_purchase_enrichment_batches.id", ondelete="CASCADE"), nullable=False, index=True
+ )
+ tally_guid: Mapped[str] = mapped_column(String(120), nullable=False, default="", index=True)
+
+ source_type: Mapped[str] = mapped_column(String(30), nullable=False, index=True)
+ source_document_key: Mapped[str] = mapped_column(String(500), nullable=False, index=True)
+ external_reference: Mapped[str] = mapped_column(String(180), nullable=False, default="", index=True)
+
+ supplier_gstin: Mapped[str] = mapped_column(String(20), nullable=False, default="", index=True)
+ supplier_name: Mapped[str] = mapped_column(String(260), nullable=False, default="", index=True)
+ recipient_gstin: Mapped[str] = mapped_column(String(20), nullable=False, default="", index=True)
+
+ document_number: Mapped[str] = mapped_column(String(160), nullable=False, default="", index=True)
+ document_date: Mapped[str] = mapped_column(String(20), nullable=False, default="", index=True)
+ document_type: Mapped[str] = mapped_column(String(40), nullable=False, default="invoice", index=True)
+
+ taxable_value: Mapped[float] = mapped_column(Float, nullable=False, default=0.0)
+ igst: Mapped[float] = mapped_column(Float, nullable=False, default=0.0)
+ cgst: Mapped[float] = mapped_column(Float, nullable=False, default=0.0)
+ sgst: Mapped[float] = mapped_column(Float, nullable=False, default=0.0)
+ cess: Mapped[float] = mapped_column(Float, nullable=False, default=0.0)
+ invoice_value: Mapped[float] = mapped_column(Float, nullable=False, default=0.0)
+
+ place_of_supply: Mapped[str] = mapped_column(String(120), nullable=False, default="")
+ transport_mode: Mapped[str] = mapped_column(String(80), nullable=False, default="")
+ vehicle_number: Mapped[str] = mapped_column(String(40), nullable=False, default="")
+ transporter_id: Mapped[str] = mapped_column(String(40), nullable=False, default="")
+
+ raw_summary_json: Mapped[str | None] = mapped_column(Text, nullable=True)
+ source_row_hash: Mapped[str] = mapped_column(String(64), nullable=False, default="", index=True)
+
+ gstr2b_purchase_id: Mapped[int | None] = mapped_column(
+ ForeignKey("accounting_gstr2b_purchases.id", ondelete="SET NULL"), nullable=True, index=True
+ )
+ match_status: Mapped[str] = mapped_column(String(30), nullable=False, default="unmatched", index=True)
+ match_method: Mapped[str] = mapped_column(String(80), nullable=False, default="")
+ match_confidence: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
+ manually_linked: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
+
+ 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,
+ )
+
+
+class AccountingPurchaseEnrichmentItem(CommonBase):
+ __tablename__ = "accounting_purchase_enrichment_items"
+
+ id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
+ record_id: Mapped[int] = mapped_column(
+ ForeignKey("accounting_purchase_enrichment_records.id", ondelete="CASCADE"), nullable=False, index=True
+ )
+ line_number: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
+ product_name: Mapped[str] = mapped_column(String(500), nullable=False, default="")
+ description_text: Mapped[str | None] = mapped_column(Text, nullable=True)
+ hsn_code: Mapped[str] = mapped_column(String(20), nullable=False, default="", index=True)
+ quantity: Mapped[float] = mapped_column(Float, nullable=False, default=0.0)
+ unit: Mapped[str] = mapped_column(String(30), nullable=False, default="")
+ unit_price: Mapped[float] = mapped_column(Float, nullable=False, default=0.0)
+ taxable_value: Mapped[float] = mapped_column(Float, nullable=False, default=0.0)
+ gst_rate: Mapped[float] = mapped_column(Float, nullable=False, default=0.0)
+ igst: Mapped[float] = mapped_column(Float, nullable=False, default=0.0)
+ cgst: Mapped[float] = mapped_column(Float, nullable=False, default=0.0)
+ sgst: Mapped[float] = mapped_column(Float, nullable=False, default=0.0)
+ cess: Mapped[float] = mapped_column(Float, nullable=False, default=0.0)
diff --git a/app/modules/accounting/purchase_enrichment_parser.py b/app/modules/accounting/purchase_enrichment_parser.py
new file mode 100644
index 0000000..2fcd9c0
--- /dev/null
+++ b/app/modules/accounting/purchase_enrichment_parser.py
@@ -0,0 +1,457 @@
+from __future__ import annotations
+
+import csv
+import hashlib
+import io
+import json
+import re
+from datetime import date, datetime
+from pathlib import Path
+
+from openpyxl import load_workbook
+
+
+MAX_UPLOAD_BYTES = 25 * 1024 * 1024
+
+
+def _s(value) -> str:
+ if value is None:
+ return ""
+ if isinstance(value, float) and value.is_integer():
+ return str(int(value))
+ return str(value).strip()
+
+
+def _norm(value) -> str:
+ text = _s(value).lower().replace("\n", " ").replace("\r", " ")
+ text = re.sub(r"[_\-]+", " ", text)
+ text = re.sub(r"[^\w/(). ]+", " ", text)
+ return re.sub(r"\s+", " ", text).strip()
+
+
+def _gstin(value) -> str:
+ return re.sub(r"[^A-Z0-9]", "", _s(value).upper())[:15]
+
+
+def _hsn(value) -> str:
+ return re.sub(r"\D", "", _s(value))[:8]
+
+
+def _money(value) -> float:
+ if value in (None, ""):
+ return 0.0
+ if isinstance(value, (int, float)):
+ return float(value)
+ text = _s(value).replace(",", "").replace("₹", "")
+ text = re.sub(r"^\((.*)\)$", r"-\1", text)
+ try:
+ return float(text)
+ except Exception:
+ return 0.0
+
+
+def _date(value) -> str:
+ if value in (None, ""):
+ return ""
+ if isinstance(value, datetime):
+ return value.date().isoformat()
+ if isinstance(value, date):
+ return value.isoformat()
+ text = _s(value)
+ for fmt in (
+ "%d/%m/%Y", "%d-%m-%Y", "%Y-%m-%d", "%d-%b-%Y", "%d %b %Y",
+ "%d.%m.%Y", "%d/%m/%y", "%d-%m-%y",
+ ):
+ try:
+ return datetime.strptime(text, fmt).date().isoformat()
+ except Exception:
+ pass
+ return text[:20]
+
+
+def _doc_type(value: str, number: str = "") -> str:
+ text = f"{value} {number}".upper()
+ if "CREDIT" in text or "CR NOTE" in text:
+ return "credit_note"
+ if "DEBIT" in text or "DR NOTE" in text:
+ return "debit_note"
+ return "invoice"
+
+
+def _sha(value: str) -> str:
+ return hashlib.sha256(value.encode("utf-8", "ignore")).hexdigest()
+
+
+def file_sha256(content: bytes) -> str:
+ return hashlib.sha256(content).hexdigest()
+
+
+def _record_key(source_type: str, row: dict) -> str:
+ external = row.get("external_reference", "")
+ if external:
+ return f"{source_type}|REF|{external}"
+ return "|".join([
+ source_type,
+ row.get("supplier_gstin", ""),
+ row.get("document_number", ""),
+ row.get("document_date", ""),
+ row.get("document_type", ""),
+ ])
+
+
+def _finalize(source_type: str, row: dict, items: list[dict]):
+ row["source_type"] = source_type
+ row["supplier_gstin"] = _gstin(row.get("supplier_gstin"))
+ row["recipient_gstin"] = _gstin(row.get("recipient_gstin"))
+ row["document_number"] = _s(row.get("document_number"))[:160]
+ row["document_date"] = _date(row.get("document_date"))
+ row["document_type"] = _doc_type(_s(row.get("document_type")), row["document_number"])
+ row["supplier_name"] = _s(row.get("supplier_name"))[:260]
+ row["external_reference"] = _s(row.get("external_reference"))[:180]
+ for key in ("taxable_value", "igst", "cgst", "sgst", "cess", "invoice_value"):
+ row[key] = _money(row.get(key))
+ row["place_of_supply"] = _s(row.get("place_of_supply"))[:120]
+ row["transport_mode"] = _s(row.get("transport_mode"))[:80]
+ row["vehicle_number"] = _s(row.get("vehicle_number"))[:40]
+ row["transporter_id"] = _s(row.get("transporter_id"))[:40]
+
+ normalized_items = []
+ for index, item in enumerate(items or [], start=1):
+ product_name = _s(item.get("product_name"))[:500]
+ description = _s(item.get("description_text"))[:4000]
+ hsn = _hsn(item.get("hsn_code"))
+ normalized_items.append({
+ "line_number": int(item.get("line_number") or index),
+ "product_name": product_name,
+ "description_text": description,
+ "hsn_code": hsn,
+ "quantity": _money(item.get("quantity")),
+ "unit": _s(item.get("unit"))[:30],
+ "unit_price": _money(item.get("unit_price")),
+ "taxable_value": _money(item.get("taxable_value")),
+ "gst_rate": _money(item.get("gst_rate")),
+ "igst": _money(item.get("igst")),
+ "cgst": _money(item.get("cgst")),
+ "sgst": _money(item.get("sgst")),
+ "cess": _money(item.get("cess")),
+ })
+
+ if not row["taxable_value"]:
+ row["taxable_value"] = sum(x["taxable_value"] for x in normalized_items)
+ if not row["igst"]:
+ row["igst"] = sum(x["igst"] for x in normalized_items)
+ if not row["cgst"]:
+ row["cgst"] = sum(x["cgst"] for x in normalized_items)
+ if not row["sgst"]:
+ row["sgst"] = sum(x["sgst"] for x in normalized_items)
+ if not row["cess"]:
+ row["cess"] = sum(x["cess"] for x in normalized_items)
+ if not row["invoice_value"]:
+ row["invoice_value"] = row["taxable_value"] + row["igst"] + row["cgst"] + row["sgst"] + row["cess"]
+
+ if not row["document_number"] or not row["document_date"]:
+ return None
+ if not row["supplier_gstin"] and not row["supplier_name"]:
+ return None
+
+ row["source_document_key"] = _record_key(source_type, row)
+ compact = json.dumps({"row": row, "items": normalized_items}, sort_keys=True, default=str, ensure_ascii=False)
+ row["source_row_hash"] = _sha(compact)
+ return row, normalized_items
+
+
+def _path(data, *keys, default=None):
+ cur = data
+ for key in keys:
+ if not isinstance(cur, dict):
+ return default
+ candidates = (key, key.lower(), key.upper(), key.capitalize())
+ found = None
+ for cand in candidates:
+ if cand in cur:
+ found = cur[cand]
+ break
+ if found is None:
+ # case-insensitive fallback
+ kmap = {str(k).lower(): k for k in cur.keys()}
+ real = kmap.get(str(key).lower())
+ if real is None:
+ return default
+ found = cur[real]
+ cur = found
+ return cur
+
+
+def _parse_einvoice_object(obj: dict):
+ doc = _path(obj, "DocDtls", default={}) or {}
+ seller = _path(obj, "SellerDtls", default={}) or {}
+ buyer = _path(obj, "BuyerDtls", default={}) or {}
+ val = _path(obj, "ValDtls", default={}) or {}
+ tran = _path(obj, "TranDtls", default={}) or {}
+ irn = _path(obj, "Irn") or _path(obj, "IRN") or _path(obj, "AckNo") or ""
+
+ row = {
+ "external_reference": irn,
+ "supplier_gstin": _path(seller, "Gstin") or _path(obj, "SellerGstin") or _path(obj, "supplier_gstin"),
+ "supplier_name": _path(seller, "TrdNm") or _path(seller, "LglNm") or _path(obj, "SellerName") or _path(obj, "supplier_name"),
+ "recipient_gstin": _path(buyer, "Gstin") or _path(obj, "BuyerGstin") or _path(obj, "recipient_gstin"),
+ "document_number": _path(doc, "No") or _path(obj, "DocNo") or _path(obj, "invoice_number"),
+ "document_date": _path(doc, "Dt") or _path(obj, "DocDt") or _path(obj, "invoice_date"),
+ "document_type": _path(doc, "Typ") or _path(obj, "DocTyp") or "invoice",
+ "taxable_value": _path(val, "AssVal") or _path(obj, "TaxableValue"),
+ "igst": _path(val, "IgstVal") or _path(obj, "IgstVal"),
+ "cgst": _path(val, "CgstVal") or _path(obj, "CgstVal"),
+ "sgst": _path(val, "SgstVal") or _path(obj, "SgstVal"),
+ "cess": _path(val, "CesVal") or _path(obj, "CessVal"),
+ "invoice_value": _path(val, "TotInvVal") or _path(obj, "TotInvVal") or _path(obj, "invoice_value"),
+ "place_of_supply": _path(buyer, "Pos") or _path(obj, "Pos"),
+ "transport_mode": _path(tran, "TransMode") or _path(obj, "TransMode"),
+ "vehicle_number": _path(obj, "VehNo") or "",
+ "transporter_id": _path(tran, "TransId") or _path(obj, "TransId"),
+ }
+ raw_items = _path(obj, "ItemList", default=[]) or _path(obj, "items", default=[]) or []
+ items = []
+ for i, item in enumerate(raw_items if isinstance(raw_items, list) else [], start=1):
+ ass = _path(item, "AssAmt") or _path(item, "TaxableValue") or 0
+ rate = _path(item, "GstRt") or _path(item, "GSTRate") or 0
+ items.append({
+ "line_number": _path(item, "SlNo") or i,
+ "product_name": _path(item, "PrdDesc") or _path(item, "ProductName") or _path(item, "Nm"),
+ "description_text": _path(item, "PrdDesc") or _path(item, "Desc"),
+ "hsn_code": _path(item, "HsnCd") or _path(item, "HSN"),
+ "quantity": _path(item, "Qty"),
+ "unit": _path(item, "Unit"),
+ "unit_price": _path(item, "UnitPrice"),
+ "taxable_value": ass,
+ "gst_rate": rate,
+ "igst": _path(item, "IgstAmt"),
+ "cgst": _path(item, "CgstAmt"),
+ "sgst": _path(item, "SgstAmt"),
+ "cess": _path(item, "CesAmt"),
+ })
+ return _finalize("e_invoice", row, items)
+
+
+def _parse_ewaybill_object(obj: dict):
+ row = {
+ "external_reference": _path(obj, "ewbNo") or _path(obj, "EwbNo") or _path(obj, "ewayBillNo"),
+ "supplier_gstin": _path(obj, "fromGstin") or _path(obj, "supplierGstin") or _path(obj, "supplier_gstin"),
+ "supplier_name": _path(obj, "fromTrdName") or _path(obj, "fromPlace") or _path(obj, "supplierName"),
+ "recipient_gstin": _path(obj, "toGstin") or _path(obj, "recipientGstin") or _path(obj, "recipient_gstin"),
+ "document_number": _path(obj, "docNo") or _path(obj, "documentNumber") or _path(obj, "invoice_number"),
+ "document_date": _path(obj, "docDate") or _path(obj, "documentDate") or _path(obj, "invoice_date"),
+ "document_type": _path(obj, "docType") or "invoice",
+ "taxable_value": _path(obj, "totalValue") or _path(obj, "taxableValue"),
+ "igst": _path(obj, "igstValue"),
+ "cgst": _path(obj, "cgstValue"),
+ "sgst": _path(obj, "sgstValue"),
+ "cess": _path(obj, "cessValue"),
+ "invoice_value": _path(obj, "totInvValue") or _path(obj, "invoiceValue"),
+ "place_of_supply": _path(obj, "toStateCode") or _path(obj, "placeOfSupply"),
+ "transport_mode": _path(obj, "transMode") or _path(obj, "transportMode"),
+ "vehicle_number": _path(obj, "vehicleNo") or _path(obj, "vehNo"),
+ "transporter_id": _path(obj, "transporterId") or _path(obj, "transporterGstin"),
+ }
+ raw_items = _path(obj, "itemList", default=[]) or _path(obj, "items", default=[]) or []
+ items = []
+ for i, item in enumerate(raw_items if isinstance(raw_items, list) else [], start=1):
+ items.append({
+ "line_number": i,
+ "product_name": _path(item, "productName") or _path(item, "productDesc"),
+ "description_text": _path(item, "productDesc") or _path(item, "description"),
+ "hsn_code": _path(item, "hsnCode") or _path(item, "hsn"),
+ "quantity": _path(item, "quantity") or _path(item, "qty"),
+ "unit": _path(item, "qtyUnit") or _path(item, "unit"),
+ "unit_price": _path(item, "unitPrice"),
+ "taxable_value": _path(item, "taxableAmount") or _path(item, "taxableValue"),
+ "gst_rate": (
+ _money(_path(item, "igstRate"))
+ or _money(_path(item, "cgstRate")) + _money(_path(item, "sgstRate"))
+ ),
+ "igst": _path(item, "igstValue") or _path(item, "igstAmount"),
+ "cgst": _path(item, "cgstValue") or _path(item, "cgstAmount"),
+ "sgst": _path(item, "sgstValue") or _path(item, "sgstAmount"),
+ "cess": _path(item, "cessValue") or _path(item, "cessAmount"),
+ })
+ return _finalize("e_way_bill", row, items)
+
+
+HEADER_ALIASES = {
+ "external_reference": {"irn", "ack no", "ackno", "e way bill no", "eway bill no", "ewb no", "ewbno"},
+ "supplier_gstin": {"supplier gstin", "seller gstin", "from gstin", "fromgstin", "gstin of supplier"},
+ "supplier_name": {"supplier name", "seller name", "from trade name", "fromtrdname", "trade/legal name"},
+ "recipient_gstin": {"recipient gstin", "buyer gstin", "to gstin", "togstin"},
+ "document_number": {"document number", "doc no", "docno", "invoice number", "invoice no"},
+ "document_date": {"document date", "doc date", "docdate", "invoice date"},
+ "document_type": {"document type", "doc type", "doctype", "invoice type"},
+ "taxable_value": {"taxable value", "taxable amount", "total value", "totalvalue"},
+ "igst": {"igst", "igst value", "igst amount"},
+ "cgst": {"cgst", "cgst value", "cgst amount"},
+ "sgst": {"sgst", "sgst value", "sgst amount"},
+ "cess": {"cess", "cess value", "cess amount"},
+ "invoice_value": {"invoice value", "total invoice value", "tot inv value", "totinvvalue"},
+ "place_of_supply": {"place of supply", "pos", "to state code"},
+ "transport_mode": {"transport mode", "trans mode", "transmode"},
+ "vehicle_number": {"vehicle number", "vehicle no", "veh no"},
+ "transporter_id": {"transporter id", "trans id", "transporter gstin"},
+ "hsn_code": {"hsn", "hsn code", "hsn/sac"},
+ "product_name": {"product name", "item name", "product"},
+ "description_text": {"description", "product description", "item description"},
+ "quantity": {"quantity", "qty"},
+ "unit": {"unit", "qty unit", "uom"},
+ "unit_price": {"unit price", "rate"},
+}
+ALIAS = {a: k for k, vals in HEADER_ALIASES.items() for a in vals}
+
+
+def _canon_header(value):
+ n = _norm(value)
+ if n in ALIAS:
+ return ALIAS[n]
+ for alias, key in ALIAS.items():
+ if len(alias) >= 8 and alias in n:
+ return key
+ return None
+
+
+def _sheet_records(content: bytes, source_type: str):
+ wb = load_workbook(io.BytesIO(content), read_only=True, data_only=True)
+ out = []
+ try:
+ for ws in wb.worksheets:
+ rows = [list(r) for r in ws.iter_rows(values_only=True)]
+ best = None
+ for idx, vals in enumerate(rows[:40]):
+ mapping = {}
+ for c, v in enumerate(vals):
+ key = _canon_header(v)
+ if key and key not in mapping:
+ mapping[key] = c
+ score = sum(k in mapping for k in ("document_number", "document_date", "supplier_gstin"))
+ if score >= 2 and (best is None or score > best[0]):
+ best = (score, idx, mapping)
+ if not best:
+ continue
+ _, header_idx, mapping = best
+ for row_no, vals in enumerate(rows[header_idx + 1:], start=header_idx + 2):
+ def get(k):
+ c = mapping.get(k)
+ return vals[c] if c is not None and c < len(vals) else None
+ row = {k: get(k) for k in (
+ "external_reference", "supplier_gstin", "supplier_name", "recipient_gstin",
+ "document_number", "document_date", "document_type", "taxable_value", "igst",
+ "cgst", "sgst", "cess", "invoice_value", "place_of_supply", "transport_mode",
+ "vehicle_number", "transporter_id",
+ )}
+ item = {
+ "line_number": row_no,
+ "product_name": get("product_name"),
+ "description_text": get("description_text"),
+ "hsn_code": get("hsn_code"),
+ "quantity": get("quantity"),
+ "unit": get("unit"),
+ "unit_price": get("unit_price"),
+ "taxable_value": get("taxable_value"),
+ }
+ finalized = _finalize(source_type, row, [item] if any(_s(v) for v in item.values()) else [])
+ if finalized:
+ out.append(finalized)
+ finally:
+ wb.close()
+ return out
+
+
+def _csv_records(content: bytes, source_type: str):
+ text = None
+ for enc in ("utf-8-sig", "utf-8", "cp1252", "latin-1"):
+ try:
+ text = content.decode(enc)
+ break
+ except UnicodeDecodeError:
+ pass
+ text = text or content.decode("utf-8", "replace")
+ try:
+ dialect = csv.Sniffer().sniff(text[:8192], delimiters=",;\t|")
+ except Exception:
+ dialect = csv.excel
+ rows = list(csv.reader(io.StringIO(text), dialect))
+ best = None
+ for idx, vals in enumerate(rows[:40]):
+ mapping = {}
+ for c, v in enumerate(vals):
+ key = _canon_header(v)
+ if key and key not in mapping:
+ mapping[key] = c
+ score = sum(k in mapping for k in ("document_number", "document_date", "supplier_gstin"))
+ if score >= 2 and (best is None or score > best[0]):
+ best = (score, idx, mapping)
+ if not best:
+ return []
+ _, header_idx, mapping = best
+ result = []
+ for row_no, vals in enumerate(rows[header_idx + 1:], start=header_idx + 2):
+ def get(k):
+ c = mapping.get(k)
+ return vals[c] if c is not None and c < len(vals) else None
+ row = {k: get(k) for k in (
+ "external_reference", "supplier_gstin", "supplier_name", "recipient_gstin",
+ "document_number", "document_date", "document_type", "taxable_value", "igst",
+ "cgst", "sgst", "cess", "invoice_value", "place_of_supply", "transport_mode",
+ "vehicle_number", "transporter_id",
+ )}
+ item = {
+ "line_number": row_no,
+ "product_name": get("product_name"),
+ "description_text": get("description_text"),
+ "hsn_code": get("hsn_code"),
+ "quantity": get("quantity"),
+ "unit": get("unit"),
+ "unit_price": get("unit_price"),
+ "taxable_value": get("taxable_value"),
+ }
+ finalized = _finalize(source_type, row, [item] if any(_s(v) for v in item.values()) else [])
+ if finalized:
+ result.append(finalized)
+ return result
+
+
+def _json_objects(payload):
+ if isinstance(payload, list):
+ return payload
+ if not isinstance(payload, dict):
+ return []
+ for key in ("data", "result", "records", "invoices", "ewayBills", "ewaybills", "items"):
+ val = payload.get(key)
+ if isinstance(val, list) and val and isinstance(val[0], dict):
+ return val
+ return [payload]
+
+
+def parse_enrichment(content: bytes, filename: str, source_type: str):
+ if source_type not in {"e_invoice", "e_way_bill"}:
+ raise ValueError("Source type must be E-Invoice or E-Way Bill.")
+ if not content:
+ raise ValueError("Uploaded enrichment file is empty.")
+ if len(content) > MAX_UPLOAD_BYTES:
+ raise ValueError("Enrichment upload exceeds the 25 MB limit.")
+
+ suffix = Path(filename or "").suffix.lower()
+ records = []
+ if suffix == ".json":
+ payload = json.loads(content.decode("utf-8-sig"))
+ for obj in _json_objects(payload):
+ if not isinstance(obj, dict):
+ continue
+ item = _parse_einvoice_object(obj) if source_type == "e_invoice" else _parse_ewaybill_object(obj)
+ if item:
+ records.append(item)
+ elif suffix in {".xlsx", ".xlsm"}:
+ records = _sheet_records(content, source_type)
+ elif suffix in {".csv", ".txt"}:
+ records = _csv_records(content, source_type)
+ else:
+ raise ValueError("Upload a .json, .xlsx, .xlsm or .csv file.")
+
+ if not records:
+ raise ValueError("No usable E-Invoice/E-Way Bill document records were detected.")
+ return records
diff --git a/app/modules/accounting/purchase_enrichment_service.py b/app/modules/accounting/purchase_enrichment_service.py
new file mode 100644
index 0000000..055dbd2
--- /dev/null
+++ b/app/modules/accounting/purchase_enrichment_service.py
@@ -0,0 +1,350 @@
+from __future__ import annotations
+
+import json
+import re
+from collections import Counter
+from datetime import datetime, timezone
+
+from sqlalchemy import func, select
+
+from app.modules.accounting.gstr2b_models import AccountingGSTR2BPurchase
+from app.modules.accounting.gstr2b_service import analyze_purchase
+from app.modules.accounting.purchase_enrichment_models import (
+ AccountingPurchaseEnrichmentBatch,
+ AccountingPurchaseEnrichmentItem,
+ AccountingPurchaseEnrichmentRecord,
+)
+from app.modules.accounting.purchase_enrichment_parser import file_sha256, parse_enrichment
+
+
+def _utcnow():
+ return datetime.now(timezone.utc)
+
+
+def _norm_invoice(value: str) -> str:
+ return re.sub(r"[^A-Z0-9]", "", str(value or "").upper())
+
+
+def _value_close(a: float, b: float) -> bool:
+ a = float(a or 0)
+ b = float(b or 0)
+ tolerance = max(2.0, abs(a) * 0.005)
+ return abs(a - b) <= tolerance
+
+
+def _candidate_rows(db, record: AccountingPurchaseEnrichmentRecord):
+ stmt = select(AccountingGSTR2BPurchase).where(
+ AccountingGSTR2BPurchase.tenant_id == record.tenant_id,
+ AccountingGSTR2BPurchase.client_id == record.client_id,
+ )
+ if record.supplier_gstin:
+ stmt = stmt.where(AccountingGSTR2BPurchase.supplier_gstin == record.supplier_gstin)
+ return list(db.execute(stmt).scalars().all())
+
+
+def match_record(db, record: AccountingPurchaseEnrichmentRecord):
+ rows = _candidate_rows(db, record)
+ inv = _norm_invoice(record.document_number)
+ date = record.document_date
+ exact = [
+ r for r in rows
+ if _norm_invoice(r.invoice_number) == inv
+ and r.invoice_date == date
+ and r.document_type == record.document_type
+ ]
+ if len(exact) == 1:
+ record.gstr2b_purchase_id = exact[0].id
+ record.match_status = "linked"
+ record.match_method = "gstin_invoice_date_document_type"
+ record.match_confidence = 100
+ return exact[0]
+ if len(exact) > 1:
+ record.match_status = "ambiguous"
+ record.match_method = "multiple_exact_candidates"
+ record.match_confidence = 70
+ record.gstr2b_purchase_id = None
+ return None
+
+ same_invoice = [r for r in rows if _norm_invoice(r.invoice_number) == inv]
+ if len(same_invoice) == 1:
+ record.gstr2b_purchase_id = same_invoice[0].id
+ record.match_status = "linked"
+ record.match_method = "gstin_invoice_number"
+ record.match_confidence = 94
+ return same_invoice[0]
+
+ date_value = [
+ r for r in rows
+ if r.invoice_date == date and _value_close(r.invoice_value, record.invoice_value)
+ ]
+ if len(date_value) == 1:
+ record.gstr2b_purchase_id = date_value[0].id
+ record.match_status = "linked"
+ record.match_method = "gstin_date_invoice_value"
+ record.match_confidence = 86
+ return date_value[0]
+
+ record.gstr2b_purchase_id = None
+ record.match_status = "ambiguous" if len(same_invoice) > 1 or len(date_value) > 1 else "unmatched"
+ record.match_method = "multiple_candidates" if record.match_status == "ambiguous" else "no_match"
+ record.match_confidence = 50 if record.match_status == "ambiguous" else 0
+ return None
+
+
+def enriched_context(db, purchase_id: int):
+ records = list(db.execute(select(AccountingPurchaseEnrichmentRecord).where(
+ AccountingPurchaseEnrichmentRecord.gstr2b_purchase_id == purchase_id,
+ AccountingPurchaseEnrichmentRecord.match_status == "linked",
+ )).scalars().all())
+ if not records:
+ return {"description": "", "hsn_code": "", "records": [], "items": []}
+
+ record_ids = [r.id for r in records]
+ items = list(db.execute(select(AccountingPurchaseEnrichmentItem).where(
+ AccountingPurchaseEnrichmentItem.record_id.in_(record_ids)
+ ).order_by(
+ AccountingPurchaseEnrichmentItem.record_id,
+ AccountingPurchaseEnrichmentItem.line_number,
+ )).scalars().all())
+
+ descriptions = []
+ hsns = []
+ for item in items:
+ if item.product_name:
+ descriptions.append(item.product_name)
+ if item.description_text and item.description_text not in descriptions:
+ descriptions.append(item.description_text)
+ if item.hsn_code:
+ hsns.append(item.hsn_code)
+
+ # Only use one HSN as the Phase 6 scalar HSN input when all enriched item lines agree.
+ unique_hsn = sorted(set(hsns))
+ hsn = unique_hsn[0] if len(unique_hsn) == 1 else ""
+ description = " | ".join(dict.fromkeys(descriptions))[:12000]
+ return {
+ "description": description,
+ "hsn_code": hsn,
+ "records": records,
+ "items": items,
+ "all_hsn_codes": unique_hsn,
+ }
+
+
+def reanalyze_linked_purchase(db, purchase: AccountingGSTR2BPurchase):
+ context = enriched_context(db, purchase.id)
+ original_desc = purchase.description_text or ""
+ original_hsn = purchase.hsn_code or ""
+
+ # Preserve the GSTR-2B source record. Temporarily enrich the classifier inputs only.
+ merged_desc = " | ".join(x for x in (original_desc, context["description"]) if x)
+ saved_desc = purchase.description_text
+ saved_hsn = purchase.hsn_code
+ try:
+ purchase.description_text = merged_desc[:12000] or None
+ if context["hsn_code"]:
+ purchase.hsn_code = context["hsn_code"]
+ analyze_purchase(db, purchase)
+ finally:
+ purchase.description_text = saved_desc
+ purchase.hsn_code = saved_hsn
+
+ reasons = []
+ try:
+ reasons = json.loads(purchase.suggestion_explanation_json or "[]")
+ except Exception:
+ reasons = []
+ if context["records"]:
+ srcs = sorted({r.source_type.replace("_", " ").title() for r in context["records"]})
+ reasons.insert(0, f"Enriched with linked {' + '.join(srcs)} source data.")
+ if context["all_hsn_codes"]:
+ reasons.insert(1, f"Enriched item HSN(s): {', '.join(context['all_hsn_codes'][:8])}.")
+ purchase.suggestion_explanation_json = json.dumps(reasons, ensure_ascii=False)
+ db.add(purchase)
+ return purchase
+
+
+def _update_batch_counts(db, batch: AccountingPurchaseEnrichmentBatch):
+ rows = list(db.execute(select(AccountingPurchaseEnrichmentRecord).where(
+ AccountingPurchaseEnrichmentRecord.batch_id == batch.id
+ )).scalars().all())
+ counter = Counter(r.match_status for r in rows)
+ batch.records_linked = counter.get("linked", 0)
+ batch.records_unmatched = counter.get("unmatched", 0)
+ batch.records_ambiguous = counter.get("ambiguous", 0)
+ batch.status = "matched"
+ batch.completed_at_utc = _utcnow()
+ db.add(batch)
+
+
+def import_enrichment(
+ db, *,
+ tenant_id: int,
+ client_id: int,
+ tally_guid: str,
+ source_type: str,
+ source_period: str,
+ filename: str,
+ content: bytes,
+ user_id: int,
+):
+ digest = file_sha256(content)
+ old = db.execute(select(AccountingPurchaseEnrichmentBatch).where(
+ AccountingPurchaseEnrichmentBatch.tenant_id == tenant_id,
+ AccountingPurchaseEnrichmentBatch.client_id == client_id,
+ AccountingPurchaseEnrichmentBatch.source_type == source_type,
+ AccountingPurchaseEnrichmentBatch.file_sha256 == digest,
+ )).scalar_one_or_none()
+ if old:
+ return old, True
+
+ parsed = parse_enrichment(content, filename, source_type)
+ batch = AccountingPurchaseEnrichmentBatch(
+ tenant_id=tenant_id,
+ client_id=client_id,
+ tally_guid=(tally_guid or "").strip(),
+ source_type=source_type,
+ original_filename=(filename or "")[:260],
+ file_sha256=digest,
+ source_period=(source_period or "")[:20],
+ status="importing",
+ records_read=len(parsed),
+ imported_by_user_id=user_id,
+ )
+ db.add(batch)
+ db.flush()
+
+ imported = duplicate = 0
+ linked_purchases = set()
+ for row, items in parsed:
+ existing = db.execute(select(AccountingPurchaseEnrichmentRecord.id).where(
+ AccountingPurchaseEnrichmentRecord.tenant_id == tenant_id,
+ AccountingPurchaseEnrichmentRecord.client_id == client_id,
+ AccountingPurchaseEnrichmentRecord.source_type == source_type,
+ AccountingPurchaseEnrichmentRecord.source_document_key == row["source_document_key"],
+ )).scalar_one_or_none()
+ if existing:
+ duplicate += 1
+ continue
+
+ raw_summary = {
+ "source_type": source_type,
+ "source_document_key": row["source_document_key"],
+ "item_count": len(items),
+ }
+ record = AccountingPurchaseEnrichmentRecord(
+ tenant_id=tenant_id,
+ client_id=client_id,
+ batch_id=batch.id,
+ tally_guid=(tally_guid or "").strip(),
+ raw_summary_json=json.dumps(raw_summary, ensure_ascii=False),
+ **row,
+ )
+ db.add(record)
+ db.flush()
+ for item in items:
+ db.add(AccountingPurchaseEnrichmentItem(record_id=record.id, **item))
+
+ purchase = match_record(db, record)
+ if purchase:
+ linked_purchases.add(purchase.id)
+ imported += 1
+
+ batch.records_imported = imported
+ batch.records_duplicate = duplicate
+ db.flush()
+ _update_batch_counts(db, batch)
+
+ for purchase_id in linked_purchases:
+ purchase = db.get(AccountingGSTR2BPurchase, purchase_id)
+ if purchase:
+ reanalyze_linked_purchase(db, purchase)
+
+ db.commit()
+ db.refresh(batch)
+ return batch, False
+
+
+def rematch_batch(db, batch: AccountingPurchaseEnrichmentBatch):
+ rows = list(db.execute(select(AccountingPurchaseEnrichmentRecord).where(
+ AccountingPurchaseEnrichmentRecord.batch_id == batch.id
+ )).scalars().all())
+ linked = set()
+ for record in rows:
+ if record.manually_linked and record.gstr2b_purchase_id:
+ linked.add(record.gstr2b_purchase_id)
+ continue
+ purchase = match_record(db, record)
+ if purchase:
+ linked.add(purchase.id)
+ _update_batch_counts(db, batch)
+ for purchase_id in linked:
+ purchase = db.get(AccountingGSTR2BPurchase, purchase_id)
+ if purchase:
+ reanalyze_linked_purchase(db, purchase)
+ db.commit()
+ return len(rows)
+
+
+def manually_link_record(db, record: AccountingPurchaseEnrichmentRecord, purchase: AccountingGSTR2BPurchase):
+ if record.tenant_id != purchase.tenant_id or record.client_id != purchase.client_id:
+ raise ValueError("The enrichment record and GSTR-2B purchase belong to different client scopes.")
+ record.gstr2b_purchase_id = purchase.id
+ record.match_status = "linked"
+ record.match_method = "manual_user_link"
+ record.match_confidence = 100
+ record.manually_linked = True
+ db.add(record)
+ batch = db.get(AccountingPurchaseEnrichmentBatch, record.batch_id)
+ if batch:
+ _update_batch_counts(db, batch)
+ reanalyze_linked_purchase(db, purchase)
+ db.commit()
+ return record
+
+
+def unlink_record(db, record: AccountingPurchaseEnrichmentRecord):
+ purchase_id = record.gstr2b_purchase_id
+ record.gstr2b_purchase_id = None
+ record.match_status = "unmatched"
+ record.match_method = "manual_unlink"
+ record.match_confidence = 0
+ record.manually_linked = False
+ db.add(record)
+ batch = db.get(AccountingPurchaseEnrichmentBatch, record.batch_id)
+ if batch:
+ _update_batch_counts(db, batch)
+ if purchase_id:
+ purchase = db.get(AccountingGSTR2BPurchase, purchase_id)
+ if purchase:
+ reanalyze_linked_purchase(db, purchase)
+ db.commit()
+ return record
+
+
+def batches_for_client(db, tenant_id: int, client_id: int, limit: int = 30):
+ return list(db.execute(select(AccountingPurchaseEnrichmentBatch).where(
+ AccountingPurchaseEnrichmentBatch.tenant_id == tenant_id,
+ AccountingPurchaseEnrichmentBatch.client_id == client_id,
+ ).order_by(AccountingPurchaseEnrichmentBatch.id.desc()).limit(limit)).scalars().all())
+
+
+def records_for_batch(db, batch_id: int, limit: int = 300):
+ return list(db.execute(select(AccountingPurchaseEnrichmentRecord).where(
+ AccountingPurchaseEnrichmentRecord.batch_id == batch_id
+ ).order_by(AccountingPurchaseEnrichmentRecord.id.desc()).limit(limit)).scalars().all())
+
+
+def items_for_records(db, record_ids):
+ ids = [int(x) for x in record_ids if x]
+ if not ids:
+ return {}
+ rows = list(db.execute(select(AccountingPurchaseEnrichmentItem).where(
+ AccountingPurchaseEnrichmentItem.record_id.in_(ids)
+ ).order_by(
+ AccountingPurchaseEnrichmentItem.record_id,
+ AccountingPurchaseEnrichmentItem.line_number,
+ )).scalars().all())
+ result = {}
+ for item in rows:
+ result.setdefault(item.record_id, []).append(item)
+ return result
diff --git a/app/modules/accounting/purchase_enrichment_ui.py b/app/modules/accounting/purchase_enrichment_ui.py
new file mode 100644
index 0000000..8f55d71
--- /dev/null
+++ b/app/modules/accounting/purchase_enrichment_ui.py
@@ -0,0 +1,261 @@
+from __future__ import annotations
+
+from urllib.parse import urlencode
+
+from fastapi import APIRouter, File, Form, Request, UploadFile
+from fastapi.responses import RedirectResponse
+from sqlalchemy import select
+
+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.gstr2b_models import AccountingGSTR2BPurchase
+from app.modules.accounting.gstr2b_service import purchases_for_client
+from app.modules.accounting.ledger_learning_service import available_tally_guids
+from app.modules.accounting.purchase_enrichment_models import (
+ AccountingPurchaseEnrichmentBatch,
+ AccountingPurchaseEnrichmentRecord,
+)
+from app.modules.accounting.purchase_enrichment_service import (
+ batches_for_client,
+ import_enrichment,
+ items_for_records,
+ manually_link_record,
+ records_for_batch,
+ rematch_batch,
+ unlink_record,
+)
+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-enrichment", tags=["accounting-purchase-enrichment-ui"])
+
+
+def _redirect(client_id: int, *, batch_id: int | None = None, message: str = "", error: str = ""):
+ params = {"client_id": client_id}
+ if batch_id:
+ params["batch_id"] = batch_id
+ if message:
+ params["message"] = message[:220]
+ if error:
+ params["error"] = error[:220]
+ return RedirectResponse(
+ url="/tools/accounting/purchase-enrichment?" + urlencode(params),
+ status_code=303,
+ )
+
+
+@router.get("")
+def page(
+ request: Request,
+ client_id: int | None = None,
+ batch_id: int | None = None,
+ 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)
+
+ batches = []
+ selected_batch = None
+ records = []
+ items = {}
+ gstr2b_rows = []
+ gstr2b_by_id = {}
+ companies = []
+
+ if selected:
+ batches = batches_for_client(db, scope.tenant_id, selected.id)
+ selected_batch = next((b for b in batches if batch_id and b.id == batch_id), None)
+ if not selected_batch and batches:
+ selected_batch = batches[0]
+ if selected_batch:
+ records = records_for_batch(db, selected_batch.id)
+ items = items_for_records(db, [r.id for r in records])
+ gstr2b_rows = purchases_for_client(db, scope.tenant_id, selected.id, limit=500)
+ gstr2b_by_id = {r.id: r for r in gstr2b_rows}
+ companies = available_tally_guids(db, scope.tenant_id, selected.id)
+
+ return templates.TemplateResponse(
+ "modules/accounting/templates/accounting/purchase_enrichment.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": "E-Invoice / E-Way Bill Enrichment",
+ "clients": clients,
+ "selected_client": selected,
+ "batches": batches,
+ "selected_batch": selected_batch,
+ "records": records,
+ "items_by_record": items,
+ "gstr2b_rows": gstr2b_rows,
+ "gstr2b_by_id": gstr2b_by_id,
+ "tally_companies": companies,
+ "message": message,
+ "error": error,
+ },
+ )
+ finally:
+ db.close()
+
+
+@router.post("/upload")
+async def upload(
+ request: Request,
+ client_id: int = Form(...),
+ source_type: str = Form(...),
+ source_period: str = Form(""),
+ tally_guid: str = Form(""),
+ upload: UploadFile = File(...),
+ 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()
+ content = await upload.read()
+ batch, duplicate_file = import_enrichment(
+ db,
+ tenant_id=scope.tenant_id,
+ client_id=client.id,
+ tally_guid=tally_guid,
+ source_type=source_type,
+ source_period=source_period,
+ filename=upload.filename or "source.json",
+ content=content,
+ user_id=user.id,
+ )
+ if duplicate_file:
+ msg = f"This exact {source_type.replace('_', ' ')} file was already imported as batch #{batch.id}."
+ else:
+ msg = (
+ f"Imported {batch.records_imported} document(s): "
+ f"{batch.records_linked} linked, {batch.records_unmatched} unmatched, "
+ f"{batch.records_ambiguous} ambiguous, {batch.records_duplicate} duplicate(s)."
+ )
+ return _redirect(client.id, batch_id=batch.id, message=msg)
+ except Exception as exc:
+ db.rollback()
+ return _redirect(client_id, error=str(exc))
+ finally:
+ db.close()
+
+
+@router.post("/batch/{batch_id}/rematch")
+def rematch(
+ request: Request,
+ batch_id: int,
+ client_id: int = Form(...),
+ 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()
+ batch = db.execute(select(AccountingPurchaseEnrichmentBatch).where(
+ AccountingPurchaseEnrichmentBatch.id == batch_id,
+ AccountingPurchaseEnrichmentBatch.tenant_id == scope.tenant_id,
+ AccountingPurchaseEnrichmentBatch.client_id == client.id,
+ )).scalar_one_or_none()
+ if not batch:
+ return _redirect(client_id, error="Enrichment batch was not found.")
+ count = rematch_batch(db, batch)
+ return _redirect(client_id, batch_id=batch.id, message=f"Re-matched {count} enrichment document(s).")
+ except Exception as exc:
+ db.rollback()
+ return _redirect(client_id, batch_id=batch_id, error=str(exc))
+ finally:
+ db.close()
+
+
+@router.post("/record/{record_id}/link")
+def manual_link(
+ request: Request,
+ record_id: int,
+ client_id: int = Form(...),
+ purchase_id: int = Form(...),
+ 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()
+ record = db.execute(select(AccountingPurchaseEnrichmentRecord).where(
+ AccountingPurchaseEnrichmentRecord.id == record_id,
+ AccountingPurchaseEnrichmentRecord.tenant_id == scope.tenant_id,
+ AccountingPurchaseEnrichmentRecord.client_id == client.id,
+ )).scalar_one_or_none()
+ purchase = db.execute(select(AccountingGSTR2BPurchase).where(
+ AccountingGSTR2BPurchase.id == purchase_id,
+ AccountingGSTR2BPurchase.tenant_id == scope.tenant_id,
+ AccountingGSTR2BPurchase.client_id == client.id,
+ )).scalar_one_or_none()
+ if not record or not purchase:
+ return _redirect(client_id, error="Record or GSTR-2B purchase was not found.")
+ manually_link_record(db, record, purchase)
+ return _redirect(client_id, batch_id=record.batch_id, message=f"Linked to GSTR-2B invoice {purchase.invoice_number}.")
+ except Exception as exc:
+ db.rollback()
+ return _redirect(client_id, error=str(exc))
+ finally:
+ db.close()
+
+
+@router.post("/record/{record_id}/unlink")
+def manual_unlink(
+ request: Request,
+ record_id: int,
+ client_id: int = Form(...),
+ 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()
+ record = db.execute(select(AccountingPurchaseEnrichmentRecord).where(
+ AccountingPurchaseEnrichmentRecord.id == record_id,
+ AccountingPurchaseEnrichmentRecord.tenant_id == scope.tenant_id,
+ AccountingPurchaseEnrichmentRecord.client_id == client.id,
+ )).scalar_one_or_none()
+ if not record:
+ return _redirect(client_id, error="Enrichment record was not found.")
+ batch_id = record.batch_id
+ unlink_record(db, record)
+ return _redirect(client_id, batch_id=batch_id, message="Enrichment source was unlinked.")
+ except Exception as exc:
+ db.rollback()
+ return _redirect(client_id, error=str(exc))
+ finally:
+ db.close()
diff --git a/app/modules/accounting/templates/accounting/purchase_enrichment.html b/app/modules/accounting/templates/accounting/purchase_enrichment.html
new file mode 100644
index 0000000..a6c3a41
--- /dev/null
+++ b/app/modules/accounting/templates/accounting/purchase_enrichment.html
@@ -0,0 +1,195 @@
+{% extends "ui/templates/base/layout.html" %}
+{% block content %}
+
+
+
+
Tools · Accounting Intelligence
+
E-Invoice / E-Way Bill Enrichment
+
Phase 8 links E-Invoice and E-Way Bill source documents to existing GSTR-2B purchases and enriches classification with item description, HSN, quantity and logistics context. Original source fields remain unchanged and nothing is posted to Tally.
+
+
+
+
+ {% if message %}
{{ message }}
{% endif %}
+ {% if error %}
{{ error }}
{% endif %}
+
+
+
+ {% if selected_client %}
+
+
+
+
Import enrichment source
+
Supported: JSON, XLSX/XLSM and CSV. JSON supports common GST E-Invoice and E-Way Bill field structures, including item lists.
+
+
Enrichment only · No voucher posting
+
+
+
+
+ {% if selected_batch %}
+
+ Source
{{ selected_batch.source_type|replace('_',' ')|title }}
+ Imported
{{ selected_batch.records_imported }}
+ Linked
{{ selected_batch.records_linked }}
+ Unmatched
{{ selected_batch.records_unmatched }}
+ Ambiguous
{{ selected_batch.records_ambiguous }}
+ Duplicates
{{ selected_batch.records_duplicate }}
+
+
+
+
+
+
{{ selected_batch.original_filename }}
+
Automatic matching uses supplier GSTIN + invoice number/date first, with a conservative value/date fallback.
+
+
+
+
+
+
+ {% for record in records %}
+ {% set linked = gstr2b_by_id.get(record.gstr2b_purchase_id) if record.gstr2b_purchase_id else None %}
+
+
+
+
+
{{ record.supplier_name or record.supplier_gstin or 'Supplier' }}
+ {{ record.source_type|replace('_',' ')|title }}
+ {{ record.match_status|title }}
+
+
+
GSTIN: {{ record.supplier_gstin or '-' }}
+
Document: {{ record.document_number }} · {{ record.document_date }}
+
Reference: {{ record.external_reference or '-' }}
+ {% if record.vehicle_number %}
Vehicle: {{ record.vehicle_number }}{% if record.transport_mode %} · {{ record.transport_mode }}{% endif %}
{% endif %}
+
+
+
+
+
+
Taxable
₹{{ '%.2f'|format(record.taxable_value) }}
+
Invoice Value
₹{{ '%.2f'|format(record.invoice_value) }}
+
IGST
₹{{ '%.2f'|format(record.igst) }}
+
CGST + SGST
₹{{ '%.2f'|format(record.cgst + record.sgst) }}
+
+
+
+
+ {% if linked %}
+
+
Linked GSTR-2B purchase · {{ record.match_confidence }}%
+
{{ linked.invoice_number }} · {{ linked.invoice_date }}
+
₹{{ '%.2f'|format(linked.taxable_value) }} taxable · {{ record.match_method|replace('_',' ') }}
+
+
+ {% else %}
+
+ {% endif %}
+
+
+
+ {% if items_by_record.get(record.id) %}
+
+
+ | Item | HSN | Qty | Unit | Taxable |
+
+ {% for item in items_by_record.get(record.id) %}
+
+ {{ item.product_name or item.description_text or '-' }} {% if item.description_text and item.description_text != item.product_name %}{{ item.description_text }} {% endif %} |
+ {{ item.hsn_code or '-' }} |
+ {{ item.quantity }} |
+ {{ item.unit or '-' }} |
+ ₹{{ '%.2f'|format(item.taxable_value) }} |
+
+ {% endfor %}
+
+
+
+ {% endif %}
+
+ {% else %}
+ No enrichment records in this batch.
+ {% endfor %}
+
+ {% endif %}
+ {% endif %}
+
+{% endblock %}
diff --git a/app/modules/accounting/templates/accounting/tally.html b/app/modules/accounting/templates/accounting/tally.html
index 27a8f07..818237a 100644
--- a/app/modules/accounting/templates/accounting/tally.html
+++ b/app/modules/accounting/templates/accounting/tally.html
@@ -17,6 +17,7 @@
{% if selected_client %}Historical Learning{% endif %}
{% if selected_client %}Ledger Learning{% endif %}
{% if selected_client %}GSTR-2B Intelligence{% endif %}
+ {% if selected_client %}E-Invoice / E-Way Bill{% endif %}
{% if selected_client %}Depreciation (IT){% endif %}
Refresh Tally Companies
diff --git a/app/ui/app.py b/app/ui/app.py
index a64160e..ac3534a 100644
--- a/app/ui/app.py
+++ b/app/ui/app.py
@@ -41,6 +41,7 @@ from app.modules.accounting.taxonomy_ui import router as accounting_taxonomy_ui_
from app.modules.accounting.historical_learning_ui import router as accounting_historical_learning_ui_router
from app.modules.accounting.ledger_learning_ui import router as accounting_ledger_learning_ui_router
from app.modules.accounting.gstr2b_ui import router as accounting_gstr2b_ui_router
+from app.modules.accounting.purchase_enrichment_ui import router as accounting_purchase_enrichment_ui_router
from app.modules.registrations.ui import router as registrations_ui_router
from app.modules.credential_vault.ui import router as credential_vault_ui_router
from app.modules.client_identity.ui import router as client_identity_ui_router
@@ -68,6 +69,7 @@ def mount_ui(app: FastAPI) -> None:
app.include_router(accounting_historical_learning_ui_router)
app.include_router(accounting_ledger_learning_ui_router)
app.include_router(accounting_gstr2b_ui_router)
+ app.include_router(accounting_purchase_enrichment_ui_router)
app.include_router(work_tracker_ui_router)
app.include_router(billing_ui_router)
app.include_router(platform_billing_ui_router)