Essay

Agent Harness 工程:后台任务与异步通知

By XiaoLeiJun

Agent Harness 工程:后台任务与异步通知

Agent 修完搜索接口,接着运行完整测试:

pytest

终端安静下来,光标不再闪动。十几分钟里,模型得不到工具结果,Agent Loop 也无法继续。

它明明还可以检查文档、整理变更说明,甚至处理另一项已经就绪的 Task,却只能守着一个尚未结束的进程。

上一篇的 Task System 解决了“工作如何跨会话保存”,但具体工具调用仍然是同步的。只要一个命令足够慢,整个 Agent 就会被它拖住。

这一章要把“等待命令”变成“等待一个有身份的 Job”:

启动命令
-> 立即得到 job_id
-> Agent 继续工作
-> 后台进程独立运行
-> 完成事件唤醒 Runtime
-> Agent 回来处理结果

Job 不会让命令跑得更快,却能让等待不再堵住整条路。

Task 描述目标,Job 记录一次执行

Task 和 Job 经常都被叫作“任务”,但它们回答的是两个问题。

假设 Task 是:

修复搜索接口,并完成验证。

它可能先后启动这些 Job:

pytest tests/test_search.py
npm run build
python scripts/check_api_contract.py

Task 关心最终目标、依赖、checkpoint 和交付结果;Job 只关心某一次执行何时开始、进程是否结束、退出码是什么、日志写在哪里。

因此:

一个 Task 可以拥有多个 Job
一个 Job 通常服务于一个 Task
Job succeeded 不等于 Task completed

测试通过,只能证明验证步骤成功。Agent 还要确认代码、文档和其它验收条件,才能完成 Task。

不是所有命令都值得放到后台

pwdgit status 和短文件查询的结果马上就要使用,后台化只会多出一次查询。

真正适合 Job 的通常是:

  • 安装依赖;
  • 完整构建和大型测试;
  • 批量数据处理;
  • 开发服务器、监听器等常驻进程;
  • 可以与其它工作并行的远程操作。

需要交互输入的命令不适合最小 Job Runner,因为它没有 PTY;修改系统或需要高权限的命令,则必须先走原有审批。

可以在 bash 工具中加入显式参数:

BASH_TOOL = {
    "type": "function",
    "function": {
        "name": "bash",
        "description": "Run a command in foreground or background.",
        "parameters": {
            "type": "object",
            "properties": {
                "command": {"type": "string"},
                "background": {"type": "boolean"},
                "timeout_seconds": {
                    "type": "integer",
                    "minimum": 1,
                    "maximum": 3600,
                },
            },
            "required": ["command"],
        },
    },
}

模型可以请求 background=true,但它不能借此绕过权限、沙箱、并发数和超时上限。后台只是调度方式,不是新的安全等级。

一条 Job 怎样走完整个生命周期

后台执行至少包含六步:

  1. 模型请求启动耗时命令;
  2. Harness 完成参数与权限校验;
  3. Job Manager 启动进程并返回 job_id
  4. Agent Loop 继续推进其它工作;
  5. Job Manager 轮询进程并持久化结果;
  6. Runtime Event Loop 收到事件,再次唤醒 Agent。
后台 Job 从启动到通知 Agent

Agent Loop 与后台进程各自向前走。Job 结束后先产生运行时事件,完整日志仍然按需读取。

这张图里有两个循环:

Agent Loop
  负责模型推理、工具调用和工具结果

Runtime Event Loop
  负责用户输入、Job 事件和其它外部事件

如果只有 Agent Loop,模型输出最终回复后,循环就结束了。后台进程即使完成,也找不到可以接收结果的模型回合。

为“不知道结果”保留一个状态

最小 Job 生命周期可以这样定义:

from typing import Literal


JobStatus = Literal[
    "running",
    "succeeded",
    "failed",
    "cancelled",
    "timed_out",
    "lost",
]
后台 Job 生命周期

退出码决定成功或失败,取消和超时是独立终态;重启后无法确认真实结果的进程进入 lost。

lost 很重要。

当前进程可以通过 Popen 句柄读取退出码。Harness 重启后,这个句柄消失了。磁盘上的 Job 仍然写着 running,但新进程无法可靠确认它是否已经成功、失败,甚至是否仍在执行。

