260 lines
8.1 KiB
C#
260 lines
8.1 KiB
C#
using System;
|
|
using System.Collections.Generic;
|
|
using System.Linq;
|
|
using System.Text;
|
|
using System.Threading.Tasks;
|
|
using Microsoft.Extensions.DependencyInjection;
|
|
using Microsoft.Extensions.Logging;
|
|
using RabbitMQ.Client.Events;
|
|
using xCommons.Extensions;
|
|
using xEventService.Configuration;
|
|
using xEventService.Extensions;
|
|
using xEventService.Interfaces;
|
|
using xEventService.Models;
|
|
using xExceptions.Constants;
|
|
|
|
namespace xEventService.Providers {
|
|
public sealed class XEventBus : IXEventBus {
|
|
//
|
|
#region Properties ...
|
|
private readonly List<Type> events;
|
|
private readonly ILogger<XEventBus> logger;
|
|
private readonly IServiceScopeFactory serviceScopeFactory;
|
|
private readonly XEventServiceConfiguration configuration;
|
|
private readonly IDictionary<string, List<Type>> handlers;
|
|
#endregion
|
|
|
|
//
|
|
#region Constructors ...
|
|
public XEventBus (
|
|
ILoggerFactory loggerFactory,
|
|
IServiceScopeFactory serviceScopeFactory,
|
|
XEventServiceConfiguration configuration
|
|
) {
|
|
//
|
|
this.configuration = configuration;
|
|
this.serviceScopeFactory = serviceScopeFactory;
|
|
|
|
//
|
|
events = new List<Type> ();
|
|
handlers = new Dictionary<string, List<Type>> ();
|
|
logger = loggerFactory.CreateLogger<XEventBus> ();
|
|
}
|
|
#endregion
|
|
|
|
//
|
|
#region Abstracts ...
|
|
public bool Publish<T> (T @event) where T : XEvent {
|
|
//
|
|
// Validate Args ...
|
|
if (@event.IsNull ()) {
|
|
return false;
|
|
}
|
|
|
|
//
|
|
// Prepare Channel ...
|
|
XChannel channel = null;
|
|
|
|
//
|
|
// Extract event name ...
|
|
var exchangeName = @event.GetType ().Name;
|
|
|
|
//
|
|
// Do Action ...
|
|
try {
|
|
//
|
|
// Get XChannel Model ...
|
|
channel = configuration.CreateChannel (exchangeName);
|
|
|
|
//
|
|
// Publish Command ...
|
|
channel.Publish (message: @event);
|
|
|
|
//
|
|
// Return Result ...
|
|
return true;
|
|
} catch (Exception ex) {
|
|
//
|
|
logger.LogError ($"XEventService Exception: Publishing Failed {exchangeName}/{ex.Message} ...");
|
|
|
|
//
|
|
return false;
|
|
} finally {
|
|
//
|
|
if (!channel.IsNull ()) {
|
|
channel.Dispose ();
|
|
}
|
|
}
|
|
}
|
|
|
|
public bool Subscribe<TEvent, THandler> ()
|
|
where TEvent : XEvent
|
|
where THandler : IXEventHandler<TEvent> {
|
|
//
|
|
// Extract required data ...
|
|
var eventType = typeof (TEvent);
|
|
var eventName = eventType.Name;
|
|
var handlerType = typeof (THandler);
|
|
|
|
//
|
|
// Check event exists in list or not ...
|
|
if (!events.Contains (eventType)) {
|
|
//
|
|
// Add it if not exists ...
|
|
events.Add (eventType);
|
|
}
|
|
|
|
//
|
|
// Check handlers Dictionary contains list or not ...
|
|
if (!handlers.ContainsKey (eventName)) {
|
|
handlers.Add (eventName, new List<Type> ());
|
|
}
|
|
|
|
//
|
|
// Check handlers Subscribed before or not ...
|
|
var isHandlerSubscribed = handlers[eventName]
|
|
.Any (h => h.GetType () == handlerType);
|
|
if (isHandlerSubscribed) {
|
|
//
|
|
logger.LogError ($"XEventService Exception: Duplicate Event Handler Subscription {eventName}/{handlerType.Name} ...");
|
|
|
|
//
|
|
return false;
|
|
}
|
|
|
|
//
|
|
// Register Handler ...
|
|
handlers[eventName].Add (handlerType);
|
|
|
|
//
|
|
// Do Consuming Handler ...
|
|
ConsumeEvent<TEvent> ();
|
|
|
|
//
|
|
return true;
|
|
}
|
|
#endregion
|
|
|
|
//
|
|
#region Private ...
|
|
/// <summary>
|
|
/// Consume Specific Event ...
|
|
/// </summary>
|
|
/// <typeparam name="TEvent"></typeparam>
|
|
private void ConsumeEvent<TEvent> ()
|
|
where TEvent : XEvent {
|
|
//
|
|
// Extract required data ...
|
|
var eventType = typeof (TEvent);
|
|
var eventName = eventType.Name;
|
|
|
|
//
|
|
// Create ChannelObject ...
|
|
var channel = new XChannel (
|
|
exchangeName: eventName,
|
|
configuration: configuration,
|
|
dispatchConsumersAsync: true
|
|
);
|
|
|
|
//
|
|
// Create Consumer ...
|
|
var consumer = new AsyncEventingBasicConsumer (channel.Channel);
|
|
|
|
//
|
|
// Set Consumer Delegate ...
|
|
consumer.Received += XConsumerReceivedDelegate;
|
|
|
|
//
|
|
// Consume Async Consumer ...
|
|
channel.ConsumeAsync (consumer);
|
|
}
|
|
|
|
/// <summary>
|
|
/// this Delegate method calls when a message recieved ...
|
|
/// </summary>
|
|
private async Task XConsumerReceivedDelegate (
|
|
object sender,
|
|
BasicDeliverEventArgs ea
|
|
) {
|
|
//
|
|
// Generate Required Data ...
|
|
var eventName = ea.Exchange;
|
|
var json = Encoding.UTF8
|
|
.GetString (
|
|
ea.Body
|
|
.ToArray ()
|
|
);
|
|
|
|
//
|
|
// Do Processing Event ...
|
|
try {
|
|
//
|
|
// Process Event in non Blocking Tasks ...
|
|
await ProcessEvent (eventName, json)
|
|
.ConfigureAwait (false);
|
|
} catch (Exception ex) {
|
|
//
|
|
logger.LogError ($"XEventService Exception: Processing Event failed {eventName}/{ex.Message} ...");
|
|
|
|
//
|
|
XException.ActionFailed.Throw ();
|
|
} finally {
|
|
//
|
|
var channel = ((EventingBasicConsumer) sender).Model;
|
|
channel.BasicAck (
|
|
deliveryTag: ea.DeliveryTag,
|
|
multiple: false
|
|
);
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Here we must call registered Event Handlers to Handle the Event ...
|
|
/// </summary>
|
|
private async Task ProcessEvent (
|
|
string eventName,
|
|
string json
|
|
) {
|
|
//
|
|
// Check Handler Exists for current Event ...
|
|
var isExistsHandler = handlers.ContainsKey (eventName);
|
|
if (!isExistsHandler) {
|
|
return;
|
|
}
|
|
|
|
//
|
|
// Using Service Scope Factory ...
|
|
using (var scope = serviceScopeFactory.CreateScope ()) {
|
|
//
|
|
// Get Handler Subscriptions ...
|
|
var subscriptions = handlers[eventName];
|
|
foreach (var subscription in subscriptions) {
|
|
//
|
|
// Inject Handler from Dependency Injections ...
|
|
var handler = scope.ServiceProvider.GetService (subscription);
|
|
if (handler.IsNull ()) {
|
|
//
|
|
logger.LogInformation ($"XEventService: there is no handler Registered in DependencyInjection for {eventName}/{subscription.Name}");
|
|
|
|
//
|
|
continue;
|
|
}
|
|
|
|
//
|
|
// Retrieve Event Type ...
|
|
var eventType = events.SingleOrDefault (t => t.Name == eventName);
|
|
var @event = json.FromJSON (eventType);
|
|
var conreteType = typeof (IXEventHandler<>).MakeGenericType (eventType);
|
|
|
|
//
|
|
// Call EventHandler 'HandleAsync' method ...
|
|
await ((Task) conreteType
|
|
.GetMethod (nameof (IXEventHandler<XEvent>.HandleAsync))
|
|
.Invoke (handler, new object[] { @event }))
|
|
.ConfigureAwait (true);
|
|
}
|
|
}
|
|
}
|
|
#endregion
|
|
}
|
|
} |