Skip to content

shuffle: use blocking reads for gapped-producer replays - #3576

Merged
jgraettinger merged 1 commit into
masterfrom
johnny/shuffle-replay-offset-retry-986219
Oct 6, 2026
Merged

jgraettinger merged 1 commit into
masterfrom
johnny/shuffle-replay-offset-retry-986219

Conversation

@jgraettinger

@jgraettinger jgraettinger commented Oct 4, 2026 •

Copy link
Copy Markdown
Member

Summary

Gapped-producer replays read [offset, end_offset), where end_offset is the begin of a trigger document the main read already saw, so the whole range is known to exist. A broker that lags, though, may not have indexed it yet. A newly-assigned replica serves reads before it has listed a fragment still being persisted by its prior topology.

Replays used non-blocking reads, which fail against such a broker in one of two ways:

  • OFFSET_NOT_YET_AVAILABLE: when the broker's write head is past the replay offset (e.g. it has newer local spool content). classify_read_failure treats this as terminal, so the task fails.
  • A clean EOF short of end_offset: when the broker's write head is exactly where the replay has read to. The gazette client ends a non-blocking read at the write head, so the replay reports "complete" and silently skips [write_head, end_offset).

This switches replays to block: true. The broker waits until it has indexed the fragment, and the read still stops at end_offset.

A replay reads [offset, end_offset), where end_offset is the begin of a
trigger document the main read already observed, so the whole range is
known to exist. But a lagging broker may not yet index it: a newly-assigned
replica serves reads before it has listed a fragment still being persisted
by its prior topology, and reports a write head at the end of the last
fragment it knows.

A non-blocking replay read routed to such a broker either fails with
OFFSET_NOT_YET_AVAILABLE (failing the task), or reads through to the stale
write head and ends cleanly short of end_offset, silently skipping the
remainder of the replay. A blocking read instead waits for the broker to
index the fragment, and still terminates at end_offset.
@strix-security

strix-security Bot commented Oct 4, 2026 •

Copy link
Copy Markdown

Strix Security Review

No security issues found.

Review summary

Reviewed the single-file diff in crates/shuffle/src/slice/replay.rs, which changes a bounded gapped-producer historical replay read from non-blocking to blocking (block: false → block: true) to prevent silent data skips and premature failures against lagging brokers. The change introduces no untrusted input, authorization, authentication, secret, or injection surface. The blocking read is bounded by end_offset, the range is guaranteed to exist, and the Gazette client surfaces stream/status errors through classify_read_failure for retry or terminal handling. No security issues were identified in the changed code.

Updated for f24c164.


Reviewed by Strix
Re-run review · Configure security review settings

@jgraettinger
jgraettinger requested a review from a team October 4, 2026 23:52

@dgreer-dev dgreer-dev left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM.

Do you think there's any value in a regression test where the replay hits a stale write head, waits for a missing fragment to become available, and then verifies that all replayed documents are delivered exactly once? Should be red with the old non-blocking read and green with this change.

@jgraettinger

jgraettinger commented Oct 6, 2026 •

Copy link
Copy Markdown
Member Author

Do you think there's any value in a regression test

Not especially. I think it would be a lot of test SLOC, and would predominantly be testing Gazette and client behavior -- behavior that's well covered under main-path reads. This change is fundamentally just making replay reads less special and like main-path blocking reads, but with an end offset, and we have existing coverage that the end-offset is respected.

@jgraettinger
jgraettinger merged commit 90caf25 into master Oct 6, 2026
12 checks passed
@jgraettinger
jgraettinger deleted the johnny/shuffle-replay-offset-retry-986219 branch October 6, 2026 14:34
@github-actions github-actions Bot added pending:agent Merged, in the control-plane-agent image, and not yet rolled to flow-agent pending:agent-api Merged, ships via Deploy agent-api, and not yet deployed pending:flowctl Merged, changes the flowctl binary, and not in a published release labels Oct 6, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

pending:agent Merged, in the control-plane-agent image, and not yet rolled to flow-agent pending:agent-api Merged, ships via Deploy agent-api, and not yet deployed pending:flowctl Merged, changes the flowctl binary, and not in a published release

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants