diff --git a/packages/client/src/Client/Internal/Http2/ConcurrentClientStream.php b/packages/client/src/Client/Internal/Http2/ConcurrentClientStream.php index 7bd4e99..c1e5777 100644 --- a/packages/client/src/Client/Internal/Http2/ConcurrentClientStream.php +++ b/packages/client/src/Client/Internal/Http2/ConcurrentClientStream.php @@ -14,6 +14,7 @@ use Thesis\Grpc\Client\Internal\CancellationError; use Thesis\Grpc\ClientStream; use Thesis\Grpc\Exception\ClientStreamIsClosed; +use Thesis\Grpc\Internal\Http2; use Thesis\Grpc\InvokeError; use Thesis\Grpc\Metadata; @@ -75,7 +76,7 @@ public function receive(): object public function headers(): Metadata { try { - return new Metadata($this->response->getHeaders()); + return Http2\decodeMetadata($this->response->getHeaders()); } catch (CancelledException $e) { throw CancellationError::from($e); } @@ -84,7 +85,7 @@ public function headers(): Metadata #[\Override] public function trailers(Cancellation $cancellation = new NullCancellation()): Metadata { - return new Metadata($this->response->getTrailers()->await($cancellation)->getHeaders()); + return Http2\decodeMetadata($this->response->getTrailers()->await($cancellation)->getHeaders()); } #[\Override] diff --git a/packages/client/src/Client/Internal/Http2/StreamFactory.php b/packages/client/src/Client/Internal/Http2/StreamFactory.php index cb53595..d72f066 100644 --- a/packages/client/src/Client/Internal/Http2/StreamFactory.php +++ b/packages/client/src/Client/Internal/Http2/StreamFactory.php @@ -21,7 +21,7 @@ use Thesis\Grpc\Compression\CompressionUnavailable; use Thesis\Grpc\Compression\Compressor; use Thesis\Grpc\Encoding\Encoder; -use Thesis\Grpc\Internal\Http2\StreamCodec; +use Thesis\Grpc\Internal\Http2; use Thesis\Grpc\InvokeError; use Thesis\Grpc\Metadata; use function Amp\async; @@ -31,7 +31,7 @@ */ final readonly class StreamFactory { - private StreamCodec $codec; + private Http2\StreamCodec $codec; /** * @param positive-int $maxReceiveMessageSize @@ -48,7 +48,7 @@ public function __construct( int $maxReceiveMessageSize, array $compressors, ) { - $this->codec = new StreamCodec( + $this->codec = new Http2\StreamCodec( $encoder, $compressor, $maxReceiveMessageSize, @@ -82,7 +82,7 @@ public function create( ), ); $request->setProtocolVersions(['2']); - $request->setHeaders($md->kv); + $request->setHeaders(Http2\encodeMetadata($md)); $request->setTransferTimeout($this->transferTimeout); $request->setInactivityTimeout($this->inactivityTimeout); // gRPC limits the size of a single message, not the stream, see {@see StreamCodec}. diff --git a/packages/grpc/src/Internal/Http2/Headers.php b/packages/grpc/src/Internal/Http2/Headers.php new file mode 100644 index 0000000..0b80070 --- /dev/null +++ b/packages/grpc/src/Internal/Http2/Headers.php @@ -0,0 +1,50 @@ +> + */ +function encodeMetadata(Metadata $md): array +{ + $headers = $md->kv; + + foreach ($headers as $key => $values) { + if (str_ends_with($key, '-bin')) { + $headers[$key] = array_map( + static fn(string $value): string => rtrim(base64_encode($value), '='), + $values, + ); + } + } + + return $headers; +} + +/** + * @internal + * @param array> $headers + * @throws InvokeError + */ +function decodeMetadata(array $headers): Metadata +{ + foreach ($headers as $key => $values) { + if (str_ends_with($key, '-bin')) { + $headers[$key] = array_map( + static fn(string $value): string => ($decoded = base64_decode($value, true)) !== false + ? $decoded + : throw new InvokeError(Code::INTERNAL, "Malformed binary metadata in header \"{$key}\""), + explode(',', implode(',', $values)), + ); + } + } + + return new Metadata($headers); +} diff --git a/packages/grpc/src/Status/Context.php b/packages/grpc/src/Status/Context.php index 69f3089..208ab04 100644 --- a/packages/grpc/src/Status/Context.php +++ b/packages/grpc/src/Status/Context.php @@ -47,7 +47,7 @@ function serializeContext(Context $context, Encoder $protobuf): Metadata\Status return new Metadata\Status( $context->code, $context->message, - base64_encode($protobuf->encode($status)), + $protobuf->encode($status), ); } @@ -61,14 +61,11 @@ function deserializeContext(Metadata $md, Decoder $protobuf): Context $details = []; - if (($bin = $status->details) !== null) { - $decoded = base64_decode($bin, true); - if ($decoded !== false) { - $details = array_map( - static fn(Protobuf\Any $detail) => Protobuf\decodeAny($detail, $protobuf), - $protobuf->decode($decoded, Rpc\Status::class)->details, - ); - } + if ($status->details !== null) { + $details = array_map( + static fn(Protobuf\Any $detail) => Protobuf\decodeAny($detail, $protobuf), + $protobuf->decode($status->details, Rpc\Status::class)->details, + ); } return new Context( diff --git a/packages/grpc/src/autoload.files.php b/packages/grpc/src/autoload.files.php index 748f4db..f7aa0e1 100644 --- a/packages/grpc/src/autoload.files.php +++ b/packages/grpc/src/autoload.files.php @@ -2,6 +2,7 @@ declare(strict_types=1); +require_once __DIR__ . '/Internal/Http2/Headers.php'; require_once __DIR__ . '/Internal/Protocol/Frame.php'; require_once __DIR__ . '/Metadata/Timeout.php'; require_once __DIR__ . '/Metadata/ContentType.php'; diff --git a/packages/server/src/Server/Internal/Http2/ConcurrentServerStream.php b/packages/server/src/Server/Internal/Http2/ConcurrentServerStream.php index bcdb409..65a9266 100644 --- a/packages/server/src/Server/Internal/Http2/ConcurrentServerStream.php +++ b/packages/server/src/Server/Internal/Http2/ConcurrentServerStream.php @@ -7,6 +7,7 @@ use Amp\DeferredFuture; use Amp\Pipeline; use Thesis\Grpc\Exception\ServerStreamIsClosed; +use Thesis\Grpc\Internal\Http2; use Thesis\Grpc\Metadata; use Thesis\Grpc\ServerStream; @@ -63,7 +64,7 @@ public function close(): void return; } - $this->trailersFuture->complete($this->trailers->kv); + $this->trailersFuture->complete(Http2\encodeMetadata($this->trailers)); $this->send->complete(); } diff --git a/packages/server/src/Server/Internal/Http2/ServerRequestHandler.php b/packages/server/src/Server/Internal/Http2/ServerRequestHandler.php index f87a6cb..bba2a1b 100644 --- a/packages/server/src/Server/Internal/Http2/ServerRequestHandler.php +++ b/packages/server/src/Server/Internal/Http2/ServerRequestHandler.php @@ -15,6 +15,8 @@ use Amp\Http\Server\Trailers; use Amp\TimeoutCancellation; use Thesis\Google\Rpc; +use Thesis\Grpc\Internal\Http2; +use Thesis\Grpc\InvokeError; use Thesis\Grpc\Metadata; use Thesis\Grpc\Server\Internal\StreamHandleInterceptor; use Thesis\Grpc\Server\Internal\StreamInterceptorComposer; @@ -90,7 +92,14 @@ public function services(): array #[\Override] public function handleRequest(Request $request): Response { - $md = new Metadata($request->getHeaders()); + try { + $md = Http2\decodeMetadata($request->getHeaders()); + } catch (InvokeError $e) { + return self::trailersOnly( + new Metadata()->withKey(new Metadata\ContentType()), + new Metadata\Status($e->statusCode, $e->statusMessage), + ); + } $headers = new Metadata(); @@ -102,7 +111,7 @@ public function handleRequest(Request $request): Response $headers = $headers->withKey($contentType ?? new Metadata\ContentType()); if ($contentType === null) { - return new Response(status: HttpStatus::UNSUPPORTED_MEDIA_TYPE, headers: $headers->kv); + return new Response(status: HttpStatus::UNSUPPORTED_MEDIA_TYPE, headers: Http2\encodeMetadata($headers)); } // For "grpc-encoding" header we follow the same approach as for "Content-Type": we should not specify "IDENTITY" by default for the response to avoid sending an unnecessary header. @@ -114,13 +123,7 @@ public function handleRequest(Request $request): Response $encoder = $this->encoderFactory->encoder($contentType->encoding ?? Metadata\ContentType::GRPC_DEFAULT_ENCODING); $rpc = $this->router->route($request); } catch (UnimplementedException $e) { - return new Response( - status: HttpStatus::OK, - headers: $headers->kv, - trailers: new Trailers(Future::complete( - new Metadata()->withKey(new Metadata\Status(Rpc\Code::UNIMPLEMENTED, $e->getMessage()))->kv, - )), - ); + return self::trailersOnly($headers, new Metadata\Status(Rpc\Code::UNIMPLEMENTED, $e->getMessage())); } // The "grpc-encoding" header should only be sent when a protobuf message is expected to be returned. @@ -137,7 +140,7 @@ public function handleRequest(Request $request): Response $this->maxReceiveMessageSize, ); - $response = new Response(status: HttpStatus::OK, headers: $headers->kv); + $response = new Response(status: HttpStatus::OK, headers: Http2\encodeMetadata($headers)); $cancellation = new DeferredCancellation(); @@ -154,6 +157,7 @@ public function handleRequest(Request $request): Response /** @var ServerStream $stream */ $stream = $factory->create( $rpc->handle, + $md, $request, $response, $streamCancellation, @@ -242,4 +246,13 @@ public function stop(Cancellation $cancellation): void Future\awaitAll($futures, $cancellation); } + + private static function trailersOnly(Metadata $headers, Metadata\Status $status): Response + { + return new Response( + status: HttpStatus::OK, + headers: Http2\encodeMetadata($headers), + trailers: new Trailers(Future::complete(Http2\encodeMetadata(new Metadata()->withKey($status)))), + ); + } } diff --git a/packages/server/src/Server/Internal/Http2/StreamFactory.php b/packages/server/src/Server/Internal/Http2/StreamFactory.php index 624fb61..d3a865e 100644 --- a/packages/server/src/Server/Internal/Http2/StreamFactory.php +++ b/packages/server/src/Server/Internal/Http2/StreamFactory.php @@ -46,6 +46,7 @@ public function __construct( */ public function create( Handle $handle, + Metadata $md, Request $request, Response $response, Cancellation $cancellation, @@ -62,7 +63,7 @@ public function create( )); return new ConcurrentServerStream( - new Metadata($request->getHeaders()), + $md, $this->codec->decode($request->getBody(), $handle->type, $cancellation), $send, $trailers, diff --git a/tests/BinaryMetadataTest.php b/tests/BinaryMetadataTest.php new file mode 100644 index 0000000..aacad03 --- /dev/null +++ b/tests/BinaryMetadataTest.php @@ -0,0 +1,105 @@ +server = new Server\Builder() + ->withServices( + new FileServiceServerRegistry(new readonly class implements FileServiceServer { + #[\Override] + public function upload(Server\ClientStreamChannel $stream, Metadata $md, Cancellation $cancellation): FileInfo + { + iterator_to_array($stream); + + return new FileInfo(); + } + }), + new EchoServiceServerRegistry(new readonly class implements EchoServiceServer { + #[\Override] + public function echo(EchoRequest $request, Metadata $md, Cancellation $cancellation): EchoResponse + { + return new EchoResponse(); + } + }), + ) + ->withStreamInterceptors(new CallableStreamInterceptor(static function ( + ServerStream $stream, + StreamInfo $info, + Metadata $md, + Cancellation $cancellation, + callable $next, + ): void { + $stream->trailers->join(new Metadata()->with('x-test-bin', ...$md['x-test-bin'])); + $next($stream, $info, $md, $cancellation); + })) + ->build(); + + $this->server->start(); + } + + protected function tearDown(): void + { + $this->server->stop(); + } + + public function testBinaryMetadata(): void + { + $stream = new Client\Builder()->build()->createStream( + new Invoke('/file.api.v1.FileService/Upload', FileInfo::class, RpcType::ClientStream), + new Metadata()->with('x-test-bin', "\xab\xab\xab", "\x00\xff"), + ); + + $stream->send(new Chunk('content')); + $stream->close(); + $stream->receive(); + + self::assertSame(["\xab\xab\xab", "\x00\xff"], $stream->trailers()['x-test-bin']); + } + + public function testMalformedBinaryMetadata(): void + { + $request = new Request('http://127.0.0.1:50051/echos.api.v1.EchoService/Echo', 'POST'); + $request->setProtocolVersions(['2']); + $request->setHeaders([ + 'content-type' => 'application/grpc', + 'x-test-bin' => 'q6u!', + ]); + + $response = HttpClientBuilder::buildDefault()->request($request); + $response->getBody()->buffer(); + + self::assertSame([ + 'grpc-status' => [(string) Code::INTERNAL->value], + 'grpc-message' => ['Malformed binary metadata in header "x-test-bin"'], + ], $response->getTrailers()->await()->getHeaders()); + } +} diff --git a/tests/Internal/Http2/HeadersTest.php b/tests/Internal/Http2/HeadersTest.php new file mode 100644 index 0000000..0e48ec4 --- /dev/null +++ b/tests/Internal/Http2/HeadersTest.php @@ -0,0 +1,99 @@ +> $headers + */ + #[DataProvider('provideEncodeCases')] + public function testEncode(Metadata $md, array $headers): void + { + self::assertSame($headers, encodeMetadata($md)); + } + + /** + * @return iterable>}> + */ + public static function provideEncodeCases(): iterable + { + yield 'binary without padding' => [ + new Metadata(['x-test-bin' => "\xab\xab\xab"]), + ['x-test-bin' => ['q6ur']], + ]; + + yield 'padding is stripped' => [ + new Metadata(['x-test-bin' => ["\xab", "\xab\xab"]]), + ['x-test-bin' => ['qw', 'q6s']], + ]; + + yield 'text is untouched' => [ + new Metadata(['x-test' => "\xab", 'x-test-bin-suffix' => 'q6ur']), + ['x-test' => ["\xab"], 'x-test-bin-suffix' => ['q6ur']], + ]; + } + + /** + * @param array> $headers + */ + #[DataProvider('provideDecodeCases')] + public function testDecode(array $headers, Metadata $md): void + { + self::assertEquals($md, decodeMetadata($headers)); + } + + /** + * @return iterable>, Metadata}> + */ + public static function provideDecodeCases(): iterable + { + yield 'without padding' => [ + ['x-test-bin' => ['q6ur', 'q6s']], + new Metadata(['x-test-bin' => ["\xab\xab\xab", "\xab\xab"]]), + ]; + + yield 'with padding' => [ + ['x-test-bin' => ['q6s=', 'qw==']], + new Metadata(['x-test-bin' => ["\xab\xab", "\xab"]]), + ]; + + yield 'comma separated values' => [ + ['x-test-bin' => ['q6ur,q6s=', 'qw']], + new Metadata(['x-test-bin' => ["\xab\xab\xab", "\xab\xab", "\xab"]]), + ]; + + yield 'text is untouched' => [ + ['x-test' => ['q6ur,q6s']], + new Metadata(['x-test' => ['q6ur,q6s']]), + ]; + } + + #[DataProvider('provideDecodeMalformedCases')] + public function testDecodeMalformed(string $value): void + { + $this->expectExceptionObject(new InvokeError(Code::INTERNAL, 'Malformed binary metadata in header "x-test-bin"')); + decodeMetadata(['x-test-bin' => [$value]]); + } + + /** + * @return iterable + */ + public static function provideDecodeMalformedCases(): iterable + { + yield 'invalid character' => ['q6u!']; + yield 'excessive padding' => ['q6ur==']; + yield 'truncated' => ['q']; + } +}