package main // spawner_test.go — clone→build→create/start pipeline, git failures, env wiring. import ( "archive/tar" "bytes" "context" "crypto/sha256" "encoding/hex" "errors" "fmt" "io" "net/http" "net/http/httptest" "os" "path/filepath" "reflect" "regexp" "strings" "testing" "time" "github.com/docker/docker/client" ) func TestSpawnerStartHappyPath(t *testing.T) { gitLog := useFakeGit(t, fakeGitModeOK) f := newFakeDocker() sp, store := newTestSpawner(t, f) res, err := sp.Start(context.Background(), "group/project", "main", "", false) if err != nil { t.Fatalf("Start: %v", err) } if res.SessionID == "" || res.ImageUsed != imageRefWorker { t.Fatalf("Start result = %+v, want sessionId set and default image", res) } job := waitJobState(t, sp, res.SessionID, stateRunning) if job.Repo != "group/project" || job.ContainerID == "" { t.Fatalf("running job = %+v, want containerId set", job) } // clone URL carries the PAT only when none is stored (it isn't here) calls := readGitLog(t, gitLog) if len(calls) != 1 || !strings.HasPrefix(calls[0], "clone --branch main -- https://gitlab.example/group/project.git ") { t.Fatalf("git calls = %v", calls) } // container created with the right image, env, labels, mounts, network creates := f.createsByName("lvmh-agent-") if len(creates) != 1 { t.Fatalf("agent containers created = %+v", creates) } c := creates[0] if c.Image != imageRefWorker { t.Fatalf("image = %q", c.Image) } if c.Labels[labelSession] != res.SessionID { t.Fatalf("labels = %v", c.Labels) } wantEnv := map[string]bool{ envProviderAPIKey + "=key-123": true, envToken + "=" + testToken: true, "LVMH_URL=" + defaultContainerURL: true, envLVMHSessionID + "=" + res.SessionID: true, envLVMHAgent + "=1": true, envLVMHRepo + "=group/project": true, } // provider key passthroughs are environment-dependent; ignore them here passthrough := map[string]bool{ "OPENAI_API_KEY": true, "GEMINI_API_KEY": true, "DEEPSEEK_KEY": true, "ANTHROPIC_API_KEY": true, "LVMH_GITEA_TOKEN": true, "PLAYWRIGHT_BROWSERS_PATH": true, } for _, e := range c.Env { if name, _, ok := strings.Cut(e, "="); ok && passthrough[name] { continue } if !wantEnv[e] { t.Fatalf("unexpected env %q", e) } delete(wantEnv, e) } if len(wantEnv) != 0 { t.Fatalf("missing env %v", wantEnv) } wantBinds := []string{ volumeRepoPrefix + repoSlug("group/project") + ":" + workspaceMount, volumeSessions + ":" + sessionsMount, volumePiCache + ":" + cacheMount, } if len(c.HostConfig.Binds) != 3 { t.Fatalf("binds = %v, want %v", c.HostConfig.Binds, wantBinds) } for i, b := range wantBinds { if c.HostConfig.Binds[i] != b { t.Fatalf("binds = %v, want %v", c.HostConfig.Binds, wantBinds) } } if c.HostConfig.NetworkMode != defaultNetwork { t.Fatalf("network = %q", c.HostConfig.NetworkMode) } if c.HostConfig.Init == nil || !*c.HostConfig.Init { t.Fatalf("agent container Init = %v, want true (tini zombie reaping)", c.HostConfig.Init) } // fresh repo volume seeded from the clone via CopyToContainer (tar) seed := f.createsByName("lvmh-seed-") if len(seed) != 1 { t.Fatalf("seed containers = %+v", seed) } if seed[0].Image != imageRefWorker { t.Fatalf("seed create = %+v", seed[0]) } wantSeedBinds := []string{volumeRepoPrefix + repoSlug("group/project") + ":" + workspaceMount} if !reflect.DeepEqual(seed[0].HostConfig.Binds, wantSeedBinds) { t.Fatalf("seed binds = %v, want %v (no host-path binds)", seed[0].HostConfig.Binds, wantSeedBinds) } if seed[0].HostConfig.Init == nil || !*seed[0].HostConfig.Init { t.Fatalf("seed container Init = %v, want true", seed[0].HostConfig.Init) } if len(f.archives) != 1 || f.archives[0] == 0 { t.Fatalf("CopyToContainer archives = %v, want one non-empty tar", f.archives) } // no build expected (fake daemon already has the image) if f.hasCall(http.MethodPost, "/build") { t.Fatal("image present but build was called") } // session volume created too if !f.volumeExists(volumeSessions) { t.Fatalf("volume %q missing", volumeSessions) } // container row persisted row, ok, err := store.GetContainer(res.SessionID) if err != nil || !ok || row.Repo != "group/project" { t.Fatalf("container row = %+v ok=%v err=%v", row, ok, err) } f.assertNoUnknown(t) } func TestSpawnerBuildsImageWhenMissing(t *testing.T) { useFakeGit(t, fakeGitModeOK) f := newFakeDocker() f.images = 0 // image absent → ensureImage must build sp, _ := newTestSpawner(t, f) res, err := sp.Start(context.Background(), "group/project", "", "", false) if err != nil { t.Fatalf("Start: %v", err) } waitJobState(t, sp, res.SessionID, stateRunning) builds := f.countCalls(http.MethodPost, "/build") if builds != 1 { t.Fatalf("build calls = %d, want 1", builds) } f.assertNoUnknown(t) } func TestSpawnerStartValidatesDockerAndDockerfile(t *testing.T) { f := newFakeDocker() f.images = 0 sp, _ := newTestSpawner(t, f) // dockerfile exists in the fake context // docker reachable but image absent and no Dockerfile anywhere → clear error t.Setenv(envWorkerDockerfile, filepath.Join(t.TempDir(), "missing.Dockerfile")) sp2, err := NewSpawner(context.Background(), sp.store, sp.hub, "https://gitlab.example/") if err != nil { t.Fatalf("NewSpawner: %v", err) } _, err = sp2.Start(context.Background(), "group/project", "", "", false) if err == nil || !strings.Contains(err.Error(), "no worker Dockerfile") { t.Fatalf("Start without dockerfile err = %v", err) } if len(sp2.JobsSnapshot()) != 0 { t.Fatalf("failed Start must not leave a job: %+v", sp2.JobsSnapshot()) } } func TestSpawnerRunJobErrorStates(t *testing.T) { cases := []struct { name string setup func(t *testing.T, f *fakeDocker) gitMode string wantMsg string }{ { name: "clone-fails", gitMode: fakeGitModeFail, wantMsg: "fatal: repository not found", }, { name: "noisy-clone-failure-truncated-to-tail", gitMode: fakeGitModeNoisy, wantMsg: "xxx", }, { name: "build-fails", setup: func(t *testing.T, f *fakeDocker) { f.images = 0 f.failBuild = true }, wantMsg: "build exploded", }, { name: "create-fails", setup: func(t *testing.T, f *fakeDocker) { f.failCreate = true }, wantMsg: "create failed", }, { name: "start-fails-and-cleans-up", setup: func(t *testing.T, f *fakeDocker) { f.failStart = true }, wantMsg: "start failed", }, { name: "volume-create-fails", setup: func(t *testing.T, f *fakeDocker) { f.failVolumeCreate = true }, wantMsg: "volume create failed", }, { name: "seed-copy-fails", setup: func(t *testing.T, f *fakeDocker) { f.failArchive = true }, wantMsg: "copy into volume", }, } for _, tc := range cases { t.Run(tc.name, func(t *testing.T) { gitMode := tc.gitMode if gitMode == "" { gitMode = fakeGitModeOK } useFakeGit(t, gitMode) f := newFakeDocker() if tc.setup != nil { tc.setup(t, f) } sp, _ := newTestSpawner(t, f) res, err := sp.Start(context.Background(), "group/project", "", "", false) if err != nil { t.Fatalf("Start: %v", err) } job := waitJobState(t, sp, res.SessionID, stateError) if !strings.Contains(job.Message, tc.wantMsg) { t.Fatalf("job message = %q, want containing %q", job.Message, tc.wantMsg) } if job.ContainerID != "" { t.Fatalf("error job has containerId %q", job.ContainerID) } }) } } func TestSpawnerCloneOrUpdatePullsExisting(t *testing.T) { gitLog := useFakeGit(t, fakeGitModeOK) f := newFakeDocker() sp, _ := newTestSpawner(t, f) // existing clone → pull --ff-only in that dir, no clone dir := filepath.Join(sp.reposDir, repoSlug("group/project")) if err := os.MkdirAll(filepath.Join(dir, ".git"), 0o755); err != nil { t.Fatalf("mkdir .git: %v", err) } if err := sp.cloneOrUpdate(context.Background(), "group/project", "", repoSlug("group/project")); err != nil { t.Fatalf("cloneOrUpdate: %v", err) } calls := readGitLog(t, gitLog) if len(calls) != 1 || calls[0] != "pull --ff-only" { t.Fatalf("git calls = %v, want pull --ff-only", calls) } } func TestSpawnerCloneURLNeverCarriesPAT(t *testing.T) { f := newFakeDocker() sp, store := newTestSpawner(t, f) if u, err := sp.cloneURL("group/project"); err != nil || u != "https://gitlab.example/group/project.git" { t.Fatalf("cloneURL without PAT = %q err=%v", u, err) } if err := store.SetSetting(settingGitLabToken, "pat-1"); err != nil { t.Fatalf("set token: %v", err) } // A stored PAT must never leak into the persisted clone URL. u, err := sp.cloneURL("group/project") if err != nil { t.Fatalf("cloneURL: %v", err) } if u != "https://gitlab.example/group/project.git" || strings.Contains(u, "pat-1") { t.Fatalf("cloneURL with stored PAT = %q, want clean URL", u) } } func TestSpawnerRemoveSession(t *testing.T) { useFakeGit(t, fakeGitModeOK) f := newFakeDocker() sp, store := newTestSpawner(t, f) if err := sp.RemoveSession(context.Background(), "nope"); err != errNoContainer { t.Fatalf("remove unknown = %v, want errNoContainer", err) } if err := store.UpsertContainer("s1", "cid-9", "group/project"); err != nil { t.Fatalf("upsert container: %v", err) } sp.setJob("s1", "group/project", stateRunning, "cid-9", "") if err := sp.RemoveSession(context.Background(), "s1"); err != nil { t.Fatalf("RemoveSession: %v", err) } if !f.hasCall(http.MethodPost, "/containers/cid-9/stop") || !f.hasCall(http.MethodDelete, "/containers/cid-9") { t.Fatal("container not stopped+removed") } if _, ok, _ := store.GetContainer("s1"); ok { t.Fatal("container row must be deleted") } for _, j := range sp.JobsSnapshot() { if j.SessionID == "s1" { t.Fatalf("job after removal = %+v, want entry deleted", j) } } // stop failure surfaces, row kept if err := store.UpsertContainer("s2", "cid-10", "group/project"); err != nil { t.Fatalf("upsert: %v", err) } f.failStop = true if err := sp.RemoveSession(context.Background(), "s2"); err == nil || !strings.Contains(err.Error(), "docker stop") { t.Fatalf("remove with stop failure = %v", err) } f.assertNoUnknown(t) } func TestSpawnerSeedVolumeCancel(t *testing.T) { useFakeGit(t, fakeGitModeOK) f := newFakeDocker() f.archiveHang = true // server never answers the archive PUT sp, _ := newTestSpawner(t, f) if err := os.MkdirAll(filepath.Join(sp.reposDir, repoSlug("group/project")), 0o755); err != nil { t.Fatal(err) } if err := os.WriteFile(filepath.Join(sp.reposDir, repoSlug("group/project"), "README.md"), []byte("x"), 0o644); err != nil { t.Fatal(err) } ctx, cancel := context.WithCancel(context.Background()) defer cancel() done := make(chan error, 1) go func() { done <- sp.seedVolume(ctx, repoSlug("group/project"), volumeRepoPrefix+repoSlug("group/project"), "s1") }() waitFor(t, 5*time.Second, func() bool { return f.hasCallSuffix(http.MethodPut, "/archive") }) cancel() select { case err := <-done: if !errors.Is(err, context.Canceled) { t.Fatalf("seedVolume after cancel = %v, want context.Canceled", err) } case <-time.After(5 * time.Second): t.Fatal("seedVolume did not return after context cancel") } } func TestSpawnerSameRepoSpawnsSerialize(t *testing.T) { useFakeGit(t, fakeGitModeOK) f := newFakeDocker() sp, _ := newTestSpawner(t, f) l1 := sp.slugLock("a") l2 := sp.slugLock("a") if l1 != l2 { t.Fatal("slugLock must return the same mutex per slug") } if sp.slugLock("b") == l1 { t.Fatal("slugLock must return distinct mutexes per slug") } res1, err := sp.Start(context.Background(), "group/project", "", "", false) if err != nil { t.Fatalf("Start 1: %v", err) } res2, err := sp.Start(context.Background(), "group/project", "", "", false) if err != nil { t.Fatalf("Start 2: %v", err) } waitJobState(t, sp, res1.SessionID, stateRunning) waitJobState(t, sp, res2.SessionID, stateRunning) if len(f.createsByName("lvmh-agent-")) != 2 { t.Fatalf("agent containers = %d, want 2", len(f.createsByName("lvmh-agent-"))) } } func TestRepoSlugAndUUID(t *testing.T) { slugRe := regexp.MustCompile(`^a--b--c-[0-9a-f]{6}$`) if got := repoSlug("a/b/c"); !slugRe.MatchString(got) { t.Fatalf("repoSlug = %q, want a--b--c-<6 hex>", got) } // sanitized path + first 6 hex of sha256(full repo path) sum := sha256.Sum256([]byte("a/b/c")) want := "a--b--c-" + hex.EncodeToString(sum[:3]) if got := repoSlug("a/b/c"); got != want { t.Fatalf("repoSlug = %q, want %q", got, want) } // "/"→"--" alone is ambiguous: distinct repo paths must never collide. seen := map[string]bool{} for _, repo := range []string{"a/b/c", "a/b-c", "a/b--c", "a/b--c/d"} { s := repoSlug(repo) if seen[s] { t.Fatalf("slug collision: %q", s) } seen[s] = true } id := newUUID() if len(id) != 36 || id[8] != '-' || id[13] != '-' || id[18] != '-' || id[23] != '-' { t.Fatalf("newUUID shape = %q", id) } if id2 := newUUID(); id2 == id { t.Fatal("newUUID must not repeat") } } func TestSpawnerCloneUsesHeaderAuthNotURLCredentials(t *testing.T) { gitLog := useFakeGit(t, fakeGitModeOK) f := newFakeDocker() sp, store := newTestSpawner(t, f) if err := store.SetSetting(settingGitLabToken, "pat-1"); err != nil { t.Fatalf("set token: %v", err) } res, err := sp.Start(context.Background(), "group/project", "", "", false) if err != nil { t.Fatalf("Start: %v", err) } waitJobState(t, sp, res.SessionID, stateRunning) calls := readGitLog(t, gitLog) if len(calls) != 1 { t.Fatalf("git calls = %v", calls) } call := calls[0] // auth travels as a per-invocation -c http.extraHeader arg… if !strings.HasPrefix(call, "-c http.extraHeader=Authorization: token pat-1 clone ") { t.Fatalf("git call = %q, want -c http.extraHeader auth before clone", call) } // …and the URL recorded into .git/config stays credential-free. if !strings.Contains(call, " -- https://gitlab.example/group/project.git ") { t.Fatalf("git call = %q, want clean clone URL", call) } if strings.Contains(call, "pat-1@") { t.Fatalf("git call = %q leaks the PAT into the URL", call) } } func TestSpawnerJobsPrunedToCap(t *testing.T) { f := newFakeDocker() sp, _ := newTestSpawner(t, f) for i := 0; i < maxSpawnJobs+10; i++ { sp.setJob(fmt.Sprintf("s%d", i), "group/project", stateCloning, "", "") } jobs := sp.JobsSnapshot() if len(jobs) != maxSpawnJobs { t.Fatalf("jobs = %d, want capped at %d", len(jobs), maxSpawnJobs) } seen := map[string]bool{} for _, j := range jobs { seen[j.SessionID] = true } if seen["s0"] { t.Fatal("oldest job must be pruned first") } for _, id := range []string{fmt.Sprintf("s%d", maxSpawnJobs), fmt.Sprintf("s%d", maxSpawnJobs+9)} { if !seen[id] { t.Fatalf("newest job %s pruned; kept = %v", id, seen) } } } func TestSpawnerFailedSeedRemovesRepoVolume(t *testing.T) { useFakeGit(t, fakeGitModeOK) f := newFakeDocker() f.failArchive = true // CopyToContainer fails → seed fails sp, _ := newTestSpawner(t, f) res, err := sp.Start(context.Background(), "group/project", "", "", false) if err != nil { t.Fatalf("Start: %v", err) } job := waitJobState(t, sp, res.SessionID, stateError) if !strings.Contains(job.Message, "seed") { t.Fatalf("job message = %q, want seed failure", job.Message) } if !f.hasCall(http.MethodDelete, "/volumes/"+volumeRepoPrefix+repoSlug("group/project")) { t.Fatal("failed seed must force-remove the repo volume for a fresh retry") } if f.volumeExists(volumeRepoPrefix + repoSlug("group/project")) { t.Fatal("repo volume must not linger half-seeded") } } func TestSpawnerFailedSeedSurvivesVolumeRemoveFailure(t *testing.T) { useFakeGit(t, fakeGitModeOK) f := newFakeDocker() f.failArchive = true f.failVolumeDelete = true // cleanup itself fails; seed error still surfaces sp, _ := newTestSpawner(t, f) res, err := sp.Start(context.Background(), "group/project", "", "", false) if err != nil { t.Fatalf("Start: %v", err) } job := waitJobState(t, sp, res.SessionID, stateError) if !strings.Contains(job.Message, "seed") { t.Fatalf("job message = %q, want seed failure despite cleanup failure", job.Message) } } func TestSpawnerCloneOrUpdateBranchFetchFails(t *testing.T) { useFakeGit(t, fakeGitModeFail) // branch probe and fetch both fail f := newFakeDocker() sp, _ := newTestSpawner(t, f) slug := repoSlug("group/project") if err := os.MkdirAll(filepath.Join(sp.reposDir, slug, ".git"), 0o755); err != nil { t.Fatalf("mkdir .git: %v", err) } err := sp.cloneOrUpdate(context.Background(), "group/project", "dev", slug) if err == nil || !strings.Contains(err.Error(), "git fetch") { t.Fatalf("cloneOrUpdate branch fetch failure = %v, want git fetch error", err) } } func TestSpawnerDeleteJobMissing(t *testing.T) { sp, _ := newTestSpawner(t, newFakeDocker()) sp.deleteJob("never-existed") // no-op, must not panic sp.setJob("s1", "group/project", stateRunning, "cid", "") sp.deleteJob("s1") for _, j := range sp.JobsSnapshot() { if j.SessionID == "s1" { t.Fatalf("job %s still present after deleteJob", j.SessionID) } } } func TestSpawnerCloneOrUpdateSwitchesBranch(t *testing.T) { cases := []struct { name string branch string headOut string wantTail []string }{ { name: "same-branch-skips-fetch-checkout", branch: "main", headOut: "main\n", wantTail: []string{"pull --ff-only"}, }, { name: "different-branch-fetches-and-checks-out", branch: "dev", headOut: "main\n", wantTail: []string{"fetch origin dev", "checkout dev", "pull --ff-only"}, }, { name: "unknown-head-falls-through-to-fetch", branch: "dev", headOut: "", wantTail: []string{"fetch origin dev", "checkout dev", "pull --ff-only"}, }, } for _, tc := range cases { t.Run(tc.name, func(t *testing.T) { gitLog := useFakeGit(t, fakeGitModeOK) t.Setenv("FAKE_GIT_HEAD", tc.headOut) f := newFakeDocker() sp, store := newTestSpawner(t, f) if err := store.SetSetting(settingGitLabToken, "pat-1"); err != nil { t.Fatalf("set token: %v", err) } slug := repoSlug("group/project") dir := filepath.Join(sp.reposDir, slug) if err := os.MkdirAll(filepath.Join(dir, ".git"), 0o755); err != nil { t.Fatalf("mkdir .git: %v", err) } if err := sp.cloneOrUpdate(context.Background(), "group/project", tc.branch, slug); err != nil { t.Fatalf("cloneOrUpdate: %v", err) } calls := readGitLog(t, gitLog) if len(calls) != len(tc.wantTail)+1 { // +1: rev-parse probe t.Fatalf("git calls = %v, want %v (+rev-parse)", calls, tc.wantTail) } if !strings.HasPrefix(calls[0], "-C ") || !strings.HasSuffix(calls[0], "rev-parse --abbrev-ref HEAD") { t.Fatalf("first call = %q, want branch probe", calls[0]) } authPrefix := "-c http.extraHeader=Authorization: token pat-1 " for i, want := range tc.wantTail { got := calls[i+1] if want == "checkout dev" { if got != want { // checkout needs no auth t.Fatalf("call %d = %q, want %q", i+1, got, want) } continue } if got != authPrefix+want { t.Fatalf("call %d = %q, want %q%q", i+1, got, authPrefix, want) } } }) } } func TestEnvOr(t *testing.T) { t.Setenv("LVMH_TEST_ENV_OR", " value ") if got := envOr("LVMH_TEST_ENV_OR", "def"); got != "value" { t.Fatalf("envOr trimmed = %q", got) } t.Setenv("LVMH_TEST_ENV_OR", " ") if got := envOr("LVMH_TEST_ENV_OR", "def"); got != "def" { t.Fatalf("envOr whitespace-only = %q", got) } if got := envOr("LVMH_TEST_ENV_OR_UNSET", "def"); got != "def" { t.Fatalf("envOr unset = %q", got) } } func TestExtractBuildError(t *testing.T) { cases := []struct { name string body string want string }{ {"empty", "", ""}, {"no-error", `{"stream":"Step 1/3"}` + "\n" + `{"stream":"done"}`, ""}, {"error-field", `{"error":"nope"}`, "nope"}, {"error-detail-wins", `{"errorDetail":{"message":"boom"},"error":"generic"}`, "boom"}, {"mixed-stream", "{\"stream\":\"...\"}\n{\"error\":\"late failure\"}", "late failure"}, {"non-json-line", "garbage line\n{\"error\":\"after garbage\"}", "after garbage"}, } for _, tc := range cases { t.Run(tc.name, func(t *testing.T) { if got := extractBuildError([]byte(tc.body)); got != tc.want { t.Fatalf("extractBuildError(%q) = %q, want %q", tc.body, got, tc.want) } }) } } func TestTarDir(t *testing.T) { dir := t.TempDir() if err := os.WriteFile(filepath.Join(dir, "one.txt"), []byte("one"), 0o644); err != nil { t.Fatalf("write: %v", err) } if err := os.MkdirAll(filepath.Join(dir, "sub"), 0o755); err != nil { t.Fatalf("mkdir: %v", err) } if err := os.WriteFile(filepath.Join(dir, "sub", "two.txt"), []byte("two"), 0o644); err != nil { t.Fatalf("write: %v", err) } var buf bytes.Buffer if err := tarDir(&buf, dir); err != nil { t.Fatalf("tarDir: %v", err) } tr := tar.NewReader(bytes.NewReader(buf.Bytes())) names := map[string]string{} for { hdr, err := tr.Next() if err == io.EOF { break } if err != nil { t.Fatalf("tar next: %v", err) } body, _ := io.ReadAll(tr) names[hdr.Name] = string(body) if hdr.Uname != "root" || hdr.Gname != "root" { t.Fatalf("header %q owner = %s/%s, want root/root", hdr.Name, hdr.Uname, hdr.Gname) } } for name, want := range map[string]string{"one.txt": "one", "sub/two.txt": "two"} { if names[name] != want { t.Fatalf("archive missing %q (have %v)", name, names) } } if _, ok := names["sub/"]; !ok { t.Fatalf("directories must be archived with trailing slash: %v", names) } if err := tarDir(&buf, filepath.Join(dir, "does-not-exist")); err == nil { t.Fatal("tarDir of missing dir must fail") } } func TestGitRunSilentFailure(t *testing.T) { useFakeGit(t, fakeGitModeSilent) err := gitRun(context.Background(), "", "clone", "x") if err == nil || strings.Contains(err.Error(), "clone x") { // silent failure surfaces the bare exec error, not a padded message t.Fatalf("gitRun silent failure = %v", err) } } func TestSpawnerNewBadDockerHost(t *testing.T) { t.Setenv("DOCKER_HOST", "http://") store := openTestStore(t) if _, err := NewSpawner(context.Background(), store, NewHub(store), "https://gitlab.example"); err == nil { t.Fatal("NewSpawner must fail on an unparseable DOCKER_HOST") } } func TestSpawnerStartDockerUnavailable(t *testing.T) { f := newFakeDocker() sp, _ := newTestSpawner(t, f) // point the spawner at a dead endpoint: last Start request kills nothing // because the client is already built; close the fake server instead. sp.cli.Close() closeDocker(t, sp) if _, err := sp.Start(context.Background(), "group/project", "", "", false); err == nil || !strings.Contains(err.Error(), "docker unavailable") { t.Fatalf("Start with dead docker = %v, want docker unavailable", err) } } // closeDocker re-aims the spawner at a closed server (dial refused). func closeDocker(t *testing.T, sp *Spawner) { t.Helper() dead := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { writeJSONNow(w, http.StatusInternalServerError, `{"message":"dead"}`) })) dead.Close() // listener gone: connections refused if err := sp.cli.Close(); err != nil { t.Fatalf("close docker client: %v", err) } cli, err := client.NewClientWithOpts(client.WithHost("tcp://" + dead.Listener.Addr().String())) if err != nil { t.Fatalf("rebuild docker client: %v", err) } sp.cli = cli } func TestSpawnerCloneOrUpdateErrors(t *testing.T) { f := newFakeDocker() sp, _ := newTestSpawner(t, f) // unparseable repo path (control char) → cloneURL parse error if err := sp.cloneOrUpdate(context.Background(), "bad\nrepo", "", repoSlug("bad\nrepo")); err == nil { t.Fatal("cloneOrUpdate with control-char repo must fail") } // reposDir path occupied by a file → MkdirAll fails file := filepath.Join(t.TempDir(), "not-a-dir") if err := os.WriteFile(file, nil, 0o644); err != nil { t.Fatalf("write file: %v", err) } sp.reposDir = file if err := sp.cloneOrUpdate(context.Background(), "group/project", "", repoSlug("group/project")); err == nil { t.Fatal("cloneOrUpdate with file reposDir must fail") } } func TestSpawnerEnsureImageErrors(t *testing.T) { f := newFakeDocker() f.images = 0 sp, _ := newTestSpawner(t, f) ctx := context.Background() // dockerfile missing on disk sp.dockerfile = filepath.Join(t.TempDir(), "gone.Dockerfile") if err := sp.ensureImage(ctx); err == nil || !strings.Contains(err.Error(), "worker Dockerfile missing") { t.Fatalf("ensureImage without dockerfile = %v", err) } // dockerfile exists but cannot be made relative to the build context sp.dockerfile = filepath.Join(t.TempDir(), "worker.Dockerfile") if err := os.WriteFile(sp.dockerfile, []byte("FROM scratch"), 0o644); err != nil { t.Fatalf("write: %v", err) } sp.buildContext = "relative-context" if err := sp.ensureImage(ctx); err == nil { t.Fatal("ensureImage with relative-context vs absolute dockerfile must fail") } // tar of a missing build context fails sp.buildContext = filepath.Join(t.TempDir(), "no-such-context") if err := sp.ensureImage(ctx); err == nil || !strings.Contains(err.Error(), "build context") { t.Fatalf("ensureImage with missing context = %v", err) } } func TestSpawnerSessionsVolumeCreateFails(t *testing.T) { useFakeGit(t, fakeGitModeOK) f := newFakeDocker() f.volume[volumeRepoPrefix+repoSlug("group/project")] = true // repo volume exists → skip seed f.failVolumeCreate = true sp, _ := newTestSpawner(t, f) res, err := sp.Start(context.Background(), "group/project", "", "", false) if err != nil { t.Fatalf("Start: %v", err) } job := waitJobState(t, sp, res.SessionID, stateError) if !strings.Contains(job.Message, volumeSessions) { t.Fatalf("job message = %q, want sessions volume failure", job.Message) } } func TestSpawnerSeedVolumeCreateStartFail(t *testing.T) { useFakeGit(t, fakeGitModeOK) f := newFakeDocker() sp, _ := newTestSpawner(t, f) ctx := context.Background() f.failCreate = true if err := sp.seedVolume(ctx, repoSlug("group/project"), volumeRepoPrefix+repoSlug("group/project"), "s1"); err == nil { t.Fatal("seedVolume with failing create must fail") } f.failCreate = false f.failStart = true if err := sp.seedVolume(ctx, repoSlug("group/project"), volumeRepoPrefix+repoSlug("group/project"), "s1"); err == nil { t.Fatal("seedVolume with failing start must fail") } } func TestSpawnerStoreFailures(t *testing.T) { useFakeGit(t, fakeGitModeOK) f := newFakeDocker() sp, _ := newTestSpawner(t, f) // closed store: GetContainer errors deadStore := openTestStore(t) _ = deadStore.Close() sp.store = deadStore if err := sp.RemoveSession(context.Background(), "s1"); err == nil || !strings.Contains(err.Error(), "closed") { t.Fatalf("RemoveSession with closed store = %v", err) } // runJob completes container start but cannot persist the row sp2, _ := newTestSpawner(t, newFakeDocker()) dead2 := openTestStore(t) _ = dead2.Close() sp2.store = dead2 res, err := sp2.Start(context.Background(), "group/project", "", "", false) if err != nil { t.Fatalf("Start: %v", err) } job := waitJobState(t, sp2, res.SessionID, stateError) if !strings.Contains(job.Message, "not persisted") { t.Fatalf("job message = %q, want persistence failure", job.Message) } } func TestSpawnerCloneOrUpdatePullFails(t *testing.T) { useFakeGit(t, fakeGitModeFail) f := newFakeDocker() sp, _ := newTestSpawner(t, f) dir := filepath.Join(sp.reposDir, repoSlug("group/project")) if err := os.MkdirAll(filepath.Join(dir, ".git"), 0o755); err != nil { t.Fatalf("mkdir .git: %v", err) } err := sp.cloneOrUpdate(context.Background(), "group/project", "", repoSlug("group/project")) if err == nil || !strings.Contains(err.Error(), "git pull") { t.Fatalf("cloneOrUpdate pull failure = %v, want git pull error", err) } } func TestSpawnerWorkerCreateStartFailures(t *testing.T) { // repo volume pre-exists → seeding skipped → failures surface at the // worker container create/start instead of at the seed container. useFakeGit(t, fakeGitModeOK) f := newFakeDocker() f.volume[volumeRepoPrefix+repoSlug("group/project")] = true sp, _ := newTestSpawner(t, f) f.failCreate = true res, err := sp.Start(context.Background(), "group/project", "", "", false) if err != nil { t.Fatalf("Start: %v", err) } if job := waitJobState(t, sp, res.SessionID, stateError); !strings.Contains(job.Message, "docker create") { t.Fatalf("job = %+v, want docker create failure", job) } f.failCreate = false f.failStart = true res2, err := sp.Start(context.Background(), "group/project", "", "", false) if err != nil { t.Fatalf("Start 2: %v", err) } if job := waitJobState(t, sp, res2.SessionID, stateError); !strings.Contains(job.Message, "docker start") { t.Fatalf("job = %+v, want docker start failure", job) } // failed worker start must remove the created container if f.countCalls(http.MethodDelete, "/containers/") == 0 { t.Fatal("failed start must clean up the created container") } } func TestSpawnerBuildHTTPErrors(t *testing.T) { useFakeGit(t, fakeGitModeOK) f := newFakeDocker() f.images = 0 f.failBuildHTTP = true // /build endpoint itself 500s sp, _ := newTestSpawner(t, f) res, err := sp.Start(context.Background(), "group/project", "", "", false) if err != nil { t.Fatalf("Start: %v", err) } if job := waitJobState(t, sp, res.SessionID, stateError); !strings.Contains(job.Message, "docker build") { t.Fatalf("job = %+v, want docker build failure", job) } // dead docker endpoint: ensureImage cannot even list images sp2, _ := newTestSpawner(t, newFakeDocker()) sp2.images0AndDead(t) if err := sp2.ensureImage(context.Background()); err == nil { t.Fatal("ensureImage with dead docker must fail") } } func (sp *Spawner) images0AndDead(t *testing.T) { t.Helper() dead := httptest.NewServer(http.NotFoundHandler()) t.Cleanup(dead.Close) dead.Close() if err := sp.cli.Close(); err != nil { t.Fatalf("close client: %v", err) } cli, err := client.NewClientWithOpts(client.WithHost("tcp://" + dead.Listener.Addr().String())) if err != nil { t.Fatalf("client: %v", err) } sp.cli = cli } // failingWriter errors on the Nth write onward. type failingWriter struct { n, failAt int } func (w *failingWriter) Write(p []byte) (int, error) { w.n++ if w.n >= w.failAt { return 0, fmt.Errorf("write boom %d", w.n) } return len(p), nil } func TestTarDirWriteErrors(t *testing.T) { dir := t.TempDir() if err := os.WriteFile(filepath.Join(dir, "f.txt"), []byte("content"), 0o644); err != nil { t.Fatalf("write: %v", err) } if err := tarDir(&failingWriter{failAt: 1}, dir); err == nil { t.Fatal("tarDir into failing writer must fail") } if err := tarDir(&failingWriter{failAt: 4}, dir); err == nil { t.Fatal("tarDir copy into failing writer must fail") } } func TestTarDirUnreadableFile(t *testing.T) { if os.Geteuid() == 0 { t.Skip("root ignores file permissions") } dir := t.TempDir() secret := filepath.Join(dir, "secret.txt") if err := os.WriteFile(secret, []byte("x"), 0o644); err != nil { t.Fatalf("write: %v", err) } if err := os.Chmod(secret, 0o000); err != nil { t.Fatalf("chmod: %v", err) } t.Cleanup(func() { _ = os.Chmod(secret, 0o644) }) var buf bytes.Buffer if err := tarDir(&buf, dir); err == nil { t.Fatal("tarDir with unreadable file must fail") } } func TestGitRunUnderivableExit(t *testing.T) { // gitRun is the exec seam: with real git present, invoking a nonexistent // subcommand must surface stderr, and trimming applies to long output. if err := gitRun(context.Background(), "", "version"); err != nil { t.Fatalf("git version: %v", err) } err := gitRun(context.Background(), "", "this-subcommand-does-not-exist") if err == nil || !strings.Contains(err.Error(), "this-subcommand-does-not-exist") { t.Fatalf("gitRun unknown subcommand = %v", err) } } // --- round-2 regression tests --- // TestTarDirSymlinks pins the symlink contract: headers carry TypeSymlink, // the real target in Linkname, and no content body (Size 0). Docker's untar // recreates symlinks via os.Symlink(Linkname, path); an empty Linkname made // every symlink-bearing repo fail to seed permanently. func TestTarDirSymlinks(t *testing.T) { dir := t.TempDir() if err := os.WriteFile(filepath.Join(dir, "file.txt"), []byte("content"), 0o644); err != nil { t.Fatalf("write: %v", err) } if err := os.Symlink("file.txt", filepath.Join(dir, "rel-link")); err != nil { t.Fatalf("rel symlink: %v", err) } if err := os.Symlink(filepath.Join(dir, "file.txt"), filepath.Join(dir, "abs-link")); err != nil { t.Fatalf("abs symlink: %v", err) } var buf bytes.Buffer if err := tarDir(&buf, dir); err != nil { t.Fatalf("tarDir: %v", err) } type linkInfo struct { linkname string size int64 } links := map[string]linkInfo{} regulars := 0 tr := tar.NewReader(bytes.NewReader(buf.Bytes())) for { hdr, err := tr.Next() if err == io.EOF { break } if err != nil { t.Fatalf("tar next: %v", err) } if hdr.Typeflag == tar.TypeSymlink { links[hdr.Name] = linkInfo{hdr.Linkname, hdr.Size} continue } if hdr.Typeflag == tar.TypeReg { regulars++ } } if regulars != 1 { t.Fatalf("regular files = %d, want 1", regulars) } want := map[string]string{ "rel-link": "file.txt", "abs-link": filepath.Join(dir, "file.txt"), } for name, target := range want { got, ok := links[name] if !ok { t.Fatalf("symlink %q missing from archive (have %v)", name, links) } if got.linkname != target { t.Fatalf("symlink %q Linkname = %q, want %q", name, got.linkname, target) } if got.size != 0 { t.Fatalf("symlink %q has Size %d, want 0 (no content body)", name, got.size) } } } // TestSpawnerRemoveSessionSurvivesClientHangup: a cancelled request context // (client hangup) must not skip the container remove — the DB row is deleted // either way, so a skipped remove would orphan the container forever. func TestSpawnerRemoveSessionSurvivesClientHangup(t *testing.T) { useFakeGit(t, fakeGitModeOK) f := newFakeDocker() sp, store := newTestSpawner(t, f) if err := store.UpsertContainer("s1", "cid-9", "group/project"); err != nil { t.Fatalf("upsert container: %v", err) } ctx, cancel := context.WithCancel(context.Background()) cancel() // client already gone if err := sp.RemoveSession(ctx, "s1"); err != nil { t.Fatalf("RemoveSession with cancelled ctx: %v", err) } if !f.hasCall(http.MethodDelete, "/containers/cid-9") { t.Fatal("remove must run on a background ctx despite the cancelled request ctx") } if _, ok, _ := store.GetContainer("s1"); ok { t.Fatal("container row must be deleted") } f.assertNoUnknown(t) } // TestSpawnerSeedVolumeUniqueNames: two seeds for the same repo must not // reuse one container name — a crash between create and the deferred remove // would then block every future spawn with a name conflict. func TestSpawnerSeedVolumeUniqueNames(t *testing.T) { useFakeGit(t, fakeGitModeOK) f := newFakeDocker() sp, _ := newTestSpawner(t, f) slug := repoSlug("group/project") repoVolume := volumeRepoPrefix + slug if err := os.MkdirAll(filepath.Join(sp.reposDir, slug), 0o755); err != nil { t.Fatalf("mkdir clone dir: %v", err) } if err := os.WriteFile(filepath.Join(sp.reposDir, slug, "README.md"), []byte("x"), 0o644); err != nil { t.Fatalf("write README: %v", err) } for _, sid := range []string{"aaaaaaaa-1111-2222-3333-444444444444", "bbbbbbbb-1111-2222-3333-444444444444"} { if err := sp.seedVolume(context.Background(), slug, repoVolume, sid); err != nil { t.Fatalf("seedVolume(%s): %v", sid, err) } } seeds := f.createsByName("lvmh-seed-") if len(seeds) != 2 { t.Fatalf("seed containers = %+v, want 2", seeds) } if seeds[0].Name == seeds[1].Name { t.Fatalf("seed names must differ per session: %q", seeds[0].Name) } for _, s := range seeds { if !strings.HasPrefix(s.Name, "lvmh-seed-"+slug+"-") { t.Fatalf("seed name %q, want prefix lvmh-seed-%s-", s.Name, slug) } } f.assertNoUnknown(t) } // TestSpawnerEnsureRepoVolumeTransientInspectError: a non-NotFound volume // inspect error must propagate, never be treated as "fresh" (which would // seed over an existing volume). func TestSpawnerEnsureRepoVolumeTransientInspectError(t *testing.T) { f := newFakeDocker() f.failVolumeInspect = true f.volume[volumeRepoPrefix+repoSlug("group/project")] = true // volume exists sp, _ := newTestSpawner(t, f) fresh, err := sp.ensureRepoVolume(context.Background(), volumeRepoPrefix+repoSlug("group/project")) if err == nil || !strings.Contains(err.Error(), "transient docker error") { t.Fatalf("ensureRepoVolume = fresh:%v err:%v, want transient error propagated", fresh, err) } if fresh { t.Fatal("transient inspect error must never report a fresh volume") } if f.countCalls(http.MethodPost, "/volumes/create") != 0 { t.Fatal("no volume create may follow a transient inspect error") } f.assertNoUnknown(t) } // TestSpawnerCloneOrUpdateRespectsCtx: git spawned via CommandContext must // die when the context is cancelled (daemon shutdown must not hang on git). func TestSpawnerCloneOrUpdateRespectsCtx(t *testing.T) { useFakeGit(t, fakeGitModeHang) // clone sleeps 30s f := newFakeDocker() sp, _ := newTestSpawner(t, f) ctx, cancel := context.WithTimeout(context.Background(), 300*time.Millisecond) defer cancel() start := time.Now() err := sp.cloneOrUpdate(ctx, "group/project", "", repoSlug("group/project")) if err == nil { t.Fatal("cloneOrUpdate must fail when ctx expires mid-clone") } if elapsed := time.Since(start); elapsed > 5*time.Second { t.Fatalf("cloneOrUpdate took %v, want prompt cancellation", elapsed) } if ctx.Err() == nil { t.Fatal("expected the context to be the failure cause") } } // TestGitBranchRespectsCtx: the branch probe must honour cancellation too. func TestGitBranchRespectsCtx(t *testing.T) { useFakeGit(t, fakeGitModeHang) ctx, cancel := context.WithTimeout(context.Background(), 300*time.Millisecond) defer cancel() start := time.Now() if got := gitBranch(ctx, t.TempDir()); got != "" { t.Fatalf("gitBranch under cancelled ctx = %q, want \"\"", got) } if elapsed := time.Since(start); elapsed > 5*time.Second { t.Fatalf("gitBranch took %v, want prompt cancellation", elapsed) } } // TestSpawnerStartDedupesImageList: Start must list images exactly once // (plus once more inside ensureImage) — a regression guard for the removed // double imageExists call. func TestSpawnerStartDedupesImageList(t *testing.T) { useFakeGit(t, fakeGitModeOK) f := newFakeDocker() sp, _ := newTestSpawner(t, f) res, err := sp.Start(context.Background(), "group/project", "", "", false) if err != nil { t.Fatalf("Start: %v", err) } waitJobState(t, sp, res.SessionID, stateRunning) if n := f.countCalls(http.MethodGet, "/images/json"); n != 2 { t.Fatalf("image list calls = %d, want 2 (Start + ensureImage)", n) } f.assertNoUnknown(t) } // TestSpawnerCustomImageUsed: a repo-registered image replaces the default // worker image in the created container. func TestSpawnerCustomImageUsed(t *testing.T) { useFakeGit(t, fakeGitModeOK) f := newFakeDocker() custom := "lvmh-worker-group--project-ab12cd" f.imageTags = map[string]bool{imageRefWorker: true, custom: true} sp, store := newTestSpawner(t, f) if err := store.SetRepoImage("group/project", custom); err != nil { t.Fatalf("set repo image: %v", err) } res, err := sp.Start(context.Background(), "group/project", "", "", false) if err != nil { t.Fatalf("Start: %v", err) } waitJobState(t, sp, res.SessionID, stateRunning) creates := f.createsByName("lvmh-agent-") if len(creates) != 1 || creates[0].Image != custom { t.Fatalf("agent create image = %+v, want %q", creates[0], custom) } // seed container still uses the default worker image seed := f.createsByName("lvmh-seed-") if len(seed) != 1 || seed[0].Image != imageRefWorker { t.Fatalf("seed create image = %+v, want default %q", seed[0], imageRefWorker) } f.assertNoUnknown(t) } // TestSpawnerCustomImageMissing: registering an image that was never built // fails the job with the ask-ops message instead of a raw docker error. func TestSpawnerCustomImageMissing(t *testing.T) { useFakeGit(t, fakeGitModeOK) f := newFakeDocker() f.imageTags = map[string]bool{imageRefWorker: true} // custom absent sp, store := newTestSpawner(t, f) custom := "lvmh-worker-group--project-ab12cd" if err := store.SetRepoImage("group/project", custom); err != nil { t.Fatalf("set repo image: %v", err) } res, err := sp.Start(context.Background(), "group/project", "", "", false) if err != nil { t.Fatalf("Start: %v", err) } job := waitJobState(t, sp, res.SessionID, stateError) want := "custom image " + custom + " not built (ask the ops agent to build it)" if job.Message != want { t.Fatalf("job message = %q, want %q", job.Message, want) } if len(f.createsByName("lvmh-agent-")) != 0 { t.Fatal("no agent container may be created for a missing image") } f.assertNoUnknown(t) } // TestSpawnerCustomImageDeleteFallsBack: deleting the registration returns // the repo to the default worker image. func TestSpawnerCustomImageDeleteFallsBack(t *testing.T) { useFakeGit(t, fakeGitModeOK) f := newFakeDocker() sp, store := newTestSpawner(t, f) if err := store.SetRepoImage("group/project", "lvmh-worker-x"); err != nil { t.Fatalf("set repo image: %v", err) } if err := store.DeleteRepoImage("group/project"); err != nil { t.Fatalf("delete repo image: %v", err) } res, err := sp.Start(context.Background(), "group/project", "", "", false) if err != nil { t.Fatalf("Start: %v", err) } waitJobState(t, sp, res.SessionID, stateRunning) creates := f.createsByName("lvmh-agent-") if len(creates) != 1 || creates[0].Image != imageRefWorker { t.Fatalf("create after deregister = %+v, want default %q", creates[0], imageRefWorker) } f.assertNoUnknown(t) } // TestSpawnerCustomImageCheckDockerDead: the custom-image presence check // wraps docker unavailability like Start does. func TestSpawnerCustomImageCheckDockerDead(t *testing.T) { f := newFakeDocker() sp, store := newTestSpawner(t, f) if err := store.SetRepoImage("group/project", "lvmh-worker-x"); err != nil { t.Fatalf("set repo image: %v", err) } sp.images0AndDead(t) _, _, err := sp.createAndStart(context.Background(), "group/project", repoSlug("group/project"), "", false, "s1") if err == nil || !strings.Contains(err.Error(), "docker unavailable") { t.Fatalf("createAndStart with dead docker = %v, want docker unavailable", err) } } func TestSpawnerModelEnv(t *testing.T) { useFakeGit(t, fakeGitModeOK) f := newFakeDocker() sp, _ := newTestSpawner(t, f) res, err := sp.Start(context.Background(), "group/project", "", "zai-renaud/glm-5.2", false) if err != nil { t.Fatalf("Start: %v", err) } waitJobState(t, sp, res.SessionID, stateRunning) creates := f.createsByName("lvmh-agent-") if len(creates) != 1 { t.Fatalf("creates = %d", len(creates)) } found := false for _, e := range creates[0].Env { if e == "LVMH_MODEL=zai-renaud/glm-5.2" { found = true } } if !found { t.Fatalf("LVMH_MODEL missing from %v", creates[0].Env) } } func TestSpawnerEmptyScratch(t *testing.T) { useFakeGit(t, fakeGitModeOK) f := newFakeDocker() sp, _ := newTestSpawner(t, f) res, err := sp.Start(context.Background(), "group/project", "", "", true) if err != nil { t.Fatalf("Start: %v", err) } waitJobState(t, sp, res.SessionID, stateRunning) creates := f.createsByName("lvmh-agent-") if len(creates) != 1 { t.Fatalf("creates = %d", len(creates)) } binds := creates[0].HostConfig.Binds for _, b := range binds { if b == "lvmh-repo-group--project-ab12cd:/workspace" { t.Fatalf("empty spawn must not mount the repo volume: %v", binds) } } hasSessions := false for _, b := range binds { if b == "lvmh-sessions:/pi-sessions" { hasSessions = true } } if !hasSessions { t.Fatalf("sessions volume missing: %v", binds) } }