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