From 2ad89f76b93e6bf98e2117f12cf1627cdcb35a1a Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Fri, 21 Aug 2026 13:53:54 -0700 Subject: [PATCH 01/12] perf(runtime): reduce activation working set entry overhead --- .../Catalog/ActivationWorkingSet.cs | 426 +++++++++--------- 1 file changed, 223 insertions(+), 203 deletions(-) diff --git a/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs b/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs index bbba7539b3a..6860e57ac3e 100644 --- a/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs +++ b/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs @@ -10,276 +10,296 @@ 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 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; } - } - - 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) - { - _logger = logger; - _scanPeriodTimer = asyncTimerFactory.Create(TimeSpan.FromMilliseconds(5_000), nameof(ActivationWorkingSet) + "." + nameof(MonitorWorkingSet), timeProvider); - _observers = observers.ToList(); - catalogInstruments.RegisterActivationWorkingSetObserve(() => Count); - } + _logger = logger; + _scanPeriodTimer = asyncTimerFactory.Create(TimeSpan.FromMilliseconds(5_000), nameof(ActivationWorkingSet) + "." + nameof(MonitorWorkingSet), timeProvider); + _observers = observers.ToList(); + catalogInstruments.RegisterActivationWorkingSetObserve(() => Count); + } - public int Count => _activeCount; + public int Count => _activeCount; - internal IEnumerable Members => EnumerateActiveMembers(); + internal IEnumerable Members => EnumerateActiveMembers(); - private IEnumerable EnumerateActiveMembers() + private IEnumerable EnumerateActiveMembers() + { + foreach (var pair in _members) { - foreach (var pair in _members) + if (!pair.Value) { - if (!pair.Value.IsIdle) - { - yield return pair.Key; - } + yield return pair.Key; } } + } - public void OnActivated(IActivationWorkingSetMember member) + public void OnActivated(IActivationWorkingSetMember member) + { + Debug.Assert(member is not ICollectibleGrainContext collectible || collectible.IsValid); + if (_members.TryAdd(member, false)) { - Debug.Assert(member is not ICollectibleGrainContext collectible || collectible.IsValid); - if (_members.TryAdd(member, new MemberState())) + Interlocked.Increment(ref _activeCount); + foreach (var observer in _observers) { - Interlocked.Increment(ref _activeCount); - foreach (var observer in _observers) - { - observer.OnAdded(member); - } - - return; + observer.OnAdded(member); } - throw new InvalidOperationException($"Member {member} is already a member of the working set"); + return; } - public void OnActive(IActivationWorkingSetMember member) + throw new InvalidOperationException($"Member {member} is already a member of the working set"); + } + + public void OnActive(IActivationWorkingSetMember member) + { + while (true) { - if (_members.TryGetValue(member, out var state)) + if (_members.TryGetValue(member, out var isIdle)) { - state.IsIdle = false; + if (!isIdle || _members.TryUpdate(member, false, comparisonValue: true)) + { + break; + } } - else if (_members.TryAdd(member, new())) + else if (_members.TryAdd(member, false)) { Interlocked.Increment(ref _activeCount); + break; } + } - foreach (var observer in _observers) - { - observer.OnActive(member); - } + foreach (var observer in _observers) + { + observer.OnActive(member); + } + } + + public void OnEvicted(IActivationWorkingSetMember member) + { + if (_members.TryRemove(member, out _)) + { + OnEvictedCore(member); } + } - public void OnEvicted(IActivationWorkingSetMember member) + private void OnEvicted(IActivationWorkingSetMember member, bool isIdle) + { + if (_members.TryRemove(KeyValuePair.Create(member, isIdle))) { - if (_members.TryRemove(member, out _)) - { - Interlocked.Decrement(ref _activeCount); - foreach (var observer in _observers) - { - observer.OnEvicted(member); - } - } + OnEvictedCore(member); } + } - public void OnDeactivating(IActivationWorkingSetMember member) + private void OnEvictedCore(IActivationWorkingSetMember member) + { + Interlocked.Decrement(ref _activeCount); + 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, pair.Value); + } + catch (Exception exception) + { + LogExceptionVisitingWorkingSetMember(exception, pair.Key); } } } + } - private void VisitMember(IActivationWorkingSetMember member, MemberState state) + private void VisitMember(IActivationWorkingSetMember member, bool isIdle) + { + var wouldRemove = isIdle; + if (member.IsCandidateForRemoval(wouldRemove)) { - var wouldRemove = state.IsIdle; - if (member.IsCandidateForRemoval(wouldRemove)) + if (wouldRemove) { - if (wouldRemove) - { - OnEvicted(member); - } - else + OnEvicted(member, isIdle); + } + else + { + if (_members.TryUpdate(member, true, comparisonValue: isIdle)) { - state.IsIdle = true; foreach (var observer in _observers) { observer.OnIdle(member); } } } - else - { - state.IsIdle = false; - foreach (var observer in _observers) - { - observer.OnActive(member); - } - } } - - void ILifecycleParticipant.Participate(ISiloLifecycle lifecycle) + else { - lifecycle.Subscribe( - nameof(ActivationWorkingSet), - ServiceLifecycleStage.BecomeActive, - StartMonitoring, - StopMonitoring); - - Task StartMonitoring(CancellationToken ct) + if (isIdle) { - using var _ = new ExecutionContextSuppressor(); - _runTask = Task.Run(MonitorWorkingSet); - return Task.CompletedTask; + _members.TryUpdate(member, false, comparisonValue: true); } - async Task StopMonitoring(CancellationToken ct) + foreach (var observer in _observers) { - _scanPeriodTimer.Dispose(); - if (_runTask is Task task) - { - await task.WaitAsync(ct).SuppressThrowing(); - } + observer.OnActive(member); } } + } - [LoggerMessage( - Level = LogLevel.Error, - Message = "Exception visiting working set member {Member}" - )] - private partial void LogExceptionVisitingWorkingSetMember(Exception exception, IActivationWorkingSetMember member); + void ILifecycleParticipant.Participate(ISiloLifecycle lifecycle) + { + lifecycle.Subscribe( + nameof(ActivationWorkingSet), + ServiceLifecycleStage.BecomeActive, + StartMonitoring, + StopMonitoring); + + 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) + { + await task.WaitAsync(ct).SuppressThrowing(); + } + } } + [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); +} + +/// +/// 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) { } } From db72d9f7d706e56bb4904731c6f8a648400f478d Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Sat, 22 Aug 2026 02:27:30 -0700 Subject: [PATCH 02/12] fix(runtime): preserve re-added working set members --- .../Catalog/ActivationWorkingSet.cs | 38 +++++++++++------- .../Runtime/ActivationCollectorTests.cs | 40 ++++++++++++++++++- 2 files changed, 63 insertions(+), 15 deletions(-) diff --git a/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs b/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs index 6860e57ac3e..26a47ca7f00 100644 --- a/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs +++ b/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs @@ -17,12 +17,14 @@ namespace Orleans.Runtime; /// internal sealed partial class ActivationWorkingSet : IActivationWorkingSet, ILifecycleParticipant { - private readonly ConcurrentDictionary _members = new(); + private const long IdleMask = 1; + private readonly ConcurrentDictionary _members = new(); private readonly ILogger _logger; private readonly IAsyncTimer _scanPeriodTimer; private readonly List _observers; private int _activeCount; + private long _nextGeneration; private Task? _runTask; public ActivationWorkingSet( @@ -46,7 +48,7 @@ private IEnumerable EnumerateActiveMembers() { foreach (var pair in _members) { - if (!pair.Value) + if (!IsIdle(pair.Value)) { yield return pair.Key; } @@ -56,7 +58,7 @@ private IEnumerable EnumerateActiveMembers() public void OnActivated(IActivationWorkingSetMember member) { Debug.Assert(member is not ICollectibleGrainContext collectible || collectible.IsValid); - if (_members.TryAdd(member, false)) + if (_members.TryAdd(member, GetNextActiveState())) { Interlocked.Increment(ref _activeCount); foreach (var observer in _observers) @@ -74,14 +76,14 @@ public void OnActive(IActivationWorkingSetMember member) { while (true) { - if (_members.TryGetValue(member, out var isIdle)) + if (_members.TryGetValue(member, out var state)) { - if (!isIdle || _members.TryUpdate(member, false, comparisonValue: true)) + if (!IsIdle(state) || _members.TryUpdate(member, GetActiveState(state), comparisonValue: state)) { break; } } - else if (_members.TryAdd(member, false)) + else if (_members.TryAdd(member, GetNextActiveState())) { Interlocked.Increment(ref _activeCount); break; @@ -102,9 +104,9 @@ public void OnEvicted(IActivationWorkingSetMember member) } } - private void OnEvicted(IActivationWorkingSetMember member, bool isIdle) + private void OnEvicted(IActivationWorkingSetMember member, long state) { - if (_members.TryRemove(KeyValuePair.Create(member, isIdle))) + if (_members.TryRemove(KeyValuePair.Create(member, state))) { OnEvictedCore(member); } @@ -155,18 +157,18 @@ private async Task MonitorWorkingSet() } } - private void VisitMember(IActivationWorkingSetMember member, bool isIdle) + private void VisitMember(IActivationWorkingSetMember member, long state) { - var wouldRemove = isIdle; + var wouldRemove = IsIdle(state); if (member.IsCandidateForRemoval(wouldRemove)) { if (wouldRemove) { - OnEvicted(member, isIdle); + OnEvicted(member, state); } else { - if (_members.TryUpdate(member, true, comparisonValue: isIdle)) + if (_members.TryUpdate(member, GetIdleState(state), comparisonValue: state)) { foreach (var observer in _observers) { @@ -177,9 +179,9 @@ private void VisitMember(IActivationWorkingSetMember member, bool isIdle) } else { - if (isIdle) + if (wouldRemove) { - _members.TryUpdate(member, false, comparisonValue: true); + _members.TryUpdate(member, GetActiveState(state), comparisonValue: state); } foreach (var observer in _observers) @@ -189,6 +191,14 @@ private void VisitMember(IActivationWorkingSetMember member, bool isIdle) } } + private long GetNextActiveState() => unchecked(Interlocked.Increment(ref _nextGeneration) << 1); + + private static bool IsIdle(long state) => (state & IdleMask) != 0; + + private static long GetActiveState(long state) => state & ~IdleMask; + + private static long GetIdleState(long state) => state | IdleMask; + void ILifecycleParticipant.Participate(ISiloLifecycle lifecycle) { lifecycle.Subscribe( diff --git a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs index 24180c647f5..6f71122f46a 100644 --- a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs +++ b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs @@ -108,7 +108,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 +380,44 @@ 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 workingSet = new ActivationWorkingSet( + timerFactory, + NullLogger.Instance, + Array.Empty(), + CreateCatalogInstruments(), + TimeProvider.System); + var member = Substitute.For(); + var scanStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + using var resumeScan = new ManualResetEventSlim(); + member.IsCandidateForRemoval(false).Returns(_ => + { + scanStarted.TrySetResult(); + resumeScan.Wait(); + return true; + }); + workingSet.OnActivated(member); + + var lifecycle = new SiloLifecycleSubject(NullLogger.Instance); + ((ILifecycleParticipant)workingSet).Participate(lifecycle); + await lifecycle.OnStart(); + await scanStarted.Task.WaitAsync(TimeSpan.FromSeconds(10)); + + workingSet.OnEvicted(member); + workingSet.OnActivated(member); + resumeScan.Set(); + await lifecycle.OnStop(); + + Assert.Equal(1, workingSet.Count); + Assert.Contains(member, workingSet.Members); + } + private IActivationWorkingSetMember PrepareActivation(int collectionAgeLimitMinutes, ActivationCollector collector) => PrepareActivation(TimeSpan.FromMinutes(collectionAgeLimitMinutes), collector); From 3529dbddbae2943bd8db57d5597c8b5d0c6dffaf Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Sun, 23 Aug 2026 07:52:15 -0700 Subject: [PATCH 03/12] test(runtime): strengthen activation working set invariants --- .../Catalog/ActivationWorkingSet.cs | 3 + .../Runtime/ActivationCollectorTests.cs | 232 +++++++++++++++++- 2 files changed, 232 insertions(+), 3 deletions(-) diff --git a/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs b/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs index 26a47ca7f00..ce483e14023 100644 --- a/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs +++ b/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs @@ -17,6 +17,9 @@ namespace Orleans.Runtime; /// internal sealed partial class ActivationWorkingSet : IActivationWorkingSet, ILifecycleParticipant { + // The low bit stores idle state and the remaining 63 bits identify the membership generation. + // Comparing the complete state prevents stale scan entries from changing a member after removal and re-addition. + // A generation is reused only after 2^63 successful additions within one process lifetime. private const long IdleMask = 1; private readonly ConcurrentDictionary _members = new(); private readonly ILogger _logger; diff --git a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs index 6f71122f46a..bbf538fa826 100644 --- a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs +++ b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs @@ -1,4 +1,5 @@ using System.Collections.Concurrent; +using System.Reflection; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging.Abstractions; using Microsoft.Extensions.Options; @@ -387,10 +388,11 @@ public async Task WorkingSetScan_DoesNotUpdateReaddedMember() 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, - Array.Empty(), + [observer], CreateCatalogInstruments(), TimeProvider.System); var member = Substitute.For(); @@ -409,11 +411,230 @@ public async Task WorkingSetScan_DoesNotUpdateReaddedMember() await lifecycle.OnStart(); await scanStarted.Task.WaitAsync(TimeSpan.FromSeconds(10)); - workingSet.OnEvicted(member); + try + { + workingSet.OnEvicted(member); + workingSet.OnActivated(member); + var stopTask = lifecycle.OnStop(); + Assert.False(stopTask.IsCompleted); + resumeScan.Set(); + 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.DidNotReceive().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(); + await removalScanStarted.Task.WaitAsync(TimeSpan.FromSeconds(10)); + + try + { + workingSet.OnEvicted(member); + workingSet.OnActivated(member); + var stopTask = lifecycle.OnStop(); + Assert.False(stopTask.IsCompleted); + resumeScan.Set(); + 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); - resumeScan.Set(); + + var lifecycle = new SiloLifecycleSubject(NullLogger.Instance); + ((ILifecycleParticipant)workingSet).Participate(lifecycle); + await lifecycle.OnStart(); await lifecycle.OnStop(); + 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); + } + }); + await firstEviction.Task.WaitAsync(TimeSpan.FromSeconds(10)); + + try + { + while (enumerator.MoveNext()) + { + enumeratedMembers.Add(enumerator.Current); + } + } + finally + { + resumeWriter.Set(); + } + + await writer.WaitAsync(TimeSpan.FromSeconds(10)); + + 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 async Task WorkingSetGenerationWrap_DoesNotAliasAdjacentMembership() + { + 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 workingSet = new ActivationWorkingSet( + timerFactory, + NullLogger.Instance, + Array.Empty(), + CreateCatalogInstruments(), + TimeProvider.System); + var nextGeneration = typeof(ActivationWorkingSet).GetField("_nextGeneration", BindingFlags.Instance | BindingFlags.NonPublic) + ?? throw new InvalidOperationException("Could not find the working-set generation field."); + nextGeneration.SetValue(workingSet, long.MaxValue - 1); + var scanStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + using var resumeScan = new ManualResetEventSlim(); + var member = new TestWorkingSetMember(_ => + { + scanStarted.TrySetResult(); + resumeScan.Wait(); + return true; + }); + workingSet.OnActivated(member); + + var lifecycle = new SiloLifecycleSubject(NullLogger.Instance); + ((ILifecycleParticipant)workingSet).Participate(lifecycle); + await lifecycle.OnStart(); + await scanStarted.Task.WaitAsync(TimeSpan.FromSeconds(10)); + + try + { + workingSet.OnEvicted(member); + workingSet.OnActivated(member); + resumeScan.Set(); + await lifecycle.OnStop(); + } + finally + { + resumeScan.Set(); + } + Assert.Equal(1, workingSet.Count); Assert.Contains(member, workingSet.Members); } @@ -441,5 +662,10 @@ private IActivationWorkingSetMember PrepareActivation(TimeSpan collectionAgeLimi return (IActivationWorkingSetMember)activation; } + + private sealed class TestWorkingSetMember(Func? isCandidateForRemoval = null) : IActivationWorkingSetMember + { + public bool IsCandidateForRemoval(bool wouldRemove) => isCandidateForRemoval?.Invoke(wouldRemove) ?? false; + } } } From 71b75f47013ba227af888af0f28107db9adf7268 Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Sun, 23 Aug 2026 15:45:46 -0700 Subject: [PATCH 04/12] perf(runtime): store working set clock state on activations --- src/Orleans.Runtime/Catalog/ActivationData.cs | 56 ++++++-- .../Catalog/ActivationWorkingSet.cs | 127 ++++++++++-------- src/api/Orleans.Runtime/Orleans.Runtime.cs | 2 + .../Runtime/ActivationCollectorTests.cs | 105 ++++++++++++--- .../DeactivatedGrainQueueTests.cs | 1 + 5 files changed, 205 insertions(+), 86 deletions(-) diff --git a/src/Orleans.Runtime/Catalog/ActivationData.cs b/src/Orleans.Runtime/Catalog/ActivationData.cs index 0ce9b85710e..c42b11c7c45 100644 --- a/src/Orleans.Runtime/Catalog/ActivationData.cs +++ b/src/Orleans.Runtime/Catalog/ActivationData.cs @@ -41,6 +41,10 @@ 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 readonly GrainTypeSharedContext _shared; private readonly IServiceScope _serviceScope; private readonly WorkItemGroup _workItemGroup; @@ -50,7 +54,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 +150,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 +171,18 @@ 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); + } + // Currently, the only supported multi-activation grain is one using the StatelessWorkerPlacement strategy. internal bool IsStatelessWorker => PlacementStrategy is StatelessWorkerPlacement; @@ -404,7 +420,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 +989,21 @@ public TExtensionInterface GetExtension() bool IActivationWorkingSetMember.IsCandidateForRemoval(bool wouldRemove) { const int IdlenessLowerBound = 10_000; - lock (this) - { - var inactive = IsInactive && _idleDuration.ElapsedMilliseconds > IdlenessLowerBound; + Debug.Assert(Monitor.IsEntered(this)); + var inactive = IsInactive && _idleDuration.ElapsedMilliseconds > IdlenessLowerBound; + + // 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; + } - // 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 IActivationWorkingSetMember.IsIdle + { + get => IsIdleInWorkingSet; + set + { + Debug.Assert(Monitor.IsEntered(this)); + IsIdleInWorkingSet = value; } } @@ -1485,10 +1516,11 @@ private void OnCompletedRequest(Message message) if (message.IsKeepAlive) { _idleDuration = CoarseStopwatch.StartNew(); + IsIdleInWorkingSet = false; - if (!_isInWorkingSet) + if (!IsInWorkingSet) { - _isInWorkingSet = true; + IsInWorkingSet = true; _shared.InternalRuntime.ActivationWorkingSet.OnActive(this); } } @@ -2034,7 +2066,7 @@ private async Task FinishDeactivating(Command.Deactivate deactivateCommand, Canc deactivationMetrics = deactivationMetrics.Migration(); _shared.CatalogInstruments.ActivationShutdownViaMigration(); } - else if (_isInWorkingSet) + else if (IsInWorkingSet) { deactivationMetrics = deactivationMetrics.DeactivateOnIdle(); _shared.CatalogInstruments.ActivationShutdownViaDeactivateOnIdle(); diff --git a/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs b/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs index ce483e14023..2680be77132 100644 --- a/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs +++ b/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs @@ -17,17 +17,12 @@ namespace Orleans.Runtime; /// internal sealed partial class ActivationWorkingSet : IActivationWorkingSet, ILifecycleParticipant { - // The low bit stores idle state and the remaining 63 bits identify the membership generation. - // Comparing the complete state prevents stale scan entries from changing a member after removal and re-addition. - // A generation is reused only after 2^63 successful additions within one process lifetime. - private const long IdleMask = 1; - private readonly ConcurrentDictionary _members = new(); + private readonly ConcurrentDictionary _members = new(); private readonly ILogger _logger; private readonly IAsyncTimer _scanPeriodTimer; private readonly List _observers; private int _activeCount; - private long _nextGeneration; private Task? _runTask; public ActivationWorkingSet( @@ -51,7 +46,7 @@ private IEnumerable EnumerateActiveMembers() { foreach (var pair in _members) { - if (!IsIdle(pair.Value)) + if (!pair.Key.IsIdle) { yield return pair.Key; } @@ -61,35 +56,31 @@ private IEnumerable EnumerateActiveMembers() public void OnActivated(IActivationWorkingSetMember member) { Debug.Assert(member is not ICollectibleGrainContext collectible || collectible.IsValid); - if (_members.TryAdd(member, GetNextActiveState())) + lock (member) { - Interlocked.Increment(ref _activeCount); - foreach (var observer in _observers) + if (!_members.TryAdd(member, 0)) { - observer.OnAdded(member); + throw new InvalidOperationException($"Member {member} is already a member of the working set"); } - return; + member.IsIdle = false; } - throw new InvalidOperationException($"Member {member} is already a member of the working set"); + Interlocked.Increment(ref _activeCount); + foreach (var observer in _observers) + { + observer.OnAdded(member); + } } public void OnActive(IActivationWorkingSetMember member) { - while (true) + lock (member) { - if (_members.TryGetValue(member, out var state)) - { - if (!IsIdle(state) || _members.TryUpdate(member, GetActiveState(state), comparisonValue: state)) - { - break; - } - } - else if (_members.TryAdd(member, GetNextActiveState())) + member.IsIdle = false; + if (_members.TryAdd(member, 0)) { Interlocked.Increment(ref _activeCount); - break; } } @@ -101,15 +92,17 @@ public void OnActive(IActivationWorkingSetMember member) public void OnEvicted(IActivationWorkingSetMember member) { - if (_members.TryRemove(member, out _)) + bool removed; + lock (member) { - OnEvictedCore(member); + removed = _members.TryRemove(member, out _); + if (removed) + { + member.IsIdle = false; + } } - } - private void OnEvicted(IActivationWorkingSetMember member, long state) - { - if (_members.TryRemove(KeyValuePair.Create(member, state))) + if (removed) { OnEvictedCore(member); } @@ -150,7 +143,7 @@ private async Task MonitorWorkingSet() { try { - VisitMember(pair.Key, pair.Value); + VisitMember(pair.Key); } catch (Exception exception) { @@ -160,47 +153,66 @@ private async Task MonitorWorkingSet() } } - private void VisitMember(IActivationWorkingSetMember member, long state) + private void VisitMember(IActivationWorkingSetMember member) { - var wouldRemove = IsIdle(state); - if (member.IsCandidateForRemoval(wouldRemove)) + 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) { - if (wouldRemove) + var wouldRemove = member.IsIdle; + if (member.IsCandidateForRemoval(wouldRemove)) { - OnEvicted(member, state); - } - else - { - if (_members.TryUpdate(member, GetIdleState(state), comparisonValue: state)) + if (wouldRemove) { - foreach (var observer in _observers) + if (_members.TryRemove(member, out _)) + { + member.IsIdle = false; + Interlocked.Decrement(ref _activeCount); + result = MemberVisitResult.Evicted; + } + else { - observer.OnIdle(member); + result = MemberVisitResult.None; } } + else + { + member.IsIdle = true; + result = MemberVisitResult.Idle; + } } - } - else - { - if (wouldRemove) + else { - _members.TryUpdate(member, GetActiveState(state), comparisonValue: state); + member.IsIdle = false; + result = MemberVisitResult.Active; } + } - foreach (var observer in _observers) + foreach (var observer in _observers) + { + switch (result) { - observer.OnActive(member); + case MemberVisitResult.Active: + observer.OnActive(member); + break; + case MemberVisitResult.Idle: + observer.OnIdle(member); + break; + case MemberVisitResult.Evicted: + observer.OnEvicted(member); + break; } } } - private long GetNextActiveState() => unchecked(Interlocked.Increment(ref _nextGeneration) << 1); - - private static bool IsIdle(long state) => (state & IdleMask) != 0; - - private static long GetActiveState(long state) => state & ~IdleMask; - - private static long GetIdleState(long state) => state | IdleMask; + private enum MemberVisitResult + { + None, + Active, + Idle, + Evicted + } void ILifecycleParticipant.Participate(ISiloLifecycle lifecycle) { @@ -271,6 +283,11 @@ public interface IActivationWorkingSet /// public interface IActivationWorkingSetMember { + /// + /// Gets or sets whether this member was idle during the previous working-set scan. + /// + bool IsIdle { get; set; } + /// /// Returns if the member is eligible for removal, otherwise. /// diff --git a/src/api/Orleans.Runtime/Orleans.Runtime.cs b/src/api/Orleans.Runtime/Orleans.Runtime.cs index b8b8d347716..12aa1ab7215 100644 --- a/src/api/Orleans.Runtime/Orleans.Runtime.cs +++ b/src/api/Orleans.Runtime/Orleans.Runtime.cs @@ -748,6 +748,8 @@ public partial interface IActivationWorkingSet public partial interface IActivationWorkingSetMember { + bool IsIdle { get; set; } + bool IsCandidateForRemoval(bool wouldRemove); } diff --git a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs index bbf538fa826..8dab0a71736 100644 --- a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs +++ b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs @@ -1,5 +1,5 @@ using System.Collections.Concurrent; -using System.Reflection; +using System.Runtime.CompilerServices; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging.Abstractions; using Microsoft.Extensions.Options; @@ -395,11 +395,11 @@ public async Task WorkingSetScan_DoesNotUpdateReaddedMember() [observer], CreateCatalogInstruments(), TimeProvider.System); - var member = Substitute.For(); var scanStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); using var resumeScan = new ManualResetEventSlim(); - member.IsCandidateForRemoval(false).Returns(_ => + var member = new TestWorkingSetMember(wouldRemove => { + Assert.False(wouldRemove); scanStarted.TrySetResult(); resumeScan.Wait(); return true; @@ -411,13 +411,21 @@ public async Task WorkingSetScan_DoesNotUpdateReaddedMember() await lifecycle.OnStart(); await scanStarted.Task.WaitAsync(TimeSpan.FromSeconds(10)); - try + var mutationStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var mutationTask = Task.Run(() => { + mutationStarted.SetResult(); workingSet.OnEvicted(member); workingSet.OnActivated(member); + }); + await mutationStarted.Task.WaitAsync(TimeSpan.FromSeconds(10)); + try + { var stopTask = lifecycle.OnStop(); Assert.False(stopTask.IsCompleted); + Assert.False(mutationTask.IsCompleted); resumeScan.Set(); + await mutationTask; await stopTask; } finally @@ -429,7 +437,7 @@ public async Task WorkingSetScan_DoesNotUpdateReaddedMember() Assert.Contains(member, workingSet.Members); observer.Received(2).OnAdded(member); observer.Received(1).OnEvicted(member); - observer.DidNotReceive().OnIdle(member); + observer.Received(1).OnIdle(member); } [Fact, TestCategory("Activation")] @@ -466,13 +474,21 @@ public async Task WorkingSetScan_DoesNotRemoveReaddedMember() await lifecycle.OnStart(); await removalScanStarted.Task.WaitAsync(TimeSpan.FromSeconds(10)); - try + var mutationStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var mutationTask = Task.Run(() => { + mutationStarted.SetResult(); workingSet.OnEvicted(member); workingSet.OnActivated(member); + }); + await mutationStarted.Task.WaitAsync(TimeSpan.FromSeconds(10)); + try + { var stopTask = lifecycle.OnStop(); Assert.False(stopTask.IsCompleted); + Assert.False(mutationTask.IsCompleted); resumeScan.Set(); + await mutationTask; await stopTask; } finally @@ -593,26 +609,29 @@ public async Task WorkingSetMembers_EnumerationToleratesConcurrentRemoveAndReadd } [Fact, TestCategory("Activation")] - public async Task WorkingSetGenerationWrap_DoesNotAliasAdjacentMembership() + public async Task WorkingSetScan_SerializesRemovalWithReactivation() { var timer = Substitute.For(); - timer.NextTick().Returns(Task.FromResult(true), Task.FromResult(false)); + 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, - Array.Empty(), + [observer], CreateCatalogInstruments(), TimeProvider.System); - var nextGeneration = typeof(ActivationWorkingSet).GetField("_nextGeneration", BindingFlags.Instance | BindingFlags.NonPublic) - ?? throw new InvalidOperationException("Could not find the working-set generation field."); - nextGeneration.SetValue(workingSet, long.MaxValue - 1); - var scanStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var removalScanStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); using var resumeScan = new ManualResetEventSlim(); - var member = new TestWorkingSetMember(_ => + var member = new TestWorkingSetMember(wouldRemove => { - scanStarted.TrySetResult(); + if (!wouldRemove) + { + return true; + } + + removalScanStarted.TrySetResult(); resumeScan.Wait(); return true; }); @@ -621,13 +640,20 @@ public async Task WorkingSetGenerationWrap_DoesNotAliasAdjacentMembership() var lifecycle = new SiloLifecycleSubject(NullLogger.Instance); ((ILifecycleParticipant)workingSet).Participate(lifecycle); await lifecycle.OnStart(); - await scanStarted.Task.WaitAsync(TimeSpan.FromSeconds(10)); + await removalScanStarted.Task.WaitAsync(TimeSpan.FromSeconds(10)); + var reactivationStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var reactivationTask = Task.Run(() => + { + reactivationStarted.SetResult(); + workingSet.OnActive(member); + }); + await reactivationStarted.Task.WaitAsync(TimeSpan.FromSeconds(10)); try { - workingSet.OnEvicted(member); - workingSet.OnActivated(member); + Assert.False(reactivationTask.IsCompleted); resumeScan.Set(); + await reactivationTask; await lifecycle.OnStop(); } finally @@ -637,6 +663,30 @@ public async Task WorkingSetGenerationWrap_DoesNotAliasAdjacentMembership() 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 ActivationStatus_PreservesLifecycleStateAcrossWorkingSetTransitions() + { + var activation = (ActivationData)RuntimeHelpers.GetUninitializedObject(typeof(ActivationData)); + var workingSetState = (IActivationWorkingSetMember)activation; + + lock (activation) + { + activation.SetState(ActivationState.Valid); + workingSetState.IsIdle = true; + Assert.Equal(ActivationState.Valid, activation.State); + Assert.True(workingSetState.IsIdle); + + activation.SetState(ActivationState.Deactivating); + workingSetState.IsIdle = false; + Assert.Equal(ActivationState.Deactivating, activation.State); + Assert.False(workingSetState.IsIdle); + } } private IActivationWorkingSetMember PrepareActivation(int collectionAgeLimitMinutes, ActivationCollector collector) @@ -665,7 +715,24 @@ private IActivationWorkingSetMember PrepareActivation(TimeSpan collectionAgeLimi private sealed class TestWorkingSetMember(Func? isCandidateForRemoval = null) : IActivationWorkingSetMember { - public bool IsCandidateForRemoval(bool wouldRemove) => isCandidateForRemoval?.Invoke(wouldRemove) ?? false; + private bool _isIdle; + + public bool IsIdle + { + get => Volatile.Read(ref _isIdle); + set + { + Assert.True(Monitor.IsEntered(this)); + _isIdle = value; + } + } + + public bool IsCandidateForRemoval(bool wouldRemove) + { + Assert.True(Monitor.IsEntered(this)); + return isCandidateForRemoval?.Invoke(wouldRemove) ?? false; + } + } } } diff --git a/test/Orleans.Placement.Tests/ActivationRepartitioningTests/DeactivatedGrainQueueTests.cs b/test/Orleans.Placement.Tests/ActivationRepartitioningTests/DeactivatedGrainQueueTests.cs index 4d040fbe0c0..a38d028528f 100644 --- a/test/Orleans.Placement.Tests/ActivationRepartitioningTests/DeactivatedGrainQueueTests.cs +++ b/test/Orleans.Placement.Tests/ActivationRepartitioningTests/DeactivatedGrainQueueTests.cs @@ -89,6 +89,7 @@ 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 IsIdle { get; set; } public bool IsCandidateForRemoval(bool wouldRemove) => false; public void SetComponent(TComponent? value) where TComponent : class => throw new NotImplementedException(); From a166295fb7568c8c1b79308a2583cdb34b0d2a9b Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Sun, 23 Aug 2026 16:13:22 -0700 Subject: [PATCH 05/12] fix(runtime): skip removed working set members --- src/Orleans.Runtime/Catalog/ActivationData.cs | 26 +++++++-- .../Catalog/ActivationWorkingSet.cs | 46 ++++++++++------ src/api/Orleans.Runtime/Orleans.Runtime.cs | 2 + .../Runtime/ActivationCollectorTests.cs | 53 +++++++++++++++++++ .../DeactivatedGrainQueueTests.cs | 1 + 5 files changed, 110 insertions(+), 18 deletions(-) diff --git a/src/Orleans.Runtime/Catalog/ActivationData.cs b/src/Orleans.Runtime/Catalog/ActivationData.cs index c42b11c7c45..93060b17d6f 100644 --- a/src/Orleans.Runtime/Catalog/ActivationData.cs +++ b/src/Orleans.Runtime/Catalog/ActivationData.cs @@ -45,6 +45,7 @@ internal sealed partial class ActivationData : 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; @@ -183,6 +184,12 @@ private bool IsIdleInWorkingSet 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; @@ -992,11 +999,24 @@ bool IActivationWorkingSetMember.IsCandidateForRemoval(bool wouldRemove) Debug.Assert(Monitor.IsEntered(this)); var inactive = IsInactive && _idleDuration.ElapsedMilliseconds > IdlenessLowerBound; - // This instance will remain in the working set if it is either not pending removal or if it is currently active. - IsInWorkingSet = !wouldRemove || !inactive; + WasRemovedByCollection = wouldRemove && inactive; return inactive; } + bool IActivationWorkingSetMember.IsInWorkingSet + { + get => IsInWorkingSet; + set + { + Debug.Assert(Monitor.IsEntered(this)); + IsInWorkingSet = value; + if (value) + { + WasRemovedByCollection = false; + } + } + } + bool IActivationWorkingSetMember.IsIdle { get => IsIdleInWorkingSet; @@ -2066,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 2680be77132..a87f49e2a49 100644 --- a/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs +++ b/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs @@ -63,6 +63,7 @@ public void OnActivated(IActivationWorkingSetMember member) throw new InvalidOperationException($"Member {member} is already a member of the working set"); } + member.IsInWorkingSet = true; member.IsIdle = false; } @@ -77,6 +78,7 @@ public void OnActive(IActivationWorkingSetMember member) { lock (member) { + member.IsInWorkingSet = true; member.IsIdle = false; if (_members.TryAdd(member, 0)) { @@ -98,6 +100,7 @@ public void OnEvicted(IActivationWorkingSetMember member) removed = _members.TryRemove(member, out _); if (removed) { + member.IsInWorkingSet = false; member.IsIdle = false; } } @@ -160,33 +163,41 @@ private void VisitMember(IActivationWorkingSetMember member) // member's current state while holding its lock instead of adding a dictionary validation to every scan. lock (member) { - var wouldRemove = member.IsIdle; - if (member.IsCandidateForRemoval(wouldRemove)) + if (!member.IsInWorkingSet) { - if (wouldRemove) + result = MemberVisitResult.None; + } + else + { + var wouldRemove = member.IsIdle; + if (member.IsCandidateForRemoval(wouldRemove)) { - if (_members.TryRemove(member, out _)) + if (wouldRemove) { - member.IsIdle = false; - Interlocked.Decrement(ref _activeCount); - result = MemberVisitResult.Evicted; + if (_members.TryRemove(member, out _)) + { + member.IsInWorkingSet = false; + member.IsIdle = false; + Interlocked.Decrement(ref _activeCount); + result = MemberVisitResult.Evicted; + } + else + { + result = MemberVisitResult.None; + } } else { - result = MemberVisitResult.None; + member.IsIdle = true; + result = MemberVisitResult.Idle; } } else { - member.IsIdle = true; - result = MemberVisitResult.Idle; + member.IsIdle = false; + result = MemberVisitResult.Active; } } - else - { - member.IsIdle = false; - result = MemberVisitResult.Active; - } } foreach (var observer in _observers) @@ -283,6 +294,11 @@ public interface IActivationWorkingSet /// public interface IActivationWorkingSetMember { + /// + /// Gets or sets whether this member is registered in the working set. + /// + bool IsInWorkingSet { get; set; } + /// /// Gets or sets whether this member was idle during the previous working-set scan. /// diff --git a/src/api/Orleans.Runtime/Orleans.Runtime.cs b/src/api/Orleans.Runtime/Orleans.Runtime.cs index 12aa1ab7215..76757ea9106 100644 --- a/src/api/Orleans.Runtime/Orleans.Runtime.cs +++ b/src/api/Orleans.Runtime/Orleans.Runtime.cs @@ -750,6 +750,8 @@ public partial interface IActivationWorkingSetMember { bool IsIdle { get; set; } + bool IsInWorkingSet { get; set; } + bool IsCandidateForRemoval(bool wouldRemove); } diff --git a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs index 8dab0a71736..8a00ceb0aac 100644 --- a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs +++ b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs @@ -1,4 +1,5 @@ using System.Collections.Concurrent; +using System.Reflection; using System.Runtime.CompilerServices; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging.Abstractions; @@ -669,6 +670,43 @@ public async Task WorkingSetScan_SerializesRemovalWithReactivation() 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 ActivationStatus_PreservesLifecycleStateAcrossWorkingSetTransitions() { @@ -678,13 +716,17 @@ public void ActivationStatus_PreservesLifecycleStateAcrossWorkingSetTransitions( lock (activation) { activation.SetState(ActivationState.Valid); + workingSetState.IsInWorkingSet = true; workingSetState.IsIdle = true; Assert.Equal(ActivationState.Valid, activation.State); + Assert.True(workingSetState.IsInWorkingSet); Assert.True(workingSetState.IsIdle); activation.SetState(ActivationState.Deactivating); + workingSetState.IsInWorkingSet = false; workingSetState.IsIdle = false; Assert.Equal(ActivationState.Deactivating, activation.State); + Assert.False(workingSetState.IsInWorkingSet); Assert.False(workingSetState.IsIdle); } } @@ -716,6 +758,7 @@ private IActivationWorkingSetMember PrepareActivation(TimeSpan collectionAgeLimi private sealed class TestWorkingSetMember(Func? isCandidateForRemoval = null) : IActivationWorkingSetMember { private bool _isIdle; + private bool _isInWorkingSet; public bool IsIdle { @@ -727,6 +770,16 @@ public bool IsIdle } } + 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)); diff --git a/test/Orleans.Placement.Tests/ActivationRepartitioningTests/DeactivatedGrainQueueTests.cs b/test/Orleans.Placement.Tests/ActivationRepartitioningTests/DeactivatedGrainQueueTests.cs index a38d028528f..dd41f42591e 100644 --- a/test/Orleans.Placement.Tests/ActivationRepartitioningTests/DeactivatedGrainQueueTests.cs +++ b/test/Orleans.Placement.Tests/ActivationRepartitioningTests/DeactivatedGrainQueueTests.cs @@ -90,6 +90,7 @@ private sealed class TestActivationWorkingSetMember(GrainId grainId) : IActivati public IWorkItemScheduler Scheduler => throw new NotImplementedException(); public Task Deactivated => Task.CompletedTask; public bool IsIdle { get; set; } + public bool IsInWorkingSet { get; set; } public bool IsCandidateForRemoval(bool wouldRemove) => false; public void SetComponent(TComponent? value) where TComponent : class => throw new NotImplementedException(); From 0cbd1314d918b3ab15feb4ea140c08c3fd931095 Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Sun, 23 Aug 2026 19:48:19 -0700 Subject: [PATCH 06/12] test(runtime): expand activation working set coverage --- .../Runtime/ActivationCollectorTests.cs | 1311 ++++++++++++++++- 1 file changed, 1287 insertions(+), 24 deletions(-) diff --git a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs index 8a00ceb0aac..441c3cd1344 100644 --- a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs +++ b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs @@ -1,4 +1,5 @@ using System.Collections.Concurrent; +using System.Diagnostics.CodeAnalysis; using System.Reflection; using System.Runtime.CompilerServices; using Microsoft.Extensions.DependencyInjection; @@ -707,30 +708,6 @@ public void WorkingSetScan_SkipsRemovedMember() observer.DidNotReceive().OnActive(member); } - [Fact, TestCategory("Activation")] - public void ActivationStatus_PreservesLifecycleStateAcrossWorkingSetTransitions() - { - var activation = (ActivationData)RuntimeHelpers.GetUninitializedObject(typeof(ActivationData)); - var workingSetState = (IActivationWorkingSetMember)activation; - - lock (activation) - { - activation.SetState(ActivationState.Valid); - workingSetState.IsInWorkingSet = true; - workingSetState.IsIdle = true; - Assert.Equal(ActivationState.Valid, activation.State); - Assert.True(workingSetState.IsInWorkingSet); - Assert.True(workingSetState.IsIdle); - - activation.SetState(ActivationState.Deactivating); - workingSetState.IsInWorkingSet = false; - workingSetState.IsIdle = false; - Assert.Equal(ActivationState.Deactivating, activation.State); - Assert.False(workingSetState.IsInWorkingSet); - Assert.False(workingSetState.IsIdle); - } - } - private IActivationWorkingSetMember PrepareActivation(int collectionAgeLimitMinutes, ActivationCollector collector) => PrepareActivation(TimeSpan.FromMinutes(collectionAgeLimitMinutes), collector); @@ -787,5 +764,1291 @@ public bool IsCandidateForRemoval(bool wouldRemove) } } + + [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_DuplicateOnActivated_ThrowsWithoutChangingCommittedState() + { + await using var harness = await WorkingSetHarness.CreateAsync(); + var member = harness.Members[0]; + harness.WorkingSet.OnActivated(member); + var expectedMembers = harness.WorkingSet.Members.ToArray(); + var expectedHistory = harness.Observer.GetHistory(0); + + var exception = Assert.Throws(() => harness.WorkingSet.OnActivated(member)); + + Assert.Contains("already a member of the working set", exception.Message); + Assert.Equal(1, harness.WorkingSet.Count); + Assert.True(member.IsInWorkingSet); + Assert.False(member.IsIdle); + Assert.Equal(expectedMembers, harness.WorkingSet.Members); + Assert.Equal(["Added"], expectedHistory); + Assert.Equal(expectedHistory, harness.Observer.GetHistory(0)); + } + + [Fact, TestCategory("Activation")] + public async Task WorkingSet_DirectObserverFailure_PropagatesAfterTransitionIsCommitted() + { + await using var harness = await WorkingSetHarness.CreateAsync(); + var member = harness.Members[0]; + var expectedException = new InvalidOperationException("observer failure"); + harness.Observer.AddedException = expectedException; + + var actualException = Record.Exception(() => harness.WorkingSet.OnActivated(member)); + + Assert.Same(expectedException, actualException); + Assert.Equal(1, harness.WorkingSet.Count); + Assert.True(member.IsInWorkingSet); + Assert.False(member.IsIdle); + Assert.Same(member, Assert.Single(harness.WorkingSet.Members)); + Assert.Equal(["Added"], harness.Observer.GetHistory(0)); + } + + [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); + })).ToArray(); + + try + { + Assert.True(ready.Wait(TimeSpan.FromSeconds(10)), "Workers did not reach the OnActive start gate."); + start.TrySetResult(); + await Task.WhenAll(workers).WaitAsync(TimeSpan.FromSeconds(10)); + } + finally + { + start.TrySetResult(); + await Task.WhenAll(workers).WaitAsync(TimeSpan.FromSeconds(10)); + } + + 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_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)); + Task? reAdd = null; + try + { + await gate.Entered.Task.WaitAsync(TimeSpan.FromSeconds(10)); + reAdd = Task.Run(() => harness.WorkingSet.OnActive(member)); + await reAdd.WaitAsync(TimeSpan.FromSeconds(10)); + 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)); + } + + await eviction.WaitAsync(TimeSpan.FromSeconds(10)); + } + + 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)); + } + + [Fact, TestCategory("Activation")] + public async Task WorkingSet_MonitorContinuesAfterMemberCandidateThrows() + { + await using var harness = await WorkingSetHarness.CreateAsync(); + harness.MemberStates[0].CandidateException = new InvalidOperationException("candidate failure"); + harness.MemberStates[1].CandidateEligible = true; + harness.WorkingSet.OnActivated(harness.Members[0]); + harness.WorkingSet.OnActivated(harness.Members[1]); + harness.WorkingSet.OnActivated(harness.Members[2]); + + await harness.ScanOnceAsync(); + + Assert.Equal(2, harness.TimerGeneration); + Assert.Equal(3, harness.WorkingSet.Count); + Assert.True(harness.Members[0].IsInWorkingSet); + Assert.False(harness.Members[0].IsIdle); + Assert.True(harness.Members[1].IsInWorkingSet); + Assert.True(harness.Members[1].IsIdle); + Assert.True(harness.Members[2].IsInWorkingSet); + Assert.False(harness.Members[2].IsIdle); + Assert.False(harness.Members[3].IsInWorkingSet); + Assert.False(harness.Members[3].IsIdle); + Assert.Equal([0, 2], harness.WorkingSet.Members.Select(harness.GetMemberId).Order()); + Assert.Equal([false], harness.MemberStates[0].GetCandidateCalls()); + Assert.Equal([false], harness.MemberStates[1].GetCandidateCalls()); + Assert.Equal([false], harness.MemberStates[2].GetCandidateCalls()); + Assert.Empty(harness.MemberStates[3].GetCandidateCalls()); + Assert.Equal(["Added"], harness.Observer.GetHistory(0)); + Assert.Equal(["Added", "Idle"], harness.Observer.GetHistory(1)); + Assert.Equal(["Added", "Active"], harness.Observer.GetHistory(2)); + Assert.Empty(harness.Observer.GetHistory(3)); + } + + [Theory, TestCategory("Activation")] + [InlineData(0x0000C0DE)] + [InlineData(0x0013579B)] + [InlineData(0x02468ACE)] + public async Task WorkingSet_SeededConcurrentOperations_PreserveTerminalInvariants(int seed) + { + await using var harness = await WorkingSetHarness.CreateAsync(); + var runner = new SeededWorkingSetStressRunner(seed, harness); + + await runner.RunAsync(); + } + + 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 Exception? _candidateException; + private int _candidateEligible; + + public Exception? CandidateException + { + get => Volatile.Read(ref _candidateException); + set => Volatile.Write(ref _candidateException, value); + } + + 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); + if (CandidateException is { } exception) + { + throw exception; + } + + 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 Exception? AddedException { get; set; } + + public void OnAdded(IActivationWorkingSetMember member) + { + Record(member, "Added"); + if (AddedException is { } exception) + { + throw exception; + } + } + + 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(); + 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().WaitAsync(TimeSpan.FromSeconds(10)); + } + 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)); + } + } + + 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); + } + } + + private sealed class SeededWorkingSetStressRunner(int seed, WorkingSetHarness harness) + { + private const int WorkerCount = 4; + private const int PhaseCount = 32; + private readonly ConcurrentQueue _trace = new(); + private readonly ConcurrentQueue _exceptions = new(); + private readonly Random _random = new(seed); + private int _currentPhase = -1; + private int _outcomeCount; + private int _selectedCount; + + public async Task RunAsync() + { + try + { + for (var phase = 0; phase < PhaseCount; phase++) + { + _currentPhase = phase; + var operations = Enumerable.Range(0, WorkerCount) + .Select(worker => CreateOperation(phase, worker)) + .ToArray(); + foreach (var operation in operations) + { + Record(operation.Worker, phase, operation.Operation, "selected"); + _selectedCount++; + } + + using var ready = new CountdownEvent(WorkerCount); + var start = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var workers = operations + .Select(operation => Task.Run(() => RunWorkerAsync(operation, ready, start.Task))) + .ToArray(); + + try + { + Assert.True( + ready.Wait(TimeSpan.FromSeconds(10)), + $"seed=0x{seed:X8}; worker=coordinator; phase={phase}; start gate timed out."); + start.TrySetResult(); + await Task.WhenAll(workers).WaitAsync(TimeSpan.FromSeconds(10)); + } + finally + { + start.TrySetResult(); + await Task.WhenAll(workers).WaitAsync(TimeSpan.FromSeconds(10)); + } + + Assert.Empty(_exceptions); + AssertTerminalInvariants(); + Assert.Equal(_selectedCount, _outcomeCount); + } + + Assert.Equal(WorkerCount * PhaseCount, _selectedCount); + Assert.Equal(WorkerCount * PhaseCount, _outcomeCount); + } + catch (Exception exception) + { + throw new InvalidOperationException( + $"seed=0x{seed:X8}; worker=all; phase={_currentPhase}; failure={exception.Message}" + + $"{Environment.NewLine}Full trace:{Environment.NewLine}{string.Join(Environment.NewLine, _trace)}", + exception); + } + } + + private (int Worker, WorkingSetOperation Operation) CreateOperation(int phase, int worker) + { + var kind = _random.Next(4) switch + { + 0 => WorkingSetOperationKind.Active, + 1 => WorkingSetOperationKind.Evict, + 2 => WorkingSetOperationKind.Deactivating, + _ => WorkingSetOperationKind.Deactivated + }; + var memberId = _random.Next(harness.Members.Count); + return (worker, new WorkingSetOperation(phase * WorkerCount + worker, kind, memberId, false)); + } + + private async Task RunWorkerAsync( + (int Worker, WorkingSetOperation Operation) work, + CountdownEvent ready, + Task start) + { + try + { + ready.Signal(); + await start; + ExecuteStressOperation(work.Operation); + Record(work.Worker, _currentPhase, work.Operation, "completed"); + } + catch (Exception exception) + { + _exceptions.Enqueue(exception); + Record(work.Worker, _currentPhase, work.Operation, $"exception={exception.GetType().Name}:{exception.Message}"); + } + finally + { + Interlocked.Increment(ref _outcomeCount); + } + } + + private void ExecuteStressOperation(WorkingSetOperation operation) + { + var member = harness.Members[operation.MemberId]; + switch (operation.Kind) + { + case WorkingSetOperationKind.Active: + harness.WorkingSet.OnActive(member); + 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; + default: + throw new ArgumentOutOfRangeException(nameof(operation)); + } + } + + private void AssertTerminalInvariants() + { + Assert.InRange(harness.WorkingSet.Count, 0, harness.Members.Count); + var inSetIds = harness.Members + .Select(static (member, id) => (member, id)) + .Where(static item => item.member.IsInWorkingSet) + .Select(static item => item.id) + .Order() + .ToArray(); + Assert.Equal(inSetIds.Length, harness.WorkingSet.Count); + Assert.All( + harness.Members, + static member => Assert.False(member.IsIdle && !member.IsInWorkingSet)); + + var visibleIds = harness.WorkingSet.Members + .Select(harness.GetMemberId) + .Order() + .ToArray(); + Assert.Equal(visibleIds.Length, visibleIds.Distinct().Count()); + var expectedVisibleIds = harness.Members + .Select(static (member, id) => (member, id)) + .Where(static item => item.member.IsInWorkingSet && !item.member.IsIdle) + .Select(static item => item.id) + .Order() + .ToArray(); + Assert.Equal(expectedVisibleIds, visibleIds); + } + + private void Record(int worker, int phase, WorkingSetOperation operation, string result) + { + _trace.Enqueue( + $"seed=0x{seed:X8}; worker={worker}; phase={phase}; operation={operation}; result={result}"); + } + } + + [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 async Task ActivationCollector_AgeCollection_RequestsActivationIdle() + { + var timeProvider = new FakeTimeProvider(DateTimeOffset.Parse("2025-01-01T00:00:00.000+00:00")); + using var serviceProvider = CreateCatalogServiceProvider(); + var collector = new ActivationCollector( + timeProvider, + Options.Create(new GrainCollectionOptions()), + NullLogger.Instance, + Substitute.For(), + serviceProvider.GetRequiredService()); + var activation = CreateCollectorActivation(TimeSpan.FromMinutes(5)); + activation.GetIdleness().Returns(TimeSpan.FromMinutes(2)); + collector.ScheduleCollection(activation, TimeSpan.FromMinutes(5), timeProvider.GetUtcNow().UtcDateTime); + + await collector.CollectActivations(TimeSpan.FromMinutes(1), CancellationToken.None); + + activation.Received(1).Deactivate( + Arg.Is(reason => reason.ReasonCode == DeactivationReasonCode.ActivationIdle), + Arg.Any()); + activation.Received(1).Deactivate(Arg.Any(), Arg.Any()); + activation.DidNotReceive().Deactivate( + Arg.Is(reason => reason.ReasonCode == DeactivationReasonCode.HighMemoryPressure), + Arg.Any()); + } + + [Fact, TestCategory("Activation")] + public async Task ActivationCollector_HighMemoryCollection_RequestsHighMemoryPressure() + { + var timeProvider = new FakeTimeProvider(DateTimeOffset.Parse("2025-01-01T00:00:00.000+00:00")); + using var serviceProvider = CreateCatalogServiceProvider(); + var collector = new ActivationCollector( + timeProvider, + Options.Create(new GrainCollectionOptions()), + NullLogger.Instance, + Substitute.For(), + serviceProvider.GetRequiredService()); + var activation = CreateCollectorActivation(TimeSpan.FromMinutes(5)); + collector.ScheduleCollection(activation, TimeSpan.FromMinutes(5), timeProvider.GetUtcNow().UtcDateTime); + collector._activationCount = 1; + + await collector.DeactivateInDueTimeOrder(1, CancellationToken.None); + + activation.Received(1).Deactivate( + Arg.Is(reason => reason.ReasonCode == DeactivationReasonCode.HighMemoryPressure), + Arg.Any()); + activation.Received(1).Deactivate(Arg.Any(), Arg.Any()); + activation.DidNotReceive().Deactivate( + Arg.Is(reason => reason.ReasonCode == DeactivationReasonCode.ActivationIdle), + Arg.Any()); + } + + 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 static ICollectibleGrainContext CreateCollectorActivation(TimeSpan collectionAgeLimit) + { + var activation = Substitute.For(); + activation.CollectionAgeLimit.Returns(collectionAgeLimit); + activation.IsValid.Returns(true); + activation.IsExemptFromCollection.Returns(false); + activation.IsInactive.Returns(true); + activation.Deactivated.Returns(Task.CompletedTask); + return activation; + } + + private static ServiceProvider CreateCatalogServiceProvider() + { + var services = new ServiceCollection(); + services.AddMetrics(); + services.AddSingleton(); + services.AddSingleton(); + return services.BuildServiceProvider(); + } + + 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 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 IActivationWorkingSetMember 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 Dispose() + { + _timer.Dispose(); + _activationScope.Dispose(); + _serviceProvider.Dispose(); + } + } } } From 390fbd8d4555eb0751f4c851d440fe80d5d789c5 Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Sun, 23 Aug 2026 19:59:19 -0700 Subject: [PATCH 07/12] fix(runtime): avoid reactivating deactivating grains --- src/Orleans.Runtime/Catalog/ActivationData.cs | 2 +- .../Runtime/ActivationCollectorTests.cs | 27 +++++++++++++++++++ 2 files changed, 28 insertions(+), 1 deletion(-) diff --git a/src/Orleans.Runtime/Catalog/ActivationData.cs b/src/Orleans.Runtime/Catalog/ActivationData.cs index 93060b17d6f..10d035a6c68 100644 --- a/src/Orleans.Runtime/Catalog/ActivationData.cs +++ b/src/Orleans.Runtime/Catalog/ActivationData.cs @@ -1533,7 +1533,7 @@ 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; diff --git a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs index 441c3cd1344..a578a29d87d 100644 --- a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs +++ b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs @@ -1853,6 +1853,26 @@ public void ActivationData_ExplicitWorkingSetDeactivation_DoesNotSetCollectionMa 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)); + } + [Fact, TestCategory("Activation")] public async Task ActivationCollector_AgeCollection_RequestsActivationIdle() { @@ -1959,6 +1979,11 @@ private sealed class ActivationDataWorkingSetFixture : IDisposable 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; @@ -2043,6 +2068,8 @@ public void AdvanceIdleDurationTo(long elapsedMilliseconds) public void ScanOnce() => VisitMemberMethod.Invoke(WorkingSet, [Member]); + public void CompleteRequest(Message message) => CompleteRequestMethod.Invoke(Activation, [message]); + public void Dispose() { _timer.Dispose(); From 8643b4668256ec8066aacdcda6c61db9800fc812 Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Sun, 23 Aug 2026 21:05:45 -0700 Subject: [PATCH 08/12] test(runtime): trim activation working set coverage --- .../Runtime/ActivationCollectorTests.cs | 335 +----------------- 1 file changed, 1 insertion(+), 334 deletions(-) diff --git a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs index a578a29d87d..8497a538676 100644 --- a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs +++ b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs @@ -792,44 +792,6 @@ public void WorkingSet_SequentialGeneratedTrace_MatchesReferenceModel() print: FormatWorkingSetTrace); } - [Fact, TestCategory("Activation")] - public async Task WorkingSet_DuplicateOnActivated_ThrowsWithoutChangingCommittedState() - { - await using var harness = await WorkingSetHarness.CreateAsync(); - var member = harness.Members[0]; - harness.WorkingSet.OnActivated(member); - var expectedMembers = harness.WorkingSet.Members.ToArray(); - var expectedHistory = harness.Observer.GetHistory(0); - - var exception = Assert.Throws(() => harness.WorkingSet.OnActivated(member)); - - Assert.Contains("already a member of the working set", exception.Message); - Assert.Equal(1, harness.WorkingSet.Count); - Assert.True(member.IsInWorkingSet); - Assert.False(member.IsIdle); - Assert.Equal(expectedMembers, harness.WorkingSet.Members); - Assert.Equal(["Added"], expectedHistory); - Assert.Equal(expectedHistory, harness.Observer.GetHistory(0)); - } - - [Fact, TestCategory("Activation")] - public async Task WorkingSet_DirectObserverFailure_PropagatesAfterTransitionIsCommitted() - { - await using var harness = await WorkingSetHarness.CreateAsync(); - var member = harness.Members[0]; - var expectedException = new InvalidOperationException("observer failure"); - harness.Observer.AddedException = expectedException; - - var actualException = Record.Exception(() => harness.WorkingSet.OnActivated(member)); - - Assert.Same(expectedException, actualException); - Assert.Equal(1, harness.WorkingSet.Count); - Assert.True(member.IsInWorkingSet); - Assert.False(member.IsIdle); - Assert.Same(member, Assert.Single(harness.WorkingSet.Members)); - Assert.Equal(["Added"], harness.Observer.GetHistory(0)); - } - [Fact, TestCategory("Activation")] public async Task WorkingSet_ConcurrentOnActiveForAbsentMember_AddsOnceAndNotifiesEveryCaller() { @@ -910,51 +872,6 @@ public async Task WorkingSet_EvictionCallbackCanOverlapReAddWithoutHoldingMember Assert.Same(member, Assert.Single(harness.WorkingSet.Members)); } - [Fact, TestCategory("Activation")] - public async Task WorkingSet_MonitorContinuesAfterMemberCandidateThrows() - { - await using var harness = await WorkingSetHarness.CreateAsync(); - harness.MemberStates[0].CandidateException = new InvalidOperationException("candidate failure"); - harness.MemberStates[1].CandidateEligible = true; - harness.WorkingSet.OnActivated(harness.Members[0]); - harness.WorkingSet.OnActivated(harness.Members[1]); - harness.WorkingSet.OnActivated(harness.Members[2]); - - await harness.ScanOnceAsync(); - - Assert.Equal(2, harness.TimerGeneration); - Assert.Equal(3, harness.WorkingSet.Count); - Assert.True(harness.Members[0].IsInWorkingSet); - Assert.False(harness.Members[0].IsIdle); - Assert.True(harness.Members[1].IsInWorkingSet); - Assert.True(harness.Members[1].IsIdle); - Assert.True(harness.Members[2].IsInWorkingSet); - Assert.False(harness.Members[2].IsIdle); - Assert.False(harness.Members[3].IsInWorkingSet); - Assert.False(harness.Members[3].IsIdle); - Assert.Equal([0, 2], harness.WorkingSet.Members.Select(harness.GetMemberId).Order()); - Assert.Equal([false], harness.MemberStates[0].GetCandidateCalls()); - Assert.Equal([false], harness.MemberStates[1].GetCandidateCalls()); - Assert.Equal([false], harness.MemberStates[2].GetCandidateCalls()); - Assert.Empty(harness.MemberStates[3].GetCandidateCalls()); - Assert.Equal(["Added"], harness.Observer.GetHistory(0)); - Assert.Equal(["Added", "Idle"], harness.Observer.GetHistory(1)); - Assert.Equal(["Added", "Active"], harness.Observer.GetHistory(2)); - Assert.Empty(harness.Observer.GetHistory(3)); - } - - [Theory, TestCategory("Activation")] - [InlineData(0x0000C0DE)] - [InlineData(0x0013579B)] - [InlineData(0x02468ACE)] - public async Task WorkingSet_SeededConcurrentOperations_PreserveTerminalInvariants(int seed) - { - await using var harness = await WorkingSetHarness.CreateAsync(); - var runner = new SeededWorkingSetStressRunner(seed, harness); - - await runner.RunAsync(); - } - private static WorkingSetOperation[] GetWorkingSetCoverageSpine() => [ new(0, WorkingSetOperationKind.Activate, 0, false), @@ -1075,15 +992,8 @@ public override string ToString() private sealed class WorkingSetMemberState { private readonly ConcurrentQueue _candidateCalls = new(); - private Exception? _candidateException; private int _candidateEligible; - public Exception? CandidateException - { - get => Volatile.Read(ref _candidateException); - set => Volatile.Write(ref _candidateException, value); - } - public bool CandidateEligible { get => Volatile.Read(ref _candidateEligible) != 0; @@ -1093,11 +1003,6 @@ public bool CandidateEligible public bool IsCandidateForRemoval(bool wouldRemove) { _candidateCalls.Enqueue(wouldRemove); - if (CandidateException is { } exception) - { - throw exception; - } - return CandidateEligible; } @@ -1115,16 +1020,7 @@ private sealed class RecordingWorkingSetObserver( private TaskCompletionSource? _evictionEntered; private TaskCompletionSource? _evictionRelease; - public Exception? AddedException { get; set; } - - public void OnAdded(IActivationWorkingSetMember member) - { - Record(member, "Added"); - if (AddedException is { } exception) - { - throw exception; - } - } + public void OnAdded(IActivationWorkingSetMember member) => Record(member, "Added"); public void OnActive(IActivationWorkingSetMember member) => Record(member, "Active"); @@ -1511,163 +1407,6 @@ public void Dispose() } } - private sealed class SeededWorkingSetStressRunner(int seed, WorkingSetHarness harness) - { - private const int WorkerCount = 4; - private const int PhaseCount = 32; - private readonly ConcurrentQueue _trace = new(); - private readonly ConcurrentQueue _exceptions = new(); - private readonly Random _random = new(seed); - private int _currentPhase = -1; - private int _outcomeCount; - private int _selectedCount; - - public async Task RunAsync() - { - try - { - for (var phase = 0; phase < PhaseCount; phase++) - { - _currentPhase = phase; - var operations = Enumerable.Range(0, WorkerCount) - .Select(worker => CreateOperation(phase, worker)) - .ToArray(); - foreach (var operation in operations) - { - Record(operation.Worker, phase, operation.Operation, "selected"); - _selectedCount++; - } - - using var ready = new CountdownEvent(WorkerCount); - var start = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); - var workers = operations - .Select(operation => Task.Run(() => RunWorkerAsync(operation, ready, start.Task))) - .ToArray(); - - try - { - Assert.True( - ready.Wait(TimeSpan.FromSeconds(10)), - $"seed=0x{seed:X8}; worker=coordinator; phase={phase}; start gate timed out."); - start.TrySetResult(); - await Task.WhenAll(workers).WaitAsync(TimeSpan.FromSeconds(10)); - } - finally - { - start.TrySetResult(); - await Task.WhenAll(workers).WaitAsync(TimeSpan.FromSeconds(10)); - } - - Assert.Empty(_exceptions); - AssertTerminalInvariants(); - Assert.Equal(_selectedCount, _outcomeCount); - } - - Assert.Equal(WorkerCount * PhaseCount, _selectedCount); - Assert.Equal(WorkerCount * PhaseCount, _outcomeCount); - } - catch (Exception exception) - { - throw new InvalidOperationException( - $"seed=0x{seed:X8}; worker=all; phase={_currentPhase}; failure={exception.Message}" - + $"{Environment.NewLine}Full trace:{Environment.NewLine}{string.Join(Environment.NewLine, _trace)}", - exception); - } - } - - private (int Worker, WorkingSetOperation Operation) CreateOperation(int phase, int worker) - { - var kind = _random.Next(4) switch - { - 0 => WorkingSetOperationKind.Active, - 1 => WorkingSetOperationKind.Evict, - 2 => WorkingSetOperationKind.Deactivating, - _ => WorkingSetOperationKind.Deactivated - }; - var memberId = _random.Next(harness.Members.Count); - return (worker, new WorkingSetOperation(phase * WorkerCount + worker, kind, memberId, false)); - } - - private async Task RunWorkerAsync( - (int Worker, WorkingSetOperation Operation) work, - CountdownEvent ready, - Task start) - { - try - { - ready.Signal(); - await start; - ExecuteStressOperation(work.Operation); - Record(work.Worker, _currentPhase, work.Operation, "completed"); - } - catch (Exception exception) - { - _exceptions.Enqueue(exception); - Record(work.Worker, _currentPhase, work.Operation, $"exception={exception.GetType().Name}:{exception.Message}"); - } - finally - { - Interlocked.Increment(ref _outcomeCount); - } - } - - private void ExecuteStressOperation(WorkingSetOperation operation) - { - var member = harness.Members[operation.MemberId]; - switch (operation.Kind) - { - case WorkingSetOperationKind.Active: - harness.WorkingSet.OnActive(member); - 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; - default: - throw new ArgumentOutOfRangeException(nameof(operation)); - } - } - - private void AssertTerminalInvariants() - { - Assert.InRange(harness.WorkingSet.Count, 0, harness.Members.Count); - var inSetIds = harness.Members - .Select(static (member, id) => (member, id)) - .Where(static item => item.member.IsInWorkingSet) - .Select(static item => item.id) - .Order() - .ToArray(); - Assert.Equal(inSetIds.Length, harness.WorkingSet.Count); - Assert.All( - harness.Members, - static member => Assert.False(member.IsIdle && !member.IsInWorkingSet)); - - var visibleIds = harness.WorkingSet.Members - .Select(harness.GetMemberId) - .Order() - .ToArray(); - Assert.Equal(visibleIds.Length, visibleIds.Distinct().Count()); - var expectedVisibleIds = harness.Members - .Select(static (member, id) => (member, id)) - .Where(static item => item.member.IsInWorkingSet && !item.member.IsIdle) - .Select(static item => item.id) - .Order() - .ToArray(); - Assert.Equal(expectedVisibleIds, visibleIds); - } - - private void Record(int worker, int phase, WorkingSetOperation operation, string result) - { - _trace.Enqueue( - $"seed=0x{seed:X8}; worker={worker}; phase={phase}; operation={operation}; result={result}"); - } - } - [Fact, TestCategory("Activation")] public void ActivationData_Constructor_InitializesWorkingSetClockStatus() { @@ -1873,58 +1612,6 @@ public void ActivationData_CompletedRequest_DoesNotReaddDeactivatingActivation() Assert.Empty(fixture.Observer.GetHistory(0)); } - [Fact, TestCategory("Activation")] - public async Task ActivationCollector_AgeCollection_RequestsActivationIdle() - { - var timeProvider = new FakeTimeProvider(DateTimeOffset.Parse("2025-01-01T00:00:00.000+00:00")); - using var serviceProvider = CreateCatalogServiceProvider(); - var collector = new ActivationCollector( - timeProvider, - Options.Create(new GrainCollectionOptions()), - NullLogger.Instance, - Substitute.For(), - serviceProvider.GetRequiredService()); - var activation = CreateCollectorActivation(TimeSpan.FromMinutes(5)); - activation.GetIdleness().Returns(TimeSpan.FromMinutes(2)); - collector.ScheduleCollection(activation, TimeSpan.FromMinutes(5), timeProvider.GetUtcNow().UtcDateTime); - - await collector.CollectActivations(TimeSpan.FromMinutes(1), CancellationToken.None); - - activation.Received(1).Deactivate( - Arg.Is(reason => reason.ReasonCode == DeactivationReasonCode.ActivationIdle), - Arg.Any()); - activation.Received(1).Deactivate(Arg.Any(), Arg.Any()); - activation.DidNotReceive().Deactivate( - Arg.Is(reason => reason.ReasonCode == DeactivationReasonCode.HighMemoryPressure), - Arg.Any()); - } - - [Fact, TestCategory("Activation")] - public async Task ActivationCollector_HighMemoryCollection_RequestsHighMemoryPressure() - { - var timeProvider = new FakeTimeProvider(DateTimeOffset.Parse("2025-01-01T00:00:00.000+00:00")); - using var serviceProvider = CreateCatalogServiceProvider(); - var collector = new ActivationCollector( - timeProvider, - Options.Create(new GrainCollectionOptions()), - NullLogger.Instance, - Substitute.For(), - serviceProvider.GetRequiredService()); - var activation = CreateCollectorActivation(TimeSpan.FromMinutes(5)); - collector.ScheduleCollection(activation, TimeSpan.FromMinutes(5), timeProvider.GetUtcNow().UtcDateTime); - collector._activationCount = 1; - - await collector.DeactivateInDueTimeOrder(1, CancellationToken.None); - - activation.Received(1).Deactivate( - Arg.Is(reason => reason.ReasonCode == DeactivationReasonCode.HighMemoryPressure), - Arg.Any()); - activation.Received(1).Deactivate(Arg.Any(), Arg.Any()); - activation.DidNotReceive().Deactivate( - Arg.Is(reason => reason.ReasonCode == DeactivationReasonCode.ActivationIdle), - Arg.Any()); - } - public static IEnumerable ActivationStatusCases() { foreach (var state in Enum.GetValues()) @@ -1936,26 +1623,6 @@ public static IEnumerable ActivationStatusCases() } } - private static ICollectibleGrainContext CreateCollectorActivation(TimeSpan collectionAgeLimit) - { - var activation = Substitute.For(); - activation.CollectionAgeLimit.Returns(collectionAgeLimit); - activation.IsValid.Returns(true); - activation.IsExemptFromCollection.Returns(false); - activation.IsInactive.Returns(true); - activation.Deactivated.Returns(Task.CompletedTask); - return activation; - } - - private static ServiceProvider CreateCatalogServiceProvider() - { - var services = new ServiceCollection(); - services.AddMetrics(); - services.AddSingleton(); - services.AddSingleton(); - return services.BuildServiceProvider(); - } - private sealed class ActivationDataWorkingSetFixture : IDisposable { private static readonly PropertyInfo WasRemovedByCollectionProperty = typeof(ActivationData).GetProperty( From 9b51e286a61adc4501ad0b19decf7adf9eaa8696 Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Fri, 28 Aug 2026 08:22:03 -0700 Subject: [PATCH 09/12] fix(runtime): preserve working set member compatibility --- src/Orleans.Runtime/Catalog/ActivationData.cs | 6 +- .../Catalog/ActivationWorkingSet.cs | 87 +++++++++---- src/api/Orleans.Runtime/Orleans.Runtime.cs | 4 - .../Runtime/ActivationCollectorTests.cs | 117 +++++++++++++++++- .../DeactivatedGrainQueueTests.cs | 3 - 5 files changed, 180 insertions(+), 37 deletions(-) diff --git a/src/Orleans.Runtime/Catalog/ActivationData.cs b/src/Orleans.Runtime/Catalog/ActivationData.cs index 10d035a6c68..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, @@ -1003,7 +1003,7 @@ bool IActivationWorkingSetMember.IsCandidateForRemoval(bool wouldRemove) return inactive; } - bool IActivationWorkingSetMember.IsInWorkingSet + bool IActivationWorkingSetMemberStatus.IsInWorkingSet { get => IsInWorkingSet; set @@ -1017,7 +1017,7 @@ bool IActivationWorkingSetMember.IsInWorkingSet } } - bool IActivationWorkingSetMember.IsIdle + bool IActivationWorkingSetMemberStatus.IsIdle { get => IsIdleInWorkingSet; set diff --git a/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs b/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs index a87f49e2a49..99e2df85a63 100644 --- a/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs +++ b/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs @@ -17,6 +17,7 @@ namespace Orleans.Runtime; /// 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; @@ -46,7 +47,9 @@ private IEnumerable EnumerateActiveMembers() { foreach (var pair in _members) { - if (!pair.Key.IsIdle) + if (pair.Key is IActivationWorkingSetMemberStatus status + ? !status.IsIdle + : (pair.Value & IsIdleMask) == 0) { yield return pair.Key; } @@ -63,8 +66,11 @@ public void OnActivated(IActivationWorkingSetMember member) throw new InvalidOperationException($"Member {member} is already a member of the working set"); } - member.IsInWorkingSet = true; - member.IsIdle = false; + if (member is IActivationWorkingSetMemberStatus status) + { + status.IsInWorkingSet = true; + status.IsIdle = false; + } } Interlocked.Increment(ref _activeCount); @@ -78,9 +84,18 @@ public void OnActive(IActivationWorkingSetMember member) { lock (member) { - member.IsInWorkingSet = true; - member.IsIdle = false; - if (_members.TryAdd(member, 0)) + var added = _members.TryAdd(member, 0); + if (member is IActivationWorkingSetMemberStatus status) + { + status.IsInWorkingSet = true; + status.IsIdle = false; + } + else if (!added) + { + _members[member] = 0; + } + + if (added) { Interlocked.Increment(ref _activeCount); } @@ -98,10 +113,10 @@ public void OnEvicted(IActivationWorkingSetMember member) lock (member) { removed = _members.TryRemove(member, out _); - if (removed) + if (removed && member is IActivationWorkingSetMemberStatus status) { - member.IsInWorkingSet = false; - member.IsIdle = false; + status.IsInWorkingSet = false; + status.IsIdle = false; } } @@ -163,21 +178,30 @@ private void VisitMember(IActivationWorkingSetMember member) // member's current state while holding its lock instead of adding a dictionary validation to every scan. lock (member) { - if (!member.IsInWorkingSet) + var status = member as IActivationWorkingSetMemberStatus; + byte dictionaryState = 0; + if ((status is null && !_members.TryGetValue(member, out dictionaryState)) + || (status is not null && !status.IsInWorkingSet)) { result = MemberVisitResult.None; } else { - var wouldRemove = member.IsIdle; + var wouldRemove = status is not null + ? status.IsIdle + : (dictionaryState & IsIdleMask) != 0; if (member.IsCandidateForRemoval(wouldRemove)) { if (wouldRemove) { if (_members.TryRemove(member, out _)) { - member.IsInWorkingSet = false; - member.IsIdle = false; + if (status is not null) + { + status.IsInWorkingSet = false; + status.IsIdle = false; + } + Interlocked.Decrement(ref _activeCount); result = MemberVisitResult.Evicted; } @@ -188,13 +212,29 @@ private void VisitMember(IActivationWorkingSetMember member) } else { - member.IsIdle = true; + if (status is not null) + { + status.IsIdle = true; + } + else + { + _members[member] = IsIdleMask; + } + result = MemberVisitResult.Idle; } } else { - member.IsIdle = false; + if (wouldRemove && status is not null) + { + status.IsIdle = false; + } + else if (wouldRemove) + { + _members[member] = 0; + } + result = MemberVisitResult.Active; } } @@ -294,16 +334,6 @@ public interface IActivationWorkingSet /// public interface IActivationWorkingSetMember { - /// - /// Gets or sets whether this member is registered in the working set. - /// - bool IsInWorkingSet { get; set; } - - /// - /// Gets or sets whether this member was idle during the previous working-set scan. - /// - bool IsIdle { get; set; } - /// /// Returns if the member is eligible for removal, otherwise. /// @@ -314,6 +344,13 @@ public interface IActivationWorkingSetMember bool IsCandidateForRemoval(bool wouldRemove); } +internal interface IActivationWorkingSetMemberStatus : IActivationWorkingSetMember +{ + bool IsInWorkingSet { get; set; } + + bool IsIdle { get; set; } +} + /// /// An observer. /// diff --git a/src/api/Orleans.Runtime/Orleans.Runtime.cs b/src/api/Orleans.Runtime/Orleans.Runtime.cs index 76757ea9106..b8b8d347716 100644 --- a/src/api/Orleans.Runtime/Orleans.Runtime.cs +++ b/src/api/Orleans.Runtime/Orleans.Runtime.cs @@ -748,10 +748,6 @@ public partial interface IActivationWorkingSet public partial interface IActivationWorkingSetMember { - bool IsIdle { get; set; } - - bool IsInWorkingSet { get; set; } - bool IsCandidateForRemoval(bool wouldRemove); } diff --git a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs index 8497a538676..9e927796999 100644 --- a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs +++ b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs @@ -708,6 +708,43 @@ public void WorkingSetScan_SkipsRemovedMember() 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); @@ -732,7 +769,7 @@ private IActivationWorkingSetMember PrepareActivation(TimeSpan collectionAgeLimi return (IActivationWorkingSetMember)activation; } - private sealed class TestWorkingSetMember(Func? isCandidateForRemoval = null) : IActivationWorkingSetMember + private sealed class TestWorkingSetMember(Func? isCandidateForRemoval = null) : IActivationWorkingSetMemberStatus { private bool _isIdle; private bool _isInWorkingSet; @@ -765,6 +802,18 @@ public bool IsCandidateForRemoval(bool wouldRemove) } + 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() { @@ -832,6 +881,70 @@ public async Task WorkingSet_ConcurrentOnActiveForAbsentMember_AddsOnceAndNotifi 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; + } + } + })).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() { @@ -1710,7 +1823,7 @@ public ActivationDataWorkingSetFixture() } public ActivationData Activation { get; } - public IActivationWorkingSetMember Member { get; } + public IActivationWorkingSetMemberStatus Member { get; } public FakeTimeProvider TimeProvider { get; } public ActivationWorkingSet WorkingSet { get; } public RecordingWorkingSetObserver Observer { get; } diff --git a/test/Orleans.Placement.Tests/ActivationRepartitioningTests/DeactivatedGrainQueueTests.cs b/test/Orleans.Placement.Tests/ActivationRepartitioningTests/DeactivatedGrainQueueTests.cs index dd41f42591e..02e1f6c4d0a 100644 --- a/test/Orleans.Placement.Tests/ActivationRepartitioningTests/DeactivatedGrainQueueTests.cs +++ b/test/Orleans.Placement.Tests/ActivationRepartitioningTests/DeactivatedGrainQueueTests.cs @@ -89,9 +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 IsIdle { get; set; } - public bool IsInWorkingSet { get; set; } - 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(); From cf833bfc529dc952897026e4cb8a86ecf8efd5f0 Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Fri, 28 Aug 2026 08:33:46 -0700 Subject: [PATCH 10/12] test(runtime): observe activation test cancellation --- .../Runtime/ActivationCollectorTests.cs | 75 ++++++++++--------- 1 file changed, 41 insertions(+), 34 deletions(-) diff --git a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs index 9e927796999..80e1f7041d4 100644 --- a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs +++ b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs @@ -410,8 +410,8 @@ public async Task WorkingSetScan_DoesNotUpdateReaddedMember() var lifecycle = new SiloLifecycleSubject(NullLogger.Instance); ((ILifecycleParticipant)workingSet).Participate(lifecycle); - await lifecycle.OnStart(); - await scanStarted.Task.WaitAsync(TimeSpan.FromSeconds(10)); + 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(() => @@ -419,11 +419,11 @@ public async Task WorkingSetScan_DoesNotUpdateReaddedMember() mutationStarted.SetResult(); workingSet.OnEvicted(member); workingSet.OnActivated(member); - }); - await mutationStarted.Task.WaitAsync(TimeSpan.FromSeconds(10)); + }, TestContext.Current.CancellationToken); + await mutationStarted.Task.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); try { - var stopTask = lifecycle.OnStop(); + var stopTask = lifecycle.OnStop(TestContext.Current.CancellationToken); Assert.False(stopTask.IsCompleted); Assert.False(mutationTask.IsCompleted); resumeScan.Set(); @@ -473,8 +473,8 @@ public async Task WorkingSetScan_DoesNotRemoveReaddedMember() var lifecycle = new SiloLifecycleSubject(NullLogger.Instance); ((ILifecycleParticipant)workingSet).Participate(lifecycle); - await lifecycle.OnStart(); - await removalScanStarted.Task.WaitAsync(TimeSpan.FromSeconds(10)); + 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(() => @@ -482,11 +482,11 @@ public async Task WorkingSetScan_DoesNotRemoveReaddedMember() mutationStarted.SetResult(); workingSet.OnEvicted(member); workingSet.OnActivated(member); - }); - await mutationStarted.Task.WaitAsync(TimeSpan.FromSeconds(10)); + }, TestContext.Current.CancellationToken); + await mutationStarted.Task.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); try { - var stopTask = lifecycle.OnStop(); + var stopTask = lifecycle.OnStop(TestContext.Current.CancellationToken); Assert.False(stopTask.IsCompleted); Assert.False(mutationTask.IsCompleted); resumeScan.Set(); @@ -533,8 +533,8 @@ public async Task WorkingSetScan_RepeatedCyclesPreserveCountAndObserverConsisten var lifecycle = new SiloLifecycleSubject(NullLogger.Instance); ((ILifecycleParticipant)workingSet).Participate(lifecycle); - await lifecycle.OnStart(); - await lifecycle.OnStop(); + await lifecycle.OnStart(TestContext.Current.CancellationToken); + await lifecycle.OnStop(TestContext.Current.CancellationToken); Assert.Equal(0, workingSet.Count); Assert.DoesNotContain(member, workingSet.Members); @@ -586,8 +586,8 @@ public async Task WorkingSetMembers_EnumerationToleratesConcurrentRemoveAndReadd workingSet.OnEvicted(member); workingSet.OnActivated(member); } - }); - await firstEviction.Task.WaitAsync(TimeSpan.FromSeconds(10)); + }, TestContext.Current.CancellationToken); + await firstEviction.Task.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); try { @@ -601,7 +601,7 @@ public async Task WorkingSetMembers_EnumerationToleratesConcurrentRemoveAndReadd resumeWriter.Set(); } - await writer.WaitAsync(TimeSpan.FromSeconds(10)); + await writer.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); Assert.All(enumeratedMembers, member => Assert.Contains(member, members)); Assert.Equal(members.Length, workingSet.Count); @@ -641,22 +641,22 @@ public async Task WorkingSetScan_SerializesRemovalWithReactivation() var lifecycle = new SiloLifecycleSubject(NullLogger.Instance); ((ILifecycleParticipant)workingSet).Participate(lifecycle); - await lifecycle.OnStart(); - await removalScanStarted.Task.WaitAsync(TimeSpan.FromSeconds(10)); + 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); - }); - await reactivationStarted.Task.WaitAsync(TimeSpan.FromSeconds(10)); + }, 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(); + await lifecycle.OnStop(TestContext.Current.CancellationToken); } finally { @@ -854,18 +854,20 @@ public async Task WorkingSet_ConcurrentOnActiveForAbsentMember_AddsOnceAndNotifi ready.Signal(); await start.Task; harness.WorkingSet.OnActive(member); - })).ToArray(); + }, TestContext.Current.CancellationToken)).ToArray(); try { - Assert.True(ready.Wait(TimeSpan.FromSeconds(10)), "Workers did not reach the OnActive start gate."); + 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)); + await Task.WhenAll(workers).WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); } finally { start.TrySetResult(); - await Task.WhenAll(workers).WaitAsync(TimeSpan.FromSeconds(10)); + await Task.WhenAll(workers).WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); } Assert.Equal(1, harness.WorkingSet.Count); @@ -922,7 +924,7 @@ public async Task WorkingSet_BoundedConcurrentTransitionsPreserveMembershipCount break; } } - })).ToArray(); + }, TestContext.Current.CancellationToken)).ToArray(); await Task.WhenAll(workers); @@ -953,13 +955,17 @@ public async Task WorkingSet_EvictionCallbackCanOverlapReAddWithoutHoldingMember harness.WorkingSet.OnActivated(member); harness.Observer.Clear(); var gate = harness.Observer.ArmEvictionGate(0); - var eviction = Task.Run(() => harness.WorkingSet.OnEvicted(member)); + var eviction = Task.Run( + () => harness.WorkingSet.OnEvicted(member), + TestContext.Current.CancellationToken); Task? reAdd = null; try { - await gate.Entered.Task.WaitAsync(TimeSpan.FromSeconds(10)); - reAdd = Task.Run(() => harness.WorkingSet.OnActive(member)); - await reAdd.WaitAsync(TimeSpan.FromSeconds(10)); + 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); @@ -971,10 +977,10 @@ public async Task WorkingSet_EvictionCallbackCanOverlapReAddWithoutHoldingMember gate.Release.TrySetResult(); if (reAdd is not null) { - await reAdd.WaitAsync(TimeSpan.FromSeconds(10)); + await reAdd.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); } - await eviction.WaitAsync(TimeSpan.FromSeconds(10)); + await eviction.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); } Assert.Equal(["EvictedStarted", "Active", "EvictedCompleted"], harness.Observer.GetHistory(0)); @@ -1390,7 +1396,7 @@ public static async Task CreateAsync() var result = new WorkingSetHarness(); try { - await result._lifecycle.OnStart(); + await result._lifecycle.OnStart(TestContext.Current.CancellationToken); await result._timer.WaitForGenerationAsync(1); return result; } @@ -1412,7 +1418,8 @@ public async ValueTask DisposeAsync() { try { - await _lifecycle.OnStop().WaitAsync(TimeSpan.FromSeconds(10)); + await _lifecycle.OnStop(TestContext.Current.CancellationToken) + .WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); } finally { @@ -1491,7 +1498,7 @@ public async Task WaitForGenerationAsync(int expectedGeneration) generationChanged = _generationChanged.Task; } - await generationChanged.WaitAsync(TimeSpan.FromSeconds(10)); + await generationChanged.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); } } From f73b4f6c8eef6a8ef27efb9be398ee38e90d891c Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Fri, 28 Aug 2026 09:06:20 -0700 Subject: [PATCH 11/12] perf(runtime): avoid redundant working set locks --- .../Catalog/ActivationWorkingSet.cs | 92 ++++++++++++++----- 1 file changed, 68 insertions(+), 24 deletions(-) diff --git a/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs b/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs index 99e2df85a63..952f2b4d7a7 100644 --- a/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs +++ b/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs @@ -59,64 +59,92 @@ private IEnumerable EnumerateActiveMembers() public void OnActivated(IActivationWorkingSetMember member) { Debug.Assert(member is not ICollectibleGrainContext collectible || collectible.IsValid); - lock (member) + if (Monitor.IsEntered(member)) + { + AddMember(); + } + else { + lock (member) + { + AddMember(); + } + } + + foreach (var observer in _observers) + { + observer.OnAdded(member); + } + + void AddMember() + { + Debug.Assert(Monitor.IsEntered(member)); if (!_members.TryAdd(member, 0)) { 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 OnActive(IActivationWorkingSetMember member) + { + if (Monitor.IsEntered(member)) + { + MarkActive(); + } + else + { + lock (member) + { + MarkActive(); + } + } - Interlocked.Increment(ref _activeCount); foreach (var observer in _observers) { - observer.OnAdded(member); + observer.OnActive(member); } - } - public void OnActive(IActivationWorkingSetMember member) - { - lock (member) + void MarkActive() { + Debug.Assert(Monitor.IsEntered(member)); var added = _members.TryAdd(member, 0); + if (added) + { + Interlocked.Increment(ref _activeCount); + } + if (member is IActivationWorkingSetMemberStatus status) { status.IsInWorkingSet = true; status.IsIdle = false; } - else if (!added) + else { _members[member] = 0; } - - if (added) - { - Interlocked.Increment(ref _activeCount); - } - } - - foreach (var observer in _observers) - { - observer.OnActive(member); } } public void OnEvicted(IActivationWorkingSetMember member) { bool removed; - lock (member) + if (Monitor.IsEntered(member)) { - removed = _members.TryRemove(member, out _); - if (removed && member is IActivationWorkingSetMemberStatus status) + removed = RemoveMember(); + } + else + { + lock (member) { - status.IsInWorkingSet = false; - status.IsIdle = false; + removed = RemoveMember(); } } @@ -124,11 +152,27 @@ public void OnEvicted(IActivationWorkingSetMember member) { OnEvictedCore(member); } + + bool RemoveMember() + { + Debug.Assert(Monitor.IsEntered(member)); + var result = _members.TryRemove(member, out _); + if (result) + { + Interlocked.Decrement(ref _activeCount); + if (member is IActivationWorkingSetMemberStatus status) + { + status.IsInWorkingSet = false; + status.IsIdle = false; + } + } + + return result; + } } private void OnEvictedCore(IActivationWorkingSetMember member) { - Interlocked.Decrement(ref _activeCount); foreach (var observer in _observers) { observer.OnEvicted(member); From c730b82820ddf38e6ecddb3e0972be9c8cf3d15b Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Sat, 29 Aug 2026 02:45:09 -0700 Subject: [PATCH 12/12] fix(runtime): exclude evicted working set members --- .../Catalog/ActivationWorkingSet.cs | 2 +- .../Runtime/ActivationCollectorTests.cs | 23 +++++++++++++++++++ 2 files changed, 24 insertions(+), 1 deletion(-) diff --git a/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs b/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs index 952f2b4d7a7..0e251c5e182 100644 --- a/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs +++ b/src/Orleans.Runtime/Catalog/ActivationWorkingSet.cs @@ -48,7 +48,7 @@ private IEnumerable EnumerateActiveMembers() foreach (var pair in _members) { if (pair.Key is IActivationWorkingSetMemberStatus status - ? !status.IsIdle + ? status.IsInWorkingSet && !status.IsIdle : (pair.Value & IsIdleMask) == 0) { yield return pair.Key; diff --git a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs index 80e1f7041d4..3448d4f62e1 100644 --- a/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs +++ b/test/Orleans.Core.Tests/Runtime/ActivationCollectorTests.cs @@ -610,6 +610,29 @@ public async Task WorkingSetMembers_EnumerationToleratesConcurrentRemoveAndReadd 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() {