diff --git a/sdk/eventhub/Microsoft.Azure.WebJobs.Extensions.EventHubs/tests/EventHubEndToEndTests.cs b/sdk/eventhub/Microsoft.Azure.WebJobs.Extensions.EventHubs/tests/EventHubEndToEndTests.cs index 68d81103e6f3..d921d59c0fef 100644 --- a/sdk/eventhub/Microsoft.Azure.WebJobs.Extensions.EventHubs/tests/EventHubEndToEndTests.cs +++ b/sdk/eventhub/Microsoft.Azure.WebJobs.Extensions.EventHubs/tests/EventHubEndToEndTests.cs @@ -31,7 +31,8 @@ public class EventHubEndToEndTests : WebJobsEventHubTestBase { private static readonly TimeSpan NoEventReadTimeout = TimeSpan.FromSeconds(5); - private static EventWaitHandle _eventWait; + private static EventWaitHandle _eventWait1; + private static EventWaitHandle _eventWait2; private static List _results; private static DateTimeOffset _initialOffsetEnqueuedTimeUTC; @@ -39,13 +40,15 @@ public class EventHubEndToEndTests : WebJobsEventHubTestBase public void SetUp() { _results = new List(); - _eventWait = new ManualResetEvent(initialState: false); + _eventWait1 = new ManualResetEvent(initialState: false); + _eventWait2 = new ManualResetEvent(initialState: false); } [TearDown] public void TearDown() { - _eventWait.Dispose(); + _eventWait1.Dispose(); + _eventWait2.Dispose(); } [Test] @@ -56,7 +59,7 @@ public async Task EventHub_PocoBinding() { await jobHost.CallAsync(nameof(EventHubTestBindToPocoJobs.SendEvent_TestHub)); - bool result = _eventWait.WaitOne(Timeout); + bool result = _eventWait1.WaitOne(Timeout); Assert.True(result); } @@ -75,7 +78,7 @@ public async Task EventHub_StringBinding() { await jobHost.CallAsync(nameof(EventHubTestBindToStringJobs.SendEvent_TestHub), new { input = "data" }); - bool result = _eventWait.WaitOne(Timeout); + bool result = _eventWait1.WaitOne(Timeout); Assert.True(result); var logs = host.GetTestLoggerProvider().GetAllLogMessages().Select(p => p.FormattedMessage); @@ -94,7 +97,7 @@ public async Task EventHub_SingleDispatch() { await jobHost.CallAsync(nameof(EventHubTestSingleDispatchJobs.SendEvent_TestHub), new { input = "data" }); - bool result = _eventWait.WaitOne(Timeout); + bool result = _eventWait1.WaitOne(Timeout); Assert.True(result); await StopWithDrainAsync(host); @@ -108,11 +111,11 @@ public async Task EventHub_SingleDispatch_Dispose() { await using var producer = new EventHubProducerClient(EventHubsTestEnvironment.Instance.EventHubsConnectionString, _eventHubScope.EventHubName); await producer.SendAsync(new EventData[] { new EventData(new BinaryData("data")) }); - var (_, host) = BuildHost(ConfigureTestEventHub); + var (jobHost, _) = BuildHost(ConfigureTestEventHub); - bool result = _eventWait.WaitOne(Timeout); - Assert.True(result); - host.Dispose(); + Assert.True(_eventWait1.WaitOne(Timeout)); + jobHost.Dispose(); + Assert.True(_eventWait2.WaitOne(Timeout)); } [Test] @@ -120,11 +123,11 @@ public async Task EventHub_SingleDispatch_StopWithoutDrain() { await using var producer = new EventHubProducerClient(EventHubsTestEnvironment.Instance.EventHubsConnectionString, _eventHubScope.EventHubName); await producer.SendAsync(new EventData[] { new EventData(new BinaryData("data")) }); - var (_, host) = BuildHost(ConfigureTestEventHub); + var (jobHost, _) = BuildHost(ConfigureTestEventHub); - bool result = _eventWait.WaitOne(Timeout); - Assert.True(result); - await host.StopAsync(); + Assert.True(_eventWait1.WaitOne(Timeout)); + await jobHost.StopAsync(); + Assert.True(_eventWait2.WaitOne(Timeout)); } [Test] @@ -146,7 +149,7 @@ public async Task EventHub_SingleDispatch_ConsumerGroup() { await jobHost.CallAsync(nameof(EventHubTestSingleDispatchWithConsumerGroupJobs.SendEvent_TestHub)); - bool result = _eventWait.WaitOne(Timeout); + bool result = _eventWait1.WaitOne(Timeout); Assert.True(result); } } @@ -159,7 +162,7 @@ public async Task EventHub_SingleDispatch_BinaryData() { await jobHost.CallAsync(nameof(EventHubTestSingleDispatchJobsBinaryData.SendEvent_TestHub), new { input = "data" }); - bool result = _eventWait.WaitOne(Timeout); + bool result = _eventWait1.WaitOne(Timeout); Assert.True(result); await StopWithDrainAsync(host); @@ -176,7 +179,7 @@ public async Task EventHub_ProducerClient() { await jobHost.CallAsync(nameof(EventHubTestClientDispatch.SendEvents)); - bool result = _eventWait.WaitOne(Timeout); + bool result = _eventWait1.WaitOne(Timeout); Assert.True(result); } } @@ -189,7 +192,7 @@ public async Task EventHub_Collector() { await jobHost.CallAsync(nameof(EventHubTestCollectorDispatch.SendEvents)); - bool result = _eventWait.WaitOne(Timeout); + bool result = _eventWait1.WaitOne(Timeout); Assert.True(result); } } @@ -202,7 +205,7 @@ public async Task EventHub_CollectorPartitionKey() { await jobHost.CallAsync(nameof(EventHubTestCollectorDispatch.SendEventsWithKey)); - bool result = _eventWait.WaitOne(Timeout); + bool result = _eventWait1.WaitOne(Timeout); Assert.True(result); } } @@ -320,7 +323,7 @@ public async Task AssertCanSendReceiveMessage(Action hostConfigura { await jobHost.CallAsync(nameof(EventHubTestSingleDispatchJobWithConnection.SendEvent_TestHub), new { input = "data" }); - bool result = _eventWait.WaitOne(Timeout); + bool result = _eventWait1.WaitOne(Timeout); Assert.True(result); } } @@ -334,7 +337,7 @@ public async Task EventHub_MultipleDispatch() int numEvents = 5; await jobHost.CallAsync(nameof(EventHubTestMultipleDispatchJobs.SendEvents_TestHub), new { numEvents = numEvents, input = "data" }); - bool result = _eventWait.WaitOne(Timeout); + bool result = _eventWait1.WaitOne(Timeout); Assert.True(result); await StopWithDrainAsync(host); @@ -352,7 +355,7 @@ public async Task EventHub_MultipleDispatch_BinaryData() int numEvents = 5; await jobHost.CallAsync(nameof(EventHubTestMultipleDispatchJobsBinaryData.SendEvents_TestHub), new { numEvents = numEvents, input = "data" }); - bool result = _eventWait.WaitOne(Timeout); + bool result = _eventWait1.WaitOne(Timeout); Assert.True(result); await StopWithDrainAsync(host); @@ -389,7 +392,7 @@ public async Task EventHub_MultipleDispatch_MinBatchSize() int numEvents = 5; await jobHost.CallAsync(nameof(EventHubTestMultipleDispatchMinBatchSizeJobs.SendEvents_TestHub), new { numEvents = numEvents, input = "data" }); - bool result = _eventWait.WaitOne(Timeout); + bool result = _eventWait1.WaitOne(Timeout); Assert.True(result); await StopWithDrainAsync(host); @@ -403,23 +406,23 @@ public async Task EventHub_MultipleDispatch_Dispose() { await using var producer = new EventHubProducerClient(EventHubsTestEnvironment.Instance.EventHubsConnectionString, _eventHubScope.EventHubName); await producer.SendAsync(new EventData[] { new EventData(new BinaryData("data")) }); - var (_, host) = BuildHost(); + var (jobHost, _) = BuildHost(); - bool result = _eventWait.WaitOne(Timeout); - Assert.True(result); - host.Dispose(); - } + Assert.True(_eventWait1.WaitOne(Timeout)); + jobHost.Dispose(); + Assert.True(_eventWait2.WaitOne(Timeout)); + } [Test] public async Task EventHub_MultipleDispatch_StopWithoutDrain() { await using var producer = new EventHubProducerClient(EventHubsTestEnvironment.Instance.EventHubsConnectionString, _eventHubScope.EventHubName); await producer.SendAsync(new EventData[] { new EventData(new BinaryData("data")) }); - var (_, host) = BuildHost(); + var (jobHost, _) = BuildHost(); - bool result = _eventWait.WaitOne(Timeout); - Assert.True(result); - await host.StopAsync(); + Assert.True(_eventWait1.WaitOne(Timeout)); + await jobHost.StopAsync(); + Assert.True(_eventWait2.WaitOne(Timeout)); } [Test] @@ -429,7 +432,7 @@ public async Task EventHub_PartitionKey() using (jobHost) { await jobHost.CallAsync(nameof(EventHubPartitionKeyTestJobs.SendEvents_TestHub), new { input = "data" }); - bool result = _eventWait.WaitOne(Timeout); + bool result = _eventWait1.WaitOne(Timeout); Assert.True(result); } @@ -455,7 +458,7 @@ public async Task EventHub_InitialOffsetFromStart() }); using (jobHost) { - bool result = _eventWait.WaitOne(Timeout); + bool result = _eventWait1.WaitOne(Timeout); Assert.True(result); } } @@ -482,7 +485,7 @@ public async Task EventHub_InitialOffsetFromEnd() using (jobHost) { // We don't expect to get signaled as there should be no messages received with a FromEnd initial offset - bool result = _eventWait.WaitOne(NoEventReadTimeout); + bool result = _eventWait1.WaitOne(NoEventReadTimeout); Assert.False(result, "An event was received while none were expected."); // send events which should be received. To ensure that the test is @@ -499,7 +502,7 @@ public async Task EventHub_InitialOffsetFromEnd() } }); - result = _eventWait.WaitOne(Timeout); + result = _eventWait1.WaitOne(Timeout); cts.Cancel(); try { await sendTask; } catch { /* Ignore, we're not testing sends */ } @@ -549,7 +552,7 @@ public async Task EventHub_InitialOffsetFromEnqueuedTime() }); using (jobHost) { - bool result = _eventWait.WaitOne(Timeout); + bool result = _eventWait1.WaitOne(Timeout); Assert.True(result); } } @@ -640,7 +643,7 @@ public static void ProcessSingleEvent([EventHubTrigger(TestHubName, Connection = Assert.AreNotEqual(default(LastEnqueuedEventProperties), triggerPartitionContext.ReadLastEnqueuedEventProperties()); Assert.True(triggerPartitionContext.IsCheckpointingAfterInvocation); - _eventWait.Set(); + _eventWait1.Set(); } } @@ -648,10 +651,13 @@ public class EventHubTestSingleDispatchJobs_Dispose { public static async Task SendEvent_TestHub([EventHubTrigger(TestHubName, Connection = TestHubName)] string evt, CancellationToken cancellationToken) { - _eventWait.Set(); - // wait a small amount of time for the host to call dispose - await Task.Delay(3000, CancellationToken.None); - Assert.IsTrue(cancellationToken.IsCancellationRequested); + _eventWait1.Set(); + // wait for the host to call dispose + while (!cancellationToken.IsCancellationRequested) + { + await Task.Delay(100, CancellationToken.None); + } + _eventWait2.Set(); } } @@ -659,10 +665,13 @@ public class EventHubTestMultipleDispatchJobs_Dispose { public static async Task SendEvent_TestHub([EventHubTrigger(TestHubName, Connection = TestHubName)] string[] evt, CancellationToken cancellationToken) { - _eventWait.Set(); - // wait a small amount of time for the host to call dispose - await Task.Delay(3000, CancellationToken.None); - Assert.IsTrue(cancellationToken.IsCancellationRequested); + _eventWait1.Set(); + // wait for the host to call dispose + while (!cancellationToken.IsCancellationRequested) + { + await Task.Delay(1000, CancellationToken.None); + } + _eventWait2.Set(); } } @@ -693,7 +702,7 @@ public static void ProcessSingleEvent([EventHubTrigger(TestHubName, Connection = Assert.AreEqual(eventData.PartitionKey, s_partitionKey); } - _eventWait.Set(); + _eventWait1.Set(); } } @@ -710,7 +719,7 @@ await producer.SendAsync(new[] public static void ProcessSingleEvent([EventHubTrigger(TestHubName, Connection = TestHubName)] EventData eventData) { Assert.AreEqual(eventData.EventBody.ToString(), "Event 1"); - _eventWait.Set(); + _eventWait1.Set(); } } @@ -725,7 +734,7 @@ public static void ProcessSingleEvent([EventHubTrigger(TestHubName, Connection = { Assert.AreEqual(evt, nameof(EventHubTestSingleDispatchWithConsumerGroupJobs)); - _eventWait.Set(); + _eventWait1.Set(); } } @@ -741,7 +750,7 @@ public static void ProcessSingleEvent([EventHubTrigger(TestHubName, Connection = IDictionary systemProperties) { Assert.AreEqual("data", evt.ToString()); - _eventWait.Set(); + _eventWait1.Set(); } } @@ -757,7 +766,7 @@ public static void BindToPoco([EventHubTrigger(TestHubName, Connection = TestHub Assert.AreEqual(input.Value, "data"); Assert.AreEqual(input.Name, "foo"); logger.LogInformation($"PocoValues(foo,data)"); - _eventWait.Set(); + _eventWait1.Set(); } } @@ -771,7 +780,7 @@ public static void SendEvent_TestHub(string input, [EventHub(TestHubName, Connec public static void BindToString([EventHubTrigger(TestHubName, Connection = TestHubName)] string input, ILogger logger) { logger.LogInformation($"Input({input})"); - _eventWait.Set(); + _eventWait1.Set(); } } @@ -817,7 +826,7 @@ public static void ProcessMultipleEvents([EventHubTrigger(TestHubName, Connectio if (s_processedEventCount == s_eventCount) { _results.AddRange(events); - _eventWait.Set(); + _eventWait1.Set(); } } } @@ -850,7 +859,7 @@ public static void ProcessMultipleEventsBinaryData([EventHubTrigger(TestHubName, // filter for the ID the current test is using if (s_processedEventCount == s_eventCount) { - _eventWait.Set(); + _eventWait1.Set(); } } } @@ -915,7 +924,7 @@ public static void ProcessMultipleEvents([EventHubTrigger(TestHubName, Connectio if (s_processedEventCount >= s_eventCount) { - _eventWait.Set(); + _eventWait1.Set(); } } } @@ -963,7 +972,7 @@ public static void ProcessMultiplePartitionEvents([EventHubTrigger(TestHubName, if (_results.Count > 0) { - _eventWait.Set(); + _eventWait1.Set(); } } } @@ -982,7 +991,7 @@ public static void ProcessSingleEvent([EventHubTrigger(TestHubName, Connection = Assert.AreEqual("value1", properties["TestProp1"]); Assert.AreEqual("value2", properties["TestProp2"]); - _eventWait.Set(); + _eventWait1.Set(); } } public class TestPoco @@ -997,7 +1006,7 @@ public static void ProcessSingleEvent([EventHubTrigger(TestHubName, Connection = string partitionKey, DateTime enqueuedTimeUtc, IDictionary properties, IDictionary systemProperties) { - _eventWait.Set(); + _eventWait1.Set(); } } @@ -1022,7 +1031,7 @@ public static void ProcessMultipleEvents([EventHubTrigger(TestHubName, Connectio { Assert.GreaterOrEqual(DateTimeOffset.Parse(result), earliestAllowedOffset); } - _eventWait.Set(); + _eventWait1.Set(); } } }