From 745c9e41b9ded2345f6196c37ad24ec666a67c34 Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Tue, 6 Oct 2026 21:40:34 +0000 Subject: [PATCH] Add bounded workflow observations for operator clients --- composer.json | 2 +- docs/quickstart-contract.json | 4 +- docs/sdk-reference.md | 21 +++++++++ src/Client.php | 44 ++++++++++++++++++ tests/ControlPlaneParityTest.php | 79 ++++++++++++++++++++++++++++++++ 5 files changed, 147 insertions(+), 3 deletions(-) diff --git a/composer.json b/composer.json index f342a3c..14fb314 100644 --- a/composer.json +++ b/composer.json @@ -82,7 +82,7 @@ } }, "durable-workflow": { - "product-train": "2.2.2", + "product-train": "2.2.3", "supported-server-versions": "2.5.0", "worker-protocol-version": "1.19", "control-plane-version": "2", diff --git a/docs/quickstart-contract.json b/docs/quickstart-contract.json index b711c99..718e96a 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.2.2", - "composer_requirement": "2.2.2", + "published_version": "2.2.3", + "composer_requirement": "2.2.3", "onboarding_requirement": "^2.2" }, "runtime_targets": { diff --git a/docs/sdk-reference.md b/docs/sdk-reference.md index ae4c151..081834a 100644 --- a/docs/sdk-reference.md +++ b/docs/sdk-reference.md @@ -37,6 +37,27 @@ additive Server route. Older Servers return their original HTTP 404 without an unfiltered fallback. The SDK rejects malformed selections and encoded query values larger than 4096 bytes before sending them. +## Bounded workflow observations + +With an operator credential, `workflowObservation('order-1', 'run-1')` reads +scalar run metadata, current waits, related run statuses and recent failure +references. Omit the run ID to observe the instance's stored current pointer. +Related runs have their own statuses. Unknown or pruned evidence stays explicit. +This read does not establish authority to issue commands. + +Application context is opt in. Pass up to 20 distinct search attribute keys as +the third argument. The response preserves their metadata without interpreting +it as a payload reference. The observation includes 200 history events by +default. `historyPageSize` selects up to 1000. Pass `history.next_page_token` as +`historyPageToken` to continue inside the original sequence boundary. Supporting +failure references also provide a cursor for direct inspection. These cursors +belong to this observation route. Application values and external payload +references in the history stay encoded. + +This additive Server route requires a runtime that supports bounded observations. +An older runtime's refusal is returned directly. Use `workflowDiagnostics()` +explicitly when full diagnostics are needed. + ## Plain PHP quickstart Create an empty Composer project and install the current published package: diff --git a/src/Client.php b/src/Client.php index 7b52032..f8827fe 100644 --- a/src/Client.php +++ b/src/Client.php @@ -517,6 +517,50 @@ public function workflowDiagnostics(string $workflowId, ?string $runId = null): ); } + /** + * Read bounded run metadata without loading full diagnostics or application payloads. + * + * @param array $searchAttributeKeys A list of explicit application context keys, at most 20. + * @return array + */ + public function workflowObservation( + string $workflowId, + ?string $runId = null, + array $searchAttributeKeys = [], + ?int $historyPageSize = null, + ?string $historyPageToken = null, + ): array { + if (!array_is_list($searchAttributeKeys) || count($searchAttributeKeys) > 20) { + throw new InvalidArgumentException('Search attribute keys must be a list of at most 20 distinct strings.'); + } + $seen = []; + foreach ($searchAttributeKeys as $key) { + if (!is_string($key) || $key === '' || strlen($key) > 255 || isset($seen[$key])) { + throw new InvalidArgumentException('Search attribute keys must be distinct nonempty strings of at most 255 bytes.'); + } + $seen[$key] = true; + } + if ($historyPageSize !== null && ($historyPageSize < 1 || $historyPageSize > 1000)) { + throw new InvalidArgumentException('Observation history page size must be from 1 to 1000.'); + } + if ($historyPageToken !== null && ($historyPageToken === '' || strlen($historyPageToken) > 4096)) { + throw new InvalidArgumentException('Observation history cursor must be nonempty and at most 4096 bytes.'); + } + $path = $this->workflowOperationPath($workflowId, $runId, 'observation'); + $parameters = $this->withoutNulls([ + 'history_page_size' => $historyPageSize, + 'history_page_token' => $historyPageToken, + ]); + if ($searchAttributeKeys !== []) { + $parameters = ['search_attribute_keys' => $searchAttributeKeys] + $parameters; + } + if ($parameters !== []) { + $path .= '?'.http_build_query($parameters, '', '&', PHP_QUERY_RFC3986); + } + + return $this->control('GET', $path); + } + /** * @param list $arguments * @return array diff --git a/tests/ControlPlaneParityTest.php b/tests/ControlPlaneParityTest.php index 396c653..3331107 100644 --- a/tests/ControlPlaneParityTest.php +++ b/tests/ControlPlaneParityTest.php @@ -12,6 +12,7 @@ use DurableWorkflow\Tests\Support\FakeTransport; use DurableWorkflow\Transport\Psr18Transport; use GuzzleHttp\Psr7\Response; +use InvalidArgumentException; use PHPUnit\Framework\Attributes\DataProvider; use PHPUnit\Framework\TestCase; use Psr\Http\Client\ClientInterface; @@ -473,6 +474,7 @@ public function testEveryAddedSurfacePreservesNonSuccessEvidence(): void $operations = [ 'workflow visibility' => static fn (Client $client) => $client->listWorkflows(), 'workflow diagnostics' => static fn (Client $client) => $client->workflowDiagnostics('order-1', 'run-1'), + 'workflow observation' => static fn (Client $client) => $client->workflowObservation('order-1', 'run-1'), 'workflow activities' => static fn (Client $client) => $client->workflowActivities('order-1', 'run-1'), 'workflow redrive' => static fn (Client $client) => $client->redriveWorkflow('order-1', 'run-1'), 'search attributes' => static fn (Client $client) => $client->listSearchAttributes(), @@ -531,6 +533,83 @@ public function testBoundedOperatorDashboardUsesTheAuthenticatedControlPlaneRout self::assertSame('2', $transport->requests[0]['headers']['X-Durable-Workflow-Control-Plane-Version']); } + public function testWorkflowObservationKeepsRunNamespaceAndContextSelectionExplicit(): void + { + $response = ['run_id' => 'run/b', 'read_mode' => 'bounded', 'current_waits' => ['waits' => []], + 'history' => ['events' => [['payload' => ['arguments' => ['codec' => 'avro', 'external_payload' => ['opaque' => true]]]]]], + 'search_attributes' => ['Order /+' => ['codec' => 'avro', 'external_payload' => ['opaque' => true]]]]; + $transport = new FakeTransport([$response, $response]); + $client = new Client('https://server.example', transport: $transport, namespace: 'ops', controlToken: 'operator-test-token'); + + self::assertSame($response, $client->workflowObservation('order/a')); + self::assertSame($response, $client->workflowObservation('order/a', 'run/b', ['Order /+', 'Account&Id'], 100, 'cursor/+')); + self::assertSame('https://server.example/api/workflows/order%2Fa/observation', $transport->requests[0]['uri']); + self::assertSame('https://server.example/api/workflows/order%2Fa/runs/run%2Fb/observation?search_attribute_keys%5B0%5D=Order%20%2F%2B&search_attribute_keys%5B1%5D=Account%26Id&history_page_size=100&history_page_token=cursor%2F%2B', $transport->requests[1]['uri']); + foreach ($transport->requests as $request) { + self::assertSame('GET', $request['method']); + self::assertNull($request['body']); + self::assertSame('ops', $request['headers']['X-Namespace']); + self::assertSame('Bearer operator-test-token', $request['headers']['Authorization']); + self::assertSame('2', $request['headers']['X-Durable-Workflow-Control-Plane-Version']); + } + self::assertCount(2, $transport->requests); + } + + /** @param array $keys */ + #[DataProvider('invalidObservationKeys')] + public function testInvalidObservationSelectionIsRefusedBeforeSending(array $keys): void + { + $transport = new FakeTransport(); + $client = new Client('https://server.example', transport: $transport); + $this->expectException(InvalidArgumentException::class); + try { + $client->workflowObservation('order-1', 'run-1', $keys); + } finally { + self::assertSame([], $transport->requests); + } + } + + /** @return array}> */ + public static function invalidObservationKeys(): array + { + return [ + 'map' => [['key' => 'OrderId']], + 'empty key' => [['']], + 'non-string key' => [[1]], + 'oversized key' => [[str_repeat('x', 256)]], + 'duplicate key' => [['OrderId', 'OrderId']], + 'too many keys' => [array_map('strval', range(1, 21))], + ]; + } + + public function testUnsupportedObservationPreservesRefusalWithoutAFullRead(): void + { + $refusal = new ServerException('Not found.', 404, 'not_found'); + $transport = new FakeTransport([$refusal]); + $client = new Client('https://server.example', transport: $transport); + try { + $client->workflowObservation('order-1', 'run-1'); + self::fail('An unsupported observation must preserve the Server refusal.'); + } catch (ServerException $exception) { + self::assertSame($refusal, $exception); + } + self::assertCount(1, $transport->requests); + } + + public function testInvalidObservationHistoryWindowCannotSendARequest(): void + { + $transport = new FakeTransport(); + $client = new Client('https://server.example', transport: $transport); + foreach ([[0, null], [1001, null], [null, ''], [null, str_repeat('x', 4097)]] as [$size, $token]) { + try { + $client->workflowObservation('order-1', 'run-1', historyPageSize: $size, historyPageToken: $token); + self::fail('Invalid history window was accepted.'); + } catch (InvalidArgumentException) { + self::assertSame([], $transport->requests); + } + } + } + public function testUnsupportedBoundedDashboardPreservesTheServerErrorWithoutAnotherRead(): void { $refusal = new ServerException('Bounded dashboards are unavailable.', 404, 'not_found');