Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
25 commits
Select commit Hold shift + click to select a range
91d3dd6
io_uring: route requests to a ring by fd, not by calling thread
adamsitnik Sep 23, 2026
f6e8e90
io_uring: make multishot receive with provided buffers mandatory
adamsitnik Sep 23, 2026
93e5469
io_uring: add native ring-owned provided-buffer pool for multishot re…
adamsitnik Sep 23, 2026
09b3434
io_uring: add managed multishot receive with provided-buffer pool
adamsitnik Sep 23, 2026
eea5879
io_uring: fix out-of-order multishot completion processing race
adamsitnik Sep 23, 2026
a65338f
io_uring sockets: add Socket.ReceiveMultishotAsync
adamsitnik Sep 23, 2026
c114df6
io_uring: fix ThreadPool deadlock, lost ENOBUFS re-arm, and cancellat…
adamsitnik Sep 23, 2026
91f5379
io_uring: avoid a redundant ThreadPool hop for in-order multishot rec…
adamsitnik Sep 23, 2026
91a8331
io_uring: avoid unnecessary cancellation delegate allocation; clarify…
adamsitnik Sep 23, 2026
19021bf
io_uring: set multishot receive channel to SingleWriter=true
adamsitnik Sep 23, 2026
b09c4bc
io_uring: replace ActiveMultishotReceives dictionary with IIoUringOpe…
adamsitnik Sep 23, 2026
fc7137c
Honor cancellation when selecting the io_uring accept path
adamsitnik Sep 24, 2026
72272df
Preserve optimistic send bytes in io_uring completions
adamsitnik Sep 24, 2026
0e097e2
Fix multishot handle lifetime and cancellation publication
adamsitnik Sep 24, 2026
c53c9d3
Isolate thread state between multishot receive callbacks
adamsitnik Sep 24, 2026
1713b72
Inline multishot consumers and drain canceled enumerations
adamsitnik Sep 24, 2026
1cd7543
Prevent lost wakeups in io_uring issuer and receive dispatch
adamsitnik Sep 24, 2026
aea5ded
Rearm multishot receives after positive terminal completions
adamsitnik Sep 24, 2026
661c724
Use SPSC storage for ordered multishot completions
adamsitnik Sep 24, 2026
5381d56
Avoid token reference churn for nonterminal multishot CQEs
adamsitnik Sep 24, 2026
c827e88
Track returned io_uring buffers with an atomic bitmap
adamsitnik Sep 24, 2026
812139d
Materialize multishot receive buffers on callback workers
adamsitnik Sep 24, 2026
caf1141
Revert cancellable-accept epoll fallback for io_uring experiment
adamsitnik Sep 24, 2026
a27d564
Reject epoll socket fallback while io_uring is enabled
adamsitnik Sep 24, 2026
b332679
Document why io_uring routing must remain descriptor-based
adamsitnik Sep 24, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,8 @@ internal enum IoRingOp : int
Connect = 5,
Recv = 6,
Send = 7,
Cancel = 8,
RecvMultishot = 9,
}

