Clean up comparison test workers
This commit is contained in:
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user