feature. See also
. The project being documented here (as the example) is the Zig library itself.
Dispatch.Mutex
const Mutex = struct
File
Code
const Mutex = struct {
state: State,
queue: c.dispatch.queue_t,
waiters: std.DoublyLinkedList,
const State = packed struct(usize) {
locked: bool,
num_waiters: NumWaiters,
const NumWaiters = @Int(.unsigned, @bitSizeOf(usize) - 1);
};
const Waiter = struct {
sleeper: Sleeper = undefined,
cancelable: Cancelable,
mutex: *Mutex,
node: std.DoublyLinkedList.Node = undefined,
fn add(context: ?*anyopaque) callconv(.c) void {
const waiter: *Waiter = @ptrCast(@alignCast(context));
waiter.cancelable.enter(waiter.sleeper.fiber) catch |err| switch (err) {
error.CancelRequested => return waiter.wake(),
};
var state = @atomicRmw(State, &waiter.mutex.state, .Add, .{
.locked = false,
.num_waiters = 1,
}, .monotonic);
state.num_waiters += 1;
while (!state.locked) {
@branchHint(.unlikely);
state = @cmpxchgWeak(State, &waiter.mutex.state, state, .{
.locked = true,
.num_waiters = state.num_waiters - 1,
}, .acquire, .monotonic) orelse break;
} else return waiter.mutex.waiters.append(&waiter.node);
waiter.cancelable.leave(waiter.sleeper.fiber) catch |err| switch (err) {
error.CancelRequested => {
waiter.node.next = &waiter.node;
return;
},
};
waiter.wake();
}
fn canceled(context: ?*anyopaque) callconv(.c) void {
const cancelable: *Cancelable = @ptrCast(@alignCast(context));
const waiter: *Waiter = @fieldParentPtr("cancelable", cancelable);
cancelable.requested(waiter.sleeper.fiber);
const mutex = waiter.mutex;
if (waiter.node.next != &waiter.node) {
@branchHint(.likely);
mutex.waiters.remove(&waiter.node);
assert(@atomicRmw(State, &mutex.state, .Sub, .{
.locked = false,
.num_waiters = 1,
}, .monotonic).num_waiters >= 1);
}
waiter.node = undefined;
waiter.wake();
}
fn remove(context: ?*anyopaque) callconv(.c) void {
const mutex: *Mutex = @ptrCast(@alignCast(context));
var state = @atomicLoad(State, &mutex.state, .monotonic);
while (!state.locked and state.num_waiters > 0) {
@branchHint(.likely);
state = @cmpxchgWeak(State, &mutex.state, state, .{
.locked = true,
.num_waiters = state.num_waiters - 1,
}, .acquire, .monotonic) orelse break;
} else return;
var num_removed: State.NumWaiters = 0;
while (mutex.waiters.popFirst()) |node| {
@branchHint(.likely);
const waiter: *Waiter = @fieldParentPtr("node", node);
node.* = undefined;
waiter.cancelable.leave(waiter.sleeper.fiber) catch |err| switch (err) {
error.CancelRequested => {
num_removed += 1;
node.next = node;
continue;
},
};
break;
}
if (num_removed > 0) {
@branchHint(.unlikely);
assert(@atomicRmw(State, &mutex.state, .Sub, .{
.locked = false,
.num_waiters = num_removed,
}, .monotonic).num_waiters >= num_removed);
}
}
fn wake(waiter: *Waiter) void {
Sleeper.wake(&waiter.sleeper);
}
};
fn init(mutex: *Mutex, queue: c.dispatch.queue_t) error{SystemResources}!void {
mutex.* = .{
.state = .{ .locked = false, .num_waiters = 0 },
.queue = c.dispatch.queue_create_with_target(
"org.ziglang.std.Io.Dispatch.Mutex",
.SERIAL(),
queue,
) orelse return error.SystemResources,
.waiters = .{},
};
}
fn deinit(mutex: *Mutex) void {
assert(mutex.state == State{ .locked = false, .num_waiters = 0 });
assert(mutex.waiters.first == null and mutex.waiters.last == null);
mutex.queue.as_object().release();
mutex.* = undefined;
}
fn tryLock(mutex: *Mutex) bool {
const state =
@atomicRmw(State, &mutex.state, .Or, .{ .locked = true, .num_waiters = 0 }, .acquire);
if (state.locked) {
@branchHint(.unlikely);
}
return !state.locked;
}
fn lock(mutex: *Mutex, ev: *Evented) Io.Cancelable!void {
if (mutex.tryLock()) return;
var waiter: Waiter = .{
.cancelable = .{ .queue = mutex.queue, .cancel = &Mutex.Waiter.canceled },
.mutex = mutex,
};
ev.yield(.{ .mutex_wait = &waiter });
try waiter.cancelable.acknowledge(waiter.sleeper.fiber);
}
fn lockUncancelable(mutex: *Mutex, ev: *Evented) void {
if (mutex.tryLock()) return;
var waiter: Waiter = .{ .cancelable = .blocked, .mutex = mutex };
ev.yield(.{ .mutex_wait = &waiter });
waiter.cancelable.acknowledge(waiter.sleeper.fiber) catch |err| switch (err) {
error.Canceled => unreachable,
};
}
fn unlock(mutex: *Mutex) void {
const state = @atomicRmw(State, &mutex.state, .And, .{
.locked = false,
.num_waiters = std.math.maxInt(State.NumWaiters),
}, .release);
if (state.num_waiters > 0) {
@branchHint(.unlikely);
mutex.queue.async(mutex, &Waiter.remove);
}
}
}