using System; using System.Text; using System.Threading; using System.Threading.Tasks; using Microsoft.Extensions.Logging; using NetMQ; using NetMQ.Sockets; using xCommons.Extensions; using xEventService.Configuration; using xEventService.Constants; using xEventService.Extensions; using xEventService.Models; using xExceptions.Constants; namespace xEventService.Providers { /// /// an Abstract ZeroMQ Consumer Service ... /// public abstract class XZeroMQConsumerServiceBase : XBaseEventConsumerBackgroundService { // private NetMQSocket socket; protected XZeroMQConsumerServiceBase( string[] allowedActions, string[] allowedSenders, IServiceProvider serviceProvider, XEventServiceConfiguration configuration, ILogger logger, string topic = XEventServiceConstants.XDefaultTopic ) : base( topic: topic, logger: logger, configuration: configuration, allowedActions: allowedActions, allowedSenders: allowedSenders, serviceProvider: serviceProvider ) { // // Validating Configuration ... if (configuration.Broker != XEventBroker.XKafka) { XException.InvalidConfiguration.Throw(); } // // Initialize ... InitializeSocket(); } /// /// Default Action Execution for Background Services ... /// /// /// protected override async Task ExecuteAsync( CancellationToken cancellationToken = default ) { // await Task.Yield(); while (!cancellationToken.IsCancellationRequested) { // try { // string receivedTopic = string.Empty; byte[] body = null; if (configuration.ZeroMQ.Pattern == XZeroMQPattern.PubSub) { // if (socket.TryReceiveFrameString(TimeSpan.FromMilliseconds(100), out receivedTopic)) { // if (socket.TryReceiveFrameBytes(TimeSpan.FromMilliseconds(100), out body)) { // Process } } } else { if (socket.TryReceiveFrameBytes(TimeSpan.FromMilliseconds(100), out body)) { // Process } } if (body != null && body.Length > 0) { // var json = Encoding.UTF8.GetString(body); var message = json.FromJSON(); if (!message.IsNullOrDefault() && message.IsValid( allowedSenders: allowedSenders, allowedActions: allowedActions) ) { // await DoWorkAsync( payload: message, cancellationToken: cancellationToken ); } } else { await Task.Delay(10, cancellationToken); } } catch (Exception) { await Task.Delay(1000, cancellationToken); } } } public override async Task StopAsync(CancellationToken cancellationToken) { // socket.Close(); await base.StopAsync(cancellationToken); } /// /// Dispose Object ... /// public override void Dispose() { // if (socket != null) { socket.Dispose(); } // GC.SuppressFinalize(this); base.Dispose(); } // #region Private ... private void InitializeSocket() { // if (configuration.ZeroMQ.Pattern == XZeroMQPattern.PubSub) { // socket = new SubscriberSocket(); socket.Options.ReceiveHighWatermark = configuration.ZeroMQ.ReceiveHighWaterMark; // // Subscribe to specific topic or all if (topic.IsNullOrEmpty() || topic == XEventServiceConstants.XWildcard) { ((SubscriberSocket)socket).SubscribeToAnyTopic(); } else { ((SubscriberSocket)socket).Subscribe(topic); } } else { // socket = new PullSocket(); socket.Options.ReceiveHighWatermark = configuration.ZeroMQ.ReceiveHighWaterMark; } // if (configuration.ZeroMQ.IsServer) { socket.Bind(configuration.Url); } else { socket.Connect(configuration.Url); } } #endregion } }