Luigit
repositories / bugabinga.net

bugabinga.net

personal infrastructure for bugabinga!

owned by admin

services/luci/internal/runner/runner.go

Raw
package runner

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

	"bugabinga.net/luci/internal/artifact"
	"bugabinga.net/luci/internal/cache"
	"bugabinga.net/luci/internal/ciconfig"
	"bugabinga.net/luci/internal/config"
	"bugabinga.net/luci/internal/gitrepo"
	"bugabinga.net/luci/internal/history"
	"bugabinga.net/luci/internal/logs"
	"bugabinga.net/luci/internal/podman"
	"bugabinga.net/luci/internal/publish"
	"bugabinga.net/luci/internal/runstate"
	"bugabinga.net/luci/internal/trigger"
)

type Runner struct {
	Config    config.Config
	Executor  podman.Executor
	Publisher publish.Executor
	Now       func() time.Time
}

var defaultExecutor = func(socket string) podman.Executor { return podman.CLI{Socket: socket} }
var defaultPublisher = func(cfg config.Config) publish.Executor {
	return publish.Service{
		RegistryAuthFile: cfg.PublishRegistryAuthFile,
		PkgRoot:          cfg.PublishPkgRoot,
		SiteRoot:         cfg.PublishSiteRoot,
	}
}

type plan struct {
	repo      string
	parentJob string
	repoPath  string
	rev       string
	ref       string
	trigger   string
	subject   string
	children  []ciconfig.ExpandedJob
}

type childResult struct {
	status string
	detail string
}

func (r Runner) Process(runID string, event trigger.Event) error {
	return r.ProcessContext(context.Background(), runID, event)
}

// ProcessContext permits cancellation while waiting for durable job mounts.
func (r Runner) ProcessContext(ctx context.Context, runID string, event trigger.Event) error {
	completed, err := r.history().Contains(runID, "run")
	if err != nil {
		return err
	}
	if completed {
		return nil
	}
	switch event.(type) {
	case trigger.ManualEvent, trigger.PushEvent, trigger.ScheduleEvent:
	default:
		return fmt.Errorf("unsupported event %T", event)
	}
	started := r.now()
	p, err := r.makePlan(event, runID)
	defer os.RemoveAll(filepath.Join(r.Config.DataDir, "workspaces", p.repo, runID))
	if err != nil {
		if p.parentJob == "" {
			p = eventPlan(event)
		}
		_, appendErr := r.history().AppendOnce(history.Event{RunID: runID, ChildID: "run", Repo: p.repo, Job: p.parentJob, Rev: p.rev, Ref: p.ref, Trigger: p.trigger, Status: "failed", Subject: p.subject, Detail: err.Error(), StartedAt: started, Time: r.now()})
		return appendErr
	}
	failed := false
	for i, child := range p.children {
		childID := fmt.Sprintf("%03d", i+1)
		if failed {
			_, err := r.history().AppendOnce(history.Event{RunID: runID, ChildID: childID, Repo: p.repo, Job: child.Name, Rev: p.rev, Ref: p.ref, Trigger: p.trigger, Status: "skipped", Detail: "skipped after earlier failure", Matrix: child.Values, Time: r.now()})
			if err != nil {
				return err
			}
			continue
		}
		result, err := r.runExpandedContext(ctx, runID, childID, p, child)
		if err != nil {
			return err
		}
		failed = result.status == "failed"
	}
	status := "success"
	if failed {
		status = "failed"
	}
	_, err = r.history().AppendOnce(history.Event{RunID: runID, ChildID: "run", Repo: p.repo, Job: p.parentJob, Rev: p.rev, Ref: p.ref, Trigger: p.trigger, Status: status, Subject: p.subject, StartedAt: started, Time: r.now()})
	return err
}

func (r Runner) makePlan(event trigger.Event, runID string) (plan, error) {
	switch event := event.(type) {
	case trigger.ManualEvent:
		return r.manualPlan(event, runID)
	case trigger.PushEvent:
		return r.pushPlan(event, runID)
	case trigger.ScheduleEvent:
		return r.schedulePlan(event, runID)
	default:
		return plan{}, fmt.Errorf("unsupported event %T", event)
	}
}

