diff --git a/pkg/atproto/sweep.go b/pkg/atproto/sweep.go index bad2c25c..efa909f5 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 42fc01e4..fef8db25 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 ce42b308..429ec82d 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) } }()