ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

WPS 事件系统深度指南:在 Wave Terminal 中实现发布订阅架构

WPS 事件系统深度指南:在 Wave Terminal 中实现发布订阅架构 WPS 事件系统深度指南在 Wave Terminal 中实现发布订阅架构【免费下载链接】wavetermAn open-source, AI-integrated, cross-platform terminal for seamless workflows项目地址: https://gitcode.com/GitHub_Trending/wa/waveterm导读本文档系统讲解 Wave Terminal 内置的 WPSWave PubSub事件系统从事件结构、Broker 路由原理、事件类型注册的三段式流程到发布/订阅的 Go 与 TypeScript 双端实现再到 Scope 通配匹配、事件持久化与调试技巧。读完本文你将掌握在 Wave Terminal 后端各模块之间乃至前后端通过 WebSocket异步通信的标准姿势并能按规范新增一个类型安全的事件类型。概述什么是 WPSWPSWave PubSub是 Wave Terminal 的发布-订阅事件系统用于让应用的不同部分异步通信。其核心是一个基于Broker 模式的全局路由器发布者Publisher把事件交给 BrokerBroker 根据事件类型Event Type和作用域Scope把事件路由给匹配的订阅者Subscriber。从架构上看WPS 处于系统总线位置各后端服务块控制器、作业控制器、连接控制器、工作区服务、AI 会话等通过wps.Broker发布状态变更前端Electron 渲染进程通过 WebSocket 订阅事件实时刷新界面独立进程如 job controller 子进程也能通过 WshRouter 路由注册订阅。这种解耦让发布者无需知道订阅者的存在双方只通过事件名 作用域约定契约。核心文件文件职责pkg/wps/wpstypes.go事件类型常量、AllEvents注册表、WaveEvent/SubscriptionRequest数据结构pkg/wps/wps.goBroker 实现订阅管理、作用域匹配、事件持久化、路由投递pkg/tsgen/tsgenevent.goWaveEventDataTypes映射驱动 TypeScript 事件类型生成pkg/tsgen/tsgen.goExtraTypes额外类型清单用于前端类型生成pkg/wshrpc/wshserver/wshserver.goRPC 层的事件订阅/退订/读历史命令实现frontend/app/store/wps.ts前端事件订阅封装与分发pkg/wcore/wcore.go核心用法示例块关闭等事件发布事件结构WaveEvent所有 WPS 事件统一使用WaveEvent结构定义于 pkg/wps/wpstypes.gotype WaveEvent struct { Event string json:event // 事件类型常量 Scopes []string json:scopes,omitempty // 可选作用域用于定向投递 Sender string json:sender,omitempty // 可选发送者标识 Persist int json:persist,omitempty // 要持久化到历史中的事件条数 Data any json:data,omitempty // 事件载荷 }各字段语义Event必填事件类型字符串如waveai:ratelimit必须对应wpstypes.go中定义的常量Scopes可选目标作用域列表。订阅者可按作用域过滤实现只关心某个对象/工作区的定向接收Sender可选标识发布来源。RPC 层发布时若为空会自动填入调用方路由标识见 pkg/wshrpc/wshserver/wshserver.goPersist可选大于 0 时表示在 Broker 内存中保留最近 N 条同类型事件供迟到的订阅者补读上限为MaxPersist 4096见 pkg/wps/wps.goData可选任意类型的载荷Go 端建议使用强类型 struct。配套的SubscriptionRequest结构用于订阅type SubscriptionRequest struct { Event string json:event Scopes []string json:scopes,omitempty AllScopes bool json:allscopes,omitempty }AllScopes为true时表示接收该事件类型的所有作用域等价于不限定作用域。内置事件类型一览当前仓库中已注册的全部事件常量pkg/wps/wpstypes.go常量名事件字符串Data 类型Event_BlockCloseblockclosestring块 IDEvent_ConnChangeconnchangewshrpc.ConnStatusEvent_SysInfosysinfowshrpc.TimeSeriesDataEvent_ControllerStatuscontrollerstatus*blockcontroller.BlockControllerRuntimeStatusEvent_BuilderStatusbuilderstatuswshrpc.BuilderStatusDataEvent_BuilderOutputbuilderoutputmap[string]anyEvent_WaveObjUpdatewaveobj:updatewaveobj.WaveObjUpdateEvent_BlockFileblockfile*WSFileEventDataEvent_Configconfigwconfig.WatcherUpdateEvent_UserInputuserinput*userinput.UserInputRequestEvent_RouteDownroute:down无noneEvent_RouteUproute:up无noneEvent_WorkspaceUpdateworkspace:update无noneEvent_WaveAIRateLimitwaveai:ratelimit*uctypes.RateLimitInfoEvent_WaveAppAppGoUpdatedwaveapp:appgoupdated无noneEvent_TsunamiUpdateMetatsunami:updatemetawshrpc.AppMetaEvent_AIModeConfigwaveai:modeconfigwconfig.AIModeConfigUpdateEvent_BlockJobStatusblock:jobstatuswshrpc.BlockJobStatusDataEvent_Badgebadgebaseds.BadgeEvent注意命名规律无冒号的如blockclose属于早期事件较新的事件采用namespace:eventname形式waveai:、waveobj:、workspace:、block:、tsunami:、waveapp:、route:。添加新事件类型四段式标准流程WPS 对新增事件有严格的注册要求遗漏任何一步都会导致前端类型缺失或事件无法投递。原文档中的检查清单在 pkg/wps/wpstypes.go 的文件头注释中也有对应说明。Step 1定义事件常量在 pkg/wps/wpstypes.go 的const块中添加const ( Event_BlockClose blockclose Event_ConnChange connchange // ... 其他事件 ... Event_YourNewEvent your:newevent // type: YourEventData无数据则写 none }命名规范常量名使用带Event_前缀的描述性 PascalCase字符串值使用小写 冒号命名空间如namespace:eventname相关事件用同一命名空间前缀分组必须添加// type: TypeName注释无数据载荷时写// type: none。该注释是给开发者与代码生成器阅读的契约说明。Step 2加入 AllEvents 注册表将新常量追加到同文件的AllEvents切片中var AllEvents []string []string{ // ... 已有事件 ... Event_YourNewEvent, }AllEvents的遍历顺序直接决定前端生成的WaveEventName联合类型中事件名的排列因此追加而非乱序插入。Step 3注册到 WaveEventDataTypes必做在 pkg/tsgen/tsgenevent.go 的WaveEventDataTypes映射中添加条目。这一步不可省略它驱动事件data字段的 TypeScript 类型生成var WaveEventDataTypes map[string]reflect.Type{ // ... 已有条目 ... wps.Event_YourNewEvent: reflect.TypeOf(YourEventData{}), // 值类型 // wps.Event_YourNewEvent: reflect.TypeOf((*YourEventData)(nil)), // 指针类型 // wps.Event_YourNewEvent: nil, // 无数据type: none }三种写法对应三种情况值类型reflect.TypeOf(YourType{})指针类型reflect.TypeOf((*YourType)(nil))如Event_WaveAIRateLimit、Event_ControllerStatus均使用指针无数据直接给nil如Event_RouteUp、Event_WorkspaceUpdate。从源码看若事件未注册该映射getWaveEventDataTSType会回退返回anypkg/tsgen/tsgenevent.go即前端类型退化为不安全类型——这正是必须注册的原因。类型生成器还会为未注册事件在测试中显式断言见 pkg/tsgen/tsgenevent_test.go。Step 4定义事件数据结构可选若事件携带结构化数据为其定义强类型 struct字段使用 JSON tag 对齐前端命名type YourEventData struct { Field1 string json:field1 Field2 int json:field2 }Step 5向前端暴露类型如需如果该数据类型尚未通过任何 RPC 调用暴露给前端需要在 pkg/tsgen/tsgen.go 的ExtraTypes中添加// add extra types to generate here var ExtraTypes []any{ waveobj.ORef{}, // ... 其他类型 ... uctypes.RateLimitInfo{}, // 示例已添加 YourEventData{}, // 在此添加你的新类型 }然后运行代码生成命令task generate该命令会更新 frontend/types/gotypes.d.ts为你的类型生成 TypeScript 定义。生成的WaveEvent是判别联合类型——event字段作为判别器data字段按事件名自动收紧为对应类型生成逻辑见 pkg/tsgen/tsgenevent.go前端处理事件时能获得完整的类型安全。发布事件基础发布使用全局 Broker 发布import github.com/wavetermdev/waveterm/pkg/wps wps.Broker.Publish(wps.WaveEvent{ Event: wps.Event_YourNewEvent, Data: yourData, })Publish的完整执行链路pkg/wps/wps.go若Persist 0先执行persistEvent写入内存历史通过GetClient()获取当前注册的投递客户端默认为全局WshRouter在 pkg/wshrpc/wshclient/barerpcclient.go 中通过wps.Broker.SetClient(wshutil.DefaultRouter)挂接计算匹配的路由 ID 集合逐个调用client.SendEvent(routeId, event)。带作用域的定向发布wps.Broker.Publish(wps.WaveEvent{ Event: wps.Event_WaveObjUpdate, Scopes: []string{oref.String()}, // 定向到某个具体对象 Data: updateData, })作用域通常使用对象的 ORef 字符串如waveobj.MakeORef(...).String()以实现只通知关心该对象的订阅者。真实示例见 pkg/wcore/block.go 中的sendBlockCloseEvent它以块 ORef 为作用域发布Event_BlockClose。在 goroutine 中异步发布为避免阻塞调用方可将发布放入 goroutinego func() { wps.Broker.Publish(wps.WaveEvent{ Event: wps.Event_YourNewEvent, Data: data, }) }()何时使用 goroutine在性能关键代码路径上发布时事件只是信息性通知、不要求立即投递时在持有锁的代码中发布时防止死锁。真实范例updateRateLimit在持有rateLimitLock期间通过 goroutine 发布见下文完整示例。事件持久化Persistwps.Broker.Publish(wps.WaveEvent{ Event: wps.Event_YourNewEvent, Persist: 100, // 保留最近 100 条 Data: data, })持久化实现细节pkg/wps/wps.goPersist 0时直接跳过持久化请求值超过MaxPersist4096时会被截断为 4096事件会同时按每个 scope 单独一份空 scope全局一份存入PersistMap采用环形截断超过numPersist时丢弃最旧事件保留最近 N 条。订阅者可用wps.Broker.ReadEventHistory(eventType, scope, maxItems)补读历史pkg/wps/wps.go注意该 API 不支持通配符取表示全局历史。设计建议原文档强调持久化通常不启用。Wave Terminal 的惯例是先通过实时 RPC 拉取当前值再订阅增量更新只有状态类事件如controllerstatus这类迟到订阅者必须拿到最新值的场景才配合Persist: 1使用。订阅事件Go 端订阅// 订阅某类型全部事件 wps.Broker.Subscribe(routeId, wps.SubscriptionRequest{ Event: wps.Event_YourNewEvent, AllScopes: true, }) // 订阅特定作用域 wps.Broker.Subscribe(routeId, wps.SubscriptionRequest{ Event: wps.Event_WaveObjUpdate, Scopes: []string{workspace:123}, }) // 退订 wps.Broker.Unsubscribe(routeId, wps.Event_YourNewEvent)订阅语义要点pkg/wps/wps.go同一路由对同一事件重复订阅属于重订阅旧订阅先被移除再按新请求替换AllScopes: true的订阅者进入AllSubs列表接收该类型全部事件否则按 scope 精确匹配进入ScopeSubs含*/**的进入StarSubs另有UnsubscribeAll(routeId)一次退订该路由的全部订阅pkg/wps/wps.go内部订阅表是map[string]*BrokerSubscription事件类型 → 订阅组BrokerSubscription内部分为全量、精确 scope、通配 scope 三组pkg/wps/wps.go。进程内订阅的真实用法在 Go 进程内先建立 bare RPC client 再订阅并同时挂接本地事件监听器rpcClient : wshclient.GetBareRpcClient() rpcClient.EventListener.On(wps.Event_RouteUp, handleRouteUpEvent) rpcClient.EventListener.On(wps.Event_BlockClose, handleBlockCloseEvent) wshclient.EventSubCommand(rpcClient, wps.SubscriptionRequest{ Event: wps.Event_RouteUp, AllScopes: true, }, nil) wshclient.EventSubCommand(rpcClient, wps.SubscriptionRequest{ Event: wps.Event_BlockClose, AllScopes: true, }, nil)摘录自 pkg/jobcontroller/jobcontroller.go 的InitJobController订阅了route:up、route:down、connchange、blockclose四类事件。RPC 层订阅命令跨进程/前端订阅走 RPC 命令实现在 pkg/wshrpc/wshserver/wshserver.goeventsubEventSubCommand以调用方路由 ID 作为订阅者身份调用Broker.SubscribeeventunsubEventUnsubCommand按事件名退订eventunsuballEventUnsubAllCommand退订该路由全部事件eventreadhistoryEventReadHistoryCommand读取持久化历史。对应客户端封装见 pkg/wshrpc/wshclient/wshclient.go。Scope 匹配与通配符Scope 支持通配符匹配*匹配单个scope 段以:分隔的一段**匹配多个scope 段仅允许出现在模式末尾。// 订阅所有 workspace 事件 wps.Broker.Subscribe(routeId, wps.SubscriptionRequest{ Event: wps.Event_WaveObjUpdate, Scopes: []string{workspace:*}, })匹配算法的底层实现是 pkg/util/utilfn/utilfn.go 的StarMatchString把模式和目标字符串按分隔符切段后逐段比对*通配任意单段**若出现在最后一段则直接匹配剩余全部。在getMatchingRouteIds中pkg/wps/wps.go发布事件时会遍历StarSubs里所有通配模式与事件的每个 scope 做匹配命中即投递精确 scope 订阅则直接查表。所有匹配到的 routeId 先去重再统一投递保证一个订阅者即使同时命中多种匹配方式也只会收到一次。前端TypeScript订阅前端通过 WebSocket 连接订阅事件底层在 frontend/app/store/wps.tsimport { waveEventSubscribeSingle } from /app/store/wps; // 订阅率更新对应 waveai:ratelimit const unsubscribe waveEventSubscribeSingle({ eventType: waveai:ratelimit, scope: undefined, // 不限定 scope 时前端会自动升级为 allscopes handler: (event) { // event.data 已被推断为 RateLimitInfo | null console.log(rate limit, event.data); }, }); // 组件卸载时退订 unsubscribe();前端实现的关键机制订阅聚合updateWaveEventSub会汇总该事件类型下所有本地订阅者若存在不限 scope 的订阅则发送allscopes: true否则合并所有 scope 列表一次性发送EventSubCommandfrontend/app/store/wps.ts重连补偿wpsReconnectHandler在 WebSocket 重连后对所有已注册事件类型重新发送订阅命令避免连接中断丢失订阅状态frontend/app/store/wps.ts本地分发handleWaveEvent收到事件后按订阅者的 scope 逐一分发不匹配则跳过frontend/app/store/wps.ts引用计数退订waveEventUnsubscribe在移除最后一个订阅者后才真正发送退订命令并清理本地 subjectfrontend/app/store/wps.ts。完整示例AI 限流信息发布下面以waveai:ratelimit事件为完整闭环示例展示从事件定义到前后端联动的全部环节。1. 定义事件类型pkg/wps/wpstypes.goconst ( // ... 其他事件 ... Event_WaveAIRateLimit waveai:ratelimit // type: *uctypes.RateLimitInfo )2. 发布事件pkg/aiusechat/usechat.go 中的真实实现——注意它在持锁期间通过 goroutine 发布以避免死锁func updateRateLimit(info *uctypes.RateLimitInfo) { if info nil { return } rateLimitLock.Lock() defer rateLimitLock.Unlock() globalRateLimitInfo info // 在 goroutine 中发布避免阻塞调用方 go func() { wps.Broker.Publish(wps.WaveEvent{ Event: wps.Event_WaveAIRateLimit, Data: info, // *uctypes.RateLimitInfo }) }() }调用链AI 对话每步完成backend.RunChatStep后返回rateLimitInfo随即调用updateRateLimitpkg/aiusechat/usechat.go。下游消费者如shouldUsePremium则通过GetGlobalRateLimit()读取当前快照pkg/aiusechat/usechat.go——这正是原文档实时 RPC 拉快照 事件订阅增量惯例的落地。3. 前端订阅前端无需手工构造 WebSocket 订阅消息使用封装后的waveEventSubscribeSinglewaveEventSubscribeSingle({ eventType: waveai:ratelimit, handler: (event) { // event.data 类型为 RateLimitInfo | null updateRateLimitDisplay(event.data); }, });eventType的类型是WaveEventName判别联合由 pkg/tsgen/tsgenevent.go 从AllEvents生成写错事件名会在编译期报错。常见事件模式状态更新配合 Persist适合迟到订阅者必须拿到最新状态的场景只保留最近一条wps.Broker.Publish(wps.WaveEvent{ Event: wps.Event_ControllerStatus, Scopes: []string{blockId}, Persist: 1, // 仅保留最新状态 Data: statusData, })对象更新对象变更统一走waveobj:updatescope 为对象 ORefdata 为WaveObjUpdate包含更新类型、对象类型、OID 与新对象体wps.Broker.Publish(wps.WaveEvent{ Event: wps.Event_WaveObjUpdate, Scopes: []string{oref.String()}, Data: waveobj.WaveObjUpdate{ UpdateType: waveobj.UpdateType_Update, OType: obj.GetOType(), OID: waveobj.GetOID(obj), Obj: obj, }, })批量更新多个对象一起变更时使用 Broker 内置的SendUpdateEvents助手批量发布pkg/wps/wps.go// 辅助函数为批量更新逐个发布事件 func (b *BrokerType) SendUpdateEvents(updates waveobj.UpdatesRtnType) { for _, update : range updates { b.Publish(WaveEvent{ Event: Event_WaveObjUpdate, Scopes: []string{waveobj.MakeORef(update.OType, update.OID).String()}, Data: update, }) } }SendUpdateEvents在工作区服务pkg/service/workspaceservice/workspaceservice.go与窗口服务pkg/service/windowservice/windowservice.go等对象写入路径上被大量使用——对象落库后立即广播前端据此刷新界面。调试技巧检查订阅表直接查看wps.Broker.SubMapmap[event]→{AllSubs, ScopeSubs, StarSubs}确认订阅是否注册成功、作用域是否合理补读持久化事件调用wps.Broker.ReadEventHistory(eventType, scope, maxItems)查看 Broker 内存中保留的历史事件打点日志Publish/Subscribe/Unsubscribe方法内都有被注释掉的log.Printf见 pkg/wps/wps.go 与 pkg/wps/wps.go调试时取消注释即可观察事件流与订阅/退订动作抓 WebSocket 流量前端事件通过 WebSocket 以eventrecv命令下发投递实现见 pkg/wshutil/wshrouter.go将事件序列化为RpcMessage发送到目标链接在浏览器 DevTools 的 Network → WS 面板中可直接观察eventrecv消息内容类型生成自检运行pkg/tsgen/tsgenevent_test.go中的TestGenerateWaveEventTypes它会断言已知事件的 TS 类型声明如route:up的data?: null、未注册事件的any回退等行为新增事件后跑一遍即可验证注册完整性。最佳实践清单使用命名空间事件名统一加前缀waveai:、workspace:、block:等避免全局命名冲突也便于检索分组不要阻塞在性能关键路径或持锁状态下发布时用 goroutine 包装Publish类型安全的数据为事件数据定义 struct 类型而非裸 map并注册到WaveEventDataTypes让前端获得完整类型推断谨慎使用作用域合理设置 scope 以限制投递范围、减少无关订阅者的无谓处理为事件写注释注明事件的触发时机、携带的数据结构// type: TypeName注释务必与WaveEventDataTypes保持一致按需持久化Persist仅用于迟到订阅者需要的状态类事件通常做法是先 RPC 拉当前值、再订阅增量更新。快速参考新增事件检查清单按原文档整理新增事件时逐项核对在 pkg/wps/wpstypes.go 添加事件常量并写// type: TypeName注释无数据用none将该常量加入同文件的AllEvents切片必做在 pkg/tsgen/tsgenevent.go 的WaveEventDataTypes注册映射——无数据事件写nil按需定义事件数据结构struct JSON tag若该类型未通过 RPC 暴露加入 pkg/tsgen/tsgen.go 的ExtraTypes运行task generate刷新 frontend/types/gotypes.d.ts使用wps.Broker.Publish()发布事件合适时用 goroutine 异步发布在相关组件中订阅事件Go 端wps.Broker.Subscribe/ 前端waveEventSubscribeSingle。遵循这套流程你就能以类型安全的方式在 Wave Terminal 的任意两个模块之间建立松耦合的异步通信通道。【免费下载链接】wavetermAn open-source, AI-integrated, cross-platform terminal for seamless workflows项目地址: https://gitcode.com/GitHub_Trending/wa/waveterm创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表