Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415package 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", Knot: "k", Rkey: "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) }}