Skip to content

multishot reads - #21

Open
adamsitnik wants to merge 25 commits into
io_uring_threadsfrom
io_uring_multishot_reads_local
Open

adamsitnik wants to merge 25 commits into
io_uring_threadsfrom
io_uring_multishot_reads_local

Conversation

@adamsitnik

Copy link
Copy Markdown
Owner

No description provided.

adamsitnik and others added 9 commits September 23, 2026 15:22
Replace the per-calling-thread round-robin ring assignment
(t_assignedRing/s_nextRingIndex, sticky per OS thread) with a
fd % ringCount mapping: GetAssignedRing() becomes GetRing(IntPtr fd),
derived directly from the fd already present on every IoRingRequest
(request.Fd), so no TrySubmit call site needs to change.

This is a prerequisite for upcoming multishot operations (e.g. a
persistent multishot receive per socket): those need a stable,
recomputable fd -> ring mapping to find and cancel a specific fd's
in-flight operation, which a per-thread assignment cannot provide (the
same fd's operations could, in principle, have been submitted from
different calling threads). fd allocation on Unix is a small, densely
packed, monotonically increasing counter, so a plain modulo still
spreads load reasonably evenly across rings.

Verified: System.Net.Sockets.Tests IoUringTests (14 cases, ring counts
1 and 3) pass unchanged against a checked System.Private.CoreLib
rebuild.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Require IORING_RECV_MULTISHOT, IORING_REGISTER_PBUF_RING,
IOSQE_BUFFER_SELECT, and struct io_uring_buf_ring to be available at
compile time for io_uring support to be considered available at all.
Kernels/headers missing this fall back to the existing non-io_uring
code paths entirely, the same as if the io_uring syscalls themselves
were unavailable.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
…ceive

Add IoRingOp_RecvMultishot and IoRingOp_Cancel opcodes, plus
SystemNative_IoRingRegisterBufferRing/SystemNative_IoRingReturnBuffers,
which allocate and own a page-aligned, natively-allocated pool of
provided buffers (IORING_REGISTER_PBUF_RING) for a ring, publish/
republish buffer ids to the kernel with no syscall on the return path,
and free the storage on ring close. Mirror the new opcodes and PAL
entry points in Interop.IoRing.cs.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Adds PortableThreadPool.IoUring.Receive.Unix.cs implementing
IoUringThreadPool.TrySubmitReceiveMultishot/TryCancelReceiveMultishot on
top of the native IORING_OP_RECV + IORING_RECV_MULTISHOT + provided-buffer
support added previously:

- ReceiveBufferPool/ReceiveBufferLease wrap the ring's natively-owned,
  page-aligned provided-buffer storage as zero-copy MemoryManager<byte>
  leases, one persistent lease object reused per buffer id; disposing a
  lease enqueues its buffer id onto Ring.PendingBufferReturns for the
  issuer to republish opportunistically (see PortableThreadPool.IoUring.
  Unix.cs's DrainReceiveBufferReturns), never waking the issuer just to
  return a single buffer.
- MultishotReceiveOperation implements IIoUringOperation. A single
  submission produces many completions over its lifetime, and because
  CompletionProcessorWorkItem can schedule a second concurrent processor
  for the same ring before the first has finished, two completions of the
  same still-active operation can race onto two different Thread Pool
  worker threads - so each completion's dispatch captures its own
  self-contained state (via ThreadPool.UnsafeQueueUserWorkItem<TState>)
  instead of reusing mutable fields on the long-lived operation, unlike
  the existing single-completion-only operations in this file.
- Cancellation reuses IoRingOp.Cancel, keyed by the target operation's
  UserData token recorded in the new Ring.ActiveMultishotReceives
  fd -> token map (populated via a new TrySubmit(..., out ulong userData)
  overload), looked up via the same fd -> ring routing used to submit the
  original request.

Exposes this as IoUring.TrySubmitRecvMultishot/TryCancelRecvMultishot
(and the corresponding System.Runtime.cs reference-assembly entries).

Widens PortableThreadPool.IoUringThreadPool.Ring and OperationSlot from
private to internal so this new partial file - and the receive-buffer
registration handshake added to the static constructor - can reference
them across the partial-class-file boundary.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Two completions belonging to the same still-active multishot operation
(e.g. IORING_OP_RECV with IORING_RECV_MULTISHOT) can be dequeued off a
ring's completion queue in order, but then *processed* concurrently, in
either order, by two independent Thread Pool workers - the completion
processor deliberately overlaps a new worker with the one already
draining a ring's queue for throughput.

This broke two separate invariants:

- Slot/token lifetime: the final completion (no CQE 'More' flag) used to
  release its operation's slot/GCHandle immediately. If an earlier,
  non-final completion for the same operation was still being processed
  on another worker at that moment, it could observe an already-freed or
  generation-bumped slot and hit
  Environment.FailFast("Stale io_uring operation slot.").
- Delivery order: MultishotReceiveOperation dispatches each completion as
  an independent Thread Pool work item, with no guarantee those work
  items run in arrival order - so a caller iterating
  Socket.ReceiveMultishotAsync could observe completions (e.g. the final
  EOF/error completion) out of order relative to earlier data.

Fix:
- Give each operation's slot/GCHandle token a reference count, taken once
  per completion by the single issuer thread strictly before that
  completion is handed off to any worker, and released once each
  completion has been processed - the token is only actually freed once
  every reference (including the two extra ones a final completion
  carries: its own, and the initial "not yet finalized" bias) has been
  released, regardless of processing order.
- Assign each completion a 0-based delivery sequence number, also by the
  issuer at the same point, and have MultishotReceiveOperation gate its
  own delivery on that sequence (lock-free spin-wait) so the caller's
  callback always observes completions in true arrival order.

IIoUringOperation.CompleteFromIoUring gained a 'long sequence' parameter
threaded through from the issuer; existing single-completion operations
(ActionIoUringOperation, ThreadPoolValueTaskSource) ignore it, since it is
always 0 for them.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Adds Socket.ReceiveMultishotAsync(CancellationToken), an
IAsyncEnumerable<IMemoryOwner<byte>>-returning API backed by
IoUring.TrySubmitRecvMultishot: a single persistent submission for the
lifetime of the connection, with the kernel delivering a completion (and
a ring-mapped provided buffer) each time it reads new data. Consumers
return buffers to the pool by disposing the yielded IMemoryOwner<byte>.

- Socket.Multishot.Linux.cs (compiled for TargetPlatformIdentifier ==
  'unix'): the actual implementation, using an unbounded
  Channel<IMemoryOwner<byte>> (SingleReader: true, SingleWriter: false -
  see note below) to bridge the io_uring completion callback to the
  async enumerator. Cancellation (via the caller's token, or enumerator
  disposal) calls IoUring.TryCancelRecvMultishot and completes the
  channel; unyielded buffers left in the channel are disposed. Throws
  InvalidOperationException when io_uring is unavailable.
- Socket.Multishot.cs (compiled for every other TargetPlatformIdentifier):
  throws PlatformNotSupportedException.
- Added the member to the ref assembly and a new SR string for the
  InvalidOperationException message.
- Added a ProjectReference to System.Threading.Channels.
- Added functional tests (IoUring.Unix.cs): streaming multiple sends in
  order, graceful EOF shutdown, cancellation, and the
  InvalidOperationException thrown when io_uring is disabled.

Note: the Channel is SingleWriter: false rather than the originally
planned SingleWriter: true - a still-active multishot operation's
completions can be *processed* by more than one Thread Pool worker at a
time (see the companion completion-ordering fix), so more than one
thread can call the channel writer concurrently.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
…ion races in multishot receive

Three additional bugs found while benchmarking ReceiveMultishotAsync under
concurrent long-lived connections:

- The out-of-order completion ordering gate in MultishotReceiveOperation.Deliver
  used to busy-spin waiting for its turn. Under load, every available Thread
  Pool worker could end up spinning on a predecessor completion that itself sat
  undispatched in the ring's completion queue waiting for a free worker -
  deadlocking the whole pool until slow Hill-Climbing thread injection rescued
  it. Deliver now requeues itself instead of ever blocking the calling thread.

- The kernel terminates (rather than pauses) a multishot receive once the
  ring's shared provided-buffer pool is momentarily exhausted (ENOBUFS) - this
  is documented, expected behavior under concurrent long-lived receives, not a
  real error. This previously surfaced to callers as a random SocketException.
  The operation now transparently resubmits a fresh multishot request instead,
  invisible to the caller.

- Resolving TryCancelReceiveMultishot's target operation via the token/slot
  machinery from an arbitrary calling thread (one holding no reference on that
  token) was unsafe, and resetting the delivery-ordering sequence counter
  inside the ENOBUFS resubmit path could race with sibling completions of the
  same submission still awaiting dispatch on other workers, permanently
  orphaning them. Ring.ActiveMultishotReceives now tracks the operation
  instance directly (safe to read from any thread), and the ENOBUFS re-arm
  decision moved into Deliver, which only makes it once its own ordering gate
  has confirmed every prior completion of that submission was already
  delivered.

Verified via a standalone concurrent-connection benchmark harness that
previously reproduced all three bugs (hangs/timeouts) reliably; re-run 30+
times after these fixes with zero failures. Also re-ran the existing
IoUringTests suite (21/21 passing).

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
…eive completions

CompleteFromIoUring previously always explicitly self-queued via
ThreadPool.UnsafeQueueUserWorkItem before calling Deliver, even when the
completion was already known to be next in arrival order. Now it returns
the completion itself as the work item (letting the existing dispatch
mechanism run it, inline or batched, exactly like every other operation)
whenever the ordering gate is already satisfied, and only falls back to
the extra explicit queue-and-retry when it genuinely is out of turn.

CompletionState now implements IThreadPoolWorkItem directly instead of
being dispatched via a static Action<CompletionState>, so it can serve
both roles without changes to Deliver's own retry path.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
… channel writer ordering comment

- Skip allocating the CancellationToken.UnsafeRegister callback delegate
  when the token can never be canceled (e.g. CancellationToken.None).
  Register/UnsafeRegister already short-circuit internally for that case,
  but the delegate passed to them is still constructed by the caller
  regardless - guarding with CanBeCanceled avoids that allocation.
- Correct a stale comment on the receive channel's SingleWriter = false:
  MultishotReceiveOperation.Deliver's sequence gate already guarantees
  completions are delivered to the channel strictly in order and never
  concurrently, so the write is already effectively single-writer today.
@adamsitnik
adamsitnik changed the base branch from main to io_uring_threads September 23, 2026 16:54
MultishotReceiveOperation.Deliver's sequence gate already guarantees
completions are handed to the channel writer strictly one at a time and
in true arrival order, regardless of which Thread Pool worker actually
executes each completion - so the channel write is already effectively
single-writer. Flip Channel.CreateUnbounded's SingleWriter option to
true to let it skip its internal writer synchronization.

Add a dedicated stress test (many concurrent connections sharing a
handful of rings, each streamed with a rapid one-byte-at-a-time burst)
to validate this holds under real concurrent-completion pressure, not
just the existing single-connection ordering test. Ran 25x locally with
no failures.

@adamsitnik adamsitnik left a comment

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Looks promising, but could be simplified a bit.

adamsitnik and others added 15 commits September 23, 2026 20:59
…ration cancellation handle

Implements three PR review comments:
- TrySubmitReceiveMultishot/TrySubmitRecvMultishot now return an
  out IIoUringOperation? operation instead of registering the
  in-flight request in a fd-keyed dictionary; callers cancel by
  calling operation.RequestCancellation() directly.
- Removed Ring.ActiveMultishotReceives entirely.
- Replaced the _deliveredThrough completion sequence-gate with a
  simpler design: since every fd (and so every operation) is bound
  to exactly one ring with exactly one issuer thread, the issuer
  now enqueues completions directly onto the operation's own
  single-producer queue and schedules a single coalesced drainer,
  guaranteeing true arrival order with no per-completion sequence
  bookkeeping.
- Hoisted IIoUringOperation to a public top-level interface (moved
  its ref declaration from System.Runtime.cs to
  System.Threading.ThreadPool.cs, which can reference
  IThreadPoolWorkItem).

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Keep accepts with a cancellable token on the existing readiness-based path. The io_uring accept fast path does not register cancellation, so selecting it could leave listener shutdown waiting for a connection that never arrives.

Uncancellable accepts continue to use io_uring.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Carry bytes already written by the optimistic Unix send into the io_uring completion callback and add the asynchronous result to that count.

Reporting only the remaining asynchronous transfer loses the successful prefix and can make callers resend data that is already on the wire. Receive completions retain their existing zero-prefix behavior.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Release each native receive's SafeHandle reference on the issuer when its terminal CQE arrives, and reacquire a reference for every rearm. Socket disposal must not wait for a reference that only its own blocked callback worker can release.

Publish cancellation and replacement tokens with full fences, recheck cancellation after rearming, and use compare-exchange for the initial token. This prevents missed cancellation and a late initial publication overwriting a replacement submission's token. Update the bridge lifetime documentation.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Reset ExecutionContext and ThreadPool thread state after each callback drained by a multishot receive work item.

Several callbacks can run inside one work item, so the ThreadPool's outer cleanup alone does not prevent AsyncLocal or SynchronizationContext changes from leaking into the next callback. Add single-worker coverage with one and three rings to exercise consecutive callbacks on the same worker.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Allow Channel continuations to execute on the existing receive callback worker and read the Channel directly instead of adding a second async iterator. This removes an avoidable worker hop and iterator layer without running application code on the issuer.

Publish cancellation visibly across threads and asynchronously drain unyielded buffers until the terminal callback. Native cancellation is not a drain barrier; a one-time dequeue can leak buffers arriving after an early break. Add repeated early-break coverage with a four-buffer pool and one worker.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Use full-fence exchanges when resetting the issuer wake flag and the receive worker's dispatch flag before checking their queues again.

Release-only stores permit the subsequent queue read to miss a racing producer whose enqueue also observed the old flag. That can strand a receive callback or leave an idle issuer asleep indefinitely. Add repeated pending-receive bursts across 64 connections and three rings to exercise the worker handoff.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Treat a positive CQE without MORE as the end of a native submission, not the end of the byte stream. Deliver its buffer with logical continuation, then rearm unless the callback canceled the receive or closed its handle.

Report rearm failure as cancellation or a closed-handle error rather than successful EOF, and reset callback context between data delivery and a possible terminal notification. Completion-queue pressure otherwise truncates streams even while unread data remains in the socket.

Add an end-to-end CQ-pressure regression and targeted cancellation/handle-disposal callback tests. Update the Socket documentation to describe transparent rearming.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Replace the per-operation ConcurrentQueue with the existing SingleProducerSingleConsumerQueue and keep dispatch scheduling in its single enqueue entry point.

The fd-bound issuer is the only producer and the dispatch flag permits one active drainer, so the general-purpose concurrent queue provides synchronization this path does not need. Add coverage that blocks the drainer while enough completions accumulate to cross a queue segment. An isolated external throughput gain has not been established.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Leave the submission's correlation-token reference in place across nonterminal multishot completions, and retire it only when the native request ends.

Multishot workers retain the operation object directly and never resolve its token. Retaining, sequencing and releasing that token for every buffer therefore adds issuer work without protecting any additional lifetime. Keep the existing token protocol for ordinary I/O completions.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Replace the ring's concurrent return queue with one bit per provided buffer. Disposing a lease atomically sets its bit; the issuer takes returned bits and republishes their IDs through the existing batch API.

A buffer can have only one outstanding return and return order is irrelevant, so a bitmap avoids general-purpose queue bookkeeping. Preserve the default pool dimensions, issuer-only publication and no per-buffer wakeup policy.

This is experimental tuning: the single external comparison did not establish a repeatable gain. A return to a word already scanned waits for a later issuer pass.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Queue compact result/flags records instead of buffer-owner references. Decode the selected buffer ID and set the lease length when the worker consumes the completion, while keeping terminal SafeHandle release on the issuer.

This removes managed references from completion storage and moves lease writes to the thread that immediately consumes the data. Preserve completion order and native buffer ownership until the consumer disposes the lease.

Strengthen the blocked-drainer test with exact lengths and retained-content checks, and add a four-buffer cycling test that holds an owner through enumeration and socket disposal. Functional validation passed; the external performance comparison is still pending, so this commit makes no throughput claim.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Revert fc7137c so cancellable tokens no longer exclude pending accepts from the io_uring path.

Accept cancellation is deliberately unsupported during this performance experiment. Routing those accepts through epoll changes the architecture being measured instead of addressing the experiment's throughput goal. Add a TODO: io_uring comment making that limitation explicit.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Skip SocketAsyncEngine creation when io_uring is active and throw explicitly if an operation attempts to register with the epoll-backed engine. Socket construction reads this type's inline-completion configuration, so registration checks alone would still create idle epoll engines.

Keep the performance experiment on its intended backend instead of silently measuring a mixture of io_uring and epoll. Add enabled/disabled tests for engine creation and rejected fallback.

Release build and all 36 io_uring tests pass. All 36 TCP cases pass with io_uring disabled. With io_uring enabled, broader TCP coverage exposes 18 failures involving existing multi-buffer and forced-nonblocking synchronous fallback paths; these are intentionally not hidden or implemented by this change.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant