// 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...) }