Agent 工程体系 · 第 73/98 篇。内容以 2026 年 9 月可验证的公开规范和稳定接口为基线;框架版本敏感能力会明确标注,不把实验行为写成通用保证。

Agent 队列与背压:会话任务、生成 Worker、发送器和过期丢弃

在一个支持流式输出的 Agent 服务中,“收到用户消息后启动一个异步任务”很快就会演变成一条生产链路:

用户请求
  │
  ▼
会话任务 Session Task
  │
  ├── 调度 Agent / 工具 / 子 Agent
  │
  ▼
生成 Worker Generation Worker
  │
  ├── 产生文本增量、工具事件、状态事件
  │
  ▼
发送器 Sender
  │
  ▼
WebSocket / SSE / IM / 回调接口

这条链路中至少存在三个不同速度:

  1. 用户提交任务的速度;
  2. Agent 生成事件的速度;
  3. 网络发送事件的速度。

如果生成速度长期高于发送速度,事件就会在内存中累积;如果发送器阻塞了生成 Worker,模型调用和工具执行又会被网络抖动拖慢;如果用户已经断开连接,继续生成完整答案只会消耗模型、工具和计算资源。

因此,队列并不是简单的“把结果放进去,再从另一边取出来”。它承担了隔离速度、传播压力、限制资源、判断新鲜度和丢弃无价值工作等职责。


一、先区分四种对象:会话、任务、Worker 和事件

1. 会话不是任务

**会话(session)**表示一个用户或客户端上下文,通常包含:

  • session_id
  • 对话历史或其引用
  • 当前连接
  • 用户级并发限制
  • 会话取消信号
  • 最近一次用户请求的版本号

**任务(task)**表示会话中的一次可执行工作,例如:

  • 回答一次用户消息;
  • 执行一次搜索;
  • 生成一份报告;
  • 处理一个用户确认后的工具调用。

同一个会话可以先后提交多个任务,但通常不能让多个任务无约束地同时修改同一份会话状态。

例如,用户快速发送:

用户:帮我查一下杭州天气
用户:顺便比较一下宁波
用户:算了,只看杭州

如果三个任务都并行运行,第三条消息可能已经使前两个任务失去价值。即使前两个任务最后生成成功,也不能任意覆盖第三个任务的结果。

因此,常见的会话级不变量是:

同一个会话内,最多允许一个“拥有会话写权限”的生成任务处于运行状态。

这不意味着系统只能全局串行。不同会话仍然可以并行;同一个任务内部,独立的检索子任务也可以并行。

任务需要拥有稳定身份

任务不能只用一个字符串描述。至少应包含:

task_id       任务唯一 ID
session_id    所属会话
generation    会话内版本号
created_at    创建时间
deadline      绝对截止时间
priority      调度优先级
cancel_token  取消信号
payload       用户输入或任务引用

其中 generation 很重要。它用于解决“旧任务晚到”的问题。

设会话 s-1 先后产生三个任务:

T1: generation = 1
T2: generation = 2
T3: generation = 3

如果 T1 因为模型或工具调用较慢,在 T3 之后才产生输出,那么发送器必须拒绝 T1 的输出。判断依据不是“这个输出是否生成成功”,而是:

event.generation == session.latest_generation

这是一种逻辑过期,不依赖系统时钟。


二、队列到底隔离了什么

队列隔离的不是“代码模块”,而是不同阶段的处理速度。

假设:

  • 生成 Worker 平均每秒产生 λ_g = 20 个事件;
  • 发送器平均每秒只能发送 μ_s = 10 个事件。

事件队列长度的期望增长速度约为:

dQdt=λgμs=10\frac{dQ}{dt} = \lambda_g - \mu_s = 10

其中:

  • QQ 是队列中等待发送的事件数;
  • λg\lambda_g 是事件到达率;
  • μs\mu_s 是事件服务率。

如果这个差值持续为正,队列最终必然耗尽内存。增加队列容量只能推迟故障:

容量 100:约 10 秒后爆满
容量 1000:约 100 秒后爆满

