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
26 changes: 26 additions & 0 deletions packages/client/src/Client/Builder.php
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,9 @@ final class Builder

private ?Compressor $compressor = null;

/** @var list<Compressor> */
private array $compressors = [];

private ?DelegateHttpClient $httpclient = null;

/** @var list<UnaryInterceptor> */
Expand Down Expand Up @@ -95,6 +98,20 @@ public function withCompression(Compressor $compressor): self
return $builder;
}

/**
* @no-named-arguments
*/
public function withCompressors(Compressor ...$compressors): self
{
$builder = clone $this;
$builder->compressors = [
...$builder->compressors,
...$compressors,
];

return $builder;
}

public function withHttpClient(DelegateHttpClient $httpclient): self
{
$builder = clone $this;
Expand Down Expand Up @@ -266,6 +283,10 @@ public function build(): Client

$encoder = $this->encoder ?? ProtobufEncoder::default();
$compressor = $this->compressor ?? IdentityCompressor::Compressor;
$compressors = [
...$this->compressors,
IdentityCompressor::Compressor,
];
$protobuf = $this->protobuf ?? Decoder\Builder::buildDefault();
$loadBalancerFactory = $this->loadBalancerFactory ?? new LoadBalancer\PickFirstFactory();
$tlsContext = $this->credentials?->createContext();
Expand All @@ -283,6 +304,10 @@ public function build(): Client
$controlMetadata = new Internal\AppendControlMetadataInterceptor(
$encoder->name(),
$compressor->name(),
array_values(array_unique(array_map(
static fn(Compressor $compressor) => $compressor->name(),
[$compressor, ...$compressors],
))),
);

// Control metadata sits innermost (closest to the transport) so every user
Expand Down Expand Up @@ -336,6 +361,7 @@ public function build(): Client
inactivityTimeout: $inactivityTimeout,
encoder: $encoder,
compressor: $compressor,
compressors: $compressors,
),
),
),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,10 +21,12 @@
/**
* @param non-empty-string $encoding
* @param non-empty-string $compression
* @param list<non-empty-string> $acceptEncoding
*/
public function __construct(
private string $encoding,
private string $compression,
private array $acceptEncoding,
) {}

#[\Override]
Expand Down Expand Up @@ -54,6 +56,7 @@ private function decorate(Metadata $md): Metadata
->withKey(new Metadata\ContentType($this->encoding))
->withKey(Metadata\UserAgent::Key)
->withKey(new Metadata\ContentEncoding($this->compression))
->withKey(new Metadata\AcceptEncoding($this->acceptEncoding))
->with('TE', 'trailers');
}
}
21 changes: 20 additions & 1 deletion packages/client/src/Client/Internal/Http2/StreamFactory.php
Original file line number Diff line number Diff line change
Expand Up @@ -14,12 +14,15 @@
use Amp\Http\Client\StreamedContent;
use Amp\NullCancellation;
use Amp\Pipeline;
use Thesis\Google\Rpc\Code;
use Thesis\Grpc\Client\Address;
use Thesis\Grpc\Client\Invoke;
use Thesis\Grpc\ClientStream;
use Thesis\Grpc\Compression\CompressionUnavailable;
use Thesis\Grpc\Compression\Compressor;
use Thesis\Grpc\Encoding\Encoder;
use Thesis\Grpc\Internal\Http2\StreamCodec;
use Thesis\Grpc\InvokeError;
use Thesis\Grpc\Metadata;
use function Amp\async;

