diff --git a/pkg/atproto/firehose_cursor.go b/pkg/atproto/firehose_cursor.go index 2cd013ffd..d76a82a7f 100644 --- a/pkg/atproto/firehose_cursor.go +++ b/pkg/atproto/firehose_cursor.go @@ -29,24 +29,20 @@ type relayCursor struct { host string model model.Model - latest atomic.Int64 // highest seq seen; 0 = nothing yet (tail from live) - flushed int64 // last persisted value; only the flush loop touches it - - // group is the highest MoQ group sequence seen on a moqt:// relay, used to - // resume replay (see connectRelayMoq). -1 = none yet, tail from the live - // edge. Persisted alongside the seq cursor, so a Streamplace restart resumes - // from the last group too — the relay assigns durable group ids across its - // own restarts, so a stored group 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). - group atomic.Int64 - flushedGroup int64 // last persisted group; only the flush loop touches it + // latest is the high-water cursor: the upstream at-sequence for a WebSocket + // relay, or the high-water MoQ group sequence for a moqt:// relay (used to + // resume replay via SubscribeFrom — see connectRelayMoq). A host is one + // transport or the other, so a single value covers both. 0 = nothing seen + // yet (tail from live). Persisted periodically so a Streamplace restart + // resumes from here — the relay 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 atomic.Int64 + flushed int64 // last persisted value; only the flush loop touches it } func (atsync *ATProtoSynchronizer) newRelayCursor(ctx context.Context, host string) *relayCursor { rc := &relayCursor{host: host, model: atsync.Model} - rc.group.Store(-1) // -1 = no MoQ group seen yet (tail from live edge) - rc.flushedGroup = -1 stored, err := atsync.Model.GetRelayCursor(host) if err != nil { log.Error(ctx, "failed to load relay cursor; tailing from live", "err", err) @@ -55,11 +51,7 @@ func (atsync *ATProtoSynchronizer) newRelayCursor(ctx context.Context, host stri if stored != nil { rc.latest.Store(stored.Cursor) rc.flushed = stored.Cursor - if stored.GroupSeq != nil { - rc.group.Store(*stored.GroupSeq) - rc.flushedGroup = *stored.GroupSeq - } - log.Log(ctx, "resuming relay from stored cursor", "cursor", stored.Cursor, "group", rc.group.Load()) + log.Log(ctx, "resuming relay from stored cursor", "cursor", stored.Cursor) } return rc } @@ -89,28 +81,23 @@ func (rc *relayCursor) param() (int64, bool) { // highSeq returns the high-water upstream sequence number observed so far. func (rc *relayCursor) highSeq() int64 { return rc.latest.Load() } -// observeGroup advances the high-water MoQ group sequence. Called on every -// frame received from a moqt:// relay (concurrency-safe). +// observeGroup advances the high-water cursor from a MoQ group sequence (the +// moqt:// transport's flavour of a cursor). Called on every frame received from +// a moqt:// relay (concurrency-safe). func (rc *relayCursor) observeGroup(seq uint64) { - s := int64(seq) - for { - cur := rc.group.Load() - if s <= cur { - return - } - if rc.group.CompareAndSwap(cur, s) { - return - } - } + rc.observe(int64(seq)) } // groupStart returns the MoQ group to resume replay from and whether to request // replay at all. Before any frame is seen we tail the live edge; after a // reconnect we resume from the last group seen so the relay replays from there -// (already-seen frames in that group are deduped downstream). +// (already-seen frames in that group are deduped downstream). Group 0 (a +// brand-new relay's very first group) reads as "none" and tails live — harmless +// and self-healing, and in practice relay group ids are large (seeded for +// durability across restarts). func (rc *relayCursor) groupStart() (uint64, bool) { - v := rc.group.Load() - if v < 0 { + v := rc.latest.Load() + if v <= 0 { return 0, false } return uint64(v), true @@ -120,18 +107,12 @@ func (rc *relayCursor) groupStart() (uint64, bool) { // Only ever called from the single flush goroutine, so flushed is unsynchronized. func (rc *relayCursor) flush(ctx context.Context) { v := rc.latest.Load() - g := rc.group.Load() - if v == rc.flushed && g == rc.flushedGroup { + if v == rc.flushed { return } - var groupPtr *int64 - if g >= 0 { - groupPtr = &g - } - if err := rc.model.UpsertRelayCursor(rc.host, v, groupPtr); err != nil { - log.Error(ctx, "failed to persist relay cursor", "err", err, "cursor", v, "group", g) + if err := rc.model.UpsertRelayCursor(rc.host, v); err != nil { + log.Error(ctx, "failed to persist relay cursor", "err", err, "cursor", v) return } rc.flushed = v - rc.flushedGroup = g } diff --git a/pkg/atproto/firehose_cursor_test.go b/pkg/atproto/firehose_cursor_test.go index aefe6f408..2cf7b8ab0 100644 --- a/pkg/atproto/firehose_cursor_test.go +++ b/pkg/atproto/firehose_cursor_test.go @@ -81,15 +81,12 @@ func TestRelayCursorGroupResume(t *testing.T) { require.True(t, ok) require.Equal(t, uint64(640), g) - // MoQ frames also carry the at-seq, so both advance; flush persists both. - rc.observe(9000) + // flush persists the high-water group as the relay's cursor. rc.flush(ctx) stored, err := mod.GetRelayCursor(host) require.NoError(t, err) require.NotNil(t, stored) - require.Equal(t, int64(9000), stored.Cursor) - require.NotNil(t, stored.GroupSeq) - require.Equal(t, int64(640), *stored.GroupSeq) + require.Equal(t, int64(640), stored.Cursor) // As if the process restarted: a new cursor resumes replay from the stored // group (connectRelayMoq calls SubscribeFrom with it). @@ -97,15 +94,4 @@ func TestRelayCursorGroupResume(t *testing.T) { g, ok = resumed.groupStart() require.True(t, ok) require.Equal(t, uint64(640), g) - - // A WebSocket relay never sets a group, so its stored group stays NULL (it - // resumes by sequence number instead). - const wsHost = "wss://relay.example" - ws := atsync.newRelayCursor(ctx, wsHost) - ws.observe(123) - ws.flush(ctx) - wsStored, err := mod.GetRelayCursor(wsHost) - require.NoError(t, err) - require.NotNil(t, wsStored) - require.Nil(t, wsStored.GroupSeq) } diff --git a/pkg/atproto/firehose_moq.go b/pkg/atproto/firehose_moq.go index 7b42bc133..507d88a13 100644 --- a/pkg/atproto/firehose_moq.go +++ b/pkg/atproto/firehose_moq.go @@ -20,10 +20,11 @@ import ( // WebSocket messages, so each frame decodes through the same indigo event types // and feeds the same scheduler + callbacks as connectRelay's WebSocket path. // -// MoQ has no cursor/replay (subscriptions always start at the publisher's latest -// group), so the stored cursor is observed for liveness but never used to -// resume. Gaps across a reconnect are covered by cross-relay redelivery and the -// idempotent handlers, exactly as for the best-effort WebSocket cursor. +// On a reconnect or restart it resumes replay from the last MoQ group it saw +// (via SubscribeFrom, served from the relay's replay window); on the first +// connect it tails the live edge. Gaps past the relay's retained window are +// covered by cross-relay redelivery and the idempotent handlers, exactly as for +// the best-effort WebSocket cursor. func (atsync *ATProtoSynchronizer) connectRelayMoq(ctx context.Context, relay string, cursor *relayCursor) error { streamCtx, cancel := context.WithCancel(ctx) defer cancel() diff --git a/pkg/model/model.go b/pkg/model/model.go index 3e801efb8..91f80686b 100644 --- a/pkg/model/model.go +++ b/pkg/model/model.go @@ -104,7 +104,7 @@ type Model interface { UpdateLabelerCursor(did string, cursor int64) error GetRelayCursor(host string) (*RelayCursor, error) - UpsertRelayCursor(host string, cursor int64, group *int64) error + UpsertRelayCursor(host string, cursor 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 e4b916f30..c139d356d 100644 --- a/pkg/model/relay_cursor.go +++ b/pkg/model/relay_cursor.go @@ -8,18 +8,18 @@ import ( ) // RelayCursor remembers how far we have consumed each relay's firehose, keyed by -// the relay's websocket URL. On reconnect or restart we resume from the stored -// sequence number instead of re-tailing from live (which would leave a gap) or -// replaying from the beginning. Cursors are per-relay because each relay -// assigns its own sequence numbers. +// the relay's URL. On reconnect or restart we resume from the stored cursor +// instead of re-tailing from live (which would leave a gap) or replaying from +// the beginning. Cursors are per-relay because each relay assigns its own +// numbering. +// +// Cursor is the at-sequence for WebSocket relays and the high-water MoQ group +// sequence for moqt:// relays (resumed via SubscribeFrom). A host is one +// transport or the other, so a single opaque monotonic int64 we hand back to +// the relay covers both — we just call a moqt:// group id a "cursor" too. type RelayCursor struct { Host string `gorm:"primaryKey;column:host"` Cursor int64 `gorm:"column:cursor"` - // GroupSeq is the high-water MoQ group sequence for moqt:// relays, used to - // resume replay via SubscribeFrom after a restart. NULL for WebSocket relays, - // which resume by the Cursor sequence number instead. (Column is not named - // "group" because that is a reserved SQL keyword.) - GroupSeq *int64 `gorm:"column:group_seq"` } // GetRelayCursor returns the stored cursor for a relay, or nil if we have never @@ -36,12 +36,11 @@ func (m *DBModel) GetRelayCursor(host string) (*RelayCursor, error) { return &rc, nil } -// UpsertRelayCursor stores the latest consumed sequence number for a relay, and -// (for moqt:// relays) the high-water MoQ group sequence; pass group=nil for -// WebSocket relays. -func (m *DBModel) UpsertRelayCursor(host string, cursor int64, group *int64) 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 { return m.DB.Clauses(clause.OnConflict{ Columns: []clause.Column{{Name: "host"}}, - DoUpdates: clause.AssignmentColumns([]string{"cursor", "group_seq"}), - }).Create(&RelayCursor{Host: host, Cursor: cursor, GroupSeq: group}).Error + DoUpdates: clause.AssignmentColumns([]string{"cursor"}), + }).Create(&RelayCursor{Host: host, Cursor: cursor}).Error }