Fix two goroutine leaks: stats handlers on timeout, streamer tickers on reconnect (closes #12)
check / check (push) Failing after 1s
check / check (push) Failing after 1s
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 was merged in pull request #19.
This commit is contained in:
@@ -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
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user