fix(quota): 回收的服务器保留资源规格,认领时按集群里的真实规格过配额并写回缓存
This commit is contained in:
6 files changed
+171
-19
No files matched your search
@@ -16,6 +16,7 @@ import (
|
|||||||
|
|
||||||
"felis.lolicon.best/internal/apis/felis/v1alpha1"
|
"felis.lolicon.best/internal/apis/felis/v1alpha1"
|
||||||
corev1 "k8s.io/api/core/v1"
|
corev1 "k8s.io/api/core/v1"
|
||||||
|
"k8s.io/apimachinery/pkg/api/resource"
|
||||||
)
|
)
|
||||||
|
|
||||||
const testRoot = "mc.example.net" // neutral; never a deployment domain
|
const testRoot = "mc.example.net" // neutral; never a deployment domain
|
||||||
@@ -47,8 +48,10 @@ type fakeRepo struct {
|
|||||||
// UpdateServerResources records the write for assertions.
|
// UpdateServerResources records the write for assertions.
|
||||||
serverResources map[string]ResourceSpec
|
serverResources map[string]ResourceSpec
|
||||||
resourceUpdates map[string]ResourceSpec
|
resourceUpdates map[string]ResourceSpec
|
||||||
audits []AuditEntry
|
// quotaChecked records the size every QuotaCheck was asked about, in order.
|
||||||
failAudit error // Audit fails with it (a store outage)
|
quotaChecked []ResourceSpec
|
||||||
|
audits []AuditEntry
|
||||||
|
failAudit error // Audit fails with it (a store outage)
|
||||||
// backupRequested mirrors the newest backup.create audit row per server,
|
// backupRequested mirrors the newest backup.create audit row per server,
|
||||||
// stamped by Audit with the wall clock (LastBackupRequest).
|
// stamped by Audit with the wall clock (LastBackupRequest).
|
||||||
backupRequested map[string]time.Time
|
backupRequested map[string]time.Time
|
||||||
@@ -323,7 +326,8 @@ 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) 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) QuotaAvailable(_ context.Context, u string) (bool, error) { return f.quota[u], nil }
|
||||||
|
|
||||||
func (f *fakeRepo) QuotaCheck(_ context.Context, userID string, _ string, _ ResourceSpec) (bool, error) {
|
func (f *fakeRepo) QuotaCheck(_ context.Context, userID string, _ string, incoming ResourceSpec) (bool, error) {
|
||||||
|
f.quotaChecked = append(f.quotaChecked, incoming)
|
||||||
// For hermetic tests, QuotaCheck delegates to the same QuotaAvailable
|
// For hermetic tests, QuotaCheck delegates to the same QuotaAvailable
|
||||||
// store — tests that care about per-dimension checks should use
|
// store — tests that care about per-dimension checks should use
|
||||||
// fakeQuotas with direct inspection.
|
// fakeQuotas with direct inspection.
|
||||||
@@ -2400,7 +2404,7 @@ func TestClaimStateMachine(t *testing.T) {
|
|||||||
t.Run("over quota -> 403", func(t *testing.T) {
|
t.Run("over quota -> 403", func(t *testing.T) {
|
||||||
repo := newFakeRepo()
|
repo := newFakeRepo()
|
||||||
repo.linked["u1"] = true
|
repo.linked["u1"] = true
|
||||||
api := newTestAPI(repo, newFakeCluster())
|
api := newTestAPI(repo, claimCluster())
|
||||||
api.External = staticExternal{p: user}
|
api.External = staticExternal{p: user}
|
||||||
w := do(api.ExternalHandler(), "POST", "/api/v1/servers/survival/claim", "", nil)
|
w := do(api.ExternalHandler(), "POST", "/api/v1/servers/survival/claim", "", nil)
|
||||||
if w.Code != http.StatusForbidden || decodeErr(t, w) != "quota_exceeded" {
|
if w.Code != http.StatusForbidden || decodeErr(t, w) != "quota_exceeded" {
|
||||||
@@ -2415,7 +2419,7 @@ func TestClaimStateMachine(t *testing.T) {
|
|||||||
repo.quota["u1"] = true
|
repo.quota["u1"] = true
|
||||||
repo.claimOK["survival"] = true
|
repo.claimOK["survival"] = true
|
||||||
repo.claimQuotaRefuse["survival"] = true
|
repo.claimQuotaRefuse["survival"] = true
|
||||||
api := newTestAPI(repo, newFakeCluster())
|
api := newTestAPI(repo, claimCluster())
|
||||||
api.External = staticExternal{p: user}
|
api.External = staticExternal{p: user}
|
||||||
w := do(api.ExternalHandler(), "POST", "/api/v1/servers/survival/claim", "", nil)
|
w := do(api.ExternalHandler(), "POST", "/api/v1/servers/survival/claim", "", nil)
|
||||||
if w.Code != http.StatusForbidden || decodeErr(t, w) != "quota_exceeded" {
|
if w.Code != http.StatusForbidden || decodeErr(t, w) != "quota_exceeded" {
|
||||||
@@ -2427,7 +2431,7 @@ func TestClaimStateMachine(t *testing.T) {
|
|||||||
repo.linked["u1"] = true
|
repo.linked["u1"] = true
|
||||||
repo.quota["u1"] = true
|
repo.quota["u1"] = true
|
||||||
repo.claimOK["survival"] = false // row exists but owner already set
|
repo.claimOK["survival"] = false // row exists but owner already set
|
||||||
api := newTestAPI(repo, newFakeCluster())
|
api := newTestAPI(repo, claimCluster())
|
||||||
api.External = staticExternal{p: user}
|
api.External = staticExternal{p: user}
|
||||||
w := do(api.ExternalHandler(), "POST", "/api/v1/servers/survival/claim", "", nil)
|
w := do(api.ExternalHandler(), "POST", "/api/v1/servers/survival/claim", "", nil)
|
||||||
if w.Code != http.StatusConflict || decodeErr(t, w) != "already_claimed" {
|
if w.Code != http.StatusConflict || decodeErr(t, w) != "already_claimed" {
|
||||||
@@ -2439,7 +2443,7 @@ func TestClaimStateMachine(t *testing.T) {
|
|||||||
repo.linked["u1"] = true
|
repo.linked["u1"] = true
|
||||||
repo.quota["u1"] = true
|
repo.quota["u1"] = true
|
||||||
repo.claimOK["survival"] = true
|
repo.claimOK["survival"] = true
|
||||||
api := newTestAPI(repo, newFakeCluster())
|
api := newTestAPI(repo, claimCluster())
|
||||||
api.External = staticExternal{p: user}
|
api.External = staticExternal{p: user}
|
||||||
w := do(api.ExternalHandler(), "POST", "/api/v1/servers/survival/claim", "", nil)
|
w := do(api.ExternalHandler(), "POST", "/api/v1/servers/survival/claim", "", nil)
|
||||||
if w.Code != http.StatusOK {
|
if w.Code != http.StatusOK {
|
||||||
@@ -2449,6 +2453,71 @@ func TestClaimStateMachine(t *testing.T) {
|
|||||||
t.Fatalf("audit not written as expected: %+v", repo.audits)
|
t.Fatalf("audit not written as expected: %+v", repo.audits)
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
|
// A world the reaper released used to have its resource cache zeroed, and the
|
||||||
|
// quota gate read the cache: the claim passed every resource cap and the server
|
||||||
|
// then counted as nothing. The claim reads the server's size off the cluster,
|
||||||
|
// gates on it, and writes it back to the cache.
|
||||||
|
t.Run("gated on the server's real size, whatever the cache says", func(t *testing.T) {
|
||||||
|
repo := newFakeRepo()
|
||||||
|
repo.linked["u1"] = true
|
||||||
|
repo.claimOK["survival"] = true
|
||||||
|
repo.serverResources["survival"] = ResourceSpec{}
|
||||||
|
api := newTestAPI(repo, claimCluster())
|
||||||
|
api.External = staticExternal{p: user}
|
||||||
|
w := do(api.ExternalHandler(), "POST", "/api/v1/servers/survival/claim", "", nil)
|
||||||
|
if w.Code != http.StatusForbidden || decodeErr(t, w) != "quota_exceeded" {
|
||||||
|
t.Fatalf("code = %d body %s, want 403 quota_exceeded", w.Code, w.Body.String())
|
||||||
|
}
|
||||||
|
want := ResourceSpec{CPUMilli: 2000, MemoryMB: 4096, StorageMB: 10240}
|
||||||
|
if len(repo.quotaChecked) != 1 || repo.quotaChecked[0] != want {
|
||||||
|
t.Errorf("quota checked against %+v, want [%+v]", repo.quotaChecked, want)
|
||||||
|
}
|
||||||
|
if got := repo.resourceUpdates["survival"]; got != want {
|
||||||
|
t.Errorf("resource cache = %+v, want %+v", got, want)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
t.Run("the internal face is gated the same way", func(t *testing.T) {
|
||||||
|
repo := newFakeRepo()
|
||||||
|
repo.links[menuUUID] = "u1"
|
||||||
|
repo.claimOK["survival"] = true
|
||||||
|
api := newTestAPI(repo, claimCluster())
|
||||||
|
w := internalClaim(api, `{"mc_uuid":"`+menuUUID+`"}`)
|
||||||
|
if w.Code != http.StatusForbidden || decodeErr(t, w) != "quota_exceeded" {
|
||||||
|
t.Fatalf("code = %d body %s, want 403 quota_exceeded", w.Code, w.Body.String())
|
||||||
|
}
|
||||||
|
want := ResourceSpec{CPUMilli: 2000, MemoryMB: 4096, StorageMB: 10240}
|
||||||
|
if len(repo.quotaChecked) != 1 || repo.quotaChecked[0] != want {
|
||||||
|
t.Errorf("quota checked against %+v, want [%+v]", repo.quotaChecked, want)
|
||||||
|
}
|
||||||
|
if got := repo.resourceUpdates["survival"]; got != want {
|
||||||
|
t.Errorf("resource cache = %+v, want %+v", got, want)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
t.Run("a server the cluster does not have -> 404", func(t *testing.T) {
|
||||||
|
repo := newFakeRepo()
|
||||||
|
repo.linked["u1"] = true
|
||||||
|
repo.quota["u1"] = true
|
||||||
|
repo.claimOK["survival"] = true
|
||||||
|
api := newTestAPI(repo, newFakeCluster())
|
||||||
|
api.External = staticExternal{p: user}
|
||||||
|
w := do(api.ExternalHandler(), "POST", "/api/v1/servers/survival/claim", "", nil)
|
||||||
|
if w.Code != http.StatusNotFound {
|
||||||
|
t.Fatalf("code = %d body %s, want 404", w.Code, w.Body.String())
|
||||||
|
}
|
||||||
|
if len(repo.audits) != 0 {
|
||||||
|
t.Fatalf("a refused claim was audited: %+v", repo.audits)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
// claimCluster holds the survival server a claim test claims: 2 CPUs, 4Gi of memory
|
||||||
|
// and a 10Gi world.
|
||||||
|
func claimCluster() *fakeCluster {
|
||||||
|
cl := newFakeCluster()
|
||||||
|
cl.byName["survival"] = &ServerInfo{Name: "survival", StorageSize: "10Gi", Resources: corev1.ResourceRequirements{
|
||||||
|
Limits: corev1.ResourceList{corev1.ResourceCPU: resource.MustParse("2"), corev1.ResourceMemory: resource.MustParse("4Gi")},
|
||||||
|
}}
|
||||||
|
return cl
|
||||||
}
|
}
|
||||||
|
|
||||||
// ---- wake (autostartPolicy gate + cooldown) ----
|
// ---- wake (autostartPolicy gate + cooldown) ----
|
||||||
|
|||||||
@@ -34,7 +34,7 @@ func TestAccountLinkVertical(t *testing.T) {
|
|||||||
repo.quota["u1"] = true
|
repo.quota["u1"] = true
|
||||||
repo.claimOK["survival"] = true
|
repo.claimOK["survival"] = true
|
||||||
|
|
||||||
api := newTestAPI(repo, newFakeCluster())
|
api := newTestAPI(repo, claimCluster())
|
||||||
api.External = staticExternal{p: user}
|
api.External = staticExternal{p: user}
|
||||||
ih := api.InternalHandler()
|
ih := api.InternalHandler()
|
||||||
eh := api.ExternalHandler()
|
eh := api.ExternalHandler()
|
||||||
|
|||||||
@@ -251,10 +251,11 @@ func (a *API) handleInternalClaim(w http.ResponseWriter, r *http.Request) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// ② quota gate, evaluated before the ownership write (mirrors handleClaim).
|
// ② quota gate, evaluated before the ownership write (mirrors handleClaim).
|
||||||
// All four dimensions (servers, CPU, memory, storage) are checked.
|
// All four dimensions (servers, CPU, memory, storage) are checked, with the
|
||||||
res, err := a.Repo.ServerResources(r.Context(), name)
|
// server at its real size.
|
||||||
|
res, err := a.claimResources(r.Context(), name)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
writeError(w, r, err)
|
a.writeLookupError(w, r, err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
ok, err := a.Repo.QuotaCheck(r.Context(), userID, "", res)
|
ok, err := a.Repo.QuotaCheck(r.Context(), userID, "", res)
|
||||||
|
|||||||
@@ -138,11 +138,12 @@ func (a *API) handleClaim(w http.ResponseWriter, r *http.Request) {
|
|||||||
|
|
||||||
// ② quota gate, evaluated before the ownership write. All four dimensions
|
// ② quota gate, evaluated before the ownership write. All four dimensions
|
||||||
// (servers, CPU, memory, storage) are checked against the user's quota caps
|
// (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
|
// (spec §9.3, §22), with the server at its real size (claimResources). The
|
||||||
// claimed has owner_id=NULL so it is not yet in the per-owner aggregate.
|
// server being claimed has owner_id=NULL so it is not yet in the per-owner
|
||||||
res, err := a.Repo.ServerResources(r.Context(), name)
|
// aggregate.
|
||||||
|
res, err := a.claimResources(r.Context(), name)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
writeError(w, r, err)
|
a.writeLookupError(w, r, err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
ok, err := a.Repo.QuotaCheck(r.Context(), p.UserID, "", res)
|
ok, err := a.Repo.QuotaCheck(r.Context(), p.UserID, "", res)
|
||||||
@@ -178,6 +179,28 @@ func (a *API) handleClaim(w http.ResponseWriter, r *http.Request) {
|
|||||||
writeJSON(w, http.StatusOK, map[string]any{"name": name, "claimed": true})
|
writeJSON(w, http.StatusOK, map[string]any{"name": name, "claimed": true})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// claimResources is the size a claim is gated on and then counted at: the server's
|
||||||
|
// spec as the cluster holds it, written through to the resource cache first. The
|
||||||
|
// cache is all the per-owner quota sums read, and a world the reaper released
|
||||||
|
// before it kept the cache holds zeros there: gated on those, a claim passed every
|
||||||
|
// resource cap and the server went uncounted for as long as its new owner kept it.
|
||||||
|
func (a *API) claimResources(ctx context.Context, name string) (ResourceSpec, error) {
|
||||||
|
info, err := a.Cluster.GetServer(ctx, name)
|
||||||
|
if err != nil {
|
||||||
|
return ResourceSpec{}, err
|
||||||
|
}
|
||||||
|
storage, _ := resource.ParseQuantity(info.StorageSize)
|
||||||
|
res := ResourceSpec{
|
||||||
|
CPUMilli: quantityToMilli(info.Resources.Limits[corev1.ResourceCPU]),
|
||||||
|
MemoryMB: quantityToMB(info.Resources.Limits[corev1.ResourceMemory]),
|
||||||
|
StorageMB: quantityToMB(storage),
|
||||||
|
}
|
||||||
|
if err := a.Repo.UpdateServerResources(ctx, name, res.CPUMilli, res.MemoryMB, res.StorageMB); err != nil {
|
||||||
|
return ResourceSpec{}, err
|
||||||
|
}
|
||||||
|
return res, nil
|
||||||
|
}
|
||||||
|
|
||||||
// handleStatus returns the CRD status view (spec §7 GET /servers/{name}/status).
|
// handleStatus returns the CRD status view (spec §7 GET /servers/{name}/status).
|
||||||
func (a *API) handleStatus(w http.ResponseWriter, r *http.Request) {
|
func (a *API) handleStatus(w http.ResponseWriter, r *http.Request) {
|
||||||
name := r.PathValue("name")
|
name := r.PathValue("name")
|
||||||
|
|||||||
@@ -5,6 +5,7 @@ package pgint
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"database/sql"
|
"database/sql"
|
||||||
|
"errors"
|
||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
@@ -392,6 +393,63 @@ func TestBackupReadBack(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestReleasedWorldKeepsItsSize: a world the reaper released keeps its resource
|
||||||
|
// cache, so the next claim is gated on the server's size and counts it. Zeroed on
|
||||||
|
// release, it let a user whose quota is spent claim the server anyway and then
|
||||||
|
// counted it as nothing.
|
||||||
|
func TestReleasedWorldKeepsItsSize(t *testing.T) {
|
||||||
|
ctx := context.Background()
|
||||||
|
st := reaper.NewPGStore(db)
|
||||||
|
name := "released-" + suffix(t)
|
||||||
|
if _, err := db.ExecContext(ctx,
|
||||||
|
`INSERT INTO servers (name, cached_cpu_milli, cached_memory_mb, cached_storage_mb) VALUES ($1, 2000, 4096, 10240)`,
|
||||||
|
name); err != nil {
|
||||||
|
t.Fatalf("seed server: %v", err)
|
||||||
|
}
|
||||||
|
first := newUser(t, "user", "released-a")
|
||||||
|
if ok, err := repo.ClaimServer(ctx, name, first.ID); err != nil || !ok {
|
||||||
|
t.Fatalf("first claim = (%v, %v)", ok, err)
|
||||||
|
}
|
||||||
|
if err := st.ReleaseWorld(ctx, name, time.Now()); err != nil {
|
||||||
|
t.Fatalf("ReleaseWorld: %v", err)
|
||||||
|
}
|
||||||
|
var owner sql.NullString
|
||||||
|
var cpu, mem, stor int
|
||||||
|
if err := db.QueryRowContext(ctx,
|
||||||
|
`SELECT owner_id, cached_cpu_milli, cached_memory_mb, cached_storage_mb FROM servers WHERE name = $1`, name).
|
||||||
|
Scan(&owner, &cpu, &mem, &stor); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if owner.Valid || cpu != 2000 || mem != 4096 || stor != 10240 {
|
||||||
|
t.Fatalf("after ReleaseWorld: owner=%v cache=%d/%d/%d; want no owner, cache 2000/4096/10240", owner, cpu, mem, stor)
|
||||||
|
}
|
||||||
|
|
||||||
|
// One CPU of quota does not cover a two-CPU server.
|
||||||
|
small := newUser(t, "user", "released-b")
|
||||||
|
oneCPU := 1000
|
||||||
|
if _, err := repo.SetQuotas(ctx, small.ID, api.QuotaInput{MaxCPUMilli: &oneCPU}, "pgint"); err != nil {
|
||||||
|
t.Fatalf("SetQuotas: %v", err)
|
||||||
|
}
|
||||||
|
if ok, err := repo.ClaimServer(ctx, name, small.ID); !errors.Is(err, api.ErrQuotaExceeded) {
|
||||||
|
t.Fatalf("claim over the CPU quota = (%v, %v), want ErrQuotaExceeded", ok, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// A claim that fits counts the server at its size.
|
||||||
|
fits := newUser(t, "user", "released-c")
|
||||||
|
if ok, err := repo.ClaimServer(ctx, name, fits.ID); err != nil || !ok {
|
||||||
|
t.Fatalf("claim within quota = (%v, %v)", ok, err)
|
||||||
|
}
|
||||||
|
var used int
|
||||||
|
if err := db.QueryRowContext(ctx,
|
||||||
|
`SELECT COALESCE(SUM(cached_cpu_milli), 0) FROM servers WHERE owner_id = $1 AND deleted_at IS NULL`, fits.ID).
|
||||||
|
Scan(&used); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if used != 2000 {
|
||||||
|
t.Fatalf("the new owner's CPU use = %d, want 2000", used)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// TestRestartClock: an idle server with no world to reclaim gets its clock and
|
// TestRestartClock: an idle server with no world to reclaim gets its clock and
|
||||||
// warnings reset, and keeps its owner and resource cache (data-durability-19).
|
// warnings reset, and keeps its owner and resource cache (data-durability-19).
|
||||||
func TestRestartClock(t *testing.T) {
|
func TestRestartClock(t *testing.T) {
|
||||||
|
|||||||
@@ -79,12 +79,13 @@ func (s *PGStore) InsertBackup(ctx context.Context, rec BackupRecord) error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
// ReleaseWorld releases ownership, zeros the resource cache, and resets the
|
// ReleaseWorld releases ownership and resets the activity clock and warnings —
|
||||||
// activity clock and warnings — without deleting the row (red line ②).
|
// without deleting the row (red line ②). The resource cache stays: the server
|
||||||
|
// keeps its spec, an ownerless row is in nobody's quota sum, and the next claim is
|
||||||
|
// gated on that size and counts it.
|
||||||
func (s *PGStore) ReleaseWorld(ctx context.Context, name string, at time.Time) error {
|
func (s *PGStore) ReleaseWorld(ctx context.Context, name string, at time.Time) error {
|
||||||
const q = `UPDATE servers
|
const q = `UPDATE servers
|
||||||
SET owner_id = NULL, cached_cpu_milli = 0, cached_memory_mb = 0, cached_storage_mb = 0,
|
SET owner_id = NULL, last_active_at = $2, warned_3d_at = NULL, warned_1d_at = NULL
|
||||||
last_active_at = $2, warned_3d_at = NULL, warned_1d_at = NULL
|
|
||||||
WHERE name = $1 AND deleted_at IS NULL`
|
WHERE name = $1 AND deleted_at IS NULL`
|
||||||
_, err := s.db.ExecContext(ctx, q, name, at)
|
_, err := s.db.ExecContext(ctx, q, name, at)
|
||||||
return err
|
return err
|
||||||
|
|||||||
Reference in new issue
Block a user