package quota import ( "context" "crypto/rand" "encoding/hex" "errors" "fmt" "strings" "sync" "time" ) var ErrDenied = errors.New("quota denied permanently") type waiter struct { ctx context.Context req ReserveRequest resource string leaseChan chan Lease errChan chan error start time.Time } type noopLease struct{} func (noopLease) ID() string { return "" } func (noopLease) Release() {} func (noopLease) ReleaseWithError() error { return nil } type managerLease struct { id string manager *Manager mu sync.Mutex released bool } func (l *managerLease) ID() string { return l.id } func (l *managerLease) Release() { l.mu.Lock() if l.released { l.mu.Unlock() return } l.released = true l.mu.Unlock() if err := l.manager.releaseLease(l, false); err != nil { l.manager.releaseOrRetry(l.id) } } func (l *managerLease) ReleaseWithError() error { l.mu.Lock() if l.released { l.mu.Unlock() return nil } l.mu.Unlock() err := l.manager.releaseLease(l, true) if err == nil { l.mu.Lock() l.released = true l.mu.Unlock() } return err } func (m *Manager) releaseLease(l *managerLease, retainOnError bool) error { m.lifecycleMu.Lock() defer m.lifecycleMu.Unlock() origID, _ := parseID(l.id) unlock := m.lockID(origID) defer unlock() m.mu.Lock() if m.activeLeases[origID] != l { m.mu.Unlock() return nil } delete(m.activeLeases, origID) m.mu.Unlock() err := m.store.Release(context.Background(), origID) if err != nil && retainOnError { m.mu.Lock() if _, exists := m.activeLeases[origID]; !exists { m.activeLeases[origID] = l } m.mu.Unlock() } if err == nil { go m.processQueue() } return err } type refLock struct { mu sync.Mutex ref int } type Manager struct { store ReservationStore observer Observer lifecycleMu sync.Mutex mu sync.Mutex queue []*waiter stop chan struct{} pendingReleases []string activeLeases map[string]*managerLease idLocks map[string]*refLock } func NewManager(store ReservationStore, retryInterval time.Duration, observer Observer) *Manager { m := &Manager{ store: store, observer: observer, stop: make(chan struct{}), activeLeases: make(map[string]*managerLease), idLocks: make(map[string]*refLock), } if retryInterval > 0 { m.startTicker(retryInterval) } return m } func (m *Manager) lockID(id string) func() { m.mu.Lock() if m.idLocks == nil { m.idLocks = make(map[string]*refLock) } lk, ok := m.idLocks[id] if !ok { lk = &refLock{} m.idLocks[id] = lk } lk.ref++ m.mu.Unlock() lk.mu.Lock() return func() { lk.mu.Unlock() m.mu.Lock() lk.ref-- if lk.ref == 0 { delete(m.idLocks, id) } m.mu.Unlock() } } func parseID(id string) (string, string) { parts := strings.SplitN(id, "#", 2) if len(parts) == 2 { return parts[0], parts[1] } return id, "" } func StorageReservationID(id string) string { storageID, _ := parseID(id) return storageID } func (m *Manager) Close() { m.mu.Lock() select { case <-m.stop: m.mu.Unlock() return default: } close(m.stop) queueCopy := m.queue m.queue = nil m.mu.Unlock() // fail queued waiters outside the lock kinds := make(map[string]bool) for _, w := range queueCopy { kinds[string(w.req.Kind)] = true select { case w.errChan <- errors.New("manager closed"): default: } m.recordDecisionForResources(w.req, false, false, "no_wait_slot_unavailable") m.recordWait(w, time.Since(w.start)) } for k := range kinds { m.setQueueDepth(k) } } func (m *Manager) startTicker(interval time.Duration) { ticker := time.NewTicker(interval) go func() { for { select { case <-ticker.C: m.processQueue() case <-m.stop: ticker.Stop() return } } }() } func (m *Manager) recordDecision(kind, resource string, allowed, temporary bool, reason string) { if obs := m.observer; obs != nil { obs.RecordDecision(kind, resource, allowed, temporary, reason) } } func (m *Manager) recordDecisionForResources(req ReserveRequest, allowed, temporary bool, reason string) { if obs := m.observer; obs != nil { for res := range req.Resources { obs.RecordDecision(string(req.Kind), res, allowed, temporary, reason) } } } func (m *Manager) recordWait(w *waiter, duration time.Duration) { if obs := m.observer; obs != nil { obs.RecordWait(string(w.req.Kind), w.resource, duration) } } func (m *Manager) setQueueDepth(kind string) { m.mu.Lock() var depth int64 for _, w := range m.queue { if string(w.req.Kind) == kind && w.ctx.Err() == nil { depth++ } } m.mu.Unlock() if obs := m.observer; obs != nil { obs.SetWaitDepth(kind, depth) } } func (m *Manager) hasConflict(req ReserveRequest) bool { for _, w := range m.queue { if w.ctx.Err() != nil { continue } if w.req.Identity.OwnerDID != req.Identity.OwnerDID && w.req.Identity.RepoDID != req.Identity.RepoDID { continue } for res := range req.Resources { if _, ok := w.req.Resources[res]; ok { return true } } } return false } func (m *Manager) TryAcquire(ctx context.Context, req ReserveRequest) (Lease, Reservation, error) { if err := Validate(req); err != nil { return nil, Reservation{}, err } m.mu.Lock() select { case <-m.stop: m.mu.Unlock() return nil, Reservation{}, errors.New("manager closed") default: } if m.hasConflict(req) { m.mu.Unlock() res := Reservation{ Allowed: false, Temporary: true, Reason: "queued", } m.recordDecisionForResources(req, false, true, "queued") return nil, res, nil } m.mu.Unlock() var unlock func() m.lifecycleMu.Lock() defer m.lifecycleMu.Unlock() if req.ID != "" { unlock = m.lockID(req.ID) } res, err := m.store.Reserve(ctx, req) if err != nil { if unlock != nil { unlock() } m.recordDecisionForResources(req, false, false, "store_error") return nil, Reservation{}, err } if res.Allowed { if res.ID == "" { if unlock != nil { unlock() } m.recordDecisionForResources(req, true, false, res.Reason) return noopLease{}, res, nil } if unlock == nil { unlock = m.lockID(res.ID) } token := make([]byte, 4) _, _ = rand.Read(token) gen := hex.EncodeToString(token) modifiedID := fmt.Sprintf("%s#%s", res.ID, gen) l := &managerLease{ id: modifiedID, manager: m, } m.mu.Lock() if _, exists := m.activeLeases[res.ID]; exists { m.mu.Unlock() if unlock != nil { unlock() } res.ID = "" res.Allowed = false res.Temporary = true res.Reason = "queued" m.recordDecisionForResources(req, false, true, res.Reason) return nil, res, nil } if m.activeLeases == nil { m.activeLeases = make(map[string]*managerLease) } m.activeLeases[res.ID] = l m.mu.Unlock() if unlock != nil { unlock() } res.ID = modifiedID m.recordDecisionForResources(req, true, false, res.Reason) return l, res, nil } if unlock != nil { unlock() } m.recordDecision(string(req.Kind), res.Resource, false, res.Temporary, res.Reason) return nil, res, nil } func (m *Manager) Acquire(ctx context.Context, req ReserveRequest) (Lease, error) { if err := Validate(req); err != nil { return nil, err } start := time.Now() m.mu.Lock() select { case <-m.stop: m.mu.Unlock() return nil, errors.New("manager closed") default: } if m.hasConflict(req) { w := &waiter{ ctx: ctx, req: req, leaseChan: make(chan Lease, 1), errChan: make(chan error, 1), start: start, } m.queue = append(m.queue, w) m.mu.Unlock() m.recordDecisionForResources(req, false, true, "queued") m.setQueueDepth(string(req.Kind)) go m.processQueue() return m.wait(ctx, w, start, req) } m.mu.Unlock() m.lifecycleMu.Lock() m.mu.Lock() res, err := m.store.Reserve(ctx, req) if err != nil { m.mu.Unlock() m.lifecycleMu.Unlock() m.recordDecisionForResources(req, false, false, "store_error") return nil, err } if res.Allowed { if res.ID == "" { m.mu.Unlock() m.lifecycleMu.Unlock() m.recordDecisionForResources(req, true, false, res.Reason) return noopLease{}, nil } if _, exists := m.activeLeases[res.ID]; exists { w := &waiter{ ctx: ctx, req: req, resource: res.Resource, leaseChan: make(chan Lease, 1), errChan: make(chan error, 1), start: start, } m.queue = append(m.queue, w) m.mu.Unlock() m.lifecycleMu.Unlock() m.recordDecisionForResources(req, false, true, "queued") m.setQueueDepth(string(req.Kind)) go m.processQueue() return m.wait(ctx, w, start, req) } token := make([]byte, 4) _, _ = rand.Read(token) gen := hex.EncodeToString(token) modifiedID := fmt.Sprintf("%s#%s", res.ID, gen) l := &managerLease{ id: modifiedID, manager: m, } if m.activeLeases == nil { m.activeLeases = make(map[string]*managerLease) } m.activeLeases[res.ID] = l m.mu.Unlock() m.lifecycleMu.Unlock() m.recordDecisionForResources(req, true, false, res.Reason) return l, nil } if !res.Temporary { m.mu.Unlock() m.lifecycleMu.Unlock() m.recordDecision(string(req.Kind), res.Resource, false, false, res.Reason) return nil, fmt.Errorf("%w: %s", ErrDenied, res.Reason) } w := &waiter{ ctx: ctx, req: req, resource: res.Resource, leaseChan: make(chan Lease, 1), errChan: make(chan error, 1), start: start, } m.queue = append(m.queue, w) m.mu.Unlock() m.lifecycleMu.Unlock() m.recordDecision(string(req.Kind), res.Resource, false, true, res.Reason) m.setQueueDepth(string(req.Kind)) go m.processQueue() return m.wait(ctx, w, start, req) } func (m *Manager) wait(ctx context.Context, w *waiter, start time.Time, req ReserveRequest) (Lease, error) { select { case <-ctx.Done(): m.mu.Lock() inQueue := false for i, q := range m.queue { if q == w { m.queue = append(m.queue[:i], m.queue[i+1:]...) inQueue = true break } } m.mu.Unlock() select { case lease := <-w.leaseChan: m.mu.Lock() if m.activeLeases != nil { origID, _ := parseID(lease.ID()) delete(m.activeLeases, origID) } m.mu.Unlock() m.releaseOrRetry(lease.ID()) m.recordDecisionForResources(w.req, false, false, "context_done_after_ready") default: if inQueue { m.recordDecisionForResources(w.req, false, false, "context_done_in_queue") } } m.recordWait(w, time.Since(start)) m.setQueueDepth(string(req.Kind)) return nil, ctx.Err() case err := <-w.errChan: return nil, err case lease := <-w.leaseChan: if err := ctx.Err(); err != nil { m.mu.Lock() if m.activeLeases != nil { origID, _ := parseID(lease.ID()) delete(m.activeLeases, origID) } m.mu.Unlock() m.releaseOrRetry(lease.ID()) m.recordDecisionForResources(w.req, false, false, "context_done_after_ready") m.recordWait(w, time.Since(start)) m.setQueueDepth(string(req.Kind)) return nil, err } return lease, nil } } func (m *Manager) releaseOrRetry(id string) { m.lifecycleMu.Lock() defer m.lifecycleMu.Unlock() origID, _ := parseID(id) unlock := m.lockID(origID) defer unlock() m.mu.Lock() if m.activeLeases != nil { if _, active := m.activeLeases[origID]; active { m.mu.Unlock() return } } m.mu.Unlock() err := m.store.Release(context.Background(), origID) if err != nil { m.mu.Lock() if m.activeLeases != nil { if _, active := m.activeLeases[origID]; active { m.mu.Unlock() return } } select { case <-m.stop: default: m.pendingReleases = append(m.pendingReleases, origID) } m.mu.Unlock() go m.processQueue() } } func (m *Manager) processQueue() { m.lifecycleMu.Lock() defer m.lifecycleMu.Unlock() m.mu.Lock() select { case <-m.stop: m.mu.Unlock() return default: } toRetry := m.pendingReleases m.pendingReleases = nil m.mu.Unlock() var failed []string for _, id := range toRetry { m.mu.Lock() var active bool if m.activeLeases != nil { _, active = m.activeLeases[id] } m.mu.Unlock() if active { continue } if err := m.store.Release(context.Background(), id); err != nil { m.mu.Lock() var activeNow bool if m.activeLeases != nil { _, activeNow = m.activeLeases[id] } if !activeNow { failed = append(failed, id) } m.mu.Unlock() } } m.mu.Lock() select { case <-m.stop: if len(failed) > 0 { m.pendingReleases = append(m.pendingReleases, failed...) } m.mu.Unlock() return default: } if len(failed) > 0 { m.pendingReleases = append(m.pendingReleases, failed...) } blockedOwners := make(map[string]bool) blockedRepos := make(map[string]bool) grantedOwners := make(map[string]bool) kinds := make(map[string]bool) var nextQueue []*waiter var releasesToMake []string for i, w := range m.queue { if w.ctx.Err() != nil { continue } kinds[string(w.req.Kind)] = true ownerBlocked := blockedOwners[w.req.Identity.OwnerDID] repoBlocked := blockedRepos[w.req.Identity.RepoDID] if ownerBlocked || repoBlocked { blockedOwners[w.req.Identity.OwnerDID] = true blockedRepos[w.req.Identity.RepoDID] = true nextQueue = append(nextQueue, w) continue } if grantedOwners[w.req.Identity.OwnerDID] { hasOtherOwner := false for _, other := range m.queue[i+1:] { if other.req.Identity.OwnerDID != w.req.Identity.OwnerDID && other.ctx.Err() == nil { hasOtherOwner = true break } } if hasOtherOwner { nextQueue = append(nextQueue, w) continue } } res, err := m.store.Reserve(w.ctx, w.req) if err != nil { select { case w.errChan <- err: default: } continue } if res.Allowed { var l Lease = noopLease{} if res.ID != "" { if _, exists := m.activeLeases[res.ID]; exists { blockedOwners[w.req.Identity.OwnerDID] = true blockedRepos[w.req.Identity.RepoDID] = true nextQueue = append(nextQueue, w) w.resource = res.Resource continue } token := make([]byte, 4) _, _ = rand.Read(token) modifiedID := fmt.Sprintf("%s#%s", res.ID, hex.EncodeToString(token)) managed := &managerLease{ id: modifiedID, manager: m, } if m.activeLeases == nil { m.activeLeases = make(map[string]*managerLease) } m.activeLeases[res.ID] = managed l = managed } select { case w.leaseChan <- l: grantedOwners[w.req.Identity.OwnerDID] = true m.recordDecisionForResources(w.req, true, false, "allowed_from_queue") m.recordWait(w, time.Since(w.start)) case <-w.ctx.Done(): if res.ID != "" { delete(m.activeLeases, res.ID) releasesToMake = append(releasesToMake, res.ID) } } } else { if !res.Temporary { select { case w.errChan <- fmt.Errorf("%w: %s", ErrDenied, res.Reason): default: } m.recordDecision(string(w.req.Kind), res.Resource, false, false, res.Reason) m.recordWait(w, time.Since(w.start)) } else { blockedOwners[w.req.Identity.OwnerDID] = true blockedRepos[w.req.Identity.RepoDID] = true nextQueue = append(nextQueue, w) w.resource = res.Resource } } } m.queue = nextQueue // update gauges after releasing the queue lock m.mu.Unlock() for k := range kinds { m.setQueueDepth(k) } // release abandoned grants without blocking the queue lock var failedReleases []string for _, id := range releasesToMake { if err := m.store.Release(context.Background(), id); err != nil { failedReleases = append(failedReleases, id) } } if len(failedReleases) > 0 { m.mu.Lock() select { case <-m.stop: default: m.pendingReleases = append(m.pendingReleases, failedReleases...) } m.mu.Unlock() } } func (m *Manager) Release(ctx context.Context, reservationID string) error { m.lifecycleMu.Lock() defer m.lifecycleMu.Unlock() origID, _ := parseID(reservationID) unlock := m.lockID(origID) defer unlock() var owned *managerLease m.mu.Lock() if m.activeLeases != nil { activeL, exists := m.activeLeases[origID] owned = activeL if exists && activeL.id != reservationID { m.mu.Unlock() return nil } if exists { delete(m.activeLeases, origID) } } m.mu.Unlock() err := m.store.Release(ctx, origID) if err != nil { if owned == nil { owned = &managerLease{id: reservationID, manager: m} } m.mu.Lock() if _, exists := m.activeLeases[origID]; !exists { m.activeLeases[origID] = owned } m.mu.Unlock() } if err == nil { go m.processQueue() } return err } func (m *Manager) Commit(ctx context.Context, reservationID string) error { m.lifecycleMu.Lock() defer m.lifecycleMu.Unlock() origID, _ := parseID(reservationID) unlock := m.lockID(origID) defer unlock() var owned *managerLease m.mu.Lock() if m.activeLeases != nil { activeL, exists := m.activeLeases[origID] owned = activeL if exists && activeL.id != reservationID { m.mu.Unlock() return nil } if exists { delete(m.activeLeases, origID) } } m.mu.Unlock() err := m.store.Commit(ctx, origID) if err != nil { if owned == nil { owned = &managerLease{id: reservationID, manager: m} } m.mu.Lock() if _, exists := m.activeLeases[origID]; !exists { m.activeLeases[origID] = owned } m.mu.Unlock() } if err == nil { go m.processQueue() } return err } func (m *Manager) BeginCommit(ctx context.Context, reservationID string) error { m.lifecycleMu.Lock() defer m.lifecycleMu.Unlock() origID, _ := parseID(reservationID) unlock := m.lockID(origID) defer unlock() m.mu.Lock() active, exists := m.activeLeases[origID] if exists && active.id != reservationID { m.mu.Unlock() return errors.New("quota reservation ownership changed") } m.mu.Unlock() return m.store.BeginCommit(ctx, origID) }