#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(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 where TEvent : IShrinkEvent { private static readonly object Gate = new(); private static ShrinkEventChannel?[] _slots = Array.Empty?>(); [MethodImpl(MethodImplOptions.AggressiveInlining)] public static ShrinkEventChannel? Get(int slot) { var snapshot = _slots; return (uint)slot < (uint)snapshot.Length ? snapshot[slot] : null; } public static void Set(int slot, ShrinkEventChannel 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?[length]; Array.Copy(current, next, current.Length); next[slot] = channel; Volatile.Write(ref _slots, next); } } public static void Clear(int slot, ShrinkEventChannel channel) { lock (Gate) { var current = _slots; if ((uint)slot >= (uint)current.Length || !ReferenceEquals(current[slot], channel)) return; var next = (ShrinkEventChannel?[])current.Clone(); next[slot] = null; Volatile.Write(ref _slots, next); } } } internal sealed class ShrinkEventChannel : 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? SyncHandler; public ShrinkAsyncEventHandler? AsyncHandler; public ShrinkSubscribeDescriptor Descriptor; } private readonly object _gate = new(); private readonly List _entries = new(); private HandlerEntry[] _snapshot = Array.Empty(); private Action? _syncDispatcher; public int Count => _snapshot.Length; public void Add(long subscriptionId, long registrationOrder, Action 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 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 PostAsync(TEvent eventData, ShrinkDispatchMode dispatchMode, CancellationToken cancellationToken) { var snapshot = _snapshot; var handled = false; if (dispatchMode == ShrinkDispatchMode.Parallel) { var tasks = new List(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(); _syncDispatcher = null; } } public void DetachSlot(int slot) => ShrinkEventChannelSlots.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? 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; } } }