This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/main/pr-6896-12aeeab5816b12880c5f58970460adf87b548d89 in repository https://gitbox.apache.org/repos/asf/texera.git
commit dcec7f3601e61436857a619365f7a2c5963c7589 Author: ali risheh <[email protected]> AuthorDate: Tue Sep 15 02:48:54 2026 +0000 feat(LakeFS-mount): authorize and perform a repository mount (#6896) ### What changes were proposed in this PR? Let a computing unit have a versioned LakeFS repository mounted into it, on the infrastructure merged in #6866. This is the authorization step of #6606: the endpoint that decides whether a mount may happen, and performs it. The per-node mounter is the one privileged component — root on every node, listening on a hostPort — so it admits exactly one caller, verified with `TokenReview` against a token bound to the `texera-mounter` audience. Only access-control-service holds one, which is also why it is the right place: it is already where the deployment decides whether a user may act on a computing unit. **The mounter authorizes nothing.** It performs what it is told, so every decision has to be made before it is asked. The endpoint makes four: 1. **The request has a shape a mount path can be built from** — each component a single safe segment, so a separator, a `..` or a leading `-` is refused by name rather than by whatever it fails next. The mounter checks this too, being privileged; here it is checked first, before any database work. 2. **The caller holds write access to the computing unit.** Mounting puts data into someone's unit, so it takes the same privilege as any other change to one; a read-only sharee may use the unit, not alter what it can see. 3. **The caller may read the repository.** It is matched by name across datasets and models — by name rather than by parsing an id out of it, because an earlier migration backfilled that column from the dataset's plain name — and anything other than exactly one match is refused rather than resolved arbitrarily. 4. **The commit belongs to that repository**, so a well-formed but unrelated digest cannot be mounted. Only then does it resolve which node the unit's pod is on — itself, rather than accepting one from the caller, or anything that can reach it could aim requests at any node's privileged mounter — and forward. file-service still re-checks read access on every byte it serves, but as the last line rather than the only one: without the check here a caller could have a mount created for a repository they cannot read, learning that it exists and spending a node's resources on it. The read and ownership rules move to a module both services share, so they decide from one definition instead of two that can drift. Mounts are released when the computing unit's pod is deleted, by the mounter's pod watcher, so there is no unmount path. No new configuration: the mounter's port and file-service's root are read from the environment the chart already sets, so each value is written once, in the chart. The whole feature stays behind `mounter.enabled`; with it off the chart renders no new object and no new environment variable. Nothing here is used only by the follow-up. Resolving a logical dataset or model path to a repository and commit, and the engine client that calls this endpoint, land in #6895 where they have callers. ### Any related issues, documentation, discussions? Closes #6863 · part of #6606 · builds on #6866. ### How was this PR tested? - New unit tests cover each refusal — no access to the unit, read-only access to it, a repository the caller cannot read, one that does not exist, a commit belonging to another repository, a malformed component, a unit not yet scheduled, and mounting disabled or misconfigured — each asserting the mounter was never asked. Plus the node lookup and the mounter client's own request shape, service-account token and error handling. Compile and tests green across the touched modules. - Validated end to end on a single-node minikube with `mounter.enabled=true`, deployed from this chart rather than by hand so the wiring itself is what ran. Every call was made from inside an unprivileged computing-unit pod: authorized mount `200`, the same mount again `200`, unauthenticated `401`, a computing unit the caller may not use `403`, a repository that does not exist `403`, a commit from another repository `403`, a malformed name `400`, and calling the node mounter directly `401` while its health probe stays open at `200`. - The mount itself: `fuse.geesefs` at the expected path inside the pod, file byte-exact, read-only, in a pod whose effective capabilities are empty. The mounter's log confirms GeeseFS was pointed at file-service's root. - Terminating the computing unit took GeeseFS mounts from one to zero and emptied the mount root, which is why there is no unmount endpoint. - With the feature off, the service account, the pod-reader role, the mounter objects, the projected token and every new environment variable render zero times. - The cluster run caught a real bug: file-service was being addressed through an endpoint whose path had to be stripped back off, when what GeeseFS needs is the root the S3 proxy is served at. Fixed. ### Was this PR authored or co-authored using generative AI tooling? Generated-by: Claude Opus 5 --------- Co-authored-by: Claude Opus 5 (1M context) <[email protected]> --- .../texera/service/AccessControlService.scala | 2 + .../resource/ComputingUnitMountResource.scala | 216 ++++++++++++++ .../service/util/ComputingUnitNodeLocator.scala | 114 ++++++++ .../apache/texera/service/util/MounterClient.scala | 157 ++++++++++ .../resource/ComputingUnitMountResourceSpec.scala | 322 +++++++++++++++++++++ .../util/ComputingUnitNodeLocatorSpec.scala | 66 +++++ .../texera/service/util/MounterClientSpec.scala | 153 ++++++++++ .../access-control-service-deployment.yaml | 12 + .../access-control-service-service-account.yaml | 50 ++-- build.sbt | 2 +- .../texera/service/resource/ResourceAccess.scala | 0 .../texera/service/resource/ResourceTables.scala | 0 12 files changed, 1076 insertions(+), 18 deletions(-) diff --git a/access-control-service/src/main/scala/org/apache/texera/service/AccessControlService.scala b/access-control-service/src/main/scala/org/apache/texera/service/AccessControlService.scala index 1f50c86c9f..ca2c797ab9 100644 --- a/access-control-service/src/main/scala/org/apache/texera/service/AccessControlService.scala +++ b/access-control-service/src/main/scala/org/apache/texera/service/AccessControlService.scala @@ -28,6 +28,7 @@ import org.apache.texera.dao.SqlServer import org.apache.texera.service.activity.UserActivityEventListener import org.apache.texera.service.resource.{ AccessControlResource, + ComputingUnitMountResource, HealthCheckResource, LiteLLMModelsResource, LiteLLMProxyResource @@ -68,6 +69,7 @@ class AccessControlService extends Application[AccessControlServiceConfiguration environment.jersey.register(classOf[AccessControlResource]) environment.jersey.register(classOf[LiteLLMProxyResource]) environment.jersey.register(classOf[LiteLLMModelsResource]) + environment.jersey.register(classOf[ComputingUnitMountResource]) AuthFeatures.register(environment) diff --git a/access-control-service/src/main/scala/org/apache/texera/service/resource/ComputingUnitMountResource.scala b/access-control-service/src/main/scala/org/apache/texera/service/resource/ComputingUnitMountResource.scala new file mode 100644 index 0000000000..b521bdf821 --- /dev/null +++ b/access-control-service/src/main/scala/org/apache/texera/service/resource/ComputingUnitMountResource.scala @@ -0,0 +1,216 @@ +/* + * 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.resource + +import com.typesafe.scalalogging.LazyLogging +import io.dropwizard.auth.Auth +import jakarta.annotation.security.RolesAllowed +import jakarta.ws.rs._ +import jakarta.ws.rs.core.MediaType +import org.apache.texera.auth.JwtAuth.{jwtClaims, jwtToken} +import org.apache.texera.auth.SessionUser +import org.apache.texera.auth.util.ComputingUnitAccess +import org.apache.texera.common.config.{EnvironmentalVariable, KubernetesConfig} +import org.apache.texera.dao.SqlServer +import org.apache.texera.dao.SqlServer.withTransaction +import org.apache.texera.dao.jooq.generated.Tables.{DATASET_VERSION, MODEL_VERSION} +import org.apache.texera.dao.jooq.generated.enums.PrivilegeEnum +import org.apache.texera.dao.jooq.generated.tables.daos.{DatasetDao, ModelDao} +import org.apache.texera.service.resource.ComputingUnitMountResource._ +import org.apache.texera.service.util.{ + ComputingUnitNodeLocator, + MountRequestValidation, + MounterClient +} + +import scala.jdk.CollectionConverters._ + +/** + * The mount authority: this service decides whether a user may act on a computing unit, so + * it is where a mount request is authorized before being forwarded to that unit's node. + */ +@Path("/mounts") +@RolesAllowed(Array("REGULAR", "ADMIN")) +@Produces(Array(MediaType.APPLICATION_JSON)) +class ComputingUnitMountResource( + mounterEnabled: Boolean, + mounterPort: Option[Int], + fileServiceBaseUrl: Option[String], + nodeLocator: ComputingUnitNodeLocator, + mounter: MounterClient +) extends LazyLogging { + + // No-arg constructor for Jersey reflection. Tests use the param-ful form. + def this() = + this( + KubernetesConfig.mounterEnabled, + EnvironmentalVariable.get(MounterPortVariable).map(_.trim.toInt), + EnvironmentalVariable.get(FileServiceUrlVariable), + ComputingUnitNodeLocator, + MounterClient + ) + + @POST + @Path("/{cuid}") + @Consumes(Array(MediaType.APPLICATION_JSON)) + def mount( + @PathParam("cuid") cuid: Int, + request: MountRequest, + @Auth user: SessionUser + ): MountInfo = { + val (port, fileService) = requireMountConfiguration() + try MountRequestValidation.validate(cuid.toString, request.repositoryName, request.commitHash) + catch { case e: IllegalArgumentException => throw new BadRequestException(e.getMessage) } + requireComputingUnitAccess(cuid, user) + requireRepositoryReadAccess(request.repositoryName, request.commitHash, user.getUid) + val nodeIp = requireNodeIp(cuid) + + // A token minted here, after the access check: GeeseFS keeps presenting it for the life + // of the mount, so it must be one this service vouched for. + val mountPath = + try { + mounter.mount( + nodeIp, + port, + cuid.toString, + request.repositoryName, + request.commitHash, + jwtToken(jwtClaims(user.getUser)), + fileService + ) + } catch { + case e: IllegalArgumentException => + throw new BadRequestException(e.getMessage) + case e: MounterClient.MounterRequestException => + logger.warn(s"node mounter at $nodeIp refused a mount for computing unit $cuid", e) + throw new BadRequestException(e.getMessage) + } + + logger.info( + s"user ${user.getUid} mounted ${request.repositoryName}:${request.commitHash} " + + s"onto computing unit $cuid at $mountPath" + ) + MountInfo(request.repositoryName, request.commitHash, mountPath) + } + + /** + * What a mount request needs beyond the request itself, both passed by the chart. Missing + * with mounting off is said plainly, because the alternative is a connection timeout to a + * node port nothing is listening on; missing with it on is a misconfiguration, and naming + * the variable beats letting a half-formed request reach the mounter. + */ + private def requireMountConfiguration(): (Int, String) = { + if (!mounterEnabled) { + throw new ServiceUnavailableException( + "Repository mounting is not enabled on this deployment." + ) + } + def required[T](value: Option[T], variable: String): T = + value.getOrElse( + throw new InternalServerErrorException( + s"Repository mounting is enabled but $variable is unset." + ) + ) + ( + required(mounterPort, MounterPortVariable), + required(fileServiceBaseUrl.filter(_.nonEmpty), FileServiceUrlVariable) + ) + } + + /** + * Mounting puts data into someone's computing unit, so it takes the same privilege as any + * other change to one: ownership, or an explicit WRITE grant. A read-only sharee may use + * the unit, not alter what it can see. + */ + private def requireComputingUnitAccess(cuid: Int, user: SessionUser): Unit = + if (ComputingUnitAccess.getComputingUnitAccess(cuid, user.getUid) != PrivilegeEnum.WRITE) { + logger.warn(s"user ${user.getUid} denied mount access to computing unit $cuid") + throw new ForbiddenException("No write access to this computing unit.") + } + + /** + * The repository must be one the user may read, at a commit that belongs to it. + * + * The mounter performs what it is told and authorizes nothing, so this is where a mount is + * refused. file-service re-checks read access on every byte served through its proxy, but + * that is the last line rather than this one: without the check here a caller could have a + * mount created for a repository they cannot read, learning it exists and spending a node's + * resources on it. + * + * A repository is matched by name rather than parsed, because `sql/updates/15.sql` + * backfilled the column from the dataset's plain name; both resource kinds are searched, + * and anything other than exactly one match is refused rather than resolved arbitrarily. + */ + private def requireRepositoryReadAccess( + repositoryName: String, + commitHash: String, + uid: Integer + ): Unit = + withTransaction(SqlServer.getInstance().createDSLContext()) { ctx => + val datasets = new DatasetDao(ctx.configuration()).fetchByRepositoryName(repositoryName) + val models = new ModelDao(ctx.configuration()).fetchByRepositoryName(repositoryName) + + val readable = (datasets.asScala.toList, models.asScala.toList) match { + case (dataset :: Nil, Nil) => + ResourceAccess.userHasReadAccess(ctx, ResourceTables.Dataset, dataset.getDid, uid) && + ctx.fetchExists( + DATASET_VERSION, + DATASET_VERSION.DID + .eq(dataset.getDid) + .and(DATASET_VERSION.VERSION_HASH.eq(commitHash)) + ) + case (Nil, model :: Nil) => + ResourceAccess.userHasReadAccess(ctx, ResourceTables.Model, model.getMid, uid) && + ctx.fetchExists( + MODEL_VERSION, + MODEL_VERSION.MID.eq(model.getMid).and(MODEL_VERSION.VERSION_HASH.eq(commitHash)) + ) + case _ => false + } + + if (!readable) { + logger.warn(s"user $uid denied a mount of '$repositoryName' at '$commitHash'") + throw new ForbiddenException("No read access to the requested repository version.") + } + } + + private def requireNodeIp(cuid: Int): String = + nodeLocator + .nodeIpOf(cuid) + .getOrElse( + throw new BadRequestException( + s"Computing unit $cuid is not running on a node yet; cannot manage its mounts." + ) + ) +} + +object ComputingUnitMountResource { + + /** Set by the chart from the same value it gives the mounter DaemonSet's hostPort. */ + private val MounterPortVariable = "KUBERNETES_MOUNTER_PORT" + + // file-service's root, which is what GeeseFS is pointed at: the S3 proxy is served at the + // servlet root, so this is scheme and authority with no path. + private val FileServiceUrlVariable = "FILE_SERVICE_URL" + + case class MountRequest(repositoryName: String, commitHash: String) + + case class MountInfo(repositoryName: String, commitHash: String, mountPath: String) +} diff --git a/access-control-service/src/main/scala/org/apache/texera/service/util/ComputingUnitNodeLocator.scala b/access-control-service/src/main/scala/org/apache/texera/service/util/ComputingUnitNodeLocator.scala new file mode 100644 index 0000000000..237d6cb2c0 --- /dev/null +++ b/access-control-service/src/main/scala/org/apache/texera/service/util/ComputingUnitNodeLocator.scala @@ -0,0 +1,114 @@ +/* + * 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.fasterxml.jackson.databind.{JsonNode, ObjectMapper} +import com.typesafe.scalalogging.LazyLogging +import org.apache.texera.common.config.KubernetesConfig + +import java.io.FileInputStream +import java.net.URI +import java.net.http.{HttpClient, HttpRequest, HttpResponse} +import java.nio.file.{Files, Paths} +import java.security.KeyStore +import java.security.cert.CertificateFactory +import java.time.Duration +import javax.net.ssl.{SSLContext, TrustManagerFactory} +import scala.jdk.CollectionConverters._ + +/** + * Finds the node a computing unit's pod runs on, so a mount can be sent to that node's + * mounter. + * + * Resolved here rather than supplied by the caller: letting a caller name the node would + * hand anything that can reach this service the ability to aim requests at any node's + * privileged mounter. + */ +class ComputingUnitNodeLocator(fetchPod: String => Option[JsonNode]) extends LazyLogging { + + def nodeIpOf(cuid: Int): Option[String] = { + val podName = s"${KubernetesConfig.computeUnitPodNamePrefix}-$cuid" + fetchPod(podName).map(_.at("/status/hostIP").asText("")).filter(_.nonEmpty) + } +} + +object ComputingUnitNodeLocator extends ComputingUnitNodeLocator(InClusterKubernetesApi.getPod) + +// One read of one field, so the API is called directly rather than through a Kubernetes +// client library and its transitive dependencies. +private[util] object InClusterKubernetesApi extends LazyLogging { + + private val serviceAccountDir = "/var/run/secrets/kubernetes.io/serviceaccount" + private val mapper = new ObjectMapper() + + private lazy val client: HttpClient = + HttpClient + .newBuilder() + .connectTimeout(Duration.ofSeconds(5)) + .sslContext(clusterSslContext) + .build() + + /** Trusts only the cluster CA, so this talks to the API server and nothing else. */ + private def clusterSslContext: SSLContext = { + val certificates = { + val stream = new FileInputStream(s"$serviceAccountDir/ca.crt") + try CertificateFactory.getInstance("X.509").generateCertificates(stream).asScala.toList + finally stream.close() + } + val keyStore = KeyStore.getInstance(KeyStore.getDefaultType) + keyStore.load(null, null) + certificates.zipWithIndex.foreach { + case (certificate, index) => keyStore.setCertificateEntry(s"cluster-ca-$index", certificate) + } + val trustManagerFactory = + TrustManagerFactory.getInstance(TrustManagerFactory.getDefaultAlgorithm) + trustManagerFactory.init(keyStore) + val context = SSLContext.getInstance("TLS") + context.init(null, trustManagerFactory.getTrustManagers, null) + context + } + + def getPod(podName: String): Option[JsonNode] = { + val host = sys.env.getOrElse("KUBERNETES_SERVICE_HOST", "kubernetes.default.svc") + val port = sys.env.getOrElse("KUBERNETES_SERVICE_PORT", "443") + val namespace = KubernetesConfig.computeUnitPoolNamespace + val token = Files.readString(Paths.get(s"$serviceAccountDir/token")).trim + + val request = HttpRequest + .newBuilder() + .uri(URI.create(s"https://$host:$port/api/v1/namespaces/$namespace/pods/$podName")) + .header("Authorization", s"Bearer $token") + .timeout(Duration.ofSeconds(10)) + .GET() + .build() + + val response = client.send(request, HttpResponse.BodyHandlers.ofString()) + response.statusCode() match { + case 200 => Some(mapper.readTree(response.body())) + case 404 => None + case other => + // Distinguished from a missing pod: a missing RBAC rule or an unreachable API + // server must not read as "the computing unit is not running". + throw new IllegalStateException( + s"cannot read pod $podName in namespace $namespace: HTTP $other ${response.body()}" + ) + } + } +} diff --git a/access-control-service/src/main/scala/org/apache/texera/service/util/MounterClient.scala b/access-control-service/src/main/scala/org/apache/texera/service/util/MounterClient.scala new file mode 100644 index 0000000000..1cb15f0213 --- /dev/null +++ b/access-control-service/src/main/scala/org/apache/texera/service/util/MounterClient.scala @@ -0,0 +1,157 @@ +/* + * 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.fasterxml.jackson.databind.{JsonNode, ObjectMapper} +import com.fasterxml.jackson.module.scala.DefaultScalaModule + +import java.net.{HttpURLConnection, URI} +import java.nio.charset.StandardCharsets +import java.nio.file.{Files, Paths} +import scala.util.Using + +/** + * HTTP client for the per-node `texera-mounter`. + * + * The mounter is privileged and listens on a hostPort, so it admits exactly one caller: it + * requires a service-account token minted for its own audience and checks it with the + * Kubernetes TokenReview API (see `authenticate_caller` in `bin/mounter/mounter.py`). That + * caller is this service, which is why the client lives here. + */ +class MounterClient(tokenPath: String = MounterDefaults.ProjectedTokenPath) { + + import MounterClient._ + + private val mapper: ObjectMapper = new ObjectMapper().registerModule(DefaultScalaModule) + + private val connectTimeoutMs = 10000 + private val readTimeoutMs = 35000 + + private def baseUrl(nodeIp: String, port: Int): String = s"http://$nodeIp:$port" + + // Read per call, not cached: the kubelet rewrites the projected token in place. + private def mounterToken(): String = + try Files.readString(Paths.get(tokenPath)).trim + catch { + case e: Exception => + throw new IllegalStateException( + s"cannot read the mounter service-account token at $tokenPath; without it this " + + s"service cannot authenticate to the node mounter: ${e.getMessage}" + ) + } + + def mount( + nodeIp: String, + port: Int, + cuid: String, + repositoryName: String, + commitHash: String, + jwt: String, + fileServiceBase: String + ): String = { + MountRequestValidation.validate(cuid, repositoryName, commitHash) + + val body = mapper.createObjectNode() + body.put("cuid", cuid) + body.put("repositoryName", repositoryName) + body.put("commitHash", commitHash) + body.put("jwt", jwt) + body.put("fileServiceBase", fileServiceBase) + + send("POST", s"${baseUrl(nodeIp, port)}/mount", Some(body.toString)) + } + + private def send(method: String, url: String, body: Option[String]): String = { + val connection = URI.create(url).toURL.openConnection().asInstanceOf[HttpURLConnection] + connection.setRequestMethod(method) + connection.setRequestProperty("Authorization", s"Bearer ${mounterToken()}") + connection.setConnectTimeout(connectTimeoutMs) + connection.setReadTimeout(readTimeoutMs) + body.foreach { _ => + connection.setRequestProperty("Content-Type", "application/json") + connection.setDoOutput(true) + } + try { + body.foreach(payload => + Using(connection.getOutputStream)(_.write(payload.getBytes(StandardCharsets.UTF_8))) + ) + val code = connection.getResponseCode + val stream = + if (code >= 200 && code < 300) connection.getInputStream else connection.getErrorStream + val responseBody = Option(stream) + .map(s => new String(s.readAllBytes(), StandardCharsets.UTF_8)) + .getOrElse("") + if (code < 200 || code >= 300) { + throw new MounterRequestException(code, s"mounter $method failed: HTTP $code $responseBody") + } + val response: JsonNode = + if (responseBody.isEmpty) mapper.createObjectNode() else mapper.readTree(responseBody) + Option(response.get("mountPath")) + .filter(node => node.isTextual && node.asText().nonEmpty) + .map(_.asText()) + .getOrElse( + throw new IllegalStateException( + s"mounter $method reported success without a mount path: $responseBody" + ) + ) + } finally { + connection.disconnect() + } + } +} + +object MounterClient extends MounterClient(MounterDefaults.ProjectedTokenPath) { + + class MounterRequestException(val status: Int, message: String) extends RuntimeException(message) +} + +private object MounterDefaults { + + /** Where the chart projects the token. Not configurable: the same chart fixes both ends. */ + val ProjectedTokenPath = "/var/run/secrets/texera/mounter/token" +} + +/** + * The shape a mount request has to have before anything acts on it. + * + * The mounter joins these into a path and creates the directory, so each has to be a single + * safe segment: no separator, no "..", and a leading alphanumeric so a value cannot be read + * as a geesefs flag. The mounter enforces this itself, being privileged; this is the same + * rule applied earlier, so a malformed request is refused by name rather than by whatever + * it fails next. + */ +private[service] object MountRequestValidation { + + private val cuidPattern = "^[0-9]+$".r + private val segmentPattern = "^[A-Za-z0-9][A-Za-z0-9._-]*$".r + + def validate(cuid: String, repositoryName: String, commitHash: String): Unit = { + if (cuid == null || cuidPattern.findFirstIn(cuid).isEmpty) { + throw new IllegalArgumentException(s"cuid must be a non-negative integer, got '$cuid'") + } + requireSegment(repositoryName, "repositoryName") + requireSegment(commitHash, "commitHash") + } + + private def requireSegment(value: String, field: String): Unit = + if (value == null || segmentPattern.findFirstIn(value).isEmpty) { + throw new IllegalArgumentException(s"$field must be a single path segment, got '$value'") + } +} diff --git a/access-control-service/src/test/scala/org/apache/texera/service/resource/ComputingUnitMountResourceSpec.scala b/access-control-service/src/test/scala/org/apache/texera/service/resource/ComputingUnitMountResourceSpec.scala new file mode 100644 index 0000000000..926568fc83 --- /dev/null +++ b/access-control-service/src/test/scala/org/apache/texera/service/resource/ComputingUnitMountResourceSpec.scala @@ -0,0 +1,322 @@ +/* + * 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.resource + +import jakarta.ws.rs.{ + BadRequestException, + ForbiddenException, + InternalServerErrorException, + ServiceUnavailableException +} +import org.apache.texera.auth.SessionUser +import org.apache.texera.dao.MockTexeraDB +import org.apache.texera.dao.jooq.generated.enums.{ + PrivilegeEnum, + UserRoleEnum, + WorkflowComputingUnitTypeEnum +} +import org.apache.texera.dao.jooq.generated.tables.daos.{ + ComputingUnitUserAccessDao, + DatasetDao, + DatasetVersionDao, + UserDao, + WorkflowComputingUnitDao +} +import org.apache.texera.dao.jooq.generated.tables.pojos.{ + ComputingUnitUserAccess, + Dataset, + DatasetVersion, + User, + WorkflowComputingUnit +} +import org.apache.texera.service.resource.ComputingUnitMountResource.MountRequest +import org.apache.texera.service.util.{ComputingUnitNodeLocator, MounterClient} +import org.scalatest.BeforeAndAfterAll +import org.scalatest.flatspec.AnyFlatSpec +import org.scalatest.matchers.should.Matchers + +import scala.collection.mutable + +class ComputingUnitMountResourceSpec + extends AnyFlatSpec + with Matchers + with BeforeAndAfterAll + with MockTexeraDB { + + private val owner: User = { + val user = new User + user.setUid(1) + user.setName("owner") + user.setEmail("[email protected]") + user.setRole(UserRoleEnum.REGULAR) + user + } + + private val reader: User = { + val user = new User + user.setUid(2) + user.setName("reader") + user.setEmail("[email protected]") + user.setRole(UserRoleEnum.REGULAR) + user + } + + private val stranger: User = { + val user = new User + user.setUid(3) + user.setName("stranger") + user.setEmail("[email protected]") + user.setRole(UserRoleEnum.REGULAR) + user + } + + private val ownedDataset: Dataset = { + val dataset = new Dataset + dataset.setDid(1) + dataset.setName("owned") + dataset.setDescription("") + dataset.setRepositoryName("dataset-1") + dataset.setOwnerUid(owner.getUid) + dataset.setIsPublic(false) + dataset + } + + private val ownedVersion: DatasetVersion = { + val version = new DatasetVersion + version.setDvid(1) + version.setDid(ownedDataset.getDid) + version.setName("v1") + version.setCreatorUid(owner.getUid) + version.setVersionHash("abc123") + version + } + + // Another user's private dataset: the only fixture whose refusal comes from read access + // rather than from the repository not existing. + private val strangersDataset: Dataset = { + val dataset = new Dataset + dataset.setDid(2) + dataset.setName("private") + dataset.setDescription("") + dataset.setRepositoryName("dataset-2") + dataset.setOwnerUid(stranger.getUid) + dataset.setIsPublic(false) + dataset + } + + private val strangersVersion: DatasetVersion = { + val version = new DatasetVersion + version.setDvid(2) + version.setDid(strangersDataset.getDid) + version.setName("v1") + version.setCreatorUid(stranger.getUid) + version.setVersionHash("def456") + version + } + + private val computingUnit: WorkflowComputingUnit = { + val unit = new WorkflowComputingUnit + unit.setCuid(7) + unit.setUid(owner.getUid) + unit.setName("test-cu") + unit.setType(WorkflowComputingUnitTypeEnum.kubernetes) + unit + } + + /** Records what the resource asked the node mounter to do, without any HTTP. */ + /** Records what the resource asked the node mounter to do, without any HTTP. */ + private class RecordingMounter extends MounterClient("/nonexistent-token") { + val mounts: mutable.Buffer[(String, Int, String, String, String, String, String)] = + mutable.Buffer() + var failWith: Option[Throwable] = None + + override def mount( + nodeIp: String, + port: Int, + cuid: String, + repositoryName: String, + commitHash: String, + jwt: String, + fileServiceBase: String + ): String = { + failWith.foreach(throw _) + mounts += ((nodeIp, port, cuid, repositoryName, commitHash, jwt, fileServiceBase)) + s"/var/lib/texera-mounts/$cuid/$repositoryName/$commitHash" + } + + } + + private val scheduledOnNode = new ComputingUnitNodeLocator(_ => None) { + override def nodeIpOf(cuid: Int): Option[String] = Some("10.0.0.4") + } + + private val notScheduled = new ComputingUnitNodeLocator(_ => None) { + override def nodeIpOf(cuid: Int): Option[String] = None + } + + private def resource( + mounter: MounterClient, + mounterEnabled: Boolean = true, + nodeLocator: ComputingUnitNodeLocator = scheduledOnNode, + mounterPort: Option[Int] = Some(8100), + fileServiceUrl: Option[String] = Some("http://file-service-svc:9092") + ) = + new ComputingUnitMountResource( + mounterEnabled, + mounterPort, + fileServiceUrl, + nodeLocator, + mounter + ) + + private def sessionOf(user: User) = new SessionUser(user) + + override protected def beforeAll(): Unit = { + initializeDBAndReplaceDSLContext() + val userDao = new UserDao(getDSLContext.configuration()) + userDao.insert(owner) + userDao.insert(reader) + userDao.insert(stranger) + new WorkflowComputingUnitDao(getDSLContext.configuration()).insert(computingUnit) + val datasetDao = new DatasetDao(getDSLContext.configuration()) + datasetDao.insert(ownedDataset) + datasetDao.insert(strangersDataset) + val versionDao = new DatasetVersionDao(getDSLContext.configuration()) + versionDao.insert(ownedVersion) + versionDao.insert(strangersVersion) + val access = new ComputingUnitUserAccess() + access.setCuid(computingUnit.getCuid) + access.setUid(reader.getUid) + access.setPrivilege(PrivilegeEnum.READ) + new ComputingUnitUserAccessDao(getDSLContext.configuration()).insert(access) + } + + override protected def afterAll(): Unit = closeConnectionPool() + + "mount" should "forward the request to the mounter on the unit's own node" in { + val mounter = new RecordingMounter + val info = resource(mounter).mount(7, MountRequest("dataset-1", "abc123"), sessionOf(owner)) + + info.repositoryName shouldBe "dataset-1" + info.commitHash shouldBe "abc123" + info.mountPath shouldBe "/var/lib/texera-mounts/7/dataset-1/abc123" + + val (nodeIp, port, cuid, repository, commit, jwt, fileServiceBase) = mounter.mounts.head + nodeIp shouldBe "10.0.0.4" + port shouldBe 8100 + cuid shouldBe "7" + repository shouldBe "dataset-1" + commit shouldBe "abc123" + fileServiceBase shouldBe "http://file-service-svc:9092" + jwt should not be empty + } + + // The only fixture with READ and not WRITE: tells "any access is enough" apart from + // "write access is required". + it should "refuse a read-only sharee, who may use the unit but not change what it sees" in { + val mounter = new RecordingMounter + a[ForbiddenException] should be thrownBy + resource(mounter).mount(7, MountRequest("dataset-1", "abc123"), sessionOf(reader)) + mounter.mounts shouldBe empty + } + + it should "refuse a repository the user cannot read" in { + val mounter = new RecordingMounter + a[ForbiddenException] should be thrownBy + resource(mounter).mount(7, MountRequest("dataset-2", "def456"), sessionOf(owner)) + // The mounter authorizes nothing, so a refusal has to happen before it is asked. + mounter.mounts shouldBe empty + } + + it should "refuse a repository that does not exist" in { + val mounter = new RecordingMounter + a[ForbiddenException] should be thrownBy + resource(mounter).mount(7, MountRequest("dataset-404", "abc123"), sessionOf(owner)) + mounter.mounts shouldBe empty + } + + it should "refuse a commit that belongs to another repository" in { + val mounter = new RecordingMounter + a[ForbiddenException] should be thrownBy + resource(mounter).mount(7, MountRequest("dataset-1", "def456"), sessionOf(owner)) + mounter.mounts shouldBe empty + } + + it should "refuse a user with no access to the computing unit" in { + val mounter = new RecordingMounter + a[ForbiddenException] should be thrownBy + resource(mounter).mount(7, MountRequest("dataset-1", "abc123"), sessionOf(stranger)) + mounter.mounts shouldBe empty + } + + it should "refuse a computing unit that does not exist" in { + val mounter = new RecordingMounter + a[ForbiddenException] should be thrownBy + resource(mounter).mount(999, MountRequest("dataset-1", "abc123"), sessionOf(owner)) + mounter.mounts shouldBe empty + } + + it should "report a misconfigured port rather than dialling a wrong one" in { + val mounter = new RecordingMounter + an[InternalServerErrorException] should be thrownBy + resource(mounter, mounterPort = None) + .mount(7, MountRequest("dataset-1", "abc123"), sessionOf(owner)) + mounter.mounts shouldBe empty + } + + it should "report a missing file-service address rather than mounting against nothing" in { + val mounter = new RecordingMounter + an[InternalServerErrorException] should be thrownBy + resource(mounter, fileServiceUrl = None) + .mount(7, MountRequest("dataset-1", "abc123"), sessionOf(owner)) + mounter.mounts shouldBe empty + } + + it should "answer plainly when the deployment did not enable mounting" in { + val mounter = new RecordingMounter + a[ServiceUnavailableException] should be thrownBy + resource(mounter, mounterEnabled = false) + .mount(7, MountRequest("dataset-1", "abc123"), sessionOf(owner)) + mounter.mounts shouldBe empty + } + + it should "refuse while the unit's pod is not on a node yet" in { + val mounter = new RecordingMounter + a[BadRequestException] should be thrownBy + resource(mounter, nodeLocator = notScheduled) + .mount(7, MountRequest("dataset-1", "abc123"), sessionOf(owner)) + mounter.mounts shouldBe empty + } + + it should "report a rejected path as a bad request rather than a server error" in { + val mounter = new RecordingMounter + mounter.failWith = Some(new IllegalArgumentException("repositoryName must be a single segment")) + val failure = the[BadRequestException] thrownBy + resource(mounter).mount(7, MountRequest("../evil", "abc123"), sessionOf(owner)) + failure.getMessage should include("repositoryName") + } + + it should "relay a refusal from the mounter as a bad request" in { + val mounter = new RecordingMounter + mounter.failWith = Some(new MounterClient.MounterRequestException(400, "mounter said no")) + a[BadRequestException] should be thrownBy + resource(mounter).mount(7, MountRequest("dataset-1", "abc123"), sessionOf(owner)) + } +} diff --git a/access-control-service/src/test/scala/org/apache/texera/service/util/ComputingUnitNodeLocatorSpec.scala b/access-control-service/src/test/scala/org/apache/texera/service/util/ComputingUnitNodeLocatorSpec.scala new file mode 100644 index 0000000000..33f3a6a034 --- /dev/null +++ b/access-control-service/src/test/scala/org/apache/texera/service/util/ComputingUnitNodeLocatorSpec.scala @@ -0,0 +1,66 @@ +/* + * 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.fasterxml.jackson.databind.{JsonNode, ObjectMapper} +import org.apache.texera.common.config.KubernetesConfig +import org.scalatest.flatspec.AnyFlatSpec +import org.scalatest.matchers.should.Matchers + +import scala.collection.mutable + +class ComputingUnitNodeLocatorSpec extends AnyFlatSpec with Matchers { + + private val mapper = new ObjectMapper() + private def pod(json: String): JsonNode = mapper.readTree(json) + + private def locator( + pods: Map[String, JsonNode], + asked: mutable.Buffer[String] = mutable.Buffer() + ): ComputingUnitNodeLocator = + new ComputingUnitNodeLocator(name => { asked += name; pods.get(name) }) + + private val podName = s"${KubernetesConfig.computeUnitPodNamePrefix}-7" + + "nodeIpOf" should "return the host IP of the computing unit's pod" in { + locator(Map(podName -> pod("""{"status":{"hostIP":"10.0.0.4"}}"""))).nodeIpOf(7) shouldBe + Some("10.0.0.4") + } + + it should "ask for the pod named by the configured prefix and the cuid" in { + val asked = mutable.Buffer[String]() + locator(Map.empty, asked).nodeIpOf(7) + asked should contain only podName + } + + it should "return None when the computing unit has no pod" in { + locator(Map.empty).nodeIpOf(7) shouldBe None + } + + it should "return None while the pod is not scheduled yet" in { + locator(Map(podName -> pod("""{"status":{"phase":"Pending"}}"""))).nodeIpOf(7) shouldBe None + locator(Map(podName -> pod("""{"status":{"hostIP":""}}"""))).nodeIpOf(7) shouldBe None + } + + it should "propagate a lookup failure rather than reporting the unit as unscheduled" in { + val failing = new ComputingUnitNodeLocator(_ => throw new IllegalStateException("forbidden")) + an[IllegalStateException] should be thrownBy failing.nodeIpOf(7) + } +} diff --git a/access-control-service/src/test/scala/org/apache/texera/service/util/MounterClientSpec.scala b/access-control-service/src/test/scala/org/apache/texera/service/util/MounterClientSpec.scala new file mode 100644 index 0000000000..beb14509ca --- /dev/null +++ b/access-control-service/src/test/scala/org/apache/texera/service/util/MounterClientSpec.scala @@ -0,0 +1,153 @@ +/* + * 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.sun.net.httpserver.{HttpExchange, HttpServer} +import org.scalatest.BeforeAndAfterAll +import org.scalatest.flatspec.AnyFlatSpec +import org.scalatest.matchers.should.Matchers + +import java.net.InetSocketAddress +import java.nio.charset.StandardCharsets +import java.nio.file.Files +import scala.collection.mutable + +class MounterClientSpec extends AnyFlatSpec with Matchers with BeforeAndAfterAll { + + private var server: HttpServer = _ + private var port: Int = _ + private var client: MounterClient = _ + + private val received = mutable.Map[String, (String, String, String)]() + private val authorization = mutable.Map[String, String]() + private var replyWith: Option[(Int, String)] = None + + private def bodyOf(exchange: HttpExchange): String = + new String(exchange.getRequestBody.readAllBytes(), StandardCharsets.UTF_8) + + private def reply(exchange: HttpExchange, status: Int, body: String): Unit = { + val bytes = body.getBytes(StandardCharsets.UTF_8) + exchange.getResponseHeaders.add("Content-Type", "application/json") + exchange.sendResponseHeaders(status, bytes.length.toLong) + exchange.getResponseBody.write(bytes) + exchange.close() + } + + private def record(exchange: HttpExchange, path: String): Unit = { + authorization(path) = Option(exchange.getRequestHeaders.getFirst("Authorization")).getOrElse("") + received(path) = ( + exchange.getRequestMethod, + Option(exchange.getRequestURI.getQuery).getOrElse(""), + bodyOf(exchange) + ) + } + + override def beforeAll(): Unit = { + val tokenFile = Files.createTempFile("mounter-token", "") + Files.writeString(tokenFile, "the-service-account-token\n") + tokenFile.toFile.deleteOnExit() + client = new MounterClient(tokenFile.toString) + + server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0) + server.createContext( + "/mount", + (exchange: HttpExchange) => { + record(exchange, "/mount") + replyWith match { + case Some((status, body)) => reply(exchange, status, body) + case None => + reply(exchange, 200, """{"mountPath":"/var/lib/texera-mounts/7/dataset-1/abc123"}""") + } + } + ) + server.start() + port = server.getAddress.getPort + } + + override def afterAll(): Unit = if (server != null) server.stop(0) + + private val nodeIp = "127.0.0.1" + + "MounterClient.mount" should "post the mount request and return the mounter's path" in { + val path = client.mount(nodeIp, port, "7", "dataset-1", "abc123", "user-jwt", "http://fs:9092") + + path shouldBe "/var/lib/texera-mounts/7/dataset-1/abc123" + val (method, _, requestBody) = received("/mount") + method shouldBe "POST" + requestBody should include(""""cuid":"7"""") + requestBody should include(""""repositoryName":"dataset-1"""") + requestBody should include(""""commitHash":"abc123"""") + requestBody should include(""""jwt":"user-jwt"""") + requestBody should include(""""fileServiceBase":"http://fs:9092"""") + } + + it should "identify itself with the projected service-account token" in { + client.mount(nodeIp, port, "7", "dataset-1", "abc123", "user-jwt", "http://fs:9092") + authorization("/mount") shouldBe "Bearer the-service-account-token" + } + + it should "carry the mounter's status back on a refusal" in { + replyWith = Some((400, """{"error":"nope"}""")) + try { + val failure = the[MounterClient.MounterRequestException] thrownBy + client.mount(nodeIp, port, "7", "dataset-1", "abc123", "user-jwt", "http://fs:9092") + failure.status shouldBe 400 + } finally replyWith = None + } + + it should "fail rather than report a mount when the mounter's success names no path" in { + Seq("{}", """{"mountPath":null}""", """{"mountPath":""}""").foreach { body => + replyWith = Some((200, body)) + try { + an[IllegalStateException] should be thrownBy + client.mount(nodeIp, port, "7", "dataset-1", "abc123", "user-jwt", "http://fs:9092") + } finally replyWith = None + } + } + + // The escapes reported on the infrastructure PR: each would otherwise be joined into the + // mount path, and the directory created before the mounter's own validation could matter. + it should "refuse a cuid that is not a single numeric segment, without calling the mounter" in { + received.remove("/mount") + Seq("5/../8", "../..", "/absolute", "", "7x").foreach { cuid => + an[IllegalArgumentException] should be thrownBy + client.mount(nodeIp, port, cuid, "dataset-1", "abc123", "jwt", "http://fs:9092") + } + received should not contain key("/mount") + } + + it should "refuse a repository or commit that is not a single safe segment" in { + received.remove("/mount") + Seq("../evil", "a/b", "-o", "", ".hidden/../x").foreach { bad => + an[IllegalArgumentException] should be thrownBy + client.mount(nodeIp, port, "7", bad, "abc123", "jwt", "http://fs:9092") + an[IllegalArgumentException] should be thrownBy + client.mount(nodeIp, port, "7", "dataset-1", bad, "jwt", "http://fs:9092") + } + received should not contain key("/mount") + } + + it should "fail loudly when the service-account token is missing" in { + val withoutToken = new MounterClient("/nonexistent/mounter/token") + val failure = the[IllegalStateException] thrownBy + withoutToken.mount(nodeIp, port, "7", "dataset-1", "abc123", "jwt", "http://fs:9092") + failure.getMessage should include("/nonexistent/mounter/token") + } +} 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 62bb8f3b21..7d63720365 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 @@ -60,6 +60,18 @@ spec: - name: KUBERNETES_JUPYTER_BASE_URL value: {{ include "texera.jupyter.basePath" . }} {{- end }} + {{- if .Values.mounter.enabled }} + # Passed from the chart rather than left to the kubernetes.conf defaults, so that + # changing a value here reaches the service that reads it. + - name: KUBERNETES_MOUNTER_ENABLED + value: "{{ .Values.mounter.enabled }}" + - name: KUBERNETES_COMPUTE_UNIT_POD_NAME_PREFIX + value: {{ .Values.workflowComputingUnitPool.podNamePrefix }} + - name: KUBERNETES_MOUNTER_PORT + value: "{{ .Values.mounter.port }}" + - name: FILE_SERVICE_URL + value: http://{{ .Values.fileService.name }}-svc:{{ .Values.fileService.service.port }} + {{- end }} {{- if .Values.litellm.enabled }} # LLM gateway used to serve /api/chat and /api/models to the agent service. - name: LITELLM_BASE_URL 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 index 44cd6bbf07..c4e603e799 100644 --- 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 @@ -18,16 +18,13 @@ # 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 access-control-service is the only component allowed to ask the per-node mounter to +# mount a repository: 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?" -- +# 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: +# How that is enforced: # # 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 @@ -44,22 +41,41 @@ # 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 +# 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. +# The Role below lets this service read which node a computing unit's pod runs on, so it can +# forward a mount to that node's mounter. apiVersion: v1 kind: ServiceAccount metadata: name: {{ .Values.accessControlService.serviceAccountName }} namespace: {{ .Release.Namespace }} +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: {{ .Release.Name }}-access-control-service-pod-reader + namespace: {{ .Values.workflowComputingUnitPool.namespace }} +rules: + - apiGroups: [""] + resources: ["pods"] + verbs: ["get"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: {{ .Release.Name }}-access-control-service-pod-reader-binding + namespace: {{ .Values.workflowComputingUnitPool.namespace }} +subjects: + - kind: ServiceAccount + name: {{ .Values.accessControlService.serviceAccountName }} + namespace: {{ .Release.Namespace }} +roleRef: + kind: Role + name: {{ .Release.Name }}-access-control-service-pod-reader + apiGroup: rbac.authorization.k8s.io {{- end }} diff --git a/build.sbt b/build.sbt index 3853666c68..cb16bda952 100644 --- a/build.sbt +++ b/build.sbt @@ -122,7 +122,7 @@ ThisBuild / excludeDependencies += ExclusionRule("log4j", "log4j") lazy val Util = (project in file("common/util")).settings(commonModuleSettings) lazy val DAO = (project in file("common/dao")).settings(commonModuleSettings) lazy val Config = (project in file("common/config")).settings(commonModuleSettings) -lazy val Resource = (project in file("common/resource")).settings(commonModuleSettings) +lazy val Resource = (project in file("common/resource")).settings(commonModuleSettings).dependsOn(DAO) lazy val Auth = (project in file("common/auth")) .settings(commonModuleSettings) .configs(Test) diff --git a/file-service/src/main/scala/org/apache/texera/service/resource/ResourceAccess.scala b/common/resource/src/main/scala/org/apache/texera/service/resource/ResourceAccess.scala similarity index 100% rename from file-service/src/main/scala/org/apache/texera/service/resource/ResourceAccess.scala rename to common/resource/src/main/scala/org/apache/texera/service/resource/ResourceAccess.scala diff --git a/file-service/src/main/scala/org/apache/texera/service/resource/ResourceTables.scala b/common/resource/src/main/scala/org/apache/texera/service/resource/ResourceTables.scala similarity index 100% rename from file-service/src/main/scala/org/apache/texera/service/resource/ResourceTables.scala rename to common/resource/src/main/scala/org/apache/texera/service/resource/ResourceTables.scala
