diff --git a/pkg/atproto/contiguity.go b/pkg/atproto/contiguity.go index f07aba14..ac99084a 100644 --- a/pkg/atproto/contiguity.go +++ b/pkg/atproto/contiguity.go @@ -13,12 +13,15 @@ import ( // revCASAttempts is how many times a commit tries to place itself on the repo // row before concluding it found a gap. // -// One attempt is the whole story when commits arrive in order. They do not: -// events are handled one goroutine each, so two commits on one repo race, and -// the older one can win the CAS after the newer one has already missed it. A -// second look then finds the row exactly where the newer commit expected it. -// Three is one more than that story needs; a lost race past it costs one -// unnecessary repair, never a missed record. +// One attempt is the whole story when commits arrive in order. Within a single +// relay they now do -- the parallel scheduler chains same-DID events through a +// per-repo queue -- but this node consumes several relays at once, each with a +// scheduler of its own, and the deduper only guarantees that a given commit is +// handled once, not that two commits on one repo are handled in order. So they +// still race, and the older one can win the CAS after the newer one has already +// missed it. A second look then finds the row exactly where the newer commit +// expected it. Three is one more than that story needs; a lost race past it +// costs one unnecessary repair, never a missed record. const revCASAttempts = 3 // repairSlack is how far before the last known rev a repair starts reading. diff --git a/pkg/atproto/firehose.go b/pkg/atproto/firehose.go index aef45112..3742d0e5 100644 --- a/pkg/atproto/firehose.go +++ b/pkg/atproto/firehose.go @@ -242,16 +242,51 @@ func (atsync *ATProtoSynchronizer) consumeRelay(ctx context.Context, relay strin } } +// firehoseReplayWindow is how stale a stored relay cursor may be and still be +// worth replaying from. Zero (or negative) disables the cap, so every connect +// resumes from whatever is stored. +func (atsync *ATProtoSynchronizer) firehoseReplayWindow() time.Duration { + if atsync.CLI == nil { + return config.DefaultFirehoseReplayWindow + } + return atsync.CLI.FirehoseReplayWindow +} + +// dropStaleCursor abandons a relay cursor too old to be worth replaying from, +// and orders a sweep to heal the gap that leaves. +// +// The kick matters: without it the gap sits there until the next scheduled +// sweep, up to --sweep-interval away, which is the one way dropping the cursor +// could be worse than replaying it. sweepOnce skips if a sweep is already +// running, so several relays deciding this at once costs nothing. +func (atsync *ATProtoSynchronizer) dropStaleCursor(ctx context.Context, cursor *relayCursor) { + if !cursor.dropIfStale(ctx, atsync.firehoseReplayWindow()) { + return + } + if atsync.StatefulDB == nil { + // No index to sweep (tests that exercise the firehose on its own). + return + } + log.Log(ctx, "kicking off a sweep to heal the skipped firehose replay") + go atsync.sweepOnce(ctx, atsync.Sweep) +} + // connectRelay dials one relay and pumps its firehose until the connection -// drops or ctx is cancelled. Event handlers are spawned on the parent ctx (not -// the per-connection one) so an in-flight commit keeps indexing across a -// reconnect — important because dedup has already claimed it, so no other relay -// will re-deliver it to us. +// drops or ctx is cancelled. Event handlers run on the parent ctx (not the +// per-connection one) so an in-flight commit keeps indexing across a reconnect +// — important because dedup has already claimed it, so no other relay will +// re-deliver it to us. func (atsync *ATProtoSynchronizer) connectRelay(ctx context.Context, relay string, cursor *relayCursor) error { u, err := url.Parse(relay) if err != nil { return fmt.Errorf("invalid relay URI %q: %w", relay, err) } + // A cursor whose newest event is hours old asks the relay to re-send hours + // of the whole network's traffic, at full speed, to a node that indexes a + // fraction of a percent of it. Checked here rather than once at load + // because a long backoff stretch ages a cursor that was fresh when we + // started. + atsync.dropStaleCursor(ctx, cursor) // MoQ relays (moqt:// and aliases) are consumed over QUIC instead of // WebSocket. Everything downstream of frame-decode — dedup, cursor, // handlers, backoff — is shared, so this is just a transport swap. @@ -328,10 +363,44 @@ func relayProtocol(relay string) string { // repoStreamCallbacks builds the event callbacks shared by the WebSocket and // MoQ transports: per-relay metrics, liveness marking, cursor advance, -// cross-relay dedup, and spawning the indexing handlers on the parent ctx. -// cancel ends the current connection when the relay sends an error frame. The -// handlers run on ctx (not the per-connection context) so an in-flight commit -// keeps indexing across a reconnect. +// cross-relay dedup, and running the indexing handlers. cancel ends the current +// connection when the relay sends an error frame. +// +// The handlers run INLINE, on the callback's own goroutine, which is one of +// indigo's parallel-scheduler workers. That is deliberate and load-bearing. +// Spawning a goroutine per event here — which this used to do — returns from +// the callback instantly and so defeats every bound the scheduler exists to +// provide. Running them inline gets all three back: +// +// - Concurrency is capped at the scheduler's worker count, instead of one +// goroutine per event. A relay replaying a backlog at full send rate is +// what turned that into millions of goroutines all queued behind a single +// sqlite write connection. +// - Events for one repo are handled in the order the relay sent them: the +// scheduler chains same-DID tasks through a per-repo queue, so at most one +// is in flight per repo. The per-event goroutines had silently dropped that +// guarantee. It is per-relay, not global — each relay has its own scheduler, +// so see revCASAttempts for the cross-relay race that remains. +// - Backpressure. The scheduler's feeder channel is unbuffered, so once every +// worker is busy, AddWork for a new repo blocks, HandleRepoStream (or the +// moq read loop) stops reading the socket, and TCP tells the relay to slow +// down. Overload degrades to "the firehose lags", which the sweep covers, +// instead of to an OOM. +// +// Both transports feed the same scheduler (connectRelayMoq → dispatchMoqFrame → +// scheduler.AddWork), so this holds for websocket and moq alike. +// +// The handlers keep using the captured ctx rather than the one the scheduler +// hands them: parallel workers call the event handler with context.TODO(), and +// the captured ctx is where our log values and shutdown cancellation live. It +// is the parent ctx and not the per-connection one so an in-flight commit keeps +// indexing across a reconnect. +// +// The tradeoff, accepted: on disconnect both transports defer +// scheduler.Shutdown(), which lets in-flight handlers finish but drops tasks +// still queued behind them — whose dedup claim has already been made, so no +// other relay will re-deliver them. That is a small silent gap, and healing +// exactly this kind of gap is what the head check and sweep are for. func (atsync *ATProtoSynchronizer) repoStreamCallbacks(ctx context.Context, relay string, cursor *relayCursor, cancel context.CancelFunc) *events.RepoStreamCallbacks { protocol := relayProtocol(relay) // observeSeq advances the cursor and republishes the per-relay high-water @@ -354,27 +423,41 @@ func (atsync *ATProtoSynchronizer) repoStreamCallbacks(ctx context.Context, rela cursor.observe(seq) spmetrics.FirehoseRelayHighSeq.WithLabelValues(relay, protocol).Set(float64(cursor.highSeq())) } + // observeTime records how current this relay is. Unlike the seq it applies + // to every transport — an event's stamp means the same thing whoever + // carried it — and it is recorded before dedup, because a duplicate is + // still proof of how far along this relay is. One small parse per event; an + // unparsable stamp is simply not an observation. + observeTime := func(t string) { + aqt, err := aqtime.FromString(t) + if err != nil { + return + } + cursor.observeTime(aqt.Time()) + } return &events.RepoStreamCallbacks{ RepoCommit: func(evt *indigoatproto.SyncSubscribeRepos_Commit) error { atsync.markSeen() observeSeq(evt.Seq) + observeTime(evt.Time) spmetrics.FirehoseEventsReceivedTotal.WithLabelValues(relay, protocol, "commit").Inc() if atsync.commitDedup.seen(evt.Commit.String()) { spmetrics.FirehoseEventsDedupedTotal.WithLabelValues(relay, protocol, "commit").Inc() return nil } - go atsync.handleCommitEventOps(ctx, evt) + atsync.handleCommitEventOps(ctx, evt) return nil }, RepoIdentity: func(evt *indigoatproto.SyncSubscribeRepos_Identity) error { atsync.markSeen() observeSeq(evt.Seq) + observeTime(evt.Time) spmetrics.FirehoseEventsReceivedTotal.WithLabelValues(relay, protocol, "identity").Inc() if atsync.identityDedup.seen(identityDedupKey(evt)) { spmetrics.FirehoseEventsDedupedTotal.WithLabelValues(relay, protocol, "identity").Inc() return nil } - go atsync.handleIdentityEventOps(ctx, evt) + atsync.handleIdentityEventOps(ctx, evt) return nil }, Error: func(evt *events.ErrorFrame) error { @@ -449,6 +532,44 @@ var CollectionFilter = []string{ constants.PLACE_STREAM_LIVE_RECOMMENDATIONS, } +// indexedCollection reports whether we index anything in this collection: the +// explicit [CollectionFilter] list plus everything under place.stream., which is +// our own namespace and always ours to index. An empty filter means "index +// everything", as it always has. +func indexedCollection(collection string) bool { + if len(CollectionFilter) == 0 { + return true + } + if slices.Contains(CollectionFilter, collection) { + return true + } + return strings.HasPrefix(collection, "place.stream.") +} + +// commitHasIndexedOps reports whether a commit touches anything we index, from +// the op paths alone — no CAR parsing. +// +// This is what lets handleCommitEventOps skip the CAR parse for the large +// majority of firehose traffic (the whole network's posts, likes and follows on +// repos we have never heard of). Handlers run on a bounded worker pool now, so +// that parse is no longer paid by a throwaway goroutine: it is paid out of the +// pool, and directly caps how fast the node can catch up. +// +// An op whose path does not parse counts as indexed, so that the main loop +// reaches it and logs the same error it always has. +func commitHasIndexedOps(evt *indigoatproto.SyncSubscribeRepos_Commit) bool { + for _, op := range evt.Ops { + collection, _, err := syntax.ParseRepoPath(op.Path) + if err != nil { + return true + } + if indexedCollection(collection.String()) { + return true + } + } + return false +} + func (atsync *ATProtoSynchronizer) handleCommitEventOps(ctx context.Context, evt *indigoatproto.SyncSubscribeRepos_Commit) { ctx = log.WithLogValues(ctx, "event", "commit", "did", evt.Repo, "rev", evt.Rev, "seq", fmt.Sprintf("%d", evt.Seq), "func", "handleCommitEventOps") @@ -457,10 +578,34 @@ func (atsync *ATProtoSynchronizer) handleCommitEventOps(ctx context.Context, evt return } + // A commit with nothing we index still has to be tracked: trackCommitRev + // follows the rev chain for repos we DO track, and a tracked repo posting a + // like is a commit whose ops are all filtered. Skipping it there would + // leave our stored rev behind, and the next commit we did care about would + // look like a hole and order a pointless repair. So the ops are what get + // skipped, not the event — and trackCommitRev never needed the CAR anyway. + if commitHasIndexedOps(evt) && !atsync.handleIndexedOps(ctx, evt) { + return + } + + // Every op in this commit is indexed, so the index can claim this commit. + // Only reached on a clean pass: an event we bailed out of half-applied + // leaves the stored rev where it was, and the next commit for that repo + // notices the hole and orders a repair. + atsync.trackCommitRev(ctx, evt) +} + +// handleIndexedOps reads the commit's CAR and applies the ops we index, +// reporting whether it got all the way through. It is split out of +// handleCommitEventOps so that a commit with nothing for us can skip the CAR +// parse entirely while still having its rev tracked; a bail-out here (an +// unreadable CAR, an unparsable path) reports false and so leaves the repo's +// stored rev where it was, exactly as before, so the next commit finds the hole. +func (atsync *ATProtoSynchronizer) handleIndexedOps(ctx context.Context, evt *indigoatproto.SyncSubscribeRepos_Commit) bool { rr, err := repo.ReadRepoFromCar(ctx, bytes.NewReader(evt.Blocks)) if err != nil { log.Error(ctx, "failed to read repo from car", "err", err) - return + return false } for _, op := range evt.Ops { @@ -468,18 +613,12 @@ func (atsync *ATProtoSynchronizer) handleCommitEventOps(ctx context.Context, evt uri := fmt.Sprintf("at://%s/%s", evt.Repo, op.Path) if err != nil { log.Error(ctx, "invalid path in repo op", "eventKind", op.Action, "path", op.Path) - return + return false } ctx := log.WithLogValues(ctx, "eventKind", op.Action, "collection", collection.String(), "rkey", rkey.String()) - if len(CollectionFilter) > 0 { - keep := slices.Contains(CollectionFilter, collection.String()) - if strings.HasPrefix(collection.String(), "place.stream.") { - keep = true - } - if !keep { - continue - } + if !indexedCollection(collection.String()) { + continue } aqt, err := aqtime.FromString(evt.Time) @@ -495,6 +634,13 @@ func (atsync *ATProtoSynchronizer) handleCommitEventOps(ctx context.Context, evt log.Error(ctx, "failed to get repo", "err", err) continue } + if r != nil { + // One enrichment here gives every line under this op a handle to go + // with the DID -- the dispatch switch below, the contiguity checks, + // and everything reviveRepo and handleCreateUpdate log through ctx + // -- without any of them doing a lookup of their own. + ctx = log.WithLogValues(ctx, "handle", r.Handle) + } atsync.reviveRepo(ctx, r) // log.Warn(ctx, "got record we care about", "collection", collection, "rkey", rkey) @@ -719,11 +865,7 @@ func (atsync *ATProtoSynchronizer) handleCommitEventOps(ctx context.Context, evt } } - // Every op in this commit is indexed, so the index can claim this commit. - // Only reached on a clean pass: an event we bailed out of half-applied - // leaves the stored rev where it was, and the next commit for that repo - // notices the hole and orders a repair. - atsync.trackCommitRev(ctx, evt) + return true } // reviveRepo un-parks a repo we had written off. A commit event is proof the @@ -737,9 +879,9 @@ func (atsync *ATProtoSynchronizer) reviveRepo(ctx context.Context, r *model.Repo if !r.TerminalStatus() { return } - log.Log(ctx, "repo committed while parked, clearing terminal status", "did", r.DID, "status", r.Status) + log.Log(ctx, "repo committed while parked, clearing terminal status", "did", r.DID, "handle", r.Handle, "status", r.Status) if err := atsync.Model.SetRepoStatus(ctx, r.DID, model.RepoStatusOK); err != nil { - log.Error(ctx, "failed to clear repo status", "did", r.DID, "err", err) + log.Error(ctx, "failed to clear repo status", "did", r.DID, "handle", r.Handle, "err", err) return } r.Status = model.RepoStatusOK diff --git a/pkg/atproto/firehose_cursor.go b/pkg/atproto/firehose_cursor.go index 1cd3d347..b3dbed3e 100644 --- a/pkg/atproto/firehose_cursor.go +++ b/pkg/atproto/firehose_cursor.go @@ -47,6 +47,11 @@ func (h *highWater) set(seq int64) { h.v.Store(seq) } // a handful of in-flight frames just below it can be skipped on resume. That is // safe here: downstream handlers are idempotent, and with several relays plus a // cold deduper after restart, those commits get re-delivered and re-indexed. +// +// Alongside the sequence it tracks the time of the newest event seen, which is +// what makes the cursor's age knowable and so gives [relayCursor.dropIfStale] +// something to judge. A sequence number on its own says nothing about how much +// history resuming from it would ask the relay to re-send. type relayCursor struct { host string model model.Model @@ -63,8 +68,17 @@ type relayCursor struct { // assigns durable ids across its own restarts, so a stored cursor stays // valid (it just ages out of the relay's replay window if we are down too // long, which is the gap PDS re-sync covers). - latest highWater - flushed int64 // last persisted value; only the flush loop touches it + latest highWater + // lastEvent is the unix-seconds stamp of the newest event seen on this + // relay, high-watered the same way and for the same reason: event times + // jitter slightly out of order across a relay's own workers, and the newest + // one is the only one that says anything about how current we are. Unlike + // latest it is transport-independent — every relay stamps its events with a + // wall-clock time — so it is fed on both the WebSocket and moq paths. + lastEvent highWater + + flushed int64 // last persisted cursor; only the flush loop touches it + flushedEvent int64 // last persisted event time; likewise } func (atsync *ATProtoSynchronizer) newRelayCursor(ctx context.Context, host string) *relayCursor { @@ -77,7 +91,10 @@ func (atsync *ATProtoSynchronizer) newRelayCursor(ctx context.Context, host stri if stored != nil { rc.latest.set(stored.Cursor) rc.flushed = stored.Cursor - log.Log(ctx, "resuming relay from stored cursor", "cursor", stored.Cursor) + rc.lastEvent.set(stored.LastEventTime) + rc.flushedEvent = stored.LastEventTime + log.Log(ctx, "resuming relay from stored cursor", "cursor", stored.Cursor, + "lastEventTime", stored.LastEventTime) } return rc } @@ -88,6 +105,63 @@ func (rc *relayCursor) observe(seq int64) { rc.latest.observe(seq) } +// observeTime records the time stamped on an event we just received, which is +// what [relayCursor.dropIfStale] later measures the cursor's age against. Safe +// for concurrent callers; max-wins, so a slightly out-of-order stamp cannot +// walk the mark backwards. +func (rc *relayCursor) observeTime(t time.Time) { + rc.lastEvent.observe(t.Unix()) +} + +// lastEventTime is the unix-seconds stamp of the newest event seen so far, or 0 +// if we have never seen one (and never loaded one from the index). +func (rc *relayCursor) lastEventTime() int64 { return rc.lastEvent.get() } + +// stale reports whether resuming from this cursor would ask the relay to replay +// more than window of history. +// +// A zero cursor is never stale: there is nothing to replay from, so the connect +// tails live regardless. A nonzero cursor with NO event time is always stale — +// that is a row written before this column existed, and it is exactly the state +// that melted production: unknown age, unbounded replay. A non-positive window +// disables the whole check. +func (rc *relayCursor) stale(window time.Duration, now time.Time) bool { + if window <= 0 { + return false + } + if rc.latest.get() == 0 { + return false + } + last := rc.lastEvent.get() + if last == 0 { + return true + } + return now.Sub(time.Unix(last, 0)) > window +} + +// dropIfStale abandons the cursor when its last observed event is older than +// window, so the connect tails live instead of replaying a backlog the sweep +// can heal more cheaply. Returns true if it dropped. +// +// Called at the top of every connect attempt rather than once at load, because +// a relay also goes stale mid-process: a long disconnect-and-backoff stretch +// leaves a cursor that was fine when we loaded it pointing hours into the past +// by the time we get back in. +func (rc *relayCursor) dropIfStale(ctx context.Context, window time.Duration) bool { + if !rc.stale(window, time.Now()) { + return false + } + seq := rc.latest.get() + age := "unknown" + if last := rc.lastEvent.get(); last != 0 { + age = time.Since(time.Unix(last, 0)).String() + } + log.Warn(ctx, "abandoning stale relay cursor; tailing live and leaving the gap to the sweep", + "relay", rc.host, "cursor", seq, "lastEventAge", age, "replayWindow", window) + rc.reset() + return true +} + // param returns the cursor to dial with and whether to send one at all. With no // progress yet we send none, so a fresh external relay tails from live instead // of backfilling its entire history. @@ -123,23 +197,32 @@ func (rc *relayCursor) groupStart() (uint64, bool) { // reset abandons the cursor so the next connect tails the live edge. For when // a resume proves the stored value can't be trusted — e.g. a row poisoned by -// upstream at-seqs before the group/at-seq split in repoStreamCallbacks — -// which observe()'s max-wins semantics could otherwise never walk back. The -// flush loop persists the zero, healing the stored row too. +// upstream at-seqs before the group/at-seq split in repoStreamCallbacks, or one +// so old that replaying it would cost more than the sweep — which observe()'s +// max-wins semantics could otherwise never walk back. The flush loop persists +// the zeroes, healing the stored row too. +// +// The event time goes with it: keeping a fresh timestamp next to a zero cursor +// would describe a relay we are current with, which is the opposite of what a +// reset means. func (rc *relayCursor) reset() { rc.latest.set(0) + rc.lastEvent.set(0) } -// flush persists the high-water mark if it has advanced since the last write. -// Only ever called from the single flush goroutine, so flushed is unsynchronized. +// flush persists the high-water marks if either has advanced since the last +// write. Only ever called from the single flush goroutine, so the flushed +// values are unsynchronized. func (rc *relayCursor) flush(ctx context.Context) { v := rc.latest.get() - if v == rc.flushed { + t := rc.lastEvent.get() + if v == rc.flushed && t == rc.flushedEvent { return } - if err := rc.model.UpsertRelayCursor(rc.host, v); err != nil { - log.Error(ctx, "failed to persist relay cursor", "err", err, "cursor", v) + if err := rc.model.UpsertRelayCursor(rc.host, v, t); err != nil { + log.Error(ctx, "failed to persist relay cursor", "err", err, "cursor", v, "lastEventTime", t) return } rc.flushed = v + rc.flushedEvent = t } diff --git a/pkg/atproto/firehose_cursor_test.go b/pkg/atproto/firehose_cursor_test.go index 20e4ea61..4ac27f76 100644 --- a/pkg/atproto/firehose_cursor_test.go +++ b/pkg/atproto/firehose_cursor_test.go @@ -3,10 +3,12 @@ package atproto import ( "context" "testing" + "time" indigoatproto "github.com/bluesky-social/indigo/api/atproto" lexutil "github.com/bluesky-social/indigo/lex/util" "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/model" ) @@ -166,3 +168,186 @@ func TestRelayCursorReset(t *testing.T) { t.Fatal("restarted cursor should tail live after reset") } } + +// TestRelayCursorEventTime covers the second half of the cursor: the newest +// event time is what makes the cursor's age knowable, so it has to high-water, +// persist and reload alongside the sequence. +func TestRelayCursorEventTime(t *testing.T) { + mod, err := model.MakeDB(":memory:") + require.NoError(t, err) + atsync := &ATProtoSynchronizer{Model: mod} + ctx := context.Background() + const host = "wss://relay.example" + + rc := atsync.newRelayCursor(ctx, host) + require.Equal(t, int64(0), rc.lastEventTime(), "a fresh relay has seen no events") + + // Max-wins, like the seq: relays stamp events from several workers, so a + // slightly older stamp arriving late must not walk the mark backwards. + base := time.Unix(1700000000, 0) + rc.observeTime(base) + rc.observeTime(base.Add(-30 * time.Second)) + rc.observeTime(base.Add(90 * time.Second)) + require.Equal(t, base.Add(90*time.Second).Unix(), rc.lastEventTime()) + + // flush writes both columns, and a flush with nothing new is a no-op. + rc.observe(120) + rc.flush(ctx) + stored, err := mod.GetRelayCursor(host) + require.NoError(t, err) + require.NotNil(t, stored) + require.Equal(t, int64(120), stored.Cursor) + require.Equal(t, base.Add(90*time.Second).Unix(), stored.LastEventTime) + + // The event time advancing on its own is enough to trigger a write -- a moq + // relay whose group cursor is unchanged still gets its freshness recorded. + rc.observeTime(base.Add(5 * time.Minute)) + rc.flush(ctx) + stored, err = mod.GetRelayCursor(host) + require.NoError(t, err) + require.Equal(t, int64(120), stored.Cursor) + require.Equal(t, base.Add(5*time.Minute).Unix(), stored.LastEventTime) + + // A restart reloads both, so the staleness check has something to judge. + resumed := atsync.newRelayCursor(ctx, host) + v, ok := resumed.param() + require.True(t, ok) + require.Equal(t, int64(120), v) + require.Equal(t, base.Add(5*time.Minute).Unix(), resumed.lastEventTime()) + + // reset drops the event time with the cursor: a fresh timestamp next to a + // zero cursor would claim we are current with a relay we just abandoned. + resumed.reset() + require.Equal(t, int64(0), resumed.lastEventTime()) + _, ok = resumed.param() + require.False(t, ok) +} + +// TestRelayCursorDropIfStale is the production incident in a unit test: a +// cursor whose newest event is hours old must not be replayed from, because +// asking a relay to re-send hours of the whole network is the thing that buried +// the node. Fresh cursors are left alone, since replay is genuinely cheaper +// across an ordinary restart. +func TestRelayCursorDropIfStale(t *testing.T) { + ctx := context.Background() + const window = 15 * time.Minute + + newCursor := func(seq int64, lastEvent time.Time) *relayCursor { + rc := &relayCursor{host: "wss://relay.example"} + rc.latest.set(seq) + if !lastEvent.IsZero() { + rc.observeTime(lastEvent) + } + return rc + } + + // Inside the window: keep the cursor and replay from it. + rc := newCursor(120, time.Now().Add(-time.Minute)) + require.False(t, rc.dropIfStale(ctx, window)) + v, ok := rc.param() + require.True(t, ok) + require.Equal(t, int64(120), v) + + // Older than the window: drop it and tail live, leaving the gap for the + // sweep's head check. + rc = newCursor(120, time.Now().Add(-2*time.Hour)) + require.True(t, rc.dropIfStale(ctx, window)) + _, ok = rc.param() + require.False(t, ok, "a dropped cursor must not be sent to the relay") + require.Equal(t, int64(0), rc.lastEventTime()) + + // A nonzero cursor with NO event time is a row written before the column + // existed -- exactly the state production was poisoned with. Unknown age + // means unbounded replay, so it counts as stale. + rc = newCursor(120, time.Time{}) + require.Equal(t, int64(0), rc.lastEventTime()) + require.True(t, rc.dropIfStale(ctx, window)) + _, ok = rc.param() + require.False(t, ok) + + // A zero cursor is already tailing live; there is nothing to drop, so + // nothing is reported (and no sweep gets kicked). + rc = newCursor(0, time.Time{}) + require.False(t, rc.dropIfStale(ctx, window)) + + // Window 0 disables the cap entirely: however old, replay from it. + rc = newCursor(120, time.Now().Add(-72*time.Hour)) + require.False(t, rc.dropIfStale(ctx, 0)) + v, ok = rc.param() + require.True(t, ok) + require.Equal(t, int64(120), v) + + // The moq flavour of the same cursor is dropped on the same terms, so + // SubscribeFrom asks for the live edge rather than a two-hour replay. + rc = &relayCursor{host: "moqt://relay.example"} + rc.observeGroup(640) + rc.observeTime(time.Now().Add(-2 * time.Hour)) + require.True(t, rc.dropIfStale(ctx, window)) + _, ok = rc.groupStart() + require.False(t, ok) +} + +// TestFirehoseReplayWindow pins the flag plumbing: unset CLI means the default, +// and an explicit zero means "never cap", which is how the escape hatch is +// spelled everywhere else in the sweep config. +func TestFirehoseReplayWindow(t *testing.T) { + require.Equal(t, config.DefaultFirehoseReplayWindow, + (&ATProtoSynchronizer{}).firehoseReplayWindow()) + require.Equal(t, 30*time.Second, + (&ATProtoSynchronizer{CLI: &config.CLI{FirehoseReplayWindow: 30 * time.Second}}).firehoseReplayWindow()) + require.Equal(t, time.Duration(0), + (&ATProtoSynchronizer{CLI: &config.CLI{}}).firehoseReplayWindow()) +} + +// TestRelayCursorObserveTimeFromCallbacks checks the wiring the staleness cap +// depends on: event times are recorded for both commit and identity events, on +// every transport (unlike the seq, which stays out of a moq relay's group +// cursor), and BEFORE dedup -- a duplicate is still evidence of how current +// this relay is. +func TestRelayCursorObserveTimeFromCallbacks(t *testing.T) { + atsync := &ATProtoSynchronizer{ + commitDedup: newFirehoseDeduper(dedupWindow), + identityDedup: newFirehoseDeduper(dedupWindow), + } + evtTime := time.Now().UTC().Add(-time.Minute).Format(time.RFC3339Nano) + commit := &indigoatproto.SyncSubscribeRepos_Commit{ + Seq: 5_000_000_000, + Repo: "did:plc:abc123", + Time: evtTime, + Commit: lexutil.LexLink(mustCID(t)), + } + // Pre-claim it so the callback stops after its bookkeeping instead of + // running the indexing handler, which would want a whole model. + atsync.commitDedup.seen(commit.Commit.String()) + + want, err := time.Parse(time.RFC3339Nano, evtTime) + require.NoError(t, err) + + moqCursor := &relayCursor{host: "moqt://relay.example"} + rsc := atsync.repoStreamCallbacks(context.Background(), "moqt://relay.example", moqCursor, func() {}) + require.NoError(t, rsc.RepoCommit(commit)) + require.Equal(t, want.Unix(), moqCursor.lastEventTime(), + "event time is transport-independent, unlike the seq") + + identity := &indigoatproto.SyncSubscribeRepos_Identity{ + Seq: 5_000_000_001, + Did: "did:plc:abc123", + Time: evtTime, + } + atsync.identityDedup.seen(identityDedupKey(identity)) + wsCursor := &relayCursor{host: "wss://relay.example"} + rsc = atsync.repoStreamCallbacks(context.Background(), "wss://relay.example", wsCursor, func() {}) + require.NoError(t, rsc.RepoIdentity(identity)) + require.Equal(t, want.Unix(), wsCursor.lastEventTime()) + + // An unparsable stamp is simply not an observation, never a crash. + wsCursor2 := &relayCursor{host: "wss://relay.example"} + rsc = atsync.repoStreamCallbacks(context.Background(), "wss://relay.example", wsCursor2, func() {}) + require.NoError(t, rsc.RepoCommit(&indigoatproto.SyncSubscribeRepos_Commit{ + Seq: 7, + Repo: "did:plc:abc123", + Time: "not a timestamp", + Commit: lexutil.LexLink(mustCID(t)), + })) + require.Equal(t, int64(0), wsCursor2.lastEventTime()) +} diff --git a/pkg/atproto/firehose_filter_test.go b/pkg/atproto/firehose_filter_test.go new file mode 100644 index 00000000..acc926a6 --- /dev/null +++ b/pkg/atproto/firehose_filter_test.go @@ -0,0 +1,121 @@ +package atproto + +import ( + "context" + "testing" + + indigoatproto "github.com/bluesky-social/indigo/api/atproto" + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/constants" +) + +func repoOp(action, path string) *indigoatproto.SyncSubscribeRepos_RepoOp { + return &indigoatproto.SyncSubscribeRepos_RepoOp{Action: action, Path: path} +} + +// TestIndexedCollection pins the keep/skip decision that now runs before the +// CAR parse. It has to answer exactly what the per-op filter answered, because +// it is the same function -- a divergence here would silently stop indexing a +// collection. +func TestIndexedCollection(t *testing.T) { + require.True(t, indexedCollection(constants.APP_BSKY_FEED_POST), "a member of CollectionFilter") + require.True(t, indexedCollection(constants.APP_BSKY_GRAPH_FOLLOW)) + require.True(t, indexedCollection(constants.PLACE_STREAM_LIVE_RECOMMENDATIONS)) + + // Everything under our own namespace is ours, listed or not: the filter + // names only the app.bsky.* collections it wants. + require.True(t, indexedCollection(constants.PLACE_STREAM_CHAT_MESSAGE)) + require.True(t, indexedCollection(constants.PLACE_STREAM_VIDEO)) + require.True(t, indexedCollection("place.stream.something.invented.tomorrow")) + + // The majority of the firehose, which we index nothing from. + require.False(t, indexedCollection("app.bsky.feed.like")) + require.False(t, indexedCollection("app.bsky.feed.repost")) + require.False(t, indexedCollection("com.whtwnd.blog.entry")) + require.False(t, indexedCollection("placestream.notours")) + require.False(t, indexedCollection("")) +} + +// TestCommitHasIndexedOps is the hoist itself: from op paths alone, without +// touching the CAR, does this commit contain anything for us? +func TestCommitHasIndexedOps(t *testing.T) { + boring := &indigoatproto.SyncSubscribeRepos_Commit{ + Ops: []*indigoatproto.SyncSubscribeRepos_RepoOp{ + repoOp("create", "app.bsky.feed.like/3lpabc"), + repoOp("delete", "app.bsky.feed.repost/3lpdef"), + }, + } + require.False(t, commitHasIndexedOps(boring)) + + // One interesting op among many boring ones is enough -- the CAR has to be + // read for the commit as a whole. + mixed := &indigoatproto.SyncSubscribeRepos_Commit{ + Ops: []*indigoatproto.SyncSubscribeRepos_RepoOp{ + repoOp("create", "app.bsky.feed.like/3lpabc"), + repoOp("create", constants.PLACE_STREAM_CHAT_MESSAGE+"/3lpghi"), + }, + } + require.True(t, commitHasIndexedOps(mixed)) + + // A commit with no ops at all has nothing for us. + require.False(t, commitHasIndexedOps(&indigoatproto.SyncSubscribeRepos_Commit{})) + + // An unparsable path counts as interesting, so the main loop reaches it and + // logs the error it always logged rather than dropping it in silence. + require.True(t, commitHasIndexedOps(&indigoatproto.SyncSubscribeRepos_Commit{ + Ops: []*indigoatproto.SyncSubscribeRepos_RepoOp{repoOp("create", "nopath")}, + })) +} + +// TestHandleCommitEventOpsTracksFilteredCommits is the subtle half of the +// hoist, and the reason the early return sits where it does. +// +// A repo we track writes plenty of records we index nothing from. Those commits +// still carry the rev chain, so trackCommitRev has to see them -- otherwise our +// stored rev falls behind, and the next commit we DO care about arrives with a +// Since we have never heard of, looks like a firehose gap, and orders a repair +// of a repo that was never damaged. trackCommitRev needs no CAR, so it runs on +// the far side of the skip. +// +// The empty Blocks here is the load-bearing part of the fixture: an empty CAR +// cannot be parsed, so a test that passes proves the parse was skipped. +func TestHandleCommitEventOpsTracksFilteredCommits(t *testing.T) { + ctx := context.Background() + atsync, mod := contiguityTestSync(t) + require.NoError(t, mod.UpdateRepo(syncedRepo("did:plc:filtered", "3lprev0000000"))) + + evt := commitEvent("did:plc:filtered", "3lprev0000000", "3lprev0000001") + evt.Ops = []*indigoatproto.SyncSubscribeRepos_RepoOp{ + repoOp("create", "app.bsky.feed.like/3lpabc"), + } + atsync.handleCommitEventOps(ctx, evt) + + got, err := mod.GetRepo("did:plc:filtered") + require.NoError(t, err) + require.NotNil(t, got) + require.Equal(t, "3lprev0000001", got.Version, + "a commit whose ops are all filtered still advances the rev chain") + require.Empty(t, got.RepairFrom, "and orders no repair") + + // A commit that DOES carry something we index cannot skip the CAR, and an + // unreadable one has to leave the rev where it was -- half-applied is the + // one state the chain must never record as clean. + evt = commitEvent("did:plc:filtered", "3lprev0000001", "3lprev0000002") + evt.Ops = []*indigoatproto.SyncSubscribeRepos_RepoOp{ + repoOp("create", constants.PLACE_STREAM_CHAT_MESSAGE+"/3lpghi"), + } + atsync.handleCommitEventOps(ctx, evt) + + got, err = mod.GetRepo("did:plc:filtered") + require.NoError(t, err) + require.Equal(t, "3lprev0000001", got.Version, + "an event that bailed out before indexing its ops must not claim the commit") + + // tooBig is unchanged: skipped entirely, rev included. + evt = commitEvent("did:plc:filtered", "3lprev0000001", "3lprev0000003") + evt.TooBig = true + atsync.handleCommitEventOps(ctx, evt) + got, err = mod.GetRepo("did:plc:filtered") + require.NoError(t, err) + require.Equal(t, "3lprev0000001", got.Version) +} diff --git a/pkg/atproto/firehose_moq.go b/pkg/atproto/firehose_moq.go index 3b5265d9..fcb8e741 100644 --- a/pkg/atproto/firehose_moq.go +++ b/pkg/atproto/firehose_moq.go @@ -48,6 +48,10 @@ func (atsync *ATProtoSynchronizer) connectRelayMoq(ctx context.Context, relay st // first connect. If that group has aged out of the relay's replay window the // relay jumps forward to the oldest it still retains, leaving a gap we accept // here (deep recovery is a PDS re-sync, tracked separately). + // + // A cursor too old to be worth replaying at all has already been dropped: + // connectRelay runs the staleness check before it dispatches here, and it is + // the only caller, so groupStart below already reflects that decision. var sub *atmoq.Subscription resumedFrom, resumed := cursor.groupStart() if resumed { diff --git a/pkg/config/config.go b/pkg/config/config.go index f549d8de..9466fd7f 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -172,6 +172,7 @@ type CLI struct { MaximumLiveBitrate int SweepConcurrency int SweepInterval time.Duration + FirehoseReplayWindow time.Duration } // DefaultSweepInterval is how often the atproto sweep re-runs when @@ -194,6 +195,21 @@ const DefaultSweepInterval = 6 * time.Hour // requests per second) against any single PDS. const DefaultSweepConcurrency = 32 +// DefaultFirehoseReplayWindow is how stale a stored relay cursor may be before +// this node stops trying to replay from it and tails the live edge instead. +// +// The firehose is a latency optimization, not the sync engine: the sweep's head +// check asks every repo's host one question and repairs the ones that have +// drifted, so a gap of hours costs a few thousand cheap requests spread across +// hundreds of hosts. Replaying that same gap costs the relay a full-rate flood +// of every commit on the network -- including the overwhelming majority from +// repos this node has never heard of -- which is how a two-hour-old cursor once +// buried a node under millions of queued events. Fifteen minutes is long enough +// to cover an ordinary restart or deploy, where replay genuinely is the cheaper +// answer, and short enough that anything worse is handed to the mechanism built +// for it. 0 disables the cap and always replays from the stored cursor. +const DefaultFirehoseReplayWindow = 15 * time.Minute + // ContentFilters represents the content filtering configuration type ContentFilters struct { ContentWarnings struct { @@ -848,6 +864,13 @@ func (cli *CLI) NewCommand(name string) *urfavecli.Command { Destination: &cli.SweepInterval, Sources: urfavecli.EnvVars("SP_SWEEP_INTERVAL"), }, + &urfavecli.DurationFlag{ + Name: "firehose-replay-window", + Usage: "how old a stored relay cursor may be and still be replayed from on connect. A cursor whose newest event is older than this is discarded and we tail the relay's live edge instead, leaving the gap for the sweep's head check to repair -- which is far cheaper than making the relay re-send every commit on the network. 0 always replays from the stored cursor", + Value: DefaultFirehoseReplayWindow, + Destination: &cli.FirehoseReplayWindow, + Sources: urfavecli.EnvVars("SP_FIREHOSE_REPLAY_WINDOW"), + }, &urfavecli.StringFlag{ Name: "maximum-live-bitrate", Usage: "maximum allowed live ingest bitrate, measured per emitted segment. Accepts a bits-per-second number or a decimal SI suffix — e.g. 30M, 30000k, or 30000000 (all 30 Mbps). A stream whose bitrate exceeds this (plus a 10% margin) is disconnected and the streamer is shown a problem. 0 = unlimited", diff --git a/pkg/model/model.go b/pkg/model/model.go index 46aa90d7..f381d7fd 100644 --- a/pkg/model/model.go +++ b/pkg/model/model.go @@ -109,7 +109,7 @@ type Model interface { UpdateLabelerCursor(did string, cursor int64) error GetRelayCursor(host string) (*RelayCursor, error) - UpsertRelayCursor(host string, cursor int64) error + UpsertRelayCursor(host string, cursor int64, lastEventTime int64) error CreateLabel(label *Label) error GetActiveLabels(uri string) ([]*comatproto.LabelDefs_Label, error) diff --git a/pkg/model/relay_cursor.go b/pkg/model/relay_cursor.go index c139d356..97331c0c 100644 --- a/pkg/model/relay_cursor.go +++ b/pkg/model/relay_cursor.go @@ -20,6 +20,11 @@ import ( type RelayCursor struct { Host string `gorm:"primaryKey;column:host"` Cursor int64 `gorm:"column:cursor"` + // LastEventTime is the unix-seconds timestamp of the newest firehose event + // observed on this relay, persisted alongside Cursor. It is what lets a + // restart decide whether the stored cursor is recent enough to replay from, + // without knowing anything about the relay's sequence numbering. + LastEventTime int64 `gorm:"column:last_event_time"` } // GetRelayCursor returns the stored cursor for a relay, or nil if we have never @@ -37,10 +42,13 @@ func (m *DBModel) GetRelayCursor(host string) (*RelayCursor, error) { } // UpsertRelayCursor stores the latest consumed cursor for a relay (the at-seq -// for WebSocket relays, the high-water MoQ group for moqt:// relays). -func (m *DBModel) UpsertRelayCursor(host string, cursor int64) error { +// for WebSocket relays, the high-water MoQ group for moqt:// relays) together +// with the timestamp of the newest event seen there, which is what makes the +// cursor's age -- and so whether it is worth replaying from at all -- knowable +// after a restart. +func (m *DBModel) UpsertRelayCursor(host string, cursor int64, lastEventTime int64) error { return m.DB.Clauses(clause.OnConflict{ Columns: []clause.Column{{Name: "host"}}, - DoUpdates: clause.AssignmentColumns([]string{"cursor"}), - }).Create(&RelayCursor{Host: host, Cursor: cursor}).Error + DoUpdates: clause.AssignmentColumns([]string{"cursor", "last_event_time"}), + }).Create(&RelayCursor{Host: host, Cursor: cursor, LastEventTime: lastEventTime}).Error } diff --git a/pkg/model/relay_cursor_test.go b/pkg/model/relay_cursor_test.go new file mode 100644 index 00000000..452a5496 --- /dev/null +++ b/pkg/model/relay_cursor_test.go @@ -0,0 +1,57 @@ +package model + +import ( + "testing" + + "github.com/stretchr/testify/require" +) + +// TestUpsertRelayCursor is the storage half of the replay-window cap: the +// cursor is only safe to resume from if we also know how old it is, so both +// columns have to survive the round trip and both have to move on an upsert. +func TestUpsertRelayCursor(t *testing.T) { + db := indexedTestDB(t) + const host = "wss://relay.example" + + // Nothing recorded yet reads as "fresh subscription", not as an error. + stored, err := db.GetRelayCursor(host) + require.NoError(t, err) + require.Nil(t, stored) + + require.NoError(t, db.UpsertRelayCursor(host, 120, 1700000000)) + stored, err = db.GetRelayCursor(host) + require.NoError(t, err) + require.NotNil(t, stored) + require.Equal(t, host, stored.Host) + require.Equal(t, int64(120), stored.Cursor) + require.Equal(t, int64(1700000000), stored.LastEventTime) + + // An upsert overwrites both columns rather than inserting a second row -- + // the event time especially, since a stale one next to a fresh cursor is + // exactly the state the staleness check exists to catch. + require.NoError(t, db.UpsertRelayCursor(host, 500, 1700009999)) + stored, err = db.GetRelayCursor(host) + require.NoError(t, err) + require.NotNil(t, stored) + require.Equal(t, int64(500), stored.Cursor) + require.Equal(t, int64(1700009999), stored.LastEventTime) + require.Equal(t, int64(1), countRows(t, db, &RelayCursor{})) + + // A reset writes zeroes to both, and that has to persist as written rather + // than being skipped as a zero value. + require.NoError(t, db.UpsertRelayCursor(host, 0, 0)) + stored, err = db.GetRelayCursor(host) + require.NoError(t, err) + require.NotNil(t, stored) + require.Equal(t, int64(0), stored.Cursor) + require.Equal(t, int64(0), stored.LastEventTime) + + // Cursors are per-relay. + require.NoError(t, db.UpsertRelayCursor("moqt://other.example", 640, 1700000042)) + stored, err = db.GetRelayCursor("moqt://other.example") + require.NoError(t, err) + require.NotNil(t, stored) + require.Equal(t, int64(640), stored.Cursor) + require.Equal(t, int64(1700000042), stored.LastEventTime) + require.Equal(t, int64(2), countRows(t, db, &RelayCursor{})) +}