| ... | ... | @@ -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, |
| ... | ... | @@ -511,6 +512,8 @@ pub const IO_Uring = struct { |
| 511 | 512 | /// which will repeatedly trigger a CQE when a connection request comes in. |
| 512 | 513 | /// While IORING_CQE_F_MORE flag is set in CQE flags accept will generate |
| 513 | 514 | /// further CQEs. |
| 515 | /// |
| 516 | /// Available since 5.19 |
| 514 | 517 | pub fn accept_multishot( |
| 515 | 518 | self: *IO_Uring, |
| 516 | 519 | user_data: u64, |
| ... | ... | @@ -534,6 +537,8 @@ pub const IO_Uring = struct { |
| 534 | 537 | /// After creation, they can be used by setting IOSQE_FIXED_FILE in the SQE |
| 535 | 538 | /// flags member, and setting the SQE fd field to the direct descriptor value |
| 536 | 539 | /// rather than the regular file descriptor. |
| 540 | /// |
| 541 | /// Available since 5.19 |
| 537 | 542 | pub fn accept_direct( |
| 538 | 543 | self: *IO_Uring, |
| 539 | 544 | user_data: u64, |
| ... | ... | @@ -549,6 +554,7 @@ pub const IO_Uring = struct { |
| 549 | 554 | } |
| 550 | 555 | |
| 551 | 556 | /// Queues an multishot accept using direct (registered) file descriptors. |
| 557 | /// Available since 5.19 |
| 552 | 558 | pub fn accept_multishot_direct( |
| 553 | 559 | self: *IO_Uring, |
| 554 | 560 | user_data: u64, |
| ... | ... | @@ -726,6 +732,7 @@ pub const IO_Uring = struct { |
| 726 | 732 | |
| 727 | 733 | /// Queues (but does not submit) an SQE to perform an `openat(2)`. |
| 728 | 734 | /// Returns a pointer to the SQE. |
| 735 | /// Available since 5.6. |
| 729 | 736 | pub fn openat( |
| 730 | 737 | self: *IO_Uring, |
| 731 | 738 | user_data: u64, |
| ... | ... | @@ -749,6 +756,8 @@ pub const IO_Uring = struct { |
| 749 | 756 | /// After creation, they can be used by setting IOSQE_FIXED_FILE in the SQE |
| 750 | 757 | /// flags member, and setting the SQE fd field to the direct descriptor value |
| 751 | 758 | /// rather than the regular file descriptor. |
| 759 | /// |
| 760 | /// Available since 5.15 |
| 752 | 761 | pub fn openat_direct( |
| 753 | 762 | self: *IO_Uring, |
| 754 | 763 | user_data: u64, |
| ... | ... | @@ -766,6 +775,7 @@ pub const IO_Uring = struct { |
| 766 | 775 | |
| 767 | 776 | /// Queues (but does not submit) an SQE to perform a `close(2)`. |
| 768 | 777 | /// Returns a pointer to the SQE. |
| 778 | /// Available since 5.6. |
| 769 | 779 | pub fn close(self: *IO_Uring, user_data: u64, fd: os.fd_t) !*linux.io_uring_sqe { |
| 770 | 780 | const sqe = try self.get_sqe(); |
| 771 | 781 | io_uring_prep_close(sqe, fd); |
| ... | ... | @@ -774,6 +784,7 @@ pub const IO_Uring = struct { |
| 774 | 784 | } |
| 775 | 785 | |
| 776 | 786 | /// Queues close of registered file descriptor. |
| 787 | /// Available since 5.15 |
| 777 | 788 | pub fn close_direct(self: *IO_Uring, user_data: u64, file_index: u32) !*linux.io_uring_sqe { |
| 778 | 789 | const sqe = try self.get_sqe(); |
| 779 | 790 | io_uring_prep_close_direct(sqe, file_index); |
| ... | ... | @@ -1232,6 +1243,7 @@ pub const IO_Uring = struct { |
| 1232 | 1243 | |
| 1233 | 1244 | /// Prepares a socket creation request. |
| 1234 | 1245 | /// New socket fd will be returned in completion result. |
| 1246 | /// Available since 5.19 |
| 1235 | 1247 | pub fn socket( |
| 1236 | 1248 | self: *IO_Uring, |
| 1237 | 1249 | user_data: u64, |
| ... | ... | @@ -1247,6 +1259,7 @@ pub const IO_Uring = struct { |
| 1247 | 1259 | } |
| 1248 | 1260 | |
| 1249 | 1261 | /// Prepares a socket creation request for registered file at index `file_index`. |
| 1262 | /// Available since 5.19 |
| 1250 | 1263 | pub fn socket_direct( |
| 1251 | 1264 | self: *IO_Uring, |
| 1252 | 1265 | user_data: u64, |
| ... | ... | @@ -1264,6 +1277,7 @@ pub const IO_Uring = struct { |
| 1264 | 1277 | |
| 1265 | 1278 | /// Prepares a socket creation request for registered file, index chosen by kernel (file index alloc). |
| 1266 | 1279 | /// File index will be returned in CQE res field. |
| 1280 | /// Available since 5.19 |
| 1267 | 1281 | pub fn socket_direct_alloc( |
| 1268 | 1282 | self: *IO_Uring, |
| 1269 | 1283 | user_data: u64, |
| ... | ... | @@ -3826,22 +3840,15 @@ test "accept/connect/send_zc/recv" { |
| 3826 | 3840 | } |
| 3827 | 3841 | |
| 3828 | 3842 | test "accept_direct" { |
| 3829 | | if (builtin.os.tag != .linux) return error.SkipZigTest; |
| 3843 | try skipKernelLessThan(.{ .major = 5, .minor = 19, .patch = 0 }); |
| 3830 | 3844 | |
| 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 | | }; |
| 3845 | var ring = try IO_Uring.init(1, 0); |
| 3836 | 3846 | defer ring.deinit(); |
| 3837 | 3847 | var address = try net.Address.parseIp4("127.0.0.1", 0); |
| 3838 | 3848 | |
| 3839 | 3849 | // register direct file descriptors |
| 3840 | 3850 | 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 | | }; |
| 3851 | try ring.register_files(registered_fds[0..]); |
| 3845 | 3852 | |
| 3846 | 3853 | const listener_socket = try createListenerSocket(&address); |
| 3847 | 3854 | defer os.closeSocket(listener_socket); |
| ... | ... | @@ -3866,10 +3873,8 @@ test "accept_direct" { |
| 3866 | 3873 | |
| 3867 | 3874 | // accept completion |
| 3868 | 3875 | const cqe_accept = try ring.copy_cqe(); |
| 3869 | | if (cqe_accept.err() == .INVAL) return error.SkipZigTest; |
| 3870 | 3876 | try testing.expectEqual(os.E.SUCCESS, cqe_accept.err()); |
| 3871 | 3877 | const fd_index = cqe_accept.res; |
| 3872 | | if (fd_index >= registered_fds.len) return error.SkipZigTest; // old kernel fallback to ordinary accept |
| 3873 | 3878 | try testing.expect(fd_index < registered_fds.len); |
| 3874 | 3879 | try testing.expect(cqe_accept.user_data == accept_userdata); |
| 3875 | 3880 | |
| ... | ... | @@ -3911,21 +3916,15 @@ test "accept_direct" { |
| 3911 | 3916 | } |
| 3912 | 3917 | |
| 3913 | 3918 | test "accept_multishot_direct" { |
| 3914 | | if (builtin.os.tag != .linux) return error.SkipZigTest; |
| 3919 | try skipKernelLessThan(.{ .major = 5, .minor = 19, .patch = 0 }); |
| 3915 | 3920 | |
| 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 | var ring = try IO_Uring.init(1, 0); |
| 3921 | 3922 | defer ring.deinit(); |
| 3923 | |
| 3922 | 3924 | var address = try net.Address.parseIp4("127.0.0.1", 0); |
| 3923 | 3925 | |
| 3924 | 3926 | 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 | | }; |
| 3927 | try ring.register_files(registered_fds[0..]); |
| 3929 | 3928 | |
| 3930 | 3929 | const listener_socket = try createListenerSocket(&address); |
| 3931 | 3930 | defer os.closeSocket(listener_socket); |
| ... | ... | @@ -3933,7 +3932,7 @@ test "accept_multishot_direct" { |
| 3933 | 3932 | const accept_userdata: u64 = 0xaaaaaaaa; |
| 3934 | 3933 | |
| 3935 | 3934 | for (0..2) |_| { |
| 3936 | | // submit accept |
| 3935 | // submit multishot accept |
| 3937 | 3936 | // Will chose registered fd and return index of the selected registered file in cqe. |
| 3938 | 3937 | _ = try ring.accept_multishot_direct(accept_userdata, listener_socket, null, null, 0); |
| 3939 | 3938 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| ... | ... | @@ -3946,7 +3945,6 @@ test "accept_multishot_direct" { |
| 3946 | 3945 | |
| 3947 | 3946 | // accept completion |
| 3948 | 3947 | const cqe_accept = try ring.copy_cqe(); |
| 3949 | | if (cqe_accept.err() == .INVAL) return error.SkipZigTest; |
| 3950 | 3948 | const fd_index = cqe_accept.res; |
| 3951 | 3949 | try testing.expect(fd_index < registered_fds.len); |
| 3952 | 3950 | try testing.expect(cqe_accept.user_data == accept_userdata); |
| ... | ... | @@ -3971,52 +3969,58 @@ test "accept_multishot_direct" { |
| 3971 | 3969 | try ring.unregister_files(); |
| 3972 | 3970 | } |
| 3973 | 3971 | |
| 3974 | | test "socket/socket_direct/socket_direct_alloc/close_direct" { |
| 3975 | | if (builtin.os.tag != .linux) return error.SkipZigTest; |
| 3972 | test "socket" { |
| 3973 | try skipKernelLessThan(.{ .major = 5, .minor = 19, .patch = 0 }); |
| 3976 | 3974 | |
| 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 | | }; |
| 3975 | var ring = try IO_Uring.init(2, 0); |
| 3982 | 3976 | defer ring.deinit(); |
| 3983 | | var address = try net.Address.parseIp4("127.0.0.1", 0); |
| 3984 | 3977 | |
| 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); |
| 3978 | // prepare, submit socket operation |
| 3979 | _ = try ring.socket(0, linux.AF.INET, os.SOCK.STREAM, 0, 0); |
| 3980 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| 3981 | |
| 3982 | // test completion |
| 3983 | var cqe = try ring.copy_cqe(); |
| 3984 | try testing.expectEqual(os.E.SUCCESS, cqe.err()); |
| 3985 | const fd: os.fd_t = @intCast(cqe.res); |
| 3986 | try testing.expect(fd > 2); |
| 3987 | |
| 3988 | os.close(fd); |
| 3989 | } |
| 3990 | |
| 3991 | test "socket_direct/socket_direct_alloc/close_direct" { |
| 3992 | try skipKernelLessThan(.{ .major = 5, .minor = 19, .patch = 0 }); |
| 3993 | |
| 3994 | var ring = try IO_Uring.init(2, 0); |
| 3995 | defer ring.deinit(); |
| 3996 | |
| 3997 | var registered_fds = [_]os.fd_t{-1} ** 3; |
| 3998 | try ring.register_files(registered_fds[0..]); |
| 3999 | |
| 4000 | // create socket in registered file descriptor at index 0 (last param) |
| 4001 | _ = try ring.socket_direct(0, linux.AF.INET, os.SOCK.STREAM, 0, 0, 0); |
| 3993 | 4002 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| 3994 | 4003 | var cqe_socket = try ring.copy_cqe(); |
| 3995 | | if (cqe_socket.err() == .INVAL) return error.SkipZigTest; |
| 3996 | 4004 | 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 | | }; |
| 4005 | try testing.expect(cqe_socket.res == 0); |
| 4004 | 4006 | |
| 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 | // create socket in registered file descriptor at index 1 (last param) |
| 4008 | _ = try ring.socket_direct(0, linux.AF.INET, os.SOCK.STREAM, 0, 0, 1); |
| 4007 | 4009 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| 4008 | 4010 | cqe_socket = try ring.copy_cqe(); |
| 4009 | 4011 | try testing.expectEqual(os.E.SUCCESS, cqe_socket.err()); |
| 4010 | | try testing.expect(cqe_socket.res == 0); |
| 4012 | try testing.expect(cqe_socket.res == 0); // res is 0 when index is specified |
| 4011 | 4013 | |
| 4012 | | // 4. io_uring socket_direct_alloc |
| 4013 | | _ = try ring.socket_direct_alloc(socket_userdata, linux.AF.INET, os.SOCK.STREAM, 0, 0); |
| 4014 | // create socket in kernel chosen file descriptor index (_alloc version) |
| 4015 | // completion res has index from registered files |
| 4016 | _ = try ring.socket_direct_alloc(0, linux.AF.INET, os.SOCK.STREAM, 0, 0); |
| 4014 | 4017 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| 4015 | 4018 | cqe_socket = try ring.copy_cqe(); |
| 4016 | 4019 | try testing.expectEqual(os.E.SUCCESS, cqe_socket.err()); |
| 4017 | | try testing.expect(cqe_socket.res == 3); |
| 4020 | try testing.expect(cqe_socket.res == 2); // returns registered file index |
| 4018 | 4021 | |
| 4019 | 4022 | // use sockets from registered_fds in connect operation |
| 4023 | var address = try net.Address.parseIp4("127.0.0.1", 0); |
| 4020 | 4024 | const listener_socket = try createListenerSocket(&address); |
| 4021 | 4025 | defer os.closeSocket(listener_socket); |
| 4022 | 4026 | const accept_userdata: u64 = 0xaaaaaaaa; |
| ... | ... | @@ -4027,14 +4031,12 @@ test "socket/socket_direct/socket_direct_alloc/close_direct" { |
| 4027 | 4031 | _ = try ring.accept(accept_userdata, listener_socket, null, null, 0); |
| 4028 | 4032 | // prepare connect with fixed socket |
| 4029 | 4033 | const connect_sqe = try ring.connect(connect_userdata, @intCast(fd_index), &address.any, address.getOsSockLen()); |
| 4030 | | connect_sqe.flags |= linux.IOSQE_FIXED_FILE; |
| 4034 | connect_sqe.flags |= linux.IOSQE_FIXED_FILE; // fd is fixed file index |
| 4031 | 4035 | // submit both |
| 4032 | 4036 | try testing.expectEqual(@as(u32, 2), try ring.submit()); |
| 4033 | 4037 | // get completions |
| 4034 | 4038 | var cqe_connect = try ring.copy_cqe(); |
| 4035 | | if (cqe_connect.err() == .INVAL) return error.SkipZigTest; |
| 4036 | 4039 | var cqe_accept = try ring.copy_cqe(); |
| 4037 | | if (cqe_accept.err() == .INVAL) return error.SkipZigTest; |
| 4038 | 4040 | // ignore order |
| 4039 | 4041 | if (cqe_connect.user_data == accept_userdata and cqe_accept.user_data == connect_userdata) { |
| 4040 | 4042 | const a = cqe_accept; |
| ... | ... | @@ -4049,7 +4051,7 @@ test "socket/socket_direct/socket_direct_alloc/close_direct" { |
| 4049 | 4051 | try testing.expect(cqe_accept.user_data == accept_userdata); |
| 4050 | 4052 | try testing.expectEqual(os.E.SUCCESS, cqe_accept.err()); |
| 4051 | 4053 | |
| 4052 | | // submit and test close completion |
| 4054 | // submit and test close_direct |
| 4053 | 4055 | _ = try ring.close_direct(close_userdata, @intCast(fd_index)); |
| 4054 | 4056 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| 4055 | 4057 | var cqe_close = try ring.copy_cqe(); |
| ... | ... | @@ -4061,20 +4063,13 @@ test "socket/socket_direct/socket_direct_alloc/close_direct" { |
| 4061 | 4063 | } |
| 4062 | 4064 | |
| 4063 | 4065 | test "openat_direct/close_direct" { |
| 4064 | | if (builtin.os.tag != .linux) return error.SkipZigTest; |
| 4066 | try skipKernelLessThan(.{ .major = 5, .minor = 19, .patch = 0 }); |
| 4065 | 4067 | |
| 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 | | }; |
| 4068 | var ring = try IO_Uring.init(2, 0); |
| 4071 | 4069 | defer ring.deinit(); |
| 4072 | 4070 | |
| 4073 | 4071 | 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 | | }; |
| 4072 | try ring.register_files(registered_fds[0..]); |
| 4078 | 4073 | |
| 4079 | 4074 | var tmp = std.testing.tmpDir(.{}); |
| 4080 | 4075 | defer tmp.cleanup(); |
| ... | ... | @@ -4087,8 +4082,6 @@ test "openat_direct/close_direct" { |
| 4087 | 4082 | _ = try ring.openat_direct(user_data, tmp.dir.fd, path, flags, mode, 0); |
| 4088 | 4083 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| 4089 | 4084 | 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 | 4085 | try testing.expectEqual(os.E.SUCCESS, cqe.err()); |
| 4093 | 4086 | try testing.expect(cqe.res == 0); |
| 4094 | 4087 | |
| ... | ... | @@ -4103,7 +4096,6 @@ test "openat_direct/close_direct" { |
| 4103 | 4096 | _ = try ring.openat_direct(user_data, tmp.dir.fd, path, flags, mode, linux.IORING_FILE_INDEX_ALLOC); |
| 4104 | 4097 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| 4105 | 4098 | cqe = try ring.copy_cqe(); |
| 4106 | | if (cqe.err() == .INVAL) return error.SkipZigTest; // kernel 5.15 bug |
| 4107 | 4099 | try testing.expectEqual(os.E.SUCCESS, cqe.err()); |
| 4108 | 4100 | try testing.expect(cqe.res == 2); // chosen index is in res |
| 4109 | 4101 | |
| ... | ... | @@ -4116,3 +4108,20 @@ test "openat_direct/close_direct" { |
| 4116 | 4108 | } |
| 4117 | 4109 | try ring.unregister_files(); |
| 4118 | 4110 | } |
| 4111 | |
| 4112 | /// For use in tests. Returns SkipZigTest is kernel version is less than required. |
| 4113 | fn skipKernelLessThan(required: std.SemanticVersion) !void { |
| 4114 | if (builtin.os.tag != .linux) return error.SkipZigTest; |
| 4115 | |
| 4116 | var uts: linux.utsname = undefined; |
| 4117 | const res = linux.uname(&uts); |
| 4118 | switch (linux.getErrno(res)) { |
| 4119 | .SUCCESS => {}, |
| 4120 | else => |errno| return os.unexpectedErrno(errno), |
| 4121 | } |
| 4122 | |
| 4123 | const release = mem.sliceTo(&uts.release, 0); |
| 4124 | var current = try std.SemanticVersion.parse(release); |
| 4125 | current.pre = null; // don't check pre field |
| 4126 | if (required.order(current) == .gt) return error.SkipZigTest; |
| 4127 | } |