Stop the streamer without sending on closed queues (closes #34) #36
@@ -23,6 +23,10 @@ runs make check on main.
|
|||||||
|
|
||||||
# Completed Steps
|
# 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
|
- 2026-09-28: `docker stop` no longer kills the daemon 2 seconds after the
|
||||||
stop signal: the entrypoint switches to the `routewatch` user with
|
stop signal: the entrypoint switches to the `routewatch` user with
|
||||||
`setpriv` instead of `runuser`, so the daemon receives the signal itself
|
`setpriv` instead of `runuser`, so the daemon receives the signal itself
|
||||||
|
|||||||
@@ -210,9 +210,14 @@ func (s *Streamer) Start() error {
|
|||||||
// the connection status in metrics. This method is safe to call multiple times.
|
// the connection status in metrics. This method is safe to call multiple times.
|
||||||
func (s *Streamer) Stop() {
|
func (s *Streamer) Stop() {
|
||||||
s.mu.Lock()
|
s.mu.Lock()
|
||||||
if s.cancel != nil {
|
if s.cancel == nil {
|
||||||
s.cancel()
|
// 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
|
// Close all handler queues to signal workers to stop
|
||||||
for _, info := range s.handlers {
|
for _, info := range s.handlers {
|
||||||
close(info.queue)
|
close(info.queue)
|
||||||
@@ -660,8 +665,15 @@ func (s *Streamer) stream(ctx context.Context) error {
|
|||||||
continue
|
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()
|
s.mu.RLock()
|
||||||
|
if ctx.Err() != nil {
|
||||||
|
s.mu.RUnlock()
|
||||||
|
|
||||||
|
return ctx.Err()
|
||||||
|
}
|
||||||
for _, info := range s.handlers {
|
for _, info := range s.handlers {
|
||||||
if !info.handler.WantsMessage(msg.Type) {
|
if !info.handler.WantsMessage(msg.Type) {
|
||||||
continue
|
continue
|
||||||
|
|||||||
@@ -2,6 +2,8 @@ package streamer
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"errors"
|
||||||
|
"io"
|
||||||
"net/http"
|
"net/http"
|
||||||
"net/http/httptest"
|
"net/http/httptest"
|
||||||
"runtime"
|
"runtime"
|
||||||
@@ -10,6 +12,7 @@ import (
|
|||||||
|
|
||||||
"git.eeqj.de/sneak/routewatch/internal/logger"
|
"git.eeqj.de/sneak/routewatch/internal/logger"
|
||||||
"git.eeqj.de/sneak/routewatch/internal/metrics"
|
"git.eeqj.de/sneak/routewatch/internal/metrics"
|
||||||
|
"git.eeqj.de/sneak/routewatch/internal/ristypes"
|
||||||
)
|
)
|
||||||
|
|
||||||
func TestNewStreamer(t *testing.T) {
|
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
|
// settledGoroutineCount lets transient goroutines finish, then reports the
|
||||||
// current count.
|
// current count.
|
||||||
func settledGoroutineCount() int {
|
func settledGoroutineCount() int {
|
||||||
|
|||||||
Reference in New Issue
Block a user