Files
ScadaBridge/src/ZB.MOM.WW.ScadaBridge.Communication/Actors/StreamRelayActor.cs
T
Joseph Doherty e7660134f2 fix(communication): drop IsConfiguredPlaceholder rows in StreamRelayActor before gRPC pack
Placeholder AlarmStateChanged rows are a DebugView snapshot-only concept emitted
by InstanceActor.BuildAlarmStatesSnapshot; they are never a real alarm transition.
Their timestamp may be DateTimeOffset.MinValue (the Protobuf Timestamp lower boundary),
which can throw when packed via Timestamp.FromDateTimeOffset.

Added early-return guard at the top of HandleAlarmStateChanged before any timestamp
pack or channel write. Updated the existing NativeBindingLinkage round-trip test to
use a real (non-placeholder) native alarm; added DropsAlarmStateChanged_WhenIsConfiguredPlaceholder
to assert placeholders are silently dropped (15/15 pass).
2026-06-17 15:44:28 -04:00

137 lines
5.6 KiB
C#

using System.Threading.Channels;
using Akka.Actor;
using Akka.Event;
using Google.Protobuf.WellKnownTypes;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Streaming;
using ZB.MOM.WW.ScadaBridge.Commons.Types;
using ZB.MOM.WW.ScadaBridge.Communication.Grpc;
using AlarmState = ZB.MOM.WW.ScadaBridge.Commons.Types.Enums.AlarmState;
using AlarmLevel = ZB.MOM.WW.ScadaBridge.Commons.Types.Enums.AlarmLevel;
namespace ZB.MOM.WW.ScadaBridge.Communication.Actors;
/// <summary>
/// Lightweight relay actor that bridges Akka domain events (AttributeValueChanged,
/// AlarmStateChanged) to a System.Threading.Channels.Channel of protobuf SiteStreamEvent
/// messages. The gRPC server method reads from the channel's reader side.
/// </summary>
public class StreamRelayActor : ReceiveActor
{
private readonly ILoggingAdapter _log = Context.GetLogger();
private readonly string _correlationId;
private readonly ChannelWriter<SiteStreamEvent> _channelWriter;
/// <summary>
/// Initializes a new <see cref="StreamRelayActor"/> for the given gRPC stream correlation.
/// </summary>
/// <param name="correlationId">Correlation id stamped on every relayed <see cref="SiteStreamEvent"/>.</param>
/// <param name="channelWriter">Channel writer to which converted events are written.</param>
public StreamRelayActor(string correlationId, ChannelWriter<SiteStreamEvent> channelWriter)
{
_correlationId = correlationId;
_channelWriter = channelWriter;
Receive<AttributeValueChanged>(HandleAttributeValueChanged);
Receive<AlarmStateChanged>(HandleAlarmStateChanged);
}
private void HandleAttributeValueChanged(AttributeValueChanged msg)
{
var protoEvent = new SiteStreamEvent
{
CorrelationId = _correlationId,
AttributeChanged = new AttributeValueUpdate
{
InstanceUniqueName = msg.InstanceUniqueName,
AttributePath = msg.AttributePath,
AttributeName = msg.AttributeName,
Value = AttributeValueCodec.Encode(msg.Value) ?? string.Empty,
Quality = MapQuality(msg.Quality),
Timestamp = Timestamp.FromDateTimeOffset(msg.Timestamp)
}
};
WriteToChannel(protoEvent);
}
private void HandleAlarmStateChanged(AlarmStateChanged msg)
{
// Placeholder rows (IsConfiguredPlaceholder) are a Debug View snapshot-only
// concept emitted by InstanceActor.BuildAlarmStatesSnapshot — they are never a
// real alarm transition and must not be relayed to the live gRPC stream (their
// timestamp may be DateTimeOffset.MinValue, the Protobuf Timestamp boundary).
if (msg.IsConfiguredPlaceholder)
{
return;
}
var protoEvent = new SiteStreamEvent
{
CorrelationId = _correlationId,
AlarmChanged = new AlarmStateUpdate
{
InstanceUniqueName = msg.InstanceUniqueName,
AlarmName = msg.AlarmName,
State = MapAlarmState(msg.State),
Priority = msg.Priority,
Timestamp = Timestamp.FromDateTimeOffset(msg.Timestamp),
Level = MapAlarmLevel(msg.Level),
Message = msg.Message ?? string.Empty,
// Native alarm enrichment (additive — computed alarms map their default condition).
Kind = msg.Kind.ToString(),
Active = msg.Condition.Active,
Acknowledged = msg.Condition.Acknowledged,
Confirmed = msg.Condition.Confirmed ?? false,
ShelveState = AlarmShelveStateCodec.ToWire(msg.Condition.Shelve),
Suppressed = msg.Condition.Suppressed,
SourceReference = msg.SourceReference ?? string.Empty,
AlarmTypeName = msg.AlarmTypeName ?? string.Empty,
Category = msg.Category ?? string.Empty,
OperatorUser = msg.OperatorUser ?? string.Empty,
OperatorComment = msg.OperatorComment ?? string.Empty,
OriginalRaiseTime = msg.OriginalRaiseTime.HasValue
? Timestamp.FromDateTimeOffset(msg.OriginalRaiseTime.Value)
: null,
CurrentValue = msg.CurrentValue ?? string.Empty,
LimitValue = msg.LimitValue ?? string.Empty,
NativeSourceCanonicalName = msg.NativeSourceCanonicalName ?? string.Empty,
IsConfiguredPlaceholder = msg.IsConfiguredPlaceholder
}
};
WriteToChannel(protoEvent);
}
private void WriteToChannel(SiteStreamEvent protoEvent)
{
if (!_channelWriter.TryWrite(protoEvent))
{
_log.Warning("Channel full, dropping event for correlation {0}", _correlationId);
}
}
private static Quality MapQuality(string quality) => quality switch
{
"Good" => Quality.Good,
"Uncertain" => Quality.Uncertain,
"Bad" => Quality.Bad,
_ => Quality.Unspecified
};
private static AlarmStateEnum MapAlarmState(AlarmState state) => state switch
{
AlarmState.Normal => AlarmStateEnum.AlarmStateNormal,
AlarmState.Active => AlarmStateEnum.AlarmStateActive,
_ => AlarmStateEnum.AlarmStateUnspecified
};
private static AlarmLevelEnum MapAlarmLevel(AlarmLevel level) => level switch
{
AlarmLevel.Low => AlarmLevelEnum.AlarmLevelLow,
AlarmLevel.LowLow => AlarmLevelEnum.AlarmLevelLowLow,
AlarmLevel.High => AlarmLevelEnum.AlarmLevelHigh,
AlarmLevel.HighHigh => AlarmLevelEnum.AlarmLevelHighHigh,
_ => AlarmLevelEnum.AlarmLevelNone
};
}