fix(api): 调整资源先在 owner 的认领锁内核配额并写缓存再 patch 集群,集群拒绝时缓存回到集群的大小
This commit is contained in:
6 files changed
+451
-71
No files matched your search
+85
-25
@@ -371,13 +371,16 @@ func (p *PGRepo) QuotaAvailable(ctx context.Context, userID string) (bool, error
|
||||
// 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. As a standalone read it is advisory — it backs the
|
||||
// handler's fast-path 403 — while the AUTHORITATIVE gate for claims is the one
|
||||
// ClaimServer re-runs atomically; the resize path (server PATCH) keeps this
|
||||
// advisory shape because its write goes through the Kubernetes API, not this
|
||||
// transaction.
|
||||
// handler's fast-path 403 — while the AUTHORITATIVE gates are the ones
|
||||
// ClaimServer and ResizeServer run in the owner's claim lane.
|
||||
func (p *PGRepo) QuotaCheck(ctx context.Context, userID string, excludeName string, incoming ResourceSpec) (bool, error) {
|
||||
return quotaFits(ctx, p.db, userID, excludeName, incoming)
|
||||
}
|
||||
|
||||
// quotaFits is QuotaCheck over either the pool or a transaction.
|
||||
func quotaFits(ctx context.Context, q rowQuerier, userID, excludeName string, incoming ResourceSpec) (bool, error) {
|
||||
var maxServers, maxCPU, maxMem, maxStor sql.NullInt64
|
||||
switch err := p.db.QueryRowContext(ctx,
|
||||
switch err := q.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); {
|
||||
@@ -388,20 +391,19 @@ func (p *PGRepo) QuotaCheck(ctx context.Context, userID string, excludeName stri
|
||||
}
|
||||
|
||||
var count, cpuSum, memSum, storSum int64
|
||||
switch err := p.db.QueryRowContext(ctx,
|
||||
if err := q.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:
|
||||
userID, excludeName).Scan(&count, &cpuSum, &memSum, &storSum); err != nil {
|
||||
return false, err
|
||||
}
|
||||
|
||||
return quotaAllows(maxServers, maxCPU, maxMem, maxStor, count, cpuSum, memSum, storSum, 1, incoming), 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.
|
||||
// UpdateServerResources overwrites a server's resource cache. The per-owner
|
||||
// aggregate used by QuotaCheck is a SQL SUM over the cached columns, so every
|
||||
// size the cluster takes must reach them, through here or ResizeServer.
|
||||
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`,
|
||||
@@ -409,18 +411,76 @@ func (p *PGRepo) UpdateServerResources(ctx context.Context, name string, cpuMill
|
||||
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
|
||||
// ResizeServer gates and records a CPU and memory resize before the cluster sees
|
||||
// it. The owner is only known from the row, and the lane must be taken before the
|
||||
// row lock (a redeem holds lanes while it moves rows), so a server that changes
|
||||
// hands in between is read again in a fresh transaction under its new owner's lane.
|
||||
func (p *PGRepo) ResizeServer(ctx context.Context, name string, cpuMilli, memoryMB int) (ResourceSpec, error) {
|
||||
for attempt := 1; ; attempt++ {
|
||||
prev, moved, err := p.resizeServer(ctx, name, cpuMilli, memoryMB)
|
||||
if !moved {
|
||||
return prev, err
|
||||
}
|
||||
if attempt == 3 {
|
||||
return ResourceSpec{}, fmt.Errorf("resize %s: its owner changed three times while it waited", name)
|
||||
}
|
||||
}
|
||||
return r, nil
|
||||
}
|
||||
|
||||
// resizeServer is one attempt of ResizeServer. moved reports that the owner it
|
||||
// locked for is no longer the row's, and nothing was written.
|
||||
func (p *PGRepo) resizeServer(ctx context.Context, name string, cpuMilli, memoryMB int) (prev ResourceSpec, moved bool, err error) {
|
||||
tx, err := p.db.BeginTx(ctx, nil)
|
||||
if err != nil {
|
||||
return ResourceSpec{}, false, err
|
||||
}
|
||||
defer tx.Rollback() //nolint:errcheck // no-op after commit
|
||||
|
||||
var owner sql.NullString
|
||||
switch err := tx.QueryRowContext(ctx,
|
||||
`SELECT owner_id FROM servers WHERE name = $1 AND deleted_at IS NULL`, name).Scan(&owner); {
|
||||
case errors.Is(err, sql.ErrNoRows):
|
||||
return ResourceSpec{}, false, nil
|
||||
case err != nil:
|
||||
return ResourceSpec{}, false, err
|
||||
}
|
||||
if owner.Valid {
|
||||
if _, err := tx.ExecContext(ctx, `SELECT pg_advisory_xact_lock(hashtext($1))`, owner.String); err != nil {
|
||||
return ResourceSpec{}, false, err
|
||||
}
|
||||
}
|
||||
var locked sql.NullString
|
||||
switch err := tx.QueryRowContext(ctx,
|
||||
`SELECT owner_id, cached_cpu_milli, cached_memory_mb, cached_storage_mb
|
||||
FROM servers WHERE name = $1 AND deleted_at IS NULL FOR UPDATE`,
|
||||
name).Scan(&locked, &prev.CPUMilli, &prev.MemoryMB, &prev.StorageMB); {
|
||||
case errors.Is(err, sql.ErrNoRows):
|
||||
return ResourceSpec{}, false, nil
|
||||
case err != nil:
|
||||
return ResourceSpec{}, false, err
|
||||
}
|
||||
if locked != owner {
|
||||
return ResourceSpec{}, true, nil
|
||||
}
|
||||
|
||||
// A server nobody owns is checked against the empty owner, who has no caps:
|
||||
// every quotas row belongs to a user.
|
||||
if cpuMilli > prev.CPUMilli || memoryMB > prev.MemoryMB {
|
||||
ok, err := quotaFits(ctx, tx, owner.String, name,
|
||||
ResourceSpec{CPUMilli: cpuMilli, MemoryMB: memoryMB, StorageMB: prev.StorageMB})
|
||||
if err != nil {
|
||||
return ResourceSpec{}, false, err
|
||||
}
|
||||
if !ok {
|
||||
return ResourceSpec{}, false, ErrQuotaExceeded
|
||||
}
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx,
|
||||
`UPDATE servers SET cached_cpu_milli = $2, cached_memory_mb = $3 WHERE name = $1`,
|
||||
name, cpuMilli, memoryMB); err != nil {
|
||||
return ResourceSpec{}, false, err
|
||||
}
|
||||
return prev, false, tx.Commit()
|
||||
}
|
||||
|
||||
// ClaimServer performs the atomic ownership transfer (spec §9.3). A missing
|
||||
@@ -2696,15 +2756,15 @@ func (p *PGRepo) OTPLockedUntil(ctx context.Context, userID, purpose string, now
|
||||
return otpLockedUntil(ctx, p.db, userID, purpose, now, false)
|
||||
}
|
||||
|
||||
// otpQuerier is the read half shared by *sql.DB and *sql.Tx.
|
||||
type otpQuerier interface {
|
||||
// rowQuerier is the read half shared by *sql.DB and *sql.Tx.
|
||||
type rowQuerier interface {
|
||||
QueryRowContext(ctx context.Context, query string, args ...any) *sql.Row
|
||||
}
|
||||
|
||||
// otpLockedUntil returns when the (user, purpose) lock ends, or zero when the
|
||||
// budget is not spent in the current window. forUpdate takes the row lock so a
|
||||
// redeem serialises its check with its own charge.
|
||||
func otpLockedUntil(ctx context.Context, q otpQuerier, userID, purpose string, now time.Time, forUpdate bool) (time.Time, error) {
|
||||
func otpLockedUntil(ctx context.Context, q rowQuerier, userID, purpose string, now time.Time, forUpdate bool) (time.Time, error) {
|
||||
query := `SELECT window_start, failures FROM otp_failure_windows WHERE user_id = $1 AND purpose = $2`
|
||||
if forUpdate {
|
||||
query += ` FOR UPDATE`
|
||||
|
||||
Reference in new issue
Block a user