package microvm import ( "bufio" "context" "fmt" "io" agentv1 "tangled.org/core/spindle/agentproto/gen" "tangled.org/core/spindle/engine" "tangled.org/core/spindle/models" "tangled.org/core/spindle/storage" ) func (e *Engine) RestoreCache(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, store storage.Storage, caches []engine.ResolvedCache, wfLogger models.WorkflowLogger) error { state, ok := wf.Data.(*workflowState) if !ok || state == nil || state.Agent == nil { return fmt.Errorf("microVM workflow is not connected to agent") } 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) decompress := engine.CacheDecompressCmd(br) var restored int64 exit, err := state.Agent.Exec(ctx, AgentExec{ ID: fmt.Sprintf("%s-cache-restore", wid.String()), ExecStart: cacheExecStart(state, fmt.Sprintf("set -o pipefail\n%s | tar -x -C /", decompress)), Stdin: &countingReader{r: br, n: &restored}, Stderr: out, }) rc.Close() if err != nil { fmt.Fprintf(out, "cache %q: restore failed: %v\n", entry.Key, err) continue } if exit != 0 { fmt.Fprintf(out, "cache %q: restore failed: guest exited %d\n", entry.Key, exit) continue } fmt.Fprintf(out, "cache %q: restored %d bytes\n", entry.Key, restored) } return nil } func (e *Engine) SaveCache(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, store storage.Storage, caches []engine.ResolvedCache, wfLogger models.WorkflowLogger) error { state, ok := wf.Data.(*workflowState) if !ok || state == nil || state.Agent == nil { return fmt.Errorf("microVM workflow is not connected to agent") } out := wfLogger.DataWriter(engine.CacheSaveStepIdx, "stdout") for _, entry := range caches { if err := ctx.Err(); err != nil { return err } script := engine.CacheSaveScript(entry.Paths, guestWorkDir, entry.CompressionLevel) up := engine.NewCacheUpload(ctx, store, entry.SaveKey) exit, execErr := state.Agent.Exec(ctx, AgentExec{ ID: fmt.Sprintf("%s-cache-save", wid.String()), ExecStart: cacheExecStart(state, script), Stdout: up.Writer, Stderr: out, }) switch { case exit == engine.CacheExitNoPaths: up.Abort(fmt.Errorf("guest exited %d", exit)) fmt.Fprintf(out, "cache %q: nothing to save\n", entry.Key) continue case exit == engine.CacheExitNoCompressor: up.Abort(fmt.Errorf("guest exited %d", exit)) fmt.Fprintf(out, "cache %q: zstd not available in image; skipping\n", entry.Key) continue case execErr != nil: up.Abort(execErr) return fmt.Errorf("save cache %q: %w", entry.Key, execErr) case exit != 0: up.Abort(fmt.Errorf("guest exited %d", exit)) return fmt.Errorf("save cache %q: save script exited %d", entry.Key, exit) } 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 } func cacheExecStart(state *workflowState, script string) *agentv1.ExecStart { return &agentv1.ExecStart{ Argv: []string{state.ImageSpec.Shell, "-c", script}, Env: guestBaseEnv(), User: guestWorkflowUser, } } type countingReader struct { r io.Reader n *int64 } func (c *countingReader) Read(p []byte) (int, error) { n, err := c.r.Read(p) *c.n += int64(n) return n, err }