Skip to content
Open
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
24 changes: 24 additions & 0 deletions readme.md
Original file line number Diff line number Diff line change
Expand Up @@ -247,6 +247,30 @@ AI assistants use these tools together:
...
```

## Kubernetes RPC discovery

RPC clients can discover services inside Kubernetes by using the `k8s` target
scheme:

```go
c := zrpc.RpcClientConf{
NonBlock: true,
Target: "k8s://dev/demo-rpc:8080",
}
```

The target format is `k8s://<namespace>/<service>:<port>`. If the namespace is
omitted, `default` is used. The resolver runs in-cluster and reads Kubernetes
EndpointSlices, so the service account used by the client must be allowed to
read EndpointSlices in the target namespace:

```yaml
rules:
- apiGroups: ["discovery.k8s.io"]
resources: ["endpointslices"]
verbs: ["list", "watch"]
```

## Benchmark

![benchmark](https://raw.githubusercontent.com/zeromicro/zero-doc/main/doc/images/benchmark.png)
Expand Down
43 changes: 35 additions & 8 deletions zrpc/resolver/internal/kubebuilder.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,17 +11,30 @@ import (
"github.com/zeromicro/go-zero/core/threading"
"github.com/zeromicro/go-zero/zrpc/resolver/internal/kube"
"google.golang.org/grpc/resolver"
apierrors "k8s.io/apimachinery/pkg/api/errors"
v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/informers"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
"k8s.io/client-go/tools/cache"
)

const (
resyncInterval = 5 * time.Minute
serviceSelector = "kubernetes.io/service-name="
)

var (
inClusterConfig = rest.InClusterConfig
newKubeClient = func(config *rest.Config) (kubernetes.Interface, error) {
return kubernetes.NewForConfig(config)
}
addEndpointSliceEventHandler = func(informer cache.SharedIndexInformer, handler cache.ResourceEventHandler) error {
_, err := informer.AddEventHandler(handler)
return err
}
)

type kubeResolver struct {
cc resolver.ClientConn
inf informers.SharedInformerFactory
Expand Down Expand Up @@ -49,12 +62,12 @@ func (b *kubeBuilder) Build(target resolver.Target, cc resolver.ClientConn,
return nil, err
}

config, err := rest.InClusterConfig()
config, err := inClusterConfig()
if err != nil {
return nil, err
}

cs, err := kubernetes.NewForConfig(config)
cs, err := newKubeClient(config)
if err != nil {
return nil, err
}
Expand All @@ -65,7 +78,7 @@ func (b *kubeBuilder) Build(target resolver.Target, cc resolver.ClientConn,
LabelSelector: serviceSelector + svc.Name,
})
if err != nil {
return nil, err
return nil, wrapEndpointSliceListError(svc, err)
}
if len(endpointSlices.Items) == 0 {
return nil, fmt.Errorf("no endpoint slices found for service %s in namespace %s",
Expand Down Expand Up @@ -106,11 +119,9 @@ func (b *kubeBuilder) Build(target resolver.Target, cc resolver.ClientConn,
})
inf := informers.NewSharedInformerFactoryWithOptions(cs, resyncInterval,
informers.WithNamespace(svc.Namespace),
informers.WithTweakListOptions(func(options *v1.ListOptions) {
options.LabelSelector = serviceSelector + svc.Name
}))
informers.WithTweakListOptions(endpointSliceTweakListOptions(svc.Name)))
in := inf.Discovery().V1().EndpointSlices()
_, err = in.Informer().AddEventHandler(handler)
err = addEndpointSliceEventHandler(in.Informer(), handler)
if err != nil {
return nil, err
}
Expand All @@ -122,7 +133,7 @@ func (b *kubeBuilder) Build(target resolver.Target, cc resolver.ClientConn,
LabelSelector: serviceSelector + svc.Name,
})
if err != nil {
return nil, err
return nil, wrapEndpointSliceListError(svc, err)
}

// Aggregate endpoints from all EndpointSlices.
Expand All @@ -144,3 +155,19 @@ func (b *kubeBuilder) Build(target resolver.Target, cc resolver.ClientConn,
func (b *kubeBuilder) Scheme() string {
return KubernetesScheme
}

func endpointSliceTweakListOptions(service string) func(*v1.ListOptions) {
return func(options *v1.ListOptions) {
options.LabelSelector = serviceSelector + service
}
}

func wrapEndpointSliceListError(svc kube.Service, err error) error {
if apierrors.IsForbidden(err) {
return fmt.Errorf("failed to list EndpointSlices for Kubernetes service %q in namespace %q: %w; "+
"the k8s resolver requires list/watch permissions on endpointslices.discovery.k8s.io",
svc.Name, svc.Namespace, err)
}

return err
}
162 changes: 162 additions & 0 deletions zrpc/resolver/internal/kubebuilder_test.go
Original file line number Diff line number Diff line change
@@ -1,12 +1,25 @@
//go:build !no_k8s

package internal

import (
"errors"
"fmt"
"net/url"
"testing"

"github.com/stretchr/testify/assert"
"github.com/zeromicro/go-zero/zrpc/resolver/internal/kube"
"google.golang.org/grpc/resolver"
apierrors "k8s.io/apimachinery/pkg/api/errors"
v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/runtime/schema"
"k8s.io/client-go/kubernetes"
k8sfake "k8s.io/client-go/kubernetes/fake"
"k8s.io/client-go/rest"
ktesting "k8s.io/client-go/testing"
"k8s.io/client-go/tools/cache"
)

func TestKubeBuilder_Scheme(t *testing.T) {
Expand All @@ -32,3 +45,152 @@ func TestKubeBuilder_Build(t *testing.T) {
}, nil, resolver.BuildOptions{})
assert.Error(t, err)
}

