package web import ( "bufio" "bytes" "context" "fmt" "net/http" "net/http/httptest" "runtime/pprof" "strings" "testing" "time" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "go.uber.org/goleak" ) // verifyNoRuntimeGoroutineLeaks complements goleak with Go 1.27's runtime // reachability analysis, which can identify permanently blocked goroutines. func verifyNoRuntimeGoroutineLeaks(t *testing.T) { t.Helper() profile := pprof.Lookup("goroutineleak") require.NotNil(t, profile, "Go 1.27 goroutineleak profile is unavailable") var report bytes.Buffer require.NoError(t, profile.WriteTo(&report, 1)) if count := profile.Count(); count != 0 { t.Fatalf("runtime detected %d leaked goroutine(s):\n%s", count, report.String()) } } // TestJobManagerLifecycleNoGoroutineLeaks runs a full job-manager lifecycle — // SSE clients, relays, and jobs reaching every terminal state — and verifies // no goroutines outlive it. func TestJobManagerLifecycleNoGoroutineLeaks(t *testing.T) { defer goleak.VerifyNone(t) t.Cleanup(func() { verifyNoRuntimeGoroutineLeaks(t) }) jm := NewJobManager() // SSE-style clients attached through the whole lifecycle. clientAll := jm.RegisterClient() clientFiltered := jm.RegisterClient() // A job that completes, with a relay streaming output into it. job := jm.CreateJob("install", "mpv-leak-test", "MPV") _, relay := newJobCommandRunner(job.Context(), jm, job.ID, 4) for i := 0; i < 10; i++ { relay.output <- fmt.Sprintf("line-%d", i) } relay.errors <- fmt.Errorf("simulated failure detail") relay.Close() relay.Close() // idempotent jm.UpdateProgress(job.ID, 50, "halfway (50%)") jm.CompleteJob(job.ID) // A job that ends in error. jobErr := jm.CreateJob("install", "mpv-leak-test-err", "MPV") jm.SetError(jobErr.ID, fmt.Errorf("boom"), "details") // A job cancelled through both the manager and the job itself. jobCancel := jm.CreateJob("uninstall", "mpv-leak-test-cancel", "MPV") require.NoError(t, jm.CancelJob(jobCancel.ID)) jobCancel.Cancel() jm.SetError(jobCancel.ID, context.Canceled, "worker stopped") // A relay whose job context is already cancelled must still close cleanly. jobCancelledCtx := jm.CreateJob("install", "mpv-leak-test-ctx", "MPV") jobCancelledCtx.Cancel() _, relay2 := newJobCommandRunner(jobCancelledCtx.Context(), jm, jobCancelledCtx.ID, 1) relay2.Close() jm.CompleteJob(jobCancelledCtx.ID) jm.UnregisterClient(clientAll) jm.UnregisterClient(clientFiltered) // Drain events buffered before unregistration; the channels are closed by // UnregisterClient so both ranges terminate. for range clientAll { } for range clientFiltered { } assert.Len(t, jm.GetRecentJobs(), 4) } // readUntilEvent reads SSE lines from body until an "event: " line is // seen, failing the test if none arrives before the timeout. func readUntilEvent(t *testing.T, body *bufio.Reader, name string) { t.Helper() done := make(chan string, 1) go func() { for { line, err := body.ReadString('\n') if err != nil { done <- "err:" + err.Error() return } if strings.TrimSpace(line) == "event: "+name { done <- strings.TrimSpace(line) return } } }() select { case got := <-done: require.Equal(t, "event: "+name, got) case <-time.After(5 * time.Second): t.Fatalf("timed out waiting for SSE event %q", name) } } // TestJobStreamSSEHandlerNoGoroutineLeaks exercises the SSE handler over real // HTTP — connect, initial status, live events, disconnect — and verifies the // handler goroutine and client registration are cleaned up afterwards. func TestJobStreamSSEHandlerNoGoroutineLeaks(t *testing.T) { defer goleak.VerifyNone(t) t.Cleanup(func() { verifyNoRuntimeGoroutineLeaks(t) }) jm := NewJobManager() srv := &Server{jobManager: jm} httpSrv := httptest.NewServer(http.HandlerFunc(srv.handleJobStreamSSE)) defer httpSrv.Close() job := jm.CreateJob("install", "mpv-sse-leak-test", "MPV") clientsRegistered := func(want int) bool { jm.clientsMu.RLock() defer jm.clientsMu.RUnlock() return len(jm.clients) == want } // connectStream opens an SSE connection and returns its cancel func and body. connectStream := func(t *testing.T, url string) (context.CancelFunc, *bufio.Reader, func()) { t.Helper() ctx, cancel := context.WithCancel(context.Background()) req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil) require.NoError(t, err) resp, err := http.DefaultClient.Do(req) require.NoError(t, err) require.Equal(t, http.StatusOK, resp.StatusCode) return cancel, bufio.NewReader(resp.Body), func() { resp.Body.Close() } } // Client 1: filtered to a specific job, receives the initial status event. cancel1, body1, closeBody1 := connectStream(t, httpSrv.URL+"/?job_id="+job.ID) defer closeBody1() require.Eventually(t, func() bool { return clientsRegistered(1) }, 2*time.Second, 5*time.Millisecond) readUntilEvent(t, body1, "status") // Client 2: unfiltered stream. The handler writes nothing until the first // matching event, so response headers only arrive with the broadcast below; // drive the request concurrently to avoid blocking on them. ctx2, cancel2 := context.WithCancel(context.Background()) req2, err := http.NewRequestWithContext(ctx2, http.MethodGet, httpSrv.URL+"/", nil) require.NoError(t, err) type streamResult struct { resp *http.Response err error } respCh := make(chan streamResult, 1) go func() { resp, err := http.DefaultClient.Do(req2) respCh <- streamResult{resp, err} }() require.Eventually(t, func() bool { return clientsRegistered(2) }, 2*time.Second, 5*time.Millisecond) // A progress broadcast reaches both clients (and sends headers to client 2). jm.UpdateProgress(job.ID, 42, "working (42%)") var body2 *bufio.Reader select { case res := <-respCh: require.NoError(t, res.err) require.Equal(t, http.StatusOK, res.resp.StatusCode) defer res.resp.Body.Close() body2 = bufio.NewReader(res.resp.Body) case <-time.After(5 * time.Second): t.Fatal("unfiltered SSE client never received response headers") } readUntilEvent(t, body1, "progress") readUntilEvent(t, body2, "progress") // Disconnect both clients; the handler goroutines must exit and unregister. cancel1() cancel2() require.Eventually(t, func() bool { return clientsRegistered(0) }, 2*time.Second, 5*time.Millisecond) jm.CompleteJob(job.ID) http.DefaultClient.CloseIdleConnections() }