Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
91 changes: 74 additions & 17 deletions src/ioxide/Connection/Tcp/TcpConnection.Write.Flush.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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();
}

/// <summary>
/// The teardown counterpart of <see cref="FlushAsync"/>: flush what is staged, but say nothing
/// and send nothing when a flush is already in flight.
/// </summary>
/// <remarks>
/// For the internal callers that flush on the application's behalf while tearing a connection
/// down - <see cref="ioxide.tls.TlsConnectionDualPipe.DisposeAsync"/> 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.
/// </remarks>
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++)
{
Expand Down
9 changes: 8 additions & 1 deletion src/ioxide/Tls/TlsConnectionDualPipe.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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
{
Expand Down
123 changes: 122 additions & 1 deletion tests/Ioxide.Tests.Tls/WriterContractTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Outcome>(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);
});
}

/// <summary>What the handler observed at the moment it completed the writer.</summary>
/// <summary>
/// What the handler observed at the moment it completed the writer, or disposed the pipe.
/// </summary>
private readonly record struct Outcome(bool FlushPending, long Staged, string Error);

/// <summary>
Expand Down Expand Up @@ -291,6 +328,90 @@ private static Func<Reactor, TcpConnection, Task> AbandonedFlushHandler(TaskComp
}
};

/// <summary>
/// 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.
/// </summary>
private static Func<Reactor, TcpConnection, Task> AbandonedFlushDisposalHandler(TaskCompletionSource<Outcome> report)
=> async (reactor, connection) =>
{
TlsSession? session = null;
TlsConnectionDualPipe? pipe = null;

bool reproducing = false;

try
{
session = await reactor.GetService<TlsService>()!.AcceptAsync(connection);
pipe = new TlsConnectionDualPipe(connection, session);

if (!await ReadHeadAsync(pipe.Input))
{
return;
}
reproducing = true;

byte[] chunk = new byte[ParkChunkBytes];
Task<FlushResult>? parked = null;

for (int attempt = 0; attempt < 64 && parked is null; attempt++)
{
pipe.Output.Write(chunk);
Task<FlushResult> 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);
}
};

/// <summary>
/// Reads until the request head is complete. False means the stream ended first, which is a
/// connection this suite has nothing to say about.
Expand Down
Loading