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
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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);
}
Expand All @@ -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]
Expand Down
8 changes: 4 additions & 4 deletions packages/client/src/Client/Internal/Http2/StreamFactory.php
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -31,7 +31,7 @@
*/
final readonly class StreamFactory
{
private StreamCodec $codec;
private Http2\StreamCodec $codec;

/**
* @param positive-int $maxReceiveMessageSize
Expand All @@ -48,7 +48,7 @@ public function __construct(
int $maxReceiveMessageSize,
array $compressors,
) {
$this->codec = new StreamCodec(
$this->codec = new Http2\StreamCodec(
$encoder,
$compressor,
$maxReceiveMessageSize,
Expand Down Expand Up @@ -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}.
Expand Down
50 changes: 50 additions & 0 deletions packages/grpc/src/Internal/Http2/Headers.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
<?php

declare(strict_types=1);

namespace Thesis\Grpc\Internal\Http2;

use Thesis\Google\Rpc\Code;
use Thesis\Grpc\InvokeError;
use Thesis\Grpc\Metadata;

/**
* @internal
* @return array<non-empty-string, list<string>>
*/
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<non-empty-string, list<string>> $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);
}
15 changes: 6 additions & 9 deletions packages/grpc/src/Status/Context.php
Original file line number Diff line number Diff line change
Expand Up @@ -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),
);
}

Expand All @@ -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(
Expand Down
1 change: 1 addition & 0 deletions packages/grpc/src/autoload.files.php
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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();
}

Expand Down
33 changes: 23 additions & 10 deletions packages/server/src/Server/Internal/Http2/ServerRequestHandler.php
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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();

Expand All @@ -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.
Expand All @@ -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.
Expand All @@ -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();

Expand All @@ -154,6 +157,7 @@ public function handleRequest(Request $request): Response
/** @var ServerStream<object, object> $stream */
$stream = $factory->create(
$rpc->handle,
$md,
$request,
$response,
$streamCancellation,
Expand Down Expand Up @@ -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)))),
);
}
}
3 changes: 2 additions & 1 deletion packages/server/src/Server/Internal/Http2/StreamFactory.php
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,7 @@ public function __construct(
*/
public function create(
Handle $handle,
Metadata $md,
Request $request,
Response $response,
Cancellation $cancellation,
Expand All @@ -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,
Expand Down
105 changes: 105 additions & 0 deletions tests/BinaryMetadataTest.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,105 @@
<?php

declare(strict_types=1);

namespace Thesis\Grpc;

use Amp\Cancellation;
use Amp\Http\Client\HttpClientBuilder;
use Amp\Http\Client\Request;
use Echos\Api\V1\EchoRequest;
use Echos\Api\V1\EchoResponse;
use Echos\Api\V1\EchoServiceServer;
use Echos\Api\V1\EchoServiceServerRegistry;
use File\Api\V1\Chunk;
use File\Api\V1\FileInfo;
use File\Api\V1\FileServiceServer;
use File\Api\V1\FileServiceServerRegistry;
use PHPUnit\Framework\Attributes\CoversClass;
use PHPUnit\Framework\TestCase;
use Thesis\Google\Rpc\Code;
use Thesis\Grpc\Client\Internal\AmphpHttpClient;
use Thesis\Grpc\Client\Invoke;
use Thesis\Grpc\Server\CallableStreamInterceptor;
use Thesis\Grpc\Server\Internal\AmphpHttpServer;
use Thesis\Grpc\Server\StreamInfo;

#[CoversClass(AmphpHttpServer::class)]
#[CoversClass(AmphpHttpClient::class)]
final class BinaryMetadataTest extends TestCase
{
private Server $server;

protected function setUp(): void
{
$this->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());
}
}
Loading
Loading