From 5b0aa10f4f2d75c704b8b116abef49df19b7d79a Mon Sep 17 00:00:00 2001 From: kafkiansky Date: Wed, 30 Sep 2026 10:01:38 +0300 Subject: [PATCH] grpc-accept-encoding on client-side --- packages/client/src/Client/Builder.php | 26 ++++++ .../AppendControlMetadataInterceptor.php | 3 + .../Client/Internal/Http2/StreamFactory.php | 21 ++++- .../grpc/src/Internal/Http2/StreamCodec.php | 26 +++++- packages/grpc/src/Metadata/AcceptEncoding.php | 28 ++++++ tests/CompressionTest.php | 63 +++++++++++++- tests/Internal/Http2/StreamCodecTest.php | 86 +++++++++++++++++++ 7 files changed, 249 insertions(+), 4 deletions(-) create mode 100644 packages/grpc/src/Metadata/AcceptEncoding.php create mode 100644 tests/Internal/Http2/StreamCodecTest.php diff --git a/packages/client/src/Client/Builder.php b/packages/client/src/Client/Builder.php index c02c416..66f865d 100644 --- a/packages/client/src/Client/Builder.php +++ b/packages/client/src/Client/Builder.php @@ -36,6 +36,9 @@ final class Builder private ?Compressor $compressor = null; + /** @var list */ + private array $compressors = []; + private ?DelegateHttpClient $httpclient = null; /** @var list */ @@ -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; @@ -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(); @@ -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 @@ -336,6 +361,7 @@ public function build(): Client inactivityTimeout: $inactivityTimeout, encoder: $encoder, compressor: $compressor, + compressors: $compressors, ), ), ), diff --git a/packages/client/src/Client/Internal/AppendControlMetadataInterceptor.php b/packages/client/src/Client/Internal/AppendControlMetadataInterceptor.php index aaade6c..4f71799 100644 --- a/packages/client/src/Client/Internal/AppendControlMetadataInterceptor.php +++ b/packages/client/src/Client/Internal/AppendControlMetadataInterceptor.php @@ -21,10 +21,12 @@ /** * @param non-empty-string $encoding * @param non-empty-string $compression + * @param list $acceptEncoding */ public function __construct( private string $encoding, private string $compression, + private array $acceptEncoding, ) {} #[\Override] @@ -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'); } } diff --git a/packages/client/src/Client/Internal/Http2/StreamFactory.php b/packages/client/src/Client/Internal/Http2/StreamFactory.php index 7f9a465..03ba812 100644 --- a/packages/client/src/Client/Internal/Http2/StreamFactory.php +++ b/packages/client/src/Client/Internal/Http2/StreamFactory.php @@ -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; @@ -30,6 +33,9 @@ { private StreamCodec $codec; + /** + * @param list $compressors + */ public function __construct( private DelegateHttpClient $http, private UriFactory $uri, @@ -38,10 +44,12 @@ public function __construct( private float $inactivityTimeout, Encoder $encoder, Compressor $compressor, + array $compressors, ) { $this->codec = new StreamCodec( $encoder, $compressor, + $compressors, ); } @@ -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(), ); diff --git a/packages/grpc/src/Internal/Http2/StreamCodec.php b/packages/grpc/src/Internal/Http2/StreamCodec.php index d28a378..d38a698 100644 --- a/packages/grpc/src/Internal/Http2/StreamCodec.php +++ b/packages/grpc/src/Internal/Http2/StreamCodec.php @@ -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; @@ -18,10 +19,24 @@ */ final readonly class StreamCodec { + /** @var array */ + private array $compressors; + + /** + * @param list $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 @@ -78,13 +93,20 @@ public function encode( /** * @template T of object * @param class-string $type + * @param ?non-empty-string $encoding * @return Pipeline\ConcurrentIterator + * @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 $out */ $out = new Pipeline\Queue(); @@ -92,7 +114,7 @@ public function decode( $out->push(...), $type, $this->encoder, - $this->compressor, + $compressor, ); EventLoop::queue(static function () use ( diff --git a/packages/grpc/src/Metadata/AcceptEncoding.php b/packages/grpc/src/Metadata/AcceptEncoding.php new file mode 100644 index 0000000..0987264 --- /dev/null +++ b/packages/grpc/src/Metadata/AcceptEncoding.php @@ -0,0 +1,28 @@ + $encodings + */ + public function __construct( + public array $encodings, + ) {} + + #[\Override] + public function append(Metadata $md): Metadata + { + return $md->replace(self::HEADER, implode(',', $this->encodings)); + } +} diff --git a/tests/CompressionTest.php b/tests/CompressionTest.php index 63acf17..4e7de86 100644 --- a/tests/CompressionTest.php +++ b/tests/CompressionTest.php @@ -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; @@ -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 @@ -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(); } @@ -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 + */ + 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( @@ -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); + } + } } diff --git a/tests/Internal/Http2/StreamCodecTest.php b/tests/Internal/Http2/StreamCodecTest.php new file mode 100644 index 0000000..af4bed1 --- /dev/null +++ b/tests/Internal/Http2/StreamCodecTest.php @@ -0,0 +1,86 @@ +encode( + Pipeline::fromIterable($messages)->getIterator(), + new NullCancellation(), + ), + preserve_keys: false, + )); + + self::assertEquals($messages, iterator_to_array( + $codec->decode(new ReadableBuffer($frames), EchoRequest::class, new NullCancellation(), $encoding), + preserve_keys: false, + )); + } + + /** + * @return iterable + */ + public static function provideDecodeCases(): iterable + { + yield 'default compressor' => [ + new GzipCompressor(), + new StreamCodec(ProtobufEncoder::default(), new GzipCompressor()), + null, + ]; + + yield 'own compressor by encoding' => [ + new GzipCompressor(), + new StreamCodec(ProtobufEncoder::default(), new GzipCompressor()), + 'gzip', + ]; + + yield 'additional compressor by encoding' => [ + new DeflateCompressor(), + new StreamCodec(ProtobufEncoder::default(), IdentityCompressor::Compressor, [new GzipCompressor(), new DeflateCompressor()]), + 'deflate', + ]; + + yield 'identity by encoding' => [ + IdentityCompressor::Compressor, + new StreamCodec(ProtobufEncoder::default(), new GzipCompressor(), [IdentityCompressor::Compressor]), + 'identity', + ]; + } + + public function testDecodeUnknownEncoding(): void + { + $codec = new StreamCodec(ProtobufEncoder::default(), IdentityCompressor::Compressor, [new GzipCompressor()]); + + $this->expectExceptionObject(new CompressionUnavailable('snappy')); + $codec->decode(new ReadableBuffer(), EchoRequest::class, new NullCancellation(), 'snappy'); + } +}