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
6 changes: 4 additions & 2 deletions src/Commands/WorkflowTaskCommand/CompleteCommand.php
Original file line number Diff line number Diff line change
Expand Up @@ -23,16 +23,18 @@ protected function configure(): void
Complete a leased workflow task through the worker protocol. By default the
command emits a single <comment>complete_workflow</comment> command; pass
<comment>--command</comment> for lower-level SDK command payloads.
<comment>--complete-result</comment> accepts a JSON string containing base64 Avro
single-object bytes or a JSON payload envelope produced by an SDK.

<comment>Examples:</comment>

<info>dw workflow-task:complete task-123 1 --lease-owner=cli-worker --complete-result='{"ok":true}'</info>
<info>dw workflow-task:complete task-123 1 --lease-owner=cli-worker</info>
<info>dw workflow-task:complete task-123 1 --command='{"type":"fail_workflow","message":"boom"}' --json</info>
HELP)
->addArgument('task-id', InputArgument::REQUIRED, 'Workflow task ID')
->addArgument('attempt', InputArgument::REQUIRED, 'Workflow task attempt number')
->addOption('lease-owner', null, InputOption::VALUE_REQUIRED, 'Lease owner identity', 'cli')
->addOption('complete-result', null, InputOption::VALUE_OPTIONAL, 'JSON result for a complete_workflow command')
->addOption('complete-result', null, InputOption::VALUE_OPTIONAL, 'JSON Avro payload for a complete_workflow command')
->addOption('command', null, InputOption::VALUE_OPTIONAL | InputOption::VALUE_IS_ARRAY, 'Raw workflow task command JSON')
->addOption('json', null, InputOption::VALUE_NONE, 'Output the command response as JSON');
}
Expand Down
34 changes: 21 additions & 13 deletions src/Support/ServerClient.php
Original file line number Diff line number Diff line change
Expand Up @@ -363,22 +363,26 @@ private function decode(ResponseInterface $response, string $method, string $pat
throw new NetworkException("Server unreachable: {$e->getMessage()}", 0, $e);
}

$body = json_decode($rawContent, true) ?? [];
$decoded = json_decode($rawContent, true);
$body = is_array($decoded) ? $decoded : [];

// Authentication and authorization rejections can be generated by a
// managed namespace gateway before a request reaches Server, so they
// do not carry Server's success-response protocol envelope.
if ($statusCode === 401 || $statusCode === 403) {
throw $this->httpException($statusCode, $body, $path);
}

$body = $this->normalizePayload($method, $path, $body, $response, $statusCode);

if ($statusCode >= 400) {
if ($statusCode !== 401 && $statusCode !== 403) {
try {
$body = $this->normalizePayload($method, $path, $body, $response, $statusCode);
} catch (\RuntimeException) {
// Valid envelopes add structured diagnostics, but an
// absent or invalid envelope must not mask an HTTP failure.
}
}

throw $this->httpException($statusCode, $body, $path);
}

return $body;
return $this->normalizePayload($method, $path, $body, $response, $statusCode);
}

private function requestUrl(string $path): string
Expand All @@ -391,11 +395,15 @@ private function requestUrl(string $path): string
*/
private function httpException(int $statusCode, array $body, string $path): ServerHttpException
{
$message = $body['message']
?? $body['error']
?? $this->firstValidationMessage($body)
?? (isset($body['rejection_reason']) ? 'Rejected: '.$body['rejection_reason'] : null)
?? "HTTP {$statusCode}";
$reason = $body['rejection_reason'] ?? null;
$message = is_string($reason) && $reason !== '' ? 'Rejected: '.$reason : "HTTP {$statusCode}";

foreach ([$body['message'] ?? null, $body['error'] ?? null, $this->firstValidationMessage($body)] as $candidate) {
if (is_string($candidate) && $candidate !== '') {
$message = $candidate;
break;
}
}

if ($statusCode === 403 && $this->isManagedNamespaceRuntimeUrl()) {
$message .= self::isWorkerProtocolPath($path)
Expand Down
4 changes: 3 additions & 1 deletion tests/Integration/ServerSmokeTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -145,7 +145,9 @@ public function test_cli_control_and_worker_plane_smoke_against_running_server()
(string) $polledTask['task']['task_id'],
(string) $polledTask['task']['workflow_task_attempt'],
'--lease-owner='.$workerId,
'--complete-result={"status":"completed-by-cli-smoke"}',
// Encoded by the published PHP SDK's AvroPayloadCodec from
// {"status":"completed-by-cli-smoke"} using the fixed Value schema.
'--complete-result="wwHioz3/VYAiNw4CDHN0YXR1cwosY29tcGxldGVkLWJ5LWNsaS1zbW9rZQA="',
'--json',
], $namespace, $this->workerToken);

Expand Down
67 changes: 67 additions & 0 deletions tests/Support/ServerClientHttpFailureTest.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,67 @@
<?php

declare(strict_types=1);

namespace Tests\Support;

use DurableWorkflow\Cli\Support\ExitCode;
use DurableWorkflow\Cli\Support\ServerClient;
use DurableWorkflow\Cli\Support\ServerHttpException;
use PHPUnit\Framework\Attributes\DataProvider;
use PHPUnit\Framework\TestCase;
use Symfony\Component\HttpClient\MockHttpClient;
use Symfony\Component\HttpClient\Response\MockResponse;

