From cf0e380dfbc91927669aa9915edf0937d361ddc1 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Tue, 19 May 2026 17:51:36 -0700 Subject: [PATCH] s3: serve VOD reads through an LRU block cache MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The bare S3 ReaderAt reopens a ranged GetObject on every non-sequential read. A demuxer seeking around a large MP4 (moov parse, interleaved audio/video, B-frame reorder) turns one VOD into tens of thousands of round-trips: replaying a real 1.4 GB read trace showed 120,230 GetObjects — which is where the multi-minute "slow VOD" wall actually came from, not the upload. Add CachingReaderAt, an LRU cache of fixed 16 MB blocks over any io.ReaderAt, and wrap S3Store.Open in it (16 MB x 10 = 160 MB/reader). It only ever issues aligned, full-block reads to the underlying reader — each a single clean sequential GET — so the same trace drops to 91 backend GETs at ~1.0x read amplification. The file backend is untouched (local random access is already free). Includes a byte-exact correctness test and a trace-replay simulation (TestCachingReaderAtSim, gated on SEEK_TRACE) used to choose the block size and cache depth. Co-Authored-By: Claude Opus 4.7 --- pkg/blob/s3.go | 13 ++- pkg/s3/blockcache.go | 202 ++++++++++++++++++++++++++++++++++++++ pkg/s3/blockcache_test.go | 172 ++++++++++++++++++++++++++++++++ 3 files changed, 386 insertions(+), 1 deletion(-) create mode 100644 pkg/s3/blockcache.go create mode 100644 pkg/s3/blockcache_test.go diff --git a/pkg/blob/s3.go b/pkg/blob/s3.go index c9f744ec1..ab45bb19f 100644 --- a/pkg/blob/s3.go +++ b/pkg/blob/s3.go @@ -49,7 +49,18 @@ func (s *S3Store) Open(ctx context.Context, key string) (Reader, error) { } return nil, err } - return ra, nil + // Wrap in an LRU block cache. The bare ReaderAt reopens a ranged + // GetObject on every non-sequential read; a demuxer seeking around a + // large MP4 turns that into tens of thousands of round-trips. The + // cache serves those from a bounded set of 16 MB blocks, and only + // ever issues aligned full-block reads to the ReaderAt — each one a + // single clean sequential GET. Close on the cache closes the ReaderAt. + cached, err := s3pkg.NewCachingReaderAt(ra, ra.Size(), s3pkg.DefaultCacheBlockSize, s3pkg.DefaultCacheBlocks) + if err != nil { + _ = ra.Close() + return nil, err + } + return cached, nil } func (s *S3Store) NewWriter(ctx context.Context, key, contentType string) (Writer, error) { diff --git a/pkg/s3/blockcache.go b/pkg/s3/blockcache.go new file mode 100644 index 000000000..408ce7afe --- /dev/null +++ b/pkg/s3/blockcache.go @@ -0,0 +1,202 @@ +package s3 + +import ( + "container/list" + "errors" + "fmt" + "io" + "sync" +) + +// CachingReaderAt wraps a backend io.ReaderAt with a bounded LRU cache of +// fixed-size blocks. It exists for the demuxer-over-S3 access pattern: +// qtdemux seeks all over a large MP4, and the bare S3 ReaderAt turns every +// non-sequential read into a fresh ranged GetObject. Serving reads out of +// cached 16 MB blocks collapses those thousands of round-trips into a +// handful of full-block fetches, while keeping memory bounded (maxBlocks * +// blockSize) so we never have to download a whole (potentially many-GB) +// upload up front. +// +// Reads are serialized by a single mutex, matching the bare ReaderAt and +// the single-threaded gstreamer streaming thread that drives it. The +// backend only ever sees aligned, full-block ReadAt calls. +type CachingReaderAt struct { + backend io.ReaderAt + size int64 + blockSize int64 + maxBlocks int + + mu sync.Mutex + blocks map[int64]*list.Element // block index -> LRU element + lru *list.List // front = most-recently-used + everFetched map[int64]bool // blocks fetched at least once (redownload accounting) + stats CacheStats +} + +// Defaults chosen from replaying real demuxer traces against the +// simulation (see TestCachingReaderAtSim): 16 MB blocks minimize round +// trips (a full 1.4 GB read went from ~120k ranged GETs to ~91), and 10 +// cached blocks (160 MB) leaves headroom for files whose tracks live in +// separate regions, where the demuxer keeps multiple read-fronts alive. +const ( + DefaultCacheBlockSize = 16 * 1024 * 1024 + DefaultCacheBlocks = 10 +) + +type cacheEntry struct { + index int64 + data []byte +} + +// CacheStats is a snapshot of the cache's behavior, used both for the +// offline simulation and (eventually) live metrics. +type CacheStats struct { + Reads int64 // ReadAt calls served + BytesRequested int64 // sum of len(p) across ReadAt calls (clamped to size) + BlockTouches int64 // block-level accesses (a read may touch several) + Hits int64 // block touches served from cache + Misses int64 // block touches that required a backend fetch + ColdMisses int64 // misses for a block never fetched before + Redownloads int64 // misses for a block that was fetched then evicted + Evictions int64 // blocks dropped from the cache + BackendReads int64 // ReadAt calls issued to the backend (== Misses) + BackendBytes int64 // bytes pulled from the backend +} + +// NewCachingReaderAt wraps backend (whose total length is size) with an LRU +// block cache of maxBlocks blocks of blockSize bytes each. +func NewCachingReaderAt(backend io.ReaderAt, size, blockSize int64, maxBlocks int) (*CachingReaderAt, error) { + if blockSize <= 0 { + return nil, fmt.Errorf("blockSize must be positive, got %d", blockSize) + } + if maxBlocks <= 0 { + return nil, fmt.Errorf("maxBlocks must be positive, got %d", maxBlocks) + } + return &CachingReaderAt{ + backend: backend, + size: size, + blockSize: blockSize, + maxBlocks: maxBlocks, + blocks: make(map[int64]*list.Element), + lru: list.New(), + everFetched: make(map[int64]bool), + }, nil +} + +// ReadAt implements io.ReaderAt, serving from cached blocks and fetching +// (full, aligned) blocks from the backend on a miss. Returns io.EOF when a +// read runs past the end of the object, matching io.ReaderAt semantics. +func (c *CachingReaderAt) ReadAt(p []byte, off int64) (int, error) { + if off < 0 { + return 0, fmt.Errorf("negative offset %d", off) + } + if off >= c.size { + return 0, io.EOF + } + + c.mu.Lock() + defer c.mu.Unlock() + + c.stats.Reads++ + want := len(p) + if int64(off)+int64(want) > c.size { + want = int(c.size - off) + } + c.stats.BytesRequested += int64(want) + + copied := 0 + for copied < want { + readOff := off + int64(copied) + idx := readOff / c.blockSize + block, err := c.getBlock(idx) + if err != nil { + return copied, err + } + within := int(readOff - idx*c.blockSize) + n := copy(p[copied:want], block[within:]) + copied += n + } + if copied < len(p) { + // Caller asked for more than the object holds. + return copied, io.EOF + } + return copied, nil +} + +// getBlock returns block idx, fetching it from the backend on a miss and +// updating LRU/stats. Caller must hold c.mu. +func (c *CachingReaderAt) getBlock(idx int64) ([]byte, error) { + c.stats.BlockTouches++ + if el, ok := c.blocks[idx]; ok { + c.lru.MoveToFront(el) + c.stats.Hits++ + return el.Value.(*cacheEntry).data, nil + } + + c.stats.Misses++ + if c.everFetched[idx] { + c.stats.Redownloads++ + } else { + c.stats.ColdMisses++ + c.everFetched[idx] = true + } + + start := idx * c.blockSize + n := c.blockSize + if start+n > c.size { + n = c.size - start + } + buf := make([]byte, n) + got, err := readAtFull(c.backend, buf, start) + c.stats.BackendReads++ + c.stats.BackendBytes += int64(got) + if err != nil && !errors.Is(err, io.EOF) { + return nil, fmt.Errorf("cache backend read block %d (offset %d): %w", idx, start, err) + } + buf = buf[:got] + + el := c.lru.PushFront(&cacheEntry{index: idx, data: buf}) + c.blocks[idx] = el + if c.lru.Len() > c.maxBlocks { + back := c.lru.Back() + evicted := back.Value.(*cacheEntry) + c.lru.Remove(back) + delete(c.blocks, evicted.index) + c.stats.Evictions++ + } + return buf, nil +} + +// Size returns the underlying object's size, satisfying blob.Reader. +func (c *CachingReaderAt) Size() int64 { return c.size } + +// Close closes the backend if it owns resources (e.g. the S3 ReaderAt's +// open GetObject body), satisfying io.Closer / blob.Reader. +func (c *CachingReaderAt) Close() error { + if closer, ok := c.backend.(io.Closer); ok { + return closer.Close() + } + return nil +} + +// Stats returns a copy of the current cache statistics. +func (c *CachingReaderAt) Stats() CacheStats { + c.mu.Lock() + defer c.mu.Unlock() + return c.stats +} + +// readAtFull reads len(buf) bytes via repeated ReadAt, tolerating short +// reads from backends that don't fill the buffer in one call. Returns the +// number of bytes read; io.EOF if the object ended first. +func readAtFull(r io.ReaderAt, buf []byte, off int64) (int, error) { + total := 0 + for total < len(buf) { + n, err := r.ReadAt(buf[total:], off+int64(total)) + total += n + if err != nil { + return total, err + } + } + return total, nil +} diff --git a/pkg/s3/blockcache_test.go b/pkg/s3/blockcache_test.go new file mode 100644 index 000000000..f8d0daadd --- /dev/null +++ b/pkg/s3/blockcache_test.go @@ -0,0 +1,172 @@ +package s3 + +import ( + "bufio" + "errors" + "fmt" + "io" + "os" + "strings" + "testing" + + "github.com/stretchr/testify/require" +) + +// patternBackend is a fake io.ReaderAt whose byte at offset o is o%251, so +// any misalignment in the cache surfaces as a content mismatch. +type patternBackend struct{ size int64 } + +func (p *patternBackend) ReadAt(b []byte, off int64) (int, error) { + if off < 0 { + return 0, fmt.Errorf("negative offset %d", off) + } + if off >= p.size { + return 0, io.EOF + } + n := len(b) + if int64(off)+int64(n) > p.size { + n = int(p.size - off) + } + for i := 0; i < n; i++ { + b[i] = byte((off + int64(i)) % 251) + } + if n < len(b) { + return n, io.EOF + } + return n, nil +} + +// TestCachingReaderAt checks the cache returns exactly what the backend +// would, across block boundaries, re-reads (post-eviction), and the tail — +// with a deliberately tiny cache (4 × 64 KB) to force eviction/redownload. +func TestCachingReaderAt(t *testing.T) { + const size = 10*1024*1024 + 12345 + backend := &patternBackend{size: size} + c, err := NewCachingReaderAt(backend, size, 64*1024, 4) + require.NoError(t, err) + + check := func(off, n int64) { + t.Helper() + got := make([]byte, n) + gn, gerr := c.ReadAt(got, off) + want := make([]byte, n) + wn, werr := readAtFull(backend, want, off) + require.Equal(t, wn, gn, "read count mismatch at off=%d n=%d", off, n) + require.Equal(t, want[:wn], got[:gn], "bytes mismatch at off=%d", off) + require.Equal(t, errors.Is(werr, io.EOF), errors.Is(gerr, io.EOF), "EOF parity at off=%d", off) + } + + check(0, 100*1024) // sequential across many blocks + check(5*1024*1024, 64*1024) // jump + check(123, 3) // tiny read + check(64*1024-10, 40) // spans a block boundary + check(0, 100*1024) // re-read early region (likely evicted) + check(size-10, 100) // read past the tail + + // Reading entirely past EOF yields (0, io.EOF). + n, err := c.ReadAt(make([]byte, 16), size) + require.Equal(t, 0, n) + require.ErrorIs(t, err, io.EOF) + + // A whole-file read through the cache equals a direct backend read. + full := make([]byte, size) + fn, ferr := readAtFull(c, full, 0) + require.NoError(t, ferr) + require.Equal(t, size, fn) + direct := make([]byte, size) + _, _ = readAtFull(backend, direct, 0) + require.Equal(t, direct, full) +} + +type readOp struct{ pos, size int64 } + +func parseTrace(path string) ([]readOp, int64, error) { + f, err := os.Open(path) + if err != nil { + return nil, 0, err + } + defer f.Close() + var ops []readOp + var size int64 + sc := bufio.NewScanner(f) + for sc.Scan() { + var pos, sz int64 + if _, err := fmt.Sscanf(strings.TrimSpace(sc.Text()), "readAt pos=%d size=%d", &pos, &sz); err != nil { + continue + } + ops = append(ops, readOp{pos, sz}) + if pos+sz > size { + size = pos + sz + } + } + return ops, size, sc.Err() +} + +// baselineGets counts what today's ReaderAt would do: a fresh ranged +// GetObject for the first read and every subsequent non-sequential read. +func baselineGets(ops []readOp) int64 { + var gets, expected int64 + first := true + for _, op := range ops { + if first || op.pos != expected { + gets++ + } + expected = op.pos + op.size + first = false + } + return gets +} + +// TestCachingReaderAtSim replays a real RandomAccessSrcBin read trace +// (captured as "readAt pos=N size=M" lines) through the cache against a +// fake backend and reports cache behavior for a few configs. Set +// SEEK_TRACE=/path/to/seek-example to run it. +func TestCachingReaderAtSim(t *testing.T) { + path := os.Getenv("SEEK_TRACE") + if path == "" { + t.Skip("set SEEK_TRACE=/path/to/trace to run the cache simulation") + } + ops, size, err := parseTrace(path) + require.NoError(t, err) + require.NotEmpty(t, ops, "no 'readAt pos=.. size=..' lines parsed from %s", path) + + const MB = 1024 * 1024 + var requested int64 + for _, op := range ops { + requested += op.size + } + t.Logf("trace: %d reads, %.2f MB requested, object size %.2f GB", + len(ops), float64(requested)/MB, float64(size)/1e9) + t.Logf("baseline (current ReaderAt): %d GetObject round-trips", baselineGets(ops)) + t.Logf("--- LRU block cache ---") + + configs := []struct { + blockSize int64 + maxBlocks int + }{ + {16 * MB, 4}, {16 * MB, 10}, {16 * MB, 20}, + {8 * MB, 10}, {4 * MB, 10}, {1 * MB, 16}, + } + for _, cfg := range configs { + backend := &patternBackend{size: size} + c, err := NewCachingReaderAt(backend, size, cfg.blockSize, cfg.maxBlocks) + require.NoError(t, err) + for _, op := range ops { + _, _ = c.ReadAt(make([]byte, op.size), op.pos) + } + s := c.Stats() + hitRate := 100 * float64(s.Hits) / float64(s.BlockTouches) + amp := float64(s.BackendBytes) / float64(s.BytesRequested) + t.Logf("block=%-4s cache=%-2d (%4d MB) -> backendGETs=%-5d backendMB=%-8.1f hit=%.1f%% cold=%-3d redownloads=%-4d evictions=%-4d amp=%.1fx", + human(cfg.blockSize), cfg.maxBlocks, int(cfg.blockSize/MB)*cfg.maxBlocks, + s.BackendReads, float64(s.BackendBytes)/MB, hitRate, s.ColdMisses, s.Redownloads, s.Evictions, amp) + } +} + +func human(b int64) string { + const MB = 1024 * 1024 + if b%MB == 0 { + return fmt.Sprintf("%dMB", b/MB) + } + return fmt.Sprintf("%dKB", b/1024) +} -- 2.51.2