From e75cd6f3d87f3cc25b41f21d3fe1a01561aeef14 Mon Sep 17 00:00:00 2001 From: Joseph Doherty Date: Wed, 8 Jul 2026 17:43:40 -0400 Subject: [PATCH] fix(communication): debug-stream teardown/failover unsubscribes via TryGet, never creates or disposes shared channels --- .../Actors/DebugStreamBridgeActor.cs | 14 ++-- .../Grpc/DebugStreamBridgeActorTests.cs | 67 +++++++++++++++++++ 2 files changed, 76 insertions(+), 5 deletions(-) diff --git a/src/ZB.MOM.WW.ScadaBridge.Communication/Actors/DebugStreamBridgeActor.cs b/src/ZB.MOM.WW.ScadaBridge.Communication/Actors/DebugStreamBridgeActor.cs index 84b27635..4b4170c7 100644 --- a/src/ZB.MOM.WW.ScadaBridge.Communication/Actors/DebugStreamBridgeActor.cs +++ b/src/ZB.MOM.WW.ScadaBridge.Communication/Actors/DebugStreamBridgeActor.cs @@ -466,8 +466,11 @@ public class DebugStreamBridgeActor : ReceiveActor, IWithTimers // stops the StreamRelayActor for this correlation ID, rather than leaving a // zombie relay actor until TCP RST / keepalive eventually detects the loss. var previousEndpoint = _useNodeA ? _grpcNodeAAddress : _grpcNodeBAddress; - var previousClient = _grpcFactory.GetOrCreate(_siteIdentifier, previousEndpoint); - previousClient.Unsubscribe(_correlationId); + // TryGet, not GetOrCreate: unsubscribing a failed stream must never open a + // fresh channel (and, with (site,endpoint) keying, must never touch another + // session's healthy channel). Absent client => the channel is already gone + // and the site-side relay will be reaped by keepalive/session-lifetime. + _grpcFactory.TryGet(_siteIdentifier, previousEndpoint)?.Unsubscribe(_correlationId); // Flip to the other node _useNodeA = !_useNodeA; @@ -489,9 +492,10 @@ public class DebugStreamBridgeActor : ReceiveActor, IWithTimers _grpcCts?.Dispose(); _grpcCts = null; - var client = _grpcFactory.GetOrCreate(_siteIdentifier, - _useNodeA ? _grpcNodeAAddress : _grpcNodeBAddress); - client.Unsubscribe(_correlationId); + // TryGet, not GetOrCreate: teardown must never open a fresh channel just to + // unsubscribe. Absent client => nothing to cancel. + var endpoint = _useNodeA ? _grpcNodeAAddress : _grpcNodeBAddress; + _grpcFactory.TryGet(_siteIdentifier, endpoint)?.Unsubscribe(_correlationId); } private void SendUnsubscribe() diff --git a/tests/ZB.MOM.WW.ScadaBridge.Communication.Tests/Grpc/DebugStreamBridgeActorTests.cs b/tests/ZB.MOM.WW.ScadaBridge.Communication.Tests/Grpc/DebugStreamBridgeActorTests.cs index 6ea2faff..6068f628 100644 --- a/tests/ZB.MOM.WW.ScadaBridge.Communication.Tests/Grpc/DebugStreamBridgeActorTests.cs +++ b/tests/ZB.MOM.WW.ScadaBridge.Communication.Tests/Grpc/DebugStreamBridgeActorTests.cs @@ -400,6 +400,60 @@ public class DebugStreamBridgeActorTests : TestKit Assert.Equal("corr-1", factory.ClientFor(GrpcNodeB).SubscribeCalls[0].CorrelationId); } + // ── Task 6 (arch review 02, High): teardown/failover unsubscribe is endpoint-safe ── + // Both paths use TryGet, never GetOrCreate, so cleanup can never open a fresh + // channel or (with (site,endpoint) keying) touch another session's channel. + + private (IActorRef Actor, EndpointTrackingGrpcClientFactory Factory) CreateBridgeWithTrackingFactory() + { + var commProbe = CreateTestProbe(); + var factory = new EndpointTrackingGrpcClientFactory(); + var props = Props.Create(typeof(DebugStreamBridgeActor), + SiteId, InstanceName, "corr-1", commProbe.Ref, + (Action)(_ => { }), (Action)(() => { }), + factory, GrpcNodeA, GrpcNodeB); + var actor = Sys.ActorOf(props); + commProbe.ExpectMsg(); // PreStart subscribe handshake + return (actor, factory); + } + + [Fact] + public void CleanupGrpc_UnsubscribesViaTryGet_WithoutOpeningAChannel() + { + var (actor, factory) = CreateBridgeWithTrackingFactory(); + AwaitCondition(() => factory.ClientFor(GrpcNodeA).SubscribeCalls.Count == 1, + TimeSpan.FromSeconds(3)); + var createdBefore = factory.CreatedCount; // node-a only + + actor.Tell(new StopDebugStream()); + + AwaitCondition(() => factory.ClientFor(GrpcNodeA).UnsubscribedCorrelationIds.Contains("corr-1"), + TimeSpan.FromSeconds(3)); + // Pre-fix CleanupGrpc called GetOrCreate → could open a fresh channel here. + Assert.Equal(createdBefore, factory.CreatedCount); + } + + [Fact] + public void SessionFailover_UnsubscribesFailedEndpointViaTryGet_OpensOnlyTheReconnectChannel() + { + var (_, factory) = CreateBridgeWithTrackingFactory(); + AwaitCondition(() => factory.ClientFor(GrpcNodeA).SubscribeCalls.Count == 1, + TimeSpan.FromSeconds(3)); + + // Fail NodeA: the bridge unsubscribes the failed NodeA stream via TryGet + // (if that path had used GetOrCreate on the post-flip cache it could have + // disposed/created the wrong client), then reconnects to NodeB. + factory.ClientFor(GrpcNodeA).SubscribeCalls[0].OnError(new Exception("NodeA down")); + + AwaitCondition(() => factory.ClientFor(GrpcNodeA).UnsubscribedCorrelationIds.Contains("corr-1"), + TimeSpan.FromSeconds(5)); + AwaitCondition(() => factory.ClientFor(GrpcNodeB).SubscribeCalls.Count == 1, + TimeSpan.FromSeconds(5)); + // Only the two real node endpoints were ever opened — the unsubscribe + // itself created nothing. + Assert.Equal(2, factory.CreatedCount); + } + // --------------------------------------------------------------------- // M2.18 (#26) — stream-first + replay/dedup // --------------------------------------------------------------------- @@ -883,6 +937,12 @@ internal class MockSiteStreamGrpcClientFactory : SiteStreamGrpcClientFactory RequestedEndpoints.Add(grpcEndpoint); return _mockClient; } + + // Mirrors the real factory's TryGet: returns the client only for an endpoint + // already opened via GetOrCreate, never creating one. Lets teardown/failover + // unsubscribe paths (now TryGet-based, Task 6) resolve the mock client. + public override SiteStreamGrpcClient? TryGet(string siteIdentifier, string grpcEndpoint) + => RequestedEndpoints.Contains(grpcEndpoint) ? _mockClient : null; } /// @@ -902,6 +962,13 @@ internal class EndpointTrackingGrpcClientFactory : SiteStreamGrpcClientFactory public MockSiteStreamGrpcClient ClientFor(string endpoint) => _byEndpoint.GetOrAdd(endpoint, _ => new MockSiteStreamGrpcClient()); + /// Number of distinct endpoint clients created so far (channels opened). + public int CreatedCount => _byEndpoint.Count; + public override SiteStreamGrpcClient GetOrCreate(string siteIdentifier, string grpcEndpoint) => ClientFor(grpcEndpoint); + + // TryGet never creates: returns the client only for an already-opened endpoint. + public override SiteStreamGrpcClient? TryGet(string siteIdentifier, string grpcEndpoint) + => _byEndpoint.TryGetValue(grpcEndpoint, out var client) ? client : null; }