| ... | @@ -505,10 +505,12 @@ pub const IO_Uring = struct { | ... | @@ -505,10 +505,12 @@ pub const IO_Uring = struct { |
| 505 | return sqe; | 505 | return sqe; |
| 506 | } | 506 | } |
| 507 | | 507 | |
| 508 | /// Queues (but does not submit) an SQE to perform an multishot `accept4(2)` on a socket. | 508 | /// Queues an multishot accept on a socket. |
| | 509 | /// |
| 509 | /// Multishot variant allows an application to issue a single accept request, | 510 | /// Multishot variant allows an application to issue a single accept request, |
| 510 | /// which will repeatedly trigger a CQE when a connection request comes in. | 511 | /// which will repeatedly trigger a CQE when a connection request comes in. |
| 511 | /// Returns a pointer to the SQE. | 512 | /// While IORING_CQE_F_MORE flag is set in CQE flags accept will generate |
| | 513 | /// further CQEs. |
| 512 | pub fn accept_multishot( | 514 | pub fn accept_multishot( |
| 513 | self: *IO_Uring, | 515 | self: *IO_Uring, |
| 514 | user_data: u64, | 516 | user_data: u64, |
| ... | @@ -523,6 +525,44 @@ pub const IO_Uring = struct { | ... | @@ -523,6 +525,44 @@ pub const IO_Uring = struct { |
| 523 | return sqe; | 525 | return sqe; |
| 524 | } | 526 | } |
| 525 | | 527 | |
| | 528 | /// Queues an accept using direct (registered) file descriptors. |
| | 529 | /// |
| | 530 | /// To use an accept direct variant, the application must first have registered |
| | 531 | /// a file table (with register_files). An unused table index will be |
| | 532 | /// dynamically chosen and returned in the CQE res field. |
| | 533 | /// |
| | 534 | /// After creation, they can be used by setting IOSQE_FIXED_FILE in the SQE |
| | 535 | /// flags member, and setting the SQE fd field to the direct descriptor value |
| | 536 | /// rather than the regular file descriptor. |
| | 537 | pub fn accept_direct( |
| | 538 | self: *IO_Uring, |
| | 539 | user_data: u64, |
| | 540 | fd: os.fd_t, |
| | 541 | addr: ?*os.sockaddr, |
| | 542 | addrlen: ?*os.socklen_t, |
| | 543 | flags: u32, |
| | 544 | ) !*linux.io_uring_sqe { |
| | 545 | const sqe = try self.get_sqe(); |
| | 546 | io_uring_prep_accept_direct(sqe, fd, addr, addrlen, flags, linux.IORING_FILE_INDEX_ALLOC); |
| | 547 | sqe.user_data = user_data; |
| | 548 | return sqe; |
| | 549 | } |
| | 550 | |
| | 551 | /// Queues an multishot accept using direct (registered) file descriptors. |
| | 552 | pub fn accept_multishot_direct( |
| | 553 | self: *IO_Uring, |
| | 554 | user_data: u64, |
| | 555 | fd: os.fd_t, |
| | 556 | addr: ?*os.sockaddr, |
| | 557 | addrlen: ?*os.socklen_t, |
| | 558 | flags: u32, |
| | 559 | ) !*linux.io_uring_sqe { |
| | 560 | const sqe = try self.get_sqe(); |
| | 561 | io_uring_prep_multishot_accept_direct(sqe, fd, addr, addrlen, flags); |
| | 562 | sqe.user_data = user_data; |
| | 563 | return sqe; |
| | 564 | } |
| | 565 | |
| 526 | /// Queue (but does not submit) an SQE to perform a `connect(2)` on a socket. | 566 | /// Queue (but does not submit) an SQE to perform a `connect(2)` on a socket. |
| 527 | /// Returns a pointer to the SQE. | 567 | /// Returns a pointer to the SQE. |
| 528 | pub fn connect( | 568 | pub fn connect( |
| ... | @@ -700,6 +740,30 @@ pub const IO_Uring = struct { | ... | @@ -700,6 +740,30 @@ pub const IO_Uring = struct { |
| 700 | return sqe; | 740 | return sqe; |
| 701 | } | 741 | } |
| 702 | | 742 | |
| | 743 | /// Queues an openat using direct (registered) file descriptors. |
| | 744 | /// |
| | 745 | /// To use an accept direct variant, the application must first have registered |
| | 746 | /// a file table (with register_files). An unused table index will be |
| | 747 | /// dynamically chosen and returned in the CQE res field. |
| | 748 | /// |
| | 749 | /// After creation, they can be used by setting IOSQE_FIXED_FILE in the SQE |
| | 750 | /// flags member, and setting the SQE fd field to the direct descriptor value |
| | 751 | /// rather than the regular file descriptor. |
| | 752 | pub fn openat_direct( |
| | 753 | self: *IO_Uring, |
| | 754 | user_data: u64, |
| | 755 | fd: os.fd_t, |
| | 756 | path: [*:0]const u8, |
| | 757 | flags: u32, |
| | 758 | mode: os.mode_t, |
| | 759 | file_index: u32, |
| | 760 | ) !*linux.io_uring_sqe { |
| | 761 | const sqe = try self.get_sqe(); |
| | 762 | io_uring_prep_openat_direct(sqe, fd, path, flags, mode, file_index); |
| | 763 | sqe.user_data = user_data; |
| | 764 | return sqe; |
| | 765 | } |
| | 766 | |
| 703 | /// Queues (but does not submit) an SQE to perform a `close(2)`. | 767 | /// Queues (but does not submit) an SQE to perform a `close(2)`. |
| 704 | /// Returns a pointer to the SQE. | 768 | /// Returns a pointer to the SQE. |
| 705 | pub fn close(self: *IO_Uring, user_data: u64, fd: os.fd_t) !*linux.io_uring_sqe { | 769 | pub fn close(self: *IO_Uring, user_data: u64, fd: os.fd_t) !*linux.io_uring_sqe { |
| ... | @@ -709,6 +773,14 @@ pub const IO_Uring = struct { | ... | @@ -709,6 +773,14 @@ pub const IO_Uring = struct { |
| 709 | return sqe; | 773 | return sqe; |
| 710 | } | 774 | } |
| 711 | | 775 | |
| | 776 | /// Queues close of registered file descriptor. |
| | 777 | pub fn close_direct(self: *IO_Uring, user_data: u64, file_index: u32) !*linux.io_uring_sqe { |
| | 778 | const sqe = try self.get_sqe(); |
| | 779 | io_uring_prep_close_direct(sqe, file_index); |
| | 780 | sqe.user_data = user_data; |
| | 781 | return sqe; |
| | 782 | } |
| | 783 | |
| 712 | /// Queues (but does not submit) an SQE to register a timeout operation. | 784 | /// Queues (but does not submit) an SQE to register a timeout operation. |
| 713 | /// Returns a pointer to the SQE. | 785 | /// Returns a pointer to the SQE. |
| 714 | /// | 786 | /// |
| ... | @@ -1157,6 +1229,54 @@ pub const IO_Uring = struct { | ... | @@ -1157,6 +1229,54 @@ pub const IO_Uring = struct { |
| 1157 | else => |errno| return os.unexpectedErrno(errno), | 1229 | else => |errno| return os.unexpectedErrno(errno), |
| 1158 | } | 1230 | } |
| 1159 | } | 1231 | } |
| | 1232 | |
| | 1233 | /// Prepares a socket creation request. |
| | 1234 | /// New socket fd will be returned in completion result. |
| | 1235 | pub fn socket( |
| | 1236 | self: *IO_Uring, |
| | 1237 | user_data: u64, |
| | 1238 | domain: u32, |
| | 1239 | socket_type: u32, |
| | 1240 | protocol: u32, |
| | 1241 | flags: u32, |
| | 1242 | ) !*linux.io_uring_sqe { |
| | 1243 | const sqe = try self.get_sqe(); |
| | 1244 | io_uring_prep_socket(sqe, domain, socket_type, protocol, flags); |
| | 1245 | sqe.user_data = user_data; |
| | 1246 | return sqe; |
| | 1247 | } |
| | 1248 | |
| | 1249 | /// Prepares a socket creation request for registered file at index `file_index`. |
| | 1250 | pub fn socket_direct( |
| | 1251 | self: *IO_Uring, |
| | 1252 | user_data: u64, |
| | 1253 | domain: u32, |
| | 1254 | socket_type: u32, |
| | 1255 | protocol: u32, |
| | 1256 | flags: u32, |
| | 1257 | file_index: u32, |
| | 1258 | ) !*linux.io_uring_sqe { |
| | 1259 | const sqe = try self.get_sqe(); |
| | 1260 | io_uring_prep_socket_direct(sqe, domain, socket_type, protocol, flags, file_index); |
| | 1261 | sqe.user_data = user_data; |
| | 1262 | return sqe; |
| | 1263 | } |
| | 1264 | |
| | 1265 | /// Prepares a socket creation request for registered file, index chosen by kernel (file index alloc). |
| | 1266 | /// File index will be returned in CQE res field. |
| | 1267 | pub fn socket_direct_alloc( |
| | 1268 | self: *IO_Uring, |
| | 1269 | user_data: u64, |
| | 1270 | domain: u32, |
| | 1271 | socket_type: u32, |
| | 1272 | protocol: u32, |
| | 1273 | flags: u32, |
| | 1274 | ) !*linux.io_uring_sqe { |
| | 1275 | const sqe = try self.get_sqe(); |
| | 1276 | io_uring_prep_socket_direct_alloc(sqe, domain, socket_type, protocol, flags); |
| | 1277 | sqe.user_data = user_data; |
| | 1278 | return sqe; |
| | 1279 | } |
| 1160 | }; | 1280 | }; |
| 1161 | | 1281 | |
| 1162 | pub const SubmissionQueue = struct { | 1282 | pub const SubmissionQueue = struct { |
| ... | @@ -1391,6 +1511,41 @@ pub fn io_uring_prep_accept( | ... | @@ -1391,6 +1511,41 @@ pub fn io_uring_prep_accept( |
| 1391 | sqe.rw_flags = flags; | 1511 | sqe.rw_flags = flags; |
| 1392 | } | 1512 | } |
| 1393 | | 1513 | |
| | 1514 | pub fn io_uring_prep_accept_direct( |
| | 1515 | sqe: *linux.io_uring_sqe, |
| | 1516 | fd: os.fd_t, |
| | 1517 | addr: ?*os.sockaddr, |
| | 1518 | addrlen: ?*os.socklen_t, |
| | 1519 | flags: u32, |
| | 1520 | file_index: u32, |
| | 1521 | ) void { |
| | 1522 | io_uring_prep_accept(sqe, fd, addr, addrlen, flags); |
| | 1523 | __io_uring_set_target_fixed_file(sqe, file_index); |
| | 1524 | } |
| | 1525 | |
| | 1526 | pub fn io_uring_prep_multishot_accept_direct( |
| | 1527 | sqe: *linux.io_uring_sqe, |
| | 1528 | fd: os.fd_t, |
| | 1529 | addr: ?*os.sockaddr, |
| | 1530 | addrlen: ?*os.socklen_t, |
| | 1531 | flags: u32, |
| | 1532 | ) void { |
| | 1533 | io_uring_prep_multishot_accept(sqe, fd, addr, addrlen, flags); |
| | 1534 | __io_uring_set_target_fixed_file(sqe, linux.IORING_FILE_INDEX_ALLOC); |
| | 1535 | } |
| | 1536 | |
| | 1537 | fn __io_uring_set_target_fixed_file(sqe: *linux.io_uring_sqe, file_index: u32) void { |
| | 1538 | const sqe_file_index: u32 = if (file_index == linux.IORING_FILE_INDEX_ALLOC) |
| | 1539 | linux.IORING_FILE_INDEX_ALLOC |
| | 1540 | else |
| | 1541 | // 0 means no fixed files, indexes should be encoded as "index + 1" |
| | 1542 | file_index + 1; |
| | 1543 | // This filed is overloaded in liburing: |
| | 1544 | // splice_fd_in: i32 |
| | 1545 | // sqe_file_index: u32 |
| | 1546 | sqe.splice_fd_in = @bitCast(sqe_file_index); |
| | 1547 | } |
| | 1548 | |
| 1394 | pub fn io_uring_prep_connect( | 1549 | pub fn io_uring_prep_connect( |
| 1395 | sqe: *linux.io_uring_sqe, | 1550 | sqe: *linux.io_uring_sqe, |
| 1396 | fd: os.fd_t, | 1551 | fd: os.fd_t, |
| ... | @@ -1474,6 +1629,18 @@ pub fn io_uring_prep_openat( | ... | @@ -1474,6 +1629,18 @@ pub fn io_uring_prep_openat( |
| 1474 | sqe.rw_flags = flags; | 1629 | sqe.rw_flags = flags; |
| 1475 | } | 1630 | } |
| 1476 | | 1631 | |
| | 1632 | pub fn io_uring_prep_openat_direct( |
| | 1633 | sqe: *linux.io_uring_sqe, |
| | 1634 | fd: os.fd_t, |
| | 1635 | path: [*:0]const u8, |
| | 1636 | flags: u32, |
| | 1637 | mode: os.mode_t, |
| | 1638 | file_index: u32, |
| | 1639 | ) void { |
| | 1640 | io_uring_prep_openat(sqe, fd, path, flags, mode); |
| | 1641 | __io_uring_set_target_fixed_file(sqe, file_index); |
| | 1642 | } |
| | 1643 | |
| 1477 | pub fn io_uring_prep_close(sqe: *linux.io_uring_sqe, fd: os.fd_t) void { | 1644 | pub fn io_uring_prep_close(sqe: *linux.io_uring_sqe, fd: os.fd_t) void { |
| 1478 | sqe.* = .{ | 1645 | sqe.* = .{ |
| 1479 | .opcode = .CLOSE, | 1646 | .opcode = .CLOSE, |
| ... | @@ -1493,6 +1660,11 @@ pub fn io_uring_prep_close(sqe: *linux.io_uring_sqe, fd: os.fd_t) void { | ... | @@ -1493,6 +1660,11 @@ pub fn io_uring_prep_close(sqe: *linux.io_uring_sqe, fd: os.fd_t) void { |
| 1493 | }; | 1660 | }; |
| 1494 | } | 1661 | } |
| 1495 | | 1662 | |
| | 1663 | pub fn io_uring_prep_close_direct(sqe: *linux.io_uring_sqe, file_index: u32) void { |
| | 1664 | io_uring_prep_close(sqe, 0); |
| | 1665 | __io_uring_set_target_fixed_file(sqe, file_index); |
| | 1666 | } |
| | 1667 | |
| 1496 | pub fn io_uring_prep_timeout( | 1668 | pub fn io_uring_prep_timeout( |
| 1497 | sqe: *linux.io_uring_sqe, | 1669 | sqe: *linux.io_uring_sqe, |
| 1498 | ts: *const os.linux.kernel_timespec, | 1670 | ts: *const os.linux.kernel_timespec, |
| ... | @@ -1720,6 +1892,40 @@ pub fn io_uring_prep_multishot_accept( | ... | @@ -1720,6 +1892,40 @@ pub fn io_uring_prep_multishot_accept( |
| 1720 | sqe.ioprio |= linux.IORING_ACCEPT_MULTISHOT; | 1892 | sqe.ioprio |= linux.IORING_ACCEPT_MULTISHOT; |
| 1721 | } | 1893 | } |
| 1722 | | 1894 | |
| | 1895 | pub fn io_uring_prep_socket( |
| | 1896 | sqe: *linux.io_uring_sqe, |
| | 1897 | domain: u32, |
| | 1898 | socket_type: u32, |
| | 1899 | protocol: u32, |
| | 1900 | flags: u32, |
| | 1901 | ) void { |
| | 1902 | io_uring_prep_rw(.SOCKET, sqe, @intCast(domain), 0, protocol, socket_type); |
| | 1903 | sqe.rw_flags = flags; |
| | 1904 | } |
| | 1905 | |
| | 1906 | pub fn io_uring_prep_socket_direct( |
| | 1907 | sqe: *linux.io_uring_sqe, |
| | 1908 | domain: u32, |
| | 1909 | socket_type: u32, |
| | 1910 | protocol: u32, |
| | 1911 | flags: u32, |
| | 1912 | file_index: u32, |
| | 1913 | ) void { |
| | 1914 | io_uring_prep_socket(sqe, domain, socket_type, protocol, flags); |
| | 1915 | __io_uring_set_target_fixed_file(sqe, file_index); |
| | 1916 | } |
| | 1917 | |
| | 1918 | pub fn io_uring_prep_socket_direct_alloc( |
| | 1919 | sqe: *linux.io_uring_sqe, |
| | 1920 | domain: u32, |
| | 1921 | socket_type: u32, |
| | 1922 | protocol: u32, |
| | 1923 | flags: u32, |
| | 1924 | ) void { |
| | 1925 | io_uring_prep_socket(sqe, domain, socket_type, protocol, flags); |
| | 1926 | __io_uring_set_target_fixed_file(sqe, linux.IORING_FILE_INDEX_ALLOC); |
| | 1927 | } |
| | 1928 | |
| 1723 | test "structs/offsets/entries" { | 1929 | test "structs/offsets/entries" { |
| 1724 | if (builtin.os.tag != .linux) return error.SkipZigTest; | 1930 | if (builtin.os.tag != .linux) return error.SkipZigTest; |
| 1725 | | 1931 | |
| ... | @@ -3582,7 +3788,8 @@ test "accept/connect/send_zc/recv" { | ... | @@ -3582,7 +3788,8 @@ test "accept/connect/send_zc/recv" { |
| 3582 | const send = try ring.send_zc(0xeeeeeeee, socket_test_harness.client, buffer_send[0..], 0, 0); | 3788 | const send = try ring.send_zc(0xeeeeeeee, socket_test_harness.client, buffer_send[0..], 0, 0); |
| 3583 | send.flags |= linux.IOSQE_IO_LINK; | 3789 | send.flags |= linux.IOSQE_IO_LINK; |
| 3584 | _ = try ring.recv(0xffffffff, socket_test_harness.server, .{ .buffer = buffer_recv[0..] }, 0); | 3790 | _ = try ring.recv(0xffffffff, socket_test_harness.server, .{ .buffer = buffer_recv[0..] }, 0); |
| 3585 | try testing.expectEqual(@as(u32, 2), try ring.submit()); | 3791 | const submitted = try ring.submit(); |
| | 3792 | if (submitted != 2) return error.SkipZigTest; // on kernel 5.8 (without zc support) |
| 3586 | | 3793 | |
| 3587 | // First completion of zero-copy send. | 3794 | // First completion of zero-copy send. |
| 3588 | // IORING_CQE_F_MORE, means that there | 3795 | // IORING_CQE_F_MORE, means that there |
| ... | @@ -3617,3 +3824,295 @@ test "accept/connect/send_zc/recv" { | ... | @@ -3617,3 +3824,295 @@ test "accept/connect/send_zc/recv" { |
| 3617 | .flags = linux.IORING_CQE_F_NOTIF, | 3824 | .flags = linux.IORING_CQE_F_NOTIF, |
| 3618 | }, cqe_send); | 3825 | }, cqe_send); |
| 3619 | } | 3826 | } |
| | 3827 | |
| | 3828 | test "accept_direct" { |
| | 3829 | if (builtin.os.tag != .linux) return error.SkipZigTest; |
| | 3830 | |
| | 3831 | var ring = IO_Uring.init(1, 0) catch |err| switch (err) { |
| | 3832 | error.SystemOutdated => return error.SkipZigTest, |
| | 3833 | error.PermissionDenied => return error.SkipZigTest, |
| | 3834 | else => return err, |
| | 3835 | }; |
| | 3836 | defer ring.deinit(); |
| | 3837 | var address = try net.Address.parseIp4("127.0.0.1", 0); |
| | 3838 | |
| | 3839 | // register direct file descriptors |
| | 3840 | var registered_fds = [_]os.fd_t{-1} ** 2; |
| | 3841 | ring.register_files(registered_fds[0..]) catch |err| switch (err) { |
| | 3842 | error.FileDescriptorInvalid => return error.SkipZigTest, |
| | 3843 | else => return err, |
| | 3844 | }; |
| | 3845 | |
| | 3846 | const listener_socket = try createListenerSocket(&address); |
| | 3847 | defer os.closeSocket(listener_socket); |
| | 3848 | |
| | 3849 | const accept_userdata: u64 = 0xaaaaaaaa; |
| | 3850 | const read_userdata: u64 = 0xbbbbbbbb; |
| | 3851 | const data = [_]u8{ 0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 0xa, 0xb, 0xc, 0xd, 0xe }; |
| | 3852 | |
| | 3853 | for (0..2) |_| { |
| | 3854 | for (registered_fds, 0..) |_, i| { |
| | 3855 | var buffer_recv = [_]u8{0} ** 16; |
| | 3856 | const buffer_send: []const u8 = data[0 .. data.len - i]; // make it different at each loop |
| | 3857 | |
| | 3858 | // submit accept, will chose registered fd and return index in cqe |
| | 3859 | _ = try ring.accept_direct(accept_userdata, listener_socket, null, null, 0); |
| | 3860 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| | 3861 | |
| | 3862 | // connect |
| | 3863 | var client = try os.socket(address.any.family, os.SOCK.STREAM | os.SOCK.CLOEXEC, 0); |
| | 3864 | try os.connect(client, &address.any, address.getOsSockLen()); |
| | 3865 | defer os.closeSocket(client); |
| | 3866 | |
| | 3867 | // accept completion |
| | 3868 | const cqe_accept = try ring.copy_cqe(); |
| | 3869 | if (cqe_accept.err() == .INVAL) return error.SkipZigTest; |
| | 3870 | try testing.expectEqual(os.E.SUCCESS, cqe_accept.err()); |
| | 3871 | const fd_index = cqe_accept.res; |
| | 3872 | if (fd_index >= registered_fds.len) return error.SkipZigTest; // old kernel fallback to ordinary accept |
| | 3873 | try testing.expect(fd_index < registered_fds.len); |
| | 3874 | try testing.expect(cqe_accept.user_data == accept_userdata); |
| | 3875 | |
| | 3876 | // send data |
| | 3877 | _ = try os.send(client, buffer_send, 0); |
| | 3878 | |
| | 3879 | // Example of how to use registered fd: |
| | 3880 | // Submit receive to fixed file returned by accept (fd_index). |
| | 3881 | // Fd field is set to registered file index, returned by accept. |
| | 3882 | // Flag linux.IOSQE_FIXED_FILE must be set. |
| | 3883 | const recv_sqe = try ring.recv(read_userdata, fd_index, .{ .buffer = &buffer_recv }, 0); |
| | 3884 | recv_sqe.flags |= linux.IOSQE_FIXED_FILE; |
| | 3885 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| | 3886 | |
| | 3887 | // accept receive |
| | 3888 | const recv_cqe = try ring.copy_cqe(); |
| | 3889 | try testing.expect(recv_cqe.user_data == read_userdata); |
| | 3890 | try testing.expect(recv_cqe.res == buffer_send.len); |
| | 3891 | try testing.expectEqualSlices(u8, buffer_send, buffer_recv[0..buffer_send.len]); |
| | 3892 | } |
| | 3893 | // no more available fds, accept will get NFILE error |
| | 3894 | { |
| | 3895 | // submit accept |
| | 3896 | _ = try ring.accept_direct(accept_userdata, listener_socket, null, null, 0); |
| | 3897 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| | 3898 | // connect |
| | 3899 | var client = try os.socket(address.any.family, os.SOCK.STREAM | os.SOCK.CLOEXEC, 0); |
| | 3900 | try os.connect(client, &address.any, address.getOsSockLen()); |
| | 3901 | defer os.closeSocket(client); |
| | 3902 | // completion with error |
| | 3903 | const cqe_accept = try ring.copy_cqe(); |
| | 3904 | try testing.expect(cqe_accept.user_data == accept_userdata); |
| | 3905 | try testing.expectEqual(os.E.NFILE, cqe_accept.err()); |
| | 3906 | } |
| | 3907 | // return file descriptors to kernel |
| | 3908 | try ring.register_files_update(0, registered_fds[0..]); |
| | 3909 | } |
| | 3910 | try ring.unregister_files(); |
| | 3911 | } |
| | 3912 | |
| | 3913 | test "accept_multishot_direct" { |
| | 3914 | if (builtin.os.tag != .linux) return error.SkipZigTest; |
| | 3915 | |
| | 3916 | var ring = IO_Uring.init(1, 0) catch |err| switch (err) { |
| | 3917 | error.SystemOutdated => return error.SkipZigTest, |
| | 3918 | error.PermissionDenied => return error.SkipZigTest, |
| | 3919 | else => return err, |
| | 3920 | }; |
| | 3921 | defer ring.deinit(); |
| | 3922 | var address = try net.Address.parseIp4("127.0.0.1", 0); |
| | 3923 | |
| | 3924 | var registered_fds = [_]os.fd_t{-1} ** 2; |
| | 3925 | ring.register_files(registered_fds[0..]) catch |err| switch (err) { |
| | 3926 | error.FileDescriptorInvalid => return error.SkipZigTest, |
| | 3927 | else => return err, |
| | 3928 | }; |
| | 3929 | |
| | 3930 | const listener_socket = try createListenerSocket(&address); |
| | 3931 | defer os.closeSocket(listener_socket); |
| | 3932 | |
| | 3933 | const accept_userdata: u64 = 0xaaaaaaaa; |
| | 3934 | |
| | 3935 | for (0..2) |_| { |
| | 3936 | // submit accept |
| | 3937 | // Will chose registered fd and return index of the selected registered file in cqe. |
| | 3938 | _ = try ring.accept_multishot_direct(accept_userdata, listener_socket, null, null, 0); |
| | 3939 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| | 3940 | |
| | 3941 | for (registered_fds) |_| { |
| | 3942 | // connect |
| | 3943 | var client = try os.socket(address.any.family, os.SOCK.STREAM | os.SOCK.CLOEXEC, 0); |
| | 3944 | try os.connect(client, &address.any, address.getOsSockLen()); |
| | 3945 | defer os.closeSocket(client); |
| | 3946 | |
| | 3947 | // accept completion |
| | 3948 | const cqe_accept = try ring.copy_cqe(); |
| | 3949 | if (cqe_accept.err() == .INVAL) return error.SkipZigTest; |
| | 3950 | const fd_index = cqe_accept.res; |
| | 3951 | try testing.expect(fd_index < registered_fds.len); |
| | 3952 | try testing.expect(cqe_accept.user_data == accept_userdata); |
| | 3953 | try testing.expect(cqe_accept.flags & linux.IORING_CQE_F_MORE > 0); // has more is set |
| | 3954 | } |
| | 3955 | // No more available fds, accept will get NFILE error. |
| | 3956 | // Multishot is terminated (more flag is not set). |
| | 3957 | { |
| | 3958 | // connect |
| | 3959 | var client = try os.socket(address.any.family, os.SOCK.STREAM | os.SOCK.CLOEXEC, 0); |
| | 3960 | try os.connect(client, &address.any, address.getOsSockLen()); |
| | 3961 | defer os.closeSocket(client); |
| | 3962 | // completion with error |
| | 3963 | const cqe_accept = try ring.copy_cqe(); |
| | 3964 | try testing.expect(cqe_accept.user_data == accept_userdata); |
| | 3965 | try testing.expectEqual(os.E.NFILE, cqe_accept.err()); |
| | 3966 | try testing.expect(cqe_accept.flags & linux.IORING_CQE_F_MORE == 0); // has more is not set |
| | 3967 | } |
| | 3968 | // return file descriptors to kernel |
| | 3969 | try ring.register_files_update(0, registered_fds[0..]); |
| | 3970 | } |
| | 3971 | try ring.unregister_files(); |
| | 3972 | } |
| | 3973 | |
| | 3974 | test "socket/socket_direct/socket_direct_alloc/close_direct" { |
| | 3975 | if (builtin.os.tag != .linux) return error.SkipZigTest; |
| | 3976 | |
| | 3977 | var ring = IO_Uring.init(2, 0) catch |err| switch (err) { |
| | 3978 | error.SystemOutdated => return error.SkipZigTest, |
| | 3979 | error.PermissionDenied => return error.SkipZigTest, |
| | 3980 | else => return err, |
| | 3981 | }; |
| | 3982 | defer ring.deinit(); |
| | 3983 | var address = try net.Address.parseIp4("127.0.0.1", 0); |
| | 3984 | |
| | 3985 | // Below are 4 different ways to get socket fd. |
| | 3986 | // Two upfront before register_files, and two after |
| | 3987 | var registered_fds = [_]os.fd_t{-1} ** 4; |
| | 3988 | // 1. sync syscall socket call |
| | 3989 | registered_fds[0] = try os.socket(address.any.family, os.SOCK.STREAM | os.SOCK.CLOEXEC, 0); |
| | 3990 | // 2. io_uring socket |
| | 3991 | const socket_userdata = 0xcccccccc; |
| | 3992 | _ = try ring.socket(socket_userdata, linux.AF.INET, os.SOCK.STREAM, 0, 0); |
| | 3993 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| | 3994 | var cqe_socket = try ring.copy_cqe(); |
| | 3995 | if (cqe_socket.err() == .INVAL) return error.SkipZigTest; |
| | 3996 | try testing.expectEqual(os.E.SUCCESS, cqe_socket.err()); |
| | 3997 | try testing.expect(cqe_socket.res > 2); |
| | 3998 | registered_fds[1] = cqe_socket.res; // set index 1 to created socket |
| | 3999 | |
| | 4000 | ring.register_files(registered_fds[0..]) catch |err| switch (err) { |
| | 4001 | error.FileDescriptorInvalid => return error.SkipZigTest, |
| | 4002 | else => return err, |
| | 4003 | }; |
| | 4004 | |
| | 4005 | // 3. io_uring socket_direct, create new socket on index 2 |
| | 4006 | _ = try ring.socket_direct(socket_userdata, linux.AF.INET, os.SOCK.STREAM, 0, 0, @intCast(2)); |
| | 4007 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| | 4008 | cqe_socket = try ring.copy_cqe(); |
| | 4009 | try testing.expectEqual(os.E.SUCCESS, cqe_socket.err()); |
| | 4010 | try testing.expect(cqe_socket.res == 0); |
| | 4011 | |
| | 4012 | // 4. io_uring socket_direct_alloc |
| | 4013 | _ = try ring.socket_direct_alloc(socket_userdata, linux.AF.INET, os.SOCK.STREAM, 0, 0); |
| | 4014 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| | 4015 | cqe_socket = try ring.copy_cqe(); |
| | 4016 | try testing.expectEqual(os.E.SUCCESS, cqe_socket.err()); |
| | 4017 | try testing.expect(cqe_socket.res == 3); |
| | 4018 | |
| | 4019 | // use sockets from registered_fds in connect operation |
| | 4020 | const listener_socket = try createListenerSocket(&address); |
| | 4021 | defer os.closeSocket(listener_socket); |
| | 4022 | const accept_userdata: u64 = 0xaaaaaaaa; |
| | 4023 | const connect_userdata: u64 = 0xbbbbbbbb; |
| | 4024 | const close_userdata: u64 = 0xcccccccc; |
| | 4025 | for (registered_fds, 0..) |_, fd_index| { |
| | 4026 | // prepare accept |
| | 4027 | _ = try ring.accept(accept_userdata, listener_socket, null, null, 0); |
| | 4028 | // prepare connect with fixed socket |
| | 4029 | const connect_sqe = try ring.connect(connect_userdata, @intCast(fd_index), &address.any, address.getOsSockLen()); |
| | 4030 | connect_sqe.flags |= linux.IOSQE_FIXED_FILE; |
| | 4031 | // submit both |
| | 4032 | try testing.expectEqual(@as(u32, 2), try ring.submit()); |
| | 4033 | // get completions |
| | 4034 | var cqe_connect = try ring.copy_cqe(); |
| | 4035 | if (cqe_connect.err() == .INVAL) return error.SkipZigTest; |
| | 4036 | var cqe_accept = try ring.copy_cqe(); |
| | 4037 | if (cqe_accept.err() == .INVAL) return error.SkipZigTest; |
| | 4038 | // ignore order |
| | 4039 | if (cqe_connect.user_data == accept_userdata and cqe_accept.user_data == connect_userdata) { |
| | 4040 | const a = cqe_accept; |
| | 4041 | const b = cqe_connect; |
| | 4042 | cqe_accept = b; |
| | 4043 | cqe_connect = a; |
| | 4044 | } |
| | 4045 | // test connect completion |
| | 4046 | try testing.expect(cqe_connect.user_data == connect_userdata); |
| | 4047 | try testing.expectEqual(os.E.SUCCESS, cqe_connect.err()); |
| | 4048 | // test accept completion |
| | 4049 | try testing.expect(cqe_accept.user_data == accept_userdata); |
| | 4050 | try testing.expectEqual(os.E.SUCCESS, cqe_accept.err()); |
| | 4051 | |
| | 4052 | // submit and test close completion |
| | 4053 | _ = try ring.close_direct(close_userdata, @intCast(fd_index)); |
| | 4054 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| | 4055 | var cqe_close = try ring.copy_cqe(); |
| | 4056 | try testing.expect(cqe_close.user_data == close_userdata); |
| | 4057 | try testing.expectEqual(os.E.SUCCESS, cqe_close.err()); |
| | 4058 | } |
| | 4059 | |
| | 4060 | try ring.unregister_files(); |
| | 4061 | } |
| | 4062 | |
| | 4063 | test "openat_direct/close_direct" { |
| | 4064 | if (builtin.os.tag != .linux) return error.SkipZigTest; |
| | 4065 | |
| | 4066 | var ring = IO_Uring.init(2, 0) catch |err| switch (err) { |
| | 4067 | error.SystemOutdated => return error.SkipZigTest, |
| | 4068 | error.PermissionDenied => return error.SkipZigTest, |
| | 4069 | else => return err, |
| | 4070 | }; |
| | 4071 | defer ring.deinit(); |
| | 4072 | |
| | 4073 | var registered_fds = [_]os.fd_t{-1} ** 3; |
| | 4074 | ring.register_files(registered_fds[0..]) catch |err| switch (err) { |
| | 4075 | error.FileDescriptorInvalid => return error.SkipZigTest, |
| | 4076 | else => return err, |
| | 4077 | }; |
| | 4078 | |
| | 4079 | var tmp = std.testing.tmpDir(.{}); |
| | 4080 | defer tmp.cleanup(); |
| | 4081 | const path = "test_io_uring_close_direct"; |
| | 4082 | const flags: u32 = os.O.RDWR | os.O.CREAT; |
| | 4083 | const mode: os.mode_t = 0o666; |
| | 4084 | const user_data: u64 = 0; |
| | 4085 | |
| | 4086 | // use registered file at index 0 (last param) |
| | 4087 | _ = try ring.openat_direct(user_data, tmp.dir.fd, path, flags, mode, 0); |
| | 4088 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| | 4089 | var cqe = try ring.copy_cqe(); |
| | 4090 | if (cqe.err() == .INVAL) return error.SkipZigTest; |
| | 4091 | if (cqe.res != 0) return error.SkipZigTest; // old kernel fallback to openat without direct |
| | 4092 | try testing.expectEqual(os.E.SUCCESS, cqe.err()); |
| | 4093 | try testing.expect(cqe.res == 0); |
| | 4094 | |
| | 4095 | // use registered file at index 1 |
| | 4096 | _ = try ring.openat_direct(user_data, tmp.dir.fd, path, flags, mode, 1); |
| | 4097 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| | 4098 | cqe = try ring.copy_cqe(); |
| | 4099 | try testing.expectEqual(os.E.SUCCESS, cqe.err()); |
| | 4100 | try testing.expect(cqe.res == 0); // res is 0 when we specify index |
| | 4101 | |
| | 4102 | // let kernel choose registered file index |
| | 4103 | _ = try ring.openat_direct(user_data, tmp.dir.fd, path, flags, mode, linux.IORING_FILE_INDEX_ALLOC); |
| | 4104 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| | 4105 | cqe = try ring.copy_cqe(); |
| | 4106 | if (cqe.err() == .INVAL) return error.SkipZigTest; // kernel 5.15 bug |
| | 4107 | try testing.expectEqual(os.E.SUCCESS, cqe.err()); |
| | 4108 | try testing.expect(cqe.res == 2); // chosen index is in res |
| | 4109 | |
| | 4110 | // close all open file descriptors |
| | 4111 | for (registered_fds, 0..) |_, fd_index| { |
| | 4112 | _ = try ring.close_direct(user_data, @intCast(fd_index)); |
| | 4113 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| | 4114 | var cqe_close = try ring.copy_cqe(); |
| | 4115 | try testing.expectEqual(os.E.SUCCESS, cqe_close.err()); |
| | 4116 | } |
| | 4117 | try ring.unregister_files(); |
| | 4118 | } |