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 }