authorgravatar for bblack@wikimedia.orgBrandon Black <bblack@wikimedia.org> 2026-05-05 12:19:01-05:00
committergravatar for alex@alexrp.comAlex Rønne Petersen <alex@alexrp.com> 2026-08-10 20:33:08+02:00
loga6c2839505c9214ca6cd8baaa52e67fcf99c9e2e
tree1973e8eda849262648691c89c613de93f1b371ed
parentb9adc659b367bef744b79f9e966c0300f863407f

Move std.Io netSend to an Operation

Also adds sendTimeout() and sendManyTimeout() functions to Io.net.Socket, and updates the UDP tests to exercise all of the combined interfaces. Deprecates `sendMany` in favor of `sendManyTimeout` with a `Timeout` of `.none`. The existing `sendMany` interface is not capable of reporting partial success, and eventually removing it will align with the current interfaces on the receive side. Seems to work for Posix + Threaded? Beyond very basic "avoid compile errors maybe", I have not tried to address non-Posix or non-Threaded parts.

6 files changed, 301 insertions(+), 108 deletions(-)

lib/std/Io.zig+45-10
...@@ -235,7 +235,6 @@ pub const VTable = struct {...@@ -235,7 +235,6 @@ pub const VTable = struct {
235 netListenUnix: *const fn (?*anyopaque, *const net.UnixAddress, net.UnixAddress.ListenOptions) net.UnixAddress.ListenError!net.Socket.Handle,235 netListenUnix: *const fn (?*anyopaque, *const net.UnixAddress, net.UnixAddress.ListenOptions) net.UnixAddress.ListenError!net.Socket.Handle,
236 netConnectUnix: *const fn (?*anyopaque, *const net.UnixAddress) net.UnixAddress.ConnectError!net.Socket.Handle,236 netConnectUnix: *const fn (?*anyopaque, *const net.UnixAddress) net.UnixAddress.ConnectError!net.Socket.Handle,
237 netSocketCreatePair: *const fn (?*anyopaque, net.Socket.CreatePairOptions) net.Socket.CreatePairError![2]net.Socket,237 netSocketCreatePair: *const fn (?*anyopaque, net.Socket.CreatePairOptions) net.Socket.CreatePairError![2]net.Socket,
238 netSend: *const fn (?*anyopaque, net.Socket.Handle, []net.OutgoingMessage, net.SendFlags) struct { ?net.Socket.SendError, usize },
239 netWrite: *const fn (?*anyopaque, dest: net.Socket.Handle, header: []const u8, data: []const []const u8, splat: usize) net.Stream.Writer.Error!usize,238 netWrite: *const fn (?*anyopaque, dest: net.Socket.Handle, header: []const u8, data: []const []const u8, splat: usize) net.Stream.Writer.Error!usize,
240 netWriteFile: *const fn (?*anyopaque, net.Socket.Handle, header: []const u8, *Io.File.Reader, Io.Limit) net.Stream.Writer.WriteFileError!usize,239 netWriteFile: *const fn (?*anyopaque, net.Socket.Handle, header: []const u8, *Io.File.Reader, Io.Limit) net.Stream.Writer.WriteFileError!usize,
241 netClose: *const fn (?*anyopaque, sockets: []const net.Socket) void,240 netClose: *const fn (?*anyopaque, sockets: []const net.Socket) void,
...@@ -252,6 +251,7 @@ pub const Operation = union(enum) {...@@ -252,6 +251,7 @@ pub const Operation = union(enum) {
252 /// other systems this tag is unreachable.251 /// other systems this tag is unreachable.
253 device_io_control: DeviceIoControl,252 device_io_control: DeviceIoControl,
254 net_receive: NetReceive,253 net_receive: NetReceive,
254 net_send: NetSend,
255 net_read: NetRead,255 net_read: NetRead,
256256
257 pub const Tag = @typeInfo(Operation).@"union".tag_type.?;257 pub const Tag = @typeInfo(Operation).@"union".tag_type.?;
...@@ -378,6 +378,49 @@ pub const Operation = union(enum) {...@@ -378,6 +378,49 @@ pub const Operation = union(enum) {
378 pub const Result = struct { ?net.Socket.ReceiveError, usize };378 pub const Result = struct { ?net.Socket.ReceiveError, usize };
379 };379 };
380380
381 pub const NetSend = struct {
382 socket_handle: net.Socket.Handle,
383 messages: []net.OutgoingMessage,
384 flags: net.SendFlags,
385
386 pub const Error = error{
387 /// The socket type requires that message be sent atomically, and the
388 /// size of the message to be sent made this impossible. The message
389 /// was not transmitted, or was partially transmitted.
390 MessageOversize,
391 /// The output queue for a network interface was full. This generally indicates that the
392 /// interface has stopped sending, but may be caused by transient congestion. (Normally,
393 /// this does not occur in Linux. Packets are just silently dropped when a device queue
394 /// overflows.)
395 ///
396 /// This is also caused when there is not enough kernel memory available.
397 SystemResources,
398 /// No route to network.
399 NetworkUnreachable,
400 /// Network reached but no route to host.
401 HostUnreachable,
402 /// The local network interface used to reach the destination is offline.
403 NetworkDown,
404 /// The destination address is not listening. Can still occur for
405 /// connectionless messages.
406 ConnectionRefused,
407 /// Operating system or protocol does not support the address family.
408 AddressFamilyUnsupported,
409 /// Another TCP Fast Open is already in progress.
410 FastOpenAlreadyInProgress,
411 /// Network session was unexpectedly closed by recipient.
412 ConnectionResetByPeer,
413 /// Local end has been shut down on a connection-oriented socket, or
414 /// the socket was never connected.
415 SocketUnconnected,
416 /// An attempt was made to send to a network/broadcast address as
417 /// though it was a unicast address.
418 AccessDenied,
419 } || Io.UnexpectedError;
420
421 pub const Result = struct { ?net.Socket.SendError, usize };
422 };
423
381 pub const NetRead = struct {424 pub const NetRead = struct {
382 socket_handle: net.Socket.Handle,425 socket_handle: net.Socket.Handle,
383 data: [][]u8,426 data: [][]u8,
...@@ -2712,7 +2755,6 @@ pub const failing: std.Io = .{...@@ -2712,7 +2755,6 @@ pub const failing: std.Io = .{
2712 .netListenUnix = failingNetListenUnix,2755 .netListenUnix = failingNetListenUnix,
2713 .netConnectUnix = failingNetConnectUnix,2756 .netConnectUnix = failingNetConnectUnix,
2714 .netSocketCreatePair = failingNetSocketCreatePair,2757 .netSocketCreatePair = failingNetSocketCreatePair,
2715 .netSend = failingNetSend,
2716 .netWrite = failingNetWrite,2758 .netWrite = failingNetWrite,
2717 .netWriteFile = failingNetWriteFile,2759 .netWriteFile = failingNetWriteFile,
2718 .netClose = unreachableNetClose,2760 .netClose = unreachableNetClose,
...@@ -2860,6 +2902,7 @@ pub fn failingOperate(userdata: ?*anyopaque, operation: Operation) Cancelable!Op...@@ -2860,6 +2902,7 @@ pub fn failingOperate(userdata: ?*anyopaque, operation: Operation) Cancelable!Op
2860 .file_write_streaming => .{ .file_write_streaming = error.InputOutput },2902 .file_write_streaming => .{ .file_write_streaming = error.InputOutput },
2861 .device_io_control => unreachable,2903 .device_io_control => unreachable,
2862 .net_receive => .{ .net_receive = .{ error.NetworkDown, 0 } },2904 .net_receive => .{ .net_receive = .{ error.NetworkDown, 0 } },
2905 .net_send => .{ .net_send = .{ error.NetworkDown, 0 } },
2863 .net_read => .{ .net_read = error.NetworkDown },2906 .net_read => .{ .net_read = error.NetworkDown },
2864 };2907 };
2865}2908}
...@@ -3456,14 +3499,6 @@ pub fn failingNetSocketCreatePair(userdata: ?*anyopaque, options: net.Socket.Cre...@@ -3456,14 +3499,6 @@ pub fn failingNetSocketCreatePair(userdata: ?*anyopaque, options: net.Socket.Cre
3456 return error.OperationUnsupported;3499 return error.OperationUnsupported;
3457}3500}
34583501
3459pub fn failingNetSend(userdata: ?*anyopaque, handle: net.Socket.Handle, messages: []net.OutgoingMessage, flags: net.SendFlags) struct { ?net.Socket.SendError, usize } {
3460 _ = userdata;
3461 _ = handle;
3462 _ = messages;
3463 _ = flags;
3464 return .{ error.NetworkDown, 0 };
3465}
3466
3467pub fn failingNetWrite(userdata: ?*anyopaque, dest: net.Socket.Handle, header: []const u8, data: []const []const u8, splat: usize) net.Stream.Writer.Error!usize {3502pub fn failingNetWrite(userdata: ?*anyopaque, dest: net.Socket.Handle, header: []const u8, data: []const []const u8, splat: usize) net.Stream.Writer.Error!usize {
3468 _ = userdata;3503 _ = userdata;
3469 _ = dest;3504 _ = dest;
lib/std/Io/Dispatch.zig+1-15
...@@ -458,7 +458,6 @@ pub fn io(ev: *Evented) Io {...@@ -458,7 +458,6 @@ pub fn io(ev: *Evented) Io {
458 .netListenUnix = netListenUnixUnavailable,458 .netListenUnix = netListenUnixUnavailable,
459 .netConnectUnix = netConnectUnixUnavailable,459 .netConnectUnix = netConnectUnixUnavailable,
460 .netSocketCreatePair = netSocketCreatePairUnavailable,460 .netSocketCreatePair = netSocketCreatePairUnavailable,
461 .netSend = netSendUnavailable,
462 .netWrite = netWriteUnavailable,461 .netWrite = netWriteUnavailable,
463 .netWriteFile = netWriteFileUnavailable,462 .netWriteFile = netWriteFileUnavailable,
464 .netClose = netClose,463 .netClose = netClose,
...@@ -1713,6 +1712,7 @@ fn operate(userdata: ?*anyopaque, operation: Io.Operation) Io.Cancelable!Io.Oper...@@ -1713,6 +1712,7 @@ fn operate(userdata: ?*anyopaque, operation: Io.Operation) Io.Cancelable!Io.Oper
1713 },1712 },
1714 .device_io_control => |*o| return .{ .device_io_control = try deviceIoControl(o) },1713 .device_io_control => |*o| return .{ .device_io_control = try deviceIoControl(o) },
1715 .net_receive => @panic("TODO implement net_receive operation"),1714 .net_receive => @panic("TODO implement net_receive operation"),
1715 .net_send => @panic("TODO implement net_send operation"),
1716 .net_read => @panic("TODO implement net_read operation"),1716 .net_read => @panic("TODO implement net_read operation"),
1717 }1717 }
1718}1718}
...@@ -4866,20 +4866,6 @@ fn netSocketCreatePairUnavailable(...@@ -4866,20 +4866,6 @@ fn netSocketCreatePairUnavailable(
4866 return error.OperationUnsupported;4866 return error.OperationUnsupported;
4867}4867}
48684868
4869fn netSendUnavailable(
4870 userdata: ?*anyopaque,
4871 handle: net.Socket.Handle,
4872 messages: []net.OutgoingMessage,
4873 flags: net.SendFlags,
4874) struct { ?net.Socket.SendError, usize } {
4875 const ev: *Evented = @ptrCast(@alignCast(userdata));
4876 _ = ev;
4877 _ = handle;
4878 _ = messages;
4879 _ = flags;
4880 return .{ error.NetworkDown, 0 };
4881}
4882
4883fn netWriteUnavailable(4869fn netWriteUnavailable(
4884 userdata: ?*anyopaque,4870 userdata: ?*anyopaque,
4885 handle: net.Socket.Handle,4871 handle: net.Socket.Handle,
lib/std/Io/Threaded.zig+75-14
...@@ -1951,10 +1951,6 @@ pub fn io(t: *Threaded) Io {...@@ -1951,10 +1951,6 @@ pub fn io(t: *Threaded) Io {
1951 else => netWritePosix,1951 else => netWritePosix,
1952 },1952 },
1953 .netWriteFile = netWriteFile,1953 .netWriteFile = netWriteFile,
1954 .netSend = switch (native_os) {
1955 .windows => netSendWindows,
1956 else => netSendPosix,
1957 },
1958 .netInterfaceNameResolve = netInterfaceNameResolve,1954 .netInterfaceNameResolve = netInterfaceNameResolve,
1959 .netInterfaceName = netInterfaceName,1955 .netInterfaceName = netInterfaceName,
1960 .netLookup = netLookup,1956 .netLookup = netLookup,
...@@ -2565,6 +2561,25 @@ fn operate(userdata: ?*anyopaque, operation: Io.Operation) Io.Cancelable!Io.Oper...@@ -2565,6 +2561,25 @@ fn operate(userdata: ?*anyopaque, operation: Io.Operation) Io.Cancelable!Io.Oper
2565 };2561 };
2566 break :o .{ null, 1 };2562 break :o .{ null, 1 };
2567 } },2563 } },
2564 .net_send => |*o| return .{
2565 .net_send = o: {
2566 if (!have_networking) break :o .{ error.NetworkDown, 0 };
2567 if (is_windows) break :o netSendWindows(t, o.socket_handle, o.messages, o.flags);
2568 const send_err, const sent = netSendPosix(t, o.socket_handle, o.messages, o.flags, false);
2569 if (send_err) |err| switch (err) {
2570 error.Canceled => |e| if (sent == 0) {
2571 return e;
2572 } else {
2573 // Leave the `error.Canceled` for later, but don't try to send any more messages.
2574 recancelInner();
2575 break :o .{ null, sent };
2576 },
2577 error.WouldBlock => unreachable,
2578 else => |e| break :o .{ e, sent },
2579 };
2580 break :o .{ null, sent };
2581 },
2582 },
2568 .net_read => |o| return .{2583 .net_read => |o| return .{
2569 .net_read = netRead(o.socket_handle, o.data) catch |err| switch (err) {2584 .net_read = netRead(o.socket_handle, o.data) catch |err| switch (err) {
2570 error.Canceled => |e| return e,2585 error.Canceled => |e| return e,
...@@ -2626,6 +2641,14 @@ fn batchAwaitAsync(userdata: ?*anyopaque, b: *Io.Batch) Io.Cancelable!void {...@@ -2626,6 +2641,14 @@ fn batchAwaitAsync(userdata: ?*anyopaque, b: *Io.Batch) Io.Cancelable!void {
2626 };2641 };
2627 poll_len += 1;2642 poll_len += 1;
2628 },2643 },
2644 .net_send => |*o| {
2645 poll_buffer[poll_len] = .{
2646 .fd = o.socket_handle,
2647 .events = posix.POLL.OUT | posix.POLL.ERR,
2648 .revents = 0,
2649 };
2650 poll_len += 1;
2651 },
2629 .net_read => |o| {2652 .net_read => |o| {
2630 poll_buffer[poll_len] = .{2653 poll_buffer[poll_len] = .{
2631 .fd = o.socket_handle,2654 .fd = o.socket_handle,
...@@ -2811,6 +2834,35 @@ fn batchAwaitConcurrent(userdata: ?*anyopaque, b: *Io.Batch, timeout: Io.Timeout...@@ -2811,6 +2834,35 @@ fn batchAwaitConcurrent(userdata: ?*anyopaque, b: *Io.Batch, timeout: Io.Timeout
2811 storage.* = .{ .completion = .{ .node = .{ .next = .none }, .result = result } };2834 storage.* = .{ .completion = .{ .node = .{ .next = .none }, .result = result } };
2812 b.completed.tail = index;2835 b.completed.tail = index;
2813 },2836 },
2837 .net_send => |*o| nb: {
2838 const result: Io.Operation.Result = .{
2839 .net_send = o: {
2840 const send_err, const sent = netSendPosix(t, o.socket_handle, o.messages, o.flags, true);
2841 if (send_err) |err| switch (err) {
2842 error.Canceled => |e| if (sent == 0) {
2843 return e;
2844 } else {
2845 // Leave the `error.Canceled` for later, but don't try to send any more messages.
2846 recancelInner();
2847 break :o .{ null, sent };
2848 },
2849 error.WouldBlock => {
2850 if (sent != 0) break :o .{ null, sent };
2851 try poll_storage.add(o.socket_handle, posix.POLL.OUT | posix.POLL.ERR);
2852 break :nb;
2853 },
2854 else => |e| break :o .{ e, sent },
2855 };
2856 break :o .{ null, sent };
2857 },
2858 };
2859 switch (b.completed.tail) {
2860 .none => b.completed.head = index,
2861 else => |tail_index| b.storage[tail_index.toIndex()].completion.node.next = index,
2862 }
2863 storage.* = .{ .completion = .{ .node = .{ .next = .none }, .result = result } };
2864 b.completed.tail = index;
2865 },
2814 .net_read => |o| try poll_storage.add(o.socket_handle, posix.POLL.IN | posix.POLL.ERR),2866 .net_read => |o| try poll_storage.add(o.socket_handle, posix.POLL.IN | posix.POLL.ERR),
2815 }2867 }
2816 index = submission.node.next;2868 index = submission.node.next;
...@@ -3007,6 +3059,7 @@ fn batchApc(...@@ -3007,6 +3059,7 @@ fn batchApc(
3007 .file_write_streaming => .{ .file_write_streaming = ntWriteFileResult(iosb) },3059 .file_write_streaming => .{ .file_write_streaming = ntWriteFileResult(iosb) },
3008 .device_io_control => .{ .device_io_control = iosb.* },3060 .device_io_control => .{ .device_io_control = iosb.* },
3009 .net_receive => unreachable,3061 .net_receive => unreachable,
3062 .net_send => unreachable,
3010 .net_read => unreachable,3063 .net_read => unreachable,
3011 };3064 };
3012 storage.* = .{ .completion = .{ .node = .{ .next = .none }, .result = result } };3065 storage.* = .{ .completion = .{ .node = .{ .next = .none }, .result = result } };
...@@ -3216,6 +3269,13 @@ fn batchDrainSubmittedWindows(t: *Threaded, b: *Io.Batch, concurrency: bool) (Io...@@ -3216,6 +3269,13 @@ fn batchDrainSubmittedWindows(t: *Threaded, b: *Io.Batch, concurrency: bool) (Io
3216 .net_receive = netReceiveWindows(t, o.socket_handle, o.message_buffer, o.data_buffer, o.flags),3269 .net_receive = netReceiveWindows(t, o.socket_handle, o.message_buffer, o.data_buffer, o.flags),
3217 });3270 });
3218 },3271 },
3272 .net_send => |*o| {
3273 // TODO integrate with overlapped I/O or equivalent to avoid this error
3274 if (concurrency) return error.ConcurrencyUnavailable;
3275 batchCompleteBlockingWindows(b, operation_userdata, .{
3276 .net_send = netSendWindows(t, o.socket_handle, o.messages, o.flags),
3277 });
3278 },
3219 .net_read => |*o| {3279 .net_read => |*o| {
3220 // TODO integrate with overlapped I/O or equivalent to avoid this error3280 // TODO integrate with overlapped I/O or equivalent to avoid this error
3221 if (concurrency) return error.ConcurrencyUnavailable;3281 if (concurrency) return error.ConcurrencyUnavailable;
...@@ -12866,13 +12926,13 @@ fn netReadWindows(socket_handle: net.Socket.Handle, data: [][]u8) net.Stream.Rea...@@ -12866,13 +12926,13 @@ fn netReadWindows(socket_handle: net.Socket.Handle, data: [][]u8) net.Stream.Rea
12866}12926}
1286712927
12868fn netSendPosix(12928fn netSendPosix(
12869 userdata: ?*anyopaque,12929 t: *Threaded,
12870 socket_handle: net.Socket.Handle,12930 socket_handle: net.Socket.Handle,
12871 messages: []net.OutgoingMessage,12931 messages: []net.OutgoingMessage,
12872 flags: net.SendFlags,12932 flags: net.SendFlags,
12873) struct { ?net.Socket.SendError, usize } {12933 nonblocking: bool,
12934) struct { ?(net.Socket.SendError || error{WouldBlock}), usize } {
12874 if (!have_networking) return .{ error.NetworkDown, 0 };12935 if (!have_networking) return .{ error.NetworkDown, 0 };
12875 const t: *Threaded = @ptrCast(@alignCast(userdata));
1287612936
12877 const posix_flags: u32 =12937 const posix_flags: u32 =
12878 @as(u32, if (@hasDecl(posix.MSG, "CONFIRM") and flags.confirm) posix.MSG.CONFIRM else 0) |12938 @as(u32, if (@hasDecl(posix.MSG, "CONFIRM") and flags.confirm) posix.MSG.CONFIRM else 0) |
...@@ -12880,6 +12940,7 @@ fn netSendPosix(...@@ -12880,6 +12940,7 @@ fn netSendPosix(
12880 @as(u32, if (@hasDecl(posix.MSG, "EOR") and flags.eor) posix.MSG.EOR else 0) |12940 @as(u32, if (@hasDecl(posix.MSG, "EOR") and flags.eor) posix.MSG.EOR else 0) |
12881 @as(u32, if (@hasDecl(posix.MSG, "OOB") and flags.oob) posix.MSG.OOB else 0) |12941 @as(u32, if (@hasDecl(posix.MSG, "OOB") and flags.oob) posix.MSG.OOB else 0) |
12882 @as(u32, if (@hasDecl(posix.MSG, "FASTOPEN") and flags.fastopen) posix.MSG.FASTOPEN else 0) |12942 @as(u32, if (@hasDecl(posix.MSG, "FASTOPEN") and flags.fastopen) posix.MSG.FASTOPEN else 0) |
12943 @as(u32, if (@hasDecl(posix.MSG, "DONTWAIT") and nonblocking) posix.MSG.DONTWAIT else 0) |
12883 posix.MSG.NOSIGNAL;12944 posix.MSG.NOSIGNAL;
1288412945
12885 var i: usize = 0;12946 var i: usize = 0;
...@@ -12895,13 +12956,12 @@ fn netSendPosix(...@@ -12895,13 +12956,12 @@ fn netSendPosix(
12895}12956}
1289612957
12897fn netSendWindows(12958fn netSendWindows(
12898 userdata: ?*anyopaque,12959 t: *Threaded,
12899 socket_handle: net.Socket.Handle,12960 socket_handle: net.Socket.Handle,
12900 messages: []net.OutgoingMessage,12961 messages: []net.OutgoingMessage,
12901 flags: net.SendFlags,12962 flags: net.SendFlags,
12902) struct { ?net.Socket.SendError, usize } {12963) struct { ?net.Socket.SendError, usize } {
12903 if (!have_networking) return .{ error.NetworkDown, 0 };12964 if (!have_networking) return .{ error.NetworkDown, 0 };
12904 const t: *Threaded = @ptrCast(@alignCast(userdata));
12905 for (messages, 0..) |*m, i| {12965 for (messages, 0..) |*m, i| {
12906 t.netSendOneWindows(socket_handle, m, flags) catch |err| return .{ err, i };12966 t.netSendOneWindows(socket_handle, m, flags) catch |err| return .{ err, i };
12907 }12967 }
...@@ -12953,7 +13013,7 @@ fn netSendOnePosix(...@@ -12953,7 +13013,7 @@ fn netSendOnePosix(
12953 socket_handle: net.Socket.Handle,13013 socket_handle: net.Socket.Handle,
12954 message: *net.OutgoingMessage,13014 message: *net.OutgoingMessage,
12955 flags: u32,13015 flags: u32,
12956) net.Socket.SendError!void {13016) (net.Socket.SendError || error{WouldBlock})!void {
12957 _ = t;13017 _ = t;
12958 var addr: PosixAddress = undefined;13018 var addr: PosixAddress = undefined;
12959 var iovec: posix.iovec_const = .{ .base = @constCast(message.data_ptr), .len = message.data_len };13019 var iovec: posix.iovec_const = .{ .base = @constCast(message.data_ptr), .len = message.data_len };
...@@ -12981,6 +13041,7 @@ fn netSendOnePosix(...@@ -12981,6 +13041,7 @@ fn netSendOnePosix(
12981 continue;13041 continue;
12982 },13042 },
12983 .ACCES => return syscall.fail(error.AccessDenied),13043 .ACCES => return syscall.fail(error.AccessDenied),
13044 .AGAIN => return syscall.fail(error.WouldBlock),
12984 .ALREADY => return syscall.fail(error.FastOpenAlreadyInProgress),13045 .ALREADY => return syscall.fail(error.FastOpenAlreadyInProgress),
12985 .CONNRESET => return syscall.fail(error.ConnectionResetByPeer),13046 .CONNRESET => return syscall.fail(error.ConnectionResetByPeer),
12986 .MSGSIZE => return syscall.fail(error.MessageOversize),13047 .MSGSIZE => return syscall.fail(error.MessageOversize),
...@@ -13008,7 +13069,7 @@ fn netSendManyPosix(...@@ -13008,7 +13069,7 @@ fn netSendManyPosix(
13008 socket_handle: net.Socket.Handle,13069 socket_handle: net.Socket.Handle,
13009 messages: []net.OutgoingMessage,13070 messages: []net.OutgoingMessage,
13010 flags: u32,13071 flags: u32,
13011) net.Socket.SendError!usize {13072) (net.Socket.SendError || error{WouldBlock})!usize {
13012 var msg_buffer: [64]posix.system.mmsghdr = undefined;13073 var msg_buffer: [64]posix.system.mmsghdr = undefined;
13013 var addr_buffer: [msg_buffer.len]PosixAddress = undefined;13074 var addr_buffer: [msg_buffer.len]PosixAddress = undefined;
13014 var iovecs_buffer: [msg_buffer.len]posix.iovec = undefined;13075 var iovecs_buffer: [msg_buffer.len]posix.iovec = undefined;
...@@ -13051,6 +13112,7 @@ fn netSendManyPosix(...@@ -13051,6 +13112,7 @@ fn netSendManyPosix(
13051 continue;13112 continue;
13052 },13113 },
13053 .ACCES => return syscall.fail(error.AccessDenied),13114 .ACCES => return syscall.fail(error.AccessDenied),
13115 .AGAIN => return syscall.fail(error.WouldBlock),
13054 .ALREADY => return syscall.fail(error.FastOpenAlreadyInProgress),13116 .ALREADY => return syscall.fail(error.FastOpenAlreadyInProgress),
13055 .CONNRESET => return syscall.fail(error.ConnectionResetByPeer),13117 .CONNRESET => return syscall.fail(error.ConnectionResetByPeer),
13056 .MSGSIZE => return syscall.fail(error.MessageOversize),13118 .MSGSIZE => return syscall.fail(error.MessageOversize),
...@@ -13063,7 +13125,6 @@ fn netSendManyPosix(...@@ -13063,7 +13125,6 @@ fn netSendManyPosix(
13063 .NOTCONN => return syscall.fail(error.SocketUnconnected),13125 .NOTCONN => return syscall.fail(error.SocketUnconnected),
13064 .NETDOWN => return syscall.fail(error.NetworkDown),13126 .NETDOWN => return syscall.fail(error.NetworkDown),
1306513127
13066 .AGAIN => |err| return syscall.errnoBug(err),
13067 .BADF => |err| return syscall.errnoBug(err), // File descriptor used after closed.13128 .BADF => |err| return syscall.errnoBug(err), // File descriptor used after closed.
13068 .DESTADDRREQ => |err| return syscall.errnoBug(err), // The socket is not connection-mode, and no peer address is set.13129 .DESTADDRREQ => |err| return syscall.errnoBug(err), // The socket is not connection-mode, and no peer address is set.
13069 .FAULT => |err| return syscall.errnoBug(err), // An invalid user space address was specified for an argument.13130 .FAULT => |err| return syscall.errnoBug(err), // An invalid user space address was specified for an argument.
...@@ -14609,7 +14670,7 @@ fn lookupDns(...@@ -14609,7 +14670,7 @@ fn lookupDns(
14609 message_i += 1;14670 message_i += 1;
14610 }14671 }
14611 }14672 }
14612 _ = netSendPosix(t, socket.handle, message_buffer[0..message_i], .{});14673 _ = netSendPosix(t, socket.handle, message_buffer[0..message_i], .{}, false);
14613 }14674 }
1461414675
14615 const timeout: Io.Timeout = .{ .deadline = .{14676 const timeout: Io.Timeout = .{ .deadline = .{
...@@ -14657,7 +14718,7 @@ fn lookupDns(...@@ -14657,7 +14718,7 @@ fn lookupDns(
14657 .data_ptr = query.ptr,14718 .data_ptr = query.ptr,
14658 .data_len = query.len,14719 .data_len = query.len,
14659 };14720 };
14660 _ = netSendPosix(t, socket.handle, (&retry_message)[0..1], .{});14721 _ = netSendPosix(t, socket.handle, (&retry_message)[0..1], .{}, false);
14661 continue;14722 continue;
14662 },14723 },
14663 else => continue,14724 else => continue,
lib/std/Io/Uring.zig+11-15
...@@ -778,7 +778,6 @@ pub fn io(ev: *Evented) Io {...@@ -778,7 +778,6 @@ pub fn io(ev: *Evented) Io {
778 .netListenUnix = netListenUnixUnavailable,778 .netListenUnix = netListenUnixUnavailable,
779 .netConnectUnix = netConnectUnixUnavailable,779 .netConnectUnix = netConnectUnixUnavailable,
780 .netSocketCreatePair = netSocketCreatePairUnavailable,780 .netSocketCreatePair = netSocketCreatePairUnavailable,
781 .netSend = netSendUnavailable,
782 .netWrite = netWriteUnavailable,781 .netWrite = netWriteUnavailable,
783 .netWriteFile = netWriteFileUnavailable,782 .netWriteFile = netWriteFileUnavailable,
784 .netClose = netClose,783 .netClose = netClose,
...@@ -2107,6 +2106,12 @@ fn operate(userdata: ?*anyopaque, operation: Io.Operation) Io.Cancelable!Io.Oper...@@ -2107,6 +2106,12 @@ fn operate(userdata: ?*anyopaque, operation: Io.Operation) Io.Cancelable!Io.Oper
2107 };2106 };
2108 },2107 },
2109 },2108 },
2109 .net_send => |o| .{
2110 .net_send = r: {
2111 _ = o;
2112 break :r .{ error.NetworkDown, 0 }; // TODO
2113 },
2114 },
2110 .net_read => |o| .{2115 .net_read => |o| .{
2111 .net_read = r: {2116 .net_read = r: {
2112 _ = o;2117 _ = o;
...@@ -2400,6 +2405,10 @@ fn batchDrainSubmitted(...@@ -2400,6 +2405,10 @@ fn batchDrainSubmitted(
2400 _ = o;2405 _ = o;
2401 @panic("TODO implement batchDrainSubmitted for net_receive");2406 @panic("TODO implement batchDrainSubmitted for net_receive");
2402 },2407 },
2408 .net_send => |o| {
2409 _ = o;
2410 @panic("TODO implement batchDrainSubmitted for net_send");
2411 },
2403 .net_read => |o| {2412 .net_read => |o| {
2404 _ = o;2413 _ = o;
2405 @panic("TODO implement batchDrainSubmitted for net_read");2414 @panic("TODO implement batchDrainSubmitted for net_read");
...@@ -2505,6 +2514,7 @@ fn batchDrainReady(batch: *Io.Batch) Io.Timeout.Error!void {...@@ -2505,6 +2514,7 @@ fn batchDrainReady(batch: *Io.Batch) Io.Timeout.Error!void {
2505 },2514 },
2506 .device_io_control => unreachable,2515 .device_io_control => unreachable,
2507 .net_receive => @panic("TODO"),2516 .net_receive => @panic("TODO"),
2517 .net_send => @panic("TODO"),
2508 .net_read => @panic("TODO"),2518 .net_read => @panic("TODO"),
2509 })) |result| {2519 })) |result| {
2510 switch (batch.completed.tail) {2520 switch (batch.completed.tail) {
...@@ -5054,20 +5064,6 @@ fn netSocketCreatePairUnavailable(...@@ -5054,20 +5064,6 @@ fn netSocketCreatePairUnavailable(
5054 return error.OperationUnsupported;5064 return error.OperationUnsupported;
5055}5065}
50565066
5057fn netSendUnavailable(
5058 userdata: ?*anyopaque,
5059 handle: net.Socket.Handle,
5060 messages: []net.OutgoingMessage,
5061 flags: net.SendFlags,
5062) struct { ?net.Socket.SendError, usize } {
5063 const ev: *Evented = @ptrCast(@alignCast(userdata));
5064 _ = ev;
5065 _ = handle;
5066 _ = messages;
5067 _ = flags;
5068 return .{ error.NetworkDown, 0 };
5069}
5070
5071fn netReceive(5067fn netReceive(
5072 ev: *Evented,5068 ev: *Evented,
5073 cancel_region: *CancelRegion,5069 cancel_region: *CancelRegion,
lib/std/Io/net.zig+59-38
...@@ -1088,52 +1088,73 @@ pub const Socket = struct {...@@ -1088,52 +1088,73 @@ pub const Socket = struct {
1088 io.vtable.netClose(io.userdata, sockets);1088 io.vtable.netClose(io.userdata, sockets);
1089 }1089 }
10901090
1091 pub const SendError = error{1091 pub const SendError = Io.Operation.NetSend.Error || Io.Cancelable;
1092 /// The socket type requires that message be sent atomically, and the
1093 /// size of the message to be sent made this impossible. The message
1094 /// was not transmitted, or was partially transmitted.
1095 MessageOversize,
1096 /// The output queue for a network interface was full. This generally indicates that the
1097 /// interface has stopped sending, but may be caused by transient congestion. (Normally,
1098 /// this does not occur in Linux. Packets are just silently dropped when a device queue
1099 /// overflows.)
1100 ///
1101 /// This is also caused when there is not enough kernel memory available.
1102 SystemResources,
1103 /// No route to network.
1104 NetworkUnreachable,
1105 /// Network reached but no route to host.
1106 HostUnreachable,
1107 /// The local network interface used to reach the destination is offline.
1108 NetworkDown,
1109 /// The destination address is not listening. Can still occur for
1110 /// connectionless messages.
1111 ConnectionRefused,
1112 /// Operating system or protocol does not support the address family.
1113 AddressFamilyUnsupported,
1114 /// Another TCP Fast Open is already in progress.
1115 FastOpenAlreadyInProgress,
1116 /// Network session was unexpectedly closed by recipient.
1117 ConnectionResetByPeer,
1118 /// Local end has been shut down on a connection-oriented socket, or
1119 /// the socket was never connected.
1120 SocketUnconnected,
1121 /// An attempt was made to send to a network/broadcast address as
1122 /// though it was a unicast address.
1123 AccessDenied,
1124 } || Io.UnexpectedError || Io.Cancelable;
11251092
1126 /// Transfers `data` to `dest`, connectionless, in one packet.1093 /// Transfers `data` to `dest`, connectionless, in one packet.
1127 pub fn send(s: *const Socket, io: Io, dest: *const IpAddress, data: []const u8) SendError!void {1094 pub fn send(s: *const Socket, io: Io, dest: *const IpAddress, data: []const u8) SendError!void {
1128 var message: OutgoingMessage = .{ .address = dest, .data_ptr = data.ptr, .data_len = data.len };1095 var message: OutgoingMessage = .{ .address = dest, .data_ptr = data.ptr, .data_len = data.len };
1129 const err, const n = io.vtable.netSend(io.userdata, s.handle, (&message)[0..1], .{});1096 const maybe_err, const count = (try io.operate(.{ .net_send = .{
1130 if (n != 1) return err.?;1097 .socket_handle = s.handle,
1098 .messages = (&message)[0..1],
1099 .flags = .{},
1100 } })).net_send;
1101 if (maybe_err) |err| {
1102 assert(count == 0);
1103 return err;
1104 } else {
1105 assert(count == 1);
1106 }
1131 if (message.data_len != data.len) return error.MessageOversize;1107 if (message.data_len != data.len) return error.MessageOversize;
1132 }1108 }
11331109
1110 pub const SendTimeoutError = SendError || Io.Timeout.Error || Io.ConcurrentError;
1111
1112 pub fn sendTimeout(
1113 s: *const Socket,
1114 io: Io,
1115 dest: *const IpAddress,
1116 data: []const u8,
1117 timeout: Io.Timeout,
1118 ) SendTimeoutError!void {
1119 var message: OutgoingMessage = .{ .address = dest, .data_ptr = data.ptr, .data_len = data.len };
1120 const maybe_err, const count = (try io.operateTimeout(.{ .net_send = .{
1121 .socket_handle = s.handle,
1122 .messages = (&message)[0..1],
1123 .flags = .{},
1124 } }, timeout)).net_send;
1125 if (maybe_err) |err| return err;
1126 assert(1 == count);
1127 if (message.data_len != data.len) return error.MessageOversize;
1128 }
1129
1130 /// Deprecated; use `sendManyTimeout` with a timeout of `.none`.
1131 ///
1132 /// If this function returns an error, some (but not all) of `messages` may
1133 /// still have been sent. This condition is not reported by this function,
1134 /// but is reported by `sendManyTimeout`.
1134 pub fn sendMany(s: *const Socket, io: Io, messages: []OutgoingMessage, flags: SendFlags) SendError!void {1135 pub fn sendMany(s: *const Socket, io: Io, messages: []OutgoingMessage, flags: SendFlags) SendError!void {
1135 const err, const n = io.vtable.netSend(io.userdata, s.handle, messages, flags);1136 const result = try io.operate(.{ .net_send = .{
1136 if (n != messages.len) return err.?;1137 .socket_handle = s.handle,
1138 .messages = messages,
1139 .flags = flags,
1140 } });
1141 const maybe_send_err, _ = result.net_send;
1142 return maybe_send_err orelse {};
1143 }
1144
1145 pub fn sendManyTimeout(
1146 s: *const Socket,
1147 io: Io,
1148 messages: []OutgoingMessage,
1149 flags: SendFlags,
1150 timeout: Io.Timeout,
1151 ) struct { ?SendTimeoutError, usize } {
1152 const result = io.operateTimeout(.{ .net_send = .{
1153 .socket_handle = s.handle,
1154 .messages = messages,
1155 .flags = flags,
1156 } }, timeout) catch |err| return .{ err, 0 };
1157 return result.net_send;
1137 }1158 }
11381159
1139 pub const ReceiveError = Io.Operation.NetReceive.Error || Io.Cancelable;1160 pub const ReceiveError = Io.Operation.NetReceive.Error || Io.Cancelable;
lib/std/Io/net/test.zig+110-16
...@@ -394,7 +394,7 @@ test "UDP send and receive" {...@@ -394,7 +394,7 @@ test "UDP send and receive" {
394 try testing.expectEqualStrings(&send_data, received.data);394 try testing.expectEqualStrings(&send_data, received.data);
395}395}
396396
397test "UDP send and receiveTimeout" {397test "UDP sendTimeout and receiveTimeout" {
398 const io = testing.io;398 const io = testing.io;
399 const localhost: net.IpAddress = .{ .ip4 = .loopback(0) };399 const localhost: net.IpAddress = .{ .ip4 = .loopback(0) };
400400
...@@ -406,15 +406,16 @@ test "UDP send and receiveTimeout" {...@@ -406,15 +406,16 @@ test "UDP send and receiveTimeout" {
406 const send_sock = try localhost.bind(io, .{ .mode = .dgram });406 const send_sock = try localhost.bind(io, .{ .mode = .dgram });
407 defer send_sock.close(io);407 defer send_sock.close(io);
408408
409 const send_data: [3]u8 = .{ '1', '2', '3' };409 const six_hours: Io.Timeout = .{ .duration = .{ .clock = .awake, .raw = .fromSeconds(21600) } };
410 try send_sock.send(io, &recv_sock.address, &send_data);
411410
412 const timeo: Io.Timeout = .{ .duration = .{ .clock = .awake, .raw = .fromSeconds(10) } };411 const send_data: [3]u8 = .{ '1', '2', '3' };
413 var recv_buf: [4]u8 = undefined;412 send_sock.sendTimeout(io, &recv_sock.address, &send_data, six_hours) catch |err| switch (err) {
414 const received = recv_sock.receiveTimeout(io, &recv_buf, timeo) catch |err| switch (err) {
415 error.ConcurrencyUnavailable => return error.SkipZigTest,413 error.ConcurrencyUnavailable => return error.SkipZigTest,
416 else => |e| return e,414 else => |e| return e,
417 };415 };
416
417 var recv_buf: [4]u8 = undefined;
418 const received = try recv_sock.receiveTimeout(io, &recv_buf, six_hours);
418 try testing.expect(received.from.eql(&send_sock.address));419 try testing.expect(received.from.eql(&send_sock.address));
419 try testing.expectEqualStrings(&send_data, received.data);420 try testing.expectEqualStrings(&send_data, received.data);
420421
...@@ -422,7 +423,7 @@ test "UDP send and receiveTimeout" {...@@ -422,7 +423,7 @@ test "UDP send and receiveTimeout" {
422 try testing.expectError(error.Timeout, recv_sock.receiveTimeout(io, &recv_buf, short));423 try testing.expectError(error.Timeout, recv_sock.receiveTimeout(io, &recv_buf, short));
423}424}
424425
425test "UDP sendMany 1 recvManyTimeout 2" {426test "UDP sendMany 2 and receive 2" {
426 const io = testing.io;427 const io = testing.io;
427 const localhost: net.IpAddress = .{ .ip4 = .loopback(0) };428 const localhost: net.IpAddress = .{ .ip4 = .loopback(0) };
428429
...@@ -434,26 +435,119 @@ test "UDP sendMany 1 recvManyTimeout 2" {...@@ -434,26 +435,119 @@ test "UDP sendMany 1 recvManyTimeout 2" {
434 const send_sock = try localhost.bind(io, .{ .mode = .dgram });435 const send_sock = try localhost.bind(io, .{ .mode = .dgram });
435 defer send_sock.close(io);436 defer send_sock.close(io);
436437
438 const send_data: [3]u8 = .{ '1', '2', '3' };
439 var send_msgs: [2]Io.net.OutgoingMessage = @splat(.{
440 .address = &recv_sock.address,
441 .data_ptr = &send_data,
442 .data_len = 3,
443 });
444 // note sendMany is deprecated, but should remain tested until removed
445 try send_sock.sendMany(io, &send_msgs, .{});
446 try testing.expectEqual(3, send_msgs[0].data_len);
447 try testing.expectEqual(3, send_msgs[1].data_len);
448
449 var recv_buf: [4]u8 = undefined;
450
451 {
452 const first = try recv_sock.receive(io, &recv_buf);
453 try testing.expect(first.from.eql(&send_sock.address));
454 try testing.expectEqualStrings(&send_data, first.data);
455 }
456
457 {
458 const second = try recv_sock.receive(io, &recv_buf);
459 try testing.expect(second.from.eql(&send_sock.address));
460 try testing.expectEqualStrings(&send_data, second.data);
461 }
462}
463
464test "UDP sendManyTimeout 1 recvManyTimeout 2" {
465 const io = testing.io;
466 const localhost: net.IpAddress = .{ .ip4 = .loopback(0) };
467
468 const recv_sock = localhost.bind(io, .{ .mode = .dgram }) catch |err| switch (err) {
469 error.NetworkDown => return error.SkipZigTest,
470 else => |e| return e,
471 };
472 defer recv_sock.close(io);
473 const send_sock = try localhost.bind(io, .{ .mode = .dgram });
474 defer send_sock.close(io);
475
476 const six_hours: Io.Timeout = .{ .duration = .{ .clock = .awake, .raw = .fromSeconds(21600) } };
437 const send_data: [3]u8 = .{ '1', '2', '3' };477 const send_data: [3]u8 = .{ '1', '2', '3' };
438 var send_msg: Io.net.OutgoingMessage = .{478 var send_msg: Io.net.OutgoingMessage = .{
439 .address = &recv_sock.address,479 .address = &recv_sock.address,
440 .data_ptr = &send_data,480 .data_ptr = &send_data,
441 .data_len = 3,481 .data_len = 3,
442 };482 };
443 try send_sock.sendMany(io, (&send_msg)[0..1], .{});
444 try testing.expectEqual(3, send_msg.data_len);
445483
446 // This should not wait 10 seconds for the absent second message, it should484 const maybe_send_err, const send_count = send_sock.sendManyTimeout(io, (&send_msg)[0..1], .{}, six_hours);
447 // complete as soon as the first one arrives485 if (maybe_send_err) |err| switch (err) {
448 var recv_msgs: [2]net.IncomingMessage = @splat(.init);
449 var recv_buf: [10]u8 = undefined;
450 const timeo: Io.Timeout = .{ .duration = .{ .clock = .awake, .raw = .fromSeconds(10) } };
451 const maybe_recv_err, const recv_count = recv_sock.receiveManyTimeout(io, &recv_msgs, &recv_buf, .{}, timeo);
452 if (maybe_recv_err) |err| switch (err) {
453 error.ConcurrencyUnavailable => return error.SkipZigTest,486 error.ConcurrencyUnavailable => return error.SkipZigTest,
454 else => |e| return e,487 else => |e| return e,
455 };488 };
489 try testing.expectEqual(1, send_count);
490 try testing.expectEqual(3, send_msg.data_len);
491
492 // This should complete as soon as the first message arrives, and not stall
493 // for the timeout waiting on the second one.
494 var recv_msgs: [2]net.IncomingMessage = @splat(.init);
495 var recv_buf: [10]u8 = undefined;
496 const maybe_recv_err, const recv_count = recv_sock.receiveManyTimeout(io, &recv_msgs, &recv_buf, .{}, six_hours);
497 if (maybe_recv_err) |err| return err;
456 try testing.expectEqual(1, recv_count);498 try testing.expectEqual(1, recv_count);
457 try testing.expect(recv_msgs[0].from.eql(&send_sock.address));499 try testing.expect(recv_msgs[0].from.eql(&send_sock.address));
458 try testing.expectEqualStrings(&send_data, recv_msgs[0].data);500 try testing.expectEqualStrings(&send_data, recv_msgs[0].data);
459}501}
502
503fn testUdpSender(io: Io, send_sock: Io.net.Socket, send_data: []const u8, dest: Io.net.IpAddress) !void {
504 try io.sleep(.fromMilliseconds(10), .boot);
505 try send_sock.send(io, &dest, send_data);
506 try send_sock.send(io, &dest, send_data);
507 try io.sleep(.fromMilliseconds(10), .boot);
508 try send_sock.send(io, &dest, send_data);
509}
510
511test "UDP concurrency and timeouts" {
512 const io = testing.io;
513 const localhost: net.IpAddress = .{ .ip4 = .loopback(0) };
514
515 const recv_sock = localhost.bind(io, .{ .mode = .dgram }) catch |err| switch (err) {
516 error.NetworkDown => return error.SkipZigTest,
517 else => |e| return e,
518 };
519 defer recv_sock.close(io);
520 const send_sock = try localhost.bind(io, .{ .mode = .dgram });
521 defer send_sock.close(io);
522
523 const send_data: [3]u8 = .{ '1', '2', '3' };
524 var sender = io.async(testUdpSender, .{ io, send_sock, &send_data, recv_sock.address });
525 defer sender.cancel(io) catch {};
526
527 // Because the sender is async (may execute serially or concurrently) and
528 // it has some 10ms timing gaps, it's likely that there will be a variety
529 // of random behaviors (total iterations, messages per iteration) on the
530 // receive end of things related to the target and runner conditions. This
531 // should still suceed so long as all 3 packets arrive in reasonable time.
532 const six_hours: Io.Timeout = .{ .duration = .{ .clock = .awake, .raw = .fromSeconds(21600) } };
533 var received: usize = 0;
534 for (0..3) |_| {
535 var recv_msgs: [3]net.IncomingMessage = @splat(.init);
536 var recv_buf: [9]u8 = undefined;
537 const maybe_recv_err, const recv_count = recv_sock.receiveManyTimeout(io, &recv_msgs, &recv_buf, .{}, six_hours);
538 if (maybe_recv_err) |err| switch (err) {
539 error.ConcurrencyUnavailable => return error.SkipZigTest,
540 else => |e| return e,
541 };
542 received += recv_count;
543 try testing.expect(received <= 3);
544 for (0..recv_count) |i| {
545 const msg = recv_msgs[i];
546 try testing.expect(msg.from.eql(&send_sock.address));
547 try testing.expectEqualStrings(&send_data, msg.data);
548 }
549 if (received == 3) break;
550 }
551 try testing.expectEqual(3, received);
552 try sender.await(io); // ensure sender didn't fail
553}