Merge pull request #15449 from fuweid/fix-15409
tests/integration: deflake TestEtcdVersionFromWALdependabot/go_modules/go.uber.org/atomic-1.10.0
commit
043525c69d
|
@ -2168,12 +2168,14 @@ func (s *EtcdServer) StorageVersion() *semver.Version {
|
|||
|
||||
// monitorClusterVersions every monitorVersionInterval checks if it's the leader and updates cluster version if needed.
|
||||
func (s *EtcdServer) monitorClusterVersions() {
|
||||
monitor := serverversion.NewMonitor(s.Logger(), NewServerVersionAdapter(s))
|
||||
lg := s.Logger()
|
||||
monitor := serverversion.NewMonitor(lg, NewServerVersionAdapter(s))
|
||||
for {
|
||||
select {
|
||||
case <-s.firstCommitInTerm.Receive():
|
||||
case <-time.After(monitorVersionInterval):
|
||||
case <-s.stopping:
|
||||
lg.Info("server has stopped; stopping cluster version's monitor")
|
||||
return
|
||||
}
|
||||
|
||||
|
@ -2189,12 +2191,14 @@ func (s *EtcdServer) monitorClusterVersions() {
|
|||
|
||||
// monitorStorageVersion every monitorVersionInterval updates storage version if needed.
|
||||
func (s *EtcdServer) monitorStorageVersion() {
|
||||
monitor := serverversion.NewMonitor(s.Logger(), NewServerVersionAdapter(s))
|
||||
lg := s.Logger()
|
||||
monitor := serverversion.NewMonitor(lg, NewServerVersionAdapter(s))
|
||||
for {
|
||||
select {
|
||||
case <-time.After(monitorVersionInterval):
|
||||
case <-s.clusterVersionChanged.Receive():
|
||||
case <-s.stopping:
|
||||
lg.Info("server has stopped; stopping storage version's monitor")
|
||||
return
|
||||
}
|
||||
monitor.UpdateStorageVersionIfNeeded()
|
||||
|
@ -2218,6 +2222,7 @@ func (s *EtcdServer) monitorKVHash() {
|
|||
for {
|
||||
select {
|
||||
case <-s.stopping:
|
||||
lg.Info("server has stopped; stopping kv hash's monitor")
|
||||
return
|
||||
case <-checkTicker.C:
|
||||
}
|
||||
|
@ -2239,6 +2244,8 @@ func (s *EtcdServer) monitorCompactHash() {
|
|||
select {
|
||||
case <-time.After(t):
|
||||
case <-s.stopping:
|
||||
lg := s.Logger()
|
||||
lg.Info("server has stopped; stopping compact hash's monitor")
|
||||
return
|
||||
}
|
||||
if !s.isLeader() {
|
||||
|
|
|
@ -28,6 +28,7 @@ import (
|
|||
"go.etcd.io/etcd/server/v3/embed"
|
||||
"go.etcd.io/etcd/server/v3/storage/wal"
|
||||
"go.etcd.io/etcd/server/v3/storage/wal/walpb"
|
||||
framecfg "go.etcd.io/etcd/tests/v3/framework/config"
|
||||
"go.etcd.io/etcd/tests/v3/framework/integration"
|
||||
)
|
||||
|
||||
|
@ -45,28 +46,62 @@ func TestEtcdVersionFromWAL(t *testing.T) {
|
|||
t.Fatalf("failed to start embed.Etcd for test")
|
||||
}
|
||||
|
||||
// When the member becomes leader, it will update the cluster version
|
||||
// with the cluster's minimum version. As it's updated asynchronously,
|
||||
// it could not be updated in time before close. Wait for it to become
|
||||
// ready.
|
||||
if err := waitForClusterVersionReady(srv); err != nil {
|
||||
srv.Close()
|
||||
t.Fatalf("failed to wait for cluster version to become ready: %v", err)
|
||||
}
|
||||
|
||||
ccfg := clientv3.Config{Endpoints: []string{cfg.ACUrls[0].String()}}
|
||||
cli, err := integration.NewClient(t, ccfg)
|
||||
if err != nil {
|
||||
srv.Close()
|
||||
t.Fatal(err)
|
||||
}
|
||||
// Get auth status to increase etcd version of proto stored in wal
|
||||
|
||||
// Once the cluster version has been updated, any entity's storage
|
||||
// version should be align with cluster version.
|
||||
ctx, cancel := context.WithTimeout(context.Background(), testutil.RequestTimeout)
|
||||
cli.AuthStatus(ctx)
|
||||
_, err = cli.AuthStatus(ctx)
|
||||
cancel()
|
||||
if err != nil {
|
||||
srv.Close()
|
||||
t.Fatalf("failed to get auth status: %v", err)
|
||||
}
|
||||
|
||||
cli.Close()
|
||||
srv.Close()
|
||||
|
||||
w, err := wal.Open(zap.NewNop(), cfg.Dir+"/member/wal", walpb.Snapshot{})
|
||||
if err != nil {
|
||||
panic(err)
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer w.Close()
|
||||
|
||||
walVersion, err := wal.ReadWALVersion(w)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
assert.Equal(t, &semver.Version{Major: 3, Minor: 6}, walVersion.MinimalEtcdVersion())
|
||||
}
|
||||
|
||||
func waitForClusterVersionReady(srv *embed.Etcd) error {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||
defer cancel()
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
default:
|
||||
}
|
||||
|
||||
if srv.Server.ClusterVersion() != nil {
|
||||
return nil
|
||||
}
|
||||
time.Sleep(framecfg.TickDuration)
|
||||
}
|
||||
}
|
||||
|
|
Loading…
Reference in New Issue