From ca5f18dcef4b876ad985b7ec40f9866de17ee4e6 Mon Sep 17 00:00:00 2001 From: Davide Montanari Date: Mon, 14 Sep 2026 09:48:42 -0700 Subject: [PATCH 01/19] Fixed executionId usage in WaitForOrchestrationAsync in ServiceBusOrchestrationService --- .../WaitForOrchestrationTests.cs | 262 ++++++++++++++++++ .../ServiceBusOrchestrationService.cs | 5 +- 2 files changed, 266 insertions(+), 1 deletion(-) create mode 100644 Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs diff --git a/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs b/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs new file mode 100644 index 000000000..6bfa660f8 --- /dev/null +++ b/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs @@ -0,0 +1,262 @@ +// ---------------------------------------------------------------------------------- +// Copyright Microsoft Corporation +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// http://www.apache.org/licenses/LICENSE-2.0 +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. +// ---------------------------------------------------------------------------------- + +namespace DurableTask.ServiceBus.Tests +{ + using System; + using System.Collections.Generic; + using System.Linq; + using System.Threading; + using System.Threading.Tasks; + using DurableTask.Core; + using DurableTask.Core.Tracking; + using DurableTask.ServiceBus.Settings; + using Microsoft.VisualStudio.TestTools.UnitTesting; + + /// + /// Unit tests for . + /// These use an in-memory instance store so they do not require a live Service Bus namespace. + /// + [TestClass] + public class WaitForOrchestrationTests + { + const string InstanceId = "instance-1"; + const string FakeConnectionString = + "Endpoint=sb://test.servicebus.windows.net/;SharedAccessKeyName=key;SharedAccessKey=dGVzdGtleQ=="; + + static readonly DateTime BaseTime = new DateTime(2024, 1, 1, 0, 0, 0, DateTimeKind.Utc); + + static ServiceBusOrchestrationService CreateService(FakeInstanceStore instanceStore) + { + return new ServiceBusOrchestrationService( + FakeConnectionString, + "testhub", + instanceStore, + null, + new ServiceBusOrchestrationServiceSettings()); + } + + static OrchestrationState CreateState( + string executionId, + OrchestrationStatus status, + DateTime createdTime, + string output = null) + { + return new OrchestrationState + { + OrchestrationInstance = new OrchestrationInstance + { + InstanceId = InstanceId, + ExecutionId = executionId + }, + OrchestrationStatus = status, + CreatedTime = createdTime, + LastUpdatedTime = createdTime, + Output = output + }; + } + + /// + /// While waiting on a specific execution, state written by an earlier run + /// of the same instance id must never be returned, even though it is the only readable state. + /// + [TestMethod] + public async Task WaitForOrchestration_PinnedExecution_IgnoresStateFromPreviousRun() + { + var store = new FakeInstanceStore(); + + // A previous run of the same instance id that already completed. + store.States.Add(CreateState("previous-run", OrchestrationStatus.Completed, BaseTime, "stale output")); + + // The new run's state only becomes readable after the first poll. + store.OnQuery = s => + { + if (s.QueryCount == 2) + { + s.States.Add(CreateState("current-run", OrchestrationStatus.Completed, BaseTime.AddMinutes(5), "fresh output")); + } + }; + + ServiceBusOrchestrationService service = CreateService(store); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + "current-run", + TimeSpan.FromSeconds(30), + CancellationToken.None); + + Assert.IsNotNull(state); + Assert.AreEqual("current-run", state.OrchestrationInstance.ExecutionId, "Returned state from the wrong run."); + Assert.AreEqual("fresh output", state.Output); + } + + /// + /// If the new run never becomes readable, the wait must time out rather than return the + /// previous run's result. + /// + [TestMethod] + public async Task WaitForOrchestration_PinnedExecution_TimesOutRatherThanReturnPreviousRun() + { + var store = new FakeInstanceStore(); + store.States.Add(CreateState("previous-run", OrchestrationStatus.Completed, BaseTime, "stale output")); + + ServiceBusOrchestrationService service = CreateService(store); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + "current-run", + TimeSpan.FromSeconds(2), + CancellationToken.None); + + Assert.IsNull(state, "State from a previous run of the same instance id must not be returned."); + } + + /// + /// Callers that do not supply an execution id keep the legacy behavior of following the + /// current generation of the instance. + /// + [DataTestMethod] + [DataRow(null)] + [DataRow("")] + [DataRow(" ")] + public async Task WaitForOrchestration_WithoutExecutionId_UsesCurrentGeneration(string executionId) + { + var store = new FakeInstanceStore(); + store.States.Add(CreateState("generation-1", OrchestrationStatus.ContinuedAsNew, BaseTime, "next input")); + store.States.Add(CreateState("generation-2", OrchestrationStatus.Completed, BaseTime.AddMinutes(1), "final output")); + + ServiceBusOrchestrationService service = CreateService(store); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + executionId, + TimeSpan.FromSeconds(30), + CancellationToken.None); + + Assert.IsNotNull(state); + Assert.AreEqual(OrchestrationStatus.Completed, state.OrchestrationStatus); + Assert.AreEqual("generation-2", state.OrchestrationInstance.ExecutionId); + Assert.AreEqual(0, store.PinnedQueryCount, "An empty execution id must not be queried as an exact execution."); + Assert.AreEqual(1, store.LatestQueryCount, "The current generation lookup should not have been used."); + } + + [TestMethod] + public async Task WaitForOrchestration_PinnedExecution_ReturnsFailedState() + { + var store = new FakeInstanceStore(); + store.States.Add(CreateState("generation-1", OrchestrationStatus.Failed, BaseTime, "boom")); + + ServiceBusOrchestrationService service = CreateService(store); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + "generation-1", + TimeSpan.FromSeconds(30), + CancellationToken.None); + + Assert.IsNotNull(state); + Assert.AreEqual(OrchestrationStatus.Failed, state.OrchestrationStatus); + Assert.AreEqual(1, store.PinnedQueryCount, "The pinned lookup should have been used."); + Assert.AreEqual(0, store.LatestQueryCount, "The current generation lookup should not have been used."); + } + + /// + /// In-memory instance store that mimics the query semantics of AzureTableInstanceStore. + /// + sealed class FakeInstanceStore : IOrchestrationServiceInstanceStore + { + public List States { get; } = new List(); + + /// + /// Mimics AzureTableInstanceStore, which excludes ContinuedAsNew rows from the + /// current generation lookup. + /// + public bool FilterContinuedAsNew { get; set; } = true; + + /// + /// Invoked before every state lookup so a test can make state readable after N polls. + /// + public Action OnQuery { get; set; } + + public int PinnedQueryCount { get; private set; } + + public int LatestQueryCount { get; private set; } + + public int QueryCount => this.PinnedQueryCount + this.LatestQueryCount; + + public int MaxHistoryEntryLength => 1024; + + public Task> GetOrchestrationStateAsync(string instanceId, bool allInstances) + { + this.LatestQueryCount++; + this.OnQuery?.Invoke(this); + + IEnumerable matches = this.States + .Where(s => s.OrchestrationInstance.InstanceId == instanceId); + + if (!allInstances && this.FilterContinuedAsNew) + { + matches = matches.Where(s => s.OrchestrationStatus != OrchestrationStatus.ContinuedAsNew); + } + + if (allInstances) + { + return Task.FromResult(matches.Select(Wrap)); + } + + OrchestrationState latest = matches.OrderByDescending(s => s.LastUpdatedTime).FirstOrDefault(); + + return Task.FromResult(latest == null + ? Enumerable.Empty() + : new[] { Wrap(latest) }.AsEnumerable()); + } + + public Task GetOrchestrationStateAsync(string instanceId, string executionId) + { + this.PinnedQueryCount++; + this.OnQuery?.Invoke(this); + + OrchestrationState match = this.States.FirstOrDefault( + s => s.OrchestrationInstance.InstanceId == instanceId && + s.OrchestrationInstance.ExecutionId == executionId); + + return Task.FromResult(match == null ? null : Wrap(match)); + } + + static OrchestrationStateInstanceEntity Wrap(OrchestrationState state) + { + return new OrchestrationStateInstanceEntity { State = state }; + } + + public Task InitializeStoreAsync(bool recreate) => throw new NotImplementedException(); + + public Task DeleteStoreAsync() => throw new NotImplementedException(); + + public Task WriteEntitiesAsync(IEnumerable entities) => throw new NotImplementedException(); + + public Task> GetEntitiesAsync(string instanceId, string executionId) => throw new NotImplementedException(); + + public Task DeleteEntitiesAsync(IEnumerable entities) => throw new NotImplementedException(); + + public Task> GetOrchestrationHistoryEventsAsync(string instanceId, string executionId) => throw new NotImplementedException(); + + public Task PurgeOrchestrationHistoryEventsAsync(DateTime thresholdDateTimeUtc, OrchestrationStateTimeRangeFilterType timeRangeFilterType) => throw new NotImplementedException(); + + public Task WriteJumpStartEntitiesAsync(IEnumerable entities) => throw new NotImplementedException(); + + public Task DeleteJumpStartEntitiesAsync(IEnumerable entities) => throw new NotImplementedException(); + + public Task> GetJumpStartEntitiesAsync(int top) => throw new NotImplementedException(); + } + } +} diff --git a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs index ae451e033..1fe0e3c7d 100644 --- a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs +++ b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs @@ -1243,7 +1243,10 @@ public async Task WaitForOrchestrationAsync( while (!cancellationToken.IsCancellationRequested && timeoutSeconds > 0) { - OrchestrationState state = (await GetOrchestrationStateAsync(instanceId, false))?.FirstOrDefault(); + OrchestrationState state = !string.IsNullOrWhiteSpace(executionId) + ? await GetOrchestrationStateAsync(instanceId, executionId) + : (await GetOrchestrationStateAsync(instanceId, false))?.FirstOrDefault(); + if (state == null || (state.OrchestrationStatus == OrchestrationStatus.Running) || (state.OrchestrationStatus == OrchestrationStatus.Pending)) From 5a232c6317be3c934564d13cc94ee63c90fadcc8 Mon Sep 17 00:00:00 2001 From: Davide Montanari Date: Mon, 14 Sep 2026 10:34:32 -0700 Subject: [PATCH 02/19] Marked ContinuedAsNew and Suspended as non-terminal states --- .../WaitForOrchestrationTests.cs | 146 +++++++++++++++++- .../ServiceBusOrchestrationService.cs | 64 ++++++-- 2 files changed, 199 insertions(+), 11 deletions(-) diff --git a/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs b/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs index 6bfa660f8..2676b649d 100644 --- a/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs +++ b/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs @@ -67,7 +67,7 @@ static OrchestrationState CreateState( } /// - /// While waiting on a specific execution, state written by an earlier run + /// The core of the fix: while waiting on a specific execution, state written by an earlier run /// of the same instance id must never be returned, even though it is the only readable state. /// [TestMethod] @@ -121,6 +121,117 @@ public async Task WaitForOrchestration_PinnedExecution_TimesOutRatherThanReturnP Assert.IsNull(state, "State from a previous run of the same instance id must not be returned."); } + /// + /// A ContinuedAsNew row is a tombstone that is never updated again, so the wait must follow the + /// new generation instead of polling the pinned execution forever. + /// + [TestMethod] + public async Task WaitForOrchestration_PinnedExecution_ContinuedAsNew_FollowsNextGeneration() + { + var store = new FakeInstanceStore(); + store.States.Add(CreateState("generation-1", OrchestrationStatus.ContinuedAsNew, BaseTime, "next input")); + store.States.Add(CreateState("generation-2", OrchestrationStatus.Completed, BaseTime.AddMinutes(1), "final output")); + + ServiceBusOrchestrationService service = CreateService(store); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + "generation-1", + TimeSpan.FromSeconds(30), + CancellationToken.None); + + Assert.IsNotNull(state); + Assert.AreEqual(OrchestrationStatus.Completed, state.OrchestrationStatus); + Assert.AreEqual("generation-2", state.OrchestrationInstance.ExecutionId); + Assert.AreEqual("final output", state.Output); + } + + /// + /// After following a continue-as-new to the current generation we must still not accept state + /// left behind by an earlier run of the same instance id. + /// + [TestMethod] + public async Task WaitForOrchestration_ContinuedAsNew_IgnoresStateFromPreviousRun() + { + var store = new FakeInstanceStore(); + + // Previous run: completed long before the current run started. + store.States.Add(CreateState("previous-run", OrchestrationStatus.Completed, BaseTime, "stale output")); + + // Current run, first generation, already continued as new. + store.States.Add(CreateState("generation-1", OrchestrationStatus.ContinuedAsNew, BaseTime.AddMinutes(5), "next input")); + + // The next generation only becomes readable later. + store.OnQuery = s => + { + if (s.QueryCount == 3) + { + s.States.Add(CreateState("generation-2", OrchestrationStatus.Completed, BaseTime.AddMinutes(6), "final output")); + } + }; + + ServiceBusOrchestrationService service = CreateService(store); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + "generation-1", + TimeSpan.FromSeconds(30), + CancellationToken.None); + + Assert.IsNotNull(state); + Assert.AreEqual("generation-2", state.OrchestrationInstance.ExecutionId, "Returned state from the wrong run."); + Assert.AreEqual("final output", state.Output); + } + + /// + /// Suspended is a pause, not a terminal state: the orchestration has no result yet. + /// + [TestMethod] + public async Task WaitForOrchestration_Suspended_KeepsWaiting() + { + var store = new FakeInstanceStore(); + store.States.Add(CreateState("generation-1", OrchestrationStatus.Suspended, BaseTime)); + + store.OnQuery = s => + { + if (s.QueryCount == 2) + { + // Resumed and completed. + s.States.Clear(); + s.States.Add(CreateState("generation-1", OrchestrationStatus.Completed, BaseTime, "final output")); + } + }; + + ServiceBusOrchestrationService service = CreateService(store); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + "generation-1", + TimeSpan.FromSeconds(30), + CancellationToken.None); + + Assert.IsNotNull(state); + Assert.AreEqual(OrchestrationStatus.Completed, state.OrchestrationStatus); + Assert.AreEqual("final output", state.Output); + } + + [TestMethod] + public async Task WaitForOrchestration_Suspended_IsNotReturnedAsTerminal() + { + var store = new FakeInstanceStore(); + store.States.Add(CreateState("generation-1", OrchestrationStatus.Suspended, BaseTime)); + + ServiceBusOrchestrationService service = CreateService(store); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + "generation-1", + TimeSpan.FromSeconds(2), + CancellationToken.None); + + Assert.IsNull(state, "A suspended orchestration has not completed and must not be returned."); + } + /// /// Callers that do not supply an execution id keep the legacy behavior of following the /// current generation of the instance. @@ -147,7 +258,38 @@ public async Task WaitForOrchestration_WithoutExecutionId_UsesCurrentGeneration( Assert.AreEqual(OrchestrationStatus.Completed, state.OrchestrationStatus); Assert.AreEqual("generation-2", state.OrchestrationInstance.ExecutionId); Assert.AreEqual(0, store.PinnedQueryCount, "An empty execution id must not be queried as an exact execution."); - Assert.AreEqual(1, store.LatestQueryCount, "The current generation lookup should not have been used."); + } + + /// + /// Hiding ContinuedAsNew rows from the current generation lookup is an implementation detail of + /// AzureTableInstanceStore, not a guarantee of IOrchestrationServiceInstanceStore. A store that + /// surfaces them must not cause a tombstone to be reported as the final state. + /// + [TestMethod] + public async Task WaitForOrchestration_StoreWithoutContinuedAsNewFilter_DoesNotReturnTombstone() + { + var store = new FakeInstanceStore { FilterContinuedAsNew = false }; + store.States.Add(CreateState("generation-1", OrchestrationStatus.ContinuedAsNew, BaseTime, "next input")); + + store.OnQuery = s => + { + if (s.QueryCount == 2) + { + s.States.Add(CreateState("generation-2", OrchestrationStatus.Completed, BaseTime.AddMinutes(1), "final output")); + } + }; + + ServiceBusOrchestrationService service = CreateService(store); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + null, + TimeSpan.FromSeconds(30), + CancellationToken.None); + + Assert.IsNotNull(state); + Assert.AreEqual(OrchestrationStatus.Completed, state.OrchestrationStatus); + Assert.AreEqual("generation-2", state.OrchestrationInstance.ExecutionId); } [TestMethod] diff --git a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs index 1fe0e3c7d..84f84b616 100644 --- a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs +++ b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs @@ -62,7 +62,7 @@ public class ServiceBusOrchestrationService : IOrchestrationService, IOrchestrat // as every fetched message also creates a tracking message which counts towards this limit. const int MaxMessageCount = 80; const int SessionStreamWarningSizeInBytes = 150 * 1024; - const int StatusPollingIntervalInSeconds = 2; + static readonly TimeSpan StatusPollingInterval = TimeSpan.FromSeconds(2); const int DuplicateDetectionWindowInHours = 4; /// @@ -1222,9 +1222,10 @@ public async Task ForceTerminateTaskOrchestrationAsync(string instanceId, string /// /// Wait for an orchestration to reach any terminal state within the given timeout /// - /// The execution id of the orchestration + /// The execution id of the orchestration. When specified, only that execution is + /// tracked; if it has , the current generation of the instance is followed instead. /// Instance to wait for - /// Max timeout to wait + /// Max timeout to wait. Only positive values, , or are allowed. /// Task cancellation token public async Task WaitForOrchestrationAsync( string instanceId, @@ -1239,20 +1240,65 @@ public async Task WaitForOrchestrationAsync( throw new ArgumentException("instanceId"); } - double timeoutSeconds = timeout.TotalSeconds; + bool isInfiniteTimeSpan = timeout == Timeout.InfiniteTimeSpan; + if (timeout < TimeSpan.Zero && !isInfiniteTimeSpan) + { + throw new ArgumentException($"The parameter {nameof(timeout)} cannot be negative." + + $" The value for {nameof(timeout)} was '{timeout}'." + + $" Please provide either a positive timeout value or Timeout.InfiniteTimeSpan."); + } + + bool pinnedToExecution = !string.IsNullOrWhiteSpace(executionId); - while (!cancellationToken.IsCancellationRequested && timeoutSeconds > 0) + // Once we stop tracking a specific execution we must not accept state left behind by an + // earlier run of the same instance id, otherwise we reintroduce the race that querying by + // execution id is meant to avoid. + DateTime minimumCreatedTime = DateTimeUtils.MinDateTime; + + while (!cancellationToken.IsCancellationRequested) { - OrchestrationState state = !string.IsNullOrWhiteSpace(executionId) + OrchestrationState state = pinnedToExecution ? await GetOrchestrationStateAsync(instanceId, executionId) : (await GetOrchestrationStateAsync(instanceId, false))?.FirstOrDefault(); + // A pinned execution that continued-as-new is only a tombstone: the live orchestration + // moved on to a new execution id, so stop pinning and follow the current generation. + if (pinnedToExecution && state?.OrchestrationStatus == OrchestrationStatus.ContinuedAsNew) + { + pinnedToExecution = false; + minimumCreatedTime = state.CreatedTime; + continue; + } + + if (state?.CreatedTime < minimumCreatedTime) + { + // State from a previous run of this instance id; the current generation is not readable yet. + state = null; + } + + // ContinuedAsNew is never a final state: a new generation always follows it. The built-in + // AzureTableInstanceStore hides these rows from the non-pinned lookup, but that is not + // guaranteed by IOrchestrationServiceInstanceStore, so keep polling if one surfaces. if (state == null || (state.OrchestrationStatus == OrchestrationStatus.Running) - || (state.OrchestrationStatus == OrchestrationStatus.Pending)) + || (state.OrchestrationStatus == OrchestrationStatus.Pending) + || (state.OrchestrationStatus == OrchestrationStatus.Suspended) + || (state.OrchestrationStatus == OrchestrationStatus.ContinuedAsNew)) { - await Task.Delay(StatusPollingIntervalInSeconds * 1000, cancellationToken); - timeoutSeconds -= StatusPollingIntervalInSeconds; + if (!isInfiniteTimeSpan) + { + timeout -= StatusPollingInterval; + + // For a user-provided timeout of `TimeSpan.Zero`, + // we want to check the status of the orchestration once and then return. + // Therefore, we check the timeout condition after the status check. + if (timeout <= TimeSpan.Zero) + { + break; + } + } + + await Task.Delay(StatusPollingInterval, cancellationToken); } else { From 556b30b9a26293751f3f79ab34f46b03a9c693d3 Mon Sep 17 00:00:00 2001 From: Davide Montanari Date: Mon, 14 Sep 2026 11:09:13 -0700 Subject: [PATCH 03/19] Added new tests --- .../WaitForOrchestrationTests.cs | 227 +++++++++++++++--- 1 file changed, 197 insertions(+), 30 deletions(-) diff --git a/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs b/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs index 2676b649d..ec778c63d 100644 --- a/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs +++ b/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs @@ -184,37 +184,9 @@ public async Task WaitForOrchestration_ContinuedAsNew_IgnoresStateFromPreviousRu } /// - /// Suspended is a pause, not a terminal state: the orchestration has no result yet. + /// Suspended is a pause, not a terminal state: an orchestration that is never resumed has no + /// result, so the wait must time out rather than report it as finished. /// - [TestMethod] - public async Task WaitForOrchestration_Suspended_KeepsWaiting() - { - var store = new FakeInstanceStore(); - store.States.Add(CreateState("generation-1", OrchestrationStatus.Suspended, BaseTime)); - - store.OnQuery = s => - { - if (s.QueryCount == 2) - { - // Resumed and completed. - s.States.Clear(); - s.States.Add(CreateState("generation-1", OrchestrationStatus.Completed, BaseTime, "final output")); - } - }; - - ServiceBusOrchestrationService service = CreateService(store); - - OrchestrationState state = await service.WaitForOrchestrationAsync( - InstanceId, - "generation-1", - TimeSpan.FromSeconds(30), - CancellationToken.None); - - Assert.IsNotNull(state); - Assert.AreEqual(OrchestrationStatus.Completed, state.OrchestrationStatus); - Assert.AreEqual("final output", state.Output); - } - [TestMethod] public async Task WaitForOrchestration_Suspended_IsNotReturnedAsTerminal() { @@ -312,6 +284,201 @@ public async Task WaitForOrchestration_PinnedExecution_ReturnsFailedState() Assert.AreEqual(0, store.LatestQueryCount, "The current generation lookup should not have been used."); } + /// + /// A suspended orchestration that is resumed and then runs to completion must return the + /// final state, not stop at the intermediate Suspended or Running rows. + /// + [TestMethod] + public async Task WaitForOrchestration_Suspended_ResumedAndCompleted_ReturnsFinalState() + { + var store = new FakeInstanceStore(); + store.States.Add(CreateState("generation-1", OrchestrationStatus.Suspended, BaseTime)); + + // Suspended -> Running (resumed) -> Completed, one transition per poll. + store.OnQuery = s => + { + OrchestrationStatus? next = s.QueryCount == 2 ? OrchestrationStatus.Running + : s.QueryCount == 3 ? OrchestrationStatus.Completed + : (OrchestrationStatus?)null; + + if (next != null) + { + s.States.Clear(); + s.States.Add(CreateState( + "generation-1", + next.Value, + BaseTime, + next == OrchestrationStatus.Completed ? "final output" : null)); + } + }; + + ServiceBusOrchestrationService service = CreateService(store); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + "generation-1", + TimeSpan.FromSeconds(30), + CancellationToken.None); + + Assert.IsNotNull(state, "The resumed orchestration completed and must be returned."); + Assert.AreEqual(OrchestrationStatus.Completed, state.OrchestrationStatus); + Assert.AreEqual("final output", state.Output); + Assert.AreEqual(3, store.QueryCount, "The wait should have polled through Suspended and Running."); + } + + /// + /// Negative timeouts are a caller bug. Timeout.InfiniteTimeSpan is itself negative (-1ms), so it + /// must be excluded from this check; values adjacent to it must still be rejected. + /// + [DataTestMethod] + [DataRow(-2)] + [DataRow(-1000)] + [DataRow(-60000)] + public async Task WaitForOrchestration_Timeout_Negative_Throws(int timeoutMilliseconds) + { + var store = new FakeInstanceStore(); + store.States.Add(CreateState("generation-1", OrchestrationStatus.Completed, BaseTime, "final output")); + + ServiceBusOrchestrationService service = CreateService(store); + + await Assert.ThrowsExceptionAsync( + () => service.WaitForOrchestrationAsync( + InstanceId, + "generation-1", + TimeSpan.FromMilliseconds(timeoutMilliseconds), + CancellationToken.None), + $"A timeout of {timeoutMilliseconds}ms should be rejected."); + + Assert.AreEqual(0, store.QueryCount, "The timeout should be validated before any lookup."); + } + + /// + /// A zero timeout means "check once and return", so a state that is already terminal is returned. + /// + [TestMethod] + public async Task WaitForOrchestration_Timeout_Zero_ChecksOnceAndReturnsTerminalState() + { + var store = new FakeInstanceStore(); + store.States.Add(CreateState("generation-1", OrchestrationStatus.Completed, BaseTime, "final output")); + + ServiceBusOrchestrationService service = CreateService(store); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + "generation-1", + TimeSpan.Zero, + CancellationToken.None); + + Assert.IsNotNull(state, "A zero timeout must still perform one status check."); + Assert.AreEqual(OrchestrationStatus.Completed, state.OrchestrationStatus); + Assert.AreEqual(1, store.QueryCount); + } + + /// + /// A zero timeout must not wait: if the orchestration is not yet terminal it returns null after + /// a single check rather than polling. + /// + [TestMethod] + public async Task WaitForOrchestration_Timeout_Zero_DoesNotPollWhenNotComplete() + { + var store = new FakeInstanceStore(); + store.States.Add(CreateState("generation-1", OrchestrationStatus.Running, BaseTime)); + + ServiceBusOrchestrationService service = CreateService(store); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + "generation-1", + TimeSpan.Zero, + CancellationToken.None); + + Assert.IsNull(state); + Assert.AreEqual(1, store.QueryCount, "A zero timeout must check exactly once and not poll."); + } + + /// + /// A positive timeout polls until it elapses, then gives up and returns null. + /// + [TestMethod] + public async Task WaitForOrchestration_Timeout_Positive_PollsThenReturnsNull() + { + var store = new FakeInstanceStore(); + store.States.Add(CreateState("generation-1", OrchestrationStatus.Running, BaseTime)); + + ServiceBusOrchestrationService service = CreateService(store); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + "generation-1", + TimeSpan.FromSeconds(4), + CancellationToken.None); + + Assert.IsNull(state, "The orchestration never completed, so the wait must time out."); + Assert.IsTrue(store.QueryCount > 1, $"A 4 second timeout should poll more than once, polled {store.QueryCount} time(s)."); + } + + /// + /// Timeout.InfiniteTimeSpan is negative, so a naive remaining-time check would treat it as already + /// elapsed and return null on the first poll instead of waiting. + /// + [TestMethod] + public async Task WaitForOrchestration_Timeout_Infinite_WaitsForCompletion() + { + var store = new FakeInstanceStore(); + store.States.Add(CreateState("generation-1", OrchestrationStatus.Running, BaseTime)); + + store.OnQuery = s => + { + if (s.QueryCount == 2) + { + s.States.Clear(); + s.States.Add(CreateState("generation-1", OrchestrationStatus.Completed, BaseTime, "final output")); + } + }; + + ServiceBusOrchestrationService service = CreateService(store); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + "generation-1", + Timeout.InfiniteTimeSpan, + CancellationToken.None); + + Assert.IsNotNull(state, "An infinite timeout must keep waiting instead of giving up immediately."); + Assert.AreEqual(OrchestrationStatus.Completed, state.OrchestrationStatus); + Assert.AreEqual("final output", state.Output); + } + + /// + /// An infinite wait must still observe cancellation, otherwise it can never be stopped. + /// + [TestMethod] + public async Task WaitForOrchestration_Timeout_Infinite_HonorsCancellation() + { + var store = new FakeInstanceStore(); + store.States.Add(CreateState("generation-1", OrchestrationStatus.Running, BaseTime)); + + using (var cts = new CancellationTokenSource(TimeSpan.FromSeconds(3))) + { + ServiceBusOrchestrationService service = CreateService(store); + + try + { + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + "generation-1", + Timeout.InfiniteTimeSpan, + cts.Token); + + Assert.IsNull(state, "A cancelled wait must not return a state."); + } + catch (OperationCanceledException) + { + // Also acceptable: the polling delay observes the token directly. + } + } + } + /// /// In-memory instance store that mimics the query semantics of AzureTableInstanceStore. /// From b417710f0c02d2af6eb15d02cfae5215409c2fa2 Mon Sep 17 00:00:00 2001 From: Davide Montanari Date: Mon, 14 Sep 2026 11:16:09 -0700 Subject: [PATCH 04/19] Fixed doc --- .../ServiceBusOrchestrationService.cs | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs index 84f84b616..0d5922132 100644 --- a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs +++ b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs @@ -1222,8 +1222,10 @@ public async Task ForceTerminateTaskOrchestrationAsync(string instanceId, string /// /// Wait for an orchestration to reach any terminal state within the given timeout /// - /// The execution id of the orchestration. When specified, only that execution is - /// tracked; if it has , the current generation of the instance is followed instead. + /// The execution id of the orchestration to wait for. When null, empty, or + /// whitespace, the current generation of the instance is followed. Otherwise only that execution is + /// tracked, except for a execution, + /// where the current generation of the instance is followed instead. /// Instance to wait for /// Max timeout to wait. Only positive values, , or are allowed. /// Task cancellation token From 176514364d970947e05d23a8b2ea261ba8c69857 Mon Sep 17 00:00:00 2001 From: Davide Montanari Date: Mon, 14 Sep 2026 11:35:31 -0700 Subject: [PATCH 05/19] Fixed check when execution is ContinuedAsNew and timeout is TimeSpan.Zero --- .../WaitForOrchestrationTests.cs | 24 +++++++++++++++++++ .../ServiceBusOrchestrationService.cs | 3 ++- 2 files changed, 26 insertions(+), 1 deletion(-) diff --git a/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs b/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs index ec778c63d..0ad3d62f8 100644 --- a/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs +++ b/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs @@ -146,6 +146,30 @@ public async Task WaitForOrchestration_PinnedExecution_ContinuedAsNew_FollowsNex Assert.AreEqual("final output", state.Output); } + /// + /// Un-pinning after a ContinuedAsNew tombstone must not re-query within the same iteration: + /// the tombstone was itself this iteration's status check, so a zero timeout still performs + /// exactly one lookup rather than reaching the next generation for free. + /// + [TestMethod] + public async Task WaitForOrchestration_Timeout_Zero_ContinuedAsNew_ChecksOnceAndDoesNotFollowNextGeneration() + { + var store = new FakeInstanceStore(); + store.States.Add(CreateState("generation-1", OrchestrationStatus.ContinuedAsNew, BaseTime, "next input")); + store.States.Add(CreateState("generation-2", OrchestrationStatus.Completed, BaseTime.AddMinutes(1), "final output")); + + ServiceBusOrchestrationService service = CreateService(store); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + "generation-1", + TimeSpan.Zero, + CancellationToken.None); + + Assert.IsNull(state, "A zero timeout must not reach the next generation after un-pinning."); + Assert.AreEqual(1, store.QueryCount, "A zero timeout must perform exactly one lookup."); + } + /// /// After following a continue-as-new to the current generation we must still not accept state /// left behind by an earlier run of the same instance id. diff --git a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs index 0d5922132..55da2b225 100644 --- a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs +++ b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs @@ -1265,11 +1265,12 @@ public async Task WaitForOrchestrationAsync( // A pinned execution that continued-as-new is only a tombstone: the live orchestration // moved on to a new execution id, so stop pinning and follow the current generation. + // The tombstone still counts as this iteration's status check, so fall through to the + // timeout accounting below rather than re-querying immediately. if (pinnedToExecution && state?.OrchestrationStatus == OrchestrationStatus.ContinuedAsNew) { pinnedToExecution = false; minimumCreatedTime = state.CreatedTime; - continue; } if (state?.CreatedTime < minimumCreatedTime) From e0a2e2c0532d250284520dddf1427c93726e4b5f Mon Sep 17 00:00:00 2001 From: Davide Montanari Date: Mon, 14 Sep 2026 12:02:28 -0700 Subject: [PATCH 06/19] Fixed timeout check --- .../WaitForOrchestrationTests.cs | 66 +++++++++++++++++++ .../ServiceBusOrchestrationService.cs | 22 +++++-- 2 files changed, 82 insertions(+), 6 deletions(-) diff --git a/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs b/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs index 0ad3d62f8..a7f61f758 100644 --- a/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs +++ b/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs @@ -441,6 +441,72 @@ public async Task WaitForOrchestration_Timeout_Positive_PollsThenReturnsNull() Assert.IsTrue(store.QueryCount > 1, $"A 4 second timeout should poll more than once, polled {store.QueryCount} time(s)."); } + /// + /// The whole timeout window must be polled. Charging a full polling interval against the budget + /// before the delay is awaited would abandon the wait about halfway through and miss an + /// orchestration that completes late in the window. + /// + [TestMethod] + public async Task WaitForOrchestration_Timeout_Positive_PollsForTheFullTimeout() + { + var store = new FakeInstanceStore(); + store.States.Add(CreateState("generation-1", OrchestrationStatus.Running, BaseTime)); + + // With a 2 second polling interval a 4 second timeout allows checks at roughly t=0, t=2 + // and t=4, so state that only becomes terminal on the third check is still observed. + store.OnQuery = s => + { + if (s.QueryCount == 3) + { + s.States.Clear(); + s.States.Add(CreateState("generation-1", OrchestrationStatus.Completed, BaseTime, "final output")); + } + }; + + ServiceBusOrchestrationService service = CreateService(store); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + "generation-1", + TimeSpan.FromSeconds(4), + CancellationToken.None); + + Assert.IsNotNull(state, "The orchestration completed within the timeout and must be returned."); + Assert.AreEqual(OrchestrationStatus.Completed, state.OrchestrationStatus); + Assert.AreEqual("final output", state.Output); + } + + /// + /// A timeout shorter than the polling interval must wait out its window rather than returning on + /// the first check, but must not overshoot it by delaying for a whole polling interval. + /// + [TestMethod] + public async Task WaitForOrchestration_Timeout_ShorterThanPollingInterval_WaitsWithoutOvershooting() + { + var store = new FakeInstanceStore(); + store.States.Add(CreateState("generation-1", OrchestrationStatus.Running, BaseTime)); + + ServiceBusOrchestrationService service = CreateService(store); + + var stopwatch = System.Diagnostics.Stopwatch.StartNew(); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + "generation-1", + TimeSpan.FromSeconds(1), + CancellationToken.None); + + stopwatch.Stop(); + + Assert.IsNull(state, "The orchestration never completed, so the wait must time out."); + Assert.IsTrue( + stopwatch.Elapsed >= TimeSpan.FromMilliseconds(900), + $"The wait must use its window instead of returning immediately, took {stopwatch.ElapsedMilliseconds}ms."); + Assert.IsTrue( + stopwatch.Elapsed < TimeSpan.FromMilliseconds(1800), + $"The wait must not delay for a whole polling interval past its timeout, took {stopwatch.ElapsedMilliseconds}ms."); + } + /// /// Timeout.InfiniteTimeSpan is negative, so a naive remaining-time check would treat it as already /// elapsed and return null on the first poll instead of waiting. diff --git a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs index 55da2b225..267bda5bd 100644 --- a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs +++ b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs @@ -1288,20 +1288,30 @@ public async Task WaitForOrchestrationAsync( || (state.OrchestrationStatus == OrchestrationStatus.Suspended) || (state.OrchestrationStatus == OrchestrationStatus.ContinuedAsNew)) { + TimeSpan delay = StatusPollingInterval; + if (!isInfiniteTimeSpan) { - timeout -= StatusPollingInterval; - - // For a user-provided timeout of `TimeSpan.Zero`, - // we want to check the status of the orchestration once and then return. - // Therefore, we check the timeout condition after the status check. + // The timeout condition is checked after the status check so that a user-provided + // timeout of `TimeSpan.Zero` still checks the status of the orchestration once + // before returning. if (timeout <= TimeSpan.Zero) { break; } + + // Only the time actually spent waiting is charged against the budget, and the + // last delay is clamped to what remains, so the full timeout window is polled + // before giving up. + if (timeout < delay) + { + delay = timeout; + } + + timeout -= delay; } - await Task.Delay(StatusPollingInterval, cancellationToken); + await Task.Delay(delay, cancellationToken); } else { From 6f416b4e5581090ef61b53203012c9516dcaab53 Mon Sep 17 00:00:00 2001 From: Davide Montanari Date: Mon, 14 Sep 2026 13:26:02 -0700 Subject: [PATCH 07/19] - Fixed timeout validation error message to include TimeSpan.Zero - Fixed assets WaitForOrchestration_Timeout_Infinite_HonorsCancellation test for robustness --- .../WaitForOrchestrationTests.cs | 27 ++++++++++++++----- .../ServiceBusOrchestrationService.cs | 7 ++--- 2 files changed, 24 insertions(+), 10 deletions(-) diff --git a/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs b/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs index a7f61f758..a632e5f01 100644 --- a/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs +++ b/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs @@ -552,20 +552,33 @@ public async Task WaitForOrchestration_Timeout_Infinite_HonorsCancellation() { ServiceBusOrchestrationService service = CreateService(store); + Task waitTask = service.WaitForOrchestrationAsync( + InstanceId, + "generation-1", + Timeout.InfiniteTimeSpan, + cts.Token); + + // An infinite wait that ignores its token would never complete, so bound the await to + // fail the test instead of hanging the run. + Task finished = await Task.WhenAny(waitTask, Task.Delay(TimeSpan.FromSeconds(30))); + + Assert.AreSame(waitTask, finished, "An infinite wait must stop once its token is cancelled."); + try { - OrchestrationState state = await service.WaitForOrchestrationAsync( - InstanceId, - "generation-1", - Timeout.InfiniteTimeSpan, - cts.Token); - - Assert.IsNull(state, "A cancelled wait must not return a state."); + Assert.IsNull(await waitTask, "A cancelled wait must not return a state."); } catch (OperationCanceledException) { // Also acceptable: the polling delay observes the token directly. } + + // Whichever path ends the wait, it must have kept polling until cancellation rather + // than treating the negative Timeout.InfiniteTimeSpan as an elapsed budget and + // bailing out after the first lookup. + Assert.IsTrue( + store.QueryCount > 1, + $"An infinite wait must poll until cancelled, polled {store.QueryCount} time(s)."); } } diff --git a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs index 267bda5bd..fa25fe4cc 100644 --- a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs +++ b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs @@ -1239,15 +1239,16 @@ public async Task WaitForOrchestrationAsync( if (string.IsNullOrWhiteSpace(instanceId)) { - throw new ArgumentException("instanceId"); + throw new ArgumentException("The instance id cannot be null, empty, or whitespace.", nameof(instanceId)); } bool isInfiniteTimeSpan = timeout == Timeout.InfiniteTimeSpan; if (timeout < TimeSpan.Zero && !isInfiniteTimeSpan) { - throw new ArgumentException($"The parameter {nameof(timeout)} cannot be negative." + + throw new ArgumentException($"The parameter {nameof(timeout)} cannot be negative unless it is Timeout.InfiniteTimeSpan." + $" The value for {nameof(timeout)} was '{timeout}'." + - $" Please provide either a positive timeout value or Timeout.InfiniteTimeSpan."); + $" Please provide a positive timeout value, TimeSpan.Zero or Timeout.InfiniteTimeSpan.", + nameof(timeout)); } bool pinnedToExecution = !string.IsNullOrWhiteSpace(executionId); From 46259872b6583a164ce056bdbd5e8364bfc8809d Mon Sep 17 00:00:00 2001 From: Davide Montanari Date: Mon, 14 Sep 2026 14:35:59 -0700 Subject: [PATCH 08/19] Renamed test file path --- .../DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs | 0 1 file changed, 0 insertions(+), 0 deletions(-) rename {Test => test}/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs (100%) diff --git a/Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs b/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs similarity index 100% rename from Test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs rename to test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs From 0578772055fe052a299cac4fe7e1fb6c93721843 Mon Sep 17 00:00:00 2001 From: Davide Montanari Date: Mon, 14 Sep 2026 14:48:07 -0700 Subject: [PATCH 09/19] Fixed test --- .../WaitForOrchestrationTests.cs | 20 +++++++++++++++---- 1 file changed, 16 insertions(+), 4 deletions(-) diff --git a/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs b/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs index a632e5f01..09e0b7c6c 100644 --- a/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs +++ b/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs @@ -488,23 +488,35 @@ public async Task WaitForOrchestration_Timeout_ShorterThanPollingInterval_WaitsW ServiceBusOrchestrationService service = CreateService(store); + // The timeout is kept far below the 2 second polling interval so the three possible + // behaviours are widely separated in time and the assertions below do not depend on + // precise scheduling: clamping the delay takes about 200ms, delaying for a whole + // interval takes about 2s, and charging the budget before the delay returns instantly. + TimeSpan timeout = TimeSpan.FromMilliseconds(200); + var stopwatch = System.Diagnostics.Stopwatch.StartNew(); OrchestrationState state = await service.WaitForOrchestrationAsync( InstanceId, "generation-1", - TimeSpan.FromSeconds(1), + timeout, CancellationToken.None); stopwatch.Stop(); Assert.IsNull(state, "The orchestration never completed, so the wait must time out."); + Assert.AreEqual(2, store.QueryCount, "The wait must check once, wait out its window, then check again."); + + // Task.Delay never returns early, so this only fails if the wait did not delay at all. Assert.IsTrue( - stopwatch.Elapsed >= TimeSpan.FromMilliseconds(900), + stopwatch.Elapsed >= TimeSpan.FromMilliseconds(150), $"The wait must use its window instead of returning immediately, took {stopwatch.ElapsedMilliseconds}ms."); + + // Well clear of the ~200ms a clamped delay needs, and well below the 2s a full polling + // interval would take. Assert.IsTrue( - stopwatch.Elapsed < TimeSpan.FromMilliseconds(1800), - $"The wait must not delay for a whole polling interval past its timeout, took {stopwatch.ElapsedMilliseconds}ms."); + stopwatch.Elapsed < TimeSpan.FromMilliseconds(1500), + $"The wait must clamp its delay to the remaining timeout instead of waiting a whole polling interval, took {stopwatch.ElapsedMilliseconds}ms."); } /// From c698660bdbd2bcfc5899076ac25a91717a05c478 Mon Sep 17 00:00:00 2001 From: Davide Montanari Date: Mon, 14 Sep 2026 16:11:15 -0700 Subject: [PATCH 10/19] Added check on LastUpdatedTime to ignore previous runs --- .../ServiceBusOrchestrationService.cs | 13 ++++- .../WaitForOrchestrationTests.cs | 54 ++++++++++++++++++- 2 files changed, 63 insertions(+), 4 deletions(-) diff --git a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs index fa25fe4cc..4bea7d9eb 100644 --- a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs +++ b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs @@ -1257,6 +1257,7 @@ public async Task WaitForOrchestrationAsync( // earlier run of the same instance id, otherwise we reintroduce the race that querying by // execution id is meant to avoid. DateTime minimumCreatedTime = DateTimeUtils.MinDateTime; + DateTime minimumLastUpdatedTime = DateTimeUtils.MinDateTime; while (!cancellationToken.IsCancellationRequested) { @@ -1272,11 +1273,19 @@ public async Task WaitForOrchestrationAsync( { pinnedToExecution = false; minimumCreatedTime = state.CreatedTime; + minimumLastUpdatedTime = state.LastUpdatedTime; } - if (state?.CreatedTime < minimumCreatedTime) + // State from a previous run of this instance id; the current generation is not readable yet. + // CreatedTime comes from a HistoryEvent timestamp and is not guaranteed unique, so it cannot + // separate the two on its own. When it ties, fall back to LastUpdatedTime: a previous run + // necessarily stopped being updated no later than the continue-as-new that wrote the + // tombstone, while the generation that follows the tombstone is updated at or after it. + if (state != null + && (state.CreatedTime < minimumCreatedTime + || (state.CreatedTime == minimumCreatedTime && state.LastUpdatedTime < minimumLastUpdatedTime))) + { - // State from a previous run of this instance id; the current generation is not readable yet. state = null; } diff --git a/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs b/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs index 09e0b7c6c..e751da077 100644 --- a/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs +++ b/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs @@ -50,7 +50,8 @@ static OrchestrationState CreateState( string executionId, OrchestrationStatus status, DateTime createdTime, - string output = null) + string output = null, + DateTime? lastUpdatedTime = null) { return new OrchestrationState { @@ -61,7 +62,7 @@ static OrchestrationState CreateState( }, OrchestrationStatus = status, CreatedTime = createdTime, - LastUpdatedTime = createdTime, + LastUpdatedTime = lastUpdatedTime ?? createdTime, Output = output }; } @@ -207,6 +208,55 @@ public async Task WaitForOrchestration_ContinuedAsNew_IgnoresStateFromPreviousRu Assert.AreEqual("final output", state.Output); } + /// + /// CreatedTime comes from a HistoryEvent timestamp and is not guaranteed unique, so a previous + /// run can happen to share the ContinuedAsNew tombstone's CreatedTime. LastUpdatedTime breaks + /// that tie: the previous run stopped being updated before the continue-as-new occurred. + /// + [TestMethod] + public async Task WaitForOrchestration_ContinuedAsNew_IgnoresPreviousRunSharingTheTombstoneCreatedTime() + { + var store = new FakeInstanceStore(); + + // Previous run: completed before the current run continued as new, but its ExecutionStarted + // timestamp collides with the current run's. + store.States.Add(CreateState( + "previous-run", + OrchestrationStatus.Completed, + BaseTime, + "stale output", + BaseTime.AddMinutes(1))); + + // Current run, first generation, already continued as new. + store.States.Add(CreateState( + "generation-1", + OrchestrationStatus.ContinuedAsNew, + BaseTime, + "next input", + BaseTime.AddMinutes(5))); + + // The next generation only becomes readable later. + store.OnQuery = s => + { + if (s.QueryCount == 3) + { + s.States.Add(CreateState("generation-2", OrchestrationStatus.Completed, BaseTime.AddMinutes(5), "final output")); + } + }; + + ServiceBusOrchestrationService service = CreateService(store); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + "generation-1", + TimeSpan.FromSeconds(30), + CancellationToken.None); + + Assert.IsNotNull(state); + Assert.AreEqual("generation-2", state.OrchestrationInstance.ExecutionId, "Returned state from the wrong run."); + Assert.AreEqual("final output", state.Output); + } + /// /// Suspended is a pause, not a terminal state: an orchestration that is never resumed has no /// result, so the wait must time out rather than report it as finished. From d8a1397e6c45017fd1a5b790bd65e65f3a0fc180 Mon Sep 17 00:00:00 2001 From: Davide Montanari Date: Mon, 14 Sep 2026 17:21:34 -0700 Subject: [PATCH 11/19] Fixed stale-result race when no executionId is passed --- .../ServiceBusOrchestrationService.cs | 51 +++++++++++++++-- .../WaitForOrchestrationTests.cs | 57 +++++++++++++++++++ 2 files changed, 102 insertions(+), 6 deletions(-) diff --git a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs index 4bea7d9eb..80df0d74f 100644 --- a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs +++ b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs @@ -1259,6 +1259,10 @@ public async Task WaitForOrchestrationAsync( DateTime minimumCreatedTime = DateTimeUtils.MinDateTime; DateTime minimumLastUpdatedTime = DateTimeUtils.MinDateTime; + // Once a generation floor is established, state left behind by an earlier run is rejected + // by the floor check alone, so the full execution history no longer has to be consulted. + bool hasGenerationFloor = false; + while (!cancellationToken.IsCancellationRequested) { OrchestrationState state = pinnedToExecution @@ -1274,6 +1278,7 @@ public async Task WaitForOrchestrationAsync( pinnedToExecution = false; minimumCreatedTime = state.CreatedTime; minimumLastUpdatedTime = state.LastUpdatedTime; + hasGenerationFloor = true; } // State from a previous run of this instance id; the current generation is not readable yet. @@ -1284,19 +1289,40 @@ public async Task WaitForOrchestrationAsync( if (state != null && (state.CreatedTime < minimumCreatedTime || (state.CreatedTime == minimumCreatedTime && state.LastUpdatedTime < minimumLastUpdatedTime))) - { state = null; } + // The current-generation lookup only sees rows the store chooses to expose. + // AzureTableInstanceStore hides ContinuedAsNew tombstones from it, so at a + // continue-as-new boundary -- the tombstone hidden, the next generation not readable + // yet -- the newest surviving row can be a terminal state left by an earlier run of + // this instance id. Before accepting a terminal state, confirm against the full + // execution history that no newer row exists. Only needed until a generation floor + // exists: past that point nothing older than the floor is accepted anyway. + if (state != null && !pinnedToExecution && !hasGenerationFloor && !IsNonTerminalStatus(state.OrchestrationStatus)) + { + // allExecutions carries no ordering guarantee, so pick the newest row here. + OrchestrationState newest = (await GetOrchestrationStateAsync(instanceId, true)) + ?.OrderByDescending(s => s.LastUpdatedTime) + .FirstOrDefault(); + + if (newest != null && newest.LastUpdatedTime > state.LastUpdatedTime) + { + // A newer generation exists but is not visible to the current generation + // lookup yet. Raise the floor so the stale row stops being a candidate, and + // keep polling until the newer generation becomes readable. + minimumCreatedTime = newest.CreatedTime; + minimumLastUpdatedTime = newest.LastUpdatedTime; + hasGenerationFloor = true; + state = null; + } + } + // ContinuedAsNew is never a final state: a new generation always follows it. The built-in // AzureTableInstanceStore hides these rows from the non-pinned lookup, but that is not // guaranteed by IOrchestrationServiceInstanceStore, so keep polling if one surfaces. - if (state == null - || (state.OrchestrationStatus == OrchestrationStatus.Running) - || (state.OrchestrationStatus == OrchestrationStatus.Pending) - || (state.OrchestrationStatus == OrchestrationStatus.Suspended) - || (state.OrchestrationStatus == OrchestrationStatus.ContinuedAsNew)) + if (state == null || IsNonTerminalStatus(state.OrchestrationStatus)) { TimeSpan delay = StatusPollingInterval; @@ -1332,6 +1358,19 @@ public async Task WaitForOrchestrationAsync( return null; } + /// + /// Determines whether an orchestration status can still change. Suspended is a pause rather + /// than an outcome, and ContinuedAsNew is only a handoff to the next generation, so neither + /// is a result a caller can wait on. + /// + static bool IsNonTerminalStatus(OrchestrationStatus status) + { + return status == OrchestrationStatus.Running + || status == OrchestrationStatus.Pending + || status == OrchestrationStatus.Suspended + || status == OrchestrationStatus.ContinuedAsNew; + } + /// /// Get a list of orchestration states from the instance storage for the most current execution (generation) of the specified instance. /// diff --git a/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs b/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs index e751da077..e80f1e2bb 100644 --- a/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs +++ b/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs @@ -306,6 +306,63 @@ public async Task WaitForOrchestration_WithoutExecutionId_UsesCurrentGeneration( Assert.AreEqual(0, store.PinnedQueryCount, "An empty execution id must not be queried as an exact execution."); } + /// + /// Without an execution id the wait depends on the store's current generation lookup, which + /// AzureTableInstanceStore builds by hiding ContinuedAsNew rows. At a continue-as-new boundary + /// the current run's tombstone is hidden and the next generation is not readable yet, so the + /// newest surviving row is a completed state left by an earlier run of the same instance id. + /// The wait must consult the full execution history before accepting that as the result. + /// + [TestMethod] + public async Task WaitForOrchestration_WithoutExecutionId_IgnoresPreviousRunWhileCurrentGenerationIsTombstoned() + { + var store = new FakeInstanceStore(); + + // Previous run of this instance id, completed long ago. Added first so that insertion + // order does not match recency: the all-executions lookup carries no ordering guarantee, + // so the wait has to sort by itself to find the tombstone. + store.States.Add(CreateState( + "previous-run", + OrchestrationStatus.Completed, + BaseTime, + "stale output", + BaseTime.AddMinutes(1))); + + // Current run, mid continue-as-new: the tombstone is hidden from the current generation + // lookup, so only the previous run's completed row is visible. + store.States.Add(CreateState( + "generation-1", + OrchestrationStatus.ContinuedAsNew, + BaseTime.AddMinutes(5), + "next input")); + + // The next generation only becomes readable later. + store.OnQuery = s => + { + if (s.QueryCount == 3) + { + s.States.Add(CreateState("generation-2", OrchestrationStatus.Completed, BaseTime.AddMinutes(6), "final output")); + } + }; + + ServiceBusOrchestrationService service = CreateService(store); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + null, + TimeSpan.FromSeconds(30), + CancellationToken.None); + + Assert.IsNotNull(state); + Assert.AreEqual("generation-2", state.OrchestrationInstance.ExecutionId, "Returned state from the wrong run."); + Assert.AreEqual("final output", state.Output); + Assert.AreEqual(0, store.PinnedQueryCount, "A null execution id must not be queried as an exact execution."); + Assert.AreEqual( + 3, + store.QueryCount, + "The full history should be consulted once, to establish the generation floor, and not again."); + } + /// /// Hiding ContinuedAsNew rows from the current generation lookup is an implementation detail of /// AzureTableInstanceStore, not a guarantee of IOrchestrationServiceInstanceStore. A store that From e16cb28f6997d9d7318d17bd7bf2e47da5faad04 Mon Sep 17 00:00:00 2001 From: Davide Montanari Date: Mon, 14 Sep 2026 17:51:31 -0700 Subject: [PATCH 12/19] Reverted escalation code --- .../ServiceBusOrchestrationService.cs | 51 ++--------------- .../WaitForOrchestrationTests.cs | 57 ------------------- 2 files changed, 6 insertions(+), 102 deletions(-) diff --git a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs index 80df0d74f..4bea7d9eb 100644 --- a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs +++ b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs @@ -1259,10 +1259,6 @@ public async Task WaitForOrchestrationAsync( DateTime minimumCreatedTime = DateTimeUtils.MinDateTime; DateTime minimumLastUpdatedTime = DateTimeUtils.MinDateTime; - // Once a generation floor is established, state left behind by an earlier run is rejected - // by the floor check alone, so the full execution history no longer has to be consulted. - bool hasGenerationFloor = false; - while (!cancellationToken.IsCancellationRequested) { OrchestrationState state = pinnedToExecution @@ -1278,7 +1274,6 @@ public async Task WaitForOrchestrationAsync( pinnedToExecution = false; minimumCreatedTime = state.CreatedTime; minimumLastUpdatedTime = state.LastUpdatedTime; - hasGenerationFloor = true; } // State from a previous run of this instance id; the current generation is not readable yet. @@ -1289,40 +1284,19 @@ public async Task WaitForOrchestrationAsync( if (state != null && (state.CreatedTime < minimumCreatedTime || (state.CreatedTime == minimumCreatedTime && state.LastUpdatedTime < minimumLastUpdatedTime))) - { - state = null; - } - // The current-generation lookup only sees rows the store chooses to expose. - // AzureTableInstanceStore hides ContinuedAsNew tombstones from it, so at a - // continue-as-new boundary -- the tombstone hidden, the next generation not readable - // yet -- the newest surviving row can be a terminal state left by an earlier run of - // this instance id. Before accepting a terminal state, confirm against the full - // execution history that no newer row exists. Only needed until a generation floor - // exists: past that point nothing older than the floor is accepted anyway. - if (state != null && !pinnedToExecution && !hasGenerationFloor && !IsNonTerminalStatus(state.OrchestrationStatus)) { - // allExecutions carries no ordering guarantee, so pick the newest row here. - OrchestrationState newest = (await GetOrchestrationStateAsync(instanceId, true)) - ?.OrderByDescending(s => s.LastUpdatedTime) - .FirstOrDefault(); - - if (newest != null && newest.LastUpdatedTime > state.LastUpdatedTime) - { - // A newer generation exists but is not visible to the current generation - // lookup yet. Raise the floor so the stale row stops being a candidate, and - // keep polling until the newer generation becomes readable. - minimumCreatedTime = newest.CreatedTime; - minimumLastUpdatedTime = newest.LastUpdatedTime; - hasGenerationFloor = true; - state = null; - } + state = null; } // ContinuedAsNew is never a final state: a new generation always follows it. The built-in // AzureTableInstanceStore hides these rows from the non-pinned lookup, but that is not // guaranteed by IOrchestrationServiceInstanceStore, so keep polling if one surfaces. - if (state == null || IsNonTerminalStatus(state.OrchestrationStatus)) + if (state == null + || (state.OrchestrationStatus == OrchestrationStatus.Running) + || (state.OrchestrationStatus == OrchestrationStatus.Pending) + || (state.OrchestrationStatus == OrchestrationStatus.Suspended) + || (state.OrchestrationStatus == OrchestrationStatus.ContinuedAsNew)) { TimeSpan delay = StatusPollingInterval; @@ -1358,19 +1332,6 @@ public async Task WaitForOrchestrationAsync( return null; } - /// - /// Determines whether an orchestration status can still change. Suspended is a pause rather - /// than an outcome, and ContinuedAsNew is only a handoff to the next generation, so neither - /// is a result a caller can wait on. - /// - static bool IsNonTerminalStatus(OrchestrationStatus status) - { - return status == OrchestrationStatus.Running - || status == OrchestrationStatus.Pending - || status == OrchestrationStatus.Suspended - || status == OrchestrationStatus.ContinuedAsNew; - } - /// /// Get a list of orchestration states from the instance storage for the most current execution (generation) of the specified instance. /// diff --git a/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs b/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs index e80f1e2bb..e751da077 100644 --- a/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs +++ b/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs @@ -306,63 +306,6 @@ public async Task WaitForOrchestration_WithoutExecutionId_UsesCurrentGeneration( Assert.AreEqual(0, store.PinnedQueryCount, "An empty execution id must not be queried as an exact execution."); } - /// - /// Without an execution id the wait depends on the store's current generation lookup, which - /// AzureTableInstanceStore builds by hiding ContinuedAsNew rows. At a continue-as-new boundary - /// the current run's tombstone is hidden and the next generation is not readable yet, so the - /// newest surviving row is a completed state left by an earlier run of the same instance id. - /// The wait must consult the full execution history before accepting that as the result. - /// - [TestMethod] - public async Task WaitForOrchestration_WithoutExecutionId_IgnoresPreviousRunWhileCurrentGenerationIsTombstoned() - { - var store = new FakeInstanceStore(); - - // Previous run of this instance id, completed long ago. Added first so that insertion - // order does not match recency: the all-executions lookup carries no ordering guarantee, - // so the wait has to sort by itself to find the tombstone. - store.States.Add(CreateState( - "previous-run", - OrchestrationStatus.Completed, - BaseTime, - "stale output", - BaseTime.AddMinutes(1))); - - // Current run, mid continue-as-new: the tombstone is hidden from the current generation - // lookup, so only the previous run's completed row is visible. - store.States.Add(CreateState( - "generation-1", - OrchestrationStatus.ContinuedAsNew, - BaseTime.AddMinutes(5), - "next input")); - - // The next generation only becomes readable later. - store.OnQuery = s => - { - if (s.QueryCount == 3) - { - s.States.Add(CreateState("generation-2", OrchestrationStatus.Completed, BaseTime.AddMinutes(6), "final output")); - } - }; - - ServiceBusOrchestrationService service = CreateService(store); - - OrchestrationState state = await service.WaitForOrchestrationAsync( - InstanceId, - null, - TimeSpan.FromSeconds(30), - CancellationToken.None); - - Assert.IsNotNull(state); - Assert.AreEqual("generation-2", state.OrchestrationInstance.ExecutionId, "Returned state from the wrong run."); - Assert.AreEqual("final output", state.Output); - Assert.AreEqual(0, store.PinnedQueryCount, "A null execution id must not be queried as an exact execution."); - Assert.AreEqual( - 3, - store.QueryCount, - "The full history should be consulted once, to establish the generation floor, and not again."); - } - /// /// Hiding ContinuedAsNew rows from the current generation lookup is an implementation detail of /// AzureTableInstanceStore, not a guarantee of IOrchestrationServiceInstanceStore. A store that From ec1079abeb5014900ad2484ab4cbe6b63813ec9f Mon Sep 17 00:00:00 2001 From: Davide Montanari Date: Mon, 14 Sep 2026 17:58:47 -0700 Subject: [PATCH 13/19] Removed checks --- .../ServiceBusOrchestrationService.cs | 21 ----- .../WaitForOrchestrationTests.cs | 91 +------------------ 2 files changed, 2 insertions(+), 110 deletions(-) diff --git a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs index 4bea7d9eb..f759ff591 100644 --- a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs +++ b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs @@ -1253,12 +1253,6 @@ public async Task WaitForOrchestrationAsync( bool pinnedToExecution = !string.IsNullOrWhiteSpace(executionId); - // Once we stop tracking a specific execution we must not accept state left behind by an - // earlier run of the same instance id, otherwise we reintroduce the race that querying by - // execution id is meant to avoid. - DateTime minimumCreatedTime = DateTimeUtils.MinDateTime; - DateTime minimumLastUpdatedTime = DateTimeUtils.MinDateTime; - while (!cancellationToken.IsCancellationRequested) { OrchestrationState state = pinnedToExecution @@ -1272,21 +1266,6 @@ public async Task WaitForOrchestrationAsync( if (pinnedToExecution && state?.OrchestrationStatus == OrchestrationStatus.ContinuedAsNew) { pinnedToExecution = false; - minimumCreatedTime = state.CreatedTime; - minimumLastUpdatedTime = state.LastUpdatedTime; - } - - // State from a previous run of this instance id; the current generation is not readable yet. - // CreatedTime comes from a HistoryEvent timestamp and is not guaranteed unique, so it cannot - // separate the two on its own. When it ties, fall back to LastUpdatedTime: a previous run - // necessarily stopped being updated no later than the continue-as-new that wrote the - // tombstone, while the generation that follows the tombstone is updated at or after it. - if (state != null - && (state.CreatedTime < minimumCreatedTime - || (state.CreatedTime == minimumCreatedTime && state.LastUpdatedTime < minimumLastUpdatedTime))) - - { - state = null; } // ContinuedAsNew is never a final state: a new generation always follows it. The built-in diff --git a/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs b/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs index e751da077..faee45921 100644 --- a/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs +++ b/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs @@ -50,8 +50,7 @@ static OrchestrationState CreateState( string executionId, OrchestrationStatus status, DateTime createdTime, - string output = null, - DateTime? lastUpdatedTime = null) + string output = null) { return new OrchestrationState { @@ -62,7 +61,7 @@ static OrchestrationState CreateState( }, OrchestrationStatus = status, CreatedTime = createdTime, - LastUpdatedTime = lastUpdatedTime ?? createdTime, + LastUpdatedTime = createdTime, Output = output }; } @@ -171,92 +170,6 @@ public async Task WaitForOrchestration_Timeout_Zero_ContinuedAsNew_ChecksOnceAnd Assert.AreEqual(1, store.QueryCount, "A zero timeout must perform exactly one lookup."); } - /// - /// After following a continue-as-new to the current generation we must still not accept state - /// left behind by an earlier run of the same instance id. - /// - [TestMethod] - public async Task WaitForOrchestration_ContinuedAsNew_IgnoresStateFromPreviousRun() - { - var store = new FakeInstanceStore(); - - // Previous run: completed long before the current run started. - store.States.Add(CreateState("previous-run", OrchestrationStatus.Completed, BaseTime, "stale output")); - - // Current run, first generation, already continued as new. - store.States.Add(CreateState("generation-1", OrchestrationStatus.ContinuedAsNew, BaseTime.AddMinutes(5), "next input")); - - // The next generation only becomes readable later. - store.OnQuery = s => - { - if (s.QueryCount == 3) - { - s.States.Add(CreateState("generation-2", OrchestrationStatus.Completed, BaseTime.AddMinutes(6), "final output")); - } - }; - - ServiceBusOrchestrationService service = CreateService(store); - - OrchestrationState state = await service.WaitForOrchestrationAsync( - InstanceId, - "generation-1", - TimeSpan.FromSeconds(30), - CancellationToken.None); - - Assert.IsNotNull(state); - Assert.AreEqual("generation-2", state.OrchestrationInstance.ExecutionId, "Returned state from the wrong run."); - Assert.AreEqual("final output", state.Output); - } - - /// - /// CreatedTime comes from a HistoryEvent timestamp and is not guaranteed unique, so a previous - /// run can happen to share the ContinuedAsNew tombstone's CreatedTime. LastUpdatedTime breaks - /// that tie: the previous run stopped being updated before the continue-as-new occurred. - /// - [TestMethod] - public async Task WaitForOrchestration_ContinuedAsNew_IgnoresPreviousRunSharingTheTombstoneCreatedTime() - { - var store = new FakeInstanceStore(); - - // Previous run: completed before the current run continued as new, but its ExecutionStarted - // timestamp collides with the current run's. - store.States.Add(CreateState( - "previous-run", - OrchestrationStatus.Completed, - BaseTime, - "stale output", - BaseTime.AddMinutes(1))); - - // Current run, first generation, already continued as new. - store.States.Add(CreateState( - "generation-1", - OrchestrationStatus.ContinuedAsNew, - BaseTime, - "next input", - BaseTime.AddMinutes(5))); - - // The next generation only becomes readable later. - store.OnQuery = s => - { - if (s.QueryCount == 3) - { - s.States.Add(CreateState("generation-2", OrchestrationStatus.Completed, BaseTime.AddMinutes(5), "final output")); - } - }; - - ServiceBusOrchestrationService service = CreateService(store); - - OrchestrationState state = await service.WaitForOrchestrationAsync( - InstanceId, - "generation-1", - TimeSpan.FromSeconds(30), - CancellationToken.None); - - Assert.IsNotNull(state); - Assert.AreEqual("generation-2", state.OrchestrationInstance.ExecutionId, "Returned state from the wrong run."); - Assert.AreEqual("final output", state.Output); - } - /// /// Suspended is a pause, not a terminal state: an orchestration that is never resumed has no /// result, so the wait must time out rather than report it as finished. From beaa47c1abcb587f269f02662f3d14b68669e4b3 Mon Sep 17 00:00:00 2001 From: Davide Montanari Date: Mon, 14 Sep 2026 18:36:24 -0700 Subject: [PATCH 14/19] Added check on state from previous run --- .../ServiceBusOrchestrationService.cs | 20 ++++ .../WaitForOrchestrationTests.cs | 91 ++++++++++++++++++- 2 files changed, 109 insertions(+), 2 deletions(-) diff --git a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs index f759ff591..aa064bd5e 100644 --- a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs +++ b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs @@ -1253,6 +1253,11 @@ public async Task WaitForOrchestrationAsync( bool pinnedToExecution = !string.IsNullOrWhiteSpace(executionId); + // Once we stop tracking a specific execution we must not accept state left behind by an + // earlier run of the same instance id. + DateTime minimumCreatedTime = DateTimeUtils.MinDateTime; + DateTime minimumLastUpdatedTime = DateTimeUtils.MinDateTime; + while (!cancellationToken.IsCancellationRequested) { OrchestrationState state = pinnedToExecution @@ -1266,6 +1271,21 @@ public async Task WaitForOrchestrationAsync( if (pinnedToExecution && state?.OrchestrationStatus == OrchestrationStatus.ContinuedAsNew) { pinnedToExecution = false; + minimumCreatedTime = state.CreatedTime; + minimumLastUpdatedTime = state.LastUpdatedTime; + } + + // State from a previous run of this instance id; the current generation is not readable yet. + // CreatedTime comes from a HistoryEvent timestamp and is not guaranteed unique, so it cannot + // separate the two on its own. When it ties, fall back to LastUpdatedTime: a previous run + // necessarily stopped being updated no later than the continue-as-new that wrote the + // tombstone, while the generation that follows the tombstone is updated at or after it. + if (state != null + && (state.CreatedTime < minimumCreatedTime + || (state.CreatedTime == minimumCreatedTime && state.LastUpdatedTime < minimumLastUpdatedTime))) + + { + state = null; } // ContinuedAsNew is never a final state: a new generation always follows it. The built-in diff --git a/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs b/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs index faee45921..e751da077 100644 --- a/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs +++ b/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs @@ -50,7 +50,8 @@ static OrchestrationState CreateState( string executionId, OrchestrationStatus status, DateTime createdTime, - string output = null) + string output = null, + DateTime? lastUpdatedTime = null) { return new OrchestrationState { @@ -61,7 +62,7 @@ static OrchestrationState CreateState( }, OrchestrationStatus = status, CreatedTime = createdTime, - LastUpdatedTime = createdTime, + LastUpdatedTime = lastUpdatedTime ?? createdTime, Output = output }; } @@ -170,6 +171,92 @@ public async Task WaitForOrchestration_Timeout_Zero_ContinuedAsNew_ChecksOnceAnd Assert.AreEqual(1, store.QueryCount, "A zero timeout must perform exactly one lookup."); } + /// + /// After following a continue-as-new to the current generation we must still not accept state + /// left behind by an earlier run of the same instance id. + /// + [TestMethod] + public async Task WaitForOrchestration_ContinuedAsNew_IgnoresStateFromPreviousRun() + { + var store = new FakeInstanceStore(); + + // Previous run: completed long before the current run started. + store.States.Add(CreateState("previous-run", OrchestrationStatus.Completed, BaseTime, "stale output")); + + // Current run, first generation, already continued as new. + store.States.Add(CreateState("generation-1", OrchestrationStatus.ContinuedAsNew, BaseTime.AddMinutes(5), "next input")); + + // The next generation only becomes readable later. + store.OnQuery = s => + { + if (s.QueryCount == 3) + { + s.States.Add(CreateState("generation-2", OrchestrationStatus.Completed, BaseTime.AddMinutes(6), "final output")); + } + }; + + ServiceBusOrchestrationService service = CreateService(store); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + "generation-1", + TimeSpan.FromSeconds(30), + CancellationToken.None); + + Assert.IsNotNull(state); + Assert.AreEqual("generation-2", state.OrchestrationInstance.ExecutionId, "Returned state from the wrong run."); + Assert.AreEqual("final output", state.Output); + } + + /// + /// CreatedTime comes from a HistoryEvent timestamp and is not guaranteed unique, so a previous + /// run can happen to share the ContinuedAsNew tombstone's CreatedTime. LastUpdatedTime breaks + /// that tie: the previous run stopped being updated before the continue-as-new occurred. + /// + [TestMethod] + public async Task WaitForOrchestration_ContinuedAsNew_IgnoresPreviousRunSharingTheTombstoneCreatedTime() + { + var store = new FakeInstanceStore(); + + // Previous run: completed before the current run continued as new, but its ExecutionStarted + // timestamp collides with the current run's. + store.States.Add(CreateState( + "previous-run", + OrchestrationStatus.Completed, + BaseTime, + "stale output", + BaseTime.AddMinutes(1))); + + // Current run, first generation, already continued as new. + store.States.Add(CreateState( + "generation-1", + OrchestrationStatus.ContinuedAsNew, + BaseTime, + "next input", + BaseTime.AddMinutes(5))); + + // The next generation only becomes readable later. + store.OnQuery = s => + { + if (s.QueryCount == 3) + { + s.States.Add(CreateState("generation-2", OrchestrationStatus.Completed, BaseTime.AddMinutes(5), "final output")); + } + }; + + ServiceBusOrchestrationService service = CreateService(store); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + "generation-1", + TimeSpan.FromSeconds(30), + CancellationToken.None); + + Assert.IsNotNull(state); + Assert.AreEqual("generation-2", state.OrchestrationInstance.ExecutionId, "Returned state from the wrong run."); + Assert.AreEqual("final output", state.Output); + } + /// /// Suspended is a pause, not a terminal state: an orchestration that is never resumed has no /// result, so the wait must time out rather than report it as finished. From c7667c3188326f0bce40c3644599a0fcbc3f03d4 Mon Sep 17 00:00:00 2001 From: Davide Montanari Date: Tue, 22 Sep 2026 15:56:14 -0700 Subject: [PATCH 15/19] Preserved the existing store retry policy for execution-specific polling --- .../Tracking/AzureTableInstanceStore.cs | 19 +++++++++++++++---- 1 file changed, 15 insertions(+), 4 deletions(-) diff --git a/src/DurableTask.ServiceBus/Tracking/AzureTableInstanceStore.cs b/src/DurableTask.ServiceBus/Tracking/AzureTableInstanceStore.cs index 818ae3cfe..c6a6d5ee9 100644 --- a/src/DurableTask.ServiceBus/Tracking/AzureTableInstanceStore.cs +++ b/src/DurableTask.ServiceBus/Tracking/AzureTableInstanceStore.cs @@ -217,17 +217,28 @@ public async Task> GetOrchestratio /// The matching orchestration state or null if not found public async Task GetOrchestrationStateAsync(string instanceId, string executionId) { + // Retries here mirror the instance-id overload above so that both lookups survive the same + // transient table failures. Without them a single failure that escapes the storage SDK's own + // retries would fault every caller that polls by execution id. AzureTableOrchestrationStateEntity result = - (await this.tableClient.QueryOrchestrationStatesAsync( - new OrchestrationStateQuery().AddInstanceFilter(instanceId, executionId)).ConfigureAwait(false)).FirstOrDefault(); + (await Utils.ExecuteWithRetries(() => this.tableClient.QueryOrchestrationStatesAsync( + new OrchestrationStateQuery().AddInstanceFilter(instanceId, executionId)), + string.Empty, + "GetOrchestrationStateAsync-stateEntity", + MaxRetriesTableStore, + IntervalBetweenRetriesSecs).ConfigureAwait(false)).FirstOrDefault(); // ReSharper disable once ConvertIfStatementToNullCoalescingExpression if (result == null) { // Query from JumpStart table - result = (await this.tableClient.QueryJumpStartOrchestrationsAsync( + result = (await Utils.ExecuteWithRetries(() => this.tableClient.QueryJumpStartOrchestrationsAsync( new OrchestrationStateQuery() - .AddInstanceFilter(instanceId, executionId)) + .AddInstanceFilter(instanceId, executionId)), + string.Empty, + "GetOrchestrationStateAsync-jumpStartEntity", + MaxRetriesTableStore, + IntervalBetweenRetriesSecs) .ConfigureAwait(false)) .FirstOrDefault(); } From 62a96d23b208deedff4b2b7a926066adbb21df36 Mon Sep 17 00:00:00 2001 From: Davide Montanari Date: Tue, 22 Sep 2026 18:19:00 -0700 Subject: [PATCH 16/19] Added tests for transient failures when passing executionId --- .../AzureTableInstanceStoreRetryTests.cs | 151 ++++++++++++++++++ .../Tracking/AzureTableClient.cs | 4 +- .../Tracking/AzureTableInstanceStore.cs | 14 ++ 3 files changed, 167 insertions(+), 2 deletions(-) create mode 100644 Test/DurableTask.ServiceBus.Tests/AzureTableInstanceStoreRetryTests.cs diff --git a/Test/DurableTask.ServiceBus.Tests/AzureTableInstanceStoreRetryTests.cs b/Test/DurableTask.ServiceBus.Tests/AzureTableInstanceStoreRetryTests.cs new file mode 100644 index 000000000..78d6ef7ba --- /dev/null +++ b/Test/DurableTask.ServiceBus.Tests/AzureTableInstanceStoreRetryTests.cs @@ -0,0 +1,151 @@ +// ---------------------------------------------------------------------------------- +// Copyright Microsoft Corporation +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// http://www.apache.org/licenses/LICENSE-2.0 +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. +// ---------------------------------------------------------------------------------- + +namespace DurableTask.ServiceBus.Tests +{ + using System; + using System.Collections.Generic; + using System.Linq; + using System.Threading.Tasks; + using global::Azure; + using DurableTask.Core; + using DurableTask.Core.Tracking; + using DurableTask.ServiceBus.Tracking; + using Microsoft.VisualStudio.TestTools.UnitTesting; + + /// + /// Verifies that the execution-id lookup on retries transient + /// table failures, matching the instance-id overload. Polling by execution id would otherwise fault + /// on the first transient error that escapes the storage SDK's own retries. + /// + [TestClass] + public class AzureTableInstanceStoreRetryTests + { + const string InstanceId = "instance-1"; + const string ExecutionId = "generation-1"; + const string HubName = "testhub"; + + // Parsed only; constructing a TableServiceClient performs no I/O. + const string FakeStorageConnectionString = "UseDevelopmentStorage=true"; + + static readonly DateTime BaseTime = new DateTime(2024, 1, 1, 0, 0, 0, DateTimeKind.Utc); + + static AzureTableOrchestrationStateEntity CreateEntity() + { + return new AzureTableOrchestrationStateEntity(new OrchestrationState + { + OrchestrationInstance = new OrchestrationInstance + { + InstanceId = InstanceId, + ExecutionId = ExecutionId + }, + OrchestrationStatus = OrchestrationStatus.Completed, + CreatedTime = BaseTime, + LastUpdatedTime = BaseTime, + Output = "final output" + }); + } + + /// + /// A transient failure on the state table must be retried rather than surfaced to the caller. + /// + [TestMethod] + public async Task GetOrchestrationStateAsync_ByExecutionId_RetriesTransientStateTableFailure() + { + var tableClient = new FailOnceTableClient { FailuresBeforeSuccess = 1 }; + var store = new AzureTableInstanceStore(tableClient); + + OrchestrationStateInstanceEntity state = await store.GetOrchestrationStateAsync(InstanceId, ExecutionId); + + Assert.IsNotNull(state, "A single transient failure must not fault the execution-id lookup."); + Assert.AreEqual(ExecutionId, state.State.OrchestrationInstance.ExecutionId); + Assert.AreEqual("final output", state.State.Output); + Assert.AreEqual(2, tableClient.StateQueryCount, "Expected one failed attempt followed by a successful retry."); + } + + /// + /// The JumpStart fallback is a separate query and needs the same protection: without it a + /// transient failure there faults the lookup even though the state table responded fine. + /// + [TestMethod] + public async Task GetOrchestrationStateAsync_ByExecutionId_RetriesTransientJumpStartTableFailure() + { + var tableClient = new FailOnceTableClient + { + // Force the JumpStart fallback by leaving the state table empty. + StateTableIsEmpty = true, + JumpStartFailuresBeforeSuccess = 1 + }; + + var store = new AzureTableInstanceStore(tableClient); + + OrchestrationStateInstanceEntity state = await store.GetOrchestrationStateAsync(InstanceId, ExecutionId); + + Assert.IsNotNull(state, "A single transient failure must not fault the JumpStart fallback."); + Assert.AreEqual(ExecutionId, state.State.OrchestrationInstance.ExecutionId); + Assert.AreEqual(2, tableClient.JumpStartQueryCount, "Expected one failed attempt followed by a successful retry."); + } + + /// + /// Table client that throws a retryable storage error a fixed number of times before returning + /// a row, so a lookup that lacks retries fails and one that retries succeeds. + /// + sealed class FailOnceTableClient : AzureTableClient + { + public FailOnceTableClient() + : base(HubName, FakeStorageConnectionString) + { + } + + public int FailuresBeforeSuccess { get; set; } + + public int JumpStartFailuresBeforeSuccess { get; set; } + + public bool StateTableIsEmpty { get; set; } + + public int StateQueryCount { get; private set; } + + public int JumpStartQueryCount { get; private set; } + + public override Task> QueryOrchestrationStatesAsync( + OrchestrationStateQuery stateQuery) + { + this.StateQueryCount++; + + if (this.StateQueryCount <= this.FailuresBeforeSuccess) + { + throw new RequestFailedException(503, "The server is busy."); + } + + IEnumerable result = this.StateTableIsEmpty + ? Enumerable.Empty() + : new[] { CreateEntity() }; + + return Task.FromResult(result); + } + + public override Task> QueryJumpStartOrchestrationsAsync( + OrchestrationStateQuery stateQuery) + { + this.JumpStartQueryCount++; + + if (this.JumpStartQueryCount <= this.JumpStartFailuresBeforeSuccess) + { + throw new RequestFailedException(503, "The server is busy."); + } + + return Task.FromResult>(new[] { CreateEntity() }); + } + } + } +} diff --git a/src/DurableTask.ServiceBus/Tracking/AzureTableClient.cs b/src/DurableTask.ServiceBus/Tracking/AzureTableClient.cs index c0e28a8df..e4baf99bf 100644 --- a/src/DurableTask.ServiceBus/Tracking/AzureTableClient.cs +++ b/src/DurableTask.ServiceBus/Tracking/AzureTableClient.cs @@ -124,7 +124,7 @@ internal async Task DeleteJumpStartTableIfExistsAsync() } } - public Task> QueryOrchestrationStatesAsync( + public virtual Task> QueryOrchestrationStatesAsync( OrchestrationStateQuery stateQuery) { var query = CreateQueryInternal(stateQuery, false); @@ -138,7 +138,7 @@ public async Task> QueryOrchestrationSt return await QueryTableSegmentAsync(this.historyTableClient, query, continuationToken, count); } - public Task> QueryJumpStartOrchestrationsAsync(OrchestrationStateQuery stateQuery) + public virtual Task> QueryJumpStartOrchestrationsAsync(OrchestrationStateQuery stateQuery) { // TODO: Enable segmented query for paging purpose var query = CreateQueryInternal(stateQuery, true); diff --git a/src/DurableTask.ServiceBus/Tracking/AzureTableInstanceStore.cs b/src/DurableTask.ServiceBus/Tracking/AzureTableInstanceStore.cs index c6a6d5ee9..d8c98a092 100644 --- a/src/DurableTask.ServiceBus/Tracking/AzureTableInstanceStore.cs +++ b/src/DurableTask.ServiceBus/Tracking/AzureTableInstanceStore.cs @@ -90,6 +90,20 @@ public AzureTableInstanceStore(string hubName, Uri endpoint, TokenCredential cre DateTimeUtils.SetMinDateTimeForStorageEmulator(); } + /// + /// Creates a new AzureTableInstanceStore over the supplied table client. Used by tests to + /// substitute the table client; production code uses the connection string or credential + /// constructors above. + /// + /// The table client to issue queries through + internal AzureTableInstanceStore(AzureTableClient tableClient) + { + this.tableClient = tableClient ?? throw new ArgumentNullException(nameof(tableClient)); + + // Workaround an issue with Storage that throws exceptions for any date < 1600 so DateTime.Min cannot be used + DateTimeUtils.SetMinDateTimeForStorageEmulator(); + } + /// /// Runs initialization to prepare the storage for use /// From 98ce63b8e6b4b24901cfd8cded353cb754279f20 Mon Sep 17 00:00:00 2001 From: Davide Montanari Date: Tue, 22 Sep 2026 18:23:43 -0700 Subject: [PATCH 17/19] Renamed test file --- .../AzureTableInstanceStoreRetryTests.cs | 0 1 file changed, 0 insertions(+), 0 deletions(-) rename {Test => test}/DurableTask.ServiceBus.Tests/AzureTableInstanceStoreRetryTests.cs (100%) diff --git a/Test/DurableTask.ServiceBus.Tests/AzureTableInstanceStoreRetryTests.cs b/test/DurableTask.ServiceBus.Tests/AzureTableInstanceStoreRetryTests.cs similarity index 100% rename from Test/DurableTask.ServiceBus.Tests/AzureTableInstanceStoreRetryTests.cs rename to test/DurableTask.ServiceBus.Tests/AzureTableInstanceStoreRetryTests.cs From 1c8f6bfd1da14f2e5bea5051044747d6293778b0 Mon Sep 17 00:00:00 2001 From: Davide Montanari Date: Tue, 22 Sep 2026 20:19:24 -0700 Subject: [PATCH 18/19] Removed check on CreatedTime and short circuit the check on LastUpdatedTime --- .../ServiceBusOrchestrationService.cs | 31 ++++-- .../WaitForOrchestrationTests.cs | 105 +++++++++++++++++- 2 files changed, 124 insertions(+), 12 deletions(-) diff --git a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs index aa064bd5e..d860d577c 100644 --- a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs +++ b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs @@ -65,6 +65,15 @@ public class ServiceBusOrchestrationService : IOrchestrationService, IOrchestrat static readonly TimeSpan StatusPollingInterval = TimeSpan.FromSeconds(2); const int DuplicateDetectionWindowInHours = 4; + /// + /// How many below-floor states to discard after a continue-as-new handoff before accepting one. + /// The floor exists only to cover the brief window where the successor row is not readable yet, + /// which resolves in well under this budget. Past it, a state below the floor is far more likely + /// a successor persisted by a worker whose clock is behind than state left by an earlier run, + /// so accepting it beats waiting forever. + /// + const int MaxPreviousRunRejections = 5; + /// /// Orchestration service settings /// @@ -1255,8 +1264,8 @@ public async Task WaitForOrchestrationAsync( // Once we stop tracking a specific execution we must not accept state left behind by an // earlier run of the same instance id. - DateTime minimumCreatedTime = DateTimeUtils.MinDateTime; DateTime minimumLastUpdatedTime = DateTimeUtils.MinDateTime; + int remainingRejections = 0; while (!cancellationToken.IsCancellationRequested) { @@ -1271,20 +1280,22 @@ public async Task WaitForOrchestrationAsync( if (pinnedToExecution && state?.OrchestrationStatus == OrchestrationStatus.ContinuedAsNew) { pinnedToExecution = false; - minimumCreatedTime = state.CreatedTime; minimumLastUpdatedTime = state.LastUpdatedTime; + remainingRejections = MaxPreviousRunRejections; } // State from a previous run of this instance id; the current generation is not readable yet. - // CreatedTime comes from a HistoryEvent timestamp and is not guaranteed unique, so it cannot - // separate the two on its own. When it ties, fall back to LastUpdatedTime: a previous run - // necessarily stopped being updated no later than the continue-as-new that wrote the - // tombstone, while the generation that follows the tombstone is updated at or after it. - if (state != null - && (state.CreatedTime < minimumCreatedTime - || (state.CreatedTime == minimumCreatedTime && state.LastUpdatedTime < minimumLastUpdatedTime))) - + // A previous run stopped being updated no later than the continue-as-new that wrote the + // tombstone, while its successor is persisted at or after it. CreatedTime is deliberately + // not used: it is stamped on the client for the first generation but on the worker for every + // continue-as-new generation, so a client running ahead would discard a valid successor + // forever. LastUpdatedTime is rewritten by the persisting worker on every state write, so it + // stays within one tier, but worker clocks still disagree: a successor that completes in a + // single episode inside that skew is never rewritten and would stay below the floor forever. + // The budget bounds that, so an infinite wait cannot hang on it. + if (state != null && remainingRejections > 0 && state.LastUpdatedTime < minimumLastUpdatedTime) { + remainingRejections--; state = null; } diff --git a/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs b/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs index e751da077..0723e3b57 100644 --- a/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs +++ b/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs @@ -210,8 +210,9 @@ public async Task WaitForOrchestration_ContinuedAsNew_IgnoresStateFromPreviousRu /// /// CreatedTime comes from a HistoryEvent timestamp and is not guaranteed unique, so a previous - /// run can happen to share the ContinuedAsNew tombstone's CreatedTime. LastUpdatedTime breaks - /// that tie: the previous run stopped being updated before the continue-as-new occurred. + /// run can happen to share the ContinuedAsNew tombstone's CreatedTime. The floor must therefore + /// not depend on CreatedTime at all: LastUpdatedTime alone separates the two, because the + /// previous run stopped being updated before the continue-as-new occurred. /// [TestMethod] public async Task WaitForOrchestration_ContinuedAsNew_IgnoresPreviousRunSharingTheTombstoneCreatedTime() @@ -257,6 +258,106 @@ public async Task WaitForOrchestration_ContinuedAsNew_IgnoresPreviousRunSharingT Assert.AreEqual("final output", state.Output); } + /// + /// The first generation's ExecutionStartedEvent is stamped on the client that scheduled the + /// orchestration, while a continue-as-new generation's is stamped on the worker that ran it, and + /// both come from DateTime.UtcNow on those machines. A client running ahead of the worker therefore + /// produces a tombstone whose CreatedTime is later than its own successor's. Rejecting state by + /// CreatedTime would discard that successor on every iteration, so the wait would never observe the + /// completion it was asked about. + /// + [TestMethod] + public async Task WaitForOrchestration_ContinuedAsNew_FollowsSuccessorWhenClientClockIsAheadOfWorker() + { + var store = new FakeInstanceStore(); + + // First generation: started by a client running one second ahead of the worker, so its + // CreatedTime is later than everything the worker subsequently stamps. Its LastUpdatedTime is + // written by the worker when the continue-as-new is persisted. + store.States.Add(CreateState( + "generation-1", + OrchestrationStatus.ContinuedAsNew, + BaseTime.AddSeconds(1), + "next input", + BaseTime.AddMilliseconds(100))); + + // Successor: created and persisted entirely on the worker, and already complete. + store.States.Add(CreateState( + "generation-2", + OrchestrationStatus.Completed, + BaseTime.AddMilliseconds(50), + "final output", + BaseTime.AddMilliseconds(200))); + + ServiceBusOrchestrationService service = CreateService(store); + + OrchestrationState state = await service.WaitForOrchestrationAsync( + InstanceId, + "generation-1", + TimeSpan.FromSeconds(30), + CancellationToken.None); + + Assert.IsNotNull(state, "Clock skew between the client and the worker discarded a valid successor."); + Assert.AreEqual("generation-2", state.OrchestrationInstance.ExecutionId); + Assert.AreEqual("final output", state.Output); + } + + /// + /// Worker clocks can disagree too, so a short successor can be created, complete and never be + /// written again while its LastUpdatedTime still sits below the tombstone's. An unbounded floor + /// would reject it on every poll, which an infinite wait can never escape, so the floor only + /// discards a limited number of below-floor states before accepting one. + /// + [TestMethod] + public async Task WaitForOrchestration_ContinuedAsNew_InfiniteWait_FollowsSuccessorPersistedBehindTheTombstoneClock() + { + var store = new FakeInstanceStore(); + + // Previous run of the same instance id, completed long ago. + store.States.Add(CreateState("previous-run", OrchestrationStatus.Completed, BaseTime, "stale output")); + + // Current run, first generation, already continued as new. + store.States.Add(CreateState( + "generation-1", + OrchestrationStatus.ContinuedAsNew, + BaseTime.AddMinutes(5), + "next input")); + + // The successor runs on a worker whose clock is behind and completes in a single episode: its + // one and only state write lands below the tombstone's LastUpdatedTime and is never revised. + store.States.Add(CreateState( + "generation-2", + OrchestrationStatus.Completed, + BaseTime.AddMinutes(4), + "final output")); + + ServiceBusOrchestrationService service = CreateService(store); + + // Guard the test itself: without the budget an infinite wait never returns. + using (var cancellation = new CancellationTokenSource(TimeSpan.FromSeconds(30))) + { + OrchestrationState state; + + try + { + state = await service.WaitForOrchestrationAsync( + InstanceId, + "generation-1", + Timeout.InfiniteTimeSpan, + cancellation.Token); + } + catch (OperationCanceledException) + { + Assert.Fail("An infinite wait never observed a successor that was persisted behind the tombstone's clock."); + throw; + } + + Assert.IsNotNull(state); + Assert.AreEqual("generation-2", state.OrchestrationInstance.ExecutionId, "Returned state from the wrong run."); + Assert.AreEqual("final output", state.Output); + } + } + /// /// Suspended is a pause, not a terminal state: an orchestration that is never resumed has no /// result, so the wait must time out rather than report it as finished. From a2c61e6844e8a56a334917f7d3d4aee2cae6ef0f Mon Sep 17 00:00:00 2001 From: Davide Montanari Date: Tue, 22 Sep 2026 21:02:16 -0700 Subject: [PATCH 19/19] Removed wall-clock dependent code --- .../ServiceBusOrchestrationService.cs | 31 --- .../WaitForOrchestrationTests.cs | 192 +----------------- 2 files changed, 2 insertions(+), 221 deletions(-) diff --git a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs index d860d577c..f759ff591 100644 --- a/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs +++ b/src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs @@ -65,15 +65,6 @@ public class ServiceBusOrchestrationService : IOrchestrationService, IOrchestrat static readonly TimeSpan StatusPollingInterval = TimeSpan.FromSeconds(2); const int DuplicateDetectionWindowInHours = 4; - /// - /// How many below-floor states to discard after a continue-as-new handoff before accepting one. - /// The floor exists only to cover the brief window where the successor row is not readable yet, - /// which resolves in well under this budget. Past it, a state below the floor is far more likely - /// a successor persisted by a worker whose clock is behind than state left by an earlier run, - /// so accepting it beats waiting forever. - /// - const int MaxPreviousRunRejections = 5; - /// /// Orchestration service settings /// @@ -1262,11 +1253,6 @@ public async Task WaitForOrchestrationAsync( bool pinnedToExecution = !string.IsNullOrWhiteSpace(executionId); - // Once we stop tracking a specific execution we must not accept state left behind by an - // earlier run of the same instance id. - DateTime minimumLastUpdatedTime = DateTimeUtils.MinDateTime; - int remainingRejections = 0; - while (!cancellationToken.IsCancellationRequested) { OrchestrationState state = pinnedToExecution @@ -1280,23 +1266,6 @@ public async Task WaitForOrchestrationAsync( if (pinnedToExecution && state?.OrchestrationStatus == OrchestrationStatus.ContinuedAsNew) { pinnedToExecution = false; - minimumLastUpdatedTime = state.LastUpdatedTime; - remainingRejections = MaxPreviousRunRejections; - } - - // State from a previous run of this instance id; the current generation is not readable yet. - // A previous run stopped being updated no later than the continue-as-new that wrote the - // tombstone, while its successor is persisted at or after it. CreatedTime is deliberately - // not used: it is stamped on the client for the first generation but on the worker for every - // continue-as-new generation, so a client running ahead would discard a valid successor - // forever. LastUpdatedTime is rewritten by the persisting worker on every state write, so it - // stays within one tier, but worker clocks still disagree: a successor that completes in a - // single episode inside that skew is never rewritten and would stay below the floor forever. - // The budget bounds that, so an infinite wait cannot hang on it. - if (state != null && remainingRejections > 0 && state.LastUpdatedTime < minimumLastUpdatedTime) - { - remainingRejections--; - state = null; } // ContinuedAsNew is never a final state: a new generation always follows it. The built-in diff --git a/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs b/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs index 0723e3b57..faee45921 100644 --- a/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs +++ b/test/DurableTask.ServiceBus.Tests/WaitForOrchestrationTests.cs @@ -50,8 +50,7 @@ static OrchestrationState CreateState( string executionId, OrchestrationStatus status, DateTime createdTime, - string output = null, - DateTime? lastUpdatedTime = null) + string output = null) { return new OrchestrationState { @@ -62,7 +61,7 @@ static OrchestrationState CreateState( }, OrchestrationStatus = status, CreatedTime = createdTime, - LastUpdatedTime = lastUpdatedTime ?? createdTime, + LastUpdatedTime = createdTime, Output = output }; } @@ -171,193 +170,6 @@ public async Task WaitForOrchestration_Timeout_Zero_ContinuedAsNew_ChecksOnceAnd Assert.AreEqual(1, store.QueryCount, "A zero timeout must perform exactly one lookup."); } - /// - /// After following a continue-as-new to the current generation we must still not accept state - /// left behind by an earlier run of the same instance id. - /// - [TestMethod] - public async Task WaitForOrchestration_ContinuedAsNew_IgnoresStateFromPreviousRun() - { - var store = new FakeInstanceStore(); - - // Previous run: completed long before the current run started. - store.States.Add(CreateState("previous-run", OrchestrationStatus.Completed, BaseTime, "stale output")); - - // Current run, first generation, already continued as new. - store.States.Add(CreateState("generation-1", OrchestrationStatus.ContinuedAsNew, BaseTime.AddMinutes(5), "next input")); - - // The next generation only becomes readable later. - store.OnQuery = s => - { - if (s.QueryCount == 3) - { - s.States.Add(CreateState("generation-2", OrchestrationStatus.Completed, BaseTime.AddMinutes(6), "final output")); - } - }; - - ServiceBusOrchestrationService service = CreateService(store); - - OrchestrationState state = await service.WaitForOrchestrationAsync( - InstanceId, - "generation-1", - TimeSpan.FromSeconds(30), - CancellationToken.None); - - Assert.IsNotNull(state); - Assert.AreEqual("generation-2", state.OrchestrationInstance.ExecutionId, "Returned state from the wrong run."); - Assert.AreEqual("final output", state.Output); - } - - /// - /// CreatedTime comes from a HistoryEvent timestamp and is not guaranteed unique, so a previous - /// run can happen to share the ContinuedAsNew tombstone's CreatedTime. The floor must therefore - /// not depend on CreatedTime at all: LastUpdatedTime alone separates the two, because the - /// previous run stopped being updated before the continue-as-new occurred. - /// - [TestMethod] - public async Task WaitForOrchestration_ContinuedAsNew_IgnoresPreviousRunSharingTheTombstoneCreatedTime() - { - var store = new FakeInstanceStore(); - - // Previous run: completed before the current run continued as new, but its ExecutionStarted - // timestamp collides with the current run's. - store.States.Add(CreateState( - "previous-run", - OrchestrationStatus.Completed, - BaseTime, - "stale output", - BaseTime.AddMinutes(1))); - - // Current run, first generation, already continued as new. - store.States.Add(CreateState( - "generation-1", - OrchestrationStatus.ContinuedAsNew, - BaseTime, - "next input", - BaseTime.AddMinutes(5))); - - // The next generation only becomes readable later. - store.OnQuery = s => - { - if (s.QueryCount == 3) - { - s.States.Add(CreateState("generation-2", OrchestrationStatus.Completed, BaseTime.AddMinutes(5), "final output")); - } - }; - - ServiceBusOrchestrationService service = CreateService(store); - - OrchestrationState state = await service.WaitForOrchestrationAsync( - InstanceId, - "generation-1", - TimeSpan.FromSeconds(30), - CancellationToken.None); - - Assert.IsNotNull(state); - Assert.AreEqual("generation-2", state.OrchestrationInstance.ExecutionId, "Returned state from the wrong run."); - Assert.AreEqual("final output", state.Output); - } - - /// - /// The first generation's ExecutionStartedEvent is stamped on the client that scheduled the - /// orchestration, while a continue-as-new generation's is stamped on the worker that ran it, and - /// both come from DateTime.UtcNow on those machines. A client running ahead of the worker therefore - /// produces a tombstone whose CreatedTime is later than its own successor's. Rejecting state by - /// CreatedTime would discard that successor on every iteration, so the wait would never observe the - /// completion it was asked about. - /// - [TestMethod] - public async Task WaitForOrchestration_ContinuedAsNew_FollowsSuccessorWhenClientClockIsAheadOfWorker() - { - var store = new FakeInstanceStore(); - - // First generation: started by a client running one second ahead of the worker, so its - // CreatedTime is later than everything the worker subsequently stamps. Its LastUpdatedTime is - // written by the worker when the continue-as-new is persisted. - store.States.Add(CreateState( - "generation-1", - OrchestrationStatus.ContinuedAsNew, - BaseTime.AddSeconds(1), - "next input", - BaseTime.AddMilliseconds(100))); - - // Successor: created and persisted entirely on the worker, and already complete. - store.States.Add(CreateState( - "generation-2", - OrchestrationStatus.Completed, - BaseTime.AddMilliseconds(50), - "final output", - BaseTime.AddMilliseconds(200))); - - ServiceBusOrchestrationService service = CreateService(store); - - OrchestrationState state = await service.WaitForOrchestrationAsync( - InstanceId, - "generation-1", - TimeSpan.FromSeconds(30), - CancellationToken.None); - - Assert.IsNotNull(state, "Clock skew between the client and the worker discarded a valid successor."); - Assert.AreEqual("generation-2", state.OrchestrationInstance.ExecutionId); - Assert.AreEqual("final output", state.Output); - } - - /// - /// Worker clocks can disagree too, so a short successor can be created, complete and never be - /// written again while its LastUpdatedTime still sits below the tombstone's. An unbounded floor - /// would reject it on every poll, which an infinite wait can never escape, so the floor only - /// discards a limited number of below-floor states before accepting one. - /// - [TestMethod] - public async Task WaitForOrchestration_ContinuedAsNew_InfiniteWait_FollowsSuccessorPersistedBehindTheTombstoneClock() - { - var store = new FakeInstanceStore(); - - // Previous run of the same instance id, completed long ago. - store.States.Add(CreateState("previous-run", OrchestrationStatus.Completed, BaseTime, "stale output")); - - // Current run, first generation, already continued as new. - store.States.Add(CreateState( - "generation-1", - OrchestrationStatus.ContinuedAsNew, - BaseTime.AddMinutes(5), - "next input")); - - // The successor runs on a worker whose clock is behind and completes in a single episode: its - // one and only state write lands below the tombstone's LastUpdatedTime and is never revised. - store.States.Add(CreateState( - "generation-2", - OrchestrationStatus.Completed, - BaseTime.AddMinutes(4), - "final output")); - - ServiceBusOrchestrationService service = CreateService(store); - - // Guard the test itself: without the budget an infinite wait never returns. - using (var cancellation = new CancellationTokenSource(TimeSpan.FromSeconds(30))) - { - OrchestrationState state; - - try - { - state = await service.WaitForOrchestrationAsync( - InstanceId, - "generation-1", - Timeout.InfiniteTimeSpan, - cancellation.Token); - } - catch (OperationCanceledException) - { - Assert.Fail("An infinite wait never observed a successor that was persisted behind the tombstone's clock."); - throw; - } - - Assert.IsNotNull(state); - Assert.AreEqual("generation-2", state.OrchestrationInstance.ExecutionId, "Returned state from the wrong run."); - Assert.AreEqual("final output", state.Output); - } - } - /// /// Suspended is a pause, not a terminal state: an orchestration that is never resumed has no /// result, so the wait must time out rather than report it as finished.