From b6eb671ecdfbc02fbdbf45a384f38fbbceb5fac7 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Wed, 1 Apr 2026 00:28:14 -0500 Subject: [PATCH] =?UTF-8?q?0.16:=20migrate=20to=20std.Io=20=E2=80=94=20rep?= =?UTF-8?q?lace=20libc=20socket=20calls=20with=20Io=20vtable?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit All networking now goes through std.Io vtable (netRead/netWrite/netClose) instead of raw libc recv/send/close. Mutex/Condition use Io variants. PriorityDequeue, DebugAllocator, and Io.Clock APIs updated for 0.16. All 31 tests pass. Co-Authored-By: Claude Opus 4.6 --- build.zig.zon | 20 +- src/buffer.zig | 16 +- src/client/client.zig | 147 +++++------- src/server/handshake.zig | 20 +- src/server/server.zig | 329 +++++++++++--------------- src/server/thread_pool.zig | 60 ++--- src/t.zig | 126 ++++------ src/testing.zig | 33 +-- support/autobahn/client/build.zig.zon | 14 +- support/autobahn/client/main.zig | 18 +- support/autobahn/server/build.zig.zon | 14 +- support/autobahn/server/main.zig | 2 +- test_runner.zig | 42 ++-- 13 files changed, 350 insertions(+), 491 deletions(-) diff --git a/build.zig.zon b/build.zig.zon index da8b7d1..b795091 100644 --- a/build.zig.zon +++ b/build.zig.zon @@ -1,12 +1,12 @@ .{ - .name = .websocket, - .version = "0.1.0", - .dependencies = .{}, - .fingerprint = 0x42ce80b97512f264, - .paths = .{ - "readme.md", - "build.zig", - "build.zig.zon", - "src", - }, + .name = .websocket, + .version = "0.1.0", + .dependencies = .{}, + .fingerprint = 0x42ce80b97512f264, + .paths = .{ + "readme.md", + "build.zig", + "build.zig.zon", + "src", + }, } diff --git a/src/buffer.zig b/src/buffer.zig index d5dc77e..71e9896 100644 --- a/src/buffer.zig +++ b/src/buffer.zig @@ -215,7 +215,8 @@ pub const Pool = struct { available: usize, buffers: [][]u8, allocator: Allocator, - mutex: std.Thread.Mutex, + mutex: std.Io.Mutex, + io: std.Io, pub fn init(allocator: Allocator, count: usize, buffer_size: usize) !Pool { const buffers = try allocator.alloc([]u8, count); @@ -225,7 +226,8 @@ pub const Pool = struct { } return .{ - .mutex = .{}, + .mutex = .init, + .io = std.Options.debug_io, .buffers = buffers, .available = count, .allocator = allocator, @@ -244,8 +246,8 @@ pub const Pool = struct { pub fn acquire(self: *Pool) ?[]u8 { const buffers = self.buffers; - self.mutex.lock(); - defer self.mutex.unlock(); + self.mutex.lockUncancelable(self.io); + defer self.mutex.unlock(self.io); const available = self.available; if (available == 0) { return null; @@ -263,16 +265,16 @@ pub const Pool = struct { pub fn release(self: *Pool, buffer: []u8) void { var buffers = self.buffers; - self.mutex.lock(); + self.mutex.lockUncancelable(self.io); const available = self.available; if (available == buffers.len) { - self.mutex.unlock(); + self.mutex.unlock(self.io); self.allocator.free(buffer); return; } buffers[available] = buffer; self.available = available + 1; - self.mutex.unlock(); + self.mutex.unlock(self.io); } }; diff --git a/src/client/client.zig b/src/client/client.zig index b12f730..0f98385 100644 --- a/src/client/client.zig +++ b/src/client/client.zig @@ -3,7 +3,6 @@ const proto = @import("../proto.zig"); const buffer = @import("../buffer.zig"); const ascii = std.ascii; -const c = std.c; const Io = std.Io; const net = Io.net; const posix = std.posix; @@ -16,11 +15,9 @@ const Bundle = std.crypto.Certificate.Bundle; const CompressionOpts = @import("../websocket.zig").Compression; const ServerHandshake = @import("../server/handshake.zig").Handshake; -// 0.16: std.time.milliTimestamp() removed, use libc -fn milliTimestamp() i64 { - var tv: c.timeval = undefined; - _ = c.gettimeofday(&tv, null); - return @as(i64, tv.sec) * 1000 + @divTrunc(@as(i64, tv.usec), 1000); +fn milliTimestamp(io: Io) i64 { + const ts = Io.Timestamp.now(io, .real); + return @intCast(@divTrunc(ts.nanoseconds, std.time.ns_per_ms)); } fn ReadLoopHandler(comptime T: type) type { @@ -100,7 +97,7 @@ pub const Client = struct { var tls_client: ?*TLSClient = null; if (config.tls) { - tls_client = try TLSClient.init(allocator, net_stream, &config); + tls_client = try TLSClient.init(allocator, io, net_stream, &config); } const stream = Stream.init(io, net_stream, tls_client); @@ -451,9 +448,6 @@ pub const Stream = struct { io: Io, stream: net.Stream, tls_client: ?*TLSClient = null, - // 0.16: net.Stream needs buffers for Reader/Writer - read_buf: [8192]u8 = undefined, - write_buf: [8192]u8 = undefined, pub fn init(io: Io, stream: net.Stream, tls_client: ?*TLSClient) Stream { return .{ @@ -465,13 +459,10 @@ pub const Stream = struct { pub fn close(self: *Stream) void { if (self.tls_client) |tls_client| { - // Shutdown the socket first, so readLoop() can exit, before tls_client's buffers are freed self.stream.shutdown(self.io, .both) catch {}; tls_client.deinit(); } - - // 0.16: use libc.close directly (debug_io doesn't support real network) - _ = c.close(self.stream.socket.handle); + self.stream.close(self.io); } pub fn read(self: *Stream, buf: []u8) !usize { @@ -484,17 +475,14 @@ pub const Stream = struct { } } } - // 0.16: use libc.recv directly - const rc = c.recv(self.stream.socket.handle, buf.ptr, buf.len, 0); - if (rc == -1) { - const err = posix.errno(-1); + var bufs = [_][]u8{buf}; + return self.io.vtable.netRead(self.io.userdata, self.stream.socket.handle, &bufs) catch |err| { return switch (err) { - .CONNRESET => error.ConnectionResetByPeer, - .AGAIN => error.WouldBlock, + error.ConnectionResetByPeer => error.ConnectionResetByPeer, + error.Timeout => error.WouldBlock, else => error.Unexpected, }; - } - return @intCast(rc); + }; } pub fn writeAll(self: *Stream, data: []const u8) !void { @@ -504,19 +492,18 @@ pub const Stream = struct { try tls_client.stream_writer.interface.flush(); return; } - // 0.16: use libc.send directly var remaining = data; while (remaining.len > 0) { - const rc = c.send(self.stream.socket.handle, remaining.ptr, remaining.len, 0); - if (rc == -1) { - const err = posix.errno(-1); + // netWrite: header is sent first, data array's last element is the splat pattern. + // Pass remaining as header, empty pattern with splat=0. + const empty = [_][]const u8{""}; + const n = self.io.vtable.netWrite(self.io.userdata, self.stream.socket.handle, remaining, &empty, 0) catch |err| { return switch (err) { - .CONNRESET => error.ConnectionResetByPeer, - .PIPE => error.BrokenPipe, + error.ConnectionResetByPeer => error.ConnectionResetByPeer, else => error.Unexpected, }; - } - const n: usize = @intCast(rc); + }; + if (n == 0) return error.Unexpected; remaining = remaining[n..]; } } @@ -543,7 +530,6 @@ pub const Stream = struct { } pub fn setsockopt(self: *const Stream, opt_name: u32, value: []const u8) !void { - // 0.16: access handle through socket return posix.setsockopt(self.stream.socket.handle, posix.SOL.SOCKET, opt_name, value); } }; @@ -555,7 +541,7 @@ const TLSClient = struct { stream_reader: net.Stream.Reader, arena: std.heap.ArenaAllocator, - fn init(allocator: Allocator, stream: net.Stream, config: *const Client.Config) !*TLSClient { + fn init(allocator: Allocator, io: Io, stream: net.Stream, config: *const Client.Config) !*TLSClient { var arena = std.heap.ArenaAllocator.init(allocator); errdefer arena.deinit(); @@ -563,9 +549,7 @@ const TLSClient = struct { const bundle = config.ca_bundle orelse blk: { var b = Bundle{}; - // 0.16: rescan takes (allocator, io, timestamp) - const tls_io = std.Options.debug_io; - try b.rescan(aa, tls_io, Io.Timestamp.zero); + try b.rescan(aa, io, Io.Timestamp.zero); break :blk b; }; @@ -579,19 +563,16 @@ const TLSClient = struct { var buf = try aa.alloc(u8, buf_len * 4); const self = try aa.create(TLSClient); - // 0.16: writer/reader need io parameter - const tls_io = std.Options.debug_io; self.* = .{ .stream = stream, .arena = arena, .client = undefined, - .stream_writer = stream.writer(tls_io, buf.ptr[0..buf_len][0..buf_len]), - .stream_reader = stream.reader(tls_io, buf.ptr[buf_len .. 2 * buf_len][0..buf_len]), + .stream_writer = stream.writer(io, buf.ptr[0..buf_len][0..buf_len]), + .stream_reader = stream.reader(io, buf.ptr[buf_len .. 2 * buf_len][0..buf_len]), }; - // 0.16: interface is a field, not a function; also need entropy and time var entropy: [tls.Client.Options.entropy_len]u8 = undefined; - tls_io.random(&entropy); + io.random(&entropy); self.client = try tls.Client.init( &self.stream_reader.interface, &self.stream_writer.interface, @@ -601,7 +582,7 @@ const TLSClient = struct { .read_buffer = buf.ptr[2 * buf_len .. 3 * buf_len][0..buf_len], .write_buffer = buf.ptr[3 * buf_len .. 4 * buf_len][0..buf_len], .entropy = &entropy, - .realtime_now_seconds = @divTrunc(milliTimestamp(), 1000), + .realtime_now_seconds = @divTrunc(milliTimestamp(io), 1000), }, ); @@ -690,7 +671,7 @@ const HandShakeReply = struct { fn read(buf: []u8, key: []const u8, opts: *const Client.HandshakeOpts, compression: bool, stream: anytype) !HandShakeReply { const timeout_ms = opts.timeout_ms; - const deadline = milliTimestamp() + timeout_ms; + const deadline = milliTimestamp(stream.io) + timeout_ms; try stream.readTimeout(timeout_ms); var pos: usize = 0; @@ -799,7 +780,7 @@ const HandShakeReply = struct { } } - if (milliTimestamp() > deadline) { + if (milliTimestamp(stream.io) > deadline) { return error.Timeout; } @@ -861,7 +842,7 @@ test "Client: handshake" { defer pair.deinit(); try pair.clientWriteAll("\r\n\r\n"); - var client = testClient(pair.server); + var client = testClient(&pair); defer client.deinit(); try t.expectError(error.InvalidHandshakeResponse, client.handshake("/", .{})); } @@ -872,7 +853,7 @@ test "Client: handshake" { defer pair.deinit(); try pair.clientWriteAll("HTTP/1.1 200 OK\r\n\r\n"); - var client = testClient(pair.server); + var client = testClient(&pair); defer client.deinit(); try t.expectError(error.InvalidHandshakeResponse, client.handshake("/", .{})); } @@ -883,7 +864,7 @@ test "Client: handshake" { defer pair.deinit(); try pair.clientWriteAll("HTTP/1.1 101 Switching Protocol\r\n\r\n"); - var client = testClient(pair.server); + var client = testClient(&pair); defer client.deinit(); try t.expectError(error.InvalidHandshakeResponse, client.handshake("/", .{})); } @@ -894,7 +875,7 @@ test "Client: handshake" { defer pair.deinit(); try pair.clientWriteAll("HTTP/1.1 101 Switching Protocol\r\nUpgrade: nope\r\n\r\n"); - var client = testClient(pair.server); + var client = testClient(&pair); defer client.deinit(); try t.expectError(error.InvalidUpgradeHeader, client.handshake("/", .{})); } @@ -905,7 +886,7 @@ test "Client: handshake" { defer pair.deinit(); try pair.clientWriteAll("HTTP/1.1 101 Switching Protocol\r\nUpgrade: websocket\r\n\r\n"); - var client = testClient(pair.server); + var client = testClient(&pair); defer client.deinit(); try t.expectError(error.InvalidHandshakeResponse, client.handshake("/", .{})); } @@ -916,7 +897,7 @@ test "Client: handshake" { defer pair.deinit(); try pair.clientWriteAll("HTTP/1.1 101 Switching Protocol\r\nupgrade: WebSocket\r\nConnection: something\r\n\r\n"); - var client = testClient(pair.server); + var client = testClient(&pair); defer client.deinit(); try t.expectError(error.InvalidConnectionHeader, client.handshake("/", .{})); } @@ -927,7 +908,7 @@ test "Client: handshake" { defer pair.deinit(); try pair.clientWriteAll("HTTP/1.1 101 Switching Protocol\r\nUpgrade: websocket\r\nConnection: upgrade\r\n\r\n"); - var client = testClient(pair.server); + var client = testClient(&pair); defer client.deinit(); try t.expectError(error.InvalidHandshakeResponse, client.handshake("/", .{})); } @@ -938,7 +919,7 @@ test "Client: handshake" { defer pair.deinit(); try pair.clientWriteAll("HTTP/1.1 101 Switching Protocol\r\nupgrade: WebSocket\r\nConnection: UPGRADE\r\nSec-Websocket-Accept: hack\r\n\r\n"); - var client = testClient(pair.server); + var client = testClient(&pair); defer client.deinit(); try t.expectError(error.InvalidWebsocketAcceptHeader, client.handshake("/", .{})); } @@ -949,7 +930,7 @@ test "Client: handshake" { defer pair.deinit(); try pair.clientWriteAll("HTTP/1.1 101 Switching Protocol\r\nupgrade: WebSocket\r\nConnection: UPGRADE\r\nSec-Websocket-Accept: C/0nmHhBztSRGR1CwL6Tf4ZjwpY=\r\n\r\n"); - var client = testClient(pair.server); + var client = testClient(&pair); defer client.deinit(); try client.handshake("/", .{}); try t.expectEqual(0, client._reader.pos); @@ -961,7 +942,7 @@ test "Client: handshake" { defer pair.deinit(); try pair.clientWriteAll("HTTP/1.1 101 Switching Protocol\r\nupgrade: WebSocket\r\nConnection: UPGRADE\r\nSec-Websocket-Accept: C/0nmHhBztSRGR1CwL6Tf4ZjwpY=\r\n\r\nSome Random Data Which is Part Of the Next Message"); - var client = testClient(pair.server); + var client = testClient(&pair); defer client.deinit(); try client.handshake("/", .{}); try t.expectEqual(50, client._reader.pos); @@ -969,9 +950,9 @@ test "Client: handshake" { } test "Client: write/read" { - // 0.16: use initWithStream with libc-based connection - const stream = try testConnectedStream("127.0.0.1", 9292); - var client = try Client.initWithStream(std.Options.debug_io, t.allocator, stream, .{ + const io = std.Options.debug_io; + const stream = try testConnectedStream(io, "127.0.0.1", 9292); + var client = try Client.initWithStream(io, t.allocator, stream, .{ .port = 9292, .host = "127.0.0.1", }); @@ -993,9 +974,9 @@ test "Client: write/read" { } test "Client: close with code" { - // 0.16: use initWithStream with libc-based connection - const stream = try testConnectedStream("127.0.0.1", 9292); - var client = try Client.initWithStream(std.Options.debug_io, t.allocator, stream, .{ + const io = std.Options.debug_io; + const stream = try testConnectedStream(io, "127.0.0.1", 9292); + var client = try Client.initWithStream(io, t.allocator, stream, .{ .port = 9292, .host = "127.0.0.1", }); @@ -1009,9 +990,9 @@ test "Client: close with code" { } test "Client: with code and reason" { - // 0.16: use initWithStream with libc-based connection - const stream = try testConnectedStream("127.0.0.1", 9292); - var client = try Client.initWithStream(std.Options.debug_io, t.allocator, stream, .{ + const io = std.Options.debug_io; + const stream = try testConnectedStream(io, "127.0.0.1", 9292); + var client = try Client.initWithStream(io, t.allocator, stream, .{ .port = 9292, .host = "127.0.0.1", }); @@ -1057,7 +1038,9 @@ test "Client: Handler" { try t.expectEqual(true, h.closed); } -fn testClient(stream: net.Stream) Client { +fn testClient(pair: *t.SocketPair) Client { + pair.server_taken = true; + const stream = pair.server; const io = std.Options.debug_io; const bp = t.allocator.create(buffer.Provider) catch unreachable; bp.* = buffer.Provider.init(t.allocator, .{ .count = 0, .size = 0, .max = 4096 }) catch unreachable; @@ -1075,33 +1058,9 @@ fn testClient(stream: net.Stream) Client { }; } -// 0.16: Create a connected socket using libc (debug_io doesn't support real network) -fn testConnectedStream(host: []const u8, port: u16) !net.Stream { - const socket = c.socket(c.AF.INET, c.SOCK.STREAM, c.IPPROTO.TCP); - if (socket == -1) return error.SocketError; - errdefer _ = c.close(socket); - - // Setup sockaddr_in - var addr: posix.sockaddr.in = .{ - .port = @byteSwap(port), - .addr = 0, - }; - addr.family = c.AF.INET; - - // Parse simple IPv4 addresses - if (std.mem.eql(u8, host, "127.0.0.1")) { - addr.addr = 0x0100007f; - } else { - return error.UnsupportedHost; - } - - if (c.connect(socket, @ptrCast(&addr), @sizeOf(posix.sockaddr.in)) == -1) { - return error.ConnectionRefused; - } - - // Wrap in net.Stream - const loopback = net.Ip4Address.loopback(port); - return .{ .socket = .{ .handle = socket, .address = .{ .ip4 = loopback } } }; +fn testConnectedStream(io: Io, host: []const u8, port: u16) !net.Stream { + const addr = try net.IpAddress.parse(host, port); + return net.IpAddress.connect(&addr, io, .{ .mode = .stream }); } const ClientHandler = struct { @@ -1112,9 +1071,9 @@ const ClientHandler = struct { client: Client, fn init(allocator: Allocator) !ClientHandler { - // 0.16: use initWithStream with libc-based connection - const stream = try testConnectedStream("127.0.0.1", 9292); - var client = try Client.initWithStream(std.Options.debug_io, allocator, stream, .{ + const io = std.Options.debug_io; + const stream = try testConnectedStream(io, "127.0.0.1", 9292); + var client = try Client.initWithStream(io, allocator, stream, .{ .port = 9292, .host = "127.0.0.1", }); diff --git a/src/server/handshake.zig b/src/server/handshake.zig index a4c56a4..579a93a 100644 --- a/src/server/handshake.zig +++ b/src/server/handshake.zig @@ -368,7 +368,8 @@ pub const Handshake = struct { }; pub const Pool = struct { - mutex: std.Thread.Mutex, + mutex: std.Io.Mutex, + io: std.Io, available: usize, allocator: Allocator, buffer_size: usize, @@ -384,7 +385,8 @@ pub const Pool = struct { errdefer allocator.destroy(pool); pool.* = .{ - .mutex = .{}, + .mutex = .init, + .io = std.Options.debug_io, .states = states, .allocator = allocator, .available = count, @@ -416,12 +418,13 @@ pub const Pool = struct { pub fn acquire(self: *Pool) !*Handshake.State { const states = self.states; + const io = self.io; - self.mutex.lock(); + self.mutex.lockUncancelable(io); const available = self.available; if (available == 0) { // dont hold the lock over factory - self.mutex.unlock(); + self.mutex.unlock(io); const allocator = self.allocator; const state = try allocator.create(Handshake.State); @@ -432,24 +435,25 @@ pub const Pool = struct { const index = available - 1; const state = states[index]; self.available = index; - self.mutex.unlock(); + self.mutex.unlock(io); return state; } fn release(self: *Pool, state: *Handshake.State) void { var states = self.states; + const io = self.io; - self.mutex.lock(); + self.mutex.lockUncancelable(io); const available = self.available; if (available == states.len) { - self.mutex.unlock(); + self.mutex.unlock(io); state.deinit(); self.allocator.destroy(state); return; } states[available] = state; self.available = available + 1; - self.mutex.unlock(); + self.mutex.unlock(io); } }; diff --git a/src/server/server.zig b/src/server/server.zig index b599f10..9bf562a 100644 --- a/src/server/server.zig +++ b/src/server/server.zig @@ -88,40 +88,36 @@ const Address = struct { } }; -// 0.16: Io.net.Stream doesn't have read method, wrap socket for proto.Reader.fill +// wraps a socket handle for proto.Reader.fill, using Io vtable const SocketReader = struct { socket: posix.socket_t, + io: Io, pub fn read(self: SocketReader, buf: []u8) !usize { - const rc = libc.recv(self.socket, buf.ptr, buf.len, 0); - if (rc == -1) { - const err = posix.errno(-1); + var bufs = [_][]u8{buf}; + return self.io.vtable.netRead(self.io.userdata, self.socket, &bufs) catch |err| { return switch (err) { - .CONNRESET => error.ConnectionResetByPeer, - .PIPE => error.BrokenPipe, - .AGAIN => error.WouldBlock, + error.ConnectionResetByPeer => error.ConnectionResetByPeer, + error.Timeout => error.WouldBlock, else => error.Unexpected, }; - } - return @intCast(rc); + }; } }; -// 0.16: Io.net.Stream doesn't have writeAll method, use libc.send -fn socketWriteAll(socket: posix.socket_t, data: []const u8) !void { +fn socketWriteAll(io: Io, socket: posix.socket_t, data: []const u8) !void { var remaining = data; while (remaining.len > 0) { - const rc = libc.send(socket, remaining.ptr, remaining.len, 0); - if (rc == -1) { - const err = posix.errno(-1); + // netWrite: header is sent first, data array's last element is the splat pattern. + // Pass remaining as header, empty pattern with splat=0. + const empty = [_][]const u8{""}; + const n = io.vtable.netWrite(io.userdata, socket, remaining, &empty, 0) catch |err| { return switch (err) { - .CONNRESET => error.ConnectionResetByPeer, - .PIPE => error.BrokenPipe, - .AGAIN => error.WouldBlock, + error.ConnectionResetByPeer => error.ConnectionResetByPeer, else => error.Unexpected, }; - } - const n: usize = @intCast(rc); + }; + if (n == 0) return error.Unexpected; remaining = remaining[n..]; } } @@ -208,11 +204,12 @@ pub fn Server(comptime H: type) type { return struct { config: Config, allocator: Allocator, + io: Io, _state: WorkerState, _signals: []posix.fd_t, - _mut: Thread.Mutex, - _cond: Thread.Condition, + _mut: Io.Mutex, + _cond: Io.Condition, const Self = @This(); @@ -237,8 +234,9 @@ pub fn Server(comptime H: type) type { errdefer state.deinit(); return .{ - ._mut = .{}, - ._cond = .{}, + .io = std.Options.debug_io, + ._mut = .init, + ._cond = .init, ._state = state, ._signals = signals, .config = config, @@ -252,20 +250,22 @@ pub fn Server(comptime H: type) type { } pub fn listenInNewThread(self: *Self, ctx: anytype) !Thread { - self._mut.lock(); - defer self._mut.unlock(); + const io = self.io; + self._mut.lockUncancelable(io); + defer self._mut.unlock(io); const thrd = try Thread.spawn(.{}, Self.listen, .{ self, ctx }); // we don't return until listen() signals us that the server is up - self._cond.wait(&self._mut); + self._cond.waitUncancelable(io, &self._mut); return thrd; } pub fn listen(self: *Self, ctx: anytype) !void { - self._mut.lock(); + const io = self.io; + self._mut.lockUncancelable(io); errdefer { - self._cond.signal(); - self._mut.unlock(); + self._cond.signal(io); + self._mut.unlock(io); } const config = &self.config; @@ -329,7 +329,7 @@ pub fn Server(comptime H: type) type { const C = @TypeOf(ctx); if (comptime blockingMode()) { - errdefer posix.close(socket); + errdefer _ = libc.close(socket); var w = try Blocking(H).init(self.allocator, &self._state); defer w.deinit(); @@ -337,14 +337,14 @@ pub fn Server(comptime H: type) type { log.info("starting blocking worker to listen on {f}", .{address}); // incase listenInNewThread was used and is waiting for us to start - self._cond.signal(); + self._cond.signal(io); // this is what we'll shutdown when stop() is called self._signals[0] = socket; - self._mut.unlock(); + self._mut.unlock(io); thrd.join(); } else { - defer posix.close(socket); + defer _ = libc.close(socket); const W = NonBlocking(H, C); const allocator = self.allocator; @@ -358,7 +358,7 @@ pub fn Server(comptime H: type) type { errdefer for (0..started) |i| { // on success, these will be closed by a call to stop(); - posix.close(signals[i]); + _ = libc.close(signals[i]); }; defer { @@ -393,31 +393,28 @@ pub fn Server(comptime H: type) type { log.info("starting nonblocking worker to listen on {f}", .{address}); // in case startInNewThread is waiting - self._cond.signal(); + self._cond.signal(io); - self._mut.unlock(); + self._mut.unlock(io); for (threads) |thrd| { thrd.join(); } } - self._cond.signal(); + self._cond.signal(io); } pub fn stop(self: *Self) void { - self._mut.lock(); - defer self._mut.unlock(); + const io = self.io; + self._mut.lockUncancelable(io); + defer self._mut.unlock(io); for (self._signals) |s| { if (blockingMode()) { - // necessary to unblock accept on linux - // (which might not be that necessary since, on Linux, - // NonBlocking should be used) - // 0.16: posix.shutdown removed, use libc _ = libc.shutdown(s, 0); // SHUT_RD = 0 } - posix.close(s); + _ = libc.close(s); } - self._cond.wait(&self._mut); + self._cond.waitUncancelable(io, &self._mut); } }; } @@ -492,7 +489,7 @@ pub fn Blocking(comptime H: type) type { log.debug("({f}) connected", .{address}); const thread = std.Thread.spawn(.{}, Self.handleConnection, .{ self, socket, address, ctx }) catch |err| { - posix.close(socket); + _ = libc.close(socket); log.err("({f}) failed to spawn connection thread: {}", .{ address, err }); continue; }; @@ -584,7 +581,8 @@ pub fn Blocking(comptime H: type) type { // called for each hc when shutting down fn shutdownCleanup(_: *Self, hc: *HandlerConn(H)) void { - posix.shutdown(hc.socket, .recv) catch {}; + const io = std.Options.debug_io; + io.vtable.netShutdown(io.userdata, hc.socket, .receive) catch {}; } }; } @@ -713,8 +711,8 @@ fn NonBlocking(comptime H: type, comptime C: type) type { // always ordered from oldest to newest, so once we find a conneciton // that isn't timed out, we can stop - cm.lock.lock(); - defer cm.lock.unlock(); + cm.lock.lockUncancelable(cm.io); + defer cm.lock.unlock(cm.io); var next_conn = cm.pending.head; while (next_conn) |hc| { @@ -766,15 +764,15 @@ fn NonBlocking(comptime H: type, comptime C: type) type { log.debug("({f}) connected", .{address}); { - errdefer posix.close(socket); + errdefer _ = libc.close(socket); // socket is _probably_ in NONBLOCKING mode (it inherits // the flag from the listening socket). - const flags = try posix.fcntl(socket, posix.F.GETFL, 0); - const nonblocking = @as(u32, @bitCast(posix.O{ .NONBLOCK = true })); - if (flags & nonblocking == nonblocking) { + const flags = libc.fcntl(socket, libc.F.GETFL); + const O_NONBLOCK: c_int = 0x0004; + if (flags & O_NONBLOCK == O_NONBLOCK) { // Yup, it's in nonblocking mode. Disable that flag to // put it in blocking mode. - _ = try posix.fcntl(socket, posix.F.SETFL, flags & ~nonblocking); + _ = libc.fcntl(socket, libc.F.SETFL, flags & ~O_NONBLOCK); } } const hc = try self.base.newConn(socket, address, now); @@ -798,8 +796,8 @@ fn NonBlocking(comptime H: type, comptime C: type) type { fn dataAvailable(self: *Self, hc: *HandlerConn(H), thread_buf: []u8) void { var success = false; { - hc.cleanup.lock(); - defer hc.cleanup.unlock(); + hc.cleanup.lockUncancelable(hc.conn.io); + defer hc.cleanup.unlock(hc.conn.io); if (hc.handler == null) { success = self.dataForHandshake(hc) catch |err| blk: { log.err("({f}) error processing handshake: {}", .{ hc.conn.address, err }); @@ -930,8 +928,8 @@ fn NonBlockingBase(comptime H: type, comptime MANAGE_HS: bool) type { fn cleanupConn(self: *Self, hc: *HandlerConn(H)) void { { - hc.cleanup.lock(); - defer hc.cleanup.unlock(); + hc.cleanup.lockUncancelable(hc.conn.io); + defer hc.cleanup.unlock(hc.conn.io); if (hc.reader) |*reader| { if (self.small_buffer_pool) |*sbp| { sbp.release(reader.static); @@ -951,8 +949,8 @@ fn NonBlockingBase(comptime H: type, comptime MANAGE_HS: bool) type { // called for each hc when shutting down fn shutdownCleanup(self: *Self, hc: *HandlerConn(H)) void { - hc.cleanup.lock(); - defer hc.cleanup.unlock(); + hc.cleanup.lockUncancelable(hc.conn.io); + defer hc.cleanup.unlock(hc.conn.io); if (hc.reader) |*reader| { if (self.small_buffer_pool == null) { self.allocator.free(reader.static); @@ -1004,7 +1002,7 @@ const KQueue = struct { } fn deinit(self: KQueue) void { - posix.close(self.q); + _ = libc.close(self.q); } fn monitorAccept(self: *KQueue, fd: c_int) !void { @@ -1093,30 +1091,32 @@ const EPoll = struct { const EpollEvent = linux.epoll_event; fn init() !EPoll { + const q = linux.epoll_create1(0); + if (q == -1) return error.EpollError; return .{ .event_list = undefined, - .q = try posix.epoll_create1(0), + .q = q, }; } fn deinit(self: EPoll) void { - posix.close(self.q); + _ = libc.close(self.q); } fn monitorAccept(self: *EPoll, fd: c_int) !void { var event = linux.epoll_event{ .events = linux.EPOLL.IN, .data = .{ .ptr = 0 } }; - return std.posix.epoll_ctl(self.q, linux.EPOLL.CTL_ADD, fd, &event); + return linux.epoll_ctl(self.q, linux.EPOLL.CTL_ADD, fd, &event); } fn monitorSignal(self: *EPoll, fd: c_int) !void { var event = linux.epoll_event{ .events = linux.EPOLL.IN, .data = .{ .ptr = 1 } }; - return std.posix.epoll_ctl(self.q, linux.EPOLL.CTL_ADD, fd, &event); + return linux.epoll_ctl(self.q, linux.EPOLL.CTL_ADD, fd, &event); } fn monitorRead(self: *EPoll, hc: anytype, comptime rearm: bool) !void { const op = if (rearm) linux.EPOLL.CTL_MOD else linux.EPOLL.CTL_ADD; var event = linux.epoll_event{ .events = linux.EPOLL.IN | linux.EPOLL.ONESHOT, .data = .{ .ptr = @intFromPtr(hc) } }; - return posix.epoll_ctl(self.q, op, hc.socket, &event); + return linux.epoll_ctl(self.q, op, hc.socket, &event); } fn wait(self: *EPoll, timeout_sec: ?i32) !Iterator { @@ -1131,7 +1131,7 @@ const EPoll = struct { } } - const event_count = posix.epoll_wait(self.q, event_list, timeout); + const event_count = linux.epoll_wait(self.q, event_list, timeout); return .{ .index = 0, .events = event_list[0..event_count], @@ -1265,7 +1265,7 @@ pub fn HandlerConn(comptime H: type) type { reader: ?Reader, socket: posix.socket_t, // denormalization from conn.stream.socket.handle handshake: ?*Handshake.State, - cleanup: Thread.Mutex = .{}, + cleanup: Io.Mutex = .init, compression: ?Compression = null, next: ?*HandlerConn(H) = null, prev: ?*HandlerConn(H) = null, @@ -1279,7 +1279,8 @@ pub fn HandlerConn(comptime H: type) type { pub fn ConnManager(comptime H: type, comptime MANAGE_HS: bool) type { return struct { - lock: Thread.Mutex, + lock: Io.Mutex, + io: Io, allocator: Allocator, active: List(HandlerConn(H)), pending: List(HandlerConn(H)), @@ -1298,7 +1299,8 @@ pub fn ConnManager(comptime H: type, comptime MANAGE_HS: bool) type { errdefer compression_pool.deinit(allocator); return .{ - .lock = .{}, + .lock = .init, + .io = std.Options.debug_io, .pool = pool, .active = .{}, .pending = .{}, @@ -1315,8 +1317,8 @@ pub fn ConnManager(comptime H: type, comptime MANAGE_HS: bool) type { } pub fn count(self: *Self) usize { - self.lock.lock(); - defer self.lock.unlock(); + self.lock.lockUncancelable(self.io); + defer self.lock.unlock(self.io); if (MANAGE_HS == false) { return self.active.len; } @@ -1324,10 +1326,10 @@ pub fn ConnManager(comptime H: type, comptime MANAGE_HS: bool) type { } pub fn create(self: *Self, socket: posix.socket_t, address: Address, now: u32) !*HandlerConn(H) { - errdefer posix.close(socket); + errdefer _ = libc.close(socket); - self.lock.lock(); - defer self.lock.unlock(); + self.lock.lockUncancelable(self.io); + defer self.lock.unlock(self.io); const hc = try self.pool.create(self.allocator); hc.* = .{ @@ -1341,8 +1343,7 @@ pub fn ConnManager(comptime H: type, comptime MANAGE_HS: bool) type { ._closed = false, .started = now, .address = address, - // 0.16: Stream requires .socket.handle + .socket.address - // The actual client address is in address field; use loopback for stream + .io = self.io, .stream = .{ .socket = .{ .handle = socket, .address = .{ .ip4 = net.Ip4Address.loopback(0) } } }, .compression = null, }, @@ -1363,8 +1364,8 @@ pub fn ConnManager(comptime H: type, comptime MANAGE_HS: bool) type { // our caller made sute this was the case std.debug.assert(hc.state == .handshake); - self.lock.lock(); - defer self.lock.unlock(); + self.lock.lockUncancelable(self.io); + defer self.lock.unlock(self.io); self.pending.remove(hc); self.active.insert(hc); hc.state = .active; @@ -1378,8 +1379,8 @@ pub fn ConnManager(comptime H: type, comptime MANAGE_HS: bool) type { std.debug.assert(hc.state == .active); hc.state = .handshake; - self.lock.lock(); - defer self.lock.unlock(); + self.lock.lockUncancelable(self.io); + defer self.lock.unlock(self.io); self.active.remove(hc); self.pending.insert(hc); } @@ -1404,7 +1405,7 @@ pub fn ConnManager(comptime H: type, comptime MANAGE_HS: bool) type { c.writer.deinit(); } - self.lock.lock(); + self.lock.lockUncancelable(self.io); if (hc.state == .active) { self.active.remove(hc); } else { @@ -1416,7 +1417,7 @@ pub fn ConnManager(comptime H: type, comptime MANAGE_HS: bool) type { } self.pool.destroy(hc); - self.lock.unlock(); + self.lock.unlock(self.io); } fn setupCompression(self: *Self, hc: *HandlerConn(H)) !void { @@ -1446,8 +1447,8 @@ pub fn ConnManager(comptime H: type, comptime MANAGE_HS: bool) type { } pub fn shutdown(self: *Self, worker: anytype) void { - self.lock.lock(); - defer self.lock.unlock(); + self.lock.lockUncancelable(self.io); + defer self.lock.unlock(self.io); shutdownList(self.active.head, worker); shutdownList(self.pending.head, worker); @@ -1462,8 +1463,8 @@ pub fn ConnManager(comptime H: type, comptime MANAGE_HS: bool) type { var next_node = head; while (next_node) |hc| { if (comptime std.meta.hasFn(H, "close")) { - hc.cleanup.lock(); - defer hc.cleanup.unlock(); + hc.cleanup.lockUncancelable(hc.conn.io); + defer hc.cleanup.unlock(hc.conn.io); if (hc.handler) |*h| { h.close(); hc.handler = null; @@ -1485,7 +1486,8 @@ pub const Conn = struct { started: u32, stream: net.Stream, address: Address, - lock: Thread.Mutex = .{}, + io: Io, + lock: Io.Mutex = .init, compression: ?*Conn.Compression = null, const Compression = struct { @@ -1591,9 +1593,9 @@ pub const Conn = struct { if (payload.len == 0) { // no body, just write the header - self.lock.lock(); - defer self.lock.unlock(); - return socketWriteAll(stream.socket.handle, header); + self.lock.lockUncancelable(self.io); + defer self.lock.unlock(self.io); + return socketWriteAll(self.io, stream.socket.handle, header); } var vec = [2]std.posix.iovec_const{ @@ -1605,38 +1607,32 @@ pub const Conn = struct { } pub fn writeFramed(self: *Conn, data: []const u8) !void { - self.lock.lock(); - defer self.lock.unlock(); - try socketWriteAll(self.stream.socket.handle, data); + self.lock.lockUncancelable(self.io); + defer self.lock.unlock(self.io); + try socketWriteAll(self.io, self.stream.socket.handle, data); } fn writeAllIOVec(self: *Conn, vec: []std.posix.iovec_const) !void { const socket = self.stream.socket.handle; - - self.lock.lock(); - defer self.lock.unlock(); - - var i: usize = 0; - while (true) { - // 0.16: posix.writev removed, use libc.writev - const remaining = vec[i..]; - const rc = libc.writev(socket, @ptrCast(remaining.ptr), @intCast(remaining.len)); - if (rc == -1) { - const err = posix.errno(-1); - return switch (err) { - .CONNRESET => error.ConnectionResetByPeer, - .PIPE => error.BrokenPipe, - else => error.Unexpected, + const io = self.io; + + self.lock.lockUncancelable(io); + defer self.lock.unlock(io); + + // write each iovec segment via Io vtable + const empty = [_][]const u8{""}; + for (vec) |*v| { + var remaining = v.base[0..v.len]; + while (remaining.len > 0) { + const n = io.vtable.netWrite(io.userdata, socket, remaining, &empty, 0) catch |err| { + return switch (err) { + error.ConnectionResetByPeer => error.ConnectionResetByPeer, + else => error.Unexpected, + }; }; + if (n == 0) return error.Unexpected; + remaining = remaining[n..]; } - var n: usize = @intCast(rc); - while (n >= vec[i].len) { - n -= vec[i].len; - i += 1; - if (i >= vec.len) return; - } - vec[i].base += n; - vec[i].len -= n; } } @@ -1655,7 +1651,7 @@ pub const Conn = struct { fn closeSocket(self: *Conn) void { if (@atomicRmw(bool, &self._closed, .Xchg, true, .monotonic) == false) { - posix.close(self.stream.socket.handle); + self.io.vtable.netClose(self.io.userdata, (&self.stream.socket.handle)[0..1]); } } @@ -1710,9 +1706,9 @@ fn _handleHandshake(comptime H: type, worker: anytype, hc: *HandlerConn(H), ctx: return .{ false, false }; } - const n = posix.read(hc.socket, buf[len..]) catch |err| { + const n = (SocketReader{ .socket = hc.socket, .io = std.Options.debug_io }).read(buf[len..]) catch |err| { switch (err) { - error.BrokenPipe, error.ConnectionResetByPeer => log.debug("({f}) handshake connection closed: {}", .{ conn.address, err }), + error.ConnectionResetByPeer => log.debug("({f}) handshake connection closed: {}", .{ conn.address, err }), error.WouldBlock => { std.debug.assert(blockingMode()); log.debug("({f}) handshake timeout", .{conn.address}); @@ -1784,10 +1780,9 @@ fn handleClientData(comptime H: type, hc: *HandlerConn(H), allocator: Allocator, fn _handleClientData(comptime H: type, hc: *HandlerConn(H), allocator: Allocator, fba: *FixedBufferAllocator) !bool { var conn = &hc.conn; var reader = &hc.reader.?; - // 0.16: use SocketReader wrapper since Io.net.Stream doesn't have .read() - reader.fill(SocketReader{ .socket = conn.stream.socket.handle }) catch |err| { + reader.fill(SocketReader{ .socket = conn.stream.socket.handle, .io = conn.io }) catch |err| { switch (err) { - error.BrokenPipe, error.Closed, error.ConnectionResetByPeer => log.debug("({f}) connection closed: {}", .{ conn.address, err }), + error.Closed, error.ConnectionResetByPeer => log.debug("({f}) connection closed: {}", .{ conn.address, err }), else => log.warn("({f}) error reading from connection: {}", .{ conn.address, err }), } return false; @@ -1951,8 +1946,8 @@ fn preHandOffWrite(conn: *Conn, response: []const u8) void { // to *Conn yet. In theory, this means we don't need to worry about thread-safety // However, it is possible for the worker to be stopped while we're doing this // which causes issues unless we lock - conn.lock.lock(); - defer conn.lock.unlock(); + conn.lock.lockUncancelable(conn.io); + defer conn.lock.unlock(conn.io); if (conn.isClosed()) { return; @@ -1962,25 +1957,13 @@ fn preHandOffWrite(conn: *Conn, response: []const u8) void { const timeout = std.mem.toBytes(posix.timeval{ .sec = 5, .usec = 0 }); posix.setsockopt(socket, posix.SOL.SOCKET, posix.SO.SNDTIMEO, &timeout) catch return; - var pos: usize = 0; - while (pos < response.len) { - // 0.16: posix.write removed, use libc.write - const remaining = response[pos..]; - const rc = libc.write(socket, remaining.ptr, remaining.len); - if (rc <= 0) { - // error or closed - return; - } - pos += @intCast(rc); - } + socketWriteAll(conn.io, socket, response) catch return; } fn timestamp() u32 { - if (comptime @hasDecl(posix, "CLOCK") == false or posix.CLOCK == void) { - return @intCast(std.time.timestamp()); - } - const ts = posix.clock_gettime(posix.CLOCK.REALTIME) catch unreachable; - return @intCast(ts.sec); + const io = std.Options.debug_io; + const ts = Io.Timestamp.now(io, .real); + return @intCast(@divTrunc(ts.nanoseconds, std.time.ns_per_s)); } // intrusive doubly-linked list with count, not thread safe @@ -2028,7 +2011,7 @@ const t = @import("../t.zig"); var test_thread: Thread = undefined; var test_server: Server(TestHandler) = undefined; -var global_test_allocator = std.heap.GeneralPurposeAllocator(.{}){}; +var global_test_allocator: std.heap.DebugAllocator(.{}) = .init; test "tests:beforeAll" { test_server = try Server(TestHandler).init(global_test_allocator.allocator(), .{ @@ -2139,30 +2122,16 @@ test "Conn: close" { } } -// 0.16: Use libc directly (debug_io doesn't support real network I/O) const TestStream = struct { socket: posix.socket_t, + io: Io, fn close(self: *const TestStream) void { - _ = libc.close(self.socket); + self.io.vtable.netClose(self.io.userdata, (&self.socket)[0..1]); } fn writeAll(self: *const TestStream, data: []const u8) !void { - var remaining = data; - while (remaining.len > 0) { - const rc = libc.send(self.socket, remaining.ptr, remaining.len, 0); - if (rc == -1) { - const err = posix.errno(-1); - return switch (err) { - .CONNRESET => error.ConnectionResetByPeer, - .PIPE => error.BrokenPipe, - .AGAIN => error.WouldBlock, - else => error.Unexpected, - }; - } - const n: usize = @intCast(rc); - remaining = remaining[n..]; - } + try socketWriteAll(self.io, self.socket, data); } fn readAtLeast(self: *const TestStream, buf: []u8, min: usize) !usize { @@ -2176,44 +2145,26 @@ const TestStream = struct { } fn read(self: *const TestStream, buf: []u8) !usize { - const rc = libc.recv(self.socket, buf.ptr, buf.len, 0); - if (rc == -1) { - const err = posix.errno(-1); - return switch (err) { - .CONNRESET => error.ConnectionResetByPeer, - .AGAIN => error.WouldBlock, - else => error.Unexpected, - }; - } - return @intCast(rc); + return (SocketReader{ .socket = self.socket, .io = self.io }).read(buf); } }; fn testStream(handshake: bool) !TestStream { - const timeout = std.mem.toBytes(std.posix.timeval{ .sec = 0, .usec = 20_000 }); + const io = std.Options.debug_io; - // Connect using libc directly (debug_io doesn't support real network) - const socket = libc.socket(libc.AF.INET, libc.SOCK.STREAM, libc.IPPROTO.TCP); - if (socket == -1) return error.SocketError; - errdefer _ = libc.close(socket); + // Connect via Io.net + var addr = net.IpAddress.parse("127.0.0.1", 9292) catch unreachable; + const stream = try net.IpAddress.connect(&addr, io, .{ .mode = .stream }); - // Setup sockaddr_in for 127.0.0.1:9292 - var addr: posix.sockaddr.in = .{ - .port = @byteSwap(@as(u16, 9292)), - .addr = 0x0100007f, // 127.0.0.1 in network byte order - }; - addr.family = libc.AF.INET; - - if (libc.connect(socket, @ptrCast(&addr), @sizeOf(posix.sockaddr.in)) == -1) { - return error.ConnectionRefused; - } + const socket = stream.socket.handle; - // Set socket options using libc.setsockopt - _ = libc.setsockopt(socket, libc.IPPROTO.TCP, 1, &std.mem.toBytes(@as(c_int, 1)), @sizeOf(c_int)); - _ = libc.setsockopt(socket, libc.SOL.SOCKET, libc.SO.RCVTIMEO, &timeout, @sizeOf(@TypeOf(timeout))); - _ = libc.setsockopt(socket, libc.SOL.SOCKET, libc.SO.SNDTIMEO, &timeout, @sizeOf(@TypeOf(timeout))); + // Socket options (timeouts) — no Io equivalent, use posix.setsockopt + const timeout = std.mem.toBytes(posix.timeval{ .sec = 0, .usec = 20_000 }); + try posix.setsockopt(socket, posix.IPPROTO.TCP, 1, &std.mem.toBytes(@as(c_int, 1))); + try posix.setsockopt(socket, posix.SOL.SOCKET, posix.SO.RCVTIMEO, &timeout); + try posix.setsockopt(socket, posix.SOL.SOCKET, posix.SO.SNDTIMEO, &timeout); - var result = TestStream{ .socket = socket }; + var result = TestStream{ .socket = socket, .io = io }; if (handshake == false) { return result; diff --git a/src/server/thread_pool.zig b/src/server/thread_pool.zig index 974b712..4feb68d 100644 --- a/src/server/thread_pool.zig +++ b/src/server/thread_pool.zig @@ -1,10 +1,11 @@ const std = @import("std"); const c = std.c; +const Io = std.Io; const Thread = std.Thread; const Allocator = std.mem.Allocator; -// 0.16: std.Thread.sleep was removed, use libc nanosleep +// nanosleep via libc — no Io-based uncancelable sleep exists fn sleep(ns: u64) void { const secs = ns / std.time.ns_per_s; const nsecs = ns % std.time.ns_per_s; @@ -53,11 +54,12 @@ pub fn ThreadPool(comptime F: anytype) type { pending: usize, queue: []Args, threads: []Thread, - mutex: Thread.Mutex, - pull_cond: Thread.Condition, - push_cond: Thread.Condition, + mutex: Io.Mutex, + pull_cond: Io.Condition, + push_cond: Io.Condition, queue_end: usize, allocator: Allocator, + io: Io, const Self = @This(); @@ -71,15 +73,18 @@ pub fn ThreadPool(comptime F: anytype) type { const thread_pool = try allocator.create(Self); errdefer allocator.destroy(thread_pool); + const io = std.Options.debug_io; + thread_pool.* = .{ .pull = 0, .push = 0, .pending = 0, - .mutex = .{}, + .io = io, + .mutex = .init, .stopped = false, .queue = queue, - .pull_cond = .{}, - .push_cond = .{}, + .pull_cond = .init, + .push_cond = .init, .threads = threads, .allocator = allocator, .queue_end = queue.len - 1, @@ -88,7 +93,7 @@ pub fn ThreadPool(comptime F: anytype) type { var started: usize = 0; errdefer { thread_pool.stopped = true; - thread_pool.pull_cond.broadcast(); + thread_pool.pull_cond.broadcast(io); for (0..started) |i| { threads[i].join(); } @@ -116,70 +121,67 @@ pub fn ThreadPool(comptime F: anytype) type { } pub fn stop(self: *Self) void { + const io = self.io; { - self.mutex.lock(); - defer self.mutex.unlock(); + self.mutex.lockUncancelable(io); + defer self.mutex.unlock(io); if (self.stopped == true) { return; } self.stopped = true; } - self.pull_cond.broadcast(); + self.pull_cond.broadcast(io); for (self.threads) |thrd| { thrd.join(); } } pub fn empty(self: *Self) bool { - self.mutex.lock(); - defer self.mutex.unlock(); + const io = self.io; + self.mutex.lockUncancelable(io); + defer self.mutex.unlock(io); return self.pull == self.push; } pub fn spawn(self: *Self, args: Args) void { const queue = self.queue; const len = queue.len; + const io = self.io; - self.mutex.lock(); + self.mutex.lockUncancelable(io); while (self.pending == len) { - self.push_cond.wait(&self.mutex); + self.push_cond.waitUncancelable(io, &self.mutex); } const push = self.push; self.queue[push] = args; self.push = if (push == self.queue_end) 0 else push + 1; self.pending += 1; - self.mutex.unlock(); + self.mutex.unlock(io); - self.pull_cond.signal(); + self.pull_cond.signal(io); } fn worker(self: *Self, buffer: []u8) void { - // Having a re-usable buffer per thread is the most efficient way - // we can do any dynamic allocations. We'll pair this later with - // a FallbackAllocator. The main issue is that some data must outlive - // the worker thread (in nonblocking mode), but this isn't something - // we need to worry about here. As far as this worker thread is - // concerned, it has a chunk of memory (buffer) which it'll pass - // to the callback function to do with as it wants. defer self.allocator.free(buffer); + const io = self.io; while (true) { - self.mutex.lock(); + self.mutex.lockUncancelable(io); while (self.pending == 0) { if (self.stopped) { - self.mutex.unlock(); + self.mutex.unlock(io); return; } - self.pull_cond.wait(&self.mutex); + self.pull_cond.waitUncancelable(io, &self.mutex); } const pull = self.pull; const args = self.queue[pull]; self.pull = if (pull == self.queue_end) 0 else pull + 1; self.pending -= 1; - self.mutex.unlock(); - self.push_cond.signal(); + self.mutex.unlock(io); + self.push_cond.signal(io); // convert Args to FullArgs, i.e. inject buffer as the last argument var full_args: FullArgs = undefined; diff --git a/src/t.zig b/src/t.zig index 0494d1e..45426f0 100644 --- a/src/t.zig +++ b/src/t.zig @@ -1,16 +1,14 @@ const std = @import("std"); const proto = @import("proto.zig"); -const c = std.c; const Io = std.Io; -const posix = std.posix; +const net = Io.net; const ArrayList = std.ArrayList; const Message = proto.Message; pub const allocator = std.testing.allocator; -// 0.16: expectEqual argument order changed (expected, actual) pub fn expectEqual(expected: anytype, actual: anytype) !void { try std.testing.expectEqual(expected, actual); } @@ -19,7 +17,6 @@ pub const expectError = std.testing.expectError; pub const expectString = std.testing.expectEqualStrings; pub const expectSlice = std.testing.expectEqualSlices; -// 0.16: posix.getrandom removed, use io.random() which fills a buffer pub fn getRandom() std.Random.DefaultPrng { const io = std.Options.debug_io; var seed_bytes: [8]u8 = undefined; @@ -114,7 +111,6 @@ pub const Writer = struct { var mask: [4]u8 = undefined; self.random.random().bytes(&mask); - // var mask = [_]u8{1, 1, 1, 1}; buf.appendSliceAssumeCapacity(&mask); for (payload, 0..) |b, i| { @@ -149,38 +145,29 @@ pub const Writer = struct { } }; -// 0.16: Io.net.Stream doesn't have read method, use libc.recv wrapper +/// reads from a net.Stream via the Io vtable (replaces raw libc recv) pub const StreamReader = struct { - socket: posix.socket_t, + io: Io, + handle: net.Socket.Handle, pub fn read(self: *StreamReader, buf: []u8) !usize { - const rc = c.recv(self.socket, buf.ptr, buf.len, 0); - if (rc == -1) { - const err = posix.errno(-1); + var bufs = [_][]u8{buf}; + return self.io.vtable.netRead(self.io.userdata, self.handle, &bufs) catch |err| { return switch (err) { - .CONNRESET => error.ConnectionResetByPeer, - .AGAIN => error.WouldBlock, + error.ConnectionResetByPeer => error.ConnectionResetByPeer, + error.Timeout => error.WouldBlock, else => error.Unexpected, }; - } - return @intCast(rc); + }; } }; -// 0.16: sockaddr_in not in std.c, define locally for darwin -const SockaddrIn = extern struct { - len: u8 = @sizeOf(SockaddrIn), - family: u8 = c.AF.INET, - port: u16, - addr: u32, - zero: [8]u8 = [_]u8{0} ** 8, -}; - pub const SocketPair = struct { writer: Writer, io: Io, - client: Io.net.Stream, - server: Io.net.Stream, + client: net.Stream, + server: net.Stream, + server_taken: bool = false, const Opts = struct { port: ?u16 = null, @@ -190,62 +177,32 @@ pub const SocketPair = struct { const io = std.Options.debug_io; const port: u16 = opts.port orelse 0; - // 0.16: use libc directly for socket operations - // Note: SOCK.CLOEXEC is not supported by darwin, set CLOEXEC via fcntl - const listener = c.socket(c.AF.INET, c.SOCK.STREAM, c.IPPROTO.TCP); - if (listener == -1) unreachable; - _ = c.fcntl(listener, c.F.SETFD, @as(c_int, c.FD_CLOEXEC)); - defer _ = c.close(listener); - - // setup sockaddr_in for 127.0.0.1 - var addr: SockaddrIn = .{ - .port = @byteSwap(port), - .addr = 0x0100007f, // 127.0.0.1 in network byte order - }; + // use Io.net to create a listener, connect a client, and accept + var listen_addr: net.IpAddress = .{ .ip4 = .loopback(port) }; + var server = net.IpAddress.listen(&listen_addr, io, .{ .reuse_address = true }) catch unreachable; - { - // setup our listener - if (c.bind(listener, @ptrCast(&addr), @sizeOf(SockaddrIn)) == -1) unreachable; - if (c.listen(listener, 1) == -1) unreachable; - // get assigned port - var addr_len: c.socklen_t = @sizeOf(SockaddrIn); - if (c.getsockname(listener, @ptrCast(&addr), &addr_len) == -1) unreachable; - } + // get the actual bound address (for ephemeral port) + const bound_addr = server.socket.address; - const client_fd = c.socket(c.AF.INET, c.SOCK.STREAM, c.IPPROTO.TCP); - if (client_fd == -1) unreachable; - { - // connect the client - const flags = c.fcntl(client_fd, c.F.GETFL); - _ = c.fcntl(client_fd, c.F.SETFL, flags | @as(c_int, c.SOCK.NONBLOCK)); - const connect_result = c.connect(client_fd, @ptrCast(&addr), @sizeOf(SockaddrIn)); - if (connect_result == -1) { - const err = std.posix.errno(connect_result); - if (err != .INPROGRESS) unreachable; - } - _ = c.fcntl(client_fd, c.F.SETFL, flags); - } + // connect client + const client = net.IpAddress.connect(&bound_addr, io, .{ .mode = .stream }) catch unreachable; - var client_addr: c.sockaddr = undefined; - var client_addr_len: c.socklen_t = @sizeOf(c.sockaddr); - const server_fd = c.accept(listener, &client_addr, &client_addr_len); - if (server_fd == -1) unreachable; + // accept server-side connection + const accepted = server.accept(io) catch unreachable; + server.deinit(io); - // 0.16: Socket struct requires address field - const loopback = Io.net.Ip4Address.loopback(@byteSwap(addr.port)); return .{ .io = io, - .client = .{ .socket = .{ .handle = client_fd, .address = .{ .ip4 = loopback } } }, - .server = .{ .socket = .{ .handle = server_fd, .address = .{ .ip4 = loopback } } }, + .client = client, + .server = accepted, .writer = Writer.init(), }; } pub fn deinit(self: *SocketPair) void { self.writer.deinit(); - // 0.16: use libc.close directly - _ = c.close(self.client.socket.handle); - _ = c.close(self.server.socket.handle); + self.client.close(self.io); + if (!self.server_taken) self.server.close(self.io); } pub fn pingPayload(self: *SocketPair, payload: []const u8) void { @@ -261,36 +218,33 @@ pub const SocketPair = struct { } pub fn sendBuf(self: *SocketPair) void { - // 0.16: use libc.send directly (debug_io doesn't support real network) - const data = self.writer.bytes(); - var remaining = data; - while (remaining.len > 0) { - const rc = c.send(self.client.socket.handle, remaining.ptr, remaining.len, 0); - if (rc <= 0) unreachable; - const n: usize = @intCast(rc); - remaining = remaining[n..]; - } + self.ioWriteAll(self.writer.bytes()) catch unreachable; self.writer.clear(); } - // 0.16: use libc.send directly pub fn clientWriteAll(self: *SocketPair, data: []const u8) !void { + try self.ioWriteAll(data); + } + + fn ioWriteAll(self: *SocketPair, data: []const u8) !void { var remaining = data; while (remaining.len > 0) { - const rc = c.send(self.client.socket.handle, remaining.ptr, remaining.len, 0); - if (rc <= 0) return error.SendError; - const n: usize = @intCast(rc); + // netWrite: header is sent first, data array's last element is the splat pattern. + // Pass remaining as header, empty pattern with splat=0. + const empty = [_][]const u8{""}; + const n = self.io.vtable.netWrite(self.io.userdata, self.client.socket.handle, remaining, &empty, 0) catch { + return error.SendError; + }; + if (n == 0) return error.SendError; remaining = remaining[n..]; } } - // 0.16: return StreamReader wrapper for server socket pub fn serverReader(self: *SocketPair) StreamReader { - return .{ .socket = self.server.socket.handle }; + return .{ .io = self.io, .handle = self.server.socket.handle }; } - // 0.16: return StreamReader wrapper for client socket pub fn clientReader(self: *SocketPair) StreamReader { - return .{ .socket = self.client.socket.handle }; + return .{ .io = self.io, .handle = self.client.socket.handle }; } }; diff --git a/src/testing.zig b/src/testing.zig index 6113429..09b619c 100644 --- a/src/testing.zig +++ b/src/testing.zig @@ -32,7 +32,6 @@ pub const Testing = struct { .sec = 0, .usec = 50_000, }); - // 0.16: access handle through socket std.posix.setsockopt(pair.client.socket.handle, std.posix.SOL.SOCKET, std.posix.SO.RCVTIMEO, &timeout) catch unreachable; const aa = arena.allocator(); @@ -54,19 +53,16 @@ pub const Testing = struct { ._closed = false, .started = 0, .stream = pair.server, - .address = std.net.Address.parseIp("127.0.0.1", port) catch unreachable, + .address = std.Io.net.IpAddress.parse("127.0.0.1", port) catch unreachable, }, .reader = reader, - .received = std.ArrayList(ws.Message).init(aa), + .received = .empty, .received_index = 0, }; } pub fn deinit(self: *Testing) void { - self.pair.writer.deinit(); - close(self.pair.client.handle); - close(self.pair.server.handle); - + self.pair.deinit(); self.arena.deinit(); t.allocator.destroy(self.arena); } @@ -107,7 +103,6 @@ pub const Testing = struct { } fn fill(self: *Testing) !void { - // 0.16: use StreamReader wrapper since Io.net.Stream doesn't have read var client_reader = self.pair.clientReader(); self.reader.fill(&client_reader) catch |err| switch (err) { error.WouldBlock => return error.NoMoreData, @@ -119,29 +114,11 @@ pub const Testing = struct { while (true) { const more, const message = (try self.reader.read()) orelse return error.NoMoreData; - try self.received.append(message); + const aa = self.arena.allocator(); + self.received.append(aa, message) catch unreachable; if (more == false) { return; } } } }; - -// std.posix.close panics on EBADF -// This is a general issue in Zig: -// https://github.com/ziglang/zig/issues/6389 -// -// For these tests, we realy don't know if the server-side of the connection -// is closed, so we try to close and ignore any errors. -fn close(fd: std.posix.fd_t) void { - const builtin = @import("builtin"); - const native_os = builtin.os.tag; - if (native_os == .windows) { - return std.os.windows.CloseHandle(fd); - } - if (native_os == .wasi and !builtin.link_libc) { - _ = std.os.wasi.fd_close(fd); - return; - } - _ = std.posix.system.close(fd); -} diff --git a/support/autobahn/client/build.zig.zon b/support/autobahn/client/build.zig.zon index 9cc6022..dcd40b8 100644 --- a/support/autobahn/client/build.zig.zon +++ b/support/autobahn/client/build.zig.zon @@ -1,11 +1,9 @@ .{ - .name = .autobahn_test_client, - .paths = .{""}, - .version = "0.0.0", - .fingerprint = 0xf155d3bf7ec3bfc, - .dependencies = .{ - .websocket = .{ - .path = "../../../" + .name = .autobahn_test_client, + .paths = .{""}, + .version = "0.0.0", + .fingerprint = 0xf155d3bf7ec3bfc, + .dependencies = .{ + .websocket = .{ .path = "../../../" }, }, - }, } diff --git a/support/autobahn/client/main.zig b/support/autobahn/client/main.zig index 4c123b8..409576f 100644 --- a/support/autobahn/client/main.zig +++ b/support/autobahn/client/main.zig @@ -28,9 +28,21 @@ pub fn main() !void { "7.9.6", "7.9.7", "7.9.8", "7.9.9", "7.13.1", "7.13.2", "9.1.1", "9.1.2", "9.1.3", "9.1.4", "9.1.5", "9.1.6", "9.2.1", "9.2.2", "9.2.3", "9.2.4", "9.2.5", "9.2.6", "9.3.1", "9.3.2", "9.3.3", "9.3.4", "9.3.5", "9.3.6", "9.3.7", "9.3.8", "9.3.9", "9.4.1", "9.4.2", "9.4.3", "9.4.4", "9.4.5", "9.4.6", "9.4.7", "9.4.8", "9.4.9", "9.5.1", "9.5.2", "9.5.3", "9.5.4", "9.5.5", "9.5.6", "9.6.1", "9.6.2", "9.6.3", "9.6.4", "9.6.5", "9.6.6", - "9.7.1", "9.7.2", "9.7.3", "9.7.4", "9.7.5", "9.7.6", "9.8.1", "9.8.2", "9.8.3", "9.8.4", "9.8.5", "9.8.6", "10.1.1", - "12.1.1", "12.1.2", "12.1.3", "12.1.4", "12.1.5", "12.1.6", "12.1.7", "12.1.8", "12.1.9", "12.1.10", "12.1.11", "12.1.12", "12.1.13", "12.1.14", "12.1.15", "12.1.16", "12.1.17", "12.1.18", "12.2.1", "12.2.2", "12.2.3", "12.2.4", "12.2.5", "12.2.6", "12.2.7", "12.2.8", "12.2.9", "12.2.10", "12.2.11", "12.2.12", "12.2.13", "12.2.14", "12.2.15", "12.2.16", "12.2.17", "12.2.18", "12.3.1", "12.3.2", "12.3.3", "12.3.4", "12.3.5", "12.3.6", "12.3.7", "12.3.8", "12.3.9", "12.3.10", "12.3.11", "12.3.12", "12.3.13", "12.3.14", "12.3.15", "12.3.16", "12.3.17", "12.3.18", "12.4.1", "12.4.2", "12.4.3", "12.4.4", "12.4.5", "12.4.6", "12.4.7", "12.4.8", "12.4.9", "12.4.10", "12.4.11", "12.4.12", "12.4.13", "12.4.14", "12.4.15", "12.4.16", "12.4.17", "12.4.18", "12.5.1", "12.5.2", "12.5.3", "12.5.4", "12.5.5", "12.5.6", "12.5.7", "12.5.8", "12.5.9", "12.5.10", "12.5.11", "12.5.12", "12.5.13", "12.5.14", "12.5.15", "12.5.16", "12.5.17", "12.5.18", - "13.1.1", "13.1.2", "13.1.3", "13.1.4", "13.1.5", "13.1.6", "13.1.7", "13.1.8", "13.1.9", "13.1.10", "13.1.11", "13.1.12", "13.1.13", "13.1.14", "13.1.15", "13.1.16", "13.1.17", "13.1.18", "13.2.1", "13.2.2", "13.2.3", "13.2.4", "13.2.5", "13.2.6", "13.2.7", "13.2.8", "13.2.9", "13.2.10", "13.2.11", "13.2.12", "13.2.13", "13.2.14", "13.2.15", "13.2.16", "13.2.17", "13.2.18", "13.3.1", "13.3.2", "13.3.3", "13.3.4", "13.3.5", "13.3.6", "13.3.7", "13.3.8", "13.3.9", "13.3.10", "13.3.11", "13.3.12", "13.3.13", "13.3.14", "13.3.15", "13.3.16", "13.3.17", "13.3.18", "13.4.1", "13.4.2", "13.4.3", "13.4.4", "13.4.5", "13.4.6", "13.4.7", "13.4.8", "13.4.9", "13.4.10", "13.4.11", "13.4.12", "13.4.13", "13.4.14", "13.4.15", "13.4.16", "13.4.17", "13.4.18", "13.5.1", "13.5.2", "13.5.3", "13.5.4", "13.5.5", "13.5.6", "13.5.7", "13.5.8", "13.5.9", "13.5.10", "13.5.11", "13.5.12", "13.5.13", "13.5.14", "13.5.15", "13.5.16", "13.5.17", "13.5.18", "13.6.1", "13.6.2", "13.6.3", "13.6.4", "13.6.5", "13.6.6", "13.6.7", "13.6.8", "13.6.9", "13.6.10", "13.6.11", "13.6.12", "13.6.13", "13.6.14", "13.6.15", "13.6.16", "13.6.17", "13.6.18", "13.7.1", "13.7.2", "13.7.3", "13.7.4", "13.7.5", "13.7.6", "13.7.7", "13.7.8", "13.7.9", "13.7.10", "13.7.11", "13.7.12", "13.7.13", "13.7.14", "13.7.15", "13.7.16", "13.7.17", "13.7.18", + "9.7.1", "9.7.2", "9.7.3", "9.7.4", "9.7.5", "9.7.6", "9.8.1", "9.8.2", "9.8.3", "9.8.4", "9.8.5", "9.8.6", "10.1.1", "12.1.1", "12.1.2", "12.1.3", + "12.1.4", "12.1.5", "12.1.6", "12.1.7", "12.1.8", "12.1.9", "12.1.10", "12.1.11", "12.1.12", "12.1.13", "12.1.14", "12.1.15", "12.1.16", "12.1.17", "12.1.18", "12.2.1", + "12.2.2", "12.2.3", "12.2.4", "12.2.5", "12.2.6", "12.2.7", "12.2.8", "12.2.9", "12.2.10", "12.2.11", "12.2.12", "12.2.13", "12.2.14", "12.2.15", "12.2.16", "12.2.17", + "12.2.18", "12.3.1", "12.3.2", "12.3.3", "12.3.4", "12.3.5", "12.3.6", "12.3.7", "12.3.8", "12.3.9", "12.3.10", "12.3.11", "12.3.12", "12.3.13", "12.3.14", "12.3.15", + "12.3.16", "12.3.17", "12.3.18", "12.4.1", "12.4.2", "12.4.3", "12.4.4", "12.4.5", "12.4.6", "12.4.7", "12.4.8", "12.4.9", "12.4.10", "12.4.11", "12.4.12", "12.4.13", + "12.4.14", "12.4.15", "12.4.16", "12.4.17", "12.4.18", "12.5.1", "12.5.2", "12.5.3", "12.5.4", "12.5.5", "12.5.6", "12.5.7", "12.5.8", "12.5.9", "12.5.10", "12.5.11", + "12.5.12", "12.5.13", "12.5.14", "12.5.15", "12.5.16", "12.5.17", "12.5.18", "13.1.1", "13.1.2", "13.1.3", "13.1.4", "13.1.5", "13.1.6", "13.1.7", "13.1.8", "13.1.9", + "13.1.10", "13.1.11", "13.1.12", "13.1.13", "13.1.14", "13.1.15", "13.1.16", "13.1.17", "13.1.18", "13.2.1", "13.2.2", "13.2.3", "13.2.4", "13.2.5", "13.2.6", "13.2.7", + "13.2.8", "13.2.9", "13.2.10", "13.2.11", "13.2.12", "13.2.13", "13.2.14", "13.2.15", "13.2.16", "13.2.17", "13.2.18", "13.3.1", "13.3.2", "13.3.3", "13.3.4", "13.3.5", + "13.3.6", "13.3.7", "13.3.8", "13.3.9", "13.3.10", "13.3.11", "13.3.12", "13.3.13", "13.3.14", "13.3.15", "13.3.16", "13.3.17", "13.3.18", "13.4.1", "13.4.2", "13.4.3", + "13.4.4", "13.4.5", "13.4.6", "13.4.7", "13.4.8", "13.4.9", "13.4.10", "13.4.11", "13.4.12", "13.4.13", "13.4.14", "13.4.15", "13.4.16", "13.4.17", "13.4.18", "13.5.1", + "13.5.2", "13.5.3", "13.5.4", "13.5.5", "13.5.6", "13.5.7", "13.5.8", "13.5.9", "13.5.10", "13.5.11", "13.5.12", "13.5.13", "13.5.14", "13.5.15", "13.5.16", "13.5.17", + "13.5.18", "13.6.1", "13.6.2", "13.6.3", "13.6.4", "13.6.5", "13.6.6", "13.6.7", "13.6.8", "13.6.9", "13.6.10", "13.6.11", "13.6.12", "13.6.13", "13.6.14", "13.6.15", + "13.6.16", "13.6.17", "13.6.18", "13.7.1", "13.7.2", "13.7.3", "13.7.4", "13.7.5", "13.7.6", "13.7.7", "13.7.8", "13.7.9", "13.7.10", "13.7.11", "13.7.12", "13.7.13", + "13.7.14", "13.7.15", "13.7.16", "13.7.17", "13.7.18", }; // wait 5 seconds for autobanh server to be up diff --git a/support/autobahn/server/build.zig.zon b/support/autobahn/server/build.zig.zon index 8a462b2..9c9c4d7 100644 --- a/support/autobahn/server/build.zig.zon +++ b/support/autobahn/server/build.zig.zon @@ -1,9 +1,9 @@ .{ - .name = .autobahn_test_server, - .paths = .{""}, - .version = "0.0.0", - .fingerprint = 0x923c8c98ca4d09ab, - .dependencies = .{ - .websocket = .{.path = "../../../"}, - }, + .name = .autobahn_test_server, + .paths = .{""}, + .version = "0.0.0", + .fingerprint = 0x923c8c98ca4d09ab, + .dependencies = .{ + .websocket = .{ .path = "../../../" }, + }, } diff --git a/support/autobahn/server/main.zig b/support/autobahn/server/main.zig index 3a80db9..e34d868 100644 --- a/support/autobahn/server/main.zig +++ b/support/autobahn/server/main.zig @@ -102,7 +102,7 @@ const Handler = struct { if (std.unicode.utf8ValidateSlice(data)) { try self.conn.writeText(data); } else { - self.conn.close(.{.code = 1007}) catch {}; + self.conn.close(.{ .code = 1007 }) catch {}; } }, } diff --git a/test_runner.zig b/test_runner.zig index 37e1f27..3208b39 100644 --- a/test_runner.zig +++ b/test_runner.zig @@ -163,18 +163,24 @@ const Status = enum { }; const SlowTracker = struct { + const Io = std.Io; const SlowestQueue = std.PriorityDequeue(TestInfo, void, compareTiming); max: usize, slowest: SlowestQueue, - timer: std.time.Timer, + start_ts: Io.Timestamp, + io: Io, + + allocator: Allocator, fn init(allocator: Allocator, count: u32) SlowTracker { - const timer = std.time.Timer.start() catch @panic("failed to start timer"); - var slowest = SlowestQueue.init(allocator, {}); - slowest.ensureTotalCapacity(count) catch @panic("OOM"); + const io = std.Options.debug_io; + var slowest = SlowestQueue.initContext({}); + slowest.ensureUnusedCapacity(allocator, count) catch @panic("OOM"); return .{ .max = count, - .timer = timer, + .io = io, + .allocator = allocator, + .start_ts = Io.Timestamp.now(io, .awake), .slowest = slowest, }; } @@ -185,47 +191,41 @@ const SlowTracker = struct { }; fn deinit(self: SlowTracker) void { - self.slowest.deinit(); + self.slowest.deinit(self.allocator); } fn startTiming(self: *SlowTracker) void { - self.timer.reset(); + self.start_ts = Io.Timestamp.now(self.io, .awake); } fn endTiming(self: *SlowTracker, test_name: []const u8) u64 { - var timer = self.timer; - const ns = timer.lap(); + const end_ts = Io.Timestamp.now(self.io, .awake); + const ns: u64 = @intCast(end_ts.nanoseconds - self.start_ts.nanoseconds); var slowest = &self.slowest; - if (slowest.count() < self.max) { - // Capacity is fixed to the # of slow tests we want to track - // If we've tracked fewer tests than this capacity, than always add - slowest.add(TestInfo{ .ns = ns, .name = test_name }) catch @panic("failed to track test timing"); + if (slowest.len < self.max) { + slowest.push(self.allocator, TestInfo{ .ns = ns, .name = test_name }) catch @panic("failed to track test timing"); return ns; } { - // Optimization to avoid shifting the dequeue for the common case - // where the test isn't one of our slowest. const fastest_of_the_slow = slowest.peekMin() orelse unreachable; if (fastest_of_the_slow.ns > ns) { - // the test was faster than our fastest slow test, don't add return ns; } } - // the previous fastest of our slow tests, has been pushed off. - _ = slowest.removeMin(); - slowest.add(TestInfo{ .ns = ns, .name = test_name }) catch @panic("failed to track test timing"); + _ = slowest.popMin(); + slowest.push(self.allocator, TestInfo{ .ns = ns, .name = test_name }) catch @panic("failed to track test timing"); return ns; } fn display(self: *SlowTracker) !void { var slowest = self.slowest; - const count = slowest.count(); + const count = slowest.len; Printer.fmt("Slowest {d} test{s}: \n", .{ count, if (count != 1) "s" else "" }); - while (slowest.removeMinOrNull()) |info| { + while (slowest.popMin()) |info| { const ms = @as(f64, @floatFromInt(info.ns)) / 1_000_000.0; Printer.fmt(" {d:.2}ms\t{s}\n", .{ ms, info.name }); } -- 2.51.2