From c53ef6a416b343bf7a5e770f16690aaf5e3360cd Mon Sep 17 00:00:00 2001 From: Color2333 <1552429809@qq.com> Date: Sat, 18 Jul 2026 21:30:21 +0800 Subject: [PATCH] =?UTF-8?q?chore:=20=E5=88=A0=E9=99=A4=E5=BF=83=E8=B7=B3?= =?UTF-8?q?=E9=82=AE=E4=BB=B6=E5=91=8A=E8=AD=A6=20+=20docker=20healthcheck?= =?UTF-8?q?=EF=BC=8C=E5=8F=AA=E7=95=99=20status=20=E9=A1=B5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 按需求简化心跳机制:worker 卡死不再自动重启/发邮件,靠 status 页 (/system/worker 端点 + Operations 面板)发现 + 人工介入。 删除: - apps/worker/main.py:heartbeat_alert_job + _send_alert + _last_alerts + _ALERT_DEDUP_SECONDS(邮件告警去重)+ scheduler 注册。保留 _write_heartbeat/_read_heartbeat/_HEALTH_FILE(status 页读)。 - docker-compose.yml:worker healthcheck 段(不再自动重启卡死 worker)。 - scripts/worker_healthcheck.py:整个文件删除(healthcheck 脚本)。 - tests/test_repositories.py:TestWorkerAlertDedup 2 测试(测已删的 _send_alert)。 调整: - frontend Operations.tsx:过期提示从"worker 每 10min 自检发告警邮件" 改"请检查 worker 容器状态与日志,必要时手动重启"。 保留: - status 页(/system/worker 端点 + Operations 心跳面板)读共享卷心跳 展示 worker 健康时效。 - last_error 面板(topic 抓取错误展示 + 重新抓取)与心跳机制独立,保留。 - worker 写心跳逻辑(各 job 成功写心跳,status 页消费)。 验证:ruff + tsc + 61 passed(删 2 个告警去重测试)。无残留引用。 --- apps/worker/main.py | 102 ++---------------------------- docker-compose.yml | 8 --- frontend/src/pages/Operations.tsx | 2 +- scripts/worker_healthcheck.py | 36 ----------- tests/test_repositories.py | 52 --------------- 5 files changed, 8 insertions(+), 192 deletions(-) delete mode 100644 scripts/worker_healthcheck.py diff --git a/apps/worker/main.py b/apps/worker/main.py index 584f85c..15fb196 100644 --- a/apps/worker/main.py +++ b/apps/worker/main.py @@ -37,20 +37,17 @@ logger = logging.getLogger(__name__) # 心跳改写共享卷 pm_data(/app/data),backend 也能读同一文件暴露 worker 状态。 -# 此前写 /tmp(容器内,后端读不到),可观测性端点无法查询 worker 健康。 +# status 页(/system/worker 端点 + Operations 面板)读此文件展示 worker 健康。 _HEALTH_FILE = Path("/app/data/worker_heartbeat.json") -# 心跳健康判定:最近一次心跳距现在超过此秒数视为不健康(捕获 worker 卡死/全部任务失败) -_HEARTBEAT_STALE_SECONDS = 1200 # 20 分钟(cron job 最小间隔 30min,留足缓冲) -# 告警去重:同一告警类型在此秒数内不重复发邮件(防刷屏) -_ALERT_DEDUP_SECONDS = 3600 # 1 小时 +# 心跳健康判定:最近一次心跳距现在超过此秒数视为不健康(status 页展示用) +_HEARTBEAT_STALE_SECONDS = 1200 # 20 分钟 def _write_heartbeat(error: str | None = None) -> None: - """写入心跳文件供外部健康检查(High 2e:记录最近一次错误,不再掩盖故障)。 + """写入心跳文件供 status 页查询(记录最近一次错误,不再掩盖故障)。 - 此前无条件写时间戳,healthcheck 仅 test -f → 即使所有 job 失败 worker 仍判健康。 - 现写入 JSON {ts, error}:健康检查读 ts 判定时效,error 字段记录最近致命错误。 - job 全部失败时不写心跳(让心跳自然过期 → healthcheck 反映故障)。 + 写入 JSON {ts, error}:status 页读 ts 判定时效,error 字段记录最近致命错误。 + job 全部失败时不写心跳(让心跳自然过期 → status 页反映故障)。 """ import json @@ -61,7 +58,7 @@ def _write_heartbeat(error: str | None = None) -> None: def _read_heartbeat() -> dict | None: - """读共享卷心跳文件,供自检告警 + 端点查询。文件缺失/损坏返回 None。""" + """读共享卷心跳文件,供 status 端点查询。文件缺失/损坏返回 None。""" import json try: @@ -70,82 +67,6 @@ def _read_heartbeat() -> dict | None: return None -# 告警去重状态:{alert_key: last_alert_ts}(进程内,重启后重置——可接受) -_last_alerts: dict[str, float] = {} - - -def _send_alert(subject: str, html: str, alert_key: str) -> None: - """发告警邮件给 notify_default_to,带 1h 去重。SMTP 未配置时静默跳过。""" - now = time.time() - last = _last_alerts.get(alert_key, 0) - if now - last < _ALERT_DEDUP_SECONDS: - logger.debug( - "告警 %s 去重中(距上次 %.0fs < %ds),跳过", - alert_key, - now - last, - _ALERT_DEDUP_SECONDS, - ) - return - from packages.config import get_settings - from packages.integrations.notifier import NotificationService - - recipient = get_settings().notify_default_to - if not recipient: - logger.debug("notify_default_to 未配置,跳过告警 %s", alert_key) - return - ok = NotificationService().send_email_html(recipient, subject, html) - if ok: - _last_alerts[alert_key] = now - logger.info("告警邮件已发送: %s -> %s", alert_key, recipient) - else: - logger.warning("告警邮件发送失败(SMTP 未配置或出错): %s", alert_key) - - -def heartbeat_alert_job() -> None: - """每 10min 自检:心跳过期 + 主题抓取错误 → 发告警邮件(可观测性闭环)。 - - 心跳过期说明 worker 卡死或全部 job 失败;主题 last_error 说明某主题抓取失败。 - 两者此前只记日志无人知,现在发邮件给 notify_default_to(带 1h 去重防刷屏)。 - """ - import html as html_lib - - # 1. 心跳过期检查 - hb = _read_heartbeat() - if hb is None or (time.time() - float(hb.get("ts", 0))) > _HEARTBEAT_STALE_SECONDS: - age = ( - "未知(文件缺失/损坏)" - if hb is None - else f"{int(time.time() - float(hb.get('ts', 0)))}s" - ) - body = ( - f"