它没有改变系统不稳定的事实。

对于长期运行的稳定队列,基本条件是:

λ<μ\lambda < \mu

也就是平均进入速率必须低于平均处理速率。瞬时突发可以由有限队列吸收,但持续过载必须通过以下一种或多种方式处理:

  • 限制生产者;
  • 拒绝新任务;
  • 降低事件粒度;
  • 丢弃过期事件;
  • 扩容消费者;
  • 降低模型生成速度;
  • 将部分任务转入离线处理。

Little 定律与队列延迟

在稳定状态下,Little 定律给出:

L=λWL = \lambda W

其中:

  • LL 是系统中平均任务数或事件数;
  • λ\lambda 是平均到达率;
  • WW 是平均等待时间加处理时间。

例如发送器每秒稳定处理 10 个事件,平均队列中有 50 个事件,则仅平均等待时间就约为:

W=Lλ=5010=5 秒W = \frac{L}{\lambda} = \frac{50}{10} = 5\text{ 秒}

这解释了一个常见现象:队列没有爆满,但用户已经感觉“流式输出卡住了”。队列长度的增长通常先表现为延迟增长,最后才表现为内存和超时故障。


三、背压不是“队列满了再报错”

**背压(backpressure)**是下游处理不过来时,将压力传回上游,使上游降低生产速度、等待、拒绝或丢弃工作。

对 Agent 系统而言,背压至少有四个层级。

1. 入口背压:限制会话任务

当一个会话已有任务运行时,可以采用以下策略:

串行等待

新任务进入会话队列,等待旧任务完成。

适合:

  • 每个任务都重要;
  • 会话操作具有顺序依赖;
  • 用户明确期望按顺序执行。

缺点是用户连续修改请求时,旧任务会阻塞新任务。

合并

把多个相邻任务合并成一次输入,例如:

“写一封邮件”
“语气正式一点”
“收件人改成张三”

合并为一次最新请求。

合并要求任务语义允许重写,否则可能改变用户意图。

取消旧任务,只保留最新任务

这是交互式 Agent 常用的策略。它适合“用户正在编辑问题”或“搜索结果只需要最新版本”的场景。

但取消必须分层处理:

  • 取消等待中的任务;
  • 取消模型流;
  • 取消工具调用;
  • 释放并发槽位;
  • 阻止迟到事件发送;
  • 清理临时文件和连接。

只把 cancelled = true 写入内存变量,并不能停止已经发出的 HTTP 请求或子进程。

2. 全局背压:限制生成 Worker

生成 Worker 通常消耗昂贵资源:

  • 模型并发额度;
  • 工具连接;
  • CPU;
  • GPU;
  • 上下文内存;
  • 任务租约。

因此,不能为每个进入请求无限创建 Worker。更安全的模型是:

入口任务队列:限制任务数量
生成 Worker 池:限制同时生成的任务数
发送队列:限制待发送事件数
发送 Worker:限制下游连接压力

生成 Worker 的槽位代表“可以执行 Agent 逻辑”,而不是简单的线程数。一个 Agent 任务可能在等待模型响应时仍占用连接、上下文和外部资源,所以并发额度应按实际瓶颈设置。

3. 事件背压:生成速度受发送队列影响

最直接的设计是有限队列:

await output_queue.put(event)

当队列满时,put 阻塞,生成 Worker 暂停产生新事件。这个机制成立的前提是:

  • 生成 Worker 可以安全暂停;
  • 暂停不会持有无法释放的锁;
  • 模型流或工具流可以被限速或取消;
  • 任务有明确的截止时间。

否则,生成 Worker 可能停在队列写入处,但模型连接仍保持打开,最终形成资源泄漏。

4. 发送背压:连接关闭必须向上游传播

如果客户端断开,发送器通常会遇到:

BrokenPipeError
ConnectionResetError
WebSocketDisconnect

此时不能只退出发送循环。正确路径应是:

发送失败
  │
  ├── 标记连接不可用
  ├── 取消当前会话任务
  ├── 关闭模型流
  ├── 取消工具调用
  ├── 清空或废弃输出队列
  └── 释放 Worker 槽位

