Files
GarminHealthLab/backend/services/jobs.py
ericwyuan a2d6a5c57f fix(ai): 队列会把共用的网关打到 502,加并发上限和间隔
排查生产上一直生成不出来,发现网关在返 502。上去看:机器好好的、systemd
说 active、5100 端口在监听——但它是 `gunicorn -w 1 --threads 4`,全部并发
就四个,而且 fam-edge 和摄像头项目也在用同一个。

我们这边一个请求占一个线程 2~5 分钟,NVIDIA 链重试起来最坏十七分钟(它自己
README 已知问题 #3)。而我写的 worker 是跑完一个立刻拉下一个,同步后还有八
个 scope 排队——等于拿满线程不撒手。这个 502 大概率是我打出来的,而且顺带
把另外两个项目也打下线了。

- 并发按整个部署计算,不是每个 gunicorn worker 一个:claim 前先数全局
  running(两个 worker 各跑「一个」就是两个并发)
- 每跑完一个任务停 20 秒,不只是空闲时才停
- 两个都可用环境变量调,注释里写清楚调大的代价是什么

补齐的历史数据晚二十分钟到没有任何人受影响;网关不响应是三个项目一起受影响。

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-01 15:46:42 +08:00

317 lines
12 KiB
Python

"""
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)
# The gateway is not ours alone. It runs `gunicorn -w 1 --threads 4` on the
# Oracle box and is shared with fam-edge and the camera project, so it can
# serve four requests at a time in total — while one of ours occupies a thread
# for two to five minutes, and up to seventeen when its NVIDIA chain retries
# (its own README, known issue #3).
#
# So this consumer deliberately runs one job at a time across the whole
# deployment, not one per Gunicorn worker, and leaves a gap between jobs. A
# backfill of eight scopes is background work; taking the shared box down to
# finish it sooner is not a trade worth making — and a 502 there is a 502 for
# the other two projects too.
MAX_CONCURRENT = int(os.environ.get("AI_JOB_CONCURRENCY") or 1)
GAP_SECONDS = float(os.environ.get("AI_JOB_GAP_SECONDS") or 20)
# 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 supersede(user_id, kind, subject):
"""Drop queued work of the same kind for a different subject.
A single-entry scope has one live subject; anything else queued under that
kind is about a snapshot that no longer exists, and running it would spend
a gateway call on an answer nothing will read.
This is a safety net, not the mechanism: subjects are supposed to be stable
(see scopes.py). It exists because they were not — a row count in the key
made every poll mint a new `trends` job, and production had 36 of them
queued before anyone noticed.
"""
execute(
"DELETE FROM ai_jobs WHERE user_id = ? AND kind = ? AND subject != ? "
"AND status = 'pending'",
[user_id, kind, subject],
)
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))
# Concurrency is counted across the deployment, not per process: two
# Gunicorn workers each running "one job" is two concurrent calls into a
# gateway that has four threads for three projects.
running = query_one(
"SELECT COUNT(*) AS n FROM ai_jobs WHERE status = 'running' "
"AND claimed_at IS NOT NULL AND claimed_at >= ?",
[cutoff],
)
if (running or {}).get("n", 0) >= MAX_CONCURRENT:
return None
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:
# A gap after each job, not just when idle. Running them
# back-to-back is what saturates the shared gateway: a backfill
# finishing twenty minutes later costs nobody anything, a gateway
# that stops answering costs all three projects.
time.sleep(GAP_SECONDS if run_once() else 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'"
)