feature. See also
. The project being documented here (as the example) is the Zig library itself.
Threaded.groupAwait
fn groupAwait(userdata: ?*anyopaque, type_erased: *Io.Group, initial_token: *anyopaque) Io.Cancelable!void
File
Code
fn groupAwait(userdata: ?*anyopaque, type_erased: *Io.Group, initial_token: *anyopaque) Io.Cancelable!void {
_ = initial_token;
if (builtin.single_threaded) unreachable;
const t: *Threaded = @ptrCast(@alignCast(userdata));
const g: Group = .{ .ptr = type_erased };
var num_completed: std.atomic.Value(u32) = .init(0);
g.awaiter().* = &num_completed;
const pre_await_status = g.status().fetchOr(.{
.num_running = 0,
.have_awaiter = true,
.canceled = false,
}, .acq_rel);
assert(!pre_await_status.have_awaiter);
assert(!pre_await_status.canceled);
if (pre_await_status.num_running == 0) {
// until we return, so we can access `g.status()` non-atomically.
g.status().raw.have_awaiter = false;
return;
}
while (Thread.futexWait(&num_completed.raw, 0, null)) {
switch (num_completed.load(.acquire)) {
0 => continue,
1 => break,
else => unreachable,
}
} else |err| switch (err) {
error.Canceled => {
const pre_cancel_status = g.status().fetchOr(.{
.num_running = 0,
.have_awaiter = false,
.canceled = true,
}, .acq_rel);
assert(pre_cancel_status.have_awaiter);
assert(!pre_cancel_status.canceled);
// because in that case the last member of the group is already trying to modify it.
// However, if we know everything is done, we *can* skip signaling blocked threads.
const skip_signals = pre_cancel_status.num_running == 0;
g.waitForCancelWithSignaling(t, &num_completed, skip_signals);
// we can access `g.status()` non-atomically.
g.status().raw.canceled = false;
g.status().raw.have_awaiter = false;
return error.Canceled;
},
}
// we can access `g.status()` non-atomically.
g.status().raw.have_awaiter = false;
}