authorgravatar for jacobly@ziglang.orgJacob Young <jacobly@ziglang.org> 2025-03-27 19:32:26-04:00
committergravatar for andrew@ziglang.orgAndrew Kelley <andrew@ziglang.org> 2025-10-02 16:30:59-07:00
log1e79f2c12fcb81bf8e8ce4bbfc3457dd56a3391e
tree056cc2ce8a1267650f2343b268ee16e3fe34e526
parent2c1ceb4c9ce1cbead655264879dd0cee16fa3fa4

EventLoop: implement main idle fiber


1 files changed, 95 insertions(+), 50 deletions(-)

lib/std/Io/EventLoop.zig+95-50
......@@ -11,19 +11,21 @@ cond: std.Thread.Condition,
1111queue: std.DoublyLinkedList(void),
1212free: std.DoublyLinkedList(void),
1313main_fiber_buffer: [@sizeOf(Fiber) + max_result_len]u8 align(@alignOf(Fiber)),
14exiting: bool,
14exit_awaiter: ?*Fiber,
1515idle_count: usize,
1616threads: std.ArrayListUnmanaged(Thread),
1717
18threadlocal var current_thread: *Thread = undefined;
19threadlocal var current_fiber: *Fiber = undefined;
18threadlocal var current_idle_context: *Context = undefined;
19threadlocal var current_fiber_context: *Context = undefined;
2020
2121const max_result_len = 64;
2222const min_stack_size = 4 * 1024 * 1024;
23const idle_stack_size = 32 * 1024;
24const stack_align = 16;
2325
2426const Thread = struct {
2527 thread: std.Thread,
26 idle_fiber: Fiber,
28 idle_context: Context,
2729};
2830
2931const Fiber = struct {
......@@ -54,6 +56,11 @@ const Fiber = struct {
5456};
5557
5658pub 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);
5764 el.* = .{
5865 .gpa = gpa,
5966 .mutex = .{},
......@@ -61,28 +68,37 @@ pub fn init(el: *EventLoop, gpa: Allocator) error{OutOfMemory}!void {
6168 .queue = .{},
6269 .free = .{},
6370 .main_fiber_buffer = undefined,
64 .exiting = false,
71 .exit_awaiter = null,
6572 .idle_count = 0,
66 .threads = try .initCapacity(gpa, @max(std.Thread.getCpuCount() catch 1, 1)),
73 .threads = .initBuffer(@ptrCast(allocated_slice[0..threads_bytes])),
6774 };
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;
7088}
7189
7290pub 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);
8093 while (el.free.pop()) |free_node| {
8194 const free_fiber: *Fiber = @fieldParentPtr("queue_node", free_node);
8295 el.gpa.free(free_fiber.allocatedSlice());
8396 }
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]);
86102}
87103
88104fn allocateFiber(el: *EventLoop, result_len: usize) error{OutOfMemory}!*Fiber {
......@@ -103,40 +119,44 @@ fn allocateFiber(el: *EventLoop, result_len: usize) error{OutOfMemory}!*Fiber {
103119}
104120
105121fn 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 };
114133 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,
117136 .register_awaiter = register_awaiter,
118137 };
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 });
123139 contextSwitch(&message).handle(el);
124140}
125141
126142fn 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;
138149 }
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 };
140160}
141161
142162fn recycle(el: *EventLoop, fiber: *Fiber) void {
......@@ -148,14 +168,28 @@ fn recycle(el: *EventLoop, fiber: *Fiber) void {
148168 el.free.append(&fiber.queue_node);
149169}
150170
171fn 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
151177fn 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
184fn idle(el: *EventLoop) *Fiber {
154185 while (true) {
155186 el.yield(null, null);
187 if (@atomicLoad(?*Fiber, &el.exit_awaiter, .acquire)) |exit_awaiter| {
188 el.cond.broadcast();
189 return exit_awaiter;
190 }
156191 el.mutex.lock();
157192 defer el.mutex.unlock();
158 if (el.exiting) return;
159193 el.idle_count += 1;
160194 defer el.idle_count -= 1;
161195 el.cond.wait(&el.mutex);
......@@ -169,7 +203,7 @@ const SwitchMessage = extern struct {
169203
170204 fn handle(message: *const SwitchMessage, el: *EventLoop) void {
171205 const prev_fiber: *Fiber = @fieldParentPtr("context", message.prev_context);
172 current_fiber = @fieldParentPtr("context", message.ready_context);
206 current_fiber_context = message.ready_context;
173207 if (message.register_awaiter) |awaiter| if (@atomicRmw(?*Fiber, awaiter, .Xchg, prev_fiber, .acq_rel) == Fiber.finished) el.schedule(prev_fiber);
174208 }
175209};
......@@ -208,6 +242,18 @@ inline fn contextSwitch(message: *const SwitchMessage) *const SwitchMessage {
208242 };
209243}
210244
245fn 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
211257fn fiberEntry() callconv(.naked) void {
212258 switch (builtin.cpu.arch) {
213259 .x86_64 => asm volatile (
......@@ -238,7 +284,7 @@ pub fn @"async"(
238284 const closure: *AsyncClosure = @ptrFromInt(std.mem.alignBackward(
239285 usize,
240286 @intFromPtr(fiber.stackEndPointer() - @sizeOf(AsyncClosure)),
241 @alignOf(AsyncClosure),
287 @max(@alignOf(AsyncClosure), stack_align),
242288 ));
243289 closure.* = .{
244290 .event_loop = event_loop,
......@@ -246,7 +292,7 @@ pub fn @"async"(
246292 .fiber = fiber,
247293 .start = start,
248294 };
249 const stack_end: [*]align(16) usize = @alignCast(@ptrCast(closure));
295 const stack_end: [*]align(stack_align) usize = @alignCast(@ptrCast(closure));
250296 fiber.context = .{
251297 .rsp = @intFromPtr(stack_end - 1),
252298 .rbp = 0,
......@@ -258,7 +304,6 @@ pub fn @"async"(
258304}
259305
260306const AsyncClosure = struct {
261 _: void align(16) = {},
262307 event_loop: *EventLoop,
263308 context: ?*anyopaque,
264309 fiber: *Fiber,