Fix two goroutine leaks: stats handlers on timeout, streamer tickers on reconnect (closes #12)

The stats handlers ran the database query in a goroutine that sent on an
unbuffered channel. When the 4s request timeout won, nothing received and
the goroutine blocked forever; the status page polls every 2s, so once the
query exceeds the timeout every poll leaked one goroutine. Give both
channels capacity 1 so the send always completes.

The streamer started two ticker goroutines per connection that exited only
with the streamer's lifetime context, leaking two on every reconnect. Scope
them to a per-connection context cancelled when the stream call returns.

Tests force the stats timeout repeatedly and drive many reconnects, then
assert the goroutine count settles back to its starting value. The streamer
gains an internal endpoint field so a test can point it at a local server.

Model: opus-4-8
This commit is contained in:
2026-09-21 13:37:45 +00:00
parent 54014c88c8
commit 7a67c3e96e
4 changed files with 191 additions and 9 deletions
+12 -3
View File
@@ -106,6 +106,7 @@ type handlerInfo struct {
type Streamer struct {
logger *logger.Logger
client *http.Client
url string
handlers []*handlerInfo
rawHandler RawMessageHandler
mu sync.RWMutex
@@ -124,6 +125,7 @@ type Streamer struct {
func New(logger *logger.Logger, metrics *metrics.Tracker) *Streamer {
return &Streamer{
logger: logger,
url: risLiveURL,
client: &http.Client{
Timeout: 0, // No timeout for streaming
Transport: &http.Transport{
@@ -463,7 +465,14 @@ func (s *Streamer) streamWithReconnect(ctx context.Context) {
}
func (s *Streamer) stream(ctx context.Context) error {
req, err := http.NewRequestWithContext(ctx, "GET", risLiveURL, nil)
// connCtx is scoped to this single connection: cancelling it when stream
// returns stops the ticker goroutines below, so a reconnect does not leak
// them. Without this they would live until the streamer's lifetime context
// is cancelled, leaking two per reconnect.
connCtx, connCancel := context.WithCancel(ctx)
defer connCancel()
req, err := http.NewRequestWithContext(ctx, "GET", s.url, nil)
if err != nil {
return fmt.Errorf("failed to create request: %w", err)
}
@@ -516,7 +525,7 @@ func (s *Streamer) stream(ctx context.Context) error {
select {
case <-metricsTicker.C:
s.logMetrics()
case <-ctx.Done():
case <-connCtx.Done():
return
}
}
@@ -536,7 +545,7 @@ func (s *Streamer) stream(ctx context.Context) error {
s.metrics.RecordWireBytes(delta)
lastWireBytes = currentBytes
}
case <-ctx.Done():
case <-connCtx.Done():
return
}
}
+70
View File
@@ -1,7 +1,12 @@
package streamer
import (
"context"
"net/http"
"net/http/httptest"
"runtime"
"testing"
"time"
"git.eeqj.de/sneak/routewatch/internal/logger"
"git.eeqj.de/sneak/routewatch/internal/metrics"
@@ -32,3 +37,68 @@ func TestNewStreamer(t *testing.T) {
t.Error("metrics tracker not set correctly")
}
}
// TestStreamDoesNotLeakTickersAcrossReconnects drives many short-lived
// connections (each stream call is one reconnect cycle) and asserts the
// goroutine count returns to its starting value. Each connection starts two
// ticker goroutines; before the fix they lived until the streamer's lifetime
// context was cancelled, so every reconnect leaked two.
func TestStreamDoesNotLeakTickersAcrossReconnects(t *testing.T) {
// The handler returns immediately, so the response body is empty and each
// stream call ends at once, standing in for a dropped connection.
srv := httptest.NewServer(http.HandlerFunc(func(_ http.ResponseWriter, _ *http.Request) {}))
defer srv.Close()
s := New(logger.New(), metrics.New())
s.url = srv.URL
// One warm-up connection so any persistent HTTP transport goroutine exists
// before we take the baseline.
if err := s.stream(context.Background()); err != nil {
t.Fatalf("warm-up stream returned error: %v", err)
}
s.client.CloseIdleConnections()
baseline := settledGoroutineCount()
const reconnects = 20
for range reconnects {
if err := s.stream(context.Background()); err != nil {
t.Fatalf("stream returned error: %v", err)
}
}
s.client.CloseIdleConnections()
if !waitForGoroutines(baseline) {
t.Fatalf("goroutines did not return to baseline %d after %d reconnects, got %d",
baseline, reconnects, runtime.NumGoroutine())
}
}
// settledGoroutineCount lets transient goroutines finish, then reports the
// current count.
func settledGoroutineCount() int {
prev := runtime.NumGoroutine()
for range 20 {
time.Sleep(10 * time.Millisecond)
cur := runtime.NumGoroutine()
if cur == prev {
return cur
}
prev = cur
}
return prev
}
// waitForGoroutines waits until the goroutine count drops to target or below.
func waitForGoroutines(target int) bool {
for range 100 {
if runtime.NumGoroutine() <= target {
return true
}
time.Sleep(10 * time.Millisecond)
}
return false
}