From 7c4230908115a6e5c606f35b621614d915ca7055 Mon Sep 17 00:00:00 2001 From: Lev Brouk Date: Thu, 2 Jul 2026 08:06:00 -0700 Subject: [PATCH] fix(orchestrator): compute chunk readiness in the fetch session's frame MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit When a sandbox resumes, chunks are fetched from a peer orchestrator while the build is still uploading (4MB uncompressed chunks) and from object storage after it finalizes (smaller compressed frames). A fetch session started under the old peer geometry can still be in flight when a read locates its chunk under the new frame table, and getOrCreateSession reuses that session. Chunker.fetch computed the "already streamed" threshold from the newly located chunk's offset, while the session's bytesReady watermark counts from the session's own offset. With the origins mixed, a session that had streamed only its first blocks could report a deeper, unwritten block as ready, and the read served zeroed mmap to the guest as valid data — silent memory corruption on resume (guest segfaults and kernel panics), specific to compression+P2P, ~2.5% of resumes in a rapid pause/resume loop. Compute the readiness threshold and the wait-loop bound in the session's frame, and add a regression test reproducing the shifted-origin reuse. --- .../pkg/sandbox/block/streaming_chunk.go | 10 +++- .../pkg/sandbox/block/streaming_chunk_test.go | 48 +++++++++++++++++++ 2 files changed, 56 insertions(+), 2 deletions(-) diff --git a/packages/orchestrator/pkg/sandbox/block/streaming_chunk.go b/packages/orchestrator/pkg/sandbox/block/streaming_chunk.go index d92b4233b7..4049ef2e70 100644 --- a/packages/orchestrator/pkg/sandbox/block/streaming_chunk.go +++ b/packages/orchestrator/pkg/sandbox/block/streaming_chunk.go @@ -158,6 +158,11 @@ func (c *Chunker) getOrCreateSession(ctx context.Context, off, length int64, ups // for every block the range spans (a span can cross block boundaries // after dedup; waiting only on the start block leaves the tail unfetched). func (c *Chunker) fetch(ctx context.Context, off, length int64, upstream storage.RangeOpener, ft *storage.FrameTable) (storage.Source, error) { + // (upstream, ft) is a per-call snapshot, so this read can be served by a + // session created before a peer→storage switch and framed on a different + // geometry (4MB uncompressed peer chunks vs smaller compressed frames). + // bytesReady counts from the session's chunkOff, so all readiness math + // below is in the session's frame. chunkOff, chunkLen, err := c.locateChunk(off, ft) if err != nil { return storage.UnknownSource, fmt.Errorf("failed to locate chunk for offset %d: %w", off, err) @@ -171,10 +176,11 @@ func (c *Chunker) fetch(ctx context.Context, off, length int64, upstream storage blockSize := c.cache.BlockSize() startBlock := (off / blockSize) * blockSize endBlock := ((off + length - 1) / blockSize) * blockSize - chunkEnd := chunkOff + chunkLen + + chunkEnd := session.chunkOff + session.chunkLen // Already streamed past every byte we need: it's in the mmap, source=mmap. - endByte := min(endBlock+blockSize, chunkEnd) - chunkOff + endByte := min(endBlock+blockSize, chunkEnd) - session.chunkOff if session.bytesReady.Load() >= endByte { return storage.SourceMmap, nil } diff --git a/packages/orchestrator/pkg/sandbox/block/streaming_chunk_test.go b/packages/orchestrator/pkg/sandbox/block/streaming_chunk_test.go index f536286634..fcd09e48dd 100644 --- a/packages/orchestrator/pkg/sandbox/block/streaming_chunk_test.go +++ b/packages/orchestrator/pkg/sandbox/block/streaming_chunk_test.go @@ -5,6 +5,7 @@ package block import ( "bytes" "context" + "errors" "fmt" "io" "math/rand/v2" @@ -337,6 +338,53 @@ func TestChunker_EarlyReturn(t *testing.T) { require.Equal(t, data[lastOff:lastOff+testBlockSize], r.data) } +// TestChunker_SessionOriginMismatch reproduces a corruption bug in fetch's +// readiness fast path. A fetch session created on one frame geometry (an +// uncompressed whole-file chunk) is reused by a later read that locates a +// different origin — as happens when a peer→storage transition swaps 4MB +// uncompressed peer chunks for smaller compressed frames, so getOrCreateSession +// returns a session framed on the old geometry. bytesReady counts from the +// session's chunkOff; the readiness check must use the same origin. The buggy +// code computed the threshold from the freshly located chunkOff, so a session +// that filled only its first block reported a not-yet-written deeper block as +// ready — serving stale (zero) mmap to the guest instead of fetching it. +// +// The scenario is built synchronously: a partially filled, then terminated, +// session stands in for the aborted peer fetch. With the fix, the shifted-origin +// read finds the block unready and falls through to the terminated session's +// error; with the bug, it short-circuits to SourceMmap and returns no error. +func TestChunker_SessionOriginMismatch(t *testing.T) { + t.Parallel() + + data := makeTestData(testFileSize) + // A compressed frame table with testFrameSize-aligned uncompressed frames: + // locateChunk for the second frame yields chunkOff=testFrameSize, an origin + // different from the uncompressed session's chunkOff of 0. + compFT, upstream := makeCompressedTestData(t, data) + shiftedOff := int64(testFrameSize) + + chunker := newTestChunker(t, int64(len(data))) + defer chunker.Close() + + // An in-flight session on the uncompressed whole-file geometry that filled + // only its first block, then terminated (as a peer fetch does when the peer + // goes away mid-transition). bytesReady covers block 0 but not shiftedOff. + sess := newFetchSession(0, int64(len(data)), chunker.cache) + sess.advance(testBlockSize) + sess.fail(errors.New("peer gone")) + + chunker.fetchMu.Lock() + chunker.fetchSessions = append(chunker.fetchSessions, sess) + chunker.fetchMu.Unlock() + + // The shifted-origin read reuses the session. shiftedOff is not ready, so + // fetch must consult the (terminated) session and surface its error rather + // than report the block ready and serve stale mmap. + _, err := chunker.fetch(t.Context(), shiftedOff, testBlockSize, upstream, compFT) + require.Error(t, err, "fetch reported an unfilled shifted-origin block as ready (served stale mmap)") + require.ErrorContains(t, err, "peer gone") +} + // TestChunker_ErrorKeepsPartialData verifies that an upstream error at the // midpoint of a chunk still allows data before the error to be served. func TestChunker_ErrorKeepsPartialData(t *testing.T) {