feature. See also
. The project being documented here (as the example) is the Zig library itself.
Dispatch.Group
const Group = struct
File
Code
const Group = struct {
ptr: *Io.Group,
const List = packed struct(usize) {
cancel_requested: bool,
awaiter_delayed: bool,
fibers: Fiber.PackedPtr,
};
fn listPtr(group: Group) *List {
return @ptrCast(&group.ptr.token);
}
const Mutex = packed struct(u32) {
locked: bool,
contended: bool,
shared2: u30,
};
fn mutexPtr(group: Group) *Group.Mutex {
return switch (comptime builtin.cpu.arch.endian()) {
.little => @ptrCast(&group.ptr.state),
.big => @ptrCast(@alignCast(
@as([*]u8, @ptrCast(&group.ptr.state)) + @sizeOf(usize) - @sizeOf(u32),
)),
};
}
const Awaiter = packed struct(usize) {
locked: bool,
contended: bool,
awaiter: Fiber.PackedPtr,
};
fn awaiterPtr(group: Group) *Awaiter {
return @ptrCast(&group.ptr.state);
}
fn lock(group: Group, ev: *Evented) void {
const mutex = group.mutexPtr();
{
const old_state = @atomicRmw(
Group.Mutex,
mutex,
.Or,
.{ .locked = true, .contended = false, .shared2 = 0 },
.acquire,
);
if (!old_state.locked) {
@branchHint(.likely);
return;
}
if (old_state.contended) {
futexWaitUncancelable(ev, @ptrCast(mutex), @bitCast(old_state));
}
}
while (true) {
var old_state = @atomicRmw(
Group.Mutex,
mutex,
.Or,
.{ .locked = true, .contended = true, .shared2 = 0 },
.acquire,
);
if (!old_state.locked) {
@branchHint(.likely);
return;
}
old_state.contended = true;
futexWaitUncancelable(ev, @ptrCast(mutex), @bitCast(old_state));
}
}
fn unlock(group: Group, ev: *Evented) void {
const mutex = group.mutexPtr();
const old_state = @atomicRmw(
Group.Mutex,
mutex,
.And,
.{ .locked = false, .contended = false, .shared2 = std.math.maxInt(u30) },
.release,
);
assert(old_state.locked);
if (old_state.contended) futexWake(ev, @ptrCast(mutex), 1);
}
fn addFiber(group: Group, ev: *Evented, fiber: *Fiber) void {
group.lock(ev);
defer group.unlock(ev);
const list_ptr = group.listPtr();
const list = @atomicLoad(List, list_ptr, .monotonic);
if (list.cancel_requested) fiber.cancel_status = .{ .requested = true, .awaiting = .nothing };
const old_head = list.fibers.unpack();
if (old_head) |head| head.link.group.prev = fiber;
fiber.link.group.next = old_head;
@atomicStore(List, list_ptr, .{
.cancel_requested = list.cancel_requested,
.awaiter_delayed = list.awaiter_delayed,
.fibers = .pack(fiber),
}, .monotonic);
}
fn removeFiber(group: Group, ev: *Evented, fiber: *Fiber) ?*Fiber {
group.lock(ev);
defer group.unlock(ev);
const list_ptr = group.listPtr();
const list = @atomicLoad(List, list_ptr, .monotonic);
if (fiber.link.group.next) |next| next.link.group.prev = fiber.link.group.prev;
if (fiber.link.group.prev) |prev| {
prev.link.group.next = fiber.link.group.next;
} else if (fiber.link.group.next) |new_head| {
@atomicStore(List, list_ptr, .{
.cancel_requested = list.cancel_requested,
.awaiter_delayed = list.awaiter_delayed,
.fibers = .pack(new_head),
}, .monotonic);
} else if (@atomicLoad(Awaiter, group.awaiterPtr(), .monotonic).awaiter.unpack()) |awaiter| {
if (!awaiter.cancel_status.changeAwaiting(.group, .nothing) or list.cancel_requested) {
@atomicStore(List, list_ptr, .{
.cancel_requested = false,
.awaiter_delayed = false,
.fibers = .null,
}, .release);
assert(awaiter.awaiting_group.ptr == group.ptr);
awaiter.awaiting_group = undefined;
return awaiter;
}
@atomicStore(List, list_ptr, .{
.cancel_requested = false,
.awaiter_delayed = true,
.fibers = .null,
}, .monotonic);
} else @atomicStore(List, list_ptr, .{
.cancel_requested = false,
.awaiter_delayed = false,
.fibers = .null,
}, .release);
return null;
}
fn await(group: Group, ev: *Evented, awaiter: *Fiber) bool {
group.lock(ev);
defer group.unlock(ev);
if (@atomicLoad(List, group.listPtr(), .monotonic).fibers.unpack()) |_| {
if (group.registerAwaiter(awaiter) and awaiter.cancel_protection.check() == .unblocked) {
// attempting to await a group, so propagate the cancelation to the group.
assert(!group.cancelLocked(ev, null));
}
return false;
}
return true;
}
fn cancel(group: Group, ev: *Evented, maybe_awaiter: ?*Fiber) bool {
group.lock(ev);
defer group.unlock(ev);
return group.cancelLocked(ev, maybe_awaiter);
}
fn cancelLocked(group: Group, ev: *Evented, maybe_awaiter: ?*Fiber) bool {
const list_ptr = group.listPtr();
const list = @atomicRmw(
List,
list_ptr,
.Add,
.{ .cancel_requested = true, .awaiter_delayed = false, .fibers = .null },
.monotonic,
);
assert(!list.cancel_requested);
if (list.fibers.unpack()) |head| {
var maybe_fiber: ?*Fiber = head;
while (maybe_fiber) |fiber| {
fiber.requestCancel(ev);
maybe_fiber = fiber.link.group.next;
}
if (maybe_awaiter) |awaiter| _ = group.registerAwaiter(awaiter);
return false;
}
@atomicStore(
List,
list_ptr,
.{ .cancel_requested = false, .awaiter_delayed = false, .fibers = .null },
.release,
);
return if (maybe_awaiter) |_| true else list.awaiter_delayed;
}
fn registerAwaiter(group: Group, awaiter: *Fiber) bool {
awaiter.awaiting_group = group;
assert(@atomicRmw(
Awaiter,
group.awaiterPtr(),
.Add,
.{ .locked = false, .contended = false, .awaiter = .pack(awaiter) },
.monotonic,
).awaiter == .null);
return awaiter.cancel_status.changeAwaiting(.nothing, .group);
}
const AsyncClosure = struct {
evented: *Evented,
group: Group,
fiber: *Fiber,
start: *const fn (context: *const anyopaque) void,
fn fromFiber(fiber: *Fiber) *Group.AsyncClosure {
return @ptrFromInt(Fiber.max_context_align.max(.of(Group.AsyncClosure)).backward(
@intFromPtr(fiber.allocatedEnd()) - Fiber.max_context_size,
) - @sizeOf(Group.AsyncClosure));
}
fn contextPointer(
closure: *Group.AsyncClosure,
) [*]align(Fiber.max_context_align.toByteUnits()) u8 {
return @alignCast(@as([*]u8, @ptrCast(closure)) + @sizeOf(Group.AsyncClosure));
}
fn entry() callconv(.naked) void {
switch (builtin.cpu.arch) {
.aarch64 => asm volatile (
\\ mov x0, sp
\\ b %[call]
:
: [call] "X" (&call),
),
.x86_64 => asm volatile (
\\ leaq 8(%%rsp), %%rdi
\\ jmp %[call:P]
:
: [call] "X" (&call),
),
else => |arch| @compileError("unimplemented architecture: " ++ @tagName(arch)),
}
}
fn call(
closure: *Group.AsyncClosure,
message: *const SwitchMessage,
) callconv(.withStackAlign(.c, @alignOf(Group.AsyncClosure))) noreturn {
const ev = closure.evented;
const fiber = closure.fiber;
message.handle(ev);
closure.start(closure.contextPointer());
if (closure.group.removeFiber(ev, fiber)) |awaiter| ev.queue.async(awaiter, &Fiber.@"resume");
ev.yield(.destroy);
unreachable;
}
};
}