此时不能直接写成 failed,更不能自动重跑。原命令可能已经完成部署、创建记录或修改文件,只是结果没有回来。

lost 表达的不是“执行失败”,而是“执行结果未知,需要对账”。

Job 记录只保存索引,不吞下整份日志

Job 元数据可以这样定义:

from dataclasses import asdict, dataclass


@dataclass
class BackgroundJob:
    id: str
    command_summary: str
    cwd: str
    status: JobStatus
    stdout_path: str
    stderr_path: str
    timeout_seconds: int

    task_id: str | None = None
    task_agent_id: str | None = None
    task_claim_token: str | None = None

    pid: int | None = None
    returncode: int | None = None
    created_at: str = ""
    started_at: str | None = None
    finished_at: str | None = None
    error: str | None = None
    schema_version: int = 1

    def to_dict(self) -> dict:
        return asdict(self)

目录中保存三类文件:

.agent/jobs/
├── job-a81d3f.json
├── job-a81d3f.stdout.log
└── job-a81d3f.stderr.log

JSON 只保存状态和日志路径。stdout、stderr 由操作系统直接写文件,避免一份大型构建日志撑坏 Job JSON,也避免它被整段塞进模型上下文。

command_summary 必须先脱敏。原始命令可能包含 Token、私有地址或密码,不能未经处理写进磁盘和审计日志。

启动进程后,立刻把控制权还给 Agent

下面的 Job Manager 复用上一章的 atomic_write_json()

import os
import subprocess
import time
import uuid
from pathlib import Path
from threading import RLock


class BackgroundJobManager:
    def __init__(self, root: Path):
        self.root = root.resolve()
        self.root.mkdir(parents=True, exist_ok=True)
        self.jobs: dict[str, BackgroundJob] = {}
        self.processes: dict[str, subprocess.Popen[bytes]] = {}
        self.deadlines: dict[str, float] = {}
        self.lock = RLock()

    def _save(self, job: BackgroundJob) -> None:
        atomic_write_json(
            self.root / f"{job.id}.json",
            job.to_dict(),
        )

    def start(
        self,
        command: str,
        command_summary: str,
        cwd: Path,
        timeout_seconds: int = 300,
        task_id: str | None = None,
        task_agent_id: str | None = None,
        task_claim_token: str | None = None,
    ) -> BackgroundJob:
        timeout_seconds = max(1, min(timeout_seconds, 3600))
        job_id = f"job-{uuid.uuid4().hex[:12]}"
        stdout_path = self.root / f"{job_id}.stdout.log"
        stderr_path = self.root / f"{job_id}.stderr.log"

        with stdout_path.open("ab", buffering=0) as stdout_file:
            with stderr_path.open("ab", buffering=0) as stderr_file:
                process = subprocess.Popen(
                    ["/bin/zsh", "-lc", command],
                    cwd=cwd,
                    stdin=subprocess.DEVNULL,
                    stdout=stdout_file,
                    stderr=stderr_file,
                    start_new_session=True,
                )

        now = utc_now_text()
        job = BackgroundJob(
            id=job_id,
            command_summary=command_summary,
            cwd=str(cwd),
            status="running",
            stdout_path=str(stdout_path),
            stderr_path=str(stderr_path),
            timeout_seconds=timeout_seconds,
            task_id=task_id,
            task_agent_id=task_agent_id,
            task_claim_token=task_claim_token,
            pid=process.pid,
            created_at=now,
            started_at=now,
        )

        with self.lock:
            self.jobs[job.id] = job
            self.processes[job.id] = process
            self.deadlines[job.id] = (
                time.monotonic() + timeout_seconds
            )
            self._save(job)
        return job

这里有三个细节值得留意。

第一,stdin=DEVNULL 明确拒绝交互输入。需要终端交互的命令要用专门的 PTY 工具。

第二,start_new_session=True 为子进程建立独立进程组。取消构建时,Harness 才能终止整棵进程树,而不是只杀掉最外层 shell。

