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.