如果只停止发送器,生成 Worker 仍会继续工作,系统就会在“无人接收”的任务上持续消耗资源。


四、为什么要拆成会话任务、生成 Worker 和发送器

把所有逻辑写成一个协程看起来简单:

async def handle_request(request):
    async for chunk in model_stream(request):
        await websocket.send_text(chunk)

但这段代码把三个责任绑定在了一起:

  1. Agent 生成;
  2. 连接写入;
  3. 连接生命周期管理。

当网络发送变慢时,model_stream 的消费速度也变慢;当连接断开时,生成逻辑和发送逻辑的异常边界混在一起;当需要丢弃旧输出时,也没有独立的过滤点。

更清晰的结构是:

flowchart LR
    A[用户请求] --> B[会话调度器]
    B --> C[会话任务队列]
    C --> D[生成 Worker]
    D --> E[有限输出队列]
    E --> F[发送器]
    F --> G[客户端]

    H[取消信号] --> B
    H --> D
    H --> F

    I[截止时间/版本号] --> D
    I --> F

    F -->|连接关闭| H

会话任务

会话任务负责:

  • 取得会话锁;
  • 检查任务是否已经过期;
  • 载入上下文;
  • 调用 Agent 运行循环;
  • 将生成结果写入输出队列;
  • 在退出时释放会话锁和资源。

它不应直接控制具体的网络发送。

生成 Worker

生成 Worker 负责:

  • 执行模型调用;
  • 处理工具调用和子 Agent;
  • 产生标准化事件;
  • 响应取消和截止时间;
  • 在输出队列满时接受背压;
  • 结束时发送终止事件或错误事件。

生成事件最好是结构化的:

{
  "task_id": "t-123",
  "session_id": "s-1",
  "generation": 7,
  "seq": 42,
  "kind": "text_delta",
  "payload": "背压",
  "expires_at": 1788249600.0
}

seq 用于保证同一任务内的事件顺序;generation 用于拒绝旧任务;expires_at 用于时间过期检查。

发送器

发送器负责:

  • 从输出队列取事件;
  • 检查连接是否仍然有效;
  • 检查事件是否过期;
  • 检查任务是否仍是会话最新版本;
  • seq 顺序发送;
  • 处理网络异常;
  • 触发上游取消。

发送器不应该重新执行 Agent,也不应该修改会话业务状态。它的职责是把“仍然有价值的事件”交付给客户端。


五、过期丢弃:为什么成功生成的结果也可能不应发送

**过期丢弃(stale-drop / expiry discard)**是指事件虽然已经生成,但由于用户意图、连接或截止时间已经失效,系统主动不再发送。

过期至少有三种含义。

1. 时间过期

任务有绝对截止时间:

nowdeadlinenow \ge deadline

则任务不再接受新的计算结果。

必须使用单调时钟计算持续时间,不能用墙上时钟判断超时。系统时间可能因为校时回拨或跳跃。

2. 版本过期

任务版本小于会话当前版本:

event.generation<session.latest_generationevent.generation < session.latest\_generation

即使事件刚刚生成,也应丢弃。

3. 连接过期

客户端连接已关闭,或者发送上下文已经失效。此时输出没有消费者,应取消生成,而不是继续积累。

过期检查应该放在哪里

只在任务开始时检查一次是不够的:

任务开始时有效
→ 等待模型
→ 调用工具
→ 排队等待发送
→ 发送时已过期

因此至少要在四个位置检查:

  1. 从任务队列取出时;
  2. 调用模型或工具前;
  3. 写入输出队列前;
  4. 从输出队列取出、准备发送前。

最后一个检查尤其重要,因为队列等待本身就是延迟来源。

丢弃增量事件与丢弃最终结果

两者不能混为一谈。

  • 对流式增量文本,丢弃单个旧事件通常是安全的,因为它本来只是中间显示;
  • 对最终结果,丢弃意味着客户端可能永远等不到 done 事件,因此应明确发送 cancelledexpired 或关闭流;
  • 对工具执行结果,不能只因为输出过期就假设副作用不存在。发送过期和撤销工具副作用是两个不同问题。

