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