package engine import ( "bufio" "bytes" "context" "errors" "io" "log/slog" "os" "os/exec" "path/filepath" "slices" "strings" "testing" "time" "tangled.org/core/spindle/db" "tangled.org/core/spindle/models" "tangled.org/core/spindle/storage" ) type fakeStorage struct { objects map[string][]byte putErr error deleteErr error // fails once, then clears onDelete func(string) } func (f *fakeStorage) Get(_ context.Context, key string) (io.ReadCloser, error) { data, ok := f.objects[key] if !ok { return nil, storage.ErrNotExist } return io.NopCloser(bytes.NewReader(data)), nil } func (f *fakeStorage) Put(_ context.Context, key string, r io.Reader) error { data, err := io.ReadAll(r) if err != nil { return err } f.objects[key] = data return f.putErr } func (f *fakeStorage) Delete(_ context.Context, key string) error { if f.onDelete != nil { f.onDelete(key) } if f.deleteErr != nil { err := f.deleteErr f.deleteErr = nil return err } delete(f.objects, key) return nil } func (f *fakeStorage) has(key string) bool { _, ok := f.objects[key] return ok } func gitRepo(t *testing.T, files map[string]string) (string, string) { t.Helper() if _, err := exec.LookPath("git"); err != nil { t.Skip("git not available") } dir := t.TempDir() run := func(args ...string) { t.Helper() cmd := exec.Command("git", append([]string{"-C", dir, "-c", "user.email=t@t", "-c", "user.name=t", "-c", "init.defaultBranch=main"}, args...)...) if out, err := cmd.CombinedOutput(); err != nil { t.Fatalf("git %v: %v\n%s", args, err, out) } } run("init") for name, content := range files { p := filepath.Join(dir, name) if err := os.MkdirAll(filepath.Dir(p), 0o755); err != nil { t.Fatal(err) } if err := os.WriteFile(p, []byte(content), 0o644); err != nil { t.Fatal(err) } } run("add", ".") run("commit", "-m", "init") return dir, "HEAD" } func newCacheTestDB(t *testing.T) *db.DB { t.Helper() d, err := db.Make(context.Background(), filepath.Join(t.TempDir(), "spindle.db")) if err != nil { t.Fatalf("db.Make: %v", err) } t.Cleanup(func() { _ = d.Close() }) return d } func cacheTestEntry(id, repoDID, engine, key, hash, state string, at time.Time) db.CacheEntry { return db.CacheEntry{ ID: id, StorageKey: "objects/" + id, OwnerDID: "did:plc:owner", RepoDID: repoDID, Engine: engine, CacheKey: key, CacheHash: hash, State: state, CreatedAt: at, LastUsedAt: at, } } func insertCacheTestEntry(t *testing.T, d *db.DB, entry db.CacheEntry) { t.Helper() if err := d.InsertCacheEntry(context.Background(), entry); err != nil { t.Fatalf("InsertCacheEntry(%s): %v", entry.ID, err) } } func cacheTestEntryState(t *testing.T, d *db.DB, id string) (string, int64, time.Time) { t.Helper() var state string var size, lastUsed int64 if err := d.QueryRow(`select state, size_bytes, last_used_at from cache_entries where id = ?`, id).Scan(&state, &size, &lastUsed); err != nil { t.Fatalf("query cache entry %s: %v", id, err) } return state, size, time.Unix(0, lastUsed) } func cacheTestEntryExists(t *testing.T, d *db.DB, id string) bool { t.Helper() var exists bool if err := d.QueryRow(`select exists(select 1 from cache_entries where id = ?)`, id).Scan(&exists); err != nil { t.Fatalf("query cache entry existence %s: %v", id, err) } return exists } var discardLogger = slog.New(slog.NewTextHandler(io.Discard, nil)) func TestResolveCachesExactUnhashed(t *testing.T) { ctx := context.Background() d := newCacheTestDB(t) base := time.Date(2026, 4, 5, 6, 7, 8, 0, time.UTC) insertCacheTestEntry(t, d, cacheTestEntry("exact", "did:plc:repo", "microvm", "deps", "", "ready", base)) insertCacheTestEntry(t, d, cacheTestEntry("other-repo", "did:plc:other", "microvm", "deps", "", "ready", base.Add(time.Hour))) insertCacheTestEntry(t, d, cacheTestEntry("other-engine", "did:plc:repo", "nixery", "deps", "", "ready", base.Add(time.Hour))) insertCacheTestEntry(t, d, cacheTestEntry("pending", "did:plc:repo", "microvm", "deps", "", "pending", base.Add(time.Hour))) insertCacheTestEntry(t, d, cacheTestEntry("hashed-only", "did:plc:repo", "microvm", "tools", "old", "ready", base)) resolved := ResolveCaches(ctx, discardLogger, d, "did:plc:repo", "microvm", "", "", []models.CacheEntry{ {Key: "deps", Paths: []string{"/x"}, CompressionLevel: 7, When: "always"}, {Key: "tools", Paths: []string{"/y"}}, }) if len(resolved) != 2 { t.Fatalf("got %d resolved entries, want 2", len(resolved)) } got := resolved[0] if got.RestoreID != "exact" || got.RestoreKey != "objects/exact" || got.RestoreName != "" { t.Fatalf("exact restore = (%q, %q, %q)", got.RestoreID, got.RestoreKey, got.RestoreName) } if got.SaveKey != "" { t.Fatalf("resolve allocated a save key %q", got.SaveKey) } if got.CompressionLevel != 7 || got.When != "always" { t.Fatalf("resolved metadata = %+v", got) } if resolved[1].RestoreKey != "" { t.Fatalf("unhashed entry used hashed fallback %q", resolved[1].RestoreKey) } } func TestResolveCachesHashExactAndFallback(t *testing.T) { ctx := context.Background() d := newCacheTestDB(t) repoPath, rev := gitRepo(t, map[string]string{"go.sum": "v1 contents", "go.mod": "module x"}) request := []models.CacheEntry{{Key: "go-mod", Hash: []string{"go.sum", "go.mod"}, Paths: []string{"/x"}}} first := ResolveCaches(ctx, discardLogger, d, "did:plc:repo", "microvm", repoPath, rev, request)[0] if first.Hash == "" { t.Fatal("hash is empty") } base := time.Date(2026, 4, 5, 6, 7, 8, 0, time.UTC) insertCacheTestEntry(t, d, cacheTestEntry("exact", "did:plc:repo", "microvm", "go-mod", first.Hash, "ready", base.Add(time.Minute))) insertCacheTestEntry(t, d, cacheTestEntry("sibling", "did:plc:repo", "microvm", "go-modules", "newer", "ready", base.Add(time.Hour))) exact := ResolveCaches(ctx, discardLogger, d, "did:plc:repo", "microvm", repoPath, rev, request)[0] if exact.RestoreID != "exact" || exact.RestoreKey != "objects/exact" || exact.RestoreName != "" { t.Fatalf("exact generation restore = %+v", exact) } if err := os.WriteFile(filepath.Join(repoPath, "go.sum"), []byte("v2 contents"), 0o644); err != nil { t.Fatal(err) } commit := exec.Command("git", "-C", repoPath, "-c", "user.email=t@t", "-c", "user.name=t", "commit", "-qam", "bump") if out, err := commit.CombinedOutput(); err != nil { t.Fatalf("commit: %v\n%s", err, out) } rotated := ResolveCaches(ctx, discardLogger, d, "did:plc:repo", "microvm", repoPath, rev, request)[0] if rotated.Hash == first.Hash { t.Fatal("lockfile change did not rotate the cache hash") } if rotated.RestoreID != "exact" || rotated.RestoreKey != "objects/exact" || rotated.RestoreName != "go-mod-"+first.Hash { t.Fatalf("fallback restore = %+v", rotated) } } func TestResolveCachesMissingHashFiles(t *testing.T) { ctx := context.Background() d := newCacheTestDB(t) repoPath, rev := gitRepo(t, map[string]string{"go.mod": "module x"}) partial := ResolveCaches(ctx, discardLogger, d, "did:plc:repo", "microvm", repoPath, rev, []models.CacheEntry{ {Key: "go-mod", Hash: []string{"go.sum", "go.mod"}, Paths: []string{"/x"}}, })[0] if len(partial.Hash) != 12 { t.Fatalf("partial hash = %q, want 12 characters", partial.Hash) } none := ResolveCaches(ctx, discardLogger, d, "did:plc:repo", "microvm", repoPath, rev, []models.CacheEntry{ {Key: "go-mod", Hash: []string{"nope.lock"}, Paths: []string{"/x"}}, })[0] if none.Hash != "" || none.RestoreKey != "" { t.Fatalf("all-missing result = hash %q, restore %q", none.Hash, none.RestoreKey) } } func TestIndexedCacheStoreSaveLifecycle(t *testing.T) { ctx := context.Background() d := newCacheTestDB(t) old := cacheTestEntry("superseded", "did:plc:repo", "microvm", "deps", "hash", "ready", time.Date(2026, 1, 2, 3, 4, 5, 0, time.UTC)) insertCacheTestEntry(t, d, old) base := &fakeStorage{objects: map[string][]byte{old.StorageKey: []byte("old")}} entries := []ResolvedCache{{Key: "deps", Hash: "hash"}} store, err := prepareCacheSaves(ctx, base, d, discardLogger, "did:plc:owner", "did:plc:repo", "microvm", entries) if err != nil { t.Fatalf("prepareCacheSaves: %v", err) } saveKey := entries[0].SaveKey id := store.entries[saveKey] payload := []byte("cache archive bytes") if err := store.Put(ctx, saveKey, bytes.NewReader(payload)); err != nil { t.Fatalf("indexed Put: %v", err) } state, size, _ := cacheTestEntryState(t, d, id) if state != "ready" || size != int64(len(payload)) { t.Fatalf("saved metadata = state %q, size %d", state, size) } if got := base.objects[saveKey]; !bytes.Equal(got, payload) { t.Fatalf("stored payload = %q, want %q", got, payload) } if base.has(old.StorageKey) || cacheTestEntryExists(t, d, old.ID) { t.Fatal("completed save left the superseded generation") } store.cleanup(ctx) if !base.has(saveKey) || !cacheTestEntryExists(t, d, id) { t.Fatal("cleanup removed a completed cache") } } func TestIndexedCacheStoreRestoreTouchesEntry(t *testing.T) { ctx := context.Background() d := newCacheTestDB(t) old := time.Date(2025, 1, 2, 3, 4, 5, 0, time.UTC) entry := cacheTestEntry("restore", "did:plc:repo", "microvm", "deps", "hash", "ready", old) insertCacheTestEntry(t, d, entry) base := &fakeStorage{objects: map[string][]byte{entry.StorageKey: []byte("archive")}} store := cacheStoreForRestore(base, d, discardLogger, []ResolvedCache{{RestoreID: entry.ID, RestoreKey: entry.StorageKey}}) r, err := store.Get(ctx, entry.StorageKey) if err != nil { t.Fatalf("indexed Get: %v", err) } data, err := io.ReadAll(r) if err != nil { t.Fatalf("read restored object: %v", err) } if err := r.Close(); err != nil { t.Fatalf("close restored object: %v", err) } if string(data) != "archive" { t.Fatalf("restored payload = %q", data) } _, _, touched := cacheTestEntryState(t, d, entry.ID) if !touched.After(old) { t.Fatalf("last used = %v, want after %v", touched, old) } } func TestIndexedCacheStoreFailedPutCleanup(t *testing.T) { ctx := context.Background() d := newCacheTestDB(t) putErr := errors.New("put failed") base := &fakeStorage{objects: make(map[string][]byte), putErr: putErr} entries := []ResolvedCache{{Key: "deps"}} store, err := prepareCacheSaves(ctx, base, d, discardLogger, "did:plc:owner", "did:plc:repo", "microvm", entries) if err != nil { t.Fatalf("prepareCacheSaves: %v", err) } saveKey := entries[0].SaveKey if err := store.Put(ctx, saveKey, strings.NewReader("partial")); !errors.Is(err, putErr) { t.Fatalf("indexed Put error = %v, want %v", err, putErr) } store.cleanup(ctx) if base.has(saveKey) { t.Fatal("cleanup left partial object") } if cacheTestEntryExists(t, d, store.entries[saveKey]) { t.Fatal("cleanup left pending metadata") } } func TestIndexedCacheStoreRejectsUnpreparedKey(t *testing.T) { store := &indexedCacheStore{} if err := store.Put(context.Background(), "objects/missing", strings.NewReader("data")); err == nil { t.Fatal("Put accepted a key without pending metadata") } } func TestAbsolutizePaths(t *testing.T) { got := AbsolutizePaths([]string{"node_modules", "/root/.cache", ".gocache"}, "/workspace/repo") want := []string{"/workspace/repo/node_modules", "/root/.cache", "/workspace/repo/.gocache"} if !slices.Equal(got, want) { t.Fatalf("got %v, want %v", got, want) } } func TestCacheDecompressCmd(t *testing.T) { cases := []struct { name string head []byte want string }{ {"zstd", []byte{0x28, 0xb5, 0x2f, 0xfd, 0x00}, "zstd -dc"}, {"gzip", []byte{0x1f, 0x8b, 0x08, 0x00}, "gzip -dc"}, {"empty", nil, "zstd -dc"}, } for _, tc := range cases { got := CacheDecompressCmd(bufio.NewReader(bytes.NewReader(tc.head))) if got != tc.want { t.Errorf("%s: got %q, want %q", tc.name, got, tc.want) } } } func TestResolvedCacheSaveOn(t *testing.T) { cases := []struct { when string onFail, onPass bool }{ {"", false, true}, {"on-success", false, true}, {"always", true, true}, } for _, tc := range cases { rc := ResolvedCache{When: tc.when} if got := rc.saveOn(true); got != tc.onFail { t.Errorf("when=%q failed run: got %v, want %v", tc.when, got, tc.onFail) } if got := rc.saveOn(false); got != tc.onPass { t.Errorf("when=%q passing run: got %v, want %v", tc.when, got, tc.onPass) } } } var _ storage.Storage = (*fakeStorage)(nil)