diff --git a/pkg/atproto/atproto.go b/pkg/atproto/atproto.go index ce2151fb..07d5f07a 100644 --- a/pkg/atproto/atproto.go +++ b/pkg/atproto/atproto.go @@ -102,7 +102,9 @@ func (atsync *ATProtoSynchronizer) SyncBlueskyRepo(ctx context.Context, handle s log.Log(ctx, "discovered new user", "did", ident.DID.String(), "handle", ident.Handle.String(), "pds", ident.PDSEndpoint()) } - log.Log(ctx, "resolved bluesky identity", "did", ident.DID, "handle", ident.Handle, "pds", ident.PDSEndpoint()) + // Debug: this fires on every sync operation, not just first contact -- + // "discovered new user" above is the first-contact line. + log.Debug(ctx, "resolved atproto identity", "did", ident.DID, "handle", ident.Handle, "pds", ident.PDSEndpoint()) xrpcc := xrpc.Client{ Host: ident.PDSEndpoint(), Client: SyncHTTPClient, @@ -219,7 +221,9 @@ func (atsync *ATProtoSynchronizer) DeepenRepo(ctx context.Context, did string) ( 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) } - log.Log(ctx, "deepened repo history", "rev", rev, "floor", window.Lo, "done", window.Genesis) + // 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 } diff --git a/pkg/atproto/sweep.go b/pkg/atproto/sweep.go index efa909f5..5a5c4b36 100644 --- a/pkg/atproto/sweep.go +++ b/pkg/atproto/sweep.go @@ -528,6 +528,10 @@ func (atsync *ATProtoSynchronizer) sweepDeepen(ctx context.Context, progress *sw 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 @@ -540,8 +544,14 @@ func (atsync *ATProtoSynchronizer) sweepDeepen(ctx context.Context, progress *sw log.Error(ctx, "failed to deepen repo history", "did", item.DID, "err", err) return } + 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 } mu.Lock() @@ -598,18 +608,20 @@ type sweepProgress struct { mu sync.Mutex phase string done int + windows int total int horizon time.Time started bool } -// begin starts a phase, resetting the completion count. +// begin starts a phase, resetting the completion counts. func (p *sweepProgress) begin(phase string, total int, horizon time.Time) { p.mu.Lock() defer p.mu.Unlock() p.phase = phase p.total = total p.done = 0 + p.windows = 0 p.horizon = horizon p.started = true } @@ -621,6 +633,15 @@ func (p *sweepProgress) finished() { p.done++ } +// 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() { + p.mu.Lock() + defer p.mu.Unlock() + p.windows++ +} + // setHorizon updates how far back the sweep has taken every repo it is working // on. func (p *sweepProgress) setHorizon(horizon time.Time) { @@ -639,7 +660,11 @@ func (p *sweepProgress) status() []any { if !p.horizon.IsZero() { horizon = p.horizon.Unix() } - return []any{"phase", p.phase, "users", p.done, "total", p.total, "horizon", horizon} + kv := []any{"phase", p.phase, "users", p.done, "total", p.total, "horizon", horizon} + if p.phase == sweepPhaseDeepen { + kv = append(kv, "windows", p.windows) + } + 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 fef8db25..9c5f3e51 100644 --- a/pkg/atproto/sweep_test.go +++ b/pkg/atproto/sweep_test.go @@ -502,18 +502,23 @@ func TestSweepProgressStatusLine(t *testing.T) { []any{"phase", "shallow", "users", 1, "total", 3, "horizon", horizon.Unix()}, progress.status()) - // A new phase resets the count and moves the horizon. + // 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) require.Equal(t, - []any{"phase", "deepen", "users", 0, "total", 2, "horizon", deeper.Unix()}, + []any{"phase", "deepen", "users", 0, "total", 2, "horizon", deeper.Unix(), "windows", 0}, progress.status()) + progress.window() + progress.window() + progress.window() progress.finished() progress.finished() deepest := time.Now().Add(-180 * 24 * time.Hour) progress.setHorizon(deepest) require.Equal(t, - []any{"phase", "deepen", "users", 2, "total", 2, "horizon", deepest.Unix()}, + []any{"phase", "deepen", "users", 2, "total", 2, "horizon", deepest.Unix(), "windows", 3}, progress.status()) // The ticker stops when told to, without leaking a goroutine.