From 1ac24669d3c211f08442bac5359153d97384a6e3 Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Mon, 28 Sep 2026 21:42:33 +0200 Subject: [PATCH] Stop the streamer without sending on closed queues (closes #34) Stopping the daemon while the RIS Live feed was flowing could panic with "send on closed channel" and skip the rest of the shutdown. Stop cancels the stream and closes the handler queues under the streamer's write lock, but the read loop checked for a stop only before parsing each line. It now checks again under the read lock it already takes just before handing a message to the queues, so it never sends to a closed queue. Stop also clears its cancel function and returns early when there is none, so a second call no longer closes the queues again. A test forces both cases. Behaviour change: Stop before Start now does nothing. Model: opus-5-5 --- TODO.md | 4 +++ internal/streamer/streamer.go | 18 ++++++++++--- internal/streamer/streamer_test.go | 42 ++++++++++++++++++++++++++++++ 3 files changed, 61 insertions(+), 3 deletions(-) diff --git a/TODO.md b/TODO.md index f74ae68..70ef510 100644 --- a/TODO.md +++ b/TODO.md @@ -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 diff --git a/internal/streamer/streamer.go b/internal/streamer/streamer.go index 42d6547..b8b75a7 100644 --- a/internal/streamer/streamer.go +++ b/internal/streamer/streamer.go @@ -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 diff --git a/internal/streamer/streamer_test.go b/internal/streamer/streamer_test.go index 3f0b4aa..5d9582b 100644 --- a/internal/streamer/streamer_test.go +++ b/internal/streamer/streamer_test.go @@ -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 {