This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/texera.git
The following commit(s) were added to refs/heads/main by this push:
new 209fc4152e feat(computing-unit): out-of-pod LakeFS repository mount
infrastructure (#6866)
209fc4152e is described below
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 ""
+ }
+}