From 60303d11dd99a1be0b885019222ead151b5fc1a2 Mon Sep 17 00:00:00 2001 From: Nikolay Edigaryev Date: Tue, 11 Nov 2025 21:16:28 +0400 Subject: [PATCH] VM specification: allow suspendable VMs (#366) --- api/openapi.yaml | 13 +++ internal/command/create/vm.go | 5 + internal/controller/api_vms.go | 10 ++ internal/tests/spec_update_test.go | 71 ++++++++++++- internal/worker/vmmanager/condition.go | 10 ++ internal/worker/vmmanager/vm.go | 136 +++++++++++++++++++------ internal/worker/worker.go | 50 ++++++--- pkg/resource/v1/v1.go | 1 + 8 files changed, 247 insertions(+), 49 deletions(-) create mode 100644 internal/worker/vmmanager/condition.go diff --git a/api/openapi.yaml b/api/openapi.yaml index 765b3c4..111dac3 100644 --- a/api/openapi.yaml +++ b/api/openapi.yaml @@ -520,6 +520,19 @@ components: - "66.66.0.0/16" items: type: string + suspendable: + type: boolean + description: | + When set, a VM will be started with an additional + `--suspendable` command-line argument to `tart run`, + which allows suspending it. + + Further generations of the VM will be `tart suspend`'ed instead of + `tart stopped`. + + For example, this allows you to prepare a VM with loose Softnet settings and then + move to the next generation by tightening the settings while preserving the VM's state. + default: false net-bridged: type: string description: Whether to use bridged network mode diff --git a/internal/command/create/vm.go b/internal/command/create/vm.go index c35126a..353f86d 100644 --- a/internal/command/create/vm.go +++ b/internal/command/create/vm.go @@ -25,6 +25,7 @@ var netSoftnetBlock []string var netBridged string var headless bool var nested bool +var suspendable bool var username string var password string var resources map[string]string @@ -58,6 +59,9 @@ func newCreateVMCommand() *cobra.Command { command.Flags().StringVar(&netBridged, "net-bridged", "", "whether to use Bridged network mode") command.Flags().BoolVar(&headless, "headless", true, "whether to run without graphics") command.Flags().BoolVar(&nested, "nested", false, "enable nested virtualization") + command.Flags().BoolVar(&suspendable, "suspendable", false, "treat the VM as suspendable, "+ + "disabling certain devices for suspendability support and issuing \"tart suspend\" instead of \"tart stop\" "+ + "when VM's specification is updated, thus preserving the VM's state between specification generations") command.Flags().StringVar(&username, "username", "admin", "SSH username to use when executing a startup script on the VM") command.Flags().StringVar(&password, "password", "admin", @@ -118,6 +122,7 @@ func runCreateVM(cmd *cobra.Command, args []string) error { NetSoftnet: netSoftnet, NetSoftnetAllow: netSoftnetAllow, NetSoftnetBlock: netSoftnetBlock, + Suspendable: suspendable, }, NetBridged: netBridged, Headless: headless, diff --git a/internal/controller/api_vms.go b/internal/controller/api_vms.go index bc46bce..cae3d2f 100644 --- a/internal/controller/api_vms.go +++ b/internal/controller/api_vms.go @@ -140,6 +140,16 @@ func (controller *Controller) updateVMSpec(ctx *gin.Context) responder.Responder userVM.NetSoftnet = true } + // Suspendable-specific sanity checks + if dbVM.Suspendable && !userVM.Suspendable { + return responder.JSON(http.StatusPreconditionFailed, NewErrorResponse("\"suspendable\" cannot be "+ + "toggled for suspendable VMs")) + } + if dbVM.Suspendable && dbVM.NetSoftnet != userVM.NetSoftnet { + return responder.JSON(http.StatusPreconditionFailed, NewErrorResponse("\"netSoftnet\" cannot be "+ + "toggled for suspendable VMs")) + } + if cmp.Equal(dbVM.VMSpec, userVM.VMSpec) { // Nothing was changed return responder.JSON(http.StatusOK, dbVM) diff --git a/internal/tests/spec_update_test.go b/internal/tests/spec_update_test.go index ebee72f..cd841a7 100644 --- a/internal/tests/spec_update_test.go +++ b/internal/tests/spec_update_test.go @@ -30,7 +30,6 @@ func TestSpecUpdateSoftnet(t *testing.T) { CPU: 4, Memory: 8 * 1024, Headless: true, - Status: v1.VMStatusPending, }) require.NoError(t, err) @@ -71,7 +70,7 @@ func TestSpecUpdateSoftnet(t *testing.T) { t.Logf("Waiting for the VM's observed generation to be updated...") return vm.ObservedGeneration == 1 - }), "failed to update a VM") + }), "failed to wait for the VM's observed generation to be updated") tartRunCmdline, err = tartRunProcessCmdline(tartVMName) require.NoError(t, err) @@ -80,6 +79,74 @@ func TestSpecUpdateSoftnet(t *testing.T) { require.True(t, sliceContainsAnotherSlice(tartRunCmdline, []string{"--net-softnet-block", "0.0.0.0/0"})) } +func TestSpecUpdateSoftnetSuspendable(t *testing.T) { + devClient, _, _ := devcontroller.StartIntegrationTestEnvironment(t) + + // Create a suspendable VM with Softnet enabled + vmName := "test" + + err := devClient.VMs().Create(t.Context(), &v1.VM{ + Meta: v1.Meta{ + Name: vmName, + }, + Image: imageconstant.DefaultMacosImage, + CPU: 4, + Memory: 8 * 1024, + Headless: true, + VMSpec: v1.VMSpec{ + Suspendable: true, + NetSoftnet: true, + }, + }) + require.NoError(t, err) + + // Wait for the VM to start + var vm *v1.VM + + require.True(t, wait.Wait(2*time.Minute, func() bool { + vm, err = devClient.VMs().Get(context.Background(), vmName) + require.NoError(t, err) + + t.Logf("Waiting for the VM to start. Current status: %s", vm.Status) + + return vm.Status == v1.VMStatusRunning + }), "failed to start a VM") + + // Ensure that the VM is using "--suspendable" and "--net-softnet" + tartVMName := ondiskname.New(vmName, vm.UID, vm.RestartCount).String() + + tartRunCmdline, err := tartRunProcessCmdline(tartVMName) + require.NoError(t, err) + require.Contains(t, tartRunCmdline, "--suspendable") + require.Contains(t, tartRunCmdline, "--net-softnet") + + // Update the VM's specification and tighten the Softnet restrictions + vm.NetSoftnetAllow = []string{"10.0.0.0/16"} + vm.NetSoftnetBlock = []string{"0.0.0.0/0"} + + vm, err = devClient.VMs().Update(t.Context(), *vm) + require.NoError(t, err) + require.EqualValues(t, 1, vm.Generation) + require.EqualValues(t, 0, vm.ObservedGeneration) + + require.True(t, wait.Wait(30*time.Second, func() bool { + vm, err = devClient.VMs().Get(context.Background(), vmName) + require.NoError(t, err) + + t.Logf("Waiting for the VM's observed generation to be updated...") + + return vm.ObservedGeneration == 1 + }), "failed to wait for the VM's observed generation to be updated") + + // Ensure that the VM is using "--suspendable", "--net-softnet" and "--net-softnet-{allow,block}" + tartRunCmdline, err = tartRunProcessCmdline(tartVMName) + require.NoError(t, err) + require.Contains(t, tartRunCmdline, "--suspendable") + require.Contains(t, tartRunCmdline, "--net-softnet") + require.True(t, sliceContainsAnotherSlice(tartRunCmdline, []string{"--net-softnet-allow", "10.0.0.0/16"})) + require.True(t, sliceContainsAnotherSlice(tartRunCmdline, []string{"--net-softnet-block", "0.0.0.0/0"})) +} + func tartRunProcessCmdline(vmName string) ([]string, error) { processes, err := process.Processes() if err != nil { diff --git a/internal/worker/vmmanager/condition.go b/internal/worker/vmmanager/condition.go new file mode 100644 index 0000000..7d52356 --- /dev/null +++ b/internal/worker/vmmanager/condition.go @@ -0,0 +1,10 @@ +package vmmanager + +type Condition string + +const ( + ConditionCloning Condition = "cloning" + ConditionReady Condition = "ready" + ConditionSuspending Condition = "suspending" + ConditionStopping Condition = "stopping" +) diff --git a/internal/worker/vmmanager/vm.go b/internal/worker/vmmanager/vm.go index 37ffc3d..d30a7d7 100644 --- a/internal/worker/vmmanager/vm.go +++ b/internal/worker/vmmanager/vm.go @@ -19,6 +19,7 @@ import ( "github.com/cirruslabs/orchard/internal/worker/tart" "github.com/cirruslabs/orchard/pkg/client" "github.com/cirruslabs/orchard/pkg/resource/v1" + mapset "github.com/deckarep/golang-set/v2" "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/metric" "go.uber.org/zap" @@ -32,16 +33,21 @@ type VM struct { Resource v1.VM logger *zap.SugaredLogger - // cloned state allows us to prevent a superfluous "tart delete" of a non-existent VM - cloned atomic.Bool - - // started state allows us to only set VM's "status" field in the API - // to "running" when we're indeed about to start the VM + // Backward compatibility with v1.VM specification's "Status" field + // + // "started" is always true after the first "tart run", + // whereas ConditionReady can be used to tell if a VM + // is really running or not. started atomic.Bool - // stopping state allows us to catch unexpected - // "tart run" terminations more correctly - stopping atomic.Bool + // A more orthogonal alternative to v1.VM specification's "Status" field, + // which allows a VM to have more than one state. + // + // For example, a VM can be both in ConditionReady and ConditionSuspending/ + // ConditionStopping states for a short time. This way in run() we know + // that we're in a process of rebooting a VM, so we can avoid throwing + // an error about unexpected VM termination. + conditions mapset.Set[Condition] // Image FQN feature, see https://github.com/cirruslabs/orchard/issues/164 imageFQN atomic.Pointer[string] @@ -75,6 +81,8 @@ func NewVM( "vm_restart_count", vmResource.RestartCount, ), + conditions: mapset.NewSet(ConditionCloning), + ctx: vmContext, cancel: vmContextCancel, @@ -122,9 +130,11 @@ func NewVM( return } - vm.cloned.Store(true) + // Backward compatibility with v1.VM specification's "Status" field vm.started.Store(true) + vm.conditions.Add(ConditionReady) + vm.run(vm.ctx, eventStreamer) }() @@ -135,10 +145,6 @@ func (vm *VM) OnDiskName() ondiskname.OnDiskName { return vm.onDiskName } -func (vm *VM) Started() bool { - return vm.started.Load() -} - func (vm *VM) ImageFQN() *string { return vm.imageFQN.Load() } @@ -152,7 +158,7 @@ func (vm *VM) Status() v1.VMStatus { return v1.VMStatusFailed } - if vm.Started() { + if vm.started.Load() { return v1.VMStatusRunning } @@ -188,6 +194,10 @@ func (vm *VM) setErr(err error) { } } +func (vm *VM) Conditions() mapset.Set[Condition] { + return vm.conditions.Clone() +} + func (vm *VM) cloneAndConfigure(ctx context.Context) error { vm.setStatusMessage("cloning VM...") @@ -196,6 +206,8 @@ func (vm *VM) cloneAndConfigure(ctx context.Context) error { return err } + vm.conditions.Remove(ConditionCloning) + // Image FQN feature, see https://github.com/cirruslabs/orchard/issues/164 fqnRaw, _, err := tart.Tart(ctx, vm.logger, "fqn", vm.Resource.Image) if err == nil { @@ -317,6 +329,8 @@ func (vm *VM) cloneAndConfigure(ctx context.Context) error { } func (vm *VM) run(ctx context.Context, eventStreamer *client.EventStreamer) { + defer vm.conditions.RemoveAll(ConditionReady, ConditionSuspending, ConditionStopping) + // Launch the startup script goroutine as close as possible // to the VM startup (below) to avoid "tart ip" timing out if vm.Resource.StartupScript != nil { @@ -350,6 +364,10 @@ func (vm *VM) run(ctx context.Context, eventStreamer *client.EventStreamer) { runArgs = append(runArgs, "--nested") } + if vm.Resource.Suspendable { + runArgs = append(runArgs, "--suspendable") + } + for _, hostDir := range vm.Resource.HostDirs { runArgs = append(runArgs, fmt.Sprintf("--dir=%s", hostDir.String())) } @@ -371,7 +389,7 @@ func (vm *VM) run(ctx context.Context, eventStreamer *client.EventStreamer) { case <-vm.ctx.Done(): // Do not return an error because it's the user's intent to cancel this VM default: - if !vm.stopping.Load() { + if !vm.conditions.ContainsAny(ConditionSuspending, ConditionStopping) { vm.setErr(fmt.Errorf("%w: VM exited unexpectedly", ErrVMFailed)) } } @@ -404,24 +422,77 @@ func (vm *VM) IP(ctx context.Context) (string, error) { return strings.TrimSpace(stdout), nil } -func (vm *VM) Stop() { - vm.logger.Debugf("stopping VM") +func (vm *VM) Suspend() <-chan error { + errCh := make(chan error, 1) - vm.stopping.Store(true) - defer vm.stopping.Store(false) + select { + case <-vm.ctx.Done(): + // VM is already suspended/stopped + errCh <- nil - // Try to gracefully terminate the VM - _, _, _ = tart.Tart(context.Background(), zap.NewNop().Sugar(), "stop", "--timeout", "5", vm.id()) + return errCh + default: + // VM is still running + } - // Terminate the VM goroutine ("tart pull", "tart clone", "tart run", etc.) via the context - vm.cancel() - vm.wg.Wait() + vm.setStatusMessage("Suspending VM") + vm.conditions.Add(ConditionSuspending) - vm.logger.Debugf("VM stopped") + go func() { + _, _, err := tart.Tart(context.Background(), zap.NewNop().Sugar(), "suspend", vm.id()) + if err != nil { + err := fmt.Errorf("failed to suspend VM: %w", err) + vm.setErr(err) + errCh <- err + + return + } + + errCh <- nil + }() + + return errCh } -func (vm *VM) Reboot(eventStreamer *client.EventStreamer) { - vm.Stop() +func (vm *VM) Stop() <-chan error { + errCh := make(chan error, 1) + + select { + case <-vm.ctx.Done(): + // VM is already suspended/stopped + errCh <- nil + + return errCh + default: + // VM is still running + } + + vm.setStatusMessage("Stopping VM") + vm.conditions.Add(ConditionStopping) + + go func() { + // Try to gracefully terminate the VM + _, _, _ = tart.Tart(context.Background(), zap.NewNop().Sugar(), "stop", "--timeout", "5", vm.id()) + + // Terminate the VM goroutine ("tart pull", "tart clone", "tart run", etc.) via the context + vm.cancel() + vm.wg.Wait() + + // We don't return an error because we always terminate a VM + errCh <- nil + }() + + return errCh +} + +func (vm *VM) Start(vmResource v1.VM, eventStreamer *client.EventStreamer) { + vm.Resource = vmResource + vm.Resource.ObservedGeneration = vmResource.Generation + + vm.setStatusMessage("Starting VM") + vm.conditions.Add(ConditionReady) + + vm.cancel() vm.ctx, vm.cancel = context.WithCancel(context.Background()) vm.wg.Add(1) @@ -434,19 +505,20 @@ func (vm *VM) Reboot(eventStreamer *client.EventStreamer) { } func (vm *VM) Delete() error { - if !vm.cloned.Load() { + // Cancel all currently running Tart invocations + // (e.g. "tart clone", "tart run", etc.) + vm.cancel() + + if vm.conditions.Contains(ConditionCloning) { + // Not cloned yet, nothing to delete return nil } - vm.logger.Debugf("deleting VM") - _, _, err := tart.Tart(context.Background(), vm.logger, "delete", vm.id()) if err != nil { return fmt.Errorf("%w: failed to delete VM: %v", ErrVMFailed, err) } - vm.logger.Debugf("deleted VM") - return nil } diff --git a/internal/worker/worker.go b/internal/worker/worker.go index 138f08e..630916c 100644 --- a/internal/worker/worker.go +++ b/internal/worker/worker.go @@ -6,6 +6,7 @@ import ( "fmt" "os" "slices" + "strings" "time" "github.com/avast/retry-go/v4" @@ -138,7 +139,7 @@ func (worker *Worker) Run(ctx context.Context) error { func (worker *Worker) Close() error { var result error for _, vm := range worker.vmm.List() { - vm.Stop() + <-vm.Stop() } for _, vm := range worker.vmm.List() { err := vm.Delete() @@ -320,14 +321,20 @@ func (worker *Worker) syncVMs(ctx context.Context, updateVM func(context.Context } localState := mo.None[v1.VMStatus]() + var localConditions []string if vm != nil { localState = mo.Some(vm.Status()) + + for condition := range vm.Conditions().Iter() { + localConditions = append(localConditions, string(condition)) + } } action := transitions[remoteState][localState] - worker.logger.Debugf("processing VM: %s, remote: %v, local: %v, action: %v\n", onDiskName, - optionToString(remoteState), optionToString(localState), action) + worker.logger.Debugf("processing VM: %s, remote state: %s, local state: %s, "+ + "local conditions: [%s], action: %v\n", onDiskName, optionToString(remoteState), + optionToString(localState), strings.Join(localConditions, ", "), action) switch action { case ActionCreate: @@ -359,24 +366,37 @@ func (worker *Worker) syncVMs(ctx context.Context, updateVM func(context.Context return err } case ActionMonitorRunning: - if vmResource.StatusMessage != vm.StatusMessage() { - vmResource.StatusMessage = vm.StatusMessage() - - if _, err := updateVM(ctx, *vmResource); err != nil { - return err + if vmResource.Generation != vm.Resource.Generation { + // VM specification changed, reboot the VM for the changes to take effect + if vm.Conditions().Contains(vmmanager.ConditionReady) { + // VM is running, suspend or stop it first + if vm.Resource.Suspendable { + vm.Suspend() + } else { + vm.Stop() + } + } else { + // VM stopped, start it with the new specification + eventStreamer := worker.client.VMs().StreamEvents(vmResource.Name) + vm.Start(*vmResource, eventStreamer) } } - if vmResource.Generation != vm.Resource.Generation { - // Something changed, reboot the VM for the changes to take effect - vm.Resource = *vmResource + var updateNeeded bool - eventStreamer := worker.client.VMs().StreamEvents(vmResource.Name) + if vmResource.StatusMessage != vm.StatusMessage() { + vmResource.StatusMessage = vm.StatusMessage() - vm.Reboot(eventStreamer) + updateNeeded = true + } - vmResource.ObservedGeneration = vm.Resource.Generation + if vmResource.ObservedGeneration != vm.Resource.ObservedGeneration { + vmResource.ObservedGeneration = vm.Resource.ObservedGeneration + updateNeeded = true + } + + if updateNeeded { if _, err := updateVM(ctx, *vmResource); err != nil { return err } @@ -487,7 +507,7 @@ func (worker *Worker) syncOnDiskVMs(ctx context.Context) error { } func (worker *Worker) deleteVM(vm *vmmanager.VM) error { - vm.Stop() + <-vm.Stop() if err := vm.Delete(); err != nil { return err diff --git a/pkg/resource/v1/v1.go b/pkg/resource/v1/v1.go index 38ae7ac..fed4415 100644 --- a/pkg/resource/v1/v1.go +++ b/pkg/resource/v1/v1.go @@ -102,6 +102,7 @@ type VMSpec struct { NetSoftnet bool `json:"netSoftnet,omitempty"` NetSoftnetAllow []string `json:"netSoftnetAllow,omitempty"` NetSoftnetBlock []string `json:"netSoftnetBlock,omitempty"` + Suspendable bool `json:"suspendable,omitempty"` } type VMState struct {