feat: implement cluster joining for server and worker nodes

Different options get applied to each depending on what is needed for the role.

This also updates how we initialize the cluster to better support how metrics exposure and cross-server networking.
This commit is contained in:
Jose Diaz-Gonzalez
2024-01-20 04:58:02 -05:00
parent af3009f0da
commit 69e9511666
7 changed files with 593 additions and 29 deletions

View File

@@ -145,10 +145,31 @@ In some cases, Dokku must be setup from scratch and joined to an existing cluste
dokku scheduler-k3s:set --global token topsecret:server:token
```
Next, run the `scheduler-k3s:join-cluster` command. This command takes a single server node in the cluster, will SSH onto the specified server, retrieve the kubeconfig, and then configure kubectl on the Dokku node to speak with the specified k3s cluster. Dokku will also periodically retrieve all server nodes in the cluster so that server nodes can be safely replaced as needed.
Next, run the `scheduler-k3s:cluster-add` command. This command takes a single server node in the cluster, will SSH onto the specified server, retrieve the kubeconfig, and then configure kubectl on the Dokku node to speak with the specified k3s cluster.
> [!TODO]
> TODO: Dokku will also periodically retrieve all server nodes in the cluster so that server nodes can be safely replaced as needed.
```shell
dokku scheduler-k3s:join-cluster root@server-1.example.com
dokku scheduler-k3s:cluster-add ssh://root@worker-1.example.com
```
If the server isn't in the `known_hosts` file, the connection will fail. This can be bypassed by setting the `--insecure-allow-unknown-hosts` flag:
```shell
dokku scheduler-k3s:cluster-add --insecure-allow-unknown-hosts ssh://root@worker-1.example.com
```
New server nodes are added as worker nodes which run application workloads but do not otherwise participate in managing the cluster. For high availability, new server nodes can be added to the cluster by specifying the `--role` flag:
```shell
dokku scheduler-kes:cluster-add --role server ssh://root@cluster-1.example.com
```
Server nodes allow any workloads to be scheduled on them by default, in addition to the control-plane, etcd, and the scheduler itself. To avoid app workloads being scheduled on your control-plane, use the `--taint-scheduling` flag:
```shell
dokku scheduler-kes:cluster-add --role server --taint-scheduling ssh://root@cluster-1.example.com
```
### Using kubectl remotely

View File

@@ -1,4 +1,4 @@
SUBCOMMANDS = subcommands/initialize subcommands/report subcommands/set subcommands/show-kubeconfig
SUBCOMMANDS = subcommands/initialize subcommands/join-cluster subcommands/report subcommands/set subcommands/show-kubeconfig subcommands/uninstall
TRIGGERS = triggers/install triggers/post-app-clone-setup triggers/post-app-rename-setup triggers/post-delete triggers/post-registry-login triggers/report triggers/scheduler-deploy triggers/scheduler-enter triggers/scheduler-logs triggers/scheduler-post-delete triggers/scheduler-run triggers/scheduler-run-list triggers/scheduler-stop
BUILD = commands subcommands triggers
PLUGIN_NAME = scheduler-k3s

View File

@@ -49,6 +49,13 @@ type StartCommandOutput struct {
Command []string
}
type WaitForNodeToExistInput struct {
Clientset KubernetesClient
Namespace string
RetryCount int
NodeName string
}
type WaitForPodBySelectorRunningInput struct {
Clientset KubernetesClient
Namespace string
@@ -188,25 +195,25 @@ func extractStartCommand(input StartCommandInput) string {
return "/start " + input.ProcessType
}
startCommandResp, err := common.CallPlugnTrigger(common.PlugnTriggerInput{
resp, err := common.CallPlugnTrigger(common.PlugnTriggerInput{
Trigger: "config-get",
Args: []string{input.AppName, "DOKKU_START_CMD"},
CaptureOutput: true,
StreamStdio: false,
})
if err == nil && startCommandResp.ExitCode == 0 && len(startCommandResp.Stdout) > 0 {
command = startCommandResp.Stdout
if err == nil && resp.ExitCode == 0 && len(resp.Stdout) > 0 {
command = strings.TrimSpace(resp.Stdout)
}
if input.ImageSourceType == "dockerfile" {
startCommandDockerfileResp, err := common.CallPlugnTrigger(common.PlugnTriggerInput{
resp, err := common.CallPlugnTrigger(common.PlugnTriggerInput{
Trigger: "config-get",
Args: []string{input.AppName, "DOKKU_DOCKERFILE_START_CMD"},
CaptureOutput: true,
StreamStdio: false,
})
if err == nil && startCommandDockerfileResp.ExitCode == 0 && len(startCommandDockerfileResp.Stdout) > 0 {
command = startCommandDockerfileResp.Stdout
if err == nil && resp.ExitCode == 0 && len(resp.Stdout) > 0 {
command = strings.TrimSpace(resp.Stdout)
}
}
@@ -278,6 +285,22 @@ func getStartCommand(input StartCommandInput) (StartCommandOutput, error) {
}, nil
}
func isK3sInstalled() error {
if !common.FileExists("/usr/local/bin/k3s") {
return fmt.Errorf("k3s binary is not available")
}
if !common.FileExists(RegistryConfigPath) {
return fmt.Errorf("k3s registry config is not available")
}
if !common.FileExists(KubeConfigPath) {
return fmt.Errorf("k3s kubeconfig is not available")
}
return nil
}
func isPodReady(ctx context.Context, clientset KubernetesClient, podName, namespace string) wait.ConditionWithContextFunc {
return func(ctx context.Context) (bool, error) {
fmt.Printf(".") // progress bar!
@@ -331,6 +354,37 @@ func waitForPodBySelectorRunning(ctx context.Context, input WaitForPodBySelector
return nil
}
func waitForNodeToExist(ctx context.Context, input WaitForNodeToExistInput) ([]v1.Node, error) {
var matchingNodes []v1.Node
var err error
for i := 0; i < input.RetryCount; i++ {
nodes, err := input.Clientset.ListNodes(ctx, ListNodesInput{})
if err != nil {
time.Sleep(1 * time.Second)
}
if input.NodeName == "" {
matchingNodes = nodes
break
}
for _, node := range nodes {
if node.Name == input.NodeName {
matchingNodes = append(matchingNodes, node)
break
}
}
if len(matchingNodes) > 0 {
break
}
time.Sleep(1 * time.Second)
}
if err != nil {
return matchingNodes, fmt.Errorf("Error listing nodes: %w", err)
}
return matchingNodes, nil
}
func waitForPodToExist(ctx context.Context, input WaitForPodToExistInput) ([]v1.Pod, error) {
var pods []v1.Pod
var err error
@@ -352,6 +406,7 @@ func waitForPodToExist(ctx context.Context, input WaitForPodToExistInput) ([]v1.
break
}
}
time.Sleep(1 * time.Second)
}
if err != nil {
return pods, fmt.Errorf("Error listing pods: %w", err)

View File

@@ -3,7 +3,10 @@ package scheduler_k3s
import (
"context"
"errors"
"fmt"
"github.com/dokku/dokku/plugins/common"
"github.com/go-openapi/jsonpointer"
appsv1 "k8s.io/api/apps/v1"
autoscalingv1 "k8s.io/api/autoscaling/v1"
batchv1 "k8s.io/api/batch/v1"
@@ -11,6 +14,7 @@ import (
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/runtime/schema"
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
"k8s.io/client-go/tools/clientcmd"
@@ -165,6 +169,37 @@ func (k KubernetesClient) GetPod(ctx context.Context, input GetPodInput) (v1.Pod
return *pod, err
}
// LabelNodeInput contains all the information needed to label a Kubernetes node
type LabelNodeInput struct {
// Name is the Kubernetes node name
Name string
// Key is the label key
Key string
// Value is the label value
Value string
}
// LabelNode labels a Kubernetes node
func (k KubernetesClient) LabelNode(ctx context.Context, input LabelNodeInput) error {
node, err := k.Client.CoreV1().Nodes().Get(ctx, input.Name, metav1.GetOptions{})
if err != nil {
return err
}
if node == nil {
return errors.New("node is nil")
}
keyPath := fmt.Sprintf("/metadata/labels/%s", jsonpointer.Escape("kubernetes.io/role"))
patch := fmt.Sprintf(`[{"op":"add", "path":"%s", "value":"%s" }]`, keyPath, "worker")
_, err = k.Client.CoreV1().Nodes().Patch(context.Background(), node.Name, types.JSONPatchType, []byte(patch), metav1.PatchOptions{})
if err != nil {
return fmt.Errorf("failed to label node: %w", err)
}
return nil
}
// ListCronJobsInput contains all the information needed to list Kubernetes cron jobs
type ListCronJobsInput struct {
// LabelSelector is the Kubernetes label selector
@@ -230,6 +265,31 @@ func (k KubernetesClient) ListNamespaces(ctx context.Context) ([]v1.Namespace, e
return namespaces.Items, nil
}
// ListNodesInput contains all the information needed to list Kubernetes nodes
type ListNodesInput struct {
// LabelSelector is the Kubernetes label selector
LabelSelector string
}
// ListNodes lists Kubernetes nodes
func (k KubernetesClient) ListNodes(ctx context.Context, input ListNodesInput) ([]v1.Node, error) {
listOptions := metav1.ListOptions{}
if input.LabelSelector != "" {
common.LogDebug(fmt.Sprintf("Using label selector: %s", input.LabelSelector))
listOptions.LabelSelector = input.LabelSelector
}
nodeList, err := k.Client.CoreV1().Nodes().List(ctx, listOptions)
if err != nil {
return []v1.Node{}, err
}
if nodeList == nil {
return []v1.Node{}, errors.New("pod list is nil")
}
return nodeList.Items, err
}
// ListPodsInput contains all the information needed to list Kubernetes pods
type ListPodsInput struct {
// Namespace is the Kubernetes namespace

View File

@@ -20,8 +20,17 @@ func main() {
switch subcommand {
case "initialize":
args := flag.NewFlagSet("scheduler-k3s:initialize", flag.ExitOnError)
taintScheduling := args.Bool("taint-scheduling", false, "taint-scheduling: add a taint against scheduling app workloads")
args.Parse(os.Args[2:])
err = scheduler_k3s.CommandInitialize()
err = scheduler_k3s.CommandInitialize(*taintScheduling)
case "cluster-add":
args := flag.NewFlagSet("scheduler-k3s:cluster-add", flag.ExitOnError)
allowUknownHosts := args.Bool("insecure-allow-unknown-hosts", false, "insecure-allow-unknown-hosts: allow unknown hosts")
taintScheduling := args.Bool("taint-scheduling", false, "taint-scheduling: add a taint against scheduling app workloads")
role := args.String("role", "worker", "role: [ server | worker ]")
args.Parse(os.Args[2:])
remoteHost := args.Arg(0)
err = scheduler_k3s.CommandClusterAdd(*role, remoteHost, *allowUknownHosts, *taintScheduling)
case "report":
args := flag.NewFlagSet("scheduler-k3s:report", flag.ExitOnError)
format := args.String("format", "stdout", "format: [ stdout | json ]")
@@ -48,6 +57,10 @@ func main() {
args := flag.NewFlagSet("scheduler-k3s:show-kubeconfig", flag.ExitOnError)
args.Parse(os.Args[2:])
err = scheduler_k3s.CommandShowKubeconfig()
case "uninstall":
args := flag.NewFlagSet("scheduler-k3s:uninstall", flag.ExitOnError)
args.Parse(os.Args[2:])
err = scheduler_k3s.CommandUninstall()
default:
err = fmt.Errorf("Invalid plugin subcommand call: %s", subcommand)
}

View File

@@ -1,7 +1,11 @@
package scheduler_k3s
import (
"context"
"crypto/rand"
"fmt"
"net"
"net/url"
"os"
"strings"
@@ -10,11 +14,77 @@ import (
)
// CommandInitialize initializes a k3s cluster on the local server
func CommandInitialize() error {
if common.FileExists("/usr/local/bin/k3s") {
func CommandInitialize(taintScheduling bool) error {
if err := isK3sInstalled(); err == nil {
return fmt.Errorf("k3s already installed, cannot re-initialize k3s")
}
networkInterface := common.PropertyGetDefault("scheduler-k3s", "--global", "network-interface", "eth0")
ifaces, err := net.Interfaces()
if err != nil {
return fmt.Errorf("Unable to get network interfaces: %w", err)
}
serverIp := ""
for _, iface := range ifaces {
if iface.Name == networkInterface {
addr, err := iface.Addrs()
if err != nil {
return fmt.Errorf("Unable to get network addresses for interface %s: %w", networkInterface, err)
}
for _, a := range addr {
if ipnet, ok := a.(*net.IPNet); ok {
if ipnet.IP.To4() != nil {
serverIp = ipnet.IP.String()
}
}
}
}
}
if len(serverIp) == 0 {
return fmt.Errorf(fmt.Sprintf("Unable to determine server ip address from network-interface %s", networkInterface))
}
common.LogInfo1Quiet("Initializing k3s")
common.LogInfo2Quiet("Updating apt")
aptUpdateCmd, err := common.CallExecCommand(common.ExecCommandInput{
Command: "apt-get",
Args: []string{
"update",
},
StreamStdio: true,
})
if err != nil {
return fmt.Errorf("Unable to call apt-get update command: %w", err)
}
if aptUpdateCmd.ExitCode != 0 {
return fmt.Errorf("Invalid exit code from apt-get update command: %d", aptUpdateCmd.ExitCode)
}
common.LogInfo2Quiet("Installing k3s dependencies")
aptInstallCmd, err := common.CallExecCommand(common.ExecCommandInput{
Command: "apt-get",
Args: []string{
"-y",
"install",
"ca-certificates",
"curl",
"open-iscsi",
"nfs-common",
"wireguard",
},
StreamStdio: true,
})
if err != nil {
return fmt.Errorf("Unable to call apt-get install command: %w", err)
}
if aptInstallCmd.ExitCode != 0 {
return fmt.Errorf("Invalid exit code from apt-get install command: %d", aptInstallCmd.ExitCode)
}
common.LogInfo2Quiet("Downloading k3s installer")
client := resty.New()
resp, err := client.R().
Get("https://get.k3s.io")
@@ -63,33 +133,358 @@ func CommandInitialize() error {
return fmt.Errorf("Unable to set k3s token: %w", err)
}
installerCmd := common.NewShellCmd(strings.Join([]string{
f.Name(),
nodeName := serverIp
n := 5
b := make([]byte, n)
if _, err := rand.Read(b); err != nil {
return fmt.Errorf("Unable to generate random node name: %w", err)
}
nodeName = strings.ReplaceAll(strings.ToLower(fmt.Sprintf("ip-%s-%s", nodeName, fmt.Sprintf("%X", b))), ".", "-")
args := []string{
// initialize the cluster
"--cluster-init",
// disable local-storage
"--disable", "local-storage",
// expose etcd metrics
"--etcd-expose-metrics",
// use wireguard for flannel
"--flannel-backend=wireguard-native",
// bind controller-manager to all interfaces
"--kube-controller-manager-arg", "bind-address=0.0.0.0",
// bind proxy metrics to all interfaces
"--kube-proxy-arg", "metrics-bind-address=0.0.0.0",
// bind scheduler to all interfaces
"--kube-scheduler-arg", "bind-address=0.0.0.0",
// gc terminated pods
"--kube-controller-manager-arg", "terminated-pod-gc-threshold=10",
// specify the node name
"--node-name", nodeName,
// allow access for the dokku user
"--write-kubeconfig-mode", "0644",
// specify a token
"--token", token,
}, " "))
if !installerCmd.Execute() {
return fmt.Errorf("Error installing k3s: %w", installerCmd.ExitError)
}
if taintScheduling {
args = append(args, "--node-taint", "node-role.kubernetes.io/master=true:NoSchedule")
}
common.LogInfo2Quiet("Running k3s installer")
installerCmd, err := common.CallExecCommand(common.ExecCommandInput{
Command: f.Name(),
Args: args,
StreamStdio: true,
})
if err != nil {
return fmt.Errorf("Unable to call k3s installer command: %w", err)
}
if installerCmd.ExitCode != 0 {
return fmt.Errorf("Invalid exit code from k3s installer command: %d", installerCmd.ExitCode)
}
if err := common.TouchFile(RegistryConfigPath); err != nil {
return fmt.Errorf("Error creating initial registries.yaml file")
}
registryAclCmd := common.NewShellCmd(strings.Join([]string{
"setfacl",
"-m",
"user:dokku:rwx",
RegistryConfigPath,
}, " "))
if !registryAclCmd.Execute() {
return fmt.Errorf("Error updating acls on k3s registries.yaml file: %w", registryAclCmd.ExitError)
common.LogInfo2Quiet("Setting registries.yaml permissions")
registryAclCmd, err := common.CallExecCommand(common.ExecCommandInput{
Command: "setfacl",
Args: []string{
"-m",
"user:dokku:rwx",
RegistryConfigPath,
},
StreamStdio: true,
})
if err != nil {
return fmt.Errorf("Unable to call setfacl command: %w", err)
}
if registryAclCmd.ExitCode != 0 {
return fmt.Errorf("Invalid exit code from setfacl command: %d", registryAclCmd.ExitCode)
}
common.LogInfo2Quiet("Installing k3s automatic upgrader")
upgradeCmd, err := common.CallExecCommand(common.ExecCommandInput{
Command: "kubectl",
Args: []string{
"apply",
"-f",
"https://github.com/rancher/system-upgrade-controller/releases/latest/download/system-upgrade-controller.yaml",
},
StreamStdio: true,
})
if err != nil {
return fmt.Errorf("Unable to call kubectl command: %w", err)
}
if upgradeCmd.ExitCode != 0 {
return fmt.Errorf("Invalid exit code from kubectl command: %d", upgradeCmd.ExitCode)
}
common.LogVerboseQuiet("Done")
return nil
}
// CommandClusterAdd adds a server to the k3s cluster
func CommandClusterAdd(role string, remoteHost string, allowUknownHosts bool, taintScheduling bool) error {
if err := isK3sInstalled(); err != nil {
return fmt.Errorf("k3s not installed, cannot join cluster")
}
if role != "server" && role != "worker" {
return fmt.Errorf("Invalid server-type: %s", role)
}
token := common.PropertyGet("scheduler-k3s", "--global", "token")
if len(token) == 0 {
return fmt.Errorf("Missing k3s token")
}
if taintScheduling && role == "worker" {
return fmt.Errorf("Taint scheduling can only be used on the server role")
}
networkInterface := common.PropertyGetDefault("scheduler-k3s", "--global", "network-interface", "eth0")
ifaces, err := net.Interfaces()
if err != nil {
return fmt.Errorf("Unable to get network interfaces: %w", err)
}
serverIp := ""
for _, iface := range ifaces {
if iface.Name == networkInterface {
addr, err := iface.Addrs()
if err != nil {
return fmt.Errorf("Unable to get network addresses for interface %s: %w", networkInterface, err)
}
for _, a := range addr {
if ipnet, ok := a.(*net.IPNet); ok {
if ipnet.IP.To4() != nil {
serverIp = ipnet.IP.String()
}
}
}
}
}
if len(serverIp) == 0 {
return fmt.Errorf(fmt.Sprintf("Unable to determine server ip address from network-interface %s", networkInterface))
}
// todo: check if k3s is installed on the remote host
k3sVersionCmd, err := common.CallExecCommand(common.ExecCommandInput{
Command: "k3s",
Args: []string{
"--version",
},
CaptureOutput: true,
})
if err != nil {
return fmt.Errorf("Unable to call k3s version command: %w", err)
}
if k3sVersionCmd.ExitCode != 0 {
return fmt.Errorf("Invalid exit code from k3s --version command: %d", k3sVersionCmd.ExitCode)
}
k3sVersion := ""
k3sVersionLines := strings.Split(string(k3sVersionCmd.Stdout), "\n")
if len(k3sVersionLines) > 0 {
k3sVersionParts := strings.Split(k3sVersionLines[0], " ")
if len(k3sVersionParts) != 4 {
return fmt.Errorf("Unable to get k3s version from k3s --version: %s", k3sVersionCmd.Stdout)
}
k3sVersion = k3sVersionParts[2]
}
common.LogDebug(fmt.Sprintf("k3s version: %s", k3sVersion))
common.LogInfo1(fmt.Sprintf("Joining %s to k3s cluster as %s", remoteHost, role))
common.LogInfo2Quiet("Updating apt")
aptUpdateCmd, err := common.CallSshCommand(common.SshCommandInput{
Command: "apt-get",
Args: []string{
"update",
},
AllowUknownHosts: allowUknownHosts,
RemoteHost: remoteHost,
StreamStdio: true,
Sudo: true,
})
if err != nil {
return fmt.Errorf("Unable to call apt-get update command over ssh: %w", err)
}
if aptUpdateCmd.ExitCode != 0 {
return fmt.Errorf("Invalid exit code from apt-get update command over ssh: %d", aptUpdateCmd.ExitCode)
}
common.LogInfo2Quiet("Installing k3s dependencies")
aptInstallCmd, err := common.CallSshCommand(common.SshCommandInput{
Command: "apt-get",
Args: []string{
"-y",
"install",
"ca-certificates",
"curl",
"open-iscsi",
"nfs-common",
"wireguard",
},
AllowUknownHosts: allowUknownHosts,
RemoteHost: remoteHost,
StreamStdio: true,
Sudo: true,
})
if err != nil {
return fmt.Errorf("Unable to call apt-get install command over ssh: %w", err)
}
if aptInstallCmd.ExitCode != 0 {
return fmt.Errorf("Invalid exit code from apt-get install command over ssh: %d", aptInstallCmd.ExitCode)
}
common.LogInfo2Quiet("Downloading k3s installer")
curlTask, err := common.CallSshCommand(common.SshCommandInput{
Command: "curl",
Args: []string{
"-o /tmp/k3s-installer.sh",
"https://get.k3s.io",
},
AllowUknownHosts: allowUknownHosts,
RemoteHost: remoteHost,
StreamStdio: true,
})
if err != nil {
return fmt.Errorf("Unable to call curl command over ssh: %w", err)
}
if curlTask.ExitCode != 0 {
return fmt.Errorf("Invalid exit code from curl command over ssh: %d", curlTask.ExitCode)
}
common.LogInfo2Quiet("Setting k3s installer permissions")
chmodCmd, err := common.CallSshCommand(common.SshCommandInput{
Command: "chmod",
Args: []string{
"0755",
"/tmp/k3s-installer.sh",
},
AllowUknownHosts: allowUknownHosts,
RemoteHost: remoteHost,
StreamStdio: true,
})
if err != nil {
return fmt.Errorf("Unable to call chmod command over ssh: %w", err)
}
if chmodCmd.ExitCode != 0 {
return fmt.Errorf("Invalid exit code from chmod command over ssh: %d", chmodCmd.ExitCode)
}
u, err := url.Parse(remoteHost)
if err != nil {
return fmt.Errorf("failed to parse remote host: %w", err)
}
nodeName := u.Hostname()
n := 5
b := make([]byte, n)
if _, err := rand.Read(b); err != nil {
return fmt.Errorf("Unable to generate random node name: %w", err)
}
nodeName = strings.ReplaceAll(strings.ToLower(fmt.Sprintf("ip-%s-%s", nodeName, fmt.Sprintf("%X", b))), ".", "-")
args := []string{
// disable local-storage
"--disable", "local-storage",
// add a node label
"--node-label=node_type=worker",
// use wireguard for flannel
"--flannel-backend=wireguard-native",
// specify the node name
"--node-name", nodeName,
// server to connect to as the main
"--server",
fmt.Sprintf("https://%s:6443", serverIp),
// specify a token
"--token",
token,
}
if role == "server" {
args = append([]string{"server"}, args...)
// expose etcd metrics
args = append(args, "--etcd-expose-metrics")
// bind controller-manager to all interfaces
args = append(args, "--kube-controller-manager-arg", "bind-address=0.0.0.0")
// bind proxy metrics to all interfaces
args = append(args, "--kube-proxy-arg", "metrics-bind-address=0.0.0.0")
// bind scheduler to all interfaces
args = append(args, "--kube-scheduler-arg", "bind-address=0.0.0.0")
// gc terminated pods
args = append(args, "--kube-controller-manager-arg", "terminated-pod-gc-threshold=10")
// allow access for the dokku user
args = append(args, "--write-kubeconfig-mode", "0644")
} else {
// disable etcd on workers
args = append(args, "--disable-etcd")
// disable apiserver on workers
args = append(args, "--disable-apiserver")
// disable controller-manager on workers
args = append(args, "--disable-controller-manager")
// disable scheduler on workers
args = append(args, "--disable-scheduler")
// bind proxy metrics to all interfaces
args = append(args, "--kube-proxy-arg", "metrics-bind-address=0.0.0.0")
}
if taintScheduling {
args = append(args, "--node-taint", "node-role.kubernetes.io/master=true:NoSchedule")
}
common.LogInfo2Quiet("Joining k3s cluster")
joinCmd, err := common.CallSshCommand(common.SshCommandInput{
Command: "/tmp/k3s-installer.sh",
Args: args,
AllowUknownHosts: allowUknownHosts,
RemoteHost: remoteHost,
StreamStdio: true,
Sudo: true,
})
if err != nil {
return fmt.Errorf("Unable to call k3s installer command over ssh: %w", err)
}
if joinCmd.ExitCode != 0 {
return fmt.Errorf("Invalid exit code from k3s installer command over ssh: %d", joinCmd.ExitCode)
}
if role == "worker" {
ctx := context.Background()
clientset, err := NewKubernetesClient()
if err != nil {
return fmt.Errorf("Unable to create kubernetes client: %w", err)
}
common.LogInfo2Quiet("Waiting for node to exist")
nodes, err := waitForNodeToExist(ctx, WaitForNodeToExistInput{
Clientset: clientset,
NodeName: nodeName,
RetryCount: 20,
})
if err != nil {
return fmt.Errorf("Error waiting for pod to exist: %w", err)
}
if len(nodes) == 0 {
return fmt.Errorf("Unable to find node after joining cluster, node will not be labeled kubernetes.io/role=worker")
}
common.LogInfo2Quiet("Labeling node kubernetes.io/role=worker")
err = clientset.LabelNode(ctx, LabelNodeInput{
Name: nodes[0].Name,
Key: "kubernetes.io/role",
Value: "worker",
})
if err != nil {
return fmt.Errorf("Unable to patch node: %w", err)
}
}
common.LogVerboseQuiet("Done")
return nil
}
@@ -132,3 +527,23 @@ func CommandShowKubeconfig() error {
return nil
}
func CommandUninstall() error {
if err := isK3sInstalled(); err != nil {
return fmt.Errorf("k3s not installed, cannot uninstall")
}
common.LogInfo1("Uninstalling k3s")
uninstallerCmd, err := common.CallExecCommand(common.ExecCommandInput{
Command: "/usr/local/bin/k3s-uninstall.sh",
StreamStdio: true,
})
if err != nil {
return fmt.Errorf("Unable to call k3s uninstaller command: %w", err)
}
if uninstallerCmd.ExitCode != 0 {
return fmt.Errorf("Invalid exit code from k3s uninstaller command: %d", uninstallerCmd.ExitCode)
}
return nil
}

View File

@@ -713,24 +713,24 @@ func TriggerSchedulerRun(scheduler string, appName string, envCount int, args []
dokkuRmContainer := os.Getenv("DOKKU_RM_CONTAINER")
if dokkuRmContainer == "" {
appRmContainer, err := common.CallPlugnTrigger(common.PlugnTriggerInput{
resp, err := common.CallPlugnTrigger(common.PlugnTriggerInput{
Trigger: "config-get",
Args: []string{appName, "DOKKU_RM_CONTAINER"},
CaptureOutput: true,
StreamStdio: false,
})
if err != nil {
globalRmContainer, err := common.CallPlugnTrigger(common.PlugnTriggerInput{
resp, err := common.CallPlugnTrigger(common.PlugnTriggerInput{
Trigger: "config-get-global",
Args: []string{"DOKKU_RM_CONTAINER"},
CaptureOutput: true,
StreamStdio: false,
})
if err == nil {
dokkuRmContainer = globalRmContainer.Stdout
dokkuRmContainer = strings.TrimSpace(resp.Stdout)
}
} else {
dokkuRmContainer = appRmContainer.Stdout
dokkuRmContainer = strings.TrimSpace(resp.Stdout)
}
}
if dokkuRmContainer == "" {