Skip to content

Commit

Permalink
Merge pull request #155 from Resgrid/develop
Browse files Browse the repository at this point in the history
CU-868a3neq0 trying to fix eventing issue.
  • Loading branch information
ucswift authored Oct 8, 2024
2 parents 9c35c06 + 0aa2ae1 commit 4c21f2c
Show file tree
Hide file tree
Showing 24 changed files with 480 additions and 290 deletions.
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ namespace Resgrid.Model.Providers
{
public interface IRabbitInboundEventProvider
{
Task Start();
Task Start(string clientName, string queueName);
void RegisterForEvents(Func<int, string, Task> personnelStatusChanged,
Func<int, string, Task> unitStatusChanged,
Func<int, string, Task> callStatusChanged,
Expand Down
24 changes: 11 additions & 13 deletions Providers/Resgrid.Providers.Bus.Rabbit/RabbitConnection.cs
Original file line number Diff line number Diff line change
Expand Up @@ -9,27 +9,26 @@ internal class RabbitConnection
{
private static IConnection _connection { get; set; }
private static ConnectionFactory _factory { get; set; }
private static object LOCK = new object();
private readonly static object LOCK = new object();


public static bool VerifyAndCreateClients()
public static bool VerifyAndCreateClients(string clientName)
{
if (_connection != null && !_connection.IsOpen)
{
_connection?.Dispose();

_connection.Dispose();
_connection = null;
_factory = null;
}

if (_connection == null)
{
lock (LOCK)
{
try
{
_factory = new ConnectionFactory() { HostName = ServiceBusConfig.RabbitHostname, UserName = ServiceBusConfig.RabbitUsername, Password = ServiceBusConfig.RabbbitPassword };
_connection = _factory.CreateConnection();
_connection = _factory.CreateConnection(clientName);
}
catch (Exception ex)
{
Expand All @@ -40,7 +39,7 @@ public static bool VerifyAndCreateClients()
try
{
_factory = new ConnectionFactory() { HostName = ServiceBusConfig.RabbitHostname2, UserName = ServiceBusConfig.RabbitUsername, Password = ServiceBusConfig.RabbbitPassword };
_connection = _factory.CreateConnection();
_connection = _factory.CreateConnection(clientName);
}
catch (Exception ex2)
{
Expand All @@ -51,7 +50,7 @@ public static bool VerifyAndCreateClients()
try
{
_factory = new ConnectionFactory() { HostName = ServiceBusConfig.RabbitHostname3, UserName = ServiceBusConfig.RabbitUsername, Password = ServiceBusConfig.RabbbitPassword };
_connection = _factory.CreateConnection();
_connection = _factory.CreateConnection(clientName);
}
catch (Exception ex3)
{
Expand Down Expand Up @@ -155,19 +154,18 @@ public static bool VerifyAndCreateClients()
return false;
}

public static IConnection CreateConnection()
public static IConnection CreateConnection(string clientName)
{
if (_connection == null)
VerifyAndCreateClients();
VerifyAndCreateClients(clientName);

if (!_connection.IsOpen)
{
_connection?.Dispose();

_connection.Dispose();
_connection = null;
_factory = null;

VerifyAndCreateClients();
VerifyAndCreateClients(clientName);
}

return _connection;
Expand Down
138 changes: 63 additions & 75 deletions Providers/Resgrid.Providers.Bus.Rabbit/RabbitInboundEventProvider.cs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
using System;
using System.Text;
using System.Threading;
using System.Threading.Channels;
using System.Threading.Tasks;
using Newtonsoft.Json;
using RabbitMQ.Client;
Expand All @@ -15,7 +16,6 @@ namespace Resgrid.Providers.Bus.Rabbit
{
public class RabbitInboundEventProvider : IRabbitInboundEventProvider
{
//private ConnectionFactory _factory;
private IConnection _connection;
private IModel _channel;

Expand All @@ -28,25 +28,25 @@ public class RabbitInboundEventProvider : IRabbitInboundEventProvider
public Func<int, PersonnelLocationUpdatedEvent, Task> PersonnelLocationUpdated;
public Func<int, UnitLocationUpdatedEvent, Task> UnitLocationUpdated;

public async Task Start()
public async Task Start(string clientName, string queueName)
{
VerifyAndCreateClients();
await StartMonitoring();
VerifyAndCreateClients(clientName);
await StartMonitoring(queueName);
}

private void VerifyAndCreateClients()
private void VerifyAndCreateClients(string clientName)
{
try
{
_connection = RabbitConnection.CreateConnection();
_connection = RabbitConnection.CreateConnection(clientName);

if (_connection != null)
{
_channel = _connection.CreateModel();

if (_channel != null)
{
_channel.ExchangeDeclare(SetQueueNameForEnv(Topics.EventingTopic), "fanout");
_channel.ExchangeDeclare(RabbitConnection.SetQueueNameForEnv(Topics.EventingTopic), "fanout");
}
}
}
Expand All @@ -56,69 +56,69 @@ private void VerifyAndCreateClients()
}
}

private async Task StartMonitoring()
private async Task StartMonitoring(string queueName)
{
if (SystemBehaviorConfig.ServiceBusType == ServiceBusTypes.Rabbit)
{
var queueName = _channel.QueueDeclare().QueueName;
//var queueName = _channel.QueueDeclare().QueueName;

_channel.QueueBind(queue: queueName,
exchange: SetQueueNameForEnv(Topics.EventingTopic),
routingKey: "");
var queue = _channel.QueueDeclare(RabbitConnection.SetQueueNameForEnv(queueName), durable: true,
autoDelete: false, exclusive: false);

var consumer = new EventingBasicConsumer(_channel);
consumer.Received += async (model, ea) =>
{
var body = ea.Body.ToArray();
var message = Encoding.UTF8.GetString(body);
_channel.QueueBind(queue: queue.QueueName,
exchange: RabbitConnection.SetQueueNameForEnv(Topics.EventingTopic),
routingKey: "");

var consumer = new EventingBasicConsumer(_channel);
consumer.Received += async (model, ea) =>
{
var body = ea.Body.ToArray();
var message = Encoding.UTF8.GetString(body);

var eventingMessage = JsonConvert.DeserializeObject<EventingMessage>(message);
var eventingMessage = JsonConvert.DeserializeObject<EventingMessage>(message);

if (eventingMessage != null)
if (eventingMessage != null)
{
switch ((EventingTypes)eventingMessage.Type)
{
switch ((EventingTypes)eventingMessage.Type)
{
case EventingTypes.PersonnelStatusUpdated:
if (ProcessPersonnelStatusChanged != null)
await ProcessPersonnelStatusChanged(eventingMessage.DepartmentId, eventingMessage.ItemId);
break;
case EventingTypes.UnitStatusUpdated:
if (ProcessUnitStatusChanged != null)
await ProcessUnitStatusChanged.Invoke(eventingMessage.DepartmentId, eventingMessage.ItemId);
break;
case EventingTypes.CallsUpdated:
if (ProcessCallStatusChanged != null)
await ProcessCallStatusChanged.Invoke(eventingMessage.DepartmentId, eventingMessage.ItemId);
break;
case EventingTypes.CallAdded:
if (ProcessCallStatusChanged != null)
await ProcessCallStatusChanged.Invoke(eventingMessage.DepartmentId, eventingMessage.ItemId);
break;
case EventingTypes.CallClosed:
if (ProcessCallStatusChanged != null)
await ProcessCallStatusChanged.Invoke(eventingMessage.DepartmentId, eventingMessage.ItemId);
break;
case EventingTypes.PersonnelStaffingUpdated:
if (ProcessPersonnelStaffingChanged != null)
await ProcessPersonnelStaffingChanged.Invoke(eventingMessage.DepartmentId, eventingMessage.ItemId);
break;
case EventingTypes.PersonnelLocationUpdated:
if (PersonnelLocationUpdated != null)
await PersonnelLocationUpdated.Invoke(eventingMessage.DepartmentId, JsonConvert.DeserializeObject<PersonnelLocationUpdatedEvent>(eventingMessage.Payload));
break;
case EventingTypes.UnitLocationUpdated:
if (UnitLocationUpdated != null)
await UnitLocationUpdated.Invoke(eventingMessage.DepartmentId, JsonConvert.DeserializeObject<UnitLocationUpdatedEvent>(eventingMessage.Payload));
break;
default:
throw new ArgumentOutOfRangeException();
}
case EventingTypes.PersonnelStatusUpdated:
if (ProcessPersonnelStatusChanged != null)
await ProcessPersonnelStatusChanged(eventingMessage.DepartmentId, eventingMessage.ItemId);
break;
case EventingTypes.UnitStatusUpdated:
if (ProcessUnitStatusChanged != null)
await ProcessUnitStatusChanged.Invoke(eventingMessage.DepartmentId, eventingMessage.ItemId);
break;
case EventingTypes.CallsUpdated:
if (ProcessCallStatusChanged != null)
await ProcessCallStatusChanged.Invoke(eventingMessage.DepartmentId, eventingMessage.ItemId);
break;
case EventingTypes.CallAdded:
if (ProcessCallAdded != null)
await ProcessCallAdded.Invoke(eventingMessage.DepartmentId, eventingMessage.ItemId);
break;
case EventingTypes.CallClosed:
if (ProcessCallClosed != null)
await ProcessCallClosed.Invoke(eventingMessage.DepartmentId, eventingMessage.ItemId);
break;
case EventingTypes.PersonnelStaffingUpdated:
if (ProcessPersonnelStaffingChanged != null)
await ProcessPersonnelStaffingChanged.Invoke(eventingMessage.DepartmentId, eventingMessage.ItemId);
break;
case EventingTypes.PersonnelLocationUpdated:
if (PersonnelLocationUpdated != null)
await PersonnelLocationUpdated.Invoke(eventingMessage.DepartmentId, JsonConvert.DeserializeObject<PersonnelLocationUpdatedEvent>(eventingMessage.Payload));
break;
case EventingTypes.UnitLocationUpdated:
if (UnitLocationUpdated != null)
await UnitLocationUpdated.Invoke(eventingMessage.DepartmentId, JsonConvert.DeserializeObject<UnitLocationUpdatedEvent>(eventingMessage.Payload));
break;
default:
throw new ArgumentOutOfRangeException();
}
};
_channel.BasicConsume(queue: queueName,
autoAck: true,
consumer: consumer);
}
}
};
_channel.BasicConsume(queue: queue.QueueName,
autoAck: true,
consumer: consumer);
}

public bool IsConnected()
Expand Down Expand Up @@ -147,17 +147,5 @@ public void RegisterForEvents(Func<int, string, Task> personnelStatusChanged,
PersonnelLocationUpdated = personnelLocationUpdated;
UnitLocationUpdated = unitLocationUpdated;
}

private static string SetQueueNameForEnv(string cacheKey)
{
if (Config.SystemBehaviorConfig.Environment == SystemEnvironment.Dev)
return $"DEV{cacheKey}";
else if (Config.SystemBehaviorConfig.Environment == SystemEnvironment.QA)
return $"QA{cacheKey}";
else if (Config.SystemBehaviorConfig.Environment == SystemEnvironment.Staging)
return $"ST{cacheKey}";

return cacheKey;
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ namespace Resgrid.Providers.Bus.Rabbit
{
public class RabbitInboundQueueProvider
{
private string _clientName;
private IModel _channel;
public Func<CallQueueItem, Task> CallQueueReceived;
public Func<MessageQueueItem, Task> MessageQueueReceived;
Expand All @@ -32,9 +33,10 @@ public RabbitInboundQueueProvider()
RabbitOutboundQueueProvider provider = new RabbitOutboundQueueProvider();
}

public async Task Start()
public async Task Start(string clientName)
{
var connection = RabbitConnection.CreateConnection();
_clientName = clientName;
var connection = RabbitConnection.CreateConnection(clientName);

if (connection != null)
{
Expand Down Expand Up @@ -584,7 +586,7 @@ private bool RetryQueueItem(BasicDeliverEventArgs ea, Exception mex)
//var factory = new ConnectionFactory() { HostName = ServiceBusConfig.RabbitHostname, UserName = ServiceBusConfig.RabbitUsername, Password = ServiceBusConfig.RabbbitPassword };
//using (var connection = RabbitConnection.CreateConnection())
//{
var connection = RabbitConnection.CreateConnection();
var connection = RabbitConnection.CreateConnection(_clientName);
if (connection != null)
{
using (var channel = connection.CreateModel())
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,8 @@ namespace Resgrid.Providers.Bus.Rabbit
{
public class RabbitOutboundQueueProvider : IRabbitOutboundQueueProvider
{
private readonly string _clientName = "Resgrid-Outbound";

public bool EnqueueCall(CallQueueItem callQueue)
{
string serializedObject = ObjectSerialization.Serialize(callQueue);
Expand Down Expand Up @@ -97,11 +99,6 @@ public bool EnqueueSecurityRefreshEvent(SecurityRefreshEvent securityRefreshEven
return SendMessage(ServiceBusConfig.SecurityRefreshQueueName, serializedObject, false, "300000");
}

public bool VerifyAndCreateClients()
{
return RabbitConnection.VerifyAndCreateClients();
}

private bool SendMessage(string queueName, string message, bool durable = true, string expiration = "36000000")
{
if (String.IsNullOrWhiteSpace(queueName))
Expand All @@ -110,15 +107,9 @@ private bool SendMessage(string queueName, string message, bool durable = true,
if (String.IsNullOrWhiteSpace(message))
throw new ArgumentNullException("message");

//if (SystemBehaviorConfig.ServiceBusType == ServiceBusTypes.Rabbit)
//{
try
{
// TODO: Maybe? https://github.com/EasyNetQ/EasyNetQ -SJ
//var factory = new ConnectionFactory() { HostName = ServiceBusConfig.RabbitHostname, UserName = ServiceBusConfig.RabbitUsername, Password = ServiceBusConfig.RabbbitPassword };
//using (var connection = RabbitConnection.CreateConnection())
//{
var connection = RabbitConnection.CreateConnection();
var connection = RabbitConnection.CreateConnection(_clientName);
if (connection != null)
{
using (var channel = connection.CreateModel())
Expand Down Expand Up @@ -157,16 +148,17 @@ private bool SendMessage(string queueName, string message, bool durable = true,
}

return false;
//}
}
catch (Exception ex)
{
Logging.LogException(ex);
return false;
}
//}
}

//return false;
bool IRabbitOutboundQueueProvider.VerifyAndCreateClients()
{
return RabbitConnection.VerifyAndCreateClients(_clientName);
}
}
}
Loading

0 comments on commit 4c21f2c

Please sign in to comment.