From 2d9f92db91a189386639ac4d0d28b239e5288416 Mon Sep 17 00:00:00 2001 From: Will Date: Thu, 8 Feb 2024 17:26:43 +0000 Subject: [PATCH] some tweaks and typos --- server/message_store.go | 19 ++++++++----------- server/server.go | 26 ++++---------------------- server/subscriber.go | 20 +++++++------------- server/topic.go | 2 +- 4 files changed, 20 insertions(+), 47 deletions(-) diff --git a/server/message_store.go b/server/message_store.go index a9439ed..ed0730c 100644 --- a/server/message_store.go +++ b/server/message_store.go @@ -1,15 +1,14 @@ package server import ( - "fmt" "sync" ) // Memory store allows messages to be stored in memory type MemoryStore struct { - mu sync.Mutex - msgs map[int]message - offset int + mu sync.Mutex + msgs map[int]message + nextOffset int } // New memory store initializes a new in memory store @@ -24,17 +23,17 @@ func (m *MemoryStore) Write(msg message) error { m.mu.Lock() defer m.mu.Unlock() - m.msgs[m.offset] = msg + m.msgs[m.nextOffset] = msg - m.offset++ + m.nextOffset++ return nil } // ReadFrom will read messages from (and including) the provided offset and pass them to the provided handler -func (m *MemoryStore) ReadFrom(offset int, handleFunc func(msg message)) error { - if offset < 0 || offset > m.offset { - return fmt.Errorf("invalid offset provided") +func (m *MemoryStore) ReadFrom(offset int, handleFunc func(msg message)) { + if offset < 0 || offset >= m.nextOffset { + return } m.mu.Lock() @@ -43,6 +42,4 @@ func (m *MemoryStore) ReadFrom(offset int, handleFunc func(msg message)) error { for i := offset; i < len(m.msgs); i++ { handleFunc(m.msgs[i]) } - - return nil } diff --git a/server/server.go b/server/server.go index dca6a58..fc7ea9f 100644 --- a/server/server.go +++ b/server/server.go @@ -27,23 +27,6 @@ const ( Nack Action = 5 ) -func (a Action) String() string { - switch a { - case Subscribe: - return "subscribe" - case Unsubscribe: - return "unsubscribe" - case Publish: - return "publish" - case Ack: - return "ack" - case Nack: - return "nack" - } - - return "" -} - // Status represents the status of a request type Status uint16 @@ -70,9 +53,9 @@ func (s Status) String() string { type StartAtType uint16 const ( - Begining StartAtType = 0 - Current StartAtType = 1 - From StartAtType = 2 + Beginning StartAtType = 0 + Current StartAtType = 1 + From StartAtType = 2 ) // Server accepts subscribe and publish connections and passes messages around @@ -234,7 +217,6 @@ func (s *Server) subscribePeerToTopic(peer *peer.Peer) { var startAt int switch startAtType { case From: - // read the from var s uint16 err = binary.Read(conn, binary.BigEndian, &s) if err != nil { @@ -243,7 +225,7 @@ func (s *Server) subscribePeerToTopic(peer *peer.Peer) { return nil } startAt = int(s) - case Begining: + case Beginning: startAt = 0 case Current: startAt = -1 diff --git a/server/subscriber.go b/server/subscriber.go index 288c094..15e12e5 100644 --- a/server/subscriber.go +++ b/server/subscriber.go @@ -41,19 +41,15 @@ func newSubscriber(peer *peer.Peer, topic *topic, ackDelay, ackTimeout time.Dura go s.sendMessages() - offset := startAt - go func() { - if startAt < 0 { - return - } - - err := topic.messageStore.ReadFrom(offset, func(msg message) { - s.messages <- msg + topic.messageStore.ReadFrom(startAt, func(msg message) { + select { + case s.messages <- msg: + return + case <-s.unsubscribeCh: + return + } }) - if err != nil { - slog.Error("failed to replay messages from offset", "error", err, "offset", offset) - } }() return s @@ -94,9 +90,7 @@ func (s *subscriber) addMessage(msg message, delay time.Duration) { case <-s.unsubscribeCh: return case <-timer.C: - fmt.Printf("waiting to put message on queue: %s\n", msg.data) s.messages <- msg - fmt.Printf("put message on queue: %s\n", msg.data) } }() } diff --git a/server/topic.go b/server/topic.go index 1fb5ebf..02987e4 100644 --- a/server/topic.go +++ b/server/topic.go @@ -8,7 +8,7 @@ import ( type Store interface { Write(msg message) error - ReadFrom(offset int, handleFunc func(msg message)) error + ReadFrom(offset int, handleFunc func(msg message)) } type topic struct { -- 2.51.2