516 строки
15 KiB
Go
516 строки
15 KiB
Go
package workercluster
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"io"
|
|
"net"
|
|
"strings"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/hashicorp/raft"
|
|
"github.com/stretchr/testify/assert"
|
|
"github.com/stretchr/testify/require"
|
|
)
|
|
|
|
// clusterFixture is a 3-node in-process cluster used by the e2e tests.
|
|
// It opens three bbolt stores, three HTTP listeners, three rafthttp
|
|
// handlers, and starts each raft.Raft. The first node bootstraps as a
|
|
// 1-voter cluster; the other two join via the raft-membership handler.
|
|
type clusterFixture struct {
|
|
Nodes []*Cluster
|
|
|
|
creds HTTPCreds
|
|
logOut io.Writer
|
|
stopOnce sync.Once
|
|
}
|
|
|
|
func newClusterFixture(t *testing.T, n int) *clusterFixture {
|
|
t.Helper()
|
|
require.GreaterOrEqual(t, n, 1, "cluster fixture requires at least one node")
|
|
|
|
creds := HTTPCreds{Login: "alice", Password: "secret"}
|
|
f := &clusterFixture{
|
|
creds: creds,
|
|
logOut: io.Discard,
|
|
}
|
|
|
|
// Pre-reserve three ports so the bootstrap node has stable
|
|
// addresses for the other nodes to dial into.
|
|
addrs := make([]string, n)
|
|
dataDirs := make([]string, n)
|
|
for i := 0; i < n; i++ {
|
|
l := newLocalListener(t)
|
|
addrs[i] = l.Addr().String()
|
|
require.NoError(t, l.Close())
|
|
dataDirs[i] = t.TempDir()
|
|
}
|
|
|
|
ctx := context.Background()
|
|
|
|
// Bootstrap the first node.
|
|
bootstrapOpts := Options{
|
|
NodeID: "node-1",
|
|
LocalAddr: addrs[0],
|
|
DataDir: dataDirs[0],
|
|
Creds: creds,
|
|
Bootstrap: true,
|
|
HeartbeatTimeout: 300 * time.Millisecond,
|
|
ElectionTimeout: 1000 * time.Millisecond,
|
|
LogOutput: f.logOut,
|
|
}
|
|
c1, err := New(&bootstrapOpts)
|
|
require.NoError(t, err)
|
|
require.NoError(t, c1.Start(ctx))
|
|
require.NoError(t, waitForLeader(ctx, c1, 5*time.Second))
|
|
f.Nodes = append(f.Nodes, c1)
|
|
|
|
// Bring up the remaining nodes sequentially, joining via node-1
|
|
// until each becomes a voter.
|
|
for i := 1; i < n; i++ {
|
|
opts := Options{
|
|
NodeID: fmt.Sprintf("node-%d", i+1),
|
|
LocalAddr: addrs[i],
|
|
DataDir: dataDirs[i],
|
|
Creds: creds,
|
|
Bootstrap: false,
|
|
Seed: Peer{WorkerID: "node-1", Address: addrs[0]},
|
|
HeartbeatTimeout: 300 * time.Millisecond,
|
|
ElectionTimeout: 1000 * time.Millisecond,
|
|
LogOutput: f.logOut,
|
|
}
|
|
c, err := New(&opts)
|
|
require.NoError(t, err)
|
|
require.NoError(t, c.Start(ctx))
|
|
// Wait for the join to be reflected in the cluster config.
|
|
require.NoError(t, waitForVoterCount(ctx, c, i+1, 15*time.Second))
|
|
f.Nodes = append(f.Nodes, c)
|
|
}
|
|
|
|
t.Cleanup(f.Close)
|
|
return f
|
|
}
|
|
|
|
func waitForVoterCount(ctx context.Context, c *Cluster, n int, timeout time.Duration) error {
|
|
deadline := time.Now().Add(timeout)
|
|
for {
|
|
s := c.Stats()
|
|
if s.NumPeers+1 >= n {
|
|
return nil
|
|
}
|
|
if time.Now().After(deadline) {
|
|
return fmt.Errorf("workercluster: only %d voters after %s (need %d)", s.NumPeers+1, timeout, n)
|
|
}
|
|
select {
|
|
case <-ctx.Done():
|
|
return ctx.Err()
|
|
case <-time.After(50 * time.Millisecond):
|
|
}
|
|
}
|
|
}
|
|
|
|
func (f *clusterFixture) Close() {
|
|
f.stopOnce.Do(func() {
|
|
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
|
defer cancel()
|
|
for _, n := range f.Nodes {
|
|
if n != nil {
|
|
_ = n.Shutdown(ctx)
|
|
}
|
|
}
|
|
})
|
|
}
|
|
|
|
// TestThreeNodeCluster_LeaderElection starts a 3-node cluster and
|
|
// verifies exactly one leader is elected within the timeout, and all
|
|
// nodes see the same leader.
|
|
func TestThreeNodeCluster_LeaderElection(t *testing.T) {
|
|
if testing.Short() {
|
|
t.Skip("e2e test skipped in -short mode")
|
|
}
|
|
f := newClusterFixture(t, 3)
|
|
|
|
// Wait until every node reports the same leader.
|
|
deadline := time.Now().Add(5 * time.Second)
|
|
var lastLeader string
|
|
for {
|
|
leaders := make(map[string]int)
|
|
for _, n := range f.Nodes {
|
|
s := n.Stats()
|
|
if s.Leader == "" {
|
|
leaders[""]++
|
|
continue
|
|
}
|
|
leaders[s.Leader]++
|
|
}
|
|
if len(leaders) == 1 {
|
|
for leader, count := range leaders {
|
|
if leader != "" && count == 3 {
|
|
lastLeader = leader
|
|
break
|
|
}
|
|
}
|
|
if lastLeader != "" {
|
|
break
|
|
}
|
|
}
|
|
require.False(t, time.Now().After(deadline), "no consensus on leader within deadline, last seen: %v", leaders)
|
|
time.Sleep(50 * time.Millisecond)
|
|
}
|
|
assert.NotEmpty(t, lastLeader)
|
|
}
|
|
|
|
// TestThreeNodeCluster_FSMApply verifies a config.adopt entry proposed
|
|
// on one node ends up in the FSM of all three.
|
|
func TestThreeNodeCluster_FSMApply(t *testing.T) {
|
|
if testing.Short() {
|
|
t.Skip("e2e test skipped in -short mode")
|
|
}
|
|
f := newClusterFixture(t, 3)
|
|
require.NoError(t, waitForLeader(context.Background(), f.Nodes[0], 5*time.Second))
|
|
|
|
leader := leaderCluster(f)
|
|
require.NotNil(t, leader, "expected a leader in the fixture")
|
|
|
|
checks := []CriticalCheckConfig{
|
|
{ID: 1, MonitorID: 11, Kind: "http", Target: "https://a", IntervalS: 30},
|
|
{ID: 2, MonitorID: 12, Kind: "ssl", Target: "b", IntervalS: 60},
|
|
}
|
|
raw, err := jsonMarshal(ConfigAdoptPayload{Version: 11, Actor: leader.opts.NodeID, Checks: checks})
|
|
require.NoError(t, err)
|
|
|
|
entry, err := EncodeEntry(&Entry{Type: EntryConfigAdopt, Adopted: raw})
|
|
require.NoError(t, err)
|
|
require.NoError(t, leader.Apply(context.Background(), entry, 5*time.Second))
|
|
|
|
// Wait until every node has applied the entry.
|
|
deadline := time.Now().Add(5 * time.Second)
|
|
for {
|
|
ok := true
|
|
for _, n := range f.Nodes {
|
|
if n.FSM().Stats().ConfigVersion != 11 {
|
|
ok = false
|
|
break
|
|
}
|
|
}
|
|
if ok {
|
|
break
|
|
}
|
|
require.False(t, time.Now().After(deadline), "config did not replicate within deadline")
|
|
time.Sleep(50 * time.Millisecond)
|
|
}
|
|
|
|
for _, n := range f.Nodes {
|
|
s := n.FSM().Stats()
|
|
assert.EqualValues(t, 11, s.ConfigVersion, "node %s", n.opts.NodeID)
|
|
assert.Equal(t, 2, s.ConfigCount, "node %s", n.opts.NodeID)
|
|
}
|
|
}
|
|
|
|
// TestThreeNodeCluster_Failover shuts down the leader and verifies a
|
|
// new leader is elected within the election-timeout window.
|
|
func TestThreeNodeCluster_Failover(t *testing.T) {
|
|
if testing.Short() {
|
|
t.Skip("e2e test skipped in -short mode")
|
|
}
|
|
f := newClusterFixture(t, 3)
|
|
require.NoError(t, waitForLeader(context.Background(), f.Nodes[0], 5*time.Second))
|
|
|
|
leader := leaderCluster(f)
|
|
require.NotNil(t, leader)
|
|
oldLeaderID := leader.opts.NodeID
|
|
|
|
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
require.NoError(t, leader.Shutdown(ctx))
|
|
|
|
// Wait until one of the remaining two nodes holds leadership.
|
|
require.NoError(t, waitForLeaderOnAny(f.Nodes[1:], 10*time.Second))
|
|
|
|
leaders := make(map[string]bool)
|
|
for _, n := range f.Nodes[1:] {
|
|
s := n.Stats()
|
|
require.NotEmpty(t, s.Leader, "node %s has no leader after failover", s.NodeID)
|
|
leaders[s.Leader] = true
|
|
}
|
|
require.True(t, leaders[oldLeaderID] == false, "old leader should not be elected after shutdown")
|
|
require.Len(t, leaders, 1, "exactly one leader expected after failover, got %v", leaders)
|
|
}
|
|
|
|
func waitForLeaderOnAny(nodes []*Cluster, timeout time.Duration) error {
|
|
deadline := time.Now().Add(timeout)
|
|
for {
|
|
for _, n := range nodes {
|
|
if n.Stats().Leader != "" {
|
|
return nil
|
|
}
|
|
}
|
|
if time.Now().After(deadline) {
|
|
return fmt.Errorf("no leader after %s", timeout)
|
|
}
|
|
time.Sleep(50 * time.Millisecond)
|
|
}
|
|
}
|
|
|
|
// TestThreeNodeCluster_JoinAndRemove boots a 1-node cluster, adds a
|
|
// second and a third voter, then removes the third. We assert the FSM
|
|
// membership cache reflects the change.
|
|
func TestThreeNodeCluster_JoinAndRemove(t *testing.T) {
|
|
if testing.Short() {
|
|
t.Skip("e2e test skipped in -short mode")
|
|
}
|
|
ctx := context.Background()
|
|
|
|
// 1-node bootstrap.
|
|
creds := HTTPCreds{Login: "alice", Password: "secret"}
|
|
addr1 := pickPort(t)
|
|
c1, err := New(&Options{
|
|
NodeID: "node-1",
|
|
LocalAddr: addr1,
|
|
DataDir: t.TempDir(),
|
|
Creds: creds,
|
|
Bootstrap: true,
|
|
HeartbeatTimeout: 300 * time.Millisecond,
|
|
ElectionTimeout: 1000 * time.Millisecond,
|
|
LogOutput: io.Discard,
|
|
})
|
|
require.NoError(t, err)
|
|
require.NoError(t, c1.Start(ctx))
|
|
require.NoError(t, waitForLeader(ctx, c1, 5*time.Second))
|
|
t.Cleanup(func() {
|
|
shutCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
for _, n := range []*Cluster{c1} {
|
|
if n != nil {
|
|
_ = n.Shutdown(shutCtx)
|
|
}
|
|
}
|
|
})
|
|
|
|
// Join node-2.
|
|
addr2 := pickPort(t)
|
|
c2, err := JoinCluster(ctx, &Options{
|
|
NodeID: "node-2",
|
|
LocalAddr: addr2,
|
|
DataDir: t.TempDir(),
|
|
Creds: creds,
|
|
HeartbeatTimeout: 300 * time.Millisecond,
|
|
ElectionTimeout: 1000 * time.Millisecond,
|
|
LogOutput: io.Discard,
|
|
}, Peer{WorkerID: "node-1", Address: addr1})
|
|
require.NoError(t, err)
|
|
t.Cleanup(func() {
|
|
shutCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
_ = c2.Shutdown(shutCtx)
|
|
})
|
|
require.NoError(t, waitForVoterCount(ctx, c2, 2, 10*time.Second))
|
|
|
|
// Join node-3.
|
|
addr3 := pickPort(t)
|
|
c3, err := JoinCluster(ctx, &Options{
|
|
NodeID: "node-3",
|
|
LocalAddr: addr3,
|
|
DataDir: t.TempDir(),
|
|
Creds: creds,
|
|
HeartbeatTimeout: 300 * time.Millisecond,
|
|
ElectionTimeout: 1000 * time.Millisecond,
|
|
LogOutput: io.Discard,
|
|
}, Peer{WorkerID: "node-1", Address: addr1})
|
|
require.NoError(t, err)
|
|
t.Cleanup(func() {
|
|
shutCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
_ = c3.Shutdown(shutCtx)
|
|
})
|
|
require.NoError(t, waitForVoterCount(ctx, c3, 3, 10*time.Second))
|
|
|
|
// All three nodes should see 3 voters.
|
|
for _, n := range []*Cluster{c1, c2, c3} {
|
|
cfgFuture := n.Raft().GetConfiguration()
|
|
require.NoError(t, cfgFuture.Error())
|
|
ids := make([]string, 0, 3)
|
|
for _, srv := range cfgFuture.Configuration().Servers {
|
|
ids = append(ids, string(srv.ID))
|
|
}
|
|
require.Equal(t, 3, len(ids), "node %s should have 3 voters, got %v", n.opts.NodeID, ids)
|
|
}
|
|
|
|
// Remove node-3 by proposing a membership change via its local
|
|
// raft. (rafthttp only ships Join + Leave; we use raft.RemoveServer
|
|
// directly here to keep the test driver straightforward.)
|
|
leader := pickLeader(t, []*Cluster{c1, c2, c3})
|
|
require.NotNil(t, leader)
|
|
require.NoError(t, leader.Raft().RemoveServer(raft.ServerID("node-3"), 0, 5*time.Second).Error())
|
|
|
|
require.NoError(t, waitForVoterCount(ctx, c1, 2, 10*time.Second))
|
|
}
|
|
|
|
// TestThreeNodeCluster_SnapshotRestore populates the FSM with a few
|
|
// entries, takes a snapshot on the leader, then restarts the leader
|
|
// from the same data dir to verify the snapshot+log restore cycle.
|
|
func TestThreeNodeCluster_SnapshotRestore(t *testing.T) {
|
|
if testing.Short() {
|
|
t.Skip("e2e test skipped in -short mode")
|
|
}
|
|
f := newClusterFixture(t, 3)
|
|
require.NoError(t, waitForLeader(context.Background(), f.Nodes[0], 5*time.Second))
|
|
|
|
leader := leaderCluster(f)
|
|
require.NotNil(t, leader)
|
|
|
|
// Push a few entries so the FSM has something worth snapshotting.
|
|
for i := 0; i < 5; i++ {
|
|
checks := []CriticalCheckConfig{
|
|
{ID: int64(i + 1), MonitorID: int64(i + 1), Kind: "http", Target: fmt.Sprintf("https://a/%d", i), IntervalS: 30},
|
|
}
|
|
raw, err := jsonMarshal(ConfigAdoptPayload{Version: uint64(i + 1), Checks: checks})
|
|
require.NoError(t, err)
|
|
entry, err := EncodeEntry(&Entry{Type: EntryConfigAdopt, Adopted: raw})
|
|
require.NoError(t, err)
|
|
require.NoError(t, leader.Apply(context.Background(), entry, 5*time.Second))
|
|
}
|
|
|
|
// Force a snapshot.
|
|
require.NoError(t, leader.Snapshot())
|
|
|
|
// All followers should see the snapshot effect.
|
|
deadline := time.Now().Add(5 * time.Second)
|
|
for {
|
|
ok := true
|
|
for _, n := range f.Nodes {
|
|
if n.FSM().Stats().ConfigVersion != 5 {
|
|
ok = false
|
|
break
|
|
}
|
|
}
|
|
if ok {
|
|
break
|
|
}
|
|
require.False(t, time.Now().After(deadline), "snapshot did not replicate within deadline")
|
|
time.Sleep(50 * time.Millisecond)
|
|
}
|
|
|
|
// Snapshot file must exist on disk for at least the leader.
|
|
store := leader.store
|
|
snaps, err := store.SnapshotStore().List()
|
|
require.NoError(t, err)
|
|
require.NotEmpty(t, snaps, "expected at least one snapshot on disk after Snapshot()")
|
|
}
|
|
|
|
// TestBootstrapSafety_TwoClusters confirms that starting two
|
|
// independent clusters with the same bootstrap token creates two
|
|
// distinct raft groups, not a single merged one.
|
|
func TestBootstrapSafety_TwoClusters(t *testing.T) {
|
|
if testing.Short() {
|
|
t.Skip("e2e test skipped in -short mode")
|
|
}
|
|
ctx := context.Background()
|
|
|
|
creds := HTTPCreds{Login: "alice", Password: "secret"}
|
|
a1, err := New(&Options{
|
|
NodeID: "a-1",
|
|
LocalAddr: pickPort(t),
|
|
DataDir: t.TempDir(),
|
|
Creds: creds,
|
|
Bootstrap: true,
|
|
HeartbeatTimeout: 300 * time.Millisecond,
|
|
ElectionTimeout: 1000 * time.Millisecond,
|
|
LogOutput: io.Discard,
|
|
})
|
|
require.NoError(t, err)
|
|
require.NoError(t, a1.Start(ctx))
|
|
t.Cleanup(func() {
|
|
shutCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
_ = a1.Shutdown(shutCtx)
|
|
})
|
|
|
|
b1, err := New(&Options{
|
|
NodeID: "b-1",
|
|
LocalAddr: pickPort(t),
|
|
DataDir: t.TempDir(),
|
|
Creds: creds,
|
|
Bootstrap: true,
|
|
HeartbeatTimeout: 300 * time.Millisecond,
|
|
ElectionTimeout: 1000 * time.Millisecond,
|
|
LogOutput: io.Discard,
|
|
})
|
|
require.NoError(t, err)
|
|
require.NoError(t, b1.Start(ctx))
|
|
t.Cleanup(func() {
|
|
shutCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
_ = b1.Shutdown(shutCtx)
|
|
})
|
|
|
|
require.NoError(t, waitForLeader(ctx, a1, 5*time.Second))
|
|
require.NoError(t, waitForLeader(ctx, b1, 5*time.Second))
|
|
|
|
require.NoError(t, waitForLeaderOnAny([]*Cluster{a1}, 1*time.Second))
|
|
|
|
// Each cluster must have exactly one voter: itself.
|
|
aConfig := a1.Raft().GetConfiguration()
|
|
require.NoError(t, aConfig.Error())
|
|
require.Len(t, aConfig.Configuration().Servers, 1)
|
|
assert.Equal(t, raft.ServerID("a-1"), aConfig.Configuration().Servers[0].ID)
|
|
|
|
bConfig := b1.Raft().GetConfiguration()
|
|
require.NoError(t, bConfig.Error())
|
|
require.Len(t, bConfig.Configuration().Servers, 1)
|
|
assert.Equal(t, raft.ServerID("b-1"), bConfig.Configuration().Servers[0].ID)
|
|
}
|
|
|
|
// pickPort reserves a free port on the loopback interface and returns
|
|
// its address. We close the listener immediately; the actual raft
|
|
// listener will rebind to the same port because no other process has
|
|
// claimed it.
|
|
func pickPort(t *testing.T) string {
|
|
t.Helper()
|
|
l := newLocalListener(t)
|
|
addr := l.Addr().String()
|
|
require.NoError(t, l.Close())
|
|
return addr
|
|
}
|
|
|
|
func leaderCluster(f *clusterFixture) *Cluster {
|
|
for _, n := range f.Nodes {
|
|
if n.Raft().State() == raft.Leader {
|
|
return n
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func leaderFromFixture(f *clusterFixture) string {
|
|
for _, n := range f.Nodes {
|
|
if n.Raft().State() == raft.Leader {
|
|
return n.opts.NodeID
|
|
}
|
|
}
|
|
return ""
|
|
}
|
|
|
|
func pickLeader(t *testing.T, nodes []*Cluster) *Cluster {
|
|
t.Helper()
|
|
for _, n := range nodes {
|
|
if n.Raft().State() == raft.Leader {
|
|
return n
|
|
}
|
|
}
|
|
t.Fatal("no leader in cluster fixture")
|
|
return nil
|
|
}
|
|
|
|
// ensureLocalAddrsAreLoopback sanity-checks the helper above produces
|
|
// 127.0.0.1 addresses so the e2e suite never accidentally opens a
|
|
// non-loopback port.
|
|
func TestPickPortIsLoopback(t *testing.T) {
|
|
addr := pickPort(t)
|
|
host, _, err := net.SplitHostPort(addr)
|
|
require.NoError(t, err)
|
|
assert.True(t, strings.HasPrefix(host, "127.") || host == "::1", "expected loopback, got %q", host)
|
|
}
|