diff --git a/src/ZB.MOM.WW.MxGateway.Contracts/Generated/MxaccessGateway.cs b/src/ZB.MOM.WW.MxGateway.Contracts/Generated/MxaccessGateway.cs index 400e221..bbaa9ff 100644 --- a/src/ZB.MOM.WW.MxGateway.Contracts/Generated/MxaccessGateway.cs +++ b/src/ZB.MOM.WW.MxGateway.Contracts/Generated/MxaccessGateway.cs @@ -8012,6 +8012,11 @@ namespace ZB.MOM.WW.MxGateway.Contracts.Proto { } + /// + /// The unary reply's statuses field carries the correlated OnWriteComplete + /// outcome when it arrives within the worker's bounded wait — see + /// MxCommandReply.statuses. + /// [global::System.Diagnostics.DebuggerDisplayAttribute("{ToString(),nq}")] public sealed partial class WriteCommand : pb::IMessage #if !GOOGLE_PROTOBUF_REFSTRUCT_COMPATIBILITY_MODE @@ -8330,6 +8335,9 @@ namespace ZB.MOM.WW.MxGateway.Contracts.Proto { } + /// + /// Same statuses correlation as WriteCommand. + /// [global::System.Diagnostics.DebuggerDisplayAttribute("{ToString(),nq}")] public sealed partial class Write2Command : pb::IMessage #if !GOOGLE_PROTOBUF_REFSTRUCT_COMPATIBILITY_MODE @@ -17215,18 +17223,19 @@ namespace ZB.MOM.WW.MxGateway.Contracts.Proto { = pb::FieldCodec.ForMessage(58, global::ZB.MOM.WW.MxGateway.Contracts.Proto.MxStatusProxy.Parser); private readonly pbc::RepeatedField statuses_ = new pbc::RepeatedField(); /// - /// Correlated per-item outcome rows. For WRITE_SECURED / WRITE_SECURED2 - /// replies the worker holds the reply for a bounded window (default 1.5 s, - /// MXGATEWAY_WORKER_WRITE_COMPLETION_WAIT_MS) waiting for the matching - /// MXAccess OnWriteComplete callback and copies its status rows here, so - /// statuses[0] carries the real MXAccess commit outcome (success OR failure) - /// while protocol_status/hresult still describe command acceptance only. - /// Empty statuses on a write reply means the completion did not arrive + /// Correlated per-item outcome rows. For WRITE / WRITE2 / WRITE_SECURED / + /// WRITE_SECURED2 replies the worker holds the reply for a bounded window + /// (default 1.5 s, MXGATEWAY_WORKER_WRITE_COMPLETION_WAIT_MS) waiting for + /// the matching MXAccess OnWriteComplete callback and copies its status rows + /// here, so statuses[0] carries the real MXAccess commit outcome (success OR + /// failure) while protocol_status/hresult still describe command acceptance + /// only. Empty statuses on a write reply means the completion did not arrive /// within the window — the write is unconfirmed, not failed. Correlation is /// best-effort per (server_handle, item_handle): MXAccess's callback carries /// no transaction id, so concurrent writes to the same item within the /// window can swap rows. The OnWriteComplete event still flows on the event - /// stream unchanged. Other command kinds leave this field as before. + /// stream unchanged. Bulk write kinds and all non-write kinds leave this + /// field as before. /// [global::System.Diagnostics.DebuggerNonUserCodeAttribute] [global::System.CodeDom.Compiler.GeneratedCode("protoc", null)] diff --git a/src/ZB.MOM.WW.MxGateway.Contracts/Protos/mxaccess_gateway.proto b/src/ZB.MOM.WW.MxGateway.Contracts/Protos/mxaccess_gateway.proto index ee923b5..db5fb4f 100644 --- a/src/ZB.MOM.WW.MxGateway.Contracts/Protos/mxaccess_gateway.proto +++ b/src/ZB.MOM.WW.MxGateway.Contracts/Protos/mxaccess_gateway.proto @@ -241,6 +241,9 @@ message ActivateCommand { int32 item_handle = 2; } +// The unary reply's statuses field carries the correlated OnWriteComplete +// outcome when it arrives within the worker's bounded wait — see +// MxCommandReply.statuses. message WriteCommand { int32 server_handle = 1; int32 item_handle = 2; @@ -248,6 +251,7 @@ message WriteCommand { int32 user_id = 4; } +// Same statuses correlation as WriteCommand. message Write2Command { int32 server_handle = 1; int32 item_handle = 2; @@ -531,18 +535,19 @@ message MxCommandReply { // transport failures. optional int32 hresult = 5; MxValue return_value = 6; - // Correlated per-item outcome rows. For WRITE_SECURED / WRITE_SECURED2 - // replies the worker holds the reply for a bounded window (default 1.5 s, - // MXGATEWAY_WORKER_WRITE_COMPLETION_WAIT_MS) waiting for the matching - // MXAccess OnWriteComplete callback and copies its status rows here, so - // statuses[0] carries the real MXAccess commit outcome (success OR failure) - // while protocol_status/hresult still describe command acceptance only. - // Empty statuses on a write reply means the completion did not arrive + // Correlated per-item outcome rows. For WRITE / WRITE2 / WRITE_SECURED / + // WRITE_SECURED2 replies the worker holds the reply for a bounded window + // (default 1.5 s, MXGATEWAY_WORKER_WRITE_COMPLETION_WAIT_MS) waiting for + // the matching MXAccess OnWriteComplete callback and copies its status rows + // here, so statuses[0] carries the real MXAccess commit outcome (success OR + // failure) while protocol_status/hresult still describe command acceptance + // only. Empty statuses on a write reply means the completion did not arrive // within the window — the write is unconfirmed, not failed. Correlation is // best-effort per (server_handle, item_handle): MXAccess's callback carries // no transaction id, so concurrent writes to the same item within the // window can swap rows. The OnWriteComplete event still flows on the event - // stream unchanged. Other command kinds leave this field as before. + // stream unchanged. Bulk write kinds and all non-write kinds leave this + // field as before. repeated MxStatusProxy statuses = 7; string diagnostic_message = 8; diff --git a/src/ZB.MOM.WW.MxGateway.Worker.Tests/MxAccess/MxAccessCommandExecutorTests.cs b/src/ZB.MOM.WW.MxGateway.Worker.Tests/MxAccess/MxAccessCommandExecutorTests.cs index 34f9406..a8ac61d 100644 --- a/src/ZB.MOM.WW.MxGateway.Worker.Tests/MxAccess/MxAccessCommandExecutorTests.cs +++ b/src/ZB.MOM.WW.MxGateway.Worker.Tests/MxAccess/MxAccessCommandExecutorTests.cs @@ -851,6 +851,9 @@ public sealed class MxAccessCommandExecutorTests FakeMxAccessComObjectFactory factory = new(fakeComObject); using StaRuntime runtime = CreateRuntime(); using MxAccessStaSession session = new(runtime, factory, new NoopEventSink()); + // No completion source in this test — disable the bounded reply wait + // so the forwarding assertions don't pay the default 1.5 s timeout. + session.WriteCompletionTimeout = TimeSpan.Zero; await session.StartAsync(workerProcessId: 1234); MxCommandReply reply = await session.DispatchAsync(CreateWriteCommand( @@ -874,6 +877,8 @@ public sealed class MxAccessCommandExecutorTests FakeMxAccessComObjectFactory factory = new(fakeComObject); using StaRuntime runtime = CreateRuntime(); using MxAccessStaSession session = new(runtime, factory, new NoopEventSink()); + // Same rationale as the Write forwarding test above. + session.WriteCompletionTimeout = TimeSpan.Zero; await session.StartAsync(workerProcessId: 1234); DateTime timestamp = new(2026, 5, 19, 12, 0, 0, DateTimeKind.Utc); @@ -952,7 +957,7 @@ public sealed class MxAccessCommandExecutorTests FakeMxAccessComObject fakeComObject = new(registerHandle: 82); FakeMxAccessComObjectFactory factory = new(fakeComObject); CompletionCacheEventSink sink = new(); - fakeComObject.OnWriteSecuredCallback = () => + fakeComObject.OnWriteCallback = () => sink.WriteCompletionCache.Record(82, 820, CreateCompletionRows(detail: 4321)); using StaRuntime runtime = CreateRuntime(); using MxAccessStaSession session = new(runtime, factory, sink); @@ -989,7 +994,7 @@ public sealed class MxAccessCommandExecutorTests // baseline is committed and a Record from the test thread is // guaranteed to be "newer" — no fixed sleep racing the STA thread. using System.Threading.ManualResetEventSlim comCallReached = new(initialState: false); - fakeComObject.OnWriteSecuredCallback = () => comCallReached.Set(); + fakeComObject.OnWriteCallback = () => comCallReached.Set(); using StaRuntime runtime = CreateRuntime(); using MxAccessStaSession session = new(runtime, factory, sink); session.WriteCompletionTimeout = TimeSpan.FromSeconds(10); @@ -1068,7 +1073,7 @@ public sealed class MxAccessCommandExecutorTests FakeMxAccessComObject fakeComObject = new(registerHandle: 86); FakeMxAccessComObjectFactory factory = new(fakeComObject); CompletionCacheEventSink sink = new(); - fakeComObject.OnWriteSecuredCallback = () => + fakeComObject.OnWriteCallback = () => sink.WriteCompletionCache.Record(86, 860, CreateCompletionRows(detail: 2222)); using StaRuntime runtime = CreateRuntime(); using MxAccessStaSession session = new(runtime, factory, sink); @@ -1085,25 +1090,110 @@ public sealed class MxAccessCommandExecutorTests } /// - /// Verifies plain Write stays fire-and-forget: even with a huge completion - /// timeout configured and no completion source, the reply returns - /// immediately (guarded well under the configured wait) with empty - /// statuses — only the secured write kinds enter the bounded wait. + /// Verifies plain Write correlates the same way as WriteSecured + /// (fast-completion edge): OtOpcUa's dominant FreeAccess path goes out + /// as MX_COMMAND_KIND_WRITE, so its reply must carry the OnWriteComplete + /// rows too. /// /// A task that represents the asynchronous operation. [Fact] - public async Task DispatchAsync_Write_DoesNotWaitForCompletion() + public async Task DispatchAsync_Write_WhenCompletionArrivesDuringComCall_ReturnsStatuses() { FakeMxAccessComObject fakeComObject = new(registerHandle: 87); FakeMxAccessComObjectFactory factory = new(fakeComObject); CompletionCacheEventSink sink = new(); + fakeComObject.OnWriteCallback = () => + sink.WriteCompletionCache.Record(87, 870, CreateCompletionRows(detail: 5555)); + using StaRuntime runtime = CreateRuntime(); + using MxAccessStaSession session = new(runtime, factory, sink); + // Hermetic: don't inherit MXGATEWAY_WORKER_WRITE_COMPLETION_WAIT_MS + // from the test runner's environment. + session.WriteCompletionTimeout = TimeSpan.FromSeconds(10); + await session.StartAsync(workerProcessId: 1234); + + MxCommandReply reply = await session.DispatchAsync(CreateWriteCommand( + "plain-write-fast", serverHandle: 87, itemHandle: 870, value: 1, userId: 5)); + + Assert.Equal(ProtocolStatusCode.Ok, reply.ProtocolStatus.Code); + Assert.True(reply.HasHresult); + Assert.Equal(0, reply.Hresult); + MxStatusProxy row = Assert.Single(reply.Statuses); + Assert.Equal(5555, row.Detail); + Assert.Equal(MxStatusCategory.Ok, row.Category); + } + + /// + /// Verifies the plain-Write timeout fallback mirrors the secured one: + /// no completion within the bounded wait returns protocol OK with EMPTY + /// statuses — unconfirmed, never a synthesized failure row. + /// + /// A task that represents the asynchronous operation. + [Fact] + public async Task DispatchAsync_Write_WhenNoCompletion_TimesOutWithEmptyStatusesAndOkProtocol() + { + FakeMxAccessComObject fakeComObject = new(registerHandle: 88); + FakeMxAccessComObjectFactory factory = new(fakeComObject); + CompletionCacheEventSink sink = new(); + using StaRuntime runtime = CreateRuntime(); + using MxAccessStaSession session = new(runtime, factory, sink); + session.WriteCompletionTimeout = TimeSpan.FromMilliseconds(100); + await session.StartAsync(workerProcessId: 1234); + + MxCommandReply reply = await session.DispatchAsync(CreateWriteCommand( + "plain-write-timeout", serverHandle: 88, itemHandle: 880, value: 1, userId: 5)); + + Assert.Equal(ProtocolStatusCode.Ok, reply.ProtocolStatus.Code); + Assert.True(reply.HasHresult); + Assert.Equal(0, reply.Hresult); + Assert.Empty(reply.Statuses); + } + + /// + /// Verifies Write2 correlates the same way as Write (fast-completion + /// edge). + /// + /// A task that represents the asynchronous operation. + [Fact] + public async Task DispatchAsync_Write2_WhenCompletionArrivesDuringComCall_ReturnsStatuses() + { + FakeMxAccessComObject fakeComObject = new(registerHandle: 89); + FakeMxAccessComObjectFactory factory = new(fakeComObject); + CompletionCacheEventSink sink = new(); + fakeComObject.OnWriteCallback = () => + sink.WriteCompletionCache.Record(89, 890, CreateCompletionRows(detail: 6666)); + using StaRuntime runtime = CreateRuntime(); + using MxAccessStaSession session = new(runtime, factory, sink); + session.WriteCompletionTimeout = TimeSpan.FromSeconds(10); + await session.StartAsync(workerProcessId: 1234); + + MxCommandReply reply = await session.DispatchAsync(CreateWrite2Command( + "plain-write2-fast", serverHandle: 89, itemHandle: 890, value: 1, + timestamp: new DateTime(2026, 8, 9, 12, 0, 0, DateTimeKind.Utc), userId: 6)); + + Assert.Equal(ProtocolStatusCode.Ok, reply.ProtocolStatus.Code); + Assert.Equal(6666, Assert.Single(reply.Statuses).Detail); + } + + /// + /// Verifies WriteBulk stays fire-and-forget: even with a huge completion + /// timeout configured and no completion source, the reply returns + /// immediately with per-entry results — bulk writes never enter the + /// bounded wait (latency for high-rate loops). + /// + /// A task that represents the asynchronous operation. + [Fact] + public async Task DispatchAsync_WriteBulk_DoesNotWaitForCompletion() + { + FakeMxAccessComObject fakeComObject = new(registerHandle: 90); + FakeMxAccessComObjectFactory factory = new(fakeComObject); + CompletionCacheEventSink sink = new(); using StaRuntime runtime = CreateRuntime(); using MxAccessStaSession session = new(runtime, factory, sink); session.WriteCompletionTimeout = TimeSpan.FromSeconds(30); await session.StartAsync(workerProcessId: 1234); - Task pending = session.DispatchAsync(CreateWriteCommand( - "plain-write-no-wait", serverHandle: 87, itemHandle: 870, value: 1, userId: 5)); + Task pending = session.DispatchAsync(CreateWriteBulkCommand( + "bulk-write-no-wait", serverHandle: 90, entries: new[] { (itemHandle: 900, value: 1, userId: 5) })); Task completed = await Task.WhenAny(pending, Task.Delay(TimeSpan.FromSeconds(5))); Assert.Same(pending, completed); @@ -2004,12 +2094,13 @@ public sealed class MxAccessCommandExecutorTests private readonly List operationNames = new(); /// - /// Invoked at the end of a successful WriteSecured/WriteSecured2 — - /// stands in for MXAccess committing synchronously and delivering - /// OnWriteComplete while the COM call is still on the stack, so - /// tests can exercise the fast-completion ordering edge. + /// Invoked at the end of a successful Write/Write2/WriteSecured/ + /// WriteSecured2 — stands in for MXAccess committing synchronously + /// and delivering OnWriteComplete while the COM call is still on + /// the stack, so tests can exercise the fast-completion ordering + /// edge on every correlated write kind. /// - public Action? OnWriteSecuredCallback { get; set; } + public Action? OnWriteCallback { get; set; } /// Initializes a fake MXAccess COM object with the given handles and optional exceptions. /// Return value for Register method. @@ -2285,6 +2376,7 @@ public sealed class MxAccessCommandExecutorTests WriteUserId = userId; WriteThreadId = Environment.CurrentManagedThreadId; ThrowIfWriteFailureConfigured(itemHandle); + OnWriteCallback?.Invoke(); } /// @@ -2303,6 +2395,7 @@ public sealed class MxAccessCommandExecutorTests WriteUserId = userId; WriteThreadId = Environment.CurrentManagedThreadId; ThrowIfWriteFailureConfigured(itemHandle); + OnWriteCallback?.Invoke(); } /// @@ -2321,7 +2414,7 @@ public sealed class MxAccessCommandExecutorTests WriteValue = value; WriteThreadId = Environment.CurrentManagedThreadId; ThrowIfWriteFailureConfigured(itemHandle); - OnWriteSecuredCallback?.Invoke(); + OnWriteCallback?.Invoke(); } /// @@ -2342,7 +2435,7 @@ public sealed class MxAccessCommandExecutorTests WriteTimestamp = timestamp; WriteThreadId = Environment.CurrentManagedThreadId; ThrowIfWriteFailureConfigured(itemHandle); - OnWriteSecuredCallback?.Invoke(); + OnWriteCallback?.Invoke(); } private void ThrowIfWriteFailureConfigured(int itemHandle) diff --git a/src/ZB.MOM.WW.MxGateway.Worker/MxAccess/MxAccessCommandExecutor.cs b/src/ZB.MOM.WW.MxGateway.Worker/MxAccess/MxAccessCommandExecutor.cs index 6b061bd..9c36df2 100644 --- a/src/ZB.MOM.WW.MxGateway.Worker/MxAccess/MxAccessCommandExecutor.cs +++ b/src/ZB.MOM.WW.MxGateway.Worker/MxAccess/MxAccessCommandExecutor.cs @@ -427,13 +427,28 @@ public sealed class MxAccessCommandExecutor : IStaCommandExecutor return CreateInvalidRequestReply(command, "Write command value is required."); } + // Same pre-call baseline rule as ExecuteWriteSecured: plain Write is + // also fire-and-forget in MXAccess, and its commit outcome only exists + // in the later OnWriteComplete callback. + MxAccessWriteCompletionCache completionCache = session.WriteCompletionCache; + ulong completionBaseline = completionCache.CurrentVersion( + writeCommand.ServerHandle, + writeCommand.ItemHandle); + session.Write( writeCommand.ServerHandle, writeCommand.ItemHandle, variantConverter.ConvertToComValue(writeCommand.Value), writeCommand.UserId); - return CreateOkReply(command); + MxCommandReply reply = CreateOkReply(command); + AwaitWriteCompletion( + reply, + completionCache, + writeCommand.ServerHandle, + writeCommand.ItemHandle, + completionBaseline); + return reply; } private MxCommandReply ExecuteWrite2(StaCommand command) @@ -454,6 +469,12 @@ public sealed class MxAccessCommandExecutor : IStaCommandExecutor return CreateInvalidRequestReply(command, "Write2 command timestamp value is required."); } + // Same pre-call baseline rule as ExecuteWriteSecured. + MxAccessWriteCompletionCache completionCache = session.WriteCompletionCache; + ulong completionBaseline = completionCache.CurrentVersion( + write2Command.ServerHandle, + write2Command.ItemHandle); + session.Write2( write2Command.ServerHandle, write2Command.ItemHandle, @@ -461,7 +482,14 @@ public sealed class MxAccessCommandExecutor : IStaCommandExecutor variantConverter.ConvertToComValue(write2Command.TimestampValue), write2Command.UserId); - return CreateOkReply(command); + MxCommandReply reply = CreateOkReply(command); + AwaitWriteCompletion( + reply, + completionCache, + write2Command.ServerHandle, + write2Command.ItemHandle, + completionBaseline); + return reply; } private MxCommandReply ExecuteWriteSecured(StaCommand command)