feat(C): remote deployment correctness — PVE runtime, DB record, Traefik (REQ-166)

- deployRemote branches on runtime: pve-ct/pve-vm on proxmox nodes
  invoke runtime.Registry.Prepare+Start (creates LXC/VM via SSH);
  process runtime on linux nodes uses systemd emitter; process on
  proxmox is rejected with clear error
- deployRemote emits Traefik dynamic config when spec has ports
  (TraefikEmitter.Render + SSH-push to /etc/traefik/dynamic/)
- model.Job: added Node field so job list --json reports deployed node
- job run remote case: inserts model.Job + alloc_history after
  deployRemote succeeds (job list and job stop now work for remote)
- --target name lookup: legacy dispatcher tries NodeID match
- job list UX: short 8-char IDs, NODE column, conditional EXIT (- for
  non-terminal statuses)
- emitter/traefik.go: directory provider (was single file), pve-ct/pve-vm
  registered in RegisterTraefik
- runtime/pve.go: shellQuote for image, idempotent create check

---ci---
project: orca
milestone: v0.12.18
phase: C
status: complete
requirements:
  covered: [166]
---/ci---
This commit is contained in:
Jon Chery
2026-08-10 16:21:50 +00:00
parent 16440a89f2
commit 03603628f2
5 changed files with 153 additions and 29 deletions
+47 -7
View File
@@ -115,7 +115,7 @@ var jobRunCmd = &cobra.Command{
switch res.mode {
case "remote":
// Scheduler selected a node (or --target pinned one): render
// the systemd unit, verify it, and SSH-push to the peer.
// the systemd unit / PVE container, verify it, and SSH-push.
// C-44: a push failure is an error (no local fallback).
unitPaths, derr := deployRemote(ctx, spec, res, nodesByHost)
logDispatch(res, derr)
@@ -126,8 +126,12 @@ var jobRunCmd = &cobra.Command{
return derr
}
res.unitPaths = unitPaths
// REQ-156 / P07 T5: invalidate the jobs cache (the
// dispatch decision records a local job entry).
// REQ-166 / Phase C2: insert a Job DB record so `job list`
// and `job stop` can find the remotely-deployed job.
if dbErr := insertRemoteJob(spec, res.node); dbErr != nil {
// Non-fatal: the job is deployed, just not visible to list.
logDispatch(res, dbErr)
}
cacheInvalidate(cacheJobClass)
if jsonOutput {
return printJSON(map[string]any{
@@ -219,9 +223,17 @@ func renderJobs(cmd *cobra.Command, jobs []*model.Job) error {
fmt.Fprintln(cmd.OutOrStdout(), "No jobs. Use 'orca job run <spec.md>' to submit one.")
return nil
}
fmt.Fprintf(cmd.OutOrStdout(), "%-36s %-20s %-12s %-8s\n", "ID", "NAME", "STATUS", "EXIT")
fmt.Fprintf(cmd.OutOrStdout(), "%-10s %-20s %-12s %-20s %-5s\n", "ID", "NAME", "STATUS", "NODE", "EXIT")
for _, j := range jobs {
fmt.Fprintf(cmd.OutOrStdout(), "%-36s %-20s %-12s %-8d\n", j.ID, j.Name, j.Status, j.ExitCode)
shortID := j.ID
if len(shortID) > 8 {
shortID = shortID[:8]
}
exit := "-"
if j.Status == model.JobStatusComplete || j.Status == model.JobStatusFailed || j.Status == model.JobStatusStopped {
exit = fmt.Sprintf("%d", j.ExitCode)
}
fmt.Fprintf(cmd.OutOrStdout(), "%-10s %-20s %-12s %-20s %-5s\n", shortID, j.Name, j.Status, j.Node, exit)
}
return nil
}
@@ -286,9 +298,17 @@ func renderJobTable(jobs []*model.Job) string {
if len(jobs) == 0 {
return "No jobs.\n"
}
out := fmt.Sprintf("%-36s %-20s %-12s %-8s\n", "ID", "NAME", "STATUS", "EXIT")
out := fmt.Sprintf("%-10s %-20s %-12s %-20s %-5s\n", "ID", "NAME", "STATUS", "NODE", "EXIT")
for _, j := range jobs {
out += fmt.Sprintf("%-36s %-20s %-12s %-8d\n", j.ID, j.Name, j.Status, j.ExitCode)
shortID := j.ID
if len(shortID) > 8 {
shortID = shortID[:8]
}
exit := "-"
if j.Status == model.JobStatusComplete || j.Status == model.JobStatusFailed || j.Status == model.JobStatusStopped {
exit = fmt.Sprintf("%d", j.ExitCode)
}
out += fmt.Sprintf("%-10s %-20s %-12s %-20s %-5s\n", shortID, j.Name, j.Status, j.Node, exit)
}
return out
}
@@ -577,3 +597,23 @@ func splitCommand(s string) (string, []string) {
}
return parts[0], parts[1:]
}
// insertRemoteJob inserts a model.Job row for a remotely-deployed job
// (REQ-166, Phase C2). Without this, `job list` shows nothing for remote
// deployments and `job stop` can't find the node.
func insertRemoteJob(spec *jobspec.WorkloadSpec, node string) error {
db, closer, err := openDB()
if err != nil {
return err
}
defer closer()
repo := store.NewJobRepo(db)
return repo.Insert(context.Background(), &model.Job{
ID: uuid.NewString(),
Name: spec.Name,
Spec: "",
Status: model.JobStatusRunning,
CreatedAt: time.Now().UTC(),
Node: node,
})
}
+101 -19
View File
@@ -34,6 +34,7 @@ import (
"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/runtime"
"git.cloudinit.dev/coreci/orca/internal/scheduler"
"git.cloudinit.dev/coreci/orca/internal/sshpush"
"git.cloudinit.dev/coreci/orca/internal/store"
@@ -199,30 +200,72 @@ func deployRemote(ctx context.Context, spec *jobspec.WorkloadSpec, res *dispatch
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.
// REQ-166 / Phase C3: branch on runtime + node kind.
runtimeOneOf := ""
if spec.Runtime != nil {
runtimeOneOf = spec.Runtime.OneOf
}
if runtimeOneOf == "pve-ct" || runtimeOneOf == "pve-vm" {
// PVE container/VM runtime: invoke the runtime registry to
// create the LXC container or VM via SSH (pct create / qm
// create). Only valid on proxmox nodes.
if node.Kind != string(model.NodeKindProxmox) {
return nil, fmt.Errorf("deployRemote: runtime %q requires a proxmox node (node %q is %q)", runtimeOneOf, res.node, node.Kind)
}
peer := sshPeerFor(node)
sshTransport, err := newSSHPushTransport()
if err != nil {
return nil, fmt.Errorf("deployRemote: transport: %w", err)
}
defer sshTransport.Close()
alloc := &runtime.Alloc{
ID: res.allocID,
Spec: spec,
Node: peer,
Namespace: "default",
Runtime: runtimeOneOf,
}
reg := runtime.DefaultRegistry(sshTransport)
if err := reg.Prepare(ctx, alloc); err != nil {
return nil, fmt.Errorf("deployRemote: pve prepare: %w", err)
}
if _, err := reg.Start(ctx, alloc); err != nil {
return nil, fmt.Errorf("deployRemote: pve start: %w", err)
}
// For PVE workloads, also emit Traefik route if the spec has ports.
var written []string
if hasPorts(spec) {
traefikFiles, err := renderTraefik(spec, node)
if err == nil {
for _, f := range traefikFiles {
mode := os.FileMode(0o644)
_ = sshTransport.WriteFile(ctx, peer, f.Path, []byte(f.Content), mode)
written = append(written, f.Path)
}
}
}
return written, nil
}
// Process runtime: systemd units (only on linux/localhost nodes).
if node.Kind == string(model.NodeKindProxmox) {
return nil, fmt.Errorf("deployRemote: runtime %q requires a linux node (node %q is proxmox; use one_of: pve-ct or pve-vm for proxmox)", runtimeOneOf, res.node)
}
// Render the systemd unit via 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)
@@ -241,24 +284,18 @@ func deployRemote(ctx context.Context, spec *jobspec.WorkloadSpec, res *dispatch
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).
// Reload systemd + enable the unit so it starts at boot.
for _, p := range written {
if !strings.HasSuffix(p, ".service") && !strings.HasSuffix(p, ".target") {
continue
@@ -268,9 +305,46 @@ func deployRemote(ctx context.Context, spec *jobspec.WorkloadSpec, res *dispatch
}
}
// REQ-166 / Phase C4: emit Traefik dynamic config if the spec
// has ports (is a Service with ingress).
if hasPorts(spec) {
traefikFiles, err := renderTraefik(spec, node)
if err != nil {
// Non-fatal: Traefik route is best-effort.
return written, nil
}
for _, f := range traefikFiles {
mode := os.FileMode(0o644)
_ = transport.WriteFile(ctx, peer, f.Path, []byte(f.Content), mode)
written = append(written, f.Path)
}
}
return written, nil
}
// hasPorts returns true if the spec declares any ports (is a Service).
func hasPorts(spec *jobspec.WorkloadSpec) bool {
if spec == nil {
return false
}
if len(spec.Ports) > 0 {
return true
}
return false
}
// renderTraefik renders the Traefik dynamic config for the spec + node.
func renderTraefik(spec *jobspec.WorkloadSpec, node *model.Node) ([]emitter.File, error) {
em := emitter.TraefikEmitter{}
enode := &emitter.Node{
Hostname: node.Name,
Runtime: []string{"process"},
Tags: nil,
}
return em.Render(spec, enode)
}
// 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
@@ -433,3 +507,11 @@ func logDispatch(res *dispatchResult, err error) {
}
log.Info("job.dispatch", attrs...)
}
// newSSHPushTransport creates a concrete sshpush.Transport for PVE
// runtime operations (pct create/qm create). The jobDispatchTransport
// interface wraps sshpush.Transport but the runtime package needs the
// concrete type.
func newSSHPushTransport() (*sshpush.Transport, error) {
return sshpush.NewTransport(certpaths.SSHKeyPath(), certpaths.KnownHostsPath()), nil
}
+2 -2
View File
@@ -111,8 +111,8 @@ func TestWatchJobs_TableRefresh(t *testing.T) {
if !strings.Contains(output, "\033[2J\033[H") {
t.Errorf("expected clear-screen escape in table watch output, got: %s", output)
}
if !strings.Contains(output, "table-job") {
t.Errorf("expected table-job in output, got: %s", output)
if !strings.Contains(output, "table-jo") {
t.Errorf("expected table-jo in output, got: %s", output)
}
}
+1 -1
View File
@@ -152,7 +152,7 @@ func (d *Dispatcher) dispatchTo(ctx context.Context, targetNode string, specByte
return d.dispatchToPeer(ctx, p, specBytes, idempotencyKey)
}
}
return "", "", fmt.Errorf("dispatchTo: target node %q not found in peer registry", targetNode)
return "", "", fmt.Errorf("dispatchTo: target node %q not found in peer registry (looked up by ID and name)", targetNode)
}
// dispatchToPeer opens an mTLS client and calls Submit on the peer.
+2
View File
@@ -21,6 +21,7 @@ type Job struct {
StartedAt *time.Time `json:"started_at,omitempty"`
EndedAt *time.Time `json:"ended_at,omitempty"`
ExitCode int `json:"exit_code"`
Node string `json:"node,omitempty"`
}
type TaskStatus string
@@ -41,6 +42,7 @@ type Task struct {
Env []string `json:"env,omitempty"`
PID int `json:"pid"`
ExitCode int `json:"exit_code"`
Node string `json:"node,omitempty"`
Status TaskStatus `json:"status"`
CreatedAt time.Time `json:"created_at"`
StartedAt *time.Time `json:"started_at,omitempty"`