package oracle // This file is compiled into the pinned upstream internal/oracle package via // tests/differential_oracle.py's Go overlay. It is intentionally a real // process/socket/disk test: only the deterministic atproto network is local. import ( "bytes" "context" "encoding/json" "fmt" "io" "log/slog" "math/rand/v2" "net" "net/http" "net/http/httptest" "os" "os/exec" "path/filepath" "slices" "strconv" "strings" "sync" "sync/atomic" "syscall" "testing" "time" jetstream "github.com/bluesky-social/jetstream" "github.com/bluesky-social/jetstream/internal/simulator/fanout" simhttp "github.com/bluesky-social/jetstream/internal/simulator/http" "github.com/bluesky-social/jetstream/internal/simulator/world" "github.com/stretchr/testify/require" ) const streamDifferentialPin = "289b0328c2e1a0ccf8c870cb45de0b2397de19fb" type synchronizedBuffer struct { mu sync.Mutex b bytes.Buffer } type listReposCountingHandler struct { next http.Handler calls atomic.Int64 mu sync.Mutex events []string } type discoveryParityHandler struct { next http.Handler mu sync.Mutex resume string resumeHits int longCursor string } type initialLongCursorHandler struct { next http.Handler mu sync.Mutex longCursor string longHits int dids [2]string } func (h *initialLongCursorHandler) ServeHTTP(rw http.ResponseWriter, r *http.Request) { if r.URL.Path != "/xrpc/com.atproto.sync.listRepos" { h.next.ServeHTTP(rw, r) return } cursor := r.URL.Query().Get("cursor") rw.Header().Set("Content-Type", "application/json") if cursor == "" { _, _ = fmt.Fprintf(rw, `{"repos":[{"did":"%s","active":true}],"cursor":"%s"}`, h.dids[0], h.longCursor) return } if cursor == h.longCursor { h.mu.Lock() h.longHits++ hit := h.longHits h.mu.Unlock() if hit == 1 { _, _ = fmt.Fprintf(rw, `{"repos":[{"did":"%s","active":true}]}`, h.dids[1]) } else { _, _ = io.WriteString(rw, `{"repos":[]}`) } return } http.Error(rw, "unexpected listRepos cursor", http.StatusBadRequest) } func (h *discoveryParityHandler) ServeHTTP(rw http.ResponseWriter, r *http.Request) { if r.URL.Path != "/xrpc/com.atproto.sync.listRepos" { h.next.ServeHTTP(rw, r) return } cursor := r.URL.Query().Get("cursor") h.mu.Lock() if cursor != "" && h.resume == "" { // The first non-empty input cursor belongs to bootstrap's second // page. Discovery later resumes from exactly that cursor. h.resume = cursor h.resumeHits = 1 h.mu.Unlock() h.next.ServeHTTP(rw, r) return } if cursor == h.resume && h.resumeHits == 1 { h.resumeHits++ longCursor := h.longCursor h.mu.Unlock() rw.Header().Set("Content-Type", "application/json") _, _ = fmt.Fprintf(rw, `{"repos":[{"did":"did:plc:discoveredinactive","active":false}],"cursor":"%s"}`, longCursor) return } if cursor == h.longCursor { h.mu.Unlock() rw.Header().Set("Content-Type", "application/json") _, _ = io.WriteString(rw, `{"repos":[{"did":"did:plc:discoveredafterlongcursor","active":true}]}`) return } h.mu.Unlock() h.next.ServeHTTP(rw, r) } type heldGetRepoHandler struct { next http.Handler mu sync.Mutex release chan struct{} released bool active int peak int total int } func (h *heldGetRepoHandler) begin(t *testing.T) { t.Helper() h.mu.Lock() defer h.mu.Unlock() require.Zero(t, h.active) h.release = make(chan struct{}) h.released = false h.peak = 0 h.total = 0 } func (h *heldGetRepoHandler) ServeHTTP(rw http.ResponseWriter, r *http.Request) { if r.URL.Path != "/xrpc/com.atproto.sync.getRepo" { h.next.ServeHTTP(rw, r) return } h.mu.Lock() release := h.release h.active++ h.total++ if h.active > h.peak { h.peak = h.active } h.mu.Unlock() <-release h.next.ServeHTTP(rw, r) h.mu.Lock() h.active-- h.mu.Unlock() } func (h *heldGetRepoHandler) waitForPeak(t *testing.T, want int) { t.Helper() require.Eventually(t, func() bool { h.mu.Lock() defer h.mu.Unlock() return h.peak >= want }, 20*time.Second, 5*time.Millisecond, "held getRepo peak never reached %d", want) } func (h *heldGetRepoHandler) releaseAll() { h.mu.Lock() defer h.mu.Unlock() if !h.released { close(h.release) h.released = true } } func (h *heldGetRepoHandler) receipt() (peak, total int) { h.mu.Lock() defer h.mu.Unlock() return h.peak, h.total } func TestStreamBackfillWorkerControlsOracle(t *testing.T) { require.Equal(t, streamDifferentialPin, os.Getenv("STREAM_ORACLE_EXPECTED_PIN"), "runner must attest the exact pinned upstream checkout") simCfg := world.DefaultConfig() simCfg.DataDir = filepath.Join(t.TempDir(), "simulator") simCfg.Seed = 0xbac4f111 simCfg.Accounts = 104 simCfg.InitialRecords = 1 simCfg.CommitsPerSec = 1 simCfg.FirehoseHistory = 256 w, err := world.New(t.Context(), simCfg) require.NoError(t, err) defer func() { require.NoError(t, w.Close()) }() _, err = w.EnsureSeed() require.NoError(t, err) require.NoError(t, w.Bootstrap(t.Context(), slog.Default())) fan := fanout.New(256) require.NoError(t, w.AttachRuntime( rand.New(rand.NewPCG(simCfg.Seed^0xfeedf00d, simCfg.Seed^0xc0ffee)), fan, )) server := httptest.NewUnstartedServer(nil) gate := &heldGetRepoHandler{next: simhttp.NewHandler(w, server.URL)} server.Config.Handler = gate server.Start() defer server.Close() run := func(t *testing.T, workerArg string, want int) { t.Helper() gate.begin(t) defer gate.releaseAll() port := freeStreamOraclePort(t) baseURL := "http://127.0.0.1:" + strconv.Itoa(port) logs := &synchronizedBuffer{} args := []string{ "--port=" + strconv.Itoa(port), "--data-dir=" + filepath.Join(t.TempDir(), "stream"), "--relay-url=" + server.URL, "--plc-url=" + server.URL, "--backfill", "--skip-merge-discovery", "--max-segment-bytes=1", "--compaction-interval=0", "--retry-interval=0", "--no-verify", } if workerArg != "" { args = append(args, "--backfill-workers="+workerArg) } cmd := exec.Command(streamOracleBinary(t), args...) cmd.Stdout = logs cmd.Stderr = logs require.NoError(t, cmd.Start()) stopped := false defer func() { if stopped { return } _ = cmd.Process.Kill() _ = waitStreamCommand(cmd, 10*time.Second) }() gate.waitForPeak(t, want) gate.releaseAll() waitForStreamOracleServing(t, cmd, baseURL, logs) require.NoError(t, cmd.Process.Signal(syscall.SIGTERM)) require.NoErrorf(t, waitStreamCommand(cmd, 20*time.Second), "Stream did not stop cleanly:\n%s", logs.String()) stopped = true peak, total := gate.receipt() require.Equal(t, want, peak) require.Equal(t, simCfg.Accounts, total) } t.Run("explicit positive worker count reaches physical admission", func(t *testing.T) { run(t, "7", 7) }) t.Run("zero selects production default", func(t *testing.T) { run(t, "0", 100) }) t.Run("omitted flag uses production default", func(t *testing.T) { run(t, "", 100) }) } func (h *listReposCountingHandler) ServeHTTP(rw http.ResponseWriter, r *http.Request) { if r.URL.Path == "/xrpc/com.atproto.sync.listRepos" { h.calls.Add(1) h.mu.Lock() h.events = append(h.events, "listRepos") h.mu.Unlock() } else if r.URL.Path == "/xrpc/com.atproto.sync.getRepo" { h.mu.Lock() h.events = append(h.events, "getRepo") h.mu.Unlock() } h.next.ServeHTTP(rw, r) } func (h *listReposCountingHandler) reset() { h.calls.Store(0) h.mu.Lock() h.events = h.events[:0] h.mu.Unlock() } func (h *listReposCountingHandler) eventLog() string { h.mu.Lock() defer h.mu.Unlock() return strings.Join(h.events, ",") } func TestStreamBootstrapControlsOracle(t *testing.T) { require.Equal(t, streamDifferentialPin, os.Getenv("STREAM_ORACLE_EXPECTED_PIN"), "runner must attest the exact pinned upstream checkout") simCfg := world.DefaultConfig() simCfg.DataDir = filepath.Join(t.TempDir(), "simulator") simCfg.Seed = 0x51ced15c0 simCfg.Accounts = 4 simCfg.InitialRecords = 1 simCfg.CommitsPerSec = 1 simCfg.FirehoseHistory = 128 w, err := world.New(t.Context(), simCfg) require.NoError(t, err) defer func() { require.NoError(t, w.Close()) }() _, err = w.EnsureSeed() require.NoError(t, err) require.NoError(t, w.Bootstrap(t.Context(), slog.Default())) fan := fanout.New(128) require.NoError(t, w.AttachRuntime( rand.New(rand.NewPCG(simCfg.Seed^0xfeedf00d, simCfg.Seed^0xc0ffee)), fan, )) server := httptest.NewUnstartedServer(nil) counter := &listReposCountingHandler{next: simhttp.NewHandlerWithOptions(w, server.URL, simhttp.HandlerOptions{ ListReposPageLimit: 2, })} server.Config.Handler = counter server.Start() defer server.Close() selected, _, err := w.ListReposPage(0, 1) require.NoError(t, err) require.Len(t, selected, 1) initialLongCursorRepos, _, err := w.ListReposPage(0, 2) require.NoError(t, err) require.Len(t, initialLongCursorRepos, 2) run := func(t *testing.T, selectionArgs []string, expectSkip bool) (int64, string) { t.Helper() counter.reset() port := freeStreamOraclePort(t) baseURL := "http://127.0.0.1:" + strconv.Itoa(port) logs := &synchronizedBuffer{} args := []string{ "--port=" + strconv.Itoa(port), "--data-dir=" + filepath.Join(t.TempDir(), "stream"), "--relay-url=" + server.URL, "--plc-url=" + server.URL, "--backfill-workers=2", "--max-segment-bytes=1", "--compaction-interval=0", "--retry-interval=0", } args = append(args, selectionArgs...) cmd := exec.Command(streamOracleBinary(t), args...) cmd.Stdout = logs cmd.Stderr = logs require.NoError(t, cmd.Start()) waitForStreamOracleServing(t, cmd, baseURL, logs) require.NoError(t, cmd.Process.Signal(syscall.SIGTERM)) require.NoErrorf(t, waitStreamCommand(cmd, 15*time.Second), "Stream did not stop cleanly:\n%s", logs.String()) if expectSkip { require.Contains(t, logs.String(), "merge discovery: skipped by configuration") } else { require.Contains(t, logs.String(), "merge discovery: 0 post-bootstrap DIDs queued for retry") } if slices.Contains(selectionArgs, "--backfill-async-flush-workers=0") { require.Contains(t, logs.String(), "backfill: async block compression off") } else if slices.Contains(selectionArgs, "--backfill-async-flush-workers=2") { require.Contains(t, logs.String(), "backfill: async block compression on (2 workers)") } else { require.Contains(t, logs.String(), "backfill: async block compression on (4 workers)") } return counter.calls.Load(), counter.eventLog() } t.Run("normal full crawl performs discovery", func(t *testing.T) { calls, _ := run(t, []string{"--backfill"}, false) require.Equal(t, int64(3), calls, "two bootstrap pages plus one final-page discovery replay") }) t.Run("explicit skip omits discovery", func(t *testing.T) { calls, _ := run(t, []string{"--backfill", "--skip-merge-discovery"}, true) require.Equal(t, int64(2), calls, "skip must leave only the two bootstrap pages") }) t.Run("max repo selection automatically skips discovery", func(t *testing.T) { calls, _ := run(t, []string{"--max-backfill-repos=4"}, true) require.Equal(t, int64(2), calls, "partial crawl must not add a discovery request") }) t.Run("explicit DID selection automatically skips discovery", func(t *testing.T) { calls, _ := run(t, []string{"--backfill-repos=" + string(selected[0].DID)}, true) require.Equal(t, int64(0), calls, "explicit DID mode must bypass both listRepos and discovery") }) t.Run("page-sized batch dispatches before second page", func(t *testing.T) { _, events := run(t, []string{"--backfill", "--skip-merge-discovery", "--backfill-batch-size=2"}, true) firstGet := strings.Index(events, "getRepo") secondList := strings.LastIndex(events, "listRepos") require.NotEqual(t, -1, firstGet) require.NotEqual(t, -1, secondList) require.Less(t, firstGet, secondList, events) }) t.Run("page-aligned batch may exceed target before dispatch", func(t *testing.T) { _, events := run(t, []string{"--backfill", "--skip-merge-discovery", "--backfill-batch-size=3"}, true) firstGet := strings.Index(events, "getRepo") secondList := strings.LastIndex(events, "listRepos") require.NotEqual(t, -1, firstGet) require.NotEqual(t, -1, secondList) require.Less(t, secondList, firstGet, events) }) t.Run("zero uses production cross-page default", func(t *testing.T) { _, events := run(t, []string{"--backfill", "--skip-merge-discovery", "--backfill-batch-size=0"}, true) firstGet := strings.Index(events, "getRepo") secondList := strings.LastIndex(events, "listRepos") require.NotEqual(t, -1, firstGet) require.NotEqual(t, -1, secondList) require.Less(t, secondList, firstGet, events) }) t.Run("zero disables async block compression", func(t *testing.T) { calls, _ := run(t, []string{"--backfill", "--skip-merge-discovery", "--backfill-async-flush-workers=0"}, true) require.Equal(t, int64(2), calls) }) t.Run("positive async block compression count is applied", func(t *testing.T) { calls, _ := run(t, []string{"--backfill", "--skip-merge-discovery", "--backfill-async-flush-workers=2"}, true) require.Equal(t, int64(2), calls) }) t.Run("discovery preserves inactive rows and follows a 4 KiB cursor", func(t *testing.T) { longCursor := strings.Repeat("x", 4096) special := httptest.NewUnstartedServer(nil) handler := &discoveryParityHandler{ next: counter.next, longCursor: longCursor, } special.Config.Handler = handler special.Start() defer special.Close() port := freeStreamOraclePort(t) baseURL := "http://127.0.0.1:" + strconv.Itoa(port) logs := &synchronizedBuffer{} cmd := exec.Command(streamOracleBinary(t), "--port="+strconv.Itoa(port), "--data-dir="+filepath.Join(t.TempDir(), "stream"), "--relay-url="+special.URL, "--plc-url="+special.URL, "--backfill", "--backfill-workers=2", "--max-segment-bytes=1", "--compaction-interval=0", "--retry-interval=0", "--no-verify", ) cmd.Stdout = logs cmd.Stderr = logs require.NoError(t, cmd.Start()) stopped := false defer func() { if stopped { return } _ = cmd.Process.Kill() _ = waitStreamCommand(cmd, 10*time.Second) }() waitForStreamOracleServing(t, cmd, baseURL, logs) accountText := func(did string) string { req, err := http.NewRequest(http.MethodGet, baseURL+"/status?tab=accounts&account="+did, nil) require.NoError(t, err) req.Header.Set("Accept", "text/plain") resp, err := http.DefaultClient.Do(req) require.NoError(t, err) defer resp.Body.Close() body, err := io.ReadAll(resp.Body) require.NoError(t, err) require.Equal(t, http.StatusOK, resp.StatusCode, string(body)) return string(body) } inactive := accountText("did:plc:discoveredinactive") require.Contains(t, inactive, "found\tyes") require.Contains(t, inactive, "active\tfalse") require.Contains(t, inactive, "backfill\tfailed") require.Contains(t, inactive, "last error\tdiscovered post-bootstrap; queued for retry") afterLong := accountText("did:plc:discoveredafterlongcursor") require.Contains(t, afterLong, "found\tyes") require.Contains(t, afterLong, "active\ttrue") require.Contains(t, afterLong, "backfill\tfailed") require.NoError(t, cmd.Process.Signal(syscall.SIGTERM)) require.NoErrorf(t, waitStreamCommand(cmd, 20*time.Second), "Stream did not stop cleanly:\n%s", logs.String()) stopped = true }) t.Run("initial crawl follows a 4 KiB cursor without truncation", func(t *testing.T) { longCursor := strings.Repeat("y", 4096) special := httptest.NewUnstartedServer(nil) handler := &initialLongCursorHandler{ next: counter.next, longCursor: longCursor, dids: [2]string{ string(initialLongCursorRepos[0].DID), string(initialLongCursorRepos[1].DID), }, } special.Config.Handler = handler special.Start() defer special.Close() port := freeStreamOraclePort(t) baseURL := "http://127.0.0.1:" + strconv.Itoa(port) logs := &synchronizedBuffer{} cmd := exec.Command(streamOracleBinary(t), "--port="+strconv.Itoa(port), "--data-dir="+filepath.Join(t.TempDir(), "stream"), "--relay-url="+special.URL, "--plc-url="+special.URL, "--backfill", "--backfill-workers=2", "--max-segment-bytes=1", "--compaction-interval=0", "--retry-interval=0", "--no-verify", ) cmd.Stdout = logs cmd.Stderr = logs require.NoError(t, cmd.Start()) stopped := false defer func() { if stopped { return } _ = cmd.Process.Kill() _ = waitStreamCommand(cmd, 10*time.Second) }() waitForStreamOracleServing(t, cmd, baseURL, logs) require.Contains(t, logs.String(), "merge discovery: 0 post-bootstrap DIDs queued for retry") require.NotContains(t, logs.String(), "ListReposCursorTooLong") require.NoError(t, cmd.Process.Signal(syscall.SIGTERM)) require.NoErrorf(t, waitStreamCommand(cmd, 20*time.Second), "Stream did not stop cleanly:\n%s", logs.String()) stopped = true }) } func (b *synchronizedBuffer) Write(p []byte) (int, error) { b.mu.Lock() defer b.mu.Unlock() return b.b.Write(p) } func (b *synchronizedBuffer) String() string { b.mu.Lock() defer b.mu.Unlock() return b.b.String() } func TestStreamDifferentialOracle(t *testing.T) { require.Equal(t, streamDifferentialPin, os.Getenv("STREAM_ORACLE_EXPECTED_PIN"), "runner must attest the exact pinned upstream checkout") streamBin := os.Getenv("STREAM_ORACLE_BIN") require.NotEmpty(t, streamBin) streamBin, err := filepath.Abs(streamBin) require.NoError(t, err) _, err = os.Stat(streamBin) require.NoError(t, err) simCfg := world.DefaultConfig() simCfg.DataDir = filepath.Join(t.TempDir(), "simulator") simCfg.Seed = 0x5eed5eed simCfg.Accounts = 8 simCfg.InitialRecords = 3 // RunTraffic is deliberately never started; a positive configured rate is // still required by the world's production config validator. simCfg.CommitsPerSec = 1 simCfg.FirehoseHistory = 1024 w, err := world.New(t.Context(), simCfg) require.NoError(t, err) defer func() { require.NoError(t, w.Close()) }() _, err = w.EnsureSeed() require.NoError(t, err) require.NoError(t, w.Bootstrap(t.Context(), slog.Default())) fan := fanout.New(2048) require.NoError(t, w.AttachRuntime( rand.New(rand.NewPCG(simCfg.Seed^0xfeedf00d, simCfg.Seed^0xc0ffee)), fan, )) simSrv := httptest.NewServer(nil) simSrv.Config.Handler = simhttp.NewHandler(w, simSrv.URL) defer simSrv.Close() port := freeStreamOraclePort(t) baseURL := "http://127.0.0.1:" + strconv.Itoa(port) dataDir := filepath.Join(t.TempDir(), "stream") logs := &synchronizedBuffer{} ctx, cancel := context.WithCancel(context.Background()) defer cancel() cmd := exec.CommandContext(ctx, streamBin, "--port="+strconv.Itoa(port), "--data-dir="+dataDir, "--relay-url="+simSrv.URL, "--plc-url="+simSrv.URL, "--max-backfill-repos="+strconv.Itoa(simCfg.Accounts), "--backfill-workers=4", "--max-segment-bytes=1", "--compaction-interval=0", "--retry-interval=0", ) cmd.Stdout = logs cmd.Stderr = logs require.NoError(t, cmd.Start()) stopped := false defer func() { if stopped { return } _ = cmd.Process.Signal(syscall.SIGTERM) waitDone := make(chan struct{}) go func() { _ = cmd.Wait() close(waitDone) }() select { case <-waitDone: case <-time.After(10 * time.Second): _ = cmd.Process.Kill() <-waitDone } }() waitForStreamOracleServing(t, cmd, baseURL, logs) initial := waitForStableStreamSegments(t, dataDir, 1, logs) require.NoError(t, CheckInvariants(initial)) require.NotEmpty(t, initial, "bootstrap must archive real rows") initialMaxSeq := EventsSortedBySeq(initial)[len(initial)-1].Seq ground, err := GroundTruthFromWorld(w) require.NoError(t, err) reconstructed, err := Reconstruct(EventsSortedBySeq(initial)) require.NoError(t, err) require.NoError(t, Compare(ground, reconstructed), "bootstrap archive must reconstruct to the independent simulator MST") // /listSegments becoming ready and the relay websocket attaching are // independent concurrent milestones. Prove the live subscriber is attached // with a real, state-neutral identity event before defining the exact // firehose comparison window; this is an acknowledgement, not a sleep. initial = waitForStreamLiveAttachment(t, w, dataDir, len(initial), baseURL, logs) initialMaxSeq = EventsSortedBySeq(initial)[len(initial)-1].Seq ground, err = GroundTruthFromWorld(w) require.NoError(t, err) reconstructed, err = Reconstruct(EventsSortedBySeq(initial)) require.NoError(t, err) require.NoError(t, Compare(ground, reconstructed), "durable live-attachment barrier must converge to simulator ground truth") startTip := w.CurrentSeq() ctxGenerate := t.Context() _, _, err = w.GenerateRecordOpForTest(ctxGenerate, 0, "create", "app.bsky.feed.post", "oracle-diff") require.NoError(t, err) _, _, err = w.GenerateRecordOpForTest(ctxGenerate, 0, "update", "app.bsky.feed.post", "oracle-diff") require.NoError(t, err) _, _, err = w.GenerateRecordOpForTest(ctxGenerate, 0, "delete", "app.bsky.feed.post", "oracle-diff") require.NoError(t, err) _, err = w.GenerateIdentityForTest(ctxGenerate, 1, false) require.NoError(t, err) _, err = w.GenerateIdentityForTest(ctxGenerate, 1, true) require.NoError(t, err) _, err = w.GenerateAccountStatusForTest(ctxGenerate, 2, false, "deactivated") require.NoError(t, err) _, err = w.GenerateSilentMutationThenSyncForTest(ctxGenerate, 3) require.NoError(t, err) _, err = w.GenerateAccountDeleteForTest(ctxGenerate, 4) require.NoError(t, err) _, err = w.GenerateAccountReactivateForTest(ctxGenerate, 4) require.NoError(t, err) endTip := w.CurrentSeq() require.Greater(t, endTip, startTip) expected, err := ExpectedEventLogFromFirehose(w, startTip, int(endTip-startTip)) require.NoError(t, err) require.NotEmpty(t, expected) zeroEventLogSeqs(expected) requireEventKinds(t, expected, "create", "update", "delete", "identity", "account", "sync", "create_resync") // First establish that every accepted live ticket reached the writer, then // exercise SIGTERM's drain+durable-flush boundary. Disk comparison below is // therefore against a quiescent archive, never a sub-block in memory. waitForStreamPipelineDrain(t, baseURL, endTip, logs) require.NoError(t, cmd.Process.Signal(syscall.SIGTERM)) require.NoErrorf(t, cmd.Wait(), "Stream did not shut down cleanly:\n%s", logs.String()) stopped = true var all []ObservedEvent deadline := time.Now().Add(30 * time.Second) for { all, err = ObserveSegments(dataDir) if err == nil { var live []ObservedEvent for _, ev := range all { if ev.Seq > initialMaxSeq { live = append(live, ev) } } got := NormalizeEventLog(live) zeroEventLogSeqs(got) if len(got) >= len(expected) { require.NoError(t, CompareEventLogMultiset(expected, got), "Stream's durable live log must equal the pinned oracle's firehose expansion") break } } if time.Now().After(deadline) { t.Fatalf("timed out waiting for %d oracle rows: err=%v logs:\n%s", len(expected), err, logs.String()) } time.Sleep(20 * time.Millisecond) } require.NoError(t, CheckInvariants(all)) ground, err = GroundTruthFromWorld(w) require.NoError(t, err) reconstructed, err = Reconstruct(EventsSortedBySeq(all)) require.NoError(t, err) require.NoError(t, Compare(ground, reconstructed), "complete durable Stream log must converge to the independent simulator MST") // Detection-power receipt: an event-log oracle must reject loss of an // intermediate update even when a later delete leaves final state unchanged. mutated := append([]EventLogRow(nil), expected...) for i, row := range mutated { if row.Kind == "update" { mutated = append(mutated[:i], mutated[i+1:]...) break } } require.Error(t, CompareEventLogMultiset(expected, mutated), "oracle must detect a missing intermediate update") // Reopen from durable steady-state metadata. Like upstream, recovery resumes // the active segment in place: the archive XRPC frontier remains the last // sealed segment, while a normal subscriber must replay both sealed and // active durable rows before joining live delivery. publicThrough := initialMaxSeqFrom(all) cmd2 := exec.CommandContext(ctx, streamBin, cmd.Args[1:]...) cmd2.Stdout = logs cmd2.Stderr = logs require.NoError(t, cmd2.Start()) stopped2 := false defer func() { if stopped2 { return } _ = cmd2.Process.Signal(syscall.SIGTERM) _ = cmd2.Wait() }() waitForStreamOracleServing(t, cmd2, baseURL, logs) sealedThrough := streamSealedHighWater(t, baseURL) require.Less(t, sealedThrough, publicThrough, "restart oracle must retain an active segment beyond the sealed XRPC frontier") sealedEvents := collectStreamOracleBackfill(t, baseURL, sealedThrough) require.NoError(t, CompareEventLogMultiset( NormalizeEventLog(observedEventsThrough(all, sealedThrough)), NormalizeEventLog(sealedEvents), ), "backfill-only public client must reproduce the sealed archive frontier") // Exercise the pinned upstream public client against Stream's actual XRPC // archive plan/getSegment/getBlock plus cold active-segment replay and // /subscribe-v2 decode path. publicEvents := collectStreamOracleThrough(t, baseURL, publicThrough) publicModel, err := Reconstruct(EventsSortedBySeq(publicEvents)) require.NoError(t, err) require.NoError(t, Compare(ground, publicModel), "pinned upstream public client must reconstruct Stream to simulator ground truth") require.Zero(t, fan.TotalDrops(), "simulator fanout must not drop oracle frames") t.Logf("semantic receipt: bootstrap=%d live_rows=%d total=%d sealed_public=%d public=%d sealed_through=%d upstream_seq=%d..%d pin=%s", len(initial), len(expected), len(all), len(sealedEvents), len(publicEvents), sealedThrough, startTip+1, endTip, streamDifferentialPin) require.NoError(t, cmd2.Process.Signal(syscall.SIGTERM)) require.NoErrorf(t, cmd2.Wait(), "restarted Stream did not shut down cleanly:\n%s", logs.String()) stopped2 = true require.NotContains(t, logs.String(), "fatal event handling failure") require.NotContains(t, logs.String(), "OutOfMemory") } func waitForStreamPipelineDrain(t *testing.T, baseURL string, upstreamThrough int64, logs *synchronizedBuffer) { t.Helper() client := &http.Client{Timeout: time.Second} deadline := time.Now().Add(20 * time.Second) lastMetrics := "unavailable" for time.Now().Before(deadline) { resp, err := client.Get(baseURL + "/metrics") if err == nil { body := new(bytes.Buffer) _, _ = body.ReadFrom(resp.Body) _ = resp.Body.Close() text := body.String() lastMetrics = text submitted, haveSubmitted := metricValue(text, `stream_pipeline_ticket{stage="submitted"}`) emitted, haveEmitted := metricValue(text, `stream_pipeline_ticket{stage="emitted"}`) upstream, haveUpstream := metricValue(text, "stream_upstream_seq") if haveSubmitted && haveEmitted && haveUpstream && emitted == submitted && upstream >= float64(upstreamThrough) { return } } time.Sleep(20 * time.Millisecond) } t.Fatalf("pipeline did not drain through upstream seq %d\nmetrics:\n%s\nlogs:\n%s", upstreamThrough, lastMetrics, logs.String()) } func metricValue(metrics, family string) (float64, bool) { for _, line := range strings.Split(metrics, "\n") { fields := strings.Fields(line) if len(fields) != 2 || fields[0] != family { continue } value, err := strconv.ParseFloat(fields[1], 64) return value, err == nil } return 0, false } func freeStreamOraclePort(t *testing.T) int { t.Helper() ln, err := net.Listen("tcp", "127.0.0.1:0") require.NoError(t, err) defer func() { require.NoError(t, ln.Close()) }() return ln.Addr().(*net.TCPAddr).Port } func waitForStreamOracleServing(t *testing.T, cmd *exec.Cmd, baseURL string, logs *synchronizedBuffer) { t.Helper() client := &http.Client{Timeout: time.Second} deadline := time.Now().Add(45 * time.Second) for time.Now().Before(deadline) { if cmd.ProcessState != nil { t.Fatalf("Stream exited before serving: %v\n%s", cmd.ProcessState, logs.String()) } resp, err := client.Get(baseURL + "/xrpc/network.bsky.jetstream.listSegments") if err == nil { _ = resp.Body.Close() if resp.StatusCode == http.StatusOK { return } } time.Sleep(25 * time.Millisecond) } t.Fatalf("timed out waiting for Stream steady state\n%s", logs.String()) } func waitForStableStreamSegments(t *testing.T, dataDir string, minimum int, logs *synchronizedBuffer) []ObservedEvent { t.Helper() deadline := time.Now().Add(20 * time.Second) lastCount := -1 stable := 0 for time.Now().Before(deadline) { events, err := ObserveSegments(dataDir) if err == nil && len(events) >= minimum { if len(events) == lastCount { stable++ } else { stable = 0 lastCount = len(events) } if stable >= 5 { return events } } time.Sleep(20 * time.Millisecond) } t.Fatalf("Stream segments never became stable\n%s", logs.String()) return nil } func waitForStreamLiveAttachment(t *testing.T, w *world.World, dataDir string, baseline int, baseURL string, logs *synchronizedBuffer) []ObservedEvent { t.Helper() for attempt := 0; attempt < 20; attempt++ { // A sync repair owns an explicit archive flush boundary. If an attempt // races ahead of websocket attachment, the next sync repairs the same // DID to its latest authoritative state, so no state is stranded. _, err := w.GenerateSilentMutationThenSyncForTest(t.Context(), 6) require.NoError(t, err) deadline := time.Now().Add(time.Second) for time.Now().Before(deadline) { events, observeErr := ObserveSegments(dataDir) if observeErr == nil && len(events) > baseline { return waitForStableStreamSegments(t, dataDir, len(events), logs) } time.Sleep(10 * time.Millisecond) } } metrics := "unavailable" if resp, err := (&http.Client{Timeout: time.Second}).Get(baseURL + "/metrics"); err == nil { defer func() { _ = resp.Body.Close() }() body := new(bytes.Buffer) _, _ = body.ReadFrom(resp.Body) metrics = body.String() } t.Fatalf("live firehose never acknowledged a durable sync-repair frame\nmetrics:\n%s\nlogs:\n%s", metrics, logs.String()) return nil } func zeroEventLogSeqs(rows []EventLogRow) { for i := range rows { rows[i].Seq = 0 } } func requireEventKinds(t *testing.T, rows []EventLogRow, kinds ...string) { t.Helper() seen := make(map[string]int) for _, row := range rows { seen[row.Kind]++ } for _, kind := range kinds { require.Positivef(t, seen[kind], "anti-vacuity: expected real %s row; counts=%v", kind, seen) } } func initialMaxSeqFrom(events []ObservedEvent) uint64 { var maxSeq uint64 for _, ev := range events { if ev.Seq > maxSeq { maxSeq = ev.Seq } } return maxSeq } func observedEventsThrough(events []ObservedEvent, through uint64) []ObservedEvent { out := make([]ObservedEvent, 0, len(events)) for _, event := range events { if event.Seq <= through { out = append(out, event) } } return out } func streamSealedHighWater(t *testing.T, baseURL string) uint64 { t.Helper() resp, err := (&http.Client{Timeout: 5 * time.Second}).Get( baseURL + "/xrpc/network.bsky.jetstream.listSegments?limit=1000", ) require.NoError(t, err) defer func() { require.NoError(t, resp.Body.Close()) }() require.Equal(t, http.StatusOK, resp.StatusCode) var payload struct { Segments []struct { MaxSeq uint64 `json:"maxSeq"` } `json:"segments"` } require.NoError(t, json.NewDecoder(resp.Body).Decode(&payload)) require.NotEmpty(t, payload.Segments) var highWater uint64 for _, segment := range payload.Segments { highWater = max(highWater, segment.MaxSeq) } require.NotZero(t, highWater) return highWater } func collectStreamOracleBackfill(t *testing.T, baseURL string, through uint64) []ObservedEvent { t.Helper() client, err := jetstream.Subscribe(baseURL, jetstream.WithAfterSeq(0), jetstream.WithBeforeSeq(through), jetstream.WithSnapshotOnly(), jetstream.WithBatchSize(32), ) require.NoError(t, err) defer func() { require.NoError(t, client.Close()) }() ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) defer cancel() var out []ObservedEvent for batch, err := range client.Events(ctx) { require.NoError(t, err) for _, event := range batch.Events() { out = append(out, observedEventFromClient(t, event)) } } require.NoError(t, ctx.Err()) require.NotEmpty(t, out) require.Equal(t, through, initialMaxSeqFrom(out), "public client must reach archive high-water") return out } func collectStreamOracleThrough(t *testing.T, baseURL string, through uint64) []ObservedEvent { t.Helper() client, err := jetstream.Subscribe(baseURL, jetstream.WithAfterSeq(0), jetstream.WithBatchSize(32), ) require.NoError(t, err) defer func() { require.NoError(t, client.Close()) }() ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) defer cancel() var out []ObservedEvent for batch, err := range client.Events(ctx) { require.NoError(t, err) for _, event := range batch.Events() { observed := observedEventFromClient(t, event) if observed.Seq <= through { out = append(out, observed) } } if initialMaxSeqFrom(out) >= through { break } } require.NotEmpty(t, out) require.Equal(t, through, initialMaxSeqFrom(out), "normal public client must replay through the active durable frontier") return out }