package quota_test import ( "context" "crypto/rand" "encoding/hex" "errors" "strings" "sync" "testing" "time" "tangled.org/core/spindle/quota" ) type mockStore struct { mu sync.Mutex limits map[string]int64 allocations map[string]int64 reservations map[string]quota.ReserveRequest phases map[string]string repoOwners map[string]string defaultUser int64 defaultRepo int64 } func newMockStore(defaultUser, defaultRepo int64) *mockStore { return &mockStore{ limits: make(map[string]int64), allocations: make(map[string]int64), reservations: make(map[string]quota.ReserveRequest), phases: make(map[string]string), repoOwners: make(map[string]string), defaultUser: defaultUser, defaultRepo: defaultRepo, } } var _ quota.Store = (*mockStore)(nil) func (m *mockStore) Reserve(ctx context.Context, req quota.ReserveRequest) (quota.Reservation, error) { if err := quota.Validate(req); err != nil { return quota.Reservation{}, err } m.mu.Lock() defer m.mu.Unlock() m.repoOwners[req.Identity.RepoDID] = req.Identity.OwnerDID // committed keys remain idempotent for res := range req.Resources { allocKey := req.Identity.RepoDID + "/" + string(res) + "/" + string(req.Kind) + "/" + req.Key if _, ok := m.allocations[allocKey]; ok { return quota.Reservation{ ID: "", Allowed: true, Temporary: false, Reason: quota.ReasonUnlimited, }, nil } } // reserved keys remain idempotent for id, resReq := range m.reservations { if resReq.Identity.RepoDID == req.Identity.RepoDID && resReq.Kind == req.Kind && resReq.Key == req.Key && (m.phases[id] == "reserved" || m.phases[id] == "publishing") { return quota.Reservation{ ID: id, Allowed: true, Temporary: false, Reason: quota.ReasonUnlimited, }, nil } } for res, amount := range req.Resources { repoLimitKey := req.Identity.RepoDID + "/" + string(res) var repoLimit int64 = -1 if lim, ok := m.limits[repoLimitKey]; ok { repoLimit = lim } else if res == quota.ResourceCacheStorageBytes && m.defaultRepo > 0 { repoLimit = m.defaultRepo } if repoLimit >= 0 { var currentUsage int64 seenKeys := make(map[string]bool) for k, amt := range m.allocations { parts := m.split(k) if parts[0] == req.Identity.RepoDID && parts[1] == string(res) { seenKeys[parts[3]] = true currentUsage += amt } } for id, r := range m.reservations { if r.Identity.RepoDID == req.Identity.RepoDID && (m.phases[id] == "reserved" || m.phases[id] == "publishing") { if amt, ok := r.Resources[res]; ok { if !seenKeys[r.Key] { seenKeys[r.Key] = true currentUsage += amt } } } } if currentUsage+amount > repoLimit { temporary := true if amount > repoLimit { temporary = false } return quota.Reservation{ ID: "", Allowed: false, Temporary: temporary, Reason: quota.ReasonRepoLimit, Resource: res, }, nil } } userLimitKey := req.Identity.OwnerDID + "/" + string(res) var userLimit int64 = -1 if lim, ok := m.limits[userLimitKey]; ok { userLimit = lim } else if res == quota.ResourceCacheStorageBytes && m.defaultUser > 0 { userLimit = m.defaultUser } if userLimit >= 0 { var currentUsage int64 seenKeys := make(map[string]bool) for k, amt := range m.allocations { parts := m.split(k) owner := m.repoOwners[parts[0]] if owner == req.Identity.OwnerDID && parts[1] == string(res) { if !seenKeys[parts[3]] { seenKeys[parts[3]] = true currentUsage += amt } } } for id, r := range m.reservations { if r.Identity.OwnerDID == req.Identity.OwnerDID && (m.phases[id] == "reserved" || m.phases[id] == "publishing") { if amt, ok := r.Resources[res]; ok { if !seenKeys[r.Key] { seenKeys[r.Key] = true currentUsage += amt } } } } // identical owner-level claims do not consume quota twice keyExists := false for k := range m.allocations { parts := m.split(k) owner := m.repoOwners[parts[0]] if owner == req.Identity.OwnerDID && parts[1] == string(res) && parts[3] == req.Key { keyExists = true break } } if !keyExists { for id, r := range m.reservations { if r.Identity.OwnerDID == req.Identity.OwnerDID && r.Key == req.Key && (m.phases[id] == "reserved" || m.phases[id] == "publishing") { if _, ok := r.Resources[res]; ok { keyExists = true break } } } } var additional int64 if !keyExists { additional = amount } if currentUsage+additional > userLimit { temporary := true if additional > userLimit { temporary = false } return quota.Reservation{ ID: "", Allowed: false, Temporary: temporary, Reason: quota.ReasonUserLimit, Resource: res, }, nil } } } id := req.ID if id == "" { var b [16]byte _, _ = rand.Read(b[:]) id = hex.EncodeToString(b[:]) } m.reservations[id] = req m.phases[id] = "reserved" repoLimitKey := req.Identity.RepoDID + "/cache_storage_bytes" userLimitKey := req.Identity.OwnerDID + "/cache_storage_bytes" _, hasRepoLimit := m.limits[repoLimitKey] _, hasUserLimit := m.limits[userLimitKey] reason := quota.ReasonUnlimited if hasRepoLimit || hasUserLimit || m.defaultRepo > 0 || m.defaultUser > 0 { reason = quota.ReasonWithinLimit } return quota.Reservation{ ID: id, Allowed: true, Temporary: false, Reason: reason, }, nil } func (m *mockStore) split(k string) []string { idx := m.findResourceIndex(k) if idx != -1 { repo := k[:idx] rest := k[idx+1:] sub := m.splitRest(rest) parts := []string{repo} return append(parts, sub...) } return []string{k, "", "", ""} } func (m *mockStore) findResourceIndex(k string) int { resources := []string{ "/cache_storage_bytes/", "/workflows/", "/vcpus/", "/memory_mib/", "/disk_mib/", } for _, res := range resources { idx := findSubstr(k, res) if idx != -1 { return idx } } return -1 } func findSubstr(s, sub string) int { n := len(s) m := len(sub) if n < m { return -1 } for i := 0; i <= n-m; i++ { if s[i:i+m] == sub { return i } } return -1 } func (m *mockStore) splitRest(r string) []string { idx1 := -1 for i := 0; i < len(r); i++ { if r[i] == '/' { idx1 = i break } } if idx1 == -1 { return []string{r, "", ""} } res := r[:idx1] rest := r[idx1+1:] idx2 := -1 for i := 0; i < len(rest); i++ { if rest[i] == '/' { idx2 = i break } } if idx2 == -1 { return []string{res, rest, ""} } kind := rest[:idx2] key := rest[idx2+1:] return []string{res, kind, key} } func (m *mockStore) BeginCommit(ctx context.Context, id string) error { m.mu.Lock() defer m.mu.Unlock() if _, ok := m.reservations[id]; ok { m.phases[id] = "publishing" } return nil } func (m *mockStore) Commit(ctx context.Context, id string) error { m.mu.Lock() defer m.mu.Unlock() req, ok := m.reservations[id] if !ok { return nil } for res, amt := range req.Resources { allocKey := req.Identity.RepoDID + "/" + string(res) + "/" + string(req.Kind) + "/" + req.Key m.allocations[allocKey] = amt } delete(m.reservations, id) delete(m.phases, id) return nil } func (m *mockStore) Release(ctx context.Context, id string) error { m.mu.Lock() defer m.mu.Unlock() delete(m.reservations, id) delete(m.phases, id) return nil } func (m *mockStore) MetricsSnapshot(ctx context.Context) (quota.MetricsSnapshot, error) { return quota.MetricsSnapshot{}, errors.New("unimplemented") } func (m *mockStore) SetLimit(ctx context.Context, did string, resource string, limit int64) error { m.mu.Lock() defer m.mu.Unlock() key := did + "/" + resource m.limits[key] = limit return nil } func (m *mockStore) UnsetLimit(ctx context.Context, did string, resource string) error { m.mu.Lock() defer m.mu.Unlock() key := did + "/" + resource delete(m.limits, key) return nil } func (m *mockStore) GetLimit(ctx context.Context, did string, resource string) (*quota.Limit, error) { m.mu.Lock() defer m.mu.Unlock() limit, ok := m.limits[did+"/"+resource] if !ok { return nil, nil } return "a.Limit{DID: did, Resource: resource, Limit: limit}, nil } func (m *mockStore) ListLimits(ctx context.Context) ([]quota.Limit, error) { return nil, errors.New("unimplemented") } func (m *mockStore) ListUsage(ctx context.Context) ([]quota.Usage, error) { return nil, errors.New("unimplemented") } func (m *mockStore) Recover(ctx context.Context, liveIDs []string) error { m.mu.Lock() defer m.mu.Unlock() liveMap := make(map[string]bool) for _, id := range liveIDs { liveMap[id] = true } for id, req := range m.reservations { if liveMap[id] { continue } if m.phases[id] == "publishing" { for res, amt := range req.Resources { allocKey := req.Identity.RepoDID + "/" + res + "/" + string(req.Kind) + "/" + req.Key m.allocations[allocKey] = amt } } delete(m.reservations, id) delete(m.phases, id) } return nil } func TestManagerValidation(t *testing.T) { req := quota.ReserveRequest{ Kind: quota.KindWorkflow, Key: "job", Identity: quota.Identity{OwnerDID: "did:web:alice", RepoDID: "did:web:repo"}, Resources: quota.Resources{"gpu_slices": 1}, } if err := quota.Validate(req); err != nil { t.Fatalf("generic resource rejected: %v", err) } for i := range quota.MaxResourcePairs { req.Resources[strings.Repeat("x", i+1)] = 1 } if err := quota.Validate(req); err == nil { t.Fatal("expected too many resources to be rejected") } req.Resources = quota.Resources{strings.Repeat("x", quota.MaxResourceKeyBytes+1): 1} if err := quota.Validate(req); err == nil { t.Fatal("expected an oversized resource name to be rejected") } for _, resource := range []string{strings.Repeat("é", 33), string([]byte{0xff}), "esc\x1b[2J", "new\nline", "tab\tname", "del\x7fname"} { req.Resources = quota.Resources{resource: 1} if err := quota.Validate(req); err == nil { t.Fatalf("expected invalid resource name %q to be rejected", resource) } } req.Resources = quota.Resources{"large": quota.MaxResourceAmount + 1} if err := quota.Validate(req); err == nil { t.Fatal("expected an oversized resource amount to be rejected") } if err := quota.ValidateOverride("did:web:alice", "gpu_slices"); err != nil { t.Fatalf("generic override rejected: %v", err) } } func TestManagerPrecedenceAndDenials(t *testing.T) { store := newMockStore(100, 50) mgr := quota.NewManager(store, 50*time.Millisecond, nil) defer mgr.Close() id1 := quota.Identity{OwnerDID: "did:web:alice", RepoDID: "did:web:alice/repo1"} _, res, err := mgr.TryAcquire(context.Background(), quota.ReserveRequest{ Kind: quota.KindNixCache, Key: "hash1", Resources: quota.Resources{ quota.ResourceCacheStorageBytes: 60, }, Identity: id1, }) if err != nil { t.Fatal(err) } if res.Allowed || res.Temporary { t.Fatalf("expected permanent denial, got %+v", res) } _ = store.SetLimit(context.Background(), id1.RepoDID, quota.ResourceCacheStorageBytes, 200) lease2, res2, err := mgr.TryAcquire(context.Background(), quota.ReserveRequest{ Kind: quota.KindNixCache, Key: "hash2", Resources: quota.Resources{ quota.ResourceCacheStorageBytes: 60, }, Identity: id1, }) if err != nil { t.Fatal(err) } if !res2.Allowed { t.Fatal("expected allowed reservation") } if lease2 == nil { t.Fatal("expected non-nil lease") } ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond) defer cancel() _, err = mgr.Acquire(ctx, quota.ReserveRequest{ Kind: quota.KindNixCache, Key: "hash3", Resources: quota.Resources{ quota.ResourceCacheStorageBytes: 50, }, Identity: id1, }) if !errors.Is(err, context.DeadlineExceeded) { t.Fatalf("expected DeadlineExceeded, got %v", err) } } func TestManagerConcurrencyAndFairness(t *testing.T) { store := newMockStore(100, 100) mgr := quota.NewManager(store, 10*time.Millisecond, nil) defer mgr.Close() id1 := quota.Identity{OwnerDID: "did:web:alice", RepoDID: "did:web:alice/repo1"} res1, err := mgr.Acquire(context.Background(), quota.ReserveRequest{ Kind: quota.KindNixCache, Key: "hash1", Resources: quota.Resources{ quota.ResourceCacheStorageBytes: 60, }, Identity: id1, }) if err != nil { t.Fatal(err) } if res1 == nil { t.Fatal("expected non-nil lease") } _ = store.SetLimit(context.Background(), id1.OwnerDID, quota.ResourceCacheStorageBytes, 90) type acquireResult struct { lease quota.Lease err error } resultA := make(chan acquireResult, 1) resultB := make(chan acquireResult, 1) go func() { lease, err := mgr.Acquire(context.Background(), quota.ReserveRequest{ Kind: quota.KindNixCache, Key: "hashA", Resources: quota.Resources{ quota.ResourceCacheStorageBytes: 50, }, Identity: id1, }) resultA <- acquireResult{lease: lease, err: err} }() time.Sleep(20 * time.Millisecond) go func() { lease, err := mgr.Acquire(context.Background(), quota.ReserveRequest{ Kind: quota.KindNixCache, Key: "hashB", Resources: quota.Resources{ quota.ResourceCacheStorageBytes: 50, }, Identity: id1, }) resultB <- acquireResult{lease: lease, err: err} }() time.Sleep(20 * time.Millisecond) err = mgr.Release(context.Background(), res1.ID()) if err != nil { t.Fatal(err) } var acquiredA acquireResult select { case acquiredA = <-resultA: if acquiredA.err != nil || acquiredA.lease == nil { t.Fatalf("expected waiter A to be allowed, got lease=%v, err=%v", acquiredA.lease, acquiredA.err) } case <-time.After(time.Second): t.Fatal("timed out waiting for waiter A") } select { case acquiredB := <-resultB: t.Fatalf("waiter B completed before waiter A released capacity: lease=%v, err=%v", acquiredB.lease, acquiredB.err) default: } var hashAResID string store.mu.Lock() for id, r := range store.reservations { if r.Key == "hashA" { hashAResID = id break } } store.mu.Unlock() if hashAResID == "" { t.Fatal("could not find reservation ID for hashA") } acquiredA.lease.Release() select { case acquiredB := <-resultB: if acquiredB.err != nil || acquiredB.lease == nil { t.Fatalf("expected waiter B to be allowed, got lease=%v, err=%v", acquiredB.lease, acquiredB.err) } acquiredB.lease.Release() case <-time.After(time.Second): t.Fatal("timed out waiting for waiter B") } } func TestManagerCancellation(t *testing.T) { store := newMockStore(100, 100) mgr := quota.NewManager(store, 50*time.Millisecond, nil) defer mgr.Close() id1 := quota.Identity{OwnerDID: "did:web:alice", RepoDID: "did:web:alice/repo1"} _, _ = mgr.Acquire(context.Background(), quota.ReserveRequest{ Kind: quota.KindNixCache, Key: "hash1", Resources: quota.Resources{ quota.ResourceCacheStorageBytes: 80, }, Identity: id1, }) ctx, cancel := context.WithCancel(context.Background()) var errA error var wg sync.WaitGroup wg.Add(1) go func() { defer wg.Done() _, errA = mgr.Acquire(ctx, quota.ReserveRequest{ Kind: quota.KindNixCache, Key: "hashA", Resources: quota.Resources{ quota.ResourceCacheStorageBytes: 50, }, Identity: id1, }) }() time.Sleep(20 * time.Millisecond) cancel() wg.Wait() if !errors.Is(errA, context.Canceled) { t.Fatalf("expected context Canceled, got %v", errA) } store.mu.Lock() found := false for _, r := range store.reservations { if r.Key == "hashA" { found = true } } store.mu.Unlock() if found { t.Fatal("expected hashA reservation to be cleaned up and not leaked") } } func TestExternalLimitWake(t *testing.T) { store := newMockStore(100, 100) mgr := quota.NewManager(store, 20*time.Millisecond, nil) defer mgr.Close() id1 := quota.Identity{OwnerDID: "did:web:alice", RepoDID: "did:web:alice/repo1"} lease0, err := mgr.Acquire(context.Background(), quota.ReserveRequest{ Kind: quota.KindNixCache, Key: "hash0", Resources: quota.Resources{ quota.ResourceCacheStorageBytes: 80, }, Identity: id1, }) if err != nil { t.Fatal(err) } defer lease0.Release() var allowed bool var eErr error var wg sync.WaitGroup wg.Add(1) go func() { defer wg.Done() lease, e := mgr.Acquire(context.Background(), quota.ReserveRequest{ Kind: quota.KindNixCache, Key: "hash1", Resources: quota.Resources{ quota.ResourceCacheStorageBytes: 30, }, Identity: id1, }) allowed = (e == nil && lease != nil) eErr = e }() time.Sleep(20 * time.Millisecond) _ = store.SetLimit(context.Background(), id1.RepoDID, quota.ResourceCacheStorageBytes, 200) _ = store.SetLimit(context.Background(), id1.OwnerDID, quota.ResourceCacheStorageBytes, 200) wg.Wait() if !allowed || eErr != nil { t.Fatalf("expected waiter to be woken up and allowed, got allowed=%t, err=%v", allowed, eErr) } } func TestWorkflowReservationID(t *testing.T) { id1 := quota.WorkflowReservationID("run-1", "owner", "repo", "rkey", "name") id2 := quota.WorkflowReservationID("run-1", "owner", "repo", "rkey", "name") if id1 != id2 { t.Errorf("expected deterministic IDs, got %q and %q", id1, id2) } if len(id1) != 73 { t.Errorf("expected length 73, got %d for %q", len(id1), id1) } // length-prefixing keeps shifted field boundaries distinct idShift1 := quota.WorkflowReservationID("run-1", "ab", "c", "rkey", "name") idShift2 := quota.WorkflowReservationID("run-1", "a", "bc", "rkey", "name") if idShift1 == idShift2 { t.Errorf("expected different IDs for shifted boundaries, both got %q", idShift1) } idDiff := quota.WorkflowReservationID("run-1", "owner", "repo", "rkey", "name-changed") if id1 == idDiff { t.Errorf("expected different IDs on changed argument, both got %q", id1) } idOtherRun := quota.WorkflowReservationID("run-2", "owner", "repo", "rkey", "name") if id1 == idOtherRun { t.Errorf("expected different IDs for distinct runs, both got %q", id1) } } type waitObserver struct { mu sync.Mutex resources []string } func (o *waitObserver) RecordDecision(string, string, bool, bool, string) {} func (o *waitObserver) SetWaitDepth(string, int64) {} func (o *waitObserver) RecordWait(kind, resource string, duration time.Duration) { o.mu.Lock() defer o.mu.Unlock() o.resources = append(o.resources, resource) } func TestManagerWaitMetricUsesBlockingResource(t *testing.T) { store := newMockStore(0, 0) id := quota.Identity{OwnerDID: "did:web:alice", RepoDID: "did:web:alice/repo"} if err := store.SetLimit(context.Background(), id.RepoDID, quota.ResourceVCPUs, 2); err != nil { t.Fatal(err) } observer := &waitObserver{} mgr := quota.NewManager(store, 5*time.Millisecond, observer) defer mgr.Close() lease, err := mgr.Acquire(context.Background(), quota.ReserveRequest{ Kind: quota.KindWorkflow, Key: "first", Identity: id, Resources: quota.Resources{ quota.ResourceWorkflows: 1, quota.ResourceVCPUs: 2, }, }) if err != nil { t.Fatal(err) } defer lease.Release() ctx, cancel := context.WithTimeout(context.Background(), 30*time.Millisecond) defer cancel() _, err = mgr.Acquire(ctx, quota.ReserveRequest{ Kind: quota.KindWorkflow, Key: "second", Identity: id, Resources: quota.Resources{ quota.ResourceWorkflows: 1, quota.ResourceVCPUs: 1, }, }) if !errors.Is(err, context.DeadlineExceeded) { t.Fatalf("expected deadline exceeded, got %v", err) } observer.mu.Lock() defer observer.mu.Unlock() if len(observer.resources) != 1 || observer.resources[0] != quota.ResourceVCPUs { t.Fatalf("wait resources = %v, want [%s]", observer.resources, quota.ResourceVCPUs) } } type failReleaseStore struct { *mockStore mu sync.Mutex failCounts map[string]int cancelFn context.CancelFunc } func (s *failReleaseStore) Reserve(ctx context.Context, req quota.ReserveRequest) (quota.Reservation, error) { res, err := s.mockStore.Reserve(ctx, req) if err == nil && req.Key == "hash2" && res.Allowed { s.mu.Lock() s.failCounts[res.ID] = 2 if s.cancelFn != nil { s.cancelFn() } s.mu.Unlock() } return res, err } func (s *failReleaseStore) Release(ctx context.Context, id string) error { s.mu.Lock() n := s.failCounts[id] if n > 0 { s.failCounts[id]-- s.mu.Unlock() return errors.New("transient release failure") } s.mu.Unlock() return s.mockStore.Release(ctx, id) } func TestManagerAbandonedReleaseRetry(t *testing.T) { mock := newMockStore(100, 100) ctx, cancel := context.WithCancel(context.Background()) store := &failReleaseStore{ mockStore: mock, failCounts: make(map[string]int), cancelFn: cancel, } mgr := quota.NewManager(store, 100*time.Millisecond, nil) defer mgr.Close() id := quota.Identity{OwnerDID: "did:web:alice", RepoDID: "did:web:alice/repo1"} req1 := quota.ReserveRequest{ Kind: quota.KindNixCache, Key: "hash1", Resources: quota.Resources{ quota.ResourceCacheStorageBytes: 100, }, Identity: id, } lease1, err := mgr.Acquire(context.Background(), req1) if err != nil { t.Fatal(err) } req2 := quota.ReserveRequest{ Kind: quota.KindNixCache, Key: "hash2", Resources: quota.Resources{ quota.ResourceCacheStorageBytes: 50, }, Identity: id, } errChan := make(chan error, 1) go func() { _, err := mgr.Acquire(ctx, req2) errChan <- err }() // wait until the request has joined the queue before releasing capacity time.Sleep(20 * time.Millisecond) // cancellation after a grant must release the abandoned reservation lease1.Release() err = <-errChan if !errors.Is(err, context.Canceled) { t.Fatalf("expected context.Canceled, got %v", err) } mock.mu.Lock() var resID string for id, req := range mock.reservations { if req.Key == "hash2" { resID = id break } } mock.mu.Unlock() if resID == "" { t.Fatal("expected reservation to be created for req2") } time.Sleep(5 * time.Millisecond) mock.mu.Lock() _, exists := mock.reservations[resID] mock.mu.Unlock() if !exists { t.Fatal("expected reservation to still exist due to failed release") } time.Sleep(40 * time.Millisecond) mock.mu.Lock() _, exists = mock.reservations[resID] mock.mu.Unlock() if !exists { t.Fatal("expected reservation to still exist before ticker runs") } deadline := time.Now().Add(1 * time.Second) released := false for time.Now().Before(deadline) { mock.mu.Lock() _, exists = mock.reservations[resID] mock.mu.Unlock() if !exists { released = true break } time.Sleep(10 * time.Millisecond) } if !released { t.Fatal("expected reservation to be released after successful ticker retry") } } type trackingStore struct { *mockStore mu sync.Mutex releaseCalls map[string]int failRelease bool } func (s *trackingStore) Release(ctx context.Context, id string) error { s.mu.Lock() s.releaseCalls[id]++ fail := s.failRelease s.mu.Unlock() if fail { return errors.New("transient release failure") } return s.mockStore.Release(ctx, id) } func TestConcurrentReleaseAttempts(t *testing.T) { mock := newMockStore(100, 100) store := &trackingStore{ mockStore: mock, releaseCalls: make(map[string]int), failRelease: true, } // keep the ticker out of the release-attempt count mgr := quota.NewManager(store, 10*time.Second, nil) defer mgr.Close() id := quota.Identity{OwnerDID: "did:web:alice", RepoDID: "did:web:alice/repo1"} lease, err := mgr.Acquire(context.Background(), quota.ReserveRequest{ Kind: quota.KindNixCache, Key: "hash1", Resources: quota.Resources{ quota.ResourceCacheStorageBytes: 10, }, Identity: id, }) if err != nil { t.Fatal(err) } // concurrent lease release remains idempotent var wg sync.WaitGroup numConcurrent := 10 wg.Add(numConcurrent) for range numConcurrent { go func() { defer wg.Done() lease.Release() }() } wg.Wait() store.mu.Lock() origID := strings.Split(lease.ID(), "#")[0] calls := store.releaseCalls[origID] store.mu.Unlock() if calls < 2 || calls > 3 { t.Fatalf("expected 2 or 3 store.Release calls (1 initial + 1 immediate retry + optional background retry), got %d", calls) } store.mu.Lock() store.failRelease = false store.mu.Unlock() err = mgr.Release(context.Background(), "dummy") if err != nil { t.Fatal(err) } time.Sleep(50 * time.Millisecond) mock.mu.Lock() origID = strings.Split(lease.ID(), "#")[0] _, exists := mock.reservations[origID] mock.mu.Unlock() if exists { t.Fatal("expected lease reservation to be successfully released after retry") } } type fairnessMockStore struct { *mockStore reserveStartChan chan struct{} reserveHoldChan chan struct{} reserveBChan chan struct{} onceStart sync.Once onceB sync.Once } func (s *fairnessMockStore) Reserve(ctx context.Context, req quota.ReserveRequest) (quota.Reservation, error) { switch req.Key { case "hashA": s.onceStart.Do(func() { close(s.reserveStartChan) <-s.reserveHoldChan }) case "hashB": s.onceB.Do(func() { close(s.reserveBChan) }) } return s.mockStore.Reserve(ctx, req) } func TestAcquireFairnessTransition(t *testing.T) { mock := newMockStore(100, 100) store := &fairnessMockStore{ mockStore: mock, reserveStartChan: make(chan struct{}), reserveHoldChan: make(chan struct{}), reserveBChan: make(chan struct{}), } mgr := quota.NewManager(store, 10*time.Millisecond, nil) defer mgr.Close() id := quota.Identity{OwnerDID: "did:web:alice", RepoDID: "did:web:alice/repo1"} lease1, err := mgr.Acquire(context.Background(), quota.ReserveRequest{ Kind: quota.KindNixCache, Key: "hash1", Resources: quota.Resources{ quota.ResourceCacheStorageBytes: 60, }, Identity: id, }) if err != nil { t.Fatal(err) } // hold the earlier request inside the store while a later request arrives errChanA := make(chan error, 1) var leaseA quota.Lease go func() { l, err := mgr.Acquire(context.Background(), quota.ReserveRequest{ Kind: quota.KindNixCache, Key: "hashA", Resources: quota.Resources{ quota.ResourceCacheStorageBytes: 100, }, Identity: id, }) leaseA = l errChanA <- err }() <-store.reserveStartChan go lease1.Release() errChanB := make(chan error, 1) go func() { _, err := mgr.Acquire(context.Background(), quota.ReserveRequest{ Kind: quota.KindNixCache, Key: "hashB", Resources: quota.Resources{ quota.ResourceCacheStorageBytes: 100, }, Identity: id, }) errChanB <- err }() select { case <-store.reserveBChan: close(store.reserveHoldChan) t.Fatal("later request reached the store before the earlier reservation completed") case <-time.After(20 * time.Millisecond): } close(store.reserveHoldChan) select { case err := <-errChanA: if err != nil { t.Fatalf("expected A to succeed, got error: %v", err) } if leaseA == nil { t.Fatal("expected non-nil lease for A") } case <-time.After(500 * time.Millisecond): t.Fatal("timeout waiting for A to finish") } select { case err := <-errChanB: t.Fatalf("later request completed while the earlier request held capacity: %v", err) default: } if leaseA != nil { leaseA.Release() } } func TestManagerFencesDuplicateReservationOwners(t *testing.T) { store := newMockStore(100, 100) mgr := quota.NewManager(store, 0, nil) defer mgr.Close() req := quota.ReserveRequest{ ID: "fixed-reservation", Kind: quota.KindNixCache, Key: "same-object", Identity: quota.Identity{OwnerDID: "did:web:alice", RepoDID: "did:web:alice/repo"}, Resources: quota.Resources{quota.ResourceCacheStorageBytes: 10}, } lease1, res1, err := mgr.TryAcquire(context.Background(), req) if err != nil || lease1 == nil || !res1.Allowed { t.Fatalf("first acquisition = lease %v, reservation %+v, error %v", lease1, res1, err) } lease2, res2, err := mgr.TryAcquire(context.Background(), req) if err != nil { t.Fatal(err) } if lease2 != nil || res2.Allowed || !res2.Temporary { t.Fatalf("duplicate active acquisition = lease %v, reservation %+v, want a temporary denial", lease2, res2) } lease1.Release() lease3, res3, err := mgr.TryAcquire(context.Background(), req) if err != nil || lease3 == nil || !res3.Allowed { t.Fatalf("replacement acquisition = lease %v, reservation %+v, error %v", lease3, res3, err) } if res1.ID == res3.ID { t.Fatalf("replacement owner reused fenced ID %q", res1.ID) } if quota.StorageReservationID(res1.ID) != quota.StorageReservationID(res3.ID) { t.Fatalf("fenced IDs refer to different storage reservations: %q and %q", res1.ID, res3.ID) } if err := mgr.BeginCommit(context.Background(), res1.ID); err == nil { t.Fatal("superseded owner began publication") } if err := mgr.BeginCommit(context.Background(), res3.ID); err != nil { t.Fatalf("current owner could not begin publication: %v", err) } lease1.Release() storageID := quota.StorageReservationID(res3.ID) store.mu.Lock() _, stillReserved := store.reservations[storageID] store.mu.Unlock() if !stillReserved { t.Fatal("superseded owner released the current reservation") } lease3.Release() } func TestManagerQueuesDuplicateBlockingAcquisition(t *testing.T) { store := newMockStore(100, 100) mgr := quota.NewManager(store, 5*time.Millisecond, nil) defer mgr.Close() req := quota.ReserveRequest{ ID: "fixed-reservation", Kind: quota.KindWorkflow, Key: "same-workflow", Identity: quota.Identity{OwnerDID: "did:web:alice", RepoDID: "did:web:alice/repo"}, Resources: quota.Resources{quota.ResourceWorkflows: 1}, } first, _, err := mgr.TryAcquire(context.Background(), req) if err != nil || first == nil { t.Fatalf("first acquisition = lease %v, error %v", first, err) } result := make(chan struct { lease quota.Lease err error }, 1) go func() { lease, err := mgr.Acquire(context.Background(), req) result <- struct { lease quota.Lease err error }{lease: lease, err: err} }() select { case got := <-result: t.Fatalf("duplicate acquisition completed before release: lease %v, error %v", got.lease, got.err) case <-time.After(30 * time.Millisecond): } first.Release() select { case got := <-result: if got.err != nil || got.lease == nil { t.Fatalf("queued acquisition = lease %v, error %v", got.lease, got.err) } got.lease.Release() case <-time.After(time.Second): t.Fatal("queued duplicate acquisition did not complete after release") } }