Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| ea00158fa5 | |||
| ae6eb5a27b |
@@ -1 +1 @@
|
||||
{ "phase": "P02", "stage": "verify", "milestone": "v0.9", "phase_role": "execution", "updated_at": "2026-08-05T03:55:00Z", "milestone_complete": false, "gates_cleared_this_phase": ["C-10"], "verify": { "build": "pass", "go_test": "22/22", "bats": "20/20", "gofmt": "clean", "verify_reqs": "90 consistent" } }
|
||||
{ "phase": "P03/P04/P08", "stage": "verify", "milestone": "v0.9", "phase_role": "execution", "updated_at": "2026-08-05T04:05:00Z", "milestone_complete": false, "verify": { "build": "pass", "go_test": "22/22", "bats": "20/20", "gofmt": "clean", "verify_reqs": "90 consistent" } }
|
||||
|
||||
@@ -0,0 +1,128 @@
|
||||
package emitter
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
"git.cloudinit.dev/coreci/orca/internal/jobspec"
|
||||
)
|
||||
|
||||
// SocketEmitter renders the systemd directives that implement the
|
||||
// R-007 socket-plumbing contract: workloads bind to
|
||||
// /run/orca/alloc-<id>/port-<name>.sock unless overridden via
|
||||
// service.bind = "127.0.0.1" (the only documented opt-in).
|
||||
//
|
||||
// The systemd side of the contract uses two directives:
|
||||
//
|
||||
// - RuntimeDirectory=orca/alloc-<alloc-id> — systemd creates
|
||||
// /run/orca/alloc-<alloc-id>/ owned by the service user (orca:orca)
|
||||
// with mode 0750. The directory is removed when the unit stops
|
||||
// (RuntimeDirectory= semantics). P08 emits one RuntimeDirectory=
|
||||
// line per port so each port's socket directory is created; the
|
||||
// alloc-id placeholder is spec.Name (the real alloc-id is assigned
|
||||
// by the scheduler at submit time — see allocIDFor).
|
||||
//
|
||||
// - ExecStartPre= — only when service.bind is "127.0.0.1" (the TCP
|
||||
// opt-in). In that case the workload binds a TCP port directly
|
||||
// (no socket), and the ExecStartPre is a placeholder that records
|
||||
// the bind (the actual bind happens in the process; the directive
|
||||
// is a no-op marker so operators can see the bind mode in the unit
|
||||
// file). When service.bind is empty (the default), the workload
|
||||
// binds the socket and no ExecStartPre is emitted for sockets.
|
||||
//
|
||||
// The socket path format is /run/orca/alloc-<alloc-id>/port-<port-name>.sock
|
||||
// where alloc-id is a PLACEHOLDER (spec.Name) — the real alloc-id is
|
||||
// assigned at submit time by the scheduler. The placeholder is
|
||||
// documented in the rendered unit via a comment so operators reading
|
||||
// the unit file understand the substitution.
|
||||
//
|
||||
// P08 is a PLAN/plumbing layer — the actual socket activation (socket
|
||||
// unit files, systemd socket-activation passing the pre-bound socket
|
||||
// fd to the process) lands in v0.10. P08 just renders the
|
||||
// RuntimeDirectory= lines and the optional TCP-bind ExecStartPre so
|
||||
// the directory exists at runtime.
|
||||
type SocketEmitter struct{}
|
||||
|
||||
// runtimeDirectoryRoot is the systemd RuntimeDirectory path root.
|
||||
// systemd joins this with the RuntimeDirectory= value to create
|
||||
// /run/orca/alloc-<id>. The leading slash is implicit in systemd
|
||||
// (RuntimeDirectory= is relative to /run).
|
||||
const runtimeDirectoryRoot = "orca"
|
||||
|
||||
// SocketPath returns the R-007 socket path for a port on the given
|
||||
// alloc-id. The alloc-id is the placeholder spec.Name when the real
|
||||
// alloc-id is not yet known (the scheduler assigns the real alloc-id
|
||||
// at submit time).
|
||||
func SocketPath(allocID, portName string) string {
|
||||
return fmt.Sprintf("/run/orca/alloc-%s/port-%s.sock", allocID, portName)
|
||||
}
|
||||
|
||||
// RenderSocketLines renders the systemd directives that implement
|
||||
// the R-007 socket plumbing for the given spec. The lines are returned
|
||||
// WITHOUT a trailing newline so the caller (the systemd emitter) can
|
||||
// append them to the [Service] block with consistent formatting.
|
||||
//
|
||||
// The returned lines are:
|
||||
//
|
||||
// - one RuntimeDirectory= line per port (so each port's socket
|
||||
// directory is created by systemd at unit start).
|
||||
// - a comment documenting the alloc-id placeholder.
|
||||
// - when service.bind is "127.0.0.1", an ExecStartPre= marker that
|
||||
// records the TCP opt-in (the actual bind is in the process).
|
||||
//
|
||||
// Returns an empty slice when the spec has no ports (no socket
|
||||
// plumbing needed — e.g. a Job or a port-less DaemonSet).
|
||||
func (SocketEmitter) RenderSocketLines(spec *jobspec.WorkloadSpec) []string {
|
||||
if spec == nil || len(spec.Ports) == 0 {
|
||||
return nil
|
||||
}
|
||||
allocID := allocIDForSocket(spec)
|
||||
var lines []string
|
||||
// One RuntimeDirectory= per port. systemd dedupes identical
|
||||
// values, but we emit one per port so the unit file is
|
||||
// self-documenting (each port maps to a directory entry).
|
||||
for _, p := range spec.Ports {
|
||||
lines = append(lines, fmt.Sprintf("RuntimeDirectory=%s/alloc-%s", runtimeDirectoryRoot, allocID))
|
||||
// Document the socket path this directory serves. systemd
|
||||
// ignores comment lines (lines starting with '#').
|
||||
lines = append(lines, fmt.Sprintf("# socket: %s", SocketPath(allocID, p.Name)))
|
||||
}
|
||||
// TCP opt-in: when service.bind is 127.0.0.1, the workload binds
|
||||
// a TCP port directly instead of the socket. We emit an
|
||||
// ExecStartPre marker so the bind mode is visible in the unit
|
||||
// file. The actual bind is in the process; the marker is a
|
||||
// no-op (echo to journald).
|
||||
if spec.Service != nil && strings.TrimSpace(spec.Service.Bind) != "" {
|
||||
if isTCPOptIn(spec.Service.Bind) {
|
||||
for _, p := range spec.Ports {
|
||||
lines = append(lines, fmt.Sprintf("ExecStartPre=/bin/echo orca: bind %s port %s (tcp, R-007 opt-in)", spec.Service.Bind, p.Name))
|
||||
}
|
||||
}
|
||||
}
|
||||
return lines
|
||||
}
|
||||
|
||||
// allocIDForSocket returns the alloc-id placeholder for the spec. The
|
||||
// real alloc-id is assigned by the scheduler at submit time; P08 uses
|
||||
// spec.Name as a deterministic placeholder so the rendered unit is
|
||||
// stable across re-renders. This mirrors the Traefik emitter's
|
||||
// allocIDFor (which uses the node hostname for the Traefik
|
||||
// dynamic-config server URL); the systemd unit is per-alloc, so
|
||||
// spec.Name is the right placeholder here.
|
||||
func allocIDForSocket(spec *jobspec.WorkloadSpec) string {
|
||||
if spec == nil || strings.TrimSpace(spec.Name) == "" {
|
||||
return "<allocID>"
|
||||
}
|
||||
return spec.Name
|
||||
}
|
||||
|
||||
// isTCPOptIn returns true when the bind value is the documented
|
||||
// 127.0.0.1 TCP opt-in (R-007). Other valid IPs (::1, etc.) are also
|
||||
// TCP opt-ins (any non-empty bind opts out of the socket default); we
|
||||
// only emit the marker for 127.0.0.1 because that is the only
|
||||
// documented opt-in per the PRD — other IPs are accepted by the
|
||||
// schema validator but are operator-specific and we do not
|
||||
// second-guess them.
|
||||
func isTCPOptIn(bind string) bool {
|
||||
return strings.TrimSpace(bind) == "127.0.0.1"
|
||||
}
|
||||
@@ -0,0 +1,265 @@
|
||||
package emitter
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"git.cloudinit.dev/coreci/orca/internal/jobspec"
|
||||
)
|
||||
|
||||
func TestSocketEmitter_RenderSocketLines_NoPorts(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Job",
|
||||
Name: "backup",
|
||||
Runtime: &jobspec.RuntimeBlock{OneOf: "process", Command: "/bin/rsync"},
|
||||
}
|
||||
lines := (SocketEmitter{}).RenderSocketLines(spec)
|
||||
if len(lines) != 0 {
|
||||
t.Errorf("got %d lines, want 0 for no ports: %v", len(lines), lines)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSocketEmitter_RenderSocketLines_NilSpec(t *testing.T) {
|
||||
lines := (SocketEmitter{}).RenderSocketLines(nil)
|
||||
if lines != nil {
|
||||
t.Errorf("nil spec should return nil, got %v", lines)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSocketEmitter_RenderSocketLines_SinglePort(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}},
|
||||
}
|
||||
lines := (SocketEmitter{}).RenderSocketLines(spec)
|
||||
// Expect: RuntimeDirectory + comment. No TCP bind (default socket).
|
||||
wantRT := "RuntimeDirectory=orca/alloc-web"
|
||||
if !contains(lines, wantRT) {
|
||||
t.Errorf("lines %v missing %q", lines, wantRT)
|
||||
}
|
||||
wantSock := "# socket: /run/orca/alloc-web/port-http.sock"
|
||||
if !contains(lines, wantSock) {
|
||||
t.Errorf("lines %v missing %q", lines, wantSock)
|
||||
}
|
||||
for _, l := range lines {
|
||||
if strings.HasPrefix(l, "ExecStartPre=") {
|
||||
t.Errorf("socket bind should not emit ExecStartPre (no TCP opt-in): %s", l)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestSocketEmitter_RenderSocketLines_MultiplePorts(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "api",
|
||||
Ports: []jobspec.PortSpec{
|
||||
{Name: "http", Port: 8080},
|
||||
{Name: "grpc", Port: 9090},
|
||||
},
|
||||
}
|
||||
lines := (SocketEmitter{}).RenderSocketLines(spec)
|
||||
// Two RuntimeDirectory lines (one per port).
|
||||
count := 0
|
||||
for _, l := range lines {
|
||||
if l == "RuntimeDirectory=orca/alloc-api" {
|
||||
count++
|
||||
}
|
||||
}
|
||||
if count != 2 {
|
||||
t.Errorf("RuntimeDirectory count = %d, want 2 (one per port)", count)
|
||||
}
|
||||
if !contains(lines, "# socket: /run/orca/alloc-api/port-http.sock") {
|
||||
t.Errorf("missing http socket comment")
|
||||
}
|
||||
if !contains(lines, "# socket: /run/orca/alloc-api/port-grpc.sock") {
|
||||
t.Errorf("missing grpc socket comment")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSocketEmitter_RenderSocketLines_TCPBind127(t *testing.T) {
|
||||
// service.bind = 127.0.0.1 → TCP opt-in → ExecStartPre marker per port.
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}},
|
||||
Service: &jobspec.ServiceBlock{Bind: "127.0.0.1"},
|
||||
}
|
||||
lines := (SocketEmitter{}).RenderSocketLines(spec)
|
||||
found := false
|
||||
for _, l := range lines {
|
||||
if strings.HasPrefix(l, "ExecStartPre=/bin/echo orca: bind 127.0.0.1 port http (tcp, R-007 opt-in)") {
|
||||
found = true
|
||||
}
|
||||
}
|
||||
if !found {
|
||||
t.Errorf("missing TCP bind ExecStartPre marker; lines: %v", lines)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSocketEmitter_RenderSocketLines_TCPBindIPv6(t *testing.T) {
|
||||
// Non-127.0.0.1 bind is accepted by schema but not the documented
|
||||
// opt-in; the marker is only emitted for 127.0.0.1. The
|
||||
// RuntimeDirectory lines are still emitted (the directory exists
|
||||
// regardless of bind mode — sockets or TCP).
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}},
|
||||
Service: &jobspec.ServiceBlock{Bind: "::1"},
|
||||
}
|
||||
lines := (SocketEmitter{}).RenderSocketLines(spec)
|
||||
for _, l := range lines {
|
||||
if strings.HasPrefix(l, "ExecStartPre=") {
|
||||
t.Errorf("::1 bind should NOT emit TCP marker (only 127.0.0.1 is documented opt-in): %s", l)
|
||||
}
|
||||
}
|
||||
if !contains(lines, "RuntimeDirectory=orca/alloc-web") {
|
||||
t.Errorf("RuntimeDirectory should still be emitted for ::1 bind")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSocketEmitter_RenderSocketLines_EmptyBindSocket(t *testing.T) {
|
||||
// Empty bind → default socket → no TCP marker, but RuntimeDirectory emitted.
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}},
|
||||
Service: &jobspec.ServiceBlock{Bind: ""},
|
||||
}
|
||||
lines := (SocketEmitter{}).RenderSocketLines(spec)
|
||||
for _, l := range lines {
|
||||
if strings.HasPrefix(l, "ExecStartPre=") {
|
||||
t.Errorf("empty bind should NOT emit TCP marker: %s", l)
|
||||
}
|
||||
}
|
||||
if !contains(lines, "RuntimeDirectory=orca/alloc-web") {
|
||||
t.Errorf("RuntimeDirectory missing for empty bind")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSocketEmitter_RenderSocketLines_NilService(t *testing.T) {
|
||||
// No service block → default socket → no TCP marker, but RuntimeDirectory emitted.
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}},
|
||||
}
|
||||
lines := (SocketEmitter{}).RenderSocketLines(spec)
|
||||
for _, l := range lines {
|
||||
if strings.HasPrefix(l, "ExecStartPre=") {
|
||||
t.Errorf("nil service should NOT emit TCP marker: %s", l)
|
||||
}
|
||||
}
|
||||
if !contains(lines, "RuntimeDirectory=orca/alloc-web") {
|
||||
t.Errorf("RuntimeDirectory missing for nil service")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSocketEmitter_SocketPath(t *testing.T) {
|
||||
got := SocketPath("alloc-123", "http")
|
||||
want := "/run/orca/alloc-alloc-123/port-http.sock"
|
||||
if got != want {
|
||||
t.Errorf("SocketPath = %q, want %q", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSocketEmitter_AllocIDPlaceholderNilSpec(t *testing.T) {
|
||||
if got := allocIDForSocket(nil); got != "<allocID>" {
|
||||
t.Errorf("allocIDForSocket(nil) = %q, want <allocID>", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSocketEmitter_AllocIDPlaceholderEmptyName(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{Name: " "}
|
||||
if got := allocIDForSocket(spec); got != "<allocID>" {
|
||||
t.Errorf("allocIDForSocket(empty name) = %q, want <allocID>", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSocketEmitter_AllocIDPlaceholderNamedSpec(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{Name: "web"}
|
||||
if got := allocIDForSocket(spec); got != "web" {
|
||||
t.Errorf("allocIDForSocket(web) = %q, want web", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSocketEmitter_SocketPathPlaceholder(t *testing.T) {
|
||||
got := SocketPath("<allocID>", "grpc")
|
||||
want := "/run/orca/alloc-<allocID>/port-grpc.sock"
|
||||
if got != want {
|
||||
t.Errorf("SocketPath = %q, want %q", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSystemdEmitter_IntegratesSocketLines(t *testing.T) {
|
||||
// End-to-end: the systemd unit for a Service with ports contains
|
||||
// the RuntimeDirectory line emitted by the SocketEmitter.
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Runtime: &jobspec.RuntimeBlock{OneOf: "process", Command: "/bin/httpd"},
|
||||
Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}},
|
||||
}
|
||||
files, err := SystemdEmitter{}.Render(spec, &Node{})
|
||||
if err != nil {
|
||||
t.Fatalf("Render: %v", err)
|
||||
}
|
||||
c := files[0].Content
|
||||
if !strings.Contains(c, "RuntimeDirectory=orca/alloc-web\n") {
|
||||
t.Errorf("unit missing RuntimeDirectory line\n%s", c)
|
||||
}
|
||||
if !strings.Contains(c, "# socket: /run/orca/alloc-web/port-http.sock\n") {
|
||||
t.Errorf("unit missing socket path comment\n%s", c)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSystemdEmitter_IntegratesSocketLinesTCPBind(t *testing.T) {
|
||||
// When service.bind = 127.0.0.1, the unit contains the ExecStartPre
|
||||
// TCP-bind marker.
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Runtime: &jobspec.RuntimeBlock{OneOf: "process", Command: "/bin/httpd"},
|
||||
Ports: []jobspec.PortSpec{{Name: "http", Port: 8080}},
|
||||
Service: &jobspec.ServiceBlock{Bind: "127.0.0.1"},
|
||||
}
|
||||
files, err := SystemdEmitter{}.Render(spec, &Node{})
|
||||
if err != nil {
|
||||
t.Fatalf("Render: %v", err)
|
||||
}
|
||||
c := files[0].Content
|
||||
if !strings.Contains(c, "ExecStartPre=/bin/echo orca: bind 127.0.0.1 port http (tcp, R-007 opt-in)\n") {
|
||||
t.Errorf("unit missing TCP bind ExecStartPre marker\n%s", c)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSystemdEmitter_NoSocketLinesForPortlessSpec(t *testing.T) {
|
||||
// A Job with no ports → no RuntimeDirectory line in the unit.
|
||||
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)
|
||||
}
|
||||
c := files[0].Content
|
||||
if strings.Contains(c, "RuntimeDirectory=") {
|
||||
t.Errorf("portless spec should not emit RuntimeDirectory\n%s", c)
|
||||
}
|
||||
if strings.Contains(c, "# socket:") {
|
||||
t.Errorf("portless spec should not emit socket comment\n%s", c)
|
||||
}
|
||||
}
|
||||
|
||||
// contains reports whether the slice contains the string s.
|
||||
func contains(lines []string, s string) bool {
|
||||
for _, l := range lines {
|
||||
if l == s {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
+88
-20
@@ -8,21 +8,30 @@ import (
|
||||
"git.cloudinit.dev/coreci/orca/internal/jobspec"
|
||||
)
|
||||
|
||||
// SystemdEmitter is a stub Emitter implementation for the "process"
|
||||
// runtime. It renders a minimal systemd unit file for the workload.
|
||||
// SystemdEmitter is the Emitter implementation for the "process"
|
||||
// runtime. It renders the systemd unit file for the workload,
|
||||
// including the lifecycle hooks (P04) and the R-007 socket plumbing
|
||||
// (P08).
|
||||
//
|
||||
// This is a STUB — the full systemd emitter (with lifecycle hooks,
|
||||
// sockets, EnvironmentFile, LoadCredential) lands in later phases:
|
||||
// Lifecycle hooks map to systemd semantics (PRD §10.1):
|
||||
//
|
||||
// - lifecycle.pre_stop → ExecStop= (the command run on stop; systemd
|
||||
// runs ExecStop, then kills the main process after the deadline).
|
||||
// - lifecycle.post_start → ExecStartPost= (runs after the main
|
||||
// process starts).
|
||||
//
|
||||
// systemd has no ExecStartPre equivalent for a "pre_start" hook; the
|
||||
// spec does not define pre_start (only pre_stop and post_start per
|
||||
// PRD §10.1), so no mapping is needed.
|
||||
//
|
||||
// The unit name carries the `orca-v1-` prefix per the dual-write
|
||||
// window (REQ-090) so the v0.9 SSH-push path does not collide with the
|
||||
// v0.8 daemon's `orca-<job>.service` units during the migration
|
||||
// window.
|
||||
//
|
||||
// Later phases extend this emitter:
|
||||
//
|
||||
// - P04: lifecycle hooks (ExecStop, ExecStartPre/Post, timeouts)
|
||||
// - P08: socket plumbing (R-007)
|
||||
// - v0.10-P03: secrets via EnvironmentFile= + LoadCredential=
|
||||
//
|
||||
// P0c ships only the minimal [Service]\nExecStart=... shape to prove
|
||||
// the Emitter interface end-to-end. The unit name carries the
|
||||
// `orca-v1-` prefix per the dual-write window (REQ-090) so the v0.9
|
||||
// SSH-push path does not collide with the v0.8 daemon's
|
||||
// `orca-<job>.service` units during the migration window.
|
||||
type SystemdEmitter struct{}
|
||||
|
||||
// unitNamePrefix is the v0.9 SSH-push unit-name prefix. The v0.8
|
||||
@@ -32,15 +41,28 @@ type SystemdEmitter struct{}
|
||||
// updating the dual-write window contract.
|
||||
const unitNamePrefix = "orca-v1-"
|
||||
|
||||
// Render renders a minimal systemd unit file for a process-runtime
|
||||
// workload. The unit name is `/etc/systemd/system/<unitNamePrefix><spec.Name>.service`
|
||||
// and the content is a minimal `[Service]` block with the runtime
|
||||
// command as ExecStart. Mode is 0644 (the lead applier chmods after
|
||||
// atomic rename).
|
||||
// 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:
|
||||
//
|
||||
// [Service]
|
||||
// ExecStart=<runtime command>
|
||||
// ExecStartPost=<post_start command 1>
|
||||
// ExecStartPost=<post_start command 2>
|
||||
// ExecStop=<pre_stop command 1>
|
||||
// ExecStop=<pre_stop command 2>
|
||||
// RuntimeDirectory=orca/alloc-<alloc-id>
|
||||
// # socket: /run/orca/alloc-<alloc-id>/port-<name>.sock
|
||||
// ExecStartPre=/bin/echo orca: bind 127.0.0.1 port <name> (tcp, R-007 opt-in)
|
||||
//
|
||||
// Returns an error if the spec is nil, the spec is missing its name,
|
||||
// or the runtime command is empty (a workload with no command has
|
||||
// nothing to ExecStart).
|
||||
// the runtime block is nil, or the runtime command is empty (a
|
||||
// workload with no command has nothing to ExecStart).
|
||||
func (SystemdEmitter) Render(spec *jobspec.WorkloadSpec, node *Node) ([]File, error) {
|
||||
if spec == nil {
|
||||
return nil, errors.New("emitter/systemd: spec is nil")
|
||||
@@ -55,6 +77,52 @@ func (SystemdEmitter) Render(spec *jobspec.WorkloadSpec, node *Node) ([]File, er
|
||||
return nil, errors.New("emitter/systemd: runtime command is empty")
|
||||
}
|
||||
path := fmt.Sprintf("/etc/systemd/system/%s%s.service", unitNamePrefix, spec.Name)
|
||||
content := fmt.Sprintf("[Service]\nExecStart=%s\n", spec.Runtime.Command)
|
||||
content := renderSystemdUnit(spec)
|
||||
return []File{{Path: path, Content: content, Mode: "0644"}}, nil
|
||||
}
|
||||
|
||||
// renderSystemdUnit renders the full [Service] block for the spec,
|
||||
// including ExecStart, lifecycle hooks (ExecStartPost, ExecStop), and
|
||||
// the R-007 socket-plumbing lines (RuntimeDirectory=, optional
|
||||
// TCP-bind ExecStartPre). The output is a single string with a
|
||||
// trailing newline per line.
|
||||
func renderSystemdUnit(spec *jobspec.WorkloadSpec) string {
|
||||
var b strings.Builder
|
||||
b.WriteString("[Service]\n")
|
||||
b.WriteString(fmt.Sprintf("ExecStart=%s\n", spec.Runtime.Command))
|
||||
// Lifecycle: post_start → ExecStartPost (runs after start).
|
||||
for _, cmd := range lifecyclePostStart(spec) {
|
||||
b.WriteString(fmt.Sprintf("ExecStartPost=%s\n", cmd))
|
||||
}
|
||||
// Lifecycle: pre_stop → ExecStop (runs before the process is killed).
|
||||
for _, cmd := range lifecyclePreStop(spec) {
|
||||
b.WriteString(fmt.Sprintf("ExecStop=%s\n", cmd))
|
||||
}
|
||||
// R-007 socket plumbing: RuntimeDirectory= per port + optional
|
||||
// TCP-bind ExecStartPre.
|
||||
for _, line := range (SocketEmitter{}).RenderSocketLines(spec) {
|
||||
b.WriteString(line)
|
||||
b.WriteString("\n")
|
||||
}
|
||||
return b.String()
|
||||
}
|
||||
|
||||
// lifecyclePostStart returns the post_start lifecycle commands for
|
||||
// the spec, or nil when the spec has no lifecycle block or no
|
||||
// post_start commands.
|
||||
func lifecyclePostStart(spec *jobspec.WorkloadSpec) []string {
|
||||
if spec.Lifecycle == nil {
|
||||
return nil
|
||||
}
|
||||
return spec.Lifecycle.PostStart
|
||||
}
|
||||
|
||||
// lifecyclePreStop returns the pre_stop lifecycle commands for the
|
||||
// spec, or nil when the spec has no lifecycle block or no pre_stop
|
||||
// commands.
|
||||
func lifecyclePreStop(spec *jobspec.WorkloadSpec) []string {
|
||||
if spec.Lifecycle == nil {
|
||||
return nil
|
||||
}
|
||||
return spec.Lifecycle.PreStop
|
||||
}
|
||||
|
||||
@@ -0,0 +1,172 @@
|
||||
package emitter
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"git.cloudinit.dev/coreci/orca/internal/jobspec"
|
||||
)
|
||||
|
||||
func TestSystemdEmitter_LifecyclePostStart(t *testing.T) {
|
||||
// lifecycle.post_start → ExecStartPost (one line per command).
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Runtime: &jobspec.RuntimeBlock{OneOf: "process", Command: "/usr/local/bin/httpd -f"},
|
||||
Lifecycle: &jobspec.LifecycleBlock{
|
||||
PostStart: []string{"/usr/bin/sleep 1", "/usr/bin/curl localhost/healthz"},
|
||||
},
|
||||
}
|
||||
files, err := SystemdEmitter{}.Render(spec, &Node{Hostname: "n1"})
|
||||
if err != nil {
|
||||
t.Fatalf("Render: %v", err)
|
||||
}
|
||||
c := files[0].Content
|
||||
if !strings.Contains(c, "ExecStartPost=/usr/bin/sleep 1\n") {
|
||||
t.Errorf("missing ExecStartPost for sleep 1\n%s", c)
|
||||
}
|
||||
if !strings.Contains(c, "ExecStartPost=/usr/bin/curl localhost/healthz\n") {
|
||||
t.Errorf("missing ExecStartPost for curl\n%s", c)
|
||||
}
|
||||
// ExecStart must still be present.
|
||||
if !strings.Contains(c, "ExecStart=/usr/local/bin/httpd -f\n") {
|
||||
t.Errorf("missing ExecStart\n%s", c)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSystemdEmitter_LifecyclePreStop(t *testing.T) {
|
||||
// lifecycle.pre_stop → ExecStop (one line per command).
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Runtime: &jobspec.RuntimeBlock{OneOf: "process", Command: "/usr/local/bin/httpd -f"},
|
||||
Lifecycle: &jobspec.LifecycleBlock{
|
||||
PreStop: []string{"/usr/local/bin/httpd -graceful", "/usr/bin/sleep 5"},
|
||||
},
|
||||
}
|
||||
files, err := SystemdEmitter{}.Render(spec, &Node{Hostname: "n1"})
|
||||
if err != nil {
|
||||
t.Fatalf("Render: %v", err)
|
||||
}
|
||||
c := files[0].Content
|
||||
if !strings.Contains(c, "ExecStop=/usr/local/bin/httpd -graceful\n") {
|
||||
t.Errorf("missing ExecStop for graceful\n%s", c)
|
||||
}
|
||||
if !strings.Contains(c, "ExecStop=/usr/bin/sleep 5\n") {
|
||||
t.Errorf("missing ExecStop for sleep 5\n%s", c)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSystemdEmitter_LifecycleBoth(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Runtime: &jobspec.RuntimeBlock{OneOf: "process", Command: "/bin/httpd"},
|
||||
Lifecycle: &jobspec.LifecycleBlock{
|
||||
PostStart: []string{"/bin/after-start"},
|
||||
PreStop: []string{"/bin/before-stop"},
|
||||
},
|
||||
}
|
||||
files, err := SystemdEmitter{}.Render(spec, &Node{})
|
||||
if err != nil {
|
||||
t.Fatalf("Render: %v", err)
|
||||
}
|
||||
c := files[0].Content
|
||||
// ExecStartPost must appear before ExecStop (post_start runs after
|
||||
// start; pre_stop runs before stop — the order in the unit file
|
||||
// reflects the lifecycle order).
|
||||
startIdx := strings.Index(c, "ExecStart=")
|
||||
postIdx := strings.Index(c, "ExecStartPost=")
|
||||
stopIdx := strings.Index(c, "ExecStop=")
|
||||
if startIdx < 0 || postIdx < 0 || stopIdx < 0 {
|
||||
t.Fatalf("missing one of ExecStart/ExecStartPost/ExecStop\n%s", c)
|
||||
}
|
||||
if !(startIdx < postIdx && postIdx < stopIdx) {
|
||||
t.Errorf("expected order ExecStart < ExecStartPost < ExecStop\n%s", c)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSystemdEmitter_LifecycleNilOmitsDirectives(t *testing.T) {
|
||||
// No lifecycle block → no ExecStartPost / ExecStop lines.
|
||||
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)
|
||||
}
|
||||
c := files[0].Content
|
||||
if strings.Contains(c, "ExecStartPost=") {
|
||||
t.Errorf("ExecStartPost should be omitted when no lifecycle\n%s", c)
|
||||
}
|
||||
if strings.Contains(c, "ExecStop=") {
|
||||
t.Errorf("ExecStop should be omitted when no lifecycle\n%s", c)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSystemdEmitter_LifecycleEmptyListsOmitted(t *testing.T) {
|
||||
// Lifecycle block present but empty lists → no ExecStartPost / ExecStop.
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Job",
|
||||
Name: "x",
|
||||
Runtime: &jobspec.RuntimeBlock{OneOf: "process", Command: "/bin/x"},
|
||||
Lifecycle: &jobspec.LifecycleBlock{},
|
||||
}
|
||||
files, err := SystemdEmitter{}.Render(spec, &Node{})
|
||||
if err != nil {
|
||||
t.Fatalf("Render: %v", err)
|
||||
}
|
||||
c := files[0].Content
|
||||
if strings.Contains(c, "ExecStartPost=") {
|
||||
t.Errorf("ExecStartPost should be omitted for empty PostStart\n%s", c)
|
||||
}
|
||||
if strings.Contains(c, "ExecStop=") {
|
||||
t.Errorf("ExecStop should be omitted for empty PreStop\n%s", c)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSystemdEmitter_LifecycleOnlyPostStart(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Job",
|
||||
Name: "x",
|
||||
Runtime: &jobspec.RuntimeBlock{OneOf: "process", Command: "/bin/x"},
|
||||
Lifecycle: &jobspec.LifecycleBlock{
|
||||
PostStart: []string{"/bin/notify-up"},
|
||||
},
|
||||
}
|
||||
files, err := SystemdEmitter{}.Render(spec, &Node{})
|
||||
if err != nil {
|
||||
t.Fatalf("Render: %v", err)
|
||||
}
|
||||
c := files[0].Content
|
||||
if !strings.Contains(c, "ExecStartPost=/bin/notify-up\n") {
|
||||
t.Errorf("missing ExecStartPost\n%s", c)
|
||||
}
|
||||
if strings.Contains(c, "ExecStop=") {
|
||||
t.Errorf("ExecStop should be omitted when only PostStart set\n%s", c)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSystemdEmitter_LifecycleOnlyPreStop(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Job",
|
||||
Name: "x",
|
||||
Runtime: &jobspec.RuntimeBlock{OneOf: "process", Command: "/bin/x"},
|
||||
Lifecycle: &jobspec.LifecycleBlock{
|
||||
PreStop: []string{"/bin/notify-down"},
|
||||
},
|
||||
}
|
||||
files, err := SystemdEmitter{}.Render(spec, &Node{})
|
||||
if err != nil {
|
||||
t.Fatalf("Render: %v", err)
|
||||
}
|
||||
c := files[0].Content
|
||||
if !strings.Contains(c, "ExecStop=/bin/notify-down\n") {
|
||||
t.Errorf("missing ExecStop\n%s", c)
|
||||
}
|
||||
if strings.Contains(c, "ExecStartPost=") {
|
||||
t.Errorf("ExecStartPost should be omitted when only PreStop set\n%s", c)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,237 @@
|
||||
package emitter
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strconv"
|
||||
"strings"
|
||||
|
||||
"git.cloudinit.dev/coreci/orca/internal/jobspec"
|
||||
)
|
||||
|
||||
// UpdatePlan is the computed update sequence for a Service (P03). It is
|
||||
// a PLAN, not an execution — the transactional execution lands in
|
||||
// v0.10-P10. Each step describes a discrete action the executor takes:
|
||||
// start a set of allocs (Action="start"), wait for them to become
|
||||
// healthy (WaitForHealthy=true), or cutover from old to new
|
||||
// (Action="cutover" for blue-green). The Allocs field carries
|
||||
// placeholder alloc names of the form "<spec.Name>-<index>" where
|
||||
// index is 1-based (the scheduler assigns the real alloc-id at submit
|
||||
// time; P03 uses spec.Name as a placeholder per the socket layer
|
||||
// contract — see SocketEmitter).
|
||||
type UpdatePlan struct {
|
||||
Steps []UpdateStep
|
||||
}
|
||||
|
||||
// UpdateStep is a single step in an UpdatePlan. Action is one of
|
||||
// "start", "wait", "cutover", "promote". Allocs is the list of
|
||||
// placeholder alloc names the step applies to. WaitForHealthy is true
|
||||
// when the executor must wait for the allocs in this step to pass
|
||||
// their health check before proceeding to the next step (driven by
|
||||
// min_healthy_time / healthy_deadline on the spec, which the executor
|
||||
// — not the plan — enforces).
|
||||
type UpdateStep struct {
|
||||
Action string
|
||||
Allocs []string
|
||||
WaitForHealthy bool
|
||||
}
|
||||
|
||||
// maxParallelFor returns the effective max_parallel for the spec,
|
||||
// defaulting to 1 when unset (0) and clamping to count (the validator
|
||||
// already rejects out-of-range values; this is a defensive clamp for
|
||||
// direct callers that bypass the validator).
|
||||
func maxParallelFor(spec *jobspec.WorkloadSpec) int {
|
||||
if spec.Update == nil {
|
||||
return 1
|
||||
}
|
||||
if spec.Update.MaxParallel < 1 {
|
||||
return 1
|
||||
}
|
||||
if spec.Count > 0 && spec.Update.MaxParallel > spec.Count {
|
||||
return spec.Count
|
||||
}
|
||||
return spec.Update.MaxParallel
|
||||
}
|
||||
|
||||
// allocName returns the placeholder alloc name for index i (1-based).
|
||||
// The real alloc-id is assigned by the scheduler at submit time; P03
|
||||
// uses spec.Name as the placeholder per the socket-layer contract.
|
||||
func allocName(spec *jobspec.WorkloadSpec, i int) string {
|
||||
return fmt.Sprintf("%s-%d", spec.Name, i)
|
||||
}
|
||||
|
||||
// allAllocs returns the placeholder alloc names for the full count of
|
||||
// the spec (1..count).
|
||||
func allAllocs(spec *jobspec.WorkloadSpec) []string {
|
||||
out := make([]string, 0, spec.Count)
|
||||
for i := 1; i <= spec.Count; i++ {
|
||||
out = append(out, allocName(spec, i))
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// canaryCount returns the integer canary count for the spec. The
|
||||
// canary field accepts an integer count or a percentage ("<n>%"). For
|
||||
// a percentage, the count is ceil(count * n / 100) with a minimum of 1
|
||||
// when n > 0 (a 10% canary of a 3-replica service is 1 alloc, not 0).
|
||||
// When the canary field is empty, the default is 1 (a single canary
|
||||
// alloc — the smallest meaningful canary).
|
||||
func canaryCount(spec *jobspec.WorkloadSpec) int {
|
||||
if spec.Update == nil {
|
||||
return 1
|
||||
}
|
||||
c := strings.TrimSpace(spec.Update.Canary)
|
||||
if c == "" {
|
||||
return 1
|
||||
}
|
||||
if strings.HasSuffix(c, "%") {
|
||||
n, err := strconv.Atoi(strings.TrimSpace(strings.TrimSuffix(c, "%")))
|
||||
if err != nil || n <= 0 {
|
||||
return 1
|
||||
}
|
||||
allocs := spec.Count * n / 100
|
||||
if allocs < 1 {
|
||||
allocs = 1
|
||||
}
|
||||
return allocs
|
||||
}
|
||||
n, err := strconv.Atoi(c)
|
||||
if err != nil || n < 1 {
|
||||
return 1
|
||||
}
|
||||
if spec.Count > 0 && n > spec.Count {
|
||||
return spec.Count
|
||||
}
|
||||
return n
|
||||
}
|
||||
|
||||
// RenderUpdatePlan computes the rolling/canary/blue-green update
|
||||
// sequence for a Service spec. Returns an *UpdatePlan describing the
|
||||
// steps; the actual transactional execution lands in v0.10-P10.
|
||||
//
|
||||
// The three strategies:
|
||||
//
|
||||
// - rolling: allocs are started in batches of max_parallel. Each
|
||||
// batch waits for healthy before the next batch starts. This is
|
||||
// the simplest strategy and the default for stateless services.
|
||||
//
|
||||
// - canary: a single canary alloc (or N per the canary field) is
|
||||
// started first and waits for healthy. After the canary is
|
||||
// healthy, the plan emits a "promote" step (manual or auto per
|
||||
// auto_promote); the remaining allocs are then started in
|
||||
// max_parallel batches.
|
||||
//
|
||||
// - blue-green: all new allocs are started in parallel (a single
|
||||
// "start" step with the full count). After they are healthy, a
|
||||
// "cutover" step swaps traffic from the old allocs to the new
|
||||
// ones. The old allocs are then stopped (the stop is implicit in
|
||||
// the cutover step for the plan; v0.10-P10 makes it explicit).
|
||||
//
|
||||
// Returns an error if the spec is nil, the update block is nil, or
|
||||
// the strategy is unknown (the validator should have caught these,
|
||||
// but RenderUpdatePlan is defensive — emitters are called from
|
||||
// render paths that may bypass the schema validator).
|
||||
func RenderUpdatePlan(spec *jobspec.WorkloadSpec) (*UpdatePlan, error) {
|
||||
if spec == nil {
|
||||
return nil, fmt.Errorf("emitter/update: spec is nil")
|
||||
}
|
||||
if spec.Update == nil {
|
||||
return nil, fmt.Errorf("emitter/update: update block is nil")
|
||||
}
|
||||
if spec.Count < 1 {
|
||||
return nil, fmt.Errorf("emitter/update: count must be ≥ 1, got %d", spec.Count)
|
||||
}
|
||||
switch spec.Update.Strategy {
|
||||
case "rolling":
|
||||
return renderRollingPlan(spec), nil
|
||||
case "canary":
|
||||
return renderCanaryPlan(spec), nil
|
||||
case "blue-green":
|
||||
return renderBlueGreenPlan(spec), nil
|
||||
default:
|
||||
return nil, fmt.Errorf("emitter/update: unknown strategy %q (want rolling, canary, or blue-green)", spec.Update.Strategy)
|
||||
}
|
||||
}
|
||||
|
||||
// renderRollingPlan emits the rolling-update plan: allocs in batches
|
||||
// of max_parallel, each batch waiting for healthy before the next.
|
||||
func renderRollingPlan(spec *jobspec.WorkloadSpec) *UpdatePlan {
|
||||
plan := &UpdatePlan{}
|
||||
batch := maxParallelFor(spec)
|
||||
allocs := allAllocs(spec)
|
||||
for i := 0; i < len(allocs); i += batch {
|
||||
end := i + batch
|
||||
if end > len(allocs) {
|
||||
end = len(allocs)
|
||||
}
|
||||
plan.Steps = append(plan.Steps, UpdateStep{
|
||||
Action: "start",
|
||||
Allocs: allocs[i:end],
|
||||
WaitForHealthy: true,
|
||||
})
|
||||
}
|
||||
return plan
|
||||
}
|
||||
|
||||
// renderCanaryPlan emits the canary-update plan: a canary batch first
|
||||
// (size per the canary field, default 1), a "promote" step, then the
|
||||
// remaining allocs in max_parallel batches.
|
||||
func renderCanaryPlan(spec *jobspec.WorkloadSpec) *UpdatePlan {
|
||||
plan := &UpdatePlan{}
|
||||
allocs := allAllocs(spec)
|
||||
canary := canaryCount(spec)
|
||||
if canary > len(allocs) {
|
||||
canary = len(allocs)
|
||||
}
|
||||
if canary < 1 {
|
||||
canary = 1
|
||||
}
|
||||
// Step 1: start the canary alloc(s) and wait for healthy.
|
||||
plan.Steps = append(plan.Steps, UpdateStep{
|
||||
Action: "start",
|
||||
Allocs: allocs[:canary],
|
||||
WaitForHealthy: true,
|
||||
})
|
||||
// Step 2: promote (manual or auto per auto_promote).
|
||||
plan.Steps = append(plan.Steps, UpdateStep{
|
||||
Action: "promote",
|
||||
Allocs: allocs[:canary],
|
||||
})
|
||||
// Step 3+: remaining allocs in max_parallel batches.
|
||||
batch := maxParallelFor(spec)
|
||||
remaining := allocs[canary:]
|
||||
for i := 0; i < len(remaining); i += batch {
|
||||
end := i + batch
|
||||
if end > len(remaining) {
|
||||
end = len(remaining)
|
||||
}
|
||||
plan.Steps = append(plan.Steps, UpdateStep{
|
||||
Action: "start",
|
||||
Allocs: remaining[i:end],
|
||||
WaitForHealthy: true,
|
||||
})
|
||||
}
|
||||
return plan
|
||||
}
|
||||
|
||||
// renderBlueGreenPlan emits the blue-green update plan: all new allocs
|
||||
// start in parallel, wait for healthy, then cutover (swap traffic).
|
||||
func renderBlueGreenPlan(spec *jobspec.WorkloadSpec) *UpdatePlan {
|
||||
plan := &UpdatePlan{}
|
||||
allocs := allAllocs(spec)
|
||||
// Step 1: start ALL new allocs in parallel (blue-green does not
|
||||
// batch — the new fleet stands up alongside the old).
|
||||
plan.Steps = append(plan.Steps, UpdateStep{
|
||||
Action: "start",
|
||||
Allocs: allocs,
|
||||
WaitForHealthy: true,
|
||||
})
|
||||
// Step 2: cutover — swap traffic from old to new. The old allocs
|
||||
// are stopped implicitly as part of the cutover (v0.10-P10 makes
|
||||
// the stop explicit in the transactional plane).
|
||||
plan.Steps = append(plan.Steps, UpdateStep{
|
||||
Action: "cutover",
|
||||
Allocs: allocs,
|
||||
WaitForHealthy: false,
|
||||
})
|
||||
return plan
|
||||
}
|
||||
@@ -0,0 +1,396 @@
|
||||
package emitter
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"git.cloudinit.dev/coreci/orca/internal/jobspec"
|
||||
)
|
||||
|
||||
func TestRenderUpdatePlan_RollingBatches(t *testing.T) {
|
||||
// count=4, max_parallel=2 → 2 batches of 2.
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 4,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "rolling",
|
||||
MaxParallel: 2,
|
||||
},
|
||||
}
|
||||
plan, err := RenderUpdatePlan(spec)
|
||||
if err != nil {
|
||||
t.Fatalf("RenderUpdatePlan: %v", err)
|
||||
}
|
||||
if len(plan.Steps) != 2 {
|
||||
t.Fatalf("got %d steps, want 2", len(plan.Steps))
|
||||
}
|
||||
for i, s := range plan.Steps {
|
||||
if s.Action != "start" {
|
||||
t.Errorf("step %d action = %q, want start", i, s.Action)
|
||||
}
|
||||
if !s.WaitForHealthy {
|
||||
t.Errorf("step %d WaitForHealthy = false, want true", i)
|
||||
}
|
||||
if len(s.Allocs) != 2 {
|
||||
t.Errorf("step %d allocs = %d, want 2", i, len(s.Allocs))
|
||||
}
|
||||
}
|
||||
if plan.Steps[0].Allocs[0] != "web-1" || plan.Steps[0].Allocs[1] != "web-2" {
|
||||
t.Errorf("step 0 allocs = %v, want [web-1 web-2]", plan.Steps[0].Allocs)
|
||||
}
|
||||
if plan.Steps[1].Allocs[0] != "web-3" || plan.Steps[1].Allocs[1] != "web-4" {
|
||||
t.Errorf("step 1 allocs = %v, want [web-3 web-4]", plan.Steps[1].Allocs)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderUpdatePlan_RollingUnevenBatches(t *testing.T) {
|
||||
// count=5, max_parallel=2 → 3 batches: 2, 2, 1.
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 5,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "rolling",
|
||||
MaxParallel: 2,
|
||||
},
|
||||
}
|
||||
plan, err := RenderUpdatePlan(spec)
|
||||
if err != nil {
|
||||
t.Fatalf("RenderUpdatePlan: %v", err)
|
||||
}
|
||||
if len(plan.Steps) != 3 {
|
||||
t.Fatalf("got %d steps, want 3", len(plan.Steps))
|
||||
}
|
||||
if len(plan.Steps[2].Allocs) != 1 {
|
||||
t.Errorf("step 2 allocs = %d, want 1 (remainder)", len(plan.Steps[2].Allocs))
|
||||
}
|
||||
if plan.Steps[2].Allocs[0] != "web-5" {
|
||||
t.Errorf("step 2 allocs = %v, want [web-5]", plan.Steps[2].Allocs)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderUpdatePlan_RollingMaxParallelUnset(t *testing.T) {
|
||||
// max_parallel unset (0) → default 1 → 4 batches of 1.
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 4,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "rolling",
|
||||
},
|
||||
}
|
||||
plan, err := RenderUpdatePlan(spec)
|
||||
if err != nil {
|
||||
t.Fatalf("RenderUpdatePlan: %v", err)
|
||||
}
|
||||
if len(plan.Steps) != 4 {
|
||||
t.Fatalf("got %d steps, want 4 (one per alloc, batch=1)", len(plan.Steps))
|
||||
}
|
||||
for _, s := range plan.Steps {
|
||||
if len(s.Allocs) != 1 {
|
||||
t.Errorf("allocs = %d, want 1", len(s.Allocs))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderUpdatePlan_CanaryDefaultOne(t *testing.T) {
|
||||
// canary unset → default 1 canary alloc, then 3 in batches of 2.
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 4,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "canary",
|
||||
MaxParallel: 2,
|
||||
},
|
||||
}
|
||||
plan, err := RenderUpdatePlan(spec)
|
||||
if err != nil {
|
||||
t.Fatalf("RenderUpdatePlan: %v", err)
|
||||
}
|
||||
// Step 0: canary (1 alloc). Step 1: promote. Steps 2..: remaining.
|
||||
if plan.Steps[0].Action != "start" || len(plan.Steps[0].Allocs) != 1 {
|
||||
t.Errorf("step 0 = %+v, want canary start with 1 alloc", plan.Steps[0])
|
||||
}
|
||||
if plan.Steps[0].Allocs[0] != "web-1" {
|
||||
t.Errorf("canary alloc = %q, want web-1", plan.Steps[0].Allocs[0])
|
||||
}
|
||||
if !plan.Steps[0].WaitForHealthy {
|
||||
t.Error("canary step should wait for healthy")
|
||||
}
|
||||
if plan.Steps[1].Action != "promote" {
|
||||
t.Errorf("step 1 action = %q, want promote", plan.Steps[1].Action)
|
||||
}
|
||||
// Remaining: web-2, web-3, web-4 in batches of 2 → [web-2,web-3], [web-4].
|
||||
if len(plan.Steps) != 4 {
|
||||
t.Fatalf("got %d steps, want 4 (canary + promote + 2 batches)", len(plan.Steps))
|
||||
}
|
||||
if len(plan.Steps[2].Allocs) != 2 || plan.Steps[2].Allocs[0] != "web-2" {
|
||||
t.Errorf("step 2 = %v, want [web-2 web-3]", plan.Steps[2].Allocs)
|
||||
}
|
||||
if len(plan.Steps[3].Allocs) != 1 || plan.Steps[3].Allocs[0] != "web-4" {
|
||||
t.Errorf("step 3 = %v, want [web-4]", plan.Steps[3].Allocs)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderUpdatePlan_CanaryPercent(t *testing.T) {
|
||||
// count=10, canary=20% → 2 canary allocs, then 8 in batches of 3.
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 10,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "canary",
|
||||
MaxParallel: 3,
|
||||
Canary: "20%",
|
||||
},
|
||||
}
|
||||
plan, err := RenderUpdatePlan(spec)
|
||||
if err != nil {
|
||||
t.Fatalf("RenderUpdatePlan: %v", err)
|
||||
}
|
||||
if len(plan.Steps[0].Allocs) != 2 {
|
||||
t.Errorf("canary step allocs = %d, want 2 (20%% of 10)", len(plan.Steps[0].Allocs))
|
||||
}
|
||||
// Remaining 8 in batches of 3 → ceil(8/3)=3 batches.
|
||||
// Steps: canary, promote, batch(3), batch(3), batch(2) = 5 steps.
|
||||
if len(plan.Steps) != 5 {
|
||||
t.Fatalf("got %d steps, want 5", len(plan.Steps))
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderUpdatePlan_CanaryIntegerCount(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 4,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "canary",
|
||||
Canary: "2",
|
||||
},
|
||||
}
|
||||
plan, err := RenderUpdatePlan(spec)
|
||||
if err != nil {
|
||||
t.Fatalf("RenderUpdatePlan: %v", err)
|
||||
}
|
||||
if len(plan.Steps[0].Allocs) != 2 {
|
||||
t.Errorf("canary step allocs = %d, want 2", len(plan.Steps[0].Allocs))
|
||||
}
|
||||
if plan.Steps[0].Allocs[0] != "web-1" || plan.Steps[0].Allocs[1] != "web-2" {
|
||||
t.Errorf("canary allocs = %v, want [web-1 web-2]", plan.Steps[0].Allocs)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderUpdatePlan_CanaryFullCount(t *testing.T) {
|
||||
// canary == count → no remaining allocs after canary.
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 3,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "canary",
|
||||
Canary: "3",
|
||||
},
|
||||
}
|
||||
plan, err := RenderUpdatePlan(spec)
|
||||
if err != nil {
|
||||
t.Fatalf("RenderUpdatePlan: %v", err)
|
||||
}
|
||||
// Steps: canary start (3), promote. No remaining batches.
|
||||
if len(plan.Steps) != 2 {
|
||||
t.Fatalf("got %d steps, want 2 (canary + promote, no remainder)", len(plan.Steps))
|
||||
}
|
||||
if plan.Steps[1].Action != "promote" {
|
||||
t.Errorf("step 1 action = %q, want promote", plan.Steps[1].Action)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderUpdatePlan_BlueGreen(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 4,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "blue-green",
|
||||
},
|
||||
}
|
||||
plan, err := RenderUpdatePlan(spec)
|
||||
if err != nil {
|
||||
t.Fatalf("RenderUpdatePlan: %v", err)
|
||||
}
|
||||
// Step 0: start all 4 in parallel. Step 1: cutover.
|
||||
if len(plan.Steps) != 2 {
|
||||
t.Fatalf("got %d steps, want 2", len(plan.Steps))
|
||||
}
|
||||
if plan.Steps[0].Action != "start" {
|
||||
t.Errorf("step 0 action = %q, want start", plan.Steps[0].Action)
|
||||
}
|
||||
if len(plan.Steps[0].Allocs) != 4 {
|
||||
t.Errorf("step 0 allocs = %d, want 4 (all new in parallel)", len(plan.Steps[0].Allocs))
|
||||
}
|
||||
if !plan.Steps[0].WaitForHealthy {
|
||||
t.Error("blue-green start step should wait for healthy")
|
||||
}
|
||||
if plan.Steps[1].Action != "cutover" {
|
||||
t.Errorf("step 1 action = %q, want cutover", plan.Steps[1].Action)
|
||||
}
|
||||
if plan.Steps[1].WaitForHealthy {
|
||||
t.Error("cutover step should NOT wait for healthy (already healthy)")
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderUpdatePlan_BlueGreenAllocs(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "api",
|
||||
Count: 3,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "blue-green",
|
||||
},
|
||||
}
|
||||
plan, err := RenderUpdatePlan(spec)
|
||||
if err != nil {
|
||||
t.Fatalf("RenderUpdatePlan: %v", err)
|
||||
}
|
||||
want := []string{"api-1", "api-2", "api-3"}
|
||||
if len(plan.Steps[0].Allocs) != 3 {
|
||||
t.Errorf("allocs = %v, want %v", plan.Steps[0].Allocs, want)
|
||||
}
|
||||
for i, a := range want {
|
||||
if plan.Steps[0].Allocs[i] != a {
|
||||
t.Errorf("alloc[%d] = %q, want %q", i, plan.Steps[0].Allocs[i], a)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderUpdatePlan_NilSpec(t *testing.T) {
|
||||
_, err := RenderUpdatePlan(nil)
|
||||
if err == nil {
|
||||
t.Fatal("expected error for nil spec")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "spec is nil") {
|
||||
t.Errorf("error = %q, want 'spec is nil'", err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderUpdatePlan_NilUpdate(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{Kind: "Service", Name: "web", Count: 1}
|
||||
_, err := RenderUpdatePlan(spec)
|
||||
if err == nil {
|
||||
t.Fatal("expected error for nil update block")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "update block is nil") {
|
||||
t.Errorf("error = %q, want 'update block is nil'", err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderUpdatePlan_CountZero(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 0,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "rolling",
|
||||
},
|
||||
}
|
||||
_, err := RenderUpdatePlan(spec)
|
||||
if err == nil {
|
||||
t.Fatal("expected error for count 0")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "count must be") {
|
||||
t.Errorf("error = %q, want 'count must be'", err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderUpdatePlan_UnknownStrategy(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 2,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "recreate",
|
||||
},
|
||||
}
|
||||
_, err := RenderUpdatePlan(spec)
|
||||
if err == nil {
|
||||
t.Fatal("expected error for unknown strategy")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "unknown strategy") {
|
||||
t.Errorf("error = %q, want 'unknown strategy'", err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderUpdatePlan_AllStrategies(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: 3,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: strat,
|
||||
MaxParallel: 1,
|
||||
},
|
||||
}
|
||||
plan, err := RenderUpdatePlan(spec)
|
||||
if err != nil {
|
||||
t.Fatalf("RenderUpdatePlan(%s): %v", strat, err)
|
||||
}
|
||||
if len(plan.Steps) == 0 {
|
||||
t.Errorf("strategy %s produced 0 steps", strat)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderUpdatePlan_PromoteStepCarriesCanaryAllocs(t *testing.T) {
|
||||
// The promote step lists the canary allocs so the executor knows
|
||||
// which allocs are being promoted from canary to stable.
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 4,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "canary",
|
||||
Canary: "2",
|
||||
},
|
||||
}
|
||||
plan, err := RenderUpdatePlan(spec)
|
||||
if err != nil {
|
||||
t.Fatalf("RenderUpdatePlan: %v", err)
|
||||
}
|
||||
promote := plan.Steps[1]
|
||||
if promote.Action != "promote" {
|
||||
t.Fatalf("step 1 action = %q, want promote", promote.Action)
|
||||
}
|
||||
if len(promote.Allocs) != 2 {
|
||||
t.Errorf("promote allocs = %d, want 2 (the canary allocs)", len(promote.Allocs))
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderUpdatePlan_MaxParallelClampedToCount(t *testing.T) {
|
||||
// max_parallel > count is clamped to count (defensive; validator
|
||||
// rejects this but the emitter is defensive against direct
|
||||
// callers).
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 2,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "rolling",
|
||||
MaxParallel: 99,
|
||||
},
|
||||
}
|
||||
plan, err := RenderUpdatePlan(spec)
|
||||
if err != nil {
|
||||
t.Fatalf("RenderUpdatePlan: %v", err)
|
||||
}
|
||||
// Clamped to 2 → single batch of 2.
|
||||
if len(plan.Steps) != 1 {
|
||||
t.Errorf("got %d steps, want 1 (clamped)", len(plan.Steps))
|
||||
}
|
||||
if len(plan.Steps[0].Allocs) != 2 {
|
||||
t.Errorf("step 0 allocs = %d, want 2", len(plan.Steps[0].Allocs))
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,126 @@
|
||||
package schema
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"git.cloudinit.dev/coreci/orca/internal/jobspec"
|
||||
)
|
||||
|
||||
// UpdateValidator validates the rolling/canary/blue-green update stanza
|
||||
// (P03). The schema validator (ServiceValidator) already enforces that
|
||||
// the strategy is one of rolling/canary/blue-green and that the update
|
||||
// block is present for a Service. UpdateValidator adds the
|
||||
// field-level validation:
|
||||
//
|
||||
// - max_parallel: integer in [1, count] (defaults to 1 when unset)
|
||||
// - min_healthy_time: a valid time.Duration when set (time.ParseDuration)
|
||||
// - healthy_deadline: a valid time.Duration when set (time.ParseDuration)
|
||||
// - canary: an integer count in [0, count] OR a percentage string of
|
||||
// the form "<n>%" where n is in [0, 100] (the parser already accepts
|
||||
// both shapes; the validator accepts them too). Only meaningful
|
||||
// for the canary strategy; ignored (but still validated for shape)
|
||||
// for rolling/blue-green.
|
||||
// - auto_promote: boolean (no validation beyond the parser's
|
||||
// true/false parse; the field is always populated)
|
||||
//
|
||||
// The validator is pure (no I/O). Violations return a clear error
|
||||
// listing every problem found, mirroring the per-field style of
|
||||
// ServiceValidator.
|
||||
type UpdateValidator struct{}
|
||||
|
||||
// Validate validates the UpdateBlock on the given spec. The spec must
|
||||
// be non-nil and carry a Count (services have count ≥ 1 per
|
||||
// ServiceValidator). When spec.Update is nil the validator returns an
|
||||
// error (the update block is required for Service; this validator
|
||||
// assumes the caller has already established the spec is a Service).
|
||||
func (UpdateValidator) Validate(spec *jobspec.WorkloadSpec) error {
|
||||
if spec == nil {
|
||||
return fmt.Errorf("schema/Update: spec is nil")
|
||||
}
|
||||
if spec.Update == nil {
|
||||
return fmt.Errorf("schema/Update: update block is nil")
|
||||
}
|
||||
var errs []string
|
||||
u := spec.Update
|
||||
|
||||
switch u.Strategy {
|
||||
case "rolling", "canary", "blue-green":
|
||||
case "":
|
||||
errs = append(errs, "update strategy required (one of rolling, canary, blue-green)")
|
||||
default:
|
||||
errs = append(errs, fmt.Sprintf("update strategy %q invalid (want one of rolling, canary, blue-green)", u.Strategy))
|
||||
}
|
||||
|
||||
// max_parallel defaults to 1 when unset (0); validate the range
|
||||
// only when the user has set it explicitly.
|
||||
if u.MaxParallel != 0 {
|
||||
if u.MaxParallel < 1 {
|
||||
errs = append(errs, fmt.Sprintf("update.max_parallel must be ≥ 1, got %d", u.MaxParallel))
|
||||
}
|
||||
if spec.Count > 0 && u.MaxParallel > spec.Count {
|
||||
errs = append(errs, fmt.Sprintf("update.max_parallel %d exceeds count %d (must be 1..count)", u.MaxParallel, spec.Count))
|
||||
}
|
||||
}
|
||||
|
||||
if u.MinHealthyTime != "" {
|
||||
if _, err := time.ParseDuration(u.MinHealthyTime); err != nil {
|
||||
errs = append(errs, fmt.Sprintf("update.min_healthy_time %q is not a valid duration: %v", u.MinHealthyTime, err))
|
||||
}
|
||||
}
|
||||
if u.HealthyDeadline != "" {
|
||||
if _, err := time.ParseDuration(u.HealthyDeadline); err != nil {
|
||||
errs = append(errs, fmt.Sprintf("update.healthy_deadline %q is not a valid duration: %v", u.HealthyDeadline, err))
|
||||
}
|
||||
}
|
||||
|
||||
// canary accepts an integer count (0..count) or a percentage
|
||||
// ("<n>%" with n in 0..100). The field is only meaningful for the
|
||||
// canary strategy but we validate the shape regardless so a typo
|
||||
// in a rolling/blue-green stanza still surfaces.
|
||||
if u.Canary != "" {
|
||||
if err := validateCanary(u.Canary, spec.Count); err != nil {
|
||||
errs = append(errs, err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
// auto_promote is a bool; no extra validation beyond the parser.
|
||||
|
||||
return composeErrors("schema/Update", errs)
|
||||
}
|
||||
|
||||
// validateCanary validates the canary field shape: either an integer
|
||||
// count (0..count) or a percentage string "<n>%" (n in 0..100). count
|
||||
// is the spec.Count; when count is 0 (e.g. a DaemonSet or unset), the
|
||||
// integer-count upper bound is not enforced (only the percentage
|
||||
// bound is enforced, since percentage does not depend on count).
|
||||
func validateCanary(canary string, count int) error {
|
||||
c := strings.TrimSpace(canary)
|
||||
if c == "" {
|
||||
return nil
|
||||
}
|
||||
if strings.HasSuffix(c, "%") {
|
||||
nStr := strings.TrimSuffix(c, "%")
|
||||
n, err := strconv.Atoi(strings.TrimSpace(nStr))
|
||||
if err != nil {
|
||||
return fmt.Errorf("update.canary %q is not a valid percentage (want \"<n>%%\")", canary)
|
||||
}
|
||||
if n < 0 || n > 100 {
|
||||
return fmt.Errorf("update.canary percentage %d out of range (want 0..100)", n)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
n, err := strconv.Atoi(c)
|
||||
if err != nil {
|
||||
return fmt.Errorf("update.canary %q is not a valid count or percentage (want integer or \"<n>%%\")", canary)
|
||||
}
|
||||
if n < 0 {
|
||||
return fmt.Errorf("update.canary count %d must be ≥ 0", n)
|
||||
}
|
||||
if count > 0 && n > count {
|
||||
return fmt.Errorf("update.canary count %d exceeds count %d (must be 0..count)", n, count)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,464 @@
|
||||
package schema
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"git.cloudinit.dev/coreci/orca/internal/jobspec"
|
||||
)
|
||||
|
||||
func TestUpdateValidator_ValidRolling(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 4,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "rolling",
|
||||
MaxParallel: 2,
|
||||
MinHealthyTime: "30s",
|
||||
HealthyDeadline: "5m",
|
||||
},
|
||||
}
|
||||
v := UpdateValidator{}
|
||||
if err := v.Validate(spec); err != nil {
|
||||
t.Fatalf("expected nil, got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_ValidCanary(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 4,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "canary",
|
||||
MaxParallel: 2,
|
||||
MinHealthyTime: "30s",
|
||||
HealthyDeadline: "5m",
|
||||
Canary: "10%",
|
||||
AutoPromote: true,
|
||||
},
|
||||
}
|
||||
if err := (UpdateValidator{}).Validate(spec); err != nil {
|
||||
t.Fatalf("expected nil, got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_ValidBlueGreen(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 3,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "blue-green",
|
||||
MinHealthyTime: "1m",
|
||||
HealthyDeadline: "10m",
|
||||
},
|
||||
}
|
||||
if err := (UpdateValidator{}).Validate(spec); err != nil {
|
||||
t.Fatalf("expected nil, got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_ValidCanaryIntegerCount(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 4,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "canary",
|
||||
Canary: "1",
|
||||
},
|
||||
}
|
||||
if err := (UpdateValidator{}).Validate(spec); err != nil {
|
||||
t.Fatalf("integer canary count 1 should be valid, got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_ValidCanaryZeroPercent(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 4,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "canary",
|
||||
Canary: "0%",
|
||||
},
|
||||
}
|
||||
if err := (UpdateValidator{}).Validate(spec); err != nil {
|
||||
t.Fatalf("0%% canary should be valid, got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_ValidCanaryHundredPercent(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 4,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "canary",
|
||||
Canary: "100%",
|
||||
},
|
||||
}
|
||||
if err := (UpdateValidator{}).Validate(spec); err != nil {
|
||||
t.Fatalf("100%% canary should be valid, got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_EmptyDurationsOK(t *testing.T) {
|
||||
// Empty min_healthy_time / healthy_deadline should be accepted
|
||||
// (they are optional; defaults are applied by the executor).
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 2,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "rolling",
|
||||
},
|
||||
}
|
||||
if err := (UpdateValidator{}).Validate(spec); err != nil {
|
||||
t.Fatalf("empty durations should be valid, got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_NilSpec(t *testing.T) {
|
||||
if err := (UpdateValidator{}).Validate(nil); err == nil {
|
||||
t.Fatal("expected error for nil spec")
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_NilUpdateBlock(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{Kind: "Service", Name: "web", Count: 2}
|
||||
err := (UpdateValidator{}).Validate(spec)
|
||||
if err == nil {
|
||||
t.Fatal("expected error for nil update block")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "update block is nil") {
|
||||
t.Errorf("error = %q, want 'update block is nil'", err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_InvalidStrategy(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 2,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "recreate",
|
||||
},
|
||||
}
|
||||
err := (UpdateValidator{}).Validate(spec)
|
||||
if err == nil {
|
||||
t.Fatal("expected error for invalid strategy")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "strategy") || !strings.Contains(err.Error(), "invalid") {
|
||||
t.Errorf("error = %q, want 'strategy ... invalid'", err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_EmptyStrategy(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 2,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "",
|
||||
},
|
||||
}
|
||||
err := (UpdateValidator{}).Validate(spec)
|
||||
if err == nil {
|
||||
t.Fatal("expected error for empty strategy")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "strategy required") {
|
||||
t.Errorf("error = %q, want 'strategy required'", err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_MaxParallelZero(t *testing.T) {
|
||||
// max_parallel=0 means "unset" → default 1; accepted.
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 2,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "rolling",
|
||||
MaxParallel: 0,
|
||||
},
|
||||
}
|
||||
if err := (UpdateValidator{}).Validate(spec); err != nil {
|
||||
t.Fatalf("max_parallel=0 (unset) should be valid, got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_MaxParallelNegative(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 2,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "rolling",
|
||||
MaxParallel: -1,
|
||||
},
|
||||
}
|
||||
err := (UpdateValidator{}).Validate(spec)
|
||||
if err == nil {
|
||||
t.Fatal("expected error for negative max_parallel")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "max_parallel") {
|
||||
t.Errorf("error = %q, want 'max_parallel'", err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_MaxParallelExceedsCount(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 2,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "rolling",
|
||||
MaxParallel: 5,
|
||||
},
|
||||
}
|
||||
err := (UpdateValidator{}).Validate(spec)
|
||||
if err == nil {
|
||||
t.Fatal("expected error for max_parallel > count")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "exceeds count") {
|
||||
t.Errorf("error = %q, want 'exceeds count'", err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_MaxParallelEqualsCount(t *testing.T) {
|
||||
// max_parallel == count is the upper bound; valid.
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 3,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "rolling",
|
||||
MaxParallel: 3,
|
||||
},
|
||||
}
|
||||
if err := (UpdateValidator{}).Validate(spec); err != nil {
|
||||
t.Fatalf("max_parallel==count should be valid, got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_InvalidMinHealthyTime(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 2,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "rolling",
|
||||
MinHealthyTime: "not-a-duration",
|
||||
},
|
||||
}
|
||||
err := (UpdateValidator{}).Validate(spec)
|
||||
if err == nil {
|
||||
t.Fatal("expected error for invalid min_healthy_time")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "min_healthy_time") {
|
||||
t.Errorf("error = %q, want 'min_healthy_time'", err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_InvalidHealthyDeadline(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 2,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "rolling",
|
||||
HealthyDeadline: "nope",
|
||||
},
|
||||
}
|
||||
err := (UpdateValidator{}).Validate(spec)
|
||||
if err == nil {
|
||||
t.Fatal("expected error for invalid healthy_deadline")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "healthy_deadline") {
|
||||
t.Errorf("error = %q, want 'healthy_deadline'", err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_CanaryNegative(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 4,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "canary",
|
||||
Canary: "-1",
|
||||
},
|
||||
}
|
||||
err := (UpdateValidator{}).Validate(spec)
|
||||
if err == nil {
|
||||
t.Fatal("expected error for negative canary count")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "canary") {
|
||||
t.Errorf("error = %q, want 'canary'", err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_CanaryExceedsCount(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 4,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "canary",
|
||||
Canary: "5",
|
||||
},
|
||||
}
|
||||
err := (UpdateValidator{}).Validate(spec)
|
||||
if err == nil {
|
||||
t.Fatal("expected error for canary > count")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "exceeds count") {
|
||||
t.Errorf("error = %q, want 'exceeds count'", err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_CanaryPercentOver100(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 4,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "canary",
|
||||
Canary: "150%",
|
||||
},
|
||||
}
|
||||
err := (UpdateValidator{}).Validate(spec)
|
||||
if err == nil {
|
||||
t.Fatal("expected error for canary > 100%")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "out of range") {
|
||||
t.Errorf("error = %q, want 'out of range'", err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_CanaryPercentNegative(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 4,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "canary",
|
||||
Canary: "-10%",
|
||||
},
|
||||
}
|
||||
err := (UpdateValidator{}).Validate(spec)
|
||||
if err == nil {
|
||||
t.Fatal("expected error for negative canary percent")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "out of range") {
|
||||
t.Errorf("error = %q, want 'out of range'", err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_CanaryNotANumber(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 4,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "canary",
|
||||
Canary: "abc",
|
||||
},
|
||||
}
|
||||
err := (UpdateValidator{}).Validate(spec)
|
||||
if err == nil {
|
||||
t.Fatal("expected error for non-numeric canary")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "not a valid") {
|
||||
t.Errorf("error = %q, want 'not a valid'", err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_CanaryPercentNotANumber(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 4,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "canary",
|
||||
Canary: "xx%",
|
||||
},
|
||||
}
|
||||
err := (UpdateValidator{}).Validate(spec)
|
||||
if err == nil {
|
||||
t.Fatal("expected error for non-numeric canary percent")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "not a valid percentage") {
|
||||
t.Errorf("error = %q, want 'not a valid percentage'", err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_CanaryCountZeroOK(t *testing.T) {
|
||||
// canary=0 is the lower bound; valid.
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 4,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "canary",
|
||||
Canary: "0",
|
||||
},
|
||||
}
|
||||
if err := (UpdateValidator{}).Validate(spec); err != nil {
|
||||
t.Fatalf("canary=0 should be valid, got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_CanaryPercentWithoutCountOK(t *testing.T) {
|
||||
// A percentage canary does not depend on count; valid even when
|
||||
// count is 0 (e.g. DaemonSet-shaped spec bypassing the Service
|
||||
// validator — defensive).
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 0,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "canary",
|
||||
Canary: "25%",
|
||||
},
|
||||
}
|
||||
if err := (UpdateValidator{}).Validate(spec); err != nil {
|
||||
t.Fatalf("percentage canary with count=0 should be valid, got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateValidator_MultipleErrors(t *testing.T) {
|
||||
spec := &jobspec.WorkloadSpec{
|
||||
Kind: "Service",
|
||||
Name: "web",
|
||||
Count: 2,
|
||||
Update: &jobspec.UpdateBlock{
|
||||
Strategy: "recreate",
|
||||
MaxParallel: 99,
|
||||
MinHealthyTime: "nope",
|
||||
HealthyDeadline: "also-nope",
|
||||
Canary: "200%",
|
||||
},
|
||||
}
|
||||
err := (UpdateValidator{}).Validate(spec)
|
||||
if err == nil {
|
||||
t.Fatal("expected error, got nil")
|
||||
}
|
||||
for _, want := range []string{
|
||||
"strategy",
|
||||
"max_parallel",
|
||||
"min_healthy_time",
|
||||
"healthy_deadline",
|
||||
"canary",
|
||||
} {
|
||||
if !strings.Contains(err.Error(), want) {
|
||||
t.Errorf("error %q missing %q", err.Error(), want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Compile-time assertion that UpdateValidator implements Validator.
|
||||
var _ Validator = UpdateValidator{}
|
||||
Reference in New Issue
Block a user