From dec69a1490da991d00012b409f042819e7a33350 Mon Sep 17 00:00:00 2001 From: Paul Stadler Date: Fri, 29 May 2026 16:05:09 -0400 Subject: [PATCH] fix(workflow): use server-blocking WaitForInstance* RPCs instead of client polling WaitForWorkflowStartAsync and WaitForWorkflowCompletionAsync polled GetInstance on a fixed 500ms/1s cadence. The sidecar already exposes server-blocking WaitForInstanceStart / WaitForInstanceCompletion RPCs that return the moment the instance reaches the target state, and the go-sdk (via durabletask-go) already uses them. This switches both methods to the blocking RPCs, wrapped in an exponential-backoff retry that re-issues on transient interruptions (DeadlineExceeded / Unavailable) and propagates cancellation, mirroring durabletask-go's client. Calls continue to flow through CreateCallOptions so the dapr-api-token is honored. Removes the now-dead client-poll tests and adds coverage for the blocking RPC path, transient-error retry, non-transient propagation, and cancellation. Fixes #1838 Co-Authored-By: Claude Opus 4.8 (1M context) Signed-off-by: Paul Stadler --- .../Client/WorkflowGrpcClient.cs | 111 ++++++---- src/Dapr.Workflow/Logging.cs | 3 + .../Client/WorkflowGrpcClientTests.cs | 208 +++++++----------- 3 files changed, 158 insertions(+), 164 deletions(-) diff --git a/src/Dapr.Workflow/Client/WorkflowGrpcClient.cs b/src/Dapr.Workflow/Client/WorkflowGrpcClient.cs index 7971bb531..c56e289f7 100644 --- a/src/Dapr.Workflow/Client/WorkflowGrpcClient.cs +++ b/src/Dapr.Workflow/Client/WorkflowGrpcClient.cs @@ -92,50 +92,43 @@ public override async Task ScheduleNewWorkflowAsync(string workflowName, public override async Task WaitForWorkflowStartAsync(string instanceId, bool getInputsAndOutputs = true, CancellationToken cancellationToken = default) { - // Poll until the workflow status (not Pending) - while (true) - { - var metadata = await GetWorkflowMetadataAsync(instanceId, getInputsAndOutputs, cancellationToken); - - if (metadata is null) - { - var ex = new InvalidOperationException($"Workflow instance '{instanceId}' does not exist"); - logger.LogWaitForStartException(ex, instanceId); - throw ex; - } + var response = await WaitForInstanceAsync( + (request, callOptions) => grpcClient.WaitForInstanceStartAsync(request, callOptions), + instanceId, + getInputsAndOutputs, + cancellationToken); - if (metadata.RuntimeStatus != WorkflowRuntimeStatus.Pending) - { - logger.LogWaitForStartCompleted(instanceId, metadata.RuntimeStatus); - return metadata; - } - - await Task.Delay(TimeSpan.FromMilliseconds(500), cancellationToken); + if (!response.Exists) + { + var ex = new InvalidOperationException($"Workflow instance '{instanceId}' does not exist"); + logger.LogWaitForStartException(ex, instanceId); + throw ex; } + + var metadata = ProtoConverters.ToWorkflowMetadata(response.WorkflowState, serializer); + logger.LogWaitForStartCompleted(instanceId, metadata.RuntimeStatus); + return metadata; } /// public override async Task WaitForWorkflowCompletionAsync(string instanceId, bool getInputsAndOutputs = true, CancellationToken cancellationToken = default) { - while (true) - { - var metadata = await GetWorkflowMetadataAsync(instanceId, getInputsAndOutputs, cancellationToken); - - if (metadata is null) - { - var ex = new InvalidOperationException($"Workflow instance '{instanceId}' does not exist"); - logger.LogWaitForCompletionException(ex, instanceId); - throw ex; - } + var response = await WaitForInstanceAsync( + (request, callOptions) => grpcClient.WaitForInstanceCompletionAsync(request, callOptions), + instanceId, + getInputsAndOutputs, + cancellationToken); - if (IsTerminalStatus(metadata.RuntimeStatus)) - { - logger.LogWaitForCompletionCompleted(instanceId, metadata.RuntimeStatus); - return metadata; - } - - await Task.Delay(TimeSpan.FromSeconds(1), cancellationToken); + if (!response.Exists) + { + var ex = new InvalidOperationException($"Workflow instance '{instanceId}' does not exist"); + logger.LogWaitForCompletionException(ex, instanceId); + throw ex; } + + var metadata = ProtoConverters.ToWorkflowMetadata(response.WorkflowState, serializer); + logger.LogWaitForCompletionCompleted(instanceId, metadata.RuntimeStatus); + return metadata; } /// @@ -333,8 +326,48 @@ public override ValueTask DisposeAsync() private CallOptions CreateCallOptions(CancellationToken cancellationToken) => DaprClientUtilities.ConfigureGrpcCallOptions(typeof(DaprWorkflowClient).Assembly, daprApiToken, cancellationToken); - private static bool IsTerminalStatus(WorkflowRuntimeStatus status) => - status is WorkflowRuntimeStatus.Completed - or WorkflowRuntimeStatus.Failed - or WorkflowRuntimeStatus.Terminated; + private static readonly TimeSpan MinWaitRetryDelay = TimeSpan.FromMilliseconds(50); + private static readonly TimeSpan MaxWaitRetryDelay = TimeSpan.FromSeconds(15); + + /// + /// Issues a server-side blocking wait RPC (WaitForInstanceStart / + /// WaitForInstanceCompletion), which returns the moment the instance reaches the + /// target state. The blocking call can be interrupted by a deadline or a transient + /// sidecar/connection drop; in that case it is re-issued with exponential backoff rather + /// than failing the wait — mirroring the durabletask-go client. The loop ends only when the + /// RPC succeeds or the caller cancels. + /// + private async Task WaitForInstanceAsync( + Func> waitCall, + string instanceId, + bool getInputsAndOutputs, + CancellationToken cancellationToken) + { + var request = new grpc.GetInstanceRequest + { + InstanceId = instanceId, + GetInputsAndOutputs = getInputsAndOutputs + }; + + var delay = MinWaitRetryDelay; + while (true) + { + cancellationToken.ThrowIfCancellationRequested(); + + try + { + var grpcCallOptions = CreateCallOptions(cancellationToken); + return await waitCall(request, grpcCallOptions); + } + catch (RpcException ex) when ( + !cancellationToken.IsCancellationRequested && + ex.StatusCode is StatusCode.DeadlineExceeded or StatusCode.Unavailable) + { + logger.LogWaitForInstanceRetry(ex, instanceId, ex.StatusCode, delay); + await Task.Delay(delay, cancellationToken); + delay = TimeSpan.FromMilliseconds( + Math.Min(delay.TotalMilliseconds * 2, MaxWaitRetryDelay.TotalMilliseconds)); + } + } + } } diff --git a/src/Dapr.Workflow/Logging.cs b/src/Dapr.Workflow/Logging.cs index 14b117fbc..72e3f3e29 100644 --- a/src/Dapr.Workflow/Logging.cs +++ b/src/Dapr.Workflow/Logging.cs @@ -155,6 +155,9 @@ public static partial void LogWaitForStartException(this ILogger logger, Invalid [LoggerMessage(LogLevel.Debug, "Workflow '{InstanceId}' completed with status '{Status}'")] public static partial void LogWaitForCompletionCompleted(this ILogger logger, string instanceId, WorkflowRuntimeStatus status); + [LoggerMessage(LogLevel.Debug, "Wait for workflow instance '{InstanceId}' was interrupted with status '{StatusCode}'; retrying in {Delay}")] + public static partial void LogWaitForInstanceRetry(this ILogger logger, RpcException ex, string instanceId, StatusCode statusCode, TimeSpan delay); + [LoggerMessage(LogLevel.Information, "Raised event '{EventName}' to workflow '{InstanceId}'")] public static partial void LogRaisedEvent(this ILogger logger, string eventName, string instanceId); diff --git a/test/Dapr.Workflow.Test/Client/WorkflowGrpcClientTests.cs b/test/Dapr.Workflow.Test/Client/WorkflowGrpcClientTests.cs index 7197f30c2..cc958995f 100644 --- a/test/Dapr.Workflow.Test/Client/WorkflowGrpcClientTests.cs +++ b/test/Dapr.Workflow.Test/Client/WorkflowGrpcClientTests.cs @@ -152,13 +152,13 @@ public async Task GetWorkflowMetadataAsync_ShouldPassThroughGetInputsAndOutputsF } [Fact] - public async Task WaitForWorkflowStartAsync_ShouldReturnImmediately_WhenStatusIsNotPending() + public async Task WaitForWorkflowStartAsync_ShouldUseServerBlockingRpc_AndReturnStartedState() { var serializer = new JsonDaprSerializer(); var grpcClientMock = CreateGrpcClientMock(); grpcClientMock - .Setup(x => x.GetInstanceAsync(It.IsAny(), It.IsAny())) + .Setup(x => x.WaitForInstanceStartAsync(It.IsAny(), It.IsAny())) .Returns(CreateAsyncUnaryCall(new GetInstanceResponse { Exists = true, @@ -176,6 +176,8 @@ public async Task WaitForWorkflowStartAsync_ShouldReturnImmediately_WhenStatusIs Assert.Equal("i", result.InstanceId); Assert.Equal(WorkflowRuntimeStatus.Running, result.RuntimeStatus); + grpcClientMock.Verify(x => x.WaitForInstanceStartAsync(It.IsAny(), It.IsAny()), Times.Once); + grpcClientMock.Verify(x => x.GetInstanceAsync(It.IsAny(), It.IsAny()), Times.Never); } [Fact] @@ -185,7 +187,7 @@ public async Task WaitForWorkflowStartAsync_ShouldThrowInvalidOperationException var grpcClientMock = CreateGrpcClientMock(); grpcClientMock - .Setup(x => x.GetInstanceAsync(It.IsAny(), It.IsAny())) + .Setup(x => x.WaitForInstanceStartAsync(It.IsAny(), It.IsAny())) .Returns(CreateAsyncUnaryCall(new GetInstanceResponse { Exists = false })); var client = new WorkflowGrpcClient(grpcClientMock.Object, NullLogger.Instance, serializer); @@ -194,13 +196,13 @@ public async Task WaitForWorkflowStartAsync_ShouldThrowInvalidOperationException } [Fact] - public async Task WaitForWorkflowCompletionAsync_ShouldReturnImmediately_WhenStatusIsTerminal() + public async Task WaitForWorkflowCompletionAsync_ShouldUseServerBlockingRpc_AndReturnTerminalState() { var serializer = new JsonDaprSerializer(); var grpcClientMock = CreateGrpcClientMock(); grpcClientMock - .Setup(x => x.GetInstanceAsync(It.IsAny(), It.IsAny())) + .Setup(x => x.WaitForInstanceCompletionAsync(It.IsAny(), It.IsAny())) .Returns(CreateAsyncUnaryCall(new GetInstanceResponse { Exists = true, @@ -220,6 +222,79 @@ public async Task WaitForWorkflowCompletionAsync_ShouldReturnImmediately_WhenSta Assert.Equal("i", result.InstanceId); Assert.Equal(WorkflowRuntimeStatus.Completed, result.RuntimeStatus); Assert.Equal("{\"ok\":true}", result.SerializedOutput); + grpcClientMock.Verify(x => x.WaitForInstanceCompletionAsync(It.IsAny(), It.IsAny()), Times.Once); + grpcClientMock.Verify(x => x.GetInstanceAsync(It.IsAny(), It.IsAny()), Times.Never); + } + + [Fact] + public async Task WaitForWorkflowCompletionAsync_ShouldRetry_WhenWaitIsInterruptedByTransientError() + { + var serializer = new JsonDaprSerializer(); + + var attempts = 0; + var grpcClientMock = CreateGrpcClientMock(); + grpcClientMock + .Setup(x => x.WaitForInstanceCompletionAsync(It.IsAny(), It.IsAny())) + .Returns(() => + { + attempts++; + if (attempts == 1) + { + return CreateAsyncUnaryCallThrows( + new RpcException(new Status(StatusCode.DeadlineExceeded, "deadline"))); + } + + return CreateAsyncUnaryCall(new GetInstanceResponse + { + Exists = true, + WorkflowState = new Dapr.DurableTask.Protobuf.WorkflowState + { + InstanceId = "i", + Name = "n", + WorkflowStatus = OrchestrationStatus.Completed + } + }); + }); + + var client = new WorkflowGrpcClient(grpcClientMock.Object, NullLogger.Instance, serializer); + + var result = await client.WaitForWorkflowCompletionAsync("i", getInputsAndOutputs: true, cancellationToken: TestContext.Current.CancellationToken); + + Assert.Equal(WorkflowRuntimeStatus.Completed, result.RuntimeStatus); + Assert.Equal(2, attempts); + } + + [Fact] + public async Task WaitForWorkflowCompletionAsync_ShouldNotRetry_WhenWaitFailsWithNonTransientError() + { + var serializer = new JsonDaprSerializer(); + + var grpcClientMock = CreateGrpcClientMock(); + grpcClientMock + .Setup(x => x.WaitForInstanceCompletionAsync(It.IsAny(), It.IsAny())) + .Returns(CreateAsyncUnaryCallThrows( + new RpcException(new Status(StatusCode.InvalidArgument, "bad")))); + + var client = new WorkflowGrpcClient(grpcClientMock.Object, NullLogger.Instance, serializer); + + var ex = await Assert.ThrowsAsync(() => client.WaitForWorkflowCompletionAsync("i", getInputsAndOutputs: true, cancellationToken: TestContext.Current.CancellationToken)); + Assert.Equal(StatusCode.InvalidArgument, ex.StatusCode); + grpcClientMock.Verify(x => x.WaitForInstanceCompletionAsync(It.IsAny(), It.IsAny()), Times.Once); + } + + [Fact] + public async Task WaitForWorkflowCompletionAsync_ShouldThrowOperationCanceled_WhenCancelled() + { + var serializer = new JsonDaprSerializer(); + + using var cts = new CancellationTokenSource(); + await cts.CancelAsync(); + + var grpcClientMock = CreateGrpcClientMock(); + + var client = new WorkflowGrpcClient(grpcClientMock.Object, NullLogger.Instance, serializer); + + await Assert.ThrowsAsync(() => client.WaitForWorkflowCompletionAsync("i", getInputsAndOutputs: true, cancellationToken: cts.Token)); } [Fact] @@ -535,100 +610,10 @@ public async Task DisposeAsync_ShouldCompleteSynchronously() await client.DisposeAsync(); } - [Fact] - public async Task WaitForWorkflowStartAsync_ShouldPollUntilStatusIsNotPending() - { - var serializer = new JsonDaprSerializer(); - - var grpcClientMock = CreateGrpcClientMock(); - - var requests = new List(); - var callCount = 0; - - grpcClientMock - .Setup(x => x.GetInstanceAsync(It.IsAny(), It.IsAny())) - .Returns((GetInstanceRequest request, CallOptions _) => - { - requests.Add(request); - callCount++; - - var status = callCount == 1 ? OrchestrationStatus.Pending : OrchestrationStatus.Running; - - return CreateAsyncUnaryCall(new GetInstanceResponse - { - Exists = true, - WorkflowState = new Dapr.DurableTask.Protobuf.WorkflowState - { - InstanceId = "i", - Name = "n", - WorkflowStatus = status - } - }); - }); - - var client = new WorkflowGrpcClient(grpcClientMock.Object, NullLogger.Instance, serializer); - - using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(5)); - var result = await client.WaitForWorkflowStartAsync("i", getInputsAndOutputs: false, cancellationToken: cts.Token); - - Assert.Equal("i", result.InstanceId); - Assert.Equal(WorkflowRuntimeStatus.Running, result.RuntimeStatus); - - Assert.True(requests.Count >= 2); - Assert.All(requests, r => Assert.Equal("i", r.InstanceId)); - Assert.All(requests, r => Assert.False(r.GetInputsAndOutputs)); - } - - [Fact] - public async Task WaitForWorkflowCompletionAsync_ShouldPollUntilTerminalStatus() - { - var serializer = new JsonDaprSerializer(); - - var grpcClientMock = CreateGrpcClientMock(); - - var requests = new List(); - var callCount = 0; - - grpcClientMock - .Setup(x => x.GetInstanceAsync(It.IsAny(), It.IsAny())) - .Returns((GetInstanceRequest request, CallOptions _) => - { - requests.Add(request); - callCount++; - - var status = callCount == 1 ? OrchestrationStatus.Running : OrchestrationStatus.Completed; - - return CreateAsyncUnaryCall(new GetInstanceResponse - { - Exists = true, - WorkflowState = new Dapr.DurableTask.Protobuf.WorkflowState - { - InstanceId = "i", - Name = "n", - WorkflowStatus = status, - Output = status == OrchestrationStatus.Completed ? "\"done\"" : string.Empty - } - }); - }); - - var client = new WorkflowGrpcClient(grpcClientMock.Object, NullLogger.Instance, serializer); - - using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(5)); - var result = await client.WaitForWorkflowCompletionAsync("i", getInputsAndOutputs: true, cancellationToken: cts.Token); - - Assert.Equal("i", result.InstanceId); - Assert.Equal(WorkflowRuntimeStatus.Completed, result.RuntimeStatus); - Assert.Equal("\"done\"", result.SerializedOutput); - - Assert.True(requests.Count >= 2); - Assert.All(requests, r => Assert.Equal("i", r.InstanceId)); - Assert.All(requests, r => Assert.True(r.GetInputsAndOutputs)); - } - [Theory] [InlineData(OrchestrationStatus.Failed, WorkflowRuntimeStatus.Failed)] [InlineData(OrchestrationStatus.Terminated, WorkflowRuntimeStatus.Terminated)] - public async Task WaitForWorkflowCompletionAsync_ShouldReturnImmediately_WhenStatusIsTerminalFailedOrTerminated( + public async Task WaitForWorkflowCompletionAsync_ShouldMapTerminalStatus_ForFailedOrTerminated( OrchestrationStatus protoStatus, WorkflowRuntimeStatus expectedStatus) { @@ -636,7 +621,7 @@ public async Task WaitForWorkflowCompletionAsync_ShouldReturnImmediately_WhenSta var grpcClientMock = CreateGrpcClientMock(); grpcClientMock - .Setup(x => x.GetInstanceAsync(It.IsAny(), It.IsAny())) + .Setup(x => x.WaitForInstanceCompletionAsync(It.IsAny(), It.IsAny())) .Returns(CreateAsyncUnaryCall(new GetInstanceResponse { Exists = true, @@ -663,7 +648,7 @@ public async Task WaitForWorkflowCompletionAsync_ShouldThrowInvalidOperationExce var grpcClientMock = CreateGrpcClientMock(); grpcClientMock - .Setup(x => x.GetInstanceAsync(It.IsAny(), It.IsAny())) + .Setup(x => x.WaitForInstanceCompletionAsync(It.IsAny(), It.IsAny())) .Returns(CreateAsyncUnaryCall(new GetInstanceResponse { Exists = false })); var client = new WorkflowGrpcClient(grpcClientMock.Object, NullLogger.Instance, serializer); @@ -674,33 +659,6 @@ public async Task WaitForWorkflowCompletionAsync_ShouldThrowInvalidOperationExce Assert.Contains("missing", ex.Message); } - [Fact] - public async Task WaitForWorkflowCompletionAsync_ShouldRespectCancellationToken_WhileWaiting() - { - var serializer = new JsonDaprSerializer(); - - var grpcClientMock = CreateGrpcClientMock(); - grpcClientMock - .Setup(x => x.GetInstanceAsync(It.IsAny(), It.IsAny())) - .Returns(CreateAsyncUnaryCall(new GetInstanceResponse - { - Exists = true, - WorkflowState = new Dapr.DurableTask.Protobuf.WorkflowState - { - InstanceId = "i", - Name = "n", - WorkflowStatus = OrchestrationStatus.Running - } - })); - - var client = new WorkflowGrpcClient(grpcClientMock.Object, NullLogger.Instance, serializer); - - using var cts = new CancellationTokenSource(TimeSpan.FromMilliseconds(50)); - - await Assert.ThrowsAnyAsync( - () => client.WaitForWorkflowCompletionAsync("i", getInputsAndOutputs: true, cancellationToken: cts.Token)); - } - private static Mock CreateGrpcClientMock() { var callInvoker = new Mock(MockBehavior.Loose);