-
Notifications
You must be signed in to change notification settings - Fork 15
Stabilize Amphora serial LB tests: per-spec teardown and sourceRanges wait fix #320
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -28,6 +28,7 @@ import ( | |
| exutil "github.com/openshift/origin/test/extended/util" | ||
| ini "gopkg.in/ini.v1" | ||
| v1 "k8s.io/api/core/v1" | ||
| apierrors "k8s.io/apimachinery/pkg/api/errors" | ||
| metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" | ||
| "k8s.io/apimachinery/pkg/fields" | ||
| "k8s.io/apimachinery/pkg/types" | ||
|
|
@@ -125,6 +126,9 @@ var _ = g.Describe("[OTP][sig-installer][Suite:openshift/openstack][lb][Serial] | |
| svcPort := int32(8082) | ||
| jig := e2eservice.NewTestJig(clientSet, oc.Namespace(), svcName) | ||
| jig.Labels = labels | ||
| // Amphora LBs are costly (slots / max-shared-lb). Register before create so a | ||
| // timed-out ensure still tears down any partial Octavia LB for the next serial [lb] case. | ||
| registerAmphoraLoadBalancerTeardown(loadBalancerClient, clientSet, oc.Namespace(), svcName, lbProviderUnderTest) | ||
| svc, err := jig.CreateLoadBalancerService(ctx, loadBalancerServiceTimeout, func(svc *v1.Service) { | ||
| svc.Spec.Ports = []v1.ServicePort{{Protocol: protocolUnderTest, Port: svcPort, TargetPort: intstr.FromInt(8081)}} | ||
| svc.Spec.Selector = labels | ||
|
|
@@ -414,6 +418,7 @@ var _ = g.Describe("[OTP][sig-installer][Suite:openshift/openstack][lb][Serial] | |
| monitorMaxRetries := 2 | ||
| jig := e2eservice.NewTestJig(clientSet, oc.Namespace(), svcName) | ||
| jig.Labels = labels | ||
| registerAmphoraLoadBalancerTeardown(loadBalancerClient, clientSet, oc.Namespace(), svcName, lbProviderUnderTest) | ||
| svc, err := jig.CreateLoadBalancerService(ctx, loadBalancerServiceTimeout, func(svc *v1.Service) { | ||
| svc.Spec.Ports = []v1.ServicePort{{Protocol: protocolUnderTest, Port: svcPort, TargetPort: intstr.FromInt(8081)}} | ||
| svc.Spec.Selector = labels | ||
|
|
@@ -659,8 +664,10 @@ var _ = g.Describe("[OTP][sig-installer][Suite:openshift/openstack][lb][Serial] | |
| o.Expect(err).NotTo(o.HaveOccurred()) | ||
|
|
||
| g.By(fmt.Sprintf("Creating Openshift LoadBalancer Service with loadBalancerSourceRanges: '%s'", allowed_sourcerange)) | ||
| jig := e2eservice.NewTestJig(clientSet, oc.Namespace(), "udp-lb-sourceranges-svc") | ||
| svcName := "udp-lb-sourceranges-svc" | ||
| jig := e2eservice.NewTestJig(clientSet, oc.Namespace(), svcName) | ||
| jig.Labels = labels | ||
| registerAmphoraLoadBalancerTeardown(loadBalancerClient, clientSet, oc.Namespace(), svcName, lbProviderUnderTest) | ||
| svc, err := jig.CreateLoadBalancerService(ctx, loadBalancerServiceTimeout, func(svc *v1.Service) { | ||
| svc.Spec.Ports = []v1.ServicePort{{Protocol: v1.ProtocolUDP, Port: svcPort, TargetPort: intstr.FromInt(8081)}} | ||
| svc.Spec.Selector = labels | ||
|
|
@@ -687,8 +694,10 @@ var _ = g.Describe("[OTP][sig-installer][Suite:openshift/openstack][lb][Serial] | |
| listenerId := listeners[0].ID | ||
| e2e.Logf("LB listener ID: %s", listenerId) | ||
|
|
||
| // Check allowed_cidrs in the LB listener matches with the svc spec | ||
| if lbProviderUnderTest == "amphora" { // it only makes sense with Amphora | ||
| // Check allowed_cidrs in the LB listener matches with the svc spec. | ||
| // Provider loop uses "Amphora"/"OVN"; compare case-insensitively. | ||
| isAmphoraProvider := strings.EqualFold(lbProviderUnderTest, "amphora") | ||
| if isAmphoraProvider { | ||
| o.Expect(listeners[0].AllowedCIDRs).Should(o.ConsistOf([]string{allowed_sourcerange}), "Unexpected allowed_cidrs in Openstack LoadBalancer Listener '%s'", listenerId) | ||
| e2e.Logf("Found expected allowed_cidrs '%v' in Openstack LoadBalancer Listener '%s'", listeners[0].AllowedCIDRs, listenerId) | ||
| } | ||
|
|
@@ -719,18 +728,18 @@ var _ = g.Describe("[OTP][sig-installer][Suite:openshift/openstack][lb][Serial] | |
| e2e.Logf("Removed LoadBalancerSourceRanges spec from the service") | ||
|
|
||
| // Wait until allowed_cidrs in the LB listener is updated to 0.0.0.0/0 (all traffic allowed) | ||
| if lbProviderUnderTest == "amphora" { // it only makes sense with Amphora | ||
| if isAmphoraProvider { | ||
| allowAllAllowedCidrs := []string{"0.0.0.0/0"} | ||
| e2e.Logf("Expected allowed_cidrs: '%v'", allowAllAllowedCidrs) | ||
| e2e.Logf("Waiting for Amphora listener allowed_cidrs to become '%v' (was restricted to '%s')", allowAllAllowedCidrs, allowed_sourcerange) | ||
| o.Eventually(func() []string { | ||
| lbListener, err := octavialisteners.Get(ctx, loadBalancerClient, listenerId).Extract() | ||
| if err != nil { | ||
| e2e.Logf("Error ocurred: %v, trying next iteration", err) | ||
| return []string{} | ||
| } | ||
| e2e.Logf("Found AllowedCIDRs: %v", lbListener.AllowedCIDRs) | ||
| e2e.Logf("Found AllowedCIDRs: %v (ProvisioningStatus=%s OperatingStatus=%s)", lbListener.AllowedCIDRs, lbListener.ProvisioningStatus, lbListener.OperatingStatus) | ||
| return lbListener.AllowedCIDRs | ||
| }, "60s", "1s").Should(o.ConsistOf(allowAllAllowedCidrs), "Didn't find the expected allowed_cidrs '%v'", allowAllAllowedCidrs) | ||
| }, "120s", "2s").Should(o.ConsistOf(allowAllAllowedCidrs), "Didn't find the expected allowed_cidrs '%v'", allowAllAllowedCidrs) | ||
| e2e.Logf("Found expected allowed_cidrs '%v' in Openstack LoadBalancer Listener '%s'", allowAllAllowedCidrs, listenerId) | ||
| } else { | ||
| // It should be fast with OVN, but no good way to check it, let's just sleep. | ||
|
|
@@ -763,11 +772,13 @@ var _ = g.Describe("[OTP][sig-installer][Suite:openshift/openstack][lb][Serial] | |
| connNumber = 10 | ||
| g.By(fmt.Sprintf("accessing the service %d times from outside and storing the name of the pods answering", connNumber)) | ||
| results = make(map[string]int) | ||
| udpErrors := make([]string, 0, connNumber) | ||
| for i := 0; i < connNumber; i++ { | ||
| // https://github.com/kubernetes/kubernetes/blob/master/test/images/agnhost/README.md#netexec | ||
| podName, err := getPodNameThroughLb(svcIp, fmt.Sprintf("%d", svcPort), v1.ProtocolUDP, "hostname") | ||
| if err != nil { | ||
| e2e.Logf("Error detected while accessing the LoadBalancer service on try %d: %q", i, err) | ||
| udpErrors = append(udpErrors, err.Error()) | ||
| } else { | ||
| results[podName]++ | ||
| } | ||
|
|
@@ -782,6 +793,10 @@ var _ = g.Describe("[OTP][sig-installer][Suite:openshift/openstack][lb][Serial] | |
| // Check the number of successful connections with some margin (80%) to avoid flakes | ||
| var connMargin float32 = 0.8 | ||
| minSuccessConn := int(float32(connNumber) * connMargin) | ||
| if successConnCount < minSuccessConn { | ||
| // Dump Octavia/K8s state so we can tell ACL race vs datapath/member failure. | ||
| logUDPSourceRangesFailureDiagnostics(ctx, loadBalancerClient, clientSet, oc.Namespace(), svcName, loadBalancerId, listenerId, svcIp, svcPort, allowed_sourcerange, results, udpErrors) | ||
| } | ||
| o.Expect(successConnCount >= minSuccessConn).To(o.BeTrue(), "Found less successful connections (%d) than the minimum expected of '%d'", successConnCount, minSuccessConn) | ||
| e2e.Logf("Found expected number of successfull connections: '%d'", connNumber) | ||
| }) | ||
|
|
@@ -798,6 +813,154 @@ func skipIfNotLbProvider(expectedLbProvider string, ini *ini.File) { | |
| } | ||
| } | ||
|
|
||
| // registerAmphoraLoadBalancerTeardown registers best-effort per-spec cleanup for Amphora-backed | ||
| // LoadBalancer Services. Serial [lb] Amphora cases that leave Octavia LBs behind starve later | ||
| // specs (max-shared-lb / amphora capacity), which showed up as ensure-LB timeouts and UDP | ||
| // sourceRanges connectivity flakes in osp_verification. OVN variants are unchanged. | ||
| func registerAmphoraLoadBalancerTeardown(loadBalancerClient *gophercloud.ServiceClient, clientSet *kubernetes.Clientset, namespace, svcName, lbProviderUnderTest string) { | ||
| if !strings.EqualFold(lbProviderUnderTest, "Amphora") { | ||
| return | ||
| } | ||
| g.DeferCleanup(func(ctx context.Context) { | ||
| cleanupAmphoraLoadBalancerService(ctx, loadBalancerClient, clientSet, namespace, svcName) | ||
| }) | ||
| } | ||
|
|
||
| // cleanupAmphoraLoadBalancerService deletes the Service (so OCCM can release the LB) and, if the | ||
| // Octavia load balancer is still present, cascade-deletes it. Errors are logged only so teardown | ||
| // does not mask the original spec result. | ||
| func cleanupAmphoraLoadBalancerService(ctx context.Context, loadBalancerClient *gophercloud.ServiceClient, clientSet *kubernetes.Clientset, namespace, svcName string) { | ||
| g.By(fmt.Sprintf("Teardown: Amphora LoadBalancer leftovers for service %s/%s", namespace, svcName)) | ||
|
|
||
| lbID := "" | ||
| svc, err := clientSet.CoreV1().Services(namespace).Get(ctx, svcName, metav1.GetOptions{}) | ||
| if err == nil { | ||
| lbID = svc.GetAnnotations()["loadbalancer.openstack.org/load-balancer-id"] | ||
| if delErr := clientSet.CoreV1().Services(namespace).Delete(ctx, svcName, metav1.DeleteOptions{}); delErr != nil && !apierrors.IsNotFound(delErr) { | ||
| e2e.Logf("Teardown: error deleting service %s/%s: %v", namespace, svcName, delErr) | ||
| } | ||
| } else if !apierrors.IsNotFound(err) { | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I would guard for whatever error this function would return. So I think you can use o.Expect(err).NotTo(o.HaveOccurred(), fmt.Sprintf("Teardown: error getting service %s/%s: %v", namespace, svcName, err))This would give us a better error message. This should apply for other error references.
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. These errors stay as logs. DeferCleanup runs after the spec, so an Expect on Get or Delete would fail a passing spec on a transient API error, or cover the original timeout when the spec already failed. NotFound is expected when the Service never got a load balancer. The UDP diagnostic dump is log-only for the same reason: it runs after the connectivity assert has already failed, and an Expect there would stop the snapshot before sourceRanges_diag_verdict is printed. |
||
| e2e.Logf("Teardown: error getting service %s/%s: %v", namespace, svcName, err) | ||
| } | ||
|
|
||
| if lbID == "" { | ||
| e2e.Logf("Teardown: no load-balancer-id for %s/%s; nothing to force-delete in Octavia", namespace, svcName) | ||
| return | ||
| } | ||
|
|
||
| // Best-effort poll: a timeout is logged and does not fail the spec. | ||
| lbGone := false | ||
| pollErr := o.InterceptGomegaFailure(func() { | ||
| o.Eventually(func() bool { | ||
| _, getErr := octavialoadbalancers.Get(ctx, loadBalancerClient, lbID).Extract() | ||
| if getErr != nil { | ||
| if gophercloud.ResponseCodeIs(getErr, http.StatusNotFound) { | ||
| lbGone = true | ||
| return true | ||
| } | ||
| e2e.Logf("Teardown: error getting Octavia LB %s: %v", lbID, getErr) | ||
| } | ||
| delErr := octavialoadbalancers.Delete(ctx, loadBalancerClient, lbID, octavialoadbalancers.DeleteOpts{Cascade: true}).ExtractErr() | ||
| if delErr != nil { | ||
| if gophercloud.ResponseCodeIs(delErr, http.StatusNotFound) { | ||
| lbGone = true | ||
| return true | ||
| } | ||
| e2e.Logf("Teardown: cascade delete of Octavia LB %s: %v", lbID, delErr) | ||
| } | ||
| return false | ||
| }, "3m", "5s").Should(o.BeTrue()) | ||
| }) | ||
| if lbGone { | ||
| e2e.Logf("Teardown: Octavia LB %s is gone", lbID) | ||
| return | ||
| } | ||
| e2e.Logf("Teardown: timed out waiting for Octavia LB %s to disappear: %v", lbID, pollErr) | ||
| } | ||
|
|
||
| // logUDPSourceRangesFailureDiagnostics dumps Octavia + Service state when the post-open | ||
| // UDP connectivity assert fails. Use the printed "sourceRanges_diag_verdict" line to classify: | ||
| // - acl_not_allow_all: listener AllowedCIDRs still restricted → OCCM/Octavia ACL update bug or race | ||
| // - acl_allow_all_but_udp_blackhole: allow-all is set but UDP still times out → datapath/member/SG | ||
| // - unknown: could not read listener state | ||
| func logUDPSourceRangesFailureDiagnostics(ctx context.Context, loadBalancerClient *gophercloud.ServiceClient, clientSet *kubernetes.Clientset, namespace, svcName, loadBalancerId, listenerId, svcIp string, svcPort int32, previouslyRestrictedCIDR string, results map[string]int, udpErrors []string) { | ||
| g.By("Diagnostics: UDP sourceRanges connectivity failure snapshot") | ||
| e2e.Logf("sourceRanges_diag svc=%s/%s vip=%s:%d previouslyRestrictedCIDR=%s results=%v", namespace, svcName, svcIp, svcPort, previouslyRestrictedCIDR, results) | ||
| for i, errMsg := range udpErrors { | ||
| e2e.Logf("sourceRanges_diag udp_error[%d]=%s", i, errMsg) | ||
| } | ||
|
|
||
| svc, err := clientSet.CoreV1().Services(namespace).Get(ctx, svcName, metav1.GetOptions{}) | ||
| if err != nil { | ||
| e2e.Logf("sourceRanges_diag service_get_error=%v", err) | ||
| } else { | ||
| e2e.Logf("sourceRanges_diag service LoadBalancerSourceRanges=%v annotations[load-balancer-id]=%q", | ||
| svc.Spec.LoadBalancerSourceRanges, svc.GetAnnotations()["loadbalancer.openstack.org/load-balancer-id"]) | ||
| } | ||
|
|
||
| allowedCIDRs := []string{} | ||
| listenerProv, listenerOper := "", "" | ||
| lbListener, err := octavialisteners.Get(ctx, loadBalancerClient, listenerId).Extract() | ||
| if err != nil { | ||
| e2e.Logf("sourceRanges_diag listener_get_error=%v", err) | ||
| } else { | ||
| allowedCIDRs = lbListener.AllowedCIDRs | ||
| listenerProv = lbListener.ProvisioningStatus | ||
| listenerOper = lbListener.OperatingStatus | ||
| e2e.Logf("sourceRanges_diag listener id=%s AllowedCIDRs=%v ProvisioningStatus=%s OperatingStatus=%s", | ||
| listenerId, allowedCIDRs, listenerProv, listenerOper) | ||
| } | ||
|
|
||
| lb, err := octavialoadbalancers.Get(ctx, loadBalancerClient, loadBalancerId).Extract() | ||
| if err != nil { | ||
| e2e.Logf("sourceRanges_diag lb_get_error=%v", err) | ||
| } else { | ||
| e2e.Logf("sourceRanges_diag lb id=%s ProvisioningStatus=%s OperatingStatus=%s VipAddress=%s provider=%s", | ||
| loadBalancerId, lb.ProvisioningStatus, lb.OperatingStatus, lb.VipAddress, lb.Provider) | ||
| if len(lb.Pools) > 0 { | ||
| pool, poolErr := pools.Get(ctx, loadBalancerClient, lb.Pools[0].ID).Extract() | ||
| if poolErr != nil { | ||
| e2e.Logf("sourceRanges_diag pool_get_error=%v", poolErr) | ||
| } else { | ||
| e2e.Logf("sourceRanges_diag pool id=%s ProvisioningStatus=%s OperatingStatus=%s Protocol=%s LBMethod=%s", | ||
| pool.ID, pool.ProvisioningStatus, pool.OperatingStatus, pool.Protocol, pool.LBMethod) | ||
| allMembers, memErr := pools.ListMembers(loadBalancerClient, pool.ID, pools.ListMembersOpts{}).AllPages(ctx) | ||
| if memErr != nil { | ||
| e2e.Logf("sourceRanges_diag members_list_error=%v", memErr) | ||
| } else { | ||
| members, extErr := pools.ExtractMembers(allMembers) | ||
| if extErr != nil { | ||
| e2e.Logf("sourceRanges_diag members_extract_error=%v", extErr) | ||
| } else { | ||
| for _, m := range members { | ||
| e2e.Logf("sourceRanges_diag member id=%s address=%s protocol_port=%d ProvisioningStatus=%s OperatingStatus=%s weight=%d", | ||
| m.ID, m.Address, m.ProtocolPort, m.ProvisioningStatus, m.OperatingStatus, m.Weight) | ||
| } | ||
| } | ||
| } | ||
| } | ||
| } | ||
| } | ||
|
|
||
| verdict := "unknown" | ||
| allowAll := len(allowedCIDRs) == 1 && allowedCIDRs[0] == "0.0.0.0/0" | ||
| stillRestricted := false | ||
| for _, c := range allowedCIDRs { | ||
| if c == previouslyRestrictedCIDR { | ||
| stillRestricted = true | ||
| break | ||
| } | ||
| } | ||
| switch { | ||
| case stillRestricted || (len(allowedCIDRs) > 0 && !allowAll): | ||
| verdict = "acl_not_allow_all" | ||
| case allowAll: | ||
| verdict = "acl_allow_all_but_udp_blackhole" | ||
| } | ||
| e2e.Logf("sourceRanges_diag_verdict=%s (AllowedCIDRs=%v listenerProvisioning=%s listenerOperating=%s)", | ||
| verdict, allowedCIDRs, listenerProv, listenerOper) | ||
| } | ||
|
|
||
| // Return the FloatingIP assigned to a provided IP and return error if it is not found. | ||
| func getFipbyFixedIP(ctx context.Context, client *gophercloud.ServiceClient, vip string) (floatingips.FloatingIP, error) { | ||
| var result floatingips.FloatingIP | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Wouldn't be better to register the function within BeforeEach since it's conditionally cleans up when the driver is amphora?
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
The cleanup stays registered inside the three Its (create, ETP:Local+monitors, and UDP sourceRanges), and it is registered before CreateLoadBalancerService so a timed-out ensure still deletes a partial load balancer. BeforeEach would run for every spec in the provider/protocol loop, including OVN and the Amphora cases we are not changing. The service name is also different in each It. The Amphora check already lives in registerAmphoraLoadBalancerTeardown.