This repository has no description
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}