Add live progress for historical Purchase and Sales sync

This commit is contained in:
A R R R Associates
2026-08-23 12:00:21 +05:30
parent 6910b5b6b5
commit 4c6546f48d
5 changed files with 348 additions and 3 deletions
@@ -4,7 +4,7 @@ import json
from datetime import date, datetime, timezone
from fastapi import APIRouter, Form, Request
from fastapi.responses import RedirectResponse
from fastapi.responses import JSONResponse, RedirectResponse
from sqlalchemy import select
from app.core.db.common import CommonSessionLocal
@@ -14,6 +14,7 @@ from app.modules.accounting.historical_learning_models import AccountingHistoric
from app.modules.accounting.historical_learning_service import active_natures, classifiable_ledgers, evidence_rows, ingest_completed_run, ledger_mappings, party_suggestions, remove_ledger_mapping, save_ledger_mapping
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.core.rbac.permission_guard import require_permission
from app.modules.documents.agent_jobs import enqueue_agent_job
@@ -60,7 +61,7 @@ def historical_learning(request: Request, client_id: int | None=None, collected:
return templates.TemplateResponse("modules/accounting/templates/accounting/historical_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":"Historical Tally Learning","clients":clients,"selected_client":selected,"workstations":workstations,"workstation_companies":{w.id:workstation_companies(w) for w in workstations},
"runs":runs,"evidence":evid,"classifiable_ledgers":ledgers,"suggestions":suggestions,"mappings":mappings,"natures":active_natures(db,scope.tenant_id),
"runs":runs,"latest_run":runs[0] if runs else None,"evidence":evid,"classifiable_ledgers":ledgers,"suggestions":suggestions,"mappings":mappings,"natures":active_natures(db,scope.tenant_id),
"date_from":df,"date_to":dt,"collected":bool(collected),"mapped":bool(mapped),"error":error,
})
finally: db.close()
@@ -128,3 +129,30 @@ def unmap_ledger(request: Request, client_id: int=Form(...), mapping_id: int=For
remove_ledger_mapping(db,tenant_id=scope.tenant_id,client_id=client.id,mapping_id=mapping_id)
return RedirectResponse(url=f"/tools/accounting/historical-learning?client_id={client_id}",status_code=303)
finally: db.close()
@router.get("/run/{run_id}/status")
def run_status(request: Request, run_id: int):
db = CommonSessionLocal()
try:
user, response = _require_partner(request, db, "accounting.learning.view")
if response:
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(AccountingHistoricalLearningRun, 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 Purchase run not found."}, status_code=404)
if run.status not in {"completed", "failed"}:
ingest_completed_run(db, run)
return JSONResponse({"ok": True, **historical_run_status(db, run)})
finally:
db.close()
@@ -0,0 +1,97 @@
from __future__ import annotations
from datetime import datetime, timezone
from app.modules.documents.models import ERPAgentJob, ERPWorkstationAgent
def _utcnow():
return datetime.now(timezone.utc)
def _iso(value):
return value.isoformat() if value else None
def _elapsed_seconds(started):
if not started:
return 0
try:
now = _utcnow()
value = started
if value.tzinfo is None:
value = value.replace(tzinfo=timezone.utc)
return max(0, int((now - value).total_seconds()))
except Exception:
return 0
def historical_run_status(db, run):
job = db.get(ERPAgentJob, run.agent_job_id) if run and run.agent_job_id else None
workstation = db.get(ERPWorkstationAgent, run.workstation_agent_id) if run and run.workstation_agent_id else None
job_status = str(job.status if job else run.status or "queued").lower()
run_status = str(run.status or job_status).lower()
if run_status == "completed":
stage = "completed"
label = "Completed"
terminal = True
ok = True
elif run_status == "failed":
stage = "failed"
label = "Failed"
terminal = True
ok = False
elif job_status == "claimed":
stage = "processing"
label = "Processing on workstation"
terminal = False
ok = None
elif job_status == "succeeded":
stage = "ingesting"
label = "Result received; updating ERP evidence"
terminal = False
ok = None
elif job_status == "failed":
stage = "failed"
label = "Failed"
terminal = True
ok = False
elif job_status == "cancelled":
stage = "failed"
label = "Cancelled"
terminal = True
ok = False
else:
stage = "queued"
label = "Queued for workstation"
terminal = False
ok = None
return {
"run_id": int(run.id),
"job_id": int(job.id) if job else None,
"status": run_status,
"job_status": job_status,
"stage": stage,
"label": label,
"terminal": terminal,
"ok": ok,
"workstation": (
str(getattr(workstation, "machine_name", "") or "").strip()
or str(getattr(workstation, "agent_instance_id", "") or "").strip()
or "ERP Local Agent"
) if workstation else "ERP Local Agent",
"company_name": str(run.company_name or run.tally_guid or ""),
"date_from": str(run.date_from or ""),
"date_to": str(run.date_to or ""),
"evidence_rows": int(run.evidence_rows or 0),
"attempts": int(job.attempts or 0) if job else 0,
"max_attempts": int(job.max_attempts or 0) if job else 0,
"error": str(run.error_message or (job.last_error if job else "") or ""),
"created_at": _iso(run.created_at_utc),
"claimed_at": _iso(job.claimed_at_utc) if job else None,
"completed_at": _iso(run.completed_at_utc or (job.completed_at_utc if job else None)),
"elapsed_seconds": _elapsed_seconds(run.created_at_utc),
}
+30 -1
View File
@@ -4,7 +4,7 @@ from datetime import date, timedelta
from urllib.parse import urlencode
from fastapi import APIRouter, Form, Request
from fastapi.responses import RedirectResponse
from fastapi.responses import JSONResponse, RedirectResponse
from sqlalchemy import select
from app.core.db.common import CommonSessionLocal
@@ -20,6 +20,7 @@ from app.modules.accounting.sales_learning_service import (
)
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
@@ -104,6 +105,7 @@ def page(
"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,
},
@@ -191,3 +193,30 @@ def collect(
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()
@@ -17,6 +17,25 @@
</section>
{% if selected_client %}
{% if latest_run %}
<section id="historicalRunProgress"
data-status-url="/tools/accounting/historical-learning/run/{{ latest_run.id }}/status"
class="rounded-2xl border border-blue-200 bg-white p-5 shadow-soft">
<div class="flex flex-wrap items-center justify-between gap-3">
<div>
<div class="text-xs font-semibold uppercase tracking-[0.14em] text-blue-700">Latest Historical Purchase Run</div>
<div id="runProgressLabel" class="mt-1 text-lg font-semibold text-slate-900">{{ latest_run.status|replace('_',' ')|title }}</div>
<div id="runProgressDetail" class="mt-1 text-sm text-slate-500">{{ latest_run.company_name or latest_run.tally_guid }}</div>
</div>
<div id="runProgressMeta" class="text-right text-xs text-slate-500">Run #{{ latest_run.id }}</div>
</div>
<div class="mt-4 h-2 overflow-hidden rounded-full bg-slate-200">
<div id="runProgressBar" class="h-full w-1/3 rounded-full bg-blue-600 transition-all duration-500"></div>
</div>
<div id="runProgressError" class="mt-3 hidden rounded-lg border border-red-200 bg-red-50 p-3 text-sm text-red-800"></div>
</section>
{% endif %}
<section class="rounded-2xl bg-white p-5 shadow-soft">
<h2 class="text-lg font-semibold text-slate-900">Collect Historical Purchase Evidence</h2>
<p class="mt-1 text-sm text-slate-500">Select the connected ERP Local Agent workstation and currently open mapped Tally company. The agent refreshes the selected period read-only before returning neutral ledger-history facts; classification stays on the ERP server.</p>
@@ -59,4 +78,80 @@
<script>
const s=document.getElementById('wsCompany'); if(s){s.addEventListener('change',()=>{const p=s.value.split('|');document.getElementById('workstationId').value=p[0]||'';document.getElementById('tallyGuid').value=p[1]||'';document.getElementById('companyName').value=p.slice(2).join('|')||'';});}
</script>
<script>
(function () {
const box = document.getElementById("historicalRunProgress");
if (!box) return;
const url = box.dataset.statusUrl;
const label = document.getElementById("runProgressLabel");
const detail = document.getElementById("runProgressDetail");
const meta = document.getElementById("runProgressMeta");
const bar = document.getElementById("runProgressBar");
const errorBox = document.getElementById("runProgressError");
let completedReloaded = false;
function humanSeconds(seconds) {
seconds = Number(seconds || 0);
if (seconds < 60) return seconds + " sec";
const mins = Math.floor(seconds / 60);
const secs = seconds % 60;
return mins + " min " + secs + " sec";
}
function paint(data) {
label.textContent = data.label || data.status || "Working";
detail.textContent = [
data.workstation,
data.company_name,
data.date_from && data.date_to ? (data.date_from + " → " + data.date_to) : ""
].filter(Boolean).join(" · ");
meta.textContent =
"Run #" + data.run_id +
" · attempt " + (data.attempts || 0) + "/" + (data.max_attempts || 0) +
" · elapsed " + humanSeconds(data.elapsed_seconds);
if (data.stage === "queued") {
bar.className = "h-full w-1/4 rounded-full bg-blue-500 animate-pulse transition-all duration-500";
} else if (data.stage === "processing") {
bar.className = "h-full w-2/3 rounded-full bg-blue-600 animate-pulse transition-all duration-500";
} else if (data.stage === "ingesting") {
bar.className = "h-full w-5/6 rounded-full bg-indigo-600 animate-pulse transition-all duration-500";
} else if (data.stage === "completed") {
bar.className = "h-full w-full rounded-full bg-emerald-600 transition-all duration-500";
label.textContent = "Completed · " + (data.evidence_rows || 0) + " evidence row(s)";
} else if (data.stage === "failed") {
bar.className = "h-full w-full rounded-full bg-red-600 transition-all duration-500";
}
if (data.error) {
errorBox.textContent = data.error;
errorBox.classList.remove("hidden");
} else {
errorBox.classList.add("hidden");
}
if (data.terminal && !completedReloaded) {
completedReloaded = true;
window.setTimeout(function () { window.location.reload(); }, 1200);
}
}
async function poll() {
try {
const response = await fetch(url, {headers: {"Accept": "application/json"}, cache: "no-store"});
if (response.ok) {
const data = await response.json();
if (data.ok) paint(data);
}
} catch (e) {
// Keep the last visible state; next poll will retry.
}
if (!completedReloaded) window.setTimeout(poll, 3000);
}
poll();
})();
</script>
{% endblock %}
@@ -32,6 +32,26 @@
{% endfor %}
</section>
{% if latest_run %}
<section id="historicalRunProgress"
data-status-url="/tools/accounting/sales-learning/run/{{ latest_run.id }}/status"
data-run-status="{{ latest_run.status }}"
class="rounded-2xl border border-blue-200 bg-white p-5 shadow-soft">
<div class="flex flex-wrap items-center justify-between gap-3">
<div>
<div class="text-xs font-semibold uppercase tracking-[0.14em] text-blue-700">Latest Historical Sales Run</div>
<div id="runProgressLabel" class="mt-1 text-lg font-semibold text-slate-900">{{ latest_run.status|replace('_',' ')|title }}</div>
<div id="runProgressDetail" class="mt-1 text-sm text-slate-500">{{ latest_run.company_name or latest_run.tally_guid }}</div>
</div>
<div id="runProgressMeta" class="text-right text-xs text-slate-500">Run #{{ latest_run.id }}</div>
</div>
<div class="mt-4 h-2 overflow-hidden rounded-full bg-slate-200">
<div id="runProgressBar" class="h-full w-1/3 rounded-full bg-blue-600 transition-all duration-500"></div>
</div>
<div id="runProgressError" class="mt-3 hidden rounded-lg border border-red-200 bg-red-50 p-3 text-sm text-red-800"></div>
</section>
{% endif %}
<section class="rounded-2xl bg-white p-5 shadow-soft">
<h2 class="font-semibold">Collect Historical Sales Evidence from Tally</h2>
<p class="mt-1 text-sm text-slate-500">Read-only. The selected workstation refreshes Sales vouchers from the currently open mapped Tally company for the chosen period, then stores neutral historical evidence. It does not create or alter Tally vouchers.</p>
@@ -109,4 +129,80 @@
syncSelection();
})();
</script>
<script>
(function () {
const box = document.getElementById("historicalRunProgress");
if (!box) return;
const url = box.dataset.statusUrl;
const label = document.getElementById("runProgressLabel");
const detail = document.getElementById("runProgressDetail");
const meta = document.getElementById("runProgressMeta");
const bar = document.getElementById("runProgressBar");
const errorBox = document.getElementById("runProgressError");
let completedReloaded = false;
function humanSeconds(seconds) {
seconds = Number(seconds || 0);
if (seconds < 60) return seconds + " sec";
const mins = Math.floor(seconds / 60);
const secs = seconds % 60;
return mins + " min " + secs + " sec";
}
function paint(data) {
label.textContent = data.label || data.status || "Working";
detail.textContent = [
data.workstation,
data.company_name,
data.date_from && data.date_to ? (data.date_from + " → " + data.date_to) : ""
].filter(Boolean).join(" · ");
meta.textContent =
"Run #" + data.run_id +
" · attempt " + (data.attempts || 0) + "/" + (data.max_attempts || 0) +
" · elapsed " + humanSeconds(data.elapsed_seconds);
if (data.stage === "queued") {
bar.className = "h-full w-1/4 rounded-full bg-blue-500 animate-pulse transition-all duration-500";
} else if (data.stage === "processing") {
bar.className = "h-full w-2/3 rounded-full bg-blue-600 animate-pulse transition-all duration-500";
} else if (data.stage === "ingesting") {
bar.className = "h-full w-5/6 rounded-full bg-indigo-600 animate-pulse transition-all duration-500";
} else if (data.stage === "completed") {
bar.className = "h-full w-full rounded-full bg-emerald-600 transition-all duration-500";
label.textContent = "Completed · " + (data.evidence_rows || 0) + " evidence row(s)";
} else if (data.stage === "failed") {
bar.className = "h-full w-full rounded-full bg-red-600 transition-all duration-500";
}
if (data.error) {
errorBox.textContent = data.error;
errorBox.classList.remove("hidden");
} else {
errorBox.classList.add("hidden");
}
if (data.terminal && !completedReloaded) {
completedReloaded = true;
window.setTimeout(function () { window.location.reload(); }, 1200);
}
}
async function poll() {
try {
const response = await fetch(url, {headers: {"Accept": "application/json"}, cache: "no-store"});
if (response.ok) {
const data = await response.json();
if (data.ok) paint(data);
}
} catch (e) {
// Keep the last visible state; next poll will retry.
}
if (!completedReloaded) window.setTimeout(poll, 3000);
}
poll();
})();
</script>
{% endblock %}