AI Agent 的安全挑战及防护策略
2026-07-22 3416492
2026-07-22 0
一个 ReAct Agent 跑完一次,外部几乎看不见内部发生了什么:

没有这些数据,生产环境的问题排查和成本控制都是瞎猜。
Eino 的 Callback 系统是一个侵入性极低的观测钩子:不修改任何组件的业务逻辑,只在组件的生命周期关键点插入回调。
Handler 是所有观察者必须实现的接口:
复制代码// 5个生命周期回调:
// OnStart / OnEnd / OnError / OnStartWithStreamInput / OnEndWithStreamOutput
type Handler interface {
OnStart(ctx context.Context, info *RunInfo, input CallbackInput) context.Context
OnEnd(ctx context.Context, info *RunInfo, output CallbackOutput) context.Context
OnError(ctx context.Context, info *RunInfo, err error) context.Context OnStartWithStreamInput(ctx context.Context, info *RunInfo,
input *schema.StreamReader[CallbackInput]) context.Context
OnEndWithStreamOutput(ctx context.Context, info *RunInfo,
output *schema.StreamReader[CallbackOutput]) context.Context
}
每个方法的返回值是 context.Context。这个设计让同一个 Handler 的不同 timing 之间可以传递状态:
复制代码// OnStart 在 context 里存一个开始时间
func OnStart(ctx, info, input) context.Context {
return context.WithValue(ctx, startKey{}, time.Now())
}// OnEnd 从 context 里取出开始时间,计算耗时
func OnEnd(ctx, info, output) context.Context {
start := ctx.Value(startKey{}).(time.Time)
log.Printf("[%s] 耗时: %v", info.Name, time.Since(start))
return ctx
}
注意:这个 context 链只在同一个 Handler 的同一次调用内流动,不会跨 Handler。
RunInfo 描述了"是谁触发了这个回调":
复制代码type RunInfo struct {
Name string // 节点名(compose.WithNodeName 指定)
Type string // 实现类型,如 "OpenAI"、"DeepSeek"
Component components.Component // 组件种类:ChatModel/Tool/Retriever...
}
这三个字段是 Callback 逻辑里最重要的过滤条件:不同种类的组件要采集不同的指标。
| Timing | 触发时机 | 用于 |
|---|---|---|
OnStart | 组件开始处理前(非流式输入) | 记录入参、开始时间 |
OnEnd | 组件成功返回后(非流式输出) | 记录输出、耗时、Token 数 |
OnError | 组件返回错误时 | 报错统计、告警 |
OnStartWithStreamInput | 组件接收流式输入时 | 较少用,Collect/Transform 场景 |
OnEndWithStreamOutput | 组件产出流式输出时 | Token 流统计、流式日志 |
大多数使用场景只需要 OnStart + OnEnd + OnError。
复制代码// 在 main 或 TestMain 里调用,一次配置全部生效
// 注意:不是线程安全的,不能在 graph 运行时调用
callbacks.AppendGlobalHandlers(myTracingHandler, myMetricsHandler)
全局 Handler 对所有组件、所有 Graph 调用都生效。适合全平台的 OTel 追踪、Token 计量。
复制代码// 只对这一次 Invoke 生效
result, err := runner.Invoke(ctx, input,
compose.WithCallbacks(myDebugHandler),
)
这种方式适合开发调试或者租户级别的独立观测,不影响其他租户。
HandlerBuilder:最快写一个 Handler不用实现全部 5 个方法,只订阅你关心的 timing:
复制代码tokenCounter := callbacks.NewHandlerBuilder().
OnEndFn(func(ctx context.Context, info *callbacks.RunInfo, output callbacks.CallbackOutput) context.Context {
// 只处理 ChatModel 的输出,用 model.ConvCallbackOutput 安全类型转换
mo := model.ConvCallbackOutput(output)
if mo != nil && mo.TokenUsage != nil {
metrics.Add("llm.tokens.total", float64(mo.TokenUsage.TotalTokens), map[string]string{
"model": info.Name,
"type": info.Type,
})
}
return ctx
}).
Build()
HandlerBuilder 内部实现了 TimingChecker——只注册了 OnEnd 的 Handler,在其他 timing 时框架会直接跳过,不分配 goroutine 也不复制 stream:
复制代码func (hb *handlerImpl) Needed(_ context.Context, _ *RunInfo, timing CallbackTiming) bool {
switch timing {
case TimingOnEnd:
return hb.onEndFn != nil // 只注册了 OnEnd → 其他 timing 返回 false
case TimingOnStart:
return hb.onStartFn != nil // 没注册 → false,跳过
// ...
}
}
HandlerHelper:按组件类型分发问题:CallbackInput / CallbackOutput 是 interface{},你在 OnEnd 里拿到的可能是 ChatModel 的输出,也可能是 Tool 的输出,需要手动判断类型。
HandlerHelper 把这个分发逻辑封装好了:
复制代码helper := utils_callbacks.NewHandlerHelper().
ChatModel(&utils_callbacks.ModelCallbackHandler{
OnEnd: func(ctx context.Context, info *callbacks.RunInfo, output *model.CallbackOutput) context.Context {
// 这里拿到的已经是强类型的 *model.CallbackOutput
if output.TokenUsage != nil {
log.Printf("[model:%s] tokens=%d", info.Name, output.TokenUsage.TotalTokens)
}
return ctx
},
OnEndWithStreamOutput: func(ctx context.Context, info *callbacks.RunInfo,
output *schema.StreamReader[*model.CallbackOutput]) context.Context {
// 流式输出:必须关闭 StreamReader,否则 goroutine 泄漏
go func() {
defer output.Close()
var totalTokens int
for chunk, err := output.Recv(); err == nil; chunk, err = output.Recv() {
if chunk.TokenUsage != nil {
totalTokens += chunk.TokenUsage.TotalTokens
}
}
log.Printf("[model:%s] stream total tokens=%d", info.Name, totalTokens)
}()
return ctx
},
}).
Tool(&utils_callbacks.ToolCallbackHandler{
OnStart: func(ctx context.Context, info *callbacks.RunInfo, input *tool.CallbackInput) context.Context {
log.Printf("[tool:%s] called with: %s", info.Name, input.ArgumentsInJSON)
return context.WithValue(ctx, toolStartKey{}, time.Now())
},
OnEnd: func(ctx context.Context, info *callbacks.RunInfo, output *tool.CallbackOutput) context.Context {
start := ctx.Value(toolStartKey{}).(time.Time)
log.Printf("[tool:%s] took=%v result=%s", info.Name, time.Since(start), output.Content)
return ctx
},
}).
Handler()callbacks.AppendGlobalHandlers(helper)
HandlerHelper 支持的组件类型:Prompt / ChatModel / Embedding / Indexer / Retriever / Loader / Transformer / Tool / ToolsNode / Agent + 对应的 Agentic 变体。
回到 E42 讲过的 taskManager.execute,它在每个节点执行前都会调用 initNodeCallbacks:
复制代码// compose/graph_manager.go
func (t *taskManager) execute(currentTask *task) {
// 注入:把节点的 RunInfo + 本次 WithCallbacks 的 Handler 注入 context
ctx := initNodeCallbacks(
currentTask.ctx,
currentTask.nodeKey,
currentTask.call.action.nodeInfo,
currentTask.call.action.meta,
t.opts...,
)
// 真正执行节点
task.output, task.err = runWrapper(ctx, call.action, ...)
}
initNodeCallbacks 的逻辑:
复制代码// compose/utils.go
func initNodeCallbacks(ctx context.Context, key string, info *nodeInfo, meta *executorMeta, opts ...Option) context.Context {
// 1. 从 meta 提取 RunInfo.Component、RunInfo.Type(组件种类 + 实现类型)
ri := &callbacks.RunInfo{}
if meta != nil {
ri.Component = meta.component
ri.Type = meta.componentImplType // e.g. "OpenAI" / "DeepSeek"
}
// 2. 从 info 提取 RunInfo.Name(节点名,WithNodeName 指定的)
if info != nil {
ri.Name = info.name
} // 3. 收集本次 Invoke 传入的 WithCallbacks(handler)(仅路径匹配的)
var cbs []callbacks.Handler
for i := range opts {
if len(opts[i].paths) > 0 {
for _, k := range opts[i].paths {
if k.path[0] == key {
cbs = append(cbs, opts[i].handler...)
}
}
}
} // 4. 把全局 Handler + 路径匹配的 Handler 注入到 context
if len(cbs) == 0 {
return icb.ReuseHandlers(ctx, ri) // 只用全局 Handler
}
return icb.AppendHandlers(ctx, ri, cbs...) // 全局 + 额外 Handler
}
这就是 Callback 的注入路径:Graph 执行前每个节点都会调用这个函数,把 RunInfo 和 Handler 写进 context,节点内部的组件实现(ChatModel/Tool)调用 callbacks.OnStart(ctx, ...) 时就能找到它们。
当 ChatModel 流式输出时,OnEndWithStreamOutput 收到的是一个 *schema.StreamReader 副本(框架自动 stream.Copy)。
复制代码// 错误做法:直接在 OnEndWithStreamOutput 同步消费
OnEndWithStreamOutput: func(ctx, info, output) context.Context {
for chunk, _ := output.Recv() { // 会阻塞,因为流还在传输
// ...
}
// 忘记 output.Close() → goroutine 泄漏
return ctx
},
正确做法:异步消费 + 一定要 Close
复制代码OnEndWithStreamOutput: func(ctx, info, output) context.Context {
go func() {
defer output.Close() // 关闭这个副本,否则底层资源不释放
for {
chunk, err := output.Recv()
if err != nil { break }
// 处理 chunk...
}
}()
return ctx // 立即返回,不阻塞主流程
},
框架注释说得很清楚:
复制代码import (
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/attribute"
)var tracer = otel.Tracer("deepflux/agent")otlHandler := callbacks.NewHandlerBuilder().
OnStartFn(func(ctx context.Context, info *callbacks.RunInfo, input callbacks.CallbackInput) context.Context {
// 每个节点开始时创建一个 OTel Span
ctx, span := tracer.Start(ctx, info.Name,
trace.WithAttributes(
attribute.String("eino.component", string(info.Component)),
attribute.String("eino.type", info.Type),
))
// 把 span 存进 context,供 OnEnd 结束它
return context.WithValue(ctx, spanKey{}, span)
}).
OnEndFn(func(ctx context.Context, info *callbacks.RunInfo, output callbacks.CallbackOutput) context.Context {
if span, ok := ctx.Value(spanKey{}).(trace.Span); ok {
// Token 信息作为 Span attribute
if mo := model.ConvCallbackOutput(output); mo != nil && mo.TokenUsage != nil {
span.SetAttributes(
attribute.Int("llm.prompt_tokens", mo.TokenUsage.PromptTokens),
attribute.Int("llm.completion_tokens", mo.TokenUsage.CompletionTokens),
)
}
span.End()
}
return ctx
}).
OnErrorFn(func(ctx context.Context, info *callbacks.RunInfo, err error) context.Context {
if span, ok := ctx.Value(spanKey{}).(trace.Span); ok {
span.RecordError(err)
span.SetStatus(codes.Error, err.Error())
span.End()
}
return ctx
}).
Build()callbacks.AppendGlobalHandlers(otlHandler)
这样每个节点的执行都会在 Jaeger/Tempo 里有对应的 Span,自动组成 Trace 树。
Eino Callback 系统的设计核心是最小侵入 + 类型安全:
| 机制 | 作用 |
|---|---|
Handler 接口 + 5 个 Timing | 覆盖组件全生命周期,统一契约 |
RunInfo 三字段 | 标识"谁在运行",用于 Metric label / Span attribute |
TimingChecker.Needed | 避免没注册的 timing 产生不必要的 stream 复制 |
HandlerBuilder | 只订阅关心的 timing,其余跳过 |
HandlerHelper | 按组件类型强类型分发,无需手写 type switch |
initNodeCallbacks | Graph 节点执行时注入 RunInfo + Handler,不需要组件感知 Graph |
流式副本 + Close | 流式 Callback 的资源管理约定,漏 Close = goroutine 泄漏 |
在 DeepFlux 的 server/internal/observability 模块里,就是用这套机制把每次 Agent 调用的 Token 消耗、工具调用延迟、Retriever 召回延迟,以 OTel Span 的形式发送给 Tempo,再通过 Grafana 展示出来,完整覆盖从 HTTP 入口到 LLM 调用的链路追踪。
代码来源:eino/callbacks/interface.go · eino/callbacks/handler_builder.go · eino/utils/callbacks/template.go · eino/compose/utils.go