from __future__ import annotations from datetime import date, timedelta from urllib.parse import urlencode from fastapi import APIRouter, Form, Request from fastapi.responses import JSONResponse, 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.historical_learning_models import AccountingHistoricalLearningRun from app.modules.accounting.sales_learning_models import AccountingSalesHistoricalRun from app.modules.accounting.sales_learning_service import ( historical_sales_rows, ingest_completed_sales_run, learning_summary, sales_mappings, ) from app.modules.accounting.ui import _accounting_storage_payload, _find_visible_client, _require_partner, _visible_clients from app.modules.accounting.historical_sync_support import eligible_tally_workstations, resolve_workstation_company, workstation_companies from app.modules.accounting.historical_progress import historical_run_status from app.modules.core.rbac.deps import get_user_permissions, get_user_roles from app.modules.documents.agent_jobs import enqueue_agent_job from app.modules.documents.models import ERPWorkstationAgent router = APIRouter( prefix="/tools/accounting/sales-learning", tags=["accounting-sales-learning-ui"], ) def _go(client_id=0, message="", error=""): q = {} if client_id: q["client_id"] = client_id if message: q["message"] = message[:300] if error: q["error"] = error[:300] return RedirectResponse( "/tools/accounting/sales-learning" + ("?" + urlencode(q) if q else ""), status_code=303, ) @router.get("") def page( request: Request, client_id: int | None = None, message: str = "", error: str = "", ): db = CommonSessionLocal() try: user, denied = _require_partner(request, db, "accounting.learning.view") if denied: return denied clients, scope = _visible_clients(db, request, user) selected = next( (row for row in clients if client_id and int(row.id) == int(client_id)), None, ) mappings = [] evidence = [] summary = {"mapping_count": 0, "historical_rows": 0, "review_confirmed": 0, "high_confidence": 0} workstations = [] runs = [] if selected: # Synchronize latest sales evidence run states opportunistically. runs = list(db.execute( select(AccountingSalesHistoricalRun).where( AccountingSalesHistoricalRun.tenant_id == scope.tenant_id, AccountingSalesHistoricalRun.client_id == selected.id, ).order_by(AccountingSalesHistoricalRun.id.desc()).limit(20) ).scalars().all()) for run in runs: if run.status not in {"completed", "failed"}: ingest_completed_sales_run(db, run) mappings = sales_mappings(db, tenant_id=scope.tenant_id, client_id=selected.id) evidence = historical_sales_rows(db, tenant_id=scope.tenant_id, client_id=selected.id)[:200] summary = learning_summary(db, tenant_id=scope.tenant_id, client_id=selected.id) workstations = eligible_tally_workstations(db, scope.tenant_id, scope.branch_id) return templates.TemplateResponse( "modules/accounting/templates/accounting/sales_learning.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": "Customer & Sales Ledger Intelligence", "clients": clients, "selected_client": selected, "mappings": mappings, "evidence": evidence, "summary": summary, "workstations": workstations, "workstation_companies": {w.id: workstation_companies(w) for w in workstations}, "runs": runs, "latest_run": runs[0] if runs else None, "message": message, "error": error, }, ) finally: db.close() @router.post("/collect") def collect( request: Request, client_id: int = Form(...), workstation_id: int = Form(...), tally_guid: str = Form(...), date_from: str = Form(""), date_to: str = 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.") workstation, company = resolve_workstation_company( db, workstation_id=workstation_id, tenant_id=scope.tenant_id, branch_id=scope.branch_id, tally_guid=tally_guid, ) if not date_from: date_from = (date.today() - timedelta(days=730)).isoformat() if not date_to: date_to = date.today().isoformat() payload = { **_accounting_storage_payload(client), "tenant_id": scope.tenant_id, "requested_by_user_id": user.id, "tally_guid": tally_guid, "company_name": str(company.get("name") or "").strip(), "date_from": date_from, "date_to": date_to, "voucher_scope": "sales", "refresh_transactions": True, } job = enqueue_agent_job( db, workstation_agent_id=workstation.id, action="accounting_historical_evidence", payload=payload, priority=6, max_attempts=2, created_by_user_id=user.id, ) run = AccountingSalesHistoricalRun( tenant_id=scope.tenant_id, client_id=client.id, tally_guid=tally_guid, company_name=str(company.get("name") or ""), workstation_agent_id=workstation.id, agent_job_id=job.id, date_from=date_from, date_to=date_to, status="queued", requested_by_user_id=user.id, ) db.add(run) db.commit() return _go( client.id, message="Historical Tally Sales evidence collection queued.", ) except Exception as exc: db.rollback() return _go(client_id, error=str(exc)) finally: db.close() @router.get("/run/{run_id}/status") def run_status(request: Request, run_id: int): db = CommonSessionLocal() try: user, denied = _require_partner(request, db, "accounting.learning.view") if denied: return JSONResponse({"ok": False, "error": "Access denied."}, status_code=403) clients, scope = _visible_clients(db, request, user) visible_client_ids = {int(c.id) for c in clients} run = db.get(AccountingSalesHistoricalRun, int(run_id)) if ( not run or int(run.tenant_id) != int(scope.tenant_id) or int(run.client_id) not in visible_client_ids ): return JSONResponse({"ok": False, "error": "Historical Sales run not found."}, status_code=404) if run.status not in {"completed", "failed"}: ingest_completed_sales_run(db, run) return JSONResponse({"ok": True, **historical_run_status(db, run)}) finally: db.close()