feature. See also
. The project being documented here (as the example) is the Zig library itself.
Io.TypeErasedQueue
pub const TypeErasedQueue = struct
File
Code
pub const TypeErasedQueue = struct {
mutex: Mutex,
closed: bool,
buffer: []u8,
start: usize,
len: usize,
putters: std.DoublyLinkedList,
getters: std.DoublyLinkedList,
const Put = struct {
remaining: []const u8,
needed: usize,
condition: Condition,
node: std.DoublyLinkedList.Node,
};
const Get = struct {
remaining: []u8,
needed: usize,
condition: Condition,
node: std.DoublyLinkedList.Node,
};
pub fn init(buffer: []u8) TypeErasedQueue {
return .{
.mutex = .init,
.closed = false,
.buffer = buffer,
.start = 0,
.len = 0,
.putters = .{},
.getters = .{},
};
}
pub fn close(q: *TypeErasedQueue, io: Io) void {
q.mutex.lockUncancelable(io);
defer q.mutex.unlock(io);
q.closed = true;
{
var it = q.getters.first;
while (it) |node| : (it = node.next) {
const getter: *Get = @alignCast(@fieldParentPtr("node", node));
getter.condition.signal(io);
}
}
{
var it = q.putters.first;
while (it) |node| : (it = node.next) {
const putter: *Put = @alignCast(@fieldParentPtr("node", node));
putter.condition.signal(io);
}
}
}
pub fn put(q: *TypeErasedQueue, io: Io, elements: []const u8, min: usize) (QueueClosedError || Cancelable)!usize {
assert(elements.len >= min);
if (elements.len == 0) return 0;
try q.mutex.lock(io);
defer q.mutex.unlock(io);
return q.putLocked(io, elements, min, false);
}
pub fn putUncancelable(q: *TypeErasedQueue, io: Io, elements: []const u8, min: usize) QueueClosedError!usize {
assert(elements.len >= min);
if (elements.len == 0) return 0;
q.mutex.lockUncancelable(io);
defer q.mutex.unlock(io);
return q.putLocked(io, elements, min, true) catch |err| switch (err) {
error.Canceled => unreachable,
error.Closed => |e| return e,
};
}
fn puttableSlice(q: *const TypeErasedQueue) ?[]u8 {
const unwrapped_index = q.start + q.len;
const wrapped_index, const overflow = @subWithOverflow(unwrapped_index, q.buffer.len);
const slice = switch (overflow) {
1 => q.buffer[unwrapped_index..],
0 => q.buffer[wrapped_index..q.start],
};
return if (slice.len > 0) slice else null;
}
fn putLocked(q: *TypeErasedQueue, io: Io, elements: []const u8, min: usize, uncancelable: bool) (QueueClosedError || Cancelable)!usize {
if (q.closed) return error.Closed;
// queue is empty do we start populating the buffer.
// The number of elements we add immediately, before possibly blocking.
var n: usize = 0;
while (q.getters.popFirst()) |getter_node| {
const getter: *Get = @alignCast(@fieldParentPtr("node", getter_node));
const copy_len = @min(getter.remaining.len, elements.len - n);
assert(copy_len > 0);
@memcpy(getter.remaining[0..copy_len], elements[n..][0..copy_len]);
getter.remaining = getter.remaining[copy_len..];
getter.needed -|= copy_len;
n += copy_len;
if (getter.needed == 0) {
getter.condition.signal(io);
} else {
assert(n == elements.len);
q.getters.prepend(getter_node);
}
if (n == elements.len) return elements.len;
}
while (q.puttableSlice()) |slice| {
const copy_len = @min(slice.len, elements.len - n);
assert(copy_len > 0);
@memcpy(slice[0..copy_len], elements[n..][0..copy_len]);
q.len += copy_len;
n += copy_len;
if (n == elements.len) return elements.len;
}
if (n >= min) return n;
var pending: Put = .{
.remaining = elements[n..],
.needed = min - n,
.condition = .init,
.node = .{},
};
q.putters.append(&pending.node);
defer if (pending.needed > 0) q.putters.remove(&pending.node);
while (pending.needed > 0 and !q.closed) {
if (uncancelable) {
pending.condition.waitUncancelable(io, &q.mutex);
continue;
}
pending.condition.wait(io, &q.mutex) catch |err| switch (err) {
error.Canceled => if (pending.remaining.len == elements.len) {
return error.Canceled;
} else {
io.recancel();
return elements.len - pending.remaining.len;
},
};
}
if (pending.remaining.len == elements.len) {
assert(q.closed);
return error.Closed;
}
return elements.len - pending.remaining.len;
}
pub fn get(q: *TypeErasedQueue, io: Io, buffer: []u8, min: usize) (QueueClosedError || Cancelable)!usize {
assert(buffer.len >= min);
if (buffer.len == 0) return 0;
try q.mutex.lock(io);
defer q.mutex.unlock(io);
return q.getLocked(io, buffer, min, false);
}
pub fn getUncancelable(q: *TypeErasedQueue, io: Io, buffer: []u8, min: usize) QueueClosedError!usize {
assert(buffer.len >= min);
if (buffer.len == 0) return 0;
q.mutex.lockUncancelable(io);
defer q.mutex.unlock(io);
return q.getLocked(io, buffer, min, true) catch |err| switch (err) {
error.Canceled => unreachable,
error.Closed => |e| return e,
};
}
fn gettableSlice(q: *const TypeErasedQueue) ?[]const u8 {
const overlong_slice = q.buffer[q.start..];
const slice = overlong_slice[0..@min(overlong_slice.len, q.len)];
return if (slice.len > 0) slice else null;
}
fn getLocked(q: *TypeErasedQueue, io: Io, buffer: []u8, min: usize, uncancelable: bool) (QueueClosedError || Cancelable)!usize {
// queued putters, then finally the ring buffer should be filled with
// data from putters so they can be resumed.
// The number of elements we received immediately, before possibly blocking.
var n: usize = 0;
while (q.gettableSlice()) |slice| {
const copy_len = @min(slice.len, buffer.len - n);
assert(copy_len > 0);
@memcpy(buffer[n..][0..copy_len], slice[0..copy_len]);
q.start += copy_len;
if (q.buffer.len - q.start == 0) q.start = 0;
q.len -= copy_len;
n += copy_len;
if (n == buffer.len) {
q.fillRingBufferFromPutters(io);
return buffer.len;
}
}
while (q.putters.popFirst()) |putter_node| {
const putter: *Put = @alignCast(@fieldParentPtr("node", putter_node));
const copy_len = @min(putter.remaining.len, buffer.len - n);
assert(copy_len > 0);
@memcpy(buffer[n..][0..copy_len], putter.remaining[0..copy_len]);
putter.remaining = putter.remaining[copy_len..];
putter.needed -|= copy_len;
n += copy_len;
if (putter.needed == 0) {
putter.condition.signal(io);
} else {
assert(n == buffer.len);
q.putters.prepend(putter_node);
}
if (n == buffer.len) {
q.fillRingBufferFromPutters(io);
return buffer.len;
}
}
// because we emptied the ring buffer *and* the putter queue!
// Don't block if we hit the min or if the queue is closed. Return how
// many elements we could get immediately, unless the queue was closed and
// empty, in which case report `error.Closed`.
if (n == 0 and q.closed) return error.Closed;
if (n >= min or q.closed) return n;
var pending: Get = .{
.remaining = buffer[n..],
.needed = min - n,
.condition = .init,
.node = .{},
};
q.getters.append(&pending.node);
defer if (pending.needed > 0) q.getters.remove(&pending.node);
while (pending.needed > 0 and !q.closed) {
if (uncancelable) {
pending.condition.waitUncancelable(io, &q.mutex);
continue;
}
pending.condition.wait(io, &q.mutex) catch |err| switch (err) {
error.Canceled => if (pending.remaining.len == buffer.len) {
return error.Canceled;
} else {
io.recancel();
return buffer.len - pending.remaining.len;
},
};
}
if (pending.remaining.len == buffer.len) {
assert(q.closed);
return error.Closed;
}
return buffer.len - pending.remaining.len;
}
fn fillRingBufferFromPutters(q: *TypeErasedQueue, io: Io) void {
while (q.putters.popFirst()) |putter_node| {
const putter: *Put = @alignCast(@fieldParentPtr("node", putter_node));
while (q.puttableSlice()) |slice| {
const copy_len = @min(slice.len, putter.remaining.len);
assert(copy_len > 0);
@memcpy(slice[0..copy_len], putter.remaining[0..copy_len]);
q.len += copy_len;
putter.remaining = putter.remaining[copy_len..];
putter.needed -|= copy_len;
if (putter.needed == 0) {
putter.condition.signal(io);
break;
}
} else {
q.putters.prepend(putter_node);
break;
}
}
}
}