diff --git a/docs/ops_center_runtime/HANDOFF_20260419_0257.md b/docs/ops_center_runtime/HANDOFF_20260419_0257.md new file mode 100644 index 0000000..e632ea4 --- /dev/null +++ b/docs/ops_center_runtime/HANDOFF_20260419_0257.md @@ -0,0 +1,66 @@ +# HANDOFF 2026-04-19 02:57 + +## 本轮新增完成 + +在上一轮“中央可接收并落库逐条事件”基础上,这一轮又补了最后一个自动触发点: + +- `app.sync_agent` 每个 tick 会自动产出 `detect_result_projection` +- 不再依赖有人去访问 mainland 检测页面 + +## 当前真实判断 + +当前看到的中央证据是: + +- `mainland-controller-01` / `mainland-worker-01` 心跳正常 +- `runtime_ingest.updated_at` 已刷新到 `2026-04-19 02:55:36` +- `detect_result_ingest.updated_at` 仍停留在 `2026-04-17 16:57:32` +- `detect_run_events` 里仍没有 mainland `domain_*` + +这说明: + +- mainland 同步链是活的 +- 但 अभी还没有跑到“自动产 detect_result_projection”的新代码 + +## 这轮后的唯一剩余动作 + +现在只剩一个最小外部动作: + +1. mainland 相关节点拉最新代码 +2. 重启 `domaincheck-sync-agent` + +建议最小命令: + +```bash +systemctl restart domaincheck-sync-agent +sleep 3 +systemctl status domaincheck-sync-agent --no-pager -l +journalctl -u domaincheck-sync-agent -n 80 --no-pager +``` + +如果 service 名不是这个,再按机器当前实际 service 名执行。 + +## 重启后我下一轮只查什么 + +下一轮只查这 3 件事: + +1. 最新 `detect_result_ingest.payload.projection.recent_domain_events` 是否非空 +2. `detect_run_events` 是否出现 mainland `domain_started/domain_completed/domain_failed/domain_blacklisted` +3. `detect/job/active` 的 `recent_events` 是否开始出现 imported mainland 逐条事件 + +## 继续不要做什么 + +- 不做 item 级最终回写 +- 不做新页面 +- 不做控制面增强 +- 不做发布动作 + +## 推荐模型与推理等级 + +继续推荐: + +- `GPT-5.4` +- `high` + +## 一句话结论 + +现在代码链路已经补到“mainland 拉新并重启 sync agent 后就该自动把逐条结果送进中央”;下一轮只差验证真正贯通。 diff --git a/docs/ops_center_runtime/IMPLEMENTATION_STATUS.md b/docs/ops_center_runtime/IMPLEMENTATION_STATUS.md index f120c7a..fdbe3f3 100644 --- a/docs/ops_center_runtime/IMPLEMENTATION_STATUS.md +++ b/docs/ops_center_runtime/IMPLEMENTATION_STATUS.md @@ -1,6 +1,6 @@ # IMPLEMENTATION_STATUS -更新时间:2026-04-19 02:51 CST +更新时间:2026-04-19 02:57 CST ## 当前真实状态 @@ -13,12 +13,13 @@ 本轮最重要的新事实: - 中央“逐条事件接收与落库”能力已经补齐 +- mainland `sync_agent` 自动产 `detect_result_projection` 的缺口也已经补齐 - 当前真正剩余阻塞变成: - - 大陆节点还没把这版代码跑起来并把逐条事件推到中央 + - 大陆节点还没把这版代码跑起来并重启 `sync_agent` ## 本轮结论收口 -本轮没有发散做新模块,而是完成了 `J1` 的最小实现: +本轮没有发散做新模块,而是完成了 `J1` 的第二轮最小实现: ### 1. `detect_result_projection` 已扩展为可带逐条事件 @@ -46,25 +47,35 @@ - `import_source_job_code` - `import_fingerprint` -### 3. 真实端到端仍待大陆节点拉新验证 +### 3. `sync_agent` 已会自动产出 `detect_result_projection` + +本轮新增: + +- `app.sync_agent` 每个 tick 会基于当前 `active_job` +- 自动调用 `append_detect_result_projection_if_changed` + +这意味着: + +- mainland 只要拉到这版代码并重启 `domaincheck-sync-agent` +- 就会自动开始补推 `detect_result_projection` +- 不再依赖额外访问 mainland 检测页面 + +### 4. 真实端到端仍待 mainland 拉新并重启 `sync_agent` 当前中央实际可见: -- `overseas-control-01` - - `domain_started` - - `domain_completed` +- `mainland-controller-01` / `mainland-worker-01` 心跳正常 +- `runtime_ingest.updated_at` + - 已刷新到 `2026-04-19 02:55:36` +- `detect_result_ingest.updated_at` + - 仍停在 `2026-04-17 16:57:32` +- `detect_run_events` + - 仍无 mainland `domain_*` -当前中央仍不可见: +说明: -- `mainland-controller-01` - - 无 `domain_*` -- `mainland-worker-01` - - 无 `domain_*` - -说明中央代码已补,但真实输入流还没跑到新版本: - -- 中央还没收到带 `recent_domain_events` 的 mainland 投影 -- 所以真实库里还没有新的 mainland `domain_*` +- mainland 同步链活着 +- 但 अभी还没有跑到会自动产新格式投影的这版 `sync_agent` ## 当前已闭合的问题 @@ -75,17 +86,19 @@ - Detect 主统计口径不一致 - Runtime / Queue / Detect 主摘要不统一 - “中央如何接并落大陆逐条事件”的代码缺口 +- “mainland 必须人工访问检测页面才会产生 detect_result_projection”的自动触发缺口 ## 当前剩余问题 当前唯一剩余问题: -- 大陆节点尚未完成“拉新代码并继续推投影”的最后验证 +- mainland 节点尚未完成“拉新代码并重启 `domaincheck-sync-agent`”的最后验证 换句话说: - 当前不是“中央不会接” -- 而是“大陆还没把新格式投影推上来” +- 也不是“代码还没写” +- 而是 mainland 还没跑到会自动产新格式投影的这一版 ## 当前是否可以继续跑检测测试 @@ -93,7 +106,7 @@ - 可以继续跑检测测试 - 主统计和节点参与已经足够可信 -- 逐条账本是否完整,取决于大陆节点是否已经拉到这版代码 +- 逐条账本是否完整,取决于 mainland 是否已经拉到这版代码并重启 `domaincheck-sync-agent` ## 当前是否建议直接上线 @@ -107,7 +120,9 @@ - 如果目标是: - 中央账本逐条反映大陆每个域名的终态 - 那现在只差最后一个落地验证动作 + 那现在只差最后一个落地验证动作: + - mainland 拉新 + - 重启 `domaincheck-sync-agent` ## 当前优先级判断 @@ -125,7 +140,7 @@ ## 完成下一轮后的预期 -如果下一轮完成大陆拉新并验证贯通: +如果下一轮完成 mainland 拉新、重启 `sync_agent` 并验证贯通: - item 级统一才会真正可落地 - 后续 `detect_job_items` 回写会变成纯工程实现问题 diff --git a/docs/ops_center_runtime/TASK_BOARD.md b/docs/ops_center_runtime/TASK_BOARD.md index c5a0585..912a238 100644 --- a/docs/ops_center_runtime/TASK_BOARD.md +++ b/docs/ops_center_runtime/TASK_BOARD.md @@ -1,6 +1,6 @@ # TASK_BOARD -更新时间:2026-04-19 02:51 CST +更新时间:2026-04-19 02:57 CST ## 当前主批次 @@ -15,11 +15,12 @@ ## 本轮已完成 -本轮完成的是最小输入通道补齐: +本轮完成的是最小输入通道第二轮收口: - 已让 `detect_result_projection` 开始携带 `recent_domain_events` - 已让中央 `detect_result_ingest` 在接收时尝试落库 `detect_run_events` - 已补幂等去重,避免重复投影重复写事件 +- 已补 `sync_agent` 自动产出 `detect_result_projection` - 已完成最小单测与中央 API 重启 ## 关键证据 @@ -34,7 +35,33 @@ - 提取 `domain_started/domain_completed/domain_failed/domain_blacklisted` - 以幂等方式写入中央 `detect_run_events` -### 证据 2:当前真实库里大陆逐条事件仍未出现 +### 证据 2:mainland 当前只在持续推 `runtime_projection` + +中央当前可见: + +- `runtime_ingest.updated_at` + - 已刷新到 `2026-04-19 02:55:36` +- `detect_result_ingest.updated_at` + - 仍停在 `2026-04-17 16:57:32` + +说明: + +- mainland 同步链是活的 +- 但还没有持续产出新的 `detect_result_projection` + +### 证据 3:自动产投影缺口已经补上 + +当前代码状态: + +- `sync_agent` 每个 tick 会基于 `active_job` 自动调用: + - `append_detect_result_projection_if_changed` + +所以 mainland 拉新并重启 `domaincheck-sync-agent` 后: + +- 不再依赖有人访问 mainland 检测页面 +- 逐条结果会自动进入中央 + +### 证据 4:当前真实库里 mainland 逐条事件仍未出现 当前中央库里可见: @@ -49,33 +76,32 @@ - `mainland-worker-01` - 无逐条 `domain_*` 事件 -### 证据 3:当前剩余阻塞已经收缩为“大陆节点尚未跑到这版代码” - -所以现在的状态是: - -- 中央接收与落库逻辑已具备 -- 但大陆节点必须拉到最新代码并继续推投影 -- 拉新前,中央库里当然还看不到新的 mainland `domain_*` - ## 当前结论 -`J1-大陆逐条结果同步前置批` 已进入实现态,不再只是结论确认: +`J1-大陆逐条结果同步前置批` 的代码闭环已经补完: -- 中央缺口已补 -- 现在转入“大陆节点拉新并验证贯通” +- 中央会接 +- 中央会落 +- mainland `sync_agent` 也会自动产投影 + +现在只差: + +- mainland 拉最新代码 +- 重启 `domaincheck-sync-agent` +- 然后验证真实事件是否进入中央 ## 下一步唯一主批次 -唯一主批次:`J1-大陆逐条结果同步前置批` +唯一主批次仍然是:`J1-大陆逐条结果同步前置批` 范围: - 不做最终 item 回写 -- 先补“大陆逐条结果进入中央”的最小数据通道 +- 先完成 mainland 真实事件进中央的端到端验证 本批次唯一目标: -- 让大陆节点把 `recent_domain_events` 真正同步到中央 +- 让 mainland 通过新版 `sync_agent` 自动把 `recent_domain_events` 同步到中央 - 并在中央确认出现 mainland `domain_*` 事件 ## 候选批次 @@ -124,6 +150,6 @@ 下一轮唯一应该继续做的事情: -- 在大陆节点拉最新代码后,验证: - - `detect_result_projection` 已含 `recent_domain_events` +- 在 mainland 节点拉最新代码并重启 `domaincheck-sync-agent` 后,验证: + - `detect_result_ingest.payload.projection.recent_domain_events` 已非空 - 中央 `detect_run_events` 已出现 mainland `domain_*` diff --git a/domain-api/app/sync_agent.py b/domain-api/app/sync_agent.py index 2e249a3..add2a66 100644 --- a/domain-api/app/sync_agent.py +++ b/domain-api/app/sync_agent.py @@ -6,12 +6,32 @@ import time from app.core.config import settings from app.services.debug_event_service import push_debug_event from app.services.detect_job_service import get_active_detect_job_summary, get_detect_queue_health, list_recent_detect_run_events +from app.services.sync_record_service import append_detect_result_projection_if_changed from app.services.sync_push_service import pull_detect_task_batch_now, push_runtime_projection_now logger = logging.getLogger("domaincheck.sync_agent") +def _append_detect_result_projection_snapshot(active_job: dict) -> None: + if not active_job: + return + append_detect_result_projection_if_changed( + detect={ + "active_job": active_job, + "progress": { + "pending": int(active_job.get("items_pending", 0) or 0), + "running": int(active_job.get("items_running", 0) or 0), + "completed": int(active_job.get("items_completed", 0) or 0), + "blacklisted": int(active_job.get("items_blacklisted", 0) or 0), + "failed": int(active_job.get("items_failed", 0) or 0), + }, + "phase_label": str(active_job.get("status") or "").strip(), + "phase_detail": f"sync-agent snapshot for {active_job.get('job_code', '')}", + } + ) + + def _emit_structured_tick( *, base_event_type: str, @@ -91,6 +111,9 @@ def main() -> None: ) while True: try: + active_job = get_active_detect_job_summary(event_limit=10) + if active_job: + _append_detect_result_projection_snapshot(active_job) ok, message, data = push_runtime_projection_now() logger.info("sync tick: ok=%s message=%s data=%s", ok, message, data) push_debug_event( @@ -112,7 +135,6 @@ def main() -> None: payload={"ok": pull_ok, "data": pull_data}, ) _emit_structured_tick(base_event_type="task_pull", ok=pull_ok, message=pull_message, data=pull_data) - active_job = get_active_detect_job_summary(event_limit=10) if active_job: queue_health = get_detect_queue_health(window_minutes=15) recent_events = list_recent_detect_run_events(limit=8) diff --git a/domain-api/tests/test_sync_agent.py b/domain-api/tests/test_sync_agent.py new file mode 100644 index 0000000..82952a3 --- /dev/null +++ b/domain-api/tests/test_sync_agent.py @@ -0,0 +1,35 @@ +import unittest +from unittest.mock import patch + +from app.sync_agent import _append_detect_result_projection_snapshot + + +class SyncAgentTests(unittest.TestCase): + @patch("app.sync_agent.append_detect_result_projection_if_changed") + def test_append_detect_result_projection_snapshot_uses_active_job_summary(self, mock_append) -> None: + _append_detect_result_projection_snapshot( + { + "job_id": 1, + "job_code": "detect-20260419030000-abc123", + "status": "running", + "items_pending": 12, + "items_running": 3, + "items_completed": 8, + "items_blacklisted": 1, + "items_failed": 2, + } + ) + + mock_append.assert_called_once() + payload = mock_append.call_args.kwargs["detect"] + self.assertEqual("running", payload["phase_label"]) + self.assertEqual("detect-20260419030000-abc123", payload["active_job"]["job_code"]) + self.assertEqual(12, payload["progress"]["pending"]) + self.assertEqual(3, payload["progress"]["running"]) + self.assertEqual(8, payload["progress"]["completed"]) + self.assertEqual(1, payload["progress"]["blacklisted"]) + self.assertEqual(2, payload["progress"]["failed"]) + + +if __name__ == "__main__": + unittest.main()