From d1086e3dbdf86f8416f1bd4c7e965d902d684d07 Mon Sep 17 00:00:00 2001 From: dawn Date: Thu, 13 Aug 2026 10:34:11 +0900 Subject: [PATCH] spindle/mill: proxy debug ssh to executors Signed-off-by: dawn --- docker-compose.mill.yml | 31 +++-- docker-compose.yml | 2 +- nix/modules/spindle.nix | 31 +++++ spindle/config/config.go | 29 +++-- spindle/config/config_test.go | 45 +++++++ spindle/engines/microvm/README.md | 20 ++- spindle/engines/microvm/debug_test.go | 4 +- spindle/mill/jump.go | 170 ++++++++++++++++++++++++++ spindle/mill/jump_test.go | 158 ++++++++++++++++++++++++ spindle/mill/mill.go | 7 ++ spindle/server.go | 4 + 11 files changed, 476 insertions(+), 25 deletions(-) create mode 100644 spindle/mill/jump.go create mode 100644 spindle/mill/jump_test.go diff --git a/docker-compose.mill.yml b/docker-compose.mill.yml index 36e151a0..6888a5f0 100644 --- a/docker-compose.mill.yml +++ b/docker-compose.mill.yml @@ -33,8 +33,8 @@ x-mill-executor: &mill-executor SPINDLE_NIX_CACHE_UPLOAD_URL: http://ncps:8501/upload SPINDLE_MICROVM_PIPELINES_DEBUG_SSH_ENABLED: "true" SPINDLE_MICROVM_PIPELINES_DEBUG_SSH_LISTEN_ADDR: 0.0.0.0:2223 - SPINDLE_MICROVM_PIPELINES_DEBUG_SSH_HOST: 127.0.0.1 - SPINDLE_MICROVM_PIPELINES_DEBUG_SSH_JUMP_HOST: chernobog + SPINDLE_MICROVM_PIPELINES_DEBUG_SSH_HOST: spindle-executor + SPINDLE_MICROVM_PIPELINES_DEBUG_SSH_JUMP_HOST: debug@spindle.tngl.boltless.dev:2224 SPINDLE_MICROVM_PIPELINES_DEBUG_SSH_GRACE_PERIOD: 10m SPINDLE_MICROVM_PIPELINES_DEBUG_SSH_HOST_KEY_PATH: /var/lib/spindle/debug_ssh_host_key # dials the mill container directly using ws @@ -78,8 +78,14 @@ services: SPINDLE_ROLE: mill SPINDLE_ARTIFACT_STORES_DISK_DIR: /var/lib/spindle/artifacts SPINDLE_MILL_ARTIFACT_STORE: disk + SPINDLE_MILL_JUMP_LISTEN_ADDR: 0.0.0.0:2224 + SPINDLE_MILL_JUMP_HOST_KEY_PATH: /var/lib/spindle/debug_jump_host_key + SPINDLE_MILL_DEBUG_EXECUTOR_PORT: "2223" + SPINDLE_MILL_MAX_JUMP_CONNECTIONS: "128" volumes: - spindle-artifacts:/var/lib/spindle/artifacts + ports: !override + - "127.0.0.1:2224:2224" mill-tokens: profiles: ["linux"] @@ -119,15 +125,16 @@ services: SPINDLE_MILL_LABELS: linux,fast SPINDLE_MILL_TOKEN_FILE: /shared/executor-a.mill-token SPINDLE_MICROVM_PIPELINES_AGENT_PORT: "11241" - SPINDLE_MICROVM_PIPELINES_DEBUG_SSH_LISTEN_ADDR: 0.0.0.0:2224 - ports: - - "127.0.0.1:2224:2224" + SPINDLE_MICROVM_PIPELINES_DEBUG_SSH_HOST: executor-a volumes: - spindle-executor-a-data:/var/lib/spindle - spindle-artifacts:/var/lib/spindle/artifacts - ./out/localinfra-spindle-images:/var/lib/spindle/images:ro - init-state:/shared:ro - ./localinfra/certs/root.crt:/usr/local/share/ca-certificates/caddy.crt:ro + networks: + tngl: + aliases: [executor-a] spindle-executor-b: <<: *mill-executor @@ -137,15 +144,16 @@ services: SPINDLE_MILL_LABELS: linux,slow SPINDLE_MILL_TOKEN_FILE: /shared/executor-b.mill-token SPINDLE_MICROVM_PIPELINES_AGENT_PORT: "11242" - SPINDLE_MICROVM_PIPELINES_DEBUG_SSH_LISTEN_ADDR: 0.0.0.0:2225 - ports: - - "127.0.0.1:2225:2225" + SPINDLE_MICROVM_PIPELINES_DEBUG_SSH_HOST: executor-b volumes: - spindle-executor-b-data:/var/lib/spindle - spindle-artifacts:/var/lib/spindle/artifacts - ./out/localinfra-spindle-images:/var/lib/spindle/images:ro - init-state:/shared:ro - ./localinfra/certs/root.crt:/usr/local/share/ca-certificates/caddy.crt:ro + networks: + tngl: + aliases: [executor-b] spindle-executor-c: <<: *mill-executor @@ -155,15 +163,16 @@ services: SPINDLE_MILL_LABELS: linux,gpu SPINDLE_MILL_TOKEN_FILE: /shared/executor-c.mill-token SPINDLE_MICROVM_PIPELINES_AGENT_PORT: "11243" - SPINDLE_MICROVM_PIPELINES_DEBUG_SSH_LISTEN_ADDR: 0.0.0.0:2226 - ports: - - "127.0.0.1:2226:2226" + SPINDLE_MICROVM_PIPELINES_DEBUG_SSH_HOST: executor-c volumes: - spindle-executor-c-data:/var/lib/spindle - spindle-artifacts:/var/lib/spindle/artifacts - ./out/localinfra-spindle-images-alpine:/var/lib/spindle/images:ro - init-state:/shared:ro - ./localinfra/certs/root.crt:/usr/local/share/ca-certificates/caddy.crt:ro + networks: + tngl: + aliases: [executor-c] volumes: spindle-executor-a-data: diff --git a/docker-compose.yml b/docker-compose.yml index 5c004a54..3e094a6b 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -203,7 +203,7 @@ services: - label=disable - seccomp=unconfined ports: - - "2223:2223" + - "127.0.0.1:2223:2223" volumes: - spindle-data:/var/lib/spindle - spindle-logs:/var/log/spindle diff --git a/nix/modules/spindle.nix b/nix/modules/spindle.nix index 1f6899a4..499d52d2 100644 --- a/nix/modules/spindle.nix +++ b/nix/modules/spindle.nix @@ -348,6 +348,33 @@ in }; }; }; + mill = { + jumpListenAddr = mkOption { + type = types.str; + default = ""; + example = "0.0.0.0:22"; + description = "Address for the mill's restricted debug SSH jump server."; + }; + + jumpHostKeyPath = mkOption { + type = with types; nullOr path; + default = null; + example = "/var/lib/spindle/debug_jump_host_key"; + description = "Path to the debug SSH jump server host key."; + }; + + debugExecutorPort = mkOption { + type = types.port; + default = 2223; + description = "Private debug SSH port shared by executors."; + }; + + maxJumpConnections = mkOption { + type = types.ints.positive; + default = 128; + description = "Maximum concurrent connections to the mill's debug SSH jump server."; + }; + }; environmentFile = mkOption { type = with types; nullOr path; @@ -458,6 +485,10 @@ in "SPINDLE_ARTIFACT_STORES_S3_BUCKET=${cfg.artifactStores.s3.bucket}" "SPINDLE_ARTIFACT_STORES_S3_REGION=${cfg.artifactStores.s3.region}" "SPINDLE_MILL_ARTIFACT_STORE=s3" + "SPINDLE_MILL_JUMP_LISTEN_ADDR=${cfg.mill.jumpListenAddr}" + "SPINDLE_MILL_JUMP_HOST_KEY_PATH=${optionalString (cfg.mill.jumpHostKeyPath != null) (toString cfg.mill.jumpHostKeyPath)}" + "SPINDLE_MILL_DEBUG_EXECUTOR_PORT=${toString cfg.mill.debugExecutorPort}" + "SPINDLE_MILL_MAX_JUMP_CONNECTIONS=${toString cfg.mill.maxJumpConnections}" ]; ExecStart = "${cfg.package}/bin/spindle"; Restart = "always"; diff --git a/spindle/config/config.go b/spindle/config/config.go index 2143aeee..fc6f61a2 100644 --- a/spindle/config/config.go +++ b/spindle/config/config.go @@ -138,13 +138,17 @@ const ( // fields are selectively active depending on the role type Mill struct { - URL string `env:"URL"` // mill websocket endpoint dialled by the executor - SharedSecret string `env:"SHARED_SECRET"` // the executor's token for dialing the mill - MaxPending int `env:"MAX_PENDING, default=100"` // mill pending job queue limit - ReconnectGrace time.Duration `env:"RECONNECT_GRACE, default=45s"` // reconnect window before leases are failed - Seats int `env:"SEATS, default=4"` // executor seats advertised to the mill - Labels []string `env:"LABELS"` // executor capability labels - ArtifactStore string `env:"ARTIFACT_STORE"` // store shared by mill and its executors + URL string `env:"URL"` // mill websocket endpoint dialled by the executor + SharedSecret string `env:"SHARED_SECRET"` // the executor's token for dialing the mill + MaxPending int `env:"MAX_PENDING, default=100"` // mill pending job queue limit + ReconnectGrace time.Duration `env:"RECONNECT_GRACE, default=45s"` // reconnect window before leases are failed + Seats int `env:"SEATS, default=4"` // executor seats advertised to the mill + Labels []string `env:"LABELS"` // executor capability labels + ArtifactStore string `env:"ARTIFACT_STORE"` // store shared by mill and its executors + JumpListenAddr string `env:"JUMP_LISTEN_ADDR"` + JumpHostKeyPath string `env:"JUMP_HOST_KEY_PATH"` + DebugExecutorPort uint32 `env:"DEBUG_EXECUTOR_PORT, default=2223"` + MaxJumpConnections int `env:"MAX_JUMP_CONNECTIONS, default=128"` } type Config struct { @@ -174,6 +178,17 @@ func (c *Config) validate() error { default: return fmt.Errorf("unknown SPINDLE_ROLE %q (want standalone, mill, or executor)", c.Role) } + if c.Mill.JumpListenAddr != "" { + if c.Role != RoleMill { + return fmt.Errorf("SPINDLE_MILL_JUMP_LISTEN_ADDR requires SPINDLE_ROLE=mill") + } + if c.Mill.JumpHostKeyPath == "" { + return fmt.Errorf("SPINDLE_MILL_JUMP_LISTEN_ADDR requires SPINDLE_MILL_JUMP_HOST_KEY_PATH") + } + if c.Mill.MaxJumpConnections <= 0 { + return fmt.Errorf("SPINDLE_MILL_MAX_JUMP_CONNECTIONS must be greater than zero") + } + } return nil } diff --git a/spindle/config/config_test.go b/spindle/config/config_test.go index 12a42e39..b23f7f0a 100644 --- a/spindle/config/config_test.go +++ b/spindle/config/config_test.go @@ -18,3 +18,48 @@ func TestLoadAllowsUnconfiguredMicroVMEngine(t *testing.T) { t.Fatalf("image directory = %q, want empty", cfg.MicroVMPipelines.ImageDir) } } + +func TestLoadRequiresJumpHostKey(t *testing.T) { + t.Setenv("SPINDLE_SERVER_HOSTNAME", "spindle.example.com") + t.Setenv("SPINDLE_SERVER_OWNER", "did:web:spindle.example.com") + t.Setenv("SPINDLE_ROLE", "mill") + t.Setenv("SPINDLE_MILL_JUMP_LISTEN_ADDR", "0.0.0.0:22") + + if _, err := Load(context.Background()); err == nil { + t.Fatal("Load accepted a jump listener without a host key path") + } +} + +func TestLoadRejectsJumpListenerOutsideMill(t *testing.T) { + t.Setenv("SPINDLE_SERVER_HOSTNAME", "spindle.example.com") + t.Setenv("SPINDLE_SERVER_OWNER", "did:web:spindle.example.com") + t.Setenv("SPINDLE_ROLE", "standalone") + t.Setenv("SPINDLE_MILL_JUMP_LISTEN_ADDR", "0.0.0.0:22") + t.Setenv("SPINDLE_MILL_JUMP_HOST_KEY_PATH", "/tmp/jump-host-key") + + if _, err := Load(context.Background()); err == nil { + t.Fatal("Load accepted a mill jump listener in standalone mode") + } +} + +func TestLoadValidatesMaxJumpConnections(t *testing.T) { + t.Setenv("SPINDLE_SERVER_HOSTNAME", "spindle.example.com") + t.Setenv("SPINDLE_SERVER_OWNER", "did:web:spindle.example.com") + t.Setenv("SPINDLE_ROLE", "mill") + t.Setenv("SPINDLE_MILL_JUMP_LISTEN_ADDR", "0.0.0.0:22") + t.Setenv("SPINDLE_MILL_JUMP_HOST_KEY_PATH", "/tmp/jump-host-key") + t.Setenv("SPINDLE_MILL_MAX_JUMP_CONNECTIONS", "0") + + if _, err := Load(context.Background()); err == nil { + t.Fatal("Load accepted a non-positive jump connection limit") + } + + t.Setenv("SPINDLE_MILL_MAX_JUMP_CONNECTIONS", "17") + cfg, err := Load(context.Background()) + if err != nil { + t.Fatal(err) + } + if cfg.Mill.MaxJumpConnections != 17 { + t.Fatalf("max jump connections = %d, want 17", cfg.Mill.MaxJumpConnections) + } +} diff --git a/spindle/engines/microvm/README.md b/spindle/engines/microvm/README.md index da346ff3..6c516818 100644 --- a/spindle/engines/microvm/README.md +++ b/spindle/engines/microvm/README.md @@ -243,15 +243,27 @@ SSH credentials or the destination store itself. ### Debug ssh When a workflow fails, spindle can keep its microVM alive for a configured grace -window (`MicroVMPipelines.SSH`) and print an `ssh` invocation so you can poke at -the failed VM interactively. Spindle terminates the ssh connection itself and -bridges a pty into the live guest over the agent's vsock; the guest stays -keyless and never runs an ssh daemon. +window (`MicroVMPipelines.DebugSSH.GracePeriod`) and print an `ssh` invocation +so you can poke at the failed VM interactively. Spindle terminates the ssh +connection itself and bridges a pty into the live guest over the agent's +vsock; the guest stays keyless and never runs an ssh daemon. Access mirrors a git push: the ssh username is the job id, and the offered public key is sent to the job's repo knot (`sh.tangled.repo.checkPushAllowed`). The session is accepted only if that key is allowed to push to the job's repo. +In a mill fleet, the printed command can use `ssh -J` through the mill's +restricted jump listener. The inner SSH connection still terminates on the +executor, so the mill only forwards an encrypted TCP stream to a live, +operator-registered executor route. + +`SPINDLE_MILL_MAX_JUMP_CONNECTIONS` limits concurrent outer SSH connections. + +Configure each executor's debug host as its registered executor name. The jump +listener checks that the name has a live authenticated mill session and dials +it on the configured private debug port. The registered name must therefore +resolve on the mill's private network. + The shell is deliberately not configurable from either end. It always: - runs as the `spindle-workflow` user (the ssh username selects the *job*, not a unix user), diff --git a/spindle/engines/microvm/debug_test.go b/spindle/engines/microvm/debug_test.go index 38ce61bb..c5a2d0e5 100644 --- a/spindle/engines/microvm/debug_test.go +++ b/spindle/engines/microvm/debug_test.go @@ -5,8 +5,8 @@ package microvm import "testing" func TestDebugSSHCommandUsesJumpHost(t *testing.T) { - got := debugSSHCommand("0.0.0.0:2224", "executor-a.internal", "127.0.0.1", "spindle.example", "job-1") - want := "ssh -tt -J spindle.example -p 2224 job-1@127.0.0.1" + got := debugSSHCommand("0.0.0.0:2224", "executor-a.internal", "executor-a", "spindle.example", "job-1") + want := "ssh -tt -J spindle.example -p 2224 job-1@executor-a" if got != want { t.Fatalf("debug ssh command = %q, want %q", got, want) } diff --git a/spindle/mill/jump.go b/spindle/mill/jump.go new file mode 100644 index 00000000..82c1a354 --- /dev/null +++ b/spindle/mill/jump.go @@ -0,0 +1,170 @@ +package mill + +import ( + "context" + "crypto/ed25519" + "crypto/rand" + "encoding/pem" + "fmt" + "net" + "os" + "path/filepath" + "sync" + "time" + + "github.com/gliderlabs/ssh" + gossh "golang.org/x/crypto/ssh" +) + +const ( + jumpIdleTimeout = 5 * time.Minute + jumpMaxTimeout = 24 * time.Hour + maxJumpConnectionsPerIP = 8 +) + +type jumpContextKey string + +const jumpRouteOpened jumpContextKey = "route-opened" + +func (m *Mill) ServeJump(ctx context.Context, listenAddr, hostKeyPath string, executorPort uint32, maxConnections int) { + if listenAddr == "" { + return + } + srv, err := m.newJumpServer(hostKeyPath, executorPort, maxConnections) + if err != nil { + m.l.Error("setup debug ssh jump server", "err", err) + return + } + + go func() { + <-ctx.Done() + _ = srv.Close() + }() + + m.l.Info("starting debug ssh jump server", "address", listenAddr) + srv.Addr = listenAddr + if err := srv.ListenAndServe(); err != nil && err != ssh.ErrServerClosed { + m.l.Error("debug ssh jump server stopped", "err", err) + } +} + +func (m *Mill) newJumpServer(hostKeyPath string, executorPort uint32, maxConnections int) (*ssh.Server, error) { + if executorPort == 0 { + executorPort = 2223 + } + if maxConnections <= 0 { + return nil, fmt.Errorf("max jump connections must be greater than zero") + } + if err := ensureJumpHostKey(hostKeyPath); err != nil { + return nil, fmt.Errorf("prepare jump host key: %w", err) + } + limiter := newJumpConnectionLimiter(maxConnections, maxJumpConnectionsPerIP) + srv := &ssh.Server{ + PublicKeyHandler: func(ctx ssh.Context, _ ssh.PublicKey) bool { + return ctx.User() == "debug" + }, + ConnCallback: func(ctx ssh.Context, conn net.Conn) net.Conn { + if !limiter.acquire(conn.RemoteAddr()) { + _ = conn.Close() + return conn + } + go func() { + <-ctx.Done() + limiter.release(conn.RemoteAddr()) + }() + return conn + }, + LocalPortForwardingCallback: func(ctx ssh.Context, host string, port uint32) bool { + if port != executorPort || !m.hasLiveExecutor(host) { + return false + } + ctx.Lock() + defer ctx.Unlock() + if opened, _ := ctx.Value(jumpRouteOpened).(bool); opened { + return false + } + ctx.SetValue(jumpRouteOpened, true) + return true + }, + ChannelHandlers: map[string]ssh.ChannelHandler{ + "direct-tcpip": ssh.DirectTCPIPHandler, + }, + IdleTimeout: jumpIdleTimeout, + MaxTimeout: jumpMaxTimeout, + } + if err := srv.SetOption(ssh.HostKeyFile(hostKeyPath)); err != nil { + return nil, fmt.Errorf("load jump host key: %w", err) + } + return srv, nil +} + +type jumpConnectionLimiter struct { + mu sync.Mutex + total int + perIP map[string]int + maxTotal int + maxPerIP int +} + +func newJumpConnectionLimiter(maxTotal, maxPerIP int) *jumpConnectionLimiter { + return &jumpConnectionLimiter{ + perIP: make(map[string]int), + maxTotal: maxTotal, + maxPerIP: maxPerIP, + } +} + +func (l *jumpConnectionLimiter) acquire(addr net.Addr) bool { + host := jumpRemoteHost(addr) + l.mu.Lock() + defer l.mu.Unlock() + if l.total >= l.maxTotal || l.perIP[host] >= l.maxPerIP { + return false + } + l.total++ + l.perIP[host]++ + return true +} + +func (l *jumpConnectionLimiter) release(addr net.Addr) { + host := jumpRemoteHost(addr) + l.mu.Lock() + defer l.mu.Unlock() + l.total-- + l.perIP[host]-- + if l.perIP[host] == 0 { + delete(l.perIP, host) + } +} + +func jumpRemoteHost(addr net.Addr) string { + if addr == nil { + return "" + } + host, _, err := net.SplitHostPort(addr.String()) + if err != nil { + return addr.String() + } + return host +} + +func ensureJumpHostKey(path string) error { + if _, err := os.Stat(path); err == nil { + return nil + } else if !os.IsNotExist(err) { + return fmt.Errorf("stat host key: %w", err) + } + + _, privateKey, err := ed25519.GenerateKey(rand.Reader) + if err != nil { + return fmt.Errorf("generate host key: %w", err) + } + block, err := gossh.MarshalPrivateKey(privateKey, "") + if err != nil { + return fmt.Errorf("marshal host key: %w", err) + } + if err := os.MkdirAll(filepath.Dir(path), 0o700); err != nil { + return fmt.Errorf("create host key directory: %w", err) + } + return os.WriteFile(path, pem.EncodeToMemory(block), 0o600) +} diff --git a/spindle/mill/jump_test.go b/spindle/mill/jump_test.go new file mode 100644 index 00000000..12244400 --- /dev/null +++ b/spindle/mill/jump_test.go @@ -0,0 +1,158 @@ +package mill + +import ( + "io" + "log/slog" + "net" + "os" + "strconv" + "testing" + "time" + + gossh "golang.org/x/crypto/ssh" +) + +func TestJumpConnectionLimiter(t *testing.T) { + limiter := newJumpConnectionLimiter(2, 1) + first := &net.TCPAddr{IP: net.ParseIP("192.0.2.1"), Port: 1000} + sameIP := &net.TCPAddr{IP: net.ParseIP("192.0.2.1"), Port: 1001} + second := &net.TCPAddr{IP: net.ParseIP("192.0.2.2"), Port: 1000} + third := &net.TCPAddr{IP: net.ParseIP("192.0.2.3"), Port: 1000} + + if !limiter.acquire(first) { + t.Fatal("rejected first connection") + } + if limiter.acquire(sameIP) { + t.Fatal("accepted a second connection from an exhausted IP") + } + if !limiter.acquire(second) { + t.Fatal("rejected connection within the global limit") + } + if limiter.acquire(third) { + t.Fatal("accepted a connection beyond the global limit") + } + limiter.release(first) + if !limiter.acquire(sameIP) { + t.Fatal("did not release the per-IP slot") + } +} + +func TestJumpServerForwardsOnlyLiveExecutorRoute(t *testing.T) { + backend := startJumpBackend(t) + _, portText, err := net.SplitHostPort(backend.Addr().String()) + if err != nil { + t.Fatal(err) + } + port, err := strconv.ParseUint(portText, 10, 32) + if err != nil { + t.Fatal(err) + } + + m := New(slog.New(slog.NewTextHandler(io.Discard, nil)), Config{}) + m.sessions["127.0.0.1"] = newSession("127.0.0.1", "epoch", nil, nil, m.l) + + hostKeyPath := t.TempDir() + "/host-key" + srv, err := m.newJumpServer(hostKeyPath, uint32(port), 2) + if err != nil { + t.Fatal(err) + } + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + go func() { _ = srv.Serve(listener) }() + t.Cleanup(func() { _ = srv.Close() }) + + if unauthorized, err := dialJump(t, listener.Addr().String(), "operator"); err == nil { + _ = unauthorized.Close() + t.Fatal("jump server accepted a non-debug user") + } + + client := jumpClient(t, listener.Addr().String()) + defer client.Close() + + if _, err := client.Dial("tcp", net.JoinHostPort("missing", portText)); err == nil { + t.Fatal("forwarded an executor without a live mill session") + } + wrongPort := strconv.FormatUint(port+1, 10) + if _, err := client.Dial("tcp", net.JoinHostPort("127.0.0.1", wrongPort)); err == nil { + t.Fatal("forwarded an executor on an unconfigured port") + } + + conn, err := client.Dial("tcp", net.JoinHostPort("127.0.0.1", portText)) + if err != nil { + t.Fatal(err) + } + defer conn.Close() + if _, err := conn.Write([]byte("hello")); err != nil { + t.Fatal(err) + } + got := make([]byte, 5) + if _, err := io.ReadFull(conn, got); err != nil { + t.Fatal(err) + } + if string(got) != "hello" { + t.Fatalf("forwarded payload = %q", got) + } + if _, err := client.Dial("tcp", net.JoinHostPort("127.0.0.1", portText)); err == nil { + t.Fatal("forwarded a second route on one jump connection") + } +} + +func startJumpBackend(t *testing.T) net.Listener { + t.Helper() + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = listener.Close() }) + go func() { + for { + conn, err := listener.Accept() + if err != nil { + return + } + go func() { + defer conn.Close() + _, _ = io.Copy(conn, conn) + }() + } + }() + return listener +} + +func jumpClient(t *testing.T, address string) *gossh.Client { + t.Helper() + client, err := dialJump(t, address, "debug") + if err != nil { + t.Fatal(err) + } + return client +} + +func dialJump(t *testing.T, address, user string) (*gossh.Client, error) { + t.Helper() + privateKey, err := gossh.ParsePrivateKey(testPrivateKey(t)) + if err != nil { + t.Fatal(err) + } + return gossh.Dial("tcp", address, &gossh.ClientConfig{ + User: user, + Auth: []gossh.AuthMethod{gossh.PublicKeys(privateKey)}, + HostKeyCallback: gossh.InsecureIgnoreHostKey(), + Timeout: 5 * time.Second, + }) +} + +func testPrivateKey(t *testing.T) []byte { + t.Helper() + path := t.TempDir() + "/key" + if err := ensureJumpHostKey(path); err != nil { + t.Fatal(err) + } + key, err := os.ReadFile(path) + if err != nil { + t.Fatal(err) + } + return key +} diff --git a/spindle/mill/mill.go b/spindle/mill/mill.go index ed44f956..80ec64d2 100644 --- a/spindle/mill/mill.go +++ b/spindle/mill/mill.go @@ -250,6 +250,13 @@ func (m *Mill) sessionReady(sess *millSession) { m.notifyChange() } +func (m *Mill) hasLiveExecutor(nodeID string) bool { + m.mu.Lock() + defer m.mu.Unlock() + sess := m.sessions[nodeID] + return sess != nil && sess.live(m.cfg.ReconnectGrace) +} + func (m *Mill) cancelledLeasesForNode(nodeID string) []*RemoteLease { m.mu.Lock() var candidates []*RemoteLease diff --git a/spindle/server.go b/spindle/server.go index 4895a1fb..1a8a2355 100644 --- a/spindle/server.go +++ b/spindle/server.go @@ -322,6 +322,10 @@ func (s *Spindle) Start(ctx context.Context) error { go s.exec.Connect(ctx) } + if s.mill != nil && s.cfg.Mill.JumpListenAddr != "" { + go s.mill.ServeJump(ctx, s.cfg.Mill.JumpListenAddr, s.cfg.Mill.JumpHostKeyPath, s.cfg.Mill.DebugExecutorPort, s.cfg.Mill.MaxJumpConnections) + } + if stopper, ok := s.vault.(secrets.Stopper); ok { defer stopper.Stop() } -- 2.51.2