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 ¶ms[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]