Merge branch 'fix/gwc-28-29-30-polish'
# Conflicts: # archreview/2026-07-12/remediation/00-tracking.md
This commit is contained in:
@@ -116,9 +116,12 @@ public sealed class MxAccessGatewayService(
|
||||
return bulkConstraintPlan.CreateDeniedReply(request);
|
||||
}
|
||||
|
||||
MxCommandRequest invokeRequest = request.Clone();
|
||||
invokeRequest.Command = commandToInvoke;
|
||||
WorkerCommand workerCommand = mapper.MapCommand(invokeRequest);
|
||||
// Map from the command alone: cloning the whole request only to overwrite its command with
|
||||
// commandToInvoke deep-cloned the (potentially large) original payload for nothing, since
|
||||
// MapCommand reads nothing but the command (GWC-29). The one clone that matters still
|
||||
// happens inside MapCommand, which is what keeps the worker-bound graph unaliased from
|
||||
// commandToInvoke — the caller still reads it below via TrackCommandReply.
|
||||
WorkerCommand workerCommand = mapper.MapCommand(commandToInvoke);
|
||||
WorkerCommandReply workerReply = await sessionManager
|
||||
.InvokeAsync(request.SessionId, workerCommand, context.CancellationToken)
|
||||
.ConfigureAwait(false);
|
||||
|
||||
@@ -29,9 +29,27 @@ public sealed class MxAccessGrpcMapper
|
||||
ArgumentNullException.ThrowIfNull(request);
|
||||
ArgumentNullException.ThrowIfNull(request.Command);
|
||||
|
||||
return MapCommand(request.Command);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Maps a gRPC MX command to a worker command. Callers that already hold the command — including
|
||||
/// the constraint pipeline, whose rewritten command is not the one on the request — use this
|
||||
/// overload rather than cloning a whole request to carry a single field (GWC-29); nothing outside
|
||||
/// the command is read here.
|
||||
/// </summary>
|
||||
/// <param name="command">Command payload.</param>
|
||||
/// <returns>The mapped <see cref="WorkerCommand"/> ready for worker dispatch.</returns>
|
||||
public WorkerCommand MapCommand(MxCommand command)
|
||||
{
|
||||
ArgumentNullException.ThrowIfNull(command);
|
||||
|
||||
// The clone is required and must stay: the caller may hand us the gRPC-owned request command,
|
||||
// and the caller keeps reading it after dispatch (TrackCommandReply). Cloning here is what makes
|
||||
// WorkerClient.CreateCommandEnvelope's no-aliasing invariant true.
|
||||
return new WorkerCommand
|
||||
{
|
||||
Command = request.Command.Clone(),
|
||||
Command = command.Clone(),
|
||||
EnqueueTimestamp = Timestamp.FromDateTimeOffset(_timeProvider.GetUtcNow()),
|
||||
};
|
||||
}
|
||||
|
||||
@@ -40,7 +40,9 @@ public sealed class WorkerClient : IWorkerClient
|
||||
private readonly ConcurrentDictionary<string, PendingCommand> _pendingCommands = new(StringComparer.Ordinal);
|
||||
private readonly SemaphoreSlim _pendingCommandSlots;
|
||||
private readonly CancellationTokenSource _stopCts = new();
|
||||
private long _nextSequence;
|
||||
// Touched only by WriteLoopAsync — the single consumer of _outboundEnvelopes — so it needs no
|
||||
// interlocking. See WriteLoopAsync for why the stamp happens there rather than at construction.
|
||||
private ulong _nextSequence;
|
||||
private WorkerClientState _state;
|
||||
private DateTimeOffset _lastHeartbeatAt;
|
||||
private int? _processId;
|
||||
@@ -404,6 +406,13 @@ public sealed class WorkerClient : IWorkerClient
|
||||
{
|
||||
await foreach (WorkerEnvelope envelope in _outboundEnvelopes.Reader.ReadAllAsync(_stopCts.Token).ConfigureAwait(false))
|
||||
{
|
||||
// GWC-28: stamp the sequence at the point of writing, not when the envelope is built.
|
||||
// Stamping at construction let two concurrent InvokeAsync callers take 1 and 2 and then
|
||||
// enqueue in the order 2, 1 — non-monotonic on the wire, breaking gateway.md's
|
||||
// "monotonic per sender" contract. This loop is the channel's single consumer
|
||||
// (SingleReader = true), so wire order and stamp order are the same thing here and
|
||||
// _nextSequence needs no interlocking. Mirrors the worker's WRK-04 fix.
|
||||
envelope.Sequence = unchecked(++_nextSequence);
|
||||
await _writer.WriteAsync(envelope, _stopCts.Token).ConfigureAwait(false);
|
||||
}
|
||||
}
|
||||
@@ -1072,11 +1081,13 @@ public sealed class WorkerClient : IWorkerClient
|
||||
string correlationId,
|
||||
Action<WorkerEnvelope> setBody)
|
||||
{
|
||||
// Sequence is deliberately left unset here: WriteLoopAsync stamps it immediately before the
|
||||
// frame goes out, so the numbers are monotonic in wire order however the callers interleave
|
||||
// between construction and enqueue (GWC-28, mirroring the worker's WRK-04 fix).
|
||||
WorkerEnvelope envelope = new()
|
||||
{
|
||||
ProtocolVersion = _connection.FrameOptions.ProtocolVersion,
|
||||
SessionId = SessionId,
|
||||
Sequence = (ulong)Interlocked.Increment(ref _nextSequence),
|
||||
CorrelationId = correlationId,
|
||||
};
|
||||
setBody(envelope);
|
||||
|
||||
@@ -5,11 +5,24 @@ using ZB.MOM.WW.MxGateway.Contracts.Proto;
|
||||
|
||||
namespace ZB.MOM.WW.MxGateway.Server.Workers;
|
||||
|
||||
/// <summary>
|
||||
/// Reads length-prefixed WorkerEnvelope protobuf frames from a stream.
|
||||
/// </summary>
|
||||
/// <remarks>
|
||||
/// <see cref="ReadAsync"/> is not reentrant: the reader keeps a per-instance length-prefix scratch
|
||||
/// buffer, so exactly one consumer may be inside a read at a time. That matches how the reader is
|
||||
/// used — a single read loop per <c>WorkerClient</c>, with handshake reads completing before the
|
||||
/// loop starts.
|
||||
/// </remarks>
|
||||
public sealed class WorkerFrameReader
|
||||
{
|
||||
private readonly WorkerFrameProtocolOptions _options;
|
||||
private readonly Stream _stream;
|
||||
|
||||
// Reused across frames rather than allocated per read (GWC-30). Safe because ReadAsync is
|
||||
// single-consumer by construction; the prefix is fully overwritten by every read.
|
||||
private readonly byte[] _lengthPrefix = new byte[sizeof(uint)];
|
||||
|
||||
/// <summary>
|
||||
/// Initializes a new instance of <see cref="WorkerFrameReader"/>.
|
||||
/// </summary>
|
||||
@@ -30,10 +43,9 @@ public sealed class WorkerFrameReader
|
||||
/// <returns>Parsed worker envelope.</returns>
|
||||
public async ValueTask<WorkerEnvelope> ReadAsync(CancellationToken cancellationToken = default)
|
||||
{
|
||||
byte[] lengthPrefix = new byte[sizeof(uint)];
|
||||
await ReadExactlyOrThrowAsync(lengthPrefix, cancellationToken).ConfigureAwait(false);
|
||||
await ReadExactlyOrThrowAsync(_lengthPrefix, cancellationToken).ConfigureAwait(false);
|
||||
|
||||
uint payloadLength = BinaryPrimitives.ReadUInt32LittleEndian(lengthPrefix);
|
||||
uint payloadLength = BinaryPrimitives.ReadUInt32LittleEndian(_lengthPrefix);
|
||||
if (payloadLength == 0)
|
||||
{
|
||||
throw new WorkerFrameProtocolException(
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
using Microsoft.Extensions.Time.Testing;
|
||||
using ZB.MOM.WW.MxGateway.Contracts.Proto;
|
||||
using ZB.MOM.WW.MxGateway.Server.Grpc;
|
||||
|
||||
@@ -38,6 +39,48 @@ public sealed class MxAccessGrpcMapperTests
|
||||
Assert.NotNull(workerCommand.EnqueueTimestamp);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// The command-only overload exists so <c>Invoke</c> does not deep-clone the whole request just to
|
||||
/// overwrite and discard its command (GWC-29). It must still perform the one clone that keeps the
|
||||
/// worker-bound graph unaliased from the caller-owned gRPC command, and must produce the same
|
||||
/// <see cref="WorkerCommand"/> as the request overload.
|
||||
/// </summary>
|
||||
[Fact]
|
||||
public void MapCommandFromCommandClonesPayload()
|
||||
{
|
||||
FakeTimeProvider timeProvider = new(new DateTimeOffset(2026, 8, 7, 12, 0, 0, TimeSpan.Zero));
|
||||
MxAccessGrpcMapper mapper = new(timeProvider);
|
||||
MxCommand command = new()
|
||||
{
|
||||
Kind = MxCommandKind.Write,
|
||||
Write = new WriteCommand
|
||||
{
|
||||
ServerHandle = 10,
|
||||
ItemHandle = 20,
|
||||
UserId = 30,
|
||||
Value = new MxValue
|
||||
{
|
||||
DataType = MxDataType.String,
|
||||
StringValue = "value",
|
||||
},
|
||||
},
|
||||
};
|
||||
MxCommandRequest request = new()
|
||||
{
|
||||
SessionId = "session-1",
|
||||
Command = command.Clone(),
|
||||
};
|
||||
|
||||
WorkerCommand fromCommand = mapper.MapCommand(command);
|
||||
WorkerCommand fromRequest = mapper.MapCommand(request);
|
||||
command.Write.Value.StringValue = "changed";
|
||||
|
||||
Assert.Equal(MxCommandKind.Write, fromCommand.Command.Kind);
|
||||
Assert.Equal("value", fromCommand.Command.Write.Value.StringValue);
|
||||
Assert.NotNull(fromCommand.EnqueueTimestamp);
|
||||
Assert.Equal(fromRequest, fromCommand);
|
||||
}
|
||||
|
||||
/// <summary>Verifies that command reply mapping preserves HRESULT and status information.</summary>
|
||||
[Fact]
|
||||
public void MapCommandReply_PreservesHresultStatusesAndPayload()
|
||||
|
||||
@@ -4,6 +4,7 @@ using ZB.MOM.WW.MxGateway.Contracts;
|
||||
using ZB.MOM.WW.MxGateway.Contracts.Proto;
|
||||
using ZB.MOM.WW.MxGateway.Server.Metrics;
|
||||
using ZB.MOM.WW.MxGateway.Server.Workers;
|
||||
using ZB.MOM.WW.MxGateway.Tests.Gateway.Workers.Fakes;
|
||||
using ZB.MOM.WW.MxGateway.Tests.TestSupport;
|
||||
|
||||
namespace ZB.MOM.WW.MxGateway.Tests.Gateway.Workers;
|
||||
@@ -29,6 +30,33 @@ public sealed class WorkerClientTests
|
||||
Assert.Equal(WorkerProcessId, client.ProcessId);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// The <c>GatewayHello</c> carries the negotiated worker-frame maximum so the worker adopts the
|
||||
/// configured limit instead of its own default (IPC-02); a regression to sending <c>0</c> means
|
||||
/// "older gateway, use default" to the worker and would silently downgrade the negotiated limit.
|
||||
/// The adoption half is asserted only in the Windows-only worker suite, so the gateway half is
|
||||
/// pinned here, in the portable suite (TST-28). Both the default and an override are covered so
|
||||
/// the assertion tracks configuration rather than a constant.
|
||||
/// </summary>
|
||||
/// <param name="maxMessageBytes">Configured worker-frame maximum to negotiate.</param>
|
||||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||||
[Theory]
|
||||
[InlineData(WorkerFrameProtocolOptions.DefaultMaxMessageBytes)]
|
||||
[InlineData(2 * 1024 * 1024)]
|
||||
public async Task StartAsync_SendsGatewayHelloWithConfiguredMaxFrameBytes(int maxMessageBytes)
|
||||
{
|
||||
await using FakeWorkerHarness harness =
|
||||
await FakeWorkerHarness.CreateConnectedPairAsync(maxMessageBytes: maxMessageBytes);
|
||||
await using WorkerClient client = harness.CreateClient();
|
||||
|
||||
Task startTask = client.StartAsync(CancellationToken.None);
|
||||
WorkerEnvelope gatewayHello = await harness.CompleteStartupAsync().WaitAsync(TestTimeout);
|
||||
await startTask.WaitAsync(TestTimeout);
|
||||
|
||||
Assert.NotEqual(0u, gatewayHello.GatewayHello.MaxFrameBytes);
|
||||
Assert.Equal((uint)maxMessageBytes, gatewayHello.GatewayHello.MaxFrameBytes);
|
||||
}
|
||||
|
||||
/// <summary>Verifies that InvokeAsync completes a pending command when a matching reply arrives.</summary>
|
||||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||||
[Fact]
|
||||
@@ -141,6 +169,48 @@ public sealed class WorkerClientTests
|
||||
Assert.Equal(MxCommandKind.GetWorkerInfo, reply.Reply.Kind);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// The envelope <c>sequence</c> is a monotonic per-sender counter (gateway.md), so the values
|
||||
/// observed on the pipe must be strictly increasing in wire order. Stamping the sequence when
|
||||
/// the envelope is constructed lets two concurrent invokes stamp 1 and 2 but enqueue 2 then 1
|
||||
/// (GWC-28); stamping on the single-consumer write loop makes wire order and sequence order the
|
||||
/// same thing by construction.
|
||||
/// </summary>
|
||||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||||
[Fact]
|
||||
public async Task ConcurrentInvokesEmitStrictlyIncreasingSequencesOnTheWire()
|
||||
{
|
||||
const int commandCount = 32;
|
||||
await using PipePair pipePair = await PipePair.CreateAsync();
|
||||
await using WorkerClient client = CreateClient(pipePair);
|
||||
await CompleteHandshakeAsync(client, pipePair);
|
||||
|
||||
Task<WorkerCommandReply>[] invokeTasks = Enumerable.Range(0, commandCount)
|
||||
.Select(_ => Task.Run(async () => await client.InvokeAsync(
|
||||
CreateCommand(MxCommandKind.Ping),
|
||||
TestTimeout,
|
||||
CancellationToken.None)))
|
||||
.ToArray();
|
||||
|
||||
ulong previousSequence = 0;
|
||||
for (int index = 0; index < commandCount; index++)
|
||||
{
|
||||
WorkerEnvelope commandEnvelope = await pipePair.WorkerReader.ReadAsync().AsTask().WaitAsync(TestTimeout);
|
||||
Assert.Equal(WorkerEnvelope.BodyOneofCase.WorkerCommand, commandEnvelope.BodyCase);
|
||||
Assert.True(
|
||||
commandEnvelope.Sequence > previousSequence,
|
||||
$"Command {index} arrived with sequence {commandEnvelope.Sequence} after {previousSequence}; "
|
||||
+ "envelope sequences must be strictly increasing in wire order.");
|
||||
previousSequence = commandEnvelope.Sequence;
|
||||
|
||||
await pipePair.WorkerWriter.WriteAsync(
|
||||
CreateCommandReplyEnvelope(commandEnvelope.CorrelationId, MxCommandKind.Ping));
|
||||
}
|
||||
|
||||
await Task.WhenAll(invokeTasks).WaitAsync(TestTimeout);
|
||||
Assert.Equal(WorkerClientState.Ready, client.State);
|
||||
}
|
||||
|
||||
/// <summary>Verifies that ReadEventsAsync yields events in pipe order from the worker.</summary>
|
||||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||||
[Fact]
|
||||
|
||||
@@ -55,6 +55,41 @@ public sealed class WorkerFrameProtocolTests
|
||||
Assert.Equal(original, parsed);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// One reader instance reads many frames in sequence. The reader reuses a single length-prefix
|
||||
/// scratch buffer across calls (GWC-30), so varying the payload length frame to frame proves the
|
||||
/// reused buffer is fully overwritten each time rather than carrying a stale prefix forward.
|
||||
/// </summary>
|
||||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||||
[Fact]
|
||||
public async Task ReadAsync_WithMultipleFramesOnOneReader_ParsesEveryFrame()
|
||||
{
|
||||
const int frameCount = 5;
|
||||
WorkerFrameProtocolOptions options = new(SessionId);
|
||||
await using MemoryStream stream = new();
|
||||
WorkerFrameWriter writer = new(stream, options);
|
||||
|
||||
List<WorkerEnvelope> originals = [];
|
||||
for (int index = 1; index <= frameCount; index++)
|
||||
{
|
||||
WorkerEnvelope envelope = CreateEnvelope();
|
||||
envelope.Sequence = (ulong)index;
|
||||
|
||||
// Differing payload lengths so a stale length prefix would misparse rather than pass.
|
||||
envelope.WorkerHello.WorkerVersion = new string('v', index * 37);
|
||||
originals.Add(envelope);
|
||||
await writer.WriteAsync(envelope);
|
||||
}
|
||||
|
||||
stream.Position = 0;
|
||||
WorkerFrameReader reader = new(stream, options);
|
||||
foreach (WorkerEnvelope original in originals)
|
||||
{
|
||||
WorkerEnvelope parsed = await reader.ReadAsync();
|
||||
Assert.Equal(original, parsed);
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>Verifies that reading a frame with partial reads reassembles the frame correctly.</summary>
|
||||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||||
[Fact]
|
||||
|
||||
Reference in New Issue
Block a user