From 8328dc2b68dd7e33e4bfdcd0d6e4a32f34314819 Mon Sep 17 00:00:00 2001 From: Diogo Martins Date: Sat, 19 Sep 2026 11:17:10 +0100 Subject: [PATCH] tls: teardown skips its flush when the application left one in flight TlsConnectionDualPipe.DisposeAsync flushes the connection unconditionally, because a handler that wrote its response and never flushed has nothing else to carry it out. Against a flush the application left in flight that call met TcpConnection.FlushAsync's one-flush-at-a-time guard and threw out of a finally, surfacing as "[r0] connection handler faulted: FlushAsync already in progress." - reported on a long-lived TLS websocket, a few times per run (#234). The guard is right and stays. The parked caller waits on _flushSignal, a single ManualResetValueTaskSourceCore that holds one continuation, so a second flush is not something the connection can serve; an application that flushes twice concurrently needs to hear about it. What is wrong is the caller: disposal is the connection's own flush, not the application's, and it has no business tripping a guard aimed at application misuse. So FlushAsync keeps the throw and grows FlushIfIdleAsync beside it, which returns instead. Skipping sends nothing because there is nothing to send: while _flushInProgress is set every GetSpan/GetMemory/Advance/Write is refused, so nothing can have entered the slab since the in-flight flush armed, and that flush snapshotted everything already there. The in-flight send owns the whole slab; a second flush would compute target == 0 and return anyway. The regression test parks a real flush the way the existing contract test does - a peer that stops draining - and stages nothing before tearing down, which is the production shape: every application write was flushed, so Complete has nothing to commit and the disposal's own flush is the only call left that can throw. Without the fix it fails with the reported message. Not covered: the same disposal throws a line earlier, out of Complete, when plaintext IS staged - GetSpan refuses it and Complete catches only IOException. That one is already tracked as the pending test above this one, and needs its own call about whether those bytes are dropped or waited for, since anything appended to the slab during an in-flight flush is zeroed by CompleteFlush. TLS suite 142 passed / 0 failed / 7 pending; Unit, Http, Chaos, File and E2E green. Tls/OpenSsl -0.6%, Tls/OpenSslPipes +1.6% at 2 reactors - inside noise. --- .../Tcp/TcpConnection.Write.Flush.cs | 91 ++++++++++--- src/ioxide/Tls/TlsConnectionDualPipe.cs | 9 +- tests/Ioxide.Tests.Tls/WriterContractTests.cs | 123 +++++++++++++++++- 3 files changed, 204 insertions(+), 19 deletions(-) diff --git a/src/ioxide/Connection/Tcp/TcpConnection.Write.Flush.cs b/src/ioxide/Connection/Tcp/TcpConnection.Write.Flush.cs index a1f1d821..060e3d2d 100644 --- a/src/ioxide/Connection/Tcp/TcpConnection.Write.Flush.cs +++ b/src/ioxide/Connection/Tcp/TcpConnection.Write.Flush.cs @@ -20,33 +20,90 @@ public sealed unsafe partial class TcpConnection : IValueTaskSource public ValueTask FlushAsync() { - // TcpConnection already torn down: complete immediately so the handler unwinds - // to its next ReadAsync, sees IsClosed, and exits. - // - // The staged bytes are dropped rather than kept. Nothing will ever send them, and leaving - // the tail where it was made made every later write append to a slab that only ever grew - - // a writer that keeps producing against a peer that has gone reaches gigabytes, because - // doubling the slab is the one thing that never fails. - // - // Segmented needs its own release, and zeroing the tail alone did nothing for it: once - // _inOverflow is set every write is routed to overflow rather than to the slab, so the - // growth continued one pooled slab per fill. Clear() has always released them on recycle; - // this is the same release on the path that never reaches recycle. if (Volatile.Read(ref _closed) == 1) { - WriteTail = 0; - if (_ovCount > 0) - { - ReleaseOverflow(); - } + DropStaged(); + return default; } + // One flush at a time, and not as a matter of taste: the parked caller waits on + // _flushSignal, a single ManualResetValueTaskSourceCore with room for exactly one + // continuation. A second flush would hand a second awaiter the same token and the source + // would throw on the second OnCompleted instead - deeper, and on the reactor thread. + // An application that flushes twice concurrently has a bug in its own write serialization, + // and hearing about it here is the point. if (Interlocked.Exchange(ref _flushInProgress, 1) == 1) { throw new InvalidOperationException("FlushAsync already in progress."); } + return FlushCore(); + } + + /// + /// The teardown counterpart of : flush what is staged, but say nothing + /// and send nothing when a flush is already in flight. + /// + /// + /// For the internal callers that flush on the application's behalf while tearing a connection + /// down - is the one that has to + /// have it. They flush unconditionally because a handler that wrote its response and never + /// flushed has nothing else to carry it out, and against a flush the application left in + /// flight that unconditional call threw out of a finally and faulted the connection handler + /// (#234). + /// + /// Skipping loses nothing, and that is a property of the write path rather than a hope: while + /// _flushInProgress is set every GetSpan/GetMemory/Advance/Write is refused + /// (TcpConnection.Write.cs), so nothing can have entered the slab since the in-flight flush + /// armed, and that flush snapshotted everything that was already there. The slab holds exactly + /// what is on its way out. A second flush would compute target == 0 and return anyway. + /// + /// Not the same as relaxing the guard above. A caller that means "send my bytes" must still + /// hear that they were not sent. + /// + internal ValueTask FlushIfIdleAsync() + { + if (Volatile.Read(ref _closed) == 1) + { + DropStaged(); + + return default; + } + + if (Interlocked.Exchange(ref _flushInProgress, 1) == 1) + { + return default; + } + + return FlushCore(); + } + + // TcpConnection already torn down: complete immediately so the handler unwinds + // to its next ReadAsync, sees IsClosed, and exits. + // + // The staged bytes are dropped rather than kept. Nothing will ever send them, and leaving + // the tail where it was made made every later write append to a slab that only ever grew - + // a writer that keeps producing against a peer that has gone reaches gigabytes, because + // doubling the slab is the one thing that never fails. + // + // Segmented needs its own release, and zeroing the tail alone did nothing for it: once + // _inOverflow is set every write is routed to overflow rather than to the slab, so the + // growth continued one pooled slab per fill. Clear() has always released them on recycle; + // this is the same release on the path that never reaches recycle. + private void DropStaged() + { + WriteTail = 0; + if (_ovCount > 0) + { + ReleaseOverflow(); + } + } + + // Both entry points past their guards: whoever gets here owns _flushInProgress and is the one + // flush the connection is allowed to have outstanding. + private ValueTask FlushCore() + { int target = WriteTail; for (int i = 0; i < _ovCount; i++) { diff --git a/src/ioxide/Tls/TlsConnectionDualPipe.cs b/src/ioxide/Tls/TlsConnectionDualPipe.cs index 13d36f4c..042080b1 100644 --- a/src/ioxide/Tls/TlsConnectionDualPipe.cs +++ b/src/ioxide/Tls/TlsConnectionDualPipe.cs @@ -99,10 +99,17 @@ public async ValueTask DisposeAsync() // unflushed plaintext into the slab, and the flush carries it out. The read side's disposal // marks the connection closed (nothing else releases a pump parked on a quiet peer), and // after that a flush is a no-op - so this order is load-bearing, not stylistic. + // + // FlushIfIdleAsync rather than FlushAsync, because this flush is the connection's, not the + // application's. A handler that left a flush in flight - wrote, stopped waiting on a peer + // that was not draining, and tore down - met the one-flush-at-a-time guard here and had + // the InvalidOperationException come out of its teardown as a faulted connection handler + // (#234). There is nothing to carry out in that state anyway: the in-flight send owns the + // whole slab, and the write guards have refused everything since it armed. try { _writer.Complete(); - await _conn.FlushAsync(); + await _conn.FlushIfIdleAsync(); } finally { diff --git a/tests/Ioxide.Tests.Tls/WriterContractTests.cs b/tests/Ioxide.Tests.Tls/WriterContractTests.cs index dfc7c970..63f2b787 100644 --- a/tests/Ioxide.Tests.Tls/WriterContractTests.cs +++ b/tests/Ioxide.Tests.Tls/WriterContractTests.cs @@ -107,9 +107,46 @@ public static void Register(Runner runner) "Complete threw out of a teardown that may not fail: " + outcome.Error); }, "Complete catches only IOException, but its commit reaches TcpConnection.GetSpan, which " + "throws InvalidOperationException(\"Cannot write while flush is in progress\")"); + + // The other half of the same two-line teardown, and the half that reached production: a + // handler that flushed everything it wrote leaves Complete nothing to commit, so it returns + // quietly and DisposeAsync's OWN flush is what meets the connection's one-flush-at-a-time + // guard. Reported as a websocket over TLS faulting with "FlushAsync already in progress" + // a few times per run (#234), where the throw escapes as a faulted connection handler. + // + // Same park as above, because the state is the same one: what differs is only that nothing + // is staged before the teardown. + runner.Test("tls writer: disposal does not throw while the connection's flush is in flight", () => + { + (string certPath, string keyPath) = TestCert.Ensure(); + var options = new TlsOptions { CertificatePath = certPath, KeyPath = keyPath }; + + var report = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + int port = TestServer.Start(AbandonedFlushDisposalHandler(report), r => TlsService.Start(r, options)); + + RequestAndStopReading(port, report.Task); + + Assert.True(report.Task.Wait(TimeSpan.FromSeconds(5)), + "the handler never reached the disposal under test"); + + Outcome outcome = report.Task.Result; + + // Non-vacuous in both directions: a flush genuinely still in flight, and nothing staged. + // Without the first there is no second flush to refuse; with the second this would be + // measuring the pending test above, which throws a line earlier and for another reason. + Assert.True(outcome.FlushPending, + $"no flush was in flight when the pipe was disposed: {outcome.Error}"); + Assert.True(outcome.Staged == 0, + $"plaintext was staged, so Complete threw before the flush under test: {outcome.Staged} B"); + + Assert.True(outcome.Error.Length == 0, + "disposal threw against a flush the application left in flight: " + outcome.Error); + }); } - /// What the handler observed at the moment it completed the writer. + /// + /// What the handler observed at the moment it completed the writer, or disposed the pipe. + /// private readonly record struct Outcome(bool FlushPending, long Staged, string Error); /// @@ -291,6 +328,90 @@ private static Func AbandonedFlushHandler(TaskComp } }; + /// + /// The handler of the production report: it writes until one flush stops coming back, gives up + /// on it, and tears the connection down with nothing staged - so the disposal's own flush is + /// the only call left that can throw. Reports whether a flush was still in flight, that nothing + /// was staged, and whatever the disposal threw. + /// + private static Func AbandonedFlushDisposalHandler(TaskCompletionSource report) + => async (reactor, connection) => + { + TlsSession? session = null; + TlsConnectionDualPipe? pipe = null; + + bool reproducing = false; + + try + { + session = await reactor.GetService()!.AcceptAsync(connection); + pipe = new TlsConnectionDualPipe(connection, session); + + if (!await ReadHeadAsync(pipe.Input)) + { + return; + } + reproducing = true; + + byte[] chunk = new byte[ParkChunkBytes]; + Task? parked = null; + + for (int attempt = 0; attempt < 64 && parked is null; attempt++) + { + pipe.Output.Write(chunk); + Task flush = pipe.Output.FlushAsync().AsTask(); + + if (await Task.WhenAny(flush, Task.Delay(TimeSpan.FromSeconds(2))) != flush) + { + parked = flush; + } + } + + if (parked is null) + { + report.TrySetResult(new Outcome(false, 0, + "the peer drained 16 MB; no flush ever stayed in flight")); + return; + } + + // Deliberately nothing written here, unlike the test above: everything this handler + // produced went into the flush that parked, which is the shape the report describes + // - the application's writes were serialized and all of them were flushed. + long staged = pipe.Output.UnflushedBytes; + + // Read on the reactor thread with nothing awaited before the disposal, so what is + // reported is the state the disposal actually ran against. + bool flushPending = !parked.IsCompleted; + + TlsConnectionDualPipe disposing = pipe; + pipe = null; // disposed here; the finally must not do it a second time + + string error = ""; + try + { + await disposing.DisposeAsync(); + } + catch (Exception e) + { + error = $"{e.GetType().Name}: {e.Message}"; + } + + report.TrySetResult(new Outcome(flushPending, staged, error)); + } + catch (Exception e) + { + if (reproducing) + { + report.TrySetResult(new Outcome(false, 0, + $"the reproduction broke before the disposal: {e.GetType().Name}: {e.Message}")); + } + } + finally + { + await ReleaseAsync(pipe, session, connection); + } + }; + /// /// Reads until the request head is complete. False means the stream ended first, which is a /// connection this suite has nothing to say about.