authorgravatar for andrew@ziglang.orgAndrew Kelley <andrew@ziglang.org> 2025-03-27 23:59:35-07:00
committergravatar for andrew@ziglang.orgAndrew Kelley <andrew@ziglang.org> 2025-07-20 10:38:38-07:00
log1d7a69cb7de7c935525428565767ab50c1198488
tree83b1d0b44d1eb68e3b4fb2d0d8523aeb4a430744
parentad3c5f02925fe939cc54619c11f687c975dccb89

update threaded fibers impl to actually storing args

sorry, something still not working correctly

1 files changed, 42 insertions(+), 18 deletions(-)

lib/std/Io/EventLoop.zig+42-18
......@@ -18,7 +18,11 @@ threads: std.ArrayListUnmanaged(Thread),
1818threadlocal var current_idle_context: *Context = undefined;
1919threadlocal var current_fiber_context: *Context = undefined;
2020
21/// Also used for context.
2122const max_result_len = 64;
23/// Also used for context.
24const max_result_align: std.mem.Alignment = .@"16";
25
2226const min_stack_size = 4 * 1024 * 1024;
2327const idle_stack_size = 32 * 1024;
2428const stack_align = 16;
......@@ -29,6 +33,8 @@ const Thread = struct {
2933};
3034
3135const Fiber = struct {
36 _: void align(max_result_align.toByteUnits()) = {},
37
3238 context: Context,
3339 awaiter: ?*Fiber,
3440 queue_node: std.DoublyLinkedList(void).Node,
......@@ -39,14 +45,27 @@ const Fiber = struct {
3945 const base: [*]align(@alignOf(Fiber)) u8 = @ptrCast(f);
4046 return base[0..std.mem.alignForward(
4147 usize,
42 @sizeOf(Fiber) + max_result_len + min_stack_size,
48 resultOffset() + max_result_len + min_stack_size,
4349 std.heap.page_size_max,
4450 )];
4551 }
4652
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
4766 fn resultSlice(f: *Fiber) []u8 {
4867 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];
5069 }
5170
5271 fn stackEndPointer(f: *Fiber) [*]u8 {
......@@ -102,7 +121,7 @@ pub fn deinit(el: *EventLoop) void {
102121 assert(el.queue.len == 0); // pending async
103122 el.yield(null, &el.exit_awaiter);
104123 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));
106125 el.gpa.free(free_fiber.allocatedSlice());
107126 }
108127 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 {
112131 el.gpa.free(allocated_ptr[0..idle_stack_end]);
113132}
114133
115fn allocateFiber(el: *EventLoop, result_len: usize) error{OutOfMemory}!*Fiber {
116 assert(result_len <= max_result_len);
134fn allocateFiber(el: *EventLoop) error{OutOfMemory}!*Fiber {
117135 const free_node = free_node: {
118136 el.mutex.lock();
119137 defer el.mutex.unlock();
......@@ -121,12 +139,12 @@ fn allocateFiber(el: *EventLoop, result_len: usize) error{OutOfMemory}!*Fiber {
121139 } orelse {
122140 const n = std.mem.alignForward(
123141 usize,
124 @sizeOf(Fiber) + max_result_len + min_stack_size,
142 Fiber.resultOffset() + max_result_len + min_stack_size,
125143 std.heap.page_size_max,
126144 );
127145 return @alignCast(@ptrCast(try el.gpa.alignedAlloc(u8, @alignOf(Fiber), n)));
128146 };
129 return @fieldParentPtr("queue_node", free_node);
147 return @alignCast(@fieldParentPtr("queue_node", free_node));
130148}
131149
132150fn 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
136154 defer el.mutex.unlock();
137155 break :ready_node el.queue.pop();
138156 }) |ready_node|
139 @fieldParentPtr("queue_node", ready_node)
157 @alignCast(@fieldParentPtr("queue_node", ready_node))
140158 else
141159 break :ready_context current_idle_context;
142160 break :ready_context &ready_fiber.context;
......@@ -213,7 +231,7 @@ const SwitchMessage = extern struct {
213231 register_awaiter: ?*?*Fiber,
214232
215233 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));
217235 current_fiber_context = message.ready_context;
218236 if (message.register_awaiter) |awaiter| if (@atomicRmw(?*Fiber, awaiter, .Xchg, prev_fiber, .acq_rel) == Fiber.finished) el.schedule(prev_fiber);
219237 }
......@@ -279,17 +297,25 @@ fn fiberEntry() callconv(.naked) void {
279297
280298pub fn @"async"(
281299 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,
285305) ?*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
286311 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);
289314 return null;
290315 };
291316 fiber.awaiter = null;
292317 fiber.queue_node = .{ .data = {} };
318 @memcpy(fiber.argsSlice()[0..context.len], context);
293319 std.log.debug("allocated {*}", .{fiber});
294320
295321 const closure: *AsyncClosure = @ptrFromInt(std.mem.alignBackward(
......@@ -299,7 +325,6 @@ pub fn @"async"(
299325 ));
300326 closure.* = .{
301327 .event_loop = event_loop,
302 .context = context,
303328 .fiber = fiber,
304329 .start = start,
305330 };
......@@ -316,14 +341,13 @@ pub fn @"async"(
316341
317342const AsyncClosure = struct {
318343 event_loop: *EventLoop,
319 context: ?*anyopaque,
320344 fiber: *Fiber,
321 start: *const fn (context: ?*anyopaque, result: *anyopaque) void,
345 start: *const fn (context: *const anyopaque, result: *anyopaque) void,
322346
323347 fn call(closure: *AsyncClosure, message: *const SwitchMessage) callconv(.c) noreturn {
324348 message.handle(closure.event_loop);
325349 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);
327351 const awaiter = @atomicRmw(?*Fiber, &closure.fiber.awaiter, .Xchg, Fiber.finished, .acq_rel);
328352 closure.event_loop.yield(awaiter, null);
329353 unreachable; // switched to dead fiber