Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
22 commits
Select commit Hold shift + click to select a range
ca5f18d
Fixed executionId usage in WaitForOrchestrationAsync in ServiceBusOrc…
davidemontanari Sep 14, 2026
5a232c6
Marked ContinuedAsNew and Suspended as non-terminal states
davidemontanari Sep 14, 2026
556b30b
Added new tests
davidemontanari Sep 14, 2026
b417710
Fixed doc
davidemontanari Sep 14, 2026
d9567bb
Merge branch 'main' into davidemontanari/dtfx-sb-orchestration-wait-e…
davidemontanari Sep 14, 2026
1765143
Fixed check when execution is ContinuedAsNew and timeout is TimeSpan.…
davidemontanari Sep 14, 2026
a31d3ba
Merge branch 'davidemontanari/dtfx-sb-orchestration-wait-executionid'…
davidemontanari Sep 14, 2026
e0a2e2c
Fixed timeout check
davidemontanari Sep 14, 2026
6f416b4
- Fixed timeout validation error message to include TimeSpan.Zero
davidemontanari Sep 14, 2026
4625987
Renamed test file path
davidemontanari Sep 14, 2026
0578772
Fixed test
davidemontanari Sep 14, 2026
c698660
Added check on LastUpdatedTime to ignore previous runs
davidemontanari Sep 14, 2026
d8a1397
Fixed stale-result race when no executionId is passed
davidemontanari Sep 15, 2026
e16cb28
Reverted escalation code
davidemontanari Sep 15, 2026
ec1079a
Removed checks
davidemontanari Sep 15, 2026
beaa47c
Added check on state from previous run
davidemontanari Sep 15, 2026
c7667c3
Preserved the existing store retry policy for execution-specific polling
davidemontanari Sep 22, 2026
62a96d2
Added tests for transient failures when passing executionId
davidemontanari Sep 23, 2026
09442fa
Merge branch 'main' into davidemontanari/dtfx-sb-orchestration-wait-e…
davidemontanari Sep 23, 2026
98ce63b
Renamed test file
davidemontanari Sep 23, 2026
1c8f6bf
Removed check on CreatedTime and short circuit the check on LastUpdat…
davidemontanari Sep 23, 2026
a2c61e6
Removed wall-clock dependent code
davidemontanari Sep 23, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
71 changes: 61 additions & 10 deletions src/DurableTask.ServiceBus/ServiceBusOrchestrationService.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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>
Expand Down Expand Up @@ -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,
Expand All @@ -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 thread
davidemontanari marked this conversation as resolved.
Comment on lines +1258 to +1260

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The 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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The 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;
}
Comment thread
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
{
Expand Down
4 changes: 2 additions & 2 deletions src/DurableTask.ServiceBus/Tracking/AzureTableClient.cs
Original file line number Diff line number Diff line change
Expand Up @@ -124,7 +124,7 @@ internal async Task DeleteJumpStartTableIfExistsAsync()
}
}

public Task<IEnumerable<AzureTableOrchestrationStateEntity>> QueryOrchestrationStatesAsync(
public virtual Task<IEnumerable<AzureTableOrchestrationStateEntity>> QueryOrchestrationStatesAsync(
OrchestrationStateQuery stateQuery)
{
var query = CreateQueryInternal(stateQuery, false);
Expand All @@ -138,7 +138,7 @@ public async Task<Page<AzureTableOrchestrationStateEntity>> QueryOrchestrationSt
return await QueryTableSegmentAsync<AzureTableOrchestrationStateEntity>(this.historyTableClient, query, continuationToken, count);
}

public Task<IEnumerable<AzureTableOrchestrationStateEntity>> QueryJumpStartOrchestrationsAsync(OrchestrationStateQuery stateQuery)
public virtual Task<IEnumerable<AzureTableOrchestrationStateEntity>> QueryJumpStartOrchestrationsAsync(OrchestrationStateQuery stateQuery)
{
// TODO: Enable segmented query for paging purpose
var query = CreateQueryInternal(stateQuery, true);
Expand Down
33 changes: 29 additions & 4 deletions src/DurableTask.ServiceBus/Tracking/AzureTableInstanceStore.cs
Original file line number Diff line number Diff line change
Expand Up @@ -90,6 +90,20 @@ public AzureTableInstanceStore(string hubName, Uri endpoint, TokenCredential cre
DateTimeUtils.SetMinDateTimeForStorageEmulator();
}

/// <summary>
/// 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.
/// </summary>
/// <param name="tableClient">The table client to issue queries through</param>
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();
}

/// <summary>
/// Runs initialization to prepare the storage for use
/// </summary>
Expand Down Expand Up @@ -217,17 +231,28 @@ public async Task<IEnumerable<OrchestrationStateInstanceEntity>> GetOrchestratio
/// <returns>The matching orchestration state or null if not found</returns>
public async Task<OrchestrationStateInstanceEntity> 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();
}
Expand Down
151 changes: 151 additions & 0 deletions test/DurableTask.ServiceBus.Tests/AzureTableInstanceStoreRetryTests.cs
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() });
}
}
}
}
Loading
Loading