Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168package nixery
import ( "bufio" "context" "fmt" "io"
"github.com/docker/docker/api/types" "github.com/docker/docker/api/types/container" "github.com/docker/docker/pkg/stdcopy"
"tangled.org/core/spindle/engine" "tangled.org/core/spindle/models" "tangled.org/core/spindle/storage")
func (e *Engine) baseEnv() EnvVars { envs := EnvVars{} envs.AddEnv("HOME", homeDir) envs.AddEnv("PATH", fmt.Sprintf("%s/.nix-profile/bin:/nix/var/nix/profiles/default/bin:/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin", homeDir)) return envs}
func (e *Engine) execAttached(ctx context.Context, containerID string, opts container.ExecOptions) (string, types.HijackedResponse, error) { execResp, err := e.docker.ContainerExecCreate(ctx, containerID, opts) if err != nil { return "", types.HijackedResponse{}, fmt.Errorf("create exec: %w", err) } attach, err := e.docker.ContainerExecAttach(ctx, execResp.ID, container.ExecAttachOptions{}) if err != nil { return "", types.HijackedResponse{}, fmt.Errorf("attach exec: %w", err) } return execResp.ID, attach, nil}
func (e *Engine) containerID(wf *models.Workflow) (string, error) { addl, ok := wf.Data.(addlFields) if !ok || addl.container == "" { return "", fmt.Errorf("nixery workflow has no container") } return addl.container, nil}
func (e *Engine) RestoreCache(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, store storage.Storage, caches []models.CacheBinding, wfLogger models.WorkflowLogger) error { containerID, err := e.containerID(wf) if err != nil { return err }
out := wfLogger.DataWriter(engine.CacheRestoreStepIdx, "stdout") for _, entry := range caches { if err := ctx.Err(); err != nil { return err } if entry.RestoreKey == "" { fmt.Fprintf(out, "cache %q: miss\n", entry.Key) continue } if entry.RestoreName != "" { fmt.Fprintf(out, "cache %q: restoring from %q\n", entry.Key, entry.RestoreName) }
rc, err := store.Get(ctx, entry.RestoreKey) if err != nil { fmt.Fprintf(out, "cache %q: fetch failed: %v\n", entry.Key, err) continue }
br := bufio.NewReader(rc) execID, attach, err := e.execAttached(ctx, containerID, container.ExecOptions{ Cmd: []string{"bash", "-c", engine.CacheDecompressCmd(br) + " | tar -x -C /"}, Env: e.baseEnv(), AttachStdin: true, AttachStdout: true, AttachStderr: true, }) if err != nil { rc.Close() return fmt.Errorf("restore cache %q: %w", entry.Key, err) }
// drain this now or tar can block on stderr before reading stdin copyDone := make(chan error, 1) go func() { _, err := io.Copy(attach.Conn, br) _ = attach.CloseWrite() copyDone <- err }() _, _ = stdcopy.StdCopy(out, out, attach.Reader) copyErr := <-copyDone rc.Close() attach.Close() if copyErr != nil { return fmt.Errorf("restore cache %q: stream archive: %w", entry.Key, copyErr) }
inspect, err := e.docker.ContainerExecInspect(ctx, execID) if err != nil { return fmt.Errorf("restore cache %q: %w", entry.Key, err) } if inspect.ExitCode != 0 { fmt.Fprintf(out, "cache %q: extract failed (exit %d)\n", entry.Key, inspect.ExitCode) continue } fmt.Fprintf(out, "cache %q: restored\n", entry.Key) } return nil}
func (e *Engine) SaveCache(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, store storage.Storage, caches []models.CacheBinding, wfLogger models.WorkflowLogger) error { containerID, err := e.containerID(wf) if err != nil { return err }
out := wfLogger.DataWriter(engine.CacheSaveStepIdx, "stdout") for _, entry := range caches { if err := ctx.Err(); err != nil { return err }
script := engine.CacheSaveScript(entry.Paths, workspaceDir, entry.CompressionLevel)
execID, attach, err := e.execAttached(ctx, containerID, container.ExecOptions{ Cmd: []string{"bash", "-c", script}, Env: e.baseEnv(), AttachStdout: true, AttachStderr: true, }) if err != nil { return fmt.Errorf("save cache %q: %w", entry.Key, err) }
up := engine.NewCacheUpload(ctx, store, entry.SaveKey)
// StdCopy only returns once the archive is fully written _, copyErr := stdcopy.StdCopy(up.Writer, out, attach.Reader) attach.Close() inspect, inspectErr := e.docker.ContainerExecInspect(ctx, execID)
switch { case inspectErr == nil && inspect.ExitCode == engine.CacheExitNoPaths: up.Abort(fmt.Errorf("no cache paths")) fmt.Fprintf(out, "cache %q: nothing to save\n", entry.Key) continue case inspectErr == nil && inspect.ExitCode == engine.CacheExitNoCompressor: up.Abort(fmt.Errorf("zstd not available")) fmt.Fprintf(out, "cache %q: zstd not available in image; skipping\n", entry.Key) continue case copyErr != nil: up.Abort(copyErr) return fmt.Errorf("save cache %q: stream archive: %w", entry.Key, copyErr) case inspectErr != nil: up.Abort(inspectErr) return fmt.Errorf("save cache %q: %w", entry.Key, inspectErr) case inspect.ExitCode != 0: up.Abort(fmt.Errorf("exited %d", inspect.ExitCode)) return fmt.Errorf("save cache %q: tar exited %d", entry.Key, inspect.ExitCode) } if err := up.Finish(); err != nil { return fmt.Errorf("save cache %q: %w", entry.Key, err) } fmt.Fprintf(out, "cache %q: saved\n", entry.Key) } return nil}