Test that idle host semaphores and .meta files are removed
New tests, failing before the fix: fetches from many hosts, and fetches that end without a connection, must leave no semaphore in hostSems once they finish; VariantStorage.Delete must remove the variant's .meta file, and must succeed when that file is already missing. Model: opus-5-5
This commit is contained in:
@@ -194,6 +194,14 @@ func semLen(f *HTTPFetcher, host string) int {
|
|||||||
return len(f.getHostSemaphore(host))
|
return len(f.getHostSemaphore(host))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// hostSemCount reports how many hosts have a semaphore in hostSems.
|
||||||
|
func hostSemCount(f *HTTPFetcher) int {
|
||||||
|
f.hostSemMu.Lock()
|
||||||
|
defer f.hostSemMu.Unlock()
|
||||||
|
|
||||||
|
return len(f.hostSems)
|
||||||
|
}
|
||||||
|
|
||||||
func TestFetchRedirectToPrivateIPBlocked(t *testing.T) {
|
func TestFetchRedirectToPrivateIPBlocked(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
|
|||||||
@@ -5,6 +5,7 @@ import (
|
|||||||
"errors"
|
"errors"
|
||||||
"net"
|
"net"
|
||||||
"strconv"
|
"strconv"
|
||||||
|
"sync"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
@@ -119,6 +120,84 @@ func TestFetchFreesHostSlotWhenContextEndsWaitingForConnection(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestFetchRemovesIdleHostSemaphores checks that a host's semaphore is
|
||||||
|
// removed once no fetch holds or waits for one of its slots: after 100
|
||||||
|
// concurrent fetches from 50 hosts have all finished, no semaphore is left.
|
||||||
|
func TestFetchRemovesIdleHostSemaphores(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
srv := startUpstream(t)
|
||||||
|
f, _ := newServerFetcher(t, srv, nil)
|
||||||
|
ctx := testContext(t)
|
||||||
|
|
||||||
|
var wg sync.WaitGroup
|
||||||
|
|
||||||
|
for i := range 100 {
|
||||||
|
wg.Go(func() {
|
||||||
|
res, err := f.Fetch(ctx, imageURLOnPort(1+i%50))
|
||||||
|
if err != nil {
|
||||||
|
t.Errorf("Fetch() error = %v", err)
|
||||||
|
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
_ = res.Content.Close()
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
wg.Wait()
|
||||||
|
|
||||||
|
if n := hostSemCount(f); n != 0 {
|
||||||
|
t.Errorf("%d host semaphores left after every fetch finished, want 0", n)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestFetchRemovesHostSemaphoreWhenNoConnection checks that a fetch that
|
||||||
|
// ends without a connection leaves no semaphore behind, both when it is
|
||||||
|
// refused after waiting for a connection shared by all hosts and when its
|
||||||
|
// context ends while it waits for its host's slot.
|
||||||
|
func TestFetchRemovesHostSemaphoreWhenNoConnection(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
srv := startUpstream(t)
|
||||||
|
|
||||||
|
cfg := DefaultConfig()
|
||||||
|
cfg.MaxConnections = 1
|
||||||
|
cfg.MaxConnectionsPerHost = 1
|
||||||
|
|
||||||
|
f, _ := newServerFetcher(t, srv, cfg)
|
||||||
|
f.connectionWaitTimeout = 100 * time.Millisecond
|
||||||
|
|
||||||
|
open, err := f.Fetch(testContext(t), imageURLOnPort(81))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("first Fetch() error = %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
_, err = f.Fetch(testContext(t), imageURLOnPort(82))
|
||||||
|
if !errors.Is(err, ErrTooManyConnections) {
|
||||||
|
t.Fatalf("Fetch() from another host: error = %v, "+
|
||||||
|
"want ErrTooManyConnections", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(t.Context(), 100*time.Millisecond)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
_, err = f.Fetch(ctx, imageURLOnPort(81))
|
||||||
|
if !errors.Is(err, context.DeadlineExceeded) {
|
||||||
|
t.Fatalf("Fetch() from the busy host: error = %v, "+
|
||||||
|
"want context.DeadlineExceeded", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
err = open.Content.Close()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("close first body: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if n := hostSemCount(f); n != 0 {
|
||||||
|
t.Errorf("%d host semaphores left after every fetch finished, want 0", n)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// TestFetchReleasesConnectionOnError checks that a fetch that fails after
|
// TestFetchReleasesConnectionOnError checks that a fetch that fails after
|
||||||
// taking its connection gives it back: with MaxConnections at 1, the slot
|
// taking its connection gives it back: with MaxConnections at 1, the slot
|
||||||
// must be free after the failure and the next fetch must succeed.
|
// must be free after the failure and the next fetch must succeed.
|
||||||
|
|||||||
@@ -438,3 +438,73 @@ func TestVariantStorage_StoreLogsFailedMetaWrite(t *testing.T) {
|
|||||||
t.Errorf("log missing %s; got %q", want, logBuf.String())
|
t.Errorf("log missing %s; got %q", want, logBuf.String())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// storeTestVariant stores one variant, with its .meta file, in a new
|
||||||
|
// VariantStorage and returns the storage and the variant's key.
|
||||||
|
func storeTestVariant(t *testing.T) (*VariantStorage, VariantKey) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
storage, err := NewVariantStorage(t.TempDir(), slog.New(slog.DiscardHandler))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("NewVariantStorage() error = %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
key := CacheKey(&ImageRequest{SourceHost: testHostCDN, SourcePath: testPathCat})
|
||||||
|
|
||||||
|
_, err = storage.Store(key, bytes.NewReader([]byte("variant data")), "image/webp")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Store() error = %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
return storage, key
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestVariantStorage_DeleteRemovesMeta verifies that Delete removes the
|
||||||
|
// variant's .meta file along with the variant file.
|
||||||
|
func TestVariantStorage_DeleteRemovesMeta(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
storage, key := storeTestVariant(t)
|
||||||
|
metaPath := storage.keyToPath(key) + ".meta"
|
||||||
|
|
||||||
|
_, err := os.Stat(metaPath)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Store() wrote no .meta file: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
err = storage.Delete(key)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Delete() error = %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if storage.Exists(key) {
|
||||||
|
t.Error("Exists() = true after delete, want false")
|
||||||
|
}
|
||||||
|
|
||||||
|
_, err = os.Stat(metaPath)
|
||||||
|
if !os.IsNotExist(err) {
|
||||||
|
t.Errorf(".meta file left after Delete() (stat err=%v)", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestVariantStorage_DeleteWithoutMeta verifies that Delete succeeds for a
|
||||||
|
// variant whose .meta file is missing.
|
||||||
|
func TestVariantStorage_DeleteWithoutMeta(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
storage, key := storeTestVariant(t)
|
||||||
|
|
||||||
|
err := os.Remove(storage.keyToPath(key) + ".meta")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("removing .meta file: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
err = storage.Delete(key)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Delete() error = %v, want nil", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if storage.Exists(key) {
|
||||||
|
t.Error("Exists() = true after delete, want false")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user