例如:

支付工具已成功扣款
但用户连接在响应前断开

这不是“丢弃输出”就能解决的。系统必须依赖工具幂等键、事务状态和补偿流程,避免因为客户端重试而重复扣款。


六、一个可运行的最小实现

下面的示例只使用 Python 标准库,模拟:

  • 会话级最新任务策略;
  • 有限输出队列;
  • 生成 Worker;
  • 发送器;
  • 任务超时;
  • 旧版本丢弃;
  • 客户端断开后的取消传播。
from __future__ import annotations

import asyncio
import time
from dataclasses import dataclass


@dataclass
class Task:
    task_id: str
    session_id: str
    generation: int
    text: str
    deadline: float
    cancel: asyncio.Event


@dataclass
class Event:
    task_id: str
    session_id: str
    generation: int
    seq: int
    kind: str
    payload: str
    created_at: float


class Session:
    def __init__(self, session_id: str, output_limit: int = 3):
        self.session_id = session_id
        self.latest_generation = 0
        self.current_cancel: asyncio.Event | None = None
        self.output: asyncio.Queue[Event] = asyncio.Queue(maxsize=output_limit)
        self.lock = asyncio.Lock()


async def generate(task: Task, session: Session) -> None:
    """模拟 Agent 生成事件。真实系统中这里可调用模型、工具和子 Agent。"""
    words = task.text.split()

    for seq, word in enumerate(words):
        if task.cancel.is_set():
            return

        if time.monotonic() >= task.deadline:
            task.cancel.set()
            return

        event = Event(
            task_id=task.task_id,
            session_id=task.session_id,
            generation=task.generation,
            seq=seq,
            kind="text_delta",
            payload=word + " ",
            created_at=time.monotonic(),
        )

        # 有限队列形成背压:发送器太慢时,生成 Worker 会在这里等待。
        await session.output.put(event)
        await asyncio.sleep(0.05)


async def sender(
    session: Session,
    client_name: str,
    disconnect_after: float | None = None,
) -> None:
    started = time.monotonic()

    while True:
        if disconnect_after is not None:
            if time.monotonic() - started >= disconnect_after:
                raise ConnectionError(f"{client_name}: client disconnected")

        try:
            event = await asyncio.wait_for(session.output.get(), timeout=0.2)
        except asyncio.TimeoutError:
            # 队列暂时为空;真实发送器可以在这里检查连接心跳。
            continue

        try:
            # 发送前再次做逻辑过期检查。
            is_latest = event.generation == session.latest_generation
            if not is_latest:
                print(
                    f"[sender] drop stale event: "
                    f"task={event.task_id}, generation={event.generation}"
                )
                continue

            print(f"[{client_name}] {event.payload!r}")

            # 模拟慢网络。
            await asyncio.sleep(0.15)
        finally:
            session.output.task_done()


async def run_task(session: Session, task: Task) -> None:
    async with session.lock:
        if task.generation != session.latest_generation:
            print(f"[worker] skip stale task: {task.task_id}")
            return

        try:
            await generate(task, session)
        except asyncio.CancelledError:
            task.cancel.set()
            raise
        finally:
            print(f"[worker] finished: {task.task_id}")


async def submit(
    session: Session,
    task_id: str,
    text: str,
    timeout: float,
) -> asyncio.Task:
    # 新任务发布后,旧任务立即失去“发送资格”。
    session.latest_generation += 1
    generation = session.latest_generation

    if session.current_cancel is not None:
        session.current_cancel.set()

    cancel = asyncio.Event()
    session.current_cancel = cancel

    task = Task(
        task_id=task_id,
        session_id=session.session_id,
        generation=generation,
        text=text,
        deadline=time.monotonic() + timeout,
        cancel=cancel,
    )

    return asyncio.create_task(run_task(session, task))


