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
11 changes: 0 additions & 11 deletions Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -133,17 +133,6 @@ fix: fixer rector composer-normalize ## Run all fixing recipes
check: fixer-check rector-check composer-validate composer-normalize-check deps-analyze phpstan test ## Run all project checks
.PHONY: check

compile-stub:
docker run --rm \
--pull always \
--user $(CONTAINER_USER) \
-v $(PWD):/workspace \
-w /workspace \
ghcr.io/thesis-php/protoc-plugin:latest \
--php-plugin_out=genproto \
protos/*.proto
.PHONY: compile-stub

compile-test-stub:
docker run --rm \
--pull always \
Expand Down
22 changes: 18 additions & 4 deletions packages/client/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -184,17 +184,16 @@ final readonly class RandomBalancerFactory implements LoadBalancerFactory

## Endpoint resolution

Resolver is selected by target scheme. You can override resolver for a specific scheme:
Resolver is selected by target scheme. You can override the resolver for a scheme:

```php
use Amp\Cache\LocalCache;
use Thesis\Grpc\Client\Builder;
use Thesis\Grpc\Client\EndpointResolver\DnsResolver;
use Thesis\Grpc\Client\Scheme;

$client = new Builder()
->withHost('dns:///my-grpc-server:50051')
->withEndpointResolver(Scheme::Dns, new DnsResolver(
->withEndpointResolver('dns', new DnsResolver(
cache: new LocalCache(),
minResolveInterval: 60,
maxResolveInterval: 600,
Expand All @@ -208,7 +207,22 @@ Default resolvers by scheme:
- `ipv4`, `ipv6`, `unix` -> `StaticResolver`
- `passthrough` -> `PassthroughResolver`

You can also implement your own `EndpointResolver` for service discovery backends like Consul/etcd.
### Custom schemes (service discovery)

Register a resolver under any scheme to plug in a service registry such as etcd or
Consul. The target is then `<scheme>://[authority]/<endpoint>`, and the endpoint is
handed to the resolver as `$target->opaque`:

```php
$client = new Builder()
->withEndpointResolver('etcd', new EtcdResolver($etcd, prefix: '/services/'))
->withLoadBalancer(new RoundRobinFactory())
->withHost('etcd:///echo')
->build();
```

An `EndpointResolver` returns the current `Resolution` and may push updates over time
through `EndpointResolverListener::onResolve()`, which the balancer picks up live.

## Error handling

Expand Down
17 changes: 8 additions & 9 deletions packages/client/src/Client/Builder.php
Original file line number Diff line number Diff line change
Expand Up @@ -66,13 +66,8 @@ final class Builder

private ?LoadBalancerFactory $loadBalancerFactory = null;

/** @var \SplObjectStorage<Scheme, EndpointResolver> */
private \SplObjectStorage $endpointResolvers;

public function __construct()
{
$this->endpointResolvers = new \SplObjectStorage();
}
/** @var array<string, EndpointResolver> */
private array $endpointResolvers = [];

public function withProtobuf(Decoder $decoder): self
{
Expand Down Expand Up @@ -204,7 +199,7 @@ public function withLoadBalancer(LoadBalancerFactory $factory): self
return $builder;
}

public function withEndpointResolver(Scheme $scheme, EndpointResolver $resolver): self
public function withEndpointResolver(string $scheme, EndpointResolver $resolver): self
{
$builder = clone $this;
$builder->endpointResolvers[$scheme] = $resolver;
Expand All @@ -217,6 +212,9 @@ public static function buildDefault(): Client
return new self()->build();
}

/**
* @throws InvalidTarget
*/
public function build(): Client
{
$target = Target::parse($this->host ?? self::DEFAULT_HOST);
Expand All @@ -230,10 +228,11 @@ public function build(): Client
$transferTimeout = $this->transferTimeout;
$inactivityTimeout = $this->inactivityTimeout;

$resolver = $this->endpointResolvers[$target->scheme] ?? match ($target->scheme) {
$resolver = $this->endpointResolvers[$target->scheme] ?? match (Scheme::tryFrom($target->scheme)) {
Scheme::Dns => new EndpointResolver\DnsResolver(),
Scheme::Passthrough => new EndpointResolver\PassthroughResolver(),
Scheme::Ipv4, Scheme::Ipv6, Scheme::Unix => new EndpointResolver\StaticResolver(),
null => throw new InvalidTarget($this->host ?? ''),
};

$controlMetadata = new Internal\AppendControlMetadataInterceptor(
Expand Down
56 changes: 48 additions & 8 deletions packages/client/src/Client/Target.php
Original file line number Diff line number Diff line change
Expand Up @@ -28,25 +28,33 @@ public static function parse(string $target): self
}

return match ($scheme) {
Scheme::Dns => self::parseDns($addr, $target, $scheme),
Scheme::Dns => self::parseDns($addr, $target),
Scheme::Passthrough => self::parsePassthrough($addr, $target),
Scheme::Ipv4, Scheme::Ipv6 => new self($scheme, self::parseAddresses($addr, $target), opaque: $addr),
Scheme::Unix => new self($scheme, [self::parseUnix($addr, $target)], opaque: $addr),
Scheme::Ipv4, Scheme::Ipv6 => new self($scheme->value, self::parseAddresses($addr, $target), opaque: $addr),
Scheme::Unix => new self($scheme->value, [self::parseUnix($addr, $target)], opaque: $addr),
};
}
}

