Initial Commit ...
This commit is contained in:
@@ -0,0 +1,260 @@
|
||||
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
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user