diff --git a/composer.json b/composer.json index c5a98fb24..7fca4e97c 100644 --- a/composer.json +++ b/composer.json @@ -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" }, diff --git a/docs/architecture/cancellation-scope.md b/docs/architecture/cancellation-scope.md index a57a3a5b3..a59751185 100644 --- a/docs/architecture/cancellation-scope.md +++ b/docs/architecture/cancellation-scope.md @@ -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. diff --git a/src/V2/TaskWatchdog.php b/src/V2/TaskWatchdog.php index c3cf9553b..f3060b693 100644 --- a/src/V2/TaskWatchdog.php +++ b/src/V2/TaskWatchdog.php @@ -72,6 +72,7 @@ public static function wake(?string $connection = null, ?string $queue = null): * missing_run_failures: list, * deadline_expired_candidates: int, * deadline_expired_tasks_created: int, + * cancellation_deadlines_enforced: int, * deadline_expired_failures: list, * activity_timeout_candidates: int, * activity_timeouts_enforced: int, @@ -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, @@ -353,6 +358,7 @@ private static function recoverMissingTask(string $runId): array * missing_run_failures: list, * deadline_expired_candidates: int, * deadline_expired_tasks_created: int, + * cancellation_deadlines_enforced: int, * deadline_expired_failures: list, * activity_timeout_candidates: int, * activity_timeouts_enforced: int, @@ -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, @@ -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 */ @@ -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]) @@ -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 !== []) { @@ -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() @@ -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; @@ -521,7 +566,7 @@ private static function createDeadlineExpiredTask(string $runId): array ); return $task; - }); + }, 3); if ($task instanceof WorkflowTask) { TaskDispatcher::dispatch($task); @@ -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, ]; } diff --git a/tests/Feature/V2/V2FinallyCleanupTest.php b/tests/Feature/V2/V2FinallyCleanupTest.php index df414118c..167d87f42 100644 --- a/tests/Feature/V2/V2FinallyCleanupTest.php +++ b/tests/Feature/V2/V2FinallyCleanupTest.php @@ -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(); } @@ -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 */ diff --git a/tests/Feature/V2/V2WorkflowTaskBridgeTest.php b/tests/Feature/V2/V2WorkflowTaskBridgeTest.php index 3be5d6dd6..9d2fee167 100644 --- a/tests/Feature/V2/V2WorkflowTaskBridgeTest.php +++ b/tests/Feature/V2/V2WorkflowTaskBridgeTest.php @@ -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, [[ @@ -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'), @@ -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() diff --git a/tests/Unit/Commands/V2RepairPassCommandTest.php b/tests/Unit/Commands/V2RepairPassCommandTest.php index 328212c7d..b6ce14598 100644 --- a/tests/Unit/Commands/V2RepairPassCommandTest.php +++ b/tests/Unit/Commands/V2RepairPassCommandTest.php @@ -4,6 +4,7 @@ namespace Tests\Unit\Commands; +use Illuminate\Support\Carbon; use Illuminate\Support\Facades\Cache; use Illuminate\Support\Facades\Queue; use Illuminate\Support\Str; @@ -12,12 +13,14 @@ use Workflow\Serializers\CodecRegistry; use Workflow\Serializers\Serializer; use Workflow\V2\Contracts\MatchingRole; +use Workflow\V2\Enums\HistoryEventType; use Workflow\V2\Enums\RunStatus; use Workflow\V2\Enums\SignalStatus; use Workflow\V2\Enums\TaskStatus; use Workflow\V2\Enums\TaskType; use Workflow\V2\Jobs\RunWorkflowTask; use Workflow\V2\Models\WorkerCompatibilityHeartbeat; +use Workflow\V2\Models\WorkflowHistoryEvent; use Workflow\V2\Models\WorkflowInstance; use Workflow\V2\Models\WorkflowRun; use Workflow\V2\Models\WorkflowRunSummary; @@ -25,6 +28,7 @@ use Workflow\V2\Models\WorkflowTask; use Workflow\V2\Support\WorkerCompatibilityFleet; use Workflow\V2\TaskWatchdog; +use Workflow\V2\WorkflowStub; final class V2RepairPassCommandTest extends TestCase { @@ -66,6 +70,7 @@ public function testItIgnoresTheLoopThrottleByDefaultAndRepairsMissingSignalTask '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, @@ -131,6 +136,7 @@ public function testRespectThrottleOptionSkipsRepairWhenTheLoopThrottleIsHeld(): '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, @@ -248,6 +254,7 @@ public function runPass( '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, @@ -275,6 +282,7 @@ public function runPass( '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, @@ -326,6 +334,7 @@ public function testRunIdScopeRepairsOnlyTheSelectedMissingTaskRun(): void '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, @@ -381,6 +390,7 @@ public function testInstanceIdScopeRepairsOnlyTheSelectedInstanceTasks(): void '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, @@ -453,6 +463,7 @@ public function testConnectionAndQueueScopeRepairsOnlyMatchingExistingAndMissing '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, @@ -551,6 +562,7 @@ public function testCommaSeparatedQueueScopeRepairsEveryRequestedQueue(): void '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, @@ -624,6 +636,7 @@ public function testConnectionAndQueueScopeRepairsOnlyMatchingDeadlineExpiredRun 'missing_run_failures' => [], 'deadline_expired_candidates' => 1, 'deadline_expired_tasks_created' => 1, + 'cancellation_deadlines_enforced' => 0, 'deadline_expired_failures' => [], 'activity_timeout_candidates' => 0, 'activity_timeouts_enforced' => 0, @@ -649,6 +662,54 @@ public function testConnectionAndQueueScopeRepairsOnlyMatchingDeadlineExpiredRun Queue::assertPushed(RunWorkflowTask::class, 1); } + public function testCleanupDeadlineClosesForeignRunsWithNoWorkerWithinRequestedQueue(): void + { + Queue::fake(); + + foreach ([TaskStatus::Completed, TaskStatus::Ready, TaskStatus::Leased] as $status) { + $run = $this->createWaitingRun('cleanup-deadline-' . $status->value, queue: 'critical'); + $request = WorkflowStub::loadRun($run->id)->requestCancellation('worker unavailable', 60); + $run->refresh() + ->forceFill([ + 'compatibility' => 'foreign-sdk-cleanup', + ])->save(); + $deadline = $run->cancellation_deadline_at; + $this->assertNotNull($deadline); + $task = $run->tasks() + ->where('task_type', TaskType::Workflow->value)->sole(); + $task->forceFill([ + 'status' => $status, + 'compatibility' => 'foreign-sdk-cleanup', + 'lease_owner' => $status === TaskStatus::Leased ? 'departed-worker' : null, + 'lease_expires_at' => $status === TaskStatus::Leased ? $deadline->copy()->addMinutes(5) : null, + ])->save(); + $other = $this->createWaitingRun('cleanup-other-queue-' . $status->value, queue: 'default'); + WorkflowStub::loadRun($other->id)->requestCancellation('different queue', 60); + + try { + Carbon::setTestNow($deadline->copy()->addSecond()); + $report = TaskWatchdog::runPass(connection: 'redis', queue: 'critical'); + $this->assertSame(1, $report['cancellation_deadlines_enforced']); + $this->assertSame(0, $report['deadline_expired_tasks_created']); + $this->assertSame([], $report['deadline_expired_failures']); + $this->assertSame(RunStatus::Cancelled, $run->refresh()->status); + $this->assertSame(RunStatus::Waiting, $other->refresh()->status); + $this->assertSame(0, $run->tasks()->whereIn('status', [ + TaskStatus::Ready->value, + TaskStatus::Leased->value, + ])->count()); + $terminal = WorkflowHistoryEvent::query()->where('workflow_run_id', $run->id) + ->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']); + $again = TaskWatchdog::runPass(connection: 'redis', queue: 'critical'); + $this->assertSame(0, $again['cancellation_deadlines_enforced']); + } finally { + Carbon::setTestNow(); + } + } + } + public function testConnectionAndQueueScopeEnforcesOnlyMatchingActivityTimeouts(): void { Queue::fake(); @@ -682,6 +743,7 @@ public function testConnectionAndQueueScopeEnforcesOnlyMatchingActivityTimeouts( 'missing_run_failures' => [], 'deadline_expired_candidates' => 0, 'deadline_expired_tasks_created' => 0, + 'cancellation_deadlines_enforced' => 0, 'deadline_expired_failures' => [], 'activity_timeout_candidates' => 1, 'activity_timeouts_enforced' => 1, @@ -741,6 +803,7 @@ public function testLoopModeRunsTheRequestedNumberOfIterationsAndForcesThrottleA '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, @@ -764,6 +827,7 @@ public function testLoopModeRunsTheRequestedNumberOfIterationsAndForcesThrottleA '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,