using System.Collections.Immutable;
using Akka.Actor;
using Akka.Cluster;
using Akka.Cluster.Tools.Client;
using Akka.Configuration;
using Microsoft.EntityFrameworkCore;
using Shouldly;
using Xunit;
using ZB.MOM.WW.OtOpcUa.Cluster;
using ZB.MOM.WW.OtOpcUa.Commons.Messages.Admin;
using ZB.MOM.WW.OtOpcUa.Commons.Messages.Deploy;
using ZB.MOM.WW.OtOpcUa.Commons.Messages.Mesh;
using ZB.MOM.WW.OtOpcUa.Commons.Types;
using ZB.MOM.WW.OtOpcUa.Configuration;
using ZB.MOM.WW.OtOpcUa.Configuration.Entities;
using ZB.MOM.WW.OtOpcUa.ControlPlane.Communication;
using ZB.MOM.WW.OtOpcUa.Runtime.Communication;
using AkkaCluster = Akka.Cluster.Cluster;
namespace ZB.MOM.WW.OtOpcUa.Host.IntegrationTests;
///
/// Drives a real ClusterClient across two genuinely separate s —
/// the only test in the repo where the mesh boundary carries a message over the wire.
///
///
///
/// The sister project substitutes an in-process relay that unwraps
/// and forwards. That tests actor logic and cannot
/// fail for any transport reason: not a missing receptionist extension, not a service that
/// was never registered, not a malformed contact path, not a serialization break. Those are
/// precisely the Phase 2 failure modes, so they are what this test exercises.
///
///
/// Two clusters, one ActorSystem name, joined ONLY by ClusterClient. The shared name
/// is mandatory — Akka.Remote matches on the full address, so a contact path naming
/// otopcua cannot reach a system called anything else. Each system seeds itself and
/// neither lists the other, and
/// asserts the
/// separation outright: if they ever gossiped into one mesh, every delivery below could be
/// plain local messaging and the whole class would be testing nothing.
///
///
/// Scope limit. With exactly one node per side, "registered per node" and "registered
/// as a cluster singleton" are indistinguishable by construction, so this class does NOT
/// cover the per-node property whose structural assertions MeshCommActorPathTests
/// dropped. That property's discriminating coverage is live-gate step 8: stop the OLDEST
/// central and the survivor's comm actor must still receive a node ack, which a singleton
/// hosted on the stopped node cannot.
///
///
public sealed class MeshClusterClientBoundaryTests : IAsyncLifetime
{
private const string SharedSystemName = "otopcua";
private static readonly TimeSpan Arrival = TimeSpan.FromSeconds(30);
///
/// How long the falsifiability control waits before concluding nothing arrived. Deliberately
/// shorter than but far longer than a successful delivery takes in
/// practice, so a pass means "the boundary is shut", not "the clock ran out first".
///
private static readonly TimeSpan Silence = TimeSpan.FromSeconds(10);
///
/// One factory for every client this class creates, because the generation counter that keeps
/// client actor names unique is per factory instance . A second factory on the same
/// system restarts at generation 1 and collides with the first factory's client — which is
/// exactly how this test first failed. Production has one instance per
/// CentralCommunicationActor , so the real code path is safe.
///
private readonly DefaultMeshClusterClientFactory _clientFactory = new();
private ActorSystem _central = null!;
private ActorSystem _node = null!;
private IActorRef _nodeComm = null!;
private IActorRef _centralClient = null!;
private TaskCompletionSource _centralInbox = null!;
///
public async ValueTask InitializeAsync()
{
_central = StartSelfSeededSystem(withReceptionist: true);
_node = StartSelfSeededSystem(withReceptionist: true);
await Task.WhenAll(AwaitSingleMemberUpAsync(_central), AwaitSingleMemberUpAsync(_node));
// Central: the ack landing pad, at the path the node's comm actor sends to.
_centralInbox = NewInbox();
var centralComm = _central.ActorOf(
CaptureActor.Props(_centralInbox), MeshPaths.CentralCommunicationName);
ClusterClientReceptionist.Get(_central).RegisterService(centralComm);
// Node: the REAL comm actor, dialling central through the REAL production client factory.
var nodeToCentral = _clientFactory.Create(_node, "central", ReceptionistOf(_central));
_nodeComm = _node.ActorOf(
NodeCommunicationActor.Props(nodeToCentral), MeshPaths.NodeCommunicationName);
ClusterClientReceptionist.Get(_node).RegisterService(_nodeComm);
_centralClient = _clientFactory.Create(_central, "site-a", ReceptionistOf(_node));
}
///
public async ValueTask DisposeAsync()
{
await ShutdownAsync(_central);
await ShutdownAsync(_node);
}
[Fact]
public void The_two_systems_are_separate_clusters_joined_only_by_the_client()
{
// The load-bearing precondition of every other test here. Two systems sharing a name in one
// process are one gossip accident away from being one cluster, at which point ClusterClient
// is no longer the thing under test.
AkkaCluster.Get(_central).State.Members.Count.ShouldBe(1,
"central must be its own single-member cluster, or the boundary is not being crossed");
AkkaCluster.Get(_node).State.Members.Count.ShouldBe(1,
"the node must be its own single-member cluster, or the boundary is not being crossed");
}
[Fact]
public async Task Central_send_to_all_reaches_the_node_comm_actor_and_lands_on_its_event_stream()
{
var inbox = NewInbox();
_node.EventStream.Subscribe(_node.ActorOf(CaptureActor.Props(inbox)), typeof(RestartDriver));
var command = new RestartDriver("C1", "driver-1", "alice", Guid.NewGuid());
var delivered = await SendUntilDeliveredAsync(
_centralClient,
() => new ClusterClient.SendToAll(MeshPaths.NodeCommunication, command),
inbox,
Arrival);
delivered.ShouldNotBeNull(
"SendToAll must reach the node's registered comm actor and be re-emitted node-locally; "
+ "this is the whole central→node leg of the mesh transport");
delivered.ShouldBeOfType().DriverInstanceId.ShouldBe("driver-1");
}
[Fact]
public async Task A_node_ack_reaches_central_through_the_node_comm_actors_own_client()
{
// Told directly rather than routed from central, so a failure here can only be the outbound
// leg. The inbound leg has its own test above.
var ack = new ApplyAck(
new DeploymentId(Guid.NewGuid()),
new NodeId("node-a"),
ApplyAckOutcome.Applied,
null,
CorrelationId.NewId());
var delivered = await SendUntilDeliveredAsync(
_nodeComm, () => ack, _centralInbox, Arrival);
delivered.ShouldNotBeNull(
"the node's ApplyAck must reach central's comm actor, or every deployment times out at "
+ "the apply deadline naming nodes that applied successfully");
delivered.ShouldBeOfType().NodeId.Value.ShouldBe("node-a");
}
[Fact]
public async Task Falsifiability_control_a_node_without_the_receptionist_extension_receives_nothing()
{
// Without this, the two tests above cannot distinguish "the ClusterClient boundary works"
// from "something else in the process delivered the message". Same message, same client
// factory, same paths, same comm actor — and it must NOT arrive.
//
// HONEST SCOPE: this removes the receptionist extension AND, unavoidably, the RegisterService
// call — resolving ClusterClientReceptionist.Get to register would materialise the very
// extension whose absence is the control. So it falsifies "delivery happens without a
// receptionist boundary", not the extension line in isolation. Verified falsifiable: giving
// this node the extension and the registration turns the test red, which also proves the send
// below is live rather than silently misaddressed.
var deaf = StartSelfSeededSystem(withReceptionist: false);
try
{
await AwaitSingleMemberUpAsync(deaf);
var inbox = NewInbox();
deaf.EventStream.Subscribe(deaf.ActorOf(CaptureActor.Props(inbox)), typeof(RestartDriver));
// Deliberately NOT ClusterClientReceptionist.Get(deaf) — resolving the extension would
// materialise the very receptionist whose absence is the control.
var comm = deaf.ActorOf(NodeCommunicationActor.Props(null), MeshPaths.NodeCommunicationName);
comm.ShouldNotBeNull();
var client = _clientFactory.Create(_central, "deaf", ReceptionistOf(deaf));
var command = new RestartDriver("C1", "driver-1", "alice", Guid.NewGuid());
var deadline = DateTime.UtcNow + Silence;
while (DateTime.UtcNow < deadline && !inbox.Task.IsCompleted)
{
client.Tell(new ClusterClient.SendToAll(MeshPaths.NodeCommunication, command));
await Task.Delay(500, TestContext.Current.CancellationToken);
}
inbox.Task.IsCompleted.ShouldBeFalse(
"delivery survived removing the receptionist extension, so the two positive tests "
+ "above are not measuring the ClusterClient boundary at all");
}
finally
{
await ShutdownAsync(deaf);
}
}
[Fact]
public async Task One_dispatch_through_central_reaches_a_node_in_every_cluster_mesh()
{
// THE Phase 6 assertion at the transport level, driving the REAL CentralCommunicationActor.
// Two genuinely separate site meshes (site-A = _node, site-B created here), each its own
// single-member cluster with its own receptionist. Central must build ONE ClusterClient per
// cluster and fan SendToAll across them, so a single MeshCommand reaches BOTH meshes. A single
// fleet-wide client would reach only the one mesh whose receptionist answered — the exact
// failure this test is red-before-green against.
var siteB = StartSelfSeededSystem(withReceptionist: true);
try
{
await AwaitSingleMemberUpAsync(siteB);
RegisterNodeComm(siteB);
var inboxA = NewInbox();
_node.EventStream.Subscribe(_node.ActorOf(CaptureActor.Props(inboxA)), typeof(RestartDriver));
var inboxB = NewInbox();
siteB.EventStream.Subscribe(siteB.ActorOf(CaptureActor.Props(inboxB)), typeof(RestartDriver));
var dbFactory = SeedTwoClusterDb(SiteAddress(_node), SiteAddress(siteB));
var actor = _central.ActorOf(CentralCommunicationActor.Props(
dbFactory, ClusterClientOptions(), new DefaultMeshClusterClientFactory(), () => null));
var command = new RestartDriver("C1", "driver-1", "alice", Guid.NewGuid());
var deadline = DateTime.UtcNow + Arrival;
while (DateTime.UtcNow < deadline && !(inboxA.Task.IsCompleted && inboxB.Task.IsCompleted))
{
actor.Tell(new MeshCommand("deployments", command));
await Task.Delay(250, TestContext.Current.CancellationToken);
}
inboxA.Task.IsCompleted.ShouldBeTrue(
"site-A's mesh must receive the dispatch through its own per-cluster ClusterClient");
inboxB.Task.IsCompleted.ShouldBeTrue(
"site-B's mesh must receive the SAME dispatch through its own per-cluster ClusterClient; "
+ "a single fleet-wide client would have reached only one of the two meshes");
}
finally
{
await ShutdownAsync(siteB);
}
}
[Fact]
public async Task A_client_built_for_one_cluster_does_not_deliver_into_another_clusters_mesh()
{
// The partition proof for the fan-out: a client dialling ONLY site-A's contacts must land in
// site-A's mesh and NEVER cross into site-B's. This is why the per-cluster fan-out is correct
// rather than accidentally N x broadcasting — each client is scoped to its own mesh.
var siteB = StartSelfSeededSystem(withReceptionist: true);
try
{
await AwaitSingleMemberUpAsync(siteB);
RegisterNodeComm(siteB);
var inboxA = NewInbox();
_node.EventStream.Subscribe(_node.ActorOf(CaptureActor.Props(inboxA)), typeof(RestartDriver));
var inboxB = NewInbox();
siteB.EventStream.Subscribe(siteB.ActorOf(CaptureActor.Props(inboxB)), typeof(RestartDriver));
var clientA = _clientFactory.Create(_central, "site-a", ReceptionistOf(_node));
var command = new RestartDriver("C1", "driver-1", "alice", Guid.NewGuid());
// Keep sending for the full Silence window so any erroneous cross-delivery into site-B has
// every chance to arrive — a pass means the boundary is shut, not that the clock ran out.
var deadline = DateTime.UtcNow + Silence;
while (DateTime.UtcNow < deadline)
{
clientA.Tell(new ClusterClient.SendToAll(MeshPaths.NodeCommunication, command));
await Task.Delay(250, TestContext.Current.CancellationToken);
}
inboxA.Task.IsCompleted.ShouldBeTrue("site-A's client must deliver into site-A's mesh");
inboxB.Task.IsCompleted.ShouldBeFalse(
"a client built for site-A's contacts must NOT deliver into site-B's mesh — the "
+ "per-cluster fan-out must stay partitioned, one client per mesh");
}
finally
{
await ShutdownAsync(siteB);
}
}
/// Registers a real node-comm actor with a system's receptionist so central can reach it.
/// The site system to register into.
private static void RegisterNodeComm(ActorSystem system)
{
var comm = system.ActorOf(NodeCommunicationActor.Props(null), MeshPaths.NodeCommunicationName);
ClusterClientReceptionist.Get(system).RegisterService(comm);
}
private static (string Host, int Port) SiteAddress(ActorSystem system)
{
var addr = AkkaCluster.Get(system).SelfAddress;
return (addr.Host!, addr.Port!.Value);
}
private static MeshTransportOptions ClusterClientOptions() => new()
{
Mode = MeshTransportOptions.ModeClusterClient,
CentralContactPoints = ["akka.tcp://otopcua@127.0.0.1:4053"],
ContactRefreshSeconds = 1,
};
private static IDbContextFactory SeedTwoClusterDb(
(string Host, int Port) siteA, (string Host, int Port) siteB)
{
var factory = new InMemoryConfigDbFactory(Guid.NewGuid().ToString("N"));
using var db = factory.CreateDbContext();
db.ClusterNodes.Add(NodeRow("SITE-A", "node-a", siteA));
db.ClusterNodes.Add(NodeRow("SITE-B", "node-b", siteB));
db.SaveChanges();
return factory;
}
private static ClusterNode NodeRow(string clusterId, string nodeId, (string Host, int Port) addr) => new()
{
NodeId = nodeId,
ClusterId = clusterId,
Host = addr.Host,
AkkaPort = addr.Port,
ApplicationUri = $"urn:{nodeId}",
Enabled = true,
MaintenanceMode = false,
CreatedBy = "test",
};
private sealed class InMemoryConfigDbFactory(string dbName)
: IDbContextFactory
{
public OtOpcUaConfigDbContext CreateDbContext()
{
var opts = new DbContextOptionsBuilder()
.UseInMemoryDatabase(dbName)
.Options;
return new OtOpcUaConfigDbContext(opts);
}
}
///
/// Starts a single-member cluster on a dynamic port, optionally without the receptionist.
///
///
/// drops the receptionist provider from akka.extensions ,
/// leaving DistributedPubSub in place so only the one variable changes.
///
/// The started system.
private static ActorSystem StartSelfSeededSystem(bool withReceptionist)
{
var extensions = withReceptionist
? string.Empty
: "akka.extensions = [\"Akka.Cluster.Tools.PublishSubscribe.DistributedPubSubExtensionProvider, "
+ "Akka.Cluster.Tools\"]\n";
// Port 0 + a post-start self-Join, rather than a seed-nodes list: the address is not known
// until the transport binds, and hard-coding ports makes the test collide with itself under
// parallel execution.
var config = ConfigurationFactory.ParseString(
"akka.remote.dot-netty.tcp.hostname = \"127.0.0.1\"\n"
+ "akka.remote.dot-netty.tcp.public-hostname = \"127.0.0.1\"\n"
+ "akka.remote.dot-netty.tcp.port = 0\n"
+ "akka.cluster.seed-nodes = []\n"
+ "akka.loglevel = \"WARNING\"\n"
+ extensions)
.WithFallback(ZB.MOM.WW.OtOpcUa.Cluster.HoconLoader.LoadBaseConfig());
var system = ActorSystem.Create(SharedSystemName, config);
var cluster = AkkaCluster.Get(system);
cluster.Join(cluster.SelfAddress);
return system;
}
private static async Task AwaitSingleMemberUpAsync(ActorSystem system)
{
var cluster = AkkaCluster.Get(system);
var deadline = DateTime.UtcNow + TimeSpan.FromSeconds(30);
while (DateTime.UtcNow < deadline)
{
if (cluster.State.Members.Any(m => m.Status == MemberStatus.Up)) return;
await Task.Delay(100, TestContext.Current.CancellationToken);
}
cluster.State.Members.Any(m => m.Status == MemberStatus.Up).ShouldBeTrue(
$"{cluster.SelfAddress} never reached Up as its own seed");
}
///
/// Re-sends until the inbox completes or the deadline passes.
///
///
/// A ClusterClient establishes contact asynchronously and buffer-size = 0 means anything
/// sent before that lands is dropped on the floor — by design. A single send would therefore
/// race the handshake and fail intermittently for a reason that is not a defect. Re-sending
/// weakens nothing: the control test re-sends on exactly the same schedule and must still
/// receive nothing.
///
/// Actor to send to.
/// Builds the message for each attempt.
/// Completion source the receiving capture actor resolves.
/// How long to keep trying.
/// The delivered message, or if none arrived.
private static async Task SendUntilDeliveredAsync(
IActorRef target,
Func message,
TaskCompletionSource inbox,
TimeSpan timeout)
{
var deadline = DateTime.UtcNow + timeout;
while (DateTime.UtcNow < deadline)
{
if (inbox.Task.IsCompleted) return await inbox.Task;
target.Tell(message());
await Task.Delay(250, TestContext.Current.CancellationToken);
}
return inbox.Task.IsCompleted ? await inbox.Task : null;
}
private static TaskCompletionSource NewInbox() =>
new(TaskCreationOptions.RunContinuationsAsynchronously);
private static ImmutableHashSet ReceptionistOf(ActorSystem system) =>
ImmutableHashSet.Create(
new RootActorPath(AkkaCluster.Get(system).SelfAddress) / "system" / "receptionist");
private static async Task ShutdownAsync(ActorSystem? system)
{
if (system is null) return;
await system.Terminate().WaitAsync(TimeSpan.FromSeconds(15));
}
/// Completes an inbox with the first message it receives.
private sealed class CaptureActor : ReceiveActor
{
public static Props Props(TaskCompletionSource sink) =>
Akka.Actor.Props.Create(() => new CaptureActor(sink));
public CaptureActor(TaskCompletionSource sink) => ReceiveAny(m => sink.TrySetResult(m));
}
}