Files
foodlinkk-command-center/tools-api/app/export_intel_sync_progress.py
T

187 lines
5.3 KiB
Python
Raw Normal View History

2026-07-30 09:33:03 +00:00
"""File-backed sync progress for Export Intel background and UI jobs."""
from __future__ import annotations
import json
import os
import time
from pathlib import Path
from typing import Any, Optional
STATUS_PATH = Path("/tmp/export-intel-sync-status.json")
LOG_PATH = Path("/tmp/export-intel-sync-progress.jsonl")
LOCK_PATH = Path("/tmp/export-intel-sync.lock")
_active_job_id: Optional[str] = None
def _now() -> float:
return time.time()
def _read_json(path: Path, default: Any) -> Any:
try:
if path.exists():
return json.loads(path.read_text())
except Exception:
pass
return default
def _write_json(path: Path, data: Any) -> None:
path.write_text(json.dumps(data, indent=2, default=str))
def log_event(event: str, **data: Any) -> None:
row = {"ts": _now(), "event": event, **data}
with LOG_PATH.open("a") as fh:
fh.write(json.dumps(row, default=str) + "\n")
def set_active_job(job_id: Optional[str]) -> None:
global _active_job_id
_active_job_id = job_id
def get_active_job() -> Optional[str]:
return _active_job_id
def _load_status() -> dict[str, Any]:
data = _read_json(STATUS_PATH, {"jobs": {}, "history": []})
if "jobs" not in data:
data["jobs"] = {}
if "history" not in data:
data["history"] = []
return data
def _save_status(data: dict[str, Any]) -> None:
_write_json(STATUS_PATH, data)
def _touch_lock(job_id: str, label: str) -> None:
LOCK_PATH.write_text(json.dumps({"pid": os.getpid(), "job_id": job_id, "label": label, "ts": _now()}))
def _clear_lock() -> None:
LOCK_PATH.unlink(missing_ok=True)
def start_job(job_type: str, label: str, total: int = 0, meta: Optional[dict[str, Any]] = None) -> str:
job_id = f"{job_type}-{_now():.0f}"
status = _load_status()
job = {
"id": job_id,
"type": job_type,
"label": label,
"status": "running",
"started_at": _now(),
"updated_at": _now(),
"current": 0,
"total": total,
"percent": 0,
"country": None,
"step": None,
"message": label,
"meta": meta or {},
"error": None,
}
status["jobs"][job_id] = job
_save_status(status)
_touch_lock(job_id, label)
set_active_job(job_id)
log_event("job_start", job_id=job_id, type=job_type, label=label, total=total)
return job_id
def update_job(job_id: str, **fields: Any) -> None:
status = _load_status()
job = status["jobs"].get(job_id)
if not job:
return
job.update(fields)
job["updated_at"] = _now()
cur = int(job.get("current") or 0)
tot = int(job.get("total") or 0)
if tot > 0:
job["percent"] = min(100, round(cur * 100 / tot))
status["jobs"][job_id] = job
_save_status(status)
if LOCK_PATH.exists():
_touch_lock(job_id, job.get("label") or job_id)
def job_step(
job_id: str,
*,
message: str,
country: Optional[str] = None,
step: Optional[str] = None,
current: Optional[int] = None,
total: Optional[int] = None,
) -> None:
fields: dict[str, Any] = {"message": message}
if country is not None:
fields["country"] = country
if step is not None:
fields["step"] = step
if current is not None:
fields["current"] = current
if total is not None:
fields["total"] = total
update_job(job_id, **fields)
log_event("step", job_id=job_id, message=message, country=country, step=step, current=current, total=total)
def finish_job(job_id: str, result: Optional[dict[str, Any]] = None, error: Optional[str] = None) -> None:
status = _load_status()
job = status["jobs"].get(job_id)
if not job:
return
job["status"] = "error" if error else "completed"
job["updated_at"] = _now()
job["finished_at"] = _now()
job["percent"] = 100 if not error else job.get("percent", 0)
job["error"] = error
if result is not None:
job["result"] = result
status["jobs"][job_id] = job
status["history"] = ([job_id] + [h for h in status.get("history", []) if h != job_id])[:20]
_save_status(status)
log_event("job_finish", job_id=job_id, status=job["status"], error=error)
if get_active_job() == job_id:
set_active_job(None)
running = [j for j in status["jobs"].values() if j.get("status") == "running"]
if not running:
_clear_lock()
def tail_log(limit: int = 40) -> list[dict[str, Any]]:
if not LOG_PATH.exists():
return []
lines = LOG_PATH.read_text().strip().splitlines()
out: list[dict[str, Any]] = []
for line in lines[-limit:]:
try:
out.append(json.loads(line))
except Exception:
continue
return out
def get_monitor_snapshot(db_counts: Optional[dict[str, Any]] = None) -> dict[str, Any]:
status = _load_status()
jobs = status.get("jobs", {})
active = [j for j in jobs.values() if j.get("status") == "running"]
recent = [jobs[jid] for jid in status.get("history", [])[:8] if jid in jobs and jobs[jid].get("status") != "running"]
lock = _read_json(LOCK_PATH, None) if LOCK_PATH.exists() else None
return {
"active_jobs": active,
"recent_jobs": recent,
"events": tail_log(50),
"lock": lock,
"any_running": bool(active) or bool(lock),
"coverage": db_counts or {},
"updated_at": _now(),
}