⚠️ Worker 心跳过期

" - f"

心跳距今 {age}(阈值 {_HEARTBEAT_STALE_SECONDS}s)。

" - f"

可能原因:worker 卡死、全部 job 失败、或容器异常。

" - f"

最近错误:{(hb or {}).get('error') or 'N/A'}

" - f"

请检查 worker 容器状态与日志。

" - ) - _send_alert("[PaperMind] Worker 心跳过期告警", body, "heartbeat_stale") - # 2. 主题抓取错误检查 - try: - with session_scope() as session: - topics = TopicRepository(session).list_topics(enabled_only=True) - errored = [t for t in topics if t.last_error] - if errored: - rows = "".join( - f"{html_lib.escape(t.name)}" - f"{html_lib.escape((t.last_error or '')[:200])}" - for t in errored - ) - body = ( - f"

⚠️ 主题抓取错误({len(errored)} 个)

" - f"" - f"{rows}
主题最近错误
" - ) - _send_alert(f"[PaperMind] {len(errored)} 个主题抓取失败", body, "topic_errors") - except Exception: - logger.exception("heartbeat_alert_job 检查主题错误失败") - - def _update_topic_run_status(topic_id: str, *, error: str | None) -> None: """记录主题抓取的最近运行时间与错误(Critical #4:失败可查可补抓)。 @@ -394,15 +315,6 @@ def run_worker() -> None: ) logger.info("✅ 已添加:每周图谱维护任务(UTC 周日 22:00)") - # 可观测性:心跳过期 + 主题抓取错误告警(每 10min 自检发邮件) - scheduler.add_job( - heartbeat_alert_job, - trigger=CronTrigger(minute="*/10"), - id="heartbeat_alert", - **_job_kwargs, - ) - logger.info("✅ 已添加:心跳告警自检任务(每 10min)") - # 优雅关闭(High 3f:等待进行中任务跑完,避免已下载 PDF 未 set_pdf_path # 的中间态丢失;wait=True + 60s 超时兜底) def _graceful_stop(*_: object) -> None: diff --git a/docker-compose.yml b/docker-compose.yml index 1e40828..b6e6298 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -79,14 +79,6 @@ services: condition: service_started postgres: condition: service_healthy - healthcheck: - # High 2e:改判时效而非仅文件存在——所有 job 失败时心跳不再写,文件过期即不健康。 - # 心跳超过 20 分钟(1200s)视为不健康(worker 卡死或全部任务失败) - test: ["CMD", "python", "-m", "scripts.worker_healthcheck"] - interval: 30s - timeout: 5s - start_period: 40s - retries: 3 deploy: resources: limits: diff --git a/frontend/src/pages/Operations.tsx b/frontend/src/pages/Operations.tsx index bb529af..af64e9c 100644 --- a/frontend/src/pages/Operations.tsx +++ b/frontend/src/pages/Operations.tsx @@ -394,7 +394,7 @@ export default function Operations() {

心跳过期({heartbeat.age_seconds}s > {staleThreshold}s)

worker 可能卡死或全部任务失败。最近错误:{heartbeat.error || "N/A"}

-

worker 每 10min 自检发告警邮件给 notify_default_to。

+

请检查 worker 容器状态与日志,必要时手动重启。

) : ( diff --git a/scripts/worker_healthcheck.py b/scripts/worker_healthcheck.py deleted file mode 100644 index 39af10a..0000000 --- a/scripts/worker_healthcheck.py +++ /dev/null @@ -1,36 +0,0 @@ -#!/usr/bin/env python3 -"""Worker 健康检查脚本(High 2e + 可观测性:改读共享卷)。 - -读取 /app/data/worker_heartbeat.json(pm_data 共享卷,backend 也可读), -判定心跳时效: -- 文件不存在 / 解析失败 → 不健康(exit 1) -- ts 距今超过 1200 秒 → 不健康(worker 卡死或全部任务失败,心跳已过期) -- 否则健康(exit 0) - -此前 healthcheck 仅 test -f 文件存在 → 即使所有 job 失败 worker 仍判健康。 -心跳改共享卷后,backend /system/worker 端点也能读同一文件暴露 worker 状态。 -""" - -from __future__ import annotations - -import json -import sys -import time -from pathlib import Path - -HEALTH_FILE = Path("/app/data/worker_heartbeat.json") -STALE_SECONDS = 1200 # 20 分钟 - - -def main() -> int: - try: - data = json.loads(HEALTH_FILE.read_text()) - ts = float(data.get("ts", 0)) - except (OSError, ValueError, TypeError): - # 文件不存在或损坏 → 视为不健康 - return 1 - return 0 if (time.time() - ts) < STALE_SECONDS else 1 - - -if __name__ == "__main__": - sys.exit(main()) diff --git a/tests/test_repositories.py b/tests/test_repositories.py index c7eadd6..53d9181 100644 --- a/tests/test_repositories.py +++ b/tests/test_repositories.py @@ -426,58 +426,6 @@ def test_link_category_isolation(self, db_session): assert t_ai.id != t_lg.id -class TestWorkerAlertDedup: - """Worker 告警去重逻辑测试(可观测性:同错误 1h 内不重复发邮件)""" - - def test_send_alert_dedup_within_window(self, monkeypatch): - from apps.worker import main as wm - - # 重置去重状态 - wm._last_alerts.clear() - sent: list[str] = [] - - def fake_send(self, recipient, subject, html): - sent.append(subject) - return True - - monkeypatch.setattr( - "packages.integrations.notifier.NotificationService.send_email_html", fake_send - ) - monkeypatch.setattr( - "packages.config.get_settings", - lambda: type("S", (), {"notify_default_to": "test@example.com"})(), - ) - - wm._send_alert("[PaperMind] Test Alert", "

x

", "test_key") - assert len(sent) == 1, "首次告警应发送" - # 窗口内再次告警 → 去重,不发送 - wm._send_alert("[PaperMind] Test Alert", "

x

", "test_key") - assert len(sent) == 1, "去重窗口内不应重复发送" - wm._last_alerts.clear() - - def test_send_alert_no_recipient_silent(self, monkeypatch): - from apps.worker import main as wm - - wm._last_alerts.clear() - sent: list[str] = [] - - def fake_send(self, recipient, subject, html): - sent.append(subject) - return True - - monkeypatch.setattr( - "packages.integrations.notifier.NotificationService.send_email_html", fake_send - ) - # notify_default_to 未配置 - monkeypatch.setattr( - "packages.config.get_settings", - lambda: type("S", (), {"notify_default_to": None})(), - ) - wm._send_alert("[PaperMind] Test", "

x

", "no_recv") - assert len(sent) == 0, "无收件人应静默跳过" - wm._last_alerts.clear() - - class TestTopicLastErrorExposed: """主题 last_error 经 repository 可读(可观测性:API 不再掩盖)"""