diff --git a/test/Extensions/Orleans.Azure.Tests/Streaming/StreamReliabilityTests.cs b/test/Extensions/Orleans.Azure.Tests/Streaming/StreamReliabilityTests.cs index c6e05907b46..6a0800b72a9 100644 --- a/test/Extensions/Orleans.Azure.Tests/Streaming/StreamReliabilityTests.cs +++ b/test/Extensions/Orleans.Azure.Tests/Streaming/StreamReliabilityTests.cs @@ -1,7 +1,6 @@ //#define USE_GENERICS //#define DELETE_AFTER_TEST -using Microsoft.Extensions.Configuration; using Microsoft.Extensions.Logging.Abstractions; using Microsoft.Extensions.Options; using Orleans.Configuration; @@ -27,7 +26,7 @@ namespace UnitTests.Streaming.Reliability { [TestCategory("Streaming"), TestCategory("Reliability")] - public class StreamReliabilityTests : TestClusterPerTest + public class StreamReliabilityTests : BaseInProcessTestClusterFixture { private readonly ITestOutputHelper _output; public const string MEMORY_STREAM_PROVIDER_NAME = StreamTestsConstants.MEMORY_STREAM_PROVIDER_NAME; @@ -36,75 +35,60 @@ public class StreamReliabilityTests : TestClusterPerTest private Guid _streamId; private string _streamProviderName; private int _numExpectedSilos; + private IInternalClusterClient InternalClient => (IInternalClusterClient)this.Client; #if DELETE_AFTER_TEST private HashSet _usedGrains; #endif - protected override void ConfigureTestCluster(TestClusterBuilder builder) + protected override void ConfigureTestCluster(InProcessTestClusterBuilder builder) { TestUtils.CheckForAzureStorage(); this._numExpectedSilos = 2; - builder.CreateSiloAsync = StandaloneSiloHandle.CreateForAssembly(this.GetType().Assembly); builder.Options.InitialSilosCount = (short) this._numExpectedSilos; - builder.Options.UseTestClusterMembership = false; - builder.AddSiloBuilderConfigurator(); - builder.AddClientBuilderConfigurator(); + builder.ConfigureSilo((_, siloBuilder) => ConfigureSilo(siloBuilder)); + builder.ConfigureClient(ConfigureClient); } - public class ClientBuilderConfigurator : IClientBuilderConfigurator + private static void ConfigureClient(IClientBuilder clientBuilder) { - public void Configure(IConfiguration configuration, IClientBuilder clientBuilder) - { - clientBuilder.UseAzureStorageClustering(gatewayOptions => + clientBuilder.AddAzureQueueStreams(AZURE_QUEUE_STREAM_PROVIDER_NAME, ob => ob.Configure>( + (options, dep) => { - gatewayOptions.ConfigureTestDefaults(); - }) - .AddAzureQueueStreams(AZURE_QUEUE_STREAM_PROVIDER_NAME, ob => ob.Configure>( - (options, dep) => - { - options.ConfigureTestDefaults(); - options.QueueNames = AzureQueueUtilities.GenerateQueueNames(dep.Value.ClusterId, QueueCount); - })) - .AddMemoryStreams(MEMORY_STREAM_PROVIDER_NAME) - .Configure(options => options.GatewayListRefreshPeriod = TimeSpan.FromSeconds(5)); - } + options.ConfigureTestDefaults(); + options.QueueNames = AzureQueueUtilities.GenerateQueueNames(dep.Value.ClusterId, QueueCount); + })) + .AddMemoryStreams(MEMORY_STREAM_PROVIDER_NAME) + .Configure(options => options.GatewayListRefreshPeriod = TimeSpan.FromSeconds(5)); } - public class SiloBuilderConfigurator : ISiloConfigurator + private static void ConfigureSilo(ISiloBuilder hostBuilder) { - public void Configure(ISiloBuilder hostBuilder) + hostBuilder.AddAzureTableGrainStorage("AzureStore", builder => builder.Configure>((options, silo) => { - hostBuilder.UseAzureStorageClustering(options => - { - options.ConfigureTestDefaults(); - }) - .AddAzureTableGrainStorage("AzureStore", builder => builder.Configure>((options, silo) => - { - options.ConfigureTestDefaults(); - options.DeleteStateOnClear = true; - })) - .AddMemoryGrainStorage("MemoryStore", options => options.NumStorageGrains = 1) - .AddMemoryStreams(MEMORY_STREAM_PROVIDER_NAME) - .AddAzureTableGrainStorage("PubSubStore", builder => builder.Configure>((options, silo) => - { - options.DeleteStateOnClear = true; - options.ConfigureTestDefaults(); - })) - .AddAzureQueueStreams(AZURE_QUEUE_STREAM_PROVIDER_NAME, ob => ob.Configure>( + options.ConfigureTestDefaults(); + options.DeleteStateOnClear = true; + })) + .AddMemoryGrainStorage("MemoryStore", options => options.NumStorageGrains = 1) + .AddMemoryStreams(MEMORY_STREAM_PROVIDER_NAME) + .AddAzureTableGrainStorage("PubSubStore", builder => builder.Configure>((options, silo) => + { + options.DeleteStateOnClear = true; + options.ConfigureTestDefaults(); + })) + .AddAzureQueueStreams(AZURE_QUEUE_STREAM_PROVIDER_NAME, ob => ob.Configure>( (options, dep) => { options.ConfigureTestDefaults(); options.QueueNames = AzureQueueUtilities.GenerateQueueNames(dep.Value.ClusterId, QueueCount); })) - .AddAzureQueueStreams("AzureQueueProvider2", ob => ob.Configure>( + .AddAzureQueueStreams("AzureQueueProvider2", ob => ob.Configure>( (options, dep) => { options.ConfigureTestDefaults(); options.QueueNames = AzureQueueUtilities.GenerateQueueNames($"{dep.Value.ClusterId}2", QueueCount); })); - } } public StreamReliabilityTests(ITestOutputHelper output) @@ -152,8 +136,8 @@ public void Baseline_StreamRel() { // This test case is just a sanity-check that the silo test config is OK. const string testName = "Baseline_StreamRel"; - StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, logger, HostedCluster); - StreamTestUtils.LogEndTest(testName, logger); + StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, Logger, HostedCluster); + StreamTestUtils.LogEndTest(testName, Logger); } [SkippableFact, TestCategory("Functional")] @@ -161,7 +145,7 @@ public async Task Baseline_StreamRel_RestartSilos() { // This test case is just a sanity-check that the silo test config is OK. const string testName = "Baseline_StreamRel_RestartSilos"; - StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, logger, HostedCluster); + StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, Logger, HostedCluster); CheckSilosRunning("Before Restart", _numExpectedSilos); var silos = this.HostedCluster.Silos; @@ -171,7 +155,7 @@ public async Task Baseline_StreamRel_RestartSilos() Assert.NotEqual(silos, this.HostedCluster.Silos); // Should be different silos after restart - StreamTestUtils.LogEndTest(testName, logger); + StreamTestUtils.LogEndTest(testName, Logger); } [SkippableFact, TestCategory("Functional")] @@ -182,7 +166,7 @@ public async Task SMS_Baseline_StreamRel() _streamId = Guid.NewGuid(); _streamProviderName = MEMORY_STREAM_PROVIDER_NAME; - StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, logger, HostedCluster); + StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, Logger, HostedCluster); // Grain Producer -> Grain Consumer @@ -191,7 +175,7 @@ public async Task SMS_Baseline_StreamRel() await Do_BaselineTest(consumerGrainId, producerGrainId); - StreamTestUtils.LogEndTest(testName, logger); + StreamTestUtils.LogEndTest(testName, Logger); } [SkippableFact, TestCategory("Functional"), TestCategory("AzureStorage")] @@ -202,14 +186,14 @@ public async Task AQ_Baseline_StreamRel() _streamId = Guid.NewGuid(); _streamProviderName = AZURE_QUEUE_STREAM_PROVIDER_NAME; - StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, logger, HostedCluster); + StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, Logger, HostedCluster); long consumerGrainId = Random.Shared.Next(); long producerGrainId = Random.Shared.Next(); await Do_BaselineTest(consumerGrainId, producerGrainId); - StreamTestUtils.LogEndTest(testName, logger); + StreamTestUtils.LogEndTest(testName, Logger); } [SkippableFact(Skip ="Ignore"), TestCategory("Failures"), TestCategory("Streaming"), TestCategory("Reliability")] @@ -343,7 +327,7 @@ private async Task> Do_BaselineTest(long consum private async Task Do_BaselineTest(long consumerGrainId, long producerGrainId) #endif { - logger.LogInformation("Initializing: ConsumerGrain={ConsumerGrainId} ProducerGrain={ProducerGrainId}", consumerGrainId, producerGrainId); + Logger.LogInformation("Initializing: ConsumerGrain={ConsumerGrainId} ProducerGrain={ProducerGrainId}", consumerGrainId, producerGrainId); var consumerGrain = GetGrain(consumerGrainId); var producerGrain = GetGrain(producerGrainId); #if DELETE_AFTER_TEST @@ -356,9 +340,9 @@ private async Task Do_BaselineTest(long consumerGra string when = "Before subscribe"; await CheckConsumerProducerStatus(when, producerGrainId, consumerGrainId, false, false); - logger.LogInformation("AddConsumer: StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); + Logger.LogInformation("AddConsumer: StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); await consumerGrain.AddConsumer(_streamId, _streamProviderName); - logger.LogInformation("BecomeProducer: StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); + Logger.LogInformation("BecomeProducer: StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); await producerGrain.BecomeProducer(_streamId, _streamProviderName); when = "After subscribe"; @@ -381,7 +365,7 @@ private async Task[]> Do_AddConsumerGrains(long private async Task Do_AddConsumerGrains(long baseId, int numGrains) #endif { - logger.LogInformation("Initializing: BaseId={BaseId} NumGrains={NumGrains}", baseId, numGrains); + Logger.LogInformation("Initializing: BaseId={BaseId} NumGrains={NumGrains}", baseId, numGrains); #if USE_GENERICS var grains = new IStreamReliabilityTestGrain[numGrains]; @@ -400,7 +384,7 @@ private async Task Do_AddConsumerGrains(long base } await Task.WhenAll(promises); - logger.LogInformation("AddConsumer: StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); + Logger.LogInformation("AddConsumer: StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); await Task.WhenAll(grains.Select(g => g.AddConsumer(_streamId, _streamProviderName))); return grains; @@ -416,7 +400,7 @@ private async Task Test_AddMany_Consumers(string testName, string streamProvider _streamId = Guid.NewGuid(); _streamProviderName = streamProviderName; - StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, logger, HostedCluster); + StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, Logger, HostedCluster); long consumerGrainId = Random.Shared.Next(); long producerGrainId = Random.Shared.Next(); @@ -464,7 +448,7 @@ await Task.WhenAll(grains1.Select(async g => // Messages received by new consumer grains await Task.WhenAll(grains2.Select(g => CheckReceivedCounts(when2, g, numLoops, 0))); - StreamTestUtils.LogEndTest(testName, logger); + StreamTestUtils.LogEndTest(testName, Logger); } private async Task Test_PubSub_MultiConsumerSameGrain(string testName, string streamProviderName) @@ -472,7 +456,7 @@ private async Task Test_PubSub_MultiConsumerSameGrain(string testName, string st _streamId = Guid.NewGuid(); _streamProviderName = streamProviderName; - StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, logger, HostedCluster); + StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, Logger, HostedCluster); // Grain Producer -> Grain 2 x Consumer @@ -480,11 +464,11 @@ private async Task Test_PubSub_MultiConsumerSameGrain(string testName, string st long producerGrainId = Random.Shared.Next(); string when; - logger.LogInformation("Initializing: ConsumerGrain={ConsumerGrainId} ProducerGrain={ProducerGrainId}", consumerGrainId, producerGrainId); + Logger.LogInformation("Initializing: ConsumerGrain={ConsumerGrainId} ProducerGrain={ProducerGrainId}", consumerGrainId, producerGrainId); var consumerGrain = GetGrain(consumerGrainId); var producerGrain = GetGrain(producerGrainId); - logger.LogInformation("BecomeProducer: StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); + Logger.LogInformation("BecomeProducer: StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); await producerGrain.BecomeProducer(_streamId, _streamProviderName); when = "After BecomeProducer"; @@ -492,7 +476,7 @@ private async Task Test_PubSub_MultiConsumerSameGrain(string testName, string st await producerGrain.SendItem(0); await StreamTestUtils.CheckPubSubCounts(this.InternalClient, _output, when, 1, 0, _streamId, _streamProviderName, StreamTestsConstants.StreamReliabilityNamespace); - logger.LogInformation("AddConsumer x 2 : StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); + Logger.LogInformation("AddConsumer x 2 : StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); await consumerGrain.AddConsumer(_streamId, _streamProviderName); when = "After first AddConsumer"; await StreamTestUtils.CheckPubSubCounts(this.InternalClient, _output, when, 1, 1, _streamId, _streamProviderName, StreamTestsConstants.StreamReliabilityNamespace); @@ -501,7 +485,7 @@ private async Task Test_PubSub_MultiConsumerSameGrain(string testName, string st when = "After second AddConsumer"; await StreamTestUtils.CheckPubSubCounts(this.InternalClient, _output, when, 1, 2, _streamId, _streamProviderName, StreamTestsConstants.StreamReliabilityNamespace); - StreamTestUtils.LogEndTest(testName, logger); + StreamTestUtils.LogEndTest(testName, Logger); } private async Task Test_PubSub_MultiProducerSameGrain(string testName, string streamProviderName) @@ -509,7 +493,7 @@ private async Task Test_PubSub_MultiProducerSameGrain(string testName, string st _streamId = Guid.NewGuid(); _streamProviderName = streamProviderName; - StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, logger, HostedCluster); + StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, Logger, HostedCluster); // Grain Producer -> Grain 2 x Consumer @@ -517,11 +501,11 @@ private async Task Test_PubSub_MultiProducerSameGrain(string testName, string st long producerGrainId = Random.Shared.Next(); string when; - logger.LogInformation("Initializing: ConsumerGrain={ConsumerGrainId} ProducerGrain={ProducerGrainId}", consumerGrainId, producerGrainId); + Logger.LogInformation("Initializing: ConsumerGrain={ConsumerGrainId} ProducerGrain={ProducerGrainId}", consumerGrainId, producerGrainId); var consumerGrain = GetGrain(consumerGrainId); var producerGrain = GetGrain(producerGrainId); - logger.LogInformation("BecomeProducer: StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); + Logger.LogInformation("BecomeProducer: StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); await producerGrain.BecomeProducer(_streamId, _streamProviderName); when = "After first BecomeProducer"; // Note: Only semantics guarenteed for producer is that they will have been registered by time that first msg is sent. @@ -533,7 +517,7 @@ private async Task Test_PubSub_MultiProducerSameGrain(string testName, string st await producerGrain.SendItem(0); await StreamTestUtils.CheckPubSubCounts(this.InternalClient, _output, when, 1, 0, _streamId, _streamProviderName, StreamTestsConstants.StreamReliabilityNamespace); - logger.LogInformation("AddConsumer x 2 : StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); + Logger.LogInformation("AddConsumer x 2 : StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); await consumerGrain.AddConsumer(_streamId, _streamProviderName); when = "After first AddConsumer"; await StreamTestUtils.CheckPubSubCounts(this.InternalClient, _output, when, 1, 1, _streamId, _streamProviderName, StreamTestsConstants.StreamReliabilityNamespace); @@ -542,7 +526,7 @@ private async Task Test_PubSub_MultiProducerSameGrain(string testName, string st when = "After second AddConsumer"; await StreamTestUtils.CheckPubSubCounts(this.InternalClient, _output, when, 1, 2, _streamId, _streamProviderName, StreamTestsConstants.StreamReliabilityNamespace); - StreamTestUtils.LogEndTest(testName, logger); + StreamTestUtils.LogEndTest(testName, Logger); } private async Task Test_PubSub_Unsubscribe(string testName, string streamProviderName) @@ -550,7 +534,7 @@ private async Task Test_PubSub_Unsubscribe(string testName, string streamProvide _streamId = Guid.NewGuid(); _streamProviderName = streamProviderName; - StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, logger, HostedCluster); + StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, Logger, HostedCluster); // Grain Producer -> Grain 2 x Consumer // Note: PubSub should only count distinct grains, even if a grain has multiple consumer handles @@ -559,11 +543,11 @@ private async Task Test_PubSub_Unsubscribe(string testName, string streamProvide long producerGrainId = Random.Shared.Next(); string when; - logger.LogInformation("Initializing: ConsumerGrain={ConsumerGrainId} ProducerGrain={ProducerGrainId}", consumerGrainId, producerGrainId); + Logger.LogInformation("Initializing: ConsumerGrain={ConsumerGrainId} ProducerGrain={ProducerGrainId}", consumerGrainId, producerGrainId); var consumerGrain = GetGrain(consumerGrainId); var producerGrain = GetGrain(producerGrainId); - logger.LogInformation("BecomeProducer: StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); + Logger.LogInformation("BecomeProducer: StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); await producerGrain.BecomeProducer(_streamId, _streamProviderName); await producerGrain.BecomeProducer(_streamId, _streamProviderName); when = "After BecomeProducer"; @@ -571,7 +555,7 @@ private async Task Test_PubSub_Unsubscribe(string testName, string streamProvide await producerGrain.SendItem(0); await StreamTestUtils.CheckPubSubCounts(this.InternalClient, _output, when, 1, 0, _streamId, _streamProviderName, StreamTestsConstants.StreamReliabilityNamespace); - logger.LogInformation("AddConsumer x 2 : StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); + Logger.LogInformation("AddConsumer x 2 : StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); var c1 = await consumerGrain.AddConsumer(_streamId, _streamProviderName); when = "After first AddConsumer"; await StreamTestUtils.CheckPubSubCounts(this.InternalClient, _output, when, 1, 1, _streamId, _streamProviderName, StreamTestsConstants.StreamReliabilityNamespace); @@ -581,19 +565,19 @@ private async Task Test_PubSub_Unsubscribe(string testName, string streamProvide await StreamTestUtils.CheckPubSubCounts(this.InternalClient, _output, when, 1, 2, _streamId, _streamProviderName, StreamTestsConstants.StreamReliabilityNamespace); await CheckConsumerCounts(when, consumerGrain, 2); - logger.LogInformation("RemoveConsumer: StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); + Logger.LogInformation("RemoveConsumer: StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); await consumerGrain.RemoveConsumer(_streamId, _streamProviderName, c1); when = "After first RemoveConsumer"; await StreamTestUtils.CheckPubSubCounts(this.InternalClient, _output, when, 1, 1, _streamId, _streamProviderName, StreamTestsConstants.StreamReliabilityNamespace); await CheckConsumerCounts(when, consumerGrain, 1); #if REMOVE_PRODUCER - logger.LogInformation("RemoveProducer: StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); + Logger.LogInformation("RemoveProducer: StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); await producerGrain.RemoveProducer(_streamId, _streamProviderName); when = "After RemoveProducer"; await CheckPubSubCounts(when, 0, 1); await CheckConsumerCounts(when, consumerGrain, 1); #endif - logger.LogInformation("RemoveConsumer: StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); + Logger.LogInformation("RemoveConsumer: StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); await consumerGrain.RemoveConsumer(_streamId, _streamProviderName, c2); when = "After second RemoveConsumer"; #if REMOVE_PRODUCER @@ -603,7 +587,7 @@ private async Task Test_PubSub_Unsubscribe(string testName, string streamProvide #endif await CheckConsumerCounts(when, consumerGrain, 0); - StreamTestUtils.LogEndTest(testName, logger); + StreamTestUtils.LogEndTest(testName, Logger); } [SkippableFact, TestCategory("Functional")] @@ -613,12 +597,12 @@ public async Task SMS_AllSilosRestart_UnsubscribeConsumer() _streamId = Guid.NewGuid(); _streamProviderName = MEMORY_STREAM_PROVIDER_NAME; - StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, logger, HostedCluster); + StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, Logger, HostedCluster); long consumerGrainId = Random.Shared.Next(); var consumerGrain = this.GrainFactory.GetGrain(consumerGrainId); - logger.LogInformation("Subscribe: StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); + Logger.LogInformation("Subscribe: StreamId={StreamId} Provider={Provider}", _streamId, _streamProviderName); await consumerGrain.Subscribe(_streamId, _streamProviderName); // Restart silos @@ -645,7 +629,7 @@ public async Task SMS_AllSilosRestart_UnsubscribeConsumer() await Task.Delay(100); } - StreamTestUtils.LogEndTest(testName, logger); + StreamTestUtils.LogEndTest(testName, Logger); } private async Task Test_AllSilosRestart(string testName, string streamProviderName) @@ -653,7 +637,7 @@ private async Task Test_AllSilosRestart(string testName, string streamProviderNa _streamId = Guid.NewGuid(); _streamProviderName = streamProviderName; - StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, logger, HostedCluster); + StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, Logger, HostedCluster); long consumerGrainId = Random.Shared.Next(); long producerGrainId = Random.Shared.Next(); @@ -671,7 +655,7 @@ private async Task Test_AllSilosRestart(string testName, string streamProviderNa await producerGrain.SendItem(1); await CheckConsumerProducerStatus(when, producerGrainId, consumerGrainId, true, true); - StreamTestUtils.LogEndTest(testName, logger); + StreamTestUtils.LogEndTest(testName, Logger); } private async Task Test_AllSilosRestart_PubSubCounts(string testName, string streamProviderName) @@ -679,7 +663,7 @@ private async Task Test_AllSilosRestart_PubSubCounts(string testName, string str _streamId = Guid.NewGuid(); _streamProviderName = streamProviderName; - StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, logger, HostedCluster); + StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, Logger, HostedCluster); long consumerGrainId = Random.Shared.Next(); long producerGrainId = Random.Shared.Next(); @@ -712,7 +696,7 @@ private async Task Test_AllSilosRestart_PubSubCounts(string testName, string str var consumerGrain = GetGrain(consumerGrainId); await CheckReceivedCounts(when, consumerGrain, 1, 0); - StreamTestUtils.LogEndTest(testName, logger); + StreamTestUtils.LogEndTest(testName, Logger); } private async Task Test_SiloDies_Consumer(string testName, string streamProviderName) @@ -721,7 +705,7 @@ private async Task Test_SiloDies_Consumer(string testName, string streamProvider _streamProviderName = streamProviderName; string when; - StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, logger, HostedCluster); + StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, Logger, HostedCluster); long consumerGrainId = Random.Shared.Next(); long producerGrainId = Random.Shared.Next(); @@ -731,16 +715,15 @@ private async Task Test_SiloDies_Consumer(string testName, string streamProvider when = "Before kill one silo"; CheckSilosRunning(when, _numExpectedSilos); - bool sameSilo = await CheckGrainCounts(); - // Find which silo the consumer grain is located on var consumerGrain = GetGrain(consumerGrainId); SiloAddress siloAddress = await consumerGrain.GetLocation(); + SiloAddress producerAddress = await producerGrain.GetLocation(); - _output.WriteLine("Consumer grain is located on silo {0} ; Producer on same silo = {1}", siloAddress, sameSilo); + _output.WriteLine("Consumer grain is located on silo {0} ; Producer grain is located on silo {1}", siloAddress, producerAddress); // Kill the silo containing the consumer grain - SiloHandle siloToKill = this.HostedCluster.Silos.First(s => s.SiloAddress.Equals(siloAddress)); + var siloToKill = this.HostedCluster.Silos.First(s => s.SiloAddress.Equals(siloAddress)); await StopSilo(siloToKill, true, false); // Note: Don't restart failed silo for this test case // Note: Don't reinitialize client @@ -752,7 +735,7 @@ private async Task Test_SiloDies_Consumer(string testName, string streamProvider await producerGrain.SendItem(1); await CheckConsumerProducerStatus(when, producerGrainId, consumerGrainId, true, true); - StreamTestUtils.LogEndTest(testName, logger); + StreamTestUtils.LogEndTest(testName, Logger); } private async Task Test_SiloDies_Producer(string testName, string streamProviderName) @@ -761,7 +744,7 @@ private async Task Test_SiloDies_Producer(string testName, string streamProvider _streamProviderName = streamProviderName; string when; - StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, logger, HostedCluster); + StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, Logger, HostedCluster); long consumerGrainId = Random.Shared.Next(); long producerGrainId = Random.Shared.Next(); @@ -771,14 +754,14 @@ private async Task Test_SiloDies_Producer(string testName, string streamProvider when = "Before kill one silo"; CheckSilosRunning(when, _numExpectedSilos); - bool sameSilo = await CheckGrainCounts(); - // Find which silo the producer grain is located on SiloAddress siloAddress = await producerGrain.GetLocation(); - _output.WriteLine("Producer grain is located on silo {0} ; Consumer on same silo = {1}", siloAddress, sameSilo); + var consumerGrain = GetGrain(consumerGrainId); + SiloAddress consumerAddress = await consumerGrain.GetLocation(); + _output.WriteLine("Producer grain is located on silo {0} ; Consumer grain is located on silo {1}", siloAddress, consumerAddress); // Kill the silo containing the producer grain - SiloHandle siloToKill = this.HostedCluster.Silos.First(s => s.SiloAddress.Equals(siloAddress)); + var siloToKill = this.HostedCluster.Silos.First(s => s.SiloAddress.Equals(siloAddress)); await StopSilo(siloToKill, true, false); // Note: Don't restart failed silo for this test case // Note: Don't reinitialize client @@ -790,7 +773,7 @@ private async Task Test_SiloDies_Producer(string testName, string streamProvider await producerGrain.SendItem(1); await CheckConsumerProducerStatus(when, producerGrainId, consumerGrainId, true, true); - StreamTestUtils.LogEndTest(testName, logger); + StreamTestUtils.LogEndTest(testName, Logger); } private async Task Test_SiloRestarts_Consumer(string testName, string streamProviderName) @@ -799,7 +782,7 @@ private async Task Test_SiloRestarts_Consumer(string testName, string streamProv _streamProviderName = streamProviderName; string when; - StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, logger, HostedCluster); + StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, Logger, HostedCluster); long consumerGrainId = Random.Shared.Next(); long producerGrainId = Random.Shared.Next(); @@ -809,16 +792,15 @@ private async Task Test_SiloRestarts_Consumer(string testName, string streamProv when = "Before restart one silo"; CheckSilosRunning(when, _numExpectedSilos); - bool sameSilo = await CheckGrainCounts(); - // Find which silo the consumer grain is located on var consumerGrain = GetGrain(consumerGrainId); SiloAddress siloAddress = await consumerGrain.GetLocation(); + SiloAddress producerAddress = await producerGrain.GetLocation(); - _output.WriteLine("Consumer grain is located on silo {0} ; Producer on same silo = {1}", siloAddress, sameSilo); + _output.WriteLine("Consumer grain is located on silo {0} ; Producer grain is located on silo {1}", siloAddress, producerAddress); // Restart the silo containing the consumer grain - SiloHandle siloToKill = this.HostedCluster.Silos.First(s => s.SiloAddress.Equals(siloAddress)); + var siloToKill = this.HostedCluster.Silos.First(s => s.SiloAddress.Equals(siloAddress)); await StopSilo(siloToKill, true, true); // Note: Don't reinitialize client @@ -829,7 +811,7 @@ private async Task Test_SiloRestarts_Consumer(string testName, string streamProv await producerGrain.SendItem(1); await CheckConsumerProducerStatus(when, producerGrainId, consumerGrainId, true, true); - StreamTestUtils.LogEndTest(testName, logger); + StreamTestUtils.LogEndTest(testName, Logger); } private async Task Test_SiloRestarts_Producer(string testName, string streamProviderName) @@ -838,7 +820,7 @@ private async Task Test_SiloRestarts_Producer(string testName, string streamProv _streamProviderName = streamProviderName; string when; - StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, logger, HostedCluster); + StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, Logger, HostedCluster); long consumerGrainId = Random.Shared.Next(); long producerGrainId = Random.Shared.Next(); @@ -848,15 +830,15 @@ private async Task Test_SiloRestarts_Producer(string testName, string streamProv when = "Before restart one silo"; CheckSilosRunning(when, _numExpectedSilos); - bool sameSilo = await CheckGrainCounts(); - // Find which silo the producer grain is located on SiloAddress siloAddress = await producerGrain.GetLocation(); + var consumerGrain = GetGrain(consumerGrainId); + SiloAddress consumerAddress = await consumerGrain.GetLocation(); - _output.WriteLine("Producer grain is located on silo {0} ; Consumer on same silo = {1}", siloAddress, sameSilo); + _output.WriteLine("Producer grain is located on silo {0} ; Consumer grain is located on silo {1}", siloAddress, consumerAddress); // Restart the silo containing the consumer grain - SiloHandle siloToKill = this.HostedCluster.Silos.First(s => s.SiloAddress.Equals(siloAddress)); + var siloToKill = this.HostedCluster.Silos.First(s => s.SiloAddress.Equals(siloAddress)); await StopSilo(siloToKill, true, true); // Note: Don't reinitialize client @@ -867,7 +849,7 @@ private async Task Test_SiloRestarts_Producer(string testName, string streamProv await producerGrain.SendItem(1); await CheckConsumerProducerStatus(when, producerGrainId, consumerGrainId, true, true); - StreamTestUtils.LogEndTest(testName, logger); + StreamTestUtils.LogEndTest(testName, Logger); } private async Task Test_SiloJoins(string testName, string streamProviderName) @@ -877,7 +859,7 @@ private async Task Test_SiloJoins(string testName, string streamProviderName) const int numLoops = 3; - StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, logger, HostedCluster); + StreamTestUtils.LogStartTest(testName, _streamId, _streamProviderName, Logger, HostedCluster); long consumerGrainId = Random.Shared.Next(); long producerGrainId = Random.Shared.Next(); @@ -906,7 +888,7 @@ private async Task Test_SiloJoins(string testName, string streamProviderName) // Add new silo //SiloHandle newSilo = StartAdditionalOrleans(); //WaitForLivenessToStabilize(); - SiloHandle newSilo = await this.HostedCluster.StartAdditionalSiloAsync(); + var newSilo = await this.HostedCluster.StartAdditionalSiloAsync(); await this.HostedCluster.WaitForLivenessToStabilizeAsync(); @@ -943,17 +925,18 @@ private async Task Test_SiloJoins(string testName, string streamProviderName) // New consumer received the newly published messages await CheckReceivedCounts(when+"-New", newConsumer, numLoops, 0); - StreamTestUtils.LogEndTest(testName, logger); + StreamTestUtils.LogEndTest(testName, Logger); } // ---------- Utility Functions ---------- private async Task RestartAllSilos() { + var oldSilos = this.HostedCluster.GetActiveSilos().Select(silo => silo.SiloAddress).ToArray(); _output.WriteLine("\n\n\n\n-----------------------------------------------------\n" + - "Restarting all silos - Old Primary={0} Secondary={1}" + + "Restarting all silos - Old Silos={0}" + "\n-----------------------------------------------------\n\n\n", - this.HostedCluster.Primary?.SiloAddress, this.HostedCluster.SecondarySilos.FirstOrDefault()?.SiloAddress); + string.Join(", ", oldSilos.Select(silo => silo.ToString()))); foreach (var silo in this.HostedCluster.GetActiveSilos().ToList()) { @@ -963,17 +946,17 @@ private async Task RestartAllSilos() // Note: Needed to reinitialize client in this test case to connect to new silos // this.HostedCluster.InitializeClient(); + var newSilos = this.HostedCluster.GetActiveSilos().Select(silo => silo.SiloAddress).ToArray(); _output.WriteLine("\n\n\n\n-----------------------------------------------------\n" + - "Restarted new silos - New Primary={0} Secondary={1}" + + "Restarted new silos - New Silos={0}" + "\n-----------------------------------------------------\n\n\n", - this.HostedCluster.Primary?.SiloAddress, this.HostedCluster.SecondarySilos.FirstOrDefault()?.SiloAddress); + string.Join(", ", newSilos.Select(silo => silo.ToString()))); } - private async Task StopSilo(SiloHandle silo, bool kill, bool restart) + private async Task StopSilo(InProcessSiloHandle silo, bool kill, bool restart) { SiloAddress oldSilo = silo.SiloAddress; - bool isPrimary = oldSilo.Equals(this.HostedCluster.Primary?.SiloAddress); - string siloType = isPrimary ? "Primary" : "Secondary"; + string siloType = silo.InstanceNumber == 0 ? "Primary" : "Secondary"; var action = (restart, kill) switch { (true, true) => "Kill and restart", @@ -982,30 +965,31 @@ private async Task StopSilo(SiloHandle silo, bool kill, bool restart) (false, false) => "Stop", }; - logger.LogWarning("{Action} {SiloType} silo {OldSilo}", action, siloType, oldSilo); + Logger.LogWarning("{Action} {SiloType} silo {OldSilo}", action, siloType, oldSilo); if (restart) { //RestartRuntime(silo, kill); - SiloHandle newSilo = await this.HostedCluster.RestartSiloAsync(silo); + var newSilo = await this.HostedCluster.RestartSiloAsync(silo); + Assert.NotNull(newSilo); - logger.LogInformation("Restarted new {SiloType} silo {SiloAddress}", siloType, newSilo.SiloAddress); + Logger.LogInformation("Restarted new {SiloType} silo {SiloAddress}", siloType, newSilo.SiloAddress); Assert.NotEqual(oldSilo, newSilo.SiloAddress); //"Should be different silo address after Restart" } else if (kill) { - await this.HostedCluster.KillSiloAsync(silo); - Assert.False(silo.IsActive); + await this.HostedCluster.KillSiloAsync(silo); + Assert.False(silo.IsActive); } else { - await this.HostedCluster.StopSiloAsync(silo); - Assert.False(silo.IsActive); + await this.HostedCluster.StopSiloAsync(silo); + Assert.False(silo.IsActive); } // WaitForLivenessToStabilize(!kill); - this.HostedCluster.WaitForLivenessToStabilizeAsync(kill).Wait(); + await this.HostedCluster.WaitForLivenessToStabilizeAsync(kill); } #if USE_GENERICS @@ -1070,36 +1054,6 @@ private void CheckSilosRunning(string when, int expectedNumSilos) { Assert.Equal(expectedNumSilos, this.HostedCluster.GetActiveSilos().Count()); } - protected async Task CheckGrainCounts() - { -#if USE_GENERICS - string grainType = RuntimeTypeNameFormatter.Format(typeof(StreamReliabilityTestGrain)); -#else - string grainType = RuntimeTypeNameFormatter.Format(typeof(StreamReliabilityTestGrain)); -#endif - IManagementGrain mgmtGrain = this.GrainFactory.GetGrain(0); - - SimpleGrainStatistic[] grainStats = await mgmtGrain.GetSimpleGrainStatistics(); - _output.WriteLine("Found grains " + Utils.EnumerableToString(grainStats)); - - var grainLocs = grainStats.Where(gs => gs.GrainType == grainType).ToArray(); - - Assert.True(grainLocs.Length > 0, "Found too few grains"); - Assert.True(grainLocs.Length <= 2, "Found too many grains " + grainLocs.Length); - - bool sameSilo = grainLocs.Length == 1; - if (sameSilo) - { - StreamTestUtils.Assert_AreEqual(_output, 2, grainLocs[0].ActivationCount, "Num grains on same Silo " + grainLocs[0].SiloAddress); - } - else - { - StreamTestUtils.Assert_AreEqual(_output, 1, grainLocs[0].ActivationCount, "Num grains on Silo " + grainLocs[0].SiloAddress); - StreamTestUtils.Assert_AreEqual(_output, 1, grainLocs[1].ActivationCount, "Num grains on Silo " + grainLocs[1].SiloAddress); - } - return sameSilo; - } - #if USE_GENERICS protected async Task CheckReceivedCounts(string when, IStreamReliabilityTestGrain consumerGrain, int expectedReceivedCount, int expectedErrorsCount) #else diff --git a/test/Orleans.Streaming.Tests/StreamingTests/StreamTestUtils.cs b/test/Orleans.Streaming.Tests/StreamingTests/StreamTestUtils.cs index 8f64fcf280e..01ad0297df7 100644 --- a/test/Orleans.Streaming.Tests/StreamingTests/StreamTestUtils.cs +++ b/test/Orleans.Streaming.Tests/StreamingTests/StreamTestUtils.cs @@ -16,6 +16,19 @@ internal static void LogStartTest(string testName, Guid streamId, string streamP { SiloAddress primSilo = siloHost.Primary?.SiloAddress; SiloAddress secSilo = siloHost.SecondarySilos.FirstOrDefault()?.SiloAddress; + LogStartTest(testName, streamId, streamProviderName, logger, primSilo, secSilo); + } + + internal static void LogStartTest(string testName, Guid streamId, string streamProviderName, ILogger logger, InProcessTestCluster siloHost) + { + var silos = siloHost.GetActiveSilos().Select(silo => silo.SiloAddress).ToArray(); + SiloAddress primSilo = silos.FirstOrDefault(); + SiloAddress secSilo = silos.Skip(1).FirstOrDefault(); + LogStartTest(testName, streamId, streamProviderName, logger, primSilo, secSilo); + } + + private static void LogStartTest(string testName, Guid streamId, string streamProviderName, ILogger logger, SiloAddress primSilo, SiloAddress secSilo) + { logger.LogInformation( "\n\n**START********************** {TestName} ********************************* \n\n" + "Running with initial silos Primary={PrimarySilo} Secondary={SecondarySilo} StreamId={StreamId} StreamProviderName={StreamProviderName} \n\n",