diff --git a/.agents/skills/gpustack-operator-docs/SKILL.md b/.agents/skills/gpustack-operator-docs/SKILL.md index 798915fc4..7cd5436fb 100644 --- a/.agents/skills/gpustack-operator-docs/SKILL.md +++ b/.agents/skills/gpustack-operator-docs/SKILL.md @@ -41,6 +41,7 @@ one, not to widen the overview. | Standing a cache up end to end, or which object comes first: the pasteable four-object sequence | `docs/kv-cache/walkthrough.md` | | How a **Pod** consumes a pool: the inject label and annotations, the injected keys per engine, a refusal, the isolation record | `docs/reference/kv-cache-injection.md` | | A `ModelArtifact`: its sources, resolution and revalidation, the manifest digest, how a `ModelDeployment` or an `Instance` mounts or downloads it, claim placement, the weight identity in KV keys | `docs/reference/model-artifact.md` | +| The `image` source of a `ModelArtifact`: the digest contract, building weights into an image, image-volume delivery, the version floors, double storage, kubelet image GC, registry mirrors | `docs/reference/model-image-source.md` | | A `ModelStore`, `ModelStoreBinding` or `ModelPrefetch`: the grant, the budget, pinning, TTL expiry, the warm-up pod and why it is label-free | `docs/reference/model-prefetch.md` | | The `v1` views of `ModelArtifact` and `NodeModelStore`, the `progress` subresource and who may read it, the GPUStack server capability map | `docs/reference/model-artifact-views.md` | | A `NodeModelStore` or the `model-manager` plugin: a field and its writer, the status guard, mount authorization, materialization, a failure reason, collection, a metric | `docs/reference/node-model-store.md` | diff --git a/.agents/skills/gpustack-operator-docs/references/page-map.md b/.agents/skills/gpustack-operator-docs/references/page-map.md index af2fc053b..d7fe3a01d 100644 --- a/.agents/skills/gpustack-operator-docs/references/page-map.md +++ b/.agents/skills/gpustack-operator-docs/references/page-map.md @@ -270,7 +270,9 @@ picks and how to switch its policy (`model-deployment-routing.md`), what a `Mode and its engine's exit (`model-deployment-shutdown.md`), and the `ModelArtifact` contract — sources, resolution and revalidation, the manifest digest and its patterns, how a deployment or an Instance consumes it under each delivery, claim placement and the weight identity in KV keys -(`model-artifact.md`), and the `NodeModelStore` resource with the node plugin behind it — its +(`model-artifact.md`), the `image` source's own page — the digest contract, building weights into +an image, image-volume delivery, the version floors, double storage, kubelet image GC and registry +mirrors (`model-image-source.md`), and the `NodeModelStore` resource with the node plugin behind it — its writers and status guard, mount authorization, materialization, failure reasons, collection and metrics (`node-model-store.md`). diff --git a/.agents/skills/gpustack-operator-e2e/SKILL.md b/.agents/skills/gpustack-operator-e2e/SKILL.md index 603291960..f5d42e8e4 100644 --- a/.agents/skills/gpustack-operator-e2e/SKILL.md +++ b/.agents/skills/gpustack-operator-e2e/SKILL.md @@ -156,6 +156,7 @@ Each case is self-contained; its header (see **Case header contract**) states go | 109 | The NodeModelStore's lifecycle: a plugin restart and a deletion of the object each end with the node's entries rewritten from disk, and deleting a Node collects its object while the node registered again gets a new one | `pkg/modelmanager/{report,store}/**`, `pkg/worker/controllers/worker/node_model_store{,_kubelet}.go` | yes (confirm; deletes a worker's Node object and restarts its kubelet) | kind with two workers and docker reaching the node containers (else the Node row SKIPs), the chart with `modelManager.enabled`, the stock python image | | 110 | ModelPrefetch: the store layer names itself in the node's spec; a warm-up pod pinned to the target node mounts label-free and the digest goes Ready; a projection past the grant is refused at admission; pinning is refused without `allowPinned` and lands under it; deleting the prefetch removes its pod and its pin while the tree stays for the grace | `pkg/worker/controllers/worker/{model_prefetch,model_store,model_store_binding}.go`, `pkg/worker/webhooks/worker/{model_store_binding,model_prefetch}.go`, `pkg/worker/settings/value.go`, `cases/_model-hub{-lib.sh,.py}` | yes (confirm) | one worker labeled `e2e.gpustack.ai/warm-pool=true` by the case (case-110 labels and unlabels it), the chart with `modelManager.enabled`, the stock python image; docker reaching the node container for the tree row | | 111 | Node-to-node sync: a cold node materializes a digest from a peer's published tree with source=Peer and its hub byte counter at zero; the seed's plugin pod dying mid-pull still ends Ready (checkpoint resume); a tenant pod cannot reach the peer port | `pkg/modelmanager/peer/**`, `pkg/modelmanager/materialize/materialize.go`, `pkg/modelmanager/config.go`, `deploy/gpustack-operator/chart/templates/model-manager/**`, `docs/reference/node-peer-sync.md` | yes (confirm) | two schedulable workers (a seed and a puller), the chart installed with `modelManager.port` non-zero and the peer NetworkPolicy enabled, `model-store-peer-sync` at its `true` default | +| 112 | Image source: a tag-only reference is refused naming the digest contract; an image artifact resolves claim-shaped (no manifest digest, no revision, no `status.nodes`); an Instance pinned to the node mounts the pinned image read-only through an image volume and reads the fixture byte for byte; a `ModelPrefetch` naming the artifact is refused | `pkg/worker/webhooks/worker/{model_artifact,model_prefetch}.go`, `pkg/worker/controllers/worker/{model_artifact_placement,model_deployment_artifact,instance,model_placement_preference}.go`, `pkg/kubediscovery/feature.go`, `cases/case-112.sh` | yes (confirm) | one worker whose kubelet/containerd serve image volumes (kubelet 1.35+, containerd 2.1+), docker on the runner, the stock `registry:2`, `python:3.12-slim` and `busybox:1.36` images pullable | Each note below is something the **lead** must act on before or around a run. What a case *does* — its goal, environment, inputs, assertions and cleanup — lives in its own header, which the **Case header contract** below requires to be readable on its own; the index never restates it. diff --git a/.agents/skills/gpustack-operator-e2e/cases/case-112.sh b/.agents/skills/gpustack-operator-e2e/cases/case-112.sh new file mode 100755 index 000000000..473cabff1 --- /dev/null +++ b/.agents/skills/gpustack-operator-e2e/cases/case-112.sh @@ -0,0 +1,346 @@ +#!/usr/bin/env bash +# +# CASE 112 — Image source: a digest-pinned OCI image as a ModelArtifact's weights resolves +# claim-shaped and mounts through a Kubernetes image volume (MUTATING, self-recovering) +# +# case-112.sh +# +# Goal: Prove the image source end to end: a tag-only reference is refused at admission, +# naming the digest contract; an image artifact resolves on its first pass with no +# network — claim-shaped status, no manifest digest, no revision; an Instance +# carrying that artifact as a model volume runs a pod whose weights are the image, +# read-only and whole, with the fixture file matching byte for byte; the node cache +# plays no part (status.nodes never appears). +# Environment: Cluster nodes run image volumes — kubelet 1.35+ (beta default-on; 1.33/1.34 would +# need the gate opened, which admission does not opt into) and containerd 2.1+. The +# stock registry:2 and python:3.12-slim images are pullable, and python3 is on the +# runner to generate and push the fixture. No GPU, no engine. There is no auto-skip +# for a cluster that cannot serve image volumes: the admission-refusal rows still +# run, and the mount row FAILs, which is the honest verdict for that cluster. +# Inputs: MOCKED: the fixture image this case generates on the runner (one gzip layer, one +# marker file, addressed by its own sha256) and pushes into an in-cluster registry +# through a port-forward; the node reaches that registry over its own loopback — a +# hostNetwork forwarder pod on the worker turns 127.0.0.1:35000 into the registry +# service, and containerd treats localhost as insecure by default, so no containerd +# configuration is touched. +# Expected: - a tag-only reference is refused, the message naming why the digest is required; +# - the image artifact resolves on its first pass: Resolved=True, resolved carrying +# only resolvedTime — no manifestDigest, no revision; +# - status.nodes never appears for the artifact; +# - the Instance's pod mounts the image read-only at its model mount path and the +# marker file's sha256 equals the fixture's; +# - a ModelPrefetch naming the image artifact is refused: nothing to warm. +# Cleanup: Trap deletes the case-labeled objects in the namespace, the in-cluster registry +# and forwarder, and the scratch dir. +set -uo pipefail + +E2E_SHIM_DIR="$(cd "$(dirname "$0")/../../_e2e-lib/scripts/kubectl-shim" 2>/dev/null && pwd)" +[ -n "$E2E_SHIM_DIR" ] && PATH="$E2E_SHIM_DIR:$PATH" +CASES_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +# shellcheck source=/dev/null +. "${CASES_DIR}/_rows-lib.sh" + +NS="${1:?usage: case-112.sh }" +P=c112 +REG_PORT=35000 # the node's loopback port the forwarder serves the registry on +PUSH_PORT=5550 # the runner's loopback port port-forward pushes through +FAILS=0 +ROWS=() +record() { ROWS+=("$1|$2|$3"); [ "$1" = FAIL ] && FAILS=$((FAILS + 1)); return 0; } + +cleanup() { + echo + echo "[case-112] cleanup" + kubectl -n "$NS" delete modelartifacts.worker.gpustack.ai -l e2e.gpustack.ai/case=112 --ignore-not-found --wait=false >/dev/null 2>&1 + kubectl -n "$NS" delete instances.worker.gpustack.ai -l e2e.gpustack.ai/case=112 --ignore-not-found --wait=false >/dev/null 2>&1 + kubectl -n "$NS" delete pods,svc,cm -l e2e.gpustack.ai/case=112 --ignore-not-found --wait=false --force --grace-period=0 >/dev/null 2>&1 + [ -n "${PF_PID:-}" ] && kill "$PF_PID" 2>/dev/null + rm -rf "$SCRATCH" +} +SCRATCH="$(mktemp -d)" +trap cleanup EXIT + +command -v python3 >/dev/null 2>&1 || { echo "[case-112] python3 is not on the runner; NOTHING WAS VERIFIED"; exit 2; } +# Workers are the nodes without the control-plane role — the shape the suite's kind clusters +# provision (a control-plane plus worker nodes); on a single-node kind cluster this exits 2. +read -r -a WORKERS <<<"$(kubectl get nodes -l '!node-role.kubernetes.io/control-plane' -o jsonpath='{.items[*].metadata.name}')" +[ "${#WORKERS[@]}" -ge 1 ] || { echo "[case-112] needs one schedulable worker, found none; NOTHING WAS VERIFIED"; exit 2; } +NODE="${WORKERS[0]}" +IT=$(kubectl get instancetypes.worker.gpustack.ai \ + -o jsonpath='{.items[?(@.spec.acceleratable==false)].metadata.name}' 2>/dev/null | tr ' ' '\n' | grep -m1 'gpustack-') +[ -n "$IT" ] || { echo "[case-112] no general InstanceType (run case-1 first); NOTHING WAS VERIFIED"; exit 2; } +kubectl create namespace "$NS" --dry-run=client -o yaml | kubectl apply -f - >/dev/null + +echo "[case-112] registry: in-cluster, reached by the node over its own loopback" +kubectl -n "$NS" apply -f - >/dev/null </dev/null)" = "Running" ] && { REG_UP=1; break; } + sleep 2 +done +[ "${REG_UP:-0}" = 1 ] || { echo "[case-112] the registry pod never reached Running; NOTHING WAS VERIFIED"; exit 2; } + +cat > "$SCRATCH/forward.py" <<'PY' +import socket, threading, sys +def pipe(a, b): + try: + while True: + d = a.recv(65536) + if not d: break + b.sendall(d) + except OSError: pass + finally: + try: a.close(); b.close() + except OSError: pass +srv = socket.socket(); srv.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) +srv.bind(("127.0.0.1", 35000)); srv.listen(64) +while True: + c, _ = srv.accept() + up = socket.create_connection((sys.argv[1], 5000)) + threading.Thread(target=pipe, args=(c, up), daemon=True).start() + threading.Thread(target=pipe, args=(up, c), daemon=True).start() +PY +kubectl -n "$NS" apply -f - >/dev/null </dev/null </dev/null)" = "Running" ] && { FWD_UP=1; break; } + sleep 2 +done +[ "${FWD_UP:-0}" = 1 ] || { echo "[case-112] the forwarder pod never reached Running; NOTHING WAS VERIFIED"; exit 2; } + +echo "[case-112] fixture: generate the marker image and push it" +printf 'gpustack e2e image-source fixture, case 112\n' > "$SCRATCH/marker.txt" +EXPECT="$(shasum -a 256 "$SCRATCH/marker.txt" | awk '{print $1}')" +# The port-forward is long-running, so it passes --request-timeout=0 explicitly: the kubectl shim +# would otherwise bound it at 30s and sever the push mid-flight. The local address is chosen with +# --address — docker resolves `localhost` to ::1 first, where nothing listens, and the IPv6 +# connection times out instead of refusing. +kubectl -n "$NS" port-forward --request-timeout=0 --address 127.0.0.1 svc/${P}-registry "${PUSH_PORT}:5000" >"$SCRATCH/pf.log" 2>&1 & +PF_PID=$! +PF_UP=0 +for _ in $(seq 1 30); do + if curl -sf "http://127.0.0.1:${PUSH_PORT}/v2/" >/dev/null 2>&1; then PF_UP=1; break; fi + sleep 1 +done +[ "$PF_UP" = 1 ] || { echo "[case-112] the registry port-forward never came up; NOTHING WAS VERIFIED"; cat "$SCRATCH/pf.log"; exit 2; } +curl -sfi -X POST "http://127.0.0.1:${PUSH_PORT}/v2/${P}-fixture/blobs/uploads/" | head -3 +# The fixture is generated and pushed in pure python through the port-forward: a docker push would +# run in the Docker Desktop daemon's VM, whose 127.0.0.1 is not the runner's. One gzip layer, one +# config, one manifest, each addressed by its own sha256; the manifest digest is what the artifact +# pins. Architecture matches the node's, so the platform check on pull is exact. +ARCH="$(kubectl get node "$NODE" -o jsonpath='{.status.nodeInfo.architecture}')" +DIGEST="$(python3 - "$SCRATCH/marker.txt" "127.0.0.1:${PUSH_PORT}" "$ARCH" <<'PY' +import gzip, hashlib, io, json, sys, tarfile, urllib.request + +# A macOS runner's urllib reads the system proxy settings, which cannot reach the runner's own +# loopback; the registry is loopback-only, so no proxy is ever wanted here. +urllib.request.install_opener(urllib.request.build_opener(urllib.request.ProxyHandler({}))) + +marker_path, registry, arch = sys.argv[1], sys.argv[2], sys.argv[3] +marker = open(marker_path, "rb").read() + +buf = io.BytesIO() +with tarfile.open(fileobj=buf, mode="w", format=tarfile.USTAR_FORMAT) as tf: + info = tarfile.TarInfo("marker.txt") + info.size, info.mtime, info.uid, info.gid, info.mode = len(marker), 0, 0, 0, 0o644 + tf.addfile(info, io.BytesIO(marker)) +uncompressed = buf.getvalue() +compressed = gzip.compress(uncompressed, mtime=0) +diff_id = "sha256:" + hashlib.sha256(uncompressed).hexdigest() +layer_digest = "sha256:" + hashlib.sha256(compressed).hexdigest() + +config = json.dumps({ + "architecture": arch, "os": "linux", + "config": {}, + "rootfs": {"type": "layers", "diff_ids": [diff_id]}, +}).encode() +config_digest = "sha256:" + hashlib.sha256(config).hexdigest() + +manifest = json.dumps({ + "schemaVersion": 2, + "mediaType": "application/vnd.oci.image.manifest.v1+json", + "config": {"mediaType": "application/vnd.oci.image.config.v1+json", "size": len(config), "digest": config_digest}, + "layers": [{"mediaType": "application/vnd.oci.image.layer.v1.tar+gzip", "size": len(compressed), "digest": layer_digest}], +}).encode() +manifest_digest = "sha256:" + hashlib.sha256(manifest).hexdigest() +repo = "c112-fixture" + +def put_blob(blob, digest): + req = urllib.request.Request(f"http://{registry}/v2/{repo}/blobs/uploads/", method="POST") + with urllib.request.urlopen(req) as r: + location = r.headers["Location"] + # registry:2 echoes the Host header, so the Location can already be absolute. + if not location.startswith("http"): + location = f"http://{registry}{location}" + sep = "&" if "?" in location else "?" + put = urllib.request.Request(f"{location}{sep}digest={digest}", data=blob, method="PUT", + headers={"Content-Type": "application/octet-stream"}) + urllib.request.urlopen(put) + +put_blob(config, config_digest) +put_blob(compressed, layer_digest) +req = urllib.request.Request(f"http://{registry}/v2/{repo}/manifests/v1", data=manifest, method="PUT", + headers={"Content-Type": "application/vnd.oci.image.manifest.v1+json"}) +urllib.request.urlopen(req) +print(manifest_digest) +PY +)" +kill "$PF_PID" 2>/dev/null; PF_PID="" +[ -n "$DIGEST" ] || { echo "[case-112] could not pin the fixture digest; NOTHING WAS VERIFIED"; exit 2; } +REFERENCE="localhost:${REG_PORT}/${P}-fixture@${DIGEST}" +echo "[case-112] fixture digest: $DIGEST (node pulls localhost:${REG_PORT} over the loopback forwarder)" + +echo "[case-112] admission: a tag-only reference is refused" +MSG="$(kubectl -n "$NS" create --dry-run=server -f - 2>&1 >/dev/null </dev/null </dev/null)" + [ "$RESOLVED" = "True" ] && break + sleep 5 +done +if [ "$RESOLVED" = "True" ]; then + record PASS "resolved" "the image artifact resolved on its first pass" +else + record FAIL "resolved" "the image artifact never reached Resolved=True" +fi +DIGEST_ECHO="$(kubectl -n "$NS" get modelartifacts.worker.gpustack.ai ${P}-weights -o jsonpath='{.status.resolved.manifestDigest}' 2>/dev/null)" +REVISION_ECHO="$(kubectl -n "$NS" get modelartifacts.worker.gpustack.ai ${P}-weights -o jsonpath='{.status.resolved.revision}' 2>/dev/null)" +if [ -z "$DIGEST_ECHO" ] && [ -z "$REVISION_ECHO" ]; then + record PASS "claim-shaped" "resolved carries neither a manifest digest nor a revision" +else + record FAIL "claim-shaped" "resolved echoes digest='$DIGEST_ECHO' revision='$REVISION_ECHO', want both absent" +fi +NODES_ECHO="$(kubectl -n "$NS" get modelartifacts.worker.gpustack.ai ${P}-weights -o jsonpath='{.status.nodes}' 2>/dev/null)" +if [ -z "$NODES_ECHO" ]; then + record PASS "no-nodes" "status.nodes never appeared: the node cache plays no part" +else + record FAIL "no-nodes" "status.nodes appeared ($NODES_ECHO)" +fi + +echo "[case-112] instance: the operator renders the image volume and the pod reads the weights" +kubectl -n "$NS" apply -f - >/dev/null </dev/null)" = "Running" ] && break + sleep 5 +done +if [ -z "$POD" ]; then + record FAIL "mount" "no Instance pod reached Running; the weights never mounted" +else + GOT="$(kubectl -n "$NS" exec "$POD" -- sh -c "sha256sum /weights/marker.txt" 2>/dev/null | awk '{print $1}')" + if [ "$GOT" = "$EXPECT" ]; then + record PASS "mount" "pod $POD reads the image's marker file with the fixture's sha256" + else + record FAIL "mount" "the marker sha256 read back '$GOT', want $EXPECT" + fi + MOUNT_RO="$(kubectl -n "$NS" get pod "$POD" -o jsonpath='{.spec.containers[0].volumeMounts[?(@.mountPath=="/weights")].readOnly}' 2>/dev/null)" + VOL_IMAGE="$(kubectl -n "$NS" get pod "$POD" -o jsonpath='{.spec.volumes[*].image.reference}' 2>/dev/null)" + if [ "$MOUNT_RO" = "true" ] && [ "$VOL_IMAGE" = "$REFERENCE" ]; then + record PASS "volume-shape" "the pod's weights are one image volume, read-only, the pinned reference" + else + record FAIL "volume-shape" "volumeMount readOnly='$MOUNT_RO', volume image reference='$VOL_IMAGE'" + fi +fi + +echo "[case-112] prefetch: an image artifact is refused" +if kubectl -n "$NS" apply -f - >/dev/null 2>&1 <= 64 { + return ErrIntOverflowGenerated + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + wire |= uint64(b&0x7F) << shift + if b < 0x80 { + break + } + } + fieldNum := int32(wire >> 3) + wireType := int(wire & 0x7) + if wireType == 4 { + return fmt.Errorf("proto: ModelArtifactImageSource: wiretype end group for non-group") + } + if fieldNum <= 0 { + return fmt.Errorf("proto: ModelArtifactImageSource: illegal tag %d (wire type %d)", fieldNum, wire) + } + switch fieldNum { + case 1: + if wireType != 2 { + return fmt.Errorf("proto: wrong wireType = %d for field Reference", wireType) + } + var stringLen uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowGenerated + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + stringLen |= uint64(b&0x7F) << shift + if b < 0x80 { + break + } + } + intStringLen := int(stringLen) + if intStringLen < 0 { + return ErrInvalidLengthGenerated + } + postIndex := iNdEx + intStringLen + if postIndex < 0 { + return ErrInvalidLengthGenerated + } + if postIndex > l { + return io.ErrUnexpectedEOF + } + m.Reference = string(dAtA[iNdEx:postIndex]) + iNdEx = postIndex + default: + iNdEx = preIndex + skippy, err := skipGenerated(dAtA[iNdEx:]) + if err != nil { + return err + } + if (skippy < 0) || (iNdEx+skippy) < 0 { + return ErrInvalidLengthGenerated + } + if (iNdEx + skippy) > l { + return io.ErrUnexpectedEOF + } + iNdEx += skippy + } + } + + if iNdEx > l { + return io.ErrUnexpectedEOF + } + return nil +} func (m *ModelArtifactList) Unmarshal(dAtA []byte) error { l := len(dAtA) iNdEx := 0 @@ -30363,6 +30513,42 @@ func (m *ModelArtifactSource) Unmarshal(dAtA []byte) error { return err } iNdEx = postIndex + case 4: + if wireType != 2 { + return fmt.Errorf("proto: wrong wireType = %d for field Image", wireType) + } + var msglen int + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowGenerated + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + msglen |= int(b&0x7F) << shift + if b < 0x80 { + break + } + } + if msglen < 0 { + return ErrInvalidLengthGenerated + } + postIndex := iNdEx + msglen + if postIndex < 0 { + return ErrInvalidLengthGenerated + } + if postIndex > l { + return io.ErrUnexpectedEOF + } + if m.Image == nil { + m.Image = &ModelArtifactImageSource{} + } + if err := m.Image.Unmarshal(dAtA[iNdEx:postIndex]); err != nil { + return err + } + iNdEx = postIndex default: iNdEx = preIndex skippy, err := skipGenerated(dAtA[iNdEx:]) diff --git a/api/worker/v1alpha1/generated.proto b/api/worker/v1alpha1/generated.proto index 8963f97ec..e78465bdf 100644 --- a/api/worker/v1alpha1/generated.proto +++ b/api/worker/v1alpha1/generated.proto @@ -2927,6 +2927,19 @@ message ModelArtifactHubSource { optional .k8s.io.api.core.v1.LocalObjectReference secretRef = 3; } +// ModelArtifactImageSource is an OCI image reference holding the weights at its root. +message ModelArtifactImageSource { + // Reference is the image reference, pinned to a digest: + // "registry/repository@sha256:<64 lowercase hex>". The digest is the artifact's whole + // identity; status.resolved stays empty for an image source, as for a claim, and the + // reference in this immutable spec is the only record of what the artifact delivers. + // + // +required + // +k8s:validation:minLength=1 + // +k8s:validation:maxLength=1024 + optional string reference = 1; +} + // ModelArtifactList holds the list of ModelArtifact. // // +k8s:deepcopy-gen:interfaces=k8s.io/apimachinery/pkg/runtime.Object @@ -3045,6 +3058,20 @@ message ModelArtifactSource { // // +optional optional ModelArtifactPersistentVolumeClaimSource persistentVolumeClaim = 3; + + // Image is an OCI image reference holding the weights. The reference must pin a digest: + // "registry/repository@sha256:<64 hex>". A tag is mutable, so one artifact could deliver + // different weights on different pulls, and admission refuses it. + // + // The digest pins the image's bytes; it is not the manifest digest a hub source resolves to. + // The operator never reads the registry, so it verifies nothing about what the image holds: + // weights live at the image's root, and putting them there is the build's contract. Delivery + // mounts the image root read-only through a Kubernetes image volume, which needs an apiserver + // and kubelet at 1.35 or above and containerd at 2.1 or above; admission refuses the source + // on an older apiserver. + // + // +optional + optional ModelArtifactImageSource image = 4; } // ModelArtifactSpec defines the desired state of ModelArtifact. diff --git a/api/worker/v1alpha1/generated.protomessage.pb.go b/api/worker/v1alpha1/generated.protomessage.pb.go index b51929971..1f14717ef 100644 --- a/api/worker/v1alpha1/generated.protomessage.pb.go +++ b/api/worker/v1alpha1/generated.protomessage.pb.go @@ -184,6 +184,8 @@ func (*ModelArtifact) ProtoMessage() {} func (*ModelArtifactHubSource) ProtoMessage() {} +func (*ModelArtifactImageSource) ProtoMessage() {} + func (*ModelArtifactList) ProtoMessage() {} func (*ModelArtifactNodes) ProtoMessage() {} diff --git a/api/worker/v1alpha1/model_artifact.go b/api/worker/v1alpha1/model_artifact.go index 1701b2a80..fe8670549 100644 --- a/api/worker/v1alpha1/model_artifact.go +++ b/api/worker/v1alpha1/model_artifact.go @@ -94,6 +94,33 @@ type ModelArtifactSource struct { // // +optional PersistentVolumeClaim *ModelArtifactPersistentVolumeClaimSource `json:"persistentVolumeClaim,omitempty" protobuf:"bytes,3,opt,name=persistentVolumeClaim"` // nolint: lll + + // Image is an OCI image reference holding the weights. The reference must pin a digest: + // "registry/repository@sha256:<64 hex>". A tag is mutable, so one artifact could deliver + // different weights on different pulls, and admission refuses it. + // + // The digest pins the image's bytes; it is not the manifest digest a hub source resolves to. + // The operator never reads the registry, so it verifies nothing about what the image holds: + // weights live at the image's root, and putting them there is the build's contract. Delivery + // mounts the image root read-only through a Kubernetes image volume, which needs an apiserver + // and kubelet at 1.35 or above and containerd at 2.1 or above; admission refuses the source + // on an older apiserver. + // + // +optional + Image *ModelArtifactImageSource `json:"image,omitempty" protobuf:"bytes,4,opt,name=image"` +} + +// ModelArtifactImageSource is an OCI image reference holding the weights at its root. +type ModelArtifactImageSource struct { + // Reference is the image reference, pinned to a digest: + // "registry/repository@sha256:<64 lowercase hex>". The digest is the artifact's whole + // identity; status.resolved stays empty for an image source, as for a claim, and the + // reference in this immutable spec is the only record of what the artifact delivers. + // + // +required + // +k8s:validation:minLength=1 + // +k8s:validation:maxLength=1024 + Reference string `json:"reference" protobuf:"bytes,1,name=reference"` } // ModelArtifactHubSource is one repository on a model hub at one revision. diff --git a/api/worker/v1alpha1/model_deployment.go b/api/worker/v1alpha1/model_deployment.go index 689eca32e..601a9f26b 100644 --- a/api/worker/v1alpha1/model_deployment.go +++ b/api/worker/v1alpha1/model_deployment.go @@ -964,6 +964,10 @@ const ( // ModelDeploymentModelDeliveryNode has the node's model-manager plugin materialize a hub // artifact's verified files into the node's cache and mount them read-only. ModelDeploymentModelDeliveryNode ModelDeploymentModelDelivery = "Node" + // ModelDeploymentModelDeliveryImage mounts an image artifact's OCI image read-only through a + // Kubernetes image volume: kubelet pulls the pinned image on the node that needs it, and the + // node's plugin cache plays no part. + ModelDeploymentModelDeliveryImage ModelDeploymentModelDelivery = "Image" ) // ModelDeploymentRoleStatus is one role's observed readiness. diff --git a/api/worker/v1alpha1/zz_generated.crds.go b/api/worker/v1alpha1/zz_generated.crds.go index 3560cfbe4..1266b6355 100644 --- a/api/worker/v1alpha1/zz_generated.crds.go +++ b/api/worker/v1alpha1/zz_generated.crds.go @@ -4220,6 +4220,22 @@ func crd_gpustack_api_worker_v1alpha1_ModelArtifact() *v1.CustomResourceDefiniti }, Nullable: true, }, + "image": { + Description: "Image is an OCI image reference holding the weights. The reference must pin a digest:\n\"registry/repository@sha256:<64 hex>\". A tag is mutable, so one artifact could deliver\ndifferent weights on different pulls, and admission refuses it.\nThe digest pins the image's bytes; it is not the manifest digest a hub source resolves to.\nThe operator never reads the registry, so it verifies nothing about what the image holds:\nweights live at the image's root, and putting them there is the build's contract. Delivery\nmounts the image root read-only through a Kubernetes image volume, which needs an apiserver\nand kubelet at 1.35 or above and containerd at 2.1 or above; admission refuses the source\non an older apiserver.", + Type: "object", + Required: []string{ + "reference", + }, + Properties: map[string]v1.JSONSchemaProps{ + "reference": { + Description: "Reference is the image reference, pinned to a digest:\n\"registry/repository@sha256:<64 lowercase hex>\". The digest is the artifact's whole\nidentity; status.resolved stays empty for an image source, as for a claim, and the\nreference in this immutable spec is the only record of what the artifact delivers.", + Type: "string", + MaxLength: ptr.To[int64](1024), + MinLength: ptr.To[int64](1), + }, + }, + Nullable: true, + }, "modelScope": { Description: "ModelScope is RESERVED AND REFUSED by admission in this version. Its shape is fixed so that\nopening it is a webhook change rather than a schema change, and the refusal names what opening\nit needs: branch resolution cross-checked against git, a listing that re-lists per directory at\nthe API's silent truncation point, errors classified by the envelope code, and an engine runner\nwhose ModelScope SDK accepts a commit as the revision.", Type: "object", diff --git a/api/worker/v1alpha1/zz_generated.deepcopy.go b/api/worker/v1alpha1/zz_generated.deepcopy.go index a2944195d..59832d5c1 100644 --- a/api/worker/v1alpha1/zz_generated.deepcopy.go +++ b/api/worker/v1alpha1/zz_generated.deepcopy.go @@ -2118,6 +2118,22 @@ func (in *ModelArtifactHubSource) DeepCopy() *ModelArtifactHubSource { return out } +// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. +func (in *ModelArtifactImageSource) DeepCopyInto(out *ModelArtifactImageSource) { + *out = *in + return +} + +// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new ModelArtifactImageSource. +func (in *ModelArtifactImageSource) DeepCopy() *ModelArtifactImageSource { + if in == nil { + return nil + } + out := new(ModelArtifactImageSource) + in.DeepCopyInto(out) + return out +} + // DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. func (in *ModelArtifactList) DeepCopyInto(out *ModelArtifactList) { *out = *in @@ -2227,6 +2243,11 @@ func (in *ModelArtifactSource) DeepCopyInto(out *ModelArtifactSource) { *out = new(ModelArtifactPersistentVolumeClaimSource) **out = **in } + if in.Image != nil { + in, out := &in.Image, &out.Image + *out = new(ModelArtifactImageSource) + **out = **in + } return } diff --git a/api/worker/v1alpha1/zz_generated.model_name.go b/api/worker/v1alpha1/zz_generated.model_name.go index c553482a1..a5e2798fa 100644 --- a/api/worker/v1alpha1/zz_generated.model_name.go +++ b/api/worker/v1alpha1/zz_generated.model_name.go @@ -440,6 +440,11 @@ func (in ModelArtifactHubSource) OpenAPIModelName() string { return "ai.gpustack.worker.v1alpha1.ModelArtifactHubSource" } +// OpenAPIModelName returns the OpenAPI model name for this type. +func (in ModelArtifactImageSource) OpenAPIModelName() string { + return "ai.gpustack.worker.v1alpha1.ModelArtifactImageSource" +} + // OpenAPIModelName returns the OpenAPI model name for this type. func (in ModelArtifactList) OpenAPIModelName() string { return "ai.gpustack.worker.v1alpha1.ModelArtifactList" diff --git a/api/worker/zz_generated.openapi.go b/api/worker/zz_generated.openapi.go index 7b38b5e6c..dd3f860b8 100644 --- a/api/worker/zz_generated.openapi.go +++ b/api/worker/zz_generated.openapi.go @@ -156,6 +156,7 @@ func GetOpenAPIDefinitions(ref common.ReferenceCallback) map[string]common.OpenA v1alpha1.KVCachePoolUsage{}.OpenAPIModelName(): schema_gpustack_api_worker_v1alpha1_KVCachePoolUsage(ref), v1alpha1.ModelArtifact{}.OpenAPIModelName(): schema_gpustack_api_worker_v1alpha1_ModelArtifact(ref), v1alpha1.ModelArtifactHubSource{}.OpenAPIModelName(): schema_gpustack_api_worker_v1alpha1_ModelArtifactHubSource(ref), + v1alpha1.ModelArtifactImageSource{}.OpenAPIModelName(): schema_gpustack_api_worker_v1alpha1_ModelArtifactImageSource(ref), v1alpha1.ModelArtifactList{}.OpenAPIModelName(): schema_gpustack_api_worker_v1alpha1_ModelArtifactList(ref), v1alpha1.ModelArtifactNodes{}.OpenAPIModelName(): schema_gpustack_api_worker_v1alpha1_ModelArtifactNodes(ref), v1alpha1.ModelArtifactPersistentVolumeClaimSource{}.OpenAPIModelName(): schema_gpustack_api_worker_v1alpha1_ModelArtifactPersistentVolumeClaimSource(ref), @@ -8164,6 +8165,30 @@ func schema_gpustack_api_worker_v1alpha1_ModelArtifactHubSource(ref common.Refer } } +func schema_gpustack_api_worker_v1alpha1_ModelArtifactImageSource(ref common.ReferenceCallback) common.OpenAPIDefinition { + return common.OpenAPIDefinition{ + Schema: spec.Schema{ + SchemaProps: spec.SchemaProps{ + Description: "ModelArtifactImageSource is an OCI image reference holding the weights at its root.", + Type: []string{"object"}, + Properties: map[string]spec.Schema{ + "reference": { + SchemaProps: spec.SchemaProps{ + Description: "Reference is the image reference, pinned to a digest: \"registry/repository@sha256:<64 lowercase hex>\". The digest is the artifact's whole identity; status.resolved stays empty for an image source, as for a claim, and the reference in this immutable spec is the only record of what the artifact delivers.", + Default: "", + MinLength: ptr.To[int64](1), + MaxLength: ptr.To[int64](1024), + Type: []string{"string"}, + Format: "", + }, + }, + }, + Required: []string{"reference"}, + }, + }, + } +} + func schema_gpustack_api_worker_v1alpha1_ModelArtifactList(ref common.ReferenceCallback) common.OpenAPIDefinition { return common.OpenAPIDefinition{ Schema: spec.Schema{ @@ -8376,11 +8401,17 @@ func schema_gpustack_api_worker_v1alpha1_ModelArtifactSource(ref common.Referenc Ref: ref(v1alpha1.ModelArtifactPersistentVolumeClaimSource{}.OpenAPIModelName()), }, }, + "image": { + SchemaProps: spec.SchemaProps{ + Description: "Image is an OCI image reference holding the weights. The reference must pin a digest: \"registry/repository@sha256:<64 hex>\". A tag is mutable, so one artifact could deliver different weights on different pulls, and admission refuses it.\n\nThe digest pins the image's bytes; it is not the manifest digest a hub source resolves to. The operator never reads the registry, so it verifies nothing about what the image holds: weights live at the image's root, and putting them there is the build's contract. Delivery mounts the image root read-only through a Kubernetes image volume, which needs an apiserver and kubelet at 1.35 or above and containerd at 2.1 or above; admission refuses the source on an older apiserver.", + Ref: ref(v1alpha1.ModelArtifactImageSource{}.OpenAPIModelName()), + }, + }, }, }, }, Dependencies: []string{ - v1alpha1.ModelArtifactHubSource{}.OpenAPIModelName(), v1alpha1.ModelArtifactPersistentVolumeClaimSource{}.OpenAPIModelName()}, + v1alpha1.ModelArtifactHubSource{}.OpenAPIModelName(), v1alpha1.ModelArtifactImageSource{}.OpenAPIModelName(), v1alpha1.ModelArtifactPersistentVolumeClaimSource{}.OpenAPIModelName()}, } } @@ -8915,11 +8946,11 @@ func schema_gpustack_api_worker_v1alpha1_ModelDeploymentModelStatus(ref common.R }, "delivery": { SchemaProps: spec.SchemaProps{ - Description: "Delivery is how the weights reach the engine: \"Pvc\", the claim mounted read-only at a fixed path; \"Engine\", the engine downloading the pinned commit itself; or \"Node\", the node's model-manager plugin materializing the verified files and mounting them read-only at the same fixed path.\n\n\nPossible enum values:\n - `\"Engine\"` has the engine download a hub artifact's resolved commit into a size-limited cache volume, with the artifact's token from its Secret.\n - `\"Node\"` has the node's model-manager plugin materialize a hub artifact's verified files into the node's cache and mount them read-only.\n - `\"Pvc\"` mounts a claim artifact read-only at a fixed path.", + Description: "Delivery is how the weights reach the engine: \"Pvc\", the claim mounted read-only at a fixed path; \"Engine\", the engine downloading the pinned commit itself; or \"Node\", the node's model-manager plugin materializing the verified files and mounting them read-only at the same fixed path.\n\n\nPossible enum values:\n - `\"Engine\"` has the engine download a hub artifact's resolved commit into a size-limited cache volume, with the artifact's token from its Secret.\n - `\"Image\"` mounts an image artifact's OCI image read-only through a Kubernetes image volume: kubelet pulls the pinned image on the node that needs it, and the node's plugin cache plays no part.\n - `\"Node\"` has the node's model-manager plugin materialize a hub artifact's verified files into the node's cache and mount them read-only.\n - `\"Pvc\"` mounts a claim artifact read-only at a fixed path.", Default: "", Type: []string{"string"}, Format: "", - Enum: []interface{}{"Engine", "Node", "Pvc"}, + Enum: []interface{}{"Engine", "Image", "Node", "Pvc"}, }, }, }, diff --git a/docs/README.md b/docs/README.md index 27c651912..d00ba4817 100644 --- a/docs/README.md +++ b/docs/README.md @@ -97,6 +97,7 @@ Everything written about GPUStack Operator, and the order to read it in. Start a | [Model Deployment Reference](reference/model-deployment.md) | The `ModelDeployment` contract: the inherited reuse domain, the three override tiers and the owned-key table, and the runner-image formula | users, operators, contributors | ~9 min | | [Model Artifact Reference](reference/model-artifact.md) | How a `ModelArtifact` names a model's weights, how it is resolved and revalidated, the manifest digest, file patterns, and how a `ModelDeployment` or an `Instance` mounts or downloads them under each delivery | users, operators | ~16 min | | [Model Artifact Views Reference](reference/model-artifact-views.md) | The `v1` views of `ModelArtifact` and `NodeModelStore`, the `progress` subresource (aggregate, live, authorized, naming no node), the tenant Role, and the capability map for GPUStack server's model files | users, operators, console developers | ~8 min | +| [Model Image Source Reference](reference/model-image-source.md) | The `ModelArtifact` image source: the digest-pinned reference and what it does and does not promise, building weights into an image, image-volume delivery to a `ModelDeployment` or an `Instance`, the version floors, double storage, kubelet image GC, and registry mirrors | users, operators | ~8 min | | [Model Prefetch Reference](reference/model-prefetch.md) | Warming a model onto nodes before any Pod asks: the `ModelStore` / `ModelStoreBinding` / `ModelPrefetch` objects, placement, the warm-up pod and why it carries no queue-name label, budgets, pinning and expiry, and the status views | users, operators | ~9 min | | [Node Model Store Reference](reference/node-model-store.md) | The `NodeModelStore` resource and the `model-manager` plugin: every field and its writer, the status guard, mount authorization, how content is downloaded, verified and published, failure reasons, collection and metrics | operators, contributors | ~12 min | | [Node-to-Node Sync Reference](reference/node-peer-sync.md) | The peer port a node's published trees are served on, the token authentication and NetworkPolicy that bound it, how a cold node pulls with in-stream checkpoints, and the `source` field and metrics that tell peer bytes from hub bytes | operators, contributors | ~6 min | diff --git a/docs/reference/model-artifact.md b/docs/reference/model-artifact.md index 9e5e83f7b..2e4ef85d3 100644 --- a/docs/reference/model-artifact.md +++ b/docs/reference/model-artifact.md @@ -39,6 +39,8 @@ spec: # immutable after creation # persistentVolumeClaim: # claimName: models # this namespace # path: qwen # directory inside the volume; empty is the root + # image: + # reference: registry.example.com/team/qwen@sha256:669ed7b1...48 # digest-pinned allowPatterns: ["*.safetensors", "*.json", "tokenizer*"] # optional; Hugging Face only ignorePatterns: ["original/"] # optional; wins over allowPatterns status: @@ -64,10 +66,15 @@ status: - **A `modelScope` member exists and is refused.** Opening it needs branch resolution checked against git, a listing that recovers from the API's silent truncation at 3000 entries, and a vLLM runner whose ModelScope SDK accepts a commit (1.39.1 or later). +- **An `image` member delivers weights already in a registry.** A digest-pinned reference is the + artifact's whole identity; kubelet pulls and mounts it through an image volume, outside the node + cache. The digest contract, the build, the floors and the costs are on the + [Model Image Source Reference](model-image-source.md). - **Patterns select the files.** They follow Python's `fnmatch.fnmatchcase`: case-sensitive, `*` and `?` cross `/`, a trailing `/` means everything under it, an empty allow list keeps every file, and an ignored file is dropped even when allowed. At most 32 per list, 1 to 256 characters each, - refused on a claim source. A filter that keeps no file is `Resolved=False`, `EmptyManifest`. A + refused on claim and image sources. A filter that keeps no file is `Resolved=False`, + `EmptyManifest`. A filtered artifact needs [Node delivery](#referencing-it-from-a-modeldeployment). - **`status.nodes` counts the content, not the artifact.** Nodes whose `NodeModelStore` lists the digest `Ready`, `Downloading` or `Failed`; artifacts with the same digest see the same nodes, and @@ -159,12 +166,15 @@ A claim source is always mounted directly. A Hugging Face source takes the deliv the node plugin ([switching it](../operation/model-store.md#switch-delivery) rolls each such deployment once): -| | Claim source (`Pvc`) | Hugging Face, `Engine` | Hugging Face, `Node` | -| --- | --- | --- | --- | -| Weights | the claim, read-only, at `/var/lib/gpustack/model`, `subPath` = `path` | downloaded by the engine into `/var/lib/gpustack/model-cache` | the node's verified copy, read-only, at `/var/lib/gpustack/model` | -| vLLM | `vllm serve /var/lib/gpustack/model` | `vllm serve --revision ` | as a claim | -| SGLang | `--model-path /var/lib/gpustack/model` | `--model-path --revision ` | as a claim | -| Both | `--served-model-name `, unless the role states it | same | same | +| | Claim source (`Pvc`) | Hugging Face, `Engine` | Hugging Face, `Node` | Image source (`Image`) | +| --- | --- | --- | --- | --- | +| Weights | the claim, read-only, at `/var/lib/gpustack/model`, `subPath` = `path` | downloaded by the engine into `/var/lib/gpustack/model-cache` | the node's verified copy, read-only, at `/var/lib/gpustack/model` | the image, read-only, at `/var/lib/gpustack/model` | +| vLLM | `vllm serve /var/lib/gpustack/model` | `vllm serve --revision ` | as a claim | as a claim | +| SGLang | `--model-path /var/lib/gpustack/model` | `--model-path --revision ` | as a claim | as a claim | +| Both | `--served-model-name `, unless the role states it | same | same | same | + +An image source takes `Image` whatever the Setting says, and kubelet pulls the pinned image on the +node that needs it — [Model Image Source Reference](model-image-source.md). `--revision` pins the weights and the tokenizer together on both engines. A take-over role (one with `command`) gets the claim or node mount and nothing else, and nothing at all under `Engine`. @@ -276,7 +286,8 @@ hexadecimal digits of the manifest digest, or of the SHA-256 of the artifact's U ## Status `status.model` echoes the artifact, its `revision` and `manifestDigest`, and the `delivery`, `Pvc`, -`Engine` or `Node`. `WeightsReady` says whether every engine role's weights are there: +`Engine`, `Node` or `Image` — an image source echoes neither a revision nor a digest, its reference +being the identity. `WeightsReady` says whether every engine role's weights are there: | Status | Reason | Meaning | | --- | --- | --- | @@ -287,7 +298,7 @@ hexadecimal digits of the manifest digest, or of the SHA-256 of the artifact's U | False | `FilterNeedsNodeDelivery` | an artifact with patterns under `Engine` delivery, which cannot honor them; no new Pod is created | | False | `Materializing` | a node Pod is not mounted yet and its node lists the digest `Downloading` | | False | `MaterializationFailed` | the same, and the node lists it `Failed`; the message carries the node's reason and retry time | -| False | `WeightsNotMounted` | a claim or node Pod's `PodReadyToStartContainers` is not True yet | +| False | `WeightsNotMounted` | a claim, node or image Pod's `PodReadyToStartContainers` is not True yet; for an image Pod the Pod's events carry the pull error | | False | `Downloading` | an engine Pod is not Ready yet; the engine reports no progress of its own | | True | `Mounted`, `Downloaded` | every Pod has its weights | @@ -312,11 +323,18 @@ the claim placement rules above apply. A Hugging Face artifact is mounted throug whatever `model-artifact-delivery-mode` says, since an Instance has no engine to download it; while the CSIDriver does not exist the Instance creates no Pod and says so in its phase message. +An image artifact mounts through an image volume, needing no plugin. An Instance pinned to a node +waits instead when that node cannot run one, naming the node and the floor +([Model Image Source Reference](model-image-source.md#delivery)). + ## Requirements and limits - **Kubernetes 1.29**, the floor the bundled Kueue already sets. `PodReadyToStartContainers` (beta, on by default since 1.29) feeds `WeightsReady`; with it off, a claim deployment's `WeightsReady` stays `WeightsNotMounted` while its replicas run. +- **Image sources need image volumes** and floors above Kubernetes's own; creation is refused + below them. The full line is on the + [Model Image Source Reference](model-image-source.md#versions-and-prerequisites). - **`sglang-gateway` fetches a tokenizer by the worker's `model_path`.** With a claim that path is local, so it logs one 404 warning and routes by text; with Engine delivery it would fetch `main` without a token (not measured). diff --git a/docs/reference/model-image-source.md b/docs/reference/model-image-source.md new file mode 100644 index 000000000..cd330d7db --- /dev/null +++ b/docs/reference/model-image-source.md @@ -0,0 +1,152 @@ +# Model Image Source Reference + +> **Purpose** — the `image` source of a `ModelArtifact`: the digest contract, how weights are built +> into an image, how it is delivered, and what it costs in disk and re-pulls. +> **Audience** users, operators · **Prerequisites** [Model Artifact Reference](model-artifact.md) · +> **Read time** ~8 min + +## Contents + +- [The source and its digest contract](#the-source-and-its-digest-contract) +- [Building the image](#building-the-image) +- [Delivery](#delivery) +- [Versions and prerequisites](#versions-and-prerequisites) +- [Disk cost](#disk-cost) +- [Image GC](#image-gc) +- [Registry mirrors and private registries](#registry-mirrors-and-private-registries) + +## The source and its digest contract + +```yaml +spec: + source: + image: + reference: registry.example.com/team/qwen@sha256:669ed7b1...48 # digest-pinned, required +``` + +The reference **must pin a digest**. A tag is mutable, so one artifact could deliver different +weights on different pulls, and the identity a frozen reference pins would be nothing; admission +refuses a tag with that reason. Patterns are refused too: an image is mounted whole, so there is no +listing to select from. + +> **Why** — the digest pins the image's manifest bytes. It is **not** the manifest digest a Hugging +> Face source resolves to, and `status.resolved` stays empty for an image source, as for a claim: +> the operator never reads the registry, so the reference in the immutable spec is the only record +> of what the artifact delivers, and the KV reuse identity falls back to the artifact's UID. Two +> artifacts naming one image never share KV blocks — safe, only not deduplicated. + +Creation is **refused on an apiserver older than 1.35**, where the image-volume field would be +dropped from the Pod silently — the refusal names the floor and the alternatives. On a supported +apiserver the artifact resolves on its first pass with no network: `Resolved=True`, `resolved` +carrying only `resolvedTime`, no `nodes` aggregation, no revalidation. + +## Building the image + +The build's contract: **the weights live at the image's root**, because delivery mounts the root +whole and offers no sub-path (a sub-path would raise the containerd floor — see +[Versions](#versions-and-prerequisites)). A recipe measured end to end, from +`Qwen/Qwen2.5-7B-Instruct` at commit `a09a3545…`, one safetensors shard per layer: + +```bash +for f in *.safetensors tokenizer.json tokenizer_config.json config.json generation_config.json; do + oras push "$IMAGE" --artifact-type application/vnd.gpustack.weights "$f:application/octet-stream" +done +``` + +The mount is a read-only overlay of the layer snapshots, and every file the engine reads matches the +Hub copy byte for byte — measured sha256-equal to the LFS oids on Kubernetes 1.35.7 / +containerd 2.2.6. That equality is a property of the build, which this recipe produces; the operator +does not verify it. + +## Delivery + +An image source always delivers `Image`, whatever `model-artifact-delivery-mode` says — the Setting +governs Hugging Face sources only. A `ModelDeployment` renders one image volume per role, mounted +read-only at `/var/lib/gpustack/model`, and nothing else: no cache `emptyDir`, no `HF_*` or proxy +environment, no `--revision`, no ephemeral-storage raise. + +The served path and the `--served-model-name` rule are the claim's, and the engine-argument +refusals apply unchanged. A take-over role gets the mount like a claim's. + +An `Instance` model volume mounts the image read-only and whole, and needs no CSIDriver — the +node's plugin plays no part. An Instance pinned to a node checks that node before rendering: kubelet +below 1.35, containerd below 2.1, a runtime that is not containerd, or a version that cannot be read +holds the Pod with a phase message naming the node and the floor. An unpinned Instance cannot know, +and waits like a `ModelDeployment` does. + +`WeightsReady` follows the claim path: `WeightsNotMounted` until every Pod has started. A failed or +slow pull keeps that reason — kubelet's own event on the Pod names the pull error, and that event is +the diagnostic path. + +## Versions and prerequisites + +| Component | Floor | Why | +| --- | --- | --- | +| apiserver | 1.35 | the `ImageVolume` feature gate is beta **on by default** from 1.35, GA in 1.36; below that the gate must be opened on the apiserver, which this version does not opt into, and a 1.32 or older apiserver refuses the source at admission | +| kubelet | 1.35 | it mounts the volume; the same gate line applies, checked per pinned node | +| containerd | 2.1 | the runtime that mounts image volumes; sub-path mounts would need 2.2, which is why none are offered | + +PSA is not an obstacle: on an apiserver ≥ 1.33, Restricted admits `image` volumes +(kubernetes#130394, fixed in 1.33, not backported); on ≤ 1.32 Restricted refuses them, and +`enforce-version` makes no difference either way. Baseline admits wherever the apiserver knows the +field. + +A private registry is reached with the role's or Instance's existing `imagePullSecrets`; the image +volume pull assembles credentials the same way a container image pull does. + +## Disk cost + +**A pulled image holds its bytes twice** while `discard_unpacked_layers` is `false` (the containerd +default): the compressed blobs in the content store and the unpacked snapshots. Measured on +containerd 2.2.6 with a 15,242,807,270-byte model: + +| Image shape | Content store | Snapshots | Total | vs raw weights | +| --- | --- | --- | --- | --- | +| gzip-compressed layers | +12,074,816,089 B | +15,242,807,270 B | ~27.3 GB | **1.79×** | +| uncompressed layers | +15,242,807,270 B | +15,242,807,270 B | ~30.5 GB | **2.00×** | + +Uncompressed layers save decompression time, not disk. **Kubelet's own accounting does not see the +doubling**: the CRI image size is the compressed size, ~12.07 GB where the disk really holds ~27.3 +GB (**~2.26×** apart). Size a node's disk from the table, never from `kubectl`'s image sizes. + +A node image that ships `discard_unpacked_layers = true` — the local kind images do — keeps only the +blobs and does not pay the snapshot copy; do not take that as the fleet default. + +## Image GC + +Three rules, verified on 1.35.7 / 2.2.6: + +1. **In use, not collected.** While any *running* container mounts the image, kubelet's image GC + passes over it (the protection needs kubelet ≥ 1.31 and containerd ≥ 2.1). A completed Pod + counts for nothing. +2. **Released, collected in the next round.** After the last referencing Pod goes, the image can be + collected on the next GC cycle (5 minutes; measured release→collection 4m2s). `imageMinimumGCAge` + does not protect it: kubelet compares against the image's first-detected time, not its release. + + **Scaling a deployment to zero and back can re-pull a dozen GB** (a first pull measured 12m43s on + a 2-vCPU node at ~20 MB/s). +3. **The 85% cascade.** A model image that pushes the disk past `imageGCHighThresholdPercent` + (default 85) drags every other unused image on the node into the same rounds — measured: 18 + unrelated images collected alongside. + +The only retention is another reference: a resident Pod mounting the image keeps it. A deployment's +replicas are themselves that reference while they run. + +## Registry mirrors and private registries + +The pull is a normal containerd pull, so it rides whatever the node's registry configuration +resolves: a Harbor pull-through cache, a Spegel peer-to-peer mirror, or Dragonfly as a configured +containerd mirror. None of these need operator support, and none are measured from here; Dragonfly +as an operator-integrated prefetcher is future work. + +Two images sharing layer digests share snapshots and the download: measured, a second image whose +layers overlapped pulled in 112 ms and added 11.8 MB where a cold node paid 2m10s and 7.89 GB. + +--- + +**See also** — [Model Artifact Reference](model-artifact.md) for the artifact contract the image +source joins · [Node Model Store Reference](node-model-store.md) for the plugin chain an image +source stays out of · [Model Prefetch Reference](model-prefetch.md) for why there is nothing to +warm. + +**Next** → [Model Prefetch Reference](model-prefetch.md) diff --git a/pkg/kubeclients/applyconfiguration/utils.go b/pkg/kubeclients/applyconfiguration/utils.go index 866dbeb38..f2486c6b8 100644 --- a/pkg/kubeclients/applyconfiguration/utils.go +++ b/pkg/kubeclients/applyconfiguration/utils.go @@ -1450,6 +1450,8 @@ func ForKind(kind schema.GroupVersionKind) interface{} { return &applyconfigurationworkerv1alpha1.ModelArtifactApplyConfiguration{} case workerv1alpha1.SchemeGroupVersion.WithKind("ModelArtifactHubSource"): return &applyconfigurationworkerv1alpha1.ModelArtifactHubSourceApplyConfiguration{} + case workerv1alpha1.SchemeGroupVersion.WithKind("ModelArtifactImageSource"): + return &applyconfigurationworkerv1alpha1.ModelArtifactImageSourceApplyConfiguration{} case workerv1alpha1.SchemeGroupVersion.WithKind("ModelArtifactNodes"): return &applyconfigurationworkerv1alpha1.ModelArtifactNodesApplyConfiguration{} case workerv1alpha1.SchemeGroupVersion.WithKind("ModelArtifactPersistentVolumeClaimSource"): diff --git a/pkg/kubeclients/applyconfiguration/worker/v1alpha1/modelartifactimagesource.go b/pkg/kubeclients/applyconfiguration/worker/v1alpha1/modelartifactimagesource.go new file mode 100644 index 000000000..3aca52aca --- /dev/null +++ b/pkg/kubeclients/applyconfiguration/worker/v1alpha1/modelartifactimagesource.go @@ -0,0 +1,29 @@ +// Code generated by api. DO NOT EDIT. + +package v1alpha1 + +// ModelArtifactImageSourceApplyConfiguration represents a declarative configuration of the ModelArtifactImageSource type for use +// with apply. +// +// ModelArtifactImageSource is an OCI image reference holding the weights at its root. +type ModelArtifactImageSourceApplyConfiguration struct { + // Reference is the image reference, pinned to a digest: + // "registry/repository@sha256:<64 lowercase hex>". The digest is the artifact's whole + // identity; status.resolved stays empty for an image source, as for a claim, and the + // reference in this immutable spec is the only record of what the artifact delivers. + Reference *string `json:"reference,omitempty"` +} + +// ModelArtifactImageSourceApplyConfiguration constructs a declarative configuration of the ModelArtifactImageSource type for use with +// apply. +func ModelArtifactImageSource() *ModelArtifactImageSourceApplyConfiguration { + return &ModelArtifactImageSourceApplyConfiguration{} +} + +// WithReference sets the Reference field in the declarative configuration to the given value +// and returns the receiver, so that objects can be built by chaining "With" function invocations. +// If called multiple times, the Reference field is set to the value of the last call. +func (b *ModelArtifactImageSourceApplyConfiguration) WithReference(value string) *ModelArtifactImageSourceApplyConfiguration { + b.Reference = &value + return b +} diff --git a/pkg/kubeclients/applyconfiguration/worker/v1alpha1/modelartifactsource.go b/pkg/kubeclients/applyconfiguration/worker/v1alpha1/modelartifactsource.go index de049744c..662bddde9 100644 --- a/pkg/kubeclients/applyconfiguration/worker/v1alpha1/modelartifactsource.go +++ b/pkg/kubeclients/applyconfiguration/worker/v1alpha1/modelartifactsource.go @@ -19,6 +19,17 @@ type ModelArtifactSourceApplyConfiguration struct { // reads the claim's content, so the artifact has no revision and no digest, and what the // directory holds, and any change to it, is the user's. PersistentVolumeClaim *ModelArtifactPersistentVolumeClaimSourceApplyConfiguration `json:"persistentVolumeClaim,omitempty"` + // Image is an OCI image reference holding the weights. The reference must pin a digest: + // "registry/repository@sha256:<64 hex>". A tag is mutable, so one artifact could deliver + // different weights on different pulls, and admission refuses it. + // + // The digest pins the image's bytes; it is not the manifest digest a hub source resolves to. + // The operator never reads the registry, so it verifies nothing about what the image holds: + // weights live at the image's root, and putting them there is the build's contract. Delivery + // mounts the image root read-only through a Kubernetes image volume, which needs an apiserver + // and kubelet at 1.35 or above and containerd at 2.1 or above; admission refuses the source + // on an older apiserver. + Image *ModelArtifactImageSourceApplyConfiguration `json:"image,omitempty"` } // ModelArtifactSourceApplyConfiguration constructs a declarative configuration of the ModelArtifactSource type for use with @@ -50,3 +61,11 @@ func (b *ModelArtifactSourceApplyConfiguration) WithPersistentVolumeClaim(value b.PersistentVolumeClaim = value return b } + +// WithImage sets the Image field in the declarative configuration to the given value +// and returns the receiver, so that objects can be built by chaining "With" function invocations. +// If called multiple times, the Image field is set to the value of the last call. +func (b *ModelArtifactSourceApplyConfiguration) WithImage(value *ModelArtifactImageSourceApplyConfiguration) *ModelArtifactSourceApplyConfiguration { + b.Image = value + return b +} diff --git a/pkg/kubediscovery/feature.go b/pkg/kubediscovery/feature.go index fc3f72611..c470177b8 100644 --- a/pkg/kubediscovery/feature.go +++ b/pkg/kubediscovery/feature.go @@ -13,6 +13,12 @@ const ( // below that the API server DROPS the field without an error, degrading a native sidecar to a // plain init container. FeatureNativeSidecar Feature = iota + 1 + + // FeatureImageVolume is the volumes[].image field an image volume is rendered with. The + // ImageVolume feature gate turns on by default in Kubernetes 1.35 (beta) and is GA in 1.36; + // below that the gate has to be opened on the apiserver, and an apiserver that does not know + // the field DROPS it without an error, so a Pod would silently lose its weights. + FeatureImageVolume ) // SupportsFeature returns whether the given Kubernetes version supports the feature. @@ -28,8 +34,11 @@ func SupportsFeature(v *Version, f Feature) bool { return false } - if f == FeatureNativeSidecar { + switch f { + case FeatureNativeSidecar: return parsed.Major() > 1 || (parsed.Major() == 1 && parsed.Minor() >= 29) + case FeatureImageVolume: + return parsed.Major() > 1 || (parsed.Major() == 1 && parsed.Minor() >= 35) } return false } diff --git a/pkg/kubediscovery/feature_test.go b/pkg/kubediscovery/feature_test.go index 0286e92f6..021e90396 100644 --- a/pkg/kubediscovery/feature_test.go +++ b/pkg/kubediscovery/feature_test.go @@ -9,24 +9,37 @@ import ( func TestSupportsFeature(t *testing.T) { for _, tc := range []struct { name string + feature Feature gitVersion string want bool }{ - {"at the floor", "v1.29.0", true}, - {"above the floor", "v1.30.2", true}, - {"below the floor", "v1.28.15", false}, - {"the chart's own floor", "v1.23.17", false}, - {"vendor suffix", "v1.29.3-gke.1000", true}, - {"pre-release at the floor", "v1.29.0-rc.1", true}, - {"unparseable", "not-a-version", false}, - {"empty", "", false}, + // FeatureNativeSidecar: on by default from 1.29. + {"sidecar at the floor", FeatureNativeSidecar, "v1.29.0", true}, + {"sidecar above the floor", FeatureNativeSidecar, "v1.30.2", true}, + {"sidecar below the floor", FeatureNativeSidecar, "v1.28.15", false}, + {"sidecar at the chart's own floor", FeatureNativeSidecar, "v1.23.17", false}, + {"sidecar vendor suffix", FeatureNativeSidecar, "v1.29.3-gke.1000", true}, + {"sidecar pre-release at the floor", FeatureNativeSidecar, "v1.29.0-rc.1", true}, + // FeatureImageVolume: on by default from 1.35 (1.36 GA with LockToDefault). + {"image volume at the floor", FeatureImageVolume, "v1.35.0", true}, + {"image volume above the floor", FeatureImageVolume, "v1.35.7", true}, + {"image volume at GA", FeatureImageVolume, "v1.36.1", true}, + {"image volume below the floor", FeatureImageVolume, "v1.34.9", false}, + {"image volume at the beta-off line", FeatureImageVolume, "v1.33.1", false}, + {"image volume at the alpha line", FeatureImageVolume, "v1.31.0", false}, + {"image volume vendor suffix", FeatureImageVolume, "v1.35.5-kind.1", true}, + // Versions nobody can read never unlock either capability. + {"image volume unparseable", FeatureImageVolume, "not-a-version", false}, + {"image volume empty", FeatureImageVolume, "", false}, } { t.Run(tc.name, func(t *testing.T) { assert.Equal(t, tc.want, - SupportsFeature(&Version{GitVersion: tc.gitVersion}, FeatureNativeSidecar)) + SupportsFeature(&Version{GitVersion: tc.gitVersion}, tc.feature)) }) } assert.False(t, SupportsFeature(nil, FeatureNativeSidecar), "no version never unlocks a capability whose floor cannot be decided") + assert.False(t, SupportsFeature(nil, FeatureImageVolume), + "no version never unlocks a capability whose floor cannot be decided") } diff --git a/pkg/worker/controllers/worker/instance.go b/pkg/worker/controllers/worker/instance.go index 5c68c7217..a7d9ad12b 100644 --- a/pkg/worker/controllers/worker/instance.go +++ b/pkg/worker/controllers/worker/instance.go @@ -14,6 +14,7 @@ import ( "k8s.io/apimachinery/pkg/api/resource" meta "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/util/intstr" + utilversion "k8s.io/apimachinery/pkg/util/version" "k8s.io/utils/ptr" ctrl "sigs.k8s.io/controller-runtime" ctrlbuilder "sigs.k8s.io/controller-runtime/pkg/builder" @@ -29,6 +30,7 @@ import ( workercore "gpustack.ai/gpustack/api/worker/v1alpha1" "gpustack.ai/gpustack/pkg/controller" "gpustack.ai/gpustack/pkg/deviceplugin" + "gpustack.ai/gpustack/pkg/kubediscovery" "gpustack.ai/gpustack/pkg/kubemeta" "gpustack.ai/gpustack/pkg/nodefeature" "gpustack.ai/gpustack/pkg/systemmeta" @@ -773,6 +775,9 @@ func convertAdditionalVolumes( case w.Render.Delivery == workercore.ModelDeploymentModelDeliveryNode: vs = w.Render.nodeVolumeSource() readOnly, subPath = true, "" + case w.Render.Delivery == workercore.ModelDeploymentModelDeliveryImage: + vs.Image = &core.ImageVolumeSource{Reference: w.Render.ImageReference} + readOnly, subPath = true, "" default: continue } @@ -830,12 +835,53 @@ func (r *InstanceReconciler) resolveInstanceModelVolumes( if w.Blocked { return nil, w.Message, nil } + // An image volume runs on kubelet and containerd, so a pinned Instance checks its node + // before rendering: a Pod whose volume the node cannot run would fail late, in a kubelet + // event the Instance status never surfaces. An Instance pinned to no node cannot know, and + // waits for the Pod's own fate like a ModelDeployment does. + if w.Render.Delivery == workercore.ModelDeploymentModelDeliveryImage && inst.Spec.NodeName != "" { + nd := new(core.Node) + if err := r.Client.Get(ctx, ctrlcli.ObjectKey{Name: inst.Spec.NodeName}, nd, ctrlclix.WithoutQuorum); err != nil { + if kerrors.IsNotFound(err) { + return nil, fmt.Sprintf("the node %q this Instance pins does not exist", inst.Spec.NodeName), nil + } + return nil, "", fmt.Errorf("read node %q: %w", inst.Spec.NodeName, err) + } + if unsupported := modelImageRuntimeUnsupported( + nd.Status.NodeInfo.KubeletVersion, nd.Status.NodeInfo.ContainerRuntimeVersion); unsupported != "" { + return nil, fmt.Sprintf("ModelArtifact %q is delivered as an image volume, and node %q cannot run one: "+ + "%s; pin another node or upgrade it", av.Model.ArtifactRef.Name, inst.Spec.NodeName, unsupported), nil + } + } models[i] = w } return models, "", nil } +// instanceModelImageContainerdFloor is the containerd version that mounts image volumes from. +var instanceModelImageContainerdFloor = utilversion.MustParseGeneric("2.1") + +// modelImageRuntimeUnsupported reports the first floor a node misses for an image volume, or "" +// when it meets them both: kubelet runs the ImageVolume feature, on by default from 1.35, and +// containerd mounts image volumes from 2.1. An unreadable version reports unsupported — the same +// fail-closed answer the apiserver's gate gives. +func modelImageRuntimeUnsupported(kubeletVersion, containerRuntimeVersion string) string { + if !kubediscovery.SupportsFeature(&kubediscovery.Version{GitVersion: kubeletVersion}, kubediscovery.FeatureImageVolume) { + return fmt.Sprintf("kubelet %q is below the 1.35 the ImageVolume feature turns on at", kubeletVersion) + } + runtime, runtimeVersion, found := strings.Cut(containerRuntimeVersion, "://") + if !found || runtime != "containerd" { + return fmt.Sprintf("the container runtime %q is not containerd, which is what mounts image volumes", + containerRuntimeVersion) + } + if v, err := utilversion.ParseGeneric(runtimeVersion); err != nil || v.LessThan(instanceModelImageContainerdFloor) { + return fmt.Sprintf("containerd %q is below the 2.1 that mounts image volumes from", runtimeVersion) + } + + return "" +} + // additionalVolumeName is the Pod volume name of the additional volume at the given index. func additionalVolumeName(i int) string { return "additional-" + strconvx.Itoa(i) diff --git a/pkg/worker/controllers/worker/model_artifact.go b/pkg/worker/controllers/worker/model_artifact.go index 7ec3bf93e..c50c92af3 100644 --- a/pkg/worker/controllers/worker/model_artifact.go +++ b/pkg/worker/controllers/worker/model_artifact.go @@ -131,6 +131,8 @@ func (r *ModelArtifactReconciler) Reconcile(ctx context.Context, req ctrl.Reques result, check = r.reconcileHuggingFace(ctx, ma) case ma.Spec.Source.PersistentVolumeClaim != nil: result = r.reconcileClaim(ctx, ma) + case ma.Spec.Source.Image != nil: + r.reconcileImage(ma) default: // Admission refuses every other shape, ModelScope included; one stored before a rule // existed stays unresolved and says why. @@ -280,6 +282,18 @@ func (r *ModelArtifactReconciler) reconcileClaim(ctx context.Context, ma *worker return ctrl.Result{} } +// reconcileImage resolves an image source: admission pinned the reference, so the resolution is +// the spec being stored. No registry is read, no revalidation is scheduled, and no node ever +// aggregates — the image's bytes are the user's, delivered by kubelet wherever a pod lands. +func (r *ModelArtifactReconciler) reconcileImage(ma *workercore.ModelArtifact) { + if ma.Status.Resolved == nil { + ma.Status.Resolved = &workercore.ModelArtifactResolved{ResolvedTime: meta.NewTime(r.now())} + } + ModelArtifactConditionResolved.True(ma, modelArtifactReasonResolved, + "the image is the user's; the operator does not read it") + ModelArtifactConditionDegraded.False(ma, modelArtifactReasonHealthy, "") +} + // modelArtifactCheck is one artifact's pacing entry. type modelArtifactCheck struct { secretVersion string diff --git a/pkg/worker/controllers/worker/model_artifact_image_test.go b/pkg/worker/controllers/worker/model_artifact_image_test.go new file mode 100644 index 000000000..b6ea77ba5 --- /dev/null +++ b/pkg/worker/controllers/worker/model_artifact_image_test.go @@ -0,0 +1,201 @@ +package worker + +import ( + "context" + "strconv" + "strings" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + core "k8s.io/api/core/v1" + meta "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/types" + ctrlcli "sigs.k8s.io/controller-runtime/pkg/client" + ctrlfake "sigs.k8s.io/controller-runtime/pkg/client/fake" + + workercore "gpustack.ai/gpustack/api/worker/v1alpha1" + "gpustack.ai/gpustack/pkg/kubeclients/kubernetes/scheme" +) + +var testImageReference = "registry.example.com/team/qwen@sha256:" + strings.Repeat("a", 64) + +func imageArtifactFixture(uid types.UID) *workercore.ModelArtifact { + ma := artifactFixture("models", true, true) + ma.Spec.Source = workercore.ModelArtifactSource{Image: &workercore.ModelArtifactImageSource{Reference: testImageReference}} + ma.UID = uid + + return ma +} + +func imageModelInstanceFixture(nodeName string) *workercore.Instance { + return &workercore.Instance{ + ObjectMeta: meta.ObjectMeta{Namespace: "team-a", Name: "inst"}, + Spec: workercore.InstanceSpec{ + NodeName: nodeName, + InstanceTemplate: workercore.InstanceTemplate{AdditionalVolumes: []workercore.InstanceAdditionalVolume{ + {MountPath: "/models", Model: &workercore.InstanceModelVolumeSource{ArtifactRef: core.LocalObjectReference{Name: "qwen"}}}, + }}, + }, + } +} + +func imageRuntimeNode(kubelet, containerd string) *core.Node { + return &core.Node{ + ObjectMeta: meta.ObjectMeta{Name: "node-1", Labels: map[string]string{core.LabelHostname: "host-1"}}, + Status: core.NodeStatus{NodeInfo: core.NodeSystemInfo{KubeletVersion: kubelet, ContainerRuntimeVersion: containerd}}, + } +} + +// TestModelImageRuntimeUnsupported walks the node floors an image volume needs. +func TestModelImageRuntimeUnsupported(t *testing.T) { + cases := []struct { + name string + kubelet string + containerd string + wantFragment string + }{ + {name: "current versions are supported", kubelet: "v1.35.5", containerd: "containerd://2.3.1"}, + {name: "an old kubelet is refused", kubelet: "v1.34.9", containerd: "containerd://2.3.1", wantFragment: "kubelet"}, + {name: "an old containerd is refused", kubelet: "v1.35.5", containerd: "containerd://2.0.0", wantFragment: "containerd"}, + {name: "a non-containerd runtime is refused", kubelet: "v1.35.5", containerd: "docker://27.0", wantFragment: "not containerd"}, + {name: "an unreadable kubelet is refused", kubelet: "", containerd: "containerd://2.3.1", wantFragment: "kubelet"}, + {name: "an unreadable containerd is refused", kubelet: "v1.35.5", containerd: "containerd://", wantFragment: "containerd"}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + unsupported := modelImageRuntimeUnsupported(c.kubelet, c.containerd) + if c.wantFragment == "" { + assert.Empty(t, unsupported) + return + } + assert.Contains(t, unsupported, c.wantFragment) + }) + } +} + +// TestInstanceModelImageDelivery walks the pinned Instance's node pre-check through the resolver. +func TestInstanceModelImageDelivery(t *testing.T) { + t.Run("a supported node lets the pod build", func(t *testing.T) { + cli := ctrlfake.NewClientBuilder().WithScheme(scheme.Scheme). + WithObjects(imageArtifactFixture("uid-qwen"), imageRuntimeNode("v1.35.5", "containerd://2.3.1")). + Build() + models, blocked, err := (&InstanceReconciler{Client: cli}). + resolveInstanceModelVolumes(context.Background(), imageModelInstanceFixture("node-1")) + require.NoError(t, err) + assert.Empty(t, blocked) + require.Contains(t, models, 0) + assert.Equal(t, workercore.ModelDeploymentModelDeliveryImage, models[0].Render.Delivery) + assert.Equal(t, testImageReference, models[0].Render.ImageReference) + }) + + t.Run("an old kubelet blocks with the node and the floor named", func(t *testing.T) { + cli := ctrlfake.NewClientBuilder().WithScheme(scheme.Scheme). + WithObjects(imageArtifactFixture("uid-qwen"), imageRuntimeNode("v1.32.9", "containerd://2.3.1")). + Build() + _, blocked, err := (&InstanceReconciler{Client: cli}). + resolveInstanceModelVolumes(context.Background(), imageModelInstanceFixture("node-1")) + require.NoError(t, err) + assert.Contains(t, blocked, "node-1") + assert.Contains(t, blocked, "cannot run one") + }) + + t.Run("a missing node blocks", func(t *testing.T) { + cli := ctrlfake.NewClientBuilder().WithScheme(scheme.Scheme). + WithObjects(imageArtifactFixture("uid-qwen")). + Build() + _, blocked, err := (&InstanceReconciler{Client: cli}). + resolveInstanceModelVolumes(context.Background(), imageModelInstanceFixture("node-1")) + require.NoError(t, err) + assert.Contains(t, blocked, "does not exist") + }) + + t.Run("an unpinned instance skips the node check", func(t *testing.T) { + cli := ctrlfake.NewClientBuilder().WithScheme(scheme.Scheme). + WithObjects(imageArtifactFixture("uid-qwen"), imageRuntimeNode("v1.32.9", "containerd://1.7.0")). + Build() + models, blocked, err := (&InstanceReconciler{Client: cli}). + resolveInstanceModelVolumes(context.Background(), imageModelInstanceFixture("")) + require.NoError(t, err) + assert.Empty(t, blocked, "no node is known, so the pod is left to the scheduler's fate") + require.Contains(t, models, 0) + }) +} + +// TestConvertAdditionalVolumesImageDelivery renders an Instance model volume for an image artifact. +func TestConvertAdditionalVolumesImageDelivery(t *testing.T) { + w := &modelArtifactWeights{Render: &ModelDeploymentArtifactRender{ + Delivery: workercore.ModelDeploymentModelDeliveryImage, ImageReference: testImageReference, + }} + avs := []workercore.InstanceAdditionalVolume{ + {MountPath: "/models", ReadOnly: false, Model: &workercore.InstanceModelVolumeSource{ArtifactRef: core.LocalObjectReference{Name: "qwen"}}}, + } + + vols, mounts := convertAdditionalVolumes(avs, map[int]*modelArtifactWeights{0: w}) + require.Len(t, vols, 1) + require.NotNil(t, vols[0].Image, "the model volume is an image volume") + assert.Equal(t, testImageReference, vols[0].Image.Reference) + require.Len(t, mounts, 1) + assert.True(t, mounts[0].ReadOnly, "a model image is read-only whatever the entry says") + assert.Empty(t, mounts[0].SubPath, "the image root is mounted whole") +} + +// TestModelPlacementImagePreference walks the soft preference built from Node.status.images. +func TestModelPlacementImagePreference(t *testing.T) { + other := "registry.example.com/team/other@sha256:" + strings.Repeat("b", 64) + node := func(name, hostname string, images ...string) *core.Node { + n := &core.Node{ObjectMeta: meta.ObjectMeta{Name: name}} + if hostname != "" { + n.Labels = map[string]string{core.LabelHostname: hostname} + } + for _, img := range images { + n.Status.Images = append(n.Status.Images, core.ContainerImage{Names: []string{img}}) + } + + return n + } + t.Run("nodes reporting the image are named, sorted and capped", func(t *testing.T) { + var objs []ctrlcli.Object + for i := 0; i < modelPlacementMaxNodes+2; i++ { + objs = append(objs, node( + strings.Repeat("n", 1)+strconv.Itoa(i), "host-"+strconv.Itoa(modelPlacementMaxNodes+1-i), testImageReference)) + } + objs = append(objs, node("unrelated", "host-x", other)) + cli := ctrlfake.NewClientBuilder().WithScheme(scheme.Scheme).WithObjects(objs...).Build() + + term := modelPlacementImagePreference(context.Background(), cli, testImageReference) + require.NotNil(t, term) + assert.Equal(t, int32(modelPlacementPreferenceWeight), term.Weight) + values := term.Preference.MatchExpressions[0].Values + assert.Len(t, values, modelPlacementMaxNodes, "the list is capped") + assert.Equal(t, "host-0", values[0], "the order is hostname-sorted, not listing order") + }) + + t.Run("a node not reporting the image and a node without a hostname are left out", func(t *testing.T) { + cli := ctrlfake.NewClientBuilder().WithScheme(scheme.Scheme).WithObjects( + node("holder", "host-1", other, testImageReference), + node("other-image", "host-2", other), + node("no-hostname", "", testImageReference), + ).Build() + + term := modelPlacementImagePreference(context.Background(), cli, testImageReference) + require.NotNil(t, term) + assert.Equal(t, []string{"host-1"}, term.Preference.MatchExpressions[0].Values) + }) + + t.Run("no node reporting yields no term", func(t *testing.T) { + cli := ctrlfake.NewClientBuilder().WithScheme(scheme.Scheme).WithObjects( + node("other-image", "host-2", other), + ).Build() + + assert.Nil(t, modelPlacementImagePreference(context.Background(), cli, testImageReference)) + }) + + t.Run("another delivery yields no term", func(t *testing.T) { + pvc := &modelArtifactWeights{Render: &ModelDeploymentArtifactRender{ + Delivery: workercore.ModelDeploymentModelDeliveryPvc, ClaimName: "models", + }} + assert.Nil(t, pvc.placementPreference(context.Background(), + ctrlfake.NewClientBuilder().WithScheme(scheme.Scheme).Build())) + }) +} diff --git a/pkg/worker/controllers/worker/model_artifact_placement.go b/pkg/worker/controllers/worker/model_artifact_placement.go index 7b9588784..3411e9446 100644 --- a/pkg/worker/controllers/worker/model_artifact_placement.go +++ b/pkg/worker/controllers/worker/model_artifact_placement.go @@ -111,6 +111,15 @@ func resolveModelArtifactWeights( ClaimName: source.PersistentVolumeClaim.ClaimName, Path: source.PersistentVolumeClaim.Path, } + case source.Image != nil: + // The image is delivered whole by kubelet, wherever the Pod lands: no node affinity, no + // plugin, no filter, and nothing to size. The reference in the immutable spec is the only + // content identity, so the resolved status is claim-shaped and the KV identity falls back + // to the artifact's UID. + w.Render = &ModelDeploymentArtifactRender{ + Delivery: workercore.ModelDeploymentModelDeliveryImage, + ImageReference: source.Image.Reference, + } case source.HuggingFace != nil && (nodeOnly || modelArtifactDeliveryMode(ctx) == settings.ModelArtifactDeliveryNode): w.Render = &ModelDeploymentArtifactRender{ Delivery: workercore.ModelDeploymentModelDeliveryNode, diff --git a/pkg/worker/controllers/worker/model_artifact_placement_test.go b/pkg/worker/controllers/worker/model_artifact_placement_test.go index 67637a283..5d6b158d4 100644 --- a/pkg/worker/controllers/worker/model_artifact_placement_test.go +++ b/pkg/worker/controllers/worker/model_artifact_placement_test.go @@ -689,3 +689,30 @@ func TestModelDeploymentDeliveryWaitClears(t *testing.T) { }) } } + +// TestModelDeploymentArtifactImageDelivery reconciles a deployment against a resolved image +// artifact: the replicas mount the image volume, the weights read as not mounted yet, and the +// status names the Image delivery. +func TestModelDeploymentArtifactImageDelivery(t *testing.T) { + imageArtifact := artifactFixture("models", true, true) + imageArtifact.Spec.Source = workercore.ModelArtifactSource{ + Image: &workercore.ModelArtifactImageSource{Reference: testImageReference}, + } + + cli := newModelDeploymentClient(artifactDeploymentFixture(1), newRenderInstanceType(), imageArtifact) + _, err := reconcileModelDeployment(t, cli) + require.NoError(t, err) + + pods := listReplicas(t, cli) + require.Len(t, pods, 1) + vol := findVolume(&pods[0], modelDeploymentModelVolumeName) + require.NotNil(t, vol) + require.NotNil(t, vol.Image, "the weights are an image volume") + assert.Equal(t, testImageReference, vol.Image.Reference) + assert.Nil(t, findVolume(&pods[0], modelDeploymentModelCacheVolumeName)) + + md := getModelDeployment(t, cli) + assert.Equal(t, "WeightsNotMounted", ModelDeploymentConditionWeightsReady.GetReason(md)) + require.NotNil(t, md.Status.Model) + assert.Equal(t, workercore.ModelDeploymentModelDeliveryImage, md.Status.Model.Delivery) +} diff --git a/pkg/worker/controllers/worker/model_artifact_test.go b/pkg/worker/controllers/worker/model_artifact_test.go index ab537fc8c..aac70ce51 100644 --- a/pkg/worker/controllers/worker/model_artifact_test.go +++ b/pkg/worker/controllers/worker/model_artifact_test.go @@ -701,3 +701,36 @@ func TestEarliestResult(t *testing.T) { // ctrlclientWithWatch is the fake client's own interface, which the interceptor wraps. type ctrlclientWithWatch = ctrlcli.WithWatch + +// TestModelArtifactReconcileImage walks the image source's resolution: the first pass resolves +// with a claim-shaped status, no network ask, and no nodes aggregation, and the artifact never +// requeues. +func TestModelArtifactReconcileImage(t *testing.T) { + ma := testHubArtifact("") + ma.Spec.Source = workercore.ModelArtifactSource{Image: &workercore.ModelArtifactImageSource{ + Reference: testImageReference, + }} + env := newTestArtifactEnv(t, ma) + + // The first pass only writes the Resolving placeholder; the second answers. + _, err := env.r.Reconcile(context.Background(), + ctrl.Request{NamespacedName: ctrlcli.ObjectKey{Namespace: "team-a", Name: "qwen"}}) + require.NoError(t, err) + result, err := env.r.Reconcile(context.Background(), + ctrl.Request{NamespacedName: ctrlcli.ObjectKey{Namespace: "team-a", Name: "qwen"}}) + require.NoError(t, err) + assert.False(t, result.Requeue) + assert.Zero(t, result.RequeueAfter, "there is nothing to ask again: no registry, no revalidation") + + got := new(workercore.ModelArtifact) + require.NoError(t, env.cli.Get(context.Background(), + ctrlcli.ObjectKey{Namespace: "team-a", Name: "qwen"}, got)) + require.NotNil(t, got.Status.Resolved) + assert.NotNil(t, got.Status.Resolved.ResolvedTime) + assert.Empty(t, got.Status.Resolved.Revision, "an image source has no commit") + assert.Empty(t, got.Status.Resolved.ManifestDigest, "the reference is the whole identity") + assert.Nil(t, got.Status.Nodes, "no node ever aggregates an image") + assert.Equal(t, "True", ModelArtifactConditionResolved.GetStatus(got)) + assert.Equal(t, "Resolved", ModelArtifactConditionResolved.GetReason(got)) + assert.Equal(t, "False", ModelArtifactConditionDegraded.GetStatus(got)) +} diff --git a/pkg/worker/controllers/worker/model_deployment_artifact.go b/pkg/worker/controllers/worker/model_deployment_artifact.go index 2ee82a2dc..4562c6ddf 100644 --- a/pkg/worker/controllers/worker/model_deployment_artifact.go +++ b/pkg/worker/controllers/worker/model_deployment_artifact.go @@ -41,7 +41,7 @@ const ( // A nil one is a deployment that names no artifact, which renders exactly what it rendered before // the field existed. type ModelDeploymentArtifactRender struct { - // Delivery is Pvc, Engine or Node. + // Delivery is Pvc, Engine, Node or Image. Delivery workercore.ModelDeploymentModelDelivery // ArtifactName, ArtifactUID and ManifestDigest are what a node-delivered volume names, as hints @@ -54,6 +54,10 @@ type ModelDeploymentArtifactRender struct { ClaimName string Path string + // ImageReference is an image artifact's digest-pinned reference, rendered as a Kubernetes image + // volume that kubelet pulls and mounts read-only. + ImageReference string + // Repository and Revision are a hub artifact's repository and resolved commit, SecretName the // Secret holding its token or "", and SizeBytes its manifest's total size. Repository string @@ -107,6 +111,18 @@ func (a *ModelDeploymentArtifactRender) volumes(takeOver bool) ([]core.Volume, [ Name: modelDeploymentModelVolumeName, MountPath: ModelDeploymentModelMountPath, SubPath: a.Path, ReadOnly: true, }} + case a.Delivery == workercore.ModelDeploymentModelDeliveryImage: + // The pull policy is left at its default: a digest-pinned reference defaults to + // IfNotPresent, so a node that already holds the image does not fetch it again. No cache + // volume, no download environment: kubelet's own image store is the whole delivery. + return []core.Volume{{ + Name: modelDeploymentModelVolumeName, + VolumeSource: core.VolumeSource{Image: &core.ImageVolumeSource{ + Reference: a.ImageReference, + }}, + }}, []core.VolumeMount{{ + Name: modelDeploymentModelVolumeName, MountPath: ModelDeploymentModelMountPath, ReadOnly: true, + }} case takeOver: return nil, nil default: diff --git a/pkg/worker/controllers/worker/model_deployment_artifact_test.go b/pkg/worker/controllers/worker/model_deployment_artifact_test.go index f6eafc675..46e057336 100644 --- a/pkg/worker/controllers/worker/model_deployment_artifact_test.go +++ b/pkg/worker/controllers/worker/model_deployment_artifact_test.go @@ -3,6 +3,7 @@ package worker import ( "context" "fmt" + "strings" "testing" "github.com/stretchr/testify/assert" @@ -37,6 +38,55 @@ func testNodeArtifactRender() *ModelDeploymentArtifactRender { } } +func testImageArtifactRender() *ModelDeploymentArtifactRender { + return &ModelDeploymentArtifactRender{ + Delivery: workercore.ModelDeploymentModelDeliveryImage, + ImageReference: "registry.example.com/team/qwen@sha256:" + strings.Repeat("a", 64), + } +} + +// TestRenderModelDeploymentArtifactImageVolume renders an image artifact's delivery: one image +// volume, the digest-pinned reference, the model path read-only and whole, and none of the engine +// download's machinery — no cache volume, no download environment, no revision argument. +func TestRenderModelDeploymentArtifactImageVolume(t *testing.T) { + render := testImageArtifactRender() + cases := []struct { + name string + takeOver bool + }{ + {name: "an image artifact is mounted whole and read-only"}, + {name: "a take-over role still gets the image"}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + md := newRenderDeployment(func(md *workercore.ModelDeployment) { + if c.takeOver { + md.Spec.Roles[0].Command = []string{"sh", "-c", "serve"} + } + }) + pod := renderWithArtifact(t, md, render) + main := &pod.Spec.Containers[0] + + vol, mount := findVolume(pod, modelDeploymentModelVolumeName), findMount(main, modelDeploymentModelVolumeName) + require.NotNil(t, vol) + require.NotNil(t, mount) + require.NotNil(t, vol.Image, "the weights volume is an image volume") + assert.Equal(t, render.ImageReference, vol.Image.Reference) + assert.Empty(t, vol.Image.PullPolicy, "a digest-pinned reference keeps the default policy") + assert.Equal(t, ModelDeploymentModelMountPath, mount.MountPath) + assert.True(t, mount.ReadOnly) + assert.Empty(t, mount.SubPath, "the image root is mounted whole") + assert.Nil(t, findVolume(pod, modelDeploymentModelCacheVolumeName), + "kubelet's image store is the whole delivery, so there is no cache to size") + for _, env := range main.Env { + assert.NotEqual(t, modelDeploymentHFHomeEnv, env.Name, "no download environment") + assert.NotEqual(t, modelDeploymentHFTokenEnv, env.Name) + } + assert.NotContains(t, main.Command, modelDeploymentRevisionArg) + }) + } +} + func renderWithArtifact( t *testing.T, md *workercore.ModelDeployment, artifact *ModelDeploymentArtifactRender, ) *core.Pod { diff --git a/pkg/worker/controllers/worker/model_deployment_kv_identity_test.go b/pkg/worker/controllers/worker/model_deployment_kv_identity_test.go index 198ccee94..3fbee6129 100644 --- a/pkg/worker/controllers/worker/model_deployment_kv_identity_test.go +++ b/pkg/worker/controllers/worker/model_deployment_kv_identity_test.go @@ -25,6 +25,14 @@ func TestModelDeploymentKVIdentity(t *testing.T) { ma.UID = types.UID(uid) return ma } + image := func(uid string) *workercore.ModelArtifact { + ma := artifactFixture("models", true, true) + ma.Spec.Source = workercore.ModelArtifactSource{Image: &workercore.ModelArtifactImageSource{ + Reference: testImageReference, + }} + ma.UID = types.UID(uid) + return ma + } cases := []struct { name string a, b *workercore.ModelArtifact @@ -34,6 +42,8 @@ func TestModelDeploymentKVIdentity(t *testing.T) { {name: "two digests are two", a: hub(testArtifactDigest), b: hub("sha256:" + strings.Repeat("2", 64))}, {name: "one claim artifact is one identity", a: claim("uid-a"), b: claim("uid-a"), same: true}, {name: "two claim artifacts are two, even on one claim", a: claim("uid-a"), b: claim("uid-b")}, + {name: "one image artifact is one identity", a: image("uid-a"), b: image("uid-a"), same: true}, + {name: "two image artifacts naming one image are two, the UID stands in for no digest", a: image("uid-a"), b: image("uid-b")}, } for _, c := range cases { t.Run(c.name, func(t *testing.T) { diff --git a/pkg/worker/controllers/worker/model_placement_preference.go b/pkg/worker/controllers/worker/model_placement_preference.go index 50c54b9c9..8efb251df 100644 --- a/pkg/worker/controllers/worker/model_placement_preference.go +++ b/pkg/worker/controllers/worker/model_placement_preference.go @@ -38,18 +38,81 @@ type modelPlacementCandidate struct { } // placementPreference returns the preference for the Pods these weights are mounted by: the nodes -// holding a node-delivered digest ready to mount, or nil for any other delivery, a blocked resolution, -// or a digest no node holds. +// holding a node-delivered digest ready to mount, the nodes kubelet reports holding an image +// artifact's image, or nil for any other delivery, a blocked resolution, or a reference no node +// holds. // // It walks every NodeModelStore, so it is computed only by a pass that creates Pods, once for all of // them: that keeps a reconcile that creates nothing from paying for it, and gives every member of a // replica the same term, which Kueue needs because it builds the replica's PodSet from one member. func (w *modelArtifactWeights) placementPreference(ctx context.Context, cli ctrlcli.Reader) *core.PreferredSchedulingTerm { - if w == nil || w.Blocked || w.Render == nil || w.Render.Delivery != workercore.ModelDeploymentModelDeliveryNode { + if w == nil || w.Blocked || w.Render == nil { return nil } - return modelPlacementPreference(ctx, cli, w.Render.ManifestDigest) + switch w.Render.Delivery { + case workercore.ModelDeploymentModelDeliveryNode: + return modelPlacementPreference(ctx, cli, w.Render.ManifestDigest) + case workercore.ModelDeploymentModelDeliveryImage: + return modelPlacementImagePreference(ctx, cli, w.Render.ImageReference) + } + return nil +} + +// modelPlacementImagePreference returns the preferred node-affinity term naming the nodes kubelet +// reports holding an image artifact's reference, or nil when none does. +// +// A PREFERENCE, NEVER A FILTER, for the same reasons the digest preference is one: Kueue's +// topology-aware scheduling reads the term as a score among the nodes that already fit the Pod, so +// a wrong or stale entry costs at most one pull on another node. And a stale entry is likely: +// the report is kubelet's Node.status.images, bounded by --node-status-max-images (default 50), so +// a node running many images may never name this one, and a collected one lingers until the next +// report. The softness is what makes that affordable: absence only forfeits this round's locality +// gain. No size is read — the CRI size is the compressed image, measured well under the real +// occupancy, and capacity is nobody's question here. +// +// The order is hostname-sorted and then capped, because Node.status.images carries no +// referenced-now or last-used signal for an image: sorting on the names is the only ordering that +// depends on nothing but the candidates, so the members of one replica built from one snapshot +// carry the same list. +func modelPlacementImagePreference(ctx context.Context, cli ctrlcli.Reader, reference string) *core.PreferredSchedulingTerm { + logger := ctrllog.FromContext(ctx) + + nodes := new(core.NodeList) + if err := cli.List(ctx, nodes); err != nil { + logger.Error(err, "list nodes for an image placement preference; none is added") + return nil + } + + var hostnames []string + for i := range nodes.Items { + node := &nodes.Items[i] + holds := slices.ContainsFunc(node.Status.Images, func(img core.ContainerImage) bool { + return slices.Contains(img.Names, reference) + }) + if !holds { + continue + } + hostname := node.Labels[core.LabelHostname] + if hostname == "" { + continue + } + hostnames = append(hostnames, hostname) + } + if len(hostnames) == 0 { + return nil + } + slices.Sort(hostnames) + if len(hostnames) > modelPlacementMaxNodes { + hostnames = hostnames[:modelPlacementMaxNodes] + } + + return &core.PreferredSchedulingTerm{ + Weight: modelPlacementPreferenceWeight, + Preference: core.NodeSelectorTerm{MatchExpressions: []core.NodeSelectorRequirement{ + {Key: core.LabelHostname, Operator: core.NodeSelectorOpIn, Values: hostnames}, + }}, + } } // modelPlacementPreference returns the preferred node-affinity term naming the nodes that hold digest diff --git a/pkg/worker/webhooks/worker/instance_model_volume_test.go b/pkg/worker/webhooks/worker/instance_model_volume_test.go index 623db2b65..f6debcfff 100644 --- a/pkg/worker/webhooks/worker/instance_model_volume_test.go +++ b/pkg/worker/webhooks/worker/instance_model_volume_test.go @@ -2,6 +2,7 @@ package worker import ( "context" + "strings" "testing" "github.com/stretchr/testify/assert" @@ -73,6 +74,12 @@ func TestInstanceModelVolumeAdmitsEveryArtifactSource(t *testing.T) { } return ma } + imageArtifact := &workercore.ModelArtifact{ + ObjectMeta: meta.ObjectMeta{Name: "qwen"}, + Spec: workercore.ModelArtifactSpec{Source: workercore.ModelArtifactSource{ + Image: &workercore.ModelArtifactImageSource{Reference: "registry.example.com/team/qwen@sha256:" + strings.Repeat("a", 64)}, + }}, + } cases := []struct { name string objs []ctrlcli.Object @@ -81,6 +88,7 @@ func TestInstanceModelVolumeAdmitsEveryArtifactSource(t *testing.T) { }{ {name: "a claim artifact", objs: []ctrlcli.Object{artifact(true)}}, {name: "a hub artifact", objs: []ctrlcli.Object{artifact(false)}}, + {name: "an image artifact", objs: []ctrlcli.Object{imageArtifact}}, {name: "an artifact that does not exist yet"}, {name: "a hub artifact with a sub path", objs: []ctrlcli.Object{artifact(false)}, subPath: "x", wantErr: true}, } diff --git a/pkg/worker/webhooks/worker/model_artifact.go b/pkg/worker/webhooks/worker/model_artifact.go index f9fcba033..17b2561da 100644 --- a/pkg/worker/webhooks/worker/model_artifact.go +++ b/pkg/worker/webhooks/worker/model_artifact.go @@ -14,7 +14,9 @@ import ( ctrladmission "sigs.k8s.io/controller-runtime/pkg/webhook/admission" workercore "gpustack.ai/gpustack/api/worker/v1alpha1" + "gpustack.ai/gpustack/pkg/kubediscovery" "gpustack.ai/gpustack/pkg/kubemeta" + "gpustack.ai/gpustack/pkg/system" "gpustack.ai/gpustack/pkg/webhook" ) @@ -135,9 +137,12 @@ func validateModelArtifact(ma, old *workercore.ModelArtifact) field.ErrorList { if source.PersistentVolumeClaim != nil { members = append(members, "persistentVolumeClaim") } + if source.Image != nil { + members = append(members, "image") + } if len(members) != 1 { return field.ErrorList{field.Invalid(sourcePath, members, - "exactly one of huggingFace or persistentVolumeClaim is required")} + "exactly one of huggingFace, persistentVolumeClaim or image is required")} } var errs field.ErrorList @@ -146,11 +151,14 @@ func validateModelArtifact(ma, old *workercore.ModelArtifact) field.ErrorList { return field.ErrorList{field.Forbidden(sourcePath.Child("modelScope"), modelArtifactModelScopeMessage)} case source.HuggingFace != nil: errs = validateModelArtifactHuggingFace(source.HuggingFace, sourcePath.Child("huggingFace")) + case source.Image != nil: + errs = validateModelArtifactImage(source.Image, sourcePath.Child("image")) + errs = append(errs, validateModelArtifactImageVolume(sourcePath.Child("image"))...) default: errs = validateModelArtifactClaim(source.PersistentVolumeClaim, sourcePath.Child("persistentVolumeClaim")) } - return append(errs, validateModelArtifactPatterns(&ma.Spec, specPath, source.HuggingFace != nil)...) + return append(errs, validateModelArtifactPatterns(&ma.Spec, specPath, source.HuggingFace != nil, source.Image != nil)...) } // modelArtifactMaxPatterns and modelArtifactMaxPatternLength bound a list of patterns and each of @@ -160,9 +168,14 @@ const ( modelArtifactMaxPatternLength = 256 ) -// validateModelArtifactPatterns accepts allow and ignore patterns on a hub source only: a claim's -// content is the user's and is mounted whole. -func validateModelArtifactPatterns(spec *workercore.ModelArtifactSpec, specPath *field.Path, hub bool) field.ErrorList { +// validateModelArtifactPatterns accepts allow and ignore patterns on a hub source only: every +// other source's content is the user's and is mounted whole. +func validateModelArtifactPatterns(spec *workercore.ModelArtifactSpec, specPath *field.Path, hub, image bool) field.ErrorList { + // whole is what the refusal says is mounted whole instead of being selected from. + whole := "a claim's directory is mounted whole" + if image { + whole = "an image is mounted whole" + } var errs field.ErrorList for _, list := range []struct { name string @@ -177,7 +190,7 @@ func validateModelArtifactPatterns(spec *workercore.ModelArtifactSpec, specPath continue case !hub: errs = append(errs, field.Forbidden(path, - "patterns select files of a hub source; a claim's directory is mounted whole")) + "patterns select files of a hub source; "+whole)) continue case len(list.patterns) > modelArtifactMaxPatterns: errs = append(errs, field.TooMany(path, len(list.patterns), modelArtifactMaxPatterns)) @@ -235,3 +248,55 @@ func validateModelArtifactClaim(claim *workercore.ModelArtifactPersistentVolumeC return errs } + +// modelArtifactClusterVersion is the Kubernetes version the capability rules read. It is a +// variable so a test can choose a version without a cluster: the snapshot's Configure ignores +// later calls, so a test cannot re-point the snapshot itself. +var modelArtifactClusterVersion = func() kubediscovery.Version { + return system.LoopbackKubeVersion.Get() +} + +// validateModelArtifactImageVolume refuses an image source where the cluster cannot serve image +// volumes. The ImageVolume feature is on by default from Kubernetes 1.35; an apiserver that does +// not know the volumes[].image field drops it without an error, so the weights would be silently +// absent. It reads the startup version snapshot, and a version it cannot read refuses: a version +// nobody can read must not unlock the capability. +func validateModelArtifactImageVolume(path *field.Path) field.ErrorList { + version := modelArtifactClusterVersion() + if kubediscovery.SupportsFeature(&version, kubediscovery.FeatureImageVolume) { + return nil + } + + return field.ErrorList{field.Forbidden(path, fmt.Sprintf( + "this cluster's Kubernetes version %q does not serve image volumes: the ImageVolume feature is "+ + "on by default from 1.35, and 1.33 or 1.34 would need the gate opened on the apiserver, which "+ + "this version does not offer; deliver the weights from a hub or a PersistentVolumeClaim, or "+ + "upgrade the cluster", version.GitVersion))} +} + +// modelArtifactDigestPattern is the digest half of a pinned image reference: "sha256:" and 64 +// lowercase hex. +var modelArtifactDigestPattern = regexp.MustCompile(`^sha256:[a-f0-9]{64}$`) + +// validateModelArtifactImage refuses a reference that does not pin a digest. A tag is mutable, so +// one artifact could deliver different weights on different pulls, and the identity a frozen +// reference pins would be nothing. +func validateModelArtifactImage(image *workercore.ModelArtifactImageSource, path *field.Path) field.ErrorList { + var errs field.ErrorList + reference := image.Reference + ref, digest, pinned := strings.Cut(reference, "@") + switch { + case !pinned: + errs = append(errs, field.Invalid(path.Child("reference"), reference, + "must pin the image to its digest: registry/repository@sha256:<64 hex>; a tag is mutable, "+ + "so one artifact could deliver different weights on different pulls")) + case !modelArtifactDigestPattern.MatchString(digest): + errs = append(errs, field.Invalid(path.Child("reference"), reference, + `must pin the digest as "sha256:" and 64 lowercase hex`)) + case ref == "" || strings.ContainsAny(ref, "@ \t\n\r"): + errs = append(errs, field.Invalid(path.Child("reference"), reference, + "must be one image reference without whitespace, followed by @ and the digest")) + } + + return errs +} diff --git a/pkg/worker/webhooks/worker/model_artifact_test.go b/pkg/worker/webhooks/worker/model_artifact_test.go index 4e45350ad..4e90129ca 100644 --- a/pkg/worker/webhooks/worker/model_artifact_test.go +++ b/pkg/worker/webhooks/worker/model_artifact_test.go @@ -12,6 +12,7 @@ import ( meta "k8s.io/apimachinery/pkg/apis/meta/v1" workercore "gpustack.ai/gpustack/api/worker/v1alpha1" + "gpustack.ai/gpustack/pkg/kubediscovery" ) func newTestHubArtifact(repository, revision string) *workercore.ModelArtifact { @@ -32,6 +33,96 @@ func newTestClaimArtifact(claim, path string) *workercore.ModelArtifact { } } +func newTestImageArtifact(reference string) *workercore.ModelArtifact { + return &workercore.ModelArtifact{ + ObjectMeta: meta.ObjectMeta{Namespace: "team-a", Name: "qwen"}, + Spec: workercore.ModelArtifactSpec{Source: workercore.ModelArtifactSource{ + Image: &workercore.ModelArtifactImageSource{Reference: reference}, + }}, + } +} + +// TestModelArtifactWebhookImageSource covers the image member's shape rules and its capability +// gate, both directions, through the swappable version seam: the snapshot's Configure ignores +// later calls, so a test cannot re-point the snapshot itself. +func TestModelArtifactWebhookImageSource(t *testing.T) { + defaultVersion := modelArtifactClusterVersion + t.Cleanup(func() { modelArtifactClusterVersion = defaultVersion }) + supported := func() kubediscovery.Version { return kubediscovery.Version{GitVersion: "v1.35.5"} } + unsupported := func() kubediscovery.Version { return kubediscovery.Version{GitVersion: "v1.32.9"} } + + cases := []struct { + name string + in *workercore.ModelArtifact + reference string + version func() kubediscovery.Version + wantField string + }{ + {name: "a digest-pinned reference", reference: "registry.example.com/team/qwen@sha256:" + strings.Repeat("a", 64), version: supported}, + {name: "a host and a port", reference: "localhost:5500/qwen@sha256:" + strings.Repeat("a", 64), version: supported}, + {name: "a bare repository", reference: "qwen@sha256:" + strings.Repeat("a", 64), version: supported}, + { + name: "an image source beside a hub source", version: supported, wantField: "spec.source", + reference: "", + in: func() *workercore.ModelArtifact { + ma := newTestImageArtifact("registry/qwen@sha256:" + strings.Repeat("a", 64)) + ma.Spec.Source.HuggingFace = &workercore.ModelArtifactHubSource{Repository: "owner/repo", Revision: "main"} + return ma + }(), + }, + {name: "the image volume gate off", reference: "registry/qwen@sha256:" + strings.Repeat("a", 64), version: unsupported, wantField: "spec.source.image"}, + { + name: "an unreadable cluster version", reference: "registry/qwen@sha256:" + strings.Repeat("a", 64), + version: func() kubediscovery.Version { return kubediscovery.Version{GitVersion: ""} }, + wantField: "spec.source.image", + }, + {name: "a tag instead of a digest", reference: "registry.example.com/team/qwen:v1", version: supported, wantField: "spec.source.image.reference"}, + {name: "no digest at all", reference: "registry.example.com/team/qwen", version: supported, wantField: "spec.source.image.reference"}, + {name: "an uppercase digest", reference: "registry/qwen@sha256:" + strings.Repeat("A", 64), version: supported, wantField: "spec.source.image.reference"}, + {name: "a short digest", reference: "registry/qwen@sha256:" + strings.Repeat("a", 63), version: supported, wantField: "spec.source.image.reference"}, + {name: "the wrong algorithm", reference: "registry/qwen@sha512:" + strings.Repeat("a", 128), version: supported, wantField: "spec.source.image.reference"}, + {name: "an empty reference", reference: "", version: supported, wantField: "spec.source.image.reference"}, + {name: "only a digest", reference: "@sha256:" + strings.Repeat("a", 64), version: supported, wantField: "spec.source.image.reference"}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + modelArtifactClusterVersion = c.version + in := c.in + if in == nil { + in = newTestImageArtifact(c.reference) + } + + _, err := new(ModelArtifactWebhook).ValidateCreate(context.Background(), in) + if c.wantField == "" { + require.NoError(t, err) + return + } + require.Error(t, err) + assertInvalidField(t, err, c.wantField) + }) + } + + t.Run("patterns are refused on an image source", func(t *testing.T) { + modelArtifactClusterVersion = supported + _, err := new(ModelArtifactWebhook).ValidateCreate(context.Background(), + withPatterns(newTestImageArtifact("registry/qwen@sha256:"+strings.Repeat("a", 64)), []string{"*.json"}, nil)) + require.Error(t, err) + assertInvalidField(t, err, "spec.allowPatterns") + }) + + t.Run("the gate refusal names the cluster's version and what is missing", func(t *testing.T) { + modelArtifactClusterVersion = unsupported + _, err := new(ModelArtifactWebhook).ValidateCreate(context.Background(), + newTestImageArtifact("registry/qwen@sha256:"+strings.Repeat("a", 64))) + require.Error(t, err) + status, ok := err.(kerrors.APIStatus) + require.True(t, ok) + message := status.Status().Details.Causes[0].Message + assert.Contains(t, message, "does not serve image volumes") + assert.Contains(t, message, "v1.32.9") + }) +} + func TestModelArtifactWebhookDefault(t *testing.T) { cases := []struct { name string diff --git a/pkg/worker/webhooks/worker/model_prefetch.go b/pkg/worker/webhooks/worker/model_prefetch.go index cdee8345b..dfcb492d5 100644 --- a/pkg/worker/webhooks/worker/model_prefetch.go +++ b/pkg/worker/webhooks/worker/model_prefetch.go @@ -109,6 +109,11 @@ func (r *ModelPrefetchWebhook) validateModelPrefetch(ctx context.Context, pf *wo return field.ErrorList{field.InternalError(field.NewPath("spec", "artifactRef", "name"), fmt.Errorf("get model artifact %q: %w", pf.Spec.ArtifactRef.Name, err))} } + if artifact.Spec.Source.Image != nil { + return field.ErrorList{field.Forbidden(field.NewPath("spec", "artifactRef"), + "an image artifact delivers on demand through kubelet and never enters the node cache: "+ + "there is nothing to warm, no budget to count against, and no pinning to grant")} + } binding := new(workercore.ModelStoreBinding) if err := r.APIReader.Get(ctx, ctrlcli.ObjectKey{Namespace: pf.Namespace, Name: pf.Spec.BindingRef.Name}, binding); err != nil { diff --git a/pkg/worker/webhooks/worker/model_prefetch_test.go b/pkg/worker/webhooks/worker/model_prefetch_test.go index 0f8593f1f..cc889156c 100644 --- a/pkg/worker/webhooks/worker/model_prefetch_test.go +++ b/pkg/worker/webhooks/worker/model_prefetch_test.go @@ -90,6 +90,15 @@ func TestModelPrefetchAdmission(t *testing.T) { assert.Contains(t, err.Error(), "Not found") }) + t.Run("a prefetch of an image artifact is refused", func(t *testing.T) { + objs := newModelPrefetchFixture(true, "10Gi", 5<<30) + objs[2].(*workercore.ModelArtifact).Spec.Source.Image = &workercore.ModelArtifactImageSource{Reference: "registry/qwen@sha256:" + strings.Repeat("a", 64)} + r := newModelPrefetchWebhook(objs...) + _, err := r.ValidateCreate(context.Background(), newModelPrefetch()) + require.Error(t, err) + assert.Contains(t, err.Error(), "never enters the node cache") + }) + t.Run("a prefetch naming a missing grant is refused", func(t *testing.T) { r := newModelPrefetchWebhook(newModelPrefetchFixture(true, "10Gi", 5<<30)...) pf := newModelPrefetch() diff --git a/specs/2026-09-27-model-image-source.md b/specs/2026-09-27-model-image-source.md new file mode 100644 index 000000000..86d83fa0c --- /dev/null +++ b/specs/2026-09-27-model-image-source.md @@ -0,0 +1,551 @@ +# Spec: OCI image source for ModelArtifact (S7) + +Status: Shipped +Type: Feature + +## Summary + +S7 adds an `image` source to `ModelArtifact`: the union's fourth field and the third member +admission accepts (after `huggingFace` and `persistentVolumeClaim`; `modelScope` stays reserved +and refused). Weights are identified by a digest-pinned OCI image reference and delivered by +mounting a Kubernetes image volume +(`volumes[].image`, the ImageVolume feature). Delivery is on demand — kubelet pulls the image on +the node that needs it — and nothing enters the model-manager plugin's cache chain: no +`NodeModelStore` entry, no manifest, no watermark accounting. The version floor for the feature is +its own, separate from the project's 1.29 functional floor, and admission refuses the source +where the cluster cannot serve it. + +## Motivation + +Weights already exist as OCI images in many estates: users build them from model repositories +(Dockerfile COPY, oras, skopeo) or receive them from a registry pipeline. Today such a user must +let the operator re-download the weights from a hub, or publish a PVC — both re-introduce a copy +step and a second source of truth. An `image` source lets the artifact point at the image the +user already trusts, and lets the platform deliver it with the runtime's own mechanics: content +addressing by digest, layer sharing across images on the node, and registry mirrors (pull-through +caches, P2P prefetchers) that the runtime already honors. + +### Goals + +- A tenant can declare an `image`-source `ModelArtifact` and mount it from a `ModelDeployment` + and an `Instance`, read-only, at the fixed model path. +- The artifact's identity is the image reference the user pinned; one artifact always delivers + the same bytes. +- Admission refuses the source where the platform cannot honor it, naming the version floor and + the reason (two-layer-floor degradation, never a silent field drop). +- The double-storage cost, the kubelet image-GC interaction, and the CRI size accounting gap are + documented with the measured numbers so operators can size disks and predict re-pulls. + +### Non-Goals + +- Image sources do **not** enter the node cache plugin chain: no `NodeModelStore` entry, no + plugin mount, no `ModelPrefetch` — S6's warm-up path serves the plugin's cache only — no + budget accounting, no peer sync (S5). Kubelet's own image store and image GC govern them. +- No registry client in the operator: no tag→digest resolution, no existence check, no content + verification of the image. The user's build is the authority (see "The digest and the + verification promise"). +- No `path`/`subPath` inside the image (would raise the runtime floor to containerd ≥ 2.2); the + image root is mounted whole, so the build must put the weights at the image root. +- No containerd node configuration changes (`discard_unpacked_layers` is documented as a node + administrator's option, never touched by the operator). +- Dragonfly integration (S8). Mirrors are documented only. +- ModelScope and external providers unchanged. + +## Proposal + +One API extension and three wiring points: + +- `ModelArtifactSource.Image *ModelArtifactImageSource` — the union's fourth field and the third + member admission accepts: + `{reference: "registry/repository@sha256:<64 hex>"}`. Admission requires the digest form (a + mutable tag would let one artifact deliver different weights on different pulls, breaking the + identity contract every frozen reference and the KV identity rely on) and refuses source + selection patterns, as for a claim. Admission also refuses the source on any cluster whose + apiserver is older than the ImageVolume default-on floor (see the capability gate). +- Delivery: a new `ModelDeploymentModelDeliveryImage` value. An image-source artifact always + delivers `Image`, whatever `model-artifact-delivery-mode` says (the Setting governs hub + sources' Engine/Node choice only). Both consumers render the image volume; `Instance` + additionally pre-checks the pinned node's runtime (below). +- Resolution: no network, like a claim. The first reconcile sets `Resolved=True` with + `status.resolved` carrying only `resolvedTime`; `revision`, `manifestDigest`, `fileCount` and + `sizeBytes` stay absent, `status.nodes` stays absent, and there is no revalidation. The + reference in the immutable spec is the only content identity. +- Placement preference: for `Image` delivery, the soft + node preference is computed from `Node.status.images` — the nodes kubelet already reports + holding the reference — reusing the digest preference's shape (weight 100, ≤16 hostnames, + deterministic order, soft only). + +### Measured facts this design relies on (PoC-K) + +All `[跑]` readings were taken on Kubernetes 1.35.7 (apiserver 1.35.6) / containerd 2.2.6, Ubuntu +24.04, CPU nodes (2 vCPU, 8 GB, ~96 GiB ext4 root, nodefs and imagefs on one disk), an in-cluster +registry exposed to the node over hostNetwork; the model was `Qwen/Qwen2.5-7B-Instruct` at commit +`a09a35458c702b33eeacc393d103063234e8bc28` (16 files, 15,242,807,270 B), built one safetensors +shard per layer: + +- **Image volumes work end to end (K-1).** A ~15 GB model image mounts through + `volumes[].image`; every file matches the Hub copy byte for byte (sha256 equals the LFS oid). + The mount is a read-only overlay of the layer snapshots. +- **Double storage (K-2).** The node keeps both the compressed layer blobs and the unpacked + snapshots (`discard_unpacked_layers=false`, the default): content store +12,074,816,089 B and + snapshots +15,242,807,270 B for the gzip image — **1.79× the raw weights**; an uncompressed-layer + image measures **2.00×** (it saves decompression time, not disk). Pull of the 12.07 GB image + took 12m43s (~20 MB/s over the in-cluster registry on a 2-vCPU node). Kubelet accounts the CRI + image size — the compressed size — so its books show ~12.07 GB where the disk really holds + ~27.3 GB (**~2.26×** under gzip). Capacity judgments must not use the CRI size. +- **PSA (K-3, supersedes the earlier "Restricted refuses image volumes" claim).** On an apiserver + ≥ 1.33, PSA Restricted **admits** `image` volumes (kubernetes#130394, merged 2025-02-24, + milestone 1.33, not backported); on ≤ 1.32 Restricted refuses them. `enforce-version` makes no + difference either way (the check carries a single `MinimumVersion: 1.0`). Baseline admits the + volume wherever the apiserver knows the field. On the clusters this feature targets (≥ 1.35 + apiserver), PSA is not an obstacle, including for restricted tenants. +- **Image GC interaction (K-4).** While a Pod mounts the image, kubelet does not collect it + (verified 1.35.7 / 2.2.6; the in-use protection is in kubelet's image GC manager since 1.31 and + needs containerd ≥ 2.1). After the last Pod referencing it goes away, the image can be collected + in the next GC round (cycle: 5 minutes; measured release→collection: 4m2s) — `imageMinimumGCAge` + does not protect it, because kubelet compares against the image's *first-detected* time, not its + release time. Scaling a deployment to zero and back can therefore re-pull a dozen GB. The only + retention is another reference (e.g. a resident Pod mounting it; a completed Pod does not count — + the in-use check reads *running* containers). A model image pushing the disk past + `imageGCHighThresholdPercent` (default 85) also drags down every other unused image on the node + in the same rounds (measured: 18 images collected alongside). +- **Layer sharing.** Two images sharing layer digests share snapshots on the node and containerd + skips the download (measured: 112 ms and +11.8 MB where a cold node paid 2m10s and +7.89 GB). +- **Version line.** ImageVolume: 1.31 alpha (off), 1.33 beta (off), **1.35 beta on by default**, + 1.36 GA (`LockToDefault`). containerd supports image volumes from **2.1** (subPath would need + ≥ 2.2 — excluded by scope). Kubelet's in-use GC protection needs kubelet ≥ 1.31 and + containerd ≥ 2.1. + +### User Stories + +#### Story 1 +As an MLOps engineer, my weights are already in our internal registry as an OCI image, built and +scanned by our pipeline. I declare a `ModelArtifact` with that digest-pinned reference and serve it +with a `ModelDeployment`; the node pulls the image it already trusts, and nothing re-downloads +anything from a hub. + +#### Story 2 +As a platform operator, I run a pull-through registry mirror (Harbor) or a P2P prefetcher +(Dragonfly, Spegel) under containerd. Image-source deliveries ride the mirror like every other +pull, so fleet-wide cold starts amortize without any operator involvement. + +#### Story 3 +As a tenant in a Restricted namespace, I mount an image-source artifact without asking anyone: +on a supported apiserver, Restricted admits the volume. + +#### Story 4 +As a platform operator on a 1.32 cluster, I try to create an image-source artifact and get a +refusal that names the 1.35 floor and the alternatives, instead of a Pod that silently mounts +nothing. + +### Core Features & Acceptance Criteria + +**F1 API: the `image` source union member.** `ModelArtifactSource.Image` +(`*ModelArtifactImageSource`) with one required field, `Reference` (1–1024 characters). The +artifact's spec stays immutable; an image source is created, not edited, like every other member. + +- AC1: the union still admits exactly one member; the exactly-one refusal message names the + three admitted members (`huggingFace`, `persistentVolumeClaim`, `image`), while `modelScope` + keeps its own reserved-and-refused refusal as today. +- AC2: a reference that is not `@sha256:<64 lowercase hex>` is refused at + admission, with the message naming why the digest is required (a tag is mutable, so one + artifact could deliver different weights on different pulls; reopening tags later is a webhook + change, not a schema change). +- AC3: `allowPatterns`/`ignorePatterns` on an image source are refused with the claim source's + wording (an image is mounted whole; there is no listing to select from). +- AC4: **capability gate.** Creating an image-source artifact on a cluster whose apiserver + version is below 1.35 is refused, with a message naming the floor, that 1.33/1.34 need the + feature gate opened by an administrator (not offered as an opt-in in this version; reopening it + later is a webhook change), and the alternatives (hub or PVC sources, or upgrading). At 1.35 or + above the source is admitted. The check reads the worker's startup version snapshot + (`system.LoopbackKubeVersion` + a new `kubediscovery.FeatureImageVolume`), costing no API call; + both directions are unit-tested through the version seam. + +**F2 Resolution without a registry client.** The reconciler resolves an image source on its first +pass, with no network I/O. + +- AC5: `Resolved` becomes True with reason `Resolved` and the condition message stating that the + image is the user's and the operator does not read it; `status.resolved` carries only + `resolvedTime`; `revision`, `manifestDigest`, `fileCount`, `sizeBytes` stay absent (the CRD + field docs are updated to say so for image sources). +- AC6: `status.nodes` stays absent; no revalidation pass is scheduled; `Degraded` is False + (`Healthy`) from the first successful pass. + +**F3 ModelDeployment delivery (`Image`).** A deployment referencing an image-source artifact +renders the weights as an image volume on every managed role's engine container. + +- AC7: the weights volume is `volumes[].image` with `reference` = the spec's reference and the + default pull policy (IfNotPresent, omitted), mounted read-only at + `/var/lib/gpustack/model` under the existing weights volume name; the role volume admission + rules around the model paths apply unchanged. +- AC8: no cache `emptyDir`, no `HF_*`/proxy env, no `--revision` argument, no ephemeral-storage + raise; the served model path is the mount path, exactly as for a claim; `--served-model-name` + rules unchanged. +- AC9: `status.model.delivery` = `Image`; `WeightsReady` follows the claim/node mounted path + (`PodReadyToStartContainers`); a failed pull keeps `WeightsNotMounted` and kubelet's own event + names the pull error (documented limitation: the condition does not distinguish a slow pull + from a failed one). +- AC10: `model-artifact-delivery-mode` does not affect image sources (always `Image`); the + Setting's change watch does not roll image-delivered deployments by itself. + +**F4 Instance delivery.** An Instance model volume referencing an image-source artifact mounts +the image volume read-only, with no `subPath`. + +- AC11: the CSIDriver check (node delivery's plugin dependency) does not run for image delivery — + an Instance's image volume needs no plugin; the Instance creates its Pod while the reference is + resolved. +- AC12: **node capability pre-check.** Before rendering, the Instance reconciler reads the pinned + node's `status.nodeInfo.kubeletVersion` and `containerRuntimeVersion`; kubelet < 1.35 or + containerd < 2.1 blocks Pod creation with a reason naming the node and the two floors (a Pod + whose image volume the node cannot run otherwise fails late, in a kubelet event the Instance + status does not surface). A node whose versions cannot be parsed blocks the same way. (Bodies + of this check are read from a Node the reconciler already reads for hostname pinning; no extra + watch.) + +**F5 Placement preference from `Node.status.images`** + +- AC13: for `Image` delivery, the soft preference names the hostnames of the nodes whose + `status.images[].names` include the exact reference, reusing the digest preference's weight + (100), cap (16 hostnames) and determinism; no required affinity is ever added. The order is + hostname-sorted: `Node.status.images` carries no referenced-now or last-used signal, so the + digest preference's serving-now ordering has nothing to sort on here, and sorting on the names + is the only order that depends on nothing but the candidates. A stale positive (kubelet not yet + reporting a collection) costs one re-pull, the same softness rationale as the digest + preference. + +**F6 ModelPrefetch refuses image sources.** + +- AC14: creating a `ModelPrefetch` whose artifact's source is `image` is refused at admission, + with a message that image sources deliver on demand through kubelet and do not enter the node + cache this version warms (no budget, no pinning, no warm-up Pod). + +**F7 Docs.** + +- AC15: a new docs page (working title `docs/reference/model-image-source.md`, final name per the + docs skill's routing) carries: the source's shape and the digest contract; the build recipe + (weights at the image root, one shard per layer, pinned upstream commit) with the PoC-K + byte-equality example; the version line and runtime floors; the PSA paragraph (≥ 1.33 Restricted + admits, ≤ 1.32 refuses, `enforce-version` irrelevant; Baseline admits); double storage + (1.79×/2.00×/2.26× with their exact conditions and why the CRI size must not be used for + capacity); the three GC rules (in-use protection; release→next-round collection and the + scale-to-zero re-pull; the 85% cascade clearing other unused images); registry mirrors + (Harbor pull-through, Spegel and Dragonfly as configured containerd mirrors — documented as + unmeasured, Dragonfly integration itself being S8 scope); private registries via the role's / + Instance's existing `imagePullSecrets`. +- AC16: `docs/reference/model-artifact.md` gains the union member (resource example, patterns + note, delivery table row, status row, requirements bullet) and links the new page; + `docs/README.md` indexes the new page. No `docs/settings.md` change: the feature adds no + Setting. + +### Notes / Constraints / Caveats + +- **Two-layer floor.** The project floors stay put: functional 1.29 (bundled Kueue), install 1.23. + The image source is opt-in and carries its own floors — apiserver ≥ 1.35 for the default-on + beta gate, kubelet ≥ 1.35 (1.33/1.34 with an admin-opened gate, not opt-in-able here), + containerd ≥ 2.1 — and admission refuses below the apiserver floor, refusing the source and + naming the reason rather than degrading silently; the core path is untouched. +- **The digest and the verification promise.** The OCI digest pins the *image's* manifest bytes. + It is not the artifact manifest digest of a hub source (`status.resolved.manifestDigest` stays + absent), and the two are never claimed equal. The operator verifies nothing about the image's + content in this version: it does not read the registry, so "this image contains the weights" + is the user's assertion, and the operator's guarantee is only that the Pod mounts exactly the + image the immutable spec names. The documented build recipe (from the PoC) shows how to produce + an image whose content matches a hub commit byte for byte, and the PoC verified that equality + for its fixture, but no code enforces or checks it. KV identity therefore falls back to the + artifact-UID hash (the claim source's behavior): two artifacts naming the same image never + share KV blocks — safe, merely not deduplicated. +- **Weights at the image root.** No `path`/`subPath` is offered (subPath would floor containerd + at 2.2); the mount is the image root, so the engine sees whatever the image's root holds. The + build contract ("weights at `/`") is documentation, not validation. +- **Gate explicitly off on a ≥ 1.35 apiserver** (possible until 1.36 GA) would make the apiserver + drop the volume field silently — the exact hazard `FeatureNativeSidecar` documents. This + version does not probe for it (a dry-run canary was considered and rejected as v1 + over-engineering): the failure mode is a Pod without the weights volume, which the engine + fails on visibly, and the consumer reports `WeightsNotMounted`. Documented as a known + limitation; 1.36 GA removes it. +- **Mutability of `Node.status.images`.** The image list is kubelet's periodic report; a collected + image can linger in it and a freshly pulled one can lag. It is also bounded: kubelet reports at + most `--node-status-max-images` (default 50) entries, so a node running many images may never + name the model image at all. The preference is soft exactly so that every one of these + absences costs only this round's locality gain — one extra re-pull, never a placement veto — + and no code compensates for them. `Node.status.images[].sizes` (CRI, compressed) is + deliberately unused — PoC-K measured it at ~44% of real occupancy. +- **Names.** The delivery value is `Image`; no KV-cache vocabulary is touched. +- **The PoC environment is gone.** S5's test cluster was scheduled for deletion at S5's close; no + reading above depends on it. The e2e environment is adjudicated: local kind, gated on + verification (see Open Questions). + +### Boundaries + +- **Always:** keep the plugin chain out of image delivery (no `NodeModelStore` entry, no CSI + volume, no budget); keep the reference digest-pinned; keep the preference soft; keep the + capability refusal's message self-sufficient (floor + alternative). +- **Ask first:** any tag-reference support (needs resolution or a mutable-identity waiver — + product behavior); any `path`/`subPath` (raises the runtime floor); any image-content + verification (a registry client — new dependency surface). +- **Never:** let an image source produce a `NodeModelStore` write, a prefetch warm-up Pod, or a + `status.resolved.manifestDigest`; put containerd configuration advice anywhere but docs; + reference the task-tracking report from repo documents. + +### Risks and Mitigations + +- **Silent volume drop (gate off / very old apiserver)** → the create-time capability gate closes + the ≤ 1.32 direction; the ≥ 1.35 gate-off direction is documented (above) and self-revealing + (the engine cannot start without weights). Revisit a dry-run probe if it bites in the field. +- **GC re-pull churn on scale-to-zero** → documented with the measured numbers; the placement + preference reduces cross-node re-pulls; retention via a resident reference is the user's + documented option. A `pinned`-style retention would need CRI pinning or plugin ownership — out + of scope (Non-Goals). +- **Disk sizing surprises (double storage)** → the docs page carries the three multipliers with + their conditions and the CRI-size warning; nothing in code assumes image size. +- **Private registry pulls fail with imagePullBackOff** → documented: set the role's / + Instance's `imagePullSecrets`; the consumer shows `WeightsNotMounted` with the Pod's events as + the diagnostic path. +- **Mixed-version fleets** → the Instance pre-check converts a late kubelet error into a blocked + Instance with a named node and floors; the MD path (node unknown at render) relies on the + same kubelet event surfacing, documented. + +## Design Details + +### Commands + +- `make generate` after the API edit (deepcopy, CRD, openapi, protobuf, applyconfigurations). + The generator fails fast unless the checkout's path ends in `gpustack.ai/gpustack`, so it runs + in a worktree laid out that way, never hand-editing generated files; a regeneration that leaves + `git status` clean against the tested commit proves the committed generated tree is current. +- `make lint` (+ `make lint docs` for the docs change). The chart is untouched (CRDs are + generated into the chart tree by `make generate` — regenerate, don't hand-edit; if the + generated chart CRD diff trips `make lint chart`, run it too). +- Package tests: `go test ./api/worker/... ./pkg/worker/webhooks/worker/... + ./pkg/worker/controllers/worker/... ./pkg/kubediscovery/...`. +- Environments. Unit tests and both lint targets run on the local host. e2e runs on a local kind + cluster built by the repository's own e2e harness (it builds the operator image for the cluster + and loads it); no remote build host is in this spec's path. The one remote-shaped possibility — + a provider CPU cluster if the kind verification below fails — would be provisioned through the + task's own channel, not by this spec's commands. + +### Project Structure + +- `api/worker/v1alpha1/model_artifact.go` — `ModelArtifactImageSource` + the union member + (field docs: identity = the pinned reference; resolved fields absent). +- `api/worker/v1alpha1/model_deployment.go` — `ModelDeploymentModelDeliveryImage`. +- `pkg/kubediscovery/feature.go` — `FeatureImageVolume` (floor 1.35). +- `pkg/worker/webhooks/worker/model_artifact.go` — union, digest-pin, patterns, capability gate. +- `pkg/worker/webhooks/worker/model_prefetch.go` — image-source refusal. +- `pkg/worker/controllers/worker/model_artifact.go` — `reconcileImage` (no network). +- `pkg/worker/controllers/worker/model_artifact_placement.go` — weights resolution `Image` branch. +- `pkg/worker/controllers/worker/model_deployment_artifact.go` — `Image` render branch (volume + + mount + model path), `ModelDeploymentArtifactRender.ImageReference`. +- `pkg/worker/controllers/worker/model_deployment.go` — pass-through wiring; status observation + reuses the claim/node mounted path. +- `pkg/worker/controllers/worker/instance.go` — `convertAdditionalVolumes` `Image` branch; the + node capability pre-check. +- `pkg/worker/controllers/worker/model_placement_preference.go` — `Node.status.images` preference. +- `docs/reference/model-image-source.md` (new), `docs/reference/model-artifact.md`, + `docs/README.md`. +- `.agents/skills/gpustack-operator-e2e/cases/case-.sh` + the SKILL.md row — `` taken + at my-ship as main's then-max + 1, after the final rebase. + +### Code Style + +Follow repo conventions; illustrative pieces only. + +```go +// The union's fourth field, the third member admission accepts +// (api/worker/v1alpha1/model_artifact.go). +// Image is an OCI image reference holding the weights. The reference must pin a digest: +// "registry/repository@sha256:<64 hex>". The digest is the artifact's whole identity — the +// operator never reads the registry, so it verifies nothing about the image's content, and +// status.resolved stays claim-shaped (no revision, no manifestDigest). +// +// +optional +Image *ModelArtifactImageSource `json:"image,omitempty" protobuf:"bytes,4,opt,name=image"` +``` + +```go +// The renderer's branch (pkg/worker/controllers/worker/model_deployment_artifact.go). +case a.Delivery == workercore.ModelDeploymentModelDeliveryImage: + return []core.Volume{{ + Name: modelDeploymentModelVolumeName, + VolumeSource: core.VolumeSource{Image: &core.ImageVolumeSource{ + Reference: a.ImageReference, + }}, + }}, []core.VolumeMount{{ + Name: modelDeploymentModelVolumeName, MountPath: ModelDeploymentModelMountPath, ReadOnly: true, + }} +``` + +```go +// The capability gate reads the startup version snapshot through a swappable function, so a test +// can choose the cluster version without a cluster — the package's own seam for settings-shaped +// inputs (precedent: modelArtifactDeliveryMode). The snapshot's Configure ignores later calls, +// so pointing the tests at the snapshot itself could not test both gate directions in one +// package. SupportsFeature already answers false on a nil or unparsable version, so the +// fail-closed default falls out of the helper. +var modelArtifactClusterVersion = func() kubediscovery.Version { + return system.LoopbackKubeVersion.Get() +} +``` + +### Implementation Plan + +Tasks run sequentially in one seat. Every task starts with +`git fetch origin && git rebase origin/main`; images and e2e rebase immediately before they run; +every task lands its own `--signoff` commit (folded per module at ship; the spec file is the last +commit). Stage gates go to the coordinator: plan gate (this document), built gate, ship gate. +T1 is a spike and is run first on purpose: if the kind environment cannot serve image volumes, +the e2e environment decision escalates while the code tasks are still ahead of it, not after. + +- [ ] **T1 · e2e environment verification (spike)** + Blocked by: — + Gate: review + Owns: nothing in the tree (a local kind cluster, a probe Pod, evidence files only). + Acceptance: the candidate kind node image reports containerd ≥ 2.1 and the ImageVolume + feature enabled; a digest-pinned image made visible to the node's containerd mounts through + an image volume in a raw probe Pod (no operator code involved), under both delivery paths + that matter for the case (a registry pull and a pre-loaded IfNotPresent image); the probe's + readings are recorded. A failed check stops the task and escalates the e2e environment + decision rather than improvising. + Verify: the probe Pod reads the image's marker file; the readings name the node image's + containerd version and the feature state. +- [ ] **T2 · API + codegen** + Blocked by: — + Owns: `api/worker/v1alpha1/model_artifact.go`, `api/worker/v1alpha1/model_deployment.go`, + generated trees (deepcopy, CRD, openapi, protobuf, applyconfigurations). + Acceptance: the union field and the delivery value carry field docs that state the + identity and the absent resolved fields; `make generate` leaves a clean tree in the gen + tree. + Verify: `go build ./... && go test ./api/worker/...`. +- [ ] **T3 · Admission** + Blocked by: T2 + Owns: `pkg/kubediscovery/feature.go`, + `pkg/worker/webhooks/worker/model_artifact.go`, + `pkg/worker/webhooks/worker/model_prefetch.go` (+ tests). + Acceptance: AC1–AC4, AC14 — the union, digest-pin and patterns refusals with the message + wording above; the capability gate refusing below the floor and admitting at/above it, + both directions unit-tested through the swappable version seam; the prefetch refusal. + Verify: `go test ./pkg/kubediscovery/... ./pkg/worker/webhooks/worker/...`. +- [ ] **T4 · Resolution** + Blocked by: T2 + Owns: `pkg/worker/controllers/worker/model_artifact.go` (+ tests). + Acceptance: AC5, AC6 — first-pass resolution with a claim-shaped `resolved`; no nodes + aggregation; no revalidation. + Verify: `go test ./pkg/worker/controllers/worker/...`. +- [ ] **T5 · Renderer + placement** + Blocked by: T3, T4 + Owns: `pkg/worker/controllers/worker/model_artifact_placement.go`, + `model_deployment_artifact.go`, `model_deployment.go`, `instance.go`, + `model_placement_preference.go` (+ tests). + Acceptance: AC7–AC13 — the volume/mount shape on both paths; no cache/env/args/ephemeral + raise; `status.model` and `WeightsReady` behavior; the Setting-independence; the Instance + node pre-check and the images-list preference as adjudicated (both in scope). + Verify: `go test ./pkg/worker/controllers/worker/...`. +- [ ] **T6 · Docs** + Blocked by: T5 + Owns: the three docs paths of AC15/AC16. + Acceptance: page routing, header/Contents/footer and the index entry per the docs skill; + every number carries its version condition. + Verify: `make lint docs`. +- [ ] **T7 · e2e** + Blocked by: T1, T5 (T6 not required) + Owns: the case script + SKILL.md row + execution evidence under the task directory's + `spec-7/raw/`. + Acceptance: the cold-mount scenario below on the environment T1 verified (a failed T1 + escalates before this task starts); the capability-gate directions stay unit-tested (an + old-version e2e cluster is out of scope). + Verify: the case script against the cluster running the image built from the rebased head. +- [ ] **T8 · Ship prep** + Blocked by: T1–T7 + Owns: rebase, commit folding, chart-matrix local runs (Go change ⇒ REQUIRED: the CI chart + matrix's 7 node images, locally, all rc=0), case number finalization, PR. + Verify: `make lint`, `make lint docs`, the matrix runs, PR opened. + +### Test Plan + +#### Prerequisite testing updates + +None: the existing suites stay green throughout; the touched packages all have table-driven +suites to extend. + +#### Unit tests + +- `pkg/kubediscovery`: `FeatureImageVolume` floor (1.34 false, 1.35 true, 1.36 true, unparsable + false). +- `pkg/worker/webhooks/worker`: the union exactly-one with four members; digest-pin (tag refused, + bare refused, uppercase hex refused, correct digest admitted); patterns refused; capability + gate both directions through the swappable version seam (a 1.34 cluster refuses, a 1.35 + cluster admits, an unparsable version refuses); prefetch source refusal; the Instance webhook's + admit-every-source baseline extended with an image artifact. +- `pkg/worker/controllers/worker`: `reconcileImage` (claim-shaped resolved, no nodes, healthy); + weights resolution `Image` branch; MD render shape (volume reference, mount path, read-only, + no cache volume, no env, no args, no ephemeral raise); Instance render shape; Setting + independence; WeightsReady mounted path; the pre-check (old kubelet blocked, old containerd + blocked, unparsable blocked, current admitted) and the preference (match, cap, order, + determinism) as adjudicated. +- Every load-bearing new test is proven able to fail by a compile-safe mutation of the mechanism + it guards (mutation → red → revert → green), recorded in the handback. + +#### Integration tests + +None as a separate tier; cluster behavior is the e2e case's. + +#### e2e tests + +One case (number taken at my-ship): on the kind environment after its verification (below), +build the fixture model +image per the documented recipe (small model, one shard per layer, weights at the root, +arch-matched builder; registry reachable from the nodes — in-cluster registry + socat forwarder +from the PoC prototype, or a kind-loaded image with IfNotPresent, per the plan's verification of +the node image's containerd/gate), then: create the image-source `ModelArtifact` → reference it +from a consumer → the Pod mounts the weights and reads them (fixture's expected sha256 checked) +→ delete the consumer. Documentation covers the double-storage and GC facts; if the environment +allows (root disk pressure control on a CPU node), one GC observation (release → collection on +the next round) is captured as evidence alongside, not as a gate. + +Chart-matrix coverage follows the standing Go-change requirement at my-ship. + +## Alternatives + +- **Tag references resolved once by the operator.** Rejected for this version: it needs a + registry client and a credentials story in the worker, a re-resolution hazard (the resolved + digest would need its own immutability machinery in status), all to save the user one + `crane digest`. Reopening later is additive. +- **Image sources through the plugin cache (materialize an image into the node cache).** + Rejected: it duplicates the runtime's own content store, breaks the budget's meaning (the + plugin's bytes would double-count the image's), and puts digest semantics the plugin verifies + (file manifests) next to a source that has none — the boundary the Non-Goals fix. +- **`path`/`subPath` inside the image.** Rejected: floors containerd at 2.2 for everyone to + spare one `COPY` in a Dockerfile. +- **No placement preference (scheduler spreads, re-pulls accepted).** Rejected: a first pull was + measured at 12m43s; the soft preference's staleness costs one re-pull, the same trade the + digest preference already made. Adjudicated: the preference is in scope (2026-09-27). +- **Capability opt-in Setting for 1.33/1.34.** Rejected: the Setting would claim a guarantee the + operator cannot check (the gate state is invisible to it), and the population it serves + (admins who opened the gate) can upgrade instead. The refusal message says so. + +## Open Questions + +None open. The five draft questions were adjudicated on 2026-09-27, each as recommended: + +1. **Renderer scope — both `ModelDeployment` and `Instance`.** The resolution and render code + paths are shared, the marginal cost of the second consumer is small, and an Instance-only + image source would strand the MD path, the primary consumption surface. +2. **Placement preference from `Node.status.images` — in scope (AC13).** The pull cost justifies + the bias and the softness bounds every staleness failure at one re-pull. The kubelet's image + list is bounded (`--node-status-max-images`, default 50), so a busy node may never name the + image; the preference stays soft and no code compensates for that absence. +3. **Digest-pinning required at admission — required (AC2).** The alternatives either import a + registry client or break the identity contract. The spec commits to the honest mapping + (image digest ≠ manifest digest; no content verification) rather than an unverifiable + guarantee. +4. **Instance node capability pre-check — included (AC12).** The reconciler already reads the + Node for the hostname pin, so the check is two parsed fields and one new blocked reason, and + it converts the most confusing failure mode (mixed fleet, late kubelet error) into a named + status. +5. **e2e environment — local kind, gated on verification.** The plan's first step verifies the + candidate node image (`kindest/node` v1.35.x) ships containerd ≥ 2.1 with ImageVolume + default-on, and that an image made visible to the node container's containerd + (`kind load`-equivalent) satisfies an IfNotPresent image volume. A failed verification + escalates to the coordinator for a CPU cluster per the task's cluster protocol, matching the + PoC-K-verified versions (kubelet 1.35.7 / containerd 2.2.6); paid resources remain the user's + call.