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

fix(setup): 系统服与副本 Secret 改用带乐观锁的 merge-patch 并冲突重试,安装器在构建变化时把登录与大厅固定到 digest

parent b384f628
Loading
Loading
Loading
Loading
+358 −0
Changes for cmd/felis/conflict_test.go: 358 added lines, 0 removed lines.
Original line number Diff line number Diff line
package main

import (
	"context"
	"net/http"
	"net/http/httptest"
	"slices"
	"strings"
	"testing"

	"felis.lolicon.best/internal/apis/felis/v1alpha1"
	"felis.lolicon.best/internal/imagepin"
	"felis.lolicon.best/internal/naming"
	corev1 "k8s.io/api/core/v1"
	metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
	"sigs.k8s.io/controller-runtime/pkg/client"
	"sigs.k8s.io/controller-runtime/pkg/client/fake"
	"sigs.k8s.io/controller-runtime/pkg/client/interceptor"
)

// racingClient lands one concurrent write (race) on the stored object just before
// setup's first write of it goes out, the way the operator's status update or
// felis-api's idle patch can, and counts setup's writes.
func racingClient(t *testing.T, race func(ctx context.Context, c client.WithWatch), objs ...client.Object) (client.Client, *int) {
	t.Helper()
	writes := 0
	before := func(ctx context.Context, c client.WithWatch) {
		writes++
		if writes == 1 {
			race(ctx, c)
		}
	}
	cl := fake.NewClientBuilder().WithScheme(newSystemServerScheme(t)).WithObjects(objs...).
		WithStatusSubresource(&v1alpha1.MinecraftServer{}).
		WithInterceptorFuncs(interceptor.Funcs{
			Update: func(ctx context.Context, c client.WithWatch, obj client.Object, opts ...client.UpdateOption) error {
				before(ctx, c)
				return c.Update(ctx, obj, opts...)
			},
			Patch: func(ctx context.Context, c client.WithWatch, obj client.Object, patch client.Patch, opts ...client.PatchOption) error {
				before(ctx, c)
				return c.Patch(ctx, obj, patch, opts...)
			},
		}).Build()
	return cl, &writes
}

func staleLoginGate(t *testing.T) *v1alpha1.MinecraftServer {
	t.Helper()
	ms, err := loginSystemServer("felis-limbo:demo", "minecraft",
		"http://old.internal:8081", "203.0.113.10.nip.io", "console.203.0.113.10.nip.io")
	if err != nil {
		t.Fatal(err)
	}
	return ms
}

func getLogin(t *testing.T, cl client.Client) *v1alpha1.MinecraftServer {
	t.Helper()
	var ms v1alpha1.MinecraftServer
	if err := cl.Get(context.Background(), client.ObjectKey{Namespace: "minecraft", Name: naming.SystemLoginServer}, &ms); err != nil {
		t.Fatal(err)
	}
	return &ms
}

func envMap(ms *v1alpha1.MinecraftServer) map[string]string {
	m := map[string]string{}
	for _, e := range ms.Spec.Env {
		m[e.Name] = e.Value
	}
	return m
}

// A hand edit that lands while setup refreshes the console hostnames costs setup a
// re-read and a second write; both the edit and the refresh survive.
func TestRefreshDerivedEnvRetriesAConcurrentEdit(t *testing.T) {
	ctx := context.Background()
	cl, writes := racingClient(t, func(ctx context.Context, c client.WithWatch) {
		ms := getLogin(t, c)
		ms.Spec.Env = append(ms.Spec.Env, v1alpha1.EnvVar{Name: "HAND_TUNED", Value: "1"})
		if err := c.Update(ctx, ms); err != nil {
			t.Fatal(err)
		}
	}, staleLoginGate(t))

	desired, err := loginSystemServer("felis-limbo:demo", "minecraft",
		"http://felis-api-internal.felis.svc.cluster.local:8081", "mc.example.net", "console.mc.example.net")
	if err != nil {
		t.Fatal(err)
	}
	refreshed, err := refreshDerivedEnv(ctx, cl, getLogin(t, cl), desired)
	if err != nil || !refreshed {
		t.Fatalf("refreshed=%v err=%v, want a refresh after the retry", refreshed, err)
	}
	env := envMap(getLogin(t, cl))
	if env[envPanelHostname] != "console.mc.example.net" || env[envRootDomain] != "mc.example.net" {
		t.Errorf("env = %v, want the new hostnames", env)
	}
	if env["HAND_TUNED"] != "1" {
		t.Errorf("env = %v: the concurrent hand edit was dropped", env)
	}
	if *writes != 2 {
		t.Errorf("writes = %d, want 2 (one conflict, one retry)", *writes)
	}
}

