Luigit
repositories / bugabinga.net

bugabinga.net

personal infrastructure for bugabinga!

owned by admin

services/luci/internal/runner/runner_test.go

Raw
package runner

import (
	"context"
	"errors"
	"fmt"
	"io"
	"os"
	"os/exec"
	"path/filepath"
	"reflect"
	"strings"
	"testing"
	"time"

	"bugabinga.net/luci/internal/ciconfig"
	"bugabinga.net/luci/internal/config"
	"bugabinga.net/luci/internal/history"
	"bugabinga.net/luci/internal/logs"
	"bugabinga.net/luci/internal/podman"
	"bugabinga.net/luci/internal/publish"
	"bugabinga.net/luci/internal/trigger"
)

type fakeExecutor struct {
	calls  [][]string
	ctxs   []context.Context
	envs   [][]string
	err    error
	failAt int
	output string
	onRun  func(args []string)
}

func (f *fakeExecutor) Run(ctx context.Context, args []string, stdout io.Writer, stderr io.Writer) error {
	return f.run(ctx, args, stdout, stderr)
}

func (f *fakeExecutor) RunWithEnv(ctx context.Context, args, env []string, stdout io.Writer, stderr io.Writer) error {
	f.envs = append(f.envs, append([]string(nil), env...))
	return f.run(ctx, args, stdout, stderr)
}

func (f *fakeExecutor) run(ctx context.Context, args []string, stdout io.Writer, stderr io.Writer) error {
	f.ctxs = append(f.ctxs, ctx)
	f.calls = append(f.calls, append([]string(nil), args...))
	if f.onRun != nil {
		f.onRun(args)
	}
	output := f.output
	if output == "" {
		output = "ran\n"
	}
	_, _ = io.WriteString(stdout, output)
	if f.failAt > 0 && len(f.calls) == f.failAt {
		return os.ErrPermission
	}
	return f.err
}

func manualEvent(t *testing.T, root, repo, job, ref string) trigger.ManualEvent {
	t.Helper()
	event, err := AdmitManual(testConfig(root), trigger.ManualEvent{Kind: "manual", Repo: repo, Job: job, Ref: &ref})
	if err != nil {
		t.Fatal(err)
	}
	return event
}

func pinnedManual(t *testing.T, root, repo, job, ref string) trigger.ManualEvent {
	t.Helper()
	rev := gitOut(t, root, "--git-dir", filepath.Join(root, repo+".git"), "rev-parse", ref)
	return trigger.ManualEvent{Kind: "manual", Repo: repo, Job: job, Ref: &ref, ResolvedRev: rev}
}

func TestDefaultExecutorFactory(t *testing.T) {
	if defaultExecutor("") == nil {
		t.Fatal("nil default executor")
	}
}

func TestManualRunExecutesRequestedJob(t *testing.T) {
	root := t.TempDir()
	bare := createBareRepo(t, root, "repo", `
job "check" {
  matrix {
    node "22" "24"
  }
  image "node:${node}"
  run "echo ${node}"
}
job "skip" {
  image "alpine"
  run "false"
}
`)
	_ = bare
	ref := "main"
	executor := &fakeExecutor{}
	r := Runner{Config: testConfig(root), Executor: executor, Now: fixedNow}
	if err := r.Process("test-run", manualEvent(t, root, "repo", "check", ref)); err != nil {
		t.Fatalf("Process() error = %v", err)
	}
	if len(executor.calls) != 2 {
		t.Fatalf("executor calls = %#v", executor.calls)
	}
	if !containsArg(executor.calls[0], "node:22") || !containsArg(executor.calls[1], "node:24") {
		t.Fatalf("matrix images not interpolated: %#v", executor.calls)
	}
}

func TestArtifactsCollectAfterSuccessAndFailure(t *testing.T) {
	root := t.TempDir()
	createBareRepo(t, root, "repo", `job "check" {
  image "alpine"
  run "true"
  artifact "always" {
    path "out.txt"
  }
  artifact "failure" {
    path "failure.txt"
    after-failure true
  }
}`)
	ref := "main"
	for _, test := range []struct {
		name string
		fail bool
		want []string
	}{
		{"success", false, []string{"always/out.txt"}},
		{"failure", true, []string{"failure/failure.txt"}},
	} {
		t.Run(test.name, func(t *testing.T) {
			executor := &fakeExecutor{failAt: 0}
			if test.fail {
				executor.failAt = 1
			}
			executor.onRun = func([]string) {
				workspace := filepath.Join(root, "data", "workspaces", "repo", test.name, "001")
				_ = os.WriteFile(filepath.Join(workspace, "out.txt"), []byte("out"), 0o644)
				_ = os.WriteFile(filepath.Join(workspace, "failure.txt"), []byte("failed"), 0o644)
			}
			r := Runner{Config: testConfig(root), Executor: executor, Now: fixedNow}
			if err := r.Process(test.name, manualEvent(t, root, "repo", "check", ref)); err != nil {
				t.Fatal(err)
			}
			for _, name := range test.want {
				if _, err := os.Stat(filepath.Join(root, "data", "artifacts", test.name, "001", name)); err != nil {
					entries, _ := os.ReadDir(filepath.Join(root, "data", "artifacts", test.name, "001"))
					t.Fatalf("artifact %q: %v entries=%v calls=%d", name, err, entries, len(executor.calls))
				}
			}
			if _, err := os.Stat(filepath.Join(root, "data", "workspaces", "repo", test.name)); !os.IsNotExist(err) {
				t.Fatalf("workspace retained: %v", err)
			}
		})
	}
}

func TestDurableMountsBindExpectedSources(t *testing.T) {
	root := t.TempDir()
	createBareRepo(t, root, "repo", `
job "deploy" {
  image "alpine"
  run "true"
  mount "auth" "registry" "/root/.config/containers"
  mount "cache" "build" "/var/cache/build"
}
`)
	auth := filepath.Join(root, "data", "mounts", "auth", "repo", "deploy", "registry")
	if err := os.MkdirAll(auth, 0o700); err != nil {
		t.Fatal(err)
	}
	ref := "main"
	executor := &fakeExecutor{}
	r := Runner{Config: testConfig(root), Executor: executor, Now: fixedNow}
	if err := r.Process("mounted", manualEvent(t, root, "repo", "deploy", ref)); err != nil {
		t.Fatal(err)
	}
	args := executor.calls[0]
	for _, want := range []string{
		auth + ":/root/.config/containers:Z",
		filepath.Join(root, "data", "mounts", "cache", "repo", "deploy", "build") + ":/var/cache/build:Z",
	} {
		if !containsText(args, want) {
			t.Fatalf("mount %q missing from %#v", want, args)
		}
	}
	if _, err := os.Stat(filepath.Join(root, "data", "workspaces", "repo", "mounted", "001", "root")); !os.IsNotExist(err) {
		t.Fatalf("auth mount entered workspace: %v", err)
	}
}

func TestDurableMountLockCancelsContention(t *testing.T) {
	root := t.TempDir()
	dataDir := filepath.Join(root, "data")
	source := filepath.Join(dataDir, "mounts", "cache", "repo", "job", "build")
	if err := os.MkdirAll(source, 0o700); err != nil {
		t.Fatal(err)
	}
	declarations := []ciconfig.Mount{{Type: "cache", Name: "build", Target: "/var/cache/build"}}
	_, release, err := durableMounts(context.Background(), dataDir, "repo", "job", declarations)
	if err != nil {
		t.Fatal(err)
	}
	defer release()
	ctx, cancel := context.WithTimeout(context.Background(), 20*time.Millisecond)
	defer cancel()
	if _, _, err := durableMounts(ctx, dataDir, "repo", "job", declarations); !errors.Is(err, context.DeadlineExceeded) {
		t.Fatalf("contention error=%v", err)
	}
}

func TestMissingAuthMountBlocksExecutor(t *testing.T) {
	root := t.TempDir()
	createBareRepo(t, root, "repo", `job "deploy" {
  image "alpine"
  run "true"
  mount "auth" "registry" "/root/.config/containers"
}`)
	ref := "main"
	executor := &fakeExecutor{}
	if err := (Runner{Config: testConfig(root), Executor: executor, Now: fixedNow}).Process("no-auth", manualEvent(t, root, "repo", "deploy", ref)); err != nil {
		t.Fatal(err)
	}
	if len(executor.calls) != 0 {
		t.Fatalf("executor called with unavailable auth mount: %#v", executor.calls)
	}
}

func TestProcessDoesNotRepeatCompletedRun(t *testing.T) {
	root := t.TempDir()
	createBareRepo(t, root, "repo", `
job "check" {
  image "alpine"
  run "echo ok"
}
`)
	ref := "main"
	executor := &fakeExecutor{}
	r := Runner{Config: testConfig(root), Executor: executor, Now: fixedNow}
	event := manualEvent(t, root, "repo", "check", ref)
	if err := r.Process("completed", event); err != nil {
		t.Fatal(err)
	}
	if err := r.Process("completed", event); err != nil {
		t.Fatal(err)
	}
	if len(executor.calls) != 1 {
		t.Fatalf("executor calls=%#v", executor.calls)
	}
}

func TestRunAppliesDefaultTimeout(t *testing.T) {
	root := t.TempDir()
	createBareRepo(t, root, "repo", `
job "check" {
  image "alpine"
  run "echo ok"
}
`)
	ref := "main"
	executor := &fakeExecutor{}
	cfg := testConfig(root)
	cfg.DefaultTimeout = time.Second
	if err := (Runner{Config: cfg, Executor: executor, Now: fixedNow}).Process("test-run", manualEvent(t, root, "repo", "check", ref)); err != nil {
		t.Fatal(err)
	}
	deadline, ok := executor.ctxs[0].Deadline()
	if !ok || time.Until(deadline) > time.Second {
		t.Fatalf("deadline=%v ok=%v", deadline, ok)
	}
}

func TestPushMatchExecutesAndExecutorErrorIsRecorded(t *testing.T) {
	root := t.TempDir()
	bare := createBareRepo(t, root, "repo", `
job "check" {
  image "alpine"
  run "echo ok"
  trigger {
    push {
      branch "main"
    }
  }
}
`)
	rev := gitOut(t, root, "--git-dir", bare, "rev-parse", "main")
	executor := &fakeExecutor{err: os.ErrPermission}
	r := Runner{Config: testConfig(root), Executor: executor, Now: fixedNow}
	if err := r.Process("test-run", trigger.PushEvent{Kind: "push", Repo: "repo", Old: "0", New: rev, Ref: "refs/heads/main"}); err != nil {
		t.Fatalf("terminal executor failure returned: %v", err)
	}
	if len(executor.calls) != 1 {
		t.Fatalf("executor calls = %#v", executor.calls)
	}
	data, err := os.ReadFile(filepath.Join(root, "data", "events", "1970-01.jsonl.zst"))
	if err != nil || len(data) == 0 {
		t.Fatalf("history not written: len=%d err=%v", len(data), err)
	}
}

func TestPushPathFilterUsesChangedTree(t *testing.T) {
	root := t.TempDir()
	bare := createBareRepo(t, root, "repo", `
job "source" {
  image "alpine"
  run "echo ok"
  trigger {
    push {
      include "README.md"
    }
  }
}
job "other" {
  image "alpine"
  run "echo no"
  trigger {
    push {
      include "src/**"
    }
  }
}
`)
	rev := gitOut(t, root, "--git-dir", bare, "rev-parse", "main")
	executor := &fakeExecutor{}
	r := Runner{Config: testConfig(root), Executor: executor, Now: fixedNow}
	if err := r.Process("paths", trigger.PushEvent{Kind: "push", Repo: "repo", Old: "0", New: rev, Ref: "refs/heads/main"}); err != nil {
		t.Fatal(err)
	}
	if len(executor.calls) != 1 {
		t.Fatalf("executor calls=%#v", executor.calls)
	}
}

func TestScheduleRunsPinnedJob(t *testing.T) {
	root := t.TempDir()
	bare := createBareRepo(t, root, "repo", `
job "nightly" {
  image "alpine"
  run "echo ok"
  trigger {
    schedule "0 3 * * *"
  }
}
`)
	rev := gitOut(t, root, "--git-dir", bare, "rev-parse", "main")
	executor := &fakeExecutor{}
	r := Runner{Config: testConfig(root), Executor: executor, Now: fixedNow}
	event := trigger.ScheduleEvent{Kind: "schedule", Repo: "repo", Job: "nightly", Rev: rev, Ref: "refs/heads/main", Minute: "2026-03-04T03:00Z"}
	if err := r.Process("scheduled", event); err != nil {
		t.Fatal(err)
	}
	if len(executor.calls) != 1 {
		t.Fatalf("executor calls=%#v", executor.calls)
	}
}

func TestPushNoMatchDoesNotExecute(t *testing.T) {
	root := t.TempDir()
	bare := createBareRepo(t, root, "repo", `
job "check" {
  image "alpine"
  run "echo ok"
  trigger {
    push {
      branch "main"
    }
  }
}
`)
	rev := gitOut(t, root, "--git-dir", bare, "rev-parse", "main")
	executor := &fakeExecutor{}
	r := Runner{Config: testConfig(root), Executor: executor, Now: fixedNow}
	if err := r.Process("test-run", trigger.PushEvent{Kind: "push", Repo: "repo", Old: "0", New: rev, Ref: "refs/heads/feature"}); err != nil {
		t.Fatalf("Process() error = %v", err)
	}
	if len(executor.calls) != 0 {
		t.Fatalf("unexpected executor calls: %#v", executor.calls)
	}
}

func TestSecretsMaterializeEnvAndFile(t *testing.T) {
	root := t.TempDir()
	createBareRepo(t, root, "repo", `
job "deploy" {
  image "alpine"
  run "test -f .secrets/key"
  secret "TOKEN" env="TOKEN"
  secret "KEY" file=".secrets/key"
}
`)
	if err := os.MkdirAll(filepath.Join(root, "data", "secrets", "repo", "deploy"), 0o755); err != nil {
		t.Fatal(err)
	}
	if err := os.WriteFile(filepath.Join(root, "data", "secrets", "repo", "deploy", "TOKEN"), []byte("secret\n"), 0o600); err != nil {
		t.Fatal(err)
	}
	if err := os.WriteFile(filepath.Join(root, "data", "secrets", "repo", "deploy", "KEY"), []byte("key"), 0o600); err != nil {
		t.Fatal(err)
	}
	ref := "main"
	executor := &fakeExecutor{onRun: func(args []string) {
		if !containsText(args, ":/work/.secrets/key:ro,Z") {
			t.Fatalf("read-only file secret bind missing: %#v", args)
		}
	}}
	r := Runner{Config: testConfig(root), Executor: executor, Now: fixedNow}
	if err := r.Process("test-run", manualEvent(t, root, "repo", "deploy", ref)); err != nil {
		t.Fatalf("Process() error = %v", err)
	}
	if len(executor.envs) != 1 || containsArg(executor.envs[0], "TOKEN=secret") || !containsArg(executor.calls[0], "--env-file") {
		t.Fatalf("secret leaked into process environment or env file missing: args=%#v env=%#v", executor.calls, executor.envs)
	}
}

func TestRunnerPassesOnlyFrozenCIEnvironment(t *testing.T) {
	root := t.TempDir()
	bare := createBareRepo(t, root, "repo", `
job "check" {
  image "alpine"
  run "echo ok"
  matrix {
    node "22"
  }
  secret "TOKEN" env="TOKEN"
}
`)
	if err := os.MkdirAll(filepath.Join(root, "data", "secrets", "repo", "check"), 0o755); err != nil {
		t.Fatal(err)
	}
	if err := os.WriteFile(filepath.Join(root, "data", "secrets", "repo", "check", "TOKEN"), []byte("secret\n"), 0o600); err != nil {
		t.Fatal(err)
	}
	t.Setenv("CI_RUN_ID", "host-run")
	t.Setenv("HOST_AMBIENT", "must-not-reach-container")
	ref := "main"
	rev := gitOut(t, root, "--git-dir", bare, "rev-parse", "main")
	executor := &fakeExecutor{}
	r := Runner{Config: testConfig(root), Executor: executor, Now: fixedNow}
	if err := r.Process("frozen-run", manualEvent(t, root, "repo", "check", ref)); err != nil {
		t.Fatal(err)
	}
	want := []string{
		"CI_CHILD_ID=001",
		"CI_JOB=check[node=22]",
		"CI_MATRIX_NODE=22",
		"CI_REF=main",
		"CI_REPO=repo",
		"CI_REV=" + rev,
		"CI_RUN_ID=frozen-run",
		"CI_TRIGGER=manual",
	}
	if len(executor.envs) != 1 || !reflect.DeepEqual(executor.envs[0], want) {
		t.Fatalf("environment=%#v want=%#v", executor.envs, want)
	}
}

func TestRunLogsMaskSecretsAndStripANSI(t *testing.T) {
	root := t.TempDir()
	createBareRepo(t, root, "repo", `
job "deploy" {
  image "alpine"
  run "echo deploy"
  secret "TOKEN" env="TOKEN"
}
`)
	secretDir := filepath.Join(root, "data", "secrets", "repo", "deploy")
	if err := os.MkdirAll(secretDir, 0o755); err != nil {
		t.Fatal(err)
	}
	if err := os.WriteFile(filepath.Join(secretDir, "TOKEN"), []byte("top-secret\n"), 0o600); err != nil {
		t.Fatal(err)
	}
	ref := "main"
	r := Runner{Config: testConfig(root), Executor: &fakeExecutor{output: "\x1b[31mtop-secret\x1b[0m\n"}, Now: fixedNow}
	if err := r.Process("masked", manualEvent(t, root, "repo", "deploy", ref)); err != nil {
		t.Fatal(err)
	}
	reader, err := (logs.Store{DataDir: r.Config.DataDir}).Open("masked-001")
	if err != nil {
		t.Fatal(err)
	}
	data, err := io.ReadAll(reader)
	_ = reader.Close()
	if err != nil || string(data) != "[MASKED]" {
		t.Fatalf("log=%q err=%v", data, err)
	}
}

type fakePublisher struct {
	requests []publish.Request
	err      error
}

func (f *fakePublisher) Publish(_ context.Context, request publish.Request, output io.Writer) error {
	f.requests = append(f.requests, request)
	_, _ = io.WriteString(output, "published\n")
	return f.err
}

func TestCachesRestoreBeforeRunAndSaveAfterSuccess(t *testing.T) {
	root := t.TempDir()
	createBareRepo(t, root, "repo", `
job "check" {
  image "alpine"
  run "build"
  cache "deps" {
    path ".cache"
    key {
      file "README.md"
    }
  }
}
`)
	ref := "main"
	calls := 0
	executor := &fakeExecutor{onRun: func(args []string) {
		calls++
		workspace := volumeWorkspace(args)
		path := filepath.Join(workspace, ".cache", "value")
		if calls == 2 {
			data, err := os.ReadFile(path)
			if err != nil || string(data) != "first" {
				t.Fatalf("restored cache=%q err=%v", data, err)
			}
		}
		if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil {
			t.Fatal(err)
		}
		if err := os.WriteFile(path, []byte("first"), 0o644); err != nil {
			t.Fatal(err)
		}
	}}
	r := Runner{Config: testConfig(root), Executor: executor, Now: fixedNow}
	for _, id := range []string{"cache-1", "cache-2"} {
		if err := r.Process(id, manualEvent(t, root, "repo", "check", ref)); err != nil {
			t.Fatal(err)
		}
	}
}

func TestPublishRunsAfterCommandsAndFailureFailsChild(t *testing.T) {
	root := t.TempDir()
	bare := createBareRepo(t, root, "repo", `
job "deploy" {
  image "alpine"
  run "build"
  publish "site" {
    from "dist"
    to "site/${rev}"
  }
}
`)
	rev := gitOut(t, root, "--git-dir", bare, "rev-parse", "main")
	publisher := &fakePublisher{err: os.ErrPermission}
	executor := &fakeExecutor{onRun: func(args []string) {
		workspace := volumeWorkspace(args)
		if err := os.MkdirAll(filepath.Join(workspace, "dist"), 0o755); err != nil {
			t.Fatal(err)
		}
	}}
	r := Runner{Config: testConfig(root), Executor: executor, Publisher: publisher, Now: fixedNow}
	ref := "main"
	if err := r.Process("publish", manualEvent(t, root, "repo", "deploy", ref)); err != nil {
		t.Fatal(err)
	}
	if len(publisher.requests) != 1 || publisher.requests[0].To != "site/"+rev {
		t.Fatalf("requests=%#v", publisher.requests)
	}
	events, err := history.ReadAll(r.Config.DataDir)
	if err != nil || len(events) != 2 || events[0].Status != "failed" || events[1].Status != "failed" {
		t.Fatalf("events=%#v err=%v", events, err)
	}
}

func TestRegistryPublishRequiresFreshJobArchive(t *testing.T) {
	root := t.TempDir()
	createBareRepo(t, root, "repo", `
job "deploy" {
  image "alpine"
  run "build"
  publish "registry" {
    from "image.oci"
    to "registry.invalid/app:latest"
  }
}
`)
	publisher := &fakePublisher{}
	executor := &fakeExecutor{onRun: func(args []string) {
		if err := os.WriteFile(filepath.Join(volumeWorkspace(args), "image.oci"), []byte("archive"), 0o600); err != nil {
			t.Fatal(err)
		}
	}}
	ref := "main"
	r := Runner{Config: testConfig(root), Executor: executor, Publisher: publisher, Now: fixedNow}
	if err := r.Process("registry-fresh", manualEvent(t, root, "repo", "deploy", ref)); err != nil {
		t.Fatal(err)
	}
	if len(publisher.requests) != 1 || publisher.requests[0].From != "image.oci" || publisher.requests[0].Image != "" {
		t.Fatalf("requests=%#v", publisher.requests)
	}
}

func TestRegistryPublishRejectsPreexistingAndConflictingOutputs(t *testing.T) {
	for name, job := range map[string]string{
		"preexisting": `publish "registry" {
  from "README.md"
  to "registry.invalid/app:latest"
}`,
		"cache": `cache "output" {
  path "image.oci"
  key {
    file "README.md"
  }
}
publish "registry" {
  from "image.oci"
  to "registry.invalid/app:latest"
}`,
		"artifact": `artifact "image" {
  path "image.oci"
}
publish "registry" {
  from "image.oci"
  to "registry.invalid/app:latest"
}`,
	} {
		t.Run(name, func(t *testing.T) {
			root := t.TempDir()
			createBareRepo(t, root, "repo", "job \"deploy\" {\n  image \"alpine\"\n  run \"build\"\n  "+job+"\n}\n")
			executor := &fakeExecutor{}
			ref := "main"
			r := Runner{Config: testConfig(root), Executor: executor, Publisher: &fakePublisher{}, Now: fixedNow}
			if err := r.Process("registry-"+name, manualEvent(t, root, "repo", "deploy", ref)); err != nil {
				t.Fatal(err)
			}
			if len(executor.calls) != 0 {
				t.Fatalf("commands ran despite rejected output: %v", executor.calls)
			}
			events, err := history.ReadAll(r.Config.DataDir)
			if err != nil || len(events) != 2 || events[0].Status != "failed" || !strings.Contains(events[0].Detail, "publish registry") {
				detail := ""
				if len(events) > 0 {
					detail = events[0].Detail
				}
				t.Fatalf("count=%d status=%q detail=%q err=%v", len(events), events[0].Status, detail, err)
			}
		})
	}
}

func TestPreflightFailuresBecomeTerminalParents(t *testing.T) {
	root := t.TempDir()
	r := Runner{Config: testConfig(root), Executor: &fakeExecutor{}, Now: fixedNow}
	ref := "main"
	if err := r.Process("missing-repo", trigger.ManualEvent{Kind: "manual", Repo: "missing", Job: "check", Ref: &ref}); err != nil {
		t.Fatal(err)
	}
	createBareRepo(t, root, "repo", `
job "check" {
  image "alpine"
  run "echo ok"
}
`)
	rev := "missing"
	for i, event := range []trigger.Event{
		trigger.ManualEvent{Kind: "manual", Repo: "repo", Job: "check", Rev: &rev},
		pinnedManual(t, root, "repo", "missing", ref),
		trigger.PushEvent{Kind: "push", Repo: "missing", New: "abc", Ref: "refs/heads/main"},
		trigger.PushEvent{Kind: "push", Repo: "repo", New: "missing", Ref: "refs/heads/main"},
	} {
		if err := r.Process(fmt.Sprintf("failure-%d", i), event); err != nil {
			t.Fatal(err)
		}
	}
	events, err := history.ReadAll(r.Config.DataDir)
	if err != nil || len(events) != 5 {
		t.Fatalf("events=%#v err=%v", events, err)
	}
	for _, event := range events {
		if event.ChildID != "run" || event.Status != "failed" || event.StartedAt.IsZero() {
			t.Fatalf("event=%#v", event)
		}
	}
}

func TestConfigFailuresBecomeTerminalParents(t *testing.T) {
	root := t.TempDir()
	bare := createBareRepo(t, root, "badci", `job "x" {`)
	rev := gitOut(t, root, "--git-dir", bare, "rev-parse", "main")
	ref := "main"
	r := Runner{Config: testConfig(root), Executor: &fakeExecutor{}, Now: fixedNow}
	if err := r.Process("manual", trigger.ManualEvent{Kind: "manual", Repo: "badci", Job: "x", Ref: &ref}); err != nil {
		t.Fatal(err)
	}
	if err := r.Process("push", trigger.PushEvent{Kind: "push", Repo: "badci", New: rev, Ref: "refs/heads/main"}); err != nil {
		t.Fatal(err)
	}
}

func TestLoadRepoConfigAndRunExpandedErrorPaths(t *testing.T) {
	root := t.TempDir()
	bare := createBareRepo(t, root, "badci", `job "x" {`)
	rev := gitOut(t, root, "--git-dir", bare, "rev-parse", "main")
	r := Runner{Config: testConfig(root), Executor: &fakeExecutor{}, Now: fixedNow}
	if _, _, err := r.loadRepoConfig(bare, rev, "badci", "run"); err == nil {
		t.Fatal("bad ci config accepted")
	}
	file := filepath.Join(root, "file")
	if err := os.WriteFile(file, []byte("x"), 0o644); err != nil {
		t.Fatal(err)
	}
	r.Config.DataDir = file
	p := plan{repo: "badci", repoPath: bare, rev: rev, parentJob: "check", trigger: "manual"}
	if _, err := r.runExpanded("run", "001", p, expandedJob("check")); err == nil {
		t.Fatal("runExpanded accepted bad log dir")
	}
	r.Config = testConfig(root)
	if _, _, err := r.loadRepoConfig(filepath.Join(root, "missing.git"), rev, "badci", "run"); err == nil {
		t.Fatal("loadRepoConfig accepted checkout failure")
	}
	p.repoPath = filepath.Join(root, "missing.git")
	result, err := r.runExpanded("missing-checkout", "001", p, expandedJob("check"))
	if err != nil || result.status != "failed" {
		t.Fatalf("result=%#v err=%v", result, err)
	}

	oldDefault := defaultExecutor
	fake := &fakeExecutor{}
	r.Executor = nil
	defaultExecutor = func(string) podman.Executor { return fake }
	defer func() { defaultExecutor = oldDefault }()
	p.repoPath = bare
	if _, err := r.runExpanded("default", "001", p, expandedJob("check")); err != nil {
		t.Fatal(err)
	}
	if len(fake.calls) != 1 {
		t.Fatalf("default executor calls = %#v", fake.calls)
	}
}

func TestMaterializeSecretsRejectsUnsafeAndFileFailures(t *testing.T) {
	root := t.TempDir()
	r := Runner{Config: testConfig(root), Executor: &fakeExecutor{}, Now: fixedNow}
	if _, _, _, _, err := r.materializeSecrets("repo", ciconfigJob("../job", "TOKEN", "TOKEN", ""), "run", "001"); err == nil {
		t.Fatal("unsafe secret path accepted")
	}
	secretDir := filepath.Join(root, "data", "secrets", "repo", "deploy")
	if err := os.MkdirAll(secretDir, 0o755); err != nil {
		t.Fatal(err)
	}
	if err := os.WriteFile(filepath.Join(secretDir, "KEY"), []byte("key"), 0o600); err != nil {
		t.Fatal(err)
	}
	mounts, _, _, cleanup, err := r.materializeSecrets("repo", ciconfigJob("deploy", "KEY", "", ".secrets/key"), "run", "001")
	if err != nil || len(mounts) != 1 || mounts[0].Target != "/work/.secrets/key" || !mounts[0].ReadOnly {
		t.Fatalf("secret mount=%#v err=%v", mounts, err)
	}
	path := mounts[0].Source
	cleanup()
	if _, err := os.Stat(path); !os.IsNotExist(err) {
		t.Fatalf("secret temp remains: %v", err)
	}
}

func TestMissingSecretFailsBeforeExecutor(t *testing.T) {
	root := t.TempDir()
	createBareRepo(t, root, "repo", `
job "deploy" {
  image "alpine"
  run "echo deploy"
  secret "TOKEN" env="TOKEN"
}
`)
	ref := "main"
	executor := &fakeExecutor{}
	r := Runner{Config: testConfig(root), Executor: executor, Now: fixedNow}
	if err := r.Process("test-run", manualEvent(t, root, "repo", "deploy", ref)); err != nil {
		t.Fatalf("terminal secret failure returned: %v", err)
	}
	if len(executor.calls) != 0 {
		t.Fatalf("executor called despite missing secret: %#v", executor.calls)
	}
}

func TestProcessRejectsUnsupportedEventAndDrainPrintsKinds(t *testing.T) {
	root := t.TempDir()
	r := Runner{Config: testConfig(root), Now: fixedNow}
	if err := r.Process("test-run", nil); err == nil {
		t.Fatal("unsupported event accepted")
	}
	ref := "main"
	var out testWriter
	Drain([]trigger.Event{
		trigger.PushEvent{Kind: "push", Repo: "repo", New: "abc", Ref: "refs/heads/main"},
		trigger.ManualEvent{Kind: "manual", Repo: "repo", Job: "job", Ref: &ref},
	}, &out)
	if out.text != "push\nmanual\n" {
		t.Fatalf("drain = %q", out.text)
	}
}

func fixedNow() time.Time { return time.Unix(42, 0).UTC() }

func expandedJob(name string) ciconfig.ExpandedJob {
	return ciconfig.ExpandedJob{Job: ciconfig.Job{Name: name, Image: "alpine", Run: []string{"echo ok"}}, Name: name, Values: map[string]string{}}
}

func ciconfigJob(name, secretName, env, file string) ciconfig.Job {
	return ciconfig.Job{Name: name, Secrets: []ciconfig.Secret{{Name: secretName, Env: env, File: file}}}
}

func testConfig(root string) config.Config {
	return config.Config{
		DataDir:        filepath.Join(root, "data"),
		RepoRoots:      []string{root},
		InboxDir:       filepath.Join(root, "inbox"),
		DefaultMemory:  "2g",
		DefaultCPUs:    "2",
		DefaultTimeout: time.Minute,
	}
}

func createBareRepo(t *testing.T, root string, name string, ci string) string {
	t.Helper()
	work := filepath.Join(root, "work-"+name)
	runGit(t, root, "init", work)
	runGit(t, work, "config", "user.email", "test@example.invalid")
	runGit(t, work, "config", "user.name", "Test")
	if err := os.MkdirAll(filepath.Join(work, ".ci"), 0o755); err != nil {
		t.Fatal(err)
	}
	if err := os.WriteFile(filepath.Join(work, ".ci", "ci.kdl"), []byte(ci), 0o644); err != nil {
		t.Fatal(err)
	}
	if err := os.WriteFile(filepath.Join(work, "README.md"), []byte("hello\n"), 0o644); err != nil {
		t.Fatal(err)
	}
	runGit(t, work, "add", ".")
	runGit(t, work, "commit", "-m", "initial")
	runGit(t, work, "branch", "-M", "main")
	bare := filepath.Join(root, name+".git")
	runGit(t, root, "clone", "--bare", work, bare)
	return bare
}

func runGit(t *testing.T, dir string, args ...string) {
	t.Helper()
	cmd := exec.Command("git", args...)
	cmd.Dir = dir
	if out, err := cmd.CombinedOutput(); err != nil {
		t.Fatalf("git %v failed: %v: %s", args, err, out)
	}
}

func gitOut(t *testing.T, dir string, args ...string) string {
	t.Helper()
	cmd := exec.Command("git", args...)
	cmd.Dir = dir
	out, err := cmd.Output()
	if err != nil {
		t.Fatal(err)
	}
	return string(out[:len(out)-1])
}

func containsArg(args []string, want string) bool {
	for _, arg := range args {
		if arg == want {
			return true
		}
	}
	return false
}

func containsText(args []string, want string) bool {
	for _, arg := range args {
		if strings.Contains(arg, want) {
			return true
		}
	}
	return false
}

func volumeWorkspace(args []string) string {
	for i := 0; i < len(args)-1; i++ {
		if args[i] == "--volume" {
			volume := args[i+1]
			return volume[:len(volume)-len(":/work:Z")]
		}
	}
	return ""
}

type testWriter struct{ text string }

func (w *testWriter) Write(p []byte) (int, error) {
	w.text += string(p)
	return len(p), nil
}

func TestStableOrdinalChildrenAndParent(t *testing.T) {
	root := t.TempDir()
	createBareRepo(t, root, "repo", `
job "check" {
  matrix {
    node "22" "24"
  }
  image "node:${node}"
  run "echo ${node}"
}
`)
	ref := "main"
	r := Runner{Config: testConfig(root), Executor: &fakeExecutor{}, Now: fixedNow}
	if err := r.Process("42", manualEvent(t, root, "repo", "check", ref)); err != nil {
		t.Fatal(err)
	}
	events, err := history.ReadAll(r.Config.DataDir)
	if err != nil {
		t.Fatal(err)
	}
	var got []string
	for _, event := range events {
		got = append(got, event.RunID+":"+event.ChildID+":"+event.Status)
	}
	want := []string{"42:001:success", "42:002:success", "42:run:success"}
	if !reflect.DeepEqual(got, want) {
		t.Fatalf("got=%v want=%v", got, want)
	}
	for _, child := range []string{"001", "002"} {
		if _, err := os.Stat(filepath.Join(r.Config.DataDir, "logs", "42-"+child+".log.zst")); err != nil {
			t.Fatalf("log %s: %v", child, err)
		}
	}
	if _, err := os.Stat(filepath.Join(r.Config.DataDir, "workspaces", "repo", "42")); !os.IsNotExist(err) {
		t.Fatalf("workspace remains: %v", err)
	}
	if _, err := os.Stat(filepath.Join(r.Config.DataDir, "active", "42.json")); !os.IsNotExist(err) {
		t.Fatalf("active remains: %v", err)
	}
	parent := events[2]
	if parent.Job != "check" || parent.Ref != "main" || parent.Trigger != "manual" || parent.Rev == "" || parent.StartedAt.IsZero() || parent.Time.IsZero() {
		t.Fatalf("parent=%#v", parent)
	}
}

func TestFailureSkipsRemainingChildrenAndCompletesParent(t *testing.T) {
	root := t.TempDir()
	createBareRepo(t, root, "repo", `
job "check" {
  matrix {
    node "20" "22" "24"
  }
  image "node:${node}"
  run "echo ${node}"
}
`)
	ref := "main"
	r := Runner{Config: testConfig(root), Executor: &fakeExecutor{failAt: 2}, Now: fixedNow}
	if err := r.Process("42", manualEvent(t, root, "repo", "check", ref)); err != nil {
		t.Fatal(err)
	}
	events, err := history.ReadAll(r.Config.DataDir)
	if err != nil {
		t.Fatal(err)
	}
	var got []string
	for _, event := range events {
		got = append(got, event.ChildID+":"+event.Status)
	}
	want := []string{"001:success", "002:failed", "003:skipped", "run:failed"}
	if !reflect.DeepEqual(got, want) {
		t.Fatalf("got=%v want=%v", got, want)
	}
	if !events[2].StartedAt.IsZero() {
		t.Fatalf("skipped start=%v", events[2].StartedAt)
	}
	if _, err := os.Stat(filepath.Join(r.Config.DataDir, "workspaces", "repo", "42")); !os.IsNotExist(err) {
		t.Fatalf("workspace remains after failure: %v", err)
	}
}

