jetstream client
atproto jetstream client
Something went wrong. Try again.
Zig
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677//! live backfill→cutover smoke: sweep the REAL sealed archive from just//! below the sealed tip, cut over to the live tail, and verify the seam.//!//! this is the production proof the loopback e2e can't give: planSnapshot//! → getBlock → jss decode against the real (token-gated) archive, then//! the live websocket continuing PAST the sealed tip, delivered through//! one handler with every seq contiguous across the seam (unfiltered//! subscribe sees every seq, so any gap or duplicate is loud).//!//! run (from a box that may consume the firehose; needs the archive key)://! zig build example-cutover-smoke -- https://stream.waow.tech $API_KEY
const std = @import("std");const jetstream_sdk = @import("jetstream");
const sweep_back: u64 = 2000; // rows of sealed archive to replayconst live_target: usize = 25; // live rows to observe past the sealed tip
const SeamHandler = struct { sealed_tip: u64, first_seq: u64 = 0, last_seq: u64 = 0, total: usize = 0, archive_rows: usize = 0, live_rows: usize = 0, gaps: usize = 0,
pub fn onEvent(self: *@This(), event: jetstream_sdk.Event) bool { if (self.total == 0) { self.first_seq = event.seq; } else if (event.seq != self.last_seq + 1) { self.gaps += 1; std.debug.print(" GAP/DUP: {d} -> {d}\n", .{ self.last_seq, event.seq }); } self.last_seq = event.seq; self.total += 1; if (event.seq <= self.sealed_tip) self.archive_rows += 1 else self.live_rows += 1; if (self.live_rows == 1 and event.seq == self.sealed_tip + 1) std.debug.print(" cutover seam crossed contiguously at seq {d}\n", .{event.seq}); return self.live_rows < live_target; }};
pub fn main(init: std.process.Init) !void { const allocator = init.gpa; const args = try init.minimal.args.toSlice(init.arena.allocator()); if (args.len != 3) { std.debug.print("usage: example-cutover-smoke <base-url> <api-key>\n", .{}); return error.MissingArgs; } const host = args[1]; const api_key = args[2];
var arena_state = std.heap.ArenaAllocator.init(allocator); defer arena_state.deinit(); const spans = try jetstream_sdk.ArchiveBackfill.fetchSegmentSpans(init.io, allocator, arena_state.allocator(), host, api_key); var sealed_tip: u64 = 0; for (spans) |span| { if (span.max_seq > sealed_tip) sealed_tip = span.max_seq; } if (sealed_tip == 0) return error.NoSealedArchive; const after = sealed_tip - sweep_back; std.debug.print("sealed tip {d}; sweeping (after_seq {d}, {d} sealed rows] then live\n\n", .{ sealed_tip, after, sweep_back });
var handler = SeamHandler{ .sealed_tip = sealed_tip }; try jetstream_sdk.subscribe(init.io, allocator, .{ .hosts = &.{host}, .after_seq = after, .api_key = api_key, }, &handler);
std.debug.print("\nseqs {d}..{d}: {d} archive rows + {d} live rows, {d} gaps\n", .{ handler.first_seq, handler.last_seq, handler.archive_rows, handler.live_rows, handler.gaps }); if (handler.gaps != 0) return error.SeamNotContiguous; if (handler.archive_rows == 0) return error.NoArchiveRows; std.debug.print("smoke PASS: archive sweep -> gapless cutover -> live tail, one handler\n", .{});}