This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/main/pr-6866-d0ab10d7d63b4e40bbea15a0ea1a285c886fd9ef in repository https://gitbox.apache.org/repos/asf/texera.git
commit 209fc4152e33eb7acd4fde7a33b8402962e508dc Author: ali risheh <[email protected]> AuthorDate: Mon Sep 7 04:55:45 2026 +0000 feat(computing-unit): out-of-pod LakeFS repository mount infrastructure (#6866) ### Abstract This PR adds infrastructure to mount LakeFS repositories to any pod, later it will be used for models and datasets be mounted on computing unit pods. The goal of this PR to add `mounter` service and `S3 proxy` to file service. We also moved computing unit prefix to configuration, in the past it was "computing-unit" by default, we just moved it to configuration to have one source of truth because mounter needs to know computing unit pod name. ### What changes were proposed in this PR? Perform the FUSE mount for dataset repositories **outside** the (unprivileged) computing-unit pod — the infrastructure foundation of the dataset-mounting feature (#6606). - **`texera-mounter` DaemonSet** — a per-node privileged agent (`bin/mounter/mounter.py` + tests, dockerfile, helm daemonset/rbac/values) that runs GeeseFS on a pod's behalf. The read-only mount reaches the CU pod via Kubernetes **mount propagation**, so the pod that runs user code stays **unprivileged**. - **Per-computing-unit isolation** — mounts land at `<mount-root>/<cuid>/<repository>/<commit>`, and a CU pod's `hostPath` volume is only its own `<cuid>` subtree. Two units mounting the same version get two separate GeeseFS mounts, each authorized with its own JWT. A pod watcher unmounts a CU's directories when its pod is deleted. - **Unprivileged CU pod wiring** (`KubernetesClient`) — the propagation volume and mount env, added only when the feature is switched on (see below). - **File-service JWT S3 proxy** (`S3ProxyServlet`) — fronts the LakeFS S3 gateway: verifies the pod's JWT, checks the user's read access, re-signs to LakeFS with credentials held only server-side. No global credential enters the pod. Foundation only — nothing triggers a mount yet; the platform integration (engine client + per-CU mount API + UI + UDF bindings) comes in the follow-up PR. **Off by default (`mounter.enabled: false`).** A reviewer running this branch on Talos could not create a computing unit at all: ``` pods "computing-unit-1" is forbidden: violates PodSecurity "baseline:latest": hostPath volumes (volume "texera-mounts") ``` The CU pod was being given a `hostPath` unconditionally, and both the `baseline` and `restricted` Pod Security Standards forbid `hostPath` — so on any cluster enforcing either on the pool namespace (Talos does so by default), *every* computing unit becomes unschedulable, whether or not anyone wants to mount a dataset. It went unnoticed locally because a default minikube enforces nothing. Since no caller requests a mount until the follow-up PR, the feature is now opt-in. `mounter.enabled` gates the DaemonSet, its RBAC, the access-control-service identity and token, and — through `kubernetes.mounter-enabled` — the CU pod's `hostPath`, its mount and its env. With the flag off the chart renders no mounter object and the CU pod spec is byte-for-byte what it was before this feature existed. Enabling it requires a cluster that admits `hostPath` in the pool namespace and a privileged pod in the release namespace, so an operator opts in once that is true for them. <img width="1241" height="423" alt="Ali Texera-geeseFS (4)" src="https://github.com/user-attachments/assets/c951d777-1285-48e1-be5a-082aa4b52967" /> ### Mount request validation and the caller the mounter trusts Addressing the review on request validation: - **Every path component is validated, not just `repo`/`commit`.** `cuid`, `repositoryName` and `commitHash` are all joined into `MOUNT_ROOT/<cuid>/<repo>/<commit>`, and the directory is created **before** geesefs — and therefore LakeFS — ever sees the request, so "LakeFS rejects a bad repository" was never a defence for the *path*. `cuid` must now match `^[0-9]+$` (it is the computing unit's integer primary key); `repositoryName` and `commitHash` must each be a single safe segment, `^[A-Za-z0-9][A-Za-z0-9._-]*$` — which admits everything the platform actually sends (`dataset-<did>` and a hex digest) while rejecting a separator, a `..`, an absolute path, or a leading `-` that geesefs might read as a flag. Anything else is a `400`, and nothing is created on disk. - **`_remove_empty_dirs` had a prefix bug.** It tested `path.startswith(MOUNT_ROOT/<cuid>)`, so a sibling whose name merely began the same way — `<root>/7x` against `<root>/7` — was treated as a child and deleted. It now compares path segments (`os.path.commonpath`). The `stop_at` name and docstring were also wrong: the loop removes `MOUNT_ROOT/<cuid>` itself and stops at its parent. That is safe for a running pod — the CU's hostPath volume is `DirectoryOrCreate`, so the next mount recreates it — and both the name and the docstring now say so. - **The mounter API is safeguarded: only access-control-service can call it.** The mounter is a privileged, per-node DaemonSet, so before this merges it must not be callable by anything else. Two things enforce that, and both are decided by the kube-apiserver rather than by the mounter: - **A dedicated ServiceAccount is the identity.** This PR adds `access-control-service-service-account.yaml` and binds ACS to it. ACS is already the JWT-authenticating routing proxy that validates the user's token and checks their computing-unit access, so it is the right — and only — place for that authorization; the mounter stays a small service that mounts what one known caller asks for. - **An audience scopes the credential to the mounter.** ACS receives a projected `serviceAccountToken` bound to the audience `texera-mounter`, separate from its ordinary kube-apiserver token. `authenticate_caller` (`bin/mounter/mounter.py`) submits it to the `TokenReview` API and requires all three of: the token verifies, `texera-mounter` is among the audiences the API server echoes back, and the username equals `system:serviceaccount:<ns>:<acs-sa>`. Anything else is a `401` and no mount happens. The audience matters because without it a TokenReview validates against the API server's own audience — so a *generic* ACS token, one that leaked into a log or a crash dump, would be accepted. Bound to `texera-mounter`, only the credential minted for the mounter works. The audience and the allowed caller are fixed in `templates/base/_helpers.tpl` and are deliberately **not** settable in `values.yaml`: both sides of the contract must agree, and widening the allow-list is a security decision, not a deployment preference. `/healthz` stays open because the kubelet probes it and holds no token for this audience; `/mount` and `/mounts` are both gated. This is deliberately not a `NetworkPolicy` — a NetworkPolicy is silently unenforced on CNIs that do not implement it (EKS's VPC CNI has it off by default) and is commonly bypassed by hostPort traffic, whereas `TokenReview` holds regardless of how the request arrived. The mounter is also reachable only in-cluster: it has no Service and no Ingress, so nothing routes to it from the gateway. The path validation above holds independently of all this, so a malformed `cuid` is refused whether or not the caller authenticates. ### Any related issues, documentation, discussions? Closes #6862 · part of #6606. ### How was this PR tested? - `sbt FileService/compile ComputingUnitManagingService/compile` green. - The mounter has its own **pytest suite** (`bin/mounter/tests`, now 93 tests) covering mount, the pod-deletion reaper, dead-mount self-heal, and — new in this revision — request validation, including the reported `cuid=5/../8`, `cuid=../..` and absolute-`cuid` escapes, the equivalent `repositoryName`/`commitHash` escapes, and the sibling-directory deletion bug, each asserted both at the function level and end to end over the mounter's real HTTP surface (`400` + no `geesefs` invocation + nothing created on disk). This suite **is** run in CI: the `build / infra` job runs `pytest bin/` on ubuntu and macos. The proxy's request-parsing helpers are unit-tested (`S3ProxyServletSpec`). - Validated **end-to-end on a single-node minikube**: a Python UDF read a ~2 GB sharded PyTorch model from a propagated mount via `torch.load` with **bit-exact** output; the proxy's JWT authorization was exercised for both an authorized user (200) and an unauthorized repository (403 + refused to mount). > **Note on patch coverage:** most of this PR is inherently integration/IO code — the S3 proxy's request **forwarding + re-signing** (needs a live LakeFS gateway) and the Kubernetes **pod wiring** (fabric8 has no mock server in this repo). codecov's patch % therefore reads low even though the behavior is covered by the end-to-end validation above and by the mounter's pytest suite — which runs in CI under `build / infra` but is not instrumented by codecov, being a `bin/` script rather than a build module. We'd appreciate reviewers weighing the patch-coverage signal in that light. ### Was this PR authored or co-authored using generative AI tooling? Generated-by: Claude Opus 4.8 --------- Co-authored-by: Claude Opus 4.8 (1M context) <[email protected]> --- .github/workflows/build-and-push-images.yml | 3 + bin/dockerfiles/mounter.dockerfile | 38 ++ bin/k8s/templates/base/_helpers.tpl | 23 + .../access-control-service-deployment.yaml | 25 + .../access-control-service-service-account.yaml | 65 ++ .../templates/base/mounter/mounter-daemonset.yaml | 84 +++ bin/k8s/templates/base/mounter/mounter-rbac.yaml | 70 +++ ...workflow-computing-unit-manager-deployment.yaml | 8 + bin/k8s/values.yaml | 31 + bin/mounter/mounter.py | 566 ++++++++++++++++++ bin/mounter/tests/conftest.py | 182 ++++++ bin/mounter/tests/test_mounter.py | 661 +++++++++++++++++++++ common/config/src/main/resources/kubernetes.conf | 16 + .../common/config/EnvironmentalVariable.scala | 15 + .../texera/common/config/KubernetesConfig.scala | 11 + .../common/config/KubernetesConfigSpec.scala | 9 + .../texera/service/util/KubernetesClient.scala | 78 ++- .../texera/service/util/KubernetesClientSpec.scala | 55 ++ .../org/apache/texera/service/FileService.scala | 7 + .../texera/service/util/S3ProxyServlet.scala | 273 +++++++++ .../texera/service/util/S3ProxyServletSpec.scala | 86 +++ 21 files changed, 2294 insertions(+), 12 deletions(-) diff --git a/.github/workflows/build-and-push-images.yml b/.github/workflows/build-and-push-images.yml index cc0d86d8f4..7fd76bed7c 100644 --- a/.github/workflows/build-and-push-images.yml +++ b/.github/workflows/build-and-push-images.yml @@ -264,6 +264,9 @@ jobs: "access-control-service") image_name="texera-access-control-service" ;; + "mounter") + image_name="texera-mounter" + ;; "config-service") image_name="texera-config-service" ;; diff --git a/bin/dockerfiles/mounter.dockerfile b/bin/dockerfiles/mounter.dockerfile new file mode 100644 index 0000000000..d385e2aa57 --- /dev/null +++ b/bin/dockerfiles/mounter.dockerfile @@ -0,0 +1,38 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +# The texera-mounter runs as a per-node privileged DaemonSet. It performs the GeeseFS +# FUSE mount that computing-unit pods used to do themselves, so that CU pods (which run +# untrusted user code) can be unprivileged. It needs python3 (the mounter), fuse3 + geesefs +# (to mount), and util-linux (umount) — all part of a minimal Debian base. +FROM debian:bookworm-slim + +ARG GEESEFS_VERSION=v0.43.8 +RUN apt-get update \ + && apt-get install -y --no-install-recommends python3 fuse3 mount ca-certificates curl \ + && curl -fsSL -o /usr/local/bin/geesefs \ + "https://github.com/yandex-cloud/geesefs/releases/download/${GEESEFS_VERSION}/geesefs-linux-$(dpkg --print-architecture)" \ + && chmod 755 /usr/local/bin/geesefs \ + && apt-get clean \ + && rm -rf /var/lib/apt/lists/* \ + # allow FUSE mounts to be accessible by other users (the unprivileged CU pod's UID) + && echo "user_allow_other" >> /etc/fuse.conf + +COPY bin/mounter/mounter.py /opt/mounter/mounter.py + +EXPOSE 8100 +ENTRYPOINT ["python3", "/opt/mounter/mounter.py"] diff --git a/bin/k8s/templates/base/_helpers.tpl b/bin/k8s/templates/base/_helpers.tpl index e044b7285a..bcb8e33a2d 100644 --- a/bin/k8s/templates/base/_helpers.tpl +++ b/bin/k8s/templates/base/_helpers.tpl @@ -54,3 +54,26 @@ services fall back to the in-cluster MinIO Service and its auto-generated {{- define "texera.s3.secretAccessKeyKey" -}} {{- if .Values.storage.s3.endpoint -}}secret-access-key{{- else -}}root-password{{- end -}} {{- end -}} + +{{/* +The audience the mounter's service-account tokens are bound to. Fixed rather than +configurable: it is one half of a credential contract between access-control-service and the +mounter -- the projected token declares it and the mounter's TokenReview requires it -- so +the two must always agree, and there is no deployment in which a different value is useful. +*/}} +{{- define "texera.mounter.audience" -}} +texera-mounter +{{- end -}} + +{{/* +The service-account username allowed to request a mount from the per-node mounter, in the +form MOUNTER_ALLOWED_CALLERS expects. This is access-control-service and nothing else: it +is the deployment's authorization authority for computing units, so it is the one component +that can decide whether a given user may mount into a given CU. Deliberately not +configurable -- widening it is a security decision, not a deployment preference, and a +values override would let an install quietly hand the privileged mounter to another caller. +*/}} +{{- define "texera.mounter.allowedCallers" -}} +{{- printf "system:serviceaccount:%s:%s" .Release.Namespace .Values.accessControlService.serviceAccountName -}} +{{- end -}} + diff --git a/bin/k8s/templates/base/access-control-service/access-control-service-deployment.yaml b/bin/k8s/templates/base/access-control-service/access-control-service-deployment.yaml index 99713e7071..85fbbdd5f7 100644 --- a/bin/k8s/templates/base/access-control-service/access-control-service-deployment.yaml +++ b/bin/k8s/templates/base/access-control-service/access-control-service-deployment.yaml @@ -32,6 +32,11 @@ spec: labels: app: {{ .Release.Name }}-{{ .Values.accessControlService.name }} spec: + {{- if .Values.mounter.enabled }} + # Only needed to authenticate to the mounter; without it this service has no reason + # for a dedicated identity. + serviceAccountName: {{ .Values.accessControlService.serviceAccountName }} + {{- end }} containers: - name: {{ .Values.accessControlService.name }} image: {{ .Values.texera.imageRegistry }}/{{ .Values.accessControlService.imageName }}:{{ .Values.texera.imageTag }} @@ -76,3 +81,23 @@ spec: port: {{ .Values.accessControlService.service.port }} initialDelaySeconds: 5 periodSeconds: 5 +{{- if .Values.mounter.enabled }} + volumeMounts: + - name: mounter-token + mountPath: /var/run/secrets/texera/mounter + readOnly: true + volumes: + # The token this service presents to the per-node mounter, which verifies it with + # TokenReview. It is bound to the mounter's audience, so it authenticates nothing + # else -- not even the Kubernetes API server -- and the kubelet rotates it in place + # before it expires. access-control-service holds it because it is the only caller + # the mounter admits: it is where the deployment already decides whether a user may + # act on a computing unit. + - name: mounter-token + projected: + sources: + - serviceAccountToken: + path: token + audience: {{ include "texera.mounter.audience" . }} + expirationSeconds: 3600 +{{- end }} diff --git a/bin/k8s/templates/base/access-control-service/access-control-service-service-account.yaml b/bin/k8s/templates/base/access-control-service/access-control-service-service-account.yaml new file mode 100644 index 0000000000..44cd6bbf07 --- /dev/null +++ b/bin/k8s/templates/base/access-control-service/access-control-service-service-account.yaml @@ -0,0 +1,65 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. +{{- if .Values.mounter.enabled }} + +# Dedicated identity for the access-control-service. +# +# The access-control-service is intended to become the only component allowed to ask the +# per-node mounter to mount a dataset: it is already the JWT and computing-unit-access +# authorization proxy, so it is the natural place for the decision "may this user mount +# onto this CU?". Giving it its own identity now is what makes that switch a config +# change later, rather than a redesign -- running as the namespace's `default` +# ServiceAccount (shared with every pod that does not name one) would make the mounter +# unable to tell this service apart from anything else. +# +# The enforcement mechanism already exists and is live in this PR; only the identity in +# the allow-list is still provisional: +# +# 1. The calling pod mounts a projected `serviceAccountToken` volume bound to the +# audience `texera-mounter` and sends that token as a Bearer header on each mounter +# request. The audience binding means the token is only usable against the mounter -- +# it is not the pod's general-purpose kube-apiserver token, and a token minted for +# any other audience will not validate here. +# 2. The mounter posts the token to the kube-apiserver's TokenReview API and accepts the +# request only when the response reports `authenticated: true`, the requested audience +# among `status.audiences`, and a `status.user.username` listed in +# `mounter.allowedCallers`. Verification happens at the API server, so it holds +# regardless of how the request reached the mounter -- unlike a NetworkPolicy, which +# is unenforced on CNIs that do not implement it and is routinely bypassed by hostPort +# traffic. See authenticate_caller in bin/mounter/mounter.py. +# 3. The mounter's own ServiceAccount is bound to the built-in `system:auth-delegator` +# ClusterRole, which is what grants it permission to create TokenReviews. +# +# TODO(dataset-mount): today `mounter.allowedCallers` defaults to the computing-unit +# manager, because that is the service which actually calls the mounter. Point it at this +# account -- and move the mount endpoints behind this service -- once access-control-service +# takes over as the mount authority. Nothing else has to change. +# +# Either way, computing-unit pods are never an accepted caller even though they can reach +# the mounter's hostPort: they hold no token for this audience, so a mount request forged +# from user code fails the TokenReview regardless of what it puts in the request body. And +# because authenticating the caller only establishes who is asking, the mounter still +# validates every path component of the request itself. +# +# This account needs no RBAC rules: it is an identity to authenticate as, not a client +# of the Kubernetes API. +apiVersion: v1 +kind: ServiceAccount +metadata: + name: {{ .Values.accessControlService.serviceAccountName }} + namespace: {{ .Release.Namespace }} +{{- end }} diff --git a/bin/k8s/templates/base/mounter/mounter-daemonset.yaml b/bin/k8s/templates/base/mounter/mounter-daemonset.yaml new file mode 100644 index 0000000000..43927260b0 --- /dev/null +++ b/bin/k8s/templates/base/mounter/mounter-daemonset.yaml @@ -0,0 +1,84 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. +{{- if .Values.mounter.enabled }} + +# Per-node privileged mounter that performs the GeeseFS FUSE mount on behalf of +# (unprivileged) computing-unit pods. It mounts under a host directory and the mount +# propagates into each CU pod via mountPropagation, so CU pods — which run untrusted user +# code — never need to be privileged. This is the ONLY privileged component of the mount +# feature, and it runs only trusted mount code (never user code). +apiVersion: apps/v1 +kind: DaemonSet +metadata: + name: {{ .Release.Name }}-mounter + namespace: {{ .Release.Namespace }} + labels: + app: {{ .Release.Name }}-mounter +spec: + selector: + matchLabels: + app: {{ .Release.Name }}-mounter + template: + metadata: + labels: + app: {{ .Release.Name }}-mounter + spec: + serviceAccountName: {{ .Values.mounter.serviceAccountName }} + tolerations: + - operator: "Exists" + containers: + - name: mounter + image: {{ .Values.texera.imageRegistry }}/{{ .Values.mounter.imageName }}:{{ .Values.texera.imageTag }} + imagePullPolicy: {{ .Values.mounter.imagePullPolicy }} + securityContext: + privileged: true + ports: + - name: mounter + containerPort: {{ .Values.mounter.port }} + hostPort: {{ .Values.mounter.port }} + env: + - name: MOUNTER_PORT + value: "{{ .Values.mounter.port }}" + - name: MOUNT_ROOT + value: "{{ .Values.mounter.hostMountRoot }}" + - name: POOL_NAMESPACE + value: "{{ .Values.workflowComputingUnitPool.namespace }}" + - name: CU_POD_NAME_PREFIX + value: "{{ .Values.workflowComputingUnitPool.podNamePrefix }}" + # Only a caller holding a service-account token minted for this audience, and + # belonging to one of these identities, may request a mount. See + # authenticate_caller in bin/mounter/mounter.py. + - name: MOUNTER_AUDIENCE + value: "{{ include "texera.mounter.audience" . }}" + - name: MOUNTER_ALLOWED_CALLERS + value: "{{ include "texera.mounter.allowedCallers" . }}" + readinessProbe: + httpGet: + path: /healthz + port: mounter + initialDelaySeconds: 3 + periodSeconds: 10 + volumeMounts: + - name: texera-mounts + mountPath: {{ .Values.mounter.hostMountRoot }} + mountPropagation: Bidirectional + volumes: + - name: texera-mounts + hostPath: + path: {{ .Values.mounter.hostMountRoot }} + type: DirectoryOrCreate +{{- end }} diff --git a/bin/k8s/templates/base/mounter/mounter-rbac.yaml b/bin/k8s/templates/base/mounter/mounter-rbac.yaml new file mode 100644 index 0000000000..41e98f2de9 --- /dev/null +++ b/bin/k8s/templates/base/mounter/mounter-rbac.yaml @@ -0,0 +1,70 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. +{{- if .Values.mounter.enabled }} + +# The mounter watches computing-unit pods (in the pool namespace) so it can unmount a CU's +# mounts as soon as its pod is deleted. It needs only read access to pods. +# +# It also verifies its callers with the TokenReview API, which is what the cluster-wide +# system:auth-delegator binding below grants. That role is the standard way to let a service +# authenticate tokens presented to it; it confers no other access. +apiVersion: v1 +kind: ServiceAccount +metadata: + name: {{ .Values.mounter.serviceAccountName }} + namespace: {{ .Release.Namespace }} +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: {{ .Release.Name }}-mounter + namespace: {{ .Values.workflowComputingUnitPool.namespace }} +rules: + - apiGroups: [""] + resources: ["pods"] + verbs: ["get", "list", "watch"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: {{ .Release.Name }}-mounter-binding + namespace: {{ .Values.workflowComputingUnitPool.namespace }} +subjects: + - kind: ServiceAccount + name: {{ .Values.mounter.serviceAccountName }} + namespace: {{ .Release.Namespace }} +roleRef: + kind: Role + name: {{ .Release.Name }}-mounter + apiGroup: rbac.authorization.k8s.io +--- +# Lets the mounter submit TokenReviews, so it can establish who is calling /mount before +# performing a privileged mount on their behalf. system:auth-delegator is a built-in +# ClusterRole covering exactly authentication.k8s.io/tokenreviews (and subjectaccessreviews). +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRoleBinding +metadata: + name: {{ .Release.Name }}-mounter-auth-delegator +subjects: + - kind: ServiceAccount + name: {{ .Values.mounter.serviceAccountName }} + namespace: {{ .Release.Namespace }} +roleRef: + kind: ClusterRole + name: system:auth-delegator + apiGroup: rbac.authorization.k8s.io +{{- end }} diff --git a/bin/k8s/templates/base/workflow-computing-unit-manager/workflow-computing-unit-manager-deployment.yaml b/bin/k8s/templates/base/workflow-computing-unit-manager/workflow-computing-unit-manager-deployment.yaml index ea61b242d1..d55eb5f10c 100644 --- a/bin/k8s/templates/base/workflow-computing-unit-manager/workflow-computing-unit-manager-deployment.yaml +++ b/bin/k8s/templates/base/workflow-computing-unit-manager/workflow-computing-unit-manager-deployment.yaml @@ -64,8 +64,16 @@ spec: value: {{ .Values.workflowComputingUnitPool.namespace }} - name: KUBERNETES_COMPUTE_UNIT_SERVICE_NAME value: {{ .Values.workflowComputingUnitPool.name }}-svc + - name: KUBERNETES_COMPUTE_UNIT_POD_NAME_PREFIX + value: {{ .Values.workflowComputingUnitPool.podNamePrefix }} - name: KUBERNETES_IMAGE_NAME value: {{ .Values.texera.imageRegistry }}/{{ .Values.workflowComputingUnitPool.imageName }}:{{ .Values.texera.imageTag }} + - name: KUBERNETES_MOUNTER_ENABLED + value: "{{ .Values.mounter.enabled }}" + {{- if .Values.mounter.enabled }} + - name: KUBERNETES_MOUNTER_HOST_ROOT + value: "{{ .Values.mounter.hostMountRoot }}" + {{- end }} # TexeraDB Access - name: STORAGE_JDBC_URL value: jdbc:postgresql://{{ .Release.Name }}-postgresql:5432/texera_db?currentSchema=texera_db,public diff --git a/bin/k8s/values.yaml b/bin/k8s/values.yaml index 597a65ae10..94f30a5b32 100644 --- a/bin/k8s/values.yaml +++ b/bin/k8s/values.yaml @@ -232,6 +232,11 @@ accessControlService: name: access-control-service numOfPods: 1 imageName: texera-access-control-service + # Dedicated identity, so the mounter can authenticate the access-control-service as the + # only component permitted to request a dataset mount. Created and attached only when + # mounter.enabled is true. See + # templates/base/access-control-service/access-control-service-service-account.yaml. + serviceAccountName: texera-access-control-service-account service: type: ClusterIP port: 9096 @@ -288,6 +293,9 @@ workflowComputingUnitPool: name: texera-workflow-computing-unit # Note: the namespace of the workflow computing unit pool might conflict when there are multiple texera deployments in the same cluster namespace: texera-workflow-computing-unit-pool + # Computing unit pods are named "<podNamePrefix>-<cuid>". The manager creates them under + # this name and the mounter reads a CU's id back out of it, so both are given this value. + podNamePrefix: computing-unit # Max number of resources allocated for computing units maxRequestedResources: cpu: 100 @@ -298,6 +306,29 @@ workflowComputingUnitPool: port: 8085 targetPort: 8085 +# Per-node privileged DaemonSet that performs the GeeseFS FUSE mount for computing-unit +# pods, so the CU pods stay unprivileged. hostMountRoot MUST match the computing-unit +# manager's kubernetes.mounter-host-root: the manager gives each CU pod a hostPath under +# that root, and the mounter mounts into the same subtree. +mounter: + # Off by default. When enabled, computing-unit pods gain a hostPath volume, which the + # `baseline` and `restricted` Pod Security Standards forbid -- on a cluster that enforces + # either on the pool namespace, every CU pod would be rejected. Mounting is also not + # reachable yet: no caller requests a mount until the platform integration lands. So an + # operator opts in once their cluster admits the hostPath, and until then this feature + # changes nothing about how computing units run. + enabled: false + imageName: texera-mounter + imagePullPolicy: IfNotPresent + port: 8100 + hostMountRoot: /var/lib/texera-mounts + serviceAccountName: texera-mounter-service-account + # The mounter authenticates its callers: a caller presents a service-account token + # projected for the mounter's audience, which the mounter verifies with TokenReview and + # matches against the one identity it serves. Neither the audience nor that identity is + # settable here -- both are fixed in templates/base/_helpers.tpl, because widening them is + # a security decision rather than a deployment preference. + texeraEnvVars: - name: USER_SYS_ADMIN_USERNAME value: "texera" diff --git a/bin/mounter/mounter.py b/bin/mounter/mounter.py new file mode 100644 index 0000000000..3dc939679f --- /dev/null +++ b/bin/mounter/mounter.py @@ -0,0 +1,566 @@ +#!/usr/bin/env python3 +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +""" +texera-mounter: a per-node privileged service that performs GeeseFS FUSE mounts on behalf +of (unprivileged) computing-unit pods. + +A permitted caller POSTs /mount with {cuid, repositoryName, commitHash, jwt, +fileServiceBase} once it has validated the JWT and the user's access to that computing +unit. "Permitted" is decided by the API server: the caller presents a service-account token +bound to MOUNTER_AUDIENCE and the mounter verifies it with TokenReview (see +authenticate_caller). Every path component of the request is validated here as well (see +CUID_PATTERN / PATH_SEGMENT_PATTERN), because authenticating the caller says who is asking, +not that what they asked for is a safe path. The mounter then runs GeeseFS +against file-service's JWT-authenticated S3 proxy (passing the user's JWT as the S3 access +key) and mounts read-only under MOUNT_ROOT/<cuid>/<repo>/<commit>. +That host directory is bind-mounted (mountPropagation: Bidirectional) into the mounter and +propagates back into the CU pod (mountPropagation: HostToContainer), so the CU pod sees +the mount without any privilege of its own. + +The mounter holds no LakeFS credentials; authorization stays entirely in file-service. +A background watcher unmounts a CU's directories as soon as its pod is deleted. +""" + +import json +import os +import re +import shutil +import ssl +import subprocess +import threading +import time +import urllib.request +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from urllib.parse import parse_qs, urlparse + +MOUNT_ROOT = os.environ.get("MOUNT_ROOT", "/var/lib/texera-mounts") +MOUNTER_PORT = int(os.environ.get("MOUNTER_PORT", "8100")) +POOL_NAMESPACE = os.environ.get("POOL_NAMESPACE", "texera-workflow-computing-unit-pool") +MOUNT_SECRET_PLACEHOLDER = "texera-jwt-mount" +MOUNT_TIMEOUT_S = 30 +# How long the API server keeps a watch open. Each expiry re-runs the reconcile, which is +# the safety net for events missed while disconnected and for orphans still busy earlier. +WATCH_TIMEOUT_S = 300 +WATCH_RETRY_S = 10 + +SA_DIR = "/var/run/secrets/kubernetes.io/serviceaccount" +# Computing-unit pods are named "<prefix>-<cuid>"; the helm chart gives this and the +# computing-unit manager's kubernetes.compute-unit-pod-name-prefix the same value. +CU_POD_NAME_PREFIX = os.environ.get("CU_POD_NAME_PREFIX", "computing-unit") +# Every component of the mount path is caller-supplied, so each is validated rather than +# trusted. os.path.join discards everything before an absolute component and happily walks +# through "..", so an unchecked component relocates the mount anywhere on the node -- and +# the directory is created before geesefs (and therefore LakeFS) ever sees the request, so +# "LakeFS will reject it" is not a defence for the path. +# +# A cuid is the computing unit's numeric primary key. Repositories are named +# "dataset-<did>" and commits are hex digests, so both fit a conservative single-segment +# rule: no separator, no "..", and a leading alphanumeric so a value can never be mistaken +# for a geesefs flag. +CUID_PATTERN = re.compile(r"^[0-9]+$") +PATH_SEGMENT_PATTERN = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._-]*$") +PROC_MOUNTS = "/proc/mounts" + +# ---- caller authentication ---- +# The mounter is the only privileged component in the mount path -- root on every node, with +# a hostPath and Bidirectional mount propagation -- and it listens on a hostPort, so anything +# routable to a node IP can reach it, including computing-unit pods running untrusted user +# code. Callers therefore present a Kubernetes service-account token bound to +# MOUNTER_AUDIENCE, which is verified here with the TokenReview API and matched against +# MOUNTER_ALLOWED_CALLERS. That check is enforced by the API server, unlike a NetworkPolicy, +# which is silently unenforced on CNIs that do not implement it and is routinely bypassed by +# hostPort traffic (DNAT'd at the node, so pod selectors never see it). +# +# An empty allow-list denies every request: the chart always sets this, and a mounter that +# cannot tell who is calling must not mount anything. +# +# The chart admits one caller, access-control-service, because that service is where the +# deployment already decides whether a user may act on a computing unit. The mounter itself +# cannot make that decision -- every request reaches it under one service identity, so it +# has no way to tell whose workflow is asking -- which is exactly why the component that +# does know must be the only one able to ask. +MOUNTER_AUDIENCE = os.environ.get("MOUNTER_AUDIENCE", "texera-mounter") +ALLOWED_CALLERS = frozenset( + caller.strip() + for caller in os.environ.get("MOUNTER_ALLOWED_CALLERS", "").split(",") + if caller.strip() +) + + +class Unauthorized(Exception): + """The caller did not present a token this mounter accepts.""" + + +def log(msg): + print(f"[mounter] {msg}", flush=True) + + +def _unescape(field): + """Decode the four characters /proc/mounts escapes octally.""" + for code, char in (("\\040", " "), ("\\011", "\t"), ("\\012", "\n"), ("\\134", "\\")): + field = field.replace(code, char) + return field + + +def mount_targets_under(path): + """Mount targets at or under `path`, deepest first. + + Read from /proc/mounts rather than using os.path.ismount(): when a FUSE server dies + (the mounter being restarted kills every GeeseFS it started) its mount entry survives + here, but stat()ing the mount point fails with ENOTCONN, which os.path.ismount() + reports as "not a mount point". Such a dead mount would then never be unmounted, and + its directory could never be removed. + """ + path = os.path.normpath(path) + prefix = path + "/" + targets = [] + try: + with open(PROC_MOUNTS) as mounts: + for line in mounts: + fields = line.split() + if len(fields) < 2: + continue + target = _unescape(fields[1]) + if target == path or target.startswith(prefix): + targets.append(target) + except OSError as e: + log(f"reading {PROC_MOUNTS} failed: {e}") + # Deepest first, so nested mounts are detached before their parents. + return sorted(targets, key=len, reverse=True) + + +def is_mounted(path): + """True if `path` itself is a mount point, alive or dead.""" + return os.path.normpath(path) in mount_targets_under(path) + + +def _responds(path): + """True if `path` can be stat()ed — false for a mount whose FUSE server is gone.""" + try: + os.stat(path) + return True + except OSError: + return False + + +def validated_cuid(cuid): + """Return `cuid` if it is a single numeric path segment, else raise ValueError.""" + cuid = str(cuid or "") + if not CUID_PATTERN.match(cuid): + raise ValueError(f"cuid must be a non-negative integer, got {cuid!r}") + return cuid + + +def validated_segment(value, field): + """Return `value` if it is a single safe path segment, else raise ValueError.""" + value = str(value or "") + if not PATH_SEGMENT_PATTERN.match(value): + raise ValueError(f"{field} must be a single path segment, got {value!r}") + return value + + +def _is_within(path, ancestor): + """True if `path` is `ancestor` or lies beneath it, compared segment-wise. + + A plain prefix test would accept a sibling whose name merely starts the same way + ("/var/lib/texera-mounts-evil" under "/var/lib/texera-mounts"). + """ + return os.path.commonpath([path, ancestor]) == ancestor + + +def _remove_empty_dirs(path, cuid): + """Remove `path` and any parents left empty, up to and including the CU's directory. + + Stops as soon as a directory is non-empty (another commit is still mounted under it) + and never walks above MOUNT_ROOT/<cuid>. Removing an empty MOUNT_ROOT/<cuid> is safe + for a running pod: its hostPath volume is DirectoryOrCreate, so the next mount + recreates it. + """ + cu_dir = os.path.normpath(os.path.join(MOUNT_ROOT, validated_cuid(cuid))) + path = os.path.normpath(path) + while _is_within(path, cu_dir): + try: + os.rmdir(path) + except OSError: + return # not empty (another commit is mounted here) or already gone + path = os.path.dirname(path) + + +def ensure_shared_root(): + """Make MOUNT_ROOT a shared mount so mounts created under it propagate to peers.""" + os.makedirs(MOUNT_ROOT, exist_ok=True) + if not os.path.ismount(MOUNT_ROOT): + subprocess.run(["mount", "--bind", MOUNT_ROOT, MOUNT_ROOT], check=False) + subprocess.run(["mount", "--make-rshared", MOUNT_ROOT], check=False) + + +def do_mount(cuid, repo, commit, jwt, file_service_base): + """Idempotently mount repo:commit for cuid. Returns the mount target path.""" + if not cuid or not repo or not commit or not jwt or not file_service_base: + raise ValueError("cuid, repositoryName, commitHash, jwt and fileServiceBase are required") + + cuid = validated_cuid(cuid) + repo = validated_segment(repo, "repositoryName") + commit = validated_segment(commit, "commitHash") + target = os.path.join(MOUNT_ROOT, cuid, repo, commit) + if is_mounted(target): + if _responds(target): + log(f"{repo}:{commit} already mounted for cu {cuid} at {target}") + return target + # The GeeseFS process backing this mount is gone — most likely the mounter was + # restarted, which kills every GeeseFS it started. The mount point survives but + # every access to it fails with ENOTCONN, so detach it and mount again. + log(f"{repo}:{commit} for cu {cuid} has a dead mount at {target}; remounting") + subprocess.run(["umount", "-l", target], capture_output=True, text=True) + + os.makedirs(target, exist_ok=True) + # allow_other: the mounter runs as root but the CU pod's UDF runs as a different + # (non-root) user, so the propagated FUSE mount must permit other users to access it. + cmd = [ + "geesefs", + "--endpoint", file_service_base, + "--memory-limit", "512", + "-o", "ro,allow_other", + f"{repo}:{commit}", + target, + ] + env = dict(os.environ) + env["AWS_ACCESS_KEY_ID"] = jwt + env["AWS_SECRET_ACCESS_KEY"] = MOUNT_SECRET_PLACEHOLDER + log(f"mounting {repo}:{commit} for cu {cuid} via: geesefs --endpoint {file_service_base} ... {target}") + result = subprocess.run(cmd, env=env, capture_output=True, text=True) + if result.returncode != 0: + # Nothing was mounted, so drop the directory just created for it rather than + # leaving an empty one behind until the next resync. A rejected mount (an + # unreadable repository, say) is a normal outcome, not a reason to litter. + _remove_empty_dirs(target, cuid) + raise RuntimeError( + f"geesefs exited {result.returncode}: {(result.stdout + result.stderr).strip()}" + ) + + # GeeseFS daemonizes after a successful mount; wait until the kernel reports it. + deadline = time.time() + MOUNT_TIMEOUT_S + while not is_mounted(target): + if time.time() > deadline: + raise RuntimeError(f"{repo}:{commit} did not appear as a mount at {target} in {MOUNT_TIMEOUT_S}s") + time.sleep(0.2) + log(f"mounted {repo}:{commit} for cu {cuid} at {target}") + return target + + +def list_mounts(cuid): + """List the datasets currently mounted for a CU as {repositoryName, commitHash, mountPath}. + + Derived entirely from /proc/mounts (via mount_targets_under) so the mounter stays + stateless: a mount at MOUNT_ROOT/<cuid>/<repo>/<commit> encodes its own identity. + """ + if not cuid: + raise ValueError("cuid is required") + cu_dir = os.path.normpath(os.path.join(MOUNT_ROOT, validated_cuid(cuid))) + mounts = [] + for target in mount_targets_under(cu_dir): + if os.path.normpath(target) == cu_dir: + continue # the shared-root bind itself, never a dataset mount + parts = os.path.relpath(target, cu_dir).split(os.sep) + if len(parts) != 2: + continue # not a <repo>/<commit> mount point + repo, commit = parts + mounts.append({"repositoryName": repo, "commitHash": commit, "mountPath": target}) + return mounts + + +def review_token(token): + """Ask the API server to validate `token` for this mounter's audience. + + Returns the TokenReview status. Submitting the audience is what makes the check + meaningful: without it any token the cluster issues -- including the one every pod gets + for the API server itself -- would come back authenticated. + """ + with _k8s_open( + "/apis/authentication.k8s.io/v1/tokenreviews", + timeout=10, + body={ + "apiVersion": "authentication.k8s.io/v1", + "kind": "TokenReview", + "spec": {"token": token, "audiences": [MOUNTER_AUDIENCE]}, + }, + ) as response: + return json.loads(response.read()).get("status", {}) + + +def authenticate_caller(auth_header): + """Return the caller's username, or raise Unauthorized. + + Three things have to hold, and all three are decided by the API server rather than here: + the token verifies, it was actually issued for this mounter's audience, and the identity + it belongs to is one this mounter serves. + """ + if not ALLOWED_CALLERS: + raise Unauthorized("mounter has no allowed callers configured") + + header = auth_header or "" + if not header.startswith("Bearer "): + raise Unauthorized("a bearer token is required") + token = header[len("Bearer "):].strip() + if not token: + raise Unauthorized("a bearer token is required") + + try: + status = review_token(token) + except Exception as e: # noqa: BLE001 + # Treated as a denial, not an outage: a mounter that cannot check who is calling + # must not fall back to serving the request. + raise Unauthorized(f"token review failed: {e}") + + if not status.get("authenticated"): + raise Unauthorized("token is not valid") + # The API server echoes back only the requested audiences the token actually carries, so + # an empty list here means the token was minted for something else. + if MOUNTER_AUDIENCE not in status.get("audiences", []): + raise Unauthorized(f"token is not scoped to audience {MOUNTER_AUDIENCE}") + username = status.get("user", {}).get("username", "") + if username not in ALLOWED_CALLERS: + raise Unauthorized(f"{username or 'caller'} is not permitted to request mounts") + return username + + +class Handler(BaseHTTPRequestHandler): + def _send(self, code, obj): + body = json.dumps(obj).encode() + self.send_response(code) + self.send_header("Content-Type", "application/json") + self.send_header("Content-Length", str(len(body))) + self.end_headers() + self.wfile.write(body) + + def _authenticate(self): + """True if the caller may use the mounter; otherwise answers 401 and returns False.""" + try: + authenticate_caller(self.headers.get("Authorization")) + return True + except Unauthorized as e: + log(f"rejected {self.command} {self.path}: {e}") + self._send(401, {"error": str(e)}) + return False + + def do_GET(self): + parsed = urlparse(self.path) + # The kubelet probes /healthz and holds no token for this audience, so readiness + # stays outside the authenticated surface. It reveals nothing about any mount. + if parsed.path == "/healthz": + self._send(200, {"status": "ok"}) + elif parsed.path == "/mounts": + if not self._authenticate(): + return + try: + cuid = parse_qs(parsed.query).get("cuid", [""])[0] + self._send(200, {"mounts": list_mounts(cuid)}) + except ValueError as e: + self._send(400, {"error": str(e)}) + except Exception as e: # noqa: BLE001 + log(f"listing mounts failed: {e}") + self._send(500, {"error": str(e)}) + else: + self._send(404, {"error": "not found"}) + + def do_POST(self): + if self.path != "/mount": + self._send(404, {"error": "not found"}) + return + if not self._authenticate(): + return + try: + length = int(self.headers.get("Content-Length", "0")) + req = json.loads(self.rfile.read(length) or b"{}") + target = do_mount( + str(req.get("cuid", "")), + str(req.get("repositoryName", "")), + str(req.get("commitHash", "")), + str(req.get("jwt", "")), + str(req.get("fileServiceBase", "")), + ) + self._send(200, {"mountPath": target}) + except ValueError as e: + self._send(400, {"error": str(e)}) + except Exception as e: # noqa: BLE001 + log(f"mount failed: {e}") + self._send(500, {"error": str(e)}) + + def log_message(self, fmt, *args): # silence default per-request stderr logging + pass + + +# ---- watcher: unmount a CU's directories once its pod is deleted ---- + +def _k8s_open(path, timeout, body=None): + """Call the in-cluster API server. Returns the response; the caller must close it. + + A `body` makes it a POST of that JSON object, which is how TokenReview is submitted. + """ + with open(os.path.join(SA_DIR, "token")) as f: + token = f.read().strip() + host = os.environ.get("KUBERNETES_SERVICE_HOST", "kubernetes.default.svc") + port = os.environ.get("KUBERNETES_SERVICE_PORT", "443") + ctx = ssl.create_default_context(cafile=os.path.join(SA_DIR, "ca.crt")) + headers = {"Authorization": f"Bearer {token}"} + data = None + if body is not None: + data = json.dumps(body).encode() + headers["Content-Type"] = "application/json" + req = urllib.request.Request(f"https://{host}:{port}{path}", data=data, headers=headers) + return urllib.request.urlopen(req, timeout=timeout, context=ctx) + + +def _cuid_of(pod_name): + """The cuid a CU pod name belongs to, or None if it is not a CU pod.""" + prefix = CU_POD_NAME_PREFIX + "-" + if not pod_name.startswith(prefix): + return None + return pod_name[len(prefix):] or None + + +def _list_cu_pods(): + """(live cuids, resourceVersion) for the pool namespace, or (None, None) if unreachable.""" + try: + with _k8s_open(f"/api/v1/namespaces/{POOL_NAMESPACE}/pods", timeout=30) as response: + pods = json.loads(response.read()) + except Exception as e: # noqa: BLE001 + log(f"listing CU pods failed: {e}") + return None, None + live = { + cuid + for cuid in (_cuid_of(p.get("metadata", {}).get("name", "")) for p in pods.get("items", [])) + if cuid + } + return live, pods.get("metadata", {}).get("resourceVersion") + + +def clean_cu_dir(cuid, quiet=False): + """Unmount everything under a departed CU's directory and remove it. + + Returns True once the directory is gone. `quiet` suppresses the messages for an orphan + already reported, so a repeatedly retried one does not re-log on every resync. + """ + cu_dir = os.path.join(MOUNT_ROOT, cuid) + # Never treat the shared root itself as a CU directory: unmounting it would break + # propagation for every CU on this node. + if os.path.normpath(cu_dir) == os.path.normpath(MOUNT_ROOT) or not os.path.isdir(cu_dir): + return True + if not quiet: + log(f"cu {cuid} pod is gone; unmounting {cu_dir}") + + for target in mount_targets_under(cu_dir): + # The mounter is root, so a plain lazy umount works (no setuid fusermount needed). + result = subprocess.run(["umount", "-l", target], capture_output=True, text=True) + if result.returncode != 0 and not quiet: + log(f"cu {cuid}: umount -l {target} failed: {(result.stdout + result.stderr).strip()}") + + # A lazy umount only detaches once the last reference to the mount goes away, so a mount + # another namespace still holds can outlive this call. Removing the directory would then + # fail, so only remove it once nothing is mounted underneath and let the next resync retry + # the rest — the mounts are read-only, so an orphan lingering a while is harmless. + remaining = mount_targets_under(cu_dir) + if remaining: + if not quiet: + log(f"cu {cuid}: {len(remaining)} mount(s) still busy, leaving {cu_dir} for the next resync") + return False + + shutil.rmtree(cu_dir, ignore_errors=True) + if os.path.exists(cu_dir): + if not quiet: + log(f"cu {cuid}: could not remove {cu_dir}, leaving it for the next resync") + return False + log(f"cu {cuid}: unmounted and removed {cu_dir}") + return True + + +# cuids whose cleanup could not finish, retried on the next resync and logged only once. +_pending = set() + + +def reconcile(live_cuids): + """Clean up every mount directory with no live CU pod behind it.""" + global _pending + if not os.path.isdir(MOUNT_ROOT): + return + still_pending = set() + for cuid in os.listdir(MOUNT_ROOT): + if cuid in live_cuids or not os.path.isdir(os.path.join(MOUNT_ROOT, cuid)): + continue + if not clean_cu_dir(cuid, quiet=cuid in _pending): + still_pending.add(cuid) + _pending = still_pending + + +def _handle_event(line): + try: + event = json.loads(line) + except ValueError: + return + if event.get("type") != "DELETED": + return + cuid = _cuid_of(event.get("object", {}).get("metadata", {}).get("name", "")) + if cuid and not clean_cu_dir(cuid, quiet=cuid in _pending): + _pending.add(cuid) + + +def watch_loop(): + """List-then-watch CU pods, unmounting a CU's mounts as soon as its pod is deleted. + + The initial LIST reconciles whatever was missed while the mounter was down; the WATCH + then reacts to deletions immediately. The API server closes the watch every + WATCH_TIMEOUT_S, and the resulting re-LIST doubles as the safety net for events missed + across a disconnect and for orphans whose unmount could not complete earlier. + """ + if not os.path.exists(os.path.join(SA_DIR, "token")): + log("no service account token; not running in-cluster, pod watch disabled") + return + while True: + live, resource_version = _list_cu_pods() + if live is None: # API unreachable → keep every mount (fail safe) and retry + time.sleep(WATCH_RETRY_S) + continue + reconcile(live) + try: + with _k8s_open( + f"/api/v1/namespaces/{POOL_NAMESPACE}/pods?watch=1" + f"&resourceVersion={resource_version}&timeoutSeconds={WATCH_TIMEOUT_S}", + timeout=WATCH_TIMEOUT_S + 30, + ) as response: + for line in response: + _handle_event(line) + except Exception as e: # noqa: BLE001 + log(f"pod watch failed ({e}); resyncing") + time.sleep(WATCH_RETRY_S) + + +def main(): + ensure_shared_root() + threading.Thread(target=watch_loop, daemon=True).start() + if not ALLOWED_CALLERS: + log("WARNING: MOUNTER_ALLOWED_CALLERS is empty; every mount request will be denied") + log( + f"listening on :{MOUNTER_PORT}, mount root {MOUNT_ROOT}, pool ns {POOL_NAMESPACE}, " + f"audience {MOUNTER_AUDIENCE}, callers {sorted(ALLOWED_CALLERS) or 'none'}" + ) + ThreadingHTTPServer(("0.0.0.0", MOUNTER_PORT), Handler).serve_forever() + + +if __name__ == "__main__": + main() diff --git a/bin/mounter/tests/conftest.py b/bin/mounter/tests/conftest.py new file mode 100644 index 0000000000..502ac8cb09 --- /dev/null +++ b/bin/mounter/tests/conftest.py @@ -0,0 +1,182 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +"""pytest fixtures for the texera-mounter tests. + +The mounter is the entrypoint of its own image rather than a Python package, so we +load it with `importlib.util` the same way `bin/local-dev/tests` loads the TUI. + +The module is loaded fresh per test because the tests mutate module-level state +(`MOUNT_ROOT`, `PROC_MOUNTS`, the pending-cleanup set), and it is pointed at a fake +mount root and a fake /proc/mounts under `tmp_path` so nothing touches the host. Test +helpers are attached to the module object so tests can stay terse: `mounter.logs`, +`mounter.runs`, `mounter.mounts()` and `mounter.set_mounts()`.""" + +from __future__ import annotations + +import importlib.util +import json +import os +import subprocess +import sys +import threading +import urllib.error +import urllib.request +from pathlib import Path + +import pytest + +REPO_ROOT = Path(__file__).resolve().parents[3] +MOUNTER_PATH = REPO_ROOT / "bin" / "mounter" / "mounter.py" + + [email protected] +def mounter(tmp_path, monkeypatch): + spec = importlib.util.spec_from_file_location("texera_mounter", MOUNTER_PATH) + assert spec is not None and spec.loader is not None + module = importlib.util.module_from_spec(spec) + sys.modules["texera_mounter"] = module + spec.loader.exec_module(module) + + module.MOUNT_ROOT = str(tmp_path / "mounts") + module.PROC_MOUNTS = str(tmp_path / "proc_mounts") + os.makedirs(module.MOUNT_ROOT) + + module.logs = [] + module.log = module.logs.append + + def set_mounts(*targets): + """Rewrite the fake /proc/mounts to hold exactly these FUSE mount targets.""" + with open(module.PROC_MOUNTS, "w") as mounts: + # One unrelated entry, so tests also prove the filtering works. + mounts.write("/dev/sda1 / ext4 rw,relatime 0 0\n") + for target in targets: + escaped = target.replace("\\", "\\134").replace(" ", "\\040") + mounts.write(f"dataset-1 {escaped} fuse.geesefs ro,nosuid,allow_other 0 0\n") + + def mounts(): + """The FUSE targets currently in the fake /proc/mounts.""" + return module.mount_targets_under(module.MOUNT_ROOT) + + module.set_mounts = set_mounts + module.mounts = mounts + set_mounts() + + # Every subprocess call is recorded rather than run. `geesefs` and `umount` also + # update the fake /proc/mounts, so the mount-table logic is exercised for real. + module.runs = [] + module.umount_succeeds = True + module.geesefs_returncode = 0 + # Set False to model a GeeseFS that exits 0 without the mount ever appearing. + module.geesefs_mounts = True + + def fake_run(cmd, **kwargs): + module.runs.append((list(cmd), kwargs)) + returncode = 0 + live = module.mount_targets_under(module.MOUNT_ROOT) + if cmd[0] == "geesefs": + returncode = module.geesefs_returncode + if returncode == 0 and module.geesefs_mounts: + set_mounts(*live, cmd[-1]) + elif cmd[0] == "umount": + if module.umount_succeeds: + set_mounts(*[t for t in live if t != os.path.normpath(cmd[-1])]) + else: + returncode = 32 + return subprocess.CompletedProcess(cmd, returncode, stdout="", stderr="umount: target is busy") + + monkeypatch.setattr(module.subprocess, "run", fake_run) + + # Caller authentication is stubbed at the TokenReview boundary: `review_token` is the + # only thing that talks to the API server, so replacing it exercises the whole of + # `authenticate_caller` -- the audience check, the allow-list and the failure modes -- + # without a cluster. `token_reviews` is what the API server would answer for a token. + module.ALLOWED_CALLER = "system:serviceaccount:texera:texera-access-control-service-account" + module.ALLOWED_CALLERS = frozenset({module.ALLOWED_CALLER}) + module.VALID_TOKEN = "a-token-for-access-control-service" + module.token_reviews = { + module.VALID_TOKEN: { + "authenticated": True, + "audiences": [module.MOUNTER_AUDIENCE], + "user": {"username": module.ALLOWED_CALLER}, + } + } + module.review_token_error = None + + def fake_review_token(token): + if module.review_token_error is not None: + raise module.review_token_error + # An unknown token is what the API server returns for one it cannot verify. + return module.token_reviews.get(token, {"authenticated": False}) + + module.review_token = fake_review_token + return module + + [email protected] +def http_client(mounter): + """A client for the mounter's real HTTP surface, served on an ephemeral port. + + The handler is exercised end to end -- routing, caller authentication, JSON parsing, + and the ValueError -> 400 mapping -- rather than by calling do_mount/list_mounts + directly, so the tests prove what a caller on the wire actually gets back from a + manipulated request. Requests carry the allowed caller's token by default; pass + `token=` to send a different one or `token=None` to send none.""" + server = mounter.ThreadingHTTPServer(("127.0.0.1", 0), mounter.Handler) + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + base = f"http://127.0.0.1:{server.server_address[1]}" + + _UNSET = object() + + def request(method, path, payload=None, token=_UNSET): + """Send a request, authenticated as the allowed caller unless `token` says otherwise. + + Pass `token=None` to send no Authorization header at all, or any other string to + present a different token.""" + body = None if payload is None else json.dumps(payload).encode() + headers = {"Content-Type": "application/json"} if body else {} + token = mounter.VALID_TOKEN if token is _UNSET else token + if token is not None: + headers["Authorization"] = f"Bearer {token}" + req = urllib.request.Request(base + path, data=body, method=method, headers=headers) + try: + with urllib.request.urlopen(req, timeout=10) as response: + return response.status, json.loads(response.read()) + except urllib.error.HTTPError as e: + # 4xx/5xx: the body carries the mounter's {"error": ...} payload. + return e.code, json.loads(e.read()) + + class _Client: + get = staticmethod(lambda path, **kw: request("GET", path, **kw)) + post = staticmethod(lambda path, payload, **kw: request("POST", path, payload, **kw)) + + yield _Client() + server.shutdown() + server.server_close() + + [email protected] +def cu_dir(mounter): + """Create a CU's mount directory tree and return (cuid, cu_dir, mount target).""" + + def make(cuid, repo="dataset-1", commit="abc123"): + target = os.path.join(mounter.MOUNT_ROOT, cuid, repo, commit) + os.makedirs(target, exist_ok=True) + return os.path.join(mounter.MOUNT_ROOT, cuid), target + + return make diff --git a/bin/mounter/tests/test_mounter.py b/bin/mounter/tests/test_mounter.py new file mode 100644 index 0000000000..9e0a5321b8 --- /dev/null +++ b/bin/mounter/tests/test_mounter.py @@ -0,0 +1,661 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +"""Unit tests for the texera-mounter (`bin/mounter/mounter.py`). + +Covers the logic that is easy to get wrong and expensive to debug in a cluster: how +mounts are discovered (from /proc/mounts, so a mount whose FUSE server died is still +seen), how a CU's mounts are torn down when its pod goes away, and how pod-watch +events are interpreted. The mounting itself is GeeseFS's job and is not re-tested +here — only the command we build for it and how we react to it failing.""" + +from __future__ import annotations + +import json +import os + +import pytest + + +# ─────────────────── mount_targets_under() / is_mounted() ─────────────────── + +def test_mount_targets_are_found_under_a_path_deepest_first(mounter, cu_dir): + _, target = cu_dir("7") + parent = os.path.dirname(target) + mounter.set_mounts(parent, target) + + found = mounter.mount_targets_under(os.path.join(mounter.MOUNT_ROOT, "7")) + + # Deepest first, so nested mounts are detached before the ones containing them. + assert found == [target, parent] + + +def test_mount_targets_exclude_unrelated_and_sibling_paths(mounter, cu_dir): + _, seven = cu_dir("7") + _, eight = cu_dir("8") + mounter.set_mounts(seven, eight, "/var/lib/something-else") + + assert mounter.mount_targets_under(os.path.join(mounter.MOUNT_ROOT, "7")) == [seven] + + +def test_mount_targets_do_not_match_a_path_that_is_only_a_string_prefix(mounter): + # "/mnt/7x" must not be treated as living under "/mnt/7". + seven = os.path.join(mounter.MOUNT_ROOT, "7") + mounter.set_mounts(seven + "x") + + assert mounter.mount_targets_under(seven) == [] + + +def test_mount_targets_decode_escaped_characters(mounter): + spaced = os.path.join(mounter.MOUNT_ROOT, "7", "data set", "abc") + mounter.set_mounts(spaced) + + assert mounter.mount_targets_under(mounter.MOUNT_ROOT) == [spaced] + + +def test_mount_targets_survive_an_unreadable_proc_mounts(mounter): + mounter.PROC_MOUNTS = os.path.join(mounter.MOUNT_ROOT, "does-not-exist") + + assert mounter.mount_targets_under(mounter.MOUNT_ROOT) == [] + assert any("failed" in line for line in mounter.logs) + + +def test_is_mounted_matches_the_path_itself_not_its_children(mounter, cu_dir): + _, target = cu_dir("7") + mounter.set_mounts(target) + + assert mounter.is_mounted(target) + assert not mounter.is_mounted(os.path.dirname(target)) + + +def test_a_dead_mount_is_still_reported_as_mounted(mounter, cu_dir): + """The regression this whole module exists for. + + Restarting the mounter kills every GeeseFS it started. The mount entry survives in + /proc/mounts but stat() on it fails, so os.path.ismount() — what this used to + use — reports False, and the mount was never cleaned up.""" + _, target = cu_dir("7") + mounter.set_mounts(target) + + assert not os.path.ismount(target) # what the buggy version asked + assert mounter.is_mounted(target) # what it should have asked + + +# ─────────────────── _cuid_of() ─────────────────── + [email protected]( + "pod_name, expected", + [ + ("computing-unit-17", "17"), + ("computing-unit-0", "0"), + ("computing-unit-", None), # no cuid at all + ("texera-file-service-abc", None), + ("", None), + ], +) +def test_cuid_is_parsed_only_from_computing_unit_pod_names(mounter, pod_name, expected): + assert mounter._cuid_of(pod_name) == expected + + +def test_the_pod_name_prefix_is_configurable(mounter, monkeypatch): + """The helm chart gives this and the CU manager the same prefix; honour it.""" + monkeypatch.setattr(mounter, "CU_POD_NAME_PREFIX", "other-deployment-cu") + + assert mounter._cuid_of("other-deployment-cu-17") == "17" + assert mounter._cuid_of("computing-unit-17") is None + + +# ─────────────────── do_mount() ─────────────────── + [email protected]("missing", ["cuid", "repo", "commit", "jwt", "base"]) +def test_do_mount_rejects_incomplete_requests(mounter, missing): + args = {"cuid": "7", "repo": "dataset-1", "commit": "abc", "jwt": "t", "base": "http://fs"} + args[missing] = "" + + with pytest.raises(ValueError): + mounter.do_mount(args["cuid"], args["repo"], args["commit"], args["jwt"], args["base"]) + + assert mounter.runs == [] + + +def test_do_mount_passes_the_jwt_as_the_s3_access_key(mounter): + mounter.do_mount("7", "dataset-1", "abc", "the-user-jwt", "http://file-service:9092") + + cmd, kwargs = mounter.runs[0] + assert cmd[0] == "geesefs" + assert "--endpoint" in cmd and "http://file-service:9092" in cmd + assert "dataset-1:abc" in cmd + # allow_other: the mounter is root but the UDF is not, so the propagated mount must + # be readable by another uid. ro: mounts are never writable. + assert "ro,allow_other" in cmd + # The pod's own JWT is the credential; no global LakeFS secret ever reaches the node. + assert kwargs["env"]["AWS_ACCESS_KEY_ID"] == "the-user-jwt" + + +def test_do_mount_is_idempotent_for_a_live_mount(mounter, cu_dir): + _, target = cu_dir("7", commit="abc") + mounter.set_mounts(target) + + assert mounter.do_mount("7", "dataset-1", "abc", "jwt", "http://fs") == target + assert mounter.runs == [] # no second GeeseFS for an already-mounted commit + + +def test_do_mount_replaces_a_dead_mount(mounter, cu_dir, monkeypatch): + _, target = cu_dir("7", commit="abc") + mounter.set_mounts(target) + monkeypatch.setattr(mounter, "_responds", lambda path: False) # FUSE server is gone + + assert mounter.do_mount("7", "dataset-1", "abc", "jwt", "http://fs") == target + + commands = [cmd[0] for cmd, _ in mounter.runs] + assert commands == ["umount", "geesefs"] # detached first, then mounted again + assert any("dead mount" in line for line in mounter.logs) + + +def test_do_mount_reports_what_geesefs_printed_when_it_fails(mounter): + mounter.geesefs_returncode = 1 + + with pytest.raises(RuntimeError, match="geesefs exited 1"): + mounter.do_mount("7", "dataset-1", "abc", "jwt", "http://fs") + + +def test_a_rejected_mount_leaves_no_empty_directory_behind(mounter): + """A mount the proxy refuses is routine; it should not litter until the next resync.""" + mounter.geesefs_returncode = 1 + + with pytest.raises(RuntimeError): + mounter.do_mount("7", "dataset-1", "abc", "jwt", "http://fs") + + assert not os.path.exists(os.path.join(mounter.MOUNT_ROOT, "7")) + + +def test_a_rejected_mount_keeps_a_sibling_commit_that_is_mounted(mounter, cu_dir): + _, live = cu_dir("7", commit="already-here") + mounter.set_mounts(live) + mounter.geesefs_returncode = 1 + + with pytest.raises(RuntimeError): + mounter.do_mount("7", "dataset-1", "new-commit", "jwt", "http://fs") + + assert not os.path.exists(os.path.join(mounter.MOUNT_ROOT, "7", "dataset-1", "new-commit")) + assert os.path.isdir(live) # the working mount is untouched + + +def test_do_mount_times_out_if_the_mount_never_appears(mounter, monkeypatch): + monkeypatch.setattr(mounter, "MOUNT_TIMEOUT_S", 0) + mounter.geesefs_mounts = False # exits 0 without the mount ever appearing + + with pytest.raises(RuntimeError, match="did not appear as a mount"): + mounter.do_mount("7", "dataset-1", "abc", "jwt", "http://fs") + + +# ─────────────────── clean_cu_dir() ─────────────────── + +def test_clean_unmounts_then_removes_the_directory(mounter, cu_dir): + directory, target = cu_dir("7") + mounter.set_mounts(target) + + assert mounter.clean_cu_dir("7") is True + assert [cmd for cmd, _ in mounter.runs] == [["umount", "-l", target]] + assert not os.path.exists(directory) + assert any("unmounted and removed" in line for line in mounter.logs) + + +def test_clean_unmounts_nested_mounts_deepest_first(mounter, cu_dir): + _, target = cu_dir("7") + parent = os.path.dirname(target) + mounter.set_mounts(parent, target) + + mounter.clean_cu_dir("7") + + assert [cmd[-1] for cmd, _ in mounter.runs] == [target, parent] + + +def test_clean_removes_a_directory_that_has_no_mounts(mounter, cu_dir): + directory, _ = cu_dir("7") + + assert mounter.clean_cu_dir("7") is True + assert mounter.runs == [] + assert not os.path.exists(directory) + + +def test_clean_keeps_a_directory_whose_mount_is_still_busy(mounter, cu_dir): + directory, target = cu_dir("7") + mounter.set_mounts(target) + mounter.umount_succeeds = False # lazy unmount has not completed yet + + assert mounter.clean_cu_dir("7") is False + # The directory must survive: removing it under a live mount is what the original + # code did, and it also made the removal fail on every cycle forever. + assert os.path.isdir(directory) + assert any("still busy" in line for line in mounter.logs) + + +def test_clean_is_silent_on_a_retry(mounter, cu_dir): + cu_dir("7") + mounter.set_mounts(os.path.join(mounter.MOUNT_ROOT, "7", "dataset-1", "abc123")) + mounter.umount_succeeds = False + + mounter.clean_cu_dir("7") + first_pass = len(mounter.logs) + mounter.clean_cu_dir("7", quiet=True) + + assert len(mounter.logs) == first_pass # a stuck orphan is reported once, not per cycle + + +def test_clean_never_touches_the_shared_root(mounter): + """A cuid of "" would resolve to MOUNT_ROOT, whose unmount would break every CU.""" + mounter.set_mounts(mounter.MOUNT_ROOT) + + assert mounter.clean_cu_dir("") is True + assert mounter.runs == [] + assert os.path.isdir(mounter.MOUNT_ROOT) + + +def test_clean_accepts_a_directory_that_is_already_gone(mounter): + assert mounter.clean_cu_dir("404") is True + + +# ─────────────────── reconcile() ─────────────────── + +def test_reconcile_removes_orphans_and_keeps_live_computing_units(mounter, cu_dir): + orphan, _ = cu_dir("7") + live, _ = cu_dir("8") + + mounter.reconcile({"8"}) + + assert not os.path.exists(orphan) + assert os.path.isdir(live) + + +def test_reconcile_retries_a_stuck_orphan_until_it_can_be_removed(mounter, cu_dir): + directory, target = cu_dir("7") + mounter.set_mounts(target) + mounter.umount_succeeds = False + + mounter.reconcile(set()) + assert mounter._pending == {"7"} + assert os.path.isdir(directory) + + mounter.umount_succeeds = True # the reference finally goes away + mounter.reconcile(set()) + + assert mounter._pending == set() + assert not os.path.exists(directory) + + +def test_reconcile_ignores_stray_files_in_the_mount_root(mounter): + stray = os.path.join(mounter.MOUNT_ROOT, "notes.txt") + open(stray, "w").close() + + mounter.reconcile(set()) + + assert os.path.exists(stray) + + +# ─────────────────── watch events ─────────────────── + +def _event(event_type, pod_name): + return json.dumps({"type": event_type, "object": {"metadata": {"name": pod_name}}}).encode() + + +def test_a_deleted_computing_unit_pod_is_cleaned_up_immediately(mounter, cu_dir): + directory, _ = cu_dir("7") + + mounter._handle_event(_event("DELETED", "computing-unit-7")) + + assert not os.path.exists(directory) + + [email protected]( + "line", + [ + _event("ADDED", "computing-unit-7"), + _event("MODIFIED", "computing-unit-7"), + _event("DELETED", "texera-file-service-xyz"), + _event("DELETED", "computing-unit-"), + b"{ not json\n", + b"\n", + ], + ids=["added", "modified", "other-pod", "no-cuid", "malformed", "blank"], +) +def test_irrelevant_watch_events_leave_every_mount_alone(mounter, cu_dir, line): + directory, _ = cu_dir("7") + + mounter._handle_event(line) + + assert os.path.isdir(directory) + assert mounter.runs == [] + + +def test_a_delete_that_cannot_finish_is_retried_on_the_next_resync(mounter, cu_dir): + directory, target = cu_dir("7") + mounter.set_mounts(target) + mounter.umount_succeeds = False + + mounter._handle_event(_event("DELETED", "computing-unit-7")) + + assert os.path.isdir(directory) + assert mounter._pending == {"7"} + + +# ─────────────────── _list_cu_pods() ─────────────────── + +def test_listing_returns_computing_unit_ids_and_the_resource_version(mounter, monkeypatch): + payload = { + "metadata": {"resourceVersion": "4242"}, + "items": [ + {"metadata": {"name": "computing-unit-7"}}, + {"metadata": {"name": "computing-unit-8"}}, + {"metadata": {"name": "some-other-pod"}}, + ], + } + monkeypatch.setattr(mounter, "_k8s_open", lambda path, timeout: _FakeResponse(payload)) + + assert mounter._list_cu_pods() == ({"7", "8"}, "4242") + + +def test_listing_failure_reports_no_answer_rather_than_an_empty_cluster(mounter, monkeypatch): + """A failed LIST must not read as "no CU pods exist" — that would unmount everything.""" + + def explode(path, timeout): + raise OSError("connection refused") + + monkeypatch.setattr(mounter, "_k8s_open", explode) + + live, resource_version = mounter._list_cu_pods() + + assert live is None and resource_version is None + assert any("listing CU pods failed" in line for line in mounter.logs) + + +# ─────────────────── request validation (path components) ─────────────────── +# +# Every component of MOUNT_ROOT/<cuid>/<repo>/<commit> arrives from the request. LakeFS +# only sees a request after geesefs runs, which is after the directory has been created, +# so the mounter cannot delegate path safety to it. These cases pin the escapes down. + [email protected]( + "cuid, escapes_to", + [ + ("5/../8", "8"), # traverses sideways into another CU's directory + ("../..", "outside"), # walks above the mount root entirely + ("/etc", "absolute"), # os.path.join drops MOUNT_ROOT for an absolute component + ("/", "absolute"), + ("8 ", "whitespace"), + ("-8", "negative"), + ("abc", "non-numeric"), + ("", "empty"), + ], +) +def test_do_mount_rejects_a_cuid_that_is_not_a_single_numeric_segment(mounter, cuid, escapes_to): + with pytest.raises(ValueError, match="cuid"): + mounter.do_mount(cuid, "dataset-1", "abc123", "jwt", "http://fs") + + assert mounter.runs == [] # geesefs never ran + # Nothing was created anywhere: not under another CU, not outside the mount root. + assert os.listdir(mounter.MOUNT_ROOT) == [], escapes_to + + +def test_a_manipulated_cuid_cannot_plant_a_mount_in_another_computing_units_directory(mounter): + """kunwp1's report: CU 5 asking for cuid='5/../8' must not land in CU 8's tree.""" + victim = os.path.join(mounter.MOUNT_ROOT, "8") + os.makedirs(victim) + + with pytest.raises(ValueError, match="cuid"): + mounter.do_mount("5/../8", "dataset-1", "abc123", "jwt", "http://fs") + + assert os.listdir(victim) == [] + + [email protected]("cuid", ["5/../8", "../..", "/etc", "abc"]) +def test_list_mounts_rejects_a_cuid_that_is_not_a_single_numeric_segment(mounter, cuid, cu_dir): + _, target = cu_dir("8") + mounter.set_mounts(target) + + # Enumerating another CU's mounts by traversal is refused rather than answered. + with pytest.raises(ValueError, match="cuid"): + mounter.list_mounts(cuid) + + [email protected]( + "repo", + ["../../8/dataset-1", "..", "a/b", "/etc", "-o", ".", ""], +) +def test_do_mount_rejects_a_repository_name_that_is_not_a_single_segment(mounter, repo): + with pytest.raises(ValueError, match="repositoryName|required"): + mounter.do_mount("7", repo, "abc123", "jwt", "http://fs") + + assert mounter.runs == [] + assert os.listdir(mounter.MOUNT_ROOT) == [] + + [email protected]("commit", ["../../..", "..", "a/b", "/etc", "-o", ".", ""]) +def test_do_mount_rejects_a_commit_hash_that_is_not_a_single_segment(mounter, commit): + with pytest.raises(ValueError, match="commitHash|required"): + mounter.do_mount("7", "dataset-1", commit, "jwt", "http://fs") + + assert mounter.runs == [] + assert os.listdir(mounter.MOUNT_ROOT) == [] + + +def test_legitimate_values_still_mount(mounter): + """The guard must not reject what the platform actually sends: dataset-<did> + a hex digest.""" + target = mounter.do_mount("42", "dataset-17", "0a1b2c3d4e5f", "jwt", "http://fs") + + assert target == os.path.join(mounter.MOUNT_ROOT, "42", "dataset-17", "0a1b2c3d4e5f") + assert [cmd[0] for cmd, _ in mounter.runs] == ["geesefs"] + + +_MOUNT_REQUEST = { + "cuid": "7", + "repositoryName": "dataset-1", + "commitHash": "abc123", + "jwt": "the-users-jwt", + "fileServiceBase": "http://file-service:9092", +} + + +# ─────────────────── caller authentication ─────────────────── +# The mounter runs privileged on every node and listens on a hostPort, so anything routable +# to a node IP -- including computing-unit pods running untrusted user code -- can open a +# connection to it. What separates the platform from user code is not reachability but a +# service-account token bound to the mounter's audience, verified by the API server. These +# tests stub only `review_token`, so everything `authenticate_caller` decides is real. + +def test_a_request_without_a_token_is_refused(mounter, http_client): + status, body = http_client.post("/mount", _MOUNT_REQUEST, token=None) + + assert status == 401 + assert "bearer token" in body["error"] + assert mounter.runs == [] + + [email protected]("header", ["", "Bearer ", "Bearer ", "the-token", "Basic dXNlcg=="]) +def test_a_malformed_authorization_header_is_refused(mounter, header): + with pytest.raises(mounter.Unauthorized): + mounter.authenticate_caller(header) + + +def test_a_token_the_api_server_cannot_verify_is_refused(mounter, http_client): + status, body = http_client.post("/mount", _MOUNT_REQUEST, token="forged") + + assert status == 401 + assert "not valid" in body["error"] + assert mounter.runs == [] + + +def test_a_token_minted_for_another_audience_is_refused(mounter): + """The pod's ordinary kube-apiserver token authenticates -- just not to the mounter.""" + mounter.token_reviews["kube-token"] = { + "authenticated": True, + "audiences": [], # the API server echoes back only the audiences the token carries + "user": {"username": mounter.ALLOWED_CALLER}, + } + + with pytest.raises(mounter.Unauthorized, match="audience"): + mounter.authenticate_caller("Bearer kube-token") + + +def test_a_valid_token_belonging_to_another_identity_is_refused(mounter): + """A CU pod can mint a token for this audience; it still is not an allowed caller.""" + mounter.token_reviews["cu-pod-token"] = { + "authenticated": True, + "audiences": [mounter.MOUNTER_AUDIENCE], + "user": {"username": "system:serviceaccount:texera-workflow-computing-unit-pool:default"}, + } + + with pytest.raises(mounter.Unauthorized, match="not permitted"): + mounter.authenticate_caller("Bearer cu-pod-token") + + +def test_an_unreachable_api_server_denies_rather_than_admits(mounter): + """Failing closed: a mounter that cannot tell who is calling must not mount.""" + mounter.review_token_error = RuntimeError("connection refused") + + with pytest.raises(mounter.Unauthorized, match="token review failed"): + mounter.authenticate_caller(f"Bearer {mounter.VALID_TOKEN}") + + +def test_no_configured_callers_denies_everything(mounter): + """An unset MOUNTER_ALLOWED_CALLERS is not an invitation to serve everyone.""" + mounter.ALLOWED_CALLERS = frozenset() + + with pytest.raises(mounter.Unauthorized, match="no allowed callers"): + mounter.authenticate_caller(f"Bearer {mounter.VALID_TOKEN}") + + +def test_the_allowed_caller_is_admitted(mounter): + assert mounter.authenticate_caller(f"Bearer {mounter.VALID_TOKEN}") == mounter.ALLOWED_CALLER + + +def test_authentication_does_not_replace_path_validation(mounter, http_client): + """A properly authenticated caller still cannot name a path outside the CU's directory.""" + status, body = http_client.post("/mount", {**_MOUNT_REQUEST, "cuid": "../.."}) + + assert status == 400 + assert "cuid" in body["error"] + assert mounter.runs == [] + + +def test_the_readiness_probe_stays_unauthenticated(mounter, http_client): + """The kubelet holds no token for this audience, and /healthz exposes no mount state.""" + assert http_client.get("/healthz", token=None) == (200, {"status": "ok"}) + + +# ─────────────────── HTTP surface ─────────────────── + +def test_a_manipulated_cuid_is_answered_with_400_over_http(mounter, http_client): + status, body = http_client.post( + "/mount", + { + "cuid": "../..", + "repositoryName": "dataset-1", + "commitHash": "abc123", + "jwt": "jwt", + "fileServiceBase": "http://fs", + }, + ) + + assert status == 400 + assert "cuid" in body["error"] + assert mounter.runs == [] + + +def test_a_manipulated_repository_name_is_answered_with_400_over_http(mounter, http_client): + status, body = http_client.post( + "/mount", + { + "cuid": "7", + "repositoryName": "../../8/dataset-1", + "commitHash": "abc123", + "jwt": "jwt", + "fileServiceBase": "http://fs", + }, + ) + + assert status == 400 + assert "repositoryName" in body["error"] + assert mounter.runs == [] + + +def test_listing_another_computing_unit_by_traversal_is_answered_with_400_over_http( + mounter, http_client, cu_dir +): + _, target = cu_dir("8") + mounter.set_mounts(target) + + assert http_client.get("/mounts?cuid=8")[1] == { + "mounts": [{"repositoryName": "dataset-1", "commitHash": "abc123", "mountPath": target}] + } + + status, body = http_client.get("/mounts?cuid=5/../8") + + assert status == 400 + assert "cuid" in body["error"] + + +def test_a_valid_mount_request_is_answered_with_the_mount_path_over_http(mounter, http_client): + status, body = http_client.post( + "/mount", + { + "cuid": "7", + "repositoryName": "dataset-1", + "commitHash": "abc123", + "jwt": "the-user-jwt", + "fileServiceBase": "http://fs", + }, + ) + + assert status == 200 + assert body == {"mountPath": os.path.join(mounter.MOUNT_ROOT, "7", "dataset-1", "abc123")} + + +# ─────────────────── _remove_empty_dirs() ─────────────────── + +def test_removing_empty_dirs_never_walks_above_the_computing_units_directory(mounter, cu_dir): + cu, target = cu_dir("7") + + mounter._remove_empty_dirs(target, "7") + + # The CU directory itself goes (DirectoryOrCreate recreates it), MOUNT_ROOT stays. + assert not os.path.exists(cu) + assert os.path.isdir(mounter.MOUNT_ROOT) + + +def test_removing_empty_dirs_does_not_treat_a_sibling_as_a_child(mounter): + """"/mounts/7x" merely starts with "/mounts/7"; a prefix test would delete it.""" + sibling = os.path.join(mounter.MOUNT_ROOT, "7x") + os.makedirs(sibling) + + mounter._remove_empty_dirs(sibling, "7") + + assert os.path.isdir(sibling) + + +class _FakeResponse: + def __init__(self, payload): + self._payload = json.dumps(payload).encode() + + def read(self): + return self._payload + + def __enter__(self): + return self + + def __exit__(self, *exc): + return False diff --git a/common/config/src/main/resources/kubernetes.conf b/common/config/src/main/resources/kubernetes.conf index 19c9796288..f423784334 100644 --- a/common/config/src/main/resources/kubernetes.conf +++ b/common/config/src/main/resources/kubernetes.conf @@ -29,6 +29,11 @@ kubernetes { compute-unit-service-name = "workflow-computing-unit-svc" compute-unit-service-name = ${?KUBERNETES_COMPUTE_UNIT_SERVICE_NAME} + # Computing-unit pods are named "<prefix>-<cuid>". The mounter derives a CU's id back + # out of the pod name, so both must be given the same value by the helm chart. + compute-unit-pod-name-prefix = "computing-unit" + compute-unit-pod-name-prefix = ${?KUBERNETES_COMPUTE_UNIT_POD_NAME_PREFIX} + image-name = "bobbai/texera-workflow-computing-unit:dev" image-name = ${?KUBERNETES_IMAGE_NAME} @@ -88,4 +93,15 @@ kubernetes { # falls back to the in-network address. jupyter-public-url-template = "" jupyter-public-url-template = ${?KUBERNETES_JUPYTER_PUBLIC_URL_TEMPLATE} + + # Per-node mounter (out-of-pod FUSE mounting). Off unless the deployment opted in: + # the mount gives each CU pod a hostPath volume, which the `baseline` and `restricted` + # Pod Security Standards forbid, so a cluster enforcing either would reject every CU pod. + mounter-enabled = false + mounter-enabled = ${?KUBERNETES_MOUNTER_ENABLED} + + # Must match the texera-mounter DaemonSet's hostPath in the helm chart: the CU pod's + # hostPath volume is the <root>/<cuid> subtree the mounter mounts into. + mounter-host-root = "/var/lib/texera-mounts" + mounter-host-root = ${?KUBERNETES_MOUNTER_HOST_ROOT} } \ No newline at end of file diff --git a/common/config/src/main/scala/org/apache/texera/common/config/EnvironmentalVariable.scala b/common/config/src/main/scala/org/apache/texera/common/config/EnvironmentalVariable.scala index 46bf207457..82c563b165 100644 --- a/common/config/src/main/scala/org/apache/texera/common/config/EnvironmentalVariable.scala +++ b/common/config/src/main/scala/org/apache/texera/common/config/EnvironmentalVariable.scala @@ -48,6 +48,21 @@ object EnvironmentalVariable { val ENV_USER_JWT_TOKEN = "USER_JWT_TOKEN" val ENV_AUTH_JWT_SECRET = "AUTH_JWT_SECRET" + /** + * Dataset-mount vars injected into the CU pod. The mount is performed by the per-node + * mounter and reaches the pod through mount propagation, so the pod only needs to know + * which computing unit it is and where the propagated mount appears. + * + * The mounter's own address is deliberately NOT among these. Only an authenticated + * platform caller may request a mount (see authenticate_caller in bin/mounter/mounter.py) + * and a pod running untrusted user code holds no token for that audience, so the address + * would serve no purpose here other than to probe the node's privileged mounter. A pod + * that wants a mount asks the platform, which authorizes the (user, cuid) pair before + * forwarding -- the cuid below is a claim to be checked, not a credential. + */ + val ENV_CU_ID = "TEXERA_CU_ID" + val ENV_MOUNT_IN_POD_ROOT = "TEXERA_MOUNT_IN_POD_ROOT" + // JDBC val ENV_JDBC_URL = "STORAGE_JDBC_URL" val ENV_JDBC_USERNAME = "STORAGE_JDBC_USERNAME" diff --git a/common/config/src/main/scala/org/apache/texera/common/config/KubernetesConfig.scala b/common/config/src/main/scala/org/apache/texera/common/config/KubernetesConfig.scala index 0812eae7ad..f52a14fa4e 100644 --- a/common/config/src/main/scala/org/apache/texera/common/config/KubernetesConfig.scala +++ b/common/config/src/main/scala/org/apache/texera/common/config/KubernetesConfig.scala @@ -30,6 +30,7 @@ object KubernetesConfig { val computeUnitServiceName: String = conf.getString("kubernetes.compute-unit-service-name") val computeUnitPoolName: String = conf.getString("kubernetes.compute-unit-pool-name") val computeUnitPoolNamespace: String = conf.getString("kubernetes.compute-unit-pool-namespace") + val computeUnitPodNamePrefix: String = conf.getString("kubernetes.compute-unit-pod-name-prefix") val computeUnitImageName: String = conf.getString("kubernetes.image-name") val computingUnitImagePullPolicy: String = conf.getString("kubernetes.image-pull-policy") @@ -78,4 +79,14 @@ object KubernetesConfig { // Browser-facing address with {uid} substituted; empty means use the in-network one. val jupyterPublicUrlTemplate: String = conf.getString("kubernetes.jupyter-public-url-template") + + // Whether the deployment opted into out-of-pod dataset mounting. When false the CU pod is + // built exactly as it was before the feature existed -- no hostPath, no mount env -- so a + // cluster enforcing a Pod Security Standard on the pool namespace is unaffected. + val mounterEnabled: Boolean = conf.getBoolean("kubernetes.mounter-enabled") + + // Root of the per-node mounter's host directory. This service never talks to the mounter + // -- access-control-service does -- but it builds the CU pod spec, and the pod's hostPath + // must be the <root>/<cuid> subtree the mounter mounts into. + val mounterHostRoot: String = conf.getString("kubernetes.mounter-host-root") } diff --git a/common/config/src/test/scala/org/apache/texera/common/config/KubernetesConfigSpec.scala b/common/config/src/test/scala/org/apache/texera/common/config/KubernetesConfigSpec.scala index 1c0d8da1b5..b52568594c 100644 --- a/common/config/src/test/scala/org/apache/texera/common/config/KubernetesConfigSpec.scala +++ b/common/config/src/test/scala/org/apache/texera/common/config/KubernetesConfigSpec.scala @@ -97,6 +97,15 @@ class KubernetesConfigSpec extends AnyFlatSpec with Matchers { ) } + "KubernetesConfig mounter settings" should "resolve to their kubernetes.conf defaults" in { + // Off by default: the mount gives each CU pod a hostPath volume, which the `baseline` + // and `restricted` Pod Security Standards forbid, so a deployment opts in. + ifUnset("KUBERNETES_MOUNTER_ENABLED")(KubernetesConfig.mounterEnabled shouldBe false) + ifUnset("KUBERNETES_MOUNTER_HOST_ROOT")( + KubernetesConfig.mounterHostRoot shouldBe "/var/lib/texera-mounts" + ) + } + "KubernetesConfig limit options" should "parse into trimmed, non-empty lists" in { ifUnset("KUBERNETES_COMPUTING_UNIT_CPU_LIMIT_OPTIONS")( KubernetesConfig.cpuLimitOptions shouldBe List("1", "2", "4") diff --git a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/KubernetesClient.scala b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/KubernetesClient.scala index 6f97a5cf6b..18328f1f9f 100644 --- a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/KubernetesClient.scala +++ b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/KubernetesClient.scala @@ -22,7 +22,7 @@ package org.apache.texera.service.util import io.fabric8.kubernetes.api.model._ import io.fabric8.kubernetes.api.model.metrics.v1beta1.PodMetrics import io.fabric8.kubernetes.client.KubernetesClientBuilder -import org.apache.texera.common.config.KubernetesConfig +import org.apache.texera.common.config.{EnvironmentalVariable, KubernetesConfig} import scala.jdk.CollectionConverters._ @@ -32,10 +32,16 @@ import scala.jdk.CollectionConverters._ * parameter (not a mutable global) so tests can construct an instance backed by a stubbed * client and exercise the passthrough wrappers without a live cluster. */ -class KubernetesClient(client: io.fabric8.kubernetes.client.KubernetesClient) { +class KubernetesClient( + client: io.fabric8.kubernetes.client.KubernetesClient, + // A constructor parameter rather than a direct KubernetesConfig read, so the spec can + // build a pod both ways: the mount contract below is security- and scheduling-sensitive + // and needs asserting on, but the default must stay the PodSecurity-safe one. + mountingEnabled: Boolean = KubernetesConfig.mounterEnabled +) { private val namespace: String = KubernetesConfig.computeUnitPoolNamespace - private val podNamePrefix = "computing-unit" + private val podNamePrefix = KubernetesConfig.computeUnitPodNamePrefix def generatePodURI(cuid: Int): String = { s"${generatePodName(cuid)}.${KubernetesConfig.computeUnitServiceName}.$namespace.svc.cluster.local:${KubernetesConfig.computeUnitPortNumber}" @@ -127,16 +133,31 @@ class KubernetesClient(client: io.fabric8.kubernetes.client.KubernetesClient) { throw new Exception(s"Pod with cuid $cuid already exists") } - val envList = envVars - .map { - case (key, value) => + val baseEnv = envVars.map { + case (key, value) => + new EnvVarBuilder().withName(key).withValue(value.toString).build() + }.toList + + // Which CU this is, and where its propagated mounts show up. The pod is deliberately + // not given the mounter's address: only an authenticated platform caller may request a + // mount, so the address would be of no use to code running here except to probe the + // node's privileged mounter. + val inPodMountRoot = "/mnt/texera-mounts" + val mounterEnv = + if (!mountingEnabled) Nil + else + List( new EnvVarBuilder() - .withName(key) - .withValue(value.toString) + .withName(EnvironmentalVariable.ENV_CU_ID) + .withValue(cuid.toString) + .build(), + new EnvVarBuilder() + .withName(EnvironmentalVariable.ENV_MOUNT_IN_POD_ROOT) + .withValue(inPodMountRoot) .build() - } - .toList - .asJava + ) + + val envList = (baseEnv ++ mounterEnv).asJava // Setup the resource requirements val resourceBuilder = new ResourceRequirementsBuilder() @@ -179,6 +200,18 @@ class KubernetesClient(client: io.fabric8.kubernetes.client.KubernetesClient) { .withEnv(envList) .withResources(resourceBuilder.build()) + // The FUSE mount is performed by the per-node texera-mounter (privileged), not here, + // so this pod stays UNPRIVILEGED. It only *receives* the mount via HostToContainer + // propagation from a host directory scoped to this CU id (see the hostPath volume below). + if (mountingEnabled) { + containerBuilder + .addNewVolumeMount() + .withName("texera-mounts") + .withMountPath(inPodMountRoot) + .withMountPropagation("HostToContainer") + .endVolumeMount() + } + // If shmSize requested, mount /dev/shm shmSize.foreach { _ => containerBuilder @@ -204,6 +237,21 @@ class KubernetesClient(client: io.fabric8.kubernetes.client.KubernetesClient) { .endVolume() } + // Per-CU host directory the mounter mounts into (DirectoryOrCreate so it exists + // before the mounter mounts). Scoped by cuid so a CU can only ever see its own mounts. + // Guarded because `baseline` and `restricted` forbid hostPath: on a cluster enforcing + // either on the pool namespace, an unconditional one makes every CU pod unschedulable. + if (mountingEnabled) { + specBuilder + .addNewVolume() + .withName("texera-mounts") + .withNewHostPath() + .withPath(s"${KubernetesConfig.mounterHostRoot}/$cuid") + .withType("DirectoryOrCreate") + .endHostPath() + .endVolume() + } + val pod = specBuilder .withHostname(podName) .withSubdomain(KubernetesConfig.computeUnitServiceName) @@ -219,4 +267,10 @@ class KubernetesClient(client: io.fabric8.kubernetes.client.KubernetesClient) { } /** Production singleton bound to a real in-cluster fabric8 client. */ -object KubernetesClient extends KubernetesClient(new KubernetesClientBuilder().build()) +object KubernetesClient + extends KubernetesClient( + new KubernetesClientBuilder().build(), + // Passed explicitly: a companion object extending its companion class may not rely on + // the class's default constructor arguments. + KubernetesConfig.mounterEnabled + ) diff --git a/computing-unit-managing-service/src/test/scala/org/apache/texera/service/util/KubernetesClientSpec.scala b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/util/KubernetesClientSpec.scala index f8e8353f7e..5e6743c953 100644 --- a/computing-unit-managing-service/src/test/scala/org/apache/texera/service/util/KubernetesClientSpec.scala +++ b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/util/KubernetesClientSpec.scala @@ -27,6 +27,7 @@ import io.fabric8.kubernetes.api.model.metrics.v1beta1.{ PodMetricsListBuilder } import io.fabric8.kubernetes.api.model.{ + Container, ContainerBuilder, Pod, PodBuilder, @@ -66,6 +67,13 @@ class KubernetesClientSpec extends AnyFlatSpec with Matchers { private val namespace: String = KubernetesConfig.computeUnitPoolNamespace + // getVolumes/getVolumeMounts are null rather than empty when nothing was added. + private def volumeNames(pod: Pod): List[String] = + Option(pod.getSpec.getVolumes).map(_.asScala.toList).getOrElse(Nil).map(_.getName) + + private def mountNames(container: Container): List[String] = + Option(container.getVolumeMounts).map(_.asScala.toList).getOrElse(Nil).map(_.getName) + // A fabric8 client stubbed just enough to answer the namespace-wide pod-list and pod-metrics // calls the wrappers make. RETURNS_DEEP_STUBS can't be used: fabric8's fluent API returns type // variables, so each step of the chain is mocked explicitly. @@ -278,6 +286,53 @@ class KubernetesClientSpec extends AnyFlatSpec with Matchers { // Env values arrive as Any and reach the container as strings. container.getEnv.asScala.map(e => e.getName -> e.getValue).toMap shouldBe Map("UID" -> "9", "MODE" -> "batch") + + // The PodSecurity-safe default: no hostPath, so a cluster enforcing `baseline` or + // `restricted` on the pool namespace still admits the pod. + volumeNames(built) should not contain "texera-mounts" + mountNames(container) should not contain "texera-mounts" + } + + it should "wire the mount contract into the pod only when mounting is enabled" in { + val name = KubernetesClient.generatePodName(5) + val (client, _) = clientWithNamedPod(name, null) + val namespaceable = mock(classOf[NamespaceableResource[Pod]]) + val resource = mock(classOf[Resource[Pod]]) + val captor = ArgumentCaptor.forClass(classOf[Pod]) + when(client.resource(any(classOf[Pod]))).thenReturn(namespaceable) + when(namespaceable.inNamespace(namespace)).thenReturn(resource) + when(resource.create()).thenReturn(null) + + new KubernetesClient(client, mountingEnabled = true) + .createPod(5, "2", "4Gi", "1", Map("UID" -> 9, "MODE" -> "batch")) + + verify(client).resource(captor.capture()) + val built = captor.getValue + val container = built.getSpec.getContainers.asScala.head + + // The host directory is scoped to this CU, so one unit can never see another's mounts. + val volume = built.getSpec.getVolumes.asScala.find(_.getName == "texera-mounts") + volume shouldBe defined + volume.get.getHostPath.getPath shouldBe s"${KubernetesConfig.mounterHostRoot}/5" + // DirectoryOrCreate: the mounter mounts into a path the kubelet has already created. + volume.get.getHostPath.getType shouldBe "DirectoryOrCreate" + + // HostToContainer, not Bidirectional: the pod receives the mount the privileged mounter + // makes and must not be able to propagate one back out to the node. + val mount = container.getVolumeMounts.asScala.find(_.getName == "texera-mounts") + mount shouldBe defined + mount.get.getMountPath shouldBe "/mnt/texera-mounts" + mount.get.getMountPropagation shouldBe "HostToContainer" + + // The pod is told which CU it is and where its mounts appear -- and deliberately not + // how to reach the mounter. + val env = container.getEnv.asScala.map(e => e.getName -> e.getValue).toMap + env should contain("TEXERA_CU_ID" -> "5") + env should contain("TEXERA_MOUNT_IN_POD_ROOT" -> "/mnt/texera-mounts") + + // Still unprivileged: the whole point of mounting out of pod. + val privileged = Option(container.getSecurityContext).flatMap(c => Option(c.getPrivileged)) + privileged.contains(true) shouldBe false } it should "mount a shared-memory volume only when a size is asked for" in { diff --git a/file-service/src/main/scala/org/apache/texera/service/FileService.scala b/file-service/src/main/scala/org/apache/texera/service/FileService.scala index 9ed86b9207..63cee25a7a 100644 --- a/file-service/src/main/scala/org/apache/texera/service/FileService.scala +++ b/file-service/src/main/scala/org/apache/texera/service/FileService.scala @@ -40,6 +40,7 @@ import org.apache.texera.service.resource.{ ModelResource } import org.apache.texera.service.util.S3StorageClient +import org.apache.texera.service.util.S3ProxyServlet import org.apache.texera.service.util.LargeBinaryManager import org.apache.texera.service.util.StagedFileCleanupJob import org.eclipse.jetty.server.session.SessionHandler @@ -95,6 +96,12 @@ class FileService extends Application[FileServiceConfiguration] with LazyLogging environment.jersey.register(classOf[ModelResource]) environment.jersey.register(classOf[ModelAccessResource]) + // Register the read-only S3 proxy servlet for in-pod GeeseFS dataset mounts. GeeseFS + // authenticates with the pod's per-user JWT (carried as its S3 access key) and issues + // path-style requests at the root (/<bucket>/<key>), while Jersey serves the REST API + // at /api/* (more specific, so it keeps taking precedence). + environment.servlets.addServlet("s3-mount-proxy", new S3ProxyServlet).addMapping("/*") + RoleAnnotationEnforcer.enforce(environment.jersey.getResourceConfig, "FileService") // Route request logs through SLF4J, controlled by TEXERA_SERVICE_LOG_LEVEL diff --git a/file-service/src/main/scala/org/apache/texera/service/util/S3ProxyServlet.scala b/file-service/src/main/scala/org/apache/texera/service/util/S3ProxyServlet.scala new file mode 100644 index 0000000000..81685ca507 --- /dev/null +++ b/file-service/src/main/scala/org/apache/texera/service/util/S3ProxyServlet.scala @@ -0,0 +1,273 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.texera.service.util + +import com.typesafe.scalalogging.LazyLogging +import jakarta.servlet.http.{HttpServlet, HttpServletRequest, HttpServletResponse} +import org.apache.texera.auth.JwtParser +import org.apache.texera.common.config.StorageConfig +import org.apache.texera.dao.SqlServer +import org.apache.texera.dao.SqlServer.withTransaction +import org.apache.texera.dao.jooq.generated.tables.daos.{DatasetDao, ModelDao} +import org.apache.texera.service.resource.{DatasetAccessResource, ModelAccessResource} +import software.amazon.awssdk.auth.credentials.AwsBasicCredentials +import software.amazon.awssdk.auth.signer.AwsS3V4Signer +import software.amazon.awssdk.auth.signer.params.AwsS3V4SignerParams +import software.amazon.awssdk.http.{SdkHttpFullRequest, SdkHttpMethod} +import software.amazon.awssdk.regions.Region + +import java.net.{URI, URLDecoder} +import java.net.http.{HttpClient, HttpRequest, HttpResponse} +import java.time.Duration +import scala.jdk.CollectionConverters._ +import scala.jdk.OptionConverters._ + +/** + * Read-only, JWT-authenticated, re-signing reverse proxy in front of the LakeFS S3 + * gateway. A computing-unit pod's GeeseFS mount talks to this servlet using the pod's + * own per-user JWT as the S3 credential: the JWT is passed to GeeseFS as + * `AWS_ACCESS_KEY_ID`, so it rides in the request's SigV4/SigV2 `Authorization` header. + * Reusing the JWT that is already present in the pod means no separate mount credential + * is ever issued, stored, or made multi-replica-consistent. The servlet: + * + * 1. reads the JWT back out of the incoming `Authorization` header (the JWT is the + * bearer capability; the pod-side S3 signature is not re-validated, and no LakeFS + * credentials ever leave this service), + * 2. verifies the JWT and checks that its user has read access to the requested + * repository (the S3 bucket), using the same `userHasReadAccess` gate as the + * dataset REST endpoints, and + * 3. re-signs the request with the global LakeFS credentials and forwards it to the + * LakeFS S3 gateway, streaming the response back. + * + * Because requests are forwarded verbatim, the proxy behaves identically to a direct + * GeeseFS -> LakeFS-gateway mount, just re-authenticated. GeeseFS mounts read-only, so + * only GET and HEAD are handled. + */ +class S3ProxyServlet extends HttpServlet with LazyLogging { + + // The LakeFS S3 gateway shares the LakeFS server address: the configured API endpoint + // with the trailing /api/v1 suffix removed. + private val gatewayEndpoint: URI = + URI.create(StorageConfig.lakefsEndpoint.stripSuffix("/").stripSuffix("/api/v1")) + + private val lakefsCredentials = + AwsBasicCredentials.create(StorageConfig.lakefsUsername, StorageConfig.lakefsPassword) + + // The S3-specific SigV4 signer adds and signs the x-amz-content-sha256 header (which + // S3 / the LakeFS gateway require) and disables path double-encoding — both essential + // for the signature to validate. The generic Aws4Signer omits x-amz-content-sha256. + private val signer = AwsS3V4Signer.create() + + private val httpClient: HttpClient = HttpClient + .newBuilder() + .followRedirects(HttpClient.Redirect.NEVER) + .connectTimeout(Duration.ofSeconds(10)) + .build() + + private val forwardedResponseHeaderPrefixes = + Seq("content-", "etag", "last-modified", "accept-ranges", "x-amz-") + + override def doGet(req: HttpServletRequest, resp: HttpServletResponse): Unit = + proxy(req, resp, SdkHttpMethod.GET, streamBody = true) + + override def doHead(req: HttpServletRequest, resp: HttpServletResponse): Unit = + proxy(req, resp, SdkHttpMethod.HEAD, streamBody = false) + + private def proxy( + req: HttpServletRequest, + resp: HttpServletResponse, + method: SdkHttpMethod, + streamBody: Boolean + ): Unit = { + val user = S3ProxyServlet + .extractCredentialToken(req.getHeader("Authorization")) + .flatMap(token => JwtParser.parseToken(token).toScala) + if (user.isEmpty) { + // GeeseFS probes the bucket unauthenticated on mount, so this is expected noise. + resp.sendError(HttpServletResponse.SC_FORBIDDEN, "missing or invalid user token") + return + } + + val uid = user.get.getUid + val repositoryName = S3ProxyServlet.bucketFromUri(req.getRequestURI) + if (repositoryName.isEmpty || !authorizedToRead(uid, repositoryName)) { + logger.warn( + s"user $uid denied mount access to repository '$repositoryName' for ${req.getRequestURI}" + ) + resp.sendError(HttpServletResponse.SC_FORBIDDEN, "no read access to the requested repository") + return + } + + try { + writeResponse(forward(req, method), resp, streamBody) + } catch { + case e: Exception => + logger.error(s"error proxying ${req.getRequestURI} to LakeFS gateway", e) + resp.sendError(HttpServletResponse.SC_BAD_GATEWAY, "upstream error") + } + } + + /** + * True iff `uid` has read access to the versioned resource backing `repositoryName`. + * + * Read access to a repository grants read to all of its commits, so no per-commit check + * is needed: a session addresses a single repository's data and any version the user may + * already read. + * + * Both resource types are searched, because a repository backs either a dataset or a + * model, and the name is matched rather than parsed: rows created today are named + * `<type>-<id>`, but `sql/updates/15.sql` backfilled the column from the dataset's plain + * `name`, so an upgraded deployment still has repositories called e.g. `my-data`. + * + * Exactly one row must match, or the request is denied. A repository name is unique in + * practice -- LakeFS will not create two repositories with one name, and dataset creation + * checked the name globally besides -- but `repository_name` carries no unique constraint + * (the only UNIQUE on either table is (owner_uid, name)), so nothing in the schema says + * so. Rather than let `.head` pick a row, and decide access arbitrarily in front of a + * proxy that re-signs with the global LakeFS credentials, an ambiguous name fails closed. + */ + private def authorizedToRead(uid: Integer, repositoryName: String): Boolean = + withTransaction(SqlServer.getInstance().createDSLContext()) { ctx => + val datasets = new DatasetDao(ctx.configuration()) + .fetchByRepositoryName(repositoryName) + .asScala + .toList + val models = new ModelDao(ctx.configuration()) + .fetchByRepositoryName(repositoryName) + .asScala + .toList + (datasets, models) match { + case (dataset :: Nil, Nil) => + DatasetAccessResource.userHasReadAccess(ctx, dataset.getDid, uid) + case (Nil, model :: Nil) => + ModelAccessResource.userHasReadAccess(ctx, model.getMid, uid) + case (Nil, Nil) => + false + case _ => + // Logged because it means the data violates an assumption the mount path rests on, + // and the owners of those rows will see an unexplained denial. + logger.error( + s"repository '$repositoryName' matches ${datasets.size} datasets and " + + s"${models.size} models; denying rather than choosing one" + ) + false + } + } + + /** Re-sign the request with LakeFS credentials and send it to the LakeFS gateway. */ + private def forward( + req: HttpServletRequest, + method: SdkHttpMethod + ): HttpResponse[java.io.InputStream] = { + val query = Option(req.getQueryString).map("?" + _).getOrElse("") + val targetUri = URI.create( + gatewayEndpoint.getScheme + "://" + gatewayEndpoint.getAuthority + req.getRequestURI + query + ) + + val builder = SdkHttpFullRequest + .builder() + .method(method) + .uri(targetUri) + // Range must be part of the signed request so the gateway accepts it. + Option(req.getHeader("Range")).foreach(r => builder.putHeader("Range", r)) + + val signed = signer.sign( + builder.build(), + AwsS3V4SignerParams + .builder() + .awsCredentials(lakefsCredentials) + .signingName("s3") + .signingRegion(Region.US_EAST_1) + .build() + ) + + var outbound = HttpRequest + .newBuilder() + .uri(targetUri) + .method(method.name(), HttpRequest.BodyPublishers.noBody()) + // Forward every signed header except Host, which the HTTP client sets itself to the + // same value we signed (the gateway authority). + signed.headers().asScala.foreach { + case (name, values) => + if (!name.equalsIgnoreCase("Host")) { + values.asScala.foreach(v => outbound = outbound.header(name, v)) + } + } + + httpClient.send(outbound.build(), HttpResponse.BodyHandlers.ofInputStream()) + } + + private def writeResponse( + upstream: HttpResponse[java.io.InputStream], + resp: HttpServletResponse, + streamBody: Boolean + ): Unit = { + resp.setStatus(upstream.statusCode()) + upstream.headers().map().asScala.foreach { + case (name, values) => + val lower = name.toLowerCase + if (forwardedResponseHeaderPrefixes.exists(lower.startsWith)) { + values.asScala.foreach(v => resp.addHeader(name, v)) + } + } + + if (streamBody) { + val in = upstream.body() + try { + in.transferTo(resp.getOutputStream) + } finally { + in.close() + } + } else { + upstream.body().close() + } + } +} + +object S3ProxyServlet { + + /** + * Extract the credential token from an AWS `Authorization` header. GeeseFS carries the + * user JWT in the access-key-id position and signs with either SigV4 + * (`AWS4-HMAC-SHA256 Credential=<token>/<date>/...`) or, against a plain-HTTP custom + * endpoint, SigV2 (`AWS <token>:<signature>`); support both. A JWT is base64url with + * `.` separators, so it never contains the `/` or `:` these formats delimit on. The + * pod-side S3 signature itself is not re-validated — the JWT is the bearer capability — + * so only the token needs to be read out. + */ + private[util] def extractCredentialToken(authHeader: String): Option[String] = { + Option(authHeader).flatMap { h => + "Credential=([^/,\\s]+)/".r + .findFirstMatchIn(h) + .map(_.group(1)) // SigV4 + .orElse("^AWS ([^:\\s]+):".r.findFirstMatchIn(h.trim).map(_.group(1))) // SigV2 + } + } + + /** + * The repository (S3 bucket) a path-style request URI `/<bucket>/<key>` targets: its + * first path segment, URL-decoded. Empty when the URI carries no bucket (root, or a + * service-level list-buckets), which is never authorized. + */ + private[util] def bucketFromUri(requestUri: String): String = { + val firstSegment = requestUri.stripPrefix("/").split("/", 2)(0) + if (firstSegment.isEmpty) "" else URLDecoder.decode(firstSegment, "UTF-8") + } +} diff --git a/file-service/src/test/scala/org/apache/texera/service/util/S3ProxyServletSpec.scala b/file-service/src/test/scala/org/apache/texera/service/util/S3ProxyServletSpec.scala new file mode 100644 index 0000000000..b7024778a6 --- /dev/null +++ b/file-service/src/test/scala/org/apache/texera/service/util/S3ProxyServletSpec.scala @@ -0,0 +1,86 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.texera.service.util + +import org.scalatest.flatspec.AnyFlatSpec +import org.scalatest.matchers.should.Matchers + +/** + * Unit tests for the pure request-parsing helpers of [[S3ProxyServlet]] — the credential + * token extraction (which must handle both the SigV4 and SigV2 `Authorization` formats + * GeeseFS emits, carrying the user JWT in the access-key position) and the path-style + * bucket/repository extraction that scopes each request. + */ +class S3ProxyServletSpec extends AnyFlatSpec with Matchers { + + // A representative JWT: base64url segments (A-Za-z0-9-_) joined by '.', so it contains + // none of the '/', ':' or whitespace that the two Authorization formats delimit on. + private val jwt = + "eyJhbGciOiJIUzI1NiJ9.eyJzdWIiOiJ0ZXhlcmEiLCJ1c2VySWQiOjF9.abc-_DEF123" + + "extractCredentialToken" should "read the access key from a SigV4 Authorization header" in { + val header = + s"AWS4-HMAC-SHA256 Credential=$jwt/20260721/us-east-1/s3/aws4_request, " + + "SignedHeaders=host;x-amz-content-sha256;x-amz-date, Signature=deadbeef" + S3ProxyServlet.extractCredentialToken(header) shouldBe Some(jwt) + } + + it should "read the access key from a SigV2 Authorization header" in { + S3ProxyServlet.extractCredentialToken(s"AWS $jwt:c2lnbmF0dXJl") shouldBe Some(jwt) + } + + it should "handle a short (non-JWT) access key in both formats" in { + S3ProxyServlet.extractCredentialToken( + "AWS4-HMAC-SHA256 Credential=AKIAEXAMPLE/20260721/us-east-1/s3/aws4_request, " + + "SignedHeaders=host, Signature=abc" + ) shouldBe Some("AKIAEXAMPLE") + S3ProxyServlet.extractCredentialToken("AWS AKIAEXAMPLE:sig") shouldBe Some("AKIAEXAMPLE") + } + + it should "return None for a null header" in { + S3ProxyServlet.extractCredentialToken(null) shouldBe None + } + + it should "return None for a header in neither AWS format" in { + S3ProxyServlet.extractCredentialToken("Bearer some.jwt.token") shouldBe None + S3ProxyServlet.extractCredentialToken("") shouldBe None + S3ProxyServlet.extractCredentialToken("AWS4-HMAC-SHA256 SignedHeaders=host") shouldBe None + } + + "bucketFromUri" should "return the first path segment of a path-style object request" in { + S3ProxyServlet.bucketFromUri( + "/dataset-1/097b4111e0ac9f46/model-00001-of-00003.pt" + ) shouldBe "dataset-1" + } + + it should "return the bucket for a bucket-only request (with or without trailing slash)" in { + S3ProxyServlet.bucketFromUri("/dataset-1") shouldBe "dataset-1" + S3ProxyServlet.bucketFromUri("/dataset-1/") shouldBe "dataset-1" + } + + it should "URL-decode the bucket segment" in { + S3ProxyServlet.bucketFromUri("/my%20dataset/commit/f.txt") shouldBe "my dataset" + } + + it should "return empty for the root URI (no bucket to authorize)" in { + S3ProxyServlet.bucketFromUri("/") shouldBe "" + S3ProxyServlet.bucketFromUri("") shouldBe "" + } +}
