diff --git a/go/internal/runner/run_sweep_test.go b/go/internal/runner/run_sweep_test.go index 36d9a786f..fc0493013 100644 --- a/go/internal/runner/run_sweep_test.go +++ b/go/internal/runner/run_sweep_test.go @@ -293,9 +293,10 @@ func TestRunReturnsNilWhenSweepCancelsContext(t *testing.T) { type blockingRemoveTestEngine struct { *pipeRuntime - listed chan struct{} - entered chan struct{} - removed chan struct{} + listed chan struct{} + entered chan struct{} + removed chan struct{} + removeErr error // the ctx error Remove unwound on; read after removed closes } func (e *blockingRemoveTestEngine) ListByOwner(context.Context, string, string) ([]runtime.WorkloadID, error) { @@ -306,8 +307,9 @@ func (e *blockingRemoveTestEngine) ListByOwner(context.Context, string, string) func (e *blockingRemoveTestEngine) Remove(ctx context.Context, _ runtime.WorkloadID) error { close(e.entered) <-ctx.Done() + e.removeErr = ctx.Err() close(e.removed) - return ctx.Err() + return e.removeErr } type sessionsStartedTestHandler struct { @@ -320,9 +322,9 @@ func (h *sessionsStartedTestHandler) Sessions(context.Context, *connect.BidiStre return nil } -func TestRunStartsSessionsBeforeStaleSweepDeadline(t *testing.T) { +func TestRunStartsSessionsWhenStaleSweepTimesOut(t *testing.T) { prev := staleContainerSweepTimeout - staleContainerSweepTimeout = 2 * time.Second + staleContainerSweepTimeout = 100 * time.Millisecond t.Cleanup(func() { staleContainerSweepTimeout = prev }) handler := &sessionsStartedTestHandler{sessions: make(chan struct{})} path, service := compassv1internalconnect.NewRunnerServiceHandler(handler) @@ -340,11 +342,13 @@ func TestRunStartsSessionsBeforeStaleSweepDeadline(t *testing.T) { removed: make(chan struct{}), } ctx, cancel := context.WithCancel(t.Context()) - done := make(chan error, 1) + runDone := make(chan struct{}) + var runErr error runtimeDir := shortRuntimeDir(t) httpClient := h2cHTTPClient(t) go func() { - done <- Run(ctx, RunnerConfig{ + defer close(runDone) + runErr = Run(ctx, RunnerConfig{ RunnerID: "runner-1", ServerAddr: server.URL, Token: "tok", @@ -353,27 +357,24 @@ func TestRunStartsSessionsBeforeStaleSweepDeadline(t *testing.T) { HTTPClient: httpClient, }, nil, discardLoggerRunner()) }() - - select { - case <-handler.sessions: - case <-time.After(staleContainerSweepTimeout + time.Second): + // Registered after the timeout restore so it runs first: Run is joined + // before the var and runtime dir it reads are torn down. + t.Cleanup(func() { cancel() select { - case <-engine.removed: + case <-runDone: case <-time.After(testTimeout): - t.Fatal("stale cleanup did not unwind after cancellation") + t.Error("Run did not return after cancellation") } - select { - case <-done: - case <-time.After(testTimeout): - t.Fatal("Run did not return after cancellation") - } - t.Fatalf("Sessions did not start within %s while stale cleanup was blocked", staleContainerSweepTimeout+time.Second) - } + }) + + // Event-gated, no wall-clock bound asserted: the blocked Remove unwinds on the + // sweep's own deadline, then Sessions must start without any cancel. select { case <-engine.entered: + case <-runDone: + t.Fatalf("Run returned %v before the stale Remove started", runErr) case <-time.After(testTimeout): - cancel() t.Fatal("stale Remove did not start") } select { @@ -381,15 +382,24 @@ func TestRunStartsSessionsBeforeStaleSweepDeadline(t *testing.T) { case <-time.After(testTimeout): t.Fatal("stale cleanup did not stop at its deadline") } + if !errors.Is(engine.removeErr, context.DeadlineExceeded) { + t.Fatalf("stale Remove unwound on %v, want the sweep deadline", engine.removeErr) + } + select { + case <-handler.sessions: + case <-runDone: + t.Fatalf("Run returned %v before Sessions started", runErr) + case <-time.After(testTimeout): + t.Fatal("Sessions did not start after the stale sweep deadline") + } select { - case err := <-done: - if err != nil { - t.Fatalf("Run after bounded startup sweep = %v, want nil", err) + case <-runDone: + if runErr != nil { + t.Fatalf("Run after bounded startup sweep = %v, want nil", runErr) } case <-time.After(testTimeout): t.Fatal("Run did not return after Sessions completed") } - cancel() } // The sweep reaches container backends only through a type assertion.