From 8f7f306ebfe25bdf6bd25e556fa4b9aacb5fe5e6 Mon Sep 17 00:00:00 2001 From: Nitin Kumar Date: Sun, 27 Sep 2026 17:44:13 +0530 Subject: [PATCH] fix(aws-eks): nodegroup scaling defaults and apply version updates when they settle --- providers/aws/eks/async_settle_test.go | 12 +- providers/aws/eks/driver/driver.go | 26 ++- providers/aws/eks/eks.go | 197 +++++++++++++++---- providers/aws/eks/eks_test.go | 12 +- providers/aws/eks/nodegroup_nodes_test.go | 16 +- providers/aws/eks/nodegroup_scaling_test.go | 136 +++++++++++++ providers/aws/eks/snapshot.go | 13 +- providers/aws/eks/version_settle_test.go | 161 +++++++++++++++ server/aws/eks/operations.go | 76 +++---- server/aws/eks/sdk_nodegroup_scaling_test.go | 100 ++++++++++ services/kubernetes/snapshot.go | 2 +- services/kubernetes/snapshot_test.go | 1 + services/kubernetes/state.go | 3 + services/kubernetes/version.go | 58 +++++- services/kubernetes/version_test.go | 52 +++++ 15 files changed, 756 insertions(+), 109 deletions(-) create mode 100644 providers/aws/eks/nodegroup_scaling_test.go create mode 100644 providers/aws/eks/version_settle_test.go create mode 100644 server/aws/eks/sdk_nodegroup_scaling_test.go diff --git a/providers/aws/eks/async_settle_test.go b/providers/aws/eks/async_settle_test.go index f600c3314..7d1134f3c 100644 --- a/providers/aws/eks/async_settle_test.go +++ b/providers/aws/eks/async_settle_test.go @@ -96,8 +96,8 @@ func TestAsyncSettleClusterCreateUpdate(t *testing.T) { t.Fatalf("update cluster version: %v", err) } - if updated.Status != "Successful" { - t.Fatalf("update record status = %q, want Successful", updated.Status) + if updated.Status != updateInProgress { + t.Fatalf("update record status = %q, want InProgress", updated.Status) } if got := clusterStatus(t, m, "c1"); got != eksdriver.ClusterStatusUpdating { @@ -133,7 +133,7 @@ func TestAsyncSettleNodegroupCreateUpdate(t *testing.T) { created, err := m.CreateNodegroup(ctx, eksdriver.NodegroupConfig{ ClusterName: "c1", NodegroupName: "ng1", - ScalingConfig: eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 3, DesiredSize: 2}, + ScalingConfig: &eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 3, DesiredSize: 2}, }) if err != nil { t.Fatalf("create nodegroup: %v", err) @@ -158,15 +158,15 @@ func TestAsyncSettleNodegroupCreateUpdate(t *testing.T) { t.Fatalf("settled status = %q, want %q", got, eksdriver.NodegroupStatusActive) } - scaling := eksdriver.NodegroupScalingConfig{MinSize: 2, MaxSize: 5, DesiredSize: 4} + scaling := eksdriver.NodegroupScalingUpdate{MinSize: intPtr(2), MaxSize: intPtr(5), DesiredSize: intPtr(4)} upd, err := m.UpdateNodegroupConfig(ctx, "c1", "ng1", eksdriver.NodegroupConfigUpdate{Scaling: &scaling}) if err != nil { t.Fatalf("update nodegroup config: %v", err) } - if upd.Status != "Successful" { - t.Fatalf("update record status = %q, want Successful", upd.Status) + if upd.Status != updateInProgress { + t.Fatalf("update record status = %q, want InProgress", upd.Status) } if got := nodegroupStatus(t, m, "c1", "ng1"); got != eksdriver.NodegroupStatusUpdating { diff --git a/providers/aws/eks/driver/driver.go b/providers/aws/eks/driver/driver.go index 1363c2444..834c6d44f 100644 --- a/providers/aws/eks/driver/driver.go +++ b/providers/aws/eks/driver/driver.go @@ -183,6 +183,14 @@ type NodegroupScalingConfig struct { DesiredSize int } +// NodegroupScalingUpdate is a partial scaling change for UpdateNodegroupConfig. +// Only the non-nil sizes change; the rest keep their current values. +type NodegroupScalingUpdate struct { + MinSize *int + MaxSize *int + DesiredSize *int +} + // Taint is a Kubernetes taint applied to a managed node group's nodes. Effect // is one of NO_SCHEDULE, PREFER_NO_SCHEDULE, or NO_EXECUTE. A taint is // identified by its Key+Effect pair. @@ -226,11 +234,13 @@ type NodegroupConfig struct { DiskSize int Version string ReleaseVersion string - ScalingConfig NodegroupScalingConfig - UpdateConfig NodegroupUpdateConfig - Labels map[string]string - Taints []Taint - Tags map[string]string + // ScalingConfig is optional; nil gets the EKS default of min 1, max 2, + // desired 2. + ScalingConfig *NodegroupScalingConfig + UpdateConfig NodegroupUpdateConfig + Labels map[string]string + Taints []Taint + Tags map[string]string // LaunchTemplate is optional; when set, it names the EC2 launch template // backing the node group's instances. LaunchTemplate *LaunchTemplateSpecification @@ -273,11 +283,11 @@ type Nodegroup struct { } // NodegroupConfigUpdate carries the mutable fields UpdateNodegroupConfig -// applies. Scaling, when non-nil, is the already-merged target sizing (the -// caller overlays partial requests). Label and taint changes are expressed as +// applies. Scaling, when non-nil, names the sizes to change; they are merged +// onto the current config and the result is validated. Label and taint changes are expressed as // add/update and remove deltas, matching the real EKS request shape. type NodegroupConfigUpdate struct { - Scaling *NodegroupScalingConfig + Scaling *NodegroupScalingUpdate UpdateConfig *NodegroupUpdateConfig AddOrUpdateLabels map[string]string RemoveLabels []string diff --git a/providers/aws/eks/eks.go b/providers/aws/eks/eks.go index 2b822ddb4..d764d4aac 100644 --- a/providers/aws/eks/eks.go +++ b/providers/aws/eks/eks.go @@ -108,8 +108,27 @@ type Mock struct { // nodegroupSettle is the nodegroup analog of clusterSettle, keyed by // nodegroupKey(clusterName, nodegroupName). nodegroupSettle *settle.Set + // updateSettle keeps an update record InProgress (keyed by update ID) + // until the cluster or nodegroup change it tracks has settled. + updateSettle *settle.Set + // clusterOldVersions and nodegroupOldVersions hold the version a cluster + // or nodegroup keeps reporting while a version update settles. The stored + // row already carries the new version. Guarded by mu. + clusterOldVersions map[string]oldVersion + nodegroupOldVersions map[string]oldVersion } +// oldVersion is what a cluster or nodegroup reports until window elapses. +type oldVersion struct { + version string + releaseVersion string + window settle.Window +} + +// updateStatusInProgress is the status of an update whose change is still +// settling; the stored record carries the final Successful status. +const updateStatusInProgress = "InProgress" + // New creates a new AWS EKS mock. func New(opts *config.Options) *Mock { return &Mock{ @@ -123,6 +142,10 @@ func New(opts *config.Options) *Mock { k8sUIDs: make(map[string]string), clusterSettle: settle.NewSet(), nodegroupSettle: settle.NewSet(), + updateSettle: settle.NewSet(), + + clusterOldVersions: make(map[string]oldVersion), + nodegroupOldVersions: make(map[string]oldVersion), } } @@ -207,6 +230,18 @@ func (m *Mock) recordUpdate(u *eksdriver.ClusterUpdate) *eksdriver.ClusterUpdate return &out } +// recordSettlingUpdate records an update that stays InProgress for d, the +// settle window of the change it tracks. With d = 0 it is recordUpdate. +// Callers must hold m.mu. +func (m *Mock) recordSettlingUpdate(u *eksdriver.ClusterUpdate, d time.Duration) *eksdriver.ClusterUpdate { + out := m.recordUpdate(u) + + m.updateSettle.Begin(u.ID, updateStatusInProgress, u.CreatedAt, d) + out.Status = m.updateSettle.State(u.ID, m.opts.Clock.Now(), out.Status) + + return out +} + func (m *Mock) clusterARN(region, name string) string { return idgen.AWSARN("eks", region, m.opts.AccountID, "cluster/"+name) } @@ -553,20 +588,67 @@ func mergeLabels(cur, addOrUpdate map[string]string, remove []string) map[string return out } -// validateScaling rejects an inconsistent scaling config the way real EKS does -// (InvalidParameterException): minSize must not exceed maxSize and desiredSize -// must fall within [minSize, maxSize]. An all-zero config (scaling omitted) is -// valid and left to defaults. -func validateScaling(s eksdriver.NodegroupScalingConfig) error { - if s.MinSize > s.MaxSize { - return cerrors.Newf(cerrors.InvalidArgument, - "minSize (%d) must not be greater than maxSize (%d)", s.MinSize, s.MaxSize) +// Nodegroup scaling defaults and limits. A CreateNodegroup without a +// scalingConfig gets min 1, max 2, desired 2. maxSize is capped at the default +// "Nodes per managed node group" service quota. +const ( + defaultNodegroupMinSize = 1 + defaultNodegroupMaxSize = 2 + defaultNodegroupDesiredSize = 2 + maxNodesPerNodegroup = 450 +) + +// resolveScaling returns the scaling config a new nodegroup starts with: the +// requested one, or the EKS default when the request omitted it. +func resolveScaling(s *eksdriver.NodegroupScalingConfig) eksdriver.NodegroupScalingConfig { + if s == nil { + return eksdriver.NodegroupScalingConfig{ + MinSize: defaultNodegroupMinSize, + MaxSize: defaultNodegroupMaxSize, + DesiredSize: defaultNodegroupDesiredSize, + } + } + + return *s +} + +// mergeScaling applies the sizes named in upd onto cur. +func mergeScaling(cur eksdriver.NodegroupScalingConfig, upd *eksdriver.NodegroupScalingUpdate) eksdriver.NodegroupScalingConfig { + if upd.MinSize != nil { + cur.MinSize = *upd.MinSize + } + + if upd.MaxSize != nil { + cur.MaxSize = *upd.MaxSize } - if s.DesiredSize < s.MinSize || s.DesiredSize > s.MaxSize { + if upd.DesiredSize != nil { + cur.DesiredSize = *upd.DesiredSize + } + + return cur +} + +// validateScaling rejects a scaling config the way real EKS does +// (InvalidParameterException): minSize and desiredSize are at least 0, maxSize +// is between 1 and the node quota, and minSize <= desiredSize <= maxSize. The +// ordering messages are the ones EKS returns. +func validateScaling(s eksdriver.NodegroupScalingConfig) error { + switch { + case s.MinSize < 0: + return cerrors.New(cerrors.InvalidArgument, "minSize must be greater than or equal to 0") + case s.DesiredSize < 0: + return cerrors.New(cerrors.InvalidArgument, "desiredSize must be greater than or equal to 0") + case s.MaxSize < 1: + return cerrors.New(cerrors.InvalidArgument, "maxSize must be greater than or equal to 1") + case s.MaxSize > maxNodesPerNodegroup: + return cerrors.Newf(cerrors.InvalidArgument, "maxSize can't be greater than %d", maxNodesPerNodegroup) + case s.MinSize > s.DesiredSize: + return cerrors.Newf(cerrors.InvalidArgument, + "Minimum capacity %d can't be greater than desired size %d", s.MinSize, s.DesiredSize) + case s.DesiredSize > s.MaxSize: return cerrors.Newf(cerrors.InvalidArgument, - "desiredSize (%d) must be between minSize (%d) and maxSize (%d)", - s.DesiredSize, s.MinSize, s.MaxSize) + "desired capacity %d can't be greater than max size %d", s.DesiredSize, s.MaxSize) } return nil @@ -784,10 +866,23 @@ func (m *Mock) clusterStatusLocked(c *eksdriver.Cluster) string { return m.clusterSettle.State(c.Name, m.opts.Clock.Now(), c.Status) } -// overlayClusterStatus mutates c.Status in place to the settle-overlaid value. +// overlayClusterStatus mutates c.Status in place to the settle-overlaid value, +// and puts back the old version while a version update is still settling. // Caller must hold at least an RLock on m.mu. func (m *Mock) overlayClusterStatus(c *eksdriver.Cluster) { c.Status = m.clusterStatusLocked(c) + c.Version = m.clusterVersionLocked(c) +} + +// clusterVersionLocked returns the version c's control plane runs right now: +// the old one while an UpdateClusterVersion is settling, else the stored one. +// Caller must hold at least an RLock on m.mu. +func (m *Mock) clusterVersionLocked(c *eksdriver.Cluster) string { + if old, ok := m.clusterOldVersions[c.Name]; ok && !old.window.Settled(m.opts.Clock.Now()) { + return old.version + } + + return c.Version } // nodegroupStatusLocked is the nodegroup analog of clusterStatusLocked. @@ -798,6 +893,12 @@ func (m *Mock) nodegroupStatusLocked(ng *eksdriver.Nodegroup) string { // overlayNodegroupStatus is the nodegroup analog of overlayClusterStatus. func (m *Mock) overlayNodegroupStatus(ng *eksdriver.Nodegroup) { ng.Status = m.nodegroupStatusLocked(ng) + + key := nodegroupKey(ng.ClusterName, ng.NodegroupName) + if old, ok := m.nodegroupOldVersions[key]; ok && !old.window.Settled(m.opts.Clock.Now()) { + ng.Version = old.version + ng.ReleaseVersion = old.releaseVersion + } } // DescribeCluster looks up a cluster by name. @@ -928,16 +1029,16 @@ func (m *Mock) UpdateClusterConfig( m.clusters.Set(name, c) - m.clusterSettle.Begin(name, eksdriver.ClusterStatusUpdating, m.opts.Clock.Now(), - m.opts.SettleDuration(settle.DefaultClusterSettle)) + d := m.opts.SettleDuration(settle.DefaultClusterSettle) + m.clusterSettle.Begin(name, eksdriver.ClusterStatusUpdating, m.opts.Clock.Now(), d) - return m.recordUpdate(&eksdriver.ClusterUpdate{ + return m.recordSettlingUpdate(&eksdriver.ClusterUpdate{ ID: newUpdateID(), Type: clusterConfigUpdateType(accessConfigChanged, loggingChanged, vpcEndpointChanged, vpcOtherChanged), Status: "Successful", CreatedAt: m.opts.Clock.Now().UTC(), ClusterName: name, - }), nil + }, d), nil } // UpdateClusterVersion moves a cluster up one minor version. A one-minor @@ -987,24 +1088,30 @@ func (m *Mock) UpdateClusterVersion( c.VersionUpgradedAt = m.opts.Clock.Now().UTC() } + now := m.opts.Clock.Now() + d := m.opts.SettleDuration(settle.DefaultClusterSettle) + + // The stored row takes the new version. Until the update settles the + // control plane still runs the old one, so describes and /version keep + // reporting it and switch when the window elapses. With AsyncSettle off + // d is 0 and the switch is immediate. + m.clusterOldVersions[name] = oldVersion{version: c.Version, window: settle.Pending("", now, d)} c.Version = version m.clusters.Set(name, c) - // The API server now reports the new version on /version. if uid, ok := m.k8sUIDs[name]; ok && m.k8sAPI != nil { - m.k8sAPI.SetClusterVersion(uid, kubernetes.DistributionEKS, serverPatchVersion(version)) + m.k8sAPI.SetClusterVersionAt(uid, kubernetes.DistributionEKS, serverPatchVersion(version), m.opts.Clock, now.Add(d)) } - m.clusterSettle.Begin(name, eksdriver.ClusterStatusUpdating, m.opts.Clock.Now(), - m.opts.SettleDuration(settle.DefaultClusterSettle)) + m.clusterSettle.Begin(name, eksdriver.ClusterStatusUpdating, now, d) - return m.recordUpdate(&eksdriver.ClusterUpdate{ + return m.recordSettlingUpdate(&eksdriver.ClusterUpdate{ ID: newUpdateID(), Type: updateType, Status: "Successful", - CreatedAt: m.opts.Clock.Now().UTC(), + CreatedAt: now.UTC(), ClusterName: name, - }), nil + }, d), nil } // DeleteCluster removes a cluster (only if no nodegroups, Fargate profiles, @@ -1045,6 +1152,7 @@ func (m *Mock) DeleteCluster(_ context.Context, name string) (*eksdriver.Cluster m.clusters.Delete(name) m.clusterSettle.Clear(name) + delete(m.clusterOldVersions, name) m.deleteClusterAccessEntriesLocked(name) // Resolve the endpoint before deregistering: the response describes the @@ -1076,6 +1184,7 @@ func (m *Mock) DescribeUpdate(_ context.Context, clusterName, updateID string) ( } out := u + out.Status = m.updateSettle.State(u.ID, m.opts.Clock.Now(), u.Status) return &out, nil } @@ -1142,7 +1251,8 @@ func (m *Mock) CreateNodegroup(_ context.Context, cfg eksdriver.NodegroupConfig) "nodegroup %q already exists in cluster %q", cfg.NodegroupName, cfg.ClusterName) } - if err := validateScaling(cfg.ScalingConfig); err != nil { + scaling := resolveScaling(cfg.ScalingConfig) + if err := validateScaling(scaling); err != nil { return nil, err } @@ -1171,7 +1281,7 @@ func (m *Mock) CreateNodegroup(_ context.Context, cfg eksdriver.NodegroupConfig) DiskSize: diskSize, Version: version, ReleaseVersion: cfg.ReleaseVersion, - ScalingConfig: cfg.ScalingConfig, + ScalingConfig: scaling, UpdateConfig: resolveUpdateConfig(cfg.UpdateConfig), Status: eksdriver.NodegroupStatusActive, Labels: copyTags(cfg.Labels), @@ -1264,11 +1374,12 @@ func (m *Mock) UpdateNodegroupConfig( } if upd.Scaling != nil { - if err := validateScaling(*upd.Scaling); err != nil { + scaling := mergeScaling(ng.ScalingConfig, upd.Scaling) + if err := validateScaling(scaling); err != nil { return nil, err } - ng.ScalingConfig = *upd.Scaling + ng.ScalingConfig = scaling } if upd.UpdateConfig != nil { @@ -1286,17 +1397,17 @@ func (m *Mock) UpdateNodegroupConfig( m.nodegroups.Set(key, ng) m.syncNodegroupNodesLocked(&ng, ng.ScalingConfig.DesiredSize) - m.nodegroupSettle.Begin(key, eksdriver.NodegroupStatusUpdating, m.opts.Clock.Now(), - m.opts.SettleDuration(settle.DefaultClusterSettle)) + d := m.opts.SettleDuration(settle.DefaultClusterSettle) + m.nodegroupSettle.Begin(key, eksdriver.NodegroupStatusUpdating, m.opts.Clock.Now(), d) - return m.recordUpdate(&eksdriver.ClusterUpdate{ + return m.recordSettlingUpdate(&eksdriver.ClusterUpdate{ ID: newUpdateID(), Type: "ConfigUpdate", Status: "Successful", CreatedAt: m.opts.Clock.Now().UTC(), ClusterName: clusterName, NodegroupName: nodegroupName, - }), nil + }, d), nil } // UpdateNodegroupVersion bumps the Kubernetes version of a nodegroup. @@ -1327,6 +1438,10 @@ func (m *Mock) UpdateNodegroupVersion( return nil, cerrors.Newf(cerrors.NotFound, "cluster %q not found", clusterName) } + // Resolve against the version the control plane runs now, which lags the + // stored one while an UpdateClusterVersion is settling. + parent.Version = m.clusterVersionLocked(&parent) + resolved, err := resolveNodegroupVersion(&parent, version) if err != nil { return nil, err @@ -1336,6 +1451,12 @@ func (m *Mock) UpdateNodegroupVersion( return nil, err } + now := m.opts.Clock.Now() + d := m.opts.SettleDuration(settle.DefaultClusterSettle) + m.nodegroupOldVersions[key] = oldVersion{ + version: ng.Version, releaseVersion: ng.ReleaseVersion, window: settle.Pending("", now, d), + } + ng.Version = resolved if releaseVersion != "" { @@ -1347,17 +1468,16 @@ func (m *Mock) UpdateNodegroupVersion( m.nodegroups.Set(key, ng) m.syncNodegroupNodesLocked(&ng, ng.ScalingConfig.DesiredSize) - m.nodegroupSettle.Begin(key, eksdriver.NodegroupStatusUpdating, m.opts.Clock.Now(), - m.opts.SettleDuration(settle.DefaultClusterSettle)) + m.nodegroupSettle.Begin(key, eksdriver.NodegroupStatusUpdating, now, d) - return m.recordUpdate(&eksdriver.ClusterUpdate{ + return m.recordSettlingUpdate(&eksdriver.ClusterUpdate{ ID: newUpdateID(), Type: "VersionUpdate", Status: "Successful", - CreatedAt: m.opts.Clock.Now().UTC(), + CreatedAt: now.UTC(), ClusterName: clusterName, NodegroupName: nodegroupName, - }), nil + }, d), nil } // DeleteNodegroup removes a nodegroup. @@ -1377,6 +1497,7 @@ func (m *Mock) DeleteNodegroup(_ context.Context, clusterName, nodegroupName str m.nodegroups.Delete(key) m.nodegroupSettle.Clear(key) + delete(m.nodegroupOldVersions, key) m.removeNodeEntryLocked(clusterName, ng.NodeRole) m.syncNodegroupNodesLocked(&ng, 0) @@ -1523,7 +1644,7 @@ func (m *Mock) CreateAddon(_ context.Context, cfg eksdriver.AddonConfig) (*eksdr "add-on %q already installed on cluster %q", cfg.AddonName, cfg.ClusterName) } - version, err := resolveAddonVersion(cfg.AddonName, parent.Version, cfg.AddonVersion) + version, err := resolveAddonVersion(cfg.AddonName, m.clusterVersionLocked(&parent), cfg.AddonVersion) if err != nil { return nil, err } @@ -1612,7 +1733,7 @@ func (m *Mock) UpdateAddon(_ context.Context, cfg eksdriver.AddonConfig) (*eksdr return nil, cerrors.Newf(cerrors.NotFound, "cluster %q not found", cfg.ClusterName) } - version, err := resolveAddonVersion(cfg.AddonName, parent.Version, cfg.AddonVersion) + version, err := resolveAddonVersion(cfg.AddonName, m.clusterVersionLocked(&parent), cfg.AddonVersion) if err != nil { return nil, err } diff --git a/providers/aws/eks/eks_test.go b/providers/aws/eks/eks_test.go index 2273bfff3..55f2c968e 100644 --- a/providers/aws/eks/eks_test.go +++ b/providers/aws/eks/eks_test.go @@ -143,7 +143,7 @@ func TestNodegroupLifecycle(t *testing.T) { NodegroupName: "ng1", NodeRole: "arn:aws:iam::123456789012:role/eks-node", Subnets: []string{"subnet-1"}, - ScalingConfig: eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 3, DesiredSize: 2}, + ScalingConfig: &eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 3, DesiredSize: 2}, }) requireNoError(t, err) assertEqual(t, "ACTIVE", ng.Status) @@ -158,7 +158,7 @@ func TestNodegroupLifecycle(t *testing.T) { assertEqual(t, 1, len(names)) upd, err := m.UpdateNodegroupConfig(ctx, "c1", "ng1", eksdriver.NodegroupConfigUpdate{ - Scaling: &eksdriver.NodegroupScalingConfig{MinSize: 2, MaxSize: 5, DesiredSize: 3}, + Scaling: &eksdriver.NodegroupScalingUpdate{MinSize: intPtr(2), MaxSize: intPtr(5), DesiredSize: intPtr(3)}, }) requireNoError(t, err) assertEqual(t, "Successful", upd.Status) @@ -199,7 +199,7 @@ func TestNodegroupTaintsAndModifiedAt(t *testing.T) { NodegroupName: "ng1", NodeRole: "arn:aws:iam::123456789012:role/eks-node", Subnets: []string{"subnet-1"}, - ScalingConfig: eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 3, DesiredSize: 2}, + ScalingConfig: &eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 3, DesiredSize: 2}, Taints: []eksdriver.Taint{{Key: "dedicated", Value: "gpu", Effect: "NO_SCHEDULE"}}, }) requireNoError(t, err) @@ -240,7 +240,7 @@ func TestNodegroupLabelMergeAndRemove(t *testing.T) { NodegroupName: "ng1", NodeRole: "arn:aws:iam::123456789012:role/eks-node", Subnets: []string{"subnet-1"}, - ScalingConfig: eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 2, DesiredSize: 1}, + ScalingConfig: &eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 2, DesiredSize: 1}, Labels: map[string]string{"a": "1", "b": "2"}, }) requireNoError(t, err) @@ -271,7 +271,7 @@ func TestNodegroupScalingValidation(t *testing.T) { NodegroupName: "bad", NodeRole: "arn:aws:iam::123456789012:role/eks-node", Subnets: []string{"subnet-1"}, - ScalingConfig: eksdriver.NodegroupScalingConfig{MinSize: 5, MaxSize: 2, DesiredSize: 1}, + ScalingConfig: &eksdriver.NodegroupScalingConfig{MinSize: 5, MaxSize: 2, DesiredSize: 1}, }) if err == nil { t.Fatal("expected InvalidArgument for minSize > maxSize") @@ -290,7 +290,7 @@ func TestCreateNodegroup_Defaults(t *testing.T) { NodegroupName: "ng1", NodeRole: "arn:aws:iam::123456789012:role/eks-node", Subnets: []string{"subnet-1"}, - ScalingConfig: eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 2, DesiredSize: 1}, + ScalingConfig: &eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 2, DesiredSize: 1}, }) requireNoError(t, err) assertEqual(t, 20, ng.DiskSize) diff --git a/providers/aws/eks/nodegroup_nodes_test.go b/providers/aws/eks/nodegroup_nodes_test.go index 80b343ab5..8afa0c5e4 100644 --- a/providers/aws/eks/nodegroup_nodes_test.go +++ b/providers/aws/eks/nodegroup_nodes_test.go @@ -105,7 +105,7 @@ func TestNodegroupNodesCreatedWithEKSShape(t *testing.T) { InstanceTypes: []string{"m5.large"}, AmiType: "AL2023_ARM_64_STANDARD", CapacityType: "SPOT", - ScalingConfig: eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 3, DesiredSize: 2}, + ScalingConfig: &eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 3, DesiredSize: 2}, Labels: map[string]string{"team": "payments"}, Taints: []eksdriver.Taint{{Key: "dedicated", Value: "batch", Effect: "NO_SCHEDULE"}}, }) @@ -192,7 +192,7 @@ func TestNodegroupNodesUSEast1Hostname(t *testing.T) { _, err := m.CreateNodegroup(context.Background(), eksdriver.NodegroupConfig{ ClusterName: "c1", NodegroupName: "ng", - ScalingConfig: eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 1, DesiredSize: 1}, + ScalingConfig: &eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 1, DesiredSize: 1}, }) requireNoError(t, err) @@ -222,7 +222,7 @@ func TestNodegroupNodesScaleUpdateAndDelete(t *testing.T) { _, err := m.CreateNodegroup(ctx, eksdriver.NodegroupConfig{ ClusterName: "c1", NodegroupName: "ng", - ScalingConfig: eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 3, DesiredSize: 2}, + ScalingConfig: &eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 3, DesiredSize: 2}, Labels: map[string]string{"old": "1"}, Taints: []eksdriver.Taint{{Key: "k", Value: "v", Effect: "NO_EXECUTE"}}, }) @@ -239,7 +239,7 @@ func TestNodegroupNodesScaleUpdateAndDelete(t *testing.T) { t.Helper() _, err := m.UpdateNodegroupConfig(ctx, "c1", "ng", eksdriver.NodegroupConfigUpdate{ - Scaling: &eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 3, DesiredSize: n}, + Scaling: &eksdriver.NodegroupScalingUpdate{MinSize: intPtr(1), MaxSize: intPtr(3), DesiredSize: intPtr(n)}, }) requireNoError(t, err) } @@ -303,7 +303,7 @@ func TestNodegroupNodesVersionUpdateRefreshesKubelet(t *testing.T) { _, err := m.CreateNodegroup(ctx, eksdriver.NodegroupConfig{ ClusterName: "c1", NodegroupName: "ng", Version: "1.29", - ScalingConfig: eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 1, DesiredSize: 1}, + ScalingConfig: &eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 1, DesiredSize: 1}, }) requireNoError(t, err) @@ -331,7 +331,7 @@ func TestNodegroupNodesSurviveSnapshotRestore(t *testing.T) { _, err := m.CreateNodegroup(ctx, eksdriver.NodegroupConfig{ ClusterName: "c1", NodegroupName: "ng", - ScalingConfig: eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 3, DesiredSize: 2}, + ScalingConfig: &eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 3, DesiredSize: 2}, }) requireNoError(t, err) @@ -359,7 +359,7 @@ func TestNodegroupNodesSurviveSnapshotRestore(t *testing.T) { assertEqual(t, 2, len(before)) _, err = m2.UpdateNodegroupConfig(ctx, "c1", "ng", eksdriver.NodegroupConfigUpdate{ - Scaling: &eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 3, DesiredSize: 3}, + Scaling: &eksdriver.NodegroupScalingUpdate{MinSize: intPtr(1), MaxSize: intPtr(3), DesiredSize: intPtr(3)}, }) requireNoError(t, err) @@ -404,7 +404,7 @@ func TestNodegroupNodesWithoutDataPlane(t *testing.T) { _, err := m.CreateNodegroup(context.Background(), eksdriver.NodegroupConfig{ ClusterName: "c1", NodegroupName: "ng", - ScalingConfig: eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 2, DesiredSize: 2}, + ScalingConfig: &eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 2, DesiredSize: 2}, }) requireNoError(t, err) diff --git a/providers/aws/eks/nodegroup_scaling_test.go b/providers/aws/eks/nodegroup_scaling_test.go new file mode 100644 index 000000000..54346e179 --- /dev/null +++ b/providers/aws/eks/nodegroup_scaling_test.go @@ -0,0 +1,136 @@ +package eks + +import ( + "context" + "testing" + + eksdriver "github.com/stackshy/cloudemu/v2/providers/aws/eks/driver" +) + +func intPtr(n int) *int { return &n } + +// TestCreateNodegroupScalingDefault checks that a nodegroup created without a +// scalingConfig gets the EKS defaults (min 1, max 2, desired 2) and that the +// data plane runs two Nodes for it. +func TestCreateNodegroupScalingDefault(t *testing.T) { + m := newTestMock() + base := nodesFixture(t, m) + + created, err := m.CreateNodegroup(context.Background(), eksdriver.NodegroupConfig{ + ClusterName: "c1", NodegroupName: "ng1", Subnets: []string{"subnet-a"}, + }) + requireNoError(t, err) + + want := eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 2, DesiredSize: 2} + if created.ScalingConfig != want { + t.Fatalf("scalingConfig = %+v, want %+v", created.ScalingConfig, want) + } + + if n := len(listNodes(t, base)); n != 2 { + t.Fatalf("nodes = %d, want 2", n) + } +} + +func TestCreateNodegroupScalingValidation(t *testing.T) { + tests := []struct { + name string + scaling eksdriver.NodegroupScalingConfig + want string + }{ + {"min above desired", eksdriver.NodegroupScalingConfig{MinSize: 2, MaxSize: 3, DesiredSize: 1}, + "Minimum capacity 2 can't be greater than desired size 1"}, + {"desired above max", eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 1, DesiredSize: 2}, + "desired capacity 2 can't be greater than max size 1"}, + {"min above max", eksdriver.NodegroupScalingConfig{MinSize: 5, MaxSize: 2, DesiredSize: 1}, + "Minimum capacity 5 can't be greater than desired size 1"}, + {"explicit zero max", eksdriver.NodegroupScalingConfig{}, + "maxSize must be greater than or equal to 1"}, + {"negative min", eksdriver.NodegroupScalingConfig{MinSize: -1, MaxSize: 2, DesiredSize: 1}, + "minSize must be greater than or equal to 0"}, + {"negative desired", eksdriver.NodegroupScalingConfig{MinSize: 0, MaxSize: 2, DesiredSize: -1}, + "desiredSize must be greater than or equal to 0"}, + {"max above quota", eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 451, DesiredSize: 1}, + "maxSize can't be greater than 450"}, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + m := newTestMock() + mustCluster(t, m, "c1") + + scaling := tc.scaling + _, err := m.CreateNodegroup(context.Background(), eksdriver.NodegroupConfig{ + ClusterName: "c1", NodegroupName: "ng1", ScalingConfig: &scaling, + }) + requireInvalidArg(t, err, tc.want) + }) + } +} + +// TestCreateNodegroupScalingBounds checks the edges that must be accepted: +// a scale-to-zero group and one at the node quota. +func TestCreateNodegroupScalingBounds(t *testing.T) { + for _, s := range []eksdriver.NodegroupScalingConfig{ + {MinSize: 0, MaxSize: 1, DesiredSize: 0}, + {MinSize: 0, MaxSize: 450, DesiredSize: 0}, + } { + m := newTestMock() + mustCluster(t, m, "c1") + + scaling := s + got, err := m.CreateNodegroup(context.Background(), eksdriver.NodegroupConfig{ + ClusterName: "c1", NodegroupName: "ng1", ScalingConfig: &scaling, + }) + requireNoError(t, err) + + if got.ScalingConfig != s { + t.Fatalf("scalingConfig = %+v, want %+v", got.ScalingConfig, s) + } + } +} + +// TestUpdateNodegroupConfigPartialScaling checks that an update merges only +// the sizes it names onto the current config, and validates the merged result. +func TestUpdateNodegroupConfigPartialScaling(t *testing.T) { + m := newTestMock() + ctx := context.Background() + mustCluster(t, m, "c1") + + _, err := m.CreateNodegroup(ctx, eksdriver.NodegroupConfig{ + ClusterName: "c1", NodegroupName: "ng1", + ScalingConfig: &eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 3, DesiredSize: 2}, + }) + requireNoError(t, err) + + _, err = m.UpdateNodegroupConfig(ctx, "c1", "ng1", eksdriver.NodegroupConfigUpdate{ + Scaling: &eksdriver.NodegroupScalingUpdate{DesiredSize: intPtr(3)}, + }) + requireNoError(t, err) + + got, err := m.DescribeNodegroup(ctx, "c1", "ng1") + requireNoError(t, err) + + want := eksdriver.NodegroupScalingConfig{MinSize: 1, MaxSize: 3, DesiredSize: 3} + if got.ScalingConfig != want { + t.Fatalf("scalingConfig = %+v, want %+v", got.ScalingConfig, want) + } + + // Lowering maxSize below the current desiredSize is rejected on the + // merged config, and nothing changes. + _, err = m.UpdateNodegroupConfig(ctx, "c1", "ng1", eksdriver.NodegroupConfigUpdate{ + Scaling: &eksdriver.NodegroupScalingUpdate{MaxSize: intPtr(2)}, + }) + requireInvalidArg(t, err, "desired capacity 3 can't be greater than max size 2") + + _, err = m.UpdateNodegroupConfig(ctx, "c1", "ng1", eksdriver.NodegroupConfigUpdate{ + Scaling: &eksdriver.NodegroupScalingUpdate{MinSize: intPtr(4)}, + }) + requireInvalidArg(t, err, "Minimum capacity 4 can't be greater than desired size 3") + + got, err = m.DescribeNodegroup(ctx, "c1", "ng1") + requireNoError(t, err) + + if got.ScalingConfig != want { + t.Fatalf("after rejected updates scalingConfig = %+v, want %+v", got.ScalingConfig, want) + } +} diff --git a/providers/aws/eks/snapshot.go b/providers/aws/eks/snapshot.go index 329de60d5..91b515153 100644 --- a/providers/aws/eks/snapshot.go +++ b/providers/aws/eks/snapshot.go @@ -5,6 +5,7 @@ import ( "encoding/json" "fmt" + "github.com/stackshy/cloudemu/v2/internal/settle" "github.com/stackshy/cloudemu/v2/internal/snapshot" ) @@ -87,12 +88,20 @@ func (m *Mock) Restore(_ context.Context, data json.RawMessage) error { return err } + m.mu.Lock() + if snap.K8sUIDs != nil { - m.mu.Lock() m.k8sUIDs = snap.K8sUIDs - m.mu.Unlock() } + // Settle windows are not persisted, so a restored cluster or nodegroup + // reports the version it was moving to and its updates read Successful. + m.updateSettle = settle.NewSet() + m.clusterOldVersions = make(map[string]oldVersion) + m.nodegroupOldVersions = make(map[string]oldVersion) + + m.mu.Unlock() + return nil } diff --git a/providers/aws/eks/version_settle_test.go b/providers/aws/eks/version_settle_test.go new file mode 100644 index 000000000..4b655fad1 --- /dev/null +++ b/providers/aws/eks/version_settle_test.go @@ -0,0 +1,161 @@ +package eks + +import ( + "context" + "testing" + "time" + + "github.com/stackshy/cloudemu/v2/internal/settle" + eksdriver "github.com/stackshy/cloudemu/v2/providers/aws/eks/driver" + "github.com/stackshy/cloudemu/v2/services/kubernetes" +) + +const ( + updateInProgress = "InProgress" + updateSuccessful = "Successful" +) + +func updateStatus(t *testing.T, m *Mock, cluster, id string) string { + t.Helper() + + u, err := m.DescribeUpdate(context.Background(), cluster, id) + requireNoError(t, err) + + return u.Status +} + +// TestAsyncSettleClusterVersionAppliesOnCompletion checks that under +// AsyncSettle an UpdateClusterVersion keeps reporting the old version, in +// DescribeCluster and on the data plane's /version, until the update settles. +func TestAsyncSettleClusterVersionAppliesOnCompletion(t *testing.T) { + m, fc := newAsyncMock() + ctx := context.Background() + api := kubernetes.NewAPIServer() + m.SetK8sAPI(api) + + _, err := m.CreateCluster(ctx, eksdriver.ClusterConfig{Name: "c1", Version: "1.31"}) + requireNoError(t, err) + fc.Advance(settle.DefaultClusterSettle) + + uid := m.k8sUIDs["c1"] + + upd, err := m.UpdateClusterVersion(ctx, "c1", toVersion("1.32")) + requireNoError(t, err) + + if upd.Status != updateInProgress { + t.Fatalf("update response status = %q, want %q", upd.Status, updateInProgress) + } + + fc.Advance(settle.DefaultClusterSettle - time.Millisecond) + + got, err := m.DescribeCluster(ctx, "c1") + requireNoError(t, err) + + if got.Version != "1.31" || got.Status != eksdriver.ClusterStatusUpdating { + t.Fatalf("mid-update describe = %s/%s, want 1.31/UPDATING", got.Version, got.Status) + } + + if _, minor := dataPlaneVersion(t, api, uid); minor != "31+" { + t.Fatalf("mid-update /version minor = %q, want 31+", minor) + } + + if s := updateStatus(t, m, "c1", upd.ID); s != updateInProgress { + t.Fatalf("mid-update DescribeUpdate status = %q, want %q", s, updateInProgress) + } + + fc.Advance(time.Millisecond) + + got, err = m.DescribeCluster(ctx, "c1") + requireNoError(t, err) + + if got.Version != "1.32" || got.Status != eksdriver.ClusterStatusActive { + t.Fatalf("settled describe = %s/%s, want 1.32/ACTIVE", got.Version, got.Status) + } + + if _, minor := dataPlaneVersion(t, api, uid); minor != "32+" { + t.Fatalf("settled /version minor = %q, want 32+", minor) + } + + if s := updateStatus(t, m, "c1", upd.ID); s != updateSuccessful { + t.Fatalf("settled DescribeUpdate status = %q, want %q", s, updateSuccessful) + } +} + +// TestAsyncSettleNodegroupVersionAppliesOnCompletion is the nodegroup analog: +// the nodegroup keeps its old version until UpdateNodegroupVersion settles, and +// while the control plane upgrade is still running a nodegroup can't move to +// the version the cluster hasn't reached yet. +func TestAsyncSettleNodegroupVersionAppliesOnCompletion(t *testing.T) { + m, fc := newAsyncMock() + ctx := context.Background() + + _, err := m.CreateCluster(ctx, eksdriver.ClusterConfig{Name: "c1", Version: "1.31"}) + requireNoError(t, err) + fc.Advance(settle.DefaultClusterSettle) + + _, err = m.CreateNodegroup(ctx, eksdriver.NodegroupConfig{ClusterName: "c1", NodegroupName: "ng1"}) + requireNoError(t, err) + fc.Advance(settle.DefaultClusterSettle) + + _, err = m.UpdateClusterVersion(ctx, "c1", toVersion("1.32")) + requireNoError(t, err) + + _, err = m.UpdateNodegroupVersion(ctx, "c1", "ng1", eksdriver.NodegroupVersionUpdate{Version: "1.32"}) + if err == nil { + t.Fatal("expected UpdateNodegroupVersion to 1.32 to fail while the control plane is still on 1.31") + } + + fc.Advance(settle.DefaultClusterSettle) + + upd, err := m.UpdateNodegroupVersion(ctx, "c1", "ng1", eksdriver.NodegroupVersionUpdate{Version: "1.32"}) + requireNoError(t, err) + + if upd.Status != updateInProgress { + t.Fatalf("update response status = %q, want %q", upd.Status, updateInProgress) + } + + ng, err := m.DescribeNodegroup(ctx, "c1", "ng1") + requireNoError(t, err) + + if ng.Version != "1.31" || ng.Status != eksdriver.NodegroupStatusUpdating { + t.Fatalf("mid-update nodegroup = %s/%s, want 1.31/UPDATING", ng.Version, ng.Status) + } + + if s := updateStatus(t, m, "c1", upd.ID); s != updateInProgress { + t.Fatalf("mid-update DescribeUpdate status = %q, want %q", s, updateInProgress) + } + + fc.Advance(settle.DefaultClusterSettle) + + ng, err = m.DescribeNodegroup(ctx, "c1", "ng1") + requireNoError(t, err) + + if ng.Version != "1.32" || ng.Status != eksdriver.NodegroupStatusActive { + t.Fatalf("settled nodegroup = %s/%s, want 1.32/ACTIVE", ng.Version, ng.Status) + } + + if s := updateStatus(t, m, "c1", upd.ID); s != updateSuccessful { + t.Fatalf("settled DescribeUpdate status = %q, want %q", s, updateSuccessful) + } +} + +// TestSyncClusterVersionImmediate checks the default path: with AsyncSettle +// off the new version and a Successful update are visible straight away. +func TestSyncClusterVersionImmediate(t *testing.T) { + m, _ := newVersionMock(t, "1.31") + ctx := context.Background() + + upd, err := m.UpdateClusterVersion(ctx, "c1", toVersion("1.32")) + requireNoError(t, err) + + if upd.Status != updateSuccessful { + t.Fatalf("update status = %q, want %q", upd.Status, updateSuccessful) + } + + got, err := m.DescribeCluster(ctx, "c1") + requireNoError(t, err) + + if got.Version != "1.32" { + t.Fatalf("version = %s, want 1.32", got.Version) + } +} diff --git a/server/aws/eks/operations.go b/server/aws/eks/operations.go index 3cc2b6f4d..001c759fd 100644 --- a/server/aws/eks/operations.go +++ b/server/aws/eks/operations.go @@ -305,7 +305,16 @@ func (h *Handler) createNodegroup(w http.ResponseWriter, r *http.Request, cluste } if body.ScalingConfig != nil { - cfg.ScalingConfig = scalingFromJSON(body.ScalingConfig) + sc, ok := scalingFromJSON(body.ScalingConfig) + if !ok { + // EKS takes all three sizes on create, or none of them. + writeError(w, http.StatusBadRequest, "InvalidParameterException", + "scalingConfig must specify all of minSize, maxSize and desiredSize, or none of them") + + return + } + + cfg.ScalingConfig = sc } if body.UpdateConfig != nil { @@ -362,19 +371,7 @@ func (h *Handler) updateNodegroupConfig(w http.ResponseWriter, r *http.Request, update := eksdriver.NodegroupConfigUpdate{} if body.ScalingConfig != nil { - // Real EKS applies only the sizes present in the request; the driver - // replaces ScalingConfig wholesale, so merge onto the current config - // here to avoid zeroing MinSize/MaxSize/DesiredSize that were omitted. - cur, err := h.eks.DescribeNodegroup(r.Context(), clusterName, ngName) - if err != nil { - writeErr(w, err) - - return - } - - merged := cur.ScalingConfig - mergeScaling(&merged, body.ScalingConfig) - update.Scaling = &merged + update.Scaling = scalingUpdateFromJSON(body.ScalingConfig) } if body.UpdateConfig != nil { @@ -669,20 +666,24 @@ func launchTemplateToJSON(l *eksdriver.LaunchTemplateSpecification) *launchTempl } } -// mergeScaling overlays only the sizes present in s onto dst, leaving omitted -// fields untouched. This is the partial-update semantics real EKS applies. -func mergeScaling(dst *eksdriver.NodegroupScalingConfig, s *nodegroupScalingConfigJSON) { - if s.MinSize != nil { - dst.MinSize = int(*s.MinSize) +// scalingUpdateFromJSON carries only the sizes present in s; the provider +// merges them onto the current config. +func scalingUpdateFromJSON(s *nodegroupScalingConfigJSON) *eksdriver.NodegroupScalingUpdate { + return &eksdriver.NodegroupScalingUpdate{ + MinSize: optInt(s.MinSize), + MaxSize: optInt(s.MaxSize), + DesiredSize: optInt(s.DesiredSize), } +} - if s.MaxSize != nil { - dst.MaxSize = int(*s.MaxSize) +func optInt(v *int32) *int { + if v == nil { + return nil } - if s.DesiredSize != nil { - dst.DesiredSize = int(*s.DesiredSize) - } + n := int(*v) + + return &n } // taintsToDriver converts wire taints to driver taints. @@ -750,22 +751,21 @@ func updateConfigToJSON(u eksdriver.NodegroupUpdateConfig) *nodegroupUpdateConfi } } -func scalingFromJSON(s *nodegroupScalingConfigJSON) eksdriver.NodegroupScalingConfig { - out := eksdriver.NodegroupScalingConfig{} - - if s.MinSize != nil { - out.MinSize = int(*s.MinSize) - } - - if s.MaxSize != nil { - out.MaxSize = int(*s.MaxSize) - } - - if s.DesiredSize != nil { - out.DesiredSize = int(*s.DesiredSize) +// scalingFromJSON converts a create-time scalingConfig. An empty object means +// none were given (nil). ok is false when only some of the sizes are set. +func scalingFromJSON(s *nodegroupScalingConfigJSON) (cfg *eksdriver.NodegroupScalingConfig, ok bool) { + switch { + case s.MinSize == nil && s.MaxSize == nil && s.DesiredSize == nil: + return nil, true + case s.MinSize == nil || s.MaxSize == nil || s.DesiredSize == nil: + return nil, false } - return out + return &eksdriver.NodegroupScalingConfig{ + MinSize: int(*s.MinSize), + MaxSize: int(*s.MaxSize), + DesiredSize: int(*s.DesiredSize), + }, true } func toClusterJSON(c *eksdriver.Cluster) clusterJSON { diff --git a/server/aws/eks/sdk_nodegroup_scaling_test.go b/server/aws/eks/sdk_nodegroup_scaling_test.go new file mode 100644 index 000000000..cb1ec651b --- /dev/null +++ b/server/aws/eks/sdk_nodegroup_scaling_test.go @@ -0,0 +1,100 @@ +package eks_test + +import ( + "context" + "errors" + "strings" + "testing" + + "github.com/aws/aws-sdk-go-v2/aws" + awseks "github.com/aws/aws-sdk-go-v2/service/eks" + ekstypes "github.com/aws/aws-sdk-go-v2/service/eks/types" +) + +func scalingCluster(t *testing.T, client *awseks.Client) { + t.Helper() + + if _, err := client.CreateCluster(context.Background(), &awseks.CreateClusterInput{ + Name: aws.String("c1"), + RoleArn: aws.String("arn:aws:iam::123456789012:role/eks"), + ResourcesVpcConfig: &ekstypes.VpcConfigRequest{SubnetIds: []string{"subnet-1"}}, + }); err != nil { + t.Fatalf("CreateCluster: %v", err) + } +} + +// TestSDKEKSCreateNodegroupScalingDefault checks the scalingConfig EKS fills +// in when CreateNodegroup omits it. +func TestSDKEKSCreateNodegroupScalingDefault(t *testing.T) { + client := newSDKClient(t) + ctx := context.Background() + scalingCluster(t, client) + + out, err := client.CreateNodegroup(ctx, &awseks.CreateNodegroupInput{ + ClusterName: aws.String("c1"), + NodegroupName: aws.String("ng1"), + NodeRole: aws.String("arn:aws:iam::123456789012:role/node"), + Subnets: []string{"subnet-1"}, + }) + if err != nil { + t.Fatalf("CreateNodegroup: %v", err) + } + + sc := out.Nodegroup.ScalingConfig + if sc == nil || aws.ToInt32(sc.MinSize) != 1 || aws.ToInt32(sc.MaxSize) != 2 || aws.ToInt32(sc.DesiredSize) != 2 { + t.Fatalf("default scalingConfig = %+v, want min 1 max 2 desired 2", sc) + } +} + +// TestSDKEKSCreateNodegroupPartialScalingRejected checks that CreateNodegroup +// needs all three sizes or none of them. +func TestSDKEKSCreateNodegroupPartialScalingRejected(t *testing.T) { + client := newSDKClient(t) + ctx := context.Background() + scalingCluster(t, client) + + _, err := client.CreateNodegroup(ctx, &awseks.CreateNodegroupInput{ + ClusterName: aws.String("c1"), + NodegroupName: aws.String("ng1"), + NodeRole: aws.String("arn:aws:iam::123456789012:role/node"), + Subnets: []string{"subnet-1"}, + ScalingConfig: &ekstypes.NodegroupScalingConfig{MaxSize: aws.Int32(3)}, + }) + + var ipe *ekstypes.InvalidParameterException + if !errors.As(err, &ipe) { + t.Fatalf("expected InvalidParameterException, got %v", err) + } + + if !strings.Contains(aws.ToString(ipe.Message), "minSize, maxSize and desiredSize") { + t.Fatalf("message = %q", aws.ToString(ipe.Message)) + } +} + +// TestSDKEKSUpdateNodegroupScalingMergedValidation checks that a partial +// update is validated against the merged config. +func TestSDKEKSUpdateNodegroupScalingMergedValidation(t *testing.T) { + client := newSDKClient(t) + ctx := context.Background() + scalingCluster(t, client) + + if _, err := client.CreateNodegroup(ctx, &awseks.CreateNodegroupInput{ + ClusterName: aws.String("c1"), + NodegroupName: aws.String("ng1"), + NodeRole: aws.String("arn:aws:iam::123456789012:role/node"), + Subnets: []string{"subnet-1"}, + }); err != nil { + t.Fatalf("CreateNodegroup: %v", err) + } + + _, err := client.UpdateNodegroupConfig(ctx, &awseks.UpdateNodegroupConfigInput{ + ClusterName: aws.String("c1"), + NodegroupName: aws.String("ng1"), + ScalingConfig: &ekstypes.NodegroupScalingConfig{DesiredSize: aws.Int32(3)}, + }) + + var ipe *ekstypes.InvalidParameterException + if !errors.As(err, &ipe) || !strings.Contains(aws.ToString(ipe.Message), "desired capacity 3 can't be greater than max size 2") { + t.Fatalf("expected InvalidParameterException for desired above max, got %v", err) + } +} diff --git a/services/kubernetes/snapshot.go b/services/kubernetes/snapshot.go index 8fbcc2809..1930d4217 100644 --- a/services/kubernetes/snapshot.go +++ b/services/kubernetes/snapshot.go @@ -124,7 +124,7 @@ func (s *ClusterState) snapshot() clusterSnapshot { s.mu.RLock() defer s.mu.RUnlock() - serverVersion := s.serverVersion + serverVersion := s.targetServerVersionLocked() cs := clusterSnapshot{ RV: s.rv, diff --git a/services/kubernetes/snapshot_test.go b/services/kubernetes/snapshot_test.go index 90dae8525..e30eb9ff6 100644 --- a/services/kubernetes/snapshot_test.go +++ b/services/kubernetes/snapshot_test.go @@ -38,6 +38,7 @@ var clusterStatePersistedFields = map[string]struct{}{ "endpoints": {}, // typed store (clusterSnapshot.Endpoints) "reg": {}, // registry store items (clusterSnapshot.Registry) "serverVersion": {}, // /version body (clusterSnapshot.ServerVersion) + "pendingVersion": {}, // folded into clusterSnapshot.ServerVersion (the target version) } // clusterStateRuntimeFields lists the ClusterState fields deliberately NOT diff --git a/services/kubernetes/state.go b/services/kubernetes/state.go index 8d3c07991..bf2b2d77e 100644 --- a/services/kubernetes/state.go +++ b/services/kubernetes/state.go @@ -67,6 +67,9 @@ type ClusterState struct { // set it from the cluster's control-plane version (SetClusterVersion); a // cluster with no cloud parent keeps defaultServerVersion. serverVersion version.Info + // pendingVersion, when set, replaces serverVersion once its time comes + // (see SetClusterVersionAt). + pendingVersion *pendingServerVersion // managedNodes turns true the first time a managed node pool (SyncNodePool) // adds a Node. The bootstrap nodes are retired at that point and scheduling diff --git a/services/kubernetes/version.go b/services/kubernetes/version.go index 990cf5e6b..33871e162 100644 --- a/services/kubernetes/version.go +++ b/services/kubernetes/version.go @@ -5,8 +5,11 @@ import ( "encoding/hex" "regexp" "runtime" + "time" "k8s.io/apimachinery/pkg/version" + + "github.com/stackshy/cloudemu/v2/config" ) // defaultKubernetesVersion is what /version reports for a cluster that no cloud @@ -117,8 +120,9 @@ func commitFor(seed string) string { // SetClusterVersion sets the Kubernetes version the cluster uid reports on // /version, formatted the way distribution d's control plane formats it. The // EKS, AKS and GKE providers call it on create and on every version change. -// It reports false, leaving the cluster as it was, when uid is unknown or v is -// not a concrete 1.x version. +// It drops any switch scheduled by SetClusterVersionAt. It reports false, +// leaving the cluster as it was, when uid is unknown or v is not a concrete +// 1.x version. func (s *APIServer) SetClusterVersion(uid string, d Distribution, v string) bool { state := s.Lookup(uid) if state == nil { @@ -132,15 +136,65 @@ func (s *APIServer) SetClusterVersion(uid string, d Distribution, v string) bool state.mu.Lock() state.serverVersion = info + state.pendingVersion = nil + state.mu.Unlock() + + return true +} + +// SetClusterVersionAt schedules a version change: the cluster keeps reporting +// its current version until clock reaches at, then reports v. It models a +// control-plane upgrade that is still in progress. When at is not in the +// future it behaves like SetClusterVersion. The return value is as for +// SetClusterVersion. +func (s *APIServer) SetClusterVersionAt(uid string, d Distribution, v string, clock config.Clock, at time.Time) bool { + if clock == nil || !clock.Now().Before(at) { + return s.SetClusterVersion(uid, d, v) + } + + state := s.Lookup(uid) + if state == nil { + return false + } + + info, ok := serverVersionFor(d, v) + if !ok { + return false + } + + state.mu.Lock() + state.pendingVersion = &pendingServerVersion{info: info, clock: clock, at: at} state.mu.Unlock() return true } +// pendingServerVersion is a version change scheduled by SetClusterVersionAt. +type pendingServerVersion struct { + info version.Info + clock config.Clock + at time.Time +} + // ServerVersion returns what this cluster reports on /version. func (s *ClusterState) ServerVersion() version.Info { s.mu.RLock() defer s.mu.RUnlock() + if p := s.pendingVersion; p != nil && !p.clock.Now().Before(p.at) { + return p.info + } + + return s.serverVersion +} + +// targetServerVersionLocked is the version the cluster reports once any +// scheduled switch has happened. A snapshot stores this, since the settle +// windows that drive the switch are not persisted. Caller holds s.mu. +func (s *ClusterState) targetServerVersionLocked() version.Info { + if s.pendingVersion != nil { + return s.pendingVersion.info + } + return s.serverVersion } diff --git a/services/kubernetes/version_test.go b/services/kubernetes/version_test.go index babfb5b15..1d9a0508d 100644 --- a/services/kubernetes/version_test.go +++ b/services/kubernetes/version_test.go @@ -7,8 +7,11 @@ import ( "net/http/httptest" "regexp" "testing" + "time" "k8s.io/apimachinery/pkg/version" + + "github.com/stackshy/cloudemu/v2/config" ) func getServerVersion(t *testing.T, api *APIServer, uid string) version.Info { @@ -189,3 +192,52 @@ func TestServerVersion_SnapshotRoundTrip(t *testing.T) { t.Errorf("restored /version = %+v, want %+v", got, want) } } + +// TestServerVersion_ScheduledSwitch checks that SetClusterVersionAt keeps the +// current /version until the given instant, then reports the new one, and that +// a snapshot taken mid-switch restores the target version. +func TestServerVersion_ScheduledSwitch(t *testing.T) { + ctx := context.Background() + fc := config.NewFakeClock(time.Date(2025, 1, 1, 0, 0, 0, 0, time.UTC)) + + api := NewAPIServer() + uid, _ := api.RegisterCluster() + api.SetClusterVersion(uid, DistributionEKS, "1.31.4") + + if !api.SetClusterVersionAt(uid, DistributionEKS, "1.32.1", fc, fc.Now().Add(time.Second)) { + t.Fatal("SetClusterVersionAt reported false for a known cluster") + } + + if got := getServerVersion(t, api, uid).Minor; got != "31+" { + t.Fatalf("before switch minor = %q, want 31+", got) + } + + data, err := api.Snapshot(ctx, false) + if err != nil { + t.Fatalf("Snapshot: %v", err) + } + + fc.Advance(time.Second) + + if got := getServerVersion(t, api, uid).Minor; got != "32+" { + t.Fatalf("after switch minor = %q, want 32+", got) + } + + dst := NewAPIServer() + if err := dst.Restore(ctx, data); err != nil { + t.Fatalf("Restore: %v", err) + } + + if got := getServerVersion(t, dst, uid).Minor; got != "32+" { + t.Fatalf("restored minor = %q, want 32+", got) + } + + // A later immediate set drops the pending switch. + api.SetClusterVersionAt(uid, DistributionEKS, "1.33.0", fc, fc.Now().Add(time.Second)) + api.SetClusterVersion(uid, DistributionEKS, "1.32.1") + fc.Advance(time.Second) + + if got := getServerVersion(t, api, uid).Minor; got != "32+" { + t.Fatalf("after overriding set minor = %q, want 32+", got) + } +}