2020-06-02 22:37:10 +00:00
|
|
|
package stream
|
|
|
|
|
2020-07-06 21:29:45 +00:00
|
|
|
// eventSnapshot represents the state of memdb for a given topic and key at some
|
2020-06-02 22:37:10 +00:00
|
|
|
// point in time. It is modelled as a buffer of events so that snapshots can be
|
|
|
|
// streamed to possibly multiple subscribers concurrently, and can be trivially
|
2020-07-06 21:29:45 +00:00
|
|
|
// cached by retaining a reference to a Snapshot. Once the reference to eventSnapshot
|
2020-06-10 23:07:58 +00:00
|
|
|
// is dropped from memory, any subscribers still reading from it may do so by following
|
2020-07-08 18:45:18 +00:00
|
|
|
// their pointers. When the last subscriber unsubscribes, the snapshot is garbage
|
2020-06-10 23:07:58 +00:00
|
|
|
// collected automatically by Go's runtime. This simplifies snapshot and buffer
|
|
|
|
// management dramatically.
|
2020-07-06 21:29:45 +00:00
|
|
|
type eventSnapshot struct {
|
2020-10-01 17:51:55 +00:00
|
|
|
// First item in the buffer. Used as the Head of a subscription, or to
|
|
|
|
// splice this snapshot onto another one.
|
|
|
|
First *bufferItem
|
2020-06-02 22:37:10 +00:00
|
|
|
|
2020-10-01 17:51:55 +00:00
|
|
|
// buffer is the Head of the snapshot buffer the fn should write to.
|
|
|
|
buffer *eventBuffer
|
2020-06-02 22:37:10 +00:00
|
|
|
}
|
|
|
|
|
2020-10-01 17:51:55 +00:00
|
|
|
// newEventSnapshot creates an empty snapshot buffer.
|
|
|
|
func newEventSnapshot() *eventSnapshot {
|
|
|
|
snapBuffer := newEventBuffer()
|
|
|
|
return &eventSnapshot{
|
|
|
|
First: snapBuffer.Head(),
|
|
|
|
buffer: snapBuffer,
|
2020-06-02 22:37:10 +00:00
|
|
|
}
|
2020-10-01 17:51:55 +00:00
|
|
|
}
|
2020-06-10 23:07:58 +00:00
|
|
|
|
2020-10-01 17:51:55 +00:00
|
|
|
// appendAndSlice populates the snapshot buffer by calling the SnapshotFunc,
|
|
|
|
// then adding an endOfSnapshot framing event, and finally by splicing in
|
|
|
|
// events from the topicBuffer.
|
|
|
|
func (s *eventSnapshot) appendAndSplice(req SubscribeRequest, fn SnapshotFunc, topicBufferHead *bufferItem) {
|
|
|
|
idx, err := fn(req, s.buffer)
|
|
|
|
if err != nil {
|
|
|
|
s.buffer.AppendItem(&bufferItem{Err: err})
|
|
|
|
return
|
|
|
|
}
|
|
|
|
s.buffer.Append([]Event{{
|
|
|
|
Topic: req.Topic,
|
|
|
|
Index: idx,
|
|
|
|
Payload: endOfSnapshot{},
|
|
|
|
}})
|
|
|
|
s.spliceFromTopicBuffer(topicBufferHead, idx)
|
2020-06-02 22:37:10 +00:00
|
|
|
}
|
|
|
|
|
2020-10-01 17:51:55 +00:00
|
|
|
// spliceFromTopicBuffer traverses the topicBuffer looking for the last item
|
|
|
|
// in the buffer, or the first item where the index is greater than idx. Once
|
|
|
|
// the item is found it is appended to the snapshot buffer.
|
2020-07-06 21:29:45 +00:00
|
|
|
func (s *eventSnapshot) spliceFromTopicBuffer(topicBufferHead *bufferItem, idx uint64) {
|
2020-06-10 23:07:58 +00:00
|
|
|
item := topicBufferHead
|
2020-06-02 22:37:10 +00:00
|
|
|
for {
|
2020-07-08 04:31:22 +00:00
|
|
|
switch {
|
2021-06-04 22:19:19 +00:00
|
|
|
case item.Err != nil:
|
2020-07-08 04:31:22 +00:00
|
|
|
// This case is not currently possible because errors can only come
|
|
|
|
// from a snapshot func, and this is consuming events from a topic
|
|
|
|
// buffer which does not contain a snapshot.
|
|
|
|
// Handle this case anyway in case errors can come from other places
|
|
|
|
// in the future.
|
2021-06-04 22:19:19 +00:00
|
|
|
s.buffer.AppendItem(item)
|
2020-06-02 22:37:10 +00:00
|
|
|
return
|
|
|
|
|
2021-06-04 22:19:19 +00:00
|
|
|
case len(item.Events) > 0 && item.Events[0].Index > idx:
|
2020-06-10 23:07:58 +00:00
|
|
|
// We've found an update in the topic buffer that happened after our
|
|
|
|
// snapshot was taken, splice it into the snapshot buffer so subscribers
|
|
|
|
// can continue to read this and others after it.
|
2021-06-04 22:19:19 +00:00
|
|
|
s.buffer.AppendItem(item)
|
2020-06-10 23:07:58 +00:00
|
|
|
return
|
2020-06-02 22:37:10 +00:00
|
|
|
}
|
2020-07-08 04:31:22 +00:00
|
|
|
|
2021-06-04 22:29:04 +00:00
|
|
|
next, ok := item.NextNoBlock()
|
|
|
|
if !ok {
|
2021-06-04 22:19:19 +00:00
|
|
|
// We reached the head of the topic buffer. We don't want any of the
|
|
|
|
// events in the topic buffer as they came before the snapshot.
|
|
|
|
// Append a link to any future items.
|
2021-06-04 22:29:04 +00:00
|
|
|
s.buffer.AppendItem(next)
|
2021-06-04 22:19:19 +00:00
|
|
|
return
|
|
|
|
}
|
|
|
|
// Proceed to the next item in the topic buffer
|
2020-06-02 22:37:10 +00:00
|
|
|
item = next
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2020-07-06 21:29:45 +00:00
|
|
|
// err returns an error if the snapshot func has failed with an error or nil
|
2020-06-02 22:37:10 +00:00
|
|
|
// otherwise. Nil doesn't necessarily mean there won't be an error but there
|
|
|
|
// hasn't been one yet.
|
2020-07-06 21:29:45 +00:00
|
|
|
func (s *eventSnapshot) err() error {
|
2020-06-02 22:37:10 +00:00
|
|
|
// Fetch the head of the buffer, this is atomic. If the snapshot func errored
|
|
|
|
// then the last event will be an error.
|
2020-10-01 17:51:55 +00:00
|
|
|
head := s.buffer.Head()
|
2020-06-02 22:37:10 +00:00
|
|
|
return head.Err
|
|
|
|
}
|