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);
}
}