// AdmitManual resolves a manual request and validates its selected job before it enters the inbox.
func AdmitManual(cfg config.Config, event trigger.ManualEvent) (trigger.ManualEvent, error) {
	if err := event.Validate(); err != nil {
		return trigger.ManualEvent{}, err
	}
	repoPath, err := gitrepo.RepoPath(cfg.RepoRoots, event.Repo)
	if err != nil {
		return trigger.ManualEvent{}, err
	}
	ref, rev := "", ""
	if event.Ref != nil {
		ref = *event.Ref
	}
	if event.Rev != nil {
		rev = *event.Rev
	}
	resolved, err := gitrepo.Resolve(repoPath, ref, rev)
	if err != nil {
		return trigger.ManualEvent{}, err
	}
	files, err := gitrepo.FilesAtRev(repoPath, resolved, ".ci")
	if err != nil {
		return trigger.ManualEvent{}, err
	}
	repository, err := ciconfig.ParseFiles(files)
	if err != nil {
		return trigger.ManualEvent{}, err
	}
	if _, ok := repository.Job(event.Job); !ok {
		return trigger.ManualEvent{}, fmt.Errorf("job %q not found", event.Job)
	}
	event.ResolvedRev = resolved
	return event, nil
}

func (r Runner) manualPlan(event trigger.ManualEvent, runID string) (plan, error) {
	p := eventPlan(event)
	if p.rev == "" {
		return p, fmt.Errorf("manual event was not admitted; resubmit the run")
	}
	repoPath, err := gitrepo.RepoPath(r.Config.RepoRoots, event.Repo)
	if err != nil {
		return p, err
	}
	p.repoPath = repoPath
	p.subject = gitrepo.Subject(repoPath, p.rev)
	cfg, _, err := r.loadRepoConfig(repoPath, p.rev, event.Repo, runID)
	if err != nil {
		return p, err
	}
	job, ok := cfg.Job(event.Job)
	if !ok {
		return p, fmt.Errorf("job %q not found", event.Job)
	}
	p.children = job.ExpandMatrix()
	return p, nil
}

func (r Runner) schedulePlan(event trigger.ScheduleEvent, runID string) (plan, error) {
	p := eventPlan(event)
	repoPath, err := gitrepo.RepoPath(r.Config.RepoRoots, event.Repo)
	if err != nil {
		return p, err
	}
	resolved, err := gitrepo.Resolve(repoPath, "", event.Rev)
	if err != nil {
		return p, err
	}
	p.repoPath, p.rev = repoPath, resolved
	p.subject = gitrepo.Subject(repoPath, resolved)
	cfg, _, err := r.loadRepoConfig(repoPath, resolved, event.Repo, runID)
	if err != nil {
		return p, err
	}
	job, ok := cfg.Job(event.Job)
	if !ok {
		return p, fmt.Errorf("job %q not found", event.Job)
	}
	p.children = job.ExpandMatrix()
	return p, nil
}

func (r Runner) pushPlan(event trigger.PushEvent, runID string) (plan, error) {
	p := eventPlan(event)
	requestedRev := p.rev
	p.rev = ""
	repoPath, err := gitrepo.RepoPath(r.Config.RepoRoots, event.Repo)
	if err != nil {
		return p, err
	}
	resolved, err := gitrepo.Resolve(repoPath, "", requestedRev)
	if err != nil {
		p.rev = ""
		return p, err
	}
	p.repoPath, p.rev = repoPath, resolved
	p.subject = gitrepo.Subject(repoPath, resolved)
	cfg, _, err := r.loadRepoConfig(repoPath, resolved, event.Repo, runID)
	if err != nil {
		return p, err
	}
	var changedPaths []string
	for _, job := range cfg.Jobs {
		if job.NeedsChangedPaths() {
			changedPaths, err = gitrepo.ChangedPaths(repoPath, event.Old, resolved)
			if err != nil {
				return p, err
			}
			break
		}
	}
	for _, job := range cfg.Jobs {
		if job.MatchesPush(event.Ref, changedPaths) {
			p.children = append(p.children, job.ExpandMatrix()...)
		}
	}
	return p, nil
}

func eventPlan(event trigger.Event) plan {
	switch event := event.(type) {
	case trigger.ManualEvent:
		p := plan{repo: event.Repo, parentJob: event.Job, rev: event.ResolvedRev, trigger: "manual"}
		if event.Ref != nil {
			p.ref = *event.Ref
		}
		return p
	case trigger.PushEvent:
		return plan{repo: event.Repo, parentJob: "push", rev: event.New, ref: event.Ref, trigger: "push"}
	case trigger.ScheduleEvent:
		return plan{repo: event.Repo, parentJob: event.Job, rev: event.Rev, ref: event.Ref, trigger: "schedule"}
	default:
		return plan{parentJob: "run", trigger: "unknown"}
	}
}

