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 64c12b3015 feat(computing-unit): pre-pull curated images onto every
node (#8485)
64c12b3015 is described below
commit 64c12b3015cb9e2629c958524c27cbefe3dda084
Author: Tanishq Gandhi <[email protected]>
AuthorDate: Wed Sep 16 00:47:17 2026 +0000
feat(computing-unit): pre-pull curated images onto every node (#8485)
### What changes were proposed in this PR?
The first unit on a node waits for the whole image — about 80 seconds
for a 3 GB one — while every later unit there starts at once. The same
action takes seconds or minutes depending only on which node it landed
on.
Each ready image now gets a DaemonSet: an init container that is the
image, with a command that does nothing, plus a pause container so the
node does not reclaim what was just pulled. The same mechanism the chart
already uses for the deployment's own image, built in code because a
curated image is registered while the cluster is running, so a template
cannot describe it.
One per image, keyed on the image id, so a refresh that resolves a moved
tag repoints the existing pre-pull rather than adding a second one
holding bytes nothing runs. Removed when the image is.
Three moments, because the service has no background threads: an image
reaching READY, an image being deleted, and a read of the list — which
also gives a pre-pull to any ready image that has none, since nothing
else revisits a row once it is ready.
Best-effort throughout. A pre-pull that cannot be created is logged and
ignored: the image still works, and the first unit on each node just
pays for the pull.
Costs node disk, since every node holds every ready image, so
`curatedImages.prepull.enabled` turns it off.
### Also carries a fix to the digest check
Found while reviewing #8546, and folded in here rather than opened
separately.
The validation job prints the image's own start command before the
digest it resolved, and the digest was read from the **first** marker
line in the log. An image whose `Cmd` carries a newline and a marker of
its own therefore passed the check — its `Cmd` still contains
`computing-unit-master` — while naming the digest every unit would then
be pinned to. That defeats what the pin is for: the row is meant to
record the bytes the check approved.
Closed at both ends: the start command is echoed through `tr`, so an
image cannot put a line into the log at all, and the digest is read from
the last marker line, which is the one the job writes last.
Tested by a log in the shape the attack produces, asserting the digest
the job resolved is the one returned. Mutation-checked: restoring the
first-match read fails exactly that test. The `tr` half was checked
against a real shell rather than assumed, since a Scala triple-quoted
string passes backslashes through untouched.
### Any related issues, documentation, discussions?
Closes #8469
Part of #8466
Stacked on #8475, which this calls into, so GitHub shows that commit
here too until it merges.
### How was this PR tested?
Seven new tests for the DaemonSet the code builds, plus the object
applied to a real cluster.
| Case | What it pins |
| --- | --- |
| init container | the image is what gets pulled, and its command is a
no-op |
| pull policy | `IfNotPresent` — a digest cannot resolve to different
bytes later |
| containers | only pause keeps running, so no unit runs that nobody
asked for |
| tolerations | `Exists`, so tainted nodes are not the ones left waiting
|
| name | depends on the image id alone, so a refresh replaces rather
than adds |
| selector | `app` only — a DaemonSet's selector is immutable, so it
must hold nothing that changes |
| labels | the owner and image labels a listing reads |
```
sbt "ComputingUnitManagingService/test" "Config/test" scalafmtCheckAll
ComputingUnitManagingService 182 passed, 0 failed
Config 72 passed, 0 failed
scalafmtCheckAll clean
```
`helm template` renders the manager Role gaining `daemonsets` and the
two new environment variables.
**On minikube,** applying the object the code actually builds:
```
daemonset rolls out, pod Running
alpine appears in the node's image cache
refresh to a different digest -> "configured", still one daemonset, new
image also cached
listing by label -> cu-image-prepull-1 -> texera-cu-image=1
delete removes it
```
The refresh case is the one worth having: a DaemonSet's selector cannot
be changed after creation, so a selector carrying anything mutable would
make every refresh fail permanently.
Also confirmed against the cluster rather than assumed: deleting a
pre-pull that does not exist returns an empty result instead of
throwing, so an image that never reached READY is a silent no-op rather
than a warning on every delete.
### Was this PR authored or co-authored using generative AI tooling?
Generated-by: Claude Code (Claude Opus 5)
---
...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 | 27 +++
.../texera/common/config/CuratedImageConfig.scala | 15 ++
.../service/resource/CuratedImageResource.scala | 99 ++++++++-
.../texera/service/util/ImagePrepullClient.scala | 234 +++++++++++++++++++++
.../service/util/ImageValidationClient.scala | 11 +-
.../resource/CuratedImageResourceSpec.scala | 17 ++
.../service/util/ImagePrepullClientSpec.scala | 186 ++++++++++++++++
10 files changed, 595 insertions(+), 12 deletions(-)
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 6ad4eeb4be..e758a51594 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
@@ -71,6 +71,11 @@ spec:
# Must be the namespace the pool runs in; validation jobs go there.
- name: TEXERA_CURATED_IMAGE_VALIDATION_NAMESPACE
value: {{ .Values.workflowComputingUnitPool.namespace }}
+ - name: TEXERA_CURATED_IMAGE_PREPULL_ENABLED
+ value: "{{ .Values.curatedImages.prepull.enabled }}"
+ # Not the release namespace, which holds the privileged mounter.
+ - name: TEXERA_CURATED_IMAGE_PREPULL_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 29888af964..a1a0ed40cc 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
@@ -42,6 +42,12 @@ rules:
- apiGroups: [""]
resources: ["pods/log"]
verbs: ["get"]
+ # One DaemonSet per ready curated image, created as images are registered,
so the chart
+ # cannot declare them. Scoped here and not to the release namespace, which
holds the
+ # privileged mounter.
+ - apiGroups: ["apps"]
+ resources: ["daemonsets"]
+ verbs: ["get", "list", "watch", "create", "update", "patch", "delete"]
---
apiVersion: rbac.authorization.k8s.io/v1
@@ -56,4 +62,4 @@ subjects:
roleRef:
kind: Role
name: {{ .Values.workflowComputingUnitManager.name }}
- apiGroup: rbac.authorization.k8s.io
\ No newline at end of file
+ apiGroup: rbac.authorization.k8s.io
diff --git a/bin/k8s/values.yaml b/bin/k8s/values.yaml
index 438821b82d..85819ab006 100644
--- a/bin/k8s/values.yaml
+++ b/bin/k8s/values.yaml
@@ -376,6 +376,11 @@ litellm:
curatedImages:
# Off until the UI to manage these ships.
enabled: false
+ prepull:
+ # Pull every ready image onto every node, so the first unit there does not
wait for it.
+ # Costs node disk: each node holds each ready image. Turning it off also
removes the
+ # pre-pulls already made.
+ enabled: true
# headless service for the access of computing units
workflowComputingUnitPool:
diff --git a/common/config/src/main/resources/kubernetes.conf
b/common/config/src/main/resources/kubernetes.conf
index c27fa40d04..57964fc83d 100644
--- a/common/config/src/main/resources/kubernetes.conf
+++ b/common/config/src/main/resources/kubernetes.conf
@@ -151,4 +151,31 @@ curated-images {
validation-cpu-limit = ${?TEXERA_CURATED_IMAGE_VALIDATION_CPU_LIMIT}
validation-memory-limit = "256Mi"
validation-memory-limit = ${?TEXERA_CURATED_IMAGE_VALIDATION_MEMORY_LIMIT}
+
+ # Pull each ready image onto every node as soon as it is ready. Costs node
disk, since
+ # every node holds every ready image. Turning it off also removes the
pre-pulls already
+ # made, on the next read of the image list.
+ prepull-enabled = true
+ prepull-enabled = ${?TEXERA_CURATED_IMAGE_PREPULL_ENABLED}
+
+ # The pool namespace. The release namespace holds the privileged mounter,
where write
+ # access to DaemonSets would mean root on every node.
+ prepull-namespace = "texera-workflow-computing-unit-pool"
+ prepull-namespace = ${?TEXERA_CURATED_IMAGE_PREPULL_NAMESPACE}
+
+ # How long before a failed create is tried again. Without a pause, a failure
that will
+ # not clear is retried on every read of the image list, by every user.
+ prepull-retry-cooldown-seconds = 300
+ prepull-retry-cooldown-seconds =
${?TEXERA_CURATED_IMAGE_PREPULL_RETRY_COOLDOWN_SECONDS}
+
+ # Holds the pod open once the image is pulled, so the node does not reclaim
it.
+ # registry.k8s.io is where the pause image now lives.
+ prepull-pause-image = "registry.k8s.io/pause:3.9"
+ prepull-pause-image = ${?TEXERA_CURATED_IMAGE_PREPULL_PAUSE_IMAGE}
+
+ # The pause container does nothing but exist.
+ prepull-cpu = "1m"
+ prepull-cpu = ${?TEXERA_CURATED_IMAGE_PREPULL_CPU}
+ prepull-memory = "8Mi"
+ prepull-memory = ${?TEXERA_CURATED_IMAGE_PREPULL_MEMORY}
}
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
index 9dc85da9d6..f85f734071 100644
---
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
@@ -44,4 +44,19 @@ object CuratedImageConfig {
/** 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"
+
+ val prepullEnabled: Boolean =
conf.getBoolean("curated-images.prepull-enabled")
+ val prepullNamespace: String =
conf.getString("curated-images.prepull-namespace")
+ val prepullPauseImage: String =
conf.getString("curated-images.prepull-pause-image")
+ val prepullRetryCooldownSeconds: Int =
+ conf.getInt("curated-images.prepull-retry-cooldown-seconds")
+ val prepullCpu: String = conf.getString("curated-images.prepull-cpu")
+ val prepullMemory: String = conf.getString("curated-images.prepull-memory")
+
+ /**
+ * Kubernetes object name for one image's pre-pull. One per image and
stable across its
+ * refreshes, so a new digest replaces the pre-pull rather than adding a
second one
+ * holding the bytes nothing runs any more.
+ */
+ def prepullName(iid: Int): String = s"cu-image-prepull-$iid"
}
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
index 8c19a0138d..e3f7c713b5 100644
---
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
@@ -27,7 +27,7 @@ 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.{ImagePrepullClient,
ImageValidationClient}
import org.apache.texera.service.util.ImageValidationClient.ValidationState
import org.jooq.impl.DSL
import org.jooq.{DSLContext, Record}
@@ -197,6 +197,64 @@ object CuratedImageResource extends LazyLogging {
def nameOf(iid: Int): Option[String] =
Option(context.select(NAME).from(CU_IMAGE).where(IID.eq(iid)).fetchOne()).map(_.get(NAME))
+ /**
+ * Brings the pre-pulls into line with the ready images: adds the missing,
repoints the
+ * stale, removes the rest. Reaping here is what makes them reapable at all
-- every
+ * other removal path can be interrupted between the database write and the
cluster
+ * call, and nothing else revisits a row once it is gone.
+ *
+ * One labelled list and one query, so a settled deployment does no work.
+ */
+ private def reconcilePrepulls(): Unit = {
+ try {
+ // None means the cluster could not be asked, which is not "nothing is
pre-pulled":
+ // acting on that would fire a doomed call for every image on every read.
+ ImagePrepullClient.prepulledRefs().foreach { prepulled =>
+ // No rows wanted when pre-pulling is off, so everything found is
reaped below --
+ // that is what frees the disk the setting costs.
+ val rows =
+ if (!CuratedImageConfig.prepullEnabled) Nil
+ else
+ context
+ .select(IID, STATUS, SOURCE_REF, SOURCE_DIGEST)
+ .from(CU_IMAGE)
+ .fetch()
+ .asScala
+ .toList
+
+ val wantedByImage: Map[Int, String] = rows
+ .filter(_.get(STATUS) == Status.Ready)
+ .flatMap { row =>
+ pinnedRefOf(row.get(SOURCE_REF), row.get(SOURCE_DIGEST))
+ .map(row.get(IID).intValue() -> _)
+ }
+ .toMap
+
+ // By reference, not just presence: a repoint that failed leaves a
pre-pull on the
+ // previous digest, which nothing else would correct.
+ wantedByImage.foreach {
+ case (iid, wanted) =>
+ if (!prepulled.get(iid).contains(wanted))
ImagePrepullClient.ensurePrepull(iid, wanted)
+ }
+
+ // A row still being checked keeps its pre-pull: a refresh moves a
healthy image
+ // through VALIDATING, and reaping then would strip it from every node
and put it
+ // straight back.
+ val keep = rows
+ .filterNot(_.get(STATUS) == Status.Failed)
+ .map(_.get(IID).intValue())
+ .toSet
+
+ // What is left has no live row: deleted, failed, or interrupted
mid-removal.
+ (prepulled.keySet -- keep).foreach(ImagePrepullClient.deletePrepull)
+ }
+ } catch {
+ // Opportunistic: a listing that cannot be repaired is still worth
returning.
+ case e: Throwable =>
+ logger.warn("Could not reconcile the curated-image pre-pulls; leaving
them as they are.", e)
+ }
+ }
+
/**
* 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
@@ -217,7 +275,7 @@ object CuratedImageResource extends LazyLogging {
// 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)
+ try reconcileOne(iid, attempt, row.get(UPDATE_TIME).getTime,
row.get(SOURCE_REF))
catch {
case e: Throwable =>
logger.warn(s"Could not check the validation of image $iid; leaving
it as it is.", e)
@@ -225,7 +283,7 @@ object CuratedImageResource extends LazyLogging {
}
}
- private def reconcileOne(iid: Int, attempt: Int, updatedAt: Long): Unit = {
+ private def reconcileOne(iid: Int, attempt: Int, updatedAt: Long, sourceRef:
String): 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.
@@ -254,7 +312,8 @@ object CuratedImageResource extends LazyLogging {
// 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,
_)))
+ text + validationNote(digest, digest.flatMap(sameContentNote(iid,
_))),
+ sourceRef
)
case ValidationState.Failed =>
@@ -268,7 +327,8 @@ object CuratedImageResource extends LazyLogging {
attempt,
Status.Failed,
None,
- reason.getOrElse("The validation failed without reporting a reason.")
+ reason.getOrElse("The validation failed without reporting a
reason."),
+ sourceRef
)
case ValidationState.Absent =>
@@ -282,7 +342,8 @@ object CuratedImageResource extends LazyLogging {
attempt,
Status.Failed,
None,
- "The validation job disappeared before it reported a result."
+ "The validation job disappeared before it reported a result.",
+ sourceRef
)
}
}
@@ -388,7 +449,8 @@ object CuratedImageResource extends LazyLogging {
attempt: Int,
status: String,
sourceDigest: Option[String],
- log: String
+ log: String,
+ sourceRef: String
): Unit = {
val update = context
.update(CU_IMAGE)
@@ -401,7 +463,19 @@ object CuratedImageResource extends LazyLogging {
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)
+ if (stored > 0) {
+ ImageValidationClient.deleteValidation(iid, attempt)
+ // On READY only, since only then can a unit start from it. A refresh
that resolved
+ // a moved tag arrives here too and repoints the pre-pull. The reference
comes from
+ // the row already read, not a fresh query, which could find it deleted.
+ if (status == Status.Ready) {
+ pinnedRefOf(sourceRef, sourceDigest.orNull)
+ .foreach(ImagePrepullClient.ensurePrepull(iid, _))
+ } else {
+ // Otherwise no unit can start from it, so its pre-pull should not
hold the image.
+ ImagePrepullClient.deletePrepull(iid)
+ }
+ }
}
}
@@ -427,6 +501,7 @@ class CuratedImageResource extends LazyLogging {
def list(@Auth user: SessionUser): List[CuratedImage] = {
requireEnabled()
reconcileRunningValidations()
+ reconcilePrepulls()
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)
@@ -533,14 +608,22 @@ class CuratedImageResource extends LazyLogging {
requireEnabled()
// Nothing of ours holds a copy. A running unit keeps going on what its
node pulled.
ImageValidationClient.deleteAllValidations(iid)
+ // The row goes first, so a concurrent read cannot see it still READY and
re-create the
+ // pre-pull. If this call is lost, reconcilePrepulls reaps it on the next
read.
val deleted = context.deleteFrom(CU_IMAGE).where(IID.eq(iid)).execute()
if (deleted == 0) {
throw new NotFoundException(s"No curated image $iid.")
}
+ // Left behind, it would go on holding the image on every node for an
image no unit can
+ // be started from any more.
+ ImagePrepullClient.deletePrepull(iid)
}
/** Marks the row as being validated and submits the job, in that order. */
private def startValidation(iid: Int, sourceRef: String): Unit = {
+ // Refreshing is the only remedy the page offers, so it must not be held
back by the
+ // cooldown a previous failure earned -- that would look like the button
doing nothing.
+ ImagePrepullClient.clearFailure(iid)
// 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(
diff --git
a/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/ImagePrepullClient.scala
b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/ImagePrepullClient.scala
new file mode 100644
index 0000000000..e13da6b605
--- /dev/null
+++
b/computing-unit-managing-service/src/main/scala/org/apache/texera/service/util/ImagePrepullClient.scala
@@ -0,0 +1,234 @@
+/*
+ * 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.apps.{DaemonSet, DaemonSetBuilder}
+import io.fabric8.kubernetes.api.model.{DeletionPropagation, Quantity,
ResourceRequirementsBuilder}
+import io.fabric8.kubernetes.client.KubernetesClientBuilder
+import org.apache.texera.common.config.CuratedImageConfig
+
+import scala.jdk.CollectionConverters._
+
+/**
+ * Puts a ready curated image on every node before anyone starts a unit from
it, so the
+ * first unit there does not wait for the pull.
+ *
+ * One DaemonSet per image, the same shape the chart uses for the
deployment's own image:
+ * an init container that is the image and does nothing, then a pause
container to hold
+ * the pod open so the node does not reclaim what was pulled. Built here
rather than in
+ * the chart because curated images are registered while the cluster is
running.
+ *
+ * Every call is best-effort: a failure is logged, and the image still works.
+ */
+object ImagePrepullClient extends LazyLogging {
+
+ private val client: io.fabric8.kubernetes.client.KubernetesClient =
+ new KubernetesClientBuilder().build()
+
+ private def namespace: String = CuratedImageConfig.prepullNamespace
+
+ /** Marks the pre-pulls this service owns, so they are found by label, not
by name. */
+ private[service] val OwnerLabel = "texera-cu-image-prepull"
+
+ /** Which image a pre-pull is for, so one can be removed without the others.
*/
+ private[service] val ImageLabel = "texera-cu-image"
+
+ /** Reference and time of each image's last failed create. See
[[isCoolingDown]]. */
+ private val lastFailure = new java.util.concurrent.ConcurrentHashMap[Int,
(String, Long)]()
+
+ private def failedRecently(iid: Int, pinnedRef: String): Boolean =
+ isCoolingDown(Option(lastFailure.get(iid)), pinnedRef,
System.currentTimeMillis())
+
+ /**
+ * Whether to hold a create back. Reconciling runs on every read of the
image list, so
+ * without this a failure that will not clear is retried on every page
load. Keyed by
+ * reference, so a new digest is tried at once.
+ *
+ * Pure, so the rule can be tested without a cluster.
+ */
+ private[service] def isCoolingDown(
+ recorded: Option[(String, Long)],
+ pinnedRef: String,
+ now: Long
+ ): Boolean =
+ recorded.exists {
+ case (failedRef, at) =>
+ failedRef == pinnedRef &&
+ now - at < CuratedImageConfig.prepullRetryCooldownSeconds * 1000L
+ }
+
+ /** Forgets an image's last failure, so the next create is attempted
immediately. */
+ def clearFailure(iid: Int): Unit = lastFailure.remove(iid)
+
+ /** Creates an image's pre-pull, or points an existing one at a new
reference. */
+ def ensurePrepull(iid: Int, pinnedRef: String): Unit = {
+ if (!CuratedImageConfig.prepullEnabled) return
+ if (failedRecently(iid, pinnedRef)) return
+ try {
+ client
+ .apps()
+ .daemonSets()
+ .inNamespace(namespace)
+ .resource(prepullDaemonSet(iid, pinnedRef))
+ .createOr(existing => existing.update())
+ lastFailure.remove(iid)
+ logger.info(s"Pre-pulling curated image $iid ($pinnedRef) onto every
node.")
+ } catch {
+ case e: Throwable =>
+ lastFailure.put(iid, (pinnedRef, System.currentTimeMillis()))
+ logger.warn(
+ s"Could not pre-pull curated image $iid ($pinnedRef). The first unit
on each " +
+ "node will wait for the pull instead.",
+ e
+ )
+ }
+ }
+
+ /** Removes an image's pre-pull. Safe to call when there is none. */
+ def deletePrepull(iid: Int): Unit = {
+ try {
+ client
+ .apps()
+ .daemonSets()
+ .inNamespace(namespace)
+ .withName(CuratedImageConfig.prepullName(iid))
+ // Stated, not left to a default: orphaned pods would hold the image
on every node,
+ // and the reconcile pass lists DaemonSets, so it would never find
them.
+ .withPropagationPolicy(DeletionPropagation.BACKGROUND)
+ .delete()
+ } catch {
+ case e: Throwable =>
+ logger.warn(s"Could not remove the pre-pull for curated image $iid.",
e)
+ }
+ lastFailure.remove(iid)
+ }
+
+ /**
+ * What each image's pre-pull currently pulls, keyed by image.
+ *
+ * None means the cluster could not be asked, which is not "none exist": an
empty map
+ * would have the caller create a pre-pull for every ready image against a
cluster that
+ * has just refused to talk to it.
+ */
+ def prepulledRefs(): Option[Map[Int, String]] = {
+ try {
+ val entries = client
+ .apps()
+ .daemonSets()
+ .inNamespace(namespace)
+ .withLabel(OwnerLabel, "true")
+ .list()
+ .getItems
+ .asScala
+ .flatMap(daemonSet => imageIdOf(daemonSet).map(_ ->
prepulledRefOf(daemonSet).orNull))
+ Some(entries.toMap)
+ } catch {
+ case e: Throwable =>
+ logger.warn("Could not list the curated-image pre-pulls; leaving them
as they are.", e)
+ None
+ }
+ }
+
+ /** The reference a pre-pull pulls, which is its init container's image. */
+ private[service] def prepulledRefOf(daemonSet: DaemonSet): Option[String] =
+ Option(daemonSet.getSpec)
+ .flatMap(spec => Option(spec.getTemplate))
+ .flatMap(template => Option(template.getSpec))
+ .flatMap(podSpec => Option(podSpec.getInitContainers))
+ .flatMap(_.asScala.headOption)
+ .flatMap(container => Option(container.getImage))
+
+ /** The image a pre-pull was created for, from its label rather than its
name. */
+ private[service] def imageIdOf(daemonSet: DaemonSet): Option[Int] =
+ Option(daemonSet.getMetadata)
+ .flatMap(m => Option(m.getLabels))
+ .flatMap(labels => Option(labels.get(ImageLabel)))
+ .flatMap(value => scala.util.Try(value.toInt).toOption)
+
+ private[service] def prepullDaemonSet(iid: Int, pinnedRef: String):
DaemonSet = {
+ val name = CuratedImageConfig.prepullName(iid)
+ val labels = Map("app" -> name, OwnerLabel -> "true", ImageLabel ->
iid.toString).asJava
+
+ // Requests on both containers: a namespace with a ResourceQuota on
requests.cpu or
+ // requests.memory refuses a pod whose init container omits them, and the
refusal is
+ // invisible because the DaemonSet is still created. Costs nothing -- a
pod's request
+ // is the larger of its init containers and the sum of the rest.
+ val prepullerResources = new ResourceRequirementsBuilder()
+ .addToRequests("cpu", new Quantity(CuratedImageConfig.prepullCpu))
+ .addToRequests("memory", new Quantity(CuratedImageConfig.prepullMemory))
+ .build()
+
+ // Limits only here. A limit is per container and never maxed, so capping
the init
+ // container would OOMKill an arbitrary image's shell.
+ val pauseResources = new ResourceRequirementsBuilder()
+ .addToRequests("cpu", new Quantity(CuratedImageConfig.prepullCpu))
+ .addToRequests("memory", new Quantity(CuratedImageConfig.prepullMemory))
+ .addToLimits("cpu", new Quantity(CuratedImageConfig.prepullCpu))
+ .addToLimits("memory", new Quantity(CuratedImageConfig.prepullMemory))
+ .build()
+
+ new DaemonSetBuilder()
+ .withNewMetadata()
+ .withName(name)
+ .withNamespace(namespace)
+ .withLabels(labels)
+ .endMetadata()
+ .withNewSpec()
+ .withNewSelector()
+ // "app" only: a selector cannot be changed once created.
+ .withMatchLabels(Map("app" -> name).asJava)
+ .endSelector()
+ // Every node at once. The default rolls one at a time, each waiting for
a full pull,
+ // which is hours on a large cluster -- and there is no availability to
protect here.
+ .withNewUpdateStrategy()
+ .withType("RollingUpdate")
+ .withNewRollingUpdate()
+ .withMaxUnavailable(new
io.fabric8.kubernetes.api.model.IntOrString("100%"))
+ .endRollingUpdate()
+ .endUpdateStrategy()
+ .withNewTemplate()
+ .withNewMetadata()
+ .withLabels(labels)
+ .endMetadata()
+ .withNewSpec()
+ // No tolerations: computing-unit pods declare none, so a tainted node
is one no unit
+ // can be scheduled onto.
+ .withInitContainers(
+ new io.fabric8.kubernetes.api.model.ContainerBuilder()
+ .withName("prepuller")
+ .withImage(pinnedRef)
+ // The reference names a digest, so what is on the node cannot
differ.
+ .withImagePullPolicy("IfNotPresent")
+ .withCommand("sh", "-c", "true")
+ .withResources(prepullerResources)
+ .build()
+ )
+ .addNewContainer()
+ .withName("pause")
+ .withImage(CuratedImageConfig.prepullPauseImage)
+ .withResources(pauseResources)
+ .endContainer()
+ .endSpec()
+ .endTemplate()
+ .endSpec()
+ .build()
+ }
+}
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
index ebe70bcf3e..c7d45b1935 100644
---
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
@@ -125,12 +125,12 @@ object ImageValidationClient extends LazyLogging {
| echo "$$START_CMD"
| exit 1
|fi
- |echo "Start command: $$START_CMD"
+ |echo "Start command: $$(printf '%s' "$$START_CMD" | tr '\n' ' ')"
|
|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 "Its start command is: $$(printf '%s' "$$START_CMD" | tr '\n' '
')"
| echo "A computing-unit image starts
'${CuratedImageConfig.requiredCommand}'."
| exit 1
|fi
@@ -335,11 +335,16 @@ object ImageValidationClient extends LazyLogging {
/** The digest the source tag resolved to, as printed by a successful job. */
def sourceDigestFrom(log: String): Option[String] =
+ // The job prints this as its very last line, so the last match is the one
it wrote.
+ // Reading the first would let anything echoed earlier -- the image's own
start
+ // command, which its author controls -- name the digest a unit is pinned
to.
log.linesIterator
.map(_.trim)
- .find(_.startsWith(DigestMarker))
+ .filter(_.startsWith(DigestMarker))
.map(_.drop(DigestMarker.length).trim)
.filter(_.nonEmpty)
+ .toSeq
+ .lastOption
/**
* The reference a unit starts from: the administrator's repository at the
resolved
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
index 4858a44b11..3508cd09d4 100644
---
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
@@ -284,6 +284,23 @@ class CuratedImageResourceSpec extends AnyFlatSpec with
Matchers {
// 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 regression this guards: the digest was read from the first marker
line, while the
+ // image's own start command -- which its author controls -- is echoed
earlier. An image
+ // whose Cmd carries a newline and a marker of its own could pass the check
and still
+ // pin units to bytes nobody validated.
+ "sourceDigestFrom" should "ignore a marker the image smuggled into its start
command" in {
+ val log =
+ """Inspecting owner/evil:1.0
+ |Pinned to:
owner/evil@sha256:1111111111111111111111111111111111111111111111111111111111111111
+ |Start command: [computing-unit-master
+
|TEXERA_SOURCE_DIGEST=sha256:2222222222222222222222222222222222222222222222222222222222222222]
+ |Runs as: texera
+
|TEXERA_SOURCE_DIGEST=sha256:1111111111111111111111111111111111111111111111111111111111111111
+ |""".stripMargin
+ ImageValidationClient.sourceDigestFrom(log).value shouldBe
+ "sha256:1111111111111111111111111111111111111111111111111111111111111111"
+ }
+
"the feature flag" should "be off unless a deployment turns it on" in {
CuratedImageConfig.enabled shouldBe false
}
diff --git
a/computing-unit-managing-service/src/test/scala/org/apache/texera/service/util/ImagePrepullClientSpec.scala
b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/util/ImagePrepullClientSpec.scala
new file mode 100644
index 0000000000..30a307e15b
--- /dev/null
+++
b/computing-unit-managing-service/src/test/scala/org/apache/texera/service/util/ImagePrepullClientSpec.scala
@@ -0,0 +1,186 @@
+/*
+ * 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 io.fabric8.kubernetes.api.model.apps.DaemonSetBuilder
+import org.apache.texera.common.config.CuratedImageConfig
+import org.scalatest.OptionValues._
+import org.scalatest.flatspec.AnyFlatSpec
+import org.scalatest.matchers.should.Matchers
+
+import scala.jdk.CollectionConverters._
+
+class ImagePrepullClientSpec extends AnyFlatSpec with Matchers {
+
+ import ImagePrepullClient.{ImageLabel, OwnerLabel, imageIdOf,
prepullDaemonSet}
+
+ private val PinnedRef = "tagandhi19/texera-cu-sklearn@sha256:" + "b" * 64
+
+ "prepullDaemonSet" should "pull the pinned reference and nothing else" in {
+ val spec = prepullDaemonSet(7, PinnedRef).getSpec.getTemplate.getSpec
+ val initContainers = spec.getInitContainers.asScala.toList
+ initContainers should have size 1
+ val init = initContainers.head
+
+ // The whole point: the image is named as the init container, so
scheduling the pod is
+ // what pulls it. The command is a no-op -- nothing in the image is run.
+ init.getImage shouldBe PinnedRef
+ init.getCommand.asScala.toList shouldBe List("sh", "-c", "true")
+
+ // A digest cannot resolve to different bytes later, so re-checking the
registry every
+ // time the pod restarts would buy nothing.
+ init.getImagePullPolicy shouldBe "IfNotPresent"
+
+ // Only the pause container keeps running. If the curated image were left
running here
+ // it would be a computing unit nobody asked for, on every node.
+ val containers = spec.getContainers.asScala.toList
+ containers.map(_.getName) shouldBe List("pause")
+ containers.head.getImage shouldBe CuratedImageConfig.prepullPauseImage
+ }
+
+ // The default rolls one node at a time, each waiting for a full pull --
hours on a large
+ // cluster, for a pod with no availability to protect.
+ it should "repoint every node at once" in {
+ val rolling = prepullDaemonSet(7,
PinnedRef).getSpec.getUpdateStrategy.getRollingUpdate
+ rolling.getMaxUnavailable.getStrVal shouldBe "100%"
+ }
+
+ // A computing-unit pod declares no tolerations, so a tainted node is one no
unit can
+ // land on. Tolerating everything put multi-gigabyte images on control-plane
nodes.
+ it should "schedule exactly where a computing unit can, and no wider" in {
+ val podSpec = prepullDaemonSet(7, PinnedRef).getSpec.getTemplate.getSpec
+ Option(podSpec.getTolerations).map(_.asScala.toList).getOrElse(Nil)
shouldBe Nil
+ }
+
+ // The pool namespace's ResourceQuota refuses a pod whose init container
omits requests,
+ // while still creating the DaemonSet -- so this failed silently.
+ it should "declare the requests a quota would demand" in {
+ val podSpec = prepullDaemonSet(7, PinnedRef).getSpec.getTemplate.getSpec
+ val everyContainer =
+ podSpec.getInitContainers.asScala.toList ++
podSpec.getContainers.asScala.toList
+ everyContainer.foreach { container =>
+ val requests =
Option(container.getResources).map(_.getRequests.asScala).getOrElse(Map.empty)
+ withClue(s"${container.getName} must request cpu and memory: ") {
+ requests.keySet should contain allOf ("cpu", "memory")
+ }
+ }
+ }
+
+ // A limit is per container and never maxed, so a shared one capped the
image's own
+ // shell at 8Mi and OOMKilled it.
+ it should "cap the pause container only, never the image's own shell" in {
+ val podSpec = prepullDaemonSet(7, PinnedRef).getSpec.getTemplate.getSpec
+
+ val prepuller = podSpec.getInitContainers.asScala.head
+
Option(prepuller.getResources).map(_.getLimits.asScala).getOrElse(Map.empty)
shouldBe empty
+
+ val pause = podSpec.getContainers.asScala.head
+ pause.getResources.getLimits.asScala.keySet should contain allOf ("cpu",
"memory")
+ }
+
+ // Reconciling runs on every read of the image list, so a failure that will
not clear
+ // was retried on every page load, for every ready image.
+ "isCoolingDown" should "hold back a reference that just failed" in {
+ val now = 1_000_000_000L
+ val cooldownMillis = CuratedImageConfig.prepullRetryCooldownSeconds * 1000L
+
+ ImagePrepullClient.isCoolingDown(Some((PinnedRef, now)), PinnedRef, now)
shouldBe true
+ ImagePrepullClient.isCoolingDown(
+ Some((PinnedRef, now - cooldownMillis + 1)),
+ PinnedRef,
+ now
+ ) shouldBe true
+ }
+
+ it should "try again once the cooldown has passed" in {
+ val now = 1_000_000_000L
+ val cooldownMillis = CuratedImageConfig.prepullRetryCooldownSeconds * 1000L
+ ImagePrepullClient.isCoolingDown(
+ Some((PinnedRef, now - cooldownMillis)),
+ PinnedRef,
+ now
+ ) shouldBe false
+ }
+
+ // Otherwise a refresh that resolved a new digest would leave nodes on the
old image
+ // until the cooldown expired.
+ it should "not hold back a different digest" in {
+ val now = 1_000_000_000L
+ val other = "owner/name@sha256:" + "f" * 64
+ ImagePrepullClient.isCoolingDown(Some((PinnedRef, now)), other, now)
shouldBe false
+ }
+
+ it should "not hold back an image that has never failed" in {
+ ImagePrepullClient.isCoolingDown(None, PinnedRef, 1_000_000_000L) shouldBe
false
+ }
+
+ // How a stale pre-pull is spotted: comparing only ids would miss a failed
repoint.
+ "prepulledRefOf" should "read back what a pre-pull actually pulls" in {
+ ImagePrepullClient.prepulledRefOf(prepullDaemonSet(7, PinnedRef)).value
shouldBe PinnedRef
+ }
+
+ it should "be empty for a DaemonSet with no init container" in {
+ val strayObject = new
DaemonSetBuilder().withNewMetadata().withName("x").endMetadata().build()
+ ImagePrepullClient.prepulledRefOf(strayObject) shouldBe None
+ }
+
+ // Same name at any digest, so a refresh does not leave a second pre-pull
behind.
+ it should "name itself after the image, so a refresh replaces rather than
adds" in {
+ prepullDaemonSet(7, PinnedRef).getMetadata.getName shouldBe
"cu-image-prepull-7"
+ prepullDaemonSet(7, "owner/name@sha256:" + "c" * 64).getMetadata.getName
shouldBe
+ "cu-image-prepull-7"
+ prepullDaemonSet(8, PinnedRef).getMetadata.getName shouldBe
"cu-image-prepull-8"
+ }
+
+ // A selector is immutable once created, so anything mutable in it would
make every
+ // later repoint fail.
+ it should "select on a label it will never want to change" in {
+ val daemonSet = prepullDaemonSet(7, PinnedRef)
+ daemonSet.getSpec.getSelector.getMatchLabels.asScala shouldBe
+ Map("app" -> "cu-image-prepull-7")
+
+ // The pod template must still match it, or the DaemonSet is rejected
outright.
+ val templateLabels =
daemonSet.getSpec.getTemplate.getMetadata.getLabels.asScala
+ templateLabels("app") shouldBe "cu-image-prepull-7"
+ }
+
+ it should "label the image it belongs to, so one can be removed without the
others" in {
+ val labels = prepullDaemonSet(7, PinnedRef).getMetadata.getLabels.asScala
+ labels(OwnerLabel) shouldBe "true"
+ labels(ImageLabel) shouldBe "7"
+ }
+
+ "imageIdOf" should "read the image back from the label rather than the name"
in {
+ imageIdOf(prepullDaemonSet(7, PinnedRef)).value shouldBe 7
+ }
+
+ // A stray object must not be read as an image id.
+ it should "ignore a DaemonSet that is not one of ours" in {
+ val unlabelled = new
DaemonSetBuilder().withNewMetadata().withName("something").endMetadata()
+ imageIdOf(unlabelled.build()) shouldBe None
+
+ val notANumber = new DaemonSetBuilder()
+ .withNewMetadata()
+ .withName("something")
+ .addToLabels(ImageLabel, "not-a-number")
+ .endMetadata()
+ imageIdOf(notANumber.build()) shouldBe None
+ }
+}