using NetMQ; using System; using System.Text; using NetMQ.Sockets; using System.Threading; using xCommons.Extensions; using xEventService.Models; using xExceptions.Constants; using System.Threading.Tasks; using xEventService.Constants; using xEventService.Extensions; using xEventService.Interfaces; using xEventService.Configuration; using Microsoft.Extensions.Logging; namespace xEventService.Providers { public class XZeroMQProducerService : XBaseEventProducerService, IXZeroMQProducerService { // private NetMQSocket socket; private readonly object socketLock = new(); public XZeroMQProducerService( XEventServiceConfiguration configuration, ILogger logger ) : base(configuration, logger) { // // Validate Configurations ... if (!configuration.IsValid() || !configuration.ZeroMQ.IsValid() || configuration.Broker != XEventBroker.XZeroMQ ) { XException.InvalidConfiguration.Throw(); } // InitializeSocket(); } /// /// Sending Message with Full Request Descriptor ... /// /// /// /// public async Task SendMessageAsync( XEventRequest request, CancellationToken cancellationToken = default ) { // var result = request.IsValid(); if (!result) { XException.InvalidArgs.Throw(); } // result = await SendMessageAsync( message: request.Message, topic: request.Topic, cancellationToken: cancellationToken ); // return result; } /// /// Sending a Message to Topic ... /// /// /// /// /// public async Task SendMessageAsync( XEventMessage message, string topic = XEventServiceConstants.XDefaultTopic, CancellationToken cancellationToken = default ) { // var result = message.IsValid(); if (!result) { XException.InvalidArgs.Throw(); } // // Normalize Topic ... if (topic.IsNullOrEmpty()) { topic = XEventServiceConstants.XDefaultTopic; } // try { // var json = message.ToJSON(camelCase: true); var body = Encoding.UTF8.GetBytes(json); await Task.Run(() => { // lock (socketLock) { // if (configuration.ZeroMQ.Pattern == XZeroMQPattern.PubSub) { socket.SendMoreFrame(topic).SendFrame(body); } else { socket.SendFrame(body); } } }, cancellationToken); // result = true; } catch (Exception) { result = false; } // return result; } /// /// Dispose Object ... /// public void Dispose() { // try { // lock (socketLock) { // if (socket != null) { // socket.Close(); socket.Dispose(); socket = null; } } } catch { } // GC.SuppressFinalize(this); } // #region Private ... /// /// Initialize Socket ... /// private void InitializeSocket() { // lock (socketLock) { // if (configuration.ZeroMQ.Pattern == XZeroMQPattern.PubSub) { // socket = new PublisherSocket(); socket.Options.SendHighWatermark = configuration.ZeroMQ.SendHighWaterMark; } else { // socket = new PushSocket(); socket.Options.SendHighWatermark = configuration.ZeroMQ.SendHighWaterMark; } // if (configuration.ZeroMQ.IsServer) { socket.Bind(configuration.Url); } else { socket.Connect(configuration.Url); } } } #endregion } }