This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/texera.git
The following commit(s) were added to refs/heads/main by this push:
new a310473663 feat(computing-unit): curate computing-unit images (#8475)
a310473663 is described below
commit a310473663c5e34b20214de252538a61fc47b8c0
Author: Tanishq Gandhi <[email protected]>
AuthorDate: Tue Sep 15 17:16:10 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/50.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 3c4ab2b238..5c6ccd79be 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 67ba2307b3..438821b82d 100644
--- a/bin/k8s/values.yaml
+++ b/bin/k8s/values.yaml
@@ -372,6 +372,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 c0a515d163..b16670669a 100644
--- a/sql/changelog.xml
+++ b/sql/changelog.xml
@@ -159,6 +159,11 @@
<sqlFile path="sql/updates/49.sql"/>
</changeSet>
+ <!-- Curated computing-unit images an administrator registers -->
+ <changeSet id="50" author="tanishqgandhi1908">
+ <sqlFile path="sql/updates/50.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 c95ff99c18..4a7e483db0 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/50.sql b/sql/updates/50.sql
new file mode 100644
index 0000000000..0a48bf464a
--- /dev/null
+++ b/sql/updates/50.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;