diff --git a/pkg/atproto/atproto.go b/pkg/atproto/atproto.go index 4d203e70e..d82c3a6bd 100644 --- a/pkg/atproto/atproto.go +++ b/pkg/atproto/atproto.go @@ -104,7 +104,7 @@ func (atsync *ATProtoSynchronizer) SyncBlueskyRepo(ctx context.Context, handle s log.Log(ctx, "resolved bluesky identity", "did", ident.DID, "handle", ident.Handle, "pds", ident.PDSEndpoint()) xrpcc := xrpc.Client{ Host: ident.PDSEndpoint(), - Client: &aqhttp.Client, + Client: SyncHTTPClient, } if xrpcc.Host == "" { return nil, fmt.Errorf("no PDS endpoint found for Bluesky identity %s", handle) @@ -188,7 +188,7 @@ func (atsync *ATProtoSynchronizer) DeepenRepo(ctx context.Context, did string) ( if err != nil { return false, fmt.Errorf("failed to resolve %s: %w", did, err) } - xrpcc := xrpc.Client{Host: ident.PDSEndpoint(), Client: &aqhttp.Client} + xrpcc := xrpc.Client{Host: ident.PDSEndpoint(), Client: SyncHTTPClient} if xrpcc.Host == "" { return false, fmt.Errorf("no PDS endpoint found for %s", did) } diff --git a/pkg/atproto/backfill_walk.go b/pkg/atproto/backfill_walk.go index cc17e1348..c4615e8e1 100644 --- a/pkg/atproto/backfill_walk.go +++ b/pkg/atproto/backfill_walk.go @@ -15,6 +15,7 @@ import ( "github.com/bluesky-social/indigo/repo" "github.com/bluesky-social/indigo/xrpc" "github.com/ipfs/go-cid" + "stream.place/streamplace/pkg/aqhttp" "stream.place/streamplace/pkg/constants" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/model" @@ -29,17 +30,32 @@ const placeStreamPrefix = "place.stream." // windowedCollections are the collections a backfill reads by time window // instead of all at once. // -// Their rkeys are TIDs, so their keys sort chronologically and "the last day of -// chat" is a key range (see [reposync.TIDForTime]). They are also the two -// collections whose volume decides how long a first sync takes: an account with -// years of Bluesky posts and streamplace chat has tens of thousands of records -// there and a few dozen everywhere else. +// The policy, because getting this list wrong is a correctness bug rather than a +// performance one: window a collection if and only if it is a high-volume +// append-only log whose rkeys are TIDs. Those sort chronologically, so "the last +// day of chat" is a key range (see [reposync.TIDForTime]), and they are the +// whole cost of a first sync -- an active account has tens of thousands of +// posts, follows and chat messages, and a few dozen records everywhere else. A +// production sweep measured an average "shallow" sync fetching 756 records +// because follows were being read in full; one account with 3516 follows took +// four minutes. +// +// Everything else must stay full, in particular every current-state registry: +// signing keys, chat gates, settings, delegations, the media catalog, +// app.bsky.graph.block (moderation). Windowing those would be wrong, not slow. +// A windowed collection's records below the floor are invisible until the +// deepening ladder reaches the genesis window, which is hours later -- and +// "hours without your block list" or "hours without your signing key" is not a +// tradeoff, it is a broken node. // // Windowing changes when a record is indexed, never whether: the ladder of // windows in [nextBackfillWindow] bottoms out at the start of the collection, // and until it does the repo row says so. A record whose rkey is not a TID is // not skipped either -- the windows are key ranges, so it simply arrives with -// whichever window its rkey sorts into. +// whichever window its rkey sorts into. That is the caveat on the TID +// requirement: a literal rkey that sorts below the floor ("!oldest", say) waits +// for the genesis window like any old record would, which is fine for a log and +// would not be fine for a registry keyed "self". // // Every entry must be a collection [backfillRanges] would otherwise walk whole // -- either under place.stream. or in [CollectionFilter]. TestBackfillRanges @@ -47,6 +63,7 @@ const placeStreamPrefix = "place.stream." var windowedCollections = []string{ constants.PLACE_STREAM_CHAT_MESSAGE, constants.APP_BSKY_FEED_POST, + constants.APP_BSKY_GRAPH_FOLLOW, } // backfillRanges is the set of MST key ranges a backfill walks: everything @@ -264,18 +281,24 @@ func (atsync *ATProtoSynchronizer) walkBackfill(ctx context.Context, ident *iden dir = CustomDirectory(atsync.CLI.PLCURL) } + // Every retry in this walk consults what the host has been telling us about + // backing off; see [pdsBackoffHints]. It only works if the calls go through + // a client with the hint-capturing transport installed, which is what + // [SyncHTTPClient] is for -- callers build xrpcc with it. + retry := reposync.RetryPolicy{Hints: pdsBackoffHints} + fetcher := &reposync.CachedFetcher{ // Bounded lifetime: one cache per backfill, so the head fetch and the // walk share blocks without holding a repo in memory afterwards. Cache: reposync.NewMemoryBlockCache(), Inner: &pdsLockedFetcher{ lock: pdsLocks.GetLock(ident.PDSEndpoint()), - inner: &reposync.XRPCBlockFetcher{Client: xrpcc, DID: did}, + inner: &reposync.XRPCBlockFetcher{Client: xrpcc, DID: did, Retry: retry}, }, } fetchHead := func(ctx context.Context) (*reposync.Head, error) { - head, err := reposync.FetchVerifiedHead(ctx, xrpcc, fetcher, dir, did) + head, err := reposync.FetchVerifiedHead(ctx, xrpcc, fetcher, dir, did, retry) if err != nil { return nil, fmt.Errorf("failed to fetch verified head for %s from PDS %s: %w", did, xrpcc.Host, err) } @@ -457,6 +480,39 @@ func (atsync *ATProtoSynchronizer) legacyBackfill(ctx context.Context, ident *id return sc.Rev, nil } +// pdsBackoffHints is this process's memory of what PDS hosts have said about +// backing off, shared by every repo sync so that one repo's 429 slows down the +// next repo on that host too. +// +// It is fed by the transport on [SyncHTTPClient] and read by the retry policies +// in [walkBackfill]. It has to be a package-level singleton for the same reason +// pdsLocks is: the unit a rate limit applies to is the host, and the sweep works +// on thousands of repos across it. +var pdsBackoffHints = reposync.NewBackoffHints() + +// SyncHTTPClient is the HTTP client every repo-sync XRPC call goes through. +// +// It is [aqhttp.Client] -- same SSRF-checking transport, same timeout, same +// redirect policy -- with one wrapper installed: the round tripper that notices +// throttled responses and records their Retry-After / ratelimit-reset headers. +// indigo's xrpc client discards response headers, so watching the transport is +// the only place those numbers can be read at all, and without them a retry is +// guessing at a wait the host already told us. +// +// The wrapper is scoped to sync traffic rather than installed on aqhttp.Client +// globally: it is only useful where something consults the registry, and the +// node makes plenty of unrelated HTTP requests that would otherwise pay for a +// map lookup and a header parse. +var SyncHTTPClient = newSyncHTTPClient() + +func newSyncHTTPClient() *http.Client { + // By value: aqhttp.Client is a struct of settings (no mutex), and copying it + // keeps the connection-pooling transport shared with the rest of the node. + c := aqhttp.Client + c.Transport = pdsBackoffHints.Transport(c.Transport) + return &c +} + // pdsLockedFetcher serializes block fetches per PDS, the way the legacy path // serializes its one big getRepo download. // @@ -466,8 +522,9 @@ func (atsync *ATProtoSynchronizer) legacyBackfill(ctx context.Context, ident *id // // 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. +// polite response -- all the more so now that the backoff is usually the number +// the host itself asked for ([pdsBackoffHints]) rather than a guess. +// reposync.DefaultRetryMaxDelay is what keeps that pause bounded. type pdsLockedFetcher struct { lock *sync.Mutex inner reposync.BlockFetcher diff --git a/pkg/atproto/backfill_walk_test.go b/pkg/atproto/backfill_walk_test.go index 142f7643a..6e0b23b42 100644 --- a/pkg/atproto/backfill_walk_test.go +++ b/pkg/atproto/backfill_walk_test.go @@ -602,6 +602,36 @@ func TestBackfillRanges(t *testing.T) { } } +// TestBackfillSyncClientBackoffHints is a wiring test, and worth having as one: +// the hint registry only does anything if the client the sync path builds its +// xrpc.Client on is the one watching responses. That is easy to undo by writing +// aqhttp.Client at a new call site, and the symptom -- retries guessing at waits +// a host was announcing -- is invisible from anywhere but a busy production +// sweep. +func TestBackfillSyncClientBackoffHints(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Retry-After", "42") + w.WriteHeader(http.StatusTooManyRequests) + })) + t.Cleanup(srv.Close) + + req, err := http.NewRequestWithContext(context.Background(), http.MethodGet, srv.URL, nil) + require.NoError(t, err) + resp, err := SyncHTTPClient.Do(req) + require.NoError(t, err) + require.NoError(t, resp.Body.Close()) + + hint, ok := pdsBackoffHints.Get(srv.URL) + require.True(t, ok, "the sync client must record what a throttling host asked for") + require.Equal(t, "retry-after", hint.Source) + require.WithinDuration(t, time.Now().Add(42*time.Second), hint.Until, 5*time.Second) + + // And it is still the shared client underneath: same transport, so the same + // connection pool and the same SSRF checks. + require.Equal(t, aqhttp.Client.Timeout, SyncHTTPClient.Timeout) + require.NotNil(t, SyncHTTPClient.Transport) +} + // TestBackfillRangesWindowed: with a floor, the high-volume collections are cut // down to their recent history and everything else stays whole -- including the // records that sort on either side of the hole cut out of place.stream. @@ -618,6 +648,8 @@ func TestBackfillRangesWindowed(t *testing.T) { require.False(t, inRange("place.stream.chat.message/"+older)) require.True(t, inRange("app.bsky.feed.post/"+newer)) require.False(t, inRange("app.bsky.feed.post/"+older)) + require.True(t, inRange("app.bsky.graph.follow/"+newer)) + require.False(t, inRange("app.bsky.graph.follow/"+older)) // A windowed collection is still walked exhaustively, just not all at once: // a non-TID rkey lands in whichever window its bytes fall into ("self" // sorts above every TID minted this century, so it comes in the first one) @@ -630,7 +662,10 @@ func TestBackfillRangesWindowed(t *testing.T) { require.True(t, inRange("place.stream.chat.profile/self")) // sorts after require.True(t, inRange("place.stream.live.livestream/3l")) // sorts after require.True(t, inRange("app.bsky.actor.profile/self")) - require.True(t, inRange("app.bsky.graph.follow/3l")) + // Current-state registries stay whole however big they get: a moderation + // list that arrives hours late is worse than a slow sync. + require.True(t, inRange("app.bsky.graph.block/"+older)) + require.True(t, inRange("place.stream.key/"+older)) // And the ranges are still bounded where they were. require.False(t, inRange("app.bsky.feed.postgate/3l")) require.False(t, inRange("app.bsky.actor.status/3l")) @@ -642,10 +677,18 @@ func TestBackfillRangesWindowed(t *testing.T) { require.NotNil(t, r.Hi, "no unbounded range should come out of here: %s", r) require.Less(t, string(r.Lo), string(r.Hi), "inverted range %s", r) } - // Cutting two collections down to a window turns two whole ranges into a - // bounded piece each, plus the pieces left on either side of the hole in - // place.stream. - require.Len(t, ranges, len(backfillRanges(""))+len(windowedCollections)) + // Windowing a collection that has a range of its own (app.bsky.*) replaces + // that range with a bounded one and changes nothing about the count. + // Windowing one under place.stream. punches a hole in the middle of that + // prefix range instead, leaving a piece on either side of it plus the window + // itself: two more ranges each. + extra := 0 + for _, nsid := range windowedCollections { + if strings.HasPrefix(nsid, placeStreamPrefix) { + extra += 2 + } + } + require.Len(t, ranges, len(backfillRanges(""))+extra) } // TestWindowRanges: a deepening step reads the windowed collections and nothing diff --git a/pkg/atproto/sweep.go b/pkg/atproto/sweep.go index b690d9051..bad2c25c7 100644 --- a/pkg/atproto/sweep.go +++ b/pkg/atproto/sweep.go @@ -5,19 +5,16 @@ import ( "fmt" "sort" "sync" + "sync/atomic" "time" "golang.org/x/sync/errgroup" + "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/reposync" ) const ( - // sweepConcurrency bounds how many repos a sweep works on at once. A boot - // used to start one goroutine per known repo, which meant a fresh node - // opened with a thundering herd at every PDS it had ever heard of. - sweepConcurrency = 6 - // sweepStatusInterval is how often a running sweep says where it is. There // is exactly one such line per interval, and none at all when no sweep is // running. @@ -38,6 +35,183 @@ const ( // pure belt and braces. var maxDeepenRounds = len(backfillSpans) + 3 +// sweepItem is one repo for a sweep to work on, tagged with the lane it belongs +// to. +type sweepItem struct { + DID string + // Lane is what work is grouped by: the repo's PDS host. Everything a sweep + // spends its time on is a remote server, so the host is the only shape of + // the work that matters. + Lane string +} + +// sweepLane is the lane a repo row belongs in: its PDS host, or -- for a repo +// whose host is not known even after [ATProtoSynchronizer.resolveLanes] tried to +// find out -- a lane of its own. +// +// A lane to itself, rather than a shared catch-all: an unplaceable repo is +// normally a resolution failure, so its sync is about to fail too, and queueing +// those behind one another would make an outage at one identity service look +// like a stalled sweep. Two of them on one host is no worse than the flat worker +// pool this replaced -- the per-host pdsLock still keeps their fetches from +// interleaving -- and once either finishes, its row names a PDS, so the next +// sweep lanes it properly. +func sweepLane(did, pds string) string { + if host := reposync.HostKey(pds); host != "" { + return host + } + return "did:" + did +} + +// identityResolveConcurrency is how many identities [ATProtoSynchronizer.resolveLanes] +// looks up at once. +// +// Deliberately smaller than the sweep's own concurrency: these are lookups +// against a handful of shared identity services (plc.directory, DNS) rather than +// against thousands of PDSes, and the whole point is to move work that the +// backfill would have done anyway, not to arrive at plc.directory with a +// thundering herd. +const identityResolveConcurrency = 8 + +// resolveLanes gives a lane to every item that has not got one, by resolving the +// repo's identity to find its PDS. +// +// This exists because the sweep's DID list and the PDS column come from +// different databases. The DIDs are the state database's set of "repos this node +// indexes"; the host is a column in the index. A node with both has a host for +// every repo (both the placeholder written when a backfill starts and the row +// that replaces it record the PDS). A node with a fresh index and an inherited +// state database -- which is exactly what `streamplace sync` is for, warming a +// new index revision before it takes traffic -- has no rows at all, and would +// put every repo in a lane of its own: sharding by host would do nothing on +// precisely the sweep it was built for. +// +// The lookup is moved rather than added. It goes through the same cached +// directory [ATProtoSynchronizer.SyncBlueskyRepo] resolves with, so the backfill +// a few seconds later reads this answer out of the cache instead of asking +// again. +// +// Failures are not fatal and are not even logged loudly: the repo keeps a lane +// of its own and its sync fails on its own terms, one repo at a time, the way it +// did before. +func (atsync *ATProtoSynchronizer) resolveLanes(ctx context.Context, items []sweepItem) { + var todo []int + for i := range items { + if items[i].Lane == "" { + todo = append(todo, i) + } + } + if len(todo) == 0 { + return + } + log.Log(ctx, "resolving PDS hosts to shard the sweep", "repos", len(todo)) + + start := time.Now() + var resolved atomic.Int64 + g, gctx := errgroup.WithContext(ctx) + g.SetLimit(identityResolveConcurrency) + for _, i := range todo { + g.Go(func() error { + if err := gctx.Err(); err != nil { + return err + } + ident, err := atsync.resolveIdent(gctx, items[i].DID, true) + if err != nil { + log.Debug(gctx, "could not resolve a repo's PDS for sharding", "did", items[i].DID, "err", err) + return nil + } + items[i].Lane = reposync.HostKey(ident.PDSEndpoint()) + resolved.Add(1) + return nil + }) + } + // A cancelled context is the only error this can produce, and the caller is + // about to notice it for itself. + _ = g.Wait() + + for _, i := range todo { + if items[i].Lane == "" { + items[i].Lane = sweepLane(items[i].DID, "") + } + } + log.Log(ctx, "resolved PDS hosts to shard the sweep", "repos", len(todo), + "resolved", resolved.Load(), "took", time.Since(start)) +} + +// hostLanes groups items into one lane per [sweepItem.Lane], keeping each lane's +// items in input order and the lanes in order of first appearance. +// +// Both orders matter. Input order is priority order (own DIDs first, see +// [prioritizeDIDs]), so the lane holding this node's own repos is the first lane +// [runLanes] starts. +func hostLanes(items []sweepItem) [][]sweepItem { + lanes := make([][]sweepItem, 0, len(items)) + index := make(map[string]int, len(items)) + for _, item := range items { + i, ok := index[item.Lane] + if !ok { + index[item.Lane] = len(lanes) + lanes = append(lanes, []sweepItem{item}) + continue + } + lanes[i] = append(lanes[i], item) + } + return lanes +} + +// runLanes works every lane, each in its own goroutine and each one item at a +// time, with at most limit lanes in flight. Lanes are started in order, so when +// there are more lanes than slots the earliest lanes go first. +// +// One worker per host is the whole point. A PDS gives a client something like +// ten requests a second and a single range walk already uses five to seven, so +// pointing several walks at one host wins nothing: they interleave their chunk +// fetches through the per-host pdsLock and every one of them crawls. Measured on +// a production sweep, repos on a contended host walked at 12-30 records/s +// against 120-144 uncontended. Lanes make that contention structurally +// impossible within a sweep, while the limit keeps the total request rate across +// the network bounded. +// +// work never fails the sweep -- a repo that errors is logged by the worker and +// its lane moves on to the next repo. Only a cancelled context stops the run, +// and it stops it between items. +func runLanes(ctx context.Context, limit int, lanes [][]sweepItem, work func(context.Context, sweepItem)) error { + if limit <= 0 { + limit = config.DefaultSweepConcurrency + } + g, gctx := errgroup.WithContext(ctx) + g.SetLimit(limit) + for _, lane := range lanes { + g.Go(func() error { + for _, item := range lane { + if err := gctx.Err(); err != nil { + return err + } + work(gctx, item) + } + return nil + }) + } + return g.Wait() +} + +// sweepConcurrency is how many host lanes this node runs at once. +func (atsync *ATProtoSynchronizer) sweepConcurrency() int { + if atsync.CLI != nil && atsync.CLI.SweepConcurrency > 0 { + return atsync.CLI.SweepConcurrency + } + return config.DefaultSweepConcurrency +} + +// sweepDIDs is the DIDs of items, for the row lookups that work in DIDs. +func sweepDIDs(items []sweepItem) []string { + dids := make([]string, len(items)) + for i, item := range items { + dids[i] = item.DID + } + return dids +} + // Sweep brings every repo this node knows about up to date, in two phases: // first a shallow sync of anything never indexed, then history deepening for // everything that is not complete. @@ -57,7 +231,7 @@ func (atsync *ATProtoSynchronizer) Sweep(ctx context.Context) error { if err != nil { return err } - log.Log(ctx, "starting backfill sweep", "totalRepos", len(dids)) + log.Log(ctx, "starting backfill sweep", "totalRepos", len(dids), "concurrency", atsync.sweepConcurrency()) progress := &sweepProgress{} stop := progress.start(ctx) @@ -165,7 +339,7 @@ func prioritizeDIDs(dids []string, first ...string) []string { // run that died, which is the same thing as far as anyone reading the index is // concerned. func (atsync *ATProtoSynchronizer) sweepShallow(ctx context.Context, progress *sweepProgress, dids []string) error { - var todo []string + var todo []sweepItem for _, did := range dids { repo, err := atsync.Model.GetRepo(did) if err != nil { @@ -174,39 +348,39 @@ func (atsync *ATProtoSynchronizer) sweepShallow(ctx context.Context, progress *s if repo != nil && repo.Version != "" { continue } - todo = append(todo, did) + pds := "" + if repo != nil { + pds = repo.PDS + } + // Left empty when the row does not name a host: resolveLanes below fills + // those in rather than letting them each become a lane. + todo = append(todo, sweepItem{DID: did, Lane: reposync.HostKey(pds)}) } progress.begin(sweepPhaseShallow, len(todo), time.Now().Add(-InitialWindow)) if len(todo) == 0 { return nil } - log.Log(ctx, "syncing repos", "phase", sweepPhaseShallow, "repos", len(todo)) - - var mu sync.Mutex - failed := 0 - g, gctx := errgroup.WithContext(ctx) - g.SetLimit(sweepConcurrency) - for _, did := range todo { - g.Go(func() error { - if err := gctx.Err(); err != nil { - return err - } - if _, err := atsync.SyncBlueskyRepoCached(gctx, did); err != nil { - log.Error(gctx, "failed to sync repo", "did", did, "err", err) - mu.Lock() - failed++ - mu.Unlock() - return nil - } - progress.finished() - return nil - }) + atsync.resolveLanes(ctx, todo) + if err := ctx.Err(); err != nil { + return err } - if err := g.Wait(); err != nil { + lanes := hostLanes(todo) + log.Log(ctx, "syncing repos", "phase", sweepPhaseShallow, "repos", len(todo), "hosts", len(lanes)) + + var failed atomic.Int64 + err := runLanes(ctx, atsync.sweepConcurrency(), lanes, func(ctx context.Context, item sweepItem) { + if _, err := atsync.SyncBlueskyRepoCached(ctx, item.DID); err != nil { + log.Error(ctx, "failed to sync repo", "did", item.DID, "err", err) + failed.Add(1) + return + } + progress.finished() + }) + if err != nil { return err } - if failed == len(todo) { - return fmt.Errorf("all %d repos failed to sync", failed) + if int(failed.Load()) == len(todo) { + return fmt.Errorf("all %d repos failed to sync", len(todo)) } return nil } @@ -215,6 +389,12 @@ func (atsync *ATProtoSynchronizer) sweepShallow(ctx context.Context, progress *s // one window per repo per round. Round-robin rather than draining each repo is // the point: it is what puts the same horizon behind every account. // +// Each round is a barrier: every repo gets its window, then the next round +// starts. Within a round the work is sharded by host the same way the shallow +// phase shards it, so the breadth-first guarantee ("everyone reaches 7d before +// anyone starts 30d") survives lanes untouched -- a fast host simply waits at +// the end of the round instead of racing ahead through the ladder. +// // A repo that fails a round drops out of this sweep and keeps its watermark, so // the next sweep picks it up exactly where it stopped. func (atsync *ATProtoSynchronizer) sweepDeepen(ctx context.Context, progress *sweepProgress, dids []string) error { @@ -231,54 +411,47 @@ func (atsync *ATProtoSynchronizer) sweepDeepen(ctx context.Context, progress *sw if len(pending) == 0 { return nil } - log.Log(ctx, "deepening repo history", "phase", sweepPhaseDeepen, "repos", len(pending)) + log.Log(ctx, "deepening repo history", "phase", sweepPhaseDeepen, "repos", len(pending), + "hosts", len(hostLanes(pending))) for round := 0; len(pending) > 0 && round < maxDeepenRounds; round++ { if err := ctx.Err(); err != nil { return err } var mu sync.Mutex - var next []string - g, gctx := errgroup.WithContext(ctx) - g.SetLimit(sweepConcurrency) - for _, did := range pending { - g.Go(func() error { - if err := gctx.Err(); err != nil { - return err - } - done, err := atsync.DeepenRepo(gctx, did) - if err != nil { - log.Error(gctx, "failed to deepen repo history", "did", did, "err", err) - return nil - } - if done { - progress.finished() - return nil - } - mu.Lock() - next = append(next, did) - mu.Unlock() - return nil - }) - } - if err := g.Wait(); err != nil { + var next []sweepItem + err := runLanes(ctx, atsync.sweepConcurrency(), hostLanes(pending), func(ctx context.Context, item sweepItem) { + done, err := atsync.DeepenRepo(ctx, item.DID) + if err != nil { + log.Error(ctx, "failed to deepen repo history", "did", item.DID, "err", err) + return + } + if done { + progress.finished() + return + } + mu.Lock() + next = append(next, item) + mu.Unlock() + }) + if err != nil { return err } // Restore the priority order the round scrambled. - sort.Slice(next, func(i, j int) bool { return rank[next[i]] < rank[next[j]] }) + sort.Slice(next, func(i, j int) bool { return rank[next[i].DID] < rank[next[j].DID] }) pending = next - if _, horizon, err := atsync.deepenPending(ctx, pending); err == nil { + if _, horizon, err := atsync.deepenPending(ctx, sweepDIDs(pending)); err == nil { progress.setHorizon(horizon) } } return nil } -// deepenPending is the subset of dids whose history is incomplete, plus the -// sweep's horizon: the most recent floor among them, which is the instant after -// which every one of these repos is fully indexed. -func (atsync *ATProtoSynchronizer) deepenPending(ctx context.Context, dids []string) ([]string, time.Time, error) { - var pending []string +// deepenPending is the subset of dids whose history is incomplete, laned by +// host, plus the sweep's horizon: the most recent floor among them, which is the +// instant after which every one of these repos is fully indexed. +func (atsync *ATProtoSynchronizer) deepenPending(ctx context.Context, dids []string) ([]sweepItem, time.Time, error) { + var pending []sweepItem var horizon time.Time for _, did := range dids { repo, err := atsync.Model.GetRepo(did) @@ -291,7 +464,7 @@ func (atsync *ATProtoSynchronizer) deepenPending(ctx context.Context, dids []str if repo == nil || repo.Version == "" || repo.TerminalStatus() || repo.BackfillDone { continue } - pending = append(pending, did) + pending = append(pending, sweepItem{DID: did, Lane: sweepLane(did, repo.PDS)}) floor := time.Now() if repo.BackfillFloor != "" { if t, err := reposync.TimeForTID(repo.BackfillFloor); err == nil { diff --git a/pkg/atproto/sweep_test.go b/pkg/atproto/sweep_test.go index e6332432b..42fc01e48 100644 --- a/pkg/atproto/sweep_test.go +++ b/pkg/atproto/sweep_test.go @@ -7,11 +7,14 @@ import ( "testing" "time" + "github.com/bluesky-social/indigo/atproto/identity" + "github.com/bluesky-social/indigo/atproto/syntax" "github.com/bluesky-social/indigo/xrpc" "github.com/ipfs/go-cid" "github.com/stretchr/testify/require" "stream.place/streamplace/pkg/aqhttp" "stream.place/streamplace/pkg/bus" + "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/devenv" "stream.place/streamplace/pkg/model" "stream.place/streamplace/pkg/placestream" @@ -214,6 +217,200 @@ func TestSweepPrioritizesOwnDIDs(t *testing.T) { require.Nil(t, prioritizeDIDs(nil, "did:web:server.example")) } +// TestSweepHostLanes: the bucketing a sweep's whole throughput rests on. Repos +// are grouped by PDS host, in the order they arrive, so the lane list starts +// with the lane holding whatever prioritizeDIDs put first. +func TestSweepHostLanes(t *testing.T) { + // A PDS is a host however its URL was written down. + require.Equal(t, "pds.example", sweepLane("did:plc:a", "https://pds.example")) + require.Equal(t, "pds.example", sweepLane("did:plc:a", "https://PDS.Example/")) + // A row that does not name one gets a lane of its own, keyed by DID so it + // can never collide with a host. + require.Equal(t, "did:did:plc:a", sweepLane("did:plc:a", "")) + require.Equal(t, "did:did:plc:a", sweepLane("did:plc:a", " ")) + + items := []sweepItem{ + {DID: "own", Lane: sweepLane("own", "https://own.example")}, + {DID: "a1", Lane: sweepLane("a1", "https://a.example")}, + {DID: "b1", Lane: sweepLane("b1", "https://b.example")}, + {DID: "a2", Lane: sweepLane("a2", "https://a.example")}, + {DID: "u1", Lane: sweepLane("u1", "")}, + {DID: "a3", Lane: sweepLane("a3", "https://A.EXAMPLE/")}, + {DID: "u2", Lane: sweepLane("u2", "")}, + } + lanes := hostLanes(items) + + require.Equal(t, [][]string{ + {"own"}, // the priority DID's host, first because it was first + {"a1", "a2", "a3"}, // one lane per host, whatever the URL looked like + {"b1"}, + {"u1"}, // and unknown-PDS rows do not queue up behind each other + {"u2"}, + }, laneDIDs(lanes)) + + require.Empty(t, hostLanes(nil)) +} + +// TestSweepResolvesUnknownHosts: the sweep's DID list and the PDS column live in +// different databases, so a node with a fresh index knows which repos to sync and +// nothing about where they live. Those repos have to be placed before the sharding +// means anything -- a lane each would be the flat worker pool all over again. +func TestSweepResolvesUnknownHosts(t *testing.T) { + dir := identity.NewMockDirectory() + insert := func(did, pds string) { + dir.Insert(identity.Identity{ + DID: syntax.DID(did), + Handle: syntax.HandleInvalid, + Services: map[string]identity.ServiceEndpoint{"atproto_pds": {Type: "AtprotoPersonalDataServer", URL: pds}}, + }) + } + insert("did:plc:one", "https://shared.example") + insert("did:plc:two", "https://shared.example/") + insert("did:plc:three", "https://elsewhere.example") + // Pre-set so resolveIdent never reaches for a real directory. + atsync := &ATProtoSynchronizer{PLCDirectory: &dir, CachedPLCDirectory: &dir} + + items := []sweepItem{ + {DID: "did:plc:known", Lane: sweepLane("did:plc:known", "https://known.example")}, + {DID: "did:plc:one"}, + {DID: "did:plc:two"}, + {DID: "did:plc:missing"}, // no DID document: nothing to place it by + {DID: "did:plc:three"}, + } + atsync.resolveLanes(context.Background(), items) + + require.Equal(t, [][]string{ + {"did:plc:known"}, // a row that named its host is left alone + {"did:plc:one", "did:plc:two"}, // and two resolving to one host share a lane + {"did:plc:missing"}, // unplaceable: its own lane, not a queue + {"did:plc:three"}, + }, laneDIDs(hostLanes(items))) + + // Nothing to do is the normal case, and it must not cost a lookup. + placed := []sweepItem{{DID: "did:plc:known", Lane: "known.example"}} + atsync.resolveLanes(context.Background(), placed) + require.Equal(t, "known.example", placed[0].Lane) +} + +// TestSweepLanesNeverShareAHost is the property the lanes exist for: a sweep +// never has two workers on one PDS at the same time, however many workers it is +// allowed. Nothing else in a sweep is worth optimizing until that holds -- walks +// that share a host interleave their chunk fetches and run four to ten times +// slower. +func TestSweepLanesNeverShareAHost(t *testing.T) { + const cap = 3 + var items []sweepItem + for i := 0; i < 20; i++ { + host := fmt.Sprintf("pds%d.example", i%4) + items = append(items, sweepItem{DID: fmt.Sprintf("did:plc:%d", i), Lane: sweepLane("", "https://"+host)}) + } + + var mu sync.Mutex + active := map[string]string{} // lane -> the DID holding it + var order []string + inFlight, maxInFlight := 0, 0 + err := runLanes(context.Background(), cap, hostLanes(items), func(ctx context.Context, item sweepItem) { + mu.Lock() + holder, busy := active[item.Lane] + require.False(t, busy, "%s and %s ran on %s at once", item.DID, holder, item.Lane) + active[item.Lane] = item.DID + inFlight++ + maxInFlight = max(maxInFlight, inFlight) + mu.Unlock() + + // Long enough that a broken limiter or a shared lane would overlap here, + // short enough to be free. + time.Sleep(2 * time.Millisecond) + + mu.Lock() + delete(active, item.Lane) + inFlight-- + order = append(order, item.DID) + mu.Unlock() + }) + require.NoError(t, err) + require.Len(t, order, len(items), "every repo ran exactly once") + require.LessOrEqual(t, maxInFlight, cap, "the cap bounds lanes in flight") + require.Greater(t, maxInFlight, 1, "and lanes really do run in parallel") + + // Four hosts, cap of three: at least one lane waited for a slot, which is + // the case that has to not deadlock. + require.Equal(t, 4, len(hostLanes(items))) +} + +// TestSweepLanesRunOwnDIDsFirst: the node's own repos hold what it serves, so +// their lane is the first one scheduled -- the priority order prioritizeDIDs +// produces has to survive the bucketing. +func TestSweepLanesRunOwnDIDsFirst(t *testing.T) { + dids := prioritizeDIDs([]string{"did:plc:a", "did:web:server.example", "did:plc:b"}, "did:web:server.example") + items := make([]sweepItem, 0, len(dids)) + for _, did := range dids { + // Every repo on its own host, so lane order is the only thing deciding. + items = append(items, sweepItem{DID: did, Lane: sweepLane(did, "https://"+did+".pds.example")}) + } + + var mu sync.Mutex + var order []string + // One slot: lanes are started in order, so the first thing that runs is the + // first lane. + require.NoError(t, runLanes(context.Background(), 1, hostLanes(items), func(ctx context.Context, item sweepItem) { + mu.Lock() + defer mu.Unlock() + order = append(order, item.DID) + })) + require.Equal(t, []string{"did:web:server.example", "did:plc:a", "did:plc:b"}, order) +} + +// TestSweepLanesStopOnCancel: a sweep is cancellable at every point, and a lane +// checks the context between repos rather than after all of them. +func TestSweepLanesStopOnCancel(t *testing.T) { + items := make([]sweepItem, 0, 40) + for i := 0; i < 40; i++ { + items = append(items, sweepItem{DID: fmt.Sprintf("did:plc:%d", i), Lane: "pds.example"}) + } + ctx, cancel := context.WithCancel(context.Background()) + var mu sync.Mutex + ran := 0 + err := runLanes(ctx, 4, hostLanes(items), func(ctx context.Context, item sweepItem) { + mu.Lock() + ran++ + if ran == 2 { + cancel() + } + mu.Unlock() + }) + require.ErrorIs(t, err, context.Canceled) + mu.Lock() + defer mu.Unlock() + require.Less(t, ran, len(items), "the run stopped instead of draining the lane") +} + +// TestSweepConcurrencyFlag: the cap comes from --sweep-concurrency, and an unset +// or nonsense value is the documented default. +func TestSweepConcurrencyFlag(t *testing.T) { + require.Equal(t, config.DefaultSweepConcurrency, (&ATProtoSynchronizer{}).sweepConcurrency(), + "a synchronizer without a CLI still sweeps") + require.Equal(t, config.DefaultSweepConcurrency, + (&ATProtoSynchronizer{CLI: &config.CLI{}}).sweepConcurrency(), "unset means the default") + require.Equal(t, config.DefaultSweepConcurrency, + (&ATProtoSynchronizer{CLI: &config.CLI{SweepConcurrency: -1}}).sweepConcurrency()) + require.Equal(t, 64, + (&ATProtoSynchronizer{CLI: &config.CLI{SweepConcurrency: 64}}).sweepConcurrency()) +} + +// laneDIDs renders lanes for comparison. +func laneDIDs(lanes [][]sweepItem) [][]string { + out := make([][]string, 0, len(lanes)) + for _, lane := range lanes { + dids := make([]string, 0, len(lane)) + for _, item := range lane { + dids = append(dids, item.DID) + } + out = append(out, dids) + } + return out +} + // TestSweepProgressStatusLine covers the one line an operator watches: it names // the phase, counts finished repos against the total, and reports the horizon // as unix seconds. diff --git a/pkg/config/config.go b/pkg/config/config.go index 646f45f0d..6bc0134c6 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -170,8 +170,18 @@ type CLI struct { ViewCountAggregateLag time.Duration VODConcurrency int MaximumLiveBitrate int + SweepConcurrency int } +// DefaultSweepConcurrency is how many PDS hosts the atproto backfill sweep +// works on at once when --sweep-concurrency is unset or zero. +// +// The sweep shards its work by host and gives each host one worker, so this +// bounds remote servers rather than repos: 32 of them is a few hundred requests +// per second spread across the whole network, and no more than one walk (5-7 +// requests per second) against any single PDS. +const DefaultSweepConcurrency = 32 + // ContentFilters represents the content filtering configuration type ContentFilters struct { ContentWarnings struct { @@ -812,6 +822,13 @@ func (cli *CLI) NewCommand(name string) *urfavecli.Command { Destination: &cli.VODConcurrency, Sources: urfavecli.EnvVars("SP_VOD_CONCURRENCY"), }, + &urfavecli.IntFlag{ + Name: "sweep-concurrency", + Usage: "how many PDS hosts the atproto backfill sweep talks to at once. Work is sharded by host and each host is walked by one worker, so this is a count of remote servers, not of repos; 0 for the default", + Value: DefaultSweepConcurrency, + Destination: &cli.SweepConcurrency, + Sources: urfavecli.EnvVars("SP_SWEEP_CONCURRENCY"), + }, &urfavecli.StringFlag{ Name: "maximum-live-bitrate", Usage: "maximum allowed live ingest bitrate, measured per emitted segment. Accepts a bits-per-second number or a decimal SI suffix — e.g. 30M, 30000k, or 30000000 (all 30 Mbps). A stream whose bitrate exceeds this (plus a 10% margin) is disconnected and the streamer is shown a problem. 0 = unlimited", diff --git a/pkg/config/config_test.go b/pkg/config/config_test.go new file mode 100644 index 000000000..9771eabcf --- /dev/null +++ b/pkg/config/config_test.go @@ -0,0 +1,31 @@ +package config + +import ( + "context" + "testing" + + "github.com/stretchr/testify/require" + urfavecli "github.com/urfave/cli/v3" +) + +// TestSweepConcurrencyFlag: the sweep's host-lane cap is settable from the +// command line and the environment, and every command built from NewCommand -- +// including `streamplace sync`, which is the one an operator uses to warm an +// index -- gets it. +func TestSweepConcurrencyFlag(t *testing.T) { + run := func(t *testing.T, args ...string) *CLI { + t.Helper() + cli := &CLI{} + cmd := cli.NewCommand("sync") + cmd.Action = func(context.Context, *urfavecli.Command) error { return nil } + require.NoError(t, cmd.Run(context.Background(), append([]string{"sync"}, args...))) + return cli + } + + require.Equal(t, DefaultSweepConcurrency, run(t).SweepConcurrency, "unset is the default") + require.Equal(t, 8, run(t, "--sweep-concurrency", "8").SweepConcurrency) + require.Equal(t, 8, run(t, "--sweep-concurrency=8").SweepConcurrency) + + t.Setenv("SP_SWEEP_CONCURRENCY", "12") + require.Equal(t, 12, run(t).SweepConcurrency) +} diff --git a/pkg/reposync/doc.go b/pkg/reposync/doc.go index 8dad20459..6a0b2e9ec 100644 --- a/pkg/reposync/doc.go +++ b/pkg/reposync/doc.go @@ -46,6 +46,14 @@ // exponential backoff ([RetryPolicy]). Everything else -- 4xx, verification // failures, a cancelled context -- fails immediately. // +// Guessing at the backoff is the last resort, not the first: a host that +// answers 429 or 503 usually says when to come back, and [BackoffHints] is how +// that gets read. Install its [BackoffHints.Transport] on the http.Client behind +// the xrpc.Client and point [RetryPolicy.Hints] at the same registry; waits then +// honor Retry-After and ratelimit-reset instead of a ladder. Without it the +// headers are simply lost -- indigo's xrpc client keeps a status code and +// discards the response headers. +// // 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 diff --git a/pkg/reposync/fetcher.go b/pkg/reposync/fetcher.go index 50b513237..29b468ce4 100644 --- a/pkg/reposync/fetcher.go +++ b/pkg/reposync/fetcher.go @@ -96,6 +96,9 @@ func (f *XRPCBlockFetcher) GetBlocks(ctx context.Context, cids []cid.Cid) (map[c if chunk <= 0 { chunk = DefaultChunkSize } + // The policy needs to know which host it is backing off from, and the + // client is the only thing that knows. + retry := f.Retry.forHost(f.Client.Host) for start := 0; start < len(want); start += chunk { end := start + chunk if end > len(want) { @@ -111,7 +114,7 @@ func (f *XRPCBlockFetcher) GetBlocks(ctx context.Context, cids []cid.Cid) (map[c // 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 { + err := retry.do(ctx, what, func() error { var err error raw, err = indigoat.SyncGetBlocks(ctx, f.Client, strs, f.DID) return err diff --git a/pkg/reposync/head.go b/pkg/reposync/head.go index 647d0341a..000228aca 100644 --- a/pkg/reposync/head.go +++ b/pkg/reposync/head.go @@ -45,6 +45,7 @@ func FetchVerifiedHead(ctx context.Context, client *xrpc.Client, f BlockFetcher, if len(retry) == 1 { policy = retry[0] } + policy = policy.forHost(client.Host) parsedDID, err := syntax.ParseDID(did) if err != nil { diff --git a/pkg/reposync/hints.go b/pkg/reposync/hints.go new file mode 100644 index 000000000..ab1268fc6 --- /dev/null +++ b/pkg/reposync/hints.go @@ -0,0 +1,237 @@ +package reposync + +import ( + "net/http" + "strconv" + "strings" + "sync" + "time" +) + +// Backoff hint registry tuning. +const ( + // maxHintHosts bounds the registry. A node syncs repos from a few thousand + // PDS hosts at most, and only the ones actively throttling us are in here, + // so this is generous; it exists so that a pathological run cannot grow the + // map without limit. + maxHintHosts = 512 + + // hintMaxAge is how long an observation is allowed to influence a wait. A + // host that said "come back in an hour" ten minutes ago is not evidence + // about the next request: rate limit windows roll, deploys finish, and the + // wait is clamped to [RetryPolicy.MaxDelay] anyway, so a stale hint can only + // make us sleep the maximum for no reason. + hintMaxAge = 5 * time.Minute +) + +// BackoffHint is what a host told us about when it wants to be asked again. +type BackoffHint struct { + // Until is the earliest instant the host said it would answer properly. + Until time.Time + // Source names the header Until came from, for logging: "retry-after" or + // "ratelimit-reset". + Source string + // Observed is when the response carrying the header arrived. + Observed time.Time +} + +// BackoffHints remembers, per host, the last backoff a host asked for. +// +// It exists because indigo's xrpc client throws response headers away: it +// returns an [xrpc.Error] carrying a status code, a decoded error body if there +// was one, and a RatelimitInfo only when the full ratelimit-* header set was +// present. It never looks at Retry-After at all. In production every observed +// 429 arrived with Ratelimit nil, so every retry fell back to a guessed ladder +// while the host was telling us exactly how long to wait. +// +// The fix is to watch the responses ourselves: install [BackoffHints.Transport] +// on the http.Client the xrpc.Client uses, point a [RetryPolicy] at the same +// registry, and a retry waits for what the host asked for instead of guessing. +// The two halves are deliberately decoupled -- the transport sees hosts, not +// repos, and the policy reads a host key -- so nothing has to thread a response +// through the walker. +// +// A nil *BackoffHints is a working no-op registry, so a policy without one +// behaves exactly as it did before. +type BackoffHints struct { + mu sync.Mutex + hints map[string]BackoffHint +} + +// NewBackoffHints returns an empty registry. +func NewBackoffHints() *BackoffHints { + return &BackoffHints{hints: map[string]BackoffHint{}} +} + +// HostKey normalizes a PDS base URL ("https://porcini.example.net/") or a bare +// host ("porcini.example.net") to the key the registry uses. It is exported +// because callers that shard work per PDS want to agree with the registry about +// what one host is. +func HostKey(hostOrURL string) string { + s := strings.TrimSpace(hostOrURL) + if i := strings.Index(s, "://"); i >= 0 { + s = s[i+3:] + } + if i := strings.IndexAny(s, "/?#"); i >= 0 { + s = s[:i] + } + return strings.ToLower(s) +} + +// Observe records what host's headers say about backing off, if anything. Only +// 429 and 503 responses are interesting: ratelimit-* headers ride along on +// perfectly good responses too, and treating those as a hint would throttle a +// healthy walk. +func (h *BackoffHints) Observe(host string, status int, header http.Header) { + h.observeAt(host, status, header, time.Now()) +} + +func (h *BackoffHints) observeAt(host string, status int, header http.Header, at time.Time) { + if h == nil { + return + } + switch status { + case http.StatusTooManyRequests, http.StatusServiceUnavailable: + default: + return + } + key := HostKey(host) + if key == "" { + return + } + until, source := parseBackoffHeaders(header, at) + if until.IsZero() || !until.After(at) { + return + } + h.mu.Lock() + defer h.mu.Unlock() + if h.hints == nil { + h.hints = map[string]BackoffHint{} + } + if _, ok := h.hints[key]; !ok { + h.makeRoom(at) + } + h.hints[key] = BackoffHint{Until: until, Source: source, Observed: at} +} + +// makeRoom keeps the map bounded, on the insert path so that no goroutine has to +// exist to do it: drop everything expired, and if that was not enough, drop the +// least recently observed host. +func (h *BackoffHints) makeRoom(now time.Time) { + if len(h.hints) < maxHintHosts { + return + } + for key, hint := range h.hints { + if hintExpired(hint, now) { + delete(h.hints, key) + } + } + for len(h.hints) >= maxHintHosts { + oldestKey, oldest := "", time.Time{} + for key, hint := range h.hints { + if oldest.IsZero() || hint.Observed.Before(oldest) { + oldestKey, oldest = key, hint.Observed + } + } + delete(h.hints, oldestKey) + } +} + +// Get returns the live hint for host, if there is one. +func (h *BackoffHints) Get(host string) (BackoffHint, bool) { + return h.get(host, time.Now()) +} + +func (h *BackoffHints) get(host string, now time.Time) (BackoffHint, bool) { + if h == nil { + return BackoffHint{}, false + } + key := HostKey(host) + if key == "" { + return BackoffHint{}, false + } + h.mu.Lock() + defer h.mu.Unlock() + hint, ok := h.hints[key] + if !ok || hintExpired(hint, now) { + return BackoffHint{}, false + } + return hint, true +} + +// Len is the number of hosts currently remembered, live or not. +func (h *BackoffHints) Len() int { + if h == nil { + return 0 + } + h.mu.Lock() + defer h.mu.Unlock() + return len(h.hints) +} + +func hintExpired(hint BackoffHint, now time.Time) bool { + return !hint.Until.After(now) || now.Sub(hint.Observed) > hintMaxAge +} + +// parseBackoffHeaders reads the two ways a host says "not yet": Retry-After, in +// either of its RFC 9110 forms (delay seconds or an HTTP-date), and +// ratelimit-reset, which the atproto reference implementation sends as unix +// seconds. When both are present the later one wins -- they are both promises +// about when the next request can succeed, and the longer promise is the one +// that has to hold. +func parseBackoffHeaders(header http.Header, at time.Time) (time.Time, string) { + var until time.Time + source := "" + consider := func(t time.Time, name string) { + if t.IsZero() || !t.After(until) { + return + } + until, source = t, name + } + if v := strings.TrimSpace(header.Get("Retry-After")); v != "" { + if secs, err := strconv.ParseInt(v, 10, 64); err == nil { + if secs > 0 { + consider(at.Add(time.Duration(secs)*time.Second), "retry-after") + } + } else if date, err := http.ParseTime(v); err == nil { + consider(date, "retry-after") + } + } + if v := strings.TrimSpace(header.Get("ratelimit-reset")); v != "" { + if secs, err := strconv.ParseInt(v, 10, 64); err == nil && secs > 0 { + consider(time.Unix(secs, 0), "ratelimit-reset") + } + } + return until, source +} + +// Transport wraps inner so that every throttled or unavailable response it sees +// lands in the registry. Nothing else about the request or response changes; in +// particular the body is untouched, so this is safe to install under any client. +// +// A nil inner means http.DefaultTransport, matching net/http. +func (h *BackoffHints) Transport(inner http.RoundTripper) http.RoundTripper { + if inner == nil { + inner = http.DefaultTransport + } + if h == nil { + return inner + } + return &hintTransport{inner: inner, hints: h} +} + +type hintTransport struct { + inner http.RoundTripper + hints *BackoffHints +} + +var _ http.RoundTripper = (*hintTransport)(nil) + +func (t *hintTransport) RoundTrip(req *http.Request) (*http.Response, error) { + resp, err := t.inner.RoundTrip(req) + if err != nil || resp == nil { + return resp, err + } + t.hints.Observe(req.URL.Host, resp.StatusCode, resp.Header) + return resp, err +} diff --git a/pkg/reposync/hints_test.go b/pkg/reposync/hints_test.go new file mode 100644 index 000000000..1c479b68d --- /dev/null +++ b/pkg/reposync/hints_test.go @@ -0,0 +1,256 @@ +package reposync + +import ( + "context" + "fmt" + "net/http" + "strconv" + "testing" + "time" + + "github.com/ipfs/go-cid" + "github.com/stretchr/testify/require" +) + +// retryAfter is a 429 that says when to come back the way a rate limiter +// actually does: a Retry-After header and no ratelimit-* set at all, which is +// the shape indigo throws away entirely. +func retryAfter(value string) failure { + return failure{ + status: http.StatusTooManyRequests, + body: `{"error":"RateLimitExceeded","message":"Rate Limit Exceeded"}`, + header: map[string]string{"Content-Type": "application/json", "Retry-After": value}, + } +} + +func TestHostKey(t *testing.T) { + for in, want := range map[string]string{ + "https://porcini.us-east.host.bsky.network": "porcini.us-east.host.bsky.network", + "https://porcini.us-east.host.bsky.network/": "porcini.us-east.host.bsky.network", + "http://PDS.Example:2583/xrpc/whatever": "pds.example:2583", + "pds.example": "pds.example", + " pds.example ": "pds.example", + "https://pds.example?x=1": "pds.example", + "": "", + } { + require.Equal(t, want, HostKey(in), "HostKey(%q)", in) + } +} + +// TestBackoffHintsObserve covers what the registry makes of a response, without +// any HTTP involved: which statuses count, both Retry-After forms, +// ratelimit-reset, and the junk a host might send instead. +func TestBackoffHintsObserve(t *testing.T) { + now := time.Now() + header := func(kv ...string) http.Header { + h := http.Header{} + for i := 0; i+1 < len(kv); i += 2 { + h.Set(kv[i], kv[i+1]) + } + return h + } + + for _, tc := range []struct { + name string + status int + header http.Header + want time.Duration // wait recorded, 0 for "nothing recorded" + source string + }{ + {"429 retry-after seconds", 429, header("Retry-After", "3"), 3 * time.Second, "retry-after"}, + {"429 retry-after http-date", 429, + header("Retry-After", now.Add(90*time.Second).UTC().Format(http.TimeFormat)), + 90 * time.Second, "retry-after"}, + {"503 retry-after", 503, header("Retry-After", "12"), 12 * time.Second, "retry-after"}, + {"429 ratelimit-reset", 429, + header("ratelimit-reset", strconv.FormatInt(now.Add(45*time.Second).Unix(), 10)), + 45 * time.Second, "ratelimit-reset"}, + // Both present: the longer promise is the one that has to hold. + {"the later of the two wins (reset)", 429, + header("Retry-After", "5", "ratelimit-reset", strconv.FormatInt(now.Add(60*time.Second).Unix(), 10)), + 60 * time.Second, "ratelimit-reset"}, + {"the later of the two wins (retry-after)", 429, + header("Retry-After", "60", "ratelimit-reset", strconv.FormatInt(now.Add(5*time.Second).Unix(), 10)), + 60 * time.Second, "retry-after"}, + // Statuses that say nothing about backing off. ratelimit-* headers ride + // along on healthy responses, and recording those would throttle a walk + // that is doing fine. + {"200 is not a backoff", 200, header("ratelimit-reset", strconv.FormatInt(now.Add(60*time.Second).Unix(), 10)), 0, ""}, + {"500 is not a backoff", 500, header("Retry-After", "30"), 0, ""}, + {"400 is not a backoff", 400, header("Retry-After", "30"), 0, ""}, + // Nothing usable in the headers. + {"429 with no headers", 429, header(), 0, ""}, + {"unparseable retry-after", 429, header("Retry-After", "soon"), 0, ""}, + {"retry-after zero", 429, header("Retry-After", "0"), 0, ""}, + {"negative retry-after", 429, header("Retry-After", "-5"), 0, ""}, + {"retry-after in the past", 429, + header("Retry-After", now.Add(-time.Hour).UTC().Format(http.TimeFormat)), 0, ""}, + {"reset in the past", 429, + header("ratelimit-reset", strconv.FormatInt(now.Add(-time.Hour).Unix(), 10)), 0, ""}, + {"reset is not a number", 429, header("ratelimit-reset", "later"), 0, ""}, + } { + t.Run(tc.name, func(t *testing.T) { + h := NewBackoffHints() + h.observeAt("https://pds.example", tc.status, tc.header, now) + hint, ok := h.get("pds.example", now) + if tc.want == 0 { + require.False(t, ok, "nothing should have been recorded") + require.Zero(t, h.Len()) + return + } + require.True(t, ok) + require.Equal(t, tc.source, hint.Source) + require.WithinDuration(t, now.Add(tc.want), hint.Until, 1500*time.Millisecond) + require.Equal(t, now, hint.Observed) + }) + } + + t.Run("a hint stops applying once its own deadline passes", func(t *testing.T) { + h := NewBackoffHints() + h.observeAt("pds.example", 429, header("Retry-After", "10"), now) + _, ok := h.get("pds.example", now.Add(9*time.Second)) + require.True(t, ok) + _, ok = h.get("pds.example", now.Add(11*time.Second)) + require.False(t, ok, "the wait it asked for has elapsed") + }) + + t.Run("a stale observation stops applying however long it asked for", func(t *testing.T) { + h := NewBackoffHints() + h.observeAt("pds.example", 429, header("Retry-After", "3600"), now) + _, ok := h.get("pds.example", now.Add(hintMaxAge-time.Second)) + require.True(t, ok) + _, ok = h.get("pds.example", now.Add(hintMaxAge+time.Second)) + require.False(t, ok) + }) + + t.Run("an unknown host has no hint", func(t *testing.T) { + h := NewBackoffHints() + h.observeAt("pds.example", 429, header("Retry-After", "10"), now) + _, ok := h.get("other.example", now) + require.False(t, ok) + _, ok = h.get("", now) + require.False(t, ok) + }) + + t.Run("the nil registry is a working no-op", func(t *testing.T) { + var h *BackoffHints + h.Observe("pds.example", 429, header("Retry-After", "10")) + _, ok := h.Get("pds.example") + require.False(t, ok) + require.Zero(t, h.Len()) + require.Equal(t, http.DefaultTransport, h.Transport(nil)) + }) + + t.Run("the map stays bounded", func(t *testing.T) { + h := NewBackoffHints() + // Twice the cap of live hints, all with the same deadline, so nothing + // can be pruned for being expired and the eviction path has to run. + for i := 0; i < maxHintHosts*2; i++ { + h.observeAt(fmt.Sprintf("pds%d.example", i), 429, header("Retry-After", "60"), now) + require.LessOrEqual(t, h.Len(), maxHintHosts) + } + require.Equal(t, maxHintHosts, h.Len()) + // The most recent observation is always the one kept. + _, ok := h.get(fmt.Sprintf("pds%d.example", maxHintHosts*2-1), now) + require.True(t, ok) + // Re-observing a host already in the map does not grow it. + before := h.Len() + h.observeAt(fmt.Sprintf("pds%d.example", maxHintHosts*2-1), 429, header("Retry-After", "90"), now) + require.Equal(t, before, h.Len()) + }) +} + +// TestBackoffHintsFromTransport is the mechanism end to end: a real HTTP round +// trip through the wrapped transport, the real getBlocks path, and a retry that +// waits for what the host asked for instead of guessing. +func TestBackoffHintsFromTransport(t *testing.T) { + ctx := context.Background() + sr := buildSignedRepo(t, testDID, exactnessPaths()) + + newFetcher := func(t *testing.T, retry RetryPolicy, script ...failure) (*XRPCBlockFetcher, *fakeHost, *BackoffHints) { + t.Helper() + hints := NewBackoffHints() + host := newFakeHost(sr) + host.blocksFailures = script + client := host.start(t) + client.Client.Transport = hints.Transport(client.Client.Transport) + retry.Hints = hints + return &XRPCBlockFetcher{Client: client, DID: testDID, Retry: retry}, host, hints + } + + // One second is the smallest Retry-After a host can express, and the point + // of the test is that we really wait it out, so this subtest costs a second. + t.Run("Retry-After is waited for", func(t *testing.T) { + f, host, hints := newFetcher(t, + // A ladder that would retry in a millisecond if left to itself, so + // the elapsed time can only have come from the header. + RetryPolicy{MaxAttempts: 3, BaseDelay: time.Millisecond, MaxDelay: 5 * time.Second}, + retryAfter("1")) + start := time.Now() + blocks, err := f.GetBlocks(ctx, []cid.Cid{sr.root}) + elapsed := time.Since(start) + require.NoError(t, err) + require.Contains(t, blocks, sr.root) + require.Equal(t, 2, host.requests) + // A ladder that short-circuits in a millisecond and a response carrying + // no ratelimit-* headers at all: a wait of a second can only have come + // from the Retry-After the transport captured. + require.GreaterOrEqual(t, elapsed, time.Second, "the host asked for a second") + require.Less(t, elapsed, 4*time.Second, "and not much more than a second") + + // And the hint retires the moment the wait it asked for has elapsed, so + // the next call to this host starts from the ladder again. + _, ok := hints.Get(f.Client.Host) + require.False(t, ok) + }) + + // The rest only need the observation, so they fail fast and never sleep. + observed := func(t *testing.T, script ...failure) (*BackoffHints, *XRPCBlockFetcher) { + t.Helper() + f, host, hints := newFetcher(t, RetryPolicy{MaxAttempts: 1}, script...) + _, err := f.GetBlocks(ctx, []cid.Cid{sr.root}) + require.Error(t, err) + require.Equal(t, 1, host.requests) + return hints, f + } + + t.Run("Retry-After as an HTTP-date", func(t *testing.T) { + hints, f := observed(t, retryAfter(time.Now().Add(30*time.Second).UTC().Format(http.TimeFormat))) + hint, ok := hints.Get(f.Client.Host) + require.True(t, ok) + require.Equal(t, "retry-after", hint.Source) + require.WithinDuration(t, time.Now().Add(30*time.Second), hint.Until, 2*time.Second) + }) + + t.Run("503 Retry-After", func(t *testing.T) { + hints, f := observed(t, failure{ + status: http.StatusServiceUnavailable, + body: "restarting", + header: map[string]string{"Retry-After": "7"}, + }) + hint, ok := hints.Get(f.Client.Host) + require.True(t, ok) + require.Equal(t, "retry-after", hint.Source) + require.WithinDuration(t, time.Now().Add(7*time.Second), hint.Until, 2*time.Second) + }) + + t.Run("ratelimit-reset without Retry-After", func(t *testing.T) { + hints, f := observed(t, throttled(time.Now().Add(20*time.Second))) + hint, ok := hints.Get(f.Client.Host) + require.True(t, ok) + require.Equal(t, "ratelimit-reset", hint.Source) + }) + + t.Run("a 429 that says nothing leaves the ladder alone", func(t *testing.T) { + hints, _ := observed(t, htmlThrottled) + require.Zero(t, hints.Len()) + }) + + t.Run("a successful walk records nothing", func(t *testing.T) { + f, host, hints := newFetcher(t, RetryPolicy{MaxAttempts: 1}) + _, err := f.GetBlocks(ctx, []cid.Cid{sr.root}) + require.NoError(t, err) + require.Equal(t, 1, host.requests) + require.Zero(t, hints.Len()) + }) +} diff --git a/pkg/reposync/retry.go b/pkg/reposync/retry.go index 12a75a497..ddf45a51c 100644 --- a/pkg/reposync/retry.go +++ b/pkg/reposync/retry.go @@ -45,6 +45,22 @@ type RetryPolicy struct { BaseDelay time.Duration // MaxDelay caps every wait. Zero means [DefaultRetryMaxDelay]. MaxDelay time.Duration + // Hints, when set, is where waits come from whenever a host has said what + // it wants: see [BackoffHints]. Nil means the ladder plus whatever indigo + // happened to parse onto the error. + Hints *BackoffHints + // Host is the PDS these calls go to -- a base URL or a bare host, either + // way -- and the key into Hints. Fetchers fill it in from their client, so + // callers only have to set Hints. + Host string +} + +// forHost returns p keyed to host, leaving an explicitly set Host alone. +func (p RetryPolicy) forHost(host string) RetryPolicy { + if p.Host == "" { + p.Host = host + } + return p } func (p RetryPolicy) withDefaults() RetryPolicy { @@ -63,17 +79,25 @@ func (p RetryPolicy) withDefaults() RetryPolicy { return p } -// delay is how long to wait after the attempt'th failure (1-based). +// hintPad is added to a wait derived from a server's clock: the second it named +// has to have actually elapsed by the time we ask again. +const hintPad = 250 * time.Millisecond + +// delay is how long to wait after the attempt'th failure (1-based), and where +// that wait came from ("" for the computed ladder). // // 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 { +// If the server said when to come back, and that is further out than the +// computed backoff, wait for what it said instead -- still clamped to MaxDelay, +// see the note there. There are two places that can come from: the ratelimit-* +// headers indigo parsed onto the error, and [BackoffHints], which is our own +// record of every header indigo discarded (notably Retry-After, which it never +// reads). Whichever reaches further out wins. +func (p RetryPolicy) delay(attempt int, err error) (time.Duration, string) { p = p.withDefaults() d := p.MaxDelay if attempt >= 1 && attempt < 31 { @@ -82,13 +106,17 @@ func (p RetryPolicy) delay(attempt int, err error) time.Duration { } } 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) + + until, source := ratelimitReset(err), "ratelimit-reset" + if hint, ok := p.Hints.Get(p.Host); ok && hint.Until.After(until) { + until, source = hint.Until, hint.Source + } + if !until.IsZero() { + if wait := time.Until(until) + hintPad; wait > d { + return min(wait, p.MaxDelay), source } } - return d + return d, "" } // do runs fn until it succeeds, fails with something not worth retrying, or @@ -108,8 +136,16 @@ func (p RetryPolicy) do(ctx context.Context, what string, fn func() error) error 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", errForLog(err)) + d, source := p.delay(attempt, err) + kv := []any{"call", what, "attempt", attempt, "wait", d} + if source != "" { + // Worth saying out loud: it is the difference between "we guessed" + // and "the host told us", which is the first thing an operator + // looking at a throttled sweep wants to know. + kv = append(kv, "waitSource", source) + } + kv = append(kv, "err", errForLog(err)) + log.Warn(ctx, "retrying transient xrpc failure", kv...) if serr := sleepCtx(ctx, d); serr != nil { return fmt.Errorf("aborted after %d attempts: %w", attempt, errors.Join(err, serr)) } diff --git a/pkg/reposync/retry_test.go b/pkg/reposync/retry_test.go index 9eff9e478..083ee0f8a 100644 --- a/pkg/reposync/retry_test.go +++ b/pkg/reposync/retry_test.go @@ -107,16 +107,17 @@ func TestRetryDelay(t *testing.T) { 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) + d, source := 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) + require.Empty(t, source, "nothing told us to wait, so nothing is reported") } } // 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) + d, _ := p.delay(10, plain) require.GreaterOrEqual(t, d, 22500*time.Millisecond) require.LessOrEqual(t, d, 30*time.Second) if d < 29*time.Second { @@ -126,29 +127,112 @@ func TestRetryDelay(t *testing.T) { 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) + zero, _ := RetryPolicy{}.delay(1, plain) + require.LessOrEqual(t, zero, DefaultRetryBaseDelay) + require.GreaterOrEqual(t, zero, 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) + d, source := p.delay(1, err) require.Greater(t, d, 2500*time.Millisecond) require.LessOrEqual(t, d, 3500*time.Millisecond) + require.Equal(t, "ratelimit-reset", source) }) 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)) + d, source := p.delay(1, err) + require.Equal(t, 30*time.Second, d) + require.Equal(t, "ratelimit-reset", source) }) 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) + d, source := p.delay(3, err) require.GreaterOrEqual(t, d, 3*time.Second) require.LessOrEqual(t, d, 4*time.Second) + require.Empty(t, source) + }) + + t.Run("a hint from the registry is honored", func(t *testing.T) { + hints := NewBackoffHints() + hints.Observe("https://pds.example/", http.StatusTooManyRequests, + http.Header{"Retry-After": []string{"3"}}) + hinted := RetryPolicy{BaseDelay: time.Second, MaxDelay: 30 * time.Second, + Hints: hints, Host: "https://pds.example"} + + d, source := hinted.delay(1, errors.New("boom")) + require.Greater(t, d, 2500*time.Millisecond) + require.LessOrEqual(t, d, 3500*time.Millisecond) + require.Equal(t, "retry-after", source) + + // A hint that is shorter than the ladder changes nothing: the ladder is + // the floor, the hint only ever pushes a wait out. + d, source = hinted.delay(5, errors.New("boom")) + require.GreaterOrEqual(t, d, 12*time.Second) + require.Empty(t, source) + + // A different host is a different budget. + other := hinted + other.Host = "other.example" + d, source = other.delay(1, errors.New("boom")) + require.LessOrEqual(t, d, time.Second) + require.Empty(t, source) + + // And so is no host at all, which is what an un-plumbed policy looks + // like. + nohost := hinted + nohost.Host = "" + _, source = nohost.delay(1, errors.New("boom")) + require.Empty(t, source) + }) + + t.Run("the further-out of the two sources wins", func(t *testing.T) { + hints := NewBackoffHints() + hints.Observe("pds.example", http.StatusTooManyRequests, + http.Header{"Retry-After": []string{"2"}}) + hinted := RetryPolicy{BaseDelay: time.Second, MaxDelay: 30 * time.Second, + Hints: hints, Host: "pds.example"} + + // indigo parsed a reset further out than the header we captured. + d, source := hinted.delay(1, ratelimited(time.Now().Add(6*time.Second))) + require.Greater(t, d, 5*time.Second) + require.Equal(t, "ratelimit-reset", source) + + // And the other way around. + d, source = hinted.delay(1, ratelimited(time.Now().Add(time.Millisecond))) + require.Greater(t, d, 1500*time.Millisecond) + require.Equal(t, "retry-after", source) + }) + + t.Run("a hint is clamped to MaxDelay", func(t *testing.T) { + hints := NewBackoffHints() + hints.Observe("pds.example", http.StatusTooManyRequests, + http.Header{"Retry-After": []string{"600"}}) + hinted := RetryPolicy{BaseDelay: time.Second, MaxDelay: 30 * time.Second, + Hints: hints, Host: "pds.example"} + d, source := hinted.delay(1, errors.New("boom")) + require.Equal(t, 30*time.Second, d) + require.Equal(t, "retry-after", source) + }) + + t.Run("a stale observation does not inflate a later wait", func(t *testing.T) { + hints := NewBackoffHints() + // An hour-long backoff, observed longer ago than hintMaxAge: the wait it + // asked for has not elapsed, but it is no longer evidence about now. + hints.observeAt("pds.example", http.StatusTooManyRequests, + http.Header{"Retry-After": []string{"3600"}}, + time.Now().Add(-hintMaxAge-time.Minute)) + hinted := RetryPolicy{BaseDelay: time.Second, MaxDelay: 30 * time.Second, + Hints: hints, Host: "pds.example"} + d, source := hinted.delay(1, errors.New("boom")) + require.LessOrEqual(t, d, time.Second) + require.Empty(t, source) + _, live := hints.Get("pds.example") + require.False(t, live) }) }