diff --git a/Dockerfile b/Dockerfile index 8ff2484..0f67df2 100644 --- a/Dockerfile +++ b/Dockerfile @@ -8,6 +8,10 @@ RUN curl -fSL https://ziglang.org/builds/zig-x86_64-linux-0.16.0-dev.3059+42e33d ENV PATH=/opt/zig-x86_64-linux-0.16.0-dev.3059+42e33db9d:$PATH WORKDIR /build +# patch Io.Uring networking (stubbed as Unavailable upstream, zig#31723) +COPY patches/ patches/ +RUN patch /opt/zig-x86_64-linux-0.16.0-dev.3059+42e33db9d/lib/std/Io/Uring.zig < patches/uring-networking.patch + # fetch dependencies first (cacheable — only changes when build.zig.zon changes) COPY build.zig build.zig.zon ./ RUN zig build --fetch-only 2>/dev/null || true diff --git a/build.zig.zon b/build.zig.zon index 3fce98c..5602e0a 100644 --- a/build.zig.zon +++ b/build.zig.zon @@ -5,12 +5,12 @@ .minimum_zig_version = "0.16.0", .dependencies = .{ .zat = .{ - .url = "https://tangled.org/zat.dev/zat/archive/v0.3.0-alpha.15.tar.gz", - .hash = "zat-0.3.0-alpha.15-5PuC7nVhBQCJEzz9LuzSbtLb68Wd0x_yjDgTP3EqV8dH", + .url = "https://tangled.org/zat.dev/zat/archive/v0.3.0-alpha.16.tar.gz", + .hash = "zat-0.3.0-alpha.15-5PuC7nVhBQC8QDphqwrmXGqng1Xvo8ua_H5MqS-smp1T", }, .websocket = .{ - .url = "https://github.com/zzstoatzz/websocket.zig/archive/ac3df25.tar.gz", - .hash = "websocket-0.1.0-ZPISdUvvAwDQN3W3AYDxmzMj5ipuTnB3vpQinQPF9LqI", + .url = "https://github.com/zzstoatzz/websocket.zig/archive/80c6434.tar.gz", + .hash = "websocket-0.1.0-ZPISdTHwAwA1d45BsYRE81Z8wNwZ3RhukgNADOma4eym", }, .pg = .{ .url = "git+https://github.com/zzstoatzz/pg.zig?ref=dev#5ce2355b1d851075523709c7d3068dcdb0224322", diff --git a/patches/uring-networking.patch b/patches/uring-networking.patch new file mode 100644 index 0000000..0eab1d1 --- /dev/null +++ b/patches/uring-networking.patch @@ -0,0 +1,530 @@ +--- a/lib/std/Io/Uring.zig ++++ b/lib/std/Io/Uring.zig +@@ -771,16 +771,16 @@ + .random = random, + .randomSecure = randomSecure, + +- .netListenIp = netListenIpUnavailable, +- .netAccept = netAcceptUnavailable, ++ .netListenIp = netListenIp, ++ .netAccept = netAccept, + .netBindIp = netBindIp, +- .netConnectIp = netConnectIpUnavailable, ++ .netConnectIp = netConnectIp, + .netListenUnix = netListenUnixUnavailable, + .netConnectUnix = netConnectUnixUnavailable, + .netSocketCreatePair = netSocketCreatePairUnavailable, +- .netSend = netSendUnavailable, +- .netRead = netReadUnavailable, +- .netWrite = netWriteUnavailable, ++ .netSend = netSend, ++ .netRead = netRead, ++ .netWrite = netWrite, + .netWriteFile = netWriteFileUnavailable, + .netClose = netClose, + .netShutdown = netShutdown, +@@ -4953,28 +4953,83 @@ + }; + } + +-fn netListenIpUnavailable( ++fn netListenIp( + userdata: ?*anyopaque, + address: *const net.IpAddress, + options: net.IpAddress.ListenOptions, + ) net.IpAddress.ListenError!net.Socket { + const ev: *Evented = @ptrCast(@alignCast(userdata)); +- _ = ev; +- _ = address; +- _ = options; +- return error.NetworkDown; ++ const family = posixAddressFamily(address); ++ var maybe_sync: CancelRegion.Sync.Maybe = .{ .cancel_region = .init() }; ++ defer maybe_sync.deinit(ev); ++ const socket_fd = try ev.socket(&maybe_sync.cancel_region, family, .{ ++ .mode = options.mode, ++ .protocol = options.protocol, ++ }); ++ errdefer ev.closeAsync(socket_fd); ++ if (options.reuse_address) { ++ try ev.setsockopt(&maybe_sync.cancel_region, socket_fd, linux.SOL.SOCKET, linux.SO.REUSEADDR, 1); ++ try ev.setsockopt(&maybe_sync.cancel_region, socket_fd, linux.SOL.SOCKET, linux.SO.REUSEPORT, 1); ++ } ++ var storage: PosixAddress = undefined; ++ var addr_len = addressToPosix(address, &storage); ++ try ev.bind(&maybe_sync.cancel_region, socket_fd, &storage.any, addr_len); ++ try ev.listen(&maybe_sync.cancel_region, socket_fd, options.kernel_backlog); ++ try ev.getsockname(try maybe_sync.enterSync(ev), socket_fd, &storage.any, &addr_len); ++ return .{ .handle = socket_fd, .address = addressFromPosix(&storage) }; + } + +-fn netAcceptUnavailable( ++fn netAccept( + userdata: ?*anyopaque, + listen_handle: net.Socket.Handle, + options: net.Server.AcceptOptions, + ) net.Server.AcceptError!net.Socket { +- const ev: *Evented = @ptrCast(@alignCast(userdata)); +- _ = ev; +- _ = listen_handle; + _ = options; +- return error.NetworkDown; ++ const ev: *Evented = @ptrCast(@alignCast(userdata)); ++ var cancel_region: CancelRegion = .init(); ++ defer cancel_region.deinit(); ++ var storage: PosixAddress = undefined; ++ var addr_len: linux.socklen_t = @sizeOf(PosixAddress); ++ const accepted_fd = while (true) { ++ const thread = try cancel_region.awaitIoUring(); ++ thread.enqueue().* = .{ ++ .opcode = .ACCEPT, ++ .flags = 0, ++ .ioprio = 0, ++ .fd = listen_handle, ++ .off = @intFromPtr(&addr_len), ++ .addr = @intFromPtr(&storage.any), ++ .len = 0, ++ .rw_flags = linux.SOCK.CLOEXEC, ++ .user_data = @intFromPtr(cancel_region.fiber), ++ .buf_index = 0, ++ .personality = 0, ++ .splice_fd_in = 0, ++ .addr3 = 0, ++ .resv = 0, ++ }; ++ ev.yield(null, .nothing); ++ const completion = cancel_region.completion(); ++ switch (completion.errno()) { ++ .SUCCESS => break completion.result, ++ .INTR, .CANCELED => {}, ++ .AGAIN => unreachable, ++ .BADF => |err| return errnoBug(err), ++ .CONNABORTED => return error.ConnectionAborted, ++ .INVAL => return error.SocketNotListening, ++ .NOTSOCK => |err| return errnoBug(err), ++ .MFILE => return error.ProcessFdQuotaExceeded, ++ .NFILE => return error.SystemFdQuotaExceeded, ++ .NOBUFS => return error.SystemResources, ++ .NOMEM => return error.SystemResources, ++ .OPNOTSUPP => |err| return errnoBug(err), ++ .PROTO => return error.ProtocolFailure, ++ .PERM => return error.BlockedByFirewall, ++ .NETDOWN => return error.NetworkDown, ++ else => |err| return unexpectedErrno(err), ++ } ++ }; ++ return .{ .handle = accepted_fd, .address = addressFromPosix(&storage) }; + } + + fn netBindIp( +@@ -4996,16 +5051,26 @@ + return .{ .handle = socket_fd, .address = addressFromPosix(&storage) }; + } + +-fn netConnectIpUnavailable( ++fn netConnectIp( + userdata: ?*anyopaque, + address: *const net.IpAddress, + options: net.IpAddress.ConnectOptions, + ) net.IpAddress.ConnectError!net.Socket { ++ if (options.timeout != .none) @panic("TODO: connect timeout for Io.Uring"); + const ev: *Evented = @ptrCast(@alignCast(userdata)); +- _ = ev; +- _ = address; +- _ = options; +- return error.NetworkDown; ++ const family = posixAddressFamily(address); ++ var maybe_sync: CancelRegion.Sync.Maybe = .{ .cancel_region = .init() }; ++ defer maybe_sync.deinit(ev); ++ const socket_fd = try ev.socket(&maybe_sync.cancel_region, family, .{ ++ .mode = options.mode, ++ .protocol = options.protocol, ++ }); ++ errdefer ev.closeAsync(socket_fd); ++ var storage: PosixAddress = undefined; ++ var addr_len = addressToPosix(address, &storage); ++ try ev.connect(&maybe_sync.cancel_region, socket_fd, &storage.any, addr_len); ++ try ev.getsockname(try maybe_sync.enterSync(ev), socket_fd, &storage.any, &addr_len); ++ return .{ .handle = socket_fd, .address = addressFromPosix(&storage) }; + } + + fn netListenUnixUnavailable( +@@ -5039,20 +5104,96 @@ + return error.OperationUnsupported; + } + +-fn netSendUnavailable( ++fn netSend( + userdata: ?*anyopaque, + handle: net.Socket.Handle, + messages: []net.OutgoingMessage, + flags: net.SendFlags, + ) struct { ?net.Socket.SendError, usize } { + const ev: *Evented = @ptrCast(@alignCast(userdata)); +- _ = ev; +- _ = handle; +- _ = messages; +- _ = flags; +- return .{ error.NetworkDown, 0 }; ++ const posix_flags: u32 = ++ @as(u32, if (flags.confirm) linux.MSG.CONFIRM else 0) | ++ @as(u32, if (flags.dont_route) linux.MSG.DONTROUTE else 0) | ++ @as(u32, if (flags.eor) linux.MSG.EOR else 0) | ++ @as(u32, if (flags.oob) linux.MSG.OOB else 0) | ++ @as(u32, if (flags.fastopen) linux.MSG.FASTOPEN else 0) | ++ linux.MSG.NOSIGNAL; ++ ++ for (messages, 0..) |*message, i| { ++ ev.netSendOne(handle, message, posix_flags) catch |err| return .{ err, i }; ++ } ++ return .{ null, messages.len }; + } + ++fn netSendOne( ++ ev: *Evented, ++ handle: net.Socket.Handle, ++ message: *net.OutgoingMessage, ++ flags: u32, ++) net.Socket.SendError!void { ++ var addr: PosixAddress = undefined; ++ var iov: iovec_const = .{ .base = @constCast(message.data_ptr), .len = message.data_len }; ++ var msg: linux.msghdr_const = .{ ++ .name = &addr.any, ++ .namelen = addressToPosix(message.address, &addr), ++ .iov = (&iov)[0..1], ++ .iovlen = 1, ++ .control = if (message.control.len == 0) null else @constCast(message.control.ptr), ++ .controllen = @intCast(message.control.len), ++ .flags = 0, ++ }; ++ var cancel_region: CancelRegion = .init(); ++ defer cancel_region.deinit(); ++ while (true) { ++ const thread = try cancel_region.awaitIoUring(); ++ thread.enqueue().* = .{ ++ .opcode = .SENDMSG, ++ .flags = 0, ++ .ioprio = 0, ++ .fd = handle, ++ .off = 0, ++ .addr = @intFromPtr(&msg), ++ .len = 0, ++ .rw_flags = flags, ++ .user_data = @intFromPtr(cancel_region.fiber), ++ .buf_index = 0, ++ .personality = 0, ++ .splice_fd_in = 0, ++ .addr3 = 0, ++ .resv = 0, ++ }; ++ ev.yield(null, .nothing); ++ const completion = cancel_region.completion(); ++ switch (completion.errno()) { ++ .SUCCESS => { ++ message.data_len = @intCast(completion.result); ++ return; ++ }, ++ .INTR, .CANCELED => {}, ++ .ACCES => return error.AccessDenied, ++ .ALREADY => return error.FastOpenAlreadyInProgress, ++ .BADF => |err| return errnoBug(err), ++ .CONNRESET => return error.ConnectionResetByPeer, ++ .DESTADDRREQ => |err| return errnoBug(err), ++ .FAULT => |err| return errnoBug(err), ++ .INVAL => |err| return errnoBug(err), ++ .ISCONN => |err| return errnoBug(err), ++ .MSGSIZE => return error.MessageOversize, ++ .NOBUFS => return error.SystemResources, ++ .NOMEM => return error.SystemResources, ++ .NOTSOCK => |err| return errnoBug(err), ++ .OPNOTSUPP => |err| return errnoBug(err), ++ .PIPE => return error.SocketUnconnected, ++ .AFNOSUPPORT => return error.AddressFamilyUnsupported, ++ .HOSTUNREACH => return error.HostUnreachable, ++ .NETUNREACH => return error.NetworkUnreachable, ++ .NOTCONN => return error.SocketUnconnected, ++ .NETDOWN => return error.NetworkDown, ++ else => |err| return unexpectedErrno(err), ++ } ++ } ++} ++ + fn netReceive( + ev: *Evented, + cancel_region: *CancelRegion, +@@ -5142,19 +5283,67 @@ + } + } + +-fn netReadUnavailable( ++fn netRead( + userdata: ?*anyopaque, + fd: net.Socket.Handle, + data: [][]u8, + ) net.Stream.Reader.Error!usize { + const ev: *Evented = @ptrCast(@alignCast(userdata)); +- _ = ev; +- _ = fd; +- _ = data; +- return error.NetworkDown; ++ var iovecs_buffer: [max_iovecs_len]iovec = undefined; ++ var i: usize = 0; ++ for (data) |buf| { ++ if (iovecs_buffer.len - i == 0) break; ++ if (buf.len > 0) { ++ iovecs_buffer[i] = .{ .base = buf.ptr, .len = buf.len }; ++ i += 1; ++ } ++ } ++ if (i == 0) return 0; ++ const dest = iovecs_buffer[0..i]; ++ assert(dest[0].len > 0); ++ var cancel_region: CancelRegion = .init(); ++ defer cancel_region.deinit(); ++ const gather = dest.len > 1 or dest[0].len > 0xfffff000; ++ while (true) { ++ const thread = try cancel_region.awaitIoUring(); ++ thread.enqueue().* = .{ ++ .opcode = if (gather) .READV else .READ, ++ .flags = 0, ++ .ioprio = 0, ++ .fd = fd, ++ .off = std.math.maxInt(u64), ++ .addr = if (gather) @intFromPtr(dest.ptr) else @intFromPtr(dest[0].base), ++ .len = @intCast(if (gather) dest.len else dest[0].len), ++ .rw_flags = 0, ++ .user_data = @intFromPtr(cancel_region.fiber), ++ .buf_index = 0, ++ .personality = 0, ++ .splice_fd_in = 0, ++ .addr3 = 0, ++ .resv = 0, ++ }; ++ ev.yield(null, .nothing); ++ const completion = cancel_region.completion(); ++ switch (completion.errno()) { ++ .SUCCESS => return @as(u32, @bitCast(completion.result)), ++ .INTR, .CANCELED => {}, ++ .INVAL => |err| return errnoBug(err), ++ .FAULT => |err| return errnoBug(err), ++ .AGAIN => unreachable, ++ .BADF => |err| return errnoBug(err), ++ .NOBUFS => return error.SystemResources, ++ .NOMEM => return error.SystemResources, ++ .NOTCONN => return error.SocketUnconnected, ++ .CONNRESET => return error.ConnectionResetByPeer, ++ .TIMEDOUT => return error.Timeout, ++ .PIPE => return error.SocketUnconnected, ++ .NETDOWN => return error.NetworkDown, ++ else => |err| return unexpectedErrno(err), ++ } ++ } + } + +-fn netWriteUnavailable( ++fn netWrite( + userdata: ?*anyopaque, + handle: net.Socket.Handle, + header: []const u8, +@@ -5162,12 +5351,94 @@ + splat: usize, + ) net.Stream.Writer.Error!usize { + const ev: *Evented = @ptrCast(@alignCast(userdata)); +- _ = ev; +- _ = handle; +- _ = header; +- _ = data; +- _ = splat; +- return error.NetworkDown; ++ var iovecs: [max_iovecs_len]iovec_const = undefined; ++ var iovlen: iovlen_t = 0; ++ addBuf(&iovecs, &iovlen, header); ++ for (data[0 .. data.len - 1]) |bytes| addBuf(&iovecs, &iovlen, bytes); ++ const pattern = data[data.len - 1]; ++ var backup_buffer: [splat_buffer_size]u8 = undefined; ++ if (iovecs.len - iovlen != 0) switch (splat) { ++ 0 => {}, ++ 1 => addBuf(&iovecs, &iovlen, pattern), ++ else => switch (pattern.len) { ++ 0 => {}, ++ 1 => { ++ const splat_buffer = &backup_buffer; ++ const memset_len = @min(splat_buffer.len, splat); ++ const buf = splat_buffer[0..memset_len]; ++ @memset(buf, pattern[0]); ++ addBuf(&iovecs, &iovlen, buf); ++ var remaining_splat = splat - buf.len; ++ while (remaining_splat > splat_buffer.len and iovecs.len - iovlen != 0) { ++ assert(buf.len == splat_buffer.len); ++ addBuf(&iovecs, &iovlen, splat_buffer); ++ remaining_splat -= splat_buffer.len; ++ } ++ addBuf(&iovecs, &iovlen, splat_buffer[0..@min(remaining_splat, splat_buffer.len)]); ++ }, ++ else => for (0..@min(splat, iovecs.len - iovlen)) |_| { ++ addBuf(&iovecs, &iovlen, pattern); ++ }, ++ }, ++ }; ++ const iov = iovecs[0..iovlen]; ++ if (iov.len == 0) return 0; ++ var msg: linux.msghdr_const = .{ ++ .name = null, ++ .namelen = 0, ++ .iov = iov.ptr, ++ .iovlen = iovlen, ++ .control = null, ++ .controllen = 0, ++ .flags = 0, ++ }; ++ var cancel_region: CancelRegion = .init(); ++ defer cancel_region.deinit(); ++ while (true) { ++ const thread = try cancel_region.awaitIoUring(); ++ thread.enqueue().* = .{ ++ .opcode = .SENDMSG, ++ .flags = 0, ++ .ioprio = 0, ++ .fd = handle, ++ .off = 0, ++ .addr = @intFromPtr(&msg), ++ .len = 0, ++ .rw_flags = linux.MSG.NOSIGNAL, ++ .user_data = @intFromPtr(cancel_region.fiber), ++ .buf_index = 0, ++ .personality = 0, ++ .splice_fd_in = 0, ++ .addr3 = 0, ++ .resv = 0, ++ }; ++ ev.yield(null, .nothing); ++ const completion = cancel_region.completion(); ++ switch (completion.errno()) { ++ .SUCCESS => return @as(u32, @bitCast(completion.result)), ++ .INTR, .CANCELED => {}, ++ .ACCES => |err| return errnoBug(err), ++ .AGAIN => unreachable, ++ .BADF => |err| return errnoBug(err), ++ .DESTADDRREQ => |err| return errnoBug(err), ++ .FAULT => |err| return errnoBug(err), ++ .INVAL => |err| return errnoBug(err), ++ .ISCONN => |err| return errnoBug(err), ++ .MSGSIZE => |err| return errnoBug(err), ++ .OPNOTSUPP => |err| return errnoBug(err), ++ .ALREADY => return error.FastOpenAlreadyInProgress, ++ .CONNRESET => return error.ConnectionResetByPeer, ++ .NOBUFS => return error.SystemResources, ++ .NOMEM => return error.SystemResources, ++ .PIPE => return error.SocketUnconnected, ++ .NOTCONN => return error.SocketUnconnected, ++ .AFNOSUPPORT => return error.AddressFamilyUnsupported, ++ .HOSTUNREACH => return error.HostUnreachable, ++ .NETUNREACH => return error.NetworkUnreachable, ++ .NETDOWN => return error.NetworkDown, ++ else => |err| return unexpectedErrno(err), ++ } ++ } + } + + fn netWriteFileUnavailable( +@@ -5303,11 +5574,100 @@ + .ADDRNOTAVAIL => return error.AddressUnavailable, + .FAULT => |err| return errnoBug(err), // invalid `addr` pointer + .NOMEM => return error.SystemResources, ++ else => |err| return unexpectedErrno(err), ++ } ++ } ++} ++ ++fn listen( ++ ev: *Evented, ++ cancel_region: *CancelRegion, ++ socket_fd: fd_t, ++ backlog: u31, ++) !void { ++ while (true) { ++ const thread = try cancel_region.awaitIoUring(); ++ thread.enqueue().* = .{ ++ .opcode = .LISTEN, ++ .flags = 0, ++ .ioprio = 0, ++ .fd = socket_fd, ++ .off = 0, ++ .addr = 0, ++ .len = backlog, ++ .rw_flags = 0, ++ .user_data = @intFromPtr(cancel_region.fiber), ++ .buf_index = 0, ++ .personality = 0, ++ .splice_fd_in = 0, ++ .addr3 = 0, ++ .resv = 0, ++ }; ++ ev.yield(null, .nothing); ++ switch (cancel_region.errno()) { ++ .SUCCESS => return, ++ .INTR, .CANCELED => {}, ++ .ADDRINUSE => return error.AddressInUse, ++ .BADF => |err| return errnoBug(err), ++ .NOTSOCK => |err| return errnoBug(err), ++ .OPNOTSUPP => |err| return errnoBug(err), + else => |err| return unexpectedErrno(err), + } + } + } + ++fn connect( ++ ev: *Evented, ++ cancel_region: *CancelRegion, ++ socket_fd: fd_t, ++ addr: *const linux.sockaddr, ++ addr_len: linux.socklen_t, ++) !void { ++ while (true) { ++ const thread = try cancel_region.awaitIoUring(); ++ thread.enqueue().* = .{ ++ .opcode = .CONNECT, ++ .flags = 0, ++ .ioprio = 0, ++ .fd = socket_fd, ++ .off = addr_len, ++ .addr = @intFromPtr(addr), ++ .len = 0, ++ .rw_flags = 0, ++ .user_data = @intFromPtr(cancel_region.fiber), ++ .buf_index = 0, ++ .personality = 0, ++ .splice_fd_in = 0, ++ .addr3 = 0, ++ .resv = 0, ++ }; ++ ev.yield(null, .nothing); ++ switch (cancel_region.errno()) { ++ .SUCCESS => return, ++ .INTR, .CANCELED => {}, ++ .ADDRNOTAVAIL => return error.AddressUnavailable, ++ .AFNOSUPPORT => return error.AddressFamilyUnsupported, ++ .AGAIN, .INPROGRESS => return, // non-blocking / TCP fast open ++ .ALREADY => return error.ConnectionPending, ++ .BADF => |err| return errnoBug(err), ++ .CONNREFUSED => return error.ConnectionRefused, ++ .CONNRESET => return error.ConnectionResetByPeer, ++ .FAULT => |err| return errnoBug(err), ++ .ISCONN => |err| return errnoBug(err), ++ .HOSTUNREACH => return error.HostUnreachable, ++ .NETUNREACH => return error.NetworkUnreachable, ++ .NOTSOCK => |err| return errnoBug(err), ++ .PROTOTYPE => |err| return errnoBug(err), ++ .TIMEDOUT => return error.Timeout, ++ .CONNABORTED => |err| return errnoBug(err), ++ .ACCES => return error.AccessDenied, ++ .PERM => |err| return errnoBug(err), ++ .NETDOWN => return error.NetworkDown, ++ else => |err| return unexpectedErrno(err), ++ } ++ } ++} ++ + fn chdir(sync: *CancelRegion.Sync, path: [*:0]const u8) ChdirError!void { + while (true) { + try sync.cancel_region.await(.nothing); diff --git a/src/main.zig b/src/main.zig index 5e9cd65..39b373c 100644 --- a/src/main.zig +++ b/src/main.zig @@ -47,13 +47,13 @@ const log = std.log.scoped(.relay); pub const default_stack_size = 8 * 1024 * 1024; // -- Io backend selection -- -// Evented (fibers): network orchestration, subscriber connections, WS server, broadcasting. -// Worker threads use a dedicated pool_io (Threaded) for their sync — they never touch Evented io. -// The broadcast queue bridges workers → broadcaster fiber (atomics only, no Io dependency). +// Io.Evented (fibers on io_uring): ~35 threads instead of ~2,800 (one per PDS). +// Networking via patched Uring.zig (patches/uring-networking.patch) — implements +// listen, accept, connect, read, write, send via io_uring opcodes. DNS (netLookup) +// is NOT patched; subscribers resolve hostnames through pool_io (Threaded) instead. // -// NOTE: Io.Uring has a fiber context-switch GPF under ReleaseSafe (optimizer + safety checks). -// Debug and ReleaseFast both work fine. Build with ReleaseFast until the stdlib bug is fixed. -// See repro_evented.zig for the minimal reproduction case. +// Known issue: Io.Uring GPFs under ReleaseSafe (optimizer + safety interaction). +// Build with ReleaseFast. See repro_evented.zig for minimal reproduction. const Backend = Io.Evented; var backend: Backend = undefined; diff --git a/src/subscriber.zig b/src/subscriber.zig index 55e7612..563b9ff 100644 --- a/src/subscriber.zig +++ b/src/subscriber.zig @@ -323,7 +323,13 @@ pub const Subscriber = struct { } const path = w.buffered(); - var client = try websocket.Client.init(self.io, self.allocator, .{ + // DNS + TCP connect through pool_io (Threaded — has working netLookup). + // The resulting fd is used by the Evented io for TLS + WebSocket I/O. + const dns_io = self.pool_io orelse self.io; + const host_name = try Io.net.HostName.init(self.options.hostname); + const net_stream = try host_name.connect(dns_io, 443, .{ .mode = .stream }); + + var client = try websocket.Client.initWithStream(self.io, self.allocator, net_stream, .{ .host = self.options.hostname, .port = 443, .tls = true,