authorgravatar for igor.anic@gmail.comIgor Anić <igor.anic@gmail.com> 2025-03-03 14:37:52+01:00
committergravatar for igor.anic@gmail.comIgor Anić <igor.anic@gmail.com> 2025-03-05 13:35:52+01:00
logc133171567fe3a81f817d0ea159bd9229d75291c
treed6565515345ccf2fe738744af06bf5cbf5cc6a91
parent4df039d235d5f77830fdc30a4c23121f6216364a

io_uring: incremental provided buffer consumption

[Incremental provided buffer consumption](https://github.com/axboe/liburing/wiki/What's-new-with-io_uring-in-6.11-and-6.12#incremental-provided-buffer-consumption) support is added in kernel 6.12. IoUring.BufferGroup will now use incremental consumption whenever kernel supports it. Before, provided buffers are wholly consumed when picked. Each cqe points to the different buffer. With this, cqe points to the part of the buffer. Multiple cqe's can reuse same buffer. Appropriate sizing of buffers becomes less important. There are slight changes in BufferGroup interface (it now needs to track current receive point for each buffer). Init requires allocator instead of buffers slice, it will allocate buffers slice and head pointers slice. Get and put now requires cqe becasue there we have information will the buffer be reused.

2 files changed, 106 insertions(+), 84 deletions(-)

lib/std/os/linux.zig+8-1
......@@ -5933,6 +5933,8 @@ pub const IORING_CQE_F_MORE = 1 << 1;
59335933pub const IORING_CQE_F_SOCK_NONEMPTY = 1 << 2;
59345934/// Set for notification CQEs. Can be used to distinct them from sends.
59355935pub const IORING_CQE_F_NOTIF = 1 << 3;
5936/// If set, the buffer ID set in the completion will get more completions.
5937pub const IORING_CQE_F_BUF_MORE = 1 << 4;
59365938
59375939pub const IORING_CQE_BUFFER_SHIFT = 16;
59385940
......@@ -6222,8 +6224,13 @@ pub const io_uring_buf_reg = extern struct {
62226224 ring_addr: u64,
62236225 ring_entries: u32,
62246226 bgid: u16,
6225 pad: u16,
6227 flags: u16,
62266228 resv: [3]u64,
6229
6230 pub const FLAG = struct {
6231 // Incremental buffer consummation.
6232 pub const INC: u16 = 2;
6233 };
62276234};
62286235
62296236pub const io_uring_getevents_arg = extern struct {
lib/std/os/linux/IoUring.zig+98-83
......@@ -1594,28 +1594,34 @@ pub const BufferGroup = struct {
15941594 buffers: []u8,
15951595 /// Size of each buffer in buffers.
15961596 buffer_size: u32,
1597 // 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.
15981598 buffers_count: u16,
1599 /// Head of unconsumed part of each buffer, if incremental consumption is enabled
1600 heads: []u32,
15991601 /// ID of this group, must be unique in ring.
16001602 group_id: u16,
16011603
16021604 pub fn init(
16031605 ring: *IoUring,
1606 allocator: mem.Allocator,
16041607 group_id: u16,
1605 buffers: []u8,
16061608 buffer_size: u32,
16071609 buffers_count: u16,
16081610 ) !BufferGroup {
1609 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);
16101615
1611 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, linux.io_uring_buf_reg.FLAG.INC);
16121617 buf_ring_init(br);
16131618
16141619 const mask = buf_ring_mask(buffers_count);
16151620 var i: u16 = 0;
16161621 while (i < buffers_count) : (i += 1) {
1617 const start = buffer_size * i;
1618 const buf = buffers[start .. start + buffer_size];
1622 const pos = buffer_size * i;
1623 const buf = buffers[pos .. pos + buffer_size];
1624 heads[i] = 0;
16191625 buf_ring_add(br, buf, i, mask, i);
16201626 }
16211627 buf_ring_advance(br, buffers_count);
......@@ -1625,11 +1631,18 @@ pub const BufferGroup = struct {
16251631 .group_id = group_id,
16261632 .br = br,
16271633 .buffers = buffers,
1634 .heads = heads,
16281635 .buffer_size = buffer_size,
16291636 .buffers_count = buffers_count,
16301637 };
16311638 }
16321639
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
16331646 // Prepare recv operation which will select buffer from this group.
16341647 pub fn recv(self: *BufferGroup, user_data: u64, fd: posix.fd_t, flags: u32) !*linux.io_uring_sqe {
16351648 var sqe = try self.ring.get_sqe();
......@@ -1649,33 +1662,34 @@ pub const BufferGroup = struct {
16491662 }
16501663
16511664 // Get buffer by id.
1652 pub fn get(self: *BufferGroup, buffer_id: u16) []u8 {
1653 const head = self.buffer_size * buffer_id;
1654 return self.buffers[head .. head + self.buffer_size];
1665 fn get_by_id(self: *BufferGroup, buffer_id: u16) []u8 {
1666 const pos = self.buffer_size * buffer_id;
1667 return self.buffers[pos .. pos + self.buffer_size][self.heads[buffer_id]..];
16551668 }
16561669
16571670 // Get buffer by CQE.
1658 pub fn get_cqe(self: *BufferGroup, cqe: linux.io_uring_cqe) ![]u8 {
1671 pub fn get(self: *BufferGroup, cqe: linux.io_uring_cqe) ![]u8 {
16591672 const buffer_id = try cqe.buffer_id();
16601673 const used_len = @as(usize, @intCast(cqe.res));
1661 return self.get(buffer_id)[0..used_len];
1662 }
1663
1664 // Release buffer to the kernel.
1665 pub fn put(self: *BufferGroup, buffer_id: u16) void {
1666 const mask = buf_ring_mask(self.buffers_count);
1667 const buffer = self.get(buffer_id);
1668 buf_ring_add(self.br, buffer, buffer_id, mask, 0);
1669 buf_ring_advance(self.br, 1);
1674 return self.get_by_id(buffer_id)[0..used_len];
16701675 }
16711676
16721677 // Release buffer from CQE to the kernel.
1673 pub fn put_cqe(self: *BufferGroup, cqe: linux.io_uring_cqe) !void {
1674 self.put(try cqe.buffer_id());
1675 }
1678 pub fn put(self: *BufferGroup, cqe: linux.io_uring_cqe) !void {
1679 const buffer_id = try cqe.buffer_id();
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;
16761688
1677 pub fn deinit(self: *BufferGroup) void {
1678 free_buf_ring(self.ring.fd, self.br, self.buffers_count, self.group_id);
1689 // Release buffer to the kernel. const mask = buf_ring_mask(self.buffers_count);
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);
16791693 }
16801694};
16811695
......@@ -1684,7 +1698,7 @@ pub const BufferGroup = struct {
16841698/// `fd` is IO_Uring.fd for which the provided buffer ring is being registered.
16851699/// `entries` is the number of entries requested in the buffer ring, must be power of 2.
16861700/// `group_id` is the chosen buffer group ID, unique in IO_Uring.
1687pub fn setup_buf_ring(fd: posix.fd_t, entries: u16, group_id: u16) !*align(page_size_min) linux.io_uring_buf_ring {
1701pub fn setup_buf_ring(fd: posix.fd_t, entries: u16, group_id: u16, flags: u16) !*align(page_size_min) linux.io_uring_buf_ring {
16881702 if (entries == 0 or entries > 1 << 15) return error.EntriesNotInRange;
16891703 if (!std.math.isPowerOfTwo(entries)) return error.EntriesNotPowerOfTwo;
16901704
......@@ -1701,22 +1715,24 @@ pub fn setup_buf_ring(fd: posix.fd_t, entries: u16, group_id: u16) !*align(page_
17011715 assert(mmap.len == mmap_size);
17021716
17031717 const br: *align(page_size_min) linux.io_uring_buf_ring = @ptrCast(mmap.ptr);
1704 try register_buf_ring(fd, @intFromPtr(br), entries, group_id);
1718 try register_buf_ring(fd, @intFromPtr(br), entries, group_id, flags);
17051719 return br;
17061720}
17071721
1708fn register_buf_ring(fd: posix.fd_t, addr: u64, entries: u32, group_id: u16) !void {
1722fn register_buf_ring(fd: posix.fd_t, addr: u64, entries: u32, group_id: u16, flags: u16) !void {
17091723 var reg = mem.zeroInit(linux.io_uring_buf_reg, .{
17101724 .ring_addr = addr,
17111725 .ring_entries = entries,
17121726 .bgid = group_id,
1727 .flags = flags,
17131728 });
1714 const res = linux.io_uring_register(
1715 fd,
1716 .REGISTER_PBUF_RING,
1717 @as(*const anyopaque, @ptrCast(&reg)),
1718 1,
1719 );
1729 var res = linux.io_uring_register(fd, .REGISTER_PBUF_RING, @as(*const anyopaque, @ptrCast(&reg)), 1);
1730 if (linux.E.init(res) == .INVAL and reg.flags & linux.io_uring_buf_reg.FLAG.INC > 0) {
1731 // Retry without incremental buffer consumption.
1732 // It is available since kernel 6.12. returns INVAL on older.
1733 reg.flags &= ~linux.io_uring_buf_reg.FLAG.INC;
1734 res = linux.io_uring_register(fd, .REGISTER_PBUF_RING, @as(*const anyopaque, @ptrCast(&reg)), 1);
1735 }
17201736 try handle_register_buf_ring_result(res);
17211737}
17221738
......@@ -4041,12 +4057,10 @@ test BufferGroup {
40414057 const group_id: u16 = 1; // buffers group id
40424058 const buffers_count: u16 = 1; // number of buffers in buffer group
40434059 const buffer_size: usize = 128; // size of each buffer in group
4044 const buffers = try testing.allocator.alloc(u8, buffers_count * buffer_size);
4045 defer testing.allocator.free(buffers);
40464060 var buf_grp = BufferGroup.init(
40474061 &ring,
4062 testing.allocator,
40484063 group_id,
4049 buffers,
40504064 buffer_size,
40514065 buffers_count,
40524066 ) catch |err| switch (err) {
......@@ -4054,7 +4068,7 @@ test BufferGroup {
40544068 error.ArgumentsInvalid => return error.SkipZigTest,
40554069 else => return err,
40564070 };
4057 defer buf_grp.deinit();
4071 defer buf_grp.deinit(testing.allocator);
40584072
40594073 // Create client/server fds
40604074 const fds = try createSocketTestHarness(&ring);
......@@ -4085,14 +4099,11 @@ test BufferGroup {
40854099 try testing.expectEqual(posix.E.SUCCESS, cqe.err());
40864100 try testing.expectEqual(data.len, @as(usize, @intCast(cqe.res))); // cqe.res holds received data len
40874101
4088 // Read buffer_id and used buffer len from cqe
4089 const buffer_id = try cqe.buffer_id();
4090 const len: usize = @intCast(cqe.res);
40914102 // Get buffer from pool
4092 const buf = buf_grp.get(buffer_id)[0..len];
4103 const buf = try buf_grp.get(cqe);
40934104 try testing.expectEqualSlices(u8, &data, buf);
40944105 // Release buffer to the kernel when application is done with it
4095 buf_grp.put(buffer_id);
4106 try buf_grp.put(cqe);
40964107 }
40974108}
40984109
......@@ -4110,12 +4121,10 @@ test "ring mapped buffers recv" {
41104121 const group_id: u16 = 1; // buffers group id
41114122 const buffers_count: u16 = 2; // number of buffers in buffer group
41124123 const buffer_size: usize = 4; // size of each buffer in group
4113 const buffers = try testing.allocator.alloc(u8, buffers_count * buffer_size);
4114 defer testing.allocator.free(buffers);
41154124 var buf_grp = BufferGroup.init(
41164125 &ring,
4126 testing.allocator,
41174127 group_id,
4118 buffers,
41194128 buffer_size,
41204129 buffers_count,
41214130 ) catch |err| switch (err) {
......@@ -4123,7 +4132,7 @@ test "ring mapped buffers recv" {
41234132 error.ArgumentsInvalid => return error.SkipZigTest,
41244133 else => return err,
41254134 };
4126 defer buf_grp.deinit();
4135 defer buf_grp.deinit(testing.allocator);
41274136
41284137 // create client/server fds
41294138 const fds = try createSocketTestHarness(&ring);
......@@ -4145,14 +4154,18 @@ test "ring mapped buffers recv" {
41454154 if (cqe_send.err() == .INVAL) return error.SkipZigTest;
41464155 try testing.expectEqual(linux.io_uring_cqe{ .user_data = user_data, .res = data.len, .flags = 0 }, cqe_send);
41474156 }
4148
4149 // server reads data into provided buffers
4150 // there are 2 buffers of size 4, so each read gets only chunk of data
4151 // we read four chunks of 4, 4, 4, 3 bytes each
4152 var chunk: []const u8 = data[0..buffer_size]; // first chunk
4153 const id1 = try expect_buf_grp_recv(&ring, &buf_grp, fds.server, rnd.int(u64), chunk);
4154 chunk = data[buffer_size .. buffer_size * 2]; // second chunk
4155 const id2 = try expect_buf_grp_recv(&ring, &buf_grp, fds.server, rnd.int(u64), chunk);
4157 var pos: usize = 0;
4158
4159 // read first chunk
4160 const cqe1 = try buf_grp_recv_submit_get_cqe(&ring, &buf_grp, fds.server, rnd.int(u64));
4161 var buf = try buf_grp.get(cqe1);
4162 try testing.expectEqualSlices(u8, data[pos..][0..buf.len], buf);
4163 pos += buf.len;
4164 // second chunk
4165 const cqe2 = try buf_grp_recv_submit_get_cqe(&ring, &buf_grp, fds.server, rnd.int(u64));
4166 buf = try buf_grp.get(cqe2);
4167 try testing.expectEqualSlices(u8, data[pos..][0..buf.len], buf);
4168 pos += buf.len;
41564169
41574170 // both buffers provided to the kernel are used so we get error
41584171 // 'no more buffers', until we put buffers to the kernel
......@@ -4169,16 +4182,17 @@ test "ring mapped buffers recv" {
41694182 }
41704183
41714184 // put buffers back to the kernel
4172 buf_grp.put(id1);
4173 buf_grp.put(id2);
4174
4175 chunk = data[buffer_size * 2 .. buffer_size * 3]; // third chunk
4176 const id3 = try expect_buf_grp_recv(&ring, &buf_grp, fds.server, rnd.int(u64), chunk);
4177 buf_grp.put(id3);
4178
4179 chunk = data[buffer_size * 3 ..]; // last chunk
4180 const id4 = try expect_buf_grp_recv(&ring, &buf_grp, fds.server, rnd.int(u64), chunk);
4181 buf_grp.put(id4);
4185 try buf_grp.put(cqe1);
4186 try buf_grp.put(cqe2);
4187
4188 // read remaining data
4189 while (pos < data.len) {
4190 const cqe = try buf_grp_recv_submit_get_cqe(&ring, &buf_grp, fds.server, rnd.int(u64));
4191 buf = try buf_grp.get(cqe);
4192 try testing.expectEqualSlices(u8, data[pos..][0..buf.len], buf);
4193 pos += buf.len;
4194 try buf_grp.put(cqe);
4195 }
41824196 }
41834197}
41844198
......@@ -4196,12 +4210,10 @@ test "ring mapped buffers multishot recv" {
41964210 const group_id: u16 = 1; // buffers group id
41974211 const buffers_count: u16 = 2; // number of buffers in buffer group
41984212 const buffer_size: usize = 4; // size of each buffer in group
4199 const buffers = try testing.allocator.alloc(u8, buffers_count * buffer_size);
4200 defer testing.allocator.free(buffers);
42014213 var buf_grp = BufferGroup.init(
42024214 &ring,
4215 testing.allocator,
42034216 group_id,
4204 buffers,
42054217 buffer_size,
42064218 buffers_count,
42074219 ) catch |err| switch (err) {
......@@ -4209,7 +4221,7 @@ test "ring mapped buffers multishot recv" {
42094221 error.ArgumentsInvalid => return error.SkipZigTest,
42104222 else => return err,
42114223 };
4212 defer buf_grp.deinit();
4224 defer buf_grp.deinit(testing.allocator);
42134225
42144226 // create client/server fds
42154227 const fds = try createSocketTestHarness(&ring);
......@@ -4222,7 +4234,7 @@ test "ring mapped buffers multishot recv" {
42224234 var round: usize = 4; // repeat send/recv cycle round times
42234235 while (round > 0) : (round -= 1) {
42244236 // client sends data
4225 const data = [_]u8{ 0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 0xa, 0xb, 0xc, 0xd, 0xe };
4237 const data = [_]u8{ 0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 0xa, 0xb, 0xc, 0xd, 0xe, 0xf };
42264238 {
42274239 const user_data = rnd.int(u64);
42284240 _ = try ring.send(user_data, fds.client, data[0..], 0);
......@@ -4239,7 +4251,7 @@ test "ring mapped buffers multishot recv" {
42394251
42404252 // server reads data into provided buffers
42414253 // there are 2 buffers of size 4, so each read gets only chunk of data
4242 // we read four chunks of 4, 4, 4, 3 bytes each
4254 // we read four chunks of 4, 4, 4, 4 bytes each
42434255 var chunk: []const u8 = data[0..buffer_size]; // first chunk
42444256 const cqe1 = try expect_buf_grp_cqe(&ring, &buf_grp, recv_user_data, chunk);
42454257 try testing.expect(cqe1.flags & linux.IORING_CQE_F_MORE > 0);
......@@ -4263,8 +4275,8 @@ test "ring mapped buffers multishot recv" {
42634275 }
42644276
42654277 // put buffers back to the kernel
4266 buf_grp.put(try cqe1.buffer_id());
4267 buf_grp.put(try cqe2.buffer_id());
4278 try buf_grp.put(cqe1);
4279 try buf_grp.put(cqe2);
42684280
42694281 // restart multishot
42704282 recv_user_data = rnd.int(u64);
......@@ -4274,12 +4286,12 @@ test "ring mapped buffers multishot recv" {
42744286 chunk = data[buffer_size * 2 .. buffer_size * 3]; // third chunk
42754287 const cqe3 = try expect_buf_grp_cqe(&ring, &buf_grp, recv_user_data, chunk);
42764288 try testing.expect(cqe3.flags & linux.IORING_CQE_F_MORE > 0);
4277 buf_grp.put(try cqe3.buffer_id());
4289 try buf_grp.put(cqe3);
42784290
42794291 chunk = data[buffer_size * 3 ..]; // last chunk
42804292 const cqe4 = try expect_buf_grp_cqe(&ring, &buf_grp, recv_user_data, chunk);
42814293 try testing.expect(cqe4.flags & linux.IORING_CQE_F_MORE > 0);
4282 buf_grp.put(try cqe4.buffer_id());
4294 try buf_grp.put(cqe4);
42834295
42844296 // cancel pending multishot recv operation
42854297 {
......@@ -4323,23 +4335,26 @@ test "ring mapped buffers multishot recv" {
43234335 }
43244336}
43254337
4326// Prepare and submit recv using buffer group.
4327// Test that buffer from group, pointed by cqe, matches expected.
4328fn expect_buf_grp_recv(
4338// Prepare, submit recv and get cqe using buffer group.
4339fn buf_grp_recv_submit_get_cqe(
43294340 ring: *IoUring,
43304341 buf_grp: *BufferGroup,
43314342 fd: posix.fd_t,
43324343 user_data: u64,
4333 expected: []const u8,
4334) !u16 {
4335 // prepare and submit read
4344) !linux.io_uring_cqe {
4345 // prepare and submit recv
43364346 const sqe = try buf_grp.recv(user_data, fd, 0);
43374347 try testing.expect(sqe.flags & linux.IOSQE_BUFFER_SELECT == linux.IOSQE_BUFFER_SELECT);
43384348 try testing.expect(sqe.buf_index == buf_grp.group_id);
43394349 try testing.expectEqual(@as(u32, 1), try ring.submit()); // submit
4350 // get cqe, expect success
4351 const cqe = try ring.copy_cqe();
4352 try testing.expectEqual(user_data, cqe.user_data);
4353 try testing.expect(cqe.res >= 0); // success
4354 try testing.expectEqual(posix.E.SUCCESS, cqe.err());
4355 try testing.expect(cqe.flags & linux.IORING_CQE_F_BUFFER == linux.IORING_CQE_F_BUFFER); // IORING_CQE_F_BUFFER flag is set
43404356
4341 const cqe = try expect_buf_grp_cqe(ring, buf_grp, user_data, expected);
4342 return try cqe.buffer_id();
4357 return cqe;
43434358}
43444359
43454360fn expect_buf_grp_cqe(
......@@ -4359,7 +4374,7 @@ fn expect_buf_grp_cqe(
43594374 // get buffer from pool
43604375 const buffer_id = try cqe.buffer_id();
43614376 const len = @as(usize, @intCast(cqe.res));
4362 const buf = buf_grp.get(buffer_id)[0..len];
4377 const buf = buf_grp.get_by_id(buffer_id)[0..len];
43634378 try testing.expectEqualSlices(u8, expected, buf);
43644379
43654380 return cqe;