package pubsub import ( "net" "redisClone/pkg/core" "sync" ) type Hub struct { mu sync.RWMutex Channels map[string]map[net.Conn]struct{} } func NewHub() *Hub { return &Hub{ Channels: make(map[string]map[net.Conn]struct{}), } } func (h *Hub) Subscribe(conn net.Conn, channel string) int { h.mu.Lock() defer h.mu.Unlock() if _, exists := h.Channels[channel]; !exists { h.Channels[channel] = make(map[net.Conn]struct{}) } h.Channels[channel][conn] = struct{}{} subCount := 0 for _, subscribers := range h.Channels { if _, ok := subscribers[conn]; ok { subCount++ } } return subCount } func (h *Hub) Unsubscribe(conn net.Conn, channels []string) map[string]int { h.mu.Lock() defer h.mu.Unlock() if len(channels) == 0 { for channel, subscribers := range h.Channels { if _, ok := subscribers[conn]; ok { channels = append(channels, channel) } } } results := make(map[string]int) for _, channel := range channels { if subscribers, exists := h.Channels[channel]; exists { delete(subscribers, conn) if len(subscribers) == 0 { delete(h.Channels, channel) } } count := 0 for _, subs := range h.Channels { if _, ok := subs[conn]; ok { count++ } } results[channel] = count } return results } func (h *Hub) Publish(channel, message string) int { h.mu.RLock() subscribers, exists := h.Channels[channel] if !exists || len(subscribers) == 0 { h.mu.RUnlock() return 0 } conns := make([]net.Conn, 0, len(subscribers)) for conn := range subscribers { conns = append(conns, conn) } h.mu.RUnlock() payload := core.SerializeRESP([]string{"message", channel, message}) count := 0 for _, conn := range conns { _, err := conn.Write(payload) if err == nil { count++ } } return count } func (h *Hub) HandleDisconnect(conn net.Conn) { h.mu.Lock() defer h.mu.Unlock() for channel, subscribers := range h.Channels { delete(subscribers, conn) if len(subscribers) == 0 { delete(h.Channels, channel) } } }