Skip to content

Commit 2ebe048

Browse files
Merge pull request #119 from open-source-cloud/feat/task-runner
feat(run): devstack-native task-graph runner (spec 31)
2 parents 7175140 + 7b8fa42 commit 2ebe048

8 files changed

Lines changed: 470 additions & 16 deletions

File tree

‎docs/guide/command-reference.md‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@ See [lifecycle.md](lifecycle.md).
2424
| `up [project...]` | Bring the workspace up (network, shared infra, generate, compose up, hooks) — idempotent saga. | `--build`, `--rebuild`, `--skip-clone`, `--health-timeout`, `--no-hooks`, `--no-preflight`, `--no-provision`, `--profile`/`-p` |
2525
| `down [project...]` | Stop this workspace's project stacks and release their refs (data preserved). | — |
2626
| `shell [service] [-- cmd...]` | Open a shell (or run a command) in a service container. | `--project` |
27+
| `run <task...>` | Run a project's `tasks:` graph (deps-ordered, parallel; host or in-container). | `--project`, `--parallel`, `--dry-run`, `--json` |
2728
| `status` | Service health + last saga outcome + shared-service ref graph. | — |
2829
| `logs [service...]` | Stream logs across project + shared stacks (color-keyed). | `--follow`/`-f`, `--tail` (200), `--since`, `--timestamps`, `--no-color` |
2930
| `dashboard` | Live TUI cockpit: services, health, log tail. | `--no-stats` |

‎docs/guide/whats-next.md‎

Lines changed: 23 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -57,22 +57,29 @@ Bubble Tea theme, same non-TTY fallback), but it hasn't been built. For now,
5757
projects are authored by editing `devstack.yaml` directly
5858
([projects.md](projects.md)) or scaffolded via `init`.
5959

60-
### "Command-runner" / task projects & monorepo orchestration (Turborepo-style)
61-
62-
**Not supported — this is the genuine, biggest conceptual gap.** devstack
63-
orchestrates **containers and shared infrastructure**, not a task graph across
64-
packages. There is:
65-
66-
- no `run:` / `task:` service kind (services are containers, not scripts),
67-
- no `devstack run <task>` verb, and
68-
- no dependency-aware, monorepo-aware script runner (nothing like
69-
Turborepo/Nx pipelines).
70-
71-
Closing this would be the single largest addition: a new **non-container
72-
"task"/"script" service kind** (or a `devstack run` verb) plus a monorepo-aware
73-
task graph. We don't want to overclaim — today, if you need
74-
`build → test → deploy` task graphs across packages, use your existing task
75-
runner alongside devstack; devstack handles the infra those tasks talk to.
60+
### "Command-runner" / task projects & monorepo orchestration — ✅ shipped
61+
62+
**Now built-in.** A `tasks:` block in `devstack.yaml` declares non-container
63+
commands with `deps:` edges; `devstack run <task>` plans the dependency graph and
64+
runs it — independent tasks in parallel, output streamed and prefixed per task.
65+
`run: host` runs on your host toolchain; `run: exec` runs inside a service
66+
container via `compose exec`. Monorepo/Turborepo pipelines are covered two ways:
67+
the `turborepo` template runs `turbo run` inside its container, or you map each
68+
package's scripts into `tasks:` so `devstack run` owns the graph.
69+
70+
```yaml
71+
# devstack.yaml
72+
tasks:
73+
build: { run: host, command: ["pnpm", "build"] }
74+
test: { run: host, command: ["pnpm", "test"], deps: [build] }
75+
lint: { run: host, command: ["pnpm", "lint"] }
76+
```
77+
78+
```bash
79+
devstack run test # runs build → test
80+
devstack run test lint # build+lint in parallel, then test
81+
devstack run test --dry-run
82+
```
7683

7784
### Framework dev servers with watch mode (Next.js, NestJS) — ✅ shipped
7885

