Agent Harness 工程:Agent Team 协议与优雅退出
Agent Harness 工程:Agent Team 协议与优雅退出
Backend Agent 正准备执行一次生产部署。
它在 Message Bus 里问 Lead:“可以开始吗?”
几秒后,邮箱里出现一条回复:
同意。
这两个字远远不够。
Harness 还要知道:批准的是哪一个工具、哪一组参数、哪个 Task、有效到什么时候,以及这份批准有没有已经用过。
同一时刻,Lead 又想让 Frontend Agent 退出。如果直接结束线程,成员可能正在写文件,可能持有尚未 checkpoint 的 Task,也可能等待一个仍会产生副作用的后台 Job。
上一篇让团队能够发消息,但“消息送达”只解决了沟通。审批和退出还需要更严格的承诺:
谁向谁发起了什么请求
请求现在走到哪一步
哪些回复有效
状态能否继续变化
崩溃后从哪里恢复
这些约束不能只写进 Prompt。我们需要在 Message Bus 之上,再放一层由普通代码维护的 Team Protocol。
Message 是信封,Protocol 是承诺
TeamMessage 能回答:
这条内容是谁发的
要送给谁
是否已经被收件人读取
Protocol Record 则要回答:
这是哪一次请求
谁有权回复
deadline 是否已过
当前状态是什么
下一步允许进入哪个状态
因此,团队里的事实至少分成四份:
- Message Bus 保存投递;
- Protocol Store 保存审批或关机握手;
- Task Store 保存工作状态;
- Member Store 保存成员生命周期。
协议记录先落盘,消息随后投递。消息可以重复发现,协议状态始终只有一份。
最重要的边界是:
消息被读取,不代表请求已经完成;消息正文也不能直接修改权限、Task 或成员状态。
一条“我做完了”不会完成 Task,一条“同意”也不会凭空成为有效审批。
不同承诺,需要不同状态机
审批和关机都长得像“请求 -> 回复”,却在管理完全不同的事情。
审批关心的是批准与消费:
requested -> approved -> consumed
| |
| +-> cancelled / expired
+-> rejected / cancelled / expired
关机则关心成员有没有排空现场:
requested -> draining -> completed
| |
| +-> timed_out -> completed
+-> cancelled / expired
它们应该分别拥有自己的状态:
from typing import Literal
ApprovalStatus = Literal[
"requested",
"approved",
"rejected",
"consumed",
"cancelled",
"expired",
]
ShutdownStatus = Literal[
"requested",
"draining",
"completed",
"cancelled",
"expired",
"timed_out",
]
APPROVAL_TRANSITIONS: dict[
ApprovalStatus,
set[ApprovalStatus],
] = {
"requested": {
"approved",
"rejected",
"cancelled",
"expired",
},
"approved": {
"consumed",
"cancelled",
"expired",
},
"rejected": set(),
"consumed": set(),
"cancelled": set(),
"expired": set(),
}
SHUTDOWN_TRANSITIONS: dict[
ShutdownStatus,
set[ShutdownStatus],
] = {
"requested": {
"draining",
"cancelled",
"expired",
},
"draining": {"completed", "timed_out"},
"timed_out": {"completed"},
"completed": set(),
"cancelled": set(),
"expired": set(),
}
这里有两个容易忽略的细节。
第一,Approval 没有 executing、completed 和 failed。批准只管理“是否允许这一次操作”,真正的工具执行状态属于 Tool Execution Record 或 Job。
第二,Shutdown 的 timed_out 仍然可以进入 completed。超时只说明没有按期排空,不代表成员已经停止;它稍后安全退出时,系统仍要记录最终完成。
与其发明一套看似通用、实际含义模糊的协议状态,不如让每种协议只表达自己的生命周期。
公共字段只保存“谁和哪一次”
两种协议可以共享身份与关联信息:
from dataclasses import asdict, dataclass, field
@dataclass(frozen=True)
class ProtocolMeta:
id: str
kind: Literal[
"approve_operation",
"shutdown_member",
]
requester_id: str
responder_id: str
request_message_id: str
task_id: str | None
timeout_seconds: int
deadline_at: str
created_at: str
@dataclass(frozen=True)
class ProtocolAuditEvent:
actor_id: str
event: str
summary: str
created_at: str
@dataclass
class OperationApproval:
meta: ProtocolMeta
status: ApprovalStatus
tool_name: str
arguments_digest: str
risk_summary: str
events: list[ProtocolAuditEvent] = field(
default_factory=list
)
def to_dict(self) -> dict:
return asdict(self)
@dataclass
class ShutdownRequest:
meta: ProtocolMeta
status: ShutdownStatus
target_agent_id: str
reason: str
waiting_job_ids: list[str] = field(
default_factory=list
)
events: list[ProtocolAuditEvent] = field(
default_factory=list
)
def to_dict(self) -> dict:
return asdict(self)
request_message_id 只负责把协议关联回 Message Bus。查询当前状态时,始终读取 Protocol Store,不能从几条自然语言回复里猜。
每一次请求都要有稳定身份
协议记录可能已经落盘,工具响应却在回模型途中中断。重试时不能生成第二份审批或第二次关机请求。
稳定 ID 可以由 Harness 注入的 requester_id、持久化工具 invocation_id 和协议类型生成。Provider 的 tool_call_id 仍用于消息闭环,但不承担跨会话幂等性:
import hashlib
import json
import uuid
from typing import Any
def canonical_json(value: Any) -> str:
return json.dumps(
value,
ensure_ascii=False,
sort_keys=True,
separators=(",", ":"),
)
def protocol_id(
requester_id: str,
invocation_id: str,
kind: str,
) -> str:
digest = uuid.uuid5(
uuid.NAMESPACE_URL,
f"{requester_id}:{invocation_id}:{kind}",
).hex[:16]
return f"protocol-{digest}"
def protocol_message_id(
protocol_id: str,
phase: str,
) -> str:
digest = uuid.uuid5(
uuid.NAMESPACE_URL,
f"{protocol_id}:{phase}",
).hex[:16]
return f"message-{digest}"
def arguments_digest(
tool_name: str,
arguments: dict[str, Any],
) -> str:
normalized = normalize_tool_arguments(
tool_name,
arguments,
)
payload = {
"tool_name": tool_name,
"arguments": normalized,
}
return hashlib.sha256(
canonical_json(payload).encode("utf-8")
).hexdigest()
协议请求遵循固定顺序:
1. Protocol Store 写入 requested
2. Message Bus 写入稳定 request_message_id
3. Queue 唤醒 responder
如果进程在第一步后退出,reconciliation 会补投消息;如果在第二步后退出,重试只会写入同一个 Message ID。
Queue 负责快,Store 负责不丢。
批准的是这一次调用,不是一个宽泛念头
高风险审批必须绑定:
- 请求成员;
- 当前 Task;
- 工具名称;
- 规范化参数的 digest;
- 审批者;
- 有效期;
- 单次消费。
Backend 第一次请求部署时,Permission Hook 不执行工具,而是创建 OperationApproval(status="requested")。
Lead 只能把它变成 approved 或 rejected。成员拿到批准后,再用完全相同的参数调用工具;执行前由 Harness 消费审批:
def ensure_approval_transition(
current: ApprovalStatus,
target: ApprovalStatus,
) -> None:
if target not in APPROVAL_TRANSITIONS[current]:
raise ValueError(
f"cannot transition {current} to {target}"
)
def consume_operation_approval(
state: HarnessState,
approval_id: str,
agent_id: str,
task_id: str,
tool_name: str,
arguments: dict,
) -> str:
actual_digest = arguments_digest(
tool_name,
arguments,
)
with state.assignment_lock:
ensure_agent_owns_task(
state,
agent_id,
task_id,
)
ensure_profile_allows_tool(
state,
agent_id,
tool_name,
)
with state.approval_store.lock:
approval = (
state.approval_store._get_unlocked(
approval_id
)
)
if approval.status != "approved":
raise PermissionError(
"approval is not active"
)
if approval.meta.requester_id != agent_id:
raise PermissionError(
"approval belongs to another agent"
)
if approval.meta.task_id != task_id:
raise PermissionError(
"approval belongs to another task"
)
if approval.tool_name != tool_name:
raise PermissionError(
"approved tool does not match"
)
if (
approval.arguments_digest
!= actual_digest
):
raise PermissionError(
"arguments changed after approval"
)
if deadline_has_passed(
approval.meta.deadline_at
):
raise PermissionError(
"approval has expired"
)
ensure_approval_transition(
approval.status,
"consumed",
)
approval.status = "consumed"
approval.events.append(
ProtocolAuditEvent(
actor_id=agent_id,
event="approval_consumed",
summary=(
"approved operation started"
),
created_at=utc_now_text(),
)
)
state.approval_store._save_unlocked(
approval
)
return approval.meta.id
审批一旦消费,无论工具成功还是失败,都不能再次使用。
返回的 Approval ID 应进入 Tool Execution Record,必要时也作为外部业务幂等键。若进程在操作开始前后崩溃,恢复器先查询执行记录和外部状态,而不是再申请一份审批盲目重做。
审批也不能增加 capability。一个本来没有部署工具或网络权限的成员,即使收到 Lead 的“批准”,仍然会被 ensure_profile_allows_tool() 拒绝。
优雅退出从停止接新工作开始
成员生命周期与 Shutdown Protocol 仍然分开:
MemberStatus = Literal[
"running",
"draining",
"stopped",
"unhealthy",
]
@dataclass
class TeamMemberRecord:
agent_id: str
status: MemberStatus
accepting_new_tasks: bool
current_task_id: str | None = None
shutdown_request_id: str | None = None
updated_at: str = ""
draining 表示成员还活着,但已经不再认领新 Task,正在收拾当前现场。
收到关机请求后,成员先停止接新工作,再在安全点处理 Job、checkpoint 和 Task。
一次正常退出按下面的顺序发生:
- Lead 创建
ShutdownRequest(requested); - Message Bus 投递稳定关机消息;
- 目标成员进入
draining; - 认领入口停止给它新 Task;
- 当前步骤到达安全点;
- 后台 Job 被等待或安全取消;
- Task 保存 checkpoint,并完成、失败或释放;
- Member 进入
stopped; - ShutdownRequest 进入
completed; - Outbox 向 Lead 补投
shutdown_ack。
协议记录和成员记录必须相互校验。即使进程在 ShutdownRequest 落盘后、MemberRecord 更新前退出,恢复时也不能再为这个成员分配新 Task。
文件型 Store 可以通过 reconciliation 补齐两份记录;多进程部署应把状态推进放进数据库事务或条件更新。
只在系统看得懂的地方停下
安全点是状态已经完整落盘、可以暂停的位置:
- 模型调用之前或返回之后;
- 一次工具调用开始之前或完成之后;
- 文件完成原子替换之后;
- Job 状态已经持久化之后;
- Task checkpoint 已经写入之后。
不要从外部在线程执行中间直接抛异常。长命令应该交给后台 Job,让成员线程可以在步骤之间观察关机状态。
class DrainRequested(Exception):
pass
def ensure_member_can_continue(
state: HarnessState,
agent_id: str,
) -> None:
record = state.member_store.get(agent_id)
if record.status == "draining":
raise DrainRequested(
record.shutdown_request_id
or "shutdown requested"
)
if record.status != "running":
raise RuntimeError(
f"member is not runnable: {record.status}"
)
def run_task_step(
state: HarnessState,
profile: AgentProfile,
session: TaskSession,
) -> None:
ensure_member_can_continue(
state,
profile.id,
)
response = call_model_for_task_step(
state,
profile,
session,
)
ensure_member_can_continue(
state,
profile.id,
)
for tool_call in response.tool_calls:
ensure_member_can_continue(
state,
profile.id,
)
execute_tool_call(
state,
profile,
session,
tool_call,
)
session.save_checkpoint_if_needed()
如果关机请求在工具执行中到达,当前工具先完成自己的原子边界,再进入 drain。
真正不可中断的外部操作仍在运行时,成员必须保持 draining。为了尽快退出而把 Task 交给另一位,很可能造成重复副作用。
Drain 时,先处理 Job,再释放 Task
排空顺序可以由 Runtime 确定,不需要再问模型“要不要保存”:
@dataclass(frozen=True)
class DrainResult:
ready_to_stop: bool
summary: str
waiting_job_ids: tuple[str, ...] = ()
def drain_member_once(
state: HarnessState,
profile: AgentProfile,
claim_token: str,
) -> DrainResult:
record = state.member_store.get(profile.id)
if record.status != "draining":
raise ValueError("member is not draining")
jobs = state.job_manager.list_running_for_agent(
profile.id
)
for job in jobs:
if state.job_manager.can_cancel(job.id):
state.job_manager.request_cancel(
job.id,
requested_by=profile.id,
)
if jobs:
return DrainResult(
ready_to_stop=False,
summary="waiting for background jobs",
waiting_job_ids=tuple(
job.id for job in jobs
),
)
if record.current_task_id is not None:
task = state.task_store.get(
record.current_task_id
)
if task.status == "in_progress":
checkpoint_task(
store=state.task_store,
task_id=task.id,
agent_id=profile.id,
claim_token=claim_token,
summary=build_recovery_checkpoint(
state,
profile.id,
task.id,
),
artifacts=list_pending_artifacts(
state,
task.id,
),
)
release_task(
store=state.task_store,
task_id=task.id,
agent_id=profile.id,
claim_token=claim_token,
)
return DrainResult(
ready_to_stop=True,
summary="jobs settled and task saved",
)
发出 cancel 请求不代表 Job 已经停止。只要 Job Store 仍显示 running,drain 就继续等待。
Job 成功也不意味着 Task 整体完成。如果成员还没有验证所有结果,应该 checkpoint 并 release,不能为了关机好看而强行写成 completed。
超时是一声警报,不是一把斧头
Protocol Monitor 可以确定性处理 deadline:
Approval requested 超时 -> expired
Shutdown requested 超时 -> expired
Shutdown draining 超时 -> timed_out
但 timed_out 不代表进程已经退出,更不代表工作区和外部系统已经安全。
合理的升级顺序是:
- 查询当前 Task、Job 和最后安全点;
- 再次唤醒成员,必要时延长 deadline;
- 请求取消可取消的 Job;
- 对隔离进程发送温和终止信号;
- 等待并重新读取真实状态;
- 最后才考虑强制终止,并把 Task 标记为需要检查。
强制终止后也不能立刻重跑。新的执行者要先检查 workspace、Job、checkpoint 和外部幂等记录。
崩溃以后,协议从 Store 继续
Harness 重启时可以按状态恢复:
Approval requested
补投审批消息,检查 deadline
Approval approved
通知请求成员,等待相同参数消费
Approval consumed
查询 Tool、Job 或外部幂等结果,不重新批准
Shutdown requested
补投关机消息
Shutdown draining / timed_out
恢复排空循环,检查 Task 与 Job
Shutdown completed
补投未送达的 ack,不重复执行动作
如果 MemberRecord 显示 draining,却找不到对应 ShutdownRequest,Harness 应把成员标记为 unhealthy 并通知 Lead,不能猜测它已经安全退出。
恢复逻辑不依赖模型记忆。投递、状态转换、过期和重复抑制都由 Runtime 完成,模型只参与真正需要判断的审批决策。
小结
这一章让 Agent Team 从“会发消息”走向“能履行协议”:
- Message Bus 管投递,Protocol Store 管请求状态;
- 审批和关机使用各自清晰的状态机;
- 稳定 ID 与 Outbox 抵抗重复请求和崩溃窗口;
- 高风险审批绑定成员、Task、工具、参数摘要和有效期;
- Approval 单次消费,执行结果交给 Tool 或 Job;
- 成员先进入 draining,再处理 Job、checkpoint 和 Task;
- 超时与强制终止分开处理;
- 重启后从持久化状态继续,而不是回放聊天猜测。
现在,团队不仅能开始工作,也知道怎样体面而可靠地停下来。
但只要多个成员仍在同一个项目目录里写代码,它们就可能在彼此看不见的地方覆盖文件、污染构建产物,甚至让一项测试验证了另一项尚未完成的修改。
下一章加入 Worktree Isolation:为每个 Task 准备独立工作目录,让并行不再只是任务图上的几条分叉,而是真正隔开的工作现场。