diff --git a/server/message_store.go b/server/message_store.go index fb6aefb..a9439ed 100644 --- a/server/message_store.go +++ b/server/message_store.go @@ -5,18 +5,21 @@ import ( "sync" ) +// Memory store allows messages to be stored in memory type MemoryStore struct { mu sync.Mutex msgs map[int]message offset int } +// New memory store initializes a new in memory store func NewMemoryStore() *MemoryStore { return &MemoryStore{ msgs: make(map[int]message), } } +// Write will write the provided message to the in memory store func (m *MemoryStore) Write(msg message) error { m.mu.Lock() defer m.mu.Unlock() @@ -28,6 +31,7 @@ func (m *MemoryStore) Write(msg message) error { 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") diff --git a/server/subscriber.go b/server/subscriber.go index 415cb8a..288c094 100644 --- a/server/subscriber.go +++ b/server/subscriber.go @@ -44,6 +44,10 @@ func newSubscriber(peer *peer.Peer, topic *topic, ackDelay, ackTimeout time.Dura offset := startAt go func() { + if startAt < 0 { + return + } + err := topic.messageStore.ReadFrom(offset, func(msg message) { s.messages <- msg })