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
2 changes: 1 addition & 1 deletion composer.json
Original file line number Diff line number Diff line change
Expand Up @@ -79,7 +79,7 @@
"dev-main": "2.0.x-dev"
},
"durable-workflow": {
"product-train": "2.3.1",
"product-train": "2.3.2",
"laravel-embedded-upgrade-contract": "resources/laravel-embedded-upgrade-contract.json",
"laravel-dependency-security-policy": "resources/laravel-dependency-security-policy.json"
},
Expand Down
1 change: 1 addition & 0 deletions scripts/ci/validate-regression-corpus.py
Original file line number Diff line number Diff line change
Expand Up @@ -160,6 +160,7 @@
"tests/Fixtures/V2/TestParallelChildReplayWorkflow.php",
"tests/Fixtures/V2/TestSequentialChildReplayWorkflow.php",
"tests/Fixtures/V2/TestServiceResponseReplayWorkflow.php",
"tests/Fixtures/V2/TestServiceGroupedConditionReopenWorkflow.php",
"tests/Fixtures/V2/TestSignalResumedParallelWorkflow.php",
"tests/Unit/V2/ReplayRegressionCorpusTest.php",
),
Expand Down
17 changes: 7 additions & 10 deletions src/V2/Jobs/RunTimerTask.php
Original file line number Diff line number Diff line change
Expand Up @@ -195,17 +195,13 @@ public function handle(): void
'lease_expires_at' => null,
])->save();

$operationKind = match (true) {
$conditionWaitId !== null => 'condition',
$signalWaitId !== null => 'signal',
default => 'timer',
};
if ($parallelMetadataPath !== []) {
ParallelChildGroup::claimSelectionWinner(
$run,
$parallelMetadataPath,
match (true) {
$conditionWaitId !== null => 'condition',
$signalWaitId !== null => 'signal',
default => 'timer',
},
$firedEvent,
);
ParallelChildGroup::claimSelectionWinner($run, $parallelMetadataPath, $operationKind, $firedEvent);
}

if (
Expand All @@ -214,6 +210,7 @@ public function handle(): void
$run,
$parallelMetadataPath,
TimerStatus::Fired,
$operationKind,
)
) {
$this->projectRun($run, self::PROJECTION_RUN_RELATIONS);
Expand Down
130 changes: 120 additions & 10 deletions src/V2/Support/DefaultWorkflowTaskBridge.php
Original file line number Diff line number Diff line change
Expand Up @@ -1053,7 +1053,7 @@ public function complete(string $taskId, array $commands): array
];
}

if (! self::parallelCommandsMatchSequences($parsed['non_terminal'], $sequence)) {
if (! self::parallelCommandsMatchSequences($parsed['non_terminal'], $sequence, $run)) {
return [
'completed' => false,
'task_id' => $taskId,
Expand All @@ -1065,7 +1065,11 @@ public function complete(string $taskId, array $commands): array
}

$this->recordAppliedSignalForSignalResume($run, $task);
$this->recordSatisfiedConditionWaitForSignalResume($run, $task);
$conditionSelectionResolution = $this->recordSatisfiedConditionWaitForSignalResume(
$run,
$task,
$parsed['non_terminal'],
);

foreach ($parsed['non_terminal'] as $command) {
$sequence = $this->applyNonTerminalCommand($run, $task, $command, $sequence, $createdTaskIds);
Expand All @@ -1082,7 +1086,13 @@ public function complete(string $taskId, array $commands): array
$this->applyWorkflowFailure($run, $task, $terminal);
}
} else {
$this->markRunWaiting($run, $task, $parsed['non_terminal'], $createdTaskIds);
$this->markRunWaiting(
$run,
$task,
$parsed['non_terminal'],
$createdTaskIds,
$conditionSelectionResolution,
);
}

return [
Expand Down Expand Up @@ -1292,6 +1302,7 @@ private function markRunWaiting(
WorkflowTask $task,
array $nonTerminalCommands,
array &$createdTaskIds,
?WorkflowHistoryEvent $conditionSelectionResolution = null,
): void {
$run->forceFill([
'status' => RunStatus::Waiting,
Expand All @@ -1313,6 +1324,25 @@ private function markRunWaiting(
$createdTaskIds[] = $nextMessageTask->id;
}

if ($conditionSelectionResolution !== null
&& ! WorkflowTask::query()->where('workflow_run_id', $run->id)
->where('task_type', TaskType::Workflow->value)
->whereIn('status', [TaskStatus::Ready->value, TaskStatus::Leased->value])
->exists()) {
$replayTask = WorkflowTask::query()->create([
'workflow_run_id' => $run->id,
'namespace' => $run->namespace,
'task_type' => TaskType::Workflow->value,
'status' => TaskStatus::Ready->value,
'available_at' => now(),
'payload' => WorkflowTaskPayload::forConditionResolution($conditionSelectionResolution),
'connection' => $run->connection,
'queue' => $run->queue,
'compatibility' => $run->compatibility,
]);
$createdTaskIds[] = $replayTask->id;
}

if (self::commandsIncludeChildWorkflowStart($nonTerminalCommands)) {
self::projectRunBestEffort($run, self::PROJECTION_RUN_RELATIONS_WITH_HISTORY, 'child_workflow_parent_wait');
} else {
Expand Down Expand Up @@ -2059,19 +2089,24 @@ private function validateUpdateCommands(WorkflowRun $run, WorkflowTask $task, ar
* either re-opening the wait or advancing to the next command. When a signal
* resume advances, make that resolution explicit in history for replay and
* Waterline instead of leaving only SignalReceived as an implicit cue.
*
* @param list<array{type: string, ...}> $commands
*/
private function recordSatisfiedConditionWaitForSignalResume(WorkflowRun $run, WorkflowTask $task): void
{
private function recordSatisfiedConditionWaitForSignalResume(
WorkflowRun $run,
WorkflowTask $task,
array $commands
): ?WorkflowHistoryEvent {
$taskPayload = is_array($task->payload) ? $task->payload : [];

if (($taskPayload['resume_source_kind'] ?? null) !== 'workflow_signal') {
return;
return null;
}

$wait = $this->latestOpenConditionWait($run);

if ($wait === null) {
return;
return null;
}

$this->markConditionWaitSignalConsumed($run, $task, $wait);
Expand Down Expand Up @@ -2100,9 +2135,21 @@ private function recordSatisfiedConditionWaitForSignalResume(WorkflowRun $run, W
'signal_wait_id' => self::nonEmptyString($taskPayload['signal_wait_id'] ?? null),
...$parallelMetadata,
], static fn (mixed $value): bool => $value !== null), $task);
ParallelChildGroup::claimSelectionWinner($run, $parallelPath, 'condition', $satisfiedEvent);
$occurrenceId = $wait['condition_wait_occurrence_id'];
$reopened = false;
foreach ($commands as $command) {
if ($occurrenceId !== null && $command['type'] === 'open_condition_wait'
&& ($command['condition_wait_occurrence_id'] ?? null) === $occurrenceId) {
$reopened = true;
break;
}
}
$selectionResolved = ! $reopened
&& ParallelChildGroup::claimSelectionWinner($run, $parallelPath, 'condition', $satisfiedEvent);

$this->cancelOpenConditionTimer($run, $task, $wait);

return $selectionResolved ? $satisfiedEvent : null;
}

/**
Expand Down Expand Up @@ -5495,16 +5542,27 @@ private static function parallelMetadataForCommand(array $command): array
/**
* @param list<array{type: string, ...}> $commands
*/
private static function parallelCommandsMatchSequences(array $commands, int $baseSequence): bool
private static function parallelCommandsMatchSequences(array $commands, int $baseSequence, WorkflowRun $run): bool
{
$commandsBySequence = [];
$reopenedOccurrences = [];
$sequence = $baseSequence;
foreach ($commands as $command) {
if (($command['type'] ?? null) === 'cancel_selection_operation') {
continue;
}
$commandsBySequence[$sequence] = $command;
$path = $command['parallel_group_path'] ?? null;
if (is_array($path) && ($command['type'] ?? null) === 'open_condition_wait'
&& self::recordedGroupedConditionReopenMatches($run, $command)) {
$occurrenceId = (string) $command['condition_wait_occurrence_id'];
if (isset($reopenedOccurrences[$occurrenceId])) {
return false;
}
$reopenedOccurrences[$occurrenceId] = true;
++$sequence;
continue;
}
$commandsBySequence[$sequence] = $command;
if (is_array($path)) {
foreach ($path as $entry) {
if (! is_array($entry)
Expand Down Expand Up @@ -5578,6 +5636,58 @@ private static function parallelCommandsMatchSequences(array $commands, int $bas
return true;
}

/**
* A physical reopen keeps the original authored group/member path. Only
* existing, unresolved condition history can authorize that exception to
* the complete new-group batch and sequence checks.
*
* @param array<string, mixed> $command
*/
private static function recordedGroupedConditionReopenMatches(WorkflowRun $run, array $command): bool
{
$occurrenceId = self::nonEmptyString($command['condition_wait_occurrence_id'] ?? null);
if ($occurrenceId === null) {
return false;
}
$opens = $run->historyEvents->filter(
static fn (WorkflowHistoryEvent $event): bool => $event->event_type === HistoryEventType::ConditionWaitOpened
&& ($event->payload['condition_wait_occurrence_id'] ?? null) === $occurrenceId,
);
$original = $opens->first();
$latest = $opens->last();
if (! $original instanceof WorkflowHistoryEvent || ! $latest instanceof WorkflowHistoryEvent) {
return false;
}
$path = ParallelChildGroup::metadataPathFromPayload($command);
if ($path === []
|| $path !== ParallelChildGroup::metadataPathFromPayload($original->payload)
|| $path !== ParallelChildGroup::metadataPathFromPayload($latest->payload)) {
return false;
}
foreach (['condition_key', 'condition_definition_fingerprint', 'timeout_seconds'] as $field) {
if (($command[$field] ?? null) !== ($original->payload[$field] ?? null)
|| ($command[$field] ?? null) !== ($latest->payload[$field] ?? null)) {
return false;
}
}
foreach ($path as $entry) {
if ($entry['parallel_group_base_sequence'] + $entry['parallel_group_index']
!== ($original->payload['sequence'] ?? null)) {
return false;
}
}

return ! $run->historyEvents->contains(
static fn (WorkflowHistoryEvent $event): bool => $event->sequence > $latest->sequence
&& in_array(
$event->event_type,
[HistoryEventType::ConditionWaitSatisfied, HistoryEventType::ConditionWaitTimedOut],
true
)
&& ($event->payload['condition_wait_occurrence_id'] ?? null) === $occurrenceId,
);
}

/**
* @param array<string, mixed> $command
* @return array{
Expand Down
Loading
Loading