| ... | @@ -1272,6 +1272,16 @@ pub fn unregister_buffers(self: *IoUring) !void { | ... | @@ -1272,6 +1272,16 @@ pub fn unregister_buffers(self: *IoUring) !void { |
| 1272 | } | 1272 | } |
| 1273 | } | 1273 | } |
| 1274 | | 1274 | |
| | 1275 | /// Returns a io_uring_probe which is used to probe the capabilities of the |
| | 1276 | /// io_uring subsystem of the running kernel. The io_uring_probe contains the |
| | 1277 | /// list of supported operations. |
| | 1278 | pub fn get_probe(self: *IoUring) !linux.io_uring_probe { |
| | 1279 | var probe = mem.zeroInit(linux.io_uring_probe, .{}); |
| | 1280 | const res = linux.io_uring_register(self.fd, .REGISTER_PROBE, &probe, probe.ops.len); |
| | 1281 | try handle_register_buf_ring_result(res); |
| | 1282 | return probe; |
| | 1283 | } |
| | 1284 | |
| 1275 | fn handle_registration_result(res: usize) !void { | 1285 | fn handle_registration_result(res: usize) !void { |
| 1276 | switch (linux.E.init(res)) { | 1286 | switch (linux.E.init(res)) { |
| 1277 | .SUCCESS => {}, | 1287 | .SUCCESS => {}, |
| ... | @@ -1356,6 +1366,102 @@ pub fn socket_direct_alloc( | ... | @@ -1356,6 +1366,102 @@ pub fn socket_direct_alloc( |
| 1356 | return sqe; | 1366 | return sqe; |
| 1357 | } | 1367 | } |
| 1358 | | 1368 | |
| | 1369 | /// Queues (but does not submit) an SQE to perform an `bind(2)` on a socket. |
| | 1370 | /// Returns a pointer to the SQE. |
| | 1371 | /// Available since 6.11 |
| | 1372 | pub fn bind( |
| | 1373 | self: *IoUring, |
| | 1374 | user_data: u64, |
| | 1375 | fd: posix.fd_t, |
| | 1376 | addr: *const posix.sockaddr, |
| | 1377 | addrlen: posix.socklen_t, |
| | 1378 | flags: u32, |
| | 1379 | ) !*linux.io_uring_sqe { |
| | 1380 | const sqe = try self.get_sqe(); |
| | 1381 | sqe.prep_bind(fd, addr, addrlen, flags); |
| | 1382 | sqe.user_data = user_data; |
| | 1383 | return sqe; |
| | 1384 | } |
| | 1385 | |
| | 1386 | /// Queues (but does not submit) an SQE to perform an `listen(2)` on a socket. |
| | 1387 | /// Returns a pointer to the SQE. |
| | 1388 | /// Available since 6.11 |
| | 1389 | pub fn listen( |
| | 1390 | self: *IoUring, |
| | 1391 | user_data: u64, |
| | 1392 | fd: posix.fd_t, |
| | 1393 | backlog: usize, |
| | 1394 | flags: u32, |
| | 1395 | ) !*linux.io_uring_sqe { |
| | 1396 | const sqe = try self.get_sqe(); |
| | 1397 | sqe.prep_listen(fd, backlog, flags); |
| | 1398 | sqe.user_data = user_data; |
| | 1399 | return sqe; |
| | 1400 | } |
| | 1401 | |
| | 1402 | /// Prepares an cmd request for a socket. |
| | 1403 | /// See: https://man7.org/linux/man-pages/man3/io_uring_prep_cmd.3.html |
| | 1404 | /// Available since 6.7. |
| | 1405 | pub fn cmd_sock( |
| | 1406 | self: *IoUring, |
| | 1407 | user_data: u64, |
| | 1408 | cmd_op: linux.IO_URING_SOCKET_OP, |
| | 1409 | fd: linux.fd_t, |
| | 1410 | level: u32, // linux.SOL |
| | 1411 | optname: u32, // linux.SO |
| | 1412 | optval: u64, // pointer to the option value |
| | 1413 | optlen: u32, // size of the option value |
| | 1414 | ) !*linux.io_uring_sqe { |
| | 1415 | const sqe = try self.get_sqe(); |
| | 1416 | sqe.prep_cmd_sock(cmd_op, fd, level, optname, optval, optlen); |
| | 1417 | sqe.user_data = user_data; |
| | 1418 | return sqe; |
| | 1419 | } |
| | 1420 | |
| | 1421 | /// Prepares set socket option for the optname argument, at the protocol |
| | 1422 | /// level specified by the level argument. |
| | 1423 | /// Available since 6.7.n |
| | 1424 | pub fn setsockopt( |
| | 1425 | self: *IoUring, |
| | 1426 | user_data: u64, |
| | 1427 | fd: linux.fd_t, |
| | 1428 | level: u32, // linux.SOL |
| | 1429 | optname: u32, // linux.SO |
| | 1430 | opt: []const u8, |
| | 1431 | ) !*linux.io_uring_sqe { |
| | 1432 | return try self.cmd_sock( |
| | 1433 | user_data, |
| | 1434 | .SETSOCKOPT, |
| | 1435 | fd, |
| | 1436 | level, |
| | 1437 | optname, |
| | 1438 | @intFromPtr(opt.ptr), |
| | 1439 | @intCast(opt.len), |
| | 1440 | ); |
| | 1441 | } |
| | 1442 | |
| | 1443 | /// Prepares get socket option to retrieve the value for the option specified by |
| | 1444 | /// the option_name argument for the socket specified by the fd argument. |
| | 1445 | /// Available since 6.7. |
| | 1446 | pub fn getsockopt( |
| | 1447 | self: *IoUring, |
| | 1448 | user_data: u64, |
| | 1449 | fd: linux.fd_t, |
| | 1450 | level: u32, // linux.SOL |
| | 1451 | optname: u32, // linux.SO |
| | 1452 | opt: []u8, |
| | 1453 | ) !*linux.io_uring_sqe { |
| | 1454 | return try self.cmd_sock( |
| | 1455 | user_data, |
| | 1456 | .GETSOCKOPT, |
| | 1457 | fd, |
| | 1458 | level, |
| | 1459 | optname, |
| | 1460 | @intFromPtr(opt.ptr), |
| | 1461 | @intCast(opt.len), |
| | 1462 | ); |
| | 1463 | } |
| | 1464 | |
| 1359 | pub const SubmissionQueue = struct { | 1465 | pub const SubmissionQueue = struct { |
| 1360 | head: *u32, | 1466 | head: *u32, |
| 1361 | tail: *u32, | 1467 | tail: *u32, |
| ... | @@ -1488,28 +1594,34 @@ pub const BufferGroup = struct { | ... | @@ -1488,28 +1594,34 @@ pub const BufferGroup = struct { |
| 1488 | buffers: []u8, | 1594 | buffers: []u8, |
| 1489 | /// Size of each buffer in buffers. | 1595 | /// Size of each buffer in buffers. |
| 1490 | buffer_size: u32, | 1596 | buffer_size: u32, |
| 1491 | // Number of buffers in `buffers`, number of `io_uring_buf structures` in br. | 1597 | /// Number of buffers in `buffers`, number of `io_uring_buf structures` in br. |
| 1492 | buffers_count: u16, | 1598 | buffers_count: u16, |
| | 1599 | /// Head of unconsumed part of each buffer, if incremental consumption is enabled |
| | 1600 | heads: []u32, |
| 1493 | /// ID of this group, must be unique in ring. | 1601 | /// ID of this group, must be unique in ring. |
| 1494 | group_id: u16, | 1602 | group_id: u16, |
| 1495 | | 1603 | |
| 1496 | pub fn init( | 1604 | pub fn init( |
| 1497 | ring: *IoUring, | 1605 | ring: *IoUring, |
| | 1606 | allocator: mem.Allocator, |
| 1498 | group_id: u16, | 1607 | group_id: u16, |
| 1499 | buffers: []u8, | | |
| 1500 | buffer_size: u32, | 1608 | buffer_size: u32, |
| 1501 | buffers_count: u16, | 1609 | buffers_count: u16, |
| 1502 | ) !BufferGroup { | 1610 | ) !BufferGroup { |
| 1503 | assert(buffers.len == buffers_count * buffer_size); | 1611 | const buffers = try allocator.alloc(u8, buffer_size * buffers_count); |
| | 1612 | errdefer allocator.free(buffers); |
| | 1613 | const heads = try allocator.alloc(u32, buffers_count); |
| | 1614 | errdefer allocator.free(heads); |
| 1504 | | 1615 | |
| 1505 | const br = try setup_buf_ring(ring.fd, buffers_count, group_id); | 1616 | const br = try setup_buf_ring(ring.fd, buffers_count, group_id, .{ .inc = true }); |
| 1506 | buf_ring_init(br); | 1617 | buf_ring_init(br); |
| 1507 | | 1618 | |
| 1508 | const mask = buf_ring_mask(buffers_count); | 1619 | const mask = buf_ring_mask(buffers_count); |
| 1509 | var i: u16 = 0; | 1620 | var i: u16 = 0; |
| 1510 | while (i < buffers_count) : (i += 1) { | 1621 | while (i < buffers_count) : (i += 1) { |
| 1511 | const start = buffer_size * i; | 1622 | const pos = buffer_size * i; |
| 1512 | const buf = buffers[start .. start + buffer_size]; | 1623 | const buf = buffers[pos .. pos + buffer_size]; |
| | 1624 | heads[i] = 0; |
| 1513 | buf_ring_add(br, buf, i, mask, i); | 1625 | buf_ring_add(br, buf, i, mask, i); |
| 1514 | } | 1626 | } |
| 1515 | buf_ring_advance(br, buffers_count); | 1627 | buf_ring_advance(br, buffers_count); |
| ... | @@ -1519,11 +1631,18 @@ pub const BufferGroup = struct { | ... | @@ -1519,11 +1631,18 @@ pub const BufferGroup = struct { |
| 1519 | .group_id = group_id, | 1631 | .group_id = group_id, |
| 1520 | .br = br, | 1632 | .br = br, |
| 1521 | .buffers = buffers, | 1633 | .buffers = buffers, |
| | 1634 | .heads = heads, |
| 1522 | .buffer_size = buffer_size, | 1635 | .buffer_size = buffer_size, |
| 1523 | .buffers_count = buffers_count, | 1636 | .buffers_count = buffers_count, |
| 1524 | }; | 1637 | }; |
| 1525 | } | 1638 | } |
| 1526 | | 1639 | |
| | 1640 | pub fn deinit(self: *BufferGroup, allocator: mem.Allocator) void { |
| | 1641 | free_buf_ring(self.ring.fd, self.br, self.buffers_count, self.group_id); |
| | 1642 | allocator.free(self.buffers); |
| | 1643 | allocator.free(self.heads); |
| | 1644 | } |
| | 1645 | |
| 1527 | // Prepare recv operation which will select buffer from this group. | 1646 | // Prepare recv operation which will select buffer from this group. |
| 1528 | pub fn recv(self: *BufferGroup, user_data: u64, fd: posix.fd_t, flags: u32) !*linux.io_uring_sqe { | 1647 | pub fn recv(self: *BufferGroup, user_data: u64, fd: posix.fd_t, flags: u32) !*linux.io_uring_sqe { |
| 1529 | var sqe = try self.ring.get_sqe(); | 1648 | var sqe = try self.ring.get_sqe(); |
| ... | @@ -1543,33 +1662,34 @@ pub const BufferGroup = struct { | ... | @@ -1543,33 +1662,34 @@ pub const BufferGroup = struct { |
| 1543 | } | 1662 | } |
| 1544 | | 1663 | |
| 1545 | // Get buffer by id. | 1664 | // Get buffer by id. |
| 1546 | pub fn get(self: *BufferGroup, buffer_id: u16) []u8 { | 1665 | fn get_by_id(self: *BufferGroup, buffer_id: u16) []u8 { |
| 1547 | const head = self.buffer_size * buffer_id; | 1666 | const pos = self.buffer_size * buffer_id; |
| 1548 | return self.buffers[head .. head + self.buffer_size]; | 1667 | return self.buffers[pos .. pos + self.buffer_size][self.heads[buffer_id]..]; |
| 1549 | } | 1668 | } |
| 1550 | | 1669 | |
| 1551 | // Get buffer by CQE. | 1670 | // Get buffer by CQE. |
| 1552 | pub fn get_cqe(self: *BufferGroup, cqe: linux.io_uring_cqe) ![]u8 { | 1671 | pub fn get(self: *BufferGroup, cqe: linux.io_uring_cqe) ![]u8 { |
| 1553 | const buffer_id = try cqe.buffer_id(); | 1672 | const buffer_id = try cqe.buffer_id(); |
| 1554 | const used_len = @as(usize, @intCast(cqe.res)); | 1673 | const used_len = @as(usize, @intCast(cqe.res)); |
| 1555 | return self.get(buffer_id)[0..used_len]; | 1674 | return self.get_by_id(buffer_id)[0..used_len]; |
| 1556 | } | | |
| 1557 | | | |
| 1558 | // Release buffer to the kernel. | | |
| 1559 | pub fn put(self: *BufferGroup, buffer_id: u16) void { | | |
| 1560 | const mask = buf_ring_mask(self.buffers_count); | | |
| 1561 | const buffer = self.get(buffer_id); | | |
| 1562 | buf_ring_add(self.br, buffer, buffer_id, mask, 0); | | |
| 1563 | buf_ring_advance(self.br, 1); | | |
| 1564 | } | 1675 | } |
| 1565 | | 1676 | |
| 1566 | // Release buffer from CQE to the kernel. | 1677 | // Release buffer from CQE to the kernel. |
| 1567 | pub fn put_cqe(self: *BufferGroup, cqe: linux.io_uring_cqe) !void { | 1678 | pub fn put(self: *BufferGroup, cqe: linux.io_uring_cqe) !void { |
| 1568 | self.put(try cqe.buffer_id()); | 1679 | const buffer_id = try cqe.buffer_id(); |
| 1569 | } | 1680 | if (cqe.flags & linux.IORING_CQE_F_BUF_MORE == linux.IORING_CQE_F_BUF_MORE) { |
| | 1681 | // Incremental consumption active, kernel will write to the this buffer again |
| | 1682 | const used_len = @as(u32, @intCast(cqe.res)); |
| | 1683 | // Track what part of the buffer is used |
| | 1684 | self.heads[buffer_id] += used_len; |
| | 1685 | return; |
| | 1686 | } |
| | 1687 | self.heads[buffer_id] = 0; |
| 1570 | | 1688 | |
| 1571 | pub fn deinit(self: *BufferGroup) void { | 1689 | // Release buffer to the kernel. const mask = buf_ring_mask(self.buffers_count); |
| 1572 | free_buf_ring(self.ring.fd, self.br, self.buffers_count, self.group_id); | 1690 | const mask = buf_ring_mask(self.buffers_count); |
| | 1691 | buf_ring_add(self.br, self.get_by_id(buffer_id), buffer_id, mask, 0); |
| | 1692 | buf_ring_advance(self.br, 1); |
| 1573 | } | 1693 | } |
| 1574 | }; | 1694 | }; |
| 1575 | | 1695 | |
| ... | @@ -1578,7 +1698,12 @@ pub const BufferGroup = struct { | ... | @@ -1578,7 +1698,12 @@ pub const BufferGroup = struct { |
| 1578 | /// `fd` is IO_Uring.fd for which the provided buffer ring is being registered. | 1698 | /// `fd` is IO_Uring.fd for which the provided buffer ring is being registered. |
| 1579 | /// `entries` is the number of entries requested in the buffer ring, must be power of 2. | 1699 | /// `entries` is the number of entries requested in the buffer ring, must be power of 2. |
| 1580 | /// `group_id` is the chosen buffer group ID, unique in IO_Uring. | 1700 | /// `group_id` is the chosen buffer group ID, unique in IO_Uring. |
| 1581 | pub fn setup_buf_ring(fd: posix.fd_t, entries: u16, group_id: u16) !*align(page_size_min) linux.io_uring_buf_ring { | 1701 | pub fn setup_buf_ring( |
| | 1702 | fd: posix.fd_t, |
| | 1703 | entries: u16, |
| | 1704 | group_id: u16, |
| | 1705 | flags: linux.io_uring_buf_reg.Flags, |
| | 1706 | ) !*align(page_size_min) linux.io_uring_buf_ring { |
| 1582 | if (entries == 0 or entries > 1 << 15) return error.EntriesNotInRange; | 1707 | if (entries == 0 or entries > 1 << 15) return error.EntriesNotInRange; |
| 1583 | if (!std.math.isPowerOfTwo(entries)) return error.EntriesNotPowerOfTwo; | 1708 | if (!std.math.isPowerOfTwo(entries)) return error.EntriesNotPowerOfTwo; |
| 1584 | | 1709 | |
| ... | @@ -1595,22 +1720,30 @@ pub fn setup_buf_ring(fd: posix.fd_t, entries: u16, group_id: u16) !*align(page_ | ... | @@ -1595,22 +1720,30 @@ pub fn setup_buf_ring(fd: posix.fd_t, entries: u16, group_id: u16) !*align(page_ |
| 1595 | assert(mmap.len == mmap_size); | 1720 | assert(mmap.len == mmap_size); |
| 1596 | | 1721 | |
| 1597 | const br: *align(page_size_min) linux.io_uring_buf_ring = @ptrCast(mmap.ptr); | 1722 | const br: *align(page_size_min) linux.io_uring_buf_ring = @ptrCast(mmap.ptr); |
| 1598 | try register_buf_ring(fd, @intFromPtr(br), entries, group_id); | 1723 | try register_buf_ring(fd, @intFromPtr(br), entries, group_id, flags); |
| 1599 | return br; | 1724 | return br; |
| 1600 | } | 1725 | } |
| 1601 | | 1726 | |
| 1602 | fn register_buf_ring(fd: posix.fd_t, addr: u64, entries: u32, group_id: u16) !void { | 1727 | fn register_buf_ring( |
| | 1728 | fd: posix.fd_t, |
| | 1729 | addr: u64, |
| | 1730 | entries: u32, |
| | 1731 | group_id: u16, |
| | 1732 | flags: linux.io_uring_buf_reg.Flags, |
| | 1733 | ) !void { |
| 1603 | var reg = mem.zeroInit(linux.io_uring_buf_reg, .{ | 1734 | var reg = mem.zeroInit(linux.io_uring_buf_reg, .{ |
| 1604 | .ring_addr = addr, | 1735 | .ring_addr = addr, |
| 1605 | .ring_entries = entries, | 1736 | .ring_entries = entries, |
| 1606 | .bgid = group_id, | 1737 | .bgid = group_id, |
| | 1738 | .flags = flags, |
| 1607 | }); | 1739 | }); |
| 1608 | const res = linux.io_uring_register( | 1740 | var res = linux.io_uring_register(fd, .REGISTER_PBUF_RING, @as(*const anyopaque, @ptrCast(&reg)), 1); |
| 1609 | fd, | 1741 | if (linux.E.init(res) == .INVAL and reg.flags.inc) { |
| 1610 | .REGISTER_PBUF_RING, | 1742 | // Retry without incremental buffer consumption. |
| 1611 | @as(*const anyopaque, @ptrCast(&reg)), | 1743 | // It is available since kernel 6.12. returns INVAL on older. |
| 1612 | 1, | 1744 | reg.flags.inc = false; |
| 1613 | ); | 1745 | res = linux.io_uring_register(fd, .REGISTER_PBUF_RING, @as(*const anyopaque, @ptrCast(&reg)), 1); |
| | 1746 | } |
| 1614 | try handle_register_buf_ring_result(res); | 1747 | try handle_register_buf_ring_result(res); |
| 1615 | } | 1748 | } |
| 1616 | | 1749 | |
| ... | @@ -3054,7 +3187,7 @@ test "provide_buffers: read" { | ... | @@ -3054,7 +3187,7 @@ test "provide_buffers: read" { |
| 3054 | const cqe = try ring.copy_cqe(); | 3187 | const cqe = try ring.copy_cqe(); |
| 3055 | switch (cqe.err()) { | 3188 | switch (cqe.err()) { |
| 3056 | // Happens when the kernel is < 5.7 | 3189 | // Happens when the kernel is < 5.7 |
| 3057 | .INVAL => return error.SkipZigTest, | 3190 | .INVAL, .BADF => return error.SkipZigTest, |
| 3058 | .SUCCESS => {}, | 3191 | .SUCCESS => {}, |
| 3059 | else => |errno| std.debug.panic("unhandled errno: {}", .{errno}), | 3192 | else => |errno| std.debug.panic("unhandled errno: {}", .{errno}), |
| 3060 | } | 3193 | } |
| ... | @@ -3181,7 +3314,7 @@ test "remove_buffers" { | ... | @@ -3181,7 +3314,7 @@ test "remove_buffers" { |
| 3181 | | 3314 | |
| 3182 | const cqe = try ring.copy_cqe(); | 3315 | const cqe = try ring.copy_cqe(); |
| 3183 | switch (cqe.err()) { | 3316 | switch (cqe.err()) { |
| 3184 | .INVAL => return error.SkipZigTest, | 3317 | .INVAL, .BADF => return error.SkipZigTest, |
| 3185 | .SUCCESS => {}, | 3318 | .SUCCESS => {}, |
| 3186 | else => |errno| std.debug.panic("unhandled errno: {}", .{errno}), | 3319 | else => |errno| std.debug.panic("unhandled errno: {}", .{errno}), |
| 3187 | } | 3320 | } |
| ... | @@ -3935,12 +4068,10 @@ test BufferGroup { | ... | @@ -3935,12 +4068,10 @@ test BufferGroup { |
| 3935 | const group_id: u16 = 1; // buffers group id | 4068 | const group_id: u16 = 1; // buffers group id |
| 3936 | const buffers_count: u16 = 1; // number of buffers in buffer group | 4069 | const buffers_count: u16 = 1; // number of buffers in buffer group |
| 3937 | const buffer_size: usize = 128; // size of each buffer in group | 4070 | const buffer_size: usize = 128; // size of each buffer in group |
| 3938 | const buffers = try testing.allocator.alloc(u8, buffers_count * buffer_size); | | |
| 3939 | defer testing.allocator.free(buffers); | | |
| 3940 | var buf_grp = BufferGroup.init( | 4071 | var buf_grp = BufferGroup.init( |
| 3941 | &ring, | 4072 | &ring, |
| | 4073 | testing.allocator, |
| 3942 | group_id, | 4074 | group_id, |
| 3943 | buffers, | | |
| 3944 | buffer_size, | 4075 | buffer_size, |
| 3945 | buffers_count, | 4076 | buffers_count, |
| 3946 | ) catch |err| switch (err) { | 4077 | ) catch |err| switch (err) { |
| ... | @@ -3948,7 +4079,7 @@ test BufferGroup { | ... | @@ -3948,7 +4079,7 @@ test BufferGroup { |
| 3948 | error.ArgumentsInvalid => return error.SkipZigTest, | 4079 | error.ArgumentsInvalid => return error.SkipZigTest, |
| 3949 | else => return err, | 4080 | else => return err, |
| 3950 | }; | 4081 | }; |
| 3951 | defer buf_grp.deinit(); | 4082 | defer buf_grp.deinit(testing.allocator); |
| 3952 | | 4083 | |
| 3953 | // Create client/server fds | 4084 | // Create client/server fds |
| 3954 | const fds = try createSocketTestHarness(&ring); | 4085 | const fds = try createSocketTestHarness(&ring); |
| ... | @@ -3979,14 +4110,11 @@ test BufferGroup { | ... | @@ -3979,14 +4110,11 @@ test BufferGroup { |
| 3979 | try testing.expectEqual(posix.E.SUCCESS, cqe.err()); | 4110 | try testing.expectEqual(posix.E.SUCCESS, cqe.err()); |
| 3980 | try testing.expectEqual(data.len, @as(usize, @intCast(cqe.res))); // cqe.res holds received data len | 4111 | try testing.expectEqual(data.len, @as(usize, @intCast(cqe.res))); // cqe.res holds received data len |
| 3981 | | 4112 | |
| 3982 | // Read buffer_id and used buffer len from cqe | | |
| 3983 | const buffer_id = try cqe.buffer_id(); | | |
| 3984 | const len: usize = @intCast(cqe.res); | | |
| 3985 | // Get buffer from pool | 4113 | // Get buffer from pool |
| 3986 | const buf = buf_grp.get(buffer_id)[0..len]; | 4114 | const buf = try buf_grp.get(cqe); |
| 3987 | try testing.expectEqualSlices(u8, &data, buf); | 4115 | try testing.expectEqualSlices(u8, &data, buf); |
| 3988 | // Release buffer to the kernel when application is done with it | 4116 | // Release buffer to the kernel when application is done with it |
| 3989 | buf_grp.put(buffer_id); | 4117 | try buf_grp.put(cqe); |
| 3990 | } | 4118 | } |
| 3991 | } | 4119 | } |
| 3992 | | 4120 | |
| ... | @@ -4004,12 +4132,10 @@ test "ring mapped buffers recv" { | ... | @@ -4004,12 +4132,10 @@ test "ring mapped buffers recv" { |
| 4004 | const group_id: u16 = 1; // buffers group id | 4132 | const group_id: u16 = 1; // buffers group id |
| 4005 | const buffers_count: u16 = 2; // number of buffers in buffer group | 4133 | const buffers_count: u16 = 2; // number of buffers in buffer group |
| 4006 | const buffer_size: usize = 4; // size of each buffer in group | 4134 | const buffer_size: usize = 4; // size of each buffer in group |
| 4007 | const buffers = try testing.allocator.alloc(u8, buffers_count * buffer_size); | | |
| 4008 | defer testing.allocator.free(buffers); | | |
| 4009 | var buf_grp = BufferGroup.init( | 4135 | var buf_grp = BufferGroup.init( |
| 4010 | &ring, | 4136 | &ring, |
| | 4137 | testing.allocator, |
| 4011 | group_id, | 4138 | group_id, |
| 4012 | buffers, | | |
| 4013 | buffer_size, | 4139 | buffer_size, |
| 4014 | buffers_count, | 4140 | buffers_count, |
| 4015 | ) catch |err| switch (err) { | 4141 | ) catch |err| switch (err) { |
| ... | @@ -4017,7 +4143,7 @@ test "ring mapped buffers recv" { | ... | @@ -4017,7 +4143,7 @@ test "ring mapped buffers recv" { |
| 4017 | error.ArgumentsInvalid => return error.SkipZigTest, | 4143 | error.ArgumentsInvalid => return error.SkipZigTest, |
| 4018 | else => return err, | 4144 | else => return err, |
| 4019 | }; | 4145 | }; |
| 4020 | defer buf_grp.deinit(); | 4146 | defer buf_grp.deinit(testing.allocator); |
| 4021 | | 4147 | |
| 4022 | // create client/server fds | 4148 | // create client/server fds |
| 4023 | const fds = try createSocketTestHarness(&ring); | 4149 | const fds = try createSocketTestHarness(&ring); |
| ... | @@ -4039,14 +4165,18 @@ test "ring mapped buffers recv" { | ... | @@ -4039,14 +4165,18 @@ test "ring mapped buffers recv" { |
| 4039 | if (cqe_send.err() == .INVAL) return error.SkipZigTest; | 4165 | if (cqe_send.err() == .INVAL) return error.SkipZigTest; |
| 4040 | try testing.expectEqual(linux.io_uring_cqe{ .user_data = user_data, .res = data.len, .flags = 0 }, cqe_send); | 4166 | try testing.expectEqual(linux.io_uring_cqe{ .user_data = user_data, .res = data.len, .flags = 0 }, cqe_send); |
| 4041 | } | 4167 | } |
| 4042 | | 4168 | var pos: usize = 0; |
| 4043 | // server reads data into provided buffers | 4169 | |
| 4044 | // there are 2 buffers of size 4, so each read gets only chunk of data | 4170 | // read first chunk |
| 4045 | // we read four chunks of 4, 4, 4, 3 bytes each | 4171 | const cqe1 = try buf_grp_recv_submit_get_cqe(&ring, &buf_grp, fds.server, rnd.int(u64)); |
| 4046 | var chunk: []const u8 = data[0..buffer_size]; // first chunk | 4172 | var buf = try buf_grp.get(cqe1); |
| 4047 | const id1 = try expect_buf_grp_recv(&ring, &buf_grp, fds.server, rnd.int(u64), chunk); | 4173 | try testing.expectEqualSlices(u8, data[pos..][0..buf.len], buf); |
| 4048 | chunk = data[buffer_size .. buffer_size * 2]; // second chunk | 4174 | pos += buf.len; |
| 4049 | const id2 = try expect_buf_grp_recv(&ring, &buf_grp, fds.server, rnd.int(u64), chunk); | 4175 | // second chunk |
| | 4176 | const cqe2 = try buf_grp_recv_submit_get_cqe(&ring, &buf_grp, fds.server, rnd.int(u64)); |
| | 4177 | buf = try buf_grp.get(cqe2); |
| | 4178 | try testing.expectEqualSlices(u8, data[pos..][0..buf.len], buf); |
| | 4179 | pos += buf.len; |
| 4050 | | 4180 | |
| 4051 | // both buffers provided to the kernel are used so we get error | 4181 | // both buffers provided to the kernel are used so we get error |
| 4052 | // 'no more buffers', until we put buffers to the kernel | 4182 | // 'no more buffers', until we put buffers to the kernel |
| ... | @@ -4063,16 +4193,17 @@ test "ring mapped buffers recv" { | ... | @@ -4063,16 +4193,17 @@ test "ring mapped buffers recv" { |
| 4063 | } | 4193 | } |
| 4064 | | 4194 | |
| 4065 | // put buffers back to the kernel | 4195 | // put buffers back to the kernel |
| 4066 | buf_grp.put(id1); | 4196 | try buf_grp.put(cqe1); |
| 4067 | buf_grp.put(id2); | 4197 | try buf_grp.put(cqe2); |
| 4068 | | 4198 | |
| 4069 | chunk = data[buffer_size * 2 .. buffer_size * 3]; // third chunk | 4199 | // read remaining data |
| 4070 | const id3 = try expect_buf_grp_recv(&ring, &buf_grp, fds.server, rnd.int(u64), chunk); | 4200 | while (pos < data.len) { |
| 4071 | buf_grp.put(id3); | 4201 | const cqe = try buf_grp_recv_submit_get_cqe(&ring, &buf_grp, fds.server, rnd.int(u64)); |
| 4072 | | 4202 | buf = try buf_grp.get(cqe); |
| 4073 | chunk = data[buffer_size * 3 ..]; // last chunk | 4203 | try testing.expectEqualSlices(u8, data[pos..][0..buf.len], buf); |
| 4074 | const id4 = try expect_buf_grp_recv(&ring, &buf_grp, fds.server, rnd.int(u64), chunk); | 4204 | pos += buf.len; |
| 4075 | buf_grp.put(id4); | 4205 | try buf_grp.put(cqe); |
| | 4206 | } |
| 4076 | } | 4207 | } |
| 4077 | } | 4208 | } |
| 4078 | | 4209 | |
| ... | @@ -4090,12 +4221,10 @@ test "ring mapped buffers multishot recv" { | ... | @@ -4090,12 +4221,10 @@ test "ring mapped buffers multishot recv" { |
| 4090 | const group_id: u16 = 1; // buffers group id | 4221 | const group_id: u16 = 1; // buffers group id |
| 4091 | const buffers_count: u16 = 2; // number of buffers in buffer group | 4222 | const buffers_count: u16 = 2; // number of buffers in buffer group |
| 4092 | const buffer_size: usize = 4; // size of each buffer in group | 4223 | const buffer_size: usize = 4; // size of each buffer in group |
| 4093 | const buffers = try testing.allocator.alloc(u8, buffers_count * buffer_size); | | |
| 4094 | defer testing.allocator.free(buffers); | | |
| 4095 | var buf_grp = BufferGroup.init( | 4224 | var buf_grp = BufferGroup.init( |
| 4096 | &ring, | 4225 | &ring, |
| | 4226 | testing.allocator, |
| 4097 | group_id, | 4227 | group_id, |
| 4098 | buffers, | | |
| 4099 | buffer_size, | 4228 | buffer_size, |
| 4100 | buffers_count, | 4229 | buffers_count, |
| 4101 | ) catch |err| switch (err) { | 4230 | ) catch |err| switch (err) { |
| ... | @@ -4103,7 +4232,7 @@ test "ring mapped buffers multishot recv" { | ... | @@ -4103,7 +4232,7 @@ test "ring mapped buffers multishot recv" { |
| 4103 | error.ArgumentsInvalid => return error.SkipZigTest, | 4232 | error.ArgumentsInvalid => return error.SkipZigTest, |
| 4104 | else => return err, | 4233 | else => return err, |
| 4105 | }; | 4234 | }; |
| 4106 | defer buf_grp.deinit(); | 4235 | defer buf_grp.deinit(testing.allocator); |
| 4107 | | 4236 | |
| 4108 | // create client/server fds | 4237 | // create client/server fds |
| 4109 | const fds = try createSocketTestHarness(&ring); | 4238 | const fds = try createSocketTestHarness(&ring); |
| ... | @@ -4116,7 +4245,7 @@ test "ring mapped buffers multishot recv" { | ... | @@ -4116,7 +4245,7 @@ test "ring mapped buffers multishot recv" { |
| 4116 | var round: usize = 4; // repeat send/recv cycle round times | 4245 | var round: usize = 4; // repeat send/recv cycle round times |
| 4117 | while (round > 0) : (round -= 1) { | 4246 | while (round > 0) : (round -= 1) { |
| 4118 | // client sends data | 4247 | // client sends data |
| 4119 | const data = [_]u8{ 0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 0xa, 0xb, 0xc, 0xd, 0xe }; | 4248 | const data = [_]u8{ 0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 0xa, 0xb, 0xc, 0xd, 0xe, 0xf }; |
| 4120 | { | 4249 | { |
| 4121 | const user_data = rnd.int(u64); | 4250 | const user_data = rnd.int(u64); |
| 4122 | _ = try ring.send(user_data, fds.client, data[0..], 0); | 4251 | _ = try ring.send(user_data, fds.client, data[0..], 0); |
| ... | @@ -4133,7 +4262,7 @@ test "ring mapped buffers multishot recv" { | ... | @@ -4133,7 +4262,7 @@ test "ring mapped buffers multishot recv" { |
| 4133 | | 4262 | |
| 4134 | // server reads data into provided buffers | 4263 | // server reads data into provided buffers |
| 4135 | // there are 2 buffers of size 4, so each read gets only chunk of data | 4264 | // there are 2 buffers of size 4, so each read gets only chunk of data |
| 4136 | // we read four chunks of 4, 4, 4, 3 bytes each | 4265 | // we read four chunks of 4, 4, 4, 4 bytes each |
| 4137 | var chunk: []const u8 = data[0..buffer_size]; // first chunk | 4266 | var chunk: []const u8 = data[0..buffer_size]; // first chunk |
| 4138 | const cqe1 = try expect_buf_grp_cqe(&ring, &buf_grp, recv_user_data, chunk); | 4267 | const cqe1 = try expect_buf_grp_cqe(&ring, &buf_grp, recv_user_data, chunk); |
| 4139 | try testing.expect(cqe1.flags & linux.IORING_CQE_F_MORE > 0); | 4268 | try testing.expect(cqe1.flags & linux.IORING_CQE_F_MORE > 0); |
| ... | @@ -4157,8 +4286,8 @@ test "ring mapped buffers multishot recv" { | ... | @@ -4157,8 +4286,8 @@ test "ring mapped buffers multishot recv" { |
| 4157 | } | 4286 | } |
| 4158 | | 4287 | |
| 4159 | // put buffers back to the kernel | 4288 | // put buffers back to the kernel |
| 4160 | buf_grp.put(try cqe1.buffer_id()); | 4289 | try buf_grp.put(cqe1); |
| 4161 | buf_grp.put(try cqe2.buffer_id()); | 4290 | try buf_grp.put(cqe2); |
| 4162 | | 4291 | |
| 4163 | // restart multishot | 4292 | // restart multishot |
| 4164 | recv_user_data = rnd.int(u64); | 4293 | recv_user_data = rnd.int(u64); |
| ... | @@ -4168,12 +4297,12 @@ test "ring mapped buffers multishot recv" { | ... | @@ -4168,12 +4297,12 @@ test "ring mapped buffers multishot recv" { |
| 4168 | chunk = data[buffer_size * 2 .. buffer_size * 3]; // third chunk | 4297 | chunk = data[buffer_size * 2 .. buffer_size * 3]; // third chunk |
| 4169 | const cqe3 = try expect_buf_grp_cqe(&ring, &buf_grp, recv_user_data, chunk); | 4298 | const cqe3 = try expect_buf_grp_cqe(&ring, &buf_grp, recv_user_data, chunk); |
| 4170 | try testing.expect(cqe3.flags & linux.IORING_CQE_F_MORE > 0); | 4299 | try testing.expect(cqe3.flags & linux.IORING_CQE_F_MORE > 0); |
| 4171 | buf_grp.put(try cqe3.buffer_id()); | 4300 | try buf_grp.put(cqe3); |
| 4172 | | 4301 | |
| 4173 | chunk = data[buffer_size * 3 ..]; // last chunk | 4302 | chunk = data[buffer_size * 3 ..]; // last chunk |
| 4174 | const cqe4 = try expect_buf_grp_cqe(&ring, &buf_grp, recv_user_data, chunk); | 4303 | const cqe4 = try expect_buf_grp_cqe(&ring, &buf_grp, recv_user_data, chunk); |
| 4175 | try testing.expect(cqe4.flags & linux.IORING_CQE_F_MORE > 0); | 4304 | try testing.expect(cqe4.flags & linux.IORING_CQE_F_MORE > 0); |
| 4176 | buf_grp.put(try cqe4.buffer_id()); | 4305 | try buf_grp.put(cqe4); |
| 4177 | | 4306 | |
| 4178 | // cancel pending multishot recv operation | 4307 | // cancel pending multishot recv operation |
| 4179 | { | 4308 | { |
| ... | @@ -4217,23 +4346,26 @@ test "ring mapped buffers multishot recv" { | ... | @@ -4217,23 +4346,26 @@ test "ring mapped buffers multishot recv" { |
| 4217 | } | 4346 | } |
| 4218 | } | 4347 | } |
| 4219 | | 4348 | |
| 4220 | // Prepare and submit recv using buffer group. | 4349 | // Prepare, submit recv and get cqe using buffer group. |
| 4221 | // Test that buffer from group, pointed by cqe, matches expected. | 4350 | fn buf_grp_recv_submit_get_cqe( |
| 4222 | fn expect_buf_grp_recv( | | |
| 4223 | ring: *IoUring, | 4351 | ring: *IoUring, |
| 4224 | buf_grp: *BufferGroup, | 4352 | buf_grp: *BufferGroup, |
| 4225 | fd: posix.fd_t, | 4353 | fd: posix.fd_t, |
| 4226 | user_data: u64, | 4354 | user_data: u64, |
| 4227 | expected: []const u8, | 4355 | ) !linux.io_uring_cqe { |
| 4228 | ) !u16 { | 4356 | // prepare and submit recv |
| 4229 | // prepare and submit read | | |
| 4230 | const sqe = try buf_grp.recv(user_data, fd, 0); | 4357 | const sqe = try buf_grp.recv(user_data, fd, 0); |
| 4231 | try testing.expect(sqe.flags & linux.IOSQE_BUFFER_SELECT == linux.IOSQE_BUFFER_SELECT); | 4358 | try testing.expect(sqe.flags & linux.IOSQE_BUFFER_SELECT == linux.IOSQE_BUFFER_SELECT); |
| 4232 | try testing.expect(sqe.buf_index == buf_grp.group_id); | 4359 | try testing.expect(sqe.buf_index == buf_grp.group_id); |
| 4233 | try testing.expectEqual(@as(u32, 1), try ring.submit()); // submit | 4360 | try testing.expectEqual(@as(u32, 1), try ring.submit()); // submit |
| | 4361 | // get cqe, expect success |
| | 4362 | const cqe = try ring.copy_cqe(); |
| | 4363 | try testing.expectEqual(user_data, cqe.user_data); |
| | 4364 | try testing.expect(cqe.res >= 0); // success |
| | 4365 | try testing.expectEqual(posix.E.SUCCESS, cqe.err()); |
| | 4366 | try testing.expect(cqe.flags & linux.IORING_CQE_F_BUFFER == linux.IORING_CQE_F_BUFFER); // IORING_CQE_F_BUFFER flag is set |
| 4234 | | 4367 | |
| 4235 | const cqe = try expect_buf_grp_cqe(ring, buf_grp, user_data, expected); | 4368 | return cqe; |
| 4236 | return try cqe.buffer_id(); | | |
| 4237 | } | 4369 | } |
| 4238 | | 4370 | |
| 4239 | fn expect_buf_grp_cqe( | 4371 | fn expect_buf_grp_cqe( |
| ... | @@ -4253,7 +4385,7 @@ fn expect_buf_grp_cqe( | ... | @@ -4253,7 +4385,7 @@ fn expect_buf_grp_cqe( |
| 4253 | // get buffer from pool | 4385 | // get buffer from pool |
| 4254 | const buffer_id = try cqe.buffer_id(); | 4386 | const buffer_id = try cqe.buffer_id(); |
| 4255 | const len = @as(usize, @intCast(cqe.res)); | 4387 | const len = @as(usize, @intCast(cqe.res)); |
| 4256 | const buf = buf_grp.get(buffer_id)[0..len]; | 4388 | const buf = buf_grp.get_by_id(buffer_id)[0..len]; |
| 4257 | try testing.expectEqualSlices(u8, expected, buf); | 4389 | try testing.expectEqualSlices(u8, expected, buf); |
| 4258 | | 4390 | |
| 4259 | return cqe; | 4391 | return cqe; |
| ... | @@ -4305,3 +4437,137 @@ test "copy_cqes with wrapping sq.cqes buffer" { | ... | @@ -4305,3 +4437,137 @@ test "copy_cqes with wrapping sq.cqes buffer" { |
| 4305 | try testing.expectEqual(2 + 4 * i, ring.cq.head.*); | 4437 | try testing.expectEqual(2 + 4 * i, ring.cq.head.*); |
| 4306 | } | 4438 | } |
| 4307 | } | 4439 | } |
| | 4440 | |
| | 4441 | test "bind/listen/connect" { |
| | 4442 | var ring = IoUring.init(4, 0) catch |err| switch (err) { |
| | 4443 | error.SystemOutdated => return error.SkipZigTest, |
| | 4444 | error.PermissionDenied => return error.SkipZigTest, |
| | 4445 | else => return err, |
| | 4446 | }; |
| | 4447 | defer ring.deinit(); |
| | 4448 | |
| | 4449 | const probe = ring.get_probe() catch return error.SkipZigTest; |
| | 4450 | // LISTEN is higher required operation |
| | 4451 | if (!probe.is_supported(.LISTEN)) return error.SkipZigTest; |
| | 4452 | |
| | 4453 | var addr = net.Address.initIp4([4]u8{ 127, 0, 0, 1 }, 0); |
| | 4454 | const proto: u32 = if (addr.any.family == linux.AF.UNIX) 0 else linux.IPPROTO.TCP; |
| | 4455 | |
| | 4456 | const listen_fd = brk: { |
| | 4457 | // Create socket |
| | 4458 | _ = try ring.socket(1, addr.any.family, linux.SOCK.STREAM | linux.SOCK.CLOEXEC, proto, 0); |
| | 4459 | try testing.expectEqual(1, try ring.submit()); |
| | 4460 | var cqe = try ring.copy_cqe(); |
| | 4461 | try testing.expectEqual(1, cqe.user_data); |
| | 4462 | try testing.expectEqual(posix.E.SUCCESS, cqe.err()); |
| | 4463 | const listen_fd: posix.fd_t = @intCast(cqe.res); |
| | 4464 | try testing.expect(listen_fd > 2); |
| | 4465 | |
| | 4466 | // Prepare: set socket option * 2, bind, listen |
| | 4467 | var optval: u32 = 1; |
| | 4468 | (try ring.setsockopt(2, listen_fd, linux.SOL.SOCKET, linux.SO.REUSEADDR, mem.asBytes(&optval))).link_next(); |
| | 4469 | (try ring.setsockopt(3, listen_fd, linux.SOL.SOCKET, linux.SO.REUSEPORT, mem.asBytes(&optval))).link_next(); |
| | 4470 | (try ring.bind(4, listen_fd, &addr.any, addr.getOsSockLen(), 0)).link_next(); |
| | 4471 | _ = try ring.listen(5, listen_fd, 1, 0); |
| | 4472 | // Submit 4 operations |
| | 4473 | try testing.expectEqual(4, try ring.submit()); |
| | 4474 | // Expect all to succeed |
| | 4475 | for (2..6) |user_data| { |
| | 4476 | cqe = try ring.copy_cqe(); |
| | 4477 | try testing.expectEqual(user_data, cqe.user_data); |
| | 4478 | try testing.expectEqual(posix.E.SUCCESS, cqe.err()); |
| | 4479 | } |
| | 4480 | |
| | 4481 | // Check that socket option is set |
| | 4482 | optval = 0; |
| | 4483 | _ = try ring.getsockopt(5, listen_fd, linux.SOL.SOCKET, linux.SO.REUSEADDR, mem.asBytes(&optval)); |
| | 4484 | try testing.expectEqual(1, try ring.submit()); |
| | 4485 | cqe = try ring.copy_cqe(); |
| | 4486 | try testing.expectEqual(5, cqe.user_data); |
| | 4487 | try testing.expectEqual(posix.E.SUCCESS, cqe.err()); |
| | 4488 | try testing.expectEqual(1, optval); |
| | 4489 | |
| | 4490 | // Read system assigned port into addr |
| | 4491 | var addr_len: posix.socklen_t = addr.getOsSockLen(); |
| | 4492 | try posix.getsockname(listen_fd, &addr.any, &addr_len); |
| | 4493 | |
| | 4494 | break :brk listen_fd; |
| | 4495 | }; |
| | 4496 | |
| | 4497 | const connect_fd = brk: { |
| | 4498 | // Create connect socket |
| | 4499 | _ = try ring.socket(6, addr.any.family, linux.SOCK.STREAM | linux.SOCK.CLOEXEC, proto, 0); |
| | 4500 | try testing.expectEqual(1, try ring.submit()); |
| | 4501 | const cqe = try ring.copy_cqe(); |
| | 4502 | try testing.expectEqual(6, cqe.user_data); |
| | 4503 | try testing.expectEqual(posix.E.SUCCESS, cqe.err()); |
| | 4504 | // Get connect socket fd |
| | 4505 | const connect_fd: posix.fd_t = @intCast(cqe.res); |
| | 4506 | try testing.expect(connect_fd > 2 and connect_fd != listen_fd); |
| | 4507 | break :brk connect_fd; |
| | 4508 | }; |
| | 4509 | |
| | 4510 | // Prepare accept/connect operations |
| | 4511 | _ = try ring.accept(7, listen_fd, null, null, 0); |
| | 4512 | _ = try ring.connect(8, connect_fd, &addr.any, addr.getOsSockLen()); |
| | 4513 | try testing.expectEqual(2, try ring.submit()); |
| | 4514 | // Get listener accepted socket |
| | 4515 | var accept_fd: posix.socket_t = 0; |
| | 4516 | for (0..2) |_| { |
| | 4517 | const cqe = try ring.copy_cqe(); |
| | 4518 | try testing.expectEqual(posix.E.SUCCESS, cqe.err()); |
| | 4519 | if (cqe.user_data == 7) { |
| | 4520 | accept_fd = @intCast(cqe.res); |
| | 4521 | } else { |
| | 4522 | try testing.expectEqual(8, cqe.user_data); |
| | 4523 | } |
| | 4524 | } |
| | 4525 | try testing.expect(accept_fd > 2 and accept_fd != listen_fd and accept_fd != connect_fd); |
| | 4526 | |
| | 4527 | // Communicate |
| | 4528 | try testSendRecv(&ring, connect_fd, accept_fd); |
| | 4529 | try testSendRecv(&ring, accept_fd, connect_fd); |
| | 4530 | |
| | 4531 | // Shutdown and close all sockets |
| | 4532 | for ([_]posix.socket_t{ connect_fd, accept_fd, listen_fd }) |fd| { |
| | 4533 | (try ring.shutdown(9, fd, posix.SHUT.RDWR)).link_next(); |
| | 4534 | _ = try ring.close(10, fd); |
| | 4535 | try testing.expectEqual(2, try ring.submit()); |
| | 4536 | for (0..2) |i| { |
| | 4537 | const cqe = try ring.copy_cqe(); |
| | 4538 | try testing.expectEqual(posix.E.SUCCESS, cqe.err()); |
| | 4539 | try testing.expectEqual(9 + i, cqe.user_data); |
| | 4540 | } |
| | 4541 | } |
| | 4542 | } |
| | 4543 | |
| | 4544 | fn testSendRecv(ring: *IoUring, send_fd: posix.socket_t, recv_fd: posix.socket_t) !void { |
| | 4545 | const buffer_send = "0123456789abcdf" ** 10; |
| | 4546 | var buffer_recv: [buffer_send.len * 2]u8 = undefined; |
| | 4547 | |
| | 4548 | // 2 sends |
| | 4549 | _ = try ring.send(1, send_fd, buffer_send, linux.MSG.WAITALL); |
| | 4550 | _ = try ring.send(2, send_fd, buffer_send, linux.MSG.WAITALL); |
| | 4551 | try testing.expectEqual(2, try ring.submit()); |
| | 4552 | for (0..2) |i| { |
| | 4553 | const cqe = try ring.copy_cqe(); |
| | 4554 | try testing.expectEqual(1 + i, cqe.user_data); |
| | 4555 | try testing.expectEqual(posix.E.SUCCESS, cqe.err()); |
| | 4556 | try testing.expectEqual(buffer_send.len, @as(usize, @intCast(cqe.res))); |
| | 4557 | } |
| | 4558 | |
| | 4559 | // receive |
| | 4560 | var recv_len: usize = 0; |
| | 4561 | while (recv_len < buffer_send.len * 2) { |
| | 4562 | _ = try ring.recv(3, recv_fd, .{ .buffer = buffer_recv[recv_len..] }, 0); |
| | 4563 | try testing.expectEqual(1, try ring.submit()); |
| | 4564 | const cqe = try ring.copy_cqe(); |
| | 4565 | try testing.expectEqual(3, cqe.user_data); |
| | 4566 | try testing.expectEqual(posix.E.SUCCESS, cqe.err()); |
| | 4567 | recv_len += @intCast(cqe.res); |
| | 4568 | } |
| | 4569 | |
| | 4570 | // inspect recv buffer |
| | 4571 | try testing.expectEqualSlices(u8, buffer_send, buffer_recv[0..buffer_send.len]); |
| | 4572 | try testing.expectEqualSlices(u8, buffer_send, buffer_recv[buffer_send.len..]); |
| | 4573 | } |