Skip to content
Merged
Show file tree
Hide file tree
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
16 changes: 16 additions & 0 deletions internal/execution/fanout.go
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,15 @@ func FanOut[T any](
return FanOutWithConcurrency(ctx, accountIDs, fn, DefaultMaxConcurrency)
}

// semBlockedHookKey is the context key for an optional test-only
// synchronization hook. When a func() is stored under this key,
// FanOutWithConcurrency invokes it each time an item is recorded as canceled at
// the semaphore-acquire boundary (the <-ctx.Done() branch below). It exists
// purely so tests can synchronize deterministically with that moment; no
// production caller sets it, so the branch is inert (a single context lookup on
// the already-canceled path) outside tests.
type semBlockedHookKey struct{}

// FanOutWithConcurrency is like FanOut but with an explicit concurrency limit.
func FanOutWithConcurrency[T any](
ctx context.Context,
Expand Down Expand Up @@ -82,6 +91,13 @@ func FanOutWithConcurrency[T any](
case <-ctx.Done():
results[i] = Result[T]{AccountID: id, Err: ctx.Err()}
wg.Done()
// Optional test-only synchronization hook (unset in production):
// signals that a queued item was recorded as canceled at the
// semaphore boundary, letting tests order a subsequent unblock
// deterministically. See semBlockedHookKey.
if hook, ok := ctx.Value(semBlockedHookKey{}).(func()); ok {
hook()
}
continue
}
go func(idx int, accountID string) {
Expand Down
35 changes: 33 additions & 2 deletions internal/execution/fanout_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import (
"errors"
"sort"
"testing"
"time"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
Expand Down Expand Up @@ -105,7 +106,20 @@ func TestFanOut_ContextCancelled_BlockedSemaphore(t *testing.T) {
started := make(chan struct{}, 1)
release := make(chan struct{}, 1)

ctx, cancel := context.WithCancel(context.Background())
// canceledRecorded is closed by the fan-out launch loop the moment the
// queued "second" item is recorded as ctx.Canceled at the semaphore
// boundary (via the semBlockedHookKey hook). This gives a deterministic
// happens-before edge: the test only releases the blocker AFTER the
// cancellation has been recorded, so the freed semaphore slot can never
// race the launch loop into running "second". Without it, cancel() and the
// slot-release both make the loop's select ready, and Go's uniform-random
// pick would intermittently launch "second" (recording a nil error).
canceledRecorded := make(chan struct{})
ctx, cancel := context.WithCancel(
context.WithValue(context.Background(), semBlockedHookKey{}, func() {
close(canceledRecorded)
}),
)
defer cancel()

// Collect results in a separate goroutine so this test goroutine can
Expand All @@ -130,9 +144,26 @@ func TestFanOut_ContextCancelled_BlockedSemaphore(t *testing.T) {
done <- outcome{res}
}()

// Wait until the first goroutine holds the semaphore, then cancel.
// Wait until the first goroutine holds the semaphore, then cancel. Once
// canceledRecorded fires, the loop has committed "second" to the
// <-ctx.Done() branch, so releasing the blocker is now race-free.
<-started
cancel()
select {
case <-canceledRecorded:
// Correct behavior: the loop honored ctx cancellation at the
// semaphore boundary. Fires within microseconds, so the deadline
// below is never reached on a healthy tree (no flakiness).
case <-time.After(10 * time.Second):
// A regression to unconditional `sem <- struct{}{}` would block the
// launch loop on the second item forever (the blocker never releases
// because we never send release), so canceledRecorded never fires.
// Fail fast with a clear message instead of hanging for the whole
// go-test timeout.
t.Fatal("timed out waiting for the queued item to be recorded as " +
"canceled; FanOutWithConcurrency likely blocked on the semaphore " +
"instead of honoring ctx cancellation")
}
release <- struct{}{} // let the first goroutine finish so FanOut can return

o := <-done
Expand Down
Loading