Skip to content

Commit 4988d38

Browse files
pedjakclaude
andauthored
🌱 Deduplicate metrics service lookup and port-forward in e2e steps (#2710)
SendMetricsRequest and withMetricsPortForward duplicated service endpoint discovery and port-forward lifecycle code. Extract shared helpers (getMetricsService, metricsPort, portForward) and rewrite both callers to use them. Port-forward cleanup is deferred to avoid process leaks if waitFor panics. Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
1 parent ef1a6b1 commit 4988d38

2 files changed

Lines changed: 80 additions & 75 deletions

File tree

‎test/e2e/steps/steps.go‎

Lines changed: 29 additions & 54 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,6 @@ import (
88
"encoding/json"
99
"fmt"
1010
"io"
11-
"net"
1211
"net/http"
1312
"os"
1413
"os/exec"
@@ -1407,36 +1406,24 @@ func httpGet(url string, token string) (*http.Response, error) {
14071406
return resp, nil
14081407
}
14091408

1410-
func randomAvailablePort() (int, error) {
1411-
l, err := net.Listen("tcp", "127.0.0.1:0")
1412-
if err != nil {
1413-
return 0, err
1414-
}
1415-
defer l.Close()
1416-
return l.Addr().(*net.TCPAddr).Port, nil
1417-
}
1418-
14191409
// SendMetricsRequest sets up port-forwarding to the controller's service pods and waits for the metrics endpoint
14201410
// to return a successful response. Stores the response body per pod in the scenario context. Polls with timeout.
14211411
func SendMetricsRequest(ctx context.Context, serviceAccount string, endpoint string, controllerName string) error {
14221412
sc := scenarioCtx(ctx)
1423-
serviceNs, err := k8sClient("get", "service", "-A", "-o", fmt.Sprintf(`jsonpath={.items[?(@.metadata.name=="%s-service")].metadata.namespace}`, controllerName))
1413+
svc, err := getMetricsService(controllerName)
14241414
if err != nil {
14251415
return err
14261416
}
1427-
v, err := k8sClient("get", "service", "-n", serviceNs, fmt.Sprintf("%s-service", controllerName), "-o", "json")
1417+
mPort, err := metricsPort(svc)
14281418
if err != nil {
14291419
return err
14301420
}
1431-
var service corev1.Service
1432-
if err := json.Unmarshal([]byte(v), &service); err != nil {
1433-
return err
1434-
}
1435-
podNameCmd := []string{"get", "pod", "-n", olmNamespace, "-o", "jsonpath={.items}"}
1436-
for k, v := range service.Spec.Selector {
1421+
1422+
podNameCmd := []string{"get", "pod", "-n", svc.Namespace, "-o", "jsonpath={.items}"}
1423+
for k, v := range svc.Spec.Selector {
14371424
podNameCmd = append(podNameCmd, fmt.Sprintf("--selector=%s=%s", k, v))
14381425
}
1439-
v, err = k8sClient(podNameCmd...)
1426+
v, err := k8sClient(podNameCmd...)
14401427
if err != nil {
14411428
return err
14421429
}
@@ -1449,51 +1436,39 @@ func SendMetricsRequest(ctx context.Context, serviceAccount string, endpoint str
14491436
if err != nil {
14501437
return err
14511438
}
1452-
var metricsPort int32
1453-
for _, p := range service.Spec.Ports {
1454-
if p.Name == "metrics" {
1455-
metricsPort = p.Port
1456-
break
1457-
}
1458-
}
1439+
14591440
sc.metricsResponse = make(map[string]string)
14601441
for _, p := range pods {
1461-
port, err := randomAvailablePort()
1462-
if err != nil {
1463-
return err
1464-
}
1465-
portForwardCmd := exec.Command(k8sCli, "port-forward", "-n", p.Namespace, fmt.Sprintf("pod/%s", p.Name), fmt.Sprintf("%d:%d", port, metricsPort)) //nolint:gosec // perfectly safe to start port-forwarder for provided controller name
1466-
logger.V(1).Info("starting port-forward", "command", strings.Join(portForwardCmd.Args, " "))
1467-
if err := portForwardCmd.Start(); err != nil {
1468-
logger.Error(err, fmt.Sprintf("failed to start port-forward for pod %s", p.Name))
1469-
return err
1470-
}
1471-
waitFor(ctx, func() bool {
1472-
resp, err := httpGet(fmt.Sprintf("https://localhost:%d%s", port, endpoint), token)
1442+
if err := func() error {
1443+
addr, cleanup, err := portForward(p.Namespace, fmt.Sprintf("pod/%s", p.Name), mPort)
14731444
if err != nil {
1474-
return false
1445+
return err
14751446
}
1476-
defer resp.Body.Close()
1447+
defer cleanup()
1448+
waitFor(ctx, func() bool {
1449+
resp, err := httpGet(fmt.Sprintf("https://%s%s", addr, endpoint), token)
1450+
if err != nil {
1451+
return false
1452+
}
1453+
defer resp.Body.Close()
14771454

1478-
if resp.StatusCode == http.StatusOK {
1455+
if resp.StatusCode == http.StatusOK {
1456+
b, err := io.ReadAll(resp.Body)
1457+
if err != nil {
1458+
return false
1459+
}
1460+
sc.metricsResponse[p.Name] = string(b)
1461+
return true
1462+
}
14791463
b, err := io.ReadAll(resp.Body)
14801464
if err != nil {
14811465
return false
14821466
}
1483-
sc.metricsResponse[p.Name] = string(b)
1484-
return true
1485-
}
1486-
b, err := io.ReadAll(resp.Body)
1487-
if err != nil {
1467+
logger.V(1).Info("failed to get metrics", "pod", p.Name, "response", string(b))
14881468
return false
1489-
}
1490-
logger.V(1).Info("failed to get metrics", "pod", p.Name, "response", string(b))
1491-
return false
1492-
})
1493-
if err := portForwardCmd.Process.Kill(); err != nil {
1494-
return err
1495-
}
1496-
if _, err := portForwardCmd.Process.Wait(); err != nil {
1469+
})
1470+
return nil
1471+
}(); err != nil {
14971472
return err
14981473
}
14991474
}

‎test/e2e/steps/tls_steps.go‎

Lines changed: 51 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -53,65 +53,95 @@ var curveIDByName = map[string]tls.CurveID{
5353
"secp521r1": tls.CurveP521,
5454
}
5555

56-
// getMetricsServiceEndpoint returns the namespace and metrics port for the named component service.
57-
func getMetricsServiceEndpoint(component string) (string, int32, error) {
56+
// getMetricsService returns the full Service object for the named component's metrics service.
57+
// The namespace is available via svc.Namespace.
58+
func getMetricsService(component string) (*corev1.Service, error) {
5859
serviceName := fmt.Sprintf("%s-service", component)
5960
serviceNs, err := k8sClient("get", "service", "-A", "-o",
6061
fmt.Sprintf(`jsonpath={.items[?(@.metadata.name=="%s")].metadata.namespace}`, serviceName))
6162
if err != nil {
62-
return "", 0, fmt.Errorf("failed to find namespace for service %s: %w", serviceName, err)
63+
return nil, fmt.Errorf("failed to find namespace for service %s: %w", serviceName, err)
6364
}
6465
serviceNs = strings.TrimSpace(serviceNs)
6566
if serviceNs == "" {
66-
return "", 0, fmt.Errorf("service %s not found in any namespace", serviceName)
67+
return nil, fmt.Errorf("service %s not found in any namespace", serviceName)
6768
}
6869

6970
raw, err := k8sClient("get", "service", "-n", serviceNs, serviceName, "-o", "json")
7071
if err != nil {
71-
return "", 0, fmt.Errorf("failed to get service %s: %w", serviceName, err)
72+
return nil, fmt.Errorf("failed to get service %s: %w", serviceName, err)
7273
}
7374
var svc corev1.Service
7475
if err := json.Unmarshal([]byte(raw), &svc); err != nil {
75-
return "", 0, fmt.Errorf("failed to unmarshal service %s: %w", serviceName, err)
76+
return nil, fmt.Errorf("failed to unmarshal service %s: %w", serviceName, err)
7677
}
78+
return &svc, nil
79+
}
80+
81+
// metricsPort returns the port number of the port named "metrics" on the given service.
82+
func metricsPort(svc *corev1.Service) (int32, error) {
7783
for _, p := range svc.Spec.Ports {
7884
if p.Name == "metrics" {
79-
return serviceNs, p.Port, nil
85+
return p.Port, nil
8086
}
8187
}
82-
return "", 0, fmt.Errorf("no port named 'metrics' found on service %s", serviceName)
88+
return 0, fmt.Errorf("no port named 'metrics' found on service %s", svc.Name)
8389
}
8490

85-
// withMetricsPortForward starts a kubectl port-forward to the component's metrics service,
86-
// waits until a basic TLS connection succeeds (confirming the port-forward is ready),
87-
// then calls fn with the local address. The port-forward is torn down when fn returns.
88-
func withMetricsPortForward(ctx context.Context, component string, fn func(addr string) error) error {
89-
ns, metricsPort, err := getMetricsServiceEndpoint(component)
91+
func randomAvailablePort() (int, error) {
92+
l, err := net.Listen("tcp", "127.0.0.1:0")
9093
if err != nil {
91-
return err
94+
return 0, err
9295
}
96+
defer l.Close()
97+
return l.Addr().(*net.TCPAddr).Port, nil
98+
}
9399

100+
// portForward starts a kubectl port-forward to target (e.g. "service/foo" or "pod/bar")
101+
// in the given namespace, mapping a random local port to remotePort. It returns the
102+
// local address and a cleanup function. The caller is responsible for calling cleanup.
103+
func portForward(ns, target string, remotePort int32) (string, func(), error) {
94104
localPort, err := randomAvailablePort()
95105
if err != nil {
96-
return fmt.Errorf("failed to find a free local port: %w", err)
106+
return "", nil, fmt.Errorf("failed to find a free local port: %w", err)
97107
}
98108

99-
serviceName := fmt.Sprintf("%s-service", component)
100109
pfCmd := exec.Command(k8sCli, "port-forward", "-n", ns, //nolint:gosec
101-
fmt.Sprintf("service/%s", serviceName),
102-
fmt.Sprintf("%d:%d", localPort, metricsPort))
110+
target, fmt.Sprintf("%d:%d", localPort, remotePort))
103111
pfCmd.Env = append(os.Environ(), fmt.Sprintf("KUBECONFIG=%s", kubeconfigPath))
104112
if err := pfCmd.Start(); err != nil {
105-
return fmt.Errorf("failed to start port-forward to %s: %w", serviceName, err)
113+
return "", nil, fmt.Errorf("failed to start port-forward to %s: %w", target, err)
106114
}
107-
defer func() {
115+
116+
cleanup := func() {
108117
if p := pfCmd.Process; p != nil {
109118
_ = p.Kill()
110119
_ = pfCmd.Wait()
111120
}
112-
}()
121+
}
113122

114123
addr := fmt.Sprintf("127.0.0.1:%d", localPort)
124+
return addr, cleanup, nil
125+
}
126+
127+
// withMetricsPortForward starts a kubectl port-forward to the component's metrics service,
128+
// waits until a basic TLS connection succeeds (confirming the port-forward is ready),
129+
// then calls fn with the local address. The port-forward is torn down when fn returns.
130+
func withMetricsPortForward(ctx context.Context, component string, fn func(addr string) error) error {
131+
svc, err := getMetricsService(component)
132+
if err != nil {
133+
return err
134+
}
135+
port, err := metricsPort(svc)
136+
if err != nil {
137+
return err
138+
}
139+
140+
addr, cleanup, err := portForward(svc.Namespace, fmt.Sprintf("service/%s", svc.Name), port)
141+
if err != nil {
142+
return err
143+
}
144+
defer cleanup()
115145

116146
// Wait until the port-forward is accepting connections. A plain TLS dial (no version
117147
// restrictions) serves as the readiness probe; any successful TLS handshake confirms

0 commit comments

Comments
 (0)