authorgravatar for andrew@ziglang.orgAndrew Kelley <andrew@ziglang.org> 2025-10-23 06:02:46-07:00
committergravatar for andrew@ziglang.orgAndrew Kelley <andrew@ziglang.org> 2025-10-29 06:20:51-07:00
log5578c760a77bd43ce13c9352f68f7e44c5440c8f
treec7ea4e56a539dc7c07e005b43951d240c2453320
parent6a64c9b7c8971486a818d8cb2ae44bb4dab4497f

std.Io.Kqueue: implement wait queue per fd

Solves the issue when one kevent() call would clobber another if they used the same file descriptor as an identifier.

1 files changed, 52 insertions(+), 21 deletions(-)

lib/std/Io/Kqueue.zig+52-21
......@@ -423,17 +423,35 @@ fn idle(k: *Kqueue, thread: *Thread) void {
423423 return;
424424 },
425425 _ => {
426 const fiber: *Fiber = @ptrFromInt(event.udata);
427 assert(fiber.queue_next == null);
428 fiber.resultPointer(Completion).* = .{
426 const event_head_fiber: *Fiber = @ptrFromInt(event.udata);
427 const event_tail_fiber = thread.wait_queues.fetchSwapRemove(.{
428 .ident = event.ident,
429 .filter = event.filter,
430 }).?.value;
431 assert(event_tail_fiber.queue_next == null);
432
433 // TODO reevaluate this logic
434 event_head_fiber.resultPointer(Completion).* = .{
429435 .flags = event.flags,
430436 .fflags = event.fflags,
431437 .data = event.data,
432438 };
433 if (maybe_ready_fiber == null) maybe_ready_fiber = fiber else if (maybe_ready_queue) |*ready_queue| {
434 ready_queue.tail.queue_next = fiber;
435 ready_queue.tail = fiber;
436 } else maybe_ready_queue = .{ .head = fiber, .tail = fiber };
439
440 queue_ready: {
441 const head: *Fiber = if (maybe_ready_fiber == null) f: {
442 maybe_ready_fiber = event_head_fiber;
443 const next = event_head_fiber.queue_next orelse break :queue_ready;
444 event_head_fiber.queue_next = null;
445 break :f next;
446 } else event_head_fiber;
447
448 if (maybe_ready_queue) |*ready_queue| {
449 ready_queue.tail.queue_next = head;
450 ready_queue.tail = event_tail_fiber;
451 } else {
452 maybe_ready_queue = .{ .head = head, .tail = event_tail_fiber };
453 }
454 }
437455 },
438456 };
439457 if (maybe_ready_queue) |ready_queue| k.schedule(thread, ready_queue);
......@@ -1477,7 +1495,6 @@ fn netRead(userdata: ?*anyopaque, fd: net.Socket.Handle, data: [][]u8) net.Strea
14771495
14781496 while (true) {
14791497 try k.checkCancel();
1480 std.debug.print("calling readv\n", .{});
14811498 const rc = posix.system.readv(fd, dest.ptr, @intCast(dest.len));
14821499 switch (posix.errno(rc)) {
14831500 .SUCCESS => return @intCast(rc),
......@@ -1486,19 +1503,33 @@ fn netRead(userdata: ?*anyopaque, fd: net.Socket.Handle, data: [][]u8) net.Strea
14861503 .AGAIN => {
14871504 const thread: *Thread = .current();
14881505 const fiber = thread.currentFiber();
1489 const changes = [_]posix.Kevent{
1490 .{
1491 .ident = @as(u32, @bitCast(fd)),
1492 .filter = std.c.EVFILT.READ,
1493 .flags = std.c.EV.ADD | std.c.EV.ONESHOT,
1494 .fflags = 0,
1495 .data = 0,
1496 .udata = @intFromPtr(fiber),
1497 },
1498 };
1499 assert(0 == (posix.kevent(thread.kq_fd, &changes, &.{}, null) catch |err| {
1500 @panic(@errorName(err)); // TODO
1501 }));
1506 const ident: u32 = @bitCast(fd);
1507 const filter = std.c.EVFILT.READ;
1508 const gop = thread.wait_queues.getOrPut(k.gpa, .{
1509 .ident = ident,
1510 .filter = filter,
1511 }) catch return error.SystemResources;
1512 if (gop.found_existing) {
1513 const tail_fiber = gop.value_ptr.*;
1514 assert(tail_fiber.queue_next == null);
1515 tail_fiber.queue_next = fiber;
1516 gop.value_ptr.* = fiber;
1517 } else {
1518 gop.value_ptr.* = fiber;
1519 const changes = [_]posix.Kevent{
1520 .{
1521 .ident = ident,
1522 .filter = filter,
1523 .flags = std.c.EV.ADD | std.c.EV.ONESHOT,
1524 .fflags = 0,
1525 .data = 0,
1526 .udata = @intFromPtr(fiber),
1527 },
1528 };
1529 assert(0 == (posix.kevent(thread.kq_fd, &changes, &.{}, null) catch |err| {
1530 @panic(@errorName(err)); // TODO
1531 }));
1532 }
15021533 yield(k, null, .nothing);
15031534 continue;
15041535 },