diff --git a/cloudbuild/kne_test.sh b/cloudbuild/kne_test.sh index 66bbefe9..7e65f3b5 100755 --- a/cloudbuild/kne_test.sh +++ b/cloudbuild/kne_test.sh @@ -20,8 +20,8 @@ gopath=$(go env GOPATH) export PATH=${PATH}:$gopath/bin # Replace existing kne repo with new version -rm -r "$HOME/kne" -cp -r /tmp/workspace "$HOME/kne" +rm -rf "$HOME/kne" +cp -rf /tmp/workspace "$HOME/kne" # Rebuild the kne cli pushd "$HOME/kne/kne_cli" @@ -94,7 +94,7 @@ $cli teardown kne/deploy/kne/kubeadm.yaml ## Create a kubeadm single node cluster sudo kubeadm init --cri-socket unix:///var/run/containerd/containerd.sock --pod-network-cidr 10.244.0.0/16 mkdir -p "$HOME"/.kube -sudo cp /etc/kubernetes/admin.conf "$HOME"/.kube/config +sudo cp -f /etc/kubernetes/admin.conf "$HOME"/.kube/config sudo chown "$(id -u)":"$(id -g)" "$HOME"/.kube/config kubectl taint nodes --all node-role.kubernetes.io/control-plane- # allows pods to be scheduled on control plane node kubectl apply -f kne/manifests/flannel/manifest.yaml diff --git a/cloudbuild/vendors_test.sh b/cloudbuild/vendors_test.sh old mode 100644 new mode 100755 index ba792450..c97f0a8f --- a/cloudbuild/vendors_test.sh +++ b/cloudbuild/vendors_test.sh @@ -19,9 +19,9 @@ export PATH=${PATH}:/usr/local/go/bin gopath=$(go env GOPATH) export PATH=${PATH}:$gopath/bin -# Replace exisiting kne repo with new version -rm -r "$HOME/kne" -cp -r /tmp/workspace "$HOME/kne" +# Replace existing kne repo with new version +rm -rf "$HOME/kne" +cp -rf /tmp/workspace "$HOME/kne" # Rebuild the kne cli pushd "$HOME/kne/kne_cli" @@ -38,10 +38,10 @@ popd # Run an ondatra test pushd "$HOME/kne/cloudbuild" go test -v vendors/vendors_test.go \ - -testbed testbed.textproto \ - -topology topology.textproto \ - -vendor_creds ARISTA/admin/admin \ - -vendor_creds JUNIPER/root/Google123 \ - -vendor_creds CISCO/cisco/cisco123 \ - -vendor_creds NOKIA/admin/NokiaSrl1! + -testbed testbed.textproto \ + -topology topology.textproto \ + -vendor_creds ARISTA/admin/admin \ + -vendor_creds JUNIPER/root/Google123 \ + -vendor_creds CISCO/cisco/cisco123 \ + -vendor_creds NOKIA/admin/NokiaSrl1! popd diff --git a/cluster/kubeadm/kubeadm.go b/cluster/kubeadm/kubeadm.go index 39731cac..a9727657 100644 --- a/cluster/kubeadm/kubeadm.go +++ b/cluster/kubeadm/kubeadm.go @@ -4,15 +4,104 @@ import ( "fmt" "os" "strings" + "time" "github.com/openconfig/kne/exec/run" log "k8s.io/klog/v2" ) var ( - kubeadmFlagPath = "/var/lib/kubelet/kubeadm-flags.env" + kubeadmFlagPath = "/var/lib/kubelet/kubeadm-flags.env" + kubeAPIServerManifest = "/etc/kubernetes/manifests/kube-apiserver.yaml" + apiserverWaitTimeout = 60 * time.Second + apiserverPollInterval = time.Second + sleepFn = time.Sleep ) +// SetServiceNodePortRange sets the service node port range in the kube-apiserver manifest. +func SetServiceNodePortRange(portRange string) error { + log.Infof("Setting service node port range to %q...", portRange) + b, err := os.ReadFile(kubeAPIServerManifest) + if err != nil { + if os.IsNotExist(err) { + return nil + } + // If read fails (e.g. due to permissions on root-owned manifest), try reading via sudo. + var sudoErr error + b, sudoErr = run.OutCommand("sudo", "cat", kubeAPIServerManifest) + if sudoErr != nil { + return fmt.Errorf("failed to read %s: %w", kubeAPIServerManifest, err) + } + } + content := string(b) + if strings.Contains(content, "--service-node-port-range=") { + return nil + } + target := " - --service-cluster-ip-range=" + idx := strings.Index(content, target) + if idx == -1 { + target = " - kube-apiserver\n" + idx = strings.Index(content, target) + if idx == -1 { + return fmt.Errorf("could not find insertion point in %s", kubeAPIServerManifest) + } + } + endOfLine := strings.Index(content[idx:], "\n") + if endOfLine == -1 { + endOfLine = len(content) - idx + } + insertPos := idx + endOfLine + 1 + flag := fmt.Sprintf(" - --service-node-port-range=%s\n", portRange) + newContent := content[:insertPos] + flag + content[insertPos:] + + f, err := os.CreateTemp("", "kne-apiserver-*.yaml") + if err != nil { + return err + } + defer func() { + if err := os.Remove(f.Name()); err != nil && !os.IsNotExist(err) { + log.Warningf("Failed to remove temp file %q: %v", f.Name(), err) + } + }() + if _, err := f.WriteString(newContent); err != nil { + _ = f.Close() + return err + } + if err := f.Close(); err != nil { + return err + } + if err := run.LogCommand("sudo", "cp", "-f", f.Name(), kubeAPIServerManifest); err != nil { + return err + } + return waitForAPIServer(portRange) +} + +func waitForAPIServer(portRange string) error { + log.Infof("Waiting for kube-apiserver to restart with service node port range %q...", portRange) + deadline := time.Now().Add(apiserverWaitTimeout) + expectedFlag := fmt.Sprintf("--service-node-port-range=%s", portRange) + for { + sleepFn(apiserverPollInterval) + out, err := run.OutCommand("kubectl", "get", "pod", "-n", "kube-system", "-l", "component=kube-apiserver", "-o", "jsonpath={.items[*].spec.containers[*].command}") + if err != nil { + log.V(1).Infof("kube-apiserver not ready yet: %v", err) + } else if !strings.Contains(string(out), expectedFlag) { + log.V(1).Infof("kube-apiserver has not restarted with %q yet", expectedFlag) + } else { + readyOut, err := run.OutCommand("kubectl", "get", "--raw", "/readyz") + if err != nil || strings.TrimSpace(string(readyOut)) != "ok" { + log.V(1).Infof("kube-apiserver /readyz not ready yet: %v", err) + } else { + log.Infof("kube-apiserver is ready with service node port range %q", portRange) + return nil + } + } + if time.Now().After(deadline) { + return fmt.Errorf("timed out waiting for kube-apiserver to become ready after setting service node port range") + } + } +} + // EnableCredentialProvider enables a credential provider according // to the specified config file on the kubelet. func EnableCredentialProvider(cfgPath string) error { @@ -36,11 +125,19 @@ func EnableCredentialProvider(cfgPath string) error { if err != nil { return err } - defer os.RemoveAll(f.Name()) + defer func() { + if err := os.Remove(f.Name()); err != nil && !os.IsNotExist(err) { + log.Warningf("Failed to remove temp file %q: %v", f.Name(), err) + } + }() if _, err := f.WriteString(s); err != nil { + _ = f.Close() + return err + } + if err := f.Close(); err != nil { return err } - if err := run.LogCommand("sudo", "cp", f.Name(), kubeadmFlagPath); err != nil { + if err := run.LogCommand("sudo", "cp", "-f", f.Name(), kubeadmFlagPath); err != nil { return err } if err := run.LogCommand("sudo", "systemctl", "restart", "kubelet"); err != nil { diff --git a/cluster/kubeadm/kubeadm_test.go b/cluster/kubeadm/kubeadm_test.go index c70d7f0c..67cad6c0 100644 --- a/cluster/kubeadm/kubeadm_test.go +++ b/cluster/kubeadm/kubeadm_test.go @@ -2,7 +2,9 @@ package kubeadm import ( "os" + "path/filepath" "testing" + "time" "github.com/openconfig/gnmi/errdiff" kexec "github.com/openconfig/kne/exec" @@ -39,12 +41,12 @@ func TestEnableCredentialProvider(t *testing.T) { cfgPath: cfg.Name(), resp: []fexec.Response{ {Cmd: "sudo", Args: []string{"kubeadm", "upgrade", "node", "phase", "kubelet-config"}}, - {Cmd: "sudo", Args: []string{"cp", ".*", kubeadmFlagPath}}, + {Cmd: "sudo", Args: []string{"cp", "-f", ".*", kubeadmFlagPath}}, {Cmd: "sudo", Args: []string{"systemctl", "restart", "kubelet"}}, }, }, { desc: "config file not found", - cfgPath: "dne", + cfgPath: "nonexistent", wantErr: "config file not found", }, { desc: "failed to upgrade kubelet", @@ -58,7 +60,7 @@ func TestEnableCredentialProvider(t *testing.T) { cfgPath: cfg.Name(), resp: []fexec.Response{ {Cmd: "sudo", Args: []string{"kubeadm", "upgrade", "node", "phase", "kubelet-config"}}, - {Cmd: "sudo", Args: []string{"cp", ".*", kubeadmFlagPath}, Err: "failed to copy"}, + {Cmd: "sudo", Args: []string{"cp", "-f", ".*", kubeadmFlagPath}, Err: "failed to copy"}, }, wantErr: "failed to copy", }, { @@ -66,7 +68,7 @@ func TestEnableCredentialProvider(t *testing.T) { cfgPath: cfg.Name(), resp: []fexec.Response{ {Cmd: "sudo", Args: []string{"kubeadm", "upgrade", "node", "phase", "kubelet-config"}}, - {Cmd: "sudo", Args: []string{"cp", ".*", kubeadmFlagPath}}, + {Cmd: "sudo", Args: []string{"cp", "-f", ".*", kubeadmFlagPath}}, {Cmd: "sudo", Args: []string{"systemctl", "restart", "kubelet"}, Err: "failed to restart kubelet"}, }, wantErr: "failed to restart kubelet", @@ -94,3 +96,169 @@ func checkCmds(t *testing.T, cmds *fexec.Command) { t.Errorf("%v", err) } } + +func TestSetServiceNodePortRange(t *testing.T) { + origKubeAPIServerManifest := kubeAPIServerManifest + origSleepFn := sleepFn + origTimeout := apiserverWaitTimeout + origPollInterval := apiserverPollInterval + defer func() { + kubeAPIServerManifest = origKubeAPIServerManifest + sleepFn = origSleepFn + apiserverWaitTimeout = origTimeout + apiserverPollInterval = origPollInterval + }() + sleepFn = func(time.Duration) {} + apiserverWaitTimeout = 50 * time.Millisecond + apiserverPollInterval = time.Millisecond + + manifestWithClusterIP := `apiVersion: v1 +kind: Pod +metadata: + name: kube-apiserver +spec: + containers: + - command: + - kube-apiserver + - --service-cluster-ip-range=10.96.0.0/12 + - --advertise-address=192.168.1.10 +` + manifestWithKubeAPIServer := `apiVersion: v1 +kind: Pod +metadata: + name: kube-apiserver +spec: + containers: + - command: + - kube-apiserver + - --advertise-address=192.168.1.10 +` + manifestWithExistingRange := `apiVersion: v1 +kind: Pod +metadata: + name: kube-apiserver +spec: + containers: + - command: + - kube-apiserver + - --service-node-port-range=10000-32767 +` + manifestNoMatch := `apiVersion: v1 +kind: Pod +metadata: + name: other-pod +spec: + containers: + - command: + - other-command +` + + tests := []struct { + desc string + manifestData string + nonExistent bool + portRange string + waitTimeout time.Duration + resp []fexec.Response + wantErr string + }{{ + desc: "success with service-cluster-ip-range target", + manifestData: manifestWithClusterIP, + portRange: "10000-32767", + resp: []fexec.Response{ + {Cmd: "sudo", Args: []string{"cp", "-f", ".*", ".*"}}, + {Cmd: "kubectl", Args: []string{"get", "pod", "-n", "kube-system", "-l", "component=kube-apiserver", "-o", ".*"}, Stdout: "--service-node-port-range=10000-32767"}, + {Cmd: "kubectl", Args: []string{"get", "--raw", "/readyz"}, Stdout: "ok"}, + }, + }, { + desc: "success with kube-apiserver fallback target", + manifestData: manifestWithKubeAPIServer, + portRange: "10000-32767", + resp: []fexec.Response{ + {Cmd: "sudo", Args: []string{"cp", "-f", ".*", ".*"}}, + {Cmd: "kubectl", Args: []string{"get", "pod", "-n", "kube-system", "-l", "component=kube-apiserver", "-o", ".*"}, Stdout: "--service-node-port-range=10000-32767"}, + {Cmd: "kubectl", Args: []string{"get", "--raw", "/readyz"}, Stdout: "ok"}, + }, + }, { + desc: "manifest not found returns nil", + nonExistent: true, + portRange: "10000-32767", + }, { + desc: "manifest already contains service-node-port-range returns nil", + manifestData: manifestWithExistingRange, + portRange: "10000-32767", + }, { + desc: "could not find insertion point", + manifestData: manifestNoMatch, + portRange: "10000-32767", + wantErr: "could not find insertion point", + }, { + desc: "failed to copy modified manifest", + manifestData: manifestWithClusterIP, + portRange: "10000-32767", + resp: []fexec.Response{ + {Cmd: "sudo", Args: []string{"cp", "-f", ".*", ".*"}, Err: "failed to copy"}, + }, + wantErr: "failed to copy", + }, { + desc: "success after retry", + manifestData: manifestWithClusterIP, + portRange: "10000-32767", + resp: []fexec.Response{ + {Cmd: "sudo", Args: []string{"cp", "-f", ".*", ".*"}}, + {Cmd: "kubectl", Args: []string{"get", "pod", "-n", "kube-system", "-l", "component=kube-apiserver", "-o", ".*"}, Err: "connection refused"}, + {Cmd: "kubectl", Args: []string{"get", "pod", "-n", "kube-system", "-l", "component=kube-apiserver", "-o", ".*"}, Stdout: "old-args"}, + {Cmd: "kubectl", Args: []string{"get", "pod", "-n", "kube-system", "-l", "component=kube-apiserver", "-o", ".*"}, Stdout: "--service-node-port-range=10000-32767"}, + {Cmd: "kubectl", Args: []string{"get", "--raw", "/readyz"}, Err: "not ready"}, + {Cmd: "kubectl", Args: []string{"get", "pod", "-n", "kube-system", "-l", "component=kube-apiserver", "-o", ".*"}, Stdout: "--service-node-port-range=10000-32767"}, + {Cmd: "kubectl", Args: []string{"get", "--raw", "/readyz"}, Stdout: "ok"}, + }, + }, { + desc: "timed out waiting for apiserver", + manifestData: manifestWithClusterIP, + portRange: "10000-32767", + waitTimeout: time.Nanosecond, + resp: []fexec.Response{ + {Cmd: "sudo", Args: []string{"cp", "-f", ".*", ".*"}}, + {Cmd: "kubectl", Args: []string{"get", "pod", "-n", "kube-system", "-l", "component=kube-apiserver", "-o", ".*"}, Err: "connection refused"}, + }, + wantErr: "timed out waiting for kube-apiserver", + }} + + for _, tt := range tests { + t.Run(tt.desc, func(t *testing.T) { + if tt.waitTimeout != 0 { + apiserverWaitTimeout = tt.waitTimeout + } else { + apiserverWaitTimeout = 50 * time.Millisecond + } + if tt.nonExistent { + kubeAPIServerManifest = filepath.Join(t.TempDir(), "nonexistent.yaml") + } else { + f, err := os.CreateTemp(t.TempDir(), "kube-apiserver-*.yaml") + if err != nil { + t.Fatalf("Failed to create temp file: %v", err) + } + if _, err := f.WriteString(tt.manifestData); err != nil { + t.Fatalf("Failed to write temp manifest: %v", err) + } + if err := f.Close(); err != nil { + t.Fatalf("Failed to close temp manifest: %v", err) + } + kubeAPIServerManifest = f.Name() + } + + fexec.LogCommand = func(s string) { + t.Logf("%s: %s", tt.desc, s) + } + cmds := fexec.Commands(tt.resp) + kexec.Command = cmds.Command + defer checkCmds(t, cmds) + + err := SetServiceNodePortRange(tt.portRange) + if s := errdiff.Substring(err, tt.wantErr); s != "" { + t.Fatalf("unexpected error: %s", s) + } + }) + } +} diff --git a/deploy/deploy.go b/deploy/deploy.go index 77c3b1f3..8b119cd6 100644 --- a/deploy/deploy.go +++ b/deploy/deploy.go @@ -434,6 +434,7 @@ type KubeadmSpec struct { Network string `yaml:"network"` AllowControlPlaneScheduling bool `yaml:"allowControlPlaneScheduling"` ImageRepository string `yaml:"imageRepository"` + ServiceNodePortRange string `yaml:"serviceNodePortRange"` } func (k *KubeadmSpec) checkDependencies() error { @@ -490,6 +491,13 @@ func (k *KubeadmSpec) Deploy(ctx context.Context) error { if err := os.WriteFile(filepath.Join(kubeDir, "config"), b, 0600); err != nil { return err } + portRange := "10000-32767" + if k.ServiceNodePortRange != "" { + portRange = k.ServiceNodePortRange + } + if err := kubeadm.SetServiceNodePortRange(portRange); err != nil { + log.Warningf("Failed to set service node port range: %v", err) + } if k.AllowControlPlaneScheduling { if err := run.LogCommand("kubectl", "taint", "nodes", "--all", "node-role.kubernetes.io/control-plane:NoSchedule-"); err != nil { return err diff --git a/deploy/deploy_test.go b/deploy/deploy_test.go index f808f521..5882c5c1 100644 --- a/deploy/deploy_test.go +++ b/deploy/deploy_test.go @@ -118,6 +118,16 @@ func TestKubeadmSpec(t *testing.T) { {Cmd: "sudo", Args: []string{"kubeadm", "init", "--image-repository", "us-west1-docker.pkg.dev/kne-external/kne"}}, {Cmd: "sudo", Args: []string{"cat", "/etc/kubernetes/admin.conf"}}, }, + }, { + desc: "custom service node port range", + k: &KubeadmSpec{ + ServiceNodePortRange: "20000-30000", + }, + resp: []fexec.Response{ + {Cmd: "sudo", Args: []string{"kubeadm", "init", "--image-repository", "us-west1-docker.pkg.dev/kne-external/kne"}}, + {Cmd: "sudo", Args: []string{"cat", "/etc/kubernetes/admin.conf"}}, + {Cmd: "docker", Args: []string{"network", "create", "kne-kubeadm-.*"}}, + }, }} for _, tt := range tests { t.Run(tt.desc, func(t *testing.T) { @@ -1050,7 +1060,7 @@ func TestMetalLBSpec(t *testing.T) { } mi, err := mfake.NewSimpleClientset(tt.mObjects...) if err != nil { - t.Fatalf("faild to create fake client: %v", err) + t.Fatalf("failed to create fake client: %v", err) } tt.m.SetKClient(ki)