From acb10fd7781d9a95ef3cc12da565f1e3ded89faa Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Fri, 24 Jul 2026 10:20:13 -0500 Subject: [PATCH] preserve complete listRepos discovery --- README.md | 5 +- docs/bootstrap-semantic-parity.md | 22 +- docs/configuration-parity.md | 5 +- docs/semantic-parity.md | 16 +- src/internal/bootstrap/engine.zig | 14 +- src/internal/bootstrap/lifecycle.zig | 17 +- .../stream_differential_test.go | 204 ++++++++++++++++++ 7 files changed, 247 insertions(+), 36 deletions(-) diff --git a/README.md b/README.md index 7483816..b500823 100644 --- a/README.md +++ b/README.md @@ -24,7 +24,7 @@ experiment. - a whole-network bootstrap implementation with concurrent live capture and deterministic merge. Repository completion is coupled to ordinary archive writer durability independently of the 100,000-entry dispatch/checkpoint - batch; several lifecycle, discovery, and recovery paths remain blocked by + batch; several recovery, pending-repair, and live-ingest paths remain blocked by the deployment-gate audit - Sync 1.1 commit verification, including PLC key rotation, MST inversion, op-CID checks, replay protection, durable repository-chain state, and @@ -345,6 +345,9 @@ process receipt observes three real `listRepos` requests for a normal two-page crawl and exactly two when discovery is skipped. The same receipt proves a page-sized batch dispatches before page two, while a cross-page batch and the zero/default setting enumerate both pages before the first real `getRepo`. +Initial and post-bootstrap enumeration retain arbitrary-length relay cursors; +the real-process receipt follows a 4 KiB cursor and verifies that discovery +durably records both active and inactive accounts. The same durable account row retains the distinct initial-backfill and latest known revisions, update time, declared handle, PDS endpoint, and reserved diff --git a/docs/bootstrap-semantic-parity.md b/docs/bootstrap-semantic-parity.md index 432e5cf..82d03f9 100644 --- a/docs/bootstrap-semantic-parity.md +++ b/docs/bootstrap-semantic-parity.md @@ -4,8 +4,8 @@ Reference: - upstream Jetstream `f29815c391fc2644f8a3dd36b899fb3697dd1ea6` - Atmos `v0.2.14` -- Stream `bdb9a796de5fe42b39b79a53e22fcff55a3f94f7` plus the - bootstrap-to-merging transition and audit changes in this commit +- Stream `a345f4054df568e364aff08a9b1726275fc48e29` plus the discovery + changes and audit updates in this commit The status words have the meanings defined in [semantic-parity.md](semantic-parity.md). A configuration knob being wired @@ -15,9 +15,9 @@ correctly does not verify the behavior it controls. **The overall lifecycle remains blocked.** The dispatch-batch completion defect, correctness-metadata fail-open reads, merge restart guard, and -bootstrap-to-merging ordering now have production-boundary fixes and focused -proof. Independent discovery, recovery, retry, live-encoding, and reconnect -defects still prevent another deployment. +bootstrap-to-merging ordering, and post-bootstrap discovery now have +production-boundary fixes and focused proof. Independent recovery, retry, +live-encoding, and reconnect defects still prevent another deployment. ## Initial enumeration and download @@ -34,7 +34,7 @@ defects still prevent another deployment. | Ordinary retry budget | **verified** | The initial attempt plus three ordinary transient retries, exponential base, jitter, and 30-second cap have mixed 503/429 fixtures. | The production firehose reconnect defect is separate from getRepo retry. | | Rate-limit retry budget | **verified** | 429 uses a separate twenty-retry budget and honors parsed server reset information with the bootstrap ceiling. Tests show 429 does not consume the ordinary budget. | Pre-request per-PDS scheduling is unavailable after relay indirection, as it is upstream. | | Terminal repository errors | **verified** | Structured RepoNotFound becomes complete with no rows; deactivated/suspended/takendown become unavailable. Focused tests inspect the durable status. | None known for the classified responses. | -| Inactive repository during initial discovery | **verified** | A new inactive entry is stored as not_started/active=false and is not dispatched; later active flips preserve the row. | Post-bootstrap discovery does not preserve the same rule; see below. | +| Inactive repository during initial discovery | **verified** | A new inactive entry is stored as not_started/active=false and is not dispatched; later active flips preserve the row. Post-bootstrap discovery likewise preserves inactive unknown entries as failed retry-state rows. | None known for this invariant. | | Explicit selected backfill | **partial** | The selected DID list bypasses listRepos, preserves caller order, conflicts with max-repos, clears stale discovery state, resolves identity metadata, maintains the handle index, and uses the same writer-gated per-repository completion machinery as whole-network bootstrap. | The exact current artifact has not rerun the full selected-repo receipt. | | Max-repos debug mode | **partial** | The flag truncates selected work and avoids committing a resumable whole-network cursor. Phase and relay-cursor reads now fail closed. | It remains a debug path and is not evidence for whole-network discovery or merge behavior. | @@ -58,7 +58,7 @@ defects still prevent another deployment. |---|---|---|---| | Relay listRepos cursor checkpoint | **verified** | Stream atomically writes the relay cursor and last non-empty discovery cursor only after the dispatch batch drains. Final empty cursor handling has unit coverage, while the straggler/reopen test proves individual repository completion can become durable with the page cursor still absent. | The 100,000-entry unit is checkpoint/re-enumeration cadence, not repository completion durability. | | Relay firehose cursor read | **verified** | Durable cursor writes remain coupled to archive metadata batches. Reads now propagate RocksDB errors and require upstream's `[version=1][uint64 LE]` encoding, rejecting wrong width, unknown version, and values above `MaxInt64`; writes reject negative cursors. Focused tests exercise every rejection, and a startup migration converts Stream's prior valid 8-byte encoding without accepting other malformed values. | This row does not close the independent Zat reconnect failure. | -| listRepos cursor read | **verified** | `loadRelayListReposCursor` propagates RocksDB errors and distinguishes a missing key. | Valid cursor length remains artificially bounded by fixed buffers during enumeration; initial enumeration fails loud, while merge discovery silently stops. | +| listRepos cursor read | **verified** | `loadRelayListReposCursor` propagates RocksDB errors and distinguishes a missing key. Initial enumeration and merge discovery retain dynamically owned cursors without a local length ceiling; real-process receipts walk 4 KiB cursors through both paths. | None known for this invariant. | | Lifecycle phase read | **verified** | Phase writes are synced. Reads distinguish a missing key from RocksDB failure and reject every value outside `bootstrap`, `merging`, and `steady_state`; focused tests inject both corruption and a real Store read error. The production orchestrator uses the same fallible read before any fresh-directory write. | Persisted phase-entry timing remains a separate gap below. | | Bootstrap → merging commit point | **verified** | After the backfill future drains successfully, Stream syncs `phase=merging` while bootstrap-live capture is still running, then requests capture shutdown and performs the archive cutover. This is the pinned upstream ordering. Named process-abort seams immediately before the write and immediately after it/before live stop both recover on the same disk through the lifecycle oracle. | This row establishes the durable phase boundary, not the separate persisted phase-entry timestamp/status contract. | | Merging → steady_state commit point | **partial** | Stream drains, compacts, reconciles the manifest, discovers, removes the source tree, syncs the data directory, deletes cursors, publishes seq metadata, then writes steady_state. Phase reads fail closed, and the restart guard distinguishes only `FileNotFound` from every other filesystem error before taking the cleanup-complete path. Named write-side crash tests cover several seams. | Exact persisted phase-entry timing remains absent, and other merge/discovery rows remain partial or blocked. | @@ -72,12 +72,12 @@ defects still prevent another deployment. | Bootstrap-live failure propagation | **partial** | The lifecycle now observes an unexpectedly completed capture future and treats it as fatal. | `HttpConnectionClosing` can still escape Zat's reconnect loop and terminate the process; the exact dependency behavior needs a regression test on Linux. | | Merge source existence guard | **verified** | Only `FileNotFound` selects the cleanup-complete restart path. Other filesystem errors propagate without cursor deletion or phase advancement; focused tests exercise a non-directory source and the normal absent-source restart. | This narrow guard does not prove every source-file read in the merge walker; those remain covered by their own rows. | | Merge source cursor read | **verified** | Successful source completion atomically commits upstream's `[version=1][uint64 LE]` cursor and latest-revision updates. Reads distinguish absence from Store failure and reject wrong width/version; focused tests inject each failure, and startup migrates Stream's prior valid 8-byte cursor. | This row proves cursor decoding and commit atomicity, not every source-file read. | -| Merge row filtering | **verified** | Rows are dropped only when the repository is complete and the source revision is at/below its backfill watermark; missing lookups are cached. Kept/dropped fixtures exercise this. | The surrounding cursor recovery remains blocked. | +| Merge row filtering | **verified** | Rows are dropped only when the repository is complete and the source revision is at/below its backfill watermark; missing lookups are cached. Kept/dropped fixtures exercise this. | This predicate proof does not establish source-file recovery behavior. | | Latest revision refresh | **partial** | Kept rev-bearing rows update latest revision in the same batch as the source cursor while preserving the backfill revision. | Missing repo rows are skipped; this matches the intended defensive behavior, but corrupt cursor state can replay updates and rows. | | Pending-repository pass | **blocked** | Pending repositories are repaired after captured live rows are merged, preserving intended row order. | Stream materializes pending DIDs and processes them through a separate path rather than upstream's bounded retry runner with the same global/per-host gates and retry semantics. | -| Merge-tail compaction and manifest reconcile | **partial** | Destination sealing precedes delete/update compaction; manifest reconciliation happens before serving is enabled. Physical rewrite tests cover named write failures; compaction watermark and merge source admission now fail closed on metadata/filesystem errors. | Other cutover and post-bootstrap discovery rows remain independent blockers. | -| Post-bootstrap discovery | **blocked** | It resumes from the last non-empty bootstrap cursor and creates failed rows for unknown active repositories on the normal path. | Unknown inactive repositories are skipped instead of recorded. A next cursor longer than 512 bytes silently exits the loop and cleanup proceeds, omitting the remaining network. | -| Cleanup durability | **verified** | On the successful path, the backfill tree is removed, the data directory is synced, and merge/discovery cursors are deleted in a synced RocksDB batch afterward. Restart-after-cleanup repeats the directory sync. | This row assumes the source was correctly classified as absent; that guard is separately blocked. | +| Merge-tail compaction and manifest reconcile | **partial** | Destination sealing precedes delete/update compaction; manifest reconciliation happens before serving is enabled. Physical rewrite tests cover named write failures; compaction watermark and merge source admission now fail closed on metadata/filesystem errors. | Archive startup/recovery and exact cutover failure coverage remain independent blockers. | +| Post-bootstrap discovery | **verified** | It resumes from the last non-empty bootstrap cursor, writes every previously unknown active or inactive DID as failed while preserving the relay's active flag, follows dynamically owned cursors, rejects cursor loops, and propagates relay/store errors before cleanup. A real-process relay fixture returns an inactive DID, a 4 KiB cursor, and a second DID; both durable account rows are inspected after steady-state admission. | This row proves discovery completeness and failure behavior; the downstream retry runner remains independently blocked. | +| Cleanup durability | **verified** | On the successful path, the backfill tree is removed, the data directory is synced, and merge/discovery cursors are deleted in a synced RocksDB batch afterward. Restart-after-cleanup repeats the directory sync, and the source-existence guard fails closed on non-not-found errors. | None known for this cleanup ordering invariant. | ## Steady failed-repository retry diff --git a/docs/configuration-parity.md b/docs/configuration-parity.md index 3e2ee0a..6512a3f 100644 --- a/docs/configuration-parity.md +++ b/docs/configuration-parity.md @@ -4,7 +4,8 @@ Authority: - upstream `cmd/jetstream/main.go` and `internal/jetstreamd/options.go` at `f29815c391fc2644f8a3dd36b899fb3697dd1ea6` -- Stream `b709fe32c72e9eee4725f92e7999eb6533bb3c8f` +- Stream `a345f4054df568e364aff08a9b1726275fc48e29` plus the discovery + changes and audit updates in this commit This document audits public configuration and command wiring. A **verified** configuration row proves only that the input is parsed/defaulted and reaches @@ -25,7 +26,7 @@ semantically correct. Status vocabulary is defined in | Backfill workers | **verified** | Explicit positive values set physical getRepo concurrency; zero/omitted default to 100. Held real requests observe 7/100/100. | The 200-worker production experiment showed that accepting a value does not mean it is efficient. | | Backfill batch size | **verified** | Explicit and zero/default values control only page-aligned listRepos accumulation and cursor-checkpoint cadence; zero/omitted select 100,000. Repository completion now follows independent writer durability boundaries. | This control does not define whole-network size or completion percentage. | | Backfill async-flush workers | **partial** | Omitted defaults to four, positive values start that many bootstrap compression workers, and zero selects synchronous compression. Focused tests observe overlap and byte-equivalent output. | Full cancellation/OOM state coverage was not re-established in this audit. | -| Skip merge discovery | **partial** | Explicit configuration and automatic max/selected enablement prevent the discovery request in focused tests. | Normal discovery is blocked by inactive-entry and long-cursor behavior. | +| Skip merge discovery | **verified** | Explicit configuration and automatic max/selected enablement prevent the discovery request in focused tests. Normal discovery independently preserves active/inactive unknown repositories and follows arbitrary-length cursors. | This is a debug/partial-selection control; enabling it for a whole-network crawl intentionally omits accounts born during bootstrap. | | Failed-repo retry interval | **partial** | Duration parsing, zero disablement, and timer-driven pass execution have real-process tests. | Pass failures are hidden from health and candidate scans are unbounded in RAM. | | Failed-repo global/per-host workers | **partial** | Explicit and zero/default values reach real 16-global/4-per-host gates in held-request tests. | Pending-merge retry uses a different path, and the steady candidate set is materialized eagerly. | | Failed-repo maximum delay | **partial** | Parsed values affect persisted exponential-jitter deadlines in focused RocksDB tests. | Host-park read corruption fails open; whole retry supervision is blocked. | diff --git a/docs/semantic-parity.md b/docs/semantic-parity.md index 504b871..65cd748 100644 --- a/docs/semantic-parity.md +++ b/docs/semantic-parity.md @@ -7,8 +7,8 @@ shape, or because a test written around Stream's implementation passes. ## Audit basis -- Stream: `bdb9a796de5fe42b39b79a53e22fcff55a3f94f7` plus the - bootstrap-to-merging transition and audit changes in this commit +- Stream: `a345f4054df568e364aff08a9b1726275fc48e29` plus the discovery + changes and audit updates in this commit - upstream Jetstream: `f29815c391fc2644f8a3dd36b899fb3697dd1ea6` - Atmos: `v0.2.14` - Zat: `649dae356c576a205ec97f2b23671af365756b55` @@ -39,11 +39,11 @@ test, and this audit must be reviewed together. **Stream is not at semantic parity and is not admitted for another whole-network experiment.** -Post-bootstrap discovery, archive recovery, cold replay, ordinary live -encoding, retry-supervisor health, and the firehose reconnect loop remain -deployment blockers. The completed bootstrap durability and correctness- -metadata work does not admit an experiment while those independent blockers -remain. The detailed bootstrap audit is in +Archive recovery, cold replay, pending-repository/retry orchestration, +ordinary live encoding, retry-supervisor health, and the firehose reconnect +loop remain deployment blockers. The completed bootstrap durability, +correctness-metadata, and discovery work does not admit an experiment while +those independent blockers remain. The detailed bootstrap audit is in [bootstrap-semantic-parity.md](bootstrap-semantic-parity.md). ## Audited checklist @@ -52,7 +52,7 @@ remain. The detailed bootstrap audit is in |---|---|---|---| | Lifecycle phase machine | **verified** | Phase reads distinguish absence from RocksDB failure, reject unknown values, and use synced writes. Bootstrap writes `merging` after backfill drains while bootstrap-live is still running, exactly at the pinned upstream boundary; process-kill crashpoints exercise both sides of that write. Cleanup writes `steady_state` only after directory sync and cursor deletion. | Exact persisted phase-entry timestamps are a separate status-surface gap; this row establishes transition state and ordering only. | | Whole-network bootstrap | **partial** | listRepos pages form a page-aligned dispatch/checkpoint unit; eligible repositories are shuffled; download concurrency and CAR preparation are real. Successful repositories become durable independently at archive-writer durability boundaries, including while a sibling remains incomplete and the listRepos cursor remains unchanged. | Resource-exhaustion cleanup at every preparation/emission ownership transfer is not yet proved, and selected/debug paths still need exact-artifact admission receipts. See the detailed audit. | -| Merge and post-bootstrap discovery | **blocked** | Source rows are filtered by the backfill revision; source cursor/latest-revision updates share a synced batch; merge cursor corruption/read errors fail closed; the restart guard treats only `FileNotFound` as cleanup-complete; successful cleanup is directory-synced. | Discovery skips inactive unknown repos and silently stops on cursors longer than 512 bytes. Those can omit repositories while allowing cleanup and phase advancement. | +| Merge and post-bootstrap discovery | **blocked** | Source rows are filtered by the backfill revision; source cursor/latest-revision updates share a synced batch; merge cursor corruption/read errors fail closed; the restart guard treats only `FileNotFound` as cleanup-complete; discovery records active and inactive unknown repositories and follows arbitrary-length cursors with loop detection; successful cleanup is directory-synced. | The pending-repository pass still materializes all DIDs and does not use upstream's bounded retry runner. Latest-revision refresh and source-file failure coverage remain partial. | | Failed-repository healing | **blocked** | A real global/per-host worker gate exists, final redirect hosts are recorded, and retry state is stored. | Pass-level RocksDB/archive failures are logged and hidden from service health. Every failed DID and host gate is materialized in RAM before work. Persisted host parking is not upstream behavior, and park-read errors mean “not parked.” | | JSS sealed format interoperability | **verified** | Upstream-produced sealed fixtures are parsed; Stream-produced sealed files are consumed by the pinned Go reader; header/footer, block index, blooms, collections, compression, and checksums have reciprocal fixtures. | This verifies sealed-format compatibility only. It does not verify startup recovery or replay completeness. | | Archive startup recovery | **partial** | Torn active tails with ordinary decode termination are truncated to the last complete block and sealed. A truncated sealed segment is subsequently rejected by manifest loading in the current composition. | `Archive.recover` swallows file-read/header errors and converts any active-iterator error, including allocation failure, into end-of-valid-data before rewriting the file. Upstream fails loud on these errors and checksum-verifies a sealed high-water mark. Add fault-specific recovery tests before claiming parity. | diff --git a/src/internal/bootstrap/engine.zig b/src/internal/bootstrap/engine.zig index 75e5f6e..624380c 100644 --- a/src/internal/bootstrap/engine.zig +++ b/src/internal/bootstrap/engine.zig @@ -706,14 +706,12 @@ pub fn runTraced( // an old full-network discovery watermark for merge to consume later. if (config.selected_repos.len > 0) try store.meta.deleteDurable(bootstrap_cursor_key); const full_network = config.max_repos == 0 and config.selected_repos.len == 0; - const cursor_owned: ?[]u8 = if (full_network) + var cursor_owned: ?[]u8 = if (full_network) try loadRelayListReposCursor(store.meta, allocator) else null; defer if (cursor_owned) |owned| allocator.free(owned); var cursor: ?[]const u8 = if (cursor_owned) |owned| owned else null; - var last_nonempty_cursor: ?[]const u8 = cursor; - var cursor_buf: [512]u8 = undefined; var batch_arena = std.heap.ArenaAllocator.init(allocator); defer batch_arena.deinit(); var jobs: std.ArrayList(Job) = .empty; @@ -806,10 +804,10 @@ pub fn runTraced( if (config.metrics) |m| _ = m.backfill_discovered_total.fetchAdd(discovered_active, .monotonic); const final_page = page.cursor == null; if (page.cursor) |next| { - if (next.len > cursor_buf.len) return error.ListReposCursorTooLong; - @memcpy(cursor_buf[0..next.len], next); - cursor = cursor_buf[0..next.len]; - last_nonempty_cursor = cursor; + const owned = try allocator.dupe(u8, next); + if (cursor_owned) |old| allocator.free(old); + cursor_owned = owned; + cursor = owned; } const max_reached = config.max_repos > 0 and selected >= config.max_repos; const batch_due = selected_page or max_reached or final_page or (full_network and batch_entries >= batch_size); @@ -819,7 +817,7 @@ pub fn runTraced( if (full_network) try saveListReposCheckpoint( store.meta, if (page.cursor) |next| next else "", - if (last_nonempty_cursor) |last| last else "", + if (cursor) |last| last else "", ); _ = batch_arena.reset(.retain_capacity); jobs = .empty; diff --git a/src/internal/bootstrap/lifecycle.zig b/src/internal/bootstrap/lifecycle.zig index f71ebd2..aaef4bf 100644 --- a/src/internal/bootstrap/lifecycle.zig +++ b/src/internal/bootstrap/lifecycle.zig @@ -407,28 +407,33 @@ fn runDiscovery( try seen.put(seen_arena.allocator(), try seen_arena.allocator().dupe(u8, start), {}); var discovered: u64 = 0; + var cursor_owned: ?[]u8 = null; + defer if (cursor_owned) |owned| allocator.free(owned); var cursor: ?[]const u8 = start; - var cursor_buf: [512]u8 = undefined; while (true) { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const page = try engine.listReposPage(arena.allocator(), &client, cursor); for (page.repos) |ref| { - if (!ref.active) continue; if (try repo_st.get(allocator, ref.did)) |state| { repo_st.freeState(allocator, state); continue; } - try repo_st.put(ref.did, .{ .status = .failed }); + try repo_st.put(ref.did, .{ + .status = .failed, + .active = ref.active, + .last_error = "discovered post-bootstrap; queued for retry", + }); discovered += 1; _ = stats.orchestrator_merge_dids_discovered_post_bootstrap_total.fetchAdd(1, .monotonic); } const next = page.cursor orelse break; if (seen.contains(next)) return error.ListReposCursorLoop; try seen.put(seen_arena.allocator(), try seen_arena.allocator().dupe(u8, next), {}); - if (next.len > cursor_buf.len) break; - @memcpy(cursor_buf[0..next.len], next); - cursor = cursor_buf[0..next.len]; + const owned = try allocator.dupe(u8, next); + if (cursor_owned) |old| allocator.free(old); + cursor_owned = owned; + cursor = owned; } try repo_st.flush(); log.info("merge discovery: {d} post-bootstrap DIDs queued for retry", .{discovered}); diff --git a/tests/upstream_oracle/stream_differential_test.go b/tests/upstream_oracle/stream_differential_test.go index b43dc3d..538b461 100644 --- a/tests/upstream_oracle/stream_differential_test.go +++ b/tests/upstream_oracle/stream_differential_test.go @@ -7,6 +7,8 @@ package oracle import ( "bytes" "context" + "fmt" + "io" "log/slog" "math/rand/v2" "net" @@ -45,6 +47,89 @@ type listReposCountingHandler struct { 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 @@ -256,6 +341,9 @@ func TestStreamBootstrapControlsOracle(t *testing.T) { 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() @@ -349,6 +437,122 @@ func TestStreamBootstrapControlsOracle(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) { -- 2.51.2