From 7283f45f43f758211b0e76784963f3b70504de75 Mon Sep 17 00:00:00 2001 From: A R R R Associates Date: Sat, 22 Aug 2026 14:31:24 +0530 Subject: [PATCH] Add Phase 7 GSTR-2B purchase intelligence --- ...822_gstr2b_purchase_intelligence_phase7.py | 99 ++++++ app/modules/accounting/gstr2b_models.py | 115 +++++++ app/modules/accounting/gstr2b_parser.py | 307 ++++++++++++++++++ app/modules/accounting/gstr2b_service.py | 226 +++++++++++++ app/modules/accounting/gstr2b_ui.py | 247 ++++++++++++++ .../templates/accounting/gstr2b.html | 202 ++++++++++++ .../templates/accounting/tally.html | 1 + app/ui/app.py | 2 + 8 files changed, 1199 insertions(+) create mode 100644 alembic/versions/20260822_gstr2b_purchase_intelligence_phase7.py create mode 100644 app/modules/accounting/gstr2b_models.py create mode 100644 app/modules/accounting/gstr2b_parser.py create mode 100644 app/modules/accounting/gstr2b_service.py create mode 100644 app/modules/accounting/gstr2b_ui.py create mode 100644 app/modules/accounting/templates/accounting/gstr2b.html diff --git a/alembic/versions/20260822_gstr2b_purchase_intelligence_phase7.py b/alembic/versions/20260822_gstr2b_purchase_intelligence_phase7.py new file mode 100644 index 0000000..ee9803d --- /dev/null +++ b/alembic/versions/20260822_gstr2b_purchase_intelligence_phase7.py @@ -0,0 +1,99 @@ +"""Phase 7 GSTR-2B purchase intelligence. + +Revision ID: 20260822_gstr2b_purchase_p7 +Revises: 20260822_ledger_learning_p6 +""" +from alembic import op +import sqlalchemy as sa + +revision = "20260822_gstr2b_purchase_p7" +down_revision = "20260822_ledger_learning_p6" +branch_labels = None +depends_on = None + + +def upgrade(): + op.create_table( + "accounting_gstr2b_import_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("original_filename", sa.String(260), nullable=False, server_default=""), + sa.Column("file_sha256", sa.String(64), nullable=False), + sa.Column("return_period", sa.String(20), nullable=False, server_default=""), + sa.Column("source_kind", sa.String(30), nullable=False, server_default="gstr2b_upload"), + sa.Column("status", sa.String(30), nullable=False, server_default="imported"), + sa.Column("rows_read", sa.Integer(), nullable=False, server_default="0"), + sa.Column("rows_imported", sa.Integer(), nullable=False, server_default="0"), + sa.Column("rows_skipped_duplicate", sa.Integer(), nullable=False, server_default="0"), + sa.Column("rows_skipped_invalid", sa.Integer(), nullable=False, server_default="0"), + sa.Column("analyzed_rows", sa.Integer(), nullable=False, server_default="0"), + sa.Column("reviewed_rows", sa.Integer(), nullable=False, server_default="0"), + sa.Column("error_message", sa.Text(), nullable=True), + 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("analyzed_at_utc", sa.DateTime(timezone=True), nullable=True), + sa.Column("completed_at_utc", sa.DateTime(timezone=True), nullable=True), + ) + for col in ("tenant_id", "client_id", "tally_guid", "file_sha256", "return_period", "source_kind", "status", "created_at_utc"): + op.create_index(f"ix_accounting_gstr2b_import_batches_{col}", "accounting_gstr2b_import_batches", [col]) + + op.create_table( + "accounting_gstr2b_purchases", + 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_gstr2b_import_batches.id", ondelete="CASCADE"), nullable=False), + sa.Column("tally_guid", sa.String(120), nullable=False, server_default=""), + sa.Column("return_period", sa.String(20), 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("invoice_number", sa.String(160), nullable=False, server_default=""), + sa.Column("invoice_date", sa.String(20), nullable=False, server_default=""), + sa.Column("document_type", sa.String(40), nullable=False, server_default="invoice"), + sa.Column("invoice_type", sa.String(80), nullable=False, server_default=""), + sa.Column("place_of_supply", sa.String(120), nullable=False, server_default=""), + sa.Column("reverse_charge", sa.String(20), nullable=False, server_default=""), + sa.Column("itc_availability", sa.String(80), nullable=False, server_default=""), + 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("hsn_code", sa.String(20), nullable=False, server_default=""), + sa.Column("description_text", sa.Text(), nullable=True), + sa.Column("source_sheet", sa.String(160), nullable=False, server_default=""), + sa.Column("source_row_number", sa.Integer(), nullable=False, server_default="0"), + sa.Column("source_row_hash", sa.String(64), nullable=False, server_default=""), + sa.Column("review_status", sa.String(30), nullable=False, server_default="pending_analysis"), + sa.Column("suggested_nature_id", sa.Integer(), sa.ForeignKey("accounting_natures.id", ondelete="SET NULL"), nullable=True), + sa.Column("suggested_ledger_name", sa.String(240), nullable=False, server_default=""), + sa.Column("suggested_confidence", sa.Integer(), nullable=False, server_default="0"), + sa.Column("suggestion_explanation_json", sa.Text(), nullable=True), + sa.Column("final_nature_id", sa.Integer(), sa.ForeignKey("accounting_natures.id", ondelete="SET NULL"), nullable=True), + sa.Column("final_ledger_name", sa.String(240), nullable=False, server_default=""), + sa.Column("reviewed_by_user_id", sa.Integer(), sa.ForeignKey("users.id", ondelete="SET NULL"), nullable=True), + sa.Column("reviewed_at_utc", sa.DateTime(timezone=True), nullable=True), + sa.Column("posting_status", sa.String(30), nullable=False, server_default="not_enabled"), + sa.Column("is_duplicate_source", 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", "supplier_gstin", "invoice_number", "invoice_date", "document_type", + name="uq_accounting_gstr2b_purchase_document", + ), + ) + for col in ( + "tenant_id", "client_id", "batch_id", "tally_guid", "return_period", "supplier_gstin", + "supplier_name", "invoice_number", "invoice_date", "document_type", "hsn_code", + "source_row_hash", "review_status", "suggested_nature_id", "final_nature_id", + "posting_status", "created_at_utc", + ): + op.create_index(f"ix_accounting_gstr2b_purchases_{col}", "accounting_gstr2b_purchases", [col]) + + +def downgrade(): + op.drop_table("accounting_gstr2b_purchases") + op.drop_table("accounting_gstr2b_import_batches") diff --git a/app/modules/accounting/gstr2b_models.py b/app/modules/accounting/gstr2b_models.py new file mode 100644 index 0000000..d6fd94f --- /dev/null +++ b/app/modules/accounting/gstr2b_models.py @@ -0,0 +1,115 @@ +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 AccountingGSTR2BImportBatch(CommonBase): + __tablename__ = "accounting_gstr2b_import_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) + + original_filename: Mapped[str] = mapped_column(String(260), nullable=False, default="") + file_sha256: Mapped[str] = mapped_column(String(64), nullable=False, index=True) + return_period: Mapped[str] = mapped_column(String(20), nullable=False, default="", index=True) + source_kind: Mapped[str] = mapped_column(String(30), nullable=False, default="gstr2b_upload", index=True) + + status: Mapped[str] = mapped_column(String(30), nullable=False, default="imported", index=True) + rows_read: Mapped[int] = mapped_column(Integer, nullable=False, default=0) + rows_imported: Mapped[int] = mapped_column(Integer, nullable=False, default=0) + rows_skipped_duplicate: Mapped[int] = mapped_column(Integer, nullable=False, default=0) + rows_skipped_invalid: Mapped[int] = mapped_column(Integer, nullable=False, default=0) + analyzed_rows: Mapped[int] = mapped_column(Integer, nullable=False, default=0) + reviewed_rows: Mapped[int] = mapped_column(Integer, nullable=False, default=0) + error_message: Mapped[str | None] = mapped_column(Text, nullable=True) + + 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 + ) + analyzed_at_utc: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + completed_at_utc: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + + +class AccountingGSTR2BPurchase(CommonBase): + __tablename__ = "accounting_gstr2b_purchases" + __table_args__ = ( + UniqueConstraint( + "tenant_id", + "client_id", + "supplier_gstin", + "invoice_number", + "invoice_date", + "document_type", + name="uq_accounting_gstr2b_purchase_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_gstr2b_import_batches.id", ondelete="CASCADE"), nullable=False, index=True + ) + tally_guid: Mapped[str] = mapped_column(String(120), nullable=False, default="", index=True) + + return_period: Mapped[str] = mapped_column(String(20), 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) + + invoice_number: Mapped[str] = mapped_column(String(160), nullable=False, default="", index=True) + invoice_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) + invoice_type: Mapped[str] = mapped_column(String(80), nullable=False, default="") + + place_of_supply: Mapped[str] = mapped_column(String(120), nullable=False, default="") + reverse_charge: Mapped[str] = mapped_column(String(20), nullable=False, default="") + itc_availability: Mapped[str] = mapped_column(String(80), nullable=False, default="") + + 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) + + hsn_code: Mapped[str] = mapped_column(String(20), nullable=False, default="", index=True) + description_text: Mapped[str | None] = mapped_column(Text, nullable=True) + source_sheet: Mapped[str] = mapped_column(String(160), nullable=False, default="") + source_row_number: Mapped[int] = mapped_column(Integer, nullable=False, default=0) + source_row_hash: Mapped[str] = mapped_column(String(64), nullable=False, default="", index=True) + + review_status: Mapped[str] = mapped_column(String(30), nullable=False, default="pending_analysis", index=True) + suggested_nature_id: Mapped[int | None] = mapped_column( + ForeignKey("accounting_natures.id", ondelete="SET NULL"), nullable=True, index=True + ) + suggested_ledger_name: Mapped[str] = mapped_column(String(240), nullable=False, default="") + suggested_confidence: Mapped[int] = mapped_column(Integer, nullable=False, default=0) + suggestion_explanation_json: Mapped[str | None] = mapped_column(Text, nullable=True) + + final_nature_id: Mapped[int | None] = mapped_column( + ForeignKey("accounting_natures.id", ondelete="SET NULL"), nullable=True, index=True + ) + final_ledger_name: Mapped[str] = mapped_column(String(240), nullable=False, default="") + reviewed_by_user_id: Mapped[int | None] = mapped_column(ForeignKey("users.id", ondelete="SET NULL"), nullable=True) + reviewed_at_utc: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + + posting_status: Mapped[str] = mapped_column(String(30), nullable=False, default="not_enabled", index=True) + is_duplicate_source: 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, + ) diff --git a/app/modules/accounting/gstr2b_parser.py b/app/modules/accounting/gstr2b_parser.py new file mode 100644 index 0000000..135735b --- /dev/null +++ b/app/modules/accounting/gstr2b_parser.py @@ -0,0 +1,307 @@ +from __future__ import annotations + +import csv +import hashlib +import io +import re +from datetime import date, datetime +from pathlib import Path +from typing import Iterable + +from openpyxl import load_workbook + + +MAX_UPLOAD_BYTES = 25 * 1024 * 1024 + +HEADER_ALIASES = { + "supplier_gstin": { + "gstin of supplier", "supplier gstin", "gstin", "gstin/uin of supplier", + "gstin / uin of supplier", "ctin", + }, + "supplier_name": { + "trade/legal name", "trade legal name", "supplier name", "legal name", + "trade name", "name of supplier", "supplier trade name", + }, + "invoice_number": { + "invoice number", "invoice no", "invoice no.", "document number", + "document no", "doc no", "inum", + }, + "invoice_date": { + "invoice date", "document date", "doc date", "invoice dt", "idt", + }, + "invoice_type": { + "invoice type", "document type", "type", "inv type", + }, + "invoice_value": { + "invoice value", "document value", "invoice amount", "total invoice value", + "val", + }, + "place_of_supply": { + "place of supply", "pos", "place of supply (state/ut)", "place of supply state/ut", + }, + "reverse_charge": { + "supply attract reverse charge", "reverse charge", "rcm", "reverse charge applicable", + }, + "taxable_value": { + "taxable value", "taxable amount", "txval", + }, + "igst": { + "integrated tax", "integrated tax amount", "igst", "igst amount", "iamt", + }, + "cgst": { + "central tax", "central tax amount", "cgst", "cgst amount", "camt", + }, + "sgst": { + "state/ut tax", "state / ut tax", "state tax", "sgst", "utgst", + "state/ut tax amount", "sgst amount", "samt", + }, + "cess": { + "cess", "cess amount", "csamt", + }, + "itc_availability": { + "itc availability", "itc available", "itc available for the tax period", + "itc eligibility", "eligibility for itc", + }, + "hsn_code": { + "hsn", "hsn code", "hsn/sac", "hsn / sac", "hsn/sac code", + }, + "description_text": { + "description", "item description", "product description", "goods/service description", + "description of goods/services", "description of goods / services", + }, +} + + +def _norm_header(value) -> str: + text = str(value or "").strip().lower() + text = text.replace("\n", " ").replace("\r", " ") + text = re.sub(r"[_\-]+", " ", text) + text = re.sub(r"[^\w/(). ]+", " ", text) + text = re.sub(r"\s+", " ", text).strip() + return text + + +ALIAS_LOOKUP = { + alias: key + for key, aliases in HEADER_ALIASES.items() + for alias in aliases +} + + +def _canonical_header(value) -> str | None: + norm = _norm_header(value) + if norm in ALIAS_LOOKUP: + return ALIAS_LOOKUP[norm] + # Conservative contains matching only for longer aliases. + for alias, key in ALIAS_LOOKUP.items(): + if len(alias) >= 10 and (norm.startswith(alias) or alias in norm): + return key + return None + + +def _text(value) -> str: + if value is None: + return "" + if isinstance(value, float) and value.is_integer(): + return str(int(value)) + return str(value).strip() + + +def _money(value) -> float: + if value in (None, ""): + return 0.0 + if isinstance(value, (int, float)): + return float(value) + text = str(value).strip().replace(",", "").replace("₹", "") + text = re.sub(r"^\((.*)\)$", r"-\1", text) + try: + return float(text) + except Exception: + return 0.0 + + +def _date_text(value) -> str: + if value in (None, ""): + return "" + if isinstance(value, datetime): + return value.date().isoformat() + if isinstance(value, date): + return value.isoformat() + text = str(value).strip() + candidates = ( + "%d-%m-%Y", "%d/%m/%Y", "%Y-%m-%d", "%d-%b-%Y", "%d %b %Y", + "%d.%m.%Y", "%m/%d/%Y", + ) + for fmt in candidates: + try: + return datetime.strptime(text, fmt).date().isoformat() + except Exception: + pass + # Preserve original if portal format is unfamiliar; duplicate identity remains stable. + return text[:20] + + +def _gstin(value) -> str: + return re.sub(r"[^A-Z0-9]", "", str(value or "").upper())[:15] + + +def _hsn(value) -> str: + text = re.sub(r"\D", "", str(value or "")) + return text[:8] + + +def _document_type(invoice_type: str, invoice_number: str) -> str: + text = f"{invoice_type} {invoice_number}".upper() + if "CREDIT" in text or "CR NOTE" in text or "CREDIT NOTE" in text: + return "credit_note" + if "DEBIT" in text or "DR NOTE" in text or "DEBIT NOTE" in text: + return "debit_note" + return "invoice" + + +def _row_hash(row: dict) -> str: + raw = "|".join([ + row.get("supplier_gstin", ""), + row.get("supplier_name", ""), + row.get("invoice_number", ""), + row.get("invoice_date", ""), + row.get("document_type", ""), + f"{row.get('taxable_value', 0):.2f}", + f"{row.get('igst', 0):.2f}", + f"{row.get('cgst', 0):.2f}", + f"{row.get('sgst', 0):.2f}", + row.get("hsn_code", ""), + row.get("description_text", ""), + ]) + return hashlib.sha256(raw.encode("utf-8", "ignore")).hexdigest() + + +def _find_header(rows: list[list], max_scan: int = 40): + best = None + for index, values in enumerate(rows[:max_scan]): + mapping = {} + for col, value in enumerate(values): + key = _canonical_header(value) + if key and key not in mapping: + mapping[key] = col + score = sum(1 for key in ("supplier_gstin", "invoice_number", "invoice_date", "taxable_value") if key in mapping) + if score >= 3 and ("supplier_gstin" in mapping or "supplier_name" in mapping): + if best is None or score > best[0]: + best = (score, index, mapping) + return best + + +def _normalized_record(values: list, mapping: dict[str, int], *, sheet: str, row_number: int): + def get(key): + idx = mapping.get(key) + return values[idx] if idx is not None and idx < len(values) else None + + supplier_gstin = _gstin(get("supplier_gstin")) + supplier_name = _text(get("supplier_name")) + invoice_number = _text(get("invoice_number")) + invoice_date = _date_text(get("invoice_date")) + invoice_type = _text(get("invoice_type")) + + if not supplier_gstin and not supplier_name: + return None + if not invoice_number or not invoice_date: + return None + + row = { + "supplier_gstin": supplier_gstin, + "supplier_name": supplier_name, + "invoice_number": invoice_number[:160], + "invoice_date": invoice_date, + "document_type": _document_type(invoice_type, invoice_number), + "invoice_type": invoice_type[:80], + "invoice_value": _money(get("invoice_value")), + "place_of_supply": _text(get("place_of_supply"))[:120], + "reverse_charge": _text(get("reverse_charge"))[:20], + "taxable_value": _money(get("taxable_value")), + "igst": _money(get("igst")), + "cgst": _money(get("cgst")), + "sgst": _money(get("sgst")), + "cess": _money(get("cess")), + "itc_availability": _text(get("itc_availability"))[:80], + "hsn_code": _hsn(get("hsn_code")), + "description_text": _text(get("description_text"))[:4000], + "source_sheet": sheet[:160], + "source_row_number": int(row_number), + } + if not row["invoice_value"]: + row["invoice_value"] = ( + row["taxable_value"] + row["igst"] + row["cgst"] + row["sgst"] + row["cess"] + ) + row["source_row_hash"] = _row_hash(row) + return row + + +def _iter_xlsx(content: bytes): + workbook = load_workbook(io.BytesIO(content), read_only=True, data_only=True) + try: + for ws in workbook.worksheets: + # Summary sheets usually won't pass header detection. + rows = [list(row) for row in ws.iter_rows(values_only=True)] + if not rows: + continue + found = _find_header(rows) + if not found: + continue + _, header_index, mapping = found + for idx, values in enumerate(rows[header_index + 1:], start=header_index + 2): + record = _normalized_record(values, mapping, sheet=ws.title, row_number=idx) + if record: + yield record + finally: + workbook.close() + + +def _decode_csv(content: bytes) -> str: + for encoding in ("utf-8-sig", "utf-8", "cp1252", "latin-1"): + try: + return content.decode(encoding) + except UnicodeDecodeError: + pass + return content.decode("utf-8", "replace") + + +def _iter_csv(content: bytes): + text = _decode_csv(content) + sample = text[:8192] + try: + dialect = csv.Sniffer().sniff(sample, delimiters=",;\t|") + except Exception: + dialect = csv.excel + rows = [list(row) for row in csv.reader(io.StringIO(text), dialect)] + found = _find_header(rows) + if not found: + return + _, header_index, mapping = found + for idx, values in enumerate(rows[header_index + 1:], start=header_index + 2): + record = _normalized_record(values, mapping, sheet="CSV", row_number=idx) + if record: + yield record + + +def parse_gstr2b(content: bytes, filename: str): + if not content: + raise ValueError("Uploaded GSTR-2B file is empty.") + if len(content) > MAX_UPLOAD_BYTES: + raise ValueError("GSTR-2B upload exceeds the 25 MB limit.") + suffix = Path(filename or "").suffix.lower() + if suffix in {".xlsx", ".xlsm"}: + records = list(_iter_xlsx(content)) + elif suffix in {".csv", ".txt"}: + records = list(_iter_csv(content)) + else: + raise ValueError("Upload an .xlsx, .xlsm or .csv GSTR-2B file.") + if not records: + raise ValueError( + "No GSTR-2B invoice rows were detected. The file must contain supplier GSTIN/name, " + "invoice number, invoice date and taxable value columns." + ) + return records + + +def file_sha256(content: bytes) -> str: + return hashlib.sha256(content).hexdigest() diff --git a/app/modules/accounting/gstr2b_service.py b/app/modules/accounting/gstr2b_service.py new file mode 100644 index 0000000..3e378c1 --- /dev/null +++ b/app/modules/accounting/gstr2b_service.py @@ -0,0 +1,226 @@ +from __future__ import annotations + +import json +from datetime import datetime, timezone + +from sqlalchemy import func, select +from sqlalchemy.exc import IntegrityError + +from app.modules.accounting.gstr2b_models import AccountingGSTR2BImportBatch, AccountingGSTR2BPurchase +from app.modules.accounting.gstr2b_parser import file_sha256, parse_gstr2b +from app.modules.accounting.ledger_learning_service import rank_suggestions, record_review +from app.modules.accounting.taxonomy_models import AccountingNature + + +def _utcnow(): + return datetime.now(timezone.utc) + + +def import_gstr2b( + db, *, + tenant_id: int, + client_id: int, + tally_guid: str, + return_period: str, + filename: str, + content: bytes, + user_id: int, +): + digest = file_sha256(content) + existing_batch = db.execute(select(AccountingGSTR2BImportBatch).where( + AccountingGSTR2BImportBatch.tenant_id == tenant_id, + AccountingGSTR2BImportBatch.client_id == client_id, + AccountingGSTR2BImportBatch.file_sha256 == digest, + )).scalar_one_or_none() + if existing_batch: + return existing_batch, True + + parsed = parse_gstr2b(content, filename) + batch = AccountingGSTR2BImportBatch( + tenant_id=tenant_id, + client_id=client_id, + tally_guid=(tally_guid or "").strip(), + original_filename=(filename or "gstr2b.xlsx")[:260], + file_sha256=digest, + return_period=(return_period or "").strip()[:20], + source_kind="gstr2b_upload", + status="importing", + rows_read=len(parsed), + imported_by_user_id=user_id, + ) + db.add(batch) + db.flush() + + imported = duplicate = invalid = 0 + for row in parsed: + if not row.get("invoice_number") or not row.get("invoice_date"): + invalid += 1 + continue + + exists = db.execute(select(AccountingGSTR2BPurchase.id).where( + AccountingGSTR2BPurchase.tenant_id == tenant_id, + AccountingGSTR2BPurchase.client_id == client_id, + AccountingGSTR2BPurchase.supplier_gstin == row.get("supplier_gstin", ""), + AccountingGSTR2BPurchase.invoice_number == row["invoice_number"], + AccountingGSTR2BPurchase.invoice_date == row["invoice_date"], + AccountingGSTR2BPurchase.document_type == row["document_type"], + )).scalar_one_or_none() + if exists: + duplicate += 1 + continue + + db.add(AccountingGSTR2BPurchase( + tenant_id=tenant_id, + client_id=client_id, + batch_id=batch.id, + tally_guid=(tally_guid or "").strip(), + return_period=(return_period or "").strip()[:20], + **row, + )) + imported += 1 + + batch.rows_imported = imported + batch.rows_skipped_duplicate = duplicate + batch.rows_skipped_invalid = invalid + batch.status = "imported" + batch.completed_at_utc = _utcnow() + db.commit() + db.refresh(batch) + return batch, False + + +def batches_for_client(db, tenant_id: int, client_id: int, limit: int = 20): + return list(db.execute(select(AccountingGSTR2BImportBatch).where( + AccountingGSTR2BImportBatch.tenant_id == tenant_id, + AccountingGSTR2BImportBatch.client_id == client_id, + ).order_by(AccountingGSTR2BImportBatch.id.desc()).limit(limit)).scalars().all()) + + +def purchases_for_client(db, tenant_id: int, client_id: int, *, batch_id: int | None = None, limit: int = 300): + stmt = select(AccountingGSTR2BPurchase).where( + AccountingGSTR2BPurchase.tenant_id == tenant_id, + AccountingGSTR2BPurchase.client_id == client_id, + ) + if batch_id: + stmt = stmt.where(AccountingGSTR2BPurchase.batch_id == batch_id) + return list(db.execute(stmt.order_by( + AccountingGSTR2BPurchase.invoice_date.desc(), + AccountingGSTR2BPurchase.id.desc(), + ).limit(limit)).scalars().all()) + + +def nature_lookup(db, ids): + ids = {int(x) for x in ids if x} + if not ids: + return {} + return {row.id: row for row in db.execute( + select(AccountingNature).where(AccountingNature.id.in_(ids)) + ).scalars().all()} + + +def analyze_purchase(db, row: AccountingGSTR2BPurchase): + description_parts = [ + row.description_text or "", + f"Invoice type {row.invoice_type}" if row.invoice_type else "", + f"Place of supply {row.place_of_supply}" if row.place_of_supply else "", + ] + suggestions = rank_suggestions( + db, + tenant_id=row.tenant_id, + client_id=row.client_id, + tally_guid=row.tally_guid, + supplier_name=row.supplier_name, + supplier_gstin=row.supplier_gstin, + hsn_code=row.hsn_code, + description=" | ".join(x for x in description_parts if x), + amount=row.taxable_value, + ) + if suggestions: + top = suggestions[0] + row.suggested_nature_id = top["nature"].id + row.suggested_ledger_name = top.get("suggested_ledger") or "" + row.suggested_confidence = int(top.get("confidence") or 0) + row.suggestion_explanation_json = json.dumps(top.get("reasons") or [], ensure_ascii=False) + row.review_status = "suggested" + else: + row.suggested_nature_id = None + row.suggested_ledger_name = "" + row.suggested_confidence = 0 + row.suggestion_explanation_json = json.dumps( + ["No reliable classification evidence is available yet."], ensure_ascii=False + ) + row.review_status = "review_required" + return suggestions + + +def analyze_batch(db, batch: AccountingGSTR2BImportBatch): + rows = list(db.execute(select(AccountingGSTR2BPurchase).where( + AccountingGSTR2BPurchase.batch_id == batch.id + )).scalars().all()) + analyzed = 0 + for row in rows: + analyze_purchase(db, row) + analyzed += 1 + batch.analyzed_rows = analyzed + batch.analyzed_at_utc = _utcnow() + batch.status = "analyzed" + db.commit() + return analyzed + + +def review_purchase( + db, *, + row: AccountingGSTR2BPurchase, + final_nature_id: int, + final_ledger_name: str, + user_id: int, +): + try: + explanation = json.loads(row.suggestion_explanation_json or "[]") + except Exception: + explanation = [] + + record_review( + db, + tenant_id=row.tenant_id, + client_id=row.client_id, + tally_guid=row.tally_guid, + supplier_name=row.supplier_name, + supplier_gstin=row.supplier_gstin, + hsn_code=row.hsn_code, + description=row.description_text or "", + amount=row.taxable_value, + suggested_nature_id=row.suggested_nature_id, + suggested_ledger_name=row.suggested_ledger_name, + suggested_confidence=row.suggested_confidence, + final_nature_id=final_nature_id, + final_ledger_name=final_ledger_name, + user_id=user_id, + explanation=explanation, + ) + + # record_review commits its learning event; refresh the purchase in this session. + row.final_nature_id = final_nature_id + row.final_ledger_name = (final_ledger_name or "").strip() + row.review_status = "reviewed" + row.reviewed_by_user_id = user_id + row.reviewed_at_utc = _utcnow() + db.add(row) + + batch = db.get(AccountingGSTR2BImportBatch, row.batch_id) + if batch: + reviewed_count = db.execute(select(func.count(AccountingGSTR2BPurchase.id)).where( + AccountingGSTR2BPurchase.batch_id == batch.id, + AccountingGSTR2BPurchase.review_status == "reviewed", + )).scalar_one() + # Include the current row if the database count was evaluated before flush. + if row.review_status == "reviewed": + db.flush() + reviewed_count = db.execute(select(func.count(AccountingGSTR2BPurchase.id)).where( + AccountingGSTR2BPurchase.batch_id == batch.id, + AccountingGSTR2BPurchase.review_status == "reviewed", + )).scalar_one() + batch.reviewed_rows = int(reviewed_count or 0) + db.commit() + db.refresh(row) + return row diff --git a/app/modules/accounting/gstr2b_ui.py b/app/modules/accounting/gstr2b_ui.py new file mode 100644 index 0000000..e2b63bb --- /dev/null +++ b/app/modules/accounting/gstr2b_ui.py @@ -0,0 +1,247 @@ +from __future__ import annotations + +import json +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 AccountingGSTR2BImportBatch, AccountingGSTR2BPurchase +from app.modules.accounting.gstr2b_service import ( + analyze_batch, + batches_for_client, + import_gstr2b, + nature_lookup, + purchases_for_client, + review_purchase, +) +from app.modules.accounting.historical_learning_service import active_natures, ledger_mappings +from app.modules.accounting.ledger_learning_service import available_tally_guids +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/gstr2b", tags=["accounting-gstr2b-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 + if error: + params["error"] = error[:180] + return RedirectResponse(url="/tools/accounting/gstr2b?" + urlencode(params), status_code=303) + + +@router.get("") +def gstr2b_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 = [] + purchases = [] + mappings = [] + companies = [] + natures = [] + nature_by_id = {} + selected_batch = None + + 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] + purchases = purchases_for_client( + db, scope.tenant_id, selected.id, + batch_id=selected_batch.id if selected_batch else None + ) + mappings = ledger_mappings( + db, scope.tenant_id, selected.id, + selected_batch.tally_guid if selected_batch else "" + ) + companies = available_tally_guids(db, scope.tenant_id, selected.id) + natures = active_natures(db, scope.tenant_id) + ids = [] + for row in purchases: + ids.extend([row.suggested_nature_id, row.final_nature_id]) + nature_by_id = nature_lookup(db, ids) + + explanations = {} + for row in purchases: + try: + explanations[row.id] = json.loads(row.suggestion_explanation_json or "[]") + except Exception: + explanations[row.id] = [] + + return templates.TemplateResponse("modules/accounting/templates/accounting/gstr2b.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": "GSTR-2B Purchase Intelligence", + "clients": clients, + "selected_client": selected, + "batches": batches, + "selected_batch": selected_batch, + "purchases": purchases, + "mappings": mappings, + "tally_companies": companies, + "natures": natures, + "nature_by_id": nature_by_id, + "explanations": explanations, + "message": message, + "error": error, + }) + finally: + db.close() + + +@router.post("/upload") +async def upload_gstr2b( + request: Request, + client_id: int = Form(...), + tally_guid: str = Form(""), + return_period: 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_gstr2b( + db, + tenant_id=scope.tenant_id, + client_id=client.id, + tally_guid=tally_guid, + return_period=return_period, + filename=upload.filename or "gstr2b.xlsx", + content=content, + user_id=user.id, + ) + if duplicate_file: + msg = f"This exact GSTR-2B file was already imported as batch #{batch.id}." + else: + msg = ( + f"Imported {batch.rows_imported} purchase document(s); " + f"{batch.rows_skipped_duplicate} duplicate(s) and " + f"{batch.rows_skipped_invalid} invalid row(s) skipped." + ) + 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}/analyze") +def analyze( + 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(AccountingGSTR2BImportBatch).where( + AccountingGSTR2BImportBatch.id == batch_id, + AccountingGSTR2BImportBatch.tenant_id == scope.tenant_id, + AccountingGSTR2BImportBatch.client_id == client.id, + )).scalar_one_or_none() + if not batch: + return _redirect(client_id, error="GSTR-2B import batch was not found.") + count = analyze_batch(db, batch) + return _redirect(client_id, batch_id=batch.id, message=f"Analyzed {count} purchase document(s).") + except Exception as exc: + db.rollback() + return _redirect(client_id, batch_id=batch_id, error=str(exc)) + finally: + db.close() + + +@router.post("/purchase/{purchase_id}/review") +def review( + request: Request, + purchase_id: int, + client_id: int = Form(...), + final_nature_id: int = Form(...), + final_ledger_name: str = 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() + row = 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 row: + return _redirect(client_id, error="GSTR-2B purchase record was not found.") + + # If a ledger is selected, ensure it is one of the Phase 5 mappings for this nature/company. + if final_ledger_name.strip(): + allowed = ledger_mappings(db, scope.tenant_id, client.id, row.tally_guid) + valid = any( + m.ledger_name == final_ledger_name.strip() and int(m.nature_id) == int(final_nature_id) + for m in allowed + ) + if not valid: + return _redirect( + client_id, batch_id=row.batch_id, + error="Selected Tally ledger is not mapped to the selected accounting nature for this client/company." + ) + + review_purchase( + db, + row=row, + final_nature_id=final_nature_id, + final_ledger_name=final_ledger_name, + user_id=user.id, + ) + return _redirect(client_id, batch_id=row.batch_id, message=f"Invoice {row.invoice_number} reviewed and learned.") + except Exception as exc: + db.rollback() + return _redirect(client_id, error=str(exc)) + finally: + db.close() diff --git a/app/modules/accounting/templates/accounting/gstr2b.html b/app/modules/accounting/templates/accounting/gstr2b.html new file mode 100644 index 0000000..60547f4 --- /dev/null +++ b/app/modules/accounting/templates/accounting/gstr2b.html @@ -0,0 +1,202 @@ +{% extends "ui/templates/base/layout.html" %} +{% block content %} +
+
+
+

