diff --git a/shuttle/src/gen/spindle/agent/v1/spindle.agent.v1.rs b/shuttle/src/gen/spindle/agent/v1/spindle.agent.v1.rs index 04a12079..a274c2a2 100644 --- a/shuttle/src/gen/spindle/agent/v1/spindle.agent.v1.rs +++ b/shuttle/src/gen/spindle/agent/v1/spindle.agent.v1.rs @@ -2,145 +2,146 @@ // This file is @generated by prost-build. #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct Hello { - #[prost(uint32, tag = "1")] + #[prost(uint32, tag="1")] pub protocol_version: u32, - #[prost(string, tag = "2")] + #[prost(string, tag="2")] pub agent_version: ::prost::alloc::string::String, - #[prost(string, tag = "3")] + #[prost(string, tag="3")] pub boot_id: ::prost::alloc::string::String, - #[prost(string, tag = "4")] + #[prost(string, tag="4")] pub nix_version: ::prost::alloc::string::String, } #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct Init { - #[prost(string, tag = "1")] + #[prost(string, tag="1")] pub job_id: ::prost::alloc::string::String, - #[prost(string, repeated, tag = "2")] + #[prost(string, repeated, tag="2")] pub cache_trusted_public_keys: ::prost::alloc::vec::Vec<::prost::alloc::string::String>, - #[prost(uint32, tag = "3")] + #[prost(uint32, tag="3")] pub cache_read_proxy_port: u32, - #[prost(uint32, tag = "4")] + #[prost(uint32, tag="4")] pub cache_upload_proxy_port: u32, - #[prost(uint32, tag = "5")] + #[prost(uint32, tag="5")] pub dns_proxy_port: u32, } #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct ExecStart { - #[prost(string, repeated, tag = "1")] + #[prost(string, repeated, tag="1")] pub argv: ::prost::alloc::vec::Vec<::prost::alloc::string::String>, - #[prost(string, repeated, tag = "2")] + #[prost(string, repeated, tag="2")] pub env: ::prost::alloc::vec::Vec<::prost::alloc::string::String>, - #[prost(string, tag = "3")] + #[prost(string, tag="3")] pub cwd: ::prost::alloc::string::String, - #[prost(string, tag = "4")] + #[prost(string, tag="4")] pub user: ::prost::alloc::string::String, - #[prost(uint32, tag = "5")] + #[prost(uint32, tag="5")] pub timeout_seconds: u32, } #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct ExecStdout { - #[prost(string, tag = "1")] + #[prost(string, tag="1")] pub data: ::prost::alloc::string::String, } #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct ExecStderr { - #[prost(string, tag = "1")] + #[prost(string, tag="1")] pub data: ::prost::alloc::string::String, } #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct ExecExit { - #[prost(int32, tag = "1")] + #[prost(int32, tag="1")] pub exit_code: i32, - #[prost(string, tag = "2")] + #[prost(string, tag="2")] pub error: ::prost::alloc::string::String, /// set when the guest killed the step on its own timeout timer, so the host /// can classify it as a timeout rather than inferring failure from exit_code. - #[prost(bool, tag = "3")] + #[prost(bool, tag="3")] pub timed_out: bool, } #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct ActivateConfig { - #[prost(string, tag = "1")] + #[prost(string, tag="1")] pub config_key: ::prost::alloc::string::String, - #[prost(string, tag = "2")] + #[prost(string, tag="2")] pub base_config_hash: ::prost::alloc::string::String, - #[prost(string, tag = "3")] + #[prost(string, tag="3")] pub user_config: ::prost::alloc::string::String, - #[prost(string, tag = "4")] + #[prost(string, tag="4")] pub toplevel: ::prost::alloc::string::String, - #[prost(uint32, tag = "5")] + #[prost(uint32, tag="5")] pub timeout_seconds: u32, } #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct ActivateConfigResult { - #[prost(string, tag = "1")] + #[prost(string, tag="1")] pub config_key: ::prost::alloc::string::String, - #[prost(string, tag = "2")] + #[prost(string, tag="2")] pub toplevel: ::prost::alloc::string::String, - #[prost(string, tag = "3")] + #[prost(string, tag="3")] pub error: ::prost::alloc::string::String, } #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct BuiltPaths { - #[prost(string, repeated, tag = "1")] + #[prost(string, repeated, tag="1")] pub paths: ::prost::alloc::vec::Vec<::prost::alloc::string::String>, - #[prost(string, tag = "2")] + #[prost(string, tag="2")] pub reason: ::prost::alloc::string::String, } #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)] pub struct CacheDrain { - #[prost(uint32, tag = "1")] + #[prost(uint32, tag="1")] pub timeout_seconds: u32, } #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct CacheDrainResult { - #[prost(string, tag = "1")] + #[prost(string, tag="1")] pub error: ::prost::alloc::string::String, - #[prost(uint32, tag = "2")] + #[prost(uint32, tag="2")] pub cache_queued: u32, - #[prost(uint32, tag = "3")] + #[prost(uint32, tag="3")] pub cache_active: u32, - #[prost(uint32, tag = "4")] + #[prost(uint32, tag="4")] pub cache_uploaded: u32, - #[prost(uint32, tag = "5")] + #[prost(uint32, tag="5")] pub cache_failed: u32, } #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)] -pub struct Poweroff {} +pub struct Poweroff { +} #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct PoweroffResult { - #[prost(string, tag = "1")] + #[prost(string, tag="1")] pub error: ::prost::alloc::string::String, } #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct Message { - #[prost(string, tag = "1")] + #[prost(string, tag="1")] pub id: ::prost::alloc::string::String, - #[prost(message, optional, tag = "2")] + #[prost(message, optional, tag="2")] pub hello: ::core::option::Option, - #[prost(message, optional, tag = "3")] + #[prost(message, optional, tag="3")] pub init: ::core::option::Option, - #[prost(message, optional, tag = "4")] + #[prost(message, optional, tag="4")] pub exec_start: ::core::option::Option, - #[prost(message, optional, tag = "5")] + #[prost(message, optional, tag="5")] pub exec_stdout: ::core::option::Option, - #[prost(message, optional, tag = "6")] + #[prost(message, optional, tag="6")] pub exec_stderr: ::core::option::Option, - #[prost(message, optional, tag = "7")] + #[prost(message, optional, tag="7")] pub exec_exit: ::core::option::Option, - #[prost(message, optional, tag = "8")] + #[prost(message, optional, tag="8")] pub activate_config: ::core::option::Option, - #[prost(message, optional, tag = "9")] + #[prost(message, optional, tag="9")] pub activate_config_result: ::core::option::Option, - #[prost(message, optional, tag = "10")] + #[prost(message, optional, tag="10")] pub built_paths: ::core::option::Option, - #[prost(message, optional, tag = "11")] + #[prost(message, optional, tag="11")] pub cache_drain: ::core::option::Option, - #[prost(message, optional, tag = "12")] + #[prost(message, optional, tag="12")] pub cache_drain_result: ::core::option::Option, - #[prost(message, optional, tag = "13")] + #[prost(message, optional, tag="13")] pub poweroff: ::core::option::Option, - #[prost(message, optional, tag = "14")] + #[prost(message, optional, tag="14")] pub poweroff_result: ::core::option::Option, } // @@protoc_insertion_point(module) diff --git a/spindle/engine/placement.go b/spindle/engine/placement.go new file mode 100644 index 00000000..eff71ab3 --- /dev/null +++ b/spindle/engine/placement.go @@ -0,0 +1,49 @@ +package engine + +import ( + "fmt" + "slices" + "strings" + "unicode" + + "tangled.org/core/api/tangled" + "tangled.org/core/spindle/models" +) + +// PlacementPlanner derives the engine capabilities a workflow needs. Planners +// run in the trusted mill process and must not inspect executor-local state. +type PlacementPlanner interface { + Requirements(twf tangled.Pipeline_Workflow, tpl tangled.Pipeline) ([]string, error) +} + +// CapabilityProvider reports the cached placement capabilities of one local +// engine instance. Capabilities are opaque to the mill. +type CapabilityProvider interface { + Capabilities() []string +} + +// WorkflowPlacementValidator performs the definitive immutable compatibility +// checks before an executor accepts and holds a remote lease. +type WorkflowPlacementValidator interface { + ValidateWorkflowPlacement(wf *models.Workflow) error +} + +// NormalizeCapabilities validates, sorts, and deduplicates an opaque +// capability set so snapshots and workflow requirements compare deterministically. +func NormalizeCapabilities(capabilities []string) ([]string, error) { + seen := make(map[string]struct{}, len(capabilities)) + normalized := make([]string, 0, len(capabilities)) + for _, capability := range capabilities { + capability = strings.TrimSpace(capability) + if capability == "" || strings.ContainsFunc(capability, unicode.IsSpace) { + return nil, fmt.Errorf("invalid empty or whitespace-containing engine capability %q", capability) + } + if _, ok := seen[capability]; ok { + continue + } + seen[capability] = struct{}{} + normalized = append(normalized, capability) + } + slices.Sort(normalized) + return normalized, nil +} diff --git a/spindle/engine/placement_test.go b/spindle/engine/placement_test.go new file mode 100644 index 00000000..14a54bab --- /dev/null +++ b/spindle/engine/placement_test.go @@ -0,0 +1,22 @@ +package engine + +import "testing" + +func TestNormalizeCapabilities(t *testing.T) { + got, err := NormalizeCapabilities([]string{" runner/qemu ", "image/nixos", "runner/qemu"}) + if err != nil { + t.Fatalf("NormalizeCapabilities() error = %v", err) + } + want := []string{"image/nixos", "runner/qemu"} + if len(got) != len(want) || got[0] != want[0] || got[1] != want[1] { + t.Fatalf("NormalizeCapabilities() = %v, want %v", got, want) + } +} + +func TestNormalizeCapabilitiesRejectsInvalidAtoms(t *testing.T) { + for _, capabilities := range [][]string{{""}, {" "}, {"image/nixos arm"}, {"image/nixos\tarm"}} { + if got, err := NormalizeCapabilities(capabilities); err == nil { + t.Errorf("NormalizeCapabilities(%q) = %v, nil; want error", capabilities, got) + } + } +} diff --git a/spindle/engine_microvm_linux.go b/spindle/engine_microvm_linux.go new file mode 100644 index 00000000..5c5551b0 --- /dev/null +++ b/spindle/engine_microvm_linux.go @@ -0,0 +1,16 @@ +//go:build linux + +package spindle + +import ( + "context" + + "tangled.org/core/spindle/config" + "tangled.org/core/spindle/db" + "tangled.org/core/spindle/engines/microvm" + "tangled.org/core/spindle/models" +) + +func newMicrovmEngine(ctx context.Context, cfg *config.Config, d *db.DB) (models.Engine, error) { + return microvm.New(ctx, cfg, d) +} diff --git a/spindle/engine_microvm_other.go b/spindle/engine_microvm_other.go new file mode 100644 index 00000000..9da4c9fc --- /dev/null +++ b/spindle/engine_microvm_other.go @@ -0,0 +1,16 @@ +//go:build !linux + +package spindle + +import ( + "context" + "fmt" + + "tangled.org/core/spindle/config" + "tangled.org/core/spindle/db" + "tangled.org/core/spindle/models" +) + +func newMicrovmEngine(context.Context, *config.Config, *db.DB) (models.Engine, error) { + return nil, fmt.Errorf("microvm engine is only supported on Linux") +} diff --git a/spindle/engines/microvm/agent.go b/spindle/engines/microvm/agent.go index 4a6e01f3..b4342e56 100644 --- a/spindle/engines/microvm/agent.go +++ b/spindle/engines/microvm/agent.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/args.go b/spindle/engines/microvm/args.go index 05b370af..0b4f1c42 100644 --- a/spindle/engines/microvm/args.go +++ b/spindle/engines/microvm/args.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/budget.go b/spindle/engines/microvm/budget.go index d4749a52..a45d834e 100644 --- a/spindle/engines/microvm/budget.go +++ b/spindle/engines/microvm/budget.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/cgroup.go b/spindle/engines/microvm/cgroup.go index 40f6b736..c6d6970f 100644 --- a/spindle/engines/microvm/cgroup.go +++ b/spindle/engines/microvm/cgroup.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/cgroup_oom_test.go b/spindle/engines/microvm/cgroup_oom_test.go index 5754b96a..07857377 100644 --- a/spindle/engines/microvm/cgroup_oom_test.go +++ b/spindle/engines/microvm/cgroup_oom_test.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/cgroup_test.go b/spindle/engines/microvm/cgroup_test.go index e8c54108..dadb4518 100644 --- a/spindle/engines/microvm/cgroup_test.go +++ b/spindle/engines/microvm/cgroup_test.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/dns_proxy.go b/spindle/engines/microvm/dns_proxy.go index 16e6183a..e6b99a0c 100644 --- a/spindle/engines/microvm/dns_proxy.go +++ b/spindle/engines/microvm/dns_proxy.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/dns_proxy_test.go b/spindle/engines/microvm/dns_proxy_test.go index a38b0f4e..cb1de34c 100644 --- a/spindle/engines/microvm/dns_proxy_test.go +++ b/spindle/engines/microvm/dns_proxy_test.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/engine.go b/spindle/engines/microvm/engine.go index 07e2e575..5fe6a1e7 100644 --- a/spindle/engines/microvm/engine.go +++ b/spindle/engines/microvm/engine.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( @@ -49,6 +51,8 @@ type Engine struct { agent *agentHub scheduler *engine.ResourceScheduler[Resources] cgroupParent *CgroupParent + budget Resources + maxWorkflow Resources cleanupMu sync.Mutex cleanup map[string][]cleanupFunc @@ -96,6 +100,8 @@ func New(ctx context.Context, cfg *config.Config, d *db.DB) (*Engine, error) { agent: agent, scheduler: engine.NewResourceScheduler(budget, max, agingThreshold), cgroupParent: cgroupParent, + budget: budget, + maxWorkflow: max, cleanup: make(map[string][]cleanupFunc), }, nil } diff --git a/spindle/engines/microvm/engine_test.go b/spindle/engines/microvm/engine_test.go index 0585121c..b5aa58d5 100644 --- a/spindle/engines/microvm/engine_test.go +++ b/spindle/engines/microvm/engine_test.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/image.go b/spindle/engines/microvm/image.go index 89db1552..5e5eca89 100644 --- a/spindle/engines/microvm/image.go +++ b/spindle/engines/microvm/image.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( @@ -7,6 +9,8 @@ import ( "os" "path/filepath" "strings" + + "tangled.org/core/spindle/engines/microvm/placement" ) const imageSpecFileName = "spec.json" @@ -202,13 +206,7 @@ func (e *Engine) resolveImage(name string) (ImageSpec, string, string, error) { // check if image name is not a path func isPlainImageName(name string) bool { - if name == "" || name == "." || name == ".." { - return false - } - if filepath.IsAbs(name) || strings.ContainsRune(name, '/') || strings.ContainsRune(name, filepath.Separator) { - return false - } - return true + return placement.IsPlainImageName(name) } // returns candidates, which is either a directory or spec file itself diff --git a/spindle/engines/microvm/image_test.go b/spindle/engines/microvm/image_test.go index 3000005c..7688f8d4 100644 --- a/spindle/engines/microvm/image_test.go +++ b/spindle/engines/microvm/image_test.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/models.go b/spindle/engines/microvm/models.go index 02f3c096..f80772c2 100644 --- a/spindle/engines/microvm/models.go +++ b/spindle/engines/microvm/models.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/models_test.go b/spindle/engines/microvm/models_test.go index b103e4f5..256fe1ec 100644 --- a/spindle/engines/microvm/models_test.go +++ b/spindle/engines/microvm/models_test.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/networking.go b/spindle/engines/microvm/networking.go index 30f12dfe..81597cbe 100644 --- a/spindle/engines/microvm/networking.go +++ b/spindle/engines/microvm/networking.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/nixos_toplevel_cache.go b/spindle/engines/microvm/nixos_toplevel_cache.go index 8333a65c..c5f8f784 100644 --- a/spindle/engines/microvm/nixos_toplevel_cache.go +++ b/spindle/engines/microvm/nixos_toplevel_cache.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/placement/placement.go b/spindle/engines/microvm/placement/placement.go new file mode 100644 index 00000000..9990673e --- /dev/null +++ b/spindle/engines/microvm/placement/placement.go @@ -0,0 +1,66 @@ +package placement + +import ( + "fmt" + "strings" + "unicode" + + "gopkg.in/yaml.v3" + + "tangled.org/core/api/tangled" +) + +const DefaultImageCapability = "image/default" + +type Planner struct{} + +func (Planner) Requirements(twf tangled.Pipeline_Workflow, _ tangled.Pipeline) ([]string, error) { + var manifest struct { + Image string `yaml:"image"` + } + if err := yaml.Unmarshal([]byte(twf.Raw), &manifest); err != nil { + return nil, fmt.Errorf("parse microVM workflow placement: %w", err) + } + + image := strings.TrimSpace(manifest.Image) + if image == "" { + return []string{DefaultImageCapability}, nil + } + capability, err := ImageCapability(image) + if err != nil { + return nil, err + } + return []string{capability}, nil +} + +func ImageCapability(name string) (string, error) { + name = strings.TrimSpace(name) + if !IsPlainImageName(name) { + return "", fmt.Errorf("invalid microVM image name %q: must be a whitespace-free plain name, not a path", name) + } + return "image/" + name, nil +} + +func IsPlainImageName(name string) bool { + if name == "" || name == "." || name == ".." { + return false + } + if strings.ContainsAny(name, `/\\`) || strings.ContainsFunc(name, unicode.IsSpace) { + return false + } + return true +} + +func IsNativeArchitecture(imageArch, goArch string) bool { + normalize := func(arch string) string { + switch arch { + case "x86_64", "amd64": + return "amd64" + case "aarch64", "arm64": + return "arm64" + default: + return arch + } + } + return imageArch != "" && normalize(imageArch) == normalize(goArch) +} diff --git a/spindle/engines/microvm/placement/placement_test.go b/spindle/engines/microvm/placement/placement_test.go new file mode 100644 index 00000000..afd15e11 --- /dev/null +++ b/spindle/engines/microvm/placement/placement_test.go @@ -0,0 +1,78 @@ +package placement + +import ( + "testing" + + "tangled.org/core/api/tangled" +) + +func TestPlannerRequirements(t *testing.T) { + tests := []struct { + name string + raw string + want string + wantErr bool + }{ + {name: "explicit image", raw: "image: nixos\n", want: "image/nixos"}, + {name: "explicit arm image", raw: "image: nixos-aarch64\n", want: "image/nixos-aarch64"}, + {name: "omitted image", raw: "steps: []\n", want: DefaultImageCapability}, + {name: "blank image", raw: "image: ' '\n", want: DefaultImageCapability}, + {name: "path", raw: "image: ../nixos\n", wantErr: true}, + {name: "whitespace", raw: "image: 'nixos arm'\n", wantErr: true}, + {name: "wrong yaml type", raw: "image: [nixos]\n", wantErr: true}, + {name: "malformed yaml", raw: "image: [\n", wantErr: true}, + } + + planner := Planner{} + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got, err := planner.Requirements(tangled.Pipeline_Workflow{Raw: tt.raw}, tangled.Pipeline{}) + if tt.wantErr { + if err == nil { + t.Fatalf("Requirements() = %v, nil; want error", got) + } + return + } + if err != nil { + t.Fatalf("Requirements() error = %v", err) + } + if len(got) != 1 || got[0] != tt.want { + t.Fatalf("Requirements() = %v, want [%s]", got, tt.want) + } + }) + } +} + +func TestIsPlainImageName(t *testing.T) { + for _, name := range []string{"nixos", "nixos-aarch64", "release.2026_07"} { + if !IsPlainImageName(name) { + t.Errorf("IsPlainImageName(%q) = false, want true", name) + } + } + for _, name := range []string{"", ".", "..", "/nixos", "../nixos", `dir\\nixos`, "nixos arm", "nixos\tarm"} { + if IsPlainImageName(name) { + t.Errorf("IsPlainImageName(%q) = true, want false", name) + } + } +} + +func TestIsNativeArchitecture(t *testing.T) { + tests := []struct { + image string + host string + want bool + }{ + {image: "x86_64", host: "amd64", want: true}, + {image: "amd64", host: "amd64", want: true}, + {image: "aarch64", host: "arm64", want: true}, + {image: "arm64", host: "arm64", want: true}, + {image: "x86_64", host: "arm64", want: false}, + {image: "aarch64", host: "amd64", want: false}, + {image: "", host: "amd64", want: false}, + } + for _, tt := range tests { + if got := IsNativeArchitecture(tt.image, tt.host); got != tt.want { + t.Errorf("IsNativeArchitecture(%q, %q) = %v, want %v", tt.image, tt.host, got, tt.want) + } + } +} diff --git a/spindle/engines/microvm/placement_linux.go b/spindle/engines/microvm/placement_linux.go new file mode 100644 index 00000000..38911f3d --- /dev/null +++ b/spindle/engines/microvm/placement_linux.go @@ -0,0 +1,109 @@ +//go:build linux + +package microvm + +import ( + "fmt" + "os" + "os/exec" + "runtime" + "slices" + "strings" + + "tangled.org/core/spindle/engines/microvm/placement" + "tangled.org/core/spindle/models" +) + +func (e *Engine) Capabilities() []string { + entries, err := os.ReadDir(e.cfg.MicroVMPipelines.ImageDir) + if err != nil { + e.l.Error("discover microVM placement capabilities", "err", err) + return nil + } + + seen := make(map[string]struct{}, len(entries)+1) + capabilities := make([]string, 0, len(entries)+1) + addImage := func(name string) bool { + name = strings.TrimSuffix(name, ".json") + if !placement.IsPlainImageName(name) { + return false + } + capability, err := placement.ImageCapability(name) + if err != nil { + return false + } + if _, ok := seen[capability]; ok { + return true + } + spec, _, _, err := e.resolveImage(name) + if err != nil { + e.l.Debug("microVM image is unavailable for placement", "image", name, "err", err) + return false + } + if err := e.validateImagePlacement(spec); err != nil { + e.l.Debug("microVM image is incompatible with executor", "image", name, "err", err) + return false + } + seen[capability] = struct{}{} + capabilities = append(capabilities, capability) + return true + } + + for _, entry := range entries { + addImage(entry.Name()) + } + if defaultImage := strings.TrimSpace(e.cfg.MicroVMPipelines.DefaultImage); defaultImage != "" && addImage(defaultImage) { + capabilities = append(capabilities, placement.DefaultImageCapability) + } + slices.Sort(capabilities) + return capabilities +} + +func (e *Engine) ValidateWorkflowPlacement(wf *models.Workflow) error { + state, ok := wf.Data.(*workflowState) + if !ok || state == nil { + return fmt.Errorf("microVM workflow state is not initialized") + } + return e.validateImagePlacement(state.ImageSpec) +} + +func (e *Engine) validateImagePlacement(spec ImageSpec) error { + if err := spec.Validate(); err != nil { + return err + } + if !placement.IsNativeArchitecture(spec.Arch, runtime.GOARCH) { + return fmt.Errorf("microVM image architecture %q is not native to executor architecture %q", spec.Arch, runtime.GOARCH) + } + if err := spec.validateImageFiles(); err != nil { + return err + } + runner, err := runnerFor(spec.RunnerType) + if err != nil { + return err + } + if err := runner.Validate(spec, e.cfg.MicroVMPipelines.EnableKVM); err != nil { + return err + } + if len(spec.Volumes) > 0 { + if _, err := exec.LookPath("mkfs.ext4"); err != nil { + return fmt.Errorf("required host command %q not found in PATH: %w", "mkfs.ext4", err) + } + for _, volume := range spec.Volumes { + if volume.ReadOnly { + return fmt.Errorf("read-only microvm volume %q is not supported yet", volume.Image) + } + if volume.FSType != "ext4" { + return fmt.Errorf("microvm volume %q uses unsupported fsType %q", volume.Image, volume.FSType) + } + if volume.ImageType != "" && volume.ImageType != "raw" { + return fmt.Errorf("microvm volume %q uses unsupported imageType %q", volume.Image, volume.ImageType) + } + } + } + + request := resourcesForImage(spec) + if !request.Fits(e.budget) || !request.Fits(e.maxWorkflow) { + return fmt.Errorf("microVM image resources exceed executor limits: request=%v budget=%v max=%v", request, e.budget, e.maxWorkflow) + } + return nil +} diff --git a/spindle/engines/microvm/qemu.go b/spindle/engines/microvm/qemu.go index c9d5c426..c4a61a68 100644 --- a/spindle/engines/microvm/qemu.go +++ b/spindle/engines/microvm/qemu.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/read_cache_proxy.go b/spindle/engines/microvm/read_cache_proxy.go index 2d1d749f..311945bc 100644 --- a/spindle/engines/microvm/read_cache_proxy.go +++ b/spindle/engines/microvm/read_cache_proxy.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/read_cache_proxy_test.go b/spindle/engines/microvm/read_cache_proxy_test.go index 9b92b137..47e8f207 100644 --- a/spindle/engines/microvm/read_cache_proxy_test.go +++ b/spindle/engines/microvm/read_cache_proxy_test.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/runner.go b/spindle/engines/microvm/runner.go index c08e3d84..4ec712d1 100644 --- a/spindle/engines/microvm/runner.go +++ b/spindle/engines/microvm/runner.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/upload_cache_http.go b/spindle/engines/microvm/upload_cache_http.go index 0310ef09..b479616f 100644 --- a/spindle/engines/microvm/upload_cache_http.go +++ b/spindle/engines/microvm/upload_cache_http.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/upload_cache_narinfo.go b/spindle/engines/microvm/upload_cache_narinfo.go index 26645fcf..775d20bc 100644 --- a/spindle/engines/microvm/upload_cache_narinfo.go +++ b/spindle/engines/microvm/upload_cache_narinfo.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/upload_cache_nix_store.go b/spindle/engines/microvm/upload_cache_nix_store.go index 70b9adde..8ec22c47 100644 --- a/spindle/engines/microvm/upload_cache_nix_store.go +++ b/spindle/engines/microvm/upload_cache_nix_store.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/upload_cache_nix_store_test.go b/spindle/engines/microvm/upload_cache_nix_store_test.go index 413fb66d..bf9d4af7 100644 --- a/spindle/engines/microvm/upload_cache_nix_store_test.go +++ b/spindle/engines/microvm/upload_cache_nix_store_test.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/upload_cache_proxy.go b/spindle/engines/microvm/upload_cache_proxy.go index 9b9385d2..4cd65f90 100644 --- a/spindle/engines/microvm/upload_cache_proxy.go +++ b/spindle/engines/microvm/upload_cache_proxy.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/upload_cache_proxy_test.go b/spindle/engines/microvm/upload_cache_proxy_test.go index ef4e6405..5697f12f 100644 --- a/spindle/engines/microvm/upload_cache_proxy_test.go +++ b/spindle/engines/microvm/upload_cache_proxy_test.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/engines/microvm/vm.go b/spindle/engines/microvm/vm.go index a27f9020..039d0532 100644 --- a/spindle/engines/microvm/vm.go +++ b/spindle/engines/microvm/vm.go @@ -1,3 +1,5 @@ +//go:build linux + package microvm import ( diff --git a/spindle/mill/capability_test.go b/spindle/mill/capability_test.go new file mode 100644 index 00000000..5b3f37e8 --- /dev/null +++ b/spindle/mill/capability_test.go @@ -0,0 +1,99 @@ +package mill + +import ( + "errors" + "io" + "log/slog" + "testing" + + "tangled.org/core/api/tangled" + "tangled.org/core/spindle/models" + + millproto "tangled.org/core/spindle/mill/proto" +) + +type testPlacementPlanner struct { + requirements []string + err error +} + +func (p testPlacementPlanner) Requirements(tangled.Pipeline_Workflow, tangled.Pipeline) ([]string, error) { + return p.requirements, p.err +} + +func TestMillEngineStoresCanonicalPlacementRequirements(t *testing.T) { + m := New(slog.New(slog.NewTextHandler(io.Discard, nil)), Config{}) + eng := NewEngineWithPlanner("microvm", m, testPlacementPlanner{ + requirements: []string{"runner/qemu", "image/nixos", "runner/qemu"}, + }) + + wf, err := eng.InitWorkflow(tangled.Pipeline_Workflow{Name: "build"}, tangled.Pipeline{}) + if err != nil { + t.Fatalf("InitWorkflow() error = %v", err) + } + got := wf.Data.(*millWorkflowState).RequiredCapabilities + want := []string{"image/nixos", "runner/qemu"} + if len(got) != len(want) || got[0] != want[0] || got[1] != want[1] { + t.Fatalf("required capabilities = %v, want %v", got, want) + } +} + +func TestMillEngineRejectsInvalidPlacementRequirements(t *testing.T) { + m := New(slog.New(slog.NewTextHandler(io.Discard, nil)), Config{}) + plannerErr := errors.New("invalid image") + eng := NewEngineWithPlanner("microvm", m, testPlacementPlanner{err: plannerErr}) + if _, err := eng.InitWorkflow(tangled.Pipeline_Workflow{}, tangled.Pipeline{}); !errors.Is(err, plannerErr) { + t.Fatalf("InitWorkflow() error = %v, want wrapped %v", err, plannerErr) + } + + eng = NewEngineWithPlanner("microvm", m, testPlacementPlanner{requirements: []string{"image/nixos arm"}}) + if _, err := eng.InitWorkflow(tangled.Pipeline_Workflow{}, tangled.Pipeline{}); err == nil { + t.Fatal("InitWorkflow() accepted whitespace-containing capability") + } +} + +func TestRankCandidatesIntersectsLabelsAndEngineCapabilities(t *testing.T) { + m := New(slog.New(slog.NewTextHandler(io.Discard, nil)), Config{}) + incapable := addCandidateSession(t, m, "low-load-incapable", []string{"region/eu"}, 0.05, nil) + incapable.snapshot.Engines["dummy"].Capabilities = []string{"image/nixos-aarch64"} + wrongRegion := addCandidateSession(t, m, "wrong-region", []string{"region/us"}, 0.10, nil) + wrongRegion.snapshot.Engines["dummy"].Capabilities = []string{"image/nixos", "runner/qemu"} + partial := addCandidateSession(t, m, "partial", []string{"region/eu"}, 0.20, nil) + partial.snapshot.Engines["dummy"].Capabilities = []string{"image/nixos"} + eligible := addCandidateSession(t, m, "eligible", []string{"region/eu"}, 0.70, nil) + eligible.snapshot.Engines["dummy"].Capabilities = []string{"runner/qemu", "image/nixos"} + + got := m.rankCandidates("dummy", []string{"region/eu"}, []string{"image/nixos", "runner/qemu"}) + assertRankedNodes(t, got, []string{"eligible"}) +} + +func TestBidWithMissingEngineCapabilitySendsNoReserve(t *testing.T) { + m := New(slog.New(slog.NewTextHandler(io.Discard, nil)), Config{}) + reserveSent := false + sess := addCandidateSession(t, m, "arm", nil, 0.05, scriptedEncoder(func(msg *millproto.Message) error { + if msg.GetReserveSeat() != nil { + reserveSent = true + } + return nil + })) + sess.snapshot.Engines["dummy"].Capabilities = []string{"image/nixos-aarch64"} + wf := testWorkflow("build") + wf.Data.(*millWorkflowState).RequiredCapabilities = []string{"image/nixos"} + + lease, err := m.bid(t.Context(), "dummy", wf) + if err != nil { + t.Fatalf("bid() error = %v", err) + } + if lease != nil { + t.Fatalf("bid() lease = %v, want nil", lease) + } + if reserveSent { + t.Fatal("bid() sent ReserveSeat to an engine missing a required capability") + } +} + +func TestRequiredCapabilitiesWithoutMillStateAreEmpty(t *testing.T) { + if got := requiredCapabilities(&models.Workflow{}); got != nil { + t.Fatalf("requiredCapabilities() = %v, want nil", got) + } +} diff --git a/spindle/mill/engine.go b/spindle/mill/engine.go index 7dcb6f89..e631cff2 100644 --- a/spindle/mill/engine.go +++ b/spindle/mill/engine.go @@ -2,6 +2,7 @@ package mill import ( "context" + "fmt" "log/slog" "time" @@ -15,38 +16,60 @@ import ( // engines use it). InitWorkflow parses nothing here: it just carries the raw // pipeline/workflow forward so the executor can run the real InitWorkflow later. type millWorkflowState struct { - TargetEngine string - RawWorkflow tangled.Pipeline_Workflow - RawPipeline tangled.Pipeline - Wid models.WorkflowId // stamped at placement (InitWorkflow can't see it) - Lease *RemoteLease + TargetEngine string + RawWorkflow tangled.Pipeline_Workflow + RawPipeline tangled.Pipeline + RequiredCapabilities []string + Wid models.WorkflowId // stamped at placement (InitWorkflow can't see it) + Lease *RemoteLease } // Engine is the mill's stand-in for a real engine, registered under the real // engine names ("microvm", "nixery"). All registered names share one Mill. type Engine struct { - name string - mill *Mill - l *slog.Logger + name string + mill *Mill + planner engine.PlacementPlanner + l *slog.Logger } -// NewEngine returns a mill engine view for one engine name. +// NewEngine returns a mill engine view for one engine name without +// engine-derived placement requirements. func NewEngine(name string, mill *Mill) *Engine { - return &Engine{name: name, mill: mill, l: mill.l.With("engine", "mill:"+name)} + return NewEngineWithPlanner(name, mill, nil) +} + +// NewEngineWithPlanner returns a mill engine view whose trusted planner +// derives requirements for capability-aware placement. +func NewEngineWithPlanner(name string, mill *Mill, planner engine.PlacementPlanner) *Engine { + return &Engine{name: name, mill: mill, planner: planner, l: mill.l.With("engine", "mill:"+name)} } // InitWorkflow returns a synthetic one-step workflow so processPipeline injects // TANGLED_* env and marks pending normally. The real InitWorkflow runs later on // the executor inside ReserveSeat (intended: it runs twice). func (e *Engine) InitWorkflow(twf tangled.Pipeline_Workflow, tpl tangled.Pipeline) (*models.Workflow, error) { + var required []string + if e.planner != nil { + var err error + required, err = e.planner.Requirements(twf, tpl) + if err != nil { + return nil, fmt.Errorf("derive %s placement requirements: %w", e.name, err) + } + required, err = engine.NormalizeCapabilities(required) + if err != nil { + return nil, fmt.Errorf("derive %s placement requirements: %w", e.name, err) + } + } return &models.Workflow{ Name: twf.Name, Environment: map[string]string{}, Steps: []models.Step{remoteStep{}}, Data: &millWorkflowState{ - TargetEngine: e.name, - RawWorkflow: twf, - RawPipeline: tpl, + TargetEngine: e.name, + RawWorkflow: twf, + RawPipeline: tpl, + RequiredCapabilities: required, }, }, nil } diff --git a/spindle/mill/executor/capability_test.go b/spindle/mill/executor/capability_test.go new file mode 100644 index 00000000..0b6b6f5d --- /dev/null +++ b/spindle/mill/executor/capability_test.go @@ -0,0 +1,106 @@ +package executor + +import ( + "context" + "encoding/json" + "errors" + "io" + "log/slog" + "testing" + + "tangled.org/core/api/tangled" + millv1 "tangled.org/core/spindle/mill/proto/gen" + "tangled.org/core/spindle/models" +) + +type placementEngine struct { + *fakeEngine + capabilities []string + validationErr error + validated bool +} + +func (e *placementEngine) Capabilities() []string { + return e.capabilities +} + +func (e *placementEngine) ValidateWorkflowPlacement(*models.Workflow) error { + e.validated = true + return e.validationErr +} + +func TestPushSnapshotAdvertisesCanonicalPerEngineCapabilities(t *testing.T) { + enc := newCaptureEncoder() + eng := &placementEngine{ + fakeEngine: &fakeEngine{}, + capabilities: []string{"runner/qemu", "image/nixos", "runner/qemu"}, + } + e := &Executor{ + l: slog.New(slog.NewTextHandler(io.Discard, nil)), + enc: enc, + seats: 2, + engines: map[string]models.Engine{"microvm": eng}, + active: make(map[string]*reservation), + } + + e.pushSnapshot() + availability := (<-enc.messages).GetNodeSnapshot().GetEngines()["microvm"] + if availability == nil { + t.Fatal("snapshot omitted microvm availability") + } + got := availability.GetCapabilities() + want := []string{"image/nixos", "runner/qemu"} + if len(got) != len(want) || got[0] != want[0] || got[1] != want[1] { + t.Fatalf("snapshot capabilities = %v, want %v", got, want) + } +} + +func TestHandleReserveValidatesPlacementBeforeAcquiringSlot(t *testing.T) { + enc := newCaptureEncoder() + validationErr := errors.New("image architecture is not native") + eng := &placementEngine{fakeEngine: &fakeEngine{}, validationErr: validationErr} + e := &Executor{ + l: slog.New(slog.NewTextHandler(io.Discard, nil)), + enc: enc, + seats: 1, + engines: map[string]models.Engine{"microvm": eng}, + active: make(map[string]*reservation), + } + twf, err := json.Marshal(tangled.Pipeline_Workflow{Name: "build"}) + if err != nil { + t.Fatal(err) + } + tpl, err := json.Marshal(tangled.Pipeline{}) + if err != nil { + t.Fatal(err) + } + + e.handleReserve(context.Background(), &millv1.ReserveSeat{ + LeaseId: "lease-1", + TargetEngine: "microvm", + RawWorkflowJson: string(twf), + RawPipelineJson: string(tpl), + Knot: "k", + Rkey: "r", + }) + + result := (<-enc.messages).GetReserveResult() + if result == nil { + t.Fatal("handleReserve() did not send ReserveResult") + } + if result.GetAccepted() { + t.Fatal("handleReserve() accepted placement validation failure") + } + if result.GetRejectClass() != millv1.RejectClass_REJECT_CLASS_INCOMPATIBLE { + t.Fatalf("reject class = %v, want incompatible", result.GetRejectClass()) + } + if !eng.validated { + t.Fatal("placement validator was not called") + } + if eng.acquireCalled { + t.Fatal("slot acquisition ran after placement validation failed") + } + if len(e.active) != 0 { + t.Fatalf("active reservations = %d, want 0", len(e.active)) + } +} diff --git a/spindle/mill/executor/executor.go b/spindle/mill/executor/executor.go index c217b69f..f93a5538 100644 --- a/spindle/mill/executor/executor.go +++ b/spindle/mill/executor/executor.go @@ -264,6 +264,12 @@ func (e *Executor) handleReserve(ctx context.Context, rs *millv1.ReserveSeat) { reject("init workflow: "+err.Error(), millv1.RejectClass_REJECT_CLASS_INCOMPATIBLE) return } + if validator, ok := realEngine.(engine.WorkflowPlacementValidator); ok { + if err := validator.ValidateWorkflowPlacement(wf); err != nil { + reject("validate workflow placement: "+err.Error(), millv1.RejectClass_REJECT_CLASS_INCOMPATIBLE) + return + } + } // the job skipped processPipeline, so inject TANGLED_* env here. if wf.Environment == nil { wf.Environment = make(map[string]string) @@ -480,10 +486,20 @@ func (e *Executor) pushSnapshot() { } available := !draining && active < e.seats engines := make(map[string]*millv1.EngineAvailability, len(e.engines)) - for name := range e.engines { + for name, realEngine := range e.engines { + var capabilities []string + if provider, ok := realEngine.(engine.CapabilityProvider); ok { + var err error + capabilities, err = engine.NormalizeCapabilities(provider.Capabilities()) + if err != nil { + e.l.Error("engine advertised invalid placement capabilities", "engine", name, "err", err) + capabilities = nil + } + } engines[name] = &millv1.EngineAvailability{ - Available: available, - Load: map[string]float64{"slots": load}, + Available: available, + Load: map[string]float64{"slots": load}, + Capabilities: capabilities, } } e.send(&millproto.Message{NodeSnapshot: &millv1.NodeSnapshot{ diff --git a/spindle/mill/mill.go b/spindle/mill/mill.go index 69a9d6ea..6c84d403 100644 --- a/spindle/mill/mill.go +++ b/spindle/mill/mill.go @@ -283,7 +283,8 @@ func (m *Mill) bid(ctx context.Context, engineName string, wf *models.Workflow) } requiredLabels := requiredLabels(wf) - candidates := m.rankCandidates(engineName, requiredLabels) + requiredCapabilities := requiredCapabilities(wf) + candidates := m.rankCandidates(engineName, requiredLabels, requiredCapabilities) if len(candidates) == 0 { return nil, nil } @@ -390,7 +391,7 @@ func (m *Mill) bid(ctx context.Context, engineName string, wf *models.Workflow) return winner.lease, nil } -func (m *Mill) rankCandidates(engineName string, requiredLabels []string) []*millSession { +func (m *Mill) rankCandidates(engineName string, requiredLabels, requiredCapabilities []string) []*millSession { m.mu.Lock() defer m.mu.Unlock() @@ -414,6 +415,9 @@ func (m *Mill) rankCandidates(engineName string, requiredLabels []string) []*mil if !hasLabels(s.labels, requiredLabels) { continue } + if !hasLabels(ea.GetCapabilities(), requiredCapabilities) { + continue + } worst, sum := loadScore(ea.GetLoad()) rs = append(rs, ranked{sess: s, worst: worst, sum: sum}) } @@ -458,6 +462,14 @@ func requiredLabels(wf *models.Workflow) []string { return st.RawWorkflow.RunsOn } +func requiredCapabilities(wf *models.Workflow) []string { + st, ok := wf.Data.(*millWorkflowState) + if !ok || st == nil { + return nil + } + return st.RequiredCapabilities +} + func hasLabels(labels []string, required []string) bool { for _, want := range required { if !slices.Contains(labels, want) { diff --git a/spindle/mill/mill_test.go b/spindle/mill/mill_test.go index 33b4b013..a26c031c 100644 --- a/spindle/mill/mill_test.go +++ b/spindle/mill/mill_test.go @@ -281,7 +281,7 @@ func TestRankCandidatesFiltersRequiredLabelsWithANDSemantics(t *testing.T) { for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - assertRankedNodes(t, m.rankCandidates("dummy", tt.requiredLabels), tt.want) + assertRankedNodes(t, m.rankCandidates("dummy", tt.requiredLabels, nil), tt.want) }) } } diff --git a/spindle/mill/proto/gen/mill.pb.go b/spindle/mill/proto/gen/mill.pb.go index 9cbf7223..a393ee13 100644 --- a/spindle/mill/proto/gen/mill.pb.go +++ b/spindle/mill/proto/gen/mill.pb.go @@ -190,7 +190,10 @@ type EngineAvailability struct { state protoimpl.MessageState `protogen:"open.v1"` Available bool `protobuf:"varint,1,opt,name=available,proto3" json:"available,omitempty"` // engine-defined load metrics. keys are opaque to the mill. - Load map[string]float64 `protobuf:"bytes,2,rep,name=load,proto3" json:"load,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"fixed64,2,opt,name=value"` + Load map[string]float64 `protobuf:"bytes,2,rep,name=load,proto3" json:"load,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"fixed64,2,opt,name=value"` + // engine-defined placement facts. keys are opaque to the mill and matched + // exactly against requirements from the trusted engine planner. + Capabilities []string `protobuf:"bytes,3,rep,name=capabilities,proto3" json:"capabilities,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } @@ -239,6 +242,13 @@ func (x *EngineAvailability) GetLoad() map[string]float64 { return nil } +func (x *EngineAvailability) GetCapabilities() []string { + if x != nil { + return x.Capabilities + } + return nil +} + // NodeSnapshot is pushed on connect, periodically, and right after any state // change (reserve, commit, terminal). type NodeSnapshot struct { @@ -1154,10 +1164,11 @@ const file_spindle_mill_v1_mill_proto_rawDesc = "" + "\x06labels\x18\x03 \x03(\tR\x06labels\"'\n" + "\x06Resume\x12\x1d\n" + "\n" + - "ack_offset\x18\x01 \x01(\x04R\tackOffset\"\xae\x01\n" + + "ack_offset\x18\x01 \x01(\x04R\tackOffset\"\xd2\x01\n" + "\x12EngineAvailability\x12\x1c\n" + "\tavailable\x18\x01 \x01(\bR\tavailable\x12A\n" + - "\x04load\x18\x02 \x03(\v2-.spindle.mill.v1.EngineAvailability.LoadEntryR\x04load\x1a7\n" + + "\x04load\x18\x02 \x03(\v2-.spindle.mill.v1.EngineAvailability.LoadEntryR\x04load\x12\"\n" + + "\fcapabilities\x18\x03 \x03(\tR\fcapabilities\x1a7\n" + "\tLoadEntry\x12\x10\n" + "\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" + "\x05value\x18\x02 \x01(\x01R\x05value:\x028\x01\"\xe0\x01\n" + diff --git a/spindle/mill/proto/spindle/mill/v1/mill.proto b/spindle/mill/proto/spindle/mill/v1/mill.proto index 330ef89c..f842d3ce 100644 --- a/spindle/mill/proto/spindle/mill/v1/mill.proto +++ b/spindle/mill/proto/spindle/mill/v1/mill.proto @@ -31,6 +31,9 @@ message EngineAvailability { bool available = 1; // engine-defined load metrics. keys are opaque to the mill. map load = 2; + // engine-defined placement facts. keys are opaque to the mill and matched + // exactly against requirements from the trusted engine planner. + repeated string capabilities = 3; } // NodeSnapshot is pushed on connect, periodically, and right after any state diff --git a/spindle/server.go b/spindle/server.go index 04d356f1..b73c7901 100644 --- a/spindle/server.go +++ b/spindle/server.go @@ -33,7 +33,7 @@ import ( "tangled.org/core/spindle/db" "tangled.org/core/spindle/engine" "tangled.org/core/spindle/engines/dummy" - "tangled.org/core/spindle/engines/microvm" + microvmplacement "tangled.org/core/spindle/engines/microvm/placement" "tangled.org/core/spindle/engines/nixery" "tangled.org/core/spindle/git" "tangled.org/core/spindle/mill" @@ -384,7 +384,7 @@ func Run(ctx context.Context) error { }) engines = map[string]models.Engine{ "nixery": mill.NewEngine("nixery", m), - "microvm": mill.NewEngine("microvm", m), + "microvm": mill.NewEngineWithPlanner("microvm", m, microvmplacement.Planner{}), "dummy": mill.NewEngine("dummy", m), } } else { @@ -393,7 +393,7 @@ func Run(ctx context.Context) error { if err != nil { return err } - microvmEng, err := microvm.New(ctx, cfg, d) + microvmEng, err := newMicrovmEngine(ctx, cfg, d) if err != nil { return err }