针对群聊机器人场景设计的消息处理引擎,采用 行为树 + 管线 双层架构。
- 行为树决策:用树状结构管理消息路由,天然支持优先级和互斥,告别 if-else 面条代码
- 子树系统:主树 + 可挂载子树,支持运行时动态注册、替换、注销
- 管线执行:多个 Pass 串联的流水线,单一职责,逻辑清晰
- 全面动态化:Pass / Pipeline / 子树均支持运行时注册、热替换、注销,原生支持插件系统
- 动态路由:RouterPass 支持管线间的动态跳转,轻松处理多轮交互
- 类型安全:MessageContext 泛型 Get/Set,编译期类型检查
- 超时降级:内置超时控制和降级管线,系统级横切关注点
- 软错误:SoftError 机制,非关键错误不中断管线
- 可插拔存储:StateStore 接口抽象,内存/Redis/数据库无缝切换,支持原子操作(CAS、计数器、SETNX)
- 执行链路追踪:可选的 TraceSpan 树,记录 BT 决策、管线执行、Pass 级耗时,支持导出到 OpenTelemetry/Jaeger
- 行为树可视化:conduit/debug 子包提供 BTDebugger,输出 Mermaid 流程图或 Graphviz DOT 格式
- 零第三方依赖:全部使用 Go 标准库实现
import (
"github.com/zrurf/conduit"
"github.com/zrurf/conduit/store/memory"
)
// 1. 创建引擎
store := memory.New()
engine := conduit.New(store,
conduit.WithWorkers(4),
conduit.WithTimeout(5*time.Second),
conduit.WithFallbackPipeline("pipeline.fallback"),
conduit.WithOnResponse(func(ctx *conduit.MessageContext, err error) {
for _, msg := range ctx.Output {
sendToPlatform(msg)
}
}),
)
// 2. 注册子树(可复用的决策分支)
engine.MustRegisterSubtree("subtree.admin",
conduit.NewSequence(
conduit.NewCondition(func(ctx *conduit.MessageContext) bool {
return strings.HasPrefix(ctx.RawMsg, "/admin")
}),
conduit.NewAction("pipeline.admin"),
),
)
// 3. 创建主树,引用子树
bt := conduit.NewBehaviorTree(
engine.NewSubtreeRef("subtree.admin"),
conduit.NewAction("pipeline.default"),
)
engine.SetBehaviorTree(bt)
// 4. 注册管线(动态模式,支持热替换)
engine.MustRegisterPass("pass.check_auth", &CheckAuthPass{})
engine.MustRegisterPass("pass.admin_cmd", &AdminCommandPass{})
engine.MustRegisterPipeline(conduit.NewPipelineFromIDs("pipeline.admin",
"pass.check_auth", "pass.admin_cmd",
))
engine.MustRegisterPipeline(conduit.NewPipeline("pipeline.default",
&DefaultReplyPass{},
))
// 5. 启动
engine.Start()
defer engine.Stop()
// 6. 投递消息
engine.Submit(&conduit.InputMessage{
UserID: "user123", GroupID: "group456",
Content: "hi", IsGroup: true,
})| 概念 | 说明 |
|---|---|
| Engine | 顶层调度器,接收消息、驱动行为树、执行管线、超时降级 |
| 主树 (MainTree) | 决策层入口,每个 Engine 有且仅有一个主树,所有消息从主树进入 |
| 子树 (Subtree) | 可复用的行为树片段,通过注册表管理,可动态挂载到主树 |
| SubtreeRef | 子树引用节点,在 Tick 时从注册表动态解析子树,支持热替换 |
| Pipeline | 执行层,由多个 Pass 串联的流水线,支持静态和动态两种模式 |
| Pass | 管线上的处理节点,单一职责,通过注册表管理支持热替换 |
| RouterPass | 特殊 Pass,可动态跳转到其他管线 |
| MessageContext | 请求级上下文(黑板模式),Pass 间传递中间结果 |
| StateStore | 全局状态存储接口,支持多种后端 |
Conduit 的所有核心组件均支持运行时动态操作,是插件系统和热更新的基础。
| 注册表 | 管理对象 | 操作 |
|---|---|---|
| PassRegistry | Pass | Register / RegisterOrReplace / Unregister / Get / IDs |
| PipelineRegistry | Pipeline | Register / RegisterOrReplace / Unregister / Get / IDs |
| SubtreeRegistry | Subtree | Register / RegisterOrReplace / Unregister / Get / IDs |
使用 NewPipelineFromIDs 创建的动态管线,在每次执行时从 PassRegistry 解析 Pass。
替换注册表中的 Pass 后,所有引用它的管线自动使用新版本:
// 注册
engine.RegisterPass("pass.ai_call", &AICallPass{Client: client})
// 热替换(立即生效)
engine.RegisterOrReplacePass("pass.ai_call", &NewAICallPass{Client: newClient})
// 注销
engine.UnregisterPass("pass.ai_call")// 替换整条管线
engine.RegisterOrReplacePipeline(newPipeline)
// 注销管线
engine.UnregisterPipeline("pipeline.old")newTree := conduit.NewBehaviorTree(...)
engine.SetBehaviorTree(newTree) // 运行时替换整个主树子树是可复用的行为树片段,通过 SubtreeRef 挂载到主树任意位置:
// 注册子树
engine.RegisterSubtree("subtree.game",
conduit.NewSequence(gameCondition, gameAction))
// 在主树中引用
bt := conduit.NewBehaviorTree(
engine.NewSubtreeRef("subtree.game"),
conduit.NewAction("pipeline.default"),
)
// 运行时替换子树(所有 SubtreeRef 自动使用新版本)
engine.RegisterOrReplaceSubtree("subtree.game",
conduit.NewSequence(newGameCondition, newGameAction))
// 注销子树(SubtreeRef 返回 BTFailure,回退到其他分支)
engine.UnregisterSubtree("subtree.game")子树支持"先引用、后注册"模式,适用于插件延迟加载的场景。
Conduit 的动态性天然支持插件系统。宿主程序可以通过统一的接口加载、升级和卸载插件:
// Plugin 是插件的统一接口
type Plugin interface {
// ID 返回插件的唯一标识
ID() string
// Install 将插件的 Pass、管线、子树注册到引擎
Install(engine *conduit.Engine) error
// Uninstall 从引擎注销插件注册的所有资源
Uninstall(engine *conduit.Engine) error
}// GamePlugin 猜谜游戏插件
type GamePlugin struct{}
func (p *GamePlugin) ID() string { return "plugin.game" }
func (p *GamePlugin) Install(e *conduit.Engine) error {
// 注册 Pass
e.RegisterPass("plugin.game.init", &GameInitPass{})
e.RegisterPass("plugin.game.play", &GamePlayPass{})
e.RegisterPass("plugin.game.result", &GameResultPass{})
// 注册管线(动态模式,支持热替换)
e.RegisterPipeline(conduit.NewPipelineFromIDs("plugin.game.play",
"plugin.game.init", "plugin.game.play", "plugin.game.result",
))
// 注册子树(决策分支)
e.RegisterSubtree("subtree.game",
conduit.NewSequence(
conduit.NewCondition(func(ctx *conduit.MessageContext) bool {
_, ok := conduit.Get[string](ctx, "game_active")
return ok
}),
conduit.NewAction("plugin.game.play"),
),
)
// 挂载子树到主树:需要重建主树以包含新子树
// 或者预先在主树中预留 SubtreeRef,插件注册子树后自动生效
return nil
}
func (p *GamePlugin) Uninstall(e *conduit.Engine) error {
// 按相反顺序注销
e.UnregisterSubtree("subtree.game")
e.UnregisterPipeline("plugin.game.play")
e.UnregisterPass("plugin.game.result")
e.UnregisterPass("plugin.game.play")
e.UnregisterPass("plugin.game.init")
return nil
}// 升级插件:只需重新 Install,RegisterOrReplace 自动替换
oldPlugin.Uninstall(engine)
newPlugin.Install(engine)主树可以预先预留子树槽位,插件注册子树后自动生效:
// 主树预留插件槽位(子树可以稍后注册)
bt := conduit.NewBehaviorTree(
engine.NewSubtreeRef("subtree.admin"), // 管理命令插件
engine.NewSubtreeRef("subtree.game"), // 游戏插件
engine.NewSubtreeRef("subtool.music"), // 音乐插件
conduit.NewAction("pipeline.default"), // 兜底
)
engine.SetBehaviorTree(bt)
// 未注册的子树:SubtreeRef 返回 BTFailure,跳过该分支
// 注册子树后:SubtreeRef 自动解析并执行
plugin.Install(engine) // 子树生效conduit/
├── behavior.go # 行为树节点(Selector, Sequence, Condition, Action)
├── context.go # MessageContext 及泛型 Get/Set
├── debug/
│ └── debugger.go # BTDebugger(Mermaid / DOT 可视化)
├── doc.go # Go 包文档
├── engine.go # Engine 核心逻辑
├── engine_options.go # 函数式选项
├── errors.go # SoftError、PassPanicError 类型
├── message.go # InputMessage、Message 类型
├── metrics.go # Metrics 接口
├── pass.go # Pass、RouterPass、PassRegistry
├── pipeline.go # Pipeline、PipelineRegistry
├── store.go # StateStore 接口
├── subtree.go # SubtreeRegistry、SubtreeRef
├── trace.go # 执行链路追踪(TraceSpan)
├── docs/ # 详细文档
├── store/
│ └── memory/
│ └── store.go # 内存存储实现
└── test/
├── behavior_test.go
├── context_test.go
├── debugger_test.go
├── engine_test.go
├── errors_test.go
├── pass_test.go
├── pipeline_test.go
├── store_memory_test.go
├── subtree_test.go
├── trace_test.go
└── yield_test.go