Skip to content

Commit def9ce2

Browse files
pedjakclaude
andcommitted
fix: make demo e2e catalog queries resilient to transient failures
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, 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 7ec6a0e commit def9ce2

1 file changed

Lines changed: 112 additions & 42 deletions

File tree

‎test/e2e/steps/demo_steps.go‎

Lines changed: 112 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,144 @@ 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+
if err != nil {
99+
logger.V(1).Info("Could not resolve catalogd leader pod, falling back to service", "error", err)
100+
target = "service/catalogd-service"
85101
}
86102

87-
addr, cleanup, err := portForward(ctx, componentNamespaces["catalogd"], "service/catalogd-service", 443)
103+
addr, cleanup, err := portForward(ctx, ns, target, 443)
88104
if err != nil {
89-
return "", fmt.Errorf("failed to start catalog port-forward: %w", err)
105+
return "", fmt.Errorf("failed to start catalog port-forward to %s: %w", target, err)
90106
}
91107
sc.catalogAddr = addr
92108
sc.catalogCleanup = cleanup
93109

94110
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
111+
return catalogPortForwardAlive(addr)
108112
})
109113
return addr, nil
110114
}
111115

116+
func catalogdLeaderPod(ctx context.Context, ns string) (string, error) {
117+
holder, err := k8sClient(ctx, "get", "lease", "catalogd-operator-lock", "-n", ns,
118+
"-o", "jsonpath={.spec.holderIdentity}")
119+
if err != nil {
120+
return "", fmt.Errorf("failed to get catalogd leader lease: %w", err)
121+
}
122+
holder = strings.TrimSpace(holder)
123+
podName := holder
124+
if idx := strings.LastIndex(holder, "_"); idx >= 0 {
125+
podName = holder[:idx]
126+
}
127+
if podName == "" {
128+
return "", fmt.Errorf("catalogd leader lease has empty holderIdentity")
129+
}
130+
logger.Info("Resolved catalogd leader pod", "holder", holder, "pod", podName)
131+
return fmt.Sprintf("pod/%s", podName), nil
132+
}
133+
134+
func catalogPortForwardAlive(addr string) bool {
135+
client := &http.Client{
136+
Timeout: 3 * time.Second,
137+
Transport: &http.Transport{
138+
TLSClientConfig: &tls.Config{InsecureSkipVerify: true}, //nolint:gosec
139+
DialContext: (&net.Dialer{Timeout: 2 * time.Second}).DialContext,
140+
},
141+
}
142+
resp, err := client.Get(fmt.Sprintf("https://%s/", addr))
143+
if err != nil {
144+
return false
145+
}
146+
resp.Body.Close()
147+
return true
148+
}
149+
150+
// resetCatalogPortForward tears down the cached port-forward so the next
151+
// call to ensureCatalogPortForward establishes a fresh connection. With
152+
// CatalogdHA, non-leader pods return 404 (empty local cache); resetting
153+
// lets the next retry potentially reach the leader pod.
154+
func resetCatalogPortForward(ctx context.Context) {
155+
sc := scenarioCtx(ctx)
156+
if sc.catalogCleanup != nil {
157+
sc.catalogCleanup()
158+
}
159+
sc.catalogAddr = ""
160+
sc.catalogCleanup = nil
161+
}
162+
112163
func catalogCurlJq(ctx context.Context, catalogName, jqFilter string) (string, error) {
113164
addr, err := ensureCatalogPortForward(ctx)
114165
if err != nil {
115166
return "", err
116167
}
117168
script := fmt.Sprintf(
118-
`curl -s -k https://%s/catalogs/%s/api/v1/all | jq -s '%s'`,
169+
`set -o pipefail; curl -sS -k --compressed --fail https://%s/catalogs/%s/api/v1/all | jq '%s'`,
119170
addr, catalogName, jqFilter,
120171
)
121-
return bash(ctx, script)
172+
out, err := bash(ctx, script)
173+
if err != nil {
174+
resetCatalogPortForward(ctx)
175+
}
176+
return out, err
122177
}
123178

124179
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-
}
180+
waitFor(ctx, func() bool {
181+
out, err := catalogCurlJq(ctx, catalogName,
182+
`objects | select(.schema == "olm.package") | .name`)
183+
if err != nil {
184+
logger.Info("Catalog query failed, retrying", "catalog", catalogName, "error", err, "stderr", stderrOutput(err))
185+
return false
186+
}
187+
if strings.TrimSpace(out) == "" {
188+
logger.Info("Catalog returned no packages, retrying", "catalog", catalogName)
189+
return false
190+
}
191+
return true
192+
})
133193
return nil
134194
}
135195

136196
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-
}
197+
waitFor(ctx, func() bool {
198+
out, err := catalogCurlJq(ctx, catalogName,
199+
fmt.Sprintf(`objects | select(.schema == "olm.channel") | select(.package == "%s") | .name`, packageName))
200+
if err != nil {
201+
logger.Info("Catalog query failed, retrying", "catalog", catalogName, "package", packageName, "error", err, "stderr", stderrOutput(err))
202+
return false
203+
}
204+
if strings.TrimSpace(out) == "" {
205+
logger.Info("Package has no channels, retrying", "catalog", catalogName, "package", packageName)
206+
return false
207+
}
208+
return true
209+
})
145210
return nil
146211
}
147212

148213
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-
}
214+
waitFor(ctx, func() bool {
215+
out, err := catalogCurlJq(ctx, catalogName,
216+
fmt.Sprintf(`objects | select(.schema == "olm.bundle") | select(.package == "%s") | .name`, packageName))
217+
if err != nil {
218+
logger.Info("Catalog query failed, retrying", "catalog", catalogName, "package", packageName, "error", err, "stderr", stderrOutput(err))
219+
return false
220+
}
221+
if strings.TrimSpace(out) == "" {
222+
logger.Info("Package has no bundles, retrying", "catalog", catalogName, "package", packageName)
223+
return false
224+
}
225+
return true
226+
})
157227
return nil
158228
}
159229

0 commit comments

Comments
 (0)