Tools · Accounting Intelligence

+

GSTR-2B Purchase Intelligence

+

Import GSTR-2B purchase documents, classify them through the Phase 6 learning engine, and confirm the accounting nature and mapped Tally ledger. Phase 7 does not post vouchers to Tally.

+
+ +
+ + {% if message %} +
{{ message }}
+ {% endif %} + {% if error %} +
{{ error }}
+ {% endif %} + +
+
+ + {% if selected_client %} + + {% endif %} +
+ +
+
+
+ + {% if selected_client %} +
+
+
+

Import GSTR-2B

+

Supported: .xlsx, .xlsm and .csv. Summary sheets are ignored automatically; invoice sheets are detected by column headers.

+
+ Suggestion & review only · No Tally posting +
+
+ + + + + +
+ +
+
+
+ + {% if selected_batch %} +
+
Batch
#{{ selected_batch.id }}
+
Rows Read
{{ selected_batch.rows_read }}
+
Imported
{{ selected_batch.rows_imported }}
+
Duplicates
{{ selected_batch.rows_skipped_duplicate }}
+
Analyzed
{{ selected_batch.analyzed_rows }}
+
Reviewed
{{ selected_batch.reviewed_rows }}
+
+ +
+
+
+

{{ selected_batch.original_filename }}

