Files
GarminHealthLab/backend/services/jobs.py
ericwyuan 6dd070ec9b fix(ai): subject 里塞了行数,每轮询一次就新建一个任务
生产上 trends 队列里积了 36 个任务,subject 是 2026-09-01:1033、:1039、
:1044……一路涨。这台账号当时正在补历史,get_summary 的行数每隔几分钟就变,
而我把 len(rows) 写进了 subject——subject 同时是缓存键和任务队列的键,一变
就是一条全新的任务,轮询几次就刷出十几条。

subject 该回答的是「这条解读是关于什么的」,不是「当时有多少行数据」。
数据变化本来就由 fingerprint 负责。

- trends 的 subject 改成快照日期;sleep 用配置的窗口常量而不是实际夜数
  (缺一晚也不该换键);challenges 用固定键
- 加了不变量测试:补一天历史数据后 subject 不许变;任何 subject 段都不许
  长得像行数

顺带加一层兜底 jobs.supersede():单实例 scope 只该有一个在跑的 subject,
队列里同 kind 的其它 pending 任务是关于已经不存在的快照的,跑完也没人看。
per_item 的 daily / activity 不受影响——它们本来就一天一条、一次运动一条。
兜底不是机制,机制是 subject 稳定;它存在只是因为这次 subject 不稳定,而
36 条任务堆在那里之前没人发现。

顺带按要求把 AiPanel 改成默认精简:只显示标题、来源和一句话结论,点「展开
详细」才出要点/建议/依据,可再收起——和今日晨报卡片一致。这些面板压在本来
就很密的图表页上面,全部默认展开会把真正的数据一次性挤到屏幕外。

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

291 lines
10 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)
# 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))
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'"
)