async def main() -> None:
    session = Session("s-1", output_limit=2)

    send_task = asyncio.create_task(
        sender(session, "client-A", disconnect_after=2.0)
    )

    first = await submit(
        session,
        task_id="t-1",
        text="旧任务 产生 很多 不再需要的 输出",
        timeout=5,
    )

    await asyncio.sleep(0.12)

    second = await submit(
        session,
        task_id="t-2",
        text="新任务 保留 这份 输出",
        timeout=5,
    )

    await asyncio.gather(first, second)

    # 让发送器继续消费一小段时间。
    await asyncio.sleep(1)

    send_task.cancel()
    try:
        await send_task
    except asyncio.CancelledError:
        pass


if __name__ == "__main__":
    asyncio.run(main())

这段代码中各个条件为什么成立

output_limit=2 表示发送队列最多缓存两个事件。发送器每 150 毫秒处理一个事件,而生成器每 50 毫秒产生一个事件,生产速度约为消费速度的三倍,因此生成器会在 await session.output.put(event) 处等待。

session.latest_generation 解决了旧任务晚到的问题。提交 t-2 时,版本从 1 变成 2;即使 t-1 已经把事件写进队列,发送器取出时也会发现:

event.generation != session.latest_generation

于是丢弃事件。

session.current_cancel.set() 只通知旧 Worker 尽快停止。它不是强制终止。如果旧 Worker 正卡在不可取消的同步调用中,仍然需要:

  • 为外部 HTTP 请求设置超时;
  • 为子进程发送终止信号;
  • 为模型流关闭连接;
  • 为工具操作设计幂等或补偿机制。

一个容易忽略的 bug

示例中的 sender 在客户端断开时抛出异常,但 main 没有把这个异常继续传播给生成任务。生产实现中,发送器异常必须触发会话级取消:

async def supervise(session, sender_task, worker_tasks):
    try:
        await sender_task
    except Exception:
        if session.current_cancel is not None:
            session.current_cancel.set()

        for task in worker_tasks:
            task.cancel()

更完整的实现通常会用任务组或 supervisor,保证任一关键阶段失败后,其余阶段得到明确处理,而不是留下孤儿任务。


七、输出队列不是所有事件都应该同等对待

如果所有事件都放进一个 FIFO 队列,慢消费者会让低价值事件阻塞高价值事件。

例如:

token_delta
token_delta
token_delta
tool_started
approval_required
done

当大量文本增量占满队列时,approval_required 可能无法及时发送。这会导致用户看不到需要确认的操作。

因此,事件应至少按语义分类:

事件类型 是否可丢弃 常见处理
文本增量 通常可合并或丢弃 按时间窗口合并
思考状态 可丢弃 只保留最新状态
工具开始/结束 通常不可静默丢弃 保留状态一致性
人工确认请求 不可丢弃 高优先级发送
最终结果 不可静默丢弃 发送完成、取消或失败
心跳 可丢弃 只保证连接检测

文本增量合并

如果模型每次只产生一个字符,队列元素数量会非常大。可以把多个增量合并成一个事件:

"背"
"压"
"可"
"以"
"限"
"制"

合并为:

"背压可以限制"

合并会牺牲一点实时性,但能降低队列压力和网络系统调用次数。合并不能跨越任务版本或消息边界,否则可能把旧任务文本拼入新任务。

优先级队列的边界

优先级队列可以让控制事件优先发送,但它会破坏普通 FIFO 顺序。可采用两级队列:

control_queue:取消、确认、失败、完成
data_queue:文本增量、普通状态

发送器每轮优先处理控制队列,再处理数据队列。

但“高优先级”不等于“可以绕过任务版本检查”。一个旧任务的失败事件也可能已经没有发送价值,过期检查仍然必须执行。


八、重试、取消和过期必须区分责任层

Agent 任务经常同时包含三种操作:

模型调用 → 工具调用 → 输出发送

这三层不能使用同一种重试策略。

模型调用

模型生成通常可以重试,但要考虑:

  • 是否已经产生部分输出;
  • 重试后是否可能生成不同文本;
  • 是否会重复计算 token;
  • 是否应复用同一个任务版本。

