From 395d033dd6b51692aaf24aa8949cefd77fdfbeff Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Sun, 1 Mar 2026 01:18:22 -0600 Subject: [PATCH] fix: use u64 for seq numbers, fix WebSocket bind, add favicon MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit upstream firehose seq numbers exceed i64 max — switch all seq types to u64 across subscriber, broadcaster, ring_buffer, and event_log. SQLite storage uses XOR-with-sign-bit encoding to preserve ordering. update zat dependency to v0.2.7 (adds cbor.getUint). also: bind WebSocket server to 0.0.0.0 (was 127.0.0.1, unreachable in k8s), add /favicon.svg route. Co-Authored-By: Claude Opus 4.6 --- build.zig.zon | 4 ++-- src/broadcaster.zig | 29 ++++++++++++++--------------- src/event_log.zig | 23 +++++++++++++++++------ src/main.zig | 8 ++++++++ src/ring_buffer.zig | 30 +++++++++++++++--------------- src/subscriber.zig | 6 +++--- 6 files changed, 59 insertions(+), 41 deletions(-) diff --git a/build.zig.zon b/build.zig.zon index bfcaba3..9f1d452 100644 --- a/build.zig.zon +++ b/build.zig.zon @@ -5,8 +5,8 @@ .minimum_zig_version = "0.15.0", .dependencies = .{ .zat = .{ - .url = "https://tangled.org/zat.dev/zat/archive/0071362.tar.gz", - .hash = "zat-0.2.6-5PuC7ka3BAD826XlwCMgmCA21Xy7TKoQBlVlfvEHEENm", + .url = "https://tangled.org/zat.dev/zat/archive/v0.2.7.tar.gz", + .hash = "zat-0.2.7-5PuC7oC4BADf6cKEjxRsLlg8SVYDbC1SAQanOLqMl2-Q", }, .websocket = .{ .url = "https://github.com/karlseguin/websocket.zig/archive/97fefafa59cc78ce177cff540b8685cd7f699276.tar.gz", diff --git a/src/broadcaster.zig b/src/broadcaster.zig index 96c214c..c082f33 100644 --- a/src/broadcaster.zig +++ b/src/broadcaster.zig @@ -19,7 +19,7 @@ const log = std.log.scoped(.relay); // --- stats --- pub const Stats = struct { - seq: std.atomic.Value(i64) = .{ .raw = 0 }, + seq: std.atomic.Value(u64) = .{ .raw = 0 }, relay_seq: std.atomic.Value(u64) = .{ .raw = 0 }, consumer_count: std.atomic.Value(u64) = .{ .raw = 0 }, frames_in: std.atomic.Value(u64) = .{ .raw = 0 }, @@ -296,7 +296,7 @@ pub const Broadcaster = struct { /// broadcast a frame to all consumers. non-blocking — just enqueues. /// drops slow consumers whose buffers are full after sending ConsumerTooSlow. - pub fn broadcast(self: *Broadcaster, seq: i64, data: []const u8) void { + pub fn broadcast(self: *Broadcaster, seq: u64, data: []const u8) void { _ = self.stats.frames_in.fetchAdd(1, .monotonic); self.stats.seq.store(seq, .release); @@ -359,17 +359,16 @@ pub const Broadcaster = struct { /// two-phase cursor replay: disk (diskpersist) first, then in-memory ring buffer. /// the consumer is already in the live broadcast list, so frames arriving /// during replay are buffered — no gap possible. - pub fn replayTo(self: *Broadcaster, consumer: *Consumer, cursor: i64) void { + pub fn replayTo(self: *Broadcaster, consumer: *Consumer, cursor: u64) void { // phase 1: disk replay from diskpersist if (self.persist) |dp| { - const cursor_u64: u64 = if (cursor > 0) @intCast(cursor) else 0; var entries: std.ArrayListUnmanaged(event_log_mod.PlaybackEntry) = .{}; defer { for (entries.items) |e| self.allocator.free(e.data); entries.deinit(self.allocator); } - dp.playback(cursor_u64, self.allocator, &entries) catch |err| { + dp.playback(cursor, self.allocator, &entries) catch |err| { log.warn("disk replay failed: {s}, falling back to memory", .{@errorName(err)}); self.replayFromMemory(consumer, cursor); return; @@ -393,7 +392,7 @@ pub const Broadcaster = struct { self.replayFromMemory(consumer, cursor); } - fn replayFromMemory(self: *Broadcaster, consumer: *Consumer, cursor: i64) void { + fn replayFromMemory(self: *Broadcaster, consumer: *Consumer, cursor: u64) void { self.history.mutex.lock(); defer self.history.mutex.unlock(); @@ -420,7 +419,7 @@ pub const Handler = struct { consumer: ?*Consumer, broadcaster: *Broadcaster, conn: *websocket.Conn, - cursor: ?i64, + cursor: ?u64, /// called during handshake — validate path, parse cursor. /// do NOT start the write thread here (handshake reply hasn't been sent yet). @@ -433,12 +432,12 @@ pub const Handler = struct { } // parse cursor (deferred until afterInit when connection is fully upgraded) - var cursor: ?i64 = null; + var cursor: ?u64 = null; if (std.mem.indexOf(u8, url, "cursor=")) |cursor_start| { const value_start = cursor_start + "cursor=".len; const value_end = std.mem.indexOfScalarPos(u8, url, value_start, '&') orelse url.len; const cursor_str = url[value_start..value_end]; - cursor = std.fmt.parseInt(i64, cursor_str, 10) catch null; + cursor = std.fmt.parseInt(u64, cursor_str, 10) catch null; } return .{ .consumer = null, .broadcaster = ctx, .conn = conn, .cursor = cursor }; @@ -550,7 +549,7 @@ test "broadcaster add and remove consumer" { var b = Broadcaster.init(std.testing.allocator); defer b.deinit(); - try std.testing.expectEqual(@as(i64, 0), b.stats.seq.load(.acquire)); + try std.testing.expectEqual(@as(u64, 0), b.stats.seq.load(.acquire)); try std.testing.expectEqual(@as(usize, 0), b.consumerCount()); } @@ -562,10 +561,10 @@ test "broadcast updates stats and history" { b.broadcast(2, "frame2"); b.broadcast(3, "frame3"); - try std.testing.expectEqual(@as(i64, 3), b.stats.seq.load(.acquire)); + try std.testing.expectEqual(@as(u64, 3), b.stats.seq.load(.acquire)); try std.testing.expectEqual(@as(u64, 3), b.stats.frames_in.load(.acquire)); - try std.testing.expectEqual(@as(i64, 1), b.history.oldestSeq().?); - try std.testing.expectEqual(@as(i64, 3), b.history.newestSeq().?); + try std.testing.expectEqual(@as(u64, 1), b.history.oldestSeq().?); + try std.testing.expectEqual(@as(u64, 3), b.history.newestSeq().?); } test "frame history supports cursor replay" { @@ -585,8 +584,8 @@ test "frame history supports cursor replay" { } try std.testing.expectEqual(@as(usize, 2), frames.len); - try std.testing.expectEqual(@as(i64, 4), frames[0].seq); - try std.testing.expectEqual(@as(i64, 5), frames[1].seq); + try std.testing.expectEqual(@as(u64, 4), frames[0].seq); + try std.testing.expectEqual(@as(u64, 5), frames[1].seq); } test "shared frame ref counting" { diff --git a/src/event_log.zig b/src/event_log.zig index 9195b0b..50a46ea 100644 --- a/src/event_log.zig +++ b/src/event_log.zig @@ -38,6 +38,17 @@ const default_events_per_file: u32 = 10_000; const default_flush_interval_ms: u64 = 100; const default_flush_threshold: usize = 400; +/// convert u64 seq to i64 for SQLite storage. +/// XOR with sign bit to preserve ordering across the full u64 range. +fn seqToSqlite(seq: u64) i64 { + return @bitCast(seq ^ (1 << 63)); +} + +/// convert i64 from SQLite back to u64 seq. +fn sqliteToSeq(val: i64) u64 { + return @as(u64, @bitCast(val)) ^ (1 << 63); +} + // --- header --- pub const EvtHeader = struct { @@ -340,12 +351,12 @@ pub const DiskPersist = struct { if (since > 0) { // find file whose seq_start is just before `since` - if (self.db.row("SELECT id, path, seq_start FROM log_file_refs WHERE seq_start <= ? ORDER BY seq_start DESC LIMIT 1", .{since})) |row| { + if (self.db.row("SELECT id, path, seq_start FROM log_file_refs WHERE seq_start <= ? ORDER BY seq_start DESC LIMIT 1", .{seqToSqlite(since)})) |row| { if (row) |r| { defer r.deinit(); try start_files.append(allocator, .{ .path = try allocator.dupe(u8, r.text(1)), - .seq_start = @intCast(r.int(2)), + .seq_start = sqliteToSeq(r.int(2)), }); } } else |_| {} @@ -353,12 +364,12 @@ pub const DiskPersist = struct { // find all subsequent files { - var rows = try self.db.rows("SELECT id, path, seq_start FROM log_file_refs WHERE seq_start > ? ORDER BY seq_start ASC", .{since}); + var rows = try self.db.rows("SELECT id, path, seq_start FROM log_file_refs WHERE seq_start > ? ORDER BY seq_start ASC", .{seqToSqlite(since)}); defer rows.deinit(); while (rows.next()) |r| { try start_files.append(allocator, .{ .path = try allocator.dupe(u8, r.text(1)), - .seq_start = @intCast(r.int(2)), + .seq_start = sqliteToSeq(r.int(2)), }); } } @@ -469,7 +480,7 @@ pub const DiskPersist = struct { if (r) |row| { defer row.deinit(); const path = row.text(1); - const seq_start: u64 = @intCast(row.int(2)); + const seq_start: u64 = sqliteToSeq(row.int(2)); var file = self.dir.openFile(path, .{ .mode = .read_write }) catch { // file missing, start fresh @@ -519,7 +530,7 @@ pub const DiskPersist = struct { // register in SQLite try self.db.exec( "INSERT INTO log_file_refs (path, seq_start) VALUES (?, ?)", - .{ name, @as(i64, @intCast(start_seq)) }, + .{ name, seqToSqlite(start_seq) }, ); self.event_counter = 0; diff --git a/src/main.zig b/src/main.zig index 06fdb80..9a195d9 100644 --- a/src/main.zig +++ b/src/main.zig @@ -116,6 +116,7 @@ pub fn main() !void { var server = try websocket.Server(broadcaster.Handler).init(allocator, .{ .port = port, + .address = "0.0.0.0", .max_conn = 4096, .max_message_size = 5 * 1024 * 1024, }); @@ -249,6 +250,13 @@ fn handleGet(stream: std.net.Stream, path: []const u8, stats: *broadcaster.Stats \\The firehose WebSocket path is at: /xrpc/com.atproto.sync.subscribeRepos \\ ); + } else if (std.mem.eql(u8, path, "/favicon.svg") or std.mem.eql(u8, path, "/favicon.ico")) { + httpRespond(stream, "200 OK", "image/svg+xml", + \\ + \\ + \\Z + \\ + ); } else { httpRespond(stream, "404 Not Found", "text/plain", "not found"); } diff --git a/src/ring_buffer.zig b/src/ring_buffer.zig index 95d0fbd..b87d4ca 100644 --- a/src/ring_buffer.zig +++ b/src/ring_buffer.zig @@ -7,7 +7,7 @@ const std = @import("std"); const Allocator = std.mem.Allocator; pub const Frame = struct { - seq: i64, + seq: u64, data: []const u8, // owned by the ring buffer pub const empty: Frame = .{ .seq = 0, .data = &.{} }; @@ -45,13 +45,13 @@ pub fn RingBuffer(comptime capacity: usize) type { } /// push a frame. if full, overwrites oldest. returns false if alloc failed. - pub fn push(self: *Self, seq: i64, data: []const u8) bool { + pub fn push(self: *Self, seq: u64, data: []const u8) bool { self.mutex.lock(); defer self.mutex.unlock(); return self.pushUnlocked(seq, data); } - fn pushUnlocked(self: *Self, seq: i64, data: []const u8) bool { + fn pushUnlocked(self: *Self, seq: u64, data: []const u8) bool { const duped = self.allocator.dupe(u8, data) catch return false; // free old entry if overwriting @@ -103,7 +103,7 @@ pub fn RingBuffer(comptime capacity: usize) type { /// get all frames with seq > cursor, ordered by seq. /// caller owns the returned slice AND frame data. - pub fn framesSince(self: *Self, allocator: Allocator, cursor: i64) ![]const Frame { + pub fn framesSince(self: *Self, allocator: Allocator, cursor: u64) ![]const Frame { self.mutex.lock(); defer self.mutex.unlock(); @@ -126,7 +126,7 @@ pub fn RingBuffer(comptime capacity: usize) type { } /// oldest seq in the buffer, or null if empty - pub fn oldestSeq(self: *Self) ?i64 { + pub fn oldestSeq(self: *Self) ?u64 { self.mutex.lock(); defer self.mutex.unlock(); if (self.len == 0) return null; @@ -134,7 +134,7 @@ pub fn RingBuffer(comptime capacity: usize) type { } /// newest seq in the buffer, or null if empty - pub fn newestSeq(self: *Self) ?i64 { + pub fn newestSeq(self: *Self) ?u64 { self.mutex.lock(); defer self.mutex.unlock(); if (self.len == 0) return null; @@ -156,12 +156,12 @@ test "push and pop" { const f1 = buf.pop().?; defer std.testing.allocator.free(f1.data); - try std.testing.expectEqual(@as(i64, 1), f1.seq); + try std.testing.expectEqual(@as(u64, 1), f1.seq); try std.testing.expectEqualStrings("hello", f1.data); const f2 = buf.pop().?; defer std.testing.allocator.free(f2.data); - try std.testing.expectEqual(@as(i64, 2), f2.seq); + try std.testing.expectEqual(@as(u64, 2), f2.seq); try std.testing.expect(buf.pop() == null); } @@ -181,7 +181,7 @@ test "overwrite when full" { const f1 = buf.pop().?; defer std.testing.allocator.free(f1.data); - try std.testing.expectEqual(@as(i64, 2), f1.seq); + try std.testing.expectEqual(@as(u64, 2), f1.seq); } test "framesSince" { @@ -199,8 +199,8 @@ test "framesSince" { } try std.testing.expectEqual(@as(usize, 2), frames.len); - try std.testing.expectEqual(@as(i64, 4), frames[0].seq); - try std.testing.expectEqual(@as(i64, 5), frames[1].seq); + try std.testing.expectEqual(@as(u64, 4), frames[0].seq); + try std.testing.expectEqual(@as(u64, 5), frames[1].seq); } test "oldestSeq and newestSeq" { @@ -214,8 +214,8 @@ test "oldestSeq and newestSeq" { try std.testing.expect(buf.push(20, "y")); try std.testing.expect(buf.push(30, "z")); - try std.testing.expectEqual(@as(i64, 10), buf.oldestSeq().?); - try std.testing.expectEqual(@as(i64, 30), buf.newestSeq().?); + try std.testing.expectEqual(@as(u64, 10), buf.oldestSeq().?); + try std.testing.expectEqual(@as(u64, 30), buf.newestSeq().?); } test "empty buffer operations" { @@ -251,6 +251,6 @@ test "wrap-around with pop and push" { try std.testing.expect(buf.push(5, "e")); try std.testing.expectEqual(@as(usize, 3), buf.count()); - try std.testing.expectEqual(@as(i64, 3), buf.oldestSeq().?); - try std.testing.expectEqual(@as(i64, 5), buf.newestSeq().?); + try std.testing.expectEqual(@as(u64, 3), buf.oldestSeq().?); + try std.testing.expectEqual(@as(u64, 5), buf.newestSeq().?); } diff --git a/src/subscriber.zig b/src/subscriber.zig index 7ba8e7f..69df85b 100644 --- a/src/subscriber.zig +++ b/src/subscriber.zig @@ -27,7 +27,7 @@ pub const Subscriber = struct { validator: *validator_mod.Validator, persist: ?*event_log_mod.DiskPersist, shutdown: *std.atomic.Value(bool), - last_upstream_seq: ?i64 = null, + last_upstream_seq: ?u64 = null, crawl_requests: std.ArrayListUnmanaged([]const u8) = .{}, crawl_mutex: std.Thread.Mutex = .{}, @@ -150,7 +150,7 @@ const FrameHandler = struct { }; // extract seq for cursor tracking (all event types have seq) - const upstream_seq = payload.getInt("seq"); + const upstream_seq = payload.getUint("seq"); if (upstream_seq) |s| { sub.last_upstream_seq = s; sub.bc.stats.seq.store(s, .release); @@ -208,7 +208,7 @@ const FrameHandler = struct { } sub.bc.stats.relay_seq.store(relay_seq, .release); - sub.bc.broadcast(@intCast(relay_seq), data); + sub.bc.broadcast(relay_seq, data); } else { sub.bc.broadcast(upstream_seq orelse 0, data); } -- 2.51.2