Files
worker/internal/workercluster/demo_test.go
Gleb Tv 2c7a0236da feat: publish standalone worker
Separate worker packaging and service lifecycle from the control plane.
2026-07-13 17:55:14 +03:00

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")
}
}