forked from temporalio/temporal
-
Notifications
You must be signed in to change notification settings - Fork 0
Reproduce backfiller starvation above shared capacity #51
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Open
Open
Changes from all commits
Commits
Show all changes
2 commits
Select commit
Hold shift + click to select a range
File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,73 @@ | ||
| package scheduler_test | ||
|
|
||
| import ( | ||
| "fmt" | ||
| "os" | ||
| "testing" | ||
| "time" | ||
|
|
||
| "github.com/stretchr/testify/require" | ||
| enumspb "go.temporal.io/api/enums/v1" | ||
| schedulepb "go.temporal.io/api/schedule/v1" | ||
| "go.temporal.io/server/chasm" | ||
| "go.temporal.io/server/chasm/lib/scheduler" | ||
| "go.temporal.io/server/common/clock" | ||
| "google.golang.org/protobuf/types/known/timestamppb" | ||
| ) | ||
|
|
||
| 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") | ||
| } | ||
| for _, n := range []int{451, 1000} { | ||
| t.Run(fmt.Sprint(n), func(t *testing.T) { testBackfillCapacityProgress(t, n) }) | ||
| } | ||
| } | ||
|
|
||
| func testBackfillCapacityProgress(t *testing.T, count int) { | ||
| t.Helper() | ||
| now := time.Date(2026, 9, 4, 12, 0, 0, 0, time.UTC) | ||
| ts := clock.NewEventTimeSource().Update(now) | ||
| spec := defaultSchedule() | ||
| spec.State.Paused = true | ||
| e := newSchedulerTestEngine(t, spec, withEngineTimeSource(ts)) | ||
| require.NoError(t, e.updateScheduler(func(s *scheduler.Scheduler, ctx chasm.MutableContext) error { | ||
| for range count { | ||
| s.NewRangeBackfiller(ctx, &schedulepb.BackfillRequest{ | ||
| StartTime: timestamppb.New(now), EndTime: timestamppb.New(now), | ||
| OverlapPolicy: enumspb.SCHEDULE_OVERLAP_POLICY_ALLOW_ALL, | ||
| }) | ||
| } | ||
| return nil | ||
| })) | ||
| seen := make(map[string]string) | ||
| remaining := count | ||
| for round := 0; round < 50 && remaining > 0; round++ { | ||
| buffered := 0 | ||
| require.NoError(t, e.updateScheduler(func(s *scheduler.Scheduler, ctx chasm.MutableContext) error { | ||
| i := s.Invoker.Get(ctx) | ||
| buffered = len(i.BufferedStarts) | ||
| require.LessOrEqual(t, buffered, 460, "shared half-buffer minus generator reserve plus retained-history allowance") | ||
| for _, start := range i.BufferedStarts { | ||
| require.NotContains(t, seen, start.RequestId) | ||
| seen[start.RequestId] = start.WorkflowId | ||
| } | ||
| i.BufferedStarts = nil | ||
| remaining = len(s.Backfillers) | ||
| return nil | ||
| })) | ||
| if remaining == 0 { | ||
| break | ||
| } | ||
| require.Positive(t, buffered, "seed=capacity-%d: empty buffer and %d ranges must make progress", count, remaining) | ||
| ts.Update(ts.Now().Add(time.Hour)) | ||
| _, err := e.engine.FirePureTasks(e.rootRef, ts.Now()) | ||
| require.NoError(t, err) | ||
| } | ||
| require.Zero(t, remaining) | ||
| require.Len(t, seen, count) | ||
| } | ||
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,34 @@ | ||
| # Backfiller capacity truncation prevents progress | ||
|
|
||
| Base: ad7b2298d. Seeds: `capacity-450`, `capacity-451`, `capacity-1000`. | ||
| Invariant: finite backfills with positive global capacity must eventually enqueue their actions without exceeding the shared budget. | ||
|
|
||
| The component-engine test installs concurrent one-instant ranges on a paused minute schedule. All ranges share their boundary and ALLOW_ALL policy. It runs creation's immediate tasks and repeatedly drains admitted starts, advancing the logical clock to fire persisted continuations. The 450-range native control completes; 451 and 1000 ranges admit zero starts on their first pass despite an empty buffer. The sequence is deterministic; random backfiller UUIDs do not affect the invariant. Distinct request identities are checked on every drain. | ||
|
|
||
| `allowedBufferedStarts` counts all range backfillers. At defaults, `backfillerBufferCapacity` computes `450 / count` for an empty buffer. At 451 the quotient becomes zero. Every task takes the capacity-stalled path, preserving its range and scheduling a backoff. Therefore the divisor never decreases and every subsequent attempt encounters the same state. | ||
|
|
||
| Impact: V1 can carry up to 1000 ongoing backfills into CHASM even though native patch admission normally limits concurrency to 100. Migrated ranges above the arithmetic threshold remain stuck indefinitely. A configurable smaller buffer can also expose this below the default native concurrency limit. | ||
|
|
||
| ```mermaid | ||
| sequenceDiagram | ||
| participant Task as Backfiller task | ||
| participant Invoker | ||
| loop all 451 ranges, indefinitely | ||
| Task->>Invoker: read shared free capacity = 450 | ||
| Task->>Task: 450 / 451 = 0 | ||
| Task->>Task: retain range and reschedule | ||
| end | ||
| ``` | ||
|
|
||
|  | ||
|
|
||
| Control: `go test -tags test_dep ./chasm/lib/scheduler -run '^TestBackfillCapacityNativeControl$' -count=1`. | ||
| Counterexamples: `TEMPORAL_RUN_MIGRATION_COUNTEREXAMPLES=1 go test -tags test_dep ./chasm/lib/scheduler -run '^TestBackfillCapacityCounterexample$' -count=1`. | ||
|
|
||
| The fix admits `min(available, max(1, available/count))`. Pure task mutations serialize on the scheduler tree, so every subsequent task recomputes capacity after the preceding enqueue. Zero global capacity still admits zero. The existing retained-history allowance is preserved, so the test's raw buffer bound is 460, equivalent to 450 pending slots after that allowance. This changes no protobuf or V1 workflow code and adds no new replay branch. | ||
|
|
||
| Upstream audit: all open public PR titles/bodies fetched on 2026-09-04, with scheduler/backfill/migration candidates inspected. No matching production capacity fix found. Imported invoker activation is separately covered by #11557, and the fresh reverse-boundary defect by #11878. | ||
|
|
||
| Ordering is separate from admission: native CHASM range task order is unspecified and reverse conversion iterates a map. Sorting ranges alone would not reproduce a guaranteed native order. Equal boundaries with distinct policies can therefore have order-dependent overlap outcomes; this report does not claim deterministic action equivalence across migration for that hypothesis. Capacity admission neither changes range cursors nor introduces a map-order guarantee. | ||
|
|
||
| Crash/failover reasoning: enqueue, cursor advancement, and continuation task are one existing CHASM transaction. The change introduces no new external side effect or unreplicated state. These tests exercise committed component transactions, not actual History crashes or namespace failover. Counting all backfillers still costs O(N) per task; this fix addresses liveness, not that existing aggregate O(N²) scan cost. |
Loading
Sorry, something went wrong. Reload?
Sorry, we cannot display this file.
Sorry, this file is invalid so it cannot be displayed.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
On the first loop iteration, both preceding
updateSchedulercalls only commit component mutations; they do not execute the immediate pure tasks created byNewRangeBackfiller. Consequentlybufferedis still zero andremainingis stillcount, so this assertion fails even for the always-enabled 450-range control. The firstFirePureTaskscall is below the assertion and is therefore unreachable; fire the initial tasks before the first drain/check.Useful? React with 👍 / 👎.