using System.Collections.Concurrent; using Grpc.Core; using Grpc.Net.Client; using Microsoft.Extensions.Logging; using ZB.MOM.WW.ScadaBridge.Commons.Messages.Streaming; using ZB.MOM.WW.ScadaBridge.Commons.Types.Alarms; using ZB.MOM.WW.ScadaBridge.Commons.Types.Enums; using Google.Protobuf.WellKnownTypes; namespace ZB.MOM.WW.ScadaBridge.Communication.Grpc; /// /// Per-site gRPC client that manages streaming subscriptions to a site's /// SiteStreamGrpcServer. The central-side DebugStreamBridgeActor uses this /// to open server-streaming calls for individual instances. /// public class SiteStreamGrpcClient : IAsyncDisposable, IDisposable { private readonly GrpcChannel? _channel; private readonly SiteStreamService.SiteStreamServiceClient? _client; private readonly ILogger? _logger; private readonly ConcurrentDictionary _subscriptions = new(); /// /// The gRPC endpoint (site node address) this client is bound to. The /// compares this against the requested /// endpoint so a NodeA→NodeB failover flip (or a site address edit) is honoured /// rather than served stale from cache. /// public virtual string Endpoint { get; } = string.Empty; /// /// The HTTP/2 keepalive ping delay actually applied to this client's channel. /// Exposed for tests verifying that is honoured. /// internal TimeSpan KeepAlivePingDelay { get; } /// /// The HTTP/2 keepalive ping timeout actually applied to this client's channel. /// Exposed for tests verifying that is honoured. /// internal TimeSpan KeepAlivePingTimeout { get; } /// /// Creates a client with default communication options. /// /// The gRPC endpoint address for the site. /// Logger for diagnostics and errors. public SiteStreamGrpcClient(string endpoint, ILogger logger) : this(endpoint, logger, new CommunicationOptions()) { } /// /// Creates a client whose HTTP/2 keepalive is taken from /// rather than hard-coded, satisfying the design doc's "gRPC Connection Keepalive" /// section which states these values are configurable. /// /// The gRPC endpoint address for the site. /// Logger for diagnostics and errors. /// Communication options including keepalive settings. public SiteStreamGrpcClient(string endpoint, ILogger logger, CommunicationOptions options) : this(endpoint, logger, options, pskProvider: null, siteIdentifier: null) { } /// /// Creates a client that authenticates every call with the site's preshared key. /// This is the production shape: SiteStreamService is gated by /// ControlPlaneAuthInterceptor on the site node, so a client without credentials /// gets on every call. /// /// The gRPC endpoint address for the site. /// Logger for diagnostics and errors. /// Communication options including keepalive settings. /// Resolves the site's preshared key; null leaves the channel unauthenticated. /// Site this channel talks to; null leaves the channel unauthenticated. public SiteStreamGrpcClient( string endpoint, ILogger logger, CommunicationOptions options, ISitePskProvider? pskProvider, string? siteIdentifier) { Endpoint = endpoint; KeepAlivePingDelay = options.GrpcKeepAlivePingDelay; KeepAlivePingTimeout = options.GrpcKeepAlivePingTimeout; _channel = GrpcChannel.ForAddress(endpoint, new GrpcChannelOptions { HttpHandler = new SocketsHttpHandler { KeepAlivePingDelay = options.GrpcKeepAlivePingDelay, KeepAlivePingTimeout = options.GrpcKeepAlivePingTimeout, KeepAlivePingPolicy = HttpKeepAlivePingPolicy.Always } }.WithSiteCredentials(pskProvider, siteIdentifier)); _client = new SiteStreamService.SiteStreamServiceClient(_channel); _logger = logger; } /// /// Protected constructor for unit testing without a real gRPC channel. /// Allows subclassing for mock implementations. /// protected SiteStreamGrpcClient() { } /// /// Protected constructor for unit testing — records the endpoint without /// opening a real gRPC channel, so endpoint-aware factory behaviour can be /// exercised by test doubles. /// /// The gRPC endpoint address for the site. protected SiteStreamGrpcClient(string endpoint) { Endpoint = endpoint; } /// /// Creates a test-only instance that has no gRPC channel. Used to test /// Unsubscribe and Dispose behavior without needing a real endpoint. /// /// A with no channel or client, for testing only. internal static SiteStreamGrpcClient CreateForTesting() => new(); /// /// Registers a CancellationTokenSource for a correlation ID. Test-only. /// /// Unique identifier for the subscription. /// CancellationTokenSource for managing the subscription lifecycle. internal void AddSubscriptionForTesting(string correlationId, CancellationTokenSource cts) { _subscriptions[correlationId] = cts; } /// /// Registers a subscription's CancellationTokenSource for a correlation ID. /// If an entry already exists for that correlation ID (a reconnect race where two /// calls briefly share an ID), the prior CTS is /// cancelled and disposed so it cannot leak. Internal for testability. /// /// Unique identifier for the subscription. /// CancellationTokenSource for managing the subscription lifecycle. internal void RegisterSubscription(string correlationId, CancellationTokenSource cts) { if (_subscriptions.TryGetValue(correlationId, out var prior) && !ReferenceEquals(prior, cts)) { prior.Cancel(); prior.Dispose(); } _subscriptions[correlationId] = cts; } /// /// Removes the subscription entry for a correlation ID only if the stored CTS is /// exactly the one supplied. A racing replacement stream may already own the slot, /// in which case this is a no-op. Internal for testability. /// /// Unique identifier for the subscription. /// CancellationTokenSource to match before removing. internal void RemoveSubscription(string correlationId, CancellationTokenSource cts) { _subscriptions.TryRemove(new KeyValuePair(correlationId, cts)); } /// /// Opens a server-streaming subscription for a specific instance. /// This is a long-running async method; the caller launches it as a background task. /// The callback delivers domain events, and /// lets the caller handle reconnection. /// /// Unique identifier for this subscription. /// Unique name of the instance to subscribe to. /// Callback invoked for each domain event received from the stream. /// Callback invoked when the subscription encounters an error. /// Cancellation token to stop the subscription. /// A task that represents the asynchronous operation. public virtual async Task SubscribeAsync( string correlationId, string instanceUniqueName, Action onEvent, Action onError, CancellationToken ct) { if (_client is null) throw new InvalidOperationException("Cannot subscribe on a test-only client."); var cts = CancellationTokenSource.CreateLinkedTokenSource(ct); RegisterSubscription(correlationId, cts); var request = new InstanceStreamRequest { CorrelationId = correlationId, InstanceUniqueName = instanceUniqueName }; try { using var call = _client.SubscribeInstance(request, cancellationToken: cts.Token); await foreach (var evt in call.ResponseStream.ReadAllAsync(cts.Token)) { var domainEvent = ConvertToDomainEvent(evt); if (domainEvent != null) onEvent(domainEvent); } } catch (RpcException ex) when (ex.StatusCode == StatusCode.Cancelled) { // Normal cancellation — not an error } catch (Exception ex) { onError(ex); } finally { // Remove only our own entry -- a racing reconnect may already own the slot. RemoveSubscription(correlationId, cts); } } /// /// Opens the site-wide, alarm-only server-streaming subscription /// (SubscribeSite) for a whole site rather than a single instance. This is /// the central per-site feed that backs the aggregated Alarm Summary live cache /// (plan #10): the site runtime's SubscribeSiteAlarms hub drops the /// per-instance filter and carries only events for /// all instances on the site (attributes are deliberately excluded — the /// summary never shows them and they are far higher-volume). /// /// The callback is deliberately typed as /// rather than the generic Action<object> used by /// : this stream is alarm-only by contract, so Task 4's /// per-site cache consumes an alarm delta directly with no downstream type test. /// The mapping is still shared via ; a non-alarm /// event (which should never appear on this stream) is defensively ignored rather /// than delivered or thrown. /// /// This is a long-running async method; the caller launches it as a background task. /// Error handling, cancellation and cleanup mirror . /// /// Unique identifier for this subscription. /// Callback invoked for each alarm delta received from the site-wide stream. /// Callback invoked when the subscription encounters an error. /// Cancellation token to stop the subscription. /// A task that represents the asynchronous operation. public virtual async Task SubscribeSiteAsync( string correlationId, Action onAlarmEvent, Action onError, CancellationToken ct) { if (_client is null) throw new InvalidOperationException("Cannot subscribe on a test-only client."); var cts = CancellationTokenSource.CreateLinkedTokenSource(ct); RegisterSubscription(correlationId, cts); var request = new SiteStreamRequest { CorrelationId = correlationId }; try { using var call = _client.SubscribeSite(request, cancellationToken: cts.Token); await foreach (var evt in call.ResponseStream.ReadAllAsync(cts.Token)) { // Site-wide stream is alarm-only by contract; defensively ignore anything else. if (ConvertToAlarmEvent(evt) is { } alarm) onAlarmEvent(alarm); } } catch (RpcException ex) when (ex.StatusCode == StatusCode.Cancelled) { // Normal cancellation — not an error } catch (Exception ex) { onError(ex); } finally { // Remove only our own entry -- a racing reconnect may already own the slot. RemoveSubscription(correlationId, cts); } } /// /// Cancels an active subscription by correlation ID. /// /// Unique identifier of the subscription to cancel. public virtual void Unsubscribe(string correlationId) { if (_subscriptions.TryRemove(correlationId, out var cts)) { cts.Cancel(); cts.Dispose(); } } /// /// Converts a proto SiteStreamEvent to the corresponding domain message. /// Internal for testability. /// /// The protobuf site stream event to convert. /// The converted domain event, or null if the event type is not recognized. internal static object? ConvertToDomainEvent(SiteStreamEvent evt) => evt.EventCase switch { SiteStreamEvent.EventOneofCase.AttributeChanged => new AttributeValueChanged( evt.AttributeChanged.InstanceUniqueName, evt.AttributeChanged.AttributePath, evt.AttributeChanged.AttributeName, evt.AttributeChanged.Value, MapQuality(evt.AttributeChanged.Quality), evt.AttributeChanged.Timestamp.ToDateTimeOffset()), SiteStreamEvent.EventOneofCase.AlarmChanged => new AlarmStateChanged( evt.AlarmChanged.InstanceUniqueName, evt.AlarmChanged.AlarmName, MapAlarmState(evt.AlarmChanged.State), evt.AlarmChanged.Priority, evt.AlarmChanged.Timestamp.ToDateTimeOffset()) { Level = MapAlarmLevel(evt.AlarmChanged.Level), Message = evt.AlarmChanged.Message ?? string.Empty, // Native alarm enrichment (additive — computed alarms carry their default condition). Kind = ParseAlarmKind(evt.AlarmChanged.Kind), Condition = new AlarmConditionState( Active: evt.AlarmChanged.Active, Acknowledged: evt.AlarmChanged.Acknowledged, Confirmed: evt.AlarmChanged.Confirmed, Shelve: AlarmShelveStateCodec.Parse(evt.AlarmChanged.ShelveState), Suppressed: evt.AlarmChanged.Suppressed, Severity: evt.AlarmChanged.Priority), SourceReference = evt.AlarmChanged.SourceReference ?? string.Empty, AlarmTypeName = evt.AlarmChanged.AlarmTypeName ?? string.Empty, Category = evt.AlarmChanged.Category ?? string.Empty, OperatorUser = evt.AlarmChanged.OperatorUser ?? string.Empty, OperatorComment = evt.AlarmChanged.OperatorComment ?? string.Empty, OriginalRaiseTime = evt.AlarmChanged.OriginalRaiseTime?.ToDateTimeOffset(), CurrentValue = evt.AlarmChanged.CurrentValue ?? string.Empty, LimitValue = evt.AlarmChanged.LimitValue ?? string.Empty, NativeSourceCanonicalName = evt.AlarmChanged.NativeSourceCanonicalName ?? string.Empty, IsConfiguredPlaceholder = evt.AlarmChanged.IsConfiguredPlaceholder }, _ => null }; /// /// Maps a proto to a domain , /// returning null for any non-alarm event. Used by /// to enforce the alarm-only contract of the site-wide stream without duplicating the /// enrichment mapping in . Internal for testability. /// /// The protobuf site stream event to convert. /// The mapped , or null if the event is not an alarm. internal static AlarmStateChanged? ConvertToAlarmEvent(SiteStreamEvent evt) => ConvertToDomainEvent(evt) as AlarmStateChanged; /// Parses the wire "kind" string back to ; defaults to Computed. /// The wire "kind" string from the gRPC payload; null or unrecognised defaults to . /// The parsed , or when the value is null or unrecognised. internal static AlarmKind ParseAlarmKind(string? kind) => System.Enum.TryParse(kind, ignoreCase: true, out var k) ? k : AlarmKind.Computed; /// /// Maps proto Quality enum to domain string. Internal for testability. /// /// The protobuf quality value to map. /// The mapped quality as a string ("Good", "Uncertain", "Bad", or "Unknown"). internal static string MapQuality(Quality quality) => quality switch { Quality.Good => "Good", Quality.Uncertain => "Uncertain", Quality.Bad => "Bad", _ => "Unknown" }; /// /// Maps proto AlarmStateEnum to domain AlarmState. Internal for testability. /// /// The protobuf alarm state to map. /// The mapped domain alarm state. internal static AlarmState MapAlarmState(AlarmStateEnum state) => state switch { AlarmStateEnum.AlarmStateNormal => AlarmState.Normal, AlarmStateEnum.AlarmStateActive => AlarmState.Active, _ => AlarmState.Normal }; /// /// Maps proto AlarmLevelEnum to domain AlarmLevel. Internal for testability. /// /// The protobuf alarm level to map. /// The mapped domain alarm level. internal static AlarmLevel MapAlarmLevel(AlarmLevelEnum level) => level switch { AlarmLevelEnum.AlarmLevelLow => AlarmLevel.Low, AlarmLevelEnum.AlarmLevelLowLow => AlarmLevel.LowLow, AlarmLevelEnum.AlarmLevelHigh => AlarmLevel.High, AlarmLevelEnum.AlarmLevelHighHigh => AlarmLevel.HighHigh, _ => AlarmLevel.None }; /// /// Releases all subscription CancellationTokenSources and the underlying /// gRPC channel. All teardown here is synchronous (CTS disposal and /// ), so a synchronous /// can release everything without sync-over-async blocking. /// private void ReleaseResources() { foreach (var cts in _subscriptions.Values) { cts.Cancel(); cts.Dispose(); } _subscriptions.Clear(); _channel?.Dispose(); } /// /// Asynchronously disposes of the gRPC client and all subscriptions. /// /// A completed after all subscriptions and the gRPC channel have been released. public virtual ValueTask DisposeAsync() { ReleaseResources(); return ValueTask.CompletedTask; } /// /// Synchronous disposal. All resources held by this client are released /// synchronously, so callers (e.g. ) /// need not block on the async disposal path. /// public virtual void Dispose() { ReleaseResources(); } }