using System;
using System.IO;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using ZB.MOM.WW.MxGateway.Contracts;
using ZB.MOM.WW.MxGateway.Contracts.Proto;
using ZB.MOM.WW.MxGateway.Worker.Ipc;
using ZB.MOM.WW.MxGateway.Worker.Tests.TestSupport;
namespace ZB.MOM.WW.MxGateway.Worker.Tests.Ipc;
public sealed class WorkerFrameProtocolTests
{
private const string SessionId = "session-1";
private const string Nonce = "nonce-secret";
/// Verifies that valid envelopes round-trip through write and read.
/// A task that represents the asynchronous operation.
[Fact]
public async Task WriteAndReadAsync_WithValidEnvelope_RoundTripsFrame()
{
WorkerFrameProtocolOptions options = CreateOptions();
using MemoryStream stream = new();
WorkerEnvelope original = CreateGatewayHelloEnvelope();
WorkerFrameWriter writer = new(stream, options);
await writer.WriteAsync(original);
stream.Position = 0;
WorkerFrameReader reader = new(stream, options);
WorkerEnvelope parsed = await reader.ReadAsync();
Assert.Equal(original, parsed);
}
/// Verifies that wrong protocol version throws mismatch error.
/// A task that represents the asynchronous operation.
[Fact]
public async Task ReadAsync_WithWrongProtocolVersion_ThrowsProtocolVersionMismatch()
{
WorkerFrameProtocolOptions options = CreateOptions();
WorkerEnvelope envelope = CreateGatewayHelloEnvelope();
envelope.ProtocolVersion++;
using MemoryStream stream = new(WorkerFrameTestHelpers.CreateFrame(envelope));
WorkerFrameReader reader = new(stream, options);
WorkerFrameProtocolException exception =
await Assert.ThrowsAsync(
async () => await reader.ReadAsync());
Assert.Equal(WorkerFrameProtocolErrorCode.ProtocolVersionMismatch, exception.ErrorCode);
}
/// Verifies that wrong session ID throws mismatch error.
/// A task that represents the asynchronous operation.
[Fact]
public async Task ReadAsync_WithWrongSessionId_ThrowsSessionMismatch()
{
WorkerFrameProtocolOptions options = CreateOptions();
WorkerEnvelope envelope = CreateGatewayHelloEnvelope();
envelope.SessionId = "different-session";
using MemoryStream stream = new(WorkerFrameTestHelpers.CreateFrame(envelope));
WorkerFrameReader reader = new(stream, options);
WorkerFrameProtocolException exception =
await Assert.ThrowsAsync(
async () => await reader.ReadAsync());
Assert.Equal(WorkerFrameProtocolErrorCode.SessionMismatch, exception.ErrorCode);
}
///
/// Verifies that a frame whose length prefix is zero is rejected before the
/// payload buffer is allocated. docs/WorkerFrameProtocol.md states the
/// reader rejects zero-length payloads as a malformed-length error. The
/// length prefix is the leading four bytes of the stream, so a four-zero-byte
/// stream is exactly a frame declaring a zero-length payload.
///
/// A task that represents the asynchronous operation.
[Fact]
public async Task ReadAsync_WithZeroLengthPayload_ThrowsMalformedLength()
{
WorkerFrameProtocolOptions options = CreateOptions();
using MemoryStream stream = new(new byte[sizeof(uint)]);
WorkerFrameReader reader = new(stream, options);
WorkerFrameProtocolException exception =
await Assert.ThrowsAsync(
async () => await reader.ReadAsync());
Assert.Equal(WorkerFrameProtocolErrorCode.MalformedLength, exception.ErrorCode);
}
///
/// Verifies that a frame whose length prefix exceeds the configured maximum
/// is rejected before the payload buffer is allocated. docs/WorkerFrameProtocol.md
/// states the reader rejects oversized payloads as a message-too-large error.
/// A small maximum is configured so the rejection is asserted without
/// allocating a multi-megabyte buffer.
///
/// A task that represents the asynchronous operation.
[Fact]
public async Task ReadAsync_WithPayloadAboveConfiguredMaximum_ThrowsMessageTooLarge()
{
const int maxMessageBytes = 64;
WorkerFrameProtocolOptions options = new(
SessionId,
GatewayContractInfo.WorkerProtocolVersion,
Nonce,
maxMessageBytes);
byte[] frame = new byte[sizeof(uint)];
WorkerFrameTestHelpers.WriteUInt32LittleEndian(frame, maxMessageBytes + 1);
using MemoryStream stream = new(frame);
WorkerFrameReader reader = new(stream, options);
WorkerFrameProtocolException exception =
await Assert.ThrowsAsync(
async () => await reader.ReadAsync());
Assert.Equal(WorkerFrameProtocolErrorCode.MessageTooLarge, exception.ErrorCode);
}
/// Verifies that malformed payload throws invalid envelope error.
/// A task that represents the asynchronous operation.
[Fact]
public async Task ReadAsync_WithMalformedPayload_ThrowsInvalidEnvelope()
{
WorkerFrameProtocolOptions options = CreateOptions();
using MemoryStream stream = new(WorkerFrameTestHelpers.CreateFrame(new byte[] { 0x80 }));
WorkerFrameReader reader = new(stream, options);
WorkerFrameProtocolException exception =
await Assert.ThrowsAsync(
async () => await reader.ReadAsync());
Assert.Equal(WorkerFrameProtocolErrorCode.InvalidEnvelope, exception.ErrorCode);
}
///
/// Pins the EndOfStream branch of
/// WorkerFrameReader.ReadExactlyOrThrowAsync. The gateway
/// closing its end of the pipe during a partial-frame read is the
/// most common production transport failure; the reader must
/// surface this as WorkerFrameProtocolErrorCode.EndOfStream
/// so the worker session can fault deterministically rather than
/// spinning on a partial buffer. The stream here declares a 100-byte
/// payload but only supplies 50 bytes, so the inner read loop sees
/// bytesRead == 0 mid-frame.
///
/// A task that represents the asynchronous operation.
[Fact]
public async Task ReadAsync_WhenStreamEndsMidFrame_ThrowsEndOfStream()
{
WorkerFrameProtocolOptions options = CreateOptions();
byte[] frame = new byte[sizeof(uint) + 50];
WorkerFrameTestHelpers.WriteUInt32LittleEndian(frame, 100);
using MemoryStream stream = new(frame);
WorkerFrameReader reader = new(stream, options);
WorkerFrameProtocolException exception =
await Assert.ThrowsAsync(
async () => await reader.ReadAsync());
Assert.Equal(WorkerFrameProtocolErrorCode.EndOfStream, exception.ErrorCode);
}
///
/// Pins the writer-side MessageTooLarge branch. A session that
/// constructs an envelope whose serialised size exceeds
/// MaxMessageBytes must be rejected by the writer before any
/// bytes are sent down the pipe, so a misbehaving producer cannot
/// push the receiver past its bounds. A small MaxMessageBytes
/// is configured so a modest GatewayHello payload — with its
/// nonce padded out to several hundred bytes — exceeds the limit
/// without allocating anything large.
///
/// A task that represents the asynchronous operation.
[Fact]
public async Task WriteAsync_WithEnvelopeAboveConfiguredMaximum_ThrowsMessageTooLarge()
{
const int maxMessageBytes = 64;
WorkerFrameProtocolOptions options = new(
SessionId,
GatewayContractInfo.WorkerProtocolVersion,
Nonce,
maxMessageBytes);
using MemoryStream stream = new();
WorkerFrameWriter writer = new(stream, options);
WorkerEnvelope envelope = CreateGatewayHelloEnvelope();
envelope.GatewayHello.GatewayVersion = new string('x', 1024);
WorkerFrameProtocolException exception =
await Assert.ThrowsAsync(
async () => await writer.WriteAsync(envelope));
Assert.Equal(WorkerFrameProtocolErrorCode.MessageTooLarge, exception.ErrorCode);
Assert.Equal(0, stream.Length);
}
///
/// Documents that the writer-side InvalidEnvelope branch
/// (raised when WorkerEnvelope.CalculateSize() returns 0) is
/// unreachable through public API. WorkerEnvelopeValidator.Validate
/// (run before the size check in WorkerFrameWriter.WriteAsync)
/// rejects any envelope whose BodyCase is None with
/// InvalidEnvelope; a body-less envelope is therefore
/// intercepted before the empty-payload branch can fire. Any
/// envelope carrying a typed body serialises at least the field
/// tag bytes, so CalculateSize() is strictly positive. This
/// test exercises the body-less path and asserts the same
/// InvalidEnvelope error code reaches the caller, pinning
/// the contract that "no body" is rejected before any size check.
/// The defensive zero-length branch in WriteAsync is left
/// in place because the cost is one comparison and removing it
/// would weaken the writer against future serialisation
/// regressions; this test makes its rationale visible.
///
/// A task that represents the asynchronous operation.
[Fact]
public async Task WriteAsync_WithEmptyEnvelope_ThrowsInvalidEnvelopeFromValidator()
{
WorkerFrameProtocolOptions options = CreateOptions();
using MemoryStream stream = new();
WorkerFrameWriter writer = new(stream, options);
WorkerEnvelope envelope = new()
{
ProtocolVersion = GatewayContractInfo.WorkerProtocolVersion,
SessionId = SessionId,
Sequence = 1,
// No body — BodyCase == None, validator rejects.
};
WorkerFrameProtocolException exception =
await Assert.ThrowsAsync(
async () => await writer.WriteAsync(envelope));
Assert.Equal(WorkerFrameProtocolErrorCode.InvalidEnvelope, exception.ErrorCode);
Assert.Equal(0, stream.Length);
}
/// Verifies that concurrent writes produce complete serialized frames.
/// A task that represents the asynchronous operation.
[Fact]
public async Task WriteAsync_WithConcurrentCalls_SerializesCompleteFrames()
{
WorkerFrameProtocolOptions options = CreateOptions();
using MemoryStream stream = new();
WorkerFrameWriter writer = new(stream, options);
await Task.WhenAll(
writer.WriteAsync(CreateGatewayHelloEnvelope(sequence: 1)),
writer.WriteAsync(CreateGatewayHelloEnvelope(sequence: 2)),
writer.WriteAsync(CreateGatewayHelloEnvelope(sequence: 3)));
stream.Position = 0;
WorkerFrameReader reader = new(stream, options);
WorkerEnvelope first = await reader.ReadAsync();
WorkerEnvelope second = await reader.ReadAsync();
WorkerEnvelope third = await reader.ReadAsync();
Assert.Equal(new ulong[] { 1, 2, 3 }, new[] { first.Sequence, second.Sequence, third.Sequence }.OrderBy(sequence => sequence));
}
///
/// The reader rents its payload buffer from a shared pool, so a rented
/// buffer can be larger than the current frame and may carry bytes from
/// a previous, larger frame. Reading frames of differing sizes
/// back-to-back through one reader must parse each frame using only its
/// own payload length, never trailing pooled bytes.
///
/// A task that represents the asynchronous operation.
[Fact]
public async Task ReadAsync_WithVaryingFrameSizes_ParsesEachFrameExactly()
{
WorkerFrameProtocolOptions options = CreateOptions();
using MemoryStream stream = new();
WorkerFrameWriter writer = new(stream, options);
// A large-payload frame followed by a small-payload frame: if the
// reader reused a pooled buffer without honouring the second frame's
// length, the small frame would parse with stale trailing bytes.
WorkerEnvelope large = CreateGatewayHelloEnvelope(sequence: 1);
large.GatewayHello.GatewayVersion = new string('x', 4096);
WorkerEnvelope small = CreateGatewayHelloEnvelope(sequence: 2);
await writer.WriteAsync(large);
await writer.WriteAsync(small);
stream.Position = 0;
WorkerFrameReader reader = new(stream, options);
WorkerEnvelope firstParsed = await reader.ReadAsync();
WorkerEnvelope secondParsed = await reader.ReadAsync();
Assert.Equal(large, firstParsed);
Assert.Equal(small, secondParsed);
}
private static WorkerFrameProtocolOptions CreateOptions()
{
return new WorkerFrameProtocolOptions(
SessionId,
GatewayContractInfo.WorkerProtocolVersion,
Nonce);
}
///
/// Verifies that under concurrent writers every frame receives a distinct, gap-free sequence in
/// strictly increasing on-wire order — the sequence is stamped by the writer at write time, so the
/// wire order and the stamped sequence always agree.
///
/// A task that represents the asynchronous operation.
[Fact]
public async Task WriteAsync_UnderConcurrentCalls_StampsGapFreeMonotonicSequence()
{
const int frameCount = 50;
WorkerFrameProtocolOptions options = CreateOptions();
using MemoryStream stream = new();
WorkerFrameWriter writer = new(stream, options);
await Task.WhenAll(
Enumerable.Range(0, frameCount).Select(_ => writer.WriteAsync(CreateEventEnvelope())));
stream.Position = 0;
WorkerFrameReader reader = new(stream, options);
ulong[] sequences = new ulong[frameCount];
for (int index = 0; index < frameCount; index++)
{
sequences[index] = (await reader.ReadAsync()).Sequence;
}
// On-wire order is strictly increasing 1..frameCount with no gaps or duplicates.
Assert.Equal(Enumerable.Range(1, frameCount).Select(value => (ulong)value), sequences);
}
///
/// Verifies that when a control frame and an event frame are both queued behind an in-progress
/// write, the draining lock-holder writes the control frame first even though the event was queued
/// earlier.
///
/// A task that represents the asynchronous operation.
[Fact]
public async Task WriteAsync_WhenControlAndEventQueuedTogether_WritesControlFirst()
{
WorkerFrameProtocolOptions options = CreateOptions();
using GatedWriteStream stream = new();
WorkerFrameWriter writer = new(stream, options);
// First write occupies the writer and blocks inside the stream, holding the write lock.
Task firstWrite = writer.WriteAsync(CreateGatewayHelloEnvelope(), WorkerFrameWritePriority.Control);
await AwaitWithTimeoutAsync(stream.FirstWriteStarted);
// Queue an event first, then a control frame, while the writer is blocked. Both wait for the lock.
Task eventWrite = writer.WriteAsync(CreateEventEnvelope(), WorkerFrameWritePriority.Event);
Task controlWrite = writer.WriteAsync(CreateGatewayHelloEnvelope(), WorkerFrameWritePriority.Control);
await Task.Delay(50);
stream.ReleaseFirstWrite();
await AwaitWithTimeoutAsync(Task.WhenAll(firstWrite, eventWrite, controlWrite));
stream.Position = 0;
WorkerFrameReader reader = new(stream, options);
WorkerEnvelope frame1 = await reader.ReadAsync();
WorkerEnvelope frame2 = await reader.ReadAsync();
WorkerEnvelope frame3 = await reader.ReadAsync();
Assert.Equal(WorkerEnvelope.BodyOneofCase.GatewayHello, frame1.BodyCase);
// The control frame jumped ahead of the earlier-queued event.
Assert.Equal(WorkerEnvelope.BodyOneofCase.GatewayHello, frame2.BodyCase);
Assert.Equal(WorkerEnvelope.BodyOneofCase.WorkerEvent, frame3.BodyCase);
}
///
/// Verifies the writer coalesces the flush across a batch of frames drained together: four frames
/// queued behind an in-progress write drain in a single pass and share one FlushAsync, not four.
/// Every frame still reaches the wire intact.
///
/// A task that represents the asynchronous operation.
[Fact]
public async Task WriteAsync_WhenBatchDrainedTogether_FlushesOnce()
{
WorkerFrameProtocolOptions options = CreateOptions();
using GatedWriteStream stream = new();
WorkerFrameWriter writer = new(stream, options);
// A blocked first write occupies the writer and holds the lock while more frames queue behind it.
Task firstWrite = writer.WriteAsync(CreateGatewayHelloEnvelope(), WorkerFrameWritePriority.Control);
await AwaitWithTimeoutAsync(stream.FirstWriteStarted);
Task eventWrite1 = writer.WriteAsync(CreateEventEnvelope(), WorkerFrameWritePriority.Event);
Task eventWrite2 = writer.WriteAsync(CreateEventEnvelope(), WorkerFrameWritePriority.Event);
Task eventWrite3 = writer.WriteAsync(CreateEventEnvelope(), WorkerFrameWritePriority.Event);
await Task.Delay(50);
stream.ReleaseFirstWrite();
await AwaitWithTimeoutAsync(Task.WhenAll(firstWrite, eventWrite1, eventWrite2, eventWrite3));
// Four frames written in one drain pass => exactly one flush.
Assert.Equal(1, stream.FlushCount);
stream.Position = 0;
WorkerFrameReader reader = new(stream, options);
for (int index = 0; index < 4; index++)
{
WorkerEnvelope frame = await reader.ReadAsync();
Assert.NotEqual(WorkerEnvelope.BodyOneofCase.None, frame.BodyCase);
}
}
///
/// Verifies a per-frame rejection does not burn a sequence number (WRK-23). The sequence is a
/// diagnostic counter, so a gap breaks nothing functionally — but an operator correlating a pipe
/// capture reads a gap as a lost frame and chases a bug that does not exist, and the gap-free
/// guarantee the concurrent-write test asserts would otherwise only hold until the first
/// rejection.
///
/// A task that represents the asynchronous operation.
[Fact]
public async Task WriteAsync_PerFrameRejection_DoesNotConsumeSequence()
{
const int maxMessageBytes = 512;
WorkerFrameProtocolOptions options = new(
SessionId,
GatewayContractInfo.WorkerProtocolVersion,
Nonce,
maxMessageBytes);
using MemoryStream stream = new();
WorkerFrameWriter writer = new(stream, options);
await writer.WriteAsync(CreateEventEnvelope());
WorkerEnvelope oversized = CreateGatewayHelloEnvelope();
oversized.GatewayHello.GatewayVersion = new string('x', maxMessageBytes * 2);
WorkerFrameProtocolException exception =
await Assert.ThrowsAsync(
async () => await writer.WriteAsync(oversized));
Assert.Equal(WorkerFrameProtocolErrorCode.MessageTooLarge, exception.ErrorCode);
await writer.WriteAsync(CreateEventEnvelope());
stream.Position = 0;
WorkerFrameReader reader = new(stream, options);
WorkerEnvelope first = await reader.ReadAsync();
WorkerEnvelope second = await reader.ReadAsync();
// Two frames reached the wire; the rejected frame in between left no gap.
Assert.Equal(1UL, first.Sequence);
Assert.Equal(2UL, second.Sequence);
Assert.Equal(stream.Length, stream.Position);
}
/// Verifies a zero negotiated frame maximum keeps the constructor default.
[Fact]
public void AdoptNegotiatedMaxMessageBytes_WithZero_KeepsDefault()
{
WorkerFrameProtocolOptions options = CreateOptions();
int original = options.MaxMessageBytes;
options.AdoptNegotiatedMaxMessageBytes(0);
Assert.Equal(original, options.MaxMessageBytes);
}
/// Verifies an in-range negotiated frame maximum is adopted.
[Fact]
public void AdoptNegotiatedMaxMessageBytes_WithInRangeValue_Adopts()
{
WorkerFrameProtocolOptions options = CreateOptions();
options.AdoptNegotiatedMaxMessageBytes(4 * 1024 * 1024);
Assert.Equal(4 * 1024 * 1024, options.MaxMessageBytes);
}
/// Verifies a negotiated frame maximum above the worker ceiling is rejected.
[Fact]
public void AdoptNegotiatedMaxMessageBytes_AboveCeiling_Throws()
{
WorkerFrameProtocolOptions options = CreateOptions();
WorkerFrameProtocolException exception = Assert.Throws(
() => options.AdoptNegotiatedMaxMessageBytes((uint)WorkerFrameProtocolOptions.MaxNegotiableFrameBytes + 1));
Assert.Equal(WorkerFrameProtocolErrorCode.InvalidConfiguration, exception.ErrorCode);
}
///
/// Verifies a negotiated frame maximum below the worker floor is rejected as
/// InvalidConfiguration (WRK-24), and that exactly the floor is adopted. The floor matches
/// the gateway's own GatewayOptionsValidator.MinimumMaxMessageBytes, so the worker never
/// rejects a value the gateway's validator accepts, yet a nonsensical tiny value faults at the
/// handshake instead of leaving a session that fails every later frame.
///
[Fact]
public void AdoptNegotiatedMaxMessageBytes_BelowFloor_ThrowsInvalidConfiguration()
{
WorkerFrameProtocolOptions belowFloor = CreateOptions();
WorkerFrameProtocolException exception = Assert.Throws(
() => belowFloor.AdoptNegotiatedMaxMessageBytes(512));
Assert.Equal(WorkerFrameProtocolErrorCode.InvalidConfiguration, exception.ErrorCode);
// Boundary: exactly the floor is accepted.
WorkerFrameProtocolOptions atFloor = CreateOptions();
atFloor.AdoptNegotiatedMaxMessageBytes((uint)WorkerFrameProtocolOptions.MinNegotiableFrameBytes);
Assert.Equal(WorkerFrameProtocolOptions.MinNegotiableFrameBytes, atFloor.MaxMessageBytes);
}
///
/// WRK-22 / IPC-26. A WriteAsync cancelled while it waits for the write lock must never
/// have its frame written by the next lock-holder. Writer A holds the lock mid-write (blocked in
/// the stream); an event write is queued and then cancelled; when A is released and a later
/// control frame drains, the wire carries A's frame and the control frame only — the cancelled
/// event envelope is tombstoned and skipped.
///
/// A task that represents the asynchronous operation.
[Fact]
public async Task WriteAsync_CancelledWhileWaitingForLock_FrameIsNeverWritten()
{
WorkerFrameProtocolOptions options = CreateOptions();
using GatedWriteStream stream = new();
WorkerFrameWriter writer = new(stream, options);
// Writer A occupies the writer and blocks inside the stream, holding the write lock.
Task firstWrite = writer.WriteAsync(CreateGatewayHelloEnvelope(), WorkerFrameWritePriority.Control);
await AwaitWithTimeoutAsync(stream.FirstWriteStarted);
// Queue an event write with its own CTS while the lock is held, then cancel it.
using CancellationTokenSource cts = new();
Task cancelledWrite = writer.WriteAsync(CreateEventEnvelope(), WorkerFrameWritePriority.Event, cts.Token);
await Task.Delay(50);
cts.Cancel();
await Assert.ThrowsAnyAsync(async () => await cancelledWrite);
// Release A, then drive a fresh control write.
stream.ReleaseFirstWrite();
await AwaitWithTimeoutAsync(firstWrite);
await writer.WriteAsync(CreateGatewayHelloEnvelope(), WorkerFrameWritePriority.Control);
stream.Position = 0;
WorkerFrameReader reader = new(stream, options);
WorkerEnvelope frame1 = await reader.ReadAsync();
WorkerEnvelope frame2 = await reader.ReadAsync();
Assert.Equal(WorkerEnvelope.BodyOneofCase.GatewayHello, frame1.BodyCase);
Assert.Equal(WorkerEnvelope.BodyOneofCase.GatewayHello, frame2.BodyCase);
// The cancelled event never reached the wire — no third frame, and sequences stay contiguous.
Assert.Equal(stream.Length, stream.Position);
Assert.Equal(1UL, frame1.Sequence);
Assert.Equal(2UL, frame2.Sequence);
}
///
/// NEXT-04. A frame claimed by the draining lock-holder before its caller's cancellation lands
/// is abandoned — the cancelled caller never awaits its completion. If the wire write then
/// faults, the tombstone path's fault-observing continuation must still observe the exception
/// so it never surfaces as .
///
/// A task that represents the asynchronous operation.
[Fact]
public async Task WriteAsync_ClaimedFrameAbandonedByCancellation_FaultIsObserved()
{
string marker = $"NEXT-04-{Guid.NewGuid():N}";
WorkerFrameProtocolOptions options = CreateOptions();
bool sawUnobservedMarkerFault = false;
EventHandler handler = (sender, args) =>
{
if (args.Exception.ToString().Contains(marker))
{
sawUnobservedMarkerFault = true;
}
};
TaskScheduler.UnobservedTaskException += handler;
try
{
using (SecondWriteFaultingGatedStream stream = new SecondWriteFaultingGatedStream(marker))
{
WorkerFrameWriter writer = new WorkerFrameWriter(stream, options);
// Writer A holds the lock, blocked mid-write of its own frame.
Task firstWrite = writer.WriteAsync(CreateGatewayHelloEnvelope(), WorkerFrameWritePriority.Control);
await AwaitWithTimeoutAsync(stream.FirstWriteStarted);
using (CancellationTokenSource cts = new CancellationTokenSource())
{
// Queue the doomed event write behind A, release A so its drain claims the
// event frame and blocks mid-write of it, then cancel the queued caller —
// the frame is claimed, so the caller unwinds without an awaiter for it.
Task abandonedWrite = writer.WriteAsync(CreateEventEnvelope(), WorkerFrameWritePriority.Event, cts.Token);
stream.ReleaseFirstWrite();
await AwaitWithTimeoutAsync(stream.SecondWriteStarted);
cts.Cancel();
await Assert.ThrowsAnyAsync(async () => await abandonedWrite);
}
// Fault the abandoned frame's wire write; observe writer A's own outcome so only
// the abandoned frame's completion could ever raise the marker unobserved.
stream.ReleaseSecondWrite();
_ = await Record.ExceptionAsync(async () => await firstWrite);
}
GC.Collect();
GC.WaitForPendingFinalizers();
GC.Collect();
}
finally
{
TaskScheduler.UnobservedTaskException -= handler;
}
Assert.False(
sawUnobservedMarkerFault,
"The abandoned frame's write fault surfaced as an unobserved-task exception.");
}
///
/// WRK-22 / IPC-26, the review's shutdown scenario. A cancelled event frame queued before a
/// shutdown-ack control frame must not trail the ack on the wire: the tombstone rule plus the
/// control-before-event scheduler keeps the ack the last frame written.
///
/// A task that represents the asynchronous operation.
[Fact]
public async Task WriteAsync_CancelledEventFrame_DoesNotTrailShutdownAck()
{
WorkerFrameProtocolOptions options = CreateOptions();
using GatedWriteStream stream = new();
WorkerFrameWriter writer = new(stream, options);
Task firstWrite = writer.WriteAsync(CreateGatewayHelloEnvelope(), WorkerFrameWritePriority.Control);
await AwaitWithTimeoutAsync(stream.FirstWriteStarted);
using CancellationTokenSource cts = new();
Task cancelledEvent = writer.WriteAsync(CreateEventEnvelope(), WorkerFrameWritePriority.Event, cts.Token);
await Task.Delay(50);
cts.Cancel();
await Assert.ThrowsAnyAsync(async () => await cancelledEvent);
// The shutdown ack (a control frame) is queued behind the still-blocked first write.
Task ackWrite = writer.WriteAsync(CreateShutdownAckEnvelope(), WorkerFrameWritePriority.Control);
await Task.Delay(50);
stream.ReleaseFirstWrite();
await AwaitWithTimeoutAsync(Task.WhenAll(firstWrite, ackWrite));
stream.Position = 0;
WorkerFrameReader reader = new(stream, options);
WorkerEnvelope frame1 = await reader.ReadAsync();
WorkerEnvelope frame2 = await reader.ReadAsync();
Assert.Equal(WorkerEnvelope.BodyOneofCase.GatewayHello, frame1.BodyCase);
// The ack is the last frame — the cancelled event did not trail it.
Assert.Equal(WorkerEnvelope.BodyOneofCase.WorkerShutdownAck, frame2.BodyCase);
Assert.Equal(stream.Length, stream.Position);
}
///
/// WRK-25. The batch entry point enqueues a whole event burst under one lock acquisition and
/// drains it together, so N events cost exactly one flush and reach the wire in batch order with
/// monotonic sequences — the coalescing WRK-12 shipped, now on the event hot path.
///
/// A task that represents the asynchronous operation.
[Fact]
public async Task WriteBatchAsync_FlushesOnceAndPreservesOrder()
{
const int count = 8;
WorkerFrameProtocolOptions options = CreateOptions();
using FlushCountingStream stream = new();
WorkerFrameWriter writer = new(stream, options);
WorkerEnvelope[] batch = new WorkerEnvelope[count];
for (int index = 0; index < count; index++)
{
batch[index] = CreateEventEnvelope(workerSequence: (ulong)(100 + index));
}
await writer.WriteBatchAsync(batch, WorkerFrameWritePriority.Event);
// The whole batch was queued before the single lock wait, so it drained in one pass => one flush.
Assert.Equal(1, stream.FlushCount);
stream.Position = 0;
WorkerFrameReader reader = new(stream, options);
for (int index = 0; index < count; index++)
{
WorkerEnvelope frame = await reader.ReadAsync();
Assert.Equal(WorkerEnvelope.BodyOneofCase.WorkerEvent, frame.BodyCase);
// Wire order matches batch order.
Assert.Equal((ulong)(100 + index), frame.WorkerEvent.Event.WorkerSequence);
// Write-time stamped sequence is monotonic 1..count.
Assert.Equal((ulong)(index + 1), frame.Sequence);
}
}
///
/// WRK-25. Control-before-event still holds mid-batch: a control frame queued while a batch is
/// draining jumps ahead of the batch's remaining events.
///
/// A task that represents the asynchronous operation.
[Fact]
public async Task WriteBatchAsync_ControlFrameQueuedDuringBatch_JumpsRemainingEvents()
{
WorkerFrameProtocolOptions options = CreateOptions();
using GatedWriteStream stream = new();
WorkerFrameWriter writer = new(stream, options);
WorkerEnvelope[] batch = new[]
{
CreateEventEnvelope(),
CreateEventEnvelope(),
CreateEventEnvelope(),
};
// The batch takes the lock and blocks writing its first event frame inside the stream.
Task batchWrite = writer.WriteBatchAsync(batch, WorkerFrameWritePriority.Event);
await AwaitWithTimeoutAsync(stream.FirstWriteStarted);
// A control frame queued mid-drain must jump the batch's remaining events.
Task controlWrite = writer.WriteAsync(CreateGatewayHelloEnvelope(), WorkerFrameWritePriority.Control);
await Task.Delay(50);
stream.ReleaseFirstWrite();
await AwaitWithTimeoutAsync(Task.WhenAll(batchWrite, controlWrite));
stream.Position = 0;
WorkerFrameReader reader = new(stream, options);
WorkerEnvelope f1 = await reader.ReadAsync();
WorkerEnvelope f2 = await reader.ReadAsync();
WorkerEnvelope f3 = await reader.ReadAsync();
WorkerEnvelope f4 = await reader.ReadAsync();
// First event was already writing when the control frame queued; the control frame then jumps
// ahead of the two remaining events.
Assert.Equal(WorkerEnvelope.BodyOneofCase.WorkerEvent, f1.BodyCase);
Assert.Equal(WorkerEnvelope.BodyOneofCase.GatewayHello, f2.BodyCase);
Assert.Equal(WorkerEnvelope.BodyOneofCase.WorkerEvent, f3.BodyCase);
Assert.Equal(WorkerEnvelope.BodyOneofCase.WorkerEvent, f4.BodyCase);
}
private static WorkerEnvelope CreateGatewayHelloEnvelope(ulong sequence = 1)
{
return new WorkerEnvelope
{
ProtocolVersion = GatewayContractInfo.WorkerProtocolVersion,
SessionId = SessionId,
Sequence = sequence,
GatewayHello = new GatewayHello
{
SupportedProtocolVersion = GatewayContractInfo.WorkerProtocolVersion,
Nonce = Nonce,
GatewayVersion = "test-gateway",
},
};
}
// net48 has no Task.WaitAsync(TimeSpan); fail the test rather than hang if the writer misbehaves.
private static async Task AwaitWithTimeoutAsync(Task task)
{
Task completed = await Task.WhenAny(task, Task.Delay(TimeSpan.FromSeconds(5)));
if (completed != task)
{
throw new TimeoutException("Timed out waiting for the frame writer.");
}
await task;
}
private static WorkerEnvelope CreateEventEnvelope()
{
return new WorkerEnvelope
{
ProtocolVersion = GatewayContractInfo.WorkerProtocolVersion,
SessionId = SessionId,
WorkerEvent = new WorkerEvent
{
Event = new MxEvent { SessionId = SessionId },
},
};
}
private static WorkerEnvelope CreateEventEnvelope(ulong workerSequence)
{
WorkerEnvelope envelope = CreateEventEnvelope();
envelope.WorkerEvent.Event.WorkerSequence = workerSequence;
return envelope;
}
private static WorkerEnvelope CreateShutdownAckEnvelope()
{
return new WorkerEnvelope
{
ProtocolVersion = GatewayContractInfo.WorkerProtocolVersion,
SessionId = SessionId,
WorkerShutdownAck = new WorkerShutdownAck
{
Status = new ProtocolStatus
{
Code = ProtocolStatusCode.Ok,
Message = "OK",
},
},
};
}
// A MemoryStream that counts FlushAsync calls without gating any write, so a batch write can be
// asserted to flush exactly once.
private sealed class FlushCountingStream : MemoryStream
{
private int _flushCount;
/// Gets the number of calls observed so far.
public int FlushCount => Volatile.Read(ref _flushCount);
///
public override Task FlushAsync(CancellationToken cancellationToken)
{
Interlocked.Increment(ref _flushCount);
return base.FlushAsync(cancellationToken);
}
}
// A MemoryStream whose first write blocks until released and whose second write blocks until
// released and then throws, so a test can abandon a claimed frame by cancellation and fault its
// wire write afterwards (NEXT-04).
private sealed class SecondWriteFaultingGatedStream : MemoryStream
{
private readonly SemaphoreSlim _firstRelease = new SemaphoreSlim(0);
private readonly SemaphoreSlim _secondRelease = new SemaphoreSlim(0);
private readonly TaskCompletionSource _firstWriteStarted =
new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
private readonly TaskCompletionSource _secondWriteStarted =
new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
private readonly string _faultMessage;
private int _writeCount;
public SecondWriteFaultingGatedStream(string faultMessage)
{
_faultMessage = faultMessage;
}
/// Gets a task that completes once the first call has started blocking.
public Task FirstWriteStarted => _firstWriteStarted.Task;
/// Gets a task that completes once the second call has started blocking.
public Task SecondWriteStarted => _secondWriteStarted.Task;
/// Releases the first blocked write so it can complete.
public void ReleaseFirstWrite() => _firstRelease.Release();
/// Releases the second blocked write so it can throw.
public void ReleaseSecondWrite() => _secondRelease.Release();
///
public override async Task WriteAsync(byte[] buffer, int offset, int count, CancellationToken cancellationToken)
{
int writeIndex = Interlocked.Increment(ref _writeCount);
if (writeIndex == 1)
{
_firstWriteStarted.TrySetResult(true);
await _firstRelease.WaitAsync(cancellationToken);
}
else if (writeIndex == 2)
{
_secondWriteStarted.TrySetResult(true);
await _secondRelease.WaitAsync(cancellationToken);
throw new IOException(_faultMessage);
}
await base.WriteAsync(buffer, offset, count, cancellationToken);
}
///
protected override void Dispose(bool disposing)
{
if (disposing)
{
_firstRelease.Dispose();
_secondRelease.Dispose();
}
base.Dispose(disposing);
}
}
// A MemoryStream whose first WriteAsync blocks until released, so a test can queue additional frames
// behind an in-progress write and observe the writer's priority ordering.
private sealed class GatedWriteStream : MemoryStream
{
private readonly SemaphoreSlim _release = new SemaphoreSlim(0);
private readonly TaskCompletionSource _firstWriteStarted =
new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
private int _writeCount;
private int _flushCount;
/// Gets a task that completes once the first call has started blocking.
public Task FirstWriteStarted => _firstWriteStarted.Task;
/// Gets the number of calls observed so far.
public int FlushCount => Volatile.Read(ref _flushCount);
/// Releases the first blocked write so it can complete.
public void ReleaseFirstWrite() => _release.Release();
///
public override async Task WriteAsync(byte[] buffer, int offset, int count, CancellationToken cancellationToken)
{
if (Interlocked.Increment(ref _writeCount) == 1)
{
_firstWriteStarted.TrySetResult(true);
await _release.WaitAsync(cancellationToken);
}
await base.WriteAsync(buffer, offset, count, cancellationToken);
}
///
public override Task FlushAsync(CancellationToken cancellationToken)
{
Interlocked.Increment(ref _flushCount);
return base.FlushAsync(cancellationToken);
}
///
protected override void Dispose(bool disposing)
{
if (disposing)
{
_release.Dispose();
}
base.Dispose(disposing);
}
}
}