diff --git a/spindle/engines/qemu/bakers/baker.go b/spindle/engines/qemu/bakers/baker.go index 293b4fbf..fd3ad1e3 100644 --- a/spindle/engines/qemu/bakers/baker.go +++ b/spindle/engines/qemu/bakers/baker.go @@ -11,6 +11,11 @@ import ( "tangled.org/core/log" ) +type ImageMetadata struct { + Cmdline string `json:"cmdline"` + Shell string `json:"shell"` +} + const ( DiskName = "disk" KernelName = "kernel" @@ -27,11 +32,6 @@ func ConfigPath(dir string) string { return filepath.Join(dir, ConfigName) } func UserDataPath(dir string) string { return filepath.Join(dir, UserDataName) } func SeedISOPath(dir string) string { return filepath.Join(dir, SeedISOName) } -type ImageMetadata struct { - Cmdline string `json:"cmdline"` - Shell string `json:"shell"` -} - type ImageBaker interface { Prepare(ctx context.Context, imageDir string) error } diff --git a/spindle/engines/qemu/engine.go b/spindle/engines/qemu/engine.go index be510ae1..79912a39 100644 --- a/spindle/engines/qemu/engine.go +++ b/spindle/engines/qemu/engine.go @@ -1,20 +1,15 @@ package qemu import ( - "bufio" "context" - "encoding/base64" "encoding/json" "fmt" "log/slog" - "net" "os" - "os/exec" "path/filepath" "sync" "time" - "github.com/digitalocean/go-qemu/qmp" "gopkg.in/yaml.v3" "tangled.org/core/api/tangled" @@ -22,6 +17,7 @@ import ( "tangled.org/core/spindle/config" "tangled.org/core/spindle/engine" "tangled.org/core/spindle/engines/qemu/bakers" + "tangled.org/core/spindle/engines/qemu/virt" "tangled.org/core/spindle/models" "tangled.org/core/spindle/secrets" ) @@ -81,13 +77,13 @@ func (e *Engine) InitWorkflow(twf tangled.Pipeline_Workflow, tpl tangled.Pipelin return nil, err } - swf.Data = vmState{img: img} + swf.Data = &virt.VMState{Img: img} return swf, nil } // discover and resolve kernel, initrd, and disk and other config from an image subfolder -func (e *Engine) resolveImage(name string) (ResolvedImage, error) { - var img ResolvedImage +func (e *Engine) resolveImage(name string) (virt.ResolvedImage, error) { + var img virt.ResolvedImage if name == "" { name = e.cfg.QemuPipelines.DefaultImage } @@ -102,36 +98,36 @@ func (e *Engine) resolveImage(name string) (ResolvedImage, error) { configPath := bakers.ConfigPath(imageDir) if _, err := os.Stat(diskPath); err == nil { - img.disk = diskPath + img.Disk = diskPath } if _, err := os.Stat(kernelPath); err == nil { - img.kernel = kernelPath + img.Kernel = kernelPath } if _, err := os.Stat(initrdPath); err == nil { - img.initrd = initrdPath + img.Initrd = initrdPath } if b, err := os.ReadFile(configPath); err == nil { var meta bakers.ImageMetadata if err := json.Unmarshal(b, &meta); err == nil { if meta.Cmdline != "" { - img.cmdline = meta.Cmdline + img.Cmdline = meta.Cmdline } if meta.Shell != "" { - img.shell = meta.Shell + img.Shell = meta.Shell } } } if b, err := os.ReadFile(filepath.Join(imageDir, bakers.UserDataName)); err == nil { - img.userData = string(b) + img.UserData = string(b) } - if img.disk == "" { + if img.Disk == "" { return img, fmt.Errorf("missing '%s' in %s", bakers.DiskName, imageDir) } - if img.kernel != "" && (img.initrd == "" || img.cmdline == "") { + if img.Kernel != "" && (img.Initrd == "" || img.Cmdline == "") { return img, fmt.Errorf("kernel requires initrd and cmdline, but 'initrd' and/or 'cmdline' is missing for %s", name) } - if img.shell == "" { + if img.Shell == "" { return img, fmt.Errorf("shell is not configured for %s", name) } return img, nil @@ -147,330 +143,50 @@ func (e *Engine) SetupWorkflow(ctx context.Context, wid models.WorkflowId, wf *m wfLogger.ControlWriter(setupStepIdx, setupStep, models.StepStatusStart).Write([]byte{0}) defer wfLogger.ControlWriter(setupStepIdx, setupStep, models.StepStatusEnd).Write([]byte{0}) - // some systems have tmpfs at /tmp which is not ideal for large workloads - targetTempDir := e.cfg.QemuPipelines.OverlayDir - if targetTempDir == "" { - targetTempDir = os.TempDir() - } - - tempDir, err := os.MkdirTemp(targetTempDir, "qemu-wf-*") - if err != nil { - return err - } - e.registerCleanup(wid, func(ctx context.Context) error { - return os.RemoveAll(tempDir) - }) - - state := wf.Data.(vmState) - img := state.img - - if err := generateSeedISO(tempDir, img.userData); err != nil { - return fmt.Errorf("generating seed iso: %w", err) - } - - qmpSock := filepath.Join(tempDir, "qmp.sock") - qgaSock := filepath.Join(tempDir, "qga.sock") - - // todo(dawn): ideally would be nice if we used qemu with the microvm enabled here... - // but that is not compatible with cloud-init since it expects real hw enumeration... - // and we would not be able to use standard cloud images, which is kind of annoying. - // we also have to manage a virtiofsd process for the filesystem, instead of having to - // manage a qcow overlay (see https://ubuntu.com/server/docs/explanation/virtualisation/qemu-microvm/). - // also also need to be able to have some scripts for generating our own images - // (though we should already do this anyway since some like alpine don't provide - // "cloud ready" image files, at least with kernel and initrd) - argv := []string{ - // todo(dawn): ideally probably have "tiers" and let the spindle concile the tier using - // what the user wants and what the user has exposed to them by the spindle operator? - "-m", e.cfg.QemuPipelines.Memory, "-smp", fmt.Sprintf("%d", e.cfg.QemuPipelines.SMP), - "-display", "none", "-monitor", "none", "-nodefaults", "-no-user-config", - // use snapshot=on to do copy-on-write without having us manage qcow overlays manually - "-drive", fmt.Sprintf("file=%s,media=disk,snapshot=on,if=virtio", img.disk), - "-drive", fmt.Sprintf("file=%s,media=cdrom", bakers.SeedISOPath(tempDir)), - "-netdev", "user,id=net0", - "-device", "virtio-net-pci,netdev=net0", - "-qmp", fmt.Sprintf("unix:%s,server,nowait", qmpSock), - "-chardev", fmt.Sprintf("socket,path=%s,server,nowait,id=qga0", qgaSock), - "-device", "virtio-serial", - "-device", "virtserialport,chardev=qga0,name=org.qemu.guest_agent.0", - } - - if e.cfg.Server.Dev { - argv = append(argv, "-serial", "stdio") - } else { - argv = append(argv, "-serial", "none") - } - - // support booting using qemu bios still, but otherwise we make it faster! - // incase someone wants to do this for whatever reason... - if img.kernel != "" { - argv = append(argv, "-kernel", img.kernel) - if img.initrd != "" { - argv = append(argv, "-initrd", img.initrd) - } - if img.cmdline != "" { - argv = append(argv, "-append", img.cmdline) - } - } else { - // booting with bios: explicitly select disk - argv = append(argv, "-boot", "order=c") - } - - enableKVM := e.cfg.QemuPipelines.EnableKVM - if _, err := os.Stat("/dev/kvm"); err != nil { - if enableKVM { - l.Warn("kvm was requested but /dev/kvm is not accessible; falling back to software emulation", "error", err) - } - enableKVM = false - } - if enableKVM { - argv = append(argv, "-enable-kvm", "-cpu", "host") - } - - // todo(dawn): same with above, we assume x86_64 here, but should allow other archs, - // probably just auto detect as a default - qemuCmd := exec.Command("qemu-system-x86_64", argv...) - qemuCmd.Env = append(os.Environ(), "TMPDIR="+tempDir) - qemuCmd.Stdout = os.Stdout - qemuCmd.Stderr = os.Stderr - - startedBootAt := time.Now() - if err := qemuCmd.Start(); err != nil { - return fmt.Errorf("starting qemu: %w", err) - } - - var mon *qmp.SocketMonitor + state := wf.Data.(*virt.VMState) + img := state.Img - // cleanup qemu if we fail to setup at any point below - setupOk := false - defer func() { - if !setupOk { - _ = qemuCmd.Process.Kill() - if mon != nil { - _ = mon.Disconnect() - } - } - }() - - qmpCtx, cancelQmp := context.WithTimeout(ctx, 10*time.Second) - defer cancelQmp() - - // wait for qmp to be ready - for { - mon, err = qmp.NewSocketMonitor("unix", qmpSock, 2*time.Second) - if err == nil { - if err = mon.Connect(); err == nil { - break - } - } - select { - case <-qmpCtx.Done(): - return fmt.Errorf("qmp connect timeout: %w", err) - case <-time.After(10 * time.Millisecond): - } - } - - status, err := e.qmpQueryStatus(mon) + actualState, err := virt.StartVM(ctx, virt.VMConfig{ + Memory: e.cfg.QemuPipelines.Memory, + SMP: e.cfg.QemuPipelines.SMP, + EnableKVM: e.cfg.QemuPipelines.EnableKVM, + Dev: e.cfg.Server.Dev, + OverlayDir: e.cfg.QemuPipelines.OverlayDir, + }, img, l) if err != nil { return err } - l.Info("qemu guest status", "status", status) - if status != "running" { - return fmt.Errorf("qemu guest not running (status: %s)", status) - } - - qgaCtx, cancelQga := context.WithTimeout(ctx, time.Minute) - defer cancelQga() - - // wait for guest agent to be ready (aka boot) - for { - if _, err := os.Stat(qgaSock); err == nil { - pingCtx, cancelPing := context.WithTimeout(qgaCtx, 5*time.Second) - err = e.qgaGuestPing(pingCtx, qgaSock) - cancelPing() - if err == nil { - l.Info("vm booted and guest-agent ready", "elapsed", time.Since(startedBootAt).Round(time.Millisecond)) - break - } - l.Debug("qga guest-ping failed", "error", err) - } else { - l.Debug("qga socket not found yet", "path", qgaSock) - } - - select { - case <-qgaCtx.Done(): - return fmt.Errorf("qga connect timeout: %w", err) - case <-time.After(100 * time.Millisecond): - } - } - e.registerCleanup(wid, func(ctx context.Context) error { // graceful powerdown so guest can sync filesystem - if err := e.qmpSystemPowerdown(mon); err != nil { + if err := actualState.QMPSystemPowerdown(); err != nil { l.Error("failed to powerdown qemu guest", "workflow", wid, "error", err) } done := make(chan error, 1) - go func() { done <- qemuCmd.Wait() }() + go func() { + _, err := actualState.Process.Wait() + done <- err + }() select { case <-done: case <-time.After(5 * time.Second): - _ = qemuCmd.Process.Kill() + _ = actualState.Process.Kill() <-done // drain to avoid zombie } - _ = mon.Disconnect() + _ = actualState.QMPMon.Disconnect() + _ = os.RemoveAll(actualState.TempDir) return nil }) - wf.Data = vmState{ - process: qemuCmd.Process, - qmpMon: mon, - qgaPath: qgaSock, - tempDir: tempDir, - img: img, - } - - setupOk = true + wf.Data = actualState return nil } -func (e *Engine) qmpRun(mon *qmp.SocketMonitor, command qmp.Command) ([]byte, error) { - b, err := json.Marshal(command) - if err != nil { - return nil, err - } - return mon.Run(b) -} - -// sends a command to the qemu guest agent and returns the response -func (e *Engine) qgaRun(ctx context.Context, sock string, command qmp.Command) ([]byte, error) { - b, err := json.Marshal(command) - if err != nil { - return nil, err - } - - conn, err := (&net.Dialer{}).DialContext(ctx, "unix", sock) - if err != nil { - return nil, err - } - defer conn.Close() - - if dl, ok := ctx.Deadline(); ok { - _ = conn.SetDeadline(dl) - } - - if _, err := conn.Write(append(b, '\n')); err != nil { - return nil, err - } - - return bufio.NewReader(conn).ReadBytes('\n') -} - -// query the qemu guest status -func (e *Engine) qmpQueryStatus(mon *qmp.SocketMonitor) (string, error) { - raw, err := e.qmpRun(mon, qmp.Command{Execute: "query-status"}) - if err != nil { - return "", fmt.Errorf("qmp query-status failed: %w", err) - } - - var resp struct { - Return struct { - Status string `json:"status"` - } `json:"return"` - } - if err := json.Unmarshal(raw, &resp); err != nil { - return "", fmt.Errorf("qmp query-status parse: %w", err) - } - return resp.Return.Status, nil -} - -// send a command to qemu to powerdown the guest gracefully -func (e *Engine) qmpSystemPowerdown(mon *qmp.SocketMonitor) error { - _, err := e.qmpRun(mon, qmp.Command{Execute: "system_powerdown"}) - return err -} - -// ping the guest agent to see if it's ready -func (e *Engine) qgaGuestPing(ctx context.Context, sock string) error { - _, err := e.qgaRun(ctx, sock, qmp.Command{Execute: "guest-ping"}) - return err -} - -// execute a command on the guest and return its pid -func (e *Engine) qgaGuestExec(ctx context.Context, sock string, shell string, command string, env []string) (int, error) { - cmdArgs := []string{shell, "-c", command} - raw, err := e.qgaRun(ctx, sock, qmp.Command{ - Execute: "guest-exec", - Args: map[string]any{ - "path": shell, - "arg": cmdArgs[1:], - "env": env, - "capture-output": true, - }, - }) - if err != nil { - return 0, fmt.Errorf("qga guest-exec: %w", err) - } - - var resp struct { - Return struct { - Pid int `json:"pid"` - } `json:"return"` - } - if err := json.Unmarshal(raw, &resp); err != nil { - return 0, fmt.Errorf("qga guest-exec parse: %w", err) - } - return resp.Return.Pid, nil -} - -type guestExecStatus struct { - Exited bool - ExitCode int - // these are since the last time we asked for status - OutData []byte - ErrData []byte -} - -// get the status of a command executed with guest-exec -func (e *Engine) qgaGuestExecStatus(ctx context.Context, sock string, pid int) (guestExecStatus, error) { - var status guestExecStatus - raw, err := e.qgaRun(ctx, sock, qmp.Command{ - Execute: "guest-exec-status", - Args: map[string]any{"pid": pid}, - }) - if err != nil { - return status, fmt.Errorf("qga guest-exec-status: %w", err) - } - - var resp struct { - Return struct { - Exited bool `json:"exited"` - ExitCode int `json:"exitcode"` - OutData string `json:"out-data"` - ErrData string `json:"err-data"` - } `json:"return"` - } - if err := json.Unmarshal(raw, &resp); err != nil { - return status, fmt.Errorf("qga guest-exec-status parse: %w", err) - } - - status.Exited = resp.Return.Exited - status.ExitCode = resp.Return.ExitCode - if resp.Return.OutData != "" { - status.OutData, _ = base64.StdEncoding.DecodeString(resp.Return.OutData) - } - if resp.Return.ErrData != "" { - status.ErrData, _ = base64.StdEncoding.DecodeString(resp.Return.ErrData) - } - - return status, nil -} - func (e *Engine) RunStep(ctx context.Context, wid models.WorkflowId, w *models.Workflow, idx int, secrets []secrets.UnlockedSecret, wfLogger models.WorkflowLogger) error { - state := w.Data.(vmState) + state := w.Data.(*virt.VMState) step := w.Steps[idx] env := make([]string, 0, len(w.Environment)+len(secrets)) @@ -486,13 +202,13 @@ func (e *Engine) RunStep(ctx context.Context, wid models.WorkflowId, w *models.W } } - pid, err := e.qgaGuestExec(ctx, state.qgaPath, state.img.shell, step.Command(), env) + pid, err := state.QGAGuestExec(ctx, step.Command(), env) if err != nil { return err } for { - status, err := e.qgaGuestExecStatus(ctx, state.qgaPath, pid) + status, err := state.QGAGuestExecStatus(ctx, pid) if err != nil { return err } diff --git a/spindle/engines/qemu/errors.go b/spindle/engines/qemu/errors.go deleted file mode 100644 index fadc2645..00000000 --- a/spindle/engines/qemu/errors.go +++ /dev/null @@ -1,8 +0,0 @@ -package qemu - -import "errors" - -var ( - ErrOOMKilled = errors.New("container died due to OOM kill") - ErrBootTimeout = errors.New("timed out waiting for VM to boot") -) diff --git a/spindle/engines/qemu/models.go b/spindle/engines/qemu/models.go index 2052c9d2..80abb168 100644 --- a/spindle/engines/qemu/models.go +++ b/spindle/engines/qemu/models.go @@ -1,29 +1,5 @@ package qemu -import ( - "os" - - "github.com/digitalocean/go-qemu/qmp" -) - -type vmState struct { - process *os.Process - qmpMon *qmp.SocketMonitor - qgaPath string - tempDir string - - img ResolvedImage -} - -type ResolvedImage struct { - kernel string - initrd string - disk string - cmdline string // kernel command line - shell string // shell to use for workflow steps - userData string // extra cloud-init user-data -} - type manifestWorkflow struct { Image string `yaml:"image"` Steps []struct { diff --git a/spindle/engines/qemu/cloudinit.go b/spindle/engines/qemu/virt/cloudinit.go similarity index 78% rename from spindle/engines/qemu/cloudinit.go rename to spindle/engines/qemu/virt/cloudinit.go index 3af708d7..ae067ada 100644 --- a/spindle/engines/qemu/cloudinit.go +++ b/spindle/engines/qemu/virt/cloudinit.go @@ -1,14 +1,18 @@ -package qemu +package virt import ( "fmt" "os" "os/exec" "path/filepath" - - "tangled.org/core/spindle/engines/qemu/bakers" ) +const SeedISOName = "seed.iso" + +func SeedISOPath(dir string) string { + return filepath.Join(dir, SeedISOName) +} + func generateSeedISO(dir string, extraUserData string) error { metaData := "instance-id: spindle-vm\nlocal-hostname: spindle\n" userData := `#cloud-config @@ -32,7 +36,7 @@ users: } // nocloud source expects volid "cidata" with joliet and rock ridge extensions - cmd := exec.Command("genisoimage", "-output", bakers.SeedISOPath(dir), "-volid", "cidata", "-joliet", "-rock", "meta-data", "user-data") + cmd := exec.Command("genisoimage", "-output", SeedISOPath(dir), "-volid", "cidata", "-joliet", "-rock", "meta-data", "user-data") cmd.Dir = dir cmd.Stdout = os.Stdout cmd.Stderr = os.Stderr diff --git a/spindle/engines/qemu/virt/models.go b/spindle/engines/qemu/virt/models.go new file mode 100644 index 00000000..523122bf --- /dev/null +++ b/spindle/engines/qemu/virt/models.go @@ -0,0 +1,33 @@ +package virt + +import ( + "os" + + "github.com/digitalocean/go-qemu/qmp" +) + +type VMConfig struct { + Memory string + SMP int + EnableKVM bool + Dev bool + OverlayDir string +} + +type ResolvedImage struct { + Kernel string + Initrd string + Disk string + Cmdline string // kernel command line + Shell string // shell to use for workflow steps + UserData string // extra cloud-init user-data +} + +type VMState struct { + Process *os.Process + QMPMon *qmp.SocketMonitor + QGAPath string + TempDir string + + Img ResolvedImage +} diff --git a/spindle/engines/qemu/virt/vm.go b/spindle/engines/qemu/virt/vm.go new file mode 100644 index 00000000..6c13a4f0 --- /dev/null +++ b/spindle/engines/qemu/virt/vm.go @@ -0,0 +1,320 @@ +package virt + +import ( + "bufio" + "context" + "encoding/base64" + "encoding/json" + "fmt" + "log/slog" + "net" + "os" + "os/exec" + "path/filepath" + "time" + + "github.com/digitalocean/go-qemu/qmp" +) + +func StartVM(ctx context.Context, cfg VMConfig, img ResolvedImage, l *slog.Logger) (*VMState, error) { + // some systems have tmpfs at /tmp which is not ideal for large workloads + targetTempDir := cfg.OverlayDir + if targetTempDir == "" { + targetTempDir = os.TempDir() + } + + tempDir, err := os.MkdirTemp(targetTempDir, "qemu-vm-*") + if err != nil { + return nil, err + } + + state := &VMState{ + TempDir: tempDir, + Img: img, + } + + cleanup := func() { + _ = state.Close() + _ = os.RemoveAll(tempDir) + } + + if err := generateSeedISO(tempDir, img.UserData); err != nil { + cleanup() + return nil, fmt.Errorf("generating seed iso: %w", err) + } + + qmpSock := filepath.Join(tempDir, "qmp.sock") + qgaSock := filepath.Join(tempDir, "qga.sock") + + state.QGAPath = qgaSock + + // todo(dawn): ideally would be nice if we used qemu with the microvm enabled here... + // but that is not compatible with cloud-init since it expects real hw enumeration... + // and we would not be able to use standard cloud images, which is kind of annoying. + // we also have to manage a virtiofsd process for the filesystem, instead of having to + // manage a qcow overlay (see https://ubuntu.com/server/docs/explanation/virtualisation/qemu-microvm/). + // also also need to be able to have some scripts for generating our own images + // (though we should already do this anyway since some like alpine don't provide + // "cloud ready" image files, at least with kernel and initrd) + argv := []string{ + // todo(dawn): ideally probably have "tiers" and let the spindle concile the tier using + // what the user wants and what the user has exposed to them by the spindle operator? + "-m", cfg.Memory, "-smp", fmt.Sprintf("%d", cfg.SMP), + "-display", "none", "-monitor", "none", "-nodefaults", "-no-user-config", + // use snapshot=on to do copy-on-write without having us manage qcow overlays manually + "-drive", fmt.Sprintf("file=%s,media=disk,snapshot=on,if=virtio", img.Disk), + "-drive", fmt.Sprintf("file=%s,media=cdrom", SeedISOPath(tempDir)), + "-netdev", "user,id=net0", + "-device", "virtio-net-pci,netdev=net0", + "-qmp", fmt.Sprintf("unix:%s,server,nowait", qmpSock), + "-chardev", fmt.Sprintf("socket,path=%s,server,nowait,id=qga0", qgaSock), + "-device", "virtio-serial", + "-device", "virtserialport,chardev=qga0,name=org.qemu.guest_agent.0", + } + + if cfg.Dev { + argv = append(argv, "-serial", "stdio") + } else { + argv = append(argv, "-serial", "none") + } + + // support booting using qemu bios still, but otherwise we make it faster! + // incase someone wants to do this for whatever reason... + if img.Kernel != "" { + argv = append(argv, "-kernel", img.Kernel) + if img.Initrd != "" { + argv = append(argv, "-initrd", img.Initrd) + } + if img.Cmdline != "" { + argv = append(argv, "-append", img.Cmdline) + } + } else { + argv = append(argv, "-boot", "order=c") + } + + enableKVM := cfg.EnableKVM + if _, err := os.Stat("/dev/kvm"); err != nil { + if enableKVM { + l.Warn("kvm was requested but /dev/kvm is not accessible; falling back to software emulation", "error", err) + } + enableKVM = false + } + if enableKVM { + argv = append(argv, "-enable-kvm", "-cpu", "host") + } + + // todo(dawn): same with above, we assume x86_64 here, but should allow other archs, + // probably just auto detect as a default + qemuCmd := exec.Command("qemu-system-x86_64", argv...) + qemuCmd.Env = append(os.Environ(), "TMPDIR="+tempDir) + qemuCmd.Stdout = os.Stdout + qemuCmd.Stderr = os.Stderr + + startedBootAt := time.Now() + if err := qemuCmd.Start(); err != nil { + cleanup() + return nil, fmt.Errorf("starting qemu: %w", err) + } + state.Process = qemuCmd.Process + + qmpCtx, cancelQmp := context.WithTimeout(ctx, 10*time.Second) + defer cancelQmp() + + var mon *qmp.SocketMonitor + // wait for qmp to be ready + for { + mon, err = qmp.NewSocketMonitor("unix", qmpSock, 2*time.Second) + if err == nil { + if err = mon.Connect(); err == nil { + state.QMPMon = mon + break + } + } + select { + case <-qmpCtx.Done(): + cleanup() + return nil, fmt.Errorf("qmp connect timeout: %w", err) + case <-time.After(10 * time.Millisecond): + } + } + + status, err := state.QMPQueryStatus() + if err != nil { + cleanup() + return nil, err + } + + l.Info("qemu guest status", "status", status) + if status != "running" { + cleanup() + return nil, fmt.Errorf("qemu guest not running (status: %s)", status) + } + + qgaCtx, cancelQga := context.WithTimeout(ctx, time.Minute) + defer cancelQga() + + // wait for guest agent to be ready (aka boot) + for { + if _, err := os.Stat(qgaSock); err == nil { + pingCtx, cancelPing := context.WithTimeout(qgaCtx, 5*time.Second) + err = state.QGAGuestPing(pingCtx) + cancelPing() + if err == nil { + l.Info("vm booted and guest-agent ready", "elapsed", time.Since(startedBootAt).Round(time.Millisecond)) + break + } + l.Debug("qga guest-ping failed", "error", err) + } + + select { + case <-qgaCtx.Done(): + cleanup() + return nil, fmt.Errorf("qga connect timeout: %w", err) + case <-time.After(100 * time.Millisecond): + } + } + + return state, nil +} + +func (s *VMState) Close() error { + if s.QMPMon != nil { + _ = s.QMPMon.Disconnect() + } + if s.Process != nil { + _ = s.Process.Kill() + } + return nil +} + +func (s *VMState) QMPRun(command qmp.Command) ([]byte, error) { + b, err := json.Marshal(command) + if err != nil { + return nil, err + } + return s.QMPMon.Run(b) +} + +// sends a command to the qemu guest agent and returns the response +func (s *VMState) QGARun(ctx context.Context, command qmp.Command) ([]byte, error) { + b, err := json.Marshal(command) + if err != nil { + return nil, err + } + + conn, err := (&net.Dialer{}).DialContext(ctx, "unix", s.QGAPath) + if err != nil { + return nil, err + } + defer conn.Close() + + if dl, ok := ctx.Deadline(); ok { + _ = conn.SetDeadline(dl) + } + + if _, err := conn.Write(append(b, '\n')); err != nil { + return nil, err + } + + return bufio.NewReader(conn).ReadBytes('\n') +} + +// query the qemu guest status +func (s *VMState) QMPQueryStatus() (string, error) { + raw, err := s.QMPRun(qmp.Command{Execute: "query-status"}) + if err != nil { + return "", fmt.Errorf("qmp query-status failed: %w", err) + } + + var resp struct { + Return struct { + Status string `json:"status"` + } `json:"return"` + } + if err := json.Unmarshal(raw, &resp); err != nil { + return "", fmt.Errorf("qmp query-status parse: %w", err) + } + return resp.Return.Status, nil +} + +// send a command to qemu to powerdown the guest gracefully +func (s *VMState) QMPSystemPowerdown() error { + _, err := s.QMPRun(qmp.Command{Execute: "system_powerdown"}) + return err +} + +// ping the guest agent to see if it's ready +func (s *VMState) QGAGuestPing(ctx context.Context) error { + _, err := s.QGARun(ctx, qmp.Command{Execute: "guest-ping"}) + return err +} + +// execute a command on the guest and return its pid +func (s *VMState) QGAGuestExec(ctx context.Context, command string, env []string) (int, error) { + cmdArgs := []string{s.Img.Shell, "-c", command} + raw, err := s.QGARun(ctx, qmp.Command{ + Execute: "guest-exec", + Args: map[string]any{ + "path": s.Img.Shell, + "arg": cmdArgs[1:], + "env": env, + "capture-output": true, + }, + }) + if err != nil { + return 0, fmt.Errorf("qga guest-exec: %w", err) + } + + var resp struct { + Return struct { + Pid int `json:"pid"` + } `json:"return"` + } + if err := json.Unmarshal(raw, &resp); err != nil { + return 0, fmt.Errorf("qga guest-exec parse: %w", err) + } + return resp.Return.Pid, nil +} + +type GuestExecStatus struct { + Exited bool + ExitCode int + // these are since the last time we asked for status + OutData []byte + ErrData []byte +} + +// get the status of a command executed with guest-exec +func (s *VMState) QGAGuestExecStatus(ctx context.Context, pid int) (GuestExecStatus, error) { + var status GuestExecStatus + raw, err := s.QGARun(ctx, qmp.Command{ + Execute: "guest-exec-status", + Args: map[string]any{"pid": pid}, + }) + if err != nil { + return status, fmt.Errorf("qga guest-exec-status: %w", err) + } + + var resp struct { + Return struct { + Exited bool `json:"exited"` + ExitCode int `json:"exitcode"` + OutData string `json:"out-data"` + ErrData string `json:"err-data"` + } `json:"return"` + } + if err := json.Unmarshal(raw, &resp); err != nil { + return status, fmt.Errorf("qga guest-exec-status parse: %w", err) + } + + status.Exited = resp.Return.Exited + status.ExitCode = resp.Return.ExitCode + if resp.Return.OutData != "" { + status.OutData, _ = base64.StdEncoding.DecodeString(resp.Return.OutData) + } + if resp.Return.ErrData != "" { + status.ErrData, _ = base64.StdEncoding.DecodeString(resp.Return.ErrData) + } + + return status, nil +}