Agent Harness 工程:后台任务与异步通知
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。
不是所有命令都值得放到后台
pwd、git 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 怎样走完整个生命周期
后台执行至少包含六步:
- 模型请求启动耗时命令;
- Harness 完成参数与权限校验;
- Job Manager 启动进程并返回
job_id; - Agent Loop 继续推进其它工作;
- Job Manager 轮询进程并持久化结果;
- Runtime Event Loop 收到事件,再次唤醒 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",
]
退出码决定成功或失败,取消和超时是独立终态;重启后无法确认真实结果的进程进入 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_id、task_agent_id 和 task_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。当前所有者决定复用结果、忽略结果还是重新执行。
还要记住三条规则:
- Job 成功不自动完成 Task;
- Job 运行时不要释放 Task;
lostJob 不自动重跑,先检查日志、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。