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-8475-5042d96ec85d98ed18bde841b2e74922c43f5c3b in repository https://gitbox.apache.org/repos/asf/texera.git
commit 50ecf40d04b7c5b819d7291508111face933b378 Author: Tanishq Gandhi <[email protected]> AuthorDate: Mon Sep 14 00:35:01 2026 +0000 feat(computing-unit): curate computing-unit images (#8475) ### What changes were proposed in this PR? Every computing unit runs the same image, fixed when the cluster is installed. ML Models needing a different Python version, a system package, or a library built from source cannot run, because a Python virtual environment only holds pip packages. This lets an administrator register an image reference from a public registry, and a computing unit can then be started from it. **Off by default** (`curatedImages.enabled: false`) until the UI to manage these ships in #8470 and #8471. **How it works.** Texera reads the image's manifest and config blob — a few kilobytes, never the layers — to check its start command runs `computing-unit-master`, which means it was built `FROM` the Texera computing-unit image, and to resolve the digest behind the reference. A misspelled, private or unsuitable image is refused in seconds, in front of the administrator, rather than when a user's unit will not start. The row records `owner/name@sha256:…`, and that is what units run, so a tag its owner moves later cannot change what already ran. **Nothing is copied and no registry is added** — units pull the reference the same way the deployment's own image is already pulled. **Uniqueness is enforced by the database**, not only checked in the service: two administrators registering the same link at the same moment both pass a read-then-write check and produce two rows for one image. **Non-root, for curated images only.** A curated image was supplied by an administrator and reviewed by nobody, so a unit started from one runs as a non-root user with no privilege escalation and no capabilities. The deployment's own image is untouched — it is its operator's choice, and a deployment that has replaced it with an image needing root would break on upgrade. **Known limitation.** The first unit on each node waits for the image to download — about 80 seconds for a 3 GB one — while later units on that node start at once. Pre-pulling ready images onto nodes is #8469. ### Any related issues, documentation, discussions? Closes #8468 Part of #8466 ### How was this PR tested? Unit tests, chart rendering, and a deployment to minikube exercising both states. ``` sbt "ComputingUnitManagingService/testOnly org.apache.texera.service.resource.CuratedImageResourceSpec" "ComputingUnitManagingService/testOnly org.apache.texera.service.util.KubernetesClientSpec" "Config/testOnly org.apache.texera.common.config.KubernetesConfigSpec" scalafmtCheckAll ``` CuratedImageResourceSpec 27 passed KubernetesClientSpec 17 passed KubernetesConfigSpec 6 passed scalafmtCheckAll clean `helm template` renders with the feature off and on; the manager Role gains `jobs` and `pods/log` only. **Deployed to minikube, feature off:** ``` GET /api/cu-image 503 "Curated images are not enabled on this deployment." POST /api/cu-image 503 create unit with iid 403 "Image 1 is not available..." create unit without iid 200 deployment's own image, no security context ``` **Feature on:** register a good image READY in 8s, pinned to @sha256:bdeadc3c... duplicate reference 400 names the existing row duplicate name 400 empty name 400 tag that does not exist FAILED, log names the tag and what to do instead alpine (not a CU image) FAILED, "Its start command is: [/bin/sh]" unit from a curated image image tagandhi19/texera-cu-sklearn@sha256:bdeadc3c... security {allowPrivilegeEscalation:false, capabilities:{drop:[ALL]}, runAsNonRoot:true, runAsUser:1001} id uid=1001(texera) delete the image while a unit runs on it 204, pod stays Running, unit still reports imageName 'Python ML' a new unit from it 403 ### Was this PR authored or co-authored using generative AI tooling? Generated-by: Claude Code (Claude Opus 5) --- bin/k8s/templates/base/gateway/gateway-routes.yaml | 9 + ...workflow-computing-unit-manager-deployment.yaml | 5 + ...low-computing-unit-manager-service-account.yaml | 8 + bin/k8s/values.yaml | 5 + common/config/src/main/resources/kubernetes.conf | 41 +- .../texera/common/config/CuratedImageConfig.scala | 47 ++ .../texera/common/config/KubernetesConfig.scala | 6 + .../service/ComputingUnitManagingService.scala | 2 + .../resource/ComputingUnitManagingResource.scala | 25 +- .../service/resource/CuratedImageResource.scala | 579 +++++++++++++++++++++ .../service/util/ImageValidationClient.scala | 444 ++++++++++++++++ .../texera/service/util/KubernetesClient.scala | 23 +- .../resource/CuratedImageResourceSpec.scala | 459 ++++++++++++++++ sql/changelog.xml | 5 + sql/texera_ddl.sql | 28 + sql/updates/49.sql | 61 +++ 16 files changed, 1742 insertions(+), 5 deletions(-) diff --git a/bin/k8s/templates/base/gateway/gateway-routes.yaml b/bin/k8s/templates/base/gateway/gateway-routes.yaml index eb88a2a09d..fda04b7ff2 100644 --- a/bin/k8s/templates/base/gateway/gateway-routes.yaml +++ b/bin/k8s/templates/base/gateway/gateway-routes.yaml @@ -35,6 +35,15 @@ spec: backendRefs: - name: workflow-computing-unit-manager-svc port: 8888 + # The curated-image catalogue lives on the computing-unit manager, not the webserver. + # Without its own rule it falls through to the /api catch-all and every request 404s. + - matches: + - path: + type: PathPrefix + value: /api/cu-image + backendRefs: + - name: workflow-computing-unit-manager-svc + port: 8888 - matches: - path: type: PathPrefix 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 d55eb5f10c..6ad4eeb4be 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 @@ -66,6 +66,11 @@ spec: value: {{ .Values.workflowComputingUnitPool.name }}-svc - name: KUBERNETES_COMPUTE_UNIT_POD_NAME_PREFIX value: {{ .Values.workflowComputingUnitPool.podNamePrefix }} + - name: TEXERA_CURATED_IMAGES_ENABLED + value: "{{ .Values.curatedImages.enabled }}" + # Must be the namespace the pool runs in; validation jobs go there. + - name: TEXERA_CURATED_IMAGE_VALIDATION_NAMESPACE + value: {{ .Values.workflowComputingUnitPool.namespace }} - name: KUBERNETES_IMAGE_NAME value: {{ .Values.texera.imageRegistry }}/{{ .Values.workflowComputingUnitPool.imageName }}:{{ .Values.texera.imageTag }} - name: KUBERNETES_MOUNTER_ENABLED diff --git a/bin/k8s/templates/base/workflow-computing-unit-manager/workflow-computing-unit-manager-service-account.yaml b/bin/k8s/templates/base/workflow-computing-unit-manager/workflow-computing-unit-manager-service-account.yaml index 44dce0cf87..29888af964 100644 --- a/bin/k8s/templates/base/workflow-computing-unit-manager/workflow-computing-unit-manager-service-account.yaml +++ b/bin/k8s/templates/base/workflow-computing-unit-manager/workflow-computing-unit-manager-service-account.yaml @@ -34,6 +34,14 @@ rules: - apiGroups: ["metrics.k8s.io"] # Added metrics permissions resources: ["pods"] verbs: ["list", "get"] # Added metrics permissions + # Validating a curated image runs as a Job here, and its outcome is read back from the + # Job and from its pod's log -- a Job cannot return a value. + - apiGroups: ["batch"] + resources: ["jobs"] + verbs: ["get", "list", "watch", "create", "delete"] + - apiGroups: [""] + resources: ["pods/log"] + verbs: ["get"] --- apiVersion: rbac.authorization.k8s.io/v1 diff --git a/bin/k8s/values.yaml b/bin/k8s/values.yaml index 46bb349fca..16f7cc2062 100644 --- a/bin/k8s/values.yaml +++ b/bin/k8s/values.yaml @@ -341,6 +341,11 @@ litellm: model: gpt-5-mini api_key: "os.environ/OPENAI_API_KEY" +# Images an administrator registers, which a computing unit can then be started from. +curatedImages: + # Off until the UI to manage these ships. + enabled: false + # headless service for the access of computing units workflowComputingUnitPool: createNamespaces: true diff --git a/common/config/src/main/resources/kubernetes.conf b/common/config/src/main/resources/kubernetes.conf index e3edf31aff..c27fa40d04 100644 --- a/common/config/src/main/resources/kubernetes.conf +++ b/common/config/src/main/resources/kubernetes.conf @@ -54,6 +54,16 @@ kubernetes { computing-unit-memory-limit-options = "1Gi,2Gi,4Gi" computing-unit-memory-limit-options = ${?KUBERNETES_COMPUTING_UNIT_MEMORY_LIMIT_OPTIONS} + # A curated image is supplied by an administrator and reviewed by nobody, so units + # started from one are pinned to a non-root user. runAsUser is stated too because + # kubelet cannot verify an image whose USER is a name, which the Texera image's is. + computing-unit-run-as-non-root = true + computing-unit-run-as-non-root = ${?KUBERNETES_COMPUTING_UNIT_RUN_AS_NON_ROOT} + + # The uid the Texera computing-unit image creates. It must exist in the image. + computing-unit-run-as-user = 1001 + computing-unit-run-as-user = ${?KUBERNETES_COMPUTING_UNIT_RUN_AS_USER} + # GPU configuration computing-unit-gpu-limit-options = "0,1,2" computing-unit-gpu-limit-options = ${?KUBERNETES_COMPUTING_UNIT_GPU_LIMIT_OPTIONS} @@ -112,4 +122,33 @@ kubernetes { # 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 +} + +curated-images { + # Off until the UI to manage these ships. + enabled = false + enabled = ${?TEXERA_CURATED_IMAGES_ENABLED} + + # Reads a manifest to check the start command and resolve the digest. No layers are + # pulled, so this is a few kilobytes. + validation-image = "quay.io/skopeo/stable:v1.16.1" + validation-image = ${?TEXERA_CURATED_IMAGE_VALIDATION_IMAGE} + + # Where validation jobs run; must be the pool namespace, which the chart sets. + validation-namespace = "texera-workflow-computing-unit-pool" + validation-namespace = ${?TEXERA_CURATED_IMAGE_VALIDATION_NAMESPACE} + + # Generous, because the wait is on a registry answering. + validation-timeout-seconds = 300 + validation-timeout-seconds = ${?TEXERA_CURATED_IMAGE_VALIDATION_TIMEOUT_SECONDS} + + # Reading a manifest is neither CPU nor memory work. + validation-cpu-request = "100m" + validation-cpu-request = ${?TEXERA_CURATED_IMAGE_VALIDATION_CPU_REQUEST} + validation-memory-request = "128Mi" + validation-memory-request = ${?TEXERA_CURATED_IMAGE_VALIDATION_MEMORY_REQUEST} + validation-cpu-limit = "1" + validation-cpu-limit = ${?TEXERA_CURATED_IMAGE_VALIDATION_CPU_LIMIT} + validation-memory-limit = "256Mi" + validation-memory-limit = ${?TEXERA_CURATED_IMAGE_VALIDATION_MEMORY_LIMIT} +} diff --git a/common/config/src/main/scala/org/apache/texera/common/config/CuratedImageConfig.scala b/common/config/src/main/scala/org/apache/texera/common/config/CuratedImageConfig.scala new file mode 100644 index 0000000000..9dc85da9d6 --- /dev/null +++ b/common/config/src/main/scala/org/apache/texera/common/config/CuratedImageConfig.scala @@ -0,0 +1,47 @@ +/* + * 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.common.config + +import com.typesafe.config.{Config, ConfigFactory} + +object CuratedImageConfig { + + private val conf: Config = ConfigFactory.parseResources("kubernetes.conf").resolve() + + val enabled: Boolean = conf.getBoolean("curated-images.enabled") + + val validationImage: String = conf.getString("curated-images.validation-image") + val validationNamespace: String = conf.getString("curated-images.validation-namespace") + val validationTimeoutSeconds: Int = conf.getInt("curated-images.validation-timeout-seconds") + + val validationCpuRequest: String = conf.getString("curated-images.validation-cpu-request") + val validationMemoryRequest: String = conf.getString("curated-images.validation-memory-request") + val validationCpuLimit: String = conf.getString("curated-images.validation-cpu-limit") + val validationMemoryLimit: String = conf.getString("curated-images.validation-memory-limit") + + /** + * What a computing unit runs, so an image that does not provide it cannot be one. The + * computing-unit image declares it as its CMD; the validation checks for it before + * the image can be used, so a wrong one fails in seconds and in front of the admin. + */ + val requiredCommand: String = "computing-unit-master" + + /** Kubernetes object name for one check. Unique per attempt so retries never collide. */ + def validationJobName(iid: Int, attempt: Int): String = s"cu-image-check-$iid-$attempt" +} 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 7cb177c6b6..7b56c41864 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 @@ -90,4 +90,10 @@ object KubernetesConfig { // -- 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") + + // See kubernetes.conf on why the uid has to be given alongside runAsNonRoot. + val computingUnitRunAsNonRoot: Boolean = + conf.getBoolean("kubernetes.computing-unit-run-as-non-root") + val computingUnitRunAsUser: Long = conf.getLong("kubernetes.computing-unit-run-as-user") + } diff --git a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/ComputingUnitManagingService.scala b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/ComputingUnitManagingService.scala index f0dffc89e1..97c9159184 100644 --- a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/ComputingUnitManagingService.scala +++ b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/ComputingUnitManagingService.scala @@ -30,6 +30,7 @@ import org.apache.texera.service.resource.{ AdminComputingUnitResource, ComputingUnitAccessResource, ComputingUnitManagingResource, + CuratedImageResource, HealthCheckResource } import java.nio.file.Path @@ -68,6 +69,7 @@ class ComputingUnitManagingService extends Application[ComputingUnitManagingServ environment.jersey().register(new ComputingUnitManagingResource) environment.jersey().register(new ComputingUnitAccessResource) environment.jersey().register(new AdminComputingUnitResource) + environment.jersey().register(new CuratedImageResource) RoleAnnotationEnforcer.enforce( environment.jersey.getResourceConfig, diff --git a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/ComputingUnitManagingResource.scala b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/ComputingUnitManagingResource.scala index c5df9078b0..a281a0768e 100644 --- a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/ComputingUnitManagingResource.scala +++ b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/ComputingUnitManagingResource.scala @@ -175,7 +175,9 @@ object ComputingUnitManagingResource { gpuLimit: String, jvmMemorySize: String, shmSize: String, - uri: Option[String] = None + uri: Option[String] = None, + /** A curated image to start this unit from, instead of the deployment's own. */ + iid: Option[Int] = None ) case class WorkflowComputingUnitResourceLimit( @@ -371,6 +373,18 @@ class ComputingUnitManagingResource { throw new ForbiddenException(s"Unsupported computing-unit type: ${param.unitType}") } + // Resolved before anything is written. Starting from an image that is not ready would + // leave a computing-unit row behind that can never run. + val curatedImage: Option[String] = param.iid.map { iid => + CuratedImageResource + .readyImageFor(iid) + .getOrElse( + throw new ForbiddenException( + s"Image $iid is not available. It must exist and have passed its check." + ) + ) + } + withTransaction(context) { ctx => val wcDao = new WorkflowComputingUnitDao(ctx.configuration()) @@ -395,6 +409,12 @@ class ComputingUnitManagingResource { "gpuLimit" -> param.gpuLimit, "jvmMemorySize" -> param.jvmMemorySize, "shmSize" -> param.shmSize, + // The name is stored with the id because a curated image can be removed + // while a unit started from it is still up, and "what is this running?" + // should still have an answer then. + "iid" -> param.iid, + "imageName" -> param.iid.flatMap(CuratedImageResource.nameOf), + "curatedImage" -> curatedImage, "nodeAddresses" -> Json.arr() // filled in later ) ) @@ -467,7 +487,8 @@ class ComputingUnitManagingResource { EnvironmentalVariable.ENV_USER_JWT_TOKEN -> userToken, EnvironmentalVariable.ENV_JAVA_OPTS -> s"-Xmx${param.jvmMemorySize}" ), - Some(param.shmSize) + Some(param.shmSize), + curatedImage ) } catch { diff --git a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/CuratedImageResource.scala b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/CuratedImageResource.scala new file mode 100644 index 0000000000..8c19a0138d --- /dev/null +++ b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/resource/CuratedImageResource.scala @@ -0,0 +1,579 @@ +/* + * 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.SessionUser +import org.apache.texera.common.config.CuratedImageConfig +import org.apache.texera.dao.SqlServer +import org.apache.texera.service.util.ImageValidationClient +import org.apache.texera.service.util.ImageValidationClient.ValidationState +import org.jooq.impl.DSL +import org.jooq.{DSLContext, Record} + +import java.sql.Timestamp +import scala.jdk.CollectionConverters._ + +object CuratedImageResource extends LazyLogging { + + private def context: DSLContext = SqlServer.getInstance().createDSLContext() + + // Plain DSL rather than generated DAOs: jOOQ sources are generated at build time against + // a live database and are not in the repository, so this keeps a clean checkout building. + private val CU_IMAGE = DSL.table(DSL.name("cu_image")) + private val IID = DSL.field(DSL.name("iid"), classOf[Integer]) + private val NAME = DSL.field(DSL.name("name"), classOf[String]) + private val SOURCE_REF = DSL.field(DSL.name("source_ref"), classOf[String]) + private val SOURCE_DIGEST = DSL.field(DSL.name("source_digest"), classOf[String]) + private val STATUS = DSL.field(DSL.name("status"), classOf[String]) + private val ATTEMPT = DSL.field(DSL.name("attempt"), classOf[Integer]) + private val VALIDATION_LOG = DSL.field(DSL.name("validation_log"), classOf[String]) + private val CREATED_BY = DSL.field(DSL.name("created_by"), classOf[Integer]) + private val CREATION_TIME = DSL.field(DSL.name("creation_time"), classOf[Timestamp]) + private val UPDATE_TIME = DSL.field(DSL.name("update_time"), classOf[Timestamp]) + + object Status { + val Pending = "PENDING" + val Validating = "VALIDATING" + val Ready = "READY" + val Failed = "FAILED" + } + + case class CuratedImage( + iid: Int, + name: String, + sourceRef: String, + sourceDigest: String, + status: String, + imageTag: String, + attempt: Int, + creationTime: Long, + updateTime: Long + ) + + case class CuratedImageRequest(name: String, sourceRef: String) + case class ValidationLog(iid: Int, status: String, attempt: Int, log: String) + + // An allowlist: no quotes, angle brackets, or path and shell metacharacters, so a name is + // safe to render and to log. Parentheses and "+" are allowed because real names use them + // -- "Python ML (sklearn)". + private val NamePattern = "^[A-Za-z0-9][A-Za-z0-9._ ()+-]*$".r + + /** Whether a display name is one an administrator is allowed to give an image. */ + private[resource] def isValidName(raw: String): Boolean = { + val name = Option(raw).map(_.trim).getOrElse("") + NamePattern.pattern.matcher(name).matches() && name.length <= MaxNameLength + } + // The characters a registry reference is made of. The reference reaches a shell only as + // an environment value, never as script text, so this is a second line rather than the + // defence -- but it refuses a malformed reference with a clear message instead of a + // puzzling failure from skopeo. + private val RefPattern = "^[A-Za-z0-9][A-Za-z0-9._:/@+-]*$".r + + // Registries require the repository path to be lowercase, so an uppercase one is the + // ordinary typo -- and the one worth catching here rather than in the validation job. + private val RepoPattern = "^[a-z0-9][a-z0-9._:/-]*$".r + + /** The repository part of a reference: what is left once any tag or digest is removed. */ + private[service] def repositoryOf(reference: String): String = { + val withoutDigest = reference.split("@").head + val lastSlash = withoutDigest.lastIndexOf('/') + val colon = withoutDigest.indexOf(':', lastSlash + 1) + if (colon >= 0) withoutDigest.substring(0, colon) else withoutDigest + } + + private val MaxNameLength = 128 + private val MaxRefLength = 512 + + /** Accepts either an image reference or the Docker Hub page address it was copied from. */ + private[resource] def normaliseRef(raw: String): String = { + val trimmed = raw.trim.stripSuffix("/") + val withoutScheme = trimmed.replaceFirst("^https?://", "") + val repo = + if (withoutScheme.startsWith(DockerHubHost + "/")) dockerHubRepo(withoutScheme) + else stripImplicitDockerHubRegistry(withoutScheme) + + // Default a missing tag to :latest, as every container tool does. Looks after the last + // slash so a registry's port is not mistaken for a tag. + val lastSegment = repo.substring(repo.lastIndexOf('/') + 1) + if (lastSegment.contains(":") || repo.contains("@sha256:")) repo else s"$repo:latest" + } + + private val DockerHubHost = "hub.docker.com" + + /** + * Reduces the ways of naming a Docker Hub image to one string, so "owner/name:1" and + * "docker.io/owner/name:1" are seen as the duplicate they are. + */ + private val DockerHubRegistries = + Seq("docker.io/", "index.docker.io/", "registry-1.docker.io/") + + private def stripImplicitDockerHubRegistry(reference: String): String = + DockerHubRegistries.find(reference.startsWith) match { + case None => reference + case Some(registry) => + val path = reference.stripPrefix(registry) + if (path.startsWith("library/")) path.stripPrefix("library/") else path + } + + /** + * The pull reference inside a Docker Hub web address, since pasting one is the easy + * mistake to make: + * + * hub.docker.com/r/<owner>/<name> the public page + * hub.docker.com/_/<name> an official image + * hub.docker.com/repository/docker/<owner>/<name> the owner's own page + * + * A trailing tab segment (/general, /tags) belongs to the page, not the reference. + * An unrecognised address is returned unchanged for validate() to reject. + */ + private def dockerHubRepo(address: String): String = { + val path = address.stripPrefix(DockerHubHost + "/") + if (path.startsWith("_/")) path.stripPrefix("_/").split("/").head + else if (path.startsWith("r/")) path.stripPrefix("r/").split("/").take(2).mkString("/") + else if (path.startsWith("repository/docker/")) + path.stripPrefix("repository/docker/").split("/").take(2).mkString("/") + else address + } + + /** + * The reference a unit runs: the registered repository at the digest validation resolved. + * Derived rather than stored, so the two can never disagree. + */ + private[service] def pinnedRefOf(sourceRef: String, sourceDigest: String): Option[String] = + Option(sourceDigest).filter(_.nonEmpty).map(ImageValidationClient.pinnedRef(sourceRef, _)) + + private def toCuratedImage(record: Record): CuratedImage = + CuratedImage( + iid = record.get(IID), + name = record.get(NAME), + sourceRef = record.get(SOURCE_REF), + sourceDigest = record.get(SOURCE_DIGEST), + status = record.get(STATUS), + imageTag = pinnedRefOf(record.get(SOURCE_REF), record.get(SOURCE_DIGEST)).orNull, + attempt = record.get(ATTEMPT), + creationTime = record.get(CREATION_TIME).getTime, + updateTime = record.get(UPDATE_TIME).getTime + ) + + /** + * The image a computing unit should start from, or None if it cannot be started from. + * No ownership check: these are offered to every user by design. + */ + def readyImageFor(iid: Int): Option[String] = { + // A disabled deployment starts nothing, including from a row left behind by an + // earlier enabled run. + if (!CuratedImageConfig.enabled) return None + val record = Option( + context.select(STATUS, SOURCE_REF, SOURCE_DIGEST).from(CU_IMAGE).where(IID.eq(iid)).fetchOne() + ) + record.flatMap { r => + if (r.get(STATUS) != Status.Ready) None + else pinnedRefOf(r.get(SOURCE_REF), r.get(SOURCE_DIGEST)) + } + } + + def nameOf(iid: Int): Option[String] = + Option(context.select(NAME).from(CU_IMAGE).where(IID.eq(iid)).fetchOne()).map(_.get(NAME)) + + /** + * Brings VALIDATING rows up to date with what the cluster did. Validation finishes on the + * cluster, so a row learns its outcome when someone reads it -- no background threads or + * leader election, at the cost of a status that is stale until that read. + */ + private def reconcileRunningValidations(): Unit = { + val running = context + .select(IID, ATTEMPT, UPDATE_TIME, SOURCE_REF) + .from(CU_IMAGE) + .where(STATUS.eq(Status.Validating)) + .fetch() + .asScala + .toList + + running.foreach { row => + val iid = row.get(IID).intValue() + val attempt = row.get(ATTEMPT).intValue() + // The cluster may be unreachable, or the Role not yet reapplied after an upgrade. + // Reconciling is opportunistic, so a row that cannot be checked is left as it is + // rather than failing a read that would otherwise return every other image. + try reconcileOne(iid, attempt, row.get(UPDATE_TIME).getTime) + catch { + case e: Throwable => + logger.warn(s"Could not check the validation of image $iid; leaving it as it is.", e) + } + } + } + + private def reconcileOne(iid: Int, attempt: Int, updatedAt: Long): Unit = { + // State first, then the log. The other order can read a log written while the job was + // still running and then judge it against a state that says it finished -- the digest + // line would be missing and a successful validation would be recorded as failed. + val state = ImageValidationClient.validationState(iid, attempt) + val log = ImageValidationClient.validationLog(iid, attempt) + + state match { + case ValidationState.Running => + // Kept fresh so the log can be watched while the copy is still going. + log.foreach(text => updateLogOnly(iid, attempt, text)) + + // Only the log says what a successful job resolved. If it cannot be read this time -- + // the pod not listed yet, or the read failing -- the row is left as it is and tried + // again on the next read. Recording FAILED on a guess would also delete the job, and + // with it the only account of what really happened. + case ValidationState.Succeeded if log.isEmpty => + logger.warn(s"Validation $iid/$attempt succeeded but its log could not be read yet.") + + case ValidationState.Succeeded => + val text = log.getOrElse("") + val digest = ImageValidationClient.sourceDigestFrom(text) + finishValidation( + iid, + attempt, + // Without a digest there is nothing to pin, so the image cannot be started + // from and calling it ready would strand a unit in ImagePullBackOff. + if (digest.isDefined) Status.Ready else Status.Failed, + digest, + text + validationNote(digest, digest.flatMap(sameContentNote(iid, _))) + ) + + case ValidationState.Failed => + // A job stopped by activeDeadlineSeconds has its pods removed by the Job + // controller, so there is no log left to explain it. The Job's own condition is + // the only account of what happened. + val reason = + log.filter(_.nonEmpty).orElse(ImageValidationClient.failureReason(iid, attempt)) + finishValidation( + iid, + attempt, + Status.Failed, + None, + reason.getOrElse("The validation failed without reporting a reason.") + ) + + case ValidationState.Absent => + // The job is created just after the row is marked VALIDATING and the two are not + // atomic, so a validation submitted moments ago legitimately has no job yet. Only a + // row that has been waiting a while is genuinely orphaned. + val age = System.currentTimeMillis() - updatedAt + if (age > AbsentGracePeriodMillis) { + finishValidation( + iid, + attempt, + Status.Failed, + None, + "The validation job disappeared before it reported a result." + ) + } + } + } + + private val AbsentGracePeriodMillis = 60_000L + + private[service] def validate(request: CuratedImageRequest): Unit = { + val name = Option(request.name).map(_.trim).getOrElse("") + if (!isValidName(name)) { + throw new BadRequestException( + "Image name must start with a letter or digit and contain only letters, digits, " + + "spaces, dots, hyphens, underscores, parentheses and plus signs, and be at " + + s"most $MaxNameLength characters." + ) + } + val ref = Option(request.sourceRef).map(_.trim).getOrElse("") + if (ref.isEmpty) { + throw new BadRequestException("Docker Hub link cannot be empty.") + } + // Measured on the normalised reference, because that is the one stored: an untagged + // reference grows by ":latest" and would otherwise overflow the column. + if (normaliseRef(ref).length > MaxRefLength) { + throw new BadRequestException(s"Docker Hub link exceeds $MaxRefLength characters.") + } + // Whitespace means two things were pasted; validation would fail less clearly. + if (normaliseRef(ref).exists(_.isWhitespace)) { + throw new BadRequestException("Docker Hub link cannot contain spaces.") + } + // An unrecognised hub.docker.com address would be pulled as if hub.docker.com were a + // registry; skopeo gets HTML back and reports it verbatim. Say so now instead. + if (normaliseRef(ref).startsWith(DockerHubHost + "/")) { + throw new BadRequestException( + s"'$ref' is a Docker Hub page, not an image reference, and its shape is not one " + + "this recognises. Use the image's own reference -- for example " + + "'acme/texera-cu-sklearn:1.0' -- or the address of its repository page." + ) + } + if (!RefPattern.pattern.matcher(normaliseRef(ref)).matches()) { + throw new BadRequestException( + "Docker Hub link must start with a letter or digit and contain only letters, " + + "digits, dots, hyphens, underscores, slashes, colons, at signs and plus signs." + ) + } + if (!RepoPattern.pattern.matcher(repositoryOf(normaliseRef(ref))).matches()) { + throw new BadRequestException( + s"'$ref' has an uppercase letter in its repository name. Registries only accept " + + "lowercase there -- a tag after the colon may have capitals, the part before it " + + "may not." + ) + } + } + + /** + * What a finished validation adds to its log: nothing in the ordinary case, the duplicate + * note when one applies, and the cannot-be-pinned explanation only when no digest was + * resolved. + */ + private[service] def validationNote( + digest: Option[String], + duplicateNote: Option[String] + ): String = + if (digest.isEmpty) "\n\nThe image could not be pinned: no digest was resolved." + else duplicateNote.getOrElse("") + + /** + * Logs when another row resolved to this same digest. Only knowable after validation, + * so it is a note rather than a rejection. + */ + private def sameContentNote(iid: Int, digest: String): Option[String] = { + val others = context + .select(NAME) + .from(CU_IMAGE) + .where(SOURCE_DIGEST.eq(digest)) + .and(IID.ne(iid)) + .and(STATUS.eq(Status.Ready)) + .fetch() + .asScala + .map(_.get(NAME)) + .toList + + Option.when(others.nonEmpty)( + s"\n\nNote: this is the same image as ${others.map("'" + _ + "'").mkString(", ")} " + + "-- same digest, reached by a different reference. The registry stores the layers " + + "once, so the duplicate costs little space, but only one of these rows is needed." + ) + } + + private def updateLogOnly(iid: Int, attempt: Int, log: String): Unit = + context + .update(CU_IMAGE) + .set(VALIDATION_LOG, log) + .where(IID.eq(iid).and(ATTEMPT.eq(attempt))) + .execute() + + /** + * Records an outcome, but only against the attempt that produced it. A refresh running + * at the same time moves the row to the next attempt, and without this guard the older + * validation's result would land on it and the refresh would never be polled again. + */ + private def finishValidation( + iid: Int, + attempt: Int, + status: String, + sourceDigest: Option[String], + log: String + ): Unit = { + val update = context + .update(CU_IMAGE) + .set(STATUS, status) + .set(VALIDATION_LOG, log) + .set(UPDATE_TIME, new Timestamp(System.currentTimeMillis())) + + // Leaves the previous digest alone, so a unit on the last good one keeps working. + val withDigest = sourceDigest.fold(update)(digest => update.set(SOURCE_DIGEST, digest)) + val stored = withDigest.where(IID.eq(iid).and(ATTEMPT.eq(attempt))).execute() + // The job is kept only until its outcome is recorded, so finished jobs do not pile up + // in the pool namespace. + if (stored > 0) ImageValidationClient.deleteValidation(iid, attempt) + } +} + +@Path("/cu-image") +@Produces(Array(MediaType.APPLICATION_JSON)) +class CuratedImageResource extends LazyLogging { + + import CuratedImageResource._ + + private def requireEnabled(): Unit = + if (!CuratedImageConfig.enabled) { + throw new ServiceUnavailableException("Curated images are not enabled on this deployment.") + } + + /** + * Any signed-in user may read the list, since the computing-unit dropdown is built from + * it. Only an administrator may change it -- that restriction is what makes these images + * trusted. + */ + @GET + @RolesAllowed(Array("REGULAR", "ADMIN")) + @Path("") + def list(@Auth user: SessionUser): List[CuratedImage] = { + requireEnabled() + reconcileRunningValidations() + context + // Not select(): validation_log is unbounded and this endpoint discards it. + .select(IID, NAME, SOURCE_REF, SOURCE_DIGEST, STATUS, ATTEMPT, CREATION_TIME, UPDATE_TIME) + .from(CU_IMAGE) + .orderBy(NAME.asc()) + .fetch() + .asScala + .map(toCuratedImage) + .toList + } + + @GET + @RolesAllowed(Array("ADMIN")) + @Path("/{iid}/log") + def log(@PathParam("iid") iid: Int, @Auth user: SessionUser): ValidationLog = { + requireEnabled() + reconcileRunningValidations() + val record = Option( + context.select(STATUS, ATTEMPT, VALIDATION_LOG).from(CU_IMAGE).where(IID.eq(iid)).fetchOne() + ).getOrElse(throw new NotFoundException(s"No curated image $iid.")) + + ValidationLog( + iid = iid, + status = record.get(STATUS), + attempt = record.get(ATTEMPT).intValue(), + log = Option(record.get(VALIDATION_LOG)).getOrElse("") + ) + } + + @POST + @RolesAllowed(Array("ADMIN")) + @Consumes(Array(MediaType.APPLICATION_JSON)) + @Path("") + def create(request: CuratedImageRequest, @Auth user: SessionUser): CuratedImage = { + requireEnabled() + validate(request) + val name = request.name.trim + val sourceRef = normaliseRef(request.sourceRef) + + if (context.fetchExists(context.selectFrom(CU_IMAGE).where(NAME.eq(name)))) { + throw new BadRequestException(s"An image named '$name' already exists.") + } + + // Naming the existing row lets the administrator use or refresh it, instead of ending up + // with two rows for one image. + val duplicate = Option( + context.select(NAME).from(CU_IMAGE).where(SOURCE_REF.eq(sourceRef)).fetchAny() + ).map(_.get(NAME)) + duplicate.foreach { existing => + throw new BadRequestException( + s"'$sourceRef' is already curated as '$existing'. Use that image, or refresh it " + + "to pick up a moved tag, instead of registering the same reference twice." + ) + } + + // The checks above are a read then a write, so two simultaneous registrations both pass + // them; the unique constraints refuse the second. Caught here to answer with the same + // explanation rather than a bare 500. + val iid = + try { + context + .insertInto(CU_IMAGE) + .set(NAME, name) + .set(SOURCE_REF, sourceRef) + .set(STATUS, Status.Pending) + .set(CREATED_BY, Integer.valueOf(user.getUid.intValue())) + .returning(IID) + .fetchOne() + .get(IID) + .intValue() + } catch { + case _: org.jooq.exception.IntegrityConstraintViolationException => + throw new BadRequestException( + s"'$sourceRef' or the name '$name' was registered a moment ago by someone " + + "else. Reload the list -- the image is already there." + ) + } + + startValidation(iid, sourceRef) + fetch(iid) + } + + /** + * Validates the source again: how a deployment picks up a moved tag, and how a validation + * that failed on the network is retried. + */ + @POST + @RolesAllowed(Array("ADMIN")) + @Path("/{iid}/refresh") + def refresh(@PathParam("iid") iid: Int, @Auth user: SessionUser): CuratedImage = { + requireEnabled() + val record = Option( + context.select(SOURCE_REF).from(CU_IMAGE).where(IID.eq(iid)).fetchOne() + ).getOrElse(throw new NotFoundException(s"No curated image $iid.")) + + startValidation(iid, record.get(SOURCE_REF)) + fetch(iid) + } + + @DELETE + @RolesAllowed(Array("ADMIN")) + @Path("/{iid}") + def delete(@PathParam("iid") iid: Int, @Auth user: SessionUser): Unit = { + requireEnabled() + // Nothing of ours holds a copy. A running unit keeps going on what its node pulled. + ImageValidationClient.deleteAllValidations(iid) + val deleted = context.deleteFrom(CU_IMAGE).where(IID.eq(iid)).execute() + if (deleted == 0) { + throw new NotFoundException(s"No curated image $iid.") + } + } + + /** Marks the row as being validated and submits the job, in that order. */ + private def startValidation(iid: Int, sourceRef: String): Unit = { + // Read and claimed in one statement. Two refreshes at the same moment would otherwise + // both compute the same attempt, and the second would delete the first's job. + val attempt = Option( + context + .update(CU_IMAGE) + .set(STATUS, Status.Validating) + .set(ATTEMPT, ATTEMPT.plus(1)) + .set(VALIDATION_LOG, "") + .set(UPDATE_TIME, new Timestamp(System.currentTimeMillis())) + .where(IID.eq(iid)) + .returningResult(ATTEMPT) + .fetchOne() + ).map(_.value1().intValue()) + .getOrElse(throw new NotFoundException(s"No curated image $iid.")) + + try { + ImageValidationClient.startValidation(iid, attempt, sourceRef) + } catch { + // Without this the row would sit in VALIDATING waiting for a job that was never + // created, and only the grace period would eventually call it failed. + case e: Throwable => + logger.error(s"Could not start validation for image $iid", e) + context + .update(CU_IMAGE) + .set(STATUS, Status.Failed) + .set(VALIDATION_LOG, ImageValidationClient.describeStartFailure(e)) + .where(IID.eq(iid).and(ATTEMPT.eq(attempt))) + .execute() + } + } + + private def fetch(iid: Int): CuratedImage = + Option(context.select().from(CU_IMAGE).where(IID.eq(iid)).fetchOne()) + .map(toCuratedImage) + .getOrElse(throw new NotFoundException(s"No curated image $iid.")) +} diff --git a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/ImageValidationClient.scala b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/ImageValidationClient.scala new file mode 100644 index 0000000000..ebe70bcf3e --- /dev/null +++ b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/ImageValidationClient.scala @@ -0,0 +1,444 @@ +/* + * 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 io.fabric8.kubernetes.api.model._ +import io.fabric8.kubernetes.api.model.batch.v1.{Job, JobBuilder} +import io.fabric8.kubernetes.client.KubernetesClientBuilder +import org.apache.texera.common.config.{CuratedImageConfig, KubernetesConfig} + +import scala.jdk.CollectionConverters._ + +/** + * Checks that a registered image looks like one a computing unit can start from, and + * resolves the digest behind the reference an administrator gave. + * + * "Looks like" is the honest word: the check is that the start command names + * `computing-unit-master`, which an unrelated image could also do. It catches the ordinary + * mistake -- the wrong image pasted in -- rather than proving provenance. Registration is + * admin-only, so that is the bar it needs to clear. + * + * The digest is resolved first and everything else is checked against it, so the image + * that was approved is exactly the one a unit is pinned to -- a tag moved midway through + * cannot slip a different image past the check. + * + * Only manifests and config blobs are read, a few kilobytes, never the layers. A reference + * that is misspelled, private, or not a computing-unit image is refused in seconds, in + * front of the administrator who typed it rather than the user whose unit would not + * start. + * + * skopeo rather than a pull: it reads the registry directly, with no Docker socket and no + * privileged pod. + */ +object ImageValidationClient extends LazyLogging { + + private val client: io.fabric8.kubernetes.client.KubernetesClient = + new KubernetesClientBuilder().build() + + private def namespace: String = CuratedImageConfig.validationNamespace + + /** Printed by the job and read back out of its log, since a Job cannot return a value. */ + private val DigestMarker = "TEXERA_SOURCE_DIGEST=" + + /** + * Looks at Cmd and Entrypoint rather than the whole config: a match anywhere in the + * config would also be satisfied by an environment variable that merely mentions the + * name. + */ + private def validationScript: String = { + // The image's own user, checked only where the deployment forces a non-root pod: an + // image that expects root would be admitted and then fail to start the user's first + // unit, which is far too late to find out. + val rootCheck = + if (!KubernetesConfig.computingUnitRunAsNonRoot) "" + else + """ + ||IMAGE_USER=$(skopeo inspect --config --format '{{.Config.User}}' "docker://$PINNED") + ||echo "Runs as: ${IMAGE_USER:-root}" + ||case "${IMAGE_USER:-root}" in + || root|0|"") + || echo "" + || echo "ERROR: this image runs as root." + || echo "Computing units are started as a non-root user here, so it would be" + || echo "admitted now and then fail to start. Rebuild it with a USER line." + || exit 1 + || ;; + ||esac + |""".stripMargin.replace("||", "|") + + s"""set -eu + | + |echo "Inspecting $$SOURCE_REF" + | + |# The digest first, so everything after this is checked against the exact image that + |# will be pinned. Resolving it last would leave room for the tag to move in between, + |# and what was approved would not be what a unit runs. + |if ! DIGEST=$$(skopeo inspect --format '{{.Digest}}' "docker://$$SOURCE_REF" 2>&1); then + | echo "" + | echo "ERROR: could not read $$SOURCE_REF from its registry." + | echo "$$DIGEST" + | echo "" + | echo "If that says the manifest is unknown, the tag does not exist. A Docker Hub" + | echo "page address carries no tag, so ':latest' was assumed -- and many images do" + | echo "not publish one. Register the reference with the tag you want, for example" + | echo "'owner/name:1.0'." + | echo "If it mentions authorisation, the image is private; only public images can" + | echo "be used." + | exit 1 + |fi + | + |# The repository without its tag, matching what the service derives from the same + |# reference. A colon after the last slash is a tag; one before it is a registry port. + |case "$$SOURCE_REF" in + | *@sha256:*) REPO="$${SOURCE_REF%@*}" ;; + | *) case "$${SOURCE_REF##*/}" in + | *:*) REPO="$${SOURCE_REF%:*}" ;; + | *) REPO="$$SOURCE_REF" ;; + | esac ;; + |esac + |PINNED="$$REPO@$$DIGEST" + |echo "Pinned to: $$PINNED" + | + |if ! START_CMD=$$(skopeo inspect --config \\ + | --format '{{.Config.Cmd}} {{.Config.Entrypoint}}' \\ + | "docker://$$PINNED" 2>&1); then + | echo "" + | echo "ERROR: could not read the image at $$PINNED." + | echo "$$START_CMD" + | exit 1 + |fi + |echo "Start command: $$START_CMD" + | + |if ! echo "$$START_CMD" | grep -qF '${CuratedImageConfig.requiredCommand}'; then + | echo "" + | echo "ERROR: $$SOURCE_REF does not look like a Texera computing-unit image." + | echo "Its start command is: $$START_CMD" + | echo "A computing-unit image starts '${CuratedImageConfig.requiredCommand}'." + | exit 1 + |fi + |$rootCheck + |echo "$DigestMarker$$DIGEST" + |""".stripMargin + } + + def startValidation(iid: Int, attempt: Int, sourceRef: String): Unit = { + val jobName = CuratedImageConfig.validationJobName(iid, attempt) + + // Only attempts below this one. Two refreshes claim their numbers atomically but reach + // the cluster in any order, so sweeping every job of this image would let the older of + // the two delete the newer one's job and strand the row waiting for a job that is gone. + deleteSupersededValidations(iid, attempt) + + val job = validationJob(jobName, iid, attempt, sourceRef) + client.batch().v1().jobs().inNamespace(namespace).resource(job).create() + logger.info(s"Started validation $jobName for $sourceRef") + } + + private def validationJob(jobName: String, iid: Int, attempt: Int, sourceRef: String): Job = { + val resources = new ResourceRequirementsBuilder() + .addToRequests("cpu", new Quantity(CuratedImageConfig.validationCpuRequest)) + .addToRequests("memory", new Quantity(CuratedImageConfig.validationMemoryRequest)) + .addToLimits("cpu", new Quantity(CuratedImageConfig.validationCpuLimit)) + .addToLimits("memory", new Quantity(CuratedImageConfig.validationMemoryLimit)) + .build() + + new JobBuilder() + .withNewMetadata() + .withName(jobName) + .withNamespace(namespace) + .addToLabels("texera-cu-image", iid.toString) + .addToLabels("texera-cu-image-attempt", attempt.toString) + .endMetadata() + .withNewSpec() + // A rejected image is rejected deterministically, and a network failure is better + // retried by an administrator who can see why. + .withBackoffLimit(0) + .withActiveDeadlineSeconds(CuratedImageConfig.validationTimeoutSeconds.toLong) + .withNewTemplate() + .withNewMetadata() + .addToLabels("texera-cu-image", iid.toString) + .addToLabels("texera-cu-image-attempt", attempt.toString) + .endMetadata() + .withNewSpec() + .withRestartPolicy("Never") + .addNewContainer() + .withName("skopeo") + .withImage(CuratedImageConfig.validationImage) + .withCommand("/bin/sh", "-c") + .withArgs(validationScript) + // Passed as a value, not spliced into the script: the shell expands it but never + // parses it, so a reference cannot carry commands of its own. + .addNewEnv() + .withName("SOURCE_REF") + .withValue(sourceRef) + .endEnv() + .withResources(resources) + .endContainer() + .endSpec() + .endTemplate() + .endSpec() + .build() + } + + sealed trait ValidationState + object ValidationState { + case object Running extends ValidationState + case object Succeeded extends ValidationState + case object Failed extends ValidationState + + /** The job is gone -- cleaned up, or never created. */ + case object Absent extends ValidationState + } + + def validationState(iid: Int, attempt: Int): ValidationState = { + val job = Option( + client + .batch() + .v1() + .jobs() + .inNamespace(namespace) + .withName(CuratedImageConfig.validationJobName(iid, attempt)) + .get() + ) + + stateOf(job) + } + + /** + * What a Job's status says about its validation. A job stopped by its deadline counts as + * failed here, because the Job controller records that as a failure. + */ + private[service] def stateOf(job: Option[Job]): ValidationState = + job match { + case None => ValidationState.Absent + case Some(j) => + val status = Option(j.getStatus) + val succeeded = status.flatMap(s => Option(s.getSucceeded)).exists(_ > 0) + val failed = status.flatMap(s => Option(s.getFailed)).exists(_ > 0) + if (succeeded) ValidationState.Succeeded + else if (failed) ValidationState.Failed + else ValidationState.Running + } + + /** The job's output, readable while it runs and after it finishes. */ + def validationLog(iid: Int, attempt: Int): Option[String] = { + val pods = client + .pods() + .inNamespace(namespace) + .withLabel("job-name", CuratedImageConfig.validationJobName(iid, attempt)) + .list() + .getItems + .asScala + .toList + + pods.headOption.flatMap { pod => + try { + Option(client.pods().inNamespace(namespace).withName(pod.getMetadata.getName).getLog(true)) + } catch { + // A container that has not started has no log -- ordinary, not an error. + case e: Throwable => + logger.debug(s"No log yet for validation $iid/$attempt: ${e.getMessage}") + None + } + } + } + + /** + * Names the API address in the error. A client with no usable kube config falls back to + * http://localhost:8080, and whatever answers there fails opaquely -- which reads as a + * Texera bug rather than a missing cluster. + */ + def describeStartFailure(e: Throwable): String = { + val cause = Option(e.getCause).filter(_ ne e) + val detail = Option(e.getMessage) + .map(_.trim) + .filter(_.nonEmpty) + .getOrElse("no message") + val master = + try Option(client.getMasterUrl).map(_.toString).getOrElse("unknown") + catch { case _: Throwable => "unknown" } + + s"""Could not start the validation job. + | + | ${e.getClass.getSimpleName}: $detail${cause + .map(c => s"\n caused by ${c.getClass.getSimpleName}: ${Option(c.getMessage).getOrElse("")}") + .getOrElse("")} + | + |Kubernetes API address: $master + |Namespace: $namespace + | + |Validation runs as a Kubernetes Job, so this needs a reachable cluster. If the + |address above is http://localhost:8080 then no kube context is set and the client + |fell back to that default -- check `kubectl config current-context`, and that the + |namespace above exists.""".stripMargin + } + + /** + * Why a job failed, taken from the Job itself. The pods of a job stopped by its deadline + * are removed by the Job controller, so its condition is all that is left to explain it. + */ + def failureReason(iid: Int, attempt: Int): Option[String] = + try { + Option( + client + .batch() + .v1() + .jobs() + .inNamespace(namespace) + .withName(CuratedImageConfig.validationJobName(iid, attempt)) + .get() + ).flatMap(failureReasonOf) + } catch { + case e: Throwable => + logger.debug(s"Could not read why validation $iid/$attempt failed: ${e.getMessage}") + None + } + + /** + * Why a Job says it failed. A deadline is spelled out, because the pods are gone by then + * and this is the only account the administrator will get. + */ + private[service] def failureReasonOf(job: Job): Option[String] = + Option(job.getStatus) + .flatMap(s => Option(s.getConditions)) + .map(_.asScala.toList) + .getOrElse(Nil) + .find(c => c.getType == "Failed") + .flatMap { c => + if (c.getReason == "DeadlineExceeded") + Some( + s"The validation gave up after ${CuratedImageConfig.validationTimeoutSeconds} " + + "seconds. The registry did not answer in time; try again, and check the " + + "reference is one this cluster can reach." + ) + else + Option(c.getMessage).filter(_.nonEmpty).orElse(Option(c.getReason).filter(_.nonEmpty)) + } + + /** The digest the source tag resolved to, as printed by a successful job. */ + def sourceDigestFrom(log: String): Option[String] = + log.linesIterator + .map(_.trim) + .find(_.startsWith(DigestMarker)) + .map(_.drop(DigestMarker.length).trim) + .filter(_.nonEmpty) + + /** + * The reference a unit starts from: the administrator's repository at the resolved + * digest. A tag can be moved by its owner, so pinning keeps the unit on the bytes that + * were approved. A reference that already names a digest is returned unchanged. + */ + private[service] def pinnedRef(sourceRef: String, digest: String): String = { + val reference = Option(sourceRef).map(_.trim).getOrElse("") + if (reference.contains("@sha256:")) reference + else { + // Strip the tag if there is one. A colon after the last slash is a tag; one before + // it belongs to a registry's port. + val lastSlash = reference.lastIndexOf('/') + val colon = reference.indexOf(':', lastSlash + 1) + val repository = if (colon >= 0) reference.substring(0, colon) else reference + s"$repository@$digest" + } + } + + /** + * Removes one validation's job once its outcome has been recorded. Safe to call when it + * does not exist. Not a TTL on the Job: outcomes are read when someone looks at the list, + * so a job reaped on a timer could vanish before it was ever read. + */ + def deleteValidation(iid: Int, attempt: Int): Unit = { + val jobName = CuratedImageConfig.validationJobName(iid, attempt) + try { + client + .batch() + .v1() + .jobs() + .inNamespace(namespace) + .withName(jobName) + .withPropagationPolicy(DeletionPropagation.BACKGROUND) + .delete() + } catch { + case e: Throwable => logger.warn(s"Could not clean up validation $jobName: ${e.getMessage}") + } + } + + /** Which of an image's jobs belong to an attempt this one has replaced. */ + private[service] def supersededJobs(jobs: List[Job], attempt: Int): List[String] = + jobs.filter(j => attemptOf(j).exists(_ < attempt)).flatMap(j => Option(j.getMetadata.getName)) + + /** + * The attempt a job belongs to, from its label, falling back to the trailing number of + * its name for a job created before the label existed. + */ + private[service] def attemptOf(job: Job): Option[Int] = { + val labelled = Option(job.getMetadata) + .flatMap(m => Option(m.getLabels)) + .flatMap(l => Option(l.get("texera-cu-image-attempt"))) + val named = Option(job.getMetadata).flatMap(m => Option(m.getName)).map(_.split('-').last) + labelled.orElse(named).flatMap(v => scala.util.Try(v.toInt).toOption) + } + + /** Removes the jobs of attempts this one has replaced. */ + def deleteSupersededValidations(iid: Int, attempt: Int): Unit = { + try { + val jobs = client + .batch() + .v1() + .jobs() + .inNamespace(namespace) + .withLabel("texera-cu-image", iid.toString) + .list() + .getItems + .asScala + .toList + supersededJobs(jobs, attempt).foreach { name => + client + .batch() + .v1() + .jobs() + .inNamespace(namespace) + .withName(name) + .withPropagationPolicy(DeletionPropagation.BACKGROUND) + .delete() + } + } catch { + case e: Throwable => + logger.warn(s"Could not clean up superseded validations for image $iid: ${e.getMessage}") + } + } + + /** Removes every validation job belonging to an image, superseded or not. */ + def deleteAllValidations(iid: Int): Unit = { + try { + client + .batch() + .v1() + .jobs() + .inNamespace(namespace) + .withLabel("texera-cu-image", iid.toString) + .withPropagationPolicy(DeletionPropagation.BACKGROUND) + .delete() + } catch { + case e: Throwable => + logger.warn(s"Could not clean up validation jobs for image $iid: ${e.getMessage}") + } + } +} 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 18328f1f9f..856a423484 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 @@ -126,7 +126,9 @@ class KubernetesClient( memoryLimit: String, gpuLimit: String, envVars: Map[String, Any], - shmSize: Option[String] = None + shmSize: Option[String] = None, + /** A curated image to run instead of the deployment's own. */ + curatedImage: Option[String] = None ): Pod = { val podName = generatePodName(cuid) if (getPodByName(podName).isDefined) { @@ -192,7 +194,7 @@ class KubernetesClient( val containerBuilder = specBuilder .addNewContainer() .withName("computing-unit-master") - .withImage(KubernetesConfig.computeUnitImageName) + .withImage(curatedImage.getOrElse(KubernetesConfig.computeUnitImageName)) .withImagePullPolicy(KubernetesConfig.computingUnitImagePullPolicy) .addNewPort() .withContainerPort(KubernetesConfig.computeUnitPortNumber) @@ -200,6 +202,23 @@ class KubernetesClient( .withEnv(envList) .withResources(resourceBuilder.build()) + // A curated image was supplied by an administrator and reviewed by nobody, so a unit + // started from one is pinned to a non-root user with no way to regain privilege. + // + // Curated images only. The deployment's own image is its operator's choice, and one + // that has replaced it with an image needing root would break on upgrade. + if (curatedImage.isDefined && KubernetesConfig.computingUnitRunAsNonRoot) { + containerBuilder + .withNewSecurityContext() + .withRunAsNonRoot(true) + .withRunAsUser(KubernetesConfig.computingUnitRunAsUser) + .withAllowPrivilegeEscalation(false) + .withNewCapabilities() + .withDrop("ALL") + .endCapabilities() + .endSecurityContext() + } + // 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). diff --git a/computing-unit-managing-service/src/test/scala/org/apache/texera/service/resource/CuratedImageResourceSpec.scala b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/resource/CuratedImageResourceSpec.scala new file mode 100644 index 0000000000..4858a44b11 --- /dev/null +++ b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/resource/CuratedImageResourceSpec.scala @@ -0,0 +1,459 @@ +/* + * 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 org.apache.texera.common.config.CuratedImageConfig +import io.fabric8.kubernetes.api.model.batch.v1.{Job, JobBuilder, JobCondition, JobConditionBuilder} +import org.apache.texera.service.util.ImageValidationClient +import org.apache.texera.service.util.ImageValidationClient.ValidationState +import org.scalatest.OptionValues._ + +import scala.jdk.CollectionConverters._ +import org.scalatest.flatspec.AnyFlatSpec +import jakarta.ws.rs.BadRequestException +import org.scalatest.matchers.should.Matchers + +class CuratedImageResourceSpec extends AnyFlatSpec with Matchers { + + import CuratedImageResource.{isValidName, normaliseRef} + + // The regression this guards: the pattern used to exclude parentheses, so the names an + // administrator actually reaches for -- and the ones the demo instructions themselves + // used -- were rejected with HTTP 400 before anything was created. + "isValidName" should "accept the names administrators actually type" in { + isValidName("Texera Default") shouldBe true + isValidName("Python ML (sklearn)") shouldBe true + isValidName("PyTorch 2.6 (CUDA 12)") shouldBe true + isValidName("cu-image_v1.0") shouldBe true + isValidName("gcc+cuda") shouldBe true + } + + it should "still refuse anything that is not a plain display name" in { + isValidName("") shouldBe false + isValidName(" ") shouldBe false + // must start alphanumeric, so no leading punctuation or whitespace-only leaders + isValidName("-leading-hyphen") shouldBe false + isValidName("(leading-paren)") shouldBe false + // no quoting, markup, path or shell metacharacters + isValidName("name\"quote") shouldBe false + isValidName("<script>") shouldBe false + isValidName("a/b") shouldBe false + isValidName("a;rm -rf") shouldBe false + isValidName("a\nb") shouldBe false + isValidName("caf\u00e9") shouldBe false + } + + it should "reject a name longer than the column allows" in { + isValidName("a" * 128) shouldBe true + isValidName("a" * 129) shouldBe false + } + + "normaliseRef" should "leave a complete image reference alone" in { + normaliseRef("texera/cu-alphafold3:1.0") shouldBe "texera/cu-alphafold3:1.0" + } + + it should "default a missing tag rather than reject it" in { + normaliseRef("texera/cu-alphafold3") shouldBe "texera/cu-alphafold3:latest" + } + + // The whole reason this function exists: an administrator curating from a browser will + // paste the page they are looking at, not a reference they had to construct. + it should "turn a Docker Hub page address into a pull reference" in { + normaliseRef("https://hub.docker.com/r/texera/cu-alphafold3") shouldBe + "texera/cu-alphafold3:latest" + normaliseRef("hub.docker.com/r/texera/cu-alphafold3") shouldBe "texera/cu-alphafold3:latest" + } + + // The address bar shows this while an owner manages their own image, so it is the one + // an administrator curating their own build is most likely to paste. + it should "turn an owner's own repository page address into a pull reference" in { + normaliseRef("https://hub.docker.com/repository/docker/acme/texera-cu-sklearn") shouldBe + "acme/texera-cu-sklearn:latest" + } + + // /general, /tags and /settings are parts of the web page, not of the reference. + it should "drop the page's tab segment" in { + normaliseRef( + "https://hub.docker.com/repository/docker/acme/texera-cu-sklearn/general" + ) shouldBe + "acme/texera-cu-sklearn:latest" + normaliseRef( + "https://hub.docker.com/repository/docker/acme/texera-cu-sklearn/tags" + ) shouldBe + "acme/texera-cu-sklearn:latest" + normaliseRef("https://hub.docker.com/r/acme/texera-cu-sklearn/tags") shouldBe + "acme/texera-cu-sklearn:latest" + } + + it should "leave a Docker Hub address it does not recognise alone, for validate to reject" in { + // Better an immediate rejection than a validation that discovers hub.docker.com is not + // a registry and prints the 404 page it got back. + normaliseRef("https://hub.docker.com/u/acme") should startWith("hub.docker.com/") + } + + it should "tolerate a trailing slash, which a copied address usually has" in { + normaliseRef("https://hub.docker.com/r/texera/cu-alphafold3/") shouldBe + "texera/cu-alphafold3:latest" + } + + // An official image lives under /_/ and is pulled by its bare name, so the path prefix + // has to come off or the reference would name a repository that does not exist. + it should "handle a Docker Hub official image" in { + normaliseRef("https://hub.docker.com/_/ubuntu") shouldBe "ubuntu:latest" + } + + it should "trim surrounding whitespace" in { + normaliseRef(" texera/img:1.0 ") shouldBe "texera/img:1.0" + } + + // The regression this guards: a registry's port contains a colon, and looking for one + // anywhere in the reference would read ":5000/team/img" as a tag and leave the image + // untagged. + it should "not mistake a registry port for a tag" in { + normaliseRef("myregistry.io:5000/team/img") shouldBe "myregistry.io:5000/team/img:latest" + normaliseRef("10.96.0.99:5000/texera/computing-unit-master:dev") shouldBe + "10.96.0.99:5000/texera/computing-unit-master:dev" + } + + // Docker Hub is what a bare reference already means, so both forms have to reduce to + // the same string. Otherwise one image registered as "owner/name:1" and again as + // "docker.io/owner/name:1" is curated twice -- the duplicate check compares references, + // so it can only catch what normalisation made equal. + it should "reduce an explicit Docker Hub registry to the bare reference" in { + normaliseRef("docker.io/acme/texera-cu-sklearn:1.0") shouldBe + "acme/texera-cu-sklearn:1.0" + normaliseRef("index.docker.io/acme/texera-cu-sklearn:1.0") shouldBe + "acme/texera-cu-sklearn:1.0" + normaliseRef("registry-1.docker.io/acme/texera-cu-sklearn:1.0") shouldBe + "acme/texera-cu-sklearn:1.0" + } + + // "library/" is Docker Hub's namespace for official images, whose reference is the name. + it should "reduce an official image's fully qualified reference" in { + normaliseRef("docker.io/library/ubuntu:22.04") shouldBe "ubuntu:22.04" + } + + // A real registry that merely starts with similar text must be left alone. + it should "not mistake another registry for Docker Hub" in { + normaliseRef("docker.io.evil.example/team/img:1") shouldBe "docker.io.evil.example/team/img:1" + normaliseRef("ghcr.io/apache/texera:latest") shouldBe "ghcr.io/apache/texera:latest" + } + + it should "leave a digest-pinned reference untagged" in { + normaliseRef("texera/img@sha256:abc123") shouldBe "texera/img@sha256:abc123" + } + + // A unit is started from the administrator's repository at the digest validation + // resolved, so the tag they typed has to come off first. + // The ordinary success -- a digest resolved, no other row sharing it -- must add nothing. + // The regression this guards said "could not be pinned" on every such READY image. + "validationNote" should "add nothing when a digest was resolved and is unique" in { + CuratedImageResource.validationNote(Some("sha256:abc"), None) shouldBe "" + } + + it should "add the duplicate note when another row is the same image" in { + CuratedImageResource.validationNote(Some("sha256:abc"), Some("\n\nSame as 'other'.")) shouldBe + "\n\nSame as 'other'." + } + + it should "explain the failure to pin only when no digest was resolved" in { + CuratedImageResource.validationNote(None, None) should include("could not be pinned") + } + + private def rejects(name: String, ref: String): String = + intercept[BadRequestException]( + CuratedImageResource.validate(CuratedImageResource.CuratedImageRequest(name, ref)) + ).getMessage + + // The reference reaches the validation job as an environment value, never as script + // text, so a shell metacharacter cannot run anything. It is still refused here so the + // administrator gets a clear message rather than a puzzling failure from skopeo. + "validate" should "refuse a reference carrying shell metacharacters" in { + rejects("ok", "acme/img$(id)") should include("letters") + rejects("ok", "acme/img\";id;\"") should include("letters") + rejects("ok", "acme/img`id`") should include("letters") + } + + it should "accept the reference shapes a registry actually uses" in { + noException should be thrownBy + CuratedImageResource.validate(CuratedImageResource.CuratedImageRequest("ok", "acme/img:1.0")) + noException should be thrownBy + CuratedImageResource.validate( + CuratedImageResource.CuratedImageRequest("ok", "registry.example:5000/team/img@sha256:abc") + ) + } + + // Measured after normalising, because ":latest" is added before the value is stored and + // the column is only so wide. + it should "measure the length of the reference it will store, not the one given" in { + val untagged = "acme/" + ("a" * 503) + untagged.length shouldBe 508 + rejects("ok", untagged) should include("exceeds") + } + + // image_tag used to be a stored column. It is now derived on read, so these guard that + // the derivation gives the same answer -- including the null-digest case the column had. + "pinnedRefOf" should "combine the registered reference with the resolved digest" in { + CuratedImageResource.pinnedRefOf("acme/img:1.0", "sha256:abc") shouldBe + Some("acme/img@sha256:abc") + } + + it should "have nothing to offer until a validation has resolved a digest" in { + CuratedImageResource.pinnedRefOf("acme/img:1.0", null) shouldBe None + CuratedImageResource.pinnedRefOf("acme/img:1.0", "") shouldBe None + } + + "pinnedRef" should "address the same repository by digest instead of by tag" in { + ImageValidationClient.pinnedRef("texera/img:1.0", "sha256:abc") shouldBe "texera/img@sha256:abc" + ImageValidationClient.pinnedRef("ghcr.io/apache/texera:latest", "sha256:def") shouldBe + "ghcr.io/apache/texera@sha256:def" + } + + it should "leave a reference that already names a digest alone" in { + ImageValidationClient.pinnedRef("texera/img@sha256:abc", "sha256:zzz") shouldBe + "texera/img@sha256:abc" + } + + // The same trap as everywhere else: a registry's port is a colon that is not a tag, and + // truncating there would pin a repository that does not exist. + it should "not mistake a registry port for a tag" in { + ImageValidationClient.pinnedRef("registry.example:5000/team/img:2", "sha256:abc") shouldBe + "registry.example:5000/team/img@sha256:abc" + ImageValidationClient.pinnedRef("registry.example:5000/team/img", "sha256:abc") shouldBe + "registry.example:5000/team/img@sha256:abc" + } + + it should "pin an untagged reference as it stands" in { + ImageValidationClient.pinnedRef("texera/img", "sha256:abc") shouldBe "texera/img@sha256:abc" + } + + "sourceDigestFrom" should "read the digest a finished validation printed" in { + val log = + """Inspecting texera/cu-alphafold3:1.0 + |Start command: [bin/computing-unit-master] [] + |TEXERA_SOURCE_DIGEST=sha256:0123abc + |Validated texera/cu-alphafold3:1.0 + |""".stripMargin + ImageValidationClient.sourceDigestFrom(log) shouldBe Some("sha256:0123abc") + } + + // The regression this guards: the bare exception message can be something as useless as + // "An error has occurred." when the client fell back to a default API address, which + // reads as a Texera bug rather than a missing cluster. + "describeStartFailure" should "name the failure type and where it was talking to" in { + val described = ImageValidationClient.describeStartFailure( + new RuntimeException("An error has occurred.") + ) + described should include("RuntimeException") + described should include("An error has occurred.") + described should include("Kubernetes API address") + described should include("reachable cluster") + } + + it should "still say something useful when the exception has no message" in { + val described = ImageValidationClient.describeStartFailure(new NullPointerException) + described should include("NullPointerException") + described should include("no message") + } + + it should "return nothing when validation failed before printing one" in { + val log = + """Inspecting texera/not-a-cu-image:1.0 + |Start command: [/bin/bash] [] + |ERROR: texera/not-a-cu-image:1.0 does not look like a Texera computing-unit image. + |""".stripMargin + ImageValidationClient.sourceDigestFrom(log) shouldBe None + } + + // Off until the UI ships, so a deployment that has not opted in starts no unit from a + // curated image -- including from a row left behind if it was enabled and turned off. + "the feature flag" should "be off unless a deployment turns it on" in { + CuratedImageConfig.enabled shouldBe false + } + + it should "start no unit from a curated image while it is off" in { + CuratedImageResource.readyImageFor(1) shouldBe None + } + + // The states below are the ones a real cluster produces; the DeadlineExceeded shape was + // taken from a job actually killed by activeDeadlineSeconds (failed=1, condition Failed + // with that reason, and no pods left to read a log from). + private def job(succeeded: Integer, failed: Integer, conditions: JobCondition*): Job = + new JobBuilder() + .withNewMetadata() + .withName("cu-image-check-1-1") + .endMetadata() + .withNewStatus() + .withSucceeded(succeeded) + .withFailed(failed) + .withConditions(conditions.toList.asJava) + .endStatus() + .build() + + private def condition(condType: String, reason: String, message: String): JobCondition = + new JobConditionBuilder() + .withType(condType) + .withReason(reason) + .withMessage(message) + .build() + + "stateOf" should "report a job that is not there as absent" in { + ImageValidationClient.stateOf(None) shouldBe ValidationState.Absent + } + + it should "report a job with no terminal count as still running" in { + ImageValidationClient.stateOf(Some(job(null, null))) shouldBe ValidationState.Running + } + + it should "report a succeeded job as succeeded" in { + ImageValidationClient.stateOf(Some(job(1, null))) shouldBe ValidationState.Succeeded + } + + it should "report a failed job as failed" in { + ImageValidationClient.stateOf(Some(job(null, 1))) shouldBe ValidationState.Failed + } + + // A job killed by its deadline reports failed=1, so the row is settled rather than left + // waiting for a result that will never come. + it should "treat a job stopped by its deadline as failed" in { + val deadline = + job( + null, + 1, + condition("Failed", "DeadlineExceeded", "Job was active longer than specified deadline") + ) + ImageValidationClient.stateOf(Some(deadline)) shouldBe ValidationState.Failed + } + + "failureReasonOf" should "explain a deadline in terms of the timeout that caused it" in { + val deadline = + job( + null, + 1, + condition("Failed", "DeadlineExceeded", "Job was active longer than specified deadline") + ) + val reason = ImageValidationClient.failureReasonOf(deadline) + reason.value should include(CuratedImageConfig.validationTimeoutSeconds.toString) + reason.value should include("did not answer in time") + } + + it should "pass on the message of any other failure" in { + val other = job( + null, + 1, + condition("Failed", "BackoffLimitExceeded", "Job has reached the specified backoff limit") + ) + ImageValidationClient.failureReasonOf(other).value should include("backoff limit") + } + + it should "fall back to the reason when the message is empty" in { + val bare = job(null, 1, condition("Failed", "BackoffLimitExceeded", "")) + ImageValidationClient.failureReasonOf(bare) shouldBe Some("BackoffLimitExceeded") + } + + // Nothing to say about a job that has not failed, so the caller keeps its own wording. + it should "have nothing to say about a job with no failure condition" in { + ImageValidationClient.failureReasonOf(job(1, null)) shouldBe None + ImageValidationClient.failureReasonOf( + job(null, 1, condition("Complete", "", "done")) + ) shouldBe None + } + + private def labelledJob(iid: Int, attempt: Int, labelled: Boolean = true): Job = { + val b = new JobBuilder().withNewMetadata().withName(s"cu-image-check-$iid-$attempt") + val withLabels = + if (labelled) + b.addToLabels("texera-cu-image", iid.toString) + .addToLabels("texera-cu-image-attempt", attempt.toString) + else b.addToLabels("texera-cu-image", iid.toString) + withLabels.endMetadata().build() + } + + // The race this guards: two refreshes claim 2 and 3 atomically but reach the cluster in + // any order. The one holding 2 must not delete the job of 3, or the row waits for a job + // that no longer exists. + "supersededJobs" should "select only attempts below the one starting" in { + val jobs = List(labelledJob(7, 1), labelledJob(7, 2), labelledJob(7, 3)) + ImageValidationClient.supersededJobs(jobs, 3) should contain theSameElementsAs + List("cu-image-check-7-1", "cu-image-check-7-2") + } + + it should "never select a newer attempt than the one starting" in { + val jobs = List(labelledJob(7, 2), labelledJob(7, 3)) + ImageValidationClient.supersededJobs(jobs, 2) shouldBe empty + } + + it should "not select the attempt that is starting" in { + ImageValidationClient.supersededJobs(List(labelledJob(7, 4)), 4) shouldBe empty + } + + // A job created before the attempt label existed still has to be reapable. + it should "fall back to the trailing number of the name when the label is absent" in { + ImageValidationClient.supersededJobs(List(labelledJob(7, 1, labelled = false)), 3) shouldBe + List("cu-image-check-7-1") + } + + "attemptOf" should "prefer the label over the name" in { + val odd = new JobBuilder() + .withNewMetadata() + .withName("cu-image-check-7-99") + .addToLabels("texera-cu-image-attempt", "4") + .endMetadata() + .build() + ImageValidationClient.attemptOf(odd) shouldBe Some(4) + } + + it should "have no answer for a name it cannot read a number from" in { + val odd = new JobBuilder().withNewMetadata().withName("something-else").endMetadata().build() + ImageValidationClient.attemptOf(odd) shouldBe None + } + + // Without a message or a reason there is nothing to report, and the caller's own wording + // must survive rather than being replaced by an empty string. + it should "leave a bare failure condition to the caller's fallback" in { + ImageValidationClient.failureReasonOf(job(null, 1, condition("Failed", "", ""))) shouldBe None + } + + // Registries reject an uppercase repository path, so it is caught here rather than in the + // job. A tag may have capitals; the part before the colon may not. + "validate" should "refuse an uppercase repository name" in { + rejects("ok", "Acme/img:1.0") should include("lowercase") + rejects("ok", "acme/MyImage:1.0") should include("lowercase") + } + + it should "allow capitals in a tag" in { + noException should be thrownBy + CuratedImageResource.validate( + CuratedImageResource.CuratedImageRequest("ok", "acme/img:V1.0-RC") + ) + } + + "repositoryOf" should "drop a tag but keep a registry port" in { + CuratedImageResource.repositoryOf("registry.example:5000/team/img:2") shouldBe + "registry.example:5000/team/img" + CuratedImageResource.repositoryOf("acme/img:1.0") shouldBe "acme/img" + CuratedImageResource.repositoryOf("acme/img") shouldBe "acme/img" + } + + it should "drop a digest" in { + CuratedImageResource.repositoryOf("acme/img@sha256:abc") shouldBe "acme/img" + } + +} diff --git a/sql/changelog.xml b/sql/changelog.xml index 677af284f0..cf83a257ad 100644 --- a/sql/changelog.xml +++ b/sql/changelog.xml @@ -154,6 +154,11 @@ <sqlFile path="sql/updates/48.sql"/> </changeSet> + <!-- Curated computing-unit images an administrator registers --> + <changeSet id="49" author="tanishqgandhi1908"> + <sqlFile path="sql/updates/49.sql"/> + </changeSet> + <!-- example changeSet <changeSet id="1" author="author"> <sqlFile path="sql/updates/1.sql"/> diff --git a/sql/texera_ddl.sql b/sql/texera_ddl.sql index f1890e6f78..7d5b148607 100644 --- a/sql/texera_ddl.sql +++ b/sql/texera_ddl.sql @@ -252,6 +252,34 @@ CREATE TABLE IF NOT EXISTS workflow_computing_unit FOREIGN KEY (uid) REFERENCES "user"(uid) ON DELETE CASCADE ); +-- does not restrict who may use the image. +CREATE TABLE IF NOT EXISTS cu_image +( + iid SERIAL PRIMARY KEY, + name VARCHAR(128) NOT NULL, + -- What the administrator supplied, normalised to an image reference. + source_ref VARCHAR(512) NOT NULL, + -- The digest source_ref resolved to when last validated; an upstream tag can move. + source_digest VARCHAR(128), + status VARCHAR(16) NOT NULL DEFAULT 'PENDING' + CONSTRAINT ck_cu_image_status + CHECK (status IN ('PENDING', 'VALIDATING', 'READY', 'FAILED')), + -- What a unit is started from: source_ref pinned to the digest above. + -- Counts validations of this row, so a retry gets a job name of its own. + attempt INT NOT NULL DEFAULT 0, + validation_log TEXT, + created_by INT, + creation_time TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + update_time TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + FOREIGN KEY (created_by) REFERENCES "user" (uid) ON DELETE SET NULL, + UNIQUE (name), + -- One row per reference; the service's own check is a read-then-write and so cannot + -- stop two simultaneous registrations of the same link. + UNIQUE (source_ref) +); + +CREATE INDEX IF NOT EXISTS idx_cu_image_source_digest ON cu_image (source_digest); + -- Per-user warehouse registrations (#6870): one row per warehouse a user registered. -- Base columns only; the assume-role (BYO-S3) columns come in a later change. CREATE TABLE IF NOT EXISTS user_warehouse diff --git a/sql/updates/49.sql b/sql/updates/49.sql new file mode 100644 index 0000000000..0a48bf464a --- /dev/null +++ b/sql/updates/49.sql @@ -0,0 +1,61 @@ +/* + * 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. + */ + +-- Computing-unit images an administrator has registered, which any user may then start a +-- computing unit from. Rows are global rather than owned -- the point is one trusted list +-- offered to everybody -- so created_by records who added a row for auditing only. +-- +-- source_ref is the reference the administrator gave; source_digest is what it resolved to +-- when last validated. A unit runs the two combined, so a tag moved upstream later cannot +-- change what already ran. + +\c texera_db + +SET search_path TO texera_db; + +BEGIN; + +CREATE TABLE IF NOT EXISTS cu_image +( + iid SERIAL PRIMARY KEY, + name VARCHAR(128) NOT NULL, + source_ref VARCHAR(512) NOT NULL, + source_digest VARCHAR(128), + -- A computing unit may only start from a READY image. + status VARCHAR(16) NOT NULL DEFAULT 'PENDING' + CONSTRAINT ck_cu_image_status + CHECK (status IN ('PENDING', 'VALIDATING', 'READY', 'FAILED')), + -- Numbers the validations of this row, so a retry gets a job name of its own. + attempt INT NOT NULL DEFAULT 0, + validation_log TEXT, + -- Nulled rather than cascaded: the image stays usable if its curator is deleted. + created_by INT, + creation_time TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + update_time TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + FOREIGN KEY (created_by) REFERENCES "user" (uid) ON DELETE SET NULL, + UNIQUE (name), + -- Enforced here, not only in the service: the check there is a read then a write, so + -- two simultaneous registrations would both pass it. + UNIQUE (source_ref) +); + +-- Two references can turn out to be one image once their digests are known. +CREATE INDEX IF NOT EXISTS idx_cu_image_source_digest ON cu_image (source_digest); + +COMMIT;
