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 @@ -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",
Expand Down
4 changes: 2 additions & 2 deletions docs/quickstart-contract.json
Original file line number Diff line number Diff line change
Expand Up @@ -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": {
Expand Down
21 changes: 21 additions & 0 deletions docs/sdk-reference.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
44 changes: 44 additions & 0 deletions src/Client.php
Original file line number Diff line number Diff line change
Expand Up @@ -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<mixed> $searchAttributeKeys A list of explicit application context keys, at most 20.
* @return array<string, mixed>
*/
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<mixed> $arguments
* @return array<string, mixed>
Expand Down
79 changes: 79 additions & 0 deletions tests/ControlPlaneParityTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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(),
Expand Down Expand Up @@ -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<mixed> $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<string, array{array<mixed>}> */
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');
Expand Down
Loading