108 lines
3.2 KiB
Python
108 lines
3.2 KiB
Python
"""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)
|