From c33f4f194b96ffed72f616941bd3eb06c6e9c705 Mon Sep 17 00:00:00 2001 From: Your Name Date: Sun, 19 Apr 2026 03:32:56 +0800 Subject: [PATCH] Fix worker control fallback for running detect workers --- .../HANDOFF_20260419_0332.md | 122 +++++++++ .../IMPLEMENTATION_STATUS.md | 194 ++++++++++--- docs/ops_center_runtime/TASK_BOARD.md | 259 ++++++++++++++---- .../app/services/worker_control_service.py | 2 + .../tests/test_worker_control_service.py | 38 +++ domainCheck/detect_worker.py | 22 ++ 6 files changed, 550 insertions(+), 87 deletions(-) create mode 100644 docs/ops_center_runtime/HANDOFF_20260419_0332.md create mode 100644 domain-api/tests/test_worker_control_service.py diff --git a/docs/ops_center_runtime/HANDOFF_20260419_0332.md b/docs/ops_center_runtime/HANDOFF_20260419_0332.md new file mode 100644 index 0000000..9cf757b --- /dev/null +++ b/docs/ops_center_runtime/HANDOFF_20260419_0332.md @@ -0,0 +1,122 @@ +# HANDOFF 2026-04-19 03:32 + +## 本轮做了什么 + +这轮没有扩功能,没有进新页面,也没有碰发布动作。 + +只做了两件事: + +1. 把检测执行停滞的判断继续收紧 +2. 对 worker 控制消息链补上最小代码修复 + +## 当前最新判断 + +当前不是接管问题,也不是日志回传问题。 + +当前真实状态是: + +- `remote_access_ready = 3/3` +- `log_sync_state = full_capture` +- 大陆首批结果同步已经成功过 +- 但检测任务没有持续前进 + +这个停滞目前有两层原因: + +### 第一层:已确认的现场阻塞 + +- `mainland-controller-01` 代理池刷新后 `0 available` +- 代理源可以拉到原始代理 +- 但抽样校验全部失败 + +这会直接影响 controller 侧吞吐。 + +### 第二层:刚补好的代码风险 + +之前存在一个运行态风险: + +- 控制面发 `start_detection` +- Redis `publish` 成功 +- 但 worker 如果漏收 pubsub 消息 +- pending 指令可能不会在运行态被再次消费 + +这会造成: + +- 页面像是发起成功了 +- 但 worker 不一定真的继续推进 + +现在这处代码缺口已经补上,但还没有完成线上部署验证。 + +## 本轮代码修复 + +### 1. `worker_control_service` + +文件: + +- `domain-api/app/services/worker_control_service.py` + +修复: + +- 每条 worker 控制消息补 `request_id` +- Redis `publish` 和 pending fallback 复用同一条消息体 + +### 2. `detect_worker` + +文件: + +- `domainCheck/detect_worker.py` + +修复: + +- 运行态心跳里周期性补偿消费 pending 控制消息 +- 真正收到控制消息后,按 `request_id` 清理 pending + +### 3. 新增测试 + +文件: + +- `domain-api/tests/test_worker_control_service.py` + +验证目标: + +- 确保 `send_worker_command(...)` 会生成并持久化 `request_id` + +## 本地验证结果 + +已通过: + +```bash +PYTHONPATH=/www/wwwroot/getDomain/domain-api /opt/domaincheck/domainCheck/.venv/bin/python -m unittest domain-api/tests/test_worker_control_service.py +python -m py_compile domainCheck/detect_worker.py domain-api/app/services/worker_control_service.py +``` + +## 下一轮唯一该做什么 + +下一轮不要发散,只做最小部署验证: + +1. 大陆节点拉最新代码 +2. 重启 `domaincheck-worker` +3. 复查: + - `/api/v1/detect/job/active` + - `/api/v1/runtime/sync-summary` + - controller / worker `scene-log` + - 必要时再采集 `domaincheck-worker` 日志 +4. 只判断两件事: + - 修复上线后,检测任务是否恢复推进 + - 如果仍不推进,是否只剩 controller 代理池 `0 available` 这个现场阻塞 + +## 现在不要做什么 + +- 不做新页面 +- 不做新模块 +- 不做控制面增强 +- 不做发布动作 +- 不把问题再泛化成“日志没回来”或“接管没完成” + +## 一句话结论 + +当前最准确的状态是: + +- 接管和同步基本已经闭合 +- 检测执行停滞仍然存在 +- worker 控制消息补偿链已补代码 +- 下一步只剩“拉代码重启 worker 后看任务是否恢复推进” diff --git a/docs/ops_center_runtime/IMPLEMENTATION_STATUS.md b/docs/ops_center_runtime/IMPLEMENTATION_STATUS.md index 3831a90..f7c4ae7 100644 --- a/docs/ops_center_runtime/IMPLEMENTATION_STATUS.md +++ b/docs/ops_center_runtime/IMPLEMENTATION_STATUS.md @@ -1,22 +1,22 @@ # IMPLEMENTATION_STATUS -更新时间:2026-04-19 03:00 CST +更新时间:2026-04-19 03:38 CST ## 当前真实状态 阶段判断: - 海外单脑接管能力:约 `95%` -- 分布式检测真实执行能力:约 `94%` -- 距离“可稳定上线并放心用后台发起检测”:约 `93%~94%` +- 分布式检测真实执行能力:约 `88%~90%` +- 距离“可稳定上线并放心用后台发起检测”:约 `86%~89%` ## 本轮最新结论 -这轮只做了只读复查,结论更明确了: +这轮结论需要更新为三段: -- 中央接收 `detect_result_projection` 的代码已经就位 -- `sync_agent` 自动产 `detect_result_projection` 的代码也已经就位 -- 但 mainland 端当前仍未真正跑到新版 `domaincheck-sync-agent` +- 接管与同步能力已经明显趋于完成 +- 检测执行面当前没有持续产出 +- 其中一处代码级风险已经完成最小修复,但尚未完成线上部署验证 ## 当前证据拆分 @@ -27,19 +27,119 @@ - `detect_result_projection` 支持 `recent_domain_events` - 中央 ingest 会把逐条事件写入 `detect_run_events` - `sync_agent` 会自动产出 `detect_result_projection` +- worker 控制消息新增 `request_id` +- worker 运行态心跳会补偿消费 pending 控制消息 +- worker 收到并处理控制消息后会按 `request_id` 清理 pending 指令 -### 2. 真实运行闭环 +本地代码验证已通过: -当前仍未完成: +- `unittest domain-api/tests/test_worker_control_service.py` +- `python -m py_compile domainCheck/detect_worker.py domain-api/app/services/worker_control_service.py` -- `runtime_ingest` 持续刷新 -- `detect_result_ingest` 没有刷新 -- mainland `domain_*` 事件仍为 0 +### 2. 接管/同步闭环 + +当前已完成: + +- `remote_access_ready = 3/3` +- `log_sync_state = full_capture` +- `mainland-controller-01` 的 `domaincheck-sync-agent` 已重启到新进程 +- 首轮 `detect_result_projection` 推送成功过一次 +- 中央已收到 mainland 首批 `domain_*` 事件 这说明: - mainland 到中央的基础同步链是活的 -- 但 mainland 目前还没有跑到新版 `sync_agent` +- 结果投影链至少成功打通过一次 + +### 3. 检测执行闭环 + +当前未完成: + +- 连续 40 秒观察窗口内: + - `progress_percent` 不变 + - `items_completed` 不变 + - `items_running` 不变 + - `items_claimed` 不变 + - `items_pending` 不变 +- mainland `domain_*` 事件数仍固定在 `10` +- 最新 mainland `detect_result_ingest` 仍停在 `5382` +- `activity-stream` 明确提示: + - `近窗刚有吞吐 0 台` + +这说明: + +- 当前不是单纯“页面没刷新” +- 而是执行现场这段时间没有继续出结果 + +## 本轮新增硬证据 + +通过中央观测面、节点现场日志和远端 `domaincheck-worker` 日志,已确认: + +- `mainland-controller-01` 现场日志显示: + - `当前可用代理数: 0` + - `最近结果: 刷新成功,可用 0 个` +- controller 新增远端日志显示: + - 代理源拉取成功 + - 抽样校验后 `共 0 个可用代理` + - 失败集中在: + - `ProxyError@https://m.baidu.com` + - `Unable to connect to proxy` + - `ConnectTimeoutError` +- worker 新增远端日志显示: + - `时光机检测 外部依赖异常,步骤降级继续执行` + - 同时仍有 `域名检测完成` + - 新收到的重复启动指令被忽略,因为任务已在运行 +- 两台大陆节点 full capture 已开启,但源日志时间没有继续前进 + +这说明: + +- 当前第一主阻塞更像 controller 代理池不可用 +- worker 的时光机异常是客观存在的次级问题 +- 不是控制面未接管 +- 不是同步链未打通 + +## 本轮新增代码修复 + +本轮不是只停留在诊断,还补了一处运行态最小修复: + +### 修复点 1:控制消息唯一标识 + +- 文件: + - `domain-api/app/services/worker_control_service.py` +- 变更: + - 每次 `send_worker_command(...)` 都附带 `request_id` + - Redis `publish` 与 pending fallback 使用同一份消息体 + +### 修复点 2:worker 运行态补偿消费 pending 指令 + +- 文件: + - `domainCheck/detect_worker.py` +- 变更: + - 心跳线程每轮会补偿尝试消费 pending 控制消息 + - 解决“worker 在线但 pubsub 消息漏收,导致 pending 指令长期不被消费”的风险 + +### 修复点 3:worker 收到指令后确认清理 pending + +- 文件: + - `domainCheck/detect_worker.py` +- 变更: + - worker 实际收到控制消息后,会按 `request_id` 清理对应 pending 指令 + - 避免修复后又产生重复回放 + +### 当前判断 + +这组修复解决的是: + +- “检测启动已经发布,但 worker 可能静默漏收”的代码风险 + +这组修复还没有证明能够解决的是: + +- `controller` 代理池校验后 `0 available` 的外部执行阻塞 + +所以当前最准确的状态是: + +- 代码级控制链缺口已补上 +- 运行态是否恢复,还要看线上节点拉代码重启后的观测结果 ## 当前已闭合的问题 @@ -54,46 +154,60 @@ ## 当前唯一剩余问题 -当前唯一剩余问题: +当前唯一主问题仍然是: -- mainland 节点尚未完成“拉新代码并重启 `domaincheck-sync-agent`” +- 检测执行面停滞 + +但其下应拆成两个层次: + +- 已补代码、待验证: + - worker 控制消息补偿消费链 +- 已确认现场阻塞: + - controller 代理池可用性为 0 + - 因此没有持续产生新的 domain 级结果 换句话说: - 现在不是逻辑未实现 - 不是中央映射失败 -- 也不是检测任务没跑 -- 而是 mainland 仍在跑旧版 `sync_agent` +- 不是节点未接管 +- 而是执行现场没有继续产出,且 controller 侧卡在代理校验失败 + +## acceptance 当前状态 + +最新接管验收结果仍然成立: + +- `pbr-9ce5c85f17`:`onboarding.acceptance` 成功 +- `pbr-389abd618c`:`onboarding.acceptance` 成功 + +旧的 `attention` run 依然是历史残留,但它们已经不是当前最真实的生产阻塞。 ## 当前是否可以继续跑检测测试 当前结论: -- 可以继续跑检测测试 -- 主统计和节点参与已经可信 -- 但中央逐条账本仍不能视为 mainland 最终真相 +- 可以继续做最小部署验证 +- 也可以继续做执行恢复类排查 +- 但还不能把当前状态视作“后台检测已经稳定恢复” ## 当前是否建议直接上线 当前结论: -- 如果目标是: - - 后台稳定发起检测 - - 看懂谁在跑 - - 看懂整体任务推进 - 那已经基本可用 +- 不建议现在按“可稳定上线”判断 -- 如果目标是: - - 中央逐条反映大陆每个域名的终态 - 那还差最后一个外部动作: - - mainland 拉新 - - 重启 `domaincheck-sync-agent` +原因不是接管面,而是执行面: + +- 三台节点都已接入 +- 但当前任务没有持续吞吐 +- controller 代理池全部验不过会直接影响检测产出 +- 并且这轮刚补上的 worker 控制消息修复还没经过线上运行验证 ## 当前优先级判断 最高优先级: -- `J1-大陆逐条结果同步前置批` +- `J2-检测执行停滞收口批` 当前不应继续推进: @@ -101,12 +215,20 @@ - 新模块 - 新页面 - 发布动作 -- 与逐条结果同步无关的工作 +- 与检测执行停滞无关的工作 ## 完成下一轮后的预期 -如果下一轮完成 mainland 拉新并重启 `domaincheck-sync-agent`,并验证贯通: +如果下一轮确认: -- mainland `domain_*` 会进入中央 -- item 级统一才会真正可落地 -- 整体可上线程度预计提升到 `95%~96%` +- 大陆节点已拉到这版控制消息修复 +- `domaincheck-worker` 重启后稳定收新指令 +- controller 代理池恢复可用 +- `domain_*` 开始继续增长 +- `items_completed` 和近窗吞吐重新前进 + +则整体可上线程度预计可回升到: + +- `92%~94%` + +在那之前,当前口径应保持保守。 diff --git a/docs/ops_center_runtime/TASK_BOARD.md b/docs/ops_center_runtime/TASK_BOARD.md index 4fbf08b..fd5c263 100644 --- a/docs/ops_center_runtime/TASK_BOARD.md +++ b/docs/ops_center_runtime/TASK_BOARD.md @@ -1,99 +1,248 @@ # TASK_BOARD -更新时间:2026-04-19 03:00 CST +更新时间:2026-04-19 03:38 CST ## 当前主批次 -唯一主批次:`J1-大陆逐条结果同步前置批` +唯一主批次:`J2-检测执行停滞收口批` 目标: - 不进入新页面 - 不扩展控制面 -- 只解决 item 级统一前的唯一前置缺口: - - 大陆节点逐条检测事件进入中央 +- 不新增发布动作 +- 只收口当前唯一真实阻塞: + - 节点已接管 + - 同步已打通 + - 但检测执行没有继续产出新结果 + +## 本轮最新状态 + +本轮在只读取证基础上,已经补了一处最小代码修复,但还没有完成线上节点拉代码后的运行验证。 + +### 已完成的最小修复 + +- `worker_control_service.send_worker_command(...)` + - 为每条 worker 控制消息补上 `request_id` + - 保证 Redis 发布与 pending fallback 使用同一份负载 +- `detect_worker` + - 在运行态心跳里周期性补偿消费 pending 控制消息 + - 在真正收到控制消息后,按 `request_id` 清理 pending 指令 + +本地验证已通过: + +- `unittest domain-api/tests/test_worker_control_service.py` +- `python -m py_compile domainCheck/detect_worker.py domain-api/app/services/worker_control_service.py` + +### 仍未完成的事情 + +- 大陆节点尚未基于这版修复完成拉代码并重启 `domaincheck-worker` +- 因此这轮还不能把“检测停滞”完全归因成单一外部代理问题 +- 当前更准确的表述应为: + - `controller` 代理池可用数为 `0` 是已确认现场阻塞 + - worker 控制消息存在“可能丢发布、pending 不被运行态补偿消费”的代码风险,现已本地修复,待线上验证 ## 本轮最新复查结果 -本轮没有新增功能扩展,只做了中央只读复查: +本轮没有新增实现,只做了中央只读复查 + 节点现场日志取证 + 远端 `domaincheck-worker` 日志采样: -- `runtime_ingest` 仍在持续刷新 -- `detect_result_ingest` 仍未刷新 -- `detect_run_events` 里仍没有 mainland `domain_*` +- 活跃任务仍是 `detect-20260417170546-96023e` +- 三台节点都在参与: + - `mainland-controller-01` + - `mainland-worker-01` + - `overseas-control-01` +- 但连续一个完整观察窗口内,任务进度没有前进: + - `completed = 21` + - `running = 14` + - `claimed = 34` + - `pending = 931` +- mainland `domain_*` 事件仍停在 `10` 条,没有继续增长 +- 最新 mainland `detect_result_ingest` 仍是: + - `id = 5382` + - `created_at = 2026-04-19 03:10:14` +- `activity-stream` 已明确显示: + - `近窗刚有吞吐 0 台` +- 新增日志采样结果: + - `mainland-controller-01` + - 代理源拉取成功 + - 代理校验后 `共 0 个可用代理` + - `mainland-worker-01` + - 时光机依赖异常会降级继续执行 + - 收到重复启动指令时会忽略 ## 关键证据 -### 证据 1:mainland 同步链仍然活着 +### 证据 1:接管与同步已通 -中央最新记录显示: +当前中央状态: -- `runtime_ingest.updated_at` - - 已刷新到 `2026-04-19 02:59:12` +- `remote_access_ready = 3/3` +- `log_sync_state = full_capture` +- mainland controller 的 `domaincheck-sync-agent` 已在新进程上运行 说明: -- mainland controller 仍在向中央推运行时投影 +- 当前不是接管问题 +- 也不是日志回传问题 +- 更不是同步链完全断开 -### 证据 2:mainland 逐条结果投影仍未开始刷新 +### 证据 2:检测任务当前没有继续出新结果 -中央最新记录显示: +连续 40 秒前后对比结果完全一致: -- `detect_result_ingest.updated_at` - - 仍停留在 `2026-04-17 16:57:32` - -而且当前最新 mainland `detect_result_ingest` 中: - -- `job_code = None` -- `latest_event_type = None` -- `recent_domain_events = 0` +- `JOB_PROGRESS = 2.1` +- `items_completed = 21` +- `items_running = 14` +- `items_claimed = 34` +- `items_pending = 931` 说明: -- mainland 还没有开始推送新版 `detect_result_projection` +- 当前不是“页面慢一拍” +- 而是执行面这段时间确实没有继续产出 -### 证据 3:中央逐条事件仍为空 +### 证据 3:中央 mainland 逐条结果没有继续增长 当前中央查询结果: - mainland `domain_started/domain_completed/domain_failed/domain_blacklisted` - - `count = 0` + - 仍为 `10` +- 最新 mainland `detect_result_ingest` + - 仍为 `5382` + +说明: + +- 首批同步成功过 +- 但后续并没有继续流入新逐条结果 + +### 证据 4:controller 现场日志已指向代理池可用性为 0 + +`mainland-controller-01` 现场日志显示: + +- `当前可用代理数: 0` +- `最近结果: 刷新成功,可用 0 个` +- 免费检测链包含: + - 注册查询 + - 百度 site + - 360 site + - 站长之家 + - 爱站 + - 时光机 + +进一步的远端 `logs.collect(domaincheck-worker)` 结果显示: + +- controller 能从 6 个代理源成功拉到原始代理 +- 但在抽样验证后: + - `代理池刷新完成,共 0 个可用代理` +- 失败原因集中在: + - `ProxyError@https://m.baidu.com` + - `Unable to connect to proxy` + - `ConnectTimeoutError` + +说明: + +- 当前不是代理源接口没返回 +- 而是“拿到的代理全部验不过” +- 主阻塞已经可以精确收紧到 controller 代理池不可用 + +### 证据 5:worker 的时光机异常存在,但不是第一主因 + +`mainland-worker-01` 的远端 `domaincheck-worker` 日志显示: + +- 存在: + - `时光机检测 外部依赖异常,步骤降级继续执行` +- 同时仍可见: + - `域名检测完成` +- 最近收到控制消息后: + - `收到启动检测指令,但检测任务已在运行,忽略重复启动` + +说明: + +- worker 并不是完全不能执行 +- 时光机异常是客观存在的次级问题 +- 但它不像 controller 代理池为 0 那样直接卡住整体吞吐 + +### 证据 6:full_capture 已开启,但源日志时间没有继续前进 + +节点现场日志最新可见记录仍停在较早时间: + +- `mainland-controller-01` + - 最近样本集中在 `01:00:23` +- `mainland-worker-01` + - 最近样本集中在 `01:00:35` + +额外 35 秒观察窗口结果: + +- `capture_at` 没有继续前进 +- `source_msg` 里的源日志时间也没有继续前进 + +说明: + +- 不是“全量日志没开” +- 也不是“日志回传没回来” +- 而是执行进程这段时间确实没有继续产生日志 + +### 证据 7:worker 控制消息补偿链已补上代码缺口 + +之前运行态存在一个真实风险: + +- `runtime.start_detection` 成功只代表“控制消息已发布到 Redis” +- 如果 worker 在线但恰好漏收了 pubsub 消息 +- pending 指令只在启动或重连时消费,就可能造成: + - 页面显示“已启动” + - 但 worker 实际没有进入新的执行流 + +当前已完成最小修复: + +- 每条控制消息带 `request_id` +- worker 心跳会主动补偿消费 pending 指令 +- worker 真正收到并处理后,会按 `request_id` 清理 pending + +这一步已经消除了“控制消息可能被静默漏掉”的代码风险,但还需要线上部署后再观察真实任务是否恢复推进。 ## 当前结论 -当前结论已经可以收敛为一句话: +当前结论要收敛成一句话: -- 代码闭环已完成 -- 真实链路未贯通 -- 唯一剩余阻塞是 mainland 端没有真正跑到新版 `domaincheck-sync-agent` +- 接管链路已基本完成 +- 同步链路已至少成功贯通一次 +- 当前主阻塞依旧是检测执行面停滞 +- 其中已确认的现场问题有两个层次: + - `mainland-controller-01` 代理池可用性为 `0` + - worker 控制消息补偿链之前存在代码缺口,现已本地修复,待部署验证 +- worker 侧时光机异常仍是次级噪音,不应继续放大 ## 下一步唯一主批次 -唯一主批次仍然是:`J1-大陆逐条结果同步前置批` +唯一主批次维持为:`J2-检测执行停滞收口批` 本批次唯一目标: -- 让 mainland 拉最新代码并重启 `domaincheck-sync-agent` -- 然后在中央确认出现: - - `detect_result_ingest.payload.projection.recent_domain_events` - - mainland `domain_*` +- 先把控制消息补偿修复部署到线上 worker +- 再复查检测任务是否恢复推进 +- 若仍停滞,再把剩余问题继续收紧到 controller 代理池可用数为 0 +- 全程不做任何功能扩展 ## 候选批次 -### Candidate J2 - -名称:item 级回写最小落地批 - -进入条件: - -- `J1` 完成后,中央已经能收到大陆逐条结果 - ### Candidate J3 -名称:raw/effective 运营文案收口批 +名称:检测执行恢复验证批 进入条件: -- 如果用户只想先补页面文案 +- `J2` 明确根因后 +- 代理或外部站点依赖恢复 +- 再观察 `domain_*` 是否重新增长 + +### Candidate J4 + +名称:上线签收证据批 + +进入条件: + +- `J3` 确认检测重新持续产出 +- `go-live-summary` 只剩历史 attention ## 暂停项 @@ -103,14 +252,22 @@ - 新模块 - 控制面增强 - 发布动作 -- 与大陆逐条结果同步无关的工作 +- 与检测执行停滞无关的工作 ## 下一轮唯一动作 下一轮唯一应该继续做的事情: -- 在 mainland 节点执行: - - `git pull` - - `systemctl restart domaincheck-sync-agent` - - `systemctl status domaincheck-sync-agent --no-pager -l` - - `journalctl -u domaincheck-sync-agent -n 80 --no-pager` +- 只做最小部署与运行态收口: + - 大陆节点拉取最新代码 + - 重启 `domaincheck-worker` + - 复查 `/api/v1/detect/job/active` + - 复查 `/api/v1/runtime/sync-summary` + - 复查 controller / worker `scene-log` + - 必要时追加 `domaincheck-worker` 运行日志取证 +- 唯一判断问题是: + - 修复部署后,任务是否重新推进 + - 若仍不推进,是否只剩 controller 代理池 `0 available` 这一条主阻塞 + - 是否仍是 `近窗吞吐 0` + - 是否仍是 `controller 可用代理数 0` + - 代理恢复后 `domain_*` 是否重新增长 diff --git a/domain-api/app/services/worker_control_service.py b/domain-api/app/services/worker_control_service.py index 470f23b..d4e78de 100644 --- a/domain-api/app/services/worker_control_service.py +++ b/domain-api/app/services/worker_control_service.py @@ -5,6 +5,7 @@ import os import subprocess from datetime import datetime from pathlib import Path +from uuid import uuid4 from app.core.config import settings from app.core.redis_client import get_redis @@ -288,6 +289,7 @@ def send_worker_command(action: str, payload: dict | None = None) -> tuple[bool, command_payload = {"action": action} if payload: command_payload.update(payload) + command_payload["request_id"] = str(command_payload.get("request_id") or f"workerctl-{uuid4().hex[:12]}") serialized = json.dumps(command_payload, ensure_ascii=False) redis_client.set(WORKER_PENDING_COMMAND_KEY, serialized, ex=120) redis_client.publish(WORKER_CONTROL_CHANNEL, serialized) diff --git a/domain-api/tests/test_worker_control_service.py b/domain-api/tests/test_worker_control_service.py new file mode 100644 index 0000000..d2a33cd --- /dev/null +++ b/domain-api/tests/test_worker_control_service.py @@ -0,0 +1,38 @@ +import json +import unittest +from unittest.mock import Mock, patch + +from app.services.worker_control_service import WORKER_CONTROL_CHANNEL, WORKER_PENDING_COMMAND_KEY, send_worker_command + + +class WorkerControlServiceTests(unittest.TestCase): + @patch("app.services.worker_control_service.get_redis") + def test_send_worker_command_persists_request_id_and_publishes_same_payload(self, mock_get_redis) -> None: + redis_client = Mock() + mock_get_redis.return_value = redis_client + + ok, message = send_worker_command( + "start_detection", + payload={"job_id": 1, "job_code": "detect-20260419030000-abc123"}, + ) + + self.assertTrue(ok) + self.assertIn("已发送 Worker 控制指令", message) + + redis_client.set.assert_called_once() + set_args = redis_client.set.call_args.args + self.assertEqual(WORKER_PENDING_COMMAND_KEY, set_args[0]) + serialized = set_args[1] + self.assertEqual(120, redis_client.set.call_args.kwargs["ex"]) + + payload = json.loads(serialized) + self.assertEqual("start_detection", payload["action"]) + self.assertEqual(1, payload["job_id"]) + self.assertEqual("detect-20260419030000-abc123", payload["job_code"]) + self.assertTrue(payload["request_id"].startswith("workerctl-")) + + redis_client.publish.assert_called_once_with(WORKER_CONTROL_CHANNEL, serialized) + + +if __name__ == "__main__": + unittest.main() diff --git a/domainCheck/detect_worker.py b/domainCheck/detect_worker.py index 2e1fe4f..df4a5b6 100644 --- a/domainCheck/detect_worker.py +++ b/domainCheck/detect_worker.py @@ -940,6 +940,10 @@ class DetectWorker: def _runtime_heartbeat_loop(self): while not self._runtime_heartbeat_stop.wait(self.runtime_heartbeat_interval): + try: + self._consume_pending_control_command() + except Exception as e: + logger.debug(f"运行态心跳补偿消费待执行控制指令失败: {e}") try: self._update_runtime_state( self._last_runtime_phase, @@ -1157,11 +1161,29 @@ class DetectWorker: ) return True + def _acknowledge_pending_control_command(self, control_payload): + if not self.use_redis or self.redis_client is None: + return + request_id = str((control_payload or {}).get("request_id", "")).strip() + if not request_id: + return + try: + raw_pending = self.redis_client.get(PENDING_CONTROL_KEY) + if not raw_pending: + return + pending_payload = json.loads(raw_pending) + pending_request_id = str((pending_payload or {}).get("request_id", "")).strip() + if pending_request_id and pending_request_id == request_id: + self.redis_client.delete(PENDING_CONTROL_KEY) + except Exception as e: + logger.debug(f"确认待执行控制指令失败: {e}") + def _handle_control_message(self, payload): try: control_payload = json.loads(payload) if isinstance(payload, str) else payload except Exception: control_payload = {"action": str(payload)} + self._acknowledge_pending_control_command(control_payload) action = str(control_payload.get("action", "")).strip() if action == "start_detection": self.start_detection_async(source="redis-control", control_payload=control_payload)