feat(sitestream): add SubscribeSiteAlarms site-wide alarm-only broadcast (plan #10 T1)

Adds SiteStreamManager.SubscribeSiteAlarms(IActorRef) — same broadcast-hub
wiring as Subscribe(...) but with no per-instance filter and forwarding only
AlarmStateChanged events (AttributeValueChanged dropped). Returns a subscription
id torn down via the existing Unsubscribe, symmetric with the per-instance
subscribe. Backs the Task 2 SubscribeSite gRPC server handler.

Claude-Session: https://claude.ai/code/session_01MtdgwpEeCUn6cUA5f1LMPj
This commit is contained in:
Joseph Doherty
2026-07-10 11:39:42 -04:00
parent 84112db344
commit b14a4dbc6b
2 changed files with 85 additions and 0 deletions
@@ -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<AlarmStateChanged>(TimeSpan.FromSeconds(3)),
probe.ExpectMsg<AlarmStateChanged>(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()
{