Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
Go
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205package 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") }}