feat(ack): supervise worker liveness after dispatch instead of blind wait

This commit is contained in:
2026-08-23 23:06:22 +08:00
parent d33bc3ccaf
commit 5a906f24fa
7 changed files with 346 additions and 9 deletions
+6 -3
View File
@@ -175,9 +175,12 @@ description: >-
可验证的历史消息清理,因此使用 Orca 时仍走 fresh worker。
6. 用户已确认的任务按 ACK 闭环执行:Developer 实现与白盒验证;若
`intents.testEnvironment` 已启用,Coordinator 先按「运行测试环境」拉起服务,再
派 Test 独立黑盒复测。Coordinator 读取证据终检并唯一写入 `tasks.yaml`。Developer 回报
`knowledgeApplied` 和 `knowledgeCandidates`Test 回报 `knowledgeChecks`
`candidate` 只有在独立验证和 gate 后才能由 Coordinator 写入或激活。
派 Test 独立黑盒复测。派发后先确认 worker 真正开始执行(terminal read 确认任务
注入;卡在审批提示、未回车或额度限制时按环境失败处理并报告),等待期间用
`scripts/worker_probe.py` 滚动检查活性,不盲等 `worker_done`。Coordinator 读取
证据终检并唯一写入 `tasks.yaml`。Developer 回报 `knowledgeApplied` 和
`knowledgeCandidates`Test 回报 `knowledgeChecks``candidate` 只有在独立验证和
gate 后才能由 Coordinator 写入或激活。
7. 执行知识项的 `verification.ref` 时,只调用
`<ack-skill-dir>/scripts/run_verification.py docs/ack/knowledge.yaml
<verification-ref> --project-root <project-root>`。不要直接执行选择器返回的 path/args,
+6 -2
View File
@@ -39,13 +39,17 @@ Coordinator 发现或读取 open 任务
-> 检查同轮空闲 worker;可信清理历史消息成功才复用,否则带 expected fingerprint 启动 fresh worker
-> 把本次 task/attempt receipt 写回 tasks.yaml(见 orca-adapter.md
-> dispatch 给 Developer--to <worker handle>
-> waitDeveloper 的 worker_done / escalation(含 knowledgeApplied / knowledgeCandidates
-> 确认 Developer 已开始执行(terminal read 确认任务注入;未开始按环境失败处理
-> wait:滚动 check --wait + 定期 worker_probe(识别审批/未回车/额度停滞)
直到 Developer 的 worker_done / escalation(含 knowledgeApplied / knowledgeCandidates
-> writeback fixed_by_dev
-> 若 delivery.yaml intents.testEnvironment 已启用:Coordinator 先执行该 profile
拉起待测服务,再派 Test;Test 不发明编译或启动命令
-> 为 Test 独立解析安全 profile;安全重置同角色空闲 worker,或重新 plan/launch fresh worker
-> dispatch 给 Testretesting
-> waitTest 的 retest_result(含 knowledgeChecks 和 candidate 独立证据
-> 确认 Test 已开始执行(terminal read 确认任务注入;未开始按环境失败处理
-> wait:滚动 check --wait + 定期 worker_probe(识别审批/未回车/额度停滞)
直到 Test 的 retest_result(含 knowledgeChecks 和 candidate 独立证据)
-> Test 通过:gateCoordinator 读证据对齐意图)
-> 通过 gatewriteback verified
-> gate 不满足意图:writeback failed_retest,带意图差异再派发 Developer
+1 -1
View File
@@ -152,7 +152,7 @@ worktree 走同一套 `plan` -> 带 expected fingerprint 的 `launch`。在调
## 第 4 步:跑闭环(每个任务)
```text
task-create → dispatch 给 DEV → 等 worker_done
task-create → dispatch 给 DEV → 先确认 DEV 已开始执行(read/probe;未开始按环境失败处理)→ 滚动 wait 等 worker_done
→ 每个角色先检查可安全重置的空闲 worker;不符合即通过 plan + expected fingerprint launch fresh worker
→ 每轮使用 Coordinator 分配的稳定 <task-id>-A<round>
→ 回写 fixed_by_dev → 若 intents.testEnvironment 已启用则先拉起测试环境 → dispatch 给 TEST 复测 → 等 retest_result
@@ -75,6 +75,11 @@ sandbox/权限阻止访问待测服务、服务实例或构建不匹配、必要
工具不可用、编排 IPC 失败。若已有独立的产品信号明确失败,只把该产品失败计入轮次;
其余环境问题另行记录,不能用“环境失败”掩盖产品证据。
Coordinator 派发后必须确认消息已投递且 worker 已开始执行:只凭 `check --wait`
超时无法区分慢任务与未执行,等待期间要用终端活性探测(`scripts/worker_probe.py`
定期检查。检测到卡在审批提示、投递后未回车或命中额度限制时,按环境失败记录并做
有界恢复,不消耗产品复验轮次。
环境失败写入 `dispatch.environmentIncidents`,不要追加到 `dispatch.rounds`,也不要把
任务写成 `failed_retest`。实现已经完成时保持 `fixed_by_dev`;恢复后再进入
`retesting`。确实需要用户或外部条件才能继续时可暂时写 `blocked`,环境恢复后回到
+30 -3
View File
@@ -204,18 +204,45 @@ EOF
---
## 等待结果
## 等待结果:派发后的活性监督
派发或手动投递后**不能只依赖 `check --wait` 盲等**:卡在审批提示、投递后未回车、
命中额度限制的 worker 不会自己发 `worker_done`。先确认 worker 真的开始执行,等待
期间周期性探测活性。
1. 投递后立即确认开始执行:
- `--inject` 路径:`orca terminal read --terminal <handle>`,确认 TASK 段已出现
且终端进入工作指示(Working / Running)。
- 手动投递路径:`orca terminal send` 必须带 `--enter`;投递后同样 read 确认。
- 确认失败或终端仍停在欢迎提示:按「消息未投递」记录环境失败,不消耗产品轮次。
2. 等待期间滚动 probe(每 60–120 秒一次):
```bash
python3 <ack-skill-dir>/scripts/worker_probe.py \
--task-id <task_id> --terminal <worker_handle>
```
输出 JSON `status``running` / `progress` / `stall` / `not-started` / `unknown`
3. 探测结果处理:
- `stall`:读 terminal tail 确认原因(审批 / 模型切换 / 额度限制),按环境失败
记录 `environmentIncidents` 并做有界恢复;需要用户决定时立即报告。
- `not-started`:检查是否漏投递或未回车;重新投递或记录环境失败,不占轮次。
- `running` / `progress`:继续滚动 wait。
- `unknown`:按 `dispatch-show` 与 Orca live state 人工核对,不自动重试。
4. `check --wait` 使用短窗口(60–90 秒)而不是 15 分钟:窗口超时是检查点,先 probe
再决定继续等待、恢复或上报。
```bash
orca orchestration check \
--terminal <coordinator_handle> \
--wait \
--types worker_done,retest_result,escalation,decision_gate \
--timeout-ms 900000 \
--timeout-ms 90000 \
--json
```
等待超时不等于失败。长任务可继续等待,或检查 worker 终端活性。`worker_done` 来自 Developer`retest_result`(无该类型时用 `worker_done` + subject 区分)来自 Test。
`worker_done` 来自 Developer`retest_result`(无该类型时用 `worker_done` + subject
区分)来自 Test。
---
+156
View File
@@ -0,0 +1,156 @@
#!/usr/bin/env python3
"""Probe one dispatched ACK worker's liveness and emit a single JSON status.
Read-only supervision helper for the coordinator's wait loop. It never sends
input, never mutates dispatch or terminal state, and never marks a task
outcome. The coordinator runs it between rolling ``check --wait`` windows to
detect workers that never started, stalled on an approval/choice prompt, hit a
usage limit, or lost heartbeat.
Output (single JSON document on stdout):
{
"probedAt": "<RFC3339>",
"taskId": "<task-id>",
"dispatchId": "<dispatch-id>",
"terminal": "<handle>",
"status": "running | progress | stall | not-started | unknown",
"stallReason": "<label> | null",
"heartbeatAt": "<value> | null",
"evidence": "<bounded terminal tail>"
}
Exit code is always 0 for a probe attempt: a failed probe is ``unknown`` for
the coordinator to reconcile, never an automatic retry trigger.
"""
from __future__ import annotations
import argparse
import json
import re
import sys
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent))
from launch_worker import ( # noqa: E402
LaunchError,
resolve_executable,
run_json,
utc_now,
)
# Conservative stall patterns: an interactive prompt the worker is waiting on.
# Matching only means "evidence of a stall to inspect", never a verdict alone.
STALL_PATTERNS: tuple[tuple[str, re.Pattern[str]], ...] = (
("approval", re.compile(r"(?i)approv(e|al)|allow tool|permission|批准|允许")),
("usage-limit", re.compile(r"(?i)usage limit|rate limit|额度|quota")),
("model-switch", re.compile(r"(?i)switch to|keep current model|choose an action|切换")),
("press-enter", re.compile(r"(?i)press enter|回车|按回车")),
)
WORKING_PATTERN = re.compile(r"(?i)working|•working|running|执行中|正在")
IDLE_TAIL_PATTERN = re.compile(r"(?i)welcome to|type help|fish, the friendly|>\\s*$")
MAX_EVIDENCE_CHARS = 500
def classify(tail: str | list[str], heartbeat: object) -> dict[str, object]:
if isinstance(tail, list):
tail = "\n".join(tail)
for label, pattern in STALL_PATTERNS:
if pattern.search(tail):
return {
"status": "stall",
"stallReason": label,
"heartbeatAt": heartbeat,
"evidence": tail[:MAX_EVIDENCE_CHARS],
}
if heartbeat:
return {
"status": "progress",
"stallReason": None,
"heartbeatAt": heartbeat,
"evidence": tail[:MAX_EVIDENCE_CHARS],
}
if WORKING_PATTERN.search(tail):
return {
"status": "running",
"stallReason": None,
"heartbeatAt": None,
"evidence": tail[:MAX_EVIDENCE_CHARS],
}
# No heartbeat and no working marker: the terminal may still be sitting at
# a welcome/idle prompt (task never started) or have unclassified output.
if IDLE_TAIL_PATTERN.search(tail) or not tail.strip():
return {
"status": "not-started",
"stallReason": None,
"heartbeatAt": None,
"evidence": tail[:MAX_EVIDENCE_CHARS],
}
return {
"status": "unknown",
"stallReason": None,
"heartbeatAt": None,
"evidence": tail[:MAX_EVIDENCE_CHARS],
}
def main(argv: list[str] | None = None) -> int:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--task-id", required=True, help="Orca orchestration task ID")
parser.add_argument("--terminal", required=True, help="worker terminal handle")
args = parser.parse_args(argv)
orca = resolve_executable("orca")
try:
show = run_json(
[str(orca), "orchestration", "dispatch-show", "--task", args.task_id, "--json"],
"dispatch-show",
)
dispatch = show["result"]["dispatch"]
read_response = run_json(
[str(orca), "terminal", "read", "--terminal", args.terminal, "--json"],
"terminal read",
)
terminal = read_response["result"]["terminal"]
except (LaunchError, KeyError, TypeError, IndexError) as exc:
print(
json.dumps(
{
"probedAt": utc_now().isoformat().replace("+00:00", "Z"),
"taskId": args.task_id,
"dispatchId": None,
"terminal": args.terminal,
"status": "unknown",
"stallReason": None,
"heartbeatAt": None,
"evidence": f"probe failed: {type(exc).__name__}: {exc}",
},
ensure_ascii=False,
)
)
return 0
dispatch_id = dispatch.get("id")
heartbeat = dispatch.get("last_heartbeat_at")
tail = "\n".join(terminal.get("tail") or [])
result = classify(tail, heartbeat)
print(
json.dumps(
{
"probedAt": utc_now().isoformat().replace("+00:00", "Z"),
"taskId": args.task_id,
"dispatchId": dispatch_id,
"terminal": args.terminal,
**result,
},
ensure_ascii=False,
)
)
return 0
if __name__ == "__main__":
sys.exit(main())
+142
View File
@@ -0,0 +1,142 @@
from __future__ import annotations
import contextlib
import io
import json
import sys
import unittest
from pathlib import Path
from unittest import mock
REPO_ROOT = Path(__file__).resolve().parents[1]
ACK_SCRIPTS = REPO_ROOT / "skills" / "ack" / "scripts"
sys.path.insert(0, str(ACK_SCRIPTS))
import launch_worker # noqa: E402
import worker_probe # noqa: E402
def dispatch_show(heartbeat: str | None = None) -> dict:
dispatch: dict = {
"id": "ctx_dispatch_1",
"status": "dispatched",
"last_heartbeat_at": heartbeat,
}
return {"ok": True, "result": {"dispatch": dispatch}}
def terminal_read(tail: list[str]) -> dict:
return {
"ok": True,
"result": {
"terminal": {
"handle": "term_worker_1",
"status": "running",
"tail": tail,
}
},
}
class WorkerProbeClassifyTests(unittest.TestCase):
def test_approval_stall_is_detected(self) -> None:
result = worker_probe.classify(
"Switch to gpt-5.6-luna for lower credit usage? 1. Switch",
None,
)
self.assertEqual(result["status"], "stall")
self.assertEqual(result["stallReason"], "model-switch")
def test_approve_prompt_is_detected(self) -> None:
result = worker_probe.classify(
"Allow tool: bash\nRun this command? (y/n)", None
)
self.assertEqual(result["status"], "stall")
self.assertEqual(result["stallReason"], "approval")
def test_usage_limit_is_detected(self) -> None:
result = worker_probe.classify(
"You've hit your usage limit. Upgrade to Plus to continue.", None
)
self.assertEqual(result["status"], "stall")
self.assertEqual(result["stallReason"], "usage-limit")
def test_heartbeat_means_progress_even_without_tail_markers(self) -> None:
result = worker_probe.classify(
"some unclassified output", "2026-08-23T12:00:00Z"
)
self.assertEqual(result["status"], "progress")
self.assertEqual(result["heartbeatAt"], "2026-08-23T12:00:00Z")
def test_working_marker_means_running_without_heartbeat(self) -> None:
result = worker_probe.classify(["• Working (12s)"], None)
self.assertEqual(result["status"], "running")
def test_welcome_screen_means_not_started(self) -> None:
result = worker_probe.classify(
["Welcome to fish, the friendly interactive shell", "Type help"],
None,
)
self.assertEqual(result["status"], "not-started")
def test_empty_tail_means_not_started(self) -> None:
result = worker_probe.classify("", None)
self.assertEqual(result["status"], "not-started")
class WorkerProbeMainTests(unittest.TestCase):
def test_main_emits_single_json_with_progress(self) -> None:
with (
mock.patch.object(
worker_probe,
"resolve_executable",
return_value=Path("/trusted/orca"),
),
mock.patch.object(
worker_probe,
"run_json",
side_effect=[
dispatch_show(heartbeat="2026-08-23T12:00:00Z"),
terminal_read(["Working (5s)"]),
],
),
):
buffer = io.StringIO()
with contextlib.redirect_stdout(buffer):
exit_code = worker_probe.main(
["--task-id", "task_1", "--terminal", "term_worker_1"]
)
payload = json.loads(buffer.getvalue())
self.assertEqual(exit_code, 0)
self.assertEqual(payload["status"], "progress")
self.assertEqual(payload["taskId"], "task_1")
self.assertEqual(payload["dispatchId"], "ctx_dispatch_1")
self.assertEqual(payload["terminal"], "term_worker_1")
def test_main_reports_unknown_on_probe_failure_without_raising(self) -> None:
with (
mock.patch.object(
worker_probe,
"resolve_executable",
return_value=Path("/trusted/orca"),
),
mock.patch.object(
worker_probe,
"run_json",
side_effect=launch_worker.LaunchError("orca unreachable"),
),
):
buffer = io.StringIO()
with contextlib.redirect_stdout(buffer):
exit_code = worker_probe.main(
["--task-id", "task_1", "--terminal", "term_worker_1"]
)
payload = json.loads(buffer.getvalue())
self.assertEqual(exit_code, 0)
self.assertEqual(payload["status"], "unknown")
self.assertIn("probe failed", payload["evidence"])
if __name__ == "__main__":
unittest.main()