diff --git a/cmd/relay/stream/eventmgr/event_manager.go b/cmd/relay/stream/eventmgr/event_manager.go index b25fd129..9d9d2519 100644 --- a/cmd/relay/stream/eventmgr/event_manager.go +++ b/cmd/relay/stream/eventmgr/event_manager.go @@ -168,6 +168,7 @@ func (em *EventManager) Subscribe(ctx context.Context, ident string, filter func } // TODO: send an error frame or something? + // NOTE: not doing em.rmSubscriber(sub) here because it hasn't been added yet close(out) return } @@ -175,6 +176,12 @@ func (em *EventManager) Subscribe(ctx context.Context, ident string, filter func // now, start buffering events from the live stream em.addSubscriber(sub) + // ensure that we clean up any return paths from here out, after having added the subscriber. Note that `out` is not `sub.output`, so needs to be closed separately. + defer func() { + close(out) + em.rmSubscriber(sub) + }() + first := <-sub.outgoing // run playback again to get us to the events that have started buffering @@ -193,10 +200,6 @@ func (em *EventManager) Subscribe(ctx context.Context, ident string, filter func }); err != nil { if !errors.Is(err, ErrCaughtUp) { em.log.Error("events playback", "err", err) - - // TODO: send an error frame or something? - close(out) - em.rmSubscriber(sub) return } } @@ -206,7 +209,6 @@ func (em *EventManager) Subscribe(ctx context.Context, ident string, filter func select { case out <- evt: case <-done: - em.rmSubscriber(sub) return } }