package engine import ( "bufio" "bytes" "context" "fmt" "io" "strings" "tangled.org/core/spindle/models" "tangled.org/core/spindle/storage" ) // cache log steps live below the setup step (-1) const ( CacheRestoreStepIdx = -2 CacheSaveStepIdx = -3 ) type cacheStep struct { name string command string } func (s cacheStep) Name() string { return s.name } func (s cacheStep) Command() string { return s.command } func (s cacheStep) Kind() models.StepKind { return models.StepKindSystem } var ( CacheRestoreStep models.Step = cacheStep{name: "restore cache", command: "restore cached paths"} CacheSaveStep models.Step = cacheStep{name: "save cache", command: "persist changed paths"} ) type CacheRunner interface { RestoreCache(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, store storage.Storage, caches []models.CacheBinding, wfLogger models.WorkflowLogger) error SaveCache(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, store storage.Storage, caches []models.CacheBinding, wfLogger models.WorkflowLogger) error } const CacheExitNoPaths = 42 // avoids storing an empty archive when zstd is missing const CacheExitNoCompressor = 43 // run tar from / so absolute paths survive extraction // tar exits nonzero on missing paths, so only existing ones reach it func CacheSaveScript(paths []string, workspaceRoot string, compressionLevel int) string { trimmed := make([]string, 0, len(paths)) for _, p := range paths { if !strings.HasPrefix(p, "/") { p = workspaceRoot + "/" + p } trimmed = append(trimmed, shellQuote(strings.TrimPrefix(p, "/"))) } tail := fmt.Sprintf(`tar -cf - -C / -- "$@" | %s`, CacheCompressCmd(compressionLevel)) return fmt.Sprintf(`set -o pipefail command -v zstd >/dev/null 2>&1 || { echo "zstd not found in image; cannot save cache" >&2; exit %d; } set -- for p in %s; do [ -e "/$p" ] && set -- "$@" "$p"; done if [ $# -eq 0 ]; then echo "no cache paths exist; skipping" >&2; exit %d; fi %s`, CacheExitNoCompressor, strings.Join(trimmed, " "), CacheExitNoPaths, tail) } func shellQuote(value string) string { return "'" + strings.ReplaceAll(value, "'", "'\"'\"'") + "'" } func CacheCompressCmd(level int) string { if level == 0 { return "zstd -T0 -5" } return fmt.Sprintf("zstd -T0 -%d", level) } // old entries might still be gzip, so detect them instead of trusting config func CacheDecompressCmd(br *bufio.Reader) string { head, _ := br.Peek(4) if bytes.HasPrefix(head, []byte{0x1f, 0x8b}) { return "gzip -dc" } return "zstd -dc" } // Put keeps draining the pipe after it returns, so the guest writer never // blocks on a full pipe. type CacheUpload struct { Writer *io.PipeWriter done chan error } func NewCacheUpload(ctx context.Context, store storage.Storage, key string) *CacheUpload { pr, pw := io.Pipe() u := &CacheUpload{Writer: pw, done: make(chan error, 1)} go func() { err := store.Put(ctx, key, pr) _, _ = io.Copy(io.Discard, pr) u.done <- err }() return u } func (u *CacheUpload) Abort(err error) { u.Writer.CloseWithError(err) <-u.done } // storage treats EOF as a complete archive, so only a clean exec may Finish func (u *CacheUpload) Finish() error { u.Writer.Close() return <-u.done }