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 519 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 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}