从零搭建 Agent Harness 系列(十八)事件背压:有界队列、尊重 Context 与可关闭的 EventSink

系列十七修好了终端审批泄漏。阶段七还剩另一件容易在生产里爆内存的事:事件比网络快时,往哪里堆?以及堆不住时,失败能不能变成 Run 的终态。

工具并发上限限制的是「同时跑几个工具」。事件背压限制的是「事件管道别无限灌」。两者管的不是同一类资源。

本文对应 go-tiny-claw 两次提交:

1
2
bb746e1  feat: 为 Channel EventSink 增加有界队列背压
85904e7 feat: 将 Reporter 错误回传 Engine,关闭时等待 EventSink 写循环退出

落地之后,Channel 侧的事件出口是:

1
2
3
4
5
6
7
有界 queue(默认 64)
唯一 writer goroutine
Publish 尊重 ctx
Close 叫醒阻塞中的 Publish
conn.Close 打断正在进行的 Write
Wait 确认 loop 已退出
Reporter 错误会停掉当前 Run

一、原来的 Publish 直接写连接

改造前,路径是:

1
2
3
4
5
6
7
Engine / Reporter

JSONEventSink.Publish

MessageWriter.Encode

TCP / WebSocket conn

Publish 会先看一眼 ctx,然后同步 Write

1
2
3
4
5
6
func (s *JSONEventSink) Publish(ctx context.Context, event reporter.Event) error {
if err := ctx.Err(); err != nil {
return err
}
return s.writer.Write(event)
}

慢客户端会把 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
2
3
for item := range s.queue {
_ = s.writer.Write(item.event)
}

range 一个 channel 会一直等到 channel 被 close。Session 关掉如果不 close(queue),这条 goroutine 会永远挂着。

close(queue)queue <- item 并发会发生 panic。Engine 的 Destroy 和迟到的 Publish 之间做不到绝对无缝。

所以关机信号不能靠 queue,要另开一个 done

1
2
done:广播「Sink 已关闭」,叫醒所有 select
queue:只承载事件,Close 时不要 close 它

loop 改成:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
func (s *JSONEventSink) loop() {
defer s.wg.Done()
for {
select {
case <-s.done:
return
case item := <-s.queue:
if item.ctx.Err() != nil {
continue
}
_ = s.writer.Write(item.event)
}
}
}

Closeclose(s.done)select 只能叫醒还在等 done / queue 的 loop。一旦已经进了 Writedone 对它就是空气。

四、Publish 必须同时听三件事

1
2
3
4
5
6
7
8
select {
case <-ctx.Done():
return ctx.Err()
case <-s.done:
return ErrEventSinkClosed
case s.queue <- item:
return nil
}

这就是「尊重 ctx」:可阻塞的等待里,取消和关闭都要能打断入队。

这里有一个 Go 的坑。Close 之后立刻 Publish,若队列还有空位:

1
2
<-s.done        就绪
s.queue <- item 也就绪

同一个 select 会随机挑一个。测试里曾经抽中入队,返回 nil,断言关闭失败。

拆成两个 select(先看 done,再阻塞)能修,但读起来绕。更直接的办法是加一个 closed 标记,先挡关闭后的调用done 继续负责叫醒已经堵在 select 里的人:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
func (s *JSONEventSink) Publish(ctx context.Context, event reporter.Event) error {
if s.closed.Load() {
return ErrEventSinkClosed
}

select {
case <-ctx.Done():
return ctx.Err()
case <-s.done:
return ErrEventSinkClosed
case s.queue <- item:
return nil
}
}

func (s *JSONEventSink) Close() {
s.closeOnce.Do(func() {
s.closed.Store(true)
close(s.done)
})
}

两个东西不是重复:

1
2
closed:已经关闭了吗?(布尔,查一眼)
done :正在等待的人,请醒来(channel,能广播)

只有布尔值,堵在满队列上的 Publish 醒不过来;只有 done 放进同一个 select,关闭后的瞬时调用可能误入队。

五、Close、关连接、Wait 不能写反

Close 立刻返回,只表示「通知已发出」。Wait 返回,才表示写事件的那条 goroutine 已经 return

因此 Close 不能在内部 Wait:Sink 不拥有连接,等 loop 时 Write 可能还堵着,自己死锁。等待是 Session 的责任,而且必须先把 Write 打断。

ChannelSession.Close 的顺序是:

1
2
3
4
Destroy Runtime     ← Cancel,Engine 停止继续生产
events.Close() ← closed + close(done),Publish 立刻失败
conn.Close() ← 打断正在进行的 Write
events.Wait() ← 确认 loop 已退出

两处对调会出问题:

1
2
3
4
5
先 Wait 再 conn.Close
→ loop 卡在 Write 上,Wait 永远回不来

先 conn.Close 再 events.Close
→ Publish 仍可能入队,再对已关闭的连接 Write,错误被丢掉

done 解决的是「别再入队 / 别再取下一条」。它解决不了「这一条已经交给内核的 Write」。打断 Write 是连接所有者的责任。

六、把 Publish 错误传回 Engine

有界队列只能让上游变慢。下游已经关掉时,还得让当前 Run 停下来,并且 Session 里的 ToolCall 仍然成对。

Reporter / StreamReporter 现在返回 errorJSONReporter 不再丢掉 Publish 的结果:

1
2
3
func (r *JSONReporter) publish(ctx context.Context, event Event) error {
return r.sink.Publish(ctx, event)
}

Engine 按阶段处理:

失败点行为
OnThinking / 流式 OnTextDelta还没有带 ToolCall 的 Assistant,直接返回
OnMessage 且已 Append 了 ToolCallsEnsureToolObservations 后再返回
OnToolCallExecute,补取消 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
2
ctx, cancel := context.WithCancel(ctx)
defer cancel()

任何返回路径都会取消这次流。Provider 侧的 sendStreamEvent 已经同时听 ctx.Done(),HTTP stream 也把同一个 ctx 传给 SDK。这样「Run 结束」和「读模型的 goroutine 退出」才是一件事。

八、测试要证明什么

test/channel/event_sink_test.go

  1. 队列有空位,Publish 立刻成功。
  2. 用会阻塞的 io.Writerloop 卡在 Write 上,填满容量为 1 的队列,第三次 Publishcancel 后返回 context.Canceled
  3. Close 两次不 panic;之后 Publish 返回 ErrEventSinkClosed
  4. Write 仍阻塞时 Wait 不返回;模拟 conn.CloseWrite 返回后,Wait 结束。

第二条说明背压:下游不消费时,上游必须能被 ctx 救出来。第四条说明关机契约:Close 叫不醒 WriteWait 必须发生在连接被关掉之后。

test/engine/reporter_error_test.go

  1. OnMessage 失败时,已入库的 ToolCall 会补 Observation,工具不会执行。
  2. OnToolCall 失败时同样跳过 Execute,并补取消 Observation。

test/engine/stream_cancel_test.go

  1. 父 ctx 取消后,流式 Provider 的 goroutine 退出。
  2. 只有 Reporter 失败、父 ctx 还活着时,generatedefer cancel 仍能让 Provider goroutine 退出。

为了测「满」,构造函数增加了容量参数;生产路径仍走默认 64。相关测试可以在 -race 下通过。

九、阶段七收口

对照系列十的阶段七清单:

问题落地
Approval 输入 goroutine 泄漏系列十七:stdin 归 REPL
取消后误写 SessionEnsureToolObservations
bash 子进程 / 工具级超时CommandContext + ErrToolTimeout
统一任务状态canceled / timed_out / failed
工具并发上限MaxToolConcurrency
Channel 背压有界 EventSink
Reporter 并发输出Terminal 加锁;Channel 单 goroutine 写
优雅关闭Closeconn.CloseWait
Provider 流取消后退出generate 派生 ctx + 测试
取消、超时、泄漏测试Session、bash、Sink、Reporter、Stream

阶段八才是 Grant 持久化和参数级权限,不在本篇范围。

总结

事件出口要同时成立五条约束:

1
2
3
4
5
6
7
8
9
队列有界

Publish 同时听 ctx 和关闭

Close 能叫醒等待入队的人

关连接才能打断正在进行的 Write

Reporter 错误变成 Run 终态,且 ToolCall 仍然成对

执行侧限制并行数,输出侧限制管道深度,失败时停止当前任务而不是把错误吞掉。到这里,路线图阶段七「资源生命周期和可靠取消」可以收口;下一阶段是审批持久化与最小权限。