diff --git a/chasm/lib/scheduler/backfill_capacity_progress_test.go b/chasm/lib/scheduler/backfill_capacity_progress_test.go index 5af7124fe60..8a8c5b2e8e4 100644 --- a/chasm/lib/scheduler/backfill_capacity_progress_test.go +++ b/chasm/lib/scheduler/backfill_capacity_progress_test.go @@ -2,7 +2,6 @@ package scheduler_test import ( "fmt" - "os" "testing" "time" @@ -19,10 +18,7 @@ func TestBackfillCapacityNativeControl(t *testing.T) { testBackfillCapacityProgress(t, 450) } -func TestBackfillCapacityCounterexample(t *testing.T) { - if os.Getenv("TEMPORAL_RUN_MIGRATION_COUNTEREXAMPLES") != "1" { - t.Skip("set TEMPORAL_RUN_MIGRATION_COUNTEREXAMPLES=1") - } +func TestBackfillCapacityProgress(t *testing.T) { for _, n := range []int{451, 1000} { t.Run(fmt.Sprint(n), func(t *testing.T) { testBackfillCapacityProgress(t, n) }) } diff --git a/chasm/lib/scheduler/backfiller_capacity_internal_test.go b/chasm/lib/scheduler/backfiller_capacity_internal_test.go index c697aa0019b..9761dbc04b3 100644 --- a/chasm/lib/scheduler/backfiller_capacity_internal_test.go +++ b/chasm/lib/scheduler/backfiller_capacity_internal_test.go @@ -28,7 +28,7 @@ import ( // // pending = max(0, bufferedCount - retainedActionCount) // available = max(0, maxBufferSize/2 - pending - generatorReserve) -// result = available / max(1, backfillerCount) +// result = min(available, max(1, available / max(1, backfillerCount))) func TestBackfillerBufferCapacity(t *testing.T) { cases := []struct { name string @@ -77,7 +77,7 @@ func TestBackfillerBufferCapacity(t *testing.T) { {"single backfiller gets the whole remainder", 0, 0, 1000, 50, 1, 450}, {"ten backfillers each get an even share (regression)", 0, 0, 1000, 50, 10, 45}, {"backfillers exactly dividing remainder get one each", 0, 0, 1000, 50, 450, 1}, - {"more backfillers than remainder truncate to zero", 0, 0, 1000, 50, 451, 0}, + {"more backfillers than remainder still make progress", 0, 0, 1000, 50, 451, 1}, {"zero backfillers are clamped to one", 0, 0, 1000, 50, 0, 450}, {"negative backfillers are clamped to one", 0, 0, 1000, 50, -3, 450}, } diff --git a/chasm/lib/scheduler/backfiller_tasks.go b/chasm/lib/scheduler/backfiller_tasks.go index 082215cbfcc..37ee41d3134 100644 --- a/chasm/lib/scheduler/backfiller_tasks.go +++ b/chasm/lib/scheduler/backfiller_tasks.go @@ -293,5 +293,6 @@ func backfillerBufferCapacity(bufferedCount, retainedActionCount, maxBufferSize, backfillerCount = max(1, backfillerCount) pending := max(0, bufferedCount-retainedActionCount) available := max(0, (maxBufferSize/2)-pending-generatorReserve) - return available / backfillerCount + // Pure tasks serialize on the tree and recompute the shared budget after each enqueue. + return min(available, max(1, available/backfillerCount)) }