2016-05-13 06:49:40 +03:00
|
|
|
// Copyright 2015 The etcd Authors
|
2015-02-12 01:03:14 +03:00
|
|
|
//
|
|
|
|
// 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 etcdserver
|
|
|
|
|
|
|
|
import (
|
2019-05-01 01:51:36 +03:00
|
|
|
"context"
|
2015-02-12 01:03:14 +03:00
|
|
|
"encoding/json"
|
|
|
|
"fmt"
|
2021-10-27 19:00:49 +03:00
|
|
|
"io"
|
2015-02-12 01:03:14 +03:00
|
|
|
"net/http"
|
|
|
|
"sort"
|
2020-10-12 20:50:35 +03:00
|
|
|
"strconv"
|
2019-05-01 01:51:36 +03:00
|
|
|
"strings"
|
2015-02-12 01:03:14 +03:00
|
|
|
"time"
|
|
|
|
|
2020-10-05 18:21:03 +03:00
|
|
|
"go.etcd.io/etcd/api/v3/version"
|
2021-04-05 23:31:07 +03:00
|
|
|
"go.etcd.io/etcd/client/pkg/v3/types"
|
2020-10-21 00:09:35 +03:00
|
|
|
"go.etcd.io/etcd/server/v3/etcdserver/api/membership"
|
2018-01-29 21:00:20 +03:00
|
|
|
|
2016-03-23 03:10:28 +03:00
|
|
|
"github.com/coreos/go-semver/semver"
|
2018-04-16 08:40:44 +03:00
|
|
|
"go.uber.org/zap"
|
2015-02-12 01:03:14 +03:00
|
|
|
)
|
|
|
|
|
2015-02-12 01:18:10 +03:00
|
|
|
// isMemberBootstrapped tries to check if the given member has been bootstrapped
|
2015-02-12 01:03:14 +03:00
|
|
|
// in the given cluster.
|
2018-04-16 08:40:44 +03:00
|
|
|
func isMemberBootstrapped(lg *zap.Logger, cl *membership.RaftCluster, member string, rt http.RoundTripper, timeout time.Duration) bool {
|
|
|
|
rcl, err := getClusterFromRemotePeers(lg, getRemotePeerURLs(cl, member), timeout, false, rt)
|
2015-02-12 01:03:14 +03:00
|
|
|
if err != nil {
|
|
|
|
return false
|
|
|
|
}
|
|
|
|
id := cl.MemberByName(member).ID
|
|
|
|
m := rcl.Member(id)
|
|
|
|
if m == nil {
|
|
|
|
return false
|
|
|
|
}
|
|
|
|
if len(m.ClientURLs) > 0 {
|
|
|
|
return true
|
|
|
|
}
|
|
|
|
return false
|
|
|
|
}
|
|
|
|
|
2015-02-14 06:05:29 +03:00
|
|
|
// GetClusterFromRemotePeers takes a set of URLs representing etcd peers, and
|
2015-02-12 01:03:14 +03:00
|
|
|
// attempts to construct a Cluster by accessing the members endpoint on one of
|
|
|
|
// these URLs. The first URL to provide a response is used. If no URLs provide
|
|
|
|
// a response, or a Cluster cannot be successfully created from a received
|
|
|
|
// response, an error is returned.
|
2015-09-01 01:12:58 +03:00
|
|
|
// Each request has a 10-second timeout. Because the upper limit of TTL is 5s,
|
|
|
|
// 10 second is enough for building connection and finishing request.
|
2018-04-16 08:40:44 +03:00
|
|
|
func GetClusterFromRemotePeers(lg *zap.Logger, urls []string, rt http.RoundTripper) (*membership.RaftCluster, error) {
|
|
|
|
return getClusterFromRemotePeers(lg, urls, 10*time.Second, true, rt)
|
2015-02-12 01:03:14 +03:00
|
|
|
}
|
|
|
|
|
|
|
|
// If logerr is true, it prints out more error messages.
|
2018-04-16 08:40:44 +03:00
|
|
|
func getClusterFromRemotePeers(lg *zap.Logger, urls []string, timeout time.Duration, logerr bool, rt http.RoundTripper) (*membership.RaftCluster, error) {
|
2020-02-11 19:54:14 +03:00
|
|
|
if lg == nil {
|
|
|
|
lg = zap.NewNop()
|
|
|
|
}
|
2015-02-12 01:03:14 +03:00
|
|
|
cc := &http.Client{
|
2015-11-04 21:49:42 +03:00
|
|
|
Transport: rt,
|
2015-09-01 01:12:58 +03:00
|
|
|
Timeout: timeout,
|
2015-02-12 01:03:14 +03:00
|
|
|
}
|
|
|
|
for _, u := range urls {
|
2018-04-16 08:40:44 +03:00
|
|
|
addr := u + "/members"
|
|
|
|
resp, err := cc.Get(addr)
|
2015-02-12 01:03:14 +03:00
|
|
|
if err != nil {
|
|
|
|
if logerr {
|
2020-02-11 19:54:14 +03:00
|
|
|
lg.Warn("failed to get cluster response", zap.String("address", addr), zap.Error(err))
|
2015-02-12 01:03:14 +03:00
|
|
|
}
|
|
|
|
continue
|
|
|
|
}
|
2021-10-27 19:00:49 +03:00
|
|
|
b, err := io.ReadAll(resp.Body)
|
2016-04-18 16:16:20 +03:00
|
|
|
resp.Body.Close()
|
2015-02-12 01:03:14 +03:00
|
|
|
if err != nil {
|
|
|
|
if logerr {
|
2020-02-11 19:54:14 +03:00
|
|
|
lg.Warn("failed to read body of cluster response", zap.String("address", addr), zap.Error(err))
|
2015-02-12 01:03:14 +03:00
|
|
|
}
|
|
|
|
continue
|
|
|
|
}
|
2016-04-07 22:02:37 +03:00
|
|
|
var membs []*membership.Member
|
2015-12-12 15:25:27 +03:00
|
|
|
if err = json.Unmarshal(b, &membs); err != nil {
|
2015-02-12 01:03:14 +03:00
|
|
|
if logerr {
|
2020-02-11 19:54:14 +03:00
|
|
|
lg.Warn("failed to unmarshal cluster response", zap.String("address", addr), zap.Error(err))
|
2015-02-12 01:03:14 +03:00
|
|
|
}
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
id, err := types.IDFromString(resp.Header.Get("X-Etcd-Cluster-ID"))
|
|
|
|
if err != nil {
|
|
|
|
if logerr {
|
2020-02-11 19:54:14 +03:00
|
|
|
lg.Warn(
|
|
|
|
"failed to parse cluster ID",
|
|
|
|
zap.String("address", addr),
|
|
|
|
zap.String("header", resp.Header.Get("X-Etcd-Cluster-ID")),
|
|
|
|
zap.Error(err),
|
|
|
|
)
|
2015-02-12 01:03:14 +03:00
|
|
|
}
|
|
|
|
continue
|
|
|
|
}
|
2016-08-10 16:53:19 +03:00
|
|
|
|
|
|
|
// check the length of membership members
|
|
|
|
// if the membership members are present then prepare and return raft cluster
|
|
|
|
// if membership members are not present then the raft cluster formed will be
|
|
|
|
// an invalid empty cluster hence return failed to get raft cluster member(s) from the given urls error
|
|
|
|
if len(membs) > 0 {
|
2021-04-02 18:22:19 +03:00
|
|
|
return membership.NewClusterFromMembers(lg, id, membs), nil
|
2016-08-10 16:53:19 +03:00
|
|
|
}
|
2018-04-16 08:40:44 +03:00
|
|
|
return nil, fmt.Errorf("failed to get raft cluster member(s) from the given URLs")
|
2015-02-12 01:03:14 +03:00
|
|
|
}
|
2018-04-16 08:40:44 +03:00
|
|
|
return nil, fmt.Errorf("could not retrieve cluster information from the given URLs")
|
2015-02-12 01:03:14 +03:00
|
|
|
}
|
|
|
|
|
2015-02-14 05:56:45 +03:00
|
|
|
// getRemotePeerURLs returns peer urls of remote members in the cluster. The
|
2015-02-12 01:03:14 +03:00
|
|
|
// returned list is sorted in ascending lexicographical order.
|
2016-04-07 22:02:37 +03:00
|
|
|
func getRemotePeerURLs(cl *membership.RaftCluster, local string) []string {
|
2015-02-12 01:03:14 +03:00
|
|
|
us := make([]string, 0)
|
|
|
|
for _, m := range cl.Members() {
|
2015-02-14 05:56:45 +03:00
|
|
|
if m.Name == local {
|
2015-02-12 01:03:14 +03:00
|
|
|
continue
|
|
|
|
}
|
|
|
|
us = append(us, m.PeerURLs...)
|
|
|
|
}
|
|
|
|
sort.Strings(us)
|
|
|
|
return us
|
|
|
|
}
|
*: add cluster version and cluster version detection.
Cluster version is the min major.minor of all members in
the etcd cluster. Cluster version is set to the min version
that a etcd member is compatible with when first bootstrapp.
During a rolling upgrades, the cluster version will be updated
automatically.
For example:
```
Cluster [a:1, b:1 ,c:1] -> clusterVersion 1
update a -> 2, b -> 2
after a detection
Cluster [a:2, b:2 ,c:1] -> clusterVersion 1, since c is still 1
update c -> 2
after a detection
Cluster [a:2, b:2 ,c:2] -> clusterVersion 2
```
The API/raft component can utilize clusterVersion to determine if
it can accept a client request or a raft RPC.
We choose polling rather than pushing since we want to use the same
logic for cluster version detection and (TODO) cluster version checking.
Before a member actually joins a etcd cluster, it should check the version
of the cluster. Push does not work since the other members cannot push
version info to it before it actually joins. Moreover, we do not want our
raft RPC system (which is doing the heartbeat pushing) to coordinate cluster version.
2015-04-29 20:56:34 +03:00
|
|
|
|
2021-10-04 16:53:54 +03:00
|
|
|
// getMembersVersions returns the versions of the members in the given cluster.
|
*: add cluster version and cluster version detection.
Cluster version is the min major.minor of all members in
the etcd cluster. Cluster version is set to the min version
that a etcd member is compatible with when first bootstrapp.
During a rolling upgrades, the cluster version will be updated
automatically.
For example:
```
Cluster [a:1, b:1 ,c:1] -> clusterVersion 1
update a -> 2, b -> 2
after a detection
Cluster [a:2, b:2 ,c:1] -> clusterVersion 1, since c is still 1
update c -> 2
after a detection
Cluster [a:2, b:2 ,c:2] -> clusterVersion 2
```
The API/raft component can utilize clusterVersion to determine if
it can accept a client request or a raft RPC.
We choose polling rather than pushing since we want to use the same
logic for cluster version detection and (TODO) cluster version checking.
Before a member actually joins a etcd cluster, it should check the version
of the cluster. Push does not work since the other members cannot push
version info to it before it actually joins. Moreover, we do not want our
raft RPC system (which is doing the heartbeat pushing) to coordinate cluster version.
2015-04-29 20:56:34 +03:00
|
|
|
// The key of the returned map is the member's ID. The value of the returned map
|
2015-05-14 03:04:46 +03:00
|
|
|
// is the semver versions string, including server and cluster.
|
|
|
|
// If it fails to get the version of a member, the key will be nil.
|
2022-03-01 06:11:09 +03:00
|
|
|
func getMembersVersions(lg *zap.Logger, cl *membership.RaftCluster, local types.ID, rt http.RoundTripper, timeout time.Duration) map[string]*version.Versions {
|
*: add cluster version and cluster version detection.
Cluster version is the min major.minor of all members in
the etcd cluster. Cluster version is set to the min version
that a etcd member is compatible with when first bootstrapp.
During a rolling upgrades, the cluster version will be updated
automatically.
For example:
```
Cluster [a:1, b:1 ,c:1] -> clusterVersion 1
update a -> 2, b -> 2
after a detection
Cluster [a:2, b:2 ,c:1] -> clusterVersion 1, since c is still 1
update c -> 2
after a detection
Cluster [a:2, b:2 ,c:2] -> clusterVersion 2
```
The API/raft component can utilize clusterVersion to determine if
it can accept a client request or a raft RPC.
We choose polling rather than pushing since we want to use the same
logic for cluster version detection and (TODO) cluster version checking.
Before a member actually joins a etcd cluster, it should check the version
of the cluster. Push does not work since the other members cannot push
version info to it before it actually joins. Moreover, we do not want our
raft RPC system (which is doing the heartbeat pushing) to coordinate cluster version.
2015-04-29 20:56:34 +03:00
|
|
|
members := cl.Members()
|
2015-05-14 03:04:46 +03:00
|
|
|
vers := make(map[string]*version.Versions)
|
*: add cluster version and cluster version detection.
Cluster version is the min major.minor of all members in
the etcd cluster. Cluster version is set to the min version
that a etcd member is compatible with when first bootstrapp.
During a rolling upgrades, the cluster version will be updated
automatically.
For example:
```
Cluster [a:1, b:1 ,c:1] -> clusterVersion 1
update a -> 2, b -> 2
after a detection
Cluster [a:2, b:2 ,c:1] -> clusterVersion 1, since c is still 1
update c -> 2
after a detection
Cluster [a:2, b:2 ,c:2] -> clusterVersion 2
```
The API/raft component can utilize clusterVersion to determine if
it can accept a client request or a raft RPC.
We choose polling rather than pushing since we want to use the same
logic for cluster version detection and (TODO) cluster version checking.
Before a member actually joins a etcd cluster, it should check the version
of the cluster. Push does not work since the other members cannot push
version info to it before it actually joins. Moreover, we do not want our
raft RPC system (which is doing the heartbeat pushing) to coordinate cluster version.
2015-04-29 20:56:34 +03:00
|
|
|
for _, m := range members {
|
2015-05-14 03:19:32 +03:00
|
|
|
if m.ID == local {
|
2015-05-14 17:57:25 +03:00
|
|
|
cv := "not_decided"
|
|
|
|
if cl.Version() != nil {
|
|
|
|
cv = cl.Version().String()
|
|
|
|
}
|
|
|
|
vers[m.ID.String()] = &version.Versions{Server: version.Version, Cluster: cv}
|
2015-05-14 03:19:32 +03:00
|
|
|
continue
|
|
|
|
}
|
2022-03-01 06:11:09 +03:00
|
|
|
ver, err := getVersion(lg, m, rt, timeout)
|
*: add cluster version and cluster version detection.
Cluster version is the min major.minor of all members in
the etcd cluster. Cluster version is set to the min version
that a etcd member is compatible with when first bootstrapp.
During a rolling upgrades, the cluster version will be updated
automatically.
For example:
```
Cluster [a:1, b:1 ,c:1] -> clusterVersion 1
update a -> 2, b -> 2
after a detection
Cluster [a:2, b:2 ,c:1] -> clusterVersion 1, since c is still 1
update c -> 2
after a detection
Cluster [a:2, b:2 ,c:2] -> clusterVersion 2
```
The API/raft component can utilize clusterVersion to determine if
it can accept a client request or a raft RPC.
We choose polling rather than pushing since we want to use the same
logic for cluster version detection and (TODO) cluster version checking.
Before a member actually joins a etcd cluster, it should check the version
of the cluster. Push does not work since the other members cannot push
version info to it before it actually joins. Moreover, we do not want our
raft RPC system (which is doing the heartbeat pushing) to coordinate cluster version.
2015-04-29 20:56:34 +03:00
|
|
|
if err != nil {
|
2020-02-11 19:54:14 +03:00
|
|
|
lg.Warn("failed to get version", zap.String("remote-member-id", m.ID.String()), zap.Error(err))
|
2015-05-14 03:04:46 +03:00
|
|
|
vers[m.ID.String()] = nil
|
*: add cluster version and cluster version detection.
Cluster version is the min major.minor of all members in
the etcd cluster. Cluster version is set to the min version
that a etcd member is compatible with when first bootstrapp.
During a rolling upgrades, the cluster version will be updated
automatically.
For example:
```
Cluster [a:1, b:1 ,c:1] -> clusterVersion 1
update a -> 2, b -> 2
after a detection
Cluster [a:2, b:2 ,c:1] -> clusterVersion 1, since c is still 1
update c -> 2
after a detection
Cluster [a:2, b:2 ,c:2] -> clusterVersion 2
```
The API/raft component can utilize clusterVersion to determine if
it can accept a client request or a raft RPC.
We choose polling rather than pushing since we want to use the same
logic for cluster version detection and (TODO) cluster version checking.
Before a member actually joins a etcd cluster, it should check the version
of the cluster. Push does not work since the other members cannot push
version info to it before it actually joins. Moreover, we do not want our
raft RPC system (which is doing the heartbeat pushing) to coordinate cluster version.
2015-04-29 20:56:34 +03:00
|
|
|
} else {
|
|
|
|
vers[m.ID.String()] = ver
|
|
|
|
}
|
|
|
|
}
|
|
|
|
return vers
|
|
|
|
}
|
|
|
|
|
2020-06-24 21:06:29 +03:00
|
|
|
// allowedVersionRange decides the available version range of the cluster that local server can join in;
|
|
|
|
// if the downgrade enabled status is true, the version window is [oneMinorHigher, oneMinorHigher]
|
|
|
|
// if the downgrade is not enabled, the version window is [MinClusterVersion, localVersion]
|
|
|
|
func allowedVersionRange(downgradeEnabled bool) (minV *semver.Version, maxV *semver.Version) {
|
|
|
|
minV = semver.Must(semver.NewVersion(version.MinClusterVersion))
|
|
|
|
maxV = semver.Must(semver.NewVersion(version.Version))
|
|
|
|
maxV = &semver.Version{Major: maxV.Major, Minor: maxV.Minor}
|
|
|
|
|
|
|
|
if downgradeEnabled {
|
|
|
|
// Todo: handle the case that downgrading from higher major version(e.g. downgrade from v4.0 to v3.x)
|
|
|
|
maxV.Minor = maxV.Minor + 1
|
|
|
|
minV = &semver.Version{Major: maxV.Major, Minor: maxV.Minor}
|
|
|
|
}
|
|
|
|
return minV, maxV
|
|
|
|
}
|
|
|
|
|
2016-01-08 11:21:19 +03:00
|
|
|
// isCompatibleWithCluster return true if the local member has a compatible version with
|
2015-05-14 17:57:25 +03:00
|
|
|
// the current running cluster.
|
2016-01-08 11:21:19 +03:00
|
|
|
// The version is considered as compatible when at least one of the other members in the cluster has a
|
2020-06-24 21:06:29 +03:00
|
|
|
// cluster version in the range of [MinV, MaxV] and no known members has a cluster version
|
2015-05-14 17:57:25 +03:00
|
|
|
// out of the range.
|
|
|
|
// We set this rule since when the local member joins, another member might be offline.
|
2022-03-01 06:11:09 +03:00
|
|
|
func isCompatibleWithCluster(lg *zap.Logger, cl *membership.RaftCluster, local types.ID, rt http.RoundTripper, timeout time.Duration) bool {
|
|
|
|
vers := getMembersVersions(lg, cl, local, rt, timeout)
|
|
|
|
minV, maxV := allowedVersionRange(getDowngradeEnabledFromRemotePeers(lg, cl, local, rt, timeout))
|
2018-04-16 08:40:44 +03:00
|
|
|
return isCompatibleWithVers(lg, vers, local, minV, maxV)
|
2015-05-14 17:57:25 +03:00
|
|
|
}
|
|
|
|
|
2018-04-16 08:40:44 +03:00
|
|
|
func isCompatibleWithVers(lg *zap.Logger, vers map[string]*version.Versions, local types.ID, minV, maxV *semver.Version) bool {
|
2015-05-14 17:57:25 +03:00
|
|
|
var ok bool
|
|
|
|
for id, v := range vers {
|
2016-01-08 11:21:19 +03:00
|
|
|
// ignore comparison with local version
|
2015-05-14 17:57:25 +03:00
|
|
|
if id == local.String() {
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
if v == nil {
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
clusterv, err := semver.NewVersion(v.Cluster)
|
|
|
|
if err != nil {
|
2020-02-11 19:54:14 +03:00
|
|
|
lg.Warn(
|
|
|
|
"failed to parse cluster version of remote member",
|
|
|
|
zap.String("remote-member-id", id),
|
|
|
|
zap.String("remote-member-cluster-version", v.Cluster),
|
|
|
|
zap.Error(err),
|
|
|
|
)
|
2015-05-14 17:57:25 +03:00
|
|
|
continue
|
|
|
|
}
|
|
|
|
if clusterv.LessThan(*minV) {
|
2020-02-11 19:54:14 +03:00
|
|
|
lg.Warn(
|
|
|
|
"cluster version of remote member is not compatible; too low",
|
|
|
|
zap.String("remote-member-id", id),
|
|
|
|
zap.String("remote-member-cluster-version", clusterv.String()),
|
|
|
|
zap.String("minimum-cluster-version-supported", minV.String()),
|
|
|
|
)
|
2015-05-14 17:57:25 +03:00
|
|
|
return false
|
|
|
|
}
|
|
|
|
if maxV.LessThan(*clusterv) {
|
2020-02-11 19:54:14 +03:00
|
|
|
lg.Warn(
|
|
|
|
"cluster version of remote member is not compatible; too high",
|
|
|
|
zap.String("remote-member-id", id),
|
|
|
|
zap.String("remote-member-cluster-version", clusterv.String()),
|
|
|
|
zap.String("minimum-cluster-version-supported", minV.String()),
|
|
|
|
)
|
2015-05-14 17:57:25 +03:00
|
|
|
return false
|
|
|
|
}
|
|
|
|
ok = true
|
|
|
|
}
|
|
|
|
return ok
|
|
|
|
}
|
|
|
|
|
2015-05-14 03:04:46 +03:00
|
|
|
// getVersion returns the Versions of the given member via its
|
2015-05-08 21:10:12 +03:00
|
|
|
// peerURLs. Returns the last error if it fails to get the version.
|
2022-03-01 06:11:09 +03:00
|
|
|
func getVersion(lg *zap.Logger, m *membership.Member, rt http.RoundTripper, timeout time.Duration) (*version.Versions, error) {
|
2015-05-08 21:10:12 +03:00
|
|
|
cc := &http.Client{
|
2015-11-04 21:49:42 +03:00
|
|
|
Transport: rt,
|
2022-03-01 06:11:09 +03:00
|
|
|
Timeout: timeout,
|
2015-05-08 21:10:12 +03:00
|
|
|
}
|
|
|
|
var (
|
|
|
|
err error
|
|
|
|
resp *http.Response
|
|
|
|
)
|
|
|
|
|
|
|
|
for _, u := range m.PeerURLs {
|
2018-04-16 08:40:44 +03:00
|
|
|
addr := u + "/version"
|
|
|
|
resp, err = cc.Get(addr)
|
2015-05-08 21:10:12 +03:00
|
|
|
if err != nil {
|
2020-02-11 19:54:14 +03:00
|
|
|
lg.Warn(
|
|
|
|
"failed to reach the peer URL",
|
|
|
|
zap.String("address", addr),
|
|
|
|
zap.String("remote-member-id", m.ID.String()),
|
|
|
|
zap.Error(err),
|
|
|
|
)
|
2015-05-08 21:10:12 +03:00
|
|
|
continue
|
|
|
|
}
|
2015-08-27 23:24:47 +03:00
|
|
|
var b []byte
|
2021-10-27 19:00:49 +03:00
|
|
|
b, err = io.ReadAll(resp.Body)
|
2015-05-08 21:10:12 +03:00
|
|
|
resp.Body.Close()
|
|
|
|
if err != nil {
|
2020-02-11 19:54:14 +03:00
|
|
|
lg.Warn(
|
|
|
|
"failed to read body of response",
|
|
|
|
zap.String("address", addr),
|
|
|
|
zap.String("remote-member-id", m.ID.String()),
|
|
|
|
zap.Error(err),
|
|
|
|
)
|
2015-05-08 21:10:12 +03:00
|
|
|
continue
|
|
|
|
}
|
|
|
|
var vers version.Versions
|
2015-12-12 15:25:27 +03:00
|
|
|
if err = json.Unmarshal(b, &vers); err != nil {
|
2020-02-11 19:54:14 +03:00
|
|
|
lg.Warn(
|
|
|
|
"failed to unmarshal response",
|
|
|
|
zap.String("address", addr),
|
|
|
|
zap.String("remote-member-id", m.ID.String()),
|
|
|
|
zap.Error(err),
|
|
|
|
)
|
2015-05-08 21:10:12 +03:00
|
|
|
continue
|
|
|
|
}
|
2015-05-14 03:04:46 +03:00
|
|
|
return &vers, nil
|
2015-05-08 21:10:12 +03:00
|
|
|
}
|
2015-05-14 03:04:46 +03:00
|
|
|
return nil, err
|
2015-05-08 21:10:12 +03:00
|
|
|
}
|
2019-05-01 01:51:36 +03:00
|
|
|
|
|
|
|
func promoteMemberHTTP(ctx context.Context, url string, id uint64, peerRt http.RoundTripper) ([]*membership.Member, error) {
|
|
|
|
cc := &http.Client{Transport: peerRt}
|
|
|
|
// TODO: refactor member http handler code
|
|
|
|
// cannot import etcdhttp, so manually construct url
|
|
|
|
requestUrl := url + "/members/promote/" + fmt.Sprintf("%d", id)
|
|
|
|
req, err := http.NewRequest("POST", requestUrl, nil)
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
req = req.WithContext(ctx)
|
|
|
|
resp, err := cc.Do(req)
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
defer resp.Body.Close()
|
2021-10-27 19:00:49 +03:00
|
|
|
b, err := io.ReadAll(resp.Body)
|
2019-05-01 01:51:36 +03:00
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
if resp.StatusCode == http.StatusRequestTimeout {
|
|
|
|
return nil, ErrTimeout
|
|
|
|
}
|
|
|
|
if resp.StatusCode == http.StatusPreconditionFailed {
|
|
|
|
// both ErrMemberNotLearner and ErrLearnerNotReady have same http status code
|
2019-05-08 01:46:24 +03:00
|
|
|
if strings.Contains(string(b), ErrLearnerNotReady.Error()) {
|
|
|
|
return nil, ErrLearnerNotReady
|
2019-05-01 01:51:36 +03:00
|
|
|
}
|
|
|
|
if strings.Contains(string(b), membership.ErrMemberNotLearner.Error()) {
|
|
|
|
return nil, membership.ErrMemberNotLearner
|
|
|
|
}
|
|
|
|
return nil, fmt.Errorf("member promote: unknown error(%s)", string(b))
|
|
|
|
}
|
|
|
|
if resp.StatusCode == http.StatusNotFound {
|
|
|
|
return nil, membership.ErrIDNotFound
|
|
|
|
}
|
|
|
|
|
2019-05-08 01:49:02 +03:00
|
|
|
if resp.StatusCode != http.StatusOK { // all other types of errors
|
|
|
|
return nil, fmt.Errorf("member promote: unknown error(%s)", string(b))
|
|
|
|
}
|
|
|
|
|
2019-05-01 01:51:36 +03:00
|
|
|
var membs []*membership.Member
|
|
|
|
if err := json.Unmarshal(b, &membs); err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
return membs, nil
|
|
|
|
}
|
2020-05-16 00:45:00 +03:00
|
|
|
|
2020-06-24 21:06:29 +03:00
|
|
|
// getDowngradeEnabledFromRemotePeers will get the downgrade enabled status of the cluster.
|
2022-03-01 06:11:09 +03:00
|
|
|
func getDowngradeEnabledFromRemotePeers(lg *zap.Logger, cl *membership.RaftCluster, local types.ID, rt http.RoundTripper, timeout time.Duration) bool {
|
2020-10-12 20:50:35 +03:00
|
|
|
members := cl.Members()
|
|
|
|
|
|
|
|
for _, m := range members {
|
|
|
|
if m.ID == local {
|
|
|
|
continue
|
|
|
|
}
|
2022-03-01 06:11:09 +03:00
|
|
|
enable, err := getDowngradeEnabled(lg, m, rt, timeout)
|
2020-10-12 20:50:35 +03:00
|
|
|
if err != nil {
|
|
|
|
lg.Warn("failed to get downgrade enabled status", zap.String("remote-member-id", m.ID.String()), zap.Error(err))
|
|
|
|
} else {
|
|
|
|
// Since the "/downgrade/enabled" serves linearized data,
|
|
|
|
// this function can return once it gets a non-error response from the endpoint.
|
|
|
|
return enable
|
|
|
|
}
|
|
|
|
}
|
2020-06-24 21:06:29 +03:00
|
|
|
return false
|
|
|
|
}
|
|
|
|
|
2020-10-12 20:50:35 +03:00
|
|
|
// getDowngradeEnabled returns the downgrade enabled status of the given member
|
|
|
|
// via its peerURLs. Returns the last error if it fails to get it.
|
2022-03-01 06:11:09 +03:00
|
|
|
func getDowngradeEnabled(lg *zap.Logger, m *membership.Member, rt http.RoundTripper, timeout time.Duration) (bool, error) {
|
2020-10-12 20:50:35 +03:00
|
|
|
cc := &http.Client{
|
|
|
|
Transport: rt,
|
2022-03-01 06:11:09 +03:00
|
|
|
Timeout: timeout,
|
2020-10-12 20:50:35 +03:00
|
|
|
}
|
|
|
|
var (
|
|
|
|
err error
|
|
|
|
resp *http.Response
|
|
|
|
)
|
|
|
|
|
|
|
|
for _, u := range m.PeerURLs {
|
|
|
|
addr := u + DowngradeEnabledPath
|
|
|
|
resp, err = cc.Get(addr)
|
|
|
|
if err != nil {
|
|
|
|
lg.Warn(
|
|
|
|
"failed to reach the peer URL",
|
|
|
|
zap.String("address", addr),
|
|
|
|
zap.String("remote-member-id", m.ID.String()),
|
|
|
|
zap.Error(err),
|
|
|
|
)
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
var b []byte
|
2021-10-27 19:00:49 +03:00
|
|
|
b, err = io.ReadAll(resp.Body)
|
2020-10-12 20:50:35 +03:00
|
|
|
resp.Body.Close()
|
|
|
|
if err != nil {
|
|
|
|
lg.Warn(
|
|
|
|
"failed to read body of response",
|
|
|
|
zap.String("address", addr),
|
|
|
|
zap.String("remote-member-id", m.ID.String()),
|
|
|
|
zap.Error(err),
|
|
|
|
)
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
var enable bool
|
|
|
|
if enable, err = strconv.ParseBool(string(b)); err != nil {
|
|
|
|
lg.Warn(
|
|
|
|
"failed to convert response",
|
|
|
|
zap.String("address", addr),
|
|
|
|
zap.String("remote-member-id", m.ID.String()),
|
|
|
|
zap.Error(err),
|
|
|
|
)
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
return enable, nil
|
|
|
|
}
|
|
|
|
return false, err
|
|
|
|
}
|
|
|
|
|
2020-05-16 00:45:00 +03:00
|
|
|
func convertToClusterVersion(v string) (*semver.Version, error) {
|
|
|
|
ver, err := semver.NewVersion(v)
|
|
|
|
if err != nil {
|
|
|
|
// allow input version format Major.Minor
|
|
|
|
ver, err = semver.NewVersion(v + ".0")
|
|
|
|
if err != nil {
|
|
|
|
return nil, ErrWrongDowngradeVersionFormat
|
|
|
|
}
|
|
|
|
}
|
|
|
|
// cluster version only keeps major.minor, remove patch version
|
|
|
|
ver = &semver.Version{Major: ver.Major, Minor: ver.Minor}
|
|
|
|
return ver, nil
|
|
|
|
}
|