-
Notifications
You must be signed in to change notification settings - Fork 335
[DurableTask.ServiceBus] Fixed executionId being ignored in ServiceBusOrchestrationService.WaitForOrchestrationAsync #1403
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
ca5f18d
5a232c6
556b30b
b417710
d9567bb
1765143
a31d3ba
e0a2e2c
6f416b4
4625987
0578772
c698660
d8a1397
e16cb28
ec1079a
beaa47c
c7667c3
62a96d2
09442fa
98ce63b
1c8f6bf
a2c61e6
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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; | ||
|
|
||
| /// <summary> | ||
|
|
@@ -1222,9 +1222,12 @@ public async Task ForceTerminateTaskOrchestrationAsync(string instanceId, string | |
| /// <summary> | ||
| /// Wait for an orchestration to reach any terminal state within the given timeout | ||
| /// </summary> | ||
| /// <param name="executionId">The execution id of the orchestration</param> | ||
| /// <param name="executionId">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 <see cref="OrchestrationStatus.ContinuedAsNew"/> execution, | ||
| /// where the current generation of the instance is followed instead.</param> | ||
| /// <param name="instanceId">Instance to wait for</param> | ||
| /// <param name="timeout">Max timeout to wait</param> | ||
| /// <param name="timeout">Max timeout to wait. Only positive <see cref="TimeSpan"/> values, <see cref="TimeSpan.Zero"/>, or <see cref="Timeout.InfiniteTimeSpan"/> are allowed.</param> | ||
| /// <param name="cancellationToken">Task cancellation token</param> | ||
| public async Task<OrchestrationState> WaitForOrchestrationAsync( | ||
| string instanceId, | ||
|
|
@@ -1236,20 +1239,68 @@ public async Task<OrchestrationState> WaitForOrchestrationAsync( | |
|
|
||
| if (string.IsNullOrWhiteSpace(instanceId)) | ||
| { | ||
| throw new ArgumentException("instanceId"); | ||
| throw new ArgumentException("The instance id cannot be null, empty, or whitespace.", nameof(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 unless it is Timeout.InfiniteTimeSpan." + | ||
| $" The value for {nameof(timeout)} was '{timeout}'." + | ||
| $" Please provide a positive timeout value, TimeSpan.Zero or Timeout.InfiniteTimeSpan.", | ||
| nameof(timeout)); | ||
| } | ||
|
|
||
| bool pinnedToExecution = !string.IsNullOrWhiteSpace(executionId); | ||
|
|
||
| while (!cancellationToken.IsCancellationRequested && timeoutSeconds > 0) | ||
| while (!cancellationToken.IsCancellationRequested) | ||
| { | ||
| OrchestrationState state = (await GetOrchestrationStateAsync(instanceId, false))?.FirstOrDefault(); | ||
| OrchestrationState state = pinnedToExecution | ||
| ? await GetOrchestrationStateAsync(instanceId, executionId) | ||
| : (await GetOrchestrationStateAsync(instanceId, false))?.FirstOrDefault(); | ||
|
Comment on lines
+1258
to
+1260
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. [P2] Preserve the existing store retry policy for execution-specific polling Switching to the execution-specific overload also changes failure handling. AzureTableInstanceStore.GetOrchestrationStateAsync(instanceId, bool) wraps its table queries in Utils.ExecuteWithRetries (AzureTableInstanceStore.cs:179-189), but GetOrchestrationStateAsync(instanceId, executionId) queries the tables directly (lines 220-233). Consequently, a transient failure escaping the SDK's own retries now immediately faults WaitForOrchestrationAsync whenever an execution ID is supplied, instead of receiving the existing outer retries. I reproduced this with the actual AzureTableInstanceStore and a table-client stub that throws one 503 and succeeds on the next query: the pinned path fails, while the unpinned path retries and returns the completed state. Both paths pass with the pre-PR implementation. The regression reproduces on net8.0 and net48. Please retain the existing retry behavior around the pinned store lookup and add coverage for a transient failure followed by a successful read.
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Thanks Bernd Verst (@berndverst). Fixed. |
||
|
|
||
| // 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; | ||
| } | ||
|
davidemontanari marked this conversation as resolved.
|
||
|
|
||
| // 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; | ||
| TimeSpan delay = StatusPollingInterval; | ||
|
|
||
| if (!isInfiniteTimeSpan) | ||
| { | ||
| // 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(delay, cancellationToken); | ||
| } | ||
| else | ||
| { | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -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; | ||
|
|
||
| /// <summary> | ||
| /// Verifies that the execution-id lookup on <see cref="AzureTableInstanceStore"/> 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. | ||
| /// </summary> | ||
| [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" | ||
| }); | ||
| } | ||
|
|
||
| /// <summary> | ||
| /// A transient failure on the state table must be retried rather than surfaced to the caller. | ||
| /// </summary> | ||
| [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."); | ||
| } | ||
|
|
||
| /// <summary> | ||
| /// 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. | ||
| /// </summary> | ||
| [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."); | ||
| } | ||
|
|
||
| /// <summary> | ||
| /// 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. | ||
| /// </summary> | ||
| 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<IEnumerable<AzureTableOrchestrationStateEntity>> QueryOrchestrationStatesAsync( | ||
| OrchestrationStateQuery stateQuery) | ||
| { | ||
| this.StateQueryCount++; | ||
|
|
||
| if (this.StateQueryCount <= this.FailuresBeforeSuccess) | ||
| { | ||
| throw new RequestFailedException(503, "The server is busy."); | ||
| } | ||
|
|
||
| IEnumerable<AzureTableOrchestrationStateEntity> result = this.StateTableIsEmpty | ||
| ? Enumerable.Empty<AzureTableOrchestrationStateEntity>() | ||
| : new[] { CreateEntity() }; | ||
|
|
||
| return Task.FromResult(result); | ||
| } | ||
|
|
||
| public override Task<IEnumerable<AzureTableOrchestrationStateEntity>> QueryJumpStartOrchestrationsAsync( | ||
| OrchestrationStateQuery stateQuery) | ||
| { | ||
| this.JumpStartQueryCount++; | ||
|
|
||
| if (this.JumpStartQueryCount <= this.JumpStartFailuresBeforeSuccess) | ||
| { | ||
| throw new RequestFailedException(503, "The server is busy."); | ||
| } | ||
|
|
||
| return Task.FromResult<IEnumerable<AzureTableOrchestrationStateEntity>>(new[] { CreateEntity() }); | ||
| } | ||
| } | ||
| } | ||
| } |
Uh oh!
There was an error while loading. Please reload this page.