return new self(Scheme::Dns, self::parseAddresses($target), opaque: $target);
// A custom scheme in URI form ("<scheme>://[authority]/<endpoint>") selects
// a user-registered resolver. Only the "//" form is treated as a scheme, so
// a bare "host:port" is still resolved as a DNS target.
if (preg_match('#^([a-z][a-z0-9+.\-]*)://#', $target, $matches) === 1) {
return self::parseCustom($matches[1], $target);
}

return new self(Scheme::Dns->value, self::parseAddresses($target), opaque: $target);
}

/**
* @internal use {@see Target::parse()} instead
* @param non-empty-string $scheme
* @param non-empty-list<TargetAddress> $addresses
* @param non-empty-string $opaque Raw value after scheme prefix
* @param ?non-empty-string $authority DNS server address (only for dns://authority/host form)
*/
public function __construct(
public Scheme $scheme,
public string $scheme,
public array $addresses,
public string $opaque,
public ?string $authority = null,
Expand All @@ -57,7 +65,7 @@ public function __construct(
* @param non-empty-string $target
* @throws InvalidTarget
*/
private static function parseDns(string $addr, string $target, Scheme $scheme): self
private static function parseDns(string $addr, string $target): self
{
$opaque = $addr;
$authority = null;
Expand All @@ -81,13 +89,45 @@ private static function parseDns(string $addr, string $target, Scheme $scheme):
}

return new self(
$scheme,
Scheme::Dns->value,
self::parseAddresses($addr, $target),
$opaque,
$authority,
);
}

/**
* A scheme with no built-in parser: "<scheme>://[authority]/<endpoint>". The
* endpoint is handed to a user-registered resolver as {@see self::$opaque}.
*
* @param non-empty-string $scheme
* @param non-empty-string $target
* @throws InvalidTarget
*/
private static function parseCustom(string $scheme, string $target): self
{
$rest = substr($target, \strlen($scheme) + 3);

$slash = strpos($rest, '/');
if ($slash === false) {
throw new InvalidTarget($target);
}

$authority = substr($rest, 0, $slash);
$endpoint = substr($rest, $slash + 1);

if ($endpoint === '') {
throw new InvalidTarget($target);
}

return new self(
$scheme,
[new TargetAddress($endpoint, 0)],
$endpoint,
$authority !== '' ? $authority : null,
);
}

/**
* @param non-empty-string $addr
* @param non-empty-string $target
Expand All @@ -110,7 +150,7 @@ private static function parsePassthrough(string $addr, string $target): self
throw new InvalidTarget($target);
}

return new self(Scheme::Passthrough, [new TargetAddress($addr, 0)], $addr);
return new self(Scheme::Passthrough->value, [new TargetAddress($addr, 0)], $addr);
}

/**
Expand Down
15 changes: 7 additions & 8 deletions tests/Client/EndpointResolver/DnsResolverTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,6 @@
use Thesis\Grpc\Client\Endpoint;
use Thesis\Grpc\Client\EndpointResolverListener;
use Thesis\Grpc\Client\Resolution;
use Thesis\Grpc\Client\Scheme;
use Thesis\Grpc\Client\Target;
use Thesis\Grpc\Client\TargetAddress;
use function Amp\delay;
Expand Down Expand Up @@ -56,13 +55,13 @@ public function testResolve(Target $target, array $records, array $endpoints): v
public static function provideResolveCases(): iterable
{
yield 'single A record' => [
new Target(Scheme::Dns, [new TargetAddress('myhost', 50_051)], 'myhost:50051'),
new Target('dns', [new TargetAddress('myhost', 50_051)], 'myhost:50051'),
[new DnsRecord('192.168.0.1', DnsRecord::A, 300)],
[new Endpoint(new Address('192.168.0.1:50051'))],
];

yield 'multiple A records' => [
new Target(Scheme::Dns, [new TargetAddress('myhost', 50_051)], 'myhost:50051'),
new Target('dns', [new TargetAddress('myhost', 50_051)], 'myhost:50051'),
[
new DnsRecord('192.168.0.1', DnsRecord::A, 300),
new DnsRecord('192.168.0.2', DnsRecord::A, 300),
Expand All @@ -74,13 +73,13 @@ public static function provideResolveCases(): iterable
];

yield 'AAAA record wraps in brackets' => [
new Target(Scheme::Dns, [new TargetAddress('myhost', 50_051)], 'myhost:50051'),
new Target('dns', [new TargetAddress('myhost', 50_051)], 'myhost:50051'),
[new DnsRecord('::1', DnsRecord::AAAA, 300)],
[new Endpoint(new Address('[::1]:50051'))],
];

yield 'mixed A and AAAA records' => [
new Target(Scheme::Dns, [new TargetAddress('myhost', 50_051)], 'myhost:50051'),
new Target('dns', [new TargetAddress('myhost', 50_051)], 'myhost:50051'),
[
new DnsRecord('192.168.0.1', DnsRecord::A, 300),
new DnsRecord('::1', DnsRecord::AAAA, 300),
Expand Down Expand Up @@ -114,7 +113,7 @@ public function testResolveListener(): void

$resolver = new DnsResolver($dnsResolver, minResolveInterval: 0.1, maxResolveInterval: 0.1);
$resolver->resolve(
new Target(Scheme::Dns, [new TargetAddress('myhost', 50_051)], 'myhost:50051'),
new Target('dns', [new TargetAddress('myhost', 50_051)], 'myhost:50051'),
$listener,
$deferredCancellation->getCancellation(),
);
Expand All @@ -125,7 +124,7 @@ public function testResolveListener(): void

public function testResolveStopOnCancellation(): void
{
$target = new Target(Scheme::Dns, [new TargetAddress('myhost', 50_051)], 'myhost:50051');
$target = new Target('dns', [new TargetAddress('myhost', 50_051)], 'myhost:50051');
$deferredCancellation = new DeferredCancellation();

$dnsResolver = self::createStub(AmphpDnsResolver::class);
Expand All @@ -147,7 +146,7 @@ public function testResolveStopOnCancellation(): void

public function testResolveThrows(): void
{
$target = new Target(Scheme::Dns, [new TargetAddress('myhost', 50_051)], 'myhost:50051');
$target = new Target('dns', [new TargetAddress('myhost', 50_051)], 'myhost:50051');
$deferredCancellation = new DeferredCancellation();

$dnsResolver = $this->createMock(AmphpDnsResolver::class);
Expand Down
9 changes: 4 additions & 5 deletions tests/Client/EndpointResolver/StaticResolverTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,6 @@
use Thesis\Grpc\Client\Address;
use Thesis\Grpc\Client\Endpoint;
use Thesis\Grpc\Client\EndpointResolverListener;
use Thesis\Grpc\Client\Scheme;
use Thesis\Grpc\Client\Target;
use Thesis\Grpc\Client\TargetAddress;

Expand Down Expand Up @@ -41,12 +40,12 @@ public function testResolve(Target $target, array $endpoints): void
public static function provideResolveCases(): iterable
{
yield 'ipv4: single address' => [
new Target(Scheme::Ipv4, [new TargetAddress('192.168.0.1', 50_051)], '192.168.0.1:50051'),
new Target('ipv4', [new TargetAddress('192.168.0.1', 50_051)], '192.168.0.1:50051'),
[new Endpoint(new Address('192.168.0.1:50051'))],
];

yield 'ipv4: multiple addresses' => [
new Target(Scheme::Ipv4, [
new Target('ipv4', [
new TargetAddress('192.168.0.1', 50_051),
new TargetAddress('192.168.0.2', 50_052),
], '192.168.0.1:50051,192.168.0.2:50052'),
Expand All @@ -57,12 +56,12 @@ public static function provideResolveCases(): iterable
];

yield 'ipv6: single address' => [
new Target(Scheme::Ipv6, [new TargetAddress('[::1]', 50_051)], '[::1]:50051'),
new Target('ipv6', [new TargetAddress('[::1]', 50_051)], '[::1]:50051'),
[new Endpoint(new Address('[::1]:50051'))],
];

yield 'unix: socket path' => [
new Target(Scheme::Unix, [new TargetAddress('/var/run/grpc.sock', 0)], '///var/run/grpc.sock'),
new Target('unix', [new TargetAddress('/var/run/grpc.sock', 0)], '///var/run/grpc.sock'),
[new Endpoint(new Address('/var/run/grpc.sock'))],
];
}
Expand Down
Loading