| ... | ... | @@ -491,6 +491,7 @@ pub const IO_Uring = struct { |
| 491 | 491 | |
| 492 | 492 | /// Queues (but does not submit) an SQE to perform an `accept4(2)` on a socket. |
| 493 | 493 | /// Returns a pointer to the SQE. |
| 494 | /// Available since 5.5 |
| 494 | 495 | pub fn accept( |
| 495 | 496 | self: *IO_Uring, |
| 496 | 497 | user_data: u64, |
| ... | ... | @@ -505,10 +506,14 @@ pub const IO_Uring = struct { |
| 505 | 506 | return sqe; |
| 506 | 507 | } |
| 507 | 508 | |
| 508 | | /// Queues (but does not submit) an SQE to perform an multishot `accept4(2)` on a socket. |
| 509 | /// Queues an multishot accept on a socket. |
| 510 | /// |
| 509 | 511 | /// Multishot variant allows an application to issue a single accept request, |
| 510 | 512 | /// which will repeatedly trigger a CQE when a connection request comes in. |
| 511 | | /// Returns a pointer to the SQE. |
| 513 | /// While IORING_CQE_F_MORE flag is set in CQE flags accept will generate |
| 514 | /// further CQEs. |
| 515 | /// |
| 516 | /// Available since 5.19 |
| 512 | 517 | pub fn accept_multishot( |
| 513 | 518 | self: *IO_Uring, |
| 514 | 519 | user_data: u64, |
| ... | ... | @@ -523,6 +528,47 @@ pub const IO_Uring = struct { |
| 523 | 528 | return sqe; |
| 524 | 529 | } |
| 525 | 530 | |
| 531 | /// Queues an accept using direct (registered) file descriptors. |
| 532 | /// |
| 533 | /// To use an accept direct variant, the application must first have registered |
| 534 | /// a file table (with register_files). An unused table index will be |
| 535 | /// dynamically chosen and returned in the CQE res field. |
| 536 | /// |
| 537 | /// After creation, they can be used by setting IOSQE_FIXED_FILE in the SQE |
| 538 | /// flags member, and setting the SQE fd field to the direct descriptor value |
| 539 | /// rather than the regular file descriptor. |
| 540 | /// |
| 541 | /// Available since 5.19 |
| 542 | pub fn accept_direct( |
| 543 | self: *IO_Uring, |
| 544 | user_data: u64, |
| 545 | fd: os.fd_t, |
| 546 | addr: ?*os.sockaddr, |
| 547 | addrlen: ?*os.socklen_t, |
| 548 | flags: u32, |
| 549 | ) !*linux.io_uring_sqe { |
| 550 | const sqe = try self.get_sqe(); |
| 551 | io_uring_prep_accept_direct(sqe, fd, addr, addrlen, flags, linux.IORING_FILE_INDEX_ALLOC); |
| 552 | sqe.user_data = user_data; |
| 553 | return sqe; |
| 554 | } |
| 555 | |
| 556 | /// Queues an multishot accept using direct (registered) file descriptors. |
| 557 | /// Available since 5.19 |
| 558 | pub fn accept_multishot_direct( |
| 559 | self: *IO_Uring, |
| 560 | user_data: u64, |
| 561 | fd: os.fd_t, |
| 562 | addr: ?*os.sockaddr, |
| 563 | addrlen: ?*os.socklen_t, |
| 564 | flags: u32, |
| 565 | ) !*linux.io_uring_sqe { |
| 566 | const sqe = try self.get_sqe(); |
| 567 | io_uring_prep_multishot_accept_direct(sqe, fd, addr, addrlen, flags); |
| 568 | sqe.user_data = user_data; |
| 569 | return sqe; |
| 570 | } |
| 571 | |
| 526 | 572 | /// Queue (but does not submit) an SQE to perform a `connect(2)` on a socket. |
| 527 | 573 | /// Returns a pointer to the SQE. |
| 528 | 574 | pub fn connect( |
| ... | ... | @@ -570,6 +616,7 @@ pub const IO_Uring = struct { |
| 570 | 616 | |
| 571 | 617 | /// Queues (but does not submit) an SQE to perform a `recv(2)`. |
| 572 | 618 | /// Returns a pointer to the SQE. |
| 619 | /// Available since 5.6 |
| 573 | 620 | pub fn recv( |
| 574 | 621 | self: *IO_Uring, |
| 575 | 622 | user_data: u64, |
| ... | ... | @@ -593,6 +640,7 @@ pub const IO_Uring = struct { |
| 593 | 640 | |
| 594 | 641 | /// Queues (but does not submit) an SQE to perform a `send(2)`. |
| 595 | 642 | /// Returns a pointer to the SQE. |
| 643 | /// Available since 5.6 |
| 596 | 644 | pub fn send( |
| 597 | 645 | self: *IO_Uring, |
| 598 | 646 | user_data: u64, |
| ... | ... | @@ -606,8 +654,56 @@ pub const IO_Uring = struct { |
| 606 | 654 | return sqe; |
| 607 | 655 | } |
| 608 | 656 | |
| 657 | /// Queues (but does not submit) an SQE to perform an async zerocopy `send(2)`. |
| 658 | /// |
| 659 | /// This operation will most likely produce two CQEs. The flags field of the |
| 660 | /// first cqe may likely contain IORING_CQE_F_MORE, which means that there will |
| 661 | /// be a second cqe with the user_data field set to the same value. The user |
| 662 | /// must not modify the data buffer until the notification is posted. The first |
| 663 | /// cqe follows the usual rules and so its res field will contain the number of |
| 664 | /// bytes sent or a negative error code. The notification's res field will be |
| 665 | /// set to zero and the flags field will contain IORING_CQE_F_NOTIF. The two |
| 666 | /// step model is needed because the kernel may hold on to buffers for a long |
| 667 | /// time, e.g. waiting for a TCP ACK. Notifications responsible for controlling |
| 668 | /// the lifetime of the buffers. Even errored requests may generate a |
| 669 | /// notification. |
| 670 | /// |
| 671 | /// Available since 6.0 |
| 672 | pub fn send_zc( |
| 673 | self: *IO_Uring, |
| 674 | user_data: u64, |
| 675 | fd: os.fd_t, |
| 676 | buffer: []const u8, |
| 677 | send_flags: u32, |
| 678 | zc_flags: u16, |
| 679 | ) !*linux.io_uring_sqe { |
| 680 | const sqe = try self.get_sqe(); |
| 681 | io_uring_prep_send_zc(sqe, fd, buffer, send_flags, zc_flags); |
| 682 | sqe.user_data = user_data; |
| 683 | return sqe; |
| 684 | } |
| 685 | |
| 686 | /// Queues (but does not submit) an SQE to perform an async zerocopy `send(2)`. |
| 687 | /// Returns a pointer to the SQE. |
| 688 | /// Available since 6.0 |
| 689 | pub fn send_zc_fixed( |
| 690 | self: *IO_Uring, |
| 691 | user_data: u64, |
| 692 | fd: os.fd_t, |
| 693 | buffer: []const u8, |
| 694 | send_flags: u32, |
| 695 | zc_flags: u16, |
| 696 | buf_index: u16, |
| 697 | ) !*linux.io_uring_sqe { |
| 698 | const sqe = try self.get_sqe(); |
| 699 | io_uring_prep_send_zc_fixed(sqe, fd, buffer, send_flags, zc_flags, buf_index); |
| 700 | sqe.user_data = user_data; |
| 701 | return sqe; |
| 702 | } |
| 703 | |
| 609 | 704 | /// Queues (but does not submit) an SQE to perform a `recvmsg(2)`. |
| 610 | 705 | /// Returns a pointer to the SQE. |
| 706 | /// Available since 5.3 |
| 611 | 707 | pub fn recvmsg( |
| 612 | 708 | self: *IO_Uring, |
| 613 | 709 | user_data: u64, |
| ... | ... | @@ -623,6 +719,7 @@ pub const IO_Uring = struct { |
| 623 | 719 | |
| 624 | 720 | /// Queues (but does not submit) an SQE to perform a `sendmsg(2)`. |
| 625 | 721 | /// Returns a pointer to the SQE. |
| 722 | /// Available since 5.3 |
| 626 | 723 | pub fn sendmsg( |
| 627 | 724 | self: *IO_Uring, |
| 628 | 725 | user_data: u64, |
| ... | ... | @@ -636,8 +733,25 @@ pub const IO_Uring = struct { |
| 636 | 733 | return sqe; |
| 637 | 734 | } |
| 638 | 735 | |
| 736 | /// Queues (but does not submit) an SQE to perform an async zerocopy `sendmsg(2)`. |
| 737 | /// Returns a pointer to the SQE. |
| 738 | /// Available since 6.1 |
| 739 | pub fn sendmsg_zc( |
| 740 | self: *IO_Uring, |
| 741 | user_data: u64, |
| 742 | fd: os.fd_t, |
| 743 | msg: *const os.msghdr_const, |
| 744 | flags: u32, |
| 745 | ) !*linux.io_uring_sqe { |
| 746 | const sqe = try self.get_sqe(); |
| 747 | io_uring_prep_sendmsg_zc(sqe, fd, msg, flags); |
| 748 | sqe.user_data = user_data; |
| 749 | return sqe; |
| 750 | } |
| 751 | |
| 639 | 752 | /// Queues (but does not submit) an SQE to perform an `openat(2)`. |
| 640 | 753 | /// Returns a pointer to the SQE. |
| 754 | /// Available since 5.6. |
| 641 | 755 | pub fn openat( |
| 642 | 756 | self: *IO_Uring, |
| 643 | 757 | user_data: u64, |
| ... | ... | @@ -652,8 +766,35 @@ pub const IO_Uring = struct { |
| 652 | 766 | return sqe; |
| 653 | 767 | } |
| 654 | 768 | |
| 769 | /// Queues an openat using direct (registered) file descriptors. |
| 770 | /// |
| 771 | /// To use an accept direct variant, the application must first have registered |
| 772 | /// a file table (with register_files). An unused table index will be |
| 773 | /// dynamically chosen and returned in the CQE res field. |
| 774 | /// |
| 775 | /// After creation, they can be used by setting IOSQE_FIXED_FILE in the SQE |
| 776 | /// flags member, and setting the SQE fd field to the direct descriptor value |
| 777 | /// rather than the regular file descriptor. |
| 778 | /// |
| 779 | /// Available since 5.15 |
| 780 | pub fn openat_direct( |
| 781 | self: *IO_Uring, |
| 782 | user_data: u64, |
| 783 | fd: os.fd_t, |
| 784 | path: [*:0]const u8, |
| 785 | flags: u32, |
| 786 | mode: os.mode_t, |
| 787 | file_index: u32, |
| 788 | ) !*linux.io_uring_sqe { |
| 789 | const sqe = try self.get_sqe(); |
| 790 | io_uring_prep_openat_direct(sqe, fd, path, flags, mode, file_index); |
| 791 | sqe.user_data = user_data; |
| 792 | return sqe; |
| 793 | } |
| 794 | |
| 655 | 795 | /// Queues (but does not submit) an SQE to perform a `close(2)`. |
| 656 | 796 | /// Returns a pointer to the SQE. |
| 797 | /// Available since 5.6. |
| 657 | 798 | pub fn close(self: *IO_Uring, user_data: u64, fd: os.fd_t) !*linux.io_uring_sqe { |
| 658 | 799 | const sqe = try self.get_sqe(); |
| 659 | 800 | io_uring_prep_close(sqe, fd); |
| ... | ... | @@ -661,6 +802,15 @@ pub const IO_Uring = struct { |
| 661 | 802 | return sqe; |
| 662 | 803 | } |
| 663 | 804 | |
| 805 | /// Queues close of registered file descriptor. |
| 806 | /// Available since 5.15 |
| 807 | pub fn close_direct(self: *IO_Uring, user_data: u64, file_index: u32) !*linux.io_uring_sqe { |
| 808 | const sqe = try self.get_sqe(); |
| 809 | io_uring_prep_close_direct(sqe, file_index); |
| 810 | sqe.user_data = user_data; |
| 811 | return sqe; |
| 812 | } |
| 813 | |
| 664 | 814 | /// Queues (but does not submit) an SQE to register a timeout operation. |
| 665 | 815 | /// Returns a pointer to the SQE. |
| 666 | 816 | /// |
| ... | ... | @@ -1109,6 +1259,57 @@ pub const IO_Uring = struct { |
| 1109 | 1259 | else => |errno| return os.unexpectedErrno(errno), |
| 1110 | 1260 | } |
| 1111 | 1261 | } |
| 1262 | |
| 1263 | /// Prepares a socket creation request. |
| 1264 | /// New socket fd will be returned in completion result. |
| 1265 | /// Available since 5.19 |
| 1266 | pub fn socket( |
| 1267 | self: *IO_Uring, |
| 1268 | user_data: u64, |
| 1269 | domain: u32, |
| 1270 | socket_type: u32, |
| 1271 | protocol: u32, |
| 1272 | flags: u32, |
| 1273 | ) !*linux.io_uring_sqe { |
| 1274 | const sqe = try self.get_sqe(); |
| 1275 | io_uring_prep_socket(sqe, domain, socket_type, protocol, flags); |
| 1276 | sqe.user_data = user_data; |
| 1277 | return sqe; |
| 1278 | } |
| 1279 | |
| 1280 | /// Prepares a socket creation request for registered file at index `file_index`. |
| 1281 | /// Available since 5.19 |
| 1282 | pub fn socket_direct( |
| 1283 | self: *IO_Uring, |
| 1284 | user_data: u64, |
| 1285 | domain: u32, |
| 1286 | socket_type: u32, |
| 1287 | protocol: u32, |
| 1288 | flags: u32, |
| 1289 | file_index: u32, |
| 1290 | ) !*linux.io_uring_sqe { |
| 1291 | const sqe = try self.get_sqe(); |
| 1292 | io_uring_prep_socket_direct(sqe, domain, socket_type, protocol, flags, file_index); |
| 1293 | sqe.user_data = user_data; |
| 1294 | return sqe; |
| 1295 | } |
| 1296 | |
| 1297 | /// Prepares a socket creation request for registered file, index chosen by kernel (file index alloc). |
| 1298 | /// File index will be returned in CQE res field. |
| 1299 | /// Available since 5.19 |
| 1300 | pub fn socket_direct_alloc( |
| 1301 | self: *IO_Uring, |
| 1302 | user_data: u64, |
| 1303 | domain: u32, |
| 1304 | socket_type: u32, |
| 1305 | protocol: u32, |
| 1306 | flags: u32, |
| 1307 | ) !*linux.io_uring_sqe { |
| 1308 | const sqe = try self.get_sqe(); |
| 1309 | io_uring_prep_socket_direct_alloc(sqe, domain, socket_type, protocol, flags); |
| 1310 | sqe.user_data = user_data; |
| 1311 | return sqe; |
| 1312 | } |
| 1112 | 1313 | }; |
| 1113 | 1314 | |
| 1114 | 1315 | pub const SubmissionQueue = struct { |
| ... | ... | @@ -1343,6 +1544,41 @@ pub fn io_uring_prep_accept( |
| 1343 | 1544 | sqe.rw_flags = flags; |
| 1344 | 1545 | } |
| 1345 | 1546 | |
| 1547 | pub fn io_uring_prep_accept_direct( |
| 1548 | sqe: *linux.io_uring_sqe, |
| 1549 | fd: os.fd_t, |
| 1550 | addr: ?*os.sockaddr, |
| 1551 | addrlen: ?*os.socklen_t, |
| 1552 | flags: u32, |
| 1553 | file_index: u32, |
| 1554 | ) void { |
| 1555 | io_uring_prep_accept(sqe, fd, addr, addrlen, flags); |
| 1556 | __io_uring_set_target_fixed_file(sqe, file_index); |
| 1557 | } |
| 1558 | |
| 1559 | pub fn io_uring_prep_multishot_accept_direct( |
| 1560 | sqe: *linux.io_uring_sqe, |
| 1561 | fd: os.fd_t, |
| 1562 | addr: ?*os.sockaddr, |
| 1563 | addrlen: ?*os.socklen_t, |
| 1564 | flags: u32, |
| 1565 | ) void { |
| 1566 | io_uring_prep_multishot_accept(sqe, fd, addr, addrlen, flags); |
| 1567 | __io_uring_set_target_fixed_file(sqe, linux.IORING_FILE_INDEX_ALLOC); |
| 1568 | } |
| 1569 | |
| 1570 | fn __io_uring_set_target_fixed_file(sqe: *linux.io_uring_sqe, file_index: u32) void { |
| 1571 | const sqe_file_index: u32 = if (file_index == linux.IORING_FILE_INDEX_ALLOC) |
| 1572 | linux.IORING_FILE_INDEX_ALLOC |
| 1573 | else |
| 1574 | // 0 means no fixed files, indexes should be encoded as "index + 1" |
| 1575 | file_index + 1; |
| 1576 | // This filed is overloaded in liburing: |
| 1577 | // splice_fd_in: i32 |
| 1578 | // sqe_file_index: u32 |
| 1579 | sqe.splice_fd_in = @bitCast(sqe_file_index); |
| 1580 | } |
| 1581 | |
| 1346 | 1582 | pub fn io_uring_prep_connect( |
| 1347 | 1583 | sqe: *linux.io_uring_sqe, |
| 1348 | 1584 | fd: os.fd_t, |
| ... | ... | @@ -1373,6 +1609,28 @@ pub fn io_uring_prep_send(sqe: *linux.io_uring_sqe, fd: os.fd_t, buffer: []const |
| 1373 | 1609 | sqe.rw_flags = flags; |
| 1374 | 1610 | } |
| 1375 | 1611 | |
| 1612 | pub fn io_uring_prep_send_zc(sqe: *linux.io_uring_sqe, fd: os.fd_t, buffer: []const u8, flags: u32, zc_flags: u16) void { |
| 1613 | io_uring_prep_rw(.SEND_ZC, sqe, fd, @intFromPtr(buffer.ptr), buffer.len, 0); |
| 1614 | sqe.rw_flags = flags; |
| 1615 | sqe.ioprio = zc_flags; |
| 1616 | } |
| 1617 | |
| 1618 | pub fn io_uring_prep_send_zc_fixed(sqe: *linux.io_uring_sqe, fd: os.fd_t, buffer: []const u8, flags: u32, zc_flags: u16, buf_index: u16) void { |
| 1619 | io_uring_prep_send_zc(sqe, fd, buffer, flags, zc_flags); |
| 1620 | sqe.ioprio |= linux.IORING_RECVSEND_FIXED_BUF; |
| 1621 | sqe.buf_index = buf_index; |
| 1622 | } |
| 1623 | |
| 1624 | pub fn io_uring_prep_sendmsg_zc( |
| 1625 | sqe: *linux.io_uring_sqe, |
| 1626 | fd: os.fd_t, |
| 1627 | msg: *const os.msghdr_const, |
| 1628 | flags: u32, |
| 1629 | ) void { |
| 1630 | io_uring_prep_sendmsg(sqe, fd, msg, flags); |
| 1631 | sqe.opcode = .SENDMSG_ZC; |
| 1632 | } |
| 1633 | |
| 1376 | 1634 | pub fn io_uring_prep_recvmsg( |
| 1377 | 1635 | sqe: *linux.io_uring_sqe, |
| 1378 | 1636 | fd: os.fd_t, |
| ... | ... | @@ -1404,6 +1662,18 @@ pub fn io_uring_prep_openat( |
| 1404 | 1662 | sqe.rw_flags = flags; |
| 1405 | 1663 | } |
| 1406 | 1664 | |
| 1665 | pub fn io_uring_prep_openat_direct( |
| 1666 | sqe: *linux.io_uring_sqe, |
| 1667 | fd: os.fd_t, |
| 1668 | path: [*:0]const u8, |
| 1669 | flags: u32, |
| 1670 | mode: os.mode_t, |
| 1671 | file_index: u32, |
| 1672 | ) void { |
| 1673 | io_uring_prep_openat(sqe, fd, path, flags, mode); |
| 1674 | __io_uring_set_target_fixed_file(sqe, file_index); |
| 1675 | } |
| 1676 | |
| 1407 | 1677 | pub fn io_uring_prep_close(sqe: *linux.io_uring_sqe, fd: os.fd_t) void { |
| 1408 | 1678 | sqe.* = .{ |
| 1409 | 1679 | .opcode = .CLOSE, |
| ... | ... | @@ -1423,6 +1693,11 @@ pub fn io_uring_prep_close(sqe: *linux.io_uring_sqe, fd: os.fd_t) void { |
| 1423 | 1693 | }; |
| 1424 | 1694 | } |
| 1425 | 1695 | |
| 1696 | pub fn io_uring_prep_close_direct(sqe: *linux.io_uring_sqe, file_index: u32) void { |
| 1697 | io_uring_prep_close(sqe, 0); |
| 1698 | __io_uring_set_target_fixed_file(sqe, file_index); |
| 1699 | } |
| 1700 | |
| 1426 | 1701 | pub fn io_uring_prep_timeout( |
| 1427 | 1702 | sqe: *linux.io_uring_sqe, |
| 1428 | 1703 | ts: *const os.linux.kernel_timespec, |
| ... | ... | @@ -1650,6 +1925,40 @@ pub fn io_uring_prep_multishot_accept( |
| 1650 | 1925 | sqe.ioprio |= linux.IORING_ACCEPT_MULTISHOT; |
| 1651 | 1926 | } |
| 1652 | 1927 | |
| 1928 | pub fn io_uring_prep_socket( |
| 1929 | sqe: *linux.io_uring_sqe, |
| 1930 | domain: u32, |
| 1931 | socket_type: u32, |
| 1932 | protocol: u32, |
| 1933 | flags: u32, |
| 1934 | ) void { |
| 1935 | io_uring_prep_rw(.SOCKET, sqe, @intCast(domain), 0, protocol, socket_type); |
| 1936 | sqe.rw_flags = flags; |
| 1937 | } |
| 1938 | |
| 1939 | pub fn io_uring_prep_socket_direct( |
| 1940 | sqe: *linux.io_uring_sqe, |
| 1941 | domain: u32, |
| 1942 | socket_type: u32, |
| 1943 | protocol: u32, |
| 1944 | flags: u32, |
| 1945 | file_index: u32, |
| 1946 | ) void { |
| 1947 | io_uring_prep_socket(sqe, domain, socket_type, protocol, flags); |
| 1948 | __io_uring_set_target_fixed_file(sqe, file_index); |
| 1949 | } |
| 1950 | |
| 1951 | pub fn io_uring_prep_socket_direct_alloc( |
| 1952 | sqe: *linux.io_uring_sqe, |
| 1953 | domain: u32, |
| 1954 | socket_type: u32, |
| 1955 | protocol: u32, |
| 1956 | flags: u32, |
| 1957 | ) void { |
| 1958 | io_uring_prep_socket(sqe, domain, socket_type, protocol, flags); |
| 1959 | __io_uring_set_target_fixed_file(sqe, linux.IORING_FILE_INDEX_ALLOC); |
| 1960 | } |
| 1961 | |
| 1653 | 1962 | test "structs/offsets/entries" { |
| 1654 | 1963 | if (builtin.os.tag != .linux) return error.SkipZigTest; |
| 1655 | 1964 | |
| ... | ... | @@ -3492,3 +3801,363 @@ test "accept multishot" { |
| 3492 | 3801 | os.closeSocket(client); |
| 3493 | 3802 | } |
| 3494 | 3803 | } |
| 3804 | |
| 3805 | test "accept/connect/send_zc/recv" { |
| 3806 | try skipKernelLessThan(.{ .major = 6, .minor = 0, .patch = 0 }); |
| 3807 | |
| 3808 | var ring = IO_Uring.init(16, 0) catch |err| switch (err) { |
| 3809 | error.SystemOutdated => return error.SkipZigTest, |
| 3810 | error.PermissionDenied => return error.SkipZigTest, |
| 3811 | else => return err, |
| 3812 | }; |
| 3813 | defer ring.deinit(); |
| 3814 | |
| 3815 | const socket_test_harness = try createSocketTestHarness(&ring); |
| 3816 | defer socket_test_harness.close(); |
| 3817 | |
| 3818 | const buffer_send = [_]u8{ 0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 0xa, 0xb, 0xc, 0xd, 0xe }; |
| 3819 | var buffer_recv = [_]u8{0} ** 10; |
| 3820 | |
| 3821 | // zero-copy send |
| 3822 | const send = try ring.send_zc(0xeeeeeeee, socket_test_harness.client, buffer_send[0..], 0, 0); |
| 3823 | send.flags |= linux.IOSQE_IO_LINK; |
| 3824 | _ = try ring.recv(0xffffffff, socket_test_harness.server, .{ .buffer = buffer_recv[0..] }, 0); |
| 3825 | try testing.expectEqual(@as(u32, 2), try ring.submit()); |
| 3826 | |
| 3827 | // First completion of zero-copy send. |
| 3828 | // IORING_CQE_F_MORE, means that there |
| 3829 | // will be a second completion event / notification for the |
| 3830 | // request, with the user_data field set to the same value. |
| 3831 | // buffer_send must be keep alive until second cqe. |
| 3832 | var cqe_send = try ring.copy_cqe(); |
| 3833 | try testing.expectEqual(linux.io_uring_cqe{ |
| 3834 | .user_data = 0xeeeeeeee, |
| 3835 | .res = buffer_send.len, |
| 3836 | .flags = linux.IORING_CQE_F_MORE, |
| 3837 | }, cqe_send); |
| 3838 | |
| 3839 | const cqe_recv = try ring.copy_cqe(); |
| 3840 | try testing.expectEqual(linux.io_uring_cqe{ |
| 3841 | .user_data = 0xffffffff, |
| 3842 | .res = buffer_recv.len, |
| 3843 | .flags = cqe_recv.flags & linux.IORING_CQE_F_SOCK_NONEMPTY, |
| 3844 | }, cqe_recv); |
| 3845 | |
| 3846 | try testing.expectEqualSlices(u8, buffer_send[0..buffer_recv.len], buffer_recv[0..]); |
| 3847 | |
| 3848 | // Second completion of zero-copy send. |
| 3849 | // IORING_CQE_F_NOTIF in flags signals that kernel is done with send_buffer |
| 3850 | cqe_send = try ring.copy_cqe(); |
| 3851 | try testing.expectEqual(linux.io_uring_cqe{ |
| 3852 | .user_data = 0xeeeeeeee, |
| 3853 | .res = 0, |
| 3854 | .flags = linux.IORING_CQE_F_NOTIF, |
| 3855 | }, cqe_send); |
| 3856 | } |
| 3857 | |
| 3858 | test "accept_direct" { |
| 3859 | try skipKernelLessThan(.{ .major = 5, .minor = 19, .patch = 0 }); |
| 3860 | |
| 3861 | var ring = IO_Uring.init(1, 0) catch |err| switch (err) { |
| 3862 | error.SystemOutdated => return error.SkipZigTest, |
| 3863 | error.PermissionDenied => return error.SkipZigTest, |
| 3864 | else => return err, |
| 3865 | }; |
| 3866 | defer ring.deinit(); |
| 3867 | var address = try net.Address.parseIp4("127.0.0.1", 0); |
| 3868 | |
| 3869 | // register direct file descriptors |
| 3870 | var registered_fds = [_]os.fd_t{-1} ** 2; |
| 3871 | try ring.register_files(registered_fds[0..]); |
| 3872 | |
| 3873 | const listener_socket = try createListenerSocket(&address); |
| 3874 | defer os.closeSocket(listener_socket); |
| 3875 | |
| 3876 | const accept_userdata: u64 = 0xaaaaaaaa; |
| 3877 | const read_userdata: u64 = 0xbbbbbbbb; |
| 3878 | const data = [_]u8{ 0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 0xa, 0xb, 0xc, 0xd, 0xe }; |
| 3879 | |
| 3880 | for (0..2) |_| { |
| 3881 | for (registered_fds, 0..) |_, i| { |
| 3882 | var buffer_recv = [_]u8{0} ** 16; |
| 3883 | const buffer_send: []const u8 = data[0 .. data.len - i]; // make it different at each loop |
| 3884 | |
| 3885 | // submit accept, will chose registered fd and return index in cqe |
| 3886 | _ = try ring.accept_direct(accept_userdata, listener_socket, null, null, 0); |
| 3887 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| 3888 | |
| 3889 | // connect |
| 3890 | var client = try os.socket(address.any.family, os.SOCK.STREAM | os.SOCK.CLOEXEC, 0); |
| 3891 | try os.connect(client, &address.any, address.getOsSockLen()); |
| 3892 | defer os.closeSocket(client); |
| 3893 | |
| 3894 | // accept completion |
| 3895 | const cqe_accept = try ring.copy_cqe(); |
| 3896 | try testing.expectEqual(os.E.SUCCESS, cqe_accept.err()); |
| 3897 | const fd_index = cqe_accept.res; |
| 3898 | try testing.expect(fd_index < registered_fds.len); |
| 3899 | try testing.expect(cqe_accept.user_data == accept_userdata); |
| 3900 | |
| 3901 | // send data |
| 3902 | _ = try os.send(client, buffer_send, 0); |
| 3903 | |
| 3904 | // Example of how to use registered fd: |
| 3905 | // Submit receive to fixed file returned by accept (fd_index). |
| 3906 | // Fd field is set to registered file index, returned by accept. |
| 3907 | // Flag linux.IOSQE_FIXED_FILE must be set. |
| 3908 | const recv_sqe = try ring.recv(read_userdata, fd_index, .{ .buffer = &buffer_recv }, 0); |
| 3909 | recv_sqe.flags |= linux.IOSQE_FIXED_FILE; |
| 3910 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| 3911 | |
| 3912 | // accept receive |
| 3913 | const recv_cqe = try ring.copy_cqe(); |
| 3914 | try testing.expect(recv_cqe.user_data == read_userdata); |
| 3915 | try testing.expect(recv_cqe.res == buffer_send.len); |
| 3916 | try testing.expectEqualSlices(u8, buffer_send, buffer_recv[0..buffer_send.len]); |
| 3917 | } |
| 3918 | // no more available fds, accept will get NFILE error |
| 3919 | { |
| 3920 | // submit accept |
| 3921 | _ = try ring.accept_direct(accept_userdata, listener_socket, null, null, 0); |
| 3922 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| 3923 | // connect |
| 3924 | var client = try os.socket(address.any.family, os.SOCK.STREAM | os.SOCK.CLOEXEC, 0); |
| 3925 | try os.connect(client, &address.any, address.getOsSockLen()); |
| 3926 | defer os.closeSocket(client); |
| 3927 | // completion with error |
| 3928 | const cqe_accept = try ring.copy_cqe(); |
| 3929 | try testing.expect(cqe_accept.user_data == accept_userdata); |
| 3930 | try testing.expectEqual(os.E.NFILE, cqe_accept.err()); |
| 3931 | } |
| 3932 | // return file descriptors to kernel |
| 3933 | try ring.register_files_update(0, registered_fds[0..]); |
| 3934 | } |
| 3935 | try ring.unregister_files(); |
| 3936 | } |
| 3937 | |
| 3938 | test "accept_multishot_direct" { |
| 3939 | try skipKernelLessThan(.{ .major = 5, .minor = 19, .patch = 0 }); |
| 3940 | |
| 3941 | var ring = IO_Uring.init(1, 0) catch |err| switch (err) { |
| 3942 | error.SystemOutdated => return error.SkipZigTest, |
| 3943 | error.PermissionDenied => return error.SkipZigTest, |
| 3944 | else => return err, |
| 3945 | }; |
| 3946 | defer ring.deinit(); |
| 3947 | |
| 3948 | var address = try net.Address.parseIp4("127.0.0.1", 0); |
| 3949 | |
| 3950 | var registered_fds = [_]os.fd_t{-1} ** 2; |
| 3951 | try ring.register_files(registered_fds[0..]); |
| 3952 | |
| 3953 | const listener_socket = try createListenerSocket(&address); |
| 3954 | defer os.closeSocket(listener_socket); |
| 3955 | |
| 3956 | const accept_userdata: u64 = 0xaaaaaaaa; |
| 3957 | |
| 3958 | for (0..2) |_| { |
| 3959 | // submit multishot accept |
| 3960 | // Will chose registered fd and return index of the selected registered file in cqe. |
| 3961 | _ = try ring.accept_multishot_direct(accept_userdata, listener_socket, null, null, 0); |
| 3962 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| 3963 | |
| 3964 | for (registered_fds) |_| { |
| 3965 | // connect |
| 3966 | var client = try os.socket(address.any.family, os.SOCK.STREAM | os.SOCK.CLOEXEC, 0); |
| 3967 | try os.connect(client, &address.any, address.getOsSockLen()); |
| 3968 | defer os.closeSocket(client); |
| 3969 | |
| 3970 | // accept completion |
| 3971 | const cqe_accept = try ring.copy_cqe(); |
| 3972 | const fd_index = cqe_accept.res; |
| 3973 | try testing.expect(fd_index < registered_fds.len); |
| 3974 | try testing.expect(cqe_accept.user_data == accept_userdata); |
| 3975 | try testing.expect(cqe_accept.flags & linux.IORING_CQE_F_MORE > 0); // has more is set |
| 3976 | } |
| 3977 | // No more available fds, accept will get NFILE error. |
| 3978 | // Multishot is terminated (more flag is not set). |
| 3979 | { |
| 3980 | // connect |
| 3981 | var client = try os.socket(address.any.family, os.SOCK.STREAM | os.SOCK.CLOEXEC, 0); |
| 3982 | try os.connect(client, &address.any, address.getOsSockLen()); |
| 3983 | defer os.closeSocket(client); |
| 3984 | // completion with error |
| 3985 | const cqe_accept = try ring.copy_cqe(); |
| 3986 | try testing.expect(cqe_accept.user_data == accept_userdata); |
| 3987 | try testing.expectEqual(os.E.NFILE, cqe_accept.err()); |
| 3988 | try testing.expect(cqe_accept.flags & linux.IORING_CQE_F_MORE == 0); // has more is not set |
| 3989 | } |
| 3990 | // return file descriptors to kernel |
| 3991 | try ring.register_files_update(0, registered_fds[0..]); |
| 3992 | } |
| 3993 | try ring.unregister_files(); |
| 3994 | } |
| 3995 | |
| 3996 | test "socket" { |
| 3997 | try skipKernelLessThan(.{ .major = 5, .minor = 19, .patch = 0 }); |
| 3998 | |
| 3999 | var ring = IO_Uring.init(1, 0) catch |err| switch (err) { |
| 4000 | error.SystemOutdated => return error.SkipZigTest, |
| 4001 | error.PermissionDenied => return error.SkipZigTest, |
| 4002 | else => return err, |
| 4003 | }; |
| 4004 | defer ring.deinit(); |
| 4005 | |
| 4006 | // prepare, submit socket operation |
| 4007 | _ = try ring.socket(0, linux.AF.INET, os.SOCK.STREAM, 0, 0); |
| 4008 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| 4009 | |
| 4010 | // test completion |
| 4011 | var cqe = try ring.copy_cqe(); |
| 4012 | try testing.expectEqual(os.E.SUCCESS, cqe.err()); |
| 4013 | const fd: os.fd_t = @intCast(cqe.res); |
| 4014 | try testing.expect(fd > 2); |
| 4015 | |
| 4016 | os.close(fd); |
| 4017 | } |
| 4018 | |
| 4019 | test "socket_direct/socket_direct_alloc/close_direct" { |
| 4020 | try skipKernelLessThan(.{ .major = 5, .minor = 19, .patch = 0 }); |
| 4021 | |
| 4022 | var ring = IO_Uring.init(2, 0) catch |err| switch (err) { |
| 4023 | error.SystemOutdated => return error.SkipZigTest, |
| 4024 | error.PermissionDenied => return error.SkipZigTest, |
| 4025 | else => return err, |
| 4026 | }; |
| 4027 | defer ring.deinit(); |
| 4028 | |
| 4029 | var registered_fds = [_]os.fd_t{-1} ** 3; |
| 4030 | try ring.register_files(registered_fds[0..]); |
| 4031 | |
| 4032 | // create socket in registered file descriptor at index 0 (last param) |
| 4033 | _ = try ring.socket_direct(0, linux.AF.INET, os.SOCK.STREAM, 0, 0, 0); |
| 4034 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| 4035 | var cqe_socket = try ring.copy_cqe(); |
| 4036 | try testing.expectEqual(os.E.SUCCESS, cqe_socket.err()); |
| 4037 | try testing.expect(cqe_socket.res == 0); |
| 4038 | |
| 4039 | // create socket in registered file descriptor at index 1 (last param) |
| 4040 | _ = try ring.socket_direct(0, linux.AF.INET, os.SOCK.STREAM, 0, 0, 1); |
| 4041 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| 4042 | cqe_socket = try ring.copy_cqe(); |
| 4043 | try testing.expectEqual(os.E.SUCCESS, cqe_socket.err()); |
| 4044 | try testing.expect(cqe_socket.res == 0); // res is 0 when index is specified |
| 4045 | |
| 4046 | // create socket in kernel chosen file descriptor index (_alloc version) |
| 4047 | // completion res has index from registered files |
| 4048 | _ = try ring.socket_direct_alloc(0, linux.AF.INET, os.SOCK.STREAM, 0, 0); |
| 4049 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| 4050 | cqe_socket = try ring.copy_cqe(); |
| 4051 | try testing.expectEqual(os.E.SUCCESS, cqe_socket.err()); |
| 4052 | try testing.expect(cqe_socket.res == 2); // returns registered file index |
| 4053 | |
| 4054 | // use sockets from registered_fds in connect operation |
| 4055 | var address = try net.Address.parseIp4("127.0.0.1", 0); |
| 4056 | const listener_socket = try createListenerSocket(&address); |
| 4057 | defer os.closeSocket(listener_socket); |
| 4058 | const accept_userdata: u64 = 0xaaaaaaaa; |
| 4059 | const connect_userdata: u64 = 0xbbbbbbbb; |
| 4060 | const close_userdata: u64 = 0xcccccccc; |
| 4061 | for (registered_fds, 0..) |_, fd_index| { |
| 4062 | // prepare accept |
| 4063 | _ = try ring.accept(accept_userdata, listener_socket, null, null, 0); |
| 4064 | // prepare connect with fixed socket |
| 4065 | const connect_sqe = try ring.connect(connect_userdata, @intCast(fd_index), &address.any, address.getOsSockLen()); |
| 4066 | connect_sqe.flags |= linux.IOSQE_FIXED_FILE; // fd is fixed file index |
| 4067 | // submit both |
| 4068 | try testing.expectEqual(@as(u32, 2), try ring.submit()); |
| 4069 | // get completions |
| 4070 | var cqe_connect = try ring.copy_cqe(); |
| 4071 | var cqe_accept = try ring.copy_cqe(); |
| 4072 | // ignore order |
| 4073 | if (cqe_connect.user_data == accept_userdata and cqe_accept.user_data == connect_userdata) { |
| 4074 | const a = cqe_accept; |
| 4075 | const b = cqe_connect; |
| 4076 | cqe_accept = b; |
| 4077 | cqe_connect = a; |
| 4078 | } |
| 4079 | // test connect completion |
| 4080 | try testing.expect(cqe_connect.user_data == connect_userdata); |
| 4081 | try testing.expectEqual(os.E.SUCCESS, cqe_connect.err()); |
| 4082 | // test accept completion |
| 4083 | try testing.expect(cqe_accept.user_data == accept_userdata); |
| 4084 | try testing.expectEqual(os.E.SUCCESS, cqe_accept.err()); |
| 4085 | |
| 4086 | // submit and test close_direct |
| 4087 | _ = try ring.close_direct(close_userdata, @intCast(fd_index)); |
| 4088 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| 4089 | var cqe_close = try ring.copy_cqe(); |
| 4090 | try testing.expect(cqe_close.user_data == close_userdata); |
| 4091 | try testing.expectEqual(os.E.SUCCESS, cqe_close.err()); |
| 4092 | } |
| 4093 | |
| 4094 | try ring.unregister_files(); |
| 4095 | } |
| 4096 | |
| 4097 | test "openat_direct/close_direct" { |
| 4098 | try skipKernelLessThan(.{ .major = 5, .minor = 19, .patch = 0 }); |
| 4099 | |
| 4100 | var ring = IO_Uring.init(2, 0) catch |err| switch (err) { |
| 4101 | error.SystemOutdated => return error.SkipZigTest, |
| 4102 | error.PermissionDenied => return error.SkipZigTest, |
| 4103 | else => return err, |
| 4104 | }; |
| 4105 | defer ring.deinit(); |
| 4106 | |
| 4107 | var registered_fds = [_]os.fd_t{-1} ** 3; |
| 4108 | try ring.register_files(registered_fds[0..]); |
| 4109 | |
| 4110 | var tmp = std.testing.tmpDir(.{}); |
| 4111 | defer tmp.cleanup(); |
| 4112 | const path = "test_io_uring_close_direct"; |
| 4113 | const flags: u32 = os.O.RDWR | os.O.CREAT; |
| 4114 | const mode: os.mode_t = 0o666; |
| 4115 | const user_data: u64 = 0; |
| 4116 | |
| 4117 | // use registered file at index 0 (last param) |
| 4118 | _ = try ring.openat_direct(user_data, tmp.dir.fd, path, flags, mode, 0); |
| 4119 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| 4120 | var cqe = try ring.copy_cqe(); |
| 4121 | try testing.expectEqual(os.E.SUCCESS, cqe.err()); |
| 4122 | try testing.expect(cqe.res == 0); |
| 4123 | |
| 4124 | // use registered file at index 1 |
| 4125 | _ = try ring.openat_direct(user_data, tmp.dir.fd, path, flags, mode, 1); |
| 4126 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| 4127 | cqe = try ring.copy_cqe(); |
| 4128 | try testing.expectEqual(os.E.SUCCESS, cqe.err()); |
| 4129 | try testing.expect(cqe.res == 0); // res is 0 when we specify index |
| 4130 | |
| 4131 | // let kernel choose registered file index |
| 4132 | _ = try ring.openat_direct(user_data, tmp.dir.fd, path, flags, mode, linux.IORING_FILE_INDEX_ALLOC); |
| 4133 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| 4134 | cqe = try ring.copy_cqe(); |
| 4135 | try testing.expectEqual(os.E.SUCCESS, cqe.err()); |
| 4136 | try testing.expect(cqe.res == 2); // chosen index is in res |
| 4137 | |
| 4138 | // close all open file descriptors |
| 4139 | for (registered_fds, 0..) |_, fd_index| { |
| 4140 | _ = try ring.close_direct(user_data, @intCast(fd_index)); |
| 4141 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| 4142 | var cqe_close = try ring.copy_cqe(); |
| 4143 | try testing.expectEqual(os.E.SUCCESS, cqe_close.err()); |
| 4144 | } |
| 4145 | try ring.unregister_files(); |
| 4146 | } |
| 4147 | |
| 4148 | /// For use in tests. Returns SkipZigTest is kernel version is less than required. |
| 4149 | inline fn skipKernelLessThan(required: std.SemanticVersion) !void { |
| 4150 | if (builtin.os.tag != .linux) return error.SkipZigTest; |
| 4151 | |
| 4152 | var uts: linux.utsname = undefined; |
| 4153 | const res = linux.uname(&uts); |
| 4154 | switch (linux.getErrno(res)) { |
| 4155 | .SUCCESS => {}, |
| 4156 | else => |errno| return os.unexpectedErrno(errno), |
| 4157 | } |
| 4158 | |
| 4159 | const release = mem.sliceTo(&uts.release, 0); |
| 4160 | var current = try std.SemanticVersion.parse(release); |
| 4161 | current.pre = null; // don't check pre field |
| 4162 | if (required.order(current) == .gt) return error.SkipZigTest; |
| 4163 | } |