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 +}