From 261dcc5041130e3ca6a8495d1c2ec2d9c7fb7b33 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Wed, 1 Apr 2026 21:16:35 -0500 Subject: [PATCH] thread io: std.Io through streaming clients, fix firehose ws bug - add `io: std.Io` field to JetstreamClient and FirehoseClient - subscribe() returns Io.Cancelable!void (enables async cancellation) - replace libc.nanosleep() with io.sleep() for reconnect backoff - pass caller-provided io to websocket.Client.init() instead of debug_io - fix bug: firehose connectAndRead() was missing required `io` param - bump version to 0.3.0-alpha.6 Co-Authored-By: Claude Opus 4.6 --- build.zig.zon | 2 +- scripts/jetstream_smoke.zig | 4 ++-- src/internal/streaming/firehose.zig | 12 +++++++----- src/internal/streaming/jetstream.zig | 18 ++++++++++-------- 4 files changed, 20 insertions(+), 16 deletions(-) diff --git a/build.zig.zon b/build.zig.zon index d46579b..7b76975 100644 --- a/build.zig.zon +++ b/build.zig.zon @@ -1,6 +1,6 @@ .{ .name = .zat, - .version = "0.3.0-alpha.5", + .version = "0.3.0-alpha.6", .fingerprint = 0x8da9db57ee82fbe4, .minimum_zig_version = "0.16.0", .dependencies = .{ diff --git a/scripts/jetstream_smoke.zig b/scripts/jetstream_smoke.zig index 16c14b9..094846e 100644 --- a/scripts/jetstream_smoke.zig +++ b/scripts/jetstream_smoke.zig @@ -9,11 +9,11 @@ pub fn main() !void { std.debug.print("smoke test starting\n", .{}); var handler = Handler{}; - var client = zat.JetstreamClient.init(allocator, .{ + var client = zat.JetstreamClient.init(std.Options.debug_io, allocator, .{ .hosts = &.{"jetstream2.us-east.bsky.network"}, .wanted_collections = &.{"app.bsky.feed.post"}, }); - client.subscribe(&handler); + try client.subscribe(&handler); } const Handler = struct { diff --git a/src/internal/streaming/firehose.zig b/src/internal/streaming/firehose.zig index 946c133..9dd92d8 100644 --- a/src/internal/streaming/firehose.zig +++ b/src/internal/streaming/firehose.zig @@ -18,7 +18,7 @@ const sync = @import("sync.zig"); const mem = std.mem; const Allocator = mem.Allocator; const posix = std.posix; -const libc = std.c; +const Io = std.Io; const log = std.log.scoped(.zat); pub const CommitAction = sync.CommitAction; @@ -410,12 +410,14 @@ fn encodeInfoPayload(allocator: Allocator, writer: anytype, info: InfoEvent) !vo } pub const FirehoseClient = struct { + io: Io, allocator: Allocator, options: Options, last_seq: ?i64 = null, - pub fn init(allocator: Allocator, options: Options) FirehoseClient { + pub fn init(io: Io, allocator: Allocator, options: Options) FirehoseClient { return .{ + .io = io, .allocator = allocator, .options = options, .last_seq = if (options.cursor) |c| c else null, @@ -429,7 +431,7 @@ pub const FirehoseClient = struct { /// optional: fn onError(*@TypeOf(handler), anyerror) void /// blocks forever — reconnects with exponential backoff on disconnect. /// rotates through hosts on each reconnect attempt. - pub fn subscribe(self: *FirehoseClient, handler: anytype) void { + pub fn subscribe(self: *FirehoseClient, handler: anytype) Io.Cancelable!void { var backoff: u64 = 1; var host_index: usize = 0; const max_backoff: u64 = 60; @@ -456,7 +458,7 @@ pub const FirehoseClient = struct { prev_host_index = effective_index; host_index += 1; - _ = libc.nanosleep(&.{ .sec = @intCast(backoff), .nsec = 0 }, null); + try self.io.sleep(Io.Duration.fromSeconds(@intCast(backoff)), .awake); backoff = @min(backoff * 2, max_backoff); } } @@ -473,7 +475,7 @@ pub const FirehoseClient = struct { log.info("connecting to wss://{s}{s}", .{ host, path }); - var client = try websocket.Client.init(self.allocator, .{ + var client = try websocket.Client.init(self.io, self.allocator, .{ .host = host, .port = 443, .tls = true, diff --git a/src/internal/streaming/jetstream.zig b/src/internal/streaming/jetstream.zig index 1b72db5..eba2007 100644 --- a/src/internal/streaming/jetstream.zig +++ b/src/internal/streaming/jetstream.zig @@ -13,8 +13,8 @@ const sync = @import("sync.zig"); const mem = std.mem; const json = std.json; const posix = std.posix; -const libc = std.c; const Allocator = mem.Allocator; +const Io = std.Io; const log = std.log.scoped(.zat); pub const CommitAction = sync.CommitAction; @@ -131,12 +131,14 @@ pub fn parseEvent(allocator: Allocator, payload: []const u8) !Event { } pub const JetstreamClient = struct { + io: Io, allocator: Allocator, options: Options, last_time_us: ?i64 = null, - pub fn init(allocator: Allocator, options: Options) JetstreamClient { + pub fn init(io: Io, allocator: Allocator, options: Options) JetstreamClient { return .{ + .io = io, .allocator = allocator, .options = options, .last_time_us = options.cursor, @@ -151,7 +153,7 @@ pub const JetstreamClient = struct { /// optional: fn onConnect(*@TypeOf(handler), []const u8) void — called with host on connect /// blocks forever — reconnects with exponential backoff on disconnect. /// rotates through hosts on each reconnect attempt. - pub fn subscribe(self: *JetstreamClient, handler: anytype) void { + pub fn subscribe(self: *JetstreamClient, handler: anytype) Io.Cancelable!void { var backoff: u64 = 1; var host_index: usize = 0; const max_backoff: u64 = 60; @@ -181,7 +183,7 @@ pub const JetstreamClient = struct { prev_host_index = effective_index; host_index += 1; - _ = libc.nanosleep(&.{ .sec = @intCast(backoff), .nsec = 0 }, null); + try self.io.sleep(Io.Duration.fromSeconds(@intCast(backoff)), .awake); backoff = @min(backoff * 2, max_backoff); } } @@ -192,7 +194,7 @@ pub const JetstreamClient = struct { log.info("connecting to wss://{s}{s}", .{ host, path }); - var client = try websocket.Client.init(std.Options.debug_io, self.allocator, .{ + var client = try websocket.Client.init(self.io, self.allocator, .{ .host = host, .port = 443, .tls = true, @@ -462,7 +464,7 @@ test "Event.timeUs works for all variants" { } test "build subscribe path" { - var client = JetstreamClient.init(std.testing.allocator, .{ + var client = JetstreamClient.init(std.Options.debug_io, std.testing.allocator, .{ .wanted_collections = &.{"app.bsky.feed.post"}, }); @@ -472,7 +474,7 @@ test "build subscribe path" { } test "build subscribe path with multiple params" { - var client = JetstreamClient.init(std.testing.allocator, .{ + var client = JetstreamClient.init(std.Options.debug_io, std.testing.allocator, .{ .wanted_collections = &.{ "app.bsky.feed.post", "app.bsky.feed.like" }, .wanted_dids = &.{"did:plc:abc123"}, .cursor = 1700000000000, @@ -487,7 +489,7 @@ test "build subscribe path with multiple params" { } test "build subscribe path no params" { - var client = JetstreamClient.init(std.testing.allocator, .{}); + var client = JetstreamClient.init(std.Options.debug_io, std.testing.allocator, .{}); var buf: [2048]u8 = undefined; const path = try client.buildSubscribePath(&buf); -- 2.51.2