// converge reports each fill once even when a status write forced a retry.
func TestConvergeRetriesAConcurrentStatusWrite(t *testing.T) {
	ctx := context.Background()
	lobby, err := lobbySystemServer("reg/lobby:1", "minecraft")
	if err != nil {
		t.Fatal(err)
	}
	lobby.Spec.Rcon = v1alpha1.RconSpec{}
	cl, writes := racingClient(t, func(ctx context.Context, c client.WithWatch) {
		var ms v1alpha1.MinecraftServer
		if err := c.Get(ctx, client.ObjectKey{Namespace: "minecraft", Name: naming.SystemLobbyServer}, &ms); err != nil {
			t.Fatal(err)
		}
		ms.Status.Phase = v1alpha1.PhaseRunning
		if err := c.Status().Update(ctx, &ms); err != nil {
			t.Fatal(err)
		}
	}, lobby)

	var got systemServerOutcome
	for _, o := range convergeSystemServers(ctx, cl, "minecraft", "", "reg/lobby:1",
		"http://felis-api-internal.felis.svc.cluster.local:8081", "mc.example.net", "console.mc.example.net") {
		if o.name == naming.SystemLobbyServer {
			got = o
		}
	}
	if got.err != nil || !got.updated || !slices.Equal(got.changes, []string{"spec.rcon"}) {
		t.Fatalf("lobby outcome = %+v, want spec.rcon filled once", got)
	}
	var ms v1alpha1.MinecraftServer
	if err := cl.Get(ctx, client.ObjectKey{Namespace: "minecraft", Name: naming.SystemLobbyServer}, &ms); err != nil {
		t.Fatal(err)
	}
	if !ms.Spec.Rcon.Enabled || ms.Status.Phase != v1alpha1.PhaseRunning {
		t.Errorf("rcon.enabled=%v phase=%q, want the fill and the status write both kept", ms.Spec.Rcon.Enabled, ms.Status.Phase)
	}
	if *writes != 2 {
		t.Errorf("writes = %d, want 2", *writes)
	}
}

// A replica refresh racing another writer keeps that writer's key.
func TestSecretReplicaRefreshRetriesAConcurrentWrite(t *testing.T) {
	secret := func(ns, body string) *corev1.Secret {
		return &corev1.Secret{ObjectMeta: metav1.ObjectMeta{Name: "felis-config", Namespace: ns},
			Data: map[string][]byte{"felis.toml": []byte(body)}}
	}
	cl, writes := racingClient(t, func(ctx context.Context, c client.WithWatch) {
		var s corev1.Secret
		if err := c.Get(ctx, client.ObjectKey{Namespace: "minecraft", Name: "felis-config"}, &s); err != nil {
			t.Fatal(err)
		}
		s.Data["extra"] = []byte("x")
		if err := c.Update(ctx, &s); err != nil {
			t.Fatal(err)
		}
	}, secret("felis", "current"), secret("minecraft", "stale"))

	out := ensureSecretReplica(context.Background(), cl, "felis", "minecraft",
		"felis-config", "felis.toml", "config", "minecraft ns", true)
	if out.err != nil || !out.updated {
		t.Fatalf("outcome = %+v, want refreshed", out)
	}
	var s corev1.Secret
	if err := cl.Get(context.Background(), client.ObjectKey{Namespace: "minecraft", Name: "felis-config"}, &s); err != nil {
		t.Fatal(err)
	}
	if string(s.Data["felis.toml"]) != "current" || string(s.Data["extra"]) != "x" {
		t.Errorf("replica data = %q, want felis.toml=current and extra=x", s.Data)
	}
	if *writes != 2 {
		t.Errorf("writes = %d, want 2", *writes)
	}
}

