repositories / bugabinga.net
bugabinga.net
personal infrastructure for bugabinga!
owned by admin
services/luci/internal/runner/runner.go
Rawpackage 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())
}
}