From 0b064307a01e5eba65a22b2d3df7e51054beaff7 Mon Sep 17 00:00:00 2001 From: alex Date: Fri, 10 Jul 2026 20:46:01 +0200 Subject: turn pubsub into separate package --- main.go | 10 +++--- pubsub.go | 116 -------------------------------------------------------------- 2 files changed, 4 insertions(+), 122 deletions(-) delete mode 100644 pubsub.go diff --git a/main.go b/main.go index 029bf84..9c94e6c 100644 --- a/main.go +++ b/main.go @@ -7,6 +7,7 @@ import ( "os" "redisClone/pkg/commands" "redisClone/pkg/core" + "redisClone/pkg/pubsub" "strings" "time" ) @@ -14,7 +15,7 @@ import ( var db core.RedisDB -var hub Hub +var hub pubsub.Hub type baseHandlerFunc func([]string)interface{} type connectionHandlerFunc func(net.Conn, []string) @@ -45,9 +46,6 @@ var baseCommandRegistry = map[string]baseHandlerFunc{ "LREM": handleLrem, } var connectionCommandRegistry = map[string]connectionHandlerFunc{ - "SUBSCRIBE": handleSubscribe, - "UNSUBSCRIBE": handleUnsubscribe, - "PUBLISH": handlePublish, "PING": handlePing, } @@ -61,7 +59,7 @@ func main() { } } db = core.RedisDB{Shards: shards} - hub = Hub{Channels: make(map[string]map[net.Conn]struct{}),} + hub = pubsub.Hub{Channels: make(map[string]map[net.Conn]struct{}),} go func() { ticker := time.NewTicker(1 * time.Second) @@ -98,7 +96,7 @@ func main() { func handleConnection(conn net.Conn) { defer func() { - handleDisconnect(conn) + hub.HandleDisconnect(conn) conn.Close() }() fmt.Println("Client connected:", conn.RemoteAddr()) diff --git a/pubsub.go b/pubsub.go deleted file mode 100644 index 5b6a747..0000000 --- a/pubsub.go +++ /dev/null @@ -1,116 +0,0 @@ -package main - -import ( - "errors" - "net" - "redisClone/pkg/core" - "sync" -) - -type Hub struct { - mu sync.RWMutex - Channels map[string]map[net.Conn]struct{} -} - -func handleSubscribe(conn net.Conn, args []string) { - if len(args) < 2 { - conn.Write(core.SerializeRESP(errors.New("wrong number of arguments for 'SUBSCRIBE'"))) - return - } - channels := args[1:] - hub.mu.Lock() - defer hub.mu.Unlock() - for _, channel := range channels { - if _, exists := hub.Channels[channel]; !exists { - hub.Channels[channel] = make(map[net.Conn]struct{}) - } - hub.Channels[channel][conn] = struct{}{} - - subCount := 0 - for _, subscribers := range hub.Channels { - if _, ok := subscribers[conn]; ok { - subCount++ - } - } - - conn.Write(core.SerializeRESP([]any{"subscribe", channel, subCount})) - } -} - -func handleUnsubscribe(conn net.Conn, args []string) { - hub.mu.Lock() - defer hub.mu.Unlock() - - var channels []string - if len(args) < 2 { - for channel, subscribers := range hub.Channels { - if _, ok := subscribers[conn]; ok { - channels = append(channels, channel) - } - } - } else { - channels = args[1:] - } - - for _, channel := range channels { - if subscribers, exists := hub.Channels[channel]; exists { - delete(subscribers, conn) - if len(subscribers) == 0 { - delete(hub.Channels, channel) - } - } - - subCount := 0 - for _, subscribers := range hub.Channels { - if _, ok := subscribers[conn]; ok { - subCount++ - } - } - - conn.Write(core.SerializeRESP([]any{"unsubscribe", channel, subCount})) - } -} - -func handlePublish(conn net.Conn, args []string) { - if len(args) != 3 { - conn.Write(core.SerializeRESP(errors.New("wrong number of arguments for 'PUBLISH'"))) - return - } - conn.Write(core.SerializeRESP(hub.Publish(args[1], args[2]))) -} - -func handleDisconnect(conn net.Conn) { - hub.mu.Lock() - defer hub.mu.Unlock() - for channel, subscribers := range hub.Channels { - delete(subscribers, conn) - if len(subscribers) == 0 { - delete(hub.Channels, channel) - } - } -} - -func (hub *Hub) Publish(channel, message string) int { - hub.mu.RLock() - subscribers, exists := hub.Channels[channel] - if !exists || len(subscribers) == 0 { - hub.mu.RUnlock() - return 0 - } - - conns := make([]net.Conn, 0, len(subscribers)) - for conn := range subscribers { - conns = append(conns, conn) - } - hub.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 -} -- cgit v1.2.3