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