diff --git a/pkg/atproto/atproto.go b/pkg/atproto/atproto.go index 82aaced57..4d203e70e 100644 --- a/pkg/atproto/atproto.go +++ b/pkg/atproto/atproto.go @@ -17,6 +17,7 @@ import ( "stream.place/streamplace/pkg/comatproto" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/model" + "stream.place/streamplace/pkg/reposync" ) var SyncGetRepo = comatproto.SyncGetRepo @@ -109,7 +110,11 @@ func (atsync *ATProtoSynchronizer) SyncBlueskyRepo(ctx context.Context, handle s return nil, fmt.Errorf("no PDS endpoint found for Bluesky identity %s", handle) } - rev, rootCID, err := atsync.backfillRepo(ctx, ident, &xrpcc) + // First contact is shallow: everything this node indexes, but only the last + // [InitialWindow] of the collections that can hold years of records. The + // account is servable in seconds; the sweep deepens its history afterwards. + floor := reposync.TIDForTime(time.Now().Add(-InitialWindow)) + result, err := atsync.backfillRepo(ctx, ident, &xrpcc, floor) if err != nil { if parked := parkTerminalRepo(ctx, mod, ident.DID.String(), err); parked != nil { return nil, parked @@ -120,12 +125,14 @@ func (atsync *ATProtoSynchronizer) SyncBlueskyRepo(ctx context.Context, handle s // A completed backfill proves the account is fine, so Status goes back to // empty -- UpdateRepo writes every column, so this happens by construction. newRepo := model.Repo{ - DID: ident.DID.String(), - PDS: ident.PDSEndpoint(), - Version: rev, - RootCID: rootCID, - Handle: ident.Handle.String(), - Status: model.RepoStatusOK, + DID: ident.DID.String(), + PDS: ident.PDSEndpoint(), + Version: result.Rev, + RootCID: result.RootCID, + Handle: ident.Handle.String(), + Status: model.RepoStatusOK, + BackfillFloor: result.Floor, + BackfillDone: result.Done, } err = mod.UpdateRepo(&newRepo) if err != nil { @@ -139,6 +146,82 @@ func (atsync *ATProtoSynchronizer) SyncBlueskyRepo(ctx context.Context, handle s return &newRepo, nil } +// DeepenRepo walks one more window of history for a repo whose recent records +// are already indexed, and reports whether that repo is now complete. +// +// 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. +// +// 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) { + repo, err := atsync.Model.GetRepo(did) + if err != nil { + 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) + case repo.TerminalStatus(): + // The account is gone; whatever we indexed is all there will be. + return true, nil + case repo.BackfillDone: + return true, nil + case repo.Version == "": + // Never synced (or wedged): that is the shallow phase'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) + } + + // The same lock a full sync takes, so the two cannot walk one repo at once. + // Nothing re-enters it: indexing a record calls SyncBlueskyRepoCached, which + // short-circuits on the Version this row already has. + handleLock := handleLocks.GetLock(did) + handleLock.Lock() + defer handleLock.Unlock() + + ident, err := atsync.resolveIdent(ctx, did, true) + if err != nil { + return false, fmt.Errorf("failed to resolve %s: %w", did, err) + } + xrpcc := xrpc.Client{Host: ident.PDSEndpoint(), Client: &aqhttp.Client} + if xrpcc.Host == "" { + return false, fmt.Errorf("no PDS endpoint found for %s", did) + } + + window := nextBackfillWindow(repo.BackfillFloor, time.Now()) + ctx = log.WithLogValues(ctx, "did", did) + log.Debug(ctx, "walking a history window", "floor", repo.BackfillFloor, "to", window.Lo, "genesis", window.Genesis) + + rev, root, err := atsync.walkBackfill(ctx, ident, &xrpcc, windowRanges(window.Lo, window.Hi)) + if err != nil && isMethodNotSupported(err) { + // No windowed walk to be had from this host. The full-CAR fallback reads + // the entire repo, so one of those finishes the job for good. + log.Warn(ctx, "host does not support sync.getBlocks, deepening with a full getRepo", + "pds", xrpcc.Host, "err", err) + // The legacy path has no verified MST root to record. + root = "" + rev, err = atsync.legacyBackfill(ctx, ident, &xrpcc) + window = backfillWindow{Genesis: true} + } + if err != nil { + if parked := parkTerminalRepo(ctx, atsync.Model, did, err); parked != nil { + return false, parked + } + 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) + } + log.Log(ctx, "deepened repo history", "rev", rev, "floor", window.Lo, "done", window.Genesis) + return window.Genesis, nil +} + // syncsInFlight holds the DIDs whose backfill is running in this process right // now. A placeholder repo row (empty Version) otherwise means "incomplete, // re-sync me", which would be wrong -- and, since indexing a record can call @@ -195,6 +278,10 @@ func (atsync *ATProtoSynchronizer) RefreshIdentity(ctx context.Context, did stri // account came back, and blanking it here would put every deactivated // repo back in the boot-time sync sweep. newRepo.Status = oldRepo.Status + // And for the backfill watermark: losing it would make the sweep walk + // this repo's whole history again from the top of the ladder. + newRepo.BackfillFloor = oldRepo.BackfillFloor + newRepo.BackfillDone = oldRepo.BackfillDone } err = atsync.Model.UpdateRepo(&newRepo) if err != nil { diff --git a/pkg/atproto/backfill_walk.go b/pkg/atproto/backfill_walk.go index e49eb77f0..cc17e1348 100644 --- a/pkg/atproto/backfill_walk.go +++ b/pkg/atproto/backfill_walk.go @@ -15,6 +15,7 @@ import ( "github.com/bluesky-social/indigo/repo" "github.com/bluesky-social/indigo/xrpc" "github.com/ipfs/go-cid" + "stream.place/streamplace/pkg/constants" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/model" "stream.place/streamplace/pkg/reposync" @@ -25,13 +26,41 @@ import ( // records live in one contiguous key range. const placeStreamPrefix = "place.stream." +// windowedCollections are the collections a backfill reads by time window +// instead of all at once. +// +// Their rkeys are TIDs, so their keys sort chronologically and "the last day of +// chat" is a key range (see [reposync.TIDForTime]). They are also the two +// collections whose volume decides how long a first sync takes: an account with +// years of Bluesky posts and streamplace chat has tens of thousands of records +// there and a few dozen everywhere else. +// +// Windowing changes when a record is indexed, never whether: the ladder of +// windows in [nextBackfillWindow] bottoms out at the start of the collection, +// and until it does the repo row says so. A record whose rkey is not a TID is +// not skipped either -- the windows are key ranges, so it simply arrives with +// whichever window its rkey sorts into. +// +// Every entry must be a collection [backfillRanges] would otherwise walk whole +// -- either under place.stream. or in [CollectionFilter]. TestBackfillRanges +// checks that. +var windowedCollections = []string{ + constants.PLACE_STREAM_CHAT_MESSAGE, + constants.APP_BSKY_FEED_POST, +} + // backfillRanges is the set of MST key ranges a backfill walks: everything // under place.stream., plus one range per non-streamplace collection the // firehose accepts. // // It is derived from CollectionFilter at runtime so the backfill and the // firehose can never drift apart about which records this node indexes. -func backfillRanges() []reposync.KeyRange { +// +// floor is a TID watermark for the windowed collections: those are walked only +// from floor forward, with the rest of their key space cut out of the ranges +// that would otherwise cover it. An empty floor means "from the beginning", +// which is every range whole. +func backfillRanges(floor string) []reposync.KeyRange { ranges := []reposync.KeyRange{reposync.PrefixRange(placeStreamPrefix)} for _, nsid := range CollectionFilter { if strings.HasPrefix(nsid, placeStreamPrefix) { @@ -41,47 +70,188 @@ func backfillRanges() []reposync.KeyRange { // "app.bsky.feed.postgate". ranges = append(ranges, reposync.PrefixRange(nsid+"/")) } + if floor == "" { + return ranges + } + for _, nsid := range windowedCollections { + whole := reposync.PrefixRange(nsid + "/") + kept := make([]reposync.KeyRange, 0, len(ranges)+1) + for _, r := range ranges { + kept = append(kept, subtractRange(r, whole)...) + } + ranges = append(kept, reposync.KeyRange{Lo: []byte(nsid + "/" + floor), Hi: whole.Hi}) + } return ranges } -// backfillRepo indexes every record we care about from ident's repo. It returns -// the repo revision the index is now consistent with, and the MST root CID that -// revision committed to (empty if the fallback path was used, which never sees -// a verified root). +// windowRanges covers [lo, hi) of every windowed collection and nothing else. +// It is what a deepening step walks. An empty lo starts at the first key of +// each collection; an empty hi runs to the last. +func windowRanges(lo, hi string) []reposync.KeyRange { + out := make([]reposync.KeyRange, 0, len(windowedCollections)) + for _, nsid := range windowedCollections { + r := reposync.PrefixRange(nsid + "/") + if lo != "" { + r.Lo = []byte(nsid + "/" + lo) + } + if hi != "" { + r.Hi = []byte(nsid + "/" + hi) + } + out = append(out, r) + } + return out +} + +// subtractRange returns r with cut removed: r itself when they do not overlap, +// the pieces of r on either side of cut when they do, and nothing when cut +// swallows r. Empty pieces are dropped, because a zero-width range is not a +// range the walker will accept. +func subtractRange(r, cut reposync.KeyRange) []reposync.KeyRange { + if !rangesOverlap(r, cut) { + return []reposync.KeyRange{r} + } + var out []reposync.KeyRange + if bytes.Compare(r.Lo, cut.Lo) < 0 { + out = append(out, reposync.KeyRange{Lo: r.Lo, Hi: cut.Lo}) + } + if cut.Hi != nil && (r.Hi == nil || bytes.Compare(cut.Hi, r.Hi) < 0) { + out = append(out, reposync.KeyRange{Lo: cut.Hi, Hi: r.Hi}) + } + return out +} + +func rangesOverlap(a, b reposync.KeyRange) bool { + if a.Hi != nil && bytes.Compare(b.Lo, a.Hi) >= 0 { + return false + } + if b.Hi != nil && bytes.Compare(a.Lo, b.Hi) >= 0 { + return false + } + return true +} + +// InitialWindow is how much history a first sync reads from the windowed +// collections. It is the whole cost difference between meeting an account and +// serving it: everything else in the repo is configuration-sized. +const InitialWindow = 24 * time.Hour + +// backfillSpans is the ladder of windows the deepening sweep walks, each one +// reaching further back than the last. After the last rung the next window runs +// to the start of the collection, and the repo is complete. +// +// Spans are wall-clock ages, not window widths: a repo at the 7d rung has +// everything from seven days ago forward, and its next window is +// [30d ago, 7d ago). +var backfillSpans = []time.Duration{ + InitialWindow, + 7 * 24 * time.Hour, + 30 * 24 * time.Hour, + 180 * 24 * time.Hour, +} + +// backfillWindow is one step of the ladder: the slice of the windowed +// collections a deepening walk should read next. +type backfillWindow struct { + // Lo and Hi are TID bounds, empty meaning the start/end of the collection. + Lo string + Hi string + // Genesis reports that this window reaches the start of the collection, so + // a repo that finishes it has no history left to fetch. + Genesis bool + // Horizon is the wall-clock instant Lo encodes, for logging. Zero for a + // genesis window, which has no horizon. + Horizon time.Time +} + +// nextBackfillWindow picks the next deepening window for a repo whose windowed +// collections are synced from floor forward. +// +// Which rung of the ladder a repo is on is read off the age of its floor rather +// than stored: a floor is a timestamp, and "how far back does this repo go" is +// the only thing that matters. That keeps the watermark a single self-describing +// column, and makes a row written by an older build (or a hand-edited one) land +// on a sensible rung by itself. +// +// An empty floor means nothing has been recorded, which is both a repo that has +// never been synced and every row written before this code existed: those start +// at the top of the ladder, with a window that is open-ended above. +func nextBackfillWindow(floor string, now time.Time) backfillWindow { + window := func(span time.Duration, hi string) backfillWindow { + horizon := now.Add(-span) + return backfillWindow{Lo: reposync.TIDForTime(horizon), Hi: hi, Horizon: horizon} + } + if floor == "" { + return window(backfillSpans[0], "") + } + floorTime, err := reposync.TimeForTID(floor) + if err != nil { + // Not a timestamp we can place on the ladder. One genesis-bounded + // window finishes the repo off rather than looping on it forever. + return backfillWindow{Hi: floor, Genesis: true} + } + age := now.Sub(floorTime) + for _, span := range backfillSpans { + if span > age { + return window(span, floor) + } + } + return backfillWindow{Hi: floor, Genesis: true} +} + +// backfillResult is what a completed backfill knows about the repo it read. +type backfillResult struct { + // Rev is the repo revision the index is now consistent with. + Rev string + // RootCID is the MST root that revision committed to, empty if the fallback + // path was used, which never sees a verified root. + RootCID string + // Floor is the TID watermark for the windowed collections: their history is + // synced from here forward. Empty means from the start of the collection. + Floor string + // Done reports that the windowed collections need no further deepening. + Done bool +} + +// backfillRepo indexes every record we care about from ident's repo. +// +// floor windows the high-volume collections: only their history from that TID +// forward is read, which is what makes first contact with a busy account cost +// seconds instead of minutes. An empty floor reads everything. // // The fast path walks only the subtrees holding records we index. Hosts that do // not implement com.atproto.sync.getBlocks fall back to downloading the whole -// repo as a CAR. -func (atsync *ATProtoSynchronizer) backfillRepo(ctx context.Context, ident *identity.Identity, xrpcc *xrpc.Client) (string, string, error) { - rev, root, err := atsync.walkBackfill(ctx, ident, xrpcc) +// repo as a CAR -- which reads all of it, window or no window, so such a repo +// comes back complete. +func (atsync *ATProtoSynchronizer) backfillRepo(ctx context.Context, ident *identity.Identity, xrpcc *xrpc.Client, floor string) (backfillResult, error) { + rev, root, err := atsync.walkBackfill(ctx, ident, xrpcc, backfillRanges(floor)) if err == nil { - return rev, root, nil + return backfillResult{Rev: rev, RootCID: root, Floor: floor, Done: floor == ""}, nil } if isStaleWalkError(err) { // walkBackfill already exhausted its restart-from-a-new-head budget // on this. isMethodNotSupported is written not to claim these either, // but say it once here rather than depend on that ordering: answering // "the repo moved" with a full getRepo download would be absurd. - return "", "", err + return backfillResult{}, err } if !isMethodNotSupported(err) { // Anything else -- a bad signature, a malformed tree, a network // failure -- must propagate. Falling back on a verification failure // would make the verification decorative. - return "", "", err + return backfillResult{}, err } log.Warn(ctx, "host does not support sync.getBlocks, falling back to full getRepo", "pds", xrpcc.Host, "did", ident.DID.String(), "err", err) rev, err = atsync.legacyBackfill(ctx, ident, xrpcc) if err != nil { - return "", "", err + return backfillResult{}, err } - return rev, "", nil + return backfillResult{Rev: rev, Done: true}, nil } -// walkBackfill does a verified, prefix-bounded walk of the remote repo, handing +// walkBackfill does a verified, range-bounded walk of the remote repo, handing // every record in range to the same indexing path the firehose uses. -func (atsync *ATProtoSynchronizer) walkBackfill(ctx context.Context, ident *identity.Identity, xrpcc *xrpc.Client) (string, string, error) { +func (atsync *ATProtoSynchronizer) walkBackfill(ctx context.Context, ident *identity.Identity, xrpcc *xrpc.Client, ranges []reposync.KeyRange) (string, string, error) { did := ident.DID.String() dir := atsync.PLCDirectory @@ -116,7 +286,7 @@ func (atsync *ATProtoSynchronizer) walkBackfill(ctx context.Context, ident *iden walk := func(ctx context.Context, root cid.Cid) error { records = 0 walker := &reposync.Walker{Fetcher: fetcher} - err := walker.WalkRanges(ctx, root, backfillRanges(), func(path string, rcid cid.Cid, rec []byte) error { + err := walker.WalkRanges(ctx, root, ranges, func(path string, rcid cid.Cid, rec []byte) error { nsid, rkey, err := syntax.ParseRepoPath(path) if err != nil { log.Warn(ctx, "failed to parse repo path", "k", path, "err", err) diff --git a/pkg/atproto/backfill_walk_test.go b/pkg/atproto/backfill_walk_test.go index 9d3962f79..142f7643a 100644 --- a/pkg/atproto/backfill_walk_test.go +++ b/pkg/atproto/backfill_walk_test.go @@ -80,7 +80,7 @@ func TestBackfillWalk(t *testing.T) { } var got []string walker := &reposync.Walker{Fetcher: fetcher} - err = walker.WalkRanges(ctx, h.Root, backfillRanges(), func(path string, rcid cid.Cid, rec []byte) error { + err = walker.WalkRanges(ctx, h.Root, backfillRanges(""), func(path string, rcid cid.Cid, rec []byte) error { got = append(got, path) return nil }) @@ -236,10 +236,10 @@ func TestBackfillFallsBackToGetRepo(t *testing.T) { ident, err := atsync.resolveIdent(ctx, user.DID, false) require.NoError(t, err) - var rev, root string + var result backfillResult err = untilNoErrors(t, func() error { var err error - rev, root, err = atsync.backfillRepo(ctx, ident, &xrpc.Client{Host: proxy.URL, Client: &aqhttp.Client}) + result, err = atsync.backfillRepo(ctx, ident, &xrpc.Client{Host: proxy.URL, Client: &aqhttp.Client}, reposync.TIDForTime(time.Now().Add(-InitialWindow))) if err != nil { return err } @@ -253,8 +253,10 @@ func TestBackfillFallsBackToGetRepo(t *testing.T) { return nil }) require.NoError(t, err, "backfill should have fallen back to getRepo") - require.NotEmpty(t, rev, "the legacy path still reports the commit rev") - require.Empty(t, root, "the legacy path has no verified MST root to record") + require.NotEmpty(t, result.Rev, "the legacy path still reports the commit rev") + require.Empty(t, result.RootCID, "the legacy path has no verified MST root to record") + require.Empty(t, result.Floor, "a full CAR download ignores the window") + require.True(t, result.Done, "a full CAR download leaves no history to deepen") } func TestIsMethodNotSupported(t *testing.T) { @@ -554,8 +556,23 @@ func testHead(t *testing.T, rev string) *reposync.Head { return &reposync.Head{Rev: rev, Root: root} } +// inRanges is the walker's own containment test, spelled out here so these +// tests check the ranges rather than trusting the code that builds them. +func inRanges(ranges []reposync.KeyRange, key string) bool { + for _, r := range ranges { + if r.Lo != nil && key < string(r.Lo) { + continue + } + if r.Hi != nil && key >= string(r.Hi) { + continue + } + return true + } + return false +} + func TestBackfillRanges(t *testing.T) { - ranges := backfillRanges() + ranges := backfillRanges("") // One for place.stream., plus one per non-streamplace collection the // firehose accepts. want := 1 @@ -566,18 +583,7 @@ func TestBackfillRanges(t *testing.T) { } require.Len(t, ranges, want) - inRange := func(key string) bool { - for _, r := range ranges { - if r.Lo != nil && key < string(r.Lo) { - continue - } - if r.Hi != nil && key >= string(r.Hi) { - continue - } - return true - } - return false - } + inRange := func(key string) bool { return inRanges(ranges, key) } require.True(t, inRange("place.stream.chat.message/3l")) require.True(t, inRange("place.stream.live.recommendations/self")) require.True(t, inRange("app.bsky.actor.profile/self")) @@ -588,6 +594,158 @@ func TestBackfillRanges(t *testing.T) { require.False(t, inRange("place.strea.thing/3l")) require.False(t, inRange("place.streamx.thing/3l")) require.False(t, inRange("zzz.example.thing/3l")) + + // Every windowed collection has to be one this node walks in the first + // place, or windowing it would widen the sweep instead of narrowing it. + for _, nsid := range windowedCollections { + require.True(t, inRange(nsid+"/3l"), "windowed collection %s is not in the full ranges", nsid) + } +} + +// TestBackfillRangesWindowed: with a floor, the high-volume collections are cut +// down to their recent history and everything else stays whole -- including the +// records that sort on either side of the hole cut out of place.stream. +func TestBackfillRangesWindowed(t *testing.T) { + floor := reposync.TIDForTime(time.Now().Add(-24 * time.Hour)) + older := reposync.TIDForTime(time.Now().Add(-48 * time.Hour)) + newer := reposync.TIDForTime(time.Now().Add(-time.Hour)) + ranges := backfillRanges(floor) + inRange := func(key string) bool { return inRanges(ranges, key) } + + // The windowed collections keep only what is at or after the floor. + require.True(t, inRange("place.stream.chat.message/"+floor)) + require.True(t, inRange("place.stream.chat.message/"+newer)) + require.False(t, inRange("place.stream.chat.message/"+older)) + require.True(t, inRange("app.bsky.feed.post/"+newer)) + require.False(t, inRange("app.bsky.feed.post/"+older)) + // A windowed collection is still walked exhaustively, just not all at once: + // a non-TID rkey lands in whichever window its bytes fall into ("self" + // sorts above every TID minted this century, so it comes in the first one) + // and the windows together cover the whole collection either way. + require.True(t, inRange("place.stream.chat.message/self")) + require.False(t, inRange("place.stream.chat.message/!oldest")) + + // Everything else is untouched, on both sides of the hole in place.stream. + require.True(t, inRange("place.stream.chat.gate/3l")) // sorts before chat.message + require.True(t, inRange("place.stream.chat.profile/self")) // sorts after + require.True(t, inRange("place.stream.live.livestream/3l")) // sorts after + require.True(t, inRange("app.bsky.actor.profile/self")) + require.True(t, inRange("app.bsky.graph.follow/3l")) + // And the ranges are still bounded where they were. + require.False(t, inRange("app.bsky.feed.postgate/3l")) + require.False(t, inRange("app.bsky.actor.status/3l")) + require.False(t, inRange("zzz.example.thing/3l")) + + // The walker rejects an inverted or empty range outright, and windowing is + // the only thing in here that builds a range out of two different strings. + for _, r := range ranges { + require.NotNil(t, r.Hi, "no unbounded range should come out of here: %s", r) + require.Less(t, string(r.Lo), string(r.Hi), "inverted range %s", r) + } + // Cutting two collections down to a window turns two whole ranges into a + // bounded piece each, plus the pieces left on either side of the hole in + // place.stream. + require.Len(t, ranges, len(backfillRanges(""))+len(windowedCollections)) +} + +// TestWindowRanges: a deepening step reads the windowed collections and nothing +// else. +func TestWindowRanges(t *testing.T) { + lo := reposync.TIDForTime(time.Now().Add(-7 * 24 * time.Hour)) + hi := reposync.TIDForTime(time.Now().Add(-24 * time.Hour)) + ranges := windowRanges(lo, hi) + require.Len(t, ranges, len(windowedCollections)) + inRange := func(key string) bool { return inRanges(ranges, key) } + + mid := reposync.TIDForTime(time.Now().Add(-3 * 24 * time.Hour)) + require.True(t, inRange("place.stream.chat.message/"+mid)) + require.True(t, inRange("app.bsky.feed.post/"+mid)) + require.False(t, inRange("place.stream.chat.message/"+hi), "the floor is exclusive above") + require.False(t, inRange("place.stream.chat.message/"+reposync.TIDForTime(time.Now().Add(-30*24*time.Hour)))) + // Nothing outside the windowed collections is read again. + require.False(t, inRange("place.stream.chat.profile/self")) + require.False(t, inRange("app.bsky.actor.profile/self")) + + // The genesis window: open below, so it sweeps up everything left, + // including rkeys that are not TIDs at all. + last := windowRanges("", hi) + require.True(t, inRanges(last, "place.stream.chat.message/!oldest")) + require.True(t, inRanges(last, "place.stream.chat.message/"+reposync.TIDForTime(time.Unix(0, 0)))) + require.False(t, inRanges(last, "place.stream.chat.message/"+hi)) + require.False(t, inRanges(last, "place.stream.chat.gate/3l")) +} + +// TestNextBackfillWindow walks the ladder a repo climbs down, checking that +// each window abuts the last one (no gap, so no record can be skipped) and that +// it terminates. +func TestNextBackfillWindow(t *testing.T) { + now := time.Now() + + // Nothing recorded: start at the top, open-ended above. + first := nextBackfillWindow("", now) + require.Equal(t, reposync.TIDForTime(now.Add(-InitialWindow)), first.Lo) + require.Empty(t, first.Hi, "a first window has no upper bound") + require.False(t, first.Genesis) + + // Then each rung, with the floor aging as if the previous window had just + // finished. Every window starts where the last one ended. + floor := first.Lo + var spans []time.Duration + for range len(backfillSpans) + 4 { + win := nextBackfillWindow(floor, now) + require.Equal(t, floor, win.Hi, "a window must pick up exactly where the last one stopped") + if win.Genesis { + require.Empty(t, win.Lo, "the last window runs to the start of the collection") + break + } + require.Less(t, win.Lo, win.Hi, "windows must not be inverted") + spans = append(spans, now.Sub(win.Horizon).Round(time.Hour)) + floor = win.Lo + } + require.Equal(t, []time.Duration{ + 7 * 24 * time.Hour, + 30 * 24 * time.Hour, + 180 * 24 * time.Hour, + }, spans, "the ladder after the initial window") + require.True(t, nextBackfillWindow(floor, now).Genesis, "the ladder terminates") + + // A floor much older than the whole ladder goes straight to the end. + ancient := reposync.TIDForTime(now.Add(-5 * 365 * 24 * time.Hour)) + require.True(t, nextBackfillWindow(ancient, now).Genesis) + + // A floor from the future (clock skew, a hand-edited row) still produces a + // usable window rather than an inverted one. + future := reposync.TIDForTime(now.Add(time.Hour)) + skewed := nextBackfillWindow(future, now) + require.Less(t, skewed.Lo, skewed.Hi) + + // A watermark that is not a TID at all: one final window, and done. + require.True(t, nextBackfillWindow("not-a-tid", now).Genesis) +} + +func TestSubtractRange(t *testing.T) { + r := func(lo, hi string) reposync.KeyRange { + out := reposync.KeyRange{Lo: []byte(lo)} + if hi != "" { + out.Hi = []byte(hi) + } + return out + } + str := func(ranges []reposync.KeyRange) []string { + out := []string{} + for _, x := range ranges { + out = append(out, string(x.Lo)+".."+string(x.Hi)) + } + return out + } + + require.Equal(t, []string{"a..b"}, str(subtractRange(r("a", "b"), r("c", "d"))), "disjoint") + require.Equal(t, []string{"a..b"}, str(subtractRange(r("a", "b"), r("b", "d"))), "abutting") + require.Equal(t, []string{"a..c", "d..z"}, str(subtractRange(r("a", "z"), r("c", "d"))), "a hole") + require.Equal(t, []string{"d..z"}, str(subtractRange(r("a", "z"), r("a", "d"))), "cut off the front") + require.Equal(t, []string{"a..d"}, str(subtractRange(r("a", "z"), r("d", "z"))), "cut off the back") + require.Empty(t, subtractRange(r("a", "z"), r("a", "z")), "cut swallows the range") + require.Empty(t, subtractRange(r("b", "c"), r("a", "z")), "cut swallows the range") } func backfillTestSynchronizer(t *testing.T, dev *devenv.DevEnv) (*ATProtoSynchronizer, model.Model) { diff --git a/pkg/atproto/migrate.go b/pkg/atproto/migrate.go deleted file mode 100644 index 89df958ec..000000000 --- a/pkg/atproto/migrate.go +++ /dev/null @@ -1,111 +0,0 @@ -package atproto - -import ( - "context" - "fmt" - "sync" - "sync/atomic" - "time" - - "golang.org/x/sync/errgroup" - "stream.place/streamplace/pkg/log" -) - -func (atsync *ATProtoSynchronizer) Migrate(ctx context.Context) error { - // Accounts that are deactivated, deleted, or taken down fail their backfill - // the same way on every boot forever. One query up front keeps them out of - // the sweep entirely, instead of one logged failure each. - terminalDIDs, err := atsync.Model.TerminalRepoDIDs(ctx) - if err != nil { - return fmt.Errorf("failed to list repos in terminal states: %w", err) - } - terminal := make(map[string]struct{}, len(terminalDIDs)) - for _, did := range terminalDIDs { - terminal[did] = struct{}{} - } - - var allDIDs []string - skipped := 0 - offset := 0 - for { - repos, err := atsync.StatefulDB.ListRepos(100, offset) - if err != nil { - return err - } - if len(repos) == 0 { - break - } - for _, repo := range repos { - if _, ok := terminal[repo.DID]; ok { - skipped++ - continue - } - allDIDs = append(allDIDs, repo.DID) - } - offset += len(repos) - } - - if skipped > 0 { - log.Log(ctx, "skipping repos with terminal status", "skipped", skipped) - } - log.Log(ctx, "starting migration sync", "totalRepos", len(allDIDs)) - - g, ctx := errgroup.WithContext(ctx) - var syncedCount int64 - - syncErrors := map[string]error{} - syncErrorMu := sync.Mutex{} - - // Start progress logging goroutine - progressCtx, cancelProgress := context.WithCancel(ctx) - defer cancelProgress() - - go func() { - ticker := time.NewTicker(10 * time.Second) - defer ticker.Stop() - - for { - select { - case <-progressCtx.Done(): - return - case <-ticker.C: - current := atomic.LoadInt64(&syncedCount) - log.Log(ctx, "migration progress", "synced", current, "total", len(allDIDs)) - } - } - }() - - for i, did := range allDIDs { - currentIndex := i - currentDID := did - g.Go(func() error { - log.Debug(ctx, "syncing repo", "did", currentDID, "progress", currentIndex+1, "total", len(allDIDs)) - _, err := atsync.SyncBlueskyRepoCached(ctx, currentDID) - if err != nil { - log.Error(ctx, "failed to sync repo", "did", currentDID, "err", err) - syncErrorMu.Lock() - syncErrors[currentDID] = err - syncErrorMu.Unlock() - } else { - atomic.AddInt64(&syncedCount, 1) - } - return nil - }) - } - - if err := g.Wait(); err != nil { - log.Error(ctx, "migration failed", "err", err, "synced", atomic.LoadInt64(&syncedCount), "total", len(allDIDs)) - return err - } - - for did, err := range syncErrors { - log.Error(ctx, "migration failed for user", "did", did, "err", err) - } - - if len(allDIDs) > 0 && len(syncErrors) == len(allDIDs) { - return fmt.Errorf("all users failed to migrate") - } - - log.Log(ctx, "migration completed", "synced", len(allDIDs)) - return nil -} diff --git a/pkg/atproto/redelivery_test.go b/pkg/atproto/redelivery_test.go index fb639fa8a..c80460ccb 100644 --- a/pkg/atproto/redelivery_test.go +++ b/pkg/atproto/redelivery_test.go @@ -239,10 +239,10 @@ func TestMigrateSkipsTerminalRepos(t *testing.T) { require.NoError(t, atsync.StatefulDB.AddRepo(did)) // Control: unparked, this DID resolves nowhere, so the sweep fails on it. - require.Error(t, atsync.Migrate(ctx), "the only repo in the sweep should have failed") + require.Error(t, atsync.Sweep(ctx), "the only repo in the sweep should have failed") require.NoError(t, mod.SetRepoStatus(ctx, did, model.RepoStatusDeactivated)) - require.NoError(t, atsync.Migrate(ctx), "a terminal repo should never be dialed") + require.NoError(t, atsync.Sweep(ctx), "a terminal repo should never be dialed") } // TestReviveRepo: a commit proves the account is back. diff --git a/pkg/atproto/sweep.go b/pkg/atproto/sweep.go new file mode 100644 index 000000000..b690d9051 --- /dev/null +++ b/pkg/atproto/sweep.go @@ -0,0 +1,387 @@ +package atproto + +import ( + "context" + "fmt" + "sort" + "sync" + "time" + + "golang.org/x/sync/errgroup" + "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/reposync" +) + +const ( + // sweepConcurrency bounds how many repos a sweep works on at once. A boot + // used to start one goroutine per known repo, which meant a fresh node + // opened with a thundering herd at every PDS it had ever heard of. + sweepConcurrency = 6 + + // 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" +) + +// 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. +var maxDeepenRounds = len(backfillSpans) + 3 + +// 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. +// +// 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. +// +// 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 +// where every shallow sync failed returns an error, because that is a broken +// node rather than a few broken accounts. +func (atsync *ATProtoSynchronizer) Sweep(ctx context.Context) error { + dids, err := atsync.sweepCandidates(ctx) + if err != nil { + return err + } + log.Log(ctx, "starting backfill sweep", "totalRepos", len(dids)) + + progress := &sweepProgress{} + stop := progress.start(ctx) + defer stop() + + if err := atsync.sweepShallow(ctx, progress, dids); err != nil { + return err + } + if err := ctx.Err(); err != nil { + return err + } + if err := atsync.sweepDeepen(ctx, progress, dids); err != nil { + return err + } + log.Log(ctx, "backfill sweep complete", "totalRepos", len(dids)) + return nil +} + +// 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 + // the same way on every boot forever. One query up front keeps them out of + // the sweep entirely, instead of one logged failure each. + terminalDIDs, err := atsync.Model.TerminalRepoDIDs(ctx) + if err != nil { + return nil, fmt.Errorf("failed to list repos in terminal states: %w", err) + } + terminal := make(map[string]struct{}, len(terminalDIDs)) + for _, did := range terminalDIDs { + terminal[did] = struct{}{} + } + + var allDIDs []string + skipped := 0 + offset := 0 + for { + repos, err := atsync.StatefulDB.ListRepos(100, offset) + if err != nil { + return nil, err + } + if len(repos) == 0 { + break + } + for _, repo := range repos { + if _, ok := terminal[repo.DID]; ok { + skipped++ + continue + } + allDIDs = append(allDIDs, repo.DID) + } + offset += len(repos) + } + + if skipped > 0 { + log.Log(ctx, "skipping repos with terminal status", "skipped", skipped) + } + return prioritizeDIDs(allDIDs, atsync.CLI.ServerDID(), atsync.CLI.BroadcasterDID()), nil +} + +// prioritizeDIDs moves the given DIDs to the front of the list, in the order +// given, keeping everything else where it was. +// +// This node's own repos go first: they hold the streams, videos and settings +// the node itself serves, so a boot that is going to spend an hour on the +// network should spend its first second on them. +func prioritizeDIDs(dids []string, first ...string) []string { + if len(dids) == 0 || len(first) == 0 { + return dids + } + present := make(map[string]struct{}, len(dids)) + for _, did := range dids { + present[did] = struct{}{} + } + head := make([]string, 0, len(first)) + inHead := make(map[string]struct{}, len(first)) + for _, did := range first { + if did == "" { + continue + } + if _, ok := present[did]; !ok { + continue + } + if _, ok := inHead[did]; ok { + continue + } + head = append(head, did) + inHead[did] = struct{}{} + } + if len(head) == 0 { + return dids + } + out := make([]string, 0, len(dids)) + out = append(out, head...) + for _, did := range dids { + if _, ok := inHead[did]; ok { + continue + } + out = append(out, did) + } + 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 []string + 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 + } + todo = append(todo, did) + } + progress.begin(sweepPhaseShallow, len(todo), time.Now().Add(-InitialWindow)) + if len(todo) == 0 { + return nil + } + log.Log(ctx, "syncing repos", "phase", sweepPhaseShallow, "repos", len(todo)) + + var mu sync.Mutex + failed := 0 + g, gctx := errgroup.WithContext(ctx) + g.SetLimit(sweepConcurrency) + for _, did := range todo { + g.Go(func() error { + if err := gctx.Err(); err != nil { + return err + } + if _, err := atsync.SyncBlueskyRepoCached(gctx, did); err != nil { + log.Error(gctx, "failed to sync repo", "did", did, "err", err) + mu.Lock() + failed++ + mu.Unlock() + return nil + } + progress.finished() + return nil + }) + } + if err := g.Wait(); err != nil { + return err + } + if failed == len(todo) { + return fmt.Errorf("all %d repos failed to sync", failed) + } + return nil +} + +// 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. +// +// 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)) + + for round := 0; len(pending) > 0 && round < maxDeepenRounds; round++ { + if err := ctx.Err(); err != nil { + return err + } + var mu sync.Mutex + var next []string + g, gctx := errgroup.WithContext(ctx) + g.SetLimit(sweepConcurrency) + for _, did := range pending { + g.Go(func() error { + if err := gctx.Err(); err != nil { + return err + } + done, err := atsync.DeepenRepo(gctx, did) + if err != nil { + log.Error(gctx, "failed to deepen repo history", "did", did, "err", err) + return nil + } + if done { + progress.finished() + return nil + } + mu.Lock() + next = append(next, did) + mu.Unlock() + return nil + }) + } + if err := g.Wait(); err != nil { + return err + } + // Restore the priority order the round scrambled. + sort.Slice(next, func(i, j int) bool { return rank[next[i]] < rank[next[j]] }) + pending = next + if _, horizon, err := atsync.deepenPending(ctx, pending); err == nil { + progress.setHorizon(horizon) + } + } + return nil +} + +// deepenPending is the subset of dids whose history is incomplete, 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) ([]string, time.Time, error) { + var pending []string + 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, did) + floor := time.Now() + if repo.BackfillFloor != "" { + if t, err := reposync.TimeForTID(repo.BackfillFloor); err == nil { + floor = t + } + } + if floor.After(horizon) { + horizon = floor + } + } + return pending, horizon, nil +} + +// 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 + total int + horizon time.Time + started bool +} + +// begin starts a phase, resetting the completion count. +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.horizon = horizon + p.started = true +} + +// finished records one repo completing the current phase. +func (p *sweepProgress) finished() { + p.mu.Lock() + defer p.mu.Unlock() + p.done++ +} + +// setHorizon updates how far back the sweep has taken every repo it is working +// on. +func (p *sweepProgress) setHorizon(horizon time.Time) { + p.mu.Lock() + defer p.mu.Unlock() + p.horizon = horizon +} + +// 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. +func (p *sweepProgress) status() []any { + p.mu.Lock() + defer p.mu.Unlock() + horizon := int64(0) + if !p.horizon.IsZero() { + horizon = p.horizon.Unix() + } + return []any{"phase", p.phase, "users", p.done, "total", p.total, "horizon", horizon} +} + +// start runs the status ticker until the returned function is called, which +// also waits for it to stop. Nothing is logged before the first tick, so a +// sweep with nothing to do is silent. +func (p *sweepProgress) start(ctx context.Context) func() { + ctx, cancel := context.WithCancel(ctx) + stopped := make(chan struct{}) + go func() { + defer close(stopped) + ticker := time.NewTicker(sweepStatusInterval) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + p.mu.Lock() + started := p.started + p.mu.Unlock() + if !started { + continue + } + log.Log(ctx, "backfill sweep", p.status()...) + } + } + }() + return func() { + cancel() + <-stopped + } +} diff --git a/pkg/atproto/sweep_test.go b/pkg/atproto/sweep_test.go new file mode 100644 index 000000000..e6332432b --- /dev/null +++ b/pkg/atproto/sweep_test.go @@ -0,0 +1,299 @@ +package atproto + +import ( + "context" + "fmt" + "sync" + "testing" + "time" + + "github.com/bluesky-social/indigo/xrpc" + "github.com/ipfs/go-cid" + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/aqhttp" + "stream.place/streamplace/pkg/bus" + "stream.place/streamplace/pkg/devenv" + "stream.place/streamplace/pkg/model" + "stream.place/streamplace/pkg/placestream" + "stream.place/streamplace/pkg/reposync" +) + +// TestBackfillWindowedHistory is the windowed backfill end to end against the +// reference PDS: a first sync reads the account's configuration and its recent +// chat, and the deepening ladder fetches the rest of its history afterwards, +// one window at a time. +// +// The old messages are planted with explicit TID rkeys, which is how they would +// have arrived months ago -- an rkey is a timestamp, so writing one is the only +// way to have an old record in a repo created a second ago. +func TestBackfillWindowedHistory(t *testing.T) { + dev := devenv.WithDevEnv(t) + ctx := context.Background() + atsync, mod := backfillTestSynchronizer(t, dev) + + user := dev.CreateAccount(t) + now := time.Now() + // Configuration-shaped records: never windowed, always synced. + createBackfillRecord(t, user, "place.stream.chat.profile", "self", &placestream.ChatProfile{}) + // Recent chat, inside the initial window (the PDS mints a TID for now). + createBackfillRecord(t, user, "place.stream.chat.message", "", chatMessageRecord(user.DID, "today")) + // History, at ages that land on distinct rungs of the ladder. + plant := func(age time.Duration, text string) { + t.Helper() + createBackfillRecord(t, user, "place.stream.chat.message", + reposync.TIDForTime(now.Add(-age)), chatMessageRecord(user.DID, text)) + } + plant(3*24*time.Hour, "three days ago") + plant(20*24*time.Hour, "twenty days ago") + plant(200*24*time.Hour, "two hundred days ago") + + countMessages := func() int { + messages, err := mod.MostRecentChatMessages(user.DID) + require.NoError(t, err) + return len(messages) + } + + // Wait for the PDS to have committed everything, by walking the whole + // (unwindowed) range until all five records are there. Doing this before + // the sync means a short count later is a windowing decision, not a race. + require.NoError(t, untilNoErrors(t, func() error { + paths, err := walkAll(ctx, dev, user.DID, backfillRanges("")) + if err != nil { + return err + } + if len(paths) != 5 { + return fmt.Errorf("PDS has %d records, want 5", len(paths)) + } + return nil + }), "waiting for the repo to settle") + + published := watchBus(t, atsync.Bus, user.DID) + + // The shallow sync: everything unwindowed, plus one day of chat. + repo, err := atsync.SyncBlueskyRepoCached(ctx, user.DID) + require.NoError(t, err) + require.NotEmpty(t, repo.Version, "a shallow sync still records the rev it read") + require.NotEmpty(t, repo.BackfillFloor, "a shallow sync records how far back it went") + require.False(t, repo.BackfillDone, "history is not synced yet") + floorTime, err := reposync.TimeForTID(repo.BackfillFloor) + require.NoError(t, err) + require.WithinDuration(t, now.Add(-InitialWindow), floorTime, time.Minute) + + profile, err := mod.GetChatProfile(ctx, user.DID) + require.NoError(t, err) + require.NotNil(t, profile, "unwindowed collections are synced in full on first contact") + require.Equal(t, 1, countMessages(), "only today's message is inside the initial window") + + // Now the ladder. Each rung reaches further back, and a message shows up + // exactly when the window covering its rkey is walked -- not before. + wantAfterRung := []int{ + 2, // [7d, 1d) -- the three-day-old message + 3, // [30d, 7d) -- the twenty-day-old message + 3, // [180d, 30d) -- nothing lives here + 4, // [genesis, 180d) -- the two-hundred-day-old message + } + var done bool + for rung, want := range wantAfterRung { + require.False(t, done, "the ladder finished early at rung %d", rung) + done, err = atsync.DeepenRepo(ctx, user.DID) + require.NoError(t, err, "rung %d", rung) + require.Equal(t, want, countMessages(), "message count after rung %d", rung) + } + require.True(t, done, "the last window bottoms out the collection") + + stored, err := mod.GetRepo(user.DID) + require.NoError(t, err) + require.True(t, stored.BackfillDone, "the watermark is durable") + require.NotEmpty(t, stored.Version) + + // Every message reached the chat bus exactly once, even though the window + // boundaries mean the walker re-emitted records it had already seen. + require.Equal(t, 4, published(), "each message should be published once") + + // And a repo that is done is done: another sweep costs nothing and says + // nothing. + again, err := atsync.DeepenRepo(ctx, user.DID) + require.NoError(t, err) + require.True(t, again) + require.NoError(t, atsync.Sweep(ctx)) + require.Equal(t, 4, countMessages(), "a second sweep must not duplicate anything") + require.Equal(t, 4, published(), "a second sweep must not re-publish anything") +} + +// TestSweepShallowThenDeepens drives the whole sweep over two accounts in the +// two states a real node has after a deploy: one it has never synced, and one +// carrying a row from before the watermark existed. +func TestSweepShallowThenDeepens(t *testing.T) { + dev := devenv.WithDevEnv(t) + ctx := context.Background() + atsync, mod := backfillTestSynchronizer(t, dev) + now := time.Now() + + fresh := dev.CreateAccount(t) + createBackfillRecord(t, fresh, "place.stream.chat.profile", "self", &placestream.ChatProfile{}) + createBackfillRecord(t, fresh, "place.stream.chat.message", "", chatMessageRecord(fresh.DID, "fresh today")) + createBackfillRecord(t, fresh, "place.stream.chat.message", + reposync.TIDForTime(now.Add(-90*24*time.Hour)), chatMessageRecord(fresh.DID, "fresh long ago")) + + legacy := dev.CreateAccount(t) + createBackfillRecord(t, legacy, "place.stream.chat.message", "", chatMessageRecord(legacy.DID, "legacy today")) + createBackfillRecord(t, legacy, "place.stream.chat.message", + reposync.TIDForTime(now.Add(-300*24*time.Hour)), chatMessageRecord(legacy.DID, "legacy ages ago")) + + require.NoError(t, untilNoErrors(t, func() error { + for did, want := range map[string]int{fresh.DID: 3, legacy.DID: 2} { + paths, err := walkAll(ctx, dev, did, backfillRanges("")) + if err != nil { + return err + } + if len(paths) != want { + return fmt.Errorf("repo %s has %d records, want %d", did, len(paths), want) + } + } + return nil + }), "waiting for the repos to settle") + + // The fresh account is known but unsynced: a placeholder row, exactly what + // the firehose writes when it first sees a record from a stranger. + require.NoError(t, atsync.StatefulDB.AddRepo(fresh.DID)) + // The legacy account has a completed sync from a build that had never heard + // of a backfill window: a version, no floor, not done. + require.NoError(t, mod.UpdateRepo(&model.Repo{ + DID: legacy.DID, + PDS: dev.PDSURL, + Handle: legacy.Handle, + Version: "3lpretend0000", + })) + require.NoError(t, atsync.StatefulDB.AddRepo(legacy.DID)) + + require.NoError(t, atsync.Sweep(ctx)) + + for _, did := range []string{fresh.DID, legacy.DID} { + stored, err := mod.GetRepo(did) + require.NoError(t, err) + require.NotEmpty(t, stored.Version, "%s should have been synced", did) + require.True(t, stored.BackfillDone, "%s should have been deepened to the end", did) + messages, err := mod.MostRecentChatMessages(did) + require.NoError(t, err, did) + require.Len(t, messages, 2, "both messages of %s should be indexed", did) + } + // The fresh account went through the full backfill, so its unwindowed + // records are there too; the legacy one was only ever deepened, which by + // design touches nothing but the windowed collections. + profile, err := mod.GetChatProfile(ctx, fresh.DID) + require.NoError(t, err) + require.NotNil(t, profile) + + // Idempotent: a second sweep is a few head fetches and nothing else. + require.NoError(t, atsync.Sweep(ctx)) + for _, did := range []string{fresh.DID, legacy.DID} { + messages, err := mod.MostRecentChatMessages(did) + require.NoError(t, err) + require.Len(t, messages, 2, "a second sweep must not duplicate anything") + } +} + +// TestSweepPrioritizesOwnDIDs: the node's own repos hold what it serves, so +// they go first. Ordering is checked directly because staging a node's own +// did:web account inside the dev environment proves nothing about the order. +func TestSweepPrioritizesOwnDIDs(t *testing.T) { + dids := []string{"did:plc:a", "did:web:server.example", "did:plc:b", "did:web:broadcaster.example", "did:plc:c"} + + require.Equal(t, + []string{"did:web:server.example", "did:web:broadcaster.example", "did:plc:a", "did:plc:b", "did:plc:c"}, + prioritizeDIDs(dids, "did:web:server.example", "did:web:broadcaster.example")) + + // A node whose server and broadcaster are the same host lists it once. + require.Equal(t, + []string{"did:web:server.example", "did:plc:a", "did:plc:b", "did:web:broadcaster.example", "did:plc:c"}, + prioritizeDIDs(dids, "did:web:server.example", "did:web:server.example")) + + // DIDs that are not in the sweep, or not configured, change nothing. + require.Equal(t, dids, prioritizeDIDs(dids, "did:web:nowhere.example", "")) + require.Equal(t, dids, prioritizeDIDs(dids)) + require.Nil(t, prioritizeDIDs(nil, "did:web:server.example")) +} + +// 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. +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()) + + horizon := time.Now().Add(-InitialWindow) + progress.begin(sweepPhaseShallow, 3, horizon) + progress.finished() + require.Equal(t, + []any{"phase", "shallow", "users", 1, "total", 3, "horizon", horizon.Unix()}, + progress.status()) + + // A new phase resets the count and moves the horizon. + 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()}, + progress.status()) + 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()}, + progress.status()) + + // The ticker stops when told to, without leaking a goroutine. + stop := progress.start(context.Background()) + stop() +} + +// walkAll walks a repo's ranges against the dev PDS and returns the paths, so +// tests can wait for the PDS to have committed what they wrote. +func walkAll(ctx context.Context, dev *devenv.DevEnv, did string, ranges []reposync.KeyRange) ([]string, error) { + xrpcc := &xrpc.Client{Host: dev.PDSURL, Client: &aqhttp.Client} + fetcher := &reposync.CachedFetcher{ + Cache: reposync.NewMemoryBlockCache(), + Inner: &reposync.XRPCBlockFetcher{Client: xrpcc, DID: did}, + } + head, err := reposync.FetchVerifiedHead(ctx, xrpcc, fetcher, dev.TestDirectory(), did) + if err != nil { + return nil, err + } + var paths []string + err = (&reposync.Walker{Fetcher: fetcher}).WalkRanges(ctx, head.Root, ranges, + func(path string, _ cid.Cid, _ []byte) error { + paths = append(paths, path) + return nil + }) + if err != nil { + return nil, err + } + return paths, nil +} + +// watchBus counts what a topic publishes, for asserting that re-walked records +// do not reach subscribers twice. +func watchBus(t *testing.T, b *bus.Bus, topic string) func() int { + t.Helper() + ch := b.Subscribe(topic) + t.Cleanup(func() { b.Unsubscribe(topic, ch) }) + var mu sync.Mutex + count := 0 + go func() { + for range ch { + mu.Lock() + count++ + mu.Unlock() + } + }() + return func() int { + // The publish is asynchronous; give it a moment to happen before + // reporting a count that a test is about to assert on. + time.Sleep(250 * time.Millisecond) + mu.Lock() + defer mu.Unlock() + return count + } +} diff --git a/pkg/cmd/streamplace.go b/pkg/cmd/streamplace.go index 43d5b2fa6..cfaf21c3c 100644 --- a/pkg/cmd/streamplace.go +++ b/pkg/cmd/streamplace.go @@ -85,6 +85,7 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { makeSplitCommand(build), makeLivepeerCommand(build), makeMigrateCommand(build), + makeSyncCommand(build), } // Add the verbosity flag // app.Flags = append(app.Flags, &urfavecli.StringFlag{ @@ -267,13 +268,16 @@ func runMain(ctx context.Context, build *config.BuildFlags, platformJobs []jobFu Noter: noter, Bus: b, } - // Sync every repo we know about, once per boot. Nothing below depends on it - // having finished: it is a repair sweep for repos left half-indexed by a - // previous run, and the firehose keeps them current afterwards. - err = atsync.Migrate(ctx) - if err != nil { - return fmt.Errorf("failed to migrate: %w", err) - } + // Sync every repo we know about, once per boot: a repair pass for repos left + // half-indexed by a previous run, and then history deepening, which on a + // fresh node runs for as long as the network is big. Nothing below depends + // on it, so it runs in the background off the serve context -- shutdown + // cancels it -- and the node is up and serving in the meantime. + go func() { + if err := atsync.Sweep(ctx); err != nil && ctx.Err() == nil { + log.Error(ctx, "backfill sweep failed", "err", err) + } + }() mm, err := media.MakeMediaManager(ctx, cli, signer, mod, b, atsync, ldb) if err != nil { @@ -1165,6 +1169,56 @@ func makeMigrateCommand(build *config.BuildFlags) *urfavecli.Command { } } +// makeSyncCommand runs the backfill sweep to completion and exits, without +// starting a node. It is for the case where a new index revision has to be warm +// before traffic reaches it: run this, wait for it to finish, then start the +// server -- rather than starting the server and serving from an index that is +// still filling in behind it. +func makeSyncCommand(build *config.BuildFlags) *urfavecli.Command { + cli := config.CLI{Build: build} + syncCmd := cli.NewCommand("sync") + syncCmd.Usage = "index every repo this node knows about, then exit" + syncCmd.Action = func(ctx context.Context, cmd *urfavecli.Command) error { + return runSync(ctx, build, cmd, &cli) + } + return syncCmd +} + +// runSync builds the smallest stack a sweep needs -- the index, the state +// database, an identity resolver -- and nothing else. No HTTP servers, no media +// manager, no firehose: this process talks to other people's PDSes and to the +// two databases, and then it is done. +func runSync(ctx context.Context, build *config.BuildFlags, cmd *urfavecli.Command, cli *config.CLI) error { + if err := cli.Validate(cmd); err != nil { + return err + } + log.SetColorLogger(cli.Color) + ctx = log.WithDebugValue(ctx, cli.Debug) + log.Log(ctx, "streamplace sync", "version", build.Version, "dataDir", cli.DataDir) + + if err := os.MkdirAll(cli.DataDir, os.ModePerm); err != nil { + return fmt.Errorf("error creating streamplace dir at %s: %w", cli.DataDir, err) + } + mod, err := model.MakeDB(cli.DataFilePath([]string{"index"})) + if err != nil { + return err + } + state, err := statedb.MakeDB(ctx, cli, nil, mod) + if err != nil { + return err + } + atsync := &atproto.ATProtoSynchronizer{ + CLI: cli, + Model: mod, + StatefulDB: state, + Bus: bus.NewBus(), + } + // A sweep that could not sync a single repo is a broken node and exits + // nonzero; anything less than that heals on the next run, so it is logged + // and forgiven. + return atsync.Sweep(ctx) +} + // resolveLiveSigningKey returns the did:key whose private half signed a // streamer's live segments, for stamping on live-to-VOD place.stream.media.track // records so playback can verify them. It picks the most recently created diff --git a/pkg/cmd/sync_test.go b/pkg/cmd/sync_test.go new file mode 100644 index 000000000..c2ccb4e3c --- /dev/null +++ b/pkg/cmd/sync_test.go @@ -0,0 +1,38 @@ +package cmd + +import ( + "testing" + + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/config" +) + +// TestSyncCommand checks the registration, which is the part that is easy to +// get wrong: the command has to carry the server's flags (its own --data-dir +// and --db-url, resolved into the same CLI the action reads) rather than being +// a bare subcommand that runs against defaults. +func TestSyncCommand(t *testing.T) { + cmd := makeSyncCommand(&config.BuildFlags{Version: "test"}) + require.Equal(t, "sync", cmd.Name) + require.NotNil(t, cmd.Action) + + names := map[string]bool{} + for _, flag := range cmd.Flags { + for _, name := range flag.Names() { + names[name] = true + } + } + require.True(t, names["data-dir"], "sync needs the data dir to find the index") + require.True(t, names["db-url"], "sync needs the state database") + require.True(t, names["plc-url"], "sync resolves identities") +} + +// TestSyncCommandRuns runs the command for real against empty databases: it +// opens the index and the state database, sweeps the zero repos in them, and +// exits successfully. No HTTP server is started and no network is touched, +// which is the point of the command. +func TestSyncCommandRuns(t *testing.T) { + cmd := makeSyncCommand(&config.BuildFlags{Version: "test"}) + err := cmd.Run(t.Context(), []string{"sync", "--data-dir", t.TempDir(), "--db-url", ":memory:"}) + require.NoError(t, err) +} diff --git a/pkg/model/model.go b/pkg/model/model.go index 75869bdb3..63b8429de 100644 --- a/pkg/model/model.go +++ b/pkg/model/model.go @@ -38,6 +38,7 @@ type Model interface { GetAllRepos() ([]Repo, error) SearchReposByHandle(query string, limit int) ([]Repo, error) UpdateRepo(repo *Repo) error + AdvanceRepoBackfill(ctx context.Context, did, version, rootCID, floor string, done bool) error SetRepoStatus(ctx context.Context, did string, status string) error TerminalRepoDIDs(ctx context.Context) ([]string, error) diff --git a/pkg/model/repo.go b/pkg/model/repo.go index 92dbc5994..e5d7c988b 100644 --- a/pkg/model/repo.go +++ b/pkg/model/repo.go @@ -27,6 +27,14 @@ type Repo struct { RootCID string `json:"rootCid"` // Status is one of the RepoStatus* constants; empty for a normal account. Status string `gorm:"column:status" json:"status,omitempty"` + // BackfillFloor is a TID watermark for the collections a backfill reads by + // time window (chat messages, feed posts): their history is contiguously + // indexed from this TID up to now. Empty means no window has been recorded + // -- either nothing is synced yet, or the row predates the watermark. + BackfillFloor string `gorm:"column:backfill_floor" json:"backfillFloor,omitempty"` + // BackfillDone reports that those windowed collections are indexed all the + // way back to the start of the repo, so there is no history left to fetch. + BackfillDone bool `gorm:"column:backfill_done" json:"backfillDone,omitempty"` } // TerminalStatus reports whether this repo is in an account state no amount of @@ -105,6 +113,26 @@ func (m *DBModel) SetRepoStatus(ctx context.Context, did string, status string) return m.DB.WithContext(ctx).Model(&Repo{}).Where("did = ?", did).Update("status", status).Error } +// AdvanceRepoBackfill records the outcome of one deepening window: the repo is +// now indexed from floor forward (empty floor meaning all the way back), at the +// revision that window was read at. +// +// It writes exactly those four columns rather than the whole row, so a +// concurrent handle change or status update cannot be rolled back by a sweep +// that read the row minutes ago. Select names the fields explicitly, which is +// also what makes the zero values -- an empty floor, a false flag -- get +// written instead of skipped. +func (m *DBModel) AdvanceRepoBackfill(ctx context.Context, did, version, rootCID, floor string, done bool) error { + return m.DB.WithContext(ctx).Model(&Repo{}).Where("did = ?", did). + Select("Version", "RootCID", "BackfillFloor", "BackfillDone"). + Updates(Repo{ + Version: version, + RootCID: rootCID, + BackfillFloor: floor, + BackfillDone: done, + }).Error +} + // TerminalRepoDIDs lists the repos parked in a terminal account state, so the // boot-time sync sweep can skip them in one query instead of failing on each. func (m *DBModel) TerminalRepoDIDs(ctx context.Context) ([]string, error) { diff --git a/pkg/reposync/retry.go b/pkg/reposync/retry.go index ce041e390..12a75a497 100644 --- a/pkg/reposync/retry.go +++ b/pkg/reposync/retry.go @@ -8,6 +8,7 @@ import ( "math/rand" "net" "net/http" + "strings" "syscall" "time" @@ -108,7 +109,7 @@ func (p RetryPolicy) do(ctx context.Context, what string, fn func() error) error return fmt.Errorf("giving up after %d attempts: %w", attempt, err) } d := p.delay(attempt, err) - log.Warn(ctx, "retrying transient xrpc failure", "call", what, "attempt", attempt, "wait", d, "err", err) + log.Warn(ctx, "retrying transient xrpc failure", "call", what, "attempt", attempt, "wait", d, "err", errForLog(err)) if serr := sleepCtx(ctx, d); serr != nil { return fmt.Errorf("aborted after %d attempts: %w", attempt, errors.Join(err, serr)) } @@ -130,6 +131,39 @@ func sleepCtx(ctx context.Context, d time.Duration) error { } } +// errForLog renders a retryable failure for the warning line above. +// +// It exists for one shape: a host that throttles or 502s us with an HTML error +// page. indigo tries to JSON-decode every non-200 body, so what surfaces is +// `XRPC ERROR 429: failed to decode xrpc error message: invalid character '<' +// looking for beginning of value` -- forty characters of JSON parser trivia in +// front of the one fact that matters, repeated for every retry of every walk. +// Say "HTTP 429 (undecodable error body)" instead. Only this log line is +// compressed; the error returned to the caller keeps the whole chain. +func errForLog(err error) any { + var xe *xrpc.Error + if !errors.As(err, &xe) || xe.StatusCode == 0 || !isUndecodableBody(xe.Wrapped) { + return err + } + return fmt.Sprintf("HTTP %d (undecodable error body)", xe.StatusCode) +} + +// undecodableBodyPrefix is indigo's wrapper around a response body that is not +// the JSON error object the lexicon promises. +const undecodableBodyPrefix = "failed to decode xrpc error message" + +func isUndecodableBody(err error) bool { + if err == nil { + return false + } + var xe *xrpc.XRPCError + if errors.As(err, &xe) { + // A decoded (if empty) error object: the host answered properly. + return false + } + return strings.HasPrefix(err.Error(), undecodableBodyPrefix) +} + // isRetryable reports whether err is the kind of failure that is likely to go // away on its own: the host throttled us, the host is briefly broken, or the // connection died under us. diff --git a/pkg/reposync/retry_test.go b/pkg/reposync/retry_test.go index 48d3f593e..9eff9e478 100644 --- a/pkg/reposync/retry_test.go +++ b/pkg/reposync/retry_test.go @@ -160,6 +160,36 @@ func ratelimited(reset time.Time) error { }) } +// TestErrForLog: the retry warning is the one line an operator sees when a host +// is throttling a walk, and for the commonest case -- a 429 with an HTML body -- +// indigo's JSON decoder failure was burying the status code in parser trivia. +func TestErrForLog(t *testing.T) { + // The genuine article, straight off the wire. + sr := buildSignedRepo(t, testDID, exactnessPaths()) + host := newFakeHost(sr) + host.blocksFailures = []failure{htmlThrottled} + client := host.start(t) + f := &XRPCBlockFetcher{Client: client, DID: testDID, Retry: RetryPolicy{MaxAttempts: 1}} + _, err := f.GetBlocks(context.Background(), []cid.Cid{sr.root}) + require.Error(t, err) + require.Contains(t, err.Error(), "failed to decode xrpc error message", + "fixture no longer produces the shape under test") + require.Equal(t, "HTTP 429 (undecodable error body)", errForLog(err)) + + // A host that answered properly keeps its whole message. + proper := xrpcErr(http.StatusTooManyRequests, "RateLimitExceeded", "Rate Limit Exceeded") + require.Equal(t, proper, errForLog(proper)) + // So does anything that never reached a host. + plain := fmt.Errorf("dialing: %w", syscall.ECONNREFUSED) + require.Equal(t, plain, errForLog(plain)) + // And a 502 from a load balancer gets the same treatment as the 429. + gateway := fmt.Errorf("getBlocks: %w", &xrpc.Error{ + StatusCode: http.StatusBadGateway, + Wrapped: fmt.Errorf("failed to decode xrpc error message: %w", errors.New("invalid character '<'")), + }) + require.Equal(t, "HTTP 502 (undecodable error body)", errForLog(gateway)) +} + // TestXRPCBlockFetcherRetries drives the retry loop through the real getBlocks // path against an HTTP host that fails on a script. func TestXRPCBlockFetcherRetries(t *testing.T) { diff --git a/pkg/reposync/tid.go b/pkg/reposync/tid.go new file mode 100644 index 000000000..de99718bc --- /dev/null +++ b/pkg/reposync/tid.go @@ -0,0 +1,41 @@ +package reposync + +import ( + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" +) + +// TIDForTime returns the TID that marks the instant t. +// +// TIDs are 13 characters of sortable base32 over (unix microseconds << 10 | +// clock id), so they sort chronologically and MST keys of the form +// "collection/" sort chronologically within a collection. That makes a +// time window a key range: everything a repo wrote to a TID-keyed collection +// since t lives in ["collection/"+TIDForTime(t), end of collection). +// +// The clock id is zero, which is the smallest one, so the result sorts at or +// below every real TID stamped at t and strictly above every TID stamped +// before it. A range starting here therefore cannot miss a record. +// +// Times before the unix epoch clamp to it: TIDs cannot express them, and the +// only sensible reading of "before 1970" for a repo window is "from the +// beginning". +func TIDForTime(t time.Time) string { + micros := t.UTC().UnixMicro() + if micros < 0 { + micros = 0 + } + return string(syntax.NewTID(micros, 0)) +} + +// TimeForTID is the inverse of [TIDForTime]: the wall-clock instant a TID +// encodes. It errors on anything that is not TID syntax, so a caller reading a +// watermark out of a database can tell a real timestamp from a stray string. +func TimeForTID(tid string) (time.Time, error) { + parsed, err := syntax.ParseTID(tid) + if err != nil { + return time.Time{}, err + } + return parsed.Time(), nil +} diff --git a/pkg/reposync/tid_test.go b/pkg/reposync/tid_test.go new file mode 100644 index 000000000..f430b19d5 --- /dev/null +++ b/pkg/reposync/tid_test.go @@ -0,0 +1,120 @@ +package reposync + +import ( + "testing" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/stretchr/testify/require" +) + +// TestTIDForTimeIsAValidTID: the whole point is producing a key that can sit in +// an MST range next to real record keys, so it has to be TID syntax. +func TestTIDForTimeIsAValidTID(t *testing.T) { + for _, ts := range []time.Time{ + time.Now(), + time.Now().Add(-24 * time.Hour), + time.Now().Add(-10 * 365 * 24 * time.Hour), + time.Unix(0, 0), + time.Date(2023, 6, 1, 12, 0, 0, 0, time.UTC), + } { + tid := TIDForTime(ts) + _, err := syntax.ParseTID(tid) + require.NoError(t, err, "TIDForTime(%s) = %q", ts, tid) + require.Len(t, tid, 13) + } +} + +// TestTIDForTimeOrdering: the sort order of the strings has to match the order +// of the instants, because that equivalence is what makes a time window a key +// range. +func TestTIDForTimeOrdering(t *testing.T) { + now := time.Now() + spans := []time.Duration{ + 0, + -time.Microsecond, + -time.Second, + -time.Hour, + -24 * time.Hour, + -7 * 24 * time.Hour, + -30 * 24 * time.Hour, + -180 * 24 * time.Hour, + -5 * 365 * 24 * time.Hour, + } + prev := TIDForTime(now.Add(spans[0])) + for _, span := range spans[1:] { + tid := TIDForTime(now.Add(span)) + require.Less(t, tid, prev, "an earlier instant must produce a smaller TID (span %s)", span) + prev = tid + } +} + +// TestTIDForTimeBoundary is the property the windowed backfill relies on: a +// record stamped at or after t is inside the range that starts at TIDForTime(t), +// and one stamped before it is outside -- for every clock id, since we have no +// say over which one a remote PDS uses. +func TestTIDForTimeBoundary(t *testing.T) { + t0 := time.Now().Add(-36 * time.Hour).Truncate(time.Microsecond) + floor := TIDForTime(t0) + + for _, clockID := range []uint{0, 1, 7, 512, 1023} { + at := string(syntax.NewTIDFromTime(t0, clockID)) + require.GreaterOrEqual(t, at, floor, + "a TID stamped exactly at the floor instant (clock %d) must be in range", clockID) + + after := string(syntax.NewTIDFromTime(t0.Add(time.Microsecond), clockID)) + require.Greater(t, after, floor, "a TID stamped after the floor must be in range") + + before := string(syntax.NewTIDFromTime(t0.Add(-time.Microsecond), clockID)) + require.Less(t, before, floor, "a TID stamped before the floor must be out of range") + } +} + +// TestTIDForTimeAgainstNewTIDNow: a TID minted right now sorts above a floor +// taken a moment ago and below one taken a moment hence. +func TestTIDForTimeAgainstNewTIDNow(t *testing.T) { + before := TIDForTime(time.Now().Add(-time.Second)) + now := string(syntax.NewTIDNow(0)) + after := TIDForTime(time.Now().Add(time.Second)) + + require.Greater(t, now, before) + require.Less(t, now, after) +} + +// TestTimeForTIDRoundTrip: the watermark stored in a repo row has to be +// readable back as a timestamp, since that is how the sweep decides which +// window to walk next. +func TestTimeForTIDRoundTrip(t *testing.T) { + for _, ts := range []time.Time{ + time.Now().Truncate(time.Microsecond), + time.Now().Add(-90 * 24 * time.Hour).Truncate(time.Microsecond), + time.Unix(0, 0), + } { + got, err := TimeForTID(TIDForTime(ts)) + require.NoError(t, err) + require.True(t, got.Equal(ts.UTC()), "round trip of %s gave %s", ts.UTC(), got) + } + + // A real TID keeps its timestamp too, clock id and all. + tid := syntax.NewTIDNow(42) + got, err := TimeForTID(string(tid)) + require.NoError(t, err) + require.True(t, got.Equal(tid.Time())) + + for _, bad := range []string{"", "not-a-tid", "3jui7kd54zh2", "3JUI7KD54ZH2Y"} { + _, err := TimeForTID(bad) + require.Error(t, err, "TimeForTID(%q) should not parse", bad) + } +} + +// TestTIDForTimeMonotonic: repeatedly slicing a window off the front never +// walks backwards, which is what keeps the deepening ladder terminating. +func TestTIDForTimeMonotonic(t *testing.T) { + now := time.Now() + last := TIDForTime(now) + for i := 1; i <= 100; i++ { + tid := TIDForTime(now.Add(-time.Duration(i) * time.Minute)) + require.Less(t, tid, last) + last = tid + } +}