mirror of
https://github.com/etcd-io/etcd.git
synced 2024-09-27 06:25:44 +00:00
Merge pull request #14397 from biosvs/backport-grpc-proxy-endpoints-autosync
Backport of pull/14354 to release-3.5
This commit is contained in:
commit
fbb14f91bf
@ -54,15 +54,16 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
var (
|
var (
|
||||||
grpcProxyListenAddr string
|
grpcProxyListenAddr string
|
||||||
grpcProxyMetricsListenAddr string
|
grpcProxyMetricsListenAddr string
|
||||||
grpcProxyEndpoints []string
|
grpcProxyEndpoints []string
|
||||||
grpcProxyDNSCluster string
|
grpcProxyEndpointsAutoSyncInterval time.Duration
|
||||||
grpcProxyDNSClusterServiceName string
|
grpcProxyDNSCluster string
|
||||||
grpcProxyInsecureDiscovery bool
|
grpcProxyDNSClusterServiceName string
|
||||||
grpcProxyDataDir string
|
grpcProxyInsecureDiscovery bool
|
||||||
grpcMaxCallSendMsgSize int
|
grpcProxyDataDir string
|
||||||
grpcMaxCallRecvMsgSize int
|
grpcMaxCallSendMsgSize int
|
||||||
|
grpcMaxCallRecvMsgSize int
|
||||||
|
|
||||||
// tls for connecting to etcd
|
// tls for connecting to etcd
|
||||||
|
|
||||||
@ -130,6 +131,7 @@ func newGRPCProxyStartCommand() *cobra.Command {
|
|||||||
cmd.Flags().StringVar(&grpcProxyMetricsListenAddr, "metrics-addr", "", "listen for endpoint /metrics requests on an additional interface")
|
cmd.Flags().StringVar(&grpcProxyMetricsListenAddr, "metrics-addr", "", "listen for endpoint /metrics requests on an additional interface")
|
||||||
cmd.Flags().BoolVar(&grpcProxyInsecureDiscovery, "insecure-discovery", false, "accept insecure SRV records")
|
cmd.Flags().BoolVar(&grpcProxyInsecureDiscovery, "insecure-discovery", false, "accept insecure SRV records")
|
||||||
cmd.Flags().StringSliceVar(&grpcProxyEndpoints, "endpoints", []string{"127.0.0.1:2379"}, "comma separated etcd cluster endpoints")
|
cmd.Flags().StringSliceVar(&grpcProxyEndpoints, "endpoints", []string{"127.0.0.1:2379"}, "comma separated etcd cluster endpoints")
|
||||||
|
cmd.Flags().DurationVar(&grpcProxyEndpointsAutoSyncInterval, "endpoints-auto-sync-interval", 0, "etcd endpoints auto sync interval (disabled by default)")
|
||||||
cmd.Flags().StringVar(&grpcProxyAdvertiseClientURL, "advertise-client-url", "127.0.0.1:23790", "advertise address to register (must be reachable by client)")
|
cmd.Flags().StringVar(&grpcProxyAdvertiseClientURL, "advertise-client-url", "127.0.0.1:23790", "advertise address to register (must be reachable by client)")
|
||||||
cmd.Flags().StringVar(&grpcProxyResolverPrefix, "resolver-prefix", "", "prefix to use for registering proxy (must be shared with other grpc-proxy members)")
|
cmd.Flags().StringVar(&grpcProxyResolverPrefix, "resolver-prefix", "", "prefix to use for registering proxy (must be shared with other grpc-proxy members)")
|
||||||
cmd.Flags().IntVar(&grpcProxyResolverTTL, "resolver-ttl", 0, "specify TTL, in seconds, when registering proxy endpoints")
|
cmd.Flags().IntVar(&grpcProxyResolverTTL, "resolver-ttl", 0, "specify TTL, in seconds, when registering proxy endpoints")
|
||||||
@ -333,8 +335,9 @@ func newProxyClientCfg(lg *zap.Logger, eps []string, tls *transport.TLSInfo) (*c
|
|||||||
func newClientCfg(lg *zap.Logger, eps []string) (*clientv3.Config, error) {
|
func newClientCfg(lg *zap.Logger, eps []string) (*clientv3.Config, error) {
|
||||||
// set tls if any one tls option set
|
// set tls if any one tls option set
|
||||||
cfg := clientv3.Config{
|
cfg := clientv3.Config{
|
||||||
Endpoints: eps,
|
Endpoints: eps,
|
||||||
DialTimeout: 5 * time.Second,
|
AutoSyncInterval: grpcProxyEndpointsAutoSyncInterval,
|
||||||
|
DialTimeout: 5 * time.Second,
|
||||||
}
|
}
|
||||||
|
|
||||||
if grpcMaxCallSendMsgSize > 0 {
|
if grpcMaxCallSendMsgSize > 0 {
|
||||||
|
177
tests/e2e/etcd_grpcproxy_test.go
Normal file
177
tests/e2e/etcd_grpcproxy_test.go
Normal file
@ -0,0 +1,177 @@
|
|||||||
|
// Copyright 2017 The etcd Authors
|
||||||
|
//
|
||||||
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||||
|
// you may not use this file except in compliance with the License.
|
||||||
|
// You may obtain a copy of the License at
|
||||||
|
//
|
||||||
|
// http://www.apache.org/licenses/LICENSE-2.0
|
||||||
|
//
|
||||||
|
// Unless required by applicable law or agreed to in writing, software
|
||||||
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||||
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||||
|
// See the License for the specific language governing permissions and
|
||||||
|
// limitations under the License.
|
||||||
|
|
||||||
|
package e2e
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"fmt"
|
||||||
|
"io"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/stretchr/testify/require"
|
||||||
|
"go.etcd.io/etcd/api/v3/etcdserverpb"
|
||||||
|
"go.etcd.io/etcd/api/v3/v3rpc/rpctypes"
|
||||||
|
"go.etcd.io/etcd/pkg/v3/expect"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestGrpcProxyAutoSync(t *testing.T) {
|
||||||
|
skipInShortMode(t)
|
||||||
|
|
||||||
|
var (
|
||||||
|
node1Name = "node1"
|
||||||
|
node1ClientURL = "http://localhost:12379"
|
||||||
|
node1PeerURL = "http://localhost:12380"
|
||||||
|
|
||||||
|
node2Name = "node2"
|
||||||
|
node2ClientURL = "http://localhost:22379"
|
||||||
|
node2PeerURL = "http://localhost:22380"
|
||||||
|
|
||||||
|
proxyClientURL = "127.0.0.1:32379"
|
||||||
|
|
||||||
|
autoSyncInterval = 1 * time.Second
|
||||||
|
)
|
||||||
|
|
||||||
|
// Run cluster of one node
|
||||||
|
proc1, err := runEtcdNode(
|
||||||
|
node1Name, t.TempDir(),
|
||||||
|
node1ClientURL, node1PeerURL,
|
||||||
|
"new", fmt.Sprintf("%s=%s", node1Name, node1PeerURL),
|
||||||
|
)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
// Run grpc-proxy instance
|
||||||
|
proxyProc, err := spawnCmd([]string{binDir + "/etcd", "grpc-proxy", "start",
|
||||||
|
"--advertise-client-url", proxyClientURL, "--listen-addr", proxyClientURL,
|
||||||
|
"--endpoints", node1ClientURL,
|
||||||
|
"--endpoints-auto-sync-interval", autoSyncInterval.String(),
|
||||||
|
}, nil)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
err = spawnWithExpect([]string{ctlBinPath, "--endpoints", proxyClientURL, "put", "k1", "v1"}, "OK")
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
err = spawnWithExpect([]string{ctlBinPath, "--endpoints", node1ClientURL, "member", "add", node2Name, "--peer-urls", node2PeerURL}, "added")
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
// Run new member
|
||||||
|
proc2, err := runEtcdNode(
|
||||||
|
node2Name, t.TempDir(),
|
||||||
|
node2ClientURL, node2PeerURL,
|
||||||
|
"existing", fmt.Sprintf("%s=%s,%s=%s", node1Name, node1PeerURL, node2Name, node2PeerURL),
|
||||||
|
)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
// Wait for auto sync of endpoints
|
||||||
|
_, err = proxyProc.Expect(strings.Replace(node2ClientURL, "http://", "", 1))
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
memberList, err := getMemberListFromEndpoint(node1ClientURL)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
node1MemberID, err := findMemberIDByEndpoint(memberList.Members, node1ClientURL)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
node2MemberID, err := findMemberIDByEndpoint(memberList.Members, node2ClientURL)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
// Remove node1 from member list and stop this node
|
||||||
|
|
||||||
|
// Second node could be not ready yet
|
||||||
|
for i := 0; i < 10; i++ {
|
||||||
|
err = spawnWithExpect([]string{ctlBinPath, "--endpoints", node2ClientURL, "member", "remove", fmt.Sprintf("%x", node1MemberID)}, "removed")
|
||||||
|
if err != nil && strings.Contains(err.Error(), rpctypes.ErrGRPCUnhealthy.Error()) {
|
||||||
|
time.Sleep(500 * time.Millisecond)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
break
|
||||||
|
}
|
||||||
|
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
_, err = proc1.Expect("the member has been permanently removed from the cluster")
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.NoError(t, proc1.Stop())
|
||||||
|
|
||||||
|
_, err = proc2.Expect(fmt.Sprintf("%x became leader", node2MemberID))
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
for i := 0; i < 10; i++ {
|
||||||
|
err = spawnWithExpect([]string{ctlBinPath, "--endpoints", proxyClientURL, "get", "k1"}, "v1")
|
||||||
|
if err != nil && (strings.Contains(err.Error(), rpctypes.ErrGRPCLeaderChanged.Error()) ||
|
||||||
|
strings.Contains(err.Error(), context.DeadlineExceeded.Error())) {
|
||||||
|
time.Sleep(500 * time.Millisecond)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
break
|
||||||
|
}
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
require.NoError(t, proc2.Stop())
|
||||||
|
require.NoError(t, proxyProc.Stop())
|
||||||
|
}
|
||||||
|
|
||||||
|
func runEtcdNode(name, dataDir, clientURL, peerURL, clusterState, initialCluster string) (*expect.ExpectProcess, error) {
|
||||||
|
proc, err := spawnCmd([]string{binPath,
|
||||||
|
"--name", name,
|
||||||
|
"--data-dir", dataDir,
|
||||||
|
"--listen-client-urls", clientURL, "--advertise-client-urls", clientURL,
|
||||||
|
"--listen-peer-urls", peerURL, "--initial-advertise-peer-urls", peerURL,
|
||||||
|
"--initial-cluster-token", "etcd-cluster",
|
||||||
|
"--initial-cluster-state", clusterState,
|
||||||
|
"--initial-cluster", initialCluster,
|
||||||
|
}, nil)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
_, err = proc.Expect("ready to serve client requests")
|
||||||
|
|
||||||
|
return proc, err
|
||||||
|
}
|
||||||
|
|
||||||
|
func findMemberIDByEndpoint(members []*etcdserverpb.Member, endpoint string) (uint64, error) {
|
||||||
|
for _, m := range members {
|
||||||
|
if m.ClientURLs[0] == endpoint {
|
||||||
|
return m.ID, nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return 0, fmt.Errorf("member not found")
|
||||||
|
}
|
||||||
|
|
||||||
|
func getMemberListFromEndpoint(endpoint string) (etcdserverpb.MemberListResponse, error) {
|
||||||
|
proc, err := spawnCmd([]string{ctlBinPath, "--endpoints", endpoint, "member", "list", "--write-out", "json"}, nil)
|
||||||
|
if err != nil {
|
||||||
|
return etcdserverpb.MemberListResponse{}, err
|
||||||
|
}
|
||||||
|
var txt string
|
||||||
|
txt, err = proc.Expect("members")
|
||||||
|
if err != nil {
|
||||||
|
return etcdserverpb.MemberListResponse{}, err
|
||||||
|
}
|
||||||
|
if err = proc.Close(); err != nil {
|
||||||
|
return etcdserverpb.MemberListResponse{}, err
|
||||||
|
}
|
||||||
|
|
||||||
|
resp := etcdserverpb.MemberListResponse{}
|
||||||
|
dec := json.NewDecoder(strings.NewReader(txt))
|
||||||
|
if err := dec.Decode(&resp); err == io.EOF {
|
||||||
|
return etcdserverpb.MemberListResponse{}, err
|
||||||
|
}
|
||||||
|
return resp, nil
|
||||||
|
}
|
Loading…
x
Reference in New Issue
Block a user