From 5d02402b9eb481783087471c23dddb120690de99 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Thu, 30 Jul 2026 14:54:56 -0700 Subject: [PATCH] atproto: dissolve the sweep's global barriers into per-host ladders A sweep had four global barriers: shallow -> deepen, then one after each deepening round. Every one of them ended in a straggler tail with most of the 32 lanes idle. Measured on a 20,747-repo overnight sweep (4h50m): the shallow phase spent 15.2 of its 28.3 minutes past its 90% mark, and the deepen rounds spent 126 of their 262 minutes in their last 10% of windows. Roughly half the wall clock was most-lanes-idle, waiting on the slowest host, four times over. The phases are gone. Each host lane now runs its own complete program -- laneProgram: a shallow queue, then a ladder of repos bucketed by how many windows they have had -- and a lane that finishes hands its slot to a host that has not started yet. Nothing waits for anything on another host. The guarantees that mattered survive, rescoped to the host they were always really about: - shallow strictly before deepening, per repo; - breadth-first within a host: taking work from the lowest non-empty ladder bucket means no repo gets its (n+1)th window while another repo on that host waits for its nth. Bucketing rather than a round-robin queue so a repo that joins late is absorbed into the order instead of being left a lap behind; - one worker per host, ever, and the global cap from --sweep-concurrency. Slots are handed out in the order lanes were added -- own DIDs first -- which now needs an explicit FIFO queue: a buffered channel would have handed them out in whatever order the worker goroutines happened to wake up; - a repo resolved late still preempts its lane's ladder; - failures never abort a lane or the sweep, cancellation stops between steps, and an all-shallow-failed sweep still errors. What is deliberately given up: the global round barrier, "everyone reaches 7d before anyone starts 30d" across hosts. That barrier was exactly the mechanism that made every lane wait on the slowest host. runLanes, hostLanes, the two-phase sweepShallow/sweepDeepen split and the per-round deepenPending scan are deleted; sweepPlan reads every row once up front and Sweep drives one scheduler. The status line loses phase= and round=, which no longer mean anything, and reports both halves at once: backfill sweep shallow=19000/20747 deepened=4300/20013 windows=41022 horizon=1753142400 horizon keeps its meaning -- the instant after which everything this node serves is indexed -- but is now maintained incrementally from the watermark each window reports (DeepenRepo returns it) rather than by rescanning every row at a round boundary, of which there are none. Also: reposync gives a host that is not there at all -- refused connection, or NXDOMAIN -- two attempts instead of five. Timeouts, 429s and 5xx keep the full ladder. The measured straggler tails are mostly switched-off PDSes, and every repo on one was costing a full backoff ladder to rediscover what the first connection attempt already said. Co-Authored-By: Claude Fable 5 --- pkg/atproto/atproto.go | 38 +- pkg/atproto/sweep.go | 782 ++++++++++++++++++++----------------- pkg/atproto/sweep_test.go | 419 +++++++++++++++----- pkg/reposync/retry.go | 35 +- pkg/reposync/retry_test.go | 112 ++++++ 5 files changed, 929 insertions(+), 457 deletions(-) diff --git a/pkg/atproto/atproto.go b/pkg/atproto/atproto.go index 07d5f07a..5c88bf7c 100644 --- a/pkg/atproto/atproto.go +++ b/pkg/atproto/atproto.go @@ -150,34 +150,36 @@ func (atsync *ATProtoSynchronizer) SyncBlueskyRepo(ctx context.Context, handle s } // DeepenRepo walks one more window of history for a repo whose recent records -// are already indexed, and reports whether that repo is now complete. +// are already indexed. It reports whether that repo is now complete, and the +// watermark it left behind: the TID from which the repo's windowed collections +// are now fully indexed, which is what the sweep's horizon is made of. // // Each call reaches one rung further back down [backfillSpans] and advances the -// row's watermark; the sweep calls it repeatedly, round-robin across repos, so -// that every account reaches a week of history before any account reaches a -// month. The records it re-emits from the window boundary are absorbed by the -// idempotent indexer. +// row's watermark; the sweep calls it repeatedly, round-robin across the repos +// on one host, so that every account there reaches a week of history before any +// of them reaches a month. The records it re-emits from the window boundary are +// absorbed by the idempotent indexer. // // It never writes a placeholder row and never blanks Version, so a repo stays // served -- and stays out of the wedge path -- for the entire time its history // is being filled in. -func (atsync *ATProtoSynchronizer) DeepenRepo(ctx context.Context, did string) (bool, error) { +func (atsync *ATProtoSynchronizer) DeepenRepo(ctx context.Context, did string) (bool, string, error) { repo, err := atsync.Model.GetRepo(did) if err != nil { - return false, fmt.Errorf("failed to get repo for %s: %w", did, err) + return false, "", fmt.Errorf("failed to get repo for %s: %w", did, err) } switch { case repo == nil: - return false, fmt.Errorf("no repo row for %s", did) + return false, "", fmt.Errorf("no repo row for %s", did) case repo.TerminalStatus(): // The account is gone; whatever we indexed is all there will be. - return true, nil + return true, repo.BackfillFloor, nil case repo.BackfillDone: - return true, nil + return true, repo.BackfillFloor, nil case repo.Version == "": - // Never synced (or wedged): that is the shallow phase's job, and doing + // Never synced (or wedged): that is the shallow sync's job, and doing // it here would skip the full collections entirely. - return false, fmt.Errorf("repo %s has no completed sync to deepen", did) + return false, "", fmt.Errorf("repo %s has no completed sync to deepen", did) } // The same lock a full sync takes, so the two cannot walk one repo at once. @@ -189,11 +191,11 @@ func (atsync *ATProtoSynchronizer) DeepenRepo(ctx context.Context, did string) ( ident, err := atsync.resolveIdent(ctx, did, true) if err != nil { - return false, fmt.Errorf("failed to resolve %s: %w", did, err) + return false, "", fmt.Errorf("failed to resolve %s: %w", did, err) } xrpcc := xrpc.Client{Host: ident.PDSEndpoint(), Client: SyncHTTPClient} if xrpcc.Host == "" { - return false, fmt.Errorf("no PDS endpoint found for %s", did) + return false, "", fmt.Errorf("no PDS endpoint found for %s", did) } window := nextBackfillWindow(repo.BackfillFloor, time.Now()) @@ -213,18 +215,18 @@ func (atsync *ATProtoSynchronizer) DeepenRepo(ctx context.Context, did string) ( } if err != nil { if parked := parkTerminalRepo(ctx, atsync.Model, did, err); parked != nil { - return false, parked + return false, "", parked } - return false, err + return false, "", err } if err := atsync.Model.AdvanceRepoBackfill(ctx, did, rev, root, window.Lo, window.Genesis); err != nil { - return false, fmt.Errorf("failed to record backfill watermark for %s: %w", did, err) + return false, "", fmt.Errorf("failed to record backfill watermark for %s: %w", did, err) } // Debug: at one line per repo per window this is thousands of lines per // sweep. The sweep logs one Info summary per repo when its ladder finishes. log.Debug(ctx, "deepened repo history", "rev", rev, "floor", window.Lo, "done", window.Genesis) - return window.Genesis, nil + return window.Genesis, window.Lo, nil } // syncsInFlight holds the DIDs whose backfill is running in this process right diff --git a/pkg/atproto/sweep.go b/pkg/atproto/sweep.go index f6a756b9..f8e7e079 100644 --- a/pkg/atproto/sweep.go +++ b/pkg/atproto/sweep.go @@ -3,7 +3,6 @@ package atproto import ( "context" "fmt" - "sort" "sync" "sync/atomic" "time" @@ -14,35 +13,38 @@ import ( "stream.place/streamplace/pkg/reposync" ) -const ( - // 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. - sweepStatusInterval = 10 * time.Second - - // sweepPhaseShallow syncs repos that have never been indexed: everything - // this node cares about, plus the last [InitialWindow] of the windowed - // collections. It is what makes an account servable. - sweepPhaseShallow = "shallow" - // sweepPhaseDeepen walks history backwards, one window at a time, for every - // repo that is not complete yet. - sweepPhaseDeepen = "deepen" -) +// 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. +const sweepStatusInterval = 10 * time.Second -// maxDeepenRounds stops the deepening loop from spinning if a repo somehow -// never reports itself finished. Every successful round moves a repo one rung -// down [backfillSpans], so the ladder is walked in len+1 rounds; the slack is -// pure belt and braces. +// maxDeepenRounds stops a repo's ladder from spinning if it somehow never +// reports itself finished. Every successful window moves a repo one rung down +// [backfillSpans], so the ladder is walked in len+1 windows; the slack is pure +// belt and braces. It is a per-repo budget for one sweep, not a global round +// count -- there are no global rounds. var maxDeepenRounds = len(backfillSpans) + 3 // sweepItem is one repo for a sweep to work on, tagged with the lane it belongs -// to. +// to and with which half of its lane's program it starts in. 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 + // Deepen is set for a repo that already has a completed sync and needs + // nothing but history: it starts in its lane's ladder rather than in its + // lane's shallow queue. + Deepen bool +} + +// sweepStep is one unit of work a lane does: either the shallow sync a repo +// needs before anything else can happen to it, or one window of its history. +type sweepStep struct { + sweepItem + // Windows is how many history windows this repo has already had this + // sweep, which for a deepening step is also the ladder rung it came off. + Windows int } // sweepLane is the lane a repo row belongs in: its PDS host, or -- for a repo @@ -63,6 +65,16 @@ func sweepLane(did, pds string) string { return "did:" + did } +// laneCount is how many distinct lanes -- hosts, in practice -- a set of items +// covers. +func laneCount(items []sweepItem) int { + lanes := make(map[string]struct{}, len(items)) + for _, item := range items { + lanes[item.Lane] = struct{}{} + } + return len(lanes) +} + // identityResolveConcurrency is how many identities [ATProtoSynchronizer.feedUnresolved] // looks up at once. // @@ -137,100 +149,250 @@ func (atsync *ATProtoSynchronizer) feedUnresolved(ctx context.Context, items []s "resolved", resolved.Load(), "took", time.Since(start)) } -// laneScheduler is [runLanes] for work that is still being discovered: it keeps -// the one-worker-per-lane guarantee and the global lane cap, but accepts items -// while it is running, so repos whose lane is already known are walked while a -// resolver is still finding hosts for the rest. +// laneProgram is one host's entire share of a sweep, and the reason a sweep has +// no phases: rather than every lane doing shallow work until the slowest host +// has finished its shallow work, a lane runs this program to completion by +// itself and then gives its slot to the next host. +// +// Shallow work always comes first, because a repo with no completed sync cannot +// be deepened at all, and because a repo that has just been discovered is not +// servable until it has one. Deepening is round-robin within the host, which is +// what the ladder buckets are for. +type laneProgram struct { + // shallow is the repos on this host with no completed sync, oldest first. + shallow []sweepItem + // ladder holds the repos with history left to fetch, bucketed by how many + // windows they have had this sweep: ladder[n] is the repos on their (n+1)th + // window. Always taking from the lowest non-empty bucket makes the + // breadth-first guarantee structural -- no repo on this host gets its + // (n+1)th window while another still wants its nth -- and it holds for + // repos that join late, which a plain round-robin queue would leave a full + // lap behind. The bucket index is also the per-repo spin guard: nothing is + // ever pushed past [maxDeepenRounds]. + ladder [][]sweepItem + // live reports that a worker is draining this lane. Guarded by the + // scheduler's lock, like everything else here. + live bool +} + +// next takes the lane's next step: the oldest waiting shallow sync if there is +// one, otherwise the least-deepened repo's next window. +func (p *laneProgram) next() (sweepStep, bool) { + if len(p.shallow) > 0 { + item := p.shallow[0] + p.shallow = p.shallow[1:] + return sweepStep{sweepItem: item}, true + } + for n, bucket := range p.ladder { + if len(bucket) == 0 { + continue + } + item := bucket[0] + p.ladder[n] = bucket[1:] + return sweepStep{sweepItem: item, Windows: n}, true + } + return sweepStep{}, false +} + +// add puts a repo into the half of the program it belongs in. +func (p *laneProgram) add(item sweepItem) { + if item.Deepen { + p.push(item, 0) + return + } + p.shallow = append(p.shallow, item) +} + +// push queues a repo for its next window, having had windows of them already. +// A repo that has used its whole budget is dropped: it keeps its watermark, so +// the next sweep carries on where this one stopped. +func (p *laneProgram) push(item sweepItem, windows int) { + if windows >= maxDeepenRounds { + return + } + item.Deepen = true + for len(p.ladder) <= windows { + p.ladder = append(p.ladder, nil) + } + p.ladder[windows] = append(p.ladder[windows], item) +} + +// laneScheduler runs one [laneProgram] per host, at most [laneScheduler.limit] +// of them at a time, one worker per host ever. +// +// 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. +// +// It accepts items while it is running, so repos whose lane is already known are +// walked while a resolver is still finding hosts for the rest, and a repo +// resolved late lands in a lane that is already deep in its ladder -- where its +// shallow sync preempts the rest of that ladder. +// +// work never fails the sweep: a repo that errors is logged by the worker and its +// lane moves on. Only a cancelled context stops the run, and it stops it between +// steps. type laneScheduler struct { - ctx context.Context - work func(context.Context, sweepItem) - // sem caps how many lane workers run at once; a worker holds a slot for the - // life of its lane. Blocked acquisitions queue in FIFO order, so lanes - // started earlier (own DIDs first) get slots first. - sem chan struct{} + ctx context.Context + // work does one step and reports whether the repo wants another: a shallow + // step returning true puts the repo at the bottom of the ladder, a + // deepening step returning true asks for one more window. + work func(context.Context, sweepStep) bool mu sync.Mutex - queue map[string][]sweepItem - live map[string]bool + lanes map[string]*laneProgram seen int - wg sync.WaitGroup + // free and waiting are the slot budget. A worker holds a slot for the life + // of its lane, and slots are handed out strictly in the order lanes were + // first added -- own DIDs first, see [prioritizeDIDs]. A queue rather than a + // buffered channel because a channel would hand slots out in the order + // worker goroutines happened to get scheduled, which is no order at all. + free int + waiting []chan struct{} + wg sync.WaitGroup } -func newLaneScheduler(ctx context.Context, limit int, work func(context.Context, sweepItem)) *laneScheduler { +func newLaneScheduler(ctx context.Context, limit int, work func(context.Context, sweepStep) bool) *laneScheduler { if limit <= 0 { limit = config.DefaultSweepConcurrency } return &laneScheduler{ ctx: ctx, work: work, - sem: make(chan struct{}, limit), - queue: map[string][]sweepItem{}, - live: map[string]bool{}, + lanes: map[string]*laneProgram{}, + free: limit, } } // add enqueues an item on its lane, starting a worker for the lane if none is -// running. Safe from any goroutine; must not be called after [laneScheduler.wait] -// returns. +// running. Safe from any goroutine, and never blocks; must not be called after +// [laneScheduler.wait] returns. func (s *laneScheduler) add(item sweepItem) { s.mu.Lock() - if _, ok := s.queue[item.Lane]; !ok { + prog, ok := s.lanes[item.Lane] + if !ok { + prog = &laneProgram{} + s.lanes[item.Lane] = prog s.seen++ } - s.queue[item.Lane] = append(s.queue[item.Lane], item) - spawn := !s.live[item.Lane] - if spawn { - s.live[item.Lane] = true + prog.add(item) + var slot chan struct{} + if !prog.live { + prog.live = true + slot = s.acquire() s.wg.Add(1) } s.mu.Unlock() - if spawn { - go s.run(item.Lane) + if slot != nil { + go s.run(prog, slot) + } +} + +// acquire takes a slot, or a promise of one: the returned channel is closed +// when the caller may run. Called with the lock held, so that slots are queued +// in the order [laneScheduler.add] is called rather than the order goroutines +// start. +func (s *laneScheduler) acquire() chan struct{} { + slot := make(chan struct{}) + if s.free > 0 { + s.free-- + close(slot) + return slot + } + s.waiting = append(s.waiting, slot) + return slot +} + +// release hands a finished lane's slot to the longest-waiting lane, or back to +// the budget if nobody is waiting. +func (s *laneScheduler) release() { + s.mu.Lock() + defer s.mu.Unlock() + if len(s.waiting) > 0 { + slot := s.waiting[0] + s.waiting = s.waiting[1:] + close(slot) + return } + s.free++ } -// run drains one lane, one item at a time, holding a semaphore slot throughout. -// It marks the lane not-live under the lock in the same instant it observes the -// queue empty, so a concurrent add either lands before that (and this worker -// picks it up) or after (and spawns a fresh worker). -func (s *laneScheduler) run(lane string) { +// giveUp drops a slot the caller was waiting for. If the slot was granted in +// the meantime it is passed on rather than lost. +func (s *laneScheduler) giveUp(slot chan struct{}) { + s.mu.Lock() + for i, w := range s.waiting { + if w == slot { + s.waiting = append(s.waiting[:i], s.waiting[i+1:]...) + s.mu.Unlock() + return + } + } + s.mu.Unlock() + s.release() +} + +// run works one lane's program to the end, one step at a time, holding a slot +// throughout. It marks the lane not-live under the lock in the same instant it +// observes the program empty, so a concurrent add either lands before that (and +// this worker picks it up) or after (and spawns a fresh worker). +func (s *laneScheduler) run(prog *laneProgram, slot chan struct{}) { defer s.wg.Done() select { - case s.sem <- struct{}{}: + case <-slot: case <-s.ctx.Done(): - s.abandon(lane) + s.giveUp(slot) + s.abandon(prog) return } - defer func() { <-s.sem }() + defer s.release() for { if s.ctx.Err() != nil { - s.abandon(lane) + s.abandon(prog) return } s.mu.Lock() - if len(s.queue[lane]) == 0 { - s.live[lane] = false + step, ok := prog.next() + if !ok { + prog.live = false s.mu.Unlock() return } - item := s.queue[lane][0] - s.queue[lane] = s.queue[lane][1:] s.mu.Unlock() - s.work(s.ctx, item) + + if !s.work(s.ctx, step) { + continue + } + // A shallow sync that worked drops the repo at the bottom of the + // ladder; a window that worked asks for the next rung. + windows := 0 + if step.Deepen { + windows = step.Windows + 1 + } + s.mu.Lock() + prog.push(step.sweepItem, windows) + s.mu.Unlock() } } -// abandon drops a lane's remaining items on cancellation. The repos keep their +// abandon drops a lane's remaining work on cancellation. The repos keep their // rows untouched, so the next sweep picks them up. -func (s *laneScheduler) abandon(lane string) { +func (s *laneScheduler) abandon(prog *laneProgram) { s.mu.Lock() - s.live[lane] = false - s.queue[lane] = nil - s.mu.Unlock() + defer s.mu.Unlock() + prog.live = false + prog.shallow = nil + prog.ladder = nil } -// wait blocks until every added item has been worked or abandoned, and reports -// how many distinct lanes the run touched. Callers must have finished adding. +// wait blocks until every lane has run its program out or been abandoned, and +// reports how many distinct lanes the run touched. Callers must have finished +// adding. func (s *laneScheduler) wait() (lanes int, err error) { s.wg.Wait() s.mu.Lock() @@ -238,63 +400,6 @@ func (s *laneScheduler) wait() (lanes int, err error) { return s.seen, s.ctx.Err() } -// 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 { @@ -303,24 +408,21 @@ func (atsync *ATProtoSynchronizer) sweepConcurrency() int { 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. +// Sweep brings every repo this node knows about up to date: a shallow sync for +// anything never indexed, then history deepening until it is complete. // -// It is breadth-first on purpose. The shallow phase makes accounts servable as -// fast as it can, and the deepening phase gives every repo one window before it -// gives any repo two, so a node coming up with ten thousand accounts reaches a -// week of history everywhere rather than five years of history for the first -// hundred DIDs in the table. +// Both happen at once, because a sweep is thousands of independent +// conversations with hundreds of servers and the slowest of them must not hold +// up the rest. Each host gets a lane, each lane runs its own program -- see +// [laneProgram] -- and a host that finishes early hands its slot to a host that +// has not started. The sweep is over when the last lane is. +// +// Within a host it is still breadth-first: shallow syncs first, so accounts +// become servable as fast as they can, and then one window per repo before any +// repo gets two, so a node coming up reaches a week of history everywhere on +// that host rather than five years for the first few DIDs. Across hosts there is +// no such guarantee, and buying it was what cost half the wall clock: it made +// every lane wait for the slowest host, once per rung of the ladder. // // Nothing on the node waits for this. Repos that fail are logged and left for // the next sweep -- their rows keep whatever they had -- except that a sweep @@ -333,23 +435,80 @@ func (atsync *ATProtoSynchronizer) Sweep(ctx context.Context) error { } log.Log(ctx, "starting backfill sweep", "totalRepos", len(dids), "concurrency", atsync.sweepConcurrency()) + plan, err := atsync.sweepPlan(dids) + if err != nil { + return err + } progress := &sweepProgress{} + progress.begin(plan.shallow, plan.floors) stop := progress.start(ctx) defer stop() - if err := atsync.sweepShallow(ctx, progress, dids); err != nil { - return err + log.Log(ctx, "sweeping repos", "shallow", plan.shallow, "deepen", len(plan.floors), + "knownHosts", laneCount(plan.ready), "unresolved", len(plan.unresolved)) + + var failed atomic.Int64 + sched := newLaneScheduler(ctx, atsync.sweepConcurrency(), func(ctx context.Context, step sweepStep) bool { + if !step.Deepen { + return atsync.sweepSync(ctx, progress, &failed, step) + } + return atsync.sweepWindow(ctx, progress, step) + }) + // Lanes whose host is already known start working immediately, in priority + // order (own DIDs first); the rest stream in as the resolver finds them. + for _, item := range plan.ready { + sched.add(item) } - if err := ctx.Err(); err != nil { + atsync.feedUnresolved(ctx, plan.unresolved, sched.add) + lanes, err := sched.wait() + if err != nil { return err } - if err := atsync.sweepDeepen(ctx, progress, dids); err != nil { - return err + if plan.shallow > 0 && int(failed.Load()) == plan.shallow { + return fmt.Errorf("all %d repos failed to sync", plan.shallow) } - log.Log(ctx, "backfill sweep complete", "totalRepos", len(dids)) + log.Log(ctx, "backfill sweep complete", + append([]any{"totalRepos", len(dids), "hosts", lanes}, progress.status()...)...) return nil } +// sweepSync gives a repo the shallow sync it has never had, and reports whether +// it now has history to deepen. +func (atsync *ATProtoSynchronizer) sweepSync(ctx context.Context, progress *sweepProgress, failed *atomic.Int64, step sweepStep) bool { + repo, err := atsync.SyncBlueskyRepoCached(ctx, step.DID) + if err != nil { + log.Error(ctx, "failed to sync repo", "did", step.DID, "err", err) + failed.Add(1) + // A repo whose shallow sync failed is left alone for the rest of the + // sweep: it has no Version, so deepening it would fetch the wrong + // ranges. The next sweep retries it from the top. + return false + } + progress.synced() + if repo == nil || repo.Version == "" || repo.TerminalStatus() || repo.BackfillDone { + return false + } + progress.laddered(step.DID, backfillFloorTime(repo.BackfillFloor)) + return true +} + +// sweepWindow walks one window of a repo's history and reports whether it wants +// another. +func (atsync *ATProtoSynchronizer) sweepWindow(ctx context.Context, progress *sweepProgress, step sweepStep) bool { + done, floor, err := atsync.DeepenRepo(ctx, step.DID) + if err != nil { + log.Error(ctx, "failed to deepen repo history", "did", step.DID, "err", err) + return false + } + progress.window(step.DID, backfillFloorTime(floor)) + if done { + progress.deepened(step.DID) + log.Log(ctx, "finished deepening repo history", "did", step.DID, "windows", step.Windows+1) + return false + } + return true +} + // sweepCandidates is every repo worth syncing, own DIDs first. func (atsync *ATProtoSynchronizer) sweepCandidates(ctx context.Context) ([]string, error) { // Accounts that are deactivated, deleted, or taken down fail their backfill @@ -434,250 +593,175 @@ func prioritizeDIDs(dids []string, first ...string) []string { return out } -// sweepShallow syncs every repo that has never completed one. A repo row with -// an empty Version is exactly that: either brand new, or left half-indexed by a -// 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 []sweepItem - for _, did := range dids { - repo, err := atsync.Model.GetRepo(did) - if err != nil { - return fmt.Errorf("failed to get repo for %s: %w", did, err) - } - if repo != nil && repo.Version != "" { - continue - } - pds := "" - if repo != nil { - pds = repo.PDS - } - // Left empty when the row does not name a host: feedUnresolved below - // finds those hosts in the background rather than letting each become a - // lane -- or worse, gating the whole sweep behind the lookups. - 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 - } - - var known, unresolved []sweepItem - for _, item := range todo { - if item.Lane == "" { - unresolved = append(unresolved, item) - } else { - known = append(known, item) - } - } - log.Log(ctx, "syncing repos", "phase", sweepPhaseShallow, "repos", len(todo), - "knownHosts", len(hostLanes(known)), "unresolved", len(unresolved)) - - var failed atomic.Int64 - sched := newLaneScheduler(ctx, atsync.sweepConcurrency(), 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() - }) - // Known lanes start working immediately, in priority order (own DIDs - // first); the rest stream in as the resolver finds their hosts. - for _, item := range known { - sched.add(item) - } - atsync.feedUnresolved(ctx, unresolved, sched.add) - lanes, err := sched.wait() - if err != nil { - return err - } - log.Log(ctx, "synced repos", "phase", sweepPhaseShallow, "repos", len(todo), "hosts", lanes) - if int(failed.Load()) == len(todo) { - return fmt.Errorf("all %d repos failed to sync", len(todo)) - } - return nil +// sweepPlan is what a sweep has to do, read off the index once at the start. +type sweepPlan struct { + // ready is every repo whose lane is already known, in priority order -- + // which is the order lanes are created in, and so the order they get slots + // in. Shallow and deepening work is interleaved here rather than split, + // because splitting it is what would put this node's own repos behind a + // thousand strangers' lanes. + ready []sweepItem + // unresolved is the repos needing a shallow sync whose host the index does + // not know; [ATProtoSynchronizer.feedUnresolved] streams them in. + unresolved []sweepItem + // shallow is how many repos in total need a shallow sync, ready and + // unresolved together. + shallow int + // floors is the backfill watermark of every repo that starts in a ladder, + // for the status line's horizon. + floors map[string]time.Time } -// sweepDeepen fills in history for every repo that has some but not all of it, -// 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. +// sweepPlan sorts every candidate into the work it needs. // -// 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 { - rank := make(map[string]int, len(dids)) - for i, did := range dids { - rank[did] = i - } - - pending, horizon, err := atsync.deepenPending(ctx, dids) - if err != nil { - return err - } - progress.begin(sweepPhaseDeepen, len(pending), horizon) - if len(pending) == 0 { - return nil - } - log.Log(ctx, "deepening repo history", "phase", sweepPhaseDeepen, "repos", len(pending), - "hosts", len(hostLanes(pending))) - - // Windows completed per DID across the whole ladder, for the one-line - // summary when a repo finishes. Rounds are sequential (each is a barrier), - // so each round's mu safely guards it in turn. - windows := make(map[string]int, len(pending)) - for round := 0; len(pending) > 0 && round < maxDeepenRounds; round++ { - if err := ctx.Err(); err != nil { - return err +// A repo row with an empty Version has never completed a sync: either brand new, +// or left half-indexed by a run that died, which is the same thing as far as +// anyone reading the index is concerned. One with a Version and no BackfillDone +// has some history and wants the rest. Anything parked or complete is left +// alone. +func (atsync *ATProtoSynchronizer) sweepPlan(dids []string) (*sweepPlan, error) { + plan := &sweepPlan{floors: map[string]time.Time{}} + for _, did := range dids { + repo, err := atsync.Model.GetRepo(did) + if err != nil { + return nil, fmt.Errorf("failed to get repo for %s: %w", did, err) } - progress.setRound(round + 1) - var mu sync.Mutex - 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 + switch { + case repo == nil || repo.Version == "": + plan.shallow++ + pds := "" + if repo != nil { + pds = repo.PDS } - progress.window() - mu.Lock() - windows[item.DID]++ - n := windows[item.DID] - mu.Unlock() - if done { - progress.finished() - log.Log(ctx, "finished deepening repo history", "did", item.DID, "windows", n) - return + // A row that does not name a host does not get a lane of its own + // here: feedUnresolved finds those hosts in the background rather + // than letting each become a lane. + if host := reposync.HostKey(pds); host != "" { + plan.ready = append(plan.ready, sweepItem{DID: did, Lane: host}) + } else { + plan.unresolved = append(plan.unresolved, sweepItem{DID: did}) } - 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].DID] < rank[next[j].DID] }) - pending = next - if _, horizon, err := atsync.deepenPending(ctx, sweepDIDs(pending)); err == nil { - progress.setHorizon(horizon) + case repo.TerminalStatus() || repo.BackfillDone: + default: + plan.ready = append(plan.ready, sweepItem{DID: did, Lane: sweepLane(did, repo.PDS), Deepen: true}) + plan.floors[did] = backfillFloorTime(repo.BackfillFloor) } } - return nil + return plan, nil } -// 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) - if err != nil { - return nil, time.Time{}, fmt.Errorf("failed to get repo for %s: %w", did, err) - } - // No row, no completed sync, parked, or already complete: nothing to - // deepen. A repo the shallow phase failed on has no Version and is left - // alone here rather than fetched with the wrong ranges. - if repo == nil || repo.Version == "" || repo.TerminalStatus() || repo.BackfillDone { - continue - } - 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 { - floor = t - } - } - if floor.After(horizon) { - horizon = floor +// backfillFloorTime reads a backfill watermark as the instant it encodes. A repo +// with no watermark -- a row from a build that had never heard of one -- has had +// no history walked at all, so its floor is now: it holds the horizon at the +// present moment until its first window lands. +func backfillFloorTime(tid string) time.Time { + if tid != "" { + if t, err := reposync.TimeForTID(tid); err == nil { + return t } } - return pending, horizon, nil + return time.Now() } // sweepProgress is the state behind the sweep's status line. It is written by // every worker and read by the ticker, so everything goes through the mutex. type sweepProgress struct { - mu sync.Mutex - phase string - done int - windows int - round int - total int - horizon time.Time - started bool -} - -// begin starts a phase, resetting the completion counts. -func (p *sweepProgress) begin(phase string, total int, horizon time.Time) { + mu sync.Mutex + started bool + shallowTotal int + shallowDone int + deepenTotal int + deepenDone int + windows int + // floors is the watermark of every repo that is servable but not fully + // indexed, which is the set the horizon is the maximum over. Repos leave it + // when they finish; ones that failed stay, holding the horizon where they + // left it, because that is the truth about what this node can serve. + floors map[string]time.Time +} + +// begin starts a sweep with the work its plan found. +func (p *sweepProgress) begin(shallow int, floors map[string]time.Time) { p.mu.Lock() defer p.mu.Unlock() - p.phase = phase - p.total = total - p.done = 0 - p.windows = 0 - p.round = 0 - p.horizon = horizon p.started = true + p.shallowTotal = shallow + p.shallowDone = 0 + p.deepenTotal = len(floors) + p.deepenDone = 0 + p.windows = 0 + p.floors = make(map[string]time.Time, len(floors)) + for did, floor := range floors { + p.floors[did] = floor + } } -// finished records one repo completing the current phase. -func (p *sweepProgress) finished() { +// synced records one repo's shallow sync completing. +func (p *sweepProgress) synced() { p.mu.Lock() defer p.mu.Unlock() - p.done++ + p.shallowDone++ } -// window records one history window completing. During deepening a repo only -// counts as done at the bottom of its ladder, so without this the status line -// reads users=0 while thousands of windows finish underneath it. -func (p *sweepProgress) window() { +// laddered records a freshly synced repo joining its lane's ladder. The number +// of repos being deepened is not known when a sweep starts -- every shallow sync +// can add one -- so the denominator grows as the sweep discovers it. +func (p *sweepProgress) laddered(did string, floor time.Time) { p.mu.Lock() defer p.mu.Unlock() - p.windows++ + p.deepenTotal++ + p.floors[did] = floor } -// setRound records which deepening round is running. Rounds are barriers -- -// every repo gets its first window before any repo gets a second -- so repos -// cannot finish before the last round, and users=0 is the expected reading for -// most of the phase. The round is what makes that legible on the status line. -func (p *sweepProgress) setRound(round int) { +// window records one history window completing. A repo only counts as deepened +// at the bottom of its ladder, so without this the status line would read +// deepened=0 while tens of thousands of windows finished underneath it. +func (p *sweepProgress) window(did string, floor time.Time) { p.mu.Lock() defer p.mu.Unlock() - p.round = round + p.windows++ + if p.floors != nil { + p.floors[did] = floor + } } -// setHorizon updates how far back the sweep has taken every repo it is working -// on. -func (p *sweepProgress) setHorizon(horizon time.Time) { +// deepened records a repo reaching the start of its history. +func (p *sweepProgress) deepened(did string) { p.mu.Lock() defer p.mu.Unlock() - p.horizon = horizon + p.deepenDone++ + delete(p.floors, did) } -// status is the status line's key/value pairs. horizon is unix seconds: the -// instant after which every repo in this phase is fully indexed, so a number -// that climbs backwards through history as the sweep works. +// horizon is the most recent watermark among repos that are servable but not +// fully indexed: the instant after which everything this node serves is indexed. +// It climbs backwards through history as the sweep works. Zero when there is +// nothing left to deepen. +func (p *sweepProgress) horizon() int64 { + var newest time.Time + for _, floor := range p.floors { + if floor.After(newest) { + newest = floor + } + } + if newest.IsZero() { + return 0 + } + return newest.Unix() +} + +// status is the status line's key/value pairs: how many repos have been made +// servable, how many have their whole history, how many windows that took, and +// the horizon as unix seconds. func (p *sweepProgress) status() []any { p.mu.Lock() defer p.mu.Unlock() - horizon := int64(0) - if !p.horizon.IsZero() { - horizon = p.horizon.Unix() - } - kv := []any{"phase", p.phase, "users", p.done, "total", p.total, "horizon", horizon} - if p.phase == sweepPhaseDeepen { - kv = append(kv, "round", p.round, "windows", p.windows) + return []any{ + "shallow", fmt.Sprintf("%d/%d", p.shallowDone, p.shallowTotal), + "deepened", fmt.Sprintf("%d/%d", p.deepenDone, p.deepenTotal), + "windows", p.windows, + "horizon", p.horizon(), } - return kv } // start runs the status ticker until the returned function is called, which diff --git a/pkg/atproto/sweep_test.go b/pkg/atproto/sweep_test.go index 93529acc..1fdfae86 100644 --- a/pkg/atproto/sweep_test.go +++ b/pkg/atproto/sweep_test.go @@ -96,11 +96,25 @@ func TestBackfillWindowedHistory(t *testing.T) { 4, // [genesis, 180d) -- the two-hundred-day-old message } var done bool + var floor string for rung, want := range wantAfterRung { require.False(t, done, "the ladder finished early at rung %d", rung) - done, err = atsync.DeepenRepo(ctx, user.DID) + previous := floor + done, floor, err = atsync.DeepenRepo(ctx, user.DID) require.NoError(t, err, "rung %d", rung) require.Equal(t, want, countMessages(), "message count after rung %d", rung) + if !done { + // The floor a window reports is what the sweep's horizon is made + // of, so it has to be the watermark that was actually written, and + // it has to keep reaching further back. + stored, err := mod.GetRepo(user.DID) + require.NoError(t, err) + require.Equal(t, stored.BackfillFloor, floor, "rung %d reports the watermark it wrote", rung) + if previous != "" { + // TIDs sort by time, so each rung's watermark is smaller. + require.Less(t, floor, previous, "rung %d reaches further back than the last", rung) + } + } } require.True(t, done, "the last window bottoms out the collection") @@ -115,7 +129,7 @@ func TestBackfillWindowedHistory(t *testing.T) { // And a repo that is done is done: another sweep costs nothing and says // nothing. - again, err := atsync.DeepenRepo(ctx, user.DID) + again, _, err := atsync.DeepenRepo(ctx, user.DID) require.NoError(t, err) require.True(t, again) require.NoError(t, atsync.Sweep(ctx)) @@ -218,8 +232,8 @@ func TestSweepPrioritizesOwnDIDs(t *testing.T) { } // 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. +// are grouped by PDS host, however the host was written down, and a row with no +// host does not queue up behind the other rows that have none. 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")) @@ -238,17 +252,10 @@ func TestSweepHostLanes(t *testing.T) { {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)) + // own.example, a.example (three repos, one lane), b.example, and one lane + // each for the two rows that name no host. + require.Equal(t, 5, laneCount(items)) + require.Equal(t, 0, laneCount(nil)) } // TestSweepResolvesUnknownHosts: the sweep's DID list and the PDS column live in @@ -308,11 +315,11 @@ func TestLaneSchedulerStreams(t *testing.T) { firstStarted := make(chan struct{}) var once sync.Once - sched := newLaneScheduler(context.Background(), 2, func(_ context.Context, item sweepItem) { + sched := newLaneScheduler(context.Background(), 2, func(_ context.Context, step sweepStep) bool { once.Do(func() { close(firstStarted) }) mu.Lock() - inflight[item.Lane]++ - require.LessOrEqual(t, inflight[item.Lane], 1, "two workers on lane %s", item.Lane) + inflight[step.Lane]++ + require.LessOrEqual(t, inflight[step.Lane], 1, "two workers on lane %s", step.Lane) total := 0 for _, n := range inflight { total += n @@ -320,12 +327,13 @@ func TestLaneSchedulerStreams(t *testing.T) { if total > maxTotal { maxTotal = total } - order = append(order, item.DID) + order = append(order, step.DID) mu.Unlock() <-release mu.Lock() - inflight[item.Lane]-- + inflight[step.Lane]-- mu.Unlock() + return false }) sched.add(sweepItem{DID: "a1", Lane: "hostA"}) @@ -350,8 +358,9 @@ func TestLaneSchedulerStreams(t *testing.T) { func TestLaneSchedulerCancelled(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) cancel() - sched := newLaneScheduler(ctx, 2, func(context.Context, sweepItem) { + sched := newLaneScheduler(ctx, 2, func(context.Context, sweepStep) bool { t.Error("work ran under a cancelled context") + return false }) sched.add(sweepItem{DID: "a1", Lane: "hostA"}) _, err := sched.wait() @@ -367,6 +376,219 @@ func indexOf(xs []string, x string) int { return -1 } +// stepLabel renders a step the way the lane-program tests compare them: which +// repo, and whether it is the shallow sync or the nth window. +func stepLabel(step sweepStep) string { + if !step.Deepen { + return step.DID + "/shallow" + } + return fmt.Sprintf("%s/window%d", step.DID, step.Windows+1) +} + +// TestSweepLaneProgramShallowFirst: a host's repos are all made servable before +// any of them is deepened, and each repo's ladder starts at the bottom rung. A +// sweep that deepened one repo's history while another on the same host had +// never been read at all would be optimizing the wrong thing. +func TestSweepLaneProgramShallowFirst(t *testing.T) { + // Nothing runs until every repo is queued, so that this is a statement + // about the program and not about who won a race to be added. + ready := make(chan struct{}) + var mu sync.Mutex + var steps []string + windows := map[string]int{} + + sched := newLaneScheduler(context.Background(), 4, func(_ context.Context, step sweepStep) bool { + <-ready + mu.Lock() + defer mu.Unlock() + steps = append(steps, stepLabel(step)) + if !step.Deepen { + return true + } + windows[step.DID]++ + return windows[step.DID] < 2 + }) + for _, did := range []string{"a", "b", "c"} { + sched.add(sweepItem{DID: did, Lane: "pds.example"}) + } + close(ready) + lanes, err := sched.wait() + require.NoError(t, err) + require.Equal(t, 1, lanes) + + require.Equal(t, []string{ + "a/shallow", "b/shallow", "c/shallow", + "a/window1", "b/window1", "c/window1", + "a/window2", "b/window2", "c/window2", + }, steps) +} + +// TestSweepLaneProgramBreadthFirst is the guarantee the global rounds used to +// buy, rescoped to one host: no repo gets its (n+1)th window while another repo +// on the same host is still waiting for its nth. That is what puts the same +// horizon behind every account a PDS serves, and it is now free -- a lane +// reaching it does not make any other lane wait. +func TestSweepLaneProgramBreadthFirst(t *testing.T) { + want := map[string]int{"a": 2, "b": 5, "c": 3, "d": 5} + ready := make(chan struct{}) + var mu sync.Mutex + windows := map[string]int{} + pending := map[string]bool{} + for did := range want { + pending[did] = true + } + + sched := newLaneScheduler(context.Background(), 4, func(_ context.Context, step sweepStep) bool { + <-ready + mu.Lock() + defer mu.Unlock() + require.True(t, step.Deepen, "these repos are already servable") + require.Equal(t, windows[step.DID], step.Windows, + "%s: a step knows how many windows its repo has had", step.DID) + for did := range pending { + require.LessOrEqual(t, step.Windows, windows[did], + "%s took window %d while %s was still waiting for window %d", + step.DID, step.Windows+1, did, windows[did]+1) + } + windows[step.DID]++ + if windows[step.DID] >= want[step.DID] { + delete(pending, step.DID) + return false + } + return true + }) + for _, did := range []string{"a", "b", "c", "d"} { + sched.add(sweepItem{DID: did, Lane: "pds.example", Deepen: true}) + } + close(ready) + _, err := sched.wait() + require.NoError(t, err) + require.Equal(t, want, windows, "every repo got exactly the ladder it asked for") +} + +// TestSweepLaneProgramLateShallowPreempts: a repo whose host is resolved after +// its lane started work joins that lane mid-ladder, and is synced before the +// lane takes another rung -- an account nobody has read yet is worth more than +// another month of history for accounts that are already being served. It then +// joins the ladder at the bottom, so the breadth-first order absorbs it instead +// of leaving it a lap behind. +func TestSweepLaneProgramLateShallowPreempts(t *testing.T) { + var mu sync.Mutex + var steps []string + windows := map[string]int{} + var once sync.Once + var sched *laneScheduler + + sched = newLaneScheduler(context.Background(), 4, func(_ context.Context, step sweepStep) bool { + mu.Lock() + steps = append(steps, stepLabel(step)) + if step.Deepen { + windows[step.DID]++ + } + n := windows[step.DID] + mu.Unlock() + + if !step.Deepen { + return true // a fresh sync always leaves history to fetch here + } + if step.DID == "a" && n == 2 { + // The resolver finally placed a repo on this host, half way + // through the ladder the lane was already running. + once.Do(func() { sched.add(sweepItem{DID: "late", Lane: "pds.example"}) }) + } + if step.DID == "late" { + return n < 2 + } + return n < 4 + }) + sched.add(sweepItem{DID: "a", Lane: "pds.example", Deepen: true}) + sched.add(sweepItem{DID: "b", Lane: "pds.example", Deepen: true}) + _, err := sched.wait() + require.NoError(t, err) + + require.Equal(t, []string{ + "a/window1", "b/window1", "a/window2", + "late/shallow", // straight away, ahead of b's second window + "late/window1", // and its first rung before anyone's third + "b/window2", "late/window2", + "a/window3", "b/window3", + "a/window4", "b/window4", + }, steps) +} + +// TestSweepLaneProgramsAreIndependent is the whole point of this design: a host +// that is not answering cannot hold up a host that is. Measured on a 20k-repo +// sweep, the global phase and round barriers spent roughly half the wall clock +// with most lanes idle behind stragglers exactly like this one. +func TestSweepLaneProgramsAreIndependent(t *testing.T) { + hold := make(chan struct{}) + finished := make(chan struct{}) + var mu sync.Mutex + var fast []string + windows := map[string]int{} + const fastSteps = 8 // two repos, each a shallow sync and three windows + + sched := newLaneScheduler(context.Background(), 4, func(_ context.Context, step sweepStep) bool { + if step.Lane == "stuck.example" { + <-hold + return false + } + mu.Lock() + defer mu.Unlock() + fast = append(fast, stepLabel(step)) + if len(fast) == fastSteps { + close(finished) + } + if !step.Deepen { + return true + } + windows[step.DID]++ + return windows[step.DID] < 3 + }) + // The stuck host goes first, so it also holds the first slot: priority + // order must not become priority blocking. + sched.add(sweepItem{DID: "stuck1", Lane: "stuck.example"}) + sched.add(sweepItem{DID: "a", Lane: "pds.example"}) + sched.add(sweepItem{DID: "b", Lane: "pds.example"}) + + select { + case <-finished: + case <-time.After(30 * time.Second): + t.Fatal("the working host's lane never finished while another host was stuck") + } + mu.Lock() + require.Equal(t, []string{ + "a/shallow", "b/shallow", + "a/window1", "b/window1", + "a/window2", "b/window2", + "a/window3", "b/window3", + }, fast, "a whole per-host program ran to the end with another host mid-sync") + mu.Unlock() + + close(hold) + lanes, err := sched.wait() + require.NoError(t, err) + require.Equal(t, 2, lanes) +} + +// TestSweepLaneProgramSpinGuard: a repo that never admits to being finished +// still costs a bounded number of windows per sweep. +func TestSweepLaneProgramSpinGuard(t *testing.T) { + var mu sync.Mutex + steps := 0 + sched := newLaneScheduler(context.Background(), 2, func(_ context.Context, step sweepStep) bool { + mu.Lock() + defer mu.Unlock() + steps++ + require.LessOrEqual(t, step.Windows, maxDeepenRounds) + return true + }) + sched.add(sweepItem{DID: "a", Lane: "pds.example"}) + _, err := sched.wait() + require.NoError(t, err) + require.Equal(t, maxDeepenRounds+1, steps, "one shallow sync and a bounded ladder") +} + // 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 @@ -383,12 +605,13 @@ func TestSweepLanesNeverShareAHost(t *testing.T) { var mu sync.Mutex active := map[string]string{} // lane -> the DID holding it var order []string + windows := map[string]int{} inFlight, maxInFlight := 0, 0 - err := runLanes(context.Background(), cap, hostLanes(items), func(ctx context.Context, item sweepItem) { + sched := newLaneScheduler(context.Background(), cap, func(_ context.Context, step sweepStep) bool { 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 + holder, busy := active[step.Lane] + require.False(t, busy, "%s and %s ran on %s at once", step.DID, holder, step.Lane) + active[step.Lane] = step.DID inFlight++ maxInFlight = max(maxInFlight, inFlight) mu.Unlock() @@ -398,66 +621,77 @@ func TestSweepLanesNeverShareAHost(t *testing.T) { time.Sleep(2 * time.Millisecond) mu.Lock() - delete(active, item.Lane) + defer mu.Unlock() + delete(active, step.Lane) inFlight-- - order = append(order, item.DID) - mu.Unlock() + order = append(order, step.DID) + if !step.Deepen { + return true + } + windows[step.DID]++ + return windows[step.DID] < 2 }) + for _, item := range items { + sched.add(item) + } + lanes, err := sched.wait() require.NoError(t, err) - require.Len(t, order, len(items), "every repo ran exactly once") + require.Equal(t, 4, lanes, "four hosts, and a cap of three: some lane waited for a slot") + require.Len(t, order, 3*len(items), "every repo got its sync and both its windows") 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. +// their lane is the first one given a slot -- the priority order prioritizeDIDs +// produces has to survive the bucketing, which means slots go out in the order +// lanes were added rather than in whatever order their goroutines woke up. 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) { + // One slot, and every repo on its own host, so lane order is the only + // thing deciding. + sched := newLaneScheduler(context.Background(), 1, func(_ context.Context, step sweepStep) bool { mu.Lock() defer mu.Unlock() - order = append(order, item.DID) - })) + order = append(order, step.DID) + return false + }) + for _, did := range dids { + sched.add(sweepItem{DID: did, Lane: sweepLane(did, "https://"+did+".pds.example")}) + } + _, err := sched.wait() + require.NoError(t, err) 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. +// checks the context between steps rather than at the end of a program that +// would otherwise run for hours. 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) { + sched := newLaneScheduler(ctx, 4, func(_ context.Context, step sweepStep) bool { mu.Lock() + defer mu.Unlock() ran++ if ran == 2 { cancel() } - mu.Unlock() + // Never finished: only the cancellation can end this lane. + return true }) + for i := 0; i < 40; i++ { + sched.add(sweepItem{DID: fmt.Sprintf("did:plc:%d", i), Lane: "pds.example"}) + } + _, err := sched.wait() 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") + require.Less(t, ran, 40, "the lane stopped instead of running its program out") } // TestSweepConcurrencyFlag: the cap comes from --sweep-concurrency, and an unset @@ -473,53 +707,60 @@ func TestSweepConcurrencyFlag(t *testing.T) { (&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. +// TestSweepProgressStatusLine covers the one line an operator watches. There +// are no phases left to name -- every host runs its own program -- so the line +// is two fractions, the windows they took, and the horizon in unix seconds: +// +// backfill sweep shallow=19000/20747 deepened=4300/20013 windows=41022 horizon=1753142400 func TestSweepProgressStatusLine(t *testing.T) { var progress sweepProgress // Before anything starts there is nothing to say. - require.Equal(t, []any{"phase", "", "users", 0, "total", 0, "horizon", int64(0)}, progress.status()) + require.Equal(t, + []any{"shallow", "0/0", "deepened", "0/0", "windows", 0, "horizon", int64(0)}, + progress.status()) - horizon := time.Now().Add(-InitialWindow) - progress.begin(sweepPhaseShallow, 3, horizon) - progress.finished() + day := time.Now().Add(-InitialWindow) + week := time.Now().Add(-7 * 24 * time.Hour) + month := time.Now().Add(-30 * 24 * time.Hour) + + // Three repos to make servable, one already servable and mid-ladder: the + // horizon is that one's watermark. + progress.begin(3, map[string]time.Time{"did:plc:old": week}) + require.Equal(t, + []any{"shallow", "0/3", "deepened", "0/1", "windows", 0, "horizon", week.Unix()}, + progress.status()) + + // A repo that has just been synced is servable, and joins the ladder: the + // denominator grows as the sweep discovers who needs deepening, and the + // horizon follows the least-deepened repo. + progress.synced() + progress.laddered("did:plc:new", day) + require.Equal(t, + []any{"shallow", "1/3", "deepened", "0/2", "windows", 0, "horizon", day.Unix()}, + progress.status()) + + // Windows count as they land, and each moves one repo's watermark. The + // horizon only moves when the laggard does. + progress.window("did:plc:new", week) + require.Equal(t, + []any{"shallow", "1/3", "deepened", "0/2", "windows", 1, "horizon", week.Unix()}, + progress.status()) + progress.window("did:plc:old", month) require.Equal(t, - []any{"phase", "shallow", "users", 1, "total", 3, "horizon", horizon.Unix()}, + []any{"shallow", "1/3", "deepened", "0/2", "windows", 2, "horizon", week.Unix()}, progress.status()) - // A new phase resets the counts and moves the horizon. Deepening also - // reports windows: repos only count as done at the bottom of their ladder, - // so windows is the number that shows the sweep moving in the meantime. - deeper := time.Now().Add(-30 * 24 * time.Hour) - progress.begin(sweepPhaseDeepen, 2, deeper) + // A repo with its whole history stops holding the horizon back, and when + // nothing is left to deepen there is no horizon at all. + progress.window("did:plc:new", month) + progress.deepened("did:plc:new") require.Equal(t, - []any{"phase", "deepen", "users", 0, "total", 2, "horizon", deeper.Unix(), "round", 0, "windows", 0}, + []any{"shallow", "1/3", "deepened", "1/2", "windows", 3, "horizon", month.Unix()}, progress.status()) - progress.setRound(1) - progress.window() - progress.window() - progress.window() - progress.finished() - progress.finished() - deepest := time.Now().Add(-180 * 24 * time.Hour) - progress.setHorizon(deepest) + progress.deepened("did:plc:old") require.Equal(t, - []any{"phase", "deepen", "users", 2, "total", 2, "horizon", deepest.Unix(), "round", 1, "windows", 3}, + []any{"shallow", "1/3", "deepened", "2/2", "windows", 3, "horizon", int64(0)}, progress.status()) // The ticker stops when told to, without leaking a goroutine. diff --git a/pkg/reposync/retry.go b/pkg/reposync/retry.go index ddf45a51..18924f2a 100644 --- a/pkg/reposync/retry.go +++ b/pkg/reposync/retry.go @@ -32,6 +32,11 @@ const ( // long sleep here stalls every other repo on that host, and a repo whose // backfill fails is simply retried later. DefaultRetryMaxDelay = 30 * time.Second + // deadHostAttempts is all a host that is not there gets: see [isDeadHost]. + // One retry, because a PDS that is restarting refuses connections for a + // second or two and that is worth waiting out; not five, because nothing + // else is. + deadHostAttempts = 2 ) // RetryPolicy bounds how hard a fetcher retries a transient XRPC failure. @@ -133,7 +138,11 @@ func (p RetryPolicy) do(ctx context.Context, what string, fn func() error) error if !isRetryable(err) { return err } - if attempt >= p.MaxAttempts { + budget := p.MaxAttempts + if isDeadHost(err) && deadHostAttempts < budget { + budget = deadHostAttempts + } + if attempt >= budget { return fmt.Errorf("giving up after %d attempts: %w", attempt, err) } d, source := p.delay(attempt, err) @@ -243,6 +252,30 @@ func isRetryable(err error) bool { errors.Is(err, io.EOF) } +// isDeadHost reports whether err says the host is not there at all, rather than +// busy, broken, or slow: nothing accepted the connection, or the name does not +// resolve. +// +// These get [deadHostAttempts] tries instead of the full ladder. A sweep of +// twenty thousand repos meets a long tail of PDSes that have been switched off, +// and every repo on one of them was costing five attempts and a minute of +// backoff to learn what the first attempt already said. That tail is most of +// what a sweep's stragglers are made of. +// +// The whole retry ladder is for hosts that might answer if asked again -- +// timeouts, 429s, 5xx -- and a refused connection or a missing DNS record is +// not that. It is checked with errors.Is/As rather than on the surface error +// because the real thing arrives wrapped several deep: net/http returns a +// *url.Error around a *net.OpError around the syscall or *net.DNSError, and +// indigo's xrpc wraps that again. +func isDeadHost(err error) bool { + if errors.Is(err, syscall.ECONNREFUSED) { + return true + } + var derr *net.DNSError + return errors.As(err, &derr) && derr.IsNotFound +} + // 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 diff --git a/pkg/reposync/retry_test.go b/pkg/reposync/retry_test.go index 083ee0f8..31048d4c 100644 --- a/pkg/reposync/retry_test.go +++ b/pkg/reposync/retry_test.go @@ -5,7 +5,9 @@ import ( "errors" "fmt" "io" + "net" "net/http" + "net/url" "strconv" "syscall" "testing" @@ -236,6 +238,116 @@ func TestRetryDelay(t *testing.T) { }) } +// nxdomain is a name that no longer resolves, wrapped the way it arrives: the +// resolver's error inside net/http's dial error inside net/http's request +// error. +func nxdomain() error { + return fmt.Errorf("getBlocks: %w", &url.Error{ + Op: "Get", + URL: "https://gone.example/xrpc/com.atproto.sync.getBlocks", + Err: &net.OpError{Op: "dial", Net: "tcp", Err: &net.DNSError{ + Err: "no such host", Name: "gone.example", IsNotFound: true, + }}, + }) +} + +// TestRetryFastFailsDeadHosts: a host that is not there at all gets two +// attempts, not five. This is what makes a sweep's straggler tail cheap -- +// switched-off PDSes were costing a full backoff ladder per repo to rediscover +// something the first connection attempt already reported. +func TestRetryFastFailsDeadHosts(t *testing.T) { + policy := RetryPolicy{MaxAttempts: 5, BaseDelay: time.Millisecond, MaxDelay: 2 * time.Millisecond} + attempts := func(t *testing.T, err error) int { + t.Helper() + calls := 0 + got := policy.do(context.Background(), "getBlocks", func() error { + calls++ + return err + }) + require.Error(t, got) + return calls + } + + t.Run("connection refused", func(t *testing.T) { + err := fmt.Errorf("request failed: %w", &url.Error{Op: "Get", URL: "https://pds.example/", + Err: &net.OpError{Op: "dial", Err: syscall.ECONNREFUSED}}) + require.True(t, isDeadHost(err)) + // Still retryable: a PDS that is restarting refuses for a moment. + require.True(t, isRetryable(err)) + require.Equal(t, deadHostAttempts, attempts(t, err)) + }) + + t.Run("no such host", func(t *testing.T) { + err := nxdomain() + require.True(t, isDeadHost(err)) + // A name that does not resolve now will not resolve in a second, so + // this never even reaches the two-attempt cap: it is not retryable at + // all, and one attempt is what it costs. + require.False(t, isRetryable(err)) + require.Equal(t, 1, attempts(t, err)) + }) + + t.Run("a timeout still gets the whole ladder", func(t *testing.T) { + // The host is there and answering slowly, which is exactly what the + // retries are for. + err := fmt.Errorf("request failed: %w", timeoutError{}) + require.False(t, isDeadHost(err)) + require.Equal(t, 5, attempts(t, err)) + }) + + t.Run("a DNS timeout is not a dead host", func(t *testing.T) { + // The resolver is struggling, not answering "no": that is transient. + err := fmt.Errorf("request failed: %w", &net.OpError{Op: "dial", Err: &net.DNSError{ + Err: "i/o timeout", Name: "pds.example", IsTimeout: true, + }}) + require.False(t, isDeadHost(err)) + require.Equal(t, 5, attempts(t, err)) + }) + + t.Run("429 still gets the whole ladder", func(t *testing.T) { + err := ratelimited(time.Time{}) + require.False(t, isDeadHost(err)) + require.Equal(t, 5, attempts(t, err)) + }) + + t.Run("503 still gets the whole ladder", func(t *testing.T) { + require.Equal(t, 5, attempts(t, xrpcErr(http.StatusServiceUnavailable, "", "restarting"))) + }) + + t.Run("a policy that asks for less keeps it", func(t *testing.T) { + single := RetryPolicy{MaxAttempts: 1, BaseDelay: time.Millisecond} + calls := 0 + err := single.do(context.Background(), "getBlocks", func() error { + calls++ + return fmt.Errorf("dialing: %w", syscall.ECONNREFUSED) + }) + require.Error(t, err) + require.Equal(t, 1, calls) + }) +} + +// TestRetryDeadHostOffTheWire: the classification above is only worth anything +// if a refused connection still looks like one after net/http and indigo have +// each wrapped it, so this one dials a port that nothing is listening on. +func TestRetryDeadHostOffTheWire(t *testing.T) { + ln, err := net.Listen("tcp", "127.0.0.1:0") + require.NoError(t, err) + addr := ln.Addr().String() + require.NoError(t, ln.Close()) + + sr := buildSignedRepo(t, testDID, exactnessPaths()) + f := &XRPCBlockFetcher{ + Client: &xrpc.Client{Host: "http://" + addr}, + DID: testDID, + Retry: RetryPolicy{MaxAttempts: 5, BaseDelay: time.Millisecond, MaxDelay: 2 * time.Millisecond}, + } + _, err = f.GetBlocks(context.Background(), []cid.Cid{sr.root}) + require.Error(t, err) + require.True(t, isRetryable(err), "a refused connection is worth one retry: %v", err) + require.True(t, isDeadHost(err), "but it must be recognisable as a dead host: %v", err) + require.Contains(t, err.Error(), fmt.Sprintf("giving up after %d attempts", deadHostAttempts)) +} + func ratelimited(reset time.Time) error { return fmt.Errorf("getBlocks: %w", &xrpc.Error{ StatusCode: http.StatusTooManyRequests, -- 2.51.2