| ... | @@ -5,6 +5,7 @@ const Io = std.Io; | ... | @@ -5,6 +5,7 @@ const Io = std.Io; |
| 5 | const EventLoop = @This(); | 5 | const EventLoop = @This(); |
| 6 | | 6 | |
| 7 | gpa: Allocator, | 7 | gpa: Allocator, |
| | 8 | mutex: std.Thread.Mutex, |
| 8 | queue: std.DoublyLinkedList(void), | 9 | queue: std.DoublyLinkedList(void), |
| 9 | free: std.DoublyLinkedList(void), | 10 | free: std.DoublyLinkedList(void), |
| 10 | main_fiber_buffer: [@sizeOf(Fiber) + max_result_len]u8 align(@alignOf(Fiber)), | 11 | main_fiber_buffer: [@sizeOf(Fiber) + max_result_len]u8 align(@alignOf(Fiber)), |
| ... | @@ -39,6 +40,7 @@ const Fiber = struct { | ... | @@ -39,6 +40,7 @@ const Fiber = struct { |
| 39 | pub fn init(el: *EventLoop, gpa: Allocator) void { | 40 | pub fn init(el: *EventLoop, gpa: Allocator) void { |
| 40 | el.* = .{ | 41 | el.* = .{ |
| 41 | .gpa = gpa, | 42 | .gpa = gpa, |
| | 43 | .mutex = .{}, |
| 42 | .queue = .{}, | 44 | .queue = .{}, |
| 43 | .free = .{}, | 45 | .free = .{}, |
| 44 | .main_fiber_buffer = undefined, | 46 | .main_fiber_buffer = undefined, |
| ... | @@ -48,7 +50,11 @@ pub fn init(el: *EventLoop, gpa: Allocator) void { | ... | @@ -48,7 +50,11 @@ pub fn init(el: *EventLoop, gpa: Allocator) void { |
| 48 | | 50 | |
| 49 | fn allocateFiber(el: *EventLoop, result_len: usize) error{OutOfMemory}!*Fiber { | 51 | fn allocateFiber(el: *EventLoop, result_len: usize) error{OutOfMemory}!*Fiber { |
| 50 | assert(result_len <= max_result_len); | 52 | assert(result_len <= max_result_len); |
| 51 | const free_node = el.free.pop() orelse { | 53 | const free_node = free_node: { |
| | 54 | el.mutex.lock(); |
| | 55 | defer el.mutex.unlock(); |
| | 56 | break :free_node el.free.pop(); |
| | 57 | } orelse { |
| 52 | const n = std.mem.alignForward( | 58 | const n = std.mem.alignForward( |
| 53 | usize, | 59 | usize, |
| 54 | @sizeOf(Fiber) + max_result_len + min_stack_size, | 60 | @sizeOf(Fiber) + max_result_len + min_stack_size, |
| ... | @@ -59,36 +65,48 @@ fn allocateFiber(el: *EventLoop, result_len: usize) error{OutOfMemory}!*Fiber { | ... | @@ -59,36 +65,48 @@ fn allocateFiber(el: *EventLoop, result_len: usize) error{OutOfMemory}!*Fiber { |
| 59 | return @fieldParentPtr("queue_node", free_node); | 65 | return @fieldParentPtr("queue_node", free_node); |
| 60 | } | 66 | } |
| 61 | | 67 | |
| 62 | fn yield(el: *EventLoop, optional_fiber: ?*Fiber) void { | 68 | fn yield(el: *EventLoop, optional_fiber: ?*Fiber, register_awaiter: ?*?*Fiber) void { |
| 63 | if (optional_fiber) |fiber| { | 69 | const message: SwitchMessage = .{ |
| 64 | const old = &current_fiber.regs; | 70 | .ready_fiber = optional_fiber orelse if (ready_node: { |
| 65 | current_fiber = fiber; | 71 | el.mutex.lock(); |
| 66 | contextSwitch(old, &fiber.regs); | 72 | defer el.mutex.unlock(); |
| 67 | return; | 73 | break :ready_node el.queue.pop(); |
| 68 | } | 74 | }) |ready_node| |
| 69 | if (el.queue.pop()) |node| { | 75 | @fieldParentPtr("queue_node", ready_node) |
| 70 | const fiber: *Fiber = @fieldParentPtr("queue_node", node); | 76 | else if (register_awaiter) |_| |
| 71 | const old = &current_fiber.regs; | 77 | @panic("no other fiber to switch to in order to be able to register this fiber as an awaiter") // time to switch to an idle fiber? |
| 72 | current_fiber = fiber; | 78 | else |
| 73 | contextSwitch(old, &fiber.regs); | 79 | return, // nothing to do |
| 74 | return; | 80 | .register_awaiter = register_awaiter, |
| 75 | } | 81 | }; |
| 76 | @panic("everything is done"); | 82 | std.log.debug("switching from {*} to {*}", .{ current_fiber, message.ready_fiber }); |
| | 83 | SwitchMessage.handle(@ptrFromInt(contextSwitch(&current_fiber.regs, &message.ready_fiber.regs, @intFromPtr(&message))), el); |
| 77 | } | 84 | } |
| 78 | | 85 | |
| 79 | /// Equivalent to calling `yield` and then giving the fiber back to the event loop. | 86 | const SwitchMessage = struct { |
| 80 | fn exit(el: *EventLoop, optional_fiber: ?*Fiber) noreturn { | 87 | ready_fiber: *Fiber, |
| 81 | yield(el, optional_fiber); | 88 | register_awaiter: ?*?*Fiber, |
| 82 | @panic("TODO recycle the fiber"); | 89 | |
| 83 | } | 90 | fn handle(message: *const SwitchMessage, el: *EventLoop) void { |
| | 91 | const prev_fiber = current_fiber; |
| | 92 | current_fiber = message.ready_fiber; |
| | 93 | if (message.register_awaiter) |awaiter| if (@atomicRmw(?*Fiber, awaiter, .Xchg, prev_fiber, .acq_rel) == Fiber.finished) el.schedule(prev_fiber); |
| | 94 | } |
| | 95 | }; |
| 84 | | 96 | |
| 85 | fn schedule(el: *EventLoop, fiber: *Fiber) void { | 97 | fn schedule(el: *EventLoop, fiber: *Fiber) void { |
| | 98 | el.mutex.lock(); |
| | 99 | defer el.mutex.unlock(); |
| 86 | el.queue.append(&fiber.queue_node); | 100 | el.queue.append(&fiber.queue_node); |
| 87 | } | 101 | } |
| 88 | | 102 | |
| 89 | fn myFiber(el: *EventLoop) *Fiber { | 103 | fn recycle(el: *EventLoop, fiber: *Fiber) void { |
| 90 | _ = el; | 104 | std.log.debug("recyling {*}", .{fiber}); |
| 91 | return current_fiber; | 105 | fiber.awaiter = undefined; |
| | 106 | @memset(fiber.resultPointer()[0..max_result_len], undefined); |
| | 107 | el.mutex.lock(); |
| | 108 | defer el.mutex.unlock(); |
| | 109 | el.free.append(&fiber.queue_node); |
| 92 | } | 110 | } |
| 93 | | 111 | |
| 94 | const Regs = extern struct { | 112 | const Regs = extern struct { |
| ... | @@ -101,7 +119,7 @@ const Regs = extern struct { | ... | @@ -101,7 +119,7 @@ const Regs = extern struct { |
| 101 | rbp: usize, | 119 | rbp: usize, |
| 102 | }; | 120 | }; |
| 103 | | 121 | |
| 104 | const contextSwitch: *const fn (old: *Regs, new: *Regs) callconv(.c) void = @ptrCast(&contextSwitch_naked); | 122 | const contextSwitch: *const fn (old: *Regs, new: *Regs, message: usize) callconv(.c) usize = @ptrCast(&contextSwitch_naked); |
| 105 | | 123 | |
| 106 | noinline fn contextSwitch_naked() callconv(.naked) void { | 124 | noinline fn contextSwitch_naked() callconv(.naked) void { |
| 107 | asm volatile ( | 125 | asm volatile ( |
| ... | @@ -121,6 +139,7 @@ noinline fn contextSwitch_naked() callconv(.naked) void { | ... | @@ -121,6 +139,7 @@ noinline fn contextSwitch_naked() callconv(.naked) void { |
| 121 | \\movq 0x28(%%rsi), %%rbx | 139 | \\movq 0x28(%%rsi), %%rbx |
| 122 | \\movq 0x30(%%rsi), %%rbp | 140 | \\movq 0x30(%%rsi), %%rbp |
| 123 | \\ | 141 | \\ |
| | 142 | \\movq %%rdx, %%rax |
| 124 | \\ret | 143 | \\ret |
| 125 | ); | 144 | ); |
| 126 | } | 145 | } |
| ... | @@ -128,6 +147,7 @@ noinline fn contextSwitch_naked() callconv(.naked) void { | ... | @@ -128,6 +147,7 @@ noinline fn contextSwitch_naked() callconv(.naked) void { |
| 128 | fn popRet() callconv(.naked) void { | 147 | fn popRet() callconv(.naked) void { |
| 129 | asm volatile ( | 148 | asm volatile ( |
| 130 | \\pop %%rdi | 149 | \\pop %%rdi |
| | 150 | \\movq %%rax, %%rsi |
| 131 | \\ret | 151 | \\ret |
| 132 | ); | 152 | ); |
| 133 | } | 153 | } |
| ... | @@ -145,6 +165,7 @@ pub fn @"async"( | ... | @@ -145,6 +165,7 @@ pub fn @"async"( |
| 145 | }; | 165 | }; |
| 146 | fiber.awaiter = null; | 166 | fiber.awaiter = null; |
| 147 | fiber.queue_node = .{ .data = {} }; | 167 | fiber.queue_node = .{ .data = {} }; |
| | 168 | std.log.debug("allocated {*}", .{fiber}); |
| 148 | | 169 | |
| 149 | const closure: *AsyncClosure = @ptrFromInt(std.mem.alignBackward( | 170 | const closure: *AsyncClosure = @ptrFromInt(std.mem.alignBackward( |
| 150 | usize, | 171 | usize, |
| ... | @@ -157,14 +178,16 @@ pub fn @"async"( | ... | @@ -157,14 +178,16 @@ pub fn @"async"( |
| 157 | .fiber = fiber, | 178 | .fiber = fiber, |
| 158 | .start = start, | 179 | .start = start, |
| 159 | }; | 180 | }; |
| 160 | const stack_end_ptr: [*]align(16) usize = @alignCast(@ptrCast(closure)); | 181 | const stack_end: [*]align(16) usize = @alignCast(@ptrCast(closure)); |
| 161 | (stack_end_ptr - 1)[0] = 0; | 182 | const stack_top = (stack_end - 4)[0..4]; |
| 162 | (stack_end_ptr - 2)[0] = @intFromPtr(&AsyncClosure.call); | 183 | stack_top.* = .{ |
| 163 | (stack_end_ptr - 3)[0] = @intFromPtr(closure); | 184 | @intFromPtr(&popRet), |
| 164 | (stack_end_ptr - 4)[0] = @intFromPtr(&popRet); | 185 | @intFromPtr(closure), |
| 165 | | 186 | @intFromPtr(&AsyncClosure.call), |
| | 187 | 0, |
| | 188 | }; |
| 166 | fiber.regs = .{ | 189 | fiber.regs = .{ |
| 167 | .rsp = @intFromPtr(stack_end_ptr - 4), | 190 | .rsp = @intFromPtr(stack_top), |
| 168 | .r15 = 0, | 191 | .r15 = 0, |
| 169 | .r14 = 0, | 192 | .r14 = 0, |
| 170 | .r13 = 0, | 193 | .r13 = 0, |
| ... | @@ -181,30 +204,24 @@ const AsyncClosure = struct { | ... | @@ -181,30 +204,24 @@ const AsyncClosure = struct { |
| 181 | _: void align(16) = {}, | 204 | _: void align(16) = {}, |
| 182 | event_loop: *EventLoop, | 205 | event_loop: *EventLoop, |
| 183 | context: ?*anyopaque, | 206 | context: ?*anyopaque, |
| 184 | fiber: *EventLoop.Fiber, | 207 | fiber: *Fiber, |
| 185 | start: *const fn (context: ?*anyopaque, result: *anyopaque) void, | 208 | start: *const fn (context: ?*anyopaque, result: *anyopaque) void, |
| 186 | | 209 | |
| 187 | fn call(closure: *AsyncClosure) callconv(.c) void { | 210 | fn call(closure: *AsyncClosure, message: *const SwitchMessage) callconv(.c) noreturn { |
| 188 | std.log.debug("wrap called in async", .{}); | 211 | message.handle(closure.event_loop); |
| | 212 | std.log.debug("{*} performing async", .{closure.fiber}); |
| 189 | closure.start(closure.context, closure.fiber.resultPointer()); | 213 | closure.start(closure.context, closure.fiber.resultPointer()); |
| 190 | const awaiter = @atomicRmw(?*EventLoop.Fiber, &closure.fiber.awaiter, .Xchg, EventLoop.Fiber.finished, .seq_cst); | 214 | const awaiter = @atomicRmw(?*Fiber, &closure.fiber.awaiter, .Xchg, Fiber.finished, .acq_rel); |
| 191 | closure.event_loop.exit(awaiter); | 215 | closure.event_loop.yield(awaiter, null); |
| | 216 | unreachable; // switched to dead fiber |
| 192 | } | 217 | } |
| 193 | }; | 218 | }; |
| 194 | | 219 | |
| 195 | pub fn @"await"(userdata: ?*anyopaque, any_future: *std.Io.AnyFuture, result: []u8) void { | 220 | pub fn @"await"(userdata: ?*anyopaque, any_future: *std.Io.AnyFuture, result: []u8) void { |
| 196 | const event_loop: *EventLoop = @alignCast(@ptrCast(userdata)); | 221 | const event_loop: *EventLoop = @alignCast(@ptrCast(userdata)); |
| 197 | const future_fiber: *EventLoop.Fiber = @alignCast(@ptrCast(any_future)); | 222 | const future_fiber: *Fiber = @alignCast(@ptrCast(any_future)); |
| 198 | const result_src = future_fiber.resultPointer()[0..result.len]; | 223 | const result_src = future_fiber.resultPointer()[0..result.len]; |
| 199 | const my_fiber = event_loop.myFiber(); | 224 | if (@atomicLoad(?*Fiber, &future_fiber.awaiter, .acquire) != Fiber.finished) event_loop.yield(null, &future_fiber.awaiter); |
| 200 | | | |
| 201 | const prev = @atomicRmw(?*EventLoop.Fiber, &future_fiber.awaiter, .Xchg, my_fiber, .seq_cst); | | |
| 202 | if (prev == EventLoop.Fiber.finished) { | | |
| 203 | @memcpy(result, result_src); | | |
| 204 | return; | | |
| 205 | } | | |
| 206 | event_loop.yield(prev); | | |
| 207 | // Resumed when the value is available. | | |
| 208 | std.log.debug("yield returned in await", .{}); | | |
| 209 | @memcpy(result, result_src); | 225 | @memcpy(result, result_src); |
| | 226 | event_loop.recycle(future_fiber); |
| 210 | } | 227 | } |