From 64cbbd543e00e4075944e7aa34cc01bb5f667f21 Mon Sep 17 00:00:00 2001 From: Jon Chery Date: Mon, 10 Aug 2026 16:21:50 +0000 Subject: [PATCH] =?UTF-8?q?feat(C):=20remote=20deployment=20correctness=20?= =?UTF-8?q?=E2=80=94=20PVE=20runtime,=20DB=20record,=20Traefik=20(REQ-166)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 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--- --- internal/cli/job.go | 54 +++++++++++++-- internal/cli/job_dispatch.go | 120 ++++++++++++++++++++++++++++------ internal/cli/watch_test.go | 4 +- internal/engine/dispatcher.go | 2 +- internal/model/job.go | 2 + 5 files changed, 153 insertions(+), 29 deletions(-) diff --git a/internal/cli/job.go b/internal/cli/job.go index 64e0b9b..d2f2da2 100644 --- a/internal/cli/job.go +++ b/internal/cli/job.go @@ -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 ' 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, + }) +} diff --git a/internal/cli/job_dispatch.go b/internal/cli/job_dispatch.go index fdc7dee..c917ac1 100644 --- a/internal/cli/job_dispatch.go +++ b/internal/cli/job_dispatch.go @@ -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 +} diff --git a/internal/cli/watch_test.go b/internal/cli/watch_test.go index a1fdadd..46471b5 100644 --- a/internal/cli/watch_test.go +++ b/internal/cli/watch_test.go @@ -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) } } diff --git a/internal/engine/dispatcher.go b/internal/engine/dispatcher.go index bce2934..adb7066 100644 --- a/internal/engine/dispatcher.go +++ b/internal/engine/dispatcher.go @@ -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. diff --git a/internal/model/job.go b/internal/model/job.go index 88d9d00..3524153 100644 --- a/internal/model/job.go +++ b/internal/model/job.go @@ -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"`