Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881package db
import ( "context" "crypto/rand" "database/sql" "encoding/hex" "errors" "fmt" "math" "sort" "strings" "time"
"tangled.org/core/spindle/quota")
type QuotaStore struct { db *DB defaults quota.Defaults}
var _ quota.Store = (*QuotaStore)(nil)
func NewQuotaStore(db *DB, defaults quota.Defaults) *QuotaStore { return &QuotaStore{ db: db, defaults: defaults, }}
func addQuotaUsage(total, amount int64) (int64, error) { if amount < 0 || amount > math.MaxInt64-total { return 0, fmt.Errorf("quota usage overflow: %d + %d", total, amount) } return total + amount, nil}
func sumQuotaUsage(ctx context.Context, conn *sql.Conn, kind, key string, amount int64, query string, args ...any) (total int64, exists bool, err error) { rows, err := conn.QueryContext(ctx, query, args...) if err != nil { return 0, false, err } defer rows.Close()
for rows.Next() { var rowKind, rowKey string var rowAmount int64 if err := rows.Scan(&rowKind, &rowKey, &rowAmount); err != nil { return 0, false, err } total, err = addQuotaUsage(total, rowAmount) if err != nil { return 0, false, err } exists = exists || rowKind == kind && rowKey == key && rowAmount == amount } return total, exists, rows.Err()}
func (d *QuotaStore) Close() error { return d.db.Close()}
func (d *QuotaStore) withTx(ctx context.Context, fn func(conn *sql.Conn) error) (retErr error) { conn, err := d.db.Conn(ctx) if err != nil { return err } defer conn.Close()
_, err = conn.ExecContext(ctx, "BEGIN IMMEDIATE") if err != nil { return err } var committed bool defer func() { if !committed { rollbackCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() if _, rollbackErr := conn.ExecContext(rollbackCtx, "ROLLBACK"); rollbackErr != nil { retErr = errors.Join(retErr, fmt.Errorf("rollback quota transaction: %w", rollbackErr)) } } }()
if err := fn(conn); err != nil { return err }
_, err = conn.ExecContext(ctx, "COMMIT") if err != nil { return err } committed = true return nil}
func (d *QuotaStore) withReadTx(ctx context.Context, fn func(tx *sql.Tx) error) error { tx, err := d.db.BeginTx(ctx, &sql.TxOptions{ReadOnly: true}) if err != nil { return err } defer tx.Rollback()
if err := fn(tx); err != nil { return err }
return tx.Commit()}
func (d *QuotaStore) Reserve(ctx context.Context, req quota.ReserveRequest) (quota.Reservation, error) { if err := quota.Validate(req); err != nil { return quota.Reservation{}, err }
var res quota.Reservation err := d.withTx(ctx, func(conn *sql.Conn) error { _, err := conn.ExecContext(ctx, ` INSERT INTO quota_repo_owners (repo_did, owner_did) VALUES (?, ?) ON CONFLICT(repo_did) DO UPDATE SET owner_did = excluded.owner_did `, req.Identity.RepoDID, req.Identity.OwnerDID) if err != nil { return err }
rows, err := conn.QueryContext(ctx, ` SELECT id, resource, amount FROM quota_reservations WHERE repo_did = ? AND kind = ? AND key = ? AND phase IN ('reserved', 'publishing', 'active') `, req.Identity.RepoDID, string(req.Kind), req.Key) if err != nil { return err } defer rows.Close()
resMapByID := make(map[string]map[string]int64) for rows.Next() { var id, resourceStr string var amount int64 if err := rows.Scan(&id, &resourceStr, &amount); err != nil { return err } if _, ok := resMapByID[id]; !ok { resMapByID[id] = make(map[string]int64) } resMapByID[id][resourceStr] = amount } for id, resources := range resMapByID { if len(resources) == len(req.Resources) { match := true for r, amt := range req.Resources { if resources[r] != amt { match = false break } } if match { res = quota.Reservation{ ID: id, Allowed: true, Temporary: false, Reason: quota.ReasonUnlimited, } return nil } } }
allocRows, err := conn.QueryContext(ctx, ` SELECT resource, amount FROM quota_allocations WHERE repo_did = ? AND kind = ? AND key = ? `, req.Identity.RepoDID, string(req.Kind), req.Key) if err != nil { return err } defer allocRows.Close()
allocResources := make(map[string]int64) for allocRows.Next() { var resourceStr string var amount int64 if err := allocRows.Scan(&resourceStr, &amount); err != nil { return err } allocResources[resourceStr] = amount } if len(allocResources) == len(req.Resources) { match := true for r, amt := range req.Resources { if allocResources[r] != amt { match = false break } } if match { res = quota.Reservation{ ID: "", Allowed: true, Temporary: false, Reason: quota.ReasonUnlimited, } return nil } }
deny := func(reason, resource string, temporary bool) { res = quota.Reservation{ Allowed: false, Temporary: temporary, Reason: reason, Resource: resource, } }
var sortedResources []string for r := range req.Resources { sortedResources = append(sortedResources, r) } sort.Strings(sortedResources)
for _, resource := range sortedResources { amount := req.Resources[resource]
var repoLimit *int64 var maxAmount *int64 err = conn.QueryRowContext(ctx, ` SELECT max_amount FROM quota_limits WHERE did = ? AND resource = ? `, req.Identity.RepoDID, resource).Scan(&maxAmount) if err == nil { if maxAmount != nil { val := *maxAmount repoLimit = &val } } else if errors.Is(err, sql.ErrNoRows) { if d.defaults != nil { if resMap, ok := d.defaults[quota.ScopeRepo]; ok { if val, ok := resMap[resource]; ok && val > 0 { repoLimit = &val } } } } else { return err }
currentRepoUsage, _, err := sumQuotaUsage(ctx, conn, string(req.Kind), req.Key, amount, ` SELECT kind, key, amount FROM ( SELECT kind, key, amount FROM quota_allocations WHERE repo_did = ? AND resource = ? UNION SELECT kind, key, amount FROM quota_reservations WHERE repo_did = ? AND resource = ? AND phase IN ('reserved', 'publishing', 'active') ) `, req.Identity.RepoDID, resource, req.Identity.RepoDID, resource) if err != nil { return err } if amount > math.MaxInt64-currentRepoUsage { deny(quota.ReasonRepoLimit, resource, true) return nil } if repoLimit != nil && *repoLimit >= 0 && amount > *repoLimit-currentRepoUsage { deny(quota.ReasonRepoLimit, resource, amount <= *repoLimit) return nil }
var userLimit *int64 var maxUserAmount *int64 err = conn.QueryRowContext(ctx, ` SELECT max_amount FROM quota_limits WHERE did = ? AND resource = ? `, req.Identity.OwnerDID, resource).Scan(&maxUserAmount) if err == nil { if maxUserAmount != nil { val := *maxUserAmount userLimit = &val } } else if errors.Is(err, sql.ErrNoRows) { if d.defaults != nil { if resMap, ok := d.defaults[quota.ScopeUser]; ok { if val, ok := resMap[resource]; ok && val > 0 { userLimit = &val } } } } else { return err }
currentUserUsage, existsOwner, err := sumQuotaUsage(ctx, conn, string(req.Kind), req.Key, amount, ` SELECT kind, key, amount FROM ( SELECT a.kind, a.key, a.amount FROM quota_allocations a JOIN quota_repo_owners r ON a.repo_did = r.repo_did WHERE r.owner_did = ? AND a.resource = ? UNION SELECT res.kind, res.key, res.amount FROM quota_reservations res WHERE res.owner_did = ? AND res.resource = ? AND res.phase IN ('reserved', 'publishing', 'active') ) `, req.Identity.OwnerDID, resource, req.Identity.OwnerDID, resource) if err != nil { return err }
var additional int64 if !existsOwner { additional = amount } if additional > math.MaxInt64-currentUserUsage { deny(quota.ReasonUserLimit, resource, true) return nil }
// existing reservations survive lower limits if userLimit != nil && *userLimit >= 0 && additional > 0 && additional > *userLimit-currentUserUsage { deny(quota.ReasonUserLimit, resource, additional <= *userLimit) return nil } }
resID := req.ID if resID == "" { var b [16]byte if _, err = rand.Read(b[:]); err != nil { return err } resID = hex.EncodeToString(b[:]) }
phase := "reserved" if req.Kind == quota.KindWorkflow { phase = "active" }
for _, resource := range sortedResources { amount := req.Resources[resource] _, err = conn.ExecContext(ctx, ` INSERT INTO quota_reservations (id, resource, kind, key, amount, repo_did, owner_did, phase, created_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) `, resID, resource, string(req.Kind), req.Key, amount, req.Identity.RepoDID, req.Identity.OwnerDID, phase, time.Now().Unix()) if err != nil { return err } }
hasFiniteLimit := func(scope quota.Scope, did string, resource string) (bool, error) { // the did is unique across repo and user axes var override sql.NullInt64 err := conn.QueryRowContext(ctx, ` SELECT max_amount FROM quota_limits WHERE did = ? AND resource = ? `, did, resource).Scan(&override) switch { case err == nil: return override.Valid, nil case !errors.Is(err, sql.ErrNoRows): return false, err case d.defaults == nil || d.defaults[scope] == nil: return false, nil default: return d.defaults[scope][resource] > 0, nil } }
reason := quota.ReasonUnlimited for _, resource := range sortedResources { repoFinite, err := hasFiniteLimit(quota.ScopeRepo, req.Identity.RepoDID, resource) if err != nil { return err } userFinite, err := hasFiniteLimit(quota.ScopeUser, req.Identity.OwnerDID, resource) if err != nil { return err } if repoFinite || userFinite { reason = quota.ReasonWithinLimit break } }
res = quota.Reservation{ ID: resID, Allowed: true, Temporary: false, Reason: reason, } return nil })
if err != nil { return quota.Reservation{}, err } return res, nil}
func (d *QuotaStore) BeginCommit(ctx context.Context, reservationID string) error { if reservationID == "" { return nil } return d.withTx(ctx, func(conn *sql.Conn) error { var phase string err := conn.QueryRowContext(ctx, ` SELECT phase FROM quota_reservations WHERE id = ? LIMIT 1 `, reservationID).Scan(&phase) if errors.Is(err, sql.ErrNoRows) { return nil } else if err != nil { return err }
if phase == "reserved" { _, err = conn.ExecContext(ctx, ` UPDATE quota_reservations SET phase = 'publishing' WHERE id = ? `, reservationID) return err } return nil })}
func (d *QuotaStore) Commit(ctx context.Context, reservationID string) error { if reservationID == "" { return nil } return d.withTx(ctx, func(conn *sql.Conn) error { rows, err := conn.QueryContext(ctx, ` SELECT resource, kind, key, amount, repo_did FROM quota_reservations WHERE id = ? `, reservationID) if err != nil { return err } defer rows.Close()
type resItem struct { resource string kind string key string amount int64 repo string } var items []resItem for rows.Next() { var item resItem if err := rows.Scan(&item.resource, &item.kind, &item.key, &item.amount, &item.repo); err != nil { return err } items = append(items, item) }
if len(items) == 0 { return nil }
for _, item := range items { _, err = conn.ExecContext(ctx, ` INSERT INTO quota_allocations (repo_did, resource, kind, key, amount) VALUES (?, ?, ?, ?, ?) ON CONFLICT(repo_did, resource, kind, key) DO UPDATE SET amount = excluded.amount `, item.repo, item.resource, item.kind, item.key, item.amount) if err != nil { return err } }
_, err = conn.ExecContext(ctx, ` DELETE FROM quota_reservations WHERE id = ? `, reservationID) return err })}
func (d *QuotaStore) Release(ctx context.Context, reservationID string) error { if reservationID == "" { return nil } return d.withTx(ctx, func(conn *sql.Conn) error { _, err := conn.ExecContext(ctx, ` DELETE FROM quota_reservations WHERE id = ? `, reservationID) return err })}
func (d *QuotaStore) SetLimit(ctx context.Context, did string, resource string, limit int64) error { if err := quota.ValidateOverride(did, resource); err != nil { return err } return d.withTx(ctx, func(conn *sql.Conn) error { var val any if limit <= 0 { val = nil } else { val = limit } _, err := conn.ExecContext(ctx, ` INSERT OR REPLACE INTO quota_limits (did, resource, max_amount) VALUES (?, ?, ?) `, did, resource, val) return err })}
func (d *QuotaStore) UnsetLimit(ctx context.Context, did string, resource string) error { if err := quota.ValidateOverride(did, resource); err != nil { return err } return d.withTx(ctx, func(conn *sql.Conn) error { _, err := conn.ExecContext(ctx, ` DELETE FROM quota_limits WHERE did = ? AND resource = ? `, did, resource) return err })}
func (d *QuotaStore) GetLimit(ctx context.Context, did string, resource string) (*quota.Limit, error) { if err := quota.ValidateOverride(did, resource); err != nil { return nil, err } var limit *quota.Limit err := d.withReadTx(ctx, func(tx *sql.Tx) error { var maxAmt *int64 err := tx.QueryRowContext(ctx, ` SELECT max_amount FROM quota_limits WHERE did = ? AND resource = ? `, did, resource).Scan(&maxAmt) if errors.Is(err, sql.ErrNoRows) { return nil } if err != nil { return err } l := quota.Limit{DID: did, Resource: resource, Limit: -1} if maxAmt != nil { l.Limit = *maxAmt } limit = &l return nil }) if err != nil { return nil, err } return limit, nil}
func (d *QuotaStore) ListLimits(ctx context.Context) ([]quota.Limit, error) { var limits []quota.Limit err := d.withReadTx(ctx, func(tx *sql.Tx) error { rows, err := tx.QueryContext(ctx, ` SELECT did, resource, max_amount FROM quota_limits ORDER BY did, resource `) if err != nil { return err } defer rows.Close()
for rows.Next() { var l quota.Limit var resStr string var maxAmt *int64 if err := rows.Scan(&l.DID, &resStr, &maxAmt); err != nil { return err } l.Resource = resStr if maxAmt != nil { l.Limit = *maxAmt } else { l.Limit = -1 } limits = append(limits, l) } return nil }) if err != nil { return nil, err } return limits, nil}
func (d *QuotaStore) ListUsage(ctx context.Context) ([]quota.Usage, error) { var usages []quota.Usage err := d.withReadTx(ctx, func(tx *sql.Tx) error { collect := func(scope quota.Scope, query string) error { rows, err := tx.QueryContext(ctx, query) if err != nil { return err } defer rows.Close()
for rows.Next() { var did, resource string var amount int64 if err := rows.Scan(&did, &resource, &amount); err != nil { return err } if amount == 0 { continue } if len(usages) == 0 || usages[len(usages)-1].Scope != scope || usages[len(usages)-1].DID != did || usages[len(usages)-1].Resource != resource { usages = append(usages, quota.Usage{Scope: scope, DID: did, Resource: resource}) } usage := &usages[len(usages)-1] usage.Used, err = addQuotaUsage(usage.Used, amount) if err != nil { return fmt.Errorf("%s %s %s: %w", scope, did, resource, err) } } return rows.Err() }
if err := collect(quota.ScopeRepo, ` SELECT repo_did, resource, amount FROM ( SELECT repo_did, resource, kind, key, amount FROM quota_allocations UNION SELECT repo_did, resource, kind, key, amount FROM quota_reservations WHERE phase IN ('reserved', 'publishing', 'active') UNION SELECT repo_did, '`+quota.ResourceWebhooks+`', 'webhook', cast(id as text), 1 FROM webhooks ) ORDER BY repo_did, resource `); err != nil { return err }
return collect(quota.ScopeUser, ` SELECT owner_did, resource, amount FROM ( SELECT owners.owner_did, allocations.resource, allocations.kind, allocations.key, allocations.amount FROM quota_allocations AS allocations JOIN quota_repo_owners AS owners ON allocations.repo_did = owners.repo_did UNION SELECT owner_did, resource, kind, key, amount FROM quota_reservations WHERE phase IN ('reserved', 'publishing', 'active') UNION SELECT coalesce(owners.owner_did, repos.owner) as owner_did, '`+quota.ResourceWebhooks+`', 'webhook', cast(webhooks.id as text), 1 FROM webhooks JOIN repos ON repos.repo_did = webhooks.repo_did LEFT JOIN quota_repo_owners AS owners ON owners.repo_did = webhooks.repo_did ) ORDER BY owner_did, resource `) }) if err != nil { return nil, err } return usages, nil}
func (d *QuotaStore) MetricsSnapshot(ctx context.Context) (quota.MetricsSnapshot, error) { snapshot := quota.MetricsSnapshot{ Usage: map[quota.Scope]map[string]float64{ quota.ScopeUser: make(map[string]float64), quota.ScopeRepo: make(map[string]float64), }, Subjects: map[quota.Scope]map[string]map[string]int64{ quota.ScopeUser: make(map[string]map[string]int64), quota.ScopeRepo: make(map[string]map[string]int64), }, }
resources := make(map[string]struct{}) for _, defaults := range d.defaults { for resource := range defaults { resources[resource] = struct{}{} } }
scopes := []quota.Scope{quota.ScopeUser, quota.ScopeRepo} statuses := []string{"unlimited", "under_limit", "near_limit", "at_limit", "over_limit"}
for _, sc := range scopes { for res := range resources { snapshot.Subjects[sc][res] = make(map[string]int64) for _, st := range statuses { snapshot.Subjects[sc][res][st] = 0 } } }
limits, err := d.ListLimits(ctx) if err != nil { return snapshot, err } usages, err := d.ListUsage(ctx) if err != nil { return snapshot, err }
rows, err := d.db.QueryContext(ctx, `SELECT repo_did, owner_did FROM quota_repo_owners`) if err != nil { return snapshot, err } repoDids := map[string]struct{}{} subjects := map[quota.Scope]map[string]struct{}{ quota.ScopeUser: {}, quota.ScopeRepo: {}, } for rows.Next() { var repoDID, ownerDID string if err := rows.Scan(&repoDID, &ownerDID); err != nil { rows.Close() return snapshot, err } if repoDID != "" { repoDids[repoDID] = struct{}{} subjects[quota.ScopeRepo][repoDID] = struct{}{} } if ownerDID != "" { subjects[quota.ScopeUser][ownerDID] = struct{}{} } } if err := rows.Close(); err != nil { return snapshot, err } if err := rows.Err(); err != nil { return snapshot, err }
type subjectResource struct { scope quota.Scope did string resource string } limitMap := make(map[[2]string]int64) for _, l := range limits { limitMap[[2]string{l.DID, l.Resource}] = l.Limit // unknown dids use the user axis scope := quota.ScopeUser if _, ok := repoDids[l.DID]; ok { scope = quota.ScopeRepo } subjects[scope][l.DID] = struct{}{} for _, sc := range scopes { if _, ok := snapshot.Subjects[sc][l.Resource]; !ok { snapshot.Subjects[sc][l.Resource] = make(map[string]int64) for _, st := range statuses { snapshot.Subjects[sc][l.Resource][st] = 0 } } } }
usageMap := make(map[subjectResource]int64) for _, u := range usages { key := subjectResource{scope: u.Scope, did: u.DID, resource: u.Resource} usageMap[key] = u.Used subjects[u.Scope][u.DID] = struct{}{} snapshot.Usage[u.Scope][u.Resource] += float64(u.Used) for _, sc := range scopes { if _, ok := snapshot.Subjects[sc][u.Resource]; !ok { snapshot.Subjects[sc][u.Resource] = make(map[string]int64) for _, st := range statuses { snapshot.Subjects[sc][u.Resource][st] = 0 } } } }
for _, scope := range scopes { for did := range subjects[scope] { for resource := range snapshot.Subjects[scope] { key := subjectResource{scope: scope, did: did, resource: resource} used := usageMap[key] limit, hasLimit := limitMap[[2]string{did, resource}] if !hasLimit && d.defaults != nil { if defaults, ok := d.defaults[scope]; ok { if value, ok := defaults[resource]; ok && value > 0 { limit = value hasLimit = true } } }
status := "unlimited" if hasLimit && limit >= 0 { status = classifyStatus(used, limit) } snapshot.Subjects[scope][resource][status]++ } } }
return snapshot, nil}
func classifyStatus(used, limit int64) string { if limit < 0 { return "unlimited" } if used > limit { return "over_limit" } if used == limit { return "at_limit" } if used >= limit-limit/5 { return "near_limit" } return "under_limit"}
func (d *QuotaStore) Recover(ctx context.Context, liveIDs []string) error { return d.withTx(ctx, func(conn *sql.Conn) error { _, err := conn.ExecContext(ctx, ` INSERT INTO quota_allocations (repo_did, resource, kind, key, amount) SELECT repo_did, resource, kind, key, amount FROM quota_reservations WHERE phase = 'publishing' ON CONFLICT(repo_did, resource, kind, key) DO UPDATE SET amount = excluded.amount `) if err != nil { return err }
_, err = conn.ExecContext(ctx, ` DELETE FROM quota_reservations WHERE phase IN ('reserved', 'publishing') `) if err != nil { return err }
if len(liveIDs) == 0 { _, err = conn.ExecContext(ctx, ` DELETE FROM quota_reservations WHERE phase = 'active' `) return err }
placeholders := make([]string, len(liveIDs)) args := make([]any, len(liveIDs)) for i, id := range liveIDs { placeholders[i] = "?" args[i] = id } inClause := strings.Join(placeholders, ", ")
query := fmt.Sprintf(` DELETE FROM quota_reservations WHERE phase = 'active' AND id NOT IN (%s) `, inClause) _, err = conn.ExecContext(ctx, query, args...) return err })}
// DeleteQuotaStateForRepo drops every quota row keyed by repo DID (cache// state lives in quota tables since quota-schema); the repo's limit// overrides go too, other dids' rows stay.func (d *DB) DeleteQuotaStateForRepo(ctx context.Context, repoDid string) error { tx, err := d.BeginTx(ctx, nil) if err != nil { return err } defer tx.Rollback()
stmts := []struct { query string args []any }{ {`delete from quota_reservations where repo_did = ?`, []any{repoDid}}, {`delete from quota_allocations where repo_did = ?`, []any{repoDid}}, {`delete from quota_limits where did = ?`, []any{repoDid}}, {`delete from quota_repo_owners where repo_did = ?`, []any{repoDid}}, } for _, s := range stmts { if _, err := tx.ExecContext(ctx, s.query, s.args...); err != nil { return err } } return tx.Commit()}