diff --git a/README.md b/README.md index cd8c6c1..813abcc 100644 --- a/README.md +++ b/README.md @@ -180,14 +180,43 @@ try { } ``` -## Callback Signatures +## Callbacks + +Verify the signature, then decode the CloudEvent. `CloudEvent` is a plain +envelope; `CallbackEvent::decode()` turns its data into the callback for its +type — `JobStart`, `JobLog`, `JobArtifact`, `JobExit`, `JobComplete`, or +`DeploymentResponse`. A callback that can fail carries an `Error` with a stable +`code` to branch on and a `message` to show: ```php +use OpenRuntimes\Orchestrator\Callback\CloudEvent; +use OpenRuntimes\Orchestrator\Callback\JobArtifact; +use OpenRuntimes\Orchestrator\Callback\JobExit; use OpenRuntimes\Orchestrator\Callback\Signature; +use OpenRuntimes\Orchestrator\Enum\CallbackEvent; +use OpenRuntimes\Orchestrator\Enum\ErrorCode; + +if (! Signature::verifyEvent($rawBody, $headers['x-signature-256'] ?? '', $secret)) { + return; +} -$valid = Signature::verifyEvent($rawBody, $headers['x-signature-256'] ?? '', $secret); +$event = CloudEvent::decode( + \json_decode($rawBody, true), + fn (string $type, array $data) => CallbackEvent::from($type)->decode($data), +); + +match (true) { + $event->data instanceof JobArtifact && $event->data->error !== null + => $log->error("{$event->data->artifactId}: {$event->data->error->message}"), + $event->data instanceof JobExit && $event->data->error?->code === ErrorCode::JobOom + => $log->error('Out of memory'), + default => null, +}; ``` +`CloudEvent::fromArray()` keeps the data as the raw array when you would rather +read it yourself. + ## Development ```sh diff --git a/rector.php b/rector.php index 3d60019..ce67b0c 100644 --- a/rector.php +++ b/rector.php @@ -4,6 +4,7 @@ use Rector\CodeQuality\Rector\Catch_\ThrowWithPreviousExceptionRector; use Rector\Config\RectorConfig; +use Rector\PHPUnit\CodeQuality\Rector\MethodCall\AssertEmptyNullableObjectToAssertInstanceofRector; use Rector\Strict\Rector\Empty_\DisallowedEmptyRuleFixerRector; return RectorConfig::configure() @@ -22,6 +23,8 @@ phpunitCodeQuality: true ) ->withSkip([ + // assertNull on a ?Object return states the contract; assertNotInstanceOf obscures it. + AssertEmptyNullableObjectToAssertInstanceofRector::class, ThrowWithPreviousExceptionRector::class, DisallowedEmptyRuleFixerRector::class, ]); diff --git a/src/Callback/Callback.php b/src/Callback/Callback.php new file mode 100644 index 0000000..493e0e7 --- /dev/null +++ b/src/Callback/Callback.php @@ -0,0 +1,17 @@ + $data + */ + public static function fromArray(array $data): static; +} diff --git a/src/Callback/CloudEvent.php b/src/Callback/CloudEvent.php index c5f7706..8f7c93a 100644 --- a/src/Callback/CloudEvent.php +++ b/src/Callback/CloudEvent.php @@ -9,10 +9,16 @@ use Exception; use OpenRuntimes\Orchestrator\Exception\ClientException; +/** + * A CloudEvents 1.0 envelope. It carries data without knowing what the data + * means: decode() hands the type and raw data to whatever does. + * + * @template T + */ final readonly class CloudEvent { /** - * @param array $data + * @param T $data */ public function __construct( public string $specVersion, @@ -22,13 +28,26 @@ public function __construct( public string $id, public DateTimeInterface $time, public string $dataContentType, - public array $data, + public mixed $data, ) {} /** * @param array $payload + * @return self> */ public static function fromArray(array $payload): self + { + return self::decode($payload, static fn (string $type, array $data): array => $data); + } + + /** + * @template U + * + * @param array $payload + * @param callable(string, array): U $decode + * @return self + */ + public static function decode(array $payload, callable $decode): self { $data = $payload['data'] ?? []; if (! isset($payload['time']) || ! \is_string($payload['time']) || $payload['time'] === '') { @@ -41,15 +60,17 @@ public static function fromArray(array $payload): self throw new ClientException('Invalid CloudEvent: malformed time.', previous: $e); } + $type = (string) ($payload['type'] ?? ''); + return new self( specVersion: (string) ($payload['specversion'] ?? ''), - type: (string) ($payload['type'] ?? ''), + type: $type, source: (string) ($payload['source'] ?? ''), subject: (string) ($payload['subject'] ?? ''), id: (string) ($payload['id'] ?? ''), time: $time, dataContentType: (string) ($payload['datacontenttype'] ?? ''), - data: \is_array($data) ? $data : [], + data: $decode($type, \is_array($data) ? $data : []), ); } } diff --git a/src/Callback/DeploymentResponse.php b/src/Callback/DeploymentResponse.php new file mode 100644 index 0000000..80e8e7e --- /dev/null +++ b/src/Callback/DeploymentResponse.php @@ -0,0 +1,59 @@ +>|null */ + public ?array $requestHeaders, + public bool $requestHeadersTruncated, + public ?float $durationSeconds, + public ?int $statusCode, + public ?string $body, + public ?string $bodyEncoding, + public bool $bodyTruncated, + public ?Error $error, + ) {} + + public static function fromArray(array $data): static + { + $headers = $data['requestHeaders'] ?? null; + if ($headers !== null && ! \is_array($headers)) { + throw new ClientException('Invalid deployment response: requestHeaders must be an object.'); + } + + /** @var array>|null $headers */ + return new self( + deploymentId: Data::string($data, 'deploymentId', 'deployment response'), + invocationId: Data::string($data, 'invocationId', 'deployment response'), + requestMethod: Data::string($data, 'requestMethod', 'deployment response'), + requestPath: Data::string($data, 'requestPath', 'deployment response'), + requestPathTruncated: Data::bool($data, 'requestPathTruncated', 'deployment response'), + requestHeaders: $headers, + requestHeadersTruncated: Data::bool($data, 'requestHeadersTruncated', 'deployment response'), + durationSeconds: Data::optionalFloat($data, 'durationSeconds', 'deployment response'), + statusCode: Data::optionalInt($data, 'statusCode', 'deployment response'), + body: Data::optionalString($data, 'body', 'deployment response'), + bodyEncoding: Data::optionalString($data, 'bodyEncoding', 'deployment response'), + bodyTruncated: Data::bool($data, 'bodyTruncated', 'deployment response'), + error: Error::fromData($data), + ); + } +} diff --git a/src/Callback/Error.php b/src/Callback/Error.php new file mode 100644 index 0000000..f4545f6 --- /dev/null +++ b/src/Callback/Error.php @@ -0,0 +1,41 @@ + $data + */ + public static function fromData(array $data): ?self + { + if (! \array_key_exists('error', $data)) { + return null; + } + if (! \is_array($data['error'])) { + throw new ClientException('Invalid callback error: must be an object.'); + } + + /** @var ErrorCode $code */ + $code = Data::enum($data['error'], 'code', ErrorCode::class, 'callback error'); + + return new self($code, Data::string($data['error'], 'message', 'callback error')); + } +} diff --git a/src/Callback/JobArtifact.php b/src/Callback/JobArtifact.php new file mode 100644 index 0000000..b0f18fd --- /dev/null +++ b/src/Callback/JobArtifact.php @@ -0,0 +1,45 @@ + */ + public array $meta, + ) {} + + public static function fromArray(array $data): static + { + return new self( + jobId: Data::string($data, 'jobId', 'job artifact'), + artifactId: Data::string($data, 'artifactId', 'job artifact'), + artifactType: Data::string($data, 'artifactType', 'job artifact'), + status: Data::string($data, 'status', 'job artifact'), + content: $data['content'] ?? null, + durationSeconds: Data::optionalFloat($data, 'durationSeconds', 'job artifact'), + format: Data::optionalString($data, 'format', 'job artifact'), + compression: Data::optionalString($data, 'compression', 'job artifact'), + error: Error::fromData($data), + meta: Data::stringMap($data, 'meta', 'job artifact'), + ); + } +} diff --git a/src/Callback/JobComplete.php b/src/Callback/JobComplete.php new file mode 100644 index 0000000..10d8a49 --- /dev/null +++ b/src/Callback/JobComplete.php @@ -0,0 +1,28 @@ + */ + public array $meta, + ) {} + + public static function fromArray(array $data): static + { + return new self( + jobId: Data::string($data, 'jobId', 'job complete'), + meta: Data::stringMap($data, 'meta', 'job complete'), + ); + } +} diff --git a/src/Callback/JobExit.php b/src/Callback/JobExit.php new file mode 100644 index 0000000..9fb8d3f --- /dev/null +++ b/src/Callback/JobExit.php @@ -0,0 +1,39 @@ + */ + public array $meta, + ) {} + + public static function fromArray(array $data): static + { + return new self( + jobId: Data::string($data, 'jobId', 'job exit'), + exitCode: Data::int($data, 'exitCode', 'job exit'), + reason: Data::optionalString($data, 'reason', 'job exit'), + image: Data::string($data, 'image', 'job exit'), + durationSeconds: Data::optionalFloat($data, 'durationSeconds', 'job exit'), + error: Error::fromData($data), + meta: Data::stringMap($data, 'meta', 'job exit'), + ); + } +} diff --git a/src/Callback/JobLog.php b/src/Callback/JobLog.php new file mode 100644 index 0000000..1b5f0eb --- /dev/null +++ b/src/Callback/JobLog.php @@ -0,0 +1,32 @@ + */ + public array $lines, + public string $stream, + /** @var array */ + public array $meta, + ) {} + + public static function fromArray(array $data): static + { + return new self( + jobId: Data::string($data, 'jobId', 'job log'), + lines: Data::strings($data, 'lines', 'job log'), + stream: Data::string($data, 'stream', 'job log'), + meta: Data::stringMap($data, 'meta', 'job log'), + ); + } +} diff --git a/src/Callback/JobStart.php b/src/Callback/JobStart.php new file mode 100644 index 0000000..e3e927b --- /dev/null +++ b/src/Callback/JobStart.php @@ -0,0 +1,27 @@ + */ + public array $meta, + ) {} + + public static function fromArray(array $data): static + { + return new self( + jobId: Data::string($data, 'jobId', 'job start'), + meta: Data::stringMap($data, 'meta', 'job start'), + ); + } +} diff --git a/src/Enum/CallbackEvent.php b/src/Enum/CallbackEvent.php index 4535c11..eb7dac6 100644 --- a/src/Enum/CallbackEvent.php +++ b/src/Enum/CallbackEvent.php @@ -4,6 +4,14 @@ namespace OpenRuntimes\Orchestrator\Enum; +use OpenRuntimes\Orchestrator\Callback\Callback; +use OpenRuntimes\Orchestrator\Callback\DeploymentResponse; +use OpenRuntimes\Orchestrator\Callback\JobArtifact; +use OpenRuntimes\Orchestrator\Callback\JobComplete; +use OpenRuntimes\Orchestrator\Callback\JobExit; +use OpenRuntimes\Orchestrator\Callback\JobLog; +use OpenRuntimes\Orchestrator\Callback\JobStart; + enum CallbackEvent: string { case Start = 'orchestrator.job.start'; @@ -12,4 +20,23 @@ enum CallbackEvent: string case Exit = 'orchestrator.job.exit'; case Complete = 'orchestrator.job.complete'; case DeploymentResponse = 'orchestrator.deployment.response'; + + /** + * The typed callback of an event of this kind. Pass to CloudEvent::decode(): + * + * CloudEvent::decode($raw, fn (string $type, array $data) => CallbackEvent::from($type)->decode($data)) + * + * @param array $data + */ + public function decode(array $data): Callback + { + return match ($this) { + self::Start => JobStart::fromArray($data), + self::Artifact => JobArtifact::fromArray($data), + self::Log => JobLog::fromArray($data), + self::Exit => JobExit::fromArray($data), + self::Complete => JobComplete::fromArray($data), + self::DeploymentResponse => DeploymentResponse::fromArray($data), + }; + } } diff --git a/src/Enum/ErrorCode.php b/src/Enum/ErrorCode.php new file mode 100644 index 0000000..a2e7555 --- /dev/null +++ b/src/Enum/ErrorCode.php @@ -0,0 +1,47 @@ + $data + */ + public static function bool(array $data, string $key, string $context, bool $default = false): bool + { + $value = $data[$key] ?? $default; + if (! \is_bool($value)) { + throw new ClientException("Invalid {$context}: {$key} must be a boolean."); + } + + return $value; + } + /** * @param array $data */ diff --git a/tests/Callback/CallbackTest.php b/tests/Callback/CallbackTest.php new file mode 100644 index 0000000..3fb6680 --- /dev/null +++ b/tests/Callback/CallbackTest.php @@ -0,0 +1,124 @@ + $data + * @return CloudEvent + */ + private function decode(string $type, array $data): CloudEvent + { + return CloudEvent::decode( + ['specversion' => '1.0', 'type' => $type, 'source' => 'orchestrator/service', 'subject' => 'job-1', 'id' => 'job-1-1', 'time' => '2026-01-15T10:30:00Z', 'data' => $data], + static fn (string $type, array $data): Callback => CallbackEvent::from($type)->decode($data), + ); + } + + /** + * @return iterable, class-string}> + */ + public static function events(): iterable + { + yield 'start' => ['orchestrator.job.start', ['jobId' => 'job-1', 'meta' => []], JobStart::class]; + yield 'log' => ['orchestrator.job.log', ['jobId' => 'job-1', 'lines' => ['a', 'b'], 'stream' => 'stdout', 'meta' => []], JobLog::class]; + yield 'artifact' => ['orchestrator.job.artifact', ['jobId' => 'job-1', 'artifactId' => 'source', 'artifactType' => 'download', 'status' => 'success', 'durationSeconds' => 0.4, 'meta' => []], JobArtifact::class]; + yield 'exit' => ['orchestrator.job.exit', ['jobId' => 'job-1', 'exitCode' => 0, 'image' => 'alpine', 'durationSeconds' => 3, 'meta' => []], JobExit::class]; + yield 'complete' => ['orchestrator.job.complete', ['jobId' => 'job-1', 'meta' => []], JobComplete::class]; + yield 'response' => ['orchestrator.deployment.response', ['deploymentId' => 'dep', 'invocationId' => 'inv', 'requestMethod' => 'POST', 'requestPath' => '/', 'statusCode' => 200, 'body' => 'ok', 'bodyTruncated' => false], DeploymentResponse::class]; + } + + /** + * @param array $data + * @param class-string $payload + */ + #[DataProvider('events')] + public function test_decodes_each_event_to_its_payload(string $type, array $data, string $payload): void + { + $event = $this->decode($type, $data); + + $this->assertSame($type, $event->type); + $this->assertInstanceOf($payload, $event->data); + } + + public function test_artifact_error(): void + { + $event = $this->decode('orchestrator.job.artifact', [ + 'jobId' => 'job-1', 'artifactId' => 'extract', 'artifactType' => 'unarchive', 'status' => 'failed', 'durationSeconds' => 0.01, + 'error' => ['code' => 'archive_unknown_format', 'message' => 'Unrecognized archive format for source.tar.gz'], + 'meta' => ['deploymentId' => 'dep'], + ]); + + $artifact = $event->data; + $this->assertInstanceOf(JobArtifact::class, $artifact); + $this->assertSame('failed', $artifact->status); + $this->assertInstanceOf(Error::class, $artifact->error); + $this->assertSame(ErrorCode::ArchiveUnknownFormat, $artifact->error->code); + $this->assertSame('Unrecognized archive format for source.tar.gz', $artifact->error->message); + $this->assertNull($artifact->format); + $this->assertSame(['deploymentId' => 'dep'], $artifact->meta); + } + + public function test_exit_before_worker_ran(): void + { + $event = $this->decode('orchestrator.job.exit', [ + 'jobId' => 'job-1', 'exitCode' => -1, 'reason' => 'init container failed', 'image' => 'alpine', 'durationSeconds' => 0, + 'error' => ['code' => 'job_failed', 'message' => 'Job failed before it could start'], 'meta' => [], + ]); + + $exit = $event->data; + $this->assertInstanceOf(JobExit::class, $exit); + $this->assertSame(-1, $exit->exitCode); + $this->assertSame('init container failed', $exit->reason); + $this->assertSame(ErrorCode::JobFailed, $exit->error?->code); + } + + public function test_response_that_never_reached_a_replica(): void + { + $event = $this->decode('orchestrator.deployment.response', [ + 'deploymentId' => 'dep', 'invocationId' => 'inv', 'requestMethod' => 'GET', 'requestPath' => '/health', + 'requestHeaders' => ['X-Trace' => ['abc']], + 'error' => ['code' => 'deployment_no_capacity', 'message' => 'The deployment had no capacity ready in time to serve the request'], + ]); + + $response = $event->data; + $this->assertInstanceOf(DeploymentResponse::class, $response); + $this->assertNull($response->statusCode); + $this->assertNull($response->body); + $this->assertSame(['X-Trace' => ['abc']], $response->requestHeaders); + $this->assertSame(ErrorCode::DeploymentNoCapacity, $response->error?->code); + } + + public function test_unknown_event_type_is_rejected(): void + { + $this->expectException(ValueError::class); + + $this->decode('orchestrator.job.teleport', []); + } + + public function test_missing_field_is_rejected(): void + { + $this->expectException(ClientException::class); + + $this->decode('orchestrator.job.exit', ['jobId' => 'job-1']); + } +} diff --git a/tests/Callback/ErrorTest.php b/tests/Callback/ErrorTest.php new file mode 100644 index 0000000..5244c62 --- /dev/null +++ b/tests/Callback/ErrorTest.php @@ -0,0 +1,48 @@ + 'failed', 'error' => ['code' => 'archive_unknown_format', 'message' => 'Unrecognized archive format for source.tar.gz']]); + + $this->assertInstanceOf(Error::class, $error); + $this->assertSame(ErrorCode::ArchiveUnknownFormat, $error->code); + $this->assertSame('Unrecognized archive format for source.tar.gz', $error->message); + } + + public function test_success_has_no_error(): void + { + $this->assertNull(Error::fromData(['status' => 'success'])); + } + + /** + * @return iterable + */ + public static function malformedErrors(): iterable + { + yield 'null' => [null]; + yield 'bare string' => ['job_oom']; + yield 'no code' => [['message' => 'no code']]; + yield 'unknown code' => [['code' => 'job_teleported', 'message' => 'Job teleported']]; + yield 'no message' => [['code' => 'job_oom']]; + } + + #[DataProvider('malformedErrors')] + public function test_rejects_malformed_error(mixed $error): void + { + $this->expectException(ClientException::class); + + Error::fromData(['error' => $error]); + } +}