func TestKubeBuilder_BuildReturnsKubeClientError(t *testing.T) {
restoreConfig := mockKubeConfig(t)
defer restoreConfig()

sentinel := errors.New("kube client failed")
oldNewClient := newKubeClient
defer func() {
newKubeClient = oldNewClient
}()
newKubeClient = func(*rest.Config) (kubernetes.Interface, error) {
return nil, sentinel
}

var b kubeBuilder
u, err := url.Parse("k8s://dev/demo-rpc:8080")
assert.NoError(t, err)

_, err = b.Build(resolver.Target{URL: *u}, nil, resolver.BuildOptions{})
assert.ErrorIs(t, err, sentinel)
}

func TestKubeBuilder_BuildReturnsAddEventHandlerError(t *testing.T) {
restoreConfig := mockKubeConfig(t)
defer restoreConfig()

sentinel := errors.New("add event handler failed")
restoreHandler := mockEndpointSliceEventHandler(t, sentinel)
defer restoreHandler()

var b kubeBuilder
u, err := url.Parse("k8s://dev/demo-rpc:8080")
assert.NoError(t, err)

_, err = b.Build(resolver.Target{URL: *u}, nil, resolver.BuildOptions{})
assert.ErrorIs(t, err, sentinel)
}

func TestKubeBuilder_BuildWrapsFirstEndpointSliceListError(t *testing.T) {
restore := mockKubeClient(t)
defer restore()

var b kubeBuilder
u, err := url.Parse("k8s://dev/demo-rpc")
assert.NoError(t, err)

_, err = b.Build(resolver.Target{URL: *u}, nil, resolver.BuildOptions{})
assert.Error(t, err)
assert.Contains(t, err.Error(), "failed to list EndpointSlices")
assert.Contains(t, err.Error(), "list/watch permissions")
}

func TestKubeBuilder_BuildWrapsSecondEndpointSliceListError(t *testing.T) {
restore := mockKubeClient(t)
defer restore()

var b kubeBuilder
u, err := url.Parse("k8s://dev/demo-rpc:8080")
assert.NoError(t, err)

_, err = b.Build(resolver.Target{URL: *u}, nil, resolver.BuildOptions{})
assert.Error(t, err)
assert.Contains(t, err.Error(), "failed to list EndpointSlices")
assert.Contains(t, err.Error(), "list/watch permissions")
}

func TestWrapEndpointSliceListError(t *testing.T) {
svc := kube.Service{
Namespace: "dev",
Name: "demo-rpc",
}
forbidden := apierrors.NewForbidden(
schema.GroupResource{Group: "discovery.k8s.io", Resource: "endpointslices"},
svc.Name,
fmt.Errorf("forbidden"),
)

err := wrapEndpointSliceListError(svc, forbidden)
assert.Error(t, err)
assert.Contains(t, err.Error(), "EndpointSlices")
assert.Contains(t, err.Error(), "endpointslices.discovery.k8s.io")
assert.Contains(t, err.Error(), "list/watch")
}

func TestWrapEndpointSliceListErrorOther(t *testing.T) {
svc := kube.Service{
Namespace: "dev",
Name: "demo-rpc",
}
err := errors.New("not forbidden")
assert.Same(t, err, wrapEndpointSliceListError(svc, err))
}

func TestEndpointSliceTweakListOptions(t *testing.T) {
var options v1.ListOptions
endpointSliceTweakListOptions("demo-rpc")(&options)
assert.Equal(t, serviceSelector+"demo-rpc", options.LabelSelector)
}

func mockKubeClient(t *testing.T) func() {
t.Helper()

restoreConfig := mockKubeConfig(t)
oldNewClient := newKubeClient

newKubeClient = func(*rest.Config) (kubernetes.Interface, error) {
cli := k8sfake.NewSimpleClientset()
cli.PrependReactor("list", "endpointslices", func(action ktesting.Action) (bool, runtime.Object, error) {
_ = action
return true, nil, apierrors.NewForbidden(
schema.GroupResource{Group: "discovery.k8s.io", Resource: "endpointslices"},
"demo-rpc",
fmt.Errorf("forbidden"),
)
})
return cli, nil
}

return func() {
restoreConfig()
newKubeClient = oldNewClient
}
}

func mockKubeConfig(t *testing.T) func() {
t.Helper()

oldConfig := inClusterConfig
inClusterConfig = func() (*rest.Config, error) {
return &rest.Config{Host: "https://127.0.0.1"}, nil
}

return func() {
inClusterConfig = oldConfig
}
}

func mockEndpointSliceEventHandler(t *testing.T, err error) func() {
t.Helper()

oldAddHandler := addEndpointSliceEventHandler
addEndpointSliceEventHandler = func(cache.SharedIndexInformer, cache.ResourceEventHandler) error {
return err
}

return func() {
addEndpointSliceEventHandler = oldAddHandler
}
}