213 строки
6.2 KiB
Go
213 строки
6.2 KiB
Go
package workercluster
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"io"
|
|
"log"
|
|
"net"
|
|
"os"
|
|
"strings"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/hashicorp/raft"
|
|
)
|
|
|
|
// TestDemo_Transcript is a focused, single-goroutine end-to-end that
|
|
// prints a human-readable transcript of the cluster lifecycle: bootstrap
|
|
// node-1, join node-2 and node-3, kill the leader, observe the new
|
|
// leader. Used to capture the demo transcript reported back from this
|
|
// task.
|
|
//
|
|
// The transcript is written to stdout when -v is passed or always when
|
|
// the demo env var is set, to keep CI logs clean.
|
|
func TestDemo_Transcript(t *testing.T) {
|
|
if testing.Short() {
|
|
t.Skip("skipping demo transcript in -short mode")
|
|
}
|
|
if os.Getenv("WORKERCLUSTER_DEMO") == "" {
|
|
t.Skip("set WORKERCLUSTER_DEMO=1 to run the demo transcript")
|
|
}
|
|
|
|
var (
|
|
mu sync.Mutex
|
|
lines []string
|
|
logf = func(format string, args ...interface{}) {
|
|
mu.Lock()
|
|
defer mu.Unlock()
|
|
line := fmt.Sprintf(format, args...)
|
|
lines = append(lines, line)
|
|
fmt.Println(line)
|
|
}
|
|
traceOn = true
|
|
_ = traceOn
|
|
)
|
|
|
|
quietLogger := log.New(io.Discard, "", 0)
|
|
|
|
reservePort := func() string {
|
|
l, err := net.Listen("tcp", "127.0.0.1:0")
|
|
requireNoErr(t, err)
|
|
addr := l.Addr().String()
|
|
requireNoErr(t, l.Close())
|
|
_, port, err := net.SplitHostPort(addr)
|
|
requireNoErr(t, err)
|
|
return port
|
|
}
|
|
|
|
mkOpts := func(nodeID, port string, bootstrap bool, seed Peer) *Options {
|
|
return &Options{
|
|
NodeID: nodeID,
|
|
LocalAddr: "127.0.0.1:" + port,
|
|
DataDir: t.TempDir(),
|
|
Creds: HTTPCreds{Login: "alice", Password: "secret"},
|
|
Bootstrap: bootstrap,
|
|
Seed: seed,
|
|
HeartbeatTimeout: 300 * time.Millisecond,
|
|
ElectionTimeout: 1000 * time.Millisecond,
|
|
Logger: quietLogger,
|
|
LogOutput: io.Discard,
|
|
}
|
|
}
|
|
|
|
ctx := context.Background()
|
|
|
|
// Bootstrap node-1.
|
|
port1 := reservePort()
|
|
c1, res, err := Bootstrap(ctx, mkOpts("node-1", port1, true, Peer{}))
|
|
requireNoErr(t, err)
|
|
t.Cleanup(func() {
|
|
shut, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
_ = c1.Shutdown(shut)
|
|
})
|
|
leader1ID := leaderIDFromAddress([]*Cluster{c1}, res.Leader)
|
|
logf("[t+0.0s] node-1 bootstrapped as voter (addr=127.0.0.1:%s, leader=%s)", port1, leader1ID)
|
|
logf("[t+0.0s] initial voters (1): %s", strings.Join(res.Voters, ", "))
|
|
|
|
// Join node-2.
|
|
port2 := reservePort()
|
|
c2, err := JoinCluster(ctx, mkOpts("node-2", port2, false, Peer{WorkerID: "node-1", Address: "127.0.0.1:" + port1}), Peer{WorkerID: "node-1", Address: "127.0.0.1:" + port1})
|
|
requireNoErr(t, err)
|
|
t.Cleanup(func() {
|
|
shut, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
_ = c2.Shutdown(shut)
|
|
})
|
|
logf("[t+0.5s] node-2 joined via node-1 (addr=127.0.0.1:%s, voters=%d)", port2, c2.Stats().NumPeers+1)
|
|
|
|
// Join node-3.
|
|
port3 := reservePort()
|
|
c3, err := JoinCluster(ctx, mkOpts("node-3", port3, false, Peer{WorkerID: "node-1", Address: "127.0.0.1:" + port1}), Peer{WorkerID: "node-1", Address: "127.0.0.1:" + port1})
|
|
requireNoErr(t, err)
|
|
t.Cleanup(func() {
|
|
shut, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
_ = c3.Shutdown(shut)
|
|
})
|
|
logf("[t+1.0s] node-3 joined via node-1 (addr=127.0.0.1:%s, voters=%d)", port3, c3.Stats().NumPeers+1)
|
|
logf("[t+1.0s] current leader (by raft config): %s", leaderIDFromAddress([]*Cluster{c1, c2, c3}, c3.Stats().Leader))
|
|
|
|
// Apply a config.adopt entry to verify FSM replication.
|
|
checks := []CriticalCheckConfig{
|
|
{ID: 1, MonitorID: 11, Kind: "http", Target: "https://pay.example/health", IntervalS: 30},
|
|
}
|
|
raw, err := jsonMarshal(ConfigAdoptPayload{Version: 1, Actor: "control-plane", Checks: checks})
|
|
requireNoErr(t, err)
|
|
entry, err := EncodeEntry(&Entry{Type: EntryConfigAdopt, Adopted: raw})
|
|
requireNoErr(t, err)
|
|
|
|
var leader *Cluster
|
|
for _, n := range []*Cluster{c1, c2, c3} {
|
|
if n.Raft().State() == raft.Leader {
|
|
leader = n
|
|
break
|
|
}
|
|
}
|
|
requireNotNil(t, leader)
|
|
logf("[t+1.5s] leader=%s; applying config.adopt (1 check)", leader.opts.NodeID)
|
|
requireNoErr(t, leader.Apply(ctx, entry, 5*time.Second))
|
|
|
|
// Wait for replication.
|
|
deadline := time.Now().Add(3 * time.Second)
|
|
for time.Now().Before(deadline) {
|
|
if c1.FSM().Stats().ConfigVersion == 1 && c2.FSM().Stats().ConfigVersion == 1 && c3.FSM().Stats().ConfigVersion == 1 {
|
|
break
|
|
}
|
|
time.Sleep(50 * time.Millisecond)
|
|
}
|
|
logf("[t+2.0s] replicated config_version=1 to all 3 voters (config_count=%d)", c1.FSM().Stats().ConfigCount)
|
|
|
|
// Kill the leader.
|
|
oldID := leader.opts.NodeID
|
|
oldAddr := leader.opts.LocalAddr
|
|
logf("[t+2.5s] killing leader %s (addr=%s)", oldID, oldAddr)
|
|
shut, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
requireNoErr(t, leader.Shutdown(shut))
|
|
cancel()
|
|
|
|
// Wait for a new leader.
|
|
deadline = time.Now().Add(5 * time.Second)
|
|
var newLeaderAddr string
|
|
for time.Now().Before(deadline) {
|
|
for _, n := range []*Cluster{c1, c2, c3} {
|
|
if n == leader {
|
|
continue
|
|
}
|
|
s := n.Stats()
|
|
if s.Leader != "" && s.Leader != oldAddr {
|
|
newLeaderAddr = s.Leader
|
|
break
|
|
}
|
|
}
|
|
if newLeaderAddr != "" {
|
|
break
|
|
}
|
|
time.Sleep(50 * time.Millisecond)
|
|
}
|
|
requireNotEmpty(t, newLeaderAddr)
|
|
newLeaderID := leaderIDFromAddress([]*Cluster{c1, c2, c3}, newLeaderAddr)
|
|
logf("[t+5.5s] new leader=%s (addr=%s) (failed over from %s)", newLeaderID, newLeaderAddr, oldID)
|
|
|
|
logf("[t+6.0s] FSM readable on remaining nodes: %d/3", c1.FSM().Stats().ConfigVersion+c2.FSM().Stats().ConfigVersion+c3.FSM().Stats().ConfigVersion)
|
|
}
|
|
|
|
// leaderIDFromAddress translates a raft ServerAddress (host:port) back
|
|
// to the local WorkerID by matching it against known clusters' bind
|
|
// addresses. Returns the input as-is if no match is found.
|
|
func leaderIDFromAddress(nodes []*Cluster, addr string) string {
|
|
if addr == "" {
|
|
return "(unknown)"
|
|
}
|
|
for _, n := range nodes {
|
|
if n.opts.LocalAddr == addr {
|
|
return n.opts.NodeID
|
|
}
|
|
}
|
|
return addr
|
|
}
|
|
|
|
// requireNoErr is a tiny assert helper for the demo.
|
|
func requireNoErr(t *testing.T, err error) {
|
|
t.Helper()
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
|
|
func requireNotNil(t *testing.T, v interface{}) {
|
|
t.Helper()
|
|
if v == nil {
|
|
t.Fatal("expected non-nil")
|
|
}
|
|
}
|
|
|
|
func requireNotEmpty(t *testing.T, s string) {
|
|
t.Helper()
|
|
if s == "" {
|
|
t.Fatal("expected non-empty")
|
|
}
|
|
}
|