Skip to content

Commit 255f5de

Browse files
pedjakclaude
andauthored
fix: make demo e2e catalog queries resilient to transient failures (#2838)
The generate-demos CI job fails ~45% of the time with jq exit status 5 on the ClusterCatalog Quickstart scenario. Several issues contribute: - jq -s (slurp mode) buffers the entire operatorhubio FBC response in memory before processing, risking system errors on large catalogs - catalog content queries run exactly once with no retry, so any transient port-forward or network hiccup fails the step immediately - bash() does not attach stderr to ExitError, making failures opaque - with CatalogdHA, kubectl port-forward to the service deterministically picks the same pod via GetFirstPod sorting; if that pod is not the leader, it returns 404 (empty local cache) for every retry Remove jq slurp mode so each JSON object is processed in constant memory, prefixing filters with 'objects' to skip non-object values in the FBC stream. Wrap CatalogContainsSomePackages, PackageHasSomeChannels, and PackageHasSomeBundles in waitFor for retry on transient errors. Add curl --compressed to handle gzip-encoded responses and --fail with pipefail to detect HTTP errors. Resolve the catalogd leader pod via its Lease and port-forward directly to it on the container port (8443), falling back to the service when the lease cannot be read. Reset port-forwards on query failure and re-establish dead ones via liveness checks. Inject stderr into ExitError in bash() to match k8sClient diagnostics. Log catalog query errors at V(0) so CI timeout failures are diagnosable. Co-authored-by: Claude <noreply@anthropic.com>
1 parent 0a6b0be commit 255f5de

1 file changed

Lines changed: 119 additions & 42 deletions

File tree

test/e2e/steps/demo_steps.go

Lines changed: 119 additions & 42 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import (
55
"context"
66
"crypto/tls"
77
"encoding/json"
8+
"errors"
89
"fmt"
910
"net"
1011
"net/http"
@@ -41,6 +42,10 @@ func bash(ctx context.Context, script string) (string, error) {
4142

4243
if err != nil {
4344
logger.V(1).Info("Failed to run", "command", script, "stderr", stderr, "error", err)
45+
var exitErr *exec.ExitError
46+
if errors.As(err, &exitErr) {
47+
exitErr.Stderr = stderrBuf.Bytes()
48+
}
4449
}
4550
logger.V(1).Info("Output", "command", script, "output", stdout)
4651

@@ -81,79 +86,151 @@ func CatalogReportsConditionWithoutReason(ctx context.Context, catalogUserName,
8186
func ensureCatalogPortForward(ctx context.Context) (string, error) {
8287
sc := scenarioCtx(ctx)
8388
if sc.catalogAddr != "" {
84-
return sc.catalogAddr, nil
89+
if catalogPortForwardAlive(sc.catalogAddr) {
90+
return sc.catalogAddr, nil
91+
}
92+
logger.V(1).Info("Catalog port-forward is dead, re-establishing", "addr", sc.catalogAddr)
93+
resetCatalogPortForward(ctx)
94+
}
95+
96+
ns := componentNamespaces["catalogd"]
97+
target, err := catalogdLeaderPod(ctx, ns)
98+
port := int32(443)
99+
if err != nil {
100+
logger.V(1).Info("Could not resolve catalogd leader pod, falling back to service", "error", err)
101+
target = "service/catalogd-service"
102+
} else {
103+
port = 8443
85104
}
86105

87-
addr, cleanup, err := portForward(ctx, componentNamespaces["catalogd"], "service/catalogd-service", 443)
106+
addr, cleanup, err := portForward(ctx, ns, target, port)
88107
if err != nil {
89-
return "", fmt.Errorf("failed to start catalog port-forward: %w", err)
108+
return "", fmt.Errorf("failed to start catalog port-forward to %s: %w", target, err)
90109
}
91110
sc.catalogAddr = addr
92111
sc.catalogCleanup = cleanup
93112

94113
waitFor(ctx, func() bool {
95-
client := &http.Client{
96-
Timeout: 3 * time.Second,
97-
Transport: &http.Transport{
98-
TLSClientConfig: &tls.Config{InsecureSkipVerify: true}, //nolint:gosec
99-
DialContext: (&net.Dialer{Timeout: 2 * time.Second}).DialContext,
100-
},
101-
}
102-
resp, err := client.Get(fmt.Sprintf("https://%s/", addr))
103-
if err != nil {
104-
return false
105-
}
106-
resp.Body.Close()
107-
return true
114+
return catalogPortForwardAlive(addr)
108115
})
109116
return addr, nil
110117
}
111118

119+
func catalogdLeaderPod(ctx context.Context, ns string) (string, error) {
120+
holder, err := k8sClient(ctx, "get", "lease", "catalogd-operator-lock", "-n", ns,
121+
"-o", "jsonpath={.spec.holderIdentity}")
122+
if err != nil {
123+
return "", fmt.Errorf("failed to get catalogd leader lease: %w", err)
124+
}
125+
holder = strings.TrimSpace(holder)
126+
podName := holder
127+
if idx := strings.LastIndex(holder, "_"); idx >= 0 {
128+
podName = holder[:idx]
129+
}
130+
if podName == "" {
131+
return "", fmt.Errorf("catalogd leader lease has empty holderIdentity")
132+
}
133+
logger.Info("Resolved catalogd leader pod", "holder", holder, "pod", podName)
134+
return fmt.Sprintf("pod/%s", podName), nil
135+
}
136+
137+
func catalogPortForwardAlive(addr string) bool {
138+
client := &http.Client{
139+
Timeout: 3 * time.Second,
140+
Transport: &http.Transport{
141+
TLSClientConfig: &tls.Config{InsecureSkipVerify: true}, //nolint:gosec
142+
DialContext: (&net.Dialer{Timeout: 2 * time.Second}).DialContext,
143+
},
144+
}
145+
resp, err := client.Get(fmt.Sprintf("https://%s/", addr))
146+
if err != nil {
147+
return false
148+
}
149+
resp.Body.Close()
150+
return true
151+
}
152+
153+
// resetCatalogPortForward tears down the cached port-forward so the next
154+
// call to ensureCatalogPortForward establishes a fresh connection. With
155+
// CatalogdHA, non-leader pods return 404 (empty local cache); resetting
156+
// lets the next retry potentially reach the leader pod.
157+
func resetCatalogPortForward(ctx context.Context) {
158+
sc := scenarioCtx(ctx)
159+
if sc.catalogCleanup != nil {
160+
sc.catalogCleanup()
161+
}
162+
sc.catalogAddr = ""
163+
sc.catalogCleanup = nil
164+
}
165+
112166
func catalogCurlJq(ctx context.Context, catalogName, jqFilter string) (string, error) {
113167
addr, err := ensureCatalogPortForward(ctx)
114168
if err != nil {
115169
return "", err
116170
}
171+
// pipefail: propagate curl exit code through the pipe (e.g. HTTP 404 from a non-leader catalogd pod).
172+
// -sS: silent but show errors on stderr. -k: skip TLS verification for the port-forward.
173+
// --compressed: request gzip and stream-decompress (catalogd uses gzhttp); saves network for the large operatorhubio catalog.
174+
// --fail: exit 22 on HTTP errors so non-JSON error bodies don't reach jq.
117175
script := fmt.Sprintf(
118-
`curl -s -k https://%s/catalogs/%s/api/v1/all | jq -s '%s'`,
176+
`set -o pipefail; curl -sS -k --compressed --fail https://%s/catalogs/%s/api/v1/all | jq '%s'`,
119177
addr, catalogName, jqFilter,
120178
)
121-
return bash(ctx, script)
179+
out, err := bash(ctx, script)
180+
if err != nil {
181+
resetCatalogPortForward(ctx)
182+
}
183+
return out, err
122184
}
123185

124186
func CatalogContainsSomePackages(ctx context.Context, catalogName string) error {
125-
out, err := catalogCurlJq(ctx, catalogName,
126-
`.[] | select(.schema == "olm.package") | .name`)
127-
if err != nil {
128-
return err
129-
}
130-
if strings.TrimSpace(out) == "" {
131-
return fmt.Errorf("catalog %q contains no packages", catalogName)
132-
}
187+
waitFor(ctx, func() bool {
188+
out, err := catalogCurlJq(ctx, catalogName,
189+
`objects | select(.schema == "olm.package") | .name`)
190+
if err != nil {
191+
logger.Info("Catalog query failed, retrying", "catalog", catalogName, "error", err, "stderr", stderrOutput(err))
192+
return false
193+
}
194+
if strings.TrimSpace(out) == "" {
195+
logger.Info("Catalog returned no packages, retrying", "catalog", catalogName)
196+
return false
197+
}
198+
return true
199+
})
133200
return nil
134201
}
135202

136203
func PackageHasSomeChannels(ctx context.Context, packageName, catalogName string) error {
137-
out, err := catalogCurlJq(ctx, catalogName,
138-
fmt.Sprintf(`.[] | select(.schema == "olm.channel") | select(.package == "%s") | .name`, packageName))
139-
if err != nil {
140-
return err
141-
}
142-
if strings.TrimSpace(out) == "" {
143-
return fmt.Errorf("package %q in catalog %q has no channels", packageName, catalogName)
144-
}
204+
waitFor(ctx, func() bool {
205+
out, err := catalogCurlJq(ctx, catalogName,
206+
fmt.Sprintf(`objects | select(.schema == "olm.channel") | select(.package == "%s") | .name`, packageName))
207+
if err != nil {
208+
logger.Info("Catalog query failed, retrying", "catalog", catalogName, "package", packageName, "error", err, "stderr", stderrOutput(err))
209+
return false
210+
}
211+
if strings.TrimSpace(out) == "" {
212+
logger.Info("Package has no channels, retrying", "catalog", catalogName, "package", packageName)
213+
return false
214+
}
215+
return true
216+
})
145217
return nil
146218
}
147219

148220
func PackageHasSomeBundles(ctx context.Context, packageName, catalogName string) error {
149-
out, err := catalogCurlJq(ctx, catalogName,
150-
fmt.Sprintf(`.[] | select(.schema == "olm.bundle") | select(.package == "%s") | .name`, packageName))
151-
if err != nil {
152-
return err
153-
}
154-
if strings.TrimSpace(out) == "" {
155-
return fmt.Errorf("package %q in catalog %q has no bundles", packageName, catalogName)
156-
}
221+
waitFor(ctx, func() bool {
222+
out, err := catalogCurlJq(ctx, catalogName,
223+
fmt.Sprintf(`objects | select(.schema == "olm.bundle") | select(.package == "%s") | .name`, packageName))
224+
if err != nil {
225+
logger.Info("Catalog query failed, retrying", "catalog", catalogName, "package", packageName, "error", err, "stderr", stderrOutput(err))
226+
return false
227+
}
228+
if strings.TrimSpace(out) == "" {
229+
logger.Info("Package has no bundles, retrying", "catalog", catalogName, "package", packageName)
230+
return false
231+
}
232+
return true
233+
})
157234
return nil
158235
}
159236

0 commit comments

Comments
 (0)