Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 64cbbd543e |
+47
-7
@@ -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
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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"`
|
||||
|
||||
Reference in New Issue
Block a user