Skip to content

Commit eeea7e4

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 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 7ec6a0e commit eeea7e4

1 file changed

Lines changed: 115 additions & 42 deletions

File tree

‎test/e2e/steps/demo_steps.go‎

Lines changed: 115 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,147 @@ 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
}
117171
script := fmt.Sprintf(
118-
`curl -s -k https://%s/catalogs/%s/api/v1/all | jq -s '%s'`,
172+
`set -o pipefail; curl -sS -k --compressed --fail https://%s/catalogs/%s/api/v1/all | jq '%s'`,
119173
addr, catalogName, jqFilter,
120174
)
121-
return bash(ctx, script)
175+
out, err := bash(ctx, script)
176+
if err != nil {
177+
resetCatalogPortForward(ctx)
178+
}
179+
return out, err
122180
}
123181

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

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

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

0 commit comments

Comments
 (0)