From a3ed87c3cc745dec66fdccc5a16e07bd9ee79d6c Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Sat, 3 Oct 2026 15:25:49 -0700 Subject: [PATCH] shutdown, broadcaster: never close a socket another fiber is parked on lalinsky/zio requires outstanding socket ops to be cancelled or awaited before the fd is closed. its epoll and kqueue backends assert it (a panic in ReleaseSafe); under io_uring the parked fiber is simply never woken. zlay did this in two places: - shutdown closed the websocket and metrics listeners to unblock accept. now the accept task is cancelled first, then the listener closed. - a consumer's writeLoop closed the connection to unblock its readLoop (slow-consumer eviction). it now sends the close frame and half-closes via util.net.interrupt; the reader wakes with EOF and owns the close. both orderings are also valid under Io.Threaded. the regression test runs on whichever backend the suite is built with. Co-Authored-By: Claude Opus 5.5 (1M context) --- src/internal/broadcaster.zig | 12 +++++---- src/internal/slurper.zig | 4 +-- src/internal/subscriber.zig | 6 ++--- src/internal/util/net.zig | 49 ++++++++++++++++++++++++++++++++++++ src/internal/util/util.zig | 1 + src/internal/ws_ping.zig | 4 +-- src/main.zig | 8 +++--- 7 files changed, 68 insertions(+), 16 deletions(-) create mode 100644 src/internal/util/net.zig diff --git a/src/internal/broadcaster.zig b/src/internal/broadcaster.zig index ddbb384..b1fcc7a 100644 --- a/src/internal/broadcaster.zig +++ b/src/internal/broadcaster.zig @@ -501,11 +501,13 @@ 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 {}; + // unblock readLoop, which is parked on this socket and owns its close. + // this MUST happen here, not from the broadcast thread, which would + // race readLoop's setsockopt/read calls on the socket. + if (!self.conn.isClosed()) { + self.conn.writeClose(.{}) catch {}; + util.net.interrupt(self.io, self.conn.stream); + } } fn maybePing(self: *Consumer) void { diff --git a/src/internal/slurper.zig b/src/internal/slurper.zig index f574d0a..57a88f7 100644 --- a/src/internal/slurper.zig +++ b/src/internal/slurper.zig @@ -598,8 +598,8 @@ pub const Slurper = struct { /// has already freed the cursor slot, dropped the map entry and /// decremented `connected_inbound` -- which is why the respawn below no /// longer trips the dedup. Cancelling an already-finished future is the - /// designed path, not a use-after-free: zio's Task is refcounted so it - /// outlives its fiber until awaited or cancelled exactly once. + /// designed path, not a use-after-free: a std.Io future outlives its + /// fiber until awaited or cancelled exactly once. /// /// `min_idle_ms` is the guard. A deaf host and a merely quiet host look /// identical from seq sampling, so a reconnect driven by silence alone diff --git a/src/internal/subscriber.zig b/src/internal/subscriber.zig index 1c8aad7..842510a 100644 --- a/src/internal/subscriber.zig +++ b/src/internal/subscriber.zig @@ -437,7 +437,7 @@ pub const Subscriber = struct { if (comptime backend_config.use_zio) { if (self.ws_ping != null) { // under zio, readLoopWithHeartbeat's SO_RCVTIMEO mechanism is - // inert: netRead parks the fiber in epoll and the socket + // inert: netRead parks the fiber in the event loop and the socket // timeout never fires (the 2026-08-18/27 half-open wedges). // run the read loop as a cancelable task and drive liveness // from this fiber instead. cancellation (not a cross-fiber @@ -540,8 +540,8 @@ pub const Subscriber = struct { }; /// per-connection liveness state shared between the read-loop task and the -/// pinger. all fields are atomics for the Threaded test path; under zio both -/// sides are fibers on the one loop thread. +/// pinger. all fields are atomics: the two sides are separate threads under +/// Threaded and fibers on possibly different executor threads under zio. const ConnState = struct { /// ms timestamp of the last frame received on THIS connection (any type — /// data, pong). distinct from Subscriber.last_progress_ms, which spans diff --git a/src/internal/util/net.zig b/src/internal/util/net.zig new file mode 100644 index 0000000..13e862a --- /dev/null +++ b/src/internal/util/net.zig @@ -0,0 +1,49 @@ +//! socket helpers that hold on every std.Io backend. + +const std = @import("std"); +const Io = std.Io; +const backend_config = @import("backend_config"); +const zio = @import("zio"); + +/// end a connection that another fiber may be parked on. half-closes the +/// socket so the parked side wakes with EOF and closes the descriptor itself. +/// closing it from here instead trips an assert in zio's epoll/kqueue backends +/// and, under io_uring, leaves the parked fiber asleep. +pub fn interrupt(io: Io, stream: Io.net.Stream) void { + stream.shutdown(io, .both) catch {}; +} + +fn parkedReader(io: Io, stream: Io.net.Stream, woke_with_eof: *std.atomic.Value(bool)) void { + var buf: [16]u8 = undefined; + var reader = stream.reader(io, &buf); + _ = reader.interface.takeByte() catch |err| { + woke_with_eof.store(err == error.EndOfStream, .release); + }; +} + +fn expectInterruptWakesReader(io: Io) !void { + var listener = try (Io.net.IpAddress{ .ip4 = .loopback(0) }).listen(io, .{}); + defer listener.deinit(io); + const client = try listener.socket.address.connect(io, .{ .mode = .stream }); + defer client.close(io); + const served = try listener.accept(io); + + var woke_with_eof: std.atomic.Value(bool) = .init(false); + var reader = try io.concurrent(parkedReader, .{ io, served, &woke_with_eof }); + try io.sleep(.fromMilliseconds(50), .awake); + + interrupt(io, served); + reader.await(io); + served.close(io); + try std.testing.expect(woke_with_eof.load(.acquire)); +} + +test "interrupt wakes a reader parked on the socket" { + if (comptime backend_config.use_zio) { + const rt = try zio.Runtime.init(std.testing.allocator, .{ .executors = .exact(2) }); + defer rt.deinit(); + try expectInterruptWakesReader(rt.io()); + } else { + try expectInterruptWakesReader(std.testing.io); + } +} diff --git a/src/internal/util/util.zig b/src/internal/util/util.zig index c6a9b7b..c4f6e46 100644 --- a/src/internal/util/util.zig +++ b/src/internal/util/util.zig @@ -4,6 +4,7 @@ pub const env = @import("env.zig"); pub const time = @import("time.zig"); +pub const net = @import("net.zig"); // flattened leaves — the common case pub const getenv = env.getenv; diff --git a/src/internal/ws_ping.zig b/src/internal/ws_ping.zig index 0865067..743e3e6 100644 --- a/src/internal/ws_ping.zig +++ b/src/internal/ws_ping.zig @@ -9,8 +9,8 @@ //! //! websocket.zig's readLoopWithHeartbeat closes that gap under Io.Threaded, //! where SO_RCVTIMEO makes read() return after interval_ms. under zio it is -//! structurally inert: netRead parks the fiber in epoll and the socket -//! timeout never fires. this module is the zio-side replacement — the +//! structurally inert: netRead parks the fiber in the event loop and the +//! socket timeout never fires. this module is the zio-side replacement — the //! subscriber fiber runs a ping loop (subscriber.pingLoop) against the //! decision function here, while the read loop runs as a cancelable task. //! diff --git a/src/main.zig b/src/main.zig index 3ce0198..171e4ee 100644 --- a/src/main.zig +++ b/src/main.zig @@ -560,9 +560,10 @@ fn runRelay() !void { log.info("shutdown signal received, stopping...", .{}); - // stop WebSocket server: close listener to unblock accept, then cancel task - ws_listener.deinit(io); + // stop WebSocket server: cancel the accept task, then close its listener. + // never the other way round — see util.net.interrupt. server_future.cancel(io); + ws_listener.deinit(io); // join GC thread — it checks shutdown_flag and will exit its sleep loop. // must complete before dp.deinit() runs (dp is stack-owned). @@ -586,9 +587,8 @@ fn runRelay() !void { // cancel broadcaster fiber (shutdown flag already set, it will drain remaining) broadcast_future.cancel(io); - // close metrics listener to unblock accept(), then cancel task - metrics_srv.server.deinit(io); metrics_future.cancel(io); + metrics_srv.server.deinit(io); log.info("relay stopped cleanly", .{}); } -- 2.51.2