Skip to content

Akka.IO TCP: hold back WriteAck until the output pipe's flush completes - #8646

Merged
Aaronontheweb merged 11 commits into
akkadotnet:devfrom
Aaronontheweb:feature/tcp-write-backpressure
Sep 29, 2026
Merged

Aaronontheweb merged 11 commits into
akkadotnet:devfrom
Aaronontheweb:feature/tcp-write-backpressure

Conversation

@Aaronontheweb

@Aaronontheweb Aaronontheweb commented Sep 25, 2026 •

Copy link
Copy Markdown
Member

Closes #8617

What was wrong

TcpConnection.EnqueueWrite threw away the ValueTask<FlushResult> from WriteAsync and acked every Write as soon as its bytes were copied into the output pipe. Nobody waited on the pipe's pause threshold, so when the socket drained slowly the pipe grew without bound and WriteAck promised nothing.

The fix

  • Write straight into the pipe; flush separately. ITransportConnection gains a write-only Write(ReadOnlySequence<byte>) that copies bytes into the output pipe without flushing, alongside the existing WriteAsync/FlushAsync. Every Write (past the existing per-write write-commands-queue-max-size cap) copies into the pipe and is disposed immediately; only its ack/failure waits for a flush to cover it. The actor no longer buffers raw write bytes for the open/registered path.
  • One flush at a time, and it batches. A flush covers whatever is currently unacked. If it completes synchronously (the common case), those writes are acked inline, with no allocation. If it doesn't, they wait for a FlushCompleted/FlushFailed self-message. Writes that arrive while a flush is already in flight are still copied into the pipe right away; their acks just ride the next flush - so several small writes landing mid-flush get acked together off one follow-up flush instead of one flush each.
  • A completed or canceled flush result fails every write still owed, oldest first, with CommandFailed, and closes the connection through HandleIoError. None of them are acked.
  • Close paths. Close, ConfirmedClose and the PeerClosed path still enter ClosingBehaviour right away and reject new writes. They start CloseAsync/ShutdownAsync only once nothing is owed and no flush is pending, because both complete the pipe. A Close that upgrades a pending ConfirmedClose (IO: handle Tcp.Close while a ConfirmedClose is draining #8636) goes through the existing upgrade handler and still drains first. Abort and handler death cancel the flush; every write still owed an ack gets CommandFailed from PostStop.
  • Writes before Register are unchanged - there's no transport yet, so they still buffer in a FIFO queue and replay through the same write path once Register arrives.
  • Output pipe size. Unchanged by default (PipeOptions.Default, pause at 64 KB). When Inet.SO.PipeBufferSize is set, the output pipe uses the same watermarks as the input pipe, so Artery (1 MiB) pauses at 2 MiB.

Out of scope: changing what write-commands-queue-max-size means (it still caps a single write). ResumeWriting stays a no-op as on dev; the unimplemented NACK-mode API (ResumeWriting, WritingResumed, useResumeWriting) is removed in a separate PR.

Breaking change

A WriteAck now means the output pipe is back below its resume mark (32 KB, or PipeBufferSize) after this write, so acks can arrive later under load. A single large write may overshoot the pause mark. While a flush is pending, later writes are copied into the pipe unflushed and their acks wait for the next flush, batching whatever arrived meanwhile. Recorded in BREAKING_CHANGES_V1.6.md.

Tests

New TcpConnectionWriteBackpressureSpec. It stalls the write pump with BlockingWriteStream, moved out of TcpConnectionBatchingSpec so both specs share it. No real-socket never-reading peers: Linux auto-tunes send buffers, which makes those flaky.

  • 8 x 16 KB writes: acks 1-3 arrive, then nothing for 300 ms, then acks 4-8 in order after the pump is released
  • three writes landing behind a pending flush: no ack until the stream is released, then all three ack in order, their bytes reaching the stream as a single follow-up write
  • ack order across senders and CompoundWrite parts
  • a flush result with IsCompleted (fake transport): CommandFailed, no ack, ErrorClosed
  • a later write's own flush failing while Closing drains (fake transport): the Close sender gets ErrorClosed
  • the write pump failing with a flush pending, in Open and in Closing
  • keepOpenOnPeerClosed then ConfirmedClose with a flush pending
  • writes buffered before Register that pass the pause threshold
  • Close, ConfirmedClose, PeerClosed and a Close upgrading a pending ConfirmedClose with a flush pending and writes owed: every ack before the close event, all bytes written, no extra acks
  • Abort and handler death with a flush pending and writes owed: CommandFailed for each write
  • owned-segment disposal: every write is freed right after its bytes are copied into the pipe, exactly once, whether its ack lands normally or Abort fails it instead

The first 12 fail on dev (the first commit adds only the tests). The Closing write-failure test fails on the second commit and passes after the third. A later commit reworked the write path to copy into the pipe immediately instead of queueing raw bytes in the actor while a flush is pending (same observable behavior, plus the batching above); all of these tests still pass against it.

Validation

  • dotnet build src/core/Akka/Akka.csproj -c Release -warnaserror: clean
  • Akka.Tests Akka.Tests.IO, 3 consecutive runs: 101/101 each (includes the new mid-flush batching test)
  • Akka.Streams.Tests ~Tcp: 23 passed, 3 skipped (existing Skips)
  • Akka.Remote.Tests ~Artery: 286/286
  • Akka.API.Tests: 23/23 - ITransportConnection/TcpTransportConnection gain Write(ReadOnlySequence<byte>), approved in the baseline

Benchmark

Measured before the follow-up commit, which changes only failure and close paths. Predates the later write-path rework described above (write-into-the-pipe-immediately plus batching); not re-measured since, so treat these numbers as dated.

TcpOperationsBenchmarks, short run (1 warmup, 3 iterations, 1 launch), mean ns/op. Caveat: the machine was loaded (load average about 9 on 8 cores) and the error bars are wide, so treat these as a smoke test, not a measurement.

MessageLength Clients dev this PR
10 1 3,921 4,142
10 3 1,330 1,495
10 5 817 995
10 7 627 902
10 10 534 557
10 20 466 618
10 30 565 645
10 40 576 493
100 1 4,241 4,374
100 3 1,494 1,511
100 5 891 790
100 7 661 602
100 10 662 598
100 20 564 455
100 30 524 433
100 40 696 564

With this PR, most 100-byte runs with 5 or more clients are faster and most 10-byte runs are slower. Most differences fall inside the error bars (StdDev up to 280 ns). This echo benchmark rarely goes past 64 KB unflushed, so it mostly exercises the synchronous-flush path, which doesn't allocate. It is worth a rerun on a quiet machine before merging.

These fail on dev: every Write is acked as soon as its bytes are copied
into the output pipe, even past the pipe's pause threshold.

Moves BlockingWriteStream and ConnectedSocketPair out of
TcpConnectionBatchingSpec so both specs can share them, and adds an EOF
switch to the stream.
akkadotnet#8617)

TcpConnection threw away the ValueTask<FlushResult> from WriteAsync and
acked every Write as soon as its bytes were copied into the output pipe,
so the pipe grew without bound when the socket drained slowly.

Now at most one flush is awaited at a time. If a flush doesn't complete
synchronously, its write's ack waits for it and later writes queue in
the actor (the pre-Register queue, reused) until it completes. A flush
that reports the output completed or canceled fails the write and closes
the connection. Close/ConfirmedClose/PeerClosed still enter Closing right
away, but start CloseAsync/ShutdownAsync only once the queue is drained.
ResumeWriting now replies WritingResumed once writes drain.

The output pipe keeps PipeOptions.Default (64 KB pause) unless
Inet.SO.PipeBufferSize is set, in which case it uses the same watermarks
as the input pipe (Artery sets 1 MiB).
…otnet#8617)

Should_send_ErrorClosed_to_Close_sender_When_a_queued_write_fails_while_closing
fails here: a queued write that fails while Closing drains goes through
HandleIoError, so the Close sender never gets a close event.
…t them (akkadotnet#8617)

A write that failed while Closing drained its queue called HandleIoError,
which notifies only the Register-time handler, so the Close sender (and a
ConfirmedClose-upgrading Close sender) got no close event. Both failure
branches in WriteToPipe now fail the write and tell Self FlushFailed,
whose per-behaviour handler closes the connection. FlushFailed also fails
the pending ack's write with the real cause.

Also: fold HandleGracefulClose/HandleConfirmedClose into HandleClose, make
the queued-bytes counter a long, shorten the old comments, and correct the
ledger row.
@Aaronontheweb
Aaronontheweb force-pushed the feature/tcp-write-backpressure branch from dff290a to 875534f Compare September 25, 2026 12:33
akkadotnet#8617)

The pre-Register spec checked the bytes written after Abort, which can
stop the write pump first. Check them before Abort instead.

Count queued bytes only before Register (the only place they're read),
share one FlushFailed handler between Open and PeerSentEof, and shorten
HandleStreamEof's comment.
@Aaronontheweb

Copy link
Copy Markdown
Member Author

Benchmark results: no regression

Performance looks good. Before (f923d5bec) vs. after (1a4fd3339), same machine, same session: i7-6700K (4c/8t), 16 GB, Ubuntu 24.04, .NET 10.0.201.

RemotePingPong (Artery, 5 interleaved runs per commit, msgs/sec)

Run Writes whose flush went pending Before After Change
A1, ping-pong, default config 0% – – every client count within noise
A2, one-way, default config 0% – – within ±2% (noise ~3%)
A2-sat, one-way, 16 clients, window 1000, pipe-buffer-size = 16k 1.1% 306.5K ±6.8K 301.1K ±6.7K −1.8% (within noise)
A2-sat, same, pipe-buffer-size = 4k 87% 296.9K 287.4K −3.2%
  • The pending-flush rates come from a throwaway instrumented build. At Artery's default 1 MiB pipe buffer, no flush ever goes pending in these workloads, because the Streams TCP stage coalesces about 63 messages into each write.
  • The only measurable cost is about 3%, and it appears only when the pipe buffer is shrunk to 4k so that nearly every write stalls. That config isn't realistic.

TcpOperationsBenchmarks (BenchmarkDotNet, LongRun)

All 16 cases (10 B and 100 B messages, 1–40 clients) are within ±2% and inside the error bars. Allocations are unchanged.

…dotnet#8617)

Writes now copy into the output pipe unconditionally and ride whatever
flush is next, instead of queueing raw bytes in the actor while one is
in flight. _pendingWrites/_flushWaiter are gone; _pendingAcks holds only
the ack/failure bookkeeping for writes already in the pipe, and a single
StartFlush/AckCovered/FailAllOwed path handles both the immediate and
batched-after-a-pending-flush cases. This also batches several writes
that arrive mid-flush into one follow-up flush instead of one each.

ITransportConnection gains Write(ReadOnlySequence<byte>) - a
write-only, non-flushing copy into the pipe - alongside the existing
WriteAsync/FlushAsync. Approved the new member in the API baseline.
…backpressure

# Conflicts:
#	BREAKING_CHANGES_V1.6.md
@Aaronontheweb
Aaronontheweb merged commit 5bb1f1c into akkadotnet:dev Sep 29, 2026
1 check passed
Aaronontheweb added a commit that referenced this pull request Sep 29, 2026
…eResumeWriting API (#8660)

Akka.NET copied JVM Akka's NACK write-mode messages but never implemented them. Ack-based writes are the backpressure model; see #8646 / #8617.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Akka.IO TCP: TcpConnection.EnqueueWrite discards the FlushAsync result, so the output pipe applies no backpressure

1 participant