"""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(), }