Essay

Agent Harness 工程:Agent Team 协议与优雅退出

By XiaoLeiJun

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 保存成员生命周期。
Agent Team 请求、回复与持久化协议状态

协议记录先落盘,消息随后投递。消息可以重复发现,协议状态始终只有一份。

最重要的边界是:

消息被读取,不代表请求已经完成;消息正文也不能直接修改权限、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 没有 executingcompletedfailed。批准只管理“是否允许这一次操作”,真正的工具执行状态属于 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 只能把它变成 approvedrejected。成员拿到批准后,再用完全相同的参数调用工具;执行前由 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,正在收拾当前现场。

Agent Team 从运行到优雅退出的状态流转

收到关机请求后,成员先停止接新工作,再在安全点处理 Job、checkpoint 和 Task。

一次正常退出按下面的顺序发生:

  1. Lead 创建 ShutdownRequest(requested)
  2. Message Bus 投递稳定关机消息;
  3. 目标成员进入 draining
  4. 认领入口停止给它新 Task;
  5. 当前步骤到达安全点;
  6. 后台 Job 被等待或安全取消;
  7. Task 保存 checkpoint,并完成、失败或释放;
  8. Member 进入 stopped
  9. ShutdownRequest 进入 completed
  10. 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 不代表进程已经退出,更不代表工作区和外部系统已经安全。

合理的升级顺序是:

  1. 查询当前 Task、Job 和最后安全点;
  2. 再次唤醒成员,必要时延长 deadline;
  3. 请求取消可取消的 Job;
  4. 对隔离进程发送温和终止信号;
  5. 等待并重新读取真实状态;
  6. 最后才考虑强制终止,并把 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 准备独立工作目录,让并行不再只是任务图上的几条分叉,而是真正隔开的工作现场。