Watch
1
0
Fork
You've already forked pkg-proxy
1
mirror of https://github.com/git-pkgs/proxy.git synced 2026-09-16 15:52:05 -04:00
pkg-proxy/internal/handler/coalesce_semantics_test.go

615 lines
21 KiB
Go
Raw Permalink Normal View History

Fix duplicate fetches and 502s on concurrent cache misses (#329) * Return the stored artifact from storeArtifact, not a reader storeArtifact returned a CacheResult holding an open file handle. A handle has one read position, so it can only ever serve a single caller, which is what blocks sharing one fetch between concurrent requests. Return the artifact and its storage path instead, and let each caller open its own reader through openStoredArtifact. Threading that type through fetchAndCache, fetchAndCacheFromURL and their error paths is mechanical; behaviour is unchanged. * Coalesce concurrent cache misses A cache miss went from checkCache straight to an upstream fetch with nothing tracking in-flight work, so N concurrent requests for one uncached artifact produced N upstream fetches and N stores to the same key. That is the CI shape: parallel jobs installing overlapping dependencies against a cold cache. The duplicate stores also fail requests, racing fileblob's per-key ".attrs" sidecar into a partial read served as a 502. Over 12 runs of 8 simultaneous requests for one uncached tarball, against bb2205a: before, 8 fetches per run and 12 of 96 responses were 502; after, 1 fetch per run and none failed. Route both miss paths through a shared in-flight map keyed on the artifact, including the download URL and upstream-declared hash so callers expecting different bytes never share a fetch. singleflight does not fit: Do gives waiters no way to leave, while DoChan lets the caller running the fetch abandon it, breaking storeArtifact's scan-on-disconnect contract. Deciding roles under a mutex gives both behaviours. The fetch runs on the first caller's context and is seen through; waiters leave when their own clients do. This removes the sidecar trigger on this path. The race is in fileblob and three writers bypass this path entirely, so it is fixed separately. Fewer failures now reach the circuit breaker, so it trips later. Sixteen concurrent callers against real file:// storage fail 10 of 10 runs on main and pass 10 of 10 here. Other tests pin key discrimination, failure propagation, resolver-path coalescing, per-caller readers, waiter cancellation, key release and panic safety. allocs/op is unchanged. mockStorage gains a mutex so concurrent tests can use it. * Normalize digest case in the coalescing key artifactHashMatches compares digests with strings.EqualFold, but the coalescing key used the hash verbatim. The same digest in two casings produced two keys, so two callers for one artifact each ran their own upstream fetch and store, which is what the coalescing is meant to prevent. * Make the panic coalescing test deterministic The test timed the second caller's arrival with a sleep, so which caller became the leader was left to the scheduler. When it lost that race the second caller ran the fetch itself, and its panic was not recovered, so the test binary died instead of the test failing. Whether a caller has reached the wait is not observable from outside: it runs a cache lookup against the database first, so releasing the leader on a timer races that query. Drive coalesceFetch directly and hold the shared entry instead, which removes the timing entirely. The panicking fetcher is no longer needed. * Recheck the cache before running a shared fetch A caller checks the cache before it reaches coalesceFetch, so a fetch that commits in that gap is invisible to it. Arriving after the sharing entry is gone, it became a new leader and fetched, stored and scanned an artifact the cache already held. The leader now rechecks the committed record first. It serves that record only if its bytes still open, because a record can outlive them, and refetching is the recovery the cache lookup already makes for that case. Waiters are unaffected: the record fills the same shared value a fetch would, and each caller opens its own reader from it. The recheck is the leader's alone. A waiter has a fetch in flight to wait on, and rechecking would race it for no gain. * Lock the mock fetcher's bookkeeping Coalescing tests call the handler from many goroutines. The key keeps the fetch itself serialized, but the mock should not lean on that: it now locks the fields it records, so any concurrency the handler applies is safe under the race detector. * Wait for the leader's fetch instead of sleeping The canceled-waiter test slept 200ms and assumed the leader had taken the key by then. On a slow scheduler the canceled call could become the leader and the test would no longer cover waiter cancellation. The fetcher now signals when its first fetch begins, which happens only once the key is held.
2026-09-15 11:59:03 -07:00
package handler
import (
"context"
"errors"
"io"
"strings"
"sync"
"testing"
"time"
"github.com/git-pkgs/artifacts"
"github.com/git-pkgs/registries/fetch"
)
// runConcurrent runs fn in n goroutines released together and returns their errors.
func runConcurrent(n int, fn func(i int) error) []error {
errs := make([]error, n)
start := make(chan struct{})
var wg sync.WaitGroup
for i := 0; i < n; i++ {
wg.Add(1)
go func(i int) {
defer wg.Done()
<-start
errs[i] = fn(i)
}(i)
}
close(start)
wg.Wait()
return errs
}
// artifactBody builds a one-shot upstream artifact carrying the given bytes.
func artifactBody(content string) *fetch.Artifact {
return &fetch.Artifact{
Body: io.NopCloser(strings.NewReader(content)),
ContentType: "application/gzip",
}
}
// drain consumes and closes a CacheResult reader, if there is one.
func drain(res *CacheResult) {
if res != nil && res.Reader != nil {
_, _ = io.Copy(io.Discard, res.Reader)
_ = res.Reader.Close()
}
}
// TestCoalesceKey_DifferentUpstreamHashDoesNotShare is the safety property that
// makes coalescing sound: callers expecting different bytes must never share a
// fetch, so a re-published version cannot serve stale bytes to a caller that
// asked for the new digest.
func TestCoalesceKey_DifferentUpstreamHashDoesNotShare(t *testing.T) {
const content = "artifact bytes"
proxy, _, _, _ := setupTestProxy(t)
fetcher := &countingFetcher{content: content, delay: fetchHoldTime}
proxy.Fetcher = fetcher
// The digest must carry the "sha256:" prefix; without it the API treats the
// value as unverifiable and clears the hash, which would legitimately let
// the two callers share one fetch.
hashes := []string{
"sha256:" + sha256Hex(content),
"sha256:" + sha256Hex("something else entirely"),
}
_ = runConcurrent(2, func(i int) error {
res, err := proxy.GetOrFetchArtifactFromURLWithDigest(context.Background(),
"npm", "pkg", "1.0.0", "pkg-1.0.0.tgz",
"https://registry.npmjs.org/pkg/-/pkg-1.0.0.tgz", hashes[i])
drain(res)
return err
})
if got := fetcher.calls.Load(); got != 2 {
t.Errorf("upstream fetches = %d, want 2: callers expecting different digests must not share a fetch", got)
}
}
// TestCoalesceKey_HashCasingSharesOneFetch is the other half of that property.
// artifactHashMatches compares digests case-insensitively, so one digest in two
// casings describes one artifact and must not split into two fetches.
func TestCoalesceKey_HashCasingSharesOneFetch(t *testing.T) {
const content = "artifact bytes"
proxy, _, _, _ := setupTestProxy(t)
fetcher := &countingFetcher{content: content, delay: fetchHoldTime}
proxy.Fetcher = fetcher
hex := sha256Hex(content)
digests := []string{"sha256:" + hex, "sha256:" + strings.ToUpper(hex)}
for i, err := range runConcurrent(2, func(i int) error {
res, err := proxy.GetOrFetchArtifactFromURLWithDigest(context.Background(),
"npm", "pkg", "1.0.0", "pkg-1.0.0.tgz",
"https://registry.npmjs.org/pkg/-/pkg-1.0.0.tgz", digests[i])
drain(res)
return err
}) {
if err != nil {
t.Fatalf("caller %d failed: %v", i, err)
}
}
if got := fetcher.calls.Load(); got != 1 {
t.Errorf("upstream fetches = %d, want 1: one digest in two casings is one artifact", got)
}
}
// TestCoalesceKey_DifferentDownloadURLDoesNotShare covers the other half of the
// key: same package, different upstream URL, must not collapse into one fetch.
func TestCoalesceKey_DifferentDownloadURLDoesNotShare(t *testing.T) {
proxy, _, _, _ := setupTestProxy(t)
fetcher := &countingFetcher{content: "artifact bytes", delay: fetchHoldTime}
proxy.Fetcher = fetcher
urls := []string{
"https://registry.npmjs.org/pkg/-/pkg-1.0.0.tgz",
"https://mirror.example.com/pkg/-/pkg-1.0.0.tgz",
}
_ = runConcurrent(2, func(i int) error {
res, err := proxy.GetOrFetchArtifactFromURL(context.Background(),
"npm", "pkg", "1.0.0", "pkg-1.0.0.tgz", urls[i])
drain(res)
return err
})
if got := fetcher.calls.Load(); got != 2 {
t.Errorf("upstream fetches = %d, want 2: different upstream URLs must not share a fetch", got)
}
}
// TestCoalesceKey_DistinctArtifactsDoNotSerialize guards against an over-broad
// key: four packages fetched at once must still produce four fetches.
func TestCoalesceKey_DistinctArtifactsDoNotSerialize(t *testing.T) {
const n = 4
proxy, _, _, _ := setupTestProxy(t)
fetcher := &countingFetcher{content: "artifact bytes", delay: fetchHoldTime}
proxy.Fetcher = fetcher
names := []string{"alpha", "beta", "gamma", "delta"}
errs := runConcurrent(n, func(i int) error {
res, err := proxy.GetOrFetchArtifactFromURL(context.Background(),
"npm", names[i], "1.0.0", names[i]+"-1.0.0.tgz",
"https://registry.npmjs.org/"+names[i]+"/-/"+names[i]+"-1.0.0.tgz")
drain(res)
return err
})
for i, err := range errs {
if err != nil {
t.Errorf("caller %d (%s): %v", i, names[i], err)
}
}
if got := fetcher.calls.Load(); got != n {
t.Errorf("upstream fetches = %d, want %d: distinct artifacts must not share a fetch", got, n)
}
}
// TestCoalesce_FailedFetchReachesEveryCallerAndIsRetriable verifies both claims
// in coalesceFetch's doc comment: a failed fetch reaches every caller sharing
// it, and the key is released so a later request retries.
func TestCoalesce_FailedFetchReachesEveryCallerAndIsRetriable(t *testing.T) {
const callers = 8
proxy, _, _, fetcher := setupTestProxy(t)
boom := errors.New("upstream unavailable")
fetcher.fetchErr = boom
errs := runConcurrent(callers, func(int) error {
res, err := proxy.GetOrFetchArtifactFromURL(context.Background(),
"npm", "pkg", "1.0.0", "pkg-1.0.0.tgz",
"https://registry.npmjs.org/pkg/-/pkg-1.0.0.tgz")
drain(res)
return err
})
for i, err := range errs {
if err == nil {
t.Errorf("caller %d: got nil error, want the shared fetch's failure", i)
} else if !errors.Is(err, boom) {
t.Errorf("caller %d: got %v, want it to wrap %v", i, err, boom)
}
}
// The key must be released: a later request retries rather than inheriting
// the failure.
fetcher.fetchErr = nil
fetcher.artifact = artifactBody("recovered bytes")
res, err := proxy.GetOrFetchArtifactFromURL(context.Background(),
"npm", "pkg", "1.0.0", "pkg-1.0.0.tgz",
"https://registry.npmjs.org/pkg/-/pkg-1.0.0.tgz")
if err != nil {
t.Fatalf("retry after failed coalesced fetch: %v", err)
}
body, _ := io.ReadAll(res.Reader)
_ = res.Reader.Close()
if string(body) != "recovered bytes" {
t.Errorf("retry body = %q, want %q", body, "recovered bytes")
}
}
// TestCoalesce_ResolverPath covers the other entry point: GetOrFetchArtifact
// resolves the URL itself, so it is keyed without one.
func TestCoalesce_ResolverPath(t *testing.T) {
const callers = 8
proxy, _, _, _ := setupTestProxy(t)
fetcher := &countingFetcher{content: "resolved artifact bytes", delay: fetchHoldTime}
proxy.Fetcher = fetcher
errs := runConcurrent(callers, func(int) error {
res, err := proxy.GetOrFetchArtifact(context.Background(),
"npm", "left-pad", "1.3.0", "left-pad-1.3.0.tgz")
drain(res)
return err
})
for i, err := range errs {
if err != nil {
t.Errorf("caller %d: %v", i, err)
}
}
if got := fetcher.calls.Load(); got != 1 {
t.Errorf("upstream fetches = %d, want 1", got)
}
}
// TestCoalesce_ResolverPathEmptyFilename exercises that path when the filename
// is left to be resolved, which the key cannot know up front.
func TestCoalesce_ResolverPathEmptyFilename(t *testing.T) {
const callers = 8
proxy, _, _, _ := setupTestProxy(t)
fetcher := &countingFetcher{content: "resolved artifact bytes", delay: fetchHoldTime}
proxy.Fetcher = fetcher
errs := runConcurrent(callers, func(int) error {
res, err := proxy.GetOrFetchArtifact(context.Background(), "npm", "left-pad", "1.3.0", "")
drain(res)
return err
})
for i, err := range errs {
if err != nil {
t.Errorf("caller %d: %v", i, err)
}
}
if got := fetcher.calls.Load(); got != 1 {
t.Errorf("upstream fetches = %d, want 1", got)
}
}
// TestCoalesce_SubsequentRequestIsACacheHit confirms the coalesced fetch was
// committed and is visible later, not just streamed to the waiting callers.
func TestCoalesce_SubsequentRequestIsACacheHit(t *testing.T) {
const callers = 8
const url = "https://registry.npmjs.org/pkg/-/pkg-1.0.0.tgz"
proxy, _, _, _ := setupTestProxy(t)
fetcher := &countingFetcher{content: "artifact bytes", delay: fetchHoldTime}
proxy.Fetcher = fetcher
_ = runConcurrent(callers, func(int) error {
res, err := proxy.GetOrFetchArtifactFromURL(context.Background(),
"npm", "pkg", "1.0.0", "pkg-1.0.0.tgz", url)
drain(res)
return err
})
res, err := proxy.GetOrFetchArtifactFromURL(context.Background(),
"npm", "pkg", "1.0.0", "pkg-1.0.0.tgz", url)
if err != nil {
t.Fatalf("follow-up request: %v", err)
}
defer func() { _ = res.Reader.Close() }()
if !res.Cached {
t.Error("follow-up request should be served from cache")
}
if got := fetcher.calls.Load(); got != 1 {
t.Errorf("upstream fetches = %d, want 1 after a follow-up cache hit", got)
}
}
// TestCoalesce_ReadersAreIndependent guards openStoredArtifact: callers sharing
// a fetch each need their own reader, or one closing early breaks the rest.
func TestCoalesce_ReadersAreIndependent(t *testing.T) {
const callers = 8
const content = "artifact bytes that every caller must receive intact"
proxy, _, _, _ := setupTestProxy(t)
fetcher := &countingFetcher{content: content, delay: fetchHoldTime}
proxy.Fetcher = fetcher
results := make([]*CacheResult, callers)
errs := runConcurrent(callers, func(i int) error {
res, err := proxy.GetOrFetchArtifactFromURL(context.Background(),
"npm", "pkg", "1.0.0", "pkg-1.0.0.tgz",
"https://registry.npmjs.org/pkg/-/pkg-1.0.0.tgz")
results[i] = res
return err
})
for i, err := range errs {
if err != nil {
t.Fatalf("caller %d: %v", i, err)
}
}
// Close the first caller's reader before anyone else has read a byte.
_ = results[0].Reader.Close()
for i := 1; i < callers; i++ {
body, err := io.ReadAll(results[i].Reader)
_ = results[i].Reader.Close()
if err != nil {
t.Errorf("caller %d read after another caller closed: %v", i, err)
continue
}
if string(body) != content {
t.Errorf("caller %d got %q, want %q", i, body, content)
}
}
}
// TestCoalesce_CanceledWaiterDoesNotWaitForTheSharedFetch checks that joining a
// coalesced fetch does not cost a caller its own cancellation. Without the
// leader/waiter split a waiter is pinned until the shared fetch resolves,
// bounded only by the artifact client timeout, so clients that have already
// gone away keep handler goroutines alive for minutes.
func TestCoalesce_CanceledWaiterDoesNotWaitForTheSharedFetch(t *testing.T) {
const leaderFetch = 2 * time.Second
const url = "https://registry.npmjs.org/pkg/-/pkg-1.0.0.tgz"
proxy, _, _, _ := setupTestProxy(t)
fetcher := &countingFetcher{content: "artifact bytes", delay: leaderFetch, entered: make(chan struct{})}
proxy.Fetcher = fetcher
leaderDone := make(chan error, 1)
go func() {
res, err := proxy.GetOrFetchArtifactFromURL(context.Background(),
"npm", "pkg", "1.0.0", "pkg-1.0.0.tgz", url)
drain(res)
leaderDone <- err
}()
select {
case <-fetcher.entered: // the leader holds the key and is inside its fetch
case <-time.After(5 * time.Second):
t.Fatal("leader never started its fetch")
}
ctx, cancel := context.WithCancel(context.Background())
cancel()
start := time.Now()
_, err := proxy.GetOrFetchArtifactFromURL(ctx, "npm", "pkg", "1.0.0", "pkg-1.0.0.tgz", url)
blocked := time.Since(start)
if !errors.Is(err, context.Canceled) {
t.Errorf("waiter error = %v, want context.Canceled", err)
}
if blocked > leaderFetch/4 {
t.Errorf("canceled waiter blocked %v, want well under %v: it is pinned to the shared fetch",
blocked, leaderFetch/4)
}
// A waiter leaving must not disturb the fetch the others share.
if err := <-leaderDone; err != nil {
t.Fatalf("leader failed after a waiter canceled: %v", err)
}
res, err := proxy.GetOrFetchArtifactFromURL(context.Background(),
"npm", "pkg", "1.0.0", "pkg-1.0.0.tgz", url)
if err != nil {
t.Fatalf("follow-up after leader completed: %v", err)
}
defer func() { _ = res.Reader.Close() }()
if !res.Cached {
t.Error("leader's fetch should have been committed to the cache")
}
if got := fetcher.calls.Load(); got != 1 {
t.Errorf("upstream fetches = %d, want 1", got)
}
}
// inFlightLen reports how many coalesced fetches are currently registered.
func inFlightLen(p *Proxy) int {
p.fetchMu.Lock()
defer p.fetchMu.Unlock()
return len(p.inFlight)
}
// TestCoalesce_KeyIsReleasedAfterFetch guards the bug this hand-rolled map can
// have that singleflight could not: a key left behind means later callers join
// a finished entry, see its closed done channel, and are served that stale
// result forever, while the map grows without bound.
func TestCoalesce_KeyIsReleasedAfterFetch(t *testing.T) {
const url = "https://registry.npmjs.org/pkg/-/pkg-1.0.0.tgz"
proxy, _, _, _ := setupTestProxy(t)
fetcher := &countingFetcher{content: "artifact bytes", delay: fetchHoldTime}
proxy.Fetcher = fetcher
_ = runConcurrent(8, func(int) error {
res, err := proxy.GetOrFetchArtifactFromURL(context.Background(),
"npm", "pkg", "1.0.0", "pkg-1.0.0.tgz", url)
drain(res)
return err
})
if n := inFlightLen(proxy); n != 0 {
t.Errorf("in-flight entries after a successful fetch = %d, want 0", n)
}
// A fresh miss for the same key must start a new fetch, not rejoin the old
// entry. Clearing the cache record forces the miss path again.
if err := proxy.ClearCachedArtifact(context.Background(), "npm", "pkg", "1.0.0", "pkg-1.0.0.tgz"); err != nil {
t.Fatalf("clear cached artifact: %v", err)
}
res, err := proxy.GetOrFetchArtifactFromURL(context.Background(),
"npm", "pkg", "1.0.0", "pkg-1.0.0.tgz", url)
if err != nil {
t.Fatalf("second miss for the same key: %v", err)
}
drain(res)
if got := fetcher.calls.Load(); got != 2 {
t.Errorf("upstream fetches = %d, want 2: the second miss must not reuse the finished entry", got)
}
if n := inFlightLen(proxy); n != 0 {
t.Errorf("in-flight entries at end = %d, want 0", n)
}
}
// missingFromCache is a recheck that always reports a miss, so the shared fetch
// runs.
func missingFromCache() (artifacts.Artifact, string, bool) {
return artifacts.Artifact{}, "", false
}
// TestCoalesce_LeaderRechecksCacheBeforeFetching covers the window between a
// caller's own cache lookup and it becoming the leader: a concurrent fetch can
// commit the artifact in that gap, and the leader must serve that rather than
// fetch it a second time.
func TestCoalesce_LeaderRechecksCacheBeforeFetching(t *testing.T) {
const content = "artifact bytes"
proxy, _, store, _ := setupTestProxy(t)
const storagePath = "npm/pkg/1.0.0/pkg-1.0.0.tgz"
if _, _, err := store.Store(context.Background(), storagePath, strings.NewReader(content)); err != nil {
t.Fatalf("seeding storage: %v", err)
}
committed := artifacts.Artifact{
PURL: "pkg:npm/pkg@1.0.0",
Filename: "pkg-1.0.0.tgz",
Size: int64(len(content)),
}
res, err := proxy.coalesceFetch(context.Background(), "any-key",
func() (artifacts.Artifact, string, bool) { return committed, storagePath, true },
func(context.Context) (artifacts.Artifact, string, error) {
t.Error("fetched an artifact that was already in the cache")
return artifacts.Artifact{}, "", errors.New("commit must not run")
})
if err != nil {
t.Fatalf("coalesceFetch failed: %v", err)
}
defer drain(res)
got, err := io.ReadAll(res.Reader)
if err != nil {
t.Fatalf("reading result: %v", err)
}
if string(got) != content {
t.Errorf("got %q, want %q", got, content)
}
if n := inFlightLen(proxy); n != 0 {
t.Errorf("in-flight entries = %d, want 0", n)
}
}
// TestCachedArtifactRecord covers the recheck itself: it must report the row a
// concurrent fetch committed, match its digest the way artifactHashMatches
// does, and report a miss for anything else.
func TestCachedArtifactRecord(t *testing.T) {
const (
content = "artifact bytes"
pkgPURL = "pkg:npm/pkg"
versionPURL = "pkg:npm/pkg@1.0.0"
filename = "pkg-1.0.0.tgz"
storagePath = "npm/pkg/1.0.0/pkg-1.0.0.tgz"
)
proxy, _, _, _ := setupTestProxy(t)
hex := sha256Hex(content)
committed := testArtifact(content, versionPURL, filename, "application/gzip")
if err := proxy.updateCacheDB("npm", "pkg", pkgPURL,
"https://registry.npmjs.org/pkg/-/pkg-1.0.0.tgz", storagePath, committed); err != nil {
t.Fatalf("seeding cache record: %v", err)
}
for _, tc := range []struct {
name, filename, hash string
want bool
}{
{"no upstream hash", filename, "", true},
{"matching hash", filename, hex, true},
{"matching hash in upper case", filename, strings.ToUpper(hex), true},
{"different hash", filename, sha256Hex("something else entirely"), false},
{"unknown filename", "pkg-1.0.0.zip", hex, false},
} {
t.Run(tc.name, func(t *testing.T) {
got, path, ok := proxy.cachedArtifactRecord(pkgPURL, versionPURL, tc.filename, tc.hash)
if ok != tc.want {
t.Fatalf("ok = %v, want %v", ok, tc.want)
}
if !ok {
return
}
if path != storagePath {
t.Errorf("storage path = %q, want %q", path, storagePath)
}
if got.Digest.Encoded() != hex {
t.Errorf("digest = %q, want %q", got.Digest.Encoded(), hex)
}
})
}
}
// TestCoalesce_LeaderFetchesWhenRecheckedBytesAreGone covers the other branch
// of the recheck: a record whose bytes no longer open is not served, and the
// shared fetch runs instead, the same recovery the cache lookup makes.
func TestCoalesce_LeaderFetchesWhenRecheckedBytesAreGone(t *testing.T) {
const content = "fetched bytes"
const storagePath = "npm/pkg/1.0.0/pkg-1.0.0.tgz"
proxy, _, store, _ := setupTestProxy(t)
stale := artifacts.Artifact{PURL: "pkg:npm/pkg@1.0.0", Filename: "pkg-1.0.0.tgz"}
if _, err := store.Open(context.Background(), storagePath); err == nil {
t.Fatal("stale bytes were present, so the test proves nothing")
}
// The leader runs commit on its own goroutine, so a plain counter is safe.
fetches := 0
res, err := proxy.coalesceFetch(context.Background(), "any-key",
func() (artifacts.Artifact, string, bool) { return stale, storagePath, true },
func(ctx context.Context) (artifacts.Artifact, string, error) {
fetches++
if _, _, err := store.Store(ctx, storagePath, strings.NewReader(content)); err != nil {
return artifacts.Artifact{}, "", err
}
return testArtifact(content, stale.PURL, stale.Filename, "application/gzip"), storagePath, nil
})
if err != nil {
t.Fatalf("coalesceFetch failed: %v", err)
}
defer drain(res)
if fetches != 1 {
t.Errorf("shared fetches = %d, want 1: a record without bytes must be refetched", fetches)
}
got, err := io.ReadAll(res.Reader)
if err != nil {
t.Fatalf("reading result: %v", err)
}
if string(got) != content {
t.Errorf("got %q, want %q", got, content)
}
if n := inFlightLen(proxy); n != 0 {
t.Errorf("in-flight entries = %d, want 0", n)
}
}
// TestCoalesce_PanicInSharedFetchDoesNotStrandWaiters checks the failure mode
// that matters most: a caller parked on a shared fetch must never be left
// blocked forever when that fetch dies.
//
// This drives coalesceFetch directly and holds the shared entry itself, because
// whether a second caller has reached the wait is not observable from outside:
// it runs a cache lookup against the database first, so releasing the leader on
// a timer races that query. Losing the race made a second caller the leader
// instead of a waiter, and its panic was unrecovered, killing the test binary
// rather than failing the test.
func TestCoalesce_PanicInSharedFetchDoesNotStrandWaiters(t *testing.T) {
proxy, _, _, _ := setupTestProxy(t)
const key = "pkg:npm/pkg@1.0.0\x00pkg-1.0.0.tgz"
inCommit := make(chan struct{})
release := make(chan struct{})
leaderPanicked := make(chan struct{})
go func() {
defer func() {
_ = recover() // the panic surfaces in the leader, as it would in a handler
close(leaderPanicked)
}()
_, _ = proxy.coalesceFetch(context.Background(), key, missingFromCache,
func(context.Context) (artifacts.Artifact, string, error) {
close(inCommit)
<-release
panic("upstream fetch exploded")
})
}()
<-inCommit // the leader holds the key and is inside the fetch
// Take the entry a waiter would park on, while the leader is still held.
proxy.fetchMu.Lock()
shared := proxy.inFlight[key]
proxy.fetchMu.Unlock()
if shared == nil {
t.Fatal("no in-flight entry registered for a running fetch")
}
close(release)
select {
case <-shared.done:
case <-time.After(5 * time.Second):
t.Fatal("waiter stranded: a panicking shared fetch never released its waiters")
}
if !errors.Is(shared.err, errSharedFetchAbandoned) {
t.Errorf("waiter error = %v, want errSharedFetchAbandoned", shared.err)
}
<-leaderPanicked
if n := inFlightLen(proxy); n != 0 {
t.Errorf("in-flight entries after a panic = %d, want 0", n)
}
}