package db import ( "context" "database/sql" "errors" "testing" "time" ) func testCacheEntry(id, hash, state string, createdAt time.Time) CacheEntry { return CacheEntry{ ID: id, StorageKey: "objects/" + id, OwnerDID: "did:plc:owner", RepoDID: "did:plc:repo", Engine: "microvm", CacheKey: "dependencies", CacheHash: hash, SizeBytes: 10, State: state, CreatedAt: createdAt, LastUsedAt: createdAt, } } func insertTestCacheEntry(t *testing.T, d *DB, entry CacheEntry) { t.Helper() if err := d.InsertCacheEntry(context.Background(), entry); err != nil { t.Fatalf("InsertCacheEntry(%s): %v", entry.ID, err) } } func TestInsertCacheEntryWithinQuota(t *testing.T) { ctx := context.Background() d := newTestDB(t) now := time.Now() first := testCacheEntry("first", "a", "pending", now) inserted, err := d.InsertCacheEntryWithinQuota(ctx, first, 1) if err != nil || !inserted { t.Fatalf("first quota insert = (%t, %v)", inserted, err) } second := testCacheEntry("second", "b", "pending", now) inserted, err = d.InsertCacheEntryWithinQuota(ctx, second, 1) if err != nil { t.Fatal(err) } if inserted { t.Fatal("entry quota accepted a second object for the same owner") } second.OwnerDID = "did:plc:other" inserted, err = d.InsertCacheEntryWithinQuota(ctx, second, 1) if err != nil || !inserted { t.Fatalf("other owner quota insert = (%t, %v)", inserted, err) } } func TestMarkCacheEntryReadyEvictsOverQuotaOwner(t *testing.T) { ctx := context.Background() d := newTestDB(t) now := time.Now() existing := testCacheEntry("existing", "a", "ready", now) existing.SizeBytes = 40 insertTestCacheEntry(t, d, existing) pending := testCacheEntry("pending", "b", "pending", now) pending.SizeBytes = 0 insertTestCacheEntry(t, d, pending) if _, err := d.MarkCacheEntryReady(ctx, pending.ID, 11, 50, now); err != nil { t.Fatalf("over-quota ready should evict LRU instead of failing: %v", err) } var state string if err := d.QueryRow(`select state from cache_entries where id = ?`, existing.ID).Scan(&state); err != nil { t.Fatal(err) } if state != "deleting" { t.Fatalf("evicted entry state = %q, want deleting", state) } if err := d.QueryRow(`select state from cache_entries where id = ?`, pending.ID).Scan(&state); err != nil { t.Fatal(err) } if state != "ready" { t.Fatalf("pending entry state = %q, want ready", state) } } func TestCacheEntryLookup(t *testing.T) { ctx := context.Background() d := newTestDB(t) base := time.Date(2026, 1, 2, 3, 4, 5, 6, time.UTC) insertTestCacheEntry(t, d, testCacheEntry("exact-old", "requested", "ready", base)) insertTestCacheEntry(t, d, testCacheEntry("exact-new", "requested", "ready", base.Add(time.Second))) insertTestCacheEntry(t, d, testCacheEntry("exact-pending", "requested", "pending", base.Add(2*time.Second))) insertTestCacheEntry(t, d, testCacheEntry("fallback-old", "old-a", "ready", base.Add(3*time.Second))) insertTestCacheEntry(t, d, testCacheEntry("fallback-new", "old-b", "ready", base.Add(4*time.Second))) insertTestCacheEntry(t, d, testCacheEntry("fallback-pending", "old-c", "pending", base.Add(5*time.Second))) exact, err := d.FindCacheEntry(ctx, "did:plc:repo", "microvm", "dependencies", "requested") if err != nil { t.Fatalf("FindCacheEntry: %v", err) } if exact.ID != "exact-new" { t.Fatalf("FindCacheEntry returned %q, want exact-new", exact.ID) } if !exact.CreatedAt.Equal(base.Add(time.Second)) || !exact.LastUsedAt.Equal(base.Add(time.Second)) { t.Fatalf("timestamps = (%v, %v), want %v", exact.CreatedAt, exact.LastUsedAt, base.Add(time.Second)) } fallback, err := d.FindFallbackCacheEntry(ctx, "did:plc:repo", "microvm", "dependencies", "requested") if err != nil { t.Fatalf("FindFallbackCacheEntry: %v", err) } if fallback.ID != "fallback-new" { t.Fatalf("FindFallbackCacheEntry returned %q, want fallback-new", fallback.ID) } if _, err := d.FindCacheEntry(ctx, "did:plc:repo", "microvm", "missing", "requested"); !errors.Is(err, sql.ErrNoRows) { t.Fatalf("missing FindCacheEntry error = %v, want sql.ErrNoRows", err) } } func TestMarkCacheEntryReadySupersedesMatchingReady(t *testing.T) { ctx := context.Background() d := newTestDB(t) base := time.Date(2026, 1, 3, 4, 5, 6, 7, time.UTC) first := testCacheEntry("first-completion", "same-hash", "pending", base) second := testCacheEntry("second-completion", "same-hash", "pending", base.Add(time.Second)) insertTestCacheEntry(t, d, first) insertTestCacheEntry(t, d, second) superseded, err := d.MarkCacheEntryReady(ctx, first.ID, 100, 0, base.Add(2*time.Second)) if err != nil { t.Fatalf("mark first ready: %v", err) } if len(superseded) != 0 { t.Fatalf("first completion superseded %d entries, want none", len(superseded)) } superseded, err = d.MarkCacheEntryReady(ctx, second.ID, 200, 0, base.Add(3*time.Second)) if err != nil { t.Fatalf("mark second ready: %v", err) } if len(superseded) != 1 || superseded[0].ID != first.ID || superseded[0].State != "deleting" { t.Fatalf("second completion superseded %+v, want deleting %s", superseded, first.ID) } ready, err := d.FindCacheEntry(ctx, second.RepoDID, second.Engine, second.CacheKey, second.CacheHash) if err != nil { t.Fatalf("find surviving ready entry: %v", err) } if ready.ID != second.ID || ready.SizeBytes != 200 || !ready.LastUsedAt.Equal(base.Add(3*time.Second)) { t.Fatalf("surviving entry = %+v, want %s (200 bytes)", ready, second.ID) } } func TestMarkCacheEntryReadyConcurrentCompletions(t *testing.T) { ctx := context.Background() d := newTestDB(t) base := time.Date(2026, 1, 4, 5, 6, 7, 8, time.UTC) first := testCacheEntry("concurrent-a", "same-hash", "pending", base) second := testCacheEntry("concurrent-b", "same-hash", "pending", base.Add(time.Second)) insertTestCacheEntry(t, d, first) insertTestCacheEntry(t, d, second) type result struct { superseded []CacheEntry err error } start := make(chan struct{}) results := make(chan result, 2) for _, entry := range []CacheEntry{first, second} { entry := entry go func() { <-start superseded, err := d.MarkCacheEntryReady(ctx, entry.ID, 100, 0, base.Add(2*time.Second)) results <- result{superseded: superseded, err: err} }() } close(start) var superseded []CacheEntry for range 2 { result := <-results if result.err != nil { t.Fatalf("concurrent MarkCacheEntryReady: %v", result.err) } superseded = append(superseded, result.superseded...) } if len(superseded) != 1 || superseded[0].State != "deleting" { t.Fatalf("concurrent completions superseded %+v, want one deleting entry", superseded) } var readyCount, deletingCount int if err := d.QueryRowContext(ctx, ` select sum(state = 'ready'), sum(state = 'deleting') from cache_entries where repo_did = ? and engine = ? and cache_key = ? and cache_hash = ?`, first.RepoDID, first.Engine, first.CacheKey, first.CacheHash).Scan(&readyCount, &deletingCount); err != nil { t.Fatalf("count completion states: %v", err) } if readyCount != 1 || deletingCount != 1 { t.Fatalf("completion states = %d ready, %d deleting; want one each", readyCount, deletingCount) } } func TestCacheEntryReadyTouchAndExpiry(t *testing.T) { ctx := context.Background() d := newTestDB(t) base := time.Date(2026, 2, 3, 4, 5, 6, 7, time.UTC) readyOld := testCacheEntry("ready-old", "a", "ready", base) readyFresh := testCacheEntry("ready-fresh", "b", "ready", base) pendingOld := testCacheEntry("pending-old", "c", "pending", base) pendingFresh := testCacheEntry("pending-fresh", "d", "pending", base) pendingFresh.LastUsedAt = base.Add(20 * time.Minute) deletingOld := testCacheEntry("deleting-old", "e", "deleting", base) deletingFresh := testCacheEntry("deleting-fresh", "f", "deleting", base) deletingFresh.LastUsedAt = base.Add(20 * time.Minute) for _, entry := range []CacheEntry{readyOld, readyFresh, pendingOld, pendingFresh, deletingOld, deletingFresh} { insertTestCacheEntry(t, d, entry) } touchedAt := base.Add(30 * time.Minute) if err := d.TouchCacheEntry(ctx, readyFresh.ID, touchedAt); err != nil { t.Fatalf("TouchCacheEntry: %v", err) } expired, err := d.ExpiredCacheEntries(ctx, base.Add(10*time.Minute), base.Add(10*time.Minute), 10) if err != nil { t.Fatalf("ExpiredCacheEntries: %v", err) } got := make(map[string]bool, len(expired)) for _, entry := range expired { got[entry.ID] = true } if len(got) != 3 || !got[readyOld.ID] || !got[pendingOld.ID] || !got[deletingOld.ID] { t.Fatalf("expired IDs = %v, want ready-old, pending-old, and deleting-old", got) } limited, err := d.ExpiredCacheEntries(ctx, base.Add(10*time.Minute), base.Add(10*time.Minute), 1) if err != nil { t.Fatalf("limited ExpiredCacheEntries: %v", err) } if len(limited) != 1 { t.Fatalf("limited expiry count = %d, want 1", len(limited)) } } func TestCacheEntryClaim(t *testing.T) { ctx := context.Background() d := newTestDB(t) base := time.Date(2026, 2, 4, 5, 6, 7, 8, time.UTC) entry := testCacheEntry("claim-me", "hash", "ready", base) insertTestCacheEntry(t, d, entry) touchedAt := base.Add(time.Minute) if err := d.TouchCacheEntry(ctx, entry.ID, touchedAt); err != nil { t.Fatalf("TouchCacheEntry: %v", err) } claimed, err := d.ClaimCacheEntry(ctx, entry.ID, "ready", base) if err != nil { t.Fatalf("stale ClaimCacheEntry: %v", err) } if claimed { t.Fatal("stale ClaimCacheEntry claimed a touched entry") } claimed, err = d.ClaimCacheEntry(ctx, entry.ID, "ready", touchedAt) if err != nil { t.Fatalf("ClaimCacheEntry: %v", err) } if !claimed { t.Fatal("ClaimCacheEntry did not claim unchanged entry") } if err := d.TouchCacheEntry(ctx, entry.ID, base.Add(2*time.Minute)); err != nil { t.Fatalf("TouchCacheEntry while deleting: %v", err) } if _, err := d.MarkCacheEntryReady(ctx, entry.ID, 999, 0, base.Add(3*time.Minute)); !errors.Is(err, sql.ErrNoRows) { t.Fatalf("MarkCacheEntryReady while deleting error = %v, want sql.ErrNoRows", err) } var state string var lastUsedAt int64 var sizeBytes int64 if err := d.QueryRowContext(ctx, ` select state, last_used_at, size_bytes from cache_entries where id = ?`, entry.ID).Scan(&state, &lastUsedAt, &sizeBytes); err != nil { t.Fatalf("query claimed entry: %v", err) } if state != "deleting" || lastUsedAt != touchedAt.UnixNano() || sizeBytes != entry.SizeBytes { t.Fatalf("claimed entry = state %q, last used %d, size %d; want deleting, %d, %d", state, lastUsedAt, sizeBytes, touchedAt.UnixNano(), entry.SizeBytes) } if err := d.RestoreCacheEntryState(ctx, entry.ID, "ready"); err != nil { t.Fatalf("RestoreCacheEntryState: %v", err) } restored, err := d.FindCacheEntry(ctx, entry.RepoDID, entry.Engine, entry.CacheKey, entry.CacheHash) if err != nil { t.Fatalf("find restored entry: %v", err) } if restored.State != "ready" || !restored.LastUsedAt.Equal(touchedAt) { t.Fatalf("restored entry = state %q, last used %v", restored.State, restored.LastUsedAt) } } func TestCacheEntryDelete(t *testing.T) { ctx := context.Background() d := newTestDB(t) entry := testCacheEntry("delete-me", "hash", "ready", time.Now()) insertTestCacheEntry(t, d, entry) if err := d.DeleteCacheEntry(ctx, entry.ID); err != nil { t.Fatalf("DeleteCacheEntry: %v", err) } if err := d.DeleteCacheEntry(ctx, entry.ID); err != nil { t.Fatalf("second DeleteCacheEntry: %v", err) } if _, err := d.FindCacheEntry(ctx, entry.RepoDID, entry.Engine, entry.CacheKey, entry.CacheHash); !errors.Is(err, sql.ErrNoRows) { t.Fatalf("find deleted error = %v, want sql.ErrNoRows", err) } } func TestMillCacheCapabilitiesAreLeaseScopedAndConsumable(t *testing.T) { ctx := context.Background() d := newTestDB(t) lease := MillLease{ LeaseID: "lease-1", NodeID: "node-1", Epoch: "epoch-1", Engine: "microvm", PipelineID: "r", Workflow: "build", State: "running", } if err := d.SaveMillLease(lease); err != nil { t.Fatal(err) } capability := MillCacheCapability{ Action: "save", CacheID: "cache-1", StorageKey: "objects/did:web:example.com/cache-1", } if err := d.SaveMillCacheCapabilities(lease.LeaseID, []MillCacheCapability{capability}); err != nil { t.Fatal(err) } if err := d.ApplyEventBatch(nil, func(tx *EventBatchTx) error { storageKey, consumed, err := tx.ConsumeMillCacheCapability( ctx, lease.LeaseID, capability.Action, capability.CacheID, ) if err != nil { return err } if !consumed { t.Fatal("planned capability was not consumed") } if storageKey != capability.StorageKey { t.Fatalf("consumed storage key = %q, want %q", storageKey, capability.StorageKey) } return nil }); err != nil { t.Fatal(err) } if err := d.ApplyEventBatch(nil, func(tx *EventBatchTx) error { _, consumed, err := tx.ConsumeMillCacheCapability( ctx, "another-lease", capability.Action, capability.CacheID, ) if err != nil { return err } if consumed { t.Fatal("foreign lease consumed capability") } return nil }); err != nil { t.Fatal(err) } } func TestCacheObjectDeletionQueuePersistsUntilCompleted(t *testing.T) { ctx := context.Background() d := newTestDB(t) key := "objects/did:web:example.com/orphan" if err := d.ApplyEventBatch(nil, func(tx *EventBatchTx) error { return tx.QueueCacheObjectDeletion(ctx, key, time.Now()) }); err != nil { t.Fatal(err) } keys, err := d.PendingCacheObjectDeletions(ctx, 10) if err != nil { t.Fatal(err) } if len(keys) != 1 || keys[0] != key { t.Fatalf("pending deletions = %v, want %q", keys, key) } if err := d.CompleteCacheObjectDeletion(ctx, key); err != nil { t.Fatal(err) } keys, err = d.PendingCacheObjectDeletions(ctx, 10) if err != nil { t.Fatal(err) } if len(keys) != 0 { t.Fatalf("completed deletion remained queued: %v", keys) } }