
WeKnora 事件系统使用指南在 Chat Pipeline 中实现监控、日志与可观测性【免费下载链接】WeKnoraOpen-source LLM knowledge platform: turn raw documents into a queryable RAG, an autonomous reasoning agent, and a self-maintaining Wiki.项目地址: https://gitcode.com/GitHub_Trending/we/WeKnora导读WeKnora 内置了一套轻量级的事件总线EventBus机制让查询处理链路Query → Rewrite → Retrieval → Rerank → Chat Completion中的每一步都能被发射事件并被独立监听从而在不侵入业务逻辑的前提下实现监控指标、结构化日志、用户行为分析和错误追踪。本文以 internal/event/usage_example.md 为主线结合 internal/event 的源码实现完整讲解事件系统的初始化、事件发送、Handler 层接入、自定义处理器、测试与异步模式帮助你快速把可观测性能力挂接到 Chat Pipeline 的每一个关键节点。一、事件系统整体设计1.1 核心组件从 internal/event/event.go 的源码可以看出事件系统由以下几部分构成EventType事件类型的字符串枚举覆盖查询处理、检索、重排、合并、聊天生成、Agent、错误等全流程Event事件载体包含ID未指定时自动生成 UUID、Type、SessionID、Data、Metadata与RequestID字段见 internal/event/event.goEventHandlerfunc(ctx context.Context, event Event) error形式的监听函数EventBus事件发布/订阅的核心内部用sync.RWMutex保证并发安全支持同步与异步两种模式全局事件总线通过sync.Once单例模式提供GetGlobalEventBus()并封装了包级便捷函数On/Emit/EmitAndWait/HasHandlers/Clear见 internal/event/global.go。1.2 完整事件类型清单事件类型在 internal/event/event.go 中集中定义主要分为类别事件类型语义查询处理EventQueryReceived/EventQueryValidated/EventQueryPreprocess/EventQueryRewrite/EventQueryRewritten查询到达、校验、预处理、改写开始与完成检索EventRetrievalStart/EventRetrievalVector/EventRetrievalKeyword/EventRetrievalEntity/EventRetrievalComplete检索开始、向量/关键词/实体检索、检索完成重排EventRerankStart/EventRerankComplete重排开始与完成合并EventMergeStart/EventMergeComplete多路结果合并聊天生成EventChatStart/EventChatComplete/EventChatStream生成开始、完成、流式输出AgentEventAgentQuery/EventAgentPlan/EventAgentStep/EventAgentTool/EventAgentComplete以及流式反馈事件EventAgentThought/EventAgentToolCall/EventAgentToolResult/EventAgentReflection/EventAgentReferences/EventAgentFinalAnswerAgent 全生命周期与实时反馈MCP 审批/OAuthEventToolApprovalRequired/EventToolApprovalResolved/EventMCPOAuthRequired/EventMCPOAuthResolvedMCP 危险工具人工审批与对话内 OAuth 授权记忆与压缩EventMemoryRecalled/EventContextCompacted/EventUserMessageInjected长期记忆召回、上下文压缩、运行中消息注入其他EventError/EventSessionTitle/EventStop错误、会话标题更新、停止生成注意EventMCPOAuthRequired等新增事件在对话中触发 OAuth 授权时Agent 会暂停等待用户授权这是事件系统支撑交互式控制流的典型场景。1.3 事件数据结构每种事件都配套了强类型的数据结构全部定义在 internal/event/event_data.goQueryDataOriginalQuery、RewrittenQuery、SessionID、UserID、ExtraRetrievalDataQuery、KnowledgeBaseID、TopK、Threshold、RetrievalTypevector/keyword/entity、ResultCount、Results、Duration毫秒RerankDataQuery、InputCount、OutputCount、ModelID、Threshold、Results、DurationMergeDataInputCount、OutputCount、MergeTypededup/fusion等、Results、DurationChatDataQuery、ModelID、Response、StreamChunk、TokenCount、Duration、IsStreamErrorDataError、ErrorCode、Stage错误发生阶段、SessionID、QueryAgent 系列AgentPlanData、AgentStepData、AgentActionData、AgentThoughtData、AgentToolCallData、AgentToolResultData、AgentReferencesData、AgentFinalAnswerData、AgentCompleteData含TotalSteps、FinalAnswer、KnowledgeRefs、Usage、TotalDurationMs等。所有数据结构的 JSON tag 均已定义如duration_ms、knowledge_base_id便于直接序列化后投递到日志、监控或分析平台。1.4 构造与链式方法NewEvent(eventType, data)创建一个携带空Metadata的事件随后可用链式方法补充上下文event.NewEvent(event.EventRetrievalComplete, event.RetrievalData{...}). WithSessionID(sessionID). WithRequestID(requestID). WithMetadata(env, prod)WithSessionID与WithRequestID是后续日志关联、问题追踪的关键务必在发送时带上。二、在服务初始化时装配事件总线事件系统推荐在应用启动阶段完成装配。文档示例建议在cmd/server/main.go或internal/container/container.go中集中初始化见 internal/event/usage_example.mdimport github.com/Tencent/WeKnora/internal/event func InitializeEventSystem() { // 获取全局事件总线单例 bus : event.GetGlobalEventBus() // 注册内置的监控处理器 event.NewMonitoringHandler(bus) // 注册内置的分析处理器 event.NewAnalyticsHandler(bus) // 或者注册自定义处理器 bus.On(event.EventQueryReceived, func(ctx context.Context, e event.Event) error { // 自定义处理逻辑 return nil }) }从源码看WeKnora 实际采用依赖注入容器来装配事件总线internal/container/container.go 中通过container.Provide(event.NewEventBus)将EventBus注入到应用容器随后注入到各服务例如chatManage.EventBus字段。因此全局单例 依赖注入两条路径都可用需要全局跨模块访问时使用GetGlobalEventBus()需要精确控制作用域、便于测试隔离时通过构造函数注入event.NewEventBus()实例。三、在 Chat Pipeline 插件中发送事件文档给出了四个典型的插件埋点示例分别对应search.go、rewrite.go、rerank.go、chat_completion.go覆盖了 RAG 问答主链路。以下按文档逐一展开并结合真实目录说明。3.1 检索阶段search.go在检索插件中围绕performSearch发送开始/完成/错误三组事件见 internal/event/usage_example.mdfunc (p *PluginSearch) OnEvent( ctx context.Context, eventType types.EventType, chatManage *types.ChatManage, next func() *PluginError, ) *PluginError { startTime : time.Now() // 检索开始 event.Emit(ctx, event.NewEvent(event.EventRetrievalStart, event.RetrievalData{ Query: chatManage.ProcessedQuery, KnowledgeBaseID: chatManage.KnowledgeBaseID, TopK: chatManage.EmbeddingTopK, RetrievalType: vector, }).WithSessionID(chatManage.SessionID)) results, err : p.performSearch(ctx, chatManage) if err ! nil { // 错误事件携带阶段信息便于监控定位 event.Emit(ctx, event.NewEvent(event.EventError, event.ErrorData{ Error: err.Error(), Stage: retrieval, SessionID: chatManage.SessionID, Query: chatManage.ProcessedQuery, }).WithSessionID(chatManage.SessionID)) return ErrSearch.WithError(err) } // 检索完成记录结果数与耗时 event.Emit(ctx, event.NewEvent(event.EventRetrievalComplete, event.RetrievalData{ Query: chatManage.ProcessedQuery, KnowledgeBaseID: chatManage.KnowledgeBaseID, TopK: chatManage.EmbeddingTopK, RetrievalType: vector, ResultCount: len(results), Duration: time.Since(startTime).Milliseconds(), Results: results, }).WithSessionID(chatManage.SessionID)) chatManage.SearchResult results return next() }要点RetrievalType用于区分vector/keyword/entity三种检索路径Duration以毫秒为单位Stage字段如retrieval、chat_completion让错误事件能精确归类到流水线中的哪一环。3.2 查询改写rewrite.go改写插件围绕rewriteQuery发送改写开始与完成事件见 internal/event/usage_example.mdfunc (p *PluginRewriteQuery) OnEvent( ctx context.Context, eventType types.EventType, chatManage *types.ChatManage, next func() *PluginError, ) *PluginError { event.Emit(ctx, event.NewEvent(event.EventQueryRewrite, event.QueryData{ OriginalQuery: chatManage.Query, SessionID: chatManage.SessionID, }).WithSessionID(chatManage.SessionID)) rewrittenQuery, err : p.rewriteQuery(ctx, chatManage) if err ! nil { return ErrRewrite.WithError(err) } event.Emit(ctx, event.NewEvent(event.EventQueryRewritten, event.QueryData{ OriginalQuery: chatManage.Query, RewrittenQuery: rewrittenQuery, SessionID: chatManage.SessionID, }).WithSessionID(chatManage.SessionID)) chatManage.RewriteQuery rewrittenQuery return next() }EventQueryRewritten携带原查询 → 改写后查询的对照是分析改写效果召回提升率、改写失败率的核心数据源。3.3 重排阶段rerank.go重排插件记录输入/输出数量与耗时见 internal/event/usage_example.mdfunc (p *PluginRerank) OnEvent( ctx context.Context, eventType types.EventType, chatManage *types.ChatManage, next func() *PluginError, ) *PluginError { startTime : time.Now() inputCount : len(chatManage.SearchResult) event.Emit(ctx, event.NewEvent(event.EventRerankStart, event.RerankData{ Query: chatManage.ProcessedQuery, InputCount: inputCount, ModelID: chatManage.RerankModelID, }).WithSessionID(chatManage.SessionID)) rerankResults, err : p.performRerank(ctx, chatManage) if err ! nil { return ErrRerank.WithError(err) } event.Emit(ctx, event.NewEvent(event.EventRerankComplete, event.RerankData{ Query: chatManage.ProcessedQuery, InputCount: inputCount, OutputCount: len(rerankResults), ModelID: chatManage.RerankModelID, Duration: time.Since(startTime).Milliseconds(), Results: rerankResults, }).WithSessionID(chatManage.SessionID)) chatManage.RerankResult rerankResults return next() }InputCount/OutputCount的差值直接反映重排模型的压缩效率配合ModelID可以做模型维度的耗时与效果对比。3.4 生成阶段chat_completion.go聊天生成插件发送开始/完成/错误事件见 internal/event/usage_example.mdfunc (p *PluginChatCompletion) OnEvent( ctx context.Context, eventType types.EventType, chatManage *types.ChatManage, next func() *PluginError, ) *PluginError { startTime : time.Now() event.Emit(ctx, event.NewEvent(event.EventChatStart, event.ChatData{ Query: chatManage.Query, ModelID: chatManage.ChatModelID, IsStream: false, }).WithSessionID(chatManage.SessionID)) chatModel, opt, err : prepareChatModel(ctx, p.modelService, chatManage) if err ! nil { return ErrGetChatModel.WithError(err) } chatMessages : prepareMessagesWithHistory(chatManage) chatResponse, err : chatModel.Chat(ctx, chatMessages, opt) if err ! nil { event.Emit(ctx, event.NewEvent(event.EventError, event.ErrorData{ Error: err.Error(), Stage: chat_completion, SessionID: chatManage.SessionID, Query: chatManage.Query, }).WithSessionID(chatManage.SessionID)) return ErrModelCall.WithError(err) } event.Emit(ctx, event.NewEvent(event.EventChatComplete, event.ChatData{ Query: chatManage.Query, ModelID: chatManage.ChatModelID, Response: chatResponse.Content, TokenCount: chatResponse.TokenCount, Duration: time.Since(startTime).Milliseconds(), IsStream: false, }).WithSessionID(chatManage.SessionID)) chatManage.ChatResponse chatResponse return next() }TokenCount与Duration是计费核算、模型成本分析和端到端延迟拆解的基础数据。3.5 源码中的真实调用形态在上述四个示例中事件都是通过包级函数event.Emit(ctx, ...)发送到全局总线。而在实际运行链路中internal/application/service/chat_pipeline/progress.go 与 internal/application/service/chat_pipeline/chat_completion_stream.go 展示了另一种形态通过注入到chatManage.EventBus的实例发送types.Event例如流式输出时按EventAgentThought/EventAgentFinalAnswer/EventError分片发射。两种形态底层都落到EventBus.Emit区别仅在于总线实例的来源全局单例 vs 依赖注入实例你可以根据模块的解耦需求任选其一。四、在 Handler 层发送请求接收事件除了 Pipeline 内部埋点HTTP 入口层也适合发送请求到达事件用于把用户查询与后端处理关联起来。文档以internal/handler/message.go的SendMessage为例见 internal/event/usage_example.mdfunc (h *MessageHandler) SendMessage(c *gin.Context) { ctx : c.Request.Context() var req types.SendMessageRequest if err : c.ShouldBindJSON(req); err ! nil { c.JSON(400, gin.H{error: err.Error()}) return } // 查询接收事件带 SessionID 与 RequestID event.Emit(ctx, event.NewEvent(event.EventQueryReceived, event.QueryData{ OriginalQuery: req.Content, SessionID: req.SessionID, UserID: c.GetString(user_id), }).WithSessionID(req.SessionID).WithRequestID(c.GetString(request_id))) // 处理消息... }这里的关键设计是RequestID贯穿全链路Handler 入口带上request_id后后续 Pipeline 各阶段的事件检索、重排、生成都能通过同一RequestID串联实现一次用户请求 → 全链路事件回放的可观测能力。五、编写自定义处理器事件系统的价值在于监听端可以完全独立地消费这些事件。5.1 Prometheus 监控处理器文档给出了完整的监控示例见 internal/event/usage_example.md用 HistogramVec 记录检索与重排耗时并以knowledge_base_id、retrieval_type、model_id作为标签维度package monitoring import ( context github.com/Tencent/WeKnora/internal/event github.com/prometheus/client_golang/prometheus ) var ( retrievalDuration prometheus.NewHistogramVec( prometheus.HistogramOpts{ Name: retrieval_duration_milliseconds, Help: Duration of retrieval operations, }, []string{knowledge_base_id, retrieval_type}, ) rerankDuration prometheus.NewHistogramVec( prometheus.HistogramOpts{ Name: rerank_duration_milliseconds, Help: Duration of rerank operations, }, []string{model_id}, ) ) func init() { prometheus.MustRegister(retrievalDuration) prometheus.MustRegister(rerankDuration) } func SetupEventMonitoring() { bus : event.GetGlobalEventBus() // 监控检索性能 bus.On(event.EventRetrievalComplete, func(ctx context.Context, e event.Event) error { data : e.Data.(event.RetrievalData) retrievalDuration.WithLabelValues( data.KnowledgeBaseID, data.RetrievalType, ).Observe(float64(data.Duration)) return nil }) // 监控排序性能 bus.On(event.EventRerankComplete, func(ctx context.Context, e event.Event) error { data : e.Data.(event.RerankData) rerankDuration.WithLabelValues(data.ModelID).Observe(float64(data.Duration)) return nil }) }监听器内部使用类型断言e.Data.(event.RetrievalData)恢复强类型数据——这正是预定义数据结构带来的类型安全优势。5.2 结构化日志处理器借助事件系统的中间件可以统一给所有关键事件加上耗时与日志包装见 internal/event/usage_example.mdpackage logging import ( context encoding/json github.com/Tencent/WeKnora/internal/event github.com/Tencent/WeKnora/internal/logger ) func SetupEventLogging() { bus : event.GetGlobalEventBus() // 对所有事件进行结构化日志记录 logHandler : event.ApplyMiddleware( func(ctx context.Context, e event.Event) error { data, _ : json.Marshal(e.Data) logger.Infof(ctx, Event: type%s, session%s, request%s, data%s, e.Type, e.SessionID, e.RequestID, string(data)) return nil }, event.WithTiming(), ) // 注册到所有关键事件 bus.On(event.EventQueryReceived, logHandler) bus.On(event.EventQueryRewritten, logHandler) bus.On(event.EventRetrievalComplete, logHandler) bus.On(event.EventRerankComplete, logHandler) bus.On(event.EventChatComplete, logHandler) bus.On(event.EventError, logHandler) }5.3 中间件机制中间件是横切关注点的载体全部实现在 internal/event/middleware.go中间件作用WithLogging()记录事件触发、处理成功/失败日志WithTiming()统计处理耗时并写入event.Metadata[duration_ms]WithRecovery()捕获 handler panic包装为PanicError返回而非崩溃Chain(...)/ApplyMiddleware(handler, ...)组合多个中间件按声明顺序执行WithRecovery值得特别关注同步模式下如果某个监听器 panic 而未恢复会中断整个事件分发链而异步模式下EventBus.Emit内部也自带了recoverpanic 只影响单个 goroutine见 internal/event/event.go。六、完整的初始化流程与测试6.1 一键装配将以上模块组合起来即得到完整的初始化入口见 internal/event/usage_example.mdfunc Initialize() { // 1. 初始化事件系统 eventBus : event.GetGlobalEventBus() // 2. 设置监控 event.NewMonitoringHandler(eventBus) // 3. 设置分析 event.NewAnalyticsHandler(eventBus) // 4. 设置 Prometheus 监控按需启用 // monitoring.SetupEventMonitoring() // 5. 设置结构化日志按需启用 // logging.SetupEventLogging() // 6. 其他初始化... }说明event.NewMonitoringHandler与event.NewAnalyticsHandler在 internal/event/usage_example.md 和 internal/event/SUMMARY.md 中均有引用属于文档描述的内置处理器入口。若当前源码树中尚未提供这两个工厂函数可直接用 5.1、5.2 节的自定义监听器 bus.On(...)实现等价能力。6.2 使用独立总线测试测试时不要依赖全局单例而是创建独立总线实例避免污染其他测试见 internal/event/usage_example.mdfunc TestMyService(t *testing.T) { ctx : context.Background() // 创建测试专用的事件总线 testBus : event.NewEventBus() // 注册测试监听器 var receivedEvents []event.Event testBus.On(event.EventQueryReceived, func(ctx context.Context, e event.Event) error { receivedEvents append(receivedEvents, e) return nil }) // 执行测试... testBus.Emit(ctx, event.NewEvent(event.EventQueryReceived, event.QueryData{ OriginalQuery: test, })) // 验证事件 if len(receivedEvents) ! 1 { t.Errorf(Expected 1 event, got %d, len(receivedEvents)) } }internal/event/example_test.go 还提供了可运行示例与基准测试ExampleEventBus_pipeline演示了完整 RAG 流程的事件时序TestEventBus_MultipleHandlers验证同一事件多监听器全部触发TestEventBus_Async与TestEventBus_EmitAndWait验证异步与等待语义BenchmarkEventBus_Emit用于性能回归文档与 SUMMARY 记录的基准约为纳秒级发送开销。七、异步处理模式对于不影响主流程的事件监控上报、行为分析等建议使用异步总线避免阻塞用户请求见 internal/event/usage_example.mdfunc SetupAsyncAnalytics() { asyncBus : event.NewAsyncEventBus() asyncBus.On(event.EventQueryReceived, func(ctx context.Context, e event.Event) error { // 异步发送到分析平台不阻塞主流程 // sendToAnalyticsPlatform(e) return nil }) // 使用异步总线发送事件 // asyncBus.Emit(ctx, event) }从 internal/event/event.go 可以看到NewAsyncEventBus()只是把asyncMode置为true。异步模式下Emit为 fire-and-forget每个 handler 在独立 goroutine 中执行返回nil调用方立即继续goroutine 内默认recover单个 handler panic 不影响其余 handler若确实需要等待全部异步 handler 完成改用EmitAndWait它借助sync.WaitGroup error channel 汇总结果见 internal/event/event.go。对于需要精确控制错误传播的业务关键事件仍应使用默认的同步模式handler 出错时会以event handler failed for %s: %w包装返回错误。八、性能优化建议综合文档与 internal/event/event.go 的实现落地事件系统时有以下几条实践原则关键路径上避免同步阻塞监控、日志、分析等与业务无关的监听器使用异步模式NewAsyncEventBus或EmitAndWait按需等待合理使用中间件中间件有叠加开销只在需要的监听器上使用WithTiming/WithRecovery等避免全员套用控制事件数据大小RetrievalData.Results、ChatData.Response可能体积很大异步上报前建议裁剪或仅保留摘要字段event_data.go中Results、Response等字段均带有omitempty可按需省略保持监听器单一职责一个监听器只做一件事记录指标、写日志、上报分析组合多个监听器而非堆叠逻辑便于独立启停与排查利用RequestID串联全链路Handler 入口设置RequestID后所有事件天然可回放配合SessionID可实现用户会话 → 单次请求 → 各阶段三级追溯。九、总结WeKnora 事件系统提供了一条把查询处理链路变成可观测数据流的低侵入路径Pipeline 各插件只负责Emit结构化事件监听端自由组合监控、日志、分析与告警。本文覆盖了从初始化装配、四个核心插件埋点、Handler 层事件发送到自定义 Prometheus 监控、结构化日志、中间件、独立总线测试与异步模式的全过程所有代码示例均可直接对照 internal/event/usage_example.md 与 internal/event 的源码验证。接入后你可以在不修改任何业务逻辑的情况下获得检索耗时、重排压缩率、模型 Token 消耗、错误阶段分布等全套运营指标。【免费下载链接】WeKnoraOpen-source LLM knowledge platform: turn raw documents into a queryable RAG, an autonomous reasoning agent, and a self-maintaining Wiki.项目地址: https://gitcode.com/GitHub_Trending/we/WeKnora创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考