+

Return period {{ selected_batch.return_period or '-' }} · Status {{ selected_batch.status|replace('_',' ')|title }}

+
+
+ + + +
+
+
+ +
+ {% for row in purchases %} + {% set suggested = nature_by_id.get(row.suggested_nature_id) if row.suggested_nature_id else None %} + {% set final = nature_by_id.get(row.final_nature_id) if row.final_nature_id else None %} +
+
+
+
+

{{ row.supplier_name or row.supplier_gstin or 'Supplier' }}

+ {{ row.document_type|replace('_',' ')|title }} + {% if row.review_status == 'reviewed' %} + Reviewed + {% elif row.review_status == 'suggested' %} + Suggested + {% else %} + Review required + {% endif %} +
+
+
GSTIN: {{ row.supplier_gstin or '-' }}
+
Invoice: {{ row.invoice_number }} · {{ row.invoice_date }}
+
HSN: {{ row.hsn_code or '-' }} · POS: {{ row.place_of_supply or '-' }}
+ {% if row.description_text %}
{{ row.description_text }}
{% endif %} +
+
+ +
+
+
Taxable
₹{{ '%.2f'|format(row.taxable_value) }}
+
Invoice Value
₹{{ '%.2f'|format(row.invoice_value) }}
+
IGST
₹{{ '%.2f'|format(row.igst) }}
+
CGST + SGST
₹{{ '%.2f'|format(row.cgst + row.sgst) }}
+
+
+ +
+ {% if suggested %} +
+
+
+
Suggested Nature
+
{{ suggested.name }}
+
+ {{ row.suggested_confidence }}% +
+ {% if row.suggested_ledger_name %}
Ledger: {{ row.suggested_ledger_name }}
{% endif %} + {% if explanations.get(row.id) %} +
    + {% for reason in explanations.get(row.id)[:4] %}
  • • {{ reason }}
  • {% endfor %} +
