""" The coach's producer/consumer queue. Producing one insight costs 40s to several minutes against the gateway, so it can never happen inside a request. Screens *enqueue*; a worker thread consumes. Priority is the whole point of the queue rather than a plain background thread: after a sync the backfill enqueues every scope at low priority, and those jobs may take half an hour to work through — but the moment the user opens a screen, that screen's job is promoted to the front and runs next. What they are looking at is always what the queue is working on. The queue lives in the database, not in memory, because Gunicorn runs several workers: a job is claimed with a holder id and re-read to confirm, the same way `scheduler.py` claims its tick, so exactly one worker runs a given job. """ import datetime import hashlib import os import threading import time from config import DB_TYPE from db import execute, query_one, query_all # Lower runs first. PRIORITY_INTERACTIVE = 0 # a screen the user has open right now PRIORITY_PREFETCH = 10 # backfill after a sync # A claim older than this is treated as abandoned: the worker holding it died # mid-generation, and without expiry that job would never run again. CLAIM_TIMEOUT_SECONDS = int(os.environ.get("AI_JOB_CLAIM_TIMEOUT") or 1800) # Generations are slow, not frequent; polling this often costs nothing and # keeps an interactive job's wait to a couple of seconds. POLL_SECONDS = float(os.environ.get("AI_JOB_POLL_SECONDS") or 2) MAX_ATTEMPTS = int(os.environ.get("AI_JOB_MAX_ATTEMPTS") or 3) # How long a job that exhausted its attempts stays given up on before it is # tried again. Without this a transient upstream outage is permanent: three # quick failures while the gateway is unreachable would retire that screen's # insight until its underlying data happened to change, which for a screen the # user is not syncing could be days. FAILED_RETRY_SECONDS = int(os.environ.get("AI_JOB_RETRY_AFTER") or 1800) ENABLED = (os.environ.get("AI_JOBS") or "true").lower() not in ("0", "false", "no") _started = False _start_lock = threading.Lock() # Set by analysis.py at import time. Injected rather than imported so this # module stays free of the feature logic it schedules — and so the two do not # import each other in a cycle. _runner = None def set_runner(fn): """Register `fn(user_id, kind, subject) -> None`, called for each job.""" global _runner _runner = fn def _now(): return datetime.datetime.utcnow() def _iso(dt): return dt.isoformat(timespec="seconds") def _parse(value): if not value: return None try: return datetime.datetime.fromisoformat(str(value).replace(" ", "T")) except ValueError: return None def job_id(user_id, kind, subject): return hashlib.sha256( f"{user_id}|{kind}|{subject}".encode("utf-8") ).hexdigest()[:64] def enqueue(user_id, kind, subject, fingerprint=None, priority=PRIORITY_PREFETCH): """Queue one generation, or promote it if it is already queued. Returns the job's current status. Idempotent by design: the screen polls every few seconds while it waits, and every one of those polls calls this. A finished job is re-queued only when the data it was derived from has changed — that is what `fingerprint` is for, and it is why a poll on unchanged data does not restart the work that just completed. """ jid = job_id(user_id, kind, subject) now = _iso(_now()) row = query_one("SELECT * FROM ai_jobs WHERE id = ?", [jid]) if row: if row["status"] == "running": claimed = _parse(row.get("claimed_at")) if claimed and (_now() - claimed).total_seconds() < CLAIM_TIMEOUT_SECONDS: # Already being generated. Promoting it now would not make the # in-flight call any faster. return "running" stale = fingerprint and row.get("fingerprint") != fingerprint if row["status"] == "done" and not stale: return "done" if row["status"] == "failed" and row["attempts"] >= MAX_ATTEMPTS and not stale: gave_up = _parse(row.get("updated_at")) if gave_up and (_now() - gave_up).total_seconds() < FAILED_RETRY_SECONDS: return "failed" # Past the cooldown: reset the attempt count so the outage that # exhausted it does not count against the retry. execute( "UPDATE ai_jobs SET status = 'pending', attempts = 0, error = NULL, " "holder = NULL, claimed_at = NULL, priority = ?, updated_at = ? " "WHERE id = ?", [min(priority, row["priority"]), now, jid], ) return "pending" # Promote (never demote): a screen the user just opened must not be # pushed back by the prefetch entry that was already sitting there. execute( "UPDATE ai_jobs SET status = 'pending', priority = ?, " "fingerprint = ?, holder = NULL, claimed_at = NULL, " "attempts = ?, updated_at = ? WHERE id = ?", [ min(priority, row["priority"]), fingerprint or row.get("fingerprint"), 0 if stale else row["attempts"], now, jid, ], ) return "pending" cols = ["id", "user_id", "kind", "subject", "fingerprint", "priority", "status", "attempts", "created_at", "updated_at"] execute( f"INSERT INTO ai_jobs ({', '.join(cols)}) " f"VALUES ({', '.join(['?'] * len(cols))})", [jid, user_id, kind, subject, fingerprint, priority, "pending", 0, now, now], ) return "pending" def status_of(user_id, kind, subject): row = query_one("SELECT * FROM ai_jobs WHERE id = ?", [job_id(user_id, kind, subject)]) if not row: return None return { "status": row["status"], "priority": row["priority"], "attempts": row["attempts"], "error": row.get("error"), "updatedAt": row.get("updated_at"), } def pending_count(user_id=None): sql = "SELECT COUNT(*) AS n FROM ai_jobs WHERE status IN ('pending', 'running')" params = [] if user_id: sql += " AND user_id = ?" params.append(user_id) row = query_one(sql, params) return (row or {}).get("n") or 0 def _claim_next(): """Take the highest-priority runnable job, or None. Ordered by priority then age so the interactive job wins and, among equals, the one that has waited longest goes first. """ cutoff = _iso(_now() - datetime.timedelta(seconds=CLAIM_TIMEOUT_SECONDS)) rows = query_all( "SELECT * FROM ai_jobs WHERE status = 'pending' " "OR (status = 'running' AND (claimed_at IS NULL OR claimed_at < ?)) " "ORDER BY priority ASC, created_at ASC", [cutoff], ) holder = f"{os.getpid()}-{threading.get_ident()}" for row in rows: if row["attempts"] >= MAX_ATTEMPTS: continue execute( "UPDATE ai_jobs SET status = 'running', holder = ?, claimed_at = ?, " "attempts = ?, updated_at = ? WHERE id = ? AND status = ?", [holder, _iso(_now()), row["attempts"] + 1, _iso(_now()), row["id"], row["status"]], ) # Re-read: another worker may have claimed it between the SELECT and # the UPDATE, in which case its holder is the one now recorded. check = query_one("SELECT holder FROM ai_jobs WHERE id = ?", [row["id"]]) if check and check.get("holder") == holder: return row return None def _finish(jid, error=None): execute( "UPDATE ai_jobs SET status = ?, error = ?, holder = NULL, " "claimed_at = NULL, updated_at = ? WHERE id = ?", ["failed" if error else "done", (error or "")[:500] if error else None, _iso(_now()), jid], ) def run_once(): """Claim and run one job. Returns True when something was run.""" if _runner is None: return False row = _claim_next() if not row: return False try: _runner(row["user_id"], row["kind"], row["subject"]) except Exception as e: # noqa: BLE001 - one bad job must not stop the queue _finish(row["id"], f"{type(e).__name__}: {e}") print(f"[ai-jobs] {row['kind']}:{row['subject']} failed: {e}") return True _finish(row["id"]) return True def _loop(): while True: try: # Straight on to the next job when one was just run: after a # sync there is a whole backfill waiting, and sleeping between # each would add hours to it for no reason. if not run_once(): time.sleep(POLL_SECONDS) except Exception as e: # noqa: BLE001 - the loop must outlive any failure print(f"[ai-jobs] worker error: {e}") time.sleep(POLL_SECONDS) def start(): """Start one consumer per process.""" global _started if not ENABLED: print("[ai-jobs] disabled by AI_JOBS") return with _start_lock: if _started: return _started = True threading.Thread(target=_loop, daemon=True, name="ai-jobs").start() print("[ai-jobs] worker started") def reset_stale_claims(): """Release jobs a previous process was running when it stopped. Without this they sit in `running` until the claim expires, which for the screen waiting on one looks exactly like a generation that never finishes. """ execute( "UPDATE ai_jobs SET status = 'pending', holder = NULL, claimed_at = NULL " "WHERE status = 'running'" )