diff --git a/server/message_store.go b/server/message_store.go index 8ebf768..fb6aefb 100644 --- a/server/message_store.go +++ b/server/message_store.go @@ -7,17 +7,17 @@ import ( type MemoryStore struct { mu sync.Mutex - msgs map[int]MessageToSend + msgs map[int]message offset int } func NewMemoryStore() *MemoryStore { return &MemoryStore{ - msgs: make(map[int]MessageToSend), + msgs: make(map[int]message), } } -func (m *MemoryStore) Write(msg MessageToSend) error { +func (m *MemoryStore) Write(msg message) error { m.mu.Lock() defer m.mu.Unlock() @@ -28,7 +28,7 @@ func (m *MemoryStore) Write(msg MessageToSend) error { return nil } -func (m *MemoryStore) ReadFrom(offset int, handleFunc func(msg MessageToSend)) error { +func (m *MemoryStore) ReadFrom(offset int, handleFunc func(msg message)) error { if offset < 0 || offset > m.offset { return fmt.Errorf("invalid offset provided") } diff --git a/server/server.go b/server/server.go index 156e37e..dca6a58 100644 --- a/server/server.go +++ b/server/server.go @@ -302,11 +302,6 @@ func (s *Server) handleUnsubscribe(peer *peer.Peer) { _ = peer.RunConnOperation(op) } -type MessageToSend struct { - topic string - data []byte -} - func (s *Server) handlePublish(peer *peer.Peer) { slog.Info("handling publisher", "peer", peer.Addr()) for { @@ -357,17 +352,14 @@ func (s *Server) handlePublish(peer *peer.Peer) { return nil } - message := MessageToSend{ - topic: topicStr, - data: dataBuf, - } - - topic := s.getTopic(message.topic) + topic := s.getTopic(topicStr) if topic == nil { - topic = newTopic(message.topic) - s.topics[message.topic] = topic + topic = newTopic(topicStr) + s.topics[topicStr] = topic } + message := newMessage(dataBuf) + err = topic.sendMessageToSubscribers(message) if err != nil { slog.Error("failed to send message to subscribers", "error", err, "peer", peer.Addr()) diff --git a/server/subscriber.go b/server/subscriber.go index ef14895..415cb8a 100644 --- a/server/subscriber.go +++ b/server/subscriber.go @@ -44,8 +44,8 @@ func newSubscriber(peer *peer.Peer, topic *topic, ackDelay, ackTimeout time.Dura offset := startAt go func() { - err := topic.messageStore.ReadFrom(offset, func(msg MessageToSend) { - s.messages <- newMessage(msg.data) + err := topic.messageStore.ReadFrom(offset, func(msg message) { + s.messages <- msg }) if err != nil { slog.Error("failed to replay messages from offset", "error", err, "offset", offset) diff --git a/server/topic.go b/server/topic.go index 9d631a5..1fb5ebf 100644 --- a/server/topic.go +++ b/server/topic.go @@ -7,8 +7,8 @@ import ( ) type Store interface { - Write(msg MessageToSend) error - ReadFrom(offset int, handleFunc func(msg MessageToSend)) error + Write(msg message) error + ReadFrom(offset int, handleFunc func(msg message)) error } type topic struct { @@ -27,7 +27,7 @@ func newTopic(name string) *topic { } } -func (t *topic) sendMessageToSubscribers(msg MessageToSend) error { +func (t *topic) sendMessageToSubscribers(msg message) error { err := t.messageStore.Write(msg) if err != nil { return fmt.Errorf("failed to write message to store: %w", err) @@ -38,7 +38,7 @@ func (t *topic) sendMessageToSubscribers(msg MessageToSend) error { t.mu.Unlock() for _, subscriber := range subscribers { - subscriber.addMessage(newMessage(msg.data), 0) + subscriber.addMessage(msg, 0) } return nil