ai insight rework #1
This commit is contained in:
1 parent
0d960c8153
commit
9d76f54d33
13 files changed
+730
-299
No files matched your search
@@ -0,0 +1,107 @@
|
||||
"""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)
|
||||
Reference in new issue
Block a user