Stop the streamer without sending on closed queues (closes #34) #36

Merged
clawbot merged 1 commits from issue-34-stop-panic into next 2026-09-28 21:42:34 +02:00
3 changed files with 61 additions and 3 deletions
+4
View File
@@ -23,6 +23,10 @@ runs make check on main.
# Completed Steps
- 2026-09-28: stopping the daemon while the feed is flowing no longer
panics with "send on closed channel": the read loop checks for a stop
just before handing a message to the handler queues, and a second
`Stop` no longer closes the queues again (closes #34)
- 2026-09-28: `docker stop` no longer kills the daemon 2 seconds after the
stop signal: the entrypoint switches to the `routewatch` user with
`setpriv` instead of `runuser`, so the daemon receives the signal itself
+15 -3
View File
@@ -210,9 +210,14 @@ func (s *Streamer) Start() error {
// the connection status in metrics. This method is safe to call multiple times.
func (s *Streamer) Stop() {
s.mu.Lock()
if s.cancel != nil {
s.cancel()
if s.cancel == nil {
// Not started, or already stopped: closing the queues again would panic.
s.mu.Unlock()
return
}
s.cancel()
s.cancel = nil
// Close all handler queues to signal workers to stop
for _, info := range s.handlers {
close(info.queue)
@@ -660,8 +665,15 @@ func (s *Streamer) stream(ctx context.Context) error {
continue
}
// Dispatch to interested handlers
// Dispatch to interested handlers. Stop cancels ctx and closes the
// queues under the write lock, so if ctx is cancelled here, under the
// read lock, the queues are closed and must not be sent to.
s.mu.RLock()
if ctx.Err() != nil {
s.mu.RUnlock()
return ctx.Err()
}
for _, info := range s.handlers {
if !info.handler.WantsMessage(msg.Type) {
continue
+42
View File
@@ -2,6 +2,8 @@ package streamer
import (
"context"
"errors"
"io"
"net/http"
"net/http/httptest"
"runtime"
@@ -10,6 +12,7 @@ import (
"git.eeqj.de/sneak/routewatch/internal/logger"
"git.eeqj.de/sneak/routewatch/internal/metrics"
"git.eeqj.de/sneak/routewatch/internal/ristypes"
)
func TestNewStreamer(t *testing.T) {
@@ -75,6 +78,45 @@ func TestStreamDoesNotLeakTickersAcrossReconnects(t *testing.T) {
}
}
// updateHandler wants UPDATE messages and does nothing with them.
type updateHandler struct{}
func (updateHandler) WantsMessage(messageType string) bool { return messageType == "UPDATE" }
func (updateHandler) HandleMessage(*ristypes.RISMessage) {}
func (updateHandler) QueueCapacity() int { return 10 }
// TestStopBeforeMessageReachesQueues stops the streamer after the read loop
// has checked for cancellation but before it hands the message to the handler
// queues. That is the gap a stop from another goroutine can land in, and it
// used to end in "send on closed channel". The raw handler runs in that gap on
// the read loop itself, so calling Stop from it hits the gap every time.
func TestStopBeforeMessageReachesQueues(t *testing.T) {
const line = `{"type":"ris_message","data":{"type":"UPDATE","peer":"192.0.2.1",` +
`"peer_asn":"64496","timestamp":1700000000}}` + "\n"
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
_, _ = io.WriteString(w, line)
}))
defer srv.Close()
s := New(logger.New(), metrics.New())
s.url = srv.URL
s.RegisterHandler(updateHandler{})
s.RegisterRawHandler(func(string) { s.Stop() })
// Start would run the stream in the background, where the test cannot
// wait for it. Setting cancel as Start does lets Stop cancel the stream
// run here instead.
ctx, cancel := context.WithCancel(context.Background())
s.cancel = cancel
if err := s.stream(ctx); !errors.Is(err, context.Canceled) {
t.Fatalf("stream returned %v, want %v", err, context.Canceled)
}
// A second Stop must not close the queues again.
s.Stop()
}
// settledGoroutineCount lets transient goroutines finish, then reports the
// current count.
func settledGoroutineCount() int {