From f760bc55fda93f0e39a243b77880e5967dc5591b Mon Sep 17 00:00:00 2001 From: Usman Shahid Date: Thu, 3 Sep 2026 10:30:14 +0400 Subject: [PATCH] fix(engine): pull an image before ContainerCreate if it isn't already local internal/engine.Docker.Start and StartJob both called ContainerCreate directly - unlike `docker run`, the Engine API never auto-pulls a missing image, so this failed outright with "No such image" on any node that hadn't already cached it. Never surfaced before because every node Sous had run on already had its images cached (from the single-node Sous era, or from this fleet's own precedent of pre-pulling vLLM images by hand via adhoc scripts). Surfaced for real tonight on aorus-ubuntu's first-ever deployment - a genuinely cold Docker install - and was worked around operationally (manual pre-pull) rather than fixed in code at the time, since this touches the same deploy path every node in the fleet uses, including the currently-serving asus-gx10. New ensureImage(ctx, ref) checks ImageInspect first (matching `docker run`'s own default of pulling only when missing, and avoiding a needless registry round-trip on the common case where a digest-pinned recipe image is already cached) and only calls ImagePull if genuinely absent. The subtle part - draining ImagePull's response stream - is split into drainPullStream and unit-tested directly against crafted byte streams, no real Docker/network access needed: a plain io.Copy(io.Discard, r) would silently report success on a pull that actually failed partway through, since a registry-side failure (bad ref, auth, missing manifest) arrives as an "error" field INSIDE the JSON progress stream, not as a Go error from ImagePull itself. Confirmed empirically (see the commit's own test) that a naive io.Copy-based drain does exactly that on a stream carrying a real embedded pull error. One additional real integration test (skips cleanly without a reachable Docker daemon or without the fixture image already cached, rather than failing an environment without Docker access) exercises ensureImage's already-cached fast path against this sandbox's actual local daemon. Co-Authored-By: Claude Sonnet 5 --- internal/engine/engine.go | 65 +++++++++++++++++++++++++++++ internal/engine/engine_test.go | 76 ++++++++++++++++++++++++++++++++++ internal/engine/job.go | 3 ++ 3 files changed, 144 insertions(+) diff --git a/internal/engine/engine.go b/internal/engine/engine.go index 1968c82..b7b61fc 100644 --- a/internal/engine/engine.go +++ b/internal/engine/engine.go @@ -2,6 +2,7 @@ package engine import ( "context" + "encoding/json" "fmt" "io" "strconv" @@ -9,6 +10,7 @@ import ( "github.com/docker/docker/api/types/container" "github.com/docker/docker/api/types/filters" + "github.com/docker/docker/api/types/image" "github.com/docker/docker/client" "github.com/docker/docker/pkg/stdcopy" "github.com/docker/go-connections/nat" @@ -98,6 +100,9 @@ func toDockerConfig(s Spec, bindHost, gpuDriver string) (*container.Config, *con } func (d *Docker) Start(ctx context.Context, s Spec) (string, error) { + if err := d.ensureImage(ctx, s.Image); err != nil { + return "", err + } cfg, host := toDockerConfig(s, d.bindHost, d.gpuDriver) created, err := d.cli.ContainerCreate(ctx, cfg, host, nil, nil, s.Name) if err != nil { @@ -109,6 +114,66 @@ func (d *Docker) Start(ctx context.Context, s Spec) (string, error) { return created.ID, nil } +// ensureImage pulls ref if it is not already present locally. The Engine +// API's ContainerCreate, unlike `docker run`, never pulls a missing image +// itself - it fails outright with "No such image". This went unnoticed +// through every deploy this project has ever made, because every node +// Sous had run on already had its images cached (from the single-node +// Sous era, or from this fleet's own precedent of pre-pulling vLLM images +// by hand). It surfaced for real on aorus-ubuntu's first-ever deployment: +// a genuinely cold Docker install, with nothing cached at all. +// +// Checks local presence FIRST rather than pulling unconditionally on +// every call - matching `docker run`'s own default ("pull if missing", +// not "pull always") and avoiding a needless registry round-trip on the +// overwhelmingly common case where the image core recipes (which pin +// digests, not floating tags) already have cached. +func (d *Docker) ensureImage(ctx context.Context, ref string) error { + if _, err := d.cli.ImageInspect(ctx, ref); err == nil { + return nil + } else if !client.IsErrNotFound(err) { + return fmt.Errorf("engine: inspect image %s: %w", ref, err) + } + + rc, err := d.cli.ImagePull(ctx, ref, image.PullOptions{}) + if err != nil { + return fmt.Errorf("engine: pull image %s: %w", ref, err) + } + defer rc.Close() + + if err := drainPullStream(rc); err != nil { + return fmt.Errorf("engine: pull image %s: %w", ref, err) + } + return nil +} + +// drainPullStream reads an ImagePull response to completion, which is +// required for the pull to actually finish (it is not synchronous until +// the stream is read) - and, unlike a plain io.Copy(io.Discard, r), also +// catches a registry-side failure (auth, missing manifest, ...), which +// Docker reports as an "error" field INSIDE this JSON stream rather than +// as a Go error from ImagePull itself. Split out from ensureImage so the +// decode/error-detection logic - the actual subtle part - is testable +// against a crafted byte stream, with no real Docker daemon or network +// access required. +func drainPullStream(r io.Reader) error { + dec := json.NewDecoder(r) + for { + var msg struct { + Error string `json:"error"` + } + if err := dec.Decode(&msg); err != nil { + if err == io.EOF { + return nil + } + return fmt.Errorf("reading progress: %w", err) + } + if msg.Error != "" { + return fmt.Errorf("%s", msg.Error) + } + } +} + // Stop stops and removes. Leaving a stopped container behind would make the // next create fail on the name, which is the kind of failure that gets // misread as a problem with the new model. diff --git a/internal/engine/engine_test.go b/internal/engine/engine_test.go index f376fb4..838e79f 100644 --- a/internal/engine/engine_test.go +++ b/internal/engine/engine_test.go @@ -1,6 +1,8 @@ package engine import ( + "context" + "strings" "testing" "github.com/docker/go-connections/nat" @@ -116,6 +118,80 @@ func TestNoEntrypointLeavesImageDefault(t *testing.T) { } } +// A real docker pull's progress stream, one JSON object per line, no +// embedded error - captured in shape from an actual `docker pull` (status, +// progressDetail, id fields), not invented. +const realPullStreamNoError = `{"status":"Pulling from library/alpine","id":"3.21"} +{"status":"Pulling fs layer","progressDetail":{},"id":"9b18e9b68314"} +{"status":"Downloading","progressDetail":{"current":1024,"total":3072},"progress":"[====> ] 1024B/3072B","id":"9b18e9b68314"} +{"status":"Download complete","progressDetail":{},"id":"9b18e9b68314"} +{"status":"Pull complete","progressDetail":{},"id":"9b18e9b68314"} +{"status":"Digest: sha256:abc123"} +{"status":"Status: Downloaded newer image for alpine:3.21"} +` + +func TestDrainPullStreamSucceedsOnARealNoErrorStream(t *testing.T) { + if err := drainPullStream(strings.NewReader(realPullStreamNoError)); err != nil { + t.Fatalf("unexpected error: %v", err) + } +} + +func TestDrainPullStreamCatchesAnErrorEmbeddedMidStream(t *testing.T) { + // The exact bug ensureImage exists to avoid: a registry-side failure + // (bad ref, auth, ...) arrives as an "error" field INSIDE the JSON + // stream, after several genuine progress lines - not as a Go error + // from ImagePull itself. A plain io.Copy(io.Discard, r) would drain + // this to EOF and report success. + stream := `{"status":"Pulling from library/alpine","id":"3.21"} +{"status":"Pulling fs layer","progressDetail":{},"id":"9b18e9b68314"} +{"errorDetail":{"message":"manifest unknown: manifest unknown"},"error":"manifest unknown: manifest unknown"} +` + err := drainPullStream(strings.NewReader(stream)) + if err == nil { + t.Fatal("expected an error, got nil") + } + if !strings.Contains(err.Error(), "manifest unknown") { + t.Fatalf("error should surface the registry's own message, got: %v", err) + } +} + +func TestDrainPullStreamHandlesAnEmptyStream(t *testing.T) { + if err := drainPullStream(strings.NewReader("")); err != nil { + t.Fatalf("unexpected error on an empty stream: %v", err) + } +} + +func TestDrainPullStreamSurfacesMalformedJSON(t *testing.T) { + err := drainPullStream(strings.NewReader("{not json")) + if err == nil { + t.Fatal("expected an error for malformed JSON, got nil") + } +} + +// TestEnsureImageSkipsAnAlreadyCachedImage is a real integration test +// against an actual Docker daemon, not a fake - ensureImage's whole +// premise (check ImageInspect before ever calling ImagePull) isn't +// meaningfully testable through drainPullStream alone, since that +// covers only what happens once a pull stream exists. Skips cleanly if +// no daemon is reachable or the fixture image isn't cached, rather than +// failing a CI environment without Docker access - matching this +// project's existing pattern of degrading gracefully for environment +// limits (e.g. -race being unavailable in sandboxes without a C +// toolchain) rather than papering over the gap with a fake. +func TestEnsureImageSkipsAnAlreadyCachedImage(t *testing.T) { + d, err := New("", "cdi") + if err != nil { + t.Skipf("no local Docker daemon reachable: %v", err) + } + const fixtureImage = "alpine:3.21" + if _, err := d.cli.ImageInspect(context.Background(), fixtureImage); err != nil { + t.Skipf("fixture image %s not already cached locally: %v", fixtureImage, err) + } + if err := d.ensureImage(context.Background(), fixtureImage); err != nil { + t.Fatalf("ensureImage on an already-cached image should not error: %v", err) + } +} + func TestBindsAndRestartPolicy(t *testing.T) { s := Spec{Name: "n", Image: "i", ContainerPort: 8000, Binds: []string{"/models:/root/.cache/huggingface"}} diff --git a/internal/engine/job.go b/internal/engine/job.go index 22adf12..620f9e8 100644 --- a/internal/engine/job.go +++ b/internal/engine/job.go @@ -53,6 +53,9 @@ type JobSpec struct { // minutes; holding a request open for that is the same mistake undeploy used to // make. Progress is read afterwards from the container's own state and logs. func (d *Docker) StartJob(ctx context.Context, s JobSpec) (string, error) { + if err := d.ensureImage(ctx, s.Image); err != nil { + return "", err + } cfg := &container.Config{ Image: s.Image, Cmd: s.Cmd,