feature. See also
. The project being documented here (as the example) is the Zig library itself.
Dispatch.Futex
const Futex = struct
File
Code
const Futex = struct {
num_waiters: usize,
queue: c.dispatch.queue_t,
waiters: std.DoublyLinkedList,
const Waiter = struct {
sleeper: Sleeper = undefined,
cancelable: Cancelable,
futex: *Futex,
node: std.DoublyLinkedList.Node = .{},
ptr: *const u32,
expected: u32,
timeout: c.dispatch.time_t = .FOREVER,
leeway: u64,
timer: ?c.dispatch.source_t = null,
const already_signaled: c.dispatch.source_t = @ptrFromInt(1);
fn add(context: ?*anyopaque) callconv(.c) void {
const waiter: *Waiter = @ptrCast(@alignCast(context));
const futex = waiter.futex;
_ = @atomicRmw(usize, &futex.num_waiters, .Add, 1, .acquire);
waiter.tryAdd() catch |err| switch (err) {
error.CancelRequested => {
wake(waiter);
assert(@atomicRmw(usize, &futex.num_waiters, .Sub, 1, .monotonic) >= 1);
},
};
}
fn tryAdd(waiter: *Waiter) Cancelable.RequestedError!void {
if (@atomicLoad(u32, waiter.ptr, .monotonic) != waiter.expected)
return error.CancelRequested;
try waiter.cancelable.enter(waiter.sleeper.fiber);
const futex = waiter.futex;
switch (waiter.timeout) {
.FOREVER => {},
else => |timeout| {
const timer = c.dispatch.source_create(.TIMER, 0, .none, futex.queue) orelse {
log.warn("failed to create timer for futex timeout", .{});
return error.CancelRequested;
};
timer.as_object().set_context(waiter);
timer.set_event_handler(&timedOut);
timer.set_cancel_handler(&wake);
timer.set_timer(timeout, c.dispatch.TIME_FOREVER, waiter.leeway);
timer.as_object().activate();
waiter.timer = timer;
},
}
futex.waiters.append(&waiter.node);
}
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 futex = waiter.futex;
waiter.remove();
assert(@atomicRmw(usize, &futex.num_waiters, .Sub, 1, .monotonic) >= 1);
}
fn timedOut(context: ?*anyopaque) callconv(.c) void {
const waiter: *Waiter = @ptrCast(@alignCast(context));
const futex = waiter.futex;
waiter.tryRemove() catch |err| switch (err) {
error.CancelRequested => return,
};
assert(@atomicRmw(usize, &futex.num_waiters, .Sub, 1, .monotonic) >= 1);
}
fn tryRemove(waiter: *Waiter) Cancelable.RequestedError!void {
try waiter.cancelable.leave(waiter.sleeper.fiber);
waiter.remove();
}
fn remove(waiter: *Waiter) void {
waiter.futex.waiters.remove(&waiter.node);
if (waiter.timer) |timer| timer.cancel() else wake(waiter);
}
fn wake(context: ?*anyopaque) callconv(.c) void {
const waiter: *Waiter = @ptrCast(@alignCast(context));
if (waiter.timer) |timer| timer.as_object().release();
Sleeper.wake(&waiter.sleeper);
}
};
const Waker = struct {
sleeper: Sleeper = undefined,
futex: *Futex,
ptr: *const u32,
max_waiters: u32,
fn remove(context: ?*anyopaque) callconv(.c) void {
const waker: *Waker = @ptrCast(@alignCast(context));
const futex = waker.futex;
const ptr = waker.ptr;
const max_waiters = waker.max_waiters;
var num_removed: usize = 0;
var next_node = futex.waiters.first;
while (num_removed < max_waiters) {
const waiter: *Waiter = @fieldParentPtr("node", next_node orelse break);
next_node = waiter.node.next;
if (waiter.ptr != ptr) {
@branchHint(.unlikely);
continue;
}
waiter.tryRemove() catch |err| switch (err) {
error.CancelRequested => continue,
};
num_removed += 1;
}
assert(@atomicRmw(usize, &futex.num_waiters, .Sub, num_removed, .monotonic) >= num_removed);
var sleeper = waker.sleeper;
waker.* = undefined;
Sleeper.wake(&sleeper);
}
};
fn init(futex: *Futex, queue: c.dispatch.queue_t) error{SystemResources}!void {
futex.* = .{
.num_waiters = 0,
.queue = c.dispatch.queue_create_with_target(
"org.ziglang.std.Io.Dispatch.Futex",
.SERIAL(),
queue,
) orelse return error.SystemResources,
.waiters = .{},
};
}
fn deinit(futex: *Futex) void {
assert(futex.num_waiters == 0 and futex.waiters.first == null and futex.waiters.last == null);
futex.queue.as_object().release();
futex.* = undefined;
}
}