This repository has no description
1package nixery
2
3import (
4 "bufio"
5 "context"
6 "errors"
7 "fmt"
8 "io"
9 "log/slog"
10 "path"
11 "runtime"
12 "sync"
13 "time"
14
15 "github.com/docker/docker/api/types/container"
16 "github.com/docker/docker/api/types/image"
17 "github.com/docker/docker/api/types/mount"
18 "github.com/docker/docker/api/types/network"
19 "github.com/docker/docker/client"
20 "github.com/docker/docker/pkg/stdcopy"
21 "gopkg.in/yaml.v3"
22 "tangled.org/core/api/tangled"
23 "tangled.org/core/log"
24 "tangled.org/core/spindle/config"
25 "tangled.org/core/spindle/engine"
26 "tangled.org/core/spindle/models"
27 "tangled.org/core/spindle/secrets"
28)
29
30const (
31 workspaceDir = "/tangled/workspace"
32 homeDir = "/tangled/home"
33)
34
35type cleanupFunc func(context.Context) error
36
37type Engine struct {
38 dockerMu sync.Mutex
39 docker client.APIClient
40 l *slog.Logger
41 cfg *config.Config
42
43 slotter engine.WorkflowSlotter
44
45 cleanupMu sync.Mutex
46 cleanup map[string][]cleanupFunc
47}
48
49type Step struct {
50 name string
51 kind models.StepKind
52 command string
53 environment map[string]string
54}
55
56func (s Step) Name() string {
57 return s.name
58}
59
60func (s Step) Command() string {
61 return s.command
62}
63
64func (s Step) Kind() models.StepKind {
65 return s.kind
66}
67
68// setupSteps get added to start of Steps
69type setupSteps []models.Step
70
71// addStep adds a step to the beginning of the workflow's steps.
72func (ss *setupSteps) addStep(step models.Step) {
73 *ss = append(*ss, step)
74}
75
76type addlFields struct {
77 image string
78 container string
79 mounts []mount.Mount
80}
81
82func (e *Engine) InitWorkflow(twf tangled.Pipeline_Workflow, tpl tangled.Pipeline) (*models.Workflow, error) {
83 swf := &models.Workflow{}
84 addl := addlFields{}
85
86 dwf := &struct {
87 Steps []struct {
88 Command string `yaml:"command"`
89 Name string `yaml:"name"`
90 Environment map[string]string `yaml:"environment"`
91 } `yaml:"steps"`
92 Dependencies map[string][]string `yaml:"dependencies"`
93 Environment map[string]string `yaml:"environment"`
94 Cache []models.CacheEntry `yaml:"cache"`
95 }{}
96 if err := engine.DescribeManifestError(twf.Raw, dwf); err != nil {
97 return nil, err
98 }
99 if err := yaml.Unmarshal([]byte(twf.Raw), &dwf); err != nil {
100 return nil, err
101 }
102
103 for _, dstep := range dwf.Steps {
104 sstep := Step{}
105 sstep.environment = dstep.Environment
106 sstep.command = dstep.Command
107 sstep.name = dstep.Name
108 sstep.kind = models.StepKindUser
109 swf.Steps = append(swf.Steps, sstep)
110 }
111 swf.Name = twf.Name
112 swf.Environment = dwf.Environment
113 for _, entry := range dwf.Cache {
114 if err := entry.Validate(); err != nil {
115 return nil, err
116 }
117 }
118 swf.Caches = dwf.Cache
119 addl.image = workflowImage(dwf.Dependencies, e.cfg.NixeryPipelines.Nixery)
120
121 if sock := e.cfg.Server.DockerSocket; sock != "" {
122 addl.mounts = append(addl.mounts, mount.Mount{
123 Type: mount.TypeBind,
124 Source: sock,
125 Target: sock,
126 ReadOnly: false,
127 })
128 }
129 setup := &setupSteps{}
130
131 setup.addStep(nixConfStep())
132 setup.addStep(models.BuildCloneStep(twf, *tpl.TriggerMetadata, e.cfg.Server.Dev))
133 // this step could be empty
134 if s := dependencyStep(dwf.Dependencies); s != nil {
135 setup.addStep(*s)
136 }
137
138 // append setup steps in order to the start of workflow steps
139 swf.Steps = append(*setup, swf.Steps...)
140 swf.Data = addl
141
142 return swf, nil
143}
144
145func (e *Engine) WorkflowTimeout() time.Duration {
146 workflowTimeoutStr := e.cfg.NixeryPipelines.WorkflowTimeout
147 workflowTimeout, err := time.ParseDuration(workflowTimeoutStr)
148 if err != nil {
149 e.l.Error("failed to parse workflow timeout", "error", err, "timeout", workflowTimeoutStr)
150 workflowTimeout = 5 * time.Minute
151 }
152
153 return workflowTimeout
154}
155
156func workflowImage(deps map[string][]string, nixery string) string {
157 var dependencies string
158 for reg, ds := range deps {
159 if reg == "nixpkgs" {
160 dependencies = path.Join(ds...)
161 }
162 }
163
164 // load defaults from somewhere else
165 dependencies = path.Join(dependencies, "bash", "git", "coreutils", "gnutar", "zstd", "nix")
166
167 if runtime.GOARCH == "arm64" {
168 dependencies = path.Join("arm64", dependencies)
169 }
170
171 return path.Join(nixery, dependencies)
172}
173
174func New(ctx context.Context, cfg *config.Config) (*Engine, error) {
175 l := log.FromContext(ctx).With("component", "spindle")
176
177 e := &Engine{
178 l: l,
179 cfg: cfg,
180 slotter: engine.NewSemaphoreSlotter(cfg.NixeryPipelines.MaxConcurrentWorkflows),
181 }
182
183 e.cleanup = make(map[string][]cleanupFunc)
184
185 return e, nil
186}
187
188func (e *Engine) ensureDocker() (client.APIClient, error) {
189 e.dockerMu.Lock()
190 defer e.dockerMu.Unlock()
191
192 if e.docker != nil {
193 return e.docker, nil
194 }
195
196 dcli, err := client.NewClientWithOpts(client.FromEnv, client.WithAPIVersionNegotiation())
197 if err != nil {
198 return nil, err
199 }
200 e.docker = dcli
201 return dcli, nil
202}
203
204func (e *Engine) AcquireWorkflowSlot(
205 ctx context.Context,
206 wid models.WorkflowId,
207 wf *models.Workflow,
208) (engine.WorkflowSlot, error) {
209 if e.slotter == nil {
210 return engine.NoopSlot{}, nil
211 }
212
213 return e.slotter.AcquireWorkflowSlot(ctx, wid, wf)
214}
215
216func (e *Engine) SetupWorkflow(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, wfLogger models.WorkflowLogger) (err error) {
217 /// -------------------------INITIAL SETUP------------------------------------------
218 l := e.l.With("workflow", wid)
219 l.Info("setting up workflow")
220
221 setupStep := Step{
222 name: "Pull image from Nixery",
223 kind: models.StepKindSystem,
224 }
225 setupStepIdx := -1
226
227 wfLogger.ControlWriter(setupStepIdx, setupStep, models.StepStatusStart).Write([]byte{0})
228 defer wfLogger.ControlWriter(setupStepIdx, setupStep, models.StepStatusEnd).Write([]byte{0})
229
230 defer func() {
231 if err != nil {
232 err = fmt.Errorf("Failed to setup container:\n%w", err)
233 }
234 }()
235
236 if _, err := e.ensureDocker(); err != nil {
237 return err
238 }
239
240 /// -------------------------NETWORK CREATION---------------------------------------
241 _, err = e.docker.NetworkCreate(ctx, networkName(wid), network.CreateOptions{
242 Driver: "bridge",
243 })
244 if err != nil {
245 return err
246 }
247
248 e.registerCleanup(wid, func(ctx context.Context) error {
249 if err := e.docker.NetworkRemove(ctx, networkName(wid)); err != nil {
250 return fmt.Errorf("removing network: %w", err)
251 }
252 return nil
253 })
254
255 /// -------------------------IMAGE PULL---------------------------------------------
256 addl := wf.Data.(addlFields)
257 l.Info("pulling image", "image", addl.image)
258 fmt.Fprintf(
259 wfLogger.DataWriter(setupStepIdx, "stdout"),
260 "Pulling image: %s",
261 addl.image,
262 )
263
264 reader, err := e.docker.ImagePull(ctx, addl.image, image.PullOptions{})
265 if err != nil {
266 l.Error("pipeline image pull failed!", "error", err.Error())
267 fmt.Fprintf(wfLogger.DataWriter(setupStepIdx, "stderr"), "image pull failed: %s", err)
268 return fmt.Errorf("pulling image: %w", err)
269 }
270 defer reader.Close()
271
272 scanner := bufio.NewScanner(reader)
273 for scanner.Scan() {
274 line := scanner.Text()
275 wfLogger.DataWriter(setupStepIdx, "stdout").Write([]byte(line))
276 l.Info("image pull progress", "stdout", line)
277 }
278
279 /// -------------------------CONTAINER CREATION-------------------------------------
280 l.Info("creating container")
281 wfLogger.DataWriter(setupStepIdx, "stdout").Write([]byte("Creating container..."))
282
283 extraHosts := []string{"host.docker.internal:host-gateway"}
284 for _, h := range e.cfg.Server.DevExtraHosts {
285 extraHosts = append(extraHosts, h+":host-gateway")
286 }
287
288 resp, err := e.docker.ContainerCreate(ctx, &container.Config{
289 Image: addl.image,
290 Cmd: []string{"cat"},
291 OpenStdin: true, // so cat stays alive :3
292 Tty: false,
293 Hostname: "spindle",
294 WorkingDir: workspaceDir,
295 Labels: map[string]string{
296 "sh.tangled.pipeline/workflow_id": wid.String(),
297 },
298 // TODO(winter): investigate whether environment variables passed here
299 // get propagated to ContainerExec processes
300 }, &container.HostConfig{
301 Mounts: append([]mount.Mount{
302 {
303 Type: mount.TypeTmpfs,
304 Target: "/tmp",
305 ReadOnly: false,
306 TmpfsOptions: &mount.TmpfsOptions{
307 Mode: 0o1777, // world-writable sticky bit
308 Options: [][]string{
309 {"exec"},
310 },
311 },
312 },
313 }, addl.mounts...),
314 ReadonlyRootfs: false,
315 CapDrop: []string{"ALL"},
316 CapAdd: []string{"CAP_DAC_OVERRIDE", "CAP_CHOWN", "CAP_FOWNER", "CAP_SETUID", "CAP_SETGID"},
317 SecurityOpt: []string{"no-new-privileges"},
318 ExtraHosts: extraHosts,
319 Resources: container.Resources{
320 Memory: e.cfg.NixeryPipelines.MaxJobMemoryMB * 1024 * 1024,
321 },
322 }, nil, nil, "")
323 if err != nil {
324 fmt.Fprintf(
325 wfLogger.DataWriter(setupStepIdx, "stderr"),
326 "container creation failed: %s",
327 err,
328 )
329 return fmt.Errorf("creating container: %w", err)
330 }
331
332 e.registerCleanup(wid, func(ctx context.Context) error {
333 if err := e.docker.ContainerStop(ctx, resp.ID, container.StopOptions{}); err != nil {
334 return fmt.Errorf("stopping container: %w", err)
335 }
336
337 err := e.docker.ContainerRemove(ctx, resp.ID, container.RemoveOptions{
338 RemoveVolumes: true,
339 RemoveLinks: false,
340 Force: false,
341 })
342 if err != nil {
343 return fmt.Errorf("removing container: %w", err)
344 }
345
346 return nil
347 })
348
349 /// -------------------------CONTAINER START----------------------------------------
350 wfLogger.DataWriter(setupStepIdx, "stdout").Write([]byte("Starting container..."))
351 if err := e.docker.ContainerStart(ctx, resp.ID, container.StartOptions{}); err != nil {
352 return fmt.Errorf("starting container: %w", err)
353 }
354
355 mkExecResp, err := e.docker.ContainerExecCreate(ctx, resp.ID, container.ExecOptions{
356 Cmd: []string{"mkdir", "-p", workspaceDir, homeDir},
357 AttachStdout: true, // NOTE(winter): pretty sure this will make it so that when stdout read is done below, mkdir is done. maybe??
358 AttachStderr: true, // for good measure, backed up by docker/cli ("If -d is not set, attach to everything by default")
359 })
360 if err != nil {
361 return err
362 }
363
364 // This actually *starts* the command. Thanks, Docker!
365 execResp, err := e.docker.ContainerExecAttach(ctx, mkExecResp.ID, container.ExecAttachOptions{})
366 if err != nil {
367 return err
368 }
369 defer execResp.Close()
370
371 // This is apparently best way to wait for the command to complete.
372 _, err = io.ReadAll(execResp.Reader)
373 if err != nil {
374 return err
375 }
376
377 /// -----------------------------------FINISH---------------------------------------
378 execInspectResp, err := e.docker.ContainerExecInspect(ctx, mkExecResp.ID)
379 if err != nil {
380 return err
381 }
382
383 if execInspectResp.ExitCode != 0 {
384 return fmt.Errorf("mkdir exited with exit code %d", execInspectResp.ExitCode)
385 } else if execInspectResp.Running {
386 return errors.New("mkdir is somehow still running??")
387 }
388
389 addl.container = resp.ID
390 wf.Data = addl
391
392 return nil
393}
394
395func (e *Engine) RunStep(ctx context.Context, wid models.WorkflowId, w *models.Workflow, idx int, secrets []secrets.UnlockedSecret, wfLogger models.WorkflowLogger) error {
396 addl := w.Data.(addlFields)
397 workflowEnvs := ConstructEnvs(w.Environment)
398 // TODO(winter): should SetupWorkflow also have secret access?
399 // IMO yes, but probably worth thinking on.
400 for _, s := range secrets {
401 workflowEnvs.AddEnv(s.Key, s.Value)
402 }
403
404 step := w.Steps[idx]
405
406 select {
407 case <-ctx.Done():
408 return ctx.Err()
409 default:
410 }
411
412 envs := append(EnvVars(nil), workflowEnvs...)
413 if nixStep, ok := step.(Step); ok {
414 for k, v := range nixStep.environment {
415 envs.AddEnv(k, v)
416 }
417 }
418
419 envs = append(envs, e.baseEnv()...)
420 if sock := e.cfg.Server.DockerSocket; sock != "" {
421 envs.AddEnv("DOCKER_HOST", fmt.Sprintf("unix://%s", sock))
422 }
423
424 mkExecResp, err := e.docker.ContainerExecCreate(ctx, addl.container, container.ExecOptions{
425 Cmd: []string{"bash", "-c", step.Command()},
426 AttachStdout: true,
427 AttachStderr: true,
428 Env: envs,
429 })
430 if err != nil {
431 return fmt.Errorf("User step error:\ncreating exec: %w", err)
432 }
433
434 // start tailing logs in background
435 tailDone := make(chan error, 1)
436 go func() {
437 tailDone <- e.tailStep(ctx, wfLogger, mkExecResp.ID, idx)
438 }()
439
440 select {
441 case <-tailDone:
442
443 case <-ctx.Done():
444 // cleanup will be handled by DestroyWorkflow, since
445 // Docker doesn't provide an API to kill an exec run
446 // (sure, we could grab the PID and kill it ourselves,
447 // but that's wasted effort)
448 e.l.Warn("step timed out", "step", step.Name())
449
450 <-tailDone
451
452 return engine.ErrTimedOut
453 }
454
455 select {
456 case <-ctx.Done():
457 return ctx.Err()
458 default:
459 }
460
461 execInspectResp, err := e.docker.ContainerExecInspect(ctx, mkExecResp.ID)
462 if err != nil {
463 return fmt.Errorf("User step error:\n%w", err)
464 }
465
466 if execInspectResp.ExitCode != 0 {
467 inspectResp, err := e.docker.ContainerInspect(ctx, addl.container)
468 if err != nil {
469 return fmt.Errorf("User step error:\n%w", err)
470 }
471
472 e.l.Error("workflow failed!", "workflow_id", wid.String(), "exit_code", execInspectResp.ExitCode, "oom_killed", inspectResp.State.OOMKilled)
473
474 if inspectResp.State.OOMKilled {
475 return fmt.Errorf("User step error:\n%w", ErrOOMKilled)
476 }
477 return fmt.Errorf("User step error: exited with code %d", execInspectResp.ExitCode)
478 }
479
480 return nil
481}
482
483func (e *Engine) tailStep(ctx context.Context, wfLogger models.WorkflowLogger, execID string, stepIdx int) error {
484 if wfLogger == nil {
485 return nil
486 }
487
488 // This actually *starts* the command. Thanks, Docker!
489 logs, err := e.docker.ContainerExecAttach(ctx, execID, container.ExecAttachOptions{})
490 if err != nil {
491 return err
492 }
493 defer logs.Close()
494
495 _, err = stdcopy.StdCopy(
496 wfLogger.DataWriter(stepIdx, "stdout"),
497 wfLogger.DataWriter(stepIdx, "stderr"),
498 logs.Reader,
499 )
500 if err != nil && err != io.EOF && !errors.Is(err, context.DeadlineExceeded) {
501 return fmt.Errorf("failed to copy logs: %w", err)
502 }
503
504 return nil
505}
506
507func (e *Engine) DestroyWorkflow(ctx context.Context, wid models.WorkflowId) error {
508 fns := e.drainCleanups(wid)
509
510 for _, fn := range fns {
511 if err := fn(ctx); err != nil {
512 e.l.Error("failed to cleanup workflow resource", "workflowId", wid, "error", err)
513 }
514 }
515 return nil
516}
517
518func (e *Engine) registerCleanup(wid models.WorkflowId, fn cleanupFunc) {
519 e.cleanupMu.Lock()
520 defer e.cleanupMu.Unlock()
521
522 key := wid.String()
523 e.cleanup[key] = append(e.cleanup[key], fn)
524}
525
526func (e *Engine) drainCleanups(wid models.WorkflowId) []cleanupFunc {
527 e.cleanupMu.Lock()
528 key := wid.String()
529
530 fns := e.cleanup[key]
531 delete(e.cleanup, key)
532 e.cleanupMu.Unlock()
533
534 return fns
535}
536
537func networkName(wid models.WorkflowId) string {
538 return fmt.Sprintf("workflow-network-%s", wid)
539}