From 9d807d48bedd41aa4108352516fb8c99efaf85f9 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Sun, 6 Sep 2026 16:35:38 -0500 Subject: [PATCH] release: v0.1.2 --- CHANGELOG.md | 17 ++++++ README.md | 2 +- build.zig.zon | 6 +- docs/client.md | 11 ++++ docs/upstream-review-2026-09-06.md | 35 +++++++++-- src/archive_backfill.zig | 94 +++++++++++++++++++++++++----- src/client.zig | 6 +- src/live.zig | 23 ++------ src/loopback_test.zig | 59 +++++++++++++++++++ 9 files changed, 204 insertions(+), 49 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index ec4067b..a63dbde 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,22 @@ # changelog +## 0.1.2 + +- adopt ZAT 0.5.1 and ZTTP 0.1.1, including legacy HTTP-date retry hints. +- finish snapshot pagination with a pinned sealed tip and a cumulative block + budget; reject stalled plans. Match collection wildcards and enforce live filters. +- report recoverable archive download/decode errors to the caller; continue + on the same archive only when accepted. Plan errors remain terminal. +- align download retry attempts, backoff, jitter, and server-hint handling + with upstream's separate block and whole-segment policies. +- authenticate archives only. Remove `live.Options.api_key`; callers setting + it must remove it. The unified client's archive `api_key` remains supported. +- report archive bytes fetched and allow archive reads without record decoding. +- add local segment-to-TSV export and network-wide collection slice examples; + authenticate Streamplace archive replay with an environment-provided key. +- fix loopback shutdown and EOF handling so tests complete promptly; add + upstream conformance regressions and a release workflow skill. + ## 0.1.1 - the plan request now carries `kinds` (upstream WithKinds parity): the diff --git a/README.md b/README.md index 4ec7cf5..d290faf 100644 --- a/README.md +++ b/README.md @@ -11,7 +11,7 @@ draws at the protocol: zat speaks atproto; this SDK speaks ## use ```sh -zig fetch --save https://tangled.org/zat.dev/jetstream/archive/v0.1.0.tar.gz +zig fetch --save https://tangled.org/zat.dev/jetstream/archive/v0.1.2.tar.gz ``` ```zig diff --git a/build.zig.zon b/build.zig.zon index c77fafe..e1349e7 100644 --- a/build.zig.zon +++ b/build.zig.zon @@ -1,6 +1,6 @@ .{ .name = .jetstream, - .version = "0.1.1", + .version = "0.1.2", .minimum_zig_version = "0.16.0-dev.3070+b22eb176b", .fingerprint = 0xee435acce6d720e8, .dependencies = .{ @@ -9,8 +9,8 @@ .hash = "N-V-__8AAPZ7fwBg4JoCzM_0o2A8wxH2hsUUeiU1iuZv53L5", }, .zat = .{ - .url = "https://tangled.org/zat.dev/zat/archive/v0.5.0.tar.gz", - .hash = "zat-0.5.0-5PuC7sUhDAAgtrkdHxFwI0vq4Qfy7PY4ZYYATu_pOmS3", + .url = "https://tangled.org/zat.dev/zat/archive/v0.5.1.tar.gz", + .hash = "zat-0.5.1-5PuC7g0iDAAsD79O5oQaB1q1M3Dq7EOONdlfvIE0sSis", }, .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 65646a7..862f93a 100644 --- a/docs/client.md +++ b/docs/client.md @@ -82,3 +82,14 @@ data frame (one bad frame never drops the tail). the v1 `/subscribe` wire client (time_us cursors, host round-robin) remains `zat.JetstreamClient` in [zat](https://tangled.org/zat.dev/zat). + +Authentication: the unified client’s `api_key` authenticates archive requests +only. Live WebSocket connections are public and never send that credential. +Standalone live options have no `api_key` field. + +Archive GET retries make up to three attempts. The fallback starts at 500ms +and doubles, capped at 30s. Block/list requests add 0–20% jitter and honor a +positive server delay in place of the fallback; whole-segment requests use +no jitter and the longer of the fallback and server delay. Absolute +`RateLimit-Reset` takes precedence over `Retry-After` (seconds or IMF-fixdate). +Past hints use the fallback, and all server-directed waits are capped at 30s. diff --git a/docs/upstream-review-2026-09-06.md b/docs/upstream-review-2026-09-06.md index 9d83d2d..f54205d 100644 --- a/docs/upstream-review-2026-09-06.md +++ b/docs/upstream-review-2026-09-06.md @@ -52,14 +52,18 @@ source comparison, not an exhaustive conformance certification. 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. + WebSocket. Fixed: the unified client no longer forwards its archive key + to live connections, and live handshakes never send Authorization. The + standalone live `api_key` field has been removed. A loopback regression failed before + the fix and now verifies authenticated archive requests followed by an + unauthenticated live handshake (51/51 tests pass). - **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. +- **Retry timing:** updated to upstream’s three-attempt default, 500ms + exponential base, and 30s cap. XRPC block/list requests use positive + 0–20% jitter and let a positive server hint replace the schedule; bulk + segments have no jitter and use the longer of hint and schedule. + Connection-level retries are disabled within the download loop. ## Assessment of the unreleased additions @@ -247,3 +251,22 @@ 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. + +## Retry policy verification + +Compared `segmentfetch.go` at upstream `58c4d7f` and its pinned +`github.com/jcalabro/atmos v0.4.0` (`xrpc/retry.go`, `client.go`, +`error.go`). The block path retries 429/500/502/503/504; bulk segments +retry 429 and all 5xx. Cancellation and allocation failure propagate. +Timing and attempt-count regressions failed before the change. Pure +calculations cover jitter, caps, header precedence, expired hints and +fractional waits; a loopback server verifies exactly three HTTP attempts. +Zig's `std.Random.float` supplies the random fraction, backed by `Io.random`; +no separate math library or dependency change is needed for jitter. + +The client now pins ZAT v0.5.1, which pins HTTP library v0.1.1. That +release adds RFC 850 and asctime dates and validates calendar dates; +33/33 HTTP-library tests passed (two failed before the fix). The client +was also verified against the local dependency chain before release. +Unsigned reset parsing remains a separate edge-case difference from Go +for negative reset values; byte-for-byte malformed-header parity is not claimed. diff --git a/src/archive_backfill.zig b/src/archive_backfill.zig index b166496..f24fbd7 100644 --- a/src/archive_backfill.zig +++ b/src/archive_backfill.zig @@ -231,7 +231,7 @@ fn runPage(io: Io, allocator: Allocator, options: Options, handler: anytype) !Re .{ options.host, segment.name }, ); defer allocator.free(url); - var fetched = fetchWithRetry(io, &transport, url, authorization) catch |err| { + var fetched = fetchWithRetry(io, &transport, url, authorization, .segment) catch |err| { reportArchiveError(handler, err) catch |reported| { if (reported != error.Stopped) return reported; result.stopped = true; @@ -284,7 +284,7 @@ fn reportArchiveError(handler: anytype, err: anyerror) !void { } fn fetchJob(io: Io, transport: *HttpTransport, url: []const u8, authorization: ?[]const u8) FetchError!HttpTransport.FetchResult { - return fetchWithRetry(io, transport, url, authorization); + return fetchWithRetry(io, transport, url, authorization, .xrpc); } /// a FIFO window of in-flight block fetches. starting is out-of-order I/O; @@ -419,35 +419,66 @@ fn skippedByResume(start_after: ?Position, segment_index: u32, block_index: u32) return segment_index == sa.segment_index and block_index <= sa.block_index; } -const fetch_attempts = 4; +const RetryKind = enum { xrpc, segment }; +fn retryDelay(kind: RetryKind, attempt: u32, hint_ns: u64, jitter: f64) u64 { + const cap = 30 * std.time.ns_per_s; + const base = @as(u64, 500 * std.time.ns_per_ms) << @as(u6, @intCast(@min(attempt - 1, 6))); + const scheduled: u64 = @min(cap, if (kind == .xrpc) + @as(u64, @intFromFloat(@as(f64, @floatFromInt(base)) * (1 + 0.2 * jitter))) + else + base); + if (kind == .xrpc and hint_ns > 0) return @min(cap, hint_ns); + return @min(cap, @max(scheduled, hint_ns)); +} + +fn fillRandom(io: *Io, bytes: []u8) void { + io.random(bytes); +} + +const fetch_attempts = 3; -/// GET with bounded retries: transport errors and retryable statuses -/// (5xx, 429) back off 1s/2s/4s; other non-200s fail immediately. a -/// multi-GB run should not be discarded over one transient failure. -fn fetchWithRetry(io: Io, transport: *HttpTransport, url: []const u8, authorization: ?[]const u8) !HttpTransport.FetchResult { +/// Match the Go bulk-segment and XRPC retry policies separately. +fn fetchWithRetry(io: Io, transport: *HttpTransport, url: []const u8, authorization: ?[]const u8, kind: RetryKind) !HttpTransport.FetchResult { + const previous_attempts = transport.dead_connection_attempts; + transport.dead_connection_attempts = 1; + defer transport.dead_connection_attempts = previous_attempts; + var random_io = io; + const random = std.Random.init(&random_io, fillRandom); var attempt: u32 = 0; - var backoff_s: u64 = 1; while (true) { attempt += 1; + var hint_ns: u64 = 0; if (transport.fetch(.{ .url = url, .authorization = authorization })) |result| { if (result.status == .ok) return result; var discarded = result; - discarded.deinit(transport.allocator); const code = @intFromEnum(result.status); - const retryable = code >= 500 or result.status == .too_many_requests; + hint_ns = retryHint(result.rate_limit, Io.Clock.real.now(io).toNanoseconds()); + discarded.deinit(transport.allocator); + const retryable = if (kind == .segment) code >= 500 or code == 429 else code == 429 or code == 500 or code == 502 or code == 503 or code == 504; if (!retryable) return error.FetchFailed; if (attempt >= fetch_attempts) return error.FetchRetriesExhausted; - log.warn("fetch {d} (attempt {d}/{d}), retrying in {d}s: {s}", .{ code, attempt, fetch_attempts, backoff_s, url }); } else |err| { - if (err == error.Canceled) return err; + if (err == error.Canceled or err == error.OutOfMemory) return err; if (attempt >= fetch_attempts) return err; - log.warn("fetch error {s} (attempt {d}/{d}), retrying in {d}s: {s}", .{ @errorName(err), attempt, fetch_attempts, backoff_s, url }); } - try io.sleep(Io.Duration.fromSeconds(@intCast(backoff_s)), .awake); - backoff_s *= 2; + const jitter = if (kind == .xrpc) random.float(f64) else 0; + const delay = retryDelay(kind, attempt, hint_ns, jitter); + try io.sleep(Io.Duration.fromNanoseconds(@intCast(delay)), .awake); } } +fn retryHint(headers: HttpTransport.RateLimitHeaders, now_ns: i96) u64 { + const nanoseconds: i128 = if (headers.reset) |reset| + @as(i128, reset) * std.time.ns_per_s - now_ns + else if (headers.retry_after) |delta| + @as(i128, delta) * std.time.ns_per_s + else if (headers.retry_after_at) |date| + @as(i128, date) * std.time.ns_per_s - now_ns + else + 0; + return @intCast(@min(@max(nanoseconds, 0), 30 * std.time.ns_per_s)); +} + fn limitReached(options: Options, result: *Result) bool { const max = options.max_blocks orelse return false; if (result.blocks_decoded < max) return false; @@ -615,7 +646,7 @@ pub fn fetchSegmentSpans(io: Io, allocator: Allocator, arena: Allocator, host: [ try std.fmt.allocPrint(arena, "{s}/xrpc/network.bsky.jetstream.listSegments?limit=1000&cursor={s}", .{ host, c }) else try std.fmt.allocPrint(arena, "{s}/xrpc/network.bsky.jetstream.listSegments?limit=1000", .{host}); - var fetched = try fetchWithRetry(io, &transport, url, authorization); + var fetched = try fetchWithRetry(io, &transport, url, authorization, .xrpc); defer fetched.deinit(allocator); const parsed = try json.parseFromSliceLeaky(json.Value, arena, fetched.body, .{}); const root = switch (parsed) { @@ -1603,3 +1634,34 @@ test "archive recovery never swallows cancellation allocation failure or absent try testing.expectError(error.FetchFailed, reportArchiveError(¬ify_only, error.FetchFailed)); try testing.expectEqual(@as(usize, 1), notify_only.errors); } + +test "download retries use upstream's three attempt default" { + try std.testing.expectEqual(@as(u32, 3), fetch_attempts); +} + +test "retry timing matches upstream segment and XRPC policies" { + const ms = std.time.ns_per_ms; + const sec = std.time.ns_per_s; + try std.testing.expectEqual(@as(u64, 500 * ms), retryDelay(.segment, 1, 0, 0)); + try std.testing.expectEqual(@as(u64, sec), retryDelay(.segment, 2, 0, 0)); + try std.testing.expectEqual(@as(u64, 30 * sec), retryDelay(.segment, 100, 0, 0)); + try std.testing.expectEqual(@as(u64, 5 * sec), retryDelay(.segment, 1, 5 * sec, 0)); + try std.testing.expectEqual(@as(u64, sec), retryDelay(.segment, 2, 100 * ms, 0)); + try std.testing.expectEqual(@as(u64, 500 * ms), retryDelay(.xrpc, 1, 0, 0)); + try std.testing.expect(retryDelay(.xrpc, 1, 0, 0.999999) < 600 * ms); + try std.testing.expectEqual(@as(u64, 550 * ms), retryDelay(.xrpc, 1, 0, 0.5)); + try std.testing.expectEqual(@as(u64, 100 * ms), retryDelay(.xrpc, 2, 100 * ms, 0)); + try std.testing.expectEqual(@as(u64, 30 * sec), retryDelay(.xrpc, 1, 60 * sec, 0)); +} + +test "retry hints prefer absolute reset and bound missing expired and large values" { + const sec = std.time.ns_per_s; + try std.testing.expectEqual(@as(u64, 0), retryHint(.{}, 100 * sec)); + try std.testing.expectEqual(@as(u64, 250 * std.time.ns_per_ms), retryHint(.{ .reset = 101 }, 100 * sec + 750 * std.time.ns_per_ms)); + try std.testing.expectEqual(@as(u64, 5 * sec), retryHint(.{ .retry_after = 5 }, 100 * sec)); + try std.testing.expectEqual(@as(u64, 3 * sec), retryHint(.{ .reset = 103, .retry_after = 9 }, 100 * sec)); + try std.testing.expectEqual(@as(u64, 0), retryHint(.{ .reset = 99, .retry_after = 9 }, 100 * sec)); + try std.testing.expectEqual(@as(u64, 2 * sec), retryHint(.{ .retry_after_at = 102 }, 100 * sec)); + try std.testing.expectEqual(@as(u64, 0), retryHint(.{ .retry_after_at = 90 }, 100 * sec)); + try std.testing.expectEqual(@as(u64, 30 * sec), retryHint(.{ .retry_after = std.math.maxInt(u64) }, 100 * sec)); +} diff --git a/src/client.zig b/src/client.zig index 8ca9bf4..afcb3c4 100644 --- a/src/client.zig +++ b/src/client.zig @@ -97,10 +97,7 @@ pub const Options = struct { kinds: []const livedecode.Kind = &.{}, collections: []const []const u8 = &.{}, dids: []const []const u8 = &.{}, - /// 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) + /// Bearer credential for archive requests only. Live connections are public. api_key: ?[]const u8 = null, /// block fetches in flight during the archive sweep (see ArchiveBackfill) concurrency: usize = 4, @@ -315,7 +312,6 @@ fn liveOptions(options: Options, host: []const u8, cursor: ?u64, dedup_floor: u6 .collections = options.collections, .dids = options.dids, .zstd_dict = dict, - .api_key = options.api_key, .stall_timeout_ns = options.failover_stall_ns, }; if (options.backoff_min_ns > 0) out.backoff_min_ns = options.backoff_min_ns; diff --git a/src/live.zig b/src/live.zig index 7b2402b..30de3db 100644 --- a/src/live.zig +++ b/src/live.zig @@ -74,12 +74,6 @@ pub const Options = struct { /// then degrades to an uncompressed tail. refetch_dict: ?*const fn (ctx: ?*anyopaque, allocator: Allocator) ?[]u8 = null, refetch_dict_ctx: ?*anyopaque = null, - /// bearer credential sent as an Authorization header on the websocket - /// handshake, matching the upstream Go client (engine.go: the - /// negotiation transport carries the key on every request including - /// the live dial). the OSS server never gates the live tail, but - /// hosted instances can at their edge. null = no auth sent. - api_key: ?[]const u8 = null, /// declare the session stalled when no EVENT arrives for this long /// (surfaced as transient error.LiveStalled → reconnect/failover). /// event-based, not frame-based: a connected-but-silent server that @@ -206,18 +200,11 @@ pub const LiveClient = struct { defer client.deinit(); var headers_buf: [1024]u8 = undefined; - const headers = if (self.options.api_key) |key| - try std.fmt.bufPrint( - &headers_buf, - "Host: {s}\r\nSec-WebSocket-Protocol: {s}\r\nAuthorization: Bearer {s}\r\n", - .{ target.host, subprotocol, key }, - ) - else - try std.fmt.bufPrint( - &headers_buf, - "Host: {s}\r\nSec-WebSocket-Protocol: {s}\r\n", - .{ target.host, subprotocol }, - ); + const headers = try std.fmt.bufPrint( + &headers_buf, + "Host: {s}\r\nSec-WebSocket-Protocol: {s}\r\n", + .{ target.host, subprotocol }, + ); client.handshake(path, .{ .headers = headers }) catch |err| { if (err == error.InvalidHandshakeResponse) { if (client.handshake_failure) |failure| return classifyRejection(failure); diff --git a/src/loopback_test.zig b/src/loopback_test.zig index ccfe88d..c381dd2 100644 --- a/src/loopback_test.zig +++ b/src/loopback_test.zig @@ -55,6 +55,10 @@ const Instance = struct { stalled_plan: bool = false, archive_fault: enum { none, fetch, decode, segment } = .none, reject_plan: bool = false, + retry_block_calls: ?usize = null, + live_authorized: bool = false, + archive_authorized: bool = false, + archive_missing_auth: 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 }); @@ -133,9 +137,13 @@ const Instance = struct { return true; } if (std.mem.indexOf(u8, target, "subscribeEvents") != null) { + self.live_authorized = headerValue(request, "Authorization") != null; try self.handleWebsocket(request, target, r, w); return false; } + const authorized = if (headerValue(request, "Authorization")) |value| std.mem.eql(u8, value, "Bearer test-archive-key") else false; + self.archive_authorized = self.archive_authorized or authorized; + self.archive_missing_auth = self.archive_missing_auth or !authorized; if (std.mem.indexOf(u8, target, "listSegments") != null) { try respondJson(w, try self.segmentListJson(arena)); return false; @@ -170,6 +178,12 @@ const Instance = struct { return false; } if (std.mem.indexOf(u8, target, "getBlock") != null) { + if (self.retry_block_calls) |count| { + self.retry_block_calls = count + 1; + try w.writeAll("HTTP/1.1 503 Service Unavailable\r\nRetry-After: 0\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"); + try w.flush(); + return false; + } 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; @@ -703,3 +717,48 @@ test "a whole segment error can be accepted or declined before the next entry" { try testing.expectEqual(!accept, result.stopped); } } + +test "archive key authenticates archive requests but not live handshakes" { + 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 }}); + 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(2){}; + try client_mod.subscribe(io, testing.allocator, .{ + .hosts = &.{host.hostUrl(&url)}, + .after_seq = 0, + .api_key = "test-archive-key", + .concurrency = 1, + }, &handler); + try testing.expectEqual(@as(usize, 2), handler.len); + try testing.expect(host.archive_authorized); + try testing.expect(!host.archive_missing_auth); + try testing.expect(!host.live_authorized); +} + +test "archive block retries stop after three HTTP attempts" { + 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 }}, &.{}); + host.retry_block_calls = 0; + 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(2){}; + try testing.expectError(error.FetchRetriesExhausted, client_mod.subscribe(io, testing.allocator, .{ + .hosts = &.{host.hostUrl(&url)}, + .after_seq = 0, + .snapshot_only = true, + .concurrency = 1, + }, &handler)); + try testing.expectEqual(@as(?usize, 3), host.retry_block_calls); + try testing.expectEqual(@as(usize, 0), handler.len); +} -- 2.51.2