diff --git a/src/ZB.MOM.WW.ScadaBridge.SiteRuntime/Streaming/SiteStreamManager.cs b/src/ZB.MOM.WW.ScadaBridge.SiteRuntime/Streaming/SiteStreamManager.cs index 7d9a0da8..a49fedb4 100644 --- a/src/ZB.MOM.WW.ScadaBridge.SiteRuntime/Streaming/SiteStreamManager.cs +++ b/src/ZB.MOM.WW.ScadaBridge.SiteRuntime/Streaming/SiteStreamManager.cs @@ -20,6 +20,9 @@ namespace ZB.MOM.WW.ScadaBridge.SiteRuntime.Streaming; /// public class SiteStreamManager : ISiteStreamSubscriber { + /// Sentinel instance name recorded for site-wide (non-instance-scoped) subscriptions. + private const string SiteWideInstanceName = "*"; + private ActorSystem? _system; private IMaterializer? _materializer; private readonly int _bufferSize; @@ -117,6 +120,44 @@ public class SiteStreamManager : ISiteStreamSubscriber return subscriptionId; } + /// + /// Subscribe to ALARM events for ALL instances on the site (no per-instance + /// filter). Only events are forwarded; + /// events are dropped (attributes are far + /// higher-volume and the aggregated Alarm Summary never shows them). Same + /// broadcast-hub wiring as , and the returned + /// subscription id is torn down via exactly like the + /// per-instance variant. + /// + /// The actor that receives forwarded events. + /// A subscription id to pass to . + public string SubscribeSiteAlarms(IActorRef subscriber) + { + if (_hubSource is null || _materializer is null) + throw new InvalidOperationException("SiteStreamManager.Initialize must be called before SubscribeSiteAlarms"); + + var subscriptionId = Guid.NewGuid().ToString(); + var capturedSubscriber = subscriber; + + var killSwitch = _hubSource + .Where(ev => ev is AlarmStateChanged) + .Buffer(_bufferSize, OverflowStrategy.DropHead) + .ViaMaterialized(KillSwitches.Single(), Keep.Right) + .To(Sink.ForEach(ev => capturedSubscriber.Tell(ev))) + .Run(_materializer); + + lock (_lock) + { + _subscriptions[subscriptionId] = new SubscriptionInfo( + SiteWideInstanceName, subscriber, killSwitch, DateTimeOffset.UtcNow); + } + + _logger.LogDebug( + "Subscriber {SubscriptionId} registered for site-wide alarm events", subscriptionId); + + return subscriptionId; + } + /// /// Unsubscribe from instance events. Shuts down the per-subscriber /// stream graph via its KillSwitch. diff --git a/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Streaming/SiteStreamManagerTests.cs b/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Streaming/SiteStreamManagerTests.cs index 64316ffb..e4dbfbe7 100644 --- a/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Streaming/SiteStreamManagerTests.cs +++ b/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Streaming/SiteStreamManagerTests.cs @@ -103,6 +103,50 @@ public class SiteStreamManagerTests : TestKit, IDisposable probe2.ExpectNoMsg(TimeSpan.FromMilliseconds(500)); } + [Fact] + public void SubscribeSiteAlarms_ForwardsAlarmsForAllInstances_ButNotAttributes() + { + var probe = CreateTestProbe(); + _streamManager.SubscribeSiteAlarms(probe.Ref); + + // Alarm events from two different instances — both should arrive. + _streamManager.PublishAlarmStateChanged(new AlarmStateChanged( + "Pump1", "HighTemp", AlarmState.Active, 1, DateTimeOffset.UtcNow)); + _streamManager.PublishAlarmStateChanged(new AlarmStateChanged( + "Pump2", "LowFlow", AlarmState.Active, 2, DateTimeOffset.UtcNow)); + + // Attribute events must be filtered out entirely. + _streamManager.PublishAttributeValueChanged(new AttributeValueChanged( + "Pump1", "Temperature", "Temperature", "100", "Good", DateTimeOffset.UtcNow)); + _streamManager.PublishAttributeValueChanged(new AttributeValueChanged( + "Pump3", "Flow", "Flow", "42", "Good", DateTimeOffset.UtcNow)); + + // Collect the two alarms (order across instances is not guaranteed). + var received = new[] + { + probe.ExpectMsg(TimeSpan.FromSeconds(3)), + probe.ExpectMsg(TimeSpan.FromSeconds(3)), + }; + + var instances = received.Select(a => a.InstanceUniqueName).OrderBy(n => n).ToArray(); + Assert.Equal(new[] { "Pump1", "Pump2" }, instances); + + // No attribute events (nor any further alarm) should be delivered. + probe.ExpectNoMsg(TimeSpan.FromMilliseconds(500)); + } + + [Fact] + public void SubscribeSiteAlarms_SubscriptionRemovableViaUnsubscribe() + { + var probe = CreateTestProbe(); + var id = _streamManager.SubscribeSiteAlarms(probe.Ref); + + Assert.NotNull(id); + Assert.Equal(1, _streamManager.SubscriptionCount); + Assert.True(_streamManager.Unsubscribe(id)); + Assert.Equal(0, _streamManager.SubscriptionCount); + } + [Fact] public void RemoveSubscriber_RemovesAllSubscriptionsForActor() {