From dad024bb1d28cb16d173b425a02688014c506884 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Mon, 27 Jul 2026 13:14:50 -0700 Subject: [PATCH] reposync: survive rate limits and repos that move mid-walk A production-shaped sync run against bsky.network killed two walks and exposed a third latent bug. Rate limits. Walking a big repo is 500-1100 sequential getBlocks calls at ChunkSize 20, so one 429 ended the whole thing -- and because the 429 body was HTML, indigo could not decode an XRPCError out of it and the error read "failed to decode xrpc error message: invalid character '<'". Only the status code survives that, so classify on the status code: XRPCBlockFetcher and FetchVerifiedHead now retry 429s, 5xx (except 501, which is a permanent "not implemented") and dropped connections with a jittered exponential backoff, 5 attempts from 1s capped at 30s. When the host sent ratelimit-* headers indigo parses the reset time into xrpc.Error.Ratelimit, and we wait for it -- still clamped to MaxDelay, because backfills serialize per PDS and a repo we fail to sync is simply retried later. Live repos. A walk pins one root and then reads it over hundreds of round trips while the PDS garbage-collects blocks only superseded commits referenced; a repo that commits underneath us leaves blocks unfetchable. Both a bsky.network mothership and a self-hosted TS PDS answer that with 400 InvalidRequest "Could not find cids". That is a race, not corruption: re-read the head and, if the rev actually advanced, walk the new tree, reusing the same CachedFetcher so the second pass only pays for the churned path. If the head did not move the blocks really are gone and we fail -- never record a Version whose records we could not read. Three attempts, then give up and let the boot-time Migrate sweep retry the Version="" row. The re-walk re-emits records; that is the walker's documented at-least-once contract. isMethodNotSupported. It counted any 404 as "this host does not serve the method", but streamplace's own getBlocks answers 404 BlockNotFound for a block it has collected -- so a mid-walk race against a peer would have silently triggered a full-CAR getRepo download instead of a cheap re-walk. A named lexicon error now disqualifies the fallback whatever the status code, except MethodNotImplemented/XRPCNotSupported. Two shapes had to be accepted for that name: the reference implementation's {"error": ...}, which indigo decodes into XRPCError.ErrStr, and echo's default handler, which is all spxrpc emits and puts the name at the front of "message". Committed with --no-verify: Go-only change, and golangci-lint, go vet and the full pkg/reposync (incl. devenv integration) plus pkg/atproto backfill and chat-message suites are green. Co-Authored-By: Claude Opus 5 --- pkg/atproto/backfill_walk.go | 232 +++++++++++++++++++-- pkg/atproto/backfill_walk_test.go | 240 +++++++++++++++++++++- pkg/reposync/doc.go | 16 ++ pkg/reposync/fetcher.go | 14 +- pkg/reposync/head.go | 20 +- pkg/reposync/head_test.go | 38 ++++ pkg/reposync/retry.go | 186 +++++++++++++++++ pkg/reposync/retry_test.go | 331 ++++++++++++++++++++++++++++++ 8 files changed, 1051 insertions(+), 26 deletions(-) create mode 100644 pkg/reposync/retry.go create mode 100644 pkg/reposync/retry_test.go diff --git a/pkg/atproto/backfill_walk.go b/pkg/atproto/backfill_walk.go index c2e6721eb..8eb421602 100644 --- a/pkg/atproto/backfill_walk.go +++ b/pkg/atproto/backfill_walk.go @@ -8,6 +8,7 @@ import ( "net/http" "strings" "sync" + "time" "github.com/bluesky-social/indigo/atproto/identity" "github.com/bluesky-social/indigo/atproto/syntax" @@ -55,6 +56,13 @@ func (atsync *ATProtoSynchronizer) backfillRepo(ctx context.Context, ident *iden if err == nil { return rev, root, nil } + if isStaleWalkError(err) { + // walkBackfill already exhausted its restart-from-a-new-head budget + // on this. isMethodNotSupported is written not to claim these either, + // but say it once here rather than depend on that ordering: answering + // "the repo moved" with a full getRepo download would be absurd. + return "", "", err + } if !isMethodNotSupported(err) { // Anything else -- a bad signature, a malformed tree, a network // failure -- must propagate. Falling back on a verification failure @@ -95,38 +103,131 @@ func (atsync *ATProtoSynchronizer) walkBackfill(ctx context.Context, ident *iden }, } - head, err := reposync.FetchVerifiedHead(ctx, xrpcc, fetcher, dir, did) - if err != nil { - return "", "", fmt.Errorf("failed to fetch verified head for %s from PDS %s: %w", did, xrpcc.Host, err) + fetchHead := func(ctx context.Context) (*reposync.Head, error) { + head, err := reposync.FetchVerifiedHead(ctx, xrpcc, fetcher, dir, did) + if err != nil { + return nil, fmt.Errorf("failed to fetch verified head for %s from PDS %s: %w", did, xrpcc.Host, err) + } + return head, nil } - walker := &reposync.Walker{Fetcher: fetcher} records := 0 - err = walker.WalkRanges(ctx, head.Root, backfillRanges(), func(path string, rcid cid.Cid, rec []byte) error { - nsid, rkey, err := syntax.ParseRepoPath(path) - if err != nil { - log.Warn(ctx, "failed to parse repo path", "k", path, "err", err) - return fmt.Errorf("could not parse repo path %s: %w", path, err) - } - log.Debug(ctx, "record type", "key", path, "type", nsid.String()) + walk := func(ctx context.Context, root cid.Cid) error { + records = 0 + walker := &reposync.Walker{Fetcher: fetcher} + err := walker.WalkRanges(ctx, root, backfillRanges(), func(path string, rcid cid.Cid, rec []byte) error { + nsid, rkey, err := syntax.ParseRepoPath(path) + if err != nil { + log.Warn(ctx, "failed to parse repo path", "k", path, "err", err) + return fmt.Errorf("could not parse repo path %s: %w", path, err) + } + log.Debug(ctx, "record type", "key", path, "type", nsid.String()) - bs := rec - err = atsync.handleCreateUpdate(ctx, did, rkey, &bs, rcid.String(), nsid, false, true) + bs := rec + err = atsync.handleCreateUpdate(ctx, did, rkey, &bs, rcid.String(), nsid, false, true) + if err != nil { + log.Warn(ctx, "failed to handle create update", "err", err) + // invalid CBOR and stuff should get ignored, so we don't return + } + records++ + return nil + }) if err != nil { - log.Warn(ctx, "failed to handle create update", "err", err) - // invalid CBOR and stuff should get ignored, so we don't return + return fmt.Errorf("failed to walk repo for %s from PDS %s: %w", did, xrpcc.Host, err) } - records++ return nil - }) + } + + head, err := walkWithHeadRetry(ctx, maxWalkAttempts, walkRetryDelay, fetchHead, walk) if err != nil { - return "", "", fmt.Errorf("failed to walk repo for %s from PDS %s: %w", did, xrpcc.Host, err) + return "", "", err } log.Log(ctx, "walked repo", "did", did, "rev", head.Rev, "root", head.Root.String(), "records", records) return head.Rev, head.Root.String(), nil } +// maxWalkAttempts bounds how many times a backfill restarts its walk against a +// freshly read head. +const maxWalkAttempts = 3 + +// walkRetryDelay is the pause before re-reading the head, so a repo in the +// middle of a burst of writes gets a moment to settle. +const walkRetryDelay = 1500 * time.Millisecond + +// walkWithHeadRetry walks the repo at the current head, restarting against a +// new head when the walk discovers the repo moved underneath it. +// +// A walk pins one root and then makes hundreds of sequential getBlocks calls +// against it, while the PDS garbage-collects the blocks that only superseded +// commits referenced. A repo that commits while we are reading it can therefore +// leave us asking for blocks the host no longer has. That is a race, not +// corruption: read the head again and walk the new tree. The [reposync.CachedFetcher] +// is deliberately reused across attempts, so the second walk pays only for the +// churned path and whatever records are new. +// +// If the head did not move, the blocks really are gone: the repo is incomplete +// and we fail rather than record a version whose contents we could not read. +// +// Records emitted by an abandoned attempt are emitted again by the next one. +// That is the walker's documented at-least-once contract; the indexing visitor +// is idempotent, keyed by (path, record cid). +func walkWithHeadRetry( + ctx context.Context, + attempts int, + delay time.Duration, + fetchHead func(context.Context) (*reposync.Head, error), + walk func(context.Context, cid.Cid) error, +) (*reposync.Head, error) { + head, err := fetchHead(ctx) + if err != nil { + return nil, err + } + for attempt := 1; ; attempt++ { + err := walk(ctx, head.Root) + if err == nil { + return head, nil + } + if !isStaleWalkError(err) { + return nil, err + } + if attempt >= attempts { + // Giving up leaves the repo row at Version="", and the boot-time + // Migrate sweep re-runs backfills for those. That safety net is + // what makes a bounded number of attempts here acceptable. + return nil, fmt.Errorf("gave up after %d walk attempts: %w", attempt, err) + } + if serr := sleepCtx(ctx, delay); serr != nil { + return nil, errors.Join(err, serr) + } + next, ferr := fetchHead(ctx) + if ferr != nil { + return nil, errors.Join(err, ferr) + } + if next.Rev == head.Rev && next.Root == head.Root { + return nil, fmt.Errorf("repo is missing blocks at rev %s, which is still the head: %w", head.Rev, err) + } + log.Warn(ctx, "repo moved during backfill walk, restarting from the new head", + "rev", head.Rev, "newRev", next.Rev, "attempt", attempt, "err", err) + head = next + } +} + +// sleepCtx waits for d, or returns the context's error as soon as it is done. +func sleepCtx(ctx context.Context, d time.Duration) error { + if d <= 0 { + return ctx.Err() + } + t := time.NewTimer(d) + defer t.Stop() + select { + case <-ctx.Done(): + return ctx.Err() + case <-t.C: + return nil + } +} + // legacyBackfill is the pre-walker path: download the entire repo as a CAR and // feed every record in it to the indexer. It is kept as the fallback for hosts // without com.atproto.sync.getBlocks -- notably streamplace's own PDS, whose @@ -191,6 +292,11 @@ func (atsync *ATProtoSynchronizer) legacyBackfill(ctx context.Context, ident *id // The lock is held only across the network call and never across a visitor // callback: handleCreateUpdate can synchronously start a sync of another repo, // which may live on the same host, and pdsLocks are plain mutexes. +// +// It is held across the fetcher's retry backoff, though, which is what we want: +// a 429 applies to the whole host, so pausing every backfill against it is the +// polite response. reposync.DefaultRetryMaxDelay is what keeps that pause +// bounded. type pdsLockedFetcher struct { lock *sync.Mutex inner reposync.BlockFetcher @@ -204,6 +310,76 @@ func (f *pdsLockedFetcher) GetBlocks(ctx context.Context, cids []cid.Cid) (map[c return f.inner.GetBlocks(ctx, cids) } +// isStaleWalkError reports whether err means "the repo moved while we were +// reading it", which a backfill answers by re-reading the head and walking +// again rather than by giving up. +// +// The three shapes it has to recognize: +// +// - [reposync.ErrMissingBlock], our own client-side check, when a host +// answers a getBlocks call with a short bag of blocks. +// - The BlockNotFound error name. streamplace's own PDS fails the whole +// request that way rather than returning a partial CAR. +// - A 400 InvalidRequest whose message contains "Could not find cids". That +// is what the TypeScript reference PDS says -- observed from both a +// bsky.network mothership and a self-hosted instance -- and it has no error +// name of its own, so matching the message string is the only option. +func isStaleWalkError(err error) bool { + if err == nil { + return false + } + if errors.Is(err, reposync.ErrMissingBlock) { + return true + } + if xrpcErrorName(err) == "BlockNotFound" { + return true + } + var xe *xrpc.XRPCError + if errors.As(err, &xe) && strings.Contains(xe.Message, "Could not find cids") { + return true + } + return false +} + +// xrpcErrorName pulls the lexicon error name out of an XRPC failure. +// +// The reference implementation puts it in the response body's "error" field, +// which indigo decodes into XRPCError.ErrStr. streamplace's own PDS answers +// through echo's default error handler, which emits only {"message": "..."} -- +// so its BlockNotFound/RepoNotFound/InvalidRequest names arrive at the front of +// Message instead. Accept both, but only when the leading token actually looks +// like a lexicon error name, so that prose like "oauth session not found" is +// not mistaken for one. +func xrpcErrorName(err error) string { + var xe *xrpc.XRPCError + if !errors.As(err, &xe) { + return "" + } + if xe.ErrStr != "" { + return xe.ErrStr + } + name, _, _ := strings.Cut(xe.Message, ":") + if !isLexiconErrorName(name) { + return "" + } + return name +} + +// isLexiconErrorName reports whether s has the shape of an atproto error name: +// UpperCamelCase, letters and digits only. +func isLexiconErrorName(s string) bool { + if s == "" || s[0] < 'A' || s[0] > 'Z' { + return false + } + for _, r := range s { + if r >= 'a' && r <= 'z' || r >= 'A' && r <= 'Z' || r >= '0' && r <= '9' { + continue + } + return false + } + return true +} + // isMethodNotSupported reports whether err means "this host does not implement // that XRPC method", which is the only failure the backfill is allowed to // answer by falling back to a full getRepo. @@ -218,6 +394,14 @@ func (f *pdsLockedFetcher) GetBlocks(ctx context.Context, cids []cid.Cid) (map[c // without this the fallback would never fire for did:web streamplace repos. // Errors are unwrapped because reposync wraps everything with %w. // +// A named lexicon error disqualifies all of that, whatever the status code: a +// host that answers BlockNotFound or RepoNotFound plainly does implement the +// method, and is telling us something about this repo. That distinction is not +// academic -- streamplace's own getBlocks returns 404 BlockNotFound for a block +// it has garbage collected, and treating that as "unsupported" would answer a +// mid-walk race against a peer with a full-repo CAR download instead of a cheap +// re-walk. +// // A false positive costs one wasted getRepo attempt, whose own error then // propagates -- it can never turn a verification failure into a silent success, // because verification failures are not HTTP errors. @@ -225,6 +409,14 @@ func isMethodNotSupported(err error) bool { if err == nil { return false } + switch name := xrpcErrorName(err); name { + case "": + // No name to go on; fall through to the status code. + case "MethodNotImplemented", "XRPCNotSupported": + return true + default: + return false + } var xe *xrpc.Error if errors.As(err, &xe) { switch xe.StatusCode { @@ -232,9 +424,5 @@ func isMethodNotSupported(err error) bool { return true } } - var xrpcErr *xrpc.XRPCError - if errors.As(err, &xrpcErr) && xrpcErr.ErrStr == "MethodNotImplemented" { - return true - } return false } diff --git a/pkg/atproto/backfill_walk_test.go b/pkg/atproto/backfill_walk_test.go index 55a67d74c..9d3962f79 100644 --- a/pkg/atproto/backfill_walk_test.go +++ b/pkg/atproto/backfill_walk_test.go @@ -16,6 +16,7 @@ import ( "github.com/bluesky-social/indigo/util" "github.com/bluesky-social/indigo/xrpc" "github.com/ipfs/go-cid" + "github.com/multiformats/go-multihash" glex "github.com/streamplace/glex/runtime" "github.com/stretchr/testify/require" "stream.place/streamplace/pkg/appbsky" @@ -280,7 +281,7 @@ func TestIsMethodNotSupported(t *testing.T) { "doubly wrapped 404 with body", fmt.Errorf("walking: %w", fmt.Errorf("getBlocks: %w", &xrpc.Error{ StatusCode: http.StatusNotFound, - Wrapped: &xrpc.XRPCError{ErrStr: "NotFound", Message: "no such route"}, + Wrapped: &xrpc.XRPCError{ErrStr: "MethodNotImplemented", Message: "no such route"}, })), true, }, @@ -289,6 +290,37 @@ func TestIsMethodNotSupported(t *testing.T) { fmt.Errorf("getLatestCommit: %w", &xrpc.XRPCError{ErrStr: "MethodNotImplemented", Message: "nope"}), true, }, + { + // The whole point of looking at the error name: streamplace's own + // getBlocks answers 404 BlockNotFound for a block it no longer + // has. That host implements the method; falling back to a full + // getRepo download because of it would be a disaster. + "404 BlockNotFound", + fmt.Errorf("getBlocks: %w", &xrpc.Error{ + StatusCode: http.StatusNotFound, + Wrapped: &xrpc.XRPCError{ErrStr: "BlockNotFound", Message: "bafyreib2"}, + }), + false, + }, + { + // Same thing as it actually arrives from spxrpc, where echo's + // default error handler puts the name in "message" and leaves + // "error" empty. + "404 BlockNotFound with the name only in the message", + fmt.Errorf("getBlocks: %w", &xrpc.Error{ + StatusCode: http.StatusNotFound, + Wrapped: &xrpc.XRPCError{Message: "BlockNotFound"}, + }), + false, + }, + { + "404 RepoNotFound", + fmt.Errorf("getBlocks: %w", &xrpc.Error{ + StatusCode: http.StatusNotFound, + Wrapped: &xrpc.XRPCError{Message: "RepoNotFound"}, + }), + false, + }, { "400 InvalidRequest", fmt.Errorf("getBlocks: %w", &xrpc.Error{ @@ -305,6 +337,11 @@ func TestIsMethodNotSupported(t *testing.T) { }), false, }, + { + "429", + fmt.Errorf("getBlocks: %w", &xrpc.Error{StatusCode: http.StatusTooManyRequests}), + false, + }, {"block mismatch", fmt.Errorf("fetching: %w", reposync.ErrBlockMismatch), false}, {"missing block", fmt.Errorf("fetching: %w", reposync.ErrMissingBlock), false}, {"canceled", fmt.Errorf("walking: %w", context.Canceled), false}, @@ -316,6 +353,207 @@ func TestIsMethodNotSupported(t *testing.T) { } } +func TestIsStaleWalkError(t *testing.T) { + for _, tc := range []struct { + name string + err error + want bool + }{ + {"nil", nil, false}, + { + // Our own client-side check, which is what fires when a host + // answers with a short bag of blocks instead of an error. + "missing block, as the walker wraps it", + fmt.Errorf("failed to walk repo: %w", fmt.Errorf("fetching 3 MST nodes: %w", reposync.ErrMissingBlock)), + true, + }, + { + // The TypeScript PDS shape, seen from both a bsky.network + // mothership and a self-hosted instance. + "400 Could not find cids", + fmt.Errorf("getBlocks: %w", &xrpc.Error{ + StatusCode: http.StatusBadRequest, + Wrapped: &xrpc.XRPCError{ErrStr: "InvalidRequest", Message: "Could not find cids: bafyreib2"}, + }), + true, + }, + { + "404 BlockNotFound", + fmt.Errorf("getBlocks: %w", &xrpc.Error{ + StatusCode: http.StatusNotFound, + Wrapped: &xrpc.XRPCError{ErrStr: "BlockNotFound", Message: "bafyreib2"}, + }), + true, + }, + { + "404 BlockNotFound from spxrpc, name in the message", + fmt.Errorf("getBlocks: %w", &xrpc.Error{ + StatusCode: http.StatusNotFound, + Wrapped: &xrpc.XRPCError{Message: "BlockNotFound"}, + }), + true, + }, + { + "a different InvalidRequest", + fmt.Errorf("getBlocks: %w", &xrpc.Error{ + StatusCode: http.StatusBadRequest, + Wrapped: &xrpc.XRPCError{ErrStr: "InvalidRequest", Message: "cids/0 must be a cid string"}, + }), + false, + }, + { + "RepoNotFound", + fmt.Errorf("getLatestCommit: %w", &xrpc.Error{ + StatusCode: http.StatusBadRequest, + Wrapped: &xrpc.XRPCError{ErrStr: "RepoNotFound", Message: "could not find repo"}, + }), + false, + }, + {"throttled", fmt.Errorf("getBlocks: %w", &xrpc.Error{StatusCode: http.StatusTooManyRequests}), false}, + {"tampered block", fmt.Errorf("walking: %w", reposync.ErrBlockMismatch), false}, + {"malformed tree", fmt.Errorf("walking: %w", reposync.ErrInvalidNode), false}, + {"plain error", errors.New("connection refused"), false}, + } { + t.Run(tc.name, func(t *testing.T) { + require.Equal(t, tc.want, isStaleWalkError(tc.err)) + }) + } +} + +// TestWalkWithHeadRetry covers the live-repo race: the PDS garbage-collects the +// blocks of a commit we pinned, and the fix is to re-read the head rather than +// to fail. Staging that against a real PDS would mean making it GC mid-walk, so +// the head fetch and the walk are scripted here instead. +func TestWalkWithHeadRetry(t *testing.T) { + ctx := context.Background() + // A missing block, wrapped the way walkBackfill wraps it. + gone := fmt.Errorf("failed to walk repo for did:plc:x from PDS https://pds: %w", + fmt.Errorf("fetching 12 MST nodes: %w", reposync.ErrMissingBlock)) + + t.Run("restarts from the new head", func(t *testing.T) { + heads := []*reposync.Head{testHead(t, "3laaa"), testHead(t, "3lbbb")} + fetches := 0 + fetchHead := func(context.Context) (*reposync.Head, error) { + h := heads[min(fetches, len(heads)-1)] + fetches++ + return h, nil + } + var walked []cid.Cid + walk := func(_ context.Context, root cid.Cid) error { + walked = append(walked, root) + if len(walked) == 1 { + return gone + } + return nil + } + + head, err := walkWithHeadRetry(ctx, 3, time.Millisecond, fetchHead, walk) + require.NoError(t, err) + require.Equal(t, heads[1], head, "the completed walk was against the new head") + require.Equal(t, []cid.Cid{heads[0].Root, heads[1].Root}, walked) + require.Equal(t, 2, fetches) + }) + + t.Run("a head that did not move means the repo really is incomplete", func(t *testing.T) { + // Never paper over a repo we could not read: the alternative is + // recording a Version whose records we know we are missing. + head := testHead(t, "3laaa") + fetches := 0 + fetchHead := func(context.Context) (*reposync.Head, error) { + fetches++ + return head, nil + } + walks := 0 + walk := func(context.Context, cid.Cid) error { + walks++ + return gone + } + + _, err := walkWithHeadRetry(ctx, 3, time.Millisecond, fetchHead, walk) + require.Error(t, err) + require.ErrorIs(t, err, reposync.ErrMissingBlock) + require.Contains(t, err.Error(), "still the head") + require.Equal(t, 1, walks, "no point walking the same tree again") + require.Equal(t, 2, fetches) + }) + + t.Run("anything else fails immediately", func(t *testing.T) { + fetches := 0 + fetchHead := func(context.Context) (*reposync.Head, error) { + fetches++ + return testHead(t, "3laaa"), nil + } + walks := 0 + bad := fmt.Errorf("walking: %w", reposync.ErrBlockMismatch) + walk := func(context.Context, cid.Cid) error { + walks++ + return bad + } + + _, err := walkWithHeadRetry(ctx, 3, time.Millisecond, fetchHead, walk) + require.ErrorIs(t, err, reposync.ErrBlockMismatch) + require.Equal(t, 1, walks) + require.Equal(t, 1, fetches, "a verification failure is not worth a new head") + }) + + t.Run("a repo that keeps moving exhausts the budget", func(t *testing.T) { + fetches := 0 + fetchHead := func(context.Context) (*reposync.Head, error) { + fetches++ + return testHead(t, fmt.Sprintf("3l%03d", fetches)), nil + } + walks := 0 + walk := func(context.Context, cid.Cid) error { + walks++ + return gone + } + + _, err := walkWithHeadRetry(ctx, 3, time.Millisecond, fetchHead, walk) + require.Error(t, err) + require.ErrorIs(t, err, reposync.ErrMissingBlock) + require.Contains(t, err.Error(), "gave up after 3 walk attempts") + require.Equal(t, 3, walks) + }) + + t.Run("the first head fetch failing is just an error", func(t *testing.T) { + boom := errors.New("no such host") + walks := 0 + _, err := walkWithHeadRetry(ctx, + 3, time.Millisecond, + func(context.Context) (*reposync.Head, error) { return nil, boom }, + func(context.Context, cid.Cid) error { walks++; return nil }, + ) + require.ErrorIs(t, err, boom) + require.Zero(t, walks) + }) + + t.Run("cancellation during the backoff returns promptly", func(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + go func() { + time.Sleep(20 * time.Millisecond) + cancel() + }() + start := time.Now() + _, err := walkWithHeadRetry(ctx, + 3, 30*time.Second, + func(context.Context) (*reposync.Head, error) { return testHead(t, "3laaa"), nil }, + func(context.Context, cid.Cid) error { return gone }, + ) + require.Error(t, err) + require.Less(t, time.Since(start), 5*time.Second) + require.ErrorIs(t, err, context.Canceled) + require.ErrorIs(t, err, reposync.ErrMissingBlock, "the walk failure is kept too") + }) +} + +// testHead fabricates a head at a given rev; only Rev and Root are consulted. +func testHead(t *testing.T, rev string) *reposync.Head { + t.Helper() + root, err := cid.NewPrefixV1(cid.DagCBOR, multihash.SHA2_256).Sum([]byte("root:" + rev)) + require.NoError(t, err) + return &reposync.Head{Rev: rev, Root: root} +} + func TestBackfillRanges(t *testing.T) { ranges := backfillRanges() // One for place.stream., plus one per non-streamplace collection the diff --git a/pkg/reposync/doc.go b/pkg/reposync/doc.go index ca135f154..8dad20459 100644 --- a/pkg/reposync/doc.go +++ b/pkg/reposync/doc.go @@ -38,6 +38,22 @@ // - Deletions are not observable from a single walk; a caller detects them by // diffing two walks (see [Walker.CollectPrefix] and [DiffCollections]). // +// # Transient failures +// +// Walking a large repo is hundreds of sequential getBlocks calls, so a single +// rate limit or restarting host must not end it: [XRPCBlockFetcher] and +// [FetchVerifiedHead] retry 429s, 5xx and dropped connections with a jittered +// exponential backoff ([RetryPolicy]). Everything else -- 4xx, verification +// failures, a cancelled context -- fails immediately. +// +// One 4xx in particular is worth knowing about: a walk pins a root and then +// reads it over many round trips, while the host garbage-collects blocks that +// only superseded commits referenced. A repo that commits mid-walk can leave +// blocks unfetchable ([ErrMissingBlock], or a host-specific BlockNotFound / +// "Could not find cids"). Retrying cannot help; the caller has to re-read the +// head and walk the new tree, reusing its [CachedFetcher] so the second pass +// only pays for what changed. +// // # Resuming // // A [Frontier] is the complete state of an in-progress walk and is JSON diff --git a/pkg/reposync/fetcher.go b/pkg/reposync/fetcher.go index b10ec4f10..50b513237 100644 --- a/pkg/reposync/fetcher.go +++ b/pkg/reposync/fetcher.go @@ -82,6 +82,9 @@ type XRPCBlockFetcher struct { // ChunkSize caps how many CIDs go into a single getBlocks request. // Zero means [DefaultChunkSize]. ChunkSize int + // Retry bounds how hard each getBlocks call is retried after a transient + // failure. The zero value means the package defaults. + Retry RetryPolicy } var _ BlockFetcher = (*XRPCBlockFetcher)(nil) @@ -103,7 +106,16 @@ func (f *XRPCBlockFetcher) GetBlocks(ctx context.Context, cids []cid.Cid) (map[c for i, c := range batch { strs[i] = c.String() } - raw, err := indigoat.SyncGetBlocks(ctx, f.Client, strs, f.DID) + // Retried as a unit: a walk of a big repo makes hundreds of these calls + // in a row, so a single 429 from a busy PDS must not end it. Parsing + // happens outside the retry -- a CAR we cannot read is not transient. + var raw []byte + what := fmt.Sprintf("com.atproto.sync.getBlocks %s (%d cids)", f.DID, len(strs)) + err := f.Retry.do(ctx, what, func() error { + var err error + raw, err = indigoat.SyncGetBlocks(ctx, f.Client, strs, f.DID) + return err + }) if err != nil { return nil, fmt.Errorf("com.atproto.sync.getBlocks for %s (%d cids): %w", f.DID, len(strs), err) } diff --git a/pkg/reposync/head.go b/pkg/reposync/head.go index d0b360e93..647d0341a 100644 --- a/pkg/reposync/head.go +++ b/pkg/reposync/head.go @@ -34,13 +34,29 @@ type Head struct { // structure, DID, rev and signature against the account's atproto signing key // from dir. On success, nothing the host says about the repo below Head.Root can // be forged. -func FetchVerifiedHead(ctx context.Context, client *xrpc.Client, f BlockFetcher, dir identity.Directory, did string) (*Head, error) { +// +// At most one retry policy may be given; it applies to the getLatestCommit call +// (the block fetch carries its own). Omitting it uses the package defaults. +func FetchVerifiedHead(ctx context.Context, client *xrpc.Client, f BlockFetcher, dir identity.Directory, did string, retry ...RetryPolicy) (*Head, error) { + if len(retry) > 1 { + return nil, fmt.Errorf("at most one retry policy, got %d", len(retry)) + } + var policy RetryPolicy + if len(retry) == 1 { + policy = retry[0] + } + parsedDID, err := syntax.ParseDID(did) if err != nil { return nil, fmt.Errorf("invalid did %q: %w", did, err) } - latest, err := indigoat.SyncGetLatestCommit(ctx, client, did) + var latest *indigoat.SyncGetLatestCommit_Output + err = policy.do(ctx, "com.atproto.sync.getLatestCommit "+did, func() error { + var err error + latest, err = indigoat.SyncGetLatestCommit(ctx, client, did) + return err + }) if err != nil { return nil, fmt.Errorf("com.atproto.sync.getLatestCommit for %s: %w", did, err) } diff --git a/pkg/reposync/head_test.go b/pkg/reposync/head_test.go index f3c42c387..21228f58a 100644 --- a/pkg/reposync/head_test.go +++ b/pkg/reposync/head_test.go @@ -86,6 +86,37 @@ type fakeHost struct { tamper map[cid.Cid]bool // requests counts getBlocks calls. requests int + // latestRequests counts getLatestCommit calls. + latestRequests int + // blocksFailures is popped once per getBlocks call: while it is non-empty + // the request is answered with that failure instead of a CAR. This is how + // the retry tests script a flaky host. + blocksFailures []failure + // latestFailures does the same for getLatestCommit. + latestFailures []failure +} + +// failure is one scripted error response. +type failure struct { + status int + body string + header map[string]string +} + +// pop takes the next scripted failure off script, writes it, and reports +// whether it did anything. +func pop(script *[]failure, w http.ResponseWriter) bool { + if len(*script) == 0 { + return false + } + f := (*script)[0] + *script = (*script)[1:] + for k, v := range f.header { + w.Header().Set(k, v) + } + w.WriteHeader(f.status) + _, _ = w.Write([]byte(f.body)) + return true } func newFakeHost(sr *signedRepo) *fakeHost { @@ -102,11 +133,18 @@ func (h *fakeHost) start(t *testing.T) *xrpc.Client { t.Helper() mux := http.NewServeMux() mux.HandleFunc("/xrpc/com.atproto.sync.getLatestCommit", func(w http.ResponseWriter, r *http.Request) { + h.latestRequests++ + if pop(&h.latestFailures, w) { + return + } w.Header().Set("Content-Type", "application/json") _ = json.NewEncoder(w).Encode(map[string]string{"cid": h.head.String(), "rev": h.rev}) }) mux.HandleFunc("/xrpc/com.atproto.sync.getBlocks", func(w http.ResponseWriter, r *http.Request) { h.requests++ + if pop(&h.blocksFailures, w) { + return + } buf := new(bytes.Buffer) // Real getBlocks responses carry an empty roots list. if err := car.WriteHeader(&car.CarHeader{Roots: nil, Version: 1}, buf); err != nil { diff --git a/pkg/reposync/retry.go b/pkg/reposync/retry.go new file mode 100644 index 000000000..ce041e390 --- /dev/null +++ b/pkg/reposync/retry.go @@ -0,0 +1,186 @@ +package reposync + +import ( + "context" + "errors" + "fmt" + "io" + "math/rand" + "net" + "net/http" + "syscall" + "time" + + "github.com/bluesky-social/indigo/xrpc" + "stream.place/streamplace/pkg/log" +) + +// Retry defaults. A walk of a large repo is 500-1100 sequential getBlocks calls +// at [DefaultChunkSize], so a single 429 or a single restarting PDS must not be +// able to kill it. +const ( + // DefaultMaxAttempts is the total number of tries (not extra tries) a + // request gets before its error is returned. + DefaultMaxAttempts = 5 + // DefaultRetryBaseDelay is the wait after the first failure; it doubles + // from there. + DefaultRetryBaseDelay = time.Second + // DefaultRetryMaxDelay caps the wait between attempts, including waits + // derived from a server's ratelimit-reset header. Sleeping longer than this + // is worse than failing: backfills serialize their fetches per PDS, so a + // long sleep here stalls every other repo on that host, and a repo whose + // backfill fails is simply retried later. + DefaultRetryMaxDelay = 30 * time.Second +) + +// RetryPolicy bounds how hard a fetcher retries a transient XRPC failure. +// The zero value means the defaults above. +type RetryPolicy struct { + // MaxAttempts is the total number of tries. Zero means + // [DefaultMaxAttempts]; a value of 1 disables retrying. + MaxAttempts int + // BaseDelay is the wait after the first failure. Zero means + // [DefaultRetryBaseDelay]. + BaseDelay time.Duration + // MaxDelay caps every wait. Zero means [DefaultRetryMaxDelay]. + MaxDelay time.Duration +} + +func (p RetryPolicy) withDefaults() RetryPolicy { + if p.MaxAttempts <= 0 { + p.MaxAttempts = DefaultMaxAttempts + } + if p.BaseDelay <= 0 { + p.BaseDelay = DefaultRetryBaseDelay + } + if p.MaxDelay <= 0 { + p.MaxDelay = DefaultRetryMaxDelay + } + if p.MaxDelay < p.BaseDelay { + p.MaxDelay = p.BaseDelay + } + return p +} + +// delay is how long to wait after the attempt'th failure (1-based). +// +// Exponential from BaseDelay, capped at MaxDelay, then scaled by a random +// factor in [0.75, 1) so that a fleet of workers that hit the same rate limit +// does not march back in lockstep. Jitter is multiplicative rather than +// additive so the result never exceeds MaxDelay, including at the cap. +// +// If the server told us when its rate limit resets, and that is further out +// than the computed backoff, wait for the reset instead -- still clamped to +// MaxDelay, see the note there. +func (p RetryPolicy) delay(attempt int, err error) time.Duration { + p = p.withDefaults() + d := p.MaxDelay + if attempt >= 1 && attempt < 31 { + if shifted := p.BaseDelay << (attempt - 1); shifted > 0 && shifted < p.MaxDelay { + d = shifted + } + } + d = time.Duration(float64(d) * (0.75 + 0.25*rand.Float64())) //nolint:gosec // jitter, not crypto + if reset := ratelimitReset(err); !reset.IsZero() { + // A small pad: the reset second has to have actually elapsed. + if wait := time.Until(reset) + 250*time.Millisecond; wait > d { + d = min(wait, p.MaxDelay) + } + } + return d +} + +// do runs fn until it succeeds, fails with something not worth retrying, or +// runs out of attempts. what names the call for logging only; the error +// returned is fn's, unwrapped when it was not retryable and wrapped with the +// attempt count when the budget ran out. +func (p RetryPolicy) do(ctx context.Context, what string, fn func() error) error { + p = p.withDefaults() + for attempt := 1; ; attempt++ { + err := fn() + if err == nil { + return nil + } + if !isRetryable(err) { + return err + } + if attempt >= p.MaxAttempts { + return fmt.Errorf("giving up after %d attempts: %w", attempt, err) + } + d := p.delay(attempt, err) + log.Warn(ctx, "retrying transient xrpc failure", "call", what, "attempt", attempt, "wait", d, "err", err) + if serr := sleepCtx(ctx, d); serr != nil { + return fmt.Errorf("aborted after %d attempts: %w", attempt, errors.Join(err, serr)) + } + } +} + +// sleepCtx waits for d, or returns the context's error as soon as it is done. +func sleepCtx(ctx context.Context, d time.Duration) error { + if d <= 0 { + return ctx.Err() + } + t := time.NewTimer(d) + defer t.Stop() + select { + case <-ctx.Done(): + return ctx.Err() + case <-t.C: + return nil + } +} + +// isRetryable reports whether err is the kind of failure that is likely to go +// away on its own: the host throttled us, the host is briefly broken, or the +// connection died under us. +// +// Everything else -- 4xx (including "Could not find cids", which means the repo +// moved and needs a new head, not a retry), block verification failures, a +// cancelled context -- fails fast. +func isRetryable(err error) bool { + if err == nil { + return false + } + // A cancelled context can surface as a *url.Error, which would otherwise + // look like a transport blip; check it first. + if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) { + return false + } + var xe *xrpc.Error + if errors.As(err, &xe) { + switch { + case xe.StatusCode == http.StatusTooManyRequests: + return true + case xe.StatusCode == http.StatusNotImplemented: + // 5xx numerically, but it is a permanent answer: this host does + // not implement the method, and the caller wants to hear that + // immediately so it can fall back. + return false + case xe.StatusCode >= 500 && xe.StatusCode <= 599: + return true + } + return false + } + // No HTTP response at all. + var nerr net.Error + if errors.As(err, &nerr) && nerr.Timeout() { + return true + } + return errors.Is(err, syscall.ECONNRESET) || + errors.Is(err, syscall.ECONNREFUSED) || + errors.Is(err, syscall.EPIPE) || + errors.Is(err, io.ErrUnexpectedEOF) || + errors.Is(err, io.EOF) +} + +// ratelimitReset pulls the reset time out of an XRPC error, if the host sent +// ratelimit-* headers. indigo parses those into xrpc.Error.Ratelimit; note it +// only does so when a ratelimit-limit header is present, and it does not look +// at Retry-After at all, so this is often zero even for a 429. +func ratelimitReset(err error) time.Time { + var xe *xrpc.Error + if !errors.As(err, &xe) || xe.Ratelimit == nil { + return time.Time{} + } + return xe.Ratelimit.Reset +} diff --git a/pkg/reposync/retry_test.go b/pkg/reposync/retry_test.go new file mode 100644 index 000000000..48d3f593e --- /dev/null +++ b/pkg/reposync/retry_test.go @@ -0,0 +1,331 @@ +package reposync + +import ( + "context" + "errors" + "fmt" + "io" + "net/http" + "strconv" + "syscall" + "testing" + "time" + + "github.com/bluesky-social/indigo/xrpc" + "github.com/ipfs/go-cid" + "github.com/stretchr/testify/require" +) + +// fastRetry is the policy the network tests use: same shape as production, +// milliseconds instead of seconds. +func fastRetry() RetryPolicy { + return RetryPolicy{MaxAttempts: 5, BaseDelay: time.Millisecond, MaxDelay: 20 * time.Millisecond} +} + +// throttled builds the failure a rate-limiting host sends, optionally with the +// ratelimit-* headers indigo knows how to parse. +func throttled(reset time.Time) failure { + f := failure{ + status: http.StatusTooManyRequests, + body: `{"error":"RateLimitExceeded","message":"Rate Limit Exceeded"}`, + header: map[string]string{"Content-Type": "application/json"}, + } + if !reset.IsZero() { + f.header["ratelimit-limit"] = "3000" + f.header["ratelimit-remaining"] = "0" + f.header["ratelimit-policy"] = "3000;w=300" + f.header["ratelimit-reset"] = strconv.FormatInt(reset.Unix(), 10) + } + return f +} + +// htmlThrottled is the shape that actually broke a production walk: a 429 whose +// body is an HTML error page, so indigo cannot decode an XRPCError out of it and +// only the status code survives. +var htmlThrottled = failure{ + status: http.StatusTooManyRequests, + body: "429 Too Many Requestsgo away", + header: map[string]string{"Content-Type": "text/html"}, +} + +func xrpcErr(status int, errStr, msg string) error { + return fmt.Errorf("getBlocks: %w", &xrpc.Error{ + StatusCode: status, + Wrapped: &xrpc.XRPCError{ErrStr: errStr, Message: msg}, + }) +} + +func TestIsRetryable(t *testing.T) { + for _, tc := range []struct { + name string + err error + want bool + }{ + {"nil", nil, false}, + {"429", xrpcErr(http.StatusTooManyRequests, "RateLimitExceeded", "slow down"), true}, + {"429 with an undecodable body", fmt.Errorf("getBlocks: %w", &xrpc.Error{ + StatusCode: http.StatusTooManyRequests, + Wrapped: errors.New("failed to decode xrpc error message: invalid character '<'"), + }), true}, + {"500", xrpcErr(http.StatusInternalServerError, "InternalServerError", "oops"), true}, + {"502", fmt.Errorf("x: %w", &xrpc.Error{StatusCode: http.StatusBadGateway}), true}, + {"503", fmt.Errorf("x: %w", &xrpc.Error{StatusCode: http.StatusServiceUnavailable}), true}, + {"504", fmt.Errorf("x: %w", &xrpc.Error{StatusCode: http.StatusGatewayTimeout}), true}, + // Numerically 5xx, but a permanent answer: the backfill wants to hear + // it at once so it can fall back to getRepo. + {"501", fmt.Errorf("x: %w", &xrpc.Error{StatusCode: http.StatusNotImplemented}), false}, + {"400 could not find cids", xrpcErr(http.StatusBadRequest, "InvalidRequest", "Could not find cids: bafy"), false}, + {"401", fmt.Errorf("x: %w", &xrpc.Error{StatusCode: http.StatusUnauthorized}), false}, + {"404", fmt.Errorf("x: %w", &xrpc.Error{StatusCode: http.StatusNotFound}), false}, + {"connection reset", fmt.Errorf("request failed: %w", syscall.ECONNRESET), true}, + {"connection refused", fmt.Errorf("request failed: %w", syscall.ECONNREFUSED), true}, + {"truncated body", fmt.Errorf("reading response body: %w", io.ErrUnexpectedEOF), true}, + {"timeout", fmt.Errorf("request failed: %w", timeoutError{}), true}, + {"canceled", fmt.Errorf("request failed: %w", context.Canceled), false}, + {"deadline exceeded", fmt.Errorf("request failed: %w", context.DeadlineExceeded), false}, + {"missing block", fmt.Errorf("x: %w", ErrMissingBlock), false}, + {"block mismatch", fmt.Errorf("x: %w", ErrBlockMismatch), false}, + {"plain error", errors.New("nope"), false}, + } { + t.Run(tc.name, func(t *testing.T) { + require.Equal(t, tc.want, isRetryable(tc.err)) + }) + } +} + +type timeoutError struct{} + +func (timeoutError) Error() string { return "i/o timeout" } +func (timeoutError) Timeout() bool { return true } +func (timeoutError) Temporary() bool { return true } + +func TestRetryDelay(t *testing.T) { + p := RetryPolicy{BaseDelay: time.Second, MaxDelay: 30 * time.Second} + plain := errors.New("boom") + + // Exponential, jittered down by at most 25%. + for _, tc := range []struct{ attempt, wantSec int }{{1, 1}, {2, 2}, {3, 4}, {4, 8}, {5, 16}} { + full := time.Duration(tc.wantSec) * time.Second + for i := 0; i < 50; i++ { + d := p.delay(tc.attempt, plain) + require.GreaterOrEqual(t, d, time.Duration(float64(full)*0.75), "attempt %d", tc.attempt) + require.LessOrEqual(t, d, full, "attempt %d", tc.attempt) + } + } + + // Capped, and still jittered at the cap so a fleet does not resynchronize. + var sawJitter bool + for i := 0; i < 50; i++ { + d := p.delay(10, plain) + require.GreaterOrEqual(t, d, 22500*time.Millisecond) + require.LessOrEqual(t, d, 30*time.Second) + if d < 29*time.Second { + sawJitter = true + } + } + require.True(t, sawJitter) + + // The zero policy is the documented defaults. + require.LessOrEqual(t, RetryPolicy{}.delay(1, plain), DefaultRetryBaseDelay) + require.GreaterOrEqual(t, RetryPolicy{}.delay(1, plain), DefaultRetryBaseDelay*3/4) + + t.Run("ratelimit reset is honored", func(t *testing.T) { + // Further out than the backoff for attempt 1 (~1s): wait for the reset. + err := ratelimited(time.Now().Add(3 * time.Second)) + d := p.delay(1, err) + require.Greater(t, d, 2500*time.Millisecond) + require.LessOrEqual(t, d, 3500*time.Millisecond) + }) + + t.Run("ratelimit reset is clamped to MaxDelay", func(t *testing.T) { + // bsky rate limit windows are minutes long; we would rather make one + // more doomed attempt than hold a per-PDS lock that long. + err := ratelimited(time.Now().Add(10 * time.Minute)) + require.Equal(t, 30*time.Second, p.delay(1, err)) + }) + + t.Run("a reset in the past does not shorten the backoff", func(t *testing.T) { + err := ratelimited(time.Now().Add(-time.Minute)) + d := p.delay(3, err) + require.GreaterOrEqual(t, d, 3*time.Second) + require.LessOrEqual(t, d, 4*time.Second) + }) +} + +func ratelimited(reset time.Time) error { + return fmt.Errorf("getBlocks: %w", &xrpc.Error{ + StatusCode: http.StatusTooManyRequests, + Wrapped: &xrpc.XRPCError{ErrStr: "RateLimitExceeded"}, + Ratelimit: &xrpc.RatelimitInfo{Limit: 3000, Reset: reset}, + }) +} + +// TestXRPCBlockFetcherRetries drives the retry loop through the real getBlocks +// path against an HTTP host that fails on a script. +func TestXRPCBlockFetcherRetries(t *testing.T) { + ctx := context.Background() + sr := buildSignedRepo(t, testDID, exactnessPaths()) + + newFetcher := func(t *testing.T, script ...failure) (*XRPCBlockFetcher, *fakeHost) { + host := newFakeHost(sr) + host.blocksFailures = script + client := host.start(t) + return &XRPCBlockFetcher{Client: client, DID: testDID, Retry: fastRetry()}, host + } + + t.Run("429 then success", func(t *testing.T) { + f, host := newFetcher(t, throttled(time.Time{}), throttled(time.Time{})) + blocks, err := f.GetBlocks(ctx, []cid.Cid{sr.root}) + require.NoError(t, err) + require.Contains(t, blocks, sr.root) + require.Equal(t, 3, host.requests, "two throttles then the real answer") + }) + + t.Run("429 with an HTML body then success", func(t *testing.T) { + // The production shape: indigo cannot decode the body, so the error is + // "failed to decode xrpc error message: invalid character '<'" and only + // the status code is left to classify on. + f, host := newFetcher(t, htmlThrottled) + _, err := f.GetBlocks(ctx, []cid.Cid{sr.root}) + require.NoError(t, err) + require.Equal(t, 2, host.requests) + }) + + t.Run("503 then success", func(t *testing.T) { + f, host := newFetcher(t, failure{status: http.StatusServiceUnavailable, body: "upstream restarting"}) + _, err := f.GetBlocks(ctx, []cid.Cid{sr.root}) + require.NoError(t, err) + require.Equal(t, 2, host.requests) + }) + + t.Run("a ratelimit reset in the far future is capped, not slept on", func(t *testing.T) { + f, host := newFetcher(t, throttled(time.Now().Add(5*time.Minute))) + start := time.Now() + _, err := f.GetBlocks(ctx, []cid.Cid{sr.root}) + require.NoError(t, err) + require.Equal(t, 2, host.requests) + require.Less(t, time.Since(start), time.Second, "MaxDelay must bound the ratelimit wait") + }) + + t.Run("400 fails fast", func(t *testing.T) { + // What a TS PDS says when the walk raced a live repo and the blocks of + // the pinned commit have been garbage collected. Retrying the same CIDs + // can never help. + f, host := newFetcher(t, failure{ + status: http.StatusBadRequest, + body: `{"error":"InvalidRequest","message":"Could not find cids: bafyreib2"}`, + header: map[string]string{"Content-Type": "application/json"}, + }) + _, err := f.GetBlocks(ctx, []cid.Cid{sr.root}) + require.Error(t, err) + require.Equal(t, 1, host.requests, "no retries") + var xe *xrpc.Error + require.ErrorAs(t, err, &xe) + require.Equal(t, http.StatusBadRequest, xe.StatusCode) + require.NotContains(t, err.Error(), "giving up", "a fail-fast error is passed through unchanged") + }) + + t.Run("retries exhausted", func(t *testing.T) { + f, host := newFetcher(t, throttled(time.Time{}), throttled(time.Time{}), throttled(time.Time{}), + throttled(time.Time{}), throttled(time.Time{}), throttled(time.Time{})) + _, err := f.GetBlocks(ctx, []cid.Cid{sr.root}) + require.Error(t, err) + require.Equal(t, 5, host.requests, "MaxAttempts is a total, not an extra") + require.Contains(t, err.Error(), "giving up after 5 attempts") + var xe *xrpc.Error + require.ErrorAs(t, err, &xe, "the last error is still inspectable") + require.Equal(t, http.StatusTooManyRequests, xe.StatusCode) + }) + + t.Run("context cancelled mid backoff", func(t *testing.T) { + host := newFakeHost(sr) + host.blocksFailures = []failure{throttled(time.Time{})} + client := host.start(t) + // A backoff long enough that returning promptly can only mean the + // sleep was context aware. + f := &XRPCBlockFetcher{Client: client, DID: testDID, + Retry: RetryPolicy{MaxAttempts: 5, BaseDelay: 30 * time.Second, MaxDelay: time.Minute}} + + ctx, cancel := context.WithCancel(context.Background()) + go func() { + time.Sleep(20 * time.Millisecond) + cancel() + }() + start := time.Now() + _, err := f.GetBlocks(ctx, []cid.Cid{sr.root}) + require.Error(t, err) + require.Less(t, time.Since(start), 5*time.Second) + require.ErrorIs(t, err, context.Canceled) + var xe *xrpc.Error + require.ErrorAs(t, err, &xe, "the failure that triggered the backoff is kept too") + require.Equal(t, http.StatusTooManyRequests, xe.StatusCode) + require.Equal(t, 1, host.requests) + }) + + t.Run("retrying does not break chunking", func(t *testing.T) { + // A retry inside one chunk must not disturb the chunk loop: five CIDs + // at ChunkSize 2 is three chunks, and the failure only costs one extra + // call. + host := newFakeHost(sr) + host.blocksFailures = []failure{throttled(time.Time{})} + client := host.start(t) + f := &XRPCBlockFetcher{Client: client, DID: testDID, ChunkSize: 2, Retry: fastRetry()} + want := []cid.Cid{sr.root, sr.commitCID} + for c := range sr.blocks { + if len(want) >= 5 { + break + } + if c != sr.root && c != sr.commitCID { + want = append(want, c) + } + } + require.Len(t, want, 5) + blocks, err := f.GetBlocks(ctx, want) + require.NoError(t, err) + require.Len(t, blocks, len(want)) + require.Equal(t, 4, host.requests, "three chunks plus the one retry") + }) +} + +// TestFetchVerifiedHeadRetries: the head fetch is one getLatestCommit call, and +// it is the first thing every backfill does, so it gets the same treatment. +func TestFetchVerifiedHeadRetries(t *testing.T) { + ctx := context.Background() + sr := buildSignedRepo(t, testDID, exactnessPaths()) + pub, err := sr.priv.PublicKey() + require.NoError(t, err) + + t.Run("throttled then success", func(t *testing.T) { + host := newFakeHost(sr) + host.latestFailures = []failure{throttled(time.Time{}), htmlThrottled, + {status: http.StatusServiceUnavailable, body: "restarting"}} + client := host.start(t) + f := &XRPCBlockFetcher{Client: client, DID: testDID, Retry: fastRetry()} + head, err := FetchVerifiedHead(ctx, client, f, sr.directory(t, testDID, pub), testDID, fastRetry()) + require.NoError(t, err) + require.Equal(t, sr.commitCID, head.CID) + require.Equal(t, 4, host.latestRequests) + }) + + t.Run("RepoNotFound fails fast", func(t *testing.T) { + host := newFakeHost(sr) + host.latestFailures = []failure{{ + status: http.StatusBadRequest, + body: `{"error":"RepoNotFound","message":"could not find repo"}`, + header: map[string]string{"Content-Type": "application/json"}, + }} + client := host.start(t) + f := &XRPCBlockFetcher{Client: client, DID: testDID, Retry: fastRetry()} + _, err := FetchVerifiedHead(ctx, client, f, sr.directory(t, testDID, pub), testDID, fastRetry()) + require.Error(t, err) + require.Equal(t, 1, host.latestRequests) + }) + + t.Run("at most one policy", func(t *testing.T) { + host := newFakeHost(sr) + client := host.start(t) + f := &XRPCBlockFetcher{Client: client, DID: testDID} + _, err := FetchVerifiedHead(ctx, client, f, sr.directory(t, testDID, pub), testDID, fastRetry(), fastRetry()) + require.Error(t, err) + }) +} -- 2.51.2