From a65b301ebee329eeacb3d8c75aa67340d72c9f1b Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Wed, 29 Jul 2026 17:02:32 -0700 Subject: [PATCH] atproto: stream lane scheduling instead of resolving every host up front The shallow sweep resolved the PDS host of every unknown repo before starting a single walk -- an 8-wide barrier over the full account list, each lookup including handle verification with multi-second DNS timeouts. On a fresh index (the exact case host sharding was built for) that meant minutes of dead air before any syncing: 20k repos at 8-wide is a quarter hour of nothing but identity resolution. resolveLanes becomes feedUnresolved + laneScheduler. The scheduler keeps runLanes' guarantees -- one worker per lane, global cap, FIFO slot order so own-DID lanes still go first -- but accepts items while running. Repos whose row already names a host start walking immediately; the resolver (now 16-wide) hands the rest over one at a time as answers land. The deepen phase keeps the static runLanes: after shallow, rows have hosts. Also fixes a pre-existing data race the new test surfaced: the chat message handler's fire-and-forget streamer sync assigned the enclosing function's err from its goroutine, racing every later use of err in the handler and able to mask or fabricate its error results. It gets its own variable. Co-Authored-By: Claude Fable 5 --- pkg/atproto/sweep.go | 202 +++++++++++++++++++++++++++++--------- pkg/atproto/sweep_test.go | 105 +++++++++++++++++--- pkg/atproto/sync.go | 5 +- 3 files changed, 251 insertions(+), 61 deletions(-) diff --git a/pkg/atproto/sweep.go b/pkg/atproto/sweep.go index bad2c25c7..efa909f54 100644 --- a/pkg/atproto/sweep.go +++ b/pkg/atproto/sweep.go @@ -46,8 +46,8 @@ type sweepItem struct { } // sweepLane is the lane a repo row belongs in: its PDS host, or -- for a repo -// whose host is not known even after [ATProtoSynchronizer.resolveLanes] tried to -// find out -- a lane of its own. +// whose host is not known even after [ATProtoSynchronizer.feedUnresolved] tried +// to find out -- a lane of its own. // // A lane to itself, rather than a shared catch-all: an unplaceable repo is // normally a resolution failure, so its sync is about to fail too, and queueing @@ -63,18 +63,21 @@ func sweepLane(did, pds string) string { return "did:" + did } -// identityResolveConcurrency is how many identities [ATProtoSynchronizer.resolveLanes] +// identityResolveConcurrency is how many identities [ATProtoSynchronizer.feedUnresolved] // looks up at once. // // Deliberately smaller than the sweep's own concurrency: these are lookups // against a handful of shared identity services (plc.directory, DNS) rather than // against thousands of PDSes, and the whole point is to move work that the // backfill would have done anyway, not to arrive at plc.directory with a -// thundering herd. -const identityResolveConcurrency = 8 - -// resolveLanes gives a lane to every item that has not got one, by resolving the -// repo's identity to find its PDS. +// thundering herd. On a fresh index every repo needs one of these, so this is +// also the ceiling on how fast a fresh node discovers work -- which is why it +// feeds a running [laneScheduler] instead of gating the sweep behind a barrier. +const identityResolveConcurrency = 16 + +// feedUnresolved lanes every item that has not got one, by resolving the repo's +// identity to find its PDS, handing each item to add as its answer lands. It +// returns when every item has been handed over. // // This exists because the sweep's DID list and the PDS column come from // different databases. The DIDs are the state database's set of "repos this node @@ -86,56 +89,153 @@ const identityResolveConcurrency = 8 // put every repo in a lane of its own: sharding by host would do nothing on // precisely the sweep it was built for. // +// Streaming, not a barrier: a fresh node has tens of thousands of these lookups +// to do, and doing them all before the first walk turned the start of a sweep +// into minutes of dead air. Feeding a running [laneScheduler] means the first +// repos are being walked while the last are still being resolved. +// // The lookup is moved rather than added. It goes through the same cached // directory [ATProtoSynchronizer.SyncBlueskyRepo] resolves with, so the backfill -// a few seconds later reads this answer out of the cache instead of asking -// again. +// moments later reads this answer out of the cache instead of asking again. // -// Failures are not fatal and are not even logged loudly: the repo keeps a lane +// Failures are not fatal and are not even logged loudly: the repo gets a lane // of its own and its sync fails on its own terms, one repo at a time, the way it // did before. -func (atsync *ATProtoSynchronizer) resolveLanes(ctx context.Context, items []sweepItem) { - var todo []int - for i := range items { - if items[i].Lane == "" { - todo = append(todo, i) - } - } - if len(todo) == 0 { +func (atsync *ATProtoSynchronizer) feedUnresolved(ctx context.Context, items []sweepItem, add func(sweepItem)) { + if len(items) == 0 { return } - log.Log(ctx, "resolving PDS hosts to shard the sweep", "repos", len(todo)) + log.Log(ctx, "resolving PDS hosts to shard the sweep", "repos", len(items)) start := time.Now() var resolved atomic.Int64 g, gctx := errgroup.WithContext(ctx) g.SetLimit(identityResolveConcurrency) - for _, i := range todo { + for _, item := range items { g.Go(func() error { - if err := gctx.Err(); err != nil { - return err + // On cancellation, still hand the item over (with a lane of its + // own): the scheduler's workers notice the dead context themselves, + // and every item accounted for exactly once is the simpler + // invariant to keep. + if gctx.Err() == nil { + if ident, err := atsync.resolveIdent(gctx, item.DID, true); err == nil { + item.Lane = reposync.HostKey(ident.PDSEndpoint()) + resolved.Add(1) + } else { + log.Debug(gctx, "could not resolve a repo's PDS for sharding", "did", item.DID, "err", err) + } } - ident, err := atsync.resolveIdent(gctx, items[i].DID, true) - if err != nil { - log.Debug(gctx, "could not resolve a repo's PDS for sharding", "did", items[i].DID, "err", err) - return nil + if item.Lane == "" { + item.Lane = sweepLane(item.DID, "") } - items[i].Lane = reposync.HostKey(ident.PDSEndpoint()) - resolved.Add(1) + add(item) return nil }) } - // A cancelled context is the only error this can produce, and the caller is - // about to notice it for itself. _ = g.Wait() + log.Log(ctx, "resolved PDS hosts to shard the sweep", "repos", len(items), + "resolved", resolved.Load(), "took", time.Since(start)) +} - for _, i := range todo { - if items[i].Lane == "" { - items[i].Lane = sweepLane(items[i].DID, "") +// 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. +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{} + + mu sync.Mutex + queue map[string][]sweepItem + live map[string]bool + seen int + wg sync.WaitGroup +} + +func newLaneScheduler(ctx context.Context, limit int, work func(context.Context, sweepItem)) *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{}, + } +} + +// 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. +func (s *laneScheduler) add(item sweepItem) { + s.mu.Lock() + if _, ok := s.queue[item.Lane]; !ok { + s.seen++ + } + s.queue[item.Lane] = append(s.queue[item.Lane], item) + spawn := !s.live[item.Lane] + if spawn { + s.live[item.Lane] = true + s.wg.Add(1) + } + s.mu.Unlock() + if spawn { + go s.run(item.Lane) + } +} + +// 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) { + defer s.wg.Done() + select { + case s.sem <- struct{}{}: + case <-s.ctx.Done(): + s.abandon(lane) + return + } + defer func() { <-s.sem }() + for { + if s.ctx.Err() != nil { + s.abandon(lane) + return } + s.mu.Lock() + if len(s.queue[lane]) == 0 { + s.live[lane] = 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) } - log.Log(ctx, "resolved PDS hosts to shard the sweep", "repos", len(todo), - "resolved", resolved.Load(), "took", time.Since(start)) +} + +// abandon drops a lane's remaining items on cancellation. The repos keep their +// rows untouched, so the next sweep picks them up. +func (s *laneScheduler) abandon(lane string) { + s.mu.Lock() + s.live[lane] = false + s.queue[lane] = nil + s.mu.Unlock() +} + +// 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. +func (s *laneScheduler) wait() (lanes int, err error) { + s.wg.Wait() + s.mu.Lock() + defer s.mu.Unlock() + return s.seen, s.ctx.Err() } // hostLanes groups items into one lane per [sweepItem.Lane], keeping each lane's @@ -352,23 +452,29 @@ func (atsync *ATProtoSynchronizer) sweepShallow(ctx context.Context, progress *s if repo != nil { pds = repo.PDS } - // Left empty when the row does not name a host: resolveLanes below fills - // those in rather than letting them each become a lane. + // 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 } - atsync.resolveLanes(ctx, todo) - if err := ctx.Err(); err != nil { - return err + + var known, unresolved []sweepItem + for _, item := range todo { + if item.Lane == "" { + unresolved = append(unresolved, item) + } else { + known = append(known, item) + } } - lanes := hostLanes(todo) - log.Log(ctx, "syncing repos", "phase", sweepPhaseShallow, "repos", len(todo), "hosts", len(lanes)) + log.Log(ctx, "syncing repos", "phase", sweepPhaseShallow, "repos", len(todo), + "knownHosts", len(hostLanes(known)), "unresolved", len(unresolved)) var failed atomic.Int64 - err := runLanes(ctx, atsync.sweepConcurrency(), lanes, func(ctx context.Context, item sweepItem) { + 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) @@ -376,9 +482,17 @@ func (atsync *ATProtoSynchronizer) sweepShallow(ctx context.Context, progress *s } 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)) } diff --git a/pkg/atproto/sweep_test.go b/pkg/atproto/sweep_test.go index 42fc01e48..fef8db258 100644 --- a/pkg/atproto/sweep_test.go +++ b/pkg/atproto/sweep_test.go @@ -270,26 +270,101 @@ func TestSweepResolvesUnknownHosts(t *testing.T) { // Pre-set so resolveIdent never reaches for a real directory. atsync := &ATProtoSynchronizer{PLCDirectory: &dir, CachedPLCDirectory: &dir} - items := []sweepItem{ - {DID: "did:plc:known", Lane: sweepLane("did:plc:known", "https://known.example")}, + var mu sync.Mutex + lanes := map[string]string{} + atsync.feedUnresolved(context.Background(), []sweepItem{ {DID: "did:plc:one"}, {DID: "did:plc:two"}, {DID: "did:plc:missing"}, // no DID document: nothing to place it by {DID: "did:plc:three"}, - } - atsync.resolveLanes(context.Background(), items) + }, func(item sweepItem) { + mu.Lock() + lanes[item.DID] = item.Lane + mu.Unlock() + }) - require.Equal(t, [][]string{ - {"did:plc:known"}, // a row that named its host is left alone - {"did:plc:one", "did:plc:two"}, // and two resolving to one host share a lane - {"did:plc:missing"}, // unplaceable: its own lane, not a queue - {"did:plc:three"}, - }, laneDIDs(hostLanes(items))) - - // Nothing to do is the normal case, and it must not cost a lookup. - placed := []sweepItem{{DID: "did:plc:known", Lane: "known.example"}} - atsync.resolveLanes(context.Background(), placed) - require.Equal(t, "known.example", placed[0].Lane) + require.Equal(t, map[string]string{ + "did:plc:one": "shared.example", // two resolving to one host share a lane + "did:plc:two": "shared.example", + "did:plc:missing": "did:did:plc:missing", // unplaceable: its own lane, not a queue + "did:plc:three": "elsewhere.example", + }, lanes) + + // Nothing to do is the normal case, and it must not add anything. + atsync.feedUnresolved(context.Background(), nil, func(sweepItem) { + t.Error("add called with no items to resolve") + }) +} + +// TestLaneSchedulerStreams is the property the scheduler exists for: work on +// lanes that are already known starts while more items are still arriving, +// without ever breaking one-worker-per-lane or the global cap. +func TestLaneSchedulerStreams(t *testing.T) { + var mu sync.Mutex + inflight := map[string]int{} + maxTotal := 0 + var order []string + release := make(chan struct{}) + firstStarted := make(chan struct{}) + var once sync.Once + + sched := newLaneScheduler(context.Background(), 2, func(_ context.Context, item sweepItem) { + once.Do(func() { close(firstStarted) }) + mu.Lock() + inflight[item.Lane]++ + require.LessOrEqual(t, inflight[item.Lane], 1, "two workers on lane %s", item.Lane) + total := 0 + for _, n := range inflight { + total += n + } + if total > maxTotal { + maxTotal = total + } + order = append(order, item.DID) + mu.Unlock() + <-release + mu.Lock() + inflight[item.Lane]-- + mu.Unlock() + }) + + sched.add(sweepItem{DID: "a1", Lane: "hostA"}) + // The first item is being worked before the rest have even been added -- + // that is the streaming property. + <-firstStarted + sched.add(sweepItem{DID: "a2", Lane: "hostA"}) + sched.add(sweepItem{DID: "b1", Lane: "hostB"}) + sched.add(sweepItem{DID: "c1", Lane: "hostC"}) + close(release) + + lanes, err := sched.wait() + require.NoError(t, err) + require.Equal(t, 3, lanes) + require.ElementsMatch(t, []string{"a1", "a2", "b1", "c1"}, order) + require.LessOrEqual(t, maxTotal, 2, "global lane cap exceeded") + require.Less(t, indexOf(order, "a1"), indexOf(order, "a2"), "lane order must be FIFO") +} + +// TestLaneSchedulerCancelled: a dead context stops a scheduler without working +// anything more and without hanging wait. +func TestLaneSchedulerCancelled(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + cancel() + sched := newLaneScheduler(ctx, 2, func(context.Context, sweepItem) { + t.Error("work ran under a cancelled context") + }) + sched.add(sweepItem{DID: "a1", Lane: "hostA"}) + _, err := sched.wait() + require.ErrorIs(t, err, context.Canceled) +} + +func indexOf(xs []string, x string) int { + for i, v := range xs { + if v == x { + return i + } + } + return -1 } // TestSweepLanesNeverShareAHost is the property the lanes exist for: a sweep diff --git a/pkg/atproto/sync.go b/pkg/atproto/sync.go index ce42b308c..429ec82dc 100644 --- a/pkg/atproto/sync.go +++ b/pkg/atproto/sync.go @@ -113,8 +113,9 @@ func (atsync *ATProtoSynchronizer) handleCreateUpdate(ctx context.Context, userD } go func() { - _, err = atsync.SyncBlueskyRepoCached(ctx, rec.Streamer) - if err != nil { + // Its own err on purpose: assigning the enclosing function's err + // from this goroutine races every later use of it. + if _, err := atsync.SyncBlueskyRepoCached(ctx, rec.Streamer); err != nil { log.Error(ctx, "failed to sync bluesky repo", "err", err) } }() -- 2.51.2