This repository has no description
0

Configure Feed

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

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