This repository has no description
0

Configure Feed

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

core / spindle / engines / nixery / engine.go
14 kB 535 lines
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}