From f95f32bcd4a774aff9ad2137f32608a710a2630c Mon Sep 17 00:00:00 2001 From: AriusII Date: Mon, 21 Sep 2026 19:05:25 +0200 Subject: [PATCH] Bound event stream consumers and completion --- .../Events/EventStreamOptions.cs | 6 +- .../Events/EventStreamOverflowPolicy.cs | 9 +- .../Events/IEventStreamLease.cs | 21 +- .../Domains/Events/BoundedEventStream.cs | 217 ++++++++++-------- .../Domains/Events/BoundedEventStreamTests.cs | 119 ++++++++-- .../Domains/Events/EventStreamLeaseTests.cs | 18 ++ 6 files changed, 273 insertions(+), 117 deletions(-) diff --git a/libs/CheatEngine.Client.Abstractions/Events/EventStreamOptions.cs b/libs/CheatEngine.Client.Abstractions/Events/EventStreamOptions.cs index c60a913..751fac3 100644 --- a/libs/CheatEngine.Client.Abstractions/Events/EventStreamOptions.cs +++ b/libs/CheatEngine.Client.Abstractions/Events/EventStreamOptions.cs @@ -1,6 +1,8 @@ namespace CheatEngine.Client.Events; -/// Configures bounded, non-blocking materialization of host callback events. +/// +/// Configures bounded materialization of copied host callback observations for one active asynchronous reader. +/// public readonly record struct EventStreamOptions { /// Creates bounded stream options. @@ -20,7 +22,7 @@ public EventStreamOptions( OverflowPolicy = overflowPolicy; } - /// Gets the mandatory maximum number of events buffered for a consumer. + /// Gets the mandatory maximum number of copied observations buffered for the one active reader. public int Capacity { get; diff --git a/libs/CheatEngine.Client.Abstractions/Events/EventStreamOverflowPolicy.cs b/libs/CheatEngine.Client.Abstractions/Events/EventStreamOverflowPolicy.cs index 37b8048..a2c95c9 100644 --- a/libs/CheatEngine.Client.Abstractions/Events/EventStreamOverflowPolicy.cs +++ b/libs/CheatEngine.Client.Abstractions/Events/EventStreamOverflowPolicy.cs @@ -3,12 +3,15 @@ namespace CheatEngine.Client.Events; /// Specifies the bounded-admission policy for a copied Client event stream. public enum EventStreamOverflowPolicy { - /// Evicts the oldest buffered event before admitting the newest event. + /// Evicts the oldest buffered observation before admitting the newest observation. DropOldest, - /// Rejects the newly raised event while retaining the buffered events. + /// Rejects the newly raised observation while retaining the buffered observations. DropNewest, - /// Terminates the subscription when its bounded buffer cannot admit an event. + /// + /// Terminates observation delivery when its bounded buffer cannot admit an observation. This does not select a + /// native callback disposition. + /// FailSubscription } diff --git a/libs/CheatEngine.Client.Abstractions/Events/IEventStreamLease.cs b/libs/CheatEngine.Client.Abstractions/Events/IEventStreamLease.cs index c166cb5..7bfc5dc 100644 --- a/libs/CheatEngine.Client.Abstractions/Events/IEventStreamLease.cs +++ b/libs/CheatEngine.Client.Abstractions/Events/IEventStreamLease.cs @@ -1,20 +1,31 @@ namespace CheatEngine.Client.Events; -/// Owns a bounded stream of copied callback events for one Client activation. +/// Owns a bounded, unicast stream of copied callback observations for one Client activation. /// The copied event snapshot type. /// -/// Implementations must never wait for a stream consumer on a Cheat Engine or Lua callback thread. Disposing a -/// lease closes admission, detaches its callback, completes , and releases host state. +/// The stream has one active reader at a time; a concurrent reader is rejected and observations are never +/// broadcast. Implementations must not wait for a stream consumer on a Cheat Engine or Lua callback thread, +/// although admission may use a bounded synchronization primitive. Cancelling or disposing an enumerator ends +/// that observation read only; it neither cancels nor decides the native callback. Disposing a lease closes +/// admission, detaches its callback, completes , and releases host state. Activation teardown +/// follows the same release contract, so an expired lease cannot admit later observations. /// public interface IEventStreamLease : IDisposable { - /// Gets the bounded stream of copied events. + /// + /// Gets the bounded, single-active-reader stream of copied observations. Buffered observations are delivered in + /// FIFO order to that reader; overflow is reported through and does not provide + /// a native callback decision. + /// public IAsyncEnumerable Events { get; } - /// Gets the number of events that were not admitted to the stream. + /// + /// Gets the number of observations not delivered because the bounded stream overflowed or terminal processing + /// discarded buffered observations. + /// public long DroppedEventCount { get; diff --git a/libs/CheatEngine.Client.Core/Domains/Events/BoundedEventStream.cs b/libs/CheatEngine.Client.Core/Domains/Events/BoundedEventStream.cs index 06fb809..97c9218 100644 --- a/libs/CheatEngine.Client.Core/Domains/Events/BoundedEventStream.cs +++ b/libs/CheatEngine.Client.Core/Domains/Events/BoundedEventStream.cs @@ -5,8 +5,9 @@ namespace CheatEngine.Client.Core.Domains.Events; /// -/// Provides a bounded, non-blocking handoff from a Cheat Engine callback to an asynchronous consumer. Callback -/// admission is synchronous and never waits for a reader; reader continuations always run asynchronously. +/// Provides a bounded handoff from a Cheat Engine callback to one asynchronous observation consumer. Callback +/// admission does not wait for reader consumption, but it does take a short lock and is not lock-free. +/// Reader continuations always run asynchronously. /// /// The copied event type. internal sealed class BoundedEventStream : IAsyncEnumerable, IDisposable @@ -14,13 +15,14 @@ internal sealed class BoundedEventStream : IAsyncEnumerable, IDisposable private readonly T[] _buffer; private readonly object _gate = new(); private readonly EventStreamOverflowPolicy _overflowPolicy; - private readonly LinkedList _pendingReads = []; + private Enumerator? _activeReader; private bool _completed; private Exception? _completionError; private int _count; private bool _disposed; private bool _isAdmissionOpen = true; private long _lostCount; + private PendingRead? _pendingRead; private int _readIndex; private int _writeIndex; @@ -73,16 +75,30 @@ internal bool IsCompleted } } - /// Returns an asynchronous enumerator over copied events. + /// + /// Returns the sole active asynchronous enumerator. A second concurrent reader is rejected rather than becoming + /// an unbounded competing subscription. + /// public IAsyncEnumerator GetAsyncEnumerator(CancellationToken cancellationToken = default) { - return new Enumerator(this, cancellationToken); + lock (_gate) + { + if (_activeReader is not null) + { + throw new InvalidOperationException( + "A bounded Client event stream supports only one active observation reader."); + } + + Enumerator enumerator = new(this, cancellationToken); + _activeReader = enumerator; + return enumerator; + } } /// Closes admission and discards buffered events during deterministic subscription teardown. public void Dispose() { - List? pendingReads = null; + PendingRead? pendingRead; lock (_gate) { if (_disposed) @@ -93,21 +109,23 @@ public void Dispose() _disposed = true; _isAdmissionOpen = false; _completed = true; - ClearBuffer(); - pendingReads = DetachAllPendingReads(); + _activeReader = null; + DiscardBufferedEvents(); + pendingRead = DetachPendingRead(); } - CompletePendingReads(pendingReads, ReadResult.End); + CompletePendingRead(pendingRead, ReadResult.End); } /// /// Attempts to publish a copied callback event without waiting for an asynchronous consumer. A - /// - /// result means admission has closed or the selected overflow policy terminated the subscription. + /// result means admission has closed or the selected overflow policy terminated the + /// subscription. /// internal bool TryPublish(T value) { - PendingRead? pendingRead = null; + Exception? completionError = null; + PendingRead? pendingRead; lock (_gate) { if (!_isAdmissionOpen) @@ -115,7 +133,7 @@ internal bool TryPublish(T value) return false; } - pendingRead = DetachNextPendingRead(); + pendingRead = DetachPendingRead(); if (pendingRead is null) { if (_count < _buffer.Length) @@ -139,9 +157,10 @@ internal bool TryPublish(T value) case EventStreamOverflowPolicy.FailSubscription: _lostCount++; - CompleteCore(new InvalidOperationException( - "The bounded callback stream overflowed and the subscription was closed."), true); - return false; + completionError = new InvalidOperationException( + "The bounded callback stream overflowed and the subscription was closed."); + pendingRead = CompleteCore(completionError, true); + break; default: throw new InvalidOperationException( @@ -150,17 +169,26 @@ internal bool TryPublish(T value) } } - pendingRead.Complete(new ReadResult(value)); - return true; + if (completionError is null) + { + CompletePendingRead(pendingRead, new ReadResult(value)); + return true; + } + + FailPendingRead(pendingRead, completionError); + return false; } /// Closes admission and lets readers drain the copied events already accepted by the stream. internal void Complete() { + PendingRead? pendingRead; lock (_gate) { - CompleteCore(null, false); + pendingRead = CompleteCore(null, false); } + + CompletePendingRead(pendingRead, ReadResult.End); } /// Closes callback admission without completing existing readers until the host callback is neutralized. @@ -176,22 +204,30 @@ internal void CloseAdmission() internal void Complete(Exception error) { ArgumentNullException.ThrowIfNull(error); + PendingRead? pendingRead; lock (_gate) { - CompleteCore(error, true); + pendingRead = CompleteCore(error, true); } + + FailPendingRead(pendingRead, error); } - private ValueTask ReadAsync(CancellationToken cancellationToken) + private ValueTask ReadAsync(Enumerator reader, CancellationToken cancellationToken) { if (cancellationToken.IsCancellationRequested) { return ValueTask.FromCanceled(cancellationToken); } - PendingRead? pendingRead = null; + PendingRead? pendingRead; lock (_gate) { + if (!ReferenceEquals(_activeReader, reader)) + { + return ValueTask.FromResult(ReadResult.End); + } + if (_count != 0) { return ValueTask.FromResult(new ReadResult(Dequeue())); @@ -207,8 +243,8 @@ private ValueTask ReadAsync(CancellationToken cancellationToken) return ValueTask.FromResult(ReadResult.End); } - pendingRead = new PendingRead(this, cancellationToken); - pendingRead.Node = _pendingReads.AddLast(pendingRead); + pendingRead = new PendingRead(this, reader, cancellationToken); + _pendingRead = pendingRead; } pendingRead.RegisterCancellation(); @@ -219,84 +255,65 @@ private void CancelPendingRead(PendingRead pendingRead, CancellationToken cancel { lock (_gate) { - if (pendingRead.Node is null || !pendingRead.TryBeginCompletion()) + if (!ReferenceEquals(_pendingRead, pendingRead) || !pendingRead.TryBeginCompletion()) { return; } - _pendingReads.Remove(pendingRead.Node); - pendingRead.Node = null; + _pendingRead = null; + if (ReferenceEquals(_activeReader, pendingRead.Reader)) + { + _activeReader = null; + } } pendingRead.Cancel(cancellationToken); } - private PendingRead? DetachNextPendingRead() + private void ReleaseReader(Enumerator reader) { - while (_pendingReads.First is LinkedListNode node) + PendingRead? pendingRead = null; + lock (_gate) { - PendingRead pendingRead = node.Value; - _pendingReads.Remove(node); - pendingRead.Node = null; - if (pendingRead.TryBeginCompletion()) + if (!ReferenceEquals(_activeReader, reader)) { - return pendingRead; + return; } - } - return null; - } - - private List? DetachAllPendingReads() - { - if (_pendingReads.Count == 0) - { - return null; - } - - List detached = new(_pendingReads.Count); - while (_pendingReads.First is LinkedListNode node) - { - PendingRead pendingRead = node.Value; - _pendingReads.RemoveFirst(); - pendingRead.Node = null; - if (pendingRead.TryBeginCompletion()) + _activeReader = null; + if (ReferenceEquals(_pendingRead?.Reader, reader)) { - detached.Add(pendingRead); + pendingRead = DetachPendingRead(); } } - return detached; + CompletePendingRead(pendingRead, ReadResult.End); } - private void CompleteCore(Exception? error, bool clearBufferedEvents) + private PendingRead? DetachPendingRead() + { + PendingRead? pendingRead = _pendingRead; + _pendingRead = null; + return pendingRead is not null && pendingRead.TryBeginCompletion() ? pendingRead : null; + } + + private PendingRead? CompleteCore(Exception? error, bool clearBufferedEvents) { if (_completed) { - return; + return null; } _isAdmissionOpen = false; _completed = true; _completionError = error; - if (clearBufferedEvents) - { - ClearBuffer(); - } - - List? pendingReads = DetachAllPendingReads(); - if (pendingReads is null) - { - return; - } - if (error is null) + if (clearBufferedEvents) { - CompletePendingReads(pendingReads, ReadResult.End); - return; + DiscardBufferedEvents(); } - CompletePendingReads(pendingReads, error); + return DetachPendingRead(); } private void Enqueue(T value) @@ -327,30 +344,25 @@ private void ClearBuffer() _count = 0; } + private void DiscardBufferedEvents() + { + _lostCount += _count; + ClearBuffer(); + } + private int NextIndex(int index) { return index == _buffer.Length - 1 ? 0 : index + 1; } - private static void CompletePendingReads(IEnumerable? pendingReads, ReadResult result) + private static void CompletePendingRead(PendingRead? pendingRead, ReadResult result) { - if (pendingReads is null) - { - return; - } - - foreach (PendingRead pendingRead in pendingReads) - { - pendingRead.Complete(result); - } + pendingRead?.Complete(result); } - private static void CompletePendingReads(IEnumerable pendingReads, Exception error) + private static void FailPendingRead(PendingRead? pendingRead, Exception error) { - foreach (PendingRead pendingRead in pendingReads) - { - pendingRead.Fail(error); - } + pendingRead?.Fail(error); } private readonly record struct ReadResult(bool HasValue, T? Value) @@ -363,7 +375,7 @@ internal ReadResult(T value) internal static ReadResult End => new(false, default); } - private sealed class PendingRead(BoundedEventStream owner, CancellationToken cancellationToken) + private sealed class PendingRead(BoundedEventStream owner, Enumerator reader, CancellationToken cancellationToken) { private readonly CancellationToken _cancellationToken = cancellationToken; private readonly BoundedEventStream _owner = owner; @@ -374,11 +386,7 @@ private sealed class PendingRead(BoundedEventStream owner, CancellationToken private int _completionStarted; private CancellationTokenRegistration _registration; - internal LinkedListNode? Node - { - get; - set; - } + internal Enumerator Reader => reader; internal Task Task => _source.Task; @@ -414,6 +422,7 @@ internal void Complete(ReadResult result) internal void Cancel(CancellationToken cancellationToken) { + _registration.Unregister(); _source.TrySetCanceled(cancellationToken); } @@ -429,6 +438,7 @@ private sealed class Enumerator(BoundedEventStream owner, CancellationToken c { private readonly CancellationToken _cancellationToken = cancellationToken; private readonly BoundedEventStream _owner = owner; + private int _completed; private int _disposed; private int _moveNextInProgress; @@ -440,7 +450,7 @@ public T Current public async ValueTask MoveNextAsync() { - if (Volatile.Read(ref _disposed) != 0) + if (Volatile.Read(ref _disposed) != 0 || Volatile.Read(ref _completed) != 0) { return false; } @@ -453,16 +463,22 @@ public async ValueTask MoveNextAsync() try { - ReadResult result = await _owner.ReadAsync(_cancellationToken).ConfigureAwait(false); + ReadResult result = await _owner.ReadAsync(this, _cancellationToken).ConfigureAwait(false); if (!result.HasValue) { Current = default!; + CompleteReader(); return false; } Current = result.Value!; return true; } + catch + { + CompleteReader(); + throw; + } finally { Volatile.Write(ref _moveNextInProgress, 0); @@ -471,8 +487,21 @@ public async ValueTask MoveNextAsync() public ValueTask DisposeAsync() { - Interlocked.Exchange(ref _disposed, 1); + if (Interlocked.Exchange(ref _disposed, 1) == 0) + { + Current = default!; + CompleteReader(); + } + return ValueTask.CompletedTask; } + + private void CompleteReader() + { + if (Interlocked.Exchange(ref _completed, 1) == 0) + { + _owner.ReleaseReader(this); + } + } } } diff --git a/tests/CheatEngine.Client.Core.Tests/Domains/Events/BoundedEventStreamTests.cs b/tests/CheatEngine.Client.Core.Tests/Domains/Events/BoundedEventStreamTests.cs index 368e408..545b70d 100644 --- a/tests/CheatEngine.Client.Core.Tests/Domains/Events/BoundedEventStreamTests.cs +++ b/tests/CheatEngine.Client.Core.Tests/Domains/Events/BoundedEventStreamTests.cs @@ -81,7 +81,7 @@ public async Task FailSubscriptionClosesTheStreamWithAnErrorAndRejectsFutureAdmi Assert.Contains("overflowed", exception.Message, StringComparison.Ordinal); Assert.True(stream.IsCompleted); Assert.False(stream.IsAdmissionOpen); - Assert.Equal(1, stream.LostCount); + Assert.Equal(2, stream.LostCount); Assert.False(stream.TryPublish(3)); } @@ -101,22 +101,25 @@ public async Task PendingReaderReceivesPublishedValueWithoutBlockingTheCallbackA } [Fact] - public async Task MultipleWaitingReadersReceiveCopiedEventsInAdmissionOrder() + public async Task ConcurrentReadersAreRejectedUntilTheActiveReaderIsDisposed() { using BoundedEventStream stream = new(new EventStreamOptions(1)); - await using IAsyncEnumerator first = stream.GetAsyncEnumerator(TestContext.Current.CancellationToken); - await using IAsyncEnumerator second = stream.GetAsyncEnumerator(TestContext.Current.CancellationToken); - - Task firstMoveNext = first.MoveNextAsync().AsTask(); - Task secondMoveNext = second.MoveNextAsync().AsTask(); + IAsyncEnumerator first = stream.GetAsyncEnumerator(TestContext.Current.CancellationToken); + for (int reader = 0; reader < 64; reader++) + { + InvalidOperationException exception = Assert.Throws(() => + { + _ = stream.GetAsyncEnumerator(TestContext.Current.CancellationToken); + }); + Assert.Contains("one active", exception.Message, StringComparison.Ordinal); + } + + await first.DisposeAsync(); + await using IAsyncEnumerator second = stream.GetAsyncEnumerator(TestContext.Current.CancellationToken); Assert.True(stream.TryPublish(10)); - Assert.True(stream.TryPublish(20)); - - Assert.True(await firstMoveNext); - Assert.Equal(10, first.Current); - Assert.True(await secondMoveNext); - Assert.Equal(20, second.Current); + Assert.True(await second.MoveNextAsync()); + Assert.Equal(10, second.Current); } [Fact] @@ -158,6 +161,23 @@ await Assert.ThrowsAnyAsync(async () => Assert.Equal(0, stream.LostCount); } + [Fact] + public async Task DisposingAWaitingEnumeratorCompletesItsReadAndFreesTheReaderSlot() + { + using BoundedEventStream stream = new(new EventStreamOptions(1)); + IAsyncEnumerator disposedEnumerator = stream.GetAsyncEnumerator(TestContext.Current.CancellationToken); + Task pendingMoveNext = disposedEnumerator.MoveNextAsync().AsTask(); + + await disposedEnumerator.DisposeAsync(); + + Assert.False(await pendingMoveNext); + await using IAsyncEnumerator activeEnumerator = + stream.GetAsyncEnumerator(TestContext.Current.CancellationToken); + Assert.True(stream.TryPublish(7)); + Assert.True(await activeEnumerator.MoveNextAsync()); + Assert.Equal(7, activeEnumerator.Current); + } + [Fact] public async Task DisposeClosesAdmissionDiscardsBufferedValuesAndCompletesReaders() { @@ -170,8 +190,81 @@ public async Task DisposeClosesAdmissionDiscardsBufferedValuesAndCompletesReader Assert.True(stream.IsCompleted); Assert.False(stream.IsAdmissionOpen); Assert.False(stream.TryPublish(2)); + Assert.Equal(1, stream.LostCount); + + stream.Dispose(); + } + + [Fact] + public async Task DisposeCompletesAWaitingReaderWithoutRetainingIt() + { + BoundedEventStream stream = new(new EventStreamOptions(1)); + await using IAsyncEnumerator enumerator = stream.GetAsyncEnumerator(TestContext.Current.CancellationToken); + Task pendingMoveNext = enumerator.MoveNextAsync().AsTask(); stream.Dispose(); + + Assert.False(await pendingMoveNext); + Assert.False(stream.TryPublish(1)); + } + + [Fact] + public async Task CloseAdmissionRejectsNewObservationsAndLetsAcceptedObservationsDrain() + { + using BoundedEventStream stream = new(new EventStreamOptions(1)); + Assert.True(stream.TryPublish(1)); + stream.CloseAdmission(); + stream.Complete(); + + Assert.False(stream.TryPublish(2)); + await using IAsyncEnumerator enumerator = stream.GetAsyncEnumerator(TestContext.Current.CancellationToken); + Assert.True(await enumerator.MoveNextAsync()); + Assert.Equal(1, enumerator.Current); + Assert.False(await enumerator.MoveNextAsync()); + Assert.Equal(0, stream.LostCount); + } + + [Fact] + public async Task FaultedCompletionDiscardsBufferedEventsReportsLossAndFailsReaders() + { + using BoundedEventStream stream = new(new EventStreamOptions(2)); + InvalidOperationException expected = new("callback failure"); + Assert.True(stream.TryPublish(1)); + stream.Complete(expected); + + await using IAsyncEnumerator enumerator = stream.GetAsyncEnumerator(TestContext.Current.CancellationToken); + InvalidOperationException actual = await Assert.ThrowsAsync(async () => + { + _ = await enumerator.MoveNextAsync(); + }); + + Assert.Same(expected, actual); + Assert.Equal(1, stream.LostCount); + Assert.False(stream.TryPublish(2)); + } + + [Fact] + public async Task PublicationDoesNotWaitForASlowConsumerContinuation() + { + CancellationToken cancellationToken = TestContext.Current.CancellationToken; + using BoundedEventStream stream = new(new EventStreamOptions(1)); + await using IAsyncEnumerator enumerator = stream.GetAsyncEnumerator(cancellationToken); + using ManualResetEventSlim consumerMayContinue = new(false); + TaskCompletionSource consumerEntered = new(TaskCreationOptions.RunContinuationsAsynchronously); + + Task consumer = Task.Run(async () => + { + Assert.True(await enumerator.MoveNextAsync()); + consumerEntered.SetResult(); + consumerMayContinue.Wait(cancellationToken); + }, cancellationToken); + + Task publish = Task.Run(() => stream.TryPublish(1)); + await consumerEntered.Task.WaitAsync(cancellationToken); + Assert.True(publish.IsCompletedSuccessfully); + + consumerMayContinue.Set(); + await consumer; } [Fact] diff --git a/tests/CheatEngine.Client.Core.Tests/Domains/Events/EventStreamLeaseTests.cs b/tests/CheatEngine.Client.Core.Tests/Domains/Events/EventStreamLeaseTests.cs index 780d6dc..a8aa61f 100644 --- a/tests/CheatEngine.Client.Core.Tests/Domains/Events/EventStreamLeaseTests.cs +++ b/tests/CheatEngine.Client.Core.Tests/Domains/Events/EventStreamLeaseTests.cs @@ -30,6 +30,24 @@ public async Task DisposeStopsAdmissionBeforeNeutralizingThenCompletesAndRelease Assert.False(await enumerator.MoveNextAsync()); } + [Fact] + public async Task DisposeCompletesAWaitingReaderAndClosesAdmissionBeforeReleasingTheHostRegistration() + { + BoundedEventStream stream = new(new EventStreamOptions(1)); + await using IAsyncEnumerator enumerator = + stream.GetAsyncEnumerator(TestContext.Current.CancellationToken); + Task pendingMoveNext = enumerator.MoveNextAsync().AsTask(); + EventStreamLease lease = new(stream, static () => + { + }, () => Assert.False(stream.IsAdmissionOpen)); + + lease.Dispose(); + + Assert.False(await pendingMoveNext); + Assert.True(lease.IsReleased); + Assert.False(lease.TryPublish(1)); + } + [Fact] public void DisposeIsIdempotentAndDoesNotReleaseTheHostTwice() {