mirror of
https://github.com/MindWorkAI/AI-Studio.git
synced 2026-10-09 22:33:47 +00:00
Hardened background work against disconnected circuits (#937)
This commit is contained in:
1 parent
99089f853e
commit
f03c2d1c88
44 files changed
+500
-126
No files matched your search
@@ -1,5 +1,7 @@
|
||||
using System.Collections.Concurrent;
|
||||
|
||||
using AIStudio.Tools.Services;
|
||||
|
||||
using Microsoft.AspNetCore.Components;
|
||||
// ReSharper disable RedundantRecordClassKeyword
|
||||
|
||||
@@ -11,6 +13,7 @@ public sealed class MessageBus
|
||||
|
||||
private readonly ConcurrentDictionary<IMessageBusReceiver, ComponentBase[]> componentFilters = new();
|
||||
private readonly ConcurrentDictionary<IMessageBusReceiver, Event[]> componentEvents = new();
|
||||
private readonly ConcurrentDictionary<IMessageBusReceiver, CircuitStateService> receiverCircuits = new();
|
||||
private readonly ConcurrentDictionary<Event, ConcurrentQueue<Message>> deferredMessages = new();
|
||||
private readonly ConcurrentQueue<Message> messageQueue = new();
|
||||
private readonly SemaphoreSlim sendingSemaphore = new(1, 1);
|
||||
@@ -39,16 +42,53 @@ public sealed class MessageBus
|
||||
this.componentEvents[receiver] = events.ToArray();
|
||||
}
|
||||
|
||||
public void RegisterComponent(IMessageBusReceiver receiver)
|
||||
/// <summary>
|
||||
/// Registers a receiver at the bus.
|
||||
/// </summary>
|
||||
/// <param name="receiver">That's you, the receiver.</param>
|
||||
/// <param name="circuitState">The circuit this receiver belongs to. Components hand over their circuit
|
||||
/// so the bus can let them go when that circuit ends. Services which live longer than any circuit,
|
||||
/// such as hosted services, hand over nothing.</param>
|
||||
public void RegisterComponent(IMessageBusReceiver receiver, CircuitStateService? circuitState = null)
|
||||
{
|
||||
this.componentFilters.TryAdd(receiver, []);
|
||||
this.componentEvents.TryAdd(receiver, []);
|
||||
|
||||
if (circuitState is not null)
|
||||
this.receiverCircuits[receiver] = circuitState;
|
||||
}
|
||||
|
||||
|
||||
public void Unregister(IMessageBusReceiver receiver)
|
||||
{
|
||||
this.componentFilters.TryRemove(receiver, out _);
|
||||
this.componentEvents.TryRemove(receiver, out _);
|
||||
this.receiverCircuits.TryRemove(receiver, out _);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Removes all receivers which belong to one circuit.
|
||||
/// </summary>
|
||||
/// <remarks>
|
||||
/// The circuit handler calls this when a circuit ends. Components deregister themselves when they get
|
||||
/// disposed, but a circuit which was retained and then dropped does not give all of them that chance.
|
||||
/// Since the bus holds a strong reference to every receiver, those leftovers would stay and would be
|
||||
/// served forever.
|
||||
/// </remarks>
|
||||
/// <param name="circuitState">The circuit whose receivers must go.</param>
|
||||
/// <returns>The number of removed receivers.</returns>
|
||||
public int UnregisterCircuit(CircuitStateService circuitState)
|
||||
{
|
||||
var numRemovedReceivers = 0;
|
||||
foreach (var (receiver, receiverCircuit) in this.receiverCircuits)
|
||||
{
|
||||
if (!ReferenceEquals(receiverCircuit, circuitState))
|
||||
continue;
|
||||
|
||||
this.Unregister(receiver);
|
||||
numRemovedReceivers++;
|
||||
}
|
||||
|
||||
return numRemovedReceivers;
|
||||
}
|
||||
|
||||
private record class Message(ComponentBase? SendingComponent, Event TriggeredEvent, object? Data);
|
||||
@@ -71,7 +111,7 @@ public sealed class MessageBus
|
||||
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);
|
||||
_ = DeliverMessage(receiver, message);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -85,6 +125,38 @@ public sealed class MessageBus
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Hands one message to one receiver and observes how that went.
|
||||
/// </summary>
|
||||
/// <remarks>
|
||||
/// The bus must not wait for a receiver, since one slow receiver would hold up everybody else. Not
|
||||
/// waiting is not the same as not caring, though: a receiver whose circuit is gone fails with a
|
||||
/// disconnect or disposal exception, and nobody would ever see where it came from. Such a task
|
||||
/// carries its fault until the finalizer reports it as an unobserved task exception — naming a task
|
||||
/// type instead of the receiver and the event. This is where we give those failures a name.
|
||||
/// </remarks>
|
||||
/// <param name="receiver">The receiver of the message.</param>
|
||||
/// <param name="message">The message to deliver.</param>
|
||||
private static async Task DeliverMessage(IMessageBusReceiver receiver, Message message)
|
||||
{
|
||||
try
|
||||
{
|
||||
await receiver.ProcessMessage(message.SendingComponent, message.TriggeredEvent, message.Data);
|
||||
}
|
||||
catch (Exception exception) when (exception is JSDisconnectedException or ObjectDisposedException or OperationCanceledException)
|
||||
{
|
||||
//
|
||||
// Expected whenever the browser connection of a receiver is gone: the app keeps circuits
|
||||
// of reloaded or sleeping windows around, and their components still receive events.
|
||||
//
|
||||
LOG?.LogDebug("The receiver '{ReceiverName}' did not process the event '{Event}' because its circuit was gone: {Reason}", receiver.GetType().Name, message.TriggeredEvent, exception.Message);
|
||||
}
|
||||
catch (Exception exception)
|
||||
{
|
||||
LOG?.LogError(exception, "The receiver '{ReceiverName}' failed while processing the event '{Event}'.", receiver.GetType().Name, message.TriggeredEvent);
|
||||
}
|
||||
}
|
||||
|
||||
public Task SendError(DataErrorMessage dataErrorMessage) => this.SendMessage(null, Event.SHOW_ERROR, dataErrorMessage);
|
||||
|
||||
public Task SendWarning(DataWarningMessage dataWarningMessage) => this.SendMessage(null, Event.SHOW_WARNING, dataWarningMessage);
|
||||
|
||||
Reference in new issue
Block a user