diff --git a/example/main.go b/example/main.go index d4dc09e..978e4a6 100644 --- a/example/main.go +++ b/example/main.go @@ -33,7 +33,7 @@ func main() { }() startAt := 0 startAtType := server.Current - if *consumeFrom >= 0-1 { + if *consumeFrom > -1 { startAtType = server.From startAt = *consumeFrom } diff --git a/server/server.go b/server/server.go index 5600f38..156e37e 100644 --- a/server/server.go +++ b/server/server.go @@ -75,11 +75,6 @@ const ( From StartAtType = 2 ) -type Store interface { - Write(msg MessageToSend) error - ReadFrom(offset int, handleFunc func(msg MessageToSend)) error -} - // Server accepts subscribe and publish connections and passes messages around type Server struct { Addr string @@ -222,7 +217,6 @@ func (s *Server) subscribePeerToTopic(peer *peer.Peer) { } var topics []string - fmt.Println(string(buf)) err = json.Unmarshal(buf, &topics) if err != nil { slog.Error("failed to unmarshal subscibers topic data", "error", err, "peer", peer.Addr()) @@ -316,8 +310,6 @@ type MessageToSend struct { func (s *Server) handlePublish(peer *peer.Peer) { slog.Info("handling publisher", "peer", peer.Addr()) for { - var message *MessageToSend - op := func(conn net.Conn) error { dataLen, err := dataLength(conn) if err != nil { @@ -365,7 +357,7 @@ func (s *Server) handlePublish(peer *peer.Peer) { return nil } - message = &MessageToSend{ + message := MessageToSend{ topic: topicStr, data: dataBuf, } @@ -376,7 +368,7 @@ func (s *Server) handlePublish(peer *peer.Peer) { s.topics[message.topic] = topic } - err = topic.sendMessageToSubscribers(*message) + err = topic.sendMessageToSubscribers(message) if err != nil { slog.Error("failed to send message to subscribers", "error", err, "peer", peer.Addr()) writeStatus(Error, "failed to send message to subscribers", conn) @@ -406,7 +398,7 @@ func (s *Server) addSubsciberToTopic(topicName string, peer *peer.Peer, startAt t = newTopic(topicName) } - t.subscriptions[peer.Addr()] = newSubscriber(peer, topicName, s.ackDelay, s.ackTimeout, t.messageStore, startAt) + t.subscriptions[peer.Addr()] = newSubscriber(peer, t, s.ackDelay, s.ackTimeout, startAt) s.topics[topicName] = t } diff --git a/server/subscriber.go b/server/subscriber.go index 318994f..ef14895 100644 --- a/server/subscriber.go +++ b/server/subscriber.go @@ -29,10 +29,10 @@ func newMessage(data []byte) message { return message{data: data, deliveryCount: 1} } -func newSubscriber(peer *peer.Peer, topic string, ackDelay, ackTimeout time.Duration, messageStore Store, startAt int) *subscriber { +func newSubscriber(peer *peer.Peer, topic *topic, ackDelay, ackTimeout time.Duration, startAt int) *subscriber { s := &subscriber{ peer: peer, - topic: topic, + topic: topic.name, messages: make(chan message), ackDelay: ackDelay, ackTimeout: ackTimeout, @@ -44,7 +44,7 @@ func newSubscriber(peer *peer.Peer, topic string, ackDelay, ackTimeout time.Dura offset := startAt go func() { - err := messageStore.ReadFrom(offset, func(msg MessageToSend) { + err := topic.messageStore.ReadFrom(offset, func(msg MessageToSend) { s.messages <- newMessage(msg.data) }) if err != nil { diff --git a/server/topic.go b/server/topic.go index f9a9236..9d631a5 100644 --- a/server/topic.go +++ b/server/topic.go @@ -6,6 +6,11 @@ import ( "sync" ) +type Store interface { + Write(msg MessageToSend) error + ReadFrom(offset int, handleFunc func(msg MessageToSend)) error +} + type topic struct { name string subscriptions map[net.Addr]*subscriber