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.