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,