diff --git a/pkg/reposync/walk.go b/pkg/reposync/walk.go index 69c8f2e04..b37bd42c1 100644 --- a/pkg/reposync/walk.go +++ b/pkg/reposync/walk.go @@ -221,7 +221,9 @@ type Walker struct { BatchSize int // Checkpoint, if set, is called with the frontier after every completed // step (that is, after that step's records have been emitted). Returning an - // error aborts the walk. + // error aborts the walk and rolls the frontier back to before the step, so + // a retried Resume re-emits the step's records; this makes it safe to + // commit visitor effects and the frontier together inside Checkpoint. Checkpoint func(*Frontier) error } @@ -248,7 +250,8 @@ func (w *Walker) WalkRanges(ctx context.Context, root cid.Cid, ranges []KeyRange // Resume continues a walk from a checkpointed frontier. fr is updated in place // as the walk progresses, so an aborted Resume leaves fr at the last completed -// step and can be called again. +// step — where a step only counts as completed once its Checkpoint call (if +// any) has succeeded — and can be called again. func (w *Walker) Resume(ctx context.Context, fr *Frontier, visit RecordVisitor) error { if fr == nil { return errors.New("nil frontier") @@ -344,9 +347,15 @@ func (w *Walker) step(ctx context.Context, fr *Frontier, visit RecordVisitor) er } } + prev := fr.Pending fr.Pending = next if w.Checkpoint != nil { if err := w.Checkpoint(fr); err != nil { + // Roll back so a retried Resume re-emits this step's records. A + // caller that commits visitor effects inside Checkpoint would + // otherwise lose them: the failed commit discards the effects and + // the advanced frontier would never emit those records again. + fr.Pending = prev return fmt.Errorf("checkpointing frontier: %w", err) } } diff --git a/pkg/reposync/walk_test.go b/pkg/reposync/walk_test.go index 0afca8735..b5c4168e6 100644 --- a/pkg/reposync/walk_test.go +++ b/pkg/reposync/walk_test.go @@ -474,6 +474,68 @@ func TestWalkResume(t *testing.T) { } } +// A failed checkpoint must leave the frontier at the last successfully +// checkpointed step. The caller modeled here commits visitor effects inside +// Checkpoint and loses the staged batch when the commit fails, so if the +// frontier stayed advanced, the failed step's records would never be emitted +// again and the final index would be incomplete. +func TestWalkFailedCheckpointRollsBack(t *testing.T) { + ctx := context.Background() + paths := exactnessPaths() + for i := 0; i < 300; i++ { + paths = append(paths, fmt.Sprintf("place.stream.media.origin/3lbmedia%06d", i)) + } + tr := buildRepo(t, paths) + want := expectedInRange(paths, "place.stream.") + + // Transactional caller: the visitor stages records, Checkpoint commits the + // stage together with the frontier. One commit fails, discarding its stage + // the way a rolled-back transaction would. + durable := map[string]cid.Cid{} + var staged []emission + errBoom := errors.New("simulated checkpoint failure") + failed := false + w := &Walker{ + Fetcher: newTestFetcher(tr), + BatchSize: 4, + Checkpoint: func(fr *Frontier) error { + if !failed && len(staged) > 0 { + failed = true + staged = nil + return errBoom + } + for _, e := range staged { + durable[e.path] = e.cid + } + staged = nil + return nil + }, + } + + fr := &Frontier{ + Root: tr.root, + Ranges: []KeyRange{PrefixRange("place.stream.")}, + Pending: []pendingEntry{{CID: tr.root}}, + } + err := w.Resume(ctx, fr, collectVisitor(&staged)) + require.ErrorIs(t, err, errBoom) + require.True(t, failed, "no checkpoint call ever had staged records") + require.False(t, fr.Done()) + + require.NoError(t, w.Resume(ctx, fr, collectVisitor(&staged))) + require.True(t, fr.Done()) + + got := make([]string, 0, len(durable)) + for p := range durable { + got = append(got, p) + } + sort.Strings(got) + require.Equal(t, want, got) + for _, p := range want { + require.Equal(t, tr.records[p], durable[p]) + } +} + // A warm cache makes a repeat walk entirely local. func TestWalkCachedFetcherWarmCacheDoesNoRemoteWork(t *testing.T) { ctx := context.Background()