From 226025ffe3e0c1af98a77c01932a50f41799d8d8 Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Thu, 1 Oct 2026 10:14:29 +0000 Subject: [PATCH 1/2] fix: renew workflow heartbeats under the ownership fence --- app/Http/Controllers/Api/WorkerController.php | 31 ++-- composer.json | 2 +- composer.lock | 2 +- docker-compose.dedicated-matching.yml | 4 +- docker-compose.memo-rolling.yml | 4 +- docker-compose.published.yml | 4 +- docker-compose.small-cluster.yml | 2 +- docker-compose.yml | 8 +- k8s/README.md | 10 +- k8s/helm/durable-workflow/Chart.yaml | 6 +- k8s/helm/durable-workflow/README.md | 2 +- .../ci/existing-secrets-values.yaml | 2 +- .../ci/ingress-and-hpa-values.yaml | 2 +- .../ci/inline-secrets-values.yaml | 2 +- .../durable-workflow/templates/_helpers.tpl | 2 +- k8s/helm/durable-workflow/values.yaml | 2 +- k8s/helm/examples/values-dev.yaml | 2 +- .../values-external-secrets-operator.yaml | 2 +- .../values-production-existing-secrets.yaml | 2 +- k8s/migration-job.yaml | 2 +- k8s/scheduler-cronjob.yaml | 2 +- k8s/secret.yaml | 2 +- k8s/server-deployment.yaml | 2 +- k8s/worker-deployment.yaml | 2 +- resources/release/source-release.json | 4 +- scripts/k8s-kind-smoke.sh | 2 +- .../Feature/WorkflowTaskHeartbeatRaceTest.php | 139 ++++++++++++++++++ .../WorkflowTaskHeartbeatRaceRequest.php | 94 ++++++++++++ 28 files changed, 291 insertions(+), 49 deletions(-) create mode 100644 tests/Feature/WorkflowTaskHeartbeatRaceTest.php create mode 100644 tests/Fixtures/WorkflowTaskHeartbeatRaceRequest.php diff --git a/app/Http/Controllers/Api/WorkerController.php b/app/Http/Controllers/Api/WorkerController.php index db0d7501..66dc4624 100644 --- a/app/Http/Controllers/Api/WorkerController.php +++ b/app/Http/Controllers/Api/WorkerController.php @@ -3197,22 +3197,27 @@ public function heartbeatWorkflowTask(Request $request, string $taskId): JsonRes 'workflow_task_attempt' => ['required', 'integer', 'min:1'], ]); - if ($response = $this->guardWorkflowTaskOwnership( - $request, - $namespace, - $taskId, - (int) $validated['workflow_task_attempt'], - $validated['lease_owner'], - )) { - return $response; - } - /** @var WorkflowTaskBridge $bridge */ $bridge = app(WorkflowTaskBridge::class); try { $status = $this->storageMutations->run( - static fn (): array => $bridge->heartbeat($taskId), + fn (): array|JsonResponse => DB::transaction(function () use ($request, $namespace, $taskId, $validated, $bridge): array|JsonResponse { + // Ownership must remain valid until renewal commits. A + // check before this lock can acknowledge a reclaimed lease. + NamespaceWorkflowScope::taskQuery($namespace)->lockForUpdate()->find($taskId); + if ($response = $this->guardWorkflowTaskOwnership( + $request, + $namespace, + $taskId, + (int) $validated['workflow_task_attempt'], + $validated['lease_owner'], + )) { + return $response; + } + + return $bridge->heartbeat($taskId); + }), ); } catch (\Throwable $exception) { if (! BackendLockPressure::is($exception)) { @@ -3226,6 +3231,10 @@ public function heartbeatWorkflowTask(Request $request, string $taskId): JsonRes ); } + if ($status instanceof JsonResponse) { + return $status; + } + return WorkerProtocol::json([ 'task_id' => $taskId, 'workflow_task_attempt' => (int) $validated['workflow_task_attempt'], diff --git a/composer.json b/composer.json index 7cf422d1..4f2a1c13 100644 --- a/composer.json +++ b/composer.json @@ -48,7 +48,7 @@ }, "extra": { "durable-workflow": { - "product-train": "2.4.35" + "product-train": "2.4.36" }, "laravel": { "dont-discover": [] diff --git a/composer.lock b/composer.lock index 05211656..5746ada0 100644 --- a/composer.lock +++ b/composer.lock @@ -4,7 +4,7 @@ "Read more about it at https://getcomposer.org/doc/01-basic-usage.md#installing-dependencies", "This file is @generated automatically" ], - "content-hash": "a5e39f71b2b8115b4f62db356199ca72", + "content-hash": "b6edb8a6e839966bb3af27c6e72d7ed3", "packages": [ { "name": "apache/avro", diff --git a/docker-compose.dedicated-matching.yml b/docker-compose.dedicated-matching.yml index 4a7681bb..b877c59a 100644 --- a/docker-compose.dedicated-matching.yml +++ b/docker-compose.dedicated-matching.yml @@ -32,13 +32,13 @@ name: durable-workflow-server # daemon reports `shape: dedicated`. # Generated by scripts/ci/sync-source-release.mjs. Do not edit the fallback. -x-server-image: &server-image ${DW_SERVER_IMAGE:-durableworkflow/server:${DW_SERVER_TAG:-2.4.35}} +x-server-image: &server-image ${DW_SERVER_IMAGE:-durableworkflow/server:${DW_SERVER_TAG:-2.4.36}} x-server-environment: &server-environment APP_NAME: "Durable Workflow Server" APP_ENV: ${APP_ENV:-local} DW_SERVER_KEY: ${DW_SERVER_KEY:-} - APP_VERSION: ${APP_VERSION:-${DW_SERVER_TAG:-2.4.35}} + APP_VERSION: ${APP_VERSION:-${DW_SERVER_TAG:-2.4.36}} APP_DEBUG: ${APP_DEBUG:-false} DB_CONNECTION: mysql DB_HOST: mysql diff --git a/docker-compose.memo-rolling.yml b/docker-compose.memo-rolling.yml index ca29104e..cc7caceb 100644 --- a/docker-compose.memo-rolling.yml +++ b/docker-compose.memo-rolling.yml @@ -49,14 +49,14 @@ services: command: ["server-bootstrap"] environment: <<: *runtime-environment - APP_VERSION: ${APP_VERSION:-2.4.35} + APP_VERSION: ${APP_VERSION:-2.4.36} successor: image: ${DW_MEMO_SUCCESSOR_IMAGE:-durable-workflow/server-memo-rolling:local} ports: !override [] environment: <<: *runtime-environment - APP_VERSION: ${APP_VERSION:-2.4.35} + APP_VERSION: ${APP_VERSION:-2.4.36} DW_SERVER_ID: memo-successor DW_SERVER_TOPOLOGY_SHAPE: standalone_server DW_SERVER_PROCESS_CLASS: server_http_node diff --git a/docker-compose.published.yml b/docker-compose.published.yml index 7a8afee4..a81dfed7 100644 --- a/docker-compose.published.yml +++ b/docker-compose.published.yml @@ -1,13 +1,13 @@ name: durable-workflow-server # Generated by scripts/ci/sync-source-release.mjs. Do not edit the fallback. -x-server-image: &server-image ${DW_SERVER_IMAGE:-durableworkflow/server:${DW_SERVER_TAG:-2.4.35}} +x-server-image: &server-image ${DW_SERVER_IMAGE:-durableworkflow/server:${DW_SERVER_TAG:-2.4.36}} x-server-environment: &server-environment APP_NAME: "Durable Workflow Server" APP_ENV: ${APP_ENV:-local} DW_SERVER_KEY: ${DW_SERVER_KEY:-} - APP_VERSION: ${APP_VERSION:-${DW_SERVER_TAG:-2.4.35}} + APP_VERSION: ${APP_VERSION:-${DW_SERVER_TAG:-2.4.36}} APP_DEBUG: ${APP_DEBUG:-false} LOG_CHANNEL: ${LOG_CHANNEL:-stderr} LOG_LEVEL: ${LOG_LEVEL:-info} diff --git a/docker-compose.small-cluster.yml b/docker-compose.small-cluster.yml index 6d9489fc..c6ecad95 100644 --- a/docker-compose.small-cluster.yml +++ b/docker-compose.small-cluster.yml @@ -12,7 +12,7 @@ x-server-build: &server-build x-server-environment: &server-environment APP_NAME: "Durable Workflow Server" APP_ENV: testing - APP_VERSION: ${APP_VERSION:-2.4.35} + APP_VERSION: ${APP_VERSION:-2.4.36} APP_DEBUG: "false" DW_SERVER_KEY: ${DW_SERVER_KEY:-base64:5Zt4nUhlCm3DD0nLXZJQdHiwPfb56yGo9gNV/g3jYbY=} DB_CONNECTION: ${DW_SMALL_CLUSTER_DB:-mysql} diff --git a/docker-compose.yml b/docker-compose.yml index e5d146d3..404c669c 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -15,7 +15,7 @@ services: DW_SERVER_KEY: "${DW_SERVER_KEY:-}" DW_SERVER_TOPOLOGY_SHAPE: standalone_server DW_SERVER_PROCESS_CLASS: server_http_node - APP_VERSION: "${APP_VERSION:-2.4.35}" + APP_VERSION: "${APP_VERSION:-2.4.36}" APP_DEBUG: "false" DB_CONNECTION: mysql DB_HOST: mysql @@ -62,7 +62,7 @@ services: APP_NAME: "Durable Workflow Server" APP_ENV: local DW_SERVER_KEY: "${DW_SERVER_KEY:-}" - APP_VERSION: "${APP_VERSION:-2.4.35}" + APP_VERSION: "${APP_VERSION:-2.4.36}" APP_DEBUG: "false" DB_CONNECTION: mysql DB_HOST: mysql @@ -124,7 +124,7 @@ services: DW_SERVER_KEY: "${DW_SERVER_KEY:-}" DW_SERVER_TOPOLOGY_SHAPE: standalone_server DW_SERVER_PROCESS_CLASS: worker_node - APP_VERSION: "${APP_VERSION:-2.4.35}" + APP_VERSION: "${APP_VERSION:-2.4.36}" DB_CONNECTION: mysql DB_HOST: mysql DB_PORT: 3306 @@ -178,7 +178,7 @@ services: DW_SERVER_KEY: "${DW_SERVER_KEY:-}" DW_SERVER_TOPOLOGY_SHAPE: standalone_server DW_SERVER_PROCESS_CLASS: scheduler_node - APP_VERSION: "${APP_VERSION:-2.4.35}" + APP_VERSION: "${APP_VERSION:-2.4.36}" DB_CONNECTION: mysql DB_HOST: mysql DB_PORT: 3306 diff --git a/k8s/README.md b/k8s/README.md index 81b88971..8901b15e 100644 --- a/k8s/README.md +++ b/k8s/README.md @@ -13,7 +13,7 @@ The checked-in manifests are synchronized with the repository's stable source release and pin its Docker Hub tag: ```text -durableworkflow/server:2.4.35 +durableworkflow/server:2.4.36 ``` Before production use, patch every workload image to the exact published tag or @@ -21,15 +21,15 @@ digest you intend to run: ```bash kubectl set image -n durable-workflow deploy/durable-workflow-server \ - server=durableworkflow/server:2.4.35 + server=durableworkflow/server:2.4.36 kubectl set image -n durable-workflow deploy/durable-workflow-worker \ - worker=durableworkflow/server:2.4.35 + worker=durableworkflow/server:2.4.36 kubectl set image -n durable-workflow cronjob/durable-workflow-scheduler \ - scheduler=durableworkflow/server:2.4.35 + scheduler=durableworkflow/server:2.4.36 ``` GitHub Container Registry publishes the same release line at -`ghcr.io/durable-workflow/server:2.4.35`. Digest pinning is preferred for strict +`ghcr.io/durable-workflow/server:2.4.36`. Digest pinning is preferred for strict change control. The manifests expect you to provide: diff --git a/k8s/helm/durable-workflow/Chart.yaml b/k8s/helm/durable-workflow/Chart.yaml index 1cd81cca..a17c14fb 100644 --- a/k8s/helm/durable-workflow/Chart.yaml +++ b/k8s/helm/durable-workflow/Chart.yaml @@ -5,11 +5,11 @@ type: application # The chart's own semver version. Bumped on every chart release; treated as # independent of the server image version (appVersion). Breaking-change rules # for this version live in docs/helm-upgrading.md alongside the chart. -version: 0.1.131 +version: 0.1.132 # The immutable Durable Workflow Server identity this chart release packages. # The onboarding default in values.yaml and appVersion are generated from the # checked-in source release record. -appVersion: "2.4.35" +appVersion: "2.4.36" kubeVersion: ">=1.27.0-0" home: https://durable-workflow.github.io/docs/2.0/deployment sources: @@ -30,7 +30,7 @@ annotations: # exact commit that most recently changed the packaged chart. org.opencontainers.image.source: https://github.com/durable-workflow/server dev.durable-workflow.source-revision: "unreleased" - dev.durable-workflow.image-reference: "docker.io/durableworkflow/server:2.4.35" + dev.durable-workflow.image-reference: "docker.io/durableworkflow/server:2.4.36" artifacthub.io/license: MIT artifacthub.io/category: integration-delivery # Free-form changelog for the current chart release shown by Artifact Hub. diff --git a/k8s/helm/durable-workflow/README.md b/k8s/helm/durable-workflow/README.md index f3e74eda..98c9d4f4 100644 --- a/k8s/helm/durable-workflow/README.md +++ b/k8s/helm/durable-workflow/README.md @@ -63,7 +63,7 @@ helm install durable-workflow ./k8s/helm/durable-workflow \ ```yaml image: - tag: "2.4.35" + tag: "2.4.36" # Pin a digest in production: # digest: "sha256:abc123..." # memoPayloadStorage: "raw-json-v1" # Required for a digest or custom image. diff --git a/k8s/helm/durable-workflow/ci/existing-secrets-values.yaml b/k8s/helm/durable-workflow/ci/existing-secrets-values.yaml index b3270fed..eece7e5f 100644 --- a/k8s/helm/durable-workflow/ci/existing-secrets-values.yaml +++ b/k8s/helm/durable-workflow/ci/existing-secrets-values.yaml @@ -1,7 +1,7 @@ # CI fixture: GitOps / externally-managed-secret path. The chart consumes # existing Secrets and renders no Secret resources of its own. image: - tag: "2.4.35" + tag: "2.4.36" externalDatabase: connection: pgsql diff --git a/k8s/helm/durable-workflow/ci/ingress-and-hpa-values.yaml b/k8s/helm/durable-workflow/ci/ingress-and-hpa-values.yaml index a3ad7c9e..fe673a37 100644 --- a/k8s/helm/durable-workflow/ci/ingress-and-hpa-values.yaml +++ b/k8s/helm/durable-workflow/ci/ingress-and-hpa-values.yaml @@ -1,6 +1,6 @@ # CI fixture: ingress + autoscaling enabled. Exercises optional templates. image: - tag: "2.4.35" + tag: "2.4.36" externalDatabase: connection: mysql diff --git a/k8s/helm/durable-workflow/ci/inline-secrets-values.yaml b/k8s/helm/durable-workflow/ci/inline-secrets-values.yaml index e2e15318..cc4ab8d4 100644 --- a/k8s/helm/durable-workflow/ci/inline-secrets-values.yaml +++ b/k8s/helm/durable-workflow/ci/inline-secrets-values.yaml @@ -2,7 +2,7 @@ # chart's render path is exercised end-to-end. Real deployments should use # existingSecret instead. image: - tag: "2.4.35" + tag: "2.4.36" externalDatabase: connection: mysql diff --git a/k8s/helm/durable-workflow/templates/_helpers.tpl b/k8s/helm/durable-workflow/templates/_helpers.tpl index bf25fc44..281ceb16 100644 --- a/k8s/helm/durable-workflow/templates/_helpers.tpl +++ b/k8s/helm/durable-workflow/templates/_helpers.tpl @@ -88,7 +88,7 @@ resolved by an explicit capability declaration or an existing workload marker. {{- define "durable-workflow.memoPayloadStorageForImage" -}} {{- $image := toString . -}} {{- $normalized := regexReplaceAll "^index\\.docker\\.io/" $image "docker.io/" -}} -{{- if eq $normalized "docker.io/durableworkflow/server:2.4.35" -}} +{{- if eq $normalized "docker.io/durableworkflow/server:2.4.36" -}} dual-v1 {{- else if regexMatch "^docker\\.io/durableworkflow/server:2\\.0\\.0-rc\\.[0-9]+$" $normalized -}} {{- $releaseCandidate := atoi (regexFind "[0-9]+$" $normalized) -}} diff --git a/k8s/helm/durable-workflow/values.yaml b/k8s/helm/durable-workflow/values.yaml index 020f6c5f..ca5ea8f8 100644 --- a/k8s/helm/durable-workflow/values.yaml +++ b/k8s/helm/durable-workflow/values.yaml @@ -21,7 +21,7 @@ image: registry: docker.io repository: durableworkflow/server # Generated by scripts/ci/sync-source-release.mjs. Do not edit this default. - tag: "2.4.35" + tag: "2.4.36" # Optional digest pin. When set, takes precedence over tag for change control. # Example: "sha256:abc123..." digest: "" diff --git a/k8s/helm/examples/values-dev.yaml b/k8s/helm/examples/values-dev.yaml index eca0bb5a..af74e930 100644 --- a/k8s/helm/examples/values-dev.yaml +++ b/k8s/helm/examples/values-dev.yaml @@ -3,7 +3,7 @@ # shape in production. image: - tag: "2.4.35" + tag: "2.4.36" externalDatabase: connection: mysql diff --git a/k8s/helm/examples/values-external-secrets-operator.yaml b/k8s/helm/examples/values-external-secrets-operator.yaml index 4902e1c8..e4228fca 100644 --- a/k8s/helm/examples/values-external-secrets-operator.yaml +++ b/k8s/helm/examples/values-external-secrets-operator.yaml @@ -5,7 +5,7 @@ # concern. image: - tag: "2.4.35" + tag: "2.4.36" externalDatabase: connection: pgsql diff --git a/k8s/helm/examples/values-production-existing-secrets.yaml b/k8s/helm/examples/values-production-existing-secrets.yaml index 9a924d83..66221cd4 100644 --- a/k8s/helm/examples/values-production-existing-secrets.yaml +++ b/k8s/helm/examples/values-production-existing-secrets.yaml @@ -10,7 +10,7 @@ image: repository: durable-workflow/server # Pin a digest in production for change-control auditability. digest: "" # e.g. "sha256:abc123..." - tag: "2.4.35" + tag: "2.4.36" externalDatabase: connection: pgsql diff --git a/k8s/migration-job.yaml b/k8s/migration-job.yaml index 6f96417b..34168987 100644 --- a/k8s/migration-job.yaml +++ b/k8s/migration-job.yaml @@ -13,7 +13,7 @@ spec: restartPolicy: OnFailure containers: - name: migrate - image: durableworkflow/server:2.4.35 + image: durableworkflow/server:2.4.36 command: ["server-entrypoint"] args: ["server-bootstrap"] envFrom: diff --git a/k8s/scheduler-cronjob.yaml b/k8s/scheduler-cronjob.yaml index 995f71a5..12e5f8a1 100644 --- a/k8s/scheduler-cronjob.yaml +++ b/k8s/scheduler-cronjob.yaml @@ -24,7 +24,7 @@ spec: restartPolicy: Never containers: - name: scheduler - image: durableworkflow/server:2.4.35 + image: durableworkflow/server:2.4.36 command: ["server-entrypoint"] args: - sh diff --git a/k8s/secret.yaml b/k8s/secret.yaml index 8e13ff6e..ec855e0b 100644 --- a/k8s/secret.yaml +++ b/k8s/secret.yaml @@ -12,7 +12,7 @@ metadata: app.kubernetes.io/name: durable-workflow data: APP_NAME: "Durable Workflow Server" - APP_VERSION: "2.4.35" + APP_VERSION: "2.4.36" APP_ENV: production APP_DEBUG: "false" DB_CONNECTION: mysql diff --git a/k8s/server-deployment.yaml b/k8s/server-deployment.yaml index 69086fe9..a2b53225 100644 --- a/k8s/server-deployment.yaml +++ b/k8s/server-deployment.yaml @@ -23,7 +23,7 @@ spec: spec: containers: - name: server - image: durableworkflow/server:2.4.35 + image: durableworkflow/server:2.4.36 ports: - containerPort: 8080 name: http diff --git a/k8s/worker-deployment.yaml b/k8s/worker-deployment.yaml index 553d5526..66e61d99 100644 --- a/k8s/worker-deployment.yaml +++ b/k8s/worker-deployment.yaml @@ -19,7 +19,7 @@ spec: spec: containers: - name: worker - image: durableworkflow/server:2.4.35 + image: durableworkflow/server:2.4.36 command: ["server-entrypoint"] args: ["php", "artisan", "queue:work", "--sleep=1", "--tries=3", "--max-time=3600"] envFrom: diff --git a/resources/release/source-release.json b/resources/release/source-release.json index 81f562fe..2ede815a 100644 --- a/resources/release/source-release.json +++ b/resources/release/source-release.json @@ -1,9 +1,9 @@ { "schema": "durable-workflow.server.source-release/v1", "server": { - "version": "2.4.35" + "version": "2.4.36" }, "helm_chart": { - "version": "0.1.131" + "version": "0.1.132" } } diff --git a/scripts/k8s-kind-smoke.sh b/scripts/k8s-kind-smoke.sh index 129a3a96..7448af88 100755 --- a/scripts/k8s-kind-smoke.sh +++ b/scripts/k8s-kind-smoke.sh @@ -7,7 +7,7 @@ cluster="${K8S_SMOKE_CLUSTER:-durable-workflow-server-smoke}" image="${K8S_SMOKE_IMAGE:-durableworkflow/server:k8s-smoke}" # Generated by scripts/ci/sync-source-release.mjs so the smoke replaces the # same default shipped by the public manifests. -manifest_image="durableworkflow/server:2.4.35" +manifest_image="durableworkflow/server:2.4.36" kind_node_image="${K8S_SMOKE_KIND_NODE_IMAGE:-kindest/node:v1.29.4}" artifact_dir="${K8S_SMOKE_ARTIFACT_DIR:-/tmp/durable-workflow-k8s-kind-smoke-artifacts}" rendered_dir="${artifact_dir}/rendered-manifests" diff --git a/tests/Feature/WorkflowTaskHeartbeatRaceTest.php b/tests/Feature/WorkflowTaskHeartbeatRaceTest.php new file mode 100644 index 00000000..617d7745 --- /dev/null +++ b/tests/Feature/WorkflowTaskHeartbeatRaceTest.php @@ -0,0 +1,139 @@ +assertIsString($path); + $this->databasePath = $path; + config([ + 'database.default' => 'sqlite', + 'database.connections.sqlite.database' => $path, + 'database.connections.sqlite.busy_timeout' => 100, + 'database.connections.sqlite.journal_mode' => 'WAL', + 'database.connections.sqlite.transaction_mode' => 'IMMEDIATE', + 'cache.default' => 'array', + 'workflows.v2.workflow_task_lease_seconds' => 1, + 'workflows.v2.types.workflows' => ['tests.external-greeting-workflow' => ExternalGreetingWorkflow::class], + ]); + DB::purge('sqlite'); + $this->artisan('migrate:fresh', ['--force' => true])->assertExitCode(0); + Queue::fake(); + WorkflowNamespace::query()->create([ + 'name' => 'default', 'description' => 'Test', 'retention_days' => 30, 'status' => 'active', + ]); + } + + protected function tearDown(): void + { + DB::disconnect('sqlite'); + foreach ([$this->databasePath, $this->databasePath.'-wal', $this->databasePath.'-shm'] as $path) { + if (is_file($path)) { + unlink($path); + } + } + parent::tearDown(); + } + + public function test_heartbeat_cannot_acknowledge_an_old_claim_after_concurrent_reclaim(): void + { + $clock = now()->startOfSecond(); + $this->travelTo($clock); + $this->withHeaders([ + 'X-Namespace' => 'default', 'X-Durable-Workflow-Control-Plane-Version' => '2', + ])->postJson('/api/workflows', [ + 'workflow_type' => 'tests.external-greeting-workflow', 'task_queue' => 'heartbeat-race', 'input' => ['Ada'], + ])->assertCreated(); + foreach (['original', 'replacement'] as $workerId) { + WorkerRegistration::query()->create([ + 'namespace' => 'default', 'worker_id' => $workerId, 'task_queue' => 'heartbeat-race', + 'runtime' => 'php', 'supported_workflow_types' => ['tests.external-greeting-workflow'], + 'capabilities' => [], 'max_concurrent_workflow_tasks' => 1, + 'last_heartbeat_at' => now(), 'status' => 'active', + ]); + } + $claim = $this->withHeaders([ + 'X-Namespace' => 'default', 'X-Durable-Workflow-Protocol-Version' => '1.19', + ])->postJson('/api/worker/workflow-tasks/poll', [ + 'worker_id' => 'original', 'task_queue' => 'heartbeat-race', 'poll_request_id' => 'original-poll', + ])->assertOk()->json('task'); + $this->assertIsArray($claim); + + $children = []; + try { + foreach (['heartbeat' => 'original', 'poll' => 'replacement'] as $operation => $workerId) { + $process = proc_open([ + PHP_BINARY, dirname(__DIR__).'/Fixtures/WorkflowTaskHeartbeatRaceRequest.php', + $this->databasePath, $operation, $claim['task_id'], $workerId, + $clock->toIso8601String(), (string) $claim['workflow_task_attempt'], + ], [0 => ['pipe', 'r'], 1 => ['pipe', 'w'], 2 => ['pipe', 'w']], $pipes); + $this->assertIsResource($process); + stream_set_timeout($pipes[1], 10); + $children[$operation] = ['process' => $process, 'pipes' => $pipes]; + $this->assertSame("ready\n", fgets($pipes[1])); + } + fwrite($children['heartbeat']['pipes'][0], "go\n"); + $this->assertSame("checked\n", fgets($children['heartbeat']['pipes'][1])); + // The guard passed before expiry. Both HTTP kernels now cross + // expiry, letting Native claim while heartbeat renewal is paused. + // No lease, owner or attempt is changed by the test harness. + fwrite($children['poll']['pipes'][0], "go\n"); + $read = [$children['poll']['pipes'][1]]; + $write = $except = null; + stream_select($read, $write, $except, 1); + fwrite($children['heartbeat']['pipes'][0], "go\n"); + + $results = []; + foreach ($children as $operation => $child) { + fclose($child['pipes'][0]); + $output = stream_get_contents($child['pipes'][1]); + $this->assertFalse(stream_get_meta_data($child['pipes'][1])['timed_out'], 'Heartbeat race timed out.'); + $results[$operation] = json_decode($output, true, flags: JSON_THROW_ON_ERROR); + } + $heartbeat = $results['heartbeat']; + $this->assertSame(200, $heartbeat['status'], json_encode($heartbeat)); + $this->assertTrue($heartbeat['body']['renewed']); + $task = WorkflowTask::query()->findOrFail($claim['task_id']); + $this->assertSame($task->lease_owner, $heartbeat['body']['lease_owner'], json_encode([ + 'heartbeat' => array_intersect_key($heartbeat['body'], array_flip([ + 'task_id', 'lease_owner', 'workflow_task_attempt', 'renewed', + ])), + 'replacement_claim' => is_array($results['poll']['body']['task'] ?? null) + ? array_intersect_key($results['poll']['body']['task'], array_flip([ + 'task_id', 'lease_owner', 'workflow_task_attempt', + ])) : null, + ])); + $this->assertSame($task->attempt_count, $heartbeat['body']['workflow_task_attempt']); + $this->assertSame('original', $task->lease_owner); + $this->assertSame($claim['workflow_task_attempt'], $task->attempt_count); + $this->assertSame(200, $results['poll']['status'], json_encode($results['poll'])); + $this->assertNull($results['poll']['body']['task']); + } finally { + foreach ($children as $child) { + if (proc_get_status($child['process'])['running']) { + proc_terminate($child['process']); + } + foreach ($child['pipes'] as $pipe) { + if (is_resource($pipe)) { + fclose($pipe); + } + } + proc_close($child['process']); + } + } + } +} diff --git a/tests/Fixtures/WorkflowTaskHeartbeatRaceRequest.php b/tests/Fixtures/WorkflowTaskHeartbeatRaceRequest.php new file mode 100644 index 00000000..c6084713 --- /dev/null +++ b/tests/Fixtures/WorkflowTaskHeartbeatRaceRequest.php @@ -0,0 +1,94 @@ +make(Kernel::class); +$app->make(Illuminate\Contracts\Console\Kernel::class)->bootstrap(); +config([ + 'app.env' => 'testing', + 'app.key' => 'base64:dGVzdGluZy10ZXN0aW5nLXRlc3RpbmctdGVzdGluZzEyMzQ1Ng==', + 'database.default' => 'sqlite', + 'database.connections.sqlite.database' => $argv[1], + 'database.connections.sqlite.busy_timeout' => 100, + 'database.connections.sqlite.journal_mode' => 'WAL', + 'database.connections.sqlite.transaction_mode' => 'IMMEDIATE', + 'cache.default' => 'array', + 'server.auth.driver' => 'none', + 'server.worker_protocol.version' => '1.19', + 'server.polling.timeout' => 0, + 'server.polling.interval_ms' => 1, + 'server.polling.cache_path' => sys_get_temp_dir().'/dw-heartbeat-race-'.getmypid(), + 'workflows.v2.workflow_task_lease_seconds' => 1, + 'workflows.v2.types.workflows' => ['tests.external-greeting-workflow' => ExternalGreetingWorkflow::class], +]); +Carbon::setTestNow(Carbon::parse($argv[5])); +DB::purge('sqlite'); +DB::connection()->getPdo(); +Queue::fake(); +register_shutdown_function(static function (): void { + (new Filesystem)->deleteDirectory(sys_get_temp_dir().'/dw-heartbeat-race-'.getmypid()); +}); + +if ($argv[2] === 'poll') { + Carbon::setTestNow(Carbon::parse($argv[5])->addSeconds(2)); +} else { + $native = $app->make(WorkflowTaskBridge::class); + // Add only an IPC barrier. Status and renewal execute the real published + // Native methods against separate database connections in these kernels. + $bridge = Mockery::mock(WorkflowTaskBridge::class); + $bridge->shouldReceive('status')->andReturnUsing($native->status(...)); + $paused = false; + $bridge->shouldReceive('heartbeat')->andReturnUsing(static function (string $taskId) use ($native, &$paused, $argv): array { + if (! $paused) { + $paused = true; + Carbon::setTestNow(Carbon::parse($argv[5])->addSeconds(2)); + fwrite(STDOUT, "checked\n"); + if (trim((string) fgets(STDIN)) !== 'go') { + exit(2); + } + } + + return $native->heartbeat($taskId); + }); + $app->instance(WorkflowTaskBridge::class, $bridge); + $app->instance(WorkflowTaskOwnership::class, new WorkflowTaskOwnership($bridge)); +} + +fwrite(STDOUT, "ready\n"); +if (trim((string) fgets(STDIN)) !== 'go') { + exit(2); +} +$path = $argv[2] === 'heartbeat' + ? '/api/worker/workflow-tasks/'.$argv[3].'/heartbeat' + : '/api/worker/workflow-tasks/poll'; +$body = $argv[2] === 'heartbeat' + ? ['lease_owner' => $argv[4], 'workflow_task_attempt' => (int) $argv[6]] + : ['worker_id' => $argv[4], 'task_queue' => 'heartbeat-race', 'poll_request_id' => 'replacement-poll']; +$attempts = []; +for ($attempt = 1; $attempt <= 3; $attempt++) { + $request = Request::create($path, 'POST', server: [ + 'CONTENT_TYPE' => 'application/json', 'HTTP_ACCEPT' => 'application/json', + 'HTTP_X_NAMESPACE' => 'default', 'HTTP_X_DURABLE_WORKFLOW_PROTOCOL_VERSION' => '1.19', + ], content: json_encode($body, JSON_THROW_ON_ERROR)); + $response = $kernel->handle($request); + $responseBody = json_decode($response->getContent(), true, flags: JSON_THROW_ON_ERROR); + $attempts[] = ['status' => $response->getStatusCode(), 'reason' => $responseBody['reason'] ?? null]; + $kernel->terminate($request, $response); + if ($response->getStatusCode() !== 503 || ($responseBody['reason'] ?? null) !== 'backend_lock_pressure' || $attempt === 3) { + break; + } + usleep(max(1, (int) ($responseBody['retry_after_seconds'] ?? 1)) * 1000000); +} +fwrite(STDOUT, json_encode([ + 'status' => $response->getStatusCode(), 'body' => $responseBody, 'attempts' => $attempts, +], JSON_THROW_ON_ERROR)); From 8840d6ddb73df0a78a83e748e38cd3688897d90e Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Thu, 1 Oct 2026 13:04:45 +0000 Subject: [PATCH 2/2] test: require truthful heartbeat when either side wins reclaim --- .../Feature/WorkflowTaskHeartbeatRaceTest.php | 32 +++++++++++++------ 1 file changed, 23 insertions(+), 9 deletions(-) diff --git a/tests/Feature/WorkflowTaskHeartbeatRaceTest.php b/tests/Feature/WorkflowTaskHeartbeatRaceTest.php index 617d7745..dbf85887 100644 --- a/tests/Feature/WorkflowTaskHeartbeatRaceTest.php +++ b/tests/Feature/WorkflowTaskHeartbeatRaceTest.php @@ -105,23 +105,37 @@ public function test_heartbeat_cannot_acknowledge_an_old_claim_after_concurrent_ $results[$operation] = json_decode($output, true, flags: JSON_THROW_ON_ERROR); } $heartbeat = $results['heartbeat']; - $this->assertSame(200, $heartbeat['status'], json_encode($heartbeat)); - $this->assertTrue($heartbeat['body']['renewed']); $task = WorkflowTask::query()->findOrFail($claim['task_id']); - $this->assertSame($task->lease_owner, $heartbeat['body']['lease_owner'], json_encode([ + $trace = json_encode([ + 'heartbeat_status' => $heartbeat['status'], 'heartbeat' => array_intersect_key($heartbeat['body'], array_flip([ - 'task_id', 'lease_owner', 'workflow_task_attempt', 'renewed', + 'task_id', 'lease_owner', 'workflow_task_attempt', 'renewed', 'reason', ])), 'replacement_claim' => is_array($results['poll']['body']['task'] ?? null) ? array_intersect_key($results['poll']['body']['task'], array_flip([ 'task_id', 'lease_owner', 'workflow_task_attempt', ])) : null, - ])); - $this->assertSame($task->attempt_count, $heartbeat['body']['workflow_task_attempt']); - $this->assertSame('original', $task->lease_owner); - $this->assertSame($claim['workflow_task_attempt'], $task->attempt_count); + ]); + $this->assertContains($heartbeat['status'], [200, 409], $trace); $this->assertSame(200, $results['poll']['status'], json_encode($results['poll'])); - $this->assertNull($results['poll']['body']['task']); + if ($heartbeat['status'] === 200) { + $this->assertTrue($heartbeat['body']['renewed']); + $this->assertSame($task->lease_owner, $heartbeat['body']['lease_owner'], $trace); + $this->assertSame($task->attempt_count, $heartbeat['body']['workflow_task_attempt']); + $this->assertSame('original', $task->lease_owner); + $this->assertSame($claim['workflow_task_attempt'], $task->attempt_count); + $this->assertNull($results['poll']['body']['task']); + } else { + // SQLite without a writer lock may let reclaim win and retry + // the heartbeat transaction. Revalidation must then fence it. + $this->assertSame('lease_owner_mismatch', $heartbeat['body']['reason'], $trace); + $this->assertFalse($heartbeat['body']['renewed'] ?? false); + $this->assertSame('replacement', $task->lease_owner); + $this->assertSame($claim['workflow_task_attempt'] + 1, $task->attempt_count); + $this->assertSame($claim['task_id'], $results['poll']['body']['task']['task_id']); + $this->assertSame($task->lease_owner, $results['poll']['body']['task']['lease_owner']); + $this->assertSame($task->attempt_count, $results['poll']['body']['task']['workflow_task_attempt']); + } } finally { foreach ($children as $child) { if (proc_get_status($child['process'])['running']) {