// Package hub 是一个极简的进程内发布/订阅:按主题(文章 ID)广播 // 「有变更」信号。订阅者只拿到一个空 struct 信号——不携带内容, // 消费方(SSE 前端)收到后自己 refetch,未审核的评论内容不会经由 // 这条通道外泄。 package hub import "sync" type Hub struct { mu sync.Mutex subs map[int64]map[chan struct{}]struct{} } func New() *Hub { return &Hub{subs: make(map[int64]map[chan struct{}]struct{})} } // Subscribe 订阅某主题;返回信号 channel 和退订函数。 // channel 容量 1:订阅者处理不过来时新信号直接丢弃(合并刷新)。 func (h *Hub) Subscribe(topic int64) (<-chan struct{}, func()) { // nil hub(测试/未装配场景)静默降级:永远收不到信号的通道 if h == nil { ch := make(chan struct{}) return ch, func() {} } ch := make(chan struct{}, 1) h.mu.Lock() if h.subs[topic] == nil { h.subs[topic] = make(map[chan struct{}]struct{}) } h.subs[topic][ch] = struct{}{} h.mu.Unlock() off := func() { h.mu.Lock() delete(h.subs[topic], ch) if len(h.subs[topic]) == 0 { delete(h.subs, topic) } h.mu.Unlock() } return ch, off } // Broadcast 唤醒某主题的全部订阅者;积压的订阅者不阻塞。 func (h *Hub) Broadcast(topic int64) { if h == nil { return } h.mu.Lock() defer h.mu.Unlock() for ch := range h.subs[topic] { select { case ch <- struct{}{}: default: } } }