Monorepo for Tangled tangled.org
1

Configure Feed

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

core / spindle / engines / microvm / engine.go
18 kB 614 lines
1//go:build linux 2 3package microvm 4 5import ( 6 "context" 7 "encoding/json" 8 "errors" 9 "fmt" 10 "io" 11 "log/slog" 12 "net/http" 13 "os" 14 "path/filepath" 15 "slices" 16 "strings" 17 "sync" 18 "sync/atomic" 19 "time" 20 21 "gopkg.in/yaml.v3" 22 23 "tangled.org/core/api/tangled" 24 "tangled.org/core/log" 25 "tangled.org/core/spindle/agentproto" 26 agentv1 "tangled.org/core/spindle/agentproto/gen" 27 "tangled.org/core/spindle/config" 28 "tangled.org/core/spindle/db" 29 "tangled.org/core/spindle/engine" 30 "tangled.org/core/spindle/models" 31 "tangled.org/core/spindle/secrets" 32) 33 34const ( 35 guestWorkDir = "/workspace/repo" 36 guestBasePATH = "/run/current-system/sw/bin:/nix/var/nix/profiles/default/bin:/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin" 37 guestDevShellEnvPath = "/run/spindle/devshell-env.sh" 38 activationStepAction = "activate-config" 39 agentAcceptTimeout = 2 * time.Minute 40 agentHandshakeTimeout = 30 * time.Second 41 cacheDrainTimeout = 5 * time.Minute 42 vmShutdownTimeout = 10 * time.Second 43 guestTimeoutGrace = 5 * time.Second 44) 45 46type cleanupFunc func(context.Context) error 47 48type Engine struct { 49 l *slog.Logger 50 cfg *config.Config 51 db *db.DB 52 agentMu sync.Mutex 53 agent *agentHub 54 scheduler *engine.ResourceScheduler[Resources] 55 cgroupParent *CgroupParent 56 57 cleanupMu sync.Mutex 58 cleanup map[string][]cleanupFunc 59} 60 61type Step struct { 62 name string 63 kind models.StepKind 64 command string 65 environment map[string]string 66 action string 67 config manifestConfig 68 configKey string 69} 70 71func (s Step) Name() string { return s.name } 72func (s Step) Command() string { return s.command } 73func (s Step) Kind() models.StepKind { return s.kind } 74 75func New(ctx context.Context, cfg *config.Config, d *db.DB) (*Engine, error) { 76 l := log.FromContext(ctx).With("component", "engine.microvm") 77 budget, max, agingThreshold := newVMBudgetConfig(cfg.MicroVMPipelines) 78 l.Info("initialized microVM workflow budget", "budget", budget.String(), "maxWorkflow", max.String(), "agingThreshold", agingThreshold) 79 80 var cgroupParent *CgroupParent 81 var err error 82 if cfg.MicroVMPipelines.EnableCgroups { 83 cgroupParent, err = initCgroupParent(cfg.MicroVMPipelines.CgroupParent, cfg.MicroVMPipelines.CgroupSupervisorMemoryMinMiB, l) 84 if err != nil { 85 return nil, err 86 } 87 } 88 89 return &Engine{ 90 l: l, 91 cfg: cfg, 92 db: d, 93 scheduler: engine.NewResourceScheduler(budget, max, agingThreshold), 94 cgroupParent: cgroupParent, 95 cleanup: make(map[string][]cleanupFunc), 96 }, nil 97} 98 99func (e *Engine) ensureAgentHub() (*agentHub, error) { 100 e.agentMu.Lock() 101 defer e.agentMu.Unlock() 102 103 if e.agent != nil { 104 return e.agent, nil 105 } 106 107 port := e.cfg.MicroVMPipelines.AgentPort 108 if port == 0 { 109 port = agentproto.DefaultPort 110 } 111 agent, err := newAgentHub(port, e.l) 112 if err != nil { 113 return nil, err 114 } 115 e.agent = agent 116 return agent, nil 117} 118 119func (e *Engine) InitWorkflow(twf tangled.Pipeline_Workflow, tpl tangled.Pipeline) (*models.Workflow, error) { 120 swf := &models.Workflow{} 121 var dwf manifestWorkflow 122 123 if err := engine.DescribeManifestError(twf.Raw, manifestWorkflow{}); err != nil { 124 return nil, err 125 } 126 if err := yaml.Unmarshal([]byte(twf.Raw), &dwf); err != nil { 127 return nil, err 128 } 129 130 for _, dstep := range dwf.Steps { 131 swf.Steps = append(swf.Steps, Step{ 132 name: dstep.Name, 133 kind: models.StepKindUser, 134 command: dstep.Command, 135 environment: dstep.Environment, 136 }) 137 } 138 swf.Name = twf.Name 139 swf.Environment = dwf.Environment 140 141 if tpl.TriggerMetadata != nil { 142 if clone := models.BuildCloneStep(twf, *tpl.TriggerMetadata, e.cfg.Server.Dev); clone.Command() != "" { 143 swf.Steps = append([]models.Step{clone}, swf.Steps...) 144 } 145 } 146 147 imageSpec, imageSpecPath, imageName, err := e.resolveImage(dwf.Image) 148 if err != nil { 149 return nil, err 150 } 151 configKey := "" 152 config := manifestConfig{ 153 Services: dwf.Services, 154 Virtualisation: dwf.Virtualisation, 155 Dependencies: dwf.Dependencies, 156 Registry: dwf.Registry, 157 } 158 if config.Enabled() { 159 if !imageSpec.SupportsConfigActivation() { 160 return nil, fmt.Errorf( 161 "microVM image %q is not a NixOS image: services, virtualisation, dependencies and registry workflow options require a NixOS image", 162 imageName, 163 ) 164 } 165 var err error 166 configKey, err = buildConfigKey(imageSpec, config) 167 if err != nil { 168 return nil, fmt.Errorf("build config key: %w", err) 169 } 170 activationStep := Step{ 171 name: "NixOS config activation", 172 kind: models.StepKindSystem, 173 command: "activate nixos config", 174 action: activationStepAction, 175 config: config, 176 configKey: configKey, 177 } 178 179 insertAt := 0 180 if len(swf.Steps) > 0 && swf.Steps[0].Kind() == models.StepKindSystem { 181 insertAt = 1 182 } 183 swf.Steps = append(swf.Steps, nil) 184 copy(swf.Steps[insertAt+1:], swf.Steps[insertAt:]) 185 swf.Steps[insertAt] = activationStep 186 } 187 188 cacheURLs, cacheKeys, err := workflowCaches(dwf.Caches) 189 if err != nil { 190 return nil, err 191 } 192 193 swf.Data = &workflowState{ 194 ImageSpec: imageSpec, 195 ImageSpecPath: imageSpecPath, 196 Config: config, 197 ConfigKey: configKey, 198 Image: imageName, 199 CacheReadURLs: cacheURLs, 200 CacheTrustedPublicKeys: cacheKeys, 201 NixOSToplevelCache: newNixOSToplevelCacheStore(e.db), 202 } 203 return swf, nil 204} 205 206func (e *Engine) SetupWorkflow(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, wfLogger models.WorkflowLogger) (err error) { 207 l := e.l.With("workflow", wid) 208 setupStep := Step{name: "microVM setup", kind: models.StepKindSystem} 209 210 wfLogger.ControlWriter(-1, setupStep, models.StepStatusStart).Write([]byte{0}) 211 defer wfLogger.ControlWriter(-1, setupStep, models.StepStatusEnd).Write([]byte{0}) 212 213 category := "Failed to setup VM" 214 defer func() { 215 if err != nil { 216 err = fmt.Errorf("%s:\n%w", category, err) 217 } 218 }() 219 220 state, ok := wf.Data.(*workflowState) 221 if !ok || state == nil { 222 return fmt.Errorf("workflow state is not initialized") 223 } 224 225 cid, err := AllocateCID() 226 if err != nil { 227 return err 228 } 229 agent, err := e.ensureAgentHub() 230 if err != nil { 231 return err 232 } 233 connCh, unregister, err := agent.expect(cid) 234 if err != nil { 235 return err 236 } 237 defer unregister() 238 239 workDirBase := e.cfg.MicroVMPipelines.OverlayDir 240 if workDirBase == "" { 241 workDirBase = os.TempDir() 242 } 243 workDir, err := os.MkdirTemp(workDirBase, "spindle-microvm-"+wid.String()+"-*") 244 if err != nil { 245 return fmt.Errorf("create workflow microVM directory: %w", err) 246 } 247 state.WorkDir = workDir 248 249 setupDone := false 250 defer func() { 251 if setupDone { 252 return 253 } 254 if detail := VMCrashLog(state.VM); detail != "" { 255 l.Error("microVM setup failed", "detail", detail) 256 } 257 if err := e.cleanupState(context.Background(), wid, state); err != nil { 258 l.Error("failed to cleanup failed setup", "error", err) 259 } 260 }() 261 262 upstreams, err := BuildCacheUpstreams(e.cfg.NixCache.ReadURLs, state.CacheReadURLs) 263 if err != nil { 264 return err 265 } 266 readCache, err := StartReadCacheProxy(ctx, cid, upstreams, l) 267 if err != nil { 268 return err 269 } 270 state.ReadCache = readCache 271 stagingDir := filepath.Join(workDir, "upload-cache") 272 uploadCache, err := StartUploadCacheProxy(ctx, cid, e.cfg.NixCache.UploadURL, upstreams, stagingDir, l) 273 if err != nil { 274 return err 275 } 276 state.UploadCache = uploadCache 277 dnsProxy, err := StartDNSProxy(ctx, cid, l) 278 if err != nil { 279 return err 280 } 281 state.DNSProxy = dnsProxy 282 283 port := e.cfg.MicroVMPipelines.AgentPort 284 if port == 0 { 285 port = agentproto.DefaultPort 286 } 287 state.ImageSpec.BootArgs = fmt.Sprintf("%s shuttle.vsock_port=%d", state.ImageSpec.BootArgs, port) 288 289 fmt.Fprintf(wfLogger.DataWriter(-1, "stdout"), "starting microVM image %s\n", state.Image) 290 l.Info("starting microVM workflow", "image", state.Image, "imageSpec", state.ImageSpecPath, "cid", cid, "workDir", workDir) 291 292 var vm VMHandle 293 vm, err = StartVM(ctx, VMConfig{ 294 Image: state.ImageSpec, 295 CID: cid, 296 EnableKVM: e.cfg.MicroVMPipelines.EnableKVM, 297 WorkDir: workDir, 298 Cgroup: e.cgroupLimits(wid, state.ImageSpec), 299 Dev: e.cfg.Server.Dev, 300 }, l) 301 if err != nil { 302 return err 303 } 304 state.VM = vm 305 306 category = "Failed to connect to agent" 307 308 acceptCtx, cancelAccept := context.WithTimeout(ctx, agentAcceptTimeout) 309 defer cancelAccept() 310 conn, err := waitAgentConn(acceptCtx, connCh) 311 if err != nil { 312 return err 313 } 314 315 agentSession := NewAgentSession(conn, l) 316 initCtx, cancelInit := context.WithTimeout(ctx, agentHandshakeTimeout) 317 defer cancelInit() 318 if err := agentSession.Init(initCtx, &agentv1.Init{ 319 JobId: wid.String(), 320 CacheTrustedPublicKeys: append(slices.Clone(e.cfg.NixCache.TrustedPublicKeys), state.CacheTrustedPublicKeys...), 321 CacheReadProxyPort: readCache.Port(), 322 CacheUploadProxyPort: uploadCache.Port(), 323 DnsProxyPort: dnsProxy.Port(), 324 }); err != nil { 325 _ = agentSession.Close() 326 return err 327 } 328 state.Agent = agentSession 329 wf.Data = state 330 331 e.registerCleanup(wid, func(ctx context.Context) error { 332 return e.cleanupState(ctx, wid, state) 333 }) 334 setupDone = true 335 336 fmt.Fprintf(wfLogger.DataWriter(-1, "stdout"), 337 "agent connected; serial log: %s\n", vm.Logs().Serial, 338 ) 339 return nil 340} 341 342func applyDepsSource(command string) string { 343 return fmt.Sprintf( 344 // check if it exists because not all images have this 345 `if [ -f %s ]; then . %s; export PATH="$PATH:%s"; fi; %s`, 346 guestDevShellEnvPath, guestDevShellEnvPath, guestBasePATH, command, 347 ) 348} 349 350func (e *Engine) RunStep(ctx context.Context, wid models.WorkflowId, w *models.Workflow, idx int, secrets []secrets.UnlockedSecret, wfLogger models.WorkflowLogger) error { 351 state, ok := w.Data.(*workflowState) 352 if !ok || state == nil || state.Agent == nil { 353 return fmt.Errorf("microVM workflow is not connected to agent") 354 } 355 356 stderr := wfLogger.DataWriter(idx, "stderr") 357 358 execCtx, vmExited, cancelWatch := watchVMExit(ctx, state.VM) 359 defer cancelWatch() 360 361 step := w.Steps[idx] 362 if s, ok := step.(Step); ok && s.action == activationStepAction { 363 err := e.activateConfig(execCtx, wid, state, s, wfLogger.DataWriter(idx, "stdout")) 364 return e.classifyStepError(ctx, wid, step, state, stderr, vmExited, "Failed to activate config", err) 365 } 366 env := []string{ 367 "HOME=/workspace", 368 "LOGNAME=" + guestWorkflowUser, 369 "PATH=" + guestBasePATH, 370 "USER=" + guestWorkflowUser, 371 } 372 for k, v := range w.Environment { 373 env = append(env, k+"="+v) 374 } 375 for _, s := range secrets { 376 env = append(env, s.Key+"="+s.Value) 377 } 378 if s, ok := step.(Step); ok { 379 for k, v := range s.environment { 380 env = append(env, k+"="+v) 381 } 382 } 383 384 stdout := wfLogger.DataWriter(idx, "stdout") 385 exitCode, err := state.Agent.Exec(execCtx, AgentExec{ 386 ID: fmt.Sprintf("%s-%d", wid.String(), idx), 387 ExecStart: &agentv1.ExecStart{ 388 Argv: []string{state.ImageSpec.Shell, "-lc", applyDepsSource(step.Command())}, 389 Env: env, 390 Cwd: guestWorkDir, 391 User: guestWorkflowUser, 392 // timeout not set here, Exec will fill it 393 }, 394 Stdout: stdout, 395 Stderr: stderr, 396 }) 397 if err != nil { 398 return e.classifyStepError(ctx, wid, step, state, stderr, vmExited, "User step error", err) 399 } 400 401 if exitCode != 0 { 402 e.l.Debug("step exited non-zero", "workflow", wid, "step", step.Name(), "exitCode", exitCode) 403 return fmt.Errorf("User step error: exited with code %d", exitCode) 404 } 405 return nil 406} 407 408// reads the vm serial logs so we report the tail of that as an error instead of 409// just "guest agent connection lost: EOF" 410func (e *Engine) classifyStepError(ctx context.Context, wid models.WorkflowId, step models.Step, state *workflowState, stderr io.Writer, vmExited *atomic.Bool, category string, err error) error { 411 if err == nil { 412 return nil 413 } 414 l := e.l.With("workflow", wid, "step", step.Name()) 415 416 if vmExited != nil && vmExited.Load() { 417 reason := "microVM exited unexpectedly" 418 oom := state.VM != nil && state.VM.OOMKilled() 419 if oom { 420 reason = "microVM killed by OOM (cgroup memory limit exceeded)" 421 } 422 if detail := VMCrashLog(state.VM); detail != "" { 423 fmt.Fprintf(stderr, "%s:\n%s\n", reason, detail) 424 l.Error(reason, "oom", oom, "detail", detail) 425 } else { 426 fmt.Fprintln(stderr, reason) 427 l.Error(reason, "oom", oom) 428 } 429 return fmt.Errorf("%s:\n%w", category, errors.New(reason+"; see workflow logs for serial output")) 430 } 431 432 if errors.Is(err, errGuestTimedOut) || ctx.Err() != nil { 433 l.Debug("step timed out", "guestReported", errors.Is(err, errGuestTimedOut)) 434 return engine.ErrTimedOut 435 } 436 437 // the agent connection dropped while qemu stayed up (eg. the guest kernel 438 // OOM-killed the agent or a guest panic), so surface serial logs, those 439 // will be more helpful. 440 var crashErr error 441 if detail := VMCrashLog(state.VM); detail != "" { 442 fmt.Fprintf(stderr, "step failed (%v):\n%s\n", err, detail) 443 l.Error("step failed", "error", err, "detail", detail) 444 if parsedErr, ok := ParseCrashLog(detail); ok { 445 crashErr = parsedErr 446 } else { 447 if strings.Contains(err.Error(), "guest exec error:") { 448 crashErr = err 449 } else { 450 crashErr = fmt.Errorf("guest agent connection lost: %w", err) 451 } 452 } 453 } else { 454 l.Error("step failed", "error", err) 455 crashErr = err 456 } 457 return fmt.Errorf("%s:\n%w", category, crashErr) 458} 459 460func (e *Engine) activateConfig(ctx context.Context, wid models.WorkflowId, state *workflowState, step Step, out io.Writer) error { 461 cfg := step.config 462 if !cfg.Enabled() { 463 return nil 464 } 465 466 configKey := step.configKey 467 if configKey == "" { 468 configKey = state.ConfigKey 469 } 470 471 userConfigJSON, err := json.Marshal(cfg) 472 if err != nil { 473 return fmt.Errorf("encode user config: %w", err) 474 } 475 476 var cachedToplevel string 477 if configKey != "" { 478 if record, ok, err := state.NixOSToplevelCache.Lookup(configKey); err != nil { 479 return err 480 } else if ok { 481 // todo(dawn): we should probably use gc roots to eliminate TOCTOU 482 // the spindle will have to manage the gc roots, and for remote we have to 483 // ssh in to the host and add / remove gc root. 484 // we need to have this check anyway since the only check http caches can 485 // use is this one, since we cant manage gc roots there... 486 if e.anyCacheHasPath(ctx, state, record.Toplevel) { 487 cachedToplevel = record.Toplevel 488 fmt.Fprintf(out, "realizing cached NixOS config %s\n", cachedToplevel) 489 } 490 } 491 } 492 if cachedToplevel == "" { 493 fmt.Fprintf(out, "building NixOS config from user config\n") 494 } 495 496 baseHash, err := BaseConfigHash(state.ImageSpec) 497 if err != nil { 498 return fmt.Errorf("calculate base config hash: %w", err) 499 } 500 501 result, err := state.Agent.ActivateConfig(ctx, fmt.Sprintf("%s-config", wid.String()), &agentv1.ActivateConfig{ 502 ConfigKey: configKey, 503 BaseConfigHash: baseHash, 504 UserConfig: string(userConfigJSON), 505 Toplevel: cachedToplevel, 506 }, out) 507 if err != nil { 508 return err 509 } 510 fmt.Fprintf(out, "activated NixOS config toplevel %s\n", result.Toplevel) 511 512 if cachedToplevel != "" || configKey == "" { 513 return nil 514 } 515 if e.cfg.NixCache.UploadURL == "" { 516 e.l.Warn("not committing config cache metadata: no upload URL configured", "workflow", wid, "configKey", configKey, "toplevel", result.Toplevel) 517 return nil 518 } 519 520 if err := e.drainNixCache(ctx, state); err != nil { 521 // a partial upload would leave the cache unable to realize this toplevel, 522 // so skip the metadata commit rather than poison it with an un-realizable 523 // key. the config still activated fine, so don't fail the workflow. 524 e.l.Warn("cache drain failed; skipping config cache metadata commit", "workflow", wid, "configKey", configKey, "toplevel", result.Toplevel, "error", err) 525 return nil 526 } 527 if err := state.NixOSToplevelCache.Commit(configKey, result.Toplevel); err != nil { 528 return err 529 } 530 fmt.Fprintf(out, "committed config cache metadata %s -> %s\n", configKey, result.Toplevel) 531 return nil 532} 533 534func (e *Engine) anyCacheHasPath(ctx context.Context, state *workflowState, storePath string) bool { 535 upstreams, err := BuildCacheUpstreams(e.cfg.NixCache.ReadURLs, state.CacheReadURLs) 536 if err != nil { 537 e.l.Warn("config cache check: build upstreams failed; treating as absent", "path", storePath, "error", err) 538 return false 539 } 540 if len(upstreams) == 0 { 541 return false 542 } 543 hash, _, err := parseStorePath(storePath) 544 if err != nil { 545 e.l.Warn("config cache check: invalid toplevel path; treating as absent", "path", storePath, "error", err) 546 return false 547 } 548 req, err := http.NewRequestWithContext(ctx, http.MethodHead, "http://upstream/"+hash+".narinfo", nil) 549 if err != nil { 550 e.l.Warn("config cache check: build request failed; treating as absent", "path", storePath, "error", err) 551 return false 552 } 553 resp, err := newNarinfoExistenceTransport(upstreams, e.l).RoundTrip(req) 554 if err != nil { 555 e.l.Warn("config cache check: narinfo probe failed; treating as absent", "path", storePath, "error", err) 556 return false 557 } 558 defer resp.Body.Close() 559 _, _ = io.Copy(io.Discard, resp.Body) 560 return resp.StatusCode == http.StatusOK 561} 562 563func (e *Engine) DestroyWorkflow(ctx context.Context, wid models.WorkflowId) error { 564 fns := e.drainCleanups(wid) 565 566 var cleanupErr error 567 for i := len(fns) - 1; i >= 0; i-- { 568 if err := fns[i](ctx); err != nil { 569 e.l.Error("failed to cleanup workflow resource", "workflowId", wid, "error", err) 570 cleanupErr = errors.Join(cleanupErr, err) 571 } 572 } 573 return cleanupErr 574} 575 576func (e *Engine) FinalizeWorkflow(ctx context.Context, wid models.WorkflowId, w *models.Workflow, wfLogger models.WorkflowLogger) error { 577 return nil 578} 579 580func (e *Engine) WorkflowTimeout() time.Duration { 581 d, err := time.ParseDuration(e.cfg.MicroVMPipelines.WorkflowTimeout) 582 if err != nil { 583 d = 5 * time.Minute 584 } 585 return d + guestTimeoutGrace 586} 587 588func (e *Engine) registerCleanup(wid models.WorkflowId, fn cleanupFunc) { 589 e.cleanupMu.Lock() 590 defer e.cleanupMu.Unlock() 591 key := wid.String() 592 e.cleanup[key] = append(e.cleanup[key], fn) 593} 594 595func (e *Engine) drainCleanups(wid models.WorkflowId) []cleanupFunc { 596 e.cleanupMu.Lock() 597 defer e.cleanupMu.Unlock() 598 key := wid.String() 599 fns := e.cleanup[key] 600 delete(e.cleanup, key) 601 return fns 602} 603 604func (e *Engine) cgroupLimits(wid models.WorkflowId, spec ImageSpec) CgroupLimits { 605 cfg := e.cfg.MicroVMPipelines 606 return CgroupLimits{ 607 Enabled: cfg.EnableCgroups, 608 Parent: e.cgroupParent, 609 Name: "workflow-" + wid.String(), 610 MemoryMaxMiB: resourcesForImage(spec).MemoryMiB, 611 SwapMaxMiB: cfg.CgroupSwapMaxMiB, 612 PidsMax: cfg.CgroupPidsMax, 613 } 614}