mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-11 16:57:45 +02:00
* filer: skip cache lock when a cacher cannot be removed yet SingleChunkCacher.readChunkAt defers removeConsumed on every read, and every unpin retries removal too. Each call acquired the ReaderCache mutex even though the removal conditions are plain atomics, so a busy cache paid a lock acquisition per read just to discover readers>0. Check the atomics before locking: when a cacher is not yet consumable the removal is impossible and the lock round trip is pure contention. Removals still run under the lock via removeConsumedLocked, so the attach-vs-remove race keeps its existing serialization. Ref #11676 * filer: guard stream position with a per-stream mutex, not the cache lock chunkStream.cacher is only ever shared by concurrent ReadAt calls on one ChunkReadAt, yet pin, unpin, releaseIfFinished, and releaseStream all mutated it under the ReaderCache mutex that serializes every reader in the process. Each read that switched chunks took that global lock to pin the new chunk, released it, then reacquired it in unpin() just to drop the old chunk's pin and retry removal; releaseIfFinished and releaseStream did the same lock-detach-unlock-relock dance. At high small-GET concurrency that turned per-request stream bookkeeping into a global convoy (#11676). Give chunkStream its own mutex and drop the cache lock from the pin lifecycle entirely: pin/detach serialize on stream.mu, the pin counter and consumable checks are already atomics, and removeConsumed only takes the cache lock when a cacher is actually removable. The map lock now guards only map membership and read registration. * filer: look up cached chunks under the read lock readChunkAt serialized every small read on the write lock even though the common paths are read-only: an existing downloader just needs its read registered, and a chunkCache hit needs no map access at all. Each GET also paid the lock a second time to reach the chunkCache check. Use the read lock for the downloader lookup and registration; the registered read keeps the buffer alive against a concurrent destroy, and only error eviction, insertion, and removal need the write lock. The chunkCache probe runs lock-free, and the insert path re-checks the map under the write lock to cover a downloader registered in between. * filer: start the chunk download outside the map lock The insert path held the ReaderCache write lock across goroutine spawn and the cacheStartedCh handshake, so every downloader miss serialized against the startup of a fetch goroutine. Register the cacher in the map under the lock, then start the download after releasing it; a fetch that fails early still lands in the map and is evicted by the next reader's completed-error check. * filer: test that stream pin lifecycle stays off the cache lock Regression coverage for the contention fix: pin and releaseStream on a non-removable cacher must complete while the ReaderCache lock is held by another goroutine. * filer: pin the stream under the map lock Between read registration and stream.pin the cacher showed zero pins, so a budget eviction in that gap could pick a chunk the stream was just attaching to and the stream's next slice refetched it. The pin counter is an atomic and stream.mu is never held while acquiring the cache lock, so pinning inside the map hold is deadlock-free and closes the window.
758 lines
23 KiB
Go
758 lines
23 KiB
Go
package filer
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"net/http"
|
|
"net/http/httptest"
|
|
"sync"
|
|
"sync/atomic"
|
|
"testing"
|
|
"time"
|
|
|
|
util_http "github.com/seaweedfs/seaweedfs/weed/util/http"
|
|
)
|
|
|
|
// mockChunkCacheForReaderCache implements chunk cache for testing
|
|
type mockChunkCacheForReaderCache struct {
|
|
data map[string][]byte
|
|
hitCount int32
|
|
mu sync.Mutex
|
|
}
|
|
|
|
type mockCacheInvalidatorForReaderCache struct {
|
|
calls int32
|
|
mu sync.Mutex
|
|
fileId string
|
|
}
|
|
|
|
func (m *mockCacheInvalidatorForReaderCache) InvalidateCache(fileId string) {
|
|
atomic.AddInt32(&m.calls, 1)
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
m.fileId = fileId
|
|
}
|
|
|
|
func (m *mockCacheInvalidatorForReaderCache) lastFileId() string {
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
return m.fileId
|
|
}
|
|
|
|
func newMockChunkCacheForReaderCache() *mockChunkCacheForReaderCache {
|
|
return &mockChunkCacheForReaderCache{
|
|
data: make(map[string][]byte),
|
|
}
|
|
}
|
|
|
|
func (m *mockChunkCacheForReaderCache) GetChunk(fileId string, minSize uint64) []byte {
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
if d, ok := m.data[fileId]; ok {
|
|
atomic.AddInt32(&m.hitCount, 1)
|
|
return d
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (m *mockChunkCacheForReaderCache) ReadChunkAt(data []byte, fileId string, offset uint64) (int, error) {
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
if d, ok := m.data[fileId]; ok && int(offset) < len(d) {
|
|
atomic.AddInt32(&m.hitCount, 1)
|
|
n := copy(data, d[offset:])
|
|
return n, nil
|
|
}
|
|
return 0, nil
|
|
}
|
|
|
|
func (m *mockChunkCacheForReaderCache) SetChunk(fileId string, data []byte) {
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
m.data[fileId] = data
|
|
}
|
|
|
|
func (m *mockChunkCacheForReaderCache) GetMaxFilePartSizeInCache() uint64 {
|
|
return 1024 * 1024 // 1MB
|
|
}
|
|
|
|
func (m *mockChunkCacheForReaderCache) IsInCache(fileId string, lockNeeded bool) bool {
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
_, ok := m.data[fileId]
|
|
return ok
|
|
}
|
|
|
|
func TestReaderCacheRetryAfterCacheInvalidation(t *testing.T) {
|
|
cache := newMockChunkCacheForReaderCache()
|
|
invalidator := &mockCacheInvalidatorForReaderCache{}
|
|
fileId := "425141,ef8914e9bdbbe8cb6838191f"
|
|
staleUrl := "http://fast-volume-1/" + fileId
|
|
freshUrl := "http://fast-volume-9/" + fileId
|
|
testData := []byte("fresh chunk data after cache invalidation")
|
|
|
|
var lookupCount int32
|
|
lookupFn := func(ctx context.Context, requestedFileId string) ([]string, error) {
|
|
if requestedFileId != fileId {
|
|
return nil, fmt.Errorf("unexpected lookup file id %s", requestedFileId)
|
|
}
|
|
atomic.AddInt32(&lookupCount, 1)
|
|
if atomic.LoadInt32(&invalidator.calls) == 0 {
|
|
return []string{staleUrl}, nil
|
|
}
|
|
return []string{freshUrl}, nil
|
|
}
|
|
|
|
var fetchCount int32
|
|
fetchFn := func(ctx context.Context, buffer []byte, urlStrings []string, cipherKey []byte, isGzipped bool, isFullChunk bool, offset int64, requestedFileId string, _ util_http.RefreshUrlsFunc) (int, error) {
|
|
if requestedFileId != fileId {
|
|
return 0, fmt.Errorf("unexpected fetch file id %s", requestedFileId)
|
|
}
|
|
switch atomic.AddInt32(&fetchCount, 1) {
|
|
case 1:
|
|
if len(urlStrings) != 1 || urlStrings[0] != staleUrl {
|
|
return 0, fmt.Errorf("first fetch should use stale url %v", urlStrings)
|
|
}
|
|
return 0, fmt.Errorf("404 Not Found: not found")
|
|
case 2:
|
|
if len(urlStrings) != 1 || urlStrings[0] != freshUrl {
|
|
return 0, fmt.Errorf("retry fetch should use fresh url %v", urlStrings)
|
|
}
|
|
return copy(buffer, testData), nil
|
|
default:
|
|
return 0, fmt.Errorf("unexpected extra fetch with urls %v", urlStrings)
|
|
}
|
|
}
|
|
|
|
rc := NewReaderCache(10, cache, lookupFn, invalidator)
|
|
rc.fetchChunkDataFn = fetchFn
|
|
defer rc.destroy()
|
|
|
|
buffer := make([]byte, len(testData))
|
|
n, err := rc.ReadChunkAt(context.Background(), buffer, fileId, nil, false, 0, len(testData), true)
|
|
if err != nil {
|
|
t.Fatalf("expected successful retry, got %v", err)
|
|
}
|
|
if got := string(buffer[:n]); got != string(testData) {
|
|
t.Fatalf("expected %q, got %q", testData, got)
|
|
}
|
|
if got := atomic.LoadInt32(&invalidator.calls); got != 1 {
|
|
t.Fatalf("expected one cache invalidation, got %d", got)
|
|
}
|
|
if got := invalidator.lastFileId(); got != fileId {
|
|
t.Fatalf("expected invalidated file id %s, got %s", fileId, got)
|
|
}
|
|
if got := atomic.LoadInt32(&lookupCount); got != 2 {
|
|
t.Fatalf("expected lookup before fetch and after invalidation, got %d", got)
|
|
}
|
|
if got := atomic.LoadInt32(&fetchCount); got != 2 {
|
|
t.Fatalf("expected stale fetch and retry fetch, got %d", got)
|
|
}
|
|
}
|
|
|
|
// TestReaderCacheRefreshesLocationsAfterPartialFailure is the mount reading
|
|
// through a cached location list that still names a replica the master has
|
|
// dropped: the read succeeds on the other replica, and the cached entry must
|
|
// be dropped and looked up again so the next read starts without the dead one.
|
|
func TestReaderCacheRefreshesLocationsAfterPartialFailure(t *testing.T) {
|
|
payload := []byte("chunk contents")
|
|
live := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
w.Write(payload)
|
|
}))
|
|
defer live.Close()
|
|
dead := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {}))
|
|
deadURL := dead.URL
|
|
dead.Close()
|
|
|
|
fileId := "3,abc"
|
|
invalidator := &mockCacheInvalidatorForReaderCache{}
|
|
var lookupCount int32
|
|
lookupFn := func(ctx context.Context, requestedFileId string) ([]string, error) {
|
|
atomic.AddInt32(&lookupCount, 1)
|
|
if atomic.LoadInt32(&invalidator.calls) == 0 {
|
|
return []string{deadURL + "/" + fileId, live.URL + "/" + fileId}, nil
|
|
}
|
|
return []string{live.URL + "/" + fileId}, nil
|
|
}
|
|
|
|
rc := NewReaderCache(10, newMockChunkCacheForReaderCache(), lookupFn, invalidator)
|
|
defer rc.destroy()
|
|
|
|
buffer := make([]byte, len(payload))
|
|
n, err := rc.ReadChunkAt(context.Background(), buffer, fileId, nil, false, 0, len(payload), false)
|
|
if err != nil || string(buffer[:n]) != string(payload) {
|
|
t.Fatalf("got %q, %v; want %q", buffer[:n], err, payload)
|
|
}
|
|
if got := atomic.LoadInt32(&invalidator.calls); got != 1 {
|
|
t.Fatalf("expected the stale entry to be invalidated once, got %d", got)
|
|
}
|
|
if got := atomic.LoadInt32(&lookupCount); got != 2 {
|
|
t.Fatalf("expected a lookup before the read and one after invalidation, got %d", got)
|
|
}
|
|
}
|
|
|
|
func TestReaderCacheRemovesFailedDownloader(t *testing.T) {
|
|
cache := newMockChunkCacheForReaderCache()
|
|
fileId := "425141,failed"
|
|
url := "http://fast-volume-1/" + fileId
|
|
|
|
var lookupCount int32
|
|
lookupFn := func(ctx context.Context, requestedFileId string) ([]string, error) {
|
|
if requestedFileId != fileId {
|
|
return nil, fmt.Errorf("unexpected lookup file id %s", requestedFileId)
|
|
}
|
|
atomic.AddInt32(&lookupCount, 1)
|
|
return []string{url}, nil
|
|
}
|
|
|
|
var fetchCount int32
|
|
fetchFn := func(ctx context.Context, buffer []byte, urlStrings []string, cipherKey []byte, isGzipped bool, isFullChunk bool, offset int64, requestedFileId string, _ util_http.RefreshUrlsFunc) (int, error) {
|
|
atomic.AddInt32(&fetchCount, 1)
|
|
return 0, fmt.Errorf("fetch failed")
|
|
}
|
|
|
|
rc := NewReaderCache(10, cache, lookupFn, nil)
|
|
rc.fetchChunkDataFn = fetchFn
|
|
defer rc.destroy()
|
|
|
|
for i := 0; i < 2; i++ {
|
|
buffer := make([]byte, 8)
|
|
_, err := rc.ReadChunkAt(context.Background(), buffer, fileId, nil, false, 0, len(buffer), true)
|
|
if err == nil {
|
|
t.Fatalf("read %d should fail", i+1)
|
|
}
|
|
}
|
|
|
|
if got := atomic.LoadInt32(&lookupCount); got != 2 {
|
|
t.Fatalf("failed downloader should be removed so the second read re-lookups, got %d lookups", got)
|
|
}
|
|
if got := atomic.LoadInt32(&fetchCount); got != 2 {
|
|
t.Fatalf("failed downloader should be removed so the second read refetches, got %d fetches", got)
|
|
}
|
|
}
|
|
|
|
// TestReaderCacheContextCancellation tests that a reader can cancel its wait
|
|
// while the download continues for other readers
|
|
func TestReaderCacheContextCancellation(t *testing.T) {
|
|
cache := newMockChunkCacheForReaderCache()
|
|
|
|
// Create a ReaderCache - we can't easily test the full flow without mocking HTTP,
|
|
// but we can test the context cancellation in readChunkAt
|
|
rc := NewReaderCache(10, cache, nil, nil)
|
|
defer rc.destroy()
|
|
|
|
// Pre-populate cache to avoid HTTP calls
|
|
testData := []byte("test data for context cancellation")
|
|
cache.SetChunk("test-file-1", testData)
|
|
|
|
// Test that context cancellation works
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
|
|
buffer := make([]byte, len(testData))
|
|
n, err := rc.ReadChunkAt(ctx, buffer, "test-file-1", nil, false, 0, len(testData), true)
|
|
if err != nil {
|
|
t.Errorf("Expected no error, got: %v", err)
|
|
}
|
|
if n != len(testData) {
|
|
t.Errorf("Expected %d bytes, got %d", len(testData), n)
|
|
}
|
|
|
|
// Cancel context and verify it doesn't affect already completed reads
|
|
cancel()
|
|
|
|
// Subsequent read with cancelled context should still work from cache
|
|
buffer2 := make([]byte, len(testData))
|
|
n2, err2 := rc.ReadChunkAt(ctx, buffer2, "test-file-1", nil, false, 0, len(testData), true)
|
|
// Note: This may or may not error depending on whether it hits cache
|
|
_ = n2
|
|
_ = err2
|
|
}
|
|
|
|
// TestReaderCacheFallbackToChunkCache tests that when a cacher returns n=0, err=nil,
|
|
// we fall back to the chunkCache
|
|
func TestReaderCacheFallbackToChunkCache(t *testing.T) {
|
|
cache := newMockChunkCacheForReaderCache()
|
|
|
|
// Pre-populate the chunk cache with data
|
|
testData := []byte("fallback test data that should be found in chunk cache")
|
|
cache.SetChunk("fallback-file", testData)
|
|
|
|
rc := NewReaderCache(10, cache, nil, nil)
|
|
defer rc.destroy()
|
|
|
|
// Read should hit the chunk cache
|
|
buffer := make([]byte, len(testData))
|
|
n, err := rc.ReadChunkAt(context.Background(), buffer, "fallback-file", nil, false, 0, len(testData), true)
|
|
|
|
if err != nil {
|
|
t.Errorf("Expected no error, got: %v", err)
|
|
}
|
|
if n != len(testData) {
|
|
t.Errorf("Expected %d bytes, got %d", len(testData), n)
|
|
}
|
|
|
|
// Verify cache was hit
|
|
if cache.hitCount == 0 {
|
|
t.Error("Expected chunk cache to be hit")
|
|
}
|
|
}
|
|
|
|
// TestReaderCacheMultipleReadersWaitForSameChunk tests that multiple readers
|
|
// can wait for the same chunk download to complete
|
|
func TestReaderCacheMultipleReadersWaitForSameChunk(t *testing.T) {
|
|
cache := newMockChunkCacheForReaderCache()
|
|
|
|
// Pre-populate cache so we don't need HTTP
|
|
testData := make([]byte, 1024)
|
|
for i := range testData {
|
|
testData[i] = byte(i % 256)
|
|
}
|
|
cache.SetChunk("shared-chunk", testData)
|
|
|
|
rc := NewReaderCache(10, cache, nil, nil)
|
|
defer rc.destroy()
|
|
|
|
// Launch multiple concurrent readers for the same chunk
|
|
numReaders := 10
|
|
var wg sync.WaitGroup
|
|
errors := make(chan error, numReaders)
|
|
bytesRead := make(chan int, numReaders)
|
|
|
|
for i := 0; i < numReaders; i++ {
|
|
wg.Add(1)
|
|
go func() {
|
|
defer wg.Done()
|
|
buffer := make([]byte, len(testData))
|
|
n, err := rc.ReadChunkAt(context.Background(), buffer, "shared-chunk", nil, false, 0, len(testData), true)
|
|
if err != nil {
|
|
errors <- err
|
|
}
|
|
bytesRead <- n
|
|
}()
|
|
}
|
|
|
|
wg.Wait()
|
|
close(errors)
|
|
close(bytesRead)
|
|
|
|
// Check for errors
|
|
for err := range errors {
|
|
t.Errorf("Reader got error: %v", err)
|
|
}
|
|
|
|
// Verify all readers got the expected data
|
|
for n := range bytesRead {
|
|
if n != len(testData) {
|
|
t.Errorf("Expected %d bytes, got %d", len(testData), n)
|
|
}
|
|
}
|
|
}
|
|
|
|
// TestReaderCachePartialRead tests reading at different offsets
|
|
func TestReaderCachePartialRead(t *testing.T) {
|
|
cache := newMockChunkCacheForReaderCache()
|
|
|
|
testData := []byte("0123456789ABCDEFGHIJ")
|
|
cache.SetChunk("partial-read-file", testData)
|
|
|
|
rc := NewReaderCache(10, cache, nil, nil)
|
|
defer rc.destroy()
|
|
|
|
tests := []struct {
|
|
name string
|
|
offset int64
|
|
size int
|
|
expected []byte
|
|
}{
|
|
{"read from start", 0, 5, []byte("01234")},
|
|
{"read from middle", 5, 5, []byte("56789")},
|
|
{"read to end", 15, 5, []byte("FGHIJ")},
|
|
{"read single byte", 10, 1, []byte("A")},
|
|
}
|
|
|
|
for _, tt := range tests {
|
|
t.Run(tt.name, func(t *testing.T) {
|
|
buffer := make([]byte, tt.size)
|
|
n, err := rc.ReadChunkAt(context.Background(), buffer, "partial-read-file", nil, false, tt.offset, len(testData), true)
|
|
|
|
if err != nil {
|
|
t.Errorf("Expected no error, got: %v", err)
|
|
}
|
|
if n != tt.size {
|
|
t.Errorf("Expected %d bytes, got %d", tt.size, n)
|
|
}
|
|
if string(buffer[:n]) != string(tt.expected) {
|
|
t.Errorf("Expected %q, got %q", tt.expected, buffer[:n])
|
|
}
|
|
})
|
|
}
|
|
}
|
|
|
|
// TestReaderCacheCleanup tests that old downloaders are cleaned up
|
|
func TestReaderCacheCleanup(t *testing.T) {
|
|
cache := newMockChunkCacheForReaderCache()
|
|
|
|
// Create cache with limit of 3
|
|
rc := NewReaderCache(3, cache, nil, nil)
|
|
defer rc.destroy()
|
|
|
|
// Add data for multiple files
|
|
for i := 0; i < 5; i++ {
|
|
fileId := string(rune('A' + i))
|
|
data := []byte("data for file " + fileId)
|
|
cache.SetChunk(fileId, data)
|
|
}
|
|
|
|
// Read from multiple files - should trigger cleanup when exceeding limit
|
|
for i := 0; i < 5; i++ {
|
|
fileId := string(rune('A' + i))
|
|
buffer := make([]byte, 20)
|
|
_, err := rc.ReadChunkAt(context.Background(), buffer, fileId, nil, false, 0, 20, true)
|
|
if err != nil {
|
|
t.Errorf("Read error for file %s: %v", fileId, err)
|
|
}
|
|
}
|
|
|
|
// Cache should still work - reads should succeed
|
|
for i := 0; i < 5; i++ {
|
|
fileId := string(rune('A' + i))
|
|
buffer := make([]byte, 20)
|
|
n, err := rc.ReadChunkAt(context.Background(), buffer, fileId, nil, false, 0, 20, true)
|
|
if err != nil {
|
|
t.Errorf("Second read error for file %s: %v", fileId, err)
|
|
}
|
|
if n == 0 {
|
|
t.Errorf("Expected data for file %s, got 0 bytes", fileId)
|
|
}
|
|
}
|
|
}
|
|
|
|
// TestSingleChunkCacherDoneSignal tests that done channel is always closed
|
|
func TestSingleChunkCacherDoneSignal(t *testing.T) {
|
|
cache := newMockChunkCacheForReaderCache()
|
|
rc := NewReaderCache(10, cache, nil, nil)
|
|
defer rc.destroy()
|
|
|
|
// Test that we can read even when data is in cache (done channel should work)
|
|
testData := []byte("done signal test")
|
|
cache.SetChunk("done-signal-test", testData)
|
|
|
|
// Multiple goroutines reading same chunk
|
|
var wg sync.WaitGroup
|
|
for i := 0; i < 5; i++ {
|
|
wg.Add(1)
|
|
go func() {
|
|
defer wg.Done()
|
|
buffer := make([]byte, len(testData))
|
|
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
|
|
n, err := rc.ReadChunkAt(ctx, buffer, "done-signal-test", nil, false, 0, len(testData), true)
|
|
if err != nil && err != context.DeadlineExceeded {
|
|
t.Errorf("Unexpected error: %v", err)
|
|
}
|
|
if n == 0 && err == nil {
|
|
t.Error("Got 0 bytes with no error")
|
|
}
|
|
}()
|
|
}
|
|
|
|
// Should complete without hanging
|
|
done := make(chan struct{})
|
|
go func() {
|
|
wg.Wait()
|
|
close(done)
|
|
}()
|
|
|
|
select {
|
|
case <-done:
|
|
// Success
|
|
case <-time.After(10 * time.Second):
|
|
t.Fatal("Test timed out - done channel may not be signaled correctly")
|
|
}
|
|
}
|
|
|
|
// ============================================================================
|
|
// Tests that exercise SingleChunkCacher concurrency logic
|
|
// ============================================================================
|
|
//
|
|
// These tests use blocking lookupFileIdFn to exercise the wait/cancellation
|
|
// logic in SingleChunkCacher without requiring HTTP calls.
|
|
|
|
// TestSingleChunkCacherLookupError tests handling of lookup errors
|
|
func TestSingleChunkCacherLookupError(t *testing.T) {
|
|
cache := newMockChunkCacheForReaderCache()
|
|
|
|
// Lookup function that returns an error
|
|
lookupFn := func(ctx context.Context, fileId string) ([]string, error) {
|
|
return nil, fmt.Errorf("lookup failed for %s", fileId)
|
|
}
|
|
|
|
rc := NewReaderCache(10, cache, lookupFn, nil)
|
|
defer rc.destroy()
|
|
|
|
buffer := make([]byte, 100)
|
|
_, err := rc.ReadChunkAt(context.Background(), buffer, "error-test", nil, false, 0, 100, true)
|
|
|
|
if err == nil {
|
|
t.Error("Expected an error, got nil")
|
|
}
|
|
}
|
|
|
|
// TestSingleChunkCacherContextCancellationDuringLookup tests that a reader can
|
|
// cancel its wait while the lookup is in progress. This exercises the actual
|
|
// SingleChunkCacher wait/cancel logic.
|
|
func TestSingleChunkCacherContextCancellationDuringLookup(t *testing.T) {
|
|
cache := newMockChunkCacheForReaderCache()
|
|
lookupStarted := make(chan struct{})
|
|
lookupCanFinish := make(chan struct{})
|
|
|
|
// Lookup function that blocks to simulate slow operation
|
|
lookupFn := func(ctx context.Context, fileId string) ([]string, error) {
|
|
close(lookupStarted)
|
|
<-lookupCanFinish // Block until test allows completion
|
|
return nil, fmt.Errorf("lookup completed but reader should have cancelled")
|
|
}
|
|
|
|
rc := NewReaderCache(10, cache, lookupFn, nil)
|
|
defer rc.destroy()
|
|
defer close(lookupCanFinish) // Ensure cleanup
|
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
readResult := make(chan error, 1)
|
|
|
|
go func() {
|
|
buffer := make([]byte, 100)
|
|
_, err := rc.ReadChunkAt(ctx, buffer, "cancel-during-lookup", nil, false, 0, 100, true)
|
|
readResult <- err
|
|
}()
|
|
|
|
// Wait for lookup to start, then cancel the reader's context
|
|
select {
|
|
case <-lookupStarted:
|
|
cancel() // Cancel the reader while lookup is blocked
|
|
case <-time.After(5 * time.Second):
|
|
t.Fatal("Lookup never started")
|
|
}
|
|
|
|
// Read should return with context.Canceled
|
|
select {
|
|
case err := <-readResult:
|
|
if err != context.Canceled {
|
|
t.Errorf("Expected context.Canceled, got: %v", err)
|
|
}
|
|
case <-time.After(5 * time.Second):
|
|
t.Fatal("Read did not complete after context cancellation")
|
|
}
|
|
}
|
|
|
|
// TestSingleChunkCacherMultipleReadersWaitForDownload tests that multiple readers
|
|
// can wait for the same SingleChunkCacher download to complete. When lookup fails,
|
|
// all readers should receive the same error.
|
|
func TestSingleChunkCacherMultipleReadersWaitForDownload(t *testing.T) {
|
|
cache := newMockChunkCacheForReaderCache()
|
|
lookupStarted := make(chan struct{})
|
|
lookupCanFinish := make(chan struct{})
|
|
var lookupStartedOnce sync.Once
|
|
|
|
// Lookup function that blocks to simulate slow operation
|
|
lookupFn := func(ctx context.Context, fileId string) ([]string, error) {
|
|
lookupStartedOnce.Do(func() { close(lookupStarted) })
|
|
<-lookupCanFinish
|
|
return nil, fmt.Errorf("simulated lookup error")
|
|
}
|
|
|
|
rc := NewReaderCache(10, cache, lookupFn, nil)
|
|
defer rc.destroy()
|
|
|
|
numReaders := 5
|
|
var wg sync.WaitGroup
|
|
errors := make(chan error, numReaders)
|
|
|
|
// Start multiple readers for the same chunk
|
|
for i := 0; i < numReaders; i++ {
|
|
wg.Add(1)
|
|
go func() {
|
|
defer wg.Done()
|
|
buffer := make([]byte, 100)
|
|
_, err := rc.ReadChunkAt(context.Background(), buffer, "shared-chunk", nil, false, 0, 100, true)
|
|
errors <- err
|
|
}()
|
|
}
|
|
|
|
// Wait for lookup to start, then allow completion
|
|
select {
|
|
case <-lookupStarted:
|
|
close(lookupCanFinish)
|
|
case <-time.After(5 * time.Second):
|
|
close(lookupCanFinish)
|
|
t.Fatal("Lookup never started")
|
|
}
|
|
|
|
wg.Wait()
|
|
close(errors)
|
|
|
|
// All readers should receive an error
|
|
errorCount := 0
|
|
for err := range errors {
|
|
if err != nil {
|
|
errorCount++
|
|
}
|
|
}
|
|
if errorCount != numReaders {
|
|
t.Errorf("Expected %d errors, got %d", numReaders, errorCount)
|
|
}
|
|
}
|
|
|
|
// TestReaderCacheDownloaderDedup tests that concurrent ReadChunkAt calls for
|
|
// the same fileId result in only one network fetch (lookup call), because
|
|
// the downloaders map deduplicates in-flight downloads.
|
|
func TestReaderCacheDownloaderDedup(t *testing.T) {
|
|
cache := newMockChunkCacheForReaderCache()
|
|
|
|
var lookupCount int32
|
|
var fetchCount int32
|
|
fetchGate := make(chan struct{})
|
|
testData := []byte("deduplicated data")
|
|
|
|
lookupFn := func(ctx context.Context, fileId string) ([]string, error) {
|
|
atomic.AddInt32(&lookupCount, 1)
|
|
return []string{"http://volume/" + fileId}, nil
|
|
}
|
|
|
|
fetchFn := func(ctx context.Context, buffer []byte, urlStrings []string, cipherKey []byte, isGzipped bool, isFullChunk bool, offset int64, fileId string, _ util_http.RefreshUrlsFunc) (int, error) {
|
|
atomic.AddInt32(&fetchCount, 1)
|
|
<-fetchGate
|
|
return copy(buffer, testData), nil
|
|
}
|
|
|
|
rc := NewReaderCache(10, cache, lookupFn, nil)
|
|
rc.fetchChunkDataFn = fetchFn
|
|
defer rc.destroy()
|
|
|
|
const numReaders = 10
|
|
var wg sync.WaitGroup
|
|
wg.Add(numReaders)
|
|
|
|
for i := 0; i < numReaders; i++ {
|
|
go func() {
|
|
defer wg.Done()
|
|
buffer := make([]byte, 50)
|
|
rc.ReadChunkAt(context.Background(), buffer, "dedup-file", nil, false, 0, 100, false)
|
|
}()
|
|
}
|
|
|
|
// Allow downloads to proceed.
|
|
close(fetchGate)
|
|
wg.Wait()
|
|
|
|
if count := atomic.LoadInt32(&lookupCount); count != 1 {
|
|
t.Errorf("expected exactly 1 lookup call, got %d", count)
|
|
}
|
|
if count := atomic.LoadInt32(&fetchCount); count != 1 {
|
|
t.Errorf("expected exactly 1 fetch call, got %d", count)
|
|
}
|
|
}
|
|
|
|
// TestSingleChunkCacherOneReaderCancelsOthersContinue tests that when one reader
|
|
// cancels, other readers waiting on the same chunk continue to wait.
|
|
func TestSingleChunkCacherOneReaderCancelsOthersContinue(t *testing.T) {
|
|
cache := newMockChunkCacheForReaderCache()
|
|
lookupStarted := make(chan struct{})
|
|
lookupCanFinish := make(chan struct{})
|
|
var lookupStartedOnce sync.Once
|
|
|
|
lookupFn := func(ctx context.Context, fileId string) ([]string, error) {
|
|
lookupStartedOnce.Do(func() { close(lookupStarted) })
|
|
<-lookupCanFinish
|
|
return nil, fmt.Errorf("simulated error after delay")
|
|
}
|
|
|
|
rc := NewReaderCache(10, cache, lookupFn, nil)
|
|
defer rc.destroy()
|
|
|
|
cancelledReaderDone := make(chan error, 1)
|
|
otherReaderDone := make(chan error, 1)
|
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
|
|
// Start reader that will be cancelled
|
|
go func() {
|
|
buffer := make([]byte, 100)
|
|
_, err := rc.ReadChunkAt(ctx, buffer, "shared-chunk-2", nil, false, 0, 100, true)
|
|
cancelledReaderDone <- err
|
|
}()
|
|
|
|
// Start reader that will NOT be cancelled
|
|
go func() {
|
|
buffer := make([]byte, 100)
|
|
_, err := rc.ReadChunkAt(context.Background(), buffer, "shared-chunk-2", nil, false, 0, 100, true)
|
|
otherReaderDone <- err
|
|
}()
|
|
|
|
// Wait for lookup to start
|
|
select {
|
|
case <-lookupStarted:
|
|
case <-time.After(5 * time.Second):
|
|
t.Fatal("Lookup never started")
|
|
}
|
|
|
|
// Cancel the first reader
|
|
cancel()
|
|
|
|
// First reader should complete with context.Canceled quickly
|
|
select {
|
|
case err := <-cancelledReaderDone:
|
|
if err != context.Canceled {
|
|
t.Errorf("Cancelled reader: expected context.Canceled, got: %v", err)
|
|
}
|
|
case <-time.After(2 * time.Second):
|
|
t.Error("Cancelled reader did not complete quickly")
|
|
}
|
|
|
|
// Allow the download to complete
|
|
close(lookupCanFinish)
|
|
|
|
// Other reader should eventually complete (with error since lookup returns error)
|
|
select {
|
|
case err := <-otherReaderDone:
|
|
if err == nil || err == context.Canceled {
|
|
t.Errorf("Other reader: expected non-nil non-cancelled error, got: %v", err)
|
|
}
|
|
// Expected: "simulated error after delay"
|
|
case <-time.After(5 * time.Second):
|
|
t.Error("Other reader did not complete")
|
|
}
|
|
}
|
|
|
|
// TestReaderCacheBookkeepingOffCacheLock guards the contention fix: pin
|
|
// lifecycle calls that cannot remove a cacher must not block on the
|
|
// ReaderCache lock.
|
|
func TestReaderCacheBookkeepingOffCacheLock(t *testing.T) {
|
|
cache := newMockChunkCacheForReaderCache()
|
|
rc := NewReaderCache(10, cache, nil, nil)
|
|
defer rc.destroy()
|
|
|
|
cacher := &SingleChunkCacher{parent: rc, chunkFileId: "pinned-chunk"}
|
|
atomic.StoreInt32(&cacher.readers, 1) // a read in flight: not consumable
|
|
|
|
rc.Lock()
|
|
started := make(chan struct{})
|
|
done := make(chan struct{})
|
|
go func() {
|
|
close(started)
|
|
defer close(done)
|
|
var stream chunkStream
|
|
stream.pin(cacher)
|
|
rc.releaseStream(&stream)
|
|
}()
|
|
<-started
|
|
select {
|
|
case <-done:
|
|
case <-time.After(time.Second):
|
|
rc.Unlock()
|
|
t.Fatal("stream pin lifecycle blocked on the cache lock")
|
|
}
|
|
rc.Unlock()
|
|
}
|