From b0c031f09f412c8f25a27a64707f82ff17c1a2e1 Mon Sep 17 00:00:00 2001 From: Joseph Doherty Date: Tue, 21 Jul 2026 02:40:14 -0400 Subject: [PATCH] fix(drivers): don't tear down field I/O when a deploy's artifact reads back empty (#485) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The driver-side half of the same mistake the address-space fix corrected. ReconcileDrivers and PushDesiredSubscriptionsFromArtifact both read the artifact as `...FirstOrDefault() ?? Array.Empty()`. Their THROW paths already returned early, but empty bytes flowed onwards as a real answer: zero driver specs, so DriverSpawnPlanner planned every running child for StopChild, and then an empty desired set dropped each surviving driver's live subscription handle. A node lost its entire field I/O and still ACKed Applied. Same rule as the address space: no bytes is no answer, not "a configuration with no drivers". Nothing legitimate produces a zero-length blob — deleting the last driver still deploys a JSON document with empty arrays. Both sites now skip and keep what is running. The two guards are independent because PushDesiredSubscriptions does its OWN ConfigDb read, so the row can go missing between them. Also corrects the ReconcileDrivers doc comment, which claimed an empty blob made it "effectively a no-op" — with children running it was the opposite. Tests: DriverHostActorUnreadableArtifactTests, RED-first (verified failing — after the empty dispatch the driver list was empty). A zero-length ArtifactBlob reproduces byte-for-byte what the missing-row case delivers, so the race does not have to be constructed. Two controls, both load-bearing: - a READABLE driverless artifact still stops the driver, so the guard keys on "the read gave us nothing", not on "fewer drivers than before"; - dropping a driver's LAST TAG still clears its subscription. This one was added after the absence assertion was caught passing on a race: the unsubscribe is an async self-tell, so with the second guard deleted the test still went green. The control observes that same unsubscribe arriving well inside the settle window, which is what makes the absence meaningful — with the guard deleted the suite now correctly goes RED. SubscribableStubDriver gains an UnsubscribeCount for that observation. Claude-Session: https://claude.ai/code/session_01GASWkNEi68FSCtvr6rLoEW --- .../2026-07-21-localdb-followups-live-gate.md | 8 + .../Drivers/DriverHostActor.cs | 32 ++- .../DriverHostActorUnreadableArtifactTests.cs | 233 ++++++++++++++++++ .../Drivers/StubDrivers.cs | 5 + 4 files changed, 275 insertions(+), 3 deletions(-) create mode 100644 tests/Server/ZB.MOM.WW.OtOpcUa.Runtime.Tests/Drivers/DriverHostActorUnreadableArtifactTests.cs diff --git a/docs/plans/2026-07-21-localdb-followups-live-gate.md b/docs/plans/2026-07-21-localdb-followups-live-gate.md index d3cc61a0..d48e11fb 100644 --- a/docs/plans/2026-07-21-localdb-followups-live-gate.md +++ b/docs/plans/2026-07-21-localdb-followups-live-gate.md @@ -76,3 +76,11 @@ failed. The log excerpt above is now impossible: the `PureRemove` never runs, an `OpcUaPublishActorRebuildTests.Rebuild_keeps_the_last_known_good_address_space_when_the_artifact_load_fails`, with a positive control proving a *readable* artifact that genuinely drops an equipment still tears its subtree down. + +The **driver-side half** of the same mistake was fixed alongside it: `DriverHostActor.ReconcileDrivers` +and `PushDesiredSubscriptionsFromArtifact` read the artifact the same way, so empty bytes meant zero +driver specs — `DriverSpawnPlanner` planned every running child for `StopChild` — and then an empty +desired set dropped each surviving driver's live subscription handle. A node lost its entire field I/O +and still ACKed Applied. Both now skip on empty. Covered by `DriverHostActorUnreadableArtifactTests`, +whose absence assertion is calibrated by a control that observes the very same unsubscribe *arriving* +inside the settle window (without it the assertion passed on a race). diff --git a/src/Server/ZB.MOM.WW.OtOpcUa.Runtime/Drivers/DriverHostActor.cs b/src/Server/ZB.MOM.WW.OtOpcUa.Runtime/Drivers/DriverHostActor.cs index e3b49eba..4c382570 100644 --- a/src/Server/ZB.MOM.WW.OtOpcUa.Runtime/Drivers/DriverHostActor.cs +++ b/src/Server/ZB.MOM.WW.OtOpcUa.Runtime/Drivers/DriverHostActor.cs @@ -1618,9 +1618,9 @@ public sealed class DriverHostActor : ReceiveActor, IWithTimers /// /// Read the deployment artifact + reconcile the set of running /// children. Spawn missing, ApplyDelta on config change, stop removed/disabled drivers. - /// When the artifact blob is empty (legacy ControlPlane tests, smoke fixtures) or the - /// configured can't materialise any of the requested - /// types, this is effectively a no-op. + /// When the artifact blob is empty (legacy ControlPlane tests, smoke fixtures) the reconcile is + /// SKIPPED outright — see the issue #485 guard below — and when the configured + /// can't materialise any of the requested types this is a no-op. /// /// /// The artifact blob that was reconciled, or when it could not be @@ -1653,6 +1653,20 @@ public sealed class DriverHostActor : ReceiveActor, IWithTimers return null; } + // Issue #485 — no bytes is no answer, not "a configuration with no drivers". Parsing zero bytes + // yields zero specs, and DriverSpawnPlanner then plans EVERY running child for StopChild: a missing + // deployment row (or a half-written one) silently takes this node's entire field I/O down while the + // apply still ACKs Applied. Nothing legitimate produces a zero-length blob — deleting the last + // driver still deploys a JSON document with empty arrays — so treat it exactly like the load + // failure above and leave the running children alone. + if (blob.Length == 0) + { + _log.Warning( + "DriverHost {Node}: artifact for {Id} read back empty; skipping reconcile and keeping the {Count} running driver(s)", + _localNode, deploymentId, _children.Count); + return null; + } + var specs = DeploymentArtifact.ParseDriverInstances(blob, _localNode.Value); var snapshots = _children.ToDictionary( kv => kv.Key, @@ -1814,6 +1828,18 @@ public sealed class DriverHostActor : ReceiveActor, IWithTimers /// private void PushDesiredSubscriptionsFromArtifact(DeploymentId deploymentId, byte[] blob) { + // Issue #485 — the same "no bytes is no answer" rule as ReconcileDrivers. An empty artifact parses + // to a composition with no raw tags, which hands every child an EMPTY desired set — and an empty set + // drops the live subscription handle rather than re-subscribing. Reachable independently of the + // reconcile guard: this method does its OWN ConfigDb read, so the row can go missing between the two. + if (blob.Length == 0) + { + _log.Warning( + "DriverHost {Node}: artifact for {Id} read back empty; skipping SubscribeBulk and keeping the live subscriptions", + _localNode, deploymentId); + return; + } + AddressSpaceComposition composition; try { diff --git a/tests/Server/ZB.MOM.WW.OtOpcUa.Runtime.Tests/Drivers/DriverHostActorUnreadableArtifactTests.cs b/tests/Server/ZB.MOM.WW.OtOpcUa.Runtime.Tests/Drivers/DriverHostActorUnreadableArtifactTests.cs new file mode 100644 index 00000000..5749b2fc --- /dev/null +++ b/tests/Server/ZB.MOM.WW.OtOpcUa.Runtime.Tests/Drivers/DriverHostActorUnreadableArtifactTests.cs @@ -0,0 +1,233 @@ +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.Deploy; +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.Core.Abstractions; +using ZB.MOM.WW.OtOpcUa.Runtime.Drivers; +using ZB.MOM.WW.OtOpcUa.Runtime.Tests.Harness; + +namespace ZB.MOM.WW.OtOpcUa.Runtime.Tests.Drivers; + +/// +/// Issue #485 (driver-side half) — an artifact the host could not obtain must not be mistaken for a +/// configuration that contains nothing. +/// +/// +/// +/// ReconcileDrivers and PushDesiredSubscriptions both read the artifact as +/// …FirstOrDefault() ?? Array.Empty<byte>(). The throw path in each already +/// returns early, but a row that is absent (or carries a zero-length blob) yields empty bytes that +/// flow onwards as a real answer: zero driver specs, so DriverSpawnPlanner plans every +/// running child for StopChild, and then an empty desired-subscription set drops each +/// surviving driver's live handle. A node loses its entire field I/O and still ACKs Applied. +/// +/// +/// The same reasoning as the address-space half applies: nothing legitimate produces a zero-length +/// blob — an operator who really deletes every driver still deploys a JSON document with empty +/// arrays — so "no bytes" can only mean the read did not answer. Seeding a row whose +/// ArtifactBlob is empty reproduces exactly the bytes the missing-row case delivers, without +/// needing to race a delete against the two reads. +/// +/// +public sealed class DriverHostActorUnreadableArtifactTests : RuntimeActorTestBase +{ + private static readonly NodeId TestNode = NodeId.Parse("driver-test"); + private static readonly RevisionHash RevA = RevisionHash.Parse(new string('a', 64)); + private static readonly RevisionHash RevB = RevisionHash.Parse(new string('b', 64)); + private static readonly TimeSpan Timeout = TimeSpan.FromSeconds(5); + + /// How long to let post-ACK subscription traffic settle before asserting it did NOT happen. + private static readonly TimeSpan SettleWindow = TimeSpan.FromSeconds(2); + + private const string SpeedRef = "Plant/Modbus/dev1/speed"; + + /// A dispatch whose artifact reads back as no bytes must leave the running drivers — and their + /// live subscriptions — exactly as they were. + [Fact] + public void Dispatch_whose_artifact_reads_back_empty_keeps_the_drivers_and_their_subscriptions() + { + var db = NewInMemoryDbFactory(); + var factory = new SubscribingDriverFactory("Modbus"); + var (actor, coordinator) = SpawnHostAndApply(db, SeedV3Deployment(db, RevA), factory); + + // Baseline: the child is up and subscribed to its one raw tag. + AwaitAssert(() => factory.LastSubscribedRefs.ShouldNotBeNull()!.ShouldContain(SpeedRef), duration: Timeout); + AskDiagnostics(actor).Drivers.Select(d => d.Name).ShouldContain("drv-1"); + + // The next dispatch's artifact comes back as no bytes at all. + actor.Tell(new DispatchDeployment(SeedUnreadableDeployment(db, RevB), RevB, CorrelationId.NewId())); + coordinator.ExpectMsg(Timeout); + + // The ACK precedes the SubscribeBulk pass, and the child's unsubscribe is a further async self-tell, + // so settle before asserting an ABSENCE. SettleWindow is calibrated by + // , which observes the very same + // unsubscribe ARRIVING well inside it — without that control this assertion would pass on a race. + AskDiagnostics(actor); // ordering barrier through the host's mailbox + Thread.Sleep(SettleWindow); + + AskDiagnostics(actor).Drivers.Select(d => d.Name).ShouldContain("drv-1"); // not stopped + factory.UnsubscribeCount.ShouldBe(0); // handle not dropped + factory.LastSubscribedRefs.ShouldNotBeNull()!.ShouldContain(SpeedRef); + } + + /// Positive control for the subscription half: a READABLE artifact that keeps the driver but + /// drops its last tag genuinely does clear the live subscription — the child receives an empty desired + /// set and unsubscribes. This proves both that the unsubscribe is observable through this factory and + /// that it lands well within , so the absence asserted above is real. + [Fact] + public void Dropping_a_drivers_last_tag_does_clear_its_subscription() + { + var db = NewInMemoryDbFactory(); + var factory = new SubscribingDriverFactory("Modbus"); + var (actor, coordinator) = SpawnHostAndApply(db, SeedV3Deployment(db, RevA), factory); + + AwaitAssert(() => factory.LastSubscribedRefs.ShouldNotBeNull()!.ShouldContain(SpeedRef), duration: Timeout); + factory.UnsubscribeCount.ShouldBe(0); + + actor.Tell(new DispatchDeployment(SeedTaglessDriverDeployment(db, RevB), RevB, CorrelationId.NewId())); + coordinator.ExpectMsg(Timeout); + + AwaitAssert(() => factory.UnsubscribeCount.ShouldBe(1), duration: SettleWindow); + AskDiagnostics(actor).Drivers.Select(d => d.Name).ShouldContain("drv-1"); // the driver itself stayed + } + + /// Positive control: the guard keys on "the read gave us nothing", not on "the new config has + /// fewer drivers". A READABLE artifact that genuinely drops the driver still stops it — otherwise the + /// test above would pass just as happily against a host that had stopped reconciling altogether. + [Fact] + public void Dispatch_whose_readable_artifact_genuinely_drops_the_driver_still_stops_it() + { + var db = NewInMemoryDbFactory(); + var factory = new SubscribingDriverFactory("Modbus"); + var (actor, coordinator) = SpawnHostAndApply(db, SeedV3Deployment(db, RevA), factory); + + AwaitAssert(() => factory.LastSubscribedRefs.ShouldNotBeNull()!.ShouldContain(SpeedRef), duration: Timeout); + + // A real artifact that simply carries no drivers — an operator deleting the last driver. + actor.Tell(new DispatchDeployment(SeedDriverlessDeployment(db, RevB), RevB, CorrelationId.NewId())); + coordinator.ExpectMsg(Timeout); + + AwaitAssert( + () => AskDiagnostics(actor).Drivers.Select(d => d.Name).ShouldNotContain("drv-1"), + duration: Timeout); + } + + /// Spawns the host with the subscribing factory, dispatches and + /// waits for the Applied ACK so the child + the initial SubscribeBulk pass are in place. + private (IActorRef Actor, Akka.TestKit.TestProbe Coordinator) SpawnHostAndApply( + IDbContextFactory db, DeploymentId deploymentId, IDriverFactory factory) + { + var coordinator = CreateTestProbe(); + var actor = Sys.ActorOf(DriverHostActor.Props( + db, TestNode, coordinator.Ref, + driverFactory: factory, + localRoles: new HashSet { "driver" })); + + actor.Tell(new DispatchDeployment(deploymentId, RevA, CorrelationId.NewId())); + coordinator.ExpectMsg(Timeout).Outcome.ShouldBe(ApplyAckOutcome.Applied); + return (actor, coordinator); + } + + private NodeDiagnosticsSnapshot AskDiagnostics(IActorRef actor) + { + var probe = CreateTestProbe(); + actor.Tell(new GetDiagnostics(CorrelationId.NewId()), probe.Ref); + return probe.ExpectMsg(Timeout); + } + + /// Seeds a Sealed v3 deployment: RawFolder "Plant" → DriverInstance "drv-1" (Modbus, ENABLED so a + /// real child spawns) → Device "dev1" → Tag "speed" (RawPath Plant/Modbus/dev1/speed). + private static DeploymentId SeedV3Deployment(IDbContextFactory db, RevisionHash rev) => + SeedDeployment(db, rev, JsonSerializer.SerializeToUtf8Bytes(new + { + RawFolders = new[] { new { RawFolderId = "rf-plant", ParentRawFolderId = (string?)null, Name = "Plant", ClusterId = "c1" } }, + DriverInstances = new[] + { + new { DriverInstanceId = "drv-1", RawFolderId = "rf-plant", Name = "Modbus", DriverType = "Modbus", DriverConfig = "{}", ClusterId = "c1", Enabled = true }, + }, + Devices = new[] { new { DeviceId = "dev-1", DriverInstanceId = "drv-1", Name = "dev1", DeviceConfig = "{}" } }, + TagGroups = Array.Empty(), + Tags = new[] + { + new { TagId = "t-speed", DeviceId = "dev-1", TagGroupId = (string?)null, Name = "speed", DataType = "Double", AccessLevel = 1, TagConfig = "{}" }, + }, + })); + + /// Seeds a Sealed deployment whose ArtifactBlob is zero-length — byte-for-byte what the + /// host's ?? Array.Empty<byte>() hands downstream when the row cannot be found. + private static DeploymentId SeedUnreadableDeployment(IDbContextFactory db, RevisionHash rev) => + SeedDeployment(db, rev, Array.Empty()); + + /// Seeds a Sealed deployment carrying a real, readable artifact that keeps drv-1 (so its child + /// survives the reconcile) but declares no tags for it — the driver's desired set becomes empty. + private static DeploymentId SeedTaglessDriverDeployment(IDbContextFactory db, RevisionHash rev) => + SeedDeployment(db, rev, JsonSerializer.SerializeToUtf8Bytes(new + { + RawFolders = new[] { new { RawFolderId = "rf-plant", ParentRawFolderId = (string?)null, Name = "Plant", ClusterId = "c1" } }, + DriverInstances = new[] + { + new { DriverInstanceId = "drv-1", RawFolderId = "rf-plant", Name = "Modbus", DriverType = "Modbus", DriverConfig = "{}", ClusterId = "c1", Enabled = true }, + }, + Devices = new[] { new { DeviceId = "dev-1", DriverInstanceId = "drv-1", Name = "dev1", DeviceConfig = "{}" } }, + TagGroups = Array.Empty(), + Tags = Array.Empty(), + })); + + /// Seeds a Sealed deployment carrying a real, readable artifact that declares no drivers. + private static DeploymentId SeedDriverlessDeployment(IDbContextFactory db, RevisionHash rev) => + SeedDeployment(db, rev, JsonSerializer.SerializeToUtf8Bytes(new + { + RawFolders = Array.Empty(), + DriverInstances = Array.Empty(), + Devices = Array.Empty(), + TagGroups = Array.Empty(), + Tags = Array.Empty(), + })); + + private static DeploymentId SeedDeployment( + IDbContextFactory db, RevisionHash rev, byte[] artifact) + { + 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 = artifact, + }); + ctx.SaveChanges(); + return id; + } + + /// Factory producing one shared for the supported type, so a + /// REAL (non-stubbed) child spawns and its subscribe/unsubscribe traffic + /// is observable. + private sealed class SubscribingDriverFactory(string supportedType) : IDriverFactory + { + private readonly SubscribableStubDriver _driver = new(); + + /// The reference set passed to the driver's most recent SubscribeAsync call. + public IReadOnlyList? LastSubscribedRefs => _driver.LastSubscribedRefs; + + /// Number of UnsubscribeAsync calls — non-zero means a live handle was torn down. + public int UnsubscribeCount => Volatile.Read(ref _driver.UnsubscribeCount); + + /// + public IDriver? TryCreate(string driverType, string driverInstanceId, string driverConfigJson) => + string.Equals(driverType, supportedType, StringComparison.Ordinal) ? _driver : null; + + /// + public IReadOnlyCollection SupportedTypes => new[] { supportedType }; + } +} diff --git a/tests/Server/ZB.MOM.WW.OtOpcUa.Runtime.Tests/Drivers/StubDrivers.cs b/tests/Server/ZB.MOM.WW.OtOpcUa.Runtime.Tests/Drivers/StubDrivers.cs index 758b0069..9e83fb46 100644 --- a/tests/Server/ZB.MOM.WW.OtOpcUa.Runtime.Tests/Drivers/StubDrivers.cs +++ b/tests/Server/ZB.MOM.WW.OtOpcUa.Runtime.Tests/Drivers/StubDrivers.cs @@ -104,6 +104,10 @@ internal sealed class SubscribableStubDriver : StubDriver, ISubscribable /// The reference set passed to the most recent call. public IReadOnlyList? LastSubscribedRefs; + /// Number of times was called — the observable for a live + /// subscription being torn down (an EMPTY desired set drops the handle rather than re-subscribing). + public int UnsubscribeCount; + /// When true, genuinely yields (await Task.Yield()) /// before completing, so a ConfigureAwait(false) continuation in the actor resumes off the /// Akka ActorContext on a thread-pool thread — reproducing the no-ActorContext race that a @@ -127,6 +131,7 @@ internal sealed class SubscribableStubDriver : StubDriver, ISubscribable /// Cancellation token for the operation. public async Task UnsubscribeAsync(ISubscriptionHandle handle, CancellationToken cancellationToken) { + Interlocked.Increment(ref UnsubscribeCount); if (UnsubscribeYields) { // Complete the awaited task from a fresh background thread that has NO Akka actor