Unverified Commit e5747498 authored by Lemon-miaow's avatar Lemon-miaow
Browse files

feat(api): enforce CPU/memory/storage quotas (spec §9.3, §22)

Add per-user resource quota enforcement across all four dimensions:
max_servers, max_cpu_milli, max_memory_mb, and max_storage_gb.

- Migration 0013: add cached_cpu_milli, cached_memory_mb, cached_storage_mb
  columns to servers table for pure-SQL per-owner aggregation
- SeedServer now writes resource cache alongside server row
- QuotaCheck replaces QuotaAvailable at claim time, checking all four caps
  against the owning user's cumulative usage
- handlePatchServer checks owner's quota before allowing memory/resource
  changes on owned servers; unowned servers skip the gate
- handleInternalClaim mirrors the full quota check
- UpdateServerResources keeps the cache in sync after spec mutations
- Reaper zeros resource cache on ReleaseWorld so released resources
  are not counted against a former owner
- quantityToMilli/quantityToMB helpers convert K8s quantities to
  quota-comparable integers

19 test packages pass.
parent 91bfa27e
Loading
Loading
Loading
Loading
+14 −1
Changes for internal/api/api_test.go: 14 added lines, 1 removed line.
Original line number Diff line number Diff line
@@ -209,6 +209,19 @@ func (f *fakeRepo) ServerByName(_ context.Context, n string) (*ServerRecord, err
}
func (f *fakeRepo) IsLinked(_ context.Context, u string) (bool, error)       { return f.linked[u], nil }
func (f *fakeRepo) QuotaAvailable(_ context.Context, u string) (bool, error) { return f.quota[u], nil }

func (f *fakeRepo) QuotaCheck(_ context.Context, userID string, _ string, _ ResourceSpec) (bool, error) {
	// For hermetic tests, QuotaCheck delegates to the same QuotaAvailable
	// store — tests that care about per-dimension checks should use
	// fakeQuotas with direct inspection.
	return f.QuotaAvailable(nil, userID)
}

func (f *fakeRepo) UpdateServerResources(_ context.Context, _ string, _, _, _ int) error { return nil }

