Files
lmxopcua/tests/Server/ZB.MOM.WW.OtOpcUa.Runtime.Tests/Drivers/DriverHostActorBootFromCacheTests.cs
T
Joseph Doherty 5cb0e72166 feat(mesh): node subscribers accept EventStream commands; _coordinatorOverride -> _ackRouter
Phase 2 Task 5. DriverHostActor and ScriptedAlarmHostActor now subscribe to the
node-local EventStream alongside their existing DPS topic subscriptions, so the
transport flag lives ENTIRELY on the publishing side -- exactly one of the two
ever publishes, and the switch can be flipped or reverted without touching a
single subscriber.

Note the alarm case is not symmetric: the DPS subscribe there still carries the
node-LOCAL publisher (OtOpcUaServerHostedService's OPC UA Part 9 method calls),
which never crosses the boundary and stays on DPS in both modes. Only the
central AdminUI leg moves.

Renamed DriverHostActor._coordinatorOverride to _ackRouter. Under ClusterClient
mode that ref is the NodeCommunicationActor, emphatically not a coordinator,
and the old name would send the next reader looking for one.

Adds two wiring tests, because the unit tests either side of this seam BOTH pass
with the subscription deleted -- the comm actor's test asserts a probe receives
the message, and the host tests drive the actor by direct Tell. Neither notices
a missing PreStart subscribe. Sabotage-verified: deleting either subscription
reddens exactly its own wiring test and nothing else.

Three things went wrong while writing those tests, all worth keeping:

- The first version raced. ActorOf returns before PreStart runs, and an
  EventStream.Publish with no subscriber is dropped silently -- no buffering, no
  dead letter. Fixed with a warm-up whose ack proves PreStart completed; the
  comment says why it is not ceremony.
- The alarm version was VACUOUS: it told the host directly INSIDE the assertion
  window alongside the relay, so the direct Tell alone produced the log the
  assertion waited for. Split into warm-up (outside) and relay (inside).
- The alarm test needs its own ActorSystem: the shared harness pins
  loglevel = WARNING and the only signal an unowned alarm gives is a Debug line.

Claude-Session: https://claude.ai/code/session_01GASWkNEi68FSCtvr6rLoEW
2026-07-22 12:19:29 -04:00

301 lines
12 KiB
C#

using System.Text.Json;
using Akka.Actor;
using Microsoft.EntityFrameworkCore;
using Shouldly;
using Xunit;
using ZB.MOM.WW.OtOpcUa.Commons.Interfaces;
using ZB.MOM.WW.OtOpcUa.Commons.Messages.Fleet;
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.Runtime.DeploymentCache;
using ZB.MOM.WW.OtOpcUa.Runtime.Drivers;
using ZB.MOM.WW.OtOpcUa.Runtime.OpcUa;
using ZB.MOM.WW.OtOpcUa.Runtime.Tests.Harness;
namespace ZB.MOM.WW.OtOpcUa.Runtime.Tests.Drivers;
/// <summary>
/// LocalDb Phase 1 (Task 8) — booting from the node-local artifact cache when central SQL is
/// unreachable.
/// </summary>
/// <remarks>
/// The behaviour that matters is asymmetric: the cache must rescue a node whose ConfigDb read
/// failed, and must be completely invisible otherwise. A cache consulted on the happy path would
/// be a route for stale configuration to override fresh configuration.
/// </remarks>
public sealed class DriverHostActorBootFromCacheTests : RuntimeActorTestBase
{
private static readonly NodeId TestNode = NodeId.Parse("driver-test");
private static readonly RevisionHash CachedRev = RevisionHash.Parse(new string('c', 64));
[Fact]
public void CentralUnreachableAndCacheHasArtifact_BootsFromCache()
{
var cached = new CachedDeploymentArtifact(
DeploymentId.NewId().ToString(),
CachedRev.Value,
EmptyArtifact(),
DateTimeOffset.UtcNow.AddHours(-1));
var cache = new StubArtifactCache(cached);
var actor = Sys.ActorOf(DriverHostActor.Props(
new ThrowingDbFactory(), TestNode, ackRouter: null,
localRoles: new HashSet<string> { "driver" },
deploymentArtifactCache: cache));
var snapshot = AskDiagnostics(actor);
snapshot.RunningFromCache.ShouldBeTrue();
snapshot.CurrentRevision.ShouldBe(CachedRev);
// Positive control for the UnkeyedReads counter the "never consults the cache" tests rely
// on. Without this, those ShouldBe(0) assertions would pass just as happily if the counter
// were never incremented at all.
cache.UnkeyedReads.ShouldBe(1);
}
[Fact]
public void BootFromCache_HandsTheCachedArtifactToTheAddressSpaceRebuild()
{
// The client-facing address space rebuild must materialise from the CACHED bytes, not from a
// ConfigDb read keyed by deployment id — that read is exactly what is unreachable during the
// outage boot-from-cache exists to survive. Without the blob, the node logs "running from
// cache", spins up its drivers, and still serves OPC UA clients an empty server. (Regression
// for the defect the docker-dev live gate caught: 16 nodes on the healthy peer, 1 here.)
var blob = EmptyArtifact();
var cached = new CachedDeploymentArtifact(
DeploymentId.NewId().ToString(), CachedRev.Value, blob, DateTimeOffset.UtcNow.AddHours(-1));
var publishProbe = CreateTestProbe();
Sys.ActorOf(DriverHostActor.Props(
new ThrowingDbFactory(), TestNode, ackRouter: null,
localRoles: new HashSet<string> { "driver" },
opcUaPublishActor: publishProbe.Ref,
deploymentArtifactCache: new StubArtifactCache(cached)));
var rebuild = publishProbe.ExpectMsg<OpcUaPublishActor.RebuildAddressSpace>(TimeSpan.FromSeconds(5));
rebuild.Artifact.ShouldNotBeNull();
rebuild.Artifact.ShouldBe(blob);
}
[Fact]
public void CentralUnreachableAndCacheEmpty_StaysStaleWithNoConfiguration()
{
// The pre-cache behaviour, unchanged. A cache miss must not invent a configuration.
var actor = Sys.ActorOf(DriverHostActor.Props(
new ThrowingDbFactory(), TestNode, ackRouter: null,
localRoles: new HashSet<string> { "driver" },
deploymentArtifactCache: new StubArtifactCache(current: null)));
var snapshot = AskDiagnostics(actor);
snapshot.RunningFromCache.ShouldBeFalse();
snapshot.CurrentRevision.ShouldBeNull();
}
[Fact]
public void CentralUnreachableAndNoCacheConfigured_StaysStale()
{
// Admin-only graphs and older harnesses pass null; behaviour is identical to a cache miss.
var actor = Sys.ActorOf(DriverHostActor.Props(
new ThrowingDbFactory(), TestNode, ackRouter: null,
localRoles: new HashSet<string> { "driver" }));
var snapshot = AskDiagnostics(actor);
snapshot.RunningFromCache.ShouldBeFalse();
snapshot.CurrentRevision.ShouldBeNull();
}
[Fact]
public void CentralReachable_NeverConsultsTheCache()
{
// THE guard against stale config winning. If the cache were read on the happy path, a node
// whose cache held an older revision could quietly serve it over the deployed one.
var cache = new StubArtifactCache(new CachedDeploymentArtifact(
DeploymentId.NewId().ToString(), CachedRev.Value, EmptyArtifact(), DateTimeOffset.UtcNow));
var actor = Sys.ActorOf(DriverHostActor.Props(
NewInMemoryDbFactory(), TestNode, ackRouter: null,
localRoles: new HashSet<string> { "driver" },
deploymentArtifactCache: cache));
var snapshot = AskDiagnostics(actor);
snapshot.RunningFromCache.ShouldBeFalse();
cache.UnkeyedReads.ShouldBe(0);
}
[Fact]
public void CentralReachableWithAPriorAppliedDeployment_NeverConsultsTheCache()
{
// The RestoreApplied path — the other way a boot can succeed against central.
var db = NewInMemoryDbFactory();
var rev = RevisionHash.Parse(new string('d', 64));
var deploymentId = SeedAppliedDeployment(db, rev);
var cache = new StubArtifactCache(new CachedDeploymentArtifact(
deploymentId.ToString(), CachedRev.Value, EmptyArtifact(), DateTimeOffset.UtcNow));
var actor = Sys.ActorOf(DriverHostActor.Props(
db, TestNode, ackRouter: null,
localRoles: new HashSet<string> { "driver" },
deploymentArtifactCache: cache));
var snapshot = AskDiagnostics(actor);
snapshot.RunningFromCache.ShouldBeFalse();
snapshot.CurrentRevision.ShouldBe(rev);
cache.UnkeyedReads.ShouldBe(0);
}
[Fact]
public void AThrowingCache_DegradesToStale_RatherThanFailingTheActor()
{
// A cache fault must leave the node exactly where it would have been without a cache.
var actor = Sys.ActorOf(DriverHostActor.Props(
new ThrowingDbFactory(), TestNode, ackRouter: null,
localRoles: new HashSet<string> { "driver" },
deploymentArtifactCache: new ThrowingReadCache()));
var snapshot = AskDiagnostics(actor);
snapshot.RunningFromCache.ShouldBeFalse();
snapshot.CurrentRevision.ShouldBeNull();
}
[Fact]
public void RestoringAnAppliedDeploymentOnBootstrap_RepopulatesTheCache()
{
// The cache invariant is "holds what the node serves". A node recovering an already-applied
// deployment (the RestoreApplied path) must (re)write its cache — otherwise a node whose cache
// was lost (fresh volume, disk failure) recovers its served state from central yet stays
// cache-less, unable to boot-from-cache on the NEXT outage. Replication does not heal it: a
// peer's already-acked rows are pruned from its oplog, so a wiped node is never back-filled.
// (Regression for the docker-dev live-gate finding.)
var db = NewInMemoryDbFactory();
var rev = RevisionHash.Parse(new string('e', 64));
var deploymentId = SeedAppliedDeployment(db, rev);
// A cache that starts empty — exactly the wiped-node state.
var cache = new StubArtifactCache(current: null);
var actor = Sys.ActorOf(DriverHostActor.Props(
db, TestNode, ackRouter: null,
localRoles: new HashSet<string> { "driver" },
deploymentArtifactCache: cache));
var snapshot = AskDiagnostics(actor);
snapshot.CurrentRevision.ShouldBe(rev); // recovered the applied deployment
snapshot.RunningFromCache.ShouldBeFalse(); // via central, not the cache
// The point: the restore path wrote the served artifact back into the cache.
AwaitAssert(
() => cache.Stores.ShouldBe(1),
duration: TimeSpan.FromSeconds(5));
cache.LastStoredDeploymentId.ShouldBe(deploymentId.ToString());
}
private NodeDiagnosticsSnapshot AskDiagnostics(IActorRef actor)
{
var probe = CreateTestProbe();
// GetDiagnostics is serviced in Steady, Applying AND Stale, so this observes the node
// whichever way the boot resolved.
AwaitAssert(
() =>
{
actor.Tell(new GetDiagnostics(CorrelationId.NewId()), probe.Ref);
probe.ExpectMsg<NodeDiagnosticsSnapshot>(TimeSpan.FromSeconds(2));
},
duration: TimeSpan.FromSeconds(5));
actor.Tell(new GetDiagnostics(CorrelationId.NewId()), probe.Ref);
return probe.ExpectMsg<NodeDiagnosticsSnapshot>(TimeSpan.FromSeconds(5));
}
/// <summary>A syntactically valid artifact with no driver instances.</summary>
private static byte[] EmptyArtifact()
=> JsonSerializer.SerializeToUtf8Bytes(new { DriverInstances = Array.Empty<object>() });
private static DeploymentId SeedAppliedDeployment(
IDbContextFactory<OtOpcUaConfigDbContext> db, RevisionHash rev)
{
var id = DeploymentId.NewId();
using var ctx = db.CreateDbContext();
ctx.Deployments.Add(new Deployment
{
DeploymentId = id.Value,
RevisionHash = rev.Value,
Status = DeploymentStatus.Sealed,
CreatedBy = "test",
SealedAtUtc = DateTime.UtcNow,
ArtifactBlob = EmptyArtifact(),
});
ctx.NodeDeploymentStates.Add(new NodeDeploymentState
{
NodeId = TestNode.Value,
DeploymentId = id.Value,
Status = NodeDeploymentStatus.Applied,
StartedAtUtc = DateTime.UtcNow,
});
ctx.SaveChanges();
return id;
}
/// <summary>Stands in for an unreachable central SQL Server.</summary>
private sealed class ThrowingDbFactory : IDbContextFactory<OtOpcUaConfigDbContext>
{
public OtOpcUaConfigDbContext CreateDbContext()
=> throw new InvalidOperationException("ConfigDb unreachable");
}
private sealed class StubArtifactCache(CachedDeploymentArtifact? current) : IDeploymentArtifactCache
{
private int _unkeyedReads;
private int _stores;
/// <summary>How many times the boot path consulted this cache.</summary>
public int UnkeyedReads => Volatile.Read(ref _unkeyedReads);
/// <summary>How many times the node wrote an applied artifact to this cache.</summary>
public int Stores => Volatile.Read(ref _stores);
/// <summary>The deployment id of the most recent store.</summary>
public string? LastStoredDeploymentId { get; private set; }
public Task StoreAsync(string clusterId, string deploymentId, string revisionHash,
byte[] artifact, CancellationToken ct = default)
{
Interlocked.Increment(ref _stores);
LastStoredDeploymentId = deploymentId;
return Task.CompletedTask;
}
public Task<CachedDeploymentArtifact?> GetCurrentAsync(
string clusterId, CancellationToken ct = default)
=> Task.FromResult(current);
public Task<CachedDeploymentArtifact?> GetCurrentUnkeyedAsync(CancellationToken ct = default)
{
Interlocked.Increment(ref _unkeyedReads);
return Task.FromResult(current);
}
}
private sealed class ThrowingReadCache : IDeploymentArtifactCache
{
public Task StoreAsync(string clusterId, string deploymentId, string revisionHash,
byte[] artifact, CancellationToken ct = default)
=> Task.CompletedTask;
public Task<CachedDeploymentArtifact?> GetCurrentAsync(
string clusterId, CancellationToken ct = default)
=> throw new InvalidOperationException("cache unreadable");
public Task<CachedDeploymentArtifact?> GetCurrentUnkeyedAsync(CancellationToken ct = default)
=> throw new InvalidOperationException("cache unreadable");
}
}