func TestActivePersistenceFailureStopsBeforeExecution(t *testing.T) {
	root := t.TempDir()
	createBareRepo(t, root, "repo", `
job "check" {
  image "alpine"
  run "echo ok"
}
`)
	if err := os.MkdirAll(filepath.Join(root, "data"), 0o755); err != nil {
		t.Fatal(err)
	}
	if err := os.WriteFile(filepath.Join(root, "data", "active"), []byte("x"), 0o644); err != nil {
		t.Fatal(err)
	}
	ref := "main"
	executor := &fakeExecutor{}
	r := Runner{Config: testConfig(root), Executor: executor, Now: fixedNow}
	if err := r.Process("42", manualEvent(t, root, "repo", "check", ref)); err == nil {
		t.Fatal("active failure hidden")
	}
	if len(executor.calls) != 0 {
		t.Fatalf("executor calls=%v", executor.calls)
	}
}

func TestPreflightRecordsOnlyResolvedRevision(t *testing.T) {
	root := t.TempDir()
	bare := createBareRepo(t, root, "badci", `job "x" {`)
	resolved := gitOut(t, root, "--git-dir", bare, "rev-parse", "main")
	ref := "main"
	r := Runner{Config: testConfig(root), Executor: &fakeExecutor{}, Now: fixedNow}
	if err := r.Process("resolved", trigger.ManualEvent{Kind: "manual", Repo: "badci", Job: "x", Ref: &ref, ResolvedRev: resolved}); err != nil {
		t.Fatal(err)
	}
	missing := "not-a-revision"
	if err := r.Process("unresolved", trigger.ManualEvent{Kind: "manual", Repo: "badci", Job: "x", Rev: &missing}); err != nil {
		t.Fatal(err)
	}
	events, err := history.ReadAll(r.Config.DataDir)
	if err != nil || len(events) != 2 {
		t.Fatalf("events=%#v err=%v", events, err)
	}
	if events[0].Rev != resolved || events[1].Rev != "" {
		t.Fatalf("revisions=%q %q want %q empty", events[0].Rev, events[1].Rev, resolved)
	}
}

func TestManualAdmissionPinsRevisionAcrossBranchMove(t *testing.T) {
	root := t.TempDir()
	bare := createBareRepo(t, root, "repo", "job \"check\" {\n  image \"alpine\"\n  run \"echo initial\"\n}\n")
	ref := "main"
	event := manualEvent(t, root, "repo", "check", ref)
	initial := event.ResolvedRev
	work := filepath.Join(root, "move")
	runGit(t, root, "clone", bare, work)
	runGit(t, work, "config", "user.email", "test@example.invalid")
	runGit(t, work, "config", "user.name", "Test")
	if err := os.WriteFile(filepath.Join(work, ".ci", "ci.kdl"), []byte("job \"check\" {\n  image \"alpine\"\n  run \"echo changed\"\n}\n"), 0o644); err != nil {
		t.Fatal(err)
	}
	runGit(t, work, "add", ".ci/ci.kdl")
	runGit(t, work, "commit", "-m", "move branch")
	runGit(t, work, "push", "origin", "HEAD:main")
	executor := &fakeExecutor{}
	r := Runner{Config: testConfig(root), Executor: executor, Now: fixedNow}
	if err := r.Process("pinned", event); err != nil {
		t.Fatal(err)
	}
	if len(executor.calls) != 1 || !containsText(executor.calls[0], "echo initial") || containsText(executor.calls[0], "echo changed") {
		t.Fatalf("executed=%v", executor.calls)
	}
	events, err := history.ReadAll(r.Config.DataDir)
	if err != nil || len(events) != 2 || events[1].Rev != initial || events[1].Ref != ref {
		t.Fatalf("events=%#v err=%v", events, err)
	}
}

func TestLegacyManualEventFailsBeforeExecution(t *testing.T) {
	root := t.TempDir()
	createBareRepo(t, root, "repo", "job \"check\" {\n  image \"alpine\"\n  run \"true\"\n}\n")
	ref := "main"
	executor := &fakeExecutor{}
	if err := (Runner{Config: testConfig(root), Executor: executor, Now: fixedNow}).Process("legacy", trigger.ManualEvent{Kind: "manual", Repo: "repo", Job: "check", Ref: &ref}); err != nil {
		t.Fatal(err)
	}
	if len(executor.calls) != 0 {
		t.Fatalf("executed=%v", executor.calls)
	}
}

func TestRepositoryLookupFailureDoesNotClaimRequestedRevision(t *testing.T) {
	root := t.TempDir()
	requested := "unverified"
	r := Runner{Config: testConfig(root), Executor: &fakeExecutor{}, Now: fixedNow}
	if err := r.Process("missing", trigger.ManualEvent{Kind: "manual", Repo: "missing", Job: "check", Rev: &requested}); err != nil {
		t.Fatal(err)
	}
	events, err := history.ReadAll(r.Config.DataDir)
	if err != nil || len(events) != 1 || events[0].Rev != "" {
		t.Fatalf("events=%#v err=%v", events, err)
	}
}