第三,Popen() 成功与 Job 落盘之间仍有一个很小的崩溃窗口。对严格恢复场景,可以先保存 starting 意图,再启动进程,最后写 PID 和 running。如果中途崩溃,恢复器仍把结果不确定的记录归入 lost

轮询进程,但不要把完整输出推给模型

先准备两个辅助函数:

import signal


def tail_text(path: str, max_bytes: int = 2000) -> str:
    try:
        with open(path, "rb") as file:
            file.seek(0, os.SEEK_END)
            size = file.tell()
            file.seek(max(0, size - max_bytes))
            return file.read().decode(
                "utf-8",
                errors="replace",
            ).strip()
    except OSError:
        return ""


def stop_process_group(
    process: subprocess.Popen[bytes],
) -> None:
    try:
        os.killpg(process.pid, signal.SIGTERM)
        process.wait(timeout=3)
    except subprocess.TimeoutExpired:
        os.killpg(process.pid, signal.SIGKILL)
        process.wait()
    except ProcessLookupError:
        pass

Job Monitor 周期性调用 poll()

def poll_jobs(
    manager: BackgroundJobManager,
) -> list[dict]:
    events: list[dict] = []

    with manager.lock:
        for job_id, process in list(
            manager.processes.items()
        ):
            job = manager.jobs[job_id]

            if (
                process.poll() is None
                and time.monotonic()
                >= manager.deadlines[job_id]
            ):
                stop_process_group(process)
                job.status = "timed_out"

            returncode = process.poll()
            if returncode is None:
                continue

            if job.status == "running":
                job.status = (
                    "succeeded"
                    if returncode == 0
                    else "failed"
                )

            job.returncode = returncode
            job.finished_at = utc_now_text()
            if job.status != "succeeded":
                job.error = redact_sensitive_text(
                    tail_text(job.stderr_path)
                )
            manager._save(job)

            summary_path = (
                job.stdout_path
                if job.status == "succeeded"
                else job.stderr_path
            )
            events.append(
                {
                    "event_id": f"job-finished:{job.id}",
                    "kind": "background_job_finished",
                    "job_id": job.id,
                    "task_id": job.task_id,
                    "status": job.status,
                    "returncode": job.returncode,
                    "summary": redact_sensitive_text(
                        tail_text(summary_path)
                    ),
                }
            )

            del manager.processes[job_id]
            del manager.deadlines[job_id]

    return events

通知只带有限长度的尾部摘要。需要更多信息时,模型再调用 job_output 分页读取。

event_id 要稳定,因为事件队列通常提供至少一次投递。同一个完成事件偶尔重复,比它悄无声息地丢失更容易恢复。

事件不是用户说的话

Job 完成来自 Runtime,不是用户输入。

如果把它追加成普通 user 消息,模型会误以为“测试完成”是用户的陈述,后续摘要和审计也会失真。

更合适的方式是把事件作为动态 Prompt Section:

Runtime events:
- job=job-a81d3f
  status=succeeded
  task=build-api
  summary=42 passed in 18.4s

模型调用成功后,再确认事件已消费;如果模型请求失败,事件留在队列中,下一次继续投递。

Runtime Event Loop 的核心可以保持很小:

from queue import Empty, Queue
from threading import Thread


def monitor_background_jobs(
    state: HarnessState,
    event_queue: Queue[dict],
) -> None:
    next_lease_heartbeat = 0.0
    while not state.shutdown_requested:
        for event in poll_jobs(state.job_manager):
            event_queue.put(event)

        now = time.monotonic()
        if now >= next_lease_heartbeat:
            heartbeat_running_job_tasks(state, event_queue)
            next_lease_heartbeat = now + 60
        state.shutdown_event.wait(timeout=0.5)


def run_runtime(state: HarnessState) -> None:
    queue = state.runtime_event_queue
    monitor = Thread(
        target=monitor_background_jobs,
        args=(state, queue),
        daemon=False,
    )
    monitor.start()

    while not state.shutdown_requested:
        try:
            event = queue.get(timeout=0.5)
        except Empty:
            continue

        state.runtime_events.append(event)
        try:
            run_agent_turn(state)
        except Exception:
            queue.put(event)
            state.shutdown_event.wait(timeout=1)
        finally:
            state.runtime_events.remove(event)

