| ... | @@ -132,6 +132,7 @@ pub const Loop = struct { | ... | @@ -132,6 +132,7 @@ pub const Loop = struct { |
| 132 | | 132 | |
| 133 | pub const EventFd = struct { | 133 | pub const EventFd = struct { |
| 134 | base: ResumeNode, | 134 | base: ResumeNode, |
| | 135 | epoll_op: u32, |
| 135 | eventfd: i32, | 136 | eventfd: i32, |
| 136 | }; | 137 | }; |
| 137 | }; | 138 | }; |
| ... | @@ -204,6 +205,7 @@ pub const Loop = struct { | ... | @@ -204,6 +205,7 @@ pub const Loop = struct { |
| 204 | .handle = undefined, | 205 | .handle = undefined, |
| 205 | }, | 206 | }, |
| 206 | .eventfd = try std.os.linuxEventFd(1, posix.EFD_CLOEXEC | posix.EFD_NONBLOCK), | 207 | .eventfd = try std.os.linuxEventFd(1, posix.EFD_CLOEXEC | posix.EFD_NONBLOCK), |
| | 208 | .epoll_op = posix.EPOLL_CTL_ADD, |
| 207 | }, | 209 | }, |
| 208 | .next = undefined, | 210 | .next = undefined, |
| 209 | }; | 211 | }; |
| ... | @@ -265,15 +267,20 @@ pub const Loop = struct { | ... | @@ -265,15 +267,20 @@ pub const Loop = struct { |
| 265 | errdefer { | 267 | errdefer { |
| 266 | _ = @atomicRmw(usize, &self.pending_event_count, AtomicRmwOp.Sub, 1, AtomicOrder.SeqCst); | 268 | _ = @atomicRmw(usize, &self.pending_event_count, AtomicRmwOp.Sub, 1, AtomicOrder.SeqCst); |
| 267 | } | 269 | } |
| 268 | try self.addFdNoCounter(fd, resume_node); | 270 | try self.modFd( |
| | 271 | fd, |
| | 272 | posix.EPOLL_CTL_ADD, |
| | 273 | std.os.linux.EPOLLIN | std.os.linux.EPOLLOUT | std.os.linux.EPOLLET, |
| | 274 | resume_node, |
| | 275 | ); |
| 269 | } | 276 | } |
| 270 | | 277 | |
| 271 | fn addFdNoCounter(self: *Loop, fd: i32, resume_node: *ResumeNode) !void { | 278 | pub fn modFd(self: *Loop, fd: i32, op: u32, events: u32, resume_node: *ResumeNode) !void { |
| 272 | var ev = std.os.linux.epoll_event{ | 279 | var ev = std.os.linux.epoll_event{ |
| 273 | .events = std.os.linux.EPOLLIN | std.os.linux.EPOLLOUT | std.os.linux.EPOLLET, | 280 | .events = events, |
| 274 | .data = std.os.linux.epoll_data{ .ptr = @ptrToInt(resume_node) }, | 281 | .data = std.os.linux.epoll_data{ .ptr = @ptrToInt(resume_node) }, |
| 275 | }; | 282 | }; |
| 276 | try std.os.linuxEpollCtl(self.os_data.epollfd, std.os.linux.EPOLL_CTL_ADD, fd, &ev); | 283 | try std.os.linuxEpollCtl(self.os_data.epollfd, op, fd, &ev); |
| 277 | } | 284 | } |
| 278 | | 285 | |
| 279 | pub fn removeFd(self: *Loop, fd: i32) void { | 286 | pub fn removeFd(self: *Loop, fd: i32) void { |
| ... | @@ -331,7 +338,8 @@ pub const Loop = struct { | ... | @@ -331,7 +338,8 @@ pub const Loop = struct { |
| 331 | const eventfd_node = &resume_stack_node.data; | 338 | const eventfd_node = &resume_stack_node.data; |
| 332 | eventfd_node.base.handle = handle; | 339 | eventfd_node.base.handle = handle; |
| 333 | // the pending count is already accounted for | 340 | // the pending count is already accounted for |
| 334 | self.addFdNoCounter(eventfd_node.eventfd, &eventfd_node.base) catch |_| { | 341 | const epoll_events = posix.EPOLLONESHOT | std.os.linux.EPOLLIN | std.os.linux.EPOLLOUT | std.os.linux.EPOLLET; |
| | 342 | self.modFd(eventfd_node.eventfd, eventfd_node.epoll_op, epoll_events, &eventfd_node.base) catch |_| { |
| 335 | // fine, we didn't need it anyway | 343 | // fine, we didn't need it anyway |
| 336 | _ = @atomicRmw(u8, &self.dispatch_lock, AtomicRmwOp.Xchg, 0, AtomicOrder.SeqCst); | 344 | _ = @atomicRmw(u8, &self.dispatch_lock, AtomicRmwOp.Xchg, 0, AtomicOrder.SeqCst); |
| 337 | self.os_data.available_eventfd_resume_nodes.push(resume_stack_node); | 345 | self.os_data.available_eventfd_resume_nodes.push(resume_stack_node); |
| ... | @@ -371,7 +379,7 @@ pub const Loop = struct { | ... | @@ -371,7 +379,7 @@ pub const Loop = struct { |
| 371 | ResumeNode.Id.Stop => return, | 379 | ResumeNode.Id.Stop => return, |
| 372 | ResumeNode.Id.EventFd => { | 380 | ResumeNode.Id.EventFd => { |
| 373 | const event_fd_node = @fieldParentPtr(ResumeNode.EventFd, "base", resume_node); | 381 | const event_fd_node = @fieldParentPtr(ResumeNode.EventFd, "base", resume_node); |
| 374 | self.removeFdNoCounter(event_fd_node.eventfd); | 382 | event_fd_node.epoll_op = posix.EPOLL_CTL_MOD; |
| 375 | const stack_node = @fieldParentPtr(std.atomic.Stack(ResumeNode.EventFd).Node, "data", event_fd_node); | 383 | const stack_node = @fieldParentPtr(std.atomic.Stack(ResumeNode.EventFd).Node, "data", event_fd_node); |
| 376 | self.os_data.available_eventfd_resume_nodes.push(stack_node); | 384 | self.os_data.available_eventfd_resume_nodes.push(stack_node); |
| 377 | }, | 385 | }, |
| ... | @@ -902,7 +910,7 @@ test "std.event.Lock" { | ... | @@ -902,7 +910,7 @@ test "std.event.Lock" { |
| 902 | defer cancel handle; | 910 | defer cancel handle; |
| 903 | loop.run(); | 911 | loop.run(); |
| 904 | | 912 | |
| 905 | assert(mem.eql(i32, shared_test_data, [1]i32{3 * 10} ** 10)); | 913 | assert(mem.eql(i32, shared_test_data, [1]i32{3 * @intCast(i32, shared_test_data.len)} ** shared_test_data.len)); |
| 906 | } | 914 | } |
| 907 | | 915 | |
| 908 | async fn testLock(loop: *Loop, lock: *Lock) void { | 916 | async fn testLock(loop: *Loop, lock: *Lock) void { |
| ... | @@ -945,7 +953,7 @@ async fn lockRunner(lock: *Lock) void { | ... | @@ -945,7 +953,7 @@ async fn lockRunner(lock: *Lock) void { |
| 945 | suspend; // resumed by onNextTick | 953 | suspend; // resumed by onNextTick |
| 946 | | 954 | |
| 947 | var i: usize = 0; | 955 | var i: usize = 0; |
| 948 | while (i < 10) : (i += 1) { | 956 | while (i < shared_test_data.len) : (i += 1) { |
| 949 | const handle = await (async lock.acquire() catch @panic("out of memory")); | 957 | const handle = await (async lock.acquire() catch @panic("out of memory")); |
| 950 | defer handle.release(); | 958 | defer handle.release(); |
| 951 | | 959 | |