From d0291f30205420dff44f5b708608ec068316a35a Mon Sep 17 00:00:00 2001 From: dawn Date: Mon, 27 Jul 2026 12:33:47 +0300 Subject: [PATCH] spindle/engines/microvm: treat exited vms as already destroyed instead of erroring out Signed-off-by: dawn --- spindle/engines/microvm/qemu.go | 8 +++- spindle/engines/microvm/vm.go | 29 +++++++++---- spindle/engines/microvm/vm_test.go | 65 ++++++++++++++++++++++++++++++ 3 files changed, 94 insertions(+), 8 deletions(-) create mode 100644 spindle/engines/microvm/vm_test.go diff --git a/spindle/engines/microvm/qemu.go b/spindle/engines/microvm/qemu.go index fec2cd2b..a0080bf1 100644 --- a/spindle/engines/microvm/qemu.go +++ b/spindle/engines/microvm/qemu.go @@ -306,7 +306,13 @@ func (h *QEMUVMHandle) Shutdown(ctx context.Context) error { } if h.QMPMon != nil { if err := h.QMPSystemPowerdown(); err != nil { - return err + // dead qmp socket means qemu exited concurrently so we wait + select { + case <-h.done: + return nil + case <-ctx.Done(): + return err + } } } if h.done == nil { diff --git a/spindle/engines/microvm/vm.go b/spindle/engines/microvm/vm.go index 1fc9e031..0ae4db49 100644 --- a/spindle/engines/microvm/vm.go +++ b/spindle/engines/microvm/vm.go @@ -205,29 +205,44 @@ func (e *Engine) shutdownVM(ctx context.Context, wid models.WorkflowId, state *w if state.VM == nil { return nil } + if vmExited(state.VM) { + return closeIO(&state.VM) + } - var err error + var poweroffErr error if state.Agent != nil { gracefulCtx, cancel := context.WithTimeout(ctx, vmShutdownTimeout) - poweredOff, poweroffErr := e.poweroffViaAgent(gracefulCtx, wid, state) + var poweredOff bool + poweredOff, poweroffErr = e.poweroffViaAgent(gracefulCtx, wid, state) cancel() - err = errors.Join(err, poweroffErr) if poweredOff { - return errors.Join(err, closeIO(&state.VM)) + return closeIO(&state.VM) + } + if vmExited(state.VM) { + return closeIO(&state.VM) } } fallbackCtx, cancel := context.WithTimeout(ctx, vmShutdownTimeout) defer cancel() - if shutdownErr := state.VM.Shutdown(fallbackCtx); shutdownErr != nil { + shutdownErr := state.VM.Shutdown(fallbackCtx) + if shutdownErr != nil && !vmExited(state.VM) { e.l.Warn("microVM shutdown fallback failed", "workflow", wid, "error", shutdownErr) - err = errors.Join(err, shutdownErr) + return errors.Join(poweroffErr, shutdownErr, closeIO(&state.VM)) } - return errors.Join(err, closeIO(&state.VM)) + return closeIO(&state.VM) +} + +func vmExited(vm VMHandle) bool { + ctx, cancel := context.WithCancel(context.Background()) + cancel() + // a cancelled wait means the process is still live + // any other result means it exited + return !errors.Is(vm.WaitContext(ctx), context.Canceled) } func (e *Engine) poweroffViaAgent(ctx context.Context, wid models.WorkflowId, state *workflowState) (bool, error) { diff --git a/spindle/engines/microvm/vm_test.go b/spindle/engines/microvm/vm_test.go new file mode 100644 index 00000000..65812e2c --- /dev/null +++ b/spindle/engines/microvm/vm_test.go @@ -0,0 +1,65 @@ +package microvm + +import ( + "context" + "errors" + "log/slog" + "testing" + + "tangled.org/core/spindle/models" +) + +type shutdownTestVM struct { + exited bool + waitErr error + shutdownErr error + closed bool +} + +func (v *shutdownTestVM) Shutdown(context.Context) error { + v.exited = true + return v.shutdownErr +} + +func (v *shutdownTestVM) WaitContext(ctx context.Context) error { + if v.exited { + return v.waitErr + } + return ctx.Err() +} + +func (v *shutdownTestVM) Close() error { + v.closed = true + return nil +} + +func (*shutdownTestVM) Logs() VMLogs { return VMLogs{} } +func (*shutdownTestVM) CID() uint32 { return 0 } +func (*shutdownTestVM) WorkDir() string { return "" } +func (*shutdownTestVM) OOMKilled() bool { return false } + +func TestShutdownVM_TreatsExitedVMAsCleanedUp(t *testing.T) { + for name, vm := range map[string]*shutdownTestVM{ + "already exited": { + exited: true, + waitErr: errors.New("qemu exited"), + shutdownErr: errors.New("qmp broken pipe"), + }, + "exits during fallback": { + waitErr: errors.New("qemu exited"), + shutdownErr: errors.New("qmp broken pipe"), + }, + } { + t.Run(name, func(t *testing.T) { + e := &Engine{l: slog.Default()} + state := &workflowState{VM: vm} + + if err := e.shutdownVM(context.Background(), models.WorkflowId{}, state); err != nil { + t.Fatalf("shutdownVM: %v", err) + } + if !vm.closed { + t.Fatal("expected vm handle to be closed") + } + }) + } +} -- 2.51.2