authorgravatar for lukas@lalinsky.comLukas Lalinsky <lukas@lalinsky.com> 2026-05-12 11:50:22+02:00
committergravatar for mlugg@mlugg.co.ukMatthew Lugg <mlugg@mlugg.co.uk> 2026-08-11 11:57:57+02:00
log9d92d2ce93c52083fe405716a6cf05151384b76e
tree13c48aa35738f4604330037d64ec4f66c4adad06
parent254c23b03875387798b8ce779223091c5f6ee444

Migrate `netWrite` to `std.Io.Operation`

Preparation for adding timeouts: https://codeberg.org/ziglang/zig/issues/32166

6 files changed, 88 insertions(+), 102 deletions(-)

lib/std/Io.zig+39-11
......@@ -235,7 +235,6 @@ pub const VTable = struct {
235235 netListenUnix: *const fn (?*anyopaque, *const net.UnixAddress, net.UnixAddress.ListenOptions) net.UnixAddress.ListenError!net.Socket.Handle,
236236 netConnectUnix: *const fn (?*anyopaque, *const net.UnixAddress) net.UnixAddress.ConnectError!net.Socket.Handle,
237237 netSocketCreatePair: *const fn (?*anyopaque, net.Socket.CreatePairOptions) net.Socket.CreatePairError![2]net.Socket,
238 netWrite: *const fn (?*anyopaque, dest: net.Socket.Handle, header: []const u8, data: []const []const u8, splat: usize) net.Stream.Writer.Error!usize,
239238 netWriteFile: *const fn (?*anyopaque, net.Socket.Handle, header: []const u8, *Io.File.Reader, Io.Limit) net.Stream.Writer.WriteFileError!usize,
240239 netClose: *const fn (?*anyopaque, sockets: []const net.Socket) void,
241240 netShutdown: *const fn (?*anyopaque, handle: net.Socket.Handle, how: net.ShutdownHow) net.ShutdownError!void,
......@@ -253,6 +252,7 @@ pub const Operation = union(enum) {
253252 net_receive: NetReceive,
254253 net_send: NetSend,
255254 net_read: NetRead,
255 net_write: NetWrite,
256256
257257 pub const Tag = @typeInfo(Operation).@"union".tag_type.?;
258258
......@@ -438,6 +438,43 @@ pub const Operation = union(enum) {
438438 pub const Result = Error!usize;
439439 };
440440
441 pub const NetWrite = struct {
442 socket_handle: net.Socket.Handle,
443 header: []const u8 = &.{},
444 data: []const []const u8,
445 splat: usize = 1,
446
447 pub const Error = error{
448 /// Another TCP Fast Open is already in progress.
449 FastOpenAlreadyInProgress,
450 /// Network session was unexpectedly closed by recipient.
451 ConnectionResetByPeer,
452 /// The output queue for a network interface was full. This generally indicates that the
453 /// interface has stopped sending, but may be caused by transient congestion. (Normally,
454 /// this does not occur in Linux. Packets are just silently dropped when a device queue
455 /// overflows.)
456 ///
457 /// This is also caused when there is not enough kernel memory available.
458 SystemResources,
459 /// No route to network.
460 NetworkUnreachable,
461 /// Network reached but no route to host.
462 HostUnreachable,
463 /// The local network interface used to reach the destination is down.
464 NetworkDown,
465 /// The destination address is not listening.
466 ConnectionRefused,
467 /// The passed address didn't have the correct address family in its sa_family field.
468 AddressFamilyUnsupported,
469 /// Local end has been shut down on a connection-oriented socket, or
470 /// the socket was never connected.
471 SocketUnconnected,
472 SocketNotBound,
473 } || Io.UnexpectedError;
474
475 pub const Result = Error!usize;
476 };
477
441478 pub const Result = Result: {
442479 const operation_info = @typeInfo(Operation).@"union";
443480 const operation_count = operation_info.field_names.len;
......@@ -2771,7 +2808,6 @@ pub const failing: std.Io = .{
27712808 .netListenUnix = failingNetListenUnix,
27722809 .netConnectUnix = failingNetConnectUnix,
27732810 .netSocketCreatePair = failingNetSocketCreatePair,
2774 .netWrite = failingNetWrite,
27752811 .netWriteFile = failingNetWriteFile,
27762812 .netClose = unreachableNetClose,
27772813 .netShutdown = failingNetShutdown,
......@@ -2920,6 +2956,7 @@ pub fn failingOperate(userdata: ?*anyopaque, operation: Operation) Cancelable!Op
29202956 .net_receive => .{ .net_receive = .{ error.NetworkDown, 0 } },
29212957 .net_send => .{ .net_send = .{ error.NetworkDown, 0 } },
29222958 .net_read => .{ .net_read = error.NetworkDown },
2959 .net_write => .{ .net_write = error.NetworkDown },
29232960 };
29242961}
29252962
......@@ -3515,15 +3552,6 @@ pub fn failingNetSocketCreatePair(userdata: ?*anyopaque, options: net.Socket.Cre
35153552 return error.OperationUnsupported;
35163553}
35173554
3518pub fn failingNetWrite(userdata: ?*anyopaque, dest: net.Socket.Handle, header: []const u8, data: []const []const u8, splat: usize) net.Stream.Writer.Error!usize {
3519 _ = userdata;
3520 _ = dest;
3521 _ = header;
3522 _ = data;
3523 _ = splat;
3524 return error.NetworkDown;
3525}
3526
35273555pub fn failingNetWriteFile(userdata: ?*anyopaque, handle: net.Socket.Handle, header: []const u8, file_reader: *Io.File.Reader, limit: Io.Limit) net.Stream.Writer.WriteFileError!usize {
35283556 _ = userdata;
35293557 _ = handle;
lib/std/Io/Dispatch.zig+3-17
......@@ -458,7 +458,6 @@ pub fn io(ev: *Evented) Io {
458458 .netListenUnix = netListenUnixUnavailable,
459459 .netConnectUnix = netConnectUnixUnavailable,
460460 .netSocketCreatePair = netSocketCreatePairUnavailable,
461 .netWrite = netWriteUnavailable,
462461 .netWriteFile = netWriteFileUnavailable,
463462 .netClose = netClose,
464463 .netShutdown = netShutdownUnavailable,
......@@ -1714,6 +1713,7 @@ fn operate(userdata: ?*anyopaque, operation: Io.Operation) Io.Cancelable!Io.Oper
17141713 .net_receive => @panic("TODO implement net_receive operation"),
17151714 .net_send => @panic("TODO implement net_send operation"),
17161715 .net_read => @panic("TODO implement net_read operation"),
1716 .net_write => @panic("TODO implement net_write operation"),
17171717 }
17181718}
17191719
......@@ -2136,6 +2136,7 @@ fn batchDrainSubmitted(
21362136 .device_io_control => {},
21372137 .net_receive => @panic("TODO implement batched net_receive"),
21382138 .net_read => @panic("TODO implement batched net_read"),
2139 .net_write => @panic("TODO implement batched net_write"),
21392140 };
21402141 if (concurrency) return error.ConcurrencyUnavailable;
21412142 break :result try operate(ev, storage.submission.operation);
......@@ -2196,6 +2197,7 @@ fn batchSourceEvent(context: ?*anyopaque) callconv(.c) void {
21962197 .device_io_control => unreachable,
21972198 .net_receive => @panic("TODO implement batched net_receive"),
21982199 .net_read => @panic("TODO implement batched net_read"),
2200 .net_write => @panic("TODO implement batched net_write"),
21992201 };
22002202
22012203 switch (pending.node.prev) {
......@@ -4866,22 +4868,6 @@ fn netSocketCreatePairUnavailable(
48664868 return error.OperationUnsupported;
48674869}
48684870
4869fn netWriteUnavailable(
4870 userdata: ?*anyopaque,
4871 handle: net.Socket.Handle,
4872 header: []const u8,
4873 data: []const []const u8,
4874 splat: usize,
4875) net.Stream.Writer.Error!usize {
4876 const ev: *Evented = @ptrCast(@alignCast(userdata));
4877 _ = ev;
4878 _ = handle;
4879 _ = header;
4880 _ = data;
4881 _ = splat;
4882 return error.NetworkDown;
4883}
4884
48854871fn netWriteFileUnavailable(
48864872 userdata: ?*anyopaque,
48874873 socket_handle: net.Socket.Handle,
lib/std/Io/Kqueue.zig-19
......@@ -650,10 +650,8 @@ pub fn io(k: *Kqueue) Io {
650650 .netBindIp = netBindIp,
651651 .netConnectIp = netConnectIp,
652652 .netConnectUnix = netConnectUnix,
653 .netClose = netClose,
654653 .netShutdown = netShutdown,
655654 .netRead = netRead,
656 .netWrite = netWrite,
657655 .netSend = netSend,
658656 .netReceive = netReceive,
659657 .netInterfaceNameResolve = netInterfaceNameResolve,
......@@ -1270,23 +1268,6 @@ fn netRead(userdata: ?*anyopaque, fd: net.Socket.Handle, data: [][]u8) net.Strea
12701268 }
12711269}
12721270
1273fn netWrite(userdata: ?*anyopaque, dest: net.Socket.Handle, header: []const u8, data: []const []const u8, splat: usize) net.Stream.Writer.Error!usize {
1274 const k: *Kqueue = @ptrCast(@alignCast(userdata));
1275 _ = k;
1276 _ = dest;
1277 _ = header;
1278 _ = data;
1279 _ = splat;
1280 @panic("TODO");
1281}
1282
1283fn netClose(userdata: ?*anyopaque, sockets: []const net.Socket) void {
1284 const k: *Kqueue = @ptrCast(@alignCast(userdata));
1285 _ = k;
1286 _ = sockets;
1287 @panic("TODO");
1288}
1289
12901271fn netShutdown(userdata: ?*anyopaque, handle: net.Socket.Handle, how: net.ShutdownHow) net.ShutdownError!void {
12911272 const k: *Kqueue = @ptrCast(@alignCast(userdata));
12921273 _ = k;
lib/std/Io/Threaded.zig+29-10
......@@ -1946,10 +1946,6 @@ pub fn io(t: *Threaded) Io {
19461946 .windows => netShutdownWindows,
19471947 else => netShutdownPosix,
19481948 },
1949 .netWrite = switch (native_os) {
1950 .windows => netWriteWindows,
1951 else => netWritePosix,
1952 },
19531949 .netWriteFile = netWriteFile,
19541950 .netInterfaceNameResolve = netInterfaceNameResolve,
19551951 .netInterfaceName = netInterfaceName,
......@@ -2586,6 +2582,15 @@ fn operate(userdata: ?*anyopaque, operation: Io.Operation) Io.Cancelable!Io.Oper
25862582 else => |e| e,
25872583 },
25882584 },
2585 .net_write => |o| return .{
2586 .net_write = (if (is_windows)
2587 netWriteWindows(o.socket_handle, o.header, o.data, o.splat)
2588 else
2589 netWritePosix(o.socket_handle, o.header, o.data, o.splat)) catch |err| switch (err) {
2590 error.Canceled => |e| return e,
2591 else => |e| e,
2592 },
2593 },
25892594 }
25902595}
25912596
......@@ -2657,6 +2662,14 @@ fn batchAwaitAsync(userdata: ?*anyopaque, b: *Io.Batch) Io.Cancelable!void {
26572662 };
26582663 poll_len += 1;
26592664 },
2665 .net_write => |o| {
2666 poll_buffer[poll_len] = .{
2667 .fd = o.socket_handle,
2668 .events = posix.POLL.OUT | posix.POLL.ERR,
2669 .revents = 0,
2670 };
2671 poll_len += 1;
2672 },
26602673 }
26612674 index = submission.node.next;
26622675 }
......@@ -2864,6 +2877,7 @@ fn batchAwaitConcurrent(userdata: ?*anyopaque, b: *Io.Batch, timeout: Io.Timeout
28642877 b.completed.tail = index;
28652878 },
28662879 .net_read => |o| try poll_storage.add(o.socket_handle, posix.POLL.IN | posix.POLL.ERR),
2880 .net_write => |o| try poll_storage.add(o.socket_handle, posix.POLL.OUT | posix.POLL.ERR),
28672881 }
28682882 index = submission.node.next;
28692883 }
......@@ -3061,6 +3075,7 @@ fn batchApc(
30613075 .net_receive => unreachable,
30623076 .net_send => unreachable,
30633077 .net_read => unreachable,
3078 .net_write => unreachable,
30643079 };
30653080 storage.* = .{ .completion = .{ .node = .{ .next = .none }, .result = result } };
30663081 },
......@@ -3286,6 +3301,16 @@ fn batchDrainSubmittedWindows(t: *Threaded, b: *Io.Batch, concurrency: bool) (Io
32863301 },
32873302 });
32883303 },
3304 .net_write => |*o| {
3305 // TODO integrate with overlapped I/O or equivalent to avoid this error
3306 if (concurrency) return error.ConcurrencyUnavailable;
3307 batchCompleteBlockingWindows(b, operation_userdata, .{
3308 .net_write = netWriteWindows(o.socket_handle, o.header, o.data, o.splat) catch |err| switch (err) {
3309 error.Canceled => |e| return e,
3310 else => |e| e,
3311 },
3312 });
3313 },
32893314 }
32903315 index = submission.node.next;
32913316 }
......@@ -13283,15 +13308,12 @@ fn netReceiveOneWindows(
1328313308}
1328413309
1328513310fn netWritePosix(
13286 userdata: ?*anyopaque,
1328713311 fd: net.Socket.Handle,
1328813312 header: []const u8,
1328913313 data: []const []const u8,
1329013314 splat: usize,
1329113315) net.Stream.Writer.Error!usize {
1329213316 if (!have_networking) return error.NetworkDown;
13293 const t: *Threaded = @ptrCast(@alignCast(userdata));
13294 _ = t;
1329513317
1329613318 var iovecs: [max_iovecs_len]posix.iovec_const = undefined;
1329713319 var msg: posix.msghdr_const = .{
......@@ -13377,15 +13399,12 @@ fn netWritePosix(
1337713399}
1337813400
1337913401fn netWriteWindows(
13380 userdata: ?*anyopaque,
1338113402 handle: net.Socket.Handle,
1338213403 header: []const u8,
1338313404 data: []const []const u8,
1338413405 splat: usize,
1338513406) net.Stream.Writer.Error!usize {
1338613407 if (!have_networking) return error.NetworkDown;
13387 const t: *Threaded = @ptrCast(@alignCast(userdata));
13388 _ = t;
1338913408
1339013409 var iovecs: [max_iovecs_len]windows.AFD.WSABUF(.@"const") = undefined;
1339113410 var len: u32 = 0;
lib/std/Io/Uring.zig+6-17
......@@ -778,7 +778,6 @@ pub fn io(ev: *Evented) Io {
778778 .netListenUnix = netListenUnixUnavailable,
779779 .netConnectUnix = netConnectUnixUnavailable,
780780 .netSocketCreatePair = netSocketCreatePairUnavailable,
781 .netWrite = netWriteUnavailable,
782781 .netWriteFile = netWriteFileUnavailable,
783782 .netClose = netClose,
784783 .netShutdown = netShutdown,
......@@ -2118,6 +2117,7 @@ fn operate(userdata: ?*anyopaque, operation: Io.Operation) Io.Cancelable!Io.Oper
21182117 break :r error.NetworkDown; // TODO
21192118 },
21202119 },
2120 .net_write => @panic("TODO implement net_write operation"),
21212121 };
21222122}
21232123
......@@ -2413,6 +2413,10 @@ fn batchDrainSubmitted(
24132413 _ = o;
24142414 @panic("TODO implement batchDrainSubmitted for net_read");
24152415 },
2416 .net_write => |o| {
2417 _ = o;
2418 @panic("TODO implement batchDrainSubmitted for net_write");
2419 },
24162420 })) |result| {
24172421 switch (batch.completed.tail) {
24182422 .none => batch.completed.head = index,
......@@ -2516,6 +2520,7 @@ fn batchDrainReady(batch: *Io.Batch) Io.Timeout.Error!void {
25162520 .net_receive => @panic("TODO"),
25172521 .net_send => @panic("TODO"),
25182522 .net_read => @panic("TODO"),
2523 .net_write => @panic("TODO"),
25192524 })) |result| {
25202525 switch (batch.completed.tail) {
25212526 .none => batch.completed.head = index,
......@@ -5153,22 +5158,6 @@ fn netReceive(
51535158 }
51545159}
51555160
5156fn netWriteUnavailable(
5157 userdata: ?*anyopaque,
5158 handle: net.Socket.Handle,
5159 header: []const u8,
5160 data: []const []const u8,
5161 splat: usize,
5162) net.Stream.Writer.Error!usize {
5163 const ev: *Evented = @ptrCast(@alignCast(userdata));
5164 _ = ev;
5165 _ = handle;
5166 _ = header;
5167 _ = data;
5168 _ = splat;
5169 return error.NetworkDown;
5170}
5171
51725161fn netWriteFileUnavailable(
51735162 userdata: ?*anyopaque,
51745163 socket_handle: net.Socket.Handle,
lib/std/Io/net.zig+11-28
......@@ -1348,33 +1348,7 @@ pub const Stream = struct {
13481348 err: ?Error = null,
13491349 write_file_err: ?WriteFileError = null,
13501350
1351 pub const Error = error{
1352 /// Another TCP Fast Open is already in progress.
1353 FastOpenAlreadyInProgress,
1354 /// Network session was unexpectedly closed by recipient.
1355 ConnectionResetByPeer,
1356 /// The output queue for a network interface was full. This generally indicates that the
1357 /// interface has stopped sending, but may be caused by transient congestion. (Normally,
1358 /// this does not occur in Linux. Packets are just silently dropped when a device queue
1359 /// overflows.)
1360 ///
1361 /// This is also caused when there is not enough kernel memory available.
1362 SystemResources,
1363 /// No route to network.
1364 NetworkUnreachable,
1365 /// Network reached but no route to host.
1366 HostUnreachable,
1367 /// The local network interface used to reach the destination is down.
1368 NetworkDown,
1369 /// The destination address is not listening.
1370 ConnectionRefused,
1371 /// The passed address didn't have the correct address family in its sa_family field.
1372 AddressFamilyUnsupported,
1373 /// Local end has been shut down on a connection-oriented socket, or
1374 /// the socket was never connected.
1375 SocketUnconnected,
1376 SocketNotBound,
1377 } || Io.UnexpectedError || Io.Cancelable;
1351 pub const Error = Io.Operation.NetWrite.Error || Io.Cancelable;
13781352
13791353 pub const WriteFileError = Error || error{
13801354 /// The `Io` implementation cannot offer a more efficient
......@@ -1407,7 +1381,16 @@ pub const Stream = struct {
14071381 const io = w.io;
14081382 const buffered = io_w.buffered();
14091383 const handle = w.stream.socket.handle;
1410 const n = io.vtable.netWrite(io.userdata, handle, buffered, data, splat) catch |err| {
1384 const result = io.operate(.{ .net_write = .{
1385 .socket_handle = handle,
1386 .header = buffered,
1387 .data = data,
1388 .splat = splat,
1389 } }) catch |err| {
1390 w.err = err;
1391 return error.WriteFailed;
1392 };
1393 const n = result.net_write catch |err| {
14111394 w.err = err;
14121395 return error.WriteFailed;
14131396 };