Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555//go:build linux
package microvm
import ( "context" "crypto/rand" "encoding/binary" "errors" "fmt" "io" "log/slog" "maps" "math" "net" "os" "os/exec" "path/filepath" "regexp" "slices" "strings" "sync" "sync/atomic" "time"
"tangled.org/core/spindle/engine" "tangled.org/core/spindle/models" "tangled.org/core/spindle/quota")
const ( minGuestCID = 3 maxGuestCID = minGuestCID + 59999 vmCrashLogTailBytes = 8192)
func AllocateCID() (uint32, error) { var data [4]byte if _, err := rand.Read(data[:]); err != nil { return 0, fmt.Errorf("allocate guest CID: %w", err) } return minGuestCID + binary.BigEndian.Uint32(data[:])%60000, nil}
type cidRegistry struct { mu sync.Mutex next uint32 inUse map[uint32]struct{}}
func newCIDRegistry() *cidRegistry { return &cidRegistry{next: minGuestCID, inUse: make(map[uint32]struct{})}}
func (r *cidRegistry) allocate() (uint32, error) { r.mu.Lock() defer r.mu.Unlock() for range maxGuestCID - minGuestCID + 1 { cid := r.next r.next++ if r.next > maxGuestCID { r.next = minGuestCID } if _, exists := r.inUse[cid]; exists { continue } r.inUse[cid] = struct{}{} return cid, nil } return 0, errors.New("no guest CIDs available")}
func (r *cidRegistry) release(cid uint32) { if cid == 0 { return } r.mu.Lock() delete(r.inUse, cid) r.mu.Unlock()}
func prepareWorkDir(workDir string) error { if workDir == "" { return fmt.Errorf("microvm work directory is required") } if err := os.MkdirAll(workDir, 0o755); err != nil { return fmt.Errorf("create microvm work directory: %w", err) } return nil}
func mkfsExt4ForVolumes(volumes []Volume, configured string) (string, error) { if len(volumes) == 0 || configured != "" { return configured, nil } path, err := exec.LookPath("mkfs.ext4") if err != nil { return "", fmt.Errorf("mkfs.ext4 command not found in PATH: %w", err) } return path, nil}
func prepareVolumes(ctx context.Context, workDir string, volumes []Volume, mkfsExt4 string) (map[string]string, error) { paths := make(map[string]string, len(volumes)) for _, volume := range volumes { if volume.ReadOnly { return nil, fmt.Errorf("read-only microvm volume %q is not supported yet", volume.Image) } if volume.FSType != "ext4" { return nil, fmt.Errorf("microvm volume %q uses unsupported fsType %q", volume.Image, volume.FSType) } if volume.ImageType != "" && volume.ImageType != "raw" { return nil, fmt.Errorf("microvm volume %q uses unsupported imageType %q", volume.Image, volume.ImageType) }
path := filepath.Join(workDir, filepath.Base(volume.Image)) if err := createSparseFile(path, volume.SizeMiB); err != nil { return nil, err } noJournal := volume.MountPoint == "/workspace" if err := runMkfsExt4(ctx, mkfsExt4, path, noJournal); err != nil { return nil, err } paths[volume.Image] = path } return paths, nil}
func createSparseFile(path string, sizeMiB int64) error { if sizeMiB <= 0 { return fmt.Errorf("sparse file %q size must be positive", path) } if sizeMiB > math.MaxInt64/(1024*1024) { return fmt.Errorf("sparse file %q size is too large", path) } file, err := os.OpenFile(path, os.O_RDWR|os.O_CREATE|os.O_EXCL, 0o600) if err != nil { return fmt.Errorf("create sparse file %q: %w", path, err) } defer file.Close()
if err := file.Truncate(sizeMiB * 1024 * 1024); err != nil { return fmt.Errorf("resize sparse file %q: %w", path, err) } return nil}
func runMkfsExt4(ctx context.Context, mkfsExt4, path string, noJournal bool) error { if mkfsExt4 == "" { return fmt.Errorf("mkfs.ext4 path is required") } args := []string{"-F"} if noJournal { args = append(args, "-O", "^has_journal") } args = append(args, path)
cmd := exec.CommandContext(ctx, mkfsExt4, args...) output, err := cmd.CombinedOutput() if err != nil { return fmt.Errorf("mkfs.ext4 %q: %w: %s", path, err, strings.TrimSpace(string(output))) } return nil}
func createParentedFile(path string) (*os.File, error) { if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil { return nil, fmt.Errorf("create log directory: %w", err) } file, err := os.OpenFile(path, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, 0o644) if err != nil { return nil, fmt.Errorf("create log file %q: %w", path, err) } return file, nil}
type VMLogs struct { Serial string Extra map[string]string}
type VMHandle interface { Shutdown(ctx context.Context) error WaitContext(ctx context.Context) error Close() error Logs() VMLogs CID() uint32 WorkDir() string OOMKilled() bool}
type VMResourceReporter interface { VolumeUsage() (map[string]int64, error)}
type VMConfig struct { Image ImageSpec CID uint32 EnableKVM bool WorkDir string Cgroup CgroupLimits
BootTimeout time.Duration MkfsExt4 string Dev bool}
type workflowState struct { ImageSpec ImageSpec ImageSpecPath string Config manifestConfig ConfigKey string Image string SubstituterReadURLs []string SubstituterTrustedPublicKeys []string VM VMHandle CID uint32 Agent *AgentSession Substituter *SubstituterProxy SubstituterUpload *SubstituterUploadProxy DNSProxy *DNSProxy WorkDir string NixOSToplevels nixosToplevelStore OwnerDID string RepoDID string QuotaStore quota.ReservationStore StartedAt time.Time // when the VM booted, for the max-lifetime cap ResourceUsage *engine.WorkflowResourceUsage ResourceUsageAvailable bool resourceUsageMu sync.Mutex}
func (e *Engine) cleanupState(ctx context.Context, wid models.WorkflowId, state *workflowState) error { if state == nil { return nil }
// stop advertising this VM for debug shells before we tear it down e.unregisterDebugTarget(wid)
ctx = context.WithoutCancel(ctx)
var err error // todo(dawn): expose this error to the user as a warning if drainErr := e.drainNixCache(ctx, state); drainErr != nil { e.l.Warn("cache drain failed during cleanup; continuing", "workflow", wid, "error", drainErr) } err = errors.Join(err, e.shutdownVM(ctx, wid, state)) err = errors.Join(err, closeIO(&state.Agent)) err = errors.Join(err, closeIO(&state.Substituter)) err = errors.Join(err, closeIO(&state.SubstituterUpload)) err = errors.Join(err, closeIO(&state.DNSProxy)) err = errors.Join(err, removeWorkDir(state)) e.cids.release(state.CID) state.CID = 0 return err}
func (e *Engine) drainNixCache(ctx context.Context, state *workflowState) error { if e.cfg.NixCache.UploadURL == "" { return nil }
drainCtx, cancel := context.WithTimeout(ctx, substituterDrainTimeout) defer cancel()
if state.Agent != nil { if _, err := state.Agent.Drain(drainCtx); err != nil { return fmt.Errorf("drain guest nix cache uploads: %w", err) } } return nil}
func (e *Engine) shutdownVM(ctx context.Context, wid models.WorkflowId, state *workflowState) error { if state.VM == nil { return nil } if vmExited(state.VM) { return e.closeVMWithUsage(wid, state) }
var poweroffErr error
if state.Agent != nil { gracefulCtx, cancel := context.WithTimeout(ctx, vmShutdownTimeout) var poweredOff bool poweredOff, poweroffErr = e.poweroffViaAgent(gracefulCtx, wid, state) cancel()
if poweredOff || vmExited(state.VM) { return e.closeVMWithUsage(wid, state) } }
fallbackCtx, cancel := context.WithTimeout(ctx, vmShutdownTimeout) defer cancel()
shutdownErr := state.VM.Shutdown(fallbackCtx) if shutdownErr != nil && !vmExited(state.VM) { e.l.Warn("microVM shutdown fallback failed", "workflow", wid, "error", shutdownErr) return errors.Join(poweroffErr, shutdownErr, e.closeVMWithUsage(wid, state)) }
return e.closeVMWithUsage(wid, state)}
func (e *Engine) closeVMWithUsage(wid models.WorkflowId, state *workflowState) error { if qvm, ok := state.VM.(*QEMUVMHandle); ok { err := qvm.close(func() { e.captureResourceUsage(wid, state) }) state.VM = nil return err } e.captureResourceUsage(wid, state) 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) { if err := state.Agent.Poweroff(ctx); err != nil { e.l.Warn("agent poweroff request failed", "workflow", wid, "error", err) return false, err }
if err := state.VM.WaitContext(ctx); err != nil { e.l.Warn("agent poweroff did not stop microVM", "workflow", wid, "error", err) return false, nil }
return true, nil}
// helper for closing io interfaces, sets to nil to prevent double-closefunc closeIO[T io.Closer](field *T) error { closer := *field var zero T *field = zero if any(closer) == any(zero) { return nil } return closer.Close()}
func removeWorkDir(state *workflowState) error { if state.WorkDir == "" { return nil }
err := os.RemoveAll(state.WorkDir) state.WorkDir = "" return err}
// returns a context derived from ctx that is cancelled either when ctx itself// is cancelled or when the microVM exits on its own. the returned flag reports// whether the VM exited (as opposed to ctx being cancelled for another reason,// e.g. the workflow timeout), letting callers tell a crash apart from a// timeout. cancel must be called to release the watcher goroutine.func watchVMExit(ctx context.Context, vm VMHandle) (context.Context, *atomic.Bool, context.CancelFunc) { exited := &atomic.Bool{} watchCtx, cancel := context.WithCancel(ctx) if vm == nil { return watchCtx, exited, cancel } go func() { _ = vm.WaitContext(watchCtx) // returns when VM exits or watchCtx is cancelled if watchCtx.Err() == nil { exited.Store(true) cancel() // don't forget to cancel the watchCtx... } }() return watchCtx, exited, cancel}
func VMCrashLog(vm VMHandle) string { if vm == nil { return "" } logs := vm.Logs()
var b strings.Builder if tail := tailFile(logs.Serial, vmCrashLogTailBytes); tail != "" { fmt.Fprintf(&b, "==== serial log ====\n%s\n", tail) } for _, name := range slices.Sorted(maps.Keys(logs.Extra)) { if tail := tailFile(logs.Extra[name], vmCrashLogTailBytes); tail != "" { fmt.Fprintf(&b, "==== %s log ====\n%s\n", name, tail) } } return strings.TrimRight(b.String(), "\n")}
func tailFile(path string, max int64) string { if path == "" { return "" } f, err := os.Open(path) if err != nil { return "" } defer f.Close() if info, err := f.Stat(); err == nil && info.Size() > max { if _, err := f.Seek(-max, io.SeekEnd); err != nil { return "" } } data, err := io.ReadAll(f) if err != nil { return "" } return strings.TrimSpace(string(data))}
func waitAgentConn(ctx context.Context, connCh <-chan net.Conn) (net.Conn, error) { select { case conn := <-connCh: if conn == nil { return nil, fmt.Errorf("agent connection closed before setup") } return conn, nil case <-ctx.Done(): return nil, fmt.Errorf("waiting for agent: %w", ctx.Err()) }}
func StartVM(ctx context.Context, cfg VMConfig, logger *slog.Logger) (VMHandle, error) { if logger == nil { logger = slog.Default() }
runner, err := runnerFor(cfg.Image.RunnerType) if err != nil { return nil, err } if err := cfg.Image.Validate(); err != nil { return nil, err } if err := cfg.Image.validateImageFiles(); err != nil { return nil, err } if err := runner.Validate(cfg.Image, cfg.EnableKVM); err != nil { return nil, err }
if err := prepareWorkDir(cfg.WorkDir); err != nil { return nil, err }
mkfsExt4, err := mkfsExt4ForVolumes(cfg.Image.Volumes, cfg.MkfsExt4) if err != nil { return nil, err } volumePaths, err := prepareVolumes(ctx, cfg.WorkDir, cfg.Image.Volumes, mkfsExt4) if err != nil { return nil, err }
return runner.Start(ctx, cfg, volumePaths, logger)}
// checks serial log for ooms or kernel panic// this is very linux specific! but these strings are stable in linux itself, see mm/oom_kill.c and kernel/panic.cfunc ParseCrashLog(detail string) (error, bool) { if strings.Contains(detail, "Out of memory:") { // we can show process name where possible re := regexp.MustCompile(`Out of memory: Killed process \d+ \(([^)]+)\)`) matches := re.FindStringSubmatch(detail) if len(matches) > 1 { return fmt.Errorf("guest out of memory (process '%s' killed by guest kernel OOM)", matches[1]), true } return errors.New("guest out of memory (OOM killer invoked)"), true } if strings.Contains(detail, "Kernel panic") { return errors.New("guest kernel panic"), true } return nil, false}
func (e *Engine) captureResourceUsage(wid models.WorkflowId, state *workflowState) { if state == nil { return }
state.resourceUsageMu.Lock() defer state.resourceUsageMu.Unlock() if state.ResourceUsageAvailable { return }
var usage engine.WorkflowResourceUsage
var cgroupErr error var volumesErr error var cgroupCaptured bool var volumesCaptured bool
if qvm, ok := state.VM.(*QEMUVMHandle); ok && qvm != nil { cgroupUsage, ok, err := qvm.cgroup.Stat() if err != nil { cgroupErr = err } else if ok { cgroupCaptured = true usage.CPUUsec = cgroupUsage.CPU.UsageUsec usage.MemoryCurrentBytes = cgroupUsage.Memory.Usage usage.MemoryPeakBytes = cgroupUsage.Memory.MaxUsage usage.SwapCurrentBytes = cgroupUsage.Memory.SwapUsage usage.SwapPeakBytes = cgroupUsage.Memory.SwapMaxUsage usage.PIDsCurrent = cgroupUsage.Pids.Current usage.IOReadBytes = cgroupUsage.IO.Rbytes usage.IOWriteBytes = cgroupUsage.IO.Wbytes usage.IOReadOps = cgroupUsage.IO.Rios usage.IOWriteOps = cgroupUsage.IO.Wios usage.CgroupAvailable = true } }
if reporter, ok := state.VM.(VMResourceReporter); ok && reporter != nil { volUsage, err := reporter.VolumeUsage() if err != nil { volumesErr = err } if len(volUsage) > 0 { volumesCaptured = true for _, allocated := range volUsage { if allocated > 0 { usage.VolumeAllocatedBytes += uint64(allocated) } } usage.VolumeAvailable = err == nil } }
if cgroupErr != nil || volumesErr != nil { e.l.Warn("resource usage measurement failed", "workflow_id", wid.String(), "cgroup_error", cgroupErr, "volumes_error", volumesErr, ) }
if cgroupCaptured || volumesCaptured { state.ResourceUsage = &usage state.ResourceUsageAvailable = true }}