func (r Runner) loadRepoConfig(repoPath, rev, repo, runID string) (ciconfig.Config, string, error) {
	workspace := filepath.Join(r.Config.DataDir, "workspaces", repo, runID, "config")
	if err := gitrepo.Checkout(repoPath, rev, workspace); err != nil {
		return ciconfig.Config{}, "", err
	}
	cfg, err := ciconfig.Load(workspace)
	return cfg, workspace, err
}

func (r Runner) runExpanded(runID, childID string, p plan, expanded ciconfig.ExpandedJob) (childResult, error) {
	return r.runExpandedContext(context.Background(), runID, childID, p, expanded)
}

func (r Runner) runExpandedContext(ctx context.Context, runID, childID string, p plan, expanded ciconfig.ExpandedJob) (childResult, error) {
	started := r.now()
	activeStore := runstate.Store{DataDir: r.Config.DataDir}
	if err := activeStore.Write(runstate.Active{RunID: runID, ChildID: childID, Repo: p.repo, Job: expanded.Name, Rev: p.rev, Ref: p.ref, Trigger: p.trigger, StartedAt: started}); err != nil {
		return childResult{}, err
	}
	logName := runID + "-" + childID
	logStore := logs.Store{DataDir: r.Config.DataDir}
	logFile, err := logStore.CreateLive(logName)
	if err != nil {
		return childResult{}, err
	}

	workspace := filepath.Join(r.Config.DataDir, "workspaces", p.repo, runID, childID)
	var logWriter *logs.Sanitizer
	var logErr error
	runErr := gitrepo.Checkout(p.repoPath, p.rev, workspace)
	if runErr == nil {
		runErr = registryOutputsAbsent(workspace, expanded)
	}
	if runErr == nil {
		var secretValues [][]byte
		var secretMounts []podman.Mount
		var envFile string
		var removeSecrets func()
		secretMounts, envFile, secretValues, removeSecrets, runErr = r.materializeSecrets(p.repo, expanded.Job, runID, childID)
		defer removeSecrets()
		logWriter = logs.NewBoundedSanitizer(logFile, secretValues)
		if runErr == nil {
			env, envErr := frozenEnvironment(runID, childID, p, expanded, nil)
			if envErr != nil {
				runErr = envErr
			}
			resolvedCaches := make([]ciconfig.Cache, len(expanded.Job.Caches))
			cacheKeys := make([]string, len(expanded.Job.Caches))
			cacheStore := cache.Store{DataDir: r.Config.DataDir, MaxVersions: r.Config.CacheMaxVersions, MaxBytes: r.Config.CacheMaxBytes}
			for i, entry := range expanded.Job.Caches {
				if runErr != nil {
					break
				}
				entry.Path = ciconfig.Interpolate(entry.Path, expanded.Values)
				entry.KeyFiles = append([]string(nil), entry.KeyFiles...)
				for j := range entry.KeyFiles {
					entry.KeyFiles[j] = ciconfig.Interpolate(entry.KeyFiles[j], expanded.Values)
				}
				resolvedCaches[i] = entry
				cacheKeys[i], runErr = cache.Key(workspace, entry)
				if runErr == nil {
					var outcome string
					outcome, runErr = cacheStore.Restore(p.repo, expanded.Job.Name, entry, cacheKeys[i], workspace)
					if runErr == nil {
						_, runErr = fmt.Fprintf(logWriter, "[luci] cache %s: %s\n", entry.Name, outcome)
						if runErr != nil {
							logErr = errors.Join(logErr, runErr)
						}
					}
				}
				if runErr != nil {
					break
				}
			}
			runCtx, cancel := context.WithTimeout(ctx, r.Config.DefaultTimeout)
			if runErr == nil {
				image := ciconfig.Interpolate(expanded.Job.Image, expanded.Values)
				commands := make([]string, 0, len(expanded.Job.Run))
				for _, command := range expanded.Job.Run {
					commands = append(commands, ciconfig.Interpolate(command, expanded.Values))
				}
				if r.Executor == nil {
					r.Executor = defaultExecutor(r.Config.PodmanSocket)
				}
				mounts, releaseMounts, mountErr := durableMounts(runCtx, r.Config.DataDir, p.repo, expanded.Job.Name, expanded.Job.Mounts)
				if mountErr != nil {
					runErr = mountErr
				} else {
					defer releaseMounts()
					spec := podman.Spec{Image: image, Workspace: workspace, Memory: r.Config.DefaultMemory, CPUs: r.Config.DefaultCPUs, Env: env, EnvFile: envFile, Mounts: append(mounts, secretMounts...), Commands: commands, Timeout: r.Config.DefaultTimeout}
					args := podman.BuildArgs(spec)
					if executor, ok := r.Executor.(podman.EnvironmentExecutor); ok {
						runErr = executor.RunWithEnv(runCtx, args, podman.Environment(spec), logWriter, logWriter)
					} else {
						runErr = r.Executor.Run(runCtx, args, logWriter, logWriter)
					}
				}
			}
			if runErr == nil {
				for i, entry := range resolvedCaches {
					if runErr = cacheStore.Save(p.repo, expanded.Job.Name, entry, cacheKeys[i], workspace); runErr != nil {
						break
					}
				}
			}
			if runErr == nil && len(expanded.Job.Publishes) > 0 {
				if r.Publisher == nil {
					r.Publisher = defaultPublisher(r.Config)
				}
				values := make(map[string]string, len(expanded.Values)+1)
				for key, value := range expanded.Values {
					values[key] = value
				}
				values["rev"] = p.rev
				for _, declaration := range expanded.Job.Publishes {
					runErr = r.Publisher.Publish(runCtx, publish.Request{Adapter: declaration.Adapter, Image: ciconfig.Interpolate(declaration.Image, values), From: ciconfig.Interpolate(declaration.From, values), To: ciconfig.Interpolate(declaration.To, values), Workspace: workspace}, logWriter)
					if runErr != nil {
						break
					}
				}
			}
			artifactStore := artifact.Store{DataDir: r.Config.DataDir}
			for _, declaration := range expanded.Job.Artifacts {
				if runErr != nil && !declaration.AfterFailure {
					continue
				}
				if err := artifactStore.Collect(workspace, runID, childID, declaration); err != nil {
					runErr = errors.Join(runErr, err)
				}
			}
			cancel()
		}
	}
	if logWriter == nil {
		logWriter = logs.NewBoundedSanitizer(logFile, nil)
	}
	if runErr != nil {
		if _, err := fmt.Fprintln(logWriter, runErr); err != nil {
			logErr = errors.Join(logErr, err)
		}
	}
	if err := logWriter.Close(); err != nil {
		logErr = errors.Join(logErr, err)
	}
	if err := logFile.Close(); err != nil {
		logErr = errors.Join(logErr, err)
	}
	if _, err := logStore.Finalize(logName); err != nil {
		return childResult{}, errors.Join(logErr, err)
	}
	if logErr != nil {
		return childResult{}, logErr
	}
	result := childResult{status: "success"}
	if runErr != nil {
		result.status, result.detail = "failed", runErr.Error()
	}
	_, err = r.history().AppendOnce(history.Event{RunID: runID, ChildID: childID, Repo: p.repo, Job: expanded.Name, Rev: p.rev, Ref: p.ref, Trigger: p.trigger, Status: result.status, Subject: p.subject, Detail: result.detail, Matrix: expanded.Values, StartedAt: started, Time: r.now()})
	if err != nil {
		return childResult{}, err
	}
	if err := activeStore.Remove(runID); err != nil {
		return childResult{}, err
	}
	return result, nil
}

