From 9c4bff75c7fb3c75bb717dfdfa719ff6498c8436 Mon Sep 17 00:00:00 2001 From: Akanksha Trehun Date: Mon, 31 Aug 2026 13:32:37 +0530 Subject: [PATCH] Run cleanup on stage failure so trials don't leak resources runTrial returned as soon as any stage (Prepare/Create/Start/WaitReady/ Stop/Delete) failed, so Cleanup was only ever reached on the all-succeed path. Since these stages manage real containers, VMs, and netns as root, a failed trial could leak host resources indefinitely. Restructure the stage loop to run per-adapter instead of over one flattened list, and once Prepare has succeeded for an adapter, run its Cleanup on a best-effort basis before returning a failure from any later stage. A cleanup failure is logged but does not replace the original stage error in the trial result. Cleanup is skipped if Prepare itself failed, since no resources exist yet to tear down. Added orchestrator tests using a fake in-package Adapter (no containerd required) covering: cleanup runs after each later-stage failure, cleanup is skipped on Prepare failure, a cleanup failure doesn't mask the original error, and the success path still calls Cleanup exactly once. Fixes #9 Signed-off-by: Akanksha Trehun --- internal/orchestrator/orchestrator.go | 38 +++- internal/orchestrator/orchestrator_test.go | 213 +++++++++++++++++++++ 2 files changed, 249 insertions(+), 2 deletions(-) create mode 100644 internal/orchestrator/orchestrator_test.go diff --git a/internal/orchestrator/orchestrator.go b/internal/orchestrator/orchestrator.go index 0d58e12..3cbfddc 100644 --- a/internal/orchestrator/orchestrator.go +++ b/internal/orchestrator/orchestrator.go @@ -3,6 +3,7 @@ package orchestrator import ( "context" "fmt" + "log" "time" "github.com/urunc-dev/evaluation_suite/internal/plan" @@ -147,6 +148,13 @@ func (o *Orchestrator) runTrial( } for i := 0; i < repetitions; i++ { + // prepared tracks whether the Prepare stage has succeeded, i.e. + // whether the adapter may own real resources (containers, tasks, + // netns, etc.) that Cleanup needs to tear down. cleanedUp tracks + // whether Cleanup has already run (as the last regular stage) so we + // don't invoke it a second time on the happy path. + prepared := false + cleanedUp := false for _, stage := range stages { // if the stage's experiment does not match the trial's experiment, skip it @@ -154,12 +162,38 @@ func (o *Orchestrator) runTrial( continue } stageResult, err := stage.fn(ctx, runtimeTC) + result.RuntimeStages = append(result.RuntimeStages, stageResult) + + if stage.name == harnessruntime.StageCleanup { + cleanedUp = err == nil + } + if err != nil { - result.RuntimeStages = append(result.RuntimeStages, stageResult) + // A later stage failed after resources were created in + // Prepare. Run cleanup on a best-effort basis so a failed + // trial doesn't leak containers/VMs/netns on the host. A + // cleanup failure here is logged, not fatal: it must not + // mask the original stage error being returned below. + if prepared && !cleanedUp && stage.name != harnessruntime.StageCleanup { + for _, adapter := range adapters { + if adapter.ExperimentName() != stage.experiment { + continue + } + if _, cleanupErr := adapter.Cleanup(ctx, runtimeTC); cleanupErr != nil { + log.Printf( + "trial %s: cleanup after %s stage failure also failed: %v", + trial.ID, stage.name, cleanupErr, + ) + } + break + } + } return failTrial(result, stage.name, err) } - result.RuntimeStages = append(result.RuntimeStages, stageResult) + if stage.name == harnessruntime.StagePrepare { + prepared = true + } } } diff --git a/internal/orchestrator/orchestrator_test.go b/internal/orchestrator/orchestrator_test.go new file mode 100644 index 0000000..6913425 --- /dev/null +++ b/internal/orchestrator/orchestrator_test.go @@ -0,0 +1,213 @@ +package orchestrator + +import ( + "context" + "errors" + "testing" + + "github.com/urunc-dev/evaluation_suite/internal/plan" + harnessruntime "github.com/urunc-dev/evaluation_suite/internal/runtime" +) + +// fakeAdapter is a minimal harnessruntime.Adapter used to drive runTrial +// without needing a real containerd daemon. Each stage method records that +// it ran and can be configured to fail. +type fakeAdapter struct { + experiment string + + failStage harnessruntime.Stage + failErr error + + cleanupErr error + + calls []harnessruntime.Stage +} + +func (f *fakeAdapter) ExperimentName() string { return f.experiment } + +func (f *fakeAdapter) stage(stage harnessruntime.Stage) (harnessruntime.StageResult, error) { + f.calls = append(f.calls, stage) + if f.failStage == stage { + return harnessruntime.StageResult{Stage: stage}, f.failErr + } + return harnessruntime.StageResult{Stage: stage}, nil +} + +func (f *fakeAdapter) Prepare(ctx context.Context, tc harnessruntime.TrialContext) (harnessruntime.StageResult, error) { + return f.stage(harnessruntime.StagePrepare) +} + +func (f *fakeAdapter) CreateTask(ctx context.Context, tc harnessruntime.TrialContext) (harnessruntime.StageResult, error) { + return f.stage(harnessruntime.StageCreate) +} + +func (f *fakeAdapter) StartTask(ctx context.Context, tc harnessruntime.TrialContext) (harnessruntime.StageResult, error) { + return f.stage(harnessruntime.StageStart) +} + +func (f *fakeAdapter) WaitReady(ctx context.Context, tc harnessruntime.TrialContext) (harnessruntime.StageResult, error) { + return f.stage(harnessruntime.StageWaitReady) +} + +func (f *fakeAdapter) Stop(ctx context.Context, tc harnessruntime.TrialContext) (harnessruntime.StageResult, error) { + return f.stage(harnessruntime.StageStop) +} + +func (f *fakeAdapter) DeleteTask(ctx context.Context, tc harnessruntime.TrialContext) (harnessruntime.StageResult, error) { + return f.stage(harnessruntime.StageDelete) +} + +func (f *fakeAdapter) Cleanup(ctx context.Context, tc harnessruntime.TrialContext) (harnessruntime.StageResult, error) { + f.calls = append(f.calls, harnessruntime.StageCleanup) + if f.failStage == harnessruntime.StageCleanup { + return harnessruntime.StageResult{Stage: harnessruntime.StageCleanup}, f.failErr + } + if f.cleanupErr != nil { + return harnessruntime.StageResult{Stage: harnessruntime.StageCleanup}, f.cleanupErr + } + return harnessruntime.StageResult{Stage: harnessruntime.StageCleanup}, nil +} + +func (f *fakeAdapter) GenerateResult(ctx context.Context, tc harnessruntime.TrialContext, results []harnessruntime.StageResult) (any, error) { + return nil, nil +} + +func countStage(calls []harnessruntime.Stage, stage harnessruntime.Stage) int { + n := 0 + for _, c := range calls { + if c == stage { + n++ + } + } + return n +} + +func trial(experiment string) plan.Trial { + return plan.Trial{ + ID: "trial-1", + ExperimentName: experiment, + RuntimeName: "urunc", + } +} + +// TestRunTrialCleansUpAfterStageFailure verifies that when a stage fails +// after Prepare has already created resources, the orchestrator still calls +// Cleanup on the adapter before returning the (unmasked) original error. +// This is the regression test for issue #9: a failed trial used to leak +// containers/VMs/netns because Cleanup was only reachable via the normal +// (all-stages-succeed) path. +func TestRunTrialCleansUpAfterStageFailure(t *testing.T) { + stagesToFail := []harnessruntime.Stage{ + harnessruntime.StageCreate, + harnessruntime.StageStart, + harnessruntime.StageWaitReady, + harnessruntime.StageStop, + harnessruntime.StageDelete, + } + + for _, failStage := range stagesToFail { + failStage := failStage + t.Run(string(failStage), func(t *testing.T) { + wantErr := errors.New("boom") + adapter := &fakeAdapter{experiment: "lifecycle", failStage: failStage, failErr: wantErr} + + o := New(func(plan.Trial) (harnessruntime.Adapter, error) { return adapter, nil }) + + result := o.runTrial(context.Background(), trial("lifecycle"), Options{}) + + if result.Status != TrialStatusFailed { + t.Fatalf("status = %s, want %s", result.Status, TrialStatusFailed) + } + if result.FailedStage != failStage { + t.Fatalf("failedStage = %s, want %s", result.FailedStage, failStage) + } + if result.Error != wantErr.Error() { + t.Fatalf("error = %q, want %q", result.Error, wantErr.Error()) + } + if countStage(adapter.calls, harnessruntime.StageCleanup) != 1 { + t.Fatalf("Cleanup calls = %d, want 1 (calls: %v)", countStage(adapter.calls, harnessruntime.StageCleanup), adapter.calls) + } + }) + } +} + +// TestRunTrialSkipsCleanupWhenPrepareFails verifies that Cleanup is not +// invoked when Prepare itself fails, since no resources were created yet +// (calling Cleanup in that state can be unsafe for adapters that assume +// Prepare has already initialized adapter state). +func TestRunTrialSkipsCleanupWhenPrepareFails(t *testing.T) { + wantErr := errors.New("prepare failed") + adapter := &fakeAdapter{experiment: "lifecycle", failStage: harnessruntime.StagePrepare, failErr: wantErr} + + o := New(func(plan.Trial) (harnessruntime.Adapter, error) { return adapter, nil }) + + result := o.runTrial(context.Background(), trial("lifecycle"), Options{}) + + if result.Status != TrialStatusFailed { + t.Fatalf("status = %s, want %s", result.Status, TrialStatusFailed) + } + if countStage(adapter.calls, harnessruntime.StageCleanup) != 0 { + t.Fatalf("Cleanup calls = %d, want 0 (calls: %v)", countStage(adapter.calls, harnessruntime.StageCleanup), adapter.calls) + } +} + +// TestRunTrialCleanupFailureDoesNotMaskOriginalError verifies that a +// Cleanup failure is surfaced only via logging: the trial's reported error +// and failed stage must still reflect the original stage failure. +func TestRunTrialCleanupFailureDoesNotMaskOriginalError(t *testing.T) { + wantErr := errors.New("create failed") + adapter := &fakeAdapter{ + experiment: "lifecycle", + failStage: harnessruntime.StageCreate, + failErr: wantErr, + cleanupErr: errors.New("cleanup also failed"), + } + + o := New(func(plan.Trial) (harnessruntime.Adapter, error) { return adapter, nil }) + + result := o.runTrial(context.Background(), trial("lifecycle"), Options{}) + + if result.FailedStage != harnessruntime.StageCreate { + t.Fatalf("failedStage = %s, want %s", result.FailedStage, harnessruntime.StageCreate) + } + if result.Error != wantErr.Error() { + t.Fatalf("error = %q, want original stage error %q (must not be masked by cleanup failure)", result.Error, wantErr.Error()) + } + if countStage(adapter.calls, harnessruntime.StageCleanup) != 1 { + t.Fatalf("Cleanup calls = %d, want 1", countStage(adapter.calls, harnessruntime.StageCleanup)) + } +} + +// TestRunTrialCallsCleanupExactlyOnceOnSuccess verifies the happy path is +// unchanged: Cleanup runs once, as the final regular stage. +func TestRunTrialCallsCleanupExactlyOnceOnSuccess(t *testing.T) { + adapter := &fakeAdapter{experiment: "lifecycle"} + + o := New(func(plan.Trial) (harnessruntime.Adapter, error) { return adapter, nil }) + + result := o.runTrial(context.Background(), trial("lifecycle"), Options{}) + + if result.Status != TrialStatusSuccess { + t.Fatalf("status = %s, want %s", result.Status, TrialStatusSuccess) + } + if countStage(adapter.calls, harnessruntime.StageCleanup) != 1 { + t.Fatalf("Cleanup calls = %d, want 1 (calls: %v)", countStage(adapter.calls, harnessruntime.StageCleanup), adapter.calls) + } +} + +// TestRunTrialSkipsAdapterForOtherExperiments verifies adapters whose +// ExperimentName doesn't match the trial are never invoked. +func TestRunTrialSkipsAdapterForOtherExperiments(t *testing.T) { + adapter := &fakeAdapter{experiment: "storage"} + + o := New(func(plan.Trial) (harnessruntime.Adapter, error) { return adapter, nil }) + + result := o.runTrial(context.Background(), trial("lifecycle"), Options{}) + + if result.Status != TrialStatusSuccess { + t.Fatalf("status = %s, want %s", result.Status, TrialStatusSuccess) + } + if len(adapter.calls) != 0 { + t.Fatalf("calls = %v, want none", adapter.calls) + } +}