+ {% endif %} +
+ {% endif %} + + {% if row.review_status != 'pending_analysis' %} +
+ + + + +
+ +
+
+ {% else %} +

Run Analyze / Refresh Suggestions before reviewing this invoice.

+ {% endif %} +
+
+
+ {% else %} +
This batch has no imported purchase documents.
+ {% endfor %} +
+ {% endif %} + {% endif %} +
+{% endblock %} diff --git a/app/modules/accounting/templates/accounting/tally.html b/app/modules/accounting/templates/accounting/tally.html index a134ec1..27a8f07 100644 --- a/app/modules/accounting/templates/accounting/tally.html +++ b/app/modules/accounting/templates/accounting/tally.html @@ -16,6 +16,7 @@ Accounting Taxonomy {% if selected_client %}Historical Learning{% endif %} {% if selected_client %}Ledger Learning{% endif %} + {% if selected_client %}GSTR-2B Intelligence{% endif %} {% if selected_client %}Depreciation (IT){% endif %} Refresh Tally Companies diff --git a/app/ui/app.py b/app/ui/app.py index c21e395..a64160e 100644 --- a/app/ui/app.py +++ b/app/ui/app.py @@ -40,6 +40,7 @@ from app.modules.accounting.ui import router as accounting_ui_router from app.modules.accounting.taxonomy_ui import router as accounting_taxonomy_ui_router 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.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 @@ -66,6 +67,7 @@ def mount_ui(app: FastAPI) -> None: app.include_router(accounting_taxonomy_ui_router) 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(work_tracker_ui_router) app.include_router(billing_ui_router) app.include_router(platform_billing_ui_router)