From 6ce5892244d49be8f7828afcd5dda8a3370ff3f9 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Sun, 6 Sep 2026 16:02:04 -0500 Subject: [PATCH] client: align archive recovery, pagination, and filtering with upstream --- build.zig.zon | 4 +- docs/client.md | 23 ++ docs/failover.md | 11 +- docs/upstream-review-2026-09-06.md | 249 ++++++++++++++++++++++ src/archive_backfill.zig | 252 ++++++++++++++++++++-- src/client.zig | 19 +- src/filter.zig | 28 +++ src/live.zig | 2 + src/loopback_test.zig | 329 +++++++++++++++++++++++++++-- 9 files changed, 871 insertions(+), 46 deletions(-) create mode 100644 docs/upstream-review-2026-09-06.md create mode 100644 src/filter.zig diff --git a/build.zig.zon b/build.zig.zon index 43ffb7f..c77fafe 100644 --- a/build.zig.zon +++ b/build.zig.zon @@ -9,8 +9,8 @@ .hash = "N-V-__8AAPZ7fwBg4JoCzM_0o2A8wxH2hsUUeiU1iuZv53L5", }, .zat = .{ - .url = "https://tangled.org/zat.dev/zat/archive/v0.4.4.tar.gz", - .hash = "zat-0.4.4-5PuC7u-qDAB8PpFGomEQCOtRxg00sXotZwh3hHcz-02K", + .url = "https://tangled.org/zat.dev/zat/archive/v0.5.0.tar.gz", + .hash = "zat-0.5.0-5PuC7sUhDAAgtrkdHxFwI0vq4Qfy7PY4ZYYATu_pOmS3", }, .websocket = .{ .url = "https://tangled.org/zzstoatzz.io/websocket.zig/archive/v0.1.12.tar.gz", diff --git a/docs/client.md b/docs/client.md index 9f1f0ad..65646a7 100644 --- a/docs/client.md +++ b/docs/client.md @@ -15,6 +15,13 @@ per-event seq dedup makes the seam at-least-once with no gap, and a `snapshot_only` stops at the sealed range. without `after_seq` it is a pure live tail (`live_cursor` resumes at a known seq). +The archive sweep follows every plan page, pinning the first page's sealed +tip as the upper bound. A plan that does not advance fails with `PlanStalled`. +Collection filters accept exact names and namespace wildcards such as +`app.bsky.feed.*`; archive and live delivery both enforce filters locally. +DID-level markers bypass collection filters and remain subject to DID and +kind filters. + multi-host failover covers the live tail: `hosts` takes a list (first entry primary, matching zat's v1 `JetstreamClient`/`FirehoseClient` contract); rotation kicks in when the current host is unhealthy, re-anchoring across per-instance seq spaces @@ -41,6 +48,22 @@ only, kept for existing consumers) and the unified client's identity/account/sync markers). `fetchSeqBounds` maps a witnessed-time window to plan seq bounds. +## archive errors + +An optional `onError(anyerror) bool` callback chooses whether to continue +after a recoverable archive error. Returning true accepts the error and +continues on the same instance; false stops cleanly. Omitting the callback +returns the error to the caller, so data is never skipped without a decision. + +A failed block download emits the entry's good prefix, reports the error, +and skips its remaining blocks before proceeding to the next entry. A +decompression failure skips that block; a malformed record preserves other +valid rows from the block and reports the decode error after them. Whole +segment download/header failures can likewise proceed to the next entry. +Plan failures, cancellation, and allocation failure remain terminal. +Accepting an error means accepting an incomplete portion of the archive; +it does not repair that portion or trigger failover. + ## records commit events carry both representations, matching upstream `Event`: diff --git a/docs/failover.md b/docs/failover.md index b9d1d01..8860e9e 100644 --- a/docs/failover.md +++ b/docs/failover.md @@ -47,11 +47,12 @@ host's own archive metadata: - **live half**: `failover_after_errors` consecutive transient failures with no events between (an event between failures resets the count — a limping host is not a dead one). -- **archive half**: never. a failing sweep (bounded internal retries via - `fetchWithRetry`, 4 attempts) is terminal — archives are not - interchangeable across hosts (coverage, retention, and holes differ), - so a mid-sweep switch could silently anchor past a gap. the caller - decides what a dead archive means. +- **archive half**: never. after bounded download retries, `onError` + reports the failure. If the caller returns true, the sweep continues on + the same host with the next entry; false stops cleanly. Decode errors can + preserve later blocks and valid rows. Plan failures remain terminal. + Archives are not interchangeable across hosts (coverage, retention, and + holes differ), so an archive failure never triggers a host switch. `failover_after_errors = 0` or a single-entry `hosts` disables rotation entirely; the single-host error surface is exactly the pre-failover one. diff --git a/docs/upstream-review-2026-09-06.md b/docs/upstream-review-2026-09-06.md new file mode 100644 index 0000000..9d83d2d --- /dev/null +++ b/docs/upstream-review-2026-09-06.md @@ -0,0 +1,249 @@ +# Upstream client review + +Compared the candidate after v0.1.1 with the official Go client at +[`58c4d7f7a9130e53b40348ad3d1f7aafed0e4843`](https://github.com/bluesky-social/jetstream/tree/58c4d7f7a9130e53b40348ad3d1f7aafed0e4843), +the upstream main tip observed September 6, 2026. The local documentation +targets `289b032`; the core pagination, filtering, raw-record, and error +contracts discussed below also exist at that older revision. This is a +source comparison, not an exhaustive conformance certification. + +## Initial findings, now covered and corrected locally + +- **Archive plan pagination.** Upstream + [`sweepSealedArchive`](https://github.com/bluesky-social/jetstream/blob/58c4d7f7a9130e53b40348ad3d1f7aafed0e4843/client_core.go) + pins the first `sealedTipSeq`, requests subsequent pages with + `afterSeq=plannedThroughSeq` and `beforeSeq=pinnedTip`, and rejects + non-advancing plans. Our `ArchiveBackfill.run` fetches one plan only. + `client.runOnHost` stops snapshot-only subscriptions after that call, + or starts live at that page's `planned_through_seq`. A multi-page + snapshot therefore reports success before all its history is delivered; + replay-to-live relies on cold replay/rebackfill rather than completing + the sealed sweep. `examples/slice_sync.zig` also makes one archive call. + Add a tiny two-page loopback fixture, including a growing server tip and + a non-advancing page, before correcting the loop. +- **Collection wildcard filtering.** Upstream + [`matcher`](https://github.com/bluesky-social/jetstream/blob/58c4d7f7a9130e53b40348ad3d1f7aafed0e4843/filter.go) + accepts namespace wildcards such as `app.bsky.feed.*`. Our archive + decoder uses exact `containsString` comparisons for collections. It + sends the wildcard to the planner but discards the matching commit rows + locally. Cover exact, wildcard, nonmatching sibling prefixes, and + DID-level markers with small in-memory blocks. +- **Live filter backstop.** Upstream applies its kind/DID/collection + matcher to live events even after requesting server-side filtering. + Our live path sends the filters but `FilteringHandler` checks only kind + and the failover time floor. Out-of-filter DIDs or collections from an + over-inclusive server can reach the caller. Share the collection/DID + predicate across archive and live, preserving unconditional DID-level + markers under collection filters. + +## Deliberate differences or decisions still needed + +- **Archive errors:** upstream's downloader emits recoverable download and + decode errors in order alongside valid rows; a consumer can continue. + Plan/context failures remain terminal. Our archive path aborts on a + download/decode error after applicable retries. The unreleased change + stops converting these failures into cross-host rotation; it does not + introduce the underlying terminal error policy. Retaining this stricter + policy is a possible API choice, but is not Go-client error parity. +- **Failover:** our multi-host feature is an extension; the official client + binds one service. Not transferring sequence cursors across instances + is correct. Witnessed-time re-anchoring is approximate, as already + documented in `failover.md`, and cannot promise lossless equivalence + across independently populated archives. +- **Authentication:** upstream `WithAPIKey` authenticates archive planning + and downloads, explicitly excluding the public dictionary and live + WebSocket. Our live WebSocket also receives the key. The local claim + that this matches upstream is incorrect. Decide whether to preserve it + as a documented hosted-service extension or separate live credentials. +- **Compression:** upstream enables live dictionary-zstd by default; ours + is opt-in. Both are valid wire modes, but defaults differ. +- **Retry hints:** upstream bulk-download retries honor server retry + deferrals. Our `fetchWithRetry` uses fixed 1/2/4-second delays and does + not consume the HTTP library's rate-limit hints. + +## Assessment of the unreleased additions + +- `decode_records=false` follows the intent of upstream + [`WithRawRecords`](https://github.com/bluesky-social/jetstream/blob/58c4d7f7a9130e53b40348ad3d1f7aafed0e4843/options.go): + skip map decoding and expose stored bytes. Our option is archive-only, + and the callback lifetime differs from Go's batch lifetime. The comment + saying upstream always decodes records is stale. Default decoding remains + enabled. Add coverage for raw bytes, null decoded records, and deletes. +- `bytes_fetched` counts successful returned block/segment body bytes. + It is useful payload accounting, not a complete network-byte measure: + it excludes HTTP overhead, plan requests, and failed/retried transfers, + and HTTP content decoding can change the returned size. +- The local segment-to-TSV example does not change the wire protocol. + Network-wide history selection is legitimate, but the history example + inherits the pagination limitation above. +- The ZAT v0.5.0 adoption resolves to ZTTP v0.1.0. The archive plan POST + now explicitly declares `application/json`, satisfying the new transport + contract. Fake-server shutdown and EOF fixes affect tests only. + +## Initial evidence and limits + +The candidate's 39 existing tests pass in about one second of execution, +and all examples compile. A bounded live failover smoke delivered ten +events from `stream.waow.tech` in 1.4 seconds. Those checks do not establish +multi-page, wildcard, or complete upstream behavioral parity. No large +archive download, upstream publication, or release was performed for this +review. The follow-up below records the fixtures and fixes; release remains unpublished. + + +## Introducing commits + +| Difference | Origin in our history | +|---|---| +| Single plan page and exact-only archive collection matching | ZAT `f69d75c9f4afa00e80dd1e6488462fd717d5b747`, August 7: initial archive client. Carried into this repo in `6832774`, August 13. | +| Snapshot stopping after one page; missing live DID/collection backstop | `90f70689c26f26f5be5849ced4d249a29725ce0b`, August 15: initial unified v2 client. | +| Terminal archive errors | Original archive implementation; unified single-host behavior inherited it August 15. `8e74b87`, August 16, added cross-host rotation for these errors; `e01e1867fc366d47d9c0fd099747201a4c3d73df`, August 19, removed that rotation. It did not introduce a recoverable per-entry error API. | +| API key sent on the live handshake | `d215db5c62b84655f4993ace12ccb62da4f4cf1d`, August 17, shipped in v0.1.1. Its claim of Go-client parity was incorrect. | +| Archive raw-record option | `30fa69f`, August 29; analogous to upstream raw mode, not a protocol violation. | + +These are omissions/extensions relative to behavior already present at our +pinned upstream revision `289b032`, not new requirements introduced by the +September upstream tip. The earliest local archive was a smaller API; +August 15's unified-client parity claim is where its incomplete pagination +became an explicit mismatch with the advertised full replay/snapshot API. + +## Tests first, then fixes + +Read upstream `filter_test.go` (`TestMatcherCollectionExactAndWildcard`, +`TestMatcherWildcardBoundary`, marker and DID cases), `engine_test.go` +(`TestEngineMultiPageBackfillCutover`, `TestEnginePinnedBeforeSeqAcrossPages`, +`TestEngineLiveOnlyAppliesCollectionFilter`, `TestEngineFilterContractAcrossArchiveAndLive`), +and `downloader_test.go` (`TestDownloadBlocksErrorStopsEntryNotPlan`). +The assertions, not stale adjacent comments, establish expected behavior. + +Before production fixes, four new tests compiled and ran: **40/43 passed**. +The three failures were the expected semantic mismatches: + +- Snapshot: expected three records, got the first one only. +- Archive wildcard: expected sequences 1, 2, 4; got 1, 4. +- Live filters: expected only sequence 3; received 1, 2, 3. +- Raw-record mode passed: undecodable raw bytes are preserved without map + decoding, deletes have no record, and default decoding rejects the same + malformed create. + +The fixes page through the first sealed tip, reject stalled continuation, +apply namespace matching with the dot boundary preserved, and reapply kind, +DID, and collection filters on live delivery. Additional fixtures check a +server whose tip grows after page one and a cumulative block budget across +pages. **45/45 tests now pass**, in Debug and ReleaseFast. Examples compile. +No tests were skipped to obtain the passing result. + +The separate archive-error continuation policy, live-auth extension, +compression default, and retry-hint differences above are still present. +The incorrect live-auth and raw-mode parity comments were corrected. + +## Performance and Zig Zen + +On this machine with Zig 0.16.0, Debug test execution is approximately one +second (about three seconds of compilation when changed), versus 31 seconds +before the fake-server EOF fix. The shutdown response fix also avoids +transport retry waits. Tests use small local fixtures; no bulk archive +transfer is part of this check. + +A temporary ReleaseFast probe performed five million collection checks per +case over four names, three trials. The old exact-only matcher measured +0.50–1.22 ns/check; the new matcher measured 1.66–2.88 ns/check for exact +patterns and 2.11–2.79 ns/check with a wildcard. Both used two patterns and +produced the same hit count on the matched workload. These synthetic numbers +show the small absolute cost of the extra checks, not whole-client +throughput or a guarantee across filter sizes. Matching performs no heap +allocations. Pagination frees each page's buffers before fetching the next, +so retaining all plan pages does not become an additional memory cost; +it does create fresh transport state per page rather than retaining pools. + +Ran `zig zen` and reviewed the diff under the local Zen skill: + +- "Communicate intent precisely": pinning the first tip and reporting + `PlanStalled` make the replay boundary explicit. +- "Edge cases matter": wildcard namespace boundaries, DID-level markers, + EOF, moving tips, and cumulative limits have concrete coverage. +- "Memory is a resource": matching allocates nothing; page lifetimes and + cleanup use the existing allocator/defer model. +- "Incremental improvements": the fixes preserve the callback API and + separate the local extensions from claims about upstream compatibility. + +This is a source review, not an automated idiomatic-style certification. + +## Rationale audit against the archive server + +Checked current client code and history alongside the local archive server +at `1765e48`, including `src/internal/serve/xrpcapi.zig`, the WebSocket +handshake in `server.zig`, `deploy/site/Caddyfile`, configuration, and gotchas. +This is a source/configuration audit, not a check of deployed runtime state. + +1. **Terminal archive errors:** the August 19 commit and current failover + document justify *not switching archives*, because their coverage and + sequence spaces differ. That does not establish why a caller cannot + receive an ordered error and elect to continue on the same archive. + No separate rationale for that restriction was found. The server's + cold-WebSocket reader deliberately fails on missing segments, but that + is a different layer from the SDK's HTTP block-download iterator. + The server's August 29 compaction note also records temporarily stale + manifest checksums and recommends re-listing later. That supports + exposing recoverable work to a caller; it does not authorize silently + skipping a hole or changing hosts. Upstream continuation is conditional + on the consumer accepting the error, not unconditional success. +2. **Multi-host live failover:** a rationale *was* found: the August 16 + front-end outage, with an explicit sequence-to-witnessed-time bridge. + Current code only requests rotation for live failures, and propagates + archive errors. Following a live rotation can perform archive reads on + the new host to re-anchor; that is distinct from rotating in response to + an archive failure. This remains an intentional local extension, not a + requirement of the archive server or exact upstream equivalence. +3. **Live API key:** August 17's introducing commit incorrectly claims + upstream Go-client parity. The server gates planSnapshot, listSegments, + getSegment, and getBlock; the live upgrade does not use that gate. The + checked Caddy configuration adds rate limiting, not a live bearer-auth + requirement. No implementation-backed need for the extension was found. +4. **Compression default:** opt-in dates to August 15's initial unified + client. The server supports dictionary-zstd and uncompressed v2 frames. + No documented measurement or server limitation explaining the default + difference was found. Before aligning it, cover negotiation and rotation: + this server currently emits a plain-text 400 for an unknown dictionary, + whereas our SDK classifies the upstream structured + `UnknownZstdDictionary` error. Do not confuse absence of a documented + rationale with proof that every compression edge already interoperates. +5. **Retry delays:** the archive server's August 17 metered-key feature + explicitly relies on the SDK's existing 1/2/4-second 429 backoff. It + returns a 429 without a computed Retry-After in that path, so a fallback + delay is needed. Nothing there justifies ignoring a valid retry hint + from another server or proxy. Honor hints with a finite cap and retain + fallback backoff when no hint is supplied. + +No behavior changes were made during this rationale audit. In particular, +archive failover was not restored. The supported extension is live failover; +three other differences lack a documented local need, and fixed retry +backoff is justified only as the fallback when a hint is absent. + + +## Archive error policy aligned after the rationale audit + +Nate requested upstream semantics for archive errors, separately from the +other remaining differences. Added failing tests first: the mixed-record +block returned `MalformedRecord` before delivering its later valid row, and +an entry download failure returned `FetchFailed` instead of invoking the +caller's error callback. The terminal-plan assertion already held. + +The client now forwards recoverable archive errors to `onError(anyerror) +bool. Accepting continues on the same server; declining stops cleanly. +Absent a callback, errors are returned. A failed block download ends that +entry after its good prefix; decode failures preserve later readable blocks, +and malformed records preserve valid rows within their block. Whole-segment +failures follow the same entry policy. Plan failures, cancellation, and +allocation failures remain terminal. Pending unused requests are canceled +and released on stop rather than awaited through their retry budgets. + +The loopback tests exercise acceptance and refusal at concurrency 1 and 4, +with a second unavailable host configured to detect accidental archive +failover. Additional cases cover whole-segment errors, malformed records, +terminal planning, absent callbacks, cancellation, and allocation errors. +**50/50 tests pass**, about two seconds of Debug execution with small local +fixtures. This supersedes the terminal-archive-error difference earlier in +this historical review. Live failover, live authentication, compression +configuration, and retry hints were not changed in this step. Ran `zig zen` +again: error decisions and resource cleanup are explicit, and the new +record-error bookkeeping adds no per-record allocation. diff --git a/src/archive_backfill.zig b/src/archive_backfill.zig index 61f35d8..b166496 100644 --- a/src/archive_backfill.zig +++ b/src/archive_backfill.zig @@ -35,6 +35,7 @@ const sync = zat.firehose; // re-exports sync.CommitAction const cbor = zat.cbor; const zstd = @import("zstd.zig"); const livedecode = @import("livedecode.zig"); +const filter = @import("filter.zig"); const multibase = zat.multibase; const HttpTransport = zat.HttpTransport; @@ -77,8 +78,8 @@ pub const Options = struct { /// (take it from the prior Result.position). segment indices are stable /// for a given archive, so a re-plan covers at least the same range. start_after: ?Position = null, - /// decode each commit's record CBOR into Event.record (upstream always - /// does). false skips the decode and leaves record null — for consumers + /// Decode each commit's record CBOR into Event.record. Like upstream's + /// raw-record mode, false skips decoding and leaves record null — for consumers /// that only want (seq, did, collection, rkey, op), it avoids the cost /// and means one undecodable stored payload cannot stop the sweep. decode_records: bool = true, @@ -125,8 +126,44 @@ pub const Result = struct { /// replay the archive through `handler` (same contract as JetstreamClient: /// `fn onEvent(*H, Event) void`). event slices are only valid during the -/// onEvent call. blocks until the plan is exhausted. +/// onEvent call. Pages through the sealed tip observed on the first plan, +/// unless the handler stops or max_blocks is reached. PlanStalled means +/// the server returned a continuation that cannot advance the sweep. +/// Optional onError(anyerror) bool accepts a recoverable archive error with +/// true or stops cleanly with false. Without it, the error is returned. pub fn run(io: Io, allocator: Allocator, options: Options, handler: anytype) !Result { + var page_options = options; + var pinned_tip: ?u64 = null; + var result: Result = .{ .planned_through_seq = options.after_seq orelse 0, .sealed_tip_seq = 0 }; + while (true) { + if (options.max_blocks) |max| { + if (result.blocks_decoded >= max) { + result.truncated = true; + return result; + } + page_options.max_blocks = max - result.blocks_decoded; + } + const page = try runPage(io, allocator, page_options, handler); + if (pinned_tip == null) pinned_tip = page.sealed_tip_seq; + result.sealed_tip_seq = pinned_tip.?; + result.planned_through_seq = page.planned_through_seq; + result.events_delivered += page.events_delivered; + result.blocks_decoded += page.blocks_decoded; + result.bytes_fetched += page.bytes_fetched; + if (page.position) |position| result.position = position; + if (page.last_time_us) |time| result.last_time_us = time; + result.stopped = page.stopped; + result.truncated = page.truncated; + if (page.stopped or page.truncated or page.planned_through_seq >= pinned_tip.?) return result; + if (page.planned_through_seq <= (page_options.after_seq orelse 0)) return error.PlanStalled; + page_options.after_seq = page.planned_through_seq; + // Snapshot membership is fixed by page one, even if the server seals + // more data while downloading. The live tail covers that later data. + page_options.before_seq = pinned_tip; + } +} + +fn runPage(io: Io, allocator: Allocator, options: Options, handler: anytype) !Result { var transport = HttpTransport.init(io, allocator); defer transport.deinit(); @@ -194,10 +231,23 @@ pub fn run(io: Io, allocator: Allocator, options: Options, handler: anytype) !Re .{ options.host, segment.name }, ); defer allocator.free(url); - var fetched = try fetchWithRetry(io, &transport, url, authorization); + var fetched = fetchWithRetry(io, &transport, url, authorization) catch |err| { + reportArchiveError(handler, err) catch |reported| { + if (reported != error.Stopped) return reported; + result.stopped = true; + return result; + }; + continue; + }; defer fetched.deinit(allocator); result.bytes_fetched += fetched.body.len; - try deliverSegment(allocator, fetched.body, options, segment, handler, &result); + deliverSegment(allocator, fetched.body, options, segment, handler, &result) catch |err| { + reportArchiveError(handler, err) catch |reported| { + if (reported != error.Stopped) return reported; + result.stopped = true; + return result; + }; + }; } } @@ -216,6 +266,23 @@ const max_concurrency = 16; const FetchError = anyerror; +/// No callback means no permission to skip data. Cancellation and allocator +/// failure are not corrupt archive entries and must never become recoverable. +fn reportArchiveError(handler: anytype, err: anyerror) !void { + switch (err) { + error.Canceled, error.OutOfMemory, error.Stopped => return err, + else => {}, + } + if (comptime @hasDecl(@TypeOf(handler.*), "onError")) { + const decision = handler.onError(err); + // Legacy archive handlers only log errors. Keep them source-compatible + // without treating a notification as permission to skip an entry. + if (@TypeOf(decision) == void) return err; + const accepted = if (@TypeOf(decision) == bool) decision else try decision; + if (!accepted) return error.Stopped; + } else return err; +} + fn fetchJob(io: Io, transport: *HttpTransport, url: []const u8, authorization: ?[]const u8) FetchError!HttpTransport.FetchResult { return fetchWithRetry(io, transport, url, authorization); } @@ -233,6 +300,7 @@ const FetchWindow = struct { authorization: ?[]const u8 = null, head: usize = 0, count: usize = 0, + failed_segment: ?u32 = null, const Inflight = struct { future: std.Io.Future(FetchError!HttpTransport.FetchResult), @@ -260,6 +328,7 @@ const FetchWindow = struct { } fn start(self: *FetchWindow, options: Options, segment_name: []const u8, position: Position) !void { + if (self.failed_segment == position.segment_index) return; std.debug.assert(!self.isFull()); const url = try std.fmt.allocPrint( self.allocator, @@ -282,10 +351,22 @@ const FetchWindow = struct { self.head = (self.head + 1) % self.slots.len; self.count -= 1; defer self.allocator.free(slot.url); - var fetched = try slot.future.await(self.io); + if (self.failed_segment == slot.position.segment_index) { + if (slot.future.cancel(self.io)) |fetched| { + var discarded = fetched; + discarded.deinit(allocator); + } else |_| {} + return; + } + var fetched = slot.future.await(self.io) catch |err| { + self.failed_segment = slot.position.segment_index; + return reportArchiveError(handler, err); + }; defer fetched.deinit(allocator); result.bytes_fetched += fetched.body.len; - try deliverFrame(allocator, fetched.body, options, slot.position, handler, result); + deliverFrame(allocator, fetched.body, options, slot.position, handler, result) catch |err| { + try reportArchiveError(handler, err); + }; } /// drain in order, stopping (and discarding the rest) if max_blocks hits @@ -299,13 +380,14 @@ const FetchWindow = struct { } } - /// await and drop everything in flight without decoding + /// Cancel and release pending requests when the caller stops or a budget + /// is reached; cleanup must not wait through unused retry schedules. fn discardAll(self: *FetchWindow) void { while (self.count > 0) { const slot = &self.slots[self.head]; self.head = (self.head + 1) % self.slots.len; self.count -= 1; - if (slot.future.await(self.io)) |fetched| { + if (slot.future.cancel(self.io)) |fetched| { var discarded = fetched; discarded.deinit(self.allocator); } else |_| {} @@ -434,11 +516,12 @@ fn fetchPlan(arena: Allocator, transport: *HttpTransport, options: Options, auth .url = url, .method = .POST, .payload = body, + .content_type = "application/json", .authorization = authorization, }); defer fetched.deinit(transport.allocator); if (fetched.status != .ok) { - log.err("planSnapshot failed: {d} {s}", .{ @intFromEnum(fetched.status), fetched.body }); + log.warn("planSnapshot failed: {d} {s}", .{ @intFromEnum(fetched.status), fetched.body }); return error.PlanFailed; } return parsePlan(arena, fetched.body); @@ -684,10 +767,12 @@ fn deliverSegment(allocator: Allocator, body: []const u8, options: Options, segm const compressed_size = mem.readInt(u32, body[entry_off + 8 ..][0..4], .little); const frame_start = offset + 8; if (frame_start + compressed_size > body.len) return error.MalformedSegment; - try deliverFrame(allocator, body[frame_start..][0..compressed_size], options, .{ + deliverFrame(allocator, body[frame_start..][0..compressed_size], options, .{ .segment_index = segment.index, .block_index = @intCast(bi), - }, handler, result); + }, handler, result) catch |err| { + try reportArchiveError(handler, err); + }; } } @@ -790,6 +875,7 @@ fn deliverBlock(arena: Allocator, raw: []const u8, options: Options, handler: an var rkey_off = did_off + did_total; var rev_off = rkey_off + rkey_total; var pay_off = rev_off + rev_total; + var record_error: ?anyerror = null; for (0..n) |i| { const col_len = raw[col_len_base + i]; @@ -829,14 +915,18 @@ fn deliverBlock(arena: Allocator, raw: []const u8, options: Options, handler: an if (options.dids.len > 0 and !containsString(options.dids, did)) continue; const row_event: ?livedecode.Event = switch (kind_byte) { kind_create, kind_create_resync, kind_update, kind_delete => blk: { - if (options.collections.len > 0 and !containsString(options.collections, collection)) break :blk null; + if (!filter.collectionMatches(options.collections, collection)) break :blk null; const op: livedecode.Operation = switch (kind_byte) { kind_update => .update, kind_delete => .delete, else => .create, }; const rec: ?json.Value = if (options.decode_records and op != .delete and payload.len > 0) rec: { - const decoded = cbor.decode(arena, payload) catch return error.MalformedRecord; + const decoded = cbor.decode(arena, payload) catch |err| { + if (err == error.OutOfMemory) return err; + record_error = error.MalformedRecord; + continue; + }; break :rec try livedecode.cborToJson(arena, decoded.value); } else null; break :blk .{ @@ -864,7 +954,11 @@ fn deliverBlock(arena: Allocator, raw: []const u8, options: Options, handler: an // (e.g. an async-resync #sync) default to zero values, // matching upstream's generated decoder const env: ?cbor.Value = if (payload.len > 0) - (cbor.decode(arena, payload) catch return error.MalformedRecord).value + (cbor.decode(arena, payload) catch |err| { + if (err == error.OutOfMemory) return err; + record_error = error.MalformedRecord; + continue; + }).value else null; break :blk switch (kind_byte) { @@ -909,11 +1003,15 @@ fn deliverBlock(arena: Allocator, raw: []const u8, options: Options, handler: an }; // planned blocks can contain other collections' rows — always filter client-side - if (options.collections.len > 0 and !containsString(options.collections, collection)) continue; + if (!filter.collectionMatches(options.collections, collection)) continue; if (options.dids.len > 0 and !containsString(options.dids, did)) continue; const record: ?json.Value = if (options.decode_records and operation != .delete and payload.len > 0) blk: { - const decoded = cbor.decode(arena, payload) catch return error.MalformedRecord; + const decoded = cbor.decode(arena, payload) catch |err| { + if (err == error.OutOfMemory) return err; + record_error = error.MalformedRecord; + continue; + }; break :blk try livedecode.cborToJson(arena, decoded.value); } else null; @@ -931,6 +1029,7 @@ fn deliverBlock(arena: Allocator, raw: []const u8, options: Options, handler: an result.events_delivered += 1; result.last_time_us = witnessed_at; } + if (record_error) |err| try reportArchiveError(handler, err); } fn containsString(haystack: []const []const u8, needle: []const u8) bool { @@ -1385,3 +1484,122 @@ test "deliverBlock enforces the exact (after_seq, before_seq] window per row" { }, &handler, &result); try testing.expectEqualSlices(u64, &.{ 11, 12 }, handler.seqs.items); } + +// Upstream filter_test.go: TestMatcherCollectionExactAndWildcard, +// TestMatcherWildcardBoundary, and marker/DID filter contracts. +test "conformance: archive wildcards preserve namespace boundaries and markers" { + var arena_state = std.heap.ArenaAllocator.init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const raw = try buildTestBlock(arena, &.{ + .{ .seq = 1, .witnessed_at = 1, .kind = kind_create, .collection = "app.bsky.feed.post", .did = "did:plc:a", .rkey = "r", .rev = "r", .payload = &test_record_cbor }, + .{ .seq = 2, .witnessed_at = 2, .kind = kind_create, .collection = "app.bsky.graph.follow", .did = "did:plc:a", .rkey = "r", .rev = "r", .payload = &test_record_cbor }, + .{ .seq = 3, .witnessed_at = 3, .kind = kind_create, .collection = "app.bsky.graphient.thing", .did = "did:plc:a", .rkey = "r", .rev = "r", .payload = &test_record_cbor }, + .{ .seq = 4, .witnessed_at = 4, .kind = kind_account, .collection = "$account", .did = "did:plc:a", .rkey = "", .rev = "", .payload = "" }, + .{ .seq = 5, .witnessed_at = 5, .kind = kind_identity, .collection = "$identity", .did = "did:plc:other", .rkey = "", .rev = "", .payload = "" }, + }); + const Handler = struct { + seqs: std.ArrayList(u64) = .empty, + allocator: Allocator, + pub fn onRow(self: *@This(), event: livedecode.Event) bool { + self.seqs.append(self.allocator, event.seq) catch @panic("out of memory"); + return true; + } + }; + var handler = Handler{ .allocator = arena }; + var result: Result = .{ .planned_through_seq = 0, .sealed_tip_seq = 0 }; + try deliverBlock(arena, raw, .{ + .host = "", + .collections = &.{ "app.bsky.feed.post", "app.bsky.graph.*" }, + .dids = &.{"did:plc:a"}, + }, &handler, &result); + try testing.expectEqualSlices(u64, &.{ 1, 2, 4 }, handler.seqs.items); +} + +// Upstream WithRawRecords/decodeCommitInto: skip map materialization, +// preserve raw bytes for creates, and provide no record for deletes. +test "conformance: raw archive records skip decoding and preserve bytes" { + var arena_state = std.heap.ArenaAllocator.init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const invalid_cbor = "\xff"; + const raw = try buildTestBlock(arena, &.{ + .{ .seq = 1, .witnessed_at = 1, .kind = kind_create, .collection = "app.bsky.feed.post", .did = "did:plc:a", .rkey = "r", .rev = "r", .payload = invalid_cbor }, + .{ .seq = 2, .witnessed_at = 2, .kind = kind_delete, .collection = "app.bsky.feed.post", .did = "did:plc:a", .rkey = "r", .rev = "r", .payload = "" }, + }); + const Handler = struct { + count: usize = 0, + valid: bool = true, + pub fn onRow(self: *@This(), event: livedecode.Event) bool { + const commit = event.payload.commit; + self.valid = self.valid and commit.record == null and mem.eql(u8, if (commit.operation == .delete) "" else invalid_cbor, commit.record_cbor); + self.count += 1; + return true; + } + }; + var handler = Handler{}; + var result: Result = .{ .planned_through_seq = 0, .sealed_tip_seq = 0 }; + try deliverBlock(arena, raw, .{ .host = "", .decode_records = false }, &handler, &result); + try testing.expectEqual(@as(usize, 2), handler.count); + try testing.expect(handler.valid); + try testing.expectError(error.MalformedRecord, deliverBlock(arena, raw, .{ .host = "" }, &handler, &result)); +} + +// Upstream TestDownloadTransformKeepsPayloadWithMalformedRecordError: +// valid rows from a mixed block survive and the decode error is reported. +test "archive malformed records preserve valid rows before reporting the error" { + var arena_state = std.heap.ArenaAllocator.init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const raw = try buildTestBlock(arena, &.{ + .{ .seq = 1, .witnessed_at = 1, .kind = kind_create, .collection = "app.bsky.feed.post", .did = "did:plc:a", .rkey = "r", .rev = "r", .payload = &test_record_cbor }, + .{ .seq = 2, .witnessed_at = 2, .kind = kind_create, .collection = "app.bsky.feed.post", .did = "did:plc:a", .rkey = "r", .rev = "r", .payload = "\xff" }, + .{ .seq = 3, .witnessed_at = 3, .kind = kind_create, .collection = "app.bsky.feed.post", .did = "did:plc:a", .rkey = "r", .rev = "r", .payload = &test_record_cbor }, + }); + const Handler = struct { + items: [4]u64 = undefined, + len: usize = 0, + err: ?anyerror = null, + pub fn onRow(self: *@This(), event: livedecode.Event) bool { + self.items[self.len] = event.seq; + self.len += 1; + return true; + } + pub fn onError(self: *@This(), err: anyerror) bool { + self.err = err; + self.items[self.len] = 0; + self.len += 1; + return true; + } + }; + var handler = Handler{}; + var result: Result = .{ .planned_through_seq = 0, .sealed_tip_seq = 0 }; + try deliverBlock(arena, raw, .{ .host = "" }, &handler, &result); + try testing.expectEqualSlices(u64, &.{ 1, 3, 0 }, handler.items[0..handler.len]); + try testing.expectEqual(@as(?anyerror, error.MalformedRecord), handler.err); +} + +test "archive recovery never swallows cancellation allocation failure or absent callbacks" { + const Handler = struct { + errors: usize = 0, + pub fn onError(self: *@This(), _: anyerror) bool { + self.errors += 1; + return true; + } + }; + var handler = Handler{}; + try testing.expectError(error.Canceled, reportArchiveError(&handler, error.Canceled)); + try testing.expectError(error.OutOfMemory, reportArchiveError(&handler, error.OutOfMemory)); + try testing.expectEqual(@as(usize, 0), handler.errors); + var no_callback: struct {} = .{}; + try testing.expectError(error.FetchFailed, reportArchiveError(&no_callback, error.FetchFailed)); + const NotifyOnly = struct { + errors: usize = 0, + pub fn onError(self: *@This(), _: anyerror) void { + self.errors += 1; + } + }; + var notify_only = NotifyOnly{}; + try testing.expectError(error.FetchFailed, reportArchiveError(¬ify_only, error.FetchFailed)); + try testing.expectEqual(@as(usize, 1), notify_only.errors); +} diff --git a/src/client.zig b/src/client.zig index d8c5dea..8ca9bf4 100644 --- a/src/client.zig +++ b/src/client.zig @@ -97,8 +97,8 @@ pub const Options = struct { kinds: []const livedecode.Kind = &.{}, collections: []const []const u8 = &.{}, dids: []const []const u8 = &.{}, - /// bearer credential for token-gated archives and, matching the - /// upstream Go client, sent on the live websocket handshake too + /// bearer credential for token-gated archives and, as a local extension, + /// sent on the live websocket handshake too /// (hosted instances can gate the tail at their edge; the OSS server /// does not) api_key: ?[]const u8 = null, @@ -136,7 +136,8 @@ const HostOutcome = enum { stopped, rotate }; /// run the unified stream through `handler`: /// fn onEvent(*H, livedecode.Event) bool — every event, archive and live, /// in seq order across the seam; false stops cleanly -/// optional fn onError(*H, anyerror) bool — recoverable live errors +/// optional fn onError(*H, anyerror) bool — recoverable archive/live errors; +/// true continues, false stops cleanly. Without it, archive errors return. /// optional fn onInfo(*H, livedecode.Info) void /// blocks until the handler stops it, the snapshot completes /// (snapshot_only), or a terminal error. @@ -209,10 +210,9 @@ pub fn subscribe(io: Io, allocator: Allocator, options: Options, handler: anytyp /// the pre-failover single-host flow, verbatim, parameterized by host and /// anchor. can_rotate=false propagates every error exactly as the /// original did; can_rotate=true converts live-tail host-health errors -/// (transient threshold, stall) into .rotate. archive-sweep errors are -/// always terminal: archives are not interchangeable across hosts -/// (coverage, retention, holes differ), so a mid-sweep host switch could -/// silently anchor past a gap — the caller decides, not the client. +/// (transient threshold, stall) into .rotate. Archive errors never request +/// rotation: accepted entry/decode errors continue on the same host, while +/// plan and resource failures propagate. Archives are not interchangeable. fn runOnHost( io: Io, allocator: Allocator, @@ -334,6 +334,11 @@ fn SweepHandler(comptime H: type) type { time_floor_us: i64 = 0, progress: *Progress, + pub fn onError(self: *@This(), err: anyerror) !bool { + if (comptime @hasDecl(H, "onError")) return self.inner.onError(err); + return err; + } + pub fn onRow(self: *@This(), event: livedecode.Event) bool { if (self.time_floor_us > 0 and event.time_us < self.time_floor_us) return true; if (!wantsKind(self.kinds, event)) return true; diff --git a/src/filter.zig b/src/filter.zig new file mode 100644 index 0000000..14abe77 --- /dev/null +++ b/src/filter.zig @@ -0,0 +1,28 @@ +const std = @import("std"); +const livedecode = @import("livedecode.zig"); + +pub fn collectionMatches(collections: []const []const u8, collection: []const u8) bool { + if (collections.len == 0 or collection.len == 0) return true; + for (collections) |pattern| { + if (std.mem.eql(u8, pattern, collection)) return true; + // Keep the dot: app.bsky.graph.* must not match app.bsky.graphient. + if (std.mem.endsWith(u8, pattern, ".*") and + std.mem.startsWith(u8, collection, pattern[0 .. pattern.len - 1])) return true; + } + return false; +} + +pub fn wants(event: livedecode.Event, kinds: []const livedecode.Kind, dids: []const []const u8, collections: []const []const u8) bool { + if (kinds.len > 0) { + for (kinds) |kind| { + if (kind == std.meta.activeTag(event.payload)) break; + } else return false; + } + if (dids.len > 0) { + for (dids) |did| { + if (std.mem.eql(u8, did, event.did)) break; + } else return false; + } + // DID-level markers remain visible under collection filters. + return event.payload != .commit or collectionMatches(collections, event.payload.commit.collection); +} diff --git a/src/live.zig b/src/live.zig index 43691b8..7b2402b 100644 --- a/src/live.zig +++ b/src/live.zig @@ -27,6 +27,7 @@ const std = @import("std"); const zat = @import("zat"); +const filter = @import("filter.zig"); const websocket = @import("websocket"); const livedecode = @import("livedecode.zig"); const zstd = @import("zstd.zig"); @@ -472,6 +473,7 @@ pub const LiveClient = struct { if (event.seq <= client.last_seq) return; client.last_seq = event.seq; client.seen_any = true; + if (!filter.wants(event, client.options.kinds, client.options.dids, client.options.collections)) return; if (!self.handler.onEvent(event)) { self.stopped = true; return error.Stop; diff --git a/src/loopback_test.zig b/src/loopback_test.zig index 5214a83..ccfe88d 100644 --- a/src/loopback_test.zig +++ b/src/loopback_test.zig @@ -28,6 +28,8 @@ const Row = struct { seq: u64, /// seconds past epoch_s (0..59: rendered into one RFC3339 minute) t: i64, + did: []const u8 = "did:plc:loopback", + collection: []const u8 = "app.bsky.feed.post", }; const Instance = struct { @@ -45,6 +47,14 @@ const Instance = struct { /// immediately (the host is down) max_ws_sessions: u32 = std.math.maxInt(u32), ws_sessions: u32 = 0, + page_size: usize = 0, + page_rows: ?[]const Row = null, + plan_calls: usize = 0, + last_before_seq: ?u64 = null, + first_tip: ?u64 = null, + stalled_plan: bool = false, + archive_fault: enum { none, fetch, decode, segment } = .none, + reject_plan: bool = false, fn init(io: Io, allocator: Allocator, sealed: []const Row, live: []const Row) !Instance { var listener = try (Io.net.IpAddress{ .ip4 = .loopback(0) }).listen(io, .{ .reuse_address = true }); @@ -85,9 +95,11 @@ const Instance = struct { fn quit(self: *Instance) void { var transport = zat.HttpTransport.init(self.io, self.allocator); defer transport.deinit(); + // A shutdown request must not retry after the server has exited. + transport.dead_connection_attempts = 1; var url_buf: [64]u8 = undefined; const url = std.fmt.bufPrint(&url_buf, "http://127.0.0.1:{d}/quit", .{self.port}) catch unreachable; - var response = transport.fetch(.{ .url = url }) catch return; + var response = transport.fetch(.{ .url = url, .deadline_ns = std.time.ns_per_s }) catch return; response.deinit(self.allocator); } @@ -116,7 +128,10 @@ const Instance = struct { defer arena_state.deinit(); const arena = arena_state.allocator(); - if (std.mem.startsWith(u8, target, "/quit")) return true; + if (std.mem.startsWith(u8, target, "/quit")) { + try respondBytes(w, "text/plain", ""); + return true; + } if (std.mem.indexOf(u8, target, "subscribeEvents") != null) { try self.handleWebsocket(request, target, r, w); return false; @@ -126,12 +141,51 @@ const Instance = struct { return false; } if (std.mem.eql(u8, method, "POST") and std.mem.indexOf(u8, target, "planSnapshot") != null) { - // drain the body so the client's write completes cleanly - if (contentLength(request)) |len| _ = try r.discardAll(len); + const body = try arena.alloc(u8, contentLength(request) orelse 0); + try r.readSliceAll(body); + self.plan_calls += 1; + if (self.reject_plan) { + try w.writeAll("HTTP/1.1 400 Bad Request\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"); + try w.flush(); + return false; + } + if (self.page_size > 0) { + const parsed = try std.json.parseFromSlice(std.json.Value, arena, body, .{}); + const fields = parsed.value.object; + const after: u64 = if (fields.get("afterSeq")) |v| @intCast(v.integer) else 0; + self.last_before_seq = if (fields.get("beforeSeq")) |v| @intCast(v.integer) else null; + var start: usize = 0; + while (start < self.sealed.len and self.sealed[start].seq <= after) : (start += 1) {} + var end = @min(start + self.page_size, self.sealed.len); + if (self.last_before_seq) |before| { + while (end > start and self.sealed[end - 1].seq > before) : (end -= 1) {} + } + self.page_rows = self.sealed[start..end]; + } try respondJson(w, try self.planJson(arena)); return false; } + if (std.mem.indexOf(u8, target, "getSegment") != null and self.archive_fault == .segment) { + try respondBytes(w, "application/octet-stream", "invalid segment header"); + return false; + } if (std.mem.indexOf(u8, target, "getBlock") != null) { + if (self.archive_fault != .none) { + const second = std.mem.indexOf(u8, target, "segment=second") != null; + const block: u64 = if (std.mem.indexOf(u8, target, "blockIndex=1") != null) 1 else if (std.mem.indexOf(u8, target, "blockIndex=2") != null) 2 else 0; + if (!second and block == 1) { + if (self.archive_fault == .fetch) { + try w.writeAll("HTTP/1.1 404 Not Found\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"); + try w.flush(); + } else try respondBytes(w, "application/octet-stream", "invalid zstd frame"); + return false; + } + var fixture = self.*; + fixture.sealed = &.{.{ .seq = if (second) 100 else block + 1, .t = if (second) 10 else @intCast(block + 1) }}; + fixture.page_rows = null; + try respondBytes(w, "application/octet-stream", try fixture.blockZstd(arena)); + return false; + } try respondBytes(w, "application/octet-stream", try self.blockZstd(arena)); return false; } @@ -153,27 +207,41 @@ const Instance = struct { } fn planJson(self: *const Instance, arena: Allocator) ![]const u8 { + if (self.archive_fault == .segment) return + \\{"plannedThroughSeq":100,"sealedTipSeq":100,"segments":[ + \\{"name":"first","index":0,"mode":"segment"}, + \\{"name":"second","index":1,"mode":"blocks","blocks":[{"first":0,"last":0}]}]} + ; + if (self.archive_fault != .none) return + \\{"plannedThroughSeq":100,"sealedTipSeq":100,"segments":[ + \\{"name":"first","index":0,"mode":"blocks","blocks":[{"first":0,"last":2}]}, + \\{"name":"second","index":1,"mode":"blocks","blocks":[{"first":0,"last":0}]}]} + ; if (self.sealed.len == 0) return "{\"plannedThroughSeq\": 0, \"sealedTipSeq\": 0, \"segments\": []}"; - const tip = self.sealed[self.sealed.len - 1].seq; + const tip = self.last_before_seq orelse (if (self.plan_calls == 1) self.first_tip else null) orelse self.sealed[self.sealed.len - 1].seq; + if (self.stalled_plan) return std.fmt.allocPrint(arena, "{{\"plannedThroughSeq\":0,\"sealedTipSeq\":{d},\"segments\":[]}}", .{tip}); + const rows = self.page_rows orelse self.sealed; + if (rows.len == 0) return std.fmt.allocPrint(arena, "{{\"plannedThroughSeq\":{d},\"sealedTipSeq\":{d},\"segments\":[]}}", .{ tip, tip }); return std.fmt.allocPrint(arena, \\{{"plannedThroughSeq": {d}, "sealedTipSeq": {d}, "segments": [ \\ {{"name": "seg_0000000001", "index": 0, "checksum": "0", "minSeq": {d}, "maxSeq": {d}, \\ "mode": "blocks", "blocks": [{{"first": 0, "last": 0}}]}}]}} - , .{ tip, tip, self.sealed[0].seq, tip }); + , .{ rows[rows.len - 1].seq, tip, rows[0].seq, rows[rows.len - 1].seq }); } /// the sealed rows as one columnar jss block wrapped in a raw-block /// zstd frame (no compressor needed: magic, single-segment FHD with /// 4-byte FCS, one raw block) fn blockZstd(self: *const Instance, arena: Allocator) ![]u8 { - const rows = try arena.alloc(BlockRow, self.sealed.len); - for (self.sealed, rows) |row, *out| out.* = .{ + const source_rows = self.page_rows orelse self.sealed; + const rows = try arena.alloc(BlockRow, source_rows.len); + for (source_rows, rows) |row, *out| out.* = .{ .seq = row.seq, .witnessed_at = witnessedUs(row.t), .kind = 1, // create - .collection = "app.bsky.feed.post", - .did = "did:plc:loopback", + .collection = row.collection, + .did = row.did, .rkey = "rkey1", .rev = "rev1", .payload = &archive.test_record_cbor, @@ -213,8 +281,8 @@ const Instance = struct { if (cursor) |c| if (row.seq <= c) continue; var frame_buf: [512]u8 = undefined; const frame = try std.fmt.bufPrint(&frame_buf, - \\{{"$type":"message","payload":{{"$type":"network.bsky.jetstream.subscribeEvents#commit","collection":"app.bsky.feed.post","did":"did:plc:loopback","operation":"create","record":{{"$type":"app.bsky.feed.post","text":"hi"}},"rev":"rev1","rkey":"rkey1","seq":{d},"time":"2026-07-13T00:00:{d:0>2}.000000Z"}}}} - , .{ row.seq, @as(u64, @intCast(row.t)) }); + \\{{"$type":"message","payload":{{"$type":"network.bsky.jetstream.subscribeEvents#commit","collection":"{s}","did":"{s}","operation":"create","record":{{"$type":"app.bsky.feed.post","text":"hi"}},"rev":"rev1","rkey":"rkey1","seq":{d},"time":"2026-07-13T00:00:{d:0>2}.000000Z"}}}} + , .{ row.collection, row.did, row.seq, @as(u64, @intCast(row.t)) }); try writeTextFrame(w, frame); } if (self.live_then_close) return; @@ -223,7 +291,8 @@ const Instance = struct { // rotation, reconnect) ends this read var discard: [256]u8 = undefined; while (true) { - _ = r.readSliceShort(&discard) catch return; + const read = r.readSliceShort(&discard) catch return; + if (read == 0) return; } } }; @@ -314,8 +383,7 @@ test "e2e failover: primary dies mid-live; time re-anchor on the fallback resume // and 1003 (t=5) deliver, then live 1004 (t=6) var fallback = try Instance.init(io, testing.allocator, &.{ .{ .seq = 1001, .t = 2 }, .{ .seq = 1002, .t = 3 }, .{ .seq = 1003, .t = 5 } }, &.{.{ .seq = 1004, .t = 6 }}); - // defer order matters (LIFO): the listener closes FIRST, unblocking - // the accept loop, THEN the future joins + // LIFO: request shutdown, join the serve task, then close the listener. var primary_future = try io.concurrent(Instance.serve, .{&primary}); defer primary.deinit(); defer _ = primary_future.cancel(io) catch {}; @@ -404,3 +472,234 @@ test "e2e live-phase clean stop: onEvent=false ends subscribe instead of reconne try testing.expectEqual(@as(usize, 1), handler.len); try testing.expectEqual(@as(u64, 1), handler.slice()[0].seq); } + +// Upstream engine_test.go: TestEngineMultiPageBackfillCutover and +// TestEnginePinnedBeforeSeqAcrossPages (289b032 and 58c4d7f). +test "conformance: snapshot consumes every plan page and pins the sealed tip" { + var threaded: Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + var host = try Instance.init(io, testing.allocator, &.{ + .{ .seq = 1, .t = 1 }, .{ .seq = 2, .t = 2 }, .{ .seq = 3, .t = 3 }, .{ .seq = 4, .t = 4 }, + }, &.{}); + host.page_size = 1; + host.first_tip = 3; + var future = try io.concurrent(Instance.serve, .{&host}); + defer host.deinit(); + defer _ = future.cancel(io) catch {}; + defer host.quit(); + var url: [64]u8 = undefined; + var handler = CollectingHandler(32){}; + try client_mod.subscribe(io, testing.allocator, .{ + .hosts = &.{host.hostUrl(&url)}, + .after_seq = 0, + .snapshot_only = true, + .concurrency = 1, + }, &handler); + try testing.expectEqualSlices(Collected, &.{ + .{ .seq = 1, .t = 1 }, .{ .seq = 2, .t = 2 }, .{ .seq = 3, .t = 3 }, + }, handler.slice()); + try testing.expectEqual(@as(usize, 3), host.plan_calls); + try testing.expectEqual(@as(?u64, 3), host.last_before_seq); +} + +// Upstream TestEngineLiveOnlyAppliesCollectionFilter and +// TestMatcherDIDFilterAllKinds: the server deliberately ignores filters. +test "conformance: live delivery reapplies collection and DID filters" { + var threaded: Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + var host = try Instance.init(io, testing.allocator, &.{}, &.{ + .{ .seq = 1, .t = 1, .collection = "app.bsky.feed.like" }, + .{ .seq = 2, .t = 2, .did = "did:plc:other" }, + .{ .seq = 3, .t = 3 }, + }); + var future = try io.concurrent(Instance.serve, .{&host}); + defer host.deinit(); + defer _ = future.cancel(io) catch {}; + defer host.quit(); + const Handler = struct { + collected: CollectingHandler(32) = .{}, + pub fn onEvent(self: *@This(), event: livedecode.Event) bool { + _ = self.collected.onEvent(event); + return event.seq < 3; + } + }; + var handler = Handler{}; + var url: [64]u8 = undefined; + try client_mod.subscribe(io, testing.allocator, .{ + .hosts = &.{host.hostUrl(&url)}, + .collections = &.{"app.bsky.feed.post"}, + .dids = &.{"did:plc:loopback"}, + }, &handler); + try testing.expectEqualSlices(Collected, &.{.{ .seq = 3, .t = 3 }}, handler.collected.slice()); +} + +// Upstream sweepSealedArchive rejects a non-advancing continuation. +test "conformance: a stalled plan fails instead of looping" { + var threaded: Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + var host = try Instance.init(io, testing.allocator, &.{.{ .seq = 3, .t = 3 }}, &.{}); + host.stalled_plan = true; + var future = try io.concurrent(Instance.serve, .{&host}); + defer host.deinit(); + defer _ = future.cancel(io) catch {}; + defer host.quit(); + var url: [64]u8 = undefined; + var handler = CollectingHandler(32){}; + try testing.expectError(error.PlanStalled, client_mod.subscribe(io, testing.allocator, .{ + .hosts = &.{host.hostUrl(&url)}, + .after_seq = 0, + .snapshot_only = true, + }, &handler)); + try testing.expectEqual(@as(usize, 1), host.plan_calls); +} + +test "archive block budget applies across plan pages" { + var threaded: Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + var host = try Instance.init(io, testing.allocator, &.{ + .{ .seq = 1, .t = 1 }, .{ .seq = 2, .t = 2 }, .{ .seq = 3, .t = 3 }, + }, &.{}); + host.page_size = 1; + var future = try io.concurrent(Instance.serve, .{&host}); + defer host.deinit(); + defer _ = future.cancel(io) catch {}; + defer host.quit(); + const Handler = struct { + count: usize = 0, + pub fn onRow(self: *@This(), _: livedecode.Event) bool { + self.count += 1; + return true; + } + }; + var handler = Handler{}; + var url: [64]u8 = undefined; + const result = try archive.run(io, testing.allocator, .{ + .host = host.hostUrl(&url), + .after_seq = 0, + .max_blocks = 2, + .concurrency = 1, + }, &handler); + try testing.expect(result.truncated); + try testing.expectEqual(@as(u64, 2), result.blocks_decoded); + try testing.expectEqual(@as(usize, 2), handler.count); + try testing.expectEqual(@as(usize, 2), host.plan_calls); + try testing.expectEqual(@as(u64, 3), result.sealed_tip_seq); +} + +// Upstream TestDownloadBlocksErrorStopsEntryNotPlan: good prefix, error, +// then the next entry; a decode error instead leaves later blocks readable. +test "archive errors are ordered and the caller chooses whether to continue" { + for ([_]bool{ true, false }) |accept_error| { + for ([_]usize{ 1, 4 }) |concurrency| { + for ([_]bool{ false, true }) |decode_failure| { + var threaded: Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + var host = try Instance.init(io, testing.allocator, &.{}, &.{}); + host.archive_fault = if (decode_failure) .decode else .fetch; + var future = try io.concurrent(Instance.serve, .{&host}); + defer host.deinit(); + defer _ = future.cancel(io) catch {}; + defer host.quit(); + const Handler = struct { + items: [8]u64 = undefined, + len: usize = 0, + accept: bool, + err: ?anyerror = null, + pub fn onEvent(self: *@This(), event: livedecode.Event) bool { + self.items[self.len] = event.seq; + self.len += 1; + return true; + } + pub fn onError(self: *@This(), err: anyerror) bool { + self.err = err; + self.items[self.len] = 0; + self.len += 1; + return self.accept; + } + }; + var handler = Handler{ .accept = accept_error }; + var url: [64]u8 = undefined; + // If an archive failure incorrectly rotates, port 1 cannot + // supply the next entry and the expected sequence is lost. + try client_mod.subscribe(io, testing.allocator, .{ + .hosts = &.{ host.hostUrl(&url), "http://127.0.0.1:1" }, + .after_seq = 0, + .snapshot_only = true, + .concurrency = concurrency, + }, &handler); + try testing.expectEqual(@as(?anyerror, if (decode_failure) error.BlockDecompressFailed else error.FetchFailed), handler.err); + const expected: []const u64 = if (!accept_error) &.{ 1, 0 } else if (decode_failure) &.{ 1, 0, 3, 100 } else &.{ 1, 0, 100 }; + try testing.expectEqualSlices(u64, expected, handler.items[0..handler.len]); + } + } + } +} + +test "archive plan failures stay terminal even when the caller accepts errors" { + var threaded: Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + var host = try Instance.init(io, testing.allocator, &.{}, &.{}); + host.reject_plan = true; + var future = try io.concurrent(Instance.serve, .{&host}); + defer host.deinit(); + defer _ = future.cancel(io) catch {}; + defer host.quit(); + const Handler = struct { + errors: usize = 0, + pub fn onEvent(_: *@This(), _: livedecode.Event) bool { + return true; + } + pub fn onError(self: *@This(), _: anyerror) bool { + self.errors += 1; + return true; + } + }; + var handler = Handler{}; + var url: [64]u8 = undefined; + try testing.expectError(error.PlanFailed, client_mod.subscribe(io, testing.allocator, .{ + .hosts = &.{host.hostUrl(&url)}, + .after_seq = 0, + .snapshot_only = true, + }, &handler)); + try testing.expectEqual(@as(usize, 0), handler.errors); +} + +test "a whole segment error can be accepted or declined before the next entry" { + for ([_]bool{ true, false }) |accept| { + var threaded: Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + var host = try Instance.init(io, testing.allocator, &.{}, &.{}); + host.archive_fault = .segment; + var future = try io.concurrent(Instance.serve, .{&host}); + defer host.deinit(); + defer _ = future.cancel(io) catch {}; + defer host.quit(); + const Handler = struct { + accept: bool, + errors: usize = 0, + last: u64 = 0, + pub fn onRow(self: *@This(), event: livedecode.Event) bool { + self.last = event.seq; + return true; + } + pub fn onError(self: *@This(), err: anyerror) bool { + if (err != error.MalformedSegment) @panic("unexpected error"); + self.errors += 1; + return self.accept; + } + }; + var handler = Handler{ .accept = accept }; + var url: [64]u8 = undefined; + const result = try archive.run(io, testing.allocator, .{ .host = host.hostUrl(&url) }, &handler); + try testing.expectEqual(@as(usize, 1), handler.errors); + try testing.expectEqual(@as(u64, if (accept) 100 else 0), handler.last); + try testing.expectEqual(!accept, result.stopped); + } +} -- 2.51.2