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: "