From 64432f0d773140ee1ac6f11f6ede3e085b6d9adf Mon Sep 17 00:00:00 2001 From: kafkiansky Date: Thu, 1 Oct 2026 21:05:19 +0300 Subject: [PATCH] Limit message size instead of stream size --- packages/client/src/Client/Builder.php | 17 ++ .../Client/Internal/Http2/StreamFactory.php | 5 + .../grpc/src/Internal/Http2/StreamCodec.php | 3 + .../grpc/src/Internal/Protocol/Parser.php | 9 ++ packages/server/src/Server/Builder.php | 20 ++- .../Internal/Http2/ServerRequestHandler.php | 5 +- .../Server/Internal/Http2/StreamFactory.php | 5 + tests/Internal/Http2/StreamCodecTest.php | 14 +- tests/Internal/Protocol/ParserTest.php | 17 ++ tests/MessageSizeTest.php | 149 ++++++++++++++++++ 10 files changed, 236 insertions(+), 8 deletions(-) create mode 100644 tests/MessageSizeTest.php diff --git a/packages/client/src/Client/Builder.php b/packages/client/src/Client/Builder.php index 66f865d..7508a11 100644 --- a/packages/client/src/Client/Builder.php +++ b/packages/client/src/Client/Builder.php @@ -30,6 +30,7 @@ final class Builder private const int DEFAULT_CONNECTION_LIMIT = \PHP_INT_MAX; private const int DEFAULT_TRANSFER_TIMEOUT = 0; // no transfer timeout private const int DEFAULT_INACTIVITY_TIMEOUT = 0; // no inactivity timeout + private const int DEFAULT_MAX_RECEIVE_MESSAGE_SIZE = 4 * 1_024 * 1_024; /** @var ?non-empty-string */ private ?string $host = null; @@ -58,6 +59,9 @@ final class Builder private float $inactivityTimeout = self::DEFAULT_INACTIVITY_TIMEOUT; + /** @var positive-int */ + private int $maxReceiveMessageSize = self::DEFAULT_MAX_RECEIVE_MESSAGE_SIZE; + private ?SocketConnector $connector = null; private ?Internal\KeepaliveSettings $keepalive = null; @@ -214,6 +218,17 @@ public function withInactivityTimeout(float $inactivityTimeout): self return $builder; } + /** + * @param positive-int $bytes + */ + public function withMaxReceiveMessageSize(int $bytes): self + { + $builder = clone $this; + $builder->maxReceiveMessageSize = $bytes; + + return $builder; + } + public function withTransportCredentials(TransportCredentials $credentials): self { $builder = clone $this; @@ -293,6 +308,7 @@ public function build(): Client $uriFactory = new Http2\UriFactory($tlsContext !== null ? Internal\HttpScheme::Https : Internal\HttpScheme::Http); $transferTimeout = $this->transferTimeout; $inactivityTimeout = $this->inactivityTimeout; + $maxReceiveMessageSize = $this->maxReceiveMessageSize; $resolver = $this->endpointResolvers[$target->scheme] ?? match (Scheme::tryFrom($target->scheme)) { Scheme::Dns => new EndpointResolver\DnsResolver(), @@ -361,6 +377,7 @@ public function build(): Client inactivityTimeout: $inactivityTimeout, encoder: $encoder, compressor: $compressor, + maxReceiveMessageSize: $maxReceiveMessageSize, compressors: $compressors, ), ), diff --git a/packages/client/src/Client/Internal/Http2/StreamFactory.php b/packages/client/src/Client/Internal/Http2/StreamFactory.php index 03ba812..cb53595 100644 --- a/packages/client/src/Client/Internal/Http2/StreamFactory.php +++ b/packages/client/src/Client/Internal/Http2/StreamFactory.php @@ -34,6 +34,7 @@ private StreamCodec $codec; /** + * @param positive-int $maxReceiveMessageSize * @param list $compressors */ public function __construct( @@ -44,11 +45,13 @@ public function __construct( private float $inactivityTimeout, Encoder $encoder, Compressor $compressor, + int $maxReceiveMessageSize, array $compressors, ) { $this->codec = new StreamCodec( $encoder, $compressor, + $maxReceiveMessageSize, $compressors, ); } @@ -82,6 +85,8 @@ public function create( $request->setHeaders($md->kv); $request->setTransferTimeout($this->transferTimeout); $request->setInactivityTimeout($this->inactivityTimeout); + // gRPC limits the size of a single message, not the stream, see {@see StreamCodec}. + $request->setBodySizeLimit(\PHP_INT_MAX); // If the program terminates after making a request, the HTTP client may not have enough time to finish sending the request body and trailers, // causing an error on the server side — after a certain timeout, the server will detect that the client unexpectedly closed the connection. diff --git a/packages/grpc/src/Internal/Http2/StreamCodec.php b/packages/grpc/src/Internal/Http2/StreamCodec.php index d38a698..bca50c3 100644 --- a/packages/grpc/src/Internal/Http2/StreamCodec.php +++ b/packages/grpc/src/Internal/Http2/StreamCodec.php @@ -23,11 +23,13 @@ private array $compressors; /** + * @param positive-int $maxReceiveMessageSize * @param list $compressors to decompress messages with, selected by the peer's "grpc-encoding" */ public function __construct( private Encoder $encoder, private Compressor $compressor, + private int $maxReceiveMessageSize, array $compressors = [], ) { $compressors = [$compressor, ...$compressors]; @@ -115,6 +117,7 @@ public function decode( $type, $this->encoder, $compressor, + $this->maxReceiveMessageSize, ); EventLoop::queue(static function () use ( diff --git a/packages/grpc/src/Internal/Protocol/Parser.php b/packages/grpc/src/Internal/Protocol/Parser.php index ff1ce49..e89b0ed 100644 --- a/packages/grpc/src/Internal/Protocol/Parser.php +++ b/packages/grpc/src/Internal/Protocol/Parser.php @@ -4,8 +4,10 @@ namespace Thesis\Grpc\Internal\Protocol; +use Thesis\Google\Rpc\Code; use Thesis\Grpc\Compression; use Thesis\Grpc\Encoding; +use Thesis\Grpc\InvokeError; /** * @internal @@ -18,17 +20,20 @@ final class Parser /** * @param \Closure(T): void $push * @param class-string $type + * @param positive-int $maxMessageSize */ public function __construct( private readonly \Closure $push, private readonly string $type, private readonly Encoding\Encoder $encoder, private readonly Compression\Compressor $compressor, + private readonly int $maxMessageSize, ) {} /** * @throws Compression\DecompressionFailed * @throws Encoding\DecodingFailed + * @throws InvokeError */ public function push(string $data): void { @@ -40,6 +45,10 @@ public function push(string $data): void substr($this->buffer, lengthOffset, 4), ); + if ($messageLength > $this->maxMessageSize) { + throw new InvokeError(Code::RESOURCE_EXHAUSTED, "Received message larger than max ({$messageLength} vs. {$this->maxMessageSize})"); + } + $frameSize = bodyOffset + $messageLength; if (\strlen($this->buffer) < $frameSize) { diff --git a/packages/server/src/Server/Builder.php b/packages/server/src/Server/Builder.php index e9f7d50..deb6814 100644 --- a/packages/server/src/Server/Builder.php +++ b/packages/server/src/Server/Builder.php @@ -42,7 +42,10 @@ final class Builder private const int DEFAULT_STREAM_TIMEOUT = HttpDriver::DEFAULT_STREAM_TIMEOUT; private const int DEFAULT_CONNECTION_TIMEOUT = HttpDriver::DEFAULT_CONNECTION_TIMEOUT; private const int DEFAULT_HEADER_SIZE_LIMIT = HttpDriver::DEFAULT_HEADER_SIZE_LIMIT; - private const int DEFAULT_BODY_SIZE_LIMIT = HttpDriver::DEFAULT_BODY_SIZE_LIMIT; + + // gRPC limits the size of a single message, not the stream. Not PHP_INT_MAX, since amphp adds one to the limit when updating the flow-control window. + private const int DEFAULT_BODY_SIZE_LIMIT = \PHP_INT_MAX - 1; + private const int DEFAULT_MAX_RECEIVE_MESSAGE_SIZE = 4 * 1_024 * 1_024; /** @var list */ private const array ALLOWED_HTTP_METHODS = ['POST']; @@ -98,6 +101,9 @@ final class Builder /** @var positive-int */ private int $bodySizeLimit = self::DEFAULT_BODY_SIZE_LIMIT; + /** @var positive-int */ + private int $maxReceiveMessageSize = self::DEFAULT_MAX_RECEIVE_MESSAGE_SIZE; + /** * Required for encoding the `grpc-status-details-bin` header, `status`, and `details`. */ @@ -344,6 +350,17 @@ public function withBodySizeLimit(int $limit): self return $builder; } + /** + * @param positive-int $bytes + */ + public function withMaxReceiveMessageSize(int $bytes): self + { + $builder = clone $this; + $builder->maxReceiveMessageSize = $bytes; + + return $builder; + } + public function build(): Server { $logger = $this->logger ?? new NullLogger(); @@ -418,6 +435,7 @@ public function build(): Server protobuf: $this->protobuf ?? Protobuf\Encoder\Builder::buildDefault(), unaryInterceptors: $this->unaryInterceptors, streamInterceptors: $this->streamInterceptors, + maxReceiveMessageSize: $this->maxReceiveMessageSize, ), errorHandler: new ServerErrorHandler(), ); diff --git a/packages/server/src/Server/Internal/Http2/ServerRequestHandler.php b/packages/server/src/Server/Internal/Http2/ServerRequestHandler.php index ed41275..f87a6cb 100644 --- a/packages/server/src/Server/Internal/Http2/ServerRequestHandler.php +++ b/packages/server/src/Server/Internal/Http2/ServerRequestHandler.php @@ -58,6 +58,7 @@ final class ServerRequestHandler implements /** * @param list $unaryInterceptors * @param list $streamInterceptors + * @param positive-int $maxReceiveMessageSize */ public function __construct( private readonly MessageEncoderFactory $encoderFactory, @@ -65,6 +66,7 @@ public function __construct( Protobuf\Encoder $protobuf, array $unaryInterceptors, array $streamInterceptors, + private readonly int $maxReceiveMessageSize, ) { $this->pending = new \WeakMap(); $this->router = new Router(); @@ -129,9 +131,10 @@ public function handleRequest(Request $request): Response /** @var array $streams */ static $streams = []; - $factory = $streams["{$encoder->name()}\0{$compressor->name()}"] ??= new StreamFactory( + $factory = $streams["{$encoder->name()}\0{$compressor->name()}\0{$this->maxReceiveMessageSize}"] ??= new StreamFactory( $encoder, $compressor, + $this->maxReceiveMessageSize, ); $response = new Response(status: HttpStatus::OK, headers: $headers->kv); diff --git a/packages/server/src/Server/Internal/Http2/StreamFactory.php b/packages/server/src/Server/Internal/Http2/StreamFactory.php index fd0f6c7..624fb61 100644 --- a/packages/server/src/Server/Internal/Http2/StreamFactory.php +++ b/packages/server/src/Server/Internal/Http2/StreamFactory.php @@ -25,13 +25,18 @@ { private StreamCodec $codec; + /** + * @param positive-int $maxReceiveMessageSize + */ public function __construct( Encoder $encoder, Compressor $compressor, + int $maxReceiveMessageSize, ) { $this->codec = new StreamCodec( $encoder, $compressor, + $maxReceiveMessageSize, ); } diff --git a/tests/Internal/Http2/StreamCodecTest.php b/tests/Internal/Http2/StreamCodecTest.php index af4bed1..275392a 100644 --- a/tests/Internal/Http2/StreamCodecTest.php +++ b/tests/Internal/Http2/StreamCodecTest.php @@ -21,6 +21,8 @@ #[CoversClass(StreamCodec::class)] final class StreamCodecTest extends TestCase { + private const int MAX_MESSAGE_SIZE = 4 * 1_024 * 1_024; + /** * @param ?non-empty-string $encoding */ @@ -33,7 +35,7 @@ public function testDecode(Compressor $sender, StreamCodec $codec, ?string $enco ]; $frames = implode('', iterator_to_array( - new StreamCodec(ProtobufEncoder::default(), $sender)->encode( + new StreamCodec(ProtobufEncoder::default(), $sender, self::MAX_MESSAGE_SIZE)->encode( Pipeline::fromIterable($messages)->getIterator(), new NullCancellation(), ), @@ -53,32 +55,32 @@ public static function provideDecodeCases(): iterable { yield 'default compressor' => [ new GzipCompressor(), - new StreamCodec(ProtobufEncoder::default(), new GzipCompressor()), + new StreamCodec(ProtobufEncoder::default(), new GzipCompressor(), self::MAX_MESSAGE_SIZE), null, ]; yield 'own compressor by encoding' => [ new GzipCompressor(), - new StreamCodec(ProtobufEncoder::default(), new GzipCompressor()), + new StreamCodec(ProtobufEncoder::default(), new GzipCompressor(), self::MAX_MESSAGE_SIZE), 'gzip', ]; yield 'additional compressor by encoding' => [ new DeflateCompressor(), - new StreamCodec(ProtobufEncoder::default(), IdentityCompressor::Compressor, [new GzipCompressor(), new DeflateCompressor()]), + new StreamCodec(ProtobufEncoder::default(), IdentityCompressor::Compressor, self::MAX_MESSAGE_SIZE, [new GzipCompressor(), new DeflateCompressor()]), 'deflate', ]; yield 'identity by encoding' => [ IdentityCompressor::Compressor, - new StreamCodec(ProtobufEncoder::default(), new GzipCompressor(), [IdentityCompressor::Compressor]), + new StreamCodec(ProtobufEncoder::default(), new GzipCompressor(), self::MAX_MESSAGE_SIZE, [IdentityCompressor::Compressor]), 'identity', ]; } public function testDecodeUnknownEncoding(): void { - $codec = new StreamCodec(ProtobufEncoder::default(), IdentityCompressor::Compressor, [new GzipCompressor()]); + $codec = new StreamCodec(ProtobufEncoder::default(), IdentityCompressor::Compressor, self::MAX_MESSAGE_SIZE, [new GzipCompressor()]); $this->expectExceptionObject(new CompressionUnavailable('snappy')); $codec->decode(new ReadableBuffer(), EchoRequest::class, new NullCancellation(), 'snappy'); diff --git a/tests/Internal/Protocol/ParserTest.php b/tests/Internal/Protocol/ParserTest.php index 8553fdc..3e02c60 100644 --- a/tests/Internal/Protocol/ParserTest.php +++ b/tests/Internal/Protocol/ParserTest.php @@ -8,9 +8,11 @@ use PHPUnit\Framework\Attributes\CoversClass; use PHPUnit\Framework\Attributes\DataProvider; use PHPUnit\Framework\TestCase; +use Thesis\Google\Rpc\Code; use Thesis\Grpc\Compression\Compressor; use Thesis\Grpc\Compression\GzipCompressor; use Thesis\Grpc\Compression\IdentityCompressor; +use Thesis\Grpc\InvokeError; use Thesis\Grpc\Protobuf\ProtobufEncoder; #[CoversClass(Parser::class)] @@ -32,6 +34,7 @@ static function (EchoRequest $request) use (&$frames): void { EchoRequest::class, ProtobufEncoder::default(), $compressor, + 4_096, ); foreach ($chunks as $chunk) { @@ -122,6 +125,20 @@ public static function providePushCases(): iterable } } + public function testMessageTooLarge(): void + { + $parser = new Parser( + static fn(EchoRequest $request) => self::fail('No message expected.'), + EchoRequest::class, + ProtobufEncoder::default(), + IdentityCompressor::Compressor, + 4_096, + ); + + $this->expectExceptionObject(new InvokeError(Code::RESOURCE_EXHAUSTED, 'Received message larger than max (4097 vs. 4096)')); + $parser->push(pack('CN', 0, 4_097)); + } + private static function frame(string $sentence, Compressor $compressor = IdentityCompressor::Compressor): string { $payload = ProtobufEncoder::default()->encode(new EchoRequest($sentence)); diff --git a/tests/MessageSizeTest.php b/tests/MessageSizeTest.php new file mode 100644 index 0000000..05d66ae --- /dev/null +++ b/tests/MessageSizeTest.php @@ -0,0 +1,149 @@ +build())->echo(new EchoRequest($sentence))->sentence); + } finally { + $server->stop(); + } + } + + public function testClientStream(): void + { + $server = self::server(new Server\Builder()); + + try { + $stream = new FileServiceClient(new Client\Builder()->build())->upload(); + + for ($i = 0; $i < 4; ++$i) { + $stream->send(new Chunk(random_bytes(64 * 1_024))); + } + + self::assertSame(4 * 64 * 1_024, $stream->close()->size); + } finally { + $server->stop(); + } + } + + public function testServerStream(): void + { + self::markTestSkipped('Corrupts and hangs until https://github.com/amphp/http-server/pull/396 is released.'); + + $server = self::server(new Server\Builder()); // @phpstan-ignore deadCode.unreachable (the test is skipped until the amphp fix is released) + + try { + $events = iterator_to_array(new TopicServiceClient(new Client\Builder()->build())->subscribe(new SubscribeRequest('11')), preserve_keys: false); + + self::assertSame(11 * self::PAYLOAD_SIZE, array_sum(array_map(static fn(Event $event): int => \strlen($event->payload), $events))); + } finally { + $server->stop(); + } + } + + #[DataProvider('provideMessageTooLargeCases')] + public function testMessageTooLarge(Server\Builder $serverBuilder, Client\Builder $clientBuilder): void + { + $server = self::server($serverBuilder); + + try { + $this->expectExceptionObject(new InvokeError(Code::RESOURCE_EXHAUSTED, 'Received message larger than max (2051 vs. 1024)')); + new EchoServiceClient($clientBuilder->build())->echo(new EchoRequest(str_repeat('a', 2_048))); + } finally { + $server->stop(); + } + } + + /** + * @return iterable + */ + public static function provideMessageTooLargeCases(): iterable + { + yield 'server receives' => [ + new Server\Builder()->withMaxReceiveMessageSize(1_024), + new Client\Builder(), + ]; + + yield 'client receives' => [ + new Server\Builder(), + new Client\Builder()->withMaxReceiveMessageSize(1_024), + ]; + } + + private static function server(Server\Builder $builder): Server + { + $server = $builder + ->withServices( + new EchoServiceServerRegistry(new readonly class implements EchoServiceServer { + #[\Override] + public function echo(EchoRequest $request, Metadata $md, Cancellation $cancellation): EchoResponse + { + return new EchoResponse($request->sentence); + } + }), + new FileServiceServerRegistry(new readonly class implements FileServiceServer { + #[\Override] + public function upload(Server\ClientStreamChannel $stream, Metadata $md, Cancellation $cancellation): FileInfo + { + $size = 0; + + /** @var Chunk $chunk */ + foreach ($stream as $chunk) { + $size += \strlen($chunk->content); + } + + return new FileInfo($size); + } + }), + new TopicServiceServerRegistry(new readonly class implements TopicServiceServer { + #[\Override] + public function subscribe(SubscribeRequest $request, Metadata $md, Cancellation $cancellation): iterable + { + for ($i = 0; $i < (int) $request->topic; ++$i) { + yield new Event(payload: str_repeat('x', MessageSizeTest::PAYLOAD_SIZE)); + } + } + }), + ) + ->build(); + $server->start(); + + return $server; + } +}