diff --git a/pkg/atproto/firehose_moq_test.go b/pkg/atproto/firehose_moq_test.go index 4e2225f21..bcf20e89b 100644 --- a/pkg/atproto/firehose_moq_test.go +++ b/pkg/atproto/firehose_moq_test.go @@ -14,6 +14,7 @@ import ( "github.com/ipfs/go-cid" atmoq "github.com/streamplace/atmoq-go" "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/model" ) // a real CIDv1 to populate the required commit/identity link fields so CBOR @@ -159,3 +160,77 @@ func TestFirehoseMoqLive(t *testing.T) { } t.Logf("decoded %d live commit frames from %s", commits, relay) } + +// TestRelayCursorPersistResumeLive exercises the whole cross-restart resume path +// against a real windowed atmoq relay: tail live, persist the group cursor to a +// file-backed DB, simulate a Streamplace restart by reopening the DB with a new +// model, and confirm the resumed cursor replays from the persisted group while a +// fresh live subscribe starts past it. Network-gated: set +// SP_MOQ_REPLAY_TEST=moqt://host to a relay run with --replay-window-secs. +func TestRelayCursorPersistResumeLive(t *testing.T) { + relay := os.Getenv("SP_MOQ_REPLAY_TEST") + if relay == "" { + t.Skip("set SP_MOQ_REPLAY_TEST=moqt://host (a windowed atmoq relay) to run") + } + ctx, cancel := context.WithTimeout(context.Background(), 40*time.Second) + defer cancel() + + dbPath := t.TempDir() + "/cursor.db" + mod, err := model.MakeDB(dbPath) + require.NoError(t, err) + atsync := &ATProtoSynchronizer{Model: mod} + + // Phase 1: fresh cursor tails live; record + persist the high-water group. + rc := atsync.newRelayCursor(ctx, relay) + _, ok := rc.groupStart() + require.False(t, ok, "fresh cursor should tail the live edge") + + sess, err := atmoq.Dial(ctx, relay, &atmoq.Options{}) + require.NoError(t, err) + sub, err := sess.Subscribe(ctx, atmoq.DefaultBroadcast, atmoq.DefaultTrack) + require.NoError(t, err) + for i := 0; i < 80; i++ { + _, g, err := sub.ReadFrame(ctx) + require.NoError(t, err) + rc.observeGroup(g) + } + persisted, ok := rc.groupStart() + require.True(t, ok) + sub.Close() + sess.Close() + rc.flush(ctx) // write the cursor to the DB + t.Logf("phase1: persisted group cursor G=%d", persisted) + + // Let the live edge advance past the persisted group. + time.Sleep(5 * time.Second) + + // Phase 2: "restart" — a new model over the same DB file reloads the cursor. + mod2, err := model.MakeDB(dbPath) + require.NoError(t, err) + resumed := (&ATProtoSynchronizer{Model: mod2}).newRelayCursor(ctx, relay) + rg, ok := resumed.groupStart() + require.True(t, ok, "a restart should resume the stored group") + require.Equal(t, persisted, rg, "resumed group must match the persisted one") + + // Separate sessions (avoid two subs on one): resume vs live. + resumeSess, err := atmoq.Dial(ctx, relay, &atmoq.Options{}) + require.NoError(t, err) + defer resumeSess.Close() + rsub, err := resumeSess.SubscribeFrom(ctx, atmoq.DefaultBroadcast, atmoq.DefaultTrack, rg) + require.NoError(t, err) + _, resumeFirst, err := rsub.ReadFrame(ctx) + require.NoError(t, err) + + liveSess, err := atmoq.Dial(ctx, relay, &atmoq.Options{}) + require.NoError(t, err) + defer liveSess.Close() + lsub, err := liveSess.Subscribe(ctx, atmoq.DefaultBroadcast, atmoq.DefaultTrack) + require.NoError(t, err) + _, liveFirst, err := lsub.ReadFrame(ctx) + require.NoError(t, err) + + t.Logf("phase2: resumed from persisted G=%d -> first group %d ; live -> first group %d", + rg, resumeFirst, liveFirst) + require.LessOrEqual(t, resumeFirst, rg, "resume should replay at/before the persisted group") + require.Greater(t, liveFirst, rg, "live edge should have advanced past the persisted group") +}