"""Phase 4 async job runner for POST /ai-insights/generate/ (202 + polling). Jobs persist as one JSON file each under `logs/insight_jobs/` so any gunicorn worker can answer the poll (ponytail: single-host only — move to Redis/DB queue if API ever spans multiple machines). Only the owning worker writes a job file, so no cross-process locking is needed; writes are atomic (tmp + os.replace). """ from __future__ import annotations import json import logging import os import threading import time import uuid from pathlib import Path from typing import Any from django.conf import settings logger = logging.getLogger(__name__) _JOBS_DIR = Path(settings.BASE_DIR) / "logs" / "insight_jobs" _MAX_AGE_SECONDS = 24 * 3600 def _job_path(job_id: str) -> Path: return _JOBS_DIR / f"{job_id}.json" def _write(job_id: str, job: dict[str, Any]) -> None: tmp = _job_path(job_id).with_suffix(".tmp") tmp.write_text(json.dumps(job, ensure_ascii=False, default=str), encoding="utf-8") os.replace(tmp, _job_path(job_id)) def _prune() -> None: """Drop job files older than 24h so the dir stays bounded.""" try: cutoff = time.time() - _MAX_AGE_SECONDS for f in _JOBS_DIR.glob("*.json"): if f.stat().st_mtime < cutoff: f.unlink(missing_ok=True) except OSError: logger.debug("insight_jobs prune failed", exc_info=True) def create_job(params: dict[str, Any]) -> str: """Persist job record and start worker thread. Returns job_id for polling.""" job_id = str(uuid.uuid4()) _JOBS_DIR.mkdir(parents=True, exist_ok=True) _prune() _write( job_id, { "id": job_id, "status": "queued", "stage": "queued", "result": None, "error": None, }, ) threading.Thread(target=_run, args=(job_id, params), daemon=True).start() return job_id def get_job(job_id: str) -> dict[str, Any] | None: """Read job from shared dir — works from any worker process.""" try: raw = _job_path(job_id).read_text(encoding="utf-8") except FileNotFoundError: return None except OSError: logger.warning("insight_jobs read failed for %s", job_id, exc_info=True) return None try: job = json.loads(raw) except json.JSONDecodeError: return None return job if isinstance(job, dict) else None def _set(job_id: str, **fields: Any) -> None: job = get_job(job_id) if job is None: return job.update(fields) try: _write(job_id, job) except OSError: logger.warning("insight_jobs write failed for %s", job_id, exc_info=True) def _run(job_id: str, params: dict[str, Any]) -> None: from apps.operations.services.insight_service import generate_insight def stage_cb(stage: str) -> None: _set(job_id, status="running", stage=stage) _set(job_id, status="running", stage="starting") try: result = generate_insight(**params, stage_cb=stage_cb) except Exception as exc: # noqa: BLE001 - surfaced to client via job status logger.exception("insight job %s failed", job_id) _set(job_id, status="error", stage="error", error=str(exc)) return _set(job_id, status="done", stage="done", result=result)