| ... | @@ -404,6 +404,25 @@ pub const IO_Uring = struct { | ... | @@ -404,6 +404,25 @@ pub const IO_Uring = struct { |
| 404 | return sqe; | 404 | return sqe; |
| 405 | } | 405 | } |
| 406 | | 406 | |
| | 407 | /// Queues (but does not submit) an SQE to perform a IORING_OP_READ_FIXED. |
| | 408 | /// The `buffer` provided must be registered with the kernel by calling `register_buffers` first. |
| | 409 | /// The `buffer_index` must be the same as its index in the array provided to `register_buffers`. |
| | 410 | /// |
| | 411 | /// Returns a pointer to the SQE so that you can further modify the SQE for advanced use cases. |
| | 412 | pub fn read_fixed( |
| | 413 | self: *IO_Uring, |
| | 414 | user_data: u64, |
| | 415 | fd: os.fd_t, |
| | 416 | buffer: *os.iovec, |
| | 417 | offset: u64, |
| | 418 | buffer_index: u16, |
| | 419 | ) !*io_uring_sqe { |
| | 420 | const sqe = try self.get_sqe(); |
| | 421 | io_uring_prep_read_fixed(sqe, fd, buffer, offset, buffer_index); |
| | 422 | sqe.user_data = user_data; |
| | 423 | return sqe; |
| | 424 | } |
| | 425 | |
| 407 | /// Queues (but does not submit) an SQE to perform a `pwritev()`. | 426 | /// Queues (but does not submit) an SQE to perform a `pwritev()`. |
| 408 | /// Returns a pointer to the SQE so that you can further modify the SQE for advanced use cases. | 427 | /// Returns a pointer to the SQE so that you can further modify the SQE for advanced use cases. |
| 409 | /// For example, if you want to do a `pwritev2()` then set `rw_flags` on the returned SQE. | 428 | /// For example, if you want to do a `pwritev2()` then set `rw_flags` on the returned SQE. |
| ... | @@ -421,6 +440,25 @@ pub const IO_Uring = struct { | ... | @@ -421,6 +440,25 @@ pub const IO_Uring = struct { |
| 421 | return sqe; | 440 | return sqe; |
| 422 | } | 441 | } |
| 423 | | 442 | |
| | 443 | /// Queues (but does not submit) an SQE to perform a IORING_OP_WRITE_FIXED. |
| | 444 | /// The `buffer` provided must be registered with the kernel by calling `register_buffers` first. |
| | 445 | /// The `buffer_index` must be the same as its index in the array provided to `register_buffers`. |
| | 446 | /// |
| | 447 | /// Returns a pointer to the SQE so that you can further modify the SQE for advanced use cases. |
| | 448 | pub fn write_fixed( |
| | 449 | self: *IO_Uring, |
| | 450 | user_data: u64, |
| | 451 | fd: os.fd_t, |
| | 452 | buffer: *os.iovec, |
| | 453 | offset: u64, |
| | 454 | buffer_index: u16, |
| | 455 | ) !*io_uring_sqe { |
| | 456 | const sqe = try self.get_sqe(); |
| | 457 | io_uring_prep_write_fixed(sqe, fd, buffer, offset, buffer_index); |
| | 458 | sqe.user_data = user_data; |
| | 459 | return sqe; |
| | 460 | } |
| | 461 | |
| 424 | /// Queues (but does not submit) an SQE to perform an `accept4(2)` on a socket. | 462 | /// Queues (but does not submit) an SQE to perform an `accept4(2)` on a socket. |
| 425 | /// Returns a pointer to the SQE. | 463 | /// Returns a pointer to the SQE. |
| 426 | pub fn accept( | 464 | pub fn accept( |
| ... | @@ -674,6 +712,29 @@ pub const IO_Uring = struct { | ... | @@ -674,6 +712,29 @@ pub const IO_Uring = struct { |
| 674 | try handle_registration_result(res); | 712 | try handle_registration_result(res); |
| 675 | } | 713 | } |
| 676 | | 714 | |
| | 715 | /// Registers an array of buffers for use with `read_fixed` and `write_fixed`. |
| | 716 | pub fn register_buffers(self: *IO_Uring, buffers: []const os.iovec) !void { |
| | 717 | assert(self.fd >= 0); |
| | 718 | const res = linux.io_uring_register( |
| | 719 | self.fd, |
| | 720 | .REGISTER_BUFFERS, |
| | 721 | buffers.ptr, |
| | 722 | @intCast(u32, buffers.len), |
| | 723 | ); |
| | 724 | try handle_registration_result(res); |
| | 725 | } |
| | 726 | |
| | 727 | /// Unregister the registered buffers. |
| | 728 | pub fn unregister_buffers(self: *IO_Uring) !void { |
| | 729 | assert(self.fd >= 0); |
| | 730 | const res = linux.io_uring_register(self.fd, .UNREGISTER_BUFFERS, null, 0); |
| | 731 | switch (linux.getErrno(res)) { |
| | 732 | .SUCCESS => {}, |
| | 733 | .NXIO => return error.BuffersNotRegistered, |
| | 734 | else => |errno| return os.unexpectedErrno(errno), |
| | 735 | } |
| | 736 | } |
| | 737 | |
| 677 | fn handle_registration_result(res: usize) !void { | 738 | fn handle_registration_result(res: usize) !void { |
| 678 | switch (linux.getErrno(res)) { | 739 | switch (linux.getErrno(res)) { |
| 679 | .SUCCESS => {}, | 740 | .SUCCESS => {}, |
| ... | @@ -905,6 +966,16 @@ pub fn io_uring_prep_writev( | ... | @@ -905,6 +966,16 @@ pub fn io_uring_prep_writev( |
| 905 | io_uring_prep_rw(.WRITEV, sqe, fd, @ptrToInt(iovecs.ptr), iovecs.len, offset); | 966 | io_uring_prep_rw(.WRITEV, sqe, fd, @ptrToInt(iovecs.ptr), iovecs.len, offset); |
| 906 | } | 967 | } |
| 907 | | 968 | |
| | 969 | pub fn io_uring_prep_read_fixed(sqe: *io_uring_sqe, fd: os.fd_t, buffer: *os.iovec, offset: u64, buffer_index: u16) void { |
| | 970 | io_uring_prep_rw(.READ_FIXED, sqe, fd, @ptrToInt(buffer.iov_base), buffer.iov_len, offset); |
| | 971 | sqe.buf_index = buffer_index; |
| | 972 | } |
| | 973 | |
| | 974 | pub fn io_uring_prep_write_fixed(sqe: *io_uring_sqe, fd: os.fd_t, buffer: *os.iovec, offset: u64, buffer_index: u16) void { |
| | 975 | io_uring_prep_rw(.WRITE_FIXED, sqe, fd, @ptrToInt(buffer.iov_base), buffer.iov_len, offset); |
| | 976 | sqe.buf_index = buffer_index; |
| | 977 | } |
| | 978 | |
| 908 | pub fn io_uring_prep_accept( | 979 | pub fn io_uring_prep_accept( |
| 909 | sqe: *io_uring_sqe, | 980 | sqe: *io_uring_sqe, |
| 910 | fd: os.fd_t, | 981 | fd: os.fd_t, |
| ... | @@ -1282,6 +1353,63 @@ test "write/read" { | ... | @@ -1282,6 +1353,63 @@ test "write/read" { |
| 1282 | try testing.expectEqualSlices(u8, buffer_write[0..], buffer_read[0..]); | 1353 | try testing.expectEqualSlices(u8, buffer_write[0..], buffer_read[0..]); |
| 1283 | } | 1354 | } |
| 1284 | | 1355 | |
| | 1356 | test "write_fixed/read_fixed" { |
| | 1357 | if (builtin.os.tag != .linux) return error.SkipZigTest; |
| | 1358 | |
| | 1359 | var ring = IO_Uring.init(2, 0) catch |err| switch (err) { |
| | 1360 | error.SystemOutdated => return error.SkipZigTest, |
| | 1361 | error.PermissionDenied => return error.SkipZigTest, |
| | 1362 | else => return err, |
| | 1363 | }; |
| | 1364 | defer ring.deinit(); |
| | 1365 | |
| | 1366 | const path = "test_io_uring_write_read_fixed"; |
| | 1367 | const file = try std.fs.cwd().createFile(path, .{ .read = true, .truncate = true }); |
| | 1368 | defer file.close(); |
| | 1369 | defer std.fs.cwd().deleteFile(path) catch {}; |
| | 1370 | const fd = file.handle; |
| | 1371 | |
| | 1372 | var raw_buffers: [2][11]u8 = undefined; |
| | 1373 | // First buffer will be written to the file. |
| | 1374 | std.mem.set(u8, &raw_buffers[0], 'z'); |
| | 1375 | std.mem.copy(u8, &raw_buffers[0], "foobar"); |
| | 1376 | |
| | 1377 | var buffers = [2]os.iovec{ |
| | 1378 | .{ .iov_base = &raw_buffers[0], .iov_len = raw_buffers[0].len }, |
| | 1379 | .{ .iov_base = &raw_buffers[1], .iov_len = raw_buffers[1].len }, |
| | 1380 | }; |
| | 1381 | try ring.register_buffers(&buffers); |
| | 1382 | |
| | 1383 | const sqe_write = try ring.write_fixed(0x45454545, fd, &buffers[0], 3, 0); |
| | 1384 | try testing.expectEqual(linux.IORING_OP.WRITE_FIXED, sqe_write.opcode); |
| | 1385 | try testing.expectEqual(@as(u64, 3), sqe_write.off); |
| | 1386 | sqe_write.flags |= linux.IOSQE_IO_LINK; |
| | 1387 | |
| | 1388 | const sqe_read = try ring.read_fixed(0x12121212, fd, &buffers[1], 0, 1); |
| | 1389 | try testing.expectEqual(linux.IORING_OP.READ_FIXED, sqe_read.opcode); |
| | 1390 | try testing.expectEqual(@as(u64, 0), sqe_read.off); |
| | 1391 | |
| | 1392 | try testing.expectEqual(@as(u32, 2), try ring.submit()); |
| | 1393 | |
| | 1394 | const cqe_write = try ring.copy_cqe(); |
| | 1395 | const cqe_read = try ring.copy_cqe(); |
| | 1396 | |
| | 1397 | try testing.expectEqual(linux.io_uring_cqe{ |
| | 1398 | .user_data = 0x45454545, |
| | 1399 | .res = @intCast(i32, buffers[0].iov_len), |
| | 1400 | .flags = 0, |
| | 1401 | }, cqe_write); |
| | 1402 | try testing.expectEqual(linux.io_uring_cqe{ |
| | 1403 | .user_data = 0x12121212, |
| | 1404 | .res = @intCast(i32, buffers[1].iov_len), |
| | 1405 | .flags = 0, |
| | 1406 | }, cqe_read); |
| | 1407 | |
| | 1408 | try testing.expectEqualSlices(u8, "\x00\x00\x00", buffers[1].iov_base[0..3]); |
| | 1409 | try testing.expectEqualSlices(u8, "foobar", buffers[1].iov_base[3..9]); |
| | 1410 | try testing.expectEqualSlices(u8, "zz", buffers[1].iov_base[9..11]); |
| | 1411 | } |
| | 1412 | |
| 1285 | test "openat" { | 1413 | test "openat" { |
| 1286 | if (builtin.os.tag != .linux) return error.SkipZigTest; | 1414 | if (builtin.os.tag != .linux) return error.SkipZigTest; |
| 1287 | | 1415 | |