diff --git a/src/Orleans.Runtime/Catalog/ActivationData.cs b/src/Orleans.Runtime/Catalog/ActivationData.cs index 0ce9b85710e..b7cb7681212 100644 --- a/src/Orleans.Runtime/Catalog/ActivationData.cs +++ b/src/Orleans.Runtime/Catalog/ActivationData.cs @@ -32,7 +32,7 @@ internal sealed partial class ActivationData : IGrainContext, ICollectibleGrainContext, IGrainExtensionBinder, - IActivationWorkingSetMember, + IActivationWorkingSetMemberStatus, IGrainTimerRegistry, IGrainManagementExtension, IGrainCallCancellationExtension, @@ -41,6 +41,11 @@ internal sealed partial class ActivationData : IDisposable { private const string GrainAddressMigrationContextKey = "sys.addr"; + // Activation lifecycle and working-set CLOCK state share one byte. All writes occur while holding the activation lock. + private const byte ActivationStateMask = 0b0000_0111; + private const byte IsInWorkingSetMask = 0b0000_1000; + private const byte IsIdleInWorkingSetMask = 0b0001_0000; + private const byte WasRemovedByCollectionMask = 0b0010_0000; private readonly GrainTypeSharedContext _shared; private readonly IServiceScope _serviceScope; private readonly WorkItemGroup _workItemGroup; @@ -50,7 +55,7 @@ internal sealed partial class ActivationData : private GrainLifecycle? _lifecycle; private Queue? _pendingOperations; private Message? _blockingRequest; - private bool _isInWorkingSet = true; + private byte _status = IsInWorkingSetMask; private CoarseStopwatch _busyDuration; private CoarseStopwatch _idleDuration; private GrainReference? _selfReference; @@ -146,7 +151,7 @@ public void Start(IGrainActivator grainActivator) public object? GrainInstance { get; private set; } public GrainAddress Address { get; private set; } public GrainReference GrainReference => _selfReference ??= _shared.GrainReferenceActivator.CreateReference(GrainId, default); - public ActivationState State { get; private set; } = ActivationState.Creating; + public ActivationState State => (ActivationState)(Volatile.Read(ref _status) & ActivationStateMask); public PlacementStrategy PlacementStrategy => _shared.PlacementStrategy; public DateTime CollectionTicket { get; set; } public IServiceProvider ActivationServices => _serviceScope.ServiceProvider; @@ -167,6 +172,24 @@ public IGrainLifecycle ObservableLifecycle public DateTime KeepAliveUntil { get; set; } = DateTime.MinValue; public bool IsValid => State is ActivationState.Valid; + private bool IsInWorkingSet + { + get => (Volatile.Read(ref _status) & IsInWorkingSetMask) != 0; + set => SetStatusFlag(IsInWorkingSetMask, value); + } + + private bool IsIdleInWorkingSet + { + get => (Volatile.Read(ref _status) & IsIdleInWorkingSetMask) != 0; + set => SetStatusFlag(IsIdleInWorkingSetMask, value); + } + + private bool WasRemovedByCollection + { + get => (Volatile.Read(ref _status) & WasRemovedByCollectionMask) != 0; + set => SetStatusFlag(WasRemovedByCollectionMask, value); + } + // Currently, the only supported multi-activation grain is one using the StatelessWorkerPlacement strategy. internal bool IsStatelessWorker => PlacementStrategy is StatelessWorkerPlacement; @@ -404,7 +427,14 @@ internal void SetGrainInstance(object grainInstance) public void SetState(ActivationState state) { - State = state; + Debug.Assert(Monitor.IsEntered(this)); + _status = (byte)((_status & ~ActivationStateMask) | (byte)state); + } + + private void SetStatusFlag(byte mask, bool value) + { + Debug.Assert(Monitor.IsEntered(this)); + _status = value ? (byte)(_status | mask) : (byte)(_status & ~mask); } /// @@ -966,13 +996,34 @@ public TExtensionInterface GetExtension() bool IActivationWorkingSetMember.IsCandidateForRemoval(bool wouldRemove) { const int IdlenessLowerBound = 10_000; - lock (this) + Debug.Assert(Monitor.IsEntered(this)); + var inactive = IsInactive && _idleDuration.ElapsedMilliseconds > IdlenessLowerBound; + + WasRemovedByCollection = wouldRemove && inactive; + return inactive; + } + + bool IActivationWorkingSetMemberStatus.IsInWorkingSet + { + get => IsInWorkingSet; + set { - var inactive = IsInactive && _idleDuration.ElapsedMilliseconds > IdlenessLowerBound; + Debug.Assert(Monitor.IsEntered(this)); + IsInWorkingSet = value; + if (value) + { + WasRemovedByCollection = false; + } + } + } - // This instance will remain in the working set if it is either not pending removal or if it is currently active. - _isInWorkingSet = !wouldRemove || !inactive; - return inactive; + bool IActivationWorkingSetMemberStatus.IsIdle + { + get => IsIdleInWorkingSet; + set + { + Debug.Assert(Monitor.IsEntered(this)); + IsIdleInWorkingSet = value; } } @@ -1482,13 +1533,14 @@ private void OnCompletedRequest(Message message) // If the message is meant to keep the activation active, reset the idle timer and ensure the activation // is in the activation working set. - if (message.IsKeepAlive) + if (message.IsKeepAlive && State is ActivationState.Valid) { _idleDuration = CoarseStopwatch.StartNew(); + IsIdleInWorkingSet = false; - if (!_isInWorkingSet) + if (!IsInWorkingSet) { - _isInWorkingSet = true; + IsInWorkingSet = true; _shared.InternalRuntime.ActivationWorkingSet.OnActive(this); } } @@ -2034,7 +2086,7 @@ private async Task FinishDeactivating(Command.Deactivate deactivateCommand, Canc deactivationMetrics = deactivationMetrics.Migration(); _shared.CatalogInstruments.ActivationShutdownViaMigration(); } - else if (_isInWorkingSet) + else if (!WasRemovedByCollection) { deactivationMetrics = deactivationMetrics.DeactivateOnIdle(); _shared.CatalogInstruments.ActivationShutdownViaDeactivateOnIdle(); diff --git a/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs b/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs index bbba7539b3a..0e251c5e182 100644 --- a/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs +++ b/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs @@ -10,276 +10,423 @@ using Orleans.Internal; using Orleans.Runtime.Internal; -namespace Orleans.Runtime +namespace Orleans.Runtime; + +/// +/// Maintains a list of activations which are recently active. +/// +internal sealed partial class ActivationWorkingSet : IActivationWorkingSet, ILifecycleParticipant { - /// - /// Maintains a list of activations which are recently active. - /// - internal sealed partial class ActivationWorkingSet : IActivationWorkingSet, ILifecycleParticipant + private const byte IsIdleMask = 0b0000_0001; + private readonly ConcurrentDictionary _members = new(); + private readonly ILogger _logger; + private readonly IAsyncTimer _scanPeriodTimer; + private readonly List _observers; + + private int _activeCount; + private Task? _runTask; + + public ActivationWorkingSet( + IAsyncTimerFactory asyncTimerFactory, + ILogger logger, + IEnumerable observers, + CatalogInstruments catalogInstruments, + [FromKeyedServices(TimeProviderNames.SystemTimers)] TimeProvider timeProvider) { - private class MemberState - { - public bool IsIdle { get; set; } - } + _logger = logger; + _scanPeriodTimer = asyncTimerFactory.Create(TimeSpan.FromMilliseconds(5_000), nameof(ActivationWorkingSet) + "." + nameof(MonitorWorkingSet), timeProvider); + _observers = observers.ToList(); + catalogInstruments.RegisterActivationWorkingSetObserve(() => Count); + } - private readonly ConcurrentDictionary _members = new(); - private readonly ILogger _logger; - private readonly IAsyncTimer _scanPeriodTimer; - private readonly List _observers; + public int Count => _activeCount; - private int _activeCount; - private Task? _runTask; + internal IEnumerable Members => EnumerateActiveMembers(); - public ActivationWorkingSet( - IAsyncTimerFactory asyncTimerFactory, - ILogger logger, - IEnumerable observers, - CatalogInstruments catalogInstruments, - [FromKeyedServices(TimeProviderNames.SystemTimers)] TimeProvider timeProvider) + private IEnumerable EnumerateActiveMembers() + { + foreach (var pair in _members) { - _logger = logger; - _scanPeriodTimer = asyncTimerFactory.Create(TimeSpan.FromMilliseconds(5_000), nameof(ActivationWorkingSet) + "." + nameof(MonitorWorkingSet), timeProvider); - _observers = observers.ToList(); - catalogInstruments.RegisterActivationWorkingSetObserve(() => Count); + if (pair.Key is IActivationWorkingSetMemberStatus status + ? status.IsInWorkingSet && !status.IsIdle + : (pair.Value & IsIdleMask) == 0) + { + yield return pair.Key; + } } + } - public int Count => _activeCount; + public void OnActivated(IActivationWorkingSetMember member) + { + Debug.Assert(member is not ICollectibleGrainContext collectible || collectible.IsValid); + if (Monitor.IsEntered(member)) + { + AddMember(); + } + else + { + lock (member) + { + AddMember(); + } + } - internal IEnumerable Members => EnumerateActiveMembers(); + foreach (var observer in _observers) + { + observer.OnAdded(member); + } - private IEnumerable EnumerateActiveMembers() + void AddMember() { - foreach (var pair in _members) + Debug.Assert(Monitor.IsEntered(member)); + if (!_members.TryAdd(member, 0)) { - if (!pair.Value.IsIdle) - { - yield return pair.Key; - } + throw new InvalidOperationException($"Member {member} is already a member of the working set"); + } + + Interlocked.Increment(ref _activeCount); + if (member is IActivationWorkingSetMemberStatus status) + { + status.IsInWorkingSet = true; + status.IsIdle = false; } } + } - public void OnActivated(IActivationWorkingSetMember member) + public void OnActive(IActivationWorkingSetMember member) + { + if (Monitor.IsEntered(member)) + { + MarkActive(); + } + else { - Debug.Assert(member is not ICollectibleGrainContext collectible || collectible.IsValid); - if (_members.TryAdd(member, new MemberState())) + lock (member) { - Interlocked.Increment(ref _activeCount); - foreach (var observer in _observers) - { - observer.OnAdded(member); - } - - return; + MarkActive(); } + } - throw new InvalidOperationException($"Member {member} is already a member of the working set"); + foreach (var observer in _observers) + { + observer.OnActive(member); } - public void OnActive(IActivationWorkingSetMember member) + void MarkActive() { - if (_members.TryGetValue(member, out var state)) + Debug.Assert(Monitor.IsEntered(member)); + var added = _members.TryAdd(member, 0); + if (added) { - state.IsIdle = false; + Interlocked.Increment(ref _activeCount); } - else if (_members.TryAdd(member, new())) + + if (member is IActivationWorkingSetMemberStatus status) { - Interlocked.Increment(ref _activeCount); + status.IsInWorkingSet = true; + status.IsIdle = false; } + else + { + _members[member] = 0; + } + } + } - foreach (var observer in _observers) + public void OnEvicted(IActivationWorkingSetMember member) + { + bool removed; + if (Monitor.IsEntered(member)) + { + removed = RemoveMember(); + } + else + { + lock (member) { - observer.OnActive(member); + removed = RemoveMember(); } } - public void OnEvicted(IActivationWorkingSetMember member) + if (removed) { - if (_members.TryRemove(member, out _)) + OnEvictedCore(member); + } + + bool RemoveMember() + { + Debug.Assert(Monitor.IsEntered(member)); + var result = _members.TryRemove(member, out _); + if (result) { Interlocked.Decrement(ref _activeCount); - foreach (var observer in _observers) + if (member is IActivationWorkingSetMemberStatus status) { - observer.OnEvicted(member); + status.IsInWorkingSet = false; + status.IsIdle = false; } } + + return result; } + } - public void OnDeactivating(IActivationWorkingSetMember member) + private void OnEvictedCore(IActivationWorkingSetMember member) + { + foreach (var observer in _observers) { - OnEvicted(member); - foreach (var observer in _observers) - { - observer.OnDeactivating(member); - } + observer.OnEvicted(member); + } + } + + public void OnDeactivating(IActivationWorkingSetMember member) + { + OnEvicted(member); + foreach (var observer in _observers) + { + observer.OnDeactivating(member); } + } - public void OnDeactivated(IActivationWorkingSetMember member) + public void OnDeactivated(IActivationWorkingSetMember member) + { + OnEvicted(member); + foreach (var observer in _observers) { - OnEvicted(member); - foreach (var observer in _observers) - { - observer.OnDeactivated(member); - } + observer.OnDeactivated(member); } + } - private async Task MonitorWorkingSet() + private async Task MonitorWorkingSet() + { + while (await _scanPeriodTimer.NextTick()) { - while (await _scanPeriodTimer.NextTick()) + foreach (var pair in _members) { - foreach (var pair in _members) + try { - try - { - VisitMember(pair.Key, pair.Value); - } - catch (Exception exception) - { - LogExceptionVisitingWorkingSetMember(exception, pair.Key); - } + VisitMember(pair.Key); + } + catch (Exception exception) + { + LogExceptionVisitingWorkingSetMember(exception, pair.Key); } } } + } - private void VisitMember(IActivationWorkingSetMember member, MemberState state) + private void VisitMember(IActivationWorkingSetMember member) + { + MemberVisitResult result; + // Enumeration can retain a member across removal and re-addition. CLOCK state is advisory, so visit the + // member's current state while holding its lock instead of adding a dictionary validation to every scan. + lock (member) { - var wouldRemove = state.IsIdle; - if (member.IsCandidateForRemoval(wouldRemove)) + var status = member as IActivationWorkingSetMemberStatus; + byte dictionaryState = 0; + if ((status is null && !_members.TryGetValue(member, out dictionaryState)) + || (status is not null && !status.IsInWorkingSet)) { - if (wouldRemove) + result = MemberVisitResult.None; + } + else + { + var wouldRemove = status is not null + ? status.IsIdle + : (dictionaryState & IsIdleMask) != 0; + if (member.IsCandidateForRemoval(wouldRemove)) { - OnEvicted(member); + if (wouldRemove) + { + if (_members.TryRemove(member, out _)) + { + if (status is not null) + { + status.IsInWorkingSet = false; + status.IsIdle = false; + } + + Interlocked.Decrement(ref _activeCount); + result = MemberVisitResult.Evicted; + } + else + { + result = MemberVisitResult.None; + } + } + else + { + if (status is not null) + { + status.IsIdle = true; + } + else + { + _members[member] = IsIdleMask; + } + + result = MemberVisitResult.Idle; + } } else { - state.IsIdle = true; - foreach (var observer in _observers) + if (wouldRemove && status is not null) { - observer.OnIdle(member); + status.IsIdle = false; } + else if (wouldRemove) + { + _members[member] = 0; + } + + result = MemberVisitResult.Active; } } - else + } + + foreach (var observer in _observers) + { + switch (result) { - state.IsIdle = false; - foreach (var observer in _observers) - { + case MemberVisitResult.Active: observer.OnActive(member); - } + break; + case MemberVisitResult.Idle: + observer.OnIdle(member); + break; + case MemberVisitResult.Evicted: + observer.OnEvicted(member); + break; } } + } - void ILifecycleParticipant.Participate(ISiloLifecycle lifecycle) - { - lifecycle.Subscribe( - nameof(ActivationWorkingSet), - ServiceLifecycleStage.BecomeActive, - StartMonitoring, - StopMonitoring); + private enum MemberVisitResult + { + None, + Active, + Idle, + Evicted + } - Task StartMonitoring(CancellationToken ct) - { - using var _ = new ExecutionContextSuppressor(); - _runTask = Task.Run(MonitorWorkingSet); - return Task.CompletedTask; - } + void ILifecycleParticipant.Participate(ISiloLifecycle lifecycle) + { + lifecycle.Subscribe( + nameof(ActivationWorkingSet), + ServiceLifecycleStage.BecomeActive, + StartMonitoring, + StopMonitoring); - async Task StopMonitoring(CancellationToken ct) + Task StartMonitoring(CancellationToken ct) + { + using var _ = new ExecutionContextSuppressor(); + _runTask = Task.Run(MonitorWorkingSet); + return Task.CompletedTask; + } + + async Task StopMonitoring(CancellationToken ct) + { + _scanPeriodTimer.Dispose(); + if (_runTask is Task task) { - _scanPeriodTimer.Dispose(); - if (_runTask is Task task) - { - await task.WaitAsync(ct).SuppressThrowing(); - } + await task.WaitAsync(ct).SuppressThrowing(); } } - - [LoggerMessage( - Level = LogLevel.Error, - Message = "Exception visiting working set member {Member}" - )] - private partial void LogExceptionVisitingWorkingSetMember(Exception exception, IActivationWorkingSetMember member); } + [LoggerMessage( + Level = LogLevel.Error, + Message = "Exception visiting working set member {Member}" + )] + private partial void LogExceptionVisitingWorkingSetMember(Exception exception, IActivationWorkingSetMember member); +} + +/// +/// Manages the set of recently active instances. +/// +public interface IActivationWorkingSet +{ /// - /// Manages the set of recently active instances. + /// Returns the number of grain activations which were recently active. /// - public interface IActivationWorkingSet - { - /// - /// Returns the number of grain activations which were recently active. - /// - public int Count { get; } - - /// - /// Adds a new member to the working set. - /// - void OnActivated(IActivationWorkingSetMember member); - - /// - /// Signals that a member is active and should be in the working set. - /// - void OnActive(IActivationWorkingSetMember member); - - /// - /// Signals that a member has begun to deactivate. - /// - /// - void OnDeactivating(IActivationWorkingSetMember member); - - /// - /// Signals that a members has deactivated. - /// - void OnDeactivated(IActivationWorkingSetMember member); - } + public int Count { get; } /// - /// Represents an activation from the perspective of . + /// Adds a new member to the working set. /// - public interface IActivationWorkingSetMember - { - /// - /// Returns if the member is eligible for removal, otherwise. - /// - /// if the member is eligible for removal, otherwise. - /// - /// If this method returns and is , the member must be removed from the working set and is eligible to be added again via a call to . - /// - bool IsCandidateForRemoval(bool wouldRemove); - } + void OnActivated(IActivationWorkingSetMember member); /// - /// An observer. + /// Signals that a member is active and should be in the working set. /// - public interface IActivationWorkingSetObserver - { - /// - /// Called when an activation is added to the working set. - /// - void OnAdded(IActivationWorkingSetMember member) { } - - /// - /// Called when an activation becomes active. - /// - void OnActive(IActivationWorkingSetMember member) { } - - /// - /// Called when an activation becomes idle. - /// - void OnIdle(IActivationWorkingSetMember member) { } - - /// - /// Called when an activation is removed from the working set. - /// - void OnEvicted(IActivationWorkingSetMember member) { } - - /// - /// Called when an activation starts deactivating. - /// - void OnDeactivating(IActivationWorkingSetMember member) { } - - /// - /// Called when an activation is deactivated. - /// - void OnDeactivated(IActivationWorkingSetMember member) { } - } + void OnActive(IActivationWorkingSetMember member); + + /// + /// Signals that a member has begun to deactivate. + /// + /// + void OnDeactivating(IActivationWorkingSetMember member); + + /// + /// Signals that a members has deactivated. + /// + void OnDeactivated(IActivationWorkingSetMember member); +} + +/// +/// Represents an activation from the perspective of . +/// +public interface IActivationWorkingSetMember +{ + /// + /// Returns if the member is eligible for removal, otherwise. + /// + /// if the member is eligible for removal, otherwise. + /// + /// If this method returns and is , the member must be removed from the working set and is eligible to be added again via a call to . + /// + bool IsCandidateForRemoval(bool wouldRemove); +} + +internal interface IActivationWorkingSetMemberStatus : IActivationWorkingSetMember +{ + bool IsInWorkingSet { get; set; } + + bool IsIdle { get; set; } +} + +/// +/// An observer. +/// +public interface IActivationWorkingSetObserver +{ + /// + /// Called when an activation is added to the working set. + /// + void OnAdded(IActivationWorkingSetMember member) { } + + /// + /// Called when an activation becomes active. + /// + void OnActive(IActivationWorkingSetMember member) { } + + /// + /// Called when an activation becomes idle. + /// + void OnIdle(IActivationWorkingSetMember member) { } + + /// + /// Called when an activation is removed from the working set. + /// + void OnEvicted(IActivationWorkingSetMember member) { } + + /// + /// Called when an activation starts deactivating. + /// + void OnDeactivating(IActivationWorkingSetMember member) { } + + /// + /// Called when an activation is deactivated. + /// + void OnDeactivated(IActivationWorkingSetMember member) { } } diff --git a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs index 24180c647f5..3448d4f62e1 100644 --- a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs +++ b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs @@ -1,4 +1,7 @@ using System.Collections.Concurrent; +using System.Diagnostics.CodeAnalysis; +using System.Reflection; +using System.Runtime.CompilerServices; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging.Abstractions; using Microsoft.Extensions.Options; @@ -108,7 +111,7 @@ public void TryRescheduleCollection_DoesNotThrow_WhenCollectionTicketIsMaxValue( [InlineData(80.0, 70.0, 1000, 250, 100, false, 0)] // Below threshold, no deactivation [InlineData(80.0, 70.0, 1000, 100, 200, true, 155)] // More activations, smaller per-activation size [InlineData(80.0, 70.0, 1000, 800, 100, false, 0)] // Well below threshold - [InlineData(80.0, 70.0, 1000, 50, 10, true, 7)] // Few activations, large per-activation size + [InlineData(80.0, 70.0, 1000, 50, 10, true, 7)] // Few activations, large per-activation size [InlineData(80.0, 70.0, 1000, 100, 0, false, 0)] // No activations public void IsMemoryOverloaded_WorksAsExpected( double memoryLoadThreshold, @@ -380,6 +383,391 @@ public async Task DeactivateInDueTimeOrder_SkipsActiveAndInvalidActivations() Assert.Equal(2, collector._activationCount); } + [Fact, TestCategory("Activation")] + public async Task WorkingSetScan_DoesNotUpdateReaddedMember() + { + var timer = Substitute.For(); + timer.NextTick().Returns(Task.FromResult(true), Task.FromResult(false)); + var timerFactory = Substitute.For(); + timerFactory.Create(Arg.Any(), Arg.Any(), Arg.Any()).Returns(timer); + var observer = Substitute.For(); + var workingSet = new ActivationWorkingSet( + timerFactory, + NullLogger.Instance, + [observer], + CreateCatalogInstruments(), + TimeProvider.System); + var scanStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + using var resumeScan = new ManualResetEventSlim(); + var member = new TestWorkingSetMember(wouldRemove => + { + Assert.False(wouldRemove); + scanStarted.TrySetResult(); + resumeScan.Wait(); + return true; + }); + workingSet.OnActivated(member); + + var lifecycle = new SiloLifecycleSubject(NullLogger.Instance); + ((ILifecycleParticipant)workingSet).Participate(lifecycle); + await lifecycle.OnStart(TestContext.Current.CancellationToken); + await scanStarted.Task.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); + + var mutationStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var mutationTask = Task.Run(() => + { + mutationStarted.SetResult(); + workingSet.OnEvicted(member); + workingSet.OnActivated(member); + }, TestContext.Current.CancellationToken); + await mutationStarted.Task.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); + try + { + var stopTask = lifecycle.OnStop(TestContext.Current.CancellationToken); + Assert.False(stopTask.IsCompleted); + Assert.False(mutationTask.IsCompleted); + resumeScan.Set(); + await mutationTask; + await stopTask; + } + finally + { + resumeScan.Set(); + } + + Assert.Equal(1, workingSet.Count); + Assert.Contains(member, workingSet.Members); + observer.Received(2).OnAdded(member); + observer.Received(1).OnEvicted(member); + observer.Received(1).OnIdle(member); + } + + [Fact, TestCategory("Activation")] + public async Task WorkingSetScan_DoesNotRemoveReaddedMember() + { + var timer = Substitute.For(); + timer.NextTick().Returns(Task.FromResult(true), Task.FromResult(true), Task.FromResult(false)); + var timerFactory = Substitute.For(); + timerFactory.Create(Arg.Any(), Arg.Any(), Arg.Any()).Returns(timer); + var observer = Substitute.For(); + var workingSet = new ActivationWorkingSet( + timerFactory, + NullLogger.Instance, + [observer], + CreateCatalogInstruments(), + TimeProvider.System); + var removalScanStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + using var resumeScan = new ManualResetEventSlim(); + var member = new TestWorkingSetMember(wouldRemove => + { + if (!wouldRemove) + { + return true; + } + + removalScanStarted.TrySetResult(); + resumeScan.Wait(); + return true; + }); + workingSet.OnActivated(member); + + var lifecycle = new SiloLifecycleSubject(NullLogger.Instance); + ((ILifecycleParticipant)workingSet).Participate(lifecycle); + await lifecycle.OnStart(TestContext.Current.CancellationToken); + await removalScanStarted.Task.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); + + var mutationStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var mutationTask = Task.Run(() => + { + mutationStarted.SetResult(); + workingSet.OnEvicted(member); + workingSet.OnActivated(member); + }, TestContext.Current.CancellationToken); + await mutationStarted.Task.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); + try + { + var stopTask = lifecycle.OnStop(TestContext.Current.CancellationToken); + Assert.False(stopTask.IsCompleted); + Assert.False(mutationTask.IsCompleted); + resumeScan.Set(); + await mutationTask; + await stopTask; + } + finally + { + resumeScan.Set(); + } + + Assert.Equal(1, workingSet.Count); + Assert.Contains(member, workingSet.Members); + observer.Received(2).OnAdded(member); + observer.Received(1).OnIdle(member); + observer.Received(1).OnEvicted(member); + } + + [Fact, TestCategory("Activation")] + public async Task WorkingSetScan_RepeatedCyclesPreserveCountAndObserverConsistency() + { + const int reactivationCycles = 64; + const int removalVisits = 2; + const int totalVisits = reactivationCycles * 2 + removalVisits; + var ticks = 0; + var timer = Substitute.For(); + timer.NextTick().Returns(_ => Task.FromResult(Interlocked.Increment(ref ticks) <= totalVisits)); + var timerFactory = Substitute.For(); + timerFactory.Create(Arg.Any(), Arg.Any(), Arg.Any()).Returns(timer); + var observer = Substitute.For(); + var workingSet = new ActivationWorkingSet( + timerFactory, + NullLogger.Instance, + [observer], + CreateCatalogInstruments(), + TimeProvider.System); + var visits = 0; + var member = new TestWorkingSetMember(_ => + { + var visit = visits++; + return visit >= reactivationCycles * 2 || (visit & 1) == 0; + }); + workingSet.OnActivated(member); + + var lifecycle = new SiloLifecycleSubject(NullLogger.Instance); + ((ILifecycleParticipant)workingSet).Participate(lifecycle); + await lifecycle.OnStart(TestContext.Current.CancellationToken); + await lifecycle.OnStop(TestContext.Current.CancellationToken); + + Assert.Equal(0, workingSet.Count); + Assert.DoesNotContain(member, workingSet.Members); + Assert.Equal(totalVisits, visits); + observer.Received(1).OnAdded(member); + observer.Received(reactivationCycles + 1).OnIdle(member); + observer.Received(reactivationCycles).OnActive(member); + observer.Received(1).OnEvicted(member); + } + + [Fact, TestCategory("Activation")] + public async Task WorkingSetMembers_EnumerationToleratesConcurrentRemoveAndReadd() + { + var timer = Substitute.For(); + timer.NextTick().Returns(Task.FromResult(false)); + var timerFactory = Substitute.For(); + timerFactory.Create(Arg.Any(), Arg.Any(), Arg.Any()).Returns(timer); + var observer = Substitute.For(); + var firstEviction = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + using var resumeWriter = new ManualResetEventSlim(); + var evictionCount = 0; + observer.When(static observer => observer.OnEvicted(Arg.Any())).Do(_ => + { + if (Interlocked.Increment(ref evictionCount) == 1) + { + firstEviction.TrySetResult(); + resumeWriter.Wait(); + } + }); + var workingSet = new ActivationWorkingSet( + timerFactory, + NullLogger.Instance, + [observer], + CreateCatalogInstruments(), + TimeProvider.System); + var members = Enumerable.Range(0, 128).Select(_ => new TestWorkingSetMember()).ToArray(); + foreach (var member in members) + { + workingSet.OnActivated(member); + } + + using var enumerator = workingSet.Members.GetEnumerator(); + Assert.True(enumerator.MoveNext()); + var enumeratedMembers = new List { enumerator.Current }; + var writer = Task.Run(() => + { + foreach (var member in members) + { + workingSet.OnEvicted(member); + workingSet.OnActivated(member); + } + }, TestContext.Current.CancellationToken); + await firstEviction.Task.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); + + try + { + while (enumerator.MoveNext()) + { + enumeratedMembers.Add(enumerator.Current); + } + } + finally + { + resumeWriter.Set(); + } + + await writer.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); + + Assert.All(enumeratedMembers, member => Assert.Contains(member, members)); + Assert.Equal(members.Length, workingSet.Count); + Assert.True(members.Cast().ToHashSet().SetEquals(workingSet.Members)); + observer.Received(members.Length * 2).OnAdded(Arg.Any()); + observer.Received(members.Length).OnEvicted(Arg.Any()); + } + + [Fact, TestCategory("Activation")] + public void WorkingSetMembers_ExcludesUnregisteredMember() + { + var timer = Substitute.For(); + timer.NextTick().Returns(Task.FromResult(false)); + var timerFactory = Substitute.For(); + timerFactory.Create(Arg.Any(), Arg.Any(), Arg.Any()).Returns(timer); + var workingSet = new ActivationWorkingSet( + timerFactory, + NullLogger.Instance, + Array.Empty(), + CreateCatalogInstruments(), + TimeProvider.System); + var member = new TestWorkingSetMember(); + workingSet.OnActivated(member); + lock (member) + { + member.IsInWorkingSet = false; + } + + Assert.Empty(workingSet.Members); + } + + [Fact, TestCategory("Activation")] + public async Task WorkingSetScan_SerializesRemovalWithReactivation() + { + var timer = Substitute.For(); + timer.NextTick().Returns(Task.FromResult(true), Task.FromResult(true), Task.FromResult(false)); + var timerFactory = Substitute.For(); + timerFactory.Create(Arg.Any(), Arg.Any(), Arg.Any()).Returns(timer); + var observer = Substitute.For(); + var workingSet = new ActivationWorkingSet( + timerFactory, + NullLogger.Instance, + [observer], + CreateCatalogInstruments(), + TimeProvider.System); + var removalScanStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + using var resumeScan = new ManualResetEventSlim(); + var member = new TestWorkingSetMember(wouldRemove => + { + if (!wouldRemove) + { + return true; + } + + removalScanStarted.TrySetResult(); + resumeScan.Wait(); + return true; + }); + workingSet.OnActivated(member); + + var lifecycle = new SiloLifecycleSubject(NullLogger.Instance); + ((ILifecycleParticipant)workingSet).Participate(lifecycle); + await lifecycle.OnStart(TestContext.Current.CancellationToken); + await removalScanStarted.Task.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); + + var reactivationStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var reactivationTask = Task.Run(() => + { + reactivationStarted.SetResult(); + workingSet.OnActive(member); + }, TestContext.Current.CancellationToken); + await reactivationStarted.Task.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); + try + { + Assert.False(reactivationTask.IsCompleted); + resumeScan.Set(); + await reactivationTask; + await lifecycle.OnStop(TestContext.Current.CancellationToken); + } + finally + { + resumeScan.Set(); + } + + Assert.Equal(1, workingSet.Count); + Assert.Contains(member, workingSet.Members); + observer.Received(1).OnAdded(member); + observer.Received(1).OnIdle(member); + observer.Received(1).OnEvicted(member); + observer.Received(1).OnActive(member); + } + + [Fact, TestCategory("Activation")] + public void WorkingSetScan_SkipsRemovedMember() + { + var timer = Substitute.For(); + var timerFactory = Substitute.For(); + timerFactory.Create(Arg.Any(), Arg.Any(), Arg.Any()).Returns(timer); + var observer = Substitute.For(); + var workingSet = new ActivationWorkingSet( + timerFactory, + NullLogger.Instance, + [observer], + CreateCatalogInstruments(), + TimeProvider.System); + var visits = 0; + var member = new TestWorkingSetMember(_ => + { + visits++; + return true; + }); + workingSet.OnActivated(member); + workingSet.OnEvicted(member); + + var visitMember = typeof(ActivationWorkingSet).GetMethod( + "VisitMember", + BindingFlags.Instance | BindingFlags.NonPublic, + [typeof(IActivationWorkingSetMember)]) + ?? throw new InvalidOperationException("Could not find the working-set scan method."); + visitMember.Invoke(workingSet, [member]); + + Assert.Equal(0, visits); + Assert.Equal(0, workingSet.Count); + observer.Received(1).OnAdded(member); + observer.Received(1).OnEvicted(member); + observer.DidNotReceive().OnIdle(member); + observer.DidNotReceive().OnActive(member); + } + + [Fact, TestCategory("Activation")] + public void WorkingSet_PublicMemberUsesDictionaryBackedClockState() + { + var timer = Substitute.For(); + var timerFactory = Substitute.For(); + timerFactory.Create(Arg.Any(), Arg.Any(), Arg.Any()).Returns(timer); + var observer = Substitute.For(); + var workingSet = new ActivationWorkingSet( + timerFactory, + NullLogger.Instance, + [observer], + CreateCatalogInstruments(), + TimeProvider.System); + var member = new PublicWorkingSetMember(); + var visitMember = typeof(ActivationWorkingSet).GetMethod( + "VisitMember", + BindingFlags.Instance | BindingFlags.NonPublic, + [typeof(IActivationWorkingSetMember)]) + ?? throw new InvalidOperationException("Could not find the working-set scan method."); + + workingSet.OnActivated(member); + visitMember.Invoke(workingSet, [member]); + visitMember.Invoke(workingSet, [member]); + + Assert.Equal([false, true], member.CandidateCalls); + Assert.Equal(0, workingSet.Count); + Assert.Empty(workingSet.Members); + observer.Received(1).OnIdle(member); + observer.Received(1).OnEvicted(member); + + workingSet.OnActive(member); + + Assert.Equal(1, workingSet.Count); + Assert.Equal([member], workingSet.Members); + observer.Received(1).OnActive(member); + } + private IActivationWorkingSetMember PrepareActivation(int collectionAgeLimitMinutes, ActivationCollector collector) => PrepareActivation(TimeSpan.FromMinutes(collectionAgeLimitMinutes), collector); @@ -403,5 +791,1101 @@ private IActivationWorkingSetMember PrepareActivation(TimeSpan collectionAgeLimi return (IActivationWorkingSetMember)activation; } + + private sealed class TestWorkingSetMember(Func? isCandidateForRemoval = null) : IActivationWorkingSetMemberStatus + { + private bool _isIdle; + private bool _isInWorkingSet; + + public bool IsIdle + { + get => Volatile.Read(ref _isIdle); + set + { + Assert.True(Monitor.IsEntered(this)); + _isIdle = value; + } + } + + public bool IsInWorkingSet + { + get => Volatile.Read(ref _isInWorkingSet); + set + { + Assert.True(Monitor.IsEntered(this)); + _isInWorkingSet = value; + } + } + + public bool IsCandidateForRemoval(bool wouldRemove) + { + Assert.True(Monitor.IsEntered(this)); + return isCandidateForRemoval?.Invoke(wouldRemove) ?? false; + } + + } + + private sealed class PublicWorkingSetMember : IActivationWorkingSetMember + { + public List CandidateCalls { get; } = []; + + public bool IsCandidateForRemoval(bool wouldRemove) + { + Assert.True(Monitor.IsEntered(this)); + CandidateCalls.Add(wouldRemove); + return true; + } + } + + [Fact, TestCategory("Activation")] + public void WorkingSet_SequentialGeneratedTrace_MatchesReferenceModel() + { + var generatedOperation = CsCheck.Gen.Select( + CsCheck.Gen.Int[0, 7], + CsCheck.Gen.Int[0, 3], + CsCheck.Gen.Bool, + static (kind, memberId, candidateEligible) => + new WorkingSetOperation(0, (WorkingSetOperationKind)kind, memberId, candidateEligible)); + var traceGenerator = CsCheck.Gen.SelectMany( + CsCheck.Gen.Int[0, 64], + length => CsCheck.Gen.Select( + generatedOperation.Array[length], + static generated => GetWorkingSetCoverageSpine() + .Concat(generated) + .Select(static (operation, index) => operation with { Index = index }) + .ToArray())); + + CsCheck.Check.Sample( + traceGenerator, + RunSequentialWorkingSetTrace, + seed: "0N0XIzNsQ0O2", + iter: 100, + threads: 1, + print: FormatWorkingSetTrace); + } + + [Fact, TestCategory("Activation")] + public async Task WorkingSet_ConcurrentOnActiveForAbsentMember_AddsOnceAndNotifiesEveryCaller() + { + const int workerCount = 8; + await using var harness = await WorkingSetHarness.CreateAsync(); + var member = harness.Members[0]; + using var ready = new CountdownEvent(workerCount); + var start = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var workers = Enumerable.Range(0, workerCount).Select(worker => Task.Run(async () => + { + ready.Signal(); + await start.Task; + harness.WorkingSet.OnActive(member); + }, TestContext.Current.CancellationToken)).ToArray(); + + try + { + Assert.True( + ready.Wait(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken), + "Workers did not reach the OnActive start gate."); + start.TrySetResult(); + await Task.WhenAll(workers).WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); + } + finally + { + start.TrySetResult(); + await Task.WhenAll(workers).WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); + } + + Assert.Equal(1, harness.WorkingSet.Count); + Assert.True(member.IsInWorkingSet); + Assert.False(member.IsIdle); + Assert.Same(member, Assert.Single(harness.WorkingSet.Members)); + Assert.Equal(Enumerable.Repeat("Active", workerCount), harness.Observer.GetHistory(0)); + Assert.Equal(0, harness.Observer.Count("Added")); + Assert.Equal(workerCount, harness.Observer.Count("Active")); + Assert.Equal(0, harness.Observer.Count("Idle")); + Assert.Equal(0, harness.Observer.Count("Evicted")); + Assert.Equal(0, harness.Observer.Count("Deactivating")); + Assert.Equal(0, harness.Observer.Count("Deactivated")); + } + + [Fact, TestCategory("Activation")] + public async Task WorkingSet_BoundedConcurrentTransitionsPreserveMembershipCount() + { + const int workerCount = 8; + const int operationsPerWorker = 10_000; + await using var harness = await WorkingSetHarness.CreateAsync(); + var visitMember = typeof(ActivationWorkingSet).GetMethod( + "VisitMember", + BindingFlags.Instance | BindingFlags.NonPublic, + [typeof(IActivationWorkingSetMember)]) + ?? throw new InvalidOperationException("Could not find the working-set scan method."); + foreach (var member in harness.Members) + { + harness.WorkingSet.OnActivated(member); + } + + var workers = Enumerable.Range(0, workerCount).Select(worker => Task.Run(() => + { + var random = new Random(42 + worker); + for (var i = 0; i < operationsPerWorker; i++) + { + var memberId = random.Next(harness.Members.Count); + var member = harness.Members[memberId]; + switch (random.Next(4)) + { + case 0: + harness.WorkingSet.OnActive(member); + break; + case 1: + harness.WorkingSet.OnEvicted(member); + break; + case 2: + harness.MemberStates[memberId].CandidateEligible = random.Next(2) == 0; + visitMember.Invoke(harness.WorkingSet, [member]); + break; + default: + harness.WorkingSet.OnActive(member); + visitMember.Invoke(harness.WorkingSet, [member]); + break; + } + } + }, TestContext.Current.CancellationToken)).ToArray(); + + await Task.WhenAll(workers); + + foreach (var memberState in harness.MemberStates) + { + memberState.CandidateEligible = false; + } + + foreach (var member in harness.Members) + { + harness.WorkingSet.OnActive(member); + } + + Assert.Equal(harness.Members.Count, harness.WorkingSet.Count); + Assert.True(harness.Members.Cast().ToHashSet().SetEquals(harness.WorkingSet.Members)); + Assert.All(harness.Members, member => + { + Assert.True(member.IsInWorkingSet); + Assert.False(member.IsIdle); + }); + } + + [Fact, TestCategory("Activation")] + public async Task WorkingSet_EvictionCallbackCanOverlapReAddWithoutHoldingMemberLock() + { + await using var harness = await WorkingSetHarness.CreateAsync(); + var member = harness.Members[0]; + harness.WorkingSet.OnActivated(member); + harness.Observer.Clear(); + var gate = harness.Observer.ArmEvictionGate(0); + var eviction = Task.Run( + () => harness.WorkingSet.OnEvicted(member), + TestContext.Current.CancellationToken); + Task? reAdd = null; + try + { + await gate.Entered.Task.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); + reAdd = Task.Run( + () => harness.WorkingSet.OnActive(member), + TestContext.Current.CancellationToken); + await reAdd.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); + Assert.Equal(1, harness.WorkingSet.Count); + Assert.True(member.IsInWorkingSet); + Assert.False(member.IsIdle); + Assert.Same(member, Assert.Single(harness.WorkingSet.Members)); + Assert.Equal(["EvictedStarted", "Active"], harness.Observer.GetHistory(0)); + } + finally + { + gate.Release.TrySetResult(); + if (reAdd is not null) + { + await reAdd.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); + } + + await eviction.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); + } + + Assert.Equal(["EvictedStarted", "Active", "EvictedCompleted"], harness.Observer.GetHistory(0)); + Assert.Equal(0, harness.Observer.Count("Added")); + Assert.Equal(1, harness.WorkingSet.Count); + Assert.True(member.IsInWorkingSet); + Assert.False(member.IsIdle); + Assert.Same(member, Assert.Single(harness.WorkingSet.Members)); + } + + private static WorkingSetOperation[] GetWorkingSetCoverageSpine() => + [ + new(0, WorkingSetOperationKind.Activate, 0, false), + new(0, WorkingSetOperationKind.Active, 0, false), + new(0, WorkingSetOperationKind.SetCandidate, 0, false), + new(0, WorkingSetOperationKind.Scan, 0, false), + new(0, WorkingSetOperationKind.SetCandidate, 0, true), + new(0, WorkingSetOperationKind.Scan, 0, false), + new(0, WorkingSetOperationKind.Scan, 0, false), + new(0, WorkingSetOperationKind.Activate, 1, false), + new(0, WorkingSetOperationKind.Evict, 1, false), + new(0, WorkingSetOperationKind.Activate, 1, false), + new(0, WorkingSetOperationKind.Deactivating, 1, false), + new(0, WorkingSetOperationKind.Deactivated, 1, false), + new(0, WorkingSetOperationKind.Active, 2, false), + new(0, WorkingSetOperationKind.DeactivatePair, 2, false), + new(0, WorkingSetOperationKind.Activate, 3, false), + new(0, WorkingSetOperationKind.Evict, 3, false) + ]; + + private static string FormatWorkingSetTrace(WorkingSetOperation[] trace) + => string.Join(Environment.NewLine, trace.Select(static operation => operation.ToString())); + + private static void RunSequentialWorkingSetTrace(WorkingSetOperation[] trace) + { + var harness = WorkingSetHarness.CreateAsync().GetAwaiter().GetResult(); + try + { + var model = new WorkingSetReferenceModel(harness.Members.Count); + model.AssertMatches(harness); + foreach (var operation in trace) + { + var expectedDuplicate = model.Apply(operation); + var actualDuplicate = ExecuteWorkingSetOperation(harness, operation, expectedDuplicate); + Assert.Equal(expectedDuplicate, actualDuplicate); + model.AssertMatches(harness); + } + } + finally + { + harness.DisposeAsync().AsTask().GetAwaiter().GetResult(); + } + } + + private static bool ExecuteWorkingSetOperation( + WorkingSetHarness harness, + WorkingSetOperation operation, + bool expectedDuplicate) + { + var member = harness.Members[operation.MemberId]; + switch (operation.Kind) + { + case WorkingSetOperationKind.Activate: + { + var exception = Record.Exception(() => harness.WorkingSet.OnActivated(member)); + if (expectedDuplicate) + { + Assert.IsType(exception); + } + else if (exception is not null) + { + System.Runtime.ExceptionServices.ExceptionDispatchInfo.Capture(exception).Throw(); + } + + return exception is not null; + } + case WorkingSetOperationKind.Active: + harness.WorkingSet.OnActive(member); + break; + case WorkingSetOperationKind.SetCandidate: + harness.MemberStates[operation.MemberId].CandidateEligible = operation.CandidateEligible; + break; + case WorkingSetOperationKind.Scan: + harness.ScanOnceAsync().GetAwaiter().GetResult(); + break; + case WorkingSetOperationKind.Evict: + harness.WorkingSet.OnEvicted(member); + break; + case WorkingSetOperationKind.Deactivating: + harness.WorkingSet.OnDeactivating(member); + break; + case WorkingSetOperationKind.Deactivated: + harness.WorkingSet.OnDeactivated(member); + break; + case WorkingSetOperationKind.DeactivatePair: + harness.WorkingSet.OnDeactivating(member); + harness.WorkingSet.OnDeactivated(member); + break; + default: + throw new ArgumentOutOfRangeException(nameof(operation)); + } + + return false; + } + + private enum WorkingSetOperationKind + { + Activate, + Active, + SetCandidate, + Scan, + Evict, + Deactivating, + Deactivated, + DeactivatePair + } + + private readonly record struct WorkingSetOperation( + int Index, + WorkingSetOperationKind Kind, + int MemberId, + bool CandidateEligible) + { + public override string ToString() + => $"#{Index}: {Kind}(member={MemberId}, candidateEligible={CandidateEligible})"; + } + + private sealed class WorkingSetMemberState + { + private readonly ConcurrentQueue _candidateCalls = new(); + private int _candidateEligible; + + public bool CandidateEligible + { + get => Volatile.Read(ref _candidateEligible) != 0; + set => Volatile.Write(ref _candidateEligible, value ? 1 : 0); + } + + public bool IsCandidateForRemoval(bool wouldRemove) + { + _candidateCalls.Enqueue(wouldRemove); + return CandidateEligible; + } + + public bool[] GetCandidateCalls() => _candidateCalls.ToArray(); + } + + private sealed class RecordingWorkingSetObserver( + Func getMemberId, + int memberCount) : IActivationWorkingSetObserver + { + private readonly ConcurrentQueue[] _history = + Enumerable.Range(0, memberCount).Select(static _ => new ConcurrentQueue()).ToArray(); + private readonly object _gateLock = new(); + private int _gatedMemberId = -1; + private TaskCompletionSource? _evictionEntered; + private TaskCompletionSource? _evictionRelease; + + public void OnAdded(IActivationWorkingSetMember member) => Record(member, "Added"); + + public void OnActive(IActivationWorkingSetMember member) => Record(member, "Active"); + + public void OnIdle(IActivationWorkingSetMember member) => Record(member, "Idle"); + + public void OnEvicted(IActivationWorkingSetMember member) + { + var memberId = getMemberId(member); + TaskCompletionSource? entered; + TaskCompletionSource? release; + lock (_gateLock) + { + entered = memberId == _gatedMemberId ? _evictionEntered : null; + release = memberId == _gatedMemberId ? _evictionRelease : null; + } + + if (entered is null || release is null) + { + Record(memberId, "Evicted"); + return; + } + + Record(memberId, "EvictedStarted"); + entered.TrySetResult(); + release.Task.GetAwaiter().GetResult(); + Record(memberId, "EvictedCompleted"); + } + + public void OnDeactivating(IActivationWorkingSetMember member) => Record(member, "Deactivating"); + + public void OnDeactivated(IActivationWorkingSetMember member) => Record(member, "Deactivated"); + + public (TaskCompletionSource Entered, TaskCompletionSource Release) ArmEvictionGate(int memberId) + { + var entered = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var release = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + lock (_gateLock) + { + _gatedMemberId = memberId; + _evictionEntered = entered; + _evictionRelease = release; + } + + return (entered, release); + } + + public string[] GetHistory(int memberId) => _history[memberId].ToArray(); + + public int Count(string eventName) + => _history.Sum(history => history.Count(item => string.Equals(item, eventName, StringComparison.Ordinal))); + + public void Clear() + { + foreach (var history in _history) + { + while (history.TryDequeue(out _)) + { + } + } + } + + private void Record(IActivationWorkingSetMember member, string eventName) + => Record(getMemberId(member), eventName); + + private void Record(int memberId, string eventName) => _history[memberId].Enqueue(eventName); + } + + private sealed class WorkingSetReferenceModel + { + private readonly ModelMember[] _members; + + public WorkingSetReferenceModel(int memberCount) + { + _members = Enumerable.Range(0, memberCount).Select(static _ => new ModelMember()).ToArray(); + } + + public bool Apply(WorkingSetOperation operation) + { + var member = _members[operation.MemberId]; + switch (operation.Kind) + { + case WorkingSetOperationKind.Activate: + if (member.Present) + { + return true; + } + + member.Present = true; + member.IsInWorkingSet = true; + member.IsIdle = false; + member.History.Add("Added"); + break; + case WorkingSetOperationKind.Active: + member.Present = true; + member.IsInWorkingSet = true; + member.IsIdle = false; + member.History.Add("Active"); + break; + case WorkingSetOperationKind.SetCandidate: + member.CandidateEligible = operation.CandidateEligible; + break; + case WorkingSetOperationKind.Scan: + foreach (var scanMember in _members) + { + Scan(scanMember); + } + + break; + case WorkingSetOperationKind.Evict: + Evict(member); + break; + case WorkingSetOperationKind.Deactivating: + Evict(member); + member.History.Add("Deactivating"); + break; + case WorkingSetOperationKind.Deactivated: + Evict(member); + member.History.Add("Deactivated"); + break; + case WorkingSetOperationKind.DeactivatePair: + Evict(member); + member.History.Add("Deactivating"); + Evict(member); + member.History.Add("Deactivated"); + break; + default: + throw new ArgumentOutOfRangeException(nameof(operation)); + } + + return false; + } + + public void AssertMatches(WorkingSetHarness harness) + { + Assert.Equal(_members.Count(static member => member.Present), harness.WorkingSet.Count); + var expectedVisibleMembers = _members + .Select(static (member, id) => (member, id)) + .Where(static item => item.member.Present && !item.member.IsIdle) + .Select(static item => item.id) + .ToArray(); + var actualVisibleMembers = harness.WorkingSet.Members + .Select(harness.GetMemberId) + .Order() + .ToArray(); + Assert.Equal(expectedVisibleMembers, actualVisibleMembers); + + for (var memberId = 0; memberId < _members.Length; memberId++) + { + var expected = _members[memberId]; + var actual = harness.Members[memberId]; + Assert.Equal(expected.IsInWorkingSet, actual.IsInWorkingSet); + Assert.Equal(expected.IsIdle, actual.IsIdle); + Assert.Equal(expected.CandidateCalls, harness.MemberStates[memberId].GetCandidateCalls()); + Assert.Equal(expected.History, harness.Observer.GetHistory(memberId)); + } + } + + private static void Scan(ModelMember member) + { + if (!member.IsInWorkingSet) + { + return; + } + + var wouldRemove = member.IsIdle; + member.CandidateCalls.Add(wouldRemove); + if (!member.CandidateEligible) + { + member.IsIdle = false; + member.History.Add("Active"); + } + else if (!wouldRemove) + { + member.IsIdle = true; + member.History.Add("Idle"); + } + else if (member.Present) + { + member.Present = false; + member.IsInWorkingSet = false; + member.IsIdle = false; + member.History.Add("Evicted"); + } + } + + private static void Evict(ModelMember member) + { + if (!member.Present) + { + return; + } + + member.Present = false; + member.IsInWorkingSet = false; + member.IsIdle = false; + member.History.Add("Evicted"); + } + + private sealed class ModelMember + { + public bool Present { get; set; } + public bool IsInWorkingSet { get; set; } + public bool IsIdle { get; set; } + public bool CandidateEligible { get; set; } + public List CandidateCalls { get; } = []; + public List History { get; } = []; + } + } + + private sealed class WorkingSetHarness : IAsyncDisposable + { + private readonly ServiceProvider _serviceProvider; + private readonly ControlledAsyncTimer _timer; + private readonly SiloLifecycleSubject _lifecycle; + + private WorkingSetHarness() + { + var services = new ServiceCollection(); + services.AddMetrics(); + services.AddSingleton(); + services.AddSingleton(); + _serviceProvider = services.BuildServiceProvider(); + MemberStates = Enumerable.Range(0, 4).Select(static _ => new WorkingSetMemberState()).ToArray(); + Members = MemberStates + .Select(state => new TestWorkingSetMember(state.IsCandidateForRemoval)) + .ToArray(); + var memberIds = Members + .Select(static (member, id) => (member, id)) + .ToDictionary(static item => (IActivationWorkingSetMember)item.member, static item => item.id); + GetMemberId = member => memberIds[member]; + Observer = new RecordingWorkingSetObserver(GetMemberId, Members.Count); + _timer = new ControlledAsyncTimer(); + var timerFactory = Substitute.For(); + timerFactory.Create(Arg.Any(), Arg.Any(), Arg.Any()).Returns(_timer); + WorkingSet = new ActivationWorkingSet( + timerFactory, + NullLogger.Instance, + [Observer], + _serviceProvider.GetRequiredService(), + TimeProvider.System); + _lifecycle = new SiloLifecycleSubject(NullLogger.Instance); + ((ILifecycleParticipant)WorkingSet).Participate(_lifecycle); + } + + public ActivationWorkingSet WorkingSet { get; } + public IReadOnlyList Members { get; } + public IReadOnlyList MemberStates { get; } + public RecordingWorkingSetObserver Observer { get; } + public Func GetMemberId { get; } + public int TimerGeneration => _timer.Generation; + + public static async Task CreateAsync() + { + var result = new WorkingSetHarness(); + try + { + await result._lifecycle.OnStart(TestContext.Current.CancellationToken); + await result._timer.WaitForGenerationAsync(1); + return result; + } + catch + { + await result.DisposeAsync(); + throw; + } + } + + public async Task ScanOnceAsync() + { + var nextGeneration = _timer.Generation + 1; + _timer.CompleteCurrent(result: true); + await _timer.WaitForGenerationAsync(nextGeneration); + } + + public async ValueTask DisposeAsync() + { + try + { + await _lifecycle.OnStop(TestContext.Current.CancellationToken) + .WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); + } + finally + { + await _serviceProvider.DisposeAsync(); + } + } + } + + private sealed class ControlledAsyncTimer : IAsyncTimer + { + private readonly object _lock = new(); + private TaskCompletionSource _generationChanged = + new(TaskCreationOptions.RunContinuationsAsynchronously); + private TaskCompletionSource? _current; + private bool _disposed; + private int _generation; + + public int Generation + { + get + { + lock (_lock) + { + return _generation; + } + } + } + + public Task NextTick(TimeSpan? overrideDelay = default) + { + TaskCompletionSource generationChanged; + TaskCompletionSource current; + lock (_lock) + { + if (_disposed) + { + return Task.FromResult(false); + } + + Assert.Null(_current); + current = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + _current = current; + _generation++; + generationChanged = _generationChanged; + _generationChanged = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + } + + generationChanged.TrySetResult(); + return current.Task; + } + + public void CompleteCurrent(bool result) + { + TaskCompletionSource current; + lock (_lock) + { + current = _current ?? throw new InvalidOperationException("The working-set timer is not awaiting a tick."); + _current = null; + } + + current.TrySetResult(result); + } + + public async Task WaitForGenerationAsync(int expectedGeneration) + { + while (true) + { + Task generationChanged; + lock (_lock) + { + if (_generation >= expectedGeneration) + { + return; + } + + generationChanged = _generationChanged.Task; + } + + await generationChanged.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); + } + } + + public bool CheckHealth(DateTime lastCheckTime, [NotNullWhen(false)] out string? reason) + { + reason = null; + return true; + } + + public void Dispose() + { + TaskCompletionSource? current; + lock (_lock) + { + if (_disposed) + { + return; + } + + _disposed = true; + current = _current; + _current = null; + } + + current?.TrySetResult(false); + } + } + + [Fact, TestCategory("Activation")] + public void ActivationData_Constructor_InitializesWorkingSetClockStatus() + { + using var fixture = new ActivationDataWorkingSetFixture(); + + Assert.Equal(ActivationState.Creating, fixture.Activation.State); + Assert.True(fixture.Member.IsInWorkingSet); + Assert.False(fixture.Member.IsIdle); + Assert.False(fixture.WasRemovedByCollection); + } + + [Theory, TestCategory("Activation")] + [MemberData(nameof(ActivationStatusCases))] + public void ActivationStatus_PackedFieldsPreserveIndependentBits( + int stateValue, + bool expectedInWorkingSet, + bool expectedIdle) + { + using var fixture = new ActivationDataWorkingSetFixture(); + var expectedState = (ActivationState)stateValue; + var otherState = expectedState == ActivationState.Invalid + ? ActivationState.Creating + : ActivationState.Invalid; + + lock (fixture.Activation) + { + fixture.Member.IsInWorkingSet = expectedInWorkingSet; + fixture.Member.IsIdle = expectedIdle; + fixture.Activation.SetState(otherState); + } + + Assert.Equal(otherState, fixture.Activation.State); + Assert.Equal(expectedInWorkingSet, fixture.Member.IsInWorkingSet); + Assert.Equal(expectedIdle, fixture.Member.IsIdle); + Assert.False(fixture.WasRemovedByCollection); + + lock (fixture.Activation) + { + fixture.Activation.SetState(expectedState); + } + + Assert.Equal(expectedState, fixture.Activation.State); + Assert.Equal(expectedInWorkingSet, fixture.Member.IsInWorkingSet); + Assert.Equal(expectedIdle, fixture.Member.IsIdle); + Assert.False(fixture.WasRemovedByCollection); + + lock (fixture.Activation) + { + fixture.Member.IsInWorkingSet = !expectedInWorkingSet; + } + + Assert.Equal(expectedState, fixture.Activation.State); + Assert.Equal(!expectedInWorkingSet, fixture.Member.IsInWorkingSet); + Assert.Equal(expectedIdle, fixture.Member.IsIdle); + Assert.False(fixture.WasRemovedByCollection); + + lock (fixture.Activation) + { + fixture.Member.IsInWorkingSet = expectedInWorkingSet; + fixture.Member.IsIdle = !expectedIdle; + } + + Assert.Equal(expectedState, fixture.Activation.State); + Assert.Equal(expectedInWorkingSet, fixture.Member.IsInWorkingSet); + Assert.Equal(!expectedIdle, fixture.Member.IsIdle); + Assert.False(fixture.WasRemovedByCollection); + + lock (fixture.Activation) + { + fixture.Member.IsIdle = expectedIdle; + } + + Assert.Equal(expectedState, fixture.Activation.State); + Assert.Equal(expectedInWorkingSet, fixture.Member.IsInWorkingSet); + Assert.Equal(expectedIdle, fixture.Member.IsIdle); + Assert.False(fixture.WasRemovedByCollection); + } + + [Fact, TestCategory("Activation")] + public void ActivationData_CollectionCandidateMarker_RequiresEligibleRemovalPass() + { + using var fixture = new ActivationDataWorkingSetFixture(); + + fixture.AdvanceIdleDurationTo(10_000); + bool isCandidateAtBoundary; + lock (fixture.Activation) + { + isCandidateAtBoundary = fixture.Member.IsCandidateForRemoval(wouldRemove: true); + } + + Assert.False(isCandidateAtBoundary); + Assert.False(fixture.WasRemovedByCollection); + + fixture.AdvanceIdleDurationTo(10_001); + bool isCandidateOnFirstPass; + lock (fixture.Activation) + { + isCandidateOnFirstPass = fixture.Member.IsCandidateForRemoval(wouldRemove: false); + } + + Assert.True(isCandidateOnFirstPass); + Assert.False(fixture.WasRemovedByCollection); + + bool isCandidateOnRemovalPass; + lock (fixture.Activation) + { + isCandidateOnRemovalPass = fixture.Member.IsCandidateForRemoval(wouldRemove: true); + } + + Assert.True(isCandidateOnRemovalPass); + Assert.True(fixture.WasRemovedByCollection); + + lock (fixture.Activation) + { + fixture.Activation.SetState(ActivationState.Valid); + fixture.Member.IsIdle = true; + } + + Assert.Equal(ActivationState.Valid, fixture.Activation.State); + Assert.True(fixture.Member.IsInWorkingSet); + Assert.True(fixture.Member.IsIdle); + Assert.True(fixture.WasRemovedByCollection); + } + + [Fact, TestCategory("Activation")] + public void ActivationData_ClockCollectionThenOnActive_ClearsCollectionMarker() + { + using var fixture = new ActivationDataWorkingSetFixture(); + lock (fixture.Activation) + { + fixture.Activation.SetState(ActivationState.Valid); + } + + fixture.WorkingSet.OnActivated(fixture.Member); + fixture.AdvanceIdleDurationTo(10_001); + + fixture.ScanOnce(); + + Assert.Equal(1, fixture.WorkingSet.Count); + Assert.True(fixture.Member.IsInWorkingSet); + Assert.True(fixture.Member.IsIdle); + Assert.False(fixture.WasRemovedByCollection); + Assert.Empty(fixture.WorkingSet.Members); + Assert.Equal(["Added", "Idle"], fixture.Observer.GetHistory(0)); + + fixture.ScanOnce(); + + Assert.Equal(0, fixture.WorkingSet.Count); + Assert.False(fixture.Member.IsInWorkingSet); + Assert.False(fixture.Member.IsIdle); + Assert.True(fixture.WasRemovedByCollection); + Assert.Empty(fixture.WorkingSet.Members); + Assert.Equal(["Added", "Idle", "Evicted"], fixture.Observer.GetHistory(0)); + + fixture.WorkingSet.OnActive(fixture.Member); + + Assert.Equal(1, fixture.WorkingSet.Count); + Assert.True(fixture.Member.IsInWorkingSet); + Assert.False(fixture.Member.IsIdle); + Assert.False(fixture.WasRemovedByCollection); + Assert.Equal([fixture.Member], fixture.WorkingSet.Members); + Assert.Equal(["Added", "Idle", "Evicted", "Active"], fixture.Observer.GetHistory(0)); + } + + [Fact, TestCategory("Activation")] + public void ActivationData_ExplicitWorkingSetDeactivation_DoesNotSetCollectionMarker() + { + using var fixture = new ActivationDataWorkingSetFixture(); + lock (fixture.Activation) + { + fixture.Activation.SetState(ActivationState.Valid); + } + + fixture.WorkingSet.OnActivated(fixture.Member); + + fixture.WorkingSet.OnDeactivating(fixture.Member); + + Assert.Equal(0, fixture.WorkingSet.Count); + Assert.False(fixture.Member.IsInWorkingSet); + Assert.False(fixture.Member.IsIdle); + Assert.False(fixture.WasRemovedByCollection); + Assert.Empty(fixture.WorkingSet.Members); + Assert.Equal(["Added", "Evicted", "Deactivating"], fixture.Observer.GetHistory(0)); + } + + [Fact, TestCategory("Activation")] + public void ActivationData_CompletedRequest_DoesNotReaddDeactivatingActivation() + { + using var fixture = new ActivationDataWorkingSetFixture(); + lock (fixture.Activation) + { + fixture.Activation.SetState(ActivationState.Deactivating); + fixture.Member.IsInWorkingSet = false; + fixture.Member.IsIdle = false; + } + + fixture.CompleteRequest(new Message()); + + Assert.Equal(ActivationState.Deactivating, fixture.Activation.State); + Assert.False(fixture.Member.IsInWorkingSet); + Assert.False(fixture.Member.IsIdle); + Assert.Equal(0, fixture.WorkingSet.Count); + Assert.Empty(fixture.Observer.GetHistory(0)); + } + + public static IEnumerable ActivationStatusCases() + { + foreach (var state in Enum.GetValues()) + { + yield return [(int)state, false, false]; + yield return [(int)state, false, true]; + yield return [(int)state, true, false]; + yield return [(int)state, true, true]; + } + } + + private sealed class ActivationDataWorkingSetFixture : IDisposable + { + private static readonly PropertyInfo WasRemovedByCollectionProperty = typeof(ActivationData).GetProperty( + "WasRemovedByCollection", + BindingFlags.Instance | BindingFlags.NonPublic) + ?? throw new InvalidOperationException("Could not find the collection-removal marker."); + private static readonly FieldInfo IdleDurationField = typeof(ActivationData).GetField( + "_idleDuration", + BindingFlags.Instance | BindingFlags.NonPublic) + ?? throw new InvalidOperationException("Could not find the activation idle-duration field."); + private static readonly FieldInfo ServiceScopeField = typeof(ActivationData).GetField( + "_serviceScope", + BindingFlags.Instance | BindingFlags.NonPublic) + ?? throw new InvalidOperationException("Could not find the activation service scope."); + private static readonly FieldInfo SharedSchedulerLoggerField = typeof(GrainTypeSharedContext).GetField( + "k__BackingField", + BindingFlags.Instance | BindingFlags.NonPublic) + ?? throw new InvalidOperationException("Could not find the shared scheduler logger."); + private static readonly MethodInfo VisitMemberMethod = typeof(ActivationWorkingSet).GetMethod( + "VisitMember", + BindingFlags.Instance | BindingFlags.NonPublic, + [typeof(IActivationWorkingSetMember)]) + ?? throw new InvalidOperationException("Could not find the working-set scan method."); + private static readonly MethodInfo CompleteRequestMethod = typeof(ActivationData).GetMethod( + "OnCompletedRequest", + BindingFlags.Instance | BindingFlags.NonPublic, + [typeof(Message)]) + ?? throw new InvalidOperationException("Could not find the completed-request method."); + + private readonly ServiceProvider _serviceProvider; + private readonly IServiceScope _activationScope; + private readonly ControlledAsyncTimer _timer; + private long _elapsedMilliseconds; + + public ActivationDataWorkingSetFixture() + { + TimeProvider = new FakeTimeProvider(DateTimeOffset.Parse("2025-01-01T00:00:00.000+00:00")); + var services = new ServiceCollection(); + services.AddOptions(); + services.AddLogging(); + services.AddMetrics(); + services.AddSingleton(TimeProvider); + services.AddSingleton(TimeProvider); + services.AddSingleton(); + services.AddSingleton(); + services.AddSingleton(); + services.Configure(options => + { + options.DelayWarningThreshold = TimeSpan.FromMilliseconds(100); + options.ActivationSchedulingQuantum = TimeSpan.FromMilliseconds(100); + options.TurnWarningLengthThreshold = TimeSpan.FromMilliseconds(100); + options.StoppedActivationWarningInterval = TimeSpan.FromMilliseconds(200); + }); + _serviceProvider = services.BuildServiceProvider(); + + var address = GrainAddress.NewActivationAddress( + SiloAddress.New(System.Net.IPAddress.Loopback, 11_111, 1), + GrainId.Create("activation-working-set", "clock-fixture")); + var shared = (GrainTypeSharedContext)RuntimeHelpers.GetUninitializedObject(typeof(GrainTypeSharedContext)); + SharedSchedulerLoggerField.SetValue( + shared, + _serviceProvider.GetRequiredService() + .CreateLogger(typeof(Orleans.Runtime.Scheduler.WorkItemGroup).FullName!)); + Activation = new ActivationData( + address, + context => new Orleans.Runtime.Scheduler.WorkItemGroup( + context, + _serviceProvider.GetRequiredService>(), + _serviceProvider.GetRequiredService()), + _serviceProvider, + shared); + Member = Activation; + _activationScope = (IServiceScope)ServiceScopeField.GetValue(Activation)!; + + Observer = new RecordingWorkingSetObserver(_ => 0, 1); + _timer = new ControlledAsyncTimer(); + var timerFactory = Substitute.For(); + timerFactory.Create(Arg.Any(), Arg.Any(), Arg.Any()).Returns(_timer); + WorkingSet = new ActivationWorkingSet( + timerFactory, + NullLogger.Instance, + [Observer], + _serviceProvider.GetRequiredService(), + TimeProvider); + } + + public ActivationData Activation { get; } + public IActivationWorkingSetMemberStatus Member { get; } + public FakeTimeProvider TimeProvider { get; } + public ActivationWorkingSet WorkingSet { get; } + public RecordingWorkingSetObserver Observer { get; } + public bool WasRemovedByCollection + => (bool)WasRemovedByCollectionProperty.GetValue(Activation)!; + + public void AdvanceIdleDurationTo(long elapsedMilliseconds) + { + var advance = elapsedMilliseconds - _elapsedMilliseconds; + Assert.True(advance >= 0); + TimeProvider.Advance(TimeSpan.FromMilliseconds(advance)); + lock (Activation) + { + IdleDurationField.SetValue( + Activation, + CoarseStopwatch.FromTimestamp(0, elapsedMilliseconds)); + } + + _elapsedMilliseconds = elapsedMilliseconds; + Assert.Equal(TimeSpan.FromMilliseconds(elapsedMilliseconds), Activation.GetIdleness()); + } + + public void ScanOnce() => VisitMemberMethod.Invoke(WorkingSet, [Member]); + + public void CompleteRequest(Message message) => CompleteRequestMethod.Invoke(Activation, [message]); + + public void Dispose() + { + _timer.Dispose(); + _activationScope.Dispose(); + _serviceProvider.Dispose(); + } + } } } diff --git a/test/Orleans.Placement.Tests/ActivationRepartitioningTests/DeactivatedGrainQueueTests.cs b/test/Orleans.Placement.Tests/ActivationRepartitioningTests/DeactivatedGrainQueueTests.cs index 4d040fbe0c0..02e1f6c4d0a 100644 --- a/test/Orleans.Placement.Tests/ActivationRepartitioningTests/DeactivatedGrainQueueTests.cs +++ b/test/Orleans.Placement.Tests/ActivationRepartitioningTests/DeactivatedGrainQueueTests.cs @@ -89,7 +89,6 @@ private sealed class TestActivationWorkingSetMember(GrainId grainId) : IActivati public IGrainLifecycle ObservableLifecycle => throw new NotImplementedException(); public IWorkItemScheduler Scheduler => throw new NotImplementedException(); public Task Deactivated => Task.CompletedTask; - public bool IsCandidateForRemoval(bool wouldRemove) => false; public void SetComponent(TComponent? value) where TComponent : class => throw new NotImplementedException(); public void ReceiveMessage(object message) => throw new NotImplementedException();