Files
arrr-erp/app/modules/accounting/internal_model_service.py
T
2026-08-22 21:52:44 +05:30

445 lines
15 KiB
Python

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())