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.2",
"product-train": "2.3.3",
"laravel-embedded-upgrade-contract": "resources/laravel-embedded-upgrade-contract.json",
"laravel-dependency-security-policy": "resources/laravel-dependency-security-policy.json"
},
Expand Down
8 changes: 6 additions & 2 deletions docs/architecture/cancellation-scope.md
Original file line number Diff line number Diff line change
Expand Up @@ -157,8 +157,12 @@ try {
The cleanup activity uses ordinary durable retries and can itself wait on a
timer. A successful cleanup closes the run as cancelled. An unhandled cleanup
failure leaves a failed run rather than claiming cleanup succeeded. If the
deadline expires, the watchdog requests an immediate terminal cancel, even
if cleanup is still waiting. `terminate()` remains immediate and can stop a
deadline expires, the runtime's next repair pass closes the run as cancelled
and revokes outstanding task authority, even when cleanup is waiting or no
compatible workflow worker is running. The repair pass does not execute the
workflow definition to enforce expiry. Its JSON report includes
`cancellation_deadlines_enforced` for these closures.
`terminate()` remains immediate and can stop a
run during cleanup; it does not wait for `finally` to finish. The shield is
not a permanent opt-out from cancellation or a protection against process
termination.
Expand Down
75 changes: 61 additions & 14 deletions src/V2/TaskWatchdog.php
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,7 @@ public static function wake(?string $connection = null, ?string $queue = null):
* missing_run_failures: list<array{run_id: string, message: string}>,
* deadline_expired_candidates: int,
* deadline_expired_tasks_created: int,
* cancellation_deadlines_enforced: int,
* deadline_expired_failures: list<array{run_id: string, message: string}>,
* activity_timeout_candidates: int,
* activity_timeouts_enforced: int,
Expand Down Expand Up @@ -166,6 +167,10 @@ public static function runPass(
$report['dispatched_tasks']++;
}

if ($result['cancellation_enforced']) {
$report['cancellation_deadlines_enforced']++;
}

if ($result['error'] !== null) {
$report['deadline_expired_failures'][] = [
'run_id' => $deadlineRunId,
Expand Down Expand Up @@ -353,6 +358,7 @@ private static function recoverMissingTask(string $runId): array
* missing_run_failures: list<array{run_id: string, message: string}>,
* deadline_expired_candidates: int,
* deadline_expired_tasks_created: int,
* cancellation_deadlines_enforced: int,
* deadline_expired_failures: list<array{run_id: string, message: string}>,
* activity_timeout_candidates: int,
* activity_timeouts_enforced: int,
Expand Down Expand Up @@ -384,6 +390,7 @@ private static function emptyReport(
'missing_run_failures' => [],
'deadline_expired_candidates' => 0,
'deadline_expired_tasks_created' => 0,
'cancellation_deadlines_enforced' => 0,
'deadline_expired_failures' => [],
'activity_timeout_candidates' => 0,
'activity_timeouts_enforced' => 0,
Expand All @@ -392,8 +399,8 @@ private static function emptyReport(
}

/**
* Find non-terminal runs with expired execution or run deadlines
* that have no open workflow task to detect the timeout.
* Find expired cleanup deadlines regardless of worker availability,
* and other expired deadlines without an open workflow task.
*
* @return list<string>
*/
Expand All @@ -403,7 +410,8 @@ private static function deadlineExpiredRunIds(
?string $connection = null,
?string $queue = null,
): array {
$now = now();
$now = now()
->format('Y-m-d H:i:s.u');

$query = WorkflowRun::query()
->whereIn('status', [RunStatus::Pending->value, RunStatus::Running->value, RunStatus::Waiting->value])
Expand All @@ -420,9 +428,15 @@ private static function deadlineExpiredRunIds(
->where('cancellation_deadline_at', '<=', $now);
});
})
->whereDoesntHave('tasks', static function ($task): void {
$task->where('task_type', TaskType::Workflow->value)
->whereIn('status', [TaskStatus::Ready->value, TaskStatus::Leased->value]);
->where(static function ($available) use ($now): void {
$available->whereDoesntHave('tasks', static function ($task): void {
$task->where('task_type', TaskType::Workflow->value)
->whereIn('status', [TaskStatus::Ready->value, TaskStatus::Leased->value]);
})->orWhere(static function ($cancellation) use ($now): void {
$cancellation->whereNotNull('cancellation_request_command_id')
->whereNotNull('cancellation_deadline_at')
->where('cancellation_deadline_at', '<=', $now);
});
});

if ($runIds !== []) {
Expand Down Expand Up @@ -462,12 +476,24 @@ private static function tablesReady(): bool
}

/**
* @return array{task: WorkflowTask|null, error: string|null}
* @return array{task: WorkflowTask|null, error: string|null, cancellation_enforced: bool}
*/
private static function createDeadlineExpiredTask(string $runId): array
{
$cancellationEnforced = false;

try {
$task = DB::transaction(static function () use ($runId): ?WorkflowTask {
$task = DB::transaction(static function () use ($runId, &$cancellationEnforced): ?WorkflowTask {
$cancellationEnforced = false;

// Worker protocol mutations lock the task before the run.
$existingWorkflowTask = WorkflowTask::query()
->where('workflow_run_id', $runId)
->where('task_type', TaskType::Workflow->value)
->whereIn('status', [TaskStatus::Ready->value, TaskStatus::Leased->value])
->lockForUpdate()
->first();

/** @var WorkflowRun|null $run */
$run = WorkflowRun::query()
->lockForUpdate()
Expand All @@ -490,11 +516,30 @@ private static function createDeadlineExpiredTask(string $runId): array
return null;
}

$existingWorkflowTask = WorkflowTask::query()
->where('workflow_run_id', $run->id)
->where('task_type', TaskType::Workflow->value)
->whereIn('status', [TaskStatus::Ready->value, TaskStatus::Leased->value])
->first();
if (is_string($run->cancellation_request_command_id)
&& $run->cancellation_deadline_at !== null
&& $now->gte($run->cancellation_deadline_at)) {
// Expiry is runtime authority. A foreign workflow definition
// or an absent SDK worker must not keep the run open.
$controlTask = $existingWorkflowTask ?? WorkflowTask::query()->create([
'workflow_run_id' => $run->id,
'namespace' => $run->namespace,
'task_type' => TaskType::Workflow->value,
'status' => TaskStatus::Ready->value,
'available_at' => $now,
'payload' => [
'reason' => 'cleanup_deadline_expired',
],
'connection' => $run->connection,
'queue' => $run->queue,
'compatibility' => $run->compatibility,
'repair_count' => 1,
]);
$cancellationEnforced = app(WorkflowExecutor::class)
->cancelIfCleanupDeadlineExpired($run, $controlTask);

return null;
}

if ($existingWorkflowTask !== null) {
return null;
Expand All @@ -521,7 +566,7 @@ private static function createDeadlineExpiredTask(string $runId): array
);

return $task;
});
}, 3);

if ($task instanceof WorkflowTask) {
TaskDispatcher::dispatch($task);
Expand All @@ -532,12 +577,14 @@ private static function createDeadlineExpiredTask(string $runId): array
return [
'task' => null,
'error' => $throwable->getMessage(),
'cancellation_enforced' => false,
];
}

return [
'task' => $task,
'error' => null,
'cancellation_enforced' => $cancellationEnforced,
];
}

Expand Down
58 changes: 56 additions & 2 deletions tests/Feature/V2/V2FinallyCleanupTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -799,8 +799,8 @@ public function testWatchdogEnforcesDeadlineWhileCleanupTimerIsPending(): void
Carbon::setTestNow($workflow->run()?->cancellation_deadline_at?->copy()->addSecond());
$report = TaskWatchdog::runPass(respectThrottle: false, runIds: [$runId]);
$this->assertSame(1, $report['deadline_expired_candidates']);
$this->assertSame(1, $report['deadline_expired_tasks_created']);
$this->runReadyTask($runId, TaskType::Workflow);
$this->assertSame(0, $report['deadline_expired_tasks_created']);
$this->assertSame(1, $report['cancellation_deadlines_enforced']);
} finally {
Carbon::setTestNow();
}
Expand All @@ -809,6 +809,60 @@ public function testWatchdogEnforcesDeadlineWhileCleanupTimerIsPending(): void
$this->assertSame([], $workflow->memo());
}

public function testWatchdogClosesExpiredCleanupWithoutACompatibleWorker(): void
{
Queue::fake();

foreach ([TaskStatus::Ready, TaskStatus::Leased] as $status) {
$workflow = WorkflowStub::make(TestFinallyCleanupWorkflow::class);
$workflow->start(false, 3600);
$runId = $workflow->runId();
$this->assertIsString($runId);
$this->runReadyTask($runId, TaskType::Workflow);
$request = $workflow->requestCancellation('stop unavailable worker', 60);
$run = WorkflowRun::query()->findOrFail($runId);
$run->forceFill([
'compatibility' => 'foreign-sdk-cleanup',
])->save();
$task = WorkflowTask::query()->where('workflow_run_id', $runId)
->where('task_type', TaskType::Workflow->value)
->where('status', TaskStatus::Ready->value)->sole();
$task->forceFill([
'status' => $status,
'compatibility' => 'foreign-sdk-cleanup',
'lease_owner' => $status === TaskStatus::Leased ? 'departed-sdk-worker' : null,
'lease_expires_at' => $status === TaskStatus::Leased
? $run->cancellation_deadline_at->copy()
->addMinutes(5) : null,
])->save();

$before = TaskWatchdog::runPass(respectThrottle: false, runIds: [$runId]);
$this->assertSame(0, $before['cancellation_deadlines_enforced']);
$this->assertFalse($workflow->refresh()->cancelled());

try {
Carbon::setTestNow($run->cancellation_deadline_at->copy()->addSecond());
$report = TaskWatchdog::runPass(respectThrottle: false, runIds: [$runId]);
$this->assertSame(1, $report['cancellation_deadlines_enforced']);
$this->assertSame([], $report['deadline_expired_failures']);
$this->assertTrue($workflow->refresh()->cancelled());
$this->assertSame(TaskStatus::Cancelled, $task->fresh()->status);
$this->assertNull($task->fresh()->lease_expires_at);
$terminal = WorkflowHistoryEvent::query()->where('workflow_run_id', $runId)
->where('event_type', HistoryEventType::WorkflowCancelled->value)->sole();
$this->assertSame($request->commandId(), $terminal->workflow_command_id);
$this->assertSame('Cooperative cancellation cleanup deadline expired.', $terminal->payload['reason']);
$this->assertSame([], $workflow->memo());
$again = TaskWatchdog::runPass(respectThrottle: false, runIds: [$runId]);
$this->assertSame(0, $again['cancellation_deadlines_enforced']);
$this->assertSame(1, WorkflowHistoryEvent::query()->where('workflow_run_id', $runId)
->where('event_type', HistoryEventType::WorkflowCancelled->value)->count());
} finally {
Carbon::setTestNow();
}
}
}

private function runReadyTask(string $runId, TaskType $type): void
{
/** @var WorkflowTask|null $task */
Expand Down
8 changes: 5 additions & 3 deletions tests/Feature/V2/V2WorkflowTaskBridgeTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -3804,6 +3804,7 @@ public function testPortableCleanupDeadlineCancelsOpenResources(string $boundary
Carbon::setTestNow($deadline);
$report = TaskWatchdog::runPass(respectThrottle: false, runIds: [$run->id]);
$this->assertSame(0, $report['deadline_expired_tasks_created']);
$this->assertSame(1, $report['cancellation_deadlines_enforced']);
$result = $boundary === 'heartbeat'
? $this->bridge->heartbeat($task->id)
: $this->bridge->complete($task->id, [[
Expand Down Expand Up @@ -3857,9 +3858,10 @@ public function testPortableCleanupDeadlineRejectsCompletionAfterWatchdogReclaim
Carbon::setTestNow($deadline);
$report = TaskWatchdog::runPass(respectThrottle: false, runIds: [$run->id]);
$this->assertSame(1, $report['repaired_existing_tasks']);
$this->assertSame(TaskStatus::Ready, $task->refresh()->status);
$this->assertSame(1, $report['cancellation_deadlines_enforced']);
$this->assertSame(TaskStatus::Cancelled, $task->refresh()->status);
$claimed = $this->bridge->claimStatus($task->id, 'replacement-cleanup-worker');
$this->assertTrue($claimed['claimed']);
$this->assertFalse($claimed['claimed']);
$result = $this->bridge->complete($task->id, [[
'type' => 'complete_workflow',
'result' => Serializer::serialize('late replacement result'),
Expand All @@ -3869,7 +3871,7 @@ public function testPortableCleanupDeadlineRejectsCompletionAfterWatchdogReclaim
}

$this->assertFalse($result['completed']);
$this->assertSame('run_cancelled', $result['reason']);
$this->assertSame('task_not_leased', $result['reason']);
$this->assertSame(RunStatus::Cancelled, $run->refresh()->status);
$this->assertNull($run->output);
$this->assertSame(1, $run->historyEvents()
Expand Down
Loading
Loading