authorgravatar for igor.anic@gmail.comIgor Anić <igor.anic@gmail.com> 2025-02-27 23:55:43+01:00
committergravatar for igor.anic@gmail.comIgor Anić <igor.anic@gmail.com> 2025-03-05 13:35:52+01:00
log2da8eff9d6b7f9d784a596836188cbede6cfb1d0
treed46dd9d1d1b38f9e918aaf5cb73a6b2786c6e7cb
parent79460d4a3eef8eb927b02a7eda8bc9999a766672

io_uring: add bind and listen


3 files changed, 243 insertions(+), 0 deletions(-)

lib/std/os/linux.zig+10
......@@ -5806,6 +5806,9 @@ pub const IORING_OP = enum(u8) {
58065806 FUTEX_WAITV,
58075807 FIXED_FD_INSTALL,
58085808 FTRUNCATE,
5809 BIND,
5810 LISTEN,
5811 RECV_ZC,
58095812
58105813 _,
58115814};
......@@ -6190,6 +6193,13 @@ pub const IORING_RESTRICTION = enum(u16) {
61906193 _,
61916194};
61926195
6196pub const IO_URING_SOCKET_OP = enum(u16) {
6197 SIOCIN = 0,
6198 SIOCOUTQ = 1,
6199 GETSOCKOPT = 2,
6200 SETSOCKOPT = 3,
6201};
6202
61936203pub const io_uring_buf = extern struct {
61946204 addr: u64,
61956205 len: u32,
lib/std/os/linux/IoUring.zig+186
......@@ -1356,6 +1356,55 @@ pub fn socket_direct_alloc(
13561356 return sqe;
13571357}
13581358
1359/// Queues (but does not submit) an SQE to perform an `bind(2)` on a socket.
1360/// Returns a pointer to the SQE.
1361/// Available since 6.11
1362pub fn bind(
1363 self: *IoUring,
1364 user_data: u64,
1365 fd: posix.fd_t,
1366 addr: *const posix.sockaddr,
1367 addrlen: posix.socklen_t,
1368 flags: u32,
1369) !*linux.io_uring_sqe {
1370 const sqe = try self.get_sqe();
1371 sqe.prep_bind(fd, addr, addrlen, flags);
1372 sqe.user_data = user_data;
1373 return sqe;
1374}
1375
1376/// Queues (but does not submit) an SQE to perform an `listen(2)` on a socket.
1377/// Returns a pointer to the SQE.
1378/// Available since 6.11
1379pub fn listen(
1380 self: *IoUring,
1381 user_data: u64,
1382 fd: posix.fd_t,
1383 backlog: usize,
1384 flags: u32,
1385) !*linux.io_uring_sqe {
1386 const sqe = try self.get_sqe();
1387 sqe.prep_listen(fd, backlog, flags);
1388 sqe.user_data = user_data;
1389 return sqe;
1390}
1391
1392fn cmd_sock(
1393 self: *IoUring,
1394 user_data: u64,
1395 cmd_op: linux.IO_URING_SOCKET_OP,
1396 fd: linux.fd_t,
1397 level: u32,
1398 optname: u32,
1399 optval: u64,
1400 optlen: u32,
1401) !*linux.io_uring_sqe {
1402 const sqe = try self.get_sqe();
1403 sqe.prep_cmd_sock(cmd_op, fd, level, optname, optval, optlen);
1404 sqe.user_data = user_data;
1405 return sqe;
1406}
1407
13591408pub const SubmissionQueue = struct {
13601409 head: *u32,
13611410 tail: *u32,
......@@ -4305,3 +4354,140 @@ test "copy_cqes with wrapping sq.cqes buffer" {
43054354 try testing.expectEqual(2 + 4 * i, ring.cq.head.*);
43064355 }
43074356}
4357
4358test "bind" {
4359 try skipKernelLessThan(.{ .major = 6, .minor = 11, .patch = 0 });
4360
4361 var ring = IoUring.init(4, 0) catch |err| switch (err) {
4362 error.SystemOutdated => return error.SkipZigTest,
4363 error.PermissionDenied => return error.SkipZigTest,
4364 else => return err,
4365 };
4366 defer ring.deinit();
4367
4368 var addr = net.Address.initIp4([4]u8{ 127, 0, 0, 1 }, 0);
4369 const proto: u32 = if (addr.any.family == linux.AF.UNIX) 0 else linux.IPPROTO.TCP;
4370
4371 const listen_fd = brk: {
4372 // Create socket
4373 _ = try ring.socket(1, addr.any.family, linux.SOCK.STREAM | linux.SOCK.CLOEXEC, proto, 0);
4374 try testing.expectEqual(1, try ring.submit());
4375 var cqe = try ring.copy_cqe();
4376 try testing.expectEqual(1, cqe.user_data);
4377 try testing.expectEqual(posix.E.SUCCESS, cqe.err());
4378 const listen_fd: posix.fd_t = @intCast(cqe.res);
4379 try testing.expect(listen_fd > 2);
4380
4381 // Prepare: set socket option * 2, bind, listen
4382 var sock_opt: u32 = 1;
4383 var sqe = try ring.cmd_sock(2, .SETSOCKOPT, listen_fd, linux.SOL.SOCKET, linux.SO.REUSEADDR, @intFromPtr(&sock_opt), @sizeOf(u32));
4384 sqe.flags |= linux.IOSQE_IO_LINK;
4385 sqe = try ring.cmd_sock(3, .SETSOCKOPT, listen_fd, linux.SOL.SOCKET, linux.SO.REUSEPORT, @intFromPtr(&sock_opt), @sizeOf(u32));
4386 sqe.flags |= linux.IOSQE_IO_LINK;
4387 sqe = try ring.bind(4, listen_fd, &addr.any, addr.getOsSockLen(), 0);
4388 sqe.flags |= linux.IOSQE_IO_LINK;
4389 _ = try ring.listen(5, listen_fd, 1, 0);
4390 // Submit 4 operations
4391 try testing.expectEqual(4, try ring.submit());
4392 // Expect all to succeed
4393 for (2..6) |user_data| {
4394 cqe = try ring.copy_cqe();
4395 try testing.expectEqual(user_data, cqe.user_data);
4396 try testing.expectEqual(posix.E.SUCCESS, cqe.err());
4397 }
4398
4399 // Check that socket option is set
4400 sock_opt = 0xff;
4401 _ = try ring.cmd_sock(5, .GETSOCKOPT, listen_fd, linux.SOL.SOCKET, linux.SO.REUSEADDR, @intFromPtr(&sock_opt), @sizeOf(u32));
4402 try testing.expectEqual(1, try ring.submit());
4403 cqe = try ring.copy_cqe();
4404 try testing.expectEqual(5, cqe.user_data);
4405 try testing.expectEqual(posix.E.SUCCESS, cqe.err());
4406 try testing.expectEqual(1, sock_opt);
4407
4408 // Read system assigned port into addr
4409 var addr_len: posix.socklen_t = addr.getOsSockLen();
4410 try posix.getsockname(listen_fd, &addr.any, &addr_len);
4411
4412 break :brk listen_fd;
4413 };
4414
4415 const connect_fd = brk: {
4416 // Create connect socket
4417 _ = try ring.socket(6, addr.any.family, linux.SOCK.STREAM | linux.SOCK.CLOEXEC, proto, 0);
4418 try testing.expectEqual(1, try ring.submit());
4419 const cqe = try ring.copy_cqe();
4420 try testing.expectEqual(6, cqe.user_data);
4421 try testing.expectEqual(posix.E.SUCCESS, cqe.err());
4422 // Get connect socket fd
4423 const connect_fd: posix.fd_t = @intCast(cqe.res);
4424 try testing.expect(connect_fd > 2 and connect_fd != listen_fd);
4425 break :brk connect_fd;
4426 };
4427
4428 // Prepare accept/connect operations
4429 _ = try ring.accept(7, listen_fd, null, null, 0);
4430 _ = try ring.connect(8, connect_fd, &addr.any, addr.getOsSockLen());
4431 try testing.expectEqual(2, try ring.submit());
4432 // Get listener accepted socket
4433 var accept_fd: posix.socket_t = 0;
4434 for (0..2) |_| {
4435 const cqe = try ring.copy_cqe();
4436 try testing.expectEqual(posix.E.SUCCESS, cqe.err());
4437 if (cqe.user_data == 7) {
4438 accept_fd = @intCast(cqe.res);
4439 } else {
4440 try testing.expectEqual(8, cqe.user_data);
4441 }
4442 }
4443 try testing.expect(accept_fd > 2 and accept_fd != listen_fd and accept_fd != connect_fd);
4444
4445 // Communicate
4446 try testSendRecv(&ring, connect_fd, accept_fd);
4447 try testSendRecv(&ring, accept_fd, connect_fd);
4448
4449 // Shutdown and close all sockets
4450 for ([_]posix.socket_t{ connect_fd, accept_fd, listen_fd }) |fd| {
4451 var sqe = try ring.shutdown(9, fd, posix.SHUT.RDWR);
4452 sqe.flags |= linux.IOSQE_IO_LINK;
4453 _ = try ring.close(10, fd);
4454 try testing.expectEqual(2, try ring.submit());
4455 for (0..2) |i| {
4456 const cqe = try ring.copy_cqe();
4457 try testing.expectEqual(posix.E.SUCCESS, cqe.err());
4458 try testing.expectEqual(9 + i, cqe.user_data);
4459 }
4460 }
4461}
4462
4463fn testSendRecv(ring: *IoUring, send_fd: posix.socket_t, recv_fd: posix.socket_t) !void {
4464 const buffer_send = "0123456789abcdf" ** 10;
4465 var buffer_recv: [buffer_send.len * 2]u8 = undefined;
4466
4467 // 2 sends
4468 var sqe = try ring.send(1, send_fd, buffer_send, linux.MSG.WAITALL);
4469 sqe.flags |= linux.IOSQE_IO_LINK;
4470 _ = try ring.send(2, send_fd, buffer_send, linux.MSG.WAITALL);
4471 try testing.expectEqual(2, try ring.submit());
4472 for (0..2) |i| {
4473 const cqe = try ring.copy_cqe();
4474 try testing.expectEqual(1 + i, cqe.user_data);
4475 try testing.expectEqual(posix.E.SUCCESS, cqe.err());
4476 try testing.expectEqual(buffer_send.len, @as(usize, @intCast(cqe.res)));
4477 }
4478
4479 // receive
4480 var recv_len: usize = 0;
4481 while (recv_len < buffer_send.len * 2) {
4482 _ = try ring.recv(3, recv_fd, .{ .buffer = buffer_recv[recv_len..] }, 0);
4483 try testing.expectEqual(1, try ring.submit());
4484 const cqe = try ring.copy_cqe();
4485 try testing.expectEqual(3, cqe.user_data);
4486 try testing.expectEqual(posix.E.SUCCESS, cqe.err());
4487 recv_len += @intCast(cqe.res);
4488 }
4489
4490 // inspect recv buffer
4491 try testing.expectEqualSlices(u8, buffer_send, buffer_recv[0..buffer_send.len]);
4492 try testing.expectEqualSlices(u8, buffer_send, buffer_recv[buffer_send.len..]);
4493}
lib/std/os/linux/io_uring_sqe.zig+47
......@@ -619,4 +619,51 @@ pub const io_uring_sqe = extern struct {
619619 sqe.rw_flags = flags;
620620 sqe.splice_fd_in = @bitCast(options);
621621 }
622
623 pub fn prep_bind(
624 sqe: *linux.io_uring_sqe,
625 fd: linux.fd_t,
626 addr: *const linux.sockaddr,
627 addrlen: linux.socklen_t,
628 flags: u32,
629 ) void {
630 sqe.prep_rw(.BIND, fd, @intFromPtr(addr), 0, addrlen);
631 sqe.rw_flags = flags;
632 }
633
634 pub fn prep_listen(
635 sqe: *linux.io_uring_sqe,
636 fd: linux.fd_t,
637 backlog: usize,
638 flags: u32,
639 ) void {
640 sqe.prep_rw(.LISTEN, fd, 0, backlog, 0);
641 sqe.rw_flags = flags;
642 }
643
644 pub fn prep_cmd_sock(
645 sqe: *linux.io_uring_sqe,
646 cmd_op: linux.IO_URING_SOCKET_OP,
647 fd: linux.fd_t,
648 level: u32,
649 optname: u32,
650 optval: u64,
651 optlen: u32,
652 ) void {
653 sqe.prep_rw(.URING_CMD, fd, 0, 0, 0);
654 // off is overloaded with cmd_op, https://github.com/axboe/liburing/blob/e1003e496e66f9b0ae06674869795edf772d5500/src/include/liburing/io_uring.h#L39
655 sqe.off = @intFromEnum(cmd_op);
656 // addr is overloaded, https://github.com/axboe/liburing/blob/e1003e496e66f9b0ae06674869795edf772d5500/src/include/liburing/io_uring.h#L46
657 sqe.addr = @bitCast(packed struct {
658 level: u32,
659 optname: u32,
660 }{
661 .level = level,
662 .optname = optname,
663 });
664 // splice_fd_in if overloaded u32 -> i32
665 sqe.splice_fd_in = @bitCast(optlen);
666 // addr3 is overloaded, https://github.com/axboe/liburing/blob/e1003e496e66f9b0ae06674869795edf772d5500/src/include/liburing/io_uring.h#L102
667 sqe.addr3 = optval;
668 }
622669};