From 8e8e0e1cfa6aa6692120f2c7dd7cd842cee63f4f Mon Sep 17 00:00:00 2001 From: Will Andrews Date: Thu, 8 Feb 2024 18:55:06 +0000 Subject: [PATCH] Refactor app structure + test improvements (#10) --- {pubsub => client}/message.go | 2 +- {pubsub => client}/publisher.go | 8 +-- {pubsub => client}/subscriber.go | 12 ++-- {pubsub => client}/subscriber_test.go | 10 +-- example/main.go | 10 +-- example/server/main.go | 2 +- .../messagestore/memory_store.go | 16 ++--- internal/messge.go | 12 ++++ {server/peer => internal/server}/peer.go | 4 +- {server => internal/server}/server.go | 63 +++++++++++-------- {server => internal/server}/server_test.go | 21 ++++--- {server => internal/server}/subscriber.go | 40 +++++------- {server => internal/server}/topic.go | 11 ++-- 13 files changed, 116 insertions(+), 95 deletions(-) rename {pubsub => client}/message.go (96%) rename {pubsub => client}/publisher.go (90%) rename {pubsub => client}/subscriber.go (96%) rename {pubsub => client}/subscriber_test.go (96%) rename server/message_store.go => internal/messagestore/memory_store.go (59%) create mode 100644 internal/messge.go rename {server/peer => internal/server}/peer.go (93%) rename {server => internal/server}/server.go (88%) rename {server => internal/server}/server_test.go (98%) rename {server => internal/server}/subscriber.go (70%) rename {server => internal/server}/topic.go (66%) diff --git a/pubsub/message.go b/client/message.go similarity index 96% rename from pubsub/message.go rename to client/message.go index 3014e14..9e13ada 100644 --- a/pubsub/message.go +++ b/client/message.go @@ -1,4 +1,4 @@ -package pubsub +package client // Message represents a message that can be published or consumed type Message struct { diff --git a/pubsub/publisher.go b/client/publisher.go similarity index 90% rename from pubsub/publisher.go rename to client/publisher.go index 4ba6614..25919af 100644 --- a/pubsub/publisher.go +++ b/client/publisher.go @@ -1,4 +1,4 @@ -package pubsub +package client import ( "encoding/binary" @@ -6,7 +6,7 @@ import ( "net" "sync" - "github.com/willdot/messagebroker/server" + "github.com/willdot/messagebroker/internal/server" ) // Publisher allows messages to be published to a server @@ -44,8 +44,8 @@ func (p *Publisher) PublishMessage(message *Message) error { // send topic first topic := fmt.Sprintf("topic:%s", message.Topic) - topicLenB := make([]byte, 4) - binary.BigEndian.PutUint32(topicLenB, uint32(len(topic))) + topicLenB := make([]byte, 2) + binary.BigEndian.PutUint16(topicLenB, uint16(len(topic))) headers := append(topicLenB, []byte(topic)...) diff --git a/pubsub/subscriber.go b/client/subscriber.go similarity index 96% rename from pubsub/subscriber.go rename to client/subscriber.go index 16ab920..1177dd8 100644 --- a/pubsub/subscriber.go +++ b/client/subscriber.go @@ -1,4 +1,4 @@ -package pubsub +package client import ( "context" @@ -10,7 +10,7 @@ import ( "sync" "time" - "github.com/willdot/messagebroker/server" + "github.com/willdot/messagebroker/internal/server" ) type connOpp func(conn net.Conn) error @@ -96,7 +96,7 @@ func subscribeToTopics(conn net.Conn, topicNames []string, startAtType server.St return nil } - var dataLen uint32 + var dataLen uint16 err = binary.Read(conn, binary.BigEndian, &dataLen) if err != nil { return fmt.Errorf("received status %s:", resp) @@ -140,7 +140,7 @@ func unsubscribeToTopics(conn net.Conn, topicNames []string) error { return nil } - var dataLen uint32 + var dataLen uint16 err = binary.Read(conn, binary.BigEndian, &dataLen) if err != nil { return fmt.Errorf("received status %s:", resp) @@ -198,12 +198,12 @@ func (s *Subscriber) consume(ctx context.Context, consumer *Consumer) { func (s *Subscriber) readMessage(ctx context.Context, msgChan chan *Message) error { op := func(conn net.Conn) error { - err := s.conn.SetReadDeadline(time.Now().Add(time.Second)) + err := s.conn.SetReadDeadline(time.Now().Add(time.Millisecond * 300)) if err != nil { return err } - var topicLen uint64 + var topicLen uint16 err = binary.Read(s.conn, binary.BigEndian, &topicLen) if err != nil { // TODO: check if this is needed elsewhere. I'm not sure where the read deadline resets.... diff --git a/pubsub/subscriber_test.go b/client/subscriber_test.go similarity index 96% rename from pubsub/subscriber_test.go rename to client/subscriber_test.go index eff376e..0e8bcca 100644 --- a/pubsub/subscriber_test.go +++ b/client/subscriber_test.go @@ -1,4 +1,4 @@ -package pubsub +package client import ( "context" @@ -8,8 +8,7 @@ import ( "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" - - "github.com/willdot/messagebroker/server" + "github.com/willdot/messagebroker/internal/server" ) const ( @@ -134,7 +133,8 @@ func TestUnsubscribesFromTopic(t *testing.T) { err = publisher.PublishMessage(msg) require.NoError(t, err) - time.Sleep(time.Second) + // give the consumer some time to read the messages -- TODO: make better! + time.Sleep(time.Millisecond * 300) cancel() select { @@ -181,7 +181,7 @@ func TestPublishAndSubscribe(t *testing.T) { } // give the consumer some time to read the messages -- TODO: make better! - time.Sleep(time.Second) + time.Sleep(time.Millisecond * 300) cancel() select { diff --git a/example/main.go b/example/main.go index 978e4a6..d9d8e5c 100644 --- a/example/main.go +++ b/example/main.go @@ -7,8 +7,8 @@ import ( "log/slog" "time" - "github.com/willdot/messagebroker/pubsub" - "github.com/willdot/messagebroker/server" + "github.com/willdot/messagebroker/client" + "github.com/willdot/messagebroker/internal/server" ) var consumeOnly *bool @@ -23,7 +23,7 @@ func main() { go sendMessages() } - sub, err := pubsub.NewSubscriber(":3000") + sub, err := client.NewSubscriber(":3000") if err != nil { panic(err) } @@ -56,7 +56,7 @@ func main() { } func sendMessages() { - publisher, err := pubsub.NewPublisher("localhost:3000") + publisher, err := client.NewPublisher("localhost:3000") if err != nil { panic(err) } @@ -69,7 +69,7 @@ func sendMessages() { i := 0 for { i++ - msg := pubsub.NewMessage("topic a", []byte(fmt.Sprintf("message %d", i))) + msg := client.NewMessage("topic a", []byte(fmt.Sprintf("message %d", i))) err = publisher.PublishMessage(msg) if err != nil { diff --git a/example/server/main.go b/example/server/main.go index 3841b19..ea31dff 100644 --- a/example/server/main.go +++ b/example/server/main.go @@ -7,7 +7,7 @@ import ( "syscall" "time" - "github.com/willdot/messagebroker/server" + "github.com/willdot/messagebroker/internal/server" ) func main() { diff --git a/server/message_store.go b/internal/messagestore/memory_store.go similarity index 59% rename from server/message_store.go rename to internal/messagestore/memory_store.go index ed0730c..ac9bd45 100644 --- a/server/message_store.go +++ b/internal/messagestore/memory_store.go @@ -1,25 +1,27 @@ -package server +package messagestore import ( "sync" + + "github.com/willdot/messagebroker/internal" ) -// Memory store allows messages to be stored in memory +// MemoryStore allows messages to be stored in memory type MemoryStore struct { mu sync.Mutex - msgs map[int]message + msgs map[int]internal.Message nextOffset int } -// New memory store initializes a new in memory store +// NewMemoryStore initializes a new in memory store func NewMemoryStore() *MemoryStore { return &MemoryStore{ - msgs: make(map[int]message), + msgs: make(map[int]internal.Message), } } // Write will write the provided message to the in memory store -func (m *MemoryStore) Write(msg message) error { +func (m *MemoryStore) Write(msg internal.Message) error { m.mu.Lock() defer m.mu.Unlock() @@ -31,7 +33,7 @@ func (m *MemoryStore) Write(msg message) error { } // 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)) { +func (m *MemoryStore) ReadFrom(offset int, handleFunc func(msg internal.Message)) { if offset < 0 || offset >= m.nextOffset { return } diff --git a/internal/messge.go b/internal/messge.go new file mode 100644 index 0000000..240dbf6 --- /dev/null +++ b/internal/messge.go @@ -0,0 +1,12 @@ +package internal + +// Message represents a message that can be sent / received +type Message struct { + Data []byte + DeliveryCount int +} + +// NewMessage intializes a new message +func NewMessage(data []byte) Message { + return Message{Data: data, DeliveryCount: 1} +} diff --git a/server/peer/peer.go b/internal/server/peer.go similarity index 93% rename from server/peer/peer.go rename to internal/server/peer.go index 5dbd3a9..951d1ea 100644 --- a/server/peer/peer.go +++ b/internal/server/peer.go @@ -1,4 +1,4 @@ -package peer +package server import ( "net" @@ -12,7 +12,7 @@ type Peer struct { } // New returns a new peer. -func New(conn net.Conn) *Peer { +func NewPeer(conn net.Conn) *Peer { return &Peer{ conn: conn, } diff --git a/server/server.go b/internal/server/server.go similarity index 88% rename from server/server.go rename to internal/server/server.go index fc7ea9f..7999c99 100644 --- a/server/server.go +++ b/internal/server/server.go @@ -13,7 +13,7 @@ import ( "syscall" "time" - "github.com/willdot/messagebroker/server/peer" + "github.com/willdot/messagebroker/internal" ) // Action represents the type of action that a peer requests to do @@ -111,7 +111,7 @@ func (s *Server) start() { } func (s *Server) handleConn(conn net.Conn) { - peer := peer.New(conn) + peer := NewPeer(conn) slog.Info("handling connection", "peer", peer.Addr()) defer slog.Info("ending connection", "peer", peer.Addr()) @@ -137,11 +137,15 @@ func (s *Server) handleConn(conn net.Conn) { } } -func (s *Server) handleSubscribe(peer *peer.Peer) { +func (s *Server) handleSubscribe(peer *Peer) { slog.Info("handling subscriber", "peer", peer.Addr()) // subscribe the peer to the topic s.subscribePeerToTopic(peer) + s.waitForPeerAction(peer) +} + +func (s *Server) waitForPeerAction(peer *Peer) { // keep handling the peers connection, getting the action from the peer when it wishes to do something else. // once the peers connection ends, it will be unsubscribed from all topics and returned for { @@ -177,10 +181,10 @@ func (s *Server) handleSubscribe(peer *peer.Peer) { } } -func (s *Server) subscribePeerToTopic(peer *peer.Peer) { +func (s *Server) subscribePeerToTopic(peer *Peer) { op := func(conn net.Conn) error { // get the topics the peer wishes to subscribe to - dataLen, err := dataLength(conn) + dataLen, err := dataLengthUint32(conn) if err != nil { slog.Error(err.Error(), "peer", peer.Addr()) writeStatus(Error, "invalid data length of topics provided", conn) @@ -244,11 +248,11 @@ func (s *Server) subscribePeerToTopic(peer *peer.Peer) { _ = peer.RunConnOperation(op) } -func (s *Server) handleUnsubscribe(peer *peer.Peer) { +func (s *Server) handleUnsubscribe(peer *Peer) { slog.Info("handling unsubscriber", "peer", peer.Addr()) op := func(conn net.Conn) error { // get the topics the peer wishes to unsubscribe from - dataLen, err := dataLength(conn) + dataLen, err := dataLengthUint32(conn) if err != nil { slog.Error(err.Error(), "peer", peer.Addr()) writeStatus(Error, "invalid data length of topics provided", conn) @@ -284,11 +288,11 @@ func (s *Server) handleUnsubscribe(peer *peer.Peer) { _ = peer.RunConnOperation(op) } -func (s *Server) handlePublish(peer *peer.Peer) { +func (s *Server) handlePublish(peer *Peer) { slog.Info("handling publisher", "peer", peer.Addr()) for { op := func(conn net.Conn) error { - dataLen, err := dataLength(conn) + topicDataLen, err := dataLengthUint16(conn) if err != nil { if errors.Is(err, io.EOF) { return nil @@ -297,10 +301,10 @@ func (s *Server) handlePublish(peer *peer.Peer) { writeStatus(Error, "invalid data length of data provided", conn) return nil } - if dataLen == 0 { + if topicDataLen == 0 { return nil } - topicBuf := make([]byte, dataLen) + topicBuf := make([]byte, topicDataLen) _, err = conn.Read(topicBuf) if err != nil { slog.Error("failed to read topic from peer", "error", err, "peer", peer.Addr()) @@ -316,17 +320,17 @@ func (s *Server) handlePublish(peer *peer.Peer) { } topicStr = strings.TrimPrefix(topicStr, "topic:") - dataLen, err = dataLength(conn) + msgDataLen, err := dataLengthUint32(conn) if err != nil { slog.Error(err.Error(), "peer", peer.Addr()) writeStatus(Error, "invalid data length of data provided", conn) return nil } - if dataLen == 0 { + if msgDataLen == 0 { return nil } - dataBuf := make([]byte, dataLen) + dataBuf := make([]byte, msgDataLen) _, err = conn.Read(dataBuf) if err != nil { slog.Error("failed to read data from peer", "error", err, "peer", peer.Addr()) @@ -340,7 +344,7 @@ func (s *Server) handlePublish(peer *peer.Peer) { s.topics[topicStr] = topic } - message := newMessage(dataBuf) + message := internal.NewMessage(dataBuf) err = topic.sendMessageToSubscribers(message) if err != nil { @@ -356,14 +360,14 @@ func (s *Server) handlePublish(peer *peer.Peer) { } } -func (s *Server) subscribeToTopics(peer *peer.Peer, topics []string, startAt int) { +func (s *Server) subscribeToTopics(peer *Peer, topics []string, startAt int) { slog.Info("subscribing peer to topics", "topics", topics, "peer", peer.Addr()) for _, topic := range topics { s.addSubsciberToTopic(topic, peer, startAt) } } -func (s *Server) addSubsciberToTopic(topicName string, peer *peer.Peer, startAt int) { +func (s *Server) addSubsciberToTopic(topicName string, peer *Peer, startAt int) { s.mu.Lock() defer s.mu.Unlock() @@ -377,14 +381,14 @@ func (s *Server) addSubsciberToTopic(topicName string, peer *peer.Peer, startAt s.topics[topicName] = t } -func (s *Server) unsubscribeToTopics(peer *peer.Peer, topics []string) { +func (s *Server) unsubscribeToTopics(peer *Peer, topics []string) { slog.Info("unsubscribing peer from topics", "topics", topics, "peer", peer.Addr()) for _, topic := range topics { s.removeSubsciberFromTopic(topic, peer) } } -func (s *Server) removeSubsciberFromTopic(topicName string, peer *peer.Peer) { +func (s *Server) removeSubsciberFromTopic(topicName string, peer *Peer) { s.mu.Lock() defer s.mu.Unlock() @@ -400,7 +404,7 @@ func (s *Server) removeSubsciberFromTopic(topicName string, peer *peer.Peer) { delete(t.subscriptions, peer.Addr()) } -func (s *Server) unsubscribePeerFromAllTopics(peer *peer.Peer) { +func (s *Server) unsubscribePeerFromAllTopics(peer *Peer) { s.mu.Lock() defer s.mu.Unlock() @@ -425,7 +429,7 @@ func (s *Server) getTopic(topicName string) *topic { return nil } -func readAction(peer *peer.Peer, timeout time.Duration) (Action, error) { +func readAction(peer *Peer, timeout time.Duration) (Action, error) { var action Action op := func(conn net.Conn) error { if timeout > 0 { @@ -454,7 +458,7 @@ func readAction(peer *peer.Peer, timeout time.Duration) (Action, error) { return action, nil } -func writeInvalidAction(peer *peer.Peer) { +func writeInvalidAction(peer *Peer) { op := func(conn net.Conn) error { writeStatus(Error, "unknown action", conn) return nil @@ -463,7 +467,7 @@ func writeInvalidAction(peer *peer.Peer) { _ = peer.RunConnOperation(op) } -func dataLength(conn net.Conn) (uint32, error) { +func dataLengthUint32(conn net.Conn) (uint32, error) { var dataLen uint32 err := binary.Read(conn, binary.BigEndian, &dataLen) if err != nil { @@ -472,6 +476,15 @@ func dataLength(conn net.Conn) (uint32, error) { return dataLen, nil } +func dataLengthUint16(conn net.Conn) (uint16, error) { + var dataLen uint16 + err := binary.Read(conn, binary.BigEndian, &dataLen) + if err != nil { + return 0, err + } + return dataLen, nil +} + func writeStatus(status Status, message string, conn net.Conn) { statusB := make([]byte, 2) binary.BigEndian.PutUint16(statusB, uint16(status)) @@ -479,8 +492,8 @@ func writeStatus(status Status, message string, conn net.Conn) { headers := statusB if len(message) > 0 { - sizeB := make([]byte, 4) - binary.BigEndian.PutUint32(sizeB, uint32(len(message))) + sizeB := make([]byte, 2) + binary.BigEndian.PutUint16(sizeB, uint16(len(message))) headers = append(headers, sizeB...) } diff --git a/server/server_test.go b/internal/server/server_test.go similarity index 98% rename from server/server_test.go rename to internal/server/server_test.go index 8b61392..2d44f96 100644 --- a/server/server_test.go +++ b/internal/server/server_test.go @@ -10,6 +10,7 @@ import ( "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "github.com/willdot/messagebroker/internal/messagestore" ) const ( @@ -39,7 +40,7 @@ func createServerWithExistingTopic(t *testing.T, topicName string) *Server { srv.topics[topicName] = &topic{ name: topicName, subscriptions: make(map[net.Addr]*subscriber), - messageStore: NewMemoryStore(), + messageStore: messagestore.NewMemoryStore(), } return srv @@ -211,7 +212,7 @@ func TestInvalidAction(t *testing.T) { expectedMessage := "unknown action" - var dataLen uint32 + var dataLen uint16 err = binary.Read(conn, binary.BigEndian, &dataLen) require.NoError(t, err) assert.Equal(t, len(expectedMessage), int(dataLen)) @@ -249,7 +250,7 @@ func TestInvalidTopicDataPublished(t *testing.T) { expectedMessage := "topic data does not contain 'topic:' prefix" - var dataLen uint32 + var dataLen uint16 err = binary.Read(publisherConn, binary.BigEndian, &dataLen) require.NoError(t, err) assert.Equal(t, len(expectedMessage), int(dataLen)) @@ -352,7 +353,7 @@ func TestSendsDataToTopicSubscriberNacksThenAcks(t *testing.T) { // check the subsribers got the data readMessage := func(conn net.Conn, ack Action) { - var topicLen uint64 + var topicLen uint16 err = binary.Read(conn, binary.BigEndian, &topicLen) require.NoError(t, err) @@ -382,7 +383,7 @@ func TestSendsDataToTopicSubscriberNacksThenAcks(t *testing.T) { readMessage(subscriberConn, Ack) // reading for another message should now timeout but give enough time for the ack delay to kick in // should the second read of the message not have been ack'd properly - var topicLen uint64 + var topicLen uint16 _ = subscriberConn.SetReadDeadline(time.Now().Add(ackDelay + time.Millisecond*100)) err = binary.Read(subscriberConn, binary.BigEndian, &topicLen) require.Error(t, err) @@ -406,7 +407,7 @@ func TestSendsDataToTopicSubscriberDoesntAckMessage(t *testing.T) { // check the subsribers got the data readMessage := func(conn net.Conn, ack bool) { - var topicLen uint64 + var topicLen uint16 err = binary.Read(conn, binary.BigEndian, &topicLen) require.NoError(t, err) @@ -440,7 +441,7 @@ func TestSendsDataToTopicSubscriberDoesntAckMessage(t *testing.T) { // reading for another message should now timeout but give enough time for the ack delay to kick in // should the second read of the message not have been ack'd properly - var topicLen uint64 + var topicLen uint16 _ = subscriberConn.SetReadDeadline(time.Now().Add(ackDelay + time.Millisecond*100)) err = binary.Read(subscriberConn, binary.BigEndian, &topicLen) require.Error(t, err) @@ -464,7 +465,7 @@ func TestSendsDataToTopicSubscriberDeliveryCountTooHighWithNoAck(t *testing.T) { // check the subsribers got the data readMessage := func(conn net.Conn, ack bool) { - var topicLen uint64 + var topicLen uint16 err = binary.Read(conn, binary.BigEndian, &topicLen) require.NoError(t, err) @@ -500,7 +501,7 @@ func TestSendsDataToTopicSubscriberDeliveryCountTooHighWithNoAck(t *testing.T) { readMessage(subscriberConn, false) // reading for the message should now timeout as we have nack'd the message too many times - var topicLen uint64 + var topicLen uint16 _ = subscriberConn.SetReadDeadline(time.Now().Add(ackDelay + time.Millisecond*100)) err = binary.Read(subscriberConn, binary.BigEndian, &topicLen) require.Error(t, err) @@ -592,7 +593,7 @@ func TestSubscribeAndReplaysFromIndex(t *testing.T) { } func readMessage(t *testing.T, subscriberConn net.Conn) []byte { - var topicLen uint64 + var topicLen uint16 err := binary.Read(subscriberConn, binary.BigEndian, &topicLen) require.NoError(t, err) diff --git a/server/subscriber.go b/internal/server/subscriber.go similarity index 70% rename from server/subscriber.go rename to internal/server/subscriber.go index 15e12e5..655a86a 100644 --- a/server/subscriber.go +++ b/internal/server/subscriber.go @@ -7,33 +7,24 @@ import ( "net" "time" - "github.com/willdot/messagebroker/server/peer" + "github.com/willdot/messagebroker/internal" ) type subscriber struct { - peer *peer.Peer + peer *Peer topic string - messages chan message + messages chan internal.Message unsubscribeCh chan struct{} ackDelay time.Duration ackTimeout time.Duration } -type message struct { - data []byte - deliveryCount int -} - -func newMessage(data []byte) message { - return message{data: data, deliveryCount: 1} -} - -func newSubscriber(peer *peer.Peer, topic *topic, ackDelay, ackTimeout time.Duration, startAt int) *subscriber { +func newSubscriber(peer *Peer, topic *topic, ackDelay, ackTimeout time.Duration, startAt int) *subscriber { s := &subscriber{ peer: peer, topic: topic.name, - messages: make(chan message), + messages: make(chan internal.Message), ackDelay: ackDelay, ackTimeout: ackTimeout, unsubscribeCh: make(chan struct{}, 1), @@ -42,7 +33,7 @@ func newSubscriber(peer *peer.Peer, topic *topic, ackDelay, ackTimeout time.Dura go s.sendMessages() go func() { - topic.messageStore.ReadFrom(startAt, func(msg message) { + topic.messageStore.ReadFrom(startAt, func(msg internal.Message) { select { case s.messages <- msg: return @@ -70,18 +61,18 @@ func (s *subscriber) sendMessages() { continue } - if msg.deliveryCount >= 5 { + if msg.DeliveryCount >= 5 { slog.Error("max delivery count for message. Dropping", "peer", s.peer.Addr()) continue } - msg.deliveryCount++ + msg.DeliveryCount++ s.addMessage(msg, s.ackDelay) } } } -func (s *subscriber) addMessage(msg message, delay time.Duration) { +func (s *subscriber) addMessage(msg internal.Message, delay time.Duration) { go func() { timer := time.NewTimer(delay) defer timer.Stop() @@ -95,27 +86,25 @@ func (s *subscriber) addMessage(msg message, delay time.Duration) { }() } -func (s *subscriber) sendMessage(topic string, msg message) (bool, error) { +func (s *subscriber) sendMessage(topic string, msg internal.Message) (bool, error) { var ack bool op := func(conn net.Conn) error { - // TODO: why did I chose uint64 for topic len? - topicB := make([]byte, 8) - binary.BigEndian.PutUint64(topicB, uint64(len(topic))) + topicB := make([]byte, 2) + binary.BigEndian.PutUint16(topicB, uint16(len(topic))) headers := topicB headers = append(headers, []byte(topic)...) // TODO: if message is empty, return error? dataLenB := make([]byte, 8) - binary.BigEndian.PutUint64(dataLenB, uint64(len(msg.data))) + binary.BigEndian.PutUint64(dataLenB, uint64(len(msg.Data))) headers = append(headers, dataLenB...) - _, err := conn.Write(append(headers, msg.data...)) + _, err := conn.Write(append(headers, msg.Data...)) if err != nil { return fmt.Errorf("failed to write to peer: %w", err) } - var ackRes Action if err := conn.SetReadDeadline(time.Now().Add(s.ackTimeout)); err != nil { slog.Error("failed to set connection read deadline", "error", err, "peer", s.peer.Addr()) } @@ -124,6 +113,7 @@ func (s *subscriber) sendMessage(topic string, msg message) (bool, error) { slog.Error("failed to reset connection read deadline", "error", err, "peer", s.peer.Addr()) } }() + var ackRes Action err = binary.Read(conn, binary.BigEndian, &ackRes) if err != nil { return fmt.Errorf("failed to read ack from peer: %w", err) diff --git a/server/topic.go b/internal/server/topic.go similarity index 66% rename from server/topic.go rename to internal/server/topic.go index 02987e4..fa5cf5d 100644 --- a/server/topic.go +++ b/internal/server/topic.go @@ -4,11 +4,14 @@ import ( "fmt" "net" "sync" + + "github.com/willdot/messagebroker/internal" + "github.com/willdot/messagebroker/internal/messagestore" ) type Store interface { - Write(msg message) error - ReadFrom(offset int, handleFunc func(msg message)) + Write(msg internal.Message) error + ReadFrom(offset int, handleFunc func(msg internal.Message)) } type topic struct { @@ -19,7 +22,7 @@ type topic struct { } func newTopic(name string) *topic { - messageStore := NewMemoryStore() + messageStore := messagestore.NewMemoryStore() return &topic{ name: name, subscriptions: make(map[net.Addr]*subscriber), @@ -27,7 +30,7 @@ func newTopic(name string) *topic { } } -func (t *topic) sendMessageToSubscribers(msg message) error { +func (t *topic) sendMessageToSubscribers(msg internal.Message) error { err := t.messageStore.Write(msg) if err != nil { return fmt.Errorf("failed to write message to store: %w", err) -- 2.51.2