1const Alignment = std.mem.Alignment;
2const Allocator = std.mem.Allocator;
3const Argv0 = Io.Threaded.Argv0;
4const assert = std.debug.assert;
5const builtin = @import("builtin");
6const c = std.c;
7const ChdirError = Io.Threaded.ChdirError;
8const clockToPosix = Io.Threaded.clockToPosix;
9const closeFd = Io.Threaded.closeFd;
10const Csprng = Io.Threaded.Csprng;
11const default_PATH = Io.Threaded.default_PATH;
12const Dir = Io.Dir;
13const Environ = Io.Threaded.Environ;
14const errnoBug = Io.Threaded.errnoBug;
15const Evented = @This();
16const fallbackSeed = Io.Threaded.fallbackSeed;
17const File = Io.File;
18const Io = std.Io;
19const iovec = std.posix.iovec;
20const iovec_const = std.posix.iovec_const;
21const log = std.log.scoped(.dispatch);
22const max_iovecs_len = Io.Threaded.max_iovecs_len;
23const nanosecondsFromPosix = Io.Threaded.nanosecondsFromPosix;
24const net = Io.net;
25const pathToPosix = Io.Threaded.pathToPosix;
26const process = std.process;
27const recoverableOsBugDetected = Io.Threaded.recoverableOsBugDetected;
28const setTimestampToPosix = Io.Threaded.setTimestampToPosix;
29const splat_buffer_size = Io.Threaded.splat_buffer_size;
30const statFromPosix = Io.Threaded.statFromPosix;
31const statusToTerm = Io.Threaded.statusToTerm;
32const std = @import("std");
33const timestampFromPosix = Io.Threaded.timestampFromPosix;
34const unexpectedErrno = std.posix.unexpectedErrno;
35const UseSendfile = Io.Threaded.UseSendfile;
36const UseFcopyfile = Io.Threaded.UseFcopyfile;
37
38/// Empirically saw >4KB being used by the llvm aarch64 backend.
39const main_loop_stack_size = 8 * 1024;
40
41queue: c.dispatch.queue_t,
42backing_allocator_needs_mutex: bool,
43backing_allocator_mutex: Mutex,
44/// Does not need to be thread-safe if not used elsewhere.
45backing_allocator: Allocator,
46main_fiber: Fiber,
47main_loop_stack: [*]align(builtin.target.stackAlignment()) u8,
48exit_semaphore: c.dispatch.semaphore_t,
49
50use_sendfile: UseSendfile,
51use_fcopyfile: UseFcopyfile,
52leeway: u64,
53
54futexes: [1 << 8]Futex,
55
56init_stderr_writer: c.dispatch.once_t,
57stderr_mutex: Mutex,
58stderr_writer: File.Writer,
59stderr_mode: Io.Terminal.Mode,
60
61scan_environ: c.dispatch.once_t,
62environ: Environ,
63
64open_dev_null: c.dispatch.once_t,
65dev_null_file: File.OpenError!File,
66
67csprng_mutex: Mutex,
68csprng: Csprng,
69
70const Thread = struct {
71 main_context: Io.fiber.Context,
72 current_context: ?*Io.fiber.Context,
73 seed_csprng: c.dispatch.once_t,
74 csprng: Csprng,
75
76 threadlocal var self: Thread = .{
77 .main_context = undefined,
78 .current_context = null,
79 .seed_csprng = .init,
80 .csprng = undefined,
81 };
82
83 noinline fn current() *Thread {
84 return &self;
85 }
86
87 fn currentFiber(thread: *Thread) *Fiber {
88 assert(thread.current_context != &thread.main_context);
89 return @fieldParentPtr("context", thread.current_context.?);
90 }
91
92 const List = struct {
93 allocated: []Thread,
94 reserved: u32,
95 active: u32,
96 };
97};
98
99const Fiber = struct {
100 required_align: void align(4),
101 evented: *Evented,
102 context: Io.fiber.Context,
103 link: union {
104 awaiter: ?*Fiber,
105 group: struct { prev: ?*Fiber, next: ?*Fiber },
106 },
107 awaiting_group: Group,
108 cancel_status: CancelStatus,
109 cancel_protection: CancelProtection,
110
111 var next_name: u64 = 0;
112
113 const CancelStatus = packed struct(usize) {
114 requested: bool,
115 awaiting: Awaiting,
116
117 const unrequested: CancelStatus = .{ .requested = false, .awaiting = .nothing };
118
119 const Awaiting = enum(@Int(.unsigned, @bitSizeOf(usize) - shift)) {
120 nothing = 0,
121 group = 1,
122 _,
123
124 const shift = 1;
125
126 fn subWrap(lhs: Awaiting, rhs: Awaiting) Awaiting {
127 return @fromBackingInt(@intCast(@backingInt(lhs) -% @backingInt(rhs)));
128 }
129
130 fn fromCancelable(cancelable: *Cancelable) Awaiting {
131 return @fromBackingInt(@intCast(@shrExact(@intFromPtr(cancelable), shift)));
132 }
133
134 fn toCancelable(awaiting: Awaiting) *Cancelable {
135 return @ptrFromInt(@shlExact(@as(usize, @backingInt(awaiting)), shift));
136 }
137 };
138
139 fn changeAwaiting(
140 cancel_status: *CancelStatus,
141 old_awaiting: Awaiting,
142 new_awaiting: Awaiting,
143 ) bool {
144 const old_cancel_status = @atomicRmw(CancelStatus, cancel_status, .Add, .{
145 .requested = false,
146 .awaiting = new_awaiting.subWrap(old_awaiting),
147 }, .release);
148 assert(old_cancel_status.awaiting == old_awaiting);
149 return old_cancel_status.requested;
150 }
151 };
152
153 const CancelProtection = packed struct {
154 user: Io.CancelProtection,
155 acknowledged: bool,
156
157 const unblocked: CancelProtection = .{ .user = .unblocked, .acknowledged = false };
158
159 fn check(cancel_protection: CancelProtection) Io.CancelProtection {
160 return @fromBackingInt(@intCast(@intFromBool(cancel_protection != unblocked)));
161 }
162
163 fn acknowledge(cancel_protection: *CancelProtection) void {
164 assert(!cancel_protection.acknowledged);
165 cancel_protection.acknowledged = true;
166 }
167
168 fn recancel(cancel_protection: *CancelProtection) void {
169 assert(cancel_protection.acknowledged);
170 cancel_protection.acknowledged = false;
171 }
172
173 test check {
174 try std.testing.expectEqual(Io.CancelProtection.unblocked, check(.unblocked));
175 try std.testing.expectEqual(Io.CancelProtection.blocked, check(.{
176 .user = .unblocked,
177 .acknowledged = true,
178 }));
179 try std.testing.expectEqual(Io.CancelProtection.blocked, check(.{
180 .user = .blocked,
181 .acknowledged = false,
182 }));
183 try std.testing.expectEqual(Io.CancelProtection.blocked, check(.{
184 .user = .blocked,
185 .acknowledged = true,
186 }));
187 }
188 };
189
190 const finished: ?*Fiber = @ptrFromInt(@alignOf(Fiber));
191
192 const max_result_align: Alignment = .@"16";
193 const max_result_size = max_result_align.forward(512);
194 /// This includes any stack realignments that need to happen, and also the
195 /// initial frame return address slot and argument frame, depending on target.
196 const min_stack_size = 60 * 1024 * 1024;
197 const max_context_align: Alignment = .@"16";
198 const max_context_size = max_context_align.forward(1024);
199 const max_closure_size: usize = @sizeOf(AsyncClosure);
200 const max_closure_align: Alignment = .of(AsyncClosure);
201 const allocation_size = std.mem.alignForward(
202 usize,
203 max_closure_align.max(max_context_align).forward(
204 max_result_align.forward(@sizeOf(Fiber)) + max_result_size + min_stack_size,
205 ) + max_closure_size + max_context_size,
206 std.heap.page_size_max,
207 );
208
209 fn create(ev: *Evented) error{OutOfMemory}!*Fiber {
210 return @ptrCast(try ev.allocator().alignedAlloc(u8, .of(Fiber), allocation_size));
211 }
212
213 fn destroy(fiber: *Fiber, ev: *Evented) void {
214 ev.allocator().free(fiber.allocatedSlice());
215 }
216
217 fn allocatedSlice(f: *Fiber) []align(@alignOf(Fiber)) u8 {
218 return @as([*]align(@alignOf(Fiber)) u8, @ptrCast(f))[0..allocation_size];
219 }
220
221 fn allocatedEnd(f: *Fiber) [*]u8 {
222 const allocated_slice = f.allocatedSlice();
223 return allocated_slice[allocated_slice.len..].ptr;
224 }
225
226 fn resultPointer(f: *Fiber, comptime Result: type) *Result {
227 return @ptrCast(@alignCast(f.resultBytes(.of(Result))));
228 }
229
230 fn resultBytes(f: *Fiber, alignment: Alignment) [*]u8 {
231 return @ptrFromInt(alignment.forward(@intFromPtr(f) + @sizeOf(Fiber)));
232 }
233
234 const Queue = struct { head: *Fiber, tail: *Fiber };
235
236 /// Like a `*Fiber`, but 2 bits smaller than a pointer (because the LSBs are always 0 due to
237 /// alignment) so that those two bits can be used in a `packed struct`.
238 const PackedPtr = enum(@Int(.unsigned, @bitSizeOf(usize) - 2)) {
239 null = 0,
240 all_ones = std.math.maxInt(@Int(.unsigned, @bitSizeOf(usize) - 2)),
241 _,
242
243 const Split = packed struct(usize) { low: u2, high: PackedPtr };
244 fn pack(ptr: ?*Fiber) PackedPtr {
245 const split: Split = @bitCast(@intFromPtr(ptr));
246 assert(split.low == 0);
247 return split.high;
248 }
249 fn unpack(ptr: PackedPtr) ?*Fiber {
250 const split: Split = .{ .low = 0, .high = ptr };
251 return @ptrFromInt(@as(usize, @bitCast(split)));
252 }
253 };
254
255 fn requestCancel(fiber: *Fiber, ev: *Evented) void {
256 const cancel_status = @atomicRmw(
257 Fiber.CancelStatus,
258 &fiber.cancel_status,
259 .Or,
260 .{ .requested = true, .awaiting = .nothing },
261 .acquire,
262 );
263 assert(!cancel_status.requested);
264 switch (cancel_status.awaiting) {
265 .nothing => {},
266 .group => {
267 // The awaiter received a cancelation request while awaiting a group,
268 // so propagate the cancelation to the group.
269 if (fiber.awaiting_group.cancel(ev, null)) {
270 fiber.awaiting_group = undefined;
271 ev.queue.async(fiber, &Fiber.@"resume");
272 }
273 },
274 _ => |awaiting| awaiting.toCancelable().async(),
275 }
276 }
277
278 fn @"resume"(context: ?*anyopaque) callconv(.c) void {
279 const fiber: *Fiber = @ptrCast(@alignCast(context));
280 const thread: *Thread = .current();
281 const message: SwitchMessage = .{
282 .contexts = .{
283 .old = &thread.main_context,
284 .new = &fiber.context,
285 },
286 .pending_task = .nothing,
287 };
288 contextSwitch(&message).handle(fiber.evented);
289 }
290};
291
292pub fn allocator(ev: *Evented) std.mem.Allocator {
293 return if (ev.backing_allocator_needs_mutex) .{
294 .ptr = ev,
295 .vtable = &.{
296 .alloc = alloc,
297 .resize = resize,
298 .remap = remap,
299 .free = free,
300 },
301 } else ev.backing_allocator;
302}
303
304fn alloc(userdata: *anyopaque, len: usize, alignment: std.mem.Alignment, ret_addr: usize) ?[*]u8 {
305 const ev: *Evented = @ptrCast(@alignCast(userdata));
306 ev.backing_allocator_mutex.lockUncancelable(ev);
307 defer ev.backing_allocator_mutex.unlock();
308 return ev.backing_allocator.rawAlloc(len, alignment, ret_addr);
309}
310
311fn resize(
312 userdata: *anyopaque,
313 memory: []u8,
314 alignment: std.mem.Alignment,
315 new_len: usize,
316 ret_addr: usize,
317) bool {
318 const ev: *Evented = @ptrCast(@alignCast(userdata));
319 ev.backing_allocator_mutex.lockUncancelable(ev);
320 defer ev.backing_allocator_mutex.unlock();
321 return ev.backing_allocator.rawResize(memory, alignment, new_len, ret_addr);
322}
323
324fn remap(
325 userdata: *anyopaque,
326 memory: []u8,
327 alignment: Alignment,
328 new_len: usize,
329 ret_addr: usize,
330) ?[*]u8 {
331 const ev: *Evented = @ptrCast(@alignCast(userdata));
332 ev.backing_allocator_mutex.lockUncancelable(ev);
333 defer ev.backing_allocator_mutex.unlock();
334 return ev.backing_allocator.rawRemap(memory, alignment, new_len, ret_addr);
335}
336
337fn free(userdata: *anyopaque, memory: []u8, alignment: std.mem.Alignment, ret_addr: usize) void {
338 const ev: *Evented = @ptrCast(@alignCast(userdata));
339 ev.backing_allocator_mutex.lockUncancelable(ev);
340 defer ev.backing_allocator_mutex.unlock();
341 return ev.backing_allocator.rawFree(memory, alignment, ret_addr);
342}
343
344pub fn io(ev: *Evented) Io {
345 return .{
346 .userdata = ev,
347 .vtable = &.{
348 .crashHandler = crashHandler,
349
350 .async = async,
351 .concurrent = concurrent,
352 .await = await,
353 .cancel = cancel,
354
355 .groupAsync = groupAsync,
356 .groupConcurrent = groupConcurrent,
357 .groupAwait = groupAwait,
358 .groupCancel = groupCancel,
359
360 .recancel = recancel,
361 .swapCancelProtection = swapCancelProtection,
362 .checkCancel = checkCancel,
363
364 .futexWait = futexWait,
365 .futexWaitUncancelable = futexWaitUncancelable,
366 .futexWake = futexWake,
367
368 .operate = operate,
369 .batchAwaitAsync = batchAwaitAsync,
370 .batchAwaitConcurrent = batchAwaitConcurrent,
371 .batchCancel = batchCancel,
372
373 .dirCreateDir = dirCreateDir,
374 .dirCreateDirPath = dirCreateDirPath,
375 .dirCreateDirPathOpen = dirCreateDirPathOpen,
376 .dirOpenDir = dirOpenDir,
377 .dirStat = dirStat,
378 .dirStatFile = dirStatFile,
379 .dirAccess = dirAccess,
380 .dirCreateFile = dirCreateFile,
381 .dirCreateFileAtomic = dirCreateFileAtomic,
382 .dirOpenFile = dirOpenFile,
383 .dirClose = dirClose,
384 .dirRead = dirRead,
385 .dirRealPath = dirRealPath,
386 .dirRealPathFile = dirRealPathFile,
387 .dirDeleteFile = dirDeleteFile,
388 .dirDeleteDir = dirDeleteDir,
389 .dirRename = dirRename,
390 .dirRenamePreserve = dirRenamePreserve,
391 .dirSymLink = dirSymLink,
392 .dirReadLink = dirReadLink,
393 .dirSetOwner = dirSetOwner,
394 .dirSetFileOwner = dirSetFileOwner,
395 .dirSetPermissions = dirSetPermissions,
396 .dirSetFilePermissions = dirSetFilePermissions,
397 .dirSetTimestamps = dirSetTimestamps,
398 .dirHardLink = dirHardLink,
399
400 .fileStat = fileStat,
401 .fileLength = fileLength,
402 .fileClose = fileClose,
403 .fileWritePositional = fileWritePositional,
404 .fileWriteFileStreaming = fileWriteFileStreaming,
405 .fileWriteFilePositional = fileWriteFilePositional,
406 .fileReadPositional = fileReadPositional,
407 .fileSeekBy = fileSeekBy,
408 .fileSeekTo = fileSeekTo,
409 .fileSync = fileSync,
410 .fileIsTty = fileIsTty,
411 .fileEnableAnsiEscapeCodes = fileEnableAnsiEscapeCodes,
412 .fileSupportsAnsiEscapeCodes = fileIsTty,
413 .fileSetLength = fileSetLength,
414 .fileSetOwner = fileSetOwner,
415 .fileSetPermissions = fileSetPermissions,
416 .fileSetTimestamps = fileSetTimestamps,
417 .fileLock = fileLock,
418 .fileTryLock = fileTryLock,
419 .fileUnlock = fileUnlock,
420 .fileDowngradeLock = fileDowngradeLock,
421 .fileRealPath = fileRealPath,
422 .fileHardLink = fileHardLink,
423
424 .fileMemoryMapCreate = fileMemoryMapCreate,
425 .fileMemoryMapDestroy = fileMemoryMapDestroy,
426 .fileMemoryMapSetLength = fileMemoryMapSetLength,
427 .fileMemoryMapRead = fileMemoryMapRead,
428 .fileMemoryMapWrite = fileMemoryMapWrite,
429
430 .processExecutableOpen = processExecutableOpen,
431 .processExecutablePath = processExecutablePath,
432 .lockStderr = lockStderr,
433 .tryLockStderr = tryLockStderr,
434 .unlockStderr = unlockStderr,
435 .processCurrentPath = processCurrentPath,
436 .processSetCurrentDir = processSetCurrentDir,
437 .processSetCurrentPath = processSetCurrentPath,
438 .processReplace = processReplace,
439 .processReplacePath = processReplacePath,
440 .processSpawn = processSpawn,
441 .processSpawnPath = processSpawnPath,
442 .childWait = childWait,
443 .childKill = childKill,
444
445 .progressParentFile = progressParentFile,
446
447 .now = now,
448 .clockResolution = clockResolution,
449 .sleep = sleep,
450
451 .random = random,
452 .randomSecure = randomSecure,
453
454 .netListenIp = netListenIpUnavailable,
455 .netAccept = netAcceptUnavailable,
456 .netBindIp = netBindIpUnavailable,
457 .netConnectIp = netConnectIpUnavailable,
458 .netListenUnix = netListenUnixUnavailable,
459 .netConnectUnix = netConnectUnixUnavailable,
460 .netSocketCreatePair = netSocketCreatePairUnavailable,
461 .netWriteFile = netWriteFileUnavailable,
462 .netClose = netClose,
463 .netShutdown = netShutdownUnavailable,
464 .netInterfaceNameResolve = netInterfaceNameResolveUnavailable,
465 .netInterfaceName = netInterfaceNameUnavailable,
466 .netLookup = netLookupUnavailable,
467 },
468 };
469}
470
471pub const InitOptions = struct {
472 backing_allocator_needs_mutex: bool = true,
473 target_queue: ?c.dispatch.queue_t = .TARGET_DEFAULT,
474 /// Upper limit on the allowable delay in processing timeouts in order to improve power
475 /// consumption and system performance.
476 leeway: Io.Duration = .fromMilliseconds(10),
477
478 /// Affects the following operations:
479 /// * `processExecutablePath` on OpenBSD and Haiku.
480 argv0: Argv0 = .empty,
481 /// Affects the following operations:
482 /// * `fileIsTty`
483 /// * `processSpawn`, `processSpawnPath`, `processReplace`, `processReplacePath`
484 environ: process.Environ = .empty,
485};
486
487pub fn init(ev: *Evented, backing_allocator: Allocator, options: InitOptions) !void {
488 const queue = c.dispatch.queue_create_with_target(
489 "org.ziglang.std.Io.Dispatch",
490 .CONCURRENT(),
491 options.target_queue,
492 ) orelse return error.SystemResources;
493 errdefer queue.as_object().release();
494 const main_loop_stack = try backing_allocator.alignedAlloc(
495 u8,
496 .fromByteUnits(builtin.target.stackAlignment()),
497 main_loop_stack_size,
498 );
499 errdefer backing_allocator.free(main_loop_stack);
500 const exit_semaphore = c.dispatch.semaphore_create(0) orelse return error.SystemResources;
501 errdefer exit_semaphore.as_object().release();
502 ev.* = .{
503 .queue = queue,
504 .backing_allocator_needs_mutex = options.backing_allocator_needs_mutex,
505 .backing_allocator_mutex = undefined,
506 .backing_allocator = backing_allocator,
507 .main_fiber = .{
508 .required_align = {},
509 .evented = ev,
510 .context = undefined,
511 .link = .{ .awaiter = null },
512 .awaiting_group = undefined,
513 .cancel_status = .unrequested,
514 .cancel_protection = .unblocked,
515 },
516 .main_loop_stack = main_loop_stack.ptr,
517 .exit_semaphore = exit_semaphore,
518
519 .use_fcopyfile = .default,
520 .use_sendfile = .default,
521 .leeway = std.math.lossyCast(u64, options.leeway.toNanoseconds()),
522
523 .futexes = undefined,
524
525 .init_stderr_writer = .init,
526 .stderr_mutex = undefined,
527 .stderr_writer = .{
528 .io = ev.io(),
529 .interface = Io.File.Writer.initInterface(&.{}),
530 .file = .stderr(),
531 .mode = .streaming,
532 },
533 .stderr_mode = .no_color,
534
535 .scan_environ = if (options.environ.block.isEmpty()) .done else .init,
536 .environ = .{ .process_environ = options.environ },
537
538 .open_dev_null = .init,
539 .dev_null_file = error.FileNotFound,
540
541 .csprng_mutex = undefined,
542 .csprng = .uninitialized,
543 };
544 try ev.backing_allocator_mutex.init(queue);
545 errdefer ev.backing_allocator_mutex.deinit();
546 var initialized_futexes: usize = 0;
547 errdefer for (ev.futexes[0..initialized_futexes]) |*futex| futex.deinit();
548 for (&ev.futexes) |*futex| {
549 try futex.init(queue);
550 initialized_futexes += 1;
551 }
552 try ev.stderr_mutex.init(queue);
553 errdefer ev.stderr_mutex.deinit();
554 try ev.csprng_mutex.init(queue);
555 errdefer ev.csprng_mutex.deinit();
556 const thread: *Thread = .current();
557 thread.main_context = switch (builtin.cpu.arch) {
558 .aarch64 => .{
559 .sp = @intFromPtr(main_loop_stack[main_loop_stack_size..].ptr),
560 .fp = @intFromPtr(ev),
561 .pc = @intFromPtr(&mainLoopEntry),
562 },
563 .x86_64 => .{
564 .rsp = @intFromPtr(main_loop_stack[main_loop_stack_size..].ptr) - 8,
565 .rbp = @intFromPtr(ev),
566 .rip = @intFromPtr(&mainLoopEntry),
567 },
568 else => |arch| @compileError("unimplemented architecture: " ++ @tagName(arch)),
569 };
570 thread.current_context = &ev.main_fiber.context;
571}
572
573pub fn deinit(ev: *Evented) void {
574 assert(Thread.current().currentFiber() == &ev.main_fiber);
575 ev.yield(.exit);
576 ev.csprng_mutex.deinit();
577 if (ev.dev_null_file) |file| fileClose(ev, &.{file}) else |_| {}
578 ev.stderr_mutex.deinit();
579 for (&ev.futexes) |*futex| futex.deinit();
580 ev.exit_semaphore.as_object().release();
581 ev.backing_allocator_mutex.deinit();
582 ev.backing_allocator.free(ev.main_loop_stack[0..main_loop_stack_size]);
583 ev.queue.as_object().release();
584}
585
586fn yield(ev: *Evented, pending_task: SwitchMessage.PendingTask) void {
587 const thread: *Thread = .current();
588 const message: SwitchMessage = .{
589 .contexts = .{
590 .old = thread.current_context.?,
591 .new = &thread.main_context,
592 },
593 .pending_task = pending_task,
594 };
595 contextSwitch(&message).handle(ev);
596}
597
598fn mainLoopEntry() callconv(.naked) void {
599 switch (builtin.cpu.arch) {
600 .aarch64 => asm volatile (
601 \\ mov x0, fp
602 \\ mov fp, #0
603 \\ b %[mainLoop]
604 :
605 : [mainLoop] "X" (&mainLoop),
606 ),
607 .x86_64 => asm volatile (
608 \\ movq %%rbp, %%rdi
609 \\ xor %%ebp, %%ebp
610 \\ jmp %[mainLoop:P]
611 :
612 : [mainLoop] "X" (&mainLoop),
613 ),
614 else => |arch| @compileError("unimplemented architecture: " ++ @tagName(arch)),
615 }
616}
617
618fn mainLoop(ev: *Evented, message: *const SwitchMessage) callconv(.c) noreturn {
619 message.handle(ev);
620 assert(ev.exit_semaphore.wait(.FOREVER) == 0);
621 Fiber.@"resume"(&ev.main_fiber);
622 unreachable; // switched to dead fiber
623}
624
625const SwitchMessage = struct {
626 contexts: Io.fiber.Switch,
627 pending_task: PendingTask,
628
629 const PendingTask = union(enum) {
630 nothing,
631 await: *Fiber,
632 activate: c.dispatch.object_t,
633 @"resume": c.dispatch.object_t,
634 group_await: Group,
635 group_cancel: Group,
636 mutex_wait: *Mutex.Waiter,
637 futex_wait: *Futex.Waiter,
638 futex_wake: *Futex.Waker,
639 sleep_wait: *SleepWaiter,
640 after: c.dispatch.time_t,
641 destroy,
642 exit,
643 };
644
645 fn handle(message: *const SwitchMessage, ev: *Evented) void {
646 const thread: *Thread = .current();
647 thread.current_context = message.contexts.new;
648 switch (message.pending_task) {
649 .nothing => {},
650 .await => |awaiting| {
651 const awaiter: *Fiber = @alignCast(@fieldParentPtr("context", message.contexts.old));
652 if (@atomicRmw(?*Fiber, &awaiting.link.awaiter, .Xchg, awaiter, .acq_rel) ==
653 Fiber.finished) ev.queue.async(awaiter, &Fiber.@"resume");
654 },
655 .activate => |object| object.activate(),
656 .@"resume" => |object| object.@"resume"(),
657 .group_await => |group| {
658 const fiber: *Fiber = @alignCast(@fieldParentPtr("context", message.contexts.old));
659 if (group.await(ev, fiber)) ev.queue.async(fiber, &Fiber.@"resume");
660 },
661 .group_cancel => |group| {
662 const fiber: *Fiber = @alignCast(@fieldParentPtr("context", message.contexts.old));
663 if (group.cancel(ev, fiber)) ev.queue.async(fiber, &Fiber.@"resume");
664 },
665 .mutex_wait => |waiter| {
666 waiter.sleeper =
667 .init(ev.queue, @alignCast(@fieldParentPtr("context", message.contexts.old)));
668 switch (waiter.sleeper.fiber.cancel_protection.check()) {
669 .unblocked => {},
670 .blocked => waiter.cancelable = .blocked,
671 }
672 waiter.mutex.queue.async(waiter, &Mutex.Waiter.add);
673 },
674 .futex_wait => |waiter| {
675 waiter.sleeper =
676 .init(ev.queue, @alignCast(@fieldParentPtr("context", message.contexts.old)));
677 switch (waiter.sleeper.fiber.cancel_protection.check()) {
678 .unblocked => {},
679 .blocked => waiter.cancelable = .blocked,
680 }
681 waiter.futex.queue.async(waiter, &Futex.Waiter.add);
682 },
683 .futex_wake => |waker| {
684 waker.sleeper =
685 .init(ev.queue, @alignCast(@fieldParentPtr("context", message.contexts.old)));
686 waker.futex.queue.async(waker, &Futex.Waker.remove);
687 },
688 .sleep_wait => |waiter| {
689 waiter.sleeper =
690 .init(ev.queue, @alignCast(@fieldParentPtr("context", message.contexts.old)));
691 const queue = waiter.cancelable.queue;
692 switch (waiter.sleeper.fiber.cancel_protection.check()) {
693 .unblocked => {},
694 .blocked => waiter.cancelable = .blocked,
695 }
696 queue.async(waiter, &SleepWaiter.start);
697 },
698 .after => |when| {
699 const fiber: *Fiber = @alignCast(@fieldParentPtr("context", message.contexts.old));
700 when.after(ev.queue, fiber, &Fiber.@"resume");
701 },
702 .destroy => {
703 const fiber: *Fiber = @alignCast(@fieldParentPtr("context", message.contexts.old));
704 fiber.destroy(ev);
705 },
706 .exit => _ = ev.exit_semaphore.signal(),
707 }
708 }
709};
710
711inline fn contextSwitch(message: *const SwitchMessage) *const SwitchMessage {
712 return @fieldParentPtr("contexts", Io.fiber.contextSwitch(&message.contexts));
713}
714
715const Cancelable = struct {
716 required_align: void align(2) = {},
717 queue: c.dispatch.queue_t,
718 cancel: c.dispatch.function_t,
719
720 const fn_ptr_align = std.meta.alignment(c.dispatch.function_t);
721 const is_blocked: c.dispatch.function_t = @ptrFromInt(fn_ptr_align * 1);
722 const is_requested: c.dispatch.function_t = @ptrFromInt(fn_ptr_align * 2);
723
724 const blocked: Cancelable = .{ .queue = undefined, .cancel = is_blocked };
725
726 const RequestedError = error{CancelRequested};
727
728 fn enter(cancelable: *Cancelable, fiber: *Fiber) RequestedError!void {
729 const function = cancelable.cancel;
730 assert(function != is_requested);
731 if (function == is_blocked) {
732 @branchHint(.unlikely);
733 return;
734 }
735 if (@cmpxchgStrong(
736 Fiber.CancelStatus,
737 &fiber.cancel_status,
738 .{ .requested = false, .awaiting = .nothing },
739 .{ .requested = false, .awaiting = .fromCancelable(cancelable) },
740 .release,
741 .monotonic,
742 )) |cancel_status| {
743 assert(cancel_status.requested and cancel_status.awaiting == .nothing);
744 cancelable.cancel = is_requested;
745 return error.CancelRequested;
746 }
747 }
748
749 fn leave(cancelable: *Cancelable, fiber: *Fiber) RequestedError!void {
750 const function = cancelable.cancel;
751 assert(function != is_requested);
752 if (function == is_blocked) {
753 @branchHint(.unlikely);
754 return;
755 }
756 const cancel_status = @atomicRmw(Fiber.CancelStatus, &fiber.cancel_status, .And, .{
757 .requested = true,
758 .awaiting = .nothing,
759 }, .monotonic);
760 assert(cancel_status.awaiting.toCancelable() == cancelable);
761 if (cancel_status.requested) return error.CancelRequested;
762 }
763
764 fn async(cancelable: *Cancelable) void {
765 const function = cancelable.cancel;
766 assert(function != is_blocked and function != is_requested);
767 cancelable.queue.async(cancelable, function);
768 }
769
770 fn requested(cancelable: *Cancelable, fiber: *Fiber) void {
771 const function = cancelable.cancel;
772 assert(function != is_blocked and function != is_requested);
773 assert(@atomicLoad(Fiber.CancelStatus, &fiber.cancel_status, .monotonic) == Fiber.CancelStatus{
774 .requested = true,
775 .awaiting = .fromCancelable(cancelable),
776 });
777 cancelable.cancel = is_requested;
778 @atomicStore(Fiber.CancelStatus, &fiber.cancel_status, .{
779 .requested = true,
780 .awaiting = .nothing,
781 }, .monotonic);
782 }
783
784 fn acknowledge(cancelable: *Cancelable, fiber: *Fiber) Io.Cancelable!void {
785 if (cancelable.cancel == is_requested) {
786 @branchHint(.unlikely);
787 fiber.cancel_protection.acknowledge();
788 return error.Canceled;
789 }
790 }
791};
792
793const Sleeper = struct {
794 queue: c.dispatch.queue_t,
795 fiber: *Fiber,
796
797 fn init(queue: c.dispatch.queue_t, fiber: *Fiber) Sleeper {
798 queue.as_object().retain();
799 return .{ .queue = queue, .fiber = fiber };
800 }
801
802 fn wake(context: ?*anyopaque) callconv(.c) void {
803 const sleeper: *Sleeper = @ptrCast(@alignCast(context));
804 const queue = sleeper.queue;
805 sleeper.queue = undefined;
806 queue.async(sleeper.fiber, &Fiber.@"resume");
807 queue.as_object().release();
808 }
809};
810
811const Mutex = struct {
812 state: State,
813 queue: c.dispatch.queue_t,
814 waiters: std.DoublyLinkedList,
815
816 const State = packed struct(usize) {
817 locked: bool,
818 num_waiters: NumWaiters,
819
820 const NumWaiters = @Int(.unsigned, @bitSizeOf(usize) - 1);
821 };
822
823 const Waiter = struct {
824 sleeper: Sleeper = undefined,
825 cancelable: Cancelable,
826 mutex: *Mutex,
827 node: std.DoublyLinkedList.Node = .{},
828
829 fn add(context: ?*anyopaque) callconv(.c) void {
830 const waiter: *Waiter = @ptrCast(@alignCast(context));
831 waiter.cancelable.enter(waiter.sleeper.fiber) catch |err| switch (err) {
832 error.CancelRequested => return waiter.wake(),
833 };
834 var state = @atomicRmw(State, &waiter.mutex.state, .Add, .{
835 .locked = false,
836 .num_waiters = 1,
837 }, .monotonic);
838 state.num_waiters += 1;
839 while (!state.locked) {
840 @branchHint(.unlikely);
841 state = @cmpxchgWeak(State, &waiter.mutex.state, state, .{
842 .locked = true,
843 .num_waiters = state.num_waiters - 1,
844 }, .acquire, .monotonic) orelse break;
845 } else return waiter.mutex.waiters.append(&waiter.node);
846 waiter.cancelable.leave(waiter.sleeper.fiber) catch |err| switch (err) {
847 error.CancelRequested => {
848 waiter.node.next = &waiter.node;
849 return;
850 },
851 };
852 waiter.wake();
853 }
854
855 fn canceled(context: ?*anyopaque) callconv(.c) void {
856 const cancelable: *Cancelable = @ptrCast(@alignCast(context));
857 const waiter: *Waiter = @fieldParentPtr("cancelable", cancelable);
858 cancelable.requested(waiter.sleeper.fiber);
859 const mutex = waiter.mutex;
860 if (waiter.node.next != &waiter.node) {
861 @branchHint(.likely);
862 mutex.waiters.remove(&waiter.node);
863 assert(@atomicRmw(State, &mutex.state, .Sub, .{
864 .locked = false,
865 .num_waiters = 1,
866 }, .monotonic).num_waiters >= 1);
867 }
868 waiter.node = undefined;
869 waiter.wake();
870 }
871
872 fn remove(context: ?*anyopaque) callconv(.c) void {
873 const mutex: *Mutex = @ptrCast(@alignCast(context));
874 var state = @atomicLoad(State, &mutex.state, .monotonic);
875 while (!state.locked and state.num_waiters > 0) {
876 @branchHint(.likely);
877 state = @cmpxchgWeak(State, &mutex.state, state, .{
878 .locked = true,
879 .num_waiters = state.num_waiters - 1,
880 }, .acquire, .monotonic) orelse break;
881 } else return;
882 var num_removed: State.NumWaiters = 0;
883 while (mutex.waiters.popFirst()) |node| {
884 @branchHint(.likely);
885 const waiter: *Waiter = @fieldParentPtr("node", node);
886 node.* = undefined;
887 waiter.cancelable.leave(waiter.sleeper.fiber) catch |err| switch (err) {
888 error.CancelRequested => {
889 num_removed += 1;
890 node.next = node;
891 continue;
892 },
893 };
894 break;
895 }
896 if (num_removed > 0) {
897 @branchHint(.unlikely);
898 assert(@atomicRmw(State, &mutex.state, .Sub, .{
899 .locked = false,
900 .num_waiters = num_removed,
901 }, .monotonic).num_waiters >= num_removed);
902 }
903 }
904
905 fn wake(waiter: *Waiter) void {
906 Sleeper.wake(&waiter.sleeper);
907 }
908 };
909
910 fn init(mutex: *Mutex, queue: c.dispatch.queue_t) error{SystemResources}!void {
911 mutex.* = .{
912 .state = .{ .locked = false, .num_waiters = 0 },
913 .queue = c.dispatch.queue_create_with_target(
914 "org.ziglang.std.Io.Dispatch.Mutex",
915 .SERIAL(),
916 queue,
917 ) orelse return error.SystemResources,
918 .waiters = .{},
919 };
920 }
921
922 fn deinit(mutex: *Mutex) void {
923 assert(mutex.state == State{ .locked = false, .num_waiters = 0 });
924 assert(mutex.waiters.first == null and mutex.waiters.last == null);
925 mutex.queue.as_object().release();
926 mutex.* = undefined;
927 }
928
929 fn tryLock(mutex: *Mutex) bool {
930 const state =
931 @atomicRmw(State, &mutex.state, .Or, .{ .locked = true, .num_waiters = 0 }, .acquire);
932 if (state.locked) {
933 @branchHint(.unlikely);
934 }
935 return !state.locked;
936 }
937
938 fn lock(mutex: *Mutex, ev: *Evented) Io.Cancelable!void {
939 if (mutex.tryLock()) return;
940 var waiter: Waiter = .{
941 .cancelable = .{ .queue = mutex.queue, .cancel = &Mutex.Waiter.canceled },
942 .mutex = mutex,
943 };
944 ev.yield(.{ .mutex_wait = &waiter });
945 try waiter.cancelable.acknowledge(waiter.sleeper.fiber);
946 }
947
948 fn lockUncancelable(mutex: *Mutex, ev: *Evented) void {
949 if (mutex.tryLock()) return;
950 var waiter: Waiter = .{ .cancelable = .blocked, .mutex = mutex };
951 ev.yield(.{ .mutex_wait = &waiter });
952 waiter.cancelable.acknowledge(waiter.sleeper.fiber) catch |err| switch (err) {
953 error.Canceled => unreachable, // blocked
954 };
955 }
956
957 fn unlock(mutex: *Mutex) void {
958 const state = @atomicRmw(State, &mutex.state, .And, .{
959 .locked = false,
960 .num_waiters = std.math.maxInt(State.NumWaiters),
961 }, .release);
962 if (state.num_waiters > 0) {
963 @branchHint(.unlikely);
964 mutex.queue.async(mutex, &Waiter.remove);
965 }
966 }
967};
968
969fn crashHandler(userdata: ?*anyopaque) void {
970 const ev: *Evented = @ptrCast(@alignCast(userdata));
971 _ = ev;
972 const thread = &Thread.self;
973 if (thread.current_context == null) std.process.abort();
974 if (thread.current_context == &thread.main_context) std.process.abort();
975 const fiber = thread.currentFiber();
976 @atomicStore(
977 Fiber.CancelStatus,
978 &fiber.cancel_status,
979 .{ .requested = true, .awaiting = .nothing },
980 .monotonic,
981 );
982 fiber.cancel_protection = .{ .user = .blocked, .acknowledged = true };
983}
984
985const AsyncClosure = struct {
986 evented: *Evented,
987 fiber: *Fiber,
988 start: *const fn (context: *const anyopaque, result: *anyopaque) void,
989 result_align: Alignment,
990
991 fn fromFiber(fiber: *Fiber) *AsyncClosure {
992 return @ptrFromInt(Fiber.max_context_align.max(.of(AsyncClosure)).backward(
993 @intFromPtr(fiber.allocatedEnd()) - Fiber.max_context_size,
994 ) - @sizeOf(AsyncClosure));
995 }
996
997 fn contextPointer(closure: *AsyncClosure) [*]align(Fiber.max_context_align.toByteUnits()) u8 {
998 return @alignCast(@as([*]u8, @ptrCast(closure)) + @sizeOf(AsyncClosure));
999 }
1000
1001 fn entry() callconv(.naked) void {
1002 switch (builtin.cpu.arch) {
1003 .aarch64 => asm volatile (
1004 \\ mov x0, sp
1005 \\ b %[call]
1006 :
1007 : [call] "X" (&call),
1008 ),
1009 .x86_64 => asm volatile (
1010 \\ leaq 8(%%rsp), %%rdi
1011 \\ jmp %[call:P]
1012 :
1013 : [call] "X" (&call),
1014 ),
1015 else => |arch| @compileError("unimplemented architecture: " ++ @tagName(arch)),
1016 }
1017 }
1018
1019 fn call(
1020 closure: *AsyncClosure,
1021 message: *const SwitchMessage,
1022 ) callconv(.withStackAlign(.c, @alignOf(AsyncClosure))) noreturn {
1023 const ev = closure.evented;
1024 const fiber = closure.fiber;
1025 message.handle(ev);
1026 closure.start(closure.contextPointer(), fiber.resultBytes(closure.result_align));
1027 if (@atomicRmw(?*Fiber, &fiber.link.awaiter, .Xchg, Fiber.finished, .acq_rel)) |awaiter|
1028 ev.queue.async(awaiter, &Fiber.@"resume");
1029 ev.yield(.nothing);
1030 unreachable; // switched to dead fiber
1031 }
1032};
1033
1034fn async(
1035 userdata: ?*anyopaque,
1036 result: []u8,
1037 result_alignment: Alignment,
1038 context: []const u8,
1039 context_alignment: Alignment,
1040 start: *const fn (context: *const anyopaque, result: *anyopaque) void,
1041) ?*std.Io.AnyFuture {
1042 const ev: *Evented = @ptrCast(@alignCast(userdata));
1043 return concurrent(ev, result.len, result_alignment, context, context_alignment, start) catch {
1044 start(context.ptr, result.ptr);
1045 return null;
1046 };
1047}
1048
1049fn concurrent(
1050 userdata: ?*anyopaque,
1051 result_len: usize,
1052 result_alignment: Alignment,
1053 context: []const u8,
1054 context_alignment: Alignment,
1055 start: *const fn (context: *const anyopaque, result: *anyopaque) void,
1056) Io.ConcurrentError!*std.Io.AnyFuture {
1057 assert(result_alignment.compare(.lte, Fiber.max_result_align)); // TODO
1058 assert(context_alignment.compare(.lte, Fiber.max_context_align)); // TODO
1059 assert(result_len <= Fiber.max_result_size); // TODO
1060 assert(context.len <= Fiber.max_context_size); // TODO
1061
1062 const ev: *Evented = @ptrCast(@alignCast(userdata));
1063 const fiber = Fiber.create(ev) catch |err| switch (err) {
1064 error.OutOfMemory => return error.ConcurrencyUnavailable,
1065 };
1066
1067 const closure: *AsyncClosure = .fromFiber(fiber);
1068 fiber.* = .{
1069 .required_align = {},
1070 .evented = ev,
1071 .context = switch (builtin.cpu.arch) {
1072 .aarch64 => .{
1073 .sp = @intFromPtr(closure),
1074 .fp = 0,
1075 .pc = @intFromPtr(&AsyncClosure.entry),
1076 },
1077 .x86_64 => .{
1078 .rsp = @intFromPtr(closure) - 8,
1079 .rbp = 0,
1080 .rip = @intFromPtr(&AsyncClosure.entry),
1081 },
1082 else => |arch| @compileError("unimplemented architecture: " ++ @tagName(arch)),
1083 },
1084 .link = .{ .awaiter = null },
1085 .awaiting_group = undefined,
1086 .cancel_status = .unrequested,
1087 .cancel_protection = .unblocked,
1088 };
1089 closure.* = .{
1090 .evented = ev,
1091 .fiber = fiber,
1092 .start = start,
1093 .result_align = result_alignment,
1094 };
1095 @memcpy(closure.contextPointer(), context);
1096
1097 ev.queue.async(fiber, &Fiber.@"resume");
1098 return @ptrCast(fiber);
1099}
1100
1101fn await(
1102 userdata: ?*anyopaque,
1103 future: *std.Io.AnyFuture,
1104 result: []u8,
1105 result_alignment: Alignment,
1106) void {
1107 const ev: *Evented = @ptrCast(@alignCast(userdata));
1108 const awaiting: *Fiber = @ptrCast(@alignCast(future));
1109 if (@atomicLoad(?*Fiber, &awaiting.link.awaiter, .acquire) != Fiber.finished)
1110 ev.yield(.{ .await = awaiting });
1111 @memcpy(result, awaiting.resultBytes(result_alignment));
1112 awaiting.destroy(ev);
1113}
1114
1115fn cancel(
1116 userdata: ?*anyopaque,
1117 future: *std.Io.AnyFuture,
1118 result: []u8,
1119 result_alignment: Alignment,
1120) void {
1121 const ev: *Evented = @ptrCast(@alignCast(userdata));
1122 const future_fiber: *Fiber = @ptrCast(@alignCast(future));
1123 future_fiber.requestCancel(ev);
1124 await(ev, future, result, result_alignment);
1125}
1126
1127const Group = struct {
1128 ptr: *Io.Group,
1129
1130 const List = packed struct(usize) {
1131 cancel_requested: bool,
1132 awaiter_delayed: bool,
1133 fibers: Fiber.PackedPtr,
1134 };
1135 fn listPtr(group: Group) *List {
1136 return @ptrCast(&group.ptr.token);
1137 }
1138
1139 const Mutex = packed struct(u32) {
1140 locked: bool,
1141 contended: bool,
1142 shared2: u30,
1143 };
1144 fn mutexPtr(group: Group) *Group.Mutex {
1145 return switch (comptime builtin.cpu.arch.endian()) {
1146 .little => @ptrCast(&group.ptr.state),
1147 .big => @ptrCast(@alignCast(
1148 @as([*]u8, @ptrCast(&group.ptr.state)) + @sizeOf(usize) - @sizeOf(u32),
1149 )),
1150 };
1151 }
1152
1153 const Awaiter = packed struct(usize) {
1154 locked: bool,
1155 contended: bool,
1156 awaiter: Fiber.PackedPtr,
1157 };
1158 fn awaiterPtr(group: Group) *Awaiter {
1159 return @ptrCast(&group.ptr.state);
1160 }
1161
1162 fn lock(group: Group, ev: *Evented) void {
1163 const mutex = group.mutexPtr();
1164 {
1165 const old_state = @atomicRmw(
1166 Group.Mutex,
1167 mutex,
1168 .Or,
1169 .{ .locked = true, .contended = false, .shared2 = 0 },
1170 .acquire,
1171 );
1172 if (!old_state.locked) {
1173 @branchHint(.likely);
1174 return;
1175 }
1176 if (old_state.contended) {
1177 futexWaitUncancelable(ev, @ptrCast(mutex), @bitCast(old_state));
1178 }
1179 }
1180 while (true) {
1181 var old_state = @atomicRmw(
1182 Group.Mutex,
1183 mutex,
1184 .Or,
1185 .{ .locked = true, .contended = true, .shared2 = 0 },
1186 .acquire,
1187 );
1188 if (!old_state.locked) {
1189 @branchHint(.likely);
1190 return;
1191 }
1192 old_state.contended = true;
1193 futexWaitUncancelable(ev, @ptrCast(mutex), @bitCast(old_state));
1194 }
1195 }
1196
1197 fn unlock(group: Group, ev: *Evented) void {
1198 const mutex = group.mutexPtr();
1199 const old_state = @atomicRmw(
1200 Group.Mutex,
1201 mutex,
1202 .And,
1203 .{ .locked = false, .contended = false, .shared2 = std.math.maxInt(u30) },
1204 .release,
1205 );
1206 assert(old_state.locked);
1207 if (old_state.contended) futexWake(ev, @ptrCast(mutex), 1);
1208 }
1209
1210 fn addFiber(group: Group, ev: *Evented, fiber: *Fiber) void {
1211 group.lock(ev);
1212 defer group.unlock(ev);
1213 const list_ptr = group.listPtr();
1214 const list = @atomicLoad(List, list_ptr, .monotonic);
1215 if (list.cancel_requested) fiber.cancel_status = .{ .requested = true, .awaiting = .nothing };
1216 const old_head = list.fibers.unpack();
1217 if (old_head) |head| head.link.group.prev = fiber;
1218 fiber.link.group.next = old_head;
1219 @atomicStore(List, list_ptr, .{
1220 .cancel_requested = list.cancel_requested,
1221 .awaiter_delayed = list.awaiter_delayed,
1222 .fibers = .pack(fiber),
1223 }, .monotonic);
1224 }
1225
1226 fn removeFiber(group: Group, ev: *Evented, fiber: *Fiber) ?*Fiber {
1227 group.lock(ev);
1228 defer group.unlock(ev);
1229 const list_ptr = group.listPtr();
1230 const list = @atomicLoad(List, list_ptr, .monotonic);
1231 if (fiber.link.group.next) |next| next.link.group.prev = fiber.link.group.prev;
1232 if (fiber.link.group.prev) |prev| {
1233 prev.link.group.next = fiber.link.group.next;
1234 } else if (fiber.link.group.next) |new_head| {
1235 @atomicStore(List, list_ptr, .{
1236 .cancel_requested = list.cancel_requested,
1237 .awaiter_delayed = list.awaiter_delayed,
1238 .fibers = .pack(new_head),
1239 }, .monotonic);
1240 } else if (@atomicLoad(Awaiter, group.awaiterPtr(), .monotonic).awaiter.unpack()) |awaiter| {
1241 if (!awaiter.cancel_status.changeAwaiting(.group, .nothing) or list.cancel_requested) {
1242 @atomicStore(List, list_ptr, .{
1243 .cancel_requested = false,
1244 .awaiter_delayed = false,
1245 .fibers = .null,
1246 }, .release);
1247 assert(awaiter.awaiting_group.ptr == group.ptr);
1248 awaiter.awaiting_group = undefined;
1249 return awaiter;
1250 }
1251 // Race with `Fiber.requestCancel`
1252 @atomicStore(List, list_ptr, .{
1253 .cancel_requested = false,
1254 .awaiter_delayed = true,
1255 .fibers = .null,
1256 }, .monotonic);
1257 } else @atomicStore(List, list_ptr, .{
1258 .cancel_requested = false,
1259 .awaiter_delayed = false,
1260 .fibers = .null,
1261 }, .release);
1262 return null;
1263 }
1264
1265 fn await(group: Group, ev: *Evented, awaiter: *Fiber) bool {
1266 group.lock(ev);
1267 defer group.unlock(ev);
1268 if (@atomicLoad(List, group.listPtr(), .monotonic).fibers.unpack()) |_| {
1269 if (group.registerAwaiter(awaiter) and awaiter.cancel_protection.check() == .unblocked) {
1270 // The awaiter already had an unacknowledged cancelation request before
1271 // attempting to await a group, so propagate the cancelation to the group.
1272 assert(!group.cancelLocked(ev, null));
1273 }
1274 return false;
1275 }
1276 return true;
1277 }
1278
1279 fn cancel(group: Group, ev: *Evented, maybe_awaiter: ?*Fiber) bool {
1280 group.lock(ev);
1281 defer group.unlock(ev);
1282 return group.cancelLocked(ev, maybe_awaiter);
1283 }
1284
1285 /// Assumes the mutex is held.
1286 fn cancelLocked(group: Group, ev: *Evented, maybe_awaiter: ?*Fiber) bool {
1287 const list_ptr = group.listPtr();
1288 const list = @atomicRmw(
1289 List,
1290 list_ptr,
1291 .Add,
1292 .{ .cancel_requested = true, .awaiter_delayed = false, .fibers = .null },
1293 .monotonic,
1294 );
1295 assert(!list.cancel_requested);
1296 if (list.fibers.unpack()) |head| {
1297 var maybe_fiber: ?*Fiber = head;
1298 while (maybe_fiber) |fiber| {
1299 fiber.requestCancel(ev);
1300 maybe_fiber = fiber.link.group.next;
1301 }
1302 if (maybe_awaiter) |awaiter| _ = group.registerAwaiter(awaiter);
1303 return false;
1304 }
1305 @atomicStore(
1306 List,
1307 list_ptr,
1308 .{ .cancel_requested = false, .awaiter_delayed = false, .fibers = .null },
1309 .release,
1310 );
1311 return if (maybe_awaiter) |_| true else list.awaiter_delayed;
1312 }
1313
1314 /// Assumes the mutex is held.
1315 fn registerAwaiter(group: Group, awaiter: *Fiber) bool {
1316 awaiter.awaiting_group = group;
1317 assert(@atomicRmw(
1318 Awaiter,
1319 group.awaiterPtr(),
1320 .Add,
1321 .{ .locked = false, .contended = false, .awaiter = .pack(awaiter) },
1322 .monotonic,
1323 ).awaiter == .null);
1324 return awaiter.cancel_status.changeAwaiting(.nothing, .group);
1325 }
1326
1327 const AsyncClosure = struct {
1328 evented: *Evented,
1329 group: Group,
1330 fiber: *Fiber,
1331 start: *const fn (context: *const anyopaque) void,
1332
1333 fn fromFiber(fiber: *Fiber) *Group.AsyncClosure {
1334 return @ptrFromInt(Fiber.max_context_align.max(.of(Group.AsyncClosure)).backward(
1335 @intFromPtr(fiber.allocatedEnd()) - Fiber.max_context_size,
1336 ) - @sizeOf(Group.AsyncClosure));
1337 }
1338
1339 fn contextPointer(
1340 closure: *Group.AsyncClosure,
1341 ) [*]align(Fiber.max_context_align.toByteUnits()) u8 {
1342 return @alignCast(@as([*]u8, @ptrCast(closure)) + @sizeOf(Group.AsyncClosure));
1343 }
1344
1345 fn entry() callconv(.naked) void {
1346 switch (builtin.cpu.arch) {
1347 .aarch64 => asm volatile (
1348 \\ mov x0, sp
1349 \\ b %[call]
1350 :
1351 : [call] "X" (&call),
1352 ),
1353 .x86_64 => asm volatile (
1354 \\ leaq 8(%%rsp), %%rdi
1355 \\ jmp %[call:P]
1356 :
1357 : [call] "X" (&call),
1358 ),
1359 else => |arch| @compileError("unimplemented architecture: " ++ @tagName(arch)),
1360 }
1361 }
1362
1363 fn call(
1364 closure: *Group.AsyncClosure,
1365 message: *const SwitchMessage,
1366 ) callconv(.withStackAlign(.c, @alignOf(Group.AsyncClosure))) noreturn {
1367 const ev = closure.evented;
1368 const fiber = closure.fiber;
1369 message.handle(ev);
1370 closure.start(closure.contextPointer());
1371 if (closure.group.removeFiber(ev, fiber)) |awaiter| ev.queue.async(awaiter, &Fiber.@"resume");
1372 ev.yield(.destroy);
1373 unreachable; // switched to dead fiber
1374 }
1375 };
1376};
1377
1378fn groupAsync(
1379 userdata: ?*anyopaque,
1380 type_erased: *Io.Group,
1381 context: []const u8,
1382 context_alignment: Alignment,
1383 start: *const fn (context: *const anyopaque) void,
1384) void {
1385 const ev: *Evented = @ptrCast(@alignCast(userdata));
1386 return groupConcurrent(ev, type_erased, context, context_alignment, start) catch {
1387 start(context.ptr);
1388 };
1389}
1390
1391fn groupConcurrent(
1392 userdata: ?*anyopaque,
1393 type_erased: *Io.Group,
1394 context: []const u8,
1395 context_alignment: Alignment,
1396 start: *const fn (context: *const anyopaque) void,
1397) Io.ConcurrentError!void {
1398 assert(context_alignment.compare(.lte, Fiber.max_context_align)); // TODO
1399 assert(context.len <= Fiber.max_context_size); // TODO
1400
1401 const ev: *Evented = @ptrCast(@alignCast(userdata));
1402 const group: Group = .{ .ptr = type_erased };
1403 const fiber = Fiber.create(ev) catch |err| switch (err) {
1404 error.OutOfMemory => return error.ConcurrencyUnavailable,
1405 };
1406
1407 const closure: *Group.AsyncClosure = .fromFiber(fiber);
1408 fiber.* = .{
1409 .required_align = {},
1410 .evented = ev,
1411 .context = switch (builtin.cpu.arch) {
1412 .aarch64 => .{
1413 .sp = @intFromPtr(closure),
1414 .fp = 0,
1415 .pc = @intFromPtr(&Group.AsyncClosure.entry),
1416 },
1417 .x86_64 => .{
1418 .rsp = @intFromPtr(closure) - 8,
1419 .rbp = 0,
1420 .rip = @intFromPtr(&Group.AsyncClosure.entry),
1421 },
1422 else => |arch| @compileError("unimplemented architecture: " ++ @tagName(arch)),
1423 },
1424 .link = .{ .group = .{ .prev = null, .next = null } },
1425 .awaiting_group = undefined,
1426 .cancel_status = .unrequested,
1427 .cancel_protection = .unblocked,
1428 };
1429 closure.* = .{
1430 .evented = ev,
1431 .group = group,
1432 .fiber = fiber,
1433 .start = start,
1434 };
1435 @memcpy(closure.contextPointer(), context);
1436 group.addFiber(ev, fiber);
1437 ev.queue.async(fiber, &Fiber.@"resume");
1438}
1439
1440fn groupAwait(
1441 userdata: ?*anyopaque,
1442 type_erased: *Io.Group,
1443 initial_token: *anyopaque,
1444) Io.Cancelable!void {
1445 const ev: *Evented = @ptrCast(@alignCast(userdata));
1446 _ = initial_token;
1447 ev.yield(.{ .group_await = .{ .ptr = type_erased } });
1448}
1449
1450fn groupCancel(userdata: ?*anyopaque, type_erased: *Io.Group, initial_token: *anyopaque) void {
1451 const ev: *Evented = @ptrCast(@alignCast(userdata));
1452 _ = initial_token;
1453 ev.yield(.{ .group_cancel = .{ .ptr = type_erased } });
1454}
1455
1456fn recancel(userdata: ?*anyopaque) void {
1457 const ev: *Evented = @ptrCast(@alignCast(userdata));
1458 _ = ev;
1459 Thread.current().currentFiber().cancel_protection.recancel();
1460}
1461
1462fn swapCancelProtection(userdata: ?*anyopaque, new: Io.CancelProtection) Io.CancelProtection {
1463 const ev: *Evented = @ptrCast(@alignCast(userdata));
1464 _ = ev;
1465 const cancel_protection = &Thread.current().currentFiber().cancel_protection;
1466 defer cancel_protection.user = new;
1467 return cancel_protection.user;
1468}
1469
1470fn checkCancel(userdata: ?*anyopaque) Io.Cancelable!void {
1471 const ev: *Evented = @ptrCast(@alignCast(userdata));
1472 _ = ev;
1473 const fiber = Thread.current().currentFiber();
1474 switch (fiber.cancel_protection.check()) {
1475 .unblocked => {
1476 const cancel_status = @atomicLoad(Fiber.CancelStatus, &fiber.cancel_status, .monotonic);
1477 assert(cancel_status.awaiting == .nothing);
1478 if (cancel_status.requested) {
1479 @branchHint(.unlikely);
1480 fiber.cancel_protection.acknowledge();
1481 return error.Canceled;
1482 }
1483 },
1484 .blocked => {},
1485 }
1486}
1487
1488const Futex = struct {
1489 num_waiters: usize,
1490 queue: c.dispatch.queue_t,
1491 waiters: std.DoublyLinkedList,
1492
1493 const Waiter = struct {
1494 sleeper: Sleeper = undefined,
1495 cancelable: Cancelable,
1496 futex: *Futex,
1497 node: std.DoublyLinkedList.Node = .{},
1498 ptr: *const u32,
1499 expected: u32,
1500 timeout: c.dispatch.time_t = .FOREVER,
1501 leeway: u64,
1502 timer: ?c.dispatch.source_t = null,
1503
1504 const already_signaled: c.dispatch.source_t = @ptrFromInt(1);
1505
1506 fn add(context: ?*anyopaque) callconv(.c) void {
1507 const waiter: *Waiter = @ptrCast(@alignCast(context));
1508 const futex = waiter.futex;
1509 _ = @atomicRmw(usize, &futex.num_waiters, .Add, 1, .acquire);
1510 waiter.tryAdd() catch |err| switch (err) {
1511 error.CancelRequested => {
1512 wake(waiter);
1513 assert(@atomicRmw(usize, &futex.num_waiters, .Sub, 1, .monotonic) >= 1);
1514 },
1515 };
1516 }
1517
1518 fn tryAdd(waiter: *Waiter) Cancelable.RequestedError!void {
1519 if (@atomicLoad(u32, waiter.ptr, .monotonic) != waiter.expected)
1520 return error.CancelRequested;
1521 try waiter.cancelable.enter(waiter.sleeper.fiber);
1522 const futex = waiter.futex;
1523 switch (waiter.timeout) {
1524 .FOREVER => {},
1525 else => |timeout| {
1526 const timer = c.dispatch.source_create(.TIMER, 0, .none, futex.queue) orelse {
1527 log.warn("failed to create timer for futex timeout", .{});
1528 return error.CancelRequested;
1529 };
1530 timer.as_object().set_context(waiter);
1531 timer.set_event_handler(&timedOut);
1532 timer.set_cancel_handler(&wake);
1533 timer.set_timer(timeout, c.dispatch.TIME_FOREVER, waiter.leeway);
1534 timer.as_object().activate();
1535 waiter.timer = timer;
1536 },
1537 }
1538 futex.waiters.append(&waiter.node);
1539 }
1540
1541 fn canceled(context: ?*anyopaque) callconv(.c) void {
1542 const cancelable: *Cancelable = @ptrCast(@alignCast(context));
1543 const waiter: *Waiter = @fieldParentPtr("cancelable", cancelable);
1544 cancelable.requested(waiter.sleeper.fiber);
1545 const futex = waiter.futex;
1546 waiter.remove();
1547 assert(@atomicRmw(usize, &futex.num_waiters, .Sub, 1, .monotonic) >= 1);
1548 }
1549
1550 fn timedOut(context: ?*anyopaque) callconv(.c) void {
1551 const waiter: *Waiter = @ptrCast(@alignCast(context));
1552 const futex = waiter.futex;
1553 waiter.tryRemove() catch |err| switch (err) {
1554 error.CancelRequested => return,
1555 };
1556 assert(@atomicRmw(usize, &futex.num_waiters, .Sub, 1, .monotonic) >= 1);
1557 }
1558
1559 fn tryRemove(waiter: *Waiter) Cancelable.RequestedError!void {
1560 try waiter.cancelable.leave(waiter.sleeper.fiber);
1561 waiter.remove();
1562 }
1563
1564 fn remove(waiter: *Waiter) void {
1565 waiter.futex.waiters.remove(&waiter.node);
1566 if (waiter.timer) |timer| timer.cancel() else wake(waiter);
1567 }
1568
1569 fn wake(context: ?*anyopaque) callconv(.c) void {
1570 const waiter: *Waiter = @ptrCast(@alignCast(context));
1571 if (waiter.timer) |timer| timer.as_object().release();
1572 Sleeper.wake(&waiter.sleeper);
1573 }
1574 };
1575
1576 const Waker = struct {
1577 sleeper: Sleeper = undefined,
1578 futex: *Futex,
1579 ptr: *const u32,
1580 max_waiters: u32,
1581
1582 fn remove(context: ?*anyopaque) callconv(.c) void {
1583 const waker: *Waker = @ptrCast(@alignCast(context));
1584 const futex = waker.futex;
1585 const ptr = waker.ptr;
1586 const max_waiters = waker.max_waiters;
1587
1588 var num_removed: usize = 0;
1589 var next_node = futex.waiters.first;
1590 while (num_removed < max_waiters) {
1591 const waiter: *Waiter = @fieldParentPtr("node", next_node orelse break);
1592 next_node = waiter.node.next;
1593 if (waiter.ptr != ptr) {
1594 @branchHint(.unlikely);
1595 continue;
1596 }
1597 waiter.tryRemove() catch |err| switch (err) {
1598 error.CancelRequested => continue,
1599 };
1600 num_removed += 1;
1601 }
1602 assert(@atomicRmw(usize, &futex.num_waiters, .Sub, num_removed, .monotonic) >= num_removed);
1603
1604 var sleeper = waker.sleeper;
1605 waker.* = undefined;
1606 Sleeper.wake(&sleeper);
1607 }
1608 };
1609
1610 fn init(futex: *Futex, queue: c.dispatch.queue_t) error{SystemResources}!void {
1611 futex.* = .{
1612 .num_waiters = 0,
1613 .queue = c.dispatch.queue_create_with_target(
1614 "org.ziglang.std.Io.Dispatch.Futex",
1615 .SERIAL(),
1616 queue,
1617 ) orelse return error.SystemResources,
1618 .waiters = .{},
1619 };
1620 }
1621
1622 fn deinit(futex: *Futex) void {
1623 assert(futex.num_waiters == 0 and futex.waiters.first == null and futex.waiters.last == null);
1624 futex.queue.as_object().release();
1625 futex.* = undefined;
1626 }
1627};
1628
1629fn futexForAddress(ev: *Evented, address: usize) *Futex {
1630 // Here we use Fibonacci hashing: the golden ratio can be used to evenly redistribute input
1631 // values across a range, giving a poor, but extremely quick to compute, hash.
1632
1633 // This literal is the rounded value of '2^64 / phi' (where 'phi' is the golden ratio). The
1634 // shift then converts it to '2^b / phi', where 'b' is the pointer bit width.
1635 const fibonacci_multiplier = 0x9E3779B97F4A7C15 >> (64 - @bitSizeOf(usize));
1636 const hashed = address *% fibonacci_multiplier;
1637 comptime assert(std.math.isPowerOfTwo(ev.futexes.len));
1638 // The high bits of `hashed` have better entropy than the low bits.
1639 return &ev.futexes[hashed >> @clz(ev.futexes.len - 1)];
1640}
1641
1642fn futexWait(
1643 userdata: ?*anyopaque,
1644 ptr: *const u32,
1645 expected: u32,
1646 timeout: Io.Timeout,
1647) Io.Cancelable!void {
1648 const ev: *Evented = @ptrCast(@alignCast(userdata));
1649 const futex = ev.futexForAddress(@intFromPtr(ptr));
1650 var waiter: Futex.Waiter = .{
1651 .cancelable = .{ .queue = futex.queue, .cancel = &Futex.Waiter.canceled },
1652 .futex = futex,
1653 .ptr = ptr,
1654 .expected = expected,
1655 .timeout = ev.timeFromTimeout(timeout),
1656 .leeway = ev.leeway,
1657 };
1658 ev.yield(.{ .futex_wait = &waiter });
1659 try waiter.cancelable.acknowledge(waiter.sleeper.fiber);
1660}
1661
1662fn futexWaitUncancelable(userdata: ?*anyopaque, ptr: *const u32, expected: u32) void {
1663 const ev: *Evented = @ptrCast(@alignCast(userdata));
1664 const futex = ev.futexForAddress(@intFromPtr(ptr));
1665 var waiter: Futex.Waiter = .{
1666 .cancelable = .blocked,
1667 .futex = futex,
1668 .ptr = ptr,
1669 .expected = expected,
1670 .leeway = ev.leeway,
1671 };
1672 ev.yield(.{ .futex_wait = &waiter });
1673 waiter.cancelable.acknowledge(waiter.sleeper.fiber) catch |err| switch (err) {
1674 error.Canceled => unreachable, // blocked
1675 };
1676}
1677
1678fn futexWake(userdata: ?*anyopaque, ptr: *const u32, max_waiters: u32) void {
1679 const ev: *Evented = @ptrCast(@alignCast(userdata));
1680 if (max_waiters == 0) return;
1681 const futex = ev.futexForAddress(@intFromPtr(ptr));
1682 switch (@atomicRmw(usize, &futex.num_waiters, .Add, 0, .release)) {
1683 0 => return,
1684 else => {
1685 @branchHint(.unlikely);
1686 var waker: Futex.Waker = .{ .futex = futex, .ptr = ptr, .max_waiters = max_waiters };
1687 ev.yield(.{ .futex_wake = &waker });
1688 },
1689 }
1690}
1691
1692fn operate(userdata: ?*anyopaque, operation: Io.Operation) Io.Cancelable!Io.Operation.Result {
1693 const ev: *Evented = @ptrCast(@alignCast(userdata));
1694 switch (operation) {
1695 .file_read_streaming => |o| return .{
1696 .file_read_streaming = ev.fileReadStreaming(o.file, o.data) catch |err| switch (err) {
1697 error.Canceled => |e| return e,
1698 else => |e| e,
1699 },
1700 },
1701 .file_write_streaming => |o| return .{
1702 .file_write_streaming = ev.fileWriteStreaming(
1703 o.file,
1704 o.header,
1705 o.data,
1706 o.splat,
1707 ) catch |err| switch (err) {
1708 error.Canceled => |e| return e,
1709 else => |e| e,
1710 },
1711 },
1712 .device_io_control => |*o| return .{ .device_io_control = try deviceIoControl(o) },
1713 .net_receive => @panic("TODO implement net_receive operation"),
1714 .net_send => @panic("TODO implement net_send operation"),
1715 .net_read => @panic("TODO implement net_read operation"),
1716 .net_write => @panic("TODO implement net_write operation"),
1717 }
1718}
1719
1720fn fileReadStreaming(ev: *Evented, file: File, data: []const []u8) File.ReadStreamingError!usize {
1721 if (file.flags.nonblocking) nonblocking: {
1722 return fileReadStreamingLimit(file.handle, data, .unlimited) catch |err| switch (err) {
1723 error.WouldBlock => break :nonblocking,
1724 else => |e| return e,
1725 };
1726 }
1727 const source = c.dispatch.source_create(
1728 .READ,
1729 @bitCast(@as(isize, file.handle)),
1730 .none,
1731 ev.queue,
1732 ) orelse return error.SystemResources;
1733 source.as_object().set_context(Thread.current().currentFiber());
1734 source.set_event_handler(&Fiber.@"resume");
1735 ev.yield(.{ .activate = source.as_object() });
1736 const limit = source.get_data();
1737 source.as_object().release();
1738 while (true) return fileReadStreamingLimit(
1739 file.handle,
1740 data,
1741 .limited(limit),
1742 ) catch |err| switch (err) {
1743 error.WouldBlock => {
1744 ev.yield(.nothing);
1745 continue;
1746 },
1747 else => |e| return e,
1748 };
1749}
1750fn fileReadStreamingLimit(
1751 handle: File.Handle,
1752 data: []const []u8,
1753 limit: Io.Limit,
1754) File.ReadStreamingError!usize {
1755 var iovecs: [max_iovecs_len]iovec = undefined;
1756 var iovlen: iovlen_t = 0;
1757 // .nothing can mean that the write side has been closed,
1758 // in which case the buffer still needs to be drained
1759 var remaining = if (limit == .nothing) .unlimited else limit;
1760 for (data) |buf| addBuf(false, &iovecs, &iovlen, &remaining, buf);
1761 if (iovlen == 0) return 0;
1762 while (true) {
1763 const rc = c.readv(handle, &iovecs, iovlen);
1764 switch (c.errno(rc)) {
1765 .SUCCESS => return if (rc == 0) error.EndOfStream else @intCast(rc),
1766 .INTR => continue,
1767 .INVAL => |err| return errnoBug(err),
1768 .FAULT => |err| return errnoBug(err),
1769 .AGAIN => return error.WouldBlock,
1770 .BADF => |err| return errnoBug(err), // File descriptor used after closed
1771 .IO => return error.InputOutput,
1772 .ISDIR => return error.IsDir,
1773 .NOBUFS => return error.SystemResources,
1774 .NOMEM => return error.SystemResources,
1775 .NOTCONN => return error.SocketUnconnected,
1776 .CONNRESET => return error.ConnectionResetByPeer,
1777 else => |err| return unexpectedErrno(err),
1778 }
1779 }
1780}
1781
1782fn fileWriteStreaming(
1783 ev: *Evented,
1784 file: File,
1785 header: []const u8,
1786 data: []const []const u8,
1787 splat: usize,
1788) File.Writer.Error!usize {
1789 if (file.flags.nonblocking) nonblocking: {
1790 return fileWriteStreamingLimit(
1791 file.handle,
1792 header,
1793 data,
1794 splat,
1795 .unlimited,
1796 ) catch |err| switch (err) {
1797 error.WouldBlock => break :nonblocking,
1798 else => |e| return e,
1799 };
1800 }
1801 const source = c.dispatch.source_create(
1802 .WRITE,
1803 @bitCast(@as(isize, file.handle)),
1804 .none,
1805 ev.queue,
1806 ) orelse return error.SystemResources;
1807 source.as_object().set_context(Thread.current().currentFiber());
1808 source.set_event_handler(&Fiber.@"resume");
1809 ev.yield(.{ .activate = source.as_object() });
1810 const limit = source.get_data();
1811 source.as_object().release();
1812 while (true) return fileWriteStreamingLimit(
1813 file.handle,
1814 header,
1815 data,
1816 splat,
1817 .limited(limit),
1818 ) catch |err| switch (err) {
1819 error.WouldBlock => {
1820 ev.yield(.nothing);
1821 continue;
1822 },
1823 else => |e| return e,
1824 };
1825}
1826fn fileWriteStreamingLimit(
1827 handle: File.Handle,
1828 header: []const u8,
1829 data: []const []const u8,
1830 splat: usize,
1831 limit: Io.Limit,
1832) File.Writer.Error!usize {
1833 if (limit == .nothing) return 0;
1834 var iovecs: [max_iovecs_len]iovec_const = undefined;
1835 var iovlen: iovlen_t = 0;
1836 var remaining = limit;
1837 addBuf(true, &iovecs, &iovlen, &remaining, header);
1838 for (data[0 .. data.len - 1]) |bytes| addBuf(true, &iovecs, &iovlen, &remaining, bytes);
1839 const pattern = data[data.len - 1];
1840 var backup_buffer: [splat_buffer_size]u8 = undefined;
1841 if (iovecs.len - iovlen != 0 and remaining != .nothing) switch (splat) {
1842 0 => {},
1843 1 => addBuf(true, &iovecs, &iovlen, &remaining, pattern),
1844 else => switch (pattern.len) {
1845 0 => {},
1846 1 => {
1847 const splat_buffer = &backup_buffer;
1848 const memset_len = @min(splat_buffer.len, splat);
1849 const buf = splat_buffer[0..memset_len];
1850 @memset(buf, pattern[0]);
1851 addBuf(true, &iovecs, &iovlen, &remaining, buf);
1852 var remaining_splat = splat - buf.len;
1853 while (remaining_splat > splat_buffer.len and iovecs.len - iovlen != 0 and remaining != .nothing) {
1854 assert(buf.len == splat_buffer.len);
1855 addBuf(true, &iovecs, &iovlen, &remaining, splat_buffer);
1856 remaining_splat -= splat_buffer.len;
1857 }
1858 addBuf(true, &iovecs, &iovlen, &remaining, splat_buffer[0..@min(remaining_splat, splat_buffer.len)]);
1859 },
1860 else => for (0..@min(splat, iovecs.len - iovlen)) |_| {
1861 if (remaining == .nothing) break;
1862 addBuf(true, &iovecs, &iovlen, &remaining, pattern);
1863 },
1864 },
1865 };
1866 if (iovlen == 0) return 0;
1867 while (true) {
1868 const rc = c.writev(handle, &iovecs, iovlen);
1869 switch (c.errno(rc)) {
1870 .SUCCESS => return @intCast(rc),
1871 .INTR => continue,
1872 .INVAL => |err| return errnoBug(err),
1873 .FAULT => |err| return errnoBug(err),
1874 .AGAIN => return error.WouldBlock,
1875 .BADF => return error.NotOpenForWriting, // Can be a race condition.
1876 .DESTADDRREQ => |err| return errnoBug(err), // `connect` was never called.
1877 .DQUOT => return error.DiskQuota,
1878 .FBIG => return error.FileTooBig,
1879 .IO => return error.InputOutput,
1880 .NOSPC => return error.NoSpaceLeft,
1881 .PERM => return error.PermissionDenied,
1882 .PIPE => return error.BrokenPipe,
1883 .CONNRESET => |err| return errnoBug(err), // Not a socket handle.
1884 .BUSY => return error.DeviceBusy,
1885 else => |err| return unexpectedErrno(err),
1886 }
1887 }
1888}
1889
1890fn deviceIoControl(o: *const Io.Operation.DeviceIoControl) Io.Cancelable!i32 {
1891 while (true) {
1892 const rc = c.ioctl(o.file.handle, @bitCast(o.code), @intFromPtr(o.arg));
1893 switch (c.errno(rc)) {
1894 .SUCCESS => return rc,
1895 .INTR => {},
1896 else => |err| return -@as(i32, @backingInt(err)),
1897 }
1898 }
1899}
1900
1901const BatchWaiter = struct {
1902 sleeper: Sleeper,
1903 queue: c.dispatch.queue_t,
1904 timer: ?c.dispatch.source_t = null,
1905
1906 const already_signaled: c.dispatch.source_t = @ptrFromInt(1);
1907
1908 fn signal(context: ?*anyopaque) callconv(.c) void {
1909 const waiter: *BatchWaiter = @ptrCast(@alignCast(context));
1910 if (waiter.timer) |timer| {
1911 if (timer != already_signaled) timer.cancel();
1912 } else {
1913 waiter.timer = already_signaled;
1914 waiter.queue.async(waiter, &@"suspend");
1915 }
1916 }
1917
1918 fn @"suspend"(context: ?*anyopaque) callconv(.c) void {
1919 const waiter: *BatchWaiter = @ptrCast(@alignCast(context));
1920 if (waiter.timer) |timer| if (timer != already_signaled) timer.as_object().release();
1921 waiter.queue.as_object().@"suspend"();
1922 waiter.wake();
1923 }
1924
1925 fn wake(waiter: *BatchWaiter) void {
1926 var sleeper = waiter.sleeper;
1927 waiter.* = undefined;
1928 Sleeper.wake(&sleeper);
1929 }
1930};
1931
1932fn batchAwaitAsync(userdata: ?*anyopaque, batch: *Io.Batch) Io.Cancelable!void {
1933 const ev: *Evented = @ptrCast(@alignCast(userdata));
1934 const queue = ev.batchDrainSubmitted(batch, false) catch |err| switch (err) {
1935 error.ConcurrencyUnavailable => unreachable, // passed concurrency=false
1936 error.Canceled => |e| return e,
1937 } orelse return;
1938 if (batch.pending.head == .none) return;
1939 var waiter: BatchWaiter = .{
1940 .sleeper = .init(ev.queue, Thread.current().currentFiber()),
1941 .queue = queue,
1942 };
1943 if (batch.completed.head != .none) BatchWaiter.signal(&waiter);
1944 queue.as_object().set_context(&waiter);
1945 ev.yield(.{ .@"resume" = queue.as_object() });
1946}
1947
1948fn batchAwaitConcurrent(
1949 userdata: ?*anyopaque,
1950 batch: *Io.Batch,
1951 timeout: Io.Timeout,
1952) Io.Batch.AwaitConcurrentError!void {
1953 const ev: *Evented = @ptrCast(@alignCast(userdata));
1954 const queue = try ev.batchDrainSubmitted(batch, true) orelse return;
1955 if (batch.pending.head == .none) return;
1956 var waiter: BatchWaiter = .{
1957 .sleeper = .init(ev.queue, Thread.current().currentFiber()),
1958 .queue = queue,
1959 };
1960 if (batch.completed.head == .none) switch (timeout) {
1961 .none => {},
1962 else => {
1963 const timer = c.dispatch.source_create(.TIMER, 0, .none, queue) orelse
1964 return error.ConcurrencyUnavailable;
1965 assert(timer != BatchWaiter.already_signaled);
1966 timer.as_object().set_context(&waiter);
1967 timer.set_event_handler(&BatchWaiter.signal);
1968 timer.set_cancel_handler(&BatchWaiter.@"suspend");
1969 timer.set_timer(ev.timeFromTimeout(timeout), c.dispatch.TIME_FOREVER, ev.leeway);
1970 timer.as_object().activate();
1971 waiter.timer = timer;
1972 },
1973 } else BatchWaiter.signal(&waiter);
1974 queue.as_object().set_context(&waiter);
1975 ev.yield(.{ .@"resume" = queue.as_object() });
1976}
1977
1978fn batchCancel(userdata: ?*anyopaque, batch: *Io.Batch) void {
1979 const ev: *Evented = @ptrCast(@alignCast(userdata));
1980 var index = batch.pending.head;
1981 while (index != .none) {
1982 const storage = &batch.storage[index.toIndex()];
1983 const pending = &storage.pending;
1984 const operation_userdata: *BatchOperationUserdata = .fromErased(&pending.userdata);
1985 assert(operation_userdata.batch == batch);
1986 operation_userdata.source.cancel();
1987 }
1988 const queue: c.dispatch.queue_t = @ptrCast(batch.userdata orelse return);
1989 if (batch.pending.head != .none) {
1990 var waiter: BatchWaiter = .{
1991 .sleeper = .init(ev.queue, Thread.current().currentFiber()),
1992 .queue = queue,
1993 .timer = BatchWaiter.already_signaled,
1994 };
1995 if (batch.pending.head == .none) queue.async(&waiter, &BatchWaiter.signal);
1996 queue.as_object().set_context(&waiter);
1997 ev.yield(.{ .@"resume" = queue.as_object() });
1998 }
1999 batch.userdata = null;
2000}
2001
2002const BatchOperationUserdata = extern struct {
2003 batch: *Io.Batch,
2004 source: c.dispatch.source_t,
2005 operation: extern union {
2006 file_read_streaming: extern struct {
2007 data_ptr: [*]const []u8,
2008 data_len: usize,
2009 },
2010 file_write_streaming: extern struct {
2011 header_ptr: [*]const u8,
2012 header_len: usize,
2013 data_ptr: [*]const []const u8,
2014 data_len: usize,
2015 splat: usize,
2016
2017 fn header(operation: *const @This()) []const u8 {
2018 return operation.header_ptr[0..operation.header_len];
2019 }
2020
2021 fn data(operation: *const @This()) []const []const u8 {
2022 return operation.data_ptr[0..operation.data_len];
2023 }
2024 },
2025 },
2026
2027 const Erased = Io.Operation.Storage.Pending.Userdata;
2028
2029 comptime {
2030 assert(@sizeOf(BatchOperationUserdata) <= @sizeOf(Erased));
2031 }
2032
2033 fn toErased(userdata: *BatchOperationUserdata) *Erased {
2034 return @ptrCast(userdata);
2035 }
2036
2037 fn fromErased(erased: *Erased) *BatchOperationUserdata {
2038 return @ptrCast(erased);
2039 }
2040};
2041
2042/// If `concurrency` is false, `error.ConcurrencyUnavailable` is unreachable.
2043fn batchDrainSubmitted(
2044 ev: *Evented,
2045 batch: *Io.Batch,
2046 concurrency: bool,
2047) (Io.ConcurrentError || Io.Cancelable)!?c.dispatch.queue_t {
2048 var index = batch.submitted.head;
2049 if (index == .none) return @ptrCast(batch.userdata);
2050 errdefer batch.submitted.head = index;
2051 const maybe_queue: ?c.dispatch.queue_t = if (batch.userdata) |batch_userdata|
2052 @ptrCast(batch_userdata)
2053 else maybe_queue: {
2054 const queue = c.dispatch.queue_create_with_target(
2055 "org.ziglang.std.Io.Dispatch.Batch",
2056 .SERIAL(),
2057 ev.queue,
2058 ) orelse if (concurrency) return error.ConcurrencyUnavailable else break :maybe_queue null;
2059 queue.as_object().@"suspend"();
2060 batch.userdata = queue;
2061 break :maybe_queue queue;
2062 };
2063 while (index != .none) {
2064 const storage = &batch.storage[index.toIndex()];
2065 const next_index = storage.submission.node.next;
2066 if (@as(?Io.Operation.Result, result: {
2067 if (maybe_queue) |queue| switch (storage.submission.operation) {
2068 .file_read_streaming => |operation| {
2069 const data = for (operation.data, 0..) |buffer, data_index| {
2070 if (buffer.len > 0) break operation.data[data_index..];
2071 } else break :result .{ .file_read_streaming = 0 };
2072 const source = c.dispatch.source_create(
2073 .READ,
2074 @bitCast(@as(isize, operation.file.handle)),
2075 .none,
2076 queue,
2077 ) orelse break :result .{ .file_read_streaming = error.SystemResources };
2078 storage.* = .{ .pending = .{
2079 .node = .{ .prev = batch.pending.tail, .next = .none },
2080 .tag = .file_read_streaming,
2081 .userdata = undefined,
2082 } };
2083 const operation_userdata: *BatchOperationUserdata =
2084 .fromErased(&storage.pending.userdata);
2085 operation_userdata.* = .{
2086 .batch = batch,
2087 .source = source,
2088 .operation = .{ .file_read_streaming = .{
2089 .data_ptr = data.ptr,
2090 .data_len = data.len,
2091 } },
2092 };
2093 source.as_object().set_context(storage);
2094 source.set_event_handler(&batchSourceEvent);
2095 source.set_cancel_handler(&batchSourceCancel);
2096 source.as_object().activate();
2097 break :result null;
2098 },
2099 .file_write_streaming => |operation| {
2100 const data = for (operation.data, 0..) |buffer, data_index| {
2101 if (buffer.len > 0) break operation.data[data_index..];
2102 } else if (operation.header.len > 0)
2103 operation.data[0..1]
2104 else
2105 break :result .{ .file_write_streaming = 0 };
2106 const source = c.dispatch.source_create(
2107 .WRITE,
2108 @bitCast(@as(isize, operation.file.handle)),
2109 .none,
2110 queue,
2111 ) orelse break :result .{ .file_write_streaming = error.SystemResources };
2112 storage.* = .{ .pending = .{
2113 .node = .{ .prev = batch.pending.tail, .next = .none },
2114 .tag = .file_write_streaming,
2115 .userdata = undefined,
2116 } };
2117 const operation_userdata: *BatchOperationUserdata =
2118 .fromErased(&storage.pending.userdata);
2119 operation_userdata.* = .{
2120 .batch = batch,
2121 .source = source,
2122 .operation = .{ .file_write_streaming = .{
2123 .header_ptr = operation.header.ptr,
2124 .header_len = operation.header.len,
2125 .data_ptr = data.ptr,
2126 .data_len = data.len,
2127 .splat = operation.splat,
2128 } },
2129 };
2130 source.as_object().set_context(storage);
2131 source.set_event_handler(&batchSourceEvent);
2132 source.set_cancel_handler(&batchSourceCancel);
2133 source.as_object().activate();
2134 break :result null;
2135 },
2136 .device_io_control => {},
2137 .net_receive => @panic("TODO implement batched net_receive"),
2138 .net_read => @panic("TODO implement batched net_read"),
2139 .net_write => @panic("TODO implement batched net_write"),
2140 };
2141 if (concurrency) return error.ConcurrencyUnavailable;
2142 break :result try operate(ev, storage.submission.operation);
2143 })) |result| {
2144 switch (batch.completed.tail) {
2145 .none => batch.completed.head = index,
2146 else => |tail_index| batch.storage[tail_index.toIndex()].completion.node.next = index,
2147 }
2148 batch.completed.tail = index;
2149 storage.* = .{ .completion = .{ .node = .{ .next = .none }, .result = result } };
2150 } else {
2151 switch (batch.pending.tail) {
2152 .none => batch.pending.head = index,
2153 else => |tail_index| batch.storage[tail_index.toIndex()].pending.node.next = index,
2154 }
2155 batch.pending.tail = index;
2156 }
2157 index = next_index;
2158 }
2159 batch.submitted = .{ .head = .none, .tail = .none };
2160 return maybe_queue;
2161}
2162
2163fn batchSourceEvent(context: ?*anyopaque) callconv(.c) void {
2164 const storage: *Io.Operation.Storage = @ptrCast(@alignCast(context));
2165 const pending = &storage.pending;
2166 const operation_userdata: *BatchOperationUserdata = .fromErased(&pending.userdata);
2167 const batch = operation_userdata.batch;
2168 const source = operation_userdata.source;
2169 const index: Io.Operation.OptionalIndex = .fromIndex(storage - batch.storage.ptr);
2170 const result: Io.Operation.Result = result: switch (pending.tag) {
2171 .file_read_streaming => {
2172 const operation = &operation_userdata.operation.file_read_streaming;
2173 break :result .{ .file_read_streaming = fileReadStreamingLimit(
2174 @intCast(source.get_handle()),
2175 operation.data_ptr[0..operation.data_len],
2176 .limited(source.get_data()),
2177 ) catch |err| switch (err) {
2178 error.Canceled => return Thread.current().currentFiber().cancel_protection.recancel(),
2179 error.WouldBlock => return,
2180 else => |e| e,
2181 } };
2182 },
2183 .file_write_streaming => {
2184 const operation = &operation_userdata.operation.file_write_streaming;
2185 break :result .{ .file_write_streaming = fileWriteStreamingLimit(
2186 @intCast(source.get_handle()),
2187 operation.header_ptr[0..operation.header_len],
2188 operation.data_ptr[0..operation.data_len],
2189 operation.splat,
2190 .limited(source.get_data()),
2191 ) catch |err| switch (err) {
2192 error.Canceled => return Thread.current().currentFiber().cancel_protection.recancel(),
2193 error.WouldBlock => return,
2194 else => |e| e,
2195 } };
2196 },
2197 .device_io_control => unreachable,
2198 .net_receive => @panic("TODO implement batched net_receive"),
2199 .net_read => @panic("TODO implement batched net_read"),
2200 .net_write => @panic("TODO implement batched net_write"),
2201 };
2202
2203 switch (pending.node.prev) {
2204 .none => batch.pending.head = pending.node.next,
2205 else => |prev_index| batch.storage[prev_index.toIndex()].pending.node.next = pending.node.next,
2206 }
2207 switch (pending.node.next) {
2208 .none => batch.pending.tail = pending.node.prev,
2209 else => |next_index| batch.storage[next_index.toIndex()].pending.node.prev = pending.node.prev,
2210 }
2211
2212 switch (batch.completed.tail) {
2213 .none => batch.completed.head = index,
2214 else => |tail_index| batch.storage[tail_index.toIndex()].completion.node.next = index,
2215 }
2216 storage.* = .{ .completion = .{ .node = .{ .next = .none }, .result = result } };
2217 batch.completed.tail = index;
2218
2219 source.as_object().release();
2220 const queue: c.dispatch.queue_t = @ptrCast(batch.userdata);
2221 const waiter: *BatchWaiter = @ptrCast(@alignCast(queue.as_object().get_context()));
2222 BatchWaiter.signal(waiter);
2223}
2224
2225fn batchSourceCancel(context: ?*anyopaque) callconv(.c) void {
2226 const storage: *Io.Operation.Storage = @ptrCast(@alignCast(context));
2227 const pending = &storage.pending;
2228 const operation_userdata: *BatchOperationUserdata = .fromErased(&pending.userdata);
2229 const batch = operation_userdata.batch;
2230 const source = operation_userdata.source;
2231 const index: Io.Operation.OptionalIndex = .fromIndex(storage - batch.storage.ptr);
2232
2233 switch (pending.node.prev) {
2234 .none => batch.pending.head = pending.node.next,
2235 else => |prev_index| batch.storage[prev_index.toIndex()].pending.node.next = pending.node.next,
2236 }
2237 switch (pending.node.next) {
2238 .none => batch.pending.tail = pending.node.prev,
2239 else => |next_index| batch.storage[next_index.toIndex()].pending.node.prev = pending.node.prev,
2240 }
2241
2242 const tail_index = batch.unused.tail;
2243 switch (tail_index) {
2244 .none => batch.unused.head = index,
2245 else => batch.storage[tail_index.toIndex()].unused.next = index,
2246 }
2247 storage.* = .{ .unused = .{ .prev = tail_index, .next = .none } };
2248 batch.unused.tail = index;
2249
2250 source.as_object().release();
2251 if (batch.pending.head != .none) return;
2252 const queue: c.dispatch.queue_t = @ptrCast(batch.userdata);
2253 const waiter: *BatchWaiter = @ptrCast(@alignCast(queue.as_object().get_context()));
2254 queue.as_object().release();
2255 waiter.wake();
2256}
2257
2258fn dirCreateDir(
2259 userdata: ?*anyopaque,
2260 dir: Dir,
2261 sub_path: []const u8,
2262 permissions: Dir.Permissions,
2263) Dir.CreateDirError!void {
2264 const ev: *Evented = @ptrCast(@alignCast(userdata));
2265 _ = ev;
2266
2267 var path_buffer: [c.PATH_MAX]u8 = undefined;
2268 const sub_path_posix = try pathToPosix(sub_path, &path_buffer);
2269
2270 while (true) {
2271 switch (c.errno(c.mkdirat(dir.handle, sub_path_posix, permissions.toMode()))) {
2272 .SUCCESS => return,
2273 .INTR => {},
2274 .ACCES => return error.AccessDenied,
2275 .PERM => return error.PermissionDenied,
2276 .DQUOT => return error.DiskQuota,
2277 .EXIST => return error.PathAlreadyExists,
2278 .LOOP => return error.SymLinkLoop,
2279 .MLINK => return error.LinkQuotaExceeded,
2280 .NAMETOOLONG => return error.NameTooLong,
2281 .NOENT => return error.FileNotFound,
2282 .NOMEM => return error.SystemResources,
2283 .NOSPC => return error.NoSpaceLeft,
2284 .NOTDIR => return error.NotDir,
2285 .ROFS => return error.ReadOnlyFileSystem,
2286 .ILSEQ => return error.BadPathName,
2287 .BADF => |err| return errnoBug(err), // File descriptor used after closed.
2288 .FAULT => |err| return errnoBug(err),
2289 else => |err| return unexpectedErrno(err),
2290 }
2291 }
2292}
2293
2294fn dirCreateDirPath(
2295 userdata: ?*anyopaque,
2296 dir: Dir,
2297 sub_path: []const u8,
2298 permissions: Dir.Permissions,
2299) Dir.CreateDirPathError!Dir.CreatePathStatus {
2300 const ev: *Evented = @ptrCast(@alignCast(userdata));
2301
2302 var it = Dir.path.componentIterator(sub_path);
2303 var status: Dir.CreatePathStatus = .existed;
2304 var component = it.last() orelse return error.BadPathName;
2305 while (true) {
2306 if (dirCreateDir(ev, dir, component.path, permissions)) |_| {
2307 status = .created;
2308 } else |err| switch (err) {
2309 error.PathAlreadyExists => {
2310 // It is important to return an error if it's not a directory
2311 // because otherwise a dangling symlink could cause an infinite
2312 // loop.
2313 const fstat = try dirStatFile(ev, dir, component.path, .{});
2314 if (fstat.kind != .directory) return error.NotDir;
2315 },
2316 error.FileNotFound => |e| {
2317 component = it.previous() orelse return e;
2318 continue;
2319 },
2320 else => |e| return e,
2321 }
2322 component = it.next() orelse return status;
2323 }
2324}
2325
2326fn dirCreateDirPathOpen(
2327 userdata: ?*anyopaque,
2328 dir: Dir,
2329 sub_path: []const u8,
2330 permissions: Dir.Permissions,
2331 options: Dir.OpenOptions,
2332) Dir.CreateDirPathOpenError!Dir {
2333 const ev: *Evented = @ptrCast(@alignCast(userdata));
2334 return dirOpenDir(ev, dir, sub_path, options) catch |err| switch (err) {
2335 error.FileNotFound => {
2336 _ = try dirCreateDirPath(ev, dir, sub_path, permissions);
2337 return dirOpenDir(ev, dir, sub_path, options);
2338 },
2339 else => |e| return e,
2340 };
2341}
2342
2343fn dirOpenDir(
2344 userdata: ?*anyopaque,
2345 dir: Dir,
2346 sub_path: []const u8,
2347 options: Dir.OpenOptions,
2348) Dir.OpenError!Dir {
2349 const ev: *Evented = @ptrCast(@alignCast(userdata));
2350 _ = ev;
2351
2352 var path_buffer: [c.PATH_MAX]u8 = undefined;
2353 const sub_path_posix = try pathToPosix(sub_path, &path_buffer);
2354
2355 const flags: c.O = .{
2356 .ACCMODE = .RDONLY,
2357 .NOFOLLOW = !options.follow_symlinks,
2358 .DIRECTORY = true,
2359 .CLOEXEC = true,
2360 };
2361
2362 while (true) {
2363 const rc = c.openat(dir.handle, sub_path_posix, flags);
2364 switch (c.errno(rc)) {
2365 .SUCCESS => return .{ .handle = @intCast(rc) },
2366 .INTR => {},
2367 .INVAL => return error.BadPathName,
2368 .ACCES => return error.AccessDenied,
2369 .LOOP => return error.SymLinkLoop,
2370 .MFILE => return error.ProcessFdQuotaExceeded,
2371 .NAMETOOLONG => return error.NameTooLong,
2372 .NFILE => return error.SystemFdQuotaExceeded,
2373 .NODEV => return error.NoDevice,
2374 .NOENT => return error.FileNotFound,
2375 .NOMEM => return error.SystemResources,
2376 .NOTDIR => return error.NotDir,
2377 .PERM => return error.PermissionDenied,
2378 .NXIO => return error.NoDevice,
2379 .ILSEQ => return error.BadPathName,
2380 .FAULT => |err| return errnoBug(err),
2381 .BADF => |err| return errnoBug(err), // File descriptor used after closed.
2382 .BUSY => |err| return errnoBug(err), // O_EXCL not passed
2383 else => |err| return unexpectedErrno(err),
2384 }
2385 }
2386}
2387
2388fn dirStat(userdata: ?*anyopaque, dir: Dir) Dir.StatError!Dir.Stat {
2389 const ev: *Evented = @ptrCast(@alignCast(userdata));
2390 return fileStat(ev, .{
2391 .handle = dir.handle,
2392 .flags = .{ .nonblocking = false },
2393 });
2394}
2395
2396fn dirStatFile(
2397 userdata: ?*anyopaque,
2398 dir: Dir,
2399 sub_path: []const u8,
2400 options: Dir.StatFileOptions,
2401) Dir.StatFileError!File.Stat {
2402 const ev: *Evented = @ptrCast(@alignCast(userdata));
2403 _ = ev;
2404
2405 var path_buffer: [c.PATH_MAX]u8 = undefined;
2406 const sub_path_posix = try pathToPosix(sub_path, &path_buffer);
2407
2408 const flags: u32 = if (options.follow_symlinks) 0 else c.AT.SYMLINK_NOFOLLOW;
2409
2410 while (true) {
2411 var stat = std.mem.zeroes(c.Stat);
2412 switch (c.errno(c.fstatat(dir.handle, sub_path_posix, &stat, flags))) {
2413 .SUCCESS => return statFromPosix(&stat),
2414 .INTR => {},
2415 .INVAL => |err| return errnoBug(err),
2416 .BADF => |err| return errnoBug(err), // File descriptor used after closed.
2417 .NOMEM => return error.SystemResources,
2418 .ACCES => return error.AccessDenied,
2419 .PERM => return error.PermissionDenied,
2420 .FAULT => |err| return errnoBug(err),
2421 .NAMETOOLONG => return error.NameTooLong,
2422 .LOOP => return error.SymLinkLoop,
2423 .NOENT => return error.FileNotFound,
2424 .NOTDIR => return error.FileNotFound,
2425 .ILSEQ => return error.BadPathName,
2426 else => |err| return unexpectedErrno(err),
2427 }
2428 }
2429}
2430
2431fn dirAccess(
2432 userdata: ?*anyopaque,
2433 dir: Dir,
2434 sub_path: []const u8,
2435 options: Dir.AccessOptions,
2436) Dir.AccessError!void {
2437 const ev: *Evented = @ptrCast(@alignCast(userdata));
2438 _ = ev;
2439
2440 var path_buffer: [c.PATH_MAX]u8 = undefined;
2441 const sub_path_posix = try pathToPosix(sub_path, &path_buffer);
2442
2443 const flags: u32 = if (options.follow_symlinks) 0 else c.AT.SYMLINK_NOFOLLOW;
2444
2445 const mode: u32 =
2446 @as(u32, if (options.read) c.R_OK else 0) |
2447 @as(u32, if (options.write) c.W_OK else 0) |
2448 @as(u32, if (options.execute) c.X_OK else 0);
2449
2450 while (true) switch (c.errno(c.faccessat(dir.handle, sub_path_posix, mode, flags))) {
2451 .SUCCESS => return,
2452 .INTR => {},
2453 .ACCES => return error.AccessDenied,
2454 .PERM => return error.PermissionDenied,
2455 .ROFS => return error.ReadOnlyFileSystem,
2456 .LOOP => return error.SymLinkLoop,
2457 .TXTBSY => return error.FileBusy,
2458 .NOTDIR => return error.FileNotFound,
2459 .NOENT => return error.FileNotFound,
2460 .NAMETOOLONG => return error.NameTooLong,
2461 .INVAL => |err| return errnoBug(err),
2462 .FAULT => |err| return errnoBug(err),
2463 .IO => return error.InputOutput,
2464 .NOMEM => return error.SystemResources,
2465 .ILSEQ => return error.BadPathName,
2466 else => |err| return unexpectedErrno(err),
2467 };
2468}
2469
2470fn dirCreateFile(
2471 userdata: ?*anyopaque,
2472 dir: Dir,
2473 sub_path: []const u8,
2474 flags: Dir.CreateFileOptions,
2475) File.OpenError!File {
2476 const ev: *Evented = @ptrCast(@alignCast(userdata));
2477 _ = ev;
2478
2479 var path_buffer: [c.PATH_MAX]u8 = undefined;
2480 const sub_path_posix = try pathToPosix(sub_path, &path_buffer);
2481
2482 const os_flags: c.O = .{
2483 .ACCMODE = if (flags.read) .RDWR else .WRONLY,
2484 .NONBLOCK = flags.lock == .none or flags.lock_nonblocking,
2485 .SHLOCK = flags.lock == .shared,
2486 .EXLOCK = flags.lock == .exclusive,
2487 .CREAT = true,
2488 .TRUNC = flags.truncate,
2489 .EXCL = flags.exclusive,
2490 .CLOEXEC = true,
2491 };
2492
2493 const fd: c.fd_t = while (true) {
2494 const rc = c.openat(dir.handle, sub_path_posix, os_flags, flags.permissions.toMode());
2495 switch (c.errno(rc)) {
2496 .SUCCESS => break @intCast(rc),
2497 .INTR => {},
2498 .FAULT => |err| return errnoBug(err),
2499 .INVAL => return error.BadPathName,
2500 .BADF => |err| return errnoBug(err), // File descriptor used after closed.
2501 .ACCES => return error.AccessDenied,
2502 .FBIG => return error.FileTooBig,
2503 .OVERFLOW => return error.FileTooBig,
2504 .ISDIR => return error.IsDir,
2505 .LOOP => return error.SymLinkLoop,
2506 .MFILE => return error.ProcessFdQuotaExceeded,
2507 .NAMETOOLONG => return error.NameTooLong,
2508 .NFILE => return error.SystemFdQuotaExceeded,
2509 .NODEV => return error.NoDevice,
2510 .NOENT => return error.FileNotFound,
2511 .NOMEM => return error.SystemResources,
2512 .NOSPC => return error.NoSpaceLeft,
2513 .NOTDIR => return error.NotDir,
2514 .PERM => return error.PermissionDenied,
2515 .EXIST => return error.PathAlreadyExists,
2516 .BUSY => return error.DeviceBusy,
2517 .OPNOTSUPP => return error.FileLocksUnsupported,
2518 .AGAIN => return error.WouldBlock,
2519 .TXTBSY => return error.FileBusy,
2520 .ROFS => return error.ReadOnlyFileSystem,
2521 .NXIO => return error.NoDevice,
2522 .ILSEQ => return error.BadPathName,
2523 else => |err| return unexpectedErrno(err),
2524 }
2525 };
2526 errdefer closeFd(fd);
2527
2528 return .{
2529 .handle = fd,
2530 .flags = .{ .nonblocking = os_flags.NONBLOCK },
2531 };
2532}
2533
2534fn dirCreateFileAtomic(
2535 userdata: ?*anyopaque,
2536 dir: Dir,
2537 dest_path: []const u8,
2538 options: Dir.CreateFileAtomicOptions,
2539) Dir.CreateFileAtomicError!File.Atomic {
2540 const ev: *Evented = @ptrCast(@alignCast(userdata));
2541 if (Dir.path.dirname(dest_path)) |dirname| {
2542 const new_dir = if (options.make_path)
2543 dirCreateDirPathOpen(ev, dir, dirname, .default_dir, .{}) catch |err| switch (err) {
2544 // None of these make sense in this context.
2545 error.IsDir,
2546 error.Streaming,
2547 error.DiskQuota,
2548 error.PathAlreadyExists,
2549 error.LinkQuotaExceeded,
2550 error.PipeBusy,
2551 error.FileTooBig,
2552 error.FileLocksUnsupported,
2553 error.DeviceBusy,
2554 => return error.Unexpected,
2555
2556 else => |e| return e,
2557 }
2558 else
2559 try dirOpenDir(ev, dir, dirname, .{});
2560 return ev.atomicFileInit(Dir.path.basename(dest_path), options.permissions, new_dir, true);
2561 }
2562 return ev.atomicFileInit(dest_path, options.permissions, dir, false);
2563}
2564
2565fn atomicFileInit(
2566 ev: *Evented,
2567 dest_basename: []const u8,
2568 permissions: File.Permissions,
2569 dir: Dir,
2570 close_dir_on_deinit: bool,
2571) Dir.CreateFileAtomicError!File.Atomic {
2572 while (true) {
2573 var random_integer: u64 = undefined;
2574 random(ev, @ptrCast(&random_integer));
2575 const tmp_sub_path = std.fmt.hex(random_integer);
2576 const file = dirCreateFile(ev, dir, &tmp_sub_path, .{
2577 .permissions = permissions,
2578 .exclusive = true,
2579 }) catch |err| switch (err) {
2580 error.PathAlreadyExists => continue,
2581 error.DeviceBusy => continue,
2582 error.FileBusy => continue,
2583
2584 error.IsDir => return error.Unexpected, // No path components.
2585 error.FileTooBig => return error.Unexpected, // Creating, not opening.
2586 error.FileLocksUnsupported => return error.Unexpected, // Not asking for locks.
2587 error.PipeBusy => return error.Unexpected, // Not opening a pipe.
2588
2589 else => |e| return e,
2590 };
2591 return .{
2592 .file = file,
2593 .file_basename_hex = random_integer,
2594 .dest_sub_path = dest_basename,
2595 .file_open = true,
2596 .file_exists = true,
2597 .close_dir_on_deinit = close_dir_on_deinit,
2598 .dir = dir,
2599 };
2600 }
2601}
2602
2603fn dirOpenFile(
2604 userdata: ?*anyopaque,
2605 dir: Dir,
2606 sub_path: []const u8,
2607 flags: Dir.OpenFileOptions,
2608) File.OpenError!File {
2609 const ev: *Evented = @ptrCast(@alignCast(userdata));
2610
2611 var path_buffer: [c.PATH_MAX]u8 = undefined;
2612 const sub_path_posix = try pathToPosix(sub_path, &path_buffer);
2613
2614 const os_flags: c.O = .{
2615 .ACCMODE = switch (flags.mode) {
2616 .read_only => .RDONLY,
2617 .write_only => .WRONLY,
2618 .read_write => .RDWR,
2619 },
2620 .NONBLOCK = flags.lock == .none or flags.lock_nonblocking,
2621 .SHLOCK = flags.lock == .shared,
2622 .EXLOCK = flags.lock == .exclusive,
2623 .NOFOLLOW = !flags.follow_symlinks,
2624 .NOCTTY = !flags.allow_ctty,
2625 .CLOEXEC = true,
2626 };
2627
2628 const fd: c.fd_t = while (true) {
2629 const rc = c.openat(dir.handle, sub_path_posix, os_flags);
2630 switch (c.errno(rc)) {
2631 .SUCCESS => break @intCast(rc),
2632 .INTR => {},
2633 .FAULT => |err| return errnoBug(err),
2634 .INVAL => return error.BadPathName,
2635 .BADF => |err| return errnoBug(err), // File descriptor used after closed.
2636 .ACCES => return error.AccessDenied,
2637 .FBIG => return error.FileTooBig,
2638 .OVERFLOW => return error.FileTooBig,
2639 .ISDIR => return error.IsDir,
2640 .LOOP => return error.SymLinkLoop,
2641 .MFILE => return error.ProcessFdQuotaExceeded,
2642 .NAMETOOLONG => return error.NameTooLong,
2643 .NFILE => return error.SystemFdQuotaExceeded,
2644 .NODEV => return error.NoDevice,
2645 .NOENT => return error.FileNotFound,
2646 .NOMEM => return error.SystemResources,
2647 .NOSPC => return error.NoSpaceLeft,
2648 .NOTDIR => return error.NotDir,
2649 .PERM => return error.PermissionDenied,
2650 .EXIST => return error.PathAlreadyExists,
2651 .BUSY => return error.DeviceBusy,
2652 .OPNOTSUPP => return error.FileLocksUnsupported,
2653 .AGAIN => return error.WouldBlock,
2654 .TXTBSY => return error.FileBusy,
2655 .NXIO => return error.NoDevice,
2656 .ROFS => return error.ReadOnlyFileSystem,
2657 .ILSEQ => return error.BadPathName,
2658 else => |err| return unexpectedErrno(err),
2659 }
2660 };
2661 errdefer closeFd(fd);
2662
2663 if (!flags.allow_directory) {
2664 const is_dir = is_dir: {
2665 const stat = fileStat(ev, .{
2666 .handle = fd,
2667 .flags = .{ .nonblocking = false },
2668 }) catch |err| switch (err) {
2669 // The directory-ness is either unknown or unknowable
2670 error.Streaming => break :is_dir false,
2671 else => |e| return e,
2672 };
2673 break :is_dir stat.kind == .directory;
2674 };
2675 if (is_dir) return error.IsDir;
2676 }
2677
2678 return .{
2679 .handle = fd,
2680 .flags = .{ .nonblocking = os_flags.NONBLOCK },
2681 };
2682}
2683
2684fn dirClose(userdata: ?*anyopaque, dirs: []const Dir) void {
2685 const ev: *Evented = @ptrCast(@alignCast(userdata));
2686 _ = ev;
2687 for (dirs) |dir| closeFd(dir.handle);
2688}
2689
2690fn dirRead(userdata: ?*anyopaque, dr: *Dir.Reader, buffer: []Dir.Entry) Dir.Reader.Error!usize {
2691 const ev: *Evented = @ptrCast(@alignCast(userdata));
2692 const Header = extern struct {
2693 seek: i64,
2694 };
2695 const header: *Header = @ptrCast(dr.buffer.ptr);
2696 const header_end: usize = @sizeOf(Header);
2697 if (dr.index < header_end) {
2698 // Initialize header.
2699 dr.index = header_end;
2700 dr.end = header_end;
2701 header.* = .{ .seek = 0 };
2702 }
2703 var buffer_index: usize = 0;
2704 while (buffer.len - buffer_index != 0) {
2705 if (dr.end - dr.index == 0) {
2706 // Refill the buffer, unless we've already created references to
2707 // buffered data.
2708 if (buffer_index != 0) break;
2709 if (dr.state == .reset) {
2710 ev.lseek(dr.dir.handle, 0, c.SEEK.SET) catch |err| switch (err) {
2711 error.Unseekable => return error.Unexpected,
2712 else => |e| return e,
2713 };
2714 dr.state = .reading;
2715 }
2716 const dents_buffer = dr.buffer[header_end..];
2717 const n: usize = while (true) {
2718 const rc = c.getdirentries(dr.dir.handle, dents_buffer.ptr, dents_buffer.len, &header.seek);
2719 switch (c.errno(rc)) {
2720 .SUCCESS => break @intCast(rc),
2721 .INTR => {},
2722 .BADF => |err| return errnoBug(err), // Dir is invalid or was opened without iteration ability.
2723 .FAULT => |err| return errnoBug(err),
2724 .NOTDIR => |err| return errnoBug(err),
2725 .INVAL => |err| return errnoBug(err),
2726 else => |err| return unexpectedErrno(err),
2727 }
2728 };
2729 if (n == 0) {
2730 dr.state = .finished;
2731 return 0;
2732 }
2733 dr.index = header_end;
2734 dr.end = header_end + n;
2735 }
2736 const darwin_entry = @as(*align(1) c.dirent, @ptrCast(&dr.buffer[dr.index]));
2737 const next_index = dr.index + darwin_entry.reclen;
2738 dr.index = next_index;
2739
2740 const name = @as([*]u8, @ptrCast(&darwin_entry.name))[0..darwin_entry.namlen];
2741 if (std.mem.eql(u8, name, ".") or std.mem.eql(u8, name, "..") or (darwin_entry.ino == 0))
2742 continue;
2743
2744 const entry_kind: File.Kind = switch (darwin_entry.type) {
2745 c.DT.BLK => .block_device,
2746 c.DT.CHR => .character_device,
2747 c.DT.DIR => .directory,
2748 c.DT.FIFO => .named_pipe,
2749 c.DT.LNK => .sym_link,
2750 c.DT.REG => .file,
2751 c.DT.SOCK => .unix_domain_socket,
2752 c.DT.WHT => .whiteout,
2753 else => .unknown,
2754 };
2755 buffer[buffer_index] = .{
2756 .name = name,
2757 .kind = entry_kind,
2758 .inode = darwin_entry.ino,
2759 };
2760 buffer_index += 1;
2761 }
2762 return buffer_index;
2763}
2764
2765fn dirRealPath(userdata: ?*anyopaque, dir: Dir, out_buffer: []u8) Dir.RealPathError!usize {
2766 const ev: *Evented = @ptrCast(@alignCast(userdata));
2767 return ev.realPath(dir.handle, out_buffer);
2768}
2769
2770fn realPath(ev: *Evented, fd: c.fd_t, out_buffer: []u8) File.RealPathError!usize {
2771 _ = ev;
2772 var buffer: [c.PATH_MAX]u8 = undefined;
2773 @memset(&buffer, 0);
2774 while (true) {
2775 switch (c.errno(c.fcntl(fd, c.F.GETPATH, &buffer))) {
2776 .SUCCESS => break,
2777 .INTR => {},
2778 .ACCES => return error.AccessDenied,
2779 .BADF => return error.FileNotFound,
2780 .NOENT => return error.FileNotFound,
2781 .NOMEM => return error.SystemResources,
2782 .NOSPC => return error.NameTooLong,
2783 .RANGE => return error.NameTooLong,
2784 else => |err| return unexpectedErrno(err),
2785 }
2786 }
2787 const n = std.mem.findScalar(u8, &buffer, 0) orelse buffer.len;
2788 if (n > out_buffer.len) return error.NameTooLong;
2789 @memcpy(out_buffer[0..n], buffer[0..n]);
2790 return n;
2791}
2792
2793fn dirRealPathFile(
2794 userdata: ?*anyopaque,
2795 dir: Dir,
2796 sub_path: []const u8,
2797 out_buffer: []u8,
2798) Dir.RealPathFileError!usize {
2799 const ev: *Evented = @ptrCast(@alignCast(userdata));
2800
2801 var path_buffer: [c.PATH_MAX]u8 = undefined;
2802 const sub_path_posix = try pathToPosix(sub_path, &path_buffer);
2803
2804 if (dir.handle == c.AT.FDCWD) {
2805 if (out_buffer.len < c.PATH_MAX) return error.NameTooLong;
2806 while (true) {
2807 if (c.realpath(sub_path_posix, out_buffer.ptr)) |redundant_pointer| {
2808 assert(redundant_pointer == out_buffer.ptr);
2809 return std.mem.findScalar(u8, out_buffer, 0) orelse out_buffer.len;
2810 }
2811 const err: c.E = @fromBackingInt(@intCast(c._errno().*));
2812 switch (err) {
2813 .INTR => {},
2814 .INVAL => return errnoBug(err),
2815 .BADF => return errnoBug(err),
2816 .FAULT => return errnoBug(err),
2817 .ACCES => return error.AccessDenied,
2818 .NOENT => return error.FileNotFound,
2819 .OPNOTSUPP => return error.OperationUnsupported,
2820 .NOTDIR => return error.NotDir,
2821 .NAMETOOLONG => return error.NameTooLong,
2822 .LOOP => return error.SymLinkLoop,
2823 .IO => return error.InputOutput,
2824 else => return unexpectedErrno(err),
2825 }
2826 }
2827 }
2828
2829 const os_flags: c.O = .{
2830 .NONBLOCK = true,
2831 .CLOEXEC = true,
2832 };
2833
2834 const fd: c.fd_t = while (true) {
2835 const rc = c.openat(dir.handle, sub_path_posix, os_flags);
2836 switch (c.errno(rc)) {
2837 .SUCCESS => break @intCast(rc),
2838 .INTR => {},
2839 .FAULT => |err| return errnoBug(err),
2840 .INVAL => return error.BadPathName,
2841 .BADF => |err| return errnoBug(err), // File descriptor used after closed.
2842 .ACCES => return error.AccessDenied,
2843 .FBIG => return error.FileTooBig,
2844 .OVERFLOW => return error.FileTooBig,
2845 .ISDIR => return error.IsDir,
2846 .LOOP => return error.SymLinkLoop,
2847 .MFILE => return error.ProcessFdQuotaExceeded,
2848 .NAMETOOLONG => return error.NameTooLong,
2849 .NFILE => return error.SystemFdQuotaExceeded,
2850 .NODEV => return error.NoDevice,
2851 .NOENT => return error.FileNotFound,
2852 .NOMEM => return error.SystemResources,
2853 .NOSPC => return error.NoSpaceLeft,
2854 .NOTDIR => return error.NotDir,
2855 .PERM => return error.PermissionDenied,
2856 .EXIST => return error.PathAlreadyExists,
2857 .BUSY => return error.DeviceBusy,
2858 .NXIO => return error.NoDevice,
2859 .ILSEQ => return error.BadPathName,
2860 else => |err| return unexpectedErrno(err),
2861 }
2862 };
2863 defer closeFd(fd);
2864 return ev.realPath(fd, out_buffer);
2865}
2866
2867fn dirDeleteFile(userdata: ?*anyopaque, dir: Dir, sub_path: []const u8) Dir.DeleteFileError!void {
2868 const ev: *Evented = @ptrCast(@alignCast(userdata));
2869 _ = ev;
2870
2871 var path_buffer: [c.PATH_MAX]u8 = undefined;
2872 const sub_path_posix = try pathToPosix(sub_path, &path_buffer);
2873
2874 while (true) switch (c.errno(c.unlinkat(dir.handle, sub_path_posix, 0))) {
2875 .SUCCESS => return,
2876 .INTR => {},
2877 // Some systems return permission errors when trying to delete a
2878 // directory, so we need to handle that case specifically and
2879 // translate the error.
2880 .PERM => {
2881 // Don't follow symlinks to match unlinkat (which acts on symlinks rather than follows them).
2882 var st = std.mem.zeroes(c.Stat);
2883 while (true) switch (c.errno(c.fstatat(
2884 dir.handle,
2885 sub_path_posix,
2886 &st,
2887 c.AT.SYMLINK_NOFOLLOW,
2888 ))) {
2889 .SUCCESS => break,
2890 .INTR => {},
2891 else => return error.PermissionDenied,
2892 };
2893 if (st.mode & c.S.IFMT == c.S.IFDIR) return error.IsDir else return error.PermissionDenied;
2894 },
2895 .ACCES => return error.AccessDenied,
2896 .BUSY => return error.FileBusy,
2897 .FAULT => |err| return errnoBug(err),
2898 .IO => return error.FileSystem,
2899 .ISDIR => return error.IsDir,
2900 .LOOP => return error.SymLinkLoop,
2901 .NAMETOOLONG => return error.NameTooLong,
2902 .NOENT => return error.FileNotFound,
2903 .NOTDIR => return error.NotDir,
2904 .NOMEM => return error.SystemResources,
2905 .ROFS => return error.ReadOnlyFileSystem,
2906 .EXIST => |err| return errnoBug(err),
2907 .NOTEMPTY => |err| return errnoBug(err), // Not passing AT.REMOVEDIR
2908 .ILSEQ => return error.BadPathName,
2909 .INVAL => |err| return errnoBug(err), // invalid flags, or pathname has . as last component
2910 .BADF => |err| return errnoBug(err), // File descriptor used after closed.
2911 else => |err| return unexpectedErrno(err),
2912 };
2913}
2914
2915fn dirDeleteDir(userdata: ?*anyopaque, dir: Dir, sub_path: []const u8) Dir.DeleteDirError!void {
2916 const ev: *Evented = @ptrCast(@alignCast(userdata));
2917 _ = ev;
2918
2919 var path_buffer: [c.PATH_MAX]u8 = undefined;
2920 const sub_path_posix = try pathToPosix(sub_path, &path_buffer);
2921
2922 while (true) switch (c.errno(c.unlinkat(dir.handle, sub_path_posix, c.AT.REMOVEDIR))) {
2923 .SUCCESS => return,
2924 .INTR => {},
2925 .ACCES => return error.AccessDenied,
2926 .PERM => return error.PermissionDenied,
2927 .BUSY => return error.FileBusy,
2928 .FAULT => |err| return errnoBug(err),
2929 .IO => return error.FileSystem,
2930 .ISDIR => |err| return errnoBug(err),
2931 .LOOP => return error.SymLinkLoop,
2932 .NAMETOOLONG => return error.NameTooLong,
2933 .NOENT => return error.FileNotFound,
2934 .NOTDIR => return error.NotDir,
2935 .NOMEM => return error.SystemResources,
2936 .ROFS => return error.ReadOnlyFileSystem,
2937 .EXIST => |err| return errnoBug(err),
2938 .NOTEMPTY => return error.DirNotEmpty,
2939 .ILSEQ => return error.BadPathName,
2940 .INVAL => |err| return errnoBug(err), // invalid flags, or pathname has . as last component
2941 .BADF => |err| return errnoBug(err), // File descriptor used after closed.
2942 else => |err| return unexpectedErrno(err),
2943 };
2944}
2945
2946fn dirRename(
2947 userdata: ?*anyopaque,
2948 old_dir: Dir,
2949 old_sub_path: []const u8,
2950 new_dir: Dir,
2951 new_sub_path: []const u8,
2952) Dir.RenameError!void {
2953 const ev: *Evented = @ptrCast(@alignCast(userdata));
2954 _ = ev;
2955
2956 var old_path_buffer: [c.PATH_MAX]u8 = undefined;
2957 var new_path_buffer: [c.PATH_MAX]u8 = undefined;
2958
2959 const old_sub_path_posix = try pathToPosix(old_sub_path, &old_path_buffer);
2960 const new_sub_path_posix = try pathToPosix(new_sub_path, &new_path_buffer);
2961
2962 while (true) switch (c.errno(c.renameat(old_dir.handle, old_sub_path_posix, new_dir.handle, new_sub_path_posix))) {
2963 .SUCCESS => return,
2964 .INTR => {},
2965 .ACCES => return error.AccessDenied,
2966 .PERM => return error.PermissionDenied,
2967 .BUSY => return error.FileBusy,
2968 .DQUOT => return error.DiskQuota,
2969 .ISDIR => return error.IsDir,
2970 .IO => return error.HardwareFailure,
2971 .LOOP => return error.SymLinkLoop,
2972 .MLINK => return error.LinkQuotaExceeded,
2973 .NAMETOOLONG => return error.NameTooLong,
2974 .NOENT => return error.FileNotFound,
2975 .NOTDIR => return error.NotDir,
2976 .NOMEM => return error.SystemResources,
2977 .NOSPC => return error.NoSpaceLeft,
2978 .EXIST => return error.DirNotEmpty,
2979 .NOTEMPTY => return error.DirNotEmpty,
2980 .ROFS => return error.ReadOnlyFileSystem,
2981 .XDEV => return error.CrossDevice,
2982 .ILSEQ => return error.BadPathName,
2983 .FAULT => |err| return errnoBug(err),
2984 .INVAL => |err| return errnoBug(err),
2985 else => |err| return unexpectedErrno(err),
2986 };
2987}
2988
2989fn dirRenamePreserve(
2990 userdata: ?*anyopaque,
2991 old_dir: Dir,
2992 old_sub_path: []const u8,
2993 new_dir: Dir,
2994 new_sub_path: []const u8,
2995) Dir.RenamePreserveError!void {
2996 const ev: *Evented = @ptrCast(@alignCast(userdata));
2997 // Make a hard link then delete the original.
2998 try dirHardLink(ev, old_dir, old_sub_path, new_dir, new_sub_path, .{ .follow_symlinks = false });
2999 const prev = swapCancelProtection(ev, .blocked);
3000 defer _ = swapCancelProtection(ev, prev);
3001 dirDeleteFile(ev, old_dir, old_sub_path) catch {};
3002}
3003
3004fn dirSymLink(
3005 userdata: ?*anyopaque,
3006 dir: Dir,
3007 target_path: []const u8,
3008 sym_link_path: []const u8,
3009 flags: Dir.SymLinkFlags,
3010) Dir.SymLinkError!void {
3011 const ev: *Evented = @ptrCast(@alignCast(userdata));
3012 _ = ev;
3013 _ = flags;
3014
3015 var target_path_buffer: [c.PATH_MAX]u8 = undefined;
3016 var sym_link_path_buffer: [c.PATH_MAX]u8 = undefined;
3017
3018 const target_path_posix = try pathToPosix(target_path, &target_path_buffer);
3019 const sym_link_path_posix = try pathToPosix(sym_link_path, &sym_link_path_buffer);
3020
3021 while (true) switch (c.errno(c.symlinkat(target_path_posix, dir.handle, sym_link_path_posix))) {
3022 .SUCCESS => return,
3023 .INTR => {},
3024 .FAULT => |err| return errnoBug(err),
3025 .INVAL => |err| return errnoBug(err),
3026 .ACCES => return error.AccessDenied,
3027 .PERM => return error.PermissionDenied,
3028 .DQUOT => return error.DiskQuota,
3029 .EXIST => return error.PathAlreadyExists,
3030 .IO => return error.FileSystem,
3031 .LOOP => return error.SymLinkLoop,
3032 .NAMETOOLONG => return error.NameTooLong,
3033 .NOENT => return error.FileNotFound,
3034 .NOTDIR => return error.NotDir,
3035 .NOMEM => return error.SystemResources,
3036 .NOSPC => return error.NoSpaceLeft,
3037 .ROFS => return error.ReadOnlyFileSystem,
3038 .ILSEQ => return error.BadPathName,
3039 else => |err| return unexpectedErrno(err),
3040 };
3041}
3042
3043fn dirReadLink(
3044 userdata: ?*anyopaque,
3045 dir: Dir,
3046 sub_path: []const u8,
3047 buffer: []u8,
3048) Dir.ReadLinkError!usize {
3049 const ev: *Evented = @ptrCast(@alignCast(userdata));
3050 _ = ev;
3051 var sub_path_buffer: [c.PATH_MAX]u8 = undefined;
3052 const sub_path_posix = try pathToPosix(sub_path, &sub_path_buffer);
3053 while (true) {
3054 const rc = c.readlinkat(dir.handle, sub_path_posix, buffer.ptr, buffer.len);
3055 switch (c.errno(rc)) {
3056 .SUCCESS => return @intCast(rc),
3057 .INTR => {},
3058 .ACCES => return error.AccessDenied,
3059 .FAULT => |err| return errnoBug(err),
3060 .INVAL => return error.NotLink,
3061 .IO => return error.FileSystem,
3062 .LOOP => return error.SymLinkLoop,
3063 .NAMETOOLONG => return error.NameTooLong,
3064 .NOENT => return error.FileNotFound,
3065 .NOMEM => return error.SystemResources,
3066 .NOTDIR => return error.NotDir,
3067 .ILSEQ => return error.BadPathName,
3068 else => |err| return unexpectedErrno(err),
3069 }
3070 }
3071}
3072
3073fn dirSetOwner(
3074 userdata: ?*anyopaque,
3075 dir: Dir,
3076 owner: ?File.Uid,
3077 group: ?File.Gid,
3078) Dir.SetOwnerError!void {
3079 const ev: *Evented = @ptrCast(@alignCast(userdata));
3080 _ = ev;
3081 return fchown(dir.handle, owner, group);
3082}
3083
3084fn fchown(fd: c.fd_t, owner: ?File.Uid, group: ?File.Gid) File.SetOwnerError!void {
3085 const uid = owner orelse std.math.maxInt(c.uid_t);
3086 const gid = group orelse std.math.maxInt(c.gid_t);
3087 while (true) switch (c.errno(c.fchown(fd, uid, gid))) {
3088 .SUCCESS => return,
3089 .INTR => {},
3090 .BADF => |err| return errnoBug(err), // likely fd refers to directory opened without `Dir.OpenOptions.iterate`
3091 .FAULT => |err| return errnoBug(err),
3092 .INVAL => |err| return errnoBug(err),
3093 .ACCES => return error.AccessDenied,
3094 .IO => return error.InputOutput,
3095 .LOOP => return error.SymLinkLoop,
3096 .NOENT => return error.FileNotFound,
3097 .NOMEM => return error.SystemResources,
3098 .NOTDIR => return error.FileNotFound,
3099 .PERM => return error.PermissionDenied,
3100 .ROFS => return error.ReadOnlyFileSystem,
3101 else => |err| return unexpectedErrno(err),
3102 };
3103}
3104
3105fn dirSetFileOwner(
3106 userdata: ?*anyopaque,
3107 dir: Dir,
3108 sub_path: []const u8,
3109 owner: ?File.Uid,
3110 group: ?File.Gid,
3111 options: Dir.SetFileOwnerOptions,
3112) Dir.SetFileOwnerError!void {
3113 const ev: *Evented = @ptrCast(@alignCast(userdata));
3114 var path_buffer: [c.PATH_MAX]u8 = undefined;
3115 const sub_path_posix = try pathToPosix(sub_path, &path_buffer);
3116 _ = ev;
3117 while (true) switch (c.errno(c.fchownat(
3118 dir.handle,
3119 sub_path_posix,
3120 owner orelse std.math.maxInt(c.uid_t),
3121 group orelse std.math.maxInt(c.gid_t),
3122 if (options.follow_symlinks) 0 else c.AT.SYMLINK_NOFOLLOW,
3123 ))) {
3124 .SUCCESS => return,
3125 .INTR => continue,
3126 .BADF => |err| return errnoBug(err), // likely fd refers to directory opened without `Dir.OpenOptions.iterate`
3127 .FAULT => |err| return errnoBug(err),
3128 .INVAL => |err| return errnoBug(err),
3129 .ACCES => return error.AccessDenied,
3130 .IO => return error.InputOutput,
3131 .LOOP => return error.SymLinkLoop,
3132 .NOENT => return error.FileNotFound,
3133 .NOMEM => return error.SystemResources,
3134 .NOTDIR => return error.FileNotFound,
3135 .PERM => return error.PermissionDenied,
3136 .ROFS => return error.ReadOnlyFileSystem,
3137 else => |err| return unexpectedErrno(err),
3138 };
3139}
3140
3141fn dirSetPermissions(
3142 userdata: ?*anyopaque,
3143 dir: Dir,
3144 permissions: Dir.Permissions,
3145) Dir.SetPermissionsError!void {
3146 const ev: *Evented = @ptrCast(@alignCast(userdata));
3147 return ev.fchmod(dir.handle, permissions.toMode());
3148}
3149
3150fn dirSetFilePermissions(
3151 userdata: ?*anyopaque,
3152 dir: Dir,
3153 sub_path: []const u8,
3154 permissions: Dir.Permissions,
3155 options: Dir.SetFilePermissionsOptions,
3156) Dir.SetFilePermissionsError!void {
3157 const ev: *Evented = @ptrCast(@alignCast(userdata));
3158 _ = ev;
3159
3160 var path_buffer: [c.PATH_MAX]u8 = undefined;
3161 const sub_path_posix = try pathToPosix(sub_path, &path_buffer);
3162
3163 const mode = permissions.toMode();
3164 const flags: u32 = if (options.follow_symlinks) 0 else c.AT.SYMLINK_NOFOLLOW;
3165
3166 while (true) switch (c.errno(c.fchmodat(dir.handle, sub_path_posix, mode, flags))) {
3167 .SUCCESS => return,
3168 .INTR => {},
3169 .BADF => |err| return errnoBug(err),
3170 .FAULT => |err| return errnoBug(err),
3171 .INVAL => |err| return errnoBug(err),
3172 .ACCES => return error.AccessDenied,
3173 .IO => return error.InputOutput,
3174 .LOOP => return error.SymLinkLoop,
3175 .MFILE => return error.ProcessFdQuotaExceeded,
3176 .NAMETOOLONG => return error.NameTooLong,
3177 .NFILE => return error.SystemFdQuotaExceeded,
3178 .NOENT => return error.FileNotFound,
3179 .NOTDIR => return error.FileNotFound,
3180 .NOMEM => return error.SystemResources,
3181 .OPNOTSUPP => return error.OperationUnsupported,
3182 .PERM => return error.PermissionDenied,
3183 .ROFS => return error.ReadOnlyFileSystem,
3184 else => |err| return unexpectedErrno(err),
3185 };
3186}
3187
3188fn dirSetTimestamps(
3189 userdata: ?*anyopaque,
3190 dir: Dir,
3191 sub_path: []const u8,
3192 options: Dir.SetTimestampsOptions,
3193) Dir.SetTimestampsError!void {
3194 const ev: *Evented = @ptrCast(@alignCast(userdata));
3195 _ = ev;
3196
3197 var times_buffer: [2]c.timespec = undefined;
3198 const times = if (options.modify_timestamp == .now and options.access_timestamp == .now) null else p: {
3199 times_buffer = .{
3200 setTimestampToPosix(options.access_timestamp),
3201 setTimestampToPosix(options.modify_timestamp),
3202 };
3203 break :p &times_buffer;
3204 };
3205
3206 const flags: u32 = if (options.follow_symlinks) 0 else c.AT.SYMLINK_NOFOLLOW;
3207
3208 var path_buffer: [c.PATH_MAX]u8 = undefined;
3209 const sub_path_posix = try pathToPosix(sub_path, &path_buffer);
3210
3211 while (true) switch (c.errno(c.utimensat(dir.handle, sub_path_posix, times, flags))) {
3212 .SUCCESS => return,
3213 .INTR => {},
3214 .BADF => |err| return errnoBug(err), // always a race condition
3215 .FAULT => |err| return errnoBug(err),
3216 .INVAL => |err| return errnoBug(err),
3217 .ACCES => return error.AccessDenied,
3218 .PERM => return error.PermissionDenied,
3219 .ROFS => return error.ReadOnlyFileSystem,
3220 else => |err| return unexpectedErrno(err),
3221 };
3222}
3223
3224fn dirHardLink(
3225 userdata: ?*anyopaque,
3226 old_dir: Dir,
3227 old_sub_path: []const u8,
3228 new_dir: Dir,
3229 new_sub_path: []const u8,
3230 options: Dir.HardLinkOptions,
3231) Dir.HardLinkError!void {
3232 const ev: *Evented = @ptrCast(@alignCast(userdata));
3233 _ = ev;
3234
3235 var old_path_buffer: [c.PATH_MAX]u8 = undefined;
3236 var new_path_buffer: [c.PATH_MAX]u8 = undefined;
3237
3238 const old_sub_path_posix = try pathToPosix(old_sub_path, &old_path_buffer);
3239 const new_sub_path_posix = try pathToPosix(new_sub_path, &new_path_buffer);
3240
3241 const flags: u32 = if (options.follow_symlinks) c.AT.SYMLINK_FOLLOW else 0;
3242 return linkat(old_dir.handle, old_sub_path_posix, new_dir.handle, new_sub_path_posix, flags);
3243}
3244
3245fn fileStat(userdata: ?*anyopaque, file: File) File.StatError!File.Stat {
3246 const ev: *Evented = @ptrCast(@alignCast(userdata));
3247 _ = ev;
3248 while (true) {
3249 var stat = std.mem.zeroes(c.Stat);
3250 switch (c.errno(c.fstat(file.handle, &stat))) {
3251 .SUCCESS => return statFromPosix(&stat),
3252 .INTR => {},
3253 .INVAL => |err| return errnoBug(err),
3254 .BADF => |err| return errnoBug(err), // File descriptor used after closed.
3255 .NOMEM => return error.SystemResources,
3256 .ACCES => return error.AccessDenied,
3257 else => |err| return unexpectedErrno(err),
3258 }
3259 }
3260}
3261
3262fn fileLength(userdata: ?*anyopaque, file: File) File.LengthError!u64 {
3263 const ev: *Evented = @ptrCast(@alignCast(userdata));
3264 const stat = try fileStat(ev, file);
3265 return stat.size;
3266}
3267
3268fn fileClose(userdata: ?*anyopaque, files: []const File) void {
3269 const ev: *Evented = @ptrCast(@alignCast(userdata));
3270 _ = ev;
3271 for (files) |file| closeFd(file.handle);
3272}
3273
3274fn fileWritePositional(
3275 userdata: ?*anyopaque,
3276 file: File,
3277 header: []const u8,
3278 data: []const []const u8,
3279 splat: usize,
3280 offset: u64,
3281) File.WritePositionalError!usize {
3282 const ev: *Evented = @ptrCast(@alignCast(userdata));
3283 _ = ev;
3284 var iovecs: [max_iovecs_len]iovec_const = undefined;
3285 var iovlen: iovlen_t = 0;
3286 var remaining: Io.Limit = .unlimited;
3287 addBuf(true, &iovecs, &iovlen, &remaining, header);
3288 for (data[0 .. data.len - 1]) |bytes| addBuf(true, &iovecs, &iovlen, &remaining, bytes);
3289 const pattern = data[data.len - 1];
3290 var backup_buffer: [splat_buffer_size]u8 = undefined;
3291 if (iovecs.len - iovlen != 0 and remaining != .nothing) switch (splat) {
3292 0 => {},
3293 1 => addBuf(true, &iovecs, &iovlen, &remaining, pattern),
3294 else => switch (pattern.len) {
3295 0 => {},
3296 1 => {
3297 const splat_buffer = &backup_buffer;
3298 const memset_len = @min(splat_buffer.len, splat);
3299 const buf = splat_buffer[0..memset_len];
3300 @memset(buf, pattern[0]);
3301 addBuf(true, &iovecs, &iovlen, &remaining, buf);
3302 var remaining_splat = splat - buf.len;
3303 while (remaining_splat > splat_buffer.len and iovecs.len - iovlen != 0 and remaining != .nothing) {
3304 assert(buf.len == splat_buffer.len);
3305 addBuf(true, &iovecs, &iovlen, &remaining, splat_buffer);
3306 remaining_splat -= splat_buffer.len;
3307 }
3308 addBuf(true, &iovecs, &iovlen, &remaining, splat_buffer[0..@min(remaining_splat, splat_buffer.len)]);
3309 },
3310 else => for (0..@min(splat, iovecs.len - iovlen)) |_| {
3311 if (remaining == .nothing) break;
3312 addBuf(true, &iovecs, &iovlen, &remaining, pattern);
3313 },
3314 },
3315 };
3316 if (iovlen == 0) return 0;
3317 while (true) {
3318 const rc = c.pwritev(file.handle, &iovecs, iovlen, @bitCast(offset));
3319 switch (c.errno(rc)) {
3320 .SUCCESS => return @intCast(rc),
3321 .INTR => {},
3322 .INVAL => |err| return errnoBug(err),
3323 .FAULT => |err| return errnoBug(err),
3324 .DESTADDRREQ => |err| return errnoBug(err), // `connect` was never called.
3325 .CONNRESET => |err| return errnoBug(err), // Not a socket handle.
3326 .BADF => return error.NotOpenForWriting,
3327 .AGAIN => return error.WouldBlock,
3328 .DQUOT => return error.DiskQuota,
3329 .FBIG => return error.FileTooBig,
3330 .IO => return error.InputOutput,
3331 .NOSPC => return error.NoSpaceLeft,
3332 .PERM => return error.PermissionDenied,
3333 .PIPE => return error.BrokenPipe,
3334 .BUSY => return error.DeviceBusy,
3335 .TXTBSY => return error.FileBusy,
3336 .NXIO => return error.Unseekable,
3337 .SPIPE => return error.Unseekable,
3338 .OVERFLOW => return error.Unseekable,
3339 else => |err| return unexpectedErrno(err),
3340 }
3341 }
3342}
3343
3344fn fileWriteFileStreaming(
3345 userdata: ?*anyopaque,
3346 file: File,
3347 header: []const u8,
3348 file_reader: *File.Reader,
3349 limit: Io.Limit,
3350) File.Writer.WriteFileError!usize {
3351 const ev: *Evented = @ptrCast(@alignCast(userdata));
3352 const reader_buffered = file_reader.interface.buffered();
3353 if (reader_buffered.len >= @backingInt(limit)) {
3354 const n = try fileWriteStreaming(ev, file, header, &.{limit.slice(reader_buffered)}, 1);
3355 file_reader.interface.toss(n -| header.len);
3356 return n;
3357 }
3358 const file_limit = @backingInt(limit) - reader_buffered.len;
3359 const out_fd = file.handle;
3360 const in_fd = file_reader.file.handle;
3361
3362 if (file_reader.size) |size| {
3363 if (size - file_reader.pos == 0) {
3364 if (reader_buffered.len != 0) {
3365 const n = try fileWriteStreaming(ev, file, header, &.{limit.slice(reader_buffered)}, 1);
3366 file_reader.interface.toss(n -| header.len);
3367 return n;
3368 } else {
3369 return error.EndOfStream;
3370 }
3371 }
3372 }
3373
3374 if (@atomicLoad(UseSendfile, &ev.use_sendfile, .monotonic) == .disabled) return error.Unimplemented;
3375 const offset = std.math.cast(c.off_t, file_reader.pos) orelse return error.Unimplemented;
3376 var hdtr_data: c.sf_hdtr = undefined;
3377 var headers: [2]iovec_const = undefined;
3378 var headers_i: u8 = 0;
3379 if (header.len != 0) {
3380 headers[headers_i] = .{ .base = header.ptr, .len = header.len };
3381 headers_i += 1;
3382 }
3383 if (reader_buffered.len != 0) {
3384 headers[headers_i] = .{ .base = reader_buffered.ptr, .len = reader_buffered.len };
3385 headers_i += 1;
3386 }
3387 const hdtr: ?*c.sf_hdtr = if (headers_i == 0) null else b: {
3388 hdtr_data = .{
3389 .headers = &headers,
3390 .hdr_cnt = headers_i,
3391 .trailers = null,
3392 .trl_cnt = 0,
3393 };
3394 break :b &hdtr_data;
3395 };
3396 const max_count = std.math.maxInt(i32); // Avoid EINVAL.
3397 var len: c.off_t = @min(file_limit, max_count);
3398 const flags = 0;
3399 while (true) switch (c.errno(c.sendfile(in_fd, out_fd, offset, &len, hdtr, flags))) {
3400 .SUCCESS => break,
3401 .OPNOTSUPP, .NOTSOCK, .NOSYS => {
3402 // Give calling code chance to observe before trying
3403 // something else.
3404 @atomicStore(UseSendfile, &ev.use_sendfile, .disabled, .monotonic);
3405 return 0;
3406 },
3407 .INTR => if (len > 0) break,
3408 .AGAIN => {
3409 if (len == 0) return error.WouldBlock;
3410 break;
3411 },
3412 else => |e| {
3413 assert(error.Unexpected == switch (e) {
3414 .NOTCONN => return error.BrokenPipe,
3415 .IO => return error.InputOutput,
3416 .PIPE => return error.BrokenPipe,
3417 .BADF => |err| errnoBug(err),
3418 .FAULT => |err| errnoBug(err),
3419 .INVAL => |err| errnoBug(err),
3420 else => |err| unexpectedErrno(err),
3421 });
3422 // Give calling code chance to observe the error before trying
3423 // something else.
3424 @atomicStore(UseSendfile, &ev.use_sendfile, .disabled, .monotonic);
3425 return 0;
3426 },
3427 };
3428 if (len == 0) {
3429 file_reader.size = file_reader.pos;
3430 return error.EndOfStream;
3431 }
3432 const u_len: usize = @bitCast(len);
3433 file_reader.interface.toss(u_len -| header.len);
3434 return u_len;
3435}
3436
3437fn fileWriteFilePositional(
3438 userdata: ?*anyopaque,
3439 file: File,
3440 header: []const u8,
3441 file_reader: *File.Reader,
3442 limit: Io.Limit,
3443 offset: u64,
3444) File.WriteFilePositionalError!usize {
3445 const ev: *Evented = @ptrCast(@alignCast(userdata));
3446 const reader_buffered = file_reader.interface.buffered();
3447 if (reader_buffered.len >= @backingInt(limit)) {
3448 const n = try fileWritePositional(
3449 ev,
3450 file,
3451 header,
3452 &.{limit.slice(reader_buffered)},
3453 1,
3454 offset,
3455 );
3456 file_reader.interface.toss(n -| header.len);
3457 return n;
3458 }
3459 const out_fd = file.handle;
3460 const in_fd = file_reader.file.handle;
3461
3462 if (file_reader.size) |size| {
3463 if (size - file_reader.pos == 0) {
3464 if (reader_buffered.len != 0) {
3465 const n = try fileWritePositional(
3466 ev,
3467 file,
3468 header,
3469 &.{limit.slice(reader_buffered)},
3470 1,
3471 offset,
3472 );
3473 file_reader.interface.toss(n -| header.len);
3474 return n;
3475 } else {
3476 return error.EndOfStream;
3477 }
3478 }
3479 }
3480
3481 if (@atomicLoad(UseFcopyfile, &ev.use_fcopyfile, .monotonic) == .disabled)
3482 return error.Unimplemented;
3483 if (file_reader.pos != 0) return error.Unimplemented;
3484 if (offset != 0) return error.Unimplemented;
3485 if (limit != .unlimited) return error.Unimplemented;
3486 const size = file_reader.getSize() catch return error.Unimplemented;
3487 if (header.len != 0 or reader_buffered.len != 0) {
3488 const n = try fileWritePositional(
3489 ev,
3490 file,
3491 header,
3492 &.{limit.slice(reader_buffered)},
3493 1,
3494 offset,
3495 );
3496 file_reader.interface.toss(n -| header.len);
3497 return n;
3498 }
3499 while (true) {
3500 const rc = c.fcopyfile(in_fd, out_fd, null, .{ .DATA = true });
3501 switch (c.errno(rc)) {
3502 .SUCCESS => break,
3503 .INTR => {},
3504 .OPNOTSUPP => {
3505 // Give calling code chance to observe before trying
3506 // something else.
3507 @atomicStore(UseFcopyfile, &ev.use_fcopyfile, .disabled, .monotonic);
3508 return 0;
3509 },
3510 else => |e| {
3511 assert(error.Unexpected == switch (e) {
3512 .NOMEM => return error.SystemResources,
3513 .INVAL => |err| errnoBug(err),
3514 else => |err| unexpectedErrno(err),
3515 });
3516 return 0;
3517 },
3518 }
3519 }
3520 file_reader.pos = size;
3521 return size;
3522}
3523
3524fn fileReadPositional(
3525 userdata: ?*anyopaque,
3526 file: File,
3527 data: []const []u8,
3528 offset: u64,
3529) File.ReadPositionalError!usize {
3530 const ev: *Evented = @ptrCast(@alignCast(userdata));
3531 _ = ev;
3532 var iovecs: [max_iovecs_len]iovec = undefined;
3533 var iovlen: iovlen_t = 0;
3534 var remaining: Io.Limit = .unlimited;
3535 for (data) |buf| addBuf(false, &iovecs, &iovlen, &remaining, buf);
3536 if (iovlen == 0) return 0;
3537 while (true) {
3538 const rc = c.preadv(file.handle, &iovecs, iovlen, @bitCast(offset));
3539 switch (c.errno(rc)) {
3540 .SUCCESS => return @intCast(rc),
3541 .INTR => {},
3542 .NXIO => return error.Unseekable,
3543 .SPIPE => return error.Unseekable,
3544 .OVERFLOW => return error.Unseekable,
3545 .NOBUFS => return error.SystemResources,
3546 .NOMEM => return error.SystemResources,
3547 .AGAIN => return error.WouldBlock,
3548 .IO => return error.InputOutput,
3549 .ISDIR => return error.IsDir,
3550 .NOTCONN => |err| return errnoBug(err), // not a socket
3551 .CONNRESET => |err| return errnoBug(err), // not a socket
3552 .INVAL => |err| return errnoBug(err),
3553 .FAULT => |err| return errnoBug(err),
3554 .BADF => return error.NotOpenForReading,
3555 else => |err| return unexpectedErrno(err),
3556 }
3557 }
3558}
3559
3560fn fileSeekBy(userdata: ?*anyopaque, file: File, offset: i64) File.SeekError!void {
3561 const ev: *Evented = @ptrCast(@alignCast(userdata));
3562 return ev.lseek(file.handle, @bitCast(offset), c.SEEK.CUR);
3563}
3564
3565fn fileSeekTo(userdata: ?*anyopaque, file: File, offset: u64) File.SeekError!void {
3566 const ev: *Evented = @ptrCast(@alignCast(userdata));
3567 return ev.lseek(file.handle, offset, c.SEEK.SET);
3568}
3569
3570fn lseek(ev: *Evented, fd: c.fd_t, offset: u64, whence: i32) File.SeekError!void {
3571 _ = ev;
3572 while (true) switch (c.errno(c.lseek(fd, @bitCast(offset), whence))) {
3573 .SUCCESS => return,
3574 .INTR => {},
3575 .BADF => |err| return errnoBug(err), // File descriptor used after closed.
3576 .INVAL => return error.Unseekable,
3577 .OVERFLOW => return error.Unseekable,
3578 .SPIPE => return error.Unseekable,
3579 .NXIO => return error.Unseekable,
3580 else => |err| return unexpectedErrno(err),
3581 };
3582}
3583
3584fn fileSync(userdata: ?*anyopaque, file: File) File.SyncError!void {
3585 const ev: *Evented = @ptrCast(@alignCast(userdata));
3586 _ = ev;
3587 while (true) switch (c.errno(c.fsync(file.handle))) {
3588 .SUCCESS => return,
3589 .INTR => {},
3590 .BADF => |err| return errnoBug(err),
3591 .INVAL => |err| return errnoBug(err),
3592 .ROFS => |err| return errnoBug(err),
3593 .IO => return error.InputOutput,
3594 .NOSPC => return error.NoSpaceLeft,
3595 .DQUOT => return error.DiskQuota,
3596 else => |err| return unexpectedErrno(err),
3597 };
3598}
3599
3600fn fileIsTty(userdata: ?*anyopaque, file: File) Io.Cancelable!bool {
3601 const ev: *Evented = @ptrCast(@alignCast(userdata));
3602 _ = ev;
3603 while (true) {
3604 const rc = c.isatty(file.handle);
3605 switch (c.errno(rc - 1)) {
3606 .SUCCESS => return true,
3607 .INTR => {},
3608 else => return false,
3609 }
3610 }
3611}
3612
3613fn fileEnableAnsiEscapeCodes(userdata: ?*anyopaque, file: File) File.EnableAnsiEscapeCodesError!void {
3614 const ev: *Evented = @ptrCast(@alignCast(userdata));
3615 if (!try fileIsTty(ev, file)) return error.NotTerminalDevice;
3616}
3617
3618fn fileSetLength(userdata: ?*anyopaque, file: File, length: u64) File.SetLengthError!void {
3619 const ev: *Evented = @ptrCast(@alignCast(userdata));
3620 _ = ev;
3621
3622 const signed_len: i64 = @bitCast(length);
3623 if (signed_len < 0) return error.FileTooBig; // Avoid ambiguous EINVAL errors.
3624
3625 while (true) switch (c.errno(c.ftruncate(file.handle, signed_len))) {
3626 .SUCCESS => return,
3627 .INTR => {},
3628 .FBIG => return error.FileTooBig,
3629 .IO => return error.InputOutput,
3630 .PERM => return error.PermissionDenied,
3631 .TXTBSY => return error.FileBusy,
3632 .BADF => |err| return errnoBug(err), // Handle not open for writing.
3633 .INVAL => return error.NonResizable, // This is returned for /dev/null for example.
3634 else => |err| return unexpectedErrno(err),
3635 };
3636}
3637
3638fn fileSetOwner(
3639 userdata: ?*anyopaque,
3640 file: File,
3641 owner: ?File.Uid,
3642 group: ?File.Gid,
3643) File.SetOwnerError!void {
3644 const ev: *Evented = @ptrCast(@alignCast(userdata));
3645 _ = ev;
3646 return fchown(file.handle, owner, group);
3647}
3648
3649fn fileSetPermissions(
3650 userdata: ?*anyopaque,
3651 file: File,
3652 permissions: File.Permissions,
3653) File.SetPermissionsError!void {
3654 const ev: *Evented = @ptrCast(@alignCast(userdata));
3655 return ev.fchmod(file.handle, permissions.toMode());
3656}
3657
3658fn fchmod(ev: *Evented, fd: c.fd_t, mode: c.mode_t) File.SetPermissionsError!void {
3659 _ = ev;
3660 while (true) switch (c.errno(c.fchmod(fd, mode))) {
3661 .SUCCESS => return,
3662 .INTR => {},
3663 .BADF => |err| return errnoBug(err),
3664 .FAULT => |err| return errnoBug(err),
3665 .INVAL => |err| return errnoBug(err),
3666 .ACCES => return error.AccessDenied,
3667 .IO => return error.InputOutput,
3668 .LOOP => return error.SymLinkLoop,
3669 .NOENT => return error.FileNotFound,
3670 .NOMEM => return error.SystemResources,
3671 .NOTDIR => return error.FileNotFound,
3672 .PERM => return error.PermissionDenied,
3673 .ROFS => return error.ReadOnlyFileSystem,
3674 else => |err| return unexpectedErrno(err),
3675 };
3676}
3677
3678fn fileSetTimestamps(
3679 userdata: ?*anyopaque,
3680 file: File,
3681 options: File.SetTimestampsOptions,
3682) File.SetTimestampsError!void {
3683 const ev: *Evented = @ptrCast(@alignCast(userdata));
3684 _ = ev;
3685
3686 var times_buffer: [2]c.timespec = undefined;
3687 const times = if (options.modify_timestamp == .now and options.access_timestamp == .now) null else p: {
3688 times_buffer = .{
3689 setTimestampToPosix(options.access_timestamp),
3690 setTimestampToPosix(options.modify_timestamp),
3691 };
3692 break :p &times_buffer;
3693 };
3694
3695 while (true) switch (c.errno(c.futimens(file.handle, times))) {
3696 .SUCCESS => return,
3697 .INTR => {},
3698 .BADF => |err| return errnoBug(err), // always a race condition
3699 .FAULT => |err| return errnoBug(err),
3700 .INVAL => |err| return errnoBug(err),
3701 .ACCES => return error.AccessDenied,
3702 .PERM => return error.PermissionDenied,
3703 .ROFS => return error.ReadOnlyFileSystem,
3704 else => |err| return unexpectedErrno(err),
3705 };
3706}
3707
3708fn fileLock(userdata: ?*anyopaque, file: File, lock: File.Lock) File.LockError!void {
3709 const ev: *Evented = @ptrCast(@alignCast(userdata));
3710 _ = ev;
3711 const operation: i32 = switch (lock) {
3712 .none => c.LOCK.UN,
3713 .shared => c.LOCK.SH,
3714 .exclusive => c.LOCK.EX,
3715 };
3716 while (true) switch (c.errno(c.flock(file.handle, operation))) {
3717 .SUCCESS => return,
3718 .INTR => {},
3719 .BADF => |err| return errnoBug(err),
3720 .INVAL => |err| return errnoBug(err), // invalid parameters
3721 .NOLCK => return error.SystemResources,
3722 .AGAIN => |err| return errnoBug(err),
3723 .OPNOTSUPP => return error.FileLocksUnsupported,
3724 else => |err| return unexpectedErrno(err),
3725 };
3726}
3727
3728fn fileTryLock(userdata: ?*anyopaque, file: File, lock: File.Lock) File.LockError!bool {
3729 const ev: *Evented = @ptrCast(@alignCast(userdata));
3730 _ = ev;
3731 const operation: i32 = switch (lock) {
3732 .none => c.LOCK.UN,
3733 .shared => c.LOCK.SH | c.LOCK.NB,
3734 .exclusive => c.LOCK.EX | c.LOCK.NB,
3735 };
3736 while (true) switch (c.errno(c.flock(file.handle, operation))) {
3737 .SUCCESS => return true,
3738 .INTR => {},
3739 .AGAIN => return false,
3740 .BADF => |err| return errnoBug(err),
3741 .INVAL => |err| return errnoBug(err), // invalid parameters
3742 .NOLCK => return error.SystemResources,
3743 .OPNOTSUPP => return error.FileLocksUnsupported,
3744 else => |err| return unexpectedErrno(err),
3745 };
3746}
3747
3748fn fileUnlock(userdata: ?*anyopaque, file: File) void {
3749 const ev: *Evented = @ptrCast(@alignCast(userdata));
3750 _ = ev;
3751 while (true) switch (c.errno(c.flock(file.handle, c.LOCK.UN))) {
3752 .SUCCESS => return,
3753 .INTR => {},
3754 .AGAIN => return recoverableOsBugDetected(), // unlocking can't block
3755 .BADF => return recoverableOsBugDetected(), // File descriptor used after closed.
3756 .INVAL => return recoverableOsBugDetected(), // invalid parameters
3757 .NOLCK => return recoverableOsBugDetected(), // Resource deallocation.
3758 .OPNOTSUPP => return recoverableOsBugDetected(), // We already got the lock.
3759 else => return recoverableOsBugDetected(), // Resource deallocation must succeed.
3760 };
3761}
3762
3763fn fileDowngradeLock(userdata: ?*anyopaque, file: File) File.DowngradeLockError!void {
3764 const ev: *Evented = @ptrCast(@alignCast(userdata));
3765 _ = ev;
3766 const operation = c.LOCK.SH | c.LOCK.NB;
3767 while (true) switch (c.errno(c.flock(file.handle, operation))) {
3768 .SUCCESS => return,
3769 .INTR => {},
3770 .AGAIN => |err| return errnoBug(err), // File was not locked in exclusive mode.
3771 .BADF => |err| return errnoBug(err),
3772 .INVAL => |err| return errnoBug(err), // invalid parameters
3773 .NOLCK => |err| return errnoBug(err), // Lock already obtained.
3774 .OPNOTSUPP => |err| return errnoBug(err), // Lock already obtained.
3775 else => |err| return unexpectedErrno(err),
3776 };
3777}
3778
3779fn fileRealPath(userdata: ?*anyopaque, file: File, out_buffer: []u8) File.RealPathError!usize {
3780 const ev: *Evented = @ptrCast(@alignCast(userdata));
3781 _ = ev;
3782 var buffer: [c.PATH_MAX]u8 = undefined;
3783 @memset(&buffer, 0);
3784 while (true) {
3785 switch (c.errno(c.fcntl(file.handle, c.F.GETPATH, &buffer))) {
3786 .SUCCESS => break,
3787 .INTR => {},
3788 .ACCES => return error.AccessDenied,
3789 .BADF => return error.FileNotFound,
3790 .NOENT => return error.FileNotFound,
3791 .NOMEM => return error.SystemResources,
3792 .NOSPC => return error.NameTooLong,
3793 .RANGE => return error.NameTooLong,
3794 else => |err| return unexpectedErrno(err),
3795 }
3796 }
3797 const n = std.mem.findScalar(u8, &buffer, 0) orelse buffer.len;
3798 if (n > out_buffer.len) return error.NameTooLong;
3799 @memcpy(out_buffer[0..n], buffer[0..n]);
3800 return n;
3801}
3802
3803fn fileHardLink(
3804 userdata: ?*anyopaque,
3805 file: File,
3806 new_dir: Dir,
3807 new_sub_path: []const u8,
3808 options: File.HardLinkOptions,
3809) File.HardLinkError!void {
3810 const ev: *Evented = @ptrCast(@alignCast(userdata));
3811 _ = ev;
3812 _ = file;
3813 _ = new_dir;
3814 _ = new_sub_path;
3815 _ = options;
3816 return error.OperationUnsupported;
3817}
3818
3819fn linkat(
3820 old_dir: c.fd_t,
3821 old_path: [*:0]const u8,
3822 new_dir: c.fd_t,
3823 new_path: [*:0]const u8,
3824 flags: u32,
3825) File.HardLinkError!void {
3826 while (true) switch (c.errno(c.linkat(old_dir, old_path, new_dir, new_path, flags))) {
3827 .SUCCESS => return,
3828 .INTR => {},
3829 .ACCES => return error.AccessDenied,
3830 .DQUOT => return error.DiskQuota,
3831 .EXIST => return error.PathAlreadyExists,
3832 .IO => return error.HardwareFailure,
3833 .LOOP => return error.SymLinkLoop,
3834 .MLINK => return error.LinkQuotaExceeded,
3835 .NAMETOOLONG => return error.NameTooLong,
3836 .NOENT => return error.FileNotFound,
3837 .NOMEM => return error.SystemResources,
3838 .NOSPC => return error.NoSpaceLeft,
3839 .NOTDIR => return error.NotDir,
3840 .PERM => return error.PermissionDenied,
3841 .ROFS => return error.ReadOnlyFileSystem,
3842 .XDEV => return error.CrossDevice,
3843 .ILSEQ => return error.BadPathName,
3844 .FAULT => |err| return errnoBug(err),
3845 .INVAL => |err| return errnoBug(err),
3846 else => |err| return unexpectedErrno(err),
3847 };
3848}
3849
3850fn fileMemoryMapCreate(
3851 userdata: ?*anyopaque,
3852 file: File,
3853 options: File.MemoryMap.CreateOptions,
3854) File.MemoryMap.CreateError!File.MemoryMap {
3855 const ev: *Evented = @ptrCast(@alignCast(userdata));
3856 _ = ev;
3857
3858 const prot: c.PROT = .{
3859 .READ = options.protection.read,
3860 .WRITE = options.protection.write,
3861 .EXEC = options.protection.execute,
3862 };
3863 const flags: c.MAP = .{
3864 .TYPE = .SHARED,
3865 };
3866
3867 const page_align = std.heap.page_size_min;
3868
3869 const contents = while (true) {
3870 const casted_offset = std.math.cast(i64, options.offset) orelse return error.Unseekable;
3871 const rc = c.mmap(null, options.len, prot, flags, file.handle, casted_offset);
3872 const err: c.E = if (rc != c.MAP_FAILED) .SUCCESS else @fromBackingInt(@intCast(c._errno().*));
3873 switch (err) {
3874 .SUCCESS => break @as([*]align(page_align) u8, @ptrCast(@alignCast(rc)))[0..options.len],
3875 .INTR => {},
3876 .ACCES => return error.AccessDenied,
3877 .AGAIN => return error.LockedMemoryLimitExceeded,
3878 .MFILE => return error.ProcessFdQuotaExceeded,
3879 .NFILE => return error.SystemFdQuotaExceeded,
3880 .NOMEM => return error.OutOfMemory,
3881 .PERM => return error.PermissionDenied,
3882 .OVERFLOW => return error.Unseekable,
3883 .BADF => return errnoBug(err), // Always a race condition.
3884 .INVAL => return errnoBug(err), // Invalid parameters to mmap()
3885 else => return unexpectedErrno(err),
3886 }
3887 };
3888 return .{
3889 .file = file,
3890 .offset = options.offset,
3891 .memory = contents,
3892 .section = {},
3893 };
3894}
3895
3896fn fileMemoryMapDestroy(userdata: ?*anyopaque, mm: *File.MemoryMap) void {
3897 const ev: *Evented = @ptrCast(@alignCast(userdata));
3898 _ = ev;
3899 const memory = mm.memory;
3900 if (memory.len == 0) return;
3901 switch (c.errno(c.munmap(memory.ptr, memory.len))) {
3902 .SUCCESS => {},
3903 else => |err| if (builtin.mode == .debug)
3904 std.log.err("failed to unmap {d} bytes at {*}: {t}", .{ memory.len, memory.ptr, err }),
3905 }
3906 mm.* = undefined;
3907}
3908
3909fn fileMemoryMapSetLength(
3910 userdata: ?*anyopaque,
3911 mm: *File.MemoryMap,
3912 new_len: usize,
3913) File.MemoryMap.SetLengthError!void {
3914 const ev: *Evented = @ptrCast(@alignCast(userdata));
3915 _ = ev;
3916
3917 const page_size = std.heap.pageSize();
3918 const alignment: Alignment = .fromByteUnits(page_size);
3919 const old_memory = mm.memory;
3920
3921 if (alignment.forward(new_len) == alignment.forward(old_memory.len)) {
3922 mm.memory.len = new_len;
3923 return;
3924 }
3925 return error.OperationUnsupported;
3926}
3927
3928fn fileMemoryMapRead(userdata: ?*anyopaque, mm: *File.MemoryMap) File.ReadPositionalError!void {
3929 const ev: *Evented = @ptrCast(@alignCast(userdata));
3930 _ = ev;
3931 _ = mm;
3932}
3933
3934fn fileMemoryMapWrite(userdata: ?*anyopaque, mm: *File.MemoryMap) File.WritePositionalError!void {
3935 const ev: *Evented = @ptrCast(@alignCast(userdata));
3936 _ = ev;
3937 _ = mm;
3938}
3939
3940fn processExecutableOpen(
3941 userdata: ?*anyopaque,
3942 flags: Dir.OpenFileOptions,
3943) process.OpenExecutableError!File {
3944 const ev: *Evented = @ptrCast(@alignCast(userdata));
3945 // _NSGetExecutablePath() returns a path that might be a symlink to
3946 // the executable. Here it does not matter since we open it.
3947 var symlink_path_buf: [c.PATH_MAX + 1]u8 = undefined;
3948 var n: u32 = symlink_path_buf.len;
3949 const rc = c._NSGetExecutablePath(&symlink_path_buf, &n);
3950 if (rc != 0) return error.NameTooLong;
3951 const symlink_path = std.mem.sliceTo(&symlink_path_buf, 0);
3952 return dirOpenFile(ev, .cwd(), symlink_path, flags);
3953}
3954
3955fn processExecutablePath(userdata: ?*anyopaque, out_buffer: []u8) process.ExecutablePathError!usize {
3956 const ev: *Evented = @ptrCast(@alignCast(userdata));
3957 // _NSGetExecutablePath() returns a path that might be a symlink to
3958 // the executable.
3959 var symlink_path_buf: [c.PATH_MAX + 1]u8 = undefined;
3960 var n: u32 = symlink_path_buf.len;
3961 const rc = c._NSGetExecutablePath(&symlink_path_buf, &n);
3962 if (rc != 0) return error.NameTooLong;
3963 const symlink_path = std.mem.sliceTo(&symlink_path_buf, 0);
3964 assert(Dir.path.isAbsolute(symlink_path));
3965 return dirRealPathFile(ev, .cwd(), symlink_path, out_buffer) catch |err| switch (err) {
3966 error.NetworkNotFound => unreachable, // Windows-only
3967 error.FileBusy => unreachable, // Windows-only
3968 else => |e| return e,
3969 };
3970}
3971
3972fn lockStderr(userdata: ?*anyopaque, terminal_mode: ?Io.Terminal.Mode) Io.Cancelable!Io.LockedStderr {
3973 const ev: *Evented = @ptrCast(@alignCast(userdata));
3974 try ev.stderr_mutex.lock(ev);
3975 errdefer ev.stderr_mutex.unlock();
3976 return ev.initLockedStderr(terminal_mode);
3977}
3978
3979fn tryLockStderr(
3980 userdata: ?*anyopaque,
3981 terminal_mode: ?Io.Terminal.Mode,
3982) Io.Cancelable!?Io.LockedStderr {
3983 const ev: *Evented = @ptrCast(@alignCast(userdata));
3984 if (!ev.stderr_mutex.tryLock()) return null;
3985 errdefer ev.stderr_mutex.unlock();
3986 return try ev.initLockedStderr(terminal_mode);
3987}
3988
3989fn initLockedStderr(ev: *Evented, terminal_mode: ?Io.Terminal.Mode) Io.Cancelable!Io.LockedStderr {
3990 ev.init_stderr_writer.once(ev, &initStderrWriter);
3991 return .{
3992 .file_writer = &ev.stderr_writer,
3993 .terminal_mode = terminal_mode orelse ev.stderr_mode,
3994 };
3995}
3996
3997fn initStderrWriter(context: ?*anyopaque) callconv(.c) void {
3998 const ev: *Evented = @ptrCast(@alignCast(context));
3999 const cancel_protection = swapCancelProtection(ev, .blocked);
4000 defer assert(swapCancelProtection(ev, cancel_protection) == .blocked);
4001 ev.scan_environ.once(ev, &scanEnviron);
4002 const NO_COLOR = ev.environ.exist.NO_COLOR;
4003 const CLICOLOR_FORCE = ev.environ.exist.CLICOLOR_FORCE;
4004 ev.stderr_mode = Io.Terminal.Mode.detect(
4005 ev.io(),
4006 ev.stderr_writer.file,
4007 NO_COLOR,
4008 CLICOLOR_FORCE,
4009 ) catch |err| switch (err) {
4010 error.Canceled => unreachable, // blocked
4011 };
4012}
4013
4014fn unlockStderr(userdata: ?*anyopaque) void {
4015 const ev: *Evented = @ptrCast(@alignCast(userdata));
4016 if (ev.stderr_writer.err == null) ev.stderr_writer.interface.flush() catch {};
4017 if (ev.stderr_writer.err) |err| {
4018 switch (err) {
4019 error.Canceled => Thread.current().currentFiber().cancel_protection.recancel(),
4020 else => {},
4021 }
4022 ev.stderr_writer.err = null;
4023 }
4024 ev.stderr_writer.interface.end = 0;
4025 ev.stderr_writer.interface.buffer.len = 0;
4026 ev.stderr_mutex.unlock();
4027}
4028
4029fn processCurrentPath(userdata: ?*anyopaque, buffer: []u8) process.CurrentPathError!usize {
4030 const ev: *Evented = @ptrCast(@alignCast(userdata));
4031 _ = ev;
4032 const err: c.E = if (c.getcwd(buffer.ptr, buffer.len)) |_| .SUCCESS else @fromBackingInt(@intCast(c._errno().*));
4033 switch (err) {
4034 .SUCCESS => return std.mem.findScalar(u8, buffer, 0).?,
4035 .NOENT => return error.CurrentDirUnlinked,
4036 .RANGE => return error.NameTooLong,
4037 .FAULT => |e| return errnoBug(e),
4038 .INVAL => |e| return errnoBug(e),
4039 else => return unexpectedErrno(err),
4040 }
4041}
4042
4043fn processSetCurrentDir(userdata: ?*anyopaque, dir: Dir) process.SetCurrentDirError!void {
4044 const ev: *Evented = @ptrCast(@alignCast(userdata));
4045 _ = ev;
4046 if (dir.handle == c.AT.FDCWD) return;
4047 while (true) switch (c.errno(c.fchdir(dir.handle))) {
4048 .SUCCESS => return,
4049 .INTR => {},
4050 .ACCES => return error.AccessDenied,
4051 .NOTDIR => return error.NotDir,
4052 .IO => return error.FileSystem,
4053 .BADF => |err| return errnoBug(err),
4054 else => |err| return unexpectedErrno(err),
4055 };
4056}
4057
4058fn processSetCurrentPath(userdata: ?*anyopaque, dir_path: []const u8) process.SetCurrentPathError!void {
4059 const ev: *Evented = @ptrCast(@alignCast(userdata));
4060 _ = ev;
4061 var path_buffer: [c.PATH_MAX]u8 = undefined;
4062 const dir_path_posix = try pathToPosix(dir_path, &path_buffer);
4063 while (true) switch (c.errno(c.chdir(dir_path_posix))) {
4064 .SUCCESS => return,
4065 .INTR => {},
4066 .ACCES => return error.AccessDenied,
4067 .IO => return error.FileSystem,
4068 .LOOP => return error.SymLinkLoop,
4069 .NAMETOOLONG => return error.NameTooLong,
4070 .NOENT => return error.FileNotFound,
4071 .NOMEM => return error.SystemResources,
4072 .NOTDIR => return error.NotDir,
4073 .ILSEQ => return error.BadPathName,
4074 .FAULT => |err| return errnoBug(err),
4075 else => |err| return unexpectedErrno(err),
4076 };
4077}
4078
4079fn processReplace(userdata: ?*anyopaque, options: process.ReplaceOptions) process.ReplaceError {
4080 const ev: *Evented = @ptrCast(@alignCast(userdata));
4081
4082 if (!process.can_replace) return error.OperationUnsupported;
4083
4084 ev.scan_environ.once(ev, &scanEnviron); // for PATH
4085 const PATH = ev.environ.string.PATH orelse default_PATH;
4086
4087 var arena_allocator = std.heap.ArenaAllocator.init(ev.allocator());
4088 defer arena_allocator.deinit();
4089 const arena = arena_allocator.allocator();
4090
4091 const argv_buf = try arena.allocSentinel(?[*:0]const u8, options.argv.len, null);
4092 for (options.argv, 0..) |arg, i| argv_buf[i] = (try arena.dupeSentinel(u8, arg, 0)).ptr;
4093
4094 const env_block = env_block: {
4095 const prog_fd: i32 = -1;
4096 if (options.environ_map) |environ_map| break :env_block try environ_map.createPosixBlock(arena, .{
4097 .zig_progress_fd = prog_fd,
4098 });
4099 break :env_block try ev.environ.process_environ.createPosixBlock(arena, .{
4100 .zig_progress_fd = prog_fd,
4101 });
4102 };
4103
4104 return ev.execv(options.expand_arg0, argv_buf.ptr[0].?, argv_buf.ptr, env_block, PATH);
4105}
4106
4107fn processReplacePath(
4108 userdata: ?*anyopaque,
4109 dir: Dir,
4110 options: process.ReplaceOptions,
4111) process.ReplaceError {
4112 const ev: *Evented = @ptrCast(@alignCast(userdata));
4113 _ = ev;
4114 _ = dir;
4115 _ = options;
4116 @panic("TODO processReplacePath");
4117}
4118
4119fn processSpawn(userdata: ?*anyopaque, options: process.SpawnOptions) process.SpawnError!process.Child {
4120 const ev: *Evented = @ptrCast(@alignCast(userdata));
4121 const spawned = try ev.spawn(options);
4122 defer fileClose(ev, &.{spawned.err_pipe});
4123
4124 // Wait for the child to report any errors in or before `execvpe`.
4125 var child_err: ForkBailError = undefined;
4126 ev.readAll(spawned.err_pipe, @ptrCast(&child_err)) catch |read_err| {
4127 switch (read_err) {
4128 error.Canceled => unreachable, // blocked
4129 error.EndOfStream => {
4130 // Write end closed by CLOEXEC at the time of the `execvpe` call,
4131 // indicating success.
4132 },
4133 else => {
4134 // Problem reading the error from the error reporting pipe. We
4135 // don't know if the child is alive or dead. Better to assume it is
4136 // alive so the resource does not risk being leaked.
4137 },
4138 }
4139 return .{
4140 .id = spawned.pid,
4141 .thread_handle = {},
4142 .stdin = spawned.stdin,
4143 .stdout = spawned.stdout,
4144 .stderr = spawned.stderr,
4145 .request_resource_usage_statistics = options.request_resource_usage_statistics,
4146 };
4147 };
4148 return child_err;
4149}
4150
4151fn processSpawnPath(
4152 userdata: ?*anyopaque,
4153 dir: Dir,
4154 options: process.SpawnOptions,
4155) process.SpawnError!process.Child {
4156 const ev: *Evented = @ptrCast(@alignCast(userdata));
4157 _ = ev;
4158 _ = dir;
4159 _ = options;
4160 @panic("TODO processSpawnPath");
4161}
4162
4163const prog_fileno = @max(c.STDIN_FILENO, c.STDOUT_FILENO, c.STDERR_FILENO) + 1;
4164
4165const Spawned = struct {
4166 pid: c.pid_t,
4167 err_pipe: File,
4168 stdin: ?File,
4169 stdout: ?File,
4170 stderr: ?File,
4171};
4172fn spawn(ev: *Evented, options: process.SpawnOptions) process.SpawnError!Spawned {
4173 // The child process does need to access (one end of) these pipes. However,
4174 // we must initially set CLOEXEC to avoid a race condition. If another thread
4175 // is racing to spawn a different child process, we don't want it to inherit
4176 // these FDs in any scenario; that would mean that, for instance, calls to
4177 // `poll` from the parent would not report the child's stdout as closing when
4178 // expected, since the other child may retain a reference to the write end of
4179 // the pipe. So, we create the pipes with CLOEXEC initially. After fork, we
4180 // need to do something in the new child to make sure we preserve the reference
4181 // we want. We could use `fcntl` to remove CLOEXEC from the FD, but as it
4182 // turns out, we `dup2` everything anyway, so there's no need!
4183 const pipe_flags: c.O = .{ .CLOEXEC = true };
4184
4185 const stdin_pipe = if (options.stdin == .pipe) try pipe2(pipe_flags) else undefined;
4186 errdefer if (options.stdin == .pipe) {
4187 destroyPipe(stdin_pipe);
4188 };
4189
4190 const stdout_pipe = if (options.stdout == .pipe) try pipe2(pipe_flags) else undefined;
4191 errdefer if (options.stdout == .pipe) {
4192 destroyPipe(stdout_pipe);
4193 };
4194
4195 const stderr_pipe = if (options.stderr == .pipe) try pipe2(pipe_flags) else undefined;
4196 errdefer if (options.stderr == .pipe) {
4197 destroyPipe(stderr_pipe);
4198 };
4199
4200 const any_ignore =
4201 options.stdin == .ignore or options.stdout == .ignore or options.stderr == .ignore;
4202 const dev_null_file = if (any_ignore) dev_null_file: {
4203 ev.open_dev_null.once(ev, &openDevNullFile);
4204 break :dev_null_file try ev.dev_null_file;
4205 } else undefined;
4206
4207 const prog_pipe: [2]c.fd_t = if (options.progress_node.index != .none)
4208 // We use CLOEXEC for the same reason as in `pipe_flags`.
4209 try pipe2(.{ .NONBLOCK = true, .CLOEXEC = true })
4210 else
4211 .{ -1, -1 };
4212 errdefer destroyPipe(prog_pipe);
4213
4214 var arena_allocator = std.heap.ArenaAllocator.init(ev.allocator());
4215 defer arena_allocator.deinit();
4216 const arena = arena_allocator.allocator();
4217
4218 // The POSIX standard does not allow malloc() between fork() and execve(),
4219 // and this allocator may be a libc allocator.
4220 // I have personally observed the child process deadlocking when it tries
4221 // to call malloc() due to a heap allocation between fork() and execve(),
4222 // in musl v1.1.24.
4223 // Additionally, we want to reduce the number of possible ways things
4224 // can fail between fork() and execve().
4225 // Therefore, we do all the allocation for the execve() before the fork().
4226 // This means we must do the null-termination of argv and env vars here.
4227 const argv_buf = try arena.allocSentinel(?[*:0]const u8, options.argv.len, null);
4228 for (options.argv, 0..) |arg, i| argv_buf[i] = (try arena.dupeSentinel(u8, arg, 0)).ptr;
4229
4230 const env_block = env_block: {
4231 const prog_fd: i32 = if (prog_pipe[1] == -1) -1 else prog_fileno;
4232 if (options.environ_map) |environ_map| break :env_block try environ_map.createPosixBlock(arena, .{
4233 .zig_progress_fd = prog_fd,
4234 });
4235 break :env_block try ev.environ.process_environ.createPosixBlock(arena, .{
4236 .zig_progress_fd = prog_fd,
4237 });
4238 };
4239
4240 // This pipe communicates to the parent errors in the child between `fork` and `execvpe`.
4241 // It is closed by the child (via CLOEXEC) without writing if `execvpe` succeeds.
4242 const err_pipe: [2]File = err_pipe: {
4243 const err_pipe = try pipe2(.{ .CLOEXEC = true });
4244 break :err_pipe .{
4245 .{ .handle = err_pipe[0], .flags = .{ .nonblocking = false } },
4246 .{ .handle = err_pipe[1], .flags = .{ .nonblocking = false } },
4247 };
4248 };
4249 errdefer fileClose(ev, &err_pipe);
4250
4251 ev.scan_environ.once(ev, &scanEnviron); // for PATH
4252 const PATH = ev.environ.string.PATH orelse default_PATH;
4253
4254 const pid_result: c.pid_t = fork: {
4255 const rc = c.fork();
4256 switch (c.errno(rc)) {
4257 .SUCCESS => break :fork @intCast(rc),
4258 .AGAIN => return error.SystemResources,
4259 .NOMEM => return error.SystemResources,
4260 .NOSYS => return error.OperationUnsupported,
4261 else => |err| return unexpectedErrno(err),
4262 }
4263 };
4264
4265 if (pid_result == 0) {
4266 defer comptime unreachable; // We are the child.
4267 const err = ev.setUpChild(.{
4268 .stdin_pipe = stdin_pipe[0],
4269 .stdout_pipe = stdout_pipe[1],
4270 .stderr_pipe = stderr_pipe[1],
4271 .dev_null_fd = dev_null_file.handle,
4272 .prog_pipe = prog_pipe[1],
4273 .argv_buf = argv_buf,
4274 .env_block = env_block,
4275 .PATH = PATH,
4276 .spawn = options,
4277 });
4278 ev.writeAll(err_pipe[1], @ptrCast(&err)) catch {};
4279 c.exit(1);
4280 }
4281
4282 const pid: c.pid_t = @intCast(pid_result); // We are the parent.
4283 errdefer comptime unreachable; // The child is forked; we must not error from now on
4284
4285 fileClose(ev, err_pipe[1..2]); // make sure only the child holds the write end open
4286
4287 if (options.stdin == .pipe) closeFd(stdin_pipe[0]);
4288 if (options.stdout == .pipe) closeFd(stdout_pipe[1]);
4289 if (options.stderr == .pipe) closeFd(stderr_pipe[1]);
4290
4291 if (prog_pipe[1] != -1) closeFd(prog_pipe[1]);
4292
4293 options.progress_node.setIpcFile(ev, .{ .handle = prog_pipe[0], .flags = .{ .nonblocking = true } });
4294
4295 return .{
4296 .pid = pid,
4297 .err_pipe = err_pipe[0],
4298 .stdin = switch (options.stdin) {
4299 .pipe => .{ .handle = stdin_pipe[1], .flags = .{ .nonblocking = false } },
4300 else => null,
4301 },
4302 .stdout = switch (options.stdout) {
4303 .pipe => .{ .handle = stdout_pipe[0], .flags = .{ .nonblocking = false } },
4304 else => null,
4305 },
4306 .stderr = switch (options.stderr) {
4307 .pipe => .{ .handle = stderr_pipe[0], .flags = .{ .nonblocking = false } },
4308 else => null,
4309 },
4310 };
4311}
4312
4313fn openDevNullFile(context: ?*anyopaque) callconv(.c) void {
4314 const ev: *Evented = @ptrCast(@alignCast(context));
4315 ev.dev_null_file = dirOpenFile(ev, .cwd(), "/dev/null", .{ .mode = .read_write });
4316}
4317
4318/// Errors that can occur between fork() and execv()
4319const ForkBailError = process.SetCurrentDirError || ChdirError ||
4320 process.SpawnError || process.ReplaceError;
4321fn setUpChild(ev: *Evented, options: struct {
4322 stdin_pipe: c.fd_t,
4323 stdout_pipe: c.fd_t,
4324 stderr_pipe: c.fd_t,
4325 dev_null_fd: c.fd_t,
4326 prog_pipe: c.fd_t,
4327 argv_buf: [:null]?[*:0]const u8,
4328 env_block: process.Environ.Block,
4329 PATH: []const u8,
4330 spawn: process.SpawnOptions,
4331}) ForkBailError {
4332 try ev.setUpChildIo(
4333 options.spawn.stdin,
4334 options.stdin_pipe,
4335 c.STDIN_FILENO,
4336 options.dev_null_fd,
4337 );
4338 try ev.setUpChildIo(
4339 options.spawn.stdout,
4340 options.stdout_pipe,
4341 c.STDOUT_FILENO,
4342 options.dev_null_fd,
4343 );
4344 try ev.setUpChildIo(
4345 options.spawn.stderr,
4346 options.stderr_pipe,
4347 c.STDERR_FILENO,
4348 options.dev_null_fd,
4349 );
4350
4351 switch (options.spawn.cwd) {
4352 .inherit => {},
4353 .dir => |cwd_dir| try processSetCurrentDir(ev, cwd_dir),
4354 .path => |cwd_path| try processSetCurrentPath(ev, cwd_path),
4355 }
4356
4357 // Must happen after fchdir above, the cwd file descriptor might be
4358 // equal to prog_fileno and be clobbered by this dup2 call.
4359 if (options.prog_pipe != -1) try ev.dup2(options.prog_pipe, prog_fileno);
4360
4361 if (options.spawn.gid) |gid| while (true) switch (c.errno(c.setregid(gid, gid))) {
4362 .SUCCESS => break,
4363 .INTR => {},
4364 .AGAIN => return error.ResourceLimitReached,
4365 .INVAL => return error.InvalidUserId,
4366 .PERM => return error.PermissionDenied,
4367 else => return error.Unexpected,
4368 };
4369
4370 if (options.spawn.uid) |uid| while (true) switch (c.errno(c.setreuid(uid, uid))) {
4371 .SUCCESS => break,
4372 .INTR => {},
4373 .AGAIN => return error.ResourceLimitReached,
4374 .INVAL => return error.InvalidUserId,
4375 .PERM => return error.PermissionDenied,
4376 else => return error.Unexpected,
4377 };
4378
4379 if (options.spawn.pgid) |pid| while (true) switch (c.errno(c.setpgid(0, pid))) {
4380 .SUCCESS => break,
4381 .INTR => {},
4382 .ACCES => return error.ProcessAlreadyExec,
4383 .INVAL => return error.InvalidProcessGroupId,
4384 .PERM => return error.PermissionDenied,
4385 else => return error.Unexpected,
4386 };
4387
4388 if (options.spawn.start_suspended) while (true) switch (c.errno(c.kill(0, .STOP))) {
4389 .SUCCESS => break,
4390 .INTR => {},
4391 .PERM => return error.PermissionDenied,
4392 else => return error.Unexpected,
4393 };
4394
4395 return ev.execv(
4396 options.spawn.expand_arg0,
4397 options.argv_buf.ptr[0].?,
4398 options.argv_buf.ptr,
4399 options.env_block,
4400 options.PATH,
4401 );
4402}
4403
4404fn setUpChildIo(
4405 ev: *Evented,
4406 stdio: process.SpawnOptions.StdIo,
4407 pipe_fd: c.fd_t,
4408 std_fileno: i32,
4409 dev_null_fd: c.fd_t,
4410) !void {
4411 switch (stdio) {
4412 .pipe => try ev.dup2(pipe_fd, std_fileno),
4413 .close => closeFd(std_fileno),
4414 .inherit => {},
4415 .ignore => try ev.dup2(dev_null_fd, std_fileno),
4416 .file => |file| try ev.dup2(file.handle, std_fileno),
4417 }
4418}
4419
4420const PipeError = error{
4421 SystemFdQuotaExceeded,
4422 ProcessFdQuotaExceeded,
4423} || Io.UnexpectedError;
4424
4425fn pipe2(flags: c.O) PipeError![2]c.fd_t {
4426 var fds: [2]c.fd_t = undefined;
4427
4428 while (true) switch (c.errno(c.pipe(&fds))) {
4429 .SUCCESS => break,
4430 .INTR => {},
4431 .NFILE => return error.SystemFdQuotaExceeded,
4432 .MFILE => return error.ProcessFdQuotaExceeded,
4433 else => |err| return unexpectedErrno(err),
4434 };
4435 errdefer {
4436 closeFd(fds[0]);
4437 closeFd(fds[1]);
4438 }
4439
4440 // https://github.com/ziglang/zig/issues/18882
4441 if (@as(u32, @bitCast(flags)) == 0) return fds;
4442
4443 // CLOEXEC is special, it's a file descriptor flag and must be set using
4444 // F.SETFD.
4445 if (flags.CLOEXEC) for (fds) |fd| while (true) switch (c.errno(c.fcntl(fd, c.F.SETFD, @as(u32, c.FD_CLOEXEC)))) {
4446 .SUCCESS => break,
4447 .INTR => {},
4448 else => |err| return unexpectedErrno(err),
4449 };
4450
4451 const new_flags: u32 = f: {
4452 var new_flags = flags;
4453 new_flags.CLOEXEC = false;
4454 break :f @bitCast(new_flags);
4455 };
4456
4457 // Set every other flag affecting the file status using F.SETFL.
4458 if (new_flags != 0) for (fds) |fd| while (true) switch (c.errno(c.fcntl(fd, c.F.SETFL, new_flags))) {
4459 .SUCCESS => break,
4460 .INTR => {},
4461 .INVAL => |err| return errnoBug(err),
4462 else => |err| return unexpectedErrno(err),
4463 };
4464
4465 return fds;
4466}
4467
4468fn destroyPipe(pipe: [2]c.fd_t) void {
4469 if (pipe[0] != -1) closeFd(pipe[0]);
4470 if (pipe[0] != pipe[1]) closeFd(pipe[1]);
4471}
4472
4473const DupError = error{
4474 ProcessFdQuotaExceeded,
4475 SystemResources,
4476} || Io.UnexpectedError || Io.Cancelable;
4477fn dup2(ev: *Evented, old_fd: c.fd_t, new_fd: c.fd_t) DupError!void {
4478 _ = ev;
4479 while (true) switch (c.errno(c.dup2(old_fd, new_fd))) {
4480 .SUCCESS => return,
4481 .BUSY, .INTR => {},
4482 .INVAL => |err| return errnoBug(err), // invalid parameters
4483 .BADF => |err| return errnoBug(err), // use after free
4484 .MFILE => return error.ProcessFdQuotaExceeded,
4485 .NOMEM => return error.SystemResources,
4486 else => |err| return unexpectedErrno(err),
4487 };
4488}
4489
4490fn execv(
4491 ev: *Evented,
4492 arg0_expand: process.ArgExpansion,
4493 file: [*:0]const u8,
4494 child_argv: [*:null]?[*:0]const u8,
4495 env_block: process.Environ.PosixBlock,
4496 PATH: []const u8,
4497) process.ReplaceError {
4498 const file_slice = std.mem.sliceTo(file, 0);
4499 if (std.mem.findScalar(u8, file_slice, '/') != null) return ev.execvPath(file, child_argv, env_block);
4500
4501 // Use of PATH_MAX here is valid as the path_buf will be passed
4502 // directly to the operating system in posixExecvPath.
4503 var path_buf: [c.PATH_MAX]u8 = undefined;
4504 var it = std.mem.tokenizeScalar(u8, PATH, ':');
4505 var seen_eacces = false;
4506 var err: process.ReplaceError = error.FileNotFound;
4507
4508 // In case of expanding arg0 we must put it back if we return with an error.
4509 const prev_arg0 = child_argv[0];
4510 defer switch (arg0_expand) {
4511 .expand => child_argv[0] = prev_arg0,
4512 .no_expand => {},
4513 };
4514
4515 while (it.next()) |search_path| {
4516 const path_len = search_path.len + file_slice.len + 1;
4517 if (path_buf.len < path_len + 1) return error.NameTooLong;
4518 @memcpy(path_buf[0..search_path.len], search_path);
4519 path_buf[search_path.len] = '/';
4520 @memcpy(path_buf[search_path.len + 1 ..][0..file_slice.len], file_slice);
4521 path_buf[path_len] = 0;
4522 const full_path = path_buf[0..path_len :0].ptr;
4523 switch (arg0_expand) {
4524 .expand => child_argv[0] = full_path,
4525 .no_expand => {},
4526 }
4527 err = ev.execvPath(full_path, child_argv, env_block);
4528 switch (err) {
4529 error.AccessDenied => seen_eacces = true,
4530 error.FileNotFound, error.NotDir => {},
4531 else => |e| return e,
4532 }
4533 }
4534 if (seen_eacces) return error.AccessDenied;
4535 return err;
4536}
4537/// This function ignores PATH environment variable.
4538fn execvPath(
4539 ev: *Evented,
4540 path: [*:0]const u8,
4541 child_argv: [*:null]const ?[*:0]const u8,
4542 env_block: process.Environ.PosixBlock,
4543) process.ReplaceError {
4544 _ = ev;
4545 switch (c.errno(c.execve(path, child_argv, env_block.slice.ptr))) {
4546 .FAULT => |err| return errnoBug(err), // Bad pointer parameter.
4547 .@"2BIG" => return error.SystemResources,
4548 .MFILE => return error.ProcessFdQuotaExceeded,
4549 .NAMETOOLONG => return error.NameTooLong,
4550 .NFILE => return error.SystemFdQuotaExceeded,
4551 .NOMEM => return error.SystemResources,
4552 .ACCES => return error.AccessDenied,
4553 .PERM => return error.PermissionDenied,
4554 .INVAL => return error.InvalidExe,
4555 .NOEXEC => return error.InvalidExe,
4556 .IO => return error.FileSystem,
4557 .LOOP => return error.FileSystem,
4558 .ISDIR => return error.IsDir,
4559 .NOENT => return error.FileNotFound,
4560 .NOTDIR => return error.NotDir,
4561 .TXTBSY => return error.FileBusy,
4562 .BADEXEC => return error.InvalidExe,
4563 .BADARCH => return error.InvalidExe,
4564 else => |err| return unexpectedErrno(err),
4565 }
4566}
4567
4568fn childWait(userdata: ?*anyopaque, child: *process.Child) process.Child.WaitError!process.Child.Term {
4569 const ev: *Evented = @ptrCast(@alignCast(userdata));
4570 defer ev.childCleanup(child);
4571 const pid = child.id.?;
4572 const source = c.dispatch.source_create(
4573 .PROC,
4574 @bitCast(@as(isize, pid)),
4575 .{ .PROC = .{ .EXIT = true } },
4576 ev.queue,
4577 ) orelse return error.Unexpected;
4578 source.as_object().set_context(Thread.current().currentFiber());
4579 source.set_event_handler(&Fiber.@"resume");
4580 ev.yield(.{ .activate = source.as_object() });
4581 source.as_object().release();
4582 var status: c_int = undefined;
4583 var ru: c.rusage = undefined;
4584 const ru_ptr = if (child.request_resource_usage_statistics) &ru else null;
4585 while (true) switch (c.errno(c.wait4(pid, &status, 0, ru_ptr))) {
4586 .SUCCESS => {
4587 if (ru_ptr) |p| child.resource_usage_statistics.rusage = p.*;
4588 return statusToTerm(@bitCast(status));
4589 },
4590 .INTR => {},
4591 .CHILD => |err| return errnoBug(err), // Double-free.
4592 else => |err| return unexpectedErrno(err),
4593 };
4594}
4595
4596fn childKill(userdata: ?*anyopaque, child: *process.Child) void {
4597 const ev: *Evented = @ptrCast(@alignCast(userdata));
4598 defer ev.childCleanup(child);
4599 const pid = child.id.?;
4600 while (true) switch (c.errno(c.kill(pid, .TERM))) {
4601 .SUCCESS => break,
4602 .INTR => {},
4603 .PERM => return,
4604 .INVAL => |err| errnoBug(err) catch return,
4605 .SRCH => |err| errnoBug(err) catch return,
4606 else => |err| unexpectedErrno(err) catch return,
4607 };
4608 var status: c_int = undefined;
4609 while (true) switch (c.errno(c.wait4(pid, &status, 0, null))) {
4610 .SUCCESS => return,
4611 .INTR => {},
4612 .CHILD => |err| errnoBug(err) catch return, // Double-free.
4613 else => |err| unexpectedErrno(err) catch return,
4614 };
4615}
4616
4617fn childCleanup(ev: *Evented, child: *process.Child) void {
4618 if (child.stdin) |stdin| {
4619 fileClose(ev, &.{stdin});
4620 child.stdin = null;
4621 }
4622 if (child.stdout) |stdout| {
4623 fileClose(ev, &.{stdout});
4624 child.stdout = null;
4625 }
4626 if (child.stderr) |stderr| {
4627 fileClose(ev, &.{stderr});
4628 child.stderr = null;
4629 }
4630 child.id = null;
4631}
4632
4633fn progressParentFile(userdata: ?*anyopaque) std.Progress.ParentFileError!File {
4634 const ev: *Evented = @ptrCast(@alignCast(userdata));
4635 ev.scan_environ.once(ev, &scanEnviron);
4636 return ev.environ.zig_progress_file;
4637}
4638
4639fn scanEnviron(context: ?*anyopaque) callconv(.c) void {
4640 const ev: *Evented = @ptrCast(@alignCast(context));
4641 ev.environ.scan(ev.allocator());
4642}
4643
4644fn now(userdata: ?*anyopaque, clock: Io.Clock) Io.Timestamp {
4645 const ev: *Evented = @ptrCast(@alignCast(userdata));
4646 _ = ev;
4647 const clock_id: c.clockid_t = clockToPosix(clock);
4648 var timespec: c.timespec = undefined;
4649 switch (c.errno(c.clock_gettime(clock_id, &timespec))) {
4650 .SUCCESS => return timestampFromPosix(&timespec),
4651 else => return .zero,
4652 }
4653}
4654
4655fn clockResolution(userdata: ?*anyopaque, clock: Io.Clock) Io.Clock.ResolutionError!Io.Duration {
4656 const ev: *Evented = @ptrCast(@alignCast(userdata));
4657 _ = ev;
4658 const clock_id: c.clockid_t = clockToPosix(clock);
4659 var timespec: c.timespec = undefined;
4660 return switch (c.errno(c.clock_getres(clock_id, &timespec))) {
4661 .SUCCESS => .fromNanoseconds(nanosecondsFromPosix(&timespec)),
4662 .INVAL => return error.ClockUnavailable,
4663 else => |err| return unexpectedErrno(err),
4664 };
4665}
4666
4667const SleepWaiter = struct {
4668 sleeper: Sleeper = undefined,
4669 cancelable: Cancelable,
4670 timer: c.dispatch.source_t,
4671 started: bool = false,
4672
4673 fn start(context: ?*anyopaque) callconv(.c) void {
4674 const waiter: *SleepWaiter = @ptrCast(@alignCast(context));
4675 waiter.cancelable.enter(waiter.sleeper.fiber) catch |err| switch (err) {
4676 error.CancelRequested => waiter.timer.cancel(),
4677 };
4678 waiter.timer.as_object().activate();
4679 }
4680
4681 fn timedOut(context: ?*anyopaque) callconv(.c) void {
4682 const waiter: *SleepWaiter = @ptrCast(@alignCast(context));
4683 waiter.cancelable.leave(waiter.sleeper.fiber) catch |err| switch (err) {
4684 error.CancelRequested => return,
4685 };
4686 waiter.timer.cancel();
4687 }
4688
4689 fn canceled(context: ?*anyopaque) callconv(.c) void {
4690 const cancelable: *Cancelable = @ptrCast(@alignCast(context));
4691 const waiter: *SleepWaiter = @fieldParentPtr("cancelable", cancelable);
4692 cancelable.requested(waiter.sleeper.fiber);
4693 waiter.timer.cancel();
4694 }
4695
4696 fn wake(context: ?*anyopaque) callconv(.c) void {
4697 const waiter: *SleepWaiter = @ptrCast(@alignCast(context));
4698 var sleeper = waiter.sleeper;
4699 waiter.* = undefined;
4700 Sleeper.wake(&sleeper);
4701 }
4702};
4703
4704fn sleep(userdata: ?*anyopaque, timeout: Io.Timeout) Io.Cancelable!void {
4705 const ev: *Evented = @ptrCast(@alignCast(userdata));
4706 const queue = c.dispatch.queue_create_with_target(
4707 "org.ziglang.std.Io.Dispatch.sleep",
4708 .SERIAL(),
4709 ev.queue,
4710 ) orelse {
4711 log.warn("failed to create serial queue for sleep", .{});
4712 return ev.yield(.{ .after = ev.timeFromTimeout(timeout) });
4713 };
4714 defer queue.as_object().release();
4715 const timer = c.dispatch.source_create(.TIMER, 0, .none, queue) orelse {
4716 log.warn("failed to create timer for sleep", .{});
4717 return ev.yield(.{ .after = ev.timeFromTimeout(timeout) });
4718 };
4719 var waiter: SleepWaiter = .{
4720 .cancelable = .{ .queue = queue, .cancel = &SleepWaiter.canceled },
4721 .timer = timer,
4722 };
4723 timer.as_object().set_context(&waiter);
4724 timer.set_event_handler(&SleepWaiter.timedOut);
4725 timer.set_cancel_handler(&SleepWaiter.wake);
4726 timer.set_timer(ev.timeFromTimeout(timeout), c.dispatch.TIME_FOREVER, ev.leeway);
4727 ev.yield(.{ .sleep_wait = &waiter });
4728 timer.as_object().release();
4729 try waiter.cancelable.acknowledge(waiter.sleeper.fiber);
4730}
4731
4732fn timeFromTimeout(ev: *Evented, timeout: Io.Timeout) c.dispatch.time_t {
4733 return timeout: switch (timeout) {
4734 .none => .FOREVER,
4735 .duration => |duration| .time(switch (duration.clock) {
4736 .real => .WALL_NOW,
4737 else => .NOW,
4738 }, std.math.lossyCast(i64, duration.raw.toNanoseconds())),
4739 .deadline => |deadline| switch (deadline.clock) {
4740 .real => .walltime(&.{
4741 .sec = @intCast(@divFloor(deadline.raw.toNanoseconds(), std.time.ns_per_s)),
4742 .nsec = @intCast(@mod(deadline.raw.toNanoseconds(), std.time.ns_per_s)),
4743 }, 0),
4744 else => continue :timeout .{ .duration = deadline.durationFromNow(ev.io()) },
4745 },
4746 };
4747}
4748
4749const Random = struct {
4750 evented: *Evented,
4751 thread: *Thread,
4752 buffer: []u8,
4753
4754 fn seed(context: ?*anyopaque) callconv(.c) void {
4755 const rand: *Random = @ptrCast(@alignCast(context));
4756 const ev = rand.evented;
4757 ev.csprng_mutex.lockUncancelable(ev);
4758 defer ev.csprng_mutex.unlock();
4759 var buffer: [Csprng.seed_len]u8 = undefined;
4760 if (!ev.csprng.isInitialized()) {
4761 @branchHint(.unlikely);
4762 const cancel_protection = swapCancelProtection(ev, .blocked);
4763 defer assert(swapCancelProtection(ev, cancel_protection) == .blocked);
4764 randomSecure(ev, &buffer) catch |err| switch (err) {
4765 error.Canceled => unreachable, // blocked
4766 error.EntropyUnavailable => fallbackSeed(ev, &buffer),
4767 };
4768 ev.csprng.rng = .init(buffer);
4769 }
4770 ev.csprng.rng.fill(&buffer);
4771 rand.thread.csprng.rng = .init(buffer);
4772 rand.thread.csprng.rng.fill(rand.buffer);
4773 rand.buffer.len = 0;
4774 }
4775};
4776
4777fn random(userdata: ?*anyopaque, buffer: []u8) void {
4778 const ev: *Evented = @ptrCast(@alignCast(userdata));
4779 if (buffer.len == 0) return;
4780 const thread: *Thread = .current();
4781 var rand: Random = .{ .evented = ev, .thread = thread, .buffer = buffer };
4782 thread.seed_csprng.once(&rand, &Random.seed);
4783 if (rand.buffer.len > 0) thread.csprng.rng.fill(buffer);
4784}
4785
4786fn randomSecure(userdata: ?*anyopaque, buffer: []u8) Io.RandomSecureError!void {
4787 const ev: *Evented = @ptrCast(@alignCast(userdata));
4788 _ = ev;
4789 if (buffer.len > 0) c.arc4random_buf(buffer.ptr, buffer.len);
4790}
4791
4792fn netListenIpUnavailable(
4793 userdata: ?*anyopaque,
4794 address: *const net.IpAddress,
4795 options: net.IpAddress.ListenOptions,
4796) net.IpAddress.ListenError!net.Socket {
4797 const ev: *Evented = @ptrCast(@alignCast(userdata));
4798 _ = ev;
4799 _ = address;
4800 _ = options;
4801 return error.NetworkDown;
4802}
4803
4804fn netAcceptUnavailable(
4805 userdata: ?*anyopaque,
4806 listen_handle: net.Socket.Handle,
4807 options: net.Server.AcceptOptions,
4808) net.Server.AcceptError!net.Socket {
4809 const ev: *Evented = @ptrCast(@alignCast(userdata));
4810 _ = ev;
4811 _ = listen_handle;
4812 _ = options;
4813 return error.NetworkDown;
4814}
4815
4816fn netBindIpUnavailable(
4817 userdata: ?*anyopaque,
4818 address: *const net.IpAddress,
4819 options: net.IpAddress.BindOptions,
4820) net.IpAddress.BindError!net.Socket {
4821 const ev: *Evented = @ptrCast(@alignCast(userdata));
4822 _ = ev;
4823 _ = address;
4824 _ = options;
4825 return error.NetworkDown;
4826}
4827
4828fn netConnectIpUnavailable(
4829 userdata: ?*anyopaque,
4830 address: *const net.IpAddress,
4831 options: net.IpAddress.ConnectOptions,
4832) net.IpAddress.ConnectError!net.Socket {
4833 const ev: *Evented = @ptrCast(@alignCast(userdata));
4834 _ = ev;
4835 _ = address;
4836 _ = options;
4837 return error.NetworkDown;
4838}
4839
4840fn netListenUnixUnavailable(
4841 userdata: ?*anyopaque,
4842 address: *const net.UnixAddress,
4843 options: net.UnixAddress.ListenOptions,
4844) net.UnixAddress.ListenError!net.Socket.Handle {
4845 const ev: *Evented = @ptrCast(@alignCast(userdata));
4846 _ = ev;
4847 _ = address;
4848 _ = options;
4849 return error.AddressFamilyUnsupported;
4850}
4851
4852fn netConnectUnixUnavailable(
4853 userdata: ?*anyopaque,
4854 address: *const net.UnixAddress,
4855) net.UnixAddress.ConnectError!net.Socket.Handle {
4856 const ev: *Evented = @ptrCast(@alignCast(userdata));
4857 _ = ev;
4858 _ = address;
4859 return error.AddressFamilyUnsupported;
4860}
4861
4862fn netSocketCreatePairUnavailable(
4863 userdata: ?*anyopaque,
4864 options: net.Socket.CreatePairOptions,
4865) net.Socket.CreatePairError![2]net.Socket {
4866 _ = userdata;
4867 _ = options;
4868 return error.OperationUnsupported;
4869}
4870
4871fn netWriteFileUnavailable(
4872 userdata: ?*anyopaque,
4873 socket_handle: net.Socket.Handle,
4874 header: []const u8,
4875 file_reader: *File.Reader,
4876 limit: Io.Limit,
4877) net.Stream.Writer.WriteFileError!usize {
4878 const ev: *Evented = @ptrCast(@alignCast(userdata));
4879 _ = ev;
4880 _ = socket_handle;
4881 _ = header;
4882 _ = file_reader;
4883 _ = limit;
4884 return error.Unimplemented;
4885}
4886
4887fn netClose(userdata: ?*anyopaque, sockets: []const net.Socket) void {
4888 const ev: *Evented = @ptrCast(@alignCast(userdata));
4889 _ = ev;
4890 for (sockets) |socket| closeFd(socket.handle);
4891}
4892
4893fn netShutdownUnavailable(
4894 userdata: ?*anyopaque,
4895 handle: net.Socket.Handle,
4896 how: net.ShutdownHow,
4897) net.ShutdownError!void {
4898 const ev: *Evented = @ptrCast(@alignCast(userdata));
4899 _ = ev;
4900 _ = handle;
4901 _ = how;
4902 unreachable; // How you gonna shutdown something that was impossible to open?
4903}
4904
4905fn netInterfaceNameResolveUnavailable(
4906 userdata: ?*anyopaque,
4907 name: *const net.Interface.Name,
4908) net.Interface.Name.ResolveError!net.Interface {
4909 const ev: *Evented = @ptrCast(@alignCast(userdata));
4910 _ = ev;
4911 _ = name;
4912 return error.InterfaceNotFound;
4913}
4914
4915fn netInterfaceNameUnavailable(
4916 userdata: ?*anyopaque,
4917 interface: net.Interface,
4918) net.Interface.NameError!net.Interface.Name {
4919 const ev: *Evented = @ptrCast(@alignCast(userdata));
4920 _ = ev;
4921 _ = interface;
4922 return error.Unexpected;
4923}
4924
4925fn netLookupUnavailable(
4926 userdata: ?*anyopaque,
4927 host_name: net.HostName,
4928 resolved: *Io.Queue(net.HostName.LookupResult),
4929 options: net.HostName.LookupOptions,
4930) net.HostName.LookupError!void {
4931 const ev: *Evented = @ptrCast(@alignCast(userdata));
4932 _ = host_name;
4933 _ = options;
4934 resolved.close(ev.io());
4935 return error.NetworkDown;
4936}
4937
4938fn readAll(ev: *Evented, file: File, buffer: []u8) File.ReadStreamingError!void {
4939 var index: usize = 0;
4940 while (buffer.len - index != 0) {
4941 const len = try ev.fileReadStreaming(file, &.{buffer[index..]});
4942 if (len == 0) return error.EndOfStream;
4943 index += len;
4944 }
4945}
4946
4947fn writeAll(ev: *Evented, file: File, buffer: []const u8) (File.Writer.Error || error{EndOfStream})!void {
4948 var index: usize = 0;
4949 while (buffer.len - index != 0) {
4950 const len = try ev.fileWriteStreaming(file, &.{}, &.{buffer[index..]}, 1);
4951 if (len == 0) return error.EndOfStream;
4952 index += len;
4953 }
4954}
4955
4956/// This is either usize or u32. Since, either is fine, let's use the same
4957/// `addBuf` function for both writing to a file and sending network messages.
4958const iovlen_t = @FieldType(c.msghdr_const, "iovlen");
4959
4960fn addConstBuf(v: []iovec_const, i: *iovlen_t, remaining: ?*usize, bytes: []const u8) void {
4961 if (v.len - i.* == 0) return;
4962 const len = @min(remaining.*, bytes.len);
4963 if (len == 0) return;
4964 v[i.*] = .{ .base = bytes.ptr, .len = len };
4965 i.* += 1;
4966 remaining.* -= len;
4967}
4968fn addBuf(
4969 comptime is_const: bool,
4970 vec: []if (is_const) iovec_const else iovec,
4971 vec_len: *iovlen_t,
4972 remaining: *Io.Limit,
4973 bytes: if (is_const) []const u8 else []u8,
4974) void {
4975 if (vec.len - vec_len.* == 0) return;
4976 const len = remaining.minInt(bytes.len);
4977 if (len == 0) return;
4978 vec[vec_len.*] = .{ .base = bytes.ptr, .len = len };
4979 vec_len.* += 1;
4980 remaining.* = remaining.subtract(len).?;
4981}
4982
4983test {
4984 _ = Fiber.CancelProtection;
4985}