const (
	sysOldDigest = "sha256:1111111111111111111111111111111111111111111111111111111111111111"
	sysNewDigest = "sha256:4444444444444444444444444444444444444444444444444444444444444444"
)

func systemPinRegistry(t *testing.T) imagepin.Resolver {
	t.Helper()
	reg := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
		switch r.URL.Path {
		case "/v2/felis/limbo/manifests/demo", "/v2/felis/lobby/manifests/demo":
			w.Header().Set("Docker-Content-Digest", sysNewDigest)
		default:
			http.NotFound(w, r)
		}
	}))
	t.Cleanup(reg.Close)
	return imagepin.Resolver{Registry: defaultRegistryURL, Endpoint: strings.TrimPrefix(reg.URL, "http://")}
}

func systemServer(name, image string) *v1alpha1.MinecraftServer {
	ms := &v1alpha1.MinecraftServer{}
	ms.Name, ms.Namespace = name, "minecraft"
	ms.Labels = map[string]string{v1alpha1.LabelSystemRole: name}
	ms.Spec.Image = image
	return ms
}

// The installer's --system pin moves a system server onto the build its tag names
// now, from the bare tag or from an earlier digest, and is a no-op the second time.
func TestPinSystemServerImage(t *testing.T) {
	ctx := context.Background()
	res := systemPinRegistry(t)
	limbo := defaultRegistryURL + "/felis/limbo:demo"
	lobby := defaultRegistryURL + "/felis/lobby:demo"
	cl := fake.NewClientBuilder().WithScheme(newSystemServerScheme(t)).WithObjects(
		systemServer(naming.SystemLoginServer, limbo),
		systemServer(naming.SystemLobbyServer, lobby+"@"+sysOldDigest),
	).Build()
	image := func(name string) string {
		var ms v1alpha1.MinecraftServer
		if err := cl.Get(ctx, client.ObjectKey{Namespace: "minecraft", Name: name}, &ms); err != nil {
			t.Fatal(err)
		}
		return ms.Spec.Image
	}

	for name, want := range map[string]string{
		naming.SystemLoginServer: limbo + "@" + sysNewDigest,
		naming.SystemLobbyServer: lobby + "@" + sysNewDigest,
	} {
		o, err := pinSystemServerImage(ctx, cl, "minecraft", name, res)
		if err != nil || o.err != nil || !o.updated {
			t.Fatalf("%s: outcome=%+v err=%v, want pinned", name, o, err)
		}
		if got := image(name); got != want {
			t.Errorf("%s image = %q, want %q", name, got, want)
		}
		again, err := pinSystemServerImage(ctx, cl, "minecraft", name, res)
		if err != nil || again.updated || again.skipped != "already runs "+want {
			t.Errorf("%s second pass = %+v, %v; want already runs %s", name, again, err, want)
		}
	}
}

