lib/sys/src/event/poll.zig

daab053ee43316e1809a84551d573ddd1e5bf3d2

   1 const std = @import("std");
   2 const sys = @import("../root.zig");
   3 const types = @import("types.zig");
   4 
   5 const fd = sys.fd;
   6 const net = sys.net;
   7 const descriptor_poll = sys.poll;
   8 const time = sys.time;
   9 
  10 pub const AcceptError = types.AcceptError;
  11 pub const CallbackAction = types.CallbackAction;
  12 pub const CloseError = types.CloseError;
  13 pub const Datagram = types.Datagram;
  14 pub const Options = types.Options;
  15 pub const PollError = types.PollError;
  16 pub const PollEvent = types.PollEvent;
  17 pub const ReadBuffer = types.ReadBuffer;
  18 pub const ReadError = types.ReadError;
  19 pub const RecvError = types.RecvError;
  20 pub const RunMode = types.RunMode;
  21 pub const SendError = types.SendError;
  22 pub const WriteBuffer = types.WriteBuffer;
  23 pub const WriteError = types.WriteError;
  24 
  25 const Callback = *const fn (?*anyopaque, *Loop, *Completion, Result) CallbackAction;
  26 
  27 const Result = union(enum) {
  28     async_wait: Async.WaitError!void,
  29     timer: Timer.RunError!void,
  30     tcp_accept: AcceptError!TCP,
  31     tcp_read: struct {
  32         socket: TCP,
  33         buffer: ReadBuffer,
  34         result: ReadError!usize,
  35     },
  36     tcp_write: struct {
  37         socket: TCP,
  38         buffer: WriteBuffer,
  39         result: WriteError!usize,
  40     },
  41     udp_recv_from: struct {
  42         socket: UDP,
  43         buffer: ReadBuffer,
  44         result: RecvError!Datagram,
  45     },
  46     udp_send_to: struct {
  47         socket: UDP,
  48         buffer: WriteBuffer,
  49         result: SendError!usize,
  50     },
  51     tcp_close: struct {
  52         socket: TCP,
  53         result: CloseError!void,
  54     },
  55     file_poll: struct {
  56         file: File,
  57         result: PollError!PollEvent,
  58     },
  59 };
  60 
  61 const Operation = union(enum) {
  62     none,
  63     async_wait: Async,
  64     timer: TimerState,
  65     tcp_accept: TCP,
  66     tcp_read: TcpReadState,
  67     tcp_write: TcpWriteState,
  68     udp_recv_from: UdpRecvFromState,
  69     udp_send_to: UdpSendToState,
  70     file_poll: FilePollState,
  71 
  72     fn descriptor(self: Operation) ?fd.Descriptor {
  73         return switch (self) {
  74             .none, .timer => null,
  75             .async_wait => |async_watcher| async_watcher.read_fd,
  76             .tcp_accept => |socket| socket.fd,
  77             .tcp_read => |state| state.socket.fd,
  78             .tcp_write => |state| state.socket.fd,
  79             .udp_recv_from => |state| state.socket.fd,
  80             .udp_send_to => |state| state.socket.fd,
  81             .file_poll => |state| state.file.fd,
  82         };
  83     }
  84 
  85     fn events(self: Operation) i16 {
  86         return switch (self) {
  87             .tcp_write, .udp_send_to => descriptor_poll.Event.output,
  88             else => descriptor_poll.Event.input,
  89         };
  90     }
  91 };
  92 
  93 const TimerState = struct {
  94     deadline: time.AwakeInstant,
  95     repeat: ?Timer.Repeat,
  96 };
  97 
  98 const TcpReadState = struct {
  99     socket: TCP,
 100     buffer: ReadBuffer,
 101 };
 102 
 103 const TcpWriteState = struct {
 104     socket: TCP,
 105     buffer: WriteBuffer,
 106 };
 107 
 108 const UdpRecvFromState = struct {
 109     socket: UDP,
 110     buffer: ReadBuffer,
 111 };
 112 
 113 const UdpSendToState = struct {
 114     socket: UDP,
 115     buffer: WriteBuffer,
 116     address: net.IpAddress,
 117 };
 118 
 119 const FilePollState = struct {
 120     file: File,
 121     event: PollEvent,
 122 };
 123 
 124 pub const Completion = struct {
 125     active: bool = false,
 126     op: Operation = .none,
 127     userdata: ?*anyopaque = null,
 128     callback: ?Callback = null,
 129 };
 130 
 131 pub const Loop = struct {
 132     allocator: std.mem.Allocator,
 133     awake_clock: time.AwakeClock,
 134     last_instant: ?time.AwakeInstant,
 135     entries: []?*Completion,
 136     pollfds: []descriptor_poll.Descriptor,
 137     poll_completions: []*Completion,
 138 
 139     pub fn init(options: Options) !Loop {
 140         const allocator = options.allocator;
 141         const count: usize = @intCast(options.entries);
 142         const entries = try allocator.alloc(?*Completion, count);
 143         errdefer allocator.free(entries);
 144         @memset(entries, null);
 145 
 146         const pollfds = try allocator.alloc(descriptor_poll.Descriptor, count);
 147         errdefer allocator.free(pollfds);
 148 
 149         const poll_completions = try allocator.alloc(*Completion, count);
 150         errdefer allocator.free(poll_completions);
 151 
 152         return .{
 153             .allocator = allocator,
 154             .awake_clock = options.awake_clock,
 155             .last_instant = null,
 156             .entries = entries,
 157             .pollfds = pollfds,
 158             .poll_completions = poll_completions,
 159         };
 160     }
 161 
 162     pub fn deinit(self: *Loop) void {
 163         self.allocator.free(self.poll_completions);
 164         self.allocator.free(self.pollfds);
 165         self.allocator.free(self.entries);
 166         self.* = undefined;
 167     }
 168 
 169     pub fn run(self: *Loop, mode: RunMode) !void {
 170         while (self.hasActive()) {
 171             const now = try self.readInstant();
 172             if (try self.fireDueTimer(now)) {
 173                 if (mode != .until_done) return;
 174                 continue;
 175             }
 176 
 177             const poll_count = self.preparePoll();
 178             const timeout_ms = self.timeoutMillis(mode, now);
 179             if (poll_count == 0) {
 180                 if (timeout_ms < 0) return;
 181                 if (timeout_ms > 0) {
 182                     const timeout_ns = @as(u64, @intCast(timeout_ms)) *
 183                         std.time.ns_per_ms;
 184                     std.Io.sleep(
 185                         std.Options.debug_io,
 186                         .fromNanoseconds(timeout_ns),
 187                         .awake,
 188                     ) catch {};
 189                 }
 190                 if (mode != .until_done) return;
 191                 continue;
 192             }
 193 
 194             const ready = descriptor_poll.wait(self.pollfds[0..poll_count], timeout_ms) catch
 195                 return error.Unexpected;
 196             if (ready == 0) {
 197                 if (mode != .until_done) return;
 198                 continue;
 199             }
 200 
 201             _ = try self.readInstant();
 202             if (try self.fireFirstReady(poll_count) and mode != .until_done) return;
 203             if (mode == .no_wait) return;
 204         }
 205     }
 206 
 207     pub fn currentInstant(self: *const Loop) time.AwakeInstant {
 208         return self.last_instant orelse @panic("sys.event loop has not sampled its awake clock");
 209     }
 210 
 211     pub fn awakeClock(self: *const Loop) time.AwakeClock {
 212         return self.awake_clock;
 213     }
 214 
 215     pub fn disarm(self: *Loop, completion: *Completion) void {
 216         self.deactivate(completion);
 217         completion.* = .{};
 218     }
 219 
 220     fn activate(self: *Loop, completion: *Completion) void {
 221         if (completion.active) {
 222             for (self.entries) |slot| {
 223                 if (slot == completion) return;
 224             }
 225             @panic("sys.event active completion belongs to another loop");
 226         }
 227         for (self.entries) |*slot| {
 228             if (slot.* == null) {
 229                 slot.* = completion;
 230                 completion.active = true;
 231                 return;
 232             }
 233         }
 234         @panic("sys.event loop completion capacity exhausted");
 235     }
 236 
 237     fn arm(
 238         self: *Loop,
 239         completion: *Completion,
 240         op: Operation,
 241         userdata: ?*anyopaque,
 242         callback: Callback,
 243     ) void {
 244         var replacement: Completion = .{
 245             .op = op,
 246             .userdata = userdata,
 247             .callback = callback,
 248         };
 249         if (completion.active) {
 250             for (self.entries) |slot| {
 251                 if (slot == completion) {
 252                     replacement.active = true;
 253                     completion.* = replacement;
 254                     return;
 255                 }
 256             }
 257             @panic("sys.event active completion belongs to another loop");
 258         }
 259         completion.* = replacement;
 260         self.activate(completion);
 261     }
 262 
 263     fn deactivate(self: *Loop, completion: *Completion) void {
 264         if (!completion.active) return;
 265         for (self.entries) |*slot| {
 266             if (slot.* == completion) {
 267                 slot.* = null;
 268                 completion.active = false;
 269                 return;
 270             }
 271         }
 272         @panic("sys.event active completion belongs to another loop");
 273     }
 274 
 275     fn deactivateDescriptor(self: *Loop, descriptor: fd.Descriptor) void {
 276         for (self.entries) |*slot| {
 277             const completion = slot.* orelse continue;
 278             if (completion.op.descriptor()) |active_descriptor| {
 279                 if (active_descriptor == descriptor) {
 280                     completion.active = false;
 281                     slot.* = null;
 282                 }
 283             }
 284         }
 285     }
 286 
 287     fn hasActive(self: *Loop) bool {
 288         for (self.entries) |slot| {
 289             if (slot != null) return true;
 290         }
 291         return false;
 292     }
 293 
 294     fn preparePoll(self: *Loop) usize {
 295         var count: usize = 0;
 296         for (self.entries) |slot| {
 297             const completion = slot orelse continue;
 298             const descriptor = completion.op.descriptor() orelse continue;
 299             self.pollfds[count] = .{
 300                 .fd = descriptor,
 301                 .events = completion.op.events(),
 302                 .revents = 0,
 303             };
 304             self.poll_completions[count] = completion;
 305             count += 1;
 306         }
 307         return count;
 308     }
 309 
 310     fn timeoutMillis(self: *Loop, mode: RunMode, now: time.AwakeInstant) i32 {
 311         if (mode == .no_wait) return 0;
 312         var best: ?u64 = null;
 313         for (self.entries) |slot| {
 314             const completion = slot orelse continue;
 315             switch (completion.op) {
 316                 .timer => |timer| {
 317                     const remaining = now.remainingUntil(timer.deadline).asMillisecondsCeil();
 318                     best = if (best) |current| @min(current, remaining) else remaining;
 319                 },
 320                 else => {},
 321             }
 322         }
 323         const value = best orelse return -1;
 324         return @intCast(@min(value, std.math.maxInt(i32)));
 325     }
 326 
 327     fn fireDueTimer(self: *Loop, now: time.AwakeInstant) time.ClockError!bool {
 328         var due: ?*Completion = null;
 329         for (self.entries) |slot| {
 330             const completion = slot orelse continue;
 331             switch (completion.op) {
 332                 .timer => |timer| {
 333                     if (!now.reached(timer.deadline)) continue;
 334                     const current = due orelse {
 335                         due = completion;
 336                         continue;
 337                     };
 338                     const current_timer = switch (current.op) {
 339                         .timer => |value| value,
 340                         else => unreachable,
 341                     };
 342                     if (timer.deadline.isBefore(current_timer.deadline)) due = completion;
 343                 },
 344                 else => {},
 345             }
 346         }
 347         const completion = due orelse return false;
 348         try self.fire(completion, .{ .timer = {} });
 349         return true;
 350     }
 351 
 352     fn fireFirstReady(self: *Loop, count: usize) time.ClockError!bool {
 353         var index: usize = 0;
 354         while (index < count) : (index += 1) {
 355             if (self.pollfds[index].revents == 0) continue;
 356             const completion = self.poll_completions[index];
 357             if (!completion.active) continue;
 358             try self.fireReady(completion);
 359             return true;
 360         }
 361         return false;
 362     }
 363 
 364     fn fireReady(self: *Loop, completion: *Completion) time.ClockError!void {
 365         switch (completion.op) {
 366             .async_wait => |async_watcher| {
 367                 async_watcher.drain();
 368                 try self.fire(completion, .{ .async_wait = {} });
 369             },
 370             .tcp_accept => |socket| {
 371                 const result = acceptSocket(socket);
 372                 try self.fire(completion, .{ .tcp_accept = result });
 373             },
 374             .tcp_read => |*state| {
 375                 const result = readSocket(state);
 376                 try self.fire(completion, .{ .tcp_read = .{
 377                     .socket = state.socket,
 378                     .buffer = state.buffer,
 379                     .result = result,
 380                 } });
 381             },
 382             .tcp_write => |*state| {
 383                 const result = writeSocket(state);
 384                 try self.fire(completion, .{ .tcp_write = .{
 385                     .socket = state.socket,
 386                     .buffer = state.buffer,
 387                     .result = result,
 388                 } });
 389             },
 390             .udp_recv_from => |*state| try self.fireUdpRecvFrom(completion, state),
 391             .udp_send_to => |*state| try self.fireUdpSendTo(completion, state),
 392             .file_poll => |state| {
 393                 try self.fire(completion, .{ .file_poll = .{
 394                     .file = state.file,
 395                     .result = state.event,
 396                 } });
 397             },
 398             .none, .timer => {},
 399         }
 400     }
 401 
 402     fn fireUdpRecvFrom(
 403         self: *Loop,
 404         completion: *Completion,
 405         state: *UdpRecvFromState,
 406     ) time.ClockError!void {
 407         const received = net.recvFromIpAddress(
 408             state.socket.fd,
 409             readSlice(&state.buffer),
 410             0,
 411         ) catch |err| {
 412             if (err == error.WouldBlock) return;
 413             const result: RecvError = switch (err) {
 414                 error.ConnectionResetByPeer => error.ConnectionResetByPeer,
 415                 else => error.Unexpected,
 416             };
 417             try self.fire(completion, .{ .udp_recv_from = .{
 418                 .socket = state.socket,
 419                 .buffer = state.buffer,
 420                 .result = result,
 421             } });
 422             return;
 423         };
 424         try self.fire(completion, .{ .udp_recv_from = .{
 425             .socket = state.socket,
 426             .buffer = state.buffer,
 427             .result = Datagram{
 428                 .address = received.address,
 429                 .bytes = received.bytes,
 430                 .truncated = received.truncated,
 431             },
 432         } });
 433     }
 434 
 435     fn fireUdpSendTo(
 436         self: *Loop,
 437         completion: *Completion,
 438         state: *const UdpSendToState,
 439     ) time.ClockError!void {
 440         const bytes = writeSlice(&state.buffer);
 441         const sent = net.sendToIpAddress(state.socket.fd, bytes, 0, state.address) catch |err| {
 442             if (err == error.WouldBlock) return;
 443             const result: SendError = switch (err) {
 444                 error.BrokenPipe => error.BrokenPipe,
 445                 error.ConnectionResetByPeer => error.ConnectionResetByPeer,
 446                 error.MessageTooLarge => error.MessageTooLarge,
 447                 else => error.Unexpected,
 448             };
 449             try self.fire(completion, .{ .udp_send_to = .{
 450                 .socket = state.socket,
 451                 .buffer = state.buffer,
 452                 .result = result,
 453             } });
 454             return;
 455         };
 456         std.debug.assert(sent == bytes.len);
 457         try self.fire(completion, .{ .udp_send_to = .{
 458             .socket = state.socket,
 459             .buffer = state.buffer,
 460             .result = sent,
 461         } });
 462     }
 463 
 464     fn fire(self: *Loop, completion: *Completion, result: Result) time.ClockError!void {
 465         const callback = completion.callback orelse return;
 466         self.deactivate(completion);
 467         const action = callback(completion.userdata, self, completion, result);
 468         if (action == .rearm) {
 469             switch (completion.op) {
 470                 .timer => |*timer| timer.deadline = try self.rearmDeadline(timer.*),
 471                 else => {},
 472             }
 473             self.activate(completion);
 474         }
 475     }
 476 
 477     fn readInstant(self: *Loop) time.ClockError!time.AwakeInstant {
 478         const current = try self.awake_clock.now();
 479         if (self.last_instant) |previous| _ = try current.elapsedSince(previous);
 480         self.last_instant = current;
 481         return current;
 482     }
 483 
 484     fn rearmDeadline(self: *Loop, timer: TimerState) time.ClockError!time.AwakeInstant {
 485         const repeat = timer.repeat orelse @panic("sys.event timer rearm needs a repeat policy");
 486         const now = try self.readInstant();
 487         const interval = switch (repeat) {
 488             .fixed_delay, .fixed_rate => |value| value,
 489         };
 490         const deadline = switch (repeat) {
 491             .fixed_delay => now.deadlineAfter(interval),
 492             .fixed_rate => fixedRateDeadline(timer.deadline, now, interval),
 493         };
 494         if (!interval.isZero() and now.reached(deadline)) return error.ClockOverflow;
 495         return deadline;
 496     }
 497 };
 498 
 499 pub const Async = struct {
 500     read_fd: fd.Descriptor,
 501     write_fd: fd.Descriptor,
 502 
 503     pub const WaitError: type = ReadError;
 504 
 505     pub fn init() !Async {
 506         const pipe_fds = try fd.pipeWithOptions(.{ .close_on_exec = true, .nonblocking = true });
 507         return .{ .read_fd = pipe_fds[0], .write_fd = pipe_fds[1] };
 508     }
 509 
 510     pub fn deinit(self: *Async) void {
 511         fd.close(self.read_fd);
 512         fd.close(self.write_fd);
 513         self.* = undefined;
 514     }
 515 
 516     pub fn wait(
 517         self: Async,
 518         loop: *Loop,
 519         completion: *Completion,
 520         comptime Userdata: type,
 521         userdata: ?*Userdata,
 522         comptime cb: *const fn (
 523             ud: ?*Userdata,
 524             loop: *Loop,
 525             completion: *Completion,
 526             result: WaitError!void,
 527         ) CallbackAction,
 528     ) void {
 529         loop.arm(
 530             completion,
 531             .{ .async_wait = self },
 532             userdata,
 533             AsyncWaitCallback(Userdata, cb).callback,
 534         );
 535     }
 536 
 537     pub fn notify(self: Async) !void {
 538         _ = fd.write(self.write_fd, &.{1}) catch |err| switch (err) {
 539             error.WouldBlock => return,
 540             else => return err,
 541         };
 542     }
 543 
 544     fn drain(self: Async) void {
 545         var buffer: [64]u8 = undefined;
 546         while (true) {
 547             _ = fd.read(self.read_fd, &buffer) catch |err| switch (err) {
 548                 error.WouldBlock => return,
 549                 else => return,
 550             };
 551         }
 552     }
 553 };
 554 
 555 pub const Timer = struct {
 556     pub const delayed_tick_slots_max: u8 = 1;
 557 
 558     pub const RunError = error{
 559         Canceled,
 560         Unexpected,
 561     };
 562 
 563     pub const Repeat = union(enum) {
 564         fixed_delay: time.Duration,
 565         fixed_rate: time.Duration,
 566     };
 567 
 568     pub const Schedule = struct {
 569         after: time.Duration,
 570         repeat: ?Repeat = null,
 571     };
 572 
 573     pub fn init() !Timer {
 574         return .{};
 575     }
 576 
 577     pub fn deinit(self: *Timer) void {
 578         self.* = undefined;
 579     }
 580 
 581     pub fn run(
 582         self: Timer,
 583         loop: *Loop,
 584         completion: *Completion,
 585         schedule: Schedule,
 586         comptime Userdata: type,
 587         userdata: ?*Userdata,
 588         comptime cb: *const fn (
 589             ud: ?*Userdata,
 590             loop: *Loop,
 591             completion: *Completion,
 592             result: RunError!void,
 593         ) CallbackAction,
 594     ) time.ClockError!void {
 595         _ = self;
 596         if (schedule.repeat) |repeat| switch (repeat) {
 597             .fixed_delay => {},
 598             .fixed_rate => |interval| std.debug.assert(!interval.isZero()),
 599         };
 600         const now = try loop.readInstant();
 601         loop.arm(
 602             completion,
 603             .{ .timer = .{
 604                 .deadline = now.deadlineAfter(schedule.after),
 605                 .repeat = schedule.repeat,
 606             } },
 607             userdata,
 608             TimerRunCallback(Userdata, cb).callback,
 609         );
 610     }
 611 };
 612 
 613 pub const TCP = struct {
 614     fd: net.Socket,
 615 
 616     pub fn initFd(socket_fd: anytype) TCP {
 617         return .{ .fd = @intCast(socket_fd) };
 618     }
 619 
 620     pub fn accept(
 621         self: TCP,
 622         loop: *Loop,
 623         completion: *Completion,
 624         comptime Userdata: type,
 625         userdata: ?*Userdata,
 626         comptime cb: *const fn (
 627             ud: ?*Userdata,
 628             loop: *Loop,
 629             completion: *Completion,
 630             result: AcceptError!TCP,
 631         ) CallbackAction,
 632     ) void {
 633         loop.arm(
 634             completion,
 635             .{ .tcp_accept = self },
 636             userdata,
 637             TcpAcceptCallback(Userdata, cb).callback,
 638         );
 639     }
 640 
 641     pub fn read(
 642         self: TCP,
 643         loop: *Loop,
 644         completion: *Completion,
 645         buffer: ReadBuffer,
 646         comptime Userdata: type,
 647         userdata: ?*Userdata,
 648         comptime cb: *const fn (
 649             ud: ?*Userdata,
 650             loop: *Loop,
 651             completion: *Completion,
 652             socket: TCP,
 653             buffer: ReadBuffer,
 654             result: ReadError!usize,
 655         ) CallbackAction,
 656     ) void {
 657         loop.arm(
 658             completion,
 659             .{ .tcp_read = .{ .socket = self, .buffer = buffer } },
 660             userdata,
 661             TcpReadCallback(Userdata, cb).callback,
 662         );
 663     }
 664 
 665     pub fn write(
 666         self: TCP,
 667         loop: *Loop,
 668         completion: *Completion,
 669         buffer: WriteBuffer,
 670         comptime Userdata: type,
 671         userdata: ?*Userdata,
 672         comptime cb: *const fn (
 673             ud: ?*Userdata,
 674             loop: *Loop,
 675             completion: *Completion,
 676             socket: TCP,
 677             buffer: WriteBuffer,
 678             result: WriteError!usize,
 679         ) CallbackAction,
 680     ) void {
 681         loop.arm(
 682             completion,
 683             .{ .tcp_write = .{ .socket = self, .buffer = buffer } },
 684             userdata,
 685             TcpWriteCallback(Userdata, cb).callback,
 686         );
 687     }
 688 
 689     pub fn close(
 690         self: TCP,
 691         loop: *Loop,
 692         completion: *Completion,
 693         comptime Userdata: type,
 694         userdata: ?*Userdata,
 695         comptime cb: *const fn (
 696             ud: ?*Userdata,
 697             loop: *Loop,
 698             completion: *Completion,
 699             socket: TCP,
 700             result: CloseError!void,
 701         ) CallbackAction,
 702     ) void {
 703         loop.deactivateDescriptor(self.fd);
 704         loop.deactivate(completion);
 705         net.close(self.fd);
 706         completion.* = .{
 707             .userdata = userdata,
 708             .callback = TcpCloseCallback(Userdata, cb).callback,
 709         };
 710         loop.fire(
 711             completion,
 712             .{ .tcp_close = .{ .socket = self, .result = {} } },
 713         ) catch unreachable;
 714     }
 715 };
 716 
 717 pub const UDP = struct {
 718     fd: net.Socket,
 719 
 720     pub const BindError = net.SocketError || net.BindError;
 721 
 722     pub fn initFd(socket_fd: anytype) UDP {
 723         return .{ .fd = @intCast(socket_fd) };
 724     }
 725 
 726     pub fn bind(address: net.IpAddress) BindError!UDP {
 727         const socket_fd = try net.udpDatagramSocketForAddress(address, .{
 728             .close_on_exec = true,
 729             .nonblocking = true,
 730         });
 731         errdefer net.close(socket_fd);
 732         try net.bindIpAddress(socket_fd, address);
 733         return initFd(socket_fd);
 734     }
 735 
 736     pub fn recvFrom(
 737         self: UDP,
 738         loop: *Loop,
 739         completion: *Completion,
 740         buffer: ReadBuffer,
 741         comptime Userdata: type,
 742         userdata: ?*Userdata,
 743         comptime cb: *const fn (
 744             ud: ?*Userdata,
 745             loop: *Loop,
 746             completion: *Completion,
 747             socket: UDP,
 748             buffer: ReadBuffer,
 749             result: RecvError!Datagram,
 750         ) CallbackAction,
 751     ) void {
 752         std.debug.assert(readBufferLength(buffer) > 0);
 753         loop.arm(
 754             completion,
 755             .{ .udp_recv_from = .{ .socket = self, .buffer = buffer } },
 756             userdata,
 757             UdpRecvFromCallback(Userdata, cb).callback,
 758         );
 759     }
 760 
 761     pub fn sendTo(
 762         self: UDP,
 763         loop: *Loop,
 764         completion: *Completion,
 765         buffer: WriteBuffer,
 766         address: net.IpAddress,
 767         comptime Userdata: type,
 768         userdata: ?*Userdata,
 769         comptime cb: *const fn (
 770             ud: ?*Userdata,
 771             loop: *Loop,
 772             completion: *Completion,
 773             socket: UDP,
 774             buffer: WriteBuffer,
 775             result: SendError!usize,
 776         ) CallbackAction,
 777     ) void {
 778         loop.arm(
 779             completion,
 780             .{ .udp_send_to = .{
 781                 .socket = self,
 782                 .buffer = buffer,
 783                 .address = address,
 784             } },
 785             userdata,
 786             UdpSendToCallback(Userdata, cb).callback,
 787         );
 788     }
 789 
 790     pub fn close(self: UDP) void {
 791         net.close(self.fd);
 792     }
 793 };
 794 
 795 pub const File = struct {
 796     fd: fd.Descriptor,
 797 
 798     pub fn initFd(file_fd: std.Io.File.Handle) File {
 799         return .{ .fd = file_fd };
 800     }
 801 
 802     pub fn poll(
 803         self: File,
 804         loop: *Loop,
 805         completion: *Completion,
 806         event: PollEvent,
 807         comptime Userdata: type,
 808         userdata: ?*Userdata,
 809         comptime cb: *const fn (
 810             ud: ?*Userdata,
 811             loop: *Loop,
 812             completion: *Completion,
 813             file: File,
 814             result: PollError!PollEvent,
 815         ) CallbackAction,
 816     ) void {
 817         loop.arm(
 818             completion,
 819             .{ .file_poll = .{ .file = self, .event = event } },
 820             userdata,
 821             FilePollCallback(Userdata, cb).callback,
 822         );
 823     }
 824 };
 825 
 826 fn acceptSocket(socket: TCP) AcceptError!TCP {
 827     const accepted = net.acceptNonBlocking(socket.fd) catch |err| switch (err) {
 828         error.WouldBlock => return error.Again,
 829         else => return error.Unexpected,
 830     };
 831     return TCP.initFd(accepted);
 832 }
 833 
 834 fn readSocket(state: *TcpReadState) ReadError!usize {
 835     const buffer = readSlice(&state.buffer);
 836     if (buffer.len == 0) return 0;
 837     const count = net.recv(state.socket.fd, buffer, 0) catch |err| switch (err) {
 838         error.ConnectionResetByPeer => return error.ConnectionResetByPeer,
 839         else => return error.Unexpected,
 840     };
 841     if (count == 0) return error.EOF;
 842     return count;
 843 }
 844 
 845 fn writeSocket(state: *const TcpWriteState) WriteError!usize {
 846     return net.sendNoSignal(state.socket.fd, writeSlice(&state.buffer)) catch |err| switch (err) {
 847         error.BrokenPipe => error.BrokenPipe,
 848         error.ConnectionResetByPeer => error.ConnectionResetByPeer,
 849         else => error.Unexpected,
 850     };
 851 }
 852 
 853 fn readSlice(buffer: *ReadBuffer) []u8 {
 854     return switch (buffer.*) {
 855         .slice => |slice| slice,
 856         .array => |*array| array,
 857     };
 858 }
 859 
 860 fn readBufferLength(buffer: ReadBuffer) usize {
 861     return switch (buffer) {
 862         .slice => |slice| slice.len,
 863         .array => |array| array.len,
 864     };
 865 }
 866 
 867 fn writeSlice(buffer: *const WriteBuffer) []const u8 {
 868     return switch (buffer.*) {
 869         .slice => |slice| slice,
 870         .array => |*array| array.array[0..array.len],
 871     };
 872 }
 873 
 874 fn fixedRateDeadline(
 875     previous: time.AwakeInstant,
 876     now: time.AwakeInstant,
 877     interval: time.Duration,
 878 ) time.AwakeInstant {
 879     std.debug.assert(!interval.isZero());
 880     const first = previous.deadlineAfter(interval);
 881     if (!now.reached(first)) return first;
 882     const late = now.elapsedSince(first) catch unreachable;
 883     const missed = late.asNanoseconds() / interval.asNanoseconds() +| 1;
 884     const advance = interval.asNanoseconds() *| missed;
 885     return first.deadlineAfter(.fromNanoseconds(advance));
 886 }
 887 
 888 fn AsyncWaitCallback(
 889     comptime Userdata: type,
 890     comptime cb: *const fn (
 891         ud: ?*Userdata,
 892         event_loop: *Loop,
 893         completion: *Completion,
 894         result: Async.WaitError!void,
 895     ) CallbackAction,
 896 ) type {
 897     return struct {
 898         fn callback(
 899             ud: ?*anyopaque,
 900             loop: *Loop,
 901             completion: *Completion,
 902             result: Result,
 903         ) CallbackAction {
 904             return cb(castUserdata(Userdata, ud), loop, completion, result.async_wait);
 905         }
 906     };
 907 }
 908 
 909 fn TimerRunCallback(
 910     comptime Userdata: type,
 911     comptime cb: *const fn (
 912         ud: ?*Userdata,
 913         event_loop: *Loop,
 914         completion: *Completion,
 915         result: Timer.RunError!void,
 916     ) CallbackAction,
 917 ) type {
 918     return struct {
 919         fn callback(
 920             ud: ?*anyopaque,
 921             loop: *Loop,
 922             completion: *Completion,
 923             result: Result,
 924         ) CallbackAction {
 925             return cb(castUserdata(Userdata, ud), loop, completion, result.timer);
 926         }
 927     };
 928 }
 929 
 930 fn TcpAcceptCallback(
 931     comptime Userdata: type,
 932     comptime cb: *const fn (
 933         ud: ?*Userdata,
 934         event_loop: *Loop,
 935         completion: *Completion,
 936         result: AcceptError!TCP,
 937     ) CallbackAction,
 938 ) type {
 939     return struct {
 940         fn callback(
 941             ud: ?*anyopaque,
 942             loop: *Loop,
 943             completion: *Completion,
 944             result: Result,
 945         ) CallbackAction {
 946             return cb(castUserdata(Userdata, ud), loop, completion, result.tcp_accept);
 947         }
 948     };
 949 }
 950 
 951 fn TcpReadCallback(
 952     comptime Userdata: type,
 953     comptime cb: *const fn (
 954         ud: ?*Userdata,
 955         event_loop: *Loop,
 956         completion: *Completion,
 957         socket: TCP,
 958         buffer: ReadBuffer,
 959         result: ReadError!usize,
 960     ) CallbackAction,
 961 ) type {
 962     return struct {
 963         fn callback(
 964             ud: ?*anyopaque,
 965             loop: *Loop,
 966             completion: *Completion,
 967             result: Result,
 968         ) CallbackAction {
 969             const read = result.tcp_read;
 970             return cb(castUserdata(Userdata, ud), loop, completion, read.socket, read.buffer, read.result);
 971         }
 972     };
 973 }
 974 
 975 fn TcpWriteCallback(
 976     comptime Userdata: type,
 977     comptime cb: *const fn (
 978         ud: ?*Userdata,
 979         event_loop: *Loop,
 980         completion: *Completion,
 981         socket: TCP,
 982         buffer: WriteBuffer,
 983         result: WriteError!usize,
 984     ) CallbackAction,
 985 ) type {
 986     return struct {
 987         fn callback(
 988             ud: ?*anyopaque,
 989             loop: *Loop,
 990             completion: *Completion,
 991             result: Result,
 992         ) CallbackAction {
 993             const write = result.tcp_write;
 994             return cb(castUserdata(Userdata, ud), loop, completion, write.socket, write.buffer, write.result);
 995         }
 996     };
 997 }
 998 
 999 fn UdpRecvFromCallback(
1000     comptime Userdata: type,
1001     comptime cb: *const fn (
1002         ud: ?*Userdata,
1003         event_loop: *Loop,
1004         completion: *Completion,
1005         socket: UDP,
1006         buffer: ReadBuffer,
1007         result: RecvError!Datagram,
1008     ) CallbackAction,
1009 ) type {
1010     return struct {
1011         fn callback(
1012             ud: ?*anyopaque,
1013             loop: *Loop,
1014             completion: *Completion,
1015             result: Result,
1016         ) CallbackAction {
1017             const received = result.udp_recv_from;
1018             return cb(
1019                 castUserdata(Userdata, ud),
1020                 loop,
1021                 completion,
1022                 received.socket,
1023                 received.buffer,
1024                 received.result,
1025             );
1026         }
1027     };
1028 }
1029 
1030 fn UdpSendToCallback(
1031     comptime Userdata: type,
1032     comptime cb: *const fn (
1033         ud: ?*Userdata,
1034         event_loop: *Loop,
1035         completion: *Completion,
1036         socket: UDP,
1037         buffer: WriteBuffer,
1038         result: SendError!usize,
1039     ) CallbackAction,
1040 ) type {
1041     return struct {
1042         fn callback(
1043             ud: ?*anyopaque,
1044             loop: *Loop,
1045             completion: *Completion,
1046             result: Result,
1047         ) CallbackAction {
1048             const sent = result.udp_send_to;
1049             return cb(
1050                 castUserdata(Userdata, ud),
1051                 loop,
1052                 completion,
1053                 sent.socket,
1054                 sent.buffer,
1055                 sent.result,
1056             );
1057         }
1058     };
1059 }
1060 
1061 fn TcpCloseCallback(
1062     comptime Userdata: type,
1063     comptime cb: *const fn (
1064         ud: ?*Userdata,
1065         event_loop: *Loop,
1066         completion: *Completion,
1067         socket: TCP,
1068         result: CloseError!void,
1069     ) CallbackAction,
1070 ) type {
1071     return struct {
1072         fn callback(
1073             ud: ?*anyopaque,
1074             loop: *Loop,
1075             completion: *Completion,
1076             result: Result,
1077         ) CallbackAction {
1078             const close = result.tcp_close;
1079             return cb(castUserdata(Userdata, ud), loop, completion, close.socket, close.result);
1080         }
1081     };
1082 }
1083 
1084 fn FilePollCallback(
1085     comptime Userdata: type,
1086     comptime cb: *const fn (
1087         ud: ?*Userdata,
1088         event_loop: *Loop,
1089         completion: *Completion,
1090         file: File,
1091         result: PollError!PollEvent,
1092     ) CallbackAction,
1093 ) type {
1094     return struct {
1095         fn callback(
1096             ud: ?*anyopaque,
1097             loop: *Loop,
1098             completion: *Completion,
1099             result: Result,
1100         ) CallbackAction {
1101             const poll = result.file_poll;
1102             return cb(castUserdata(Userdata, ud), loop, completion, poll.file, poll.result);
1103         }
1104     };
1105 }
1106 
1107 fn castUserdata(comptime Userdata: type, userdata: ?*anyopaque) ?*Userdata {
1108     if (Userdata == void) return null;
1109     return if (userdata) |ptr| @ptrCast(@alignCast(ptr)) else null;
1110 }
1111 
1112 fn timer_test_callback(
1113     state: ?*bool,
1114     _: *Loop,
1115     _: *Completion,
1116     result: Timer.RunError!void,
1117 ) CallbackAction {
1118     _ = result catch return .disarm;
1119     state.?.* = true;
1120     return .disarm;
1121 }
1122 
1123 test "event loop timer uses owned sys event types" {
1124     var loop = try Loop.init(.{ .allocator = std.testing.allocator });
1125     defer loop.deinit();
1126 
1127     var timer = try Timer.init();
1128     defer timer.deinit();
1129 
1130     var completion: Completion = .{};
1131     var fired = false;
1132     try timer.run(
1133         &loop,
1134         &completion,
1135         .{ .after = .zero },
1136         bool,
1137         &fired,
1138         timer_test_callback,
1139     );
1140 
1141     try loop.run(.until_done);
1142     try std.testing.expect(fired);
1143 }
1144 
1145 const TimerCount = struct {
1146     hits: usize = 0,
1147 };
1148 
1149 const TimerOrder = struct {
1150     id: u8,
1151     log: *[2]u8,
1152     count: *u8,
1153 };
1154 
1155 fn timer_order_callback(
1156     state: ?*TimerOrder,
1157     _: *Loop,
1158     _: *Completion,
1159     result: Timer.RunError!void,
1160 ) CallbackAction {
1161     _ = result catch return .disarm;
1162     const order = state.?;
1163     order.log[order.count.*] = order.id;
1164     order.count.* += 1;
1165     return .disarm;
1166 }
1167 
1168 const TimerRearmState = struct {
1169     clock: *time.FakeClock,
1170     late: time.Duration,
1171     hits: u8 = 0,
1172 };
1173 
1174 fn timer_rearm_callback(
1175     state: ?*TimerRearmState,
1176     _: *Loop,
1177     _: *Completion,
1178     result: Timer.RunError!void,
1179 ) CallbackAction {
1180     _ = result catch return .disarm;
1181     const rearm = state.?;
1182     rearm.hits += 1;
1183     if (rearm.hits != 1) return .disarm;
1184     rearm.clock.advance(rearm.late);
1185     return .rearm;
1186 }
1187 
1188 fn timer_count_callback(
1189     state: ?*TimerCount,
1190     _: *Loop,
1191     _: *Completion,
1192     result: Timer.RunError!void,
1193 ) CallbackAction {
1194     _ = result catch return .disarm;
1195     state.?.hits += 1;
1196     return .disarm;
1197 }
1198 
1199 test "event loop replaces an active completion without consuming another entry" {
1200     var loop = try Loop.init(.{ .allocator = std.testing.allocator, .entries = 1 });
1201     defer loop.deinit();
1202 
1203     var timer = try Timer.init();
1204     defer timer.deinit();
1205     var completion: Completion = .{};
1206     var first: TimerCount = .{};
1207     var second: TimerCount = .{};
1208     try timer.run(
1209         &loop,
1210         &completion,
1211         .{ .after = .zero },
1212         TimerCount,
1213         &first,
1214         timer_count_callback,
1215     );
1216     try timer.run(
1217         &loop,
1218         &completion,
1219         .{ .after = .zero },
1220         TimerCount,
1221         &second,
1222         timer_count_callback,
1223     );
1224 
1225     try loop.run(.until_done);
1226     try std.testing.expectEqual(@as(usize, 0), first.hits);
1227     try std.testing.expectEqual(@as(usize, 1), second.hits);
1228     try std.testing.expect(!completion.active);
1229 }
1230 
1231 test "event loop disarm releases an active completion" {
1232     var loop = try Loop.init(.{ .allocator = std.testing.allocator, .entries = 1 });
1233     defer loop.deinit();
1234     var timer = try Timer.init();
1235     defer timer.deinit();
1236     var completion: Completion = .{};
1237     var state: TimerCount = .{};
1238     try timer.run(
1239         &loop,
1240         &completion,
1241         .{ .after = .fromMilliseconds(std.math.maxInt(u64)) },
1242         TimerCount,
1243         &state,
1244         timer_count_callback,
1245     );
1246     try std.testing.expect(completion.active);
1247 
1248     loop.disarm(&completion);
1249     try std.testing.expect(!completion.active);
1250     try timer.run(
1251         &loop,
1252         &completion,
1253         .{ .after = .zero },
1254         TimerCount,
1255         &state,
1256         timer_count_callback,
1257     );
1258     try loop.run(.until_done);
1259     try std.testing.expectEqual(@as(usize, 1), state.hits);
1260 }
1261 
1262 test "event loop fake clock orders and delays timers" {
1263     var fake = time.FakeClock.zero();
1264     var loop = try Loop.init(.{
1265         .allocator = std.testing.allocator,
1266         .entries = 2,
1267         .awake_clock = fake.awakeClock(),
1268     });
1269     defer loop.deinit();
1270     var timer = try Timer.init();
1271     defer timer.deinit();
1272     var completions: [2]Completion = @splat(.{});
1273     var log: [2]u8 = undefined;
1274     var count: u8 = 0;
1275     var late = TimerOrder{ .id = 2, .log = &log, .count = &count };
1276     var early = TimerOrder{ .id = 1, .log = &log, .count = &count };
1277     try timer.run(
1278         &loop,
1279         &completions[0],
1280         .{ .after = .fromMilliseconds(2) },
1281         TimerOrder,
1282         &late,
1283         timer_order_callback,
1284     );
1285     try timer.run(
1286         &loop,
1287         &completions[1],
1288         .{ .after = .fromMilliseconds(1) },
1289         TimerOrder,
1290         &early,
1291         timer_order_callback,
1292     );
1293     try loop.run(.no_wait);
1294     try std.testing.expectEqual(@as(u8, 0), count);
1295     fake.advance(.fromMilliseconds(2));
1296     try loop.run(.once);
1297     try loop.run(.once);
1298     try std.testing.expectEqualSlices(u8, &.{ 1, 2 }, &log);
1299 }
1300 
1301 test "event loop fake clock cancellation prevents a delayed timer" {
1302     var fake = time.FakeClock.zero();
1303     var loop = try Loop.init(.{
1304         .allocator = std.testing.allocator,
1305         .entries = 1,
1306         .awake_clock = fake.awakeClock(),
1307     });
1308     defer loop.deinit();
1309     var timer = try Timer.init();
1310     defer timer.deinit();
1311     var completion: Completion = .{};
1312     var state: TimerCount = .{};
1313     try timer.run(
1314         &loop,
1315         &completion,
1316         .{ .after = .fromMilliseconds(1) },
1317         TimerCount,
1318         &state,
1319         timer_count_callback,
1320     );
1321     loop.disarm(&completion);
1322     fake.advance(.fromMilliseconds(1));
1323     try loop.run(.no_wait);
1324     try std.testing.expectEqual(@as(usize, 0), state.hits);
1325 }
1326 
1327 test "event loop fixed rate skips delayed ticks within one slot" {
1328     var fake = time.FakeClock.zero();
1329     var loop = try Loop.init(.{
1330         .allocator = std.testing.allocator,
1331         .entries = 1,
1332         .awake_clock = fake.awakeClock(),
1333     });
1334     defer loop.deinit();
1335     var timer = try Timer.init();
1336     defer timer.deinit();
1337     var completion: Completion = .{};
1338     var state = TimerRearmState{ .clock = &fake, .late = .fromMilliseconds(25) };
1339     const interval = time.Duration.fromMilliseconds(10);
1340     try timer.run(
1341         &loop,
1342         &completion,
1343         .{ .after = interval, .repeat = .{ .fixed_rate = interval } },
1344         TimerRearmState,
1345         &state,
1346         timer_rearm_callback,
1347     );
1348     fake.advance(interval);
1349     try loop.run(.once);
1350     try std.testing.expectEqual(@as(u8, 1), state.hits);
1351     try std.testing.expectEqual(@as(u8, 1), Timer.delayed_tick_slots_max);
1352     try loop.run(.no_wait);
1353     try std.testing.expectEqual(@as(u8, 1), state.hits);
1354     fake.advance(.fromMilliseconds(5));
1355     try loop.run(.once);
1356     try std.testing.expectEqual(@as(u8, 2), state.hits);
1357 }
1358 
1359 test "event loop fixed delay starts after a late callback" {
1360     var fake = time.FakeClock.zero();
1361     var loop = try Loop.init(.{
1362         .allocator = std.testing.allocator,
1363         .entries = 1,
1364         .awake_clock = fake.awakeClock(),
1365     });
1366     defer loop.deinit();
1367     var timer = try Timer.init();
1368     defer timer.deinit();
1369     var completion: Completion = .{};
1370     var state = TimerRearmState{ .clock = &fake, .late = .fromMilliseconds(25) };
1371     const interval = time.Duration.fromMilliseconds(10);
1372     try timer.run(
1373         &loop,
1374         &completion,
1375         .{ .after = interval, .repeat = .{ .fixed_delay = interval } },
1376         TimerRearmState,
1377         &state,
1378         timer_rearm_callback,
1379     );
1380     fake.advance(interval);
1381     try loop.run(.once);
1382     fake.advance(.fromMilliseconds(9));
1383     try loop.run(.no_wait);
1384     try std.testing.expectEqual(@as(u8, 1), state.hits);
1385     fake.advance(.fromMilliseconds(1));
1386     try loop.run(.once);
1387     try std.testing.expectEqual(@as(u8, 2), state.hits);
1388 }
1389 
1390 test "event loop rejects an awake clock regression" {
1391     var fake = time.FakeClock.zero();
1392     fake.advance(.fromMilliseconds(2));
1393     var loop = try Loop.init(.{
1394         .allocator = std.testing.allocator,
1395         .entries = 1,
1396         .awake_clock = fake.awakeClock(),
1397     });
1398     defer loop.deinit();
1399     var timer = try Timer.init();
1400     defer timer.deinit();
1401     var completion: Completion = .{};
1402     var state: TimerCount = .{};
1403     try timer.run(
1404         &loop,
1405         &completion,
1406         .{ .after = .fromMilliseconds(1) },
1407         TimerCount,
1408         &state,
1409         timer_count_callback,
1410     );
1411     fake.regress(.fromMilliseconds(1));
1412     try std.testing.expectError(error.ClockRegressed, loop.run(.no_wait));
1413 }
1414 
1415 const ReadResultState = struct {
1416     hits: usize = 0,
1417     count: ?usize = null,
1418     err: ?ReadError = null,
1419 };
1420 
1421 fn read_result_callback(
1422     state: ?*ReadResultState,
1423     _: *Loop,
1424     _: *Completion,
1425     _: TCP,
1426     _: ReadBuffer,
1427     result: ReadError!usize,
1428 ) CallbackAction {
1429     const read_state = state.?;
1430     read_state.hits += 1;
1431     if (result) |count| {
1432         read_state.count = count;
1433     } else |err| {
1434         read_state.err = err;
1435     }
1436     return .disarm;
1437 }
1438 
1439 test "event TCP read reports peer shutdown as EOF" {
1440     const sockets = net.socketPairUnixStream() catch |err| switch (err) {
1441         error.UnsupportedPlatform => return error.SkipZigTest,
1442         else => return err,
1443     };
1444     defer net.close(sockets[0]);
1445     net.close(sockets[1]);
1446 
1447     var loop = try Loop.init(.{ .allocator = std.testing.allocator, .entries = 1 });
1448     defer loop.deinit();
1449     var completion: Completion = .{};
1450     var state: ReadResultState = .{};
1451     var buffer: [1]u8 = undefined;
1452     TCP.initFd(sockets[0]).read(
1453         &loop,
1454         &completion,
1455         .{ .slice = &buffer },
1456         ReadResultState,
1457         &state,
1458         read_result_callback,
1459     );
1460 
1461     try loop.run(.until_done);
1462     try std.testing.expectEqual(@as(usize, 1), state.hits);
1463     try std.testing.expectEqual(@as(?usize, null), state.count);
1464     try std.testing.expectEqual(@as(?ReadError, error.EOF), state.err);
1465     try std.testing.expect(!completion.active);
1466 }
1467 
1468 test "event TCP zero-length read preserves a ready byte" {
1469     const sockets = net.socketPairUnixStream() catch |err| switch (err) {
1470         error.UnsupportedPlatform => return error.SkipZigTest,
1471         else => return err,
1472     };
1473     defer net.close(sockets[0]);
1474     defer net.close(sockets[1]);
1475     try std.testing.expectEqual(@as(usize, 1), try net.sendNoSignal(sockets[0], "x"));
1476 
1477     var loop = try Loop.init(.{ .allocator = std.testing.allocator, .entries = 1 });
1478     defer loop.deinit();
1479     var completion: Completion = .{};
1480     var state: ReadResultState = .{};
1481     var buffer: [1]u8 = undefined;
1482     TCP.initFd(sockets[1]).read(
1483         &loop,
1484         &completion,
1485         .{ .slice = buffer[0..0] },
1486         ReadResultState,
1487         &state,
1488         read_result_callback,
1489     );
1490 
1491     try loop.run(.until_done);
1492     try std.testing.expectEqual(@as(usize, 1), state.hits);
1493     try std.testing.expectEqual(@as(?usize, 0), state.count);
1494     try std.testing.expectEqual(@as(?ReadError, null), state.err);
1495     try std.testing.expectEqual(@as(usize, 1), try net.recv(sockets[1], &buffer, 0));
1496     try std.testing.expectEqual(@as(u8, 'x'), buffer[0]);
1497 }
1498 
1499 const AcceptResultState = struct {
1500     hits: usize = 0,
1501     accepted: ?TCP = null,
1502     err: ?AcceptError = null,
1503 };
1504 
1505 fn accept_result_callback(
1506     state: ?*AcceptResultState,
1507     _: *Loop,
1508     _: *Completion,
1509     result: AcceptError!TCP,
1510 ) CallbackAction {
1511     const accept_state = state.?;
1512     accept_state.hits += 1;
1513     if (result) |accepted| {
1514         accept_state.accepted = accepted;
1515     } else |err| {
1516         accept_state.err = err;
1517     }
1518     return .disarm;
1519 }
1520 
1521 const CloseResultState = struct {
1522     hits: usize = 0,
1523     err: ?CloseError = null,
1524 };
1525 
1526 fn close_result_callback(
1527     state: ?*CloseResultState,
1528     _: *Loop,
1529     _: *Completion,
1530     _: TCP,
1531     result: CloseError!void,
1532 ) CallbackAction {
1533     const close_state = state.?;
1534     close_state.hits += 1;
1535     _ = result catch |err| {
1536         close_state.err = err;
1537     };
1538     return .disarm;
1539 }
1540 
1541 test "event TCP accepts and closes a loopback connection" {
1542     const listener = net.tcpStreamSocket(.{
1543         .close_on_exec = true,
1544         .nonblocking = true,
1545     }) catch |err| switch (err) {
1546         error.UnsupportedPlatform => return error.SkipZigTest,
1547         else => return err,
1548     };
1549     defer net.close(listener);
1550     try net.bindIp4(listener, net.ip4Address(.{ 127, 0, 0, 1 }, 0));
1551     try net.listen(listener, 1);
1552     const port = try net.socketPort(listener);
1553 
1554     const client = try net.tcpStreamSocket(.{ .close_on_exec = true });
1555     defer net.close(client);
1556     try net.connectIp4(client, net.ip4Address(.{ 127, 0, 0, 1 }, port));
1557 
1558     var loop = try Loop.init(.{ .allocator = std.testing.allocator, .entries = 1 });
1559     defer loop.deinit();
1560     var completion: Completion = .{};
1561     var accept_state: AcceptResultState = .{};
1562     TCP.initFd(listener).accept(
1563         &loop,
1564         &completion,
1565         AcceptResultState,
1566         &accept_state,
1567         accept_result_callback,
1568     );
1569     try loop.run(.until_done);
1570     try std.testing.expectEqual(@as(usize, 1), accept_state.hits);
1571     try std.testing.expectEqual(@as(?AcceptError, null), accept_state.err);
1572     const accepted = accept_state.accepted orelse return error.MissingAcceptedSocket;
1573     var accepted_open = true;
1574     defer if (accepted_open) net.close(accepted.fd);
1575 
1576     var close_state: CloseResultState = .{};
1577     accepted.close(
1578         &loop,
1579         &completion,
1580         CloseResultState,
1581         &close_state,
1582         close_result_callback,
1583     );
1584     accepted_open = false;
1585     try std.testing.expectEqual(@as(usize, 1), close_state.hits);
1586     try std.testing.expectEqual(@as(?CloseError, null), close_state.err);
1587     try std.testing.expect(!completion.active);
1588 }
1589 
1590 test "event TCP close replaces an unrelated active completion" {
1591     const sockets = net.socketPairUnixStream() catch |err| switch (err) {
1592         error.UnsupportedPlatform => return error.SkipZigTest,
1593         else => return err,
1594     };
1595     var first_open = true;
1596     defer if (first_open) net.close(sockets[0]);
1597     defer net.close(sockets[1]);
1598 
1599     var loop = try Loop.init(.{ .allocator = std.testing.allocator, .entries = 1 });
1600     defer loop.deinit();
1601     var timer = try Timer.init();
1602     defer timer.deinit();
1603     var completion: Completion = .{};
1604     var timer_state: TimerCount = .{};
1605     try timer.run(
1606         &loop,
1607         &completion,
1608         .{ .after = .fromMilliseconds(std.math.maxInt(u64)) },
1609         TimerCount,
1610         &timer_state,
1611         timer_count_callback,
1612     );
1613 
1614     var close_state: CloseResultState = .{};
1615     TCP.initFd(sockets[0]).close(
1616         &loop,
1617         &completion,
1618         CloseResultState,
1619         &close_state,
1620         close_result_callback,
1621     );
1622     first_open = false;
1623     try std.testing.expectEqual(@as(usize, 0), timer_state.hits);
1624     try std.testing.expectEqual(@as(usize, 1), close_state.hits);
1625     try std.testing.expect(!completion.active);
1626     try std.testing.expect(!loop.hasActive());
1627 }