diff --git a/Makefile b/Makefile index bd50b240..bf8f7884 100644 --- a/Makefile +++ b/Makefile @@ -94,17 +94,44 @@ e2e: ## Run e2e tests (requires: make deploy-bink). V=1 for verbose. RUN= # each package, even though we just have one here--but I really like streaming # output...). rm -rf $(ARTIFACTS) - cd test/e2e && KUBECONFIG=$(abspath $(KUBECONFIG_BINK)) BINK_CLUSTER_NAME=$(BINK_CLUSTER_NAME) \ - $(if $(BINK_NODE_IMAGE),BINK_NODE_IMAGE=$(BINK_NODE_IMAGE)) \ + cd test/e2e && KUBECONFIG=$(abspath $(KUBECONFIG_BINK)) \ + E2E_PROVIDER=bink \ + BINK_CLUSTER_NAME=$(BINK_CLUSTER_NAME) \ BINK_NODE_DISK_IMAGE=$(BINK_NODE_DISK_IMAGE) \ - BINK_LOCAL_REGISTRY_NODE_IMAGE=$(BINK_LOCAL_REGISTRY_NODE_IMAGE) \ + E2E_NODE_IMAGE_REGISTRY=$(BINK_LOCAL_REGISTRY_NODE_IMAGE) \ ARTIFACTS=$(ARTIFACTS) \ - BINK_NODE_IMAGE_DIGEST=$$(skopeo inspect --tls-verify=false --format '{{.Digest}}' docker://localhost:5000/node:latest) \ - BINK_NODE_IMAGE_UPDATE_DIGEST=$$(skopeo inspect --tls-verify=false docker://localhost:5000/node:update | jq -r '.Digest') \ - BINK_NODE_IMAGE_UPDATE2_DIGEST=$$(skopeo inspect --tls-verify=false docker://localhost:5000/node:update2 | jq -r '.Digest') \ + E2E_NODE_IMAGE_DIGEST=$$(skopeo inspect --tls-verify=false --format '{{.Digest}}' docker://localhost:5000/node:latest) \ + E2E_NODE_IMAGE_UPDATE_DIGEST=$$(skopeo inspect --tls-verify=false docker://localhost:5000/node:update | jq -r '.Digest') \ + E2E_NODE_IMAGE_UPDATE2_DIGEST=$$(skopeo inspect --tls-verify=false docker://localhost:5000/node:update2 | jq -r '.Digest') \ E2E_REGISTRY_USER=$(E2E_REGISTRY_USER) E2E_REGISTRY_PASSWORD=$(E2E_REGISTRY_PASSWORD) \ go test -timeout 30m -count=1 $(if $(V),-v) $(if $(RUN),-run $(RUN)) . +# EKS e2e settings +EKS_CLUSTER_NAME ?= +EKS_NODE_GROUP ?= +AWS_REGION ?= +EKS_NODE_IMAGE_REF ?= +EKS_NODE_IMAGE_UPDATE_REF ?= +EKS_NODE_IMAGE_UPDATE2_REF ?= +EKS_REGISTRY_USER ?= +EKS_REGISTRY_PASSWORD ?= + +.PHONY: e2e-eks +e2e-eks: ## Run e2e tests against EKS. V=1 for verbose. RUN= to filter. + rm -rf $(ARTIFACTS) + cd test/e2e && KUBECONFIG="$(KUBECONFIG)" \ + E2E_PROVIDER=eks \ + EKS_CLUSTER_NAME="$(EKS_CLUSTER_NAME)" \ + EKS_NODE_GROUP="$(EKS_NODE_GROUP)" \ + AWS_REGION="$(AWS_REGION)" \ + E2E_NODE_IMAGE_REF="$(EKS_NODE_IMAGE_REF)" \ + E2E_NODE_IMAGE_UPDATE_REF="$(EKS_NODE_IMAGE_UPDATE_REF)" \ + E2E_NODE_IMAGE_UPDATE2_REF="$(EKS_NODE_IMAGE_UPDATE2_REF)" \ + E2E_REGISTRY_USER="$(EKS_REGISTRY_USER)" \ + E2E_REGISTRY_PASSWORD="$(EKS_REGISTRY_PASSWORD)" \ + ARTIFACTS="$(ARTIFACTS)" \ + go test -timeout 30m -count=1 $(if $(V),-v) $(if $(RUN),-run $(RUN)) . + ##@ Build .PHONY: build @@ -202,7 +229,7 @@ deploy-bink: start-bink build-update-image kustomize ## Deploy to a bink cluster .PHONY: gather-bink gather-bink: ## Gather diagnostic logs from the bink cluster. - KUBECONFIG=$(abspath $(KUBECONFIG_BINK)) BINK_CLUSTER_NAME=$(BINK_CLUSTER_NAME) \ + KUBECONFIG=$(abspath $(KUBECONFIG_BINK)) \ hack/gather-logs.sh $(ARTIFACTS)/gather-bink controller .PHONY: teardown-bink diff --git a/go.mod b/go.mod index 53c04916..1a87620c 100644 --- a/go.mod +++ b/go.mod @@ -3,6 +3,11 @@ module github.com/bootc-dev/bootc-operator go 1.26.0 require ( + github.com/aws/aws-sdk-go-v2 v1.47.0 + github.com/aws/aws-sdk-go-v2/config v1.33.5 + github.com/aws/aws-sdk-go-v2/service/autoscaling v1.78.0 + github.com/aws/aws-sdk-go-v2/service/ec2 v1.332.0 + github.com/aws/aws-sdk-go-v2/service/eks v1.99.0 github.com/distribution/reference v0.6.0 github.com/fsnotify/fsnotify v1.10.1 github.com/go-logr/logr v1.4.4 @@ -18,6 +23,18 @@ require ( require ( github.com/Azure/go-ansiterm v0.0.0-20250102033503-faa5f7b0171c // indirect github.com/MakeNowJust/heredoc v1.0.0 // indirect + github.com/aws/aws-sdk-go-v2/credentials v1.20.5 // indirect + github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.20.0 // indirect + github.com/aws/aws-sdk-go-v2/internal/configsources v1.5.3 // indirect + github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.8.3 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.5.3 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.19 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.14.3 // indirect + github.com/aws/aws-sdk-go-v2/service/signin v1.10.0 // indirect + github.com/aws/aws-sdk-go-v2/service/sso v1.38.0 // indirect + github.com/aws/aws-sdk-go-v2/service/ssooidc v1.43.0 // indirect + github.com/aws/aws-sdk-go-v2/service/sts v1.51.0 // indirect + github.com/aws/smithy-go v1.28.1 // indirect github.com/beorn7/perks v1.0.1 // indirect github.com/blang/semver/v4 v4.0.0 // indirect github.com/cespare/xxhash/v2 v2.3.0 // indirect diff --git a/go.sum b/go.sum index ebfc6f46..381f02e6 100644 --- a/go.sum +++ b/go.sum @@ -4,6 +4,40 @@ github.com/MakeNowJust/heredoc v1.0.0 h1:cXCdzVdstXyiTqTvfqk9SDHpKNjxuom+DOlyEeQ github.com/MakeNowJust/heredoc v1.0.0/go.mod h1:mG5amYoWBHf8vpLOuehzbGGw0EHxpZZ6lCpQ4fNJ8LE= github.com/Masterminds/semver/v3 v3.4.0 h1:Zog+i5UMtVoCU8oKka5P7i9q9HgrJeGzI9SA1Xbatp0= github.com/Masterminds/semver/v3 v3.4.0/go.mod h1:4V+yj/TJE1HU9XfppCwVMZq3I84lprf4nC11bSS5beM= +github.com/aws/aws-sdk-go-v2 v1.47.0 h1:0jsHallhJCeaU0Ko48c/3FK1ctOQ7NpzggxriJOQ8MQ= +github.com/aws/aws-sdk-go-v2 v1.47.0/go.mod h1:bttEH6JqnUL8LepvDVfdrds/fZ5bCIxzpe3abyUrhDU= +github.com/aws/aws-sdk-go-v2/config v1.33.5 h1:UA1dmokBFOLFoOyVBhO6HjM6edy0MIk5AZSkJVcksQw= +github.com/aws/aws-sdk-go-v2/config v1.33.5/go.mod h1:Dop8axzz0xx38GExIYWXdeyc8QQ7Cr+nPsxpD/LYy4U= +github.com/aws/aws-sdk-go-v2/credentials v1.20.5 h1:wklUVvHMc9xTQ3rcp49/ISpiMnhbCicJcA6n6S8m7J8= +github.com/aws/aws-sdk-go-v2/credentials v1.20.5/go.mod h1:fyEdrn6ccLFOkoK84j5bQyGTxp9zPt5l2XMhxf4DVZs= +github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.20.0 h1:AM4hHjww+PSFtt6E+UrBrPlZkWsePCLEt9AjkfQX+yM= +github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.20.0/go.mod h1:3x/yXezeQjpOvBb4jEMxrS8SXvpdvJ5abv6l5c1gWM8= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.5.3 h1:Hp/VgjP0BysR3OgLlR057Vz2LcbbVnoWeJ+3qWiS/fY= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.5.3/go.mod h1:nwGV5qw7F1IZPgxCvA/ph8N2TAuz+BkRG/bXn808qMA= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.8.3 h1:MUaM4f+kj1ZIBPZfUS8cxP1GKXXZtHJjAthy93AN7SM= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.8.3/go.mod h1:6YmVmEVRI5ZZzRjCSsb9SryKH0hAlMRdgA7kG9aDvBU= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.5.3 h1:fuSCw4Z2qfRCztMPO3GXJNSiEp6Wee+WOLwrHHUMy9c= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.5.3/go.mod h1:6SxcHheD1pPR5+kWm1wGvjlL/YqUsh267sAfEmN4K7A= +github.com/aws/aws-sdk-go-v2/service/autoscaling v1.78.0 h1:CN7ZkNEZb5Ob0DtntBQLE7cdpT13gzS1Gn+QMoZOjHA= +github.com/aws/aws-sdk-go-v2/service/autoscaling v1.78.0/go.mod h1:nkWNnRTHDlkZZrZzhmcPrOkoF+werzJzCOHGpbIpcfA= +github.com/aws/aws-sdk-go-v2/service/ec2 v1.332.0 h1:T9rFYUxhZEBmnIgrzKZjd/0xMsWRF+brY/Gpgf5rc40= +github.com/aws/aws-sdk-go-v2/service/ec2 v1.332.0/go.mod h1:2o5yJcnWuaBOsnNqlO1reYs1OQffFkqf2xJ2mmMWNH4= +github.com/aws/aws-sdk-go-v2/service/eks v1.99.0 h1:DJhHBuKLyXU5krDg94uUwAbU4828vmurQd0O/VBanVk= +github.com/aws/aws-sdk-go-v2/service/eks v1.99.0/go.mod h1:flOCIr3poFqMcmCXl5ZZIhtDm0/n8tofP/Tio422uO0= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.19 h1:bAdDl/HkGCcGPoe25ToSHEw23VIxt6CT5fLcg111BKg= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.19/go.mod h1:KaUzbLxv4CeSxh6ZCl9B4m7CuFenS8kUEaDs+f/DQr4= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.14.3 h1:bON1rJf67TSTDCKg816AAIE4xSTtoo9tl0XRkO72R+I= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.14.3/go.mod h1:c5BBpjJcQXpfeq9iASyVKA3T6vX6B6LEXY4mL/gklDY= +github.com/aws/aws-sdk-go-v2/service/signin v1.10.0 h1:ZD5qFpWcaOKdTuhBi431pIDkCgrMkMlMT6jlpSPoIRI= +github.com/aws/aws-sdk-go-v2/service/signin v1.10.0/go.mod h1:8Nuuf+tR346PjJ3MvZPh9pekbLiLQFWJhzMXfwy7alA= +github.com/aws/aws-sdk-go-v2/service/sso v1.38.0 h1:JGeeBcMlhg1xtOXYpeCaTQBZObtXMPQCUqBcmr65NRA= +github.com/aws/aws-sdk-go-v2/service/sso v1.38.0/go.mod h1:XwteswG9EOMRFm73UT0t+MbTwyLxMrEXkU6e+v92Lzo= +github.com/aws/aws-sdk-go-v2/service/ssooidc v1.43.0 h1:obhahQXDEdVEv8y5bTKXR30LVaxYe1kyYM0L7l2Iq+k= +github.com/aws/aws-sdk-go-v2/service/ssooidc v1.43.0/go.mod h1:6twZZ/aXHNy1vXUO8koUbp++MYzMASkOgEBdkbJYmO0= +github.com/aws/aws-sdk-go-v2/service/sts v1.51.0 h1:Zpnqa6XtrNzXZnwbdCqHOXpXhMsa01ql/pcRQ1sb4hk= +github.com/aws/aws-sdk-go-v2/service/sts v1.51.0/go.mod h1:/8JRcdTt//hG0Q4BTmGbuOplT7ABe+5rdtqUHqXvYIM= +github.com/aws/smithy-go v1.28.1 h1:R/nXH00c8qcfCzQVELtRw+eLQWtzv+VAIEFJ1/xxXlQ= +github.com/aws/smithy-go v1.28.1/go.mod h1:YE2RhdIuDbA5E5bTdciG9KrW3+TiEONeUWCqxX9i1Fc= github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= github.com/blang/semver/v4 v4.0.0 h1:1PFHFE6yCCTv8C1TeyNNarDzntLi7wMI5i/pzqYIsAM= diff --git a/hack/gather-logs.sh b/hack/gather-logs.sh index a84102df..a580b12a 100755 --- a/hack/gather-logs.sh +++ b/hack/gather-logs.sh @@ -1,9 +1,9 @@ #!/bin/bash -# Gather diagnostic logs from a bink cluster. +# Gather diagnostic logs from a cluster. # # Usage: hack/gather-logs.sh [node-names...] # -# Expects KUBECONFIG and BINK_CLUSTER_NAME from environment. +# Expects KUBECONFIG from environment. # Each command's output is written to a separate file in . # Individual command failures are non-fatal. @@ -15,7 +15,6 @@ if [[ $# -lt 1 ]]; then fi : "${KUBECONFIG:?must be set}" -: "${BINK_CLUSTER_NAME:?must be set}" output_dir="$1" shift @@ -32,10 +31,6 @@ run() { "$@" > "${output_dir}/${filename}" 2>&1 || true } -# Host diagnostics -run "host-journal.txt" journalctl --no-pager -run "host-dmesg.txt" dmesg - # Cluster-wide commands run "k-get-pods.txt" kubectl get pods -n bootc-operator -o wide run "k-describe-pods.txt" kubectl describe pods -n bootc-operator @@ -50,10 +45,19 @@ for pod in $(kubectl get pods -n bootc-operator -o jsonpath='{.items[*].metadata run "k-logs-${pod}-previous.log" kubectl logs -n bootc-operator "${pod}" --all-containers --previous done -# Per-node commands +# Per-node commands via daemon pod exec for node in "${nodes[@]}"; do run "k-describe-node-${node}.txt" kubectl describe node "${node}" - run "journal-${node}.txt" bink node ssh "${node}" --cluster-name "${BINK_CLUSTER_NAME}" -- journalctl --no-pager + + daemon_pod=$(kubectl get pods -n bootc-operator \ + -l app.kubernetes.io/name=bootc-operator,app.kubernetes.io/component=daemon \ + --field-selector "spec.nodeName=${node}" \ + -o jsonpath='{.items[0].metadata.name}' 2>/dev/null || true) + + if [[ -n "${daemon_pod}" ]]; then + run "journal-${node}.txt" kubectl exec -n bootc-operator "${daemon_pod}" -- \ + journalctl --no-pager + fi done echo "Done." diff --git a/test/e2e/bootcnode_test.go b/test/e2e/bootcnode_test.go index fb5ab7b3..eb08180f 100644 --- a/test/e2e/bootcnode_test.go +++ b/test/e2e/bootcnode_test.go @@ -36,6 +36,7 @@ const ( // BootcNodePool selecting it, and verifies that a BootcNode is created // and the node is labeled bootc.dev/managed. func TestControllerMembership(t *testing.T) { + e2eutil.Providers(t, "bink") g := NewWithT(t) g.SetDefaultEventuallyTimeout(pollTimeout) g.SetDefaultEventuallyPollingInterval(pollInterval) @@ -110,6 +111,7 @@ func TestControllerMembership(t *testing.T) { // original image, then updates the pool to a new image and verifies the // full update lifecycle: staging, reboot, and idle with the new image. func TestUpdateReboot(t *testing.T) { + e2eutil.Providers(t, "bink") g := NewWithT(t) g.SetDefaultEventuallyTimeout(pollTimeout) g.SetDefaultEventuallyPollingInterval(pollInterval) @@ -304,6 +306,7 @@ func TestUpdateReboot(t *testing.T) { // the controller resolves the tag to a digest, then retags the image // and verifies re-resolution triggers a rollout. func TestTagResolution(t *testing.T) { + e2eutil.Providers(t, "bink") g := NewWithT(t) g.SetDefaultEventuallyTimeout(pollTimeout) g.SetDefaultEventuallyPollingInterval(pollInterval) @@ -411,6 +414,7 @@ func TestTagResolution(t *testing.T) { // image and that the non-rebooting node does not wastefully reboot into the // first update image. func TestMidRolloutImageChange(t *testing.T) { + e2eutil.Providers(t, "bink") g := NewWithT(t) g.SetDefaultEventuallyTimeout(pollTimeout) g.SetDefaultEventuallyPollingInterval(pollInterval) @@ -608,6 +612,7 @@ func getBootCount(t *testing.T, env *e2eutil.Env, ctx context.Context, nodeName // pool paused, verifies the node stages but does not reboot, then resumes // and verifies the update completes. func TestPauseResume(t *testing.T) { + e2eutil.Providers(t, "bink") g := NewWithT(t) g.SetDefaultEventuallyTimeout(pollTimeout) g.SetDefaultEventuallyPollingInterval(pollInterval) @@ -732,6 +737,7 @@ func TestPauseResume(t *testing.T) { // original image, then updates to a non-existing image and verifies the // node enters degraded state and the update does not proceed. func TestNonExistingImage(t *testing.T) { + e2eutil.Providers(t, "bink") g := NewWithT(t) g.SetDefaultEventuallyTimeout(pollTimeout) g.SetDefaultEventuallyPollingInterval(pollInterval) @@ -818,6 +824,7 @@ func TestNonExistingImage(t *testing.T) { // registry shares storage with the unauthenticated one (port 5000), // so the update image is already available at both endpoints. func TestPullSecretAuth(t *testing.T) { + e2eutil.Providers(t, "bink", "eks") g := NewWithT(t) g.SetDefaultEventuallyTimeout(pollTimeout) g.SetDefaultEventuallyPollingInterval(pollInterval) @@ -830,19 +837,20 @@ func TestPullSecretAuth(t *testing.T) { ctx := context.Background() nodeName := env.AddNode(t) - // The auth registry shares storage with the unauthenticated - // registry, so the update image pushed to localhost:5000 is - // already visible at auth-registry.cluster.local:5001. + // For bink, the auth registry (port 5001) shares storage with the + // unauthenticated registry (port 5000), so we rewrite the image ref + // to go through the auth endpoint. For EKS, the update image already + // requires authentication. + authImageRef := env.AuthImageRef() digest := env.NodeImageUpdateDigest() + registryHost := env.RegistryHost() - // Create a dockerconfigjson Secret with credentials for the - // in-cluster auth registry hostname. authStr := base64.StdEncoding.EncodeToString( []byte(env.RegistryUser() + ":" + env.RegistryPassword()), ) dockerCfg := fmt.Sprintf( - `{"auths":{"auth-registry.cluster.local:5001":{"auth":"%s"}}}`, - authStr, + `{"auths":{%q:{"auth":"%s"}}}`, + registryHost, authStr, ) secret := &corev1.Secret{ ObjectMeta: metav1.ObjectMeta{ @@ -857,8 +865,6 @@ func TestPullSecretAuth(t *testing.T) { g.Expect(env.Client.Create(ctx, secret)).To(Succeed()) t.Cleanup(func() { _ = env.Client.Delete(ctx, secret) }) - // Create a pool targeting the auth registry with the pull secret. - authImageRef := "auth-registry.cluster.local:5001/node@" + digest pool := env.NewPool("pullsecret", authImageRef, testutil.WithPullSecret(secret.Name, secret.Namespace), ) @@ -876,6 +882,9 @@ func TestPullSecretAuth(t *testing.T) { t.Logf("BootcNode %q has pullSecretRef set", nodeName) // Wait for the node to stage and reboot into the update image. + // Check Booted.Image (the full ref with manifest digest) rather + // than Booted.ImageDigest (the content digest) because they can + // differ with remote registries. g.Eventually(func() (bootcv1alpha1.BootcNodeStatus, error) { var bn bootcv1alpha1.BootcNode err := env.Client.Get(ctx, client.ObjectKey{Name: nodeName}, &bn) @@ -883,7 +892,7 @@ func TestPullSecretAuth(t *testing.T) { }).WithTimeout(5 * time.Minute).Should(And( HaveField("Booted", And( Not(BeNil()), - HaveField("ImageDigest", Equal(digest)), + HaveField("Image", Equal(authImageRef)), )), HaveField("Conditions", ContainElement(And( HaveField("Type", bootcv1alpha1.NodeIdle), diff --git a/test/e2e/crd_smoke_test.go b/test/e2e/crd_smoke_test.go index a902351c..fde50611 100644 --- a/test/e2e/crd_smoke_test.go +++ b/test/e2e/crd_smoke_test.go @@ -18,6 +18,7 @@ import ( // e2e-worthy flows once we have more of the controller and daemon implemented. // Note more comprehensive CRD round-trip tests exist in the unit tests. func TestCRDSmoke(t *testing.T) { + e2eutil.Providers(t, "bink") env := e2eutil.New(t) ctx := context.Background() diff --git a/test/e2e/e2eutil/env.go b/test/e2e/e2eutil/env.go index 6efebdb6..7492245c 100644 --- a/test/e2e/e2eutil/env.go +++ b/test/e2e/e2eutil/env.go @@ -1,15 +1,13 @@ // SPDX-License-Identifier: Apache-2.0 // Package e2eutil provides helpers for running end-to-end tests against -// a bink-managed Kubernetes cluster. The cluster and operator are -// expected to be already running (via `make deploy-bink`). Each test -// provisions its own worker nodes for isolation. +// a Kubernetes cluster. The cluster and operator are expected to be +// already running. Each test provisions its own worker nodes for +// isolation. package e2eutil import ( "context" - "crypto/rand" - "encoding/hex" "fmt" "os" "os/exec" @@ -45,44 +43,46 @@ type Env struct { // registered. Client client.Client - // clusterName is the bink cluster name. - clusterName string - // testID is the sanitized test name, used as the value for // LabelE2ETest on nodes and in pool selectors. testID string + // providerName identifies the active provider ("bink" or "eks"). + providerName string + + // provider handles node provisioning and removal. + provider NodeProvider + // nodes tracks node names added via AddNode for cleanup. nodes []string - // nodeImageDigest is the manifest digest of the bootc image seeded - // into the bink registry (e.g. "sha256:abc123..."). Empty when not seeded. - nodeImageDigest string + // nodeImageRef is the full digest-qualified reference for the base + // node image (e.g. "registry.example.com/node@sha256:abc123"). + nodeImageRef string - // nodeImageRegistry is the in-cluster registry path for the seeded node image - // (e.g. "registry.cluster.local:5000/node"). Empty when not seeded. - nodeImageRegistry string + // nodeImageUpdateRef is the full digest-qualified reference for the + // update image. May point to a different registry/image than the base. + nodeImageUpdateRef string - // nodeImageUpdateDigest is the manifest digest of the update image - // (e.g. "sha256:def456..."). Empty when not built. - nodeImageUpdateDigest string + // nodeImageUpdate2Ref is the full digest-qualified reference for the + // second update image. Empty when not needed. + nodeImageUpdate2Ref string - // nodeImageUpdate2Digest is the manifest digest of the second update - // image (e.g. "sha256:789abc..."). Used by mid-rollout image change tests. - nodeImageUpdate2Digest string + // registryHost is the hostname (with optional port) of the + // authenticated registry used by TestPullSecretAuth. + registryHost string - // registryUser is the username for the authenticated e2e registry - // on port 5001. Empty when not configured. + // registryUser is the username for the authenticated e2e registry. + // Empty when not configured. registryUser string // registryPassword is the password for the authenticated e2e - // registry on port 5001. Empty when not configured. + // registry. Empty when not configured. registryPassword string } -// New connects to an existing bink cluster and returns an Env ready -// for testing. The cluster must be running with the operator deployed -// (via `make deploy-bink`). +// New connects to an existing cluster and returns an Env ready for +// testing. The cluster must be running with the operator deployed. func New(t *testing.T) *Env { t.Helper() @@ -91,40 +91,69 @@ func New(t *testing.T) *Env { t.Fatal("KUBECONFIG must be set") } - clusterName := os.Getenv("BINK_CLUSTER_NAME") - if clusterName == "" { - t.Fatal("BINK_CLUSTER_NAME must be set") - } + k8sClient := buildClient(t, kubeconfigPath) - nodeImageDigest := os.Getenv("BINK_NODE_IMAGE_DIGEST") - if nodeImageDigest == "" { - t.Fatal("BINK_NODE_IMAGE_DIGEST must be set") - } - nodeImageRegistry := os.Getenv("BINK_LOCAL_REGISTRY_NODE_IMAGE") - if nodeImageRegistry == "" { - t.Fatal("BINK_LOCAL_REGISTRY_NODE_IMAGE must be set") - } - nodeImageUpdateDigest := os.Getenv("BINK_NODE_IMAGE_UPDATE_DIGEST") - if nodeImageUpdateDigest == "" { - t.Fatal("BINK_NODE_IMAGE_UPDATE_DIGEST must be set") - } - nodeImageUpdate2Digest := os.Getenv("BINK_NODE_IMAGE_UPDATE2_DIGEST") - if nodeImageUpdate2Digest == "" { - t.Fatal("BINK_NODE_IMAGE_UPDATE2_DIGEST must be set") + providerName := os.Getenv("E2E_PROVIDER") + if providerName == "" { + providerName = "bink" + } + + var ( + provider NodeProvider + nodeImageRef string + updateRef string + update2Ref string + registryHost string + ) + + switch providerName { + case "bink": + nodeImageRegistry := requireEnv(t, "E2E_NODE_IMAGE_REGISTRY") + nodeImageDigest := requireEnv(t, "E2E_NODE_IMAGE_DIGEST") + updateDigest := requireEnv(t, "E2E_NODE_IMAGE_UPDATE_DIGEST") + update2Digest := requireEnv(t, "E2E_NODE_IMAGE_UPDATE2_DIGEST") + + nodeImageRef = nodeImageRegistry + "@" + nodeImageDigest + updateRef = nodeImageRegistry + "@" + updateDigest + update2Ref = nodeImageRegistry + "@" + update2Digest + + registryHost = "auth-registry.cluster.local:5001" + + clusterName := requireEnv(t, "BINK_CLUSTER_NAME") + diskImage := os.Getenv("BINK_NODE_DISK_IMAGE") + provider = NewBinkProvider(clusterName, nodeImageRef, diskImage) + case "eks": + nodeImageRef = requireEnv(t, "E2E_NODE_IMAGE_REF") + updateRef = requireEnv(t, "E2E_NODE_IMAGE_UPDATE_REF") + update2Ref = os.Getenv("E2E_NODE_IMAGE_UPDATE2_REF") + registryHost = extractRegistryHost(updateRef) + + eksClusterName := requireEnv(t, "EKS_CLUSTER_NAME") + nodeGroup := requireEnv(t, "EKS_NODE_GROUP") + region := requireEnv(t, "AWS_REGION") + + var err error + provider, err = NewEKSProvider( + eksClusterName, nodeGroup, region, k8sClient, + ) + if err != nil { + t.Fatalf("creating EKS provider: %v", err) + } + default: + t.Fatalf("unknown E2E_PROVIDER %q (supported: bink, eks)", providerName) } - k8sClient := buildClient(t, kubeconfigPath) - env := &Env{ - Client: k8sClient, - clusterName: clusterName, - testID: sanitizeTestName(t.Name()), - nodeImageDigest: nodeImageDigest, - nodeImageRegistry: nodeImageRegistry, - nodeImageUpdateDigest: nodeImageUpdateDigest, - nodeImageUpdate2Digest: nodeImageUpdate2Digest, - registryUser: os.Getenv("E2E_REGISTRY_USER"), - registryPassword: os.Getenv("E2E_REGISTRY_PASSWORD"), + Client: k8sClient, + testID: sanitizeTestName(t.Name()), + providerName: providerName, + provider: provider, + nodeImageRef: nodeImageRef, + nodeImageUpdateRef: updateRef, + nodeImageUpdate2Ref: update2Ref, + registryHost: registryHost, + registryUser: os.Getenv("E2E_REGISTRY_USER"), + registryPassword: os.Getenv("E2E_REGISTRY_PASSWORD"), } t.Cleanup(func() { @@ -138,16 +167,7 @@ func New(t *testing.T) *Env { type NodeOption func(*nodeConfig) type nodeConfig struct { - memory int - labels map[string]string - targetImgRef string -} - -// WithMemory sets the VM memory in MB for the node. -func WithMemory(mb int) NodeOption { - return func(c *nodeConfig) { - c.memory = mb - } + labels map[string]string } // WithLabel adds a label to the provisioned node. This is in addition @@ -161,18 +181,9 @@ func WithLabel(key, value string) NodeOption { } } -// WithTargetImgRef sets the target image reference for the node, -// passed as --target-imgref to bink node add. Overrides the automatic -// default that AddNode applies when registry metadata is available. -func WithTargetImgRef(ref string) NodeOption { - return func(c *nodeConfig) { - c.targetImgRef = ref - } -} - -// AddNode provisions a worker node via bink, waits for it to be Ready, -// and returns the node name. The node is labeled with LabelE2ETest -// (and any extra labels from WithLabel). +// AddNode provisions a worker node via the configured provider, waits +// for it to be Ready, and returns the node name. The node is labeled +// with LabelE2ETest (and any extra labels from WithLabel). func (e *Env) AddNode(t *testing.T, opts ...NodeOption) string { t.Helper() @@ -181,38 +192,20 @@ func (e *Env) AddNode(t *testing.T, opts ...NodeOption) string { o(cfg) } - if cfg.targetImgRef == "" { - if e.nodeImageRegistry == "" || e.nodeImageDigest == "" { - t.Fatal( - "BINK_LOCAL_REGISTRY_NODE_IMAGE and NODE_IMAGE_DIGEST must be set (or use WithTargetImgRef)", - ) - } - cfg.targetImgRef = e.nodeImageRegistry + "@" + e.nodeImageDigest - } - - nodeName := e.generateNodeName(t) - - // Provision the node with labels applied at join time. - args := []string{"node", "add", nodeName, "--cluster-name", e.clusterName} - args = append(args, "--label", LabelE2ETest+"="+e.testID) + labels := map[string]string{LabelE2ETest: e.testID} for k, v := range cfg.labels { - args = append(args, "--label", k+"="+v) - } - if cfg.memory > 0 { - args = append(args, "--memory", fmt.Sprintf("%d", cfg.memory)) + labels[k] = v } - if img := os.Getenv("BINK_NODE_DISK_IMAGE"); img != "" { - args = append(args, "--node-image", img) - } - args = append(args, "--target-imgref", cfg.targetImgRef) - t.Logf("Adding node %q...", nodeName) - if err := runBink(t, args...); err != nil { - t.Fatalf("adding node %q: %v", nodeName, err) + + ctx := context.Background() + t.Logf("Adding node...") + nodeName, err := e.provider.AddNode(ctx, labels) + if err != nil { + t.Fatalf("adding node: %v", err) } + t.Logf("Added node %q", nodeName) e.nodes = append(e.nodes, nodeName) - - // Wait for Ready. waitForNodeReady(t, e.Client, nodeName) return nodeName @@ -246,52 +239,69 @@ func (e *Env) TestLabels() map[string]string { return map[string]string{LabelE2ETest: e.testID} } -// digestedPullSpec builds a digest-qualified image reference from the -// registry and the given digest. Returns "" if either is empty. -func (e *Env) digestedPullSpec(digest string) string { - if e.nodeImageRegistry == "" || digest == "" { - return "" - } - return e.nodeImageRegistry + "@" + digest -} - -// NodeImageDigestedPullSpec returns the digest-qualified reference for the -// seeded node image (e.g. "registry.cluster.local:5000/node@sha256:abc123"). +// NodeImageDigestedPullSpec returns the full digest-qualified reference +// for the base node image. func (e *Env) NodeImageDigestedPullSpec() string { - return e.digestedPullSpec(e.nodeImageDigest) + return e.nodeImageRef } -// NodeImageTagRef returns the tag-based reference for the seeded node +// NodeImageTagRef returns the tag-based reference for the base node // image (e.g. "registry.cluster.local:5000/node:latest"). func (e *Env) NodeImageTagRef() string { - return e.nodeImageRegistry + ":latest" + repo, _, _ := strings.Cut(e.nodeImageRef, "@") + return repo + ":latest" } -// NodeImageDigest returns the manifest digest of the seeded node image. +// NodeImageDigest returns the manifest digest portion of the base +// node image reference. func (e *Env) NodeImageDigest() string { - return e.nodeImageDigest + _, digest, _ := strings.Cut(e.nodeImageRef, "@") + return digest } -// NodeImageUpdateDigestedPullSpec returns the digest-qualified reference for the -// update image (e.g. "registry.cluster.local:5000/node@sha256:def456"). +// NodeImageUpdateDigestedPullSpec returns the full digest-qualified +// reference for the update image. func (e *Env) NodeImageUpdateDigestedPullSpec() string { - return e.digestedPullSpec(e.nodeImageUpdateDigest) + return e.nodeImageUpdateRef } -// NodeImageUpdateDigest returns the manifest digest of the update image. +// NodeImageUpdateDigest returns the manifest digest portion of the +// update image reference. func (e *Env) NodeImageUpdateDigest() string { - return e.nodeImageUpdateDigest + _, digest, _ := strings.Cut(e.nodeImageUpdateRef, "@") + return digest } -// NodeImageUpdate2DigestedPullSpec returns the digest-qualified reference for the -// second update image (e.g. "registry.cluster.local:5000/node@sha256:789abc"). +// NodeImageUpdate2DigestedPullSpec returns the full digest-qualified +// reference for the second update image. func (e *Env) NodeImageUpdate2DigestedPullSpec() string { - return e.digestedPullSpec(e.nodeImageUpdate2Digest) + return e.nodeImageUpdate2Ref } -// NodeImageUpdate2Digest returns the manifest digest of the second update image. +// NodeImageUpdate2Digest returns the manifest digest portion of the +// second update image reference. func (e *Env) NodeImageUpdate2Digest() string { - return e.nodeImageUpdate2Digest + _, digest, _ := strings.Cut(e.nodeImageUpdate2Ref, "@") + return digest +} + +// AuthImageRef returns the image reference to use for pull secret +// tests. For bink, this rewrites the update digest to go through the +// auth registry (port 5001). For EKS, the update image already +// requires auth, so it is returned as-is. +func (e *Env) AuthImageRef() string { + switch e.providerName { + case "bink": + return e.registryHost + "/node@" + e.NodeImageUpdateDigest() + default: + return e.nodeImageUpdateRef + } +} + +// RegistryHost returns the hostname (with optional port) of the +// authenticated registry used for pull secret tests. +func (e *Env) RegistryHost() string { + return e.registryHost } // RegistryUser returns the authenticated registry username, or empty @@ -334,7 +344,7 @@ func RetagImage(t *testing.T, srcRef, dstTag string) { } // cleanup gathers diagnostic logs, then deletes test-scoped resources -// and bink nodes. +// and removes nodes via the provider. func (e *Env) cleanup(t *testing.T) { e.gatherLogs(t) @@ -349,15 +359,7 @@ func (e *Env) cleanup(t *testing.T) { } for _, name := range e.nodes { t.Logf("Removing node %q...", name) - if err := runBink( - t, - "node", - "remove", - name, - "--force", - "--cluster-name", - e.clusterName, - ); err != nil { + if err := e.provider.RemoveNode(ctx, name); err != nil { t.Logf("WARNING: failed to remove node %q: %v", name, err) } } @@ -387,6 +389,23 @@ func (e *Env) gatherLogs(t *testing.T) { } } +func requireEnv(t *testing.T, key string) string { + t.Helper() + v := os.Getenv(key) + if v == "" { + t.Fatalf("%s must be set", key) + } + return v +} + +// extractRegistryHost returns the registry hostname from a full image +// reference like "registry.example.com/path/image@sha256:...". +func extractRegistryHost(ref string) string { + withoutDigest, _, _ := strings.Cut(ref, "@") + host, _, _ := strings.Cut(withoutDigest, "/") + return host +} + // sanitizeTestName lowercases a test name for use in k8s object names. // Panics if the result exceeds 63 characters (k8s label value limit). func sanitizeTestName(name string) string { @@ -397,17 +416,6 @@ func sanitizeTestName(name string) string { return name } -// generateNodeName creates a unique node name derived from the test name. -func (e *Env) generateNodeName(t *testing.T) string { - t.Helper() - - b := make([]byte, 3) - if _, err := rand.Read(b); err != nil { - t.Fatalf("generating random suffix: %v", err) - } - return e.testID + "-" + hex.EncodeToString(b) -} - // buildClient creates a controller-runtime client from the kubeconfig // with the bootc CRD scheme registered. func buildClient(t *testing.T, kubeconfigPath string) client.Client { @@ -447,13 +455,3 @@ func waitForNodeReady(t *testing.T, c client.Client, nodeName string) { t.Logf(" node %q is Ready", nodeName) }).WithTimeout(5 * time.Minute).WithPolling(5 * time.Second).Should(Succeed()) } - -// runBink executes a bink command and returns any error. -func runBink(t *testing.T, args ...string) error { - t.Helper() - - cmd := exec.Command("bink", args...) - cmd.Stdout = os.Stdout - cmd.Stderr = os.Stderr - return cmd.Run() -} diff --git a/test/e2e/e2eutil/provider.go b/test/e2e/e2eutil/provider.go new file mode 100644 index 00000000..3154c992 --- /dev/null +++ b/test/e2e/e2eutil/provider.go @@ -0,0 +1,33 @@ +// SPDX-License-Identifier: Apache-2.0 + +package e2eutil + +import ( + "context" + "os" + "testing" +) + +// NodeProvider abstracts node provisioning so the same tests can run +// against different infrastructure (bink, EKS, etc.). +type NodeProvider interface { + AddNode(ctx context.Context, labels map[string]string) (nodeName string, err error) + RemoveNode(ctx context.Context, nodeName string) error +} + +// Providers skips the test if the current E2E_PROVIDER is not in the +// supported list. Call at the top of each test function. +func Providers(t *testing.T, supported ...string) { + t.Helper() + + provider := os.Getenv("E2E_PROVIDER") + if provider == "" { + provider = "bink" + } + for _, s := range supported { + if s == provider { + return + } + } + t.Skipf("test requires provider %v, got %q", supported, provider) +} diff --git a/test/e2e/e2eutil/provider_bink.go b/test/e2e/e2eutil/provider_bink.go new file mode 100644 index 00000000..4cb36ed5 --- /dev/null +++ b/test/e2e/e2eutil/provider_bink.go @@ -0,0 +1,84 @@ +// SPDX-License-Identifier: Apache-2.0 + +package e2eutil + +import ( + "context" + "crypto/rand" + "encoding/hex" + "fmt" + "os" + "os/exec" +) + +// BinkProvider provisions nodes using bink (KVM-based local clusters). +type BinkProvider struct { + clusterName string + targetImgRef string + diskImage string +} + +// NewBinkProvider creates a BinkProvider for the given bink cluster. +func NewBinkProvider( + clusterName, targetImgRef, diskImage string, +) *BinkProvider { + return &BinkProvider{ + clusterName: clusterName, + targetImgRef: targetImgRef, + diskImage: diskImage, + } +} + +func (p *BinkProvider) AddNode( + ctx context.Context, + labels map[string]string, +) (string, error) { + prefix := labels[LabelE2ETest] + if prefix == "" { + prefix = "node" + } + + b := make([]byte, 3) + if _, err := rand.Read(b); err != nil { + return "", fmt.Errorf("generating random suffix: %w", err) + } + nodeName := prefix + "-" + hex.EncodeToString(b) + + args := []string{ + "node", "add", nodeName, + "--cluster-name", p.clusterName, + } + for k, v := range labels { + args = append(args, "--label", k+"="+v) + } + if p.diskImage != "" { + args = append(args, "--node-image", p.diskImage) + } + args = append(args, "--target-imgref", p.targetImgRef) + + cmd := exec.CommandContext(ctx, "bink", args...) + cmd.Stdout = os.Stdout + cmd.Stderr = os.Stderr + if err := cmd.Run(); err != nil { + return "", fmt.Errorf("bink node add %q: %w", nodeName, err) + } + return nodeName, nil +} + +func (p *BinkProvider) RemoveNode( + ctx context.Context, + nodeName string, +) error { + cmd := exec.CommandContext( + ctx, + "bink", "node", "remove", nodeName, + "--force", + "--cluster-name", p.clusterName, + ) + cmd.Stdout = os.Stdout + cmd.Stderr = os.Stderr + if err := cmd.Run(); err != nil { + return fmt.Errorf("bink node remove %q: %w", nodeName, err) + } + return nil +} diff --git a/test/e2e/e2eutil/provider_eks.go b/test/e2e/e2eutil/provider_eks.go new file mode 100644 index 00000000..3ff05f35 --- /dev/null +++ b/test/e2e/e2eutil/provider_eks.go @@ -0,0 +1,392 @@ +// SPDX-License-Identifier: Apache-2.0 + +package e2eutil + +import ( + "context" + "fmt" + "log/slog" + "strings" + "time" + + "github.com/aws/aws-sdk-go-v2/aws" + awsconfig "github.com/aws/aws-sdk-go-v2/config" + "github.com/aws/aws-sdk-go-v2/service/autoscaling" + autoscalingtypes "github.com/aws/aws-sdk-go-v2/service/autoscaling/types" + "github.com/aws/aws-sdk-go-v2/service/ec2" + ec2types "github.com/aws/aws-sdk-go-v2/service/ec2/types" + corev1 "k8s.io/api/core/v1" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +type ec2RunInstancesAPI interface { + RunInstances( + ctx context.Context, + params *ec2.RunInstancesInput, + optFns ...func(*ec2.Options), + ) (*ec2.RunInstancesOutput, error) +} + +type ec2TerminateInstancesAPI interface { + TerminateInstances( + ctx context.Context, + params *ec2.TerminateInstancesInput, + optFns ...func(*ec2.Options), + ) (*ec2.TerminateInstancesOutput, error) +} + +type ec2DescribeLaunchTemplateVersionsAPI interface { + DescribeLaunchTemplateVersions( + ctx context.Context, + params *ec2.DescribeLaunchTemplateVersionsInput, + optFns ...func(*ec2.Options), + ) (*ec2.DescribeLaunchTemplateVersionsOutput, error) +} + +// EKSProvider provisions nodes by launching standalone EC2 instances +// that join an existing EKS cluster using the node group's launch +// template. +type EKSProvider struct { + ec2Runner ec2RunInstancesAPI + ec2Terminator ec2TerminateInstancesAPI + k8s client.Client + ltID string + ltVersion string + subnetID string + securityGroups []string + // nodes maps K8s node names to EC2 instance IDs so that + // RemoveNode can terminate instances even if the K8s node + // object is already gone. + nodes map[string]string +} + +// NewEKSProvider creates an EKSProvider by discovering the Auto Scaling +// group for the given eksctl node group (via tags) and extracting its +// launch template. +func NewEKSProvider( + clusterName, nodeGroup, region string, + k8sClient client.Client, +) (*EKSProvider, error) { + ctx := context.Background() + cfg, err := awsconfig.LoadDefaultConfig(ctx, awsconfig.WithRegion(region)) + if err != nil { + return nil, fmt.Errorf("loading AWS config: %w", err) + } + + asgClient := autoscaling.NewFromConfig(cfg) + ec2Client := ec2.NewFromConfig(cfg) + + asg, err := findASGByTags(ctx, asgClient, clusterName, nodeGroup) + if err != nil { + return nil, err + } + + asgName := aws.ToString(asg.AutoScalingGroupName) + if len(asg.AvailabilityZones) == 0 { + return nil, fmt.Errorf("ASG %q has no availability zones", asgName) + } + + subnetID, err := subnetFromASG(asg) + if err != nil { + return nil, err + } + + lt, err := discoverLaunchTemplate(ctx, ec2Client, asg) + if err != nil { + return nil, err + } + + return &EKSProvider{ + ec2Runner: ec2Client, + ec2Terminator: ec2Client, + k8s: k8sClient, + ltID: lt.id, + ltVersion: lt.version, + subnetID: subnetID, + securityGroups: lt.securityGroups, + nodes: make(map[string]string), + }, nil +} + +func (p *EKSProvider) AddNode( + ctx context.Context, + labels map[string]string, +) (string, error) { + out, err := p.ec2Runner.RunInstances(ctx, &ec2.RunInstancesInput{ + LaunchTemplate: &ec2types.LaunchTemplateSpecification{ + LaunchTemplateId: &p.ltID, + Version: &p.ltVersion, + }, + NetworkInterfaces: []ec2types.InstanceNetworkInterfaceSpecification{ + { + DeviceIndex: aws.Int32(0), + SubnetId: &p.subnetID, + AssociatePublicIpAddress: aws.Bool(true), + Groups: p.securityGroups, + }, + }, + MinCount: aws.Int32(1), + MaxCount: aws.Int32(1), + }) + if err != nil { + return "", fmt.Errorf("launching instance: %w", err) + } + if len(out.Instances) == 0 { + return "", fmt.Errorf("RunInstances returned no instances") + } + instanceID := aws.ToString(out.Instances[0].InstanceId) + + nodeName, err := p.waitForNode(ctx, instanceID) + if err != nil { + p.terminateInstance(ctx, instanceID) + return "", err + } + + if err := p.patchLabels(ctx, nodeName, labels); err != nil { + p.terminateInstance(ctx, instanceID) + return "", err + } + + p.nodes[nodeName] = instanceID + return nodeName, nil +} + +func (p *EKSProvider) RemoveNode( + ctx context.Context, + nodeName string, +) error { + instanceID, ok := p.nodes[nodeName] + if !ok { + node := &corev1.Node{} + if err := p.k8s.Get(ctx, client.ObjectKey{Name: nodeName}, node); err != nil { + return fmt.Errorf("getting node %q: %w", nodeName, err) + } + var err error + instanceID, err = instanceIDFromProviderID(node.Spec.ProviderID) + if err != nil { + return err + } + } + + _, err := p.ec2Terminator.TerminateInstances(ctx, &ec2.TerminateInstancesInput{ + InstanceIds: []string{instanceID}, + }) + if err != nil { + return fmt.Errorf("terminating instance %q: %w", instanceID, err) + } + + delete(p.nodes, nodeName) + + node := &corev1.Node{} + if err := p.k8s.Get(ctx, client.ObjectKey{Name: nodeName}, node); err != nil { + return nil + } + if err := p.k8s.Delete(ctx, node); err != nil { + return fmt.Errorf("deleting k8s node %q: %w", nodeName, err) + } + return nil +} + +func (p *EKSProvider) terminateInstance(ctx context.Context, instanceID string) { + _, err := p.ec2Terminator.TerminateInstances(ctx, &ec2.TerminateInstancesInput{ + InstanceIds: []string{instanceID}, + }) + if err != nil { + slog.Error("failed to terminate leaked EC2 instance", + "instanceID", instanceID, "error", err) + } +} + +func (p *EKSProvider) waitForNode( + ctx context.Context, + instanceID string, +) (string, error) { + var lastErr error + for attempt := 0; attempt < 120; attempt++ { + var nodeList corev1.NodeList + if err := p.k8s.List(ctx, &nodeList); err != nil { + lastErr = err + } else { + for j := range nodeList.Items { + if strings.HasSuffix(nodeList.Items[j].Spec.ProviderID, "/"+instanceID) { + return nodeList.Items[j].Name, nil + } + } + } + select { + case <-ctx.Done(): + return "", ctx.Err() + case <-time.After(5 * time.Second): + } + } + if lastErr != nil { + return "", fmt.Errorf( + "timed out waiting for k8s node with instance ID %s (last error: %w)", + instanceID, lastErr, + ) + } + return "", fmt.Errorf( + "timed out waiting for k8s node with instance ID %s", instanceID, + ) +} + +func (p *EKSProvider) patchLabels( + ctx context.Context, + nodeName string, + labels map[string]string, +) error { + node := &corev1.Node{} + if err := p.k8s.Get(ctx, client.ObjectKey{Name: nodeName}, node); err != nil { + return fmt.Errorf("getting node %q for label patch: %w", nodeName, err) + } + + modified := node.DeepCopy() + if modified.Labels == nil { + modified.Labels = make(map[string]string) + } + for k, v := range labels { + modified.Labels[k] = v + } + + if err := p.k8s.Patch(ctx, modified, client.MergeFrom(node)); err != nil { + return fmt.Errorf("patching labels on node %q: %w", nodeName, err) + } + return nil +} + +// instanceIDFromProviderID extracts the EC2 instance ID from a k8s +// node's providerID (format: "aws:///ZONE/INSTANCE-ID"). +func instanceIDFromProviderID(providerID string) (string, error) { + parts := strings.Split(providerID, "/") + if len(parts) < 2 { + return "", fmt.Errorf( + "unexpected providerID format: %q", providerID, + ) + } + id := parts[len(parts)-1] + if !strings.HasPrefix(id, "i-") { + return "", fmt.Errorf( + "providerID %q does not contain an EC2 instance ID", providerID, + ) + } + return id, nil +} + +func findASGByTags( + ctx context.Context, + asgClient *autoscaling.Client, + clusterName, nodeGroup string, +) (autoscalingtypes.AutoScalingGroup, error) { + out, err := asgClient.DescribeAutoScalingGroups( + ctx, + &autoscaling.DescribeAutoScalingGroupsInput{ + Filters: []autoscalingtypes.Filter{ + { + Name: aws.String("tag:eksctl.cluster.k8s.io/v1alpha1/cluster-name"), + Values: []string{clusterName}, + }, + { + Name: aws.String("tag:eksctl.io/v1alpha2/nodegroup-name"), + Values: []string{nodeGroup}, + }, + }, + }, + ) + if err != nil { + return autoscalingtypes.AutoScalingGroup{}, fmt.Errorf( + "finding ASG for cluster %q node group %q: %w", + clusterName, nodeGroup, err, + ) + } + if len(out.AutoScalingGroups) == 0 { + return autoscalingtypes.AutoScalingGroup{}, fmt.Errorf( + "no ASG found for cluster %q node group %q", + clusterName, nodeGroup, + ) + } + return out.AutoScalingGroups[0], nil +} + +func subnetFromASG(asg autoscalingtypes.AutoScalingGroup) (string, error) { + vpc := aws.ToString(asg.VPCZoneIdentifier) + if vpc == "" { + return "", fmt.Errorf( + "ASG %q has no VPCZoneIdentifier", + aws.ToString(asg.AutoScalingGroupName), + ) + } + subnet, _, _ := strings.Cut(vpc, ",") + return subnet, nil +} + +type launchTemplateInfo struct { + id string + version string + securityGroups []string +} + +func discoverLaunchTemplate( + ctx context.Context, + ec2Client ec2DescribeLaunchTemplateVersionsAPI, + asg autoscalingtypes.AutoScalingGroup, +) (launchTemplateInfo, error) { + asgName := aws.ToString(asg.AutoScalingGroupName) + + var ltID, ltVersion string + switch { + case asg.LaunchTemplate != nil: + ltID = aws.ToString(asg.LaunchTemplate.LaunchTemplateId) + ltVersion = aws.ToString(asg.LaunchTemplate.Version) + case asg.MixedInstancesPolicy != nil && + asg.MixedInstancesPolicy.LaunchTemplate != nil && + asg.MixedInstancesPolicy.LaunchTemplate.LaunchTemplateSpecification != nil: + spec := asg.MixedInstancesPolicy.LaunchTemplate.LaunchTemplateSpecification + ltID = aws.ToString(spec.LaunchTemplateId) + ltVersion = aws.ToString(spec.Version) + default: + return launchTemplateInfo{}, fmt.Errorf("ASG %q has no launch template", asgName) + } + + resolved, err := ec2Client.DescribeLaunchTemplateVersions( + ctx, + &ec2.DescribeLaunchTemplateVersionsInput{ + LaunchTemplateId: <ID, + Versions: []string{ltVersion}, + }, + ) + if err != nil { + return launchTemplateInfo{}, fmt.Errorf( + "describing launch template %q version %q: %w", + ltID, ltVersion, err, + ) + } + if len(resolved.LaunchTemplateVersions) == 0 { + return launchTemplateInfo{}, fmt.Errorf( + "launch template %q version %q not found", ltID, ltVersion, + ) + } + v := resolved.LaunchTemplateVersions[0].VersionNumber + if v == nil { + return launchTemplateInfo{}, fmt.Errorf("launch template version number is nil") + } + + lt := launchTemplateInfo{ + id: ltID, + version: fmt.Sprintf("%d", *v), + } + + ltData := resolved.LaunchTemplateVersions[0].LaunchTemplateData + if ltData != nil { + for _, ni := range ltData.NetworkInterfaces { + if ni.DeviceIndex != nil && *ni.DeviceIndex == 0 && len(ni.Groups) > 0 { + lt.securityGroups = ni.Groups + break + } + } + if len(lt.securityGroups) == 0 { + lt.securityGroups = ltData.SecurityGroupIds + } + } + + return lt, nil +} diff --git a/test/e2e/e2eutil/provider_eks_test.go b/test/e2e/e2eutil/provider_eks_test.go new file mode 100644 index 00000000..47526802 --- /dev/null +++ b/test/e2e/e2eutil/provider_eks_test.go @@ -0,0 +1,201 @@ +// SPDX-License-Identifier: Apache-2.0 + +package e2eutil + +import ( + "context" + "testing" + + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/aws/aws-sdk-go-v2/service/ec2" + ec2types "github.com/aws/aws-sdk-go-v2/service/ec2/types" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" +) + +type mockEC2Runner struct { + instanceID string + lastInput *ec2.RunInstancesInput +} + +func (m *mockEC2Runner) RunInstances( + _ context.Context, + input *ec2.RunInstancesInput, + _ ...func(*ec2.Options), +) (*ec2.RunInstancesOutput, error) { + m.lastInput = input + return &ec2.RunInstancesOutput{ + Instances: []ec2types.Instance{ + {InstanceId: aws.String(m.instanceID)}, + }, + }, nil +} + +type mockEC2Terminator struct { + terminated []string +} + +func (m *mockEC2Terminator) TerminateInstances( + _ context.Context, + input *ec2.TerminateInstancesInput, + _ ...func(*ec2.Options), +) (*ec2.TerminateInstancesOutput, error) { + m.terminated = append(m.terminated, input.InstanceIds...) + return &ec2.TerminateInstancesOutput{}, nil +} + +func fakeK8sClient(objs ...client.Object) client.Client { + scheme := runtime.NewScheme() + _ = corev1.AddToScheme(scheme) + + return fake.NewClientBuilder(). + WithScheme(scheme). + WithObjects(objs...). + Build() +} + +func TestEKSAddNode(t *testing.T) { + const ( + instanceID = "i-0123456789abcdef0" + nodeName = "ip-10-0-1-100.ec2.internal" + providerID = "aws:///us-east-1a/" + instanceID + ) + + node := &corev1.Node{ + ObjectMeta: metav1.ObjectMeta{Name: nodeName}, + Spec: corev1.NodeSpec{ProviderID: providerID}, + } + + runner := &mockEC2Runner{instanceID: instanceID} + provider := &EKSProvider{ + ec2Runner: runner, + ec2Terminator: &mockEC2Terminator{}, + k8s: fakeK8sClient(node), + ltID: "lt-012345", + ltVersion: "1", + securityGroups: []string{"sg-aaa", "sg-bbb"}, + } + + labels := map[string]string{LabelE2ETest: "testid"} + got, err := provider.AddNode(context.Background(), labels) + if err != nil { + t.Fatalf("AddNode: %v", err) + } + if got != nodeName { + t.Fatalf("expected node name %q, got %q", nodeName, got) + } + + var updated corev1.Node + if err := provider.k8s.Get( + context.Background(), + client.ObjectKey{Name: nodeName}, + &updated, + ); err != nil { + t.Fatalf("getting node: %v", err) + } + if v := updated.Labels[LabelE2ETest]; v != "testid" { + t.Fatalf("expected label %s=testid, got %q", LabelE2ETest, v) + } + + if len(runner.lastInput.NetworkInterfaces) == 0 { + t.Fatal("expected NetworkInterfaces in RunInstances input") + } + gotGroups := runner.lastInput.NetworkInterfaces[0].Groups + if len(gotGroups) != 2 || gotGroups[0] != "sg-aaa" || gotGroups[1] != "sg-bbb" { + t.Fatalf("expected security groups [sg-aaa sg-bbb], got %v", gotGroups) + } +} + +func TestEKSRemoveNode(t *testing.T) { + const ( + instanceID = "i-abcdef0123456789a" + nodeName = "ip-10-0-2-50.ec2.internal" + providerID = "aws:///us-west-2b/" + instanceID + ) + + node := &corev1.Node{ + ObjectMeta: metav1.ObjectMeta{Name: nodeName}, + Spec: corev1.NodeSpec{ProviderID: providerID}, + } + + terminator := &mockEC2Terminator{} + provider := &EKSProvider{ + ec2Runner: &mockEC2Runner{}, + ec2Terminator: terminator, + k8s: fakeK8sClient(node), + ltID: "lt-012345", + ltVersion: "1", + } + + if err := provider.RemoveNode(context.Background(), nodeName); err != nil { + t.Fatalf("RemoveNode: %v", err) + } + if len(terminator.terminated) != 1 || terminator.terminated[0] != instanceID { + t.Fatalf( + "expected TerminateInstances(%q), got %v", + instanceID, terminator.terminated, + ) + } + + // Verify the k8s Node object was deleted. + var deleted corev1.Node + err := provider.k8s.Get( + context.Background(), + client.ObjectKey{Name: nodeName}, + &deleted, + ) + if err == nil { + t.Fatal("expected node to be deleted, but it still exists") + } +} + +func TestInstanceIDFromProviderID(t *testing.T) { + tests := []struct { + name string + providerID string + wantID string + wantErr bool + }{ + { + name: "standard format", + providerID: "aws:///us-east-1a/i-0123456789abcdef0", + wantID: "i-0123456789abcdef0", + }, + { + name: "different zone", + providerID: "aws:///eu-west-1c/i-abcdef0123456789a", + wantID: "i-abcdef0123456789a", + }, + { + name: "no instance id prefix", + providerID: "aws:///us-east-1a/not-an-instance", + wantErr: true, + }, + { + name: "empty", + providerID: "", + wantErr: true, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got, err := instanceIDFromProviderID(tt.providerID) + if tt.wantErr { + if err == nil { + t.Fatal("expected error, got nil") + } + return + } + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if got != tt.wantID { + t.Fatalf("expected %q, got %q", tt.wantID, got) + } + }) + } +}