d5b5cb6ede
Phase 2 Task 3. Central-side end of the boundary: routes every central->node command out over DPS or ClusterClient per MeshTransport:Mode, and forwards inbound ApplyAcks to the deploy coordinator singleton (Forward, not Tell, so the coordinator sees the node rather than this relay). Three things here are load-bearing and easy to get wrong later: SendToAll, never Send. Today's DPS publish reaches EVERY DriverHostActor and the node side has no ClusterId or node filter to compensate -- scoping happens later, inside the artifact. ClusterClient.Send delivers to exactly ONE registered actor, so it would deploy to a single node while every other node silently kept its old config, and the deployment could still seal green. Sabotage-verified: swapping to Send reddens exactly that one test. Exactly ONE ClusterClient, fleet-wide. A receptionist serves its whole cluster, so SendToAll reaches every registered node-comm actor in the mesh regardless of which contact point was dialled. One client per application Cluster -- the shape Phase 6 wants -- would fan each command out once per cluster while the fleet is still one mesh: N x duplicate DispatchDeployment and ApplyAck. Marked TODO(Phase 6). The contact set does NOT scope delivery. Excluding a maintenance-mode node from the contacts does not stop it receiving commands; it still receives them, as it does under today's broadcast. The filter is about which receptionists are worth dialling, not about who gets the message. Stated in the code because the opposite is the natural assumption. Also carries the sister project's shipped fixes: per-row address parsing so one malformed ClusterNode cannot abort the refresh and leave the contact set half-built (regression test orders the bad row FIRST), the cache entry cleared before recreate so a failed create cannot leave commands routing into a stopping actor, and a Status.Failure handler so a DB outage is a Warning rather than a silent debug-level unhandled message. Two test-harness bugs found and fixed while writing this, both mine: SubscribeAck goes to the SENDER (so the mediator subscribe needed an explicit sender), and EventFilter matches case-INSENSITIVELY, so a "was NOT" filter also caught this actor's own "was not created" warning. Claude-Session: https://claude.ai/code/session_01GASWkNEi68FSCtvr6rLoEW
272 lines
11 KiB
C#
272 lines
11 KiB
C#
using System.Collections.Immutable;
|
|
using Akka.Actor;
|
|
using Akka.Cluster.Tools.Client;
|
|
using Akka.Cluster.Tools.PublishSubscribe;
|
|
using Akka.TestKit;
|
|
using Microsoft.EntityFrameworkCore;
|
|
using Shouldly;
|
|
using Xunit;
|
|
using ZB.MOM.WW.OtOpcUa.Cluster;
|
|
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.Configuration.Enums;
|
|
using ZB.MOM.WW.OtOpcUa.ControlPlane.Communication;
|
|
using ZB.MOM.WW.OtOpcUa.ControlPlane.Tests.Harness;
|
|
|
|
namespace ZB.MOM.WW.OtOpcUa.ControlPlane.Tests.Communication;
|
|
|
|
public class CentralCommunicationActorTests : ControlPlaneActorTestBase
|
|
{
|
|
private const string DeploymentsTopic = "deployments";
|
|
|
|
/// <summary>Records every client the actor asks for, and hands back a probe in its place.</summary>
|
|
private sealed class RecordingClientFactory : IMeshClusterClientFactory
|
|
{
|
|
private readonly Func<TestProbe> _newProbe;
|
|
|
|
public RecordingClientFactory(Func<TestProbe> newProbe) => _newProbe = newProbe;
|
|
|
|
public List<ImmutableHashSet<ActorPath>> Calls { get; } = [];
|
|
|
|
public List<TestProbe> Probes { get; } = [];
|
|
|
|
public IActorRef Create(ActorSystem system, ImmutableHashSet<ActorPath> contacts)
|
|
{
|
|
Calls.Add(contacts);
|
|
var probe = _newProbe();
|
|
Probes.Add(probe);
|
|
return probe.Ref;
|
|
}
|
|
}
|
|
|
|
private static MeshTransportOptions ClusterClientMode(int refreshSeconds = 60) => new()
|
|
{
|
|
Mode = MeshTransportOptions.ModeClusterClient,
|
|
CentralContactPoints = ["akka.tcp://otopcua@central-1:4053"],
|
|
ContactRefreshSeconds = refreshSeconds,
|
|
};
|
|
|
|
private static void SeedCluster(
|
|
IDbContextFactory<OtOpcUaConfigDbContext> factory,
|
|
string clusterId,
|
|
params (string NodeId, string Host, int AkkaPort, bool Enabled, bool Maintenance)[] nodes)
|
|
{
|
|
using var db = factory.CreateDbContext();
|
|
db.ServerClusters.Add(new ServerCluster
|
|
{
|
|
ClusterId = clusterId,
|
|
Name = clusterId,
|
|
Enterprise = "ZB",
|
|
Site = "Site",
|
|
// Warm/Hot require NodeCount = 2 (CK_ServerCluster_RedundancyMode_NodeCount); these
|
|
// fixtures only need the row to exist, so use the standalone shape.
|
|
RedundancyMode = RedundancyMode.None,
|
|
NodeCount = 1,
|
|
CreatedBy = "test",
|
|
});
|
|
foreach (var n in nodes)
|
|
{
|
|
db.ClusterNodes.Add(new ClusterNode
|
|
{
|
|
NodeId = n.NodeId,
|
|
ClusterId = clusterId,
|
|
Host = n.Host,
|
|
AkkaPort = n.AkkaPort,
|
|
ApplicationUri = $"urn:{n.NodeId}",
|
|
Enabled = n.Enabled,
|
|
MaintenanceMode = n.Maintenance,
|
|
CreatedBy = "test",
|
|
});
|
|
}
|
|
|
|
db.SaveChanges();
|
|
}
|
|
|
|
private IActorRef Spawn(
|
|
IDbContextFactory<OtOpcUaConfigDbContext> factory,
|
|
MeshTransportOptions options,
|
|
IMeshClusterClientFactory clientFactory,
|
|
IActorRef? coordinator = null) =>
|
|
Sys.ActorOf(CentralCommunicationActor.Props(factory, options, clientFactory, () => coordinator));
|
|
|
|
[Fact]
|
|
public void Dps_mode_publishes_on_the_carried_topic_and_creates_no_client()
|
|
{
|
|
var factory = NewInMemoryDbFactory();
|
|
var clients = new RecordingClientFactory(() => CreateTestProbe());
|
|
var mediatorProbe = CreateTestProbe();
|
|
// Subscribe on the probe's behalf: SubscribeAck goes to the SENDER, so telling the mediator
|
|
// without an explicit sender would ack to TestActor and this wait would hang.
|
|
DistributedPubSub.Get(Sys).Mediator.Tell(
|
|
new Subscribe(DeploymentsTopic, mediatorProbe.Ref), mediatorProbe.Ref);
|
|
mediatorProbe.ExpectMsg<SubscribeAck>(TimeSpan.FromSeconds(5));
|
|
|
|
var actor = Spawn(factory, new MeshTransportOptions(), clients);
|
|
var dispatch = new DispatchDeployment(
|
|
new DeploymentId(Guid.NewGuid()), new RevisionHash("abc"), new CorrelationId(Guid.NewGuid()));
|
|
|
|
actor.Tell(new MeshCommand(DeploymentsTopic, dispatch));
|
|
|
|
mediatorProbe.ExpectMsg<DispatchDeployment>(TimeSpan.FromSeconds(5)).ShouldBe(dispatch);
|
|
clients.Calls.ShouldBeEmpty();
|
|
}
|
|
|
|
[Fact]
|
|
public void ClusterClient_mode_sends_to_all_not_to_one()
|
|
{
|
|
// THE assertion of this phase. ClusterClient.Send delivers to exactly ONE registered actor;
|
|
// today's DPS publish reaches EVERY DriverHostActor, and the node side has no ClusterId or
|
|
// node filter to compensate. Send would deploy to a single node while every other node
|
|
// silently kept its old configuration.
|
|
var factory = NewInMemoryDbFactory();
|
|
SeedCluster(factory, "C1", ("host-a:4053", "host-a", 4053, true, false));
|
|
var clients = new RecordingClientFactory(() => CreateTestProbe());
|
|
|
|
var actor = Spawn(factory, ClusterClientMode(), clients);
|
|
AwaitCondition(() => clients.Probes.Count == 1, TimeSpan.FromSeconds(5));
|
|
|
|
var dispatch = new DispatchDeployment(
|
|
new DeploymentId(Guid.NewGuid()), new RevisionHash("abc"), new CorrelationId(Guid.NewGuid()));
|
|
actor.Tell(new MeshCommand(DeploymentsTopic, dispatch));
|
|
|
|
var sent = clients.Probes[0].ExpectMsg<ClusterClient.SendToAll>(TimeSpan.FromSeconds(5));
|
|
sent.Path.ShouldBe(MeshPaths.NodeCommunication);
|
|
sent.Message.ShouldBe(dispatch);
|
|
}
|
|
|
|
[Fact]
|
|
public void Exactly_one_client_is_created_for_a_multi_cluster_fleet()
|
|
{
|
|
// A receptionist serves its whole mesh, so one client per Cluster would fan each command out
|
|
// once per cluster while the fleet is still a single mesh — N x duplicate deployments.
|
|
var factory = NewInMemoryDbFactory();
|
|
SeedCluster(factory, "C1", ("host-a:4053", "host-a", 4053, true, false));
|
|
SeedCluster(factory, "C2", ("host-b:4053", "host-b", 4053, true, false));
|
|
var clients = new RecordingClientFactory(() => CreateTestProbe());
|
|
|
|
Spawn(factory, ClusterClientMode(), clients);
|
|
AwaitCondition(() => clients.Calls.Count == 1, TimeSpan.FromSeconds(5));
|
|
|
|
// ...and that one client dials BOTH clusters' nodes.
|
|
clients.Calls[0].Count.ShouldBe(2);
|
|
// Give a second refresh a chance to (wrongly) create another.
|
|
ExpectNoMsg(TimeSpan.FromMilliseconds(300));
|
|
clients.Calls.Count.ShouldBe(1);
|
|
}
|
|
|
|
[Fact]
|
|
public void Disabled_and_maintenance_nodes_are_not_dialled()
|
|
{
|
|
var factory = NewInMemoryDbFactory();
|
|
SeedCluster(
|
|
factory,
|
|
"C1",
|
|
("host-a:4053", "host-a", 4053, true, false),
|
|
("host-b:4053", "host-b", 4053, false, false),
|
|
("host-c:4053", "host-c", 4053, true, true));
|
|
var clients = new RecordingClientFactory(() => CreateTestProbe());
|
|
|
|
Spawn(factory, ClusterClientMode(), clients);
|
|
AwaitCondition(() => clients.Calls.Count == 1, TimeSpan.FromSeconds(5));
|
|
|
|
clients.Calls[0].Count.ShouldBe(1);
|
|
clients.Calls[0].Single().ToString().ShouldContain("host-a");
|
|
}
|
|
|
|
[Fact]
|
|
public void A_malformed_node_row_does_not_abort_the_refresh()
|
|
{
|
|
// The bad row is ordered FIRST on purpose: the failure this guards against is an exception
|
|
// that aborts the loop and leaves the contact set half-built.
|
|
var factory = NewInMemoryDbFactory();
|
|
SeedCluster(
|
|
factory,
|
|
"C1",
|
|
("bad", "not a host!!", 4053, true, false),
|
|
("host-a:4053", "host-a", 4053, true, false));
|
|
var clients = new RecordingClientFactory(() => CreateTestProbe());
|
|
|
|
Spawn(factory, ClusterClientMode(), clients);
|
|
AwaitCondition(() => clients.Calls.Count == 1, TimeSpan.FromSeconds(5));
|
|
|
|
clients.Calls[0].ShouldContain(p => p.ToString().Contains("host-a"));
|
|
}
|
|
|
|
[Fact]
|
|
public void A_command_with_no_client_is_dropped_with_a_warning()
|
|
{
|
|
// No ClusterNode rows at all: nothing to dial, so the command must be dropped loudly rather
|
|
// than queued. buffer-size = 0 covers the case where a client exists but cannot connect;
|
|
// this covers the case where there is no client at all.
|
|
var factory = NewInMemoryDbFactory();
|
|
var clients = new RecordingClientFactory(() => CreateTestProbe());
|
|
var actor = Spawn(factory, ClusterClientMode(), clients);
|
|
|
|
EventFilter.Warning(contains: "dropping").ExpectOne(TimeSpan.FromSeconds(5), () =>
|
|
actor.Tell(new MeshCommand(DeploymentsTopic, "anything")));
|
|
}
|
|
|
|
[Fact]
|
|
public void ApplyAck_is_forwarded_to_the_coordinator_preserving_the_sender()
|
|
{
|
|
var factory = NewInMemoryDbFactory();
|
|
var clients = new RecordingClientFactory(() => CreateTestProbe());
|
|
var coordinator = CreateTestProbe();
|
|
var node = CreateTestProbe();
|
|
var actor = Spawn(factory, ClusterClientMode(), clients, coordinator.Ref);
|
|
|
|
var ack = new ApplyAck(
|
|
new DeploymentId(Guid.NewGuid()),
|
|
new NodeId("host-a:4053"),
|
|
ApplyAckOutcome.Applied,
|
|
null,
|
|
new CorrelationId(Guid.NewGuid()));
|
|
|
|
actor.Tell(ack, node.Ref);
|
|
|
|
coordinator.ExpectMsg<ApplyAck>(TimeSpan.FromSeconds(5)).ShouldBe(ack);
|
|
// Forward, not Tell: the coordinator must see the node, not this relay.
|
|
coordinator.LastSender.ShouldBe(node.Ref);
|
|
}
|
|
|
|
[Fact]
|
|
public void An_unresolvable_coordinator_logs_rather_than_losing_the_ack_silently()
|
|
{
|
|
var factory = NewInMemoryDbFactory();
|
|
var clients = new RecordingClientFactory(() => CreateTestProbe());
|
|
var actor = Spawn(factory, ClusterClientMode(), clients, coordinator: null);
|
|
|
|
var ack = new ApplyAck(
|
|
new DeploymentId(Guid.NewGuid()),
|
|
new NodeId("host-a:4053"),
|
|
ApplyAckOutcome.Applied,
|
|
null,
|
|
new CorrelationId(Guid.NewGuid()));
|
|
|
|
EventFilter.Warning(contains: "not resolvable").ExpectOne(TimeSpan.FromSeconds(5), () =>
|
|
actor.Tell(ack));
|
|
}
|
|
|
|
[Fact]
|
|
public void A_faulted_contact_load_warns_and_leaves_the_actor_alive()
|
|
{
|
|
var factory = NewInMemoryDbFactory();
|
|
var clients = new RecordingClientFactory(() => CreateTestProbe());
|
|
var actor = Spawn(factory, ClusterClientMode(), clients);
|
|
|
|
// Filter on a phrase unique to the DB-failure warning. EventFilter matches case-INSENSITIVELY,
|
|
// so a looser "was NOT" also caught this actor's other warning ("the mesh ClusterClient was
|
|
// not created") and the assertion failed with two matches instead of one.
|
|
EventFilter.Warning(contains: "contact set was NOT refreshed").ExpectOne(
|
|
TimeSpan.FromSeconds(5),
|
|
() => actor.Tell(new Status.Failure(new InvalidOperationException("db down"))));
|
|
|
|
// Still serving: a transient DB outage must not take the command boundary down with it.
|
|
actor.Tell(new MeshCommand(DeploymentsTopic, "anything"));
|
|
ExpectNoMsg(TimeSpan.FromMilliseconds(200));
|
|
}
|
|
}
|