From 3e2499c14fde90b1a824df75e33e8d3183a0e7b9 Mon Sep 17 00:00:00 2001 From: Aaron Stannard Date: Tue, 19 May 2026 14:22:25 +0000 Subject: [PATCH 1/3] fix: implement IAsyncLifetime on Akka.TestKit.Xunit (v3) TestKit (#8191) xUnit v3 disposes a test class instance via IAsyncDisposable.DisposeAsync() in preference to IDisposable.Dispose() whenever the type implements both (Xunit.IAsyncLifetime derives from IAsyncDisposable). Once a derived spec implemented IAsyncLifetime, the v3 Akka.TestKit.Xunit.TestKit's Dispose() was never invoked, so AfterAll()/Shutdown() never ran and the ActorSystem was silently leaked. TestKit now implements IAsyncLifetime with public virtual InitializeAsync() and DisposeAsync(). DisposeAsync() drives the existing synchronous dispose chain, so a derived override can chain to base.DisposeAsync(). Converts the in-repo specs that hit the trap (Bugfix8144Spec, ParallelAmbientContext*, EventFilterTestBase) to the override form, and adds Bugfix8191Spec as a regression test. The EventFilterTestBase conversion also restores its AfterAll() -> EnsureNoMoreLoggedMessages() verification, which was being skipped. --- .../Bugfix8144Spec.cs | 9 +-- .../Bugfix8191Spec.cs | 73 +++++++++++++++++++ .../ParallelAmbientContextSpec.cs | 12 +-- .../testkits/Akka.TestKit.Xunit/TestKit.cs | 29 +++++++- .../EventFilterTestBase.cs | 13 +--- 5 files changed, 112 insertions(+), 24 deletions(-) create mode 100644 src/contrib/testkits/Akka.TestKit.Xunit.Tests/Bugfix8191Spec.cs diff --git a/src/contrib/testkits/Akka.TestKit.Xunit.Tests/Bugfix8144Spec.cs b/src/contrib/testkits/Akka.TestKit.Xunit.Tests/Bugfix8144Spec.cs index 06b5f2d4799..86eee5bf4c2 100644 --- a/src/contrib/testkits/Akka.TestKit.Xunit.Tests/Bugfix8144Spec.cs +++ b/src/contrib/testkits/Akka.TestKit.Xunit.Tests/Bugfix8144Spec.cs @@ -17,15 +17,12 @@ namespace Akka.TestKit.Tests; /// Verifies that implicit sender (TestActor) is preserved across async boundaries. /// Regression test for https://github.com/akkadotnet/akka.net/issues/8144 /// -public class Bugfix8144Spec: Xunit.TestKit, IAsyncLifetime +public class Bugfix8144Spec: Xunit.TestKit { - public ValueTask DisposeAsync() + public override async ValueTask InitializeAsync() { - return new ValueTask(Task.CompletedTask); - } + await base.InitializeAsync(); - public async ValueTask InitializeAsync() - { // Any await here can cause a thread switch, which previously // lost the [ThreadStatic] InternalCurrentActorCellKeeper.Current await Task.Delay(TimeSpan.FromMilliseconds(100)); diff --git a/src/contrib/testkits/Akka.TestKit.Xunit.Tests/Bugfix8191Spec.cs b/src/contrib/testkits/Akka.TestKit.Xunit.Tests/Bugfix8191Spec.cs new file mode 100644 index 00000000000..bc9c56d076e --- /dev/null +++ b/src/contrib/testkits/Akka.TestKit.Xunit.Tests/Bugfix8191Spec.cs @@ -0,0 +1,73 @@ +// ----------------------------------------------------------------------- +// +// Copyright (C) 2009-2022 Lightbend Inc. +// Copyright (C) 2013-2026 .NET Foundation +// +// ----------------------------------------------------------------------- + +using System; +using System.Threading.Tasks; +using Xunit; + +namespace Akka.TestKit.Xunit.Tests; + +/// +/// Regression tests for https://github.com/akkadotnet/akka.net/issues/8191 +/// +/// xUnit v3 tears a test class instance down via +/// in preference to whenever the type implements both. +/// must therefore drive its dispose chain — and shut the +/// ActorSystem down — from DisposeAsync, and a derived DisposeAsync +/// override must be able to chain to base.DisposeAsync(). Before the fix the +/// ActorSystem was silently leaked once a derived spec implemented +/// IAsyncLifetime. +/// +public class Bugfix8191Spec +{ + private sealed class TrackingTestKit : TestKit + { + public bool AfterAllRan { get; private set; } + + protected override void AfterAll() + { + AfterAllRan = true; + base.AfterAll(); + } + } + + private sealed class AsyncTeardownTestKit : TestKit + { + public bool DisposeAsyncOverrideRan { get; private set; } + + public override async ValueTask DisposeAsync() + { + DisposeAsyncOverrideRan = true; + await base.DisposeAsync(); + } + } + + [Fact(DisplayName = "TestKit.DisposeAsync should run the dispose chain and shut down the ActorSystem")] + public async Task Should_run_dispose_chain_and_shut_down_system_When_disposed_via_DisposeAsync() + { + var testKit = new TrackingTestKit(); + var system = testKit.Sys; + + // Exercise the exact path xUnit v3 uses to tear down a test class instance. + await ((IAsyncDisposable)testKit).DisposeAsync(); + + Assert.True(testKit.AfterAllRan, "AfterAll() should run as part of the DisposeAsync chain"); + Assert.True(system.WhenTerminated.IsCompleted, "the ActorSystem should be shut down"); + } + + [Fact(DisplayName = "A derived DisposeAsync override chaining to base should shut down the ActorSystem")] + public async Task Should_shut_down_system_When_derived_DisposeAsync_chains_to_base() + { + var testKit = new AsyncTeardownTestKit(); + var system = testKit.Sys; + + await ((IAsyncDisposable)testKit).DisposeAsync(); + + Assert.True(testKit.DisposeAsyncOverrideRan, "the derived DisposeAsync override should run"); + Assert.True(system.WhenTerminated.IsCompleted, "base.DisposeAsync() should shut the ActorSystem down"); + } +} diff --git a/src/contrib/testkits/Akka.TestKit.Xunit.Tests/ParallelAmbientContextSpec.cs b/src/contrib/testkits/Akka.TestKit.Xunit.Tests/ParallelAmbientContextSpec.cs index bcc60b7298f..83a6e3fe82b 100644 --- a/src/contrib/testkits/Akka.TestKit.Xunit.Tests/ParallelAmbientContextSpec.cs +++ b/src/contrib/testkits/Akka.TestKit.Xunit.Tests/ParallelAmbientContextSpec.cs @@ -20,14 +20,12 @@ namespace Akka.TestKit.Xunit.Tests; // This project provides its own xunit.runner.json with parallel collections // enabled so CI and local runs exercise the reported failure mode by default. -public abstract class ParallelAmbientContextSpecBase : TestKit, IAsyncLifetime +public abstract class ParallelAmbientContextSpecBase : TestKit { // Forces the post-ctor continuation onto a different SC worker — the // thread pollution only manifests when the body thread differs from the // ctor thread. - public async ValueTask InitializeAsync() => await Task.Yield(); - - public ValueTask DisposeAsync() => default; + public override async ValueTask InitializeAsync() => await Task.Yield(); [Fact] public async Task Implicit_sender_should_resolve_to_own_TestActor() @@ -67,11 +65,9 @@ public class ParallelAmbientContextSpec16 : ParallelAmbientContextSpecBase { } // Current must be null both at body entry (pre-await prefix) and across any // await continuations resumed on a reused worker thread that a sibling's // Before hook may have pinned with a non-null cell. -public abstract class ParallelNoImplicitSenderSpecBase : TestKit, IAsyncLifetime, INoImplicitSender +public abstract class ParallelNoImplicitSenderSpecBase : TestKit, INoImplicitSender { - public async ValueTask InitializeAsync() => await Task.Yield(); - - public ValueTask DisposeAsync() => default; + public override async ValueTask InitializeAsync() => await Task.Yield(); [Fact] public async Task Current_should_be_null_both_pre_and_post_await() diff --git a/src/contrib/testkits/Akka.TestKit.Xunit/TestKit.cs b/src/contrib/testkits/Akka.TestKit.Xunit/TestKit.cs index 213442df4c0..242c6736529 100644 --- a/src/contrib/testkits/Akka.TestKit.Xunit/TestKit.cs +++ b/src/contrib/testkits/Akka.TestKit.Xunit/TestKit.cs @@ -6,6 +6,7 @@ //----------------------------------------------------------------------- using System; +using System.Threading.Tasks; using Akka.Actor; using Akka.Actor.Internal; using Akka.Actor.Setup; @@ -22,7 +23,7 @@ namespace Akka.TestKit.Xunit; /// as its testing framework. /// [AkkaCleanAmbientContext] -public class TestKit : TestKitBase, IDisposable +public class TestKit : TestKitBase, IDisposable, IAsyncLifetime { private class PrefixedOutput : ITestOutputHelper { @@ -249,6 +250,14 @@ protected void InitializeLogger(ActorSystem system, string prefix) logger.Tell(new InitializeLogger(system.EventStream), ActorRefs.NoSender); } + /// + /// xUnit lifecycle hook, invoked once before the test method runs. The default + /// implementation does nothing. Override to perform asynchronous test setup, and + /// call await base.InitializeAsync() from your override. + /// + public virtual ValueTask InitializeAsync() + => default; + /// /// Performs application-defined tasks associated with freeing, releasing, or resetting unmanaged resources. /// @@ -279,4 +288,22 @@ public void Dispose() { Dispose(true); } + + /// + /// xUnit lifecycle hook, invoked once after the test method completes. The default + /// implementation drives the synchronous dispose chain + /// ( -> -> actor system shutdown). + /// + /// Override this for asynchronous teardown, and always call await base.DisposeAsync() + /// from your override. xUnit v3 invokes in + /// preference to for any type that implements both + /// interfaces, so an override that does not chain to the base will skip shutdown and leak + /// the . + /// + /// + public virtual ValueTask DisposeAsync() + { + Dispose(true); + return default; + } } \ No newline at end of file diff --git a/src/core/Akka.TestKit.Tests/TestEventListenerTests/EventFilterTestBase.cs b/src/core/Akka.TestKit.Tests/TestEventListenerTests/EventFilterTestBase.cs index ce72ebdd457..f12a005e595 100644 --- a/src/core/Akka.TestKit.Tests/TestEventListenerTests/EventFilterTestBase.cs +++ b/src/core/Akka.TestKit.Tests/TestEventListenerTests/EventFilterTestBase.cs @@ -12,7 +12,7 @@ namespace Akka.TestKit.Tests.TestEventListenerTests { - public abstract class EventFilterTestBase : TestKit.Xunit.TestKit, IAsyncLifetime + public abstract class EventFilterTestBase : TestKit.Xunit.TestKit { /// /// Used to signal that the test was successful and that we should ensure no more messages were logged @@ -24,22 +24,17 @@ protected EventFilterTestBase(string config) { } - public ValueTask InitializeAsync() + public override ValueTask InitializeAsync() { //We send a ForwardAllEventsTo containing message to the TestEventListenerToForwarder logger (configured as a logger above). //It should respond with an "OK" message when it has received the message. var initLoggerMessage = new ForwardAllEventsTestEventListener.ForwardAllEventsTo(TestActor); - + SendRawLogEventMessage(initLoggerMessage); ExpectMsg("OK"); //From now on we know that all messages will be forwarded to TestActor - - return new ValueTask(Task.CompletedTask); - } - public ValueTask DisposeAsync() - { - return new ValueTask(Task.CompletedTask); + return base.InitializeAsync(); } protected abstract void SendRawLogEventMessage(object message); From 376cfd8811853d87bd478bfad583f149465e20b2 Mon Sep 17 00:00:00 2001 From: Aaron Stannard Date: Tue, 19 May 2026 14:55:05 +0000 Subject: [PATCH 2/3] fix: convert remaining IAsyncLifetime specs to override (CS0114) BugFixSpec and StreamRefsSpec consume the xUnit v3 Akka.TestKit.Xunit transitively via the shared xunit3 AkkaSpec. Now that TestKit declares public virtual InitializeAsync()/DisposeAsync(), their own implicit implementations hide the inherited members (CS0114), which the projects' TreatWarningsAsErrors promotes to a build error. Converts both to the override form, dropping the redundant IAsyncLifetime declaration and chaining base.InitializeAsync(). Removing the no-op DisposeAsync overrides also lets the base dispose chain run, fixing the ActorSystem leak these specs had under xUnit v3. --- .../Akka.DependencyInjection.Tests/BugFixSpec.cs | 10 +++------- src/core/Akka.Streams.Tests/Dsl/StreamRefsSpec.cs | 11 ++++------- 2 files changed, 7 insertions(+), 14 deletions(-) diff --git a/src/contrib/dependencyinjection/Akka.DependencyInjection.Tests/BugFixSpec.cs b/src/contrib/dependencyinjection/Akka.DependencyInjection.Tests/BugFixSpec.cs index 6672e4700fc..353ec020177 100644 --- a/src/contrib/dependencyinjection/Akka.DependencyInjection.Tests/BugFixSpec.cs +++ b/src/contrib/dependencyinjection/Akka.DependencyInjection.Tests/BugFixSpec.cs @@ -19,7 +19,7 @@ namespace Akka.DependencyInjection.Tests { - public class BugFixSpec: AkkaSpec, IAsyncLifetime + public class BugFixSpec: AkkaSpec { private readonly IServiceProvider _serviceProvider; private readonly AkkaService _akkaService; @@ -59,8 +59,9 @@ internal class NotInServices } - public async ValueTask InitializeAsync() + public override async ValueTask InitializeAsync() { + await base.InitializeAsync(); await _akkaService.StartAsync(default); InitializeLogger(_akkaService.ActorSystem); } @@ -72,11 +73,6 @@ protected override void AfterAll() base.AfterAll(); } - public ValueTask DisposeAsync() - { - return new ValueTask(Task.CompletedTask); - } - internal class AkkaService : IHostedService { public ActorSystem ActorSystem { get; private set; } diff --git a/src/core/Akka.Streams.Tests/Dsl/StreamRefsSpec.cs b/src/core/Akka.Streams.Tests/Dsl/StreamRefsSpec.cs index cad265cc7ae..17eb49d4da0 100644 --- a/src/core/Akka.Streams.Tests/Dsl/StreamRefsSpec.cs +++ b/src/core/Akka.Streams.Tests/Dsl/StreamRefsSpec.cs @@ -189,7 +189,7 @@ public BulkSinkMsg(ISinkRef dataSink) } } - public class StreamRefsSpec : AkkaSpec, IAsyncLifetime + public class StreamRefsSpec : AkkaSpec { public static Config Config() { @@ -220,8 +220,10 @@ protected StreamRefsSpec(Config config, ITestOutputHelper output = null) : base( _probe = CreateTestProbe(); } - public async ValueTask InitializeAsync() + public override async ValueTask InitializeAsync() { + await base.InitializeAsync(); + var it = RemoteSystem.ActorOf(DataSourceActor.Props(_probe.Ref), "remoteActor"); var remoteAddress = ((ActorSystemImpl)RemoteSystem).Provider.DefaultAddress; Sys.ActorSelection(it.Path.ToStringWithAddress(remoteAddress)).Tell(new Identify("hi")); @@ -229,11 +231,6 @@ public async ValueTask InitializeAsync() _remoteActor = (await ExpectMsgAsync(TimeSpan.FromSeconds(30))).Subject; } - public ValueTask DisposeAsync() - { - return new ValueTask(Task.CompletedTask); - } - protected readonly ActorSystem RemoteSystem; protected readonly ActorMaterializer Materializer; private readonly TestProbe _probe; From c142966fbff8ad397b3b8a9771acdd26e39ae122 Mon Sep 17 00:00:00 2001 From: Aaron Stannard Date: Tue, 19 May 2026 14:55:14 +0000 Subject: [PATCH 3/3] improve: make Akka.TestKit.Xunit (v3) DisposeAsync natively async MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Addresses PR review feedback about sync-over-async teardown. Adds TestKitBase.ShutdownAsync — an async counterpart to Shutdown() that awaits ActorSystem termination instead of blocking on Terminate().Wait(). This is purely additive (extend-only); the existing synchronous Shutdown() is unchanged. Akka.TestKit.Xunit.TestKit.DisposeAsync() now runs the synchronous dispose chain (Dispose(bool) -> AfterAll, plus any derived Dispose(bool) override) with the blocking Shutdown() suppressed via a flag, then awaits ShutdownAsync(). This keeps the dispose-pattern extensibility point intact and terminates the ActorSystem without blocking the xUnit teardown thread. Regenerates the ApproveTestKit API approval baselines for the two new ShutdownAsync overloads. --- .../testkits/Akka.TestKit.Xunit/TestKit.cs | 31 +++++++++--- ...APISpec.ApproveTestKit.DotNet.verified.txt | 22 +++++---- ...oreAPISpec.ApproveTestKit.Net.verified.txt | 22 +++++---- src/core/Akka.TestKit/TestKitBase.cs | 47 +++++++++++++++++++ 4 files changed, 96 insertions(+), 26 deletions(-) diff --git a/src/contrib/testkits/Akka.TestKit.Xunit/TestKit.cs b/src/contrib/testkits/Akka.TestKit.Xunit/TestKit.cs index 242c6736529..7ea0a8954df 100644 --- a/src/contrib/testkits/Akka.TestKit.Xunit/TestKit.cs +++ b/src/contrib/testkits/Akka.TestKit.Xunit/TestKit.cs @@ -85,6 +85,7 @@ public void WriteLine(string format, params object[] args) private bool _disposed; private bool _disposing; + private bool _disposingAsync; /// /// @@ -279,7 +280,10 @@ protected virtual void Dispose(bool disposing) } finally { - Shutdown(); + // DisposeAsync() terminates the ActorSystem asynchronously and sets this flag, + // so we don't also block here with a synchronous Shutdown(). + if (!_disposingAsync) + Shutdown(); _disposed = true; } } @@ -291,8 +295,9 @@ public void Dispose() /// /// xUnit lifecycle hook, invoked once after the test method completes. The default - /// implementation drives the synchronous dispose chain - /// ( -> -> actor system shutdown). + /// implementation runs the synchronous dispose chain ( and + /// therefore ) and then asynchronously terminates the + /// , without blocking the calling thread. /// /// Override this for asynchronous teardown, and always call await base.DisposeAsync() /// from your override. xUnit v3 invokes in @@ -301,9 +306,23 @@ public void Dispose() /// the . /// /// - public virtual ValueTask DisposeAsync() + public virtual async ValueTask DisposeAsync() { - Dispose(true); - return default; + if (_disposing || _disposed) + return; + + // Run the synchronous dispose chain — AfterAll() plus any overridden Dispose(bool) — + // but suppress its blocking Shutdown() call; the ActorSystem is terminated + // asynchronously below instead. The finally guarantees shutdown still runs even if + // AfterAll() throws, matching the synchronous Dispose() behavior. + _disposingAsync = true; + try + { + Dispose(true); + } + finally + { + await ShutdownAsync(); + } } } \ No newline at end of file diff --git a/src/core/Akka.API.Tests/verify/CoreAPISpec.ApproveTestKit.DotNet.verified.txt b/src/core/Akka.API.Tests/verify/CoreAPISpec.ApproveTestKit.DotNet.verified.txt index 8dd5a2fe34d..e1ad3712459 100644 --- a/src/core/Akka.API.Tests/verify/CoreAPISpec.ApproveTestKit.DotNet.verified.txt +++ b/src/core/Akka.API.Tests/verify/CoreAPISpec.ApproveTestKit.DotNet.verified.txt @@ -405,17 +405,17 @@ namespace Akka.TestKit public System.Collections.Generic.IReadOnlyCollection ExpectMsgAllOf(System.Collections.Generic.IReadOnlyCollection messages, System.Threading.CancellationToken cancellationToken = null) { } public System.Collections.Generic.IReadOnlyCollection ExpectMsgAllOf(System.TimeSpan max, params T[] messages) { } public System.Collections.Generic.IReadOnlyCollection ExpectMsgAllOf(System.TimeSpan max, System.Collections.Generic.IReadOnlyCollection messages, System.Threading.CancellationToken cancellationToken = null) { } - [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__158))] + [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__160))] public System.Collections.Generic.IAsyncEnumerable ExpectMsgAllOfAsync(System.Collections.Generic.IReadOnlyCollection messages, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } - [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__161))] + [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__163))] public System.Collections.Generic.IAsyncEnumerable ExpectMsgAllOfAsync(System.TimeSpan max, System.Collections.Generic.IReadOnlyCollection messages, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } public System.Collections.Generic.IReadOnlyCollection ExpectMsgAllOfMatchingPredicates(params Akka.TestKit.PredicateInfo[] predicates) { } public System.Collections.Generic.IReadOnlyCollection ExpectMsgAllOfMatchingPredicates(System.Collections.Generic.IReadOnlyCollection predicates, System.Threading.CancellationToken cancellationToken) { } public System.Collections.Generic.IReadOnlyCollection ExpectMsgAllOfMatchingPredicates(System.TimeSpan max, params Akka.TestKit.PredicateInfo[] predicates) { } public System.Collections.Generic.IReadOnlyCollection ExpectMsgAllOfMatchingPredicates(System.TimeSpan max, System.Collections.Generic.IReadOnlyCollection predicates, System.Threading.CancellationToken cancellationToken) { } - [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__168))] + [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__170))] public System.Collections.Generic.IAsyncEnumerable ExpectMsgAllOfMatchingPredicatesAsync(System.Collections.Generic.IReadOnlyCollection predicates, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } - [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__171))] + [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__173))] public System.Collections.Generic.IAsyncEnumerable ExpectMsgAllOfMatchingPredicatesAsync(System.TimeSpan max, System.Collections.Generic.IReadOnlyCollection predicates, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } public T ExpectMsgAnyOf(System.Collections.Generic.IEnumerable messages, System.Threading.CancellationToken cancellationToken = null) { } public System.Threading.Tasks.ValueTask ExpectMsgAnyOfAsync(System.Collections.Generic.IEnumerable messages, System.Threading.CancellationToken cancellationToken = null) { } @@ -470,9 +470,9 @@ namespace Akka.TestKit public System.Threading.Tasks.ValueTask PeekOneAsync(System.Threading.CancellationToken cancellationToken) { } public System.Collections.Generic.IReadOnlyCollection ReceiveN(int numberOfMessages, System.Threading.CancellationToken cancellationToken = null) { } public System.Collections.Generic.IReadOnlyCollection ReceiveN(int numberOfMessages, System.TimeSpan max, System.Threading.CancellationToken cancellationToken = null) { } - [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__217))] - public System.Collections.Generic.IAsyncEnumerable ReceiveNAsync(int numberOfMessages, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__219))] + public System.Collections.Generic.IAsyncEnumerable ReceiveNAsync(int numberOfMessages, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } + [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__221))] public System.Collections.Generic.IAsyncEnumerable ReceiveNAsync(int numberOfMessages, System.TimeSpan max, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } public object ReceiveOne(System.Nullable max = null, System.Threading.CancellationToken cancellationToken = null) { } public System.Threading.Tasks.ValueTask ReceiveOneAsync(System.Nullable max = null, System.Threading.CancellationToken cancellationToken = null) { } @@ -480,19 +480,21 @@ namespace Akka.TestKit public System.Collections.Generic.IReadOnlyList ReceiveWhile(System.Nullable max, System.Nullable idle, System.Func filter, int msgs = 2147483647, System.Threading.CancellationToken cancellationToken = null) { } public System.Collections.Generic.IReadOnlyList ReceiveWhile(System.Func filter, System.Nullable max = null, System.Nullable idle = null, int msgs = 2147483647, System.Threading.CancellationToken cancellationToken = null) { } public System.Collections.Generic.IReadOnlyList ReceiveWhile(System.Predicate shouldContinue, System.Nullable max = null, System.Nullable idle = null, int msgs = 2147483647, bool shouldIgnoreOtherMessageTypes = True, System.Threading.CancellationToken cancellationToken = null) { } - [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__208))] - public System.Collections.Generic.IAsyncEnumerable ReceiveWhileAsync(System.Nullable max, System.Func filter, int msgs = 2147483647, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__210))] - public System.Collections.Generic.IAsyncEnumerable ReceiveWhileAsync(System.Nullable max, System.Nullable idle, System.Func filter, int msgs = 2147483647, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } + public System.Collections.Generic.IAsyncEnumerable ReceiveWhileAsync(System.Nullable max, System.Func filter, int msgs = 2147483647, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__212))] - public System.Collections.Generic.IAsyncEnumerable ReceiveWhileAsync(System.Func filter, System.Nullable max = null, System.Nullable idle = null, int msgs = 2147483647, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } + public System.Collections.Generic.IAsyncEnumerable ReceiveWhileAsync(System.Nullable max, System.Nullable idle, System.Func filter, int msgs = 2147483647, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__214))] + public System.Collections.Generic.IAsyncEnumerable ReceiveWhileAsync(System.Func filter, System.Nullable max = null, System.Nullable idle = null, int msgs = 2147483647, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } + [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__216))] public System.Collections.Generic.IAsyncEnumerable ReceiveWhileAsync(System.Predicate shouldContinue, System.Nullable max = null, System.Nullable idle = null, int msgs = 2147483647, bool shouldIgnoreOtherMessageTypes = True, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } protected System.TimeSpan RemainingOr(System.TimeSpan duration) { } public System.TimeSpan RemainingOrDilated(System.Nullable duration) { } public void SetAutoPilot(Akka.TestKit.AutoPilot pilot) { } public virtual void Shutdown(System.Nullable duration = null, bool verifySystemShutdown = False) { } protected virtual void Shutdown(Akka.Actor.ActorSystem system, System.Nullable duration = null, bool verifySystemShutdown = False) { } + public virtual System.Threading.Tasks.Task ShutdownAsync(System.Nullable duration = null, bool verifySystemShutdown = False) { } + protected virtual System.Threading.Tasks.Task ShutdownAsync(Akka.Actor.ActorSystem system, System.Nullable duration = null, bool verifySystemShutdown = False) { } public bool TryPeekOne(out Akka.TestKit.MessageEnvelope envelope, System.Nullable max, System.Threading.CancellationToken cancellationToken) { } [return: System.Runtime.CompilerServices.TupleElementNamesAttribute(new string[] { "success", diff --git a/src/core/Akka.API.Tests/verify/CoreAPISpec.ApproveTestKit.Net.verified.txt b/src/core/Akka.API.Tests/verify/CoreAPISpec.ApproveTestKit.Net.verified.txt index f575332de5a..ca50775f28a 100644 --- a/src/core/Akka.API.Tests/verify/CoreAPISpec.ApproveTestKit.Net.verified.txt +++ b/src/core/Akka.API.Tests/verify/CoreAPISpec.ApproveTestKit.Net.verified.txt @@ -405,17 +405,17 @@ namespace Akka.TestKit public System.Collections.Generic.IReadOnlyCollection ExpectMsgAllOf(System.Collections.Generic.IReadOnlyCollection messages, System.Threading.CancellationToken cancellationToken = null) { } public System.Collections.Generic.IReadOnlyCollection ExpectMsgAllOf(System.TimeSpan max, params T[] messages) { } public System.Collections.Generic.IReadOnlyCollection ExpectMsgAllOf(System.TimeSpan max, System.Collections.Generic.IReadOnlyCollection messages, System.Threading.CancellationToken cancellationToken = null) { } - [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__158))] + [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__160))] public System.Collections.Generic.IAsyncEnumerable ExpectMsgAllOfAsync(System.Collections.Generic.IReadOnlyCollection messages, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } - [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__161))] + [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__163))] public System.Collections.Generic.IAsyncEnumerable ExpectMsgAllOfAsync(System.TimeSpan max, System.Collections.Generic.IReadOnlyCollection messages, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } public System.Collections.Generic.IReadOnlyCollection ExpectMsgAllOfMatchingPredicates(params Akka.TestKit.PredicateInfo[] predicates) { } public System.Collections.Generic.IReadOnlyCollection ExpectMsgAllOfMatchingPredicates(System.Collections.Generic.IReadOnlyCollection predicates, System.Threading.CancellationToken cancellationToken) { } public System.Collections.Generic.IReadOnlyCollection ExpectMsgAllOfMatchingPredicates(System.TimeSpan max, params Akka.TestKit.PredicateInfo[] predicates) { } public System.Collections.Generic.IReadOnlyCollection ExpectMsgAllOfMatchingPredicates(System.TimeSpan max, System.Collections.Generic.IReadOnlyCollection predicates, System.Threading.CancellationToken cancellationToken) { } - [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__168))] + [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__170))] public System.Collections.Generic.IAsyncEnumerable ExpectMsgAllOfMatchingPredicatesAsync(System.Collections.Generic.IReadOnlyCollection predicates, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } - [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__171))] + [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__173))] public System.Collections.Generic.IAsyncEnumerable ExpectMsgAllOfMatchingPredicatesAsync(System.TimeSpan max, System.Collections.Generic.IReadOnlyCollection predicates, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } public T ExpectMsgAnyOf(System.Collections.Generic.IEnumerable messages, System.Threading.CancellationToken cancellationToken = null) { } public System.Threading.Tasks.ValueTask ExpectMsgAnyOfAsync(System.Collections.Generic.IEnumerable messages, System.Threading.CancellationToken cancellationToken = null) { } @@ -470,9 +470,9 @@ namespace Akka.TestKit public System.Threading.Tasks.ValueTask PeekOneAsync(System.Threading.CancellationToken cancellationToken) { } public System.Collections.Generic.IReadOnlyCollection ReceiveN(int numberOfMessages, System.Threading.CancellationToken cancellationToken = null) { } public System.Collections.Generic.IReadOnlyCollection ReceiveN(int numberOfMessages, System.TimeSpan max, System.Threading.CancellationToken cancellationToken = null) { } - [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__217))] - public System.Collections.Generic.IAsyncEnumerable ReceiveNAsync(int numberOfMessages, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__219))] + public System.Collections.Generic.IAsyncEnumerable ReceiveNAsync(int numberOfMessages, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } + [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__221))] public System.Collections.Generic.IAsyncEnumerable ReceiveNAsync(int numberOfMessages, System.TimeSpan max, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } public object ReceiveOne(System.Nullable max = null, System.Threading.CancellationToken cancellationToken = null) { } public System.Threading.Tasks.ValueTask ReceiveOneAsync(System.Nullable max = null, System.Threading.CancellationToken cancellationToken = null) { } @@ -480,19 +480,21 @@ namespace Akka.TestKit public System.Collections.Generic.IReadOnlyList ReceiveWhile(System.Nullable max, System.Nullable idle, System.Func filter, int msgs = 2147483647, System.Threading.CancellationToken cancellationToken = null) { } public System.Collections.Generic.IReadOnlyList ReceiveWhile(System.Func filter, System.Nullable max = null, System.Nullable idle = null, int msgs = 2147483647, System.Threading.CancellationToken cancellationToken = null) { } public System.Collections.Generic.IReadOnlyList ReceiveWhile(System.Predicate shouldContinue, System.Nullable max = null, System.Nullable idle = null, int msgs = 2147483647, bool shouldIgnoreOtherMessageTypes = True, System.Threading.CancellationToken cancellationToken = null) { } - [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__208))] - public System.Collections.Generic.IAsyncEnumerable ReceiveWhileAsync(System.Nullable max, System.Func filter, int msgs = 2147483647, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__210))] - public System.Collections.Generic.IAsyncEnumerable ReceiveWhileAsync(System.Nullable max, System.Nullable idle, System.Func filter, int msgs = 2147483647, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } + public System.Collections.Generic.IAsyncEnumerable ReceiveWhileAsync(System.Nullable max, System.Func filter, int msgs = 2147483647, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__212))] - public System.Collections.Generic.IAsyncEnumerable ReceiveWhileAsync(System.Func filter, System.Nullable max = null, System.Nullable idle = null, int msgs = 2147483647, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } + public System.Collections.Generic.IAsyncEnumerable ReceiveWhileAsync(System.Nullable max, System.Nullable idle, System.Func filter, int msgs = 2147483647, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__214))] + public System.Collections.Generic.IAsyncEnumerable ReceiveWhileAsync(System.Func filter, System.Nullable max = null, System.Nullable idle = null, int msgs = 2147483647, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } + [System.Runtime.CompilerServices.AsyncIteratorStateMachineAttribute(typeof(Akka.TestKit.TestKitBase.d__216))] public System.Collections.Generic.IAsyncEnumerable ReceiveWhileAsync(System.Predicate shouldContinue, System.Nullable max = null, System.Nullable idle = null, int msgs = 2147483647, bool shouldIgnoreOtherMessageTypes = True, [System.Runtime.CompilerServices.EnumeratorCancellationAttribute()] System.Threading.CancellationToken cancellationToken = null) { } protected System.TimeSpan RemainingOr(System.TimeSpan duration) { } public System.TimeSpan RemainingOrDilated(System.Nullable duration) { } public void SetAutoPilot(Akka.TestKit.AutoPilot pilot) { } public virtual void Shutdown(System.Nullable duration = null, bool verifySystemShutdown = False) { } protected virtual void Shutdown(Akka.Actor.ActorSystem system, System.Nullable duration = null, bool verifySystemShutdown = False) { } + public virtual System.Threading.Tasks.Task ShutdownAsync(System.Nullable duration = null, bool verifySystemShutdown = False) { } + protected virtual System.Threading.Tasks.Task ShutdownAsync(Akka.Actor.ActorSystem system, System.Nullable duration = null, bool verifySystemShutdown = False) { } public bool TryPeekOne(out Akka.TestKit.MessageEnvelope envelope, System.Nullable max, System.Threading.CancellationToken cancellationToken) { } [return: System.Runtime.CompilerServices.TupleElementNamesAttribute(new string[] { "success", diff --git a/src/core/Akka.TestKit/TestKitBase.cs b/src/core/Akka.TestKit/TestKitBase.cs index 190022246a1..9bb474dcb07 100644 --- a/src/core/Akka.TestKit/TestKitBase.cs +++ b/src/core/Akka.TestKit/TestKitBase.cs @@ -590,6 +590,53 @@ protected virtual void Shutdown( } } + /// + /// Asynchronously shuts down this system. Unlike Shutdown, this does not block the + /// calling thread while waiting for the to terminate. + /// On failure debug output will be logged about the remaining actors in the system. + /// If verifySystemShutdown is true, then an exception will be thrown on failure. + /// + /// Optional. The duration to wait for shutdown. Default is 5 seconds multiplied with the config value "akka.test.timefactor". + /// if set to true an exception will be thrown on failure. + /// TBD + public virtual Task ShutdownAsync( + TimeSpan? duration = null, + bool verifySystemShutdown = false) + => ShutdownAsync(_testState.System, duration, verifySystemShutdown); + + /// + /// Asynchronously shuts down the specified system. Unlike Shutdown, this does not block + /// the calling thread while waiting for the to terminate. + /// On failure debug output will be logged about the remaining actors in the system. + /// If verifySystemShutdown is true, then an exception will be thrown on failure. + /// + /// The system to shutdown. + /// The duration to wait for shutdown. Default is 5 seconds multiplied with the config value "akka.test.timefactor" + /// if set to true an exception will be thrown on failure. + /// TBD + protected virtual async Task ShutdownAsync( + ActorSystem system, + TimeSpan? duration = null, + bool verifySystemShutdown = false) + { + system ??= _testState.System; + + var durationValue = duration.GetValueOrDefault(Dilated(TimeSpan.FromSeconds(5)).Min(TimeSpan.FromSeconds(10))); + + var terminateTask = system.Terminate(); + var wasShutdownDuringWait = await Task.WhenAny(terminateTask, Task.Delay(durationValue)) == terminateTask; + if(!wasShutdownDuringWait) + { + // Forcefully close the ActorSystem to make sure we exit the test cleanly + ((ExtendedActorSystem) system).Guardian.Stop(); + + const string msg = "Failed to stop [{0}] within [{1}]. ActorSystem is being forcefully shut down.\n{2}"; + if(verifySystemShutdown) + throw new TimeoutException(string.Format(msg, system.Name, durationValue, ((ExtendedActorSystem) system).PrintTree())); + system.Log.Warning(msg, system.Name, durationValue, ((ExtendedActorSystem) system).PrintTree()); + } + } + /// /// Spawns an actor as a child of this test actor, and returns the child's IActorRef ///