--- 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,102 @@ }; } -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); + // bind + listen: sync syscalls (IORING_OP_BIND/LISTEN require kernel 6.11+) + switch (linux.errno(linux.bind(socket_fd, &storage.any, addr_len))) { + .SUCCESS => {}, + .ADDRINUSE => return error.AddressInUse, + .AFNOSUPPORT => return error.AddressFamilyUnsupported, + .ADDRNOTAVAIL => return error.AddressUnavailable, + .BADF => |err| return errnoBug(err), + .INVAL => |err| return errnoBug(err), + .NOTSOCK => |err| return errnoBug(err), + .FAULT => |err| return errnoBug(err), + .NOMEM => return error.SystemResources, + else => |err| return unexpectedErrno(err), + } + switch (linux.errno(linux.listen(socket_fd, options.kernel_backlog))) { + .SUCCESS => {}, + .ADDRINUSE => return error.AddressInUse, + .BADF => |err| return errnoBug(err), + .NOTSOCK => |err| return errnoBug(err), + .OPNOTSUPP => |err| return errnoBug(err), + else => |err| return unexpectedErrno(err), + } + 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 +5070,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,18 +5123,94 @@ 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( @@ -5142,19 +5302,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 +5370,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( @@ -5306,6 +5596,58 @@ 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, + .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 {