jetstream client
atproto jetstream client
Something went wrong. Try again.
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374//! live failover smoke: a dead primary rotates to a real fallback.//!//! the loopback e2e (src/loopback_test.zig) proves the mechanics against//! fake instances; this proves the wild-network half — DNS, TLS, and a//! production v2 server as the fallback. the primary is a blackhole by//! construction (connection refused), so the run exercises the error-//! threshold rotation path end to end://!//! refused primary ×N → rotate → fresh live tip on the fallback →//! deliver real events → clean stop//!//! run (from a box that may consume the firehose)://! zig build example-failover-smoke -- https://stream.waow.tech [--rewind]
const std = @import("std");const jetstream_sdk = @import("jetstream");
const stop_after = 10;
const SmokeHandler = struct { events: usize = 0, first_seq: u64 = 0, last_seq: u64 = 0,
pub fn onEvent(self: *@This(), event: jetstream_sdk.Event) bool { if (self.events == 0) self.first_seq = event.seq; self.events += 1; self.last_seq = event.seq; std.debug.print(" event {d}/{d}: seq={d} did={s} kind={s}\n", .{ self.events, stop_after, event.seq, event.did, @tagName(event.payload), }); return self.events < stop_after; }
pub fn onError(self: *@This(), err: anyerror) bool { _ = self; std.debug.print(" (transient: {s})\n", .{@errorName(err)}); return true; }};
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 < 2 or args.len > 3) { std.debug.print("usage: zig build example-failover-smoke -- <fallback-base-url> [--rewind]\n", .{}); return error.MissingFallback; } const fallback = args[1]; const rewind = args.len == 3; if (rewind and !std.mem.eql(u8, args[2], "--rewind")) return error.InvalidArgument; const cursor: ?u64 = if (rewind) @intCast(std.Io.Timestamp.now(init.io, .real).toMicroseconds() - 10 * std.time.us_per_s) else null; if (cursor) |value| std.debug.print("timestamp cursor: {d} (10s rewind)\n", .{value});
// loopback port 1: nothing listens there — instant connection refused const dead_primary = "http://127.0.0.1:1";
std.debug.print("primary: {s} (dead by construction)\nfallback: {s}\n\n", .{ dead_primary, fallback }); const started = std.Io.Timestamp.now(init.io, .awake).nanoseconds;
var handler = SmokeHandler{}; try jetstream_sdk.subscribe(init.io, allocator, .{ .hosts = &.{ dead_primary, fallback }, .live_cursor = cursor, .failover_after_errors = 3, .backoff_min_ns = 100 * std.time.ns_per_ms, .backoff_max_ns = 1 * std.time.ns_per_s, }, &handler);
if (handler.events != stop_after) return error.IncompleteSmoke; const elapsed_ms = @divFloor(std.Io.Timestamp.now(init.io, .awake).nanoseconds - started, std.time.ns_per_ms); std.debug.print("\nsmoke PASS: rotated off the dead primary, {d} live events from the fallback (seq {d}..{d}) in {d}ms\n", .{ handler.events, handler.first_seq, handler.last_seq, elapsed_ms });}