From 5202dfb44be4705250baeb352800b67d53030fec Mon Sep 17 00:00:00 2001 From: Thomas Kosiewski Date: Thu, 24 Sep 2026 11:59:02 +0000 Subject: [PATCH 1/3] fix(mcp): require a bearer token and stop serving MCP by default MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The MCP HTTP server acted with the operator's Kubernetes authority but accepted unauthenticated requests on :8090, and --app=all started it by default. - --app=all no longer starts the MCP server; it runs only with --app=mcp-http. - --app=mcp-http listens on 127.0.0.1:8090 only. - --mcp-token-file is mandatory for --app=mcp-http. A missing, unreadable, empty, short, or multi-token file fails MCP startup before any cluster access; other modes reject the flag. - Every request except exact /healthz and /readyz must send "Authorization: Bearer " (SHA-256 digest, constant-time compare). This covers initialize, requests on an existing session (POST, GET stream reconnect/resume, DELETE), and unknown paths. - Remove deploy/mcp-service.yaml and the 8090 container port; document the trust boundary (the token is a shared administrative credential for the server's tool powers, not per-caller RBAC). Signed-off-by: Thomas Kosiewski --- _Generated with `xum` • Model: `anthropic:claude-opus-5-5` • Thinking: `high`_ Change-Id: Ide3808b6ba40c83dce6df2379ef55ca30f7b6b7d --- .cspell.json | 3 +- README.md | 4 +- app_dispatch.go | 13 +- config/default/controller-mode-patch.yaml | 2 - deploy/deployment.yaml | 2 - deploy/mcp-service.yaml | 13 - docs/explanation/architecture.md | 11 +- docs/how-to/deploy-controller.md | 4 +- docs/how-to/mcp-server.md | 60 +++-- docs/how-to/troubleshooting.md | 2 +- docs/index.md | 2 +- examples/cloudnativepg/README.md | 2 +- internal/app/allapp/allapp.go | 61 ++--- internal/app/allapp/allapp_test.go | 37 +++ internal/app/mcpapp/http.go | 106 +++++++-- internal/app/mcpapp/http_auth_test.go | 278 ++++++++++++++++++++++ internal/app/mcpapp/http_test.go | 20 +- main_test.go | 20 +- 18 files changed, 527 insertions(+), 113 deletions(-) delete mode 100644 deploy/mcp-service.yaml create mode 100644 internal/app/mcpapp/http_auth_test.go diff --git a/.cspell.json b/.cspell.json index ec8265e8..bd8d7d7e 100644 --- a/.cspell.json +++ b/.cspell.json @@ -52,7 +52,8 @@ "pooler", "finalizer", "superfences", - "tolerations" + "tolerations", + "portforward" ], "ignorePaths": [ ".git/**", diff --git a/README.md b/README.md index 2e60a865..b9e70b2b 100644 --- a/README.md +++ b/README.md @@ -23,10 +23,10 @@ Pick what runs with `--app`: | `--app` | Runs | | --- | --- | -| `all` (default) | Everything in one process | +| `all` (default) | Operator and aggregated API server in one process | | `controller` | Operator only | | `aggregated-apiserver` | Aggregated API server only | -| `mcp-http` | MCP server only | +| `mcp-http` | MCP server only (never part of `all`; needs `--mcp-token-file`) | ## Quick start diff --git a/app_dispatch.go b/app_dispatch.go index 872b4478..8d0231fd 100644 --- a/app_dispatch.go +++ b/app_dispatch.go @@ -36,6 +36,7 @@ func run(args []string) error { coderSessionToken string coderNamespace string coderRequestTimeout time.Duration + mcpTokenFile string ) fs.StringVar(&appMode, "app", "all", "Application mode (all, controller, aggregated-apiserver, mcp-http)") fs.StringVar( @@ -62,6 +63,12 @@ func run(args []string) error { 30*time.Second, "Timeout for Coder SDK API requests", ) + fs.StringVar( + &mcpTokenFile, + "mcp-token-file", + "", + "Path to a file holding the bearer token that every MCP HTTP request must present (required for --app=mcp-http)", + ) if err := fs.Parse(args); err != nil { return err } @@ -83,6 +90,10 @@ func run(args []string) error { } } + if mcpTokenFile != "" && appMode != "mcp-http" { + return fmt.Errorf("--mcp-token-file is only used with --app=mcp-http; --app=%s does not run the MCP server", appMode) + } + switch appMode { case "all": return runAllApp(setupSignalHandler(), coderRequestTimeout) @@ -97,7 +108,7 @@ func run(args []string) error { } return runAggregatedAPIServerApp(setupSignalHandler(), opts) case "mcp-http": - return runMCPHTTPApp(setupSignalHandler()) + return runMCPHTTPApp(setupSignalHandler(), mcpTokenFile) default: return fmt.Errorf("assertion failed: unsupported --app value %q; must be one of: %s", appMode, supportedAppModes) } diff --git a/config/default/controller-mode-patch.yaml b/config/default/controller-mode-patch.yaml index 98b58cfd..1f85510b 100644 --- a/config/default/controller-mode-patch.yaml +++ b/config/default/controller-mode-patch.yaml @@ -14,5 +14,3 @@ spec: ports: - containerPort: 6443 $patch: delete - - containerPort: 8090 - $patch: delete diff --git a/deploy/deployment.yaml b/deploy/deployment.yaml index c21de71e..8e6682f1 100644 --- a/deploy/deployment.yaml +++ b/deploy/deployment.yaml @@ -22,8 +22,6 @@ spec: name: health - containerPort: 6443 name: https - - containerPort: 8090 - name: mcp livenessProbe: httpGet: path: /healthz diff --git a/deploy/mcp-service.yaml b/deploy/mcp-service.yaml deleted file mode 100644 index aceb30dc..00000000 --- a/deploy/mcp-service.yaml +++ /dev/null @@ -1,13 +0,0 @@ -apiVersion: v1 -kind: Service -metadata: - name: coder-k8s - namespace: coder-system -spec: - selector: - app: coder-k8s - ports: - - name: mcp - port: 8090 - protocol: TCP - targetPort: 8090 diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index c98360dc..b98ad007 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -4,25 +4,23 @@ | `--app` | Runs | | --- | --- | -| `all` (default) | Controller, aggregated API server, and MCP server in one process | +| `all` (default) | Controller and aggregated API server in one process | | `controller` | Controller-runtime manager and reconcilers | | `aggregated-apiserver` | Aggregated API server (`aggregation.coder.com/v1alpha1`) | -| `mcp-http` | MCP HTTP server | +| `mcp-http` | MCP HTTP server (only in this mode; requires `--mcp-token-file`) | ## All-in-one mode -In `all` mode, `internal/app/allapp` creates one controller-runtime manager with one shared cache. It registers the reconcilers, then starts the aggregated API server and MCP server as non-leader runnables. One process, one cache, coordinated startup. +In `all` mode, `internal/app/allapp` creates one controller-runtime manager with one shared cache. It registers the reconcilers, then starts the aggregated API server as a non-leader runnable. The MCP server is not part of `all` mode because it acts with the operator's authority; run it on purpose with `--app=mcp-http`. One process, one cache, coordinated startup. ```mermaid graph TD entry["coder-k8s (--app=all)"] --> mgr["controller-runtime manager"] mgr --> ctrl["Controller reconcilers"] mgr --> agg["Aggregated API server runnable"] - mgr --> mcp["MCP HTTP runnable"] ctrl --> crds["coder.com/v1alpha1 CRDs"] agg --> api["aggregation.coder.com/v1alpha1"] - mcp --> tools["MCP tools over /mcp"] ``` ## Components @@ -30,7 +28,7 @@ graph TD | | Controller | Aggregated API server | MCP server | | --- | --- | --- | --- | | **Code** | `internal/app/controllerapp/`, `internal/controller/` | `internal/app/apiserverapp/`, `internal/aggregated/storage/`, `internal/aggregated/coder/` | `internal/app/mcpapp/` | -| **Listens on** | `:8081` (`/healthz`, `/readyz`) | `:6443` HTTPS (default) | `:8090` (`/mcp`, `/healthz`, `/readyz`) | +| **Listens on** | `:8081` (`/healthz`, `/readyz`) | `:6443` HTTPS (default) | `127.0.0.1:8090` only (`/mcp` needs the bearer token; `/healthz`, `/readyz` do not) | | **Resources** | `CoderControlPlane`, `CoderProvisioner`, `CoderWorkspaceProxy` | `coderworkspaces`, `codertemplates` | Tools for control planes, templates, workspaces, events, pod logs, and run state | ### Controller @@ -64,4 +62,3 @@ graph TD | `config/rbac/` | ServiceAccount, `manager-role`, and bindings (including auth-delegator) | | `deploy/deployment.yaml` | The `coder-k8s` Deployment (defaults to `--app=all`) | | `deploy/apiserver-service.yaml`, `deploy/apiserver-apiservice.yaml` | Expose the aggregated API | -| `deploy/mcp-service.yaml` | MCP Service on port `8090` | diff --git a/docs/how-to/deploy-controller.md b/docs/how-to/deploy-controller.md index 8b4fb176..429580ed 100644 --- a/docs/how-to/deploy-controller.md +++ b/docs/how-to/deploy-controller.md @@ -41,10 +41,10 @@ kubectl get codercontrolplanes -A ## Want everything instead? -Skip the `kubectl patch` step to keep `--app=all`, and apply the extra Services: +Skip the `kubectl patch` step to keep `--app=all` (operator plus aggregated API server), and register the aggregated API: ```bash -kubectl apply -f deploy/apiserver-service.yaml -f deploy/apiserver-apiservice.yaml -f deploy/mcp-service.yaml +kubectl apply -f deploy/apiserver-service.yaml -f deploy/apiserver-apiservice.yaml ``` ## Connect an external PostgreSQL database diff --git a/docs/how-to/mcp-server.md b/docs/how-to/mcp-server.md index 399ad70e..e9819e8a 100644 --- a/docs/how-to/mcp-server.md +++ b/docs/how-to/mcp-server.md @@ -2,29 +2,60 @@ The MCP server gives MCP clients (such as AI agents) tools to inspect and operate `coder-k8s` resources over HTTP. -It runs by default: `--app=all` in `deploy/deployment.yaml` includes it. To run it alone, use `--app=mcp-http`. +It is off by default. `--app=all` does not include it. It runs only with `--app=mcp-http`, and only with a bearer token file. -!!! danger "No authentication" - Any client that can reach the endpoint can call every tool with the server's Kubernetes permissions. RBAC limits what the server can do; it does not identify callers. Keep it on trusted networks. Never expose port `8090` to untrusted clients. Remote access needs its own authentication layer and network restrictions. +!!! danger "What the token grants" + The MCP server calls every tool with its own Kubernetes ServiceAccount's permissions (in the default install, the operator's). The token is a shared administrative credential for those tool powers. It does not identify callers and there is no per-caller RBAC: anyone who holds the token can use every tool. -## 1. Deploy and connect +The server listens on `127.0.0.1:8090` inside the Pod only. To reach it, a client needs the token **and** a way into the Pod's loopback interface, such as `kubectl port-forward` (which Kubernetes authorizes with `pods/portforward` on that Pod). Other containers in the same Pod can reach it too. Do not add a proxy that exposes it on the Pod network. -From a clone of this repository: +## 1. Create the token + +The token file must hold one token of at least 32 bytes without whitespace (a trailing newline is ignored). + +The server reads the token file once, at startup. To rotate the token, update the Secret, then restart the pod (for example, `kubectl -n coder-system rollout restart deployment/coder-k8s`). Until the restart, the old token keeps working. ```bash -kubectl apply -f config/rbac/ -f deploy/deployment.yaml -f deploy/mcp-service.yaml -kubectl port-forward svc/coder-k8s -n coder-system 8090:8090 +openssl rand -hex 32 > mcp-token +chmod 600 mcp-token +kubectl -n coder-system create secret generic coder-k8s-mcp-token --from-file=token=mcp-token ``` -Point your MCP client at: +## 2. Run the MCP server next to the operator + +From a clone of this repository, deploy the operator, then append an `mcp` container that runs `--app=mcp-http`. The JSON patch keeps the operator as the first container, so other patches that address `containers/0` still target it: -```text -http://127.0.0.1:8090/mcp +```bash +kubectl apply -f config/rbac/ -f deploy/deployment.yaml + +kubectl -n coder-system patch deployment coder-k8s --type=json -p '[ + {"op": "add", "path": "/spec/template/spec/volumes", "value": [ + {"name": "mcp-token", "secret": {"secretName": "coder-k8s-mcp-token"}} + ]}, + {"op": "add", "path": "/spec/template/spec/containers/-", "value": { + "name": "mcp", + "image": "ghcr.io/coder/coder-k8s:latest", + "args": ["--app=mcp-http", "--mcp-token-file=/etc/coder-k8s-mcp/token"], + "volumeMounts": [{"name": "mcp-token", "mountPath": "/etc/coder-k8s-mcp", "readOnly": true}] + }} +]' ``` -In-cluster clients use `coder-k8s.coder-system.svc` on port `8090`. +Use the same image as the operator container. + +The `mcp` container exits at startup if `--mcp-token-file` is missing, unreadable, empty, or too short. The operator container is not affected. + +## 3. Connect + +```bash +kubectl -n coder-system port-forward deploy/coder-k8s 8090:8090 +``` + +Point your MCP client at `http://127.0.0.1:8090/mcp` and send `Authorization: Bearer ` on every request. Requests without the correct token get `401 Unauthorized`, including requests that carry an existing `Mcp-Session-Id`. + +## 4. Check health -## 2. Check health +Only these two paths answer without the token: ```bash curl -fsS http://127.0.0.1:8090/healthz @@ -42,7 +73,7 @@ curl -fsS http://127.0.0.1:8090/readyz ## Request rules -The MCP Go SDK (1.4.1) enforces these transport checks. They are protections, not authentication, and they can cause surprising errors: +The bearer token check runs first. After it, the MCP Go SDK (1.4.1) enforces these transport checks. They are protections, not authentication, and they can cause surprising errors: | Rule | Error when broken | | --- | --- | @@ -52,7 +83,6 @@ The MCP Go SDK (1.4.1) enforces these transport checks. They are protections, no Also note: -- The loopback check uses the address the server sees. It is not a hostname allowlist for Pod or Service traffic, and not every port-forward or proxy path arrives from loopback. - JSON field names are case-sensitive. A key with an appended null character does not alias the original key. -- Service names and client-supplied headers are not credentials. +- Service names, session IDs, and other client-supplied headers are not credentials. Only the bearer token is. - Do not disable the SDK's localhost or cross-origin protections to work around routing problems. diff --git a/docs/how-to/troubleshooting.md b/docs/how-to/troubleshooting.md index 58e1ab04..7cb9bf82 100644 --- a/docs/how-to/troubleshooting.md +++ b/docs/how-to/troubleshooting.md @@ -11,7 +11,7 @@ kubectl logs -n coder-system deploy/coder-k8s `--app` is optional and defaults to `all` (every component). To isolate one component, set it explicitly: ```bash -GOFLAGS=-mod=vendor go run . --app=controller # or aggregated-apiserver, mcp-http +GOFLAGS=-mod=vendor go run . --app=controller # or aggregated-apiserver, or mcp-http --mcp-token-file= ``` An unknown value fails at startup with `assertion failed: unsupported --app value ...`. diff --git a/docs/index.md b/docs/index.md index 1687fa9d..274a8168 100644 --- a/docs/index.md +++ b/docs/index.md @@ -11,7 +11,7 @@ Run and manage [Coder](https://coder.com) with native Kubernetes APIs. | `--app` | Runs | Resources | | --- | --- | --- | -| `all` (default) | All three components in one process | Everything below | +| `all` (default) | Operator and aggregated API server in one process (not the MCP server) | Everything below except MCP | | `controller` | Operator | `CoderControlPlane`, `CoderProvisioner`, `CoderWorkspaceProxy` (`coder.com/v1alpha1`) | | `aggregated-apiserver` | Aggregated API server | `CoderWorkspace`, `CoderTemplate` (`aggregation.coder.com/v1alpha1`) | | `mcp-http` | MCP server | Operational tools over HTTP | diff --git a/examples/cloudnativepg/README.md b/examples/cloudnativepg/README.md index a9e47875..ccd78c57 100644 --- a/examples/cloudnativepg/README.md +++ b/examples/cloudnativepg/README.md @@ -36,7 +36,7 @@ kubectl apply -f deploy/deployment.yaml kubectl rollout status deployment/coder-k8s -n coder-system ``` -`deploy/deployment.yaml` defaults to `--app=all`, which runs the controller, aggregated API server, and MCP server in a single pod. For split deployments, you can set `--app=controller`, `--app=aggregated-apiserver`, or `--app=mcp-http` in the Deployment args. +`deploy/deployment.yaml` defaults to `--app=all`, which runs the controller and aggregated API server in a single pod. The MCP server never runs in `--app=all`; see [Run the MCP server](../../docs/how-to/mcp-server.md). For split deployments, you can set `--app=controller`, `--app=aggregated-apiserver`, or `--app=mcp-http` in the Deployment args. ## 3. Deploy this example diff --git a/internal/app/allapp/allapp.go b/internal/app/allapp/allapp.go index d7b68acb..3f4b4ade 100644 --- a/internal/app/allapp/allapp.go +++ b/internal/app/allapp/allapp.go @@ -1,4 +1,7 @@ -// Package allapp composes controller, aggregated API server, and MCP app modes in one process. +// Package allapp composes the controller and aggregated API server app modes in one process. +// +// The MCP HTTP server is deliberately not part of this mode: it acts with the operator's authority and +// only runs when explicitly requested with --app=mcp-http and a bearer token file. package allapp import ( @@ -6,14 +9,12 @@ import ( "fmt" "time" - "k8s.io/client-go/kubernetes" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/manager" "github.com/coder/coder-k8s/internal/aggregated/coder" "github.com/coder/coder-k8s/internal/app/apiserverapp" "github.com/coder/coder-k8s/internal/app/controllerapp" - "github.com/coder/coder-k8s/internal/app/mcpapp" "github.com/coder/coder-k8s/internal/app/sharedscheme" ) @@ -29,8 +30,6 @@ var ( runAggregatedAPIServer = func(ctx context.Context, opts apiserverapp.Options) error { return apiserverapp.RunWithOptions(ctx, opts) } - runMCPHTTPWithClients = mcpapp.RunHTTPWithClients - newClientset = kubernetes.NewForConfig ) var _ manager.LeaderElectionRunnable = nonLeaderRunnable{} @@ -89,6 +88,22 @@ func Run(ctx context.Context, coderRequestTimeout time.Duration) error { return err } + if err := addRunnables(mgr, requestTimeout); err != nil { + return err + } + + return mgr.Start(ctx) +} + +// addRunnables registers the non-controller runnables that share the manager's cache. +func addRunnables(mgr manager.Manager, requestTimeout time.Duration) error { + if mgr == nil { + return fmt.Errorf("assertion failed: manager must not be nil") + } + if requestTimeout <= 0 { + return fmt.Errorf("assertion failed: request timeout must be positive: %s", requestTimeout) + } + if err := mgr.Add(nonLeaderRunnable{ run: func(runnableCtx context.Context) error { if runnableCtx == nil { @@ -126,41 +141,7 @@ func Run(ctx context.Context, coderRequestTimeout time.Duration) error { return fmt.Errorf("add aggregated-apiserver runnable: %w", err) } - if err := mgr.Add(nonLeaderRunnable{ - run: func(runnableCtx context.Context) error { - if runnableCtx == nil { - return fmt.Errorf("assertion failed: context must not be nil") - } - - if err := waitForCacheSync(runnableCtx, mgr, "mcp-http"); err != nil { - return err - } - - managerClient := mgr.GetClient() - if managerClient == nil { - return fmt.Errorf("assertion failed: manager client is nil") - } - - managerConfig := mgr.GetConfig() - if managerConfig == nil { - return fmt.Errorf("assertion failed: manager config is nil") - } - - clientset, err := newClientset(managerConfig) - if err != nil { - return fmt.Errorf("build Kubernetes clientset: %w", err) - } - if clientset == nil { - return fmt.Errorf("assertion failed: Kubernetes clientset is nil after successful construction") - } - - return runMCPHTTPWithClients(runnableCtx, managerClient, clientset) - }, - }); err != nil { - return fmt.Errorf("add mcp-http runnable: %w", err) - } - - return mgr.Start(ctx) + return nil } func waitForCacheSync(ctx context.Context, mgr manager.Manager, runnableName string) error { diff --git a/internal/app/allapp/allapp_test.go b/internal/app/allapp/allapp_test.go index 24fea250..ad221ca6 100644 --- a/internal/app/allapp/allapp_test.go +++ b/internal/app/allapp/allapp_test.go @@ -6,6 +6,8 @@ import ( "strings" "testing" "time" + + "sigs.k8s.io/controller-runtime/pkg/manager" ) func TestRunRejectsNilContext(t *testing.T) { @@ -65,3 +67,38 @@ func TestNonLeaderRunnableStartRequiresRunFunction(t *testing.T) { t.Fatalf("unexpected error: %v", err) } } + +type recordingManager struct { + manager.Manager + added []manager.Runnable +} + +func (m *recordingManager) Add(r manager.Runnable) error { + m.added = append(m.added, r) + return nil +} + +// TestAddRunnablesRegistersOnlyAggregatedAPIServer guards the security decision that --app=all never +// starts the MCP HTTP server. +func TestAddRunnablesRegistersOnlyAggregatedAPIServer(t *testing.T) { + mgr := &recordingManager{} + if err := addRunnables(mgr, 30*time.Second); err != nil { + t.Fatal(err) + } + if len(mgr.added) != 1 { + t.Fatalf("expected exactly one runnable (aggregated-apiserver), got %d", len(mgr.added)) + } + runnable, ok := mgr.added[0].(nonLeaderRunnable) + if !ok || runnable.run == nil { + t.Fatalf("expected a nonLeaderRunnable with a run function, got %T", mgr.added[0]) + } +} + +func TestAddRunnablesRejectsInvalidArguments(t *testing.T) { + if err := addRunnables(nil, time.Second); err == nil || !strings.Contains(err.Error(), "manager must not be nil") { + t.Fatalf("expected nil manager assertion, got %v", err) + } + if err := addRunnables(&recordingManager{}, 0); err == nil || !strings.Contains(err.Error(), "request timeout must be positive") { + t.Fatalf("expected request timeout assertion, got %v", err) + } +} diff --git a/internal/app/mcpapp/http.go b/internal/app/mcpapp/http.go index 6b3b1117..a9428dbb 100644 --- a/internal/app/mcpapp/http.go +++ b/internal/app/mcpapp/http.go @@ -1,10 +1,15 @@ package mcpapp import ( + "bytes" "context" + "crypto/sha256" + "crypto/subtle" "errors" "fmt" "net/http" + "os" + "strings" "time" "github.com/modelcontextprotocol/go-sdk/mcp" @@ -14,8 +19,11 @@ import ( ) const ( - // DefaultHTTPAddr is the default listen address used by MCP HTTP mode. - DefaultHTTPAddr = ":8090" + // DefaultHTTPAddr is the listen address used by MCP HTTP mode. It is loopback-only on purpose: + // the server acts with its ServiceAccount's authority, so it must never listen on the Pod network. + DefaultHTTPAddr = "127.0.0.1:8090" + // MinTokenLength is the minimum accepted length of the MCP bearer token, in bytes. + MinTokenLength = 32 // streamableHTTPSessionTimeout ensures abandoned MCP streamable HTTP sessions are reclaimed. streamableHTTPSessionTimeout = 15 * time.Minute ) @@ -34,22 +42,94 @@ func newMCPHTTPHandler(server *mcp.Server) *mcp.StreamableHTTPHandler { }) } +// newMCPHTTPMux builds the production HTTP routing. Only the exact health paths are unauthenticated; +// every other request, including every /mcp request that carries an existing session ID, must present +// the bearer token. +func newMCPHTTPMux(mcpHandler http.Handler, token []byte) http.Handler { + if mcpHandler == nil { + panic("assertion failed: MCP handler must not be nil") + } + if len(token) < MinTokenLength { + panic("assertion failed: MCP bearer token must be validated before building the handler") + } + + mux := http.NewServeMux() + mux.Handle("/mcp", mcpHandler) + health := func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusOK) + _, _ = w.Write([]byte("ok")) + } + + tokenDigest := sha256.Sum256(token) + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path == "/healthz" || r.URL.Path == "/readyz" { + health(w, r) + return + } + if !bearerTokenMatches(r.Header.Get("Authorization"), tokenDigest) { + w.Header().Set("WWW-Authenticate", `Bearer realm="coder-k8s-mcp"`) + http.Error(w, "unauthorized", http.StatusUnauthorized) + return + } + mux.ServeHTTP(w, r) + }) +} + +// bearerTokenMatches compares SHA-256 digests in constant time so neither the token bytes nor its +// length leak through timing. +func bearerTokenMatches(header string, want [sha256.Size]byte) bool { + scheme, credential, ok := strings.Cut(header, " ") + if !ok || !strings.EqualFold(scheme, "Bearer") || credential == "" { + return false + } + got := sha256.Sum256([]byte(credential)) + return subtle.ConstantTimeCompare(got[:], want[:]) == 1 +} + +// ReadTokenFile reads and validates the MCP bearer token file. +func ReadTokenFile(path string) ([]byte, error) { + if strings.TrimSpace(path) == "" { + return nil, fmt.Errorf("--mcp-token-file is required for --app=mcp-http") + } + data, err := os.ReadFile(path) //nolint:gosec // G304: the path is the operator-supplied --mcp-token-file flag. + if err != nil { + return nil, fmt.Errorf("read --mcp-token-file: %w", err) + } + token := bytes.TrimSuffix(bytes.TrimSuffix(data, []byte("\n")), []byte("\r")) + if len(token) == 0 { + return nil, fmt.Errorf("--mcp-token-file %q is empty", path) + } + if len(token) < MinTokenLength { + return nil, fmt.Errorf("--mcp-token-file %q holds %d bytes; at least %d are required", path, len(token), MinTokenLength) + } + if bytes.ContainsAny(token, " \t\r\n") { + return nil, fmt.Errorf("--mcp-token-file %q must hold a single token without whitespace", path) + } + return token, nil +} + // RunHTTP starts the MCP server using streamable HTTP transport. -func RunHTTP(ctx context.Context) error { +func RunHTTP(ctx context.Context, tokenFile string) error { if ctx == nil { return fmt.Errorf("assertion failed: context must not be nil") } + // Validate the credential before touching the cluster so a bad token fails fast. + token, err := ReadTokenFile(tokenFile) + if err != nil { + return err + } + k8sClient, clientset, err := newClients() if err != nil { return err } - return RunHTTPWithClients(ctx, k8sClient, clientset) + return RunHTTPWithClients(ctx, k8sClient, clientset, token) } // RunHTTPWithClients starts the MCP server using streamable HTTP transport and the provided Kubernetes clients. -func RunHTTPWithClients(ctx context.Context, k8sClient client.Client, clientset kubernetes.Interface) error { +func RunHTTPWithClients(ctx context.Context, k8sClient client.Client, clientset kubernetes.Interface, token []byte) error { if ctx == nil { return fmt.Errorf("assertion failed: context must not be nil") } @@ -59,6 +139,9 @@ func RunHTTPWithClients(ctx context.Context, k8sClient client.Client, clientset if clientset == nil { return fmt.Errorf("assertion failed: Kubernetes clientset must not be nil") } + if len(token) < MinTokenLength { + return fmt.Errorf("assertion failed: MCP bearer token must hold at least %d bytes", MinTokenLength) + } server := NewServer(k8sClient, clientset) if server == nil { @@ -70,20 +153,9 @@ func RunHTTPWithClients(ctx context.Context, k8sClient client.Client, clientset return fmt.Errorf("assertion failed: MCP HTTP handler is nil after successful construction") } - mux := http.NewServeMux() - mux.Handle("/mcp", mcpHandler) - mux.HandleFunc("/healthz", func(w http.ResponseWriter, _ *http.Request) { - w.WriteHeader(http.StatusOK) - _, _ = w.Write([]byte("ok")) - }) - mux.HandleFunc("/readyz", func(w http.ResponseWriter, _ *http.Request) { - w.WriteHeader(http.StatusOK) - _, _ = w.Write([]byte("ok")) - }) - httpServer := &http.Server{ Addr: DefaultHTTPAddr, - Handler: mux, + Handler: newMCPHTTPMux(mcpHandler, token), ReadHeaderTimeout: 5 * time.Second, } diff --git a/internal/app/mcpapp/http_auth_test.go b/internal/app/mcpapp/http_auth_test.go new file mode 100644 index 00000000..11c6f432 --- /dev/null +++ b/internal/app/mcpapp/http_auth_test.go @@ -0,0 +1,278 @@ +package mcpapp + +import ( + "context" + "io" + "net" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/modelcontextprotocol/go-sdk/mcp" +) + +const initializeRequest = `{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":"2025-11-25","capabilities":{},"clientInfo":{"name":"auth-test","version":"1"}}}` + +func TestHTTPBearerTokenRequiredForInitialize(t *testing.T) { + tests := []struct { + name string + authorization string + wantStatus int + }{ + {name: "missing header", wantStatus: http.StatusUnauthorized}, + {name: "basic scheme", authorization: "Basic " + testToken, wantStatus: http.StatusUnauthorized}, + {name: "bearer without credential", authorization: "Bearer ", wantStatus: http.StatusUnauthorized}, + {name: "bare token", authorization: testToken, wantStatus: http.StatusUnauthorized}, + {name: "wrong token", authorization: "Bearer " + strings.Repeat("x", len(testToken)), wantStatus: http.StatusUnauthorized}, + {name: "token prefix", authorization: "Bearer " + testToken[:len(testToken)-1], wantStatus: http.StatusUnauthorized}, + {name: "token with suffix", authorization: "Bearer " + testToken + "x", wantStatus: http.StatusUnauthorized}, + {name: "token with extra field", authorization: "Bearer " + testToken + " extra", wantStatus: http.StatusUnauthorized}, + {name: "correct token", authorization: "Bearer " + testToken, wantStatus: http.StatusOK}, + {name: "case-insensitive scheme", authorization: "bearer " + testToken, wantStatus: http.StatusOK}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + k8sClient, server := newHTTPTestServer(t, false) + ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) + defer cancel() + + status, headers, body := postMCPWithAuth(ctx, t, server, tt.authorization, "", "", "", "", "application/json", initializeRequest) + if status != tt.wantStatus { + t.Fatalf("status=%d, want %d; body=%s", status, tt.wantStatus, body) + } + if tt.wantStatus == http.StatusUnauthorized { + if headers.Get("Mcp-Session-Id") != "" { + t.Fatalf("unauthenticated initialize must not create a session, got %q", headers.Get("Mcp-Session-Id")) + } + if !strings.HasPrefix(headers.Get("WWW-Authenticate"), "Bearer") { + t.Fatalf("expected Bearer challenge, got %q", headers.Get("WWW-Authenticate")) + } + } + if got := k8sClient.gets.Load(); got != 0 { + t.Fatalf("initialize must not read resources, got %d reads", got) + } + }) + } +} + +// TestHTTPSessionReuseRequiresToken proves that an existing session ID is never a credential: +// every POST, GET (stream reconnect/resume), and DELETE on an established session needs the token. +func TestHTTPSessionReuseRequiresToken(t *testing.T) { + k8sClient, server := newHTTPTestServer(t, false) + ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) + defer cancel() + + status, headers, body := postMCP(ctx, t, server, "", "", "", "", "application/json", initializeRequest) + sessionID := headers.Get("Mcp-Session-Id") + if status != http.StatusOK || sessionID == "" { + t.Fatalf("initialize: status=%d session=%q body=%s", status, sessionID, body) + } + status, _, body = postMCP(ctx, t, server, sessionID, "", "", "", "application/json", `{"jsonrpc":"2.0","method":"notifications/initialized"}`) + if status != http.StatusAccepted { + t.Fatalf("initialized: status=%d body=%s", status, body) + } + + for _, authorization := range []string{"", "Bearer wrong-" + testToken, "Basic " + testToken} { + status, _, body = postMCPWithAuth(ctx, t, server, authorization, sessionID, "", "", "", "application/json", workspaceCall) + if status != http.StatusUnauthorized { + t.Fatalf("tools/call with session and authorization %q: status=%d body=%s", authorization, status, body) + } + for _, method := range []string{http.MethodGet, http.MethodDelete} { + if status := doSessionRequest(ctx, t, server, method, sessionID, authorization); status != http.StatusUnauthorized { + t.Fatalf("%s with session and authorization %q: status=%d", method, authorization, status) + } + } + } + if got := k8sClient.gets.Load(); got != 0 { + t.Fatalf("unauthenticated session requests read %d resources", got) + } + assertWorkspaceRunning(t, k8sClient.Client, false) + + // The rejected DELETE must not have ended the session: the token holder can still use it. + status, _, body = postMCP(ctx, t, server, sessionID, "", "", "", "application/json", workspaceCall) + if status != http.StatusOK || k8sClient.gets.Load() != 1 { + t.Fatalf("authenticated follow-up: status=%d reads=%d body=%s", status, k8sClient.gets.Load(), body) + } + assertWorkspaceRunning(t, k8sClient.Client, true) + + if status := doSessionRequest(ctx, t, server, http.MethodDelete, sessionID, "Bearer "+testToken); status != http.StatusNoContent { + t.Fatalf("authenticated delete: status=%d", status) + } +} + +func TestHTTPOnlyExactHealthPathsAreUnauthenticated(t *testing.T) { + _, server := newHTTPTestServer(t, false) + for path, want := range map[string]int{ + "/healthz": http.StatusOK, + "/readyz": http.StatusOK, + "/healthz/": http.StatusUnauthorized, + "/readyz/x": http.StatusUnauthorized, + "/mcp": http.StatusUnauthorized, + "/mcp/": http.StatusUnauthorized, + "/MCP": http.StatusUnauthorized, + "/": http.StatusUnauthorized, + "/healthz/../mcp": http.StatusUnauthorized, + } { + t.Run(path, func(t *testing.T) { + req, err := http.NewRequestWithContext(t.Context(), http.MethodGet, server.URL, nil) + if err != nil { + t.Fatal(err) + } + req.URL.Path = path + req.URL.RawPath = path + resp, err := server.Client().Do(req) + if err != nil { + t.Fatal(err) + } + _, _ = io.Copy(io.Discard, resp.Body) + _ = resp.Body.Close() + if resp.StatusCode != want { + t.Fatalf("GET %s without token: status=%d, want %d", path, resp.StatusCode, want) + } + }) + } +} + +func TestHTTPTransportSDKClientWithoutTokenFails(t *testing.T) { + k8sClient, server := newHTTPTestServer(t, false) + ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) + defer cancel() + mcpClient := mcp.NewClient(&mcp.Implementation{Name: "test", Version: "1"}, nil) + session, err := mcpClient.Connect(ctx, &mcp.StreamableClientTransport{Endpoint: server.URL + "/mcp", HTTPClient: server.Client()}, nil) + if err == nil { + _ = session.Close() + t.Fatal("expected Connect without token to fail") + } + if got := k8sClient.gets.Load(); got != 0 { + t.Fatalf("unauthenticated SDK client read %d resources", got) + } +} + +func TestDefaultHTTPAddrIsLoopback(t *testing.T) { + host, port, err := net.SplitHostPort(DefaultHTTPAddr) + if err != nil { + t.Fatal(err) + } + ip := net.ParseIP(host) + if ip == nil || !ip.IsLoopback() { + t.Fatalf("DefaultHTTPAddr host must be a loopback IP literal, got %q", host) + } + if port != "8090" { + t.Fatalf("DefaultHTTPAddr port changed unexpectedly: %q", port) + } +} + +func TestReadTokenFile(t *testing.T) { + dir := t.TempDir() + write := func(name, content string) string { + path := filepath.Join(dir, name) + if err := os.WriteFile(path, []byte(content), 0o600); err != nil { + t.Fatal(err) + } + return path + } + tests := []struct { + name string + path string + want string + wantError string + }{ + {name: "flag missing", path: "", wantError: "--mcp-token-file is required"}, + {name: "file missing", path: filepath.Join(dir, "absent"), wantError: "read --mcp-token-file"}, + {name: "directory", path: dir, wantError: "read --mcp-token-file"}, + {name: "empty file", path: write("empty", ""), wantError: "is empty"}, + {name: "newline only", path: write("newline", "\n"), wantError: "is empty"}, + {name: "too short", path: write("short", "abc\n"), wantError: "at least 32 are required"}, + {name: "inner whitespace", path: write("space", strings.Repeat("a", 20)+" "+strings.Repeat("b", 20)), wantError: "without whitespace"}, + {name: "two lines", path: write("lines", testToken+"\n"+testToken+"\n"), wantError: "without whitespace"}, + {name: "trailing newline", path: write("lf", testToken+"\n"), want: testToken}, + {name: "trailing CRLF", path: write("crlf", testToken+"\r\n"), want: testToken}, + {name: "no newline", path: write("plain", testToken), want: testToken}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + token, err := ReadTokenFile(tt.path) + if tt.wantError != "" { + if err == nil || !strings.Contains(err.Error(), tt.wantError) { + t.Fatalf("expected error containing %q, got token=%q err=%v", tt.wantError, token, err) + } + return + } + if err != nil { + t.Fatal(err) + } + if string(token) != tt.want { + t.Fatalf("token=%q, want %q", token, tt.want) + } + }) + } +} + +func TestReadTokenFileUnreadable(t *testing.T) { + if os.Geteuid() == 0 { + t.Skip("root can read mode 000 files") + } + path := filepath.Join(t.TempDir(), "token") + if err := os.WriteFile(path, []byte(testToken), 0o000); err != nil { + t.Fatal(err) + } + if _, err := ReadTokenFile(path); err == nil || !strings.Contains(err.Error(), "read --mcp-token-file") { + t.Fatalf("expected unreadable token file error, got %v", err) + } +} + +// TestRunHTTPFailsBeforeClusterAccessWithoutToken runs the real entry point: token validation must +// fail before any Kubernetes client is built. +func TestRunHTTPFailsBeforeClusterAccessWithoutToken(t *testing.T) { + t.Setenv("KUBECONFIG", filepath.Join(t.TempDir(), "must-not-be-read")) + if err := RunHTTP(t.Context(), ""); err == nil || !strings.Contains(err.Error(), "--mcp-token-file is required") { + t.Fatalf("expected missing token error, got %v", err) + } + if err := RunHTTPWithClients(t.Context(), nil, nil, nil); err == nil { + t.Fatal("expected assertion error for nil clients") + } +} + +func doSessionRequest(ctx context.Context, t *testing.T, server *httptest.Server, method, sessionID, authorization string) int { + t.Helper() + req, err := http.NewRequestWithContext(ctx, method, server.URL+"/mcp", nil) + if err != nil { + t.Fatal(err) + } + req.Header.Set("Mcp-Session-Id", sessionID) + req.Header.Set("Mcp-Protocol-Version", "2025-11-25") + req.Header.Set("Accept", "text/event-stream") + req.Header.Set("Last-Event-ID", "0") + if authorization != "" { + req.Header.Set("Authorization", authorization) + } + resp, err := server.Client().Do(req) + if err != nil { + t.Fatal(err) + } + _, _ = io.Copy(io.Discard, io.LimitReader(resp.Body, 4096)) + if err := resp.Body.Close(); err != nil { + t.Error(err) + } + return resp.StatusCode +} + +type tokenTransport struct { + token string + base http.RoundTripper +} + +func (tr tokenTransport) RoundTrip(req *http.Request) (*http.Response, error) { + clone := req.Clone(req.Context()) + clone.Header.Set("Authorization", "Bearer "+tr.token) + return tr.base.RoundTrip(clone) +} + +func tokenHTTPClient(server *httptest.Server, token string) *http.Client { + base := server.Client() + return &http.Client{Transport: tokenTransport{token: token, base: base.Transport}, Timeout: base.Timeout} +} diff --git a/internal/app/mcpapp/http_test.go b/internal/app/mcpapp/http_test.go index e11e89ff..5a28f266 100644 --- a/internal/app/mcpapp/http_test.go +++ b/internal/app/mcpapp/http_test.go @@ -18,6 +18,9 @@ import ( "sigs.k8s.io/controller-runtime/pkg/client" ) +// testToken is a fixed, non-secret test credential of the minimum accepted length. +const testToken = "test-token-0123456789abcdef-0123456789" + const workspaceCall = `{"jsonrpc":"2.0","id":2,"method":"tools/call","params":{"name":"set_workspace_running","arguments":{"namespace":"default","name":"dev","running":true}}}` func TestHTTPTransportSecurity(t *testing.T) { @@ -76,6 +79,7 @@ func TestHTTPTransportSecurity(t *testing.T) { t.Fatal(err) } req.Header.Set("Mcp-Session-Id", sessionID) + req.Header.Set("Authorization", "Bearer "+testToken) resp, err := server.Client().Do(req) if err != nil { t.Errorf("delete session: %v", err) @@ -134,7 +138,7 @@ func TestHTTPTransportSDKClient(t *testing.T) { ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() mcpClient := mcp.NewClient(&mcp.Implementation{Name: "test", Version: "1"}, nil) - session, err := mcpClient.Connect(ctx, &mcp.StreamableClientTransport{Endpoint: server.URL + "/mcp", HTTPClient: server.Client()}, nil) + session, err := mcpClient.Connect(ctx, &mcp.StreamableClientTransport{Endpoint: server.URL + "/mcp", HTTPClient: tokenHTTPClient(server, testToken)}, nil) if err != nil { t.Fatal(err) } @@ -161,12 +165,8 @@ func newHTTPTestServer(t *testing.T, podAddress bool) (*observedHTTPClient, *htt k8sClient := &observedHTTPClient{Client: mustNewFakeClient(t, &aggregationv1alpha1.CoderWorkspace{ ObjectMeta: metav1.ObjectMeta{Name: "dev", Namespace: "default"}, })} - handler := newMCPHTTPHandler(NewServer(k8sClient, k8sfake.NewClientset())) + handler := newMCPHTTPMux(newMCPHTTPHandler(NewServer(k8sClient, k8sfake.NewClientset())), []byte(testToken)) server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - if r.URL.Path != "/mcp" { - http.NotFound(w, r) - return - } if podAddress { // Model net/http's Pod-side local address without opening a non-loopback // listener or contacting a cluster. This is not a Kubernetes network test. @@ -183,11 +183,19 @@ func newHTTPTestServer(t *testing.T, podAddress bool) (*observedHTTPClient, *htt } func postMCP(ctx context.Context, t *testing.T, server *httptest.Server, sessionID, host, origin, fetchSite, contentType, body string) (int, http.Header, string) { + t.Helper() + return postMCPWithAuth(ctx, t, server, "Bearer "+testToken, sessionID, host, origin, fetchSite, contentType, body) +} + +func postMCPWithAuth(ctx context.Context, t *testing.T, server *httptest.Server, authorization, sessionID, host, origin, fetchSite, contentType, body string) (int, http.Header, string) { t.Helper() req, err := http.NewRequestWithContext(ctx, http.MethodPost, server.URL+"/mcp", strings.NewReader(body)) if err != nil { t.Fatal(err) } + if authorization != "" { + req.Header.Set("Authorization", authorization) + } req.Host = host req.Header.Set("Accept", "application/json, text/event-stream") if contentType != "omit" { diff --git a/main_test.go b/main_test.go index 118869b3..20b46a61 100644 --- a/main_test.go +++ b/main_test.go @@ -263,15 +263,18 @@ func TestRunDispatchesMCPHTTPMode(t *testing.T) { expectedErr := errors.New("sentinel mcp-http error") called := false - runMCPHTTPApp = func(ctx context.Context) error { + runMCPHTTPApp = func(ctx context.Context, path string) error { called = true if ctx == nil { t.Fatal("expected non-nil context passed to MCP HTTP runner") } + if path != "/etc/mcp/token" { + t.Fatalf("expected token file path to be forwarded, got %q", path) + } return expectedErr } - err := run([]string{"--app=mcp-http"}) + err := run([]string{"--app=mcp-http", "--mcp-token-file=/etc/mcp/token"}) if !called { t.Fatal("expected MCP HTTP runner to be called") } @@ -279,3 +282,16 @@ func TestRunDispatchesMCPHTTPMode(t *testing.T) { t.Fatalf("expected sentinel error %v, got %v", expectedErr, err) } } + +func TestRunRejectsMCPTokenFileOutsideMCPHTTPMode(t *testing.T) { + installMockSignalHandler(t) + + for _, mode := range []string{"all", "controller", "aggregated-apiserver"} { + t.Run(mode, func(t *testing.T) { + err := run([]string{"--app=" + mode, "--mcp-token-file=/etc/mcp/token"}) + if err == nil || !strings.Contains(err.Error(), "--mcp-token-file is only used with --app=mcp-http") { + t.Fatalf("expected --mcp-token-file rejection for --app=%s, got %v", mode, err) + } + }) + } +} From 73046cea3edee2735b2a9b748eb54989adbeb8b2 Mon Sep 17 00:00:00 2001 From: Thomas Kosiewski Date: Thu, 24 Sep 2026 12:33:10 +0000 Subject: [PATCH 2/3] fix(apiserver): require delegated authentication and authorization MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The aggregated API server used an anonymous authenticator and an allow-all authorizer, so any client that reached port 6443 directly was served with the control plane's Coder credentials. - Use the vendored k8s.io/apiserver DelegatingAuthenticationOptions and DelegatingAuthorizationOptions: front-proxy request-header client certificates from kube-system/extension-apiserver-authentication, cluster client certificates, TokenReview for bearer tokens, and SubjectAccessReview for every request. Lookup failures other than NotFound abort startup; delegation outages deny requests. - Refuse the anonymous identity except on exact /healthz, /livez and /readyz. The vendored options always add an anonymous fallback, so without this "no data for anonymous callers" would depend on cluster RBAC. - Resolve one Kubernetes API authority for both authentication and authorization: KUBECONFIG (exactly one valid file, never a fallback), otherwise the in-cluster ServiceAccount, otherwise ~/.kube/config. Missing configuration fails startup. - Keep the vendored system:masters bypass and 10s caches; document them, and the retained trust of the last loaded front-proxy CA. - Tests: fake Kubernetes API (TokenReview, SAR, ConfigMap watch-list) and a real envtest kube-apiserver with request-header flags and RBAC cover forged headers, unknown CAs, SAR attributes for every verb on both resources, fail-closed outages and missing bindings, with zero Coder backend calls asserted on every denial. - CI E2E: check RBAC through kube-apiserver for a non-admin identity, and that direct requests to 6443 from another pod (anonymous, forged headers, ServiceAccount token without RBAC) are rejected and 8090 is not served. Signed-off-by: Thomas Kosiewski --- _Generated with `xum` • Model: `anthropic:claude-opus-5-5` • Thinking: `high`_ Change-Id: If7a288381565727f87f07bc37d678d741d31b6d2 --- .cspell.json | 5 +- .github/workflows/ci.yaml | 44 ++ AGENTS.md | 2 +- docs/explanation/architecture.md | 4 + docs/how-to/deploy-aggregated-apiserver.md | 30 + docs/how-to/troubleshooting.md | 16 +- internal/app/apiserverapp/apiserverapp.go | 68 +- .../app/apiserverapp/apiserverapp_test.go | 12 +- internal/app/apiserverapp/auth.go | 138 ++++ .../app/apiserverapp/auth_envtest_test.go | 310 +++++++++ .../app/apiserverapp/auth_helpers_test.go | 542 ++++++++++++++++ internal/app/apiserverapp/auth_test.go | 591 ++++++++++++++++++ internal/app/apiserverapp/integration_test.go | 34 +- 13 files changed, 1772 insertions(+), 24 deletions(-) create mode 100644 internal/app/apiserverapp/auth.go create mode 100644 internal/app/apiserverapp/auth_envtest_test.go create mode 100644 internal/app/apiserverapp/auth_helpers_test.go create mode 100644 internal/app/apiserverapp/auth_test.go diff --git a/.cspell.json b/.cspell.json index bd8d7d7e..e200d6fb 100644 --- a/.cspell.json +++ b/.cspell.json @@ -53,7 +53,10 @@ "finalizer", "superfences", "tolerations", - "portforward" + "portforward", + "livez", + "requestheader", + "subjectaccessreviews" ], "ignorePaths": [ ".git/**", diff --git a/.github/workflows/ci.yaml b/.github/workflows/ci.yaml index 39c7dc6d..820937d6 100644 --- a/.github/workflows/ci.yaml +++ b/.github/workflows/ci.yaml @@ -372,6 +372,50 @@ jobs: test -n "$TOKEN_SECRET" kubectl -n coder get secret "$TOKEN_SECRET" + # Delegated authentication/authorization of the aggregated API server: + # RBAC via kube-apiserver for a non-admin identity, and direct requests to + # port 6443 from another pod (anonymous, forged front-proxy headers, and a + # ServiceAccount token without RBAC) are rejected. MCP is not served in --app=all. + - name: Verify aggregated API authentication and authorization + if: env.E2E_FULL == 'true' + env: + PROBE_IMAGE: curlimages/curl:8.16.0@sha256:463eaf6072688fe96ac64fa623fe73e1dbe25d8ad6c34404a669ad3ce1f104b6 + run: | + set -euo pipefail + kubectl -n coder create serviceaccount e2e-reader + kubectl -n coder create role e2e-template-reader --verb=get,list --resource=codertemplates.aggregation.coder.com + kubectl -n coder create rolebinding e2e-template-reader --role=e2e-template-reader --serviceaccount=coder:e2e-reader + reader=(--as=system:serviceaccount:coder:e2e-reader) + kubectl "${reader[@]}" -n coder get codertemplates.aggregation.coder.com + expect_forbidden() { + local out + if out=$("$@" 2>&1); then + echo "expected Forbidden, command succeeded: $*" >&2 + return 1 + fi + grep -q Forbidden <<<"$out" || { echo "expected Forbidden, got: $out" >&2; return 1; } + } + expect_forbidden kubectl "${reader[@]}" -n default get codertemplates.aggregation.coder.com + expect_forbidden kubectl "${reader[@]}" -n coder get coderworkspaces.aggregation.coder.com + expect_forbidden kubectl "${reader[@]}" -n coder delete codertemplates.aggregation.coder.com --all --dry-run=server + + kubectl -n default run e2e-authn-probe --image="$PROBE_IMAGE" --restart=Never --command -- sleep 600 + kubectl -n default wait --for=condition=Ready pod/e2e-authn-probe --timeout=180s + probe() { kubectl -n default exec e2e-authn-probe -- sh -c "$1"; } + api=https://coder-k8s-apiserver.coder-system.svc + list="$api/apis/aggregation.coder.com/v1alpha1/namespaces/coder/codertemplates" + code() { probe "curl -sk -o /dev/null -w '%{http_code}' $1"; } + test "$(code "$list")" = 401 + test "$(code "-H 'X-Remote-User: kubernetes-admin' -H 'X-Remote-Group: system:masters' $list")" = 401 + test "$(code "-H \"Authorization: Bearer \$(cat /var/run/secrets/kubernetes.io/serviceaccount/token)\" $list")" = 403 + test "$(code "$api/healthz")" = 200 + pod_ip=$(kubectl -n coder-system get pod -l app=coder-k8s -o jsonpath='{.items[0].status.podIP}') + if probe "curl -s --max-time 5 http://$pod_ip:8090/healthz"; then + echo "MCP must not listen on the pod network" >&2 + exit 1 + fi + kubectl -n default delete pod e2e-authn-probe --wait=false + # The driver creates the CoderTemplate itself (kubectl apply of config/e2e/codertemplate.yaml) after its setup checks. - name: Template and workspace lifecycle (apply, re-apply, watch, rename, delete, recreate) if: env.E2E_FULL == 'true' diff --git a/AGENTS.md b/AGENTS.md index a4518410..cf8813b1 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -73,7 +73,7 @@ Run from repository root. - **Workspace lifecycle E2E driver tests (offline, stubbed tools):** `bash ./hack/e2e-workspace-lifecycle_test.sh` (also run by `make test-scripts`) - **Lint (workflows):** `go run github.com/rhysd/actionlint/cmd/actionlint@v1.7.10` - **Development run (controller mode):** `GOFLAGS=-mod=vendor go run . --app=controller` (requires Kubernetes config via your env, e.g. `KUBECONFIG`) -- **Development run (aggregated API mode):** `GOFLAGS=-mod=vendor go run . --app=aggregated-apiserver` +- **Development run (aggregated API mode):** `GOFLAGS=-mod=vendor go run . --app=aggregated-apiserver` (needs `KUBECONFIG` or `~/.kube/config` with a cluster that can serve TokenReview/SubjectAccessReview and `kube-system/extension-apiserver-authentication`; the server does not start without it) - **Vendor consistency:** `make verify-vendor` - **Manifest generation:** `make manifests` (or `bash ./hack/update-manifests.sh`) - **Code generation:** `make codegen` (or `bash ./hack/update-codegen.sh`) diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index b98ad007..982e9a62 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -39,6 +39,10 @@ Uses controller-runtime with leader election. For each `CoderControlPlane`, it c Storage is backed by the Coder SDK, not memory or etcd: each request becomes a Coder API call. See [Aggregated API behavior](../reference/aggregated-api-behavior.md) for the consequences. +The Coder calls use the control plane's operator credentials, so the server checks every Kubernetes caller first. It uses delegated authentication (front-proxy client certificates from kube-apiserver, TokenReview for bearer tokens) and delegated authorization (SubjectAccessReview), and fails closed when those checks are unavailable. Only exact `/healthz`, `/livez`, and `/readyz` answer anonymous callers. See [How callers are checked](../how-to/deploy-aggregated-apiserver.md#how-callers-are-checked). + +Kubernetes users are not mapped to Coder users. Kubernetes RBAC on `aggregation.coder.com` in a namespace therefore grants owner-equivalent access in the Coder deployment of the control plane that serves that namespace. + How it finds its Coder backend: - **`all` mode:** `ControlPlaneClientProvider` discovers eligible `CoderControlPlane` resources and reads their operator token Secrets dynamically. diff --git a/docs/how-to/deploy-aggregated-apiserver.md b/docs/how-to/deploy-aggregated-apiserver.md index 03e4a80d..e37a8daf 100644 --- a/docs/how-to/deploy-aggregated-apiserver.md +++ b/docs/how-to/deploy-aggregated-apiserver.md @@ -12,6 +12,8 @@ kubectl apply -f config/rbac/ kubectl apply -f deploy/apiserver-service.yaml -f deploy/apiserver-apiservice.yaml ``` +`config/rbac/` includes two bindings the aggregated API server needs to check callers: `auth-delegator-binding.yaml` (create TokenReviews and SubjectAccessReviews) and `authentication-reader-binding.yaml` (read `kube-system/extension-apiserver-authentication`; the default `manager-role` also grants cluster-wide ConfigMap reads). Both name the `coder-k8s` ServiceAccount in `coder-system`; edit them if you install elsewhere. The server fails closed without these permissions: without read access to that ConfigMap it does not start, and without permission to create SubjectAccessReviews it answers every request with an error (members of `system:masters` excepted). + ## 2. Deploy ### Option A: all-in-one (recommended) @@ -67,6 +69,34 @@ kubectl -n coder-system patch deployment coder-k8s --type=strategic -p '{ }' ``` +## How callers are checked + +The aggregated API server authenticates and authorizes every request with the Kubernetes API: + +| Caller | Authentication | Authorization | +| --- | --- | --- | +| `kubectl` and other clients through kube-apiserver (the normal path) | kube-apiserver's front-proxy client certificate, verified against `requestheader-client-ca-file` from `kube-system/extension-apiserver-authentication`; the user comes from the `X-Remote-*` headers | SubjectAccessReview for that user, verb, resource, and namespace | +| Direct requests to port `6443` with a bearer token | TokenReview | SubjectAccessReview | +| Direct requests with a client certificate signed by the cluster client CA | Certificate subject | SubjectAccessReview | +| Anything else | None | Rejected with `401`, except exact `/healthz`, `/livez`, `/readyz` | + +`X-Remote-*` headers are trusted only on connections that present a valid front-proxy client certificate. Access to the resources is controlled with Kubernetes RBAC on `aggregation.coder.com` (`codertemplates`, `coderworkspaces`), but read the warning below before you grant it. + +!!! warning "RBAC on these resources is owner access in Coder" + Treat Kubernetes RBAC on `codertemplates` and `coderworkspaces` as owner-equivalent inside Coder. Grant it only to subjects you would trust as Coder owners, and prefer namespaced Roles over ClusterRoles. + +The server does not map Kubernetes users to Coder users. After Kubernetes RBAC allows a request, the server calls Coder with the control plane's operator token, which has owner rights in Coder. Each request uses only the control plane that serves the request's namespace (in standalone mode, the server is pinned to `--coder-namespace`). + +So a subject allowed to create, update, patch, or delete `codertemplates` or `coderworkspaces` in a namespace acts in that Coder deployment as an owner, and `get`, `list`, or `watch` shows that deployment's templates (including their source files) and workspaces as an owner sees them. A ClusterRole binding grants this for every namespace that has a control plane. + +Kept Kubernetes defaults, so you know what to expect: + +- Members of `system:masters` are authorized without a SubjectAccessReview, as in kube-apiserver. +- TokenReview results are cached for 10 seconds, allowed SubjectAccessReview results for 10 seconds, and denied results for 10 seconds. A permission change can take effect only after the matching cache entry expires. +- The server reads `extension-apiserver-authentication` at startup and keeps watching it. If the ConfigMap is deleted later, the server keeps trusting the CA it last loaded (retained trust) until it restarts or sees a new one. If the ConfigMap or its request-header CA is missing at startup, front-proxy requests fail with `401`. + +The server uses the same Kubernetes API for all of these checks: `KUBECONFIG` if set (exactly one file, never a fallback), otherwise the in-cluster ServiceAccount, otherwise `~/.kube/config`. With none, it does not start. + ## 3. Verify ```bash diff --git a/docs/how-to/troubleshooting.md b/docs/how-to/troubleshooting.md index 7cb9bf82..03f21089 100644 --- a/docs/how-to/troubleshooting.md +++ b/docs/how-to/troubleshooting.md @@ -8,7 +8,7 @@ kubectl logs -n coder-system deploy/coder-k8s ## A component I didn't expect is running -`--app` is optional and defaults to `all` (every component). To isolate one component, set it explicitly: +`--app` is optional and defaults to `all` (controller and aggregated API server; the MCP server never runs in `all`). To isolate one component, set it explicitly: ```bash GOFLAGS=-mod=vendor go run . --app=controller # or aggregated-apiserver, or mcp-http --mcp-token-file= @@ -64,6 +64,20 @@ kubectl get apiservice v1alpha1.aggregation.coder.com -o yaml Do not install CRDs for the same resources (`coderworkspaces.aggregation.coder.com`, `codertemplates.aggregation.coder.com`); they conflict with the aggregated API. +## The pod exits with `configure delegated authentication` + +The aggregated API server checks every caller with the Kubernetes API and refuses to start without it. + +- `... configmaps "extension-apiserver-authentication" is forbidden`: apply `config/rbac/authentication-reader-binding.yaml`. It must name the ServiceAccount the pod runs as (for example after installing into another namespace). +- `no Kubernetes configuration for delegated authentication and authorization` (outside a cluster): set `KUBECONFIG` to one kubeconfig file, or create `~/.kube/config`. +- `load kubeconfig ...` or `invalid kubeconfig ...`: the file named by `KUBECONFIG` is missing or incomplete. The server does not fall back to another configuration. + +## Aggregated requests fail with `401 Unauthorized` or `403 Forbidden` + +- **`401`:** the request has no valid credential. Requests sent straight to port `6443` need a Kubernetes bearer token; anonymous requests only reach `/healthz`, `/livez`, and `/readyz`. Use `kubectl`, which goes through kube-apiserver. +- **`403`:** the caller lacks RBAC for `aggregation.coder.com` in that namespace. Check with `kubectl auth can-i list codertemplates.aggregation.coder.com -n --as=`. Before you grant it, note that this RBAC is owner-equivalent inside Coder (see [How callers are checked](deploy-aggregated-apiserver.md#how-callers-are-checked)). +- **`500` mentioning `subjectaccessreviews`:** the server's ServiceAccount cannot create SubjectAccessReviews. Apply `config/rbac/auth-delegator-binding.yaml`. + ## Aggregated reads return `ServiceUnavailable` - **`all` mode:** no eligible `CoderControlPlane` exists yet, or its operator access is not ready. diff --git a/internal/app/apiserverapp/apiserverapp.go b/internal/app/apiserverapp/apiserverapp.go index 6ba30741..feb92b89 100644 --- a/internal/app/apiserverapp/apiserverapp.go +++ b/internal/app/apiserverapp/apiserverapp.go @@ -3,10 +3,12 @@ package apiserverapp import ( "context" + "errors" "fmt" "log" "net" "net/url" + "os" "strings" "time" @@ -18,8 +20,6 @@ import ( "k8s.io/apimachinery/pkg/runtime/schema" "k8s.io/apimachinery/pkg/runtime/serializer" utilruntime "k8s.io/apimachinery/pkg/util/runtime" - "k8s.io/apiserver/pkg/authentication/request/anonymous" - "k8s.io/apiserver/pkg/authorization/authorizerfactory" apiserveropenapi "k8s.io/apiserver/pkg/endpoints/openapi" "k8s.io/apiserver/pkg/registry/rest" genericapiserver "k8s.io/apiserver/pkg/server" @@ -60,6 +60,11 @@ type Options struct { // ClientProvider overrides the default static provider. // When set, CoderURL/CoderSessionToken/CoderNamespace flags are ignored. ClientProvider coder.ClientProvider + // Authentication and Authorization are test seams. When both are nil, the server uses + // delegated authentication and authorization against the Kubernetes API resolved by + // resolveDelegationKubeconfig. There is no anonymous or allow-all mode. + Authentication *genericoptions.DelegatingAuthenticationOptions + Authorization *genericoptions.DelegatingAuthorizationOptions } type errClientProvider struct { @@ -168,11 +173,14 @@ func NewScheme() *runtime.Scheme { return scheme } -// NewRecommendedConfig builds a recommended generic API server config. +// NewRecommendedConfig builds a recommended generic API server config with delegated +// authentication and authorization. func NewRecommendedConfig( scheme *runtime.Scheme, codecs serializer.CodecFactory, secureServingOptions *genericoptions.SecureServingOptions, + authenticationOptions *genericoptions.DelegatingAuthenticationOptions, + authorizationOptions *genericoptions.DelegatingAuthorizationOptions, ) (*genericapiserver.RecommendedConfig, error) { if scheme == nil { return nil, fmt.Errorf("assertion failed: scheme must not be nil") @@ -180,6 +188,12 @@ func NewRecommendedConfig( if secureServingOptions == nil { return nil, fmt.Errorf("assertion failed: secure serving options must not be nil") } + if authenticationOptions == nil || authorizationOptions == nil { + return nil, fmt.Errorf("assertion failed: delegated authentication and authorization options must not be nil") + } + if authenticationOptions.RemoteKubeConfigFile != authorizationOptions.RemoteKubeConfigFile { + return nil, fmt.Errorf("assertion failed: delegated authentication and authorization must use the same Kubernetes API authority") + } recommendedConfig := genericapiserver.NewRecommendedConfig(codecs) if recommendedConfig == nil { @@ -199,22 +213,35 @@ func NewRecommendedConfig( return nil, fmt.Errorf("assertion failed: loopback client config is nil after successful ApplyTo") } - authz := authorizerfactory.NewAlwaysAllowAuthorizer() - recommendedConfig.Authentication = genericapiserver.AuthenticationInfo{ - Authenticator: anonymous.NewAuthenticator(nil), + definitionNamer := apiserveropenapi.NewDefinitionNamer(scheme) + recommendedConfig.OpenAPIConfig = genericapiserver.DefaultOpenAPIConfig(getOpenAPIDefinitions, definitionNamer) + recommendedConfig.OpenAPIV3Config = genericapiserver.DefaultOpenAPIV3Config(getOpenAPIDefinitions, definitionNamer) + + if errs := authenticationOptions.Validate(); len(errs) > 0 { + return nil, fmt.Errorf("validate delegated authentication options: %w", errors.Join(errs...)) + } + if errs := authorizationOptions.Validate(); len(errs) > 0 { + return nil, fmt.Errorf("validate delegated authorization options: %w", errors.Join(errs...)) + } + if err := authenticationOptions.ApplyTo(&recommendedConfig.Authentication, recommendedConfig.SecureServing, recommendedConfig.OpenAPIConfig); err != nil { + return nil, fmt.Errorf("configure delegated authentication: %w", err) + } + if err := authorizationOptions.ApplyTo(&recommendedConfig.Authorization); err != nil { + return nil, fmt.Errorf("configure delegated authorization: %w", err) + } + if recommendedConfig.Authentication.Authenticator == nil { + return nil, fmt.Errorf("assertion failed: authenticator is nil after successful delegated authentication setup") } - recommendedConfig.Authorization = genericapiserver.AuthorizationInfo{ - Authorizer: authz, + if recommendedConfig.Authorization.Authorizer == nil { + return nil, fmt.Errorf("assertion failed: authorizer is nil after successful delegated authorization setup") } - recommendedConfig.RuleResolver = authz + // Applied before Complete(), which prepends the in-memory loopback token authenticator. + recommendedConfig.Authentication.Authenticator = anonymousHealthOnly{delegate: recommendedConfig.Authentication.Authenticator} + recommendedConfig.EffectiveVersion = apiservercompatibility.DefaultBuildEffectiveVersion() recommendedConfig.SkipOpenAPIInstallation = true recommendedConfig.RequestTimeout = defaultRequestTimeout - definitionNamer := apiserveropenapi.NewDefinitionNamer(scheme) - recommendedConfig.OpenAPIConfig = genericapiserver.DefaultOpenAPIConfig(getOpenAPIDefinitions, definitionNamer) - recommendedConfig.OpenAPIV3Config = genericapiserver.DefaultOpenAPIV3Config(getOpenAPIDefinitions, definitionNamer) - return recommendedConfig, nil } @@ -333,7 +360,20 @@ func RunWithOptions(ctx context.Context, opts Options) error { secureServingOptions.ServerCert.CertDirectory = "" secureServingOptions.ServerCert.PairName = "" - recommendedConfig, err := NewRecommendedConfig(scheme, codecs, secureServingOptions) + authenticationOptions, authorizationOptions := opts.Authentication, opts.Authorization + switch { + case authenticationOptions == nil && authorizationOptions == nil: + homeDir, _ := os.UserHomeDir() + kubeconfigPath, err := resolveDelegationKubeconfig(os.Getenv, homeDir) + if err != nil { + return fmt.Errorf("configure aggregated API server: %w", err) + } + authenticationOptions, authorizationOptions = newDelegatedAuthOptions(kubeconfigPath) + case authenticationOptions == nil || authorizationOptions == nil: + return fmt.Errorf("assertion failed: authentication and authorization options must be set together") + } + + recommendedConfig, err := NewRecommendedConfig(scheme, codecs, secureServingOptions, authenticationOptions, authorizationOptions) if err != nil { return fmt.Errorf("configure aggregated API server: %w", err) } diff --git a/internal/app/apiserverapp/apiserverapp_test.go b/internal/app/apiserverapp/apiserverapp_test.go index d0083ab3..02c13c96 100644 --- a/internal/app/apiserverapp/apiserverapp_test.go +++ b/internal/app/apiserverapp/apiserverapp_test.go @@ -215,7 +215,8 @@ func TestInstallAPIGroupRegistersDiscovery(t *testing.T) { secureServingOptions.ServerCert.CertDirectory = "" secureServingOptions.ServerCert.PairName = "" - recommendedConfig, err := NewRecommendedConfig(scheme, codecs, secureServingOptions) + authn, authz := newFakeKubeAPI(t).options(false) + recommendedConfig, err := NewRecommendedConfig(scheme, codecs, secureServingOptions, authn, authz) if err != nil { t.Fatalf("build recommended config: %v", err) } @@ -304,7 +305,8 @@ func TestNewRecommendedConfigSetsExtendedRequestTimeout(t *testing.T) { secureServingOptions.ServerCert.CertDirectory = "" secureServingOptions.ServerCert.PairName = "" - recommendedConfig, err := NewRecommendedConfig(scheme, codecs, secureServingOptions) + authn, authz := newFakeKubeAPI(t).options(false) + recommendedConfig, err := NewRecommendedConfig(scheme, codecs, secureServingOptions, authn, authz) if err != nil { t.Fatalf("build recommended config: %v", err) } @@ -417,6 +419,7 @@ func TestRunWithOptionsUsesClientProviderOverride(t *testing.T) { t.Fatalf("build static client provider: %v", err) } + authn, authz := newFakeKubeAPI(t).options(false) ctx, cancel := context.WithCancel(context.Background()) defer cancel() @@ -426,6 +429,8 @@ func TestRunWithOptionsUsesClientProviderOverride(t *testing.T) { Listener: listener, CoderURL: "https://coder.example.com", ClientProvider: provider, + Authentication: authn, + Authorization: authz, }) }() @@ -491,9 +496,10 @@ func TestRunWithOptionsStartsWithMissingCoderConfig(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() + authn, authz := newFakeKubeAPI(t).options(false) errCh := make(chan error, 1) go func() { - errCh <- RunWithOptions(ctx, Options{Listener: listener}) + errCh <- RunWithOptions(ctx, Options{Listener: listener, Authentication: authn, Authorization: authz}) }() select { diff --git a/internal/app/apiserverapp/auth.go b/internal/app/apiserverapp/auth.go new file mode 100644 index 00000000..76cb5af4 --- /dev/null +++ b/internal/app/apiserverapp/auth.go @@ -0,0 +1,138 @@ +package apiserverapp + +import ( + "errors" + "fmt" + "io/fs" + "net/http" + "os" + "path/filepath" + "strings" + + "k8s.io/apiserver/pkg/authentication/authenticator" + "k8s.io/apiserver/pkg/authentication/user" + genericoptions "k8s.io/apiserver/pkg/server/options" + "k8s.io/client-go/tools/clientcmd" +) + +// anonymousAllowedPaths are the only paths an unauthenticated caller may reach. They are exact +// matches: kubelet probes use them, and they return no Coder data. +var anonymousAllowedPaths = map[string]struct{}{ + "/healthz": {}, + "/livez": {}, + "/readyz": {}, +} + +// anonymousHealthOnly refuses the anonymous identity on every path except anonymousAllowedPaths. +// +// The vendored DelegatingAuthenticationOptions.ApplyTo always appends an anonymous fallback +// (k8s.io/apiserver/pkg/server/options/authentication.go), so without this wrapper "no data for +// anonymous callers" would depend on cluster RBAC never granting system:anonymous or +// system:unauthenticated anything on aggregation.coder.com. +type anonymousHealthOnly struct { + delegate authenticator.Request +} + +var _ authenticator.Request = anonymousHealthOnly{} + +func (a anonymousHealthOnly) AuthenticateRequest(req *http.Request) (*authenticator.Response, bool, error) { + if a.delegate == nil { + return nil, false, fmt.Errorf("assertion failed: delegate authenticator must not be nil") + } + resp, ok, err := a.delegate.AuthenticateRequest(req) + // Errors and "not authenticated" pass through unchanged; both end in 401. + if err != nil || !ok { + return resp, ok, err + } + if resp == nil || resp.User == nil { + return nil, false, fmt.Errorf("assertion failed: authenticated response must carry a user") + } + if resp.User.GetName() != user.Anonymous { + return resp, ok, nil + } + if _, allowed := anonymousAllowedPaths[req.URL.Path]; allowed { + return resp, ok, nil + } + return nil, false, nil +} + +// newDelegatedAuthOptions returns the production delegated authentication and authorization +// options. Both use the same Kubernetes API authority, resolved by resolveDelegationKubeconfig. +// Every other setting is the vendored default: lookup failures other than NotFound are fatal, +// system:masters bypasses SubjectAccessReview, and decisions are cached (see docs). +func newDelegatedAuthOptions(kubeconfigPath string) (*genericoptions.DelegatingAuthenticationOptions, *genericoptions.DelegatingAuthorizationOptions) { + authn := genericoptions.NewDelegatingAuthenticationOptions() + authn.RemoteKubeConfigFile = kubeconfigPath + authn.RemoteKubeConfigFileOptional = false + authn.SkipInClusterLookup = false + authn.TolerateInClusterLookupFailure = false + + authz := genericoptions.NewDelegatingAuthorizationOptions() + authz.RemoteKubeConfigFile = kubeconfigPath + authz.RemoteKubeConfigFileOptional = false + + return authn, authz +} + +// resolveDelegationKubeconfig picks the Kubernetes API authority used for TokenReview, +// SubjectAccessReview, and the extension-apiserver-authentication ConfigMap. It follows the same +// order as controller-runtime's config loading: +// +// 1. KUBECONFIG, if set, must name exactly one valid kubeconfig file. It never falls back. +// 2. In a cluster (KUBERNETES_SERVICE_HOST set): the in-cluster ServiceAccount (returns ""). +// 3. Otherwise $HOME/.kube/config, which must be valid. +// +// With none of these, it returns an error so the server does not start. +func resolveDelegationKubeconfig(getenv func(string) string, homeDir string) (string, error) { + if getenv == nil { + return "", fmt.Errorf("assertion failed: getenv must not be nil") + } + + if raw := getenv(clientcmd.RecommendedConfigPathEnvVar); raw != "" { + var paths []string + for _, path := range filepath.SplitList(raw) { + if strings.TrimSpace(path) != "" { + paths = append(paths, path) + } + } + if len(paths) != 1 { + return "", fmt.Errorf("KUBECONFIG must name exactly one file for delegated authentication and authorization, got %d", len(paths)) + } + if err := validateKubeconfigFile(paths[0]); err != nil { + return "", err + } + return paths[0], nil + } + + if getenv("KUBERNETES_SERVICE_HOST") != "" { + return "", nil + } + + if homeDir != "" { + path := filepath.Join(homeDir, clientcmd.RecommendedHomeDir, clientcmd.RecommendedFileName) + if _, err := os.Stat(path); err == nil { + if err := validateKubeconfigFile(path); err != nil { + return "", err + } + return path, nil + } else if !errors.Is(err, fs.ErrNotExist) { + return "", fmt.Errorf("stat %s: %w", path, err) + } + } + + return "", fmt.Errorf("no Kubernetes configuration for delegated authentication and authorization: run in a cluster or set KUBECONFIG") +} + +// validateKubeconfigFile loads the file with the non-deferred client config, which errors on +// empty or incomplete configs instead of silently falling back to in-cluster configuration the +// way the vendored options' deferred loader would. +func validateKubeconfigFile(path string) error { + cfg, err := clientcmd.LoadFromFile(path) + if err != nil { + return fmt.Errorf("load kubeconfig %s for delegated authentication: %w", path, err) + } + if _, err := clientcmd.NewDefaultClientConfig(*cfg, &clientcmd.ConfigOverrides{}).ClientConfig(); err != nil { + return fmt.Errorf("invalid kubeconfig %s for delegated authentication: %w", path, err) + } + return nil +} diff --git a/internal/app/apiserverapp/auth_envtest_test.go b/internal/app/apiserverapp/auth_envtest_test.go new file mode 100644 index 00000000..d2b162fb --- /dev/null +++ b/internal/app/apiserverapp/auth_envtest_test.go @@ -0,0 +1,310 @@ +package apiserverapp + +import ( + "context" + "crypto/tls" + "fmt" + "net/http" + "os" + "path/filepath" + "strings" + "sync" + "testing" + "time" + + authenticationv1 "k8s.io/api/authentication/v1" + corev1 "k8s.io/api/core/v1" + rbacv1 "k8s.io/api/rbac/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/client-go/kubernetes" + "sigs.k8s.io/controller-runtime/pkg/envtest" +) + +// A real kube-apiserver (envtest) configured with front-proxy request-header flags, so it +// publishes the real kube-system/extension-apiserver-authentication ConfigMap and serves real +// TokenReview/SubjectAccessReview with RBAC. The aggregated server under test uses the +// production delegated options; only the kubeconfig path is test-specific. +var ( + envtestOnce sync.Once + envtestErr error + envtestEnv *envtest.Environment + envtestDir string + envtestAdmin kubernetes.Interface + envtestFrontProxyCA *testCA +) + +func TestMain(m *testing.M) { + code := m.Run() + if envtestEnv != nil { + if err := envtestEnv.Stop(); err != nil { + fmt.Fprintf(os.Stderr, "stop envtest: %v\n", err) + } + } + if envtestDir != "" { + _ = os.RemoveAll(envtestDir) + } + os.Exit(code) +} + +func startSharedEnvtest(t *testing.T) { + t.Helper() + envtestOnce.Do(func() { + if os.Getenv("KUBEBUILDER_ASSETS") == "" { + envtestErr = fmt.Errorf("KUBEBUILDER_ASSETS is not set; run via `make test`") + return + } + envtestDir, envtestErr = os.MkdirTemp("", "apiserverapp-envtest-") + if envtestErr != nil { + return + } + envtestFrontProxyCA, envtestErr = generateTestCA(envtestDir, "front-proxy") + if envtestErr != nil { + return + } + env := &envtest.Environment{} + args := env.ControlPlane.GetAPIServer().Configure() + args.Set("requestheader-client-ca-file", envtestFrontProxyCA.pemPath) + args.Set("requestheader-allowed-names", frontProxyName) + args.Set("requestheader-username-headers", "X-Remote-User") + args.Set("requestheader-group-headers", "X-Remote-Group") + args.Set("requestheader-extra-headers-prefix", "X-Remote-Extra-") + cfg, err := env.Start() + if err != nil { + envtestErr = fmt.Errorf("start envtest: %w", err) + return + } + envtestEnv = env + envtestAdmin, envtestErr = kubernetes.NewForConfig(cfg) + if envtestErr != nil { + return + } + ctx := context.Background() + for _, ns := range []string{"test-ns", "other-ns"} { + if _, err := envtestAdmin.CoreV1().Namespaces().Create(ctx, &corev1.Namespace{ObjectMeta: metav1.ObjectMeta{Name: ns}}, metav1.CreateOptions{}); err != nil { + envtestErr = err + return + } + } + // kube-apiserver publishes the ConfigMap from a post-start hook; wait for the request-header keys. + deadline := time.Now().Add(30 * time.Second) + for { + cm, err := envtestAdmin.CoreV1().ConfigMaps("kube-system").Get(ctx, "extension-apiserver-authentication", metav1.GetOptions{}) + if err == nil && cm.Data["requestheader-client-ca-file"] != "" && cm.Data["client-ca-file"] != "" { + return + } + if time.Now().After(deadline) { + envtestErr = fmt.Errorf("extension-apiserver-authentication not published: %w", err) + return + } + time.Sleep(200 * time.Millisecond) + } + }) + if envtestErr != nil { + t.Fatalf("envtest unavailable: %v", envtestErr) + } +} + +// envtestUser creates a certificate-authenticated envtest user, binds the given roles, and +// returns a kubeconfig path plus the user's client certificate. +func envtestUser(t *testing.T, name string, authDelegator, authReader bool) (string, tls.Certificate) { + t.Helper() + ctx := t.Context() + u, err := envtestEnv.AddUser(envtest.User{Name: name}, nil) + if err != nil { + t.Fatal(err) + } + kubeconfig, err := u.KubeConfig() + if err != nil { + t.Fatal(err) + } + path := filepath.Join(t.TempDir(), name+".kubeconfig") + if err := os.WriteFile(path, kubeconfig, 0o600); err != nil { + t.Fatal(err) + } + subject := []rbacv1.Subject{{Kind: rbacv1.UserKind, APIGroup: rbacv1.GroupName, Name: name}} + if authDelegator { + mustCreate(t, func() error { + _, err := envtestAdmin.RbacV1().ClusterRoleBindings().Create(ctx, &rbacv1.ClusterRoleBinding{ + ObjectMeta: metav1.ObjectMeta{Name: name + "-auth-delegator"}, + RoleRef: rbacv1.RoleRef{APIGroup: rbacv1.GroupName, Kind: "ClusterRole", Name: "system:auth-delegator"}, + Subjects: subject, + }, metav1.CreateOptions{}) + return err + }) + } + if authReader { + mustCreate(t, func() error { + _, err := envtestAdmin.RbacV1().RoleBindings("kube-system").Create(ctx, &rbacv1.RoleBinding{ + ObjectMeta: metav1.ObjectMeta{Name: name + "-authentication-reader"}, + RoleRef: rbacv1.RoleRef{APIGroup: rbacv1.GroupName, Kind: "Role", Name: "extension-apiserver-authentication-reader"}, + Subjects: subject, + }, metav1.CreateOptions{}) + return err + }) + } + cert, err := tls.X509KeyPair(u.Config().CertData, u.Config().KeyData) + if err != nil { + t.Fatal(err) + } + return path, cert +} + +func mustCreate(t *testing.T, create func() error) { + t.Helper() + if err := create(); err != nil && !apierrors.IsAlreadyExists(err) { + t.Fatal(err) + } +} + +func grantTemplateReader(t *testing.T, namespace string, subject rbacv1.Subject, name string) { + t.Helper() + ctx := t.Context() + mustCreate(t, func() error { + _, err := envtestAdmin.RbacV1().Roles(namespace).Create(ctx, &rbacv1.Role{ + ObjectMeta: metav1.ObjectMeta{Name: name}, + Rules: []rbacv1.PolicyRule{{APIGroups: []string{aggGroup}, Resources: []string{"codertemplates"}, Verbs: []string{"get", "list"}}}, + }, metav1.CreateOptions{}) + return err + }) + mustCreate(t, func() error { + _, err := envtestAdmin.RbacV1().RoleBindings(namespace).Create(ctx, &rbacv1.RoleBinding{ + ObjectMeta: metav1.ObjectMeta{Name: name}, + RoleRef: rbacv1.RoleRef{APIGroup: rbacv1.GroupName, Kind: "Role", Name: name}, + Subjects: []rbacv1.Subject{subject}, + }, metav1.CreateOptions{}) + return err + }) +} + +func startEnvtestAuthServer(t *testing.T, kubeconfigPath string) authTestServer { + t.Helper() + // Production options: in-cluster ConfigMap lookup, fatal lookup errors, default caches. + authn, authz := newDelegatedAuthOptions(kubeconfigPath) + return startAuthTestServer(t, authn, authz) +} + +func TestEnvtestDelegatedAuthWithRealRBAC(t *testing.T) { + startSharedEnvtest(t) + kubeconfigPath, _ := envtestUser(t, "coder-k8s-delegator", true, true) + server := startEnvtestAuthServer(t, kubeconfigPath) + grantTemplateReader(t, "test-ns", rbacv1.Subject{Kind: rbacv1.UserKind, APIGroup: rbacv1.GroupName, Name: "alice"}, "alice-template-reader") + + frontProxy := envtestFrontProxyCA.clientCert(t, frontProxyName) + alice := remoteUser("alice", "team-a") + + // Front-proxy relayed identity, authorized by real RBAC. The request-header CA is loaded + // asynchronously from the ConfigMap after startup, so poll briefly. + var status int + var body string + for deadline := time.Now().Add(15 * time.Second); ; { + status, body = server.do(t, &frontProxy, http.MethodGet, templatesTestNS, alice, "") + if status == http.StatusOK || time.Now().After(deadline) { + break + } + time.Sleep(200 * time.Millisecond) + } + if status != http.StatusOK || !strings.Contains(body, testTemplateName) { + t.Fatalf("alice list test-ns: status=%d body=%.300s", status, body) + } + allowedCalls := server.provider.calls.Load() + for _, tc := range []struct { + method, path string + cert *tls.Certificate + headers map[string]string + want int + }{ + {method: http.MethodGet, path: templatesOtherNS, cert: &frontProxy, headers: alice, want: http.StatusForbidden}, + {method: http.MethodGet, path: workspacesTestNS, cert: &frontProxy, headers: alice, want: http.StatusForbidden}, + {method: http.MethodDelete, path: templatesTestNS + "/" + testTemplateName, cert: &frontProxy, headers: alice, want: http.StatusForbidden}, + {method: http.MethodGet, path: templatesTestNS, cert: &frontProxy, headers: remoteUser("mallory"), want: http.StatusForbidden}, + {method: http.MethodGet, path: templatesTestNS, want: http.StatusUnauthorized}, + {method: http.MethodGet, path: templatesTestNS, headers: alice, want: http.StatusUnauthorized}, + {method: http.MethodGet, path: templatesTestNS, headers: remoteUser("kubernetes-admin", "system:masters"), want: http.StatusUnauthorized}, + } { + if status, body := server.do(t, tc.cert, tc.method, tc.path, tc.headers, ""); status != tc.want { + t.Errorf("%s %s headers=%v: status=%d, want %d; body=%.200s", tc.method, tc.path, tc.headers, status, tc.want, body) + } + } + + // An ordinary cluster client certificate: identity is the certificate subject, headers ignored. + _, carolCert := envtestUser(t, "carol", false, false) + if status, _ := server.do(t, &carolCert, http.MethodGet, templatesTestNS, alice, ""); status != http.StatusForbidden { + t.Errorf("carol cert + forged alice headers: status=%d, want 403", status) + } + + // ServiceAccount bearer token via real TokenReview; RBAC granted after the first denial. + ctx := t.Context() + mustCreate(t, func() error { + _, err := envtestAdmin.CoreV1().ServiceAccounts("test-ns").Create(ctx, &corev1.ServiceAccount{ObjectMeta: metav1.ObjectMeta{Name: "reader"}}, metav1.CreateOptions{}) + return err + }) + tokenRequest, err := envtestAdmin.CoreV1().ServiceAccounts("test-ns").CreateToken(ctx, "reader", &authenticationv1.TokenRequest{ + Spec: authenticationv1.TokenRequestSpec{ExpirationSeconds: ptrInt64(600)}, + }, metav1.CreateOptions{}) + if err != nil { + t.Fatal(err) + } + saToken := tokenRequest.Status.Token + if status, _ := server.do(t, nil, http.MethodGet, templatesTestNS, bearer(saToken), ""); status != http.StatusForbidden { + t.Fatalf("SA token without RBAC: status=%d, want 403", status) + } + if status, _ := server.do(t, nil, http.MethodGet, templatesTestNS, mergeHeaders(bearer(saToken), remoteUser("kubernetes-admin", "system:masters")), ""); status != http.StatusForbidden { + t.Fatalf("SA token + forged masters headers: status=%d, want 403", status) + } + if calls := server.provider.calls.Load(); calls != allowedCalls { + t.Fatalf("denied requests reached the Coder backend: calls %d -> %d", allowedCalls, calls) + } + grantTemplateReader(t, "test-ns", rbacv1.Subject{Kind: rbacv1.ServiceAccountKind, Name: "reader", Namespace: "test-ns"}, "sa-template-reader") + // The deny decision is cached for up to DenyCacheTTL (10s); poll past it. + deadline := time.Now().Add(30 * time.Second) + for { + status, _ := server.do(t, nil, http.MethodGet, templatesTestNS, bearer(saToken), "") + if status == http.StatusOK { + break + } + if time.Now().After(deadline) { + t.Fatalf("SA token after RoleBinding: status=%d, want 200", status) + } + time.Sleep(time.Second) + } +} + +func TestEnvtestDelegatedAuthFailsClosedWithoutBindings(t *testing.T) { + startSharedEnvtest(t) + + // Without the authentication-reader binding the ConfigMap lookup is forbidden: no startup. + noReader, _ := envtestUser(t, "delegator-no-reader", true, false) + authn, authz := newDelegatedAuthOptions(noReader) + secureServingOptions, _ := newTestSecureServing(t) + if _, err := NewRecommendedConfig(NewScheme(), codecsFor(), secureServingOptions, authn, authz); err == nil || !strings.Contains(strings.ToLower(err.Error()), "forbidden") { + t.Fatalf("expected forbidden ConfigMap startup failure, got %v", err) + } + + // Without auth-delegator the server starts, but SubjectAccessReview is forbidden: no access. + noDelegator, _ := envtestUser(t, "delegator-no-auth-delegator", false, true) + server := startEnvtestAuthServer(t, noDelegator) + grantTemplateReader(t, "test-ns", rbacv1.Subject{Kind: rbacv1.UserKind, APIGroup: rbacv1.GroupName, Name: "alice"}, "alice-template-reader") + frontProxy := envtestFrontProxyCA.clientCert(t, frontProxyName) + // Wait until the front-proxy CA is loaded (401 turns into an authorization outcome). + var status int + var body string + for deadline := time.Now().Add(15 * time.Second); ; { + status, body = server.do(t, &frontProxy, http.MethodGet, templatesTestNS, remoteUser("alice"), "") + if status != http.StatusUnauthorized || time.Now().After(deadline) { + break + } + time.Sleep(200 * time.Millisecond) + } + if status == http.StatusUnauthorized { + t.Fatal("front-proxy CA never loaded; cannot observe the authorization outcome") + } + if status < 400 || strings.Contains(body, testTemplateName) { + t.Fatalf("SAR forbidden for the delegator: status=%d body=%.200s", status, body) + } + if calls := server.provider.calls.Load(); calls != 0 { + t.Fatalf("requests reached the Coder backend %d times without authorization", calls) + } +} + +func ptrInt64(v int64) *int64 { return &v } diff --git a/internal/app/apiserverapp/auth_helpers_test.go b/internal/app/apiserverapp/auth_helpers_test.go new file mode 100644 index 00000000..9c8f0cdb --- /dev/null +++ b/internal/app/apiserverapp/auth_helpers_test.go @@ -0,0 +1,542 @@ +package apiserverapp + +import ( + "context" + "crypto/ecdsa" + "crypto/elliptic" + "crypto/rand" + "crypto/tls" + "crypto/x509" + "crypto/x509/pkix" + "encoding/json" + "encoding/pem" + "errors" + "fmt" + "io" + "math/big" + "net" + "net/http" + "net/http/httptest" + "net/url" + "os" + "path/filepath" + "strings" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/coder/coder/v2/codersdk" + authenticationv1 "k8s.io/api/authentication/v1" + authorizationv1 "k8s.io/api/authorization/v1" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime/serializer" + "k8s.io/apimachinery/pkg/util/wait" + genericoptions "k8s.io/apiserver/pkg/server/options" + "k8s.io/client-go/tools/clientcmd" + clientcmdapi "k8s.io/client-go/tools/clientcmd/api" + + "github.com/coder/coder-k8s/internal/aggregated/coder" +) + +// ---- PKI ---- + +type testCA struct { + cert *x509.Certificate + key *ecdsa.PrivateKey + pemPath string +} + +func newTestCA(t *testing.T, name string) *testCA { + t.Helper() + ca, err := generateTestCA(t.TempDir(), name) + if err != nil { + t.Fatal(err) + } + return ca +} + +func generateTestCA(dir, name string) (*testCA, error) { + key, err := ecdsa.GenerateKey(elliptic.P256(), rand.Reader) + if err != nil { + return nil, err + } + template := &x509.Certificate{ + SerialNumber: big.NewInt(time.Now().UnixNano()), + Subject: pkix.Name{CommonName: name}, + NotBefore: time.Now().Add(-time.Hour), + NotAfter: time.Now().Add(24 * time.Hour), + KeyUsage: x509.KeyUsageCertSign | x509.KeyUsageDigitalSignature, + BasicConstraintsValid: true, + IsCA: true, + } + der, err := x509.CreateCertificate(rand.Reader, template, template, &key.PublicKey, key) + if err != nil { + return nil, err + } + cert, err := x509.ParseCertificate(der) + if err != nil { + return nil, err + } + path := filepath.Join(dir, name+"-ca.crt") + if err := os.WriteFile(path, pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: der}), 0o600); err != nil { + return nil, err + } + return &testCA{cert: cert, key: key, pemPath: path}, nil +} + +func (ca *testCA) pem() []byte { + return pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: ca.cert.Raw}) +} + +// clientCert issues a client certificate. Kubernetes maps CN to the user name and O to groups. +func (ca *testCA) clientCert(t *testing.T, commonName string, groups ...string) tls.Certificate { + t.Helper() + key, err := ecdsa.GenerateKey(elliptic.P256(), rand.Reader) + if err != nil { + t.Fatal(err) + } + template := &x509.Certificate{ + SerialNumber: big.NewInt(time.Now().UnixNano()), + Subject: pkix.Name{CommonName: commonName, Organization: groups}, + NotBefore: time.Now().Add(-time.Hour), + NotAfter: time.Now().Add(24 * time.Hour), + KeyUsage: x509.KeyUsageDigitalSignature, + ExtKeyUsage: []x509.ExtKeyUsage{x509.ExtKeyUsageClientAuth}, + } + der, err := x509.CreateCertificate(rand.Reader, template, ca.cert, &key.PublicKey, ca.key) + if err != nil { + t.Fatal(err) + } + return tls.Certificate{Certificate: [][]byte{der}, PrivateKey: key} +} + +// ---- Fake Kubernetes API for TokenReview, SubjectAccessReview and the auth ConfigMap ---- + +type fakeKubeAPI struct { + server *httptest.Server + kubeconfigPath string + + mu sync.Mutex + tokens map[string]authenticationv1.UserInfo + decide func(authorizationv1.SubjectAccessReviewSpec) bool + sars []authorizationv1.SubjectAccessReviewSpec + tokenReviews int + failTokenReview bool + failSAR bool + configMap *corev1.ConfigMap + configMapForbidden bool +} + +func newFakeKubeAPI(t *testing.T) *fakeKubeAPI { + t.Helper() + f := &fakeKubeAPI{ + tokens: map[string]authenticationv1.UserInfo{}, + decide: func(authorizationv1.SubjectAccessReviewSpec) bool { return false }, + } + f.server = httptest.NewTLSServer(http.HandlerFunc(f.serveHTTP)) + t.Cleanup(f.server.Close) + + cfg := clientcmdapi.NewConfig() + cfg.Clusters["fake"] = &clientcmdapi.Cluster{ + Server: f.server.URL, + CertificateAuthorityData: pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: f.server.Certificate().Raw}), + } + cfg.AuthInfos["delegator"] = &clientcmdapi.AuthInfo{Token: "delegator-credential"} + cfg.Contexts["fake"] = &clientcmdapi.Context{Cluster: "fake", AuthInfo: "delegator"} + cfg.CurrentContext = "fake" + f.kubeconfigPath = filepath.Join(t.TempDir(), "kubeconfig") + if err := clientcmd.WriteToFile(*cfg, f.kubeconfigPath); err != nil { + t.Fatal(err) + } + return f +} + +func (f *fakeKubeAPI) setToken(token string, info authenticationv1.UserInfo) { + f.mu.Lock() + defer f.mu.Unlock() + f.tokens[token] = info +} + +func (f *fakeKubeAPI) setDecide(decide func(authorizationv1.SubjectAccessReviewSpec) bool) { + f.mu.Lock() + defer f.mu.Unlock() + f.decide = decide +} + +func (f *fakeKubeAPI) setFailures(tokenReview, sar bool) { + f.mu.Lock() + defer f.mu.Unlock() + f.failTokenReview, f.failSAR = tokenReview, sar +} + +func (f *fakeKubeAPI) recordedSARs() []authorizationv1.SubjectAccessReviewSpec { + f.mu.Lock() + defer f.mu.Unlock() + return append([]authorizationv1.SubjectAccessReviewSpec(nil), f.sars...) +} + +func (f *fakeKubeAPI) resetSARs() { + f.mu.Lock() + defer f.mu.Unlock() + f.sars = nil +} + +func (f *fakeKubeAPI) serveHTTP(w http.ResponseWriter, r *http.Request) { + const configMapCollection = "/api/v1/namespaces/kube-system/configmaps" + switch { + case r.Method == http.MethodPost && r.URL.Path == "/apis/authentication.k8s.io/v1/tokenreviews": + var review authenticationv1.TokenReview + if err := json.NewDecoder(r.Body).Decode(&review); err != nil { + writeStatus(w, http.StatusBadRequest, "BadRequest") + return + } + f.mu.Lock() + f.tokenReviews++ + fail := f.failTokenReview + info, ok := f.tokens[review.Spec.Token] + f.mu.Unlock() + if fail { + writeStatus(w, http.StatusInternalServerError, "InternalError") + return + } + review.Status = authenticationv1.TokenReviewStatus{Authenticated: ok} + if ok { + review.Status.User = info + } + review.APIVersion, review.Kind = "authentication.k8s.io/v1", "TokenReview" + writeKubeJSON(w, http.StatusCreated, review) + case r.Method == http.MethodPost && r.URL.Path == "/apis/authorization.k8s.io/v1/subjectaccessreviews": + var review authorizationv1.SubjectAccessReview + if err := json.NewDecoder(r.Body).Decode(&review); err != nil { + writeStatus(w, http.StatusBadRequest, "BadRequest") + return + } + f.mu.Lock() + f.sars = append(f.sars, review.Spec) + fail := f.failSAR + decide := f.decide + f.mu.Unlock() + if fail { + writeStatus(w, http.StatusInternalServerError, "InternalError") + return + } + allowed := decide(review.Spec) + review.Status = authorizationv1.SubjectAccessReviewStatus{Allowed: allowed, Denied: !allowed} + review.APIVersion, review.Kind = "authorization.k8s.io/v1", "SubjectAccessReview" + writeKubeJSON(w, http.StatusCreated, review) + case r.Method == http.MethodGet && r.URL.Path == configMapCollection+"/extension-apiserver-authentication": + f.mu.Lock() + cm, forbidden := f.configMap, f.configMapForbidden + f.mu.Unlock() + switch { + case forbidden: + writeStatus(w, http.StatusForbidden, "Forbidden") + case cm == nil: + writeStatus(w, http.StatusNotFound, "NotFound") + default: + writeKubeJSON(w, http.StatusOK, cm) + } + case r.Method == http.MethodGet && r.URL.Path == configMapCollection: + if r.URL.Query().Get("watch") == "true" { + // client-go reflectors use streaming watch-list: send the current object (if any) and the + // initial-events-end bookmark, then keep the watch open until the client goes away. + f.mu.Lock() + cm := f.configMap + f.mu.Unlock() + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusOK) + encoder := json.NewEncoder(w) + if cm != nil { + _ = encoder.Encode(map[string]any{"type": "ADDED", "object": cm}) + } + if r.URL.Query().Get("sendInitialEvents") == "true" { + _ = encoder.Encode(map[string]any{"type": "BOOKMARK", "object": map[string]any{ + "apiVersion": "v1", "kind": "ConfigMap", + "metadata": map[string]any{"resourceVersion": "1", "annotations": map[string]string{"k8s.io/initial-events-end": "true"}}, + }}) + } + if flusher, ok := w.(http.Flusher); ok { + flusher.Flush() + } + <-r.Context().Done() + return + } + f.mu.Lock() + cm, forbidden := f.configMap, f.configMapForbidden + f.mu.Unlock() + if forbidden { + writeStatus(w, http.StatusForbidden, "Forbidden") + return + } + list := corev1.ConfigMapList{TypeMeta: metav1.TypeMeta{APIVersion: "v1", Kind: "ConfigMapList"}, ListMeta: metav1.ListMeta{ResourceVersion: "1"}} + if cm != nil { + list.Items = append(list.Items, *cm) + } + writeKubeJSON(w, http.StatusOK, list) + default: + writeStatus(w, http.StatusNotFound, "NotFound") + } +} + +func writeKubeJSON(w http.ResponseWriter, status int, body any) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(status) + _ = json.NewEncoder(w).Encode(body) +} + +func writeStatus(w http.ResponseWriter, code int, reason metav1.StatusReason) { + writeKubeJSON(w, code, metav1.Status{ + TypeMeta: metav1.TypeMeta{APIVersion: "v1", Kind: "Status"}, + Status: metav1.StatusFailure, + Code: int32(code), //nolint:gosec // G115: HTTP status codes fit in int32. + Reason: reason, + }) +} + +// fastBackoff keeps outage tests fast; production uses the vendored default backoff. +var fastBackoff = wait.Backoff{Duration: time.Millisecond, Factor: 1, Steps: 1} + +// options returns delegated options pointed at the fake API. The ConfigMap lookup is skipped +// unless useConfigMap is true; CA files can then be supplied explicitly. +func (f *fakeKubeAPI) options(useConfigMap bool) (*genericoptions.DelegatingAuthenticationOptions, *genericoptions.DelegatingAuthorizationOptions) { + authn, authz := newDelegatedAuthOptions(f.kubeconfigPath) + authn.SkipInClusterLookup = !useConfigMap + authn.WithCustomRetryBackoff(fastBackoff) + authz.WithCustomRetryBackoff(fastBackoff) + return authn, authz +} + +// ---- Backend spy ---- + +// countingProvider counts every backend lookup so denial tests can assert that no Coder call +// (read or write) happened. +type countingProvider struct { + inner coder.ClientProvider + calls atomic.Int32 +} + +func (p *countingProvider) ClientForNamespace(ctx context.Context, namespace string) (*codersdk.Client, error) { + p.calls.Add(1) + return p.inner.ClientForNamespace(ctx, namespace) +} + +// ---- Server harness ---- + +type authTestServer struct { + baseURL string + provider *countingProvider + mock *integrationMockCoderServer +} + +// startAuthTestServer boots the production server configuration (NewRecommendedConfig, +// NewGenericAPIServer, InstallAPIGroup) with the given delegated options against a mock Coder +// backend serving namespace test-ns. +func startAuthTestServer( + t *testing.T, + authn *genericoptions.DelegatingAuthenticationOptions, + authz *genericoptions.DelegatingAuthorizationOptions, +) authTestServer { + t.Helper() + + mockCoder := newIntegrationMockCoderServer("test-token") + t.Cleanup(mockCoder.Close) + mockCoderURL, err := url.Parse(mockCoder.URL()) + if err != nil { + t.Fatal(err) + } + sdkClient := codersdk.New(mockCoderURL) + sdkClient.SetSessionToken("test-token") + provider := &countingProvider{inner: &coder.StaticClientProvider{Client: sdkClient, Namespace: "test-ns"}} + + scheme := NewScheme() + codecs := serializer.NewCodecFactory(scheme) + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = listener.Close() }) + + secureServingOptions := genericoptions.NewSecureServingOptions() + secureServingOptions.Listener = listener + secureServingOptions.BindPort = 0 + secureServingOptions.ServerCert.CertDirectory = "" + secureServingOptions.ServerCert.PairName = "" + + recommendedConfig, err := NewRecommendedConfig(scheme, codecs, secureServingOptions, authn, authz) + if err != nil { + t.Fatalf("build recommended config: %v", err) + } + server, err := NewGenericAPIServer(recommendedConfig) + if err != nil { + t.Fatalf("build generic API server: %v", err) + } + t.Cleanup(server.Destroy) + apiGroupInfo, err := NewAPIGroupInfo(scheme, codecs, provider) + if err != nil { + t.Fatal(err) + } + if err := InstallAPIGroup(server, apiGroupInfo); err != nil { + t.Fatal(err) + } + + ctx, cancel := context.WithCancel(context.Background()) + errCh := make(chan error, 1) + go func() { errCh <- server.PrepareRun().RunWithContext(ctx) }() + t.Cleanup(func() { + cancel() + select { + case runErr := <-errCh: + if runErr != nil && !errors.Is(runErr, context.Canceled) { + t.Errorf("aggregated API server exited with error: %v", runErr) + } + case <-time.After(10 * time.Second): + t.Error("timed out waiting for aggregated API server to stop") + } + }) + + s := authTestServer{ + baseURL: strings.TrimSuffix(recommendedConfig.LoopbackClientConfig.Host, "/"), + provider: provider, + mock: mockCoder, + } + // Anonymous /readyz is one of the allowed health paths, so it doubles as a startup probe. + deadline := time.Now().Add(15 * time.Second) + for { + status, _ := s.do(t, nil, http.MethodGet, "/readyz", nil, "") + if status == http.StatusOK { + break + } + select { + case runErr := <-errCh: + t.Fatalf("server exited during startup: %v", runErr) + default: + } + if time.Now().After(deadline) { + t.Fatalf("server not ready, last /readyz status %d", status) + } + time.Sleep(50 * time.Millisecond) + } + return s +} + +// do sends one request. cert is an optional TLS client certificate. +func (s authTestServer) do(t *testing.T, cert *tls.Certificate, method, path string, headers map[string]string, body string) (int, string) { + t.Helper() + tlsConfig := &tls.Config{InsecureSkipVerify: true} //nolint:gosec // Test server uses an ephemeral self-signed cert. + if cert != nil { + tlsConfig.Certificates = []tls.Certificate{*cert} + } + client := &http.Client{Timeout: 15 * time.Second, Transport: &http.Transport{TLSClientConfig: tlsConfig}} + defer client.CloseIdleConnections() + + var reader io.Reader + if body != "" { + reader = strings.NewReader(body) + } + req, err := http.NewRequestWithContext(t.Context(), method, s.baseURL+path, reader) + if err != nil { + t.Fatal(err) + } + for k, v := range headers { + req.Header.Set(k, v) + } + if body != "" && req.Header.Get("Content-Type") == "" { + req.Header.Set("Content-Type", "application/json") + } + resp, err := client.Do(req) + if err != nil { + t.Fatalf("%s %s: %v", method, path, err) + } + defer func() { _ = resp.Body.Close() }() + data, err := io.ReadAll(resp.Body) + if err != nil { + t.Fatal(err) + } + return resp.StatusCode, string(data) +} + +func bearer(token string) map[string]string { + return map[string]string{"Authorization": "Bearer " + token} +} + +func mergeHeaders(maps ...map[string]string) map[string]string { + out := map[string]string{} + for _, m := range maps { + for k, v := range m { + out[k] = v + } + } + return out +} + +func mustContain(t *testing.T, list []string, want string) { + t.Helper() + for _, v := range list { + if v == want { + return + } + } + t.Fatalf("expected %q in %v", want, list) +} + +func mustNotContain(t *testing.T, list []string, unwanted string) { + t.Helper() + for _, v := range list { + if v == unwanted { + t.Fatalf("did not expect %q in %v", unwanted, list) + } + } +} + +func describeSAR(spec authorizationv1.SubjectAccessReviewSpec) string { + if spec.ResourceAttributes != nil { + ra := spec.ResourceAttributes + return fmt.Sprintf("user=%s verb=%s group=%s resource=%s ns=%s name=%s", spec.User, ra.Verb, ra.Group, ra.Resource, ra.Namespace, ra.Name) + } + if spec.NonResourceAttributes != nil { + return fmt.Sprintf("user=%s verb=%s path=%s", spec.User, spec.NonResourceAttributes.Verb, spec.NonResourceAttributes.Path) + } + return "user=" + spec.User +} + +func newTestSecureServing(t *testing.T) (*genericoptions.SecureServingOptions, net.Listener) { + t.Helper() + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = listener.Close() }) + opts := genericoptions.NewSecureServingOptions() + opts.Listener = listener + opts.BindPort = 0 + opts.ServerCert.CertDirectory = "" + opts.ServerCert.PairName = "" + return opts, listener +} + +func codecsFor() serializer.CodecFactory { + return serializer.NewCodecFactory(NewScheme()) +} + +func mustWrite(t *testing.T, path, content string) { + t.Helper() + if err := os.MkdirAll(filepath.Dir(path), 0o700); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(path, []byte(content), 0o600); err != nil { + t.Fatal(err) + } +} + +func httptestRequest(t *testing.T) *http.Request { + t.Helper() + return httptest.NewRequest(http.MethodGet, "https://example.test/healthz", nil) +} + +func newFakeKubeAPIOptionsOnly(t *testing.T) *genericoptions.DelegatingAuthenticationOptions { + t.Helper() + authn, _ := newFakeKubeAPI(t).options(false) + return authn +} diff --git a/internal/app/apiserverapp/auth_test.go b/internal/app/apiserverapp/auth_test.go new file mode 100644 index 00000000..384eee7e --- /dev/null +++ b/internal/app/apiserverapp/auth_test.go @@ -0,0 +1,591 @@ +package apiserverapp + +import ( + "context" + "crypto/tls" + "encoding/json" + "errors" + "net/http" + "os" + "path/filepath" + "strings" + "testing" + "time" + + authenticationv1 "k8s.io/api/authentication/v1" + authorizationv1 "k8s.io/api/authorization/v1" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apiserver/pkg/authentication/authenticator" + "k8s.io/apiserver/pkg/authentication/user" +) + +const ( + aggGroup = "aggregation.coder.com" + templatesTestNS = "/apis/aggregation.coder.com/v1alpha1/namespaces/test-ns/codertemplates" + workspacesTestNS = "/apis/aggregation.coder.com/v1alpha1/namespaces/test-ns/coderworkspaces" + templatesOtherNS = "/apis/aggregation.coder.com/v1alpha1/namespaces/other-ns/codertemplates" + templatesAllNS = "/apis/aggregation.coder.com/v1alpha1/codertemplates" + frontProxyName = "front-proxy-client" + testTemplateName = "default.my-template" + testWorkspaceName = "default.testuser.my-workspace" +) + +// authFixture wires the production server to a fake Kubernetes API plus explicit CAs: +// frontProxyCA signs front-proxy (request-header) client certs, clusterCA signs ordinary +// Kubernetes client certs. +type authFixture struct { + kube *fakeKubeAPI + frontProxyCA *testCA + clusterCA *testCA + server authTestServer +} + +func newAuthFixture(t *testing.T, tune func(*fakeKubeAPI)) *authFixture { + t.Helper() + f := &authFixture{ + kube: newFakeKubeAPI(t), + frontProxyCA: newTestCA(t, "front-proxy"), + clusterCA: newTestCA(t, "cluster"), + } + if tune != nil { + tune(f.kube) + } + authn, authz := f.kube.options(false) + authn.RequestHeader.ClientCAFile = f.frontProxyCA.pemPath + authn.RequestHeader.AllowedNames = []string{frontProxyName} + authn.ClientCert.ClientCA = f.clusterCA.pemPath + f.server = startAuthTestServer(t, authn, authz) + return f +} + +func (f *authFixture) frontProxyCert(t *testing.T) *tls.Certificate { + cert := f.frontProxyCA.clientCert(t, frontProxyName) + return &cert +} + +func remoteUser(name string, groups ...string) map[string]string { + h := map[string]string{"X-Remote-User": name} + if len(groups) > 0 { + h["X-Remote-Group"] = strings.Join(groups, ",") + } + return h +} + +func allowAll(authorizationv1.SubjectAccessReviewSpec) bool { return true } + +// TestDelegatedAuthRejectsUnauthenticatedRequests: SubjectAccessReview allows everything here, so +// every 401 below comes from authentication, not from RBAC. +func TestDelegatedAuthRejectsUnauthenticatedRequests(t *testing.T) { + f := newAuthFixture(t, func(k *fakeKubeAPI) { k.setDecide(allowAll) }) + rogueCA := newTestCA(t, "rogue") + rogueFrontProxy := rogueCA.clientCert(t, frontProxyName) + wrongName := f.frontProxyCA.clientCert(t, "not-the-front-proxy") + forged := remoteUser("kubernetes-admin", "system:masters") + + credentials := []struct { + name string + cert *tls.Certificate + headers map[string]string + }{ + {name: "no credentials"}, + {name: "forged remote user and masters group", headers: forged}, + {name: "forged masters group only", headers: map[string]string{"X-Remote-Group": "system:masters"}}, + {name: "unknown CA front-proxy cert with forged headers", cert: &rogueFrontProxy, headers: forged}, + {name: "front-proxy CA cert with disallowed CN", cert: &wrongName, headers: forged}, + {name: "unknown bearer token", headers: bearer("not-a-real-token")}, + {name: "unknown bearer token with forged headers", headers: mergeHeaders(bearer("not-a-real-token"), forged)}, + } + paths := []string{ + templatesTestNS, workspacesTestNS, templatesAllNS, templatesTestNS + "/" + testTemplateName, + templatesTestNS + "?watch=true", "/apis", "/apis/aggregation.coder.com/v1alpha1", "/version", + "/openapi/v2", "/healthz/ping", "/readyz/informer-sync", "/livezX", "/metrics", "/", + } + for _, c := range credentials { + for _, path := range paths { + status, body := f.server.do(t, c.cert, http.MethodGet, path, c.headers, "") + if status != http.StatusUnauthorized { + t.Errorf("%s GET %s: status=%d, want 401; body=%.200s", c.name, path, status, body) + } + } + status, _ := f.server.do(t, c.cert, http.MethodDelete, templatesTestNS+"/"+testTemplateName, c.headers, "") + if status != http.StatusUnauthorized { + t.Errorf("%s DELETE: status=%d, want 401", c.name, status) + } + } + if calls := f.server.provider.calls.Load(); calls != 0 { + t.Fatalf("unauthenticated requests reached the Coder backend %d times", calls) + } + if sars := f.kube.recordedSARs(); len(sars) != 0 { + t.Fatalf("unauthenticated requests must be rejected before authorization, got SARs: %v", sars) + } +} + +func TestDelegatedAuthAnonymousHealthPathsOnly(t *testing.T) { + f := newAuthFixture(t, nil) + for _, path := range []string{"/healthz", "/livez", "/readyz"} { + if status, body := f.server.do(t, nil, http.MethodGet, path, nil, ""); status != http.StatusOK { + t.Errorf("anonymous GET %s: status=%d body=%s", path, status, body) + } + } +} + +// TestDelegatedAuthFrontProxyIdentityAndSARAttributes checks that a request relayed by the +// front proxy is authorized for the relayed user with the exact verb/resource/namespace, and that +// every denial happens before the Coder backend is called. +func TestDelegatedAuthFrontProxyIdentityAndSARAttributes(t *testing.T) { + f := newAuthFixture(t, func(k *fakeKubeAPI) { + k.setDecide(func(spec authorizationv1.SubjectAccessReviewSpec) bool { + ra := spec.ResourceAttributes + return spec.User == "alice" && ra != nil && ra.Verb == "list" && ra.Group == aggGroup && + ra.Resource == "codertemplates" && ra.Namespace == "test-ns" + }) + }) + cert := f.frontProxyCert(t) + alice := mergeHeaders(remoteUser("alice", "team-a"), map[string]string{"X-Remote-Extra-Scopes": "read"}) + + status, body := f.server.do(t, cert, http.MethodGet, templatesTestNS, alice, "") + if status != http.StatusOK || !strings.Contains(body, testTemplateName) { + t.Fatalf("allowed list: status=%d body=%.300s", status, body) + } + sars := f.kube.recordedSARs() + if len(sars) != 1 { + t.Fatalf("expected one SAR, got %d: %v", len(sars), sars) + } + spec := sars[0] + if spec.User != "alice" || spec.ResourceAttributes == nil { + t.Fatalf("unexpected SAR: %s", describeSAR(spec)) + } + mustContain(t, spec.Groups, "team-a") + mustContain(t, spec.Groups, user.AllAuthenticated) + if got := spec.Extra["scopes"]; len(got) != 1 || got[0] != "read" { + t.Fatalf("expected extra scopes=[read], got %v", spec.Extra) + } + if ra := spec.ResourceAttributes; ra.Version != "v1alpha1" || ra.Name != "" { + t.Fatalf("unexpected resource attributes: %+v", ra) + } + callsAfterAllowed := f.server.provider.calls.Load() + if callsAfterAllowed == 0 { + t.Fatal("allowed list must reach the backend") + } + + denied := []struct { + method, path, body, contentType string + verb, resource, namespace, name string + }{ + {method: http.MethodGet, path: templatesOtherNS, verb: "list", resource: "codertemplates", namespace: "other-ns"}, + {method: http.MethodGet, path: templatesAllNS, verb: "list", resource: "codertemplates"}, + {method: http.MethodGet, path: templatesTestNS + "/" + testTemplateName, verb: "get", resource: "codertemplates", namespace: "test-ns", name: testTemplateName}, + {method: http.MethodGet, path: templatesTestNS + "?watch=true", verb: "watch", resource: "codertemplates", namespace: "test-ns"}, + {method: http.MethodPost, path: templatesTestNS, body: `{"apiVersion":"aggregation.coder.com/v1alpha1","kind":"CoderTemplate","metadata":{"name":"default.x"}}`, verb: "create", resource: "codertemplates", namespace: "test-ns"}, + {method: http.MethodPut, path: templatesTestNS + "/" + testTemplateName, body: `{"apiVersion":"aggregation.coder.com/v1alpha1","kind":"CoderTemplate","metadata":{"name":"default.my-template"}}`, verb: "update", resource: "codertemplates", namespace: "test-ns", name: testTemplateName}, + {method: http.MethodPatch, path: templatesTestNS + "/" + testTemplateName, body: `{"spec":{"running":false}}`, contentType: "application/merge-patch+json", verb: "patch", resource: "codertemplates", namespace: "test-ns", name: testTemplateName}, + {method: http.MethodDelete, path: templatesTestNS + "/" + testTemplateName, verb: "delete", resource: "codertemplates", namespace: "test-ns", name: testTemplateName}, + {method: http.MethodGet, path: workspacesTestNS, verb: "list", resource: "coderworkspaces", namespace: "test-ns"}, + {method: http.MethodGet, path: workspacesTestNS + "?watch=true", verb: "watch", resource: "coderworkspaces", namespace: "test-ns"}, + {method: http.MethodPost, path: workspacesTestNS, body: `{"apiVersion":"aggregation.coder.com/v1alpha1","kind":"CoderWorkspace","metadata":{"name":"default.testuser.x"}}`, verb: "create", resource: "coderworkspaces", namespace: "test-ns"}, + {method: http.MethodPatch, path: workspacesTestNS + "/" + testWorkspaceName, body: `{"spec":{"running":true}}`, contentType: "application/merge-patch+json", verb: "patch", resource: "coderworkspaces", namespace: "test-ns", name: testWorkspaceName}, + {method: http.MethodDelete, path: workspacesTestNS + "/" + testWorkspaceName, verb: "delete", resource: "coderworkspaces", namespace: "test-ns", name: testWorkspaceName}, + } + for _, d := range denied { + f.kube.resetSARs() + headers := alice + if d.contentType != "" { + headers = mergeHeaders(alice, map[string]string{"Content-Type": d.contentType}) + } + status, body := f.server.do(t, cert, d.method, d.path, headers, d.body) + if status != http.StatusForbidden { + t.Errorf("%s %s: status=%d, want 403; body=%.200s", d.method, d.path, status, body) + continue + } + sars := f.kube.recordedSARs() + if len(sars) != 1 || sars[0].ResourceAttributes == nil { + t.Errorf("%s %s: expected one resource SAR, got %v", d.method, d.path, sars) + continue + } + ra := sars[0].ResourceAttributes + if sars[0].User != "alice" || ra.Verb != d.verb || ra.Group != aggGroup || ra.Resource != d.resource || ra.Namespace != d.namespace || ra.Name != d.name { + t.Errorf("%s %s: SAR %s, want verb=%s resource=%s ns=%q name=%q", d.method, d.path, describeSAR(sars[0]), d.verb, d.resource, d.namespace, d.name) + } + } + if calls := f.server.provider.calls.Load(); calls != callsAfterAllowed { + t.Fatalf("denied requests reached the Coder backend: calls %d -> %d", callsAfterAllowed, calls) + } +} + +// TestDelegatedAuthHeadersCannotOverrideCredentialIdentity covers forged privileged-group headers +// sent alongside a real (unprivileged) credential. +func TestDelegatedAuthHeadersCannotOverrideCredentialIdentity(t *testing.T) { + f := newAuthFixture(t, func(k *fakeKubeAPI) { + k.setToken("sa-token", authenticationv1.UserInfo{Username: "system:serviceaccount:test-ns:reader", Groups: []string{"system:serviceaccounts"}}) + k.setDecide(func(spec authorizationv1.SubjectAccessReviewSpec) bool { + ra := spec.ResourceAttributes + return spec.User == "system:serviceaccount:test-ns:reader" && ra != nil && ra.Verb == "list" && ra.Namespace == "test-ns" + }) + }) + forged := remoteUser("kubernetes-admin", "system:masters") + + // Valid, unprivileged bearer token alone: allowed for its own RBAC. + if status, body := f.server.do(t, nil, http.MethodGet, templatesTestNS, bearer("sa-token"), ""); status != http.StatusOK { + t.Fatalf("bearer list: status=%d body=%.200s", status, body) + } + callsAfterAllowed := f.server.provider.calls.Load() + + f.kube.resetSARs() + status, _ := f.server.do(t, nil, http.MethodDelete, templatesTestNS+"/"+testTemplateName, mergeHeaders(bearer("sa-token"), forged), "") + if status != http.StatusForbidden { + t.Fatalf("bearer + forged masters delete: status=%d, want 403", status) + } + sars := f.kube.recordedSARs() + if len(sars) != 1 || sars[0].User != "system:serviceaccount:test-ns:reader" { + t.Fatalf("expected SAR for the token user, got %v", sars) + } + mustNotContain(t, sars[0].Groups, "system:masters") + + // Ordinary cluster client cert plus forged headers: identity is the certificate subject. + bob := f.clusterCA.clientCert(t, "bob", "team-b") + f.kube.resetSARs() + status, _ = f.server.do(t, &bob, http.MethodGet, templatesTestNS, forged, "") + if status != http.StatusForbidden { + t.Fatalf("client cert + forged masters: status=%d, want 403", status) + } + sars = f.kube.recordedSARs() + if len(sars) != 1 || sars[0].User != "bob" { + t.Fatalf("expected SAR for cert user bob, got %v", sars) + } + mustContain(t, sars[0].Groups, "team-b") + mustNotContain(t, sars[0].Groups, "system:masters") + + if calls := f.server.provider.calls.Load(); calls != callsAfterAllowed { + t.Fatalf("denied requests reached the Coder backend: calls %d -> %d", callsAfterAllowed, calls) + } +} + +// TestDelegatedAuthPrivilegedGroupBypassIsVendoredBehavior documents the retained vendored +// default: an authenticated system:masters member is authorized locally without SAR. +func TestDelegatedAuthPrivilegedGroupBypassIsVendoredBehavior(t *testing.T) { + f := newAuthFixture(t, func(k *fakeKubeAPI) { k.setFailures(false, true) }) + status, body := f.server.do(t, f.frontProxyCert(t), http.MethodGet, templatesTestNS, remoteUser("admin", "system:masters"), "") + if status != http.StatusOK { + t.Fatalf("front-proxied system:masters list during SAR outage: status=%d body=%.200s", status, body) + } + if sars := f.kube.recordedSARs(); len(sars) != 0 { + t.Fatalf("system:masters must not need SAR, got %v", sars) + } +} + +func TestDelegatedAuthFailsClosedWhenDelegationUnavailable(t *testing.T) { + f := newAuthFixture(t, func(k *fakeKubeAPI) { + k.setToken("fresh-token", authenticationv1.UserInfo{Username: "fresh"}) + k.setDecide(allowAll) + k.setFailures(true, true) + }) + cert := f.frontProxyCert(t) + + status, body := f.server.do(t, cert, http.MethodGet, templatesTestNS, remoteUser("carol", "team-c"), "") + if status < 400 || strings.Contains(body, testTemplateName) { + t.Fatalf("SAR outage (cold cache): status=%d body=%.200s", status, body) + } + if status, _ := f.server.do(t, nil, http.MethodGet, templatesTestNS, bearer("fresh-token"), ""); status != http.StatusUnauthorized { + t.Fatalf("TokenReview outage: status=%d, want 401", status) + } + + // Kubernetes API completely unreachable. + f.kube.server.Close() + status, body = f.server.do(t, cert, http.MethodDelete, templatesTestNS+"/"+testTemplateName, remoteUser("dave"), "") + if status < 400 { + t.Fatalf("Kubernetes API unreachable: status=%d body=%.200s", status, body) + } + if calls := f.server.provider.calls.Load(); calls != 0 { + t.Fatalf("requests during a delegation outage reached the Coder backend %d times", calls) + } +} + +// TestDelegatedAuthCacheTTLs documents that each delegated decision is cached for its own TTL and +// that an outage denies once the cached entry expires. +func TestDelegatedAuthCacheTTLs(t *testing.T) { + authn, authz := newDelegatedAuthOptions("unused") + if authn.CacheTTL != 10*time.Second || authz.AllowCacheTTL != 10*time.Second || authz.DenyCacheTTL != 10*time.Second { + t.Fatalf("documented default TTLs changed: token=%s allow=%s deny=%s", authn.CacheTTL, authz.AllowCacheTTL, authz.DenyCacheTTL) + } + + const ttl = 300 * time.Millisecond + kube := newFakeKubeAPI(t) + kube.setDecide(allowAll) + kube.setToken("cached-token", authenticationv1.UserInfo{Username: "erin"}) + authn, authz = kube.options(false) + authn.CacheTTL, authz.AllowCacheTTL, authz.DenyCacheTTL = ttl, ttl, ttl + server := startAuthTestServer(t, authn, authz) + + if status, _ := server.do(t, nil, http.MethodGet, templatesTestNS, bearer("cached-token"), ""); status != http.StatusOK { + t.Fatalf("warm-up: status=%d", status) + } + kube.setFailures(true, true) + if status, _ := server.do(t, nil, http.MethodGet, templatesTestNS, bearer("cached-token"), ""); status != http.StatusOK { + t.Fatalf("within TTL the cached decisions apply: status=%d", status) + } + time.Sleep(2 * ttl) + if status, _ := server.do(t, nil, http.MethodGet, templatesTestNS, bearer("cached-token"), ""); status != http.StatusUnauthorized { + t.Fatalf("after TTL expiry during an outage: status=%d, want 401", status) + } +} + +func TestDelegatedAuthConfigMapTrustMaterial(t *testing.T) { + frontProxyCA := newTestCA(t, "front-proxy") + validData := func() map[string]string { + names, _ := json.Marshal([]string{frontProxyName}) + users, _ := json.Marshal([]string{"X-Remote-User"}) + groups, _ := json.Marshal([]string{"X-Remote-Group"}) + extra, _ := json.Marshal([]string{"X-Remote-Extra-"}) + return map[string]string{ + "requestheader-client-ca-file": string(frontProxyCA.pem()), + "requestheader-allowed-names": string(names), + "requestheader-username-headers": string(users), + "requestheader-group-headers": string(groups), + "requestheader-extra-headers-prefix": string(extra), + } + } + tests := []struct { + name string + data func() map[string]string // nil: ConfigMap absent + wantOK bool + }{ + {name: "valid", data: validData, wantOK: true}, + {name: "absent"}, + {name: "CA key missing", data: func() map[string]string { d := validData(); delete(d, "requestheader-client-ca-file"); return d }}, + {name: "CA empty", data: func() map[string]string { d := validData(); d["requestheader-client-ca-file"] = ""; return d }}, + {name: "CA malformed", data: func() map[string]string { + d := validData() + d["requestheader-client-ca-file"] = "not a certificate" + return d + }}, + {name: "username headers missing", data: func() map[string]string { d := validData(); delete(d, "requestheader-username-headers"); return d }}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + kube := newFakeKubeAPI(t) + kube.setDecide(allowAll) + if tt.data != nil { + kube.configMap = &corev1.ConfigMap{ + TypeMeta: metav1.TypeMeta{APIVersion: "v1", Kind: "ConfigMap"}, + ObjectMeta: metav1.ObjectMeta{Name: "extension-apiserver-authentication", Namespace: "kube-system", ResourceVersion: "1"}, + Data: tt.data(), + } + } + authn, authz := kube.options(true) + server := startAuthTestServer(t, authn, authz) + cert := frontProxyCA.clientCert(t, frontProxyName) + + // The CA controller loads asynchronously after startup; poll for the positive case. + deadline := time.Now().Add(10 * time.Second) + var status int + for { + status, _ = server.do(t, &cert, http.MethodGet, templatesTestNS, remoteUser("alice"), "") + if status == http.StatusOK || !tt.wantOK || time.Now().After(deadline) { + break + } + time.Sleep(100 * time.Millisecond) + } + if tt.wantOK && status != http.StatusOK { + t.Fatalf("valid ConfigMap: front-proxy status=%d, want 200", status) + } + if !tt.wantOK { + // Give the async loader the same chance it had in the positive case. + time.Sleep(time.Second) + if status, _ = server.do(t, &cert, http.MethodGet, templatesTestNS, remoteUser("alice"), ""); status != http.StatusUnauthorized { + t.Fatalf("front-proxy status=%d, want 401", status) + } + if calls := server.provider.calls.Load(); calls != 0 { + t.Fatalf("rejected requests reached the Coder backend %d times", calls) + } + } + }) + } +} + +func TestDelegatedAuthConfigMapForbiddenFailsStartup(t *testing.T) { + kube := newFakeKubeAPI(t) + kube.configMapForbidden = true + authn, authz := kube.options(true) + secureServingOptions, _ := newTestSecureServing(t) + _, err := NewRecommendedConfig(NewScheme(), codecsFor(), secureServingOptions, authn, authz) + if err == nil || !strings.Contains(err.Error(), "configure delegated authentication") { + t.Fatalf("expected startup failure on forbidden ConfigMap, got %v", err) + } +} + +func TestDelegatedAuthUnreachableKubernetesAPIFailsStartup(t *testing.T) { + kube := newFakeKubeAPI(t) + kube.server.Close() + authn, authz := kube.options(true) + secureServingOptions, _ := newTestSecureServing(t) + _, err := NewRecommendedConfig(NewScheme(), codecsFor(), secureServingOptions, authn, authz) + if err == nil { + t.Fatal("expected startup failure when the Kubernetes API is unreachable") + } +} + +func TestNewRecommendedConfigRejectsMissingOrMismatchedAuthOptions(t *testing.T) { + kube := newFakeKubeAPI(t) + authn, authz := kube.options(false) + secureServingOptions, _ := newTestSecureServing(t) + if _, err := NewRecommendedConfig(NewScheme(), codecsFor(), secureServingOptions, nil, authz); err == nil || !strings.Contains(err.Error(), "must not be nil") { + t.Fatalf("expected nil authentication assertion, got %v", err) + } + if _, err := NewRecommendedConfig(NewScheme(), codecsFor(), secureServingOptions, authn, nil); err == nil || !strings.Contains(err.Error(), "must not be nil") { + t.Fatalf("expected nil authorization assertion, got %v", err) + } + authz.RemoteKubeConfigFile = filepath.Join(t.TempDir(), "other") + if _, err := NewRecommendedConfig(NewScheme(), codecsFor(), secureServingOptions, authn, authz); err == nil || !strings.Contains(err.Error(), "same Kubernetes API authority") { + t.Fatalf("expected authority mismatch assertion, got %v", err) + } +} + +type stubAuthenticator struct { + resp *authenticator.Response + ok bool + err error +} + +func (s stubAuthenticator) AuthenticateRequest(*http.Request) (*authenticator.Response, bool, error) { + return s.resp, s.ok, s.err +} + +func TestAnonymousHealthOnly(t *testing.T) { + anonymous := &authenticator.Response{User: &user.DefaultInfo{Name: user.Anonymous, Groups: []string{user.AllUnauthenticated}}} + alice := &authenticator.Response{User: &user.DefaultInfo{Name: "alice"}} + boom := errors.New("verify failed") + tests := []struct { + name string + delegate stubAuthenticator + path string + wantOK bool + wantErr bool + }{ + {name: "error passes through", delegate: stubAuthenticator{err: boom}, path: "/healthz", wantErr: true}, + {name: "not authenticated passes through", delegate: stubAuthenticator{}, path: "/healthz"}, + {name: "user on resource path", delegate: stubAuthenticator{resp: alice, ok: true}, path: templatesTestNS, wantOK: true}, + {name: "anonymous healthz", delegate: stubAuthenticator{resp: anonymous, ok: true}, path: "/healthz", wantOK: true}, + {name: "anonymous livez", delegate: stubAuthenticator{resp: anonymous, ok: true}, path: "/livez", wantOK: true}, + {name: "anonymous readyz", delegate: stubAuthenticator{resp: anonymous, ok: true}, path: "/readyz", wantOK: true}, + {name: "anonymous healthz subpath", delegate: stubAuthenticator{resp: anonymous, ok: true}, path: "/healthz/ping"}, + {name: "anonymous readyz suffix", delegate: stubAuthenticator{resp: anonymous, ok: true}, path: "/readyzX"}, + {name: "anonymous traversal", delegate: stubAuthenticator{resp: anonymous, ok: true}, path: "/healthz/../apis"}, + {name: "anonymous resource", delegate: stubAuthenticator{resp: anonymous, ok: true}, path: templatesTestNS}, + {name: "anonymous version", delegate: stubAuthenticator{resp: anonymous, ok: true}, path: "/version"}, + {name: "authenticated without user", delegate: stubAuthenticator{resp: &authenticator.Response{}, ok: true}, path: "/healthz", wantErr: true}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + req, err := http.NewRequest(http.MethodGet, "https://example.test/", nil) + if err != nil { + t.Fatal(err) + } + req.URL.Path = tt.path + resp, ok, err := anonymousHealthOnly{delegate: tt.delegate}.AuthenticateRequest(req) + if (err != nil) != tt.wantErr || ok != tt.wantOK { + t.Fatalf("ok=%v err=%v, want ok=%v wantErr=%v", ok, err, tt.wantOK, tt.wantErr) + } + if ok && resp == nil { + t.Fatal("authenticated result must carry a response") + } + }) + } + if _, _, err := (anonymousHealthOnly{}).AuthenticateRequest(httptestRequest(t)); err == nil { + t.Fatal("expected assertion for nil delegate") + } +} + +func TestResolveDelegationKubeconfig(t *testing.T) { + dir := t.TempDir() + valid := filepath.Join(dir, "valid") + kube := newFakeKubeAPI(t) + data, err := os.ReadFile(kube.kubeconfigPath) + if err != nil { + t.Fatal(err) + } + mustWrite(t, valid, string(data)) + empty := filepath.Join(dir, "empty") + mustWrite(t, empty, "apiVersion: v1\nkind: Config\n") + garbage := filepath.Join(dir, "garbage") + mustWrite(t, garbage, "{not yaml") + + homeValid := t.TempDir() + mustWrite(t, filepath.Join(homeValid, ".kube", "config"), string(data)) + homeInvalid := t.TempDir() + mustWrite(t, filepath.Join(homeInvalid, ".kube", "config"), "apiVersion: v1\nkind: Config\n") + + tests := []struct { + name string + env map[string]string + home string + want string + wantError string + }{ + {name: "explicit valid", env: map[string]string{"KUBECONFIG": valid}, want: valid}, + {name: "explicit wins over in-cluster", env: map[string]string{"KUBECONFIG": valid, "KUBERNETES_SERVICE_HOST": "10.0.0.1"}, want: valid}, + {name: "explicit missing never falls back", env: map[string]string{"KUBECONFIG": filepath.Join(dir, "absent"), "KUBERNETES_SERVICE_HOST": "10.0.0.1"}, home: homeValid, wantError: "load kubeconfig"}, + {name: "explicit empty never falls back", env: map[string]string{"KUBECONFIG": empty, "KUBERNETES_SERVICE_HOST": "10.0.0.1"}, home: homeValid, wantError: "invalid kubeconfig"}, + {name: "explicit garbage", env: map[string]string{"KUBECONFIG": garbage}, wantError: "load kubeconfig"}, + {name: "explicit list", env: map[string]string{"KUBECONFIG": valid + string(os.PathListSeparator) + valid}, wantError: "exactly one file"}, + {name: "in cluster", env: map[string]string{"KUBERNETES_SERVICE_HOST": "10.0.0.1"}, home: homeValid, want: ""}, + {name: "home config", home: homeValid, want: filepath.Join(homeValid, ".kube", "config")}, + {name: "invalid home config", home: homeInvalid, wantError: "invalid kubeconfig"}, + {name: "nothing configured", home: t.TempDir(), wantError: "no Kubernetes configuration"}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got, err := resolveDelegationKubeconfig(func(key string) string { return tt.env[key] }, tt.home) + if tt.wantError != "" { + if err == nil || !strings.Contains(err.Error(), tt.wantError) { + t.Fatalf("expected error containing %q, got path=%q err=%v", tt.wantError, got, err) + } + return + } + if err != nil || got != tt.want { + t.Fatalf("got path=%q err=%v, want %q", got, err, tt.want) + } + }) + } +} + +// TestRunWithOptionsFailsClosedWithoutKubernetesConfig exercises the production default path. +func TestRunWithOptionsFailsClosedWithoutKubernetesConfig(t *testing.T) { + t.Setenv("KUBECONFIG", "") + t.Setenv("KUBERNETES_SERVICE_HOST", "") + t.Setenv("HOME", t.TempDir()) + _, listener := newTestSecureServing(t) + ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) + defer cancel() + err := RunWithOptions(ctx, Options{Listener: listener}) + if err == nil || !strings.Contains(err.Error(), "no Kubernetes configuration") { + t.Fatalf("expected startup failure without Kubernetes config, got %v", err) + } + + t.Setenv("KUBECONFIG", filepath.Join(t.TempDir(), "absent")) + if err := RunWithOptions(ctx, Options{Listener: listener}); err == nil || !strings.Contains(err.Error(), "load kubeconfig") { + t.Fatalf("expected startup failure with an invalid explicit KUBECONFIG, got %v", err) + } +} + +// TestRunWithOptionsFailsClosedInClusterWithoutServiceAccount: the in-cluster path must not start +// without delegation credentials (guards RemoteKubeConfigFileOptional=false). +func TestRunWithOptionsFailsClosedInClusterWithoutServiceAccount(t *testing.T) { + if _, err := os.Stat("/var/run/secrets/kubernetes.io/serviceaccount/token"); err == nil { + t.Skip("running inside a Pod with a ServiceAccount token") + } + t.Setenv("KUBECONFIG", "") + t.Setenv("KUBERNETES_SERVICE_HOST", "10.0.0.1") + t.Setenv("KUBERNETES_SERVICE_PORT", "443") + _, listener := newTestSecureServing(t) + // Bounded: a regression that starts the server anyway must fail the test, not hang it. + ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) + defer cancel() + err := RunWithOptions(ctx, Options{Listener: listener}) + if err == nil || !strings.Contains(err.Error(), "delegated") { + t.Fatalf("expected startup failure without in-cluster credentials, got %v", err) + } + if err := RunWithOptions(ctx, Options{Listener: listener, Authentication: newFakeKubeAPIOptionsOnly(t)}); err == nil || !strings.Contains(err.Error(), "set together") { + t.Fatalf("expected assertion when only one option set is provided, got %v", err) + } +} diff --git a/internal/app/apiserverapp/integration_test.go b/internal/app/apiserverapp/integration_test.go index 93e02588..9880306b 100644 --- a/internal/app/apiserverapp/integration_test.go +++ b/internal/app/apiserverapp/integration_test.go @@ -18,6 +18,8 @@ import ( "time" "github.com/google/uuid" + authenticationv1 "k8s.io/api/authentication/v1" + authorizationv1 "k8s.io/api/authorization/v1" "k8s.io/apimachinery/pkg/runtime/serializer" genericoptions "k8s.io/apiserver/pkg/server/options" @@ -68,6 +70,19 @@ func TestIntegrationAggregatedAPIServerBootstrapAndList(t *testing.T) { } } +const integrationBearerToken = "integration-bearer-token" + +type bearerRoundTripper struct { + token string + base http.RoundTripper +} + +func (rt bearerRoundTripper) RoundTrip(req *http.Request) (*http.Response, error) { + clone := req.Clone(req.Context()) + clone.Header.Set("Authorization", "Bearer "+rt.token) + return rt.base.RoundTrip(clone) +} + // integrationAggregatedAPIServer is a running in-process aggregated API server // backed by the integration mock Coder server. type integrationAggregatedAPIServer struct { @@ -125,7 +140,15 @@ func startIntegrationAggregatedAPIServer(t *testing.T) integrationAggregatedAPIS secureServingOptions.ServerCert.CertDirectory = "" secureServingOptions.ServerCert.PairName = "" - recommendedConfig, err := NewRecommendedConfig(scheme, codecs, secureServingOptions) + // Every request authenticates with a bearer token through the real delegated + // TokenReview/SubjectAccessReview path; the fake Kubernetes API allows this user everything. + kubeAPI := newFakeKubeAPI(t) + kubeAPI.setToken(integrationBearerToken, authenticationv1.UserInfo{Username: "integration-tester"}) + kubeAPI.setDecide(func(spec authorizationv1.SubjectAccessReviewSpec) bool { + return spec.User == "integration-tester" + }) + authn, authz := kubeAPI.options(false) + recommendedConfig, err := NewRecommendedConfig(scheme, codecs, secureServingOptions, authn, authz) if err != nil { t.Fatalf("build recommended config: %v", err) } @@ -178,9 +201,12 @@ func startIntegrationAggregatedAPIServer(t *testing.T) integrationAggregatedAPIS httpClient := &http.Client{ Timeout: 5 * time.Second, - Transport: &http.Transport{ - //nolint:gosec // Integration test uses ephemeral self-signed certs. - TLSClientConfig: &tls.Config{InsecureSkipVerify: true}, + Transport: bearerRoundTripper{ + token: integrationBearerToken, + base: &http.Transport{ + //nolint:gosec // Integration test uses ephemeral self-signed certs. + TLSClientConfig: &tls.Config{InsecureSkipVerify: true}, + }, }, } From ee50e841d07179881db6b1cb2687dda9e4ef1b79 Mon Sep 17 00:00:00 2001 From: Thomas Kosiewski Date: Fri, 25 Sep 2026 07:39:12 +0000 Subject: [PATCH 3/3] fix(apiserver): resolve kubeconfig paths relative to the file; fix E2E delete check MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Addresses the PR #136 review and the E2E failure. - validateKubeconfigFile now calls clientcmd.ResolveLocalPaths, so a relative tokenFile or certificate path is read relative to the kubeconfig's directory, as the vendored options' loader does, instead of the process working directory. Regression test with a relative tokenFile (fails without the fix). - The E2E RBAC step checked delete with "delete --all --dry-run=server", which succeeds without any delete request when the namespace has no templates yet. It now deletes a named template, so kube-apiserver's authorization of the delete verb is always exercised. Signed-off-by: Thomas Kosiewski --- _Generated with [`xum`](https://github.com/coder/xum) • Model: `anthropic:claude-opus-5-5` • Thinking: `medium`_ Change-Id: I1cb37dd205dc1eed0a60d4817dafd14e19eba470 --- .github/workflows/ci.yaml | 2 +- internal/app/apiserverapp/auth.go | 5 +++++ internal/app/apiserverapp/auth_test.go | 21 +++++++++++++++++++++ 3 files changed, 27 insertions(+), 1 deletion(-) diff --git a/.github/workflows/ci.yaml b/.github/workflows/ci.yaml index 820937d6..81d451e0 100644 --- a/.github/workflows/ci.yaml +++ b/.github/workflows/ci.yaml @@ -397,7 +397,7 @@ jobs: } expect_forbidden kubectl "${reader[@]}" -n default get codertemplates.aggregation.coder.com expect_forbidden kubectl "${reader[@]}" -n coder get coderworkspaces.aggregation.coder.com - expect_forbidden kubectl "${reader[@]}" -n coder delete codertemplates.aggregation.coder.com --all --dry-run=server + expect_forbidden kubectl "${reader[@]}" -n coder delete codertemplates.aggregation.coder.com coder.e2e-authz-probe --dry-run=server kubectl -n default run e2e-authn-probe --image="$PROBE_IMAGE" --restart=Never --command -- sleep 600 kubectl -n default wait --for=condition=Ready pod/e2e-authn-probe --timeout=180s diff --git a/internal/app/apiserverapp/auth.go b/internal/app/apiserverapp/auth.go index 76cb5af4..025522a1 100644 --- a/internal/app/apiserverapp/auth.go +++ b/internal/app/apiserverapp/auth.go @@ -131,6 +131,11 @@ func validateKubeconfigFile(path string) error { if err != nil { return fmt.Errorf("load kubeconfig %s for delegated authentication: %w", path, err) } + // Resolve relative file references (tokenFile, certificate paths) against the kubeconfig's own + // directory, as the vendored options' loader does, instead of the process working directory. + if err := clientcmd.ResolveLocalPaths(cfg); err != nil { + return fmt.Errorf("resolve paths in kubeconfig %s for delegated authentication: %w", path, err) + } if _, err := clientcmd.NewDefaultClientConfig(*cfg, &clientcmd.ConfigOverrides{}).ClientConfig(); err != nil { return fmt.Errorf("invalid kubeconfig %s for delegated authentication: %w", path, err) } diff --git a/internal/app/apiserverapp/auth_test.go b/internal/app/apiserverapp/auth_test.go index 384eee7e..be775733 100644 --- a/internal/app/apiserverapp/auth_test.go +++ b/internal/app/apiserverapp/auth_test.go @@ -512,6 +512,26 @@ func TestResolveDelegationKubeconfig(t *testing.T) { homeValid := t.TempDir() mustWrite(t, filepath.Join(homeValid, ".kube", "config"), string(data)) + // A kubeconfig whose tokenFile is relative to the kubeconfig's directory, not the working directory. + relDir := t.TempDir() + relative := filepath.Join(relDir, "relative") + mustWrite(t, filepath.Join(relDir, "token"), "test-token\n") + mustWrite(t, relative, `apiVersion: v1 +kind: Config +clusters: +- name: c + cluster: + server: https://127.0.0.1:1 + insecure-skip-tls-verify: true +users: +- name: u + user: + tokenFile: token +contexts: +- name: x + context: {cluster: c, user: u} +current-context: x +`) homeInvalid := t.TempDir() mustWrite(t, filepath.Join(homeInvalid, ".kube", "config"), "apiVersion: v1\nkind: Config\n") @@ -526,6 +546,7 @@ func TestResolveDelegationKubeconfig(t *testing.T) { {name: "explicit wins over in-cluster", env: map[string]string{"KUBECONFIG": valid, "KUBERNETES_SERVICE_HOST": "10.0.0.1"}, want: valid}, {name: "explicit missing never falls back", env: map[string]string{"KUBECONFIG": filepath.Join(dir, "absent"), "KUBERNETES_SERVICE_HOST": "10.0.0.1"}, home: homeValid, wantError: "load kubeconfig"}, {name: "explicit empty never falls back", env: map[string]string{"KUBECONFIG": empty, "KUBERNETES_SERVICE_HOST": "10.0.0.1"}, home: homeValid, wantError: "invalid kubeconfig"}, + {name: "explicit with tokenFile relative to the kubeconfig", env: map[string]string{"KUBECONFIG": relative}, want: relative}, {name: "explicit garbage", env: map[string]string{"KUBECONFIG": garbage}, wantError: "load kubeconfig"}, {name: "explicit list", env: map[string]string{"KUBECONFIG": valid + string(os.PathListSeparator) + valid}, wantError: "exactly one file"}, {name: "in cluster", env: map[string]string{"KUBERNETES_SERVICE_HOST": "10.0.0.1"}, home: homeValid, want: ""},