Skip to content
Closed
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
62 changes: 36 additions & 26 deletions go/internal/runner/run_sweep_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand All @@ -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 {
Expand All @@ -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)
Expand All @@ -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",
Expand All @@ -353,43 +357,49 @@ 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 {
case <-engine.removed:
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.
Expand Down
Loading