From c69a286f428aa0d03caaffb399254a415b8103b8 Mon Sep 17 00:00:00 2001 From: Shirly Radco Date: Mon, 20 Jul 2026 17:15:36 +0300 Subject: [PATCH] test: add RBAC e2e tests for create and delete endpoints Add e2e tests verifying that the management API enforces Kubernetes RBAC for create and delete alert rule operations. Three user profiles are tested: unprivileged (expects 403), namespace-scoped (succeeds in own namespace, denied elsewhere), and cluster-admin (succeeds everywhere). Also fixes a critical bug in newUserScopedClientsets: when the base rest.Config uses client certificates (common in CI kubeconfigs), CopyConfig preserved them. Since Kubernetes authenticates via client certs when both certs and bearer token are present, user RBAC was bypassed entirely. Use AnonymousClientConfig to strip all auth so the API server authenticates exclusively via the user's bearer token. Signed-off-by: Shirly Radco Co-authored-by: AI Assistant --- internal/managementrouter/router.go | 38 +++-- internal/managementrouter/router_test.go | 83 ++++++++++ pkg/k8s/user_scoped_client.go | 15 +- pkg/k8s/user_scoped_client_test.go | 64 ++++++++ test/e2e/create_alert_rule_test.go | 116 ++++++++++++++ test/e2e/delete_alert_rule_test.go | 196 +++++++++++++++++++++++ test/e2e/framework/framework.go | 154 ++++++++++++++++++ 7 files changed, 649 insertions(+), 17 deletions(-) create mode 100644 internal/managementrouter/router_test.go create mode 100644 pkg/k8s/user_scoped_client_test.go diff --git a/internal/managementrouter/router.go b/internal/managementrouter/router.go index 34c468ba6..12bbf4934 100644 --- a/internal/managementrouter/router.go +++ b/internal/managementrouter/router.go @@ -11,6 +11,7 @@ import ( "github.com/gorilla/mux" "github.com/sirupsen/logrus" + apierrors "k8s.io/apimachinery/pkg/api/errors" "github.com/openshift/monitoring-plugin/pkg/k8s" "github.com/openshift/monitoring-plugin/pkg/management" @@ -62,6 +63,7 @@ func authMiddleware(next http.Handler) http.Handler { }) } +// writeError sends a JSON {"error": message} response with the given status code. func writeError(w http.ResponseWriter, statusCode int, message string) { w.Header().Set("Content-Type", "application/json") w.WriteHeader(statusCode) @@ -75,28 +77,38 @@ func writeError(w http.ResponseWriter, statusCode int, message string) { } } +// handleError maps err to an HTTP status via parseError and writes the response. func handleError(w http.ResponseWriter, err error) { status, message := parseError(err) writeError(w, status, message) } +// parseError inspects err and returns a (statusCode, userMessage) pair. +// Kubernetes auth errors are checked first to prevent information leakage; +// domain errors are then mapped to 4xx codes. func parseError(err error) (int, string) { - var nf *management.NotFoundError - if errors.As(err, &nf) { + var ( + nf *management.NotFoundError + ve *management.ValidationError + na *management.NotAllowedError + ce *management.ConflictError + ) + + switch { + case apierrors.IsUnauthorized(err): + return http.StatusUnauthorized, "authentication failed" + case apierrors.IsForbidden(err): + return http.StatusForbidden, "insufficient permissions" + case errors.As(err, &nf): return http.StatusNotFound, err.Error() - } - var ve *management.ValidationError - if errors.As(err, &ve) { + case errors.As(err, &ve): return http.StatusBadRequest, err.Error() - } - var na *management.NotAllowedError - if errors.As(err, &na) { + case errors.As(err, &na): return http.StatusMethodNotAllowed, err.Error() - } - var ce *management.ConflictError - if errors.As(err, &ce) { + case errors.As(err, &ce): return http.StatusConflict, err.Error() + default: + log.WithError(err).Error("unexpected management API error") + return http.StatusInternalServerError, "An unexpected error occurred" } - log.WithError(err).Error("unexpected management API error") - return http.StatusInternalServerError, "An unexpected error occurred" } diff --git a/internal/managementrouter/router_test.go b/internal/managementrouter/router_test.go new file mode 100644 index 000000000..787d3c7e5 --- /dev/null +++ b/internal/managementrouter/router_test.go @@ -0,0 +1,83 @@ +package managementrouter + +import ( + "fmt" + "net/http" + "testing" + + apierrors "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/runtime/schema" + + "github.com/openshift/monitoring-plugin/pkg/management" +) + +func TestParseError(t *testing.T) { + tests := []struct { + name string + err error + expectedStatus int + expectedMsg string + }{ + { + name: "NotFoundError", + err: &management.NotFoundError{Resource: "AlertRule", Id: "abc"}, + expectedStatus: http.StatusNotFound, + }, + { + name: "ValidationError", + err: &management.ValidationError{Message: "bad input"}, + expectedStatus: http.StatusBadRequest, + }, + { + name: "NotAllowedError", + err: &management.NotAllowedError{Message: "not allowed"}, + expectedStatus: http.StatusMethodNotAllowed, + }, + { + name: "ConflictError", + err: &management.ConflictError{Message: "conflict"}, + expectedStatus: http.StatusConflict, + }, + { + name: "Kubernetes Forbidden", + err: apierrors.NewForbidden(schema.GroupResource{ + Group: "monitoring.coreos.com", Resource: "prometheusrules", + }, "test-pr", fmt.Errorf("access denied")), + expectedStatus: http.StatusForbidden, + expectedMsg: "insufficient permissions", + }, + { + name: "Kubernetes Forbidden wrapped", + err: fmt.Errorf("failed to get PrometheusRule: %w", + apierrors.NewForbidden(schema.GroupResource{ + Group: "monitoring.coreos.com", Resource: "prometheusrules", + }, "test-pr", fmt.Errorf("access denied"))), + expectedStatus: http.StatusForbidden, + expectedMsg: "insufficient permissions", + }, + { + name: "Kubernetes Unauthorized", + err: apierrors.NewUnauthorized("token expired"), + expectedStatus: http.StatusUnauthorized, + expectedMsg: "authentication failed", + }, + { + name: "unknown error", + err: fmt.Errorf("something unexpected"), + expectedStatus: http.StatusInternalServerError, + expectedMsg: "An unexpected error occurred", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + status, msg := parseError(tt.err) + if status != tt.expectedStatus { + t.Errorf("expected status %d, got %d", tt.expectedStatus, status) + } + if tt.expectedMsg != "" && msg != tt.expectedMsg { + t.Errorf("expected message %q, got %q", tt.expectedMsg, msg) + } + }) + } +} diff --git a/pkg/k8s/user_scoped_client.go b/pkg/k8s/user_scoped_client.go index 448ba4d8e..d180945d5 100644 --- a/pkg/k8s/user_scoped_client.go +++ b/pkg/k8s/user_scoped_client.go @@ -13,14 +13,21 @@ type userScopedClientsets struct { osmV1 *osmv1client.Clientset } +// buildUserScopedConfig creates a rest.Config that authenticates exclusively +// with the given bearer token. It uses AnonymousClientConfig to strip all +// existing auth (certs, basic auth, auth/exec providers, impersonation) while +// preserving the server connection settings (host, TLS CA, proxy). +func buildUserScopedConfig(baseConfig *rest.Config, userToken string) *rest.Config { + cfg := rest.AnonymousClientConfig(baseConfig) + cfg.BearerToken = userToken + return cfg +} + // newUserScopedClientsets creates clientsets that carry the supplied bearer // token so that Kubernetes RBAC is enforced for the requesting user on all // mutating API calls. func newUserScopedClientsets(baseConfig *rest.Config, userToken string) (*userScopedClientsets, error) { - cfg := rest.CopyConfig(baseConfig) - // Override any SA token loaded from the file system with the user's token. - cfg.BearerToken = userToken - cfg.BearerTokenFile = "" + cfg := buildUserScopedConfig(baseConfig, userToken) monClient, err := monitoringv1client.NewForConfig(cfg) if err != nil { diff --git a/pkg/k8s/user_scoped_client_test.go b/pkg/k8s/user_scoped_client_test.go new file mode 100644 index 000000000..4de96ccd8 --- /dev/null +++ b/pkg/k8s/user_scoped_client_test.go @@ -0,0 +1,64 @@ +package k8s + +import ( + "testing" + + "k8s.io/client-go/rest" +) + +func TestBuildUserScopedConfig(t *testing.T) { + base := &rest.Config{ + Host: "https://api.example.com:6443", + BearerToken: "sa-token", + BearerTokenFile: "/var/run/secrets/kubernetes.io/serviceaccount/token", + TLSClientConfig: rest.TLSClientConfig{ + Insecure: true, + CertData: []byte("admin-cert"), + KeyData: []byte("admin-key"), + CertFile: "/path/to/cert", + KeyFile: "/path/to/key", + }, + } + + cfg := buildUserScopedConfig(base, "user-token") + + // Derived config uses the user token exclusively. + if cfg.BearerToken != "user-token" { + t.Errorf("derived BearerToken = %q, want %q", cfg.BearerToken, "user-token") + } + if cfg.BearerTokenFile != "" { + t.Errorf("derived BearerTokenFile = %q, want empty", cfg.BearerTokenFile) + } + if cfg.CertData != nil { + t.Error("derived CertData should be nil") + } + if cfg.KeyData != nil { + t.Error("derived KeyData should be nil") + } + if cfg.CertFile != "" { + t.Errorf("derived CertFile = %q, want empty", cfg.CertFile) + } + if cfg.KeyFile != "" { + t.Errorf("derived KeyFile = %q, want empty", cfg.KeyFile) + } + if !cfg.Insecure { + t.Error("derived Insecure should be preserved as true") + } + if cfg.Host != base.Host { + t.Errorf("derived Host = %q, want %q", cfg.Host, base.Host) + } + + // Base config must not be mutated. + if base.CertData == nil { + t.Error("base CertData was mutated") + } + if base.KeyData == nil { + t.Error("base KeyData was mutated") + } + if base.BearerToken != "sa-token" { + t.Errorf("base BearerToken = %q, want %q", base.BearerToken, "sa-token") + } + if base.BearerTokenFile != "/var/run/secrets/kubernetes.io/serviceaccount/token" { + t.Errorf("base BearerTokenFile = %q, was mutated", base.BearerTokenFile) + } +} diff --git a/test/e2e/create_alert_rule_test.go b/test/e2e/create_alert_rule_test.go index ee312e9de..58ce542dd 100644 --- a/test/e2e/create_alert_rule_test.go +++ b/test/e2e/create_alert_rule_test.go @@ -3,9 +3,14 @@ package e2e import ( + "bytes" "context" + "encoding/json" "errors" "fmt" + "io" + "net/http" + "net/url" "testing" "time" @@ -87,3 +92,114 @@ func TestCreateUserDefinedAlertRule(t *testing.T) { }) require.NoError(t, err) } + +// TestRBAC_CreateAlertRule verifies that the create endpoint enforces Kubernetes +// RBAC across three user profiles: anonymous (403), namespace-scoped (201 in +// own namespace, 403 elsewhere), and cluster-admin (201 everywhere). +func TestRBAC_CreateAlertRule(t *testing.T) { + f, err := framework.New() + if err != nil { + t.Fatalf("Failed to create framework: %v", err) + } + + ctx := context.Background() + + nsY, cleanupY, err := f.CreateUserNamespace(ctx, "test-rbac-create-y") + if err != nil { + t.Fatalf("Failed to create namespace Y: %v", err) + } + defer func() { _ = cleanupY() }() + + nsZ, cleanupZ, err := f.CreateUserNamespace(ctx, "test-rbac-create-z") + if err != nil { + t.Fatalf("Failed to create namespace Z: %v", err) + } + defer func() { _ = cleanupZ() }() + + anonymousUser, err := f.CreateAnonymousUser(ctx, "e2e-rbac-user-a", "default") + if err != nil { + t.Fatalf("Failed to create anonymous user: %v", err) + } + defer func() { _ = anonymousUser.Cleanup() }() + + userScopedToNamespaceY, err := f.CreateScopedUser(ctx, "e2e-rbac-user-b", nsY, + "monitoring.coreos.com", []string{"prometheusrules"}, []string{"get", "create", "update", "patch"}) + if err != nil { + t.Fatalf("Failed to create scoped user for namespace Y: %v", err) + } + defer func() { _ = userScopedToNamespaceY.Cleanup() }() + + cases := []struct { + name string + token string + namespace string + alertName string + wantStatus int + }{ + {"AnonymousUser_FailsNamespaceY", anonymousUser.Token, nsY, "RBACAlertA", http.StatusForbidden}, + {"ScopedUser_SucceedsNamespaceY", userScopedToNamespaceY.Token, nsY, "RBACAlertBY", http.StatusCreated}, + {"ScopedUser_FailsNamespaceZ", userScopedToNamespaceY.Token, nsZ, "RBACAlertBZ", http.StatusForbidden}, + {"ClusterAdmin_SucceedsNamespaceY", f.BearerToken, nsY, "RBACAlertCY", http.StatusCreated}, + {"ClusterAdmin_SucceedsNamespaceZ", f.BearerToken, nsZ, "RBACAlertCZ", http.StatusCreated}, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + status := createAlertRuleWithToken(t, f, ctx, tc.token, tc.namespace, tc.alertName) + if status != tc.wantStatus { + t.Fatalf("Expected status %d, got %d", tc.wantStatus, status) + } + }) + } +} + +// createAlertRuleWithToken sends a create alert rule request using the given +// bearer token and returns the HTTP status code. +func createAlertRuleWithToken(t *testing.T, f *framework.Framework, ctx context.Context, token, namespace, alertName string) int { + t.Helper() + + expr := fmt.Sprintf("absent(nonexistent{e2e_rbac_create=%q})", alertName) + payload := managementrouter.CreateAlertRuleRequest{ + AlertingRule: &managementrouter.AlertRuleSpec{ + Alert: &alertName, + Expr: &expr, + Labels: &map[string]string{ + "severity": "info", + }, + }, + PrometheusRule: &managementrouter.PrometheusRuleTarget{ + PrometheusRuleName: "e2e-rbac-pr", + PrometheusRuleNamespace: namespace, + }, + } + + reqBody, err := json.Marshal(payload) + if err != nil { + t.Fatalf("Failed to marshal create request: %v", err) + } + + createURL, err := url.JoinPath(f.PluginURL, "api/v1/alerting/rules") + if err != nil { + t.Fatalf("Failed to build URL: %v", err) + } + + req, err := http.NewRequestWithContext(ctx, http.MethodPost, createURL, bytes.NewBuffer(reqBody)) + if err != nil { + t.Fatalf("Failed to create HTTP request: %v", err) + } + req.Header.Set("Content-Type", "application/json") + req.Header.Set("Authorization", "Bearer "+token) + + resp, err := f.HTTPClient().Do(req) + if err != nil { + t.Fatalf("Failed to make create request: %v", err) + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusCreated { + body, _ := io.ReadAll(resp.Body) + t.Logf("Create %s in %s: status %d, body: %s", alertName, namespace, resp.StatusCode, string(body)) + } + + return resp.StatusCode +} diff --git a/test/e2e/delete_alert_rule_test.go b/test/e2e/delete_alert_rule_test.go index 672e27753..e7c897e91 100644 --- a/test/e2e/delete_alert_rule_test.go +++ b/test/e2e/delete_alert_rule_test.go @@ -9,6 +9,7 @@ import ( "fmt" "io" "net/http" + "net/url" "testing" "time" @@ -142,3 +143,198 @@ func TestDeleteAlertRule(t *testing.T) { }) require.NoError(t, err) } + +// TestRBAC_DeleteAlertRule verifies that the bulk-delete endpoint enforces +// Kubernetes RBAC across three user profiles: anonymous (403), +// namespace-scoped (204 in own namespace, 403 elsewhere), and cluster-admin +// (204 everywhere). +func TestRBAC_DeleteAlertRule(t *testing.T) { + f, err := framework.New() + if err != nil { + t.Fatalf("Failed to create framework: %v", err) + } + + ctx := context.Background() + + nsY, cleanupY, err := f.CreateUserNamespace(ctx, "test-rbac-del-y") + if err != nil { + t.Fatalf("Failed to create namespace Y: %v", err) + } + defer func() { _ = cleanupY() }() + + nsZ, cleanupZ, err := f.CreateUserNamespace(ctx, "test-rbac-del-z") + if err != nil { + t.Fatalf("Failed to create namespace Z: %v", err) + } + defer func() { _ = cleanupZ() }() + + anonymousUser, err := f.CreateAnonymousUser(ctx, "e2e-rbac-del-a", "default") + if err != nil { + t.Fatalf("Failed to create anonymous user: %v", err) + } + defer func() { _ = anonymousUser.Cleanup() }() + + userScopedToNamespaceY, err := f.CreateScopedUser(ctx, "e2e-rbac-del-b", nsY, + "monitoring.coreos.com", []string{"prometheusrules"}, []string{"get", "create", "update", "patch", "delete"}) + if err != nil { + t.Fatalf("Failed to create scoped user for namespace Y: %v", err) + } + defer func() { _ = userScopedToNamespaceY.Cleanup() }() + + ruleInY, err := createRuleViaAPI(ctx, f, managementrouter.CreateAlertRuleRequest{ + AlertingRule: &managementrouter.AlertRuleSpec{ + Alert: new("RBACDelAlertY"), + Expr: new(fmt.Sprintf("absent(nonexistent{e2e_rbac_del=%q})", "y")), + Labels: &map[string]string{ + "severity": "info", + }, + }, + PrometheusRule: &managementrouter.PrometheusRuleTarget{ + PrometheusRuleName: "e2e-rbac-del-pr", + PrometheusRuleNamespace: nsY, + }, + }) + if err != nil { + t.Fatalf("Failed to create rule in namespace Y: %v", err) + } + t.Logf("Created rule in namespace Y: %s", ruleInY) + + ruleInZ, err := createRuleViaAPI(ctx, f, managementrouter.CreateAlertRuleRequest{ + AlertingRule: &managementrouter.AlertRuleSpec{ + Alert: new("RBACDelAlertZ"), + Expr: new(fmt.Sprintf("absent(nonexistent{e2e_rbac_del=%q})", "z")), + Labels: &map[string]string{ + "severity": "info", + }, + }, + PrometheusRule: &managementrouter.PrometheusRuleTarget{ + PrometheusRuleName: "e2e-rbac-del-pr", + PrometheusRuleNamespace: nsZ, + }, + }) + if err != nil { + t.Fatalf("Failed to create rule in namespace Z: %v", err) + } + t.Logf("Created rule in namespace Z: %s", ruleInZ) + + ruleInY2, err := createRuleViaAPI(ctx, f, managementrouter.CreateAlertRuleRequest{ + AlertingRule: &managementrouter.AlertRuleSpec{ + Alert: new("RBACDelAlertY2"), + Expr: new(fmt.Sprintf("absent(nonexistent{e2e_rbac_del=%q})", "y2")), + Labels: &map[string]string{ + "severity": "info", + }, + }, + PrometheusRule: &managementrouter.PrometheusRuleTarget{ + PrometheusRuleName: "e2e-rbac-del-pr", + PrometheusRuleNamespace: nsY, + }, + }) + if err != nil { + t.Fatalf("Failed to create second rule in namespace Y: %v", err) + } + t.Logf("Created second rule in namespace Y: %s", ruleInY2) + + waitForCacheSync(t, f, ctx, anonymousUser.Token, ruleInY) + + cases := []struct { + name string + token string + ruleID string + wantStatus int + }{ + {"AnonymousUser_DeniedNamespaceY", anonymousUser.Token, ruleInY, http.StatusForbidden}, + {"ScopedUser_SucceedsNamespaceY", userScopedToNamespaceY.Token, ruleInY, http.StatusNoContent}, + {"ScopedUser_DeniedNamespaceZ", userScopedToNamespaceY.Token, ruleInZ, http.StatusForbidden}, + {"ClusterAdmin_SucceedsNamespaceZ", f.BearerToken, ruleInZ, http.StatusNoContent}, + {"ClusterAdmin_SucceedsNamespaceY", f.BearerToken, ruleInY2, http.StatusNoContent}, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + status := deleteAlertRuleWithToken(t, f, ctx, tc.token, tc.ruleID) + if status != tc.wantStatus { + t.Fatalf("Expected per-rule status %d, got %d", tc.wantStatus, status) + } + }) + } +} + +// waitForCacheSync polls until the relabeled-rules cache has synced by +// attempting a bulk-delete probe. A 403 (Forbidden) or 204 (NoContent) +// per-rule status indicates the rule was found in cache and RBAC was evaluated. +func waitForCacheSync(t *testing.T, f *framework.Framework, ctx context.Context, token, ruleID string) { + t.Helper() + const timeout = 30 * time.Second + const interval = time.Second + deadline := time.Now().Add(timeout) + for { + status, err := tryDeleteAlertRule(f, ctx, token, ruleID) + if err == nil && (status == http.StatusForbidden || status == http.StatusNoContent) { + return + } + if time.Now().After(deadline) { + t.Fatalf("Cache sync timed out after %v (last status=%d, err=%v)", timeout, status, err) + } + if err != nil { + t.Logf("Cache sync: %v, retrying...", err) + } else { + t.Logf("Cache sync: per-rule status %d, retrying...", status) + } + time.Sleep(interval) + } +} + +// tryDeleteAlertRule attempts a single-rule bulk-delete and returns the per-rule +// status code without calling t.Fatal, making it suitable for polling loops. +func tryDeleteAlertRule(f *framework.Framework, ctx context.Context, token, ruleID string) (int, error) { + payload := managementrouter.BulkDeleteAlertRulesRequest{ + RuleIds: []string{ruleID}, + } + reqBody, err := json.Marshal(payload) + if err != nil { + return 0, fmt.Errorf("marshal delete request: %w", err) + } + deleteURL, err := url.JoinPath(f.PluginURL, "api/v1/alerting/rules") + if err != nil { + return 0, fmt.Errorf("build URL: %w", err) + } + req, err := http.NewRequestWithContext(ctx, http.MethodDelete, deleteURL, bytes.NewBuffer(reqBody)) + if err != nil { + return 0, fmt.Errorf("create HTTP request: %w", err) + } + req.Header.Set("Content-Type", "application/json") + req.Header.Set("Authorization", "Bearer "+token) + + resp, err := f.HTTPClient().Do(req) + if err != nil { + return 0, fmt.Errorf("make delete request: %w", err) + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + body, _ := io.ReadAll(resp.Body) + return resp.StatusCode, fmt.Errorf("expected bulk response 200, got %d: %s", resp.StatusCode, string(body)) + } + + var deleteResp managementrouter.BulkDeleteAlertRulesResponse + if err := json.NewDecoder(resp.Body).Decode(&deleteResp); err != nil { + return 0, fmt.Errorf("decode delete response: %w", err) + } + if len(deleteResp.Rules) != 1 { + return 0, fmt.Errorf("expected 1 per-rule result, got %d", len(deleteResp.Rules)) + } + return deleteResp.Rules[0].StatusCode, nil +} + +// deleteAlertRuleWithToken sends a bulk-delete request for a single rule ID +// using the given bearer token and returns the per-rule HTTP status code. +func deleteAlertRuleWithToken(t *testing.T, f *framework.Framework, ctx context.Context, token, ruleID string) int { + t.Helper() + + status, err := tryDeleteAlertRule(f, ctx, token, ruleID) + if err != nil { + t.Fatalf("Delete request for rule %s failed: %v", ruleID, err) + } + return status +} diff --git a/test/e2e/framework/framework.go b/test/e2e/framework/framework.go index 3e007152c..1c870abae 100644 --- a/test/e2e/framework/framework.go +++ b/test/e2e/framework/framework.go @@ -223,3 +223,157 @@ func createServiceAccountToken(clientset *kubernetes.Clientset) (string, error) } return resp.Status.Token, nil } + +// ScopedUser represents a ServiceAccount with specific RBAC permissions for testing. +type ScopedUser struct { + Token string + Cleanup CleanupFunc +} + +// requestServiceAccountToken creates a short-lived (1 hour) bearer token for +// the named ServiceAccount via the TokenRequest API. The call is retried to +// tolerate transient API failures. +func (f *Framework) requestServiceAccountToken(ctx context.Context, namespace, name string) (string, error) { + expSeconds := int64(3600) + treq := &authv1.TokenRequest{ + Spec: authv1.TokenRequestSpec{ExpirationSeconds: &expSeconds}, + } + var token string + err := retry(3, func() error { + tokenResp, err := f.Clientset.CoreV1().ServiceAccounts(namespace).CreateToken(ctx, name, treq, metav1.CreateOptions{}) + if err != nil { + return err + } + token = tokenResp.Status.Token + return nil + }) + if err != nil { + return "", fmt.Errorf("requesting token for %s/%s: %w", namespace, name, err) + } + return token, nil +} + +// retry calls fn up to maxAttempts times with a 1-second pause between attempts. +// It returns nil on the first successful call or the last error after exhaustion. +func retry(maxAttempts int, fn func() error) error { + var err error + for i := range maxAttempts { + if err = fn(); err == nil { + return nil + } + if i < maxAttempts-1 { + time.Sleep(time.Second) + } + } + return err +} + +// CreateScopedUser creates a ServiceAccount in the given namespace with a Role +// granting the specified verbs on the specified resources. Returns a bearer token +// and a cleanup function. The apiGroup should be e.g. "monitoring.coreos.com". +// API calls are retried to tolerate transient failures. +func (f *Framework) CreateScopedUser(ctx context.Context, name, namespace, apiGroup string, resources, verbs []string) (*ScopedUser, error) { + rollback := func() { + _ = f.Clientset.RbacV1().RoleBindings(namespace).Delete(ctx, name, metav1.DeleteOptions{}) + _ = f.Clientset.RbacV1().Roles(namespace).Delete(ctx, name, metav1.DeleteOptions{}) + _ = f.Clientset.CoreV1().ServiceAccounts(namespace).Delete(ctx, name, metav1.DeleteOptions{}) + } + + sa := &corev1.ServiceAccount{ + ObjectMeta: metav1.ObjectMeta{Name: name}, + } + err := retry(3, func() error { + _, err := f.Clientset.CoreV1().ServiceAccounts(namespace).Create(ctx, sa, metav1.CreateOptions{}) + if apierrors.IsAlreadyExists(err) { + return nil + } + return err + }) + if err != nil { + return nil, fmt.Errorf("creating service account %s/%s: %w", namespace, name, err) + } + + role := &rbacv1.Role{ + ObjectMeta: metav1.ObjectMeta{Name: name}, + Rules: []rbacv1.PolicyRule{{ + APIGroups: []string{apiGroup}, + Resources: resources, + Verbs: verbs, + }}, + } + err = retry(3, func() error { + _, err := f.Clientset.RbacV1().Roles(namespace).Create(ctx, role, metav1.CreateOptions{}) + if apierrors.IsAlreadyExists(err) { + return nil + } + return err + }) + if err != nil { + rollback() + return nil, fmt.Errorf("creating role %s/%s: %w", namespace, name, err) + } + + rb := &rbacv1.RoleBinding{ + ObjectMeta: metav1.ObjectMeta{Name: name}, + Subjects: []rbacv1.Subject{{ + Kind: rbacv1.ServiceAccountKind, + Name: name, + Namespace: namespace, + }}, + RoleRef: rbacv1.RoleRef{ + APIGroup: rbacv1.GroupName, + Kind: "Role", + Name: name, + }, + } + err = retry(3, func() error { + _, err := f.Clientset.RbacV1().RoleBindings(namespace).Create(ctx, rb, metav1.CreateOptions{}) + if apierrors.IsAlreadyExists(err) { + return nil + } + return err + }) + if err != nil { + rollback() + return nil, fmt.Errorf("creating role binding %s/%s: %w", namespace, name, err) + } + + token, err := f.requestServiceAccountToken(ctx, namespace, name) + if err != nil { + rollback() + return nil, err + } + + return &ScopedUser{Token: token, Cleanup: func() error { rollback(); return nil }}, nil +} + +// CreateAnonymousUser creates a ServiceAccount with no RBAC permissions. +// API calls are retried to tolerate transient failures. +func (f *Framework) CreateAnonymousUser(ctx context.Context, name, namespace string) (*ScopedUser, error) { + sa := &corev1.ServiceAccount{ + ObjectMeta: metav1.ObjectMeta{Name: name}, + } + err := retry(3, func() error { + _, err := f.Clientset.CoreV1().ServiceAccounts(namespace).Create(ctx, sa, metav1.CreateOptions{}) + if apierrors.IsAlreadyExists(err) { + return nil + } + return err + }) + if err != nil { + return nil, fmt.Errorf("creating service account %s/%s: %w", namespace, name, err) + } + + token, err := f.requestServiceAccountToken(ctx, namespace, name) + if err != nil { + _ = f.Clientset.CoreV1().ServiceAccounts(namespace).Delete(ctx, name, metav1.DeleteOptions{}) + return nil, err + } + + cleanup := func() error { + _ = f.Clientset.CoreV1().ServiceAccounts(namespace).Delete(ctx, name, metav1.DeleteOptions{}) + return nil + } + + return &ScopedUser{Token: token, Cleanup: cleanup}, nil +}