feat(gateway): decouple worker event enqueue from the read loop (GWC-04)
The read loop awaited EnqueueWorkerEventAsync inline, which blocks in the bounded event channel's timed WriteAsync (up to EventChannelFullModeTimeout, default 5s) when the channel is full with no/slow consumer. A WorkerCommandReply or heartbeat queued behind an event frame was then stalled, so an in-flight InvokeAsync could hit CommandTimeout even though the worker replied in time. Mirror the existing outbound WriteLoopAsync: add an unbounded event staging channel and a dedicated EventWriteLoopAsync. DispatchEnvelope is now fully synchronous — the WorkerEvent branch hands the event off with a non-blocking TryWrite and the read loop never awaits. The event write loop owns the timed WriteAsync into the bounded channel and the sustained-overflow ProtocolViolation fault (unchanged contract). Registered in WaitForBackgroundTasks + completed on close/fault/dispose. Test: reply arriving after events with a full, consumer-less event channel is dispatched promptly (no CommandTimeout) — pipe-harness, verified on windev. Docs: GatewayProcessDesign read/write/event-loop section.
This commit is contained in:
@@ -218,6 +218,51 @@ public sealed class WorkerClientTests
|
||||
Assert.Equal(WorkerClientState.Faulted, client.State);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Verifies that a command reply arriving on the pipe after events is dispatched promptly even
|
||||
/// when the event channel is full and has no consumer — event enqueue is decoupled from the read
|
||||
/// loop, so a blocked event writer cannot delay a reply (GWC-04). The event full-mode timeout is
|
||||
/// set far above the command timeout: without the decoupling the read loop would block behind the
|
||||
/// full event channel and the in-flight InvokeAsync would hit CommandTimeout.
|
||||
/// </summary>
|
||||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||||
[Fact]
|
||||
public async Task ReadLoop_WhenEventChannelFull_DispatchesCommandReplyWithoutBlocking()
|
||||
{
|
||||
await using PipePair pipePair = await PipePair.CreateAsync();
|
||||
await using WorkerClient client = CreateClient(
|
||||
pipePair,
|
||||
new WorkerClientOptions
|
||||
{
|
||||
EventChannelCapacity = 1,
|
||||
// Much larger than TestTimeout: the read loop must not block for this long behind the
|
||||
// full event channel while a reply waits.
|
||||
EventChannelFullModeTimeout = TimeSpan.FromSeconds(30),
|
||||
HeartbeatGrace = TimeSpan.FromSeconds(60),
|
||||
HeartbeatCheckInterval = TimeSpan.FromSeconds(60),
|
||||
});
|
||||
await CompleteHandshakeAsync(client, pipePair);
|
||||
|
||||
// No StreamEvents consumer is attached, so the capacity-1 event channel fills and the event
|
||||
// writer blocks on the second event's timed WriteAsync.
|
||||
await pipePair.WorkerWriter.WriteAsync(
|
||||
CreateEventEnvelope(sequence: 11, MxEventFamily.OnDataChange));
|
||||
await pipePair.WorkerWriter.WriteAsync(
|
||||
CreateEventEnvelope(sequence: 12, MxEventFamily.OnDataChange));
|
||||
|
||||
Task<WorkerCommandReply> invokeTask = client.InvokeAsync(
|
||||
CreateCommand(MxCommandKind.Ping),
|
||||
TestTimeout,
|
||||
CancellationToken.None);
|
||||
WorkerEnvelope commandEnvelope = await pipePair.WorkerReader.ReadAsync().AsTask().WaitAsync(TestTimeout);
|
||||
await pipePair.WorkerWriter.WriteAsync(
|
||||
CreateCommandReplyEnvelope(commandEnvelope.CorrelationId, MxCommandKind.Ping));
|
||||
|
||||
WorkerCommandReply reply = await invokeTask.WaitAsync(TestTimeout);
|
||||
Assert.Equal(MxCommandKind.Ping, reply.Reply.Kind);
|
||||
Assert.Equal(WorkerClientState.Ready, client.State);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Verifies that when the client faults it kills the owned worker process.
|
||||
/// The assertion waits on <see cref="FakeWorkerProcess.WaitForExitAsync"/>, which
|
||||
|
||||
Reference in New Issue
Block a user