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
10 changes: 8 additions & 2 deletions packages/orchestrator/pkg/sandbox/block/streaming_chunk.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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
}
Expand Down
48 changes: 48 additions & 0 deletions packages/orchestrator/pkg/sandbox/block/streaming_chunk_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ package block
import (
"bytes"
"context"
"errors"
"fmt"
"io"
"math/rand/v2"
Expand Down Expand Up @@ -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) {
Expand Down
Loading