diff --git a/src/broadcaster.zig b/src/broadcaster.zig index 1035825..5f058ac 100644 --- a/src/broadcaster.zig +++ b/src/broadcaster.zig @@ -353,7 +353,7 @@ const ping_interval_ns: u64 = 30 * std.time.ns_per_s; const ping_timeout_ns: u64 = 5 * std.time.ns_per_s; pub const Consumer = struct { - const BUFFER_CAP = 8192; + const BUFFER_CAP = 65536; conn: *websocket.Conn, allocator: Allocator, @@ -427,7 +427,7 @@ pub const Consumer = struct { defer f.release(); self.conn.writeBin(f.data) catch { self.alive.store(false, .release); - return; + break; }; self.last_send_time = Io.Timestamp.now(self.io, .real).nanoseconds; } else { @@ -444,6 +444,11 @@ pub const Consumer = struct { f.release(); } else break; } + // close the connection from the writeLoop thread to unblock readLoop. + // this MUST happen here, not from the broadcast thread — closing from + // broadcast races with readLoop's setsockopt/read calls on the socket, + // causing EBADF → unreachable panic under ReleaseSafe. + self.conn.close(.{}) catch {}; } fn maybePing(self: *Consumer) void { @@ -622,13 +627,14 @@ pub const Broadcaster = struct { if (self.error_frame) |ef| { consumer.conn.writeBin(ef) catch {}; } - // signal consumer to stop — writeLoop will exit, then Handler.close → - // removeConsumer handles actual shutdown + destroy. - // IMPORTANT: do NOT spawn a plain std.Thread for cleanup — it would call - // Evented future cancel / mutex ops from a non-Uring thread → SIGSEGV. + // signal consumer to stop — writeLoop checks alive and exits, then + // Handler.close → removeConsumer handles shutdown + destroy. + // do NOT close the socket here: readLoop (websocket server thread) may + // be calling setsockopt/read on it concurrently. closing from a third + // thread causes EBADF → unreachable panic in posix.setsockopt under + // ReleaseSafe. let writeLoop's exit trigger the natural close sequence. consumer.alive.store(false, .release); consumer.cond.signal(consumer.io); - consumer.conn.close(.{}) catch {}; // consumer is removed from list by broadcast() caller (swapRemove). // cleanup deferred to removeConsumer to avoid double-destroy. }