Expand All @@ -30,6 +33,9 @@
{
private StreamCodec $codec;

/**
* @param list<Compressor> $compressors
*/
public function __construct(
private DelegateHttpClient $http,
private UriFactory $uri,
Expand All @@ -38,10 +44,12 @@ public function __construct(
private float $inactivityTimeout,
Encoder $encoder,
Compressor $compressor,
array $compressors,
) {
$this->codec = new StreamCodec(
$encoder,
$compressor,
$compressors,
);
}

Expand Down Expand Up @@ -95,7 +103,18 @@ public function create(
return new ConcurrentClientStream(
responseFuture: $response,
send: $send,
decode: fn(Response $response) => $this->codec->decode($response->getBody(), $invoke->output, $cancellation),
decode: function (Response $response) use ($invoke, $cancellation): Pipeline\ConcurrentIterator {
try {
return $this->codec->decode(
$response->getBody(),
$invoke->output,
$cancellation,
Metadata\parseContentEncoding(new Metadata($response->getHeaders()))->encoding ?? Metadata\ContentEncoding::GRPC_DEFAULT_COMPRESSION,
);
} catch (CompressionUnavailable $e) {
throw new InvokeError(Code::INTERNAL, $e->getMessage(), previous: $e);
}
},
errors: $this->errors,
complete: $deferred->getFuture(),
);
Expand Down
26 changes: 24 additions & 2 deletions packages/grpc/src/Internal/Http2/StreamCodec.php
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
use Amp\CancelledException;
use Amp\Pipeline;
use Revolt\EventLoop;
use Thesis\Grpc\Compression\CompressionUnavailable;
use Thesis\Grpc\Compression\Compressor;
use Thesis\Grpc\Encoding\Encoder;
use Thesis\Grpc\Internal\Protocol;
Expand All @@ -18,10 +19,24 @@
*/
final readonly class StreamCodec
{
/** @var array<non-empty-string, Compressor> */
private array $compressors;

/**
* @param list<Compressor> $compressors to decompress messages with, selected by the peer's "grpc-encoding"
*/
public function __construct(
private Encoder $encoder,
private Compressor $compressor,
) {}
array $compressors = [],
) {
$compressors = [$compressor, ...$compressors];

$this->compressors = array_combine(
array_map(static fn(Compressor $compressor) => $compressor->name(), $compressors),
$compressors,
);
}

/**
* @template T of object
Expand Down Expand Up @@ -78,21 +93,28 @@ public function encode(
/**
* @template T of object
* @param class-string<T> $type
* @param ?non-empty-string $encoding
* @return Pipeline\ConcurrentIterator<T>
* @throws CompressionUnavailable
*/
public function decode(
ReadableStream $in,
string $type,
Cancellation $cancellation,
?string $encoding = null,
): Pipeline\ConcurrentIterator {
$compressor = $encoding === null
? $this->compressor
: $this->compressors[$encoding] ?? throw new CompressionUnavailable($encoding);

/** @var Pipeline\Queue<T> $out */
$out = new Pipeline\Queue();

$parser = new Protocol\Parser(
$out->push(...),
$type,
$this->encoder,
$this->compressor,
$compressor,
);

EventLoop::queue(static function () use (
Expand Down
28 changes: 28 additions & 0 deletions packages/grpc/src/Metadata/AcceptEncoding.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
<?php

declare(strict_types=1);

namespace Thesis\Grpc\Metadata;

use Thesis\Grpc\Metadata;

/**
* @api
*/
final readonly class AcceptEncoding implements MetadataKey
{
public const string HEADER = 'grpc-accept-encoding';

/**
* @param list<non-empty-string> $encodings
*/
public function __construct(
public array $encodings,
) {}

#[\Override]
public function append(Metadata $md): Metadata
{
return $md->replace(self::HEADER, implode(',', $this->encodings));
}
}
63 changes: 62 additions & 1 deletion tests/CompressionTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -5,15 +5,22 @@
namespace Thesis\Grpc;

use Amp\Cancellation;
use Amp\Http\Server\Middleware\ClosureMiddleware;
use Amp\Http\Server\Request;
use Amp\Http\Server\RequestHandler;
use Amp\Http\Server\Response;
use Echos\Api\V1\EchoRequest;
use Echos\Api\V1\EchoResponse;
use Echos\Api\V1\EchoServiceClient;
use Echos\Api\V1\EchoServiceServer;
use Echos\Api\V1\EchoServiceServerRegistry;
use PHPUnit\Framework\Attributes\CoversClass;
use PHPUnit\Framework\Attributes\DataProvider;
use PHPUnit\Framework\TestCase;
use Thesis\Google\Rpc\Code;
use Thesis\Grpc\Client\Internal\AmphpHttpClient;
use Thesis\Grpc\Compression\Compressor;
use Thesis\Grpc\Compression\DeflateCompressor;
use Thesis\Grpc\Compression\GzipCompressor;
use Thesis\Grpc\Server\Internal\AmphpHttpServer;

Expand All @@ -22,6 +29,8 @@
#[CoversClass(Compressor::class)]
final class CompressionTest extends TestCase
{
private const string RESPONSE_ENCODING_HEADER = 'x-response-encoding';

private Server $server;

protected function setUp(): void
Expand All @@ -31,10 +40,19 @@ protected function setUp(): void
#[\Override]
public function echo(EchoRequest $request, Metadata $md, Cancellation $cancellation): EchoResponse
{
return new EchoResponse($request->sentence);
return new EchoResponse($request->sentence === Metadata\AcceptEncoding::HEADER ? ($md->value($request->sentence) ?? '') : $request->sentence);
}
}))
->withCompressors(new GzipCompressor())
->withMiddlewares(new ClosureMiddleware(static function (Request $request, RequestHandler $handler): Response {
$response = $handler->handleRequest($request);

if (($encoding = $request->getHeader(self::RESPONSE_ENCODING_HEADER)) !== null) {
$response->setHeader(Metadata\ContentEncoding::HEADER, $encoding);
}

return $response;
}))
->build();
$this->server->start();
}
Expand All @@ -60,6 +78,36 @@ public function testGzipCompressionUsed(): void
self::assertSame('Hello, gRPC', $client->echo(new EchoRequest('Hello, gRPC'))->sentence);
}

#[DataProvider('provideAcceptEncodingSentCases')]
public function testAcceptEncodingSent(Client\Builder $builder, string $expected): void
{
$client = new EchoServiceClient($builder->build());
self::assertSame($expected, $client->echo(new EchoRequest(Metadata\AcceptEncoding::HEADER))->sentence);
}

/**
* @return iterable<string, array{Client\Builder, string}>
*/
public static function provideAcceptEncodingSentCases(): iterable
{
yield 'default' => [
new Client\Builder(),
'identity',
];

yield 'compression' => [
new Client\Builder()->withCompression(new GzipCompressor()),
'gzip,identity',
];

yield 'compression and compressors' => [
new Client\Builder()
->withCompression(new GzipCompressor())
->withCompressors(new DeflateCompressor(), new GzipCompressor()),
'gzip,deflate,identity',
];
}

public function testUnknownForServerCompressionUsed(): void
{
$client = new EchoServiceClient(
Expand Down Expand Up @@ -89,4 +137,17 @@ public function decompress(string $buffer): string
$this->expectExceptionMessage('A grpc error with status code "UNIMPLEMENTED" and message "Decompression is not supported by server: strrev" occurred');
$client->echo(new EchoRequest('Hello, gRPC'));
}

public function testUnknownForClientResponseCompression(): void
{
$client = new EchoServiceClient(new Client\Builder()->build());

try {
$client->echo(new EchoRequest('Hello, gRPC'), new Metadata()->with(self::RESPONSE_ENCODING_HEADER, 'snappy'));
self::fail('InvokeError expected.');
} catch (InvokeError $e) {
self::assertSame(Code::INTERNAL, $e->statusCode);
self::assertSame("Compression algorithm 'snappy' is unavailable.", $e->statusMessage);
}
}
}
Loading
Loading