func registryOutputsAbsent(workspace string, expanded ciconfig.ExpandedJob) error {
	for _, declaration := range expanded.Job.Publishes {
		if declaration.Adapter != "registry" {
			continue
		}
		from := ciconfig.Interpolate(declaration.From, expanded.Values)
		clean := filepath.Clean(from)
		if strings.TrimSpace(from) == "" || filepath.IsAbs(from) || clean == "." || clean == ".." || strings.HasPrefix(clean, ".."+string(filepath.Separator)) {
			return fmt.Errorf("publish registry has unsafe from path %q", from)
		}
		for _, cache := range expanded.Job.Caches {
			path := ciconfig.Interpolate(cache.Path, expanded.Values)
			if pathsOverlap(clean, filepath.Clean(path)) {
				return fmt.Errorf("publish registry from %q conflicts with cache path %q", from, path)
			}
		}
		for _, artifact := range expanded.Job.Artifacts {
			for _, path := range artifact.Paths {
				if pathsOverlap(clean, filepath.Clean(path)) {
					return fmt.Errorf("publish registry from %q conflicts with artifact path %q", from, path)
				}
			}
		}
		if _, err := os.Lstat(filepath.Join(workspace, clean)); err == nil {
			return fmt.Errorf("publish registry output %q already exists before commands", from)
		} else if !os.IsNotExist(err) {
			return err
		}
	}
	return nil
}