final class ServerClientHttpFailureTest extends TestCase
{
/**
* @param array<string, mixed>|null $body
*/
#[DataProvider('errorResponses')]
public function test_http_failures_keep_their_status_and_diagnostics(
int $status,
string $content,
string $message,
?array $body,
int $exitCode,
): void {
foreach (['/workflows', '/worker/register'] as $path) {
$client = new ServerClient(
baseUrl: 'http://example.test',
namespace: 'default',
http: new MockHttpClient(new MockResponse($content, ['http_code' => $status])),
);

try {
$client->post($path);
self::fail('An HTTP failure must not return a successful response.');
} catch (ServerHttpException $exception) {
self::assertSame($status, $exception->statusCode);
self::assertSame($status, $exception->getCode());
self::assertSame('Server error: '.$message, $exception->getMessage());
self::assertSame($body, $exception->body);
self::assertSame($exitCode, $exception->exitCode());
self::assertSame($body['reason'] ?? null, $exception->reason());
self::assertSame($body['validation_errors'] ?? $body['errors'] ?? null, $exception->validationErrors());
}
}
}

public static function errorResponses(): iterable
{
yield 'Laravel storage failure' => [500, '{"message":"Server Error"}', 'Server Error', ['message' => 'Server Error'], ExitCode::SERVER];
yield 'HTML proxy failure' => [502, '<html><body>Bad Gateway</body></html>', 'HTTP 502', null, ExitCode::SERVER];
yield 'empty proxy failure' => [503, '', 'HTTP 503', null, ExitCode::SERVER];
yield 'JSON scalar failure' => [500, '"Server Error"', 'HTTP 500', null, ExitCode::SERVER];
yield 'JSON null failure' => [500, 'null', 'HTTP 500', null, ExitCode::SERVER];
yield 'malformed JSON failure' => [502, '{"message":', 'HTTP 502', null, ExitCode::SERVER];
yield 'structured proxy error' => [502, '{"error":{"code":"upstream_unavailable"}}', 'HTTP 502', ['error' => ['code' => 'upstream_unavailable']], ExitCode::SERVER];
yield 'authentication' => [401, '{"message":"Unauthenticated."}', 'Unauthenticated.', ['message' => 'Unauthenticated.'], ExitCode::AUTH];
yield 'authorization' => [403, '{"message":"Forbidden."}', 'Forbidden.', ['message' => 'Forbidden.'], ExitCode::AUTH];
yield 'not found' => [404, '{"message":"Workflow not found.","reason":"instance_not_found"}', 'Workflow not found.', ['message' => 'Workflow not found.', 'reason' => 'instance_not_found'], ExitCode::NOT_FOUND];
yield 'validation' => [422, '{"errors":{"input":["The input is required."]}}', 'The input is required.', ['errors' => ['input' => ['The input is required.']]], ExitCode::INVALID];
yield 'rejection' => [409, '{"rejection_reason":"workflow_already_running"}', 'Rejected: workflow_already_running', ['rejection_reason' => 'workflow_already_running'], ExitCode::INVALID];
yield 'retryable backend' => [503, '{"message":"A required backend is temporarily unavailable.","reason":"backend_unavailable","retryable":true}', 'A required backend is temporarily unavailable.', ['message' => 'A required backend is temporarily unavailable.', 'reason' => 'backend_unavailable', 'retryable' => true], ExitCode::SERVER];
yield 'broken error contract' => [500, '{"message":"Server Error","control_plane":{"operation":"start"}}', 'Server Error', ['message' => 'Server Error', 'control_plane' => ['operation' => 'start']], ExitCode::SERVER];
}
}
8 changes: 5 additions & 3 deletions tests/Support/ServerClientTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
use DurableWorkflow\Cli\Support\ControlPlaneRequestContract;
use DurableWorkflow\Cli\Support\CompatibilityException;
use DurableWorkflow\Cli\Support\ServerClient;
use DurableWorkflow\Cli\Support\ServerHttpException;
use PHPUnit\Framework\TestCase;
use Symfony\Component\HttpClient\MockHttpClient;
use Symfony\Component\HttpClient\Response\MockResponse;
Expand Down Expand Up @@ -1144,7 +1145,7 @@ public function test_it_keeps_control_plane_error_responses_usable_without_succe
$client->post('/workflows/wf-123/signal/advance');
}

public function test_it_rejects_control_plane_command_error_responses_without_the_shared_contract(): void
public function test_it_preserves_control_plane_command_errors_without_the_shared_contract(): void
{
$response = new MockResponse(json_encode([
'message' => 'Workflow not found.',
Expand All @@ -1162,8 +1163,9 @@ public function test_it_rejects_control_plane_command_error_responses_without_th
http: new MockHttpClient($response, 'http://example.test'),
);

$this->expectException(\RuntimeException::class);
$this->expectExceptionMessage('missing the shared control-plane contract');
$this->expectException(ServerHttpException::class);
$this->expectExceptionCode(404);
$this->expectExceptionMessage('Server error: Workflow not found.');

$client->post('/workflows/wf-123/signal/advance');
}
Expand Down
Loading