diff --git a/api/ops/react_loop.py b/api/ops/react_loop.py index 71ab695c..9d99dca2 100644 --- a/api/ops/react_loop.py +++ b/api/ops/react_loop.py @@ -3,9 +3,12 @@ from __future__ import annotations import json +import logging import os from typing import Any +logger = logging.getLogger(__name__) + from api.ops.chat_context import load_chat_transcript from api.ops.events_schema import handoff_payload, review_payload from api.ops.llm import chat_completion @@ -15,6 +18,11 @@ from api.ops.react_tools import _build_v0_registry, _truncate_summary from api.ops.review.rules import review_result from api.ops.store.artifacts import save_artifact_with_failure_event +from api.ops.store.checkpoints import ( + CheckpointStoreError, + find_latest_checkpoint_for_session, + save_checkpoint, +) from api.ops.store.runs import OpsRunStore, append_event from api.ops.tracing import trace_span, traceable, update_current_span_metadata @@ -22,6 +30,143 @@ MAX_RETRIES_DEFAULT = 2 +def _build_react_state( + query: str, + session_id: str | None, + messages: list[dict[str, str]], + step: int, + tool_evidence: list[dict[str, Any]], + final_answer: str, + final_verdict: str, + llm_calls: int, + llm_usages: list[LlmUsage], +) -> dict[str, Any]: + """构造可序列化的 ReAct checkpoint 状态。""" + return { + "route": "react", + "query": query, + "session_id": session_id, + "step": step, + "messages": list(messages), + "tool_evidence": list(tool_evidence), + "final_answer": final_answer, + "final_verdict": final_verdict, + "llm_calls": llm_calls, + "llm_usages": [u.to_dict() for u in llm_usages], + } + + +def _try_save_checkpoint( + run_id: str, + session_id: str, + query: str, + messages: list[dict[str, str]], + step: int, + tool_evidence: list[dict[str, Any]], + final_answer: str, + final_verdict: str, + llm_calls: int, + llm_usages: list[LlmUsage], + store: OpsRunStore, +) -> None: + """保存 checkpoint;失败时记录事件但不中断 ReAct 循环。""" + state = _build_react_state( + query=query, + session_id=session_id, + messages=messages, + step=step, + tool_evidence=tool_evidence, + final_answer=final_answer, + final_verdict=final_verdict, + llm_calls=llm_calls, + llm_usages=llm_usages, + ) + try: + save_checkpoint(run_id, session_id, state, store=store) + except Exception as exc: # pragma: no cover - 防御性降级 + logger.warning("checkpoint.save_failed: %s", exc) + store.append_event( + run_id, + "orchestrator", + "checkpoint.save_failed", + payload={"error": str(exc), "session_id": session_id}, + node_id="react.checkpoint.save_failed", + ) + + +def _resume_react_state( + run_id: str, + query: str, + session_id: str, + cp_row: dict[str, Any], + store: OpsRunStore, +) -> dict[str, Any] | None: + """尝试从 checkpoint 行恢复 ReAct 状态。 + + 成功返回状态字典;失败时记录 checkpoint.corrupted 并返回 None。 + """ + try: + state_json = cp_row.get("state_json") + state = _validate_react_checkpoint(state_json) + except CheckpointStoreError as exc: + logger.warning("checkpoint.corrupted: %s", exc) + store.append_event( + run_id, + "orchestrator", + "checkpoint.corrupted", + payload={ + "error": str(exc), + "session_id": session_id, + "from_run_id": str(cp_row.get("run_id", "")), + }, + node_id="react.checkpoint.corrupted", + ) + return None + + prev_run_id = str(cp_row.get("run_id", "")) + store.append_event( + run_id, + "orchestrator", + "checkpoint.resume", + payload={ + "from_run_id": prev_run_id, + "step": state["step"], + "session_id": session_id, + }, + node_id="react.checkpoint.resume", + ) + + messages: list[dict[str, str]] = list(state.get("messages", [])) + if state.get("query") != query: + messages.append({"role": "user", "content": query}) + + return { + "messages": messages, + "step": int(state.get("step", 0)), + "tool_evidence": list(state.get("tool_evidence", [])), + "final_answer": str(state.get("final_answer", "")), + "final_verdict": str(state.get("final_verdict", "partial")), + "llm_calls": int(state.get("llm_calls", 0)), + "llm_usages": [LlmUsage.from_dict(u) for u in state.get("llm_usages", [])], + } + + +def _validate_react_checkpoint(state_json: Any) -> dict[str, Any]: + """校验 checkpoint 状态;失败抛出 CheckpointStoreError。""" + if not isinstance(state_json, dict): + raise CheckpointStoreError("checkpoint state_json is not a dict") + for key in ("route", "query", "step", "messages", "tool_evidence"): + if key not in state_json: + raise CheckpointStoreError(f"checkpoint state missing key: {key}") + if state_json.get("route") != "react": + raise CheckpointStoreError("checkpoint route is not 'react'") + if not isinstance(state_json["messages"], list): + raise CheckpointStoreError("checkpoint state messages is not a list") + if not isinstance(state_json["step"], int): + raise CheckpointStoreError("checkpoint state step is not an int") + return state_json + + @traceable(capture_input=False, capture_output=False) def run_react_fallback( run_id: str, @@ -37,8 +182,6 @@ def run_react_fallback( 与 FSM 路径共用 ops_runs / ops_run_events / Review 闸。 超限 → status partial + 仍 synthesize(非 500)。 """ - transcript = load_chat_transcript(session_id, store=store) - update_current_span_metadata( { "ops_run_id": run_id, @@ -85,21 +228,38 @@ def run_react_fallback( store=store, ) - # System prompt for ReAct - system_prompt = _build_react_system_prompt(tools_json) - messages: list[dict[str, str]] = [ - {"role": "system", "content": system_prompt}, - ] - if transcript: - messages.extend(transcript) - messages.append({"role": "user", "content": query}) - - step = 0 - final_answer = "" - llm_calls = 0 - llm_usages: list[LlmUsage] = [] - tool_evidence: list[dict[str, Any]] = [] - final_verdict = "partial" + # Try to resume from a previous checkpoint for this session + resumed_state: dict[str, Any] | None = None + if session_id: + cp_row = find_latest_checkpoint_for_session(session_id, store=store) + if cp_row: + resumed_state = _resume_react_state(run_id, query, session_id, cp_row, store) + + if resumed_state is None: + # Cold start + transcript = load_chat_transcript(session_id, store=store) + system_prompt = _build_react_system_prompt(tools_json) + messages: list[dict[str, str]] = [ + {"role": "system", "content": system_prompt}, + ] + if transcript: + messages.extend(transcript) + messages.append({"role": "user", "content": query}) + + step = 0 + final_answer = "" + llm_calls = 0 + llm_usages: list[LlmUsage] = [] + tool_evidence: list[dict[str, Any]] = [] + final_verdict = "partial" + else: + messages = resumed_state["messages"] + step = resumed_state["step"] + final_answer = resumed_state["final_answer"] + llm_calls = resumed_state["llm_calls"] + llm_usages = resumed_state["llm_usages"] + tool_evidence = resumed_state["tool_evidence"] + final_verdict = resumed_state["final_verdict"] while step < max_steps: step += 1 @@ -190,6 +350,22 @@ def run_react_fallback( messages.append({"role": "assistant", "content": raw_content}) messages.append({"role": "user", "content": f"Tool result:\n{tool_msg}"}) + # Save checkpoint after each non-final step so that crashes can resume + if session_id: + _try_save_checkpoint( + run_id, + session_id, + query, + messages, + step, + tool_evidence, + final_answer, + final_verdict, + llm_calls, + llm_usages, + store, + ) + else: # max_steps exceeded final_verdict = "partial" diff --git a/api/ops/store/checkpoints.py b/api/ops/store/checkpoints.py new file mode 100644 index 00000000..54d02bbd --- /dev/null +++ b/api/ops/store/checkpoints.py @@ -0,0 +1,97 @@ +"""Ops Desk Checkpoint 存储适配层(P1-2)。 + +复用 `ops_run_checkpoints` 表: +- `checkpoint_id` 字段存放 thread/session 标识。 +- `state_json` 存放 ReAct 运行时状态。 +""" + +from __future__ import annotations + +import logging +from typing import Any + +from api.ops.store.runs import OpsRunStore +from api.rag_env import supabase_client + +logger = logging.getLogger(__name__) + + +class CheckpointStoreError(RuntimeError): + """Checkpoint 读写或校验失败。""" + + +REQUIRED_STATE_KEYS = ("route", "query", "step", "messages", "tool_evidence") + + +def _validate_react_state(state_json: Any) -> dict[str, Any]: + """校验 checkpoint 状态是否足够恢复 ReAct 循环。 + + 校验通过返回原字典;失败抛出 CheckpointStoreError。 + """ + if not isinstance(state_json, dict): + raise CheckpointStoreError("checkpoint state_json is not a dict") + missing = [k for k in REQUIRED_STATE_KEYS if k not in state_json] + if missing: + raise CheckpointStoreError(f"checkpoint state missing keys: {missing}") + if state_json.get("route") != "react": + raise CheckpointStoreError("checkpoint route is not 'react'") + if not isinstance(state_json.get("messages"), list): + raise CheckpointStoreError("checkpoint state messages is not a list") + if not isinstance(state_json.get("step"), int): + raise CheckpointStoreError("checkpoint state step is not an int") + return state_json + + +def save_checkpoint( + run_id: str, + thread_id: str, + state_json: dict[str, Any], + store: OpsRunStore | None = None, +) -> dict[str, Any]: + """保存 ReAct 运行时 checkpoint。 + + 参数: + run_id: 当前 run id。 + thread_id: session/thread 标识;与 `checkpoint_id` 同义。 + state_json: 运行状态字典。 + store: 可选 OpsRunStore;默认使用全局 supabase_client() 构造。 + """ + target = store if store is not None else OpsRunStore(supabase_client()) + if not hasattr(target, "save_checkpoint"): + raise CheckpointStoreError("store does not support save_checkpoint") + return target.save_checkpoint(run_id, thread_id, state_json) + + +def find_latest_checkpoint_for_session( + session_id: str, + store: OpsRunStore | None = None, +) -> dict[str, Any] | None: + """按 session_id 查找最新的有效 checkpoint(跨 run)。 + + 返回整行(含 run_id / checkpoint_id / state_json / created_at); + 不存在时返回 None。 + """ + target = store if store is not None else OpsRunStore(supabase_client()) + # 防御:部分测试 double 未实现 checkpoint 方法时直接返回 None + if not hasattr(target, "find_latest_checkpoint_for_session"): + return None + return target.find_latest_checkpoint_for_session(session_id) + + +def load_checkpoint( + run_id: str, + thread_id: str, + store: OpsRunStore | None = None, +) -> dict[str, Any] | None: + """读取指定 run + thread 的 checkpoint。 + + 返回状态字典;不存在时返回 None。 + 注意:返回前不做结构校验,由调用方 `resume_react_state` 处理。 + """ + target = store if store is not None else OpsRunStore(supabase_client()) + if not hasattr(target, "load_checkpoint"): + return None + row = target.load_checkpoint(run_id, thread_id) + if not row: + return None + return row.get("state_json") \ No newline at end of file diff --git a/api/ops/store/runs.py b/api/ops/store/runs.py index b1fc321d..9504b1ee 100644 --- a/api/ops/store/runs.py +++ b/api/ops/store/runs.py @@ -216,6 +216,42 @@ def _once() -> dict[str, Any]: return supabase_execute_with_retry(_once) + def load_checkpoint(self, run_id: str, checkpoint_id: str) -> dict[str, Any] | None: + def _once() -> dict[str, Any] | None: + res = ( + self.client.table("ops_run_checkpoints") + .select("*") + .eq("run_id", run_id) + .eq("checkpoint_id", checkpoint_id) + .limit(1) + .execute() + ) + rows = res.data if isinstance(res.data, list) else [] + if rows and isinstance(rows[0], dict): + return rows[0] + return None + + return supabase_execute_with_retry(_once) + + def find_latest_checkpoint_for_session(self, session_id: str) -> dict[str, Any] | None: + """按 session_id(即 checkpoint_id)查找最新的 checkpoint(跨 run)。""" + + def _once() -> dict[str, Any] | None: + res = ( + self.client.table("ops_run_checkpoints") + .select("*") + .eq("checkpoint_id", session_id) + .order("created_at", desc=True) + .limit(1) + .execute() + ) + rows = res.data if isinstance(res.data, list) else [] + if rows and isinstance(rows[0], dict): + return rows[0] + return None + + return supabase_execute_with_retry(_once) + def save_artifact( self, run_id: str, kind: str, payload: dict[str, Any] ) -> dict[str, Any]: diff --git a/docs/harness/invokes/by-task/ops-chat-session-sink-p0-p1/invoke_20260709_1733_30_ops_chat_session_sink_p0_p1_P1-2.md b/docs/harness/invokes/by-task/ops-chat-session-sink-p0-p1/invoke_20260709_1733_30_ops_chat_session_sink_p0_p1_P1-2.md new file mode 100644 index 00000000..e5b23551 --- /dev/null +++ b/docs/harness/invokes/by-task/ops-chat-session-sink-p0-p1/invoke_20260709_1733_30_ops_chat_session_sink_p0_p1_P1-2.md @@ -0,0 +1,89 @@ +--- +hat: "30-execute-code" +task: "ops-chat-session-sink-p0-p1" +phase: "P1-2" +subproject: "ai-ink-brain-api-python" +branch: "task/ops-chat-session-sink-p0-p1" +worktree_root: "ai-ink-brain-api-python/" +date: "2026-07-09" +time: "17:33" +--- + +| 字段 | 值 | +| --- | --- | +| **hat** | 30-execute-code | +| **task** | ops-chat-session-sink-p0-p1 | +| **phase** | P1-2 Checkpoint | +| **subproject** | ai-ink-brain-api-python | +| **branch** | task/ops-chat-session-sink-p0-p1 | +| **worktree_root** | ai-ink-brain-api-python/ | +| **date** | 2026-07-09 | +| **time** | 17:33 | + +## 用户消息快照 + +```text +你正在扮演工作区 Harness「30-execute-code · 执行编码帽」,严格遵循 docs/harness/prompts/30-execute-code.md。 + +**输入(已替换占位符)** +- 主 task 路径(相对 Projects/):`docs/harness/tasks/active/task_ops_chat_session_sink_p0_p1_v1.md` +- 逻辑子仓(相对 Projects/):`ai-ink-brain-api-python` +- Worktree 研发目录(所有 git/pytest/ruff 默认 cwd):`ai-ink-brain-api-python` +- 当前分支:`task/ops-chat-session-sink-p0-p1`(已基于 main fast-forward,包含 P1-1 merge) +- 合并前须跑通的验证命令: + ```bash + pytest tests/ops tests/ops_desk -m "not intent_eval and not intent_benchmark" -q && ruff check api/ops + ``` +- 关联任务审核书面结论路径:`ai-ink-brain-api-python/docs/harness/reviews/task_ops_chat_session_sink_p0_p1_v1_audit_R2_20260708.md` +- 关联 PLAN / 总规:`docs/harness/guides/PLAN_ops_chat_session_sink_p0_p1_v1_zh.md` + +**本棒目标:P1-2 Checkpoint 续跑 · `ops_run_checkpoints` 适配** + +P1-2 具体要求(来自 PLAN §2、§3.1 D7、task §失败路径、§实现备忘): +- 复用已有 `ops_run_checkpoints` 表(见 `supabase/sql/ops_desk_p1_run_schema.sql`)。 +- 新增 `api/ops/store/checkpoints.py` 或扩展 `api/ops/store/runs.py`:实现 `save_checkpoint(run_id, thread_id, state_json)` 与 `load_checkpoint(run_id, thread_id)`。 +- 在 `api/ops/react_loop.py` 的 ReAct 路径中: + - 每步(或关键超步)后保存 checkpoint; + - 同 session 续问时,若存在有效 checkpoint,隐式恢复并续跑; + - 损坏的 checkpoint 按 task §失败路径处理:新 run 冷启动、不 500、记录 `checkpoint.corrupted` event。 +- 优先保证 ReAct 续跑完成;deep 路径可暂不支持 checkpoint(按 PLAN 范围)。 + +**范围限制** +- 只做 P1-2;不改 P1-3 clarify、P1-4 LLM router +- 不改 P1-1 artifact 已交付行为 +- 不改 `harness_runtime` 生产图 +- 不改 Agently lab +- 不改前端代码 + +**test_strategy: required** +- 先写/调整可失败的自动化测试,再改实现 +- 新增/扩展 `tests/ops/test_checkpoint.py` 覆盖: + - save/load checkpoint 成功 + - checkpoint 损坏时冷启动且不 500 + - 同 session 续跑完成(可 mock 超步中断) + - `checkpoint.corrupted` event 被记录 +- 最终验证命令必须绿 + +**失败路径硬性检查** +- task §失败路径已列 `Checkpoint 损坏`:行为 = 新 run 冷启动 · 不 500;可观测 = `checkpoint.corrupted` 日志 / `ops_run_events.kind=checkpoint.corrupted`;可重试 = 否;验证命令 = `pytest tests/ops/test_checkpoint.py -k corrupted` + +**你必须完成** +0. **Invoke 快照(开帽起点)**:将本用户消息全文落盘到 `ai-ink-brain-api-python/docs/harness/invokes/by-task/ops-chat-session-sink-p0-p1/invoke_YYYYMMDD_HHMM_30_ops_chat_session_sink_p0_p1_P1-2.md`(含元数据表 + 快照 fenced code)。同一会话内追问 **不** 再新增快照文件。 +0b. **人工闸**:扫描 task / 关联 reviews 的 human_gate。若任一对本帽(30)为 pending → 仅输出须人改的 gate_id 与路径,拒开工;禁止代填 approved。 +1. 通读 task 全文:头部 gates_before_code、audit_profile、orchestration、chain_prompt、test_strategy / test_strategy_note、failure_paths、验收标准、必读列表、非范围。 +2. 阅读 PLAN §2、§3.1 与关联 SNAPSHOT/gap matrix。 +3. 先读现有代码:`api/ops/store/runs.py`、`api/ops/react_loop.py`、`api/ops/events_schema.py`、`api/ops/chat_service.py`,以及 `supabase/sql/ops_desk_p1_run_schema.sql` 中 `ops_run_checkpoints` 表结构。 +4. 先写失败可复现的测试(`tests/ops/test_checkpoint.py`),再实现 checkpoint save/load/resume/corrupted 处理。 +5. 在 ReAct 路径中合适位置集成 checkpoint。 +6. 执行验证命令,保留可核对输出要点;修复直至通过。 +7. 按 40-self-check.md 将结论与命令摘要回填至 task 正文「### 自检结论(执行者)· P1-2」小节(不要覆盖 P0 或 P1-1 已有结论)。 +8. 对话回复:生成可以完整复制的 Prompt,用于直接交给下一棒 40 自检执行。 +9. **自动 commit**:在输出下一棒 Prompt 且本轮代码/测试/task 自检回填已落盘后,按 HANDOFF_AUTO_COMMIT.md 在 ai-ink-brain-api-python/ commit(仅本轮路径;禁止 git add -A;对话报 short-hash)。 +10. **禁止**自行 push;由 Lead 合并。 + +**输出要求** +- 若拒开工:仅 Markdown 阻塞清单 +- 若执行:diff 摘要、验证命令输出、commit short-hash、下一棒 40 Prompt + +**Judgment(本帽 · 对话末尾必填)**:experience_capture / gate/risk / hat_self +``` diff --git a/docs/harness/invokes/by-task/ops-chat-session-sink-p0-p1/invoke_20260709_1733_40_ops_chat_session_sink_p0_p1_P1-2.md b/docs/harness/invokes/by-task/ops-chat-session-sink-p0-p1/invoke_20260709_1733_40_ops_chat_session_sink_p0_p1_P1-2.md new file mode 100644 index 00000000..d3fd6914 --- /dev/null +++ b/docs/harness/invokes/by-task/ops-chat-session-sink-p0-p1/invoke_20260709_1733_40_ops_chat_session_sink_p0_p1_P1-2.md @@ -0,0 +1,68 @@ +--- +hat: "40-self-check" +task: "ops-chat-session-sink-p0-p1" +phase: "P1-2" +subproject: "ai-ink-brain-api-python" +branch: "task/ops-chat-session-sink-p0-p1" +worktree_root: "ai-ink-brain-api-python/" +date: "2026-07-09" +time: "17:33" +--- + +| 字段 | 值 | +| --- | --- | +| **hat** | 40-self-check | +| **task** | ops-chat-session-sink-p0-p1 | +| **phase** | P1-2 Checkpoint | +| **subproject** | ai-ink-brain-api-python | +| **branch** | task/ops-chat-session-sink-p0-p1 | +| **worktree_root** | ai-ink-brain-api-python/ | +| **date** | 2026-07-09 | +| **time** | 17:33 | + +## 用户消息快照 + +```text +你正在扮演工作区 Harness「40-self-check · 执行者自检帽」,严格遵循 docs/harness/prompts/40-self-check.md。 + +**输入(已替换占位符)** +- 主 task 路径(相对 Projects/):`docs/harness/tasks/active/task_ops_chat_session_sink_p0_p1_v1.md` +- 逻辑子仓(相对 Projects/):`ai-ink-brain-api-python` +- Worktree 研发目录(所有 git/pytest/ruff 默认 cwd):`ai-ink-brain-api-python` +- 当前分支:`task/ops-chat-session-sink-p0-p1`(已基于 main fast-forward,包含 P1-1 merge) +- 合并前须跑通的验证命令: + ```bash + pytest tests/ops tests/ops_desk -m "not intent_eval and not intent_benchmark" -q && ruff check api/ops + ``` +- 上一棒 30 commit:待本 Prompt 落盘后从 30 输出获取 +- 关联任务审核书面结论路径:`ai-ink-brain-api-python/docs/harness/reviews/task_ops_chat_session_sink_p0_p1_v1_audit_R2_20260708.md` +- 关联 PLAN / 总规:`docs/harness/guides/PLAN_ops_chat_session_sink_p0_p1_v1_zh.md` + +**本棒目标:P1-2 自检复核** + +你必须完成: +0. **Invoke 快照(开帽起点)**:将本用户消息全文落盘到 `ai-ink-brain-api-python/docs/harness/invokes/by-task/ops-chat-session-sink-p0-p1/invoke_YYYYMMDD_HHMM_40_ops_chat_session_sink_p0_p1_P1-2.md`(含元数据表 + 快照 fenced code)。同一会话内追问 **不** 再新增快照文件。 +0b. **人工闸**:扫描 task / 关联 reviews 的 human_gate。若任一对本帽(40)为 pending → 仅输出须人改的 gate_id 与路径,拒开工;禁止代填 approved。 +1. 独立阅读 task 正文「### 自检结论(执行者)· P1-2」小节与上一棒 30 invoke 快照 `ai-ink-brain-api-python/docs/harness/invokes/by-task/ops-chat-session-sink-p0-p1/invoke_20260709_1733_30_ops_chat_session_sink_p0_p1_P1-2.md`。 +2. 独立阅读本轮 P1-2 改动代码: + - `api/ops/store/checkpoints.py` + - `api/ops/store/runs.py` + - `api/ops/react_loop.py` + - `tests/ops/test_checkpoint.py` +3. 在 `ai-ink-brain-api-python/` 内完整执行 30 声明的验证命令: + ```bash + pytest tests/ops tests/ops_desk -m "not intent_eval and not intent_benchmark" -q && ruff check api/ops + ``` +4. 额外执行 task §失败路径所列的 `pytest tests/ops/test_checkpoint.py -k corrupted -q`。 +5. 通过 `git diff origin/main...HEAD --stat`(在 `ai-ink-brain-api-python` 内)核对全量变更路径,确认未扩 scope 到 P1-3/4、Session 生产图、Agently lab、前端。 +6. 按 40-self-check.md 将结论与命令摘要回填至 task 正文「### 自检结论(40 复核)· P1-2」小节(不要覆盖 P0、P1-1 或 30 已有结论)。 +7. 对话回复:生成可以完整复制的 Prompt,用于直接交给下一棒 50 独立复检执行。 +8. 自动 commit:在输出下一棒 Prompt 且本轮 task 自检回填已落盘后,按 HANDOFF_AUTO_COMMIT.md 在 `ai-ink-brain-api-python/` commit(仅本轮路径;禁止 git add -A;对话报 short-hash)。 +9. **禁止**自行 push;由 Lead 合并。 + +**输出要求** +- 若拒开工:仅 Markdown 阻塞清单 +- 若执行:diff 摘要、验证命令输出、commit short-hash、下一棒 50 Prompt + +**Judgment(本帽 · 对话末尾必填)**:experience_capture / gate/risk / hat_self +``` diff --git a/docs/harness/invokes/by-task/ops-chat-session-sink-p0-p1/invoke_20260709_1758_40_ops_chat_session_sink_p0_p1_P1-2.md b/docs/harness/invokes/by-task/ops-chat-session-sink-p0-p1/invoke_20260709_1758_40_ops_chat_session_sink_p0_p1_P1-2.md new file mode 100644 index 00000000..2f5d5da2 --- /dev/null +++ b/docs/harness/invokes/by-task/ops-chat-session-sink-p0-p1/invoke_20260709_1758_40_ops_chat_session_sink_p0_p1_P1-2.md @@ -0,0 +1,68 @@ +--- +hat: "40-self-check" +task: "ops-chat-session-sink-p0-p1" +phase: "P1-2" +subproject: "ai-ink-brain-api-python" +branch: "task/ops-chat-session-sink-p0-p1" +worktree_root: "ai-ink-brain-api-python/" +date: "2026-07-09" +time: "17:58" +--- + +| 字段 | 值 | +| --- | --- | +| **hat** | 40-self-check | +| **task** | ops-chat-session-sink-p0-p1 | +| **phase** | P1-2 | +| **subproject** | ai-ink-brain-api-python | +| **branch** | task/ops-chat-session-sink-p0-p1 | +| **worktree_root** | ai-ink-brain-api-python/ | +| **date** | 2026-07-09 | +| **time** | 17:58 | + +## 用户消息快照 + +```text +你正在扮演工作区 Harness「40-self-check · 执行者自检帽」,严格遵循 docs/harness/prompts/40-self-check.md。 + +**输入(已替换占位符)** +- 主 task 路径(相对 Projects/):`docs/harness/tasks/active/task_ops_chat_session_sink_p0_p1_v1.md` +- 逻辑子仓(相对 Projects/):`ai-ink-brain-api-python` +- Worktree 研发目录(所有 git/pytest/ruff 默认 cwd):`ai-ink-brain-api-python` +- 当前分支:`task/ops-chat-session-sink-p0-p1`(已基于 main fast-forward,包含 P1-1 merge) +- 合并前须跑通的验证命令: + ```bash + pytest tests/ops tests/ops_desk -m "not intent_eval and not intent_benchmark" -q && ruff check api/ops + ``` +- 上一棒 30 commit:`api-python@35279435` · `Projects@63efbba` +- 关联任务审核书面结论路径:`ai-ink-brain-api-python/docs/harness/reviews/task_ops_chat_session_sink_p0_p1_v1_audit_R2_20260708.md` +- 关联 PLAN / 总规:`docs/harness/guides/PLAN_ops_chat_session_sink_p0_p1_v1_zh.md` + +**本棒目标:P1-2 自检复核** + +你必须完成: +0. **Invoke 快照(开帽起点)**:将本用户消息全文落盘到 `ai-ink-brain-api-python/docs/harness/invokes/by-task/ops-chat-session-sink-p0-p1/invoke_YYYYMMDD_HHMM_40_ops_chat_session_sink_p0_p1_P1-2.md`(含元数据表 + 快照 fenced code)。同一会话内追问 **不** 再新增快照文件。 +0b. **人工闸**:扫描 task / 关联 reviews 的 human_gate。若任一对本帽(40)为 pending → 仅输出须人改的 gate_id 与路径,拒开工;禁止代填 approved。 +1. 独立阅读 task 正文「### 自检结论(执行者)· P1-2」小节与上一棒 30 invoke 快照 `ai-ink-brain-api-python/docs/harness/invokes/by-task/ops-chat-session-sink-p0-p1/invoke_20260709_1733_30_ops_chat_session_sink_p0_p1_P1-2.md`。 +2. 独立阅读本轮 P1-2 改动代码: + - `api/ops/store/checkpoints.py` + - `api/ops/store/runs.py` + - `api/ops/react_loop.py` + - `tests/ops/test_checkpoint.py` +3. 在 `ai-ink-brain-api-python/` 内完整执行 30 声明的验证命令: + ```bash + pytest tests/ops tests/ops_desk -m "not intent_eval and not intent_benchmark" -q && ruff check api/ops + ``` +4. 额外执行 task §失败路径所列的 `pytest tests/ops/test_checkpoint.py -k corrupted -q`。 +5. 通过 `git diff origin/main...HEAD --stat`(在 `ai-ink-brain-api-python` 内)核对全量变更路径,确认未扩 scope 到 P1-3/4、Session 生产图、Agently lab、前端。 +6. 按 40-self-check.md 将结论与命令摘要回填至 task 正文「### 自检结论(40 复核)· P1-2」小节(不要覆盖 P0、P1-1 或 30 已有结论)。 +7. 对话回复:生成可以完整复制的 Prompt,用于直接交给下一棒 50 独立复检执行。 +8. 自动 commit:在输出下一棒 Prompt 且本轮 task 自检回填已落盘后,按 HANDOFF_AUTO_COMMIT.md 在 `ai-ink-brain-api-python/` commit(仅本轮路径;禁止 git add -A;对话报 short-hash)。 +9. **禁止**自行 push;由 Lead 合并。 + +**输出要求** +- 若拒开工:仅 Markdown 阻塞清单 +- 若执行:diff 摘要、验证命令输出、commit short-hash、下一棒 50 Prompt + +**Judgment(本帽 · 对话末尾必填)**:experience_capture / gate/risk / hat_self +``` diff --git a/docs/harness/invokes/by-task/ops-chat-session-sink-p0-p1/invoke_20260709_1758_50_ops_chat_session_sink_p0_p1_P1-2.md b/docs/harness/invokes/by-task/ops-chat-session-sink-p0-p1/invoke_20260709_1758_50_ops_chat_session_sink_p0_p1_P1-2.md new file mode 100644 index 00000000..3e4583c8 --- /dev/null +++ b/docs/harness/invokes/by-task/ops-chat-session-sink-p0-p1/invoke_20260709_1758_50_ops_chat_session_sink_p0_p1_P1-2.md @@ -0,0 +1,67 @@ +# Invoke Snapshot · 50-independent-reinspect · P1-2 + +| 项 | 内容 | +| --- | --- | +| **task** | `docs/harness/tasks/active/task_ops_chat_session_sink_p0_p1_v1.md` | +| **subproject** | `ai-ink-brain-api-python` | +| **worktree_root** | `ai-ink-brain-api-python/` | +| **hat** | `50-independent-reinspect` | +| **phase** | `P1-2` | +| **date** | `2026-07-09` | +| **timestamp** | `20260709_1758` | +| **branch** | `task/ops-chat-session-sink-p0-p1` | +| **human_gate** | `HG-TASK-DRAFT: approved`, `HG-AUDIT-R1: approved` | + +## 原始 Prompt 快照 + +```text +你正在扮演工作区 Harness「独立复检 + 全局验收帽」,严格遵循: +- docs/harness/prompts/50-independent-reinspect.md(§一 独立复检;§二 全局验收) +- docs/harness/HARNESS_V2_PLAN.md §5(test_strategy: required 时关注测试与实现关系) +- 根目录 AGENTS.md §8、docs/harness/HARNESS_V2_P0_ACCEPTANCE.md(若本次变更触及合并前必绿子仓) + +输入(已由人工替换占位符;若你仍看到 {{…}} 或 REINSPECT_MODE 非三选一字面,须先追问用户,不得开工): +- 主 task 路径(相对 Projects/): + docs/harness/tasks/active/task_ops_chat_session_sink_p0_p1_v1.md +- 子仓根(相对 Projects/;用于理解 diff 与路径): + ai-ink-brain-api-python +- 模式(必须恰好为以下之一:独立复检 / 全局验收 / 两者): + 两者 +- diff 或变更范围说明(全局验收单独模式可写「无」): + git diff origin/main...HEAD +- 任务审核书面结论路径(无则「无」): + ai-ink-brain-api-python/docs/harness/reviews/task_ops_chat_session_sink_p0_p1_v1_audit_R2_20260708.md + +你必须完成: +0. **Invoke 快照(开帽起点)**:在输出下列分节实质性结果之前,先将 **本用户消息全文**(= 本 Prompt、占位符已全部替换)落盘到 `ai-ink-brain-api-python/docs/harness/invokes/by-task/ops-chat-session-sink-p0-p1/invoke_YYYYMMDD_HHMM_50_ops_chat_session_sink_p0_p1_P1-2.md`(含元数据表 + 快照 fenced code)。同一会话内追问 **不** 再新增快照文件。 + +【当模式为「独立复检」或「两者」时 — 对应 hat §一】 +1. 读取 task 内「### 自检结论(执行者)· P1-2」及「### 自检结论(40 复核)· P1-2」;若缺失 → 阻塞首条:要求先跑 TEMPLATE-self-check-invoke + 40。 +2. 输入裁剪:以 diff、命令输出要点、自检验收表为主;避免执行过程长文。 +3. 对 P1-2 每条验收项输出表格:验收项 | pass/fail | 证据(文件:行 / 测试名 / 日志片段)| 备注;fail 须写复现步骤或缺失证据。 + 重点核验: + - `api/ops/store/checkpoints.py` 是否实现 `save_checkpoint` / `load_checkpoint` / `find_latest_checkpoint_for_session` / `_validate_react_state`。 + - `api/ops/store/runs.py` 是否为 `OpsRunStore` 新增 `load_checkpoint` / `find_latest_checkpoint_for_session`。 + - `api/ops/react_loop.py` 是否每步后保存 checkpoint、同 session 续问时隐式恢复、损坏 checkpoint 冷启动且不 500。 + - 损坏 checkpoint 是否记录 `checkpoint.corrupted` event(payload 含 error / session_id / from_run_id)。 + - `tests/ops/test_checkpoint.py` 是否覆盖 save/load、按 session 查找、损坏冷启动、无效 schema、同 session 续跑、失败路径验证。 + - task §失败路径验证命令 `pytest tests/ops/test_checkpoint.py -k corrupted -q` 是否独立通过。 + - 是否未改 P1-3 clarify、P1-4 LLM router、Session 生产图、Agently lab、前端代码。 +4. 执行 `pytest tests/ops tests/ops_desk -m "not intent_eval and not intent_benchmark" -q && ruff check api/ops` 并粘贴输出要点;与 30 / 40 结论交叉核对。 +5. 额外执行 `pytest tests/ops/test_checkpoint.py -k corrupted -q` 并记录结果。 +6. 通过 `git diff --name-status origin/main...HEAD` 核对全量变更路径,确认未扩 scope 到 P1-3/4、Session 生产图、Agently lab、前端。 +7. 汇总阻塞合并项;给出是否建议合并(供维护者决策)。 +8. 禁止:替执行者改代码(除非用户明确要求复检提交 patch);缺口退回需求/审查帽。 + +【当模式为「全局验收」或「两者」时 — 对应 hat §二】 +9. 核对本次 PR 变更是否在 P1-2 声明范围内;行为变更(ReAct checkpoint 续跑 / 损坏冷启动)是否在 task 行为变更节显式记录。 +10. 输出 checklist 表(项 / 状态 / 签注栏「待人工」);不伪造已签核;不跳过 CI 红灯叙事。 + +对话回复:若建议合并且无返工 → 输出「执行路线与 Commit 回溯」(docs/harness/prompts/HANDOFF_CLOSE_TRACE.md),勿编造下一棒 Prompt;若须打回 → 输出下一棒可复制 Prompt(按 50 打回路由表选最短回路;含打回、二次审查、上一棒修复)。 +11. **自动 commit**:落盘后按 docs/harness/prompts/HANDOFF_AUTO_COMMIT.md 分仓 commit。仅对话、零文件变更则不必空提交;用户写明「不要 commit」则跳过。 + +Judgment(本帽 · 对话末尾必填;任一项 warn/fail 须写 judgment_notes): +- experience_capture: 维持 | 建议升级 required | 建议降 n/a | 维持 n/a(≤1 行理由) +- gate/risk: 无 | 须人审: | 证据不足 +- hat_self: pass | pass-with-notes | blocked +``` diff --git a/docs/harness/reviews/task_ops_chat_session_sink_p0_p1_v1_audit_R1_20260709_50_P1-2.md b/docs/harness/reviews/task_ops_chat_session_sink_p0_p1_v1_audit_R1_20260709_50_P1-2.md new file mode 100644 index 00000000..04009a63 --- /dev/null +++ b/docs/harness/reviews/task_ops_chat_session_sink_p0_p1_v1_audit_R1_20260709_50_P1-2.md @@ -0,0 +1,148 @@ +# 50 独立复检 · Ops Chat ← Session 能力下沉 · P1-2 Checkpoint + +| 字段 | 值 | +| --- | --- | +| **task** | `ops-chat-session-sink-p0-p1` | +| **phase** | P1-2 Checkpoint | +| **subproject** | `ai-ink-brain-api-python` | +| **audit_round** | R1 | +| **date** | 2026-07-09 | +| **hat** | 50-reinspect | +| **30 commit** | `api-python@35279435` | +| **40 commit** | `api-python@3ee3c0ee` · `Projects@c8f1410` | + +--- + +## 复核方法 + +1. 独立阅读 task 正文「### 自检结论(执行者)· P1-2」与「### 自检结论(40 复核)· P1-2」小节。 +2. 独立阅读本轮 P1-2 改动代码: + - `api/ops/store/checkpoints.py` + - `api/ops/store/runs.py` + - `api/ops/react_loop.py` + - `tests/ops/test_checkpoint.py` +3. 在 `ai-ink-brain-api-python/` 内完整执行 30 声明的验证命令: + ```bash + pytest tests/ops tests/ops_desk -m "not intent_eval and not intent_benchmark" -q && ruff check api/ops + ``` +4. 额外执行 task §失败路径所列的 `pytest tests/ops/test_checkpoint.py -k corrupted -q`。 +5. 执行 `git diff origin/main...HEAD --stat`(在 `ai-ink-brain-api-python` 内)核对全量变更路径,确认未扩 scope 到 P1-3/4、Session 生产图、Agently lab、前端。 +6. 与 30 commit `api-python@35279435`、40 commit `api-python@3ee3c0ee` / `Projects@c8f1410`、R2 任务审核书面结论逐条核对。 + +--- + +## 命令输出 + +### 主验证命令 + +```text +pytest tests/ops tests/ops_desk -m "not intent_eval and not intent_benchmark" -q +.................................s...................................... [ 23%] +........................................................................ [ 46%] +........................................................................ [ 69%] +........................ss......................sssssss................. [ 93%] +..................... [100%] +=============================== warnings summary ================================ +../../../miniconda3/lib/python3.13/site-packages/fastapi/testclient.py:1 + /Users/cyning/miniconda3/lib/python3.13/site-packages/fastapi/testclient.py:1: StarletteDeprecationWarning: Using `httpx` with `starlette.testclient` is deprecated; install `httpx2` instead. + from starlette.testclient import TestClient as TestClient # noqa + +-- Docs: https://docs.pytest.org/en/stable/howto/capture-warnings.html +=========================== short test summary info ============================ +SKIPPED [1] tests/ops/test_events_schema.py:181: 需要真实 Supabase 连接;本地/CI 环境缺失时跳过 +SKIPPED [2] tests/ops_desk/test_run_schema_p1.py:102: public 中表已存在,跳过写测试以避免破坏数据 +SKIPPED [7] tests/ops_desk/test_schema_p0.py:102: public 中表已存在,跳过写测试以避免破坏数据 +299 passed, 10 skipped, 1 warning in 50.93s + +ruff check api/ops +All checks passed! +``` + +- pytest 退出码:`0` +- ruff 退出码:`0` +- 10 skipped 中:1 个为 `tests/ops/test_events_schema.py::test_append_event_integration_with_real_store`(显式 skip,需真实 Supabase 连接);其余 9 个为 `tests/ops_desk/test_run_schema_p1.py` / `tests/ops_desk/test_schema_p0.py` 中环境感知跳过(表已存在),与 P1-2 改动无关。 + +### 失败路径额外验证 + +```text +pytest tests/ops/test_checkpoint.py -k corrupted -q +... [100%] +3 passed, 4 deselected in 0.50s +``` + +- pytest 退出码:`0` + +--- + +## 与 30 / 40 结论差异核对 + +| 30 / 40 声称项 | 50 独立复核 | 结果 | +| --- | --- | --- | +| 复用已有 `ops_run_checkpoints` 表 | `api/ops/store/runs.py` 已存在 `save_checkpoint`;新增 `load_checkpoint` / `find_latest_checkpoint_for_session`;未改表结构 | 一致 | +| 新增 `api/ops/store/checkpoints.py` | 文件存在;`save_checkpoint` / `load_checkpoint` / `find_latest_checkpoint_for_session` / `CheckpointStoreError` 均实现 | 一致 | +| ReAct 路径每步后保存 checkpoint | `react_loop.py` 在每次非最终 step 后调用 `_try_save_checkpoint`(`react_loop.py:353-367`) | 一致 | +| 同 session 续问隐式恢复并续跑 | `run_react_fallback` 启动时若 `session_id` 存在则调用 `find_latest_checkpoint_for_session`;有效 checkpoint 通过 `_resume_react_state` 恢复并 emit `checkpoint.resume`;单测 `test_same_session_resumes_from_checkpoint` 通过 | 一致 | +| checkpoint 损坏时冷启动且不 500 | `_resume_react_state` 捕获 `CheckpointStoreError`,emit `checkpoint.corrupted` 并返回 None,走冷启动分支;3 个 corrupted 测例均通过 | 一致 | +| `checkpoint.corrupted` event 被记录 | 事件 `event_type=checkpoint.corrupted`,payload 含 `error` / `session_id` / `from_run_id` | 一致 | +| deep 路径未支持 checkpoint | `origin/main...HEAD` diff 未包含 `api/ops/orchestrator/core.py` | 一致 | +| P1-1 artifact 行为未被破坏 | `save_artifact_with_failure_event` 与 `artifact.write_failed` 保留;artifact 测例全绿 | 一致 | +| P0-3/P0-4 行为未被破坏 | transcript / tracing 相关测例全绿 | 一致 | +| 新增单测覆盖 | `tests/ops/test_checkpoint.py` 7 测例全绿 | 一致 | +| 验证命令绿 | 本轮独立跑通 `299 passed, 10 skipped` + ruff 全绿 | 一致 | +| 未扩 scope | 全量 diff 仅 P1-2 checkpoint 相关文件 + invoke;未涉及 P1-3/4、Session 生产图、Agently lab、前端 | 一致 | + +**差异项**:无。 + +**非阻塞观察**: +- `api/ops/store/checkpoints.py` 中定义的 `_validate_react_state` 当前未被引用,实际校验由 `api/ops/react_loop.py` 内 `_validate_react_checkpoint` 完成;属轻微代码冗余,不影响功能与测试。 +- `checkpoints.py` 文件末尾缺少换行符;不影响 ruff / pytest。 + +--- + +## 验收项复核表 + +| 验收项 | 状态 | 证据 | 备注 | +| --- | --- | --- | --- | +| `api/ops/store/checkpoints.py` 存在且实现 `save_checkpoint` / `load_checkpoint` / `find_latest_checkpoint_for_session` | pass | `checkpoints.py:45-97` | `CheckpointStoreError` 同文件定义 | +| `api/ops/store/runs.py` 为 `OpsRunStore` 新增 `load_checkpoint` / `find_latest_checkpoint_for_session` | pass | `runs.py:219-253` | 按 `checkpoint_id=session_id` 跨 run 取最新一条 | +| ReAct 路径每步后保存 checkpoint | pass | `react_loop.py:353-367` | 仅非最终 step、且 `session_id` 存在时保存 | +| 同 session 续问隐式恢复并续跑 | pass | `react_loop.py:231-262` | 查找并恢复 checkpoint;`test_same_session_resumes_from_checkpoint` 断言续跑完成 | +| checkpoint 损坏时冷启动且不 500 | pass | `_resume_react_state` 捕获 `CheckpointStoreError` 并返回 None;3 个 corrupted 测例通过 | — | +| `checkpoint.corrupted` event 被记录 | pass | `react_loop.py:112-123` emit `checkpoint.corrupted`,payload 含 `error` / `session_id` / `from_run_id` | — | +| deep 路径未支持 checkpoint(符合范围) | pass | diff 未含 `api/ops/orchestrator/core.py` | — | +| P1-1 artifact 行为未被破坏 | pass | `react_loop.py:495-506` 仍调用 `save_artifact_with_failure_event`;artifact 测例全绿 | — | +| P0-3/P0-4 行为未被破坏 | pass | transcript / tracing 相关测例全绿;`load_chat_transcript` / `trace_span` 调用路径保留 | — | +| `tests/ops/test_checkpoint.py` 覆盖完整 | pass | 7 测例全绿:save/load、按 session 查找、损坏冷启动、无效 schema、同 session 续跑、失败路径验证 | — | +| 最终验证命令绿 | pass | pytest `299 passed, 10 skipped` + ruff `All checks passed!` | 退出码均为 `0` | +| 失败路径验证命令绿 | pass | `pytest tests/ops/test_checkpoint.py -k corrupted -q` → `3 passed, 4 deselected` | — | +| 未静默扩大 scope | pass | 全量 diff 仅 `api/ops/react_loop.py`、`api/ops/store/checkpoints.py`、`api/ops/store/runs.py`、`tests/ops/test_checkpoint.py` + invoke 文档;未涉及 P1-3/4、Session 生产图、Agently lab、前端 | — | + +--- + +## 阻塞项清单 + +无。 + +--- + +## 合并建议 + +**建议合并**。P1-2 Checkpoint 续跑实现、测试、30 执行、40 自检、50 独立复检均通过,无 scope creep,人工闸 `HG-TASK-DRAFT` / `HG-AUDIT-R1` 已 approved。 + +--- + +## 执行路线与 Commit 回溯 + +| 阶段 | 帽子 | 关键动作 | 落盘工件 | 对应 commit | +|------|------|----------|----------|-------------| +| P1-2 | 30 execute | ReAct checkpoint save/load/resume + corrupted 处理 + 测试 | `api/ops/store/checkpoints.py`, `api/ops/store/runs.py`, `api/ops/react_loop.py`, `tests/ops/test_checkpoint.py` | `api-python@35279435` | +| P1-2 | 40 self-check | 复核 P1-2 验收 | task 内 P1-2 30/40 自检结论 | `api-python@3ee3c0ee` · `Projects@c8f1410` | +| P1-2 | 50 reinspect R1 | 独立复检 + 全局验收 | 本文件 | 待本审查落盘后 commit | + +--- + +## Judgment(50) + +- **experience_capture**: `required` — P1-2 checkpoint save/resume/corrupted 事件模式与状态序列化经验可复用到后续 P1-3/P1-4 及 Session 相关能力。 +- **gate/risk**: 无 — `HG-TASK-DRAFT` / `HG-AUDIT-R1` 均为 `approved`;50 未遇 pending 人工闸。 +- **hat_self**: `pass` — 独立复检与 30/40 结论一致,验证命令绿,输出 pass/fail 表、阻塞项清单、合并建议与执行路线回溯。 diff --git a/tests/ops/test_checkpoint.py b/tests/ops/test_checkpoint.py new file mode 100644 index 00000000..18bec203 --- /dev/null +++ b/tests/ops/test_checkpoint.py @@ -0,0 +1,412 @@ +"""P1-2: ops_run_checkpoints 与 ReAct 续跑单测。""" + +from __future__ import annotations + +from typing import Any + +import pytest + +from api.ops.react_loop import run_react_fallback +from tests.ops_desk._llm_mocks import patch_ops_llm_imports + + +class FakeCheckpointStore: + """内存版 Checkpoint 存储,支持 save/load/按 session 查找。""" + + def __init__(self) -> None: + self.checkpoints: dict[tuple[str, str], dict[str, Any]] = {} + self.events: dict[str, list[dict[str, Any]]] = {} + self.runs: dict[str, dict[str, Any]] = {} + self._checkpoint_counter = 0 + + def save_checkpoint( + self, run_id: str, checkpoint_id: str, state_json: dict[str, Any] + ) -> dict[str, Any]: + self._checkpoint_counter += 1 + row = { + "id": f"chk-{self._checkpoint_counter}", + "run_id": run_id, + "checkpoint_id": checkpoint_id, + "state_json": dict(state_json), + "created_at": f"2026-07-09T17:33:{self._checkpoint_counter:02d}Z", + } + self.checkpoints[(run_id, checkpoint_id)] = row + return row + + def load_checkpoint(self, run_id: str, checkpoint_id: str) -> dict[str, Any] | None: + row = self.checkpoints.get((run_id, checkpoint_id)) + if not row: + return None + return dict(row) + + def find_latest_checkpoint_for_session(self, session_id: str) -> dict[str, Any] | None: + """模拟按 checkpoint_id=session_id 取最新一条(跨 run)。""" + candidates = [ + row for (_run_id, cp_id), row in self.checkpoints.items() if cp_id == session_id + ] + if not candidates: + return None + # 按 created_at 降序(字符串即可,时间戳格式一致) + latest = max(candidates, key=lambda r: r["created_at"]) + return dict(latest) + + def get_run(self, run_id: str) -> dict[str, Any] | None: + return self.runs.get(run_id) + + def append_event( + self, + run_id: str, + agent_role: str, + event_type: str, + payload: dict[str, Any] | None = None, + node_id: str | None = None, + seq: int | None = None, + ) -> dict[str, Any]: + evt: dict[str, Any] = { + "run_id": run_id, + "agent_role": agent_role, + "event_type": event_type, + "payload": payload or {}, + "node_id": node_id, + } + self.events.setdefault(run_id, []).append(evt) + return evt + + def update_run(self, run_id: str, **fields: Any) -> None: + self.runs.setdefault(run_id, {}).update(fields) + + def update_run_metrics_json(self, run_id: str, metrics_json: dict[str, Any]) -> None: + self.update_run(run_id, metrics_json=metrics_json) + + def list_runs_by_session_id(self, session_id: str, *, limit: int = 50) -> list[dict[str, Any]]: + return [r for r in self.runs.values() if r.get("session_id") == session_id][:limit] + + +class FakeReactQueries: + def __init__(self) -> None: + self.issues = { + 545: { + "number": 545, + "title": "Deep demo issue", + "state": "open", + "labels": ["bug"], + "html_url": "https://github.com/MoonshotAI/kimi-code/issues/545", + } + } + self.pulls: dict[int, dict[str, Any]] = {} + + def fetch_issue_by_number(self, number: int) -> dict[str, Any] | None: + return self.issues.get(number) + + def fetch_pull_by_number(self, number: int) -> dict[str, Any] | None: + return self.pulls.get(number) + + def fetch_issues( + self, + days: int = 30, + state: str | None = None, + label: str | None = None, + module: str | None = None, + age: str | None = None, + limit: int = 20, + offset: int = 0, + ) -> tuple[list[dict[str, Any]], int]: + rows = list(self.issues.values()) + if state: + rows = [r for r in rows if r.get("state") == state] + return rows, len(rows) + + def fetch_pulls( + self, + days: int = 30, + state: str | None = None, + ci: str | None = None, + author: str | None = None, + limit: int = 20, + offset: int = 0, + ) -> tuple[list[dict[str, Any]], int]: + return [], 0 + + def cycle_time_metric(self, days: int = 30) -> dict[str, Any]: + return {"metric": "cycle-time", "summary": {"avg_hours": 48.0}} + + def review_time_metric(self, days: int = 30) -> dict[str, Any]: + return {"metric": "review-time", "summary": {"avg_hours": 12.0}} + + def issue_throughput_metric(self, days: int = 30) -> dict[str, Any]: + return {"metric": "issue-throughput", "summary": {"total": 2}} + + def sync_status(self) -> dict[str, Any]: + return {"status": "ok", "cursor": "c1", "as_of": "2026-06-25T00:00:00Z"} + + +@pytest.fixture +def store() -> FakeCheckpointStore: + return FakeCheckpointStore() + + +# --------------------------------------------------------------------------- +# save_checkpoint / load_checkpoint 单元行为 +# --------------------------------------------------------------------------- + + +def test_save_checkpoint_success_and_read(store: FakeCheckpointStore) -> None: + from api.ops.store.checkpoints import load_checkpoint, save_checkpoint + + state = {"step": 1, "messages": [{"role": "user", "content": "hi"}]} + row = save_checkpoint("run-1", "thread-1", state, store=store) + + assert row["run_id"] == "run-1" + assert row["checkpoint_id"] == "thread-1" + assert row["state_json"] == state + assert "id" in row + assert "created_at" in row + + loaded = load_checkpoint("run-1", "thread-1", store=store) + assert loaded == state + + +def test_load_checkpoint_returns_none_when_missing(store: FakeCheckpointStore) -> None: + from api.ops.store.checkpoints import load_checkpoint + + assert load_checkpoint("run-x", "thread-x", store=store) is None + + +def test_find_latest_checkpoint_for_session(store: FakeCheckpointStore) -> None: + from api.ops.store.checkpoints import find_latest_checkpoint_for_session + + save = store.save_checkpoint + save("run-a", "sess-1", {"step": 1}) + save("run-b", "sess-1", {"step": 2}) + save("run-c", "sess-2", {"step": 3}) + + latest = find_latest_checkpoint_for_session("sess-1", store=store) + assert latest is not None + assert latest["run_id"] == "run-b" + assert latest["state_json"] == {"step": 2} + + assert find_latest_checkpoint_for_session("sess-none", store=store) is None + + +# --------------------------------------------------------------------------- +# ReAct 路径:checkpoint 损坏时冷启动且不 500 +# --------------------------------------------------------------------------- + + +def _final_answer_chat(content: str = '{"thought":"直接回答","final_answer":"完成"}') -> Any: + from api.ops.llm.types import LlmCompletionResult, LlmUsage + + def _chat(*args: Any, **kwargs: Any) -> Any: + return LlmCompletionResult( + content=content, + usage=LlmUsage( + provider="siliconflow", + model="Qwen/Qwen2.5-72B-Instruct", + prompt_tokens=10, + completion_tokens=5, + total_tokens=15, + latency_ms=100, + step="react", + ), + ) + + return _chat + + +def test_checkpoint_corrupted_cold_start_records_event( + monkeypatch: pytest.MonkeyPatch, store: FakeCheckpointStore +) -> None: + """损坏 checkpoint → 新 run 冷启动、不抛 500、记录 checkpoint.corrupted event。""" + + def _bad_checkpoint(*args: Any, **kwargs: Any) -> dict[str, Any] | None: + return { + "run_id": "run-prev", + "checkpoint_id": "sess-corrupt", + "state_json": {"foo": "bar"}, + "created_at": "2026-07-09T17:33:00Z", + } + + monkeypatch.setattr( + "api.ops.react_loop.find_latest_checkpoint_for_session", _bad_checkpoint + ) + monkeypatch.setattr("api.ops.react_loop.load_chat_transcript", lambda *args, **kwargs: []) + + patch_ops_llm_imports( + monkeypatch, + chat_completion=_final_answer_chat(), + synthesize_answer=lambda *args, **kwargs: ("综合建议。", None), + synthesize=lambda *args, **kwargs: ("综合建议。", None), + ) + + result = run_react_fallback( + "run-corrupt", "hello", store, FakeReactQueries(), session_id="sess-corrupt", max_steps=2 + ) + + assert result["answer"] + assert result["status"] in ("done", "partial") + events = store.events["run-corrupt"] + corrupted = [e for e in events if e["event_type"] == "checkpoint.corrupted"] + assert len(corrupted) == 1 + assert corrupted[0]["payload"].get("session_id") == "sess-corrupt" + assert "error" in corrupted[0]["payload"] + + # 冷启动应出现 router.decision / run.start 等正常事件 + assert any(e["event_type"] == "run.start" for e in events) + + +def test_checkpoint_corrupted_with_invalid_state_schema( + monkeypatch: pytest.MonkeyPatch, store: FakeCheckpointStore +) -> None: + """state_json 缺少关键字段同样视为 corrupted。""" + + def _bad_checkpoint(*args: Any, **kwargs: Any) -> dict[str, Any] | None: + return { + "run_id": "run-prev", + "checkpoint_id": "sess-missing", + "state_json": { + "route": "react", + "query": "hello", + # 缺少 step / messages / tool_evidence + }, + "created_at": "2026-07-09T17:33:00Z", + } + + monkeypatch.setattr( + "api.ops.react_loop.find_latest_checkpoint_for_session", _bad_checkpoint + ) + monkeypatch.setattr("api.ops.react_loop.load_chat_transcript", lambda *args, **kwargs: []) + + patch_ops_llm_imports( + monkeypatch, + chat_completion=_final_answer_chat(), + synthesize_answer=lambda *args, **kwargs: ("综合建议。", None), + synthesize=lambda *args, **kwargs: ("综合建议。", None), + ) + + result = run_react_fallback( + "run-missing", "hello", store, FakeReactQueries(), session_id="sess-missing", max_steps=2 + ) + + assert result["answer"] + corrupted = [e for e in store.events["run-missing"] if e["event_type"] == "checkpoint.corrupted"] + assert len(corrupted) == 1 + + +# --------------------------------------------------------------------------- +# ReAct 路径:同 session 续跑完成 +# --------------------------------------------------------------------------- + + +def test_same_session_resumes_from_checkpoint( + monkeypatch: pytest.MonkeyPatch, store: FakeCheckpointStore +) -> None: + """第一跑被 max_steps 中断并落 checkpoint;同 session 第二跑隐式恢复并完成。""" + + responses = [ + # 第一跑 step 1:调用工具 + '{"thought":"搜索相关 issue","tool":"fetch_issues","arguments":{}}', + # 第二跑恢复后:直接给出最终答案 + '{"thought":"基于结果回答","final_answer":"已完成续跑"}', + ] + call_index = 0 + + def _sequenced_chat(*args: Any, **kwargs: Any) -> Any: + nonlocal call_index + from api.ops.llm.types import LlmCompletionResult, LlmUsage + + content = responses[call_index] + call_index += 1 + return LlmCompletionResult( + content=content, + usage=LlmUsage( + provider="siliconflow", + model="Qwen/Qwen2.5-72B-Instruct", + prompt_tokens=10, + completion_tokens=5, + total_tokens=15, + latency_ms=100, + step="react", + ), + ) + + patch_ops_llm_imports( + monkeypatch, + chat_completion=_sequenced_chat, + synthesize_answer=lambda *args, **kwargs: ("综合建议。", None), + synthesize=lambda *args, **kwargs: ("综合建议。", None), + ) + + # 为了隔离 checkpoint 行为,不加载历史 transcript + import api.ops.react_loop + + original_load_transcript = api.ops.react_loop.load_chat_transcript + api.ops.react_loop.load_chat_transcript = lambda *args, **kwargs: [] # type: ignore[assignment] + + try: + # 第一跑:max_steps=1,执行一步工具后中断,留下 checkpoint + result1 = run_react_fallback( + "run-1", "hello", store, FakeReactQueries(), session_id="sess-resume", max_steps=1 + ) + assert result1["answer"] + assert any( + e["event_type"] == "agent.tool_call" for e in store.events["run-1"] + ) + assert ("run-1", "sess-resume") in store.checkpoints + + # 第二跑:同 session,应恢复 checkpoint 并一步完成 + result2 = run_react_fallback( + "run-2", "hello", store, FakeReactQueries(), session_id="sess-resume", max_steps=3 + ) + assert result2["answer"] == "已完成续跑" + + resume_events = [e for e in store.events["run-2"] if e["event_type"] == "checkpoint.resume"] + assert len(resume_events) == 1 + assert resume_events[0]["payload"].get("from_run_id") == "run-1" + assert resume_events[0]["payload"].get("step") == 1 + + # 续跑应直接使用 checkpoint 中的 tool_evidence,无需再次 tool_call + assert not any( + e["event_type"] == "agent.tool_call" for e in store.events["run-2"] + ) + finally: + api.ops.react_loop.load_chat_transcript = original_load_transcript # type: ignore[assignment] + + +# --------------------------------------------------------------------------- +# 失败路径硬性检查:checkpoint.corrupted 验证命令 +# --------------------------------------------------------------------------- + + +def test_checkpoint_corrupted_failure_path_verification( + monkeypatch: pytest.MonkeyPatch, store: FakeCheckpointStore +) -> None: + """task §失败路径验证命令:`pytest tests/ops/test_checkpoint.py -k corrupted`。""" + + def _bad_checkpoint(*args: Any, **kwargs: Any) -> dict[str, Any] | None: + return { + "run_id": "run-prev", + "checkpoint_id": "sess-verify", + "state_json": "not-a-dict", + "created_at": "2026-07-09T17:33:00Z", + } + + monkeypatch.setattr( + "api.ops.react_loop.find_latest_checkpoint_for_session", _bad_checkpoint + ) + monkeypatch.setattr("api.ops.react_loop.load_chat_transcript", lambda *args, **kwargs: []) + + patch_ops_llm_imports( + monkeypatch, + chat_completion=_final_answer_chat(), + synthesize_answer=lambda *args, **kwargs: ("综合建议。", None), + synthesize=lambda *args, **kwargs: ("综合建议。", None), + ) + + result = run_react_fallback( + "run-verify", "hello", store, FakeReactQueries(), session_id="sess-verify", max_steps=2 + ) + + assert result["answer"] + corrupted = [e for e in store.events["run-verify"] if e["event_type"] == "checkpoint.corrupted"] + assert len(corrupted) == 1 + assert corrupted[0]["payload"].get("session_id") == "sess-verify"