authorgravatar for andrew@ziglang.orgAndrew Kelley <andrew@ziglang.org> 2026-02-15 19:29:21+01:00
committergravatar for andrew@ziglang.orgAndrew Kelley <andrew@ziglang.org> 2026-02-15 19:29:21+01:00
logc8f54a2f07fd4b1bc6ba83c618d0e6f6de265b5c
tree460566d2e6d92c5c1916493b87e48d13a1e5b1d2
parentc6eeae8a8c2314fd35b5baddb18f33b0122d319f
parent5763f7dbcc65288d7d13236f82eec9a23d97bee7

Merge pull request 'std.Io: remove select function' (#31223) from remove-select into master

Reviewed-on: https://codeberg.org/ziglang/zig/pulls/31223

7 files changed, 209 insertions(+), 483 deletions(-)

lib/std/Io.zig-38
......@@ -144,10 +144,6 @@ pub const VTable = struct {
144144 swapCancelProtection: *const fn (?*anyopaque, new: CancelProtection) CancelProtection,
145145 checkCancel: *const fn (?*anyopaque) Cancelable!void,
146146
147 /// Blocks until one of the futures from the list has a result ready, such
148 /// that awaiting it will not block. Returns that index.
149 select: *const fn (?*anyopaque, futures: []const *AnyFuture) Cancelable!usize,
150
151147 futexWait: *const fn (?*anyopaque, ptr: *const u32, expected: u32, Timeout) Cancelable!void,
152148 futexWaitUncancelable: *const fn (?*anyopaque, ptr: *const u32, expected: u32) void,
153149 futexWake: *const fn (?*anyopaque, ptr: *const u32, max_waiters: u32) void,
......@@ -2120,40 +2116,6 @@ pub fn sleep(io: Io, duration: Duration, clock: Clock) Cancelable!void {
21202116 } });
21212117}
21222118
2123/// Given a struct with each field a `*Future`, returns a union with the same
2124/// fields, each field type the future's result.
2125pub fn SelectUnion(S: type) type {
2126 const struct_fields = @typeInfo(S).@"struct".fields;
2127 var names: [struct_fields.len][]const u8 = undefined;
2128 var types: [struct_fields.len]type = undefined;
2129 for (struct_fields, &names, &types) |struct_field, *union_field_name, *UnionFieldType| {
2130 const FieldFuture = @typeInfo(struct_field.type).pointer.child;
2131 union_field_name.* = struct_field.name;
2132 UnionFieldType.* = @FieldType(FieldFuture, "result");
2133 }
2134 return @Union(.auto, std.meta.FieldEnum(S), &names, &types, &@splat(.{}));
2135}
2136
2137/// `s` is a struct with every field a `*Future(T)`, where `T` can be any type,
2138/// and can be different for each field.
2139pub fn select(io: Io, s: anytype) Cancelable!SelectUnion(@TypeOf(s)) {
2140 const U = SelectUnion(@TypeOf(s));
2141 const S = @TypeOf(s);
2142 const fields = @typeInfo(S).@"struct".fields;
2143 var futures: [fields.len]*AnyFuture = undefined;
2144 inline for (fields, &futures) |field, *any_future| {
2145 const future = @field(s, field.name);
2146 any_future.* = future.any_future orelse return @unionInit(U, field.name, future.result);
2147 }
2148 switch (try io.vtable.select(io.userdata, &futures)) {
2149 inline 0...(fields.len - 1) => |selected_index| {
2150 const field_name = fields[selected_index].name;
2151 return @unionInit(U, field_name, @field(s, field_name).await(io));
2152 },
2153 else => unreachable,
2154 }
2155}
2156
21572119pub const LockedStderr = struct {
21582120 file_writer: *File.Writer,
21592121 terminal_mode: Terminal.Mode,
lib/std/Io/Dispatch.zig+33-103
......@@ -100,15 +100,11 @@ const Fiber = struct {
100100 required_align: void align(4),
101101 evented: *Evented,
102102 context: Io.fiber.Context,
103 await_count: i32,
104103 link: union {
105104 awaiter: ?*Fiber,
106105 group: struct { prev: ?*Fiber, next: ?*Fiber },
107106 },
108 status: union(enum) {
109 queue_next: ?*Fiber,
110 awaiting_group: Group,
111 },
107 awaiting_group: Group,
112108 cancel_status: CancelStatus,
113109 cancel_protection: CancelProtection,
114110
......@@ -123,7 +119,6 @@ const Fiber = struct {
123119 const Awaiting = enum(@Int(.unsigned, @bitSizeOf(usize) - shift)) {
124120 nothing = 0,
125121 group = 1,
126 select = 2,
127122 _,
128123
129124 const shift = 1;
......@@ -216,7 +211,6 @@ const Fiber = struct {
216211 }
217212
218213 fn destroy(fiber: *Fiber, ev: *Evented) void {
219 assert(fiber.status.queue_next == null);
220214 ev.allocator().free(fiber.allocatedSlice());
221215 }
222216
......@@ -272,14 +266,11 @@ const Fiber = struct {
272266 .group => {
273267 // The awaiter received a cancelation request while awaiting a group,
274268 // so propagate the cancelation to the group.
275 if (fiber.status.awaiting_group.cancel(ev, null)) {
276 fiber.status = .{ .queue_next = null };
269 if (fiber.awaiting_group.cancel(ev, null)) {
270 fiber.awaiting_group = undefined;
277271 ev.queue.async(fiber, &Fiber.@"resume");
278272 }
279273 },
280 .select => if (@atomicRmw(i32, &fiber.await_count, .Add, 1, .monotonic) == -1) {
281 ev.queue.async(fiber, &Fiber.@"resume");
282 },
283274 _ => |awaiting| awaiting.toCancelable().async(),
284275 }
285276 }
......@@ -370,8 +361,6 @@ pub fn io(ev: *Evented) Io {
370361 .swapCancelProtection = swapCancelProtection,
371362 .checkCancel = checkCancel,
372363
373 .select = select,
374
375364 .futexWait = futexWait,
376365 .futexWaitUncancelable = futexWaitUncancelable,
377366 .futexWake = futexWake,
......@@ -522,9 +511,8 @@ pub fn init(ev: *Evented, backing_allocator: Allocator, options: InitOptions) !v
522511 .required_align = {},
523512 .evented = ev,
524513 .context = undefined,
525 .await_count = 0,
526514 .link = .{ .awaiter = null },
527 .status = .{ .queue_next = null },
515 .awaiting_group = undefined,
528516 .cancel_status = .unrequested,
529517 .cancel_protection = .unblocked,
530518 },
......@@ -642,7 +630,7 @@ const SwitchMessage = struct {
642630
643631 const PendingTask = union(enum) {
644632 nothing,
645 await: u31,
633 await: *Fiber,
646634 activate: c.dispatch.object_t,
647635 @"resume": c.dispatch.object_t,
648636 group_await: Group,
......@@ -661,10 +649,10 @@ const SwitchMessage = struct {
661649 thread.current_context = message.contexts.new;
662650 switch (message.pending_task) {
663651 .nothing => {},
664 .await => |count| {
665 const fiber: *Fiber = @alignCast(@fieldParentPtr("context", message.contexts.old));
666 if (@atomicRmw(i32, &fiber.await_count, .Sub, count, .monotonic) > 0)
667 ev.queue.async(fiber, &Fiber.@"resume");
652 .await => |awaiting| {
653 const awaiter: *Fiber = @alignCast(@fieldParentPtr("context", message.contexts.old));
654 if (@atomicRmw(?*Fiber, &awaiting.link.awaiter, .Xchg, awaiter, .acq_rel) ==
655 Fiber.finished) ev.queue.async(awaiter, &Fiber.@"resume");
668656 },
669657 .activate => |object| object.activate(),
670658 .@"resume" => |object| object.@"resume"(),
......@@ -998,7 +986,7 @@ fn crashHandler(userdata: ?*anyopaque) void {
998986}
999987
1000988const AsyncClosure = struct {
1001 ev: *Evented,
989 evented: *Evented,
1002990 fiber: *Fiber,
1003991 start: *const fn (context: *const anyopaque, result: *anyopaque) void,
1004992 result_align: Alignment,
......@@ -1035,13 +1023,13 @@ const AsyncClosure = struct {
10351023 closure: *AsyncClosure,
10361024 message: *const SwitchMessage,
10371025 ) callconv(.withStackAlign(.c, @alignOf(AsyncClosure))) noreturn {
1038 message.handle(closure.ev);
1026 const ev = closure.evented;
10391027 const fiber = closure.fiber;
1028 message.handle(ev);
10401029 closure.start(closure.contextPointer(), fiber.resultBytes(closure.result_align));
10411030 if (@atomicRmw(?*Fiber, &fiber.link.awaiter, .Xchg, Fiber.finished, .acq_rel)) |awaiter|
1042 if (@atomicRmw(i32, &awaiter.await_count, .Add, 1, .monotonic) == -1)
1043 closure.ev.queue.async(awaiter, &Fiber.@"resume");
1044 closure.ev.yield(.nothing);
1031 ev.queue.async(awaiter, &Fiber.@"resume");
1032 ev.yield(.nothing);
10451033 unreachable; // switched to dead fiber
10461034 }
10471035};
......@@ -1096,14 +1084,13 @@ fn concurrent(
10961084 },
10971085 else => |arch| @compileError("unimplemented architecture: " ++ @tagName(arch)),
10981086 },
1099 .await_count = 0,
11001087 .link = .{ .awaiter = null },
1101 .status = .{ .queue_next = null },
1088 .awaiting_group = undefined,
11021089 .cancel_status = .unrequested,
11031090 .cancel_protection = .unblocked,
11041091 };
11051092 closure.* = .{
1106 .ev = ev,
1093 .evented = ev,
11071094 .fiber = fiber,
11081095 .start = start,
11091096 .result_align = result_alignment,
......@@ -1121,18 +1108,11 @@ fn await(
11211108 result_alignment: Alignment,
11221109) void {
11231110 const ev: *Evented = @ptrCast(@alignCast(userdata));
1124 const fiber = Thread.current().currentFiber();
1125 const future_fiber: *Fiber = @ptrCast(@alignCast(future));
1126 if (@atomicRmw(?*Fiber, &future_fiber.link.awaiter, .Xchg, fiber, .acq_rel)) |awaiter| {
1127 assert(awaiter == Fiber.finished);
1128 } else while (true) {
1129 ev.yield(.{ .await = 1 });
1130 const awaiter = @atomicLoad(?*Fiber, &future_fiber.link.awaiter, .acquire);
1131 if (awaiter == Fiber.finished) break;
1132 assert(awaiter == fiber); // spurious wakeup
1133 }
1134 @memcpy(result, future_fiber.resultBytes(result_alignment));
1135 future_fiber.destroy(ev);
1111 const awaiting: *Fiber = @ptrCast(@alignCast(future));
1112 if (@atomicLoad(?*Fiber, &awaiting.link.awaiter, .acquire) != Fiber.finished)
1113 ev.yield(.{ .await = awaiting });
1114 @memcpy(result, awaiting.resultBytes(result_alignment));
1115 awaiting.destroy(ev);
11361116}
11371117
11381118fn cancel(
......@@ -1267,8 +1247,8 @@ const Group = struct {
12671247 .awaiter_delayed = false,
12681248 .fibers = .null,
12691249 }, .release);
1270 assert(awaiter.status.awaiting_group.ptr == group.ptr);
1271 awaiter.status = .{ .queue_next = null };
1250 assert(awaiter.awaiting_group.ptr == group.ptr);
1251 awaiter.awaiting_group = undefined;
12721252 return awaiter;
12731253 }
12741254 // Race with `Fiber.requestCancel`
......@@ -1336,8 +1316,7 @@ const Group = struct {
13361316
13371317 /// Assumes the mutex is held.
13381318 fn registerAwaiter(group: Group, awaiter: *Fiber) bool {
1339 assert(awaiter.status.queue_next == null);
1340 awaiter.status = .{ .awaiting_group = group };
1319 awaiter.awaiting_group = group;
13411320 assert(@atomicRmw(
13421321 Awaiter,
13431322 group.awaiterPtr(),
......@@ -1349,7 +1328,7 @@ const Group = struct {
13491328 }
13501329
13511330 const AsyncClosure = struct {
1352 ev: *Evented,
1331 evented: *Evented,
13531332 group: Group,
13541333 fiber: *Fiber,
13551334 start: *const fn (context: *const anyopaque) Io.Cancelable!void,
......@@ -1388,19 +1367,15 @@ const Group = struct {
13881367 closure: *Group.AsyncClosure,
13891368 message: *const SwitchMessage,
13901369 ) callconv(.withStackAlign(.c, @alignOf(Group.AsyncClosure))) noreturn {
1391 message.handle(closure.ev);
1392 assert(closure.fiber.status.queue_next == null);
1393 const result = closure.start(closure.contextPointer());
1394 const ev = closure.ev;
1395 const group = closure.group;
1370 const ev = closure.evented;
13961371 const fiber = closure.fiber;
1397 const cancel_acknowledged = fiber.cancel_protection.acknowledged;
1398 if (result) {
1399 assert(!cancel_acknowledged); // group task acknowledged cancelation but did not return `error.Canceled`
1372 message.handle(ev);
1373 if (closure.start(closure.contextPointer())) {
1374 assert(!fiber.cancel_protection.acknowledged); // group task acknowledged cancelation but did not return `error.Canceled`
14001375 } else |err| switch (err) {
1401 error.Canceled => assert(cancel_acknowledged), // group task returned `error.Canceled` but was never canceled
1376 error.Canceled => assert(fiber.cancel_protection.acknowledged), // group task returned `error.Canceled` but was never canceled
14021377 }
1403 if (group.removeFiber(ev, fiber)) |awaiter| ev.queue.async(awaiter, &Fiber.@"resume");
1378 if (closure.group.removeFiber(ev, fiber)) |awaiter| ev.queue.async(awaiter, &Fiber.@"resume");
14041379 ev.yield(.destroy);
14051380 unreachable; // switched to dead fiber
14061381 }
......@@ -1470,14 +1445,13 @@ fn groupConcurrent(
14701445 },
14711446 else => |arch| @compileError("unimplemented architecture: " ++ @tagName(arch)),
14721447 },
1473 .await_count = 0,
14741448 .link = .{ .group = .{ .prev = null, .next = null } },
1475 .status = .{ .queue_next = null },
1449 .awaiting_group = undefined,
14761450 .cancel_status = .unrequested,
14771451 .cancel_protection = .unblocked,
14781452 };
14791453 closure.* = .{
1480 .ev = ev,
1454 .evented = ev,
14811455 .group = group,
14821456 .fiber = fiber,
14831457 .start = start,
......@@ -1689,50 +1663,6 @@ fn futexForAddress(ev: *Evented, address: usize) *Futex {
16891663 return &ev.futexes[hashed >> @clz(ev.futexes.len - 1)];
16901664}
16911665
1692fn select(userdata: ?*anyopaque, futures: []const *Io.AnyFuture) Io.Cancelable!usize {
1693 const ev: *Evented = @ptrCast(@alignCast(userdata));
1694 const fiber = Thread.current().currentFiber();
1695 var await_count: u31, var result = for (futures, 0..) |future, future_index| {
1696 const future_fiber: *Fiber = @ptrCast(@alignCast(future));
1697 if (@atomicRmw(
1698 ?*Fiber,
1699 &future_fiber.link.awaiter,
1700 .Xchg,
1701 fiber,
1702 .acq_rel,
1703 )) |awaiter| {
1704 assert(awaiter == Fiber.finished);
1705 break .{ @intCast(future_index), future_index };
1706 }
1707 } else result: {
1708 const await_count: u31 = @intCast(futures.len);
1709 ev.yield(.{ .await = 1 });
1710 break :result .{ await_count - 1, futures.len };
1711 };
1712 for (futures[0..result], 0..) |future, future_index| {
1713 const future_fiber: *Fiber = @ptrCast(@alignCast(future));
1714 const awaiter = @atomicRmw(?*Fiber, &future_fiber.link.awaiter, .Xchg, null, .monotonic);
1715 if (awaiter == Fiber.finished) {
1716 @atomicStore(?*Fiber, &future_fiber.link.awaiter, Fiber.finished, .monotonic);
1717 result = @min(future_index, result);
1718 } else {
1719 assert(awaiter == fiber);
1720 await_count -= 1;
1721 }
1722 }
1723 // Equivalent to `ev.yield(null, .{ .await = await_count });`,
1724 // but avoiding a context switch in the common case.
1725 switch (std.math.order(
1726 @atomicRmw(i32, &fiber.await_count, .Sub, await_count, .monotonic),
1727 await_count,
1728 )) {
1729 .lt => ev.yield(.{ .await = 0 }),
1730 .eq => {},
1731 .gt => unreachable,
1732 }
1733 return result;
1734}
1735
17361666fn futexWait(
17371667 userdata: ?*anyopaque,
17381668 ptr: *const u32,
lib/std/Io/Kqueue.zig-22
......@@ -491,7 +491,6 @@ const SwitchMessage = struct {
491491 reschedule,
492492 recycle: *Fiber,
493493 register_awaiter: *?*Fiber,
494 register_select: []const *Io.AnyFuture,
495494 exit,
496495 };
497496
......@@ -514,19 +513,6 @@ const SwitchMessage = struct {
514513 if (@atomicRmw(?*Fiber, awaiter, .Xchg, prev_fiber, .acq_rel) == Fiber.finished)
515514 k.schedule(thread, .{ .head = prev_fiber, .tail = prev_fiber });
516515 },
517 .register_select => |futures| {
518 const prev_fiber: *Fiber = @alignCast(@fieldParentPtr("context", message.contexts.old));
519 assert(prev_fiber.queue_next == null);
520 for (futures) |any_future| {
521 const future_fiber: *Fiber = @ptrCast(@alignCast(any_future));
522 if (@atomicRmw(?*Fiber, &future_fiber.awaiter, .Xchg, prev_fiber, .acq_rel) == Fiber.finished) {
523 const closure: *AsyncClosure = .fromFiber(future_fiber);
524 if (!@atomicRmw(bool, &closure.already_awaited, .Xchg, true, .seq_cst)) {
525 k.schedule(thread, .{ .head = prev_fiber, .tail = prev_fiber });
526 }
527 }
528 }
529 },
530516 .exit => for (k.threads.allocated[0..@atomicLoad(u32, &k.threads.active, .acquire)]) |*each_thread| {
531517 const changes = [_]posix.Kevent{
532518 .{
......@@ -628,7 +614,6 @@ pub fn io(k: *Kqueue) Io {
628614 .concurrent = concurrent,
629615 .await = await,
630616 .cancel = cancel,
631 .select = select,
632617
633618 .groupAsync = groupAsync,
634619 .groupConcurrent = groupConcurrent,
......@@ -824,13 +809,6 @@ fn groupCancel(userdata: ?*anyopaque, group: *Io.Group, token: *anyopaque) void
824809 @panic("TODO");
825810}
826811
827fn select(userdata: ?*anyopaque, futures: []const *Io.AnyFuture) Io.Cancelable!usize {
828 const k: *Kqueue = @ptrCast(@alignCast(userdata));
829 _ = k;
830 _ = futures;
831 @panic("TODO");
832}
833
834812fn dirCreateDir(userdata: ?*anyopaque, dir: Dir, sub_path: []const u8, permissions: Dir.Permissions) Dir.CreateDirError!void {
835813 const k: *Kqueue = @ptrCast(@alignCast(userdata));
836814 _ = k;
lib/std/Io/Threaded.zig-70
......@@ -1772,7 +1772,6 @@ pub fn io(t: *Threaded) Io {
17721772 .concurrent = concurrent,
17731773 .await = await,
17741774 .cancel = cancel,
1775 .select = select,
17761775
17771776 .groupAsync = groupAsync,
17781777 .groupConcurrent = groupConcurrent,
......@@ -1938,7 +1937,6 @@ pub fn ioBasic(t: *Threaded) Io {
19381937 .concurrent = concurrent,
19391938 .await = await,
19401939 .cancel = cancel,
1941 .select = select,
19421940
19431941 .groupAsync = groupAsync,
19441942 .groupConcurrent = groupConcurrent,
......@@ -11727,74 +11725,6 @@ fn sleepNanosleep(t: *Threaded, timeout: Io.Timeout) Io.Cancelable!void {
1172711725 }
1172811726}
1172911727
11730fn select(userdata: ?*anyopaque, futures: []const *Io.AnyFuture) Io.Cancelable!usize {
11731 const t: *Threaded = @ptrCast(@alignCast(userdata));
11732 _ = t;
11733
11734 var num_completed: std.atomic.Value(u32) = .init(0);
11735
11736 for (futures, 0..) |any_future, i| {
11737 const future: *Future = @ptrCast(@alignCast(any_future));
11738 future.awaiter = &num_completed;
11739 const old_status = future.status.fetchOr(
11740 .{ .tag = .pending_awaited, .thread = .null },
11741 .release, // release `future.awaiter`
11742 );
11743 switch (old_status.tag) {
11744 .pending => {},
11745 .pending_awaited => unreachable, // `await` raced with `select`
11746 .pending_canceled => unreachable, // `cancel` raced with `select`
11747 .done => {
11748 future.status.store(old_status, .monotonic);
11749 _ = finishSelect(&num_completed, futures[0..i]);
11750 return i;
11751 },
11752 }
11753 }
11754
11755 errdefer _ = finishSelect(&num_completed, futures);
11756
11757 while (true) {
11758 const n = num_completed.load(.acquire);
11759 if (n > 0) break;
11760 assert(n < futures.len);
11761 try Thread.futexWait(&num_completed.raw, n, null);
11762 }
11763 return finishSelect(&num_completed, futures).?;
11764}
11765fn finishSelect(
11766 num_completed: *std.atomic.Value(u32),
11767 futures: []const *Io.AnyFuture,
11768) ?usize {
11769 var completed_index: ?usize = null;
11770 var expect_completed: u32 = 0;
11771 for (futures, 0..) |any_future, i| {
11772 const future: *Future = @ptrCast(@alignCast(any_future));
11773 // This operation will convert `.pending_awaited` to `.pending`, or leave `.done` untouched.
11774 switch (future.status.fetchAnd(
11775 .{ .tag = @enumFromInt(0b10), .thread = .all_ones },
11776 .monotonic,
11777 ).tag) {
11778 .pending_awaited => {},
11779 .pending => unreachable,
11780 .pending_canceled => unreachable,
11781 .done => {
11782 expect_completed += 1;
11783 completed_index = i;
11784 },
11785 }
11786 }
11787 // If any future has just finished, wait for it to signal `num_completed` to avoid dangling
11788 // references to stack memory.
11789 while (true) {
11790 const n = num_completed.load(.acquire);
11791 if (n == expect_completed) break;
11792 assert(n < expect_completed);
11793 Thread.futexWaitUncancelable(&num_completed.raw, n, null);
11794 }
11795 return completed_index;
11796}
11797
1179811728fn netListenIpPosix(
1179911729 userdata: ?*anyopaque,
1180011730 address: IpAddress,
lib/std/Io/Uring.zig+172-216
......@@ -150,7 +150,6 @@ const Thread = struct {
150150const Fiber = struct {
151151 required_align: void align(4),
152152 context: Io.fiber.Context,
153 await_count: i32,
154153 link: union {
155154 awaiter: ?*Fiber,
156155 group: struct { prev: ?*Fiber, next: ?*Fiber },
......@@ -175,7 +174,6 @@ const Fiber = struct {
175174 const Awaiting = enum(u31) {
176175 nothing = std.math.maxInt(u31),
177176 group = std.math.maxInt(u31) - 1,
178 select = std.math.maxInt(u31) - 2,
179177 /// An io_uring fd.
180178 _,
181179
......@@ -186,14 +184,14 @@ const Fiber = struct {
186184 fn fromIoUringFd(fd: fd_t) Awaiting {
187185 const awaiting: Awaiting = @enumFromInt(fd);
188186 switch (awaiting) {
189 .nothing, .group, .select => unreachable,
187 .nothing, .group => unreachable,
190188 _ => return awaiting,
191189 }
192190 }
193191
194192 fn toIoUringFd(awaiting: Awaiting) fd_t {
195193 switch (awaiting) {
196 .nothing, .group, .select => unreachable,
194 .nothing, .group => unreachable,
197195 _ => return @intFromEnum(awaiting),
198196 }
199197 }
......@@ -376,9 +374,6 @@ const Fiber = struct {
376374 _ = ev.schedule(.current(), .{ .head = fiber, .tail = fiber });
377375 }
378376 },
379 .select => if (@atomicRmw(i32, &fiber.await_count, .Add, 1, .monotonic) == -1) {
380 _ = ev.schedule(.current(), .{ .head = fiber, .tail = fiber });
381 },
382377 _ => |awaiting| {
383378 const awaiting_io_uring_fd = awaiting.toIoUringFd();
384379 const thread: *Thread = .current();
......@@ -684,8 +679,6 @@ pub fn io(ev: *Evented) Io {
684679 .swapCancelProtection = swapCancelProtection,
685680 .checkCancel = checkCancel,
686681
687 .select = select,
688
689682 .futexWait = futexWait,
690683 .futexWaitUncancelable = futexWaitUncancelable,
691684 .futexWake = futexWake,
......@@ -862,7 +855,6 @@ pub fn init(ev: *Evented, backing_allocator: Allocator, options: InitOptions) !v
862855 main_fiber.* = .{
863856 .required_align = {},
864857 .context = undefined,
865 .await_count = 0,
866858 .link = .{ .awaiter = null },
867859 .status = .{ .queue_next = null },
868860 .cancel_status = .unrequested,
......@@ -1266,7 +1258,7 @@ const SwitchMessage = struct {
12661258 const PendingTask = union(enum) {
12671259 nothing,
12681260 reschedule,
1269 await: u31,
1261 await: *Fiber,
12701262 group_await: Group,
12711263 group_cancel: Group,
12721264 batch_await: *Io.Batch,
......@@ -1290,10 +1282,11 @@ const SwitchMessage = struct {
12901282 assert(fiber.status.queue_next == null);
12911283 _ = ev.schedule(thread, .{ .head = fiber, .tail = fiber });
12921284 },
1293 .await => |count| {
1294 const fiber: *Fiber = @alignCast(@fieldParentPtr("context", message.contexts.old));
1295 if (@atomicRmw(i32, &fiber.await_count, .Sub, count, .monotonic) > 0)
1296 _ = ev.schedule(thread, .{ .head = fiber, .tail = fiber });
1285 .await => |awaiting| {
1286 const awaiter: *Fiber = @alignCast(@fieldParentPtr("context", message.contexts.old));
1287 assert(awaiter.status.queue_next == null);
1288 if (@atomicRmw(?*Fiber, &awaiting.link.awaiter, .Xchg, awaiter, .acq_rel) ==
1289 Fiber.finished) _ = ev.schedule(thread, .{ .head = awaiter, .tail = awaiter });
12971290 },
12981291 .group_await => |group| {
12991292 const fiber: *Fiber = @alignCast(@fieldParentPtr("context", message.contexts.old));
......@@ -1367,7 +1360,7 @@ fn crashHandler(userdata: ?*anyopaque) void {
13671360}
13681361
13691362const AsyncClosure = struct {
1370 ev: *Evented,
1363 evented: *Evented,
13711364 fiber: *Fiber,
13721365 start: *const fn (context: *const anyopaque, result: *anyopaque) void,
13731366 result_align: Alignment,
......@@ -1404,16 +1397,11 @@ const AsyncClosure = struct {
14041397 closure: *AsyncClosure,
14051398 message: *const SwitchMessage,
14061399 ) callconv(.withStackAlign(.c, @alignOf(AsyncClosure))) noreturn {
1407 message.handle(closure.ev);
1400 const ev = closure.evented;
14081401 const fiber = closure.fiber;
1402 message.handle(ev);
14091403 closure.start(closure.contextPointer(), fiber.resultBytes(closure.result_align));
1410 closure.ev.yield(
1411 if (@atomicRmw(?*Fiber, &fiber.link.awaiter, .Xchg, Fiber.finished, .acq_rel)) |awaiter|
1412 if (@atomicRmw(i32, &awaiter.await_count, .Add, 1, .monotonic) == -1) awaiter else null
1413 else
1414 null,
1415 .nothing,
1416 );
1404 ev.yield(@atomicRmw(?*Fiber, &fiber.link.awaiter, .Xchg, Fiber.finished, .acq_rel), .nothing);
14171405 unreachable; // switched to dead fiber
14181406 }
14191407};
......@@ -1467,7 +1455,6 @@ fn concurrent(
14671455 },
14681456 else => |arch| @compileError("unimplemented architecture: " ++ @tagName(arch)),
14691457 },
1470 .await_count = 0,
14711458 .link = .{ .awaiter = null },
14721459 .status = .{ .queue_next = null },
14731460 .cancel_status = .unrequested,
......@@ -1485,7 +1472,7 @@ fn concurrent(
14851472 },
14861473 };
14871474 closure.* = .{
1488 .ev = ev,
1475 .evented = ev,
14891476 .fiber = fiber,
14901477 .start = start,
14911478 .result_align = result_alignment,
......@@ -1504,18 +1491,11 @@ fn await(
15041491 result_alignment: Alignment,
15051492) void {
15061493 const ev: *Evented = @ptrCast(@alignCast(userdata));
1507 const fiber = Thread.current().currentFiber();
1508 const future_fiber: *Fiber = @ptrCast(@alignCast(future));
1509 if (@atomicRmw(?*Fiber, &future_fiber.link.awaiter, .Xchg, fiber, .acq_rel)) |awaiter| {
1510 assert(awaiter == Fiber.finished);
1511 } else while (true) {
1512 ev.yield(null, .{ .await = 1 });
1513 const awaiter = @atomicLoad(?*Fiber, &future_fiber.link.awaiter, .acquire);
1514 if (awaiter == Fiber.finished) break;
1515 assert(awaiter == fiber); // spurious wakeup
1516 }
1517 @memcpy(result, future_fiber.resultBytes(result_alignment));
1518 future_fiber.destroy();
1494 const awaiting: *Fiber = @ptrCast(@alignCast(future));
1495 if (@atomicLoad(?*Fiber, &awaiting.link.awaiter, .acquire) != Fiber.finished)
1496 ev.yield(null, .{ .await = awaiting });
1497 @memcpy(result, awaiting.resultBytes(result_alignment));
1498 awaiting.destroy();
15191499}
15201500
15211501fn cancel(
......@@ -1732,7 +1712,7 @@ const Group = struct {
17321712 }
17331713
17341714 const AsyncClosure = struct {
1735 ev: *Evented,
1715 evented: *Evented,
17361716 group: Group,
17371717 fiber: *Fiber,
17381718 start: *const fn (context: *const anyopaque) Io.Cancelable!void,
......@@ -1771,19 +1751,16 @@ const Group = struct {
17711751 closure: *Group.AsyncClosure,
17721752 message: *const SwitchMessage,
17731753 ) callconv(.withStackAlign(.c, @alignOf(Group.AsyncClosure))) noreturn {
1774 message.handle(closure.ev);
1775 assert(closure.fiber.status.queue_next == null);
1776 const result = closure.start(closure.contextPointer());
1777 const ev = closure.ev;
1778 const group = closure.group;
1754 const ev = closure.evented;
17791755 const fiber = closure.fiber;
1780 const cancel_acknowledged = fiber.cancel_protection.acknowledged;
1781 if (result) {
1782 assert(!cancel_acknowledged); // group task acknowledged cancelation but did not return `error.Canceled`
1756 message.handle(ev);
1757 assert(fiber.status.queue_next == null);
1758 if (closure.start(closure.contextPointer())) {
1759 assert(!fiber.cancel_protection.acknowledged); // group task acknowledged cancelation but did not return `error.Canceled`
17831760 } else |err| switch (err) {
1784 error.Canceled => assert(cancel_acknowledged), // group task returned `error.Canceled` but was never canceled
1761 error.Canceled => assert(fiber.cancel_protection.acknowledged), // group task returned `error.Canceled` but was never canceled
17851762 }
1786 ev.yield(group.removeFiber(ev, fiber), .destroy);
1763 ev.yield(closure.group.removeFiber(ev, fiber), .destroy);
17871764 unreachable; // switched to dead fiber
17881765 }
17891766 };
......@@ -1851,7 +1828,6 @@ fn groupConcurrent(
18511828 },
18521829 else => |arch| @compileError("unimplemented architecture: " ++ @tagName(arch)),
18531830 },
1854 .await_count = 0,
18551831 .link = .{ .group = .{ .prev = null, .next = null } },
18561832 .status = .{ .queue_next = null },
18571833 .cancel_status = .unrequested,
......@@ -1869,7 +1845,7 @@ fn groupConcurrent(
18691845 },
18701846 };
18711847 closure.* = .{
1872 .ev = ev,
1848 .evented = ev,
18731849 .group = group,
18741850 .fiber = fiber,
18751851 .start = start,
......@@ -1928,57 +1904,6 @@ fn checkCancel(userdata: ?*anyopaque) Io.Cancelable!void {
19281904 }
19291905}
19301906
1931fn select(userdata: ?*anyopaque, futures: []const *Io.AnyFuture) Io.Cancelable!usize {
1932 const ev: *Evented = @ptrCast(@alignCast(userdata));
1933 var cancel_region: CancelRegion = .init();
1934 defer cancel_region.deinit();
1935 var await_count: u31, var result = for (futures, 0..) |future, future_index| {
1936 const future_fiber: *Fiber = @ptrCast(@alignCast(future));
1937 if (@atomicRmw(
1938 ?*Fiber,
1939 &future_fiber.link.awaiter,
1940 .Xchg,
1941 cancel_region.fiber,
1942 .acq_rel,
1943 )) |awaiter| {
1944 assert(awaiter == Fiber.finished);
1945 break .{ @intCast(future_index), future_index };
1946 }
1947 } else result: {
1948 const await_count: u31 = @intCast(futures.len);
1949 cancel_region.await(.select) catch |err| switch (err) {
1950 error.Canceled => |e| break :result .{ await_count + 1, e },
1951 };
1952 ev.yield(null, .{ .await = 1 });
1953 cancel_region.await(.nothing) catch |err| switch (err) {
1954 error.Canceled => |e| break :result .{ await_count, e },
1955 };
1956 break :result .{ await_count - 1, futures.len };
1957 };
1958 for (futures[0 .. result catch futures.len], 0..) |future, future_index| {
1959 const future_fiber: *Fiber = @ptrCast(@alignCast(future));
1960 const awaiter = @atomicRmw(?*Fiber, &future_fiber.link.awaiter, .Xchg, null, .monotonic);
1961 if (awaiter == Fiber.finished) {
1962 @atomicStore(?*Fiber, &future_fiber.link.awaiter, Fiber.finished, .monotonic);
1963 result = if (result) |finished_index| @min(future_index, finished_index) else |e| e;
1964 } else {
1965 assert(awaiter == cancel_region.fiber);
1966 await_count -= 1;
1967 }
1968 }
1969 // Equivalent to `ev.yield(null, .{ .await = await_count });`,
1970 // but avoiding a context switch in the common case.
1971 switch (std.math.order(
1972 @atomicRmw(i32, &cancel_region.fiber.await_count, .Sub, await_count, .monotonic),
1973 await_count,
1974 )) {
1975 .lt => ev.yield(null, .{ .await = 0 }),
1976 .eq => {},
1977 .gt => unreachable,
1978 }
1979 return result;
1980}
1981
19821907fn futexWait(
19831908 userdata: ?*anyopaque,
19841909 ptr: *const u32,
......@@ -2230,7 +2155,7 @@ fn deviceIoControl(
22302155 const rc = linux.ioctl(o.file.handle, @bitCast(o.code), @intFromPtr(o.arg));
22312156 switch (linux.errno(rc)) {
22322157 .SUCCESS => return @bitCast(@as(u32, @truncate(rc))),
2233 .INTR => continue,
2158 .INTR => {},
22342159 else => |err| return -@as(i32, @intFromEnum(err)),
22352160 }
22362161 }
......@@ -2628,7 +2553,7 @@ fn dirCreateDir(
26282553 ev.yield(null, .nothing);
26292554 switch (cancel_region.errno()) {
26302555 .SUCCESS => return,
2631 .INTR, .CANCELED => continue,
2556 .INTR, .CANCELED => {},
26322557 .ACCES => return error.AccessDenied,
26332558 .BADF => |err| return errnoBug(err), // File descriptor used after closed.
26342559 .PERM => return error.PermissionDenied,
......@@ -2711,7 +2636,7 @@ fn filePathKind(ev: *Evented, dir: Dir, sub_path: []const u8) !File.Kind {
27112636 if (!statx_buf.mask.TYPE) return error.Unexpected;
27122637 return statxKind(statx_buf.mode);
27132638 },
2714 .INTR, .CANCELED => continue,
2639 .INTR, .CANCELED => {},
27152640 .ACCES => |err| return errnoBug(err),
27162641 .BADF => |err| return errnoBug(err), // File descriptor used after closed.
27172642 .FAULT => |err| return errnoBug(err),
......@@ -2824,7 +2749,7 @@ fn dirAccess(
28242749 try sync.cancel_region.await(.nothing);
28252750 switch (linux.errno(linux.faccessat(dir.handle, sub_path_posix, mode, flags))) {
28262751 .SUCCESS => return,
2827 .INTR => continue,
2752 .INTR => {},
28282753 .ACCES => return error.AccessDenied,
28292754 .PERM => return error.PermissionDenied,
28302755 .ROFS => return error.ReadOnlyFileSystem,
......@@ -2863,7 +2788,7 @@ fn dirCreateFile(
28632788 .EXCL = flags.exclusive,
28642789 .CLOEXEC = true,
28652790 }, flags.permissions.toMode());
2866 errdefer ev.close(fd);
2791 errdefer ev.closeAsync(fd);
28672792
28682793 switch (flags.lock) {
28692794 .none => {},
......@@ -3051,7 +2976,7 @@ fn dirOpenFile(
30512976 .CLOEXEC = true,
30522977 .PATH = flags.path_only,
30532978 }, 0);
3054 errdefer ev.close(fd);
2979 errdefer ev.closeAsync(fd);
30552980
30562981 if (!flags.allow_directory) {
30572982 const is_dir = is_dir: {
......@@ -3105,7 +3030,7 @@ fn dirRead(userdata: ?*anyopaque, dr: *Dir.Reader, buffer: []Dir.Entry) Dir.Read
31053030 const rc = linux.getdents64(dr.dir.handle, dr.buffer.ptr, dr.buffer.len);
31063031 switch (linux.errno(rc)) {
31073032 .SUCCESS => break rc,
3108 .INTR => continue,
3033 .INTR => {},
31093034 .BADF => |err| return errnoBug(err), // Dir is invalid or was opened without iteration ability.
31103035 .FAULT => |err| return errnoBug(err),
31113036 .NOTDIR => |err| return errnoBug(err),
......@@ -3198,7 +3123,7 @@ fn dirRealPathFile(
31983123 error.FileLocksUnsupported => return errnoBug(.OPNOTSUPP), // Not asking for locks.
31993124 else => |e| return e,
32003125 };
3201 defer ev.close(fd);
3126 defer ev.closeAsync(fd);
32023127 return ev.realPath(try maybe_sync.enterSync(ev), fd, out_buffer);
32033128}
32043129
......@@ -3231,7 +3156,7 @@ fn dirDeleteFile(userdata: ?*anyopaque, dir: Dir, sub_path: []const u8) Dir.Dele
32313156 ev.yield(null, .nothing);
32323157 switch (cancel_region.errno()) {
32333158 .SUCCESS => return,
3234 .INTR, .CANCELED => continue,
3159 .INTR, .CANCELED => {},
32353160 .PERM => return error.PermissionDenied,
32363161 .ACCES => return error.AccessDenied,
32373162 .BUSY => return error.FileBusy,
......@@ -3283,7 +3208,7 @@ fn dirDeleteDir(userdata: ?*anyopaque, dir: Dir, sub_path: []const u8) Dir.Delet
32833208 ev.yield(null, .nothing);
32843209 switch (cancel_region.errno()) {
32853210 .SUCCESS => return,
3286 .INTR, .CANCELED => continue,
3211 .INTR, .CANCELED => {},
32873212 .ACCES => return error.AccessDenied,
32883213 .PERM => return error.PermissionDenied,
32893214 .BUSY => return error.FileBusy,
......@@ -3399,7 +3324,7 @@ fn dirSymLink(
33993324 ev.yield(null, .nothing);
34003325 switch (cancel_region.errno()) {
34013326 .SUCCESS => return,
3402 .INTR, .CANCELED => continue,
3327 .INTR, .CANCELED => {},
34033328 .FAULT => |err| return errnoBug(err),
34043329 .INVAL => |err| return errnoBug(err),
34053330 .ACCES => return error.AccessDenied,
......@@ -3438,7 +3363,7 @@ fn dirReadLink(
34383363 const rc = linux.readlinkat(dir.handle, sub_path_posix, buffer.ptr, buffer.len);
34393364 switch (linux.errno(rc)) {
34403365 .SUCCESS => return @bitCast(rc),
3441 .INTR => continue,
3366 .INTR => {},
34423367 .ACCES => return error.AccessDenied,
34433368 .FAULT => |err| return errnoBug(err),
34443369 .INVAL => return error.NotLink,
......@@ -3628,7 +3553,7 @@ fn fileLength(userdata: ?*anyopaque, file: File) File.LengthError!u64 {
36283553 if (!statx_buf.mask.SIZE) return error.Unexpected;
36293554 return statx_buf.size;
36303555 },
3631 .INTR, .CANCELED => continue,
3556 .INTR, .CANCELED => {},
36323557 .ACCES => |err| return errnoBug(err),
36333558 .BADF => |err| return errnoBug(err), // File descriptor used after closed.
36343559 .FAULT => |err| return errnoBug(err),
......@@ -3811,7 +3736,7 @@ fn fileSync(userdata: ?*anyopaque, file: File) File.SyncError!void {
38113736 ev.yield(null, .nothing);
38123737 switch (cancel_region.errno()) {
38133738 .SUCCESS => return,
3814 .INTR, .CANCELED => continue,
3739 .INTR, .CANCELED => {},
38153740 .BADF => |err| return errnoBug(err),
38163741 .INVAL => |err| return errnoBug(err),
38173742 .ROFS => |err| return errnoBug(err),
......@@ -3833,7 +3758,7 @@ fn fileIsTty(userdata: ?*anyopaque, file: File) Io.Cancelable!bool {
38333758 const rc = linux.ioctl(file.handle, linux.T.IOCGWINSZ, @intFromPtr(&wsz));
38343759 switch (linux.errno(rc)) {
38353760 .SUCCESS => return true,
3836 .INTR => continue,
3761 .INTR => {},
38373762 else => return false,
38383763 }
38393764 }
......@@ -3869,7 +3794,7 @@ fn fileSetLength(userdata: ?*anyopaque, file: File, length: u64) File.SetLengthE
38693794 ev.yield(null, .nothing);
38703795 switch (cancel_region.errno()) {
38713796 .SUCCESS => return,
3872 .INTR, .CANCELED => continue,
3797 .INTR, .CANCELED => {},
38733798 .FBIG => return error.FileTooBig,
38743799 .IO => return error.InputOutput,
38753800 .PERM => return error.PermissionDenied,
......@@ -4051,7 +3976,7 @@ fn fileMemoryMapCreate(
40513976 const rc = linux.mmap(null, options.len, prot, flags, file.handle, casted_offset);
40523977 switch (linux.errno(rc)) {
40533978 .SUCCESS => break @as([*]align(page_align) u8, @ptrFromInt(rc))[0..options.len],
4054 .INTR => continue,
3979 .INTR => {},
40553980 .ACCES => return error.AccessDenied,
40563981 .AGAIN => return error.LockedMemoryLimitExceeded,
40573982 .MFILE => return error.ProcessFdQuotaExceeded,
......@@ -4111,7 +4036,7 @@ fn fileMemoryMapSetLength(
41114036 const rc = linux.mremap(old_memory.ptr, old_memory.len, new_len, flags, addr_hint);
41124037 switch (linux.errno(rc)) {
41134038 .SUCCESS => break @as([*]align(page_align) u8, @ptrFromInt(rc))[0..new_len],
4114 .INTR => continue,
4039 .INTR => {},
41154040 .AGAIN => return error.LockedMemoryLimitExceeded,
41164041 .NOMEM => return error.OutOfMemory,
41174042 .INVAL => |err| return errnoBug(err),
......@@ -4220,7 +4145,7 @@ fn processCurrentPath(userdata: ?*anyopaque, buffer: []u8) process.CurrentPathEr
42204145 try sync.cancel_region.await(.nothing);
42214146 switch (linux.errno(linux.getcwd(buffer.ptr, buffer.len))) {
42224147 .SUCCESS => return std.mem.findScalar(u8, buffer, 0).?,
4223 .INTR => continue,
4148 .INTR => {},
42244149 .NOENT => return error.CurrentDirUnlinked,
42254150 .RANGE => return error.NameTooLong,
42264151 .FAULT => |err| return errnoBug(err),
......@@ -4235,7 +4160,7 @@ fn processSetCurrentDir(userdata: ?*anyopaque, dir: Dir) process.SetCurrentDirEr
42354160 if (dir.handle == linux.AT.FDCWD) return;
42364161 var sync: CancelRegion.Sync = try .init(ev);
42374162 defer sync.deinit(ev);
4238 return ev.fchdir(&sync, dir.handle);
4163 return fchdir(&sync, dir.handle);
42394164}
42404165
42414166fn processSetCurrentPath(userdata: ?*anyopaque, dir_path: []const u8) ChdirError!void {
......@@ -4272,7 +4197,7 @@ fn processReplace(userdata: ?*anyopaque, options: process.ReplaceOptions) proces
42724197
42734198 var sync: CancelRegion.Sync = try .init(ev);
42744199 defer sync.deinit(ev);
4275 return ev.execv(&sync, options.expand_arg0, argv_buf.ptr[0].?, argv_buf.ptr, env_block, PATH);
4200 return execv(&sync, options.expand_arg0, argv_buf.ptr[0].?, argv_buf.ptr, env_block, PATH);
42764201}
42774202
42784203fn processReplacePath(
......@@ -4292,7 +4217,7 @@ fn processSpawn(userdata: ?*anyopaque, options: process.SpawnOptions) process.Sp
42924217 const spawned = try ev.spawn(options);
42934218 var cancel_region: CancelRegion = .initBlocked();
42944219 defer cancel_region.deinit();
4295 defer ev.close(spawned.err_fd);
4220 defer ev.closeAsync(spawned.err_fd);
42964221
42974222 // Wait for the child to report any errors in or before `execvpe`.
42984223 var child_err: ForkBailError = undefined;
......@@ -4434,8 +4359,9 @@ fn spawn(ev: *Evented, options: process.SpawnOptions) process.SpawnError!Spawned
44344359
44354360 if (pid_result == 0) {
44364361 defer comptime unreachable; // We are the child.
4362 // Note that the parent uring is no longer accessible, so we must no longer reference `ev`.
44374363 var sync: CancelRegion.Sync = .{ .cancel_region = .initBlocked() };
4438 const err = ev.setUpChild(&sync, .{
4364 const err = setUpChild(&sync, .{
44394365 .stdin_pipe = stdin_pipe[0],
44404366 .stdout_pipe = stdout_pipe[1],
44414367 .stderr_pipe = stderr_pipe[1],
......@@ -4446,7 +4372,7 @@ fn spawn(ev: *Evented, options: process.SpawnOptions) process.SpawnError!Spawned
44464372 .PATH = PATH,
44474373 .spawn = options,
44484374 });
4449 ev.writeAll(&sync.cancel_region, err_pipe[1], @ptrCast(&err)) catch {};
4375 writeAllSync(&sync, err_pipe[1], @ptrCast(&err)) catch {};
44504376 const exit = if (builtin.single_threaded) linux.exit else linux.exit_group;
44514377 exit(1);
44524378 }
......@@ -4454,13 +4380,13 @@ fn spawn(ev: *Evented, options: process.SpawnOptions) process.SpawnError!Spawned
44544380 const pid: pid_t = @intCast(pid_result); // We are the parent.
44554381 errdefer comptime unreachable; // The child is forked; we must not error from now on
44564382
4457 ev.close(err_pipe[1]); // make sure only the child holds the write end open
4383 ev.closeAsync(err_pipe[1]); // make sure only the child holds the write end open
44584384
4459 if (options.stdin == .pipe) ev.close(stdin_pipe[0]);
4460 if (options.stdout == .pipe) ev.close(stdout_pipe[1]);
4461 if (options.stderr == .pipe) ev.close(stderr_pipe[1]);
4385 if (options.stdin == .pipe) ev.closeAsync(stdin_pipe[0]);
4386 if (options.stdout == .pipe) ev.closeAsync(stdout_pipe[1]);
4387 if (options.stderr == .pipe) ev.closeAsync(stderr_pipe[1]);
44624388
4463 if (prog_pipe[1] != -1) ev.close(prog_pipe[1]);
4389 if (prog_pipe[1] != -1) ev.closeAsync(prog_pipe[1]);
44644390
44654391 options.progress_node.setIpcFile(ev, .{ .handle = prog_pipe[0], .flags = .{ .nonblocking = true } });
44664392
......@@ -4497,14 +4423,14 @@ pub fn pipe2(flags: linux.O) PipeError![2]fd_t {
44974423 }
44984424}
44994425fn destroyPipe(ev: *Evented, pipe: [2]fd_t) void {
4500 if (pipe[0] != -1) ev.close(pipe[0]);
4501 if (pipe[0] != pipe[1]) ev.close(pipe[1]);
4426 if (pipe[0] != -1) ev.closeAsync(pipe[0]);
4427 if (pipe[0] != pipe[1]) ev.closeAsync(pipe[1]);
45024428}
45034429
45044430/// Errors that can occur between fork() and execv()
45054431const ForkBailError = process.SetCurrentDirError || ChdirError ||
45064432 process.SpawnError || process.ReplaceError;
4507fn setUpChild(ev: *Evented, sync: *CancelRegion.Sync, options: struct {
4433fn setUpChild(sync: *CancelRegion.Sync, options: struct {
45084434 stdin_pipe: fd_t,
45094435 stdout_pipe: fd_t,
45104436 stderr_pipe: fd_t,
......@@ -4515,21 +4441,21 @@ fn setUpChild(ev: *Evented, sync: *CancelRegion.Sync, options: struct {
45154441 PATH: []const u8,
45164442 spawn: process.SpawnOptions,
45174443}) ForkBailError {
4518 try ev.setUpChildIo(
4444 try setUpChildIo(
45194445 sync,
45204446 options.spawn.stdin,
45214447 options.stdin_pipe,
45224448 linux.STDIN_FILENO,
45234449 options.dev_null_fd,
45244450 );
4525 try ev.setUpChildIo(
4451 try setUpChildIo(
45264452 sync,
45274453 options.spawn.stdout,
45284454 options.stdout_pipe,
45294455 linux.STDOUT_FILENO,
45304456 options.dev_null_fd,
45314457 );
4532 try ev.setUpChildIo(
4458 try setUpChildIo(
45334459 sync,
45344460 options.spawn.stderr,
45354461 options.stderr_pipe,
......@@ -4539,17 +4465,17 @@ fn setUpChild(ev: *Evented, sync: *CancelRegion.Sync, options: struct {
45394465
45404466 switch (options.spawn.cwd) {
45414467 .inherit => {},
4542 .dir => |cwd_dir| try ev.fchdir(sync, cwd_dir.handle),
4468 .dir => |cwd_dir| try fchdir(sync, cwd_dir.handle),
45434469 .path => |cwd_path| {
45444470 var cwd_path_buffer: [PATH_MAX]u8 = undefined;
45454471 const cwd_path_posix = try pathToPosix(cwd_path, &cwd_path_buffer);
4546 try ev.chdir(sync, cwd_path_posix);
4472 try chdir(sync, cwd_path_posix);
45474473 },
45484474 }
45494475
45504476 // Must happen after fchdir above, the cwd file descriptor might be
45514477 // equal to prog_fileno and be clobbered by this dup2 call.
4552 if (options.prog_pipe != -1) try ev.dup2(sync, options.prog_pipe, prog_fileno);
4478 if (options.prog_pipe != -1) try dup2(sync, options.prog_pipe, prog_fileno);
45534479
45544480 if (options.spawn.gid) |gid| {
45554481 switch (linux.errno(linux.setregid(gid, gid))) {
......@@ -4589,7 +4515,7 @@ fn setUpChild(ev: *Evented, sync: *CancelRegion.Sync, options: struct {
45894515 }
45904516 }
45914517
4592 return ev.execv(
4518 return execv(
45934519 sync,
45944520 options.spawn.expand_arg0,
45954521 options.argv_buf.ptr[0].?,
......@@ -4600,7 +4526,6 @@ fn setUpChild(ev: *Evented, sync: *CancelRegion.Sync, options: struct {
46004526}
46014527
46024528fn setUpChildIo(
4603 ev: *Evented,
46044529 sync: *CancelRegion.Sync,
46054530 stdio: process.SpawnOptions.StdIo,
46064531 pipe_fd: fd_t,
......@@ -4608,13 +4533,13 @@ fn setUpChildIo(
46084533 dev_null_fd: fd_t,
46094534) !void {
46104535 switch (stdio) {
4611 .pipe => try ev.dup2(sync, pipe_fd, std_fileno),
4536 .pipe => try dup2(sync, pipe_fd, std_fileno),
46124537 .close => _ = linux.close(std_fileno),
46134538 .inherit => {},
4614 .ignore => try ev.dup2(sync, dev_null_fd, std_fileno),
4539 .ignore => try dup2(sync, dev_null_fd, std_fileno),
46154540 .file => |file| {
46164541 if (file.flags.nonblocking) @panic("TODO implement setUpChildIo when nonblocking file is used");
4617 try ev.dup2(sync, file.handle, std_fileno);
4542 try dup2(sync, file.handle, std_fileno);
46184543 },
46194544 }
46204545}
......@@ -4623,13 +4548,12 @@ pub const DupError = error{
46234548 ProcessFdQuotaExceeded,
46244549 SystemResources,
46254550} || Io.UnexpectedError || Io.Cancelable;
4626pub fn dup2(ev: *Evented, sync: *CancelRegion.Sync, old_fd: fd_t, new_fd: fd_t) DupError!void {
4627 _ = ev;
4551pub fn dup2(sync: *CancelRegion.Sync, old_fd: fd_t, new_fd: fd_t) DupError!void {
46284552 while (true) {
46294553 try sync.cancel_region.await(.nothing);
46304554 switch (linux.errno(linux.dup2(old_fd, new_fd))) {
46314555 .SUCCESS => {},
4632 .BUSY, .INTR => continue,
4556 .BUSY, .INTR => {},
46334557 .INVAL => |err| return errnoBug(err), // invalid parameters
46344558 .BADF => |err| return errnoBug(err), // use after free
46354559 .MFILE => return error.ProcessFdQuotaExceeded,
......@@ -4640,7 +4564,6 @@ pub fn dup2(ev: *Evented, sync: *CancelRegion.Sync, old_fd: fd_t, new_fd: fd_t)
46404564}
46414565
46424566fn execv(
4643 ev: *Evented,
46444567 sync: *CancelRegion.Sync,
46454568 arg0_expand: process.ArgExpansion,
46464569 file: [*:0]const u8,
......@@ -4649,7 +4572,8 @@ fn execv(
46494572 PATH: []const u8,
46504573) process.ReplaceError {
46514574 const file_slice = std.mem.sliceTo(file, 0);
4652 if (std.mem.findScalar(u8, file_slice, '/') != null) return ev.execvPath(sync, file, child_argv, env_block);
4575 if (std.mem.findScalar(u8, file_slice, '/') != null)
4576 return execvPath(sync, file, child_argv, env_block);
46534577
46544578 // Use of PATH_MAX here is valid as the path_buf will be passed
46554579 // directly to the operating system in posixExecvPath.
......@@ -4677,7 +4601,7 @@ fn execv(
46774601 .expand => child_argv[0] = full_path,
46784602 .no_expand => {},
46794603 }
4680 err = ev.execvPath(sync, full_path, child_argv, env_block);
4604 err = execvPath(sync, full_path, child_argv, env_block);
46814605 switch (err) {
46824606 error.AccessDenied => seen_eacces = true,
46834607 error.FileNotFound, error.NotDir => {},
......@@ -4689,13 +4613,11 @@ fn execv(
46894613}
46904614/// This function ignores PATH environment variable.
46914615pub fn execvPath(
4692 ev: *Evented,
46934616 sync: *CancelRegion.Sync,
46944617 path: [*:0]const u8,
46954618 child_argv: [*:null]const ?[*:0]const u8,
46964619 env_block: process.Environ.PosixBlock,
46974620) process.ReplaceError {
4698 _ = ev;
46994621 try sync.cancel_region.await(.nothing);
47004622 switch (linux.errno(linux.execve(path, child_argv, env_block.slice.ptr))) {
47014623 .FAULT => |err| return errnoBug(err), // Bad pointer parameter.
......@@ -4766,7 +4688,7 @@ fn childWait(userdata: ?*anyopaque, child: *process.Child) process.Child.WaitErr
47664688 child.resource_usage_statistics.rusage = rusage;
47674689 break;
47684690 },
4769 .INTR, .CANCELED => continue,
4691 .INTR, .CANCELED => {},
47704692 .CHILD => |err| return errnoBug(err), // Double-free.
47714693 else => |err| return unexpectedErrno(err),
47724694 }
......@@ -4781,7 +4703,7 @@ fn childWait(userdata: ?*anyopaque, child: *process.Child) process.Child.WaitErr
47814703 _, .CONTINUED => .{ .unknown = status },
47824704 };
47834705 },
4784 .INTR, .CANCELED => continue,
4706 .INTR, .CANCELED => {},
47854707 .CHILD => |err| return errnoBug(err), // Double-free.
47864708 else => |err| return unexpectedErrno(err),
47874709 }
......@@ -4798,7 +4720,7 @@ fn childKill(userdata: ?*anyopaque, child: *process.Child) void {
47984720 const pid = child.id.?;
47994721 while (true) switch (linux.errno(linux.kill(pid, .TERM))) {
48004722 .SUCCESS => break,
4801 .INTR => continue,
4723 .INTR => {},
48024724 .PERM => return,
48034725 .INVAL => |err| return errnoBug(err) catch {},
48044726 .SRCH => |err| return errnoBug(err) catch {},
......@@ -4830,7 +4752,7 @@ fn childKill(userdata: ?*anyopaque, child: *process.Child) void {
48304752 ev.yield(null, .nothing);
48314753 switch (maybe_sync.cancel_region.errno()) {
48324754 .SUCCESS => return,
4833 .INTR, .CANCELED => continue,
4755 .INTR, .CANCELED => {},
48344756 .CHILD => |err| return errnoBug(err) catch {}, // Double-free.
48354757 else => |err| return unexpectedErrno(err) catch {},
48364758 }
......@@ -4839,15 +4761,15 @@ fn childKill(userdata: ?*anyopaque, child: *process.Child) void {
48394761
48404762fn childCleanup(ev: *Evented, child: *process.Child) void {
48414763 if (child.stdin) |*stdin| {
4842 ev.close(stdin.handle);
4764 ev.closeAsync(stdin.handle);
48434765 child.stdin = null;
48444766 }
48454767 if (child.stdout) |*stdout| {
4846 ev.close(stdout.handle);
4768 ev.closeAsync(stdout.handle);
48474769 child.stdout = null;
48484770 }
48494771 if (child.stderr) |*stderr| {
4850 ev.close(stderr.handle);
4772 ev.closeAsync(stderr.handle);
48514773 child.stderr = null;
48524774 }
48534775 child.id = null;
......@@ -4953,14 +4875,10 @@ fn sleep(userdata: ?*anyopaque, timeout: Io.Timeout) Io.Cancelable!void {
49534875 .resv = 0,
49544876 };
49554877 ev.yield(null, .nothing);
4956 switch (cancel_region.errno()) {
4957 // Handles SUCCESS as well as clock not available and unexpected
4958 // errors. The user had a chance to check clock resolution before
4959 // getting here, which would have reported 0, making this a legal
4960 // amount of time to sleep.
4961 else => return,
4962 .INTR, .CANCELED => return error.Canceled,
4963 }
4878 // Handles SUCCESS as well as clock not available and unexpected
4879 // errors. The user had a chance to check clock resolution before
4880 // getting here, which would have reported 0, making this a legal
4881 // amount of time to sleep.
49644882}
49654883
49664884fn random(userdata: ?*anyopaque, buffer: []u8) void {
......@@ -5037,7 +4955,7 @@ fn netBindIp(
50374955 var maybe_sync: CancelRegion.Sync.Maybe = .{ .cancel_region = .init() };
50384956 defer maybe_sync.deinit(ev);
50394957 const socket_fd = try ev.socket(&maybe_sync.cancel_region, family, options);
5040 errdefer ev.close(socket_fd);
4958 errdefer ev.closeAsync(socket_fd);
50414959 var storage: PosixAddress = undefined;
50424960 var addr_len = addressToPosix(address, &storage);
50434961 try ev.bind(&maybe_sync.cancel_region, socket_fd, &storage.any, addr_len);
......@@ -5204,23 +5122,19 @@ fn netReceive(
52045122 .data = data,
52055123 .control = if (msg.control) |ptr| @as([*]u8, @ptrCast(ptr))[0..msg.controllen] else message.control,
52065124 .flags = .{
5207 .eor = (msg.flags & linux.MSG.EOR) != 0,
5208 .trunc = (msg.flags & linux.MSG.TRUNC) != 0,
5209 .ctrunc = (msg.flags & linux.MSG.CTRUNC) != 0,
5210 .oob = (msg.flags & linux.MSG.OOB) != 0,
5211 .errqueue = if (@hasDecl(linux.MSG, "ERRQUEUE")) (msg.flags & linux.MSG.ERRQUEUE) != 0 else false,
5125 .eor = msg.flags & linux.MSG.EOR != 0,
5126 .trunc = msg.flags & linux.MSG.TRUNC != 0,
5127 .ctrunc = msg.flags & linux.MSG.CTRUNC != 0,
5128 .oob = msg.flags & linux.MSG.OOB != 0,
5129 .errqueue = msg.flags & linux.MSG.ERRQUEUE != 0,
52125130 },
52135131 };
52145132 message_i += 1;
52155133 continue;
52165134 },
52175135 .AGAIN => unreachable,
5218 .INTR, .CANCELED => {
5219 if (deadline) |d| {
5220 if (now(ev, d.clock).nanoseconds >= d.raw.nanoseconds) return .{ error.Timeout, message_i };
5221 }
5222 continue;
5223 },
5136 .INTR, .CANCELED => if (deadline) |d| if (now(ev, d.clock).nanoseconds >= d.raw.nanoseconds)
5137 return .{ error.Timeout, message_i },
52245138
52255139 .BADF => |err| return .{ errnoBug(err), message_i },
52265140 .NFILE => return .{ error.SystemFdQuotaExceeded, message_i },
......@@ -5323,7 +5237,7 @@ fn netShutdown(
53235237 ev.yield(null, .nothing);
53245238 switch (cancel_region.errno()) {
53255239 .SUCCESS => return,
5326 .INTR, .CANCELED => continue,
5240 .INTR, .CANCELED => {},
53275241 .BADF, .NOTSOCK, .INVAL => |err| return errnoBug(err),
53285242 .NOTCONN => return error.SocketUnconnected,
53295243 .NOBUFS => return error.SystemResources,
......@@ -5393,7 +5307,7 @@ fn bind(
53935307 ev.yield(null, .nothing);
53945308 switch (cancel_region.errno()) {
53955309 .SUCCESS => return,
5396 .INTR, .CANCELED => continue,
5310 .INTR, .CANCELED => {},
53975311 .ADDRINUSE => return error.AddressInUse,
53985312 .BADF => |err| return errnoBug(err), // File descriptor used after closed.
53995313 .INVAL => |err| return errnoBug(err), // invalid parameters
......@@ -5407,13 +5321,12 @@ fn bind(
54075321 }
54085322}
54095323
5410fn chdir(ev: *Evented, sync: *CancelRegion.Sync, path: [*:0]const u8) ChdirError!void {
5411 _ = ev;
5324fn chdir(sync: *CancelRegion.Sync, path: [*:0]const u8) ChdirError!void {
54125325 while (true) {
54135326 try sync.cancel_region.await(.nothing);
54145327 switch (linux.errno(linux.chdir(path))) {
54155328 .SUCCESS => return,
5416 .INTR => continue,
5329 .INTR => {},
54175330 .ACCES => return error.AccessDenied,
54185331 .IO => return error.FileSystem,
54195332 .LOOP => return error.SymLinkLoop,
......@@ -5429,6 +5342,36 @@ fn chdir(ev: *Evented, sync: *CancelRegion.Sync, path: [*:0]const u8) ChdirError
54295342}
54305343
54315344fn close(ev: *Evented, fd: fd_t) void {
5345 var cancel_region: CancelRegion = .initBlocked();
5346 defer cancel_region.deinit();
5347 const thread = cancel_region.awaitIoUring() catch |err| switch (err) {
5348 error.Canceled => unreachable, // blocked
5349 };
5350 thread.enqueue().* = .{
5351 .opcode = .CLOSE,
5352 .flags = 0,
5353 .ioprio = 0,
5354 .fd = fd,
5355 .off = 0,
5356 .addr = 0,
5357 .len = 0,
5358 .rw_flags = 0,
5359 .user_data = @intFromPtr(cancel_region.fiber),
5360 .buf_index = 0,
5361 .personality = 0,
5362 .splice_fd_in = 0,
5363 .addr3 = 0,
5364 .resv = 0,
5365 };
5366 ev.yield(null, .nothing);
5367 switch (cancel_region.errno()) {
5368 .BADF => recoverableOsBugDetected(), // Always a race condition.
5369 .INTR => {}, // This is still a success. See https://github.com/ziglang/zig/issues/2425
5370 else => {},
5371 }
5372}
5373
5374fn closeAsync(ev: *Evented, fd: fd_t) void {
54325375 _ = ev;
54335376 const thread: *Thread = .current();
54345377 thread.enqueue().* = .{
......@@ -5449,14 +5392,13 @@ fn close(ev: *Evented, fd: fd_t) void {
54495392 };
54505393}
54515394
5452fn fchdir(ev: *Evented, sync: *CancelRegion.Sync, dir: fd_t) process.SetCurrentDirError!void {
5453 _ = ev;
5395fn fchdir(sync: *CancelRegion.Sync, dir: fd_t) process.SetCurrentDirError!void {
54545396 if (dir == linux.AT.FDCWD) return;
54555397 while (true) {
54565398 try sync.cancel_region.await(.nothing);
54575399 switch (linux.errno(linux.fchdir(dir))) {
54585400 .SUCCESS => return,
5459 .INTR => continue,
5401 .INTR => {},
54605402 .ACCES => return error.AccessDenied,
54615403 .NOTDIR => return error.NotDir,
54625404 .IO => return error.FileSystem,
......@@ -5479,7 +5421,7 @@ fn fchmodat(
54795421 try sync.cancel_region.await(.nothing);
54805422 switch (linux.errno(linux.fchmodat2(dir, path, mode, flags))) {
54815423 .SUCCESS => return,
5482 .INTR => continue,
5424 .INTR => {},
54835425 .BADF => |err| return errnoBug(err),
54845426 .FAULT => |err| return errnoBug(err),
54855427 .INVAL => |err| return errnoBug(err),
......@@ -5511,7 +5453,7 @@ fn fchownat(
55115453 try sync.cancel_region.await(.nothing);
55125454 switch (linux.errno(linux.fchownat(dir, path, owner, group, flags))) {
55135455 .SUCCESS => return,
5514 .INTR => continue,
5456 .INTR => {},
55155457 .BADF => |err| return errnoBug(err), // likely fd refers to directory opened without `Dir.OpenOptions.iterate`
55165458 .FAULT => |err| return errnoBug(err),
55175459 .INVAL => |err| return errnoBug(err),
......@@ -5543,7 +5485,7 @@ fn flock(
55435485 .exclusive => LOCK.EX,
55445486 })))) {
55455487 .SUCCESS => return,
5546 .INTR => continue,
5488 .INTR => {},
55475489 .BADF => |err| return errnoBug(err),
55485490 .INVAL => |err| return errnoBug(err), // invalid parameters
55495491 .NOLCK => return error.SystemResources,
......@@ -5593,7 +5535,7 @@ fn getsockname(
55935535 try sync.cancel_region.await(.nothing);
55945536 switch (linux.errno(linux.getsockname(socket_fd, addr, addr_len))) {
55955537 .SUCCESS => return,
5596 .INTR => continue,
5538 .INTR => {},
55975539 .BADF => |err| return errnoBug(err), // File descriptor used after closed.
55985540 .FAULT => |err| return errnoBug(err),
55995541 .INVAL => |err| return errnoBug(err), // invalid parameters
......@@ -5634,7 +5576,7 @@ fn linkat(
56345576 ev.yield(null, .nothing);
56355577 switch (cancel_region.errno()) {
56365578 .SUCCESS => return,
5637 .INTR, .CANCELED => continue,
5579 .INTR, .CANCELED => {},
56385580 .ACCES => return error.AccessDenied,
56395581 .DQUOT => return error.DiskQuota,
56405582 .EXIST => return error.PathAlreadyExists,
......@@ -5674,7 +5616,7 @@ fn lseek(
56745616 8 => linux.lseek(fd, @bitCast(offset), whence),
56755617 })) {
56765618 .SUCCESS => return,
5677 .INTR => continue,
5619 .INTR => {},
56785620 .BADF => |err| return errnoBug(err), // File descriptor used after closed.
56795621 .INVAL => return error.Unseekable,
56805622 .OVERFLOW => return error.Unseekable,
......@@ -5717,7 +5659,7 @@ fn openat(
57175659 const completion = cancel_region.completion();
57185660 switch (completion.errno()) {
57195661 .SUCCESS => return completion.result,
5720 .INTR, .CANCELED => continue,
5662 .INTR, .CANCELED => {},
57215663 .FAULT => |err| return errnoBug(err),
57225664 .INVAL => return error.BadPathName,
57235665 .BADF => |err| return errnoBug(err), // File descriptor used after closed.
......@@ -5779,7 +5721,7 @@ fn preadv(
57795721 const completion = cancel_region.completion();
57805722 switch (completion.errno()) {
57815723 .SUCCESS => return @as(u32, @bitCast(completion.result)),
5782 .INTR, .CANCELED => continue,
5724 .INTR, .CANCELED => {},
57835725 .INVAL => |err| return errnoBug(err),
57845726 .FAULT => |err| return errnoBug(err),
57855727 .AGAIN => return error.WouldBlock,
......@@ -5826,7 +5768,7 @@ fn pwritev(
58265768 const completion = cancel_region.completion();
58275769 switch (completion.errno()) {
58285770 .SUCCESS => return @as(u32, @bitCast(completion.result)),
5829 .INTR, .CANCELED => continue,
5771 .INTR, .CANCELED => {},
58305772 .INVAL => |err| return errnoBug(err),
58315773 .FAULT => |err| return errnoBug(err),
58325774 .AGAIN => return error.WouldBlock,
......@@ -5876,7 +5818,7 @@ fn realPath(
58765818 const rc = linux.readlink(proc_path, out_buffer.ptr, out_buffer.len);
58775819 switch (linux.errno(rc)) {
58785820 .SUCCESS => return rc,
5879 .INTR => continue,
5821 .INTR => {},
58805822 .ACCES => return error.AccessDenied,
58815823 .FAULT => |err| return errnoBug(err),
58825824 .IO => return error.FileSystem,
......@@ -5921,7 +5863,7 @@ fn renameat(
59215863 ev.yield(null, .nothing);
59225864 switch (cancel_region.errno()) {
59235865 .SUCCESS => return,
5924 .INTR, .CANCELED => continue,
5866 .INTR, .CANCELED => {},
59255867 .ACCES => return error.AccessDenied,
59265868 .PERM => return error.PermissionDenied,
59275869 .BUSY => return error.FileBusy,
......@@ -5988,7 +5930,7 @@ fn setsockopt(
59885930 ev.yield(null, .nothing);
59895931 switch (cancel_region.errno()) {
59905932 .SUCCESS => return,
5991 .INTR, .CANCELED => continue,
5933 .INTR, .CANCELED => {},
59925934 .BADF => |err| return errnoBug(err), // File descriptor used after closed.
59935935 .NOTSOCK => |err| return errnoBug(err),
59945936 .INVAL => |err| return errnoBug(err),
......@@ -6039,7 +5981,7 @@ fn socket(
60395981 const completion = cancel_region.completion();
60405982 switch (completion.errno()) {
60415983 .SUCCESS => break completion.result,
6042 .INTR, .CANCELED => continue,
5984 .INTR, .CANCELED => {},
60435985 .AFNOSUPPORT => return error.AddressFamilyUnsupported,
60445986 .INVAL => return error.ProtocolUnsupportedBySystem,
60455987 .MFILE => return error.ProcessFdQuotaExceeded,
......@@ -6051,7 +5993,7 @@ fn socket(
60515993 else => |err| return unexpectedErrno(err),
60525994 }
60535995 };
6054 errdefer ev.close(socket_fd);
5996 errdefer ev.closeAsync(socket_fd);
60555997
60565998 if (options.ip6_only) {
60575999 if (linux.IPV6 == void) return error.OptionUnsupported;
......@@ -6101,7 +6043,7 @@ fn statx(
61016043 ev.yield(null, .nothing);
61026044 switch (cancel_region.errno()) {
61036045 .SUCCESS => return statFromLinux(&statx_buf),
6104 .INTR, .CANCELED => continue,
6046 .INTR, .CANCELED => {},
61056047 .ACCES => return error.AccessDenied,
61066048 .BADF => |err| return errnoBug(err), // File descriptor used after closed.
61076049 .FAULT => |err| return errnoBug(err),
......@@ -6140,7 +6082,7 @@ fn utimensat(
61406082 try sync.cancel_region.await(.nothing);
61416083 switch (linux.errno(linux.utimensat(dir, path, times, flags))) {
61426084 .SUCCESS => return,
6143 .INTR => continue,
6085 .INTR => {},
61446086 .BADF => |err| return errnoBug(err), // always a race condition
61456087 .FAULT => |err| return errnoBug(err),
61466088 .INVAL => |err| return errnoBug(err),
......@@ -6152,19 +6094,33 @@ fn utimensat(
61526094 }
61536095}
61546096
6155fn writeAll(
6156 ev: *Evented,
6157 cancel_region: *CancelRegion,
6158 fd: fd_t,
6159 buffer: []const u8,
6160) (File.Writer.Error || error{EndOfStream})!void {
6097fn writeAllSync(sync: *CancelRegion.Sync, fd: fd_t, buffer: []const u8) File.Writer.Error!void {
61616098 var index: usize = 0;
6162 while (buffer.len - index != 0) {
6163 const len = try ev.pwritev(cancel_region, fd, &.{
6164 .{ .base = buffer[index..].ptr, .len = buffer.len - index },
6165 }, null);
6166 if (len == 0) return error.EndOfStream;
6167 index += len;
6099 while (buffer.len - index != 0) index += try writeSync(sync, fd, buffer[index..]);
6100}
6101
6102fn writeSync(sync: *CancelRegion.Sync, fd: fd_t, buffer: []const u8) File.Writer.Error!usize {
6103 while (true) {
6104 try sync.cancel_region.await(.nothing);
6105 const rc = linux.write(fd, buffer.ptr, buffer.len);
6106 switch (linux.errno(rc)) {
6107 .SUCCESS => return @intCast(rc),
6108 .INTR => {},
6109 .INVAL => |err| return errnoBug(err),
6110 .FAULT => |err| return errnoBug(err),
6111 .AGAIN => return error.WouldBlock,
6112 .BADF => return error.NotOpenForWriting, // Can be a race condition.
6113 .DESTADDRREQ => |err| return errnoBug(err), // `connect` was never called.
6114 .DQUOT => return error.DiskQuota,
6115 .FBIG => return error.FileTooBig,
6116 .IO => return error.InputOutput,
6117 .NOSPC => return error.NoSpaceLeft,
6118 .PERM => return error.PermissionDenied,
6119 .PIPE => return error.BrokenPipe,
6120 .CONNRESET => |err| return errnoBug(err), // Not a socket handle.
6121 .BUSY => return error.DeviceBusy,
6122 else => |err| return unexpectedErrno(err),
6123 }
61686124 }
61696125}
61706126
lib/std/Io/test.zig-34
......@@ -282,40 +282,6 @@ test "Group.concurrent" {
282282 try testing.expectEqualSlices(usize, &.{ 45, 245 }, &results);
283283}
284284
285test "select" {
286 const io = testing.io;
287
288 var queue: Io.Queue(u8) = .init(&.{});
289
290 var get_a = io.concurrent(Io.Queue(u8).getOne, .{ &queue, io }) catch |err| switch (err) {
291 error.ConcurrencyUnavailable => {
292 try testing.expect(builtin.single_threaded);
293 return;
294 },
295 };
296 defer _ = get_a.cancel(io) catch {};
297
298 var get_b = try io.concurrent(Io.Queue(u8).getOne, .{ &queue, io });
299 defer _ = get_b.cancel(io) catch {};
300
301 var timeout = io.async(Io.sleep, .{ io, .fromMilliseconds(1), .awake });
302 defer timeout.cancel(io) catch {};
303
304 switch (try io.select(.{
305 .get_a = &get_a,
306 .get_b = &get_b,
307 .timeout = &timeout,
308 })) {
309 .get_a => return error.TestFailure,
310 .get_b => return error.TestFailure,
311 .timeout => {
312 queue.close(io);
313 try testing.expectError(error.Closed, get_a.await(io));
314 try testing.expectError(error.Closed, get_b.await(io));
315 },
316 }
317}
318
319285fn testQueue(comptime len: usize) !void {
320286 const io = testing.io;
321287 var buf: [len]usize = undefined;
lib/std/os/linux/IoUring.zig+4
......@@ -201,6 +201,10 @@ pub fn enter(self: *IoUring, to_submit: u32, min_complete: u32, flags: u32) !u32
201201 // The kernel believes our `self.fd` does not refer to an io_uring instance,
202202 // or the opcode is valid but not supported by this kernel (more likely):
203203 .OPNOTSUPP => return error.OpcodeNotSupported,
204 // The thread submitting the work is invalid. This may occur if IORING_ENTER_GETEVENTS
205 // and IORING_SETUP_DEFER_TASKRUN is set, but the submitting thread is not the thread
206 // that initially created or enabled the io_uring associated with fd.
207 .EXIST => return error.InvalidThread,
204208 // The operation was interrupted by a delivery of a signal before it could complete.
205209 // This can happen while waiting for events with IORING_ENTER_GETEVENTS:
206210 .INTR => return error.SignalInterrupt,