‎internal/cli/root.go‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -85,6 +85,7 @@ func NewRootCmd(opts Options) *cobra.Command {
8585
newStatusCmd(g),
8686
newUseCmd(g),
8787
newContextCmd(g),
88+
newRunCmd(g),
8889
newExposeCmd(g),
8990
newPortsCmd(g),
9091
newShellInitCmd(g),

‎internal/cli/run.go‎

Lines changed: 210 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,210 @@
1+
package cli
2+
3+
import (
4+
"bytes"
5+
"context"
6+
"fmt"
7+
"io"
8+
"os"
9+
"os/exec"
10+
"path/filepath"
11+
"runtime"
12+
"sort"
13+
"strings"
14+
"sync"
15+
16+
"github.com/spf13/cobra"
17+
"golang.org/x/sync/errgroup"
18+
19+
"github.com/open-source-cloud/devstack/internal/config"
20+
"github.com/open-source-cloud/devstack/internal/generate"
21+
"github.com/open-source-cloud/devstack/internal/task"
22+
)
23+
24+
// newRunCmd wires `run <task...>` (spec 31): execute a project's task graph. Tasks
25+
// are non-container commands with `deps` edges; they run in dependency order,
26+
// independent tasks in parallel (bounded by --parallel). `run: host` runs on the
27+
// host; `run: exec` runs inside a service container via compose exec. Streams each
28+
// task's output live, prefixed by task name. `devstack run` takes no flock — it
29+
// mutates no shared/ledger state.
30+
func newRunCmd(g *GlobalOpts) *cobra.Command {
31+
var project string
32+
var parallel int
33+
var dryRun bool
34+
cmd := &cobra.Command{
35+
Use: "run <task> [task2 ...]",
36+
Short: "Run a project's task graph (deps-ordered, parallel; host or in-container)",
37+
Args: cobra.MinimumNArgs(1),
38+
RunE: func(cmd *cobra.Command, args []string) error {
39+
mgr, closeFn, err := buildManager(cmd)
40+
if err != nil {
41+
return err
42+
}
43+
defer closeFn()
44+
45+
proj := project
46+
if proj == "" {
47+
proj = resolveActiveProject(mgr.Model, mgr.DB)
48+
}
49+
if proj == "" {
50+
return fmt.Errorf("no project selected (pass --project)")
51+
}
52+
p, ok := mgr.Model.Projects[proj]
53+
if !ok {
54+
return fmt.Errorf("unknown project %q", proj)
55+
}
56+
if len(p.Tasks) == 0 {
57+
return fmt.Errorf("project %q declares no tasks: (add a tasks: block to devstack.yaml)", proj)
58+
}
59+
layers, err := task.Plan(p.Tasks, args)
60+
if err != nil {
61+
return err
62+
}
63+
if dryRun {
64+
return renderRunPlan(cmd, g, proj, layers)
65+
}
66+
if parallel <= 0 {
67+
parallel = min(8, 2*runtime.NumCPU())
68+
}
69+
projDir := mgr.Model.ProjectDir(proj)
70+
composeFile := filepath.Join(projDir, generate.GenDir, generate.ComposeFile)
71+
r := &taskExec{out: cmd.OutOrStdout(), projDir: projDir, composeFile: composeFile}
72+
return runLayers(cmd.Context(), r, p.Tasks, layers, parallel)
73+
},
74+
}
75+
cmd.Flags().StringVar(&project, "project", "", "target project (default: the active/first project)")
76+
cmd.Flags().IntVar(&parallel, "parallel", 0, "max concurrent tasks (default min(8, 2*CPUs))")
77+
cmd.Flags().BoolVar(&dryRun, "dry-run", false, "print the resolved task DAG and exit")
78+
return cmd
79+
}
80+
81+
func renderRunPlan(cmd *cobra.Command, g *GlobalOpts, project string, layers [][]string) error {
82+
if g.JSON {
83+
return writeJSON(cmd, map[string]any{"project": project, "layers": layers})
84+
}
85+
w := cmd.OutOrStdout()
86+
fmt.Fprintf(w, "run plan for %q (%d layer(s)):\n", project, len(layers))
87+
for i, l := range layers {
88+
fmt.Fprintf(w, " %d: %s\n", i+1, strings.Join(l, ", "))
89+
}
90+
return nil
91+
}
92+
93+
// runLayers executes each layer in order; tasks within a layer run concurrently
94+
// up to `parallel`. A failing task fails the run (its layer's siblings finish).
95+
func runLayers(ctx context.Context, r *taskExec, tasks map[string]config.Task, layers [][]string, parallel int) error {
96+
for _, layer := range layers {
97+
eg, ectx := errgroup.WithContext(ctx)
98+
eg.SetLimit(parallel)
99+
for _, name := range layer {
100+
name := name
101+
t := tasks[name]
102+
eg.Go(func() error {
103+
if err := r.run(ectx, name, t); err != nil {
104+
return fmt.Errorf("task %q: %w", name, err)
105+
}
106+
return nil
107+
})
108+
}
109+
if err := eg.Wait(); err != nil {
110+
return err
111+
}
112+
}
113+
return nil
114+
}
115+
116+
// taskExec runs a single task with live, name-prefixed output. Concurrent tasks
117+
// share one output writer, so emit serializes whole (prefix+line) writes under a
118+
// mutex to keep lines intact and race-free.
119+
type taskExec struct {
120+
mu sync.Mutex
121+
out io.Writer
122+
projDir string
123+
composeFile string
124+
}
125+
126+
// emit writes prefix+data as one atomic unit under the lock.
127+
func (r *taskExec) emit(prefix string, data []byte) {
128+
r.mu.Lock()
129+
defer r.mu.Unlock()
130+
_, _ = io.WriteString(r.out, prefix)
131+
_, _ = r.out.Write(data)
132+
}
133+
134+
func (r *taskExec) run(ctx context.Context, name string, t config.Task) error {
135+
prefix := name + " | "
136+
pw := &prefixWriter{emit: r.emit, prefix: prefix}
137+
env := append(os.Environ(), envKV(t.Env)...)
138+
139+
var c *exec.Cmd
140+
if t.Run == "exec" {
141+
if t.Service == "" {
142+
return fmt.Errorf("run: exec requires a service")
143+
}
144+
args := []string{"compose", "-f", r.composeFile, "exec", "-T"}
145+
if t.Workdir != "" {
146+
args = append(args, "-w", t.Workdir)
147+
}
148+
for _, kv := range envKV(t.Env) {
149+
args = append(args, "-e", kv)
150+
}
151+
args = append(args, t.Service)
152+
args = append(args, t.Command...)
153+
c = exec.CommandContext(ctx, "docker", args...)
154+
c.Dir = r.projDir
155+
c.Env = env
156+
} else {
157+
c = exec.CommandContext(ctx, t.Command[0], t.Command[1:]...)
158+
c.Dir = taskWorkdir(r.projDir, t.Workdir)
159+
c.Env = env
160+
}
161+
c.Stdout = pw
162+
c.Stderr = pw
163+
r.emit(prefix, []byte("→ "+strings.Join(t.Command, " ")+"\n"))
164+
return c.Run()
165+
}
166+
167+
func taskWorkdir(projDir, workdir string) string {
168+
if workdir == "" {
169+
return projDir
170+
}
171+
if filepath.IsAbs(workdir) {
172+
return workdir
173+
}
174+
return filepath.Join(projDir, workdir)
175+
}
176+
177+
func envKV(m map[string]string) []string {
178+
if len(m) == 0 {
179+
return nil
180+
}
181+
out := make([]string, 0, len(m))
182+
for k, v := range m {
183+
out = append(out, k+"="+v)
184+
}
185+
sort.Strings(out)
186+
return out
187+
}
188+
189+
// prefixWriter buffers bytes and flushes each complete line through emit, so a
190+
// line and its prefix are written atomically even under concurrent tasks.
191+
type prefixWriter struct {
192+
emit func(prefix string, data []byte)
193+
prefix string
194+
buf []byte
195+
}
196+
197+
func (p *prefixWriter) Write(b []byte) (int, error) {
198+
p.buf = append(p.buf, b...)
199+
for {
200+
i := bytes.IndexByte(p.buf, '\n')
201+
if i < 0 {
202+
break
203+
}
204+
line := make([]byte, i+1)
205+
copy(line, p.buf[:i+1])
206+
p.emit(p.prefix, line)
207+
p.buf = p.buf[i+1:]
208+
}
209+
return len(b), nil
210+
}

‎internal/cli/run_test.go‎

Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,53 @@
1+
package cli
2+
3+
import (
4+
"bytes"
5+
"context"
6+
"strings"
7+
"testing"
8+
9+
"github.com/open-source-cloud/devstack/internal/config"
10+
"github.com/open-source-cloud/devstack/internal/task"
11+
)
12+
13+
func TestRunRegistered(t *testing.T) {
14+
if !findCmd(t, "run") {
15+
t.Fatal("run must be a real RunE command")
16+
}
17+
}
18+
19+
func TestRunLayersOrder(t *testing.T) {
20+
tasks := map[string]config.Task{
21+
"build": {Run: "host", Command: []string{"sh", "-lc", "echo BUILD"}},
22+
"lint": {Run: "host", Command: []string{"sh", "-lc", "echo LINT"}},
23+
"test": {Run: "host", Command: []string{"sh", "-lc", "echo TEST"}, Deps: []string{"build"}},
24+
"ci": {Run: "host", Command: []string{"sh", "-lc", "echo CI"}, Deps: []string{"test", "lint"}},
25+
}
26+
layers, err := task.Plan(tasks, []string{"ci"})
27+
if err != nil {
28+
t.Fatal(err)
29+
}
30+
var buf bytes.Buffer
31+
r := &taskExec{out: &buf, projDir: t.TempDir()}
32+
if err := runLayers(context.Background(), r, tasks, layers, 4); err != nil {
33+
t.Fatalf("run: %v", err)
34+
}
35+
out := buf.String()
36+
// Dependency order: BUILD before TEST, TEST before CI, LINT before CI.
37+
for _, pair := range [][2]string{{"BUILD", "TEST"}, {"TEST", "CI"}, {"LINT", "CI"}} {
38+
if strings.Index(out, pair[0]) >= strings.Index(out, pair[1]) {
39+
t.Errorf("%s should run before %s\n%s", pair[0], pair[1], out)
40+
}
41+
}
42+
}
43+
44+
func TestRunLayersPropagatesFailure(t *testing.T) {
45+
tasks := map[string]config.Task{
46+
"boom": {Run: "host", Command: []string{"sh", "-lc", "exit 3"}},
47+
}
48+
layers, _ := task.Plan(tasks, []string{"boom"})
49+
r := &taskExec{out: &bytes.Buffer{}, projDir: t.TempDir()}
50+
if err := runLayers(context.Background(), r, tasks, layers, 1); err == nil {
51+
t.Fatal("a failing task must fail the run")
52+
}
53+
}

‎internal/config/model.go‎

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -143,6 +143,23 @@ type Project struct {
143143
Services map[string]Service `yaml:"services" validate:"required,dive"`
144144
Hooks Hooks `yaml:"hooks"` // spec 11 — project-scope lifecycle hooks
145145
Resources []ResourceDecl `yaml:"resources" validate:"dive"` // spec 27 — declarative data-plane resources
146+
Tasks map[string]Task `yaml:"tasks" validate:"dive"` // spec 31 — non-container task graph (`devstack run`)
147+
}
148+
149+
// Task is one node in a project's task graph (spec 31): a short-lived command run
150+
// on demand by `devstack run`, NOT a container. `run: host` executes on the host
151+
// (inheriting your toolchain); `run: exec` runs inside a service container via
152+
// `compose exec`. `deps` are other task names that must complete first; the graph
153+
// is executed in dependency order (cycles are rejected at run time). `watch` marks
154+
// long-running dev-server tasks that `--watch` keeps alive.
155+
type Task struct {
156+
Command []string `yaml:"command" validate:"required,min=1"`
157+
Run string `yaml:"run" validate:"omitempty,oneof=host exec"` // default host
158+
Service string `yaml:"service"` // target for run:exec
159+
Deps []string `yaml:"deps"`
160+
Workdir string `yaml:"workdir"`
161+
Env map[string]string `yaml:"env"`
162+
Watch bool `yaml:"watch"`
146163
}
147164

148165
// ResourceDecl is one declarative data-plane resource a project needs INSIDE a

0 commit comments

Comments
 (0)