| ... | ... | @@ -18,7 +18,11 @@ threads: std.ArrayListUnmanaged(Thread), |
| 18 | 18 | threadlocal var current_idle_context: *Context = undefined; |
| 19 | 19 | threadlocal var current_fiber_context: *Context = undefined; |
| 20 | 20 | |
| 21 | /// Also used for context. |
| 21 | 22 | const max_result_len = 64; |
| 23 | /// Also used for context. |
| 24 | const max_result_align: std.mem.Alignment = .@"16"; |
| 25 | |
| 22 | 26 | const min_stack_size = 4 * 1024 * 1024; |
| 23 | 27 | const idle_stack_size = 32 * 1024; |
| 24 | 28 | const stack_align = 16; |
| ... | ... | @@ -29,6 +33,8 @@ const Thread = struct { |
| 29 | 33 | }; |
| 30 | 34 | |
| 31 | 35 | const Fiber = struct { |
| 36 | _: void align(max_result_align.toByteUnits()) = {}, |
| 37 | |
| 32 | 38 | context: Context, |
| 33 | 39 | awaiter: ?*Fiber, |
| 34 | 40 | queue_node: std.DoublyLinkedList(void).Node, |
| ... | ... | @@ -39,14 +45,27 @@ const Fiber = struct { |
| 39 | 45 | const base: [*]align(@alignOf(Fiber)) u8 = @ptrCast(f); |
| 40 | 46 | return base[0..std.mem.alignForward( |
| 41 | 47 | usize, |
| 42 | | @sizeOf(Fiber) + max_result_len + min_stack_size, |
| 48 | resultOffset() + max_result_len + min_stack_size, |
| 43 | 49 | std.heap.page_size_max, |
| 44 | 50 | )]; |
| 45 | 51 | } |
| 46 | 52 | |
| 53 | fn argsOffset() usize { |
| 54 | return max_result_align.forward(@sizeOf(Fiber)); |
| 55 | } |
| 56 | |
| 57 | fn resultOffset() usize { |
| 58 | return max_result_align.forward(argsOffset() + max_result_len); |
| 59 | } |
| 60 | |
| 61 | fn argsSlice(f: *Fiber) []u8 { |
| 62 | const base: [*]align(@alignOf(Fiber)) u8 = @ptrCast(f); |
| 63 | return base[argsOffset()..][0..max_result_len]; |
| 64 | } |
| 65 | |
| 47 | 66 | fn resultSlice(f: *Fiber) []u8 { |
| 48 | 67 | const base: [*]align(@alignOf(Fiber)) u8 = @ptrCast(f); |
| 49 | | return base[@sizeOf(Fiber)..][0..max_result_len]; |
| 68 | return base[resultOffset()..][0..max_result_len]; |
| 50 | 69 | } |
| 51 | 70 | |
| 52 | 71 | fn stackEndPointer(f: *Fiber) [*]u8 { |
| ... | ... | @@ -102,7 +121,7 @@ pub fn deinit(el: *EventLoop) void { |
| 102 | 121 | assert(el.queue.len == 0); // pending async |
| 103 | 122 | el.yield(null, &el.exit_awaiter); |
| 104 | 123 | while (el.free.pop()) |free_node| { |
| 105 | | const free_fiber: *Fiber = @fieldParentPtr("queue_node", free_node); |
| 124 | const free_fiber: *Fiber = @alignCast(@fieldParentPtr("queue_node", free_node)); |
| 106 | 125 | el.gpa.free(free_fiber.allocatedSlice()); |
| 107 | 126 | } |
| 108 | 127 | const idle_context_offset = std.mem.alignForward(usize, el.threads.items.len * @sizeOf(Thread), @alignOf(Context)); |
| ... | ... | @@ -112,8 +131,7 @@ pub fn deinit(el: *EventLoop) void { |
| 112 | 131 | el.gpa.free(allocated_ptr[0..idle_stack_end]); |
| 113 | 132 | } |
| 114 | 133 | |
| 115 | | fn allocateFiber(el: *EventLoop, result_len: usize) error{OutOfMemory}!*Fiber { |
| 116 | | assert(result_len <= max_result_len); |
| 134 | fn allocateFiber(el: *EventLoop) error{OutOfMemory}!*Fiber { |
| 117 | 135 | const free_node = free_node: { |
| 118 | 136 | el.mutex.lock(); |
| 119 | 137 | defer el.mutex.unlock(); |
| ... | ... | @@ -121,12 +139,12 @@ fn allocateFiber(el: *EventLoop, result_len: usize) error{OutOfMemory}!*Fiber { |
| 121 | 139 | } orelse { |
| 122 | 140 | const n = std.mem.alignForward( |
| 123 | 141 | usize, |
| 124 | | @sizeOf(Fiber) + max_result_len + min_stack_size, |
| 142 | Fiber.resultOffset() + max_result_len + min_stack_size, |
| 125 | 143 | std.heap.page_size_max, |
| 126 | 144 | ); |
| 127 | 145 | return @alignCast(@ptrCast(try el.gpa.alignedAlloc(u8, @alignOf(Fiber), n))); |
| 128 | 146 | }; |
| 129 | | return @fieldParentPtr("queue_node", free_node); |
| 147 | return @alignCast(@fieldParentPtr("queue_node", free_node)); |
| 130 | 148 | } |
| 131 | 149 | |
| 132 | 150 | fn yield(el: *EventLoop, optional_fiber: ?*Fiber, register_awaiter: ?*?*Fiber) void { |
| ... | ... | @@ -136,7 +154,7 @@ fn yield(el: *EventLoop, optional_fiber: ?*Fiber, register_awaiter: ?*?*Fiber) v |
| 136 | 154 | defer el.mutex.unlock(); |
| 137 | 155 | break :ready_node el.queue.pop(); |
| 138 | 156 | }) |ready_node| |
| 139 | | @fieldParentPtr("queue_node", ready_node) |
| 157 | @alignCast(@fieldParentPtr("queue_node", ready_node)) |
| 140 | 158 | else |
| 141 | 159 | break :ready_context current_idle_context; |
| 142 | 160 | break :ready_context &ready_fiber.context; |
| ... | ... | @@ -213,7 +231,7 @@ const SwitchMessage = extern struct { |
| 213 | 231 | register_awaiter: ?*?*Fiber, |
| 214 | 232 | |
| 215 | 233 | fn handle(message: *const SwitchMessage, el: *EventLoop) void { |
| 216 | | const prev_fiber: *Fiber = @fieldParentPtr("context", message.prev_context); |
| 234 | const prev_fiber: *Fiber = @alignCast(@fieldParentPtr("context", message.prev_context)); |
| 217 | 235 | current_fiber_context = message.ready_context; |
| 218 | 236 | if (message.register_awaiter) |awaiter| if (@atomicRmw(?*Fiber, awaiter, .Xchg, prev_fiber, .acq_rel) == Fiber.finished) el.schedule(prev_fiber); |
| 219 | 237 | } |
| ... | ... | @@ -279,17 +297,25 @@ fn fiberEntry() callconv(.naked) void { |
| 279 | 297 | |
| 280 | 298 | pub fn @"async"( |
| 281 | 299 | userdata: ?*anyopaque, |
| 282 | | eager_result: []u8, |
| 283 | | context: ?*anyopaque, |
| 284 | | start: *const fn (context: ?*anyopaque, result: *anyopaque) void, |
| 300 | result: []u8, |
| 301 | result_alignment: std.mem.Alignment, |
| 302 | context: []const u8, |
| 303 | context_alignment: std.mem.Alignment, |
| 304 | start: *const fn (context: *const anyopaque, result: *anyopaque) void, |
| 285 | 305 | ) ?*std.Io.AnyFuture { |
| 306 | assert(result_alignment.compare(.lte, max_result_align)); // TODO |
| 307 | assert(context_alignment.compare(.lte, max_result_align)); // TODO |
| 308 | assert(result.len <= max_result_len); // TODO |
| 309 | assert(context.len <= max_result_len); // TODO |
| 310 | |
| 286 | 311 | const event_loop: *EventLoop = @alignCast(@ptrCast(userdata)); |
| 287 | | const fiber = event_loop.allocateFiber(eager_result.len) catch { |
| 288 | | start(context, eager_result.ptr); |
| 312 | const fiber = event_loop.allocateFiber() catch { |
| 313 | start(context.ptr, result.ptr); |
| 289 | 314 | return null; |
| 290 | 315 | }; |
| 291 | 316 | fiber.awaiter = null; |
| 292 | 317 | fiber.queue_node = .{ .data = {} }; |
| 318 | @memcpy(fiber.argsSlice()[0..context.len], context); |
| 293 | 319 | std.log.debug("allocated {*}", .{fiber}); |
| 294 | 320 | |
| 295 | 321 | const closure: *AsyncClosure = @ptrFromInt(std.mem.alignBackward( |
| ... | ... | @@ -299,7 +325,6 @@ pub fn @"async"( |
| 299 | 325 | )); |
| 300 | 326 | closure.* = .{ |
| 301 | 327 | .event_loop = event_loop, |
| 302 | | .context = context, |
| 303 | 328 | .fiber = fiber, |
| 304 | 329 | .start = start, |
| 305 | 330 | }; |
| ... | ... | @@ -316,14 +341,13 @@ pub fn @"async"( |
| 316 | 341 | |
| 317 | 342 | const AsyncClosure = struct { |
| 318 | 343 | event_loop: *EventLoop, |
| 319 | | context: ?*anyopaque, |
| 320 | 344 | fiber: *Fiber, |
| 321 | | start: *const fn (context: ?*anyopaque, result: *anyopaque) void, |
| 345 | start: *const fn (context: *const anyopaque, result: *anyopaque) void, |
| 322 | 346 | |
| 323 | 347 | fn call(closure: *AsyncClosure, message: *const SwitchMessage) callconv(.c) noreturn { |
| 324 | 348 | message.handle(closure.event_loop); |
| 325 | 349 | std.log.debug("{*} performing async", .{closure.fiber}); |
| 326 | | closure.start(closure.context, closure.fiber.resultSlice().ptr); |
| 350 | closure.start(closure.fiber.argsSlice().ptr, closure.fiber.resultSlice().ptr); |
| 327 | 351 | const awaiter = @atomicRmw(?*Fiber, &closure.fiber.awaiter, .Xchg, Fiber.finished, .acq_rel); |
| 328 | 352 | closure.event_loop.yield(awaiter, null); |
| 329 | 353 | unreachable; // switched to dead fiber |