Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions cloudbuild/kne_test.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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
Expand Down
18 changes: 9 additions & 9 deletions cloudbuild/vendors_test.sh
100644 → 100755
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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
103 changes: 100 additions & 3 deletions cluster/kubeadm/kubeadm.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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 {
Expand Down
176 changes: 172 additions & 4 deletions cluster/kubeadm/kubeadm_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,9 @@ package kubeadm

import (
"os"
"path/filepath"
"testing"
"time"

"github.com/openconfig/gnmi/errdiff"
kexec "github.com/openconfig/kne/exec"
Expand Down Expand Up @@ -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",
Expand All @@ -58,15 +60,15 @@ 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",
}, {
desc: "failed to restart kubelet",
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",
Expand Down Expand Up @@ -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)
}
})
}
}
Loading
Loading