从零搭建 Agent Harness 系列(十八)事件背压:有界队列、尊重 Context 与可关闭的 EventSink
系列十七修好了终端审批泄漏。阶段七还剩另一件容易在生产里爆内存的事:事件比网络快时,往哪里堆?以及堆不住时,失败能不能变成 Run 的终态。
工具并发上限限制的是「同时跑几个工具」。事件背压限制的是「事件管道别无限灌」。两者管的不是同一类资源。
本文对应 go-tiny-claw 两次提交:
1 | bb746e1 feat: 为 Channel EventSink 增加有界队列背压 |
落地之后,Channel 侧的事件出口是:
1 | 有界 queue(默认 64) |
一、原来的 Publish 直接写连接
改造前,路径是:
1 | Engine / Reporter |
Publish 会先看一眼 ctx,然后同步 Write:
1 | func (s *JSONEventSink) Publish(ctx context.Context, event reporter.Event) error { |
慢客户端会把 Write 堵住。更糟的是:堵住之后再按 Ctrl-C,这次 Write 不一定听 ctx,Engine 可能一直卡在某次 OnTextDelta 上。
当时 JSONReporter 还是:
1 | _ = r.sink.Publish(ctx, event) |
错误被丢掉。背压这一层先让 Publish 自己能阻塞、能取消、能关闭;随后再把错误从 Reporter 传回 Engine。
二、背压是什么
把中间换成有界队列:
1 | Publish ──► queue(cap=64) ──► loop goroutine ──► Write(conn) |
| 情况 | 行为 |
|---|---|
| 队列有空位 | Publish 立刻返回 |
| 队列满、任务还在 | Publish 堵住,Engine 变慢 |
| 队列满、用户取消 | ctx.Done() 返回,不永久卡住 |
| Session 关闭 | 不再入队,等 Write 返回后 loop 退出 |
「满了让上游慢下来」就是背压。没有界的 channel 或无限 goroutine,只是把压力变成内存。
三、为什么不能 for range queue
第一版 writer 是:
1 | for item := range s.queue { |
range 一个 channel 会一直等到 channel 被 close。Session 关掉如果不 close(queue),这条 goroutine 会永远挂着。
但 close(queue) 和 queue <- item 并发会发生 panic。Engine 的 Destroy 和迟到的 Publish 之间做不到绝对无缝。
所以关机信号不能靠 queue,要另开一个 done:
1 | done:广播「Sink 已关闭」,叫醒所有 select |
loop 改成:
1 | func (s *JSONEventSink) loop() { |
Close 只 close(s.done)。select 只能叫醒还在等 done / queue 的 loop。一旦已经进了 Write,done 对它就是空气。
四、Publish 必须同时听三件事
1 | select { |
这就是「尊重 ctx」:可阻塞的等待里,取消和关闭都要能打断入队。
这里有一个 Go 的坑。Close 之后立刻 Publish,若队列还有空位:
1 | <-s.done 就绪 |
同一个 select 会随机挑一个。测试里曾经抽中入队,返回 nil,断言关闭失败。
拆成两个 select(先看 done,再阻塞)能修,但读起来绕。更直接的办法是加一个 closed 标记,先挡关闭后的调用,done 继续负责叫醒已经堵在 select 里的人:
1 | func (s *JSONEventSink) Publish(ctx context.Context, event reporter.Event) error { |
两个东西不是重复:
1 | closed:已经关闭了吗?(布尔,查一眼) |
只有布尔值,堵在满队列上的 Publish 醒不过来;只有 done 放进同一个 select,关闭后的瞬时调用可能误入队。
五、Close、关连接、Wait 不能写反
Close 立刻返回,只表示「通知已发出」。Wait 返回,才表示写事件的那条 goroutine 已经 return。
因此 Close 不能在内部 Wait:Sink 不拥有连接,等 loop 时 Write 可能还堵着,自己死锁。等待是 Session 的责任,而且必须先把 Write 打断。
ChannelSession.Close 的顺序是:
1 | Destroy Runtime ← Cancel,Engine 停止继续生产 |
两处对调会出问题:
1 | 先 Wait 再 conn.Close |
done 解决的是「别再入队 / 别再取下一条」。它解决不了「这一条已经交给内核的 Write」。打断 Write 是连接所有者的责任。
六、把 Publish 错误传回 Engine
有界队列只能让上游变慢。下游已经关掉时,还得让当前 Run 停下来,并且 Session 里的 ToolCall 仍然成对。
Reporter / StreamReporter 现在返回 error。JSONReporter 不再丢掉 Publish 的结果:
1 | func (r *JSONReporter) publish(ctx context.Context, event Event) error { |
Engine 按阶段处理:
| 失败点 | 行为 |
|---|---|
OnThinking / 流式 OnTextDelta | 还没有带 ToolCall 的 Assistant,直接返回 |
OnMessage 且已 Append 了 ToolCalls | EnsureToolObservations 后再返回 |
OnToolCall | 不 Execute,补取消 Observation |
OnToolResult | 工具已经跑完,仍写入真实结果,再记下 reportErr |
工具 goroutine 里用 setReportErr 只保留第一个错误。Wait 之后先看 ctx.Err(),再看 reportErr:用户取消时返回 Canceled,不会被 Sink 错误盖住。
RunSub 用同一套逻辑。它没有持久化 Session,失败时把 Observation 补进局部 contextHistory 再返回。
七、generate 返回时必须取消流式 ctx
只让 Engine 提前 return 不够。GenerateStream 那条 goroutine 可能还堵在往 events 里送。
OnTextDelta 失败时,父 ctx 不一定已经被取消——Channel 正常关闭会先 Destroy,但 Reporter 自己失败时不会。所以 generate 进入时派生一个可取消的 ctx:
1 | ctx, cancel := context.WithCancel(ctx) |
任何返回路径都会取消这次流。Provider 侧的 sendStreamEvent 已经同时听 ctx.Done(),HTTP stream 也把同一个 ctx 传给 SDK。这样「Run 结束」和「读模型的 goroutine 退出」才是一件事。
八、测试要证明什么
test/channel/event_sink_test.go:
- 队列有空位,
Publish立刻成功。 - 用会阻塞的
io.Writer让loop卡在Write上,填满容量为 1 的队列,第三次Publish在cancel后返回context.Canceled。 Close两次不 panic;之后Publish返回ErrEventSinkClosed。Write仍阻塞时Wait不返回;模拟conn.Close让Write返回后,Wait结束。
第二条说明背压:下游不消费时,上游必须能被 ctx 救出来。第四条说明关机契约:Close 叫不醒 Write,Wait 必须发生在连接被关掉之后。
test/engine/reporter_error_test.go:
OnMessage失败时,已入库的 ToolCall 会补 Observation,工具不会执行。OnToolCall失败时同样跳过Execute,并补取消 Observation。
test/engine/stream_cancel_test.go:
- 父 ctx 取消后,流式 Provider 的 goroutine 退出。
- 只有 Reporter 失败、父 ctx 还活着时,
generate的defer cancel仍能让 Provider goroutine 退出。
为了测「满」,构造函数增加了容量参数;生产路径仍走默认 64。相关测试可以在 -race 下通过。
九、阶段七收口
对照系列十的阶段七清单:
| 问题 | 落地 |
|---|---|
| Approval 输入 goroutine 泄漏 | 系列十七:stdin 归 REPL |
| 取消后误写 Session | EnsureToolObservations |
| bash 子进程 / 工具级超时 | CommandContext + ErrToolTimeout |
| 统一任务状态 | canceled / timed_out / failed |
| 工具并发上限 | MaxToolConcurrency |
| Channel 背压 | 有界 EventSink |
| Reporter 并发输出 | Terminal 加锁;Channel 单 goroutine 写 |
| 优雅关闭 | Close → conn.Close → Wait |
| Provider 流取消后退出 | generate 派生 ctx + 测试 |
| 取消、超时、泄漏测试 | Session、bash、Sink、Reporter、Stream |
阶段八才是 Grant 持久化和参数级权限,不在本篇范围。
总结
事件出口要同时成立五条约束:
1 | 队列有界 |
执行侧限制并行数,输出侧限制管道深度,失败时停止当前任务而不是把错误吞掉。到这里,路线图阶段七「资源生命周期和可靠取消」可以收口;下一阶段是审批持久化与最小权限。