| ... | ... | @@ -118,6 +118,11 @@ pub const Loop = struct { |
| 118 | 118 | extra_threads: []*std.os.Thread, |
| 119 | 119 | final_resume_node: ResumeNode, |
| 120 | 120 | |
| 121 | // pre-allocated eventfds. all permanently active. |
| 122 | // this is how we send promises to be resumed on other threads. |
| 123 | available_eventfd_resume_nodes: std.atomic.Stack(ResumeNode.EventFd), |
| 124 | eventfd_resume_nodes: []std.atomic.Stack(ResumeNode.EventFd).Node, |
| 125 | |
| 121 | 126 | pub const NextTickNode = std.atomic.QueueMpsc(promise).Node; |
| 122 | 127 | |
| 123 | 128 | pub const ResumeNode = struct { |
| ... | ... | @@ -130,10 +135,17 @@ pub const Loop = struct { |
| 130 | 135 | EventFd, |
| 131 | 136 | }; |
| 132 | 137 | |
| 133 | | pub const EventFd = struct { |
| 134 | | base: ResumeNode, |
| 135 | | epoll_op: u32, |
| 136 | | eventfd: i32, |
| 138 | pub const EventFd = switch (builtin.os) { |
| 139 | builtin.Os.macosx => struct { |
| 140 | base: ResumeNode, |
| 141 | kevent: posix.Kevent, |
| 142 | }, |
| 143 | builtin.Os.linux => struct { |
| 144 | base: ResumeNode, |
| 145 | epoll_op: u32, |
| 146 | eventfd: i32, |
| 147 | }, |
| 148 | else => @compileError("unsupported OS"), |
| 137 | 149 | }; |
| 138 | 150 | }; |
| 139 | 151 | |
| ... | ... | @@ -168,36 +180,41 @@ pub const Loop = struct { |
| 168 | 180 | .id = ResumeNode.Id.Stop, |
| 169 | 181 | .handle = undefined, |
| 170 | 182 | }, |
| 183 | .available_eventfd_resume_nodes = std.atomic.Stack(ResumeNode.EventFd).init(), |
| 184 | .eventfd_resume_nodes = undefined, |
| 171 | 185 | }; |
| 172 | | try self.initOsData(thread_count); |
| 186 | const extra_thread_count = thread_count - 1; |
| 187 | self.eventfd_resume_nodes = try self.allocator.alloc( |
| 188 | std.atomic.Stack(ResumeNode.EventFd).Node, |
| 189 | extra_thread_count, |
| 190 | ); |
| 191 | errdefer self.allocator.free(self.eventfd_resume_nodes); |
| 192 | |
| 193 | self.extra_threads = try self.allocator.alloc(*std.os.Thread, extra_thread_count); |
| 194 | errdefer self.allocator.free(self.extra_threads); |
| 195 | |
| 196 | try self.initOsData(extra_thread_count); |
| 173 | 197 | errdefer self.deinitOsData(); |
| 174 | 198 | } |
| 175 | 199 | |
| 176 | 200 | /// must call stop before deinit |
| 177 | 201 | pub fn deinit(self: *Loop) void { |
| 178 | 202 | self.deinitOsData(); |
| 203 | self.allocator.free(self.extra_threads); |
| 179 | 204 | } |
| 180 | 205 | |
| 181 | 206 | const InitOsDataError = std.os.LinuxEpollCreateError || mem.Allocator.Error || std.os.LinuxEventFdError || |
| 182 | | std.os.SpawnThreadError || std.os.LinuxEpollCtlError; |
| 207 | std.os.SpawnThreadError || std.os.LinuxEpollCtlError || std.os.BsdKEventError; |
| 183 | 208 | |
| 184 | 209 | const wakeup_bytes = []u8{0x1} ** 8; |
| 185 | 210 | |
| 186 | | fn initOsData(self: *Loop, thread_count: usize) InitOsDataError!void { |
| 211 | fn initOsData(self: *Loop, extra_thread_count: usize) InitOsDataError!void { |
| 187 | 212 | switch (builtin.os) { |
| 188 | 213 | builtin.Os.linux => { |
| 189 | | const extra_thread_count = thread_count - 1; |
| 190 | | self.os_data.available_eventfd_resume_nodes = std.atomic.Stack(ResumeNode.EventFd).init(); |
| 191 | | self.os_data.eventfd_resume_nodes = try self.allocator.alloc( |
| 192 | | std.atomic.Stack(ResumeNode.EventFd).Node, |
| 193 | | extra_thread_count, |
| 194 | | ); |
| 195 | | errdefer self.allocator.free(self.os_data.eventfd_resume_nodes); |
| 196 | | |
| 197 | 214 | errdefer { |
| 198 | | while (self.os_data.available_eventfd_resume_nodes.pop()) |node| std.os.close(node.data.eventfd); |
| 215 | while (self.available_eventfd_resume_nodes.pop()) |node| std.os.close(node.data.eventfd); |
| 199 | 216 | } |
| 200 | | for (self.os_data.eventfd_resume_nodes) |*eventfd_node| { |
| 217 | for (self.eventfd_resume_nodes) |*eventfd_node| { |
| 201 | 218 | eventfd_node.* = std.atomic.Stack(ResumeNode.EventFd).Node{ |
| 202 | 219 | .data = ResumeNode.EventFd{ |
| 203 | 220 | .base = ResumeNode{ |
| ... | ... | @@ -209,7 +226,7 @@ pub const Loop = struct { |
| 209 | 226 | }, |
| 210 | 227 | .next = undefined, |
| 211 | 228 | }; |
| 212 | | self.os_data.available_eventfd_resume_nodes.push(eventfd_node); |
| 229 | self.available_eventfd_resume_nodes.push(eventfd_node); |
| 213 | 230 | } |
| 214 | 231 | |
| 215 | 232 | self.os_data.epollfd = try std.os.linuxEpollCreate(posix.EPOLL_CLOEXEC); |
| ... | ... | @@ -228,15 +245,84 @@ pub const Loop = struct { |
| 228 | 245 | self.os_data.final_eventfd, |
| 229 | 246 | &self.os_data.final_eventfd_event, |
| 230 | 247 | ); |
| 231 | | self.extra_threads = try self.allocator.alloc(*std.os.Thread, extra_thread_count); |
| 232 | | errdefer self.allocator.free(self.extra_threads); |
| 233 | 248 | |
| 234 | 249 | var extra_thread_index: usize = 0; |
| 235 | 250 | errdefer { |
| 251 | // writing 8 bytes to an eventfd cannot fail |
| 252 | std.os.posixWrite(self.os_data.final_eventfd, wakeup_bytes) catch unreachable; |
| 253 | while (extra_thread_index != 0) { |
| 254 | extra_thread_index -= 1; |
| 255 | self.extra_threads[extra_thread_index].wait(); |
| 256 | } |
| 257 | } |
| 258 | while (extra_thread_index < extra_thread_count) : (extra_thread_index += 1) { |
| 259 | self.extra_threads[extra_thread_index] = try std.os.spawnThread(self, workerRun); |
| 260 | } |
| 261 | }, |
| 262 | builtin.Os.macosx => { |
| 263 | self.os_data.kqfd = try std.os.bsdKQueue(); |
| 264 | errdefer std.os.close(self.os_data.kqfd); |
| 265 | |
| 266 | self.os_data.kevents = try self.allocator.alloc(posix.Kevent, extra_thread_count); |
| 267 | errdefer self.allocator.free(self.os_data.kevents); |
| 268 | |
| 269 | const eventlist = ([*]posix.Kevent)(undefined)[0..0]; |
| 270 | |
| 271 | for (self.eventfd_resume_nodes) |*eventfd_node, i| { |
| 272 | eventfd_node.* = std.atomic.Stack(ResumeNode.EventFd).Node{ |
| 273 | .data = ResumeNode.EventFd{ |
| 274 | .base = ResumeNode{ |
| 275 | .id = ResumeNode.Id.EventFd, |
| 276 | .handle = undefined, |
| 277 | }, |
| 278 | // this one is for sending events |
| 279 | .kevent = posix.Kevent { |
| 280 | .ident = i, |
| 281 | .filter = posix.EVFILT_USER, |
| 282 | .flags = posix.EV_CLEAR|posix.EV_ADD|posix.EV_DISABLE, |
| 283 | .fflags = 0, |
| 284 | .data = 0, |
| 285 | .udata = @ptrToInt(&eventfd_node.data.base), |
| 286 | }, |
| 287 | }, |
| 288 | .next = undefined, |
| 289 | }; |
| 290 | self.available_eventfd_resume_nodes.push(eventfd_node); |
| 291 | const kevent_array = (*[1]posix.Kevent)(&eventfd_node.data.kevent); |
| 292 | _ = try std.os.bsdKEvent(self.os_data.kqfd, kevent_array, eventlist, null); |
| 293 | eventfd_node.data.kevent.flags = posix.EV_CLEAR|posix.EV_ENABLE; |
| 294 | eventfd_node.data.kevent.fflags = posix.NOTE_TRIGGER; |
| 295 | // this one is for waiting for events |
| 296 | self.os_data.kevents[i] = posix.Kevent { |
| 297 | .ident = i, |
| 298 | .filter = posix.EVFILT_USER, |
| 299 | .flags = 0, |
| 300 | .fflags = 0, |
| 301 | .data = 0, |
| 302 | .udata = @ptrToInt(&eventfd_node.data.base), |
| 303 | }; |
| 304 | } |
| 305 | |
| 306 | // Pre-add so that we cannot get error.SystemResources |
| 307 | // later when we try to activate it. |
| 308 | self.os_data.final_kevent = posix.Kevent{ |
| 309 | .ident = extra_thread_count, |
| 310 | .filter = posix.EVFILT_USER, |
| 311 | .flags = posix.EV_ADD | posix.EV_DISABLE, |
| 312 | .fflags = 0, |
| 313 | .data = 0, |
| 314 | .udata = @ptrToInt(&self.final_resume_node), |
| 315 | }; |
| 316 | const kevent_array = (*[1]posix.Kevent)(&self.os_data.final_kevent); |
| 317 | _ = try std.os.bsdKEvent(self.os_data.kqfd, kevent_array, eventlist, null); |
| 318 | self.os_data.final_kevent.flags = posix.EV_ENABLE; |
| 319 | self.os_data.final_kevent.fflags = posix.NOTE_TRIGGER; |
| 320 | |
| 321 | var extra_thread_index: usize = 0; |
| 322 | errdefer { |
| 323 | _ = std.os.bsdKEvent(self.os_data.kqfd, kevent_array, eventlist, null) catch unreachable; |
| 236 | 324 | while (extra_thread_index != 0) { |
| 237 | 325 | extra_thread_index -= 1; |
| 238 | | // writing 8 bytes to an eventfd cannot fail |
| 239 | | std.os.posixWrite(self.os_data.final_eventfd, wakeup_bytes) catch unreachable; |
| 240 | 326 | self.extra_threads[extra_thread_index].wait(); |
| 241 | 327 | } |
| 242 | 328 | } |
| ... | ... | @@ -252,10 +338,12 @@ pub const Loop = struct { |
| 252 | 338 | switch (builtin.os) { |
| 253 | 339 | builtin.Os.linux => { |
| 254 | 340 | std.os.close(self.os_data.final_eventfd); |
| 255 | | while (self.os_data.available_eventfd_resume_nodes.pop()) |node| std.os.close(node.data.eventfd); |
| 341 | while (self.available_eventfd_resume_nodes.pop()) |node| std.os.close(node.data.eventfd); |
| 256 | 342 | std.os.close(self.os_data.epollfd); |
| 257 | | self.allocator.free(self.os_data.eventfd_resume_nodes); |
| 258 | | self.allocator.free(self.extra_threads); |
| 343 | self.allocator.free(self.eventfd_resume_nodes); |
| 344 | }, |
| 345 | builtin.Os.macosx => { |
| 346 | self.allocator.free(self.os_data.kevents); |
| 259 | 347 | }, |
| 260 | 348 | else => {}, |
| 261 | 349 | } |
| ... | ... | @@ -332,21 +420,38 @@ pub const Loop = struct { |
| 332 | 420 | continue :start_over; |
| 333 | 421 | } |
| 334 | 422 | |
| 335 | | // non-last node, stick it in the epoll set so that |
| 423 | // non-last node, stick it in the epoll/kqueue set so that |
| 336 | 424 | // other threads can get to it |
| 337 | | if (self.os_data.available_eventfd_resume_nodes.pop()) |resume_stack_node| { |
| 425 | if (self.available_eventfd_resume_nodes.pop()) |resume_stack_node| { |
| 338 | 426 | const eventfd_node = &resume_stack_node.data; |
| 339 | 427 | eventfd_node.base.handle = handle; |
| 340 | | // the pending count is already accounted for |
| 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 |_| { |
| 343 | | // fine, we didn't need it anyway |
| 344 | | _ = @atomicRmw(u8, &self.dispatch_lock, AtomicRmwOp.Xchg, 0, AtomicOrder.SeqCst); |
| 345 | | self.os_data.available_eventfd_resume_nodes.push(resume_stack_node); |
| 346 | | resume handle; |
| 347 | | _ = @atomicRmw(usize, &self.pending_event_count, AtomicRmwOp.Sub, 1, AtomicOrder.SeqCst); |
| 348 | | continue :start_over; |
| 349 | | }; |
| 428 | switch (builtin.os) { |
| 429 | builtin.Os.macosx => { |
| 430 | const kevent_array = (*[1]posix.Kevent)(&eventfd_node.kevent); |
| 431 | const eventlist = ([*]posix.Kevent)(undefined)[0..0]; |
| 432 | _ = std.os.bsdKEvent(self.os_data.kqfd, kevent_array, eventlist, null) catch |_| { |
| 433 | // fine, we didn't need it anyway |
| 434 | _ = @atomicRmw(u8, &self.dispatch_lock, AtomicRmwOp.Xchg, 0, AtomicOrder.SeqCst); |
| 435 | self.available_eventfd_resume_nodes.push(resume_stack_node); |
| 436 | resume handle; |
| 437 | _ = @atomicRmw(usize, &self.pending_event_count, AtomicRmwOp.Sub, 1, AtomicOrder.SeqCst); |
| 438 | continue :start_over; |
| 439 | }; |
| 440 | }, |
| 441 | builtin.Os.linux => { |
| 442 | // the pending count is already accounted for |
| 443 | const epoll_events = posix.EPOLLONESHOT | std.os.linux.EPOLLIN | std.os.linux.EPOLLOUT | std.os.linux.EPOLLET; |
| 444 | self.modFd(eventfd_node.eventfd, eventfd_node.epoll_op, epoll_events, &eventfd_node.base) catch |_| { |
| 445 | // fine, we didn't need it anyway |
| 446 | _ = @atomicRmw(u8, &self.dispatch_lock, AtomicRmwOp.Xchg, 0, AtomicOrder.SeqCst); |
| 447 | self.available_eventfd_resume_nodes.push(resume_stack_node); |
| 448 | resume handle; |
| 449 | _ = @atomicRmw(usize, &self.pending_event_count, AtomicRmwOp.Sub, 1, AtomicOrder.SeqCst); |
| 450 | continue :start_over; |
| 451 | }; |
| 452 | }, |
| 453 | else => @compileError("unsupported OS"), |
| 454 | } |
| 350 | 455 | } else { |
| 351 | 456 | // threads are too busy, can't add another eventfd to wake one up |
| 352 | 457 | _ = @atomicRmw(u8, &self.dispatch_lock, AtomicRmwOp.Xchg, 0, AtomicOrder.SeqCst); |
| ... | ... | @@ -359,35 +464,74 @@ pub const Loop = struct { |
| 359 | 464 | const pending_event_count = @atomicLoad(usize, &self.pending_event_count, AtomicOrder.SeqCst); |
| 360 | 465 | if (pending_event_count == 0) { |
| 361 | 466 | // cause all the threads to stop |
| 362 | | // writing 8 bytes to an eventfd cannot fail |
| 363 | | std.os.posixWrite(self.os_data.final_eventfd, wakeup_bytes) catch unreachable; |
| 364 | | return; |
| 467 | switch (builtin.os) { |
| 468 | builtin.Os.linux => { |
| 469 | // writing 8 bytes to an eventfd cannot fail |
| 470 | std.os.posixWrite(self.os_data.final_eventfd, wakeup_bytes) catch unreachable; |
| 471 | return; |
| 472 | }, |
| 473 | builtin.Os.macosx => { |
| 474 | const final_kevent = (*[1]posix.Kevent)(&self.os_data.final_kevent); |
| 475 | const eventlist = ([*]posix.Kevent)(undefined)[0..0]; |
| 476 | // cannot fail because we already added it and this just enables it |
| 477 | _ = std.os.bsdKEvent(self.os_data.kqfd, final_kevent, eventlist, null) catch unreachable; |
| 478 | return; |
| 479 | }, |
| 480 | else => @compileError("unsupported OS"), |
| 481 | } |
| 365 | 482 | } |
| 366 | 483 | |
| 367 | 484 | _ = @atomicRmw(u8, &self.dispatch_lock, AtomicRmwOp.Xchg, 0, AtomicOrder.SeqCst); |
| 368 | 485 | } |
| 369 | 486 | |
| 370 | | // only process 1 event so we don't steal from other threads |
| 371 | | var events: [1]std.os.linux.epoll_event = undefined; |
| 372 | | const count = std.os.linuxEpollWait(self.os_data.epollfd, events[0..], -1); |
| 373 | | for (events[0..count]) |ev| { |
| 374 | | const resume_node = @intToPtr(*ResumeNode, ev.data.ptr); |
| 375 | | const handle = resume_node.handle; |
| 376 | | const resume_node_id = resume_node.id; |
| 377 | | switch (resume_node_id) { |
| 378 | | ResumeNode.Id.Basic => {}, |
| 379 | | ResumeNode.Id.Stop => return, |
| 380 | | ResumeNode.Id.EventFd => { |
| 381 | | const event_fd_node = @fieldParentPtr(ResumeNode.EventFd, "base", resume_node); |
| 382 | | event_fd_node.epoll_op = posix.EPOLL_CTL_MOD; |
| 383 | | const stack_node = @fieldParentPtr(std.atomic.Stack(ResumeNode.EventFd).Node, "data", event_fd_node); |
| 384 | | self.os_data.available_eventfd_resume_nodes.push(stack_node); |
| 385 | | }, |
| 386 | | } |
| 387 | | resume handle; |
| 388 | | if (resume_node_id == ResumeNode.Id.EventFd) { |
| 389 | | _ = @atomicRmw(usize, &self.pending_event_count, AtomicRmwOp.Sub, 1, AtomicOrder.SeqCst); |
| 390 | | } |
| 487 | switch (builtin.os) { |
| 488 | builtin.Os.linux => { |
| 489 | // only process 1 event so we don't steal from other threads |
| 490 | var events: [1]std.os.linux.epoll_event = undefined; |
| 491 | const count = std.os.linuxEpollWait(self.os_data.epollfd, events[0..], -1); |
| 492 | for (events[0..count]) |ev| { |
| 493 | const resume_node = @intToPtr(*ResumeNode, ev.data.ptr); |
| 494 | const handle = resume_node.handle; |
| 495 | const resume_node_id = resume_node.id; |
| 496 | switch (resume_node_id) { |
| 497 | ResumeNode.Id.Basic => {}, |
| 498 | ResumeNode.Id.Stop => return, |
| 499 | ResumeNode.Id.EventFd => { |
| 500 | const event_fd_node = @fieldParentPtr(ResumeNode.EventFd, "base", resume_node); |
| 501 | event_fd_node.epoll_op = posix.EPOLL_CTL_MOD; |
| 502 | const stack_node = @fieldParentPtr(std.atomic.Stack(ResumeNode.EventFd).Node, "data", event_fd_node); |
| 503 | self.available_eventfd_resume_nodes.push(stack_node); |
| 504 | }, |
| 505 | } |
| 506 | resume handle; |
| 507 | if (resume_node_id == ResumeNode.Id.EventFd) { |
| 508 | _ = @atomicRmw(usize, &self.pending_event_count, AtomicRmwOp.Sub, 1, AtomicOrder.SeqCst); |
| 509 | } |
| 510 | } |
| 511 | }, |
| 512 | builtin.Os.macosx => { |
| 513 | var eventlist: [1]posix.Kevent = undefined; |
| 514 | const count = std.os.bsdKEvent(self.os_data.kqfd, self.os_data.kevents, eventlist[0..], null) catch unreachable; |
| 515 | for (eventlist[0..count]) |ev| { |
| 516 | const resume_node = @intToPtr(*ResumeNode, ev.udata); |
| 517 | const handle = resume_node.handle; |
| 518 | const resume_node_id = resume_node.id; |
| 519 | switch (resume_node_id) { |
| 520 | ResumeNode.Id.Basic => {}, |
| 521 | ResumeNode.Id.Stop => return, |
| 522 | ResumeNode.Id.EventFd => { |
| 523 | const event_fd_node = @fieldParentPtr(ResumeNode.EventFd, "base", resume_node); |
| 524 | const stack_node = @fieldParentPtr(std.atomic.Stack(ResumeNode.EventFd).Node, "data", event_fd_node); |
| 525 | self.available_eventfd_resume_nodes.push(stack_node); |
| 526 | }, |
| 527 | } |
| 528 | resume handle; |
| 529 | if (resume_node_id == ResumeNode.Id.EventFd) { |
| 530 | _ = @atomicRmw(usize, &self.pending_event_count, AtomicRmwOp.Sub, 1, AtomicOrder.SeqCst); |
| 531 | } |
| 532 | } |
| 533 | }, |
| 534 | else => @compileError("unsupported OS"), |
| 391 | 535 | } |
| 392 | 536 | } |
| 393 | 537 | } |
| ... | ... | @@ -395,12 +539,13 @@ pub const Loop = struct { |
| 395 | 539 | const OsData = switch (builtin.os) { |
| 396 | 540 | builtin.Os.linux => struct { |
| 397 | 541 | epollfd: i32, |
| 398 | | // pre-allocated eventfds. all permanently active. |
| 399 | | // this is how we send promises to be resumed on other threads. |
| 400 | | available_eventfd_resume_nodes: std.atomic.Stack(ResumeNode.EventFd), |
| 401 | | eventfd_resume_nodes: []std.atomic.Stack(ResumeNode.EventFd).Node, |
| 402 | 542 | final_eventfd: i32, |
| 403 | | final_eventfd_event: posix.epoll_event, |
| 543 | final_eventfd_event: std.os.linux.epoll_event, |
| 544 | }, |
| 545 | builtin.Os.macosx => struct { |
| 546 | kqfd: i32, |
| 547 | final_kevent: posix.Kevent, |
| 548 | kevents: []posix.Kevent, |
| 404 | 549 | }, |
| 405 | 550 | else => struct {}, |
| 406 | 551 | }; |