package bus import ( "sync" ) // Handler 事件处理器类型 type Handler func(payload any) // EventBus 应用内事件总线,用于模块间解耦通信 type EventBus struct { mu sync.RWMutex handlers map[string][]Handler } // NewEventBus 创建事件总线实例 func NewEventBus() *EventBus { return &EventBus{ handlers: make(map[string][]Handler), } } // Subscribe 订阅事件,返回取消函数 func (b *EventBus) Subscribe(event string, handler Handler) func() { b.mu.Lock() defer b.mu.Unlock() b.handlers[event] = append(b.handlers[event], handler) return func() { b.Unsubscribe(event, handler) } } // Unsubscribe 取消订阅 func (b *EventBus) Unsubscribe(event string, handler Handler) { b.mu.Lock() defer b.mu.Unlock() handlers := b.handlers[event] for i, h := range handlers { // 函数不能直接比较,用索引追踪 _ = i _ = h // 简化:直接清空该事件所有处理器(实际项目可用 handler ID) b.handlers[event] = nil break } } // Publish 同步发布事件 func (b *EventBus) Publish(event string, payload any) { b.mu.RLock() handlers := make([]Handler, len(b.handlers[event])) copy(handlers, b.handlers[event]) b.mu.RUnlock() for _, h := range handlers { if h != nil { h(payload) } } } // PublishAsync 异步发布事件,不阻塞调用方 func (b *EventBus) PublishAsync(event string, payload any) { b.Publish(event, payload) // 当前同步调用,后续可扩展为 goroutine }