diff --git a/alembic/versions/20260822_internal_accounting_model_phase13.py b/alembic/versions/20260822_internal_accounting_model_phase13.py new file mode 100644 index 0000000..052507c --- /dev/null +++ b/alembic/versions/20260822_internal_accounting_model_phase13.py @@ -0,0 +1,76 @@ +"""Phase 13 internal accounting model. + +Revision ID: 20260822_internal_accounting_model_p13 +Revises: 20260822_accounting_ai_p12 +""" +from alembic import op +import sqlalchemy as sa + + +revision = "20260822_internal_accounting_model_p13" +down_revision = "20260822_accounting_ai_p12" +branch_labels = None +depends_on = None + + +def upgrade(): + op.create_table( + "accounting_internal_models", + sa.Column("id", sa.Integer(), primary_key=True), + sa.Column("tenant_id", sa.Integer(), sa.ForeignKey("tenants.id", ondelete="CASCADE"), nullable=False), + sa.Column("version_label", sa.String(80), nullable=False), + sa.Column("algorithm", sa.String(80), nullable=False, server_default="multinomial_nb_v1"), + sa.Column("status", sa.String(30), nullable=False, server_default="shadow"), + sa.Column("is_active", sa.Boolean(), nullable=False, server_default=sa.false()), + sa.Column("training_examples", sa.Integer(), nullable=False, server_default="0"), + sa.Column("validation_examples", sa.Integer(), nullable=False, server_default="0"), + sa.Column("class_count", sa.Integer(), nullable=False, server_default="0"), + sa.Column("vocabulary_size", sa.Integer(), nullable=False, server_default="0"), + sa.Column("validation_accuracy", sa.Float(), nullable=False, server_default="0"), + sa.Column("validation_macro_recall", sa.Float(), nullable=False, server_default="0"), + sa.Column("validation_top2_accuracy", sa.Float(), nullable=False, server_default="0"), + sa.Column("model_json", sa.Text(), nullable=False), + sa.Column("metrics_json", sa.Text(), nullable=True), + sa.Column("training_summary_json", sa.Text(), nullable=True), + sa.Column("trained_by_user_id", sa.Integer(), sa.ForeignKey("users.id", ondelete="SET NULL"), nullable=True), + sa.Column("trained_at_utc", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()), + ) + for name in ("tenant_id", "version_label", "status", "is_active", "trained_at_utc"): + op.create_index(f"ix_accounting_internal_models_{name}", "accounting_internal_models", [name]) + + op.create_table( + "accounting_internal_predictions", + 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("model_id", sa.Integer(), sa.ForeignKey("accounting_internal_models.id", ondelete="CASCADE"), nullable=False), + sa.Column("source_type", sa.String(30), nullable=False), + sa.Column("source_record_id", sa.Integer(), nullable=False), + sa.Column("source_fingerprint", sa.String(80), nullable=False, server_default=""), + sa.Column("predicted_nature_id", sa.Integer(), sa.ForeignKey("accounting_natures.id", ondelete="SET NULL"), nullable=True), + sa.Column("predicted_probability", sa.Float(), nullable=False, server_default="0"), + sa.Column("top2_json", sa.Text(), nullable=True), + sa.Column("explanation_json", sa.Text(), nullable=True), + sa.Column("shadow_mode", sa.Boolean(), nullable=False, server_default=sa.true()), + sa.Column("applied_to_source", sa.Boolean(), nullable=False, server_default=sa.false()), + sa.Column("final_nature_id", sa.Integer(), sa.ForeignKey("accounting_natures.id", ondelete="SET NULL"), nullable=True), + sa.Column("prediction_correct", sa.Boolean(), nullable=True), + 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("created_at_utc", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()), + ) + for name in ( + "tenant_id", "client_id", "model_id", "source_type", "source_record_id", + "source_fingerprint", "predicted_nature_id", "shadow_mode", "applied_to_source", + "final_nature_id", "prediction_correct", "created_at_utc", + ): + op.create_index( + f"ix_accounting_internal_predictions_{name}", + "accounting_internal_predictions", + [name], + ) + + +def downgrade(): + op.drop_table("accounting_internal_predictions") + op.drop_table("accounting_internal_models") diff --git a/app/modules/accounting/bank_service.py b/app/modules/accounting/bank_service.py index 8edebbe..ecdf3f2 100644 --- a/app/modules/accounting/bank_service.py +++ b/app/modules/accounting/bank_service.py @@ -10,6 +10,7 @@ from sqlalchemy import func, select from app.modules.accounting.bank_models import AccountingBankTransaction from app.modules.accounting.ai_service import ai_assist_bank as run_ai_assist_bank, mark_review_outcome +from app.modules.accounting.internal_model_service import predict_bank, record_prediction_review from app.modules.accounting.ledger_learning_service import ( active_natures, available_tally_guids, @@ -249,6 +250,14 @@ def confirm_review(db, *, tx_id, tenant_id, client_id, nature_id, ledger_name, v final_ledger_name=tx.final_ledger_name, user_id=user_id, ) + record_prediction_review( + db, + tenant_id=tenant_id, + source_type="bank", + source_record_id=tx.id, + final_nature_id=tx.final_nature_id, + user_id=user_id, + ) return tx @@ -260,6 +269,14 @@ def ai_assist_transaction(db, *, tx_id, tenant_id, client_id, user_id): +def internal_model_predict_transaction(db, *, tx_id, tenant_id, client_id, force_shadow=True): + tx = db.get(AccountingBankTransaction, int(tx_id)) + if not tx or int(tx.tenant_id) != int(tenant_id) or int(tx.client_id) != int(client_id): + raise ValueError("Bank transaction was not found.") + return predict_bank(db, row=tx, force_shadow=force_shadow) + + + def visible_workstations(db, tenant_id, branch_id=None): stmt = select(ERPWorkstationAgent).where( ERPWorkstationAgent.tenant_id == tenant_id, diff --git a/app/modules/accounting/bank_ui.py b/app/modules/accounting/bank_ui.py index c7fe337..7ab0fd7 100644 --- a/app/modules/accounting/bank_ui.py +++ b/app/modules/accounting/bank_ui.py @@ -9,7 +9,7 @@ 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.bank_service import ( - ai_assist_transaction, confirm_review, import_completed_job, preflight_choices, queue_post, queue_preflight, + ai_assist_transaction, confirm_review, import_completed_job, internal_model_predict_transaction, preflight_choices, queue_post, queue_preflight, queue_rows, sync_posting, visible_workstations, ) from app.modules.accounting.ledger_learning_service import active_natures @@ -155,3 +155,37 @@ def ai_assist( return _go(client_id, error=str(exc)) finally: db.close() + + +@router.post("/{tx_id}/internal-predict") +def internal_predict( + request: Request, + tx_id: int, + client_id: int = Form(...), + csrf_token: str = Form(...), +): + validate_csrf(request, csrf_token) + db = CommonSessionLocal() + try: + user, denied = _require_partner(request, db, "accounting.learning.manage") + if denied: + return denied + client, _, scope = _find_visible_client(db, request, user, client_id) + if not client: + return _go(error="Client is not visible.") + prediction = internal_model_predict_transaction( + db, + tx_id=tx_id, + tenant_id=scope.tenant_id, + client_id=client.id, + force_shadow=True, + ) + return _go( + client.id, + message=f"Internal model shadow prediction recorded at {prediction.predicted_probability * 100:.1f}%.", + ) + except Exception as exc: + db.rollback() + return _go(client_id, error=str(exc)) + finally: + db.close() diff --git a/app/modules/accounting/internal_model_engine.py b/app/modules/accounting/internal_model_engine.py new file mode 100644 index 0000000..d155b9b --- /dev/null +++ b/app/modules/accounting/internal_model_engine.py @@ -0,0 +1,223 @@ +from __future__ import annotations + +import hashlib +import json +import math +import re +from collections import Counter, defaultdict +from dataclasses import dataclass +from typing import Any + + +_TOKEN_RE = re.compile(r"[A-Z0-9]{2,}") +_STOP = { + "THE","AND","FOR","FROM","WITH","THIS","THAT","OF","TO","IN","ON","AT","BY", + "PVT","PRIVATE","LTD","LIMITED","LLP","INDIA","GST","GSTIN","INVOICE","BILL", + "PAYMENT","PAID","RECEIPT","TRANSFER","BANK","ACCOUNT","AC","A","AN", +} + + +def _s(value) -> str: + return str(value or "").strip() + + +def amount_bucket(value) -> str: + try: + amount = abs(float(value or 0)) + except Exception: + amount = 0 + if amount < 1_000: + return "AMT_LT_1K" + if amount < 10_000: + return "AMT_1K_10K" + if amount < 50_000: + return "AMT_10K_50K" + if amount < 200_000: + return "AMT_50K_2L" + if amount < 1_000_000: + return "AMT_2L_10L" + return "AMT_GE_10L" + + +def text_tokens(value: str, prefix: str = "") -> list[str]: + tokens = [] + for token in _TOKEN_RE.findall(_s(value).upper()): + if token in _STOP or len(token) < 2: + continue + if token.isdigit() and len(token) > 6: + continue + tokens.append(f"{prefix}{token}" if prefix else token) + return tokens[:120] + + +def feature_tokens(context: dict[str, Any]) -> list[str]: + tokens: list[str] = [] + + source = _s(context.get("source_type")).upper() + if source: + tokens.append("SRC_" + source) + + direction = _s(context.get("direction")).upper() + if direction: + tokens.append("DIR_" + direction) + + hsn = re.sub(r"\D", "", _s(context.get("hsn_code"))) + if hsn: + tokens.append("HSN2_" + hsn[:2]) + if len(hsn) >= 4: + tokens.append("HSN4_" + hsn[:4]) + + tokens.extend(text_tokens(context.get("supplier_name", ""), "SUP_")) + tokens.extend(text_tokens(context.get("supplier_gstin", ""), "GSTIN_")) + tokens.extend(text_tokens(context.get("party_name", ""), "PTY_")) + tokens.extend(text_tokens(context.get("description", ""), "TXT_")) + tokens.extend(text_tokens(context.get("narration", ""), "TXT_")) + tokens.extend(text_tokens(context.get("primary_industry", ""), "IND_")) + tokens.extend(text_tokens(context.get("primary_business_activity", ""), "BUS_")) + tokens.extend(text_tokens(context.get("main_products", ""), "PROD_")) + tokens.extend(text_tokens(context.get("main_services", ""), "SERV_")) + + if context.get("inventory_maintained") is True: + tokens.append("PROFILE_INVENTORY") + if context.get("capital_intensive") is True: + tokens.append("PROFILE_CAPITAL_INTENSIVE") + if context.get("vehicle_intensive") is True: + tokens.append("PROFILE_VEHICLE_INTENSIVE") + if context.get("project_job_based") is True: + tokens.append("PROFILE_PROJECT_BASED") + + tokens.append(amount_bucket(context.get("amount"))) + + # Preserve multiplicity modestly because NB benefits from repeated semantic cues. + return tokens[:240] + + +@dataclass +class Prediction: + class_id: int + probability: float + ranked: list[tuple[int, float]] + evidence: list[str] + + +class MultinomialNBClassifier: + algorithm = "multinomial_nb_v1" + + @staticmethod + def train(examples: list[dict[str, Any]], *, alpha: float = 1.0) -> dict[str, Any]: + if len(examples) < 20: + raise ValueError("At least 20 reviewed accounting examples are required to train the first internal model.") + + class_docs = Counter() + class_token_counts: dict[int, Counter] = defaultdict(Counter) + class_total_tokens = Counter() + vocabulary = set() + + for example in examples: + class_id = int(example["nature_id"]) + features = feature_tokens(example["context"]) + if not features: + continue + counts = Counter(features) + class_docs[class_id] += 1 + class_token_counts[class_id].update(counts) + class_total_tokens[class_id] += sum(counts.values()) + vocabulary.update(counts) + + classes = sorted(class_docs) + if len(classes) < 2: + raise ValueError("At least two different reviewed Accounting Natures are required for training.") + + usable_docs = sum(class_docs.values()) + if usable_docs < 20: + raise ValueError("At least 20 reviewed examples with usable accounting context are required.") + + model = { + "algorithm": MultinomialNBClassifier.algorithm, + "alpha": float(alpha), + "document_count": usable_docs, + "classes": classes, + "class_docs": {str(k): int(v) for k, v in class_docs.items()}, + "class_total_tokens": {str(k): int(v) for k, v in class_total_tokens.items()}, + "token_counts": { + str(class_id): dict(class_token_counts[class_id]) + for class_id in classes + }, + "vocabulary_size": len(vocabulary), + } + return model + + @staticmethod + def predict(model: dict[str, Any], context: dict[str, Any]) -> Prediction: + features = Counter(feature_tokens(context)) + if not features: + raise ValueError("The transaction does not contain enough usable context for internal-model prediction.") + + classes = [int(value) for value in model["classes"]] + class_docs = {int(k): int(v) for k, v in model["class_docs"].items()} + class_total = {int(k): int(v) for k, v in model["class_total_tokens"].items()} + token_counts = { + int(k): {token: int(count) for token, count in values.items()} + for k, values in model["token_counts"].items() + } + alpha = float(model.get("alpha") or 1.0) + vocab_size = max(1, int(model.get("vocabulary_size") or 1)) + total_docs = max(1, sum(class_docs.values())) + class_count = len(classes) + + scores: dict[int, float] = {} + evidence_by_class: dict[int, list[tuple[float, str]]] = defaultdict(list) + + for class_id in classes: + prior = (class_docs[class_id] + alpha) / (total_docs + alpha * class_count) + score = math.log(prior) + denominator = class_total.get(class_id, 0) + alpha * vocab_size + counts = token_counts.get(class_id, {}) + + for token, multiplicity in features.items(): + numerator = counts.get(token, 0) + alpha + contribution = multiplicity * math.log(numerator / denominator) + score += contribution + if counts.get(token, 0) > 0: + evidence_by_class[class_id].append((counts[token], token)) + scores[class_id] = score + + maximum = max(scores.values()) + exp_scores = {key: math.exp(value - maximum) for key, value in scores.items()} + total = sum(exp_scores.values()) or 1.0 + probabilities = {key: value / total for key, value in exp_scores.items()} + ranked = sorted(probabilities.items(), key=lambda item: (-item[1], item[0])) + winner = ranked[0][0] + evidence = [ + token for _, token in sorted(evidence_by_class[winner], reverse=True)[:8] + ] + return Prediction( + class_id=winner, + probability=float(ranked[0][1]), + ranked=[(int(k), float(v)) for k, v in ranked[:3]], + evidence=evidence, + ) + + +def deterministic_split(examples: list[dict[str, Any]], validation_ratio: float = 0.20): + train = [] + validation = [] + by_class: dict[int, list[dict[str, Any]]] = defaultdict(list) + for example in examples: + by_class[int(example["nature_id"])].append(example) + + for class_id, rows in by_class.items(): + rows = sorted( + rows, + key=lambda row: hashlib.sha256( + f"{row['source_type']}|{row['source_record_id']}|{class_id}".encode() + ).hexdigest(), + ) + if len(rows) < 4: + train.extend(rows) + continue + validation_count = max(1, int(round(len(rows) * validation_ratio))) + validation.extend(rows[:validation_count]) + train.extend(rows[validation_count:]) + + return train, validation diff --git a/app/modules/accounting/internal_model_models.py b/app/modules/accounting/internal_model_models.py new file mode 100644 index 0000000..ca71a0b --- /dev/null +++ b/app/modules/accounting/internal_model_models.py @@ -0,0 +1,78 @@ +from __future__ import annotations + +from datetime import datetime, timezone + +from sqlalchemy import Boolean, DateTime, Float, ForeignKey, Integer, String, Text +from sqlalchemy.orm import Mapped, mapped_column + +from app.core.db.common import CommonBase + + +class AccountingInternalModel(CommonBase): + __tablename__ = "accounting_internal_models" + + 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) + + version_label: Mapped[str] = mapped_column(String(80), nullable=False, index=True) + algorithm: Mapped[str] = mapped_column(String(80), nullable=False, default="multinomial_nb_v1") + status: Mapped[str] = mapped_column(String(30), nullable=False, default="shadow", index=True) + is_active: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False, index=True) + + training_examples: Mapped[int] = mapped_column(Integer, nullable=False, default=0) + validation_examples: Mapped[int] = mapped_column(Integer, nullable=False, default=0) + class_count: Mapped[int] = mapped_column(Integer, nullable=False, default=0) + vocabulary_size: Mapped[int] = mapped_column(Integer, nullable=False, default=0) + + validation_accuracy: Mapped[float] = mapped_column(Float, nullable=False, default=0) + validation_macro_recall: Mapped[float] = mapped_column(Float, nullable=False, default=0) + validation_top2_accuracy: Mapped[float] = mapped_column(Float, nullable=False, default=0) + + model_json: Mapped[str] = mapped_column(Text, nullable=False) + metrics_json: Mapped[str | None] = mapped_column(Text, nullable=True) + training_summary_json: Mapped[str | None] = mapped_column(Text, nullable=True) + + trained_by_user_id: Mapped[int | None] = mapped_column( + ForeignKey("users.id", ondelete="SET NULL"), nullable=True + ) + trained_at_utc: Mapped[datetime] = mapped_column( + DateTime(timezone=True), nullable=False, default=lambda: datetime.now(timezone.utc), index=True + ) + + +class AccountingInternalPrediction(CommonBase): + __tablename__ = "accounting_internal_predictions" + + 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) + model_id: Mapped[int] = mapped_column( + ForeignKey("accounting_internal_models.id", ondelete="CASCADE"), nullable=False, index=True + ) + + source_type: Mapped[str] = mapped_column(String(30), nullable=False, index=True) + source_record_id: Mapped[int] = mapped_column(Integer, nullable=False, index=True) + source_fingerprint: Mapped[str] = mapped_column(String(80), nullable=False, default="", index=True) + + predicted_nature_id: Mapped[int | None] = mapped_column( + ForeignKey("accounting_natures.id", ondelete="SET NULL"), nullable=True, index=True + ) + predicted_probability: Mapped[float] = mapped_column(Float, nullable=False, default=0) + top2_json: Mapped[str | None] = mapped_column(Text, nullable=True) + explanation_json: Mapped[str | None] = mapped_column(Text, nullable=True) + + shadow_mode: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True, index=True) + applied_to_source: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False, index=True) + + final_nature_id: Mapped[int | None] = mapped_column( + ForeignKey("accounting_natures.id", ondelete="SET NULL"), nullable=True, index=True + ) + prediction_correct: Mapped[bool | None] = mapped_column(Boolean, nullable=True, index=True) + 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) + + created_at_utc: Mapped[datetime] = mapped_column( + DateTime(timezone=True), nullable=False, default=lambda: datetime.now(timezone.utc), index=True + ) diff --git a/app/modules/accounting/internal_model_service.py b/app/modules/accounting/internal_model_service.py new file mode 100644 index 0000000..f9ca94d --- /dev/null +++ b/app/modules/accounting/internal_model_service.py @@ -0,0 +1,444 @@ +from __future__ import annotations + +import hashlib +import json +import os +from collections import defaultdict +from datetime import datetime, timezone +from typing import Any + +from sqlalchemy import func, select, update + +from app.modules.accounting.bank_models import AccountingBankTransaction +from app.modules.accounting.gstr2b_models import AccountingGSTR2BPurchase +from app.modules.accounting.internal_model_engine import MultinomialNBClassifier, deterministic_split +from app.modules.accounting.internal_model_models import AccountingInternalModel, AccountingInternalPrediction +from app.modules.accounting.taxonomy_models import AccountingNature +from app.modules.clients.models import ClientBusinessProfile + + +def _utcnow(): + return datetime.now(timezone.utc) + + +def _s(value): + return str(value or "").strip() + + +def internal_model_enabled() -> bool: + return (os.getenv("ACCOUNTING_INTERNAL_MODEL_ENABLED") or "false").strip().lower() in { + "1", "true", "yes", "on" + } + + +def minimum_apply_probability() -> float: + try: + value = float(os.getenv("ACCOUNTING_INTERNAL_MODEL_MIN_PROBABILITY", "0.80")) + except Exception: + value = 0.80 + return max(0.50, min(0.99, value)) + + +def _profile(db, client_id: int) -> dict[str, Any]: + profile = db.execute( + select(ClientBusinessProfile).where(ClientBusinessProfile.client_id == int(client_id)) + ).scalar_one_or_none() + if not profile: + return {} + return { + "primary_industry": _s(profile.primary_industry), + "primary_business_activity": _s(profile.primary_business_activity), + "main_products": _s(profile.main_products), + "main_services": _s(profile.main_services), + "inventory_maintained": profile.inventory_maintained, + "capital_intensive": profile.capital_intensive, + "vehicle_intensive": profile.vehicle_intensive, + "project_job_based": profile.project_job_based, + } + + +def _purchase_context(db, row: AccountingGSTR2BPurchase) -> dict[str, Any]: + return { + "source_type": "gstr2b", + "supplier_name": row.supplier_name, + "supplier_gstin": row.supplier_gstin, + "hsn_code": row.hsn_code, + "description": row.description_text or "", + "amount": row.taxable_value or row.invoice_value or 0, + **_profile(db, row.client_id), + } + + +def _bank_context(db, row: AccountingBankTransaction) -> dict[str, Any]: + return { + "source_type": "bank", + "party_name": row.auto_party, + "narration": row.narration, + "direction": row.direction, + "amount": row.amount, + **_profile(db, row.client_id), + } + + +def reviewed_examples(db, *, tenant_id: int) -> list[dict[str, Any]]: + examples = [] + + purchases = list(db.execute( + select(AccountingGSTR2BPurchase).where( + AccountingGSTR2BPurchase.tenant_id == int(tenant_id), + AccountingGSTR2BPurchase.review_status == "reviewed", + AccountingGSTR2BPurchase.final_nature_id.is_not(None), + ) + ).scalars().all()) + for row in purchases: + examples.append({ + "source_type": "gstr2b", + "source_record_id": int(row.id), + "client_id": int(row.client_id), + "nature_id": int(row.final_nature_id), + "context": _purchase_context(db, row), + }) + + bank_rows = list(db.execute( + select(AccountingBankTransaction).where( + AccountingBankTransaction.tenant_id == int(tenant_id), + AccountingBankTransaction.review_status == "reviewed", + AccountingBankTransaction.final_nature_id.is_not(None), + AccountingBankTransaction.final_voucher_type != "Contra", + ) + ).scalars().all()) + for row in bank_rows: + examples.append({ + "source_type": "bank", + "source_record_id": int(row.id), + "client_id": int(row.client_id), + "nature_id": int(row.final_nature_id), + "context": _bank_context(db, row), + }) + + return examples + + +def _metrics(model, validation): + if not validation: + return { + "accuracy": 0.0, + "macro_recall": 0.0, + "top2_accuracy": 0.0, + "validation_examples": 0, + "confusion": {}, + } + + correct = 0 + top2 = 0 + totals = defaultdict(int) + class_correct = defaultdict(int) + confusion = defaultdict(lambda: defaultdict(int)) + + for example in validation: + pred = MultinomialNBClassifier.predict(model, example["context"]) + actual = int(example["nature_id"]) + predicted = int(pred.class_id) + totals[actual] += 1 + confusion[actual][predicted] += 1 + if predicted == actual: + correct += 1 + class_correct[actual] += 1 + if actual in [class_id for class_id, _ in pred.ranked[:2]]: + top2 += 1 + + recalls = [ + class_correct[class_id] / totals[class_id] + for class_id in totals + if totals[class_id] > 0 + ] + return { + "accuracy": correct / len(validation), + "macro_recall": sum(recalls) / len(recalls) if recalls else 0.0, + "top2_accuracy": top2 / len(validation), + "validation_examples": len(validation), + "confusion": { + str(actual): {str(pred): count for pred, count in values.items()} + for actual, values in confusion.items() + }, + } + + +def train_model(db, *, tenant_id: int, user_id: int) -> AccountingInternalModel: + examples = reviewed_examples(db, tenant_id=tenant_id) + train, validation = deterministic_split(examples) + + model_blob = MultinomialNBClassifier.train(train) + metrics = _metrics(model_blob, validation) + + timestamp = _utcnow().strftime("%Y%m%d-%H%M%S") + version = f"ARRR-ACC-{timestamp}" + + row = AccountingInternalModel( + tenant_id=int(tenant_id), + version_label=version, + algorithm=MultinomialNBClassifier.algorithm, + status="shadow", + is_active=False, + training_examples=len(train), + validation_examples=len(validation), + class_count=len(model_blob["classes"]), + vocabulary_size=int(model_blob["vocabulary_size"]), + validation_accuracy=float(metrics["accuracy"]), + validation_macro_recall=float(metrics["macro_recall"]), + validation_top2_accuracy=float(metrics["top2_accuracy"]), + model_json=json.dumps(model_blob, separators=(",", ":")), + metrics_json=json.dumps(metrics, separators=(",", ":")), + training_summary_json=json.dumps({ + "total_reviewed_examples": len(examples), + "gstr2b_examples": sum(1 for x in examples if x["source_type"] == "gstr2b"), + "bank_examples": sum(1 for x in examples if x["source_type"] == "bank"), + }), + trained_by_user_id=user_id, + ) + db.add(row) + db.commit() + db.refresh(row) + return row + + +def active_model(db, *, tenant_id: int): + return db.execute( + select(AccountingInternalModel).where( + AccountingInternalModel.tenant_id == int(tenant_id), + AccountingInternalModel.is_active.is_(True), + ).order_by(AccountingInternalModel.id.desc()).limit(1) + ).scalar_one_or_none() + + +def latest_model(db, *, tenant_id: int): + return db.execute( + select(AccountingInternalModel).where( + AccountingInternalModel.tenant_id == int(tenant_id) + ).order_by(AccountingInternalModel.id.desc()).limit(1) + ).scalar_one_or_none() + + +def activate_model(db, *, tenant_id: int, model_id: int): + model = db.get(AccountingInternalModel, int(model_id)) + if not model or int(model.tenant_id) != int(tenant_id): + raise ValueError("Internal accounting model was not found.") + if model.validation_examples and model.validation_accuracy < 0.55: + raise ValueError( + "Validation accuracy is below 55%. Keep this model in shadow mode and collect more reviewed examples." + ) + + db.execute( + update(AccountingInternalModel).where( + AccountingInternalModel.tenant_id == int(tenant_id) + ).values(is_active=False) + ) + model.is_active = True + model.status = "active" + db.add(model) + db.commit() + db.refresh(model) + return model + + +def deactivate_models(db, *, tenant_id: int): + db.execute( + update(AccountingInternalModel).where( + AccountingInternalModel.tenant_id == int(tenant_id) + ).values(is_active=False, status="shadow") + ) + db.commit() + + +def _fingerprint(source_type, source_record_id, context): + raw = f"{source_type}|{source_record_id}|{json.dumps(context, sort_keys=True, ensure_ascii=False)}" + return hashlib.sha256(raw.encode("utf-8", "ignore")).hexdigest() + + +def _save_prediction( + db, *, + model: AccountingInternalModel, + tenant_id: int, + client_id: int, + source_type: str, + source_record_id: int, + context: dict[str, Any], + apply_to_source: bool, +): + blob = json.loads(model.model_json) + pred = MultinomialNBClassifier.predict(blob, context) + + existing = db.execute( + select(AccountingInternalPrediction).where( + AccountingInternalPrediction.model_id == model.id, + AccountingInternalPrediction.source_type == source_type, + AccountingInternalPrediction.source_record_id == int(source_record_id), + ).order_by(AccountingInternalPrediction.id.desc()).limit(1) + ).scalar_one_or_none() + if existing: + return existing, pred + + row = AccountingInternalPrediction( + tenant_id=int(tenant_id), + client_id=int(client_id), + model_id=model.id, + source_type=source_type, + source_record_id=int(source_record_id), + source_fingerprint=_fingerprint(source_type, source_record_id, context), + predicted_nature_id=pred.class_id, + predicted_probability=float(pred.probability), + top2_json=json.dumps( + [{"nature_id": class_id, "probability": probability} for class_id, probability in pred.ranked], + separators=(",", ":"), + ), + explanation_json=json.dumps({"evidence_tokens": pred.evidence}, separators=(",", ":")), + shadow_mode=not apply_to_source, + applied_to_source=False, + ) + db.add(row) + db.commit() + db.refresh(row) + return row, pred + + +def predict_purchase(db, *, row: AccountingGSTR2BPurchase, force_shadow: bool = True): + model = active_model(db, tenant_id=row.tenant_id) or latest_model(db, tenant_id=row.tenant_id) + if not model: + raise ValueError("Train an internal accounting model first.") + + may_apply = ( + internal_model_enabled() + and model.is_active + and not force_shadow + and row.review_status != "reviewed" + ) + prediction, pred = _save_prediction( + db, + model=model, + tenant_id=row.tenant_id, + client_id=row.client_id, + source_type="gstr2b", + source_record_id=row.id, + context=_purchase_context(db, row), + apply_to_source=may_apply, + ) + if may_apply and pred.probability >= minimum_apply_probability(): + # Never invent/alter ledgers; only nature is proposed. + if row.suggested_nature_id != pred.class_id: + row.suggested_ledger_name = "" + row.suggested_nature_id = pred.class_id + row.suggested_confidence = int(round(pred.probability * 100)) + row.review_status = "suggested" + prediction.applied_to_source = True + prediction.shadow_mode = False + db.add(row) + db.add(prediction) + db.commit() + return prediction + + +def predict_bank(db, *, row: AccountingBankTransaction, force_shadow: bool = True): + if _s(row.contra_pair_id): + raise ValueError("Matched Contra transactions do not require internal semantic classification.") + model = active_model(db, tenant_id=row.tenant_id) or latest_model(db, tenant_id=row.tenant_id) + if not model: + raise ValueError("Train an internal accounting model first.") + + may_apply = ( + internal_model_enabled() + and model.is_active + and not force_shadow + and row.review_status != "reviewed" + ) + prediction, pred = _save_prediction( + db, + model=model, + tenant_id=row.tenant_id, + client_id=row.client_id, + source_type="bank", + source_record_id=row.id, + context=_bank_context(db, row), + apply_to_source=may_apply, + ) + if may_apply and pred.probability >= minimum_apply_probability(): + if row.suggested_nature_id != pred.class_id: + row.suggested_ledger_name = "" + row.suggested_nature_id = pred.class_id + row.suggested_confidence = int(round(pred.probability * 100)) + prediction.applied_to_source = True + prediction.shadow_mode = False + db.add(row) + db.add(prediction) + db.commit() + return prediction + + +def record_prediction_review( + db, *, + tenant_id: int, + source_type: str, + source_record_id: int, + final_nature_id: int | None, + user_id: int, +): + row = db.execute( + select(AccountingInternalPrediction).where( + AccountingInternalPrediction.tenant_id == int(tenant_id), + AccountingInternalPrediction.source_type == source_type, + AccountingInternalPrediction.source_record_id == int(source_record_id), + ).order_by(AccountingInternalPrediction.id.desc()).limit(1) + ).scalar_one_or_none() + if not row: + return None + row.final_nature_id = final_nature_id + row.prediction_correct = bool( + final_nature_id and row.predicted_nature_id and int(final_nature_id) == int(row.predicted_nature_id) + ) + row.reviewed_by_user_id = user_id + row.reviewed_at_utc = _utcnow() + db.add(row) + db.commit() + return row + + +def model_dashboard(db, *, tenant_id: int): + models = list(db.execute( + select(AccountingInternalModel).where( + AccountingInternalModel.tenant_id == int(tenant_id) + ).order_by(AccountingInternalModel.id.desc()).limit(20) + ).scalars().all()) + + predictions = int(db.scalar(select(func.count(AccountingInternalPrediction.id)).where( + AccountingInternalPrediction.tenant_id == int(tenant_id) + )) or 0) + reviewed = int(db.scalar(select(func.count(AccountingInternalPrediction.id)).where( + AccountingInternalPrediction.tenant_id == int(tenant_id), + AccountingInternalPrediction.prediction_correct.is_not(None), + )) or 0) + correct = int(db.scalar(select(func.count(AccountingInternalPrediction.id)).where( + AccountingInternalPrediction.tenant_id == int(tenant_id), + AccountingInternalPrediction.prediction_correct.is_(True), + )) or 0) + applied = int(db.scalar(select(func.count(AccountingInternalPrediction.id)).where( + AccountingInternalPrediction.tenant_id == int(tenant_id), + AccountingInternalPrediction.applied_to_source.is_(True), + )) or 0) + + return { + "models": models, + "predictions": predictions, + "reviewed_predictions": reviewed, + "correct_predictions": correct, + "live_accuracy": round(correct * 100 / reviewed, 1) if reviewed else 0, + "applied_predictions": applied, + "enabled": internal_model_enabled(), + "minimum_apply_probability": minimum_apply_probability(), + } + + +def recent_predictions(db, *, tenant_id: int, limit: int = 100): + return list(db.execute( + select(AccountingInternalPrediction).where( + AccountingInternalPrediction.tenant_id == int(tenant_id) + ).order_by(AccountingInternalPrediction.id.desc()).limit(max(1, min(500, int(limit)))) + ).scalars().all()) diff --git a/app/modules/accounting/internal_model_ui.py b/app/modules/accounting/internal_model_ui.py new file mode 100644 index 0000000..9b532fe --- /dev/null +++ b/app/modules/accounting/internal_model_ui.py @@ -0,0 +1,122 @@ +from __future__ import annotations + +from fastapi import APIRouter, Form, Request +from fastapi.responses import RedirectResponse + +from app.core.db.common import CommonSessionLocal +from app.core.security.csrf import get_or_create_csrf_token, validate_csrf +from app.core.templating import templates +from app.modules.accounting.internal_model_service import ( + activate_model, + deactivate_models, + model_dashboard, + recent_predictions, + train_model, +) +from app.modules.accounting.ui import _require_partner +from app.modules.core.rbac.deps import get_user_permissions, get_user_roles + +router = APIRouter(prefix="/tools/accounting/internal-model", tags=["accounting-internal-model-ui"]) + + +def _redirect(message="", error=""): + from urllib.parse import urlencode + q = {} + if message: + q["message"] = message[:300] + if error: + q["error"] = error[:300] + return RedirectResponse( + "/tools/accounting/internal-model" + ("?" + urlencode(q) if q else ""), + status_code=303, + ) + + +@router.get("") +def page(request: Request, message: str = "", error: str = ""): + db = CommonSessionLocal() + try: + user, denied = _require_partner(request, db, "accounting.learning.view") + if denied: + return denied + scope = request.state.workspace_scope + summary = model_dashboard(db, tenant_id=scope.tenant_id) + predictions = recent_predictions(db, tenant_id=scope.tenant_id, limit=100) + return templates.TemplateResponse( + "modules/accounting/templates/accounting/internal_model.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": "Internal Accounting Model", + "summary": summary, + "predictions": predictions, + "message": message, + "error": error, + }, + ) + finally: + db.close() + + +@router.post("/train") +def train(request: Request, csrf_token: str = Form(...)): + validate_csrf(request, csrf_token) + db = CommonSessionLocal() + try: + user, denied = _require_partner(request, db, "accounting.learning.manage") + if denied: + return denied + scope = request.state.workspace_scope + model = train_model(db, tenant_id=scope.tenant_id, user_id=user.id) + return _redirect( + message=( + f"Internal model {model.version_label} trained. " + f"Validation accuracy {model.validation_accuracy * 100:.1f}%." + ) + ) + except Exception as exc: + db.rollback() + return _redirect(error=str(exc)) + finally: + db.close() + + +@router.post("/{model_id}/activate") +def activate(request: Request, model_id: int, csrf_token: str = Form(...)): + validate_csrf(request, csrf_token) + db = CommonSessionLocal() + try: + user, denied = _require_partner(request, db, "accounting.learning.manage") + if denied: + return denied + scope = request.state.workspace_scope + model = activate_model(db, tenant_id=scope.tenant_id, model_id=model_id) + return _redirect( + message=( + f"{model.version_label} marked active. " + "It still remains shadow-only unless ACCOUNTING_INTERNAL_MODEL_ENABLED=true on the server." + ) + ) + except Exception as exc: + db.rollback() + return _redirect(error=str(exc)) + finally: + db.close() + + +@router.post("/deactivate") +def deactivate(request: Request, csrf_token: str = Form(...)): + validate_csrf(request, csrf_token) + db = CommonSessionLocal() + try: + user, denied = _require_partner(request, db, "accounting.learning.manage") + if denied: + return denied + scope = request.state.workspace_scope + deactivate_models(db, tenant_id=scope.tenant_id) + return _redirect(message="All internal accounting models returned to shadow mode.") + finally: + db.close() diff --git a/app/modules/accounting/purchase_review_service.py b/app/modules/accounting/purchase_review_service.py index 1aef241..bf8b479 100644 --- a/app/modules/accounting/purchase_review_service.py +++ b/app/modules/accounting/purchase_review_service.py @@ -10,6 +10,7 @@ from app.modules.accounting.gstr2b_service import nature_lookup, review_purchase from app.modules.accounting.historical_learning_service import active_natures, ledger_mappings from app.modules.accounting.purchase_enrichment_models import AccountingPurchaseEnrichmentRecord from app.modules.accounting.ai_service import ai_assist_purchase as run_ai_assist_purchase, mark_review_outcome +from app.modules.accounting.internal_model_service import predict_purchase, record_prediction_review VALID_STATUSES = {"all", "pending_analysis", "suggested", "review_required", "reviewed"} VALID_CONFIDENCE = {"all", "high", "medium", "low"} @@ -144,6 +145,14 @@ def review_one(db, *, tenant_id, client_id, purchase_id, final_nature_id, final_ final_ledger_name=reviewed.final_ledger_name, user_id=user_id, ) + record_prediction_review( + db, + tenant_id=tenant_id, + source_type="gstr2b", + source_record_id=reviewed.id, + final_nature_id=reviewed.final_nature_id, + user_id=user_id, + ) return reviewed @@ -169,6 +178,14 @@ def bulk_confirm(db, *, tenant_id, client_id, purchase_ids, user_id): final_ledger_name=reviewed.final_ledger_name, user_id=user_id, ) + record_prediction_review( + db, + tenant_id=tenant_id, + source_type="gstr2b", + source_record_id=reviewed.id, + final_nature_id=reviewed.final_nature_id, + user_id=user_id, + ) completed+=1 return completed @@ -182,3 +199,14 @@ def ai_assist_one(db, *, tenant_id, client_id, purchase_id, user_id): if not row: raise ValueError("Purchase record was not found.") return run_ai_assist_purchase(db, row=row, user_id=user_id) + + +def internal_model_predict_one(db, *, tenant_id, client_id, purchase_id, force_shadow=True): + row = db.execute(select(AccountingGSTR2BPurchase).where( + AccountingGSTR2BPurchase.id == int(purchase_id), + AccountingGSTR2BPurchase.tenant_id == int(tenant_id), + AccountingGSTR2BPurchase.client_id == int(client_id), + )).scalar_one_or_none() + if not row: + raise ValueError("Purchase record was not found.") + return predict_purchase(db, row=row, force_shadow=force_shadow) diff --git a/app/modules/accounting/purchase_review_ui.py b/app/modules/accounting/purchase_review_ui.py index 7de8bdd..1dd4316 100644 --- a/app/modules/accounting/purchase_review_ui.py +++ b/app/modules/accounting/purchase_review_ui.py @@ -11,7 +11,7 @@ from app.modules.accounting.historical_learning_service import active_natures from app.modules.accounting.ledger_learning_service import available_tally_guids from app.modules.accounting.purchase_review_service import ( review_queue, review_counts, return_periods, nature_maps_for_rows, - explanations_for_rows, mappings_by_nature, review_one, bulk_confirm, ai_assist_one, + explanations_for_rows, mappings_by_nature, review_one, bulk_confirm, ai_assist_one, internal_model_predict_one, ) 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 @@ -139,3 +139,52 @@ def ai_assist_purchase_row( return _redirect(client_id, filters, error=str(exc)) finally: db.close() + + +@router.post("/purchase/{purchase_id}/internal-predict") +def internal_predict_purchase_row( + request: Request, + purchase_id: int, + client_id: int = Form(...), + tally_guid: str = Form(""), + status: str = Form("all"), + confidence: str = Form("all"), + source: str = Form("all"), + supplier: str = Form(""), + return_period: str = Form(""), + page: int = Form(1), + per_page: int = Form(25), + csrf_token: str = Form(...), +): + validate_csrf(request, csrf_token) + db = CommonSessionLocal() + filters = { + "tally_guid": tally_guid, "status": status, "confidence": confidence, + "source": source, "supplier": supplier, "return_period": return_period, + "page": page, "per_page": per_page, + } + 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() + prediction = internal_model_predict_one( + db, + tenant_id=scope.tenant_id, + client_id=client.id, + purchase_id=purchase_id, + force_shadow=True, + ) + return _redirect( + client.id, + filters, + message=f"Internal model shadow prediction recorded at {prediction.predicted_probability * 100:.1f}%.", + ) + except Exception as exc: + db.rollback() + return _redirect(client_id, filters, error=str(exc)) + finally: + db.close() diff --git a/app/modules/accounting/templates/accounting/ai_dashboard.html b/app/modules/accounting/templates/accounting/ai_dashboard.html index 6ac45d8..00ec2ee 100644 --- a/app/modules/accounting/templates/accounting/ai_dashboard.html +++ b/app/modules/accounting/templates/accounting/ai_dashboard.html @@ -7,7 +7,7 @@
AI is a fallback only. Confirmed mappings, historical evidence and deterministic rules remain primary. AI can select only an existing Accounting Nature; it cannot create ledgers or post to Tally.
- Back to Tally +