如果已经向客户端发送了一部分文本,重试可能造成重复内容。可采用:

  • 给事件增加 seq
  • 客户端按 seq 去重;
  • 重试时发出 generation_reset
  • 或直接结束本次流,由上层重新发起任务。

工具调用

工具调用可能有副作用。以下操作不能默认安全重试:

扣款
发邮件
创建工单
修改数据库
提交代码

工具接口应支持幂等键:

idempotency_key = task_id + ":" + tool_call_id

服务端通过幂等键保证同一业务操作最多生效一次,或者明确返回已执行结果。

发送

网络发送失败时,不能简单重试所有事件。因为客户端可能已经收到但服务端没有收到确认,重试会导致重复显示。

可靠发送需要协议层支持:

  • 事件序号;
  • 客户端确认;
  • 断点续传;
  • 去重;
  • 或接受“至少一次投递”。

如果系统只要求实时展示,通常选择“尽力而为的流式输出”,而把最终结果持久化到任务记录中。客户端重连后读取最终状态,而不是依赖重新发送全部 token。


九、任务状态机要能解释每一种故障

一个可诊断的任务状态机可以是:

stateDiagram-v2
    [*] --> QUEUED
    QUEUED --> RUNNING: worker 获取任务
    QUEUED --> EXPIRED: 超过 deadline
    QUEUED --> CANCELLED: 新任务替代/用户取消

    RUNNING --> EMITTING: 产生事件
    RUNNING --> CANCELLED: 取消信号
    RUNNING --> FAILED: 模型或工具失败
    RUNNING --> EXPIRED: 达到截止时间
    RUNNING --> SUCCEEDED: 无需发送或结果已保存

    EMITTING --> RUNNING: 继续生成
    EMITTING --> SUCCEEDED: 发送 done
    EMITTING --> CANCELLED: 连接关闭
    EMITTING --> EXPIRED: 事件已失效
    EMITTING --> FAILED: 发送失败

状态变化必须满足几个约束:

  1. 终态不能回到运行态;
  2. 任务版本失效后,不能再产生对外可见的新事件;
  3. CANCELLED 不一定代表外部副作用被撤销;
  4. SUCCEEDED 只代表业务任务成功,不代表客户端一定收到;
  5. FAILED 应包含责任层,例如 model_errortool_errorsend_error

最后一点很关键。发送失败不应被误记成 Agent 生成失败,否则运维人员会错误地去排查模型。


十、生产环境的诊断指标

只看“队列长度”不够。至少应分别记录:

任务层

task_queued_total
task_running
task_wait_seconds
task_execution_seconds
task_expired_total
task_cancelled_total
task_replaced_total

事件层

event_produced_total
event_dropped_stale_total
event_dropped_expired_total
event_queue_size
event_queue_wait_seconds
event_batch_size

发送层

sender_active
send_success_total
send_failure_total
client_disconnect_total
send_latency

资源层

model_inflight
tool_inflight
session_active
worker_slots_used

诊断时要区分以下现象:

现象 更可能的原因
任务等待时间高,生成时间正常 入口任务队列或 Worker 槽位不足
生成速度正常,发送延迟持续升高 网络发送器过慢
队列长度高但发送吞吐低 下游连接、批量策略或客户端消费慢
旧事件丢弃数高 用户频繁改写请求,或 Worker 取消不及时
取消数高且模型费用仍上升 取消没有传播到模型或工具层
任务成功数高但用户收不到 发送与业务成功状态没有分离
队列为空但延迟高 可能阻塞在模型、工具、锁或连接建立阶段

OpenAI 的 Agents SDK 将 Agent run、运行循环、流式处理、handoff、guardrails、状态和可观测性分成不同能力面;这类划分说明 Agent 运行本身与外围队列、连接和存储并不是同一个责任边界。若应用需要完全自定义循环、工具路由、分支和状态,应该由应用自己掌握这些控制点;若使用 SDK 管理 Agent loop,则仍需在应用层管理部署、工具实现、状态存储和审批决定。(developers.openai.com)