func pathsOverlap(left, right string) bool {
	return left == right || strings.HasPrefix(left, right+string(filepath.Separator)) || strings.HasPrefix(right, left+string(filepath.Separator))
}

func (r Runner) history() history.Store { return history.Store{DataDir: r.Config.DataDir, Now: r.Now} }
func (r Runner) now() time.Time {
	if r.Now != nil {
		return r.Now()
	}
	return time.Now()
}

func frozenEnvironment(runID, childID string, p plan, expanded ciconfig.ExpandedJob, secrets map[string]string) (map[string]string, error) {
	env := map[string]string{
		"CI_RUN_ID":   runID,
		"CI_CHILD_ID": childID,
		"CI_REPO":     p.repo,
		"CI_JOB":      expanded.Name,
		"CI_REV":      p.rev,
		"CI_REF":      p.ref,
		"CI_TRIGGER":  p.trigger,
	}
	for key, value := range expanded.Values {
		name := "CI_MATRIX_" + strings.ToUpper(key)
		if _, exists := env[name]; exists {
			return nil, fmt.Errorf("matrix env collision %q", name)
		}
		env[name] = value
	}
	for name, value := range secrets {
		if _, exists := env[name]; exists {
			return nil, fmt.Errorf("secret env %q collides with CI metadata", name)
		}
		env[name] = value
	}
	return env, nil
}

func (r Runner) materializeSecrets(repo string, job ciconfig.Job, runID, childID string) ([]podman.Mount, string, [][]byte, func(), error) {
	paths, err := job.SecretPaths(r.Config.DataDir, repo)
	if err != nil {
		return nil, "", nil, func() {}, err
	}
	tmpRoot := filepath.Join(r.Config.DataDir, "tmp")
	if err := os.MkdirAll(tmpRoot, 0o700); err != nil {
		return nil, "", nil, func() {}, err
	}
	dir, err := os.MkdirTemp(tmpRoot, "secrets-"+runID+"-"+childID+"-")
	if err != nil {
		return nil, "", nil, func() {}, err
	}
	cleanup := func() { _ = os.RemoveAll(dir) }
	var mounts []podman.Mount
	var values [][]byte
	var envLines []string
	for _, secret := range job.Secrets {
		data, err := os.ReadFile(paths[secret.Name])
		if err != nil {
			cleanup()
			return nil, "", nil, func() {}, fmt.Errorf("missing secret %s", secret.Name)
		}
		if len(data) > 0 {
			values = append(values, data)
		}
		if secret.Env != "" {
			value := strings.TrimRight(string(data), "\r\n")
			if strings.ContainsAny(value, "\r\n") {
				cleanup()
				return nil, "", nil, func() {}, fmt.Errorf("secret %s cannot be an environment value containing a newline", secret.Name)
			}
			envLines = append(envLines, secret.Env+"="+value)
		}
		if secret.File != "" {
			file := filepath.Join(dir, secret.Name)
			if err := os.WriteFile(file, data, 0o600); err != nil {
				cleanup()
				return nil, "", nil, func() {}, err
			}
			mounts = append(mounts, podman.Mount{Source: file, Target: "/work/" + filepath.ToSlash(secret.File), ReadOnly: true})
		}
	}
	envFile := ""
	if len(envLines) > 0 {
		envFile = filepath.Join(dir, "env")
		if err := os.WriteFile(envFile, []byte(strings.Join(envLines, "\n")+"\n"), 0o600); err != nil {
			cleanup()
			return nil, "", nil, func() {}, err
		}
	}
	return mounts, envFile, values, cleanup, nil
}

func Drain(events []trigger.Event, writer io.Writer) {
	for _, event := range events {
		fmt.Fprintln(writer, event.EventKind())
	}
}