package engine import ( "bufio" "bytes" "context" "errors" "io" "log/slog" "os" "os/exec" "path/filepath" "strings" "testing" "time" "github.com/bluesky-social/indigo/atproto/syntax" "tangled.org/core/api/tangled" "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 addCacheTestRepo(t *testing.T, d *db.DB) (syntax.DID, syntax.DID) { t.Helper() repoDid := syntax.DID("did:plc:repo") ownerDid := syntax.DID("did:plc:owner") if err := d.AddRepo(db.Repo{ Knot: "knot.example", Owner: ownerDid, Rkey: syntax.RecordKey("repo"), RepoDid: repoDid, }); err != nil { t.Fatal(err) } return repoDid, ownerDid } 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 TestCachePlanPreservesRestoreAtEntryQuota(t *testing.T) { ctx := context.Background() d := newCacheTestDB(t) repoDid, ownerDid := addCacheTestRepo(t, d) old := cacheTestEntry( "restore", repoDid.String(), "microvm/"+models.CacheNamespace(), "deps", "", "ready", time.Now(), ) old.OwnerDID = ownerDid.String() insertCacheTestEntry(t, d, old) controller := NewLocalCacheController( d, &fakeStorage{objects: map[string][]byte{old.StorageKey: []byte("archive")}}, "", 0, 1, discardLogger, ) bindings, err := controller.Plan(ctx, &models.Pipeline{RepoDid: repoDid}, &models.Workflow{ Engine: "microvm", Caches: []models.CacheEntry{{Key: "deps", Paths: []string{"deps"}}}, }) if err != nil { t.Fatal(err) } if len(bindings) != 1 || bindings[0].RestoreID != old.ID || bindings[0].SaveKey != "" { t.Fatalf("quota binding = %+v, want restore-only %q", bindings, old.ID) } } func TestScheduledCacheRestoresPushCacheOnSameRef(t *testing.T) { ctx := context.Background() d := newCacheTestDB(t) repoDid, _ := addCacheTestRepo(t, d) base := &fakeStorage{objects: make(map[string][]byte)} controller := NewLocalCacheController(d, base, "", 0, 0, discardLogger) wf := &models.Workflow{ Name: "build.yml", Engine: "microvm", Caches: []models.CacheEntry{{Key: "deps", Paths: []string{"deps"}}}, } sha := strings.Repeat("a", 40) push, err := controller.Plan(ctx, &models.Pipeline{ RepoDid: repoDid, TriggerMetadata: &tangled.Pipeline_TriggerMetadata{ Kind: "push", Push: &tangled.Pipeline_PushTriggerData{Ref: "refs/heads/main", NewSha: sha}, }, }, wf) if err != nil { t.Fatal(err) } if len(push) != 1 || push[0].SaveKey == "" { t.Fatalf("push did not reserve a cache: %+v", push) } store := newTrackedCacheStore(base, controller, discardLogger, push) if err := store.Put(ctx, push[0].SaveKey, strings.NewReader("push cache")); err != nil { t.Fatal(err) } for _, ref := range []string{"refs/heads/main", "refs/heads/other"} { t.Run(ref, func(t *testing.T) { bindings, err := controller.Plan(ctx, &models.Pipeline{ RepoDid: repoDid, TriggerMetadata: &tangled.Pipeline_TriggerMetadata{ Kind: "schedule", Schedule: &tangled.Pipeline_ScheduleTriggerData{Ref: ref, Sha: sha}, }, }, wf) if err != nil { t.Fatal(err) } if len(bindings) != 1 { t.Fatalf("cache bindings = %+v", bindings) } want := "" if ref == "refs/heads/main" { want = push[0].SaveKey } if bindings[0].RestoreKey != want { t.Fatalf("restore key = %q, want %q", bindings[0].RestoreKey, want) } }) } } func TestCachePlanKeepsPendingEntryThroughWorkflowDeadline(t *testing.T) { d := newCacheTestDB(t) repoDid, _ := addCacheTestRepo(t, d) deadline := time.Now().Add(6 * time.Hour) ctx, cancel := context.WithDeadline(context.Background(), deadline) defer cancel() controller := NewLocalCacheController( d, &fakeStorage{objects: make(map[string][]byte)}, "", 0, 0, discardLogger, ) bindings, err := controller.Plan(ctx, &models.Pipeline{RepoDid: repoDid}, &models.Workflow{ Engine: "microvm", Caches: []models.CacheEntry{{Key: "deps", Paths: []string{"deps"}}}, }) if err != nil { t.Fatal(err) } var pendingUntil int64 if err := d.QueryRow( `select last_used_at from cache_entries where id = ?`, bindings[0].SaveID, ).Scan(&pendingUntil); err != nil { t.Fatal(err) } if pendingUntil != deadline.UnixNano() { t.Fatalf("pending deadline = %d, want %d", pendingUntil, deadline.UnixNano()) } } 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 TestTrackedCacheStoreSaveLifecycle(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) pending := cacheTestEntry("save", "did:plc:repo", "microvm", "deps", "hash", "pending", time.Now()) insertCacheTestEntry(t, d, pending) base := &fakeStorage{objects: map[string][]byte{old.StorageKey: []byte("old")}} controller := NewLocalCacheController(d, base, "", 0, 0, discardLogger) store := newTrackedCacheStore(base, controller, discardLogger, []models.CacheBinding{{ Key: "deps", Hash: "hash", SaveID: pending.ID, SaveKey: pending.StorageKey, }}) payload := []byte("cache archive bytes") if err := store.Put(ctx, pending.StorageKey, bytes.NewReader(payload)); err != nil { t.Fatalf("tracked Put: %v", err) } state, size, _ := cacheTestEntryState(t, d, pending.ID) if state != "ready" || size != int64(len(payload)) { t.Fatalf("saved metadata = state %q, size %d", state, size) } if got := base.objects[pending.StorageKey]; !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(pending.StorageKey) || !cacheTestEntryExists(t, d, pending.ID) { t.Fatal("cleanup removed a completed cache") } } func TestTrackedCacheStoreRejectsOverQuotaSave(t *testing.T) { ctx := context.Background() d := newCacheTestDB(t) existing := cacheTestEntry("existing", "did:plc:repo", "microvm", "deps", "old", "ready", time.Now()) existing.SizeBytes = 4 insertCacheTestEntry(t, d, existing) pending := cacheTestEntry("save", "did:plc:repo", "microvm", "deps", "new", "pending", time.Now()) pending.SizeBytes = 0 insertCacheTestEntry(t, d, pending) base := &fakeStorage{objects: map[string][]byte{existing.StorageKey: []byte("kept")}} controller := NewLocalCacheController(d, base, "", 4, 0, discardLogger) store := newTrackedCacheStore(base, controller, discardLogger, []models.CacheBinding{{ Key: "deps", SaveID: pending.ID, SaveKey: pending.StorageKey, }}) if err := store.Put(ctx, pending.StorageKey, strings.NewReader("x")); err == nil { t.Fatal("over-quota cache save succeeded") } store.cleanup(ctx) if base.has(pending.StorageKey) || cacheTestEntryExists(t, d, pending.ID) { t.Fatal("over-quota save left an object or metadata") } if !base.has(existing.StorageKey) || !cacheTestEntryExists(t, d, existing.ID) { t.Fatal("over-quota save removed the existing cache") } } func TestTrackedCacheStoreRestoreTouchesEntry(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")}} controller := NewLocalCacheController(d, base, "", 0, 0, discardLogger) store := newTrackedCacheStore(base, controller, discardLogger, []models.CacheBinding{{ RestoreID: entry.ID, RestoreKey: entry.StorageKey, }}) r, err := store.Get(ctx, entry.StorageKey) if err != nil { t.Fatalf("tracked 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 TestTrackedCacheStoreFailedPutCleanup(t *testing.T) { ctx := context.Background() d := newCacheTestDB(t) putErr := errors.New("put failed") base := &fakeStorage{objects: make(map[string][]byte), putErr: putErr} pending := cacheTestEntry("save", "did:plc:repo", "microvm", "deps", "", "pending", time.Now()) insertCacheTestEntry(t, d, pending) controller := NewLocalCacheController(d, base, "", 0, 0, discardLogger) store := newTrackedCacheStore(base, controller, discardLogger, []models.CacheBinding{{ Key: "deps", SaveID: pending.ID, SaveKey: pending.StorageKey, }}) if err := store.Put(ctx, pending.StorageKey, strings.NewReader("partial")); !errors.Is(err, putErr) { t.Fatalf("tracked Put error = %v, want %v", err, putErr) } store.cleanup(ctx) if base.has(pending.StorageKey) { t.Fatal("cleanup left partial object") } if cacheTestEntryExists(t, d, pending.ID) { t.Fatal("cleanup left pending metadata") } } func TestTrackedCacheStoreRejectsUnpreparedKeys(t *testing.T) { store := newTrackedCacheStore( &fakeStorage{objects: map[string][]byte{"objects/missing": []byte("foreign")}}, &preplannedCacheController{}, discardLogger, nil, ) if _, err := store.Get(context.Background(), "objects/missing"); err == nil { t.Fatal("Get accepted a key without restore metadata") } if err := store.Put(context.Background(), "objects/missing", strings.NewReader("data")); err == nil { t.Fatal("Put accepted a key without pending metadata") } } func TestCacheSaveScriptQuotesPaths(t *testing.T) { workspace := t.TempDir() cachePath := filepath.Join(workspace, "cache;name") if err := os.Mkdir(cachePath, 0o755); err != nil { t.Fatal(err) } if err := os.WriteFile(filepath.Join(cachePath, "value"), []byte("cached"), 0o644); err != nil { t.Fatal(err) } bin := t.TempDir() zstd := filepath.Join(bin, "zstd") if err := os.WriteFile(zstd, []byte("#!/bin/sh\nset --\nexec cat\n"), 0o755); err != nil { t.Fatal(err) } cmd := exec.Command("bash", "-c", CacheSaveScript([]string{"cache;name"}, workspace, 0)) cmd.Env = append(os.Environ(), "PATH="+bin+":"+os.Getenv("PATH")) if output, err := cmd.Output(); err != nil || len(output) == 0 { t.Fatalf("cache save script = %d bytes, %v", len(output), err) } } 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 TestCacheBindingSaveOn(t *testing.T) { cases := []struct { when string onFail, onPass bool }{ {"", false, true}, {"on-success", false, true}, {"always", true, true}, } for _, tc := range cases { binding := models.CacheBinding{When: tc.when} if got := binding.SaveOn(true); got != tc.onFail { t.Errorf("when=%q failed run: got %v, want %v", tc.when, got, tc.onFail) } if got := binding.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)