This repository has no description
0

Configure Feed

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

core / spindle / engines / microvm / cache.go
3.8 kB 125 lines
1package microvm 2 3import ( 4 "bufio" 5 "context" 6 "fmt" 7 "io" 8 9 agentv1 "tangled.org/core/spindle/agentproto/gen" 10 "tangled.org/core/spindle/engine" 11 "tangled.org/core/spindle/models" 12 "tangled.org/core/spindle/storage" 13) 14 15func (e *Engine) RestoreCache(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, store storage.Storage, caches []engine.ResolvedCache, wfLogger models.WorkflowLogger) error { 16 state, ok := wf.Data.(*workflowState) 17 if !ok || state == nil || state.Agent == nil { 18 return fmt.Errorf("microVM workflow is not connected to agent") 19 } 20 21 out := wfLogger.DataWriter(engine.CacheRestoreStepIdx, "stdout") 22 for _, entry := range caches { 23 if err := ctx.Err(); err != nil { 24 return err 25 } 26 if entry.RestoreKey == "" { 27 fmt.Fprintf(out, "cache %q: miss\n", entry.Key) 28 continue 29 } 30 if entry.RestoreName != "" { 31 fmt.Fprintf(out, "cache %q: restoring from %q\n", entry.Key, entry.RestoreName) 32 } 33 34 rc, err := store.Get(ctx, entry.RestoreKey) 35 if err != nil { 36 fmt.Fprintf(out, "cache %q: fetch failed: %v\n", entry.Key, err) 37 continue 38 } 39 br := bufio.NewReader(rc) 40 decompress := engine.CacheDecompressCmd(br) 41 42 var restored int64 43 exit, err := state.Agent.Exec(ctx, AgentExec{ 44 ID: fmt.Sprintf("%s-cache-restore", wid.String()), 45 ExecStart: cacheExecStart(state, fmt.Sprintf("set -o pipefail\n%s | tar -x -C /", decompress)), 46 Stdin: &countingReader{r: br, n: &restored}, 47 Stderr: out, 48 }) 49 rc.Close() 50 if err != nil { 51 fmt.Fprintf(out, "cache %q: restore failed: %v\n", entry.Key, err) 52 continue 53 } 54 if exit != 0 { 55 fmt.Fprintf(out, "cache %q: restore failed: guest exited %d\n", entry.Key, exit) 56 continue 57 } 58 fmt.Fprintf(out, "cache %q: restored %d bytes\n", entry.Key, restored) 59 } 60 return nil 61} 62 63func (e *Engine) SaveCache(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, store storage.Storage, caches []engine.ResolvedCache, wfLogger models.WorkflowLogger) error { 64 state, ok := wf.Data.(*workflowState) 65 if !ok || state == nil || state.Agent == nil { 66 return fmt.Errorf("microVM workflow is not connected to agent") 67 } 68 69 out := wfLogger.DataWriter(engine.CacheSaveStepIdx, "stdout") 70 for _, entry := range caches { 71 if err := ctx.Err(); err != nil { 72 return err 73 } 74 75 script := engine.CacheSaveScript(entry.Paths, guestWorkDir, entry.CompressionLevel) 76 77 up := engine.NewCacheUpload(ctx, store, entry.SaveKey) 78 exit, execErr := state.Agent.Exec(ctx, AgentExec{ 79 ID: fmt.Sprintf("%s-cache-save", wid.String()), 80 ExecStart: cacheExecStart(state, script), 81 Stdout: up.Writer, 82 Stderr: out, 83 }) 84 switch { 85 case exit == engine.CacheExitNoPaths: 86 up.Abort(fmt.Errorf("guest exited %d", exit)) 87 fmt.Fprintf(out, "cache %q: nothing to save\n", entry.Key) 88 continue 89 case exit == engine.CacheExitNoCompressor: 90 up.Abort(fmt.Errorf("guest exited %d", exit)) 91 fmt.Fprintf(out, "cache %q: zstd not available in image; skipping\n", entry.Key) 92 continue 93 case execErr != nil: 94 up.Abort(execErr) 95 return fmt.Errorf("save cache %q: %w", entry.Key, execErr) 96 case exit != 0: 97 up.Abort(fmt.Errorf("guest exited %d", exit)) 98 return fmt.Errorf("save cache %q: save script exited %d", entry.Key, exit) 99 } 100 if err := up.Finish(); err != nil { 101 return fmt.Errorf("save cache %q: %w", entry.Key, err) 102 } 103 fmt.Fprintf(out, "cache %q: saved\n", entry.Key) 104 } 105 return nil 106} 107 108func cacheExecStart(state *workflowState, script string) *agentv1.ExecStart { 109 return &agentv1.ExecStart{ 110 Argv: []string{state.ImageSpec.Shell, "-c", script}, 111 Env: guestBaseEnv(), 112 User: guestWorkflowUser, 113 } 114} 115 116type countingReader struct { 117 r io.Reader 118 n *int64 119} 120 121func (c *countingReader) Read(p []byte) (int, error) { 122 n, err := c.r.Read(p) 123 *c.n += int64(n) 124 return n, err 125}