监控线程只检查 Job、写入持久化状态并投递事件,不直接修改 Agent 对话。主循环一次仍只运行一个模型回合,避免两个事件并发改写同一份 session。

需要跨 Harness 重启保证事件不丢失时,内存 Queue 还不够,应把待投递事件写入持久化 Outbox。

Job 等待期间,Task 租约不能悄悄过期

后台测试仍然属于当前 Task。Job 启动后,Task 应保持 in_progress,而不是立刻 release。

Job Manager 记录 task_idtask_agent_idtask_claim_token,Monitor 再定期调用上一章的续租逻辑:

def heartbeat_running_job_tasks(
    state: HarnessState,
    event_queue: Queue[dict],
) -> None:
    with state.job_manager.lock:
        running = [
            job
            for job in state.job_manager.jobs.values()
            if job.status == "running"
        ]

    for job in running:
        if not (
            job.task_id
            and job.task_agent_id
            and job.task_claim_token
        ):
            continue
        try:
            renew_task_lease(
                store=state.task_store,
                task_id=job.task_id,
                agent_id=job.task_agent_id,
                claim_token=job.task_claim_token,
                renew_seconds=300,
            )
        except TaskStoreError as exc:
            event_queue.put(
                {
                    "event_id": (
                        f"job-lease-conflict:{job.id}"
                    ),
                    "kind": "job_task_lease_conflict",
                    "job_id": job.id,
                    "task_id": job.task_id,
                    "summary": str(exc),
                }
            )

如果所有权校验失败,旧 Job 可以继续记录自己的进程结果,却不能再更新已经被别人接管的 Task。当前所有者决定复用结果、忽略结果还是重新执行。

还要记住三条规则:

  1. Job 成功不自动完成 Task;
  2. Job 运行时不要释放 Task;
  3. lost Job 不自动重跑,先检查日志、artifact 和外部状态。

暴露给模型的 Job 工具

除了 bash(background=true),模型只需要少量工具:

job_status   查看状态、退出码和关联 Task
job_output   有限读取 stdout 或 stderr
job_list     查看当前运行和最近完成的 Job
job_cancel   请求终止当前工作空间里的进程组

job_cancel 只改变 Job,不改变 Task。Monitor 在下一次轮询时补齐退出码和结束时间。

模型也不能只凭 job_id 取消任意进程。Harness 要校验 Job 是否属于当前用户、Task 和工作区;返回 Job 信息时,还应移除内部的 task_claim_token

权限检查发生在 Job 启动之前。命令如果不允许执行,background=true 不会让它突然合法。

后台执行最容易失控的地方

Agent 不再直接感受等待成本之后,Harness 更需要资源边界:

  • 每个 Agent 和全局的并发 Job 数;
  • 单个 Job 最大运行时间;
  • stdout 与 stderr 的磁盘配额;
  • 可用工作目录、环境变量和端口;
  • 取消与强制终止权限;
  • Harness 退出时如何处理仍在运行的进程;
  • 日志保留时间与敏感信息脱敏。

尤其不要犯三类错误:

只杀 shell,不杀进程组
把整份日志塞进上下文
重启后自动重做结果未知的写操作

Job 的价值不是让执行“更随意”,而是让长时间执行仍然拥有状态、所有者和可观测结果。

小结

这一章把一次同步命令升级成了可管理的后台 Job:

  • job_id 让 Agent 可以先离开,再回来查询;
  • stdout、stderr 单独落盘,通知只携带摘要;
  • 进程组、超时和取消让生命周期可以收束;
  • lost 明确表达重启后的未知结果;
  • Runtime Event Loop 把完成事件送回正确的 Agent;
  • Task 租约在等待期间继续续期;
  • 权限、沙箱和资源预算仍然位于执行之前。

现在,Agent 不必再守着测试进程发呆。

但 Job 仍然需要某个人在某一刻主动启动。如果我们希望每天早上八点检查服务器,每晚十一点运行自动化测试,就还缺少一只不会睡过头、重启后也记得约定的钟。

下一章加入 Scheduler:让时间规则可靠地产生一次 ScheduleRun,再把它交给 Task System 和后台 Job。