package engine import ( "bytes" "context" "fmt" "log/slog" "os/exec" "sync" "time" "github.com/google/uuid" "git.cloudinit.dev/coreci/orca/internal/model" "git.cloudinit.dev/coreci/orca/internal/store" ) type Executor struct { jobs *store.JobRepo tasks *store.TaskRepo log *slog.Logger mu sync.Mutex } func NewExecutor(jobs *store.JobRepo, tasks *store.TaskRepo, log *slog.Logger) *Executor { if log == nil { log = slog.Default() } return &Executor{jobs: jobs, tasks: tasks, log: log} } type TaskSpec struct { Name string Command string Args []string Env []string } func (e *Executor) Run(ctx context.Context, job *model.Job, specs []TaskSpec) error { e.mu.Lock() defer e.mu.Unlock() // Insert the job first so tasks can reference it via foreign key. if err := e.jobs.Insert(ctx, job); err != nil { return err } if err := e.jobs.UpdateStatus(ctx, job.ID, model.JobStatusRunning, 0); err != nil { return err } var ( wg sync.WaitGroup failedCount int exitCode int mu sync.Mutex ) for _, ts := range specs { wg.Add(1) go func(ts TaskSpec) { defer wg.Done() if err := e.runOne(ctx, job, ts); err != nil { mu.Lock() failedCount++ e.log.Error("task failed", slog.String("job_id", job.ID), slog.String("task", ts.Name), slog.String("error", err.Error())) mu.Unlock() } }(ts) } wg.Wait() if failedCount > 0 { exitCode = 1 if err := e.jobs.UpdateStatus(ctx, job.ID, model.JobStatusFailed, exitCode); err != nil { return err } return fmt.Errorf("%d/%d tasks failed", failedCount, len(specs)) } if err := e.jobs.UpdateStatus(ctx, job.ID, model.JobStatusComplete, 0); err != nil { return err } return nil } func (e *Executor) runOne(ctx context.Context, job *model.Job, ts TaskSpec) error { task := &model.Task{ ID: uuid.NewString(), JobID: job.ID, Command: ts.Command, Args: ts.Args, Env: ts.Env, Status: model.TaskStatusPending, } if err := e.tasks.Insert(ctx, task); err != nil { return err } cmd := exec.CommandContext(ctx, ts.Command, ts.Args...) cmd.Env = append(cmd.Environ(), ts.Env...) // WaitDelay (Go 1.25+) bounds the time spent waiting on a child process // that fails to exit after the context is canceled. cmd.WaitDelay = 5 * time.Second var stdout, stderr bytes.Buffer cmd.Stdout = &stdout cmd.Stderr = &stderr if err := cmd.Start(); err != nil { _ = e.tasks.UpdateKilled(ctx, task.ID) return fmt.Errorf("start: %w", err) } if err := e.tasks.UpdateRunning(ctx, task.ID, cmd.Process.Pid); err != nil { e.log.Warn("update running failed", slog.String("error", err.Error())) } e.log.Info("task started", slog.String("job_id", job.ID), slog.String("task", ts.Name), slog.Int("pid", cmd.Process.Pid)) done := make(chan error, 1) go func() { done <- cmd.Wait() }() select { case err := <-done: exitCode := 0 if err != nil { if ee, ok := err.(*exec.ExitError); ok { exitCode = ee.ExitCode() } else { exitCode = 1 } } _ = e.tasks.UpdateDone(ctx, task.ID, exitCode, stdout.String(), stderr.String()) if err != nil { return err } return nil case <-ctx.Done(): // WaitDelay (set above) gives the process a grace period to exit // cleanly before being killed. _ = e.tasks.UpdateKilled(ctx, task.ID) return ctx.Err() } }