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; /// /// 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. /// public class StreamRelayActor : ReceiveActor { private readonly ILoggingAdapter _log = Context.GetLogger(); private readonly string _correlationId; private readonly ChannelWriter _channelWriter; /// /// Initializes a new for the given gRPC stream correlation. /// /// Correlation id stamped on every relayed . /// Channel writer to which converted events are written. public StreamRelayActor(string correlationId, ChannelWriter channelWriter) { _correlationId = correlationId; _channelWriter = channelWriter; Receive(HandleAttributeValueChanged); Receive(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 = ValueFormatter.FormatDisplayValue(msg.Value), Quality = MapQuality(msg.Quality), Timestamp = Timestamp.FromDateTimeOffset(msg.Timestamp) } }; WriteToChannel(protoEvent); } private void HandleAlarmStateChanged(AlarmStateChanged msg) { 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 } }; 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 }; }