#!/usr/bin/env dub
/+ dub.sdl:
name "io_uring_multishot_recv"
dependency "during" version="~>0.5.0"
platforms "linux"
targetPath "build"
+/
/**
* `io_uring` — multishot RECV into a buffer ring (`IORING_RECV_MULTISHOT`, Linux 6.0).
*
* A plain `RECV` (5.6) consumes one SQE per received segment: you re-arm it for
* every read. Multishot RECV (6.0) flips that — a *single* armed SQE stays live
* and posts a fresh CQE for each incoming segment, each one selecting a buffer
* from a provided-buffer ring. The kernel keeps the operation armed (signalled by
* `CQEFlags.MORE` on every non-final CQE) so a server can drain a busy socket with
* one submission instead of N.
*
* This program builds directly on the 5.19 buffer-ring example:
* 1. registers a small buffer ring (`registerBufRing`) for group id `BGID`,
* 2. publishes several buffers into that ring,
* 3. arms ONE multishot RECV that selects from the group
* (`prepRecvMultishot(fd, gid, len)` sets `IOSQE_BUFFER_SELECT`, `buf_group`,
* and `IORING_RECV_MULTISHOT` for us),
* 4. writes TWO separate messages into the peer end of a socketpair,
* 5. waits and asserts it gets TWO CQEs from that one SQE, each carrying
* `CQEFlags.MORE` (still armed) + `CQEFlags.BUFFER` (a buffer was selected),
* landing in two *distinct* ring buffers with the right bytes.
*
* The MORE flag is the multishot contract: it stays set while the op is armed and
* clears on the terminal CQE. The selected buffer id is in the upper 16 bits of
* `cqe.flags` (`>> CQE_BUFFER_SHIFT`), exactly as for single-shot buffer-select.
*
* Companion to the io_uring chronology:
* see docs/research/async-io/io-uring/timeline.md
* § "6.0 — Zero-copy send, single-issuer, sync cancel (October 2022)".
*
* Run with: `dub run --single multishot-recv.d`
*
* Portability: prints `SKIP:` and exits 0 when io_uring is unavailable or the
* running kernel predates multishot RECV (the RECV CQE comes back -EINVAL on a
* pre-6.0 kernel, or registerBufRing -> -EINVAL/-EOPNOTSUPP/-ENOSYS pre-5.19).
* Exits nonzero only on a genuinely unexpected syscall failure.
*/
module (module) io_uring_multishot_recvio_uring — multishot RECV into a buffer ring (IORING_RECV_MULTISHOT, Linux 6.0).
A plain RECV (5.6) consumes one SQE per received segment: you re-arm it for
every read. Multishot RECV (6.0) flips that — a single armed SQE stays live
and posts a fresh CQE for each incoming segment, each one selecting a buffer
from a provided-buffer ring. The kernel keeps the operation armed (signalled by
CQEFlags.MORE on every non-final CQE) so a server can drain a busy socket with
one submission instead of N.
This program builds directly on the 5.19 buffer-ring example:
registers a small buffer ring (registerBufRing) for group id BGID,
publishes several buffers into that ring,
arms ONE multishot RECV that selects from the group
(prepRecvMultishot(fd, gid, len) sets IOSQE_BUFFER_SELECT, buf_group,
and IORING_RECV_MULTISHOT for us),
writes TWO separate messages into the peer end of a socketpair,
waits and asserts it gets TWO CQEs from that one SQE, each carrying
CQEFlags.MORE (still armed) + CQEFlags.BUFFER (a buffer was selected),
landing in two distinct ring buffers with the right bytes.
The MORE flag is the multishot contract: it stays set while the op is armed and
clears on the terminal CQE. The selected buffer id is in the upper 16 bits of
cqe.flags (>> CQE_BUFFER_SHIFT), exactly as for single-shot buffer-select.
Companion to the io_uring chronology:
see docs/research/async-io/io-uring/timeline.md
§ "6.0 — Zero-copy send, single-issuer, sync cancel (October 2022)".
Run with: dub run --single multishot-recv.d
Portability
prints SKIP: and exits 0 when io_uring is unavailable or the
running kernel predates multishot RECV (the RECV CQE comes back -EINVAL on a
pre-6.0 kernel, or registerBufRing -> -EINVAL/-EOPNOTSUPP/-ENOSYS pre-5.19).
Exits nonzero only on a genuinely unexpected syscall failure.
io_uring_multishot_recv;
import (module) duringSimple idiomatic dlang wrapper around linux io_uring
(see: https://kernel.dk/io_uring.pdf) asynchronous API.
during;
import (package) corecore.(package) core.stdcstdc.(module) core.stdc.errnoD header file for C99.
pubs.opengroup.org/onlinepubs/009695399/basedefs/errno.h.html, errno.h
Source
core/stdc/errno.d
errno : (alias constant) io_uring_multishot_recv.EINVAL = int core.stdc.errno.EINVAL = 22EINVAL, (alias constant) io_uring_multishot_recv.EOPNOTSUPP = int core.stdc.errno.EOPNOTSUPP = 95EOPNOTSUPP, (alias constant) io_uring_multishot_recv.ENOSYS = int core.stdc.errno.ENOSYS = 38ENOSYS;
import (package) corecore.(package) core.stdcstdc.(module) core.stdc.stdlibD header file for C99.
pubs.opengroup.org/onlinepubs/009695399/basedefs/stdlib.h.html, stdlib.h
Source
core/stdc/stdlib.d
stdlib : (alias) io_uring_multishot_recv.free = void core.stdc.stdlib.free(void* ptr) nothrow @nogcfree;
import (package) corecore.(package) core.syssys.(package) core.sys.posixposix.(module) core.sys.posix.stdlibD header file for POSIX.
stdlib : (alias) io_uring_multishot_recv.posix_memalign = int core.sys.posix.stdlib.posix_memalign(scope void**, ulong, ulong) pure nothrow @nogcposix_memalign;
import (package) corecore.(package) core.syssys.(package) core.sys.posixposix.(package) core.sys.posix.syssys.(module) core.sys.posix.sys.socketD header file for POSIX.
socket : (alias enum value) io_uring_multishot_recv.AF_UNIX = core.sys.posix.sys.socket.AF_UNIX = 1AF_UNIX, (alias enum value) io_uring_multishot_recv.SOCK_STREAM = core.sys.posix.sys.socket.SOCK_STREAM = 1SOCK_STREAM, (alias) io_uring_multishot_recv.socketpair = int core.sys.posix.sys.socket.socketpair(int, int, int, ref int[2]) nothrow @nogc @safesocketpair;
import (package) corecore.(package) core.syssys.(package) core.sys.posixposix.(module) core.sys.posix.unistdD header file for POSIX.
unistd : (alias) io_uring_multishot_recv.close = int core.sys.posix.unistd.close(int) nothrow @nogc @trustedclose, (alias) io_uring_multishot_recv.write = long core.sys.posix.unistd.write(int, scope const(void*), ulong) nothrow @nogcwrite;
import (package) stdstd.(module) std.stdioCategory Symbols File handles _popen File isFileHandle openNetwork stderr stdin stdout Reading chunks lines readf readfln readln Writing toFile write writef writefln writeln Misc KeepTerminator LockType StdioException
Standard I/O functions that extend core.stdc.stdio. core.stdc.stdio
is publically imported when importing std.stdio.
There are three layers of I/O:
The lowest layer is the operating system layer. The two main schemes are Windows and Posix.
C's stdio.h which unifies the two operating system schemes.
std.stdio, this module, unifies the various stdio.h implementations into
a high level package for D programs.
Source
std/stdio.d
stdio : stderr, (alias template) io_uring_multishot_recv.writefln = std.stdio.writefln(alias fmt, A...)(A args) if (isSomeString!(typeof(fmt)))Equivalent to writef(fmt, args, '\n').
writefln;
// Group id for our buffer ring, and the buffer geometry. RING_ENTRIES must be a
// power of two — the kernel masks the tail with `ring_entries - 1`.
enum ushort (constant) ushort io_uring_multishot_recv.BGID = cast(ushort)7uBGID = 7;
enum uint (constant) uint io_uring_multishot_recv.RING_ENTRIES = 8uRING_ENTRIES = 8;
enum uint (constant) uint io_uring_multishot_recv.BUF_SIZE = 64uBUF_SIZE = 64;
int int D main()main()
{
(struct) during.UringMain entry point to work with io_uring.
It hides SubmissionQueue and CompletionQueue behind standard range interface.
We put in SubmissionEntry entries and take out CompletionEntry entries.
Use predefined prepXX methods to fill required fields of SubmissionEntry before put or during putWith.
Note
prepXX functions doesn't touch previous entry state, just fills in operation properties. This is because for
less error prone interface it is cleared automatically when prepared using putWith. So when using on own SubmissionEntry
(outside submission queue), that would be added to the submission queue using put, be sure its cleared if it's
reused for multiple operations.
Uring (local variable) during.Uring ioio;
const (local variable) const(int) setupRetsetupRet = (local variable) during.Uring ioio.int during.setup(ref during.Uring uring, uint entries = 128u, during.io_uring.SetupFlags flags = SetupFlags.NONE) nothrow @nogc @safeSetup new instance of io_uring into provided Uring structure.
setup(8);
if ((local variable) const(int) setupRetsetupRet < 0)
{
void std.stdio.writefln!(char, const(int))(in char[] fmt, const(int) __param_1) @safeEquivalent to writef(fmt, args, '\n').
writefln("SKIP: io_uring_setup failed (errno %d) — io_uring unavailable on this host", -(local variable) const(int) setupRetsetupRet);
return 0;
}
// A connected pair of local sockets: we recv on sv[0] via io_uring and write
// the payloads into sv[1] with ordinary write(2). Loopback-only, no network.
int[2] (local variable) int[2] svsv;
if (int core.sys.posix.sys.socket.socketpair(int, int, int, ref int[2]) nothrow @nogc @safesocketpair((enum value) core.sys.posix.sys.socket.AF_UNIX = 1AF_UNIX, (enum value) core.sys.posix.sys.socket.SOCK_STREAM = 1SOCK_STREAM, 0, (local variable) int[2] svsv) != 0)
{
stderr.std.stdio.File std.stdio.makeGlobal!"core.stdc.stdio.stderr"() nothrow @nogc @property ref @systemwritefln("socketpair failed");
return 1;
}
scope (exit) { int core.sys.posix.unistd.close(int) nothrow @nogc @trustedclose((local variable) int[2] svsv[0]); int core.sys.posix.unistd.close(int) nothrow @nogc @trustedclose((local variable) int[2] svsv[1]); }
// Allocate the buffer ring page-aligned (the kernel requires page alignment for
// IORING_REGISTER_PBUF_RING). It is a flat array of RING_ENTRIES `io_uring_buf`
// slots; slot 0's resv field doubles as the ring's producer tail.
enum (alias) object.size_t = ulongsize_t (constant) ulong io_uring_multishot_recv.main.ringBytes = 128LUringBytes = (struct) during.io_uring.io_uring_bufio_uring_buf.(constant) ulong during.io_uring.io_uring_buf.sizeof = 16LUsizeof * (constant) uint io_uring_multishot_recv.RING_ENTRIES = 8uRING_ENTRIES;
void* (local variable) void* ringPtrringPtr;
if (int core.sys.posix.stdlib.posix_memalign(scope void**, ulong, ulong) pure nothrow @nogcposix_memalign(&(local variable) void* ringPtrringPtr, 4096, (constant) ulong io_uring_multishot_recv.main.ringBytes = 128LUringBytes) != 0 || (local variable) void* ringPtrringPtr is null)
{
stderr.std.stdio.File std.stdio.makeGlobal!"core.stdc.stdio.stderr"() nothrow @nogc @property ref @systemwritefln("posix_memalign failed");
return 1;
}
scope (exit) void core.stdc.stdlib.free(void* ptr) nothrow @nogcfree((local variable) void* ringPtrringPtr);
auto (local variable) during.io_uring.io_uring_buf* ringring = cast((struct) during.io_uring.io_uring_bufio_uring_buf*) (local variable) void* ringPtrringPtr;
(local variable) during.io_uring.io_uring_buf* ringring[0 .. (constant) uint io_uring_multishot_recv.RING_ENTRIES = 8uRING_ENTRIES] = (struct) during.io_uring.io_uring_bufio_uring_buf.(constant) during.io_uring.io_uring_buf during.io_uring.io_uring_buf.init = io_uring_buf(0LU, 0u, cast(ushort)0u, cast(ushort)0u)init; // zero the whole ring (incl. tail=0)
// Backing storage for the buffers themselves (separate from the ring slots).
auto (local variable) ubyte[] storestore = new ubyte[(constant) uint io_uring_multishot_recv.BUF_SIZE = 64uBUF_SIZE * (constant) uint io_uring_multishot_recv.RING_ENTRIES = 8uRING_ENTRIES];
// Register the ring with the kernel for group `BGID`.
(struct) during.io_uring.io_uring_buf_regio_uring_buf_reg (local variable) during.io_uring.io_uring_buf_reg regreg;
(local variable) during.io_uring.io_uring_buf_reg regreg.(field) ulong during.io_uring.io_uring_buf_reg.ring_addrring_addr = cast(ulong) (local variable) void* ringPtrringPtr;
(local variable) during.io_uring.io_uring_buf_reg regreg.(field) uint during.io_uring.io_uring_buf_reg.ring_entriesring_entries = (constant) uint io_uring_multishot_recv.RING_ENTRIES = 8uRING_ENTRIES;
(local variable) during.io_uring.io_uring_buf_reg regreg.(field) ushort during.io_uring.io_uring_buf_reg.bgidbgid = (constant) ushort io_uring_multishot_recv.BGID = cast(ushort)7uBGID;
const (local variable) const(int) regRetregRet = (local variable) during.Uring ioio.int during.Uring.registerBufRing(ref scope during.io_uring.io_uring_buf_reg reg, uint flags = 0u) nothrow @nogc @trustedRegisters a kernel-side provided-buffer ring (IORING_REGISTER_PBUF_RING). reg must be
filled with the buffer ring address, entry count and group id; flags is OR'd into
reg`.`flags (e.g. IOU_PBUF_RING_INC). Mirrors liburing's io_uring_register_buf_ring.
registerBufRing((local variable) during.io_uring.io_uring_buf_reg regreg);
if ((local variable) const(int) regRetregRet == -(constant) int core.stdc.errno.EINVAL = 22EINVAL || (local variable) const(int) regRetregRet == -(constant) int core.stdc.errno.EOPNOTSUPP = 95EOPNOTSUPP || (local variable) const(int) regRetregRet == -(constant) int core.stdc.errno.ENOSYS = 38ENOSYS)
{
void std.stdio.writefln!(char, const(int))(in char[] fmt, const(int) __param_1) @safeEquivalent to writef(fmt, args, '\n').
writefln("SKIP: IORING_REGISTER_PBUF_RING unsupported (errno %d) — needs Linux 5.19+", -(local variable) const(int) regRetregRet);
return 0;
}
if ((local variable) const(int) regRetregRet < 0)
{
stderr.std.stdio.File std.stdio.makeGlobal!"core.stdc.stdio.stderr"() nothrow @nogc @property ref @systemwritefln("registerBufRing failed: errno %d", -(local variable) const(int) regRetregRet);
return 1;
}
scope (exit) (local variable) during.Uring ioio.int during.Uring.unregisterBufRing(int bgid) nothrow @nogc @trustedUnregisters the provided-buffer ring with group id bgid.
Mirrors liburing's io_uring_unregister_buf_ring.
unregisterBufRing((constant) ushort io_uring_multishot_recv.BGID = cast(ushort)7uBGID);
// Publish all RING_ENTRIES buffers into the ring. Each slot points at its
// BUF_SIZE chunk of `store` and carries a distinct buffer id (`bid`). The
// kernel returns the chosen `bid` in the CQE's upper 16 bits.
enum uint (constant) uint io_uring_multishot_recv.main.mask = 7umask = (constant) uint io_uring_multishot_recv.RING_ENTRIES = 8uRING_ENTRIES - 1;
foreach (ushort (local variable) ushort ii; 0 .. cast(ushort) (constant) uint io_uring_multishot_recv.RING_ENTRIES = 8uRING_ENTRIES)
{
auto (local variable) during.io_uring.io_uring_buf* slotslot = &(local variable) during.io_uring.io_uring_buf* ringring[(local variable) ushort ii & (constant) uint io_uring_multishot_recv.main.mask = 7umask];
(local variable) during.io_uring.io_uring_buf* slotslot.(field) ulong during.io_uring.io_uring_buf.addraddr = cast(ulong) &(local variable) ubyte[] storestore[(local variable) ushort ii * (constant) uint io_uring_multishot_recv.BUF_SIZE = 64uBUF_SIZE];
(local variable) during.io_uring.io_uring_buf* slotslot.(field) uint during.io_uring.io_uring_buf.lenlen = (constant) uint io_uring_multishot_recv.BUF_SIZE = 64uBUF_SIZE;
(local variable) during.io_uring.io_uring_buf* slotslot.(field) ushort during.io_uring.io_uring_buf.bidbid = (local variable) ushort ii; // buffer id == index, so we can recover the chunk later
}
// Publish: advance the producer tail by the number of buffers we added. The
// tail lives in slot 0 (it overlays io_uring_buf.resv there).
(local variable) during.io_uring.io_uring_buf* ringring[0].(field) ushort during.io_uring.io_uring_buf.resvresv = cast(ushort) (constant) uint io_uring_multishot_recv.RING_ENTRIES = 8uRING_ENTRIES;
// Arm ONE multishot RECV that selects buffers from group BGID. The gid overload
// sets IOSQE_BUFFER_SELECT + sqe->buf_group, and the multishot wrapper adds
// IORING_RECV_MULTISHOT — so this single SQE will keep posting a CQE per segment.
enum ulong (constant) ulong io_uring_multishot_recv.main.recvCookie = 186015461LUrecvCookie = 0xB16_5EE5;
// Pass `sv[0]` as an explicit arg rather than capturing it: a capturing lambda
// would force a GC closure and break `putWith`'s `@nogc`.
(local variable) during.Uring ioio.putWith!((ref SubmissionEntry e, int recvFd) {
e.prepRecvMultishot(recvFd, BGID, BUF_SIZE);
e.user_data = recvCookie;
})(during.Uring during.Uring.putWith!(function (ref during.io_uring.SubmissionEntry e, int recvFd) nothrow @nogc @safe
{
prepRecvMultishot(e, recvFd, cast(ushort)7u, 64u, MsgFlags.NONE);
e.user_data = 186015461LU;
}
, int)(ref int __param_0) nothrow @nogc return ref @safeAdds new entry to the SubmissionQueue.
Note that this just adds entry to the queue and doesn't advance the tail
marker kernel sees. For that finishSq() is needed to be called next.
Also note that to actually enter new entries to kernel,
it's needed to call submit().
sv[0]);
const (local variable) const(int) submittedsubmitted = (local variable) during.Uring ioio.int during.Uring.submit(uint want) nothrow @nogc @safeSubmits qued SubmissionEntry to be processed by kernel.
submit(0); // flush the SQ without blocking on a count
if ((local variable) const(int) submittedsubmitted < 0)
{
stderr.std.stdio.File std.stdio.makeGlobal!"core.stdc.stdio.stderr"() nothrow @nogc @property ref @systemwritefln("submit failed: errno %d", -(local variable) const(int) submittedsubmitted);
return 1;
}
// Two separate messages to feed the peer end. Each becomes its own readable
// segment on sv[0], so the armed multishot RECV should post one CQE per message.
immutable ubyte[][2] (local variable) immutable(ubyte[][2]) payloadspayloads = [
cast(immutable ubyte[]) "first multishot segment",
cast(immutable ubyte[]) "second multishot segment!!",
];
// Drive one segment at a time: write a message, then drain the CQE the armed SQE
// produces for it, *before* writing the next. On a stream socket two back-to-back
// writes can coalesce into a single readable chunk (and thus a single CQE);
// interleaving write-then-drain guarantees two distinct segments — two CQEs from
// the one armed SQE, each selecting a fresh ring buffer — deterministically,
// without relying on kernel scheduling. We never re-arm: the multishot SQE stays
// live across both completions (signalled by CQE_F_MORE).
ushort[2] (local variable) ushort[2] gotBidsgotBids;
foreach ((parameter) ulong idxidx, (parameter) immutable(ubyte[]) pp; (local variable) immutable(ubyte[][2]) payloadspayloads)
{
const (local variable) const(long) wrotewrote = long core.sys.posix.unistd.write(int, scope const(void*), ulong) nothrow @nogcwrite((local variable) int[2] svsv[1], &(local variable) immutable(ubyte[]) pp[0], (local variable) immutable(ubyte[]) pp.(field) ulong immutable(ubyte[]).lengthlength);
if ((local variable) const(long) wrotewrote != cast((alias) object.ptrdiff_t = longptrdiff_t) (local variable) immutable(ubyte[]) pp.(field) ulong immutable(ubyte[]).lengthlength)
{
stderr.std.stdio.File std.stdio.makeGlobal!"core.stdc.stdio.stderr"() nothrow @nogc @property ref @systemwritefln("write to peer failed: %d", (local variable) const(long) wrotewrote);
return 1;
}
(local variable) during.Uring ioio.int during.Uring.wait(uint want = 1u) nothrow @nogcSimmilar to submit but with this method we just wait for required number
of CompletionEntries.
wait(1); // one CQE is guaranteed imminent (data already written)
const (local variable) const(during.io_uring.CompletionEntry) cc = (local variable) during.Uring ioio.during.io_uring.CompletionEntry during.Uring.front() pure nothrow @nogc return ref @safeGet first CompletionEntry from cq ring
front;
const (local variable) const(int) resres = (local variable) const(during.io_uring.CompletionEntry) cc.(field) int during.io_uring.CompletionEntry.resresult code for this event
res;
const (local variable) const(during.io_uring.CQEFlags) flagsflags = (local variable) const(during.io_uring.CompletionEntry) cc.(field) during.io_uring.CQEFlags during.io_uring.CompletionEntry.flagsflags;
const (local variable) const(ulong) echoedechoed = (local variable) const(during.io_uring.CompletionEntry) cc.(field) ulong during.io_uring.CompletionEntry.user_datasqe->data submission passed back
user_data;
(local variable) during.Uring ioio.void during.Uring.popFront() pure nothrow @nogc @safeMove to next CompletionEntry
popFront();
if ((local variable) const(ulong) echoedechoed != (constant) ulong io_uring_multishot_recv.main.recvCookie = 186015461LUrecvCookie)
{
stderr.std.stdio.File std.stdio.makeGlobal!"core.stdc.stdio.stderr"() nothrow @nogc @property ref @systemwritefln("user_data mismatch: expected 0x%X, got 0x%X", (constant) ulong io_uring_multishot_recv.main.recvCookie = 186015461LUrecvCookie, (local variable) const(ulong) echoedechoed);
return 1;
}
// A pre-6.0 kernel rejects the multishot bit on the very first CQE.
if ((local variable) const(int) resres == -(constant) int core.stdc.errno.EINVAL = 22EINVAL || (local variable) const(int) resres == -(constant) int core.stdc.errno.EOPNOTSUPP = 95EOPNOTSUPP || (local variable) const(int) resres == -(constant) int core.stdc.errno.ENOSYS = 38ENOSYS)
{
void std.stdio.writefln!(char, const(int))(in char[] fmt, const(int) __param_1) @safeEquivalent to writef(fmt, args, '\n').
writefln("SKIP: IORING_RECV_MULTISHOT unsupported (errno %d) — needs Linux 6.0+", -(local variable) const(int) resres);
return 0;
}
if ((local variable) const(int) resres < 0)
{
// -ENOBUFS would mean the ring ran out of published buffers — a real bug
// in our bookkeeping, not an unsupported-feature case.
stderr.std.stdio.File std.stdio.makeGlobal!"core.stdc.stdio.stderr"() nothrow @nogc @property ref @systemwritefln("multishot RECV CQE #%d failed: errno %d", (local variable) ulong idxidx, -(local variable) const(int) resres);
return 1;
}
// Multishot contract: each non-terminal CQE carries MORE (still armed). With
// two segments and 8 buffers the op cannot exhaust the ring, so both of our
// CQEs must be armed.
if (!((local variable) const(during.io_uring.CQEFlags) flagsflags & (enum) during.io_uring.CQEFlagsFlags used with CompletionEntry
CQEFlags.(enum value) during.io_uring.CQEFlags.MORE = 2uIORING_CQE_F_MORE (from Linux 5.13)
If set, parent SQE will generate more CQE entries
MORE))
{
stderr.std.stdio.File std.stdio.makeGlobal!"core.stdc.stdio.stderr"() nothrow @nogc @property ref @systemwritefln("multishot RECV CQE #%d lost CQE_F_MORE (flags=0x%X) — op disarmed early",
(local variable) ulong idxidx, cast(uint) (local variable) const(during.io_uring.CQEFlags) flagsflags);
return 1;
}
// Buffer-select contract: a buffer must have been chosen, with its id in the
// upper 16 bits of flags.
if (!((local variable) const(during.io_uring.CQEFlags) flagsflags & (enum) during.io_uring.CQEFlagsFlags used with CompletionEntry
CQEFlags.(enum value) during.io_uring.CQEFlags.BUFFER = 1uIORING_CQE_F_BUFFER (from Linux 5.7)
If set, the upper 16 bits are the buffer ID
BUFFER))
{
stderr.std.stdio.File std.stdio.makeGlobal!"core.stdc.stdio.stderr"() nothrow @nogc @property ref @systemwritefln("multishot RECV CQE #%d has no CQE_F_BUFFER (flags=0x%X)",
(local variable) ulong idxidx, cast(uint) (local variable) const(during.io_uring.CQEFlags) flagsflags);
return 1;
}
const (local variable) const(ushort) bidbid = cast(ushort)(cast(uint) (local variable) const(during.io_uring.CQEFlags) flagsflags >> (enum value) during.io_uring.CQE_BUFFER_SHIFT = 16Note
available from Linux 5.7
CQE_BUFFER_SHIFT);
if ((local variable) const(ushort) bidbid >= (constant) uint io_uring_multishot_recv.RING_ENTRIES = 8uRING_ENTRIES)
{
stderr.std.stdio.File std.stdio.makeGlobal!"core.stdc.stdio.stderr"() nothrow @nogc @property ref @systemwritefln("kernel returned out-of-range buffer id %d", (local variable) const(ushort) bidbid);
return 1;
}
// Confirm the bytes landed in exactly the buffer the kernel selected.
auto (local variable) ubyte[] gotgot = (local variable) ubyte[] storestore[(local variable) const(ushort) bidbid * (constant) uint io_uring_multishot_recv.BUF_SIZE = 64uBUF_SIZE .. (local variable) const(ushort) bidbid * (constant) uint io_uring_multishot_recv.BUF_SIZE = 64uBUF_SIZE + (local variable) const(int) resres];
if ((local variable) ubyte[] gotgot != (local variable) immutable(ubyte[][2]) payloadspayloads[(local variable) ulong idxidx])
{
stderr.std.stdio.File std.stdio.makeGlobal!"core.stdc.stdio.stderr"() nothrow @nogc @property ref @systemwritefln("payload #%d mismatch in selected buffer %d", (local variable) ulong idxidx, (local variable) const(ushort) bidbid);
return 1;
}
(local variable) ushort[2] gotBidsgotBids[(local variable) ulong idxidx] = (local variable) const(ushort) bidbid;
}
// The whole point of multishot: two segments, two CQEs, from ONE armed SQE — and
// each landed in a *distinct* ring buffer (the kernel pops a fresh slot per CQE).
if ((local variable) ushort[2] gotBidsgotBids[0] == (local variable) ushort[2] gotBidsgotBids[1])
{
stderr.std.stdio.File std.stdio.makeGlobal!"core.stdc.stdio.stderr"() nothrow @nogc @property ref @systemwritefln("expected distinct buffers, both CQEs used buffer id %d", (local variable) ushort[2] gotBidsgotBids[0]);
return 1;
}
void std.stdio.writefln!(char, ushort, ushort)(in char[] fmt, ushort __param_1, ushort __param_2) @safeEquivalent to writef(fmt, args, '\n').
writefln("ok: one armed multishot RECV posted 2 CQEs (MORE+BUFFER) into distinct buffers %d and %d",
(local variable) ushort[2] gotBidsgotBids[0], (local variable) ushort[2] gotBidsgotBids[1]);
return 0;
}