| ... | ... | @@ -505,6 +505,24 @@ pub const IO_Uring = struct { |
| 505 | 505 | return sqe; |
| 506 | 506 | } |
| 507 | 507 | |
| 508 | /// Queues (but does not submit) an SQE to perform an multishot `accept4(2)` on a socket. |
| 509 | /// 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 | /// Returns a pointer to the SQE. |
| 512 | pub fn accept_multishot( |
| 513 | self: *IO_Uring, |
| 514 | user_data: u64, |
| 515 | fd: os.fd_t, |
| 516 | addr: ?*os.sockaddr, |
| 517 | addrlen: ?*os.socklen_t, |
| 518 | flags: u32, |
| 519 | ) !*linux.io_uring_sqe { |
| 520 | const sqe = try self.get_sqe(); |
| 521 | io_uring_prep_multishot_accept(sqe, fd, addr, addrlen, flags); |
| 522 | sqe.user_data = user_data; |
| 523 | return sqe; |
| 524 | } |
| 525 | |
| 508 | 526 | /// Queue (but does not submit) an SQE to perform a `connect(2)` on a socket. |
| 509 | 527 | /// Returns a pointer to the SQE. |
| 510 | 528 | pub fn connect( |
| ... | ... | @@ -1621,6 +1639,17 @@ pub fn io_uring_prep_remove_buffers( |
| 1621 | 1639 | sqe.buf_index = @intCast(group_id); |
| 1622 | 1640 | } |
| 1623 | 1641 | |
| 1642 | pub fn io_uring_prep_multishot_accept( |
| 1643 | sqe: *linux.io_uring_sqe, |
| 1644 | fd: os.fd_t, |
| 1645 | addr: ?*os.sockaddr, |
| 1646 | addrlen: ?*os.socklen_t, |
| 1647 | flags: u32, |
| 1648 | ) void { |
| 1649 | io_uring_prep_accept(sqe, fd, addr, addrlen, flags); |
| 1650 | sqe.ioprio |= linux.IORING_ACCEPT_MULTISHOT; |
| 1651 | } |
| 1652 | |
| 1624 | 1653 | test "structs/offsets/entries" { |
| 1625 | 1654 | if (builtin.os.tag != .linux) return error.SkipZigTest; |
| 1626 | 1655 | |
| ... | ... | @@ -3353,20 +3382,10 @@ const SocketTestHarness = struct { |
| 3353 | 3382 | |
| 3354 | 3383 | fn createSocketTestHarness(ring: *IO_Uring) !SocketTestHarness { |
| 3355 | 3384 | // Create a TCP server socket |
| 3356 | | |
| 3357 | 3385 | var address = try net.Address.parseIp4("127.0.0.1", 0); |
| 3358 | | const kernel_backlog = 1; |
| 3359 | | const listener_socket = try os.socket(address.any.family, os.SOCK.STREAM | os.SOCK.CLOEXEC, 0); |
| 3386 | const listener_socket = try createListenerSocket(&address); |
| 3360 | 3387 | errdefer os.closeSocket(listener_socket); |
| 3361 | 3388 | |
| 3362 | | try os.setsockopt(listener_socket, os.SOL.SOCKET, os.SO.REUSEADDR, &mem.toBytes(@as(c_int, 1))); |
| 3363 | | try os.bind(listener_socket, &address.any, address.getOsSockLen()); |
| 3364 | | try os.listen(listener_socket, kernel_backlog); |
| 3365 | | |
| 3366 | | // set address to the OS-chosen IP/port. |
| 3367 | | var slen: os.socklen_t = address.getOsSockLen(); |
| 3368 | | try os.getsockname(listener_socket, &address.any, &slen); |
| 3369 | | |
| 3370 | 3389 | // Submit 1 accept |
| 3371 | 3390 | var accept_addr: os.sockaddr = undefined; |
| 3372 | 3391 | var accept_addr_len: os.socklen_t = @sizeOf(@TypeOf(accept_addr)); |
| ... | ... | @@ -3410,3 +3429,58 @@ fn createSocketTestHarness(ring: *IO_Uring) !SocketTestHarness { |
| 3410 | 3429 | .client = client, |
| 3411 | 3430 | }; |
| 3412 | 3431 | } |
| 3432 | |
| 3433 | fn createListenerSocket(address: *net.Address) !os.socket_t { |
| 3434 | const kernel_backlog = 1; |
| 3435 | const listener_socket = try os.socket(address.any.family, os.SOCK.STREAM | os.SOCK.CLOEXEC, 0); |
| 3436 | errdefer os.closeSocket(listener_socket); |
| 3437 | |
| 3438 | try os.setsockopt(listener_socket, os.SOL.SOCKET, os.SO.REUSEADDR, &mem.toBytes(@as(c_int, 1))); |
| 3439 | try os.bind(listener_socket, &address.any, address.getOsSockLen()); |
| 3440 | try os.listen(listener_socket, kernel_backlog); |
| 3441 | |
| 3442 | // set address to the OS-chosen IP/port. |
| 3443 | var slen: os.socklen_t = address.getOsSockLen(); |
| 3444 | try os.getsockname(listener_socket, &address.any, &slen); |
| 3445 | |
| 3446 | return listener_socket; |
| 3447 | } |
| 3448 | |
| 3449 | test "accept multishot" { |
| 3450 | if (builtin.os.tag != .linux) return error.SkipZigTest; |
| 3451 | |
| 3452 | var ring = IO_Uring.init(16, 0) catch |err| switch (err) { |
| 3453 | error.SystemOutdated => return error.SkipZigTest, |
| 3454 | error.PermissionDenied => return error.SkipZigTest, |
| 3455 | else => return err, |
| 3456 | }; |
| 3457 | defer ring.deinit(); |
| 3458 | |
| 3459 | var address = try net.Address.parseIp4("127.0.0.1", 0); |
| 3460 | const listener_socket = try createListenerSocket(&address); |
| 3461 | defer os.closeSocket(listener_socket); |
| 3462 | |
| 3463 | // submit multishot accept operation |
| 3464 | var addr: os.sockaddr = undefined; |
| 3465 | var addr_len: os.socklen_t = @sizeOf(@TypeOf(addr)); |
| 3466 | const userdata: u64 = 0xaaaaaaaa; |
| 3467 | _ = try ring.accept_multishot(userdata, listener_socket, &addr, &addr_len, 0); |
| 3468 | try testing.expectEqual(@as(u32, 1), try ring.submit()); |
| 3469 | |
| 3470 | var nr: usize = 4; // number of clients to connect |
| 3471 | while (nr > 0) : (nr -= 1) { |
| 3472 | // connect client |
| 3473 | var client = try os.socket(address.any.family, os.SOCK.STREAM | os.SOCK.CLOEXEC, 0); |
| 3474 | errdefer os.closeSocket(client); |
| 3475 | try os.connect(client, &address.any, address.getOsSockLen()); |
| 3476 | |
| 3477 | // test accept completion |
| 3478 | var cqe = try ring.copy_cqe(); |
| 3479 | if (cqe.err() == .INVAL) return error.SkipZigTest; |
| 3480 | try testing.expect(cqe.res > 0); |
| 3481 | try testing.expect(cqe.user_data == userdata); |
| 3482 | try testing.expect(cqe.flags & linux.IORING_CQE_F_MORE > 0); // more flag is set |
| 3483 | |
| 3484 | os.closeSocket(client); |
| 3485 | } |
| 3486 | } |