func TestPinSystemServerImageLeavesOthersAlone(t *testing.T) {
	ctx := context.Background()
	res := systemPinRegistry(t)
	pin := func(t *testing.T, ms *v1alpha1.MinecraftServer) (systemServerOutcome, string) {
		t.Helper()
		cl := fake.NewClientBuilder().WithScheme(newSystemServerScheme(t)).WithObjects(ms).Build()
		o, err := pinSystemServerImage(ctx, cl, "minecraft", naming.SystemLoginServer, res)
		if err != nil {
			t.Fatal(err)
		}
		var got v1alpha1.MinecraftServer
		if err := cl.Get(ctx, client.ObjectKeyFromObject(ms), &got); err != nil {
			t.Fatal(err)
		}
		return o, got.Spec.Image
	}

	t.Run("external image", func(t *testing.T) {
		o, img := pin(t, systemServer(naming.SystemLoginServer, "docker.io/example/limbo:1.2"))
		if o.err != nil || o.updated || img != "docker.io/example/limbo:1.2" ||
			o.skipped != "runs docker.io/example/limbo:1.2, which names no platform registry tag to follow; left alone" {
			t.Fatalf("outcome=%+v image=%q, want left alone", o, img)
		}
	})
	t.Run("digest without a tag", func(t *testing.T) {
		ref := defaultRegistryURL + "/felis/limbo@" + sysOldDigest
		o, img := pin(t, systemServer(naming.SystemLoginServer, ref))
		if o.err != nil || o.updated || img != ref {
			t.Fatalf("outcome=%+v image=%q, want left alone", o, img)
		}
	})
	t.Run("not a system server", func(t *testing.T) {
		ms := systemServer(naming.SystemLoginServer, defaultRegistryURL+"/felis/limbo:demo")
		ms.Labels = nil
		o, img := pin(t, ms)
		if o.err == nil || img != defaultRegistryURL+"/felis/limbo:demo" {
			t.Fatalf("outcome=%+v image=%q, want refused", o, img)
		}
	})
	t.Run("tag the registry lost", func(t *testing.T) {
		o, img := pin(t, systemServer(naming.SystemLoginServer, defaultRegistryURL+"/felis/limbo:gone"))
		if o.err == nil || img != defaultRegistryURL+"/felis/limbo:gone" {
			t.Fatalf("outcome=%+v image=%q, want an error the installer falls back on", o, img)
		}
	})
	t.Run("absent", func(t *testing.T) {
		cl := fake.NewClientBuilder().WithScheme(newSystemServerScheme(t)).Build()
		o, err := pinSystemServerImage(ctx, cl, "minecraft", naming.SystemLoginServer, res)
		if err != nil || o.err != nil || o.skipped != "not present yet; sudo felis setup creates it" {
			t.Fatalf("outcome=%+v err=%v, want a skip", o, err)
		}
	})
	t.Run("retargeted while pinning", func(t *testing.T) {
		cl, _ := racingClient(t, func(ctx context.Context, c client.WithWatch) {
			ms := getLogin(t, c)
			ms.Spec.Image = "docker.io/example/limbo:1.2"
			if err := c.Update(ctx, ms); err != nil {
				t.Fatal(err)
			}
		}, systemServer(naming.SystemLoginServer, defaultRegistryURL+"/felis/limbo:demo"))
		o, err := pinSystemServerImage(ctx, cl, "minecraft", naming.SystemLoginServer, res)
		if err != nil || o.err != nil || o.updated {
			t.Fatalf("outcome=%+v err=%v, want nothing written", o, err)
		}
		if img := getLogin(t, cl).Spec.Image; img != "docker.io/example/limbo:1.2" {
			t.Errorf("image = %q, want the admin's retarget kept", img)
		}
	})
}

// --system names one of the two system servers; anything else is a usage error
// before any cluster is touched.
func TestPinImagesSystemFlagTakesOnlySystemServers(t *testing.T) {
	var stdout, stderr strings.Builder
	if code := cmdPinImages([]string{"--system", "survival"}, &stdout, &stderr); code != 2 {
		t.Fatalf("exit = %d, want 2", code)
	}
	if got := stderr.String(); got != "felis pin-images: --system takes login or lobby, not \"survival\"\n" {
		t.Errorf("stderr = %q", got)
	}
}

