Agent Harness 工程:Agent Team 协议与优雅退出
Agent Harness 工程:Agent Team 协议与优雅退出
上一篇建立了一个最小 Agent Team。Lead、Backend、Frontend、QA 等成员拥有稳定身份和独立上下文,可以从 Task Store 认领任务,并通过持久化 Message Bus 交换请求、进度和结果。
但“能够发消息”还不等于“已经有协议”。
假设 Lead 想停止 Backend Agent。如果直接结束线程,它可能正在写文件、持有尚未保存 checkpoint 的 Task,或者等待一个仍有副作用的后台 Job。
再考虑另一种情况:Backend Agent 准备执行生产部署。即使 Lead 发来一句“同意”,Harness 仍然要确认:批准的是哪一次操作,参数有没有变化,批准是否过期,以及它是否已经被使用。
这些约束不能只写进 Prompt。我们需要在 Message Bus 之上增加一层由普通代码维护的 Team Protocol。
本章解决四个问题:
- Message 与 Protocol 分别负责什么;
- 如何可靠匹配请求与回复;
- 高风险操作如何获得不可挪用的一次性审批;
- Team Member 如何停止认领、保存现场并优雅退出。
Message 负责投递,Protocol 负责约束
上一篇的 TeamMessage 已经有稳定消息 ID、发送者、收件人、reply_to 和确认时间。它能回答“这条消息有没有送到”,却不能单独回答:
- 这是谁发起的哪一次请求;
- 谁有权回复;
- 回复是否仍在有效期内;
- 这个回复会触发什么状态变化;
- Harness 重启后应该继续等待还是判定超时。
可以把团队里的共享事实分成四类:
| 记录 | 回答的问题 | 保存位置 |
|---|---|---|
TeamMessage | 哪段信息需要投递给谁 | Message Bus |
| Protocol Record | 某次审批或关机请求进行到了哪里 | Protocol Store |
Task | 团队正在完成什么工作 | Task Store |
TeamMemberRecord | 某个成员是否还能认领和执行任务 | Member Store |
最重要的边界是:
消息被读取,不代表请求已经完成;消息正文也不能直接修改 Task、权限或成员状态。
先持久化协议记录,再通过 Message Bus 投递。消息可以重复发现,协议状态只有一份。
不要给所有协议套同一个状态机
审批和关机看起来都是“请求 -> 回复”,但它们管理的对象不同。
高风险审批关心的是:请求是否被批准,以及这份批准是否已经消费。
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(),
}
这里没有通用的 executing 和 failed:
- 工具是否正在执行、成功或失败,由 Tool Execution Record 或后台 Job 管理;
- 成员是否仍在运行,由
TeamMemberRecord管理; - Protocol 只保存本次审批或关机握手自身的状态。
这样每个状态只有一个明确含义,不会把“请求被接受”和“操作执行完成”混在一起。
公共部分只保存身份与关联关系
虽然状态机不同,两种协议仍然共享一些元数据:
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)
ProtocolMeta 负责“谁向谁发起了什么请求”,具体记录负责“这个请求现在是什么状态”。
request_message_id 用来关联 Message Bus。回复消息的 reply_to 必须指向它,但状态查询始终读取 Protocol Store,不能从聊天内容猜测。
稳定 ID 让重试不会复制请求
协议请求可能已经落盘,却在结果返回模型前中断。重试时必须得到同一个 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,
tool_call_id: str,
kind: str,
) -> str:
key = f"{requester_id}:{tool_call_id}:{kind}"
digest = uuid.uuid5(uuid.NAMESPACE_URL, key).hex[:16]
return f"protocol-{digest}"
def protocol_message_id(protocol_id: str, phase: str) -> str:
key = f"{protocol_id}:{phase}"
digest = uuid.uuid5(uuid.NAMESPACE_URL, key).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()
tool_call_id 和 requester_id 由 Harness 注入,不接受模型填写。put_if_absent() 在 Store 锁内比较参与者、Task、timeout 和业务内容;相同请求返回原记录,内容冲突则拒绝。
协议记录要先于消息写入:
1. Protocol Store 写入 requested
2. Message Bus 写入稳定 request_message_id
3. 本地 Queue 唤醒 responder
如果进程在第 1 步后退出,reconciliation 会补投消息;如果在第 2 步后退出,重试写入同一个 Message ID,不会制造第二份请求。Queue 只负责降低延迟,不承担持久化。
回复必须同时匹配身份和请求
收到一句“同意”远远不够。Harness 至少要验证:
reply_to对应哪条请求消息;- Protocol Record 是否存在;
- 当前成员是否正是
responder_id; - 请求是否仍在等待回复;
- deadline 是否已经过去;
- 目标状态是否属于该协议的合法转换。
可以把状态检查写成一个很小的辅助函数:
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 ensure_shutdown_transition(
current: ShutdownStatus,
target: ShutdownStatus,
) -> None:
if target not in SHUTDOWN_TRANSITIONS[current]:
raise ValueError(f"cannot transition {current} to {target}")
这两个函数只检查各自的图结构。身份、deadline、Task 所有权等规则仍然由对应 handler 检查。代码稍有重复,却让两套生命周期保持清楚。
高风险操作:批准的是这一次调用
高风险审批应该绑定:
- 请求成员;
- 当前 Task;
- 工具名称;
- 规范化参数的 digest;
- 审批者;
- 有效期。
第一次调用受保护工具时,Permission Hook 不执行操作,而是创建 OperationApproval(status="requested"),再通知 Lead。
Lead 只能把它变成 approved 或 rejected:
def review_operation_approval(
store: "ApprovalStore",
approval_id: str,
approver_id: str,
approve: bool,
summary: str,
) -> OperationApproval:
with store.lock:
approval = store._get_unlocked(approval_id)
if approval.meta.responder_id != approver_id:
raise PermissionError("current agent cannot review this approval")
target: ApprovalStatus = "approved" if approve else "rejected"
if approval.status == target:
return approval
if deadline_has_passed(approval.meta.deadline_at):
raise ValueError("approval request has expired")
ensure_approval_transition(approval.status, target)
approval.status = target
approval.events.append(
ProtocolAuditEvent(
actor_id=approver_id,
event="approval_reviewed",
summary=redact_sensitive_text(summary),
created_at=utc_now_text(),
)
)
store._save_unlocked(approval)
return approval
成员收到批准后,使用相同参数和 approval_id 再次调用工具。执行前,Harness 重新校验并消费这份批准:
def consume_operation_approval(
state: HarnessState,
approval_id: str,
agent_id: str,
task_id: str,
tool_name: str,
arguments: dict[str, Any],
) -> 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("tool 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.meta.id 可以作为 Tool Execution Record 或外部系统的幂等键。
批准被消费后,无论工具成功还是失败,都不能再次使用。工具结果由现有工具记录或后台 Job 保存,不要把 Approval 又改成 completed 或 failed。如果进程在发出外部请求后崩溃,恢复时先按幂等键查询结果,不能重新申请一次就盲目重做。
审批也不能增加 capability。成员原本没有部署工具或网络权限,即使 Lead 批准,Harness 仍然会在 ensure_profile_allows_tool() 处拒绝。
优雅退出:先进入 draining
成员生命周期保持独立:
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 表示成员仍然存活,但已经停止认领新任务,正在整理退出前的现场。
收到关机请求后,成员先进入 draining,在安全点保存 checkpoint、处理 Job 和释放 Task,最后确认退出。
一次正常关机按下面的顺序进行:
- Lead 创建
ShutdownRequest(status="requested"); - Message Bus 投递稳定的
shutdown_request; - 目标成员在安全点进入
draining; - 认领入口不再给它分配新 Task;
- 当前 Task 保存 checkpoint;
- 正在运行的 Job 被等待或安全取消;
- 没有未决副作用后,Task 被完成、失败或释放;
- MemberRecord 进入
stopped,ShutdownRequest 进入completed; - Runtime 补投稳定的
shutdown_ack,Supervisor 等待线程退出。
接受关机请求时,Protocol 与 MemberRecord 要同步推进:
def begin_member_shutdown(
state: HarnessState,
request_id: str,
member_id: str,
) -> TeamMemberRecord:
with state.assignment_lock:
request = state.shutdown_store.get(request_id)
record = state.member_store.get(member_id)
if request.target_agent_id != member_id:
raise PermissionError("shutdown request targets another member")
if request.meta.responder_id != member_id:
raise PermissionError("current member cannot accept this request")
if deadline_has_passed(request.meta.deadline_at):
raise ValueError("shutdown request has expired")
if record.status == "draining":
if record.shutdown_request_id != request_id:
raise ValueError("member is draining for another request")
return record
if record.status != "running":
raise ValueError(f"member cannot drain from {record.status}")
with state.shutdown_store.lock:
request = state.shutdown_store._get_unlocked(request_id)
ensure_shutdown_transition(request.status, "draining")
request.status = "draining"
request.events.append(
ProtocolAuditEvent(
actor_id=member_id,
event="shutdown_draining",
summary="member stopped accepting new tasks",
created_at=utc_now_text(),
)
)
state.shutdown_store._save_unlocked(request)
record.status = "draining"
record.accepting_new_tasks = False
record.shutdown_request_id = request_id
record.updated_at = utc_now_text()
state.member_store.save(record)
return record
claim_next_compatible_task() 也要持有 assignment_lock,同时确认:
(
record.status == "running"
and record.accepting_new_tasks
and not state.shutdown_store.has_active_request(record.agent_id)
)
这样即使进程在 ShutdownRequest 落盘后、MemberRecord 更新前退出,恢复期间也不会再认领任务。文件型 Store 通过 reconciliation 补齐两份记录;多进程部署应使用数据库事务或条件更新。
只在安全点停止
安全点是 Harness 知道当前状态完整、可以暂停的位置,例如:
- 模型调用开始前或返回后;
- 工具调用开始前或完成后;
- 文件完成一次原子替换后;
- 后台 Job 状态已经落盘后;
- Task checkpoint 已经写入后。
不要在线程执行一半时从外部直接抛异常。长时间运行的命令应该使用第 14 章的后台 Job,而不是阻塞成员线程。
Task Session 在每个步骤之间检查成员状态:
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 交给另一个 Agent。
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")
running_jobs = state.job_manager.list_running_for_agent(profile.id)
for job in running_jobs:
if state.job_manager.can_cancel(job.id):
state.job_manager.request_cancel(job.id, requested_by=profile.id)
if running_jobs:
return DrainResult(
ready_to_stop=False,
summary="waiting for background jobs to reach a terminal state",
waiting_job_ids=tuple(job.id for job in running_jobs),
)
if record.current_task_id is not None:
task = state.task_store.get(record.current_task_id)
if task.status == "in_progress":
state.task_store.checkpoint(
task_id=task.id,
agent_id=profile.id,
claim_token=claim_token,
checkpoint=build_recovery_checkpoint(state, profile.id, task.id),
artifacts=list_pending_artifacts(state, task.id),
)
state.task_store.release(
task_id=task.id,
agent_id=profile.id,
claim_token=claim_token,
)
return DrainResult(
ready_to_stop=True,
summary="jobs settled and task state saved",
)
发出 Job cancel 请求不等于 Job 已经停止。因此,只要 running_jobs 仍然非空,本轮 drain 就继续等待,下一轮再读取持久化状态。
Job 成功也不代表 Task 整体完成。如果成员还没有验证全部结果,应该 checkpoint 并 release,不能为了退出而把 Task 强行标记为 completed。
完成关机握手
当没有运行中的 Job,Task 现场也已保存后,Runtime 完成两份状态:
def finish_member_shutdown(
state: HarnessState,
member_id: str,
request_id: str,
summary: str,
) -> None:
with state.assignment_lock:
record = state.member_store.get(member_id)
if record.status == "stopped":
return
if record.status != "draining":
raise ValueError("member must be draining before it can stop")
if record.shutdown_request_id != request_id:
raise ValueError("shutdown request does not match member state")
with state.shutdown_store.lock:
request = state.shutdown_store._get_unlocked(request_id)
ensure_shutdown_transition(request.status, "completed")
request.status = "completed"
request.events.append(
ProtocolAuditEvent(
actor_id=member_id,
event="shutdown_completed",
summary=summary,
created_at=utc_now_text(),
)
)
state.shutdown_store._save_unlocked(request)
record.status = "stopped"
record.accepting_new_tasks = False
record.current_task_id = None
record.updated_at = utc_now_text()
state.member_store.save(record)
state.protocol_outbox.enqueue(
protocol_id=request_id,
phase="shutdown_ack",
)
Outbox 使用 protocol_message_id(request_id, "shutdown_ack") 写入 Lead 邮箱。如果成员在线程退出前来不及投递,恢复扫描仍会补发;真正的完成状态已经保存在 Shutdown Store 中。
超时不等于可以立即强杀
协议监控器可以确定性处理 deadline:
requested超时,进入expired;draining超时,进入timed_out。
timed_out 只说明关机握手没有按时完成,不代表进程已经退出,也不代表工作区和外部系统已经安全。成员后来完成排空时,ShutdownRequest 仍然可以从 timed_out 进入 completed,保留“曾经超时”和“最终完成”两个事实。
升级顺序通常是:
- 查询成员当前 Task、Job 和最后一个安全点;
- 再次唤醒成员,必要时延长 deadline;
- 请求取消可取消的 Job;
- 对隔离子进程发送可处理的终止信号;
- 等待并再次检查真实状态;
- 最后才考虑强制终止,并把 Task 标记为需要检查。
强制终止后不要立即重跑。新的执行者应该先检查 workspace、后台进程、Task checkpoint 和外部幂等记录。
重启后如何恢复
Harness 启动时根据持久化状态执行固定动作:
| 记录状态 | 恢复动作 |
|---|---|
Approval requested | 补投审批消息,检查 deadline |
Approval approved | 通知请求成员;等待相同参数调用消费 |
Approval consumed | 查询 Tool/Job 或外部幂等结果,不重新批准 |
Shutdown requested | 补投关机消息 |
Shutdown draining / timed_out | 恢复成员排空循环,检查 Task、Job 和升级动作 |
Shutdown completed / cancelled / expired | 补投尚未送达的回复,不再次执行动作 |
如果 MemberRecord 显示 draining,却找不到对应的 ShutdownRequest,Harness 应把成员标记为 unhealthy 并通知 Lead,不能自行猜测它已经完成关机。
恢复逻辑不依赖模型“记得发生了什么”。消息投递、状态转换、过期和重复抑制都由 Runtime 完成,只有审批判断等业务决策才需要模型参与。
审计只记录决策事实
协议审计建议保存:
- Protocol ID、kind 和关联 Task;
- 请求者与响应者的稳定身份;
- 参数 digest 和脱敏风险摘要;
- 每次合法状态转换与时间;
- Message ID、Job ID 和 artifact 路径;
- 审批是否被消费;
- 关机时释放了哪个 Task、等待了哪些 Job。
不应该保存模型隐藏推理、完整 messages、访问令牌或未经脱敏的环境变量。大段日志继续放进 Artifact Store,协议记录只保存路径和摘要。
文件型实现可以把少量 events 与协议记录一起原子保存。需要集中审计时,再通过持久化 outbox 异步导出,避免审计服务暂时不可用时阻塞本地状态转换。
暴露给 Agent 的工具
工具应该按具体协议命名,不提供一个可以随意填写目标状态的“万能回复”接口:
| 工具 | 作用 |
|---|---|
team_request_shutdown | Lead 请求某个成员优雅退出 |
team_request_approval | 为确定的高风险工具调用申请审批 |
team_review_approval | Lead 批准或拒绝审批请求 |
team_protocol_status | 查看自己参与的协议记录和审计摘要 |
team_cancel_protocol | 在动作开始前取消自己发起的请求 |
actor_id、tool_call_id 和当前 Task 所有权都由 Harness 注入。模型不能声称自己代表 Lead,也不能直接把状态改成 consumed、completed 或 timed_out。
常见坑
第一,把消息确认当成请求完成。 acknowledged_at 只说明消息被处理,协议状态仍要读取对应 Store。
第二,让所有协议共用一套状态。 审批消费、工具执行和成员退出是不同生命周期,应该分别建模。
第三,只靠自然语言匹配回复。 Harness 还要校验 request ID、responder、deadline 和合法转换。
第四,审批只绑定工具名。 “允许执行 bash”范围过大,必须绑定参数 digest、Task、成员、有效期和单次消费。
第五,审批能够增加 capability。 审批只能确认已有权限范围内的高风险动作。
第六,收到 shutdown 就立即释放 Task。 如果后台副作用还在运行,另一个 Agent 可能重复执行。
第七,把请求取消等同于操作取消。 Approval 已经 consumed 后,取消必须交给具体 Tool 或 Job。
第八,关机超时后直接强杀。 先检查安全点、Task、Job 和外部系统,强制终止只能是最后手段。
第九,只把回复保存在 Message Bus。 消息可能延迟,协议记录才是状态事实。
第十,恢复时直接重做已消费操作。 应先根据幂等键查询工具或外部系统结果。
小结
本文为 Agent Team 增加了两类最小协议:
- Message Bus 负责投递,Protocol Store 负责请求状态;
- 公共元数据只保存身份、关联 Task、消息 ID 和 deadline;
- 审批与关机分别维护自己的状态机;
- 高风险批准绑定成员、Task、工具、参数 digest 和有效期;
- Approval 被消费后不再复用,工具结果由 Tool 或 Job 管理;
- 关机通过
running -> draining -> stopped保存成员生命周期; - draining 成员停止认领,在安全点处理 Job、checkpoint 和 Task;
- 协议超时与强制终止分开处理;
- 稳定 ID、持久化 outbox 和 reconciliation 负责重复投递与重启恢复;
- 审计只保存决策事实,不保存完整对话和模型推理。
现在,Agent Team 不仅能够通信,还能围绕审批和生命周期执行可验证、可恢复的协作。
但多个成员即使认领了不同 Task,只要仍在同一个项目目录里写文件,就可能互相覆盖修改、污染构建产物,甚至让一个 Task 读到另一个 Task 尚未完成的代码。
下一章讲 Worktree Isolation:为每个 Task 创建独立的 Git worktree 和工作目录,让不同 Agent 的文件修改、依赖安装与测试产物互不干扰,再通过 commit、diff、冲突检测和集成流程把结果安全地合并回来。