diff --git a/CHANGELOG.md b/CHANGELOG.md index fea9112..3ea684a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,8 +2,15 @@ ## unreleased -`zat.ArchiveBackfill` works against post-proposal-0015 jetstream archives, -and can bound a plan to a time window instead of draining the world. +`zat.ArchiveBackfill` works against post-proposal-0015 jetstream archives — +including Bluesky's token-gated hosted instances (Jetstream v2 launched +officially 2026-08-13 as Bluesky Protocol Services) — and can bound a plan +to a time window instead of draining the world. + +- **feat**: `Options.api_key` sends `Authorization: Bearer ` on + planSnapshot/listSegments/getSegment/getBlock (mirrors the Go client's + `WithAPIKey`); `fetchSeqBounds` takes the same optional key. The live + tail never needs one, matching upstream's "live stays open" contract. - **fix**: the plan endpoint is `network.bsky.jetstream.planSnapshot` — upstream renamed it from `planBackfill` (bluesky-social/jetstream PR #330) diff --git a/examples/07_streamplace_chat.zig b/examples/07_streamplace_chat.zig index 2ce656d..80714dc 100644 --- a/examples/07_streamplace_chat.zig +++ b/examples/07_streamplace_chat.zig @@ -209,7 +209,7 @@ pub fn main(init: std.process.Init) !void { // capture began replays completely — including messages later deleted. // The witnessed-time window maps to plan seq bounds, which keeps the // fetch to the handful of blocks the session actually touches. - const bounds = try zat.ArchiveBackfill.fetchSeqBounds(init.io, allocator, archive_host, start_us, end_cursor); + const bounds = try zat.ArchiveBackfill.fetchSeqBounds(init.io, allocator, archive_host, start_us, end_cursor, null); var archive_boundary_us: ?i64 = null; if (bounds.covered) { std.debug.print("replaying from the archive ({s})...\n\n", .{archive_host}); diff --git a/src/internal/streaming/archive_backfill.zig b/src/internal/streaming/archive_backfill.zig index 5cf1904..9cb0b2a 100644 --- a/src/internal/streaming/archive_backfill.zig +++ b/src/internal/streaming/archive_backfill.zig @@ -51,6 +51,11 @@ pub const Options = struct { collections: []const []const u8 = &.{}, /// dids to replay; empty = all dids: []const []const u8 = &.{}, + /// bearer credential for token-gated archives (Bluesky's hosted + /// instances 401 archive endpoints without one; the live tail never + /// needs it). mirrors the Go client's WithAPIKey. sent as + /// "Authorization: Bearer " on plan/list/segment/block requests. + api_key: ?[]const u8 = null, /// exclusive lower seq bound for the plan; null = from the beginning. /// derive from a time window with fetchSeqBounds. after_seq: ?u64 = null, @@ -107,10 +112,13 @@ pub fn run(io: Io, allocator: Allocator, options: Options, handler: anytype) !Re var transport = HttpTransport.init(io, allocator); defer transport.deinit(); + const authorization = try bearerHeader(allocator, options.api_key); + defer if (authorization) |a| allocator.free(a); + var plan_arena = std.heap.ArenaAllocator.init(allocator); defer plan_arena.deinit(); - const plan = try fetchPlan(plan_arena.allocator(), &transport, options); + const plan = try fetchPlan(plan_arena.allocator(), &transport, options, authorization); var result: Result = .{ .planned_through_seq = plan.planned_through_seq, @@ -120,6 +128,7 @@ pub fn run(io: Io, allocator: Allocator, options: Options, handler: anytype) !Re const concurrency = @min(@max(options.concurrency, 1), max_concurrency); var window = try FetchWindow.init(io, allocator, concurrency); defer window.deinit(); + window.authorization = authorization; for (plan.segments) |segment| { if (segment.blocks) |ranges| { @@ -159,7 +168,7 @@ 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); + var fetched = try fetchWithRetry(io, &transport, url, authorization); defer fetched.deinit(allocator); try deliverSegment(allocator, fetched.body, options, segment, handler, &result); } @@ -176,8 +185,8 @@ const max_concurrency = 16; const FetchError = anyerror; -fn fetchJob(io: Io, transport: *HttpTransport, url: []const u8) FetchError!HttpTransport.FetchResult { - return fetchWithRetry(io, transport, url); +fn fetchJob(io: Io, transport: *HttpTransport, url: []const u8, authorization: ?[]const u8) FetchError!HttpTransport.FetchResult { + return fetchWithRetry(io, transport, url, authorization); } /// a FIFO window of in-flight block fetches. starting is out-of-order I/O; @@ -189,6 +198,8 @@ const FetchWindow = struct { allocator: Allocator, transports: []HttpTransport, slots: []Inflight, + /// borrowed from run(); outlives every in-flight fetch + authorization: ?[]const u8 = null, head: usize = 0, count: usize = 0, @@ -227,7 +238,7 @@ const FetchWindow = struct { errdefer self.allocator.free(url); const slot = (self.head + self.count) % self.slots.len; self.slots[slot] = .{ - .future = self.io.async(fetchJob, .{ self.io, &self.transports[slot], url }), + .future = self.io.async(fetchJob, .{ self.io, &self.transports[slot], url, self.authorization }), .url = url, .position = position, }; @@ -299,12 +310,12 @@ const fetch_attempts = 4; /// 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) !HttpTransport.FetchResult { +fn fetchWithRetry(io: Io, transport: *HttpTransport, url: []const u8, authorization: ?[]const u8) !HttpTransport.FetchResult { var attempt: u32 = 0; var backoff_s: u64 = 1; while (true) { attempt += 1; - if (transport.fetch(.{ .url = url })) |result| { + if (transport.fetch(.{ .url = url, .authorization = authorization })) |result| { if (result.status == .ok) return result; var discarded = result; discarded.deinit(transport.allocator); @@ -373,13 +384,19 @@ fn planRequestBody(arena: Allocator, options: Options) ![]const u8 { return body.written(); } -fn fetchPlan(arena: Allocator, transport: *HttpTransport, options: Options) !Plan { +fn bearerHeader(allocator: Allocator, api_key: ?[]const u8) !?[]u8 { + const key = api_key orelse return null; + return try std.fmt.allocPrint(allocator, "Bearer {s}", .{key}); +} + +fn fetchPlan(arena: Allocator, transport: *HttpTransport, options: Options, authorization: ?[]const u8) !Plan { const body = try planRequestBody(arena, options); const url = try std.fmt.allocPrint(arena, "{s}/xrpc/network.bsky.jetstream.planSnapshot", .{options.host}); var fetched = try transport.fetch(.{ .url = url, .method = .POST, .payload = body, + .authorization = authorization, }); defer fetched.deinit(transport.allocator); if (fetched.status != .ok) { @@ -436,9 +453,11 @@ pub fn seqBoundsForWindow(spans: []const SegmentSpan, start_us: i64, end_us: ?i6 /// fetch every sealed segment's metadata from `host` and derive the seq /// bounds covering the witnessed window [start_us, end_us]. end_us == null /// means "through the sealed tip". -pub fn fetchSeqBounds(io: Io, allocator: Allocator, host: []const u8, start_us: i64, end_us: ?i64) !SeqBounds { +pub fn fetchSeqBounds(io: Io, allocator: Allocator, host: []const u8, start_us: i64, end_us: ?i64, api_key: ?[]const u8) !SeqBounds { var transport = HttpTransport.init(io, allocator); defer transport.deinit(); + const authorization = try bearerHeader(allocator, api_key); + defer if (authorization) |a| allocator.free(a); var arena_state = std.heap.ArenaAllocator.init(allocator); defer arena_state.deinit(); const arena = arena_state.allocator(); @@ -450,7 +469,7 @@ pub fn fetchSeqBounds(io: Io, allocator: Allocator, host: []const u8, start_us: 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); + var fetched = try fetchWithRetry(io, &transport, url, authorization); defer fetched.deinit(allocator); const parsed = try json.parseFromSliceLeaky(json.Value, arena, fetched.body, .{}); const root = switch (parsed) {