feat(eventbus): add platform runtimes and benchmarks

Add the Entities NativeQueue adapter, standalone .NET runtime and source generator, reproducible smoke coverage, and Unity benchmark assets for EventBus 2.0.
This commit is contained in:
2026-08-26 01:17:39 +08:00
parent ad5a7b68a3
commit 724e0bc8d8
32 changed files with 1813 additions and 0 deletions
+205
View File
@@ -0,0 +1,205 @@
#nullable enable
using System;
using System.Collections.Concurrent;
using System.Threading;
using System.Threading.Tasks;
namespace ShrinkEventBus
{
internal interface IDotNetScheduler : IDisposable
{
bool IsOnSchedulerThread { get; }
bool TryPost(Action action);
ValueTask PostAsync(Func<ValueTask> action, CancellationToken cancellationToken);
}
internal sealed class QueueFullException : InvalidOperationException { }
internal static class DotNetSchedulerFactory
{
public static IDotNetScheduler Create(string name, ShrinkBusOptions options)
{
return options.Scheduler switch
{
ShrinkBusSchedulerKind.MainThread => new SynchronizationContextScheduler(options),
ShrinkBusSchedulerKind.DedicatedThread => new DedicatedScheduler(name, options),
ShrinkBusSchedulerKind.TaskPool => new TaskPoolScheduler(options),
_ => new InlineScheduler()
};
}
}
internal sealed class InlineScheduler : IDotNetScheduler
{
public bool IsOnSchedulerThread => true;
public bool TryPost(Action action) { action(); return true; }
public ValueTask PostAsync(Func<ValueTask> action, CancellationToken cancellationToken)
{
cancellationToken.ThrowIfCancellationRequested();
return action();
}
public void Dispose() { }
}
internal sealed class SynchronizationContextScheduler : IDotNetScheduler
{
private readonly SynchronizationContext _context;
private readonly int _threadId;
private int _disposed;
public SynchronizationContextScheduler(ShrinkBusOptions options)
{
_context = SynchronizationContext.Current ?? throw new InvalidOperationException(
"MainThread scheduler requires a current SynchronizationContext in a pure .NET host.");
_threadId = Thread.CurrentThread.ManagedThreadId;
}
public bool IsOnSchedulerThread => Thread.CurrentThread.ManagedThreadId == _threadId;
public bool TryPost(Action action)
{
if (Volatile.Read(ref _disposed) != 0)
return false;
if (IsOnSchedulerThread)
action();
else
_context.Post(_ => action(), null);
return true;
}
public ValueTask PostAsync(Func<ValueTask> action, CancellationToken cancellationToken)
{
if (IsOnSchedulerThread)
return action();
return new ValueTask(PostCoreAsync(action, cancellationToken));
}
private Task PostCoreAsync(Func<ValueTask> action, CancellationToken cancellationToken)
{
var completion = new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronously);
_context.Post(async _ =>
{
try { cancellationToken.ThrowIfCancellationRequested(); await action(); completion.SetResult(true); }
catch (OperationCanceledException) { completion.SetCanceled(); }
catch (Exception ex) { completion.SetException(ex); }
}, null);
return completion.Task;
}
public void Dispose() => Interlocked.Exchange(ref _disposed, 1);
}
internal sealed class TaskPoolScheduler : IDotNetScheduler
{
private readonly SemaphoreSlim _concurrency;
private readonly int _capacity;
private int _pending;
private int _disposed;
public TaskPoolScheduler(ShrinkBusOptions options)
{
var concurrency = options.DispatchMode == ShrinkDispatchMode.Ordered ? 1 : options.MaxConcurrency;
_concurrency = new SemaphoreSlim(concurrency, concurrency);
_capacity = options.QueueCapacity;
}
public bool IsOnSchedulerThread => false;
public bool TryPost(Action action)
{
if (Volatile.Read(ref _disposed) != 0 || Interlocked.Increment(ref _pending) > _capacity)
{
Interlocked.Decrement(ref _pending);
return false;
}
_ = Task.Run(async () =>
{
await _concurrency.WaitAsync().ConfigureAwait(false);
try { action(); }
finally { _concurrency.Release(); Interlocked.Decrement(ref _pending); }
});
return true;
}
public async ValueTask PostAsync(Func<ValueTask> action, CancellationToken cancellationToken)
{
if (Volatile.Read(ref _disposed) != 0 || Interlocked.Increment(ref _pending) > _capacity)
{
Interlocked.Decrement(ref _pending);
throw new QueueFullException();
}
await _concurrency.WaitAsync(cancellationToken).ConfigureAwait(false);
try { await action().ConfigureAwait(false); }
finally { _concurrency.Release(); Interlocked.Decrement(ref _pending); }
}
public void Dispose()
{
Interlocked.Exchange(ref _disposed, 1);
_concurrency.Dispose();
}
}
internal sealed class DedicatedScheduler : IDotNetScheduler
{
private sealed class Work
{
public Func<ValueTask> Action = null!;
public TaskCompletionSource<bool>? Completion;
}
private readonly BlockingCollection<Work> _queue;
private readonly Thread _thread;
private int _threadId;
private int _disposed;
public DedicatedScheduler(string name, ShrinkBusOptions options)
{
_queue = new BlockingCollection<Work>(options.QueueCapacity);
_thread = new Thread(Run) { IsBackground = true, Name = $"ShrinkBus:{name}" };
_thread.Start();
}
public bool IsOnSchedulerThread => Thread.CurrentThread.ManagedThreadId == Volatile.Read(ref _threadId);
public bool TryPost(Action action)
{
if (IsOnSchedulerThread) { action(); return true; }
return Volatile.Read(ref _disposed) == 0 && _queue.TryAdd(new Work
{
Action = () => { action(); return default; }
});
}
public ValueTask PostAsync(Func<ValueTask> action, CancellationToken cancellationToken)
{
if (IsOnSchedulerThread)
return action();
var completion = new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronously);
if (Volatile.Read(ref _disposed) != 0 || !_queue.TryAdd(new Work
{ Action = action, Completion = completion }))
throw new QueueFullException();
return new ValueTask(completion.Task);
}
public void Dispose()
{
if (Interlocked.Exchange(ref _disposed, 1) != 0)
return;
_queue.CompleteAdding();
}
private void Run()
{
Volatile.Write(ref _threadId, Thread.CurrentThread.ManagedThreadId);
foreach (var work in _queue.GetConsumingEnumerable())
{
try { work.Action().AsTask().GetAwaiter().GetResult(); work.Completion?.SetResult(true); }
catch (OperationCanceledException) { work.Completion?.SetCanceled(); }
catch (Exception ex) { work.Completion?.SetException(ex); }
}
_queue.Dispose();
}
}
}
@@ -0,0 +1,16 @@
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
<TargetFramework>netstandard2.1</TargetFramework>
<LangVersion>latest</LangVersion>
<Nullable>enable</Nullable>
<ImplicitUsings>disable</ImplicitUsings>
<AssemblyName>ShrinkEventBus.Core</AssemblyName>
<RootNamespace>ShrinkEventBus</RootNamespace>
</PropertyGroup>
<ItemGroup>
<Compile Include="..\..\Assets\Modules\ShrinkEventBus\Runtime\ShrinkEventContracts.cs"
Link="Contracts\ShrinkEventContracts.cs" />
<Compile Include="..\..\Assets\Modules\ShrinkEventBus\Runtime\ShrinkEventAttributes.cs"
Link="Contracts\ShrinkEventAttributes.cs" />
</ItemGroup>
</Project>
@@ -0,0 +1,494 @@
#nullable enable
using System;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
namespace ShrinkEventBus
{
public delegate ValueTask ShrinkAsyncEventHandler<in TEvent>(TEvent eventData,
CancellationToken cancellationToken) where TEvent : IShrinkEvent;
public interface IShrinkEventBus : IDisposable
{
ShrinkBusKey Key { get; }
ShrinkBusOptions Options { get; }
ShrinkPostResult Post<TEvent>(in TEvent eventData) where TEvent : IShrinkEvent;
ValueTask<ShrinkPostResult> PostAsync<TEvent>(TEvent eventData,
CancellationToken cancellationToken = default) where TEvent : IShrinkEvent;
IDisposable Attach(object target);
}
public interface IShrinkBusResolver
{
IShrinkEventBus GetBus(ShrinkBusKey key);
bool TryGetBus(ShrinkBusKey key, out IShrinkEventBus bus);
}
public interface IShrinkGeneratedSubscriber
{
IDisposable AttachGenerated(IShrinkBusResolver resolver, ShrinkBusKey? defaultBus = null);
}
public sealed class ShrinkEventBinding : IDisposable
{
private readonly List<IDisposable> _items = new List<IDisposable>();
private bool _disposed;
public void Add(IDisposable item)
{
if (item == null)
throw new ArgumentNullException(nameof(item));
if (_disposed)
{
item.Dispose();
throw new ObjectDisposedException(nameof(ShrinkEventBinding));
}
_items.Add(item);
}
public void Dispose()
{
if (_disposed)
return;
_disposed = true;
for (var i = _items.Count - 1; i >= 0; i--)
_items[i].Dispose();
_items.Clear();
}
}
public static class ShrinkStaticBindingRegistry
{
internal sealed class Entry
{
public long Id;
public ShrinkBusKey Key;
public Func<IShrinkBusResolver, IDisposable> Factory = null!;
}
private static readonly object Gate = new object();
private static readonly List<Entry> Entries = new List<Entry>();
private static readonly List<WeakReference<ShrinkEventBusHost>> Hosts =
new List<WeakReference<ShrinkEventBusHost>>();
private static long _nextId;
public static void Register(ShrinkBusKey key,
Func<IShrinkBusResolver, IDisposable> factory)
{
if (factory == null)
throw new ArgumentNullException(nameof(factory));
ShrinkEventBusHost[] hosts;
Entry entry;
lock (Gate)
{
entry = new Entry
{
Id = Interlocked.Increment(ref _nextId),
Key = key,
Factory = factory
};
Entries.Add(entry);
hosts = LiveHostsLocked();
}
for (var i = 0; i < hosts.Length; i++)
hosts[i].TryAttachStatic(entry);
}
internal static void TrackHost(ShrinkEventBusHost host)
{
lock (Gate)
{
LiveHostsLocked();
Hosts.Add(new WeakReference<ShrinkEventBusHost>(host));
}
}
internal static Entry[] GetEntries(ShrinkBusKey key)
{
lock (Gate)
return Entries.FindAll(item => item.Key == key).ToArray();
}
private static ShrinkEventBusHost[] LiveHostsLocked()
{
var result = new List<ShrinkEventBusHost>(Hosts.Count);
for (var i = Hosts.Count - 1; i >= 0; i--)
{
if (Hosts[i].TryGetTarget(out var host))
result.Add(host);
else
Hosts.RemoveAt(i);
}
return result.ToArray();
}
}
public sealed class ShrinkEventBusHost : IShrinkBusResolver, IDisposable
{
private readonly ConcurrentDictionary<ShrinkBusKey, IShrinkEventBus> _buses =
new ConcurrentDictionary<ShrinkBusKey, IShrinkEventBus>();
private readonly object _staticGate = new object();
private readonly Dictionary<long, IDisposable> _staticBindings =
new Dictionary<long, IDisposable>();
private bool _disposed;
public ShrinkEventBusHost()
{
ShrinkStaticBindingRegistry.TrackHost(this);
}
public IShrinkEventBus CreateBus(ShrinkBusKey key, ShrinkBusOptions options)
{
if (_disposed)
throw new ObjectDisposedException(nameof(ShrinkEventBusHost));
var bus = new DotNetShrinkEventBus(key, options);
if (!_buses.TryAdd(key, bus))
{
bus.Dispose();
throw new InvalidOperationException($"Bus '{key}' is already registered.");
}
var entries = ShrinkStaticBindingRegistry.GetEntries(key);
for (var i = 0; i < entries.Length; i++)
TryAttachStatic(entries[i]);
return bus;
}
public IShrinkEventBus GetBus(ShrinkBusKey key)
{
if (_buses.TryGetValue(key, out var bus))
return bus;
throw new KeyNotFoundException($"Bus '{key}' is not registered.");
}
public bool TryGetBus(ShrinkBusKey key, out IShrinkEventBus bus) =>
_buses.TryGetValue(key, out bus!);
public IDisposable Attach(object target, ShrinkBusKey? defaultBus = null)
{
if (target is not IShrinkGeneratedSubscriber generated)
throw new InvalidOperationException(
$"Type {target?.GetType().FullName ?? "<null>"} has no generated event binding.");
return generated.AttachGenerated(this, defaultBus);
}
public void Dispose()
{
if (_disposed)
return;
_disposed = true;
lock (_staticGate)
{
foreach (var binding in _staticBindings.Values)
binding.Dispose();
_staticBindings.Clear();
}
foreach (var bus in _buses.Values)
bus.Dispose();
_buses.Clear();
}
internal void TryAttachStatic(ShrinkStaticBindingRegistry.Entry entry)
{
lock (_staticGate)
{
if (_disposed || _staticBindings.ContainsKey(entry.Id) || !_buses.ContainsKey(entry.Key))
return;
_staticBindings.Add(entry.Id, entry.Factory(this));
}
}
}
public static class ShrinkGeneratedBinding
{
public static IDisposable Subscribe<TEvent>(IShrinkBusResolver resolver, ShrinkBusKey? defaultBus,
string? configuredBus, object? owner, Action<TEvent> handler, ShrinkEventPriority priority,
int numericPriority, bool receiveCanceled) where TEvent : IShrinkEvent
{
return Resolve(resolver, defaultBus, configuredBus).Subscribe(
owner, handler, priority, numericPriority, receiveCanceled);
}
public static IDisposable SubscribeAsync<TEvent>(IShrinkBusResolver resolver, ShrinkBusKey? defaultBus,
string? configuredBus, object? owner, ShrinkAsyncEventHandler<TEvent> handler,
ShrinkEventPriority priority, int numericPriority, bool receiveCanceled) where TEvent : IShrinkEvent
{
return Resolve(resolver, defaultBus, configuredBus).SubscribeAsync(
owner, handler, priority, numericPriority, receiveCanceled);
}
private static DotNetShrinkEventBus Resolve(IShrinkBusResolver resolver, ShrinkBusKey? defaultBus,
string? configuredBus)
{
var key = string.IsNullOrWhiteSpace(configuredBus)
? defaultBus ?? ShrinkBusKey.Game
: ShrinkBusKey.Parse(configuredBus);
return resolver.GetBus(key) as DotNetShrinkEventBus
?? throw new InvalidOperationException("Generated bindings require the built-in bus implementation.");
}
}
internal sealed class DotNetShrinkEventBus : IShrinkEventBus
{
private readonly ConcurrentDictionary<Type, IEventChannel> _channels =
new ConcurrentDictionary<Type, IEventChannel>();
private readonly IDotNetScheduler _scheduler;
private long _nextSubscriptionId;
private bool _disposed;
public DotNetShrinkEventBus(ShrinkBusKey key, ShrinkBusOptions options)
{
Key = key;
Options = (options ?? throw new ArgumentNullException(nameof(options))).CloneValidated();
_scheduler = DotNetSchedulerFactory.Create(key.Id, Options);
}
public ShrinkBusKey Key { get; }
public ShrinkBusOptions Options { get; }
public ShrinkPostResult Post<TEvent>(in TEvent eventData) where TEvent : IShrinkEvent
{
if (_disposed)
return ShrinkPostResult.Rejected(ShrinkPostFailure.BusStopped);
var copied = eventData;
if (_scheduler.IsOnSchedulerThread)
return Channel<TEvent>().Post(copied);
return _scheduler.TryPost(() => Channel<TEvent>().Post(copied))
? new ShrinkPostResult(true, false, false)
: ShrinkPostResult.Rejected(ShrinkPostFailure.QueueFull);
}
public async ValueTask<ShrinkPostResult> PostAsync<TEvent>(TEvent eventData,
CancellationToken cancellationToken = default) where TEvent : IShrinkEvent
{
if (_disposed)
return ShrinkPostResult.Rejected(ShrinkPostFailure.BusStopped);
ShrinkPostResult result = default;
try
{
await _scheduler.PostAsync(async () =>
{
result = await Channel<TEvent>().PostAsync(eventData, Options.DispatchMode, cancellationToken)
.ConfigureAwait(false);
}, cancellationToken).ConfigureAwait(false);
return result;
}
catch (QueueFullException)
{
return ShrinkPostResult.Rejected(ShrinkPostFailure.QueueFull);
}
catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
{
return new ShrinkPostResult(true, false, false, ShrinkPostFailure.Canceled);
}
}
public IDisposable Attach(object target)
{
if (target is not IShrinkGeneratedSubscriber generated)
throw new InvalidOperationException(
$"Type {target?.GetType().FullName ?? "<null>"} has no generated event binding.");
return generated.AttachGenerated(new SingleBusResolver(this), Key);
}
public IDisposable Subscribe<TEvent>(object? owner, Action<TEvent> handler,
ShrinkEventPriority priority, int numericPriority, bool receiveCanceled) where TEvent : IShrinkEvent
{
return Channel<TEvent>().Add(Interlocked.Increment(ref _nextSubscriptionId), owner,
handler, null, priority, numericPriority, receiveCanceled);
}
public IDisposable SubscribeAsync<TEvent>(object? owner, ShrinkAsyncEventHandler<TEvent> handler,
ShrinkEventPriority priority, int numericPriority, bool receiveCanceled) where TEvent : IShrinkEvent
{
return Channel<TEvent>().Add(Interlocked.Increment(ref _nextSubscriptionId), owner,
null, handler, priority, numericPriority, receiveCanceled);
}
public void Dispose()
{
if (_disposed)
return;
_disposed = true;
_scheduler.Dispose();
foreach (var channel in _channels.Values)
channel.Clear();
_channels.Clear();
}
private EventChannel<TEvent> Channel<TEvent>() where TEvent : IShrinkEvent =>
(EventChannel<TEvent>)_channels.GetOrAdd(typeof(TEvent), _ => new EventChannel<TEvent>());
private sealed class SingleBusResolver : IShrinkBusResolver
{
private readonly IShrinkEventBus _bus;
public SingleBusResolver(IShrinkEventBus bus) => _bus = bus;
public IShrinkEventBus GetBus(ShrinkBusKey key) => key == _bus.Key
? _bus
: throw new KeyNotFoundException($"Bus '{key}' is not available.");
public bool TryGetBus(ShrinkBusKey key, out IShrinkEventBus bus)
{
bus = _bus;
return key == _bus.Key;
}
}
}
internal interface IEventChannel
{
void Clear();
}
internal sealed class EventChannel<TEvent> : IEventChannel where TEvent : IShrinkEvent
{
private sealed class Entry
{
public long Id = 0;
public object? Owner;
public Action<TEvent>? Sync;
public ShrinkAsyncEventHandler<TEvent>? Async;
public ShrinkEventPriority Priority;
public int NumericPriority;
public bool ReceiveCanceled;
}
private readonly object _gate = new object();
private readonly List<Entry> _entries = new List<Entry>();
private Entry[] _snapshot = Array.Empty<Entry>();
public IDisposable Add(long id, object? owner, Action<TEvent>? sync,
ShrinkAsyncEventHandler<TEvent>? asyncHandler, ShrinkEventPriority priority,
int numericPriority, bool receiveCanceled)
{
lock (_gate)
{
_entries.Add(new Entry
{
Id = id,
Owner = owner,
Sync = sync,
Async = asyncHandler,
Priority = priority,
NumericPriority = numericPriority,
ReceiveCanceled = receiveCanceled
});
Rebuild();
}
return new Subscription(this, id);
}
public ShrinkPostResult Post(TEvent eventData)
{
var handled = false;
var snapshot = _snapshot;
for (var i = 0; i < snapshot.Length; i++)
{
var entry = snapshot[i];
if (IsCanceled(eventData) && !entry.ReceiveCanceled)
continue;
if (entry.Sync != null)
entry.Sync(eventData);
else if (entry.Async != null)
_ = entry.Async(eventData, CancellationToken.None).AsTask();
else
continue;
handled = true;
}
return ShrinkPostResult.Completed(handled, IsCanceled(eventData));
}
public async ValueTask<ShrinkPostResult> PostAsync(TEvent eventData, ShrinkDispatchMode mode,
CancellationToken cancellationToken)
{
var handled = false;
var snapshot = _snapshot;
if (mode == ShrinkDispatchMode.Parallel)
{
var tasks = new List<Task>();
for (var i = 0; i < snapshot.Length; i++)
{
var entry = snapshot[i];
if (IsCanceled(eventData) && !entry.ReceiveCanceled)
continue;
if (entry.Sync != null)
entry.Sync(eventData);
else if (entry.Async != null)
tasks.Add(entry.Async(eventData, cancellationToken).AsTask());
else
continue;
handled = true;
}
if (tasks.Count > 0)
await Task.WhenAll(tasks).ConfigureAwait(false);
}
else
{
for (var i = 0; i < snapshot.Length; i++)
{
cancellationToken.ThrowIfCancellationRequested();
var entry = snapshot[i];
if (IsCanceled(eventData) && !entry.ReceiveCanceled)
continue;
if (entry.Sync != null)
entry.Sync(eventData);
else if (entry.Async != null)
await entry.Async(eventData, cancellationToken).ConfigureAwait(false);
else
continue;
handled = true;
}
}
return ShrinkPostResult.Completed(handled, IsCanceled(eventData));
}
public void Clear()
{
lock (_gate)
{
_entries.Clear();
_snapshot = Array.Empty<Entry>();
}
}
private void Remove(long id)
{
lock (_gate)
{
_entries.RemoveAll(item => item.Id == id);
Rebuild();
}
}
private void Rebuild()
{
_entries.Sort((left, right) =>
{
var leftValue = left.NumericPriority == 0 ? (int)left.Priority * 1000 : -left.NumericPriority;
var rightValue = right.NumericPriority == 0 ? (int)right.Priority * 1000 : -right.NumericPriority;
var result = leftValue.CompareTo(rightValue);
return result != 0 ? result : left.Id.CompareTo(right.Id);
});
_snapshot = _entries.ToArray();
}
private static bool IsCanceled(TEvent value) =>
value is IShrinkCancelableEvent cancelable && cancelable.IsCanceled;
private sealed class Subscription : IDisposable
{
private EventChannel<TEvent>? _owner;
private readonly long _id;
public Subscription(EventChannel<TEvent> owner, long id)
{
_owner = owner;
_id = id;
}
public void Dispose() => Interlocked.Exchange(ref _owner, null)?.Remove(_id);
}
}
}