From f9a98733411cfa8657e82636e0c55671086ebe46 Mon Sep 17 00:00:00 2001 From: Jon Chery Date: Wed, 3 Jun 2026 12:45:20 +0000 Subject: [PATCH 1/2] feat(P03): task execution engine with HCL specs, jobs, tasks, WaitDelay Implements Phase 3 of v0.1 Foundation: - internal/model/job.go: Job + Task models with status state machines - internal/store/migrations/0002_jobs_tasks.sql: jobs + tasks tables with FK - internal/store/job_task_repo.go: JobRepo + TaskRepo with CRUD and lifecycle updates - internal/jobspec/spec.go: HCL parser using hashicorp/hcl/v2 hclsimple - internal/jobspec/spec_test.go: 4 tests for parser - internal/engine/executor.go: parallel task executor using os/exec with Go 1.25 WaitDelay for clean process shutdown - internal/cli/job.go: orca job {run,list,stop,logs} wired to executor - testdata/hello.hcl, testdata/fail.hcl: smoke test fixtures Verified: job run executes commands, captures stdout/stderr, persists state, job stop transitions status, job logs displays captured output. All tests pass with -race. ---ci--- project: orca phase: 3 milestone: v0.1 status: execute req_covered: - REQ-004 - REQ-006 - REQ-009 - REQ-018 - REQ-020 - REQ-021 ---/ci--- --- go.mod | 11 +- go.sum | 32 +++ internal/cli/job.go | 197 ++++++++++++++-- internal/engine/executor.go | 150 ++++++++++++ internal/jobspec/spec.go | 69 ++++++ internal/jobspec/spec_test.go | 60 +++++ internal/model/job.go | 50 ++++ internal/store/job_task_repo.go | 223 ++++++++++++++++++ internal/store/migrations/0002_jobs_tasks.sql | 34 +++ testdata/fail.hcl | 7 + testdata/hello.hcl | 7 + 11 files changed, 812 insertions(+), 28 deletions(-) create mode 100644 internal/engine/executor.go create mode 100644 internal/jobspec/spec.go create mode 100644 internal/jobspec/spec_test.go create mode 100644 internal/model/job.go create mode 100644 internal/store/job_task_repo.go create mode 100644 internal/store/migrations/0002_jobs_tasks.sql create mode 100644 testdata/fail.hcl create mode 100644 testdata/hello.hcl diff --git a/go.mod b/go.mod index 8ea842a..fc43e00 100644 --- a/go.mod +++ b/go.mod @@ -2,14 +2,18 @@ module git.cloudinit.dev/coreci/orca go 1.25.0 -require github.com/spf13/cobra v1.8.1 +require ( + github.com/google/uuid v1.6.0 + github.com/hashicorp/hcl/v2 v2.24.0 + github.com/spf13/cobra v1.8.1 + modernc.org/sqlite v1.51.0 +) require ( github.com/agext/levenshtein v1.2.1 // indirect github.com/apparentlymart/go-textseg/v15 v15.0.0 // indirect github.com/dustin/go-humanize v1.0.1 // indirect - github.com/google/uuid v1.6.0 // indirect - github.com/hashicorp/hcl/v2 v2.24.0 // indirect + github.com/google/go-cmp v0.7.0 // indirect github.com/inconshreveable/mousetrap v1.1.0 // indirect github.com/mattn/go-isatty v0.0.20 // indirect github.com/mitchellh/go-wordwrap v1.0.1 // indirect @@ -25,5 +29,4 @@ require ( modernc.org/libc v1.72.3 // indirect modernc.org/mathutil v1.7.1 // indirect modernc.org/memory v1.11.0 // indirect - modernc.org/sqlite v1.51.0 // indirect ) diff --git a/go.sum b/go.sum index 6078fd6..eb3eff7 100644 --- a/go.sum +++ b/go.sum @@ -3,10 +3,20 @@ github.com/agext/levenshtein v1.2.1/go.mod h1:JEDfjyjHDjOF/1e4FlBE/PkbqA9OfWu2ki github.com/apparentlymart/go-textseg/v15 v15.0.0 h1:uYvfpb3DyLSCGWnctWKGj857c6ew1u1fNQOlOtuGxQY= github.com/apparentlymart/go-textseg/v15 v15.0.0/go.mod h1:K8XmNZdhEBkdlyDdvbmmsvpAG721bKi0joRfFdHIWJ4= github.com/cpuguy83/go-md2man/v2 v2.0.4/go.mod h1:tgQtvFlXSQOSOSIRvRPT7W67SCa46tRHOmNcaadrF8o= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY= github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto= +github.com/go-test/deep v1.0.3 h1:ZrJSEWsXzPOxaZnFteGEfooLba+ju3FYIbOrS+rQd68= +github.com/go-test/deep v1.0.3/go.mod h1:wGDj63lr65AM2AQyKZd/NYHGb0R+1RLqB8NKt3aSFNA= +github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= +github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= +github.com/google/pprof v0.0.0-20250317173921-a4b03ec1a45e h1:ijClszYn+mADRFY17kjQEVQ1XRhq2/JR1M3sGqeJoxs= +github.com/google/pprof v0.0.0-20250317173921-a4b03ec1a45e/go.mod h1:boTsfXsheKC2y+lKOCMpSfarhxDeIzfZG1jqGcPl3cA= github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/hashicorp/golang-lru/v2 v2.0.7 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs4luLUK2k= +github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM= github.com/hashicorp/hcl/v2 v2.24.0 h1:2QJdZ454DSsYGoaE6QheQZjtKZSUs9Nh2izTWiwQxvE= github.com/hashicorp/hcl/v2 v2.24.0/go.mod h1:oGoO1FIQYfn/AgyOhlg9qLC6/nOJPX3qGbkZpYAcqfM= github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8= @@ -26,6 +36,8 @@ github.com/spf13/pflag v1.0.5 h1:iy+VFUOCP1a+8yFto/drg2CJ5u0yRoB7fZw3DKv/JXA= github.com/spf13/pflag v1.0.5/go.mod h1:McXfInJRrz4CZXVZOBLb0bTZqETkiAhM9Iw0y3An2Bg= github.com/zclconf/go-cty v1.16.3 h1:osr++gw2T61A8KVYHoQiFbFd1Lh3JOCXc/jFLJXKTxk= github.com/zclconf/go-cty v1.16.3/go.mod h1:VvMs5i0vgZdhYawQNq5kePSpLAoz8u1xvZgrPIxfnZE= +github.com/zclconf/go-cty-debug v0.0.0-20240509010212-0d6042c53940 h1:4r45xpDWB6ZMSMNJFMOjqrGHynW3DIBuR2H9j0ug+Mo= +github.com/zclconf/go-cty-debug v0.0.0-20240509010212-0d6042c53940/go.mod h1:CmBdvvj3nqzfzJ6nTCIwDTPZ56aVGvDrmztiO5g3qrM= golang.org/x/mod v0.33.0 h1:tHFzIWbBifEmbwtGz65eaWyGiGZatSrT9prnU8DbVL8= golang.org/x/mod v0.33.0/go.mod h1:swjeQEj+6r7fODbD2cqrnje9PnziFuw4bmLbBZFrQ5w= golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4= @@ -39,11 +51,31 @@ golang.org/x/tools v0.42.0 h1:uNgphsn75Tdz5Ji2q36v/nsFSfR/9BRFvqhGBaJGd5k= golang.org/x/tools v0.42.0/go.mod h1:Ma6lCIwGZvHK6XtgbswSoWroEkhugApmsXyrUmBhfr0= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +modernc.org/cc/v4 v4.28.2 h1:3tQ0lf2ADtoby2EtSP+J7IE2SHwEJdP8ioR59wx7XpY= +modernc.org/cc/v4 v4.28.2/go.mod h1:OnovgIhbbMXMu1aISnJ0wvVD1KnW+cAUJkIrAWh+kVI= +modernc.org/ccgo/v4 v4.34.0 h1:yRLPFZieg532OT4rp4JFNIVcquwalMX26G95WQDqwCQ= +modernc.org/ccgo/v4 v4.34.0/go.mod h1:AS5WYMyBakQ+fhsHhtP8mWB82KTGPkNNJDGfGQCe0/A= +modernc.org/fileutil v1.4.0 h1:j6ZzNTftVS054gi281TyLjHPp6CPHr2KCxEXjEbD6SM= +modernc.org/fileutil v1.4.0/go.mod h1:EqdKFDxiByqxLk8ozOxObDSfcVOv/54xDs/DUHdvCUU= +modernc.org/gc/v2 v2.6.5 h1:nyqdV8q46KvTpZlsw66kWqwXRHdjIlJOhG6kxiV/9xI= +modernc.org/gc/v2 v2.6.5/go.mod h1:YgIahr1ypgfe7chRuJi2gD7DBQiKSLMPgBQe9oIiito= +modernc.org/gc/v3 v3.1.2 h1:ZtDCnhonXSZexk/AYsegNRV1lJGgaNZJuKjJSWKyEqo= +modernc.org/gc/v3 v3.1.2/go.mod h1:HFK/6AGESC7Ex+EZJhJ2Gni6cTaYpSMmU/cT9RmlfYY= +modernc.org/goabi0 v0.2.0 h1:HvEowk7LxcPd0eq6mVOAEMai46V+i7Jrj13t4AzuNks= +modernc.org/goabi0 v0.2.0/go.mod h1:CEFRnnJhKvWT1c1JTI3Avm+tgOWbkOu5oPA8eH8LnMI= modernc.org/libc v1.72.3 h1:ZnDF4tXn4NBXFutMMQC4vtbTFSXhhKzR73fv0beZEAU= modernc.org/libc v1.72.3/go.mod h1:dn0dZNnnn1clLyvRxLxYExxiKRZIRENOfqQ8XEeg4Qs= modernc.org/mathutil v1.7.1 h1:GCZVGXdaN8gTqB1Mf/usp1Y/hSqgI2vAGGP4jZMCxOU= modernc.org/mathutil v1.7.1/go.mod h1:4p5IwJITfppl0G4sUEDtCr4DthTaT47/N3aT6MhfgJg= modernc.org/memory v1.11.0 h1:o4QC8aMQzmcwCK3t3Ux/ZHmwFPzE6hf2Y5LbkRs+hbI= modernc.org/memory v1.11.0/go.mod h1:/JP4VbVC+K5sU2wZi9bHoq2MAkCnrt2r98UGeSK7Mjw= +modernc.org/opt v0.2.0 h1:tGyef5ApycA7FSEOMraay9SaTk5zmbx7Tu+cJs4QKZg= +modernc.org/opt v0.2.0/go.mod h1:03fq9lsNfvkYSfxrfUhZCWPk1lm4cq4N+Bh//bEtgns= +modernc.org/sortutil v1.2.1 h1:+xyoGf15mM3NMlPDnFqrteY07klSFxLElE2PVuWIJ7w= +modernc.org/sortutil v1.2.1/go.mod h1:7ZI3a3REbai7gzCLcotuw9AC4VZVpYMjDzETGsSMqJE= modernc.org/sqlite v1.51.0 h1:aH/MMSoayAIhozZ7uJbVTT9QO/VhzBf0J9tymmmuC/U= modernc.org/sqlite v1.51.0/go.mod h1:tcNzv5p84E0skkmJn038y+hWJbLQXQqEnQfeh5r2JLM= +modernc.org/strutil v1.2.1 h1:UneZBkQA+DX2Rp35KcM69cSsNES9ly8mQWD71HKlOA0= +modernc.org/strutil v1.2.1/go.mod h1:EHkiggD70koQxjVdSBM3JKM7k6L0FbGE5eymy9i3B9A= +modernc.org/token v1.1.0 h1:Xl7Ap9dKaEs5kLoOQeQmPWevfnk/DM5qcLcYlA8ys6Y= +modernc.org/token v1.1.0/go.mod h1:UGzOrNV1mAFSEB63lOFHIpNRUVMvYTc6yu1SMY/XTDM= diff --git a/internal/cli/job.go b/internal/cli/job.go index ea221b0..70ea8cd 100644 --- a/internal/cli/job.go +++ b/internal/cli/job.go @@ -1,9 +1,18 @@ package cli import ( + "context" + "errors" "fmt" + "time" + "github.com/google/uuid" "github.com/spf13/cobra" + + "git.cloudinit.dev/coreci/orca/internal/engine" + "git.cloudinit.dev/coreci/orca/internal/jobspec" + "git.cloudinit.dev/coreci/orca/internal/model" + "git.cloudinit.dev/coreci/orca/internal/store" ) var jobCmd = &cobra.Command{ @@ -12,46 +21,187 @@ var jobCmd = &cobra.Command{ Long: "Run, list, stop, and inspect orca jobs.", } +func jobExecutor() (*engine.Executor, func() error, error) { + db, closer, err := openDB() + if err != nil { + return nil, nil, err + } + jobs := store.NewJobRepo(db) + tasks := store.NewTaskRepo(db) + return engine.NewExecutor(jobs, tasks, newLogger()), closer, nil +} + var jobRunCmd = &cobra.Command{ Use: "run ", Short: "Run a job from an HCL spec file", - Long: "Submit a job spec and execute it. Implemented in Phase 3.", + Long: "Submit a job spec, execute its tasks, and persist the result.", Args: cobra.ExactArgs(1), RunE: func(cmd *cobra.Command, args []string) error { - return notImplemented("orca job run " + args[0]) + spec, err := jobspec.ParseFile(args[0]) + if err != nil { + return err + } + + ctx, cancel := context.WithTimeout(cmd.Context(), 5*time.Minute) + defer cancel() + + exec, closer, err := jobExecutor() + if err != nil { + return err + } + defer closer() + + job := &model.Job{ + ID: uuid.NewString(), + Name: spec.Job.Name, + Spec: args[0], + Status: model.JobStatusPending, + } + if err := exec.Run(ctx, job, toTaskSpecs(spec.Tasks)); err != nil { + if jsonOutput { + _ = printJSON(map[string]any{"id": job.ID, "status": "failed", "error": err.Error()}) + return err + } + fmt.Fprintf(cmd.ErrOrStderr(), "✗ Job %s failed: %v\n", job.ID, err) + return err + } + if jsonOutput { + return printJSON(map[string]any{"id": job.ID, "name": job.Name, "status": "complete"}) + } + fmt.Fprintf(cmd.OutOrStdout(), "✓ Job complete: %s (%s)\n", job.ID, job.Name) + return nil }, } var jobListCmd = &cobra.Command{ Use: "list", Short: "List all jobs", - Long: "Display all jobs and their status. Implemented in Phase 3.", + Long: "Display all jobs and their status.", RunE: func(cmd *cobra.Command, args []string) error { - return notImplemented("orca job list") + ctx, cancel := context.WithTimeout(cmd.Context(), 5*time.Second) + defer cancel() + + db, closer, err := openDB() + if err != nil { + return err + } + defer closer() + + jobs, err := store.NewJobRepo(db).List(ctx) + if err != nil { + return err + } + if jsonOutput { + return printJSON(jobs) + } + if len(jobs) == 0 { + fmt.Fprintln(cmd.OutOrStdout(), "No jobs. Use 'orca job run ' to submit one.") + return nil + } + fmt.Fprintf(cmd.OutOrStdout(), "%-36s %-20s %-12s %-8s\n", "ID", "NAME", "STATUS", "EXIT") + for _, j := range jobs { + fmt.Fprintf(cmd.OutOrStdout(), "%-36s %-20s %-12s %-8d\n", j.ID, j.Name, j.Status, j.ExitCode) + } + return nil }, } +var ( + stopID string +) + var jobStopCmd = &cobra.Command{ - Use: "stop ", + Use: "stop [job-id]", Short: "Stop a running job", - Long: "Stop a job by ID. Implemented in Phase 3.", - Args: cobra.ExactArgs(1), + Long: "Mark a job as stopped. Note: this is a soft stop (cancel context for the daemon).", + Args: cobra.MaximumNArgs(1), RunE: func(cmd *cobra.Command, args []string) error { - return notImplemented("orca job stop " + args[0]) + id := stopID + if id == "" && len(args) > 0 { + id = args[0] + } + if id == "" { + return fmt.Errorf("job id required (--id or argument)") + } + ctx, cancel := context.WithTimeout(cmd.Context(), 5*time.Second) + defer cancel() + + db, closer, err := openDB() + if err != nil { + return err + } + defer closer() + + repo := store.NewJobRepo(db) + job, err := repo.Get(ctx, id) + if err != nil { + if errors.Is(err, store.ErrNotFound) { + return fmt.Errorf("job not found: %s", id) + } + return err + } + if err := repo.UpdateStatus(ctx, id, model.JobStatusStopped, 130); err != nil { + return err + } + if jsonOutput { + return printJSON(map[string]any{"id": id, "status": "stopped", "previous_status": job.Status}) + } + fmt.Fprintf(cmd.OutOrStdout(), "✓ Job stopped: %s\n", id) + return nil }, } var jobLogsCmd = &cobra.Command{ - Use: "logs ", - Short: "Show logs for a job", - Long: "Display the logs for a job by ID. Implemented in Phase 3.", - Args: cobra.ExactArgs(1), + Use: "logs [job-id]", + Short: "Show task output for a job", + Long: "Display captured stdout/stderr for all tasks in a job.", + Args: cobra.MaximumNArgs(1), RunE: func(cmd *cobra.Command, args []string) error { - return notImplemented("orca job logs " + args[0]) + id := stopID + if id == "" && len(args) > 0 { + id = args[0] + } + if id == "" { + return fmt.Errorf("job id required (--id or argument)") + } + ctx, cancel := context.WithTimeout(cmd.Context(), 5*time.Second) + defer cancel() + + db, closer, err := openDB() + if err != nil { + return err + } + defer closer() + + taskRepo := store.NewTaskRepo(db) + tasks, err := taskRepo.ListByJob(ctx, id) + if err != nil { + return err + } + if jsonOutput { + return printJSON(tasks) + } + if len(tasks) == 0 { + fmt.Fprintln(cmd.OutOrStdout(), "No tasks for this job.") + return nil + } + for i, t := range tasks { + fmt.Fprintf(cmd.OutOrStdout(), "--- task[%d] %s (%s) exit=%d ---\n", i, t.Command, t.Status, t.ExitCode) + if t.Stdout != "" { + fmt.Fprintln(cmd.OutOrStdout(), t.Stdout) + } + if t.Stderr != "" { + fmt.Fprintln(cmd.OutOrStderr(), t.Stderr) + } + } + return nil }, } func init() { + jobStopCmd.Flags().StringVar(&stopID, "id", "", "job id") + jobLogsCmd.Flags().StringVar(&stopID, "id", "", "job id") + jobCmd.AddCommand(jobRunCmd) jobCmd.AddCommand(jobListCmd) jobCmd.AddCommand(jobStopCmd) @@ -59,16 +209,15 @@ func init() { rootCmd.AddCommand(jobCmd) } -func notImplemented(cmd string) error { - if jsonOutput { - return printJSON(map[string]any{ - "command": cmd, - "status": "not_implemented", - "phase": "1-cli-skeleton", - "next": "Phase 2-6 will implement this", - }) +func toTaskSpecs(in []jobspec.TaskSpec) []engine.TaskSpec { + out := make([]engine.TaskSpec, len(in)) + for i, t := range in { + out[i] = engine.TaskSpec{ + Name: t.Name, + Command: t.Command, + Args: t.Args, + Env: t.Env, + } } - fmt.Fprintf(rootCmd.ErrOrStderr(), "✗ %s: not yet implemented (Phase 1: CLI skeleton only)\n", cmd) - fmt.Fprintf(rootCmd.ErrOrStderr(), " see .ciagent/ROADMAP.md for the full 6-phase plan\n") - return fmt.Errorf("not implemented: %s", cmd) + return out } diff --git a/internal/engine/executor.go b/internal/engine/executor.go new file mode 100644 index 0000000..5ddc1e6 --- /dev/null +++ b/internal/engine/executor.go @@ -0,0 +1,150 @@ +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() + } +} diff --git a/internal/jobspec/spec.go b/internal/jobspec/spec.go new file mode 100644 index 0000000..4cfa80e --- /dev/null +++ b/internal/jobspec/spec.go @@ -0,0 +1,69 @@ +package jobspec + +import ( + "fmt" + "os" + "strings" + + "github.com/hashicorp/hcl/v2" + "github.com/hashicorp/hcl/v2/gohcl" + "github.com/hashicorp/hcl/v2/hclsimple" +) + +type Spec struct { + Job JobSpec `hcl:"job,block"` + Tasks []TaskSpec `hcl:"task,block"` +} + +type JobSpec struct { + Name string `hcl:"name,label"` + Type string `hcl:"type,optional"` +} + +type TaskSpec struct { + Name string `hcl:"name,label"` + Command string `hcl:"command"` + Args []string `hcl:"args,optional"` + Env []string `hcl:"env,optional"` +} + +func ParseFile(path string) (*Spec, error) { + data, err := os.ReadFile(path) + if err != nil { + return nil, fmt.Errorf("read spec file: %w", err) + } + return Parse(data, path) +} + +func Parse(data []byte, filename string) (*Spec, error) { + var spec Spec + err := hclsimple.Decode(filename, data, nil, &spec) + if err != nil { + return nil, fmt.Errorf("decode hcl: %w", err) + } + if spec.Job.Name == "" { + return nil, fmt.Errorf("spec missing job name") + } + if len(spec.Tasks) == 0 { + return nil, fmt.Errorf("spec must have at least one task") + } + for i, t := range spec.Tasks { + if t.Command == "" { + return nil, fmt.Errorf("task[%d] (%s) missing command", i, t.Name) + } + } + return &spec, nil +} + +func (s *Spec) Validate() error { + if strings.TrimSpace(s.Job.Name) == "" { + return fmt.Errorf("job name is required") + } + if len(s.Tasks) == 0 { + return fmt.Errorf("at least one task is required") + } + return nil +} + +var _ = hcl.Diagnostics{} +var _ = gohcl.DecodeBody diff --git a/internal/jobspec/spec_test.go b/internal/jobspec/spec_test.go new file mode 100644 index 0000000..5a75228 --- /dev/null +++ b/internal/jobspec/spec_test.go @@ -0,0 +1,60 @@ +package jobspec + +import ( + "testing" +) + +func TestParseValid(t *testing.T) { + hcl := ` +job "demo" { +} + +task "build" { + command = "/bin/echo" + args = ["hello", "world"] +} +` + spec, err := Parse([]byte(hcl), "test.hcl") + if err != nil { + t.Fatalf("parse: %v", err) + } + if spec.Job.Name != "demo" { + t.Errorf("expected job name 'demo', got %q", spec.Job.Name) + } + if len(spec.Tasks) != 1 { + t.Fatalf("expected 1 task, got %d", len(spec.Tasks)) + } + if spec.Tasks[0].Command != "/bin/echo" { + t.Errorf("expected command '/bin/echo', got %q", spec.Tasks[0].Command) + } + if len(spec.Tasks[0].Args) != 2 { + t.Errorf("expected 2 args, got %d", len(spec.Tasks[0].Args)) + } +} + +func TestParseMissingJob(t *testing.T) { + hcl := `task "x" { command = "/bin/echo" }` + _, err := Parse([]byte(hcl), "test.hcl") + if err == nil { + t.Fatal("expected error for missing job name") + } +} + +func TestParseNoTasks(t *testing.T) { + hcl := `job "empty" {}` + _, err := Parse([]byte(hcl), "test.hcl") + if err == nil { + t.Fatal("expected error for no tasks") + } +} + +func TestParseTaskMissingCommand(t *testing.T) { + hcl := ` +job "x" {} +task "no-cmd" {} +` + _, err := Parse([]byte(hcl), "test.hcl") + if err == nil { + t.Fatal("expected error for missing command") + } +} diff --git a/internal/model/job.go b/internal/model/job.go new file mode 100644 index 0000000..88d9d00 --- /dev/null +++ b/internal/model/job.go @@ -0,0 +1,50 @@ +package model + +import "time" + +type JobStatus string + +const ( + JobStatusPending JobStatus = "pending" + JobStatusRunning JobStatus = "running" + JobStatusComplete JobStatus = "complete" + JobStatusFailed JobStatus = "failed" + JobStatusStopped JobStatus = "stopped" +) + +type Job struct { + ID string `json:"id"` + Name string `json:"name"` + Spec string `json:"spec"` + Status JobStatus `json:"status"` + CreatedAt time.Time `json:"created_at"` + StartedAt *time.Time `json:"started_at,omitempty"` + EndedAt *time.Time `json:"ended_at,omitempty"` + ExitCode int `json:"exit_code"` +} + +type TaskStatus string + +const ( + TaskStatusPending TaskStatus = "pending" + TaskStatusRunning TaskStatus = "running" + TaskStatusComplete TaskStatus = "complete" + TaskStatusFailed TaskStatus = "failed" + TaskStatusKilled TaskStatus = "killed" +) + +type Task struct { + ID string `json:"id"` + JobID string `json:"job_id"` + Command string `json:"command"` + Args []string `json:"args"` + Env []string `json:"env,omitempty"` + PID int `json:"pid"` + ExitCode int `json:"exit_code"` + Status TaskStatus `json:"status"` + CreatedAt time.Time `json:"created_at"` + StartedAt *time.Time `json:"started_at,omitempty"` + EndedAt *time.Time `json:"ended_at,omitempty"` + Stdout string `json:"stdout,omitempty"` + Stderr string `json:"stderr,omitempty"` +} diff --git a/internal/store/job_task_repo.go b/internal/store/job_task_repo.go new file mode 100644 index 0000000..71a0f34 --- /dev/null +++ b/internal/store/job_task_repo.go @@ -0,0 +1,223 @@ +package store + +import ( + "context" + "database/sql" + "encoding/json" + "errors" + "fmt" + "time" + + "git.cloudinit.dev/coreci/orca/internal/model" +) + +type JobRepo struct { + db *sql.DB +} + +func NewJobRepo(db *sql.DB) *JobRepo { + return &JobRepo{db: db} +} + +func (r *JobRepo) Insert(ctx context.Context, j *model.Job) error { + if j.CreatedAt.IsZero() { + j.CreatedAt = time.Now().UTC() + } + if j.Status == "" { + j.Status = model.JobStatusPending + } + _, err := r.db.ExecContext(ctx, + `INSERT INTO jobs (id, name, spec, status, exit_code, created_at) VALUES (?, ?, ?, ?, ?, ?)`, + j.ID, j.Name, j.Spec, string(j.Status), j.ExitCode, j.CreatedAt) + if err != nil { + return fmt.Errorf("insert job: %w", err) + } + return nil +} + +func (r *JobRepo) Get(ctx context.Context, id string) (*model.Job, error) { + row := r.db.QueryRowContext(ctx, + `SELECT id, name, spec, status, exit_code, created_at, started_at, ended_at FROM jobs WHERE id = ?`, id) + return scanJob(row) +} + +func (r *JobRepo) List(ctx context.Context) ([]*model.Job, error) { + rows, err := r.db.QueryContext(ctx, + `SELECT id, name, spec, status, exit_code, created_at, started_at, ended_at FROM jobs ORDER BY created_at DESC`) + if err != nil { + return nil, fmt.Errorf("list jobs: %w", err) + } + defer rows.Close() + var jobs []*model.Job + for rows.Next() { + j, err := scanJob(rows) + if err != nil { + return nil, err + } + jobs = append(jobs, j) + } + return jobs, rows.Err() +} + +func (r *JobRepo) UpdateStatus(ctx context.Context, id string, status model.JobStatus, exitCode int) error { + now := time.Now().UTC() + var startedAt, endedAt *time.Time + switch status { + case model.JobStatusRunning: + startedAt = &now + case model.JobStatusComplete, model.JobStatusFailed, model.JobStatusStopped: + endedAt = &now + } + _, err := r.db.ExecContext(ctx, + `UPDATE jobs SET status = ?, exit_code = ?, started_at = COALESCE(?, started_at), ended_at = COALESCE(?, ended_at) WHERE id = ?`, + string(status), exitCode, startedAt, endedAt, id) + if err != nil { + return fmt.Errorf("update job: %w", err) + } + return nil +} + +func scanJob(s scanner) (*model.Job, error) { + var ( + j model.Job + status string + startedAt sql.NullTime + endedAt sql.NullTime + ) + err := s.Scan(&j.ID, &j.Name, &j.Spec, &status, &j.ExitCode, &j.CreatedAt, &startedAt, &endedAt) + if err == sql.ErrNoRows { + return nil, ErrNotFound + } + if err != nil { + return nil, fmt.Errorf("scan job: %w", err) + } + j.Status = model.JobStatus(status) + if startedAt.Valid { + j.StartedAt = &startedAt.Time + } + if endedAt.Valid { + j.EndedAt = &endedAt.Time + } + return &j, nil +} + +type TaskRepo struct { + db *sql.DB +} + +func NewTaskRepo(db *sql.DB) *TaskRepo { + return &TaskRepo{db: db} +} + +func (r *TaskRepo) Insert(ctx context.Context, t *model.Task) error { + if t.CreatedAt.IsZero() { + t.CreatedAt = time.Now().UTC() + } + if t.Status == "" { + t.Status = model.TaskStatusPending + } + argsJSON, _ := json.Marshal(t.Args) + envJSON, _ := json.Marshal(t.Env) + _, err := r.db.ExecContext(ctx, + `INSERT INTO tasks (id, job_id, command, args, env, pid, exit_code, status, created_at, stdout, stderr) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, + t.ID, t.JobID, t.Command, string(argsJSON), string(envJSON), + t.PID, t.ExitCode, string(t.Status), t.CreatedAt, t.Stdout, t.Stderr) + if err != nil { + return fmt.Errorf("insert task: %w", err) + } + return nil +} + +func (r *TaskRepo) Get(ctx context.Context, id string) (*model.Task, error) { + row := r.db.QueryRowContext(ctx, + `SELECT id, job_id, command, args, env, pid, exit_code, status, created_at, started_at, ended_at, stdout, stderr FROM tasks WHERE id = ?`, id) + return scanTask(row) +} + +func (r *TaskRepo) ListByJob(ctx context.Context, jobID string) ([]*model.Task, error) { + rows, err := r.db.QueryContext(ctx, + `SELECT id, job_id, command, args, env, pid, exit_code, status, created_at, started_at, ended_at, stdout, stderr FROM tasks WHERE job_id = ? ORDER BY created_at ASC`, jobID) + if err != nil { + return nil, fmt.Errorf("list tasks: %w", err) + } + defer rows.Close() + var tasks []*model.Task + for rows.Next() { + t, err := scanTask(rows) + if err != nil { + return nil, err + } + tasks = append(tasks, t) + } + return tasks, rows.Err() +} + +func (r *TaskRepo) UpdateRunning(ctx context.Context, id string, pid int) error { + now := time.Now().UTC() + _, err := r.db.ExecContext(ctx, + `UPDATE tasks SET pid = ?, status = ?, started_at = ? WHERE id = ?`, + pid, string(model.TaskStatusRunning), now, id) + if err != nil { + return fmt.Errorf("update task running: %w", err) + } + return nil +} + +func (r *TaskRepo) UpdateDone(ctx context.Context, id string, exitCode int, stdout, stderr string) error { + now := time.Now().UTC() + status := model.TaskStatusComplete + if exitCode != 0 { + status = model.TaskStatusFailed + } + _, err := r.db.ExecContext(ctx, + `UPDATE tasks SET status = ?, exit_code = ?, ended_at = ?, stdout = ?, stderr = ? WHERE id = ?`, + string(status), exitCode, now, stdout, stderr, id) + if err != nil { + return fmt.Errorf("update task done: %w", err) + } + return nil +} + +func (r *TaskRepo) UpdateKilled(ctx context.Context, id string) error { + now := time.Now().UTC() + _, err := r.db.ExecContext(ctx, + `UPDATE tasks SET status = ?, ended_at = ? WHERE id = ?`, + string(model.TaskStatusKilled), now, id) + if err != nil { + return fmt.Errorf("update task killed: %w", err) + } + return nil +} + +var _ = errors.New +var _ = json.Marshal + +func scanTask(s scanner) (*model.Task, error) { + var ( + t model.Task + status string + argsJSON string + envJSON string + startedAt sql.NullTime + endedAt sql.NullTime + ) + err := s.Scan(&t.ID, &t.JobID, &t.Command, &argsJSON, &envJSON, + &t.PID, &t.ExitCode, &status, &t.CreatedAt, &startedAt, &endedAt, &t.Stdout, &t.Stderr) + if err == sql.ErrNoRows { + return nil, ErrNotFound + } + if err != nil { + return nil, fmt.Errorf("scan task: %w", err) + } + t.Status = model.TaskStatus(status) + if startedAt.Valid { + t.StartedAt = &startedAt.Time + } + if endedAt.Valid { + t.EndedAt = &endedAt.Time + } + _ = json.Unmarshal([]byte(argsJSON), &t.Args) + _ = json.Unmarshal([]byte(envJSON), &t.Env) + return &t, nil +} diff --git a/internal/store/migrations/0002_jobs_tasks.sql b/internal/store/migrations/0002_jobs_tasks.sql new file mode 100644 index 0000000..24540fb --- /dev/null +++ b/internal/store/migrations/0002_jobs_tasks.sql @@ -0,0 +1,34 @@ +-- Jobs and tasks +CREATE TABLE IF NOT EXISTS jobs ( + id TEXT PRIMARY KEY, + name TEXT NOT NULL, + spec TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'pending', + exit_code INTEGER NOT NULL DEFAULT 0, + created_at DATETIME NOT NULL, + started_at DATETIME, + ended_at DATETIME +); + +CREATE INDEX IF NOT EXISTS idx_jobs_status ON jobs(status); +CREATE INDEX IF NOT EXISTS idx_jobs_created ON jobs(created_at); + +CREATE TABLE IF NOT EXISTS tasks ( + id TEXT PRIMARY KEY, + job_id TEXT NOT NULL, + command TEXT NOT NULL, + args TEXT NOT NULL DEFAULT '[]', + env TEXT NOT NULL DEFAULT '[]', + pid INTEGER NOT NULL DEFAULT 0, + exit_code INTEGER NOT NULL DEFAULT 0, + status TEXT NOT NULL DEFAULT 'pending', + created_at DATETIME NOT NULL, + started_at DATETIME, + ended_at DATETIME, + stdout TEXT NOT NULL DEFAULT '', + stderr TEXT NOT NULL DEFAULT '', + FOREIGN KEY (job_id) REFERENCES jobs(id) ON DELETE CASCADE +); + +CREATE INDEX IF NOT EXISTS idx_tasks_job ON tasks(job_id); +CREATE INDEX IF NOT EXISTS idx_tasks_status ON tasks(status); diff --git a/testdata/fail.hcl b/testdata/fail.hcl new file mode 100644 index 0000000..48ceafd --- /dev/null +++ b/testdata/fail.hcl @@ -0,0 +1,7 @@ +job "failing-job" { +} + +task "fail" { + command = "/bin/sh" + args = ["-c", "echo oops 1>&2; exit 1"] +} diff --git a/testdata/hello.hcl b/testdata/hello.hcl new file mode 100644 index 0000000..f125cab --- /dev/null +++ b/testdata/hello.hcl @@ -0,0 +1,7 @@ +job "hello-orca" { +} + +task "greet" { + command = "/bin/echo" + args = ["hello", "from", "orca"] +} From 857f7563190e7703f97c607a50d6b0a897d250e9 Mon Sep 17 00:00:00 2001 From: Jon Chery Date: Wed, 3 Jun 2026 12:45:30 +0000 Subject: [PATCH 2/2] docs(P03): verification - 4 layers pass - Structural: go build, go vet, gofmt all clean - Behavioral: job run/list/stop/logs work end-to-end, JSON output valid, tests pass with -race - Security: gosec/govulncheck deferred to CI - Quality: tests pass, no formatting issues ---ci--- project: orca phase: 3 milestone: v0.1 status: verify verification: structural: pass behavioral: pass security: deferred_to_ci quality: pass ---/ci---