fix(panel): diagnose runtime failures and add owner recovery controls

This commit is contained in:
Lemon-miaow committed 2026-10-05 18:06:55 +08:00
1 parent d7ae1e66c1
commit a7d3ca4c9f
47 files changed
+1590 -102

No files matched your search

+9 -9
View File
@@ -177,12 +177,9 @@ type API struct {
SubmitCreateCooldown time.Duration
SubmitUploadCooldown time.Duration
// MaxRunningServers caps how many servers may be desired-Running cluster-wide
// (spec §9.1: the concurrency-上限 lever hanging on the same wake chokepoint as
// cooldown and autostartPolicy). Zero — the default — disables it: §9.2 wires
// only autostartPolicy + cooldown as active wake gates, so this lever ships
// inert, exactly like a zero WakeCooldown, and a deployment opts in by setting
// a positive value. Enforced via withinRunningCap on the wake path.
// MaxRunningServers is the default desired-Running admission limit (spec §9.1).
// Zero disables it. A saved Owner wake policy overrides this and WakeCooldown;
// panel, in-game and scheduled starts share withinRunningCap.
MaxRunningServers int
// MaxStreamsPerPrincipal caps how many concurrent Server-Sent Event streams
@@ -538,6 +535,7 @@ func (a *API) externalAPIRoutes() []apiRoute {
// App-auth tier: operations on your own servers (spec §14).
{Method: "POST", Pattern: "/api/v1/servers/{name}/wake", h: a.handleWake},
{Method: "POST", Pattern: "/api/v1/servers/{name}/stop", h: a.handleStop},
{Method: "POST", Pattern: "/api/v1/servers/{name}/emergency-stop", h: a.handleEmergencyStop, Admin: true, Owner: true},
{Method: "POST", Pattern: "/api/v1/servers/{name}/claim", h: a.handleClaim},
// Console write (spec §8 写=RCON): owner/admin-gated inside the handler, so
// it sits in the app-tier block (操作自己服 → app 鉴权), not behind adminOnly.
@@ -652,6 +650,8 @@ func (a *API) externalAPIRoutes() []apiRoute {
// Staff can designate their own game identity after panel setup. Players
// retain the in-game proof flow above.
{Method: "GET", Pattern: "/api/v1/account/link/sources", Admin: true, h: a.handleLinkSources},
{Method: "GET", Pattern: "/api/v1/settings/wake-policy", Owner: true, Admin: true, h: a.handleGetWakePolicy},
{Method: "PUT", Pattern: "/api/v1/settings/wake-policy", Owner: true, Admin: true, h: a.handleSetWakePolicy},
{Method: "GET", Pattern: "/api/v1/settings/auth-sources", Owner: true, Admin: true, h: a.handleGetAuthSources},
{Method: "PUT", Pattern: "/api/v1/settings/auth-sources", Owner: true, Admin: true, h: a.handleSetAuthSources},
{Method: "POST", Pattern: "/api/v1/settings/auth-sources/test", Owner: true, Admin: true, h: a.handleTestAuthSource},
@@ -1205,8 +1205,8 @@ func (l *streamLimiter) release(key string) {
// can momentarily exceed the cap. That is acceptable because the operator
// reconcile is idempotent and the §18 reaper / §9.3 quota bound steady-state
// load; the cap exists to refuse an obvious flood, not to hold a hard ceiling.
func (a *API) withinRunningCap(ctx context.Context, info *ServerInfo) (bool, error) {
if a.MaxRunningServers <= 0 || naming.IsSystemServer(info.Name) {
func (a *API) withinRunningCap(ctx context.Context, info *ServerInfo, cap int) (bool, error) {
if cap <= 0 || naming.IsSystemServer(info.Name) {
return true, nil
}
if info.DesiredState == string(v1alpha1.DesiredRunning) {
@@ -1222,5 +1222,5 @@ func (a *API) withinRunningCap(ctx context.Context, info *ServerInfo) (bool, err
running++
}
}
return running < a.MaxRunningServers, nil
return running < cap, nil
}
+2
View File
@@ -57,6 +57,8 @@ type ServerInfo struct {
// Resources is the spec's pod resource block. It stays off the wire; a spec
// patch reads it so the fields the admin left out keep their values.
Resources corev1.ResourceRequirements `json:"-"`
// Startup explains scheduling and container state independently of the CR phase.
Startup *StartupStatus `json:"startup,omitempty"`
}
// CreateServerInput is the validated, structured create-server form (spec §15).
+92
View File
@@ -0,0 +1,92 @@
package api
import (
"context"
"fmt"
"net/http"
"time"
"felis.lolicon.best/internal/apis/felis/v1alpha1"
appsv1 "k8s.io/api/apps/v1"
autoscalingv1 "k8s.io/api/autoscaling/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/util/retry"
"sigs.k8s.io/controller-runtime/pkg/client"
)
// EmergencyStop records the stop intent before scaling the owned workload. It
// bypasses RCON and operator reconciliation, retaining normal Pod shutdown grace.
func (k *K8sCluster) EmergencyStop(ctx context.Context, name string) error {
if err := k.SetDesiredState(ctx, name, v1alpha1.DesiredStopped); err != nil {
return err
}
return retry.RetryOnConflict(retry.DefaultRetry, func() error {
var ms v1alpha1.MinecraftServer
if err := k.getServer(ctx, name, &ms); err != nil {
return err
}
if ms.Spec.DesiredState != v1alpha1.DesiredStopped {
return newError(http.StatusConflict, "conflict", "stop intent changed; retry emergency stop")
}
var sts appsv1.StatefulSet
if err := k.c.Get(ctx, types.NamespacedName{Namespace: k.namespace, Name: name}, &sts); err != nil {
if apierrors.IsNotFound(err) {
return nil
}
return err
}
owner := metav1.GetControllerOf(&sts)
if owner == nil || ms.UID == "" || owner.UID != ms.UID || owner.Kind != "MinecraftServer" || owner.APIVersion != v1alpha1.GroupVersion.String() {
return fmt.Errorf("stop intent saved, but workload ownership could not be verified")
}
scale := &autoscalingv1.Scale{ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: k.namespace, ResourceVersion: sts.ResourceVersion}, Spec: autoscalingv1.ScaleSpec{Replicas: 0}}
if err := k.c.SubResource("scale").Update(ctx, &sts, client.WithSubResourceBody(scale)); err != nil {
return fmt.Errorf("stop intent saved, but workload scale-down was not confirmed: %w", err)
}
return nil
})
}
func (a *API) handleEmergencyStop(w http.ResponseWriter, r *http.Request) {
if !a.requireReauth(w, r, principalFromContext(r.Context())) {
return
}
name := r.PathValue("name")
if err := validateManagedServerName(r, name); err != nil {
writeError(w, r, newError(http.StatusBadRequest, "bad_name", "invalid server name"))
return
}
if err := requireJSONContentType(r); err != nil {
writeError(w, r, err)
return
}
var body struct {
Confirm string `json:"confirm"`
}
if err := decodeJSON(w, r, &body); err != nil {
writeError(w, r, err)
return
}
if body.Confirm != name {
writeError(w, r, newError(http.StatusBadRequest, "bad_request", "type the exact server name to confirm"))
return
}
stopper, ok := a.Cluster.(interface {
EmergencyStop(context.Context, string) error
})
if !ok {
writeError(w, r, newError(http.StatusServiceUnavailable, "unavailable", "this cluster does not support direct emergency shutdown"))
return
}
ctx, cancel := context.WithTimeout(r.Context(), 15*time.Second)
defer cancel()
if err := stopper.EmergencyStop(ctx, name); err != nil {
a.audit(r, "emergency_stop.failed", name)
a.writeLookupError(w, r, err)
return
}
a.audit(r, "emergency_stop", name)
writeJSON(w, http.StatusAccepted, map[string]any{"name": name, "desiredState": "Stopped"})
}
+123
View File
@@ -0,0 +1,123 @@
package api
import (
"context"
"errors"
"testing"
"felis.lolicon.best/internal/apis/felis/v1alpha1"
appsv1 "k8s.io/api/apps/v1"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/types"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/client/fake"
"sigs.k8s.io/controller-runtime/pkg/client/interceptor"
)
func TestEmergencyStopIndependentOfOperator(t *testing.T) {
for _, mode := range []string{"success", "wrong owner", "scale unavailable", "missing workload"} {
t.Run(mode, func(t *testing.T) {
scheme := runtime.NewScheme()
v1alpha1.AddToScheme(scheme)
appsv1.AddToScheme(scheme)
corev1.AddToScheme(scheme)
ms := &v1alpha1.MinecraftServer{ObjectMeta: metav1.ObjectMeta{Name: "survival", Namespace: "minecraft", UID: "server-uid"}, Spec: v1alpha1.MinecraftServerSpec{DesiredState: v1alpha1.DesiredRunning}, Status: v1alpha1.MinecraftServerStatus{Phase: v1alpha1.PhaseRunning}}
replicas := int32(1)
sts := &appsv1.StatefulSet{ObjectMeta: metav1.ObjectMeta{Name: ms.Name, Namespace: ms.Namespace, OwnerReferences: []metav1.OwnerReference{*metav1.NewControllerRef(ms, v1alpha1.GroupVersion.WithKind("MinecraftServer"))}}, Spec: appsv1.StatefulSetSpec{Replicas: &replicas}}
if mode == "wrong owner" {
sts.OwnerReferences[0].UID = "unrelated"
}
builder := fake.NewClientBuilder().WithScheme(scheme).WithObjects(ms)
if mode != "missing workload" {
builder = builder.WithObjects(sts)
}
if mode == "scale unavailable" {
builder = builder.WithInterceptorFuncs(interceptor.Funcs{SubResourceUpdate: func(context.Context, client.Client, string, client.Object, ...client.SubResourceUpdateOption) error {
return errors.New("kubernetes unavailable")
}})
}
c := builder.Build()
k := NewK8sCluster(c, "minecraft")
ctx := context.Background()
err := k.EmergencyStop(ctx, ms.Name)
if (err != nil) != (mode == "wrong owner" || mode == "scale unavailable") {
t.Fatal(mode, err)
}
if err := c.Get(ctx, client.ObjectKeyFromObject(ms), ms); err != nil {
t.Fatal(err)
}
if ms.Spec.DesiredState != v1alpha1.DesiredStopped || ms.Status.Phase != v1alpha1.PhaseRunning {
t.Fatal("stop intent missing or status falsely rewritten", ms)
}
if mode != "success" && mode != "missing workload" {
if err := c.Get(ctx, client.ObjectKeyFromObject(sts), sts); err != nil {
t.Fatal(err)
}
if *sts.Spec.Replicas != 1 {
t.Fatal("unsafe scaling")
}
return
}
if mode == "success" {
pod := &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Name: ms.Name + "-0", Namespace: ms.Namespace, Labels: map[string]string{v1alpha1.LabelServer: ms.Name, v1alpha1.LabelComponent: gamePodComponent}}, Status: corev1.PodStatus{Phase: corev1.PodRunning}}
if err := c.Create(ctx, pod); err != nil {
t.Fatal(err)
}
info, err := k.GetServer(ctx, ms.Name)
if err != nil || info.Phase != "Stopping" {
t.Fatal("running process falsely stopped", info, err)
}
if err := c.Delete(ctx, pod); err != nil {
t.Fatal(err)
}
if err := c.Get(ctx, types.NamespacedName{Namespace: ms.Namespace, Name: ms.Name}, sts); err != nil || *sts.Spec.Replicas != 0 {
t.Fatal(sts, err)
}
}
info, err := k.GetServer(ctx, ms.Name)
if err != nil || info.Phase != "Stopped" {
t.Fatal("operator-independent confirmation missing", info, err)
}
})
}
}
type emergencyFakeCluster struct {
*fakeCluster
calls int
}
func (c *emergencyFakeCluster) EmergencyStop(context.Context, string) error { c.calls++; return nil }
func TestEmergencyStopAuthorization(t *testing.T) {
cl := &emergencyFakeCluster{fakeCluster: newFakeCluster()}
a := newTestAPI(newFakeRepo(), cl)
path := "/api/v1/servers/survival/emergency-stop"
for _, p := range []*Principal{nil, {UserID: "admin", Role: "admin", ViaAdminAccess: true}, {UserID: "owner", Role: "owner"}} {
a.External = staticExternal{p: p}
w := do(a.ExternalHandler(), "POST", path, `{"confirm":"survival"}`, jsonHeader)
if w.Code != 401 && w.Code != 403 {
t.Fatal(w.Code)
}
}
p := &Principal{UserID: "owner", Role: "owner", ViaAdminAccess: true, ViaSession: true, EmailVerified: true}
a.External = staticExternal{p: p}
a.Repo.(*fakeRepo).passkeyCreds["key"] = PasskeyCredential{ID: "key", UserID: p.UserID, UserVerified: true}
if w := do(a.ExternalHandler(), "POST", path, `{"confirm":"survival"}`, jsonHeader); w.Code != 403 || decodeErr(t, w) != "reauth_required" {
t.Fatal(w.Code, w.Body.String())
}
p.ReauthAt = a.now()
if w := do(a.ExternalHandler(), "POST", path, `{"confirm":"wrong"}`, jsonHeader); w.Code != 400 {
t.Fatal(w.Code)
}
if cl.calls != 0 {
t.Fatal("unconfirmed stop ran")
}
if w := do(a.ExternalHandler(), "POST", path, `{"confirm":"survival"}`, jsonHeader); w.Code != 202 {
t.Fatal(w.Code, w.Body.String())
}
if cl.calls != 1 {
t.Fatal(cl.calls)
}
}
+7 -2
View File
@@ -173,7 +173,12 @@ func (a *API) handleInternalWake(w http.ResponseWriter, r *http.Request) {
"the server failed to start and its automatic retries are spent; its owner can retry from the panel"))
return
}
if !a.limiter().allowed(name, a.WakeCooldown) {
policy, _, err := a.readWakePolicy(r.Context())
if err != nil {
writeError(w, r, err)
return
}
if !a.limiter().allowed(name, policy.cooldown(a.WakeCooldown)) {
writeError(w, r, newError(http.StatusTooManyRequests, "cooldown", "wake is cooling down, retry shortly"))
return
}
@@ -181,7 +186,7 @@ func (a *API) handleInternalWake(w http.ResponseWriter, r *http.Request) {
// treats 503 at_capacity as "cluster full, tell the player to try later" and does
// NOT enqueue them (nothing is coming up, so waiting would only strand them),
// distinct from the 429 cooldown's "already waking, keep waiting".
ok, err := a.withinRunningCap(r.Context(), info)
ok, err := a.withinRunningCap(r.Context(), info, policy.MaxRunningServers)
if err != nil {
writeError(w, r, err)
return
+7 -2
View File
@@ -46,13 +46,18 @@ func (a *API) handleWake(w http.ResponseWriter, r *http.Request) {
writeError(w, r, errServerRetiring)
return
}
if !a.limiter().allowed(name, a.WakeCooldown) {
policy, _, err := a.readWakePolicy(r.Context())
if err != nil {
writeError(w, r, err)
return
}
if !a.limiter().allowed(name, policy.cooldown(a.WakeCooldown)) {
writeError(w, r, newError(http.StatusTooManyRequests, "cooldown", "wake is cooling down, retry shortly"))
return
}
// Global running-server cap (spec §9.1). Distinct from the per-server cooldown:
// 503 at_capacity means the cluster is full, not that this server is throttled.
ok, err := a.withinRunningCap(r.Context(), info)
ok, err := a.withinRunningCap(r.Context(), info, policy.MaxRunningServers)
if err != nil {
writeError(w, r, err)
return
+109
View File
@@ -0,0 +1,109 @@
package api
import (
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"net/http"
"time"
)
const wakePolicyKey = "wake_policy"
type wakePolicy struct {
MaxRunningServers int `json:"maxRunningServers"`
WakeCooldownSeconds int `json:"wakeCooldownSeconds"`
}
type wakePolicyView struct {
wakePolicy
Revision string `json:"revision"`
Managed bool `json:"managed"`
}
func (a *API) readWakePolicy(ctx context.Context) (wakePolicyView, []byte, error) {
view := wakePolicyView{wakePolicy: wakePolicy{a.MaxRunningServers, int(a.WakeCooldown / time.Second)}}
raw, err := a.Repo.GetSetting(ctx, wakePolicyKey)
switch {
case errors.Is(err, ErrNotFound):
raw = nil
case err != nil:
return view, nil, err
default:
if err := json.Unmarshal(raw, &view.wakePolicy); err != nil {
return view, nil, err
}
if !view.wakePolicy.valid() {
return view, nil, errors.New("invalid persisted wake policy")
}
view.Managed = true
}
canonical, _ := json.Marshal(view.wakePolicy)
sum := sha256.Sum256(canonical)
view.Revision = hex.EncodeToString(sum[:])
return view, raw, nil
}
func (a *API) handleGetWakePolicy(w http.ResponseWriter, r *http.Request) {
view, _, err := a.readWakePolicy(r.Context())
if err != nil {
writeError(w, r, err)
return
}
writeJSON(w, http.StatusOK, view)
}
func (a *API) handleSetWakePolicy(w http.ResponseWriter, r *http.Request) {
if !a.requireReauth(w, r, principalFromContext(r.Context())) {
return
}
if err := requireJSONContentType(r); err != nil {
writeError(w, r, err)
return
}
var body struct {
wakePolicy
Revision string `json:"revision"`
}
if err := decodeJSON(w, r, &body); err != nil {
writeError(w, r, err)
return
}
if !body.wakePolicy.valid() {
writeError(w, r, newError(http.StatusBadRequest, "bad_request", "running limit must be 0–10000 and wake cooldown 0–3600 seconds"))
return
}
current, expected, err := a.readWakePolicy(r.Context())
if err != nil {
writeError(w, r, err)
return
}
if body.Revision != current.Revision {
writeError(w, r, newError(http.StatusConflict, "conflict", "platform policy changed; reload before saving"))
return
}
raw, _ := json.Marshal(body.wakePolicy)
if err := a.Repo.CompareAndSetSetting(r.Context(), wakePolicyKey, expected, raw); err != nil {
writeError(w, r, err)
return
}
a.audit(r, "platform.wake_policy", "platform")
view, _, err := a.readWakePolicy(r.Context())
if err != nil {
writeError(w, r, err)
return
}
writeJSON(w, http.StatusOK, view)
}
func (p wakePolicyView) cooldown(fallback time.Duration) time.Duration {
if !p.Managed {
return fallback
}
return time.Duration(p.WakeCooldownSeconds) * time.Second
}
func (p wakePolicy) valid() bool {
return p.MaxRunningServers >= 0 && p.MaxRunningServers <= 10000 && p.WakeCooldownSeconds >= 0 && p.WakeCooldownSeconds <= 3600
}
+97
View File
@@ -0,0 +1,97 @@
package api
import (
"context"
"encoding/json"
"errors"
"testing"
"time"
"felis.lolicon.best/internal/apis/felis/v1alpha1"
)
const wakePolicyPath = "/api/v1/settings/wake-policy"
func TestWakePolicyPersistenceAndAdmission(t *testing.T) {
repo := newFakeRepo()
cl := newFakeCluster()
cl.byName["survival"] = &ServerInfo{Name: "survival", AutostartPolicy: "public", Phase: "Stopped", DesiredState: "Stopped"}
cl.list = []ServerInfo{{Name: "other", DesiredState: "Running"}}
repo.byName["survival"] = &ServerRecord{Name: "survival", OwnerID: "owner"}
a := newTestAPI(repo, cl)
a.External = staticExternal{p: &Principal{UserID: "owner", Role: "owner", ViaAdminAccess: true}}
view, _, err := a.readWakePolicy(context.Background())
if err != nil || view.Managed {
t.Fatal(view, err)
}
save := func(view wakePolicyView) int {
body, _ := json.Marshal(map[string]any{"maxRunningServers": view.MaxRunningServers, "wakeCooldownSeconds": view.WakeCooldownSeconds, "revision": view.Revision})
return do(a.ExternalHandler(), "PUT", wakePolicyPath, string(body), jsonHeader).Code
}
view.MaxRunningServers = 1
view.WakeCooldownSeconds = 60
if code := save(view); code != 200 {
t.Fatal("save", code)
}
if code := save(view); code != 409 {
t.Fatal("stale write", code)
}
replica := newTestAPI(repo, cl)
replica.External = a.External
persisted, _, err := replica.readWakePolicy(context.Background())
if err != nil || !persisted.Managed || persisted.MaxRunningServers != 1 || persisted.cooldown(0) != time.Minute {
t.Fatal(persisted, err)
}
if w := do(replica.ExternalHandler(), "POST", "/api/v1/servers/survival/wake", "", nil); w.Code != 503 || decodeErr(t, w) != "at_capacity" {
t.Fatal("panel cap", w.Code, w.Body.String())
}
if w := internalWake(replica, `{"mc_uuid":"`+wakeUUID+`"}`); w.Code != 503 || decodeErr(t, w) != "at_capacity" {
t.Fatal("game cap", w.Code, w.Body.String())
}
if why, retry, err := replica.startScheduled(context.Background(), "survival"); err != nil || !retry || why == "" {
t.Fatal("scheduled cap", why, retry, err)
}
persisted.MaxRunningServers = 2
if code := save(persisted); code != 200 {
t.Fatal("increase cap", code)
}
if w := do(replica.ExternalHandler(), "POST", "/api/v1/servers/survival/wake", "", nil); w.Code != 202 || cl.desired["survival"] != v1alpha1.DesiredRunning {
t.Fatal("wake", w.Code, w.Body.String())
}
if w := internalWake(replica, `{"mc_uuid":"`+wakeUUID+`"}`); w.Code != 429 {
t.Fatal("shared saved cooldown", w.Code, w.Body.String())
}
repo.failGetSetting = errors.New("database down")
if w := do(replica.ExternalHandler(), "POST", "/api/v1/servers/survival/wake", "", nil); w.Code != 500 {
t.Fatal("outage failed open", w.Code)
}
repo.failGetSetting = nil
repo.settings[wakePolicyKey] = []byte(`{"maxRunningServers":-1}`)
if _, _, err := replica.readWakePolicy(context.Background()); err == nil {
t.Fatal("invalid persisted policy failed open")
}
}
func TestWakePolicyAuthorizationAndValidation(t *testing.T) {
a := newTestAPI(newFakeRepo(), newFakeCluster())
for _, p := range []*Principal{nil, {UserID: "admin", Role: "admin", ViaAdminAccess: true}, {UserID: "owner", Role: "owner"}} {
a.External = staticExternal{p: p}
for _, method := range []string{"GET", "PUT"} {
if w := do(a.ExternalHandler(), method, wakePolicyPath, `{}`, jsonHeader); w.Code != 401 && w.Code != 403 {
t.Fatalf("unauthorized %+v: %d", p, w.Code)
}
}
}
p := &Principal{UserID: "owner", Role: "owner", ViaAdminAccess: true, ViaSession: true, EmailVerified: true}
a.External = staticExternal{p: p}
a.Repo.(*fakeRepo).passkeyCreds["owner-key"] = PasskeyCredential{ID: "owner-key", UserID: p.UserID, UserVerified: true}
if w := do(a.ExternalHandler(), "PUT", wakePolicyPath, `{}`, jsonHeader); w.Code != 403 || decodeErr(t, w) != "reauth_required" {
t.Fatal("reauth not enforced", w.Code)
}
p.ReauthAt = a.now()
for _, body := range []string{`{"maxRunningServers":-1}`, `{"wakeCooldownSeconds":3601}`, `{"maxRunningServers":1.5}`} {
if w := do(a.ExternalHandler(), "PUT", wakePolicyPath, body, jsonHeader); w.Code != 400 {
t.Fatal("invalid policy", w.Code, w.Body.String())
}
}
}
+29 -1
View File
@@ -11,6 +11,7 @@ import (
"felis.lolicon.best/internal/maintenance"
"felis.lolicon.best/internal/naming"
"felis.lolicon.best/internal/placement"
appsv1 "k8s.io/api/apps/v1"
batchv1 "k8s.io/api/batch/v1"
corev1 "k8s.io/api/core/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
@@ -98,7 +99,34 @@ func (k *K8sCluster) GetServer(ctx context.Context, name string) (*ServerInfo, e
}
return nil, err
}
return serverInfo(&ms), nil
info := serverInfo(&ms)
if ms.Spec.DesiredState != v1alpha1.DesiredStopped && ms.Status.Phase != v1alpha1.PhaseStarting && ms.Status.Phase != v1alpha1.PhaseFailed && ms.Status.Phase != v1alpha1.PhaseRunning {
return info, nil
}
var pods corev1.PodList
if err := k.c.List(ctx, &pods, client.InNamespace(k.namespace), client.MatchingLabels{v1alpha1.LabelServer: name, v1alpha1.LabelComponent: gamePodComponent}); err != nil {
return nil, err
}
if ms.Spec.DesiredState == v1alpha1.DesiredStopped {
info.Phase, info.Ready = string(v1alpha1.PhaseStopped), false
var sts appsv1.StatefulSet
err := k.c.Get(ctx, types.NamespacedName{Namespace: k.namespace, Name: name}, &sts)
if err != nil && !apierrors.IsNotFound(err) {
return nil, err
}
if len(pods.Items) > 0 || (err == nil && (sts.Spec.Replicas == nil || *sts.Spec.Replicas != 0)) {
info.Phase = string(v1alpha1.PhaseStopping)
}
return info, nil
}
diagnostic := startupStatus(&ms, pods.Items)
if ms.Status.Phase != v1alpha1.PhaseRunning || diagnostic.Stage == "failed" || diagnostic.Stage == "scheduling" || diagnostic.Stage == "creating" {
info.Startup = diagnostic
if ms.Status.Phase == v1alpha1.PhaseRunning {
info.Ready, info.Phase = false, string(v1alpha1.PhaseFailed)
}
}
return info, nil
}
// WorldVolumeExists reads the world PVC the operator's StatefulSet
+3
View File
@@ -12,6 +12,7 @@ import (
"felis.lolicon.best/internal/apis/felis/v1alpha1"
"felis.lolicon.best/internal/maintenance"
"felis.lolicon.best/internal/naming"
appsv1 "k8s.io/api/apps/v1"
batchv1 "k8s.io/api/batch/v1"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/resource"
@@ -114,6 +115,8 @@ func TestCreateServerDefaultsIdleStop(t *testing.T) {
if err := v1alpha1.AddToScheme(scheme); err != nil {
t.Fatalf("scheme: %v", err)
}
corev1.AddToScheme(scheme)
appsv1.AddToScheme(scheme)
c := fake.NewClientBuilder().WithScheme(scheme).Build()
k := NewK8sCluster(c, "minecraft")
if err := k.CreateServer(context.Background(), CreateServerInput{
+1 -1
View File
@@ -359,7 +359,7 @@ func (k *K8sLogStreamer) StreamLogs(ctx context.Context, name string) (io.ReadCl
// the <name>-0 pod naming. The name is DNS-1123-validated upstream, so it is a
// safe label-selector value.
pods, err := k.clientset.CoreV1().Pods(k.namespace).List(ctx, metav1.ListOptions{
LabelSelector: v1alpha1.LabelServer + "=" + name,
LabelSelector: v1alpha1.LabelServer + "=" + name + "," + v1alpha1.LabelComponent + "=" + gamePodComponent,
})
if err != nil {
return nil, ErrConsoleUnavailable
+5 -1
View File
@@ -431,7 +431,11 @@ func (a *API) startScheduled(ctx context.Context, name string) (why string, retr
if err != nil {
return "", false, err
}
ok, err := a.withinRunningCap(ctx, info)
policy, _, err := a.readWakePolicy(ctx)
if err != nil {
return "", false, err
}
ok, err := a.withinRunningCap(ctx, info, policy.MaxRunningServers)
if err != nil {
return "", false, err
}
+82
View File
@@ -0,0 +1,82 @@
package api
import (
"fmt"
"slices"
"felis.lolicon.best/internal/apis/felis/v1alpha1"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/meta"
)
type StartupStatus struct {
Stage string `json:"stage"`
Reason string `json:"reason,omitempty"`
Message string `json:"message,omitempty"`
StartedAt string `json:"startedAt,omitempty"`
LogsAvailable bool `json:"logsAvailable"`
}
func startupStatus(ms *v1alpha1.MinecraftServer, pods []corev1.Pod) *StartupStatus {
s := &StartupStatus{Stage: "creating"}
if ms.Status.Phase == v1alpha1.PhaseFailed {
s.Stage = "failed"
}
if ms.Status.StartRequestedAt != nil {
s.StartedAt = ms.Status.StartRequestedAt.UTC().Format("2006-01-02T15:04:05Z")
}
if c := meta.FindStatusCondition(ms.Status.Conditions, v1alpha1.ConditionReady); c != nil {
s.Reason, s.Message = c.Reason, c.Message
}
var p *corev1.Pod
for i := range pods {
if pods[i].DeletionTimestamp == nil && (p == nil || pods[i].CreationTimestamp.After(p.CreationTimestamp.Time)) {
p = &pods[i]
}
}
if p == nil {
if ms.Status.Phase == v1alpha1.PhaseRunning {
s.Stage, s.Reason, s.Message = "failed", "GamePodMissing", "game Pod is missing; recorded Running status is no longer confirmed"
}
return s
}
s.Stage = "preparing"
if p.Status.Phase == corev1.PodFailed {
s.Stage, s.Reason, s.Message = "failed", p.Status.Reason, p.Status.Message
}
for _, c := range p.Status.Conditions {
if c.Type == corev1.PodReady && (c.Status == corev1.ConditionUnknown || (ms.Status.Phase == v1alpha1.PhaseRunning && c.Status == corev1.ConditionFalse && p.Status.Phase == corev1.PodRunning)) {
s.Stage, s.Reason, s.Message = "failed", c.Reason, c.Message
return s
}
if c.Type == corev1.PodScheduled && c.Status == corev1.ConditionFalse {
s.Stage, s.Reason, s.Message = "scheduling", c.Reason, c.Message
return s
}
}
s.LogsAvailable = p.Status.Phase == corev1.PodRunning
for _, c := range slices.Concat(p.Status.InitContainerStatuses, p.Status.ContainerStatuses) {
if c.State.Waiting != nil {
switch c.State.Waiting.Reason {
case "CrashLoopBackOff", "ImagePullBackOff", "ErrImagePull", "CreateContainerConfigError", "RunContainerError":
s.Stage, s.Reason, s.Message = "failed", c.State.Waiting.Reason, c.State.Waiting.Message
default:
if s.Stage != "failed" {
s.Reason, s.Message = c.State.Waiting.Reason, c.State.Waiting.Message
}
}
}
if c.State.Terminated != nil && (c.State.Terminated.ExitCode != 0 || c.Name == serverLogContainer) {
s.Stage, s.Reason, s.Message = "failed", c.State.Terminated.Reason, c.State.Terminated.Message
if s.Message == "" {
s.Message = fmt.Sprintf("Container %s exited with code %d", c.Name, c.State.Terminated.ExitCode)
}
}
if c.Name == serverLogContainer && c.State.Running != nil {
if s.Stage != "failed" {
s.Stage = "booting"
}
}
}
return s
}
+78
View File
@@ -0,0 +1,78 @@
package api
import (
"context"
"strings"
"testing"
"time"
"felis.lolicon.best/internal/apis/felis/v1alpha1"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"sigs.k8s.io/controller-runtime/pkg/client/fake"
)
func TestStartupDiagnostics(t *testing.T) {
started := metav1.NewTime(time.Now().UTC().Truncate(time.Second))
ms := &v1alpha1.MinecraftServer{ObjectMeta: metav1.ObjectMeta{Name: "survival", Namespace: "minecraft"}, Status: v1alpha1.MinecraftServerStatus{Phase: v1alpha1.PhaseStarting, StartRequestedAt: &started}}
cases := []struct {
name, stage, reason string
status corev1.PodStatus
logs bool
}{
{name: "scheduling memory", stage: "scheduling", reason: "Unschedulable", status: corev1.PodStatus{Conditions: []corev1.PodCondition{{Type: corev1.PodScheduled, Status: corev1.ConditionFalse, Reason: "Unschedulable", Message: "Insufficient memory"}}}},
{name: "image failure", stage: "failed", reason: "ImagePullBackOff", status: corev1.PodStatus{ContainerStatuses: []corev1.ContainerStatus{{Name: serverLogContainer, State: corev1.ContainerState{Waiting: &corev1.ContainerStateWaiting{Reason: "ImagePullBackOff", Message: "image unavailable"}}}}}},
{name: "process crash", stage: "failed", reason: "OOMKilled", status: corev1.PodStatus{ContainerStatuses: []corev1.ContainerStatus{{Name: serverLogContainer, State: corev1.ContainerState{Terminated: &corev1.ContainerStateTerminated{Reason: "OOMKilled", ExitCode: 137}}}}}},
{name: "evicted", stage: "failed", reason: "Evicted", status: corev1.PodStatus{Phase: corev1.PodFailed, Reason: "Evicted", Message: "node memory pressure"}},
{name: "init failure", stage: "failed", reason: "CrashLoopBackOff", status: corev1.PodStatus{InitContainerStatuses: []corev1.ContainerStatus{{Name: "prepare", State: corev1.ContainerState{Waiting: &corev1.ContainerStateWaiting{Reason: "CrashLoopBackOff", Message: "preparation failed"}}}}, ContainerStatuses: []corev1.ContainerStatus{{Name: serverLogContainer, State: corev1.ContainerState{Waiting: &corev1.ContainerStateWaiting{Reason: "PodInitializing"}}}}}},
{name: "node lost", stage: "failed", reason: "NodeNotReady", status: corev1.PodStatus{Phase: corev1.PodRunning, Conditions: []corev1.PodCondition{{Type: corev1.PodReady, Status: corev1.ConditionUnknown, Reason: "NodeNotReady", Message: "node heartbeat lost"}}}},
{name: "booting", stage: "booting", logs: true, status: corev1.PodStatus{Phase: corev1.PodRunning, ContainerStatuses: []corev1.ContainerStatus{{Name: serverLogContainer, State: corev1.ContainerState{Running: &corev1.ContainerStateRunning{}}}}}},
}
scheme := runtime.NewScheme()
v1alpha1.AddToScheme(scheme)
corev1.AddToScheme(scheme)
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
pod := &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Name: "survival-0", Namespace: "minecraft", Labels: map[string]string{v1alpha1.LabelServer: "survival", v1alpha1.LabelComponent: gamePodComponent}}, Status: tc.status}
worker := &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Name: "file-worker", Namespace: "minecraft", CreationTimestamp: started, Labels: map[string]string{v1alpha1.LabelServer: "survival", v1alpha1.LabelComponent: "file-op"}}, Status: corev1.PodStatus{Phase: corev1.PodRunning}}
c := fake.NewClientBuilder().WithScheme(scheme).WithObjects(ms, pod, worker).Build()
info, err := NewK8sCluster(c, "minecraft").GetServer(context.Background(), "survival")
if err != nil {
t.Fatal(err)
}
s := info.Startup
if s == nil || s.Stage != tc.stage || s.Reason != tc.reason || s.LogsAvailable != tc.logs || s.StartedAt != started.UTC().Format(time.RFC3339) {
t.Fatalf("diagnostics=%+v", s)
}
if tc.reason == "OOMKilled" && !strings.Contains(s.Message, "137") {
t.Fatal("missing exit code", s.Message)
}
if tc.stage == "failed" {
msRunning := ms.DeepCopy()
msRunning.Status.Phase = v1alpha1.PhaseRunning
cRunning := fake.NewClientBuilder().WithScheme(scheme).WithObjects(msRunning, pod).Build()
actual, err := NewK8sCluster(cRunning, "minecraft").GetServer(context.Background(), "survival")
if err != nil || actual.Phase != "Failed" || actual.Ready || actual.Startup == nil {
t.Fatal("stale operator Running concealed runtime failure", actual, err)
}
}
if publicServerInfo(info).Startup != nil {
t.Fatal("public status leaked diagnostics")
}
})
}
running := ms.DeepCopy()
running.Status.Phase = v1alpha1.PhaseRunning
if s := startupStatus(running, nil); s.Stage != "failed" || s.Reason != "GamePodMissing" {
t.Fatal("missing game runtime concealed", s)
}
s := startupStatus(ms, nil)
if s.Stage != "creating" || s.LogsAvailable {
t.Fatalf("missing pod=%+v", s)
}
deleted := corev1.Pod{ObjectMeta: metav1.ObjectMeta{DeletionTimestamp: &started}, Status: cases[0].status}
if s := startupStatus(ms, []corev1.Pod{deleted}); s.Stage != "creating" {
t.Fatalf("terminating pod reused=%+v", s)
}
}