// An idle setting the panel saves while converge fills the default is kept.
func TestConvergeIdleKeepsAConcurrentPanelEdit(t *testing.T) {
	ctx := context.Background()
	srv := &v1alpha1.MinecraftServer{}
	srv.Name, srv.Namespace = "survival", "minecraft"
	cl, writes := racingClient(t, func(ctx context.Context, c client.WithWatch) {
		var ms v1alpha1.MinecraftServer
		if err := c.Get(ctx, client.ObjectKey{Namespace: "minecraft", Name: "survival"}, &ms); err != nil {
			t.Fatal(err)
		}
		ms.Spec.Idle = v1alpha1.IdleSpec{AutoStopEnabled: false, EmptySecondsBeforeStop: 1800}
		if err := c.Update(ctx, &ms); err != nil {
			t.Fatal(err)
		}
	}, srv)

	if out := convergeUserServerIdle(ctx, cl, "minecraft"); len(out) != 0 {
		t.Fatalf("outcomes = %+v, want none: the server has a setting by the time converge writes", out)
	}
	var ms v1alpha1.MinecraftServer
	if err := cl.Get(ctx, client.ObjectKey{Namespace: "minecraft", Name: "survival"}, &ms); err != nil {
		t.Fatal(err)
	}
	if want := (v1alpha1.IdleSpec{AutoStopEnabled: false, EmptySecondsBeforeStop: 1800}); ms.Spec.Idle != want {
		t.Errorf("idle = %+v, want the panel's %+v", ms.Spec.Idle, want)
	}
	if *writes != 1 {
		t.Errorf("writes = %d, want 1 (the conflicted attempt; the retry sends nothing)", *writes)
	}
}
+12 −2
Changes for cmd/felis/converge.go: 12 added lines, 2 removed lines.
Original line number Diff line number Diff line
@@ -92,12 +92,22 @@ func convergeUserServerIdle(ctx context.Context, cl client.Client, namespace str
		if ms.Labels[v1alpha1.LabelSystemRole] != "" || ms.Spec.Idle != (v1alpha1.IdleSpec{}) {
			continue
		}
		patch := client.MergeFrom(ms.DeepCopy())
		// Re-checked on the copy each attempt reads: an idle setting the panel saved
		// meanwhile is the user's, and the default must not land over it.
		changed, err := patchOnConflictRetry(ctx, cl, ms, func() bool {
			if ms.Spec.Idle != (v1alpha1.IdleSpec{}) {
				return false
			}
			ms.Spec.Idle = v1alpha1.DefaultIdle()
		if err := cl.Patch(ctx, ms, patch); err != nil {
			return true
		})
		if err != nil {
			out = append(out, systemServerOutcome{name: ms.Name, err: fmt.Errorf("converge %s: %w", ms.Name, err)})
			continue
		}
		if !changed {
			continue
		}
		out = append(out, systemServerOutcome{name: ms.Name, available: true, updated: true,
			changes: []string{fmt.Sprintf("spec.idle (stop after %ds empty)", v1alpha1.DefaultEmptySecondsBeforeStop)}})
	}
+91 −2
Changes for cmd/felis/pinimages.go: 91 added lines, 2 removed lines.
Original line number Diff line number Diff line
@@ -11,7 +11,9 @@ import (

	"felis.lolicon.best/internal/apis/felis/v1alpha1"
	"felis.lolicon.best/internal/imagepin"
	"felis.lolicon.best/internal/naming"
	"felis.lolicon.best/internal/platform"
	apierrors "k8s.io/apimachinery/pkg/api/errors"
	"k8s.io/apimachinery/pkg/api/meta"
	"sigs.k8s.io/controller-runtime/pkg/client"
)
@@ -35,12 +37,17 @@ func cmdPinImages(args []string, stdout, stderr io.Writer) int {
	namespace := fs.String("namespace", platform.DefaultMinecraftNamespace, "namespace the MinecraftServers live in")
	registry := fs.String("registry", defaultRegistryURL, "registry host[:port] the image refs spell")
	endpoint := fs.String("endpoint", "", "host[:port] to reach the registry at (default: 127.0.0.1 on the registry's port, its hostPort on this node)")
	system := fs.String("system", "", "pin this system server ("+naming.SystemLoginServer+" or "+naming.SystemLobbyServer+") to the build its tag names now, instead of the user servers")
	if err := fs.Parse(args); err != nil {
		if errors.Is(err, flag.ErrHelp) {
			return 0
		}
		return 2
	}
	if *system != "" && *system != naming.SystemLoginServer && *system != naming.SystemLobbyServer {
		fmt.Fprintf(stderr, "felis pin-images: --system takes %s or %s, not %q\n", naming.SystemLoginServer, naming.SystemLobbyServer, *system)
		return 2
	}
	if *endpoint == "" {
		*endpoint = loopbackEndpoint(*registry)
	}
@@ -51,6 +58,9 @@ func cmdPinImages(args []string, stdout, stderr io.Writer) int {
	}
	ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute)
	defer cancel()
	if *system != "" {
		return reportSystemPin(ctx, cl, *namespace, *system, imagepin.Resolver{Registry: *registry, Endpoint: *endpoint}, stdout, stderr)
	}
	outcomes, err := pinUserServerImages(ctx, cl, *namespace, imagepin.Resolver{Registry: *registry, Endpoint: *endpoint})
	if meta.IsNoMatchError(err) {
		fmt.Fprintln(stdout, "felis pin-images: no MinecraftServer CRD yet, so no server to pin")
@@ -77,6 +87,84 @@ func cmdPinImages(args []string, stdout, stderr io.Writer) int {
	return exit
}

// reportSystemPin runs pinSystemServerImage for `felis pin-images --system` and
// prints what it did. Only a failed pin exits non-zero: the installer falls back
// to restarting the pod on its tag then.
func reportSystemPin(ctx context.Context, cl client.Client, namespace, name string, r imagepin.Resolver, stdout, stderr io.Writer) int {
	o, err := pinSystemServerImage(ctx, cl, namespace, name, r)
	switch {
	case meta.IsNoMatchError(err):
		fmt.Fprintln(stdout, "felis pin-images: no MinecraftServer CRD yet, so no server to pin")
		return 0
	case err == nil && o.err != nil:
		err = o.err
	}
	if err != nil {
		fmt.Fprintf(stderr, "felis pin-images: %s: %v\n", name, err)
		return 1
	}
	switch {
	case o.updated:
		fmt.Fprintf(stdout, "felis pin-images: %s: %s; the operator rolls it onto that build\n", name, strings.Join(o.changes, ", "))
	default:
		fmt.Fprintf(stdout, "felis pin-images: %s: %s\n", name, o.skipped)
	}
	return 0
}

// pinSystemServerImage fixes a system server to the build its image tag names now,
// replacing the digest of an earlier build. The installer runs it after pushing a
// rebuilt login or lobby image, and the operator rolls the StatefulSet onto the new
// ref, so the build a system server runs is written in its spec and moves only when
// a build did. An image outside the platform registry, or one naming no tag to
// follow, is the admin's choice and is left alone; so is a server whose image an
// admin retargets while this runs.
func pinSystemServerImage(ctx context.Context, cl client.Client, namespace, name string, r imagepin.Resolver) (systemServerOutcome, error) {
	var ms v1alpha1.MinecraftServer
	if err := cl.Get(ctx, client.ObjectKey{Namespace: namespace, Name: name}, &ms); err != nil {
		if apierrors.IsNotFound(err) {
			return systemServerOutcome{name: name, skipped: "not present yet; sudo felis setup creates it"}, nil
		}
		return systemServerOutcome{}, err
	}
	if ms.Labels[v1alpha1.LabelSystemRole] != name {
		return systemServerOutcome{name: name, err: fmt.Errorf(
			"MinecraftServer %s/%s is not marked as the Felis %q system server; left alone", namespace, name, name)}, nil
	}
	tagged := withoutDigest(ms.Spec.Image)
	if !r.Covers(tagged) || !strings.Contains(tagged[strings.LastIndex(tagged, "/")+1:], ":") {
		return systemServerOutcome{name: name, available: true, skipped: "runs " + ms.Spec.Image +
			", which names no platform registry tag to follow; left alone"}, nil
	}
	pinned, err := r.Pin(ctx, tagged)
	if err != nil {
		return systemServerOutcome{name: name, err: fmt.Errorf("resolve %s: %w", tagged, err)}, nil
	}
	changed, err := patchOnConflictRetry(ctx, cl, &ms, func() bool {
		if withoutDigest(ms.Spec.Image) != tagged || ms.Spec.Image == pinned {
			return false
		}
		ms.Spec.Image = pinned
		return true
	})
	if err != nil {
		return systemServerOutcome{name: name, err: fmt.Errorf("patch %s: %w", name, err)}, nil
	}
	if !changed {
		return systemServerOutcome{name: name, available: true, skipped: "already runs " + ms.Spec.Image}, nil
	}
	return systemServerOutcome{name: name, available: true, updated: true,
		changes: []string{"spec.image pinned to " + pinned}}, nil
}

// withoutDigest drops the @sha256:… of a pinned ref, leaving the tag it came from.
func withoutDigest(ref string) string {
	if i := strings.Index(ref, "@"); i >= 0 {
		return ref[:i]
	}
	return ref
}

// loopbackEndpoint is the registry's port on 127.0.0.1: the registry Deployment
// binds it as a hostPort, and containerd's mirror and the installer's pushes use
// the same address.
@@ -88,8 +176,9 @@ func loopbackEndpoint(registry string) string {
}

// pinUserServerImages patches spec.image of every user server whose image the
// resolver covers and is not yet pinned. System servers are left on their tags:
// the installer rebuilds and restarts them on purpose (restart_existing_system_servers).
// resolver covers and is not yet pinned. System servers are the installer's to move:
// it re-pins them with --system when it rolls them onto a new build
// (restart_existing_system_servers).
// A server that is already pinned, or runs an image from elsewhere, produces no
// outcome, so a pinned fleet reports nothing. A running server restarts once as
// the operator rolls its StatefulSet onto the pinned ref, which is the build it
+50 −13
Changes for cmd/felis/systemservers.go: 50 added lines, 13 removed lines.
Original line number Diff line number Diff line
@@ -16,6 +16,7 @@ import (
	"k8s.io/apimachinery/pkg/runtime"
	clientgoscheme "k8s.io/client-go/kubernetes/scheme"
	"k8s.io/client-go/tools/clientcmd"
	"k8s.io/client-go/util/retry"
	ctrl "sigs.k8s.io/controller-runtime"
	"sigs.k8s.io/controller-runtime/pkg/client"
)
@@ -392,7 +393,7 @@ var derivedSystemEnv = map[string]bool{
// operator every run.
func refreshDerivedEnv(ctx context.Context, cl client.Client, existing, desired *v1alpha1.MinecraftServer) (bool, error) {
	want := derivedEnvWanted(desired)

	changed, err := patchOnConflictRetry(ctx, cl, existing, func() bool {
		changed := false
		for i, e := range existing.Spec.Env {
			if v, ok := want[e.Name]; ok && v != e.Value {
@@ -400,13 +401,40 @@ func refreshDerivedEnv(ctx context.Context, cl client.Client, existing, desired
				changed = true
			}
		}
	if !changed {
		return false, nil
	}
	if err := cl.Update(ctx, existing); err != nil {
		return changed
	})
	if err != nil {
		return false, fmt.Errorf("refresh %s env: %w", existing.Name, err)
	}
	return true, nil
	return changed, nil
}

// patchOnConflictRetry applies mutate to obj and sends only the difference, as a
// merge patch that carries the resourceVersion obj was read at. The operator writes
// status and felis-api patches spec.idle on these same objects, so a write can land
// between setup's read and its patch: the apiserver then answers 409, and this
// re-reads obj and runs mutate again on the fresh copy, up to retry.DefaultRetry's
// five attempts. The pinned resourceVersion is what keeps a list field such as
// spec.env safe — a merge patch replaces a list whole, and without the lock an
// entry added concurrently would be dropped. mutate reports whether it changed
// anything; nothing is sent when it did not. obj holds the stored object after.
func patchOnConflictRetry(ctx context.Context, cl client.Client, obj client.Object, mutate func() bool) (bool, error) {
	key := client.ObjectKeyFromObject(obj)
	changed, reread := false, false
	err := retry.RetryOnConflict(retry.DefaultRetry, func() error {
		if reread {
			if err := cl.Get(ctx, key, obj); err != nil {
				return err
			}
		}
		reread = true
		base := obj.DeepCopyObject().(client.Object)
		if changed = mutate(); !changed {
			return nil
		}
		return cl.Patch(ctx, obj, client.MergeFromWithOptions(base, client.MergeFromWithOptimisticLock{}))
	})
	return changed, err
}

// derivedEnvWanted maps the derived env keys of desired onto their values.
@@ -469,6 +497,8 @@ func convergeSystemServers(ctx context.Context, cl client.Client, namespace, log
		}

		var changes []string
		changed, err := patchOnConflictRetry(ctx, cl, &existing, func() bool {
			changes = nil
			if existing.Spec.Rcon == (v1alpha1.RconSpec{}) && desired.Spec.Rcon != (v1alpha1.RconSpec{}) {
				existing.Spec.Rcon = desired.Spec.Rcon
				changes = append(changes, "spec.rcon")
@@ -478,13 +508,14 @@ func convergeSystemServers(ctx context.Context, cl client.Client, namespace, log
				changes = append(changes, "spec.startup.healthHTTPPort")
			}
			changes = append(changes, convergeDerivedEnv(&existing, desired)...)

		if len(changes) == 0 {
			outcomes = append(outcomes, systemServerOutcome{name: p.name, available: true, skipped: "already converged"})
			return len(changes) > 0
		})
		if err != nil {
			outcomes = append(outcomes, systemServerOutcome{name: p.name, err: fmt.Errorf("converge %s: %w", p.name, err)})
			continue
		}
		if err := cl.Update(ctx, &existing); err != nil {
			outcomes = append(outcomes, systemServerOutcome{name: p.name, err: fmt.Errorf("converge %s: %w", p.name, err)})
		if !changed {
			outcomes = append(outcomes, systemServerOutcome{name: p.name, available: true, skipped: "already converged"})
			continue
		}
		outcomes = append(outcomes, systemServerOutcome{name: p.name, available: true, updated: true, changes: changes})
@@ -677,16 +708,22 @@ func ensureSecretReplica(ctx context.Context, cl client.Client, controlNamespace
		if out := validate(&src, controlNamespace, ""); !out.available {
			return out
		}
		changed, err := patchOnConflictRetry(ctx, cl, existing, func() bool {
			if bytes.Equal(existing.Data[secretKey], src.Data[secretKey]) {
			return validate(existing, minecraftNamespace, "already current")
				return false
			}
			if existing.Data == nil {
				existing.Data = map[string][]byte{}
			}
			existing.Data[secretKey] = src.Data[secretKey]
		if err := cl.Update(ctx, existing); err != nil {
			return true
		})
		if err != nil {
			return systemServerOutcome{name: name, err: err}
		}
		if !changed {
			return validate(existing, minecraftNamespace, "already current")
		}
		return systemServerOutcome{name: name, updated: true, available: true}
	}
	if controlNamespace == minecraftNamespace {
+19 −12

File changed.

Preview size limit exceeded, changes collapsed.

Loading