Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 2c53ad6213 | |||
| c3819dde12 |
@@ -1 +1 @@
|
|||||||
{ "phase": "P05", "stage": "verify", "milestone": "v0.9", "phase_role": "execution", "updated_at": "2026-08-05T04:15:00Z", "milestone_complete": false, "verify": { "build": "pass", "go_test": "23/23", "bats": "20/20", "gofmt": "clean", "verify_reqs": "90 consistent" } }
|
{ "phase": "P06", "stage": "verify", "milestone": "v0.9", "phase_role": "execution", "updated_at": "2026-08-05T04:30:00Z", "milestone_complete": false, "verify": { "build": "pass", "go_test": "22/22", "gofmt": "clean", "verify_reqs": "pending" } }
|
||||||
|
|||||||
+118
-7
@@ -42,13 +42,26 @@ type SystemdEmitter struct{}
|
|||||||
const unitNamePrefix = "orca-v1-"
|
const unitNamePrefix = "orca-v1-"
|
||||||
|
|
||||||
// Render renders the systemd unit file for a process-runtime workload.
|
// Render renders the systemd unit file for a process-runtime workload.
|
||||||
// The unit name is /etc/systemd/system/<unitNamePrefix><spec.Name>.service
|
|
||||||
// and the content is a [Service] block with ExecStart, optional
|
|
||||||
// ExecStartPost (lifecycle.post_start), optional ExecStop
|
|
||||||
// (lifecycle.pre_stop), and the R-007 socket-plumbing lines
|
|
||||||
// (RuntimeDirectory=, optional TCP-bind ExecStartPre). Mode is 0644.
|
|
||||||
//
|
//
|
||||||
// The rendered shape is:
|
// When the spec has no Tasks (the single-process case, the historical
|
||||||
|
// shape), the unit name is
|
||||||
|
// /etc/systemd/system/<unitNamePrefix><spec.Name>.service and the
|
||||||
|
// content is a [Service] block with ExecStart, optional ExecStartPost
|
||||||
|
// (lifecycle.post_start), optional ExecStop (lifecycle.pre_stop), and
|
||||||
|
// the R-007 socket-plumbing lines (RuntimeDirectory=, optional
|
||||||
|
// TCP-bind ExecStartPre). Mode is 0644.
|
||||||
|
//
|
||||||
|
// When the spec has a task group (P06, spec.Tasks non-empty), the
|
||||||
|
// alloc is multi-process and Render emits one systemd unit per task
|
||||||
|
// (`orca-v1-alloc-<alloc-id>-<task-name>.service`) plus a single
|
||||||
|
// grouping target unit (`orca-v1-alloc-<alloc-id>.target`) that
|
||||||
|
// starts/stops all tasks together. Each per-task unit carries
|
||||||
|
// `PartOf=orca-v1-alloc-<alloc-id>.target` and is
|
||||||
|
// `WantedBy=multi-user.target` so the task starts at boot. Tasks
|
||||||
|
// that omit their own runtime inherit the top-level spec.Runtime as
|
||||||
|
// the per-group default.
|
||||||
|
//
|
||||||
|
// The rendered shape (single-process) is:
|
||||||
//
|
//
|
||||||
// [Service]
|
// [Service]
|
||||||
// ExecStart=<runtime command>
|
// ExecStart=<runtime command>
|
||||||
@@ -62,7 +75,8 @@ const unitNamePrefix = "orca-v1-"
|
|||||||
//
|
//
|
||||||
// Returns an error if the spec is nil, the spec is missing its name,
|
// Returns an error if the spec is nil, the spec is missing its name,
|
||||||
// the runtime block is nil, or the runtime command is empty (a
|
// the runtime block is nil, or the runtime command is empty (a
|
||||||
// workload with no command has nothing to ExecStart).
|
// workload with no command has nothing to ExecStart). For task groups,
|
||||||
|
// returns an error if any task has no resolvable runtime command.
|
||||||
func (SystemdEmitter) Render(spec *jobspec.WorkloadSpec, node *Node) ([]File, error) {
|
func (SystemdEmitter) Render(spec *jobspec.WorkloadSpec, node *Node) ([]File, error) {
|
||||||
if spec == nil {
|
if spec == nil {
|
||||||
return nil, errors.New("emitter/systemd: spec is nil")
|
return nil, errors.New("emitter/systemd: spec is nil")
|
||||||
@@ -70,6 +84,9 @@ func (SystemdEmitter) Render(spec *jobspec.WorkloadSpec, node *Node) ([]File, er
|
|||||||
if strings.TrimSpace(spec.Name) == "" {
|
if strings.TrimSpace(spec.Name) == "" {
|
||||||
return nil, errors.New("emitter/systemd: spec name is empty")
|
return nil, errors.New("emitter/systemd: spec name is empty")
|
||||||
}
|
}
|
||||||
|
if len(spec.Tasks) > 0 {
|
||||||
|
return renderTaskGroup(spec, node)
|
||||||
|
}
|
||||||
if spec.Runtime == nil {
|
if spec.Runtime == nil {
|
||||||
return nil, errors.New("emitter/systemd: runtime block is nil")
|
return nil, errors.New("emitter/systemd: runtime block is nil")
|
||||||
}
|
}
|
||||||
@@ -81,6 +98,100 @@ func (SystemdEmitter) Render(spec *jobspec.WorkloadSpec, node *Node) ([]File, er
|
|||||||
return []File{{Path: path, Content: content, Mode: "0644"}}, nil
|
return []File{{Path: path, Content: content, Mode: "0644"}}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// renderTaskGroup renders one systemd unit per task plus the grouping
|
||||||
|
// target unit. Each task's runtime falls back to the top-level
|
||||||
|
// spec.Runtime when the task omits its own. Tasks with no resolvable
|
||||||
|
// command (no task.Command, no task.Runtime.Command, no top-level
|
||||||
|
// Runtime) return an error.
|
||||||
|
func renderTaskGroup(spec *jobspec.WorkloadSpec, node *Node) ([]File, error) {
|
||||||
|
allocID := spec.Name
|
||||||
|
targetUnit := fmt.Sprintf("%salloc-%s.target", unitNamePrefix, allocID)
|
||||||
|
targetPath := fmt.Sprintf("/etc/systemd/system/%s", targetUnit)
|
||||||
|
var files []File
|
||||||
|
for _, task := range spec.Tasks {
|
||||||
|
rt := taskRuntime(spec, &task)
|
||||||
|
if rt == nil {
|
||||||
|
return nil, fmt.Errorf("emitter/systemd: task %q has no runtime (set tasks[].runtime or top-level runtime)", task.Name)
|
||||||
|
}
|
||||||
|
cmd := taskCommand(spec, &task, rt)
|
||||||
|
if strings.TrimSpace(cmd) == "" {
|
||||||
|
return nil, fmt.Errorf("emitter/systemd: task %q command is empty", task.Name)
|
||||||
|
}
|
||||||
|
unitName := fmt.Sprintf("%salloc-%s-%s.service", unitNamePrefix, allocID, task.Name)
|
||||||
|
path := fmt.Sprintf("/etc/systemd/system/%s", unitName)
|
||||||
|
content := renderTaskUnit(spec, &task, rt, cmd, targetUnit)
|
||||||
|
files = append(files, File{Path: path, Content: content, Mode: "0644"})
|
||||||
|
}
|
||||||
|
files = append(files, File{
|
||||||
|
Path: targetPath,
|
||||||
|
Content: renderTargetUnit(targetUnit, spec, allocID),
|
||||||
|
Mode: "0644",
|
||||||
|
})
|
||||||
|
return files, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// taskRuntime returns the effective runtime for a task: the task's own
|
||||||
|
// runtime when set, otherwise the top-level spec.Runtime (the per-group
|
||||||
|
// default). Returns nil when neither is set.
|
||||||
|
func taskRuntime(spec *jobspec.WorkloadSpec, task *jobspec.TaskGroupTask) *jobspec.RuntimeBlock {
|
||||||
|
if task.Runtime != nil {
|
||||||
|
return task.Runtime
|
||||||
|
}
|
||||||
|
return spec.Runtime
|
||||||
|
}
|
||||||
|
|
||||||
|
// taskCommand returns the ExecStart command for a task. A task-level
|
||||||
|
// Command takes precedence; otherwise the task's runtime command is
|
||||||
|
// used; otherwise the top-level runtime command is used. Returns an
|
||||||
|
// empty string when none is set.
|
||||||
|
func taskCommand(spec *jobspec.WorkloadSpec, task *jobspec.TaskGroupTask, rt *jobspec.RuntimeBlock) string {
|
||||||
|
if strings.TrimSpace(task.Command) != "" {
|
||||||
|
return task.Command
|
||||||
|
}
|
||||||
|
if rt != nil && strings.TrimSpace(rt.Command) != "" {
|
||||||
|
return rt.Command
|
||||||
|
}
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
|
||||||
|
// renderTaskUnit renders a single per-task systemd [Unit]+[Service]
|
||||||
|
// block. The unit is `PartOf=` the alloc target and
|
||||||
|
// `WantedBy=multi-user.target` so it starts at boot and stops with the
|
||||||
|
// group. The [Service] block carries the task's ExecStart and the
|
||||||
|
// socket-plumbing lines derived from the spec's ports.
|
||||||
|
func renderTaskUnit(spec *jobspec.WorkloadSpec, task *jobspec.TaskGroupTask, rt *jobspec.RuntimeBlock, cmd, targetUnit string) string {
|
||||||
|
var b strings.Builder
|
||||||
|
b.WriteString("[Unit]\n")
|
||||||
|
b.WriteString(fmt.Sprintf("Description=orca alloc task %s\n", task.Name))
|
||||||
|
b.WriteString(fmt.Sprintf("PartOf=%s\n", targetUnit))
|
||||||
|
b.WriteString("\n[Service]\n")
|
||||||
|
b.WriteString(fmt.Sprintf("ExecStart=%s\n", cmd))
|
||||||
|
for _, line := range (SocketEmitter{}).RenderSocketLines(spec) {
|
||||||
|
b.WriteString(line)
|
||||||
|
b.WriteString("\n")
|
||||||
|
}
|
||||||
|
b.WriteString("\n[Install]\n")
|
||||||
|
b.WriteString("WantedBy=multi-user.target\n")
|
||||||
|
return b.String()
|
||||||
|
}
|
||||||
|
|
||||||
|
// renderTargetUnit renders the grouping target unit
|
||||||
|
// (`orca-v1-alloc-<alloc-id>.target`) that starts/stops all tasks
|
||||||
|
// together. The [Unit] block lists every per-task unit under Wants=
|
||||||
|
// so `systemctl start <target>` brings them all up, and
|
||||||
|
// `systemctl stop <target>` tears them down (PartOf= propagates stop).
|
||||||
|
func renderTargetUnit(targetUnit string, spec *jobspec.WorkloadSpec, allocID string) string {
|
||||||
|
var b strings.Builder
|
||||||
|
b.WriteString("[Unit]\n")
|
||||||
|
b.WriteString(fmt.Sprintf("Description=orca alloc %s task group\n", allocID))
|
||||||
|
for _, task := range spec.Tasks {
|
||||||
|
b.WriteString(fmt.Sprintf("Wants=%salloc-%s-%s.service\n", unitNamePrefix, allocID, task.Name))
|
||||||
|
}
|
||||||
|
b.WriteString("\n[Install]\n")
|
||||||
|
b.WriteString("WantedBy=multi-user.target\n")
|
||||||
|
return b.String()
|
||||||
|
}
|
||||||
|
|
||||||
// renderSystemdUnit renders the full [Service] block for the spec,
|
// renderSystemdUnit renders the full [Service] block for the spec,
|
||||||
// including ExecStart, lifecycle hooks (ExecStartPost, ExecStop), and
|
// including ExecStart, lifecycle hooks (ExecStartPost, ExecStop), and
|
||||||
// the R-007 socket-plumbing lines (RuntimeDirectory=, optional
|
// the R-007 socket-plumbing lines (RuntimeDirectory=, optional
|
||||||
|
|||||||
@@ -154,3 +154,190 @@ func TestSystemdEmitter_UnitNamePrefix(t *testing.T) {
|
|||||||
t.Errorf("unitNamePrefix = %q, want orca-v1-", unitNamePrefix)
|
t.Errorf("unitNamePrefix = %q, want orca-v1-", unitNamePrefix)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestSystemdEmitter_TaskGroupTwoTasks(t *testing.T) {
|
||||||
|
// P06: a task group with two tasks renders one unit per task plus
|
||||||
|
// a grouping target unit. Each per-task unit is
|
||||||
|
// `orca-v1-alloc-<alloc-id>-<task-name>.service`, carries
|
||||||
|
// `PartOf=orca-v1-alloc-<alloc-id>.target`, and is
|
||||||
|
// `WantedBy=multi-user.target`. The target unit lists every
|
||||||
|
// per-task unit under Wants=.
|
||||||
|
spec := &jobspec.WorkloadSpec{
|
||||||
|
Kind: "Service",
|
||||||
|
Name: "web",
|
||||||
|
Tasks: []jobspec.TaskGroupTask{
|
||||||
|
{
|
||||||
|
Name: "app",
|
||||||
|
Command: "/usr/bin/httpd -f",
|
||||||
|
Runtime: &jobspec.RuntimeBlock{OneOf: "process"},
|
||||||
|
},
|
||||||
|
{
|
||||||
|
Name: "sidecar",
|
||||||
|
Command: "/bin/wasm-runner sidecar.wasm",
|
||||||
|
Runtime: &jobspec.RuntimeBlock{OneOf: "wasm"},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
}
|
||||||
|
files, err := SystemdEmitter{}.Render(spec, &Node{Hostname: "n1"})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Render: %v", err)
|
||||||
|
}
|
||||||
|
// 2 per-task units + 1 target unit.
|
||||||
|
if len(files) != 3 {
|
||||||
|
t.Fatalf("got %d files, want 3 (2 per-task units + 1 target)", len(files))
|
||||||
|
}
|
||||||
|
wantApp := "/etc/systemd/system/orca-v1-alloc-web-app.service"
|
||||||
|
wantSide := "/etc/systemd/system/orca-v1-alloc-web-sidecar.service"
|
||||||
|
wantTarget := "/etc/systemd/system/orca-v1-alloc-web.target"
|
||||||
|
paths := make(map[string]*File, len(files))
|
||||||
|
for i := range files {
|
||||||
|
paths[files[i].Path] = &files[i]
|
||||||
|
}
|
||||||
|
if _, ok := paths[wantApp]; !ok {
|
||||||
|
t.Errorf("missing per-task unit %q; got paths %v", wantApp, filePaths(files))
|
||||||
|
}
|
||||||
|
if _, ok := paths[wantSide]; !ok {
|
||||||
|
t.Errorf("missing per-task unit %q; got paths %v", wantSide, filePaths(files))
|
||||||
|
}
|
||||||
|
if _, ok := paths[wantTarget]; !ok {
|
||||||
|
t.Errorf("missing target unit %q; got paths %v", wantTarget, filePaths(files))
|
||||||
|
}
|
||||||
|
if _, ok := paths[wantTarget]; !ok {
|
||||||
|
t.Errorf("missing target unit %q; got paths %v", wantTarget, filePaths(files))
|
||||||
|
}
|
||||||
|
// Verify PartOf relations and ExecStart on per-task units.
|
||||||
|
app := paths[wantApp]
|
||||||
|
if !strings.Contains(app.Content, "PartOf=orca-v1-alloc-web.target") {
|
||||||
|
t.Errorf("app unit missing PartOf=orca-v1-alloc-web.target\n%s", app.Content)
|
||||||
|
}
|
||||||
|
if !strings.Contains(app.Content, "ExecStart=/usr/bin/httpd -f") {
|
||||||
|
t.Errorf("app unit missing ExecStart=/usr/bin/httpd -f\n%s", app.Content)
|
||||||
|
}
|
||||||
|
if !strings.Contains(app.Content, "WantedBy=multi-user.target") {
|
||||||
|
t.Errorf("app unit missing WantedBy=multi-user.target\n%s", app.Content)
|
||||||
|
}
|
||||||
|
side := paths[wantSide]
|
||||||
|
if !strings.Contains(side.Content, "PartOf=orca-v1-alloc-web.target") {
|
||||||
|
t.Errorf("sidecar unit missing PartOf=orca-v1-alloc-web.target\n%s", side.Content)
|
||||||
|
}
|
||||||
|
if !strings.Contains(side.Content, "ExecStart=/bin/wasm-runner sidecar.wasm") {
|
||||||
|
t.Errorf("sidecar unit missing ExecStart\n%s", side.Content)
|
||||||
|
}
|
||||||
|
// Verify the target unit Wants= both per-task units.
|
||||||
|
target := paths[wantTarget]
|
||||||
|
if !strings.Contains(target.Content, "Wants=orca-v1-alloc-web-app.service") {
|
||||||
|
t.Errorf("target missing Wants=...app.service\n%s", target.Content)
|
||||||
|
}
|
||||||
|
if !strings.Contains(target.Content, "Wants=orca-v1-alloc-web-sidecar.service") {
|
||||||
|
t.Errorf("target missing Wants=...sidecar.service\n%s", target.Content)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestSystemdEmitter_TaskGroupInheritsTopLevelRuntime(t *testing.T) {
|
||||||
|
// P06: a task that omits its own runtime inherits the top-level
|
||||||
|
// spec.Runtime as the per-group default. The per-task unit's
|
||||||
|
// ExecStart must come from the top-level runtime command when
|
||||||
|
// the task has no own command and no own runtime.
|
||||||
|
spec := &jobspec.WorkloadSpec{
|
||||||
|
Kind: "Service",
|
||||||
|
Name: "web",
|
||||||
|
Runtime: &jobspec.RuntimeBlock{OneOf: "process", Command: "/bin/default"},
|
||||||
|
Tasks: []jobspec.TaskGroupTask{
|
||||||
|
{Name: "app"},
|
||||||
|
{Name: "sidecar", Command: "/bin/override"},
|
||||||
|
},
|
||||||
|
}
|
||||||
|
files, err := SystemdEmitter{}.Render(spec, &Node{})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Render: %v", err)
|
||||||
|
}
|
||||||
|
if len(files) != 3 {
|
||||||
|
t.Fatalf("got %d files, want 3", len(files))
|
||||||
|
}
|
||||||
|
appContent := findUnitContent(files, "/etc/systemd/system/orca-v1-alloc-web-app.service")
|
||||||
|
if appContent == "" {
|
||||||
|
t.Fatalf("missing app unit; paths %v", filePaths(files))
|
||||||
|
}
|
||||||
|
if !strings.Contains(appContent, "ExecStart=/bin/default") {
|
||||||
|
t.Errorf("app unit should inherit top-level command /bin/default\n%s", appContent)
|
||||||
|
}
|
||||||
|
sideContent := findUnitContent(files, "/etc/systemd/system/orca-v1-alloc-web-sidecar.service")
|
||||||
|
if sideContent == "" {
|
||||||
|
t.Fatalf("missing sidecar unit; paths %v", filePaths(files))
|
||||||
|
}
|
||||||
|
if !strings.Contains(sideContent, "ExecStart=/bin/override") {
|
||||||
|
t.Errorf("sidecar unit should use its own command /bin/override\n%s", sideContent)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestSystemdEmitter_TaskGroupNoCommandError(t *testing.T) {
|
||||||
|
// P06: a task with no resolvable command (no task.Command, no
|
||||||
|
// task.Runtime, no top-level Runtime) is an error.
|
||||||
|
spec := &jobspec.WorkloadSpec{
|
||||||
|
Kind: "Service",
|
||||||
|
Name: "web",
|
||||||
|
Tasks: []jobspec.TaskGroupTask{{Name: "app"}},
|
||||||
|
}
|
||||||
|
_, err := SystemdEmitter{}.Render(spec, &Node{})
|
||||||
|
if err == nil {
|
||||||
|
t.Fatal("expected error for task with no runtime, got nil")
|
||||||
|
}
|
||||||
|
if !strings.Contains(err.Error(), "no runtime") {
|
||||||
|
t.Errorf("error = %q, want 'no runtime'", err.Error())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestSystemdEmitter_TaskGroupEmptyCommandError(t *testing.T) {
|
||||||
|
// P06: a task whose resolved runtime command is empty/whitespace
|
||||||
|
// is an error (mirrors the single-process rule).
|
||||||
|
spec := &jobspec.WorkloadSpec{
|
||||||
|
Kind: "Service",
|
||||||
|
Name: "web",
|
||||||
|
Runtime: &jobspec.RuntimeBlock{OneOf: "process", Command: " "},
|
||||||
|
Tasks: []jobspec.TaskGroupTask{{Name: "app"}},
|
||||||
|
}
|
||||||
|
_, err := SystemdEmitter{}.Render(spec, &Node{})
|
||||||
|
if err == nil {
|
||||||
|
t.Fatal("expected error for empty command, got nil")
|
||||||
|
}
|
||||||
|
if !strings.Contains(err.Error(), "command is empty") {
|
||||||
|
t.Errorf("error = %q, want 'command is empty'", err.Error())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestSystemdEmitter_NoTasksBackwardCompat(t *testing.T) {
|
||||||
|
// Backward compat: a spec with no Tasks renders exactly one unit
|
||||||
|
// (the historical single-process shape).
|
||||||
|
spec := &jobspec.WorkloadSpec{
|
||||||
|
Kind: "Job",
|
||||||
|
Name: "backup",
|
||||||
|
Runtime: &jobspec.RuntimeBlock{OneOf: "process", Command: "/bin/rsync"},
|
||||||
|
}
|
||||||
|
files, err := SystemdEmitter{}.Render(spec, &Node{})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Render: %v", err)
|
||||||
|
}
|
||||||
|
if len(files) != 1 {
|
||||||
|
t.Fatalf("got %d files, want 1 (backward compat)", len(files))
|
||||||
|
}
|
||||||
|
if files[0].Path != "/etc/systemd/system/orca-v1-backup.service" {
|
||||||
|
t.Errorf("Path = %q, want /etc/systemd/system/orca-v1-backup.service", files[0].Path)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func filePaths(files []File) []string {
|
||||||
|
out := make([]string, len(files))
|
||||||
|
for i, f := range files {
|
||||||
|
out[i] = f.Path
|
||||||
|
}
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
|
||||||
|
func findUnitContent(files []File, path string) string {
|
||||||
|
for _, f := range files {
|
||||||
|
if f.Path == path {
|
||||||
|
return f.Content
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
|||||||
@@ -74,6 +74,27 @@ type WorkloadSpec struct {
|
|||||||
// Timeout is an optional execution timeout (duration string) for
|
// Timeout is an optional execution timeout (duration string) for
|
||||||
// Job. Populated by P04.
|
// Job. Populated by P04.
|
||||||
Timeout string
|
Timeout string
|
||||||
|
|
||||||
|
// Tasks is the task-group list for multi-process services (P06,
|
||||||
|
// PRD §9.1). When non-empty, the alloc runs one systemd unit per
|
||||||
|
// task (`orca-v1-alloc-<alloc-id>-<task-name>.service`) all
|
||||||
|
// grouped under a single `<alloc-id>.target`. When nil/empty,
|
||||||
|
// the alloc is a single-process alloc driven by the top-level
|
||||||
|
// Runtime block (backward compat). Tasks that omit their own
|
||||||
|
// runtime inherit the top-level Runtime as the per-group default.
|
||||||
|
Tasks []TaskGroupTask
|
||||||
|
}
|
||||||
|
|
||||||
|
// TaskGroupTask is a single task within a task group (P06, PRD §9.1).
|
||||||
|
// Each task has its own runtime (a wasm task + a process sidecar is
|
||||||
|
// allowed), its own command, and an optional env overlay. When
|
||||||
|
// Runtime is nil, the task inherits the top-level
|
||||||
|
// WorkloadSpec.Runtime (the per-group default).
|
||||||
|
type TaskGroupTask struct {
|
||||||
|
Name string
|
||||||
|
Runtime *RuntimeBlock
|
||||||
|
Env map[string]string
|
||||||
|
Command string
|
||||||
}
|
}
|
||||||
|
|
||||||
// RuntimeBlock is a minimal runtime abstraction surface populated by the
|
// RuntimeBlock is a minimal runtime abstraction surface populated by the
|
||||||
@@ -384,12 +405,18 @@ func parseFrontmatterBlock(block string) (*WorkloadSpec, error) {
|
|||||||
secLifecycle
|
secLifecycle
|
||||||
secAffinity
|
secAffinity
|
||||||
secConstraints
|
secConstraints
|
||||||
|
secTasks
|
||||||
|
secTaskEnv
|
||||||
|
secTaskRuntime
|
||||||
)
|
)
|
||||||
cur := secNone
|
cur := secNone
|
||||||
var curPort *PortSpec
|
var curPort *PortSpec
|
||||||
var curVol *VolumeSpec
|
var curVol *VolumeSpec
|
||||||
var curAffinity *AffinityRule
|
var curAffinity *AffinityRule
|
||||||
var lifecycleCur string
|
var lifecycleCur string
|
||||||
|
var curTask *TaskGroupTask
|
||||||
|
var taskIndent int
|
||||||
|
var taskFieldIndent int
|
||||||
|
|
||||||
flushPort := func() {
|
flushPort := func() {
|
||||||
if curPort != nil {
|
if curPort != nil {
|
||||||
@@ -409,6 +436,29 @@ func parseFrontmatterBlock(block string) (*WorkloadSpec, error) {
|
|||||||
curAffinity = nil
|
curAffinity = nil
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
flushTask := func() {
|
||||||
|
if curTask != nil {
|
||||||
|
spec.Tasks = append(spec.Tasks, *curTask)
|
||||||
|
curTask = nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
// taskSubBlock returns the sub-section to switch to when the
|
||||||
|
// given `key: value` line opens a nested block under a task
|
||||||
|
// (`env:` → secTaskEnv, `runtime:` → secTaskRuntime). Returns
|
||||||
|
// secTasks for non-block keys (no switch).
|
||||||
|
taskSubBlock := func(kvLine string) section {
|
||||||
|
key, _, ok := splitKV(kvLine)
|
||||||
|
if !ok {
|
||||||
|
return secTasks
|
||||||
|
}
|
||||||
|
switch key {
|
||||||
|
case "env":
|
||||||
|
return secTaskEnv
|
||||||
|
case "runtime":
|
||||||
|
return secTaskRuntime
|
||||||
|
}
|
||||||
|
return secTasks
|
||||||
|
}
|
||||||
|
|
||||||
for lineNo, raw := range lines {
|
for lineNo, raw := range lines {
|
||||||
line := stripComment(raw)
|
line := stripComment(raw)
|
||||||
@@ -423,6 +473,7 @@ func parseFrontmatterBlock(block string) (*WorkloadSpec, error) {
|
|||||||
flushPort()
|
flushPort()
|
||||||
flushVol()
|
flushVol()
|
||||||
flushAffinity()
|
flushAffinity()
|
||||||
|
flushTask()
|
||||||
cur = secNone
|
cur = secNone
|
||||||
|
|
||||||
key, val, ok := splitKV(trimmed)
|
key, val, ok := splitKV(trimmed)
|
||||||
@@ -500,6 +551,10 @@ func parseFrontmatterBlock(block string) (*WorkloadSpec, error) {
|
|||||||
} else {
|
} else {
|
||||||
cur = secAffinity
|
cur = secAffinity
|
||||||
}
|
}
|
||||||
|
case "tasks":
|
||||||
|
cur = secTasks
|
||||||
|
taskIndent = -1
|
||||||
|
taskFieldIndent = -1
|
||||||
default:
|
default:
|
||||||
// Unknown top-level key are ignored (forward-compat).
|
// Unknown top-level key are ignored (forward-compat).
|
||||||
cur = secNone
|
cur = secNone
|
||||||
@@ -720,11 +775,119 @@ func parseFrontmatterBlock(block string) (*WorkloadSpec, error) {
|
|||||||
spec.Constraints = append(spec.Constraints, unquote(item))
|
spec.Constraints = append(spec.Constraints, unquote(item))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
case secTasks:
|
||||||
|
// Tasks is a list of task objects. A `- ` at the list
|
||||||
|
// indent opens a new task; deeper-indented lines belong
|
||||||
|
// to the current task's fields (name, command) or
|
||||||
|
// nested sub-blocks (runtime, env).
|
||||||
|
if strings.HasPrefix(trimmed, "- ") || trimmed == "-" {
|
||||||
|
if taskIndent < 0 {
|
||||||
|
taskIndent = indent
|
||||||
|
taskFieldIndent = indent + 2
|
||||||
|
}
|
||||||
|
if indent == taskIndent {
|
||||||
|
flushTask()
|
||||||
|
t := TaskGroupTask{}
|
||||||
|
curTask = &t
|
||||||
|
rest := strings.TrimSpace(strings.TrimPrefix(trimmed, "-"))
|
||||||
|
if rest != "" {
|
||||||
|
if applyTaskKV(curTask, rest) {
|
||||||
|
cur = taskSubBlock(rest)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if curTask != nil {
|
||||||
|
if applyTaskKV(curTask, trimmed) {
|
||||||
|
cur = taskSubBlock(trimmed)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
case secTaskEnv:
|
||||||
|
if curTask == nil {
|
||||||
|
cur = secTasks
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
// Pop back to the task field level when the indent
|
||||||
|
// returns to taskFieldIndent (the next sibling
|
||||||
|
// field or a new `- ` list item). The line is then
|
||||||
|
// reprocessed as a task field.
|
||||||
|
if taskFieldIndent > 0 && indent <= taskFieldIndent {
|
||||||
|
cur = secTasks
|
||||||
|
if indent == taskIndent && (strings.HasPrefix(trimmed, "- ") || trimmed == "-") {
|
||||||
|
flushTask()
|
||||||
|
t := TaskGroupTask{}
|
||||||
|
curTask = &t
|
||||||
|
rest := strings.TrimSpace(strings.TrimPrefix(trimmed, "-"))
|
||||||
|
if rest != "" {
|
||||||
|
if applyTaskKV(curTask, rest) {
|
||||||
|
cur = taskSubBlock(rest)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if applyTaskKV(curTask, trimmed) {
|
||||||
|
cur = taskSubBlock(trimmed)
|
||||||
|
}
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if curTask.Env == nil {
|
||||||
|
curTask.Env = map[string]string{}
|
||||||
|
}
|
||||||
|
key, val, ok := splitKV(trimmed)
|
||||||
|
if !ok {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if val == "" {
|
||||||
|
curTask.Env[key] = ""
|
||||||
|
} else if strings.HasPrefix(val, "{") && strings.HasSuffix(val, "}") {
|
||||||
|
curTask.Env[key] = val
|
||||||
|
} else {
|
||||||
|
curTask.Env[key] = unquote(val)
|
||||||
|
}
|
||||||
|
case secTaskRuntime:
|
||||||
|
if curTask == nil || curTask.Runtime == nil {
|
||||||
|
cur = secTasks
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
// Pop back to the task field level (see secTaskEnv).
|
||||||
|
if taskFieldIndent > 0 && indent <= taskFieldIndent {
|
||||||
|
cur = secTasks
|
||||||
|
if indent == taskIndent && (strings.HasPrefix(trimmed, "- ") || trimmed == "-") {
|
||||||
|
flushTask()
|
||||||
|
t := TaskGroupTask{}
|
||||||
|
curTask = &t
|
||||||
|
rest := strings.TrimSpace(strings.TrimPrefix(trimmed, "-"))
|
||||||
|
if rest != "" {
|
||||||
|
if applyTaskKV(curTask, rest) {
|
||||||
|
cur = taskSubBlock(rest)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if applyTaskKV(curTask, trimmed) {
|
||||||
|
cur = taskSubBlock(trimmed)
|
||||||
|
}
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
key, val, ok := splitKV(trimmed)
|
||||||
|
if !ok {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
switch key {
|
||||||
|
case "one_of":
|
||||||
|
curTask.Runtime.OneOf = unquote(val)
|
||||||
|
case "image":
|
||||||
|
curTask.Runtime.Image = unquote(val)
|
||||||
|
case "command":
|
||||||
|
curTask.Runtime.Command = unquote(val)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
flushPort()
|
flushPort()
|
||||||
flushVol()
|
flushVol()
|
||||||
flushAffinity()
|
flushAffinity()
|
||||||
|
flushTask()
|
||||||
return spec, nil
|
return spec, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -791,6 +954,34 @@ func applyAffinityKV(r *AffinityRule, s string) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// applyTaskKV applies a `key: value` pair to the current TaskGroupTask.
|
||||||
|
// The returned bool reports whether the key opened a nested sub-block
|
||||||
|
// (`env` or `runtime`); when true the caller switches the parser
|
||||||
|
// section to the corresponding sub-block handler.
|
||||||
|
func applyTaskKV(t *TaskGroupTask, s string) (openedSubBlock bool) {
|
||||||
|
key, val, ok := splitKV(s)
|
||||||
|
if !ok {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
switch key {
|
||||||
|
case "name":
|
||||||
|
t.Name = unquote(val)
|
||||||
|
case "command":
|
||||||
|
t.Command = unquote(val)
|
||||||
|
case "env":
|
||||||
|
if t.Env == nil {
|
||||||
|
t.Env = map[string]string{}
|
||||||
|
}
|
||||||
|
return true
|
||||||
|
case "runtime":
|
||||||
|
if t.Runtime == nil {
|
||||||
|
t.Runtime = &RuntimeBlock{}
|
||||||
|
}
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
// appendLifecycleCmd appends a command to the named lifecycle hook list
|
// appendLifecycleCmd appends a command to the named lifecycle hook list
|
||||||
// (pre_stop or post_start) on the given LifecycleBlock.
|
// (pre_stop or post_start) on the given LifecycleBlock.
|
||||||
func appendLifecycleCmd(lb *LifecycleBlock, name, cmd string) {
|
func appendLifecycleCmd(lb *LifecycleBlock, name, cmd string) {
|
||||||
|
|||||||
@@ -701,3 +701,152 @@ func TestParseMarkdown_FullServiceSpec(t *testing.T) {
|
|||||||
t.Errorf("Body = %q, want %q (R-015)", spec.Body, "# body\n")
|
t.Errorf("Body = %q, want %q (R-015)", spec.Body, "# body\n")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestParseMarkdown_TasksBlock(t *testing.T) {
|
||||||
|
// P06: a task group with two tasks, each carrying its own runtime
|
||||||
|
// and command. The parser must populate spec.Tasks with two
|
||||||
|
// entries preserving name, runtime (one_of/image/command), and
|
||||||
|
// the task-level command.
|
||||||
|
input := "---\n" +
|
||||||
|
"kind: Service\n" +
|
||||||
|
"name: web\n" +
|
||||||
|
"tasks:\n" +
|
||||||
|
" - name: app\n" +
|
||||||
|
" runtime:\n" +
|
||||||
|
" one_of: process\n" +
|
||||||
|
" image: docker.io/nginx:latest\n" +
|
||||||
|
" command: /usr/bin/httpd -f\n" +
|
||||||
|
" command: /usr/bin/httpd -f\n" +
|
||||||
|
" - name: sidecar\n" +
|
||||||
|
" runtime:\n" +
|
||||||
|
" one_of: wasm\n" +
|
||||||
|
" command: /bin/wasm-runner sidecar.wasm\n" +
|
||||||
|
" command: /bin/wasm-runner sidecar.wasm\n" +
|
||||||
|
"---\nbody\n"
|
||||||
|
spec, err := ParseMarkdown([]byte(input))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("ParseMarkdown: %v", err)
|
||||||
|
}
|
||||||
|
if len(spec.Tasks) != 2 {
|
||||||
|
t.Fatalf("Tasks = %d, want 2", len(spec.Tasks))
|
||||||
|
}
|
||||||
|
app := spec.Tasks[0]
|
||||||
|
if app.Name != "app" {
|
||||||
|
t.Errorf("Tasks[0].Name = %q, want app", app.Name)
|
||||||
|
}
|
||||||
|
if app.Runtime == nil {
|
||||||
|
t.Fatal("Tasks[0].Runtime is nil")
|
||||||
|
}
|
||||||
|
if app.Runtime.OneOf != "process" {
|
||||||
|
t.Errorf("Tasks[0].Runtime.OneOf = %q, want process", app.Runtime.OneOf)
|
||||||
|
}
|
||||||
|
if app.Runtime.Image != "docker.io/nginx:latest" {
|
||||||
|
t.Errorf("Tasks[0].Runtime.Image = %q", app.Runtime.Image)
|
||||||
|
}
|
||||||
|
if app.Runtime.Command != "/usr/bin/httpd -f" {
|
||||||
|
t.Errorf("Tasks[0].Runtime.Command = %q", app.Runtime.Command)
|
||||||
|
}
|
||||||
|
if app.Command != "/usr/bin/httpd -f" {
|
||||||
|
t.Errorf("Tasks[0].Command = %q", app.Command)
|
||||||
|
}
|
||||||
|
side := spec.Tasks[1]
|
||||||
|
if side.Name != "sidecar" {
|
||||||
|
t.Errorf("Tasks[1].Name = %q, want sidecar", side.Name)
|
||||||
|
}
|
||||||
|
if side.Runtime == nil || side.Runtime.OneOf != "wasm" {
|
||||||
|
t.Errorf("Tasks[1].Runtime = %+v, want one_of=wasm", side.Runtime)
|
||||||
|
}
|
||||||
|
if side.Command != "/bin/wasm-runner sidecar.wasm" {
|
||||||
|
t.Errorf("Tasks[1].Command = %q", side.Command)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestParseMarkdown_TasksBlockWithEnv(t *testing.T) {
|
||||||
|
// P06: a task group task carrying an env overlay.
|
||||||
|
input := "---\n" +
|
||||||
|
"kind: Service\n" +
|
||||||
|
"name: web\n" +
|
||||||
|
"tasks:\n" +
|
||||||
|
" - name: app\n" +
|
||||||
|
" command: /usr/bin/httpd\n" +
|
||||||
|
" env:\n" +
|
||||||
|
" LOG_LEVEL: debug\n" +
|
||||||
|
" REGION: us\n" +
|
||||||
|
"---\nbody\n"
|
||||||
|
spec, err := ParseMarkdown([]byte(input))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("ParseMarkdown: %v", err)
|
||||||
|
}
|
||||||
|
if len(spec.Tasks) != 1 {
|
||||||
|
t.Fatalf("Tasks = %d, want 1", len(spec.Tasks))
|
||||||
|
}
|
||||||
|
task := spec.Tasks[0]
|
||||||
|
if task.Env == nil {
|
||||||
|
t.Fatal("Tasks[0].Env is nil")
|
||||||
|
}
|
||||||
|
if got := task.Env["LOG_LEVEL"]; got != "debug" {
|
||||||
|
t.Errorf("Env[LOG_LEVEL] = %q, want debug", got)
|
||||||
|
}
|
||||||
|
if got := task.Env["REGION"]; got != "us" {
|
||||||
|
t.Errorf("Env[REGION] = %q, want us", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestParseMarkdown_TasksBlockInheritsTopLevelRuntime(t *testing.T) {
|
||||||
|
// P06: when a task omits its own runtime, the top-level runtime
|
||||||
|
// is the per-group default. The parser must NOT create a task
|
||||||
|
// runtime when the task block lacks a `runtime:` sub-block; the
|
||||||
|
// emitter/validator resolve the default from spec.Runtime.
|
||||||
|
input := "---\n" +
|
||||||
|
"kind: Service\n" +
|
||||||
|
"name: web\n" +
|
||||||
|
"runtime:\n" +
|
||||||
|
" one_of: process\n" +
|
||||||
|
" command: /bin/default\n" +
|
||||||
|
"tasks:\n" +
|
||||||
|
" - name: app\n" +
|
||||||
|
" command: /bin/app\n" +
|
||||||
|
" - name: sidecar\n" +
|
||||||
|
" command: /bin/sidecar\n" +
|
||||||
|
"---\nbody\n"
|
||||||
|
spec, err := ParseMarkdown([]byte(input))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("ParseMarkdown: %v", err)
|
||||||
|
}
|
||||||
|
if spec.Runtime == nil || spec.Runtime.OneOf != "process" {
|
||||||
|
t.Fatalf("top-level runtime not parsed: %+v", spec.Runtime)
|
||||||
|
}
|
||||||
|
if len(spec.Tasks) != 2 {
|
||||||
|
t.Fatalf("Tasks = %d, want 2", len(spec.Tasks))
|
||||||
|
}
|
||||||
|
for i, task := range spec.Tasks {
|
||||||
|
if task.Runtime != nil {
|
||||||
|
t.Errorf("Tasks[%d].Runtime should be nil (inherit top-level), got %+v", i, task.Runtime)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if spec.Tasks[0].Name != "app" || spec.Tasks[1].Name != "sidecar" {
|
||||||
|
t.Errorf("task names = %q, %q", spec.Tasks[0].Name, spec.Tasks[1].Name)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestParseMarkdown_NoTasksBackwardCompat(t *testing.T) {
|
||||||
|
// Backward compat: a spec with no `tasks:` block parses as a
|
||||||
|
// single-process alloc; spec.Tasks must be empty/nil.
|
||||||
|
input := "---\n" +
|
||||||
|
"kind: Job\n" +
|
||||||
|
"name: backup\n" +
|
||||||
|
"runtime:\n" +
|
||||||
|
" one_of: process\n" +
|
||||||
|
" command: /bin/rsync\n" +
|
||||||
|
"---\nbody\n"
|
||||||
|
spec, err := ParseMarkdown([]byte(input))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("ParseMarkdown: %v", err)
|
||||||
|
}
|
||||||
|
if len(spec.Tasks) != 0 {
|
||||||
|
t.Fatalf("Tasks = %d, want 0 (backward compat)", len(spec.Tasks))
|
||||||
|
}
|
||||||
|
if spec.Runtime == nil || spec.Runtime.Command != "/bin/rsync" {
|
||||||
|
t.Errorf("Runtime = %+v, want command=/bin/rsync", spec.Runtime)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -37,6 +37,8 @@ type Validator interface {
|
|||||||
// - count must be 1 (or unset → 1); count > 1 is an error for Job
|
// - count must be 1 (or unset → 1); count > 1 is an error for Job
|
||||||
// (use a Service for replicas)
|
// (use a Service for replicas)
|
||||||
// - no Traefik route (a ServiceBlock is rejected)
|
// - no Traefik route (a ServiceBlock is rejected)
|
||||||
|
// - task group (spec.Tasks) optional; when present, each task must
|
||||||
|
// have a unique name and a resolvable command (P06).
|
||||||
type JobValidator struct{}
|
type JobValidator struct{}
|
||||||
|
|
||||||
// ServiceValidator validates the Service workload kind (R-012).
|
// ServiceValidator validates the Service workload kind (R-012).
|
||||||
@@ -46,11 +48,14 @@ type JobValidator struct{}
|
|||||||
// - count ≥ 1
|
// - count ≥ 1
|
||||||
// - restart required (mode must be service)
|
// - restart required (mode must be service)
|
||||||
// - update required (strategy must be rolling/canary/blue-green)
|
// - update required (strategy must be rolling/canary/blue-green)
|
||||||
// - runtime required
|
// - runtime required (unless a task group is present; each task
|
||||||
|
// can carry its own runtime — P06)
|
||||||
// - health block required (Traefik routing depends on health checks)
|
// - health block required (Traefik routing depends on health checks)
|
||||||
// - service block, if present, must have a valid bind (127.0.0.1
|
// - service block, if present, must have a valid bind (127.0.0.1
|
||||||
// opt-in per R-007; default is socket — empty bind is OK)
|
// opt-in per R-007; default is socket — empty bind is OK)
|
||||||
// - service block implied (Traefik route YES)
|
// - service block implied (Traefik route YES)
|
||||||
|
// - task group (spec.Tasks) optional; when present, each task must
|
||||||
|
// have a unique name and a resolvable command (P06).
|
||||||
type ServiceValidator struct{}
|
type ServiceValidator struct{}
|
||||||
|
|
||||||
// DaemonSetValidator validates the DaemonSet workload kind (R-012).
|
// DaemonSetValidator validates the DaemonSet workload kind (R-012).
|
||||||
@@ -60,6 +65,8 @@ type ServiceValidator struct{}
|
|||||||
// - no ports (no Traefik route by default D-175)
|
// - no ports (no Traefik route by default D-175)
|
||||||
// - no count (implicit = nodes matching condition)
|
// - no count (implicit = nodes matching condition)
|
||||||
// - restart required
|
// - restart required
|
||||||
|
// - task group (spec.Tasks) optional; when present, each task must
|
||||||
|
// have a unique name and a resolvable command (P06).
|
||||||
type DaemonSetValidator struct{}
|
type DaemonSetValidator struct{}
|
||||||
|
|
||||||
// ValidatorFor returns the Validator for the given workload kind, or an
|
// ValidatorFor returns the Validator for the given workload kind, or an
|
||||||
@@ -93,6 +100,7 @@ func (JobValidator) Validate(spec *jobspec.WorkloadSpec) error {
|
|||||||
if spec.Service != nil {
|
if spec.Service != nil {
|
||||||
errs = append(errs, "service block (Traefik route) is not allowed for Job (D-175)")
|
errs = append(errs, "service block (Traefik route) is not allowed for Job (D-175)")
|
||||||
}
|
}
|
||||||
|
errs = append(errs, validateTaskGroup(spec)...)
|
||||||
return composeErrors("schema/Job", errs)
|
return composeErrors("schema/Job", errs)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -137,8 +145,8 @@ func (ServiceValidator) Validate(spec *jobspec.WorkloadSpec) error {
|
|||||||
errs = append(errs, fmt.Sprintf("update strategy %q invalid (want one of rolling, canary, blue-green)", spec.Update.Strategy))
|
errs = append(errs, fmt.Sprintf("update strategy %q invalid (want one of rolling, canary, blue-green)", spec.Update.Strategy))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if spec.Runtime == nil {
|
if spec.Runtime == nil && len(spec.Tasks) == 0 {
|
||||||
errs = append(errs, "runtime block required for Service")
|
errs = append(errs, "runtime block required for Service (or a task group with per-task runtimes)")
|
||||||
}
|
}
|
||||||
if spec.Health == nil {
|
if spec.Health == nil {
|
||||||
errs = append(errs, "health block required for Service (Traefik routing requires health checks)")
|
errs = append(errs, "health block required for Service (Traefik routing requires health checks)")
|
||||||
@@ -148,6 +156,7 @@ func (ServiceValidator) Validate(spec *jobspec.WorkloadSpec) error {
|
|||||||
errs = append(errs, err.Error())
|
errs = append(errs, err.Error())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
errs = append(errs, validateTaskGroup(spec)...)
|
||||||
return composeErrors("schema/Service", errs)
|
return composeErrors("schema/Service", errs)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -195,6 +204,7 @@ func (DaemonSetValidator) Validate(spec *jobspec.WorkloadSpec) error {
|
|||||||
if spec.Restart == nil {
|
if spec.Restart == nil {
|
||||||
errs = append(errs, "restart block required for DaemonSet")
|
errs = append(errs, "restart block required for DaemonSet")
|
||||||
}
|
}
|
||||||
|
errs = append(errs, validateTaskGroup(spec)...)
|
||||||
return composeErrors("schema/DaemonSet", errs)
|
return composeErrors("schema/DaemonSet", errs)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -207,3 +217,38 @@ func composeErrors(name string, errs []string) error {
|
|||||||
}
|
}
|
||||||
return fmt.Errorf("%s: %s", name, strings.Join(errs, "; "))
|
return fmt.Errorf("%s: %s", name, strings.Join(errs, "; "))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// validateTaskGroup validates the task-group list shared by all kinds
|
||||||
|
// (P06, PRD §9.1). When the spec carries a task group (spec.Tasks
|
||||||
|
// non-empty), each task must have a unique name and a resolvable
|
||||||
|
// command (the task's own Command, the task's runtime command, or the
|
||||||
|
// top-level runtime command as the per-group default). The top-level
|
||||||
|
// runtime is optional when tasks is present (each task can carry its
|
||||||
|
// own runtime). Returns nil when the spec has no task group.
|
||||||
|
func validateTaskGroup(spec *jobspec.WorkloadSpec) []string {
|
||||||
|
if len(spec.Tasks) == 0 {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
var errs []string
|
||||||
|
seen := make(map[string]bool, len(spec.Tasks))
|
||||||
|
for i, task := range spec.Tasks {
|
||||||
|
if strings.TrimSpace(task.Name) == "" {
|
||||||
|
errs = append(errs, fmt.Sprintf("tasks[%d]: name is required", i))
|
||||||
|
} else if seen[task.Name] {
|
||||||
|
errs = append(errs, fmt.Sprintf("tasks[%d]: duplicate task name %q (names must be unique within the group)", i, task.Name))
|
||||||
|
} else {
|
||||||
|
seen[task.Name] = true
|
||||||
|
}
|
||||||
|
cmd := task.Command
|
||||||
|
if strings.TrimSpace(cmd) == "" && task.Runtime != nil {
|
||||||
|
cmd = task.Runtime.Command
|
||||||
|
}
|
||||||
|
if strings.TrimSpace(cmd) == "" && spec.Runtime != nil {
|
||||||
|
cmd = spec.Runtime.Command
|
||||||
|
}
|
||||||
|
if strings.TrimSpace(cmd) == "" {
|
||||||
|
errs = append(errs, fmt.Sprintf("tasks[%d]: command is required (set tasks[].command, tasks[].runtime.command, or top-level runtime.command)", i))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return errs
|
||||||
|
}
|
||||||
|
|||||||
@@ -689,3 +689,183 @@ var (
|
|||||||
_ Validator = ServiceValidator{}
|
_ Validator = ServiceValidator{}
|
||||||
_ Validator = DaemonSetValidator{}
|
_ Validator = DaemonSetValidator{}
|
||||||
)
|
)
|
||||||
|
|
||||||
|
func TestTaskGroup_Valid(t *testing.T) {
|
||||||
|
// P06: a valid task group — two tasks, each with a unique name
|
||||||
|
// and a resolvable command (own command). The top-level runtime
|
||||||
|
// is optional when each task carries its own.
|
||||||
|
spec := &jobspec.WorkloadSpec{
|
||||||
|
Kind: "Service",
|
||||||
|
Name: "web",
|
||||||
|
Count: 1,
|
||||||
|
Tasks: []jobspec.TaskGroupTask{
|
||||||
|
{Name: "app", Command: "/usr/bin/httpd"},
|
||||||
|
{Name: "sidecar", Command: "/bin/wasm-runner sidecar.wasm"},
|
||||||
|
},
|
||||||
|
Restart: &jobspec.RestartBlock{Mode: "service"},
|
||||||
|
Update: &jobspec.UpdateBlock{Strategy: "rolling"},
|
||||||
|
Health: &jobspec.HealthBlock{CheckType: "http"},
|
||||||
|
Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}},
|
||||||
|
}
|
||||||
|
if err := (ServiceValidator{}).Validate(spec); err != nil {
|
||||||
|
t.Fatalf("expected nil, got %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestTaskGroup_ValidInheritsTopLevelRuntime(t *testing.T) {
|
||||||
|
// P06: tasks without their own runtime inherit the top-level
|
||||||
|
// runtime command. The validator accepts this as long as the
|
||||||
|
// resolved command is non-empty.
|
||||||
|
spec := &jobspec.WorkloadSpec{
|
||||||
|
Kind: "Service",
|
||||||
|
Name: "web",
|
||||||
|
Count: 1,
|
||||||
|
Runtime: &jobspec.RuntimeBlock{OneOf: "process", Command: "/bin/default"},
|
||||||
|
Tasks: []jobspec.TaskGroupTask{
|
||||||
|
{Name: "app"},
|
||||||
|
{Name: "sidecar"},
|
||||||
|
},
|
||||||
|
Restart: &jobspec.RestartBlock{Mode: "service"},
|
||||||
|
Update: &jobspec.UpdateBlock{Strategy: "rolling"},
|
||||||
|
Health: &jobspec.HealthBlock{CheckType: "http"},
|
||||||
|
Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}},
|
||||||
|
}
|
||||||
|
if err := (ServiceValidator{}).Validate(spec); err != nil {
|
||||||
|
t.Fatalf("expected nil, got %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestTaskGroup_ValidTaskRuntimeCommand(t *testing.T) {
|
||||||
|
// P06: a task whose command is provided via the task's own
|
||||||
|
// runtime.command (no top-level runtime) is valid.
|
||||||
|
spec := &jobspec.WorkloadSpec{
|
||||||
|
Kind: "Service",
|
||||||
|
Name: "web",
|
||||||
|
Count: 1,
|
||||||
|
Tasks: []jobspec.TaskGroupTask{
|
||||||
|
{Name: "app", Runtime: &jobspec.RuntimeBlock{OneOf: "process", Command: "/usr/bin/httpd"}},
|
||||||
|
},
|
||||||
|
Restart: &jobspec.RestartBlock{Mode: "service"},
|
||||||
|
Update: &jobspec.UpdateBlock{Strategy: "rolling"},
|
||||||
|
Health: &jobspec.HealthBlock{CheckType: "http"},
|
||||||
|
Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}},
|
||||||
|
}
|
||||||
|
if err := (ServiceValidator{}).Validate(spec); err != nil {
|
||||||
|
t.Fatalf("expected nil, got %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestTaskGroup_MissingTaskName(t *testing.T) {
|
||||||
|
spec := &jobspec.WorkloadSpec{
|
||||||
|
Kind: "Service",
|
||||||
|
Name: "web",
|
||||||
|
Tasks: []jobspec.TaskGroupTask{
|
||||||
|
{Command: "/usr/bin/httpd"},
|
||||||
|
{Name: "sidecar", Command: "/bin/wasm-runner"},
|
||||||
|
},
|
||||||
|
Restart: &jobspec.RestartBlock{Mode: "service"},
|
||||||
|
Update: &jobspec.UpdateBlock{Strategy: "rolling"},
|
||||||
|
Health: &jobspec.HealthBlock{CheckType: "http"},
|
||||||
|
Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}},
|
||||||
|
}
|
||||||
|
err := ServiceValidator{}.Validate(spec)
|
||||||
|
if err == nil {
|
||||||
|
t.Fatal("expected error for missing task name, got nil")
|
||||||
|
}
|
||||||
|
if !strings.Contains(err.Error(), "name is required") {
|
||||||
|
t.Errorf("error = %q, want 'name is required'", err.Error())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestTaskGroup_DuplicateTaskNames(t *testing.T) {
|
||||||
|
spec := &jobspec.WorkloadSpec{
|
||||||
|
Kind: "Service",
|
||||||
|
Name: "web",
|
||||||
|
Tasks: []jobspec.TaskGroupTask{
|
||||||
|
{Name: "app", Command: "/usr/bin/httpd"},
|
||||||
|
{Name: "app", Command: "/bin/other"},
|
||||||
|
},
|
||||||
|
Restart: &jobspec.RestartBlock{Mode: "service"},
|
||||||
|
Update: &jobspec.UpdateBlock{Strategy: "rolling"},
|
||||||
|
Health: &jobspec.HealthBlock{CheckType: "http"},
|
||||||
|
Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}},
|
||||||
|
}
|
||||||
|
err := ServiceValidator{}.Validate(spec)
|
||||||
|
if err == nil {
|
||||||
|
t.Fatal("expected error for duplicate task names, got nil")
|
||||||
|
}
|
||||||
|
if !strings.Contains(err.Error(), "duplicate task name") {
|
||||||
|
t.Errorf("error = %q, want 'duplicate task name'", err.Error())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestTaskGroup_MissingCommand(t *testing.T) {
|
||||||
|
// P06: a task with no resolvable command (no task.Command, no
|
||||||
|
// task.Runtime, no top-level Runtime) is rejected.
|
||||||
|
spec := &jobspec.WorkloadSpec{
|
||||||
|
Kind: "Service",
|
||||||
|
Name: "web",
|
||||||
|
Tasks: []jobspec.TaskGroupTask{
|
||||||
|
{Name: "app"},
|
||||||
|
},
|
||||||
|
Restart: &jobspec.RestartBlock{Mode: "service"},
|
||||||
|
Update: &jobspec.UpdateBlock{Strategy: "rolling"},
|
||||||
|
Health: &jobspec.HealthBlock{CheckType: "http"},
|
||||||
|
Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}},
|
||||||
|
}
|
||||||
|
err := ServiceValidator{}.Validate(spec)
|
||||||
|
if err == nil {
|
||||||
|
t.Fatal("expected error for missing task command, got nil")
|
||||||
|
}
|
||||||
|
if !strings.Contains(err.Error(), "command is required") {
|
||||||
|
t.Errorf("error = %q, want 'command is required'", err.Error())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestTaskGroup_JobAcceptsTaskGroup(t *testing.T) {
|
||||||
|
// P06: task groups apply to all kinds, not just Service. Job
|
||||||
|
// accepts a task group with unique names + resolvable commands.
|
||||||
|
spec := &jobspec.WorkloadSpec{
|
||||||
|
Kind: "Job",
|
||||||
|
Name: "batch",
|
||||||
|
Count: 1,
|
||||||
|
Tasks: []jobspec.TaskGroupTask{
|
||||||
|
{Name: "step1", Command: "/bin/extract"},
|
||||||
|
{Name: "step2", Command: "/bin/transform"},
|
||||||
|
},
|
||||||
|
}
|
||||||
|
if err := (JobValidator{}).Validate(spec); err != nil {
|
||||||
|
t.Fatalf("expected nil, got %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestTaskGroup_DaemonSetAcceptsTaskGroup(t *testing.T) {
|
||||||
|
// P06: DaemonSet accepts a task group.
|
||||||
|
spec := &jobspec.WorkloadSpec{
|
||||||
|
Kind: "DaemonSet",
|
||||||
|
Name: "log-shipper",
|
||||||
|
Schedule: &jobspec.ScheduleBlock{Mode: "every-node"},
|
||||||
|
Restart: &jobspec.RestartBlock{Mode: "on-failure"},
|
||||||
|
Tasks: []jobspec.TaskGroupTask{
|
||||||
|
{Name: "collector", Command: "/bin/collect"},
|
||||||
|
{Name: "forwarder", Command: "/bin/forward"},
|
||||||
|
},
|
||||||
|
}
|
||||||
|
if err := (DaemonSetValidator{}).Validate(spec); err != nil {
|
||||||
|
t.Fatalf("expected nil, got %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestTaskGroup_NoTasksBackwardCompat(t *testing.T) {
|
||||||
|
// Backward compat: a spec with no Tasks is validated by the
|
||||||
|
// existing kind-specific rules (no task-group check fires).
|
||||||
|
spec := &jobspec.WorkloadSpec{
|
||||||
|
Kind: "Job",
|
||||||
|
Name: "backup",
|
||||||
|
Count: 1,
|
||||||
|
Runtime: &jobspec.RuntimeBlock{OneOf: "process", Command: "/bin/rsync"},
|
||||||
|
}
|
||||||
|
if err := (JobValidator{}).Validate(spec); err != nil {
|
||||||
|
t.Fatalf("expected nil, got %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user