Files
orca/internal/cli/job_dispatch.go
T
Jon Chery cf3d98eb2b feat(P03): wire scheduler into job run + fix jobspec parser (REQ-151, REQ-152)
R-022: orca job run now deploys to remote nodes via scheduler -> emitter
-> SSH-push. Local exec fallback only when no remote nodes registered.

jobspec parser (REQ-152):
- schedule: and timeout: now parsed (were silently dropped)
- DaemonSet Count no longer defaults to 1 (was breaking DaemonSet)
- restart: policy translated to systemd Restart=/StartLimitBurst
- job lint: advisory warnings for cron/health/update/affinity (honest)

scheduler wiring (REQ-151, C-44):
- new internal/cli/job_dispatch.go: dispatchDecision + deployRemote
- scheduler.Schedule evaluates constraints/capacity/affinity
- --target overrides scheduler (manual pinning)
- local fallback only when len(ready non-localhost nodes)==0
- C-44: SSH-push failure returns error (no silent local fallback)
- systemd-analyze verify on rendered unit before deploy

Tests: 22 new test functions covering scheduler, parser, C-44, local
fallback, target override, systemd-analyze skip, restart directives.

---ci---
project: orca
phase: 3
milestone: v0.13
status: complete
requirements:
  covered: [151, 152]
---/ci---
2026-08-07 19:59:31 +00:00

436 lines
16 KiB
Go

// Package cli: job_dispatch.go wires the v0.9 CLI-side scheduler
// (internal/scheduler), the systemd emitter (internal/emitter), and the
// SSH-push transport (internal/sshpush) into `orca job run`
// (REQ-151, binding condition C-44, v0.13 milestone phase 03).
//
// The dispatch flow (replacing the deprecated mTLS Dispatcher path) is:
//
// 1. Load registered nodes from the orca registry (DB) and project them
// into scheduler.NodeInfo + a hostname->model.Node map for SSH-push.
// 2. If --target is set, pin to that node directly (manual override).
// 3. If no --target and no remote nodes are registered (only localhost
// or none), fall back to local exec (backward compat for dev mode).
// 4. If no --target and remote nodes ARE registered, invoke
// scheduler.Schedule -> pick the best node -> render the systemd unit
// via internal/emitter -> systemd-analyze verify (when available) ->
// SSH-push the unit to the target via internal/sshpush.
//
// C-44 (binding condition): if the scheduler selects a node but the
// SSH-push FAILS, return an error. Do NOT silently fall back to local
// execution. Local fallback is ONLY when len(registeredRemoteNodes)==0.
package cli
import (
"context"
"errors"
"fmt"
"log/slog"
"os"
"os/exec"
"strings"
"time"
"git.cloudinit.dev/coreci/orca/internal/certpaths"
"git.cloudinit.dev/coreci/orca/internal/emitter"
"git.cloudinit.dev/coreci/orca/internal/jobspec"
"git.cloudinit.dev/coreci/orca/internal/model"
"git.cloudinit.dev/coreci/orca/internal/scheduler"
"git.cloudinit.dev/coreci/orca/internal/sshpush"
"git.cloudinit.dev/coreci/orca/internal/store"
)
// jobDispatchTransport is the SSH-push surface `job run` needs for
// remote deployment. *sshpush.Transport satisfies it; tests substitute
// a mock (same pattern as txn.go / job_verify.go).
type jobDispatchTransport interface {
WriteFile(ctx context.Context, peer string, path string, content []byte, mode os.FileMode) error
Exec(ctx context.Context, peer string, cmd string) ([]byte, error)
Close() error
}
// jobDispatchTransportOverride is the package-level seam. When non-nil
// it replaces the production transport; tests set it and restore nil.
var jobDispatchTransportOverride jobDispatchTransport
// jobDispatchTransportFromCtx returns the active SSH-push transport.
// Tests override via jobDispatchTransportOverride; production builds a
// real *sshpush.Transport from the orca SSH key + known_hosts paths.
func jobDispatchTransportFromCtx() (jobDispatchTransport, error) {
if jobDispatchTransportOverride != nil {
return jobDispatchTransportOverride, nil
}
keyPath := certpaths.SSHKeyPath()
khPath := certpaths.KnownHostsPath()
return sshpush.NewTransport(keyPath, khPath), nil
}
// dispatchResult is the outcome of a `job run` dispatch decision.
type dispatchResult struct {
// mode is "local" (local exec fallback) or "remote" (scheduled +
// SSH-pushed to a peer).
mode string
// node is the hostname of the selected/pinned node (remote only).
node string
// allocID is the scheduler allocation id (remote only).
allocID string
// unitPaths is the list of systemd unit paths written (remote only).
unitPaths []string
}
// dispatchDecision decides how `job run` should execute the spec:
//
// - "local" -> run via the local executor (dev mode / no remote nodes)
// - "remote" -> render + SSH-push the systemd unit to the chosen node
//
// It loads registered nodes from the DB, projects them into
// scheduler.NodeInfo, and consults the scheduler when no --target is
// set. Returns a dispatchResult describing the chosen path; the caller
// performs the actual execution.
//
// C-44: when remote nodes are registered, a scheduling failure returns
// an error (no local fallback). The local fallback ONLY happens when
// there are zero remote nodes registered (only localhost or none).
func dispatchDecision(ctx context.Context, spec *jobspec.WorkloadSpec, target string) (*dispatchResult, map[string]*model.Node, error) {
if spec == nil {
return nil, nil, errors.New("dispatch: nil spec")
}
db, closer, err := openDB()
if err != nil {
return nil, nil, fmt.Errorf("dispatch: open db: %w", err)
}
defer closer()
nodeRepo := store.NewNodeRepo(db)
capRepo := store.NewCapacityRepo(db)
nodes, err := nodeRepo.List(ctx)
if err != nil {
return nil, nil, fmt.Errorf("dispatch: list nodes: %w", err)
}
caps, err := capRepo.List(ctx)
if err != nil {
return nil, nil, fmt.Errorf("dispatch: list capacity: %w", err)
}
capByNode := make(map[string]*store.NodeCapacity, len(caps))
for _, c := range caps {
capByNode[c.NodeID] = c
}
// Project registered nodes into scheduler.NodeInfo. A node counts
// as a "remote" scheduling candidate when it is ready and is NOT
// the localhost node (kind=localhost). localhost is excluded from
// the candidate set so the scheduler only considers real peers;
// when the candidate set is empty we fall back to local exec.
var candidates []scheduler.NodeInfo
remoteNodes := make(map[string]*model.Node) // hostname -> node
for _, n := range nodes {
if n.State != model.NodeStateReady {
continue
}
if n.Kind == string(model.NodeKindLocalhost) {
continue
}
ni := nodeToNodeInfo(n, capByNode[n.ID])
candidates = append(candidates, ni)
remoteNodes[ni.Hostname] = n
}
// --target override: pin to the named node. The target may be a
// node ID, name, or hostname. We resolve it against the registered
// nodes (including localhost when explicitly targeted).
if strings.TrimSpace(target) != "" {
chosen, err := resolveTargetNode(ctx, nodeRepo, target)
if err != nil {
return nil, nil, err
}
hostname := chosen.Name
if hostname == "" {
hostname = chosen.ID
}
// Even a localhost target goes through the remote push path
// when explicitly pinned (the operator asked for it).
remoteNodes[hostname] = chosen
return &dispatchResult{
mode: "remote",
node: hostname,
allocID: allocIDFor(spec, 0),
}, remoteNodes, nil
}
// No remote nodes registered -> local exec fallback (dev mode).
if len(candidates) == 0 {
return &dispatchResult{mode: "local"}, remoteNodes, nil
}
// Remote nodes registered -> invoke the scheduler. A scheduling
// failure is an error (C-44: no silent local fallback).
placements, err := scheduler.Schedule(candidates, scheduler.WorkloadRequest{
Spec: spec,
Namespace: "default",
})
if err != nil {
return nil, nil, fmt.Errorf("dispatch: schedule: %w", err)
}
if len(placements) == 0 {
return nil, nil, fmt.Errorf("dispatch: scheduler returned no placements for %q", spec.Name)
}
// Job/DaemonSet produce one-or-many placements; for `job run` we
// deploy the first placement (the best-fit node). Multi-replica
// Service fan-out is handled by the txn/apply path, not job run.
p := placements[0]
return &dispatchResult{
mode: "remote",
node: p.Node,
allocID: p.AllocID,
}, remoteNodes, nil
}
// deployRemote renders the systemd unit for the spec on the chosen
// node, runs systemd-analyze verify (when available), and SSH-pushes
// the unit files to the peer. Returns the list of unit paths written.
//
// C-44: any render/verify/push failure is returned as an error; the
// caller must NOT fall back to local exec.
func deployRemote(ctx context.Context, spec *jobspec.WorkloadSpec, res *dispatchResult, nodesByHost map[string]*model.Node) ([]string, error) {
if res == nil || res.mode != "remote" {
return nil, errors.New("deployRemote: not a remote dispatch")
}
node, ok := nodesByHost[res.node]
if !ok {
return nil, fmt.Errorf("deployRemote: selected node %q not found in registry", res.node)
}
// Render the systemd unit via the emitter. The runtime is required
// for the process emitter; a spec with no runtime has nothing to
// ExecStart and is rejected by the emitter.
em := emitter.SystemdEmitter{}
enode := &emitter.Node{
Hostname: node.Name,
Runtime: []string{"process"},
Tags: nil,
}
// Advertise the node kind as a runtime so the emitter can branch
// (proxmox nodes expose pve-* runtimes). For process workloads
// this is informational.
if node.Kind == string(model.NodeKindProxmox) {
enode.Runtime = append(enode.Runtime, "proxmox")
}
files, err := em.Render(spec, enode)
if err != nil {
return nil, fmt.Errorf("deployRemote: render unit: %w", err)
}
// T9: systemd-analyze verify on the rendered unit before deploy.
// Run it locally (the unit is a portable text file); if
// systemd-analyze is not installed, skip silently (dev boxes
// without systemd). A verification FAILURE is an error.
for _, f := range files {
if err := verifySystemdUnit(ctx, f.Path, f.Content); err != nil {
return nil, fmt.Errorf("deployRemote: systemd-analyze verify %s: %w", f.Path, err)
}
}
// SSH-push the unit files to the peer.
peer := sshPeerFor(node)
transport, err := jobDispatchTransportFromCtx()
if err != nil {
return nil, fmt.Errorf("deployRemote: transport: %w", err)
}
defer transport.Close()
var written []string
for _, f := range files {
mode := os.FileMode(0o644)
if f.Mode != "" {
// f.Mode is an octal string like "0644".
var m uint64
if _, perr := fmt.Sscanf(f.Mode, "%o", &m); perr == nil {
mode = os.FileMode(m)
}
}
if err := transport.WriteFile(ctx, peer, f.Path, []byte(f.Content), mode); err != nil {
// C-44: SSH-push failure -> error, NOT local fallback.
return nil, fmt.Errorf("deployRemote: push %s to %s (%s): %w", f.Path, res.node, peer, err)
}
written = append(written, f.Path)
}
// Reload systemd + enable the unit so it starts at boot. These are
// best-effort; a failure here is surfaced but does not undo the
// push (the unit is on disk). We use systemctl daemon-reload +
// enable --now for each .service unit (.target units for task
// groups are also enabled).
for _, p := range written {
if !strings.HasSuffix(p, ".service") && !strings.HasSuffix(p, ".target") {
continue
}
if _, err := transport.Exec(ctx, peer, fmt.Sprintf("systemctl daemon-reload && systemctl enable --now %s", shellQuoteSystemd(p))); err != nil {
return written, fmt.Errorf("deployRemote: enable %s on %s: %w", p, res.node, err)
}
}
return written, nil
}
// verifySystemdUnit runs `systemd-analyze verify` on the rendered unit
// content. The unit is written to a temp file (with its real basename)
// so systemd-analyze resolves fragment paths correctly. When
// systemd-analyze is not on PATH, the check is skipped (dev boxes
// without systemd). A non-zero exit from systemd-analyze is an error.
func verifySystemdUnit(ctx context.Context, unitPath, content string) error {
bin, err := exec.LookPath("systemd-analyze")
if err != nil {
// systemd-analyze not available (e.g. macOS dev box, minimal
// container). Skip verification rather than failing — the
// render layer already validates the spec shape.
return nil
}
base := unitPath
if idx := strings.LastIndex(unitPath, "/"); idx >= 0 {
base = unitPath[idx+1:]
}
// os.CreateTemp appends a random suffix that would strip the
// .service/.target extension systemd-analyze needs to recognize the
// unit. Create the temp file in a dedicated temp dir with the exact
// basename so the extension is preserved.
tmpDir, err := os.MkdirTemp("", "orca-verify-")
if err != nil {
return fmt.Errorf("temp dir: %w", err)
}
defer os.RemoveAll(tmpDir)
tmpPath := tmpDir + "/" + base
if err := os.WriteFile(tmpPath, []byte(content), 0o644); err != nil {
return fmt.Errorf("write temp unit: %w", err)
}
vctx, cancel := context.WithTimeout(ctx, 10*time.Second)
defer cancel()
cmd := exec.CommandContext(vctx, bin, "verify", tmpPath)
out, err := cmd.CombinedOutput()
if err != nil {
// Trim the temp path from the output so the error reads with
// the real unit path.
msg := strings.TrimSpace(string(out))
msg = strings.ReplaceAll(msg, tmpPath, unitPath)
return fmt.Errorf("systemd-analyze verify failed: %s", msg)
}
return nil
}
// nodeToNodeInfo projects a registered model.Node (+ its capacity
// declaration) into a scheduler.NodeInfo. Runtimes are derived from the
// node kind (proxmox -> "proxmox"; else "process"). Tags are sourced
// from node metadata["tags"] (comma-separated) when present. Capacity
// is sourced from the NodeCapacity row when present (else zero, which
// the scheduler treats as always-fits on the capacity axis).
func nodeToNodeInfo(n *model.Node, cap *store.NodeCapacity) scheduler.NodeInfo {
ni := scheduler.NodeInfo{
Hostname: n.Name,
Kind: n.Kind,
}
if ni.Kind == "" {
ni.Kind = string(model.NodeKindLinux)
}
switch n.Kind {
case string(model.NodeKindProxmox):
ni.Runtimes = []string{"process", "proxmox"}
default:
ni.Runtimes = []string{"process"}
}
if tags := nodeMetadataTag(n, "tags"); tags != "" {
for _, t := range strings.Split(tags, ",") {
t = strings.TrimSpace(t)
if t != "" {
ni.Tags = append(ni.Tags, t)
}
}
}
if cap != nil {
ni.CPU = cap.CPUMillicores
ni.Memory = cap.MemoryMiB
ni.FreeCPU = cap.CPUMillicores
ni.FreeMem = cap.MemoryMiB
}
return ni
}
// nodeMetadataTag reads a key from the node's metadata map. Returns ""
// when the metadata is nil or the key is absent.
func nodeMetadataTag(n *model.Node, key string) string {
if n == nil || n.Metadata == nil {
return ""
}
return n.Metadata[key]
}
// resolveTargetNode resolves a --target value (node ID, name, or
// hostname) to a registered *model.Node. Returns an error when the
// target is not found.
func resolveTargetNode(ctx context.Context, repo *store.NodeRepo, target string) (*model.Node, error) {
target = strings.TrimSpace(target)
if target == "" {
return nil, errors.New("resolveTargetNode: empty target")
}
// Try by ID first.
if n, err := repo.Get(ctx, target); err == nil {
return n, nil
}
// Then by name.
if n, err := repo.GetByName(ctx, target); err == nil {
return n, nil
}
return nil, fmt.Errorf("resolveTargetNode: target node %q not found in registry", target)
}
// sshPeerFor returns the host:port SSH peer address for a node. The
// node's orca Address is the mTLS daemon port (host:8443); SSH uses a
// different port. We derive the host from the orca Address and use the
// SSH port from node metadata["ssh_port"] when present, else 22.
func sshPeerFor(n *model.Node) string {
host := n.Address
if idx := strings.LastIndex(host, ":"); idx >= 0 {
host = host[:idx]
}
// Strip an ipv6 bracket if present.
host = strings.TrimPrefix(host, "[")
host = strings.TrimSuffix(host, "]")
port := "22"
if n != nil && n.Metadata != nil {
if p, ok := n.Metadata["ssh_port"]; ok && strings.TrimSpace(p) != "" {
port = strings.TrimSpace(p)
}
}
return host + ":" + port
}
// allocIDFor renders a stable allocation id for a spec index, matching
// the scheduler's allocID format (ns/name-idx).
func allocIDFor(spec *jobspec.WorkloadSpec, idx int) string {
return fmt.Sprintf("default/%s-%d", spec.Name, idx)
}
// shellQuoteSystemd single-quotes a path for safe shell interpolation
// in the remote systemctl command. Mirrors sshpush.shellQuote.
func shellQuoteSystemd(s string) string {
return "'" + strings.ReplaceAll(s, "'", "'\\''") + "'"
}
// logDispatch records the dispatch decision to the structured logger.
func logDispatch(res *dispatchResult, err error) {
log := slog.Default()
if res == nil {
log.Info("job.dispatch", slog.String("event", "job.dispatch"), slog.String("mode", "error"), slog.Any("error", err))
return
}
attrs := []any{slog.String("event", "job.dispatch"), slog.String("mode", res.mode)}
if res.node != "" {
attrs = append(attrs, slog.String("node", res.node))
}
if res.allocID != "" {
attrs = append(attrs, slog.String("alloc_id", res.allocID))
}
if err != nil {
attrs = append(attrs, slog.Any("error", err))
}
log.Info("job.dispatch", attrs...)
}