187 lines
5.3 KiB
Python
187 lines
5.3 KiB
Python
|
|
"""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(),
|
||
|
|
}
|