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
17 changes: 17 additions & 0 deletions packages/client/src/Client/Builder.php
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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(),
Expand Down Expand Up @@ -361,6 +377,7 @@ public function build(): Client
inactivityTimeout: $inactivityTimeout,
encoder: $encoder,
compressor: $compressor,
maxReceiveMessageSize: $maxReceiveMessageSize,
compressors: $compressors,
),
),
Expand Down
5 changes: 5 additions & 0 deletions packages/client/src/Client/Internal/Http2/StreamFactory.php
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@
private StreamCodec $codec;

/**
* @param positive-int $maxReceiveMessageSize
* @param list<Compressor> $compressors
*/
public function __construct(
Expand All @@ -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,
);
}
Expand Down Expand Up @@ -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.
Expand Down
3 changes: 3 additions & 0 deletions packages/grpc/src/Internal/Http2/StreamCodec.php
Original file line number Diff line number Diff line change
Expand Up @@ -23,11 +23,13 @@
private array $compressors;

/**
* @param positive-int $maxReceiveMessageSize
* @param list<Compressor> $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];
Expand Down Expand Up @@ -115,6 +117,7 @@ public function decode(
$type,
$this->encoder,
$compressor,
$this->maxReceiveMessageSize,
);

EventLoop::queue(static function () use (
Expand Down
9 changes: 9 additions & 0 deletions packages/grpc/src/Internal/Protocol/Parser.php
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -18,17 +20,20 @@ final class Parser
/**
* @param \Closure(T): void $push
* @param class-string<T> $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
{
Expand All @@ -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) {
Expand Down
20 changes: 19 additions & 1 deletion packages/server/src/Server/Builder.php
Original file line number Diff line number Diff line change
Expand Up @@ -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<non-empty-string> */
private const array ALLOWED_HTTP_METHODS = ['POST'];
Expand Down Expand Up @@ -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`.
*/
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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(),
);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -58,13 +58,15 @@ final class ServerRequestHandler implements
/**
* @param list<UnaryInterceptor> $unaryInterceptors
* @param list<StreamInterceptor> $streamInterceptors
* @param positive-int $maxReceiveMessageSize
*/
public function __construct(
private readonly MessageEncoderFactory $encoderFactory,
private readonly MessageCompressorFactory $compressorFactory,
Protobuf\Encoder $protobuf,
array $unaryInterceptors,
array $streamInterceptors,
private readonly int $maxReceiveMessageSize,
) {
$this->pending = new \WeakMap();
$this->router = new Router();
Expand Down Expand Up @@ -129,9 +131,10 @@ public function handleRequest(Request $request): Response

/** @var array<non-empty-string, StreamFactory> $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);
Expand Down
5 changes: 5 additions & 0 deletions packages/server/src/Server/Internal/Http2/StreamFactory.php
Original file line number Diff line number Diff line change
Expand Up @@ -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,
);
}

Expand Down
14 changes: 8 additions & 6 deletions tests/Internal/Http2/StreamCodecTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -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
*/
Expand All @@ -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(),
),
Expand All @@ -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');
Expand Down
17 changes: 17 additions & 0 deletions tests/Internal/Protocol/ParserTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -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)]
Expand All @@ -32,6 +34,7 @@ static function (EchoRequest $request) use (&$frames): void {
EchoRequest::class,
ProtobufEncoder::default(),
$compressor,
4_096,
);

foreach ($chunks as $chunk) {
Expand Down Expand Up @@ -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));
Expand Down
Loading
Loading