diff --git a/cmd/kured/Dockerfile b/cmd/kured/Dockerfile index 1163995..6ba08e1 100644 --- a/cmd/kured/Dockerfile +++ b/cmd/kured/Dockerfile @@ -1,7 +1,4 @@ FROM alpine:3.12 RUN apk add --no-cache ca-certificates tzdata -# NB: you may need to update RBAC permissions when upgrading kubectl - see kured-rbac.yaml for details -ADD https://storage.googleapis.com/kubernetes-release/release/v1.18.8/bin/linux/amd64/kubectl /usr/bin/kubectl -RUN chmod 0755 /usr/bin/kubectl COPY ./kured /usr/bin/kured ENTRYPOINT ["/usr/bin/kured"] diff --git a/cmd/kured/main.go b/cmd/kured/main.go index 2bebf93..6489cf6 100644 --- a/cmd/kured/main.go +++ b/cmd/kured/main.go @@ -12,9 +12,11 @@ import ( log "github.com/sirupsen/logrus" "github.com/spf13/cobra" + v1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/client-go/kubernetes" "k8s.io/client-go/rest" + kubectldrain "k8s.io/kubectl/pkg/drain" "github.com/prometheus/client_golang/prometheus" "github.com/prometheus/client_golang/prometheus/promhttp" @@ -232,28 +234,33 @@ func release(lock *daemonsetlock.DaemonSetLock) { } } -func drain(nodeID string) { - log.Infof("Draining node %s", nodeID) +func drain(drainer *kubectldrain.Helper, node *v1.Node) { + nodename := node.GetName() + + log.Infof("Draining node %s", nodename) if slackHookURL != "" { - if err := slack.NotifyDrain(slackHookURL, slackUsername, slackChannel, nodeID); err != nil { + if err := slack.NotifyDrain(slackHookURL, slackUsername, slackChannel, nodename); err != nil { log.Warnf("Error notifying slack: %v", err) } } - drainCmd := newCommand("/usr/bin/kubectl", "drain", - "--ignore-daemonsets", "--delete-local-data", "--force", nodeID) + // The drainer helper already has DeleteLocalData, + // IgnoreAllDaemonSets, and Force flags set to true. + if err := kubectldrain.RunCordonOrUncordon(drainer, node, true); err != nil { + log.Fatal("Error cordonning %s: %v", nodename, err) + } - if err := drainCmd.Run(); err != nil { - log.Fatalf("Error invoking drain command: %v", err) + if err := kubectldrain.RunNodeDrain(drainer, nodename); err != nil { + log.Fatal("Error draining %s: %v", nodename, err) } } -func uncordon(nodeID string) { - log.Infof("Uncordoning node %s", nodeID) - uncordonCmd := newCommand("/usr/bin/kubectl", "uncordon", nodeID) - if err := uncordonCmd.Run(); err != nil { - log.Fatalf("Error invoking uncordon command: %v", err) +func uncordon(drainer *kubectldrain.Helper, node *v1.Node) { + nodename := node.GetName() + log.Infof("Uncordoning node %s", nodename) + if err := kubectldrain.RunCordonOrUncordon(drainer, node, false); err != nil { + log.Fatal("Error uncordonning %s: %v", nodename, err) } } @@ -300,12 +307,25 @@ func rebootAsRequired(nodeID string, window *timewindow.TimeWindow, TTL time.Dur log.Fatal(err) } + drainer := &kubectldrain.Helper{ + Client: client, + Force: true, + DeleteLocalData: true, + IgnoreAllDaemonSets: true, + ErrOut: os.Stderr, + Out: os.Stdout, + } + lock := daemonsetlock.New(client, nodeID, dsNamespace, dsName, lockAnnotation) nodeMeta := nodeMeta{} if holding(lock, &nodeMeta) { if !nodeMeta.Unschedulable { - uncordon(nodeID) + node, err := client.CoreV1().Nodes().Get(context.TODO(), nodeID, metav1.GetOptions{}) + if err != nil { + log.Fatal(err) + } + uncordon(drainer, node) } release(lock) } @@ -322,7 +342,7 @@ func rebootAsRequired(nodeID string, window *timewindow.TimeWindow, TTL time.Dur if acquire(lock, &nodeMeta, TTL) { if !nodeMeta.Unschedulable { - drain(nodeID) + drain(drainer, node) } commandReboot(nodeID) for { diff --git a/go.mod b/go.mod index 13b85cd..bf7eb6d 100644 --- a/go.mod +++ b/go.mod @@ -12,4 +12,5 @@ require ( github.com/spf13/cobra v0.0.0-20181127133106-d2d81d9a96e2 k8s.io/apimachinery v0.18.8 k8s.io/client-go v0.18.8 + k8s.io/kubectl v0.18.8 )