func (f *fakeRepo) ServerResources(_ context.Context, _ string) (ResourceSpec, error) {
	return ResourceSpec{}, nil
}
func (f *fakeRepo) CreateLinkCode(_ context.Context, code, mcUUID, authSource string, expiresAt time.Time) error {
	f.linkCodes[code] = fakeLinkCode{mcUUID: mcUUID, authSource: authSource, expiresAt: expiresAt}
	return nil
@@ -532,7 +545,7 @@ func (f *fakeRepo) ServerOwners(_ context.Context) (map[string]string, error) {
	}
	return f.owners, nil
}
func (f *fakeRepo) SeedServer(_ context.Context, name, subdomain string) error {
func (f *fakeRepo) SeedServer(_ context.Context, name, subdomain string, _, _, _ int) error {
	if f.seedErr != nil {
		return f.seedErr
	}
+8 −4
Changes for internal/api/handlers_internal.go: 8 added lines, 4 removed lines.
Original line number Diff line number Diff line
@@ -232,10 +232,14 @@ func (a *API) handleInternalClaim(w http.ResponseWriter, r *http.Request) {
		return
	}

	// ② quota gate, evaluated before the ownership write (mirrors handleClaim). It
	// shares handleClaim's quota TOCTOU KNOWN-LIMITATION — see QuotaAvailable (audit
	// #4, ENV-blocked).
	ok, err := a.Repo.QuotaAvailable(r.Context(), userID)
	// ② quota gate, evaluated before the ownership write (mirrors handleClaim).
	// All four dimensions (servers, CPU, memory, storage) are checked.
	res, err := a.Repo.ServerResources(r.Context(), name)
	if err != nil {
		writeError(w, r, err)
		return
	}
	ok, err := a.Repo.QuotaCheck(r.Context(), userID, "", res)
	if err != nil {
		writeError(w, r, err)
		return
+1 −0
Changes for internal/api/handlers_patch_test.go: 1 added line, 0 removed lines.
Original line number Diff line number Diff line
@@ -17,6 +17,7 @@ func newPatchAPI() (*API, *fakeRepo, *fakeCluster, *fakeBuilder) {
	cl.byName["survival"] = &ServerInfo{Name: "survival", Subdomain: "survival",
		AutostartPolicy: string(v1alpha1.AutostartOwnerOnly),
		DesiredState:    string(v1alpha1.DesiredStopped), Phase: string(v1alpha1.PhaseStopped)}
	repo.byName["survival"] = &ServerRecord{Name: "survival", Subdomain: "survival"}
	return api, repo, cl, fb
}

+79 −5
Changes for internal/api/handlers_user.go: 79 added lines, 5 removed lines.
Original line number Diff line number Diff line
@@ -120,10 +120,16 @@ func (a *API) handleClaim(w http.ResponseWriter, r *http.Request) {
		return
	}

	// ② quota gate, evaluated before the ownership write. This gate and ③ are two
	// separate statements, not one transaction — see the quota TOCTOU KNOWN-LIMITATION
	// on QuotaAvailable (audit #4, ENV-blocked: needs real Postgres to close/verify).
	ok, err := a.Repo.QuotaAvailable(r.Context(), p.UserID)
	// ② quota gate, evaluated before the ownership write. All four dimensions
	// (servers, CPU, memory, storage) are checked against the user's quota caps
	// using the PG-resident resource cache (spec §9.3, §22). The server being
	// claimed has owner_id=NULL so it is not yet in the per-owner aggregate.
	res, err := a.Repo.ServerResources(r.Context(), name)
	if err != nil {
		writeError(w, r, err)
		return
	}
	ok, err := a.Repo.QuotaCheck(r.Context(), p.UserID, "", res)
	if err != nil {
		writeError(w, r, err)
		return
@@ -376,7 +382,13 @@ func (a *API) handleCreateServer(w http.ResponseWriter, r *http.Request) {
	// Seed the business rows FIRST (servers + alias). ClaimServer needs the row,
	// so a CRD-only server would be unclaimable. PG-first means a later CRD
	// failure leaves a claimable ghost row — acceptable, not transactional.
	if err := a.Repo.SeedServer(r.Context(), body.Name, body.Subdomain); err != nil {
	// The resource cache (cpuMilli, memoryMB, storageMB) is seeded alongside so
	// QuotaCheck can aggregate per-owner usage without cross-system CRD reads.
	cpuMilli := quantityToMilli(resources.Limits[corev1.ResourceCPU])
	memMB := quantityToMB(resources.Limits[corev1.ResourceMemory])
	storQ, _ := resource.ParseQuantity(storage)
	storMB := quantityToMB(storQ)
	if err := a.Repo.SeedServer(r.Context(), body.Name, body.Subdomain, cpuMilli, memMB, storMB); err != nil {
		if errors.Is(err, ErrConflict) {
			writeError(w, r, newError(http.StatusConflict, "subdomain_taken",
				"subdomain %q is already in use", body.Subdomain))
@@ -676,6 +688,10 @@ func (a *API) handlePatchServer(w http.ResponseWriter, r *http.Request) {
	// override block is meaningless without that base. A resources-only patch has no
	// base ceiling to widen (this endpoint does not read the current spec back), so
	// it is rejected rather than guessed.
	var (
		newResources   corev1.ResourceRequirements
		resUpdated     bool
	)
	if body.Memory != nil {
		javaMemory, resources, err := resolveResources(*body.Memory, body.Resources)
		if err != nil {
@@ -688,16 +704,52 @@ func (a *API) handlePatchServer(w http.ResponseWriter, r *http.Request) {
		if body.Resources != nil {
			changed = append(changed, "resources")
		}
		newResources = resources
		resUpdated = true
	} else if body.Resources != nil {
		writeError(w, r, newError(http.StatusBadRequest, "bad_request",
			"resources overrides require memory to be set in the same patch"))
		return
	}

	// Resource-cache consistency + quota enforcement (spec §9.3 / §22): every
	// resource-mutating patch must update the cached columns so QuotaCheck
	// can aggregate per-owner usage without cross-system CRD reads. For OWNED
	// servers the owner's cumulative usage must also stay within their quota caps.
	if resUpdated {
		newCPU := quantityToMilli(newResources.Limits[corev1.ResourceCPU])
		newMemMB := quantityToMB(newResources.Limits[corev1.ResourceMemory])

		rec, err := a.Repo.ServerByName(r.Context(), name)
		if err != nil && !errors.Is(err, ErrNotFound) {
			writeError(w, r, err)
			return
		}
		if rec != nil && rec.OwnerID != "" {
			ok, err := a.Repo.QuotaCheck(r.Context(), rec.OwnerID, name,
				ResourceSpec{CPUMilli: newCPU, MemoryMB: newMemMB})
			if err != nil {
				writeError(w, r, err)
				return
			}
			if !ok {
				writeError(w, r, newError(http.StatusForbidden, "quota_exceeded",
					"this change would exceed the server owner's resource quota"))
				return
			}
		}

		if err := a.Cluster.PatchServerSpec(r.Context(), name, patch); err != nil {
			a.writeLookupError(w, r, err)
			return
		}
		_ = a.Repo.UpdateServerResources(r.Context(), name, newCPU, newMemMB, 0)
	} else {
		if err := a.Cluster.PatchServerSpec(r.Context(), name, patch); err != nil {
			a.writeLookupError(w, r, err)
			return
		}
	}

	a.audit(r, p.Email, "server.patch", name)
	writeJSON(w, http.StatusOK, map[string]any{
@@ -754,3 +806,25 @@ func (a *API) audit(r *http.Request, actor, action, server string) {
		RequestID:  requestIDFromContext(r.Context()),
	})
}

// quantityToMilli converts a K8s resource.Quantity to millicores (e.g. "2"→2000,
// "500m"→500). A zero/unset quantity returns 0.
func quantityToMilli(q resource.Quantity) int {
	if q.IsZero() {
		return 0
	}
	return int(q.MilliValue())
}

// quantityToMB converts a K8s resource.Quantity to whole megabytes, rounding up
// (e.g. "4Gi"→4096, "1G"→1000). A zero/unset quantity returns 0.
func quantityToMB(q resource.Quantity) int {
	if q.IsZero() {
		return 0
	}
	mb := q.Value() / (1024 * 1024)
	if mb < 1 {
		return 1
	}
	return int(mb)
}
+70 −2
Changes for internal/api/pgrepo.go: 70 added lines, 2 removed lines.
Original line number Diff line number Diff line
@@ -294,6 +294,74 @@ func (p *PGRepo) QuotaAvailable(ctx context.Context, userID string) (bool, error
	return n < maxServers.Int64, nil
}

// QuotaCheck reports whether accepting a server with resource spec `incoming`
// would push userID over any quota cap. excludeName is the server row whose own
// cached resources should be excluded ("" for a fresh claim where the row
// doesn't exist yet). Four dimensions are checked: server count, CPU millicores,
// memory MB, and storage MB. A NULL or missing quota row/column means unlimited
// for that dimension. Like QuotaAvailable, the count check and the write are not
// serialized — see the QuotaAvailable TOCTOU docstring.
func (p *PGRepo) QuotaCheck(ctx context.Context, userID string, excludeName string, incoming ResourceSpec) (bool, error) {
	var maxServers, maxCPU, maxMem, maxStor sql.NullInt64
	switch err := p.db.QueryRowContext(ctx,
		`SELECT max_servers, max_cpu_milli, max_memory_mb, max_storage_gb
		 FROM quotas WHERE user_id = $1`, userID).Scan(
		&maxServers, &maxCPU, &maxMem, &maxStor); {
	case errors.Is(err, sql.ErrNoRows):
		return true, nil // no quota row → unlimited
	case err != nil:
		return false, err
	}

	var count int64
	var cpuSum, memSum, storSum sql.NullInt64
	switch err := p.db.QueryRowContext(ctx,
		`SELECT COUNT(*), COALESCE(SUM(cached_cpu_milli), 0), COALESCE(SUM(cached_memory_mb), 0), COALESCE(SUM(cached_storage_mb), 0)
		 FROM servers WHERE owner_id = $1 AND deleted_at IS NULL AND name != $2`,
		userID, excludeName).Scan(&count, &cpuSum, &memSum, &storSum); {
	case err != nil:
		return false, err
	}

	if maxServers.Valid && count >= maxServers.Int64 {
		return false, nil
	}
	if maxCPU.Valid && cpuSum.Int64+int64(incoming.CPUMilli) > maxCPU.Int64 {
		return false, nil
	}
	if maxMem.Valid && memSum.Int64+int64(incoming.MemoryMB) > maxMem.Int64 {
		return false, nil
	}
	if maxStor.Valid && storSum.Int64+int64(incoming.StorageMB) > maxStor.Int64*1024 {
		return false, nil
	}
	return true, nil
}

// UpdateServerResources updates the resource cache for a server after a spec
// mutation (spec §7 PATCH). The per-owner aggregate used by QuotaCheck is a
// SQL SUM over the cached columns, so every mutation must write through here.
func (p *PGRepo) UpdateServerResources(ctx context.Context, name string, cpuMilli, memoryMB, storageMB int) error {
	_, err := p.db.ExecContext(ctx,
		`UPDATE servers SET cached_cpu_milli = $2, cached_memory_mb = $3, cached_storage_mb = $4 WHERE name = $1 AND deleted_at IS NULL`,
		name, cpuMilli, memoryMB, storageMB)
	return err
}

// ServerResources returns the cached resource spec for a server.
func (p *PGRepo) ServerResources(ctx context.Context, name string) (ResourceSpec, error) {
	var r ResourceSpec
	switch err := p.db.QueryRowContext(ctx,
		`SELECT cached_cpu_milli, cached_memory_mb, cached_storage_mb FROM servers WHERE name = $1 AND deleted_at IS NULL`,
		name).Scan(&r.CPUMilli, &r.MemoryMB, &r.StorageMB); {
	case errors.Is(err, sql.ErrNoRows):
		return r, nil
	case err != nil:
		return r, err
	}
	return r, nil
}

// ClaimServer performs the atomic ownership transfer (spec §9.3). A missing
// server is ErrNotFound; an existing-but-owned server yields claimed=false so the
// handler can answer 409.
@@ -440,7 +508,7 @@ func (p *PGRepo) ServerOwners(ctx context.Context) (map[string]string, error) {
// is a PRIMARY KEY, so a no-op insert means it was already bound; we then
// confirm it resolves to this server and return ErrConflict otherwise, letting
// the create handler answer 409 before it touches the CRD.
func (p *PGRepo) SeedServer(ctx context.Context, name, subdomain string) error {
func (p *PGRepo) SeedServer(ctx context.Context, name, subdomain string, cpuMilli, memoryMB, storageMB int) error {
	tx, err := p.db.BeginTx(ctx, nil)
	if err != nil {
		return err
@@ -448,7 +516,7 @@ func (p *PGRepo) SeedServer(ctx context.Context, name, subdomain string) error {
	defer tx.Rollback() //nolint:errcheck // no-op after commit

	if _, err := tx.ExecContext(ctx,
		`INSERT INTO servers (name) VALUES ($1) ON CONFLICT DO NOTHING`, name); err != nil {
		`INSERT INTO servers (name, cached_cpu_milli, cached_memory_mb, cached_storage_mb) VALUES ($1, $2, $3, $4) ON CONFLICT (name) DO UPDATE SET cached_cpu_milli = EXCLUDED.cached_cpu_milli, cached_memory_mb = EXCLUDED.cached_memory_mb, cached_storage_mb = EXCLUDED.cached_storage_mb`, name, cpuMilli, memoryMB, storageMB); err != nil {
		return fmt.Errorf("seed server row: %w", err)
	}
	if _, err := tx.ExecContext(ctx,
Loading