From 436641782c50efc5853e55ead220b82de4643a0e Mon Sep 17 00:00:00 2001 From: Jon Chery Date: Wed, 5 Aug 2026 17:48:04 +0000 Subject: [PATCH] feat(P02): Service block + Traefik emitter + atomic reload (REQ-077, gate C-10) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit P02 — Traefik dynamic config generation + atomic reload protocol. Parser (internal/jobspec/markdown.go): - Extended WorkloadSpec with Health, Constraints, Affinity, Lifecycle fields. Parsed restart/update/service/health/lifecycle/affinity/ constraints blocks. HealthBlock, AffinityRule, LifecycleBlock types. Schema (internal/spec/schema/schema.go): - ServiceValidator: restart.mode enum (service/on-failure/never), update.strategy enum (rolling/canary/blue-green), health required, service.bind IP validation (R-007 loopback opt-in). 98.5% coverage. Traefik emitter (internal/emitter/traefik.go, REQ-077): - TraefikEmitter renders /etc/traefik/dynamic/orca-.yaml with http.routers, http.services (servers = R-007 socket paths), TLS (certResolver=orca, trust domain), healthCheck. RenderDrain sets weight:0 per backend. RegisterTraefik wires process/podman/wasm. Atomic reload (internal/emitter/traefik_atomic.go, gate C-10): - WriteTraefikDynamic: write to path.tmp via WriteFileIdempotent, then mv -f path.tmp path (atomic POSIX rename, Traefik fsnotify observes IN_MOVED_TO). Traefik holds-last-good on malformed config. C-10 PASS. 22 packages pass, 20 bats pass, gofmt clean, verify-reqs 90 consistent. Coverage: emitter 96.5%, jobspec 88.8%, schema 98.5%, sshpush 93.0%. ---ci--- project: orca phase: P02 milestone: v0.9 status: execute ---/ci--- --- internal/emitter/traefik.go | 225 +++++++++++++++ internal/emitter/traefik_atomic.go | 93 +++++++ internal/emitter/traefik_atomic_test.go | 190 +++++++++++++ internal/emitter/traefik_test.go | 353 ++++++++++++++++++++++++ internal/jobspec/markdown.go | 326 +++++++++++++++++++++- internal/jobspec/markdown_test.go | 346 +++++++++++++++++++++++ internal/spec/schema/schema.go | 51 +++- internal/spec/schema/schema_test.go | 263 +++++++++++++++++- internal/sshpush/atomic_writer_test.go | 30 ++ 9 files changed, 1856 insertions(+), 21 deletions(-) create mode 100644 internal/emitter/traefik.go create mode 100644 internal/emitter/traefik_atomic.go create mode 100644 internal/emitter/traefik_atomic_test.go create mode 100644 internal/emitter/traefik_test.go create mode 100644 internal/sshpush/atomic_writer_test.go diff --git a/internal/emitter/traefik.go b/internal/emitter/traefik.go new file mode 100644 index 0000000..0ed8f84 --- /dev/null +++ b/internal/emitter/traefik.go @@ -0,0 +1,225 @@ +package emitter + +import ( + "errors" + "fmt" + "net" + "strings" + + "git.cloudinit.dev/coreci/orca/internal/jobspec" +) + +// TraefikEmitter is the Layer-4 emitter for the Traefik dynamic-config +// file (REQ-077). It renders /etc/traefik/dynamic/orca-.yaml +// — a single Traefik dynamic-config file describing the routers, +// services (servers = the R-007 socket paths), TLS config pointing at +// the step-ca root CA, and the service health check. +// +// Registered on the emitter.Registry under the service-kind keys: +// +// - service:process +// - service:podman +// - service:wasm +// +// RegisterTraefik wires all three; callers can also call Register +// directly with TraefikEmitter{} for a single runtime. +// +// Atomic reload (gate C-10): the Traefik dynamic-config file is written +// atomically via the SSH-push transport (sshpush.WriteFileIdempotent +// performs temp-file + fsync + rename, and WriteTraefikDynamic wraps +// it with an explicit tmp+mv so fsnotify sees a single rename event). +// Traefik watches the dynamic dir with fsnotify; the rename triggers a +// reload. On a malformed config Traefik logs an error and holds the +// last-good config (documented Traefik behavior; the C-10 test +// verifies the tmp+rename sequence so a half-written file is never +// observed by Traefik). Drain is rendered by setting the backend +// server's weight to 0 (or removing it) — see RenderDrain. +// +// The orca-v1- prefix is NOT applied to Traefik dynamic-config paths +// (the prefix is only for systemd unit names; the Traefik file is named +// orca-.yaml and is the single source of truth for the +// service route — there is no dual-write window for Traefik configs). +type TraefikEmitter struct{} + +// traefikDynamicDir is the canonical Traefik dynamic-config directory +// (R-006). The emitter writes one file per service at +// /etc/traefik/dynamic/orca-.yaml. +const traefikDynamicDir = "/etc/traefik/dynamic" + +// traefikRouterTLSCertResolver is the Traefik cert-resolver name that +// the orca step-ca integration configures on the Traefik static config +// (P10 / v0.10 wires the step-ca root into this resolver). The +// dynamic-config file references it by name. +const traefikRouterTLSCertResolver = "orca" + +// defaultTrustDomain is the SPIFFE trust domain used in the rendered +// TLS stanza when the spec does not carry an explicit trust domain. +// The step-ca provisioner (P10) overrides this at render time via the +// node argument; for P02 the emitter renders the placeholder. +const defaultTrustDomain = "cluster.orca.local" + +// Render renders the Traefik dynamic-config YAML for a Service +// workload. The output is a single File whose Path is +// /etc/traefik/dynamic/orca-.yaml, Content is the rendered +// YAML, and Mode is 0644. +// +// Returns an error if the spec is nil, the name is empty, the spec has +// no ports (a Service with no ports has no backends to route to), or a +// service.bind value (when present) is not a valid IP address (R-007). +func (TraefikEmitter) Render(spec *jobspec.WorkloadSpec, node *Node) ([]File, error) { + if spec == nil { + return nil, errors.New("emitter/traefik: spec is nil") + } + if strings.TrimSpace(spec.Name) == "" { + return nil, errors.New("emitter/traefik: spec name is empty") + } + if len(spec.Ports) == 0 { + return nil, errors.New("emitter/traefik: service has no ports (no backends to route to)") + } + if spec.Service != nil { + if b := strings.TrimSpace(spec.Service.Bind); b != "" && net.ParseIP(b) == nil { + return nil, fmt.Errorf("emitter/traefik: service.bind %q is not a valid IP (R-007)", b) + } + } + content, err := renderTraefikYAML(spec, node) + if err != nil { + return nil, err + } + path := fmt.Sprintf("%s/orca-%s.yaml", traefikDynamicDir, spec.Name) + return []File{{Path: path, Content: content, Mode: "0644"}}, nil +} + +// RenderDrain renders a Traefik dynamic-config that drains the service +// by setting every backend server's weight to 0 (I-B-005 drain). The +// path matches the live config so the atomic rename overwrites the +// routing config with the drained config (Traefik reloads and stops +// sending traffic). The caller writes the result via +// WriteTraefikDynamic for the C-10 atomicity protocol. +func (e TraefikEmitter) RenderDrain(spec *jobspec.WorkloadSpec, node *Node) ([]File, error) { + if spec == nil { + return nil, errors.New("emitter/traefik: spec is nil") + } + if strings.TrimSpace(spec.Name) == "" { + return nil, errors.New("emitter/traefik: spec name is empty") + } + if len(spec.Ports) == 0 { + return nil, errors.New("emitter/traefik: service has no ports (no backends to drain)") + } + content, err := renderTraefikYAMLDrain(spec, node) + if err != nil { + return nil, err + } + path := fmt.Sprintf("%s/orca-%s.yaml", traefikDynamicDir, spec.Name) + return []File{{Path: path, Content: content, Mode: "0644"}}, nil +} + +// RegisterTraefik registers the TraefikEmitter on the given Registry +// under the three service-kind runtime keys (service:process, +// service:podman, service:wasm). The emitter is the same instance for +// all three runtimes — the rendered Traefik config is runtime-agnostic +// (the backend server URL is the R-007 socket path, which the runtime +// layer binds regardless of process/wasm/podman). +func RegisterTraefik(reg *Registry) { + e := TraefikEmitter{} + reg.Register("service:process", e) + reg.Register("service:podman", e) + reg.Register("service:wasm", e) +} + +// renderTraefikYAML renders the Traefik dynamic-config YAML for the +// given spec + node. The shape (verified by the Traefik docs) is: +// +// http: +// routers: +// orca-: +// rule: PathPrefix("/") +// service: orca- +// tls: +// certResolver: orca +// domains: +// - main: "" +// services: +// orca-: +// loadBalancer: +// servers: +// - url: "unix:///run/orca/alloc-/port-.sock" +// healthCheck: +// path: /healthz +// interval: +// timeout: +// +// The alloc-id placeholder is "" pending the P08 socket +// layer; Traefik will reject the URL until a real alloc-id is +// substituted. For P02 the emitter renders the placeholder so the +// C-10 atomicity protocol is testable end-to-end; the socket layer +// (P08) replaces the placeholder with the live alloc-id. +func renderTraefikYAML(spec *jobspec.WorkloadSpec, node *Node) (string, error) { + return renderTraefikYAMLWeighted(spec, node, false) +} + +// renderTraefikYAMLDrain renders the drained Traefik dynamic-config +// (every backend server has weight: 0). The shape mirrors the live +// config so the rename overwrites the live route with the drain. +func renderTraefikYAMLDrain(spec *jobspec.WorkloadSpec, node *Node) (string, error) { + return renderTraefikYAMLWeighted(spec, node, true) +} + +// renderTraefikYAMLWeighted renders the Traefik dynamic-config YAML. +// When drain is true, every server entry is emitted with `weight: 0` +// (I-B-005). When drain is false, no weight is emitted (Traefik +// defaults to 1 — equal weighting across servers). +func renderTraefikYAMLWeighted(spec *jobspec.WorkloadSpec, node *Node, drain bool) (string, error) { + var b strings.Builder + routerName := "orca-" + spec.Name + serviceName := "orca-" + spec.Name + rule := fmt.Sprintf("PathPrefix(\"/%s\")", spec.Name) + trustDomain := defaultTrustDomain + + b.WriteString("http:\n") + b.WriteString(" routers:\n") + b.WriteString(fmt.Sprintf(" %s:\n", routerName)) + b.WriteString(fmt.Sprintf(" rule: %s\n", rule)) + b.WriteString(fmt.Sprintf(" service: %s\n", serviceName)) + b.WriteString(" tls:\n") + b.WriteString(fmt.Sprintf(" certResolver: %s\n", traefikRouterTLSCertResolver)) + b.WriteString(" domains:\n") + b.WriteString(fmt.Sprintf(" - main: %q\n", trustDomain)) + b.WriteString(" services:\n") + b.WriteString(fmt.Sprintf(" %s:\n", serviceName)) + b.WriteString(" loadBalancer:\n") + b.WriteString(" servers:\n") + allocID := allocIDFor(node) + for _, p := range spec.Ports { + sock := fmt.Sprintf("unix:///run/orca/alloc-%s/port-%s.sock", allocID, p.Name) + b.WriteString(" - url: ") + b.WriteString(fmt.Sprintf("%q\n", sock)) + if drain { + b.WriteString(" weight: 0\n") + } + } + if spec.Health != nil { + b.WriteString(" healthCheck:\n") + path := "/healthz" + b.WriteString(fmt.Sprintf(" path: %s\n", path)) + if spec.Health.Interval != "" { + b.WriteString(fmt.Sprintf(" interval: %s\n", spec.Health.Interval)) + } + if spec.Health.Timeout != "" { + b.WriteString(fmt.Sprintf(" timeout: %s\n", spec.Health.Timeout)) + } + } + return b.String(), nil +} + +// allocIDFor returns the alloc-id placeholder for the node. P08 will +// substitute the live alloc-id from the socket layer; for P02 we use a +// deterministic placeholder derived from the node hostname so the +// rendered config is stable across re-renders (the C-10 idempotency +// check depends on a stable hash). When the node is nil or has no +// hostname, the literal placeholder "" is emitted. +func allocIDFor(node *Node) string { + if node == nil || strings.TrimSpace(node.Hostname) == "" { + return "" + } + return node.Hostname +} diff --git a/internal/emitter/traefik_atomic.go b/internal/emitter/traefik_atomic.go new file mode 100644 index 0000000..6e63c10 --- /dev/null +++ b/internal/emitter/traefik_atomic.go @@ -0,0 +1,93 @@ +package emitter + +import ( + "context" + "fmt" + "os" +) + +// AtomicWriter is the SSH-push transport surface that +// WriteTraefikDynamic uses to write the Traefik dynamic-config file +// atomically. It is the subset of *sshpush.Transport that the +// atomicity protocol depends on. Tests substitute a mock to assert +// the tmp+rename sequence (gate C-10) without a real SSH server. +// +// *sshpush.Transport satisfies this interface (the compile-time +// assertion lives in internal/sshpush to avoid an import cycle — the +// sshpush package imports emitter for fan-out, so this package cannot +// import sshpush). +type AtomicWriter interface { + // WriteFileIdempotent writes content to peer:path atomically with + // mode, returning written=true if the file was actually written + // (content hash differed). Used by WriteTraefikDynamic to write + // the .tmp sibling. + WriteFileIdempotent(ctx context.Context, peer string, path string, content []byte, mode os.FileMode) (bool, error) + // Exec runs a command on peer and returns its combined output. + // Used by WriteTraefikDynamic to perform the atomic `mv -f + // path.tmp path`. + Exec(ctx context.Context, peer string, cmd string) ([]byte, error) +} + +// WriteTraefikDynamic writes a Traefik dynamic-config file atomically +// (gate C-10: tmpfile + fsync + rename). The protocol is: +// +// 1. Write content to .tmp via WriteFileIdempotent. The +// underlying sshpush transport writes the tmp file in the same +// directory as the target with mode-appended naming, fsyncs, and +// renames — but we add an extra hop here so the *Traefik* file is +// only ever observed at its final path after a single atomic +// rename event that Traefik's fsnotify watcher sees. +// 2. `mv -f .tmp ` on the peer (atomic rename on POSIX). +// Traefik's fsnotify watcher picks up the rename → reload. +// +// On a malformed config Traefik logs an error and holds the +// last-good config (documented Traefik behavior; the C-10 test +// verifies the tmp+rename sequence so a half-written file is never +// observed by Traefik — the only window where Traefik can read the +// file is after the rename, which is atomic on POSIX). +// +// The mode is 0644 (Traefik reads the dynamic dir as root; the lead +// applier chmods after the rename). +func WriteTraefikDynamic(ctx context.Context, t AtomicWriter, peer string, path string, content []byte) error { + if t == nil { + return fmt.Errorf("traefik: atomic writer is nil") + } + if path == "" { + return fmt.Errorf("traefik: path is empty") + } + tmpPath := path + ".tmp" + if _, err := t.WriteFileIdempotent(ctx, peer, tmpPath, content, 0o644); err != nil { + return fmt.Errorf("traefik: write tmp %s: %w", tmpPath, err) + } + // Atomic rename on POSIX. `mv -f` overwrites an existing target + // without prompting. The rename is atomic; Traefik's fsnotify + // watcher observes a single IN_MOVED_TO event. + renameCmd := fmt.Sprintf("mv -f %s %s", shellQuoteLocal(tmpPath), shellQuoteLocal(path)) + if _, err := t.Exec(ctx, peer, renameCmd); err != nil { + return fmt.Errorf("traefik: rename %s -> %s: %w", tmpPath, path, err) + } + return nil +} + +// shellQuoteLocal single-quotes a path for safe shell interpolation on +// the peer. It escapes embedded single-quotes via the standard '\” +// idiom (close the single-quoted string, escape the literal single +// quote, reopen the single-quoted string). This is a local +// re-implementation (the sshpush package has its own) so the emitter +// layer does not depend on the transport package's private helpers — +// the AtomicWriter interface keeps the boundary clean for testing. +func shellQuoteLocal(s string) string { + var b []byte + b = append(b, '\'') + for i := 0; i < len(s); i++ { + c := s[i] + if c == '\'' { + // close quote, escape the literal single-quote, reopen. + b = append(b, '\'', '\\', '\'', '\'') + continue + } + b = append(b, c) + } + b = append(b, '\'') + return string(b) +} diff --git a/internal/emitter/traefik_atomic_test.go b/internal/emitter/traefik_atomic_test.go new file mode 100644 index 0000000..af80a16 --- /dev/null +++ b/internal/emitter/traefik_atomic_test.go @@ -0,0 +1,190 @@ +package emitter + +import ( + "context" + "errors" + "os" + "strings" + "testing" +) + +// mockAtomicWriter is a test-only AtomicWriter that records calls so +// the C-10 atomicity protocol (tmp + rename) can be asserted. +type mockAtomicWriter struct { + written []writeCall + execed []execCall + writeErr error + writeWrote bool + execErr error +} + +type writeCall struct { + peer string + path string + mode os.FileMode + bytes []byte +} + +type execCall struct { + peer string + cmd string +} + +func (m *mockAtomicWriter) WriteFileIdempotent(ctx context.Context, peer string, path string, content []byte, mode os.FileMode) (bool, error) { + m.written = append(m.written, writeCall{peer: peer, path: path, mode: mode, bytes: append([]byte(nil), content...)}) + if m.writeErr != nil { + return false, m.writeErr + } + return m.writeWrote, nil +} + +func (m *mockAtomicWriter) Exec(ctx context.Context, peer string, cmd string) ([]byte, error) { + m.execed = append(m.execed, execCall{peer: peer, cmd: cmd}) + if m.execErr != nil { + return nil, m.execErr + } + return []byte("ok"), nil +} + +func TestWriteTraefikDynamic_TmpThenRename(t *testing.T) { + // Gate C-10: the Traefik dynamic-config write must be a tmp + + // rename sequence so Traefik's fsnotify watcher never observes a + // half-written file. + mock := &mockAtomicWriter{writeWrote: true} + path := "/etc/traefik/dynamic/orca-web.yaml" + peer := "node-1:22" + content := []byte("http:\n routers: {}\n") + + if err := WriteTraefikDynamic(context.Background(), mock, peer, path, content); err != nil { + t.Fatalf("WriteTraefikDynamic: %v", err) + } + + if len(mock.written) != 1 { + t.Fatalf("WriteFileIdempotent calls = %d, want 1", len(mock.written)) + } + w := mock.written[0] + if w.peer != peer { + t.Errorf("write peer = %q, want %q", w.peer, peer) + } + // The tmp path is the target path + ".tmp". + if w.path != path+".tmp" { + t.Errorf("write path = %q, want %q (.tmp suffix is the C-10 atomicity protocol)", w.path, path+".tmp") + } + if string(w.bytes) != string(content) { + t.Errorf("write content = %q, want %q", string(w.bytes), string(content)) + } + if w.mode != 0o644 { + t.Errorf("write mode = %o, want 0644", w.mode) + } + + if len(mock.execed) != 1 { + t.Fatalf("Exec calls = %d, want 1 (the rename)", len(mock.execed)) + } + e := mock.execed[0] + if e.peer != peer { + t.Errorf("exec peer = %q, want %q", e.peer, peer) + } + // The rename command must `mv -f` the .tmp file to the final path. + if !strings.Contains(e.cmd, "mv -f") { + t.Errorf("exec cmd = %q, want it to contain 'mv -f' (atomic rename)", e.cmd) + } + if !strings.Contains(e.cmd, path+".tmp") { + t.Errorf("exec cmd = %q, want it to contain the .tmp path as source", e.cmd) + } + if !strings.Contains(e.cmd, path) { + t.Errorf("exec cmd = %q, want it to contain the final path as destination", e.cmd) + } + // Sanity: the source must come before the destination in the + // mv command. + srcIdx := strings.Index(e.cmd, path+".tmp") + dstIdx := strings.Index(e.cmd, "'"+path+"'") + if srcIdx < 0 || dstIdx < 0 || srcIdx > dstIdx { + t.Errorf("exec cmd %q: source .tmp must come before destination %s", e.cmd, path) + } +} + +func TestWriteTraefikDynamic_WriteTmpError(t *testing.T) { + mock := &mockAtomicWriter{writeErr: errors.New("disk full")} + err := WriteTraefikDynamic(context.Background(), mock, "p", "/etc/traefik/dynamic/orca-x.yaml", []byte("x")) + if err == nil { + t.Fatal("expected error from WriteFileIdempotent, got nil") + } + if !strings.Contains(err.Error(), "write tmp") { + t.Errorf("error = %q, want 'write tmp'", err.Error()) + } + if !strings.Contains(err.Error(), "disk full") { + t.Errorf("error = %q, want underlying 'disk full'", err.Error()) + } + if len(mock.execed) != 0 { + t.Errorf("on tmp write failure, no rename should happen; execed = %v", mock.execed) + } +} + +func TestWriteTraefikDynamic_RenameError(t *testing.T) { + mock := &mockAtomicWriter{writeWrote: true, execErr: errors.New("permission denied")} + err := WriteTraefikDynamic(context.Background(), mock, "p", "/etc/traefik/dynamic/orca-x.yaml", []byte("x")) + if err == nil { + t.Fatal("expected error from rename, got nil") + } + if !strings.Contains(err.Error(), "rename") { + t.Errorf("error = %q, want 'rename'", err.Error()) + } + if !strings.Contains(err.Error(), "permission denied") { + t.Errorf("error = %q, want underlying 'permission denied'", err.Error()) + } +} + +func TestWriteTraefikDynamic_NilWriter(t *testing.T) { + err := WriteTraefikDynamic(context.Background(), nil, "p", "/x", []byte("x")) + if err == nil { + t.Fatal("expected error for nil writer") + } + if !strings.Contains(err.Error(), "nil") { + t.Errorf("error = %q, want 'nil'", err.Error()) + } +} + +func TestWriteTraefikDynamic_EmptyPath(t *testing.T) { + mock := &mockAtomicWriter{writeWrote: true} + err := WriteTraefikDynamic(context.Background(), mock, "p", "", []byte("x")) + if err == nil { + t.Fatal("expected error for empty path") + } + if !strings.Contains(err.Error(), "path is empty") { + t.Errorf("error = %q, want 'path is empty'", err.Error()) + } +} + +func TestWriteTraefikDynamic_SkipWhenContentMatches(t *testing.T) { + // When the .tmp file already matches (writeWrote=false), the + // protocol still proceeds with the rename — the idempotency + // check is per-file, not per-protocol. The rename still happens + // so the final path reflects the (unchanged) content. + mock := &mockAtomicWriter{writeWrote: false} + err := WriteTraefikDynamic(context.Background(), mock, "p", "/etc/traefik/dynamic/orca-x.yaml", []byte("x")) + if err != nil { + t.Fatalf("WriteTraefikDynamic: %v", err) + } + if len(mock.execed) != 1 { + t.Errorf("rename should still happen on idempotent skip; execed = %v", mock.execed) + } +} + +func TestShellQuoteLocal(t *testing.T) { + cases := []struct { + in, want string + }{ + {"/etc/traefik/dynamic/orca-web.yaml", "'/etc/traefik/dynamic/orca-web.yaml'"}, + {"", "''"}, + {"/path with space/x", "'/path with space/x'"}, + {"a'b", "'a'\\''b'"}, + } + for _, tc := range cases { + t.Run(tc.in, func(t *testing.T) { + got := shellQuoteLocal(tc.in) + if got != tc.want { + t.Errorf("shellQuoteLocal(%q) = %q, want %q", tc.in, got, tc.want) + } + }) + } +} diff --git a/internal/emitter/traefik_test.go b/internal/emitter/traefik_test.go new file mode 100644 index 0000000..cf5552c --- /dev/null +++ b/internal/emitter/traefik_test.go @@ -0,0 +1,353 @@ +package emitter + +import ( + "strings" + "testing" + + "git.cloudinit.dev/coreci/orca/internal/jobspec" +) + +func TestTraefikEmitter_RenderBasic(t *testing.T) { + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Runtime: &jobspec.RuntimeBlock{OneOf: "process"}, + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + Health: &jobspec.HealthBlock{CheckType: "http", Interval: "5s", Timeout: "1s"}, + } + node := &Node{Hostname: "node-1", Runtime: []string{"process"}} + files, err := TraefikEmitter{}.Render(spec, node) + if err != nil { + t.Fatalf("Render: %v", err) + } + if len(files) != 1 { + t.Fatalf("got %d files, want 1", len(files)) + } + f := files[0] + wantPath := "/etc/traefik/dynamic/orca-web.yaml" + if f.Path != wantPath { + t.Errorf("Path = %q, want %q", f.Path, wantPath) + } + if f.Mode != "0644" { + t.Errorf("Mode = %q, want 0644", f.Mode) + } + c := f.Content + if !strings.Contains(c, "http:") { + t.Errorf("content missing 'http:'\n%s", c) + } + if !strings.Contains(c, "routers:") { + t.Errorf("content missing 'routers:'\n%s", c) + } + if !strings.Contains(c, "orca-web:") { + t.Errorf("content missing 'orca-web:' router/service key\n%s", c) + } + if !strings.Contains(c, `rule: PathPrefix("/web")`) { + t.Errorf("content missing PathPrefix rule\n%s", c) + } + if !strings.Contains(c, "services:") { + t.Errorf("content missing 'services:'\n%s", c) + } + if !strings.Contains(c, "loadBalancer:") { + t.Errorf("content missing 'loadBalancer:'\n%s", c) + } + if !strings.Contains(c, "unix:///run/orca/alloc-node-1/port-http.sock") { + t.Errorf("content missing socket server URL\n%s", c) + } + if !strings.Contains(c, "certResolver: orca") { + t.Errorf("content missing 'certResolver: orca'\n%s", c) + } + if !strings.Contains(c, "domains:") { + t.Errorf("content missing TLS domains\n%s", c) + } + if !strings.Contains(c, "healthCheck:") { + t.Errorf("content missing 'healthCheck:'\n%s", c) + } + if !strings.Contains(c, "interval: 5s") { + t.Errorf("content missing 'interval: 5s'\n%s", c) + } + if !strings.Contains(c, "timeout: 1s") { + t.Errorf("content missing 'timeout: 1s'\n%s", c) + } +} + +func TestTraefikEmitter_RenderMultiplePorts(t *testing.T) { + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "api", + Runtime: &jobspec.RuntimeBlock{OneOf: "process"}, + Ports: []jobspec.PortSpec{ + {Name: "http", Port: 8080}, + {Name: "grpc", Port: 9090}, + }, + } + node := &Node{Hostname: "n1"} + files, err := TraefikEmitter{}.Render(spec, node) + if err != nil { + t.Fatalf("Render: %v", err) + } + c := files[0].Content + if !strings.Contains(c, "port-http.sock") { + t.Errorf("missing http socket: %s", c) + } + if !strings.Contains(c, "port-grpc.sock") { + t.Errorf("missing grpc socket: %s", c) + } +} + +func TestTraefikEmitter_RenderDrain(t *testing.T) { + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + } + node := &Node{Hostname: "n1"} + files, err := TraefikEmitter{}.RenderDrain(spec, node) + if err != nil { + t.Fatalf("RenderDrain: %v", err) + } + if len(files) != 1 { + t.Fatalf("got %d files, want 1", len(files)) + } + c := files[0].Content + if !strings.Contains(c, "weight: 0") { + t.Errorf("drain config missing 'weight: 0'\n%s", c) + } + if !strings.Contains(c, "unix:///run/orca/alloc-n1/port-http.sock") { + t.Errorf("drain config missing socket URL\n%s", c) + } +} + +func TestTraefikEmitter_RenderLiveHasNoWeightZero(t *testing.T) { + // Sanity: the live (non-drain) render must NOT emit `weight: 0`. + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + } + node := &Node{Hostname: "n1"} + files, err := TraefikEmitter{}.Render(spec, node) + if err != nil { + t.Fatalf("Render: %v", err) + } + if strings.Contains(files[0].Content, "weight: 0") { + t.Errorf("live config should not contain 'weight: 0'\n%s", files[0].Content) + } +} + +func TestTraefikEmitter_RenderNoHealthOmitsHealthCheck(t *testing.T) { + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + } + node := &Node{Hostname: "n1"} + files, err := TraefikEmitter{}.Render(spec, node) + if err != nil { + t.Fatalf("Render: %v", err) + } + if strings.Contains(files[0].Content, "healthCheck:") { + t.Errorf("config without Health should omit 'healthCheck:'\n%s", files[0].Content) + } +} + +func TestTraefikEmitter_NilSpec(t *testing.T) { + _, err := TraefikEmitter{}.Render(nil, &Node{}) + if err == nil { + t.Fatal("expected error for nil spec") + } +} + +func TestTraefikEmitter_EmptyName(t *testing.T) { + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: " ", + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + } + _, err := TraefikEmitter{}.Render(spec, &Node{}) + if err == nil { + t.Fatal("expected error for empty name") + } +} + +func TestTraefikEmitter_NoPorts(t *testing.T) { + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + } + _, err := TraefikEmitter{}.Render(spec, &Node{}) + if err == nil { + t.Fatal("expected error for missing ports") + } + if !strings.Contains(err.Error(), "no ports") { + t.Errorf("error = %q, want 'no ports'", err.Error()) + } +} + +func TestTraefikEmitter_NoPortsDrain(t *testing.T) { + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + } + _, err := TraefikEmitter{}.RenderDrain(spec, &Node{}) + if err == nil { + t.Fatal("expected error for missing ports on drain") + } +} + +func TestTraefikEmitter_InvalidBind(t *testing.T) { + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + Service: &jobspec.ServiceBlock{Bind: "not-an-ip"}, + } + _, err := TraefikEmitter{}.Render(spec, &Node{}) + if err == nil { + t.Fatal("expected error for invalid service.bind") + } + if !strings.Contains(err.Error(), "valid IP") { + t.Errorf("error = %q, want 'valid IP'", err.Error()) + } +} + +func TestTraefikEmitter_ValidBindLoopback(t *testing.T) { + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + Service: &jobspec.ServiceBlock{Bind: "127.0.0.1"}, + } + _, err := TraefikEmitter{}.Render(spec, &Node{}) + if err != nil { + t.Fatalf("127.0.0.1 should be accepted, got %v", err) + } +} + +func TestTraefikEmitter_NilNodeAllocPlaceholder(t *testing.T) { + // With a nil node, the alloc-id placeholder is the literal + // "" sentinel so the rendered config is still valid YAML + // (the P08 socket layer substitutes the real alloc-id). + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + } + files, err := TraefikEmitter{}.Render(spec, nil) + if err != nil { + t.Fatalf("Render: %v", err) + } + if !strings.Contains(files[0].Content, "alloc-") { + t.Errorf("nil node should render alloc- placeholder\n%s", files[0].Content) + } +} + +func TestTraefikEmitter_EmptyHostnameAllocPlaceholder(t *testing.T) { + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + } + files, err := TraefikEmitter{}.Render(spec, &Node{Hostname: " "}) + if err != nil { + t.Fatalf("Render: %v", err) + } + if !strings.Contains(files[0].Content, "alloc-") { + t.Errorf("empty hostname should render alloc- placeholder\n%s", files[0].Content) + } +} + +func TestTraefikEmitter_PathNotOrcaV1Prefixed(t *testing.T) { + // REQ-090: the orca-v1- prefix is only for systemd units; Traefik + // dynamic-config paths are named orca-.yaml (single + // source of truth — no dual-write window for Traefik configs). + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + } + files, err := TraefikEmitter{}.Render(spec, &Node{Hostname: "n1"}) + if err != nil { + t.Fatalf("Render: %v", err) + } + if strings.Contains(files[0].Path, "orca-v1-") { + t.Errorf("Path %q should NOT contain the orca-v1- prefix (systemd-only)", files[0].Path) + } + if !strings.HasPrefix(files[0].Path, "/etc/traefik/dynamic/orca-") { + t.Errorf("Path %q should start with /etc/traefik/dynamic/orca-", files[0].Path) + } + if !strings.HasSuffix(files[0].Path, ".yaml") { + t.Errorf("Path %q should end with .yaml", files[0].Path) + } +} + +func TestTraefikEmitter_RenderYAMLHasRoutersServicesTLS(t *testing.T) { + // Aggregate structural assertion: the rendered YAML has the four + // top-level Traefik concepts (routers, services, tls, servers). + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + Health: &jobspec.HealthBlock{CheckType: "http"}, + } + files, err := TraefikEmitter{}.Render(spec, &Node{Hostname: "n1"}) + if err != nil { + t.Fatalf("Render: %v", err) + } + c := files[0].Content + for _, want := range []string{"routers:", "services:", "tls:", "servers:", "url:"} { + if !strings.Contains(c, want) { + t.Errorf("rendered YAML missing %q\n%s", want, c) + } + } +} + +func TestRegisterTraefik_AllServiceRuntimes(t *testing.T) { + r := NewRegistry() + RegisterTraefik(r) + spec := func(runtime string) *jobspec.WorkloadSpec { + return &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Runtime: &jobspec.RuntimeBlock{OneOf: runtime}, + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + } + } + for _, runtime := range []string{"process", "podman", "wasm"} { + t.Run(runtime, func(t *testing.T) { + files, err := r.Render(spec(runtime), &Node{Hostname: "n1"}) + if err != nil { + t.Fatalf("Render(service:%s): %v", runtime, err) + } + if len(files) != 1 { + t.Fatalf("got %d files, want 1", len(files)) + } + if !strings.Contains(files[0].Path, "/etc/traefik/dynamic/orca-web.yaml") { + t.Errorf("Path = %q", files[0].Path) + } + }) + } +} + +func TestRegisterTraefik_OverwritesExisting(t *testing.T) { + // RegisterTraefik should overwrite any prior registration (the + // Registry documents last-wins). + r := NewRegistry() + r.Register("service:process", mockEmitter{files: []File{{Path: "/old"}}}) + RegisterTraefik(r) + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Runtime: &jobspec.RuntimeBlock{OneOf: "process"}, + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + } + files, err := r.Render(spec, &Node{Hostname: "n1"}) + if err != nil { + t.Fatalf("Render: %v", err) + } + if files[0].Path == "/old" { + t.Errorf("RegisterTraefik did not overwrite the prior registration") + } +} + +// Compile-time assertion that TraefikEmitter implements Emitter. +var _ Emitter = TraefikEmitter{} diff --git a/internal/jobspec/markdown.go b/internal/jobspec/markdown.go index a62b437..e43f6fa 100644 --- a/internal/jobspec/markdown.go +++ b/internal/jobspec/markdown.go @@ -26,14 +26,13 @@ type WorkloadSpec struct { Body string // Kind-specific blocks consumed by the P0c schema validators - // (internal/spec/schema). The Markdown parser does not populate - // these yet; later phases (P02 service block, P03 update stanza, - // P04 lifecycle hooks) extend the parser. P0c only defines the - // struct shape so validators can reference the fields. + // (internal/spec/schema). P02 populates Restart, Update, Service, + // Health, Constraints, Affinity, Lifecycle from the Markdown + // frontmatter (the rest are still populated by later phases). // Restart is the restart policy block. Required for Service and // DaemonSet; optional for Job (defaults to never/on-failure). - // Populated by the P04 lifecycle phase. + // P02 populates it from the `restart:` frontmatter block. Restart *RestartBlock // Schedule is the schedule block. For Job it carries an optional @@ -43,15 +42,35 @@ type WorkloadSpec struct { Schedule *ScheduleBlock // Update is the rolling/canary update stanza. Required for - // Service. Populated by the P03 update-stanza phase. + // Service. P02 populates it from the `update:` frontmatter block; + // the rolling/canary semantics land in P03. Update *UpdateBlock // Service is the service block (Traefik route definition). For // Service kind it is implied; Job and DaemonSet do not carry a - // Traefik route by default (D-175). Populated by the P02 service - // block phase. + // Traefik route by default (D-175). P02 populates it from the + // `service:` frontmatter block. Service *ServiceBlock + // Health is the health-check block. Required for Service (Traefik + // routing depends on it). P02 populates it from the `health:` + // frontmatter block (R-012). + Health *HealthBlock + + // Constraints is the CEL expression list for placement. P02 + // populates it from the `constraints:` frontmatter array; P05 + // consumes it for the CLI-side scheduler (REQ-083). + Constraints []string + + // Affinity is the affinity rule list for placement. P02 populates + // it from the `affinity:` frontmatter array; P05 consumes it. + Affinity []AffinityRule + + // Lifecycle is the lifecycle hook block (pre_stop, post_start). + // P02 populates it from the `lifecycle:` frontmatter block; P04 + // wires it into the systemd unit (ExecStop / ExecStartPost). + Lifecycle *LifecycleBlock + // Timeout is an optional execution timeout (duration string) for // Job. Populated by P04. Timeout string @@ -67,7 +86,8 @@ type RuntimeBlock struct { } // RestartBlock is the restart policy block. Mode is one of never, -// on-failure, service (REQ-074 schema validators). Populated by P04. +// on-failure, service (REQ-074 schema validators). P02 populates it from +// the frontmatter `restart:` block. type RestartBlock struct { Mode string MaxRetries int @@ -83,18 +103,56 @@ type ScheduleBlock struct { } // UpdateBlock is the rolling/canary update stanza. Required for Service. -// Populated by P03. +// P02 populates it from the frontmatter `update:` block; the +// rolling/canary/blue-green semantics land in P03. type UpdateBlock struct { - Strategy string - MaxSurge int + Strategy string + MaxSurge int + MaxParallel int + MinHealthyTime string + HealthyDeadline string + Canary string + AutoPromote bool } // ServiceBlock is the Traefik route definition. For Service it is // implied (Traefik route YES); Job and DaemonSet do not carry one by -// default (D-175). Populated by P02. +// default (D-175). P02 populates it from the frontmatter `service:` +// block. type ServiceBlock struct { Host string RouteID string + Name string + Port int + Bind string +} + +// HealthBlock is the health-check block. P02 populates it from the +// frontmatter `health:` block (R-012). The Traefik emitter (REQ-077) +// renders it as the service's health-check stanza; ServiceValidator +// requires it for Traefik routing. +type HealthBlock struct { + CheckType string + Interval string + Timeout string + UnhealthyThreshold int +} + +// AffinityRule is a single affinity entry: target CEL expression + +// integer weight. P02 populates it from the `affinity:` frontmatter +// array; P05 consumes it for the CLI-side scheduler (REQ-083). +type AffinityRule struct { + Target string + Weight int +} + +// LifecycleBlock is the lifecycle hook block. PreStop and PostStart +// are command lists run before stop / after start. P02 populates it +// from the `lifecycle:` frontmatter block; P04 wires it into the +// systemd unit (ExecStop / ExecStartPost). +type LifecycleBlock struct { + PreStop []string + PostStart []string } // PortSpec is a minimal port binding entry. HostIP is optional. @@ -319,10 +377,19 @@ func parseFrontmatterBlock(block string) (*WorkloadSpec, error) { secEnv secSecrets secVolumes + secRestart + secUpdate + secService + secHealth + secLifecycle + secAffinity + secConstraints ) cur := secNone var curPort *PortSpec var curVol *VolumeSpec + var curAffinity *AffinityRule + var lifecycleCur string flushPort := func() { if curPort != nil { @@ -336,6 +403,12 @@ func parseFrontmatterBlock(block string) (*WorkloadSpec, error) { curVol = nil } } + flushAffinity := func() { + if curAffinity != nil { + spec.Affinity = append(spec.Affinity, *curAffinity) + curAffinity = nil + } + } for lineNo, raw := range lines { line := stripComment(raw) @@ -349,6 +422,7 @@ func parseFrontmatterBlock(block string) (*WorkloadSpec, error) { // Flush any pending nested entry before switching sections. flushPort() flushVol() + flushAffinity() cur = secNone key, val, ok := splitKV(trimmed) @@ -392,8 +466,42 @@ func parseFrontmatterBlock(block string) (*WorkloadSpec, error) { } case "volumes": cur = secVolumes + case "restart": + spec.Restart = &RestartBlock{} + cur = secRestart + case "update": + spec.Update = &UpdateBlock{} + cur = secUpdate + case "service": + spec.Service = &ServiceBlock{} + cur = secService + case "health": + spec.Health = &HealthBlock{} + cur = secHealth + case "lifecycle": + spec.Lifecycle = &LifecycleBlock{} + cur = secLifecycle + case "constraints": + if strings.TrimSpace(val) != "" { + arr, err := parseStringArray(val) + if err != nil { + return nil, fmt.Errorf("parse markdown: line %d: constraints: %w", lineNo+1, err) + } + spec.Constraints = append(spec.Constraints, arr...) + cur = secNone + } else { + cur = secConstraints + } + case "affinity": + if strings.TrimSpace(val) != "" { + // Inline form not supported for affinity objects; + // require the block form. Ignore inline values. + cur = secNone + } else { + cur = secAffinity + } default: - // Unknown top-level keys are ignored (forward-compat). + // Unknown top-level key are ignored (forward-compat). cur = secNone } continue @@ -464,10 +572,159 @@ func parseFrontmatterBlock(block string) (*WorkloadSpec, error) { } else if curVol != nil { applyVolumeKV(curVol, trimmed) } + case secRestart: + if spec.Restart == nil { + spec.Restart = &RestartBlock{} + } + key, val, ok := splitKV(trimmed) + if !ok { + continue + } + switch key { + case "mode": + spec.Restart.Mode = unquote(val) + case "attempts", "max_retries": + if n, err := strconv.Atoi(strings.TrimSpace(unquote(val))); err == nil { + spec.Restart.MaxRetries = n + } + case "delay": + spec.Restart.Delay = unquote(val) + } + case secUpdate: + if spec.Update == nil { + spec.Update = &UpdateBlock{} + } + key, val, ok := splitKV(trimmed) + if !ok { + continue + } + switch key { + case "strategy": + spec.Update.Strategy = unquote(val) + case "max_parallel": + if n, err := strconv.Atoi(strings.TrimSpace(unquote(val))); err == nil { + spec.Update.MaxParallel = n + } + case "max_surge": + if n, err := strconv.Atoi(strings.TrimSpace(unquote(val))); err == nil { + spec.Update.MaxSurge = n + } + case "min_healthy_time": + spec.Update.MinHealthyTime = unquote(val) + case "healthy_deadline": + spec.Update.HealthyDeadline = unquote(val) + case "canary": + spec.Update.Canary = unquote(val) + case "auto_promote": + spec.Update.AutoPromote = parseBool(val) + } + case secService: + if spec.Service == nil { + spec.Service = &ServiceBlock{} + } + key, val, ok := splitKV(trimmed) + if !ok { + continue + } + switch key { + case "name": + spec.Service.Name = unquote(val) + case "port": + if n, err := strconv.Atoi(strings.TrimSpace(unquote(val))); err == nil { + spec.Service.Port = n + } + case "bind": + spec.Service.Bind = unquote(val) + case "host": + spec.Service.Host = unquote(val) + case "route_id": + spec.Service.RouteID = unquote(val) + } + case secHealth: + if spec.Health == nil { + spec.Health = &HealthBlock{} + } + key, val, ok := splitKV(trimmed) + if !ok { + continue + } + switch key { + case "check_type": + spec.Health.CheckType = unquote(val) + case "interval": + spec.Health.Interval = unquote(val) + case "timeout": + spec.Health.Timeout = unquote(val) + case "unhealthy_threshold": + if n, err := strconv.Atoi(strings.TrimSpace(unquote(val))); err == nil { + spec.Health.UnhealthyThreshold = n + } + } + case secLifecycle: + if spec.Lifecycle == nil { + spec.Lifecycle = &LifecycleBlock{} + } + // pre_stop / post_start are string arrays. The block form + // is: + // lifecycle: + // pre_stop: + // - cmd1 + // - cmd2 + // post_start: + // - cmd3 + // We track which sub-list we are appending to via a local + // cursor that is reset on every top-level section change. + key, val, ok := splitKV(trimmed) + if !ok { + // Could be a list item under pre_stop/post_start. + if strings.HasPrefix(trimmed, "- ") || trimmed == "-" { + item := strings.TrimSpace(strings.TrimPrefix(trimmed, "-")) + if item != "" && lifecycleCur != "" { + appendLifecycleCmd(spec.Lifecycle, lifecycleCur, unquote(item)) + } + } + continue + } + switch key { + case "pre_stop", "post_start": + lifecycleCur = key + if strings.TrimSpace(val) != "" { + // Inline list form: `pre_stop: [cmd1, cmd2]`. + arr, err := parseStringArray(val) + if err == nil { + for _, s := range arr { + appendLifecycleCmd(spec.Lifecycle, key, s) + } + } + lifecycleCur = "" + } + default: + lifecycleCur = "" + } + case secAffinity: + if strings.HasPrefix(trimmed, "- ") || trimmed == "-" { + flushAffinity() + r := AffinityRule{} + curAffinity = &r + rest := strings.TrimSpace(strings.TrimPrefix(trimmed, "-")) + if rest != "" { + applyAffinityKV(curAffinity, rest) + } + } else if curAffinity != nil { + applyAffinityKV(curAffinity, trimmed) + } + case secConstraints: + if strings.HasPrefix(trimmed, "- ") || trimmed == "-" { + item := strings.TrimSpace(strings.TrimPrefix(trimmed, "-")) + if item != "" { + spec.Constraints = append(spec.Constraints, unquote(item)) + } + } } } flushPort() flushVol() + flushAffinity() return spec, nil } @@ -518,6 +775,47 @@ func applyVolumeKV(v *VolumeSpec, s string) { } } +// applyAffinityKV applies a `key: value` pair to an AffinityRule entry. +func applyAffinityKV(r *AffinityRule, s string) { + key, val, ok := splitKV(s) + if !ok { + return + } + switch key { + case "target": + r.Target = unquote(val) + case "weight": + if n, err := strconv.Atoi(strings.TrimSpace(unquote(val))); err == nil { + r.Weight = n + } + } +} + +// appendLifecycleCmd appends a command to the named lifecycle hook list +// (pre_stop or post_start) on the given LifecycleBlock. +func appendLifecycleCmd(lb *LifecycleBlock, name, cmd string) { + if lb == nil || cmd == "" { + return + } + switch name { + case "pre_stop": + lb.PreStop = append(lb.PreStop, cmd) + case "post_start": + lb.PostStart = append(lb.PostStart, cmd) + } +} + +// parseBool parses a YAML-ish boolean value (true/yes/on/1 → true). The +// comparison is case-insensitive. Empty and unrecognized values return +// false (forward-compatible with future strict-mode validation). +func parseBool(s string) bool { + switch strings.ToLower(strings.TrimSpace(unquote(s))) { + case "true", "yes", "on", "1": + return true + } + return false +} + // validateWorkload enforces required fields and kind validity (R-012). func validateWorkload(spec *WorkloadSpec) error { if spec.Kind == "" { diff --git a/internal/jobspec/markdown_test.go b/internal/jobspec/markdown_test.go index 5d29ace..3b8b6b9 100644 --- a/internal/jobspec/markdown_test.go +++ b/internal/jobspec/markdown_test.go @@ -355,3 +355,349 @@ func TestParseMarkdown_UnknownKeyIgnored(t *testing.T) { t.Fatalf("ParseMarkdown should ignore unknown keys: %v", err) } } + +func TestParseMarkdown_RestartBlock(t *testing.T) { + input := "---\n" + + "kind: Service\n" + + "name: web\n" + + "restart:\n" + + " mode: service\n" + + " attempts: 5\n" + + " delay: 3s\n" + + "---\nbody\n" + spec, err := ParseMarkdown([]byte(input)) + if err != nil { + t.Fatalf("ParseMarkdown: %v", err) + } + if spec.Restart == nil { + t.Fatal("Restart is nil") + } + if spec.Restart.Mode != "service" { + t.Errorf("Restart.Mode = %q, want service", spec.Restart.Mode) + } + if spec.Restart.MaxRetries != 5 { + t.Errorf("Restart.MaxRetries = %d, want 5", spec.Restart.MaxRetries) + } + if spec.Restart.Delay != "3s" { + t.Errorf("Restart.Delay = %q, want 3s", spec.Restart.Delay) + } +} + +func TestParseMarkdown_RestartBlockMaxRetriesAlias(t *testing.T) { + // max_retries is the canonical key; attempts is an accepted alias. + input := "---\n" + + "kind: Service\n" + + "name: web\n" + + "restart:\n" + + " mode: on-failure\n" + + " max_retries: 3\n" + + " delay: 1s\n" + + "---\nbody\n" + spec, err := ParseMarkdown([]byte(input)) + if err != nil { + t.Fatalf("ParseMarkdown: %v", err) + } + if spec.Restart == nil || spec.Restart.MaxRetries != 3 { + t.Fatalf("Restart.MaxRetries = %d, want 3 (max_retries alias)", spec.Restart.MaxRetries) + } +} + +func TestParseMarkdown_UpdateBlock(t *testing.T) { + input := "---\n" + + "kind: Service\n" + + "name: web\n" + + "update:\n" + + " strategy: canary\n" + + " max_parallel: 2\n" + + " min_healthy_time: 30s\n" + + " healthy_deadline: 5m\n" + + " canary: 10%\n" + + " auto_promote: true\n" + + "---\nbody\n" + spec, err := ParseMarkdown([]byte(input)) + if err != nil { + t.Fatalf("ParseMarkdown: %v", err) + } + if spec.Update == nil { + t.Fatal("Update is nil") + } + if spec.Update.Strategy != "canary" { + t.Errorf("Update.Strategy = %q, want canary", spec.Update.Strategy) + } + if spec.Update.MaxParallel != 2 { + t.Errorf("Update.MaxParallel = %d, want 2", spec.Update.MaxParallel) + } + if spec.Update.MinHealthyTime != "30s" { + t.Errorf("Update.MinHealthyTime = %q, want 30s", spec.Update.MinHealthyTime) + } + if spec.Update.HealthyDeadline != "5m" { + t.Errorf("Update.HealthyDeadline = %q, want 5m", spec.Update.HealthyDeadline) + } + if spec.Update.Canary != "10%" { + t.Errorf("Update.Canary = %q, want 10%%", spec.Update.Canary) + } + if !spec.Update.AutoPromote { + t.Errorf("Update.AutoPromote = false, want true") + } +} + +func TestParseMarkdown_ServiceBlock(t *testing.T) { + input := "---\n" + + "kind: Service\n" + + "name: web\n" + + "service:\n" + + " name: web\n" + + " port: 8080\n" + + " bind: 127.0.0.1\n" + + "---\nbody\n" + spec, err := ParseMarkdown([]byte(input)) + if err != nil { + t.Fatalf("ParseMarkdown: %v", err) + } + if spec.Service == nil { + t.Fatal("Service is nil") + } + if spec.Service.Name != "web" { + t.Errorf("Service.Name = %q, want web", spec.Service.Name) + } + if spec.Service.Port != 8080 { + t.Errorf("Service.Port = %d, want 8080", spec.Service.Port) + } + if spec.Service.Bind != "127.0.0.1" { + t.Errorf("Service.Bind = %q, want 127.0.0.1", spec.Service.Bind) + } +} + +func TestParseMarkdown_HealthBlock(t *testing.T) { + input := "---\n" + + "kind: Service\n" + + "name: web\n" + + "health:\n" + + " check_type: http\n" + + " interval: 10s\n" + + " timeout: 2s\n" + + " unhealthy_threshold: 3\n" + + "---\nbody\n" + spec, err := ParseMarkdown([]byte(input)) + if err != nil { + t.Fatalf("ParseMarkdown: %v", err) + } + if spec.Health == nil { + t.Fatal("Health is nil") + } + if spec.Health.CheckType != "http" { + t.Errorf("Health.CheckType = %q, want http", spec.Health.CheckType) + } + if spec.Health.Interval != "10s" { + t.Errorf("Health.Interval = %q, want 10s", spec.Health.Interval) + } + if spec.Health.Timeout != "2s" { + t.Errorf("Health.Timeout = %q, want 2s", spec.Health.Timeout) + } + if spec.Health.UnhealthyThreshold != 3 { + t.Errorf("Health.UnhealthyThreshold = %d, want 3", spec.Health.UnhealthyThreshold) + } +} + +func TestParseMarkdown_ConstraintsInlineArray(t *testing.T) { + // Inline flow-array form: the parser does NOT unescape YAML + // escapes (consistent with the secrets inline parser). Use + // single-quoted scalars inside the flow array so the CEL strings + // are preserved verbatim. + input := "---\nkind: Service\nname: web\nconstraints: ['node.role == \"web\"', 'region == \"us\"']\n---\nbody\n" + spec, err := ParseMarkdown([]byte(input)) + if err != nil { + t.Fatalf("ParseMarkdown: %v", err) + } + if len(spec.Constraints) != 2 { + t.Fatalf("Constraints = %d, want 2", len(spec.Constraints)) + } + if spec.Constraints[0] != `node.role == "web"` { + t.Errorf("Constraints[0] = %q", spec.Constraints[0]) + } + if spec.Constraints[1] != `region == "us"` { + t.Errorf("Constraints[1] = %q", spec.Constraints[1]) + } +} + +func TestParseMarkdown_ConstraintsBlockArray(t *testing.T) { + input := "---\n" + + "kind: Service\n" + + "name: web\n" + + "constraints:\n" + + " - node.role == \"web\"\n" + + " - region == \"us\"\n" + + "---\nbody\n" + spec, err := ParseMarkdown([]byte(input)) + if err != nil { + t.Fatalf("ParseMarkdown: %v", err) + } + if len(spec.Constraints) != 2 { + t.Fatalf("Constraints = %d, want 2", len(spec.Constraints)) + } + if spec.Constraints[0] != `node.role == "web"` { + t.Errorf("Constraints[0] = %q", spec.Constraints[0]) + } + if spec.Constraints[1] != `region == "us"` { + t.Errorf("Constraints[1] = %q", spec.Constraints[1]) + } +} + +func TestParseMarkdown_AffinityBlock(t *testing.T) { + input := "---\n" + + "kind: Service\n" + + "name: web\n" + + "affinity:\n" + + " - target: node.role == \"web\"\n" + + " weight: 100\n" + + " - target: region == \"us\"\n" + + " weight: 50\n" + + "---\nbody\n" + spec, err := ParseMarkdown([]byte(input)) + if err != nil { + t.Fatalf("ParseMarkdown: %v", err) + } + if len(spec.Affinity) != 2 { + t.Fatalf("Affinity = %d, want 2", len(spec.Affinity)) + } + if spec.Affinity[0].Target != `node.role == "web"` { + t.Errorf("Affinity[0].Target = %q", spec.Affinity[0].Target) + } + if spec.Affinity[0].Weight != 100 { + t.Errorf("Affinity[0].Weight = %d, want 100", spec.Affinity[0].Weight) + } + if spec.Affinity[1].Target != `region == "us"` { + t.Errorf("Affinity[1].Target = %q", spec.Affinity[1].Target) + } + if spec.Affinity[1].Weight != 50 { + t.Errorf("Affinity[1].Weight = %d, want 50", spec.Affinity[1].Weight) + } +} + +func TestParseMarkdown_LifecycleBlock(t *testing.T) { + input := "---\n" + + "kind: Service\n" + + "name: web\n" + + "lifecycle:\n" + + " pre_stop:\n" + + " - /bin/sh -c 'sleep 5'\n" + + " - /usr/local/bin/drain.sh\n" + + " post_start:\n" + + " - /usr/local/bin/warm-cache.sh\n" + + "---\nbody\n" + spec, err := ParseMarkdown([]byte(input)) + if err != nil { + t.Fatalf("ParseMarkdown: %v", err) + } + if spec.Lifecycle == nil { + t.Fatal("Lifecycle is nil") + } + if len(spec.Lifecycle.PreStop) != 2 { + t.Fatalf("PreStop = %d, want 2", len(spec.Lifecycle.PreStop)) + } + if spec.Lifecycle.PreStop[0] != "/bin/sh -c 'sleep 5'" { + t.Errorf("PreStop[0] = %q", spec.Lifecycle.PreStop[0]) + } + if spec.Lifecycle.PreStop[1] != "/usr/local/bin/drain.sh" { + t.Errorf("PreStop[1] = %q", spec.Lifecycle.PreStop[1]) + } + if len(spec.Lifecycle.PostStart) != 1 { + t.Fatalf("PostStart = %d, want 1", len(spec.Lifecycle.PostStart)) + } + if spec.Lifecycle.PostStart[0] != "/usr/local/bin/warm-cache.sh" { + t.Errorf("PostStart[0] = %q", spec.Lifecycle.PostStart[0]) + } +} + +func TestParseMarkdown_LifecycleInlineArray(t *testing.T) { + input := "---\n" + + "kind: Service\n" + + "name: web\n" + + "lifecycle:\n" + + " pre_stop: [\"/bin/true\"]\n" + + " post_start: [\"/bin/warmup\", \"/bin/check\"]\n" + + "---\nbody\n" + spec, err := ParseMarkdown([]byte(input)) + if err != nil { + t.Fatalf("ParseMarkdown: %v", err) + } + if spec.Lifecycle == nil { + t.Fatal("Lifecycle is nil") + } + if len(spec.Lifecycle.PreStop) != 1 || spec.Lifecycle.PreStop[0] != "/bin/true" { + t.Errorf("PreStop = %v, want [/bin/true]", spec.Lifecycle.PreStop) + } + if len(spec.Lifecycle.PostStart) != 2 { + t.Fatalf("PostStart = %v, want 2 entries", spec.Lifecycle.PostStart) + } + if spec.Lifecycle.PostStart[0] != "/bin/warmup" || spec.Lifecycle.PostStart[1] != "/bin/check" { + t.Errorf("PostStart = %v, want [/bin/warmup /bin/check]", spec.Lifecycle.PostStart) + } +} + +func TestParseMarkdown_FullServiceSpec(t *testing.T) { + // A complete Service spec exercising every P02-parsed block together. + input := "---\n" + + "kind: Service\n" + + "name: web\n" + + "count: 3\n" + + "runtime:\n" + + " one_of: process\n" + + " command: /usr/bin/httpd\n" + + "ports:\n" + + " - name: http\n" + + " port: 8080\n" + + "restart:\n" + + " mode: service\n" + + " attempts: 5\n" + + " delay: 2s\n" + + "update:\n" + + " strategy: rolling\n" + + " max_parallel: 1\n" + + " auto_promote: false\n" + + "service:\n" + + " name: web\n" + + " port: 8080\n" + + "health:\n" + + " check_type: http\n" + + " interval: 5s\n" + + " timeout: 1s\n" + + " unhealthy_threshold: 2\n" + + "constraints:\n" + + " - node.role == \"web\"\n" + + "affinity:\n" + + " - target: zone == \"a\"\n" + + " weight: 80\n" + + "lifecycle:\n" + + " post_start:\n" + + " - /bin/ready.sh\n" + + "---\n# body\n" + spec, err := ParseMarkdown([]byte(input)) + if err != nil { + t.Fatalf("ParseMarkdown: %v", err) + } + if spec.Restart == nil || spec.Restart.Mode != "service" { + t.Errorf("Restart not parsed: %+v", spec.Restart) + } + if spec.Update == nil || spec.Update.Strategy != "rolling" { + t.Errorf("Update not parsed: %+v", spec.Update) + } + if spec.Service == nil || spec.Service.Port != 8080 { + t.Errorf("Service not parsed: %+v", spec.Service) + } + if spec.Health == nil || spec.Health.CheckType != "http" { + t.Errorf("Health not parsed: %+v", spec.Health) + } + if len(spec.Constraints) != 1 { + t.Errorf("Constraints = %v", spec.Constraints) + } + if len(spec.Affinity) != 1 || spec.Affinity[0].Weight != 80 { + t.Errorf("Affinity = %v", spec.Affinity) + } + if spec.Lifecycle == nil || len(spec.Lifecycle.PostStart) != 1 { + t.Errorf("Lifecycle not parsed: %+v", spec.Lifecycle) + } + if spec.Body != "# body\n" { + t.Errorf("Body = %q, want %q (R-015)", spec.Body, "# body\n") + } +} diff --git a/internal/spec/schema/schema.go b/internal/spec/schema/schema.go index 571c88b..a75eda4 100644 --- a/internal/spec/schema/schema.go +++ b/internal/spec/schema/schema.go @@ -13,6 +13,7 @@ package schema import ( "errors" "fmt" + "net" "strings" "git.cloudinit.dev/coreci/orca/internal/jobspec" @@ -44,8 +45,11 @@ type JobValidator struct{} // - ports required (at least one) // - count ≥ 1 // - restart required (mode must be service) -// - update required +// - update required (strategy must be rolling/canary/blue-green) // - runtime required +// - health block required (Traefik routing depends on health checks) +// - 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) // - service block implied (Traefik route YES) type ServiceValidator struct{} @@ -109,18 +113,59 @@ func (ServiceValidator) Validate(spec *jobspec.WorkloadSpec) error { } if spec.Restart == nil { errs = append(errs, "restart block required for Service") - } else if spec.Restart.Mode != "service" { - errs = append(errs, fmt.Sprintf("restart mode must be %q for Service, got %q", "service", spec.Restart.Mode)) + } else { + switch spec.Restart.Mode { + case "service", "on-failure", "never": + // Valid per R-012 (default for Service is "service", + // but the validator accepts the full enum; the + // Service-specific "must be service" rule is enforced + // below for the default case where mode is empty). + case "": + errs = append(errs, "restart mode required for Service (one of service, on-failure, never; default is service)") + default: + errs = append(errs, fmt.Sprintf("restart mode %q invalid (want one of service, on-failure, never)", spec.Restart.Mode)) + } } if spec.Update == nil { errs = append(errs, "update block required for Service") + } else { + switch spec.Update.Strategy { + case "rolling", "canary", "blue-green": + case "": + errs = append(errs, "update strategy required for Service (one of rolling, canary, blue-green)") + default: + errs = append(errs, fmt.Sprintf("update strategy %q invalid (want one of rolling, canary, blue-green)", spec.Update.Strategy)) + } } if spec.Runtime == nil { errs = append(errs, "runtime block required for Service") } + if spec.Health == nil { + errs = append(errs, "health block required for Service (Traefik routing requires health checks)") + } + if spec.Service != nil { + if err := validateServiceBind(spec.Service.Bind); err != nil { + errs = append(errs, err.Error()) + } + } return composeErrors("schema/Service", errs) } +// validateServiceBind validates the service.bind field (R-007). Empty +// is OK (default = socket). When set, it must be a valid IPv4/IPv6 +// address (the only opt-in to bind on a non-loopback address); the +// loopback 127.0.0.1 is the documented opt-in. Anything that is not +// parseable as an IP address is rejected. +func validateServiceBind(bind string) error { + if strings.TrimSpace(bind) == "" { + return nil + } + if net.ParseIP(bind) == nil { + return fmt.Errorf("service.bind %q is not a valid IP address (R-007: 127.0.0.1 opt-in; default is socket)", bind) + } + return nil +} + // Validate validates a DaemonSet spec. See DaemonSetValidator for the rules. func (DaemonSetValidator) Validate(spec *jobspec.WorkloadSpec) error { if spec == nil { diff --git a/internal/spec/schema/schema_test.go b/internal/spec/schema/schema_test.go index 47492b9..af6bbed 100644 --- a/internal/spec/schema/schema_test.go +++ b/internal/spec/schema/schema_test.go @@ -90,6 +90,7 @@ func TestServiceValidator_ValidFull(t *testing.T) { Runtime: &jobspec.RuntimeBlock{OneOf: "process", Command: "/bin/http"}, Restart: &jobspec.RestartBlock{Mode: "service"}, Update: &jobspec.UpdateBlock{Strategy: "rolling", MaxSurge: 1}, + Health: &jobspec.HealthBlock{CheckType: "http", Interval: "5s"}, Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, } v := ServiceValidator{} @@ -159,16 +160,43 @@ func TestServiceValidator_WrongRestartMode(t *testing.T) { Name: "web", Count: 1, Runtime: &jobspec.RuntimeBlock{OneOf: "process"}, - Restart: &jobspec.RestartBlock{Mode: "on-failure"}, + Restart: &jobspec.RestartBlock{Mode: "always"}, 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 wrong restart mode, got nil") + t.Fatal("expected error for invalid restart mode, got nil") } - if !strings.Contains(err.Error(), "restart mode must be") { - t.Errorf("error = %q, want 'restart mode must be'", err.Error()) + if !strings.Contains(err.Error(), "restart mode") { + t.Errorf("error = %q, want 'restart mode'", err.Error()) + } +} + +func TestServiceValidator_AcceptedRestartModes(t *testing.T) { + // R-012: restart.mode accepts service / on-failure / never for + // Service; the default per R-012 is "service" but the validator + // accepts the full enum (a Service that wants on-failure is + // unusual but not invalid — only "always" and unknown modes are + // rejected). + for _, mode := range []string{"service", "on-failure", "never"} { + t.Run(mode, func(t *testing.T) { + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Count: 1, + Runtime: &jobspec.RuntimeBlock{OneOf: "process"}, + Restart: &jobspec.RestartBlock{Mode: mode}, + Update: &jobspec.UpdateBlock{Strategy: "rolling"}, + Health: &jobspec.HealthBlock{CheckType: "http"}, + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + } + v := ServiceValidator{} + if err := v.Validate(spec); err != nil { + t.Errorf("mode %q should be accepted, got: %v", mode, err) + } + }) } } @@ -215,6 +243,233 @@ func TestServiceValidator_NilSpec(t *testing.T) { } } +func TestServiceValidator_MissingHealth(t *testing.T) { + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Count: 1, + Runtime: &jobspec.RuntimeBlock{OneOf: "process"}, + Restart: &jobspec.RestartBlock{Mode: "service"}, + Update: &jobspec.UpdateBlock{Strategy: "rolling"}, + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + } + err := ServiceValidator{}.Validate(spec) + if err == nil { + t.Fatal("expected error for missing health block, got nil") + } + if !strings.Contains(err.Error(), "health block required") { + t.Errorf("error = %q, want 'health block required'", err.Error()) + } +} + +func TestServiceValidator_InvalidRestartMode(t *testing.T) { + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Count: 1, + Runtime: &jobspec.RuntimeBlock{OneOf: "process"}, + Restart: &jobspec.RestartBlock{Mode: "always"}, + 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 invalid restart mode, got nil") + } + if !strings.Contains(err.Error(), "restart mode") || !strings.Contains(err.Error(), "invalid") { + t.Errorf("error = %q, want 'restart mode ... invalid'", err.Error()) + } +} + +func TestServiceValidator_EmptyRestartMode(t *testing.T) { + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Count: 1, + Runtime: &jobspec.RuntimeBlock{OneOf: "process"}, + Restart: &jobspec.RestartBlock{Mode: ""}, + 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 empty restart mode, got nil") + } + if !strings.Contains(err.Error(), "restart mode required") { + t.Errorf("error = %q, want 'restart mode required'", err.Error()) + } +} + +func TestServiceValidator_InvalidUpdateStrategy(t *testing.T) { + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Count: 1, + Runtime: &jobspec.RuntimeBlock{OneOf: "process"}, + Restart: &jobspec.RestartBlock{Mode: "service"}, + Update: &jobspec.UpdateBlock{Strategy: "recreate"}, + Health: &jobspec.HealthBlock{CheckType: "http"}, + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + } + err := ServiceValidator{}.Validate(spec) + if err == nil { + t.Fatal("expected error for invalid update strategy, got nil") + } + if !strings.Contains(err.Error(), "update strategy") || !strings.Contains(err.Error(), "invalid") { + t.Errorf("error = %q, want 'update strategy ... invalid'", err.Error()) + } +} + +func TestServiceValidator_EmptyUpdateStrategy(t *testing.T) { + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Count: 1, + Runtime: &jobspec.RuntimeBlock{OneOf: "process"}, + Restart: &jobspec.RestartBlock{Mode: "service"}, + Update: &jobspec.UpdateBlock{Strategy: ""}, + Health: &jobspec.HealthBlock{CheckType: "http"}, + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + } + err := ServiceValidator{}.Validate(spec) + if err == nil { + t.Fatal("expected error for empty update strategy, got nil") + } + if !strings.Contains(err.Error(), "update strategy required") { + t.Errorf("error = %q, want 'update strategy required'", err.Error()) + } +} + +func TestServiceValidator_AcceptedUpdateStrategies(t *testing.T) { + for _, strat := range []string{"rolling", "canary", "blue-green"} { + t.Run(strat, func(t *testing.T) { + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Count: 1, + Runtime: &jobspec.RuntimeBlock{OneOf: "process"}, + Restart: &jobspec.RestartBlock{Mode: "service"}, + Update: &jobspec.UpdateBlock{Strategy: strat}, + Health: &jobspec.HealthBlock{CheckType: "http"}, + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + } + v := ServiceValidator{} + if err := v.Validate(spec); err != nil { + t.Errorf("strategy %q should be accepted, got: %v", strat, err) + } + }) + } +} + +func TestServiceValidator_InvalidServiceBind(t *testing.T) { + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Count: 1, + Runtime: &jobspec.RuntimeBlock{OneOf: "process"}, + Restart: &jobspec.RestartBlock{Mode: "service"}, + Update: &jobspec.UpdateBlock{Strategy: "rolling"}, + Health: &jobspec.HealthBlock{CheckType: "http"}, + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + Service: &jobspec.ServiceBlock{Bind: "not-an-ip"}, + } + err := ServiceValidator{}.Validate(spec) + if err == nil { + t.Fatal("expected error for invalid service.bind, got nil") + } + if !strings.Contains(err.Error(), "service.bind") || !strings.Contains(err.Error(), "valid IP") { + t.Errorf("error = %q, want 'service.bind ... valid IP'", err.Error()) + } +} + +func TestServiceValidator_ValidServiceBindLoopback(t *testing.T) { + // R-007: 127.0.0.1 is the documented opt-in for a non-socket bind. + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Count: 1, + Runtime: &jobspec.RuntimeBlock{OneOf: "process"}, + Restart: &jobspec.RestartBlock{Mode: "service"}, + Update: &jobspec.UpdateBlock{Strategy: "rolling"}, + Health: &jobspec.HealthBlock{CheckType: "http"}, + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + Service: &jobspec.ServiceBlock{Bind: "127.0.0.1"}, + } + v := ServiceValidator{} + if err := v.Validate(spec); err != nil { + t.Fatalf("127.0.0.1 should be accepted, got: %v", err) + } +} + +func TestServiceValidator_ValidServiceBindIPv6(t *testing.T) { + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Count: 1, + Runtime: &jobspec.RuntimeBlock{OneOf: "process"}, + Restart: &jobspec.RestartBlock{Mode: "service"}, + Update: &jobspec.UpdateBlock{Strategy: "rolling"}, + Health: &jobspec.HealthBlock{CheckType: "http"}, + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + Service: &jobspec.ServiceBlock{Bind: "::1"}, + } + v := ServiceValidator{} + if err := v.Validate(spec); err != nil { + t.Fatalf("::1 should be accepted, got: %v", err) + } +} + +func TestServiceValidator_EmptyServiceBindOK(t *testing.T) { + // R-007: empty bind = default = socket; valid. + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "web", + Count: 1, + Runtime: &jobspec.RuntimeBlock{OneOf: "process"}, + Restart: &jobspec.RestartBlock{Mode: "service"}, + Update: &jobspec.UpdateBlock{Strategy: "rolling"}, + Health: &jobspec.HealthBlock{CheckType: "http"}, + Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}}, + Service: &jobspec.ServiceBlock{Bind: ""}, + } + v := ServiceValidator{} + if err := v.Validate(spec); err != nil { + t.Fatalf("empty bind should default to socket (valid), got: %v", err) + } +} + +func TestServiceValidator_MultipleErrors(t *testing.T) { + // Multiple violations should all surface in the composed error. + spec := &jobspec.WorkloadSpec{ + Kind: "Service", + Name: "", + Count: 0, + Restart: &jobspec.RestartBlock{Mode: "always"}, + Update: &jobspec.UpdateBlock{Strategy: "recreate"}, + Service: &jobspec.ServiceBlock{Bind: "not-an-ip"}, + } + err := ServiceValidator{}.Validate(spec) + if err == nil { + t.Fatal("expected error, got nil") + } + for _, want := range []string{ + "name is required", + "ports required", + "count must be", + "restart mode", + "update strategy", + "runtime block required", + "health block required", + "service.bind", + } { + if !strings.Contains(err.Error(), want) { + t.Errorf("error %q missing %q", err.Error(), want) + } + } +} + func TestDaemonSetValidator_Valid(t *testing.T) { spec := &jobspec.WorkloadSpec{ Kind: "DaemonSet", diff --git a/internal/sshpush/atomic_writer_test.go b/internal/sshpush/atomic_writer_test.go new file mode 100644 index 0000000..55b1eb4 --- /dev/null +++ b/internal/sshpush/atomic_writer_test.go @@ -0,0 +1,30 @@ +// Package sshpush_test contains compile-time assertions that *Transport +// satisfies the emitter.AtomicWriter interface (the Traefik C-10 +// atomicity protocol — internal/emitter/traefik_atomic.go). The +// assertion lives here (not in internal/emitter) to avoid an import +// cycle: internal/emitter is imported by this package (fanout.go), so +// internal/emitter cannot import this package. +package sshpush_test + +import ( + "context" + "testing" + + "git.cloudinit.dev/coreci/orca/internal/emitter" + "git.cloudinit.dev/coreci/orca/internal/sshpush" +) + +// Compile-time assertion: *sshpush.Transport satisfies +// emitter.AtomicWriter. WriteTraefikDynamic relies on this so the +// Traefik dynamic-config file is written atomically (gate C-10). +var _ emitter.AtomicWriter = (*sshpush.Transport)(nil) + +func TestTransportSatisfiesAtomicWriter(t *testing.T) { + // A trivial runtime check that the type conversion is valid; the + // compile-time assertion above is the real test, but this gives + // `go test` a function to run. + tr := sshpush.NewTransport("/nonexistent", "/nonexistent") + var w emitter.AtomicWriter = tr + _ = w + _ = context.Background() +}