// Mirrors the native IoRingRequest struct in pal_io.h.
Expand All @@ -44,6 +46,15 @@ internal unsafe struct IoRingRequest
[StructLayout(LayoutKind.Sequential)]
internal struct IoRingCompletion
{
// IORING_CQE_F_MORE: this is not the final completion for the request that produced
// it (e.g. a RecvMultishot request that is still active and will keep completing).
public const uint More = 1 << 1;
// IORING_CQE_F_BUFFER: Flags encodes the selected provided-buffer id, shifted left by
// BufferShift (IORING_CQE_BUFFER_SHIFT) - only set for ops that use provided buffers
// (RecvMultishot).
public const uint Buffer = 1;
public const int BufferShift = 16;

public ulong UserData;
public int Result;
public uint Flags;
Expand Down Expand Up @@ -74,6 +85,16 @@ internal struct IoRingCompletion
[LibraryImport(Libraries.SystemNative, EntryPoint = "SystemNative_IoRingRegisterEventFd", SetLastError = true)]
internal static partial int IoRingRegisterEventFd(IntPtr ringHandle);

// Allocates bufferCount page-aligned buffers of bufferSize bytes each, natively owned by
// (and freed together with) the ring, registers them as provided-buffer group zero, and
// publishes all of them. On success, *bufferStorage points at the base of that storage
// (buffer i occupies [bufferStorage + i * bufferSize, bufferStorage + (i + 1) * bufferSize)).
[LibraryImport(Libraries.SystemNative, EntryPoint = "SystemNative_IoRingRegisterBufferRing", SetLastError = true)]
internal static unsafe partial int IoRingRegisterBufferRing(IntPtr ringHandle, int bufferSize, int bufferCount, byte** bufferStorage);

[LibraryImport(Libraries.SystemNative, EntryPoint = "SystemNative_IoRingReturnBuffers", SetLastError = true)]
internal static unsafe partial int IoRingReturnBuffers(IntPtr ringHandle, ushort* bufferIds, int count);

[LibraryImport(Libraries.SystemNative, EntryPoint = "SystemNative_EventFdWrite", SetLastError = true)]
internal static partial int EventFdWrite(int eventFd);

Expand Down
2 changes: 2 additions & 0 deletions src/libraries/System.Net.Sockets/ref/System.Net.Sockets.cs
Original file line number Diff line number Diff line change
Expand Up @@ -400,6 +400,8 @@ public void Listen(int backlog) { }
public System.Threading.Tasks.ValueTask<int> ReceiveAsync(System.Memory<byte> buffer, System.Net.Sockets.SocketFlags socketFlags, System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) { throw null; }
public System.Threading.Tasks.ValueTask<int> ReceiveAsync(System.Memory<byte> buffer, System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) { throw null; }
public bool ReceiveAsync(System.Net.Sockets.SocketAsyncEventArgs e) { throw null; }
[System.Runtime.Versioning.SupportedOSPlatformAttribute("linux")]
public System.Collections.Generic.IAsyncEnumerable<System.Buffers.IMemoryOwner<byte>> ReceiveMultishotAsync(System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) { throw null; }
public int ReceiveFrom(byte[] buffer, int offset, int size, System.Net.Sockets.SocketFlags socketFlags, ref System.Net.EndPoint remoteEP) { throw null; }
public int ReceiveFrom(byte[] buffer, int size, System.Net.Sockets.SocketFlags socketFlags, ref System.Net.EndPoint remoteEP) { throw null; }
public int ReceiveFrom(byte[] buffer, ref System.Net.EndPoint remoteEP) { throw null; }
Expand Down
6 changes: 6 additions & 0 deletions src/libraries/System.Net.Sockets/src/Resources/Strings.resx
Original file line number Diff line number Diff line change
Expand Up @@ -174,6 +174,12 @@
<data name="net_sockets_zerolist" xml:space="preserve">
<value>The parameter {0} must contain one or more elements.</value>
</data>
<data name="net_sockets_io_uring_epoll_fallback" xml:space="preserve">
<value>Unexpected epoll fallback while io_uring is enabled.</value>
</data>
<data name="net_sockets_multishot_not_supported" xml:space="preserve">
<value>Multishot receive is not supported because the io_uring Thread Pool integration is unavailable on this system.</value>
</data>
<data name="net_sockets_blocking" xml:space="preserve">
<value>The operation is not allowed on a non-blocking Socket.</value>
</data>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -206,6 +206,14 @@
Link="Common\Interop\Unix\System.Native\Interop.SocketEvent.cs" />
</ItemGroup>

<ItemGroup Condition="'$(TargetPlatformIdentifier)' == 'unix'">
<Compile Include="System\Net\Sockets\Socket.Multishot.Linux.cs" />
</ItemGroup>

<ItemGroup Condition="'$(TargetPlatformIdentifier)' == 'windows' or '$(TargetPlatformIdentifier)' == 'osx' or '$(TargetPlatformIdentifier)' == 'ios' or '$(TargetPlatformIdentifier)' == 'tvos' or '$(TargetPlatformIdentifier)' == 'wasi'">
<Compile Include="System\Net\Sockets\Socket.Multishot.cs" />
</ItemGroup>

<ItemGroup Condition="'$(TargetPlatformIdentifier)' == 'unix' or '$(TargetPlatformIdentifier)' == 'wasi' or '$(TargetPlatformIdentifier)' == 'osx' or '$(TargetPlatformIdentifier)' == 'ios' or '$(TargetPlatformIdentifier)' == 'tvos'">
<Compile Include="System\Net\Sockets\SafeSocketHandle.Unix.cs" />
<Compile Include="System\Net\Sockets\SafeSocketHandle.Unix.OptionTracking.cs" />
Expand Down Expand Up @@ -323,6 +331,7 @@
<ProjectReference Include="$(LibrariesProjectRoot)System.Runtime\src\System.Runtime.csproj" />
<ProjectReference Include="$(LibrariesProjectRoot)System.Runtime.InteropServices\src\System.Runtime.InteropServices.csproj" />
<ProjectReference Include="$(LibrariesProjectRoot)System.Threading\src\System.Threading.csproj" />
<ProjectReference Include="$(LibrariesProjectRoot)System.Threading.Channels\src\System.Threading.Channels.csproj" />
<ProjectReference Include="$(LibrariesProjectRoot)System.Threading.Overlapped\src\System.Threading.Overlapped.csproj" />
<ProjectReference Include="$(LibrariesProjectRoot)System.Threading.ThreadPool\src\System.Threading.ThreadPool.csproj" />
</ItemGroup>
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,134 @@
// Licensed to the .NET Foundation under one or more agreements.
// The .NET Foundation licenses this file to you under the MIT license.

using System.Buffers;
using System.Collections.Generic;
using System.Runtime.CompilerServices;
using System.Threading;
using System.Threading.Channels;
using System.Threading.Tasks;

namespace System.Net.Sockets
{
public partial class Socket
{
/// <summary>
/// Streams received data as an <see cref="IAsyncEnumerable{T}"/>, backed by a single persistent
/// multishot io_uring receive submitted for the enumeration - see
/// <see cref="System.Threading.IoUring.TrySubmitRecvMultishot"/>. Unlike repeatedly calling
/// <see cref="ReceiveAsync(Memory{byte}, CancellationToken)"/> in a loop, there is normally one
/// submission for as long as the caller keeps enumerating: the kernel delivers data into
/// its own pool of buffers as it arrives, without this socket needing to re-arm a new read after
/// each one. Native submissions terminated by buffer exhaustion or completion-queue pressure are rearmed.
/// Each yielded <see cref="IMemoryOwner{Byte}"/> must be disposed once the caller is
/// done with it - this returns its buffer to the pool so the kernel can reuse it. Enumeration
/// ends (without an exception) on graceful peer shutdown; stopping enumeration early (e.g.
/// <c>break</c>, or disposing the enumerator) or triggering <paramref name="cancellationToken"/>
/// requests cancellation of the underlying receive.
/// </summary>
/// <exception cref="InvalidOperationException">
/// The io_uring Thread Pool integration is unavailable on this system (see
/// <see cref="System.Threading.IoUring.IsSupported"/>).
/// </exception>
[System.Runtime.Versioning.SupportedOSPlatform("linux")]
public IAsyncEnumerable<IMemoryOwner<byte>> ReceiveMultishotAsync(CancellationToken cancellationToken = default)
{
ThrowIfDisposed();

if (!System.Threading.IoUring.IsSupported)
{
throw new InvalidOperationException(SR.net_sockets_multishot_not_supported);
}

return ReceiveMultishotAsyncCore(cancellationToken);
}

private async IAsyncEnumerable<IMemoryOwner<byte>> ReceiveMultishotAsyncCore([EnumeratorCancellation] CancellationToken cancellationToken)
{
SafeSocketHandle handle = _handle;
bool cancellationRequested = false;

// Completion callbacks have one active worker drainer, and enumeration has one consumer.
// Inline continuations avoid another worker hop; the finite provided-buffer pool bounds
// how many leases can be queued even though the channel itself is unbounded.
Channel<IMemoryOwner<byte>> channel = Channel.CreateUnbounded<IMemoryOwner<byte>>(new UnboundedChannelOptions
{
SingleReader = true,
SingleWriter = true,
AllowSynchronousContinuations = true,
});

void OnCompleted(int result, IMemoryOwner<byte>? buffer, bool hasMore)
{
if (buffer is not null)
{
channel.Writer.TryWrite(buffer);
}

if (!hasMore)
{
// result == 0 here means graceful EOF - not an error. A negative result while we
// are the ones who requested cancellation is reported as OperationCanceledException
// instead of the raw (usually -ECANCELED) SocketException.
Exception? error = result >= 0
? null
: Volatile.Read(ref cancellationRequested)
? new OperationCanceledException(cancellationToken)
: new SocketException((int)SocketPal.GetSocketErrorForErrorCode(new Interop.ErrorInfo(-result).Error));
channel.Writer.TryComplete(error);
}
}

if (!System.Threading.IoUring.TrySubmitRecvMultishot(handle, OnCompleted, out System.Threading.IIoUringOperation? operation))
{
throw new InvalidOperationException(SR.net_sockets_multishot_not_supported);
}

// Skip allocating the callback delegate entirely when the token can never be canceled
// (e.g. CancellationToken.None) - UnsafeRegister would end up being a no-op internally,
// but the delegate passed to it is still allocated by the caller regardless.
using CancellationTokenRegistration registration = cancellationToken.CanBeCanceled
? cancellationToken.UnsafeRegister(_ =>
{
Volatile.Write(ref cancellationRequested, true);
operation!.RequestCancellation();
}, null)
: default;

try
{
while (await channel.Reader.WaitToReadAsync(CancellationToken.None).ConfigureAwait(false))
{
while (channel.Reader.TryRead(out IMemoryOwner<byte>? buffer))
{
yield return buffer;
}
}
}
finally
{
// The caller may stop enumerating (break, or dispose the enumerator) while the
// underlying receive is still active - request its cancellation so the kernel
// eventually stops producing completions for it instead of leaking an in-flight
// multishot receive. A no-op if it already completed on its own.
Volatile.Write(ref cancellationRequested, true);
operation!.RequestCancellation();
try
{
// Cancellation is asynchronous. Keep returning unyielded leases until the
// terminal callback, including completions queued behind this continuation.
while (await channel.Reader.WaitToReadAsync(CancellationToken.None).ConfigureAwait(false))
{
while (channel.Reader.TryRead(out IMemoryOwner<byte>? leftover))
{
leftover.Dispose();
}
}
}
catch (OperationCanceledException) when (Volatile.Read(ref cancellationRequested))
{
}
}
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,19 @@
// Licensed to the .NET Foundation under one or more agreements.
// The .NET Foundation licenses this file to you under the MIT license.

using System.Buffers;
using System.Collections.Generic;
using System.Threading;

namespace System.Net.Sockets
{
public partial class Socket
{
/// <summary>
/// Not supported on this platform - multishot io_uring receives are only available on Linux.
/// </summary>
[System.Runtime.Versioning.SupportedOSPlatform("linux")]
public IAsyncEnumerable<IMemoryOwner<byte>> ReceiveMultishotAsync(CancellationToken cancellationToken = default) =>
throw new PlatformNotSupportedException();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -92,7 +92,7 @@ private void Return()
/// instead of registering the socket for epoll-based readiness notification. See
/// <see cref="TryReceiveViaIoUring"/> for the submission/callback contract.
/// </summary>
private unsafe bool TrySendViaIoUring(Memory<byte> buffer, int offset, int count, SocketFlags flags, Action<int, Memory<byte>, SocketFlags, SocketError> callback)
private unsafe bool TrySendViaIoUring(Memory<byte> buffer, int offset, int count, SocketFlags flags, int bytesSent, Action<int, Memory<byte>, SocketFlags, SocketError> callback)
{
if (!System.Threading.IoUring.IsSupported || flags != SocketFlags.None)
{
Expand All @@ -106,7 +106,7 @@ private unsafe bool TrySendViaIoUring(Memory<byte> buffer, int offset, int count
bufferPtr,
count,
0,
result => CompleteReceiveOrSend(pin, callback, result));
result => CompleteReceiveOrSend(pin, callback, result, bytesSent));

if (!submitted)
{
Expand All @@ -116,11 +116,11 @@ private unsafe bool TrySendViaIoUring(Memory<byte> buffer, int offset, int count
return submitted;
}

private static void CompleteReceiveOrSend(MemoryHandle pin, Action<int, Memory<byte>, SocketFlags, SocketError> callback, int result)
private static void CompleteReceiveOrSend(MemoryHandle pin, Action<int, Memory<byte>, SocketFlags, SocketError> callback, int result, int bytesAlreadyTransferred = 0)
{
pin.Dispose();

int bytesTransferred = result >= 0 ? result : 0;
int bytesTransferred = bytesAlreadyTransferred + (result >= 0 ? result : 0);
SocketError errorCode = result >= 0
? SocketError.Success
: SocketPal.GetSocketErrorForErrorCode(new Interop.ErrorInfo(-result).Error);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1510,6 +1510,8 @@ public SocketError AcceptAsync(Memory<byte> socketAddress, out int socketAddress
return errorCode;
}

// TODO: io_uring: Add accept cancellation support. Cancellation is intentionally
// unsupported while we compare io_uring performance against epoll.
if (ready && TryAcceptViaIoUring(socketAddress, callback))
{
acceptedFd = (IntPtr)(-1);
Expand Down Expand Up @@ -2089,7 +2091,7 @@ public SocketError SendToAsync(Memory<byte> buffer, int offset, int count, Socke
}

if (ready && socketAddress.Length == 0 &&
TrySendViaIoUring(buffer, offset, count, flags, callback))
TrySendViaIoUring(buffer, offset, count, flags, bytesSent, callback))
{
return SocketError.IOPending;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -84,6 +84,11 @@ private static int GetEngineCount()

private static SocketAsyncEngine[] CreateEngines()
{
if (IoUring.IsSupported)
{
return [];
}

int engineCount = GetEngineCount();

var engines = new SocketAsyncEngine[engineCount];
Expand Down Expand Up @@ -123,6 +128,11 @@ private static SocketAsyncEngine[] CreateEngines()
//
public static bool TryRegisterSocket(IntPtr socketHandle, SocketAsyncContext context, out SocketAsyncEngine? engine, out Interop.Error error)
{
if (IoUring.IsSupported)
{
throw new InvalidOperationException(SR.net_sockets_io_uring_epoll_fallback);
}

int engineIndex = Math.Abs(Interlocked.Increment(ref s_allocateFromEngine) % s_engines.Length);
SocketAsyncEngine nextEngine = s_engines[engineIndex];
bool registered = nextEngine.TryRegisterCore(socketHandle, context, out error);
Expand Down
Loading