const std = @import("std"); const Io = std.Io; const Allocator = std.mem.Allocator; pub fn Queue(comptime T: type) type { return struct { const Self = @This(); allocator: Allocator, io: Io, buffer: []T, head: usize = 0, tail: usize = 0, len: usize = 0, closed: bool = false, mutex: Io.Mutex = .init, not_empty: Io.Condition = .init, not_full: Io.Condition = .init, pub const Error = error{Closed} || Allocator.Error; pub fn init( allocator: Allocator, io: Io, capacity: usize, ) Allocator.Error!Self { std.debug.assert(capacity > 0); return .{ .allocator = allocator, .io = io, .buffer = try allocator.alloc(T, capacity), }; } pub fn deinit(self: *Self) void { self.allocator.free(self.buffer); self.* = undefined; } pub fn close(self: *Self) void { self.mutex.lockUncancelable(self.io); defer self.mutex.unlock(self.io); self.closed = true; self.not_empty.broadcast(self.io); self.not_full.broadcast(self.io); } /// Blocks until space is available or the queue is closed. pub fn push(self: *Self, item: T) error{Closed}!void { self.mutex.lockUncancelable(self.io); defer self.mutex.unlock(self.io); while (self.len == self.buffer.len and !self.closed) { self.not_full.waitUncancelable(self.io, &self.mutex); } if (self.closed) return error.Closed; self.buffer[self.tail] = item; self.tail = (self.tail + 1) % self.buffer.len; self.len += 1; self.not_empty.signal(self.io); } /// Blocks until an item is available. Returns null only if the queue /// is closed and drained. pub fn pop(self: *Self) ?T { self.mutex.lockUncancelable(self.io); defer self.mutex.unlock(self.io); while (self.len == 0 and !self.closed) { self.not_empty.waitUncancelable(self.io, &self.mutex); } if (self.len == 0) return null; const item = self.buffer[self.head]; self.head = (self.head + 1) % self.buffer.len; self.len -= 1; self.not_full.signal(self.io); return item; } }; }