Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,13 @@ internal sealed class WorkflowOrchestrationContext : WorkflowContext
private readonly Dictionary<string, Queue<TaskCompletionSource<HistoryEvent>>> _externalEventSources = new(StringComparer.OrdinalIgnoreCase);
private readonly Dictionary<int, TaskCompletionSource<HistoryEvent>> _openTasks = [];
private readonly Dictionary<int, string> _taskIdToExecutionId = [];
// TaskExecutionIds as persisted on replayed TaskScheduled events: the
// scheduling-time truth. Re-derivation can drift across SDK upgrades (the
// seed gained per-execution inputs), and completions that arrive with a
// mismatched TaskScheduledId correlate by execution ID alone, so replayed
// tasks must use what history recorded rather than what the current
// derivation produces.
private readonly Dictionary<int, string> _persistedTaskExecutionIds = [];
private readonly Dictionary<string, int> _executionIdToTaskId = new(StringComparer.Ordinal);
private readonly SortedDictionary<int, WorkflowAction> _pendingActions = [];
private readonly IDaprSerializer _workflowSerializer;
Expand All @@ -68,6 +75,7 @@ internal sealed class WorkflowOrchestrationContext : WorkflowContext

// Parse execution/instance ID as GUID or derive a deterministic namespace from the ID
private readonly Guid _instanceGuid;
private readonly Guid _taskExecutionGuid;
private static readonly Guid InstanceIdNamespace = new("6f927a2e-9c7e-4a1d-9b8d-7a86f2e7f62f");

private readonly string? _appId;
Expand All @@ -84,7 +92,8 @@ public WorkflowOrchestrationContext(string name, string instanceId, DateTime cur
IDaprSerializer workflowSerializer, ILoggerFactory loggerFactory, WorkflowVersionTracker versionTracker,
string? appId = null, string? executionId = null,
IReadOnlyList<HistoryEvent>? ownHistory = null,
IEnumerable<PropagatedHistoryChunk>? incomingPropagatedHistory = null)
IEnumerable<PropagatedHistoryChunk>? incomingPropagatedHistory = null,
string? historyExecutionId = null)
{
_workflowSerializer = workflowSerializer;
_loggerFactory = loggerFactory;
Expand All @@ -94,6 +103,22 @@ public WorkflowOrchestrationContext(string name, string instanceId, DateTime cur
_instanceGuid = Guid.TryParse(guidSeed, out var guid)
? guid
: CreateGuidFromName(InstanceIdNamespace, Encoding.UTF8.GetBytes(guidSeed));

// TaskExecutionId derivation gets its own namespace so it can prefer the
// execution ID persisted in history (historyExecutionId) when the request
// carries none. Deliberately NOT folded into _instanceGuid: that namespace
// also feeds NewGuid() and autogenerated child instance IDs, and an
// execution started before an SDK upgrade must keep replaying those from
// the seed it originally used (the instance ID, for old sidecars).
// TaskExecutionId re-derivation is safe across that upgrade because
// completions match by task ID first and fall back to execution-ID
// matching only for not-yet-registered tasks.
var taskExecutionSeed = !string.IsNullOrWhiteSpace(executionId) ? executionId
: !string.IsNullOrWhiteSpace(historyExecutionId) ? historyExecutionId
: instanceId;
Comment on lines +116 to +118

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@JoshVanL This seems reasonable - could you make this change?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks @WhitWaldo, done!

_taskExecutionGuid = Guid.TryParse(taskExecutionSeed, out var taskGuid)
? taskGuid
: CreateGuidFromName(InstanceIdNamespace, Encoding.UTF8.GetBytes(taskExecutionSeed));
Name = name;
InstanceId = instanceId;
_currentUtcDateTime = currentUtcDateTime;
Expand Down Expand Up @@ -160,6 +185,15 @@ private async Task<T> CallActivityInternalAsync<T>(string name, object? input, W
var taskId = _sequenceNumber++;
var taskExecutionId = CreateTaskExecutionId(taskId, name);

// If history already scheduled this task, its persisted TaskExecutionId is
// authoritative: an execution scheduled under a different derivation seed
// (older SDK, no per-execution seeding) must keep correlating by the ID its
// completions actually carry.
if (_persistedTaskExecutionIds.TryGetValue(taskId, out var persistedExecutionId))
{
taskExecutionId = persistedExecutionId;
}

var router = CreateRouter(options?.TargetAppId);

// If the completion arrived before we registered the task, consume it now
Expand Down Expand Up @@ -660,6 +694,24 @@ private void OnTaskScheduled(HistoryEvent historyEvent)
var eventId = historyEvent.EventId;
TryDropOptionalTimerAt(eventId);
_pendingActions.Remove(eventId);

// Record the scheduling-time TaskExecutionId and, when the workflow code
// already re-registered this task under a re-derived ID, re-point the
// mapping at the persisted value so legacy completions still match.
var persisted = historyEvent.TaskScheduled?.TaskExecutionId;
if (string.IsNullOrEmpty(persisted))
{
return;
}
_persistedTaskExecutionIds[eventId] = persisted;

if (_taskIdToExecutionId.TryGetValue(eventId, out var derived) &&
!string.Equals(derived, persisted, StringComparison.Ordinal))
{
_executionIdToTaskId.Remove(derived);
_taskIdToExecutionId[eventId] = persisted;
_executionIdToTaskId[persisted] = eventId;
}
}

/// <summary>
Expand Down Expand Up @@ -889,7 +941,7 @@ private Task<T> HandleFailedActivityFromHistory<T>(string activityName, TaskFail
private string CreateTaskExecutionId(int taskId, string name)
{
var seed = $"{InstanceId}|activity|{taskId}|{name}";
return CreateGuidFromName(_instanceGuid, Encoding.UTF8.GetBytes(seed)).ToString("N");
return CreateGuidFromName(_taskExecutionGuid, Encoding.UTF8.GetBytes(seed)).ToString("N");
}

private static bool TryGetTaskExecutionId(HistoryEvent historyEvent, out string taskExecutionId)
Expand Down
20 changes: 19 additions & 1 deletion src/Dapr.Workflow/Worker/WorkflowWorker.cs
Original file line number Diff line number Diff line change
Expand Up @@ -96,6 +96,21 @@ private async Task<WorkflowResponse> HandleWorkflowResponseAsync(WorkflowRequest
string? serializedInput = null;
string? appId = null;

// TaskExecutionId derivation needs a per-execution seed. Older
// sidecars leave ExecutionId unset on the request, so recover the
// runtime-minted execution ID persisted in the
// ExecutionStartedEvent: it is stable across replays of one
// execution and fresh for each execution of a recreated instance
// ID. Without it, two executions of the same instance ID derive
// colliding TaskExecutionIds, and the runtime's activity
// de-duplication can treat the second execution's activity as a
// duplicate of the first, leaving the workflow stuck in Running.
// It seeds ONLY the TaskExecutionId namespace (see the context
// constructor): the NewGuid / child-instance-ID namespace keeps
// its historical seeding so in-flight executions replay
// identically across an SDK upgrade.
string? historyExecutionId = null;

if (request.RequiresHistoryStreaming)
{
var streamRequest = new GetInstanceHistoryRequest { InstanceId = request.InstanceId };
Expand Down Expand Up @@ -158,6 +173,8 @@ private async Task<WorkflowResponse> HandleWorkflowResponseAsync(WorkflowRequest
workflowName = e.ExecutionStarted.Name;
serializedInput = e.ExecutionStarted.Input;

historyExecutionId = e.ExecutionStarted.WorkflowInstance?.ExecutionId;

// Try pulling the app ID out of the target first, then the source if not available
if (!string.IsNullOrEmpty(e.Router?.TargetAppID))
{
Expand Down Expand Up @@ -328,7 +345,8 @@ private async Task<WorkflowResponse> HandleWorkflowResponseAsync(WorkflowRequest
var context = new WorkflowOrchestrationContext(workflowName, request.InstanceId, currentUtcDateTime,
_serializer, loggerFactory, versionTracker, appId, request.ExecutionId,
allPastEvents,
incomingPropagatedHistory);
incomingPropagatedHistory,
historyExecutionId);

// Deserialize the input
object? input = string.IsNullOrEmpty(serializedInput)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1164,6 +1164,96 @@ public void NewGuid_ShouldVary_ForDifferentExecutionIds()
Assert.NotEqual(g1, g2);
}

[Fact]
public void HistoryExecutionId_ShouldSeedTaskExecutionIdOnly_NotTheGuidNamespace()
{
// An execution started under an old sidecar (no request-level execution ID)
// replays NewGuid and autogenerated child instance IDs from the instance-ID
// seed. Recovering the execution ID from history must give TaskExecutionId a
// per-execution namespace WITHOUT changing that guid namespace, or an
// in-flight execution would stop replaying deterministically across an SDK
// upgrade.
var serializer = new JsonDaprSerializer(new JsonSerializerOptions(JsonSerializerDefaults.Web));
var now = new DateTime(2025, 01, 01, 0, 0, 0, DateTimeKind.Utc);

WorkflowOrchestrationContext Context(string? historyExecutionId) => new(
"wf", "same-instance", now, serializer, NullLoggerFactory.Instance,
new WorkflowVersionTracker([]), historyExecutionId: historyExecutionId);

// NewGuid namespace is untouched by the history-recovered execution ID.
Assert.Equal(Context(null).NewGuid(), Context("exec-1").NewGuid());

static string FirstTaskExecutionId(WorkflowOrchestrationContext ctx)
{
_ = ctx.CallActivityAsync<string>("act");
return Assert.Single(ctx.PendingActions).ScheduleTask!.TaskExecutionId;
}

var legacy = FirstTaskExecutionId(Context(null));
var exec1 = FirstTaskExecutionId(Context("exec-1"));
var exec1Replay = FirstTaskExecutionId(Context("exec-1"));
var exec2 = FirstTaskExecutionId(Context("exec-2"));

// Per-execution and replay-stable: distinct across executions of the same
// instance ID (the collision that stuck recreated workflows in Running),
// identical when the same execution replays.
Assert.NotEqual(legacy, exec1);
Assert.NotEqual(exec1, exec2);
Assert.Equal(exec1, exec1Replay);
}

[Fact]
public async Task CallActivityAsync_ShouldMatchLegacySeededCompletion_AfterSdkUpgrade()
{
// Upgrade-replay: an execution scheduled its activity under the old
// instance-seeded derivation, then the SDK upgraded and now seeds
// TaskExecutionId per execution (historyExecutionId). Replay re-derives a
// DIFFERENT id, but history's TaskScheduled event carries the id the
// completion actually references, so correlation must follow history: the
// mismatched-TaskScheduledId path matches by execution ID alone.
var serializer = new JsonDaprSerializer(new JsonSerializerOptions(JsonSerializerDefaults.Web));
var now = new DateTime(2025, 01, 01, 0, 0, 0, DateTimeKind.Utc);

WorkflowOrchestrationContext Context(string? historyExecutionId) => new(
"wf", "same-instance", now, serializer, NullLoggerFactory.Instance,
new WorkflowVersionTracker([]), historyExecutionId: historyExecutionId);

// The id the legacy (pre-upgrade, instance-seeded) SDK persisted at
// scheduling time.
var legacyContext = Context(null);
_ = legacyContext.CallActivityAsync<string>("Any");
var legacyId = Assert.Single(legacyContext.PendingActions).ScheduleTask!.TaskExecutionId;

// Post-upgrade replay of that same execution.
var upgraded = Context("exec-1");
var task = upgraded.CallActivityAsync<string>("Any");

var history = new[]
{
new HistoryEvent
{
EventId = 0,
TaskScheduled = new TaskScheduledEvent { Name = "Any", TaskExecutionId = legacyId }
},
new HistoryEvent
{
TaskCompleted = new TaskCompletedEvent
{
TaskScheduledId = 123,
TaskExecutionId = legacyId,
Result = "\"ok\""
}
}
};

upgraded.ProcessEvents(history, true);

// Bounded await: without history-authoritative correlation the completion
// never matches and the task parks forever, exactly the production hang.
var result = await task.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken);
Assert.Equal("ok", result);
}

[Fact]
public void NewGuid_ShouldBeDeterministic_ForNonGuidInstanceId()
{
Expand Down
67 changes: 67 additions & 0 deletions test/Dapr.Workflow.Test/Worker/WorkflowWorkerTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -2086,6 +2086,73 @@ public async Task HandleWorkflowResponseAsync_ShouldYield_WhenWorkflowAwaitsActi
Assert.Contains(response.Actions, a => a.ScheduleTask != null);
}

[Fact]
public async Task HandleWorkflowResponseAsync_ShouldDeriveDistinctTaskExecutionIds_AcrossExecutionsOfSameInstanceId()
{
// Models a completed-then-recreated instance: same instance ID, two
// executions. The request-level ExecutionId is unset (older sidecars never
// populate it), so the per-execution guid seed must come from the
// runtime-minted execution ID persisted in the ExecutionStartedEvent.
// If both executions derive the same TaskExecutionId for task 0, the
// runtime's activity de-duplication treats the second execution's activity
// as a redelivery of the first (already completed) one and never runs it,
// stranding the recreated workflow in Running.
var sp = new ServiceCollection().BuildServiceProvider();
var serializer = new JsonDaprSerializer(new JsonSerializerOptions(JsonSerializerDefaults.Web));

var factory = new StubWorkflowsFactory();
factory.AddWorkflow("wf", new InlineWorkflow(
inputType: typeof(object),
run: async (ctx, _) =>
{
await ctx.CallActivityAsync<int>("step");
return null;
}));

var worker = new WorkflowWorker(
CreateGrpcClientMock().Object,
factory,
NullLoggerFactory.Instance,
serializer,
sp);

static WorkflowRequest RequestForExecution(string executionId) => new()
{
InstanceId = "i",
PastEvents =
{
new HistoryEvent
{
ExecutionStarted = new ExecutionStartedEvent
{
Name = "wf",
Input = "",
WorkflowInstance = new WorkflowInstance
{
InstanceId = "i",
ExecutionId = executionId
}
}
}
}
};

var firstRun = await InvokeHandleWorkflowResponseAsync(worker, RequestForExecution("exec-1"));
var secondRun = await InvokeHandleWorkflowResponseAsync(worker, RequestForExecution("exec-2"));
var firstRunReplay = await InvokeHandleWorkflowResponseAsync(worker, RequestForExecution("exec-1"));

var firstTask = Assert.Single(firstRun.Actions, a => a.ScheduleTask != null).ScheduleTask!;
var secondTask = Assert.Single(secondRun.Actions, a => a.ScheduleTask != null).ScheduleTask!;
var replayTask = Assert.Single(firstRunReplay.Actions, a => a.ScheduleTask != null).ScheduleTask!;

Assert.False(string.IsNullOrWhiteSpace(firstTask.TaskExecutionId));
Assert.False(string.IsNullOrWhiteSpace(secondTask.TaskExecutionId));

// Distinct executions must not collide; replays of one execution must be stable.
Assert.NotEqual(firstTask.TaskExecutionId, secondTask.TaskExecutionId);
Assert.Equal(firstTask.TaskExecutionId, replayTask.TaskExecutionId);
}

// -------------------------------------------------------------------------
// Inner exception path (workflow returns a faulted Task)
// -------------------------------------------------------------------------
Expand Down
Loading