Anthropic 对 Agent 的描述也强调:Agent 通常是在工具和环境反馈驱动下循环执行,并且应设置最大迭代次数等停止条件。这个停止条件不仅用于防止模型无限循环,也应成为队列系统的资源边界:一个任务不能因为一直有“下一步计划”就无限占用会话、Worker 和输出队列。(anthropic.com)


十一、几个典型反例

反例一:无限队列加无限任务

asyncio.create_task(handle_request(request))

每个请求都创建任务,任务内部再创建模型流。短时间内吞吐看起来很好,但过载时会同时扩大:

  • 内存占用;
  • 模型并发;
  • 工具连接;
  • 数据库连接;
  • 输出队列;
  • 超时任务数量。

这不是并发控制,而是把排队位置从显式队列转移到了运行时和外部服务。

反例二:只取消 Python 任务,不取消外部调用

worker_task.cancel()

如果外部 SDK、线程池或子进程没有响应取消,实际工作仍可能继续。最终表现为:

应用认为任务已取消
模型供应商仍在生成
工具仍在执行
资源计数无法归零

取消必须覆盖资源的真实拥有者。

反例三:发送器按“任务完成”而不是按“事件有效”发送

生成 Worker 完成并不意味着结果仍然有效。任务可能在排队期间已经被新版本替代,也可能已经超过截止时间。

正确判断顺序是:

连接有效?
任务未取消?
任务未过期?
任务仍是最新版本?
事件序号合法?

全部成立后,才发送。

反例四:丢弃旧任务,却保留旧任务的最终副作用

把旧任务标记为 stale 只能阻止输出,不能撤销它已经完成的工具副作用。对于有副作用的 Agent,必须把:

输出生命周期

与:

业务操作生命周期

分开建模。


十二、一个可落地的最小设计

对于交互式、多 Agent、支持流式输出的服务,可以先建立以下边界:

每个会话:
  一个当前任务
  一个取消信号
  一个最新 generation
  一个有限输出队列
  一个发送器

全局:
  一个有界任务池
  一个有限生成 Worker 池
  一个工具并发配额
  一个统一 supervisor

任务提交时:

1. 分配 task_id 和 generation
2. 判断会话是否已有任务
3. 按产品语义选择排队、合并或替代
4. 设置绝对 deadline
5. 放入有界任务池
6. 队列满时明确拒绝或降级

Worker 执行时:

1. 获取会话所有权
2. 检查任务版本和 deadline
3. 执行 Agent 循环
4. 在模型、工具和输出边界响应取消
5. 产生带 generation、seq 的结构化事件
6. 写入有限输出队列
7. finally 中释放所有资源

发送器执行时:

1. 取出事件
2. 检查连接
3. 检查取消和 deadline
4. 检查 generation
5. 按 seq 发送
6. 处理断连
7. 将断连传播给 Worker

这套结构的核心不是“使用某一种队列库”,而是明确三条不变量:

队列容量有限\text{队列容量有限}

任务拥有明确的截止时间和取消路径\text{任务拥有明确的截止时间和取消路径}

过期事件不能重新获得发送资格\text{过期事件不能重新获得发送资格}

Agent 运行本身可以由 SDK、Responses API 或自定义循环实现。OpenAI 文档区分了“由应用掌握循环的 Responses API”和“由 SDK 管理 Agent loop 的 Agents SDK”;无论选择哪一种,队列、背压、取消和发送一致性仍然是应用运行时的责任。(developers.openai.com)

最终,可靠的 Agent 队列系统并不是尽量把所有工作都完成,而是在资源有限、用户意图变化、网络不稳定和任务持续失败的条件下,仍然能回答四个问题:

当前任务是谁?
它是否仍然有价值?
谁在消耗资源?
什么时候必须停止?

能准确回答这四个问题,队列才不仅是缓存,而是真正的 Agent 运行时边界。


系列导航与关联阅读

官方资料

本文依据 Agent、模型、协议与框架官方资料重新梳理;正文、示例与生产清单由 WR BLOG 编写。