2024-06-30 15:26:28 +02:00
using System.Collections.Concurrent ;
using Microsoft.AspNetCore.Components ;
// ReSharper disable RedundantRecordClassKeyword
namespace AIStudio.Tools ;
public sealed class MessageBus
{
public static readonly MessageBus INSTANCE = new ();
private readonly ConcurrentDictionary < IMessageBusReceiver , ComponentBase []> componentFilters = new ();
private readonly ConcurrentDictionary < IMessageBusReceiver , Event []> componentEvents = new ();
2024-08-18 12:32:18 +02:00
private readonly ConcurrentDictionary < Event , ConcurrentQueue < Message >> deferredMessages = new ();
2024-06-30 15:26:28 +02:00
private readonly ConcurrentQueue < Message > messageQueue = new ();
private readonly SemaphoreSlim sendingSemaphore = new ( 1 , 1 );
2025-06-02 20:08:25 +02:00
private static ILogger < MessageBus >? LOG ;
2024-06-30 15:26:28 +02:00
private MessageBus ()
{
}
2025-06-02 20:08:25 +02:00
public void Initialize ( ILogger < MessageBus > logger )
{
LOG = logger ;
LOG . LogInformation ( "Message bus initialized." );
}
2024-06-30 15:26:28 +02:00
2025-04-24 13:50:14 +02:00
/// <summary>
/// Define for which components and events you want to receive messages.
/// </summary>
/// <param name="receiver">That's you, the receiver.</param>
/// <param name="filterComponents">A list of components for which you want to receive messages. Use an empty list to receive messages from all components.</param>
/// <param name="events">A list of events for which you want to receive messages.</param>
2026-03-23 13:57:26 +01:00
public void ApplyFilters ( IMessageBusReceiver receiver , ComponentBase [] filterComponents , HashSet < Event > events )
2024-06-30 15:26:28 +02:00
{
2025-04-24 13:50:14 +02:00
this . componentFilters [ receiver ] = filterComponents ;
2026-03-23 13:57:26 +01:00
this . componentEvents [ receiver ] = events . ToArray ();
2024-06-30 15:26:28 +02:00
}
public void RegisterComponent ( IMessageBusReceiver receiver )
{
this . componentFilters . TryAdd ( receiver , []);
this . componentEvents . TryAdd ( receiver , []);
}
public void Unregister ( IMessageBusReceiver receiver )
{
this . componentFilters . TryRemove ( receiver , out _ );
this . componentEvents . TryRemove ( receiver , out _ );
}
private record class Message ( ComponentBase ? SendingComponent , Event TriggeredEvent , object? Data );
public async Task SendMessage < T >( ComponentBase ? sendingComponent , Event triggeredEvent , T ? data = default )
{
this . messageQueue . Enqueue ( new Message ( sendingComponent , triggeredEvent , data ));
2025-06-02 20:08:25 +02:00
2024-06-30 15:26:28 +02:00
try
{
await this . sendingSemaphore . WaitAsync ();
while ( this . messageQueue . TryDequeue ( out var message ))
{
foreach ( var ( receiver , componentFilter ) in this . componentFilters )
{
2026-05-25 20:48:26 +02:00
if ( componentFilter . Length > 0 && message . SendingComponent is not null && ! componentFilter . Contains ( message . SendingComponent ))
2024-06-30 15:26:28 +02:00
continue ;
var eventFilter = this . componentEvents [ receiver ];
2026-05-25 20:48:26 +02:00
if ( eventFilter . Length == 0 || eventFilter . Contains ( message . TriggeredEvent ))
2025-06-02 20:08:25 +02:00
2024-06-30 15:26:28 +02:00
// We don't await the task here because we don't want to block the message bus:
_ = receiver . ProcessMessage ( message . SendingComponent , message . TriggeredEvent , message . Data );
}
}
}
2025-06-02 20:08:25 +02:00
catch ( Exception e )
{
LOG ?. LogError ( e , "Error while sending message." );
}
2024-06-30 15:26:28 +02:00
finally
{
this . sendingSemaphore . Release ();
}
}
2025-04-04 16:52:33 +02:00
2025-05-29 14:01:56 +02:00
public Task SendError ( DataErrorMessage dataErrorMessage ) => this . SendMessage ( null , Event . SHOW_ERROR , dataErrorMessage );
public Task SendWarning ( DataWarningMessage dataWarningMessage ) => this . SendMessage ( null , Event . SHOW_WARNING , dataWarningMessage );
public Task SendSuccess ( DataSuccessMessage dataSuccessMessage ) => this . SendMessage ( null , Event . SHOW_SUCCESS , dataSuccessMessage );
2025-04-04 16:52:33 +02:00
2026-08-02 21:41:12 +02:00
public Task SendInfo ( DataInfoMessage dataInfoMessage ) => this . SendMessage ( null , Event . SHOW_INFO , dataInfoMessage );
2024-08-18 12:32:18 +02:00
public void DeferMessage < T >( ComponentBase ? sendingComponent , Event triggeredEvent , T ? data = default )
{
if ( this . deferredMessages . TryGetValue ( triggeredEvent , out var queue ))
queue . Enqueue ( new Message ( sendingComponent , triggeredEvent , data ));
else
{
this . deferredMessages [ triggeredEvent ] = new ();
this . deferredMessages [ triggeredEvent ]. Enqueue ( new Message ( sendingComponent , triggeredEvent , data ));
}
}
public IEnumerable < T ?> CheckDeferredMessages < T >( Event triggeredEvent )
{
if ( this . deferredMessages . TryGetValue ( triggeredEvent , out var queue ))
while ( queue . TryDequeue ( out var message ))
yield return message . Data is T data ? data : default ;
}
2024-07-13 10:37:57 +02:00
public async Task < TResult ?> SendMessageUseFirstResult < TPayload , TResult >( ComponentBase ? sendingComponent , Event triggeredEvent , TPayload ? data = default )
{
foreach ( var ( receiver , componentFilter ) in this . componentFilters )
{
if ( componentFilter . Length > 0 && sendingComponent is not null && ! componentFilter . Contains ( sendingComponent ))
continue ;
var eventFilter = this . componentEvents [ receiver ];
if ( eventFilter . Length == 0 || eventFilter . Contains ( triggeredEvent ))
{
var result = await receiver . ProcessMessageWithResult < TPayload , TResult >( sendingComponent , triggeredEvent , data );
if ( result is not null )
return ( TResult ) result ;
}
}
return default ;
}
2024-06-30 15:26:28 +02:00
}