From 57aa27c9dec482402085a3a5b135c86f84715048 Mon Sep 17 00:00:00 2001 From: Eric Rakestraw Date: Thu, 13 Aug 2026 03:34:08 +0000 Subject: [PATCH] Clean up comparison test workers --- internal/app/comparison_execution_test.go | 61 +++++++++++++---------- 1 file changed, 36 insertions(+), 25 deletions(-) diff --git a/internal/app/comparison_execution_test.go b/internal/app/comparison_execution_test.go index b44f1c9..d792385 100644 --- a/internal/app/comparison_execution_test.go +++ b/internal/app/comparison_execution_test.go @@ -20,12 +20,9 @@ func TestExecuteComparisonProfilesRunsOrderedProfilesConcurrently(t *testing.T) prepared, prompt := preparedDailyProfile(t) profiles := comparisonProfiles(10) executor := newBarrierExecutor(profiles) - results := make(chan comparisonExecutionResult, 1) - go func() { - results <- executeComparisonProfiles(context.Background(), comparisonExecutionRequest{ - Prepared: prepared, Inspection: comparisonInspection(prompt, profiles), ComparisonID: "comparison_daily", Executor: executor, - }) - }() + results := startComparisonExecution(t, context.Background(), comparisonExecutionRequest{ + Prepared: prepared, Inspection: comparisonInspection(prompt, profiles), ComparisonID: "comparison_daily", Executor: executor, + }, executor) waitForProfileStarts(t, executor, profiles, results) if executor.maximumInFlight() < 2 { t.Fatalf("maximum in-flight executions = %d, want overlap", executor.maximumInFlight()) @@ -58,12 +55,9 @@ func TestExecuteComparisonProfilesContinuesAfterProfileFailure(t *testing.T) { profiles := comparisonProfiles(3) executor := newBarrierExecutor(profiles) executor.setError(profiles[1].ProfileID, errors.New("provider response body must not escape")) - results := make(chan comparisonExecutionResult, 1) - go func() { - results <- executeComparisonProfiles(context.Background(), comparisonExecutionRequest{ - Prepared: prepared, Inspection: comparisonInspection(prompt, profiles), ComparisonID: "comparison_daily", Executor: executor, - }) - }() + results := startComparisonExecution(t, context.Background(), comparisonExecutionRequest{ + Prepared: prepared, Inspection: comparisonInspection(prompt, profiles), ComparisonID: "comparison_daily", Executor: executor, + }, executor) waitForProfileStarts(t, executor, profiles, results) for _, profile := range profiles { executor.release(profile.ProfileID) @@ -84,12 +78,9 @@ func TestExecuteComparisonProfilesPropagatesCancellationAndJoins(t *testing.T) { executor := newBarrierExecutor(profiles) ctx, cancel := context.WithCancel(context.Background()) defer cancel() - results := make(chan comparisonExecutionResult, 1) - go func() { - results <- executeComparisonProfiles(ctx, comparisonExecutionRequest{ - Prepared: prepared, Inspection: comparisonInspection(prompt, profiles), ComparisonID: "comparison_daily", Executor: executor, - }) - }() + results := startComparisonExecution(t, ctx, comparisonExecutionRequest{ + Prepared: prepared, Inspection: comparisonInspection(prompt, profiles), ComparisonID: "comparison_daily", Executor: executor, + }, executor) waitForProfileStarts(t, executor, profiles, results) cancel() result := <-results @@ -114,12 +105,9 @@ func TestExecuteComparisonProfilesUsesDistinctDeterministicDebugReferences(t *te t.Fatalf("NewPromptDebugWriter() error = %v", err) } executor := newBarrierExecutor(profiles) - results := make(chan comparisonExecutionResult, 1) - go func() { - results <- executeComparisonProfiles(context.Background(), comparisonExecutionRequest{ - Prepared: prepared, Inspection: comparisonInspection(prompt, profiles), ComparisonID: "comparison_daily", DebugWriter: debugWriter, Executor: executor, - }) - }() + results := startComparisonExecution(t, context.Background(), comparisonExecutionRequest{ + Prepared: prepared, Inspection: comparisonInspection(prompt, profiles), ComparisonID: "comparison_daily", DebugWriter: debugWriter, Executor: executor, + }, executor) waitForProfileStarts(t, executor, profiles, results) for _, profile := range profiles { executor.release(profile.ProfileID) @@ -254,9 +242,32 @@ func (e *barrierExecutor) inFlightCount() int { return e.inFlight } +const comparisonExecutionTestTimeout = 5 * time.Second + +func startComparisonExecution(t *testing.T, ctx context.Context, request comparisonExecutionRequest, executor *barrierExecutor) <-chan comparisonExecutionResult { + t.Helper() + results := make(chan comparisonExecutionResult, 1) + finished := make(chan struct{}) + t.Cleanup(func() { + executor.releaseAll() + timeout := time.NewTimer(comparisonExecutionTestTimeout) + defer timeout.Stop() + select { + case <-finished: + case <-timeout.C: + t.Error("comparison execution workers did not finish after release") + } + }) + go func() { + defer close(finished) + results <- executeComparisonProfiles(ctx, request) + }() + return results +} + func waitForProfileStarts(t *testing.T, executor *barrierExecutor, profiles []ComparisonProfileInspection, results <-chan comparisonExecutionResult) { t.Helper() - timeout := time.NewTimer(5 * time.Second) + timeout := time.NewTimer(comparisonExecutionTestTimeout) defer timeout.Stop() seen := map[string]struct{}{} for range profiles {