diff --git a/api/tangled/quotadefs.go b/api/tangled/quotadefs.go new file mode 100644 index 00000000..1c0a51c6 --- /dev/null +++ b/api/tangled/quotadefs.go @@ -0,0 +1,29 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +package tangled + +// schema: sh.tangled.spindle.quota.defs + +const () + +// SpindleQuotaDefs_Limit is a "limit" in the sh.tangled.spindle.quota.defs schema. +type SpindleQuotaDefs_Limit struct { + // did: DID of the user or repository. + Did string `json:"did" cborgen:"did"` + // limit: Limit amount in raw units (-1 if unlimited). + Limit int64 `json:"limit" cborgen:"limit"` + // resource: Resource name. + Resource string `json:"resource" cborgen:"resource"` +} + +// SpindleQuotaDefs_Usage is a "usage" in the sh.tangled.spindle.quota.defs schema. +type SpindleQuotaDefs_Usage struct { + // did: DID of the user or repository. + Did string `json:"did" cborgen:"did"` + // resource: Resource name. + Resource string `json:"resource" cborgen:"resource"` + // scope: Measurement axis the row belongs to. + Scope string `json:"scope" cborgen:"scope"` + // used: Current used amount in raw units. + Used int64 `json:"used" cborgen:"used"` +} diff --git a/api/tangled/quotaget.go b/api/tangled/quotaget.go new file mode 100644 index 00000000..6b3309ff --- /dev/null +++ b/api/tangled/quotaget.go @@ -0,0 +1,37 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +package tangled + +// schema: sh.tangled.spindle.quota.get + +import ( + "context" + + "github.com/bluesky-social/indigo/lex/util" +) + +const ( + SpindleQuotaGetNSID = "sh.tangled.spindle.quota.get" +) + +// SpindleQuotaGet_Output is the output of a sh.tangled.spindle.quota.get call. +type SpindleQuotaGet_Output struct { + Limit *SpindleQuotaDefs_Limit `json:"limit" cborgen:"limit"` +} + +// SpindleQuotaGet calls the XRPC method "sh.tangled.spindle.quota.get". +// +// did: DID of the user or repository. +// resource: Resource name. +func SpindleQuotaGet(ctx context.Context, c util.LexClient, did string, resource string) (*SpindleQuotaGet_Output, error) { + var out SpindleQuotaGet_Output + + params := map[string]interface{}{} + params["did"] = did + params["resource"] = resource + if err := c.LexDo(ctx, util.Query, "", "sh.tangled.spindle.quota.get", params, nil, &out); err != nil { + return nil, err + } + + return &out, nil +} diff --git a/api/tangled/quotalist.go b/api/tangled/quotalist.go new file mode 100644 index 00000000..acf17ea2 --- /dev/null +++ b/api/tangled/quotalist.go @@ -0,0 +1,30 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +package tangled + +// schema: sh.tangled.spindle.quota.list + +import ( + "context" + + "github.com/bluesky-social/indigo/lex/util" +) + +const ( + SpindleQuotaListNSID = "sh.tangled.spindle.quota.list" +) + +// SpindleQuotaList_Output is the output of a sh.tangled.spindle.quota.list call. +type SpindleQuotaList_Output struct { + Limits []*SpindleQuotaDefs_Limit `json:"limits" cborgen:"limits"` +} + +// SpindleQuotaList calls the XRPC method "sh.tangled.spindle.quota.list". +func SpindleQuotaList(ctx context.Context, c util.LexClient) (*SpindleQuotaList_Output, error) { + var out SpindleQuotaList_Output + if err := c.LexDo(ctx, util.Query, "", "sh.tangled.spindle.quota.list", nil, nil, &out); err != nil { + return nil, err + } + + return &out, nil +} diff --git a/api/tangled/quotaset.go b/api/tangled/quotaset.go new file mode 100644 index 00000000..eb72c795 --- /dev/null +++ b/api/tangled/quotaset.go @@ -0,0 +1,36 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +package tangled + +// schema: sh.tangled.spindle.quota.set + +import ( + "context" + + "github.com/bluesky-social/indigo/lex/util" +) + +const ( + SpindleQuotaSetNSID = "sh.tangled.spindle.quota.set" +) + +// SpindleQuotaSet_Input is the input argument to a sh.tangled.spindle.quota.set call. +type SpindleQuotaSet_Input struct { + // did: DID of the user or repository. + Did string `json:"did" cborgen:"did"` + // limit: Override limit in raw units (bytes for cache_storage_bytes, counts or MiB for others). Must be positive; unlimited=true is the way to clear enforcement. Conflicting with unlimited. + Limit *int64 `json:"limit,omitempty" cborgen:"limit,omitempty"` + // resource: Resource name (workflows, vcpus, memory_mib, disk_mib, cache_storage_bytes). + Resource string `json:"resource" cborgen:"resource"` + // unlimited: Set override limit to unlimited. Conflicting with limit. + Unlimited *bool `json:"unlimited,omitempty" cborgen:"unlimited,omitempty"` +} + +// SpindleQuotaSet calls the XRPC method "sh.tangled.spindle.quota.set". +func SpindleQuotaSet(ctx context.Context, c util.LexClient, input *SpindleQuotaSet_Input) error { + if err := c.LexDo(ctx, util.Procedure, "application/json", "sh.tangled.spindle.quota.set", nil, input, nil); err != nil { + return err + } + + return nil +} diff --git a/api/tangled/quotaunset.go b/api/tangled/quotaunset.go new file mode 100644 index 00000000..f09429f7 --- /dev/null +++ b/api/tangled/quotaunset.go @@ -0,0 +1,32 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +package tangled + +// schema: sh.tangled.spindle.quota.unset + +import ( + "context" + + "github.com/bluesky-social/indigo/lex/util" +) + +const ( + SpindleQuotaUnsetNSID = "sh.tangled.spindle.quota.unset" +) + +// SpindleQuotaUnset_Input is the input argument to a sh.tangled.spindle.quota.unset call. +type SpindleQuotaUnset_Input struct { + // did: DID of the user or repository. + Did string `json:"did" cborgen:"did"` + // resource: Resource name to unset. + Resource string `json:"resource" cborgen:"resource"` +} + +// SpindleQuotaUnset calls the XRPC method "sh.tangled.spindle.quota.unset". +func SpindleQuotaUnset(ctx context.Context, c util.LexClient, input *SpindleQuotaUnset_Input) error { + if err := c.LexDo(ctx, util.Procedure, "application/json", "sh.tangled.spindle.quota.unset", nil, input, nil); err != nil { + return err + } + + return nil +} diff --git a/api/tangled/quotausage.go b/api/tangled/quotausage.go new file mode 100644 index 00000000..9614e4ef --- /dev/null +++ b/api/tangled/quotausage.go @@ -0,0 +1,41 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +package tangled + +// schema: sh.tangled.spindle.quota.usage + +import ( + "context" + + "github.com/bluesky-social/indigo/lex/util" +) + +const ( + SpindleQuotaUsageNSID = "sh.tangled.spindle.quota.usage" +) + +// SpindleQuotaUsage_Output is the output of a sh.tangled.spindle.quota.usage call. +type SpindleQuotaUsage_Output struct { + Usages []*SpindleQuotaDefs_Usage `json:"usages" cborgen:"usages"` +} + +// SpindleQuotaUsage calls the XRPC method "sh.tangled.spindle.quota.usage". +// +// did: Filter usage rows by DID. +// scope: Filter usage rows by scope (user or repo). +func SpindleQuotaUsage(ctx context.Context, c util.LexClient, did string, scope string) (*SpindleQuotaUsage_Output, error) { + var out SpindleQuotaUsage_Output + + params := map[string]interface{}{} + if did != "" { + params["did"] = did + } + if scope != "" { + params["scope"] = scope + } + if err := c.LexDo(ctx, util.Query, "", "sh.tangled.spindle.quota.usage", params, nil, &out); err != nil { + return nil, err + } + + return &out, nil +} diff --git a/lexicons/spindle/quota/defs.json b/lexicons/spindle/quota/defs.json new file mode 100644 index 00000000..896edbe0 --- /dev/null +++ b/lexicons/spindle/quota/defs.json @@ -0,0 +1,63 @@ +{ + "lexicon": 1, + "id": "sh.tangled.spindle.quota.defs", + "defs": { + "limit": { + "description": "a per-subject override for one quota resource", + "type": "object", + "required": [ + "did", + "resource", + "limit" + ], + "properties": { + "did": { + "type": "string", + "format": "did", + "description": "DID whose quota is overridden, either a repository or an account owner" + }, + "resource": { + "type": "string", + "description": "quota resource name" + }, + "limit": { + "type": "integer", + "description": "maximum amount in the resource's native units; -1 means unlimited" + } + } + }, + "usage": { + "description": "current usage aggregated for one subject and resource", + "type": "object", + "required": [ + "scope", + "did", + "resource", + "used" + ], + "properties": { + "scope": { + "type": "string", + "knownValues": [ + "user", + "repo" + ], + "description": "aggregation axis for this row: user or repo" + }, + "did": { + "type": "string", + "format": "did", + "description": "DID for the user or repository scope" + }, + "resource": { + "type": "string", + "description": "quota resource name" + }, + "used": { + "type": "integer", + "description": "amount currently allocated or reserved in native resource units" + } + } + } + } +} diff --git a/lexicons/spindle/quota/get.json b/lexicons/spindle/quota/get.json new file mode 100644 index 00000000..46a96755 --- /dev/null +++ b/lexicons/spindle/quota/get.json @@ -0,0 +1,44 @@ +{ + "lexicon": 1, + "id": "sh.tangled.spindle.quota.get", + "defs": { + "main": { + "type": "query", + "description": "return a per-subject quota override for one resource", + "parameters": { + "type": "params", + "required": [ + "did", + "resource" + ], + "properties": { + "did": { + "type": "string", + "format": "did", + "description": "DID of the repository or account owner" + }, + "resource": { + "type": "string", + "description": "quota resource to look up" + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": [ + "limit" + ], + "properties": { + "limit": { + "description": "quota override for the requested DID and resource", + "type": "ref", + "ref": "sh.tangled.spindle.quota.defs#limit" + } + } + } + } + } + } +} diff --git a/lexicons/spindle/quota/list.json b/lexicons/spindle/quota/list.json new file mode 100644 index 00000000..cc9a2197 --- /dev/null +++ b/lexicons/spindle/quota/list.json @@ -0,0 +1,29 @@ +{ + "lexicon": 1, + "id": "sh.tangled.spindle.quota.list", + "defs": { + "main": { + "type": "query", + "description": "list all per-subject quota overrides", + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": [ + "limits" + ], + "properties": { + "limits": { + "description": "quota overrides ordered by DID and resource", + "type": "array", + "items": { + "type": "ref", + "ref": "sh.tangled.spindle.quota.defs#limit" + } + } + } + } + } + } + } +} diff --git a/lexicons/spindle/quota/set.json b/lexicons/spindle/quota/set.json new file mode 100644 index 00000000..27df160d --- /dev/null +++ b/lexicons/spindle/quota/set.json @@ -0,0 +1,39 @@ +{ + "lexicon": 1, + "id": "sh.tangled.spindle.quota.set", + "defs": { + "main": { + "type": "procedure", + "description": "set a quota override for a repository or account owner; values use the resource's native units", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": [ + "did", + "resource" + ], + "properties": { + "did": { + "type": "string", + "format": "did", + "description": "DID of the repository or account owner" + }, + "resource": { + "type": "string", + "description": "quota resource name, such as workflows, vcpus, memory_mib, disk_mib, or cache_storage_bytes" + }, + "limit": { + "type": "integer", + "description": "positive maximum in the resource's native units; cannot be combined with unlimited" + }, + "unlimited": { + "type": "boolean", + "description": "store an unlimited override; cannot be combined with limit" + } + } + } + } + } + } +} diff --git a/lexicons/spindle/quota/unset.json b/lexicons/spindle/quota/unset.json new file mode 100644 index 00000000..d3151dff --- /dev/null +++ b/lexicons/spindle/quota/unset.json @@ -0,0 +1,31 @@ +{ + "lexicon": 1, + "id": "sh.tangled.spindle.quota.unset", + "defs": { + "main": { + "type": "procedure", + "description": "remove the quota override for a repository or account owner and resource", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": [ + "did", + "resource" + ], + "properties": { + "did": { + "type": "string", + "format": "did", + "description": "DID of the repository or account owner" + }, + "resource": { + "type": "string", + "description": "quota resource whose override to remove" + } + } + } + } + } + } +} diff --git a/lexicons/spindle/quota/usage.json b/lexicons/spindle/quota/usage.json new file mode 100644 index 00000000..7637bf47 --- /dev/null +++ b/lexicons/spindle/quota/usage.json @@ -0,0 +1,47 @@ +{ + "lexicon": 1, + "id": "sh.tangled.spindle.quota.usage", + "defs": { + "main": { + "type": "query", + "description": "list non-zero quota usage aggregated by repository and account owner", + "parameters": { + "type": "params", + "properties": { + "scope": { + "type": "string", + "knownValues": [ + "user", + "repo" + ], + "description": "limit results to one aggregation axis: user or repo" + }, + "did": { + "type": "string", + "format": "did", + "description": "limit results to rows for this subject DID" + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": [ + "usages" + ], + "properties": { + "usages": { + "description": "non-zero usage rows, including active reservations and committed allocations", + "type": "array", + "items": { + "type": "ref", + "ref": "sh.tangled.spindle.quota.defs#usage" + } + } + } + } + } + } + } +} diff --git a/spindle/server.go b/spindle/server.go index 86308735..9ec553e1 100644 --- a/spindle/server.go +++ b/spindle/server.go @@ -741,6 +741,7 @@ func (s *Spindle) XrpcRouter() http.Handler { Notifier: s.Notifier(), ServiceAuth: serviceAuth, Trigger: s, + QuotaStore: db.NewQuotaStore(s.db, quota.Defaults{}), } return x.Router() diff --git a/spindle/xrpc/quota.go b/spindle/xrpc/quota.go new file mode 100644 index 00000000..b889649e --- /dev/null +++ b/spindle/xrpc/quota.go @@ -0,0 +1,308 @@ +package xrpc + +import ( + "encoding/json" + "errors" + "fmt" + "net/http" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/api/tangled" + "tangled.org/core/rbac" + "tangled.org/core/spindle/db" + "tangled.org/core/spindle/quota" + xrpcerr "tangled.org/core/xrpc/errors" +) + +func (x *Xrpc) checkSpindleOwner(r *http.Request) (syntax.DID, *xrpcerr.XrpcError, int) { + actorDid, ok := r.Context().Value(ActorDid).(syntax.DID) + if !ok { + err := xrpcerr.MissingActorDidError + return "", &err, http.StatusUnauthorized + } + + isOwner, err := x.Enforcer.IsSpindleOwner(actorDid.String(), rbac.ThisServer) + if err != nil { + x.Logger.ErrorContext(r.Context(), "failed to check spindle owner", "did", actorDid.String(), "err", err) + errObj := xrpcerr.GenericError(err) + return "", &errObj, http.StatusInternalServerError + } + if !isOwner { + x.Logger.ErrorContext(r.Context(), "insufficient permissions", "did", actorDid.String()) + errObj := xrpcerr.AccessControlError(actorDid.String()) + return "", &errObj, http.StatusUnauthorized + } + + return actorDid, nil, 0 +} + +func (x *Xrpc) getQuotaStore() (quota.Store, error) { + if x.QuotaStore != nil { + return x.QuotaStore, nil + } + if x.Db != nil { + return db.NewQuotaStore(x.Db, quota.Defaults{}), nil + } + return nil, errors.New("quota store unavailable") +} + +func (x *Xrpc) SetLimit(w http.ResponseWriter, r *http.Request) { + l := x.Logger + fail := func(e xrpcerr.XrpcError, status int) { + l.ErrorContext(r.Context(), "quota set limit failed", "kind", e.Tag, "error", e.Message) + writeError(w, e, status) + } + + if _, errObj, status := x.checkSpindleOwner(r); errObj != nil { + fail(*errObj, status) + return + } + + var data tangled.SpindleQuotaSet_Input + if err := json.NewDecoder(r.Body).Decode(&data); err != nil { + fail(xrpcerr.GenericError(err), http.StatusBadRequest) + return + } + + if _, err := syntax.ParseDID(data.Did); err != nil { + fail(xrpcerr.GenericError(fmt.Errorf("invalid DID: %s", data.Did)), http.StatusBadRequest) + return + } + + if err := quota.ValidateOverride(data.Did, data.Resource); err != nil { + fail(xrpcerr.GenericError(err), http.StatusBadRequest) + return + } + + hasLimit := data.Limit != nil + hasUnlimited := data.Unlimited != nil && *data.Unlimited + + if !hasLimit && !hasUnlimited { + fail(xrpcerr.GenericError(errors.New("must specify either limit or unlimited")), http.StatusBadRequest) + return + } + if hasLimit && hasUnlimited { + fail(xrpcerr.GenericError(errors.New("cannot specify both limit and unlimited")), http.StatusBadRequest) + return + } + + var limitVal int64 = -1 + if hasLimit { + limitInt := *data.Limit + if limitInt <= 0 { + fail(xrpcerr.GenericError(fmt.Errorf("limit must be positive: %d", limitInt)), http.StatusBadRequest) + return + } + limitVal = limitInt + } + + qs, err := x.getQuotaStore() + if err != nil { + fail(xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + + if err := qs.SetLimit(r.Context(), data.Did, data.Resource, limitVal); err != nil { + fail(xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + + w.WriteHeader(http.StatusOK) +} + +func (x *Xrpc) UnsetLimit(w http.ResponseWriter, r *http.Request) { + l := x.Logger + fail := func(e xrpcerr.XrpcError, status int) { + l.ErrorContext(r.Context(), "quota unset limit failed", "kind", e.Tag, "error", e.Message) + writeError(w, e, status) + } + + if _, errObj, status := x.checkSpindleOwner(r); errObj != nil { + fail(*errObj, status) + return + } + + var data tangled.SpindleQuotaUnset_Input + if err := json.NewDecoder(r.Body).Decode(&data); err != nil { + fail(xrpcerr.GenericError(err), http.StatusBadRequest) + return + } + + if _, err := syntax.ParseDID(data.Did); err != nil { + fail(xrpcerr.GenericError(fmt.Errorf("invalid DID: %s", data.Did)), http.StatusBadRequest) + return + } + + if err := quota.ValidateOverride(data.Did, data.Resource); err != nil { + fail(xrpcerr.GenericError(err), http.StatusBadRequest) + return + } + + qs, err := x.getQuotaStore() + if err != nil { + fail(xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + + if err := qs.UnsetLimit(r.Context(), data.Did, data.Resource); err != nil { + fail(xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + + w.WriteHeader(http.StatusOK) +} + +func (x *Xrpc) GetLimit(w http.ResponseWriter, r *http.Request) { + l := x.Logger + fail := func(e xrpcerr.XrpcError, status int) { + l.ErrorContext(r.Context(), "quota get limit failed", "kind", e.Tag, "error", e.Message) + writeError(w, e, status) + } + + if _, errObj, status := x.checkSpindleOwner(r); errObj != nil { + fail(*errObj, status) + return + } + + did := r.URL.Query().Get("did") + resource := r.URL.Query().Get("resource") + if err := quota.ValidateOverride(did, resource); err != nil { + fail(xrpcerr.GenericError(err), http.StatusBadRequest) + return + } + + qs, err := x.getQuotaStore() + if err != nil { + fail(xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + + lim, err := qs.GetLimit(r.Context(), did, resource) + if err != nil { + fail(xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + if lim == nil { + fail(xrpcerr.NewXrpcError( + xrpcerr.WithTag("LimitNotFound"), + xrpcerr.WithMessage("no override set for this did and resource"), + ), http.StatusNotFound) + return + } + + out := tangled.SpindleQuotaGet_Output{ + Limit: &tangled.SpindleQuotaDefs_Limit{ + Did: lim.DID, + Resource: lim.Resource, + Limit: lim.Limit, + }, + } + + if err := writeJson(w, http.StatusOK, out); err != nil { + fail(xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } +} + +func (x *Xrpc) ListLimits(w http.ResponseWriter, r *http.Request) { + l := x.Logger + fail := func(e xrpcerr.XrpcError, status int) { + l.ErrorContext(r.Context(), "quota list limits failed", "kind", e.Tag, "error", e.Message) + writeError(w, e, status) + } + + if _, errObj, status := x.checkSpindleOwner(r); errObj != nil { + fail(*errObj, status) + return + } + + qs, err := x.getQuotaStore() + if err != nil { + fail(xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + + limits, err := qs.ListLimits(r.Context()) + if err != nil { + fail(xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + + out := tangled.SpindleQuotaList_Output{ + Limits: make([]*tangled.SpindleQuotaDefs_Limit, 0, len(limits)), + } + for _, lim := range limits { + out.Limits = append(out.Limits, &tangled.SpindleQuotaDefs_Limit{ + Did: lim.DID, + Resource: lim.Resource, + Limit: lim.Limit, + }) + } + + if err := writeJson(w, http.StatusOK, out); err != nil { + fail(xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } +} + +func (x *Xrpc) GetUsage(w http.ResponseWriter, r *http.Request) { + l := x.Logger + fail := func(e xrpcerr.XrpcError, status int) { + l.ErrorContext(r.Context(), "quota usage failed", "kind", e.Tag, "error", e.Message) + writeError(w, e, status) + } + + if _, errObj, status := x.checkSpindleOwner(r); errObj != nil { + fail(*errObj, status) + return + } + + scopeFilter := r.URL.Query().Get("scope") + didFilter := r.URL.Query().Get("did") + + if scopeFilter != "" && scopeFilter != string(quota.ScopeUser) && scopeFilter != string(quota.ScopeRepo) { + fail(xrpcerr.GenericError(fmt.Errorf("invalid scope: %s", scopeFilter)), http.StatusBadRequest) + return + } + if didFilter != "" { + if _, err := syntax.ParseDID(didFilter); err != nil { + fail(xrpcerr.GenericError(fmt.Errorf("invalid DID: %s", didFilter)), http.StatusBadRequest) + return + } + } + + qs, err := x.getQuotaStore() + if err != nil { + fail(xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + + usages, err := qs.ListUsage(r.Context()) + if err != nil { + fail(xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + + out := tangled.SpindleQuotaUsage_Output{ + Usages: make([]*tangled.SpindleQuotaDefs_Usage, 0), + } + for _, u := range usages { + if scopeFilter != "" && string(u.Scope) != scopeFilter { + continue + } + if didFilter != "" && u.DID != didFilter { + continue + } + out.Usages = append(out.Usages, &tangled.SpindleQuotaDefs_Usage{ + Scope: string(u.Scope), + Did: u.DID, + Resource: u.Resource, + Used: u.Used, + }) + } + + if err := writeJson(w, http.StatusOK, out); err != nil { + fail(xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } +} diff --git a/spindle/xrpc/quota_test.go b/spindle/xrpc/quota_test.go new file mode 100644 index 00000000..9aee95ce --- /dev/null +++ b/spindle/xrpc/quota_test.go @@ -0,0 +1,571 @@ +package xrpc + +import ( + "bytes" + "context" + "encoding/json" + "log/slog" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/api/tangled" + "tangled.org/core/rbac" + "tangled.org/core/spindle/config" + "tangled.org/core/spindle/db" + "tangled.org/core/spindle/quota" +) + +func setupQuotaTestXrpc(t *testing.T) (*Xrpc, syntax.DID, syntax.DID) { + t.Helper() + d, e := newTestXrpcDB(t) + + ownerDid := syntax.DID("did:plc:spindleowner") + nonOwnerDid := syntax.DID("did:plc:otheruser") + + if err := e.AddSpindle(rbac.ThisServer); err != nil { + t.Fatalf("AddSpindle: %v", err) + } + if err := e.AddSpindleOwner(rbac.ThisServer, ownerDid.String()); err != nil { + t.Fatalf("AddSpindleOwner: %v", err) + } + + qs := db.NewQuotaStore(d, quota.Defaults{}) + + x := &Xrpc{ + Logger: slog.Default(), + Db: d, + Enforcer: e, + Config: &config.Config{}, + QuotaStore: qs, + } + + return x, ownerDid, nonOwnerDid +} + +func sendSetLimit(x *Xrpc, actor syntax.DID, input tangled.SpindleQuotaSet_Input) (*httptest.ResponseRecorder, int) { + body, _ := json.Marshal(input) + req := httptest.NewRequest(http.MethodPost, "/"+tangled.SpindleQuotaSetNSID, bytes.NewReader(body)) + if actor != "" { + ctx := context.WithValue(req.Context(), ActorDid, actor) + req = req.WithContext(ctx) + } + w := httptest.NewRecorder() + x.SetLimit(w, req) + return w, w.Code +} + +func sendUnsetLimit(x *Xrpc, actor syntax.DID, input tangled.SpindleQuotaUnset_Input) (*httptest.ResponseRecorder, int) { + body, _ := json.Marshal(input) + req := httptest.NewRequest(http.MethodPost, "/"+tangled.SpindleQuotaUnsetNSID, bytes.NewReader(body)) + if actor != "" { + ctx := context.WithValue(req.Context(), ActorDid, actor) + req = req.WithContext(ctx) + } + w := httptest.NewRecorder() + x.UnsetLimit(w, req) + return w, w.Code +} + +func sendListLimits(x *Xrpc, actor syntax.DID) (*httptest.ResponseRecorder, int) { + req := httptest.NewRequest(http.MethodGet, "/"+tangled.SpindleQuotaListNSID, nil) + if actor != "" { + ctx := context.WithValue(req.Context(), ActorDid, actor) + req = req.WithContext(ctx) + } + w := httptest.NewRecorder() + x.ListLimits(w, req) + return w, w.Code +} + +func sendGetLimit(x *Xrpc, actor syntax.DID, query string) (*httptest.ResponseRecorder, int) { + path := "/" + tangled.SpindleQuotaGetNSID + if query != "" { + path += "?" + query + } + req := httptest.NewRequest(http.MethodGet, path, nil) + if actor != "" { + ctx := context.WithValue(req.Context(), ActorDid, actor) + req = req.WithContext(ctx) + } + w := httptest.NewRecorder() + x.GetLimit(w, req) + return w, w.Code +} + +func sendGetUsage(x *Xrpc, actor syntax.DID, query string) (*httptest.ResponseRecorder, int) { + path := "/" + tangled.SpindleQuotaUsageNSID + if query != "" { + path += "?" + query + } + req := httptest.NewRequest(http.MethodGet, path, nil) + if actor != "" { + ctx := context.WithValue(req.Context(), ActorDid, actor) + req = req.WithContext(ctx) + } + w := httptest.NewRecorder() + x.GetUsage(w, req) + return w, w.Code +} + +func TestQuota_OwnerAuth(t *testing.T) { + x, ownerDid, nonOwnerDid := setupQuotaTestXrpc(t) + + limitVal := int64(10) + setInput := tangled.SpindleQuotaSet_Input{ + Did: "did:plc:alice", + Resource: "workflows", + Limit: &limitVal, + } + + w, code := sendSetLimit(x, ownerDid, setInput) + if code != http.StatusOK { + t.Fatalf("expected 200 for owner set limit, got %d (body: %s)", code, w.Body.String()) + } + + w, code = sendSetLimit(x, nonOwnerDid, setInput) + if code != http.StatusUnauthorized { + t.Fatalf("expected 401 for non-owner set limit, got %d", code) + } + + w, code = sendSetLimit(x, "", setInput) + if code != http.StatusUnauthorized { + t.Fatalf("expected 401 for missing actor set limit, got %d", code) + } + + w, code = sendListLimits(x, ownerDid) + if code != http.StatusOK { + t.Fatalf("expected 200 for owner list limits, got %d", code) + } + w, code = sendListLimits(x, nonOwnerDid) + if code != http.StatusUnauthorized { + t.Fatalf("expected 401 for non-owner list limits, got %d", code) + } + w, code = sendListLimits(x, "") + if code != http.StatusUnauthorized { + t.Fatalf("expected 401 for missing actor list limits, got %d", code) + } + + w, code = sendGetUsage(x, ownerDid, "") + if code != http.StatusOK { + t.Fatalf("expected 200 for owner usage, got %d", code) + } + w, code = sendGetUsage(x, nonOwnerDid, "") + if code != http.StatusUnauthorized { + t.Fatalf("expected 401 for non-owner usage, got %d", code) + } + w, code = sendGetUsage(x, "", "") + if code != http.StatusUnauthorized { + t.Fatalf("expected 401 for missing actor usage, got %d", code) + } + + unsetInput := tangled.SpindleQuotaUnset_Input{ + Did: "did:plc:alice", + Resource: "workflows", + } + w, code = sendUnsetLimit(x, nonOwnerDid, unsetInput) + if code != http.StatusUnauthorized { + t.Fatalf("expected 401 for non-owner unset limit, got %d", code) + } + w, code = sendUnsetLimit(x, "", unsetInput) + if code != http.StatusUnauthorized { + t.Fatalf("expected 401 for missing actor unset limit, got %d", code) + } + + limits, err := x.QuotaStore.ListLimits(t.Context()) + if err != nil { + t.Fatal(err) + } + var aliceLimit *quota.Limit + for i, l := range limits { + if l.DID == "did:plc:alice" && l.Resource == "workflows" { + aliceLimit = &limits[i] + } + } + if aliceLimit == nil || aliceLimit.Limit != limitVal { + t.Fatalf("alice's limit changed by rejected unsets: %+v", aliceLimit) + } + + w, code = sendUnsetLimit(x, ownerDid, unsetInput) + if code != http.StatusOK { + t.Fatalf("expected 200 for owner unset limit, got %d", code) + } +} + +func TestQuota_ValidationErrors(t *testing.T) { + x, ownerDid, _ := setupQuotaTestXrpc(t) + + limitVal := int64(10) + unlimitedTrue := true + + tests := []struct { + name string + input tangled.SpindleQuotaSet_Input + }{ + { + name: "invalid DID", + input: tangled.SpindleQuotaSet_Input{ + Did: "alice", + Resource: "workflows", + Limit: &limitVal, + }, + }, + { + name: "malformed DID did:", + input: tangled.SpindleQuotaSet_Input{ + Did: "did:", + Resource: "workflows", + Limit: &limitVal, + }, + }, + { + name: "malformed DID notadid", + input: tangled.SpindleQuotaSet_Input{ + Did: "notadid", + Resource: "workflows", + Limit: &limitVal, + }, + }, + { + name: "empty resource", + input: tangled.SpindleQuotaSet_Input{ + Did: "did:plc:alice", + Resource: "", + Limit: &limitVal, + }, + }, + { + name: "resource too long", + input: tangled.SpindleQuotaSet_Input{ + Did: "did:plc:alice", + Resource: strings.Repeat("a", 65), + Limit: &limitVal, + }, + }, + { + name: "missing limit and unlimited", + input: tangled.SpindleQuotaSet_Input{ + Did: "did:plc:alice", + Resource: "workflows", + }, + }, + { + name: "conflicting limit and unlimited", + input: tangled.SpindleQuotaSet_Input{ + Did: "did:plc:alice", + Resource: "workflows", + Limit: &limitVal, + Unlimited: &unlimitedTrue, + }, + }, + { + name: "negative limit", + input: func() tangled.SpindleQuotaSet_Input { + neg := int64(-5) + return tangled.SpindleQuotaSet_Input{ + Did: "did:plc:alice", + Resource: "workflows", + Limit: &neg, + } + }(), + }, + { + name: "zero limit", + input: func() tangled.SpindleQuotaSet_Input { + zero := int64(0) + return tangled.SpindleQuotaSet_Input{ + Did: "did:plc:alice", + Resource: "workflows", + Limit: &zero, + } + }(), + }, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + w, code := sendSetLimit(x, ownerDid, tc.input) + if code != http.StatusBadRequest { + t.Fatalf("expected 400 for %s, got %d (body: %s)", tc.name, code, w.Body.String()) + } + }) + } + + unsetTests := []struct { + name string + input tangled.SpindleQuotaUnset_Input + }{ + { + name: "invalid did unset", + input: tangled.SpindleQuotaUnset_Input{ + Did: "not-a-did", + Resource: "workflows", + }, + }, + { + name: "malformed DID did: unset", + input: tangled.SpindleQuotaUnset_Input{ + Did: "did:", + Resource: "workflows", + }, + }, + { + name: "malformed DID notadid unset", + input: tangled.SpindleQuotaUnset_Input{ + Did: "notadid", + Resource: "workflows", + }, + }, + { + name: "empty resource unset", + input: tangled.SpindleQuotaUnset_Input{ + Did: "did:plc:alice", + Resource: "", + }, + }, + } + + for _, tc := range unsetTests { + t.Run(tc.name, func(t *testing.T) { + w, code := sendUnsetLimit(x, ownerDid, tc.input) + if code != http.StatusBadRequest { + t.Fatalf("expected 400 for %s, got %d (body: %s)", tc.name, code, w.Body.String()) + } + }) + } + + usageTests := []struct { + name string + query string + }{ + { + name: "invalid scope query", + query: "scope=invalid", + }, + { + name: "invalid did query", + query: "did=notadid", + }, + { + name: "malformed DID did: query", + query: "did=did:", + }, + { + name: "malformed DID notadid query", + query: "did=notadid", + }, + } + + for _, tc := range usageTests { + t.Run(tc.name, func(t *testing.T) { + w, code := sendGetUsage(x, ownerDid, tc.query) + if code != http.StatusBadRequest { + t.Fatalf("expected 400 for %s, got %d (body: %s)", tc.name, code, w.Body.String()) + } + }) + } +} + +func TestQuota_RoundTrip(t *testing.T) { + x, ownerDid, _ := setupQuotaTestXrpc(t) + + limit1 := int64(10) + limit2 := int64(524288000) // 500 MiB in raw bytes + unlimited := true + + w, code := sendSetLimit(x, ownerDid, tangled.SpindleQuotaSet_Input{ + Did: "did:plc:alice", + Resource: "workflows", + Limit: &limit1, + }) + if code != http.StatusOK { + t.Fatalf("expected 200 for set workflows, got %d (body: %s)", code, w.Body.String()) + } + + w, code = sendSetLimit(x, ownerDid, tangled.SpindleQuotaSet_Input{ + Did: "did:plc:alice", + Resource: "cache_storage_bytes", + Limit: &limit2, + }) + if code != http.StatusOK { + t.Fatalf("expected 200 for set cache_storage_bytes, got %d", code) + } + + w, code = sendSetLimit(x, ownerDid, tangled.SpindleQuotaSet_Input{ + Did: "did:plc:repo123", + Resource: "vcpus", + Unlimited: &unlimited, + }) + if code != http.StatusOK { + t.Fatalf("expected 200 for set repo vcpus unlimited, got %d", code) + } + + w, code = sendListLimits(x, ownerDid) + if code != http.StatusOK { + t.Fatalf("expected 200 for list limits, got %d", code) + } + var listOut tangled.SpindleQuotaList_Output + if err := json.Unmarshal(w.Body.Bytes(), &listOut); err != nil { + t.Fatalf("failed to unmarshal list limits output: %v", err) + } + if len(listOut.Limits) != 3 { + t.Fatalf("expected 3 limits, got %d (%+v)", len(listOut.Limits), listOut.Limits) + } + + foundUserWorkflows := false + foundUserCache := false + foundRepoVcpus := false + for _, l := range listOut.Limits { + if l.Did == "did:plc:alice" && l.Resource == "workflows" { + if l.Limit != 10 { + t.Errorf("expected limit 10, got %d", l.Limit) + } + foundUserWorkflows = true + } + if l.Did == "did:plc:alice" && l.Resource == "cache_storage_bytes" { + if l.Limit != 524288000 { + t.Errorf("expected limit 524288000, got %d", l.Limit) + } + foundUserCache = true + } + if l.Did == "did:plc:repo123" && l.Resource == "vcpus" { + if l.Limit != -1 { + t.Errorf("expected limit -1 for unlimited, got %d", l.Limit) + } + foundRepoVcpus = true + } + } + if !foundUserWorkflows || !foundUserCache || !foundRepoVcpus { + t.Fatalf("did not find all expected limits: userWorkflows=%v, userCache=%v, repoVcpus=%v", + foundUserWorkflows, foundUserCache, foundRepoVcpus) + } + + res, err := x.QuotaStore.Reserve(context.Background(), quota.ReserveRequest{ + ID: "res-1", + Kind: quota.KindWorkflow, + Key: "job-1", + Identity: quota.Identity{ + OwnerDID: "did:plc:alice", + RepoDID: "did:plc:repo123", + }, + Resources: quota.Resources{ + quota.ResourceWorkflows: 2, + }, + }) + if err != nil || !res.Allowed { + t.Fatalf("failed to reserve test usage: %v (res: %+v)", err, res) + } + + w, code = sendGetUsage(x, ownerDid, "") + if code != http.StatusOK { + t.Fatalf("expected 200 for usage query, got %d", code) + } + var usageOut tangled.SpindleQuotaUsage_Output + if err := json.Unmarshal(w.Body.Bytes(), &usageOut); err != nil { + t.Fatalf("failed to unmarshal usage output: %v", err) + } + if len(usageOut.Usages) == 0 { + t.Fatalf("expected at least 1 usage row, got 0") + } + + w, code = sendGetUsage(x, ownerDid, "scope=user&did=did:plc:alice") + if code != http.StatusOK { + t.Fatalf("expected 200 for filtered usage, got %d", code) + } + var userUsageOut tangled.SpindleQuotaUsage_Output + if err := json.Unmarshal(w.Body.Bytes(), &userUsageOut); err != nil { + t.Fatalf("failed to unmarshal filtered usage: %v", err) + } + if len(userUsageOut.Usages) != 1 || userUsageOut.Usages[0].Used != 2 { + t.Fatalf("unexpected filtered user usage: %+v", userUsageOut.Usages) + } + + w, code = sendGetUsage(x, ownerDid, "scope=user&did=did:plc:nonexistent") + if code != http.StatusOK { + t.Fatalf("expected 200 for empty filtered usage, got %d", code) + } + var emptyUsageOut tangled.SpindleQuotaUsage_Output + if err := json.Unmarshal(w.Body.Bytes(), &emptyUsageOut); err != nil { + t.Fatalf("failed to unmarshal empty usage output: %v", err) + } + if len(emptyUsageOut.Usages) != 0 { + t.Fatalf("expected 0 usages for nonexistent did, got %d", len(emptyUsageOut.Usages)) + } + + w, code = sendUnsetLimit(x, ownerDid, tangled.SpindleQuotaUnset_Input{ + Did: "did:plc:alice", + Resource: "workflows", + }) + if code != http.StatusOK { + t.Fatalf("expected 200 for unset limit, got %d", code) + } + + w, code = sendListLimits(x, ownerDid) + if code != http.StatusOK { + t.Fatalf("expected 200 for list limits after unset, got %d", code) + } + var afterUnsetOut tangled.SpindleQuotaList_Output + if err := json.Unmarshal(w.Body.Bytes(), &afterUnsetOut); err != nil { + t.Fatalf("failed to unmarshal list limits output: %v", err) + } + if len(afterUnsetOut.Limits) != 2 { + t.Fatalf("expected 2 limits after unset, got %d", len(afterUnsetOut.Limits)) + } + for _, l := range afterUnsetOut.Limits { + if l.Did == "did:plc:alice" && l.Resource == "workflows" { + t.Fatalf("expected workflows limit to be removed, but still present") + } + } +} + +// unset and unlimited are different answers +func TestGetLimit(t *testing.T) { + x, ownerDid, _ := setupQuotaTestXrpc(t) + did := "did:plc:alice" + + if _, code := sendGetLimit(x, ownerDid, "did="+did+"&resource=workflows"); code != http.StatusNotFound { + t.Fatalf("expected 404 for an unset override, got %d", code) + } + + limit := int64(7) + if _, code := sendSetLimit(x, ownerDid, tangled.SpindleQuotaSet_Input{Did: did, Resource: "workflows", Limit: &limit}); code != http.StatusOK { + t.Fatal("set limit failed") + } + + w, code := sendGetLimit(x, ownerDid, "did="+did+"&resource=workflows") + if code != http.StatusOK { + t.Fatalf("expected 200 for a set override, got %d", code) + } + var out tangled.SpindleQuotaGet_Output + if err := json.Unmarshal(w.Body.Bytes(), &out); err != nil { + t.Fatalf("decode: %v", err) + } + if out.Limit == nil || out.Limit.Limit != 7 || out.Limit.Did != did || out.Limit.Resource != "workflows" { + t.Fatalf("expected the override echoed back, got %+v", out.Limit) + } + + unlimited := true + if _, code := sendSetLimit(x, ownerDid, tangled.SpindleQuotaSet_Input{Did: did, Resource: "workflows", Unlimited: &unlimited}); code != http.StatusOK { + t.Fatal("set unlimited failed") + } + w, code = sendGetLimit(x, ownerDid, "did="+did+"&resource=workflows") + if code != http.StatusOK { + t.Fatalf("expected 200 for an unlimited override, got %d", code) + } + out = tangled.SpindleQuotaGet_Output{} + if err := json.Unmarshal(w.Body.Bytes(), &out); err != nil { + t.Fatalf("decode: %v", err) + } + if out.Limit == nil || out.Limit.Limit != -1 { + t.Fatalf("expected unlimited reported as -1, got %+v", out.Limit) + } + + if _, code := sendGetLimit(x, ownerDid, "did=notadid&resource=workflows"); code != http.StatusBadRequest { + t.Fatalf("expected 400 for a malformed did, got %d", code) + } + if _, code := sendGetLimit(x, ownerDid, "did="+did); code != http.StatusBadRequest { + t.Fatalf("expected 400 for a missing resource, got %d", code) + } +} diff --git a/spindle/xrpc/xrpc.go b/spindle/xrpc/xrpc.go index 4d390d96..5bddf7d7 100644 --- a/spindle/xrpc/xrpc.go +++ b/spindle/xrpc/xrpc.go @@ -20,6 +20,7 @@ import ( "tangled.org/core/spindle/config" "tangled.org/core/spindle/db" "tangled.org/core/spindle/models" + "tangled.org/core/spindle/quota" "tangled.org/core/spindle/secrets" xrpcerr "tangled.org/core/xrpc/errors" "tangled.org/core/xrpc/serviceauth" @@ -62,6 +63,7 @@ type Xrpc struct { Notifier *notifier.Notifier ServiceAuth *serviceauth.ServiceAuth Trigger PipelineTrigger + QuotaStore quota.Store } func (x *Xrpc) Router() http.Handler { @@ -75,6 +77,11 @@ func (x *Xrpc) Router() http.Handler { r.Get("/"+tangled.RepoListSecretsNSID, x.ListSecrets) r.Post("/"+tangled.CiCancelPipelineNSID, x.CancelPipeline) r.Post("/"+tangled.CiTriggerPipelineNSID, x.TriggerPipeline) + r.Post("/"+tangled.SpindleQuotaSetNSID, x.SetLimit) + r.Post("/"+tangled.SpindleQuotaUnsetNSID, x.UnsetLimit) + r.Get("/"+tangled.SpindleQuotaGetNSID, x.GetLimit) + r.Get("/"+tangled.SpindleQuotaListNSID, x.ListLimits) + r.Get("/"+tangled.SpindleQuotaUsageNSID, x.GetUsage) }) // service query endpoints (no auth required)