273 lines
9.4 KiB
C#
273 lines
9.4 KiB
C#
#nullable enable
|
|
|
|
using System;
|
|
using System.Collections.Generic;
|
|
using System.Runtime.CompilerServices;
|
|
using System.Threading;
|
|
using Cysharp.Threading.Tasks;
|
|
|
|
namespace ShrinkEventBus
|
|
{
|
|
public delegate UniTask ShrinkAsyncEventHandler<in TEvent>(TEvent eventData,
|
|
CancellationToken cancellationToken) where TEvent : IShrinkEvent;
|
|
|
|
internal interface IShrinkEventChannel
|
|
{
|
|
int Count { get; }
|
|
bool RemoveSubscription(long subscriptionId);
|
|
void Clear();
|
|
void DetachSlot(int slot);
|
|
}
|
|
|
|
internal static class ShrinkEventChannelSlots<TEvent> where TEvent : IShrinkEvent
|
|
{
|
|
private static readonly object Gate = new();
|
|
private static ShrinkEventChannel<TEvent>?[] _slots =
|
|
Array.Empty<ShrinkEventChannel<TEvent>?>();
|
|
|
|
[MethodImpl(MethodImplOptions.AggressiveInlining)]
|
|
public static ShrinkEventChannel<TEvent>? Get(int slot)
|
|
{
|
|
var snapshot = _slots;
|
|
return (uint)slot < (uint)snapshot.Length ? snapshot[slot] : null;
|
|
}
|
|
|
|
public static void Set(int slot, ShrinkEventChannel<TEvent> channel)
|
|
{
|
|
lock (Gate)
|
|
{
|
|
var current = _slots;
|
|
var length = current.Length;
|
|
if (length <= slot)
|
|
{
|
|
length = Math.Max(4, length);
|
|
while (length <= slot)
|
|
length *= 2;
|
|
}
|
|
|
|
var next = new ShrinkEventChannel<TEvent>?[length];
|
|
Array.Copy(current, next, current.Length);
|
|
next[slot] = channel;
|
|
Volatile.Write(ref _slots, next);
|
|
}
|
|
}
|
|
|
|
public static void Clear(int slot, ShrinkEventChannel<TEvent> channel)
|
|
{
|
|
lock (Gate)
|
|
{
|
|
var current = _slots;
|
|
if ((uint)slot >= (uint)current.Length || !ReferenceEquals(current[slot], channel))
|
|
return;
|
|
var next = (ShrinkEventChannel<TEvent>?[])current.Clone();
|
|
next[slot] = null;
|
|
Volatile.Write(ref _slots, next);
|
|
}
|
|
}
|
|
}
|
|
|
|
internal sealed class ShrinkEventChannel<TEvent> : IShrinkEventChannel
|
|
where TEvent : IShrinkEvent
|
|
{
|
|
private static readonly bool SupportsCancellation =
|
|
typeof(IShrinkCancelableEvent).IsAssignableFrom(typeof(TEvent));
|
|
|
|
private sealed class HandlerEntry
|
|
{
|
|
public long SubscriptionId;
|
|
public long RegistrationOrder;
|
|
public Action<TEvent>? SyncHandler;
|
|
public ShrinkAsyncEventHandler<TEvent>? AsyncHandler;
|
|
public ShrinkSubscribeDescriptor Descriptor;
|
|
}
|
|
|
|
private readonly object _gate = new();
|
|
private readonly List<HandlerEntry> _entries = new();
|
|
private HandlerEntry[] _snapshot = Array.Empty<HandlerEntry>();
|
|
private Action<TEvent>? _syncDispatcher;
|
|
|
|
public int Count => _snapshot.Length;
|
|
|
|
public void Add(long subscriptionId, long registrationOrder, Action<TEvent> handler,
|
|
ShrinkSubscribeDescriptor descriptor)
|
|
{
|
|
if (handler == null)
|
|
throw new ArgumentNullException(nameof(handler));
|
|
|
|
lock (_gate)
|
|
{
|
|
_entries.Add(new HandlerEntry
|
|
{
|
|
SubscriptionId = subscriptionId,
|
|
RegistrationOrder = registrationOrder,
|
|
SyncHandler = handler,
|
|
Descriptor = descriptor
|
|
});
|
|
RebuildSnapshot();
|
|
}
|
|
}
|
|
|
|
public void Add(long subscriptionId, long registrationOrder,
|
|
ShrinkAsyncEventHandler<TEvent> handler, ShrinkSubscribeDescriptor descriptor)
|
|
{
|
|
if (handler == null)
|
|
throw new ArgumentNullException(nameof(handler));
|
|
|
|
lock (_gate)
|
|
{
|
|
_entries.Add(new HandlerEntry
|
|
{
|
|
SubscriptionId = subscriptionId,
|
|
RegistrationOrder = registrationOrder,
|
|
AsyncHandler = handler,
|
|
Descriptor = descriptor
|
|
});
|
|
RebuildSnapshot();
|
|
}
|
|
}
|
|
|
|
[MethodImpl(MethodImplOptions.AggressiveInlining)]
|
|
public ShrinkPostResult Post(in TEvent eventData)
|
|
{
|
|
var syncDispatcher = _syncDispatcher;
|
|
if (syncDispatcher != null)
|
|
{
|
|
syncDispatcher(eventData);
|
|
return new ShrinkPostResult(true, true, false);
|
|
}
|
|
|
|
var snapshot = _snapshot;
|
|
var handled = false;
|
|
for (var i = 0; i < snapshot.Length; i++)
|
|
{
|
|
var entry = snapshot[i];
|
|
if (IsCanceled(eventData) && !entry.Descriptor.ReceiveCanceled)
|
|
continue;
|
|
|
|
if (entry.SyncHandler != null)
|
|
entry.SyncHandler(eventData);
|
|
else if (entry.AsyncHandler != null)
|
|
entry.AsyncHandler(eventData, CancellationToken.None).Forget(ShrinkEventDiagnostics.LogException);
|
|
else
|
|
continue;
|
|
handled = true;
|
|
}
|
|
|
|
return ShrinkPostResult.Completed(handled, IsCanceled(eventData));
|
|
}
|
|
|
|
public async UniTask<ShrinkPostResult> PostAsync(TEvent eventData, ShrinkDispatchMode dispatchMode,
|
|
CancellationToken cancellationToken)
|
|
{
|
|
var snapshot = _snapshot;
|
|
var handled = false;
|
|
if (dispatchMode == ShrinkDispatchMode.Parallel)
|
|
{
|
|
var tasks = new List<UniTask>(snapshot.Length);
|
|
for (var i = 0; i < snapshot.Length; i++)
|
|
{
|
|
cancellationToken.ThrowIfCancellationRequested();
|
|
var entry = snapshot[i];
|
|
if (IsCanceled(eventData) && !entry.Descriptor.ReceiveCanceled)
|
|
continue;
|
|
|
|
if (entry.SyncHandler != null)
|
|
entry.SyncHandler(eventData);
|
|
else if (entry.AsyncHandler != null)
|
|
tasks.Add(entry.AsyncHandler(eventData, cancellationToken));
|
|
else
|
|
continue;
|
|
handled = true;
|
|
}
|
|
|
|
if (tasks.Count > 0)
|
|
await UniTask.WhenAll(tasks);
|
|
}
|
|
else
|
|
{
|
|
for (var i = 0; i < snapshot.Length; i++)
|
|
{
|
|
cancellationToken.ThrowIfCancellationRequested();
|
|
var entry = snapshot[i];
|
|
if (IsCanceled(eventData) && !entry.Descriptor.ReceiveCanceled)
|
|
continue;
|
|
|
|
if (entry.SyncHandler != null)
|
|
entry.SyncHandler(eventData);
|
|
else if (entry.AsyncHandler != null)
|
|
await entry.AsyncHandler(eventData, cancellationToken);
|
|
else
|
|
continue;
|
|
handled = true;
|
|
}
|
|
}
|
|
|
|
return ShrinkPostResult.Completed(handled, IsCanceled(eventData));
|
|
}
|
|
|
|
public bool RemoveSubscription(long subscriptionId)
|
|
{
|
|
lock (_gate)
|
|
{
|
|
var removed = _entries.RemoveAll(entry => entry.SubscriptionId == subscriptionId) > 0;
|
|
RebuildSnapshot();
|
|
return removed;
|
|
}
|
|
}
|
|
|
|
public void Clear()
|
|
{
|
|
lock (_gate)
|
|
{
|
|
_entries.Clear();
|
|
_snapshot = Array.Empty<HandlerEntry>();
|
|
_syncDispatcher = null;
|
|
}
|
|
}
|
|
|
|
public void DetachSlot(int slot) => ShrinkEventChannelSlots<TEvent>.Clear(slot, this);
|
|
|
|
private void RebuildSnapshot()
|
|
{
|
|
_entries.Sort(static (left, right) =>
|
|
{
|
|
var leftOrder = left.Descriptor.NumericPriority == 0
|
|
? (int)left.Descriptor.Priority * 1000
|
|
: -left.Descriptor.NumericPriority;
|
|
var rightOrder = right.Descriptor.NumericPriority == 0
|
|
? (int)right.Descriptor.Priority * 1000
|
|
: -right.Descriptor.NumericPriority;
|
|
var priority = leftOrder.CompareTo(rightOrder);
|
|
return priority != 0
|
|
? priority
|
|
: left.RegistrationOrder.CompareTo(right.RegistrationOrder);
|
|
});
|
|
_snapshot = _entries.ToArray();
|
|
|
|
Action<TEvent>? dispatcher = null;
|
|
if (!SupportsCancellation)
|
|
{
|
|
for (var i = 0; i < _snapshot.Length; i++)
|
|
{
|
|
var handler = _snapshot[i].SyncHandler;
|
|
if (handler == null)
|
|
{
|
|
dispatcher = null;
|
|
break;
|
|
}
|
|
dispatcher += handler;
|
|
}
|
|
}
|
|
Volatile.Write(ref _syncDispatcher, dispatcher);
|
|
}
|
|
|
|
[MethodImpl(MethodImplOptions.AggressiveInlining)]
|
|
private static bool IsCanceled(TEvent eventData)
|
|
{
|
|
if (!SupportsCancellation)
|
|
return false;
|
|
return ((IShrinkCancelableEvent)(object)eventData!).IsCanceled;
|
|
}
|
|
}
|
|
}
|