diff --git a/README.md b/README.md index 9d2d93a8..9cb6c17e 100644 --- a/README.md +++ b/README.md @@ -178,6 +178,74 @@ list, which briefly interrupts traffic on that port. Setting it to an empty value (`""`) sends an empty CIDR list to CloudStack — it does not block all traffic. +#### `service.beta.kubernetes.io/cloudstack-load-balancer-stickiness-method-name` + +**Type:** String + +**Default:** Not set (no stickiness policy) + +**Description:** Creates a CloudStack **LB stickiness policy** on every load balancer rule belonging +to the service, making the load balancer keep a client on the same backend node between requests. + +The value is the CloudStack stickiness method name and is passed through to CloudStack, which +validates it against the methods the network's load balancer provider offers. The VirtualRouter +(HAProxy) provider supports `LbCookie`, `AppCookie` and `SourceBased`. An unsupported method makes +the service fail to sync with an `error creating stickiness policy` error. + +Each service port has its own load balancer rule, so a service exposing several ports gets one +policy per port, all with the same method and parameters. + +**Use Case:** Applications that keep per-client state in the backend — a session held in process +memory, for example — and therefore need successive requests from one client to land on the same +node. + +**Example:** +```yaml +apiVersion: v1 +kind: Service +metadata: + name: my-service + annotations: + service.beta.kubernetes.io/cloudstack-load-balancer-stickiness-method-name: "LbCookie" + service.beta.kubernetes.io/cloudstack-load-balancer-stickiness-method-param: "cookie-name=SERVERID" +spec: + type: LoadBalancer +``` + +**Note:** Removing the annotation deletes the stickiness policy and leaves the load balancer rule in +place. Changing either the method name or the parameters replaces the policy: the controller deletes +the existing policy and creates a new one, which resets whatever affinity state the load balancer +was holding. CloudStack allows only one policy per rule, so the old policy has to go first; if the +replacement is rejected, the controller puts the previous policy back and reports the failure, so a +running service is never left without the stickiness it already had. + +#### `service.beta.kubernetes.io/cloudstack-load-balancer-stickiness-method-param` + +**Type:** String (comma-separated `key=value` list) + +**Default:** Not set (no parameters) + +**Description:** Parameters for the stickiness method selected by +`service.beta.kubernetes.io/cloudstack-load-balancer-stickiness-method-name`. Which keys are +accepted depends on the method. With the VirtualRouter provider, `LbCookie` takes `cookie-name`, +`mode`, `nocache`, `indirect`, `postonly` and `domain`; `AppCookie` takes `cookie-name`, `length`, +`holdtime` and `mode`; `SourceBased` takes `tablesize` and `expire`. CloudStack validates the keys +against the methods the network's provider advertises, so an unknown parameter makes the service +fail to sync. + +This annotation has no effect on its own: without a method name no policy is created. + +**Format:** Comma-separated `key=value` pairs. Spaces around entries, keys and values are trimmed, +and only the first `=` separates key from value, so a value may itself contain `=`. Entries without +a `=` are ignored rather than rejected, an empty value (`key=`) is passed through as an empty +string, and if a key repeats, the last occurrence wins. + +**Example:** +```yaml + service.beta.kubernetes.io/cloudstack-load-balancer-stickiness-method-name: "AppCookie" + service.beta.kubernetes.io/cloudstack-load-balancer-stickiness-method-param: "cookie-name=JSESSIONID,length=52" +``` + #### `service.beta.kubernetes.io/cloudstack-load-balancer-ip-associated-by-controller` **Type:** Boolean (`"true"` or `"false"`) @@ -267,6 +335,14 @@ annotation for it. Any other value makes the service fail to sync with `unsupported load balancer affinity`. Other CloudStack algorithms, such as `leastconn`, cannot currently be selected. +The algorithm is separate from stickiness. `spec.sessionAffinity: ClientIP` picks the load balancer +algorithm, while a CloudStack stickiness policy — cookie-based affinity, for instance — is +configured with the +[`stickiness-method-name`](#servicebetakubernetesiocloudstack-load-balancer-stickiness-method-name) +and +[`stickiness-method-param`](#servicebetakubernetesiocloudstack-load-balancer-stickiness-method-param) +annotations. The two can be used together. + ### VPC Networks VPC networks are supported. VPC networks normally do not offer the Firewall service, so the diff --git a/cloudstack_loadbalancer.go b/cloudstack_loadbalancer.go index 69d16f51..c07c5d2b 100644 --- a/cloudstack_loadbalancer.go +++ b/cloudstack_loadbalancer.go @@ -56,6 +56,9 @@ const ( // associated the IP address. This annotation is set by the controller when it associates // an unallocated IP, and is used to determine if the IP should be disassociated on deletion. ServiceAnnotationLoadBalancerIPAssociatedByController = "service.beta.kubernetes.io/cloudstack-load-balancer-ip-associated-by-controller" //nolint:gosec + + ServiceAnnotationLoadBalancerStickinessMethodName = "service.beta.kubernetes.io/cloudstack-load-balancer-stickiness-method-name" + ServiceAnnotationLoadBalancerStickinessParam = "service.beta.kubernetes.io/cloudstack-load-balancer-stickiness-method-param" ) type loadBalancer struct { @@ -70,6 +73,7 @@ type loadBalancer struct { projectID string rules map[string]*cloudstack.LoadBalancerRule duplicateRules []*cloudstack.LoadBalancerRule + stickinessPolicies map[string]*cloudstack.LBStickinessPolicyStickinesspolicy ipAssociatedByController bool } @@ -114,6 +118,10 @@ func (cs *CSCloud) EnsureLoadBalancer(ctx context.Context, clusterName string, s return nil, err } + if err := lb.loadStickinessPolicies(); err != nil { + return nil, err + } + // Set the load balancer algorithm. switch service.Spec.SessionAffinity { case corev1.ServiceAffinityNone: @@ -186,15 +194,18 @@ func (cs *CSCloud) EnsureLoadBalancer(ctx context.Context, clusterName string, s // Delete the rule from the map, to prevent it being deleted. delete(lb.rules, lbRuleName) } - } else { - klog.V(4).Infof("Creating load balancer rule: %v", lbRuleName) - lbRule, err = lb.createLoadBalancerRule(lbRuleName, port, protocol, service) - if err != nil { + + if err := lb.syncRuleHosts(lbRule); err != nil { return nil, err } - klog.V(4).Infof("Assigning hosts (%v) to load balancer rule: %v", lb.hostIDs, lbRuleName) - if err = lb.assignHostsToRule(lbRule, lb.hostIDs); err != nil { + if err := lb.reconcileStickinessPolicy(lbRuleName, lbRule, service); err != nil { + return nil, err + } + } else { + klog.V(4).Infof("Creating load balancer rule: %v", lbRuleName) + lbRule, err = lb.createConfiguredLoadBalancerRule(lbRuleName, port, protocol, service) + if err != nil { return nil, err } } @@ -279,28 +290,38 @@ func (cs *CSCloud) UpdateLoadBalancer(ctx context.Context, clusterName string, s } for _, lbRule := range lb.rules { - p := lb.LoadBalancer.NewListLoadBalancerRuleInstancesParams(lbRule.Id) - - // Retrieve all VMs currently associated to this load balancer rule. - l, err := lb.LoadBalancer.ListLoadBalancerRuleInstances(p) - if err != nil { - return fmt.Errorf("error retrieving associated instances: %v", err) + if err := lb.syncRuleHosts(lbRule); err != nil { + return err } + } - assign, remove := symmetricDifference(lb.hostIDs, l.LoadBalancerRuleInstances) + return nil +} - if len(assign) > 0 { - klog.V(4).Infof("Assigning new hosts (%v) to load balancer rule: %v", assign, lbRule.Name) - if err := lb.assignHostsToRule(lbRule, assign); err != nil { - return err - } +// syncRuleHosts brings the VMs assigned to a rule in line with the current +// hosts, so a rule that was adopted in a half-configured state, or whose +// membership drifted, converges on the next sync instead of waiting for a node +// change. +func (lb *loadBalancer) syncRuleHosts(lbRule *cloudstack.LoadBalancerRule) error { + p := lb.LoadBalancer.NewListLoadBalancerRuleInstancesParams(lbRule.Id) + l, err := lb.LoadBalancer.ListLoadBalancerRuleInstances(p) + if err != nil { + return fmt.Errorf("error retrieving associated instances: %v", err) + } + + assign, remove := symmetricDifference(lb.hostIDs, l.LoadBalancerRuleInstances) + + if len(assign) > 0 { + klog.V(4).Infof("Assigning new hosts (%v) to load balancer rule: %v", assign, lbRule.Name) + if err := lb.assignHostsToRule(lbRule, assign); err != nil { + return err } + } - if len(remove) > 0 { - klog.V(4).Infof("Removing old hosts (%v) from load balancer rule: %v", assign, lbRule.Name) - if err := lb.removeHostsFromRule(lbRule, remove); err != nil { - return err - } + if len(remove) > 0 { + klog.V(4).Infof("Removing old hosts (%v) from load balancer rule: %v", remove, lbRule.Name) + if err := lb.removeHostsFromRule(lbRule, remove); err != nil { + return err } } @@ -446,10 +467,11 @@ func (cs *CSCloud) GetLoadBalancerName(ctx context.Context, clusterName string, // getLoadBalancer retrieves the IP address and ID and all the existing rules it can find. func (cs *CSCloud) getLoadBalancer(service *corev1.Service) (*loadBalancer, error) { lb := &loadBalancer{ - CloudStackClient: cs.client, - name: cs.GetLoadBalancerName(context.TODO(), "", service), - projectID: cs.projectID, - rules: make(map[string]*cloudstack.LoadBalancerRule), + CloudStackClient: cs.client, + name: cs.GetLoadBalancerName(context.TODO(), "", service), + projectID: cs.projectID, + rules: make(map[string]*cloudstack.LoadBalancerRule), + stickinessPolicies: make(map[string]*cloudstack.LBStickinessPolicyStickinesspolicy), } p := cs.client.LoadBalancer.NewListLoadBalancerRulesParams() @@ -499,6 +521,26 @@ func (cs *CSCloud) getLoadBalancer(service *corev1.Service) (*loadBalancer, erro return lb, nil } +// loadStickinessPolicies fetches the stickiness policy of every rule the load +// balancer keeps. Only EnsureLoadBalancer reconciles policies, so the lookups +// stay off the host-update and delete paths. CloudStack answers with a wrapper +// per rule even when the rule has no policy, with an empty list inside it. +func (lb *loadBalancer) loadStickinessPolicies() error { + for _, lbRule := range lb.rules { + p := lb.LoadBalancer.NewListLBStickinessPoliciesParams() + p.SetLbruleid(lbRule.Id) + l, err := lb.LoadBalancer.ListLBStickinessPolicies(p) + if err != nil { + return fmt.Errorf("error retrieving stickiness policies for rule %v: %v", lbRule.Name, err) + } + if len(l.LBStickinessPolicies) > 0 && len(l.LBStickinessPolicies[0].Stickinesspolicy) > 0 { + lb.stickinessPolicies[lbRule.Id] = &l.LBStickinessPolicies[0].Stickinesspolicy[0] + } + } + + return nil +} + // Get network ID from Public IP Address // Every failure returns an error: GetNetworkByID does not reject an empty ID but // matches an unfiltered network list, so ("", nil) would resolve to any network. @@ -681,6 +723,52 @@ func (lb *loadBalancer) getCIDRList(service *corev1.Service) ([]string, error) { return cidrList, nil } +// checkStickinessPolicy returns the rule's current stickiness policy, if it has +// one, and whether that policy differs from what the service annotations ask for. +func (lb *loadBalancer) checkStickinessPolicy(lbRule *cloudstack.LoadBalancerRule, service *corev1.Service) (*cloudstack.LBStickinessPolicyStickinesspolicy, bool) { + policy := lb.stickinessPolicies[lbRule.Id] + wantedMethod := wantedStickinessMethod(service) + wantedParams := parseStickinessParams(getStringFromServiceAnnotation(service, ServiceAnnotationLoadBalancerStickinessParam, "")) + + if policy == nil { + if wantedMethod == "" { + return nil, false + } + klog.V(4).Infof("Stickiness policy missing from rule %v", lbRule.Name) + return nil, true + } + if wantedMethod == "" { + klog.V(4).Infof("Stickiness annotation removed from rule %v", lbRule.Name) + return policy, true + } + if !strings.EqualFold(policy.Methodname, wantedMethod) { + klog.V(4).Infof("Stickiness method of rule %v changed to %v", lbRule.Name, wantedMethod) + return policy, true + } + if !stickinessParamsMatch(policy.Params, wantedParams) { + klog.V(4).Infof("Stickiness parameters of rule %v changed to %v", lbRule.Name, wantedParams) + return policy, true + } + + return policy, false +} + +// stickinessParamsMatch reports whether a policy carries exactly the parameters +// asked for. A key that is absent differs from one set to an empty value, so the +// two maps have to agree on keys as well as values. +func stickinessParamsMatch(current, wanted map[string]string) bool { + if len(current) != len(wanted) { + return false + } + for key, value := range current { + if wantedValue, ok := wanted[key]; !ok || wantedValue != value { + return false + } + } + + return true +} + // checkLoadBalancerRule checks if the rule already exists and if it does, if it can be updated. If // it does exist but cannot be updated, it will delete the existing rule so it can be created again. func (lb *loadBalancer) checkLoadBalancerRule(lbRuleName string, port corev1.ServicePort, protocol LoadBalancerProtocol, service *corev1.Service, version semver.Version) (*cloudstack.LoadBalancerRule, bool, error) { @@ -747,6 +835,100 @@ func (lb *loadBalancer) updateLoadBalancerRule(lbRuleName string, protocol LoadB return err } +// wantedStickinessMethod returns the stickiness method the service asks for, or +// an empty string when the service does not want a stickiness policy. +func wantedStickinessMethod(service *corev1.Service) string { + return getStringFromServiceAnnotation(service, ServiceAnnotationLoadBalancerStickinessMethodName, "") +} + +// reconcileStickinessPolicy brings a rule's stickiness policy in line with the +// service annotations. +func (lb *loadBalancer) reconcileStickinessPolicy(lbRuleName string, lbRule *cloudstack.LoadBalancerRule, service *corev1.Service) error { + existing, needsUpdate := lb.checkStickinessPolicy(lbRule, service) + if !needsUpdate { + return nil + } + if existing == nil { + klog.V(4).Infof("Creating stickiness policy on rule %v", lbRuleName) + _, err := lb.createStickinessPolicy(lbRuleName, lbRule.Id, service) + return err + } + return lb.replaceStickinessPolicy(lbRuleName, lbRule, existing, service) +} + +// replaceStickinessPolicy swaps a rule's policy for the one the annotations ask +// for, or drops it when the annotation is gone. CloudStack keeps at most one +// policy per rule, so the old one has to go first, and a rejected replacement +// is rolled back rather than leaving a live rule with no stickiness at all. +func (lb *loadBalancer) replaceStickinessPolicy(lbRuleName string, lbRule *cloudstack.LoadBalancerRule, existing *cloudstack.LBStickinessPolicyStickinesspolicy, service *corev1.Service) error { + if wantedStickinessMethod(service) == "" { + klog.V(4).Infof("Removing stickiness policy from rule %v", lbRuleName) + } else { + klog.V(4).Infof("Replacing stickiness policy on rule %v", lbRuleName) + } + if err := lb.deleteStickinessPolicy(existing.Id); err != nil { + return err + } + + if _, err := lb.createStickinessPolicy(lbRuleName, lbRule.Id, service); err != nil { + return lb.restoreStickinessPolicy(lbRuleName, lbRule.Id, existing, err) + } + + return nil +} + +// restoreStickinessPolicy puts a deleted policy back after its replacement could +// not be created, and returns the failure that triggered the rollback. +func (lb *loadBalancer) restoreStickinessPolicy(lbRuleName string, lbRuleID string, policy *cloudstack.LBStickinessPolicyStickinesspolicy, cause error) error { + klog.V(4).Infof("Restoring the previous stickiness policy on rule %v after a failed replacement: %v", lbRuleName, cause) + p := lb.LoadBalancer.NewCreateLBStickinessPolicyParams(lbRuleID, policy.Methodname, policy.Name) + p.SetParam(policy.Params) + if _, err := lb.LoadBalancer.CreateLBStickinessPolicy(p); err != nil { + return fmt.Errorf("%v (restoring the previous stickiness policy failed too: %v)", cause, err) + } + + return cause +} + +// createStickinessPolicy creates a new stickiness policy and returns it. +func (lb *loadBalancer) createStickinessPolicy(lbRuleName string, lbRuleId string, service *corev1.Service) (*cloudstack.LBStickinessPolicyStickinesspolicy, error) { + stickinessMethodName := wantedStickinessMethod(service) + stickinessMethodParam := getStringFromServiceAnnotation(service, ServiceAnnotationLoadBalancerStickinessParam, "") + // If the stickiness method name is not set, we don't need to create a stickiness policy. + if stickinessMethodName == "" { + return nil, nil + } + p := lb.LoadBalancer.NewCreateLBStickinessPolicyParams(lbRuleId, stickinessMethodName, lbRuleName) + + params := parseStickinessParams(stickinessMethodParam) + p.SetParam(params) + + stickinessPolicy, err := lb.LoadBalancer.CreateLBStickinessPolicy(p) + if err != nil { + return nil, fmt.Errorf("error creating stickiness policy: %v", err) + } + if len(stickinessPolicy.Stickinesspolicy) == 0 { + return nil, fmt.Errorf("error creating stickiness policy: no policy returned for load balancer rule %v", lbRuleName) + } + return &cloudstack.LBStickinessPolicyStickinesspolicy{ + Methodname: stickinessPolicy.Stickinesspolicy[0].Methodname, + Params: stickinessPolicy.Stickinesspolicy[0].Params, + Id: stickinessPolicy.Stickinesspolicy[0].Id, + Name: stickinessPolicy.Stickinesspolicy[0].Name, + State: stickinessPolicy.Stickinesspolicy[0].State, + }, nil +} + +// deleteStickinessPolicy deletes a stickiness policy. +func (lb *loadBalancer) deleteStickinessPolicy(stickinessPolicyId string) error { + p := lb.LoadBalancer.NewDeleteLBStickinessPolicyParams(stickinessPolicyId) + + if _, err := lb.LoadBalancer.DeleteLBStickinessPolicy(p); err != nil { + return fmt.Errorf("error deleting stickiness policy %v: %v", stickinessPolicyId, err) + } + return nil +} + // createLoadBalancerRule creates a new load balancer rule and returns it's ID. func (lb *loadBalancer) createLoadBalancerRule(lbRuleName string, port corev1.ServicePort, protocol LoadBalancerProtocol, service *corev1.Service) (*cloudstack.LoadBalancerRule, error) { p := lb.LoadBalancer.NewCreateLoadBalancerRuleParams( @@ -794,6 +976,38 @@ func (lb *loadBalancer) createLoadBalancerRule(lbRuleName string, port corev1.Se return lbRule, nil } +// createConfiguredLoadBalancerRule creates a load balancer rule together with +// its stickiness policy and host assignments. A failure after the rule exists +// removes the rule again: a later sync would otherwise find it, take the +// existing-rule path, and never assign its hosts. +func (lb *loadBalancer) createConfiguredLoadBalancerRule(lbRuleName string, port corev1.ServicePort, protocol LoadBalancerProtocol, service *corev1.Service) (*cloudstack.LoadBalancerRule, error) { + lbRule, err := lb.createLoadBalancerRule(lbRuleName, port, protocol, service) + if err != nil { + return nil, err + } + if _, err := lb.createStickinessPolicy(lbRuleName, lbRule.Id, service); err != nil { + return nil, lb.rollBackLoadBalancerRule(lbRule, err) + } + + klog.V(4).Infof("Assigning hosts (%v) to load balancer rule: %v", lb.hostIDs, lbRuleName) + if err := lb.assignHostsToRule(lbRule, lb.hostIDs); err != nil { + return nil, lb.rollBackLoadBalancerRule(lbRule, err) + } + + return lbRule, nil +} + +// rollBackLoadBalancerRule deletes a rule whose configuration failed part way +// through and returns the original failure, extended with the cleanup error +// when the rule could not be removed either. +func (lb *loadBalancer) rollBackLoadBalancerRule(lbRule *cloudstack.LoadBalancerRule, cause error) error { + klog.V(4).Infof("Rolling back load balancer rule %v after a configuration failure: %v", lbRule.Name, cause) + if err := lb.deleteLoadBalancerRule(lbRule); err != nil { + return fmt.Errorf("%v (rolling back the rule failed too: %v)", cause, err) + } + return cause +} + // deleteLoadBalancerRule deletes a load balancer rule. func (lb *loadBalancer) deleteLoadBalancerRule(lbRule *cloudstack.LoadBalancerRule) error { p := lb.LoadBalancer.NewDeleteLoadBalancerRuleParams(lbRule.Id) @@ -806,6 +1020,7 @@ func (lb *loadBalancer) deleteLoadBalancerRule(lbRule *cloudstack.LoadBalancerRu if kept, ok := lb.rules[lbRule.Name]; ok && kept.Id == lbRule.Id { delete(lb.rules, lbRule.Name) } + delete(lb.stickinessPolicies, lbRule.Id) return nil } @@ -1236,6 +1451,21 @@ func getStringFromServiceAnnotation(service *corev1.Service, annotationKey strin return defaultSetting } +// parseStickinessParams parses a comma-separated string of key=value pairs into a map. +// Whitespace around entries, keys and values is trimmed. Empty entries and entries +// without a "=" are skipped, while an empty value ("key=") is kept as an empty string. +func parseStickinessParams(paramString string) map[string]string { + params := make(map[string]string) + for _, param := range strings.Split(paramString, ",") { + parts := strings.SplitN(param, "=", 2) + if len(parts) != 2 { + continue + } + params[strings.TrimSpace(parts[0])] = strings.TrimSpace(parts[1]) + } + return params +} + // getBoolFromServiceAnnotation searches a given v1.Service for a specific annotationKey and either returns the annotation's boolean value or a specified defaultSetting func getBoolFromServiceAnnotation(service *corev1.Service, annotationKey string, defaultSetting bool) bool { klog.V(4).Infof("getBoolFromServiceAnnotation(%s/%s, %v, %v)", service.Namespace, service.Name, annotationKey, defaultSetting) diff --git a/cloudstack_loadbalancer_test.go b/cloudstack_loadbalancer_test.go index 39b229d5..6327560f 100644 --- a/cloudstack_loadbalancer_test.go +++ b/cloudstack_loadbalancer_test.go @@ -3992,3 +3992,732 @@ func TestVerifyHosts(t *testing.T) { } }) } + +func TestParseStickinessParams(t *testing.T) { + tests := []struct { + name string + input string + want map[string]string + }{ + { + name: "empty string returns empty map", + input: "", + want: map[string]string{}, + }, + { + name: "single pair", + input: "cookie-name=SERVERID", + want: map[string]string{"cookie-name": "SERVERID"}, + }, + { + name: "multiple pairs", + input: "cookie-name=SERVERID,mode=insert", + want: map[string]string{"cookie-name": "SERVERID", "mode": "insert"}, + }, + { + name: "whitespace around entries is trimmed", + input: " cookie-name=SERVERID , mode=insert ", + want: map[string]string{"cookie-name": "SERVERID", "mode": "insert"}, + }, + { + name: "whitespace around keys and values is trimmed", + input: "cookie-name = SERVERID, mode =insert", + want: map[string]string{"cookie-name": "SERVERID", "mode": "insert"}, + }, + { + name: "entry without separator is ignored", + input: "cookie-name=SERVERID,bogus", + want: map[string]string{"cookie-name": "SERVERID"}, + }, + { + name: "trailing comma is ignored", + input: "mode=insert,", + want: map[string]string{"mode": "insert"}, + }, + { + name: "only separators returns empty map", + input: ",,,", + want: map[string]string{}, + }, + { + name: "empty value is preserved", + input: "cookie-name=", + want: map[string]string{"cookie-name": ""}, + }, + { + name: "value containing separator is kept intact", + input: "expr=a=b", + want: map[string]string{"expr": "a=b"}, + }, + { + name: "duplicate key keeps last value", + input: "mode=insert,mode=rewrite", + want: map[string]string{"mode": "rewrite"}, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got := parseStickinessParams(tt.input) + if !reflect.DeepEqual(got, tt.want) { + t.Errorf("parseStickinessParams(%q) = %v, want %v", tt.input, got, tt.want) + } + }) + } +} + +// stickinessTestService builds a service carrying the stickiness annotations. An +// empty method or param string omits that annotation entirely. +func stickinessTestService(methodName, params string) *corev1.Service { + annotations := map[string]string{} + if methodName != "" { + annotations[ServiceAnnotationLoadBalancerStickinessMethodName] = methodName + } + if params != "" { + annotations[ServiceAnnotationLoadBalancerStickinessParam] = params + } + return &corev1.Service{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test-service", + Namespace: "default", + Annotations: annotations, + }, + } +} + +func TestCheckStickinessPolicy(t *testing.T) { + lbRule := &cloudstack.LoadBalancerRule{Id: "rule-id", Name: "test-service-tcp-80"} + + tests := []struct { + name string + existingPolicy *cloudstack.LBStickinessPolicyStickinesspolicy + methodName string + params string + wantPolicy bool // true when the existing policy is expected back + wantNeedsUpdate bool + }{ + { + name: "no policy and no annotation is a no-op", + existingPolicy: nil, + methodName: "", + wantPolicy: false, + wantNeedsUpdate: false, + }, + { + name: "no policy with annotation needs creation", + existingPolicy: nil, + methodName: "LbCookie", + wantPolicy: false, + wantNeedsUpdate: true, + }, + { + name: "policy with annotation removed needs deletion", + existingPolicy: &cloudstack.LBStickinessPolicyStickinesspolicy{ + Id: "policy-id", + Methodname: "LbCookie", + }, + methodName: "", + wantPolicy: true, + wantNeedsUpdate: true, + }, + { + name: "matching method with no params is up-to-date", + existingPolicy: &cloudstack.LBStickinessPolicyStickinesspolicy{ + Id: "policy-id", + Methodname: "LbCookie", + Params: map[string]string{}, + }, + methodName: "LbCookie", + params: "", + wantPolicy: true, + wantNeedsUpdate: false, + }, + { + name: "matching method and params is up-to-date", + existingPolicy: &cloudstack.LBStickinessPolicyStickinesspolicy{ + Id: "policy-id", + Methodname: "AppCookie", + Params: map[string]string{"cookie-name": "SERVERID", "mode": "insert"}, + }, + methodName: "AppCookie", + params: "cookie-name=SERVERID,mode=insert", + wantPolicy: true, + wantNeedsUpdate: false, + }, + { + name: "method name mismatch needs recreation", + existingPolicy: &cloudstack.LBStickinessPolicyStickinesspolicy{ + Id: "policy-id", + Methodname: "LbCookie", + }, + methodName: "AppCookie", + wantPolicy: true, + wantNeedsUpdate: true, + }, + { + name: "extra desired param needs recreation", + existingPolicy: &cloudstack.LBStickinessPolicyStickinesspolicy{ + Id: "policy-id", + Methodname: "AppCookie", + Params: map[string]string{"cookie-name": "SERVERID"}, + }, + methodName: "AppCookie", + params: "cookie-name=SERVERID,mode=insert", + wantPolicy: true, + wantNeedsUpdate: true, + }, + { + name: "removed desired param needs recreation", + existingPolicy: &cloudstack.LBStickinessPolicyStickinesspolicy{ + Id: "policy-id", + Methodname: "AppCookie", + Params: map[string]string{"cookie-name": "SERVERID", "mode": "insert"}, + }, + methodName: "AppCookie", + params: "cookie-name=SERVERID", + wantPolicy: true, + wantNeedsUpdate: true, + }, + { + name: "param value mismatch needs recreation", + existingPolicy: &cloudstack.LBStickinessPolicyStickinesspolicy{ + Id: "policy-id", + Methodname: "AppCookie", + Params: map[string]string{"cookie-name": "SERVERID"}, + }, + methodName: "AppCookie", + params: "cookie-name=JSESSIONID", + wantPolicy: true, + wantNeedsUpdate: true, + }, + { + name: "renamed param key of equal count needs recreation", + existingPolicy: &cloudstack.LBStickinessPolicyStickinesspolicy{ + Id: "policy-id", + Methodname: "AppCookie", + Params: map[string]string{"cookie-name": ""}, + }, + methodName: "AppCookie", + params: "mode=", + wantPolicy: true, + wantNeedsUpdate: true, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + lb := &loadBalancer{ + stickinessPolicies: map[string]*cloudstack.LBStickinessPolicyStickinesspolicy{}, + } + if tt.existingPolicy != nil { + lb.stickinessPolicies[lbRule.Id] = tt.existingPolicy + } + + policy, needsUpdate := lb.checkStickinessPolicy(lbRule, stickinessTestService(tt.methodName, tt.params)) + if tt.wantPolicy && policy != tt.existingPolicy { + t.Errorf("policy = %v, want the existing policy %v", policy, tt.existingPolicy) + } + if !tt.wantPolicy && policy != nil { + t.Errorf("policy = %v, want nil", policy) + } + if needsUpdate != tt.wantNeedsUpdate { + t.Errorf("needsUpdate = %v, want %v", needsUpdate, tt.wantNeedsUpdate) + } + }) + } +} + +func TestCreateStickinessPolicy(t *testing.T) { + t.Run("no method annotation is a no-op", func(t *testing.T) { + ctrl := gomock.NewController(t) + t.Cleanup(ctrl.Finish) + + // No expectations on the mock; any API call would fail the test. + mockLB := cloudstack.NewMockLoadBalancerServiceIface(ctrl) + + lb := &loadBalancer{ + CloudStackClient: &cloudstack.CloudStackClient{ + LoadBalancer: mockLB, + }, + } + + policy, err := lb.createStickinessPolicy("test-service-tcp-80", "rule-id", stickinessTestService("", "")) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if policy != nil { + t.Errorf("policy = %v, want nil", policy) + } + }) + + t.Run("creates policy from annotations", func(t *testing.T) { + ctrl := gomock.NewController(t) + t.Cleanup(ctrl.Finish) + + mockLB := cloudstack.NewMockLoadBalancerServiceIface(ctrl) + createParams := &cloudstack.CreateLBStickinessPolicyParams{} + createResp := &cloudstack.CreateLBStickinessPolicyResponse{ + Lbruleid: "rule-id", + Stickinesspolicy: []cloudstack.CreateLBStickinessPolicyResponseStickinesspolicy{ + { + Id: "policy-id", + Name: "test-service-tcp-80", + Methodname: "AppCookie", + Params: map[string]string{"cookie-name": "SERVERID"}, + State: "Active", + }, + }, + } + + gomock.InOrder( + mockLB.EXPECT().NewCreateLBStickinessPolicyParams("rule-id", "AppCookie", "test-service-tcp-80").Return(createParams), + mockLB.EXPECT().CreateLBStickinessPolicy(createParams).Return(createResp, nil), + ) + + lb := &loadBalancer{ + CloudStackClient: &cloudstack.CloudStackClient{ + LoadBalancer: mockLB, + }, + } + + policy, err := lb.createStickinessPolicy("test-service-tcp-80", "rule-id", stickinessTestService("AppCookie", "cookie-name=SERVERID")) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if policy == nil { + t.Fatal("expected a policy, got nil") + } + if policy.Id != "policy-id" { + t.Errorf("policy ID = %q, want %q", policy.Id, "policy-id") + } + if policy.Methodname != "AppCookie" { + t.Errorf("policy method = %q, want %q", policy.Methodname, "AppCookie") + } + if policy.Name != "test-service-tcp-80" { + t.Errorf("policy name = %q, want %q", policy.Name, "test-service-tcp-80") + } + if policy.State != "Active" { + t.Errorf("policy state = %q, want %q", policy.State, "Active") + } + if !reflect.DeepEqual(policy.Params, map[string]string{"cookie-name": "SERVERID"}) { + t.Errorf("policy params = %v, want %v", policy.Params, map[string]string{"cookie-name": "SERVERID"}) + } + }) + + t.Run("API error is returned", func(t *testing.T) { + ctrl := gomock.NewController(t) + t.Cleanup(ctrl.Finish) + + mockLB := cloudstack.NewMockLoadBalancerServiceIface(ctrl) + createParams := &cloudstack.CreateLBStickinessPolicyParams{} + apiErr := fmt.Errorf("create policy API error") + + gomock.InOrder( + mockLB.EXPECT().NewCreateLBStickinessPolicyParams("rule-id", "LbCookie", "test-service-tcp-80").Return(createParams), + mockLB.EXPECT().CreateLBStickinessPolicy(createParams).Return(nil, apiErr), + ) + + lb := &loadBalancer{ + CloudStackClient: &cloudstack.CloudStackClient{ + LoadBalancer: mockLB, + }, + } + + policy, err := lb.createStickinessPolicy("test-service-tcp-80", "rule-id", stickinessTestService("LbCookie", "")) + if err == nil { + t.Fatal("expected an error, got nil") + } + if policy != nil { + t.Errorf("policy = %v, want nil", policy) + } + if !strings.Contains(err.Error(), "error creating stickiness policy") { + t.Errorf("error = %q, want it to mention creating the stickiness policy", err.Error()) + } + }) +} + +func TestDeleteStickinessPolicy(t *testing.T) { + t.Run("deletes policy by ID", func(t *testing.T) { + ctrl := gomock.NewController(t) + t.Cleanup(ctrl.Finish) + + mockLB := cloudstack.NewMockLoadBalancerServiceIface(ctrl) + deleteParams := &cloudstack.DeleteLBStickinessPolicyParams{} + + gomock.InOrder( + mockLB.EXPECT().NewDeleteLBStickinessPolicyParams("policy-id").Return(deleteParams), + mockLB.EXPECT().DeleteLBStickinessPolicy(deleteParams).Return(&cloudstack.DeleteLBStickinessPolicyResponse{Success: true}, nil), + ) + + lb := &loadBalancer{ + CloudStackClient: &cloudstack.CloudStackClient{ + LoadBalancer: mockLB, + }, + } + + if err := lb.deleteStickinessPolicy("policy-id"); err != nil { + t.Fatalf("unexpected error: %v", err) + } + }) + + t.Run("API error is returned", func(t *testing.T) { + ctrl := gomock.NewController(t) + t.Cleanup(ctrl.Finish) + + mockLB := cloudstack.NewMockLoadBalancerServiceIface(ctrl) + deleteParams := &cloudstack.DeleteLBStickinessPolicyParams{} + apiErr := fmt.Errorf("delete policy API error") + + gomock.InOrder( + mockLB.EXPECT().NewDeleteLBStickinessPolicyParams("policy-id").Return(deleteParams), + mockLB.EXPECT().DeleteLBStickinessPolicy(deleteParams).Return(nil, apiErr), + ) + + lb := &loadBalancer{ + CloudStackClient: &cloudstack.CloudStackClient{ + LoadBalancer: mockLB, + }, + } + + err := lb.deleteStickinessPolicy("policy-id") + if err == nil { + t.Fatal("expected an error, got nil") + } + if !strings.Contains(err.Error(), "error deleting stickiness policy policy-id") { + t.Errorf("error = %q, want it to mention deleting the stickiness policy", err.Error()) + } + }) +} + +// The shape CloudStack actually returns for a rule with no stickiness policy is +// a policy wrapper with an empty inner Stickinesspolicy list. +func TestLoadStickinessPolicies(t *testing.T) { + newLB := func(mockLB *cloudstack.MockLoadBalancerServiceIface, rules ...*cloudstack.LoadBalancerRule) *loadBalancer { + lb := &loadBalancer{ + CloudStackClient: &cloudstack.CloudStackClient{LoadBalancer: mockLB}, + rules: map[string]*cloudstack.LoadBalancerRule{}, + stickinessPolicies: map[string]*cloudstack.LBStickinessPolicyStickinesspolicy{}, + } + for _, rule := range rules { + lb.rules[rule.Name] = rule + } + return lb + } + withPolicy := &cloudstack.ListLBStickinessPoliciesResponse{ + Count: 1, + LBStickinessPolicies: []*cloudstack.LBStickinessPolicy{{ + Lbruleid: "rule-1", + Stickinesspolicy: []cloudstack.LBStickinessPolicyStickinesspolicy{ + {Id: "policy-1", Name: "test-service-tcp-80", Methodname: "LbCookie", Params: map[string]string{"cookie-name": "SERVERID"}}, + }, + }}, + } + emptyWrapper := &cloudstack.ListLBStickinessPoliciesResponse{ + Count: 1, + LBStickinessPolicies: []*cloudstack.LBStickinessPolicy{{Lbruleid: "rule-2", Stickinesspolicy: []cloudstack.LBStickinessPolicyStickinesspolicy{}}}, + } + + t.Run("policies are keyed by rule ID and rules without one are skipped", func(t *testing.T) { + ctrl := gomock.NewController(t) + t.Cleanup(ctrl.Finish) + mockLB := cloudstack.NewMockLoadBalancerServiceIface(ctrl) + mockLB.EXPECT().NewListLBStickinessPoliciesParams().Return(&cloudstack.ListLBStickinessPoliciesParams{}).Times(2) + mockLB.EXPECT().ListLBStickinessPolicies(gomock.Any()).DoAndReturn( + func(p *cloudstack.ListLBStickinessPoliciesParams) (*cloudstack.ListLBStickinessPoliciesResponse, error) { + if id, _ := p.GetLbruleid(); id == "rule-1" { + return withPolicy, nil + } + return emptyWrapper, nil + }).Times(2) + + lb := newLB(mockLB, + &cloudstack.LoadBalancerRule{Id: "rule-1", Name: "test-service-tcp-80"}, + &cloudstack.LoadBalancerRule{Id: "rule-2", Name: "test-service-tcp-443"}) + if err := lb.loadStickinessPolicies(); err != nil { + t.Fatalf("unexpected error: %v", err) + } + if len(lb.stickinessPolicies) != 1 { + t.Fatalf("stickinessPolicies = %v, want only rule-1", lb.stickinessPolicies) + } + if policy := lb.stickinessPolicies["rule-1"]; policy == nil || policy.Id != "policy-1" { + t.Errorf("policy for rule-1 = %v, want policy-1", policy) + } + }) + + t.Run("an empty response leaves the map empty", func(t *testing.T) { + ctrl := gomock.NewController(t) + t.Cleanup(ctrl.Finish) + mockLB := cloudstack.NewMockLoadBalancerServiceIface(ctrl) + mockLB.EXPECT().NewListLBStickinessPoliciesParams().Return(&cloudstack.ListLBStickinessPoliciesParams{}) + mockLB.EXPECT().ListLBStickinessPolicies(gomock.Any()).Return(&cloudstack.ListLBStickinessPoliciesResponse{}, nil) + + lb := newLB(mockLB, &cloudstack.LoadBalancerRule{Id: "rule-1", Name: "test-service-tcp-80"}) + if err := lb.loadStickinessPolicies(); err != nil { + t.Fatalf("unexpected error: %v", err) + } + if len(lb.stickinessPolicies) != 0 { + t.Errorf("stickinessPolicies = %v, want none", lb.stickinessPolicies) + } + }) + + t.Run("API error is returned", func(t *testing.T) { + ctrl := gomock.NewController(t) + t.Cleanup(ctrl.Finish) + mockLB := cloudstack.NewMockLoadBalancerServiceIface(ctrl) + mockLB.EXPECT().NewListLBStickinessPoliciesParams().Return(&cloudstack.ListLBStickinessPoliciesParams{}) + mockLB.EXPECT().ListLBStickinessPolicies(gomock.Any()).Return(nil, fmt.Errorf("list API error")) + + lb := newLB(mockLB, &cloudstack.LoadBalancerRule{Id: "rule-1", Name: "test-service-tcp-80"}) + err := lb.loadStickinessPolicies() + if err == nil || !strings.Contains(err.Error(), "list API error") { + t.Errorf("error = %v, want the list failure", err) + } + }) +} + +func TestSyncRuleHosts(t *testing.T) { + rule := &cloudstack.LoadBalancerRule{Id: "rule-123", Name: "test-service-tcp-80"} + instances := func(ids ...string) *cloudstack.ListLoadBalancerRuleInstancesResponse { + resp := &cloudstack.ListLoadBalancerRuleInstancesResponse{Count: len(ids)} + for _, id := range ids { + resp.LoadBalancerRuleInstances = append(resp.LoadBalancerRuleInstances, &cloudstack.VirtualMachine{Id: id}) + } + return resp + } + newLB := func(t *testing.T, current *cloudstack.ListLoadBalancerRuleInstancesResponse, listErr error, hostIDs ...string) (*loadBalancer, *cloudstack.MockLoadBalancerServiceIface) { + ctrl := gomock.NewController(t) + t.Cleanup(ctrl.Finish) + mockLB := cloudstack.NewMockLoadBalancerServiceIface(ctrl) + mockLB.EXPECT().NewListLoadBalancerRuleInstancesParams("rule-123").Return(&cloudstack.ListLoadBalancerRuleInstancesParams{}) + mockLB.EXPECT().ListLoadBalancerRuleInstances(gomock.Any()).Return(current, listErr) + lb := &loadBalancer{CloudStackClient: &cloudstack.CloudStackClient{LoadBalancer: mockLB}, hostIDs: hostIDs} + return lb, mockLB + } + + t.Run("missing hosts are assigned", func(t *testing.T) { + lb, mockLB := newLB(t, instances(), nil, "vm-1", "vm-2") + mockLB.EXPECT().NewAssignToLoadBalancerRuleParams("rule-123").Return(&cloudstack.AssignToLoadBalancerRuleParams{}) + mockLB.EXPECT().AssignToLoadBalancerRule(gomock.Any()).Return(&cloudstack.AssignToLoadBalancerRuleResponse{}, nil) + if err := lb.syncRuleHosts(rule); err != nil { + t.Fatalf("unexpected error: %v", err) + } + }) + + t.Run("stale hosts are removed", func(t *testing.T) { + lb, mockLB := newLB(t, instances("vm-1", "vm-gone"), nil, "vm-1") + mockLB.EXPECT().NewRemoveFromLoadBalancerRuleParams("rule-123").Return(&cloudstack.RemoveFromLoadBalancerRuleParams{}) + mockLB.EXPECT().RemoveFromLoadBalancerRule(gomock.Any()).Return(&cloudstack.RemoveFromLoadBalancerRuleResponse{}, nil) + if err := lb.syncRuleHosts(rule); err != nil { + t.Fatalf("unexpected error: %v", err) + } + }) + + t.Run("a rule already in sync is left alone", func(t *testing.T) { + lb, _ := newLB(t, instances("vm-1"), nil, "vm-1") + if err := lb.syncRuleHosts(rule); err != nil { + t.Fatalf("unexpected error: %v", err) + } + }) + + t.Run("list error is returned", func(t *testing.T) { + lb, _ := newLB(t, nil, fmt.Errorf("list API error"), "vm-1") + err := lb.syncRuleHosts(rule) + if err == nil || !strings.Contains(err.Error(), "list API error") { + t.Errorf("error = %v, want the list failure", err) + } + }) +} + +// Replacing a policy deletes the old one first, so a rejected replacement has to +// put the original back instead of leaving the live rule without stickiness. +func TestReplaceStickinessPolicy(t *testing.T) { + lbRule := &cloudstack.LoadBalancerRule{Id: "rule-123", Name: "test-service-tcp-80"} + existing := &cloudstack.LBStickinessPolicyStickinesspolicy{ + Id: "policy-old", Name: "test-service-tcp-80", Methodname: "LbCookie", + Params: map[string]string{"cookie-name": "SERVERID"}, + } + newLB := func(t *testing.T) (*loadBalancer, *cloudstack.MockLoadBalancerServiceIface) { + ctrl := gomock.NewController(t) + t.Cleanup(ctrl.Finish) + mockLB := cloudstack.NewMockLoadBalancerServiceIface(ctrl) + mockLB.EXPECT().NewDeleteLBStickinessPolicyParams("policy-old").Return(&cloudstack.DeleteLBStickinessPolicyParams{}) + mockLB.EXPECT().DeleteLBStickinessPolicy(gomock.Any()).Return(&cloudstack.DeleteLBStickinessPolicyResponse{}, nil) + lb := &loadBalancer{ + CloudStackClient: &cloudstack.CloudStackClient{LoadBalancer: mockLB}, + stickinessPolicies: map[string]*cloudstack.LBStickinessPolicyStickinesspolicy{"rule-123": existing}, + } + return lb, mockLB + } + + t.Run("a rejected replacement restores the previous policy", func(t *testing.T) { + lb, mockLB := newLB(t) + mockLB.EXPECT().NewCreateLBStickinessPolicyParams("rule-123", "NoSuchMethod", "test-service-tcp-80").Return(&cloudstack.CreateLBStickinessPolicyParams{}) + mockLB.EXPECT().CreateLBStickinessPolicy(gomock.Any()).Return(nil, fmt.Errorf("Failed to match Stickiness method name")) + mockLB.EXPECT().NewCreateLBStickinessPolicyParams("rule-123", "LbCookie", "test-service-tcp-80").Return(&cloudstack.CreateLBStickinessPolicyParams{}) + mockLB.EXPECT().CreateLBStickinessPolicy(gomock.Any()).Return(&cloudstack.CreateLBStickinessPolicyResponse{ + Stickinesspolicy: []cloudstack.CreateLBStickinessPolicyResponseStickinesspolicy{{Id: "policy-old"}}, + }, nil) + + err := lb.replaceStickinessPolicy("test-service-tcp-80", lbRule, existing, stickinessTestService("NoSuchMethod", "")) + if err == nil || !strings.Contains(err.Error(), "Failed to match Stickiness method name") { + t.Errorf("error = %v, want the rejection that triggered the rollback", err) + } + if strings.Contains(err.Error(), "restoring the previous stickiness policy failed") { + t.Errorf("error = %v, want it to report a successful restore", err) + } + }) + + t.Run("a failed restore is reported alongside the cause", func(t *testing.T) { + lb, mockLB := newLB(t) + mockLB.EXPECT().NewCreateLBStickinessPolicyParams("rule-123", "NoSuchMethod", "test-service-tcp-80").Return(&cloudstack.CreateLBStickinessPolicyParams{}) + mockLB.EXPECT().CreateLBStickinessPolicy(gomock.Any()).Return(nil, fmt.Errorf("method rejected")) + mockLB.EXPECT().NewCreateLBStickinessPolicyParams("rule-123", "LbCookie", "test-service-tcp-80").Return(&cloudstack.CreateLBStickinessPolicyParams{}) + mockLB.EXPECT().CreateLBStickinessPolicy(gomock.Any()).Return(nil, fmt.Errorf("restore API down")) + + err := lb.replaceStickinessPolicy("test-service-tcp-80", lbRule, existing, stickinessTestService("NoSuchMethod", "")) + if err == nil || !strings.Contains(err.Error(), "method rejected") || !strings.Contains(err.Error(), "restore API down") { + t.Errorf("error = %v, want both the cause and the restore failure", err) + } + }) + + t.Run("removing the annotation deletes the policy without creating one", func(t *testing.T) { + lb, _ := newLB(t) + if err := lb.replaceStickinessPolicy("test-service-tcp-80", lbRule, existing, &corev1.Service{}); err != nil { + t.Fatalf("unexpected error: %v", err) + } + }) +} + +// A method name that differs only in case names the same CloudStack method, so +// the policy must be left alone instead of being recreated on every sync. +func TestCheckStickinessPolicyIgnoresMethodNameCase(t *testing.T) { + lbRule := &cloudstack.LoadBalancerRule{Id: "rule-123", Name: "test-service-tcp-80"} + lb := &loadBalancer{stickinessPolicies: map[string]*cloudstack.LBStickinessPolicyStickinesspolicy{ + "rule-123": {Id: "policy-1", Methodname: "LbCookie", Params: map[string]string{"cookie-name": "SERVERID"}}, + }} + + policy, needsUpdate := lb.checkStickinessPolicy(lbRule, stickinessTestService("lbcookie", "cookie-name=SERVERID")) + if needsUpdate { + t.Errorf("needsUpdate = true for a method name differing only in case, want false") + } + if policy == nil || policy.Id != "policy-1" { + t.Errorf("policy = %v, want the existing policy-1", policy) + } +} + +// A rule that cannot be fully configured is removed again, so the next sync +// recreates it instead of adopting a rule that never gets its hosts assigned. +func TestCreateConfiguredLoadBalancerRuleRollsBackOnFailure(t *testing.T) { + port := corev1.ServicePort{Port: 80, NodePort: 30000, Protocol: corev1.ProtocolTCP} + created := &cloudstack.CreateLoadBalancerRuleResponse{Id: "rule-123", Name: "test-service-tcp-80"} + policyCreated := &cloudstack.CreateLBStickinessPolicyResponse{ + Stickinesspolicy: []cloudstack.CreateLBStickinessPolicyResponseStickinesspolicy{{Id: "policy-1", Methodname: "LbCookie"}}, + } + + newLB := func(t *testing.T) (*loadBalancer, *cloudstack.MockLoadBalancerServiceIface) { + ctrl := gomock.NewController(t) + t.Cleanup(ctrl.Finish) + mockLB := cloudstack.NewMockLoadBalancerServiceIface(ctrl) + mockLB.EXPECT().NewCreateLoadBalancerRuleParams("roundrobin", "test-service-tcp-80", 30000, 80).Return(&cloudstack.CreateLoadBalancerRuleParams{}) + mockLB.EXPECT().CreateLoadBalancerRule(gomock.Any()).Return(created, nil) + mockLB.EXPECT().NewCreateLBStickinessPolicyParams("rule-123", "LbCookie", "test-service-tcp-80").Return(&cloudstack.CreateLBStickinessPolicyParams{}) + lb := &loadBalancer{ + CloudStackClient: &cloudstack.CloudStackClient{LoadBalancer: mockLB}, + algorithm: "roundrobin", + hostIDs: []string{"vm-1"}, + rules: map[string]*cloudstack.LoadBalancerRule{}, + } + return lb, mockLB + } + expectRollback := func(mockLB *cloudstack.MockLoadBalancerServiceIface, deleteErr error) { + mockLB.EXPECT().NewDeleteLoadBalancerRuleParams("rule-123").Return(&cloudstack.DeleteLoadBalancerRuleParams{}) + mockLB.EXPECT().DeleteLoadBalancerRule(gomock.Any()).Return(&cloudstack.DeleteLoadBalancerRuleResponse{}, deleteErr) + } + + t.Run("policy creation failure deletes the rule", func(t *testing.T) { + lb, mockLB := newLB(t) + mockLB.EXPECT().CreateLBStickinessPolicy(gomock.Any()).Return(nil, fmt.Errorf("Failed to match Stickiness method name")) + expectRollback(mockLB, nil) + + rule, err := lb.createConfiguredLoadBalancerRule("test-service-tcp-80", port, LoadBalancerProtocolTCP, stickinessTestService("LbCookie", "")) + if rule != nil { + t.Errorf("rule = %v, want nil", rule) + } + if err == nil || !strings.Contains(err.Error(), "Failed to match Stickiness method name") { + t.Errorf("error = %v, want the policy creation failure", err) + } + }) + + t.Run("a failed rollback is reported alongside the cause", func(t *testing.T) { + lb, mockLB := newLB(t) + mockLB.EXPECT().CreateLBStickinessPolicy(gomock.Any()).Return(nil, fmt.Errorf("policy rejected")) + expectRollback(mockLB, fmt.Errorf("delete API down")) + + _, err := lb.createConfiguredLoadBalancerRule("test-service-tcp-80", port, LoadBalancerProtocolTCP, stickinessTestService("LbCookie", "")) + if err == nil || !strings.Contains(err.Error(), "policy rejected") || !strings.Contains(err.Error(), "delete API down") { + t.Errorf("error = %v, want both the cause and the rollback failure", err) + } + }) + + t.Run("host assignment failure deletes the rule", func(t *testing.T) { + lb, mockLB := newLB(t) + mockLB.EXPECT().CreateLBStickinessPolicy(gomock.Any()).Return(policyCreated, nil) + mockLB.EXPECT().NewAssignToLoadBalancerRuleParams("rule-123").Return(&cloudstack.AssignToLoadBalancerRuleParams{}) + mockLB.EXPECT().AssignToLoadBalancerRule(gomock.Any()).Return(nil, fmt.Errorf("assign API error")) + expectRollback(mockLB, nil) + + rule, err := lb.createConfiguredLoadBalancerRule("test-service-tcp-80", port, LoadBalancerProtocolTCP, stickinessTestService("LbCookie", "")) + if rule != nil || err == nil || !strings.Contains(err.Error(), "assign API error") { + t.Errorf("rule, error = %v, %v; want nil and the assignment failure", rule, err) + } + }) + + t.Run("a fully configured rule is kept", func(t *testing.T) { + lb, mockLB := newLB(t) + mockLB.EXPECT().CreateLBStickinessPolicy(gomock.Any()).Return(policyCreated, nil) + mockLB.EXPECT().NewAssignToLoadBalancerRuleParams("rule-123").Return(&cloudstack.AssignToLoadBalancerRuleParams{}) + mockLB.EXPECT().AssignToLoadBalancerRule(gomock.Any()).Return(&cloudstack.AssignToLoadBalancerRuleResponse{}, nil) + + rule, err := lb.createConfiguredLoadBalancerRule("test-service-tcp-80", port, LoadBalancerProtocolTCP, stickinessTestService("LbCookie", "")) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if rule == nil || rule.Id != "rule-123" { + t.Errorf("rule = %v, want rule-123", rule) + } + }) +} + +func TestCreateStickinessPolicyEmptyResponse(t *testing.T) { + ctrl := gomock.NewController(t) + t.Cleanup(ctrl.Finish) + + mockLB := cloudstack.NewMockLoadBalancerServiceIface(ctrl) + createParams := &cloudstack.CreateLBStickinessPolicyParams{} + createResp := &cloudstack.CreateLBStickinessPolicyResponse{ + Lbruleid: "rule-id", + Stickinesspolicy: []cloudstack.CreateLBStickinessPolicyResponseStickinesspolicy{}, + } + + gomock.InOrder( + mockLB.EXPECT().NewCreateLBStickinessPolicyParams("rule-id", "LbCookie", "test-service-tcp-80").Return(createParams), + mockLB.EXPECT().CreateLBStickinessPolicy(createParams).Return(createResp, nil), + ) + + lb := &loadBalancer{ + CloudStackClient: &cloudstack.CloudStackClient{ + LoadBalancer: mockLB, + }, + } + + policy, err := lb.createStickinessPolicy("test-service-tcp-80", "rule-id", stickinessTestService("LbCookie", "")) + if err == nil { + t.Fatal("expected an error, got nil") + } + if policy != nil { + t.Errorf("policy = %v, want nil", policy) + } + if !strings.Contains(err.Error(), "no policy returned") { + t.Errorf("error = %q, want it to mention that no policy was returned", err.Error()) + } +} diff --git a/go.mod b/go.mod index 24e177ad..be0b8af9 100644 --- a/go.mod +++ b/go.mod @@ -3,13 +3,14 @@ module github.com/apache/cloudstack-kubernetes-provider go 1.23.0 require ( - github.com/apache/cloudstack-go/v2 v2.19.0 + github.com/apache/cloudstack-go/v2 v2.19.1 github.com/blang/semver/v4 v4.0.0 github.com/spf13/pflag v1.0.5 go.uber.org/mock v0.5.0 gopkg.in/gcfg.v1 v1.2.3 k8s.io/api v0.24.17 k8s.io/apimachinery v0.24.17 + k8s.io/client-go v0.24.17 k8s.io/cloud-provider v0.24.17 k8s.io/component-base v0.24.17 k8s.io/klog/v2 v2.80.1 @@ -95,7 +96,6 @@ require ( gopkg.in/yaml.v2 v2.4.0 // indirect gopkg.in/yaml.v3 v3.0.1 // indirect k8s.io/apiserver v0.24.17 // indirect - k8s.io/client-go v0.24.17 // indirect k8s.io/component-helpers v0.24.17 // indirect k8s.io/controller-manager v0.24.17 // indirect k8s.io/kube-openapi v0.0.0-20220328201542-3ee0da9b0b42 // indirect diff --git a/go.sum b/go.sum index 85fd3e0f..423e8ea9 100644 --- a/go.sum +++ b/go.sum @@ -52,8 +52,8 @@ github.com/alecthomas/units v0.0.0-20151022065526-2efee857e7cf/go.mod h1:ybxpYRF github.com/alecthomas/units v0.0.0-20190717042225-c3de453c63f4/go.mod h1:ybxpYRFXyAe+OPACYpWeL0wqObRcbAqCMya13uyzqw0= github.com/alecthomas/units v0.0.0-20190924025748-f65c72e2690d/go.mod h1:rBZYJk541a8SKzHPHnH3zbiI+7dagKZ0cgpgrD7Fyho= github.com/antihax/optional v1.0.0/go.mod h1:uupD/76wgC+ih3iEmQUL+0Ugr19nfwCT1kdvxnR2qWY= -github.com/apache/cloudstack-go/v2 v2.19.0 h1:YHLw770MmgiqXx6NRFYw2Nr7DpnylLhLG2KYNCftgnc= -github.com/apache/cloudstack-go/v2 v2.19.0/go.mod h1:p/YBUwIEkQN6CQxFhw8Ff0wzf1MY0qRRRuGYNbcb1F8= +github.com/apache/cloudstack-go/v2 v2.19.1 h1:1K5O4NZpdWzOZUN6XuaVNsdX+QnoFRc5VE/oc1nUckQ= +github.com/apache/cloudstack-go/v2 v2.19.1/go.mod h1:p/YBUwIEkQN6CQxFhw8Ff0wzf1MY0qRRRuGYNbcb1F8= github.com/asaskevich/govalidator v0.0.0-20190424111038-f61b66f89f4a/go.mod h1:lB+ZfQJz7igIIfQNfa7Ml4HSf2uFQQRzpGGRXenZAgY= github.com/benbjohnson/clock v1.0.3/go.mod h1:bGMdMPoPVvcYyt1gHDf4J2KE153Yf9BuiUKYMaxlTDM= github.com/benbjohnson/clock v1.1.0 h1:Q92kusRqC1XV2MjkWETPvjJVqKetz1OzxZB7mHJLju8= diff --git a/test/e2e/annotations_test.go b/test/e2e/annotations_test.go index 05548522..685bb4e3 100644 --- a/test/e2e/annotations_test.go +++ b/test/e2e/annotations_test.go @@ -33,9 +33,11 @@ import ( ) const ( - annotationSourceCidrs = "service.beta.kubernetes.io/cloudstack-load-balancer-source-cidrs" - annotationHostname = "service.beta.kubernetes.io/cloudstack-load-balancer-hostname" - annotationIPAssociated = "service.beta.kubernetes.io/cloudstack-load-balancer-ip-associated-by-controller" //nolint:gosec + annotationSourceCidrs = "service.beta.kubernetes.io/cloudstack-load-balancer-source-cidrs" + annotationHostname = "service.beta.kubernetes.io/cloudstack-load-balancer-hostname" + annotationIPAssociated = "service.beta.kubernetes.io/cloudstack-load-balancer-ip-associated-by-controller" //nolint:gosec + annotationStickinessMethod = "service.beta.kubernetes.io/cloudstack-load-balancer-stickiness-method-name" + annotationStickinessParams = "service.beta.kubernetes.io/cloudstack-load-balancer-stickiness-method-param" ) func TestAnnot_SourceCIDRs(t *testing.T) { @@ -133,6 +135,168 @@ func TestAnnot_SessionAffinity(t *testing.T) { }) } +func TestAnnot_Stickiness(t *testing.T) { + f := NewFramework(t) + svc := f.CreateLBService(func(s *corev1.Service) { + s.Annotations = map[string]string{ + annotationStickinessMethod: "LbCookie", + annotationStickinessParams: "cookie-name=SERVERID", + } + }) + lbName := defaultLoadBalancerName(svc) + + f.WaitForIngressIP(svc) + rules := f.WaitForLBRules(lbName, 1) + ruleID := rules[0].Id + policy := f.WaitForStickinessPolicy(ruleID, "LbCookie", map[string]string{"cookie-name": "SERVERID"}) + + // A parameter change replaces the policy instead of editing it in place. + f.UpdateService(svc, func(s *corev1.Service) { + s.Annotations[annotationStickinessParams] = "cookie-name=JSESSIONID" + }) + replaced := f.WaitForStickinessPolicy(ruleID, "LbCookie", map[string]string{"cookie-name": "JSESSIONID"}) + if replaced.Id == policy.Id { + t.Errorf("policy %s was kept across a parameter change, want it recreated", policy.Id) + } + + // Dropping the annotations removes the policy but keeps the rule. + f.UpdateService(svc, func(s *corev1.Service) { + delete(s.Annotations, annotationStickinessMethod) + delete(s.Annotations, annotationStickinessParams) + }) + f.Eventually(lbSyncTimeout, lbSyncInterval, "stickiness policy to be removed", + func() (bool, error) { + current, err := f.StickinessPolicy(ruleID) + if err != nil { + return false, err + } + return current == nil, nil + }) + current := f.WaitForLBRules(lbName, 1) + if current[0].Id != ruleID { + t.Errorf("rule ID changed %s -> %s, want the rule to survive policy removal", ruleID, current[0].Id) + } +} + +// Replacing a policy has to delete the old one first, so a rejected replacement +// must leave the live rule with the stickiness it already had. +func TestAnnot_StickinessRejectedReplacement(t *testing.T) { + f := NewFramework(t) + svc := f.CreateLBService(func(s *corev1.Service) { + s.Annotations = map[string]string{ + annotationStickinessMethod: "LbCookie", + annotationStickinessParams: "cookie-name=SERVERID", + } + }) + lbName := defaultLoadBalancerName(svc) + + f.WaitForIngressIP(svc) + rules := f.WaitForLBRules(lbName, 1) + ruleID := rules[0].Id + f.WaitForStickinessPolicy(ruleID, "LbCookie", map[string]string{"cookie-name": "SERVERID"}) + + f.UpdateService(svc, func(s *corev1.Service) { + s.Annotations[annotationStickinessMethod] = "NoSuchMethod" + }) + f.Eventually(lbSyncTimeout, lbSyncInterval, "the sync to fail on the rejected replacement", + func() (bool, error) { + events, err := f.K8s.CoreV1().Events(svc.Namespace).List(context.Background(), metav1.ListOptions{ + FieldSelector: "involvedObject.name=" + svc.Name + ",reason=SyncLoadBalancerFailed", + }) + if err != nil { + return false, err + } + for _, event := range events.Items { + if strings.Contains(event.Message, "stickiness") { + return true, nil + } + } + return false, nil + }) + + // The rollback put the original policy back, so the rule never goes unprotected. + policy, err := f.StickinessPolicy(ruleID) + if err != nil { + t.Fatalf("reading the stickiness policy: %v", err) + } + if policy == nil { + t.Fatal("rule lost its stickiness policy after a rejected replacement") + } + if !strings.EqualFold(policy.Methodname, "LbCookie") { + t.Errorf("policy method = %q, want the original LbCookie", policy.Methodname) + } + + f.UpdateService(svc, func(s *corev1.Service) { + s.Annotations[annotationStickinessMethod] = "SourceBased" + s.Annotations[annotationStickinessParams] = "tablesize=200k" + }) + f.WaitForStickinessPolicy(ruleID, "SourceBased", map[string]string{"tablesize": "200k"}) +} + +// A rule whose stickiness policy CloudStack rejects must not survive as a +// half-configured rule: once the annotation is corrected, every rule has to end +// up with both its policy and its backend hosts. +func TestAnnot_StickinessInvalidMethod(t *testing.T) { + f := NewFramework(t) + svc := f.CreateLBService(nil) + lbName := defaultLoadBalancerName(svc) + + f.WaitForIngressIP(svc) + rules := f.WaitForLBRules(lbName, 1) + var wantHosts int + f.Eventually(lbSyncTimeout, lbSyncInterval, "hosts to be assigned to the first rule", + func() (bool, error) { + var err error + wantHosts, err = f.RuleInstanceCount(rules[0].Id) + return wantHosts > 0, err + }) + + // The new port is listed first so its rule is created, and its policy + // rejected, before the existing rule is reconciled. + f.UpdateService(svc, func(s *corev1.Service) { + s.Annotations = map[string]string{annotationStickinessMethod: "NoSuchMethod"} + s.Spec.Ports = []corev1.ServicePort{ + {Name: "alt", Port: 8080, Protocol: corev1.ProtocolTCP}, + {Name: "http", Port: 80, Protocol: corev1.ProtocolTCP}, + } + }) + f.Eventually(lbSyncTimeout, lbSyncInterval, "the sync to fail on the rejected stickiness method", + func() (bool, error) { + events, err := f.K8s.CoreV1().Events(svc.Namespace).List(context.Background(), metav1.ListOptions{ + FieldSelector: "involvedObject.name=" + svc.Name + ",reason=SyncLoadBalancerFailed", + }) + if err != nil { + return false, err + } + for _, event := range events.Items { + if strings.Contains(event.Message, "stickiness") { + return true, nil + } + } + return false, nil + }) + + f.UpdateService(svc, func(s *corev1.Service) { + s.Annotations[annotationStickinessMethod] = "LbCookie" + s.Annotations[annotationStickinessParams] = "cookie-name=SERVERID" + }) + rules = f.WaitForLBRules(lbName, 2) + for _, rule := range rules { + f.WaitForStickinessPolicy(rule.Id, "LbCookie", map[string]string{"cookie-name": "SERVERID"}) + f.Eventually(lbSyncTimeout, lbSyncInterval, fmt.Sprintf("rule %s to have %d hosts", rule.Name, wantHosts), + func() (bool, error) { + got, err := f.RuleInstanceCount(rule.Id) + if err != nil { + return false, err + } + if got != wantHosts { + return false, fmt.Errorf("rule %s has %d hosts", rule.Name, got) + } + return true, nil + }) + } +} + func TestAnnot_ExplicitLoadBalancerIP(t *testing.T) { f := NewFramework(t) diff --git a/test/e2e/framework.go b/test/e2e/framework.go index dbf7ecb9..045cde9e 100644 --- a/test/e2e/framework.go +++ b/test/e2e/framework.go @@ -40,6 +40,7 @@ import ( "encoding/hex" "fmt" "os" + "reflect" "strconv" "strings" "testing" @@ -364,6 +365,55 @@ func (f *Framework) WaitForLBRules(lbName string, want int) []*cloudstack.LoadBa return rules } +// StickinessPolicy returns the stickiness policy attached to a load balancer +// rule, or nil when the rule has none. CloudStack answers with a wrapper per +// rule whose inner list is empty for a rule without a policy. +func (f *Framework) StickinessPolicy(ruleID string) (*cloudstack.LBStickinessPolicyStickinesspolicy, error) { + p := f.CS.LoadBalancer.NewListLBStickinessPoliciesParams() + p.SetLbruleid(ruleID) + resp, err := f.CS.LoadBalancer.ListLBStickinessPolicies(p) + if err != nil { + return nil, err + } + for _, wrapper := range resp.LBStickinessPolicies { + if len(wrapper.Stickinesspolicy) > 0 { + return &wrapper.Stickinesspolicy[0], nil + } + } + return nil, nil +} + +// WaitForStickinessPolicy waits until the rule carries a stickiness policy +// with the given method and parameters and returns it. +func (f *Framework) WaitForStickinessPolicy(ruleID, method string, params map[string]string) *cloudstack.LBStickinessPolicyStickinesspolicy { + f.T.Helper() + var policy *cloudstack.LBStickinessPolicyStickinesspolicy + f.Eventually(lbSyncTimeout, lbSyncInterval, + fmt.Sprintf("rule %s to carry a %s stickiness policy with params %v", ruleID, method, params), + func() (bool, error) { + current, err := f.StickinessPolicy(ruleID) + if err != nil || current == nil { + return false, err + } + if current.Methodname != method || !reflect.DeepEqual(current.Params, params) { + return false, fmt.Errorf("policy is %s %v", current.Methodname, current.Params) + } + policy = current + return true, nil + }) + return policy +} + +// RuleInstanceCount returns how many VMs are assigned to a load balancer rule. +func (f *Framework) RuleInstanceCount(ruleID string) (int, error) { + p := f.CS.LoadBalancer.NewListLoadBalancerRuleInstancesParams(ruleID) + resp, err := f.CS.LoadBalancer.ListLoadBalancerRuleInstances(p) + if err != nil { + return 0, err + } + return len(resp.LoadBalancerRuleInstances), nil +} + // FirewallRules lists the firewall rules on a public IP. func (f *Framework) FirewallRules(publicIPID string) ([]*cloudstack.FirewallRule, error) { p := f.CS.Firewall.NewListFirewallRulesParams() diff --git a/test/e2e/loadbalancer_test.go b/test/e2e/loadbalancer_test.go index b4918b17..deaf33f1 100644 --- a/test/e2e/loadbalancer_test.go +++ b/test/e2e/loadbalancer_test.go @@ -161,6 +161,54 @@ func TestLB_NodeMembership(t *testing.T) { }) } +// A rule whose backend membership was lost, for instance because host +// assignment failed after the rule was created, is repaired by the next +// service sync rather than waiting for a node change. +func TestLB_HostsRestoredOnResync(t *testing.T) { + f := NewFramework(t) + svc := f.CreateLBService(nil) + lbName := defaultLoadBalancerName(svc) + + f.WaitForIngressIP(svc) + rules := f.WaitForLBRules(lbName, 1) + var assigned []string + f.Eventually(lbSyncTimeout, lbSyncInterval, "hosts to be assigned to the rule", + func() (bool, error) { + p := f.CS.LoadBalancer.NewListLoadBalancerRuleInstancesParams(rules[0].Id) + resp, err := f.CS.LoadBalancer.ListLoadBalancerRuleInstances(p) + if err != nil { + return false, err + } + assigned = assigned[:0] + for _, inst := range resp.LoadBalancerRuleInstances { + assigned = append(assigned, inst.Id) + } + return len(assigned) > 0, nil + }) + + remove := f.CS.LoadBalancer.NewRemoveFromLoadBalancerRuleParams(rules[0].Id) + remove.SetVirtualmachineids(assigned) + if _, err := f.CS.LoadBalancer.RemoveFromLoadBalancerRule(remove); err != nil { + t.Fatalf("removing hosts from rule %s: %v", rules[0].Id, err) + } + if got, err := f.RuleInstanceCount(rules[0].Id); err != nil || got != 0 { + t.Fatalf("rule still has %d hosts after removal (err %v)", got, err) + } + + forceReconcile(f, svc, lbName) + f.Eventually(lbSyncTimeout, lbSyncInterval, "the sync to restore the rule's hosts", + func() (bool, error) { + got, err := f.RuleInstanceCount(rules[0].Id) + if err != nil { + return false, err + } + if got != len(assigned) { + return false, fmt.Errorf("rule has %d hosts, want %d", got, len(assigned)) + } + return true, nil + }) +} + func TestLB_PortChange(t *testing.T) { f := NewFramework(t) svc := f.CreateLBService(nil)