| ... | ... | @@ -11,19 +11,21 @@ cond: std.Thread.Condition, |
| 11 | 11 | queue: std.DoublyLinkedList(void), |
| 12 | 12 | free: std.DoublyLinkedList(void), |
| 13 | 13 | main_fiber_buffer: [@sizeOf(Fiber) + max_result_len]u8 align(@alignOf(Fiber)), |
| 14 | | exiting: bool, |
| 14 | exit_awaiter: ?*Fiber, |
| 15 | 15 | idle_count: usize, |
| 16 | 16 | threads: std.ArrayListUnmanaged(Thread), |
| 17 | 17 | |
| 18 | | threadlocal var current_thread: *Thread = undefined; |
| 19 | | threadlocal var current_fiber: *Fiber = undefined; |
| 18 | threadlocal var current_idle_context: *Context = undefined; |
| 19 | threadlocal var current_fiber_context: *Context = undefined; |
| 20 | 20 | |
| 21 | 21 | const max_result_len = 64; |
| 22 | 22 | const min_stack_size = 4 * 1024 * 1024; |
| 23 | const idle_stack_size = 32 * 1024; |
| 24 | const stack_align = 16; |
| 23 | 25 | |
| 24 | 26 | const Thread = struct { |
| 25 | 27 | thread: std.Thread, |
| 26 | | idle_fiber: Fiber, |
| 28 | idle_context: Context, |
| 27 | 29 | }; |
| 28 | 30 | |
| 29 | 31 | const Fiber = struct { |
| ... | ... | @@ -54,6 +56,11 @@ const Fiber = struct { |
| 54 | 56 | }; |
| 55 | 57 | |
| 56 | 58 | pub fn init(el: *EventLoop, gpa: Allocator) error{OutOfMemory}!void { |
| 59 | const threads_bytes = ((std.Thread.getCpuCount() catch 1) -| 1) * @sizeOf(Thread); |
| 60 | const idle_context_offset = std.mem.alignForward(usize, threads_bytes, @alignOf(Context)); |
| 61 | const idle_stack_end_offset = std.mem.alignForward(usize, idle_context_offset + idle_stack_size, std.heap.page_size_max); |
| 62 | const allocated_slice = try gpa.alignedAlloc(u8, @max(@alignOf(Thread), @alignOf(Context), stack_align), idle_stack_end_offset); |
| 63 | errdefer gpa.free(allocated_slice); |
| 57 | 64 | el.* = .{ |
| 58 | 65 | .gpa = gpa, |
| 59 | 66 | .mutex = .{}, |
| ... | ... | @@ -61,28 +68,37 @@ pub fn init(el: *EventLoop, gpa: Allocator) error{OutOfMemory}!void { |
| 61 | 68 | .queue = .{}, |
| 62 | 69 | .free = .{}, |
| 63 | 70 | .main_fiber_buffer = undefined, |
| 64 | | .exiting = false, |
| 71 | .exit_awaiter = null, |
| 65 | 72 | .idle_count = 0, |
| 66 | | .threads = try .initCapacity(gpa, @max(std.Thread.getCpuCount() catch 1, 1)), |
| 73 | .threads = .initBuffer(@ptrCast(allocated_slice[0..threads_bytes])), |
| 67 | 74 | }; |
| 68 | | current_thread = el.threads.addOneAssumeCapacity(); |
| 69 | | current_fiber = @ptrCast(&el.main_fiber_buffer); |
| 75 | const main_idle_context: *Context = @alignCast(std.mem.bytesAsValue(Context, allocated_slice[idle_context_offset..][0..@sizeOf(Context)])); |
| 76 | const idle_stack_end: [*]align(stack_align) usize = @alignCast(@ptrCast(allocated_slice[idle_stack_end_offset..].ptr)); |
| 77 | (idle_stack_end - 1)[0..1].* = .{@intFromPtr(el)}; |
| 78 | main_idle_context.* = .{ |
| 79 | .rsp = @intFromPtr(idle_stack_end - 1), |
| 80 | .rbp = 0, |
| 81 | .rip = @intFromPtr(&mainIdleEntry), |
| 82 | }; |
| 83 | std.log.debug("created main idle {*}", .{main_idle_context}); |
| 84 | current_idle_context = main_idle_context; |
| 85 | const current_fiber: *Fiber = @ptrCast(&el.main_fiber_buffer); |
| 86 | std.log.debug("created main fiber {*}", .{current_fiber}); |
| 87 | current_fiber_context = &current_fiber.context; |
| 70 | 88 | } |
| 71 | 89 | |
| 72 | 90 | pub fn deinit(el: *EventLoop) void { |
| 73 | | { |
| 74 | | el.mutex.lock(); |
| 75 | | defer el.mutex.unlock(); |
| 76 | | assert(el.queue.len == 0); // pending async |
| 77 | | el.exiting = true; |
| 78 | | } |
| 79 | | el.cond.broadcast(); |
| 91 | assert(el.queue.len == 0); // pending async |
| 92 | el.yield(null, &el.exit_awaiter); |
| 80 | 93 | while (el.free.pop()) |free_node| { |
| 81 | 94 | const free_fiber: *Fiber = @fieldParentPtr("queue_node", free_node); |
| 82 | 95 | el.gpa.free(free_fiber.allocatedSlice()); |
| 83 | 96 | } |
| 84 | | for (el.threads.items[1..]) |*thread| thread.thread.join(); |
| 85 | | el.threads.deinit(el.gpa); |
| 97 | const idle_context_offset = std.mem.alignForward(usize, el.threads.items.len * @sizeOf(Thread), @alignOf(Context)); |
| 98 | const idle_stack_end = std.mem.alignForward(usize, idle_context_offset + idle_stack_size, std.heap.page_size_max); |
| 99 | const allocated_ptr: [*]align(@max(@alignOf(Thread), @alignOf(Context), stack_align)) u8 = @alignCast(@ptrCast(el.threads.items.ptr)); |
| 100 | for (el.threads.items) |*thread| thread.thread.join(); |
| 101 | el.gpa.free(allocated_ptr[0..idle_stack_end]); |
| 86 | 102 | } |
| 87 | 103 | |
| 88 | 104 | fn allocateFiber(el: *EventLoop, result_len: usize) error{OutOfMemory}!*Fiber { |
| ... | ... | @@ -103,40 +119,44 @@ fn allocateFiber(el: *EventLoop, result_len: usize) error{OutOfMemory}!*Fiber { |
| 103 | 119 | } |
| 104 | 120 | |
| 105 | 121 | fn yield(el: *EventLoop, optional_fiber: ?*Fiber, register_awaiter: ?*?*Fiber) void { |
| 106 | | const ready_fiber: *Fiber = optional_fiber orelse if (ready_node: { |
| 107 | | el.mutex.lock(); |
| 108 | | defer el.mutex.unlock(); |
| 109 | | break :ready_node el.queue.pop(); |
| 110 | | }) |ready_node| |
| 111 | | @fieldParentPtr("queue_node", ready_node) |
| 112 | | else |
| 113 | | &current_thread.idle_fiber; |
| 122 | const ready_context: *Context = ready_context: { |
| 123 | const ready_fiber: *Fiber = optional_fiber orelse if (ready_node: { |
| 124 | el.mutex.lock(); |
| 125 | defer el.mutex.unlock(); |
| 126 | break :ready_node el.queue.pop(); |
| 127 | }) |ready_node| |
| 128 | @fieldParentPtr("queue_node", ready_node) |
| 129 | else |
| 130 | break :ready_context current_idle_context; |
| 131 | break :ready_context &ready_fiber.context; |
| 132 | }; |
| 114 | 133 | const message: SwitchMessage = .{ |
| 115 | | .prev_context = &current_fiber.context, |
| 116 | | .ready_context = &ready_fiber.context, |
| 134 | .prev_context = current_fiber_context, |
| 135 | .ready_context = ready_context, |
| 117 | 136 | .register_awaiter = register_awaiter, |
| 118 | 137 | }; |
| 119 | | std.log.debug("switching from {*} to {*}", .{ |
| 120 | | @as(*Fiber, @fieldParentPtr("context", message.prev_context)), |
| 121 | | @as(*Fiber, @fieldParentPtr("context", message.ready_context)), |
| 122 | | }); |
| 138 | std.log.debug("switching from {*} to {*}", .{ message.prev_context, message.ready_context }); |
| 123 | 139 | contextSwitch(&message).handle(el); |
| 124 | 140 | } |
| 125 | 141 | |
| 126 | 142 | fn schedule(el: *EventLoop, fiber: *Fiber) void { |
| 127 | | signal: { |
| 128 | | el.mutex.lock(); |
| 129 | | defer el.mutex.unlock(); |
| 130 | | el.queue.append(&fiber.queue_node); |
| 131 | | if (el.idle_count > 0) break :signal; |
| 132 | | if (el.threads.items.len == el.threads.capacity) return; |
| 133 | | const thread = el.threads.addOneAssumeCapacity(); |
| 134 | | thread.thread = std.Thread.spawn(.{ |
| 135 | | .stack_size = min_stack_size, |
| 136 | | .allocator = el.gpa, |
| 137 | | }, threadEntry, .{ el, thread }) catch return; |
| 143 | el.mutex.lock(); |
| 144 | el.queue.append(&fiber.queue_node); |
| 145 | if (el.idle_count > 0) { |
| 146 | el.mutex.unlock(); |
| 147 | el.cond.signal(); |
| 148 | return; |
| 138 | 149 | } |
| 139 | | el.cond.signal(); |
| 150 | defer el.mutex.unlock(); |
| 151 | if (el.threads.items.len == el.threads.capacity) return; |
| 152 | const thread = el.threads.addOneAssumeCapacity(); |
| 153 | thread.thread = std.Thread.spawn(.{ |
| 154 | .stack_size = idle_stack_size, |
| 155 | .allocator = el.gpa, |
| 156 | }, threadEntry, .{ el, thread }) catch { |
| 157 | el.threads.items.len -= 1; |
| 158 | return; |
| 159 | }; |
| 140 | 160 | } |
| 141 | 161 | |
| 142 | 162 | fn recycle(el: *EventLoop, fiber: *Fiber) void { |
| ... | ... | @@ -148,14 +168,28 @@ fn recycle(el: *EventLoop, fiber: *Fiber) void { |
| 148 | 168 | el.free.append(&fiber.queue_node); |
| 149 | 169 | } |
| 150 | 170 | |
| 171 | fn mainIdle(el: *EventLoop, message: *const SwitchMessage) callconv(.c) noreturn { |
| 172 | message.handle(el); |
| 173 | el.yield(el.idle(), null); |
| 174 | unreachable; // switched to dead fiber |
| 175 | } |
| 176 | |
| 151 | 177 | fn threadEntry(el: *EventLoop, thread: *Thread) void { |
| 152 | | current_thread = thread; |
| 153 | | current_fiber = &thread.idle_fiber; |
| 178 | std.log.debug("created thread idle {*}", .{&thread.idle_context}); |
| 179 | current_idle_context = &thread.idle_context; |
| 180 | current_fiber_context = &thread.idle_context; |
| 181 | _ = el.idle(); |
| 182 | } |
| 183 | |
| 184 | fn idle(el: *EventLoop) *Fiber { |
| 154 | 185 | while (true) { |
| 155 | 186 | el.yield(null, null); |
| 187 | if (@atomicLoad(?*Fiber, &el.exit_awaiter, .acquire)) |exit_awaiter| { |
| 188 | el.cond.broadcast(); |
| 189 | return exit_awaiter; |
| 190 | } |
| 156 | 191 | el.mutex.lock(); |
| 157 | 192 | defer el.mutex.unlock(); |
| 158 | | if (el.exiting) return; |
| 159 | 193 | el.idle_count += 1; |
| 160 | 194 | defer el.idle_count -= 1; |
| 161 | 195 | el.cond.wait(&el.mutex); |
| ... | ... | @@ -169,7 +203,7 @@ const SwitchMessage = extern struct { |
| 169 | 203 | |
| 170 | 204 | fn handle(message: *const SwitchMessage, el: *EventLoop) void { |
| 171 | 205 | const prev_fiber: *Fiber = @fieldParentPtr("context", message.prev_context); |
| 172 | | current_fiber = @fieldParentPtr("context", message.ready_context); |
| 206 | current_fiber_context = message.ready_context; |
| 173 | 207 | if (message.register_awaiter) |awaiter| if (@atomicRmw(?*Fiber, awaiter, .Xchg, prev_fiber, .acq_rel) == Fiber.finished) el.schedule(prev_fiber); |
| 174 | 208 | } |
| 175 | 209 | }; |
| ... | ... | @@ -208,6 +242,18 @@ inline fn contextSwitch(message: *const SwitchMessage) *const SwitchMessage { |
| 208 | 242 | }; |
| 209 | 243 | } |
| 210 | 244 | |
| 245 | fn mainIdleEntry() callconv(.naked) void { |
| 246 | switch (builtin.cpu.arch) { |
| 247 | .x86_64 => asm volatile ( |
| 248 | \\ movq (%%rsp), %%rdi |
| 249 | \\ jmp %[mainIdle:P] |
| 250 | : |
| 251 | : [mainIdle] "X" (&mainIdle), |
| 252 | ), |
| 253 | else => |arch| @compileError("unimplemented architecture: " ++ @tagName(arch)), |
| 254 | } |
| 255 | } |
| 256 | |
| 211 | 257 | fn fiberEntry() callconv(.naked) void { |
| 212 | 258 | switch (builtin.cpu.arch) { |
| 213 | 259 | .x86_64 => asm volatile ( |
| ... | ... | @@ -238,7 +284,7 @@ pub fn @"async"( |
| 238 | 284 | const closure: *AsyncClosure = @ptrFromInt(std.mem.alignBackward( |
| 239 | 285 | usize, |
| 240 | 286 | @intFromPtr(fiber.stackEndPointer() - @sizeOf(AsyncClosure)), |
| 241 | | @alignOf(AsyncClosure), |
| 287 | @max(@alignOf(AsyncClosure), stack_align), |
| 242 | 288 | )); |
| 243 | 289 | closure.* = .{ |
| 244 | 290 | .event_loop = event_loop, |
| ... | ... | @@ -246,7 +292,7 @@ pub fn @"async"( |
| 246 | 292 | .fiber = fiber, |
| 247 | 293 | .start = start, |
| 248 | 294 | }; |
| 249 | | const stack_end: [*]align(16) usize = @alignCast(@ptrCast(closure)); |
| 295 | const stack_end: [*]align(stack_align) usize = @alignCast(@ptrCast(closure)); |
| 250 | 296 | fiber.context = .{ |
| 251 | 297 | .rsp = @intFromPtr(stack_end - 1), |
| 252 | 298 | .rbp = 0, |
| ... | ... | @@ -258,7 +304,6 @@ pub fn @"async"( |
| 258 | 304 | } |
| 259 | 305 | |
| 260 | 306 | const AsyncClosure = struct { |
| 261 | | _: void align(16) = {}, |
| 262 | 307 | event_loop: *EventLoop, |
| 263 | 308 | context: ?*anyopaque, |
| 264 | 309 | fiber: *Fiber, |