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 componentFilters = new(); private readonly ConcurrentDictionary componentEvents = new(); private readonly ConcurrentDictionary> deferredMessages = new(); private readonly ConcurrentQueue messageQueue = new(); private readonly SemaphoreSlim sendingSemaphore = new(1, 1); private static ILogger? LOG; private MessageBus() { } public void Initialize(ILogger logger) { LOG = logger; LOG.LogInformation("Message bus initialized."); } /// /// Define for which components and events you want to receive messages. /// /// That's you, the receiver. /// A list of components for which you want to receive messages. Use an empty list to receive messages from all components. /// A list of events for which you want to receive messages. public void ApplyFilters(IMessageBusReceiver receiver, ComponentBase[] filterComponents, HashSet events) { this.componentFilters[receiver] = filterComponents; this.componentEvents[receiver] = events.ToArray(); } 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(ComponentBase? sendingComponent, Event triggeredEvent, T? data = default) { this.messageQueue.Enqueue(new Message(sendingComponent, triggeredEvent, data)); try { await this.sendingSemaphore.WaitAsync(); while (this.messageQueue.TryDequeue(out var message)) { foreach (var (receiver, componentFilter) in this.componentFilters) { if (componentFilter.Length > 0 && message.SendingComponent is not null && !componentFilter.Contains(message.SendingComponent)) continue; var eventFilter = this.componentEvents[receiver]; if (eventFilter.Length == 0 || eventFilter.Contains(message.TriggeredEvent)) // 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); } } } catch (Exception e) { LOG?.LogError(e, "Error while sending message."); } finally { this.sendingSemaphore.Release(); } } 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); public Task SendInfo(DataInfoMessage dataInfoMessage) => this.SendMessage(null, Event.SHOW_INFO, dataInfoMessage); /// /// Stores a message until someone asks for it, cf. TakeDeferredMessages. This is how a /// component hands data to a component which does not exist yet, e.g. an assistant which /// sends its result to the chat before the user gets there. /// /// That's you, the sender. /// The event this message belongs to. /// The data to hand over. public void DeferMessage(ComponentBase? sendingComponent, Event triggeredEvent, T? data = default) { var queue = this.deferredMessages.GetOrAdd(triggeredEvent, _ => new()); queue.Enqueue(new Message(sendingComponent, triggeredEvent, data)); } /// /// Takes all deferred messages of an event out of the bus. /// /// /// This empties the queue and returns what was in it. It used to be a lazy iterator, which /// meant that a caller stopping after the first message left the rest of the queue behind: /// those messages were never delivered, and the data they carry — a complete chat thread, for /// instance — stayed alive for as long as the app ran. Returning a list makes that impossible. /// Callers who expect a single message take the last one, since that is the most recent thing /// the user asked for. /// /// The event whose messages you want. /// The deferred messages, oldest first. Empty when there are none. public IReadOnlyList TakeDeferredMessages(Event triggeredEvent) { // // Removing the queue along with its messages is what keeps the dictionary from growing: // otherwise, every event which ever deferred a message would keep an empty queue forever. // if (!this.deferredMessages.TryRemove(triggeredEvent, out var queue)) return []; var messages = new List(); while (queue.TryDequeue(out var message)) messages.Add(message.Data is T data ? data : default); return messages; } public async Task SendMessageUseFirstResult(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(sendingComponent, triggeredEvent, data); if (result is not null) return (TResult) result; } } return default; } }