乐于分享
好东西不私藏

Callback 源码:aspect_inject 切面注入(第87篇-E73)

Callback 源码:aspect_inject 切面注入(第87篇-E73)

上一篇留了个缺口:LLM 那一跳没有 OTel span,补法是 Eino 的 callback 系统。但要理解它,得先回答一个更一般的问题 —— 一个编译好的 Runnable(Graph 已经编译成了执行计划),怎么在不改任何节点代码的前提下,给每个节点的执行前后插上日志、追踪、计量?

传统答案是中间件:包一层 HTTP handler 那样的洋葱。但 Agent 图不是一条链——它是 DAG,节点可能并发执行,有的节点吃流式输入、有的吐流式输出。"包一层"在这里意味着什么?切面插在哪、上下文怎么传、多个切面什么顺序?Eino 的答案全在一个不到 400 行的包里:callbacks。这个包有个源码文件叫 aspect_inject.go——切面注入,名字就是答案。

(一)两条注入路径:框架包,还是组件自己埋

先看全景。Eino 的切面注入有两条路,由组件自己选:

路径 ① 框架包装(默认)  Graph 编译时,把每个节点的执行函数包一层:  onStart(切面)  节点逻辑  onError 或 onEnd(切面)  组件代码一行不改路径 ② 组件自管(实现 Checker 接口)  组件声明 IsCallbacksEnabled() = true   框架跳过包装(component_to_graph_node.go:41 的 !isComponentCallbackEnabled)   组件在自己的 Generate/Stream 内部手动调 callbacks.OnStart / OnEnd / OnError

路径 ① 的全部实现就 16 行(compose/utils.go):

funcrunWithCallbacks[IOTOptionany](r func(...) (O, error),	onStart on[I], onEnd on[O], onError on[error]) func(...) (O, error) {return func(ctx context.Context, input I, opts ...TOption) (output O, err error) {		ctx, input = onStart(ctx, input)      // 切面:进		output, err = r(ctx, input, opts...)  // 真正的节点逻辑if err != nil {			ctx, err = onError(ctx, err)      // 切面:出错return output, err}		ctx, output = onEnd(ctx, output)      // 切面:出return output, nil}}

这就是切面的物理形态:一个高阶函数。编译期把节点函数 r 换成包装后的版本,运行期调用方毫无感知。你的 Lambda、你的 ChatModel、你的 Tool——只要没声明自管,进图时都被这么包过。

那为什么还需要路径 ②?看一个真实的例子(eino-ext/libs/acl/openai/chat_model.go,OpenAI 适配器的 Generate):

func(c *Client) Generate(ctx context.Context, in []*schema.Message, ...) (	outMsg *schema.Message, err error) {	req, cbInput, reqOpts, specOptions, err := c.genRequest(ctx, in, opts...)...	ctx = callbacks.OnStart(ctx, cbInput)      // 手动触发 OnStart,输入是结构化的 model.CallbackInputdefer func() {if err != nil {			callbacks.OnError(ctx, err)         // defer 里接住错误}}()	resp, err := c.cli.CreateChatCompletion(ctx, *req, reqOpts...)...	callbacks.OnEnd(ctx, &model.CallbackOutput{ // 手动触发 OnEnd,带 TokenUsage		Message:    outMsg,		Config:     cbInput.Config,		TokenUsage: toModelCallbackUsage(outMsg.ResponseMeta),})return outMsg, nil}

组件自管换来三样框架包装给不了的东西:

  1. 结构化的输入输出
    。框架包装只能拿到 any,组件自己埋点可以传 model.CallbackInput——messages、tools、TokenUsage 都是强类型。观测系统最需要的 token 数就藏在 TokenUsage 里。
  2. 流式的精确时机
    。框架包装只能在整个流结束时触发 OnEnd;组件自管可以在流结束的准确位置触发 OnEndWithStreamOutput
  3. 埋点位置即语义位置
    OnEnd 在 buildGenerateResponse 之后、modifier 之前——观测到的就是真正要返回的东西。

两条路径的切换开关是个单方法接口(components/types.go):

type Checker interface {IsCallbacksEnabled() bool}

组件实现它并返回 true,框架就让位(component_to_graph_node.go:40):

run := runnableLambda(invoke, stream, collect, transform,!meta.isComponentCallbackEnabled, // 组件自管 → 框架不包)

(二)注册表不在全局变量里,在 ctx 里

下一个问题:切面(handler)注册在哪,节点执行时怎么找到它们?

答案有点反直觉:不在全局 map,在 context 里internal/callbacks/manager.go):

type manager struct {	globalHandlers []Handler // 进程级:AppendGlobalHandlers 注册,全局所有图共享	handlers       []Handler // 本次运行级:compose.WithCallbacks(h) 传入	runInfo        *RunInfo  // 当前是谁在执行:节点名 + 组件类别}var GlobalHandlers []Handler // 唯一的全局状态funcmanagerFromCtx(ctx context.Context) (*manager, bool) {	m, ok := ctx.Value(CtxManagerKey{}).(*manager)if !ok || m == nil {return nilfalse}	n := *m        // 关键:拷贝再改return &n, true}

三层 handler 来源,作用域从大到小:

callbacks.AppendGlobalHandlers(tracingHandler)              // 进程级:所有图、所有节点runnable.Invoke(ctx, input, compose.WithCallbacks(h2))      // 运行级:这张图这次运行compose.WithCallbacks(h3).DesignateNode("model")            // 节点级:只挂 model 节点

用 ctx 当注册表有个直接好处:子图自动继承。父图把 handler 放进 ctx,子图节点从同一个 ctx 取——嵌套图的观测不用额外接线(上一篇的 trace_id 靠 ctx 传播,这里切面靠同一个 ctx 传播,是同一个设计判断)。

managerFromCtx 里那行拷贝也不是随手写的:Go 的 context.Value 是不可变链,任何修改都必须造新值。拷贝保证了并发的两个节点互不污染——各自拿到自己的 manager 副本,往里加各自的 runInfo。

图运行时每个节点执行前重新装配一次(graph_manager.go:296):

func(t *taskManager) execute(currentTask *task) {...	ctx := initNodeCallbacks(currentTask.ctx, currentTask.nodeKey, ...) // 装配:挂 runInfo + 过滤定向 handler	currentTask.output, currentTask.err = t.runWrapper(ctx, ...)        // 执行:包装过的节点函数}

(三)On() 的调度:三个 30 行以内的细节

所有时机触发最终走到同一个函数(internal/callbacks/inject.go 的 On())。它短,但藏着三个值得逐行读的细节。

细节 1:RunInfo 借道 ctx 传递。 OnStart 触发时 manager 里的 runInfo 被取出来塞进 ctx;OnEnd 触发时优先从 manager 取、取不到再从 ctx 里捞回来:

if start {	info = nMgr.runInfo	nMgr.runInfo = nil	ctx = context.WithValue(ctx, CtxRunInfoKey{}, info) // 塞进 ctxelse {if nMgr.runInfo != nil {		info = nMgr.runInfoelse {		info, _ = ctx.Value(CtxRunInfoKey{}).(*RunInfo) // 从 ctx 捞回}}

为什么绕这一圈?因为组件自管路径里,组件调 callbacks.OnStart(ctx, ...) 之后,ctx 可能被组件自己再包过几层再传给自己的 OnEnd——manager 早就换了好几代,但 ctx 链没断。RunInfo 跟着 ctx 走,谁触发都能拿到同一个身份

细节 2:TimingChecker 过滤,未注册零开销。 handler 是可选实现这个接口的:

type TimingChecker interface {Needed(ctx context.Context, info *RunInfo, timing CallbackTiming) bool}

On() 在调用任何 handler 前先问一遍 Needed,没注册的时机直接不进列表。这不是锦上添花——流式时机(OnStartWithStreamInput / OnEndWithStreamOutput)触发时框架要给每个 handler 复制一份流(下文讲)。过滤掉不需要的 handler,就少 copy 一份流、少起一个 goroutine。NewHandlerBuilder 造出来的 handler 自动实现 TimingChecker——你只注册了 OnStart,其他四个时机框架连碰都不碰。

细节 3:OnStart 逆序、OnEnd 正序——洋葱。 handler 列表的拼接顺序是 run 级在前、global 在后,然后:

funcOnStartHandle[Tany](ctx context.Context, input T, runInfo *RunInfo, handlers []Handler) ... {for i := len(handlers) - 1; i >= 0; i-- {   // 逆序:列表末尾的先进		ctx = handlers[i].OnStart(ctx, runInfo, input)}return ctx, input}funcOnEndHandle[Tany](ctx context.Context, output T, runInfo *RunInfo, handlers []Handler) ... {for _, handler := range handlers {          // 正序:列表开头的先出		ctx = handler.OnEnd(ctx, runInfo, output)}return ctx, output}

demo(复刻 On() 全调度逻辑,纯标准库)场景 A 跑出来的真实顺序:

====== 场景 A:global + run 双 handler,洋葱顺序 ======  事件流(节选 model 节点):    global.OnStart(ChatModel.model)    run.OnStart(ChatModel.model)    run.OnEnd(ChatModel.model)    global.OnEnd(ChatModel.model)  ↑ OnStart 逆序(global 先进)→ 节点执行 → OnEnd 正序(global 后出)= 洋葱

global 在最外层,run 级在里面包着节点。源码注释说得很直白:global handlers 先跑,给分布式追踪、指标这类"必须看到一切"的埋点更高优先级。逆序+正序组合出来的正是 HTTP 中间件的洋葱语义——最先进去的最后出来。

顺带一提,demo 写自检时我在这里栽了第一个跟头:断言把洋葱顺序写反了(以为 run 级在最外层),跑出来红的——demo 忠实复刻了源码,错的是我对源码的预期。这种红恰恰是自检的价值:它逼你回去重读那两行 for 循环。

(四)DesignateNode:定向切面怎么实现

"只给 model 节点加计量"这种需求,靠 Option.paths 过滤。注册侧:

runnable.Invoke(ctx, input, compose.WithCallbacks(myHandler).DesignateNode("model"))

装配侧(initNodeCallbacks,每个节点执行前都跑一遍):

var cbs []callbacks.Handlerfor i := range opts {if len(opts[i].handler) != 0 {if len(opts[i].paths) != 0 {for _, k := range opts[i].paths {if len(k.path) == 1 && k.path[0] == key { // 当前节点 key 命中定向路径					cbs = append(cbs, opts[i].handler...)break}}}}}if len(cbs) == 0 {return icb.ReuseHandlers(ctx, ri)      // 没命中:只换 runInfo,handler 原样继承}return icb.AppendHandlers(ctx, ri, cbs...) // 命中:追加定向 handler

每个节点执行前都问一遍"这次轮到的节点在不在你的定向名单里"——在就挂上,不在就只换 RunInfo 继续。demo 场景 B 验证了行为:

====== 场景 B:DesignateNode("model") 只命中目标节点 ======  designated 触发 2 次(都应在 model 节点):    designated.OnStart(ChatModel.model)    designated.OnEnd(ChatModel.model)

三节点图跑下来只触发 2 次(OnStart + OnEnd 各一次),prompt 和 tool 节点完全无感。

(五)错误路径,和流式的代价

错误路径的语义在 runWithCallbacks 里一目了然:出错走 OnError、不再走 OnEnd。demo 场景 D:

====== 场景 D:节点报错 → OnError 触发、OnEnd 不触发 ======  graph 返回错误: rate limited    global.OnStart(ChatModel.model)    global.OnError(ChatModel.model): rate limited    errOnly.OnError(ChatModel.model): rate limited

注意 errOnly——一个只实现了 OnError 的 handler(TimingChecker 对 OnStart/OnEnd 返回 false),它确实没出现在 OnStart 事件里。过滤生效。

流式则是这套系统代价最重的部分,官方 doc.go 把坑列得明明白白,三条最致命:

  1. 流的 N+1 拷贝
    。OnEndWithStreamOutput 触发时,框架给每个注册了该时机的 handler 复制一份流(output.Copy(len(handlers)+1)——N 个 handler 拷 N+1 份,最后一份留给下游)。handler 必须 close 自己那份,任何一份没 close,原始流就无法释放,整条管道的 goroutine 全部泄漏。这就是细节 2 的 TimingChecker 如此重要的原因:少一个注册流式时机的 handler,就少一份拷贝。
  2. 流中的错误不走 OnError
    。组件已经把 StreamReader 返回给调用方之后,流中途出错(网络断、上游超时)是在 Recv 里暴露的——OnError 看不到。想在流式场景做错误告警,得在流的消费侧做,不能只依赖 callback。
  3. 不要修改 Input/Output
    。所有下游节点和 handler 共享同一个指针,不是深拷贝。在 handler 里顺手改了 message 内容,并发执行的图里就是数据竞争。

还有一条管理面的:AppendGlobalHandlers不是线程安全的——只能在进程启动时调一次,和图的执行并发调用就是数据竞争。命名从旧版的 InitCallbackHandlers 改成 Append,但"启动时调一次"的约束没变。

小结

问题
机制
关键源码
在哪插
两条路:框架包装(默认高阶函数)/ 组件自管(Checker 让位)
runWithCallbacks
 + chat_model.go
注册表在哪
manager 挂 ctx:global + run 两层,子图自动继承
manager.go
怎么定向
Option.paths 逐节点过滤,命中追加、不命中只换 RunInfo
initNodeCallbacks
什么顺序
OnStart 逆序、OnEnd 正序,global 最外层的洋葱
OnStartHandle/OnEndHandle
怎么省开销
TimingChecker 未注册不进列表 → 流式少拷贝
On()
 过滤段

几条贯穿的设计判断:

  • 切面 = 高阶函数 + ctx
    。不需要代码织入、不需要反射注册中心,一个 16 行的包装函数加一个挂在 ctx 里的注册表,就把 DAG 上的横切关注点解决了。
  • 自管是逃生舱,不是常规路径
    。默认框架包(零侵入),需要结构化数据或流式精确时机时组件才接管——开关只有一个 bool 方法。
  • 身份跟着 ctx 走
    。RunInfo 借道 ctx 传递,和上一篇的 trace_id 是同一个哲学:ctx 断了,切面就瞎了。
  • 开销和注册面成正比
    。未注册的时机不是"空调用"而是"零调用"——流式拷贝按注册数付费,TimingChecker 是省钱的那道闸。
  • 流式永远是例外路径
    。拷贝、泄漏、错误不可见——三个坑全部集中在流式时机,非流式的 OnStart/OnEnd/OnError 干净得像教科书。

下一篇把这套切面用起来:实战接入 Langfuse,给 Agent 装"行车记录仪"——五分钟跑通,再看 handler 怎么消费每个时机。