From 312ed98b8217fb9fb35f0ede1b76b838dbbcd05e Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Fri, 11 Sep 2026 15:23:41 +0000 Subject: [PATCH] Prevent poll capacity backpressure from starving other task kinds --- CHANGELOG.md | 10 ++ composer.json | 2 +- docs/quickstart-contract.json | 4 +- src/Worker.php | 7 ++ tests/DependencyBoundaryTest.php | 2 +- tests/WorkerPollFairnessTest.php | 175 +++++++++++++++++++++++++++++++ 6 files changed, 196 insertions(+), 4 deletions(-) create mode 100644 tests/WorkerPollFairnessTest.php diff --git a/CHANGELOG.md b/CHANGELOG.md index 908d1d1..e97cf91 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,16 @@ project follows [Semantic Versioning](https://semver.org/). ## [Unreleased] +## [2.0.10] - 2026-09-11 + +### Fixed + +- A repeated long-poll capacity refusal no longer prevents the managed worker + from polling other task kinds. Explicit empty-task refusals honor backoff, + heartbeats and shutdown before yielding to activity or query polling. +- Ambiguous transport, database and storage-admission failures retain their + existing retry and poll-identity behavior; lease fencing is unchanged. + ## [2.0.9] - 2026-09-08 ### Fixed diff --git a/composer.json b/composer.json index 3788d95..da02d4e 100644 --- a/composer.json +++ b/composer.json @@ -74,7 +74,7 @@ } }, "durable-workflow": { - "product-train": "2.0.9", + "product-train": "2.0.10", "supported-server-versions": "2.1.0", "worker-protocol-version": "1.19", "control-plane-version": "2", diff --git a/docs/quickstart-contract.json b/docs/quickstart-contract.json index 86ca7ee..267c661 100644 --- a/docs/quickstart-contract.json +++ b/docs/quickstart-contract.json @@ -3,8 +3,8 @@ "schema_version": 2, "package": { "name": "durable-workflow/sdk", - "published_version": "2.0.9", - "composer_requirement": "2.0.9", + "published_version": "2.0.10", + "composer_requirement": "2.0.10", "onboarding_requirement": "^2.0" }, "runtime_targets": { diff --git a/src/Worker.php b/src/Worker.php index e192416..dbd2b7d 100644 --- a/src/Worker.php +++ b/src/Worker.php @@ -636,6 +636,13 @@ private function pollWithRetry(string $taskKind, \Closure $poll): ?array 'exception' => $exception, ], 'warning'); $this->waitForTransientRetry($delaySeconds); + if ($exception->status === 429 + && $exception->reason === 'long_poll_capacity_exhausted' + && ($exception->details['poll_status'] ?? null) === 'long_poll_capacity_exhausted') { + // This explicit empty-task refusal acquired no lease. Give + // other task kinds a turn instead of retrying only this one. + return $this->shutdownRequested ? null : $exception->details; + } } } diff --git a/tests/DependencyBoundaryTest.php b/tests/DependencyBoundaryTest.php index 9422c6b..9e48a24 100644 --- a/tests/DependencyBoundaryTest.php +++ b/tests/DependencyBoundaryTest.php @@ -44,7 +44,7 @@ public function testStableMetadataDeclaresExactQualifiedArtifacts(): void $metadata = $this->manifest()['extra']['durable-workflow']; $quickstart = $this->quickstartContract(); - self::assertSame('2.0.9', $metadata['product-train']); + self::assertSame('2.0.10', $metadata['product-train']); self::assertSame('2.1.0', $metadata['supported-server-versions']); self::assertSame('1.19', $metadata['worker-protocol-version']); self::assertTrue($metadata['durable-selection']); diff --git a/tests/WorkerPollFairnessTest.php b/tests/WorkerPollFairnessTest.php new file mode 100644 index 0000000..461a317 --- /dev/null +++ b/tests/WorkerPollFairnessTest.php @@ -0,0 +1,175 @@ + $blockedKinds */ + #[DataProvider('blockedKinds')] + public function testWaitCapacityRefusalsGiveEveryTaskKindATurn(array $blockedKinds): void + { + $now = 0.0; + $worker = null; + $polls = []; + $transport = new FakeTransport(handler: static function ( + string $method, + string $uri, + array $headers, + ?array $body, + ) use ($blockedKinds, &$polls): array { + preg_match('#/worker/(workflow|activity|query)-tasks/poll$#', $uri, $match); + self::assertNotEmpty($match); + $kind = $match[1]; + $polls[] = ['kind' => $kind, 'id' => $body['poll_request_id']]; + if (in_array($kind, $blockedKinds, true)) { + throw self::capacityRefusal($kind); + } + + return ['task' => null, 'poll_status' => 'empty']; + }); + $worker = new Worker( + new Client('https://server.example', transport: $transport), + 'orders', + clock: static function () use (&$now): float { + return $now; + }, + sleeper: static function (int $microseconds) use (&$now, &$worker): void { + $now += $microseconds / 1_000_000; + if ($now >= 20) { + $worker->requestShutdown(); + } + }, + ); + + self::assertFalse($worker->tick(5)); + self::assertFalse($worker->tick(5)); + + self::assertSame(['workflow', 'activity', 'query', 'workflow', 'activity', 'query'], array_column($polls, 'kind')); + self::assertCount(6, array_unique(array_column($polls, 'id'))); + self::assertGreaterThanOrEqual(count($blockedKinds) * 2, $now); + self::assertEqualsWithDelta(count($blockedKinds) * 2, $now, 0.000_01); + } + + /** @return iterable}> */ + public static function blockedKinds(): iterable + { + yield 'workflow' => [['workflow']]; + yield 'activity' => [['activity']]; + yield 'query' => [['query']]; + yield 'all' => [['workflow', 'activity', 'query']]; + } + + public function testManagedWorkerCompletesActivityWhileWorkflowWaitsAreFull(): void + { + $now = 0.0; + $worker = null; + $calls = 0; + $completions = 0; + $heartbeats = 0; + $transport = new FakeTransport(handler: static function ( + string $method, + string $uri, + array $headers, + ?array $body, + ) use (&$heartbeats, &$completions): array { + if (str_ends_with($uri, '/register')) { + return ['registered' => true, 'heartbeat_interval_seconds' => 1]; + } + if (str_ends_with($uri, '/heartbeat')) { + ++$heartbeats; + return ['acknowledged' => true, 'heartbeat_interval_seconds' => 1]; + } + if (str_ends_with($uri, '/workflow-tasks/poll')) { + throw self::capacityRefusal('workflow'); + } + if (str_ends_with($uri, '/activity-tasks/poll')) { + return ['poll_status' => 'leased', 'task' => [ + 'task_id' => 'activity-1', 'activity_attempt_id' => 'attempt-1', + 'lease_owner' => 'worker-1', 'activity_type' => 'orders.charge', 'payload_codec' => 'avro', + ]]; + } + if (str_ends_with($uri, '/activity-tasks/activity-1/complete')) { + ++$completions; + self::assertSame('worker-1', $body['lease_owner']); + self::assertSame('attempt-1', $body['activity_attempt_id']); + return ['completed' => true]; + } + if (str_ends_with($uri, '/query-tasks/poll')) { + return ['task' => null, 'poll_status' => 'stopped', 'reason' => 'worker_stopped']; + } + if ($method === 'DELETE' && str_ends_with($uri, '/registrations/worker-1')) { + return ['deregistered' => true]; + } + self::fail("Unexpected request: {$method} {$uri}"); + }); + $worker = new Worker( + new Client('https://server.example', transport: $transport), + 'orders', + workerId: 'worker-1', + clock: static function () use (&$now): float { + return $now; + }, + sleeper: static function (int $microseconds) use (&$now, &$worker): void { + $now += $microseconds / 1_000_000; + if ($now >= 20) { + $worker->requestShutdown(); + } + }, + ); + $worker->registerActivity('orders.charge', static function (ActivityContext $context) use (&$calls): string { + ++$calls; + return 'charged'; + }); + $worker->run(5); + + self::assertSame(1, $calls); + self::assertSame(1, $completions); + self::assertGreaterThanOrEqual(1, $heartbeats); + self::assertLessThan(2, $now); + } + + public function testShutdownDuringCapacityBackoffDoesNotPollAnotherKind(): void + { + $now = 0.0; + $worker = null; + $transport = new FakeTransport([self::capacityRefusal('workflow')]); + $worker = new Worker( + new Client('https://server.example', transport: $transport), + 'orders', + clock: static function () use (&$now): float { + return $now; + }, + sleeper: static function (int $microseconds) use (&$now, &$worker): void { + $now += $microseconds / 1_000_000; + $worker->requestShutdown(); + }, + ); + + self::assertFalse($worker->tick(5)); + self::assertCount(1, $transport->requests); + self::assertStringEndsWith('/workflow-tasks/poll', $transport->requests[0]['uri']); + self::assertGreaterThan(0, $now); + self::assertLessThan(1, $now); + } + + private static function capacityRefusal(string $kind): TransportException + { + $response = [ + 'task' => null, 'poll_status' => 'long_poll_capacity_exhausted', + 'reason' => 'long_poll_capacity_exhausted', 'retryable' => true, + 'retry_after_seconds' => 1, 'task_kind' => $kind.'_task', + ]; + + return TransportException::fromResponse(429, $response, json_encode($response, JSON_THROW_ON_ERROR)); + } +}