This repository has no description
1package engine
2
3import (
4 "bufio"
5 "bytes"
6 "context"
7 "errors"
8 "io"
9 "log/slog"
10 "os"
11 "os/exec"
12 "path/filepath"
13 "slices"
14 "strings"
15 "testing"
16 "time"
17
18 "tangled.org/core/spindle/db"
19 "tangled.org/core/spindle/models"
20 "tangled.org/core/spindle/storage"
21)
22
23type fakeStorage struct {
24 objects map[string][]byte
25 putErr error
26 deleteErr error // fails once, then clears
27 onDelete func(string)
28}
29
30func (f *fakeStorage) Get(_ context.Context, key string) (io.ReadCloser, error) {
31 data, ok := f.objects[key]
32 if !ok {
33 return nil, storage.ErrNotExist
34 }
35 return io.NopCloser(bytes.NewReader(data)), nil
36}
37
38func (f *fakeStorage) Put(_ context.Context, key string, r io.Reader) error {
39 data, err := io.ReadAll(r)
40 if err != nil {
41 return err
42 }
43 f.objects[key] = data
44 return f.putErr
45}
46
47func (f *fakeStorage) Delete(_ context.Context, key string) error {
48 if f.onDelete != nil {
49 f.onDelete(key)
50 }
51 if f.deleteErr != nil {
52 err := f.deleteErr
53 f.deleteErr = nil
54 return err
55 }
56 delete(f.objects, key)
57 return nil
58}
59
60func (f *fakeStorage) has(key string) bool {
61 _, ok := f.objects[key]
62 return ok
63}
64
65func gitRepo(t *testing.T, files map[string]string) (string, string) {
66 t.Helper()
67 if _, err := exec.LookPath("git"); err != nil {
68 t.Skip("git not available")
69 }
70 dir := t.TempDir()
71 run := func(args ...string) {
72 t.Helper()
73 cmd := exec.Command("git", append([]string{"-C", dir, "-c", "user.email=t@t", "-c", "user.name=t", "-c", "init.defaultBranch=main"}, args...)...)
74 if out, err := cmd.CombinedOutput(); err != nil {
75 t.Fatalf("git %v: %v\n%s", args, err, out)
76 }
77 }
78 run("init")
79 for name, content := range files {
80 p := filepath.Join(dir, name)
81 if err := os.MkdirAll(filepath.Dir(p), 0o755); err != nil {
82 t.Fatal(err)
83 }
84 if err := os.WriteFile(p, []byte(content), 0o644); err != nil {
85 t.Fatal(err)
86 }
87 }
88 run("add", ".")
89 run("commit", "-m", "init")
90 return dir, "HEAD"
91}
92
93func newCacheTestDB(t *testing.T) *db.DB {
94 t.Helper()
95 d, err := db.Make(context.Background(), filepath.Join(t.TempDir(), "spindle.db"))
96 if err != nil {
97 t.Fatalf("db.Make: %v", err)
98 }
99 t.Cleanup(func() { _ = d.Close() })
100 return d
101}
102
103func cacheTestEntry(id, repoDID, engine, key, hash, state string, at time.Time) db.CacheEntry {
104 return db.CacheEntry{
105 ID: id,
106 StorageKey: "objects/" + id,
107 OwnerDID: "did:plc:owner",
108 RepoDID: repoDID,
109 Engine: engine,
110 CacheKey: key,
111 CacheHash: hash,
112 State: state,
113 CreatedAt: at,
114 LastUsedAt: at,
115 }
116}
117
118func insertCacheTestEntry(t *testing.T, d *db.DB, entry db.CacheEntry) {
119 t.Helper()
120 if err := d.InsertCacheEntry(context.Background(), entry); err != nil {
121 t.Fatalf("InsertCacheEntry(%s): %v", entry.ID, err)
122 }
123}
124
125func cacheTestEntryState(t *testing.T, d *db.DB, id string) (string, int64, time.Time) {
126 t.Helper()
127 var state string
128 var size, lastUsed int64
129 if err := d.QueryRow(`select state, size_bytes, last_used_at from cache_entries where id = ?`, id).Scan(&state, &size, &lastUsed); err != nil {
130 t.Fatalf("query cache entry %s: %v", id, err)
131 }
132 return state, size, time.Unix(0, lastUsed)
133}
134
135func cacheTestEntryExists(t *testing.T, d *db.DB, id string) bool {
136 t.Helper()
137 var exists bool
138 if err := d.QueryRow(`select exists(select 1 from cache_entries where id = ?)`, id).Scan(&exists); err != nil {
139 t.Fatalf("query cache entry existence %s: %v", id, err)
140 }
141 return exists
142}
143
144var discardLogger = slog.New(slog.NewTextHandler(io.Discard, nil))
145
146func TestResolveCachesExactUnhashed(t *testing.T) {
147 ctx := context.Background()
148 d := newCacheTestDB(t)
149 base := time.Date(2026, 4, 5, 6, 7, 8, 0, time.UTC)
150
151 insertCacheTestEntry(t, d, cacheTestEntry("exact", "did:plc:repo", "microvm", "deps", "", "ready", base))
152 insertCacheTestEntry(t, d, cacheTestEntry("other-repo", "did:plc:other", "microvm", "deps", "", "ready", base.Add(time.Hour)))
153 insertCacheTestEntry(t, d, cacheTestEntry("other-engine", "did:plc:repo", "nixery", "deps", "", "ready", base.Add(time.Hour)))
154 insertCacheTestEntry(t, d, cacheTestEntry("pending", "did:plc:repo", "microvm", "deps", "", "pending", base.Add(time.Hour)))
155 insertCacheTestEntry(t, d, cacheTestEntry("hashed-only", "did:plc:repo", "microvm", "tools", "old", "ready", base))
156
157 resolved := ResolveCaches(ctx, discardLogger, d, "did:plc:repo", "microvm", "", "", []models.CacheEntry{
158 {Key: "deps", Paths: []string{"/x"}, CompressionLevel: 7, When: "always"},
159 {Key: "tools", Paths: []string{"/y"}},
160 })
161 if len(resolved) != 2 {
162 t.Fatalf("got %d resolved entries, want 2", len(resolved))
163 }
164 got := resolved[0]
165 if got.RestoreID != "exact" || got.RestoreKey != "objects/exact" || got.RestoreName != "" {
166 t.Fatalf("exact restore = (%q, %q, %q)", got.RestoreID, got.RestoreKey, got.RestoreName)
167 }
168 if got.SaveKey != "" {
169 t.Fatalf("resolve allocated a save key %q", got.SaveKey)
170 }
171 if got.CompressionLevel != 7 || got.When != "always" {
172 t.Fatalf("resolved metadata = %+v", got)
173 }
174 if resolved[1].RestoreKey != "" {
175 t.Fatalf("unhashed entry used hashed fallback %q", resolved[1].RestoreKey)
176 }
177}
178
179func TestResolveCachesHashExactAndFallback(t *testing.T) {
180 ctx := context.Background()
181 d := newCacheTestDB(t)
182 repoPath, rev := gitRepo(t, map[string]string{"go.sum": "v1 contents", "go.mod": "module x"})
183 request := []models.CacheEntry{{Key: "go-mod", Hash: []string{"go.sum", "go.mod"}, Paths: []string{"/x"}}}
184
185 first := ResolveCaches(ctx, discardLogger, d, "did:plc:repo", "microvm", repoPath, rev, request)[0]
186 if first.Hash == "" {
187 t.Fatal("hash is empty")
188 }
189
190 base := time.Date(2026, 4, 5, 6, 7, 8, 0, time.UTC)
191 insertCacheTestEntry(t, d, cacheTestEntry("exact", "did:plc:repo", "microvm", "go-mod", first.Hash, "ready", base.Add(time.Minute)))
192 insertCacheTestEntry(t, d, cacheTestEntry("sibling", "did:plc:repo", "microvm", "go-modules", "newer", "ready", base.Add(time.Hour)))
193
194 exact := ResolveCaches(ctx, discardLogger, d, "did:plc:repo", "microvm", repoPath, rev, request)[0]
195 if exact.RestoreID != "exact" || exact.RestoreKey != "objects/exact" || exact.RestoreName != "" {
196 t.Fatalf("exact generation restore = %+v", exact)
197 }
198
199 if err := os.WriteFile(filepath.Join(repoPath, "go.sum"), []byte("v2 contents"), 0o644); err != nil {
200 t.Fatal(err)
201 }
202 commit := exec.Command("git", "-C", repoPath, "-c", "user.email=t@t", "-c", "user.name=t", "commit", "-qam", "bump")
203 if out, err := commit.CombinedOutput(); err != nil {
204 t.Fatalf("commit: %v\n%s", err, out)
205 }
206 rotated := ResolveCaches(ctx, discardLogger, d, "did:plc:repo", "microvm", repoPath, rev, request)[0]
207 if rotated.Hash == first.Hash {
208 t.Fatal("lockfile change did not rotate the cache hash")
209 }
210 if rotated.RestoreID != "exact" || rotated.RestoreKey != "objects/exact" || rotated.RestoreName != "go-mod-"+first.Hash {
211 t.Fatalf("fallback restore = %+v", rotated)
212 }
213}
214
215func TestResolveCachesMissingHashFiles(t *testing.T) {
216 ctx := context.Background()
217 d := newCacheTestDB(t)
218 repoPath, rev := gitRepo(t, map[string]string{"go.mod": "module x"})
219
220 partial := ResolveCaches(ctx, discardLogger, d, "did:plc:repo", "microvm", repoPath, rev, []models.CacheEntry{
221 {Key: "go-mod", Hash: []string{"go.sum", "go.mod"}, Paths: []string{"/x"}},
222 })[0]
223 if len(partial.Hash) != 12 {
224 t.Fatalf("partial hash = %q, want 12 characters", partial.Hash)
225 }
226
227 none := ResolveCaches(ctx, discardLogger, d, "did:plc:repo", "microvm", repoPath, rev, []models.CacheEntry{
228 {Key: "go-mod", Hash: []string{"nope.lock"}, Paths: []string{"/x"}},
229 })[0]
230 if none.Hash != "" || none.RestoreKey != "" {
231 t.Fatalf("all-missing result = hash %q, restore %q", none.Hash, none.RestoreKey)
232 }
233}
234
235func TestIndexedCacheStoreSaveLifecycle(t *testing.T) {
236 ctx := context.Background()
237 d := newCacheTestDB(t)
238 old := cacheTestEntry("superseded", "did:plc:repo", "microvm", "deps", "hash", "ready", time.Date(2026, 1, 2, 3, 4, 5, 0, time.UTC))
239 insertCacheTestEntry(t, d, old)
240 base := &fakeStorage{objects: map[string][]byte{old.StorageKey: []byte("old")}}
241 entries := []ResolvedCache{{Key: "deps", Hash: "hash"}}
242
243 store, err := prepareCacheSaves(ctx, base, d, discardLogger, "did:plc:owner", "did:plc:repo", "microvm", entries)
244 if err != nil {
245 t.Fatalf("prepareCacheSaves: %v", err)
246 }
247 saveKey := entries[0].SaveKey
248 id := store.entries[saveKey]
249
250 payload := []byte("cache archive bytes")
251 if err := store.Put(ctx, saveKey, bytes.NewReader(payload)); err != nil {
252 t.Fatalf("indexed Put: %v", err)
253 }
254 state, size, _ := cacheTestEntryState(t, d, id)
255 if state != "ready" || size != int64(len(payload)) {
256 t.Fatalf("saved metadata = state %q, size %d", state, size)
257 }
258 if got := base.objects[saveKey]; !bytes.Equal(got, payload) {
259 t.Fatalf("stored payload = %q, want %q", got, payload)
260 }
261 if base.has(old.StorageKey) || cacheTestEntryExists(t, d, old.ID) {
262 t.Fatal("completed save left the superseded generation")
263 }
264 store.cleanup(ctx)
265 if !base.has(saveKey) || !cacheTestEntryExists(t, d, id) {
266 t.Fatal("cleanup removed a completed cache")
267 }
268}
269
270func TestIndexedCacheStoreRestoreTouchesEntry(t *testing.T) {
271 ctx := context.Background()
272 d := newCacheTestDB(t)
273 old := time.Date(2025, 1, 2, 3, 4, 5, 0, time.UTC)
274 entry := cacheTestEntry("restore", "did:plc:repo", "microvm", "deps", "hash", "ready", old)
275 insertCacheTestEntry(t, d, entry)
276 base := &fakeStorage{objects: map[string][]byte{entry.StorageKey: []byte("archive")}}
277 store := cacheStoreForRestore(base, d, discardLogger, []ResolvedCache{{RestoreID: entry.ID, RestoreKey: entry.StorageKey}})
278
279 r, err := store.Get(ctx, entry.StorageKey)
280 if err != nil {
281 t.Fatalf("indexed Get: %v", err)
282 }
283 data, err := io.ReadAll(r)
284 if err != nil {
285 t.Fatalf("read restored object: %v", err)
286 }
287 if err := r.Close(); err != nil {
288 t.Fatalf("close restored object: %v", err)
289 }
290 if string(data) != "archive" {
291 t.Fatalf("restored payload = %q", data)
292 }
293 _, _, touched := cacheTestEntryState(t, d, entry.ID)
294 if !touched.After(old) {
295 t.Fatalf("last used = %v, want after %v", touched, old)
296 }
297}
298
299func TestIndexedCacheStoreFailedPutCleanup(t *testing.T) {
300 ctx := context.Background()
301 d := newCacheTestDB(t)
302 putErr := errors.New("put failed")
303 base := &fakeStorage{objects: make(map[string][]byte), putErr: putErr}
304 entries := []ResolvedCache{{Key: "deps"}}
305 store, err := prepareCacheSaves(ctx, base, d, discardLogger, "did:plc:owner", "did:plc:repo", "microvm", entries)
306 if err != nil {
307 t.Fatalf("prepareCacheSaves: %v", err)
308 }
309 saveKey := entries[0].SaveKey
310
311 if err := store.Put(ctx, saveKey, strings.NewReader("partial")); !errors.Is(err, putErr) {
312 t.Fatalf("indexed Put error = %v, want %v", err, putErr)
313 }
314 store.cleanup(ctx)
315 if base.has(saveKey) {
316 t.Fatal("cleanup left partial object")
317 }
318 if cacheTestEntryExists(t, d, store.entries[saveKey]) {
319 t.Fatal("cleanup left pending metadata")
320 }
321}
322
323func TestIndexedCacheStoreRejectsUnpreparedKey(t *testing.T) {
324 store := &indexedCacheStore{}
325 if err := store.Put(context.Background(), "objects/missing", strings.NewReader("data")); err == nil {
326 t.Fatal("Put accepted a key without pending metadata")
327 }
328}
329
330func TestAbsolutizePaths(t *testing.T) {
331 got := AbsolutizePaths([]string{"node_modules", "/root/.cache", ".gocache"}, "/workspace/repo")
332 want := []string{"/workspace/repo/node_modules", "/root/.cache", "/workspace/repo/.gocache"}
333 if !slices.Equal(got, want) {
334 t.Fatalf("got %v, want %v", got, want)
335 }
336}
337
338func TestCacheDecompressCmd(t *testing.T) {
339 cases := []struct {
340 name string
341 head []byte
342 want string
343 }{
344 {"zstd", []byte{0x28, 0xb5, 0x2f, 0xfd, 0x00}, "zstd -dc"},
345 {"gzip", []byte{0x1f, 0x8b, 0x08, 0x00}, "gzip -dc"},
346 {"empty", nil, "zstd -dc"},
347 }
348 for _, tc := range cases {
349 got := CacheDecompressCmd(bufio.NewReader(bytes.NewReader(tc.head)))
350 if got != tc.want {
351 t.Errorf("%s: got %q, want %q", tc.name, got, tc.want)
352 }
353 }
354}
355
356func TestResolvedCacheSaveOn(t *testing.T) {
357 cases := []struct {
358 when string
359 onFail, onPass bool
360 }{
361 {"", false, true},
362 {"on-success", false, true},
363 {"always", true, true},
364 }
365 for _, tc := range cases {
366 rc := ResolvedCache{When: tc.when}
367 if got := rc.saveOn(true); got != tc.onFail {
368 t.Errorf("when=%q failed run: got %v, want %v", tc.when, got, tc.onFail)
369 }
370 if got := rc.saveOn(false); got != tc.onPass {
371 t.Errorf("when=%q passing run: got %v, want %v", tc.when, got, tc.onPass)
372 }
373 }
374}
375
376var _ storage.Storage = (*fakeStorage)(nil)