using Akka.Actor; using Akka.TestKit.Xunit2; using ZB.MOM.WW.ScadaBridge.Commons.Messages.Artifacts; using ZB.MOM.WW.ScadaBridge.Commons.Messages.Integration; using ZB.MOM.WW.ScadaBridge.Commons.Messages.RemoteQuery; using ZB.MOM.WW.ScadaBridge.Commons.Types; using ZB.MOM.WW.ScadaBridge.Communication.Actors; using ZB.MOM.WW.ScadaBridge.Communication.Grpc; namespace ZB.MOM.WW.ScadaBridge.Communication.Tests; /// /// The single routing truth for the 28 migrated central→site commands. Proves each command /// resolves to EXACTLY the target the actor used inline before the extraction — including the /// node-local parked handler (never the singleton proxy) and the local failover path — so the /// Akka actor and the gRPC service can share one table without drift. /// public class SiteCommandDispatcherTests : TestKit { private const string SiteId = "site1"; private SiteCommandDispatcher Build(IActorRef dmProxy, Func? failover = null) => new(SiteId, dmProxy, failover ?? ((_, _) => null)); // ── The 19 commands that forward to the Deployment Manager singleton proxy ── /// The proxy-routed commands (lifecycle, OPC UA, debug snapshot/subscribe, route). public static IEnumerable ProxyCommandTypes() => new[] { typeof(Commons.Messages.Deployment.RefreshDeploymentCommand), typeof(Commons.Messages.Lifecycle.EnableInstanceCommand), typeof(Commons.Messages.Lifecycle.DisableInstanceCommand), typeof(Commons.Messages.Lifecycle.DeleteInstanceCommand), typeof(Commons.Messages.Deployment.DeploymentStateQueryRequest), typeof(Commons.Messages.Management.BrowseNodeCommand), typeof(Commons.Messages.Management.SearchAddressSpaceCommand), typeof(Commons.Messages.Management.ReadTagValuesCommand), typeof(Commons.Messages.Management.VerifyEndpointCommand), typeof(Commons.Messages.Management.TrustServerCertCommand), typeof(Commons.Messages.Management.ListServerCertsCommand), typeof(Commons.Messages.Management.RemoveServerCertCommand), typeof(Commons.Messages.DataConnection.WriteTagRequest), typeof(Commons.Messages.DebugView.DebugSnapshotRequest), typeof(Commons.Messages.DebugView.SubscribeDebugViewRequest), typeof(Commons.Messages.InboundApi.RouteToCallRequest), typeof(Commons.Messages.InboundApi.RouteToGetAttributesRequest), typeof(Commons.Messages.InboundApi.RouteToSetAttributesRequest), typeof(Commons.Messages.InboundApi.RouteToWaitForAttributeRequest), }.Select(t => new object[] { t }); [Theory] [MemberData(nameof(ProxyCommandTypes))] public void ProxyCommands_ForwardToDeploymentManager(Type commandType) { var dm = CreateTestProbe(); var dispatcher = Build(dm.Ref); var command = SiteCommandSamples.All[commandType][0]; var route = dispatcher.ResolveRoute(command); Assert.Equal(SiteCommandDispatcher.RouteDisposition.Forward, route.Disposition); Assert.Same(dm.Ref, route.Target); } [Fact] public void UnsubscribeDebugView_IsFireAndForget_ToDeploymentManager_WithSyntheticAck() { // Fire-and-forget: the Deployment Manager never acks an unsubscribe. The actor Forwards it; // the gRPC transport Tells it and returns this synthetic ack so a unary RPC still answers. var dm = CreateTestProbe(); var dispatcher = Build(dm.Ref); var command = SiteCommandSamples.All[typeof(Commons.Messages.DebugView.UnsubscribeDebugViewRequest)][0]; var route = dispatcher.ResolveRoute(command); Assert.Equal(SiteCommandDispatcher.RouteDisposition.TellFireAndForget, route.Disposition); Assert.Same(dm.Ref, route.Target); Assert.Same(UnsubscribeDebugViewAck.Instance, route.Reply); } // ── Artifact handler (null-guarded) ── [Fact] public void DeployArtifacts_WithHandler_ForwardsToArtifactHandler_NotTheProxy() { var dm = CreateTestProbe(); var artifact = CreateTestProbe(); var dispatcher = Build(dm.Ref); dispatcher.RegisterArtifactHandler(artifact.Ref); var route = dispatcher.ResolveRoute(SiteCommandSamples.All[typeof(DeployArtifactsCommand)][0]); Assert.Equal(SiteCommandDispatcher.RouteDisposition.Forward, route.Disposition); Assert.Same(artifact.Ref, route.Target); Assert.NotSame(dm.Ref, route.Target); } [Fact] public void DeployArtifacts_WithoutHandler_RepliesHandlerNotAvailable() { var dm = CreateTestProbe(); var dispatcher = Build(dm.Ref); var command = (DeployArtifactsCommand)SiteCommandSamples.All[typeof(DeployArtifactsCommand)][0]; var route = dispatcher.ResolveRoute(command); Assert.Equal(SiteCommandDispatcher.RouteDisposition.ImmediateReply, route.Disposition); var reply = Assert.IsType(route.Reply); Assert.False(reply.Success); Assert.Equal("Artifact handler not available", reply.ErrorMessage); Assert.Equal(command.DeploymentId, reply.DeploymentId); Assert.Equal(SiteId, reply.SiteId); } // ── Event-log handler (null-guarded) ── [Fact] public void EventLogQuery_WithHandler_ForwardsToEventLogHandler() { var dm = CreateTestProbe(); var eventLog = CreateTestProbe(); var dispatcher = Build(dm.Ref); dispatcher.RegisterEventLogHandler(eventLog.Ref); var route = dispatcher.ResolveRoute(SiteCommandSamples.All[typeof(EventLogQueryRequest)][0]); Assert.Equal(SiteCommandDispatcher.RouteDisposition.Forward, route.Disposition); Assert.Same(eventLog.Ref, route.Target); } [Fact] public void EventLogQuery_WithoutHandler_RepliesHandlerNotAvailable() { var dm = CreateTestProbe(); var dispatcher = Build(dm.Ref); var command = (EventLogQueryRequest)SiteCommandSamples.All[typeof(EventLogQueryRequest)][0]; var route = dispatcher.ResolveRoute(command); Assert.Equal(SiteCommandDispatcher.RouteDisposition.ImmediateReply, route.Disposition); var reply = Assert.IsType(route.Reply); Assert.False(reply.Success); Assert.Equal(command.CorrelationId, reply.CorrelationId); } // ── Parked handler (null-guarded) — stays NODE-LOCAL on purpose (replicated store) ── /// The five parked commands, each of which routes to the per-node parked handler. public static IEnumerable ParkedCommandTypes() => new[] { typeof(ParkedMessageQueryRequest), typeof(ParkedMessageRetryRequest), typeof(ParkedMessageDiscardRequest), typeof(RetryParkedOperation), typeof(DiscardParkedOperation), }.Select(t => new object[] { t }); [Theory] [MemberData(nameof(ParkedCommandTypes))] public void ParkedCommands_WithHandler_RouteToNodeLocalParkedHandler_NeverTheProxy(Type commandType) { // Node-locality proof: a parked retry/discard must run on the node holding the replicated // store row, so it goes to the per-node parked handler — NEVER the singleton proxy. This is // the constraint the extraction must not "fix". var dm = CreateTestProbe(); var parked = CreateTestProbe(); var dispatcher = Build(dm.Ref); dispatcher.RegisterParkedMessageHandler(parked.Ref); var route = dispatcher.ResolveRoute(SiteCommandSamples.All[commandType][0]); Assert.Equal(SiteCommandDispatcher.RouteDisposition.Forward, route.Disposition); Assert.Same(parked.Ref, route.Target); Assert.NotSame(dm.Ref, route.Target); } [Fact] public void ParkedMessageQuery_WithoutHandler_RepliesHandlerNotAvailable() { var dispatcher = Build(CreateTestProbe().Ref); var command = (ParkedMessageQueryRequest)SiteCommandSamples.All[typeof(ParkedMessageQueryRequest)][0]; var route = dispatcher.ResolveRoute(command); var reply = Assert.IsType(route.Reply); Assert.False(reply.Success); Assert.Equal(command.CorrelationId, reply.CorrelationId); Assert.Equal(command.PageNumber, reply.PageNumber); } [Fact] public void ParkedMessageRetry_WithoutHandler_RepliesHandlerNotAvailable() { var dispatcher = Build(CreateTestProbe().Ref); var command = (ParkedMessageRetryRequest)SiteCommandSamples.All[typeof(ParkedMessageRetryRequest)][0]; var reply = Assert.IsType(dispatcher.ResolveRoute(command).Reply); Assert.False(reply.Success); Assert.Equal(command.CorrelationId, reply.CorrelationId); } [Fact] public void ParkedMessageDiscard_WithoutHandler_RepliesHandlerNotAvailable() { var dispatcher = Build(CreateTestProbe().Ref); var command = (ParkedMessageDiscardRequest)SiteCommandSamples.All[typeof(ParkedMessageDiscardRequest)][0]; var reply = Assert.IsType(dispatcher.ResolveRoute(command).Reply); Assert.False(reply.Success); Assert.Equal(command.CorrelationId, reply.CorrelationId); } [Fact] public void RetryParkedOperation_WithoutHandler_RepliesNotAppliedAck() { var dispatcher = Build(CreateTestProbe().Ref); var command = new RetryParkedOperation("corr-x", TrackedOperationId.New()); var reply = Assert.IsType(dispatcher.ResolveRoute(command).Reply); Assert.False(reply.Applied); Assert.Equal("corr-x", reply.CorrelationId); Assert.NotNull(reply.ErrorMessage); } [Fact] public void DiscardParkedOperation_WithoutHandler_RepliesNotAppliedAck() { var dispatcher = Build(CreateTestProbe().Ref); var command = new DiscardParkedOperation("corr-y", TrackedOperationId.New()); var reply = Assert.IsType(dispatcher.ResolveRoute(command).Reply); Assert.False(reply.Applied); Assert.Equal("corr-y", reply.CorrelationId); } // ── Commands that must NOT enter ResolveRoute ── [Fact] public void ResolveRoute_RejectsFailover_ItGoesThroughPrepareFailover() { var dispatcher = Build(CreateTestProbe().Ref); Assert.Throws(() => dispatcher.ResolveRoute(new TriggerSiteFailover("c", SiteId))); } [Fact] public void ResolveRoute_RejectsTheExcludedIntegrationCommand() { // IntegrationCallRequest is the 29th command, dead at both ends and deliberately excluded // (28 of 29 migrate). It never enters the dispatcher. var dispatcher = Build(CreateTestProbe().Ref); var command = new IntegrationCallRequest( "c", SiteId, "inst", "es", "m", new Dictionary(), DateTimeOffset.UtcNow); Assert.Throws(() => dispatcher.ResolveRoute(command)); } // ── Failover (local path) ── [Fact] public void HandleFailover_ResolvesAndLeaves_ThenAcksWithTheTarget() { string? roleAsked = null; var leaveIssued = false; Func resolve = (role, dryRun) => { roleAsked = role; leaveIssued = !dryRun; return "akka.tcp://scadabridge@site1-a:8082"; }; var dispatcher = Build(CreateTestProbe().Ref, resolve); var ack = dispatcher.HandleFailover(new TriggerSiteFailover("corr-1", SiteId)); Assert.True(ack.Accepted); Assert.Equal("corr-1", ack.CorrelationId); Assert.Equal("akka.tcp://scadabridge@site1-a:8082", ack.TargetAddress); Assert.Null(ack.ErrorMessage); // Site singletons are scoped to the site-specific role. Assert.Equal("site-site1", roleAsked); // The actor path leaves in one step (dryRun:false). Assert.True(leaveIssued); } [Fact] public void HandleFailover_RefusesWhenThereIsNoPeer() { var dispatcher = Build(CreateTestProbe().Ref, (_, _) => null); var ack = dispatcher.HandleFailover(new TriggerSiteFailover("corr-2", SiteId)); Assert.False(ack.Accepted); Assert.Null(ack.TargetAddress); Assert.NotNull(ack.ErrorMessage); } [Fact] public void HandleFailover_RefusesACommandAddressedToAnotherSite_WithoutTouchingTheResolver() { var invoked = false; var dispatcher = Build(CreateTestProbe().Ref, (_, _) => { invoked = true; return "addr"; }); var ack = dispatcher.HandleFailover(new TriggerSiteFailover("corr-3", "site2")); Assert.False(ack.Accepted); Assert.Contains("site2", ack.ErrorMessage); Assert.False(invoked); } [Fact] public void HandleFailover_FaultInTheLeave_IsReportedNotThrown() { var dispatcher = Build(CreateTestProbe().Ref, (_, _) => throw new InvalidOperationException("cluster unavailable")); var ack = dispatcher.HandleFailover(new TriggerSiteFailover("corr-4", SiteId)); Assert.False(ack.Accepted); Assert.Contains("cluster unavailable", ack.ErrorMessage); } [Fact] public void PrepareFailover_BuildsTheAckFromADryRun_WithoutLeaving_UntilCommitLeaveRuns() { // The gRPC ack-before-Leave proof at the routing-truth level: PrepareFailover resolves the // standby with a DRY-RUN only (so the ack is built without leaving), and hands back a // deferred CommitLeave. Only invoking CommitLeave — what the gRPC service does AFTER the // ack is on the wire — performs the real leave. var events = new List(); Func resolve = (_, dryRun) => { events.Add(dryRun ? "resolve" : "leave"); return "akka.tcp://scadabridge@site1-a:8082"; }; var dispatcher = Build(CreateTestProbe().Ref, resolve); var outcome = dispatcher.PrepareFailover(new TriggerSiteFailover("corr-5", SiteId)); Assert.True(outcome.Ack.Accepted); Assert.Equal("akka.tcp://scadabridge@site1-a:8082", outcome.Ack.TargetAddress); Assert.NotNull(outcome.CommitLeave); // Building the ack did NOT leave. Assert.Equal(new[] { "resolve" }, events); outcome.CommitLeave!(); // The leave runs strictly after — the caller sends the ack first. Assert.Equal(new[] { "resolve", "leave" }, events); } [Fact] public void PrepareFailover_WhenRefused_HasNoDeferredLeave() { var dispatcher = Build(CreateTestProbe().Ref, (_, _) => null); var outcome = dispatcher.PrepareFailover(new TriggerSiteFailover("corr-6", SiteId)); Assert.False(outcome.Ack.Accepted); Assert.Null(outcome.CommitLeave); } }