diff --git a/src/Kqueue.zig b/src/Kqueue.zig index 1fd3cea..46b4736 100644 --- a/src/Kqueue.zig +++ b/src/Kqueue.zig @@ -344,6 +344,12 @@ fn prepTask(self: *Kqueue, task: *io.Task) !void { } }, + .readv => |req| { + self.in_flight.push(task); + const kevent = evSet(@intCast(req.fd), EVFILT.READ, EV.ADD | EV.ONESHOT, task); + try self.submission_queue.append(self.gpa, kevent); + }, + .recv => |req| { self.in_flight.push(task); const kevent = evSet(@intCast(req.fd), EVFILT.READ, EV.ADD | EV.ONESHOT, task); @@ -487,6 +493,18 @@ fn cancelTask(self: *Kqueue, task: *io.Task) !void { } }, + .readv => |cancel_req| { + self.in_flight.remove(task); + task.result = .{ .readv = error.Canceled }; + const kevent = evSet( + @intCast(cancel_req.fd), + EVFILT.READ, + EV.DELETE, + task, + ); + try self.submission_queue.append(self.gpa, kevent); + }, + .recv => |cancel_req| { self.in_flight.remove(task); task.result = .{ .recv = error.Canceled }; @@ -654,6 +672,7 @@ fn handleSynchronousCompletion( // async tasks. These can be handled synchronously in a cancel all .accept, .poll, + .readv, .recv, .write, .writev, @@ -700,6 +719,7 @@ fn handleSynchronousCompletion( .msg_ring => .{ .msg_ring = error.Canceled }, .noop => unreachable, .poll => .{ .poll = error.Canceled }, + .readv => .{ .readv = error.Canceled }, .recv => .{ .recv = error.Canceled }, .socket => .{ .socket = error.Canceled }, .statx => .{ .statx = error.Canceled }, @@ -786,6 +806,22 @@ fn handleCompletion( return task.callback(rt, task.*); }, + .readv => |req| { + defer self.releaseTask(rt, task); + self.in_flight.remove(task); + if (event.flags & EV.ERROR != 0) { + // Interpret data as an errno + const err = unexpectedError(dataToE(event.data)); + task.result = .{ .readv = err }; + return task.callback(rt, task.*); + } + if (posix.readv(req.fd, req.vecs)) |n| + task.result = .{ .readv = n } + else |_| + task.result = .{ .readv = error.Unexpected }; + return task.callback(rt, task.*); + }, + .recv => |req| { defer self.releaseTask(rt, task); self.in_flight.remove(task); diff --git a/src/Mock.zig b/src/Mock.zig index 60a0248..909def3 100644 --- a/src/Mock.zig +++ b/src/Mock.zig @@ -19,6 +19,7 @@ deadline_cb: ?*const fn (*io.Task) io.Result = null, msg_ring_cb: ?*const fn (*io.Task) io.Result = null, noop_cb: ?*const fn (*io.Task) io.Result = null, poll_cb: ?*const fn (*io.Task) io.Result = null, +readv_cb: ?*const fn (*io.Task) io.Result = null, recv_cb: ?*const fn (*io.Task) io.Result = null, socket_cb: ?*const fn (*io.Task) io.Result = null, statx_cb: ?*const fn (*io.Task) io.Result = null, @@ -69,6 +70,7 @@ pub fn submit(self: *Mock, queue: *Queue(io.Task, .in_flight)) !void { .msg_ring => if (self.msg_ring_cb) |cb| cb(task) else return error.NoMockCallback, .noop => if (self.noop_cb) |cb| cb(task) else return error.NoMockCallback, .poll => if (self.poll_cb) |cb| cb(task) else return error.NoMockCallback, + .readv => if (self.readv_cb) |cb| cb(task) else return error.NoMockCallback, .recv => if (self.recv_cb) |cb| cb(task) else return error.NoMockCallback, .socket => if (self.socket_cb) |cb| cb(task) else return error.NoMockCallback, .statx => if (self.statx_cb) |cb| cb(task) else return error.NoMockCallback, diff --git a/src/Uring.zig b/src/Uring.zig index 33c678a..f86ffc1 100644 --- a/src/Uring.zig +++ b/src/Uring.zig @@ -225,6 +225,13 @@ fn prepTask(self: *Uring, task: *io.Task) void { self.prepDeadline(task, sqe); }, + .readv => |req| { + const sqe = self.getSqe(); + sqe.prep_readv(req.fd, req.vecs, 0); + sqe.user_data = @intFromPtr(task); + self.prepDeadline(task, sqe); + }, + // user* is only sent internally between rings and higher level wrappers .userfd, .usermsg, .userptr => unreachable, } @@ -354,6 +361,13 @@ pub fn reapCompletions(self: *Uring, rt: *io.Ring) anyerror!void { else => |e| unexpectedError(e), } }, + .readv => .{ .readv = switch (cqeToE(cqe.res)) { + .SUCCESS => @intCast(cqe.res), + .INVAL => io.ResultError.Invalid, + .CANCELED => io.ResultError.Canceled, + else => |e| unexpectedError(e), + } }, + .usermsg => .{ .usermsg = @intCast(cqe.res) }, // userfd should never reach the runtime diff --git a/src/main.zig b/src/main.zig index 88b235a..68ef2b9 100644 --- a/src/main.zig +++ b/src/main.zig @@ -493,6 +493,7 @@ pub const Op = enum { socket, connect, statx, + readv, /// userfd is meant to send file descriptors between Ring instances (using msgRing) userfd, @@ -546,6 +547,10 @@ pub const Request = union(Op) { path: [:0]const u8, result: *Statx, // this will be filled in by the op }, + readv: struct { + fd: posix.fd_t, + vecs: []const posix.iovec, + }, userfd, usermsg, @@ -567,6 +572,7 @@ pub const Result = union(Op) { socket: ResultError!posix.fd_t, connect: ResultError!void, statx: ResultError!*Statx, + readv: ResultError!usize, userfd: anyerror!posix.fd_t, usermsg: u16,