This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

1package microvm 2 3import ( 4 "context" 5 "crypto/sha256" 6 _ "embed" 7 "encoding/hex" 8 "encoding/json" 9 "errors" 10 "fmt" 11 "log/slog" 12 "os" 13 "os/exec" 14 "path/filepath" 15 "strconv" 16 "strings" 17 "sync" 18 "time" 19 20 "github.com/digitalocean/go-qemu/qmp" 21 "github.com/google/uuid" 22) 23 24const ( 25 defaultQMPTimeout = 10 * time.Second 26 outerSlirpCIDR = "10.0.2.0/24" 27 innerSlirpNet = "10.0.3.0/24" 28 innerSlirpHost = "10.0.3.2" 29 innerSlirpDNS = "10.0.3.3" 30 innerSlirpDHCP = "10.0.3.15" 31 netnsTapName = "tap0" 32 netnsMTU = "65520" 33) 34 35type QEMUConfig struct { 36 Image ImageSpec 37 BootTimeout time.Duration 38 CID uint32 39 EnableKVM bool 40 QEMULogPath string 41 QMPPath string 42 SerialLogPath string 43 WorkDir string 44 VolumePaths map[string]string 45 VolumeBaseName string 46 Cgroup CgroupLimits 47 Dev bool 48} 49 50type QEMUVMHandle struct { 51 cid uint32 52 Process *os.Process 53 qemuLogPath string 54 QMPMon *qmp.SocketMonitor 55 QMPPath string 56 serialLogPath string 57 workDir string 58 59 qmpSocketPath string 60 61 cmd *exec.Cmd 62 done chan struct{} 63 qemuLogFile *os.File 64 cgroup *CgroupHandle 65 slirpCmd *exec.Cmd 66 slirpExit *os.File 67 waitErr error 68 waitErrMu sync.Mutex 69} 70 71type qemuRunner struct{} 72 73func (qemuRunner) Validate(spec ImageSpec, enableKVM bool) error { 74 if _, err := exec.LookPath(spec.RunnerCmd()); err != nil { 75 return fmt.Errorf("required host command %q not found in PATH: %w", spec.RunnerCmd(), err) 76 } 77 if _, err := os.Stat("/dev/vhost-vsock"); err != nil { 78 return fmt.Errorf("microvm requires /dev/vhost-vsock for vhost-vsock-device: %w", err) 79 } 80 if enableKVM { 81 if _, err := os.Stat("/dev/kvm"); err != nil { 82 return fmt.Errorf("microvm KVM was requested but /dev/kvm is not accessible: %w", err) 83 } 84 } 85 if len(spec.NetworkInterfaces) > 0 { 86 if _, err := os.Stat("/dev/net/tun"); err != nil { 87 return fmt.Errorf("microvm slirp4netns networking requires /dev/net/tun: %w", err) 88 } 89 for _, cmd := range []string{"ip", "mount", "slirp4netns", "unshare"} { 90 if _, err := exec.LookPath(cmd); err != nil { 91 return fmt.Errorf("required host command %q not found in PATH: %w", cmd, err) 92 } 93 } 94 } 95 return nil 96} 97 98func (qemuRunner) Start(ctx context.Context, cfg VMConfig, volumePaths map[string]string, logger *slog.Logger) (VMHandle, error) { 99 bootTimeout := cfg.BootTimeout 100 if bootTimeout == 0 { 101 bootTimeout = 10 * time.Second 102 } 103 return StartQEMU(ctx, QEMUConfig{ 104 Image: cfg.Image, 105 BootTimeout: bootTimeout, 106 CID: cfg.CID, 107 EnableKVM: cfg.EnableKVM, 108 WorkDir: cfg.WorkDir, 109 VolumePaths: volumePaths, 110 Cgroup: cfg.Cgroup, 111 Dev: cfg.Dev, 112 }, logger) 113} 114 115// we hash the workDir to get a deterministic qmpSock path 116// since they are AF_UNIX sockets, they are bound by a 108 char long path limit... 117func qmpSocketPath(workDir string) string { 118 base := filepath.Dir(workDir) 119 if workDir == "" { 120 base = os.TempDir() 121 } 122 sum := sha256.Sum256([]byte(workDir)) 123 return filepath.Join(base, hex.EncodeToString(sum[:8])+".qmp.sock") 124} 125 126func StartQEMU(ctx context.Context, cfg QEMUConfig, logger *slog.Logger) (VMHandle, error) { 127 if logger == nil { 128 logger = slog.Default() 129 } 130 131 workDir := cfg.WorkDir 132 133 handle := &QEMUVMHandle{ 134 workDir: workDir, 135 } 136 137 var ok bool 138 defer func() { 139 if !ok { 140 _ = handle.Close() 141 } 142 }() 143 144 cid := cfg.CID 145 if cid == 0 { 146 var err error 147 cid, err = AllocateCID() 148 if err != nil { 149 return nil, err 150 } 151 } 152 if cid < minGuestCID { 153 return nil, fmt.Errorf("guest CID must be >= %d", minGuestCID) 154 } 155 handle.cid = cid 156 157 volumePaths := cfg.VolumePaths 158 159 qemuLogPath := cfg.QEMULogPath 160 if qemuLogPath == "" { 161 qemuLogPath = filepath.Join(workDir, "qemu.log") 162 } 163 qemuLogFile, err := createParentedFile(qemuLogPath) 164 if err != nil { 165 return nil, err 166 } 167 handle.qemuLogPath = qemuLogPath 168 handle.qemuLogFile = qemuLogFile 169 170 serialLogPath := cfg.SerialLogPath 171 if serialLogPath == "" { 172 serialLogPath = filepath.Join(workDir, "serial.log") 173 } 174 if err := os.MkdirAll(filepath.Dir(serialLogPath), 0o755); err != nil { 175 return nil, fmt.Errorf("create serial log directory: %w", err) 176 } 177 handle.serialLogPath = serialLogPath 178 179 qmpPath := cfg.QMPPath 180 if qmpPath == "" { 181 qmpPath = qmpSocketPath(workDir) 182 handle.qmpSocketPath = qmpPath 183 } 184 handle.QMPPath = qmpPath 185 186 qemuCmd := cfg.Image.RunnerCmd() 187 qemuBinary, err := exec.LookPath(qemuCmd) 188 if err != nil { 189 return nil, fmt.Errorf("%s command not found in PATH: %w", qemuCmd, err) 190 } 191 192 args, err := qemuArgs(qemuArgsConfig{ 193 Image: cfg.Image, 194 CID: cid, 195 EnableKVM: cfg.EnableKVM, 196 QMPPath: qmpPath, 197 SerialLogPath: serialLogPath, 198 VolumePaths: volumePaths, 199 }) 200 if err != nil { 201 return nil, err 202 } 203 204 cmd, slirpNet, err := qemuCommand(ctx, qemuBinary, args, cfg.Image, workDir, cfg.Dev) 205 if err != nil { 206 return nil, err 207 } 208 cmd.Env = append(os.Environ(), "TMPDIR="+workDir) 209 cmd.Stdout = qemuLogFile 210 cmd.Stderr = qemuLogFile 211 212 cgroup, err := prepareCgroup(cfg.Cgroup, logger) 213 if err != nil { 214 return nil, err 215 } 216 handle.cgroup = cgroup 217 218 logger.Info("starting qemu microvm", "cid", cid, "workDir", workDir, "serialLog", serialLogPath, "qmp", qmpPath) 219 if err := cmd.Start(); err != nil { 220 return nil, fmt.Errorf("starting qemu: %w", err) 221 } 222 handle.cmd = cmd 223 handle.Process = cmd.Process 224 handle.done = make(chan struct{}) 225 go func() { 226 err := cmd.Wait() 227 handle.waitErrMu.Lock() 228 handle.waitErr = err 229 handle.waitErrMu.Unlock() 230 close(handle.done) 231 }() 232 233 if err := cgroup.AddProcess(cmd.Process.Pid, logger); err != nil { 234 return nil, err 235 } 236 237 if slirpNet != nil { 238 handle.slirpCmd, handle.slirpExit, err = slirpNet.Start(ctx, qemuLogFile, logger) 239 if err != nil { 240 return nil, err 241 } 242 if handle.slirpCmd != nil && handle.slirpCmd.Process != nil { 243 if err := cgroup.AddProcess(handle.slirpCmd.Process.Pid, logger); err != nil { 244 return nil, err 245 } 246 } 247 } 248 249 qmpTimeout := cfg.BootTimeout 250 if qmpTimeout == 0 { 251 qmpTimeout = defaultQMPTimeout 252 } 253 if err := handle.waitForQMP(ctx, qmpTimeout); err != nil { 254 return nil, err 255 } 256 257 status, err := handle.QMPQueryStatus() 258 if err != nil { 259 return nil, err 260 } 261 if status != "running" { 262 return nil, fmt.Errorf("qemu guest not running (status: %s)", status) 263 } 264 logger.Info("qemu microvm running", "cid", cid, "status", status) 265 266 ok = true 267 return handle, nil 268} 269 270func (h *QEMUVMHandle) Wait() error { 271 if h == nil || h.done == nil { 272 return nil 273 } 274 <-h.done 275 h.waitErrMu.Lock() 276 defer h.waitErrMu.Unlock() 277 return h.waitErr 278} 279 280func (h *QEMUVMHandle) WaitContext(ctx context.Context) error { 281 if h == nil || h.done == nil { 282 return nil 283 } 284 select { 285 case <-h.done: 286 h.waitErrMu.Lock() 287 defer h.waitErrMu.Unlock() 288 return h.waitErr 289 case <-ctx.Done(): 290 return ctx.Err() 291 } 292} 293 294func (h *QEMUVMHandle) Kill() error { 295 if h == nil || h.Process == nil { 296 return nil 297 } 298 return h.Process.Kill() 299} 300 301func (h *QEMUVMHandle) Shutdown(ctx context.Context) error { 302 if h == nil { 303 return nil 304 } 305 if h.QMPMon != nil { 306 if err := h.QMPSystemPowerdown(); err != nil { 307 return err 308 } 309 } 310 if h.done == nil { 311 return nil 312 } 313 select { 314 case <-h.done: 315 return h.Wait() 316 case <-ctx.Done(): 317 _ = h.Kill() 318 _ = h.Wait() 319 return ctx.Err() 320 } 321} 322 323func (h *QEMUVMHandle) Close() error { 324 if h == nil { 325 return nil 326 } 327 328 var closeErr error 329 if h.QMPMon != nil { 330 closeErr = errors.Join(closeErr, h.QMPMon.Disconnect()) 331 h.QMPMon = nil 332 } 333 if h.Process != nil { 334 _ = h.Process.Kill() 335 _ = h.Wait() 336 } 337 if h.slirpExit != nil { 338 _ = h.slirpExit.Close() 339 h.slirpExit = nil 340 } 341 if h.slirpCmd != nil && h.slirpCmd.Process != nil { 342 _ = h.slirpCmd.Process.Kill() 343 _ = h.slirpCmd.Wait() 344 h.slirpCmd = nil 345 } 346 if h.qemuLogFile != nil { 347 closeErr = errors.Join(closeErr, h.qemuLogFile.Close()) 348 h.qemuLogFile = nil 349 } 350 if h.cgroup != nil { 351 closeErr = errors.Join(closeErr, h.cgroup.Close()) 352 h.cgroup = nil 353 } 354 if h.qmpSocketPath != "" { 355 if err := os.Remove(h.qmpSocketPath); err != nil && !os.IsNotExist(err) { 356 closeErr = errors.Join(closeErr, err) 357 } 358 h.qmpSocketPath = "" 359 } 360 return closeErr 361} 362 363func (h *QEMUVMHandle) QMPRun(command qmp.Command) ([]byte, error) { 364 if h == nil || h.QMPMon == nil { 365 return nil, fmt.Errorf("qmp monitor is not connected") 366 } 367 data, err := json.Marshal(command) 368 if err != nil { 369 return nil, err 370 } 371 return h.QMPMon.Run(data) 372} 373 374func (h *QEMUVMHandle) QMPQueryStatus() (string, error) { 375 raw, err := h.QMPRun(qmp.Command{Execute: "query-status"}) 376 if err != nil { 377 return "", fmt.Errorf("qmp query-status failed: %w", err) 378 } 379 380 var resp struct { 381 Return struct { 382 Status string `json:"status"` 383 } `json:"return"` 384 } 385 if err := json.Unmarshal(raw, &resp); err != nil { 386 return "", fmt.Errorf("qmp query-status parse: %w", err) 387 } 388 return resp.Return.Status, nil 389} 390 391func (h *QEMUVMHandle) QMPSystemPowerdown() error { 392 _, err := h.QMPRun(qmp.Command{Execute: "system_powerdown"}) 393 return err 394} 395 396func (h *QEMUVMHandle) Logs() VMLogs { 397 if h == nil { 398 return VMLogs{} 399 } 400 return VMLogs{ 401 Serial: h.serialLogPath, 402 Extra: map[string]string{ 403 "qemu": h.qemuLogPath, 404 }, 405 } 406} 407 408func (h *QEMUVMHandle) CID() uint32 { 409 if h == nil { 410 return 0 411 } 412 return h.cid 413} 414 415func (h *QEMUVMHandle) WorkDir() string { 416 if h == nil { 417 return "" 418 } 419 return h.workDir 420} 421 422func (h *QEMUVMHandle) OOMKilled() bool { 423 if h == nil { 424 return false 425 } 426 return h.cgroup.OOMKilled() 427} 428 429func (h *QEMUVMHandle) waitForQMP(ctx context.Context, timeout time.Duration) error { 430 qmpCtx, cancel := context.WithTimeout(ctx, timeout) 431 defer cancel() 432 433 var lastErr error 434 for { 435 mon, err := qmp.NewSocketMonitor("unix", h.QMPPath, 2*time.Second) 436 if err == nil { 437 if err = mon.Connect(); err == nil { 438 h.QMPMon = mon 439 return nil 440 } 441 _ = mon.Disconnect() 442 } 443 lastErr = err 444 445 select { 446 case <-qmpCtx.Done(): 447 return fmt.Errorf("qmp connect timeout: %w", lastErr) 448 case <-h.done: 449 return fmt.Errorf("qemu exited before qmp was ready: %w", h.Wait()) 450 case <-time.After(25 * time.Millisecond): 451 } 452 } 453} 454 455func qemuCommand( 456 ctx context.Context, 457 qemuBinary string, 458 args []string, 459 spec ImageSpec, 460 workDir string, 461 dev bool, 462) (*exec.Cmd, *slirpNamespace, error) { 463 if len(spec.NetworkInterfaces) == 0 { 464 return exec.CommandContext(ctx, qemuBinary, args...), nil, nil 465 } 466 467 ipPath, err := exec.LookPath("ip") 468 if err != nil { 469 return nil, nil, fmt.Errorf("ip command not found in PATH: %w", err) 470 } 471 mountPath, err := exec.LookPath("mount") 472 if err != nil { 473 return nil, nil, fmt.Errorf("mount command not found in PATH: %w", err) 474 } 475 unsharePath, err := exec.LookPath("unshare") 476 if err != nil { 477 return nil, nil, fmt.Errorf("unshare command not found in PATH: %w", err) 478 } 479 480 pidFile, resolvPath, wrapperPath, err := prepareQEMUNetnsFiles(workDir, dev) 481 if err != nil { 482 return nil, nil, err 483 } 484 485 cmdArgs := append([]string{ 486 "--user", 487 "--map-root-user", 488 "--net", 489 "--mount", 490 "--propagation", "private", 491 "--", 492 wrapperPath, 493 pidFile, 494 ipPath, 495 mountPath, 496 resolvPath, 497 qemuBinary, 498 }, args...) 499 500 cmd := exec.CommandContext(ctx, unsharePath, cmdArgs...) 501 502 return cmd, &slirpNamespace{ 503 spec: spec, 504 pidFile: pidFile, 505 dev: dev, 506 }, nil 507} 508 509func prepareQEMUNetnsFiles(workDir string, dev bool) (pidFile, resolvPath, wrapperPath string, err error) { 510 pidFile = filepath.Join(workDir, "qemu-netns.pid") 511 resolvPath = filepath.Join(workDir, "qemu-netns-resolv.conf") 512 wrapperPath = filepath.Join(workDir, "qemu-netns-wrapper") 513 514 // the guest resolves through shuttle on 127.0.0.1. keep qemu's slirp DNS 515 // pointed at an unroutable local resolver inside this network namespace so 516 // direct guest queries to 10.0.3.3 don't bypass the shuttle dns policy. 517 if err := os.WriteFile(resolvPath, []byte("nameserver 127.0.0.1\n"), 0o644); err != nil { 518 return "", "", "", fmt.Errorf("write qemu network namespace resolv.conf: %w", err) 519 } 520 521 if err := writeNetnsWrapper(wrapperPath, dev); err != nil { 522 return "", "", "", fmt.Errorf("write qemu network namespace wrapper: %w", err) 523 } 524 525 return pidFile, resolvPath, wrapperPath, nil 526} 527 528type qemuArgsConfig struct { 529 Image ImageSpec 530 CID uint32 531 EnableKVM bool 532 QMPPath string 533 SerialLogPath string 534 VolumePaths map[string]string 535} 536 537func qemuArgs(cfg qemuArgsConfig) ([]string, error) { 538 uuid := uuid.New() 539 540 b := newArgBuilder(64) 541 542 addQEMUMachineArgs(&b, cfg, uuid) 543 addQEMUStoreArgs(&b, cfg) 544 545 if cfg.EnableKVM { 546 addQEMUKVMArgs(&b, cfg.Image) 547 } 548 549 if err := addQEMUVolumeArgs(&b, cfg); err != nil { 550 return nil, err 551 } 552 553 if err := addQEMUNetworkArgs(&b, cfg.Image.NetworkInterfaces); err != nil { 554 return nil, err 555 } 556 557 b.Optf("-device", "vhost-vsock-device,guest-cid=%d", cfg.CID) 558 559 if len(cfg.Image.RunnerConfig.ExtraArgs) > 0 { 560 b.Add(cfg.Image.RunnerConfig.ExtraArgs...) 561 } 562 563 return b.Args(), nil 564} 565 566func addQEMUMachineArgs(b *argBuilder, cfg qemuArgsConfig, uuid uuid.UUID) { 567 if cfg.Image.RunnerConfig.Machine != "" { 568 b.Opt("-M", cfg.Image.RunnerConfig.Machine) 569 } 570 b.Optf("-m", "%dM", cfg.Image.MemoryMiB) 571 b.Opt("-smp", strconv.Itoa(cfg.Image.VCPUs)) 572 573 b.Add( 574 "-nodefaults", 575 "-no-user-config", 576 "-no-reboot", 577 ) 578 579 b.Opt("-kernel", cfg.Image.Kernel) 580 b.Opt("-initrd", cfg.Image.Initrd) 581 582 b.Opt("-device", "virtio-rng-device") 583 584 b.Optf("-smbios", "type=1,uuid=%s", uuid) 585 b.Opt("-serial", "file:"+cfg.SerialLogPath) 586 587 // use virtio console if requsted. this is faster than the serial UART logging 588 // because serial has a higher cost when being accesssed. we still have to 589 // support serial itself for early kernel boot but thats OK. 590 if cfg.Image.RunnerConfig.Console == "hvc0" { 591 b.Optf("-chardev", "file,id=virtiocon0,path=%s,append=on", cfg.SerialLogPath) 592 b.Add("-device", "virtio-serial-device") 593 b.Opt("-device", "virtconsole,chardev=virtiocon0") 594 } 595 b.Opt("-display", "none") 596 b.Opt("-monitor", "none") 597 b.Opt("-append", cfg.Image.BootArgs) 598 599 b.Opt("-sandbox", "on") 600 b.Optf("-qmp", "unix:%s,server,nowait", cfg.QMPPath) 601} 602 603func addQEMUStoreArgs(b *argBuilder, cfg qemuArgsConfig) { 604 drive := newOptionBuilder(8) 605 drive.KV("id", "store") 606 drive.KV("format", "raw") 607 drive.Add("read-only=on") 608 drive.KV("file", cfg.Image.StoreDisk) 609 drive.Add("if=none") 610 drive.Add("aio=io_uring") 611 612 b.Opt("-drive", drive.String()) 613 b.Opt("-device", "virtio-blk-device,drive=store") 614} 615 616func addQEMUKVMArgs(b *argBuilder, image ImageSpec) { 617 b.Flag("-enable-kvm") 618 if image.RunnerConfig.CPU != "" { 619 b.Opt("-cpu", image.RunnerConfig.CPU) 620 } 621 b.Opt("-device", "i8042") 622} 623 624func addQEMUVolumeArgs(b *argBuilder, cfg qemuArgsConfig) error { 625 for index, volume := range cfg.Image.Volumes { 626 path := cfg.VolumePaths[volume.Image] 627 if path == "" { 628 return fmt.Errorf("missing prepared path for volume %q", volume.Image) 629 } 630 631 driveID := fmt.Sprintf("volume%d", index) 632 633 drive := newOptionBuilder(10) 634 drive.KV("id", driveID) 635 drive.KV("format", "raw") 636 drive.Add("read-only=off") 637 drive.KV("file", path) 638 drive.Add("if=none") 639 drive.Add("aio=io_uring") 640 drive.Add("discard=unmap") 641 drive.Add("cache=none") 642 643 b.Opt("-drive", drive.String()) 644 b.Optf("-device", "virtio-blk-device,drive=%s", driveID) 645 } 646 647 return nil 648} 649 650func addQEMUNetworkArgs(b *argBuilder, interfaces []NetworkInterface) error { 651 for _, networkInterface := range interfaces { 652 if networkInterface.Type != "slirp4netns" { 653 return fmt.Errorf("unsupported microvm network interface type %q", networkInterface.Type) 654 } 655 656 netdevOpts := newOptionBuilder(6) 657 netdevOpts.Add("user") 658 netdevOpts.KV("id", networkInterface.ID) 659 netdevOpts.KV("net", innerSlirpNet) 660 netdevOpts.KV("host", innerSlirpHost) 661 netdevOpts.KV("dns", innerSlirpDNS) 662 netdevOpts.KV("dhcpstart", innerSlirpDHCP) 663 664 b.Opt("-netdev", netdevOpts.String()) 665 b.Optf( 666 "-device", "virtio-net-device,netdev=%s,mac=%s", 667 networkInterface.ID, networkInterface.MAC, 668 ) 669 } 670 671 return nil 672} 673 674func waitForPIDFile(ctx context.Context, path string) (string, error) { 675 waitCtx, cancel := context.WithTimeout(ctx, 5*time.Second) 676 defer cancel() 677 678 ticker := time.NewTicker(25 * time.Millisecond) 679 defer ticker.Stop() 680 681 for { 682 data, err := os.ReadFile(path) 683 if err == nil { 684 pid := strings.TrimSpace(string(data)) 685 if pid != "" { 686 return pid, nil 687 } 688 } else if !errors.Is(err, os.ErrNotExist) { 689 return "", fmt.Errorf("read qemu network namespace pid: %w", err) 690 } 691 692 select { 693 case <-waitCtx.Done(): 694 return "", fmt.Errorf("waiting for qemu network namespace pid: %w", waitCtx.Err()) 695 case <-ticker.C: 696 } 697 } 698}