Copilot commented on code in PR #88:
URL: 
https://github.com/apache/cloudstack-kubernetes-provider/pull/88#discussion_r4193613443


##########
cloudstack_loadbalancer.go:
##########
@@ -1038,31 +1073,453 @@ func (lb *loadBalancer) ensureLoadBalancerRule(d 
desiredLBRule, service *corev1.
        return lbRule, nil
 }
 
+// repairEmptyRule assigns the hosts to an adopted rule that has none, which 
happens when a create
+// fails before assigning them. Rules that have hosts are left to 
UpdateLoadBalancer, because node
+// sync runs at the same time as this sync and may have newer node data.
+func (lb *loadBalancer) repairEmptyRule(lbRule *cloudstack.LoadBalancerRule) 
error {
+       p := lb.LoadBalancer.NewListLoadBalancerRuleInstancesParams(lbRule.Id)
+       instances, err := lb.LoadBalancer.ListLoadBalancerRuleInstances(p)
+       if err != nil {
+               return fmt.Errorf("error retrieving associated instances: %v", 
err)
+       }
+       if instances.Count > 0 {
+               return nil
+       }
+
+       klog.V(4).Infof("Assigning hosts (%v) to load balancer rule %v, which 
has none", lb.hostIDs, lbRule.Name)
+       return lb.assignHostsToRule(lbRule, lb.hostIDs)
+}
+
 // applyLoadBalancerRules creates or updates the load balancer rule of every 
desired service
-// port and reconciles the firewall or network ACL rules it needs.
-func (lb *loadBalancer) applyLoadBalancerRules(desired []desiredLBRule, 
service *corev1.Service, network *cloudstack.Network, version semver.Version) 
error {
+// port and reconciles the firewall or network ACL rules it needs. It returns 
the rule of each
+// desired port, in the same order.
+func (lb *loadBalancer) applyLoadBalancerRules(desired []desiredLBRule, 
service *corev1.Service, network *cloudstack.Network, version semver.Version) 
([]*cloudstack.LoadBalancerRule, error) {
+       rules := make([]*cloudstack.LoadBalancerRule, 0, len(desired))
        for _, d := range desired {
                lbRule, err := lb.ensureLoadBalancerRule(d, service, version)
                if err != nil {
-                       return err
+                       return nil, err
                }
+               rules = append(rules, lbRule)
 
                if isFirewallSupported(network.Service) {
                        klog.V(4).Infof("Creating firewall rules for load 
balancer rule: %v (%v:%v:%v)", d.name, d.protocol, lbRule.Publicip, d.port.Port)
                        if _, err := lb.updateFirewallRule(lbRule.Publicipid, 
int(d.port.Port), d.protocol, service.Spec.LoadBalancerSourceRanges); err != 
nil {
-                               return err
+                               return nil, err
                        }
                } else if isNetworkACLSupported(network.Service) {
                        klog.V(4).Infof("Creating ACL rules for load balancer 
rule: %v (%v:%v:%v)", d.name, d.protocol, lbRule.Publicip, d.port.Port)
                        if _, err := lb.updateNetworkACL(int(d.port.Port), 
d.protocol, network.Id); err != nil {
-                               return err
+                               return nil, err
+                       }
+               }
+       }
+
+       return rules, nil
+}
+
+const (
+       // stickinessPolicyDescription marks the policies the controller 
creates. Only policies with
+       // this description are deleted, so policies added by hand are kept.
+       stickinessPolicyDescription = "Managed by the CloudStack Kubernetes 
Provider"
+
+       lbCookieMethod  = "LbCookie"
+       appCookieMethod = "AppCookie"
+
+       // stickinessParamForbidden lists the characters a parameter may not 
contain. CloudStack splits
+       // stored parameters on "=" and "&", and copies them unquoted into the 
HAProxy config, where
+       // quotes, backslashes and "#" have a special meaning.
+       stickinessParamForbidden = "=&\"'\\#"
+)
+
+// lbCookieModes are the cookie modes HAProxy accepts. CloudStack does not 
check the mode.
+var lbCookieModes = []string{"insert", "rewrite", "prefix"}
+
+// stickinessSpec is the stickiness policy a Service asks for. ports holds the 
selected TCP ports.
+type stickinessSpec struct {
+       method string
+       params map[string]string
+       ports  map[int32]bool
+}
+
+// stickinessMethod is a method in a network's SupportedStickinessMethods 
capability.
+type stickinessMethod struct {
+       Name   string                  `json:"methodname"`
+       Params []stickinessMethodParam `json:"paramlist"`
+}
+
+// stickinessMethodParam is a parameter of a stickiness method.
+type stickinessMethodParam struct {
+       Name     string `json:"paramname"`
+       Required bool   `json:"required"`
+       IsFlag   bool   `json:"isflag"`
+}
+
+// parseStickinessSpec reads the stickiness annotations. It returns nil when 
the Service asks for
+// no stickiness. It makes no CloudStack calls, so bad annotations fail before 
anything changes.
+func parseStickinessSpec(service *corev1.Service) (*stickinessSpec, error) {
+       method := strings.TrimSpace(getStringFromServiceAnnotation(service, 
ServiceAnnotationLoadBalancerStickinessMethodName, ""))
+       if method == "" {
+               return nil, nil
+       }
+       if strings.EqualFold(method, appCookieMethod) {
+               return nil, fmt.Errorf("unsupported stickiness method %s in 
annotation %s: HAProxy removed the appsession option it needs", method, 
ServiceAnnotationLoadBalancerStickinessMethodName)
+       }
+
+       params, err := 
parseStickinessParams(getStringFromServiceAnnotation(service, 
ServiceAnnotationLoadBalancerStickinessMethodParam, ""))
+       if err != nil {
+               return nil, err
+       }
+       ports, err := selectStickinessPorts(service, method)
+       if err != nil {
+               return nil, err
+       }
+
+       return &stickinessSpec{method: method, params: params, ports: ports}, 
nil
+}
+
+// parseStickinessParams parses the key=value list of the method-param 
annotation. Empty values,
+// spaces and the characters in stickinessParamForbidden are rejected, because 
CloudStack would
+// read them back differently or write them into a broken HAProxy config.
+func parseStickinessParams(value string) (map[string]string, error) {
+       params := make(map[string]string)
+       seen := make(map[string]bool)
+       for _, entry := range strings.Split(value, ",") {
+               entry = strings.TrimSpace(entry)
+               if entry == "" {
+                       continue
+               }
+               key, val, found := strings.Cut(entry, "=")
+               key, val = strings.TrimSpace(key), strings.TrimSpace(val)
+               switch {
+               case !found || key == "":
+                       return nil, stickinessParamError(entry, "expected 
key=value")
+               case val == "":
+                       return nil, stickinessParamError(entry, 
fmt.Sprintf("missing value, write flags as %s=true", key))
+               case !safeStickinessText(key) || !safeStickinessText(val):
+                       return nil, stickinessParamError(entry, `spaces and the 
characters = & " ' \ # are not allowed`)
+               }
+
+               lowered := strings.ToLower(key)
+               if seen[lowered] {
+                       return nil, stickinessParamError(entry, key+" is set 
twice")
+               }
+               seen[lowered] = true
+               params[key] = val
+       }
+
+       return params, nil
+}
+
+// stickinessParamError reports an invalid entry of the method-param 
annotation.
+func stickinessParamError(entry, reason string) error {
+       return fmt.Errorf("invalid stickiness parameter %q in annotation %s: 
%s", entry, ServiceAnnotationLoadBalancerStickinessMethodParam, reason)
+}
+
+// safeStickinessText reports whether s has no spaces and no forbidden 
characters.
+func safeStickinessText(s string) bool {
+       return !strings.ContainsAny(s, stickinessParamForbidden) && 
strings.IndexFunc(s, unicode.IsSpace) < 0
+}
+
+// selectStickinessPorts returns the TCP ports that get the policy. LbCookie 
switches its ports to
+// HTTP mode, which breaks TLS and other non-HTTP traffic, so it only goes on 
ports with
+// appProtocol: http. Other methods go on every TCP port.
+func selectStickinessPorts(service *corev1.Service, method string) 
(map[int32]bool, error) {
+       tcpPorts := make(map[int32]bool)
+       httpPorts := make(map[int32]bool)
+       for _, port := range service.Spec.Ports {
+               if port.Protocol != corev1.ProtocolTCP {
+                       continue
+               }
+               tcpPorts[port.Port] = true
+               if port.AppProtocol != nil && 
strings.EqualFold(*port.AppProtocol, "http") {
+                       httpPorts[port.Port] = true
+               }
+       }
+
+       if !strings.EqualFold(method, lbCookieMethod) {
+               return tcpPorts, nil
+       }
+       if len(httpPorts) == 0 {
+               return nil, fmt.Errorf("stickiness method %s only works on 
plain HTTP ports: set appProtocol: http on them", method)
+       }
+
+       return httpPorts, nil
+}
+
+// wants reports whether a Service port should get the policy.
+func (s *stickinessSpec) wants(port corev1.ServicePort) bool {
+       return s != nil && port.Protocol == corev1.ProtocolTCP && 
s.ports[port.Port]
+}
+
+// matches reports whether a policy is the controller's and matches the spec.
+func (s *stickinessSpec) matches(policy 
cloudstack.LBStickinessPolicyStickinesspolicy) bool {
+       return policy.Description == stickinessPolicyDescription &&
+               strings.EqualFold(policy.Methodname, s.method) &&
+               maps.Equal(policy.Params, s.params)
+}
+
+// supportedStickinessMethods returns the stickiness methods the network's 
load balancer
+// supports, or nil if the network does not list them.
+func supportedStickinessMethods(network *cloudstack.Network) 
[]stickinessMethod {
+       for _, svc := range network.Service {
+               if svc.Name != "Lb" {
+                       continue
+               }
+               for _, capability := range svc.Capability {
+                       if capability.Name != "SupportedStickinessMethods" {
+                               continue
+                       }
+                       var methods []stickinessMethod
+                       if err := json.Unmarshal([]byte(capability.Value), 
&methods); err != nil {
+                               klog.Warningf("Cannot read the stickiness 
methods of network %v: %v", network.Id, err)
+                               return nil
                        }
+                       return methods
                }
        }
 
        return nil
 }
 
+// normalise checks the spec against the methods the network supports, as 
CloudStack does, and
+// also checks the LbCookie mode, which CloudStack does not. It writes names 
and flags the way
+// CloudStack stores them, so a policy read back compares equal. It does 
nothing if the network
+// does not list its methods.
+func (s *stickinessSpec) normalise(methods []stickinessMethod) error {
+       if s == nil || len(methods) == 0 {
+               return nil
+       }
+       method := findStickinessMethod(methods, s.method)
+       if method == nil {
+               return fmt.Errorf("stickiness method %q in annotation %s is not 
supported on this network; supported methods: %s", s.method, 
ServiceAnnotationLoadBalancerStickinessMethodName, 
stickinessMethodNames(methods))
+       }
+
+       params := make(map[string]string, len(s.params))
+       for key, value := range s.params {
+               param := findStickinessParam(method.Params, key)
+               if param == nil {
+                       return fmt.Errorf("stickiness parameter %q in 
annotation %s is not supported by %s; supported parameters: %s", key, 
ServiceAnnotationLoadBalancerStickinessMethodParam, method.Name, 
stickinessParamNames(method.Params))
+               }
+               normalised, send, err := normaliseStickinessValue(method.Name, 
*param, value)
+               if err != nil {
+                       return err
+               }
+               if send {
+                       params[param.Name] = normalised
+               }
+       }
+       for _, param := range method.Params {
+               if _, set := params[param.Name]; param.Required && !set {
+                       return fmt.Errorf("stickiness method %s needs parameter 
%s in annotation %s", method.Name, param.Name, 
ServiceAnnotationLoadBalancerStickinessMethodParam)
+               }
+       }
+
+       s.method, s.params = method.Name, params
+       return nil
+}
+
+// normaliseStickinessValue returns the value to send for a parameter, and 
whether to send it.
+// HAProxy turns a flag on whenever it is set, whatever the value, so flags 
are parsed as booleans
+// and sent as true or left out, like the CloudStack UI does.
+func normaliseStickinessValue(method string, param stickinessMethodParam, 
value string) (string, bool, error) {
+       if param.IsFlag {
+               on, err := strconv.ParseBool(value)
+               if err != nil {
+                       return "", false, fmt.Errorf("stickiness flag %s in 
annotation %s must be true or false, not %q", param.Name, 
ServiceAnnotationLoadBalancerStickinessMethodParam, value)
+               }
+               return "true", on, nil
+       }
+       if strings.EqualFold(method, lbCookieMethod) && param.Name == "mode" {
+               mode := strings.ToLower(value)
+               if !slices.Contains(lbCookieModes, mode) {
+                       return "", false, fmt.Errorf("stickiness parameter mode 
in annotation %s must be one of %s, not %q", 
ServiceAnnotationLoadBalancerStickinessMethodParam, strings.Join(lbCookieModes, 
", "), value)
+               }
+               return mode, true, nil
+       }
+
+       return value, true, nil
+}
+
+// findStickinessMethod returns the method with the given name, ignoring case 
like CloudStack.
+func findStickinessMethod(methods []stickinessMethod, name string) 
*stickinessMethod {
+       for i := range methods {
+               if strings.EqualFold(methods[i].Name, name) {
+                       return &methods[i]
+               }
+       }
+
+       return nil
+}
+
+// findStickinessParam returns the parameter with the given name, ignoring 
case.
+func findStickinessParam(params []stickinessMethodParam, name string) 
*stickinessMethodParam {
+       for i := range params {
+               if strings.EqualFold(params[i].Name, name) {
+                       return &params[i]
+               }
+       }
+
+       return nil
+}
+
+// stickinessMethodNames lists the supported methods for an error message. 
AppCookie is left out
+// because the controller rejects it.
+func stickinessMethodNames(methods []stickinessMethod) string {
+       names := make([]string, 0, len(methods))
+       for _, method := range methods {
+               if !strings.EqualFold(method.Name, appCookieMethod) {
+                       names = append(names, method.Name)
+               }
+       }
+
+       return strings.Join(names, ", ")
+}
+
+// stickinessParamNames lists a method's parameters for an error message.
+func stickinessParamNames(params []stickinessMethodParam) string {
+       names := make([]string, 0, len(params))
+       for _, param := range params {
+               names = append(names, param.Name)
+       }
+
+       return strings.Join(names, ", ")
+}
+
+// applyStickinessPolicies updates the stickiness policy of every desired 
rule, once the rules are
+// in place. It tries every rule and returns all errors together.
+func (lb *loadBalancer) applyStickinessPolicies(desired []desiredLBRule, rules 
[]*cloudstack.LoadBalancerRule, spec *stickinessSpec) error {
+       var errs []error
+       for i, d := range desired {
+               if err := lb.applyStickinessPolicy(d, rules[i], spec); err != 
nil {
+                       errs = append(errs, fmt.Errorf("load balancer rule %v: 
%w", d.name, err))
+               }
+       }
+
+       return errors.Join(errs...)
+}
+
+// applyStickinessPolicy updates the stickiness policy of one rule. A rule 
created in this sync has
+// no policy yet, so it is not looked up. If the account may not call the 
stickiness APIs and the
+// rule needs no policy, nothing is done. The account CloudStack's Kubernetes 
service creates for
+// clusters in projects is one such account.
+func (lb *loadBalancer) applyStickinessPolicy(d desiredLBRule, lbRule 
*cloudstack.LoadBalancerRule, spec *stickinessSpec) error {
+       want := spec.wants(d.port)
+       var live []cloudstack.LBStickinessPolicyStickinesspolicy
+       if !d.createsRule() {
+               var err error
+               live, err = lb.liveStickinessPolicies(lbRule.Id)
+               switch {
+               case err != nil && isNotAllowed(err) && !want:
+                       klog.V(4).Infof("Skipping the stickiness policies of 
load balancer rule %v: %v", d.name, err)
+                       return nil
+               case err != nil && isNotAllowed(err):
+                       return fmt.Errorf("%w; stickiness needs the 
listLBStickinessPolicies, createLBStickinessPolicy and deleteLBStickinessPolicy 
APIs", err)
+               case err != nil:
+                       return err
+               }
+       }
+
+       create, stale, err := stickinessChanges(live, spec, want)
+       if err != nil {
+               return err
+       }
+       for _, id := range stale {
+               if err := lb.deleteStickinessPolicy(d.name, id); err != nil {
+                       return err

Review Comment:
   When a rule has both a controller-owned policy and a foreign policy, this 
deletes the current controller policy before attempting the replacement. If 
`createStickinessPolicy` then fails (for example because the router rejects the 
new parameters), the old policy is already gone and the rule is left with only 
the foreign policy, contradicting the documented guarantee that a rejected 
change leaves the current policy intact. Preserve the old policy until 
replacement is known to succeed, or perform a best-effort rollback/return 
without deleting it when CloudStack requires this precondition.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to