Files
at-container-registry/pkg/appview/holdpurge/queue_test.go
T
Evan JarrettandClaude Opus 5 3589473feb appview: run the hold purge on a worker pool, not the request context
Deleting a tag while over quota could leave the user in the worst available
state. DeleteTagHandler deleted the tag and manifest rows first, then called
PurgeOnHold, which bounded itself at 10s against the *request* context. The
UpCloud load balancer in front of the appview cuts at its default backend
timeout at about the same moment, wins the race, hands the client a 504 and
cancels that context, killing the purge partway. The appview logged a warning
and returned 200.

So: gateway error, nothing freed, still locked out, the image gone from the UI
so the purge cannot be retried through it, blobs orphaned until the hold's GC,
and the appview considering it a success. Measured on production at
10.002367218s.

Purges now go to a fixed pool of 4 workers rooted at context.Background(), so
they survive the request ending. Following the shape of the hold's startJob
helper, minus the progress fragment, since nobody is watching a purge.

The buffer is bounded at 256 and sheds with an ERROR rather than growing: an
unbounded queue turns a slow hold into an appview memory leak. Submissions are
deduplicated on holdDID|manifestURI so a double-clicked delete does one purge
and one service-token fetch. The channel send happens under the mutex that
guards close, so a concurrent drain cannot send on a closed channel, and the
drain is wired into both exit paths before logging shuts down.

Failures are now classified and surfaced instead of swallowed: transient ones
retry three times under a 90s budget (the hold's purge is idempotent), an
unauthorized third-party hold logs at DEBUG since it is expected, and anything
else that exhausts its retries logs at ERROR naming the manifest and hold, which
is enough to re-drive by hand.

Deliberately not reordered. Purge-first-then-delete requires waiting for the
purge to know whether to delete, which puts the 10s call straight back on the
request. So the orphaned-blob window remains, materially narrower but real: a
purge that exhausts its retries still leaves blobs referenced by nothing until
the hold's GC, and there is no row left to say so. Closing that needs a durable
pending-purge record, which was judged out of proportion here.

server.go in this commit also carries one line belonging to the next one, the
token handler's display-name wiring, since the two changes landed in the same
file concurrently.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01PDqoCE1j3njokkZ9b1C5n9
2026-09-02 22:31:27 -05:00

368 lines
9.9 KiB
Go

package holdpurge
import (
"bytes"
"context"
"errors"
"fmt"
"log/slog"
"strings"
"sync"
"testing"
"time"
)
func testRequest(manifest string) Request {
return Request{
UserDID: "did:plc:testuser",
PDSEndpoint: "https://pds.example",
HoldDID: "did:web:hold.example",
ManifestURI: manifest,
}
}
// captureLogs swaps the default slog handler for the duration of a test and
// returns a function yielding everything logged so far.
func captureLogs(t *testing.T, level slog.Level) func() string {
t.Helper()
var mu sync.Mutex
buf := &bytes.Buffer{}
prev := slog.Default()
slog.SetDefault(slog.New(slog.NewTextHandler(&lockedWriter{mu: &mu, buf: buf}, &slog.HandlerOptions{Level: level})))
t.Cleanup(func() { slog.SetDefault(prev) })
return func() string {
mu.Lock()
defer mu.Unlock()
return buf.String()
}
}
type lockedWriter struct {
mu *sync.Mutex
buf *bytes.Buffer
}
func (w *lockedWriter) Write(p []byte) (int, error) {
w.mu.Lock()
defer w.mu.Unlock()
return w.buf.Write(p)
}
// TestSubmitReturnsWithoutWaitingForPurge is the shape of the fix: the delete
// handler hands the purge off and returns. Before, it blocked on the hold for
// up to AttemptTimeout, which is longer than the proxy in front of the appview
// waits — so the user got a 504 for a delete that had succeeded.
func TestSubmitReturnsWithoutWaitingForPurge(t *testing.T) {
release := make(chan struct{})
started := make(chan struct{})
done := make(chan struct{})
q := newQueue(1, 4, func(ctx context.Context, req Request) error {
close(started)
<-release
close(done)
return nil
})
start := time.Now()
if !q.Submit(testRequest("at://did:plc:testuser/io.atcr.manifest/abc")) {
t.Fatal("Submit rejected the purge")
}
elapsed := time.Since(start)
// Generous bound: the point is "does not wait for the hold", and the
// worker below is blocked indefinitely until we release it.
if elapsed > time.Second {
t.Fatalf("Submit blocked for %v; it must hand off and return", elapsed)
}
select {
case <-started:
case <-time.After(5 * time.Second):
t.Fatal("purge never started on a worker")
}
close(release)
select {
case <-done:
case <-time.After(5 * time.Second):
t.Fatal("purge never finished")
}
q.Wait(5 * time.Second)
}
// TestPurgeRunsAfterRequestContextCancelled is the defect itself: the proxy
// cutting the request cancelled the context the purge was running on, so the
// purge died mid-flight. The queued purge must run on a context rooted at
// context.Background(), unaffected by the request ending.
func TestPurgeRunsAfterRequestContextCancelled(t *testing.T) {
type observation struct {
errAtStart error
errLater error
}
observed := make(chan observation, 1)
q := newQueue(1, 4, func(ctx context.Context, req Request) error {
o := observation{errAtStart: ctx.Err()}
// Give the cancelled request context every chance to propagate.
time.Sleep(50 * time.Millisecond)
o.errLater = ctx.Err()
observed <- o
return nil
})
// Stand in for the request goroutine: enqueue, then die. Submit takes no
// context at all, which is the structural half of the fix — there is no
// longer a way for a handler to hand the purge its request deadline.
_, cancelRequest := context.WithCancel(context.Background())
if !q.Submit(testRequest("at://did:plc:testuser/io.atcr.manifest/cancelled")) {
t.Fatal("Submit rejected the purge")
}
cancelRequest()
select {
case o := <-observed:
if o.errAtStart != nil {
t.Fatalf("purge ran on an already-cancelled context: %v", o.errAtStart)
}
if o.errLater != nil {
t.Fatalf("purge context was cancelled by the request ending: %v", o.errLater)
}
case <-time.After(5 * time.Second):
t.Fatal("purge did not run after the request context was cancelled")
}
q.Wait(5 * time.Second)
}
// TestPurgeFailureIsSurfaced: a purge that never succeeds must not disappear.
// It retries, and the final failure is logged at ERROR with the manifest URI,
// because the appview has no durable record of work still owed to the hold.
func TestPurgeFailureIsSurfaced(t *testing.T) {
logs := captureLogs(t, slog.LevelDebug)
var attempts int
var mu sync.Mutex
finished := make(chan struct{})
q := newQueue(1, 4, func(ctx context.Context, req Request) error {
mu.Lock()
attempts++
n := attempts
mu.Unlock()
if n == maxAttempts {
defer close(finished)
}
return errors.New("hold unreachable")
})
q.backoff = time.Millisecond
if !q.Submit(testRequest("at://did:plc:testuser/io.atcr.manifest/doomed")) {
t.Fatal("Submit rejected the purge")
}
select {
case <-finished:
case <-time.After(5 * time.Second):
t.Fatal("purge did not exhaust its attempts")
}
q.Wait(5 * time.Second)
mu.Lock()
got := attempts
mu.Unlock()
if got != maxAttempts {
t.Errorf("attempts = %d, want %d", got, maxAttempts)
}
out := logs()
if !strings.Contains(out, "level=ERROR") {
t.Errorf("failed purge was not logged at ERROR:\n%s", out)
}
if !strings.Contains(out, "purge failed after retries") {
t.Errorf("failed purge did not name itself in the log:\n%s", out)
}
if !strings.Contains(out, "io.atcr.manifest/doomed") {
t.Errorf("failed purge log does not identify the manifest:\n%s", out)
}
}
// A permanent failure (no OAuth refresher, a malformed request, a 4xx from the
// hold) must not burn retries, but must still be surfaced.
func TestPermanentFailureIsNotRetried(t *testing.T) {
logs := captureLogs(t, slog.LevelDebug)
var mu sync.Mutex
var attempts int
q := newQueue(1, 4, func(ctx context.Context, req Request) error {
mu.Lock()
attempts++
mu.Unlock()
return permanent(errors.New("malformed"))
})
q.backoff = time.Millisecond
if !q.Submit(testRequest("at://did:plc:testuser/io.atcr.manifest/permanent")) {
t.Fatal("Submit rejected the purge")
}
q.Wait(5 * time.Second)
mu.Lock()
got := attempts
mu.Unlock()
if got != 1 {
t.Errorf("attempts = %d, want 1 for a permanent failure", got)
}
if !strings.Contains(logs(), "level=ERROR") {
t.Errorf("permanent failure was not surfaced:\n%s", logs())
}
}
// A sailor purging on a third-party hold has no right to; that is expected and
// handled by the hold's own GC, so it must not page anyone.
func TestNotAuthorizedIsNotAnError(t *testing.T) {
logs := captureLogs(t, slog.LevelDebug)
q := newQueue(1, 4, func(ctx context.Context, req Request) error {
return ErrNotAuthorized
})
q.backoff = time.Millisecond
if !q.Submit(testRequest("at://did:plc:testuser/io.atcr.manifest/thirdparty")) {
t.Fatal("Submit rejected the purge")
}
q.Wait(5 * time.Second)
if strings.Contains(logs(), "level=ERROR") {
t.Errorf("an unauthorized third-party purge should not log at ERROR:\n%s", logs())
}
}
// Rapid repeat deletes of the same manifest must not each spawn work.
func TestSubmitDeduplicatesInFlightManifest(t *testing.T) {
release := make(chan struct{})
var mu sync.Mutex
var runs int
q := newQueue(1, 8, func(ctx context.Context, req Request) error {
mu.Lock()
runs++
mu.Unlock()
<-release
return nil
})
req := testRequest("at://did:plc:testuser/io.atcr.manifest/dupe")
if !q.Submit(req) {
t.Fatal("first Submit rejected")
}
if q.Submit(req) {
t.Error("second Submit for an in-flight manifest should be rejected")
}
close(release)
q.Wait(5 * time.Second)
mu.Lock()
defer mu.Unlock()
if runs != 1 {
t.Errorf("runs = %d, want 1", runs)
}
}
// A saturated queue sheds load rather than growing without bound, and says so.
func TestSubmitShedsWhenQueueIsFull(t *testing.T) {
logs := captureLogs(t, slog.LevelDebug)
release := make(chan struct{})
q := newQueue(1, 1, func(ctx context.Context, req Request) error {
<-release
return nil
})
// One job occupies the worker, one fills the single buffer slot, the
// third has nowhere to go.
accepted := 0
for i := 0; i < 3; i++ {
if q.Submit(testRequest(fmt.Sprintf("at://did:plc:testuser/io.atcr.manifest/%d", i))) {
accepted++
}
// Let the worker pick the first job up so the buffer is the limit.
if i == 0 {
time.Sleep(50 * time.Millisecond)
}
}
if accepted > 2 {
t.Errorf("accepted %d submissions into a queue of depth 1", accepted)
}
if !strings.Contains(logs(), "purge queue full") {
t.Errorf("shedding was silent:\n%s", logs())
}
close(release)
q.Wait(5 * time.Second)
}
// Wait must drain the queued purges rather than let SIGTERM drop them, and it
// must return even when a purge is wedged.
func TestWaitDrainsAndDoesNotLeak(t *testing.T) {
var mu sync.Mutex
var completed int
q := newQueue(2, 8, func(ctx context.Context, req Request) error {
time.Sleep(20 * time.Millisecond)
mu.Lock()
completed++
mu.Unlock()
return nil
})
for i := 0; i < 4; i++ {
if !q.Submit(testRequest(fmt.Sprintf("at://did:plc:testuser/io.atcr.manifest/drain%d", i))) {
t.Fatalf("Submit %d rejected", i)
}
}
q.Wait(5 * time.Second)
mu.Lock()
got := completed
mu.Unlock()
if got != 4 {
t.Errorf("completed = %d, want 4 (shutdown dropped queued purges)", got)
}
// Post-drain submissions are refused rather than panicking on a closed
// channel, and a second Wait is a no-op.
if q.Submit(testRequest("at://did:plc:testuser/io.atcr.manifest/after")) {
t.Error("Submit after Wait should be rejected")
}
q.Wait(time.Second)
}
// Wait must not hang forever on a wedged purge; it gives up and cancels the
// worker's context so the goroutine cannot outlive shutdown.
func TestWaitGivesUpOnWedgedPurge(t *testing.T) {
observedCancel := make(chan struct{})
q := newQueue(1, 2, func(ctx context.Context, req Request) error {
<-ctx.Done()
close(observedCancel)
return ctx.Err()
})
q.backoff = time.Millisecond
if !q.Submit(testRequest("at://did:plc:testuser/io.atcr.manifest/wedged")) {
t.Fatal("Submit rejected")
}
start := time.Now()
q.Wait(100 * time.Millisecond)
if elapsed := time.Since(start); elapsed > 3*time.Second {
t.Fatalf("Wait blocked for %v past its grace period", elapsed)
}
select {
case <-observedCancel:
case <-time.After(5 * time.Second):
t.Fatal("wedged purge was never cancelled after the drain grace period")
}
}