From 296c7bb274631f37c3d075e74215f86b7bd918f1 Mon Sep 17 00:00:00 2001 From: Lemon-miaow Date: Wed, 7 Oct 2026 00:46:17 +0800 Subject: [PATCH] feat: manage distributed node operations from Owner panel --- cmd/felis/api.go | 5 + cmd/felis/manifests.go | 8 + cmd/felis/node_approve.go | 2 +- cmd/felis/node_control.go | 581 ++++++++++++++++++ cmd/felis/node_control_test.go | 40 ++ cmd/felis/run.go | 2 + deploy/bootstrap.sh | 52 ++ docs/distributed.md | 14 + docs/openapi.yaml | 123 ++++ internal/api/api.go | 5 + internal/api/handlers_distributed.go | 14 + internal/api/handlers_node_control.go | 115 ++++ internal/api/handlers_node_control_test.go | 122 ++++ internal/api/k8scluster.go | 26 +- internal/api/k8scluster_test.go | 23 + internal/nodecontrol/control.go | 447 ++++++++++++++ internal/nodecontrol/control_test.go | 192 ++++++ internal/platform/distributed_test.go | 40 ++ internal/platform/identities.go | 2 + internal/platform/workloads.go | 14 +- panel/dev/mockApi.ts | 4 + panel/e2e/startup-platform.smoke.spec.ts | 29 + .../src/components/NodeControlPanel.test.tsx | 63 ++ panel/src/components/NodeControlPanel.tsx | 90 +++ panel/src/i18n/resources/en-US/admin.json | 54 +- panel/src/i18n/resources/en-US/errors.json | 4 +- panel/src/i18n/resources/zh-CN/admin.json | 54 +- panel/src/i18n/resources/zh-CN/errors.json | 4 +- panel/src/lib/api.ts | 10 + panel/src/lib/openapi.gen.ts | 261 ++++++++ panel/src/lib/types.ts | 20 + .../pages/admin/PlatformSettingsPage.test.tsx | 10 +- .../src/pages/admin/PlatformSettingsPage.tsx | 13 +- 33 files changed, 2417 insertions(+), 26 deletions(-) create mode 100644 cmd/felis/node_control.go create mode 100644 cmd/felis/node_control_test.go create mode 100644 internal/api/handlers_node_control.go create mode 100644 internal/api/handlers_node_control_test.go create mode 100644 internal/nodecontrol/control.go create mode 100644 internal/nodecontrol/control_test.go create mode 100644 panel/src/components/NodeControlPanel.test.tsx create mode 100644 panel/src/components/NodeControlPanel.tsx diff --git a/cmd/felis/api.go b/cmd/felis/api.go index 76eeda8..9d6cc52 100644 --- a/cmd/felis/api.go +++ b/cmd/felis/api.go @@ -26,6 +26,7 @@ import ( "felis.lolicon.best/internal/mail" "felis.lolicon.best/internal/metrics" "felis.lolicon.best/internal/naming" + "felis.lolicon.best/internal/nodecontrol" "felis.lolicon.best/internal/panel" "felis.lolicon.best/internal/passkey" "felis.lolicon.best/internal/placement" @@ -549,6 +550,10 @@ func cmdAPI(args []string, stdout, stderr io.Writer) int { a.Distribution = distribution go reconcileDistribution(ctx, distribution, stderr) } + if socket := os.Getenv("FELIS_NODE_CONTROL_SOCKET"); socket != "" { + a.NodeControl = nodecontrol.NewClient(socket) + cluster.NodeMaintenanceGuard = api.NodeMaintenanceGuard(a.NodeControl) + } servers := []*http.Server{internalSrv, externalSrv} if httpsSrv != nil { servers = append(servers, httpsSrv) diff --git a/cmd/felis/manifests.go b/cmd/felis/manifests.go index 764287c..f607def 100644 --- a/cmd/felis/manifests.go +++ b/cmd/felis/manifests.go @@ -5,6 +5,7 @@ import ( "fmt" "io" "net" + "path/filepath" "strings" "felis.lolicon.best/internal/platform" @@ -48,6 +49,8 @@ func cmdManifests(args []string, stdout, stderr io.Writer) int { probe := fs.String("egress-probe", "", "reachable controller host:port denied to game Pods") var registryNodes multiFlag fs.Var(®istryNodes, "registry-node-cidr", "exact node pull source for the registry (repeatable)") + socket := fs.String("node-control-socket", "", "host node-control Unix socket (optional, API only)") + nodeControlNode := fs.String("node-control-node", "", "controller hostname hosting the socket") controlNS := fs.String("control-namespace", platform.DefaultControlNamespace, "namespace the control plane (api/operator/reaper) runs in") minecraftNS := fs.String("minecraft-namespace", platform.DefaultMinecraftNamespace, "namespace MinecraftServer workloads run in") buildNS := fs.String("build-namespace", platform.DefaultBuildNamespace, "namespace image-build Jobs run in") @@ -181,7 +184,12 @@ func cmdManifests(args []string, stdout, stderr io.Writer) int { } } + if *socket != "" && (!filepath.IsAbs(*socket) || *nodeControlNode == "") { + fmt.Fprintln(stderr, "node-control requires an absolute socket and its controller hostname") + return 2 + } params := platform.Params{ + NodeControlSocket: *socket, NodeControlNode: *nodeControlNode, Distributed: *distributed, ControllerNode: *controller, EgressProbe: *probe, RegistryNodeCIDRs: registryNodes, ControlNamespace: *controlNS, MinecraftNamespace: *minecraftNS, diff --git a/cmd/felis/node_approve.go b/cmd/felis/node_approve.go index 59bc312..f6257cd 100644 --- a/cmd/felis/node_approve.go +++ b/cmd/felis/node_approve.go @@ -173,7 +173,7 @@ func approveNode(ctx context.Context, cl client.Client, cs kubernetes.Interface, "proof=$(mktemp); trap 'rm -f \"$proof\"' EXIT\n" + "if /usr/local/bin/k3s kubectl --kubeconfig \"$kubeconfig\" label node " + shellQuote(name) + " felis.node-restriction.kubernetes.io/probe=controller --overwrite 2>\"$proof\"; then echo 'NodeRestriction failed' >&2;exit 1;fi\ngrep -qi forbidden \"$proof\"\n" + "image=" + shellQuote(image) + "\nif [ -n \"$(/usr/local/bin/k3s crictl images -q \"$image\")\" ]; then /usr/local/bin/k3s crictl rmi \"$image\" >/dev/null; fi\ntest -z \"$(/usr/local/bin/k3s crictl images -q \"$image\")\"\n/usr/local/bin/k3s crictl pull \"$image\" >/dev/null\n" - cmd := exec.CommandContext(ctx, "ssh", "--", remote, "if [ \"$(id -u)\" = 0 ]; then bash -s; else sudo -n bash -s; fi") + cmd := exec.CommandContext(ctx, "ssh", "-o", "BatchMode=yes", "-o", "StrictHostKeyChecking=yes", "-o", "ConnectTimeout=10", "--", remote, "if [ \"$(id -u)\" = 0 ]; then bash -s; else sudo -n bash -s; fi") cmd.Stdin = strings.NewReader(script) cmd.Stdout = stdout cmd.Stderr = stdout diff --git a/cmd/felis/node_control.go b/cmd/felis/node_control.go new file mode 100644 index 0000000..bbe8df7 --- /dev/null +++ b/cmd/felis/node_control.go @@ -0,0 +1,581 @@ +package main + +import ( + "context" + "encoding/json" + "errors" + "flag" + "fmt" + "io" + "net" + "net/http" + "os" + "os/exec" + "os/signal" + "path/filepath" + "runtime" + "slices" + "strings" + "syscall" + "time" + + felis "felis.lolicon.best" + "felis.lolicon.best/internal/apis/felis/v1alpha1" + "felis.lolicon.best/internal/nodecontrol" + "felis.lolicon.best/internal/placement" + "felis.lolicon.best/internal/platform" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/types" + "k8s.io/client-go/kubernetes" + "k8s.io/client-go/tools/clientcmd" +) + +func cmdNodeControl(args []string, stdout, stderr io.Writer) int { + fs := flag.NewFlagSet("node-control", flag.ContinueOnError) + fs.SetOutput(stderr) + socket := fs.String("socket", nodecontrol.Socket, "local API-only Unix socket") + dir := fs.String("state", "/var/lib/felis/node-control", "root-owned persistent task state") + kubeconfig := fs.String("kubeconfig", "/etc/rancher/k3s/k3s.yaml", "controller kubeconfig") + config := fs.String("config", "/etc/felis/felis.host.toml", "host configuration") + ns := fs.String("namespace", platform.DefaultMinecraftNamespace, "world namespace") + controlNS := fs.String("control-namespace", platform.DefaultControlNamespace, "control namespace") + if err := fs.Parse(args); err != nil { + return 2 + } + if os.Geteuid() != 0 { + fmt.Fprintln(stderr, "node-control requires root") + return 1 + } + cfg, err := clientcmd.BuildConfigFromFlags("", *kubeconfig) + if err != nil { + fmt.Fprintln(stderr, err) + return 1 + } + cs, err := kubernetes.NewForConfig(cfg) + if err != nil { + fmt.Fprintln(stderr, err) + return 1 + } + exe, err := os.Executable() + if err != nil { + fmt.Fprintln(stderr, err) + return 1 + } + ctx, cancel := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) + defer cancel() + execute := nodeExecutor{cs: cs, binary: exe, kubeconfig: *kubeconfig, config: *config, namespace: *ns, controlNamespace: *controlNS, stateDir: *dir} + socketDir := filepath.Dir(*socket) + if err = os.MkdirAll(socketDir, 0750); err != nil { + fmt.Fprintln(stderr, err) + return 1 + } + if err = os.Chown(socketDir, 0, 65532); err != nil { + fmt.Fprintln(stderr, err) + return 1 + } + if err = os.Chmod(socketDir, 0750); err != nil { + fmt.Fprintln(stderr, err) + return 1 + } + + if existing, err := os.Lstat(*socket); err == nil { + if existing.Mode()&os.ModeSocket == 0 { + fmt.Fprintln(stderr, "refusing to replace non-socket path") + return 1 + } + conn, err := net.DialTimeout("unix", *socket, time.Second) + if err == nil { + conn.Close() + fmt.Fprintln(stderr, "node-control is already running") + return 1 + } + if err = os.Remove(*socket); err != nil { + fmt.Fprintln(stderr, err) + return 1 + } + } + listener, err := net.Listen("unix", *socket) + if err != nil { + fmt.Fprintln(stderr, err) + return 1 + } + defer listener.Close() + if err = os.Chown(*socket, 0, 65532); err == nil { + err = os.Chmod(*socket, 0660) + } + if err != nil { + fmt.Fprintln(stderr, err) + return 1 + } + // /run is recreated at boot; label both the directory and the new socket inode. + if os.Getenv("FELIS_NODE_CONTROL_SELINUX") == "1" { + if err := exec.Command("chcon", "-R", "-t", "felis_node_control_socket_t", socketDir).Run(); err != nil { + fmt.Fprintln(stderr, "node-control socket labeling failed:", err) + return 1 + } + } + manager, err := nodecontrol.Open(ctx, *dir, execute.run) + if err != nil { + fmt.Fprintln(stderr, err) + return 1 + } + server := &http.Server{Handler: manager.Handler(), ReadHeaderTimeout: 5 * time.Second, ReadTimeout: 15 * time.Second, WriteTimeout: 15 * time.Second} + go func() { <-ctx.Done(); server.Close() }() + fmt.Fprintln(stdout, "node-control listening on", *socket) + err = server.Serve(listener) + cancel() + manager.Wait() + if err != nil && !errors.Is(err, http.ErrServerClosed) { + fmt.Fprintln(stderr, err) + return 1 + } + return 0 +} + +type nodeExecutor struct { + cs kubernetes.Interface + binary, kubeconfig, config, namespace, controlNamespace, stateDir string +} + +func (e nodeExecutor) run(ctx context.Context, r nodecontrol.Request, stage func(string) error, out io.Writer) (resultErr error) { + // Host-wide changes and cache-removing admission probes are serialized by the manager. + if r.Action == "approve" { + node, err := e.cs.CoreV1().Nodes().Get(ctx, r.Name, metav1.GetOptions{}) + if err != nil { + return err + } + if node.Labels[placement.LabelRole] != placement.RoleWorker { + return errors.New("only worker nodes can be approved") + } + patch := []byte(`{"spec":{"unschedulable":true}}`) + if !node.Spec.Unschedulable { + patch = []byte(`{"spec":{"unschedulable":true},"metadata":{"annotations":{"felis.lolicon.best/node-control-cordon":"true"}}}`) + } + if _, err = e.cs.CoreV1().Nodes().Patch(ctx, r.Name, types.MergePatchType, patch, metav1.PatchOptions{}); err != nil { + return err + } + } + if err := e.requireStopped(ctx, r); err != nil { + return err + } + nodes, err := e.cs.CoreV1().Nodes().List(ctx, metav1.ListOptions{}) + if err != nil { + return err + } + controller, peers, err := nodeController(nodes.Items, r.ExternalIP, r.Peers) + if err != nil { + return err + } + deployment, err := e.cs.AppsV1().Deployments(e.controlNamespace).Get(ctx, platform.SAAPI, metav1.GetOptions{}) + if err != nil { + return err + } + if len(deployment.Spec.Template.Spec.Containers) == 0 { + return errors.New("Felis API image is unavailable") + } + image := deployment.Spec.Template.Spec.Containers[0].Image + if r.Action == "enable" { + if err := stage("pause_operator"); err != nil { + return err + } + scale, err := e.cs.AppsV1().Deployments(e.controlNamespace).GetScale(ctx, platform.SAOperator, metav1.GetOptions{}) + if err != nil { + return err + } + replicas := scale.Spec.Replicas + if replicas < 1 { + replicas = 1 + } + scale.Spec.Replicas = 0 + if _, err = e.cs.AppsV1().Deployments(e.controlNamespace).UpdateScale(ctx, platform.SAOperator, scale, metav1.UpdateOptions{}); err != nil { + return err + } + defer func() { + cleanup, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + latest, err := e.cs.AppsV1().Deployments(e.controlNamespace).GetScale(cleanup, platform.SAOperator, metav1.GetOptions{}) + if err == nil { + latest.Spec.Replicas = replicas + _, err = e.cs.AppsV1().Deployments(e.controlNamespace).UpdateScale(cleanup, platform.SAOperator, latest, metav1.UpdateOptions{}) + } + if err != nil { + fmt.Fprintln(out, "Operator restoration failed:", err) + resultErr = fmt.Errorf("operator restoration failed: %w", err) + } + }() + // Drain the old reconciler before the second stopped-state check. + for { + pods, err := e.cs.CoreV1().Pods(e.controlNamespace).List(ctx, metav1.ListOptions{LabelSelector: platform.LabelComponent + "=" + platform.ComponentOperator}) + if err != nil { + return err + } + if len(pods.Items) == 0 { + break + } + select { + case <-ctx.Done(): + return ctx.Err() + case <-time.After(time.Second): + } + } + if err = e.requireStopped(ctx, r); err != nil { + return err + } + return e.enable(ctx, r, controller, peers, stage, out) + } + if controller.Labels[placement.LabelRole] != placement.RoleController { + return errors.New("distributed mode is not configured on the controller") + } + if r.Name == controller.Name { + return errors.New("worker name must differ from the controller") + } + // SSH trust and keys are configured on A; no password, key or bootstrap token crosses the API. + if err := stage("ssh_check"); err != nil { + return err + } + arch := map[string]string{"amd64": "x86_64", "arm64": "aarch64"}[runtime.GOARCH] + if arch == "" { + return errors.New("unsupported host architecture") + } + check := "set -eu\n[ \"$(uname -s)\" = Linux ] || { echo 'worker requires Linux' >&2; exit 1; }\n[ \"$(uname -m)\" = " + shellQuote(arch) + " ] || { echo 'worker architecture differs from controller' >&2; exit 1; }\n" + if r.Action == "join" { + check += "[ ! -e /etc/rancher/k3s/k3s.yaml ] && [ ! -e /var/lib/rancher/k3s/agent ] || { echo 'k3s is already installed; enrollment will not overwrite an existing node' >&2; exit 1; }\n" + check += "ip -o address show | awk '{print $4}' | cut -d/ -f1 | grep -Fx -- " + shellQuote(r.ExternalIP) + " >/dev/null || { echo 'fixed node IP is not assigned to the worker' >&2; exit 1; }\n" + } + + if err = e.ssh(ctx, r.SSHTarget, "if [ \"$(id -u)\" = 0 ]; then bash -s; else sudo -n bash -s; fi", strings.NewReader(check), out); err != nil { + return fmt.Errorf("worker preflight or SSH connection failed: %w", err) + } + if r.Action == "join" { + for _, n := range nodes.Items { + if n.Name == r.Name { + return errors.New("node already exists; use approval instead of reinstalling it") + } + } + if err = e.join(ctx, r, controller, peers, stage, out); err != nil { + return err + } + } + if err = stage("approval"); err != nil { + return err + } + cmd := exec.CommandContext(ctx, e.binary, "node", "approve", "--name", r.Name, "--ssh-target", r.SSHTarget, "--image", image, "--kubeconfig", e.kubeconfig, "--namespace", e.namespace, "--control-namespace", e.controlNamespace) + cmd.Stdout = out + cmd.Stderr = out + if err = cmd.Run(); err != nil { + return fmt.Errorf("node approval failed; node remains quarantined: %w", err) + } + // Admission succeeded; release only the task's scheduling quarantine. + admitted, err := e.cs.CoreV1().Nodes().Get(ctx, r.Name, metav1.GetOptions{}) + if err != nil { + return err + } + if admitted.Annotations["felis.lolicon.best/node-control-cordon"] == "true" { + if _, err = e.cs.CoreV1().Nodes().Patch(ctx, r.Name, types.MergePatchType, []byte(`{"spec":{"unschedulable":false},"metadata":{"annotations":{"felis.lolicon.best/node-control-cordon":null}}}`), metav1.PatchOptions{}); err != nil { + return fmt.Errorf("node approved but scheduling remains disabled: %w", err) + } + } + return stage("complete") +} +func (e nodeExecutor) requireStopped(ctx context.Context, r nodecontrol.Request) error { + // Even pre-existing worker names may contain worlds. Never run installer/admission on an occupied worker. + ss, err := e.cs.AppsV1().StatefulSets(e.namespace).List(ctx, metav1.ListOptions{}) + if err != nil { + return err + } + for _, s := range ss.Items { + if r.Action != "enable" && s.Spec.Template.Spec.NodeSelector[placement.LabelIdentity] != r.Name { + continue + } + if s.Spec.Replicas != nil && *s.Spec.Replicas > 0 { + return fmt.Errorf("StatefulSet %s still requests %d replicas; stop the server before changing nodes", s.Name, *s.Spec.Replicas) + } + } + pods, err := e.cs.CoreV1().Pods(e.namespace).List(ctx, metav1.ListOptions{}) + if err != nil { + return err + } + for _, p := range pods.Items { + if r.Action != "enable" && p.Spec.NodeName != r.Name && p.Spec.NodeSelector[placement.LabelIdentity] != r.Name { + continue + } + if p.Status.Phase != corev1.PodSucceeded && p.Status.Phase != corev1.PodFailed { + return fmt.Errorf("Pod %s has not exited (phase %s); wait for game and maintenance workloads to stop", p.Name, p.Status.Phase) + } + } + // Prevent a pending wake intent racing admission or an installation. + raw, err := e.cs.CoreV1().RESTClient().Get().AbsPath("/apis/" + v1alpha1.GroupVersion.Group + "/" + v1alpha1.GroupVersion.Version + "/namespaces/" + e.namespace + "/minecraftservers").DoRaw(ctx) + if err != nil { + return err + } + var servers v1alpha1.MinecraftServerList + if err = json.Unmarshal(raw, &servers); err != nil { + return err + } + for _, s := range servers.Items { + if r.Action != "enable" && s.Spec.NodeName != r.Name && s.Status.NodeName != r.Name { + continue + } + if s.Spec.DesiredState != "" && string(s.Spec.DesiredState) != "Stopped" { + return fmt.Errorf("server %s must have stopped intent", s.Name) + } + } + return nil +} + +func nodeController(nodes []corev1.Node, ip string, extra []string) (*corev1.Node, []string, error) { + var controller *corev1.Node + peers := append([]string{}, extra...) + for i := range nodes { + n := &nodes[i] + _, controlPlane := n.Labels["node-role.kubernetes.io/control-plane"] + if controlPlane || n.Labels[placement.LabelRole] == placement.RoleController { + if controller != nil { + return nil, nil, errors.New("exactly one controller is required") + } + controller = n + } + for _, a := range n.Status.Addresses { + if a.Type == corev1.NodeExternalIP || a.Type == corev1.NodeInternalIP { + cidr, err := exactCIDR(a.Address) + if err != nil { + return nil, nil, err + } + peers = append(peers, cidr) + } + } + } + if controller == nil || !placement.Online(controller) { + return nil, nil, errors.New("the sole controller must be online") + } + if ip != "" { + cidr, err := exactCIDR(ip) + if err != nil { + return nil, nil, err + } + peers = append(peers, cidr) + } + slices.Sort(peers) + peers = slices.Compact(peers) + return controller, peers, nil +} +func controllerIP(n *corev1.Node) string { + for _, kind := range []corev1.NodeAddressType{corev1.NodeExternalIP, corev1.NodeInternalIP} { + for _, a := range n.Status.Addresses { + if a.Type == kind { + return a.Address + } + } + } + return "" +} +func (e nodeExecutor) ssh(ctx context.Context, target, command string, in io.Reader, out io.Writer) error { + cmd := exec.CommandContext(ctx, "ssh", "-o", "BatchMode=yes", "-o", "StrictHostKeyChecking=yes", "-o", "ConnectTimeout=10", "--", target, command) + cmd.Stdin = in + cmd.Stdout = out + cmd.Stderr = out + return cmd.Run() +} +func (e nodeExecutor) join(ctx context.Context, r nodecontrol.Request, controller *corev1.Node, peers []string, stage func(string) error, out io.Writer) error { + registry, err := e.cs.CoreV1().Services(e.controlNamespace).Get(ctx, "registry", metav1.GetOptions{}) + if err != nil { + return err + } + if err = stage("firewall"); err != nil { + return err + } + // All peers must be updated before the new agent is admitted. Existing worker SSH targets must be configured by stable node name. + list, err := e.cs.CoreV1().Nodes().List(ctx, metav1.ListOptions{}) + if err != nil { + return err + } + for _, n := range list.Items { + if n.Name == controller.Name || n.Name == r.Name { + continue + } + if n.Labels[placement.LabelRole] != placement.RoleWorker { + return fmt.Errorf("node %s has an unknown role", n.Name) + } + script := shellQuote(e.binary) + " node firewall --peers " + shellQuote(strings.Join(peers, ",")) + " --controller-ip " + shellQuote(controllerIP(controller)) + script += " --namespace " + shellQuote(e.namespace) + " --control-namespace " + shellQuote(e.controlNamespace) + if err = e.ssh(ctx, n.Name, "if [ \"$(id -u)\" = 0 ]; then "+script+"; else sudo -n "+script+"; fi", nil, out); err != nil { + return fmt.Errorf("update peer firewall on %s: %w", n.Name, err) + } + } + cmd := exec.CommandContext(ctx, e.binary, "node", "firewall", "--controller", "--controller-ip", controllerIP(controller), "--peers", strings.Join(peers, ","), "--namespace", e.namespace, "--control-namespace", e.controlNamespace) + cmd.Stdout = out + cmd.Stderr = out + if err = cmd.Run(); err != nil { + return err + } + if err = stage("bootstrap_token"); err != nil { + return err + } + temp, err := os.MkdirTemp("", "felis-node-join-") + if err != nil { + return err + } + defer os.RemoveAll(temp) + tokenPath := filepath.Join(temp, "bootstrap") + tokenCmd := exec.CommandContext(ctx, e.binary, "node", "token", "--name", r.Name, "--ttl", "1h", "--out", tokenPath) + tokenCmd.Stdout = out + tokenCmd.Stderr = out + if err = tokenCmd.Run(); err != nil { + return err + } + // Revoke even on failure. The worker exchanges bootstrap credentials for its node identity. + token, err := os.ReadFile(tokenPath) + if err != nil { + return err + } + defer func() { + cleanup, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + exec.CommandContext(cleanup, "k3s", "token", "delete", bootstrapTokenID(string(token))).Run() + }() + remoteDir := "/var/tmp/felis-node-" + filepath.Base(temp) + if err = stage("transfer"); err != nil { + return err + } + script := "set -eu; umask 077; mkdir " + shellQuote(remoteDir) + "; cat > " + shellQuote(remoteDir+"/bootstrap") + rootCommand := "if [ \"$(id -u)\" = 0 ]; then bash -s; else sudo -n bash -s; fi" + defer func() { + cleanup, cancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cancel() + e.ssh(cleanup, r.SSHTarget, rootCommand, strings.NewReader("rm -rf -- "+shellQuote(remoteDir)), io.Discard) + }() + // stdin carries the token; it is never in command arguments, task records or logs. + if err = e.ssh(ctx, r.SSHTarget, "if [ \"$(id -u)\" = 0 ]; then "+script+"; else sudo -n sh -c "+shellQuote(script)+"; fi", strings.NewReader(string(token)), out); err != nil { + return err + } + + binary, err := os.Open(e.binary) + if err != nil { + return err + } + defer binary.Close() + script = "set -eu; cat > " + shellQuote(remoteDir+"/felis") + "; chmod 0700 " + shellQuote(remoteDir+"/felis") + if err = e.ssh(ctx, r.SSHTarget, "if [ \"$(id -u)\" = 0 ]; then "+script+"; else sudo -n sh -c "+shellQuote(script)+"; fi", binary, out); err != nil { + return err + } + if err = stage("install_worker"); err != nil { + return err + } + script = "set -eu\nexport FELIS_K3S_VERSION=" + shellQuote(controller.Status.NodeInfo.KubeletVersion) + "\n" + shellQuote(remoteDir+"/felis") + " node join --name " + shellQuote(r.Name) + " --server " + shellQuote("https://"+net.JoinHostPort(controllerIP(controller), "6443")) + " --external-ip " + shellQuote(r.ExternalIP) + " --token-file " + shellQuote(remoteDir+"/bootstrap") + " --registry-ip " + shellQuote(registry.Spec.ClusterIP) + " --peers " + shellQuote(strings.Join(peers, ",")) + "\n" + if err = e.ssh(ctx, r.SSHTarget, rootCommand, strings.NewReader(script), out); err != nil { + return fmt.Errorf("worker installation failed: %w", err) + } + if err = stage("wait_ready"); err != nil { + return err + } + deadline := time.NewTimer(5 * time.Minute) + defer deadline.Stop() + ticker := time.NewTicker(3 * time.Second) + defer ticker.Stop() + for { + n, err := e.cs.CoreV1().Nodes().Get(ctx, r.Name, metav1.GetOptions{}) + if err == nil && placement.Online(n) { + return nil + } + select { + case <-ctx.Done(): + return ctx.Err() + case <-deadline.C: + return errors.New("worker did not become Ready within five minutes") + case <-ticker.C: + } + } +} +func (e nodeExecutor) enable(ctx context.Context, r nodecontrol.Request, controller *corev1.Node, peers []string, stage func(string) error, out io.Writer) error { + if controllerIP(controller) != r.ExternalIP { + return errors.New("controller IP must match the registered fixed node address") + } + if len(peers) == 0 { + return errors.New("peer addresses are required") + } + if err := stage("database_backup"); err != nil { + return err + } + backup := exec.CommandContext(ctx, e.binary, "db", "backup", "-config", e.config, "-dir", "/var/lib/felis/db-backups", "-label", "pre-migrate") + backup.Stdout = out + backup.Stderr = out + if err := backup.Run(); err != nil { + return err + } + if err := stage("cluster_backup"); err != nil { + return err + } + if err := e.backupCluster(ctx, out); err != nil { + return err + } + if err := stage("configure_controller"); err != nil { + return err + } + patch := fmt.Sprintf(`{"metadata":{"labels":{"%s":"%s","%s":"controller"}}}`, placement.LabelIdentity, controller.Name, placement.LabelRole) + if _, err := e.cs.CoreV1().Nodes().Patch(ctx, controller.Name, types.MergePatchType, []byte(patch), metav1.PatchOptions{}); err != nil { + return err + } + for _, name := range []string{platform.SAAPI, platform.SAOperator, platform.PostgresName, "registry"} { + body := fmt.Sprintf(`{"spec":{"template":{"spec":{"nodeSelector":{"%s":"%s"}}}}}`, placement.LabelIdentity, controller.Name) + if _, err := e.cs.AppsV1().Deployments(e.controlNamespace).Patch(ctx, name, types.MergePatchType, []byte(body), metav1.PatchOptions{}); err != nil { + return err + } + } + if err := stage("install_controller"); err != nil { + return err + } + cmd := exec.CommandContext(ctx, "bash", "-s") + cmd.Stdin = strings.NewReader(felis.BootstrapScript()) + cmd.Stdout = out + cmd.Stderr = out + cmd.Env = append(os.Environ(), "FELIS_DISTRIBUTED=1", "FELIS_NODE_EXTERNAL_IP="+r.ExternalIP, "FELIS_PEER_CIDRS="+strings.Join(peers, ","), "FELIS_BOOTSTRAP_FROM_TUI=1", "FELIS_BOOTSTRAP_BINARY="+e.binary, "FELIS_NO_SETUP=1", "FELIS_NODE_CONTROL_TASK=1", "FELIS_INSTALL_MODE=full") + if err := cmd.Run(); err != nil { + return fmt.Errorf("controller installation failed; retain the database backup: %w", err) + } + return stage("complete") +} + +func bootstrapTokenID(token string) string { + token = strings.TrimSpace(token) + if _, short, ok := strings.Cut(token, "::"); ok { + token = short + } + id, _, _ := strings.Cut(token, ".") + return id +} + +func (e nodeExecutor) backupCluster(ctx context.Context, out io.Writer) (err error) { + dir := filepath.Join(filepath.Dir(e.stateDir), "cluster-backups") + if err = os.MkdirAll(dir, 0700); err != nil { + return err + } + path := filepath.Join(dir, "pre-distributed-"+time.Now().UTC().Format("20060102T150405.000000000Z")+".tar") + stop := exec.CommandContext(ctx, "systemctl", "stop", "k3s") + stop.Stdout = out + stop.Stderr = out + if err = stop.Run(); err != nil { + return err + } + defer func() { + restart, cancel := context.WithTimeout(context.Background(), 2*time.Minute) + defer cancel() + cmd := exec.CommandContext(restart, "systemctl", "start", "k3s") + cmd.Stdout = out + cmd.Stderr = out + if restartErr := cmd.Run(); restartErr != nil { + err = fmt.Errorf("k3s restart failed after snapshot: %w", restartErr) + } + }() + cmd := exec.CommandContext(ctx, "tar", "-cf", path, "-C", "/var/lib/rancher/k3s/server", "db", "token", "tls") + cmd.Stdout = out + cmd.Stderr = out + if err = cmd.Run(); err != nil { + return fmt.Errorf("cluster snapshot failed: %w", err) + } + if err = os.Chmod(path, 0600); err != nil { + return err + } + fmt.Fprintln(out, "Cluster snapshot:", path) + return nil +} diff --git a/cmd/felis/node_control_test.go b/cmd/felis/node_control_test.go new file mode 100644 index 0000000..8436e6f --- /dev/null +++ b/cmd/felis/node_control_test.go @@ -0,0 +1,40 @@ +package main + +import ( + "context" + "strings" + "testing" + + "felis.lolicon.best/internal/nodecontrol" + "felis.lolicon.best/internal/platform" + appsv1 "k8s.io/api/apps/v1" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/client-go/kubernetes/fake" +) + +func TestNodeControlTopologyAndBootstrapID(t *testing.T) { + node := corev1.Node{ObjectMeta: metav1.ObjectMeta{Name: "controller", Labels: map[string]string{"node-role.kubernetes.io/control-plane": "true"}}, Status: corev1.NodeStatus{Conditions: []corev1.NodeCondition{{Type: corev1.NodeReady, Status: corev1.ConditionTrue}}, Addresses: []corev1.NodeAddress{{Type: corev1.NodeInternalIP, Address: "192.0.2.1"}}}} + controller, peers, err := nodeController([]corev1.Node{node}, "192.0.2.2", []string{"192.0.2.1/32"}) + if err != nil || controller.Name != "controller" || strings.Join(peers, ",") != "192.0.2.1/32,192.0.2.2/32" { + t.Fatal(controller, peers, err) + } + duplicate := node + duplicate.Name = "second" + if _, _, err = nodeController([]corev1.Node{node, duplicate}, "", nil); err == nil { + t.Fatal("multiple controllers accepted") + } + for _, token := range []string{"abcdef.0123456789abcdef", "K10hash::abcdef.0123456789abcdef\n"} { + if bootstrapTokenID(token) != "abcdef" { + t.Fatal("token secret passed to revocation") + } + } +} +func TestNodeControlRejectsActiveWorkloadsBeforeHostCommands(t *testing.T) { + replicas := int32(1) + set := &appsv1.StatefulSet{ObjectMeta: metav1.ObjectMeta{Name: "survival", Namespace: "minecraft"}, Spec: appsv1.StatefulSetSpec{Replicas: &replicas}} + executor := nodeExecutor{cs: fake.NewSimpleClientset(set), namespace: platform.DefaultMinecraftNamespace} + if err := executor.requireStopped(context.Background(), nodecontrol.Request{Action: "enable"}); err == nil { + t.Fatal("host changes permitted with active worlds") + } +} diff --git a/cmd/felis/run.go b/cmd/felis/run.go index 310988d..ab6fe9c 100644 --- a/cmd/felis/run.go +++ b/cmd/felis/run.go @@ -30,6 +30,7 @@ Commands: mirror-build-tools Copy kaniko, trivy and Trivy's DBs into the registry (run by felis-build-tools.timer) server-migrate Move a stopped world between approved nodes (start|status|retry; requires root) node Join and approve trusted daemon nodes (list|token|join|approve|firewall; requires root) + node-control Run the host node-task service on an API-only Unix socket (requires root) node-probe Verify reachability and observed sources (internal admission probe) archive-serve Serve scoped one-use archive transfers on the controller (internal entrypoint) registry-gate Authorize registry writes in front of registry:2 (internal sidecar entrypoint) @@ -80,6 +81,7 @@ var commands = map[string]func(args []string, stdout, stderr io.Writer) int{ "registry-gate": cmdRegistryGate, "archive-serve": cmdArchiveServe, "node": cmdNode, + "node-control": cmdNodeControl, "server-migrate": cmdServerMigrate, "node-probe": cmdNodeProbe, "manifests": cmdManifests, diff --git a/deploy/bootstrap.sh b/deploy/bootstrap.sh index 2d67680..4222124 100644 --- a/deploy/bootstrap.sh +++ b/deploy/bootstrap.sh @@ -5217,6 +5217,56 @@ EOF # The daily database backup. The first run happens now, so a broken pipeline (pg_dump # missing, directory unwritable) shows up in this install rather than in the first # restore someone needs. +# The API mounts only this Unix socket, never the host kubeconfig or SSH keys. +install_node_control_service() { + local selinux_environment="" + if command -v selinuxenabled >/dev/null 2>&1 && selinuxenabled; then + command -v semodule >/dev/null 2>&1 || die "node control on SELinux requires semodule" + cat > "${STATE_DIR}/felis-node-control.cil" <<'EOF_NODE_SELINUX' +(type felis_node_control_socket_t) +(typeattributeset file_type (felis_node_control_socket_t)) +(typeattributeset non_auth_file_type (felis_node_control_socket_t)) +(allow container_t felis_node_control_socket_t (dir (search getattr open read))) +(allow container_t felis_node_control_socket_t (sock_file (write open read getattr))) +(allow container_t unconfined_service_t (unix_stream_socket (connectto))) +EOF_NODE_SELINUX + semodule -i "${STATE_DIR}/felis-node-control.cil" || die "could not install the node-control SELinux policy" + selinux_environment="Environment=FELIS_NODE_CONTROL_SELINUX=1" + fi + install -d -o root -g 65532 -m 0750 /run/felis-node-control + cat > /etc/tmpfiles.d/felis-node-control.conf <<'EOF_NODE_TMP' +d /run/felis-node-control 0750 root 65532 - +EOF_NODE_TMP + if [ -n "$selinux_environment" ] && command -v semanage >/dev/null 2>&1; then + semanage fcontext -a -t felis_node_control_socket_t '/run/felis-node-control(/.*)?' 2>/dev/null \ + || semanage fcontext -m -t felis_node_control_socket_t '/run/felis-node-control(/.*)?' + fi + if [ -n "$selinux_environment" ] && command -v chcon >/dev/null 2>&1; then + chcon -R -t felis_node_control_socket_t /run/felis-node-control 2>/dev/null || true + fi + cat > /etc/systemd/system/felis-node-control.service < "$DB_BACKUP_SERVICE" < 63 { + return false + } + for _, label := range strings.Split(name, ".") { + if !namePattern.MatchString(label) { + return false + } + } + return true +} + +type Request struct { + Action string `json:"action"` + Name string `json:"name,omitempty"` + SSHTarget string `json:"sshTarget,omitempty"` + ExternalIP string `json:"externalIP,omitempty"` + Peers []string `json:"peers,omitempty"` + ConfirmMaintenance bool `json:"confirmMaintenance"` +} + +func (r Request) Validate() error { + if !r.ConfirmMaintenance { + return errors.New("maintenance impact confirmation is required") + } + switch r.Action { + case "join", "approve": + if !validNodeName(r.Name) || !targetPattern.MatchString(r.SSHTarget) { + return errors.New("invalid worker name or SSH target") + } + case "enable": + default: + return errors.New("unsupported node operation") + } + if r.Action != "approve" { + if a, err := netip.ParseAddr(r.ExternalIP); err != nil || a.IsUnspecified() || a.IsLoopback() || a.IsMulticast() || a.Zone() != "" { + return errors.New("a fixed node IP is required") + } + } + for _, peer := range r.Peers { + p, err := netip.ParsePrefix(peer) + if err != nil || p.Bits() != p.Addr().BitLen() || p.Addr().IsLoopback() || p.Addr().IsMulticast() || p.Addr().IsUnspecified() { + return errors.New("peer addresses must use exact /32 or /128 prefixes") + } + } + if len(r.Peers) > 100 { + return errors.New("too many peer addresses") + } + return nil +} + +type Task struct { + ID string `json:"id"` + Request Request `json:"request"` + Actor string `json:"actor"` + State string `json:"state"` + Stage string `json:"stage"` + StartedAt time.Time `json:"startedAt"` + FinishedAt *time.Time `json:"finishedAt,omitempty"` + Error string `json:"error,omitempty"` + Log string `json:"log,omitempty"` +} + +type Executor func(context.Context, Request, func(string) error, io.Writer) error + +type Manager struct { + dir string + execute Executor + mu sync.Mutex + tasks map[string]Task + active string + ctx context.Context + wg sync.WaitGroup +} + +func Open(ctx context.Context, dir string, execute Executor) (*Manager, error) { + if err := os.MkdirAll(dir, 0700); err != nil { + return nil, err + } + m := &Manager{dir: dir, execute: execute, tasks: map[string]Task{}, ctx: ctx} + paths, err := filepath.Glob(filepath.Join(dir, "*.json")) + if err != nil { + return nil, err + } + for _, path := range paths { + raw, err := os.ReadFile(path) + if err != nil { + return nil, err + } + var t Task + if err = json.Unmarshal(raw, &t); err != nil { + return nil, fmt.Errorf("read node task: %w", err) + } + if !idPattern.MatchString(t.ID) || filepath.Base(path) != t.ID+".json" { + return nil, errors.New("invalid stored task ID") + } + if t.State == "running" { + t.State, t.Error = "failed", "Host execution service restarted. Verify the host state before retrying." + now := time.Now().UTC() + t.FinishedAt = &now + if err = m.persist(t); err != nil { + return nil, err + } + } + m.tasks[t.ID] = t + } + return m, nil +} +func (m *Manager) persist(t Task) error { + t.Log = "" + raw, err := json.Marshal(t) + if err != nil { + return err + } + path := filepath.Join(m.dir, t.ID+".json") + f, err := os.OpenFile(path+".tmp", os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0600) + if err != nil { + return err + } + if _, err = f.Write(raw); err == nil { + err = f.Sync() + } + closeErr := f.Close() + if err != nil { + return err + } + if closeErr != nil { + return closeErr + } + if err = os.Rename(path+".tmp", path); err != nil { + return err + } + d, err := os.Open(m.dir) + if err != nil { + return err + } + defer d.Close() + return d.Sync() +} +func (m *Manager) List() []Task { + m.mu.Lock() + defer m.mu.Unlock() + out := make([]Task, 0, len(m.tasks)) + for _, t := range m.tasks { + out = append(out, t) + } + sort.Slice(out, func(i, j int) bool { return out[i].StartedAt.After(out[j].StartedAt) }) + if len(out) > 100 { + out = out[:100] + } + return out +} +func (m *Manager) Get(id string) (Task, error) { + m.mu.Lock() + defer m.mu.Unlock() + t, ok := m.tasks[id] + if !ok { + return Task{}, ErrNotFound + } + f, err := os.Open(filepath.Join(m.dir, id+".log")) + if errors.Is(err, os.ErrNotExist) { + return t, nil + } + if err != nil { + return t, err + } + defer f.Close() + info, err := f.Stat() + if err != nil { + return t, err + } + if info.Size() > MaxLog { + if _, err = f.Seek(-MaxLog, io.SeekEnd); err != nil { + return t, err + } + } + raw, err := io.ReadAll(f) + t.Log = string(raw) + return t, err +} +func (m *Manager) Start(r Request, actor string) (Task, error) { + if err := r.Validate(); err != nil { + return Task{}, err + } + m.mu.Lock() + defer m.mu.Unlock() + if m.active != "" { + return Task{}, ErrBusy + } + if err := m.ctx.Err(); err != nil { + return Task{}, err + } + var id [16]byte + if _, err := rand.Read(id[:]); err != nil { + return Task{}, err + } + t := Task{ID: hex.EncodeToString(id[:]), Request: r, Actor: actor, State: "running", Stage: "preflight", StartedAt: time.Now().UTC()} + if err := m.persist(t); err != nil { + return Task{}, err + } + m.tasks[t.ID] = t + m.active = t.ID + m.wg.Add(1) + go m.run(t) + return t, nil +} +func (m *Manager) run(t Task) { + defer m.wg.Done() + ctx, cancel := context.WithTimeout(m.ctx, 45*time.Minute) + defer cancel() + log, err := os.OpenFile(filepath.Join(m.dir, t.ID+".log"), os.O_CREATE|os.O_RDWR|os.O_APPEND, 0600) + if err == nil { + stage := func(s string) error { + m.mu.Lock() + defer m.mu.Unlock() + t.Stage = s + if err := m.persist(t); err != nil { + return err + } + m.tasks[t.ID] = t + _, err := fmt.Fprintln(log, "[felis]", s) + return err + } + err = m.execute(ctx, t.Request, stage, &boundedLog{file: log}) + if ctx.Err() != nil { + err = ctx.Err() + } + if closeErr := log.Close(); err == nil { + err = closeErr + } + } + m.mu.Lock() + defer m.mu.Unlock() + now := time.Now().UTC() + t.FinishedAt = &now + t.State = "succeeded" + if err != nil { + t.State = "failed" + t.Error = err.Error() + if errors.Is(err, context.DeadlineExceeded) { + t.Error = "Host operation exceeded the 45-minute execution limit. Verify the host state before retrying." + } + if errors.Is(err, context.Canceled) { + t.Error = "Host operation was interrupted. Verify the host state before retrying." + } + } + if persistErr := m.persist(t); persistErr != nil { + t.State = "failed" + t.Error = "Task result persistence failed: " + persistErr.Error() + } + m.tasks[t.ID] = t + m.active = "" + for len(m.tasks) > 100 { + oldest := "" + for id, candidate := range m.tasks { + if id == t.ID { + continue + } + if oldest == "" || candidate.StartedAt.Before(m.tasks[oldest].StartedAt) { + oldest = id + } + } + if oldest != "" { + if err := os.Remove(filepath.Join(m.dir, oldest+".json")); err == nil { + os.Remove(filepath.Join(m.dir, oldest+".log")) + delete(m.tasks, oldest) + } else { + break + } + } + } +} +func (m *Manager) Wait() { m.wg.Wait() } + +// Handler is reachable only through the root-owned socket mounted into felis-api. +func (m *Manager) Handler() http.Handler { + mux := http.NewServeMux() + mux.HandleFunc("GET /tasks", func(w http.ResponseWriter, r *http.Request) { json.NewEncoder(w).Encode(m.List()) }) + mux.HandleFunc("GET /tasks/{id}", func(w http.ResponseWriter, r *http.Request) { + t, err := m.Get(r.PathValue("id")) + if err != nil { + http.Error(w, err.Error(), 404) + return + } + json.NewEncoder(w).Encode(t) + }) + mux.HandleFunc("POST /tasks", func(w http.ResponseWriter, r *http.Request) { + var body struct { + Request Request `json:"request"` + Actor string `json:"actor"` + } + dec := json.NewDecoder(http.MaxBytesReader(w, r.Body, 8192)) + dec.DisallowUnknownFields() + if err := dec.Decode(&body); err != nil { + http.Error(w, "invalid task request", 400) + return + } + if err := dec.Decode(&struct{}{}); err != io.EOF { + http.Error(w, "invalid trailing task data", 400) + return + } + if err := body.Request.Validate(); err != nil { + http.Error(w, err.Error(), 400) + return + } + t, err := m.Start(body.Request, body.Actor) + if err != nil { + code := 500 + if errors.Is(err, ErrBusy) { + code = 409 + } + http.Error(w, err.Error(), code) + return + } + w.WriteHeader(http.StatusAccepted) + json.NewEncoder(w).Encode(t) + }) + return mux +} + +type Client struct{ http *http.Client } + +func NewClient(socket string) *Client { + return &Client{http: &http.Client{Timeout: 10 * time.Second, Transport: &http.Transport{DialContext: func(ctx context.Context, _, _ string) (net.Conn, error) { + return (&net.Dialer{}).DialContext(ctx, "unix", socket) + }}}} +} +func (c *Client) call(ctx context.Context, method, path string, body any, out any) error { + var reader io.Reader + if body != nil { + raw, err := json.Marshal(body) + if err != nil { + return err + } + reader = strings.NewReader(string(raw)) + } + req, err := http.NewRequestWithContext(ctx, method, "http://node-control"+path, reader) + if err != nil { + return err + } + res, err := c.http.Do(req) + if err != nil { + return err + } + defer res.Body.Close() + if res.StatusCode == 409 { + return ErrBusy + } + if res.StatusCode == 404 { + return ErrNotFound + } + if res.StatusCode >= 400 { + return fmt.Errorf("host service returned HTTP %d", res.StatusCode) + } + return json.NewDecoder(io.LimitReader(res.Body, 1<<20)).Decode(out) +} +func (c *Client) List(ctx context.Context) ([]Task, error) { + var out []Task + err := c.call(ctx, "GET", "/tasks", nil, &out) + return out, err +} +func (c *Client) Get(ctx context.Context, id string) (Task, error) { + if !idPattern.MatchString(id) { + return Task{}, ErrNotFound + } + var out Task + err := c.call(ctx, "GET", "/tasks/"+id, nil, &out) + return out, err +} +func (c *Client) Start(ctx context.Context, r Request, actor string) (Task, error) { + var out Task + err := c.call(ctx, "POST", "/tasks", map[string]any{"request": r, "actor": actor}, &out) + return out, err +} + +// Bound disk use while retaining the most recent output of long builds. +type boundedLog struct { + file *os.File + mu sync.Mutex +} + +func (w *boundedLog) Write(p []byte) (int, error) { + w.mu.Lock() + defer w.mu.Unlock() + const limit = 4 << 20 + info, err := w.file.Stat() + if err != nil { + return 0, err + } + if info.Size()+int64(len(p)) > limit { + n := int64(MaxLog) + if info.Size() < n { + n = info.Size() + } + tail := make([]byte, n) + if _, err = w.file.ReadAt(tail, info.Size()-n); err != nil { + return 0, err + } + if err = w.file.Truncate(0); err != nil { + return 0, err + } + if _, err = w.file.Write(tail); err != nil { + return 0, err + } + } + original := len(p) + if len(p) > limit { + p = p[len(p)-limit:] + } + _, err = w.file.Write(p) + if err != nil { + return 0, err + } + return original, nil +} diff --git a/internal/nodecontrol/control_test.go b/internal/nodecontrol/control_test.go new file mode 100644 index 0000000..0c3dafa --- /dev/null +++ b/internal/nodecontrol/control_test.go @@ -0,0 +1,192 @@ +package nodecontrol + +import ( + "context" + "encoding/json" + "errors" + "io" + "net" + "net/http" + "os" + "path/filepath" + "strings" + "testing" + "time" +) + +func validRequest() Request { + return Request{Action: "join", Name: "worker-01", SSHTarget: "root@192.0.2.10", ExternalIP: "192.0.2.10", ConfirmMaintenance: true} +} +func TestRequestRejectsUnsafeInputs(t *testing.T) { + for _, change := range []func(*Request){func(r *Request) { r.ConfirmMaintenance = false }, func(r *Request) { r.SSHTarget = "-oProxyCommand=sh" }, func(r *Request) { r.SSHTarget = "root@host;id" }, func(r *Request) { r.Name = "../../escape" }, func(r *Request) { r.ExternalIP = "127.0.0.1" }, func(r *Request) { r.Peers = []string{"10.0.0.0/8"} }, func(r *Request) { r.Action = "exec" }} { + r := validRequest() + change(&r) + if r.Validate() == nil { + t.Fatalf("accepted %+v", r) + } + } + if err := validRequest().Validate(); err != nil { + t.Fatal(err) + } +} +func TestNodeNames(t *testing.T) { + for _, name := range []string{"worker-01", "worker.localdomain"} { + r := validRequest() + r.Name = name + if err := r.Validate(); err != nil { + t.Fatalf("%s: %v", name, err) + } + } + for _, name := range []string{"worker..localdomain", "worker.-domain", "worker.", strings.Repeat("a", 64)} { + r := validRequest() + r.Name = name + if r.Validate() == nil { + t.Fatalf("accepted %s", name) + } + } +} + +func TestPersistentSerializedTasksAndLogs(t *testing.T) { + dir := t.TempDir() + entered := make(chan struct{}) + finish := make(chan struct{}) + manager, err := Open(context.Background(), dir, func(ctx context.Context, r Request, stage func(string) error, out io.Writer) error { + if err := stage("approval"); err != nil { + return err + } + io.WriteString(out, "probe failed\n") + close(entered) + <-finish + return errors.New("quarantine retained") + }) + if err != nil { + t.Fatal(err) + } + task, err := manager.Start(validRequest(), "owner") + if err != nil { + t.Fatal(err) + } + <-entered + if _, err = manager.Start(validRequest(), "owner"); !errors.Is(err, ErrBusy) { + t.Fatal("concurrent task", err) + } + live, err := manager.Get(task.ID) + if err != nil || live.Stage != "approval" || !strings.Contains(live.Log, "probe failed") { + t.Fatal(live, err) + } + close(finish) + manager.Wait() + reopened, err := Open(context.Background(), dir, nil) + if err != nil { + t.Fatal(err) + } + got, err := reopened.Get(task.ID) + if err != nil || got.State != "failed" || got.Error != "quarantine retained" || got.FinishedAt == nil || got.Actor != "owner" { + t.Fatal(got, err) + } + if reopened.List()[0].Log != "" { + t.Fatal("task list leaked logs") + } +} +func TestInterruptedTaskIsNotReportedRunning(t *testing.T) { + dir := t.TempDir() + task := Task{ID: strings.Repeat("a", 32), Request: validRequest(), State: "running", StartedAt: time.Now()} + raw, _ := json.Marshal(task) + if err := os.WriteFile(filepath.Join(dir, task.ID+".json"), raw, 0600); err != nil { + t.Fatal(err) + } + manager, err := Open(context.Background(), dir, nil) + if err != nil { + t.Fatal(err) + } + got, err := manager.Get(task.ID) + if err != nil || got.State != "failed" || got.FinishedAt == nil || !strings.Contains(got.Error, "restarted") { + t.Fatal(got, err) + } +} +func TestLocalClientAndBoundedOutput(t *testing.T) { + dir := t.TempDir() + manager, err := Open(context.Background(), dir, func(ctx context.Context, r Request, stage func(string) error, out io.Writer) error { + for range 90 { + if _, err := io.WriteString(out, strings.Repeat("x", MaxLog)); err != nil { + return err + } + } + io.WriteString(out, "last output") + return nil + }) + if err != nil { + t.Fatal(err) + } + // Avoid macOS' short Unix socket path limit under t.TempDir. + sock, err := os.CreateTemp("/tmp", "felis-node-test-") + if err != nil { + t.Fatal(err) + } + path := sock.Name() + sock.Close() + os.Remove(path) + defer os.Remove(path) + listener, err := net.Listen("unix", path) + if err != nil { + t.Fatal(err) + } + server := &http.Server{Handler: manager.Handler()} + go server.Serve(listener) + defer server.Close() + client := NewClient(path) + task, err := client.Start(context.Background(), validRequest(), "owner") + if err != nil { + t.Fatal(err) + } + manager.Wait() + got, err := client.Get(context.Background(), task.ID) + if err != nil || got.State != "succeeded" || len(got.Log) > MaxLog || !strings.HasSuffix(got.Log, "last output") { + t.Fatal("task/log", got.State, len(got.Log), err) + } + stat, err := os.Stat(filepath.Join(dir, task.ID+".log")) + if err != nil || stat.Size() > 4<<20 { + t.Fatal("unbounded file", stat, err) + } + if _, err := client.Get(context.Background(), "../secret"); !errors.Is(err, ErrNotFound) { + t.Fatal(err) + } +} +func TestCancellationFinishesTask(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + entered := make(chan struct{}) + manager, err := Open(ctx, t.TempDir(), func(ctx context.Context, r Request, stage func(string) error, out io.Writer) error { + close(entered) + <-ctx.Done() + return ctx.Err() + }) + if err != nil { + t.Fatal(err) + } + task, err := manager.Start(validRequest(), "owner") + if err != nil { + t.Fatal(err) + } + <-entered + cancel() + manager.Wait() + got, err := manager.Get(task.ID) + if err != nil || got.State != "failed" { + t.Fatal(got, err) + } +} + +// Optional live check under the API container's UID and SELinux domain. +func TestHostSocketConnectivity(t *testing.T) { + socket := os.Getenv("FELIS_TEST_NODE_CONTROL_SOCKET") + if socket == "" { + t.Skip("live host socket not supplied") + } + tasks, err := NewClient(socket).List(context.Background()) + if err != nil { + t.Fatal(err) + } + if tasks == nil { + t.Fatal("host returned no task list") + } +} diff --git a/internal/platform/distributed_test.go b/internal/platform/distributed_test.go index e96351e..afc3977 100644 --- a/internal/platform/distributed_test.go +++ b/internal/platform/distributed_test.go @@ -59,3 +59,43 @@ func TestDistributedControllerPinningAndArchivePowers(t *testing.T) { t.Fatal("broad Velocity source accepted") } } + +func TestNodeControlSocketIsAPIOnly(t *testing.T) { + p := testParams() + p.NodeControlSocket = "/run/felis-node-control/control.sock" + p.NodeControlNode = "controller" + for _, object := range Objects(p) { + d, ok := object.(*appsv1.Deployment) + if !ok { + continue + } + for _, volume := range d.Spec.Template.Spec.Volumes { + if volume.Name == "node-control" && d.Name != SAAPI { + t.Fatal("host socket leaked to", d.Name) + } + } + if d.Name != SAAPI { + continue + } + if d.Spec.Template.Spec.NodeSelector["kubernetes.io/hostname"] != "controller" { + t.Fatal("API can run away from socket") + } + found := false + for _, container := range d.Spec.Template.Spec.Containers { + for _, mount := range container.VolumeMounts { + if mount.Name == "node-control" { + found = true + if !mount.ReadOnly { + t.Fatal("host directory writable") + } + } + } + if d.Spec.Template.Spec.SecurityContext.RunAsUser == nil || *d.Spec.Template.Spec.SecurityContext.RunAsUser == 0 { + t.Fatal("API elevated") + } + } + if !found { + t.Fatal("API has no node socket") + } + } +} diff --git a/internal/platform/identities.go b/internal/platform/identities.go index d85e8e6..48112a8 100644 --- a/internal/platform/identities.go +++ b/internal/platform/identities.go @@ -90,6 +90,8 @@ const ( // Params parameterises the install bundle. Namespaces and the registry location // have safe defaults; VelocityCIDRs has none — see the field comment. type Params struct { + NodeControlSocket string + NodeControlNode string Distributed bool ControllerNode string EgressProbe string diff --git a/internal/platform/workloads.go b/internal/platform/workloads.go index 565d37b..34f1264 100644 --- a/internal/platform/workloads.go +++ b/internal/platform/workloads.go @@ -4,6 +4,7 @@ import ( "crypto/sha256" "encoding/hex" "fmt" + "path/filepath" "time" "felis.lolicon.best/internal/naming" @@ -480,7 +481,18 @@ func APIDeployment(p Params) *appsv1.Deployment { }, } - return controlPlaneDeployment(p, SAAPI, container, volumes) + if p.NodeControlSocket != "" { + directory := filepath.Dir(p.NodeControlSocket) + hostType := corev1.HostPathDirectory + volumes = append(volumes, corev1.Volume{Name: "node-control", VolumeSource: corev1.VolumeSource{HostPath: &corev1.HostPathVolumeSource{Path: directory, Type: &hostType}}}) + container.VolumeMounts = append(container.VolumeMounts, corev1.VolumeMount{Name: "node-control", MountPath: directory, ReadOnly: true}) + container.Env = append(container.Env, corev1.EnvVar{Name: "FELIS_NODE_CONTROL_SOCKET", Value: p.NodeControlSocket}) + } + deployment := controlPlaneDeployment(p, SAAPI, container, volumes) + if p.NodeControlNode != "" && p.ControllerNode == "" { + deployment.Spec.Template.Spec.NodeSelector = map[string]string{"kubernetes.io/hostname": p.NodeControlNode} + } + return deployment } // apiService exposes the built-in HTTPS panel/API origin as a stable NodePort. diff --git a/panel/dev/mockApi.ts b/panel/dev/mockApi.ts index 6d4834f..9aa13dd 100644 --- a/panel/dev/mockApi.ts +++ b/panel/dev/mockApi.ts @@ -1019,6 +1019,10 @@ async function handlePublic(ctx: RequestContext): Promise { async function handleSession(ctx: SessionContext): Promise { switch (route(ctx)) { + case "GET settings/node-control": + if (!isOwner(ctx.account.role)) sendError(ctx.res, 403, "forbidden", "Owner account required"); + else sendJSON(ctx.res, 200, { available: false, tasks: [] }); + return true; case "GET settings/wake-policy": case "PUT settings/wake-policy": { if (!isOwner(ctx.account.role)) { diff --git a/panel/e2e/startup-platform.smoke.spec.ts b/panel/e2e/startup-platform.smoke.spec.ts index 9ef5282..2171d8c 100644 --- a/panel/e2e/startup-platform.smoke.spec.ts +++ b/panel/e2e/startup-platform.smoke.spec.ts @@ -84,3 +84,32 @@ test("Owner emergency stop stays secondary, confirms name, and distinguishes acc const links = await ownerNav.evaluateAll((elements) => elements.map((el) => el.getAttribute("href"))); expect(links.at(-1)).toBe("/admin/platform"); }); + +test("Owner executes node management and receives stage, failure logs and retry", async ({ page, signIn }) => { + await signIn("owner"); + let task: Record = {}; + let submitted: Record | null = null; + await page.route("**/api/v1/settings/node-control", (route) => route.fulfill({ json: { available: true, tasks: task.id ? [task] : [] } })); + await page.route("**/api/v1/settings/node-control/tasks", async (route) => { + submitted = route.request().postDataJSON(); + task = { id: "test-task", request: submitted, actor: "owner", state: "running", stage: "database_backup", startedAt: new Date().toISOString(), log: "[felis] database_backup" }; + await route.fulfill({ status: 202, json: task }); + }); + await page.route("**/api/v1/settings/node-control/tasks/test-task", (route) => route.fulfill({ json: task })); + await page.route("**/api/v1/settings/node-control/tasks/test-task/retry", (route) => { task = { ...task, state: "running", error: undefined }; return route.fulfill({ status: 202, json: task }); }); + await page.goto("/admin/platform"); + await page.getByLabel(t("admin:node_control_ip")).fill("192.0.2.10"); + const submit = page.getByRole("button", { name: t("admin:node_control_submit"), exact: true }); + await expect(submit).toBeDisabled(); + await page.getByRole("switch").click(); + await submit.click(); + expect(submitted).toMatchObject({ action: "enable", externalIP: "192.0.2.10", confirmMaintenance: true }); + await expect(page.getByText(t("admin:node_control_stage_database_backup"), { exact: true })).toBeVisible(); + task = { ...task, state: "failed", error: "Database backup failed", log: "[felis] pg_dump failed: connection refused" }; + await page.getByRole("button", { name: t("admin:node_control_reload"), exact: true }).click(); + await expect(page.getByText("Database backup failed", { exact: true })).toBeVisible(); + await expect(page.getByLabel(t("admin:node_control_log"))).toContainText("pg_dump failed: connection refused"); + await page.getByRole("button", { name: t("admin:node_control_retry"), exact: true }).click(); + await expect(page.getByText(t("admin:node_control_state_running"), { exact: true })).toBeVisible(); + await expectFitsScreen(page); +}); diff --git a/panel/src/components/NodeControlPanel.test.tsx b/panel/src/components/NodeControlPanel.test.tsx new file mode 100644 index 0000000..db15102 --- /dev/null +++ b/panel/src/components/NodeControlPanel.test.tsx @@ -0,0 +1,63 @@ +// @vitest-environment jsdom +import { act, fireEvent, render, screen, waitFor } from "@testing-library/react"; +import { beforeEach, describe, expect, it, vi } from "vitest"; +import userEvent from "@testing-library/user-event"; +import { NodeControlPanel } from "./NodeControlPanel"; +import type { NodeControlTask } from "@/lib/types"; +const calls = vi.hoisted(() => ({ nodeTasks: vi.fn(), nodeTask: vi.fn(), startNodeTask: vi.fn(), retryNodeTask: vi.fn() })); +vi.mock("@/lib/api", async (original) => ({ ...await original(), api: calls })); +const failed: NodeControlTask = { id: "task-1", request: { action: "approve", name: "worker-01", sshTarget: "worker-01", confirmMaintenance: true }, actor: "owner", state: "failed", stage: "approval", startedAt: "2026-10-06T10:00:00Z", error: "Network isolation probe failed", log: "probe: controller port rejected" }; +beforeEach(() => { + vi.clearAllMocks(); calls.nodeTasks.mockResolvedValue({ available: true, tasks: [] }); calls.nodeTask.mockResolvedValue(failed); +}); +describe("host node management", () => { + it("keeps enrollment form visible and disabled while availability is unknown", async () => { + let resolve!: (value: { available: boolean; tasks: [] }) => void; + calls.nodeTasks.mockReturnValue(new Promise((r) => { resolve = r; })); + render(); + expect(screen.getByLabelText("Node name")).toHaveProperty("disabled", true); + expect(screen.getByRole("button", { name: "Execute operation" })).toHaveProperty("disabled", true); + await act(async () => resolve({ available: true, tasks: [] })); + expect(screen.getByLabelText("Node name")).toHaveProperty("disabled", false); + }); + it("updates the default action after the deployment mode loads", async () => { + const view = render(); + await waitFor(() => expect(calls.nodeTasks).toHaveBeenCalled()); + expect(screen.queryByLabelText("Node name")).toBeNull(); + view.rerender(); + await waitFor(() => expect(screen.getByLabelText("Node name")).toHaveProperty("disabled", false)); + expect(screen.getByRole("combobox", { name: "Operation" }).textContent).toContain("Enroll and approve worker"); + }); + it("requires maintenance acknowledgement and submits structured enrollment without credentials", async () => { + const running = { ...failed, state: "running", request: { ...failed.request, action: "join" } }; + calls.startNodeTask.mockResolvedValue(running); calls.nodeTask.mockResolvedValue(running); + render(); + await waitFor(() => expect(screen.getByLabelText("Node name")).toHaveProperty("disabled", false)); + fireEvent.change(screen.getByLabelText("Node name"), { target: { value: "worker-01" } }); + fireEvent.change(screen.getByLabelText("SSH target"), { target: { value: "root@192.0.2.10" } }); + fireEvent.change(screen.getByLabelText("Fixed node IP"), { target: { value: "192.0.2.10" } }); + expect(screen.getByRole("button", { name: "Execute operation" })).toHaveProperty("disabled", true); + await userEvent.click(screen.getByRole("switch")); + await userEvent.click(screen.getByRole("button", { name: "Execute operation" })); + expect(calls.startNodeTask).toHaveBeenCalledWith({ action: "join", name: "worker-01", sshTarget: "root@192.0.2.10", externalIP: "192.0.2.10", peers: [], confirmMaintenance: true }); + await screen.findByRole("region", { name: "Task progress" }); + expect(screen.getByRole("button", { name: "Execute operation" })).toHaveProperty("disabled", true); + }); + it("restores failed tasks with concrete errors and logs and retries as a new attempt", async () => { + calls.nodeTasks.mockResolvedValue({ available: true, tasks: [failed] }); + calls.retryNodeTask.mockResolvedValue({ ...failed, id: "task-2", state: "running" }); + render(); + expect(await screen.findByText("Network isolation probe failed")).toBeTruthy(); + expect(screen.getByLabelText("Execution log").textContent).toContain("probe: controller port rejected"); + await userEvent.click(screen.getByRole("button", { name: "Retry operation" })); + expect(calls.retryNodeTask).toHaveBeenCalledWith("task-1"); + await waitFor(() => expect(calls.nodeTask).toHaveBeenCalledWith("task-2")); + }); + it("reports a host connection failure and disables changes", async () => { + calls.nodeTasks.mockRejectedValue({ status: 503, code: "node_control_unavailable" }); + render(); + expect(await screen.findByRole("alert")).toBeTruthy(); + expect(screen.getByRole("button", { name: "Execute operation" })).toHaveProperty("disabled", true); + expect(calls.startNodeTask).not.toHaveBeenCalled(); + }); +}); diff --git a/panel/src/components/NodeControlPanel.tsx b/panel/src/components/NodeControlPanel.tsx new file mode 100644 index 0000000..12ab785 --- /dev/null +++ b/panel/src/components/NodeControlPanel.tsx @@ -0,0 +1,90 @@ +import { useEffect, useState } from "react"; +import { Loader2, Play, RefreshCw } from "lucide-react"; +import { useTranslation } from "react-i18next"; +import { Badge } from "@/components/ui/badge"; +import { Button } from "@/components/ui/button"; +import { Card, CardContent, CardHeader, CardTitle } from "@/components/ui/card"; +import { Input } from "@/components/ui/input"; +import { Label } from "@/components/ui/label"; +import { Switch } from "@/components/ui/switch"; +import { Select, SelectContent, SelectItem, SelectTrigger, SelectValue } from "@/components/ui/select"; +import { MessageLine } from "@/components/MessageLine"; +import { isReauthCancelled, useReauth } from "@/components/ReauthDialog"; +import { api, humanizeError } from "@/lib/api"; +import { useAsync, usePolling } from "@/lib/hooks"; +import type { NodeControlRequest, NodeControlTask } from "@/lib/types"; + +export function NodeControlPanel({ distributed }: { distributed: boolean | null }) { + const { t } = useTranslation("admin"); + const reauth = useReauth(); + const tasks = useAsync(api.nodeTasks, [], { coalesce: true }); + const [id, setId] = useState(null); + const [action, setAction] = useState(distributed ? "join" : "enable"); + const [name, setName] = useState(""); + const [ssh, setSSH] = useState(""); + const [ip, setIP] = useState(""); + const [peers, setPeers] = useState(""); + const [confirmed, setConfirmed] = useState(false); + const [submitting, setSubmitting] = useState(false); + const [error, setError] = useState(null); + const progress = useAsync(() => id ? api.nodeTask(id) : Promise.resolve(null), [id], { coalesce: true }); + const task = progress.data; + const running = tasks.data?.tasks.some((entry) => entry.state === "running") || task?.state === "running"; + const busy = submitting || !!running; + const available = tasks.data?.available === true && !tasks.error && distributed !== null; + useEffect(() => { if (distributed !== null) { setAction(distributed ? "join" : "enable"); setConfirmed(false); } }, [distributed]); + useEffect(() => { if (!id && tasks.data?.tasks[0]) setId(tasks.data.tasks[0].id); }, [id, tasks.data]); + usePolling(() => { tasks.reload(); if (id) progress.reload(); }, running || progress.error || tasks.error ? 3000 : 10000); + async function run(retry = false) { + if (!available || busy || (!retry && !confirmed)) return; + setSubmitting(true); setError(null); + try { + const result: NodeControlTask = await reauth.guard(() => retry && id ? api.retryNodeTask(id) : api.startNodeTask({ + action, ...(action !== "enable" ? { name: name.trim(), sshTarget: ssh.trim() } : {}), + ...(action !== "approve" ? { externalIP: ip.trim(), peers: peers.split(/[\s,]+/).filter(Boolean) } : {}), + confirmMaintenance: confirmed, + })); + setId(result.id); setConfirmed(false); tasks.reload(); + } catch (e) { if (!isReauthCancelled(e)) setError(humanizeError(e)); } + finally { setSubmitting(false); } + } + const valid = confirmed && (action === "enable" || !!name.trim() && !!ssh.trim()) && (action === "approve" || !!ip.trim()); + return + {t("node_control_title")} + +

{t("node_control_description")}

+ {tasks.loading && !tasks.data &&

{t("node_control_loading")}

} + {tasks.error != null && } + {tasks.data?.available === false &&

{t("node_control_unavailable")}

} + {running && task?.state !== "running" &&

{t("node_control_other_running")}

} +
+
+ {action !== "enable" && <> +
setName(e.target.value)} placeholder="worker-01" />
+
setSSH(e.target.value)} placeholder="root@192.168.1.20" />

{t("node_control_ssh_hint")}

+ } + {action !== "approve" && <> +
setIP(e.target.value)} placeholder="192.168.1.20" />
+
setPeers(e.target.value)} placeholder="192.168.1.21/32, 192.168.1.22/32" />

{t("node_control_peers_hint")}

+ } +
+
+ {error && } + +
{t("node_control_prerequisites")}

{t("node_control_prerequisites_body")}

+ {!!tasks.data?.tasks.length &&
} + {progress.error != null && } + {id && progress.loading && !task &&

{t("node_control_loading")}

} + {task &&
+
{t(`node_control_state_${task.state === "running" && (progress.error || tasks.error) ? "unknown" : task.state}`)}{t(`node_control_stage_${task.stage}`, { defaultValue: task.stage })}
+

{t("node_control_task_id")}: {task.id}

+

{t("node_control_started", { time: new Date(task.startedAt).toLocaleString() })} · {t("node_control_timeout")}

+ {task.error && } +
{task.log || t("node_control_log_pending")}
+ {task.state === "failed" && <>

{t("node_control_retry_hint")}

} + {task.state === "succeeded" && task.request.action === "enable" && } +
} +
+ {reauth.dialog} +
; +} diff --git a/panel/src/i18n/resources/en-US/admin.json b/panel/src/i18n/resources/en-US/admin.json index 8bec0ea..f6022ef 100644 --- a/panel/src/i18n/resources/en-US/admin.json +++ b/panel/src/i18n/resources/en-US/admin.json @@ -335,5 +335,57 @@ "platform_distribution_approve": "Approve the node from the controller. The command verifies identity, network isolation, image pulling, and cross-node connectivity. Failed checks preserve isolation. Approval removes the specified cached image and must run before the node hosts game workloads.", "platform_distribution_runbook": "Deployment and acceptance runbook", "platform_distribution_management": "Online approved workers can be selected when creating a server. Existing servers must be fully stopped before migration. Successful migration leaves the server stopped and retains its source world volume.", - "platform_distribution_servers": "Manage server placement and migration" + "platform_distribution_servers": "Manage server placement and migration", + "node_control_title": "Node operations", + "node_control_description": "Enable distributed deployment, enroll workers and run admission checks. Tasks execute on the controller host; panel disconnection does not terminate a task.", + "node_control_loading": "Loading host task status.", + "node_control_unavailable": "Host node service is not installed. Update the controller installation to enable panel management.", + "node_control_action": "Operation", + "node_control_enable": "Enable distributed deployment", + "node_control_join": "Enroll and approve worker", + "node_control_approve": "Approve existing worker", + "node_control_name": "Node name", + "node_control_ssh": "SSH target", + "node_control_ssh_hint": "Use an alias from the controller root SSH configuration or user@host. Noninteractive authentication and verified host fingerprints are required.", + "node_control_ip": "Fixed node IP", + "node_control_peers": "Additional peer addresses (optional)", + "node_control_peers_hint": "Enter exact /32 or /128 addresses separated by commas. Existing nodes and the new node are included automatically.", + "node_control_confirm_enable": "All game servers (including Login and Lobby) and maintenance tasks are stopped. Database and cluster backups and controller network changes are authorized. Panel access may be temporarily interrupted.", + "node_control_confirm_worker": "Game and maintenance workloads on the target node are stopped. Node installation, peer firewall updates, image pulls and network isolation checks are authorized.", + "node_control_submit": "Execute operation", + "node_control_prerequisites": "Environment and SSH requirements", + "node_control_prerequisites_body": "Controller and workers require matching architecture and k3s versions and fixed addresses. SSH accounts require root or sudo -n access. Configure keys and verify host fingerprints on the controller host before first use; unknown fingerprints are rejected. Existing worker SSH aliases must match node names for peer firewall updates. Enrollment requires a host without k3s. Nodes already joined after an interrupted enrollment must use “Approve existing worker”. Failed admission retains quarantine and does not start game servers.", + "node_control_history": "Task history", + "node_control_controller": "Controller", + "node_control_state_running": "Running", + "node_control_state_succeeded": "Completed", + "node_control_state_failed": "Failed", + "node_control_connection_lost": "Task status is temporarily unavailable: {{error}}. Polling resumes when connectivity is restored; the last confirmed status remains visible.", + "node_control_progress": "Task progress", + "node_control_task_id": "Task ID", + "node_control_log": "Execution log", + "node_control_log_pending": "No execution output yet.", + "node_control_retry_hint": "Verify current host state and resolve the failure before retrying. Retry creates a new task and repeats checks. Already enrolled workers require approval instead of enrollment.", + "node_control_retry": "Retry operation", + "node_control_refresh_deployment": "Refresh deployment status", + "node_control_stage_preflight": "Checking prerequisites", + "node_control_stage_pause_operator": "Pausing operator reconciliation", + "node_control_stage_database_backup": "Backing up database", + "node_control_stage_cluster_backup": "Backing up cluster state", + "node_control_stage_configure_controller": "Configuring controller identity", + "node_control_stage_install_controller": "Deploying controller components and network", + "node_control_stage_ssh_check": "Checking SSH and host environment", + "node_control_stage_firewall": "Updating peer firewalls", + "node_control_stage_bootstrap_token": "Creating temporary enrollment token", + "node_control_stage_transfer": "Transferring installation files", + "node_control_stage_install_worker": "Installing worker components", + "node_control_stage_wait_ready": "Waiting for node readiness", + "node_control_stage_approval": "Running node admission checks", + "node_control_stage_complete": "Operation completed", + "node_control_reload": "Refresh node tasks", + "node_control_state_unknown": "State unconfirmed", + "node_control_started": "Started: {{time}}", + "node_control_timeout": "Operation timeout: 45 minutes", + "node_control_other_running": "Another node task is running; new node operations are temporarily blocked.", + "node_control_show_running": "View running task" } diff --git a/panel/src/i18n/resources/en-US/errors.json b/panel/src/i18n/resources/en-US/errors.json index b3b553c..6cdca21 100644 --- a/panel/src/i18n/resources/en-US/errors.json +++ b/panel/src/i18n/resources/en-US/errors.json @@ -145,5 +145,7 @@ "auth_sources_unavailable": "Authentication source management is unavailable.", "auth_sources_changed": "Authentication sources changed. Discard your edits and reload before saving.", "auth_source_tag_locked": "Saved source IDs cannot be renamed or removed. Disable the source instead.", - "distributed_unavailable": "Distributed deployment is not configured. Configure the controller and worker nodes using the deployment runbook first." + "distributed_unavailable": "Distributed deployment is not configured. Configure the controller and worker nodes using the deployment runbook first.", + "node_operation_busy": "Node maintenance is running. Server starts and additional node tasks are temporarily blocked.", + "node_control_unavailable": "Host node service is unavailable. Inspect felis-node-control.service on the controller." } diff --git a/panel/src/i18n/resources/zh-CN/admin.json b/panel/src/i18n/resources/zh-CN/admin.json index 65ffe01..76efcc3 100644 --- a/panel/src/i18n/resources/zh-CN/admin.json +++ b/panel/src/i18n/resources/zh-CN/admin.json @@ -331,5 +331,57 @@ "platform_distribution_approve": "在主控批准节点。该命令验证节点身份、网络隔离、镜像拉取和跨节点连通性;任何检查失败均保留隔离状态。批准检查会删除指定缓存镜像,须在该节点尚无游戏任务时执行。", "platform_distribution_runbook": "查看完整部署与验收流程", "platform_distribution_management": "在线且已批准的 worker 可在创建服务器时选择。已有服务器须完全停止后才能迁移;迁移完成后仍保持停止,源世界卷保留。", - "platform_distribution_servers": "管理服务器节点与迁移" + "platform_distribution_servers": "管理服务器节点与迁移", + "node_control_title": "节点操作", + "node_control_description": "启用分布式部署、接入 worker 和执行节点批准。任务在主控宿主机执行;面板连接中断不代表任务终止。", + "node_control_loading": "正在读取主机任务状态。", + "node_control_unavailable": "主机节点执行服务尚未安装。更新主控安装后可启用面板管理。", + "node_control_action": "操作类型", + "node_control_enable": "启用分布式部署", + "node_control_join": "接入并批准 worker", + "node_control_approve": "批准已有 worker", + "node_control_name": "节点名称", + "node_control_ssh": "SSH 目标", + "node_control_ssh_hint": "使用主控 root SSH 配置中的别名或 user@host。要求免交互认证及已核验的主机指纹。", + "node_control_ip": "节点固定 IP", + "node_control_peers": "额外对等节点地址(选填)", + "node_control_peers_hint": "填写精确的 /32 或 /128 地址,以逗号分隔。系统自动包含现有节点及本次接入节点的地址。", + "node_control_confirm_enable": "已停止全部游戏服务器(含登录空间与正式大厅)及维护任务,确认允许备份数据库与集群状态,并变更主控网络配置。期间面板可能暂时不可访问。", + "node_control_confirm_worker": "已停止目标节点上的游戏与维护任务。确认允许安装节点组件、更新对等防火墙并执行镜像拉取及网络隔离验收。", + "node_control_submit": "执行操作", + "node_control_prerequisites": "环境与 SSH 要求", + "node_control_prerequisites_body": "主控与 worker 需使用相同架构及 k3s 版本,节点使用固定地址。SSH 账户须为 root 或具备 sudo -n 权限。首次使用前,应在主控宿主机配置密钥并核验各主机指纹;系统不会自动接受未知指纹。已有 worker 的 SSH 别名应与节点名称一致,以便更新对等防火墙。接入仅用于尚未安装 k3s 的主机;中断后已加入集群的节点应使用“批准已有 worker”。批准失败时保留隔离状态,不自动启动游戏服务器。", + "node_control_history": "任务记录", + "node_control_controller": "主控", + "node_control_state_running": "执行中", + "node_control_state_succeeded": "已完成", + "node_control_state_failed": "执行失败", + "node_control_connection_lost": "任务状态暂时无法读取:{{error}}。连接恢复后自动继续查询;当前展示最后一次已确认的状态。", + "node_control_progress": "任务进度", + "node_control_task_id": "任务 ID", + "node_control_log": "执行日志", + "node_control_log_pending": "尚未产生执行日志。", + "node_control_retry_hint": "重试前需核对主机当前状态并修复失败原因。重试创建新任务并重新执行检查;已接入的 worker 应改用批准操作。", + "node_control_retry": "重试操作", + "node_control_refresh_deployment": "刷新部署状态", + "node_control_stage_preflight": "检查环境", + "node_control_stage_pause_operator": "暂停控制器调谐", + "node_control_stage_database_backup": "备份数据库", + "node_control_stage_cluster_backup": "备份集群状态", + "node_control_stage_configure_controller": "配置主控身份", + "node_control_stage_install_controller": "部署主控组件与网络", + "node_control_stage_ssh_check": "检查 SSH 与主机环境", + "node_control_stage_firewall": "更新对等防火墙", + "node_control_stage_bootstrap_token": "创建限时接入令牌", + "node_control_stage_transfer": "传输节点安装文件", + "node_control_stage_install_worker": "安装 worker 组件", + "node_control_stage_wait_ready": "等待节点就绪", + "node_control_stage_approval": "执行节点验收", + "node_control_stage_complete": "操作完成", + "node_control_reload": "刷新节点任务", + "node_control_state_unknown": "状态未确认", + "node_control_started": "开始时间:{{time}}", + "node_control_timeout": "操作超时上限为 45 分钟", + "node_control_other_running": "另一项节点任务正在执行,当前暂不允许新的节点操作。", + "node_control_show_running": "查看执行中任务" } diff --git a/panel/src/i18n/resources/zh-CN/errors.json b/panel/src/i18n/resources/zh-CN/errors.json index c113042..14b2280 100644 --- a/panel/src/i18n/resources/zh-CN/errors.json +++ b/panel/src/i18n/resources/zh-CN/errors.json @@ -145,5 +145,7 @@ "auth_sources_unavailable": "认证源管理暂不可用。", "auth_sources_changed": "认证源已被其他操作修改。请放弃当前更改并重新读取后再保存。", "auth_source_tag_locked": "已保存的认证源标识不能改名或移除,请使用停用。", - "distributed_unavailable": "分布式部署未配置。请先按部署流程配置主控与 worker 节点。" + "distributed_unavailable": "分布式部署未配置。请先按部署流程配置主控与 worker 节点。", + "node_operation_busy": "节点维护正在执行,暂不允许启动服务器或创建另一项节点任务。", + "node_control_unavailable": "主机节点执行服务不可用。请检查主控上的 felis-node-control.service。" } diff --git a/panel/src/lib/api.ts b/panel/src/lib/api.ts index 2cd26b5..1fd3a29 100644 --- a/panel/src/lib/api.ts +++ b/panel/src/lib/api.ts @@ -1,4 +1,6 @@ import type { + NodeControlRequest, + NodeControlTask, AccessResult, ExecutionNode, WorldMigration, @@ -610,6 +612,10 @@ export const api = rejectingSync({ URL.revokeObjectURL(url); }, + nodeTasks: () => request<{ available: boolean; tasks: NodeControlTask[] }>("GET", "/settings/node-control"), + nodeTask: (id: string) => request("GET", `/settings/node-control/tasks/${encodeURIComponent(id)}`), + startNodeTask: (body: NodeControlRequest) => request("POST", "/settings/node-control/tasks", body), + retryNodeTask: (id: string) => request("POST", `/settings/node-control/tasks/${encodeURIComponent(id)}/retry`), nodes: () => request<{ nodes: ExecutionNode[] }>("GET", "/nodes").then((r) => r.nodes), migration: (name: string) => request("GET", `/servers/${encodeURIComponent(name)}/migrations`), migrateServer: (name: string, targetNode: string) => @@ -1253,6 +1259,10 @@ export function humanizeError(e: unknown): string { switch (err.code) { // Session doors (spec §B): every passwordless door 403s this when local // sessions are disabled on a Zero-Trust-only deployment. + case "node_operation_busy": + return t("node_operation_busy"); + case "node_control_unavailable": + return t("node_control_unavailable"); case "distributed_unavailable": return t("distributed_unavailable"); case "local_auth_disabled": diff --git a/panel/src/lib/openapi.gen.ts b/panel/src/lib/openapi.gen.ts index af76cc7..b59f8b6 100644 --- a/panel/src/lib/openapi.gen.ts +++ b/panel/src/lib/openapi.gen.ts @@ -1263,6 +1263,77 @@ export interface paths { patch?: never; trace?: never; }; + "/api/v1/settings/node-control": { + parameters: { + query?: never; + header?: never; + path?: never; + cookie?: never; + }; + /** Read host node-service availability and persistent task history (Owner). */ + get: operations["nodeTasks"]; + put?: never; + post?: never; + delete?: never; + options?: never; + head?: never; + patch?: never; + trace?: never; + }; + "/api/v1/settings/node-control/tasks": { + parameters: { + query?: never; + header?: never; + path?: never; + cookie?: never; + }; + get?: never; + put?: never; + /** + * Execute fixed host node operation (Owner, fresh reauthentication). + * @description Tasks are serialized, audited and persisted on the controller. The host verifies stopped workloads and SSH trust before enrollment or admission. Bootstrap credentials never cross the external API. Server starts fail closed during node operations and when the configured host service cannot report its state. + */ + post: operations["startNodeTask"]; + delete?: never; + options?: never; + head?: never; + patch?: never; + trace?: never; + }; + "/api/v1/settings/node-control/tasks/{id}": { + parameters: { + query?: never; + header?: never; + path?: never; + cookie?: never; + }; + /** Read task stage, terminal error and bounded execution log (Owner). */ + get: operations["nodeTask"]; + put?: never; + post?: never; + delete?: never; + options?: never; + head?: never; + patch?: never; + trace?: never; + }; + "/api/v1/settings/node-control/tasks/{id}/retry": { + parameters: { + query?: never; + header?: never; + path?: never; + cookie?: never; + }; + get?: never; + put?: never; + /** Create a new attempt from a failed task (Owner, fresh reauthentication). */ + post: operations["retryNodeTask"]; + delete?: never; + options?: never; + head?: never; + patch?: never; + trace?: never; + }; "/api/v1/settings/wake-policy": { parameters: { query?: never; @@ -3128,6 +3199,35 @@ export interface components { attempt: number; error?: string; }; + NodeControlRequest: { + /** @enum {string} */ + action: "enable" | "join" | "approve"; + /** @description Stable worker node name; required for join and approve. */ + name?: string; + /** @description Verified root SSH alias or user@host; required for workers. */ + sshTarget?: string; + /** @description Fixed node IP; required for enable and join. */ + externalIP?: string; + /** @description Additional exact /32 or /128 peer addresses. */ + peers?: string[]; + /** @constant */ + confirmMaintenance: true; + }; + NodeControlTask: { + id: string; + request: components["schemas"]["NodeControlRequest"]; + actor: string; + /** @enum {string} */ + state: "running" | "succeeded" | "failed"; + stage: string; + /** Format: date-time */ + startedAt: string; + /** Format: date-time */ + finishedAt?: string; + error?: string; + /** @description Most recent 64 KiB of host execution output; absent in task lists. */ + log?: string; + }; /** @description Status projection of one server (internal/api/cluster.go ServerInfo). */ ServerInfo: { /** @description Private runtime and startup diagnostics for authorized server managers. */ @@ -6825,6 +6925,167 @@ export interface operations { 401: components["responses"]["Unauthorized"]; }; }; + nodeTasks: { + parameters: { + query?: never; + header?: never; + path?: never; + cookie?: never; + }; + requestBody?: never; + responses: { + /** @description Availability and most recent task records; no credentials or logs. */ + 200: { + headers: { + [name: string]: unknown; + }; + content: { + "application/json": { + available: boolean; + tasks: components["schemas"]["NodeControlTask"][]; + }; + }; + }; + 400: components["responses"]["BadRequest"]; + 401: components["responses"]["Unauthorized"]; + 403: components["responses"]["Forbidden"]; + /** @description Another node task is running or this task cannot be retried. */ + 409: { + headers: { + [name: string]: unknown; + }; + content?: never; + }; + /** @description Host node execution service is unavailable. */ + 503: { + headers: { + [name: string]: unknown; + }; + content?: never; + }; + }; + }; + startNodeTask: { + parameters: { + query?: never; + header?: never; + path?: never; + cookie?: never; + }; + requestBody: { + content: { + "application/json": components["schemas"]["NodeControlRequest"]; + }; + }; + responses: { + /** @description Persistent task created. */ + 202: { + headers: { + [name: string]: unknown; + }; + content: { + "application/json": components["schemas"]["NodeControlTask"]; + }; + }; + 400: components["responses"]["BadRequest"]; + 401: components["responses"]["Unauthorized"]; + 403: components["responses"]["Forbidden"]; + /** @description Another node task is running or this task cannot be retried. */ + 409: { + headers: { + [name: string]: unknown; + }; + content?: never; + }; + /** @description Host node execution service is unavailable. */ + 503: { + headers: { + [name: string]: unknown; + }; + content?: never; + }; + }; + }; + nodeTask: { + parameters: { + query?: never; + header?: never; + path: { + id: string; + }; + cookie?: never; + }; + requestBody?: never; + responses: { + /** @description Current task and recent output. */ + 200: { + headers: { + [name: string]: unknown; + }; + content: { + "application/json": components["schemas"]["NodeControlTask"]; + }; + }; + 400: components["responses"]["BadRequest"]; + 401: components["responses"]["Unauthorized"]; + 403: components["responses"]["Forbidden"]; + 404: components["responses"]["NotFound"]; + /** @description Another node task is running or this task cannot be retried. */ + 409: { + headers: { + [name: string]: unknown; + }; + content?: never; + }; + /** @description Host node execution service is unavailable. */ + 503: { + headers: { + [name: string]: unknown; + }; + content?: never; + }; + }; + }; + retryNodeTask: { + parameters: { + query?: never; + header?: never; + path: { + id: string; + }; + cookie?: never; + }; + requestBody?: never; + responses: { + /** @description New task created; original task is retained. */ + 202: { + headers: { + [name: string]: unknown; + }; + content: { + "application/json": components["schemas"]["NodeControlTask"]; + }; + }; + 400: components["responses"]["BadRequest"]; + 401: components["responses"]["Unauthorized"]; + 403: components["responses"]["Forbidden"]; + 404: components["responses"]["NotFound"]; + /** @description Another node task is running or this task cannot be retried. */ + 409: { + headers: { + [name: string]: unknown; + }; + content?: never; + }; + /** @description Host node execution service is unavailable. */ + 503: { + headers: { + [name: string]: unknown; + }; + content?: never; + }; + }; + }; getWakePolicy: { parameters: { query?: never; diff --git a/panel/src/lib/types.ts b/panel/src/lib/types.ts index 999555e..ef828c2 100644 --- a/panel/src/lib/types.ts +++ b/panel/src/lib/types.ts @@ -776,3 +776,23 @@ export interface WakePolicySettings { revision: string; managed: boolean; } + +export interface NodeControlRequest { + action: "enable" | "join" | "approve"; + name?: string; + sshTarget?: string; + externalIP?: string; + peers?: string[]; + confirmMaintenance: boolean; +} +export interface NodeControlTask { + id: string; + request: NodeControlRequest; + actor: string; + state: "running" | "succeeded" | "failed"; + stage: string; + startedAt: string; + finishedAt?: string; + error?: string; + log?: string; +} diff --git a/panel/src/pages/admin/PlatformSettingsPage.test.tsx b/panel/src/pages/admin/PlatformSettingsPage.test.tsx index 6054aba..1580ce5 100644 --- a/panel/src/pages/admin/PlatformSettingsPage.test.tsx +++ b/panel/src/pages/admin/PlatformSettingsPage.test.tsx @@ -5,7 +5,7 @@ import userEvent from "@testing-library/user-event"; import { MemoryRouter } from "react-router-dom"; import { PlatformSettingsPage } from "./PlatformSettingsPage"; import type { WakePolicySettings } from "@/lib/types"; -const calls = vi.hoisted(() => ({ getWakePolicy: vi.fn(), setWakePolicy: vi.fn(), nodes: vi.fn() })); +const calls = vi.hoisted(() => ({ getWakePolicy: vi.fn(), setWakePolicy: vi.fn(), nodes: vi.fn(), nodeTasks: vi.fn(), nodeTask: vi.fn(), startNodeTask: vi.fn(), retryNodeTask: vi.fn() })); vi.mock("@/lib/api", async (original) => ({ ...await original(), api: calls })); const runtime = vi.hoisted(() => ({ distributed: false, fallback: false, apiBase: "/api/v1", rootDomain: "example.test" })); vi.mock("@/lib/hooks", async (original) => ({ ...await original(), useConfig: () => runtime })); @@ -18,6 +18,7 @@ beforeEach(() => { runtime.distributed = false; runtime.fallback = false; calls.nodes.mockResolvedValue([]); + calls.nodeTasks.mockResolvedValue({ available: false, tasks: [] }); calls.getWakePolicy.mockResolvedValue(policy); calls.setWakePolicy.mockImplementation(async (body) => ({ ...body, revision: "saved", managed: true })); }); @@ -28,7 +29,7 @@ describe("platform policy", () => { page(); expect(limit().disabled).toBe(true); expect(limit().value).toBe(""); - expect(screen.getByRole("status").textContent).toContain("Loading the saved policy"); + expect(screen.getAllByRole("status").some((node) => node.textContent?.includes("Loading the saved policy"))).toBe(true); await act(async () => resolve(policy)); expect(limit().disabled).toBe(false); expect(limit().value).toBe("0"); @@ -57,9 +58,8 @@ describe("distributed deployment entry", () => { page(); await waitFor(() => expect(limit().disabled).toBe(false)); expect(screen.getByText("Single-node mode")).toBeTruthy(); - expect(screen.getByText(/cannot be switched live from the panel/)).toBeTruthy(); - expect(screen.getByText("Connect and approve worker nodes")).toBeTruthy(); - expect(screen.getByRole("link", { name: "Deployment and acceptance runbook" }).getAttribute("href")).toContain("docs/distributed.md"); + expect(screen.getByText("Node operations")).toBeTruthy(); + expect(await screen.findByText(/Host node service is not installed/)).toBeTruthy(); expect(calls.nodes).not.toHaveBeenCalled(); }); diff --git a/panel/src/pages/admin/PlatformSettingsPage.tsx b/panel/src/pages/admin/PlatformSettingsPage.tsx index da3f626..6965e1b 100644 --- a/panel/src/pages/admin/PlatformSettingsPage.tsx +++ b/panel/src/pages/admin/PlatformSettingsPage.tsx @@ -2,6 +2,7 @@ import { Link } from "react-router-dom"; import { useEffect, useState } from "react"; import { Loader2, RefreshCw, Save, Settings } from "lucide-react"; import { useTranslation } from "react-i18next"; +import { NodeControlPanel } from "@/components/NodeControlPanel"; import { DistributedNodes } from "@/components/DistributedNodes"; import { Badge } from "@/components/ui/badge"; import { PageHeader } from "@/components/PageHeader"; @@ -75,21 +76,11 @@ export function PlatformSettingsPage() { {t(!config || config.fallback ? "platform_distribution_unknown" : config.distributed ? "platform_distribution_enabled" : "platform_distribution_disabled")}

{t("platform_distribution_architecture")}

- {config && !config.fallback && !config.distributed &&

{t("platform_distribution_enable_hint")}

} -
- {t("platform_distribution_join")} -
    -
  1. {t("platform_distribution_prepare")}
  2. -
  3. {t("platform_distribution_token")}
    {"felis node token --name  --ttl 10m --out /root/worker.bootstrap"}
  4. -
  5. {t("platform_distribution_worker")}
  6. -
  7. {t("platform_distribution_approve")}
    {"felis node approve --name  --ssh-target  --image \nfelis node list"}
  8. -
- -

{t("platform_distribution_management")}

+ {config?.distributed && !config.fallback && } {reauth.dialog} ;