Files
ScadaBridge/tests/ZB.MOM.WW.ScadaBridge.Communication.Tests/SiteCommandDispatcherTests.cs
T
Joseph Doherty a5256e9b12 chore(comm): delete dead IntegrationCallRequest routing (Gitea #32)
'Pattern 4: Integration Routing' — RouteIntegrationCallAsync →
SiteEnvelope(IntegrationCallRequest) → SiteCommunicationActor integration
handler — was plumbed end to end but connected at neither end: no
producer (zero callers) and no handler (AkkaHostedService never
registered LocalHandlerType.Integration). It was an early scaffold the
architecture routed around — the brokered External→Central→Site→Central
round-trip is served by the Inbound API's routed-site-script path (the
RouteTo* verbs, driven from CommunicationServiceInstanceRouter), which is
live, tested, and shares IntegrationTimeout.

Decision (#32): delete. Removed the IntegrationCall{Request,Response}
messages, RouteIntegrationCallAsync, the SiteCommunicationActor receive
block + _integrationHandler field + LocalHandlerType.Integration, and the
four tests that covered them (2 actor, 2 message-contract, 1 dispatcher-
reject, 1 mapper-reject). KEPT IntegrationTimeout — it is the live
timeout for the RouteTo* verbs. Updated the exclusion-narrative comments
(proto/mapper/dispatcher), design §4, the components doc timeout table,
and marked the known-issue RESOLVED.

Full solution build clean (0/0); Communication 634 + Host 421 green.
Net -172/+20 across 14 files. Not in the gRPC proto (was deliberately
excluded there), so no wire-format change.
2026-07-23 14:27:32 -04:00

348 lines
14 KiB
C#

using Akka.Actor;
using Akka.TestKit.Xunit2;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Artifacts;
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;
/// <summary>
/// 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.
/// </summary>
public class SiteCommandDispatcherTests : TestKit
{
private const string SiteId = "site1";
private SiteCommandDispatcher Build(IActorRef dmProxy, Func<string, bool, string?>? failover = null)
=> new(SiteId, dmProxy, failover ?? ((_, _) => null));
// ── The 19 commands that forward to the Deployment Manager singleton proxy ──
/// <summary>The proxy-routed commands (lifecycle, OPC UA, debug snapshot/subscribe, route).</summary>
public static IEnumerable<object[]> 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<ArtifactDeploymentResponse>(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<EventLogQueryResponse>(route.Reply);
Assert.False(reply.Success);
Assert.Equal(command.CorrelationId, reply.CorrelationId);
}
// ── Parked handler (null-guarded) — stays NODE-LOCAL on purpose (replicated store) ──
/// <summary>The five parked commands, each of which routes to the per-node parked handler.</summary>
public static IEnumerable<object[]> 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<ParkedMessageQueryResponse>(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<ParkedMessageRetryResponse>(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<ParkedMessageDiscardResponse>(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<ParkedOperationActionAck>(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<ParkedOperationActionAck>(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<ArgumentException>(() =>
dispatcher.ResolveRoute(new TriggerSiteFailover("c", SiteId)));
}
// ── Failover (local path) ──
[Fact]
public void HandleFailover_ResolvesAndLeaves_ThenAcksWithTheTarget()
{
string? roleAsked = null;
var leaveIssued = false;
Func<string, bool, string?> 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<string>();
Func<string, bool, string?> 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);
}
}