Monorepo for Tangled
tangled.org
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}