authorgravatar for andrew@ziglang.orgAndrew Kelley <andrew@ziglang.org> 2026-01-07 18:00:36-08:00
committergravatar for andrew@ziglang.orgAndrew Kelley <andrew@ziglang.org> 2026-01-07 18:00:36-08:00
logce890060350a848590a7371f11316c2040a89a2b
treea45b1b40f192f43d7b7678bba1d092af03af2f50
parentc0092f5394b918f4a4407b73820574d18edfa1dd

std.Io.Kqueue: fix bitrot


2 files changed, 91 insertions(+), 181 deletions(-)

lib/std/Io.zig+1-1
......@@ -1358,7 +1358,7 @@ pub const Mutex = extern struct {
13581358
13591359 pub const init: Mutex = .{ .state = .init(.unlocked) };
13601360
1361 const State = enum(u32) {
1361 pub const State = enum(u32) {
13621362 unlocked,
13631363 locked_once,
13641364 contended,
lib/std/Io/Kqueue.zig+90-180
......@@ -239,7 +239,7 @@ pub const CreateFileDescriptorError = error{
239239 ProcessFdQuotaExceeded,
240240 /// The system-wide limit on the total number of open files has been reached.
241241 SystemFdQuotaExceeded,
242} || Io.Unexpected;
242} || Io.UnexpectedError;
243243
244244pub fn createFileDescriptor() CreateFileDescriptorError!posix.fd_t {
245245 const rc = posix.system.kqueue();
......@@ -494,14 +494,6 @@ const SwitchMessage = struct {
494494 recycle: *Fiber,
495495 register_awaiter: *?*Fiber,
496496 register_select: []const *Io.AnyFuture,
497 mutex_lock: struct {
498 prev_state: Io.Mutex.State,
499 mutex: *Io.Mutex,
500 },
501 condition_wait: struct {
502 cond: *Io.Condition,
503 mutex: *Io.Mutex,
504 },
505497 exit,
506498 };
507499
......@@ -537,59 +529,6 @@ const SwitchMessage = struct {
537529 }
538530 }
539531 },
540 .mutex_lock => |mutex_lock| {
541 const prev_fiber: *Fiber = @alignCast(@fieldParentPtr("context", message.contexts.prev));
542 assert(prev_fiber.queue_next == null);
543 var prev_state = mutex_lock.prev_state;
544 while (switch (prev_state) {
545 else => next_state: {
546 prev_fiber.queue_next = @ptrFromInt(@intFromEnum(prev_state));
547 break :next_state @cmpxchgWeak(
548 Io.Mutex.State,
549 &mutex_lock.mutex.state,
550 prev_state,
551 @enumFromInt(@intFromPtr(prev_fiber)),
552 .release,
553 .acquire,
554 );
555 },
556 .unlocked => @cmpxchgWeak(
557 Io.Mutex.State,
558 &mutex_lock.mutex.state,
559 .unlocked,
560 .locked_once,
561 .acquire,
562 .acquire,
563 ) orelse {
564 prev_fiber.queue_next = null;
565 k.schedule(thread, .{ .head = prev_fiber, .tail = prev_fiber });
566 return;
567 },
568 }) |next_state| prev_state = next_state;
569 },
570 .condition_wait => |condition_wait| {
571 const prev_fiber: *Fiber = @alignCast(@fieldParentPtr("context", message.contexts.prev));
572 assert(prev_fiber.queue_next == null);
573 const cond_impl = prev_fiber.resultPointer(Condition);
574 cond_impl.* = .{
575 .tail = prev_fiber,
576 .event = .queued,
577 };
578 if (@cmpxchgStrong(
579 ?*Fiber,
580 @as(*?*Fiber, @ptrCast(&condition_wait.cond.state)),
581 null,
582 prev_fiber,
583 .release,
584 .acquire,
585 )) |waiting_fiber| {
586 const waiting_cond_impl = waiting_fiber.?.resultPointer(Condition);
587 assert(waiting_cond_impl.tail.queue_next == null);
588 waiting_cond_impl.tail.queue_next = prev_fiber;
589 waiting_cond_impl.tail = prev_fiber;
590 }
591 condition_wait.mutex.unlock(k.io());
592 },
593532 .exit => for (k.threads.allocated[0..@atomicLoad(u32, &k.threads.active, .acquire)]) |*each_thread| {
594533 const changes = [_]posix.Kevent{
595534 .{
......@@ -878,21 +817,13 @@ pub fn io(k: *Kqueue) Io {
878817 .concurrent = concurrent,
879818 .await = await,
880819 .cancel = cancel,
881 .cancelRequested = cancelRequested,
882820 .select = select,
883821
884822 .groupAsync = groupAsync,
885 .groupWait = groupWait,
823 .groupConcurrent = groupConcurrent,
824 .groupAwait = groupAwait,
886825 .groupCancel = groupCancel,
887826
888 .mutexLock = mutexLock,
889 .mutexLockUncancelable = mutexLockUncancelable,
890 .mutexUnlock = mutexUnlock,
891
892 .conditionWait = conditionWait,
893 .conditionWaitUncancelable = conditionWaitUncancelable,
894 .conditionWake = conditionWake,
895
896827 .dirCreateDir = dirCreateDir,
897828 .dirCreateDirPath = dirCreateDirPath,
898829 .dirCreateDirPathOpen = dirCreateDirPathOpen,
......@@ -912,7 +843,6 @@ pub fn io(k: *Kqueue) Io {
912843 .fileReadPositional = fileReadPositional,
913844 .fileSeekBy = fileSeekBy,
914845 .fileSeekTo = fileSeekTo,
915 .openExecutable = openExecutable,
916846
917847 .now = now,
918848 .sleep = sleep,
......@@ -1037,139 +967,111 @@ fn cancelRequested(userdata: ?*anyopaque) bool {
1037967
1038968fn groupAsync(
1039969 userdata: ?*anyopaque,
1040 group: *Io.Group,
970 type_erased: *Io.Group,
1041971 context: []const u8,
1042 context_alignment: std.mem.Alignment,
1043 start: *const fn (*Io.Group, context: *const anyopaque) void,
972 context_alignment: Alignment,
973 start: *const fn (context: *const anyopaque) Io.Cancelable!void,
1044974) void {
1045975 const k: *Kqueue = @ptrCast(@alignCast(userdata));
1046976 _ = k;
1047 _ = group;
977 _ = type_erased;
1048978 _ = context;
1049979 _ = context_alignment;
1050980 _ = start;
1051981 @panic("TODO");
1052982}
1053983
1054fn groupWait(userdata: ?*anyopaque, group: *Io.Group, token: *anyopaque) void {
1055 const k: *Kqueue = @ptrCast(@alignCast(userdata));
1056 _ = k;
1057 _ = group;
1058 _ = token;
1059 @panic("TODO");
1060}
1061
1062fn groupCancel(userdata: ?*anyopaque, group: *Io.Group, token: *anyopaque) void {
984fn groupConcurrent(
985 userdata: ?*anyopaque,
986 type_erased: *Io.Group,
987 context: []const u8,
988 context_alignment: Alignment,
989 start: *const fn (context: *const anyopaque) Io.Cancelable!void,
990) Io.ConcurrentError!void {
1063991 const k: *Kqueue = @ptrCast(@alignCast(userdata));
1064992 _ = k;
1065 _ = group;
1066 _ = token;
993 _ = type_erased;
994 _ = context;
995 _ = context_alignment;
996 _ = start;
1067997 @panic("TODO");
1068998}
1069999
1070fn select(userdata: ?*anyopaque, futures: []const *Io.AnyFuture) Io.Cancelable!usize {
1000fn groupAwait(userdata: ?*anyopaque, type_erased: *Io.Group, initial_token: *anyopaque) Io.Cancelable!void {
10711001 const k: *Kqueue = @ptrCast(@alignCast(userdata));
10721002 _ = k;
1073 _ = futures;
1003 _ = type_erased;
1004 _ = initial_token;
10741005 @panic("TODO");
10751006}
10761007
1077fn mutexLock(userdata: ?*anyopaque, prev_state: Io.Mutex.State, mutex: *Io.Mutex) Io.Cancelable!void {
1078 const k: *Kqueue = @ptrCast(@alignCast(userdata));
1079 _ = k;
1080 _ = prev_state;
1081 _ = mutex;
1082 @panic("TODO");
1083}
1084fn mutexLockUncancelable(userdata: ?*anyopaque, prev_state: Io.Mutex.State, mutex: *Io.Mutex) void {
1085 const k: *Kqueue = @ptrCast(@alignCast(userdata));
1086 _ = k;
1087 _ = prev_state;
1088 _ = mutex;
1089 @panic("TODO");
1090}
1091fn mutexUnlock(userdata: ?*anyopaque, prev_state: Io.Mutex.State, mutex: *Io.Mutex) void {
1008fn groupCancel(userdata: ?*anyopaque, group: *Io.Group, token: *anyopaque) void {
10921009 const k: *Kqueue = @ptrCast(@alignCast(userdata));
10931010 _ = k;
1094 _ = prev_state;
1095 _ = mutex;
1011 _ = group;
1012 _ = token;
10961013 @panic("TODO");
10971014}
10981015
1099fn conditionWait(userdata: ?*anyopaque, cond: *Io.Condition, mutex: *Io.Mutex) Io.Cancelable!void {
1100 const k: *Kqueue = @ptrCast(@alignCast(userdata));
1101 k.yield(null, .{ .condition_wait = .{ .cond = cond, .mutex = mutex } });
1102 const thread = Thread.current();
1103 const fiber = thread.currentFiber();
1104 const cond_impl = fiber.resultPointer(Condition);
1105 try mutex.lock(k.io());
1106 switch (cond_impl.event) {
1107 .queued => {},
1108 .wake => |wake| if (fiber.queue_next) |next_fiber| switch (wake) {
1109 .one => if (@cmpxchgStrong(
1110 ?*Fiber,
1111 @as(*?*Fiber, @ptrCast(&cond.state)),
1112 null,
1113 next_fiber,
1114 .release,
1115 .acquire,
1116 )) |old_fiber| {
1117 const old_cond_impl = old_fiber.?.resultPointer(Condition);
1118 assert(old_cond_impl.tail.queue_next == null);
1119 old_cond_impl.tail.queue_next = next_fiber;
1120 old_cond_impl.tail = cond_impl.tail;
1121 },
1122 .all => k.schedule(thread, .{ .head = next_fiber, .tail = cond_impl.tail }),
1123 },
1124 }
1125 fiber.queue_next = null;
1126}
1127
1128fn conditionWaitUncancelable(userdata: ?*anyopaque, cond: *Io.Condition, mutex: *Io.Mutex) void {
1016fn select(userdata: ?*anyopaque, futures: []const *Io.AnyFuture) Io.Cancelable!usize {
11291017 const k: *Kqueue = @ptrCast(@alignCast(userdata));
11301018 _ = k;
1131 _ = cond;
1132 _ = mutex;
1019 _ = futures;
11331020 @panic("TODO");
11341021}
1135fn conditionWake(userdata: ?*anyopaque, cond: *Io.Condition, wake: Io.Condition.Wake) void {
1136 const k: *Kqueue = @ptrCast(@alignCast(userdata));
1137 const waiting_fiber = @atomicRmw(?*Fiber, @as(*?*Fiber, @ptrCast(&cond.state)), .Xchg, null, .acquire) orelse return;
1138 waiting_fiber.resultPointer(Condition).event = .{ .wake = wake };
1139 k.yield(waiting_fiber, .reschedule);
1140}
11411022
1142fn dirCreateDir(userdata: ?*anyopaque, dir: Dir, sub_path: []const u8, mode: Dir.Mode) Dir.CreateDirError!void {
1023fn dirCreateDir(userdata: ?*anyopaque, dir: Dir, sub_path: []const u8, permissions: Dir.Permissions) Dir.CreateDirError!void {
11431024 const k: *Kqueue = @ptrCast(@alignCast(userdata));
11441025 _ = k;
11451026 _ = dir;
11461027 _ = sub_path;
1147 _ = mode;
1028 _ = permissions;
11481029 @panic("TODO");
11491030}
1150fn dirCreateDirPath(userdata: ?*anyopaque, dir: Dir, sub_path: []const u8, mode: Dir.Mode) Dir.CreateDirError!void {
1031
1032fn dirCreateDirPath(
1033 userdata: ?*anyopaque,
1034 dir: Dir,
1035 sub_path: []const u8,
1036 permissions: Dir.Permissions,
1037) Dir.CreateDirPathError!Dir.CreatePathStatus {
11511038 const k: *Kqueue = @ptrCast(@alignCast(userdata));
11521039 _ = k;
11531040 _ = dir;
11541041 _ = sub_path;
1155 _ = mode;
1042 _ = permissions;
11561043 @panic("TODO");
11571044}
1158fn dirCreateDirPathOpen(userdata: ?*anyopaque, dir: Dir, sub_path: []const u8, options: Dir.OpenOptions) Dir.CreateDirPathOpenError!Dir {
1045
1046fn dirCreateDirPathOpen(
1047 userdata: ?*anyopaque,
1048 dir: Dir,
1049 sub_path: []const u8,
1050 permissions: Dir.Permissions,
1051 options: Dir.OpenOptions,
1052) Dir.CreateDirPathOpenError!Dir {
11591053 const k: *Kqueue = @ptrCast(@alignCast(userdata));
11601054 _ = k;
11611055 _ = dir;
11621056 _ = sub_path;
1057 _ = permissions;
11631058 _ = options;
11641059 @panic("TODO");
11651060}
1061
11661062fn dirStat(userdata: ?*anyopaque, dir: Dir) Dir.StatError!Dir.Stat {
11671063 const k: *Kqueue = @ptrCast(@alignCast(userdata));
11681064 _ = k;
11691065 _ = dir;
11701066 @panic("TODO");
11711067}
1172fn dirStatFile(userdata: ?*anyopaque, dir: Dir, sub_path: []const u8, options: Dir.StatPathOptions) Dir.StatFileError!File.Stat {
1068
1069fn dirStatFile(
1070 userdata: ?*anyopaque,
1071 dir: Dir,
1072 sub_path: []const u8,
1073 options: Dir.StatFileOptions,
1074) Dir.StatFileError!File.Stat {
11731075 const k: *Kqueue = @ptrCast(@alignCast(userdata));
11741076 _ = k;
11751077 _ = dir;
......@@ -1209,10 +1111,10 @@ fn dirOpenDir(userdata: ?*anyopaque, dir: Dir, sub_path: []const u8, options: Di
12091111 _ = options;
12101112 @panic("TODO");
12111113}
1212fn dirClose(userdata: ?*anyopaque, dir: Dir) void {
1114fn dirClose(userdata: ?*anyopaque, dirs: []const Dir) void {
12131115 const k: *Kqueue = @ptrCast(@alignCast(userdata));
12141116 _ = k;
1215 _ = dir;
1117 _ = dirs;
12161118 @panic("TODO");
12171119}
12181120fn fileStat(userdata: ?*anyopaque, file: File) File.StatError!File.Stat {
......@@ -1221,35 +1123,57 @@ fn fileStat(userdata: ?*anyopaque, file: File) File.StatError!File.Stat {
12211123 _ = file;
12221124 @panic("TODO");
12231125}
1224fn fileClose(userdata: ?*anyopaque, file: File) void {
1126
1127fn fileClose(userdata: ?*anyopaque, files: []const File) void {
12251128 const k: *Kqueue = @ptrCast(@alignCast(userdata));
12261129 _ = k;
1227 _ = file;
1130 _ = files;
12281131 @panic("TODO");
12291132}
1230fn fileWriteStreaming(userdata: ?*anyopaque, file: File, buffer: [][]const u8) File.WriteStreamingError!usize {
1133
1134fn fileWriteStreaming(
1135 userdata: ?*anyopaque,
1136 file: File,
1137 header: []const u8,
1138 data: []const []const u8,
1139 splat: usize,
1140) File.Writer.Error!usize {
12311141 const k: *Kqueue = @ptrCast(@alignCast(userdata));
12321142 _ = k;
12331143 _ = file;
1234 _ = buffer;
1144 _ = header;
1145 _ = data;
1146 _ = splat;
12351147 @panic("TODO");
12361148}
1237fn fileWritePositional(userdata: ?*anyopaque, file: File, buffer: [][]const u8, offset: u64) File.WritePositionalError!usize {
1149
1150fn fileWritePositional(
1151 userdata: ?*anyopaque,
1152 file: File,
1153 header: []const u8,
1154 data: []const []const u8,
1155 splat: usize,
1156 offset: u64,
1157) File.WritePositionalError!usize {
12381158 const k: *Kqueue = @ptrCast(@alignCast(userdata));
12391159 _ = k;
12401160 _ = file;
1241 _ = buffer;
1161 _ = header;
1162 _ = data;
1163 _ = splat;
12421164 _ = offset;
12431165 @panic("TODO");
12441166}
1245fn fileReadStreaming(userdata: ?*anyopaque, file: File, data: [][]u8) File.Reader.Error!usize {
1167
1168fn fileReadStreaming(userdata: ?*anyopaque, file: File, data: []const []u8) File.Reader.Error!usize {
12461169 const k: *Kqueue = @ptrCast(@alignCast(userdata));
12471170 _ = k;
12481171 _ = file;
12491172 _ = data;
12501173 @panic("TODO");
12511174}
1252fn fileReadPositional(userdata: ?*anyopaque, file: File, data: [][]u8, offset: u64) File.ReadPositionalError!usize {
1175
1176fn fileReadPositional(userdata: ?*anyopaque, file: File, data: []const []u8, offset: u64) File.ReadPositionalError!usize {
12531177 const k: *Kqueue = @ptrCast(@alignCast(userdata));
12541178 _ = k;
12551179 _ = file;
......@@ -1271,12 +1195,6 @@ fn fileSeekTo(userdata: ?*anyopaque, file: File, absolute_offset: u64) File.Seek
12711195 _ = absolute_offset;
12721196 @panic("TODO");
12731197}
1274fn openExecutable(userdata: ?*anyopaque, file: File.OpenFlags) File.OpenExecutableError!File {
1275 const k: *Kqueue = @ptrCast(@alignCast(userdata));
1276 _ = k;
1277 _ = file;
1278 @panic("TODO");
1279}
12801198
12811199fn now(userdata: ?*anyopaque, clock: Io.Clock) Io.Clock.Error!Io.Timestamp {
12821200 const k: *Kqueue = @ptrCast(@alignCast(userdata));
......@@ -1576,10 +1494,10 @@ fn netWrite(userdata: ?*anyopaque, dest: net.Socket.Handle, header: []const u8,
15761494 @panic("TODO");
15771495}
15781496
1579fn netClose(userdata: ?*anyopaque, handle: net.Socket.Handle) void {
1497fn netClose(userdata: ?*anyopaque, handles: []const net.Socket.Handle) void {
15801498 const k: *Kqueue = @ptrCast(@alignCast(userdata));
15811499 _ = k;
1582 _ = handle;
1500 _ = handles;
15831501 @panic("TODO");
15841502}
15851503
......@@ -1611,13 +1529,13 @@ fn netInterfaceName(userdata: ?*anyopaque, interface: net.Interface) net.Interfa
16111529fn netLookup(
16121530 userdata: ?*anyopaque,
16131531 host_name: net.HostName,
1614 result: *Io.Queue(net.HostName.LookupResult),
1532 resolved: *Io.Queue(net.HostName.LookupResult),
16151533 options: net.HostName.LookupOptions,
1616) void {
1534) net.HostName.LookupError!void {
16171535 const k: *Kqueue = @ptrCast(@alignCast(userdata));
16181536 _ = k;
16191537 _ = host_name;
1620 _ = result;
1538 _ = resolved;
16211539 _ = options;
16221540 @panic("TODO");
16231541}
......@@ -1772,14 +1690,6 @@ fn checkCancel(k: *Kqueue) error{Canceled}!void {
17721690 if (cancelRequested(k)) return error.Canceled;
17731691}
17741692
1775const Condition = struct {
1776 tail: *Fiber,
1777 event: union(enum) {
1778 queued,
1779 wake: Io.Condition.Wake,
1780 },
1781};
1782
17831693pub const KEventError = error{
17841694 /// The process does not have permission to register a filter.
17851695 AccessDenied,