feat(distributed): 支持单主控多节点部署和停服迁移

复用现有 k3s 调度和 Job 生命周期,增加 worker 接入与批准、受保护节点身份、归档传输、持久迁移锁及活动 PVC 切换;同步管理员 API、CLI、面板和隔离规则。分布式模式默认关闭,保持单机兼容。

验证:Go 全量测试与 vet;面板 874 个测试、lint/build;Linux VM 安装器测试、清单服务端 dry-run、网络命名空间防火墙实测。A/B/C 三机 WireGuard、Velocity 和迁移验收仍待完成。
This commit is contained in:
Lemon-miaow committed 2026-10-01 19:48:37 +08:00
1 parent 26e817f2c9
commit adf563ffe3
77 files changed
+5224 -73

No files matched your search

+278
View File
@@ -0,0 +1,278 @@
// Package distributed extends the existing Job executors with archive transport and stopped migration.
package distributed
import (
"context"
"encoding/json"
"errors"
"fmt"
"net/http"
"time"
"felis.lolicon.best/internal/apis/felis/v1alpha1"
"felis.lolicon.best/internal/archivetransfer"
"felis.lolicon.best/internal/backup"
"felis.lolicon.best/internal/backupjob"
"felis.lolicon.best/internal/maintenance"
"felis.lolicon.best/internal/placement"
"felis.lolicon.best/internal/restore"
"felis.lolicon.best/internal/worldexport"
batchv1 "k8s.io/api/batch/v1"
corev1 "k8s.io/api/core/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/types"
"sigs.k8s.io/controller-runtime/pkg/client"
)
const LabelTransfer = "felis.lolicon.best/archive-job"
type Backup struct {
ID string `json:"id"`
Server string `json:"server"`
Owner string `json:"owner"`
Reason string `json:"reason"`
Protect string `json:"protect,omitempty"`
Receipt archivetransfer.Receipt `json:"receipt"`
}
type pendingBackup struct {
Ticket archivetransfer.Ticket
Owner, Reason, Protect string
}
type Manager struct {
Client client.Client
Namespace, Image, Controller string
Archive archivetransfer.Client
Resolve placement.Resolver
// Record is idempotent by transfer ID and runs on A after durable archive commit.
Record func(context.Context, Backup) error
}
func (m *Manager) CreateBackupJob(ctx context.Context, p backupjob.JobParams) error {
world, err := m.Resolve(ctx, p.Server)
if err != nil {
return err
}
job, t, err := m.uploadJob(p.Server, world.Claim, world.Node, p.JobName)
if err != nil {
return err
}
reason := "manual"
if p.Scheduled {
reason = backupjob.ReasonScheduled
}
if p.RestoreRef != "" {
reason = backupjob.ReasonPreRestore
}
meta := pendingBackup{Ticket: t, Owner: p.FormerOwner, Reason: reason, Protect: p.RestoreBackupID}
raw, _ := json.Marshal(meta)
job.Annotations = map[string]string{archivetransfer.Annotation: string(raw)}
job.Labels[maintenance.LabelManagedBy] = "felis-backup"
job.Labels[archivetransfer.LabelPending] = "true"
if p.RestoreRef != "" {
job.Labels[maintenance.LabelThenRestore] = maintenance.ThenRestorePending
job.Annotations[maintenance.AnnotationRestoreRef] = p.RestoreRef
job.Annotations[maintenance.AnnotationRestoreBackupID] = p.RestoreBackupID
}
job.Spec.TTLSecondsAfterFinished = nil
err = m.Client.Create(ctx, job)
if apierrors.IsAlreadyExists(err) {
return backupjob.ErrAlreadyExists
}
return err
}
func (m *Manager) uploadJob(server, claim, node, name string) (*batchv1.Job, archivetransfer.Ticket, error) {
t, url, token, err := m.Archive.Issue(server, http.MethodPut, "", "", 2*time.Hour)
if err != nil {
return nil, t, err
}
job, err := worldexport.ExportJob(worldexport.JobParams{Server: server, ID: t.ID[:16], Mode: worldexport.ModeWorld, WorldPVC: claim, TargetURL: url, Token: token, Namespace: m.Namespace, ServiceAccount: "felis-restore", Image: m.Image, WorldsRoot: "/world", Deadline: 2 * time.Hour})
if err != nil {
return nil, t, err
}
job.Name = name
job.Labels[LabelTransfer] = "true"
job.Spec.Template.Labels[LabelTransfer] = "true"
job.Spec.Template.Spec.Containers[0].Args = append(job.Spec.Template.Spec.Containers[0].Args, "--archive-raw")
m.pin(&job.Spec.Template.Spec, node)
return job, t, nil
}
func (m *Manager) pin(p *corev1.PodSpec, node string) {
if node == "" {
node = m.Controller
}
p.NodeSelector = map[string]string{placement.LabelIdentity: node}
}
func (m *Manager) restoreParams(ctx context.Context, p restore.JobParams) (restore.JobParams, error) {
rec, err := m.Archive.Inspect(ctx, p.BackupRef)
if err != nil {
return p, err
}
_, url, token, err := m.Archive.Issue(p.Server, http.MethodGet, p.BackupRef, rec.SHA256, 2*time.Hour)
if err != nil {
return p, err
}
p.SourceURL, p.Token, p.SHA256, p.MaxBytes = url, token, rec.SHA256, m.Archive.Limit
if p.MaxBytes <= 0 {
p.MaxBytes = archivetransfer.DefaultLimit
}
p.BackupPVC = ""
return p, nil
}
func (m *Manager) CreateRestoreJob(ctx context.Context, p restore.JobParams) error {
world, err := m.Resolve(ctx, p.Server)
if err != nil {
return err
}
p.WorldPVC = world.Claim
p, err = m.restoreParams(ctx, p)
if err != nil {
return err
}
// Preserve the existing deterministic-name conflict and finished-Job retry rules.
return restore.NewK8sJobs(&pinnedClient{Client: m.Client, node: world.Node, manager: m}).CreateRestoreJob(ctx, p)
}
type pinnedClient struct {
client.Client
node string
manager *Manager
}
func (c *pinnedClient) Create(ctx context.Context, o client.Object, opts ...client.CreateOption) error {
if j, ok := o.(*batchv1.Job); ok {
c.manager.pin(&j.Spec.Template.Spec, c.node)
j.Labels[LabelTransfer] = "true"
j.Spec.Template.Labels[LabelTransfer] = "true"
}
return c.Client.Create(ctx, o, opts...)
}
// SettleBackups is retried after A restarts. Jobs remain durable until the row is recorded.
func (m *Manager) SettleBackups(ctx context.Context) error {
var list batchv1.JobList
if err := m.Client.List(ctx, &list, client.InNamespace(m.Namespace), client.MatchingLabels{archivetransfer.LabelPending: "true"}); err != nil {
return err
}
var errs []error
for i := range list.Items {
j := &list.Items[i]
var p pendingBackup
if err := json.Unmarshal([]byte(j.Annotations[archivetransfer.Annotation]), &p); err != nil {
errs = append(errs, err)
continue
}
rec, err := m.Archive.Receipt(ctx, p.Ticket.ID)
if err != nil {
if !maintenance.JobFinished(j) {
continue
}
// A terminal upload failure holds nothing, but must never start its restore chain.
if jobSucceeded(j) {
errs = append(errs, fmt.Errorf("backup %s awaits durable receipt: %w", j.Name, err))
continue
}
} else {
if rec.Ref != p.Ticket.Ref || rec.SHA256 == "" || rec.Size <= 0 || rec.Size > p.Ticket.Limit {
errs = append(errs, fmt.Errorf("invalid receipt for %s", j.Name))
continue
}
if m.Record == nil {
errs = append(errs, errors.New("backup recorder unavailable"))
continue
}
if err := m.Record(ctx, Backup{ID: p.Ticket.ID, Server: p.Ticket.Server, Owner: p.Owner, Reason: p.Reason, Protect: p.Protect, Receipt: rec}); err != nil {
errs = append(errs, err)
continue
}
}
before := j.DeepCopy()
delete(j.Labels, archivetransfer.LabelPending)
ttl := int32(600)
j.Spec.TTLSecondsAfterFinished = &ttl
if err := m.Client.Patch(ctx, j, client.MergeFromWithOptions(before, client.MergeFromWithOptimisticLock{})); err != nil {
errs = append(errs, err)
}
}
return errors.Join(errs...)
}
func jobSucceeded(j *batchv1.Job) bool {
for _, c := range j.Status.Conditions {
if c.Type == batchv1.JobComplete && c.Status == corev1.ConditionTrue {
return true
}
}
return false
}
// RemoteArchiver lets the existing reaper make every decision on A and snapshot only one remote PVC.
// Local verification, retention and offsite continue to use the same tarLocal paths.
type RemoteArchiver struct {
*backup.TarLocal
Manager *Manager
}
// A failed migration may wait for operator intervention longer than normal
// backup retention. Keep its safety archive until the persistent lock releases.
func (a *RemoteArchiver) Delete(ctx context.Context, ref backup.ArchiveRef) error {
var servers v1alpha1.MinecraftServerList
if err := a.Manager.Client.List(ctx, &servers, client.InNamespace(a.Manager.Namespace)); err != nil {
return err
}
for i := range servers.Items {
op, err := readOperation(&servers.Items[i])
if err != nil && !errors.Is(err, ErrNotFound) {
return err
}
if err == nil && op.State != "succeeded" && op.Backup.Ref == string(ref) {
return fmt.Errorf("%w: migration retains its safety archive", ErrBusy)
}
}
return a.TarLocal.Delete(ctx, ref)
}
func (a *RemoteArchiver) Archive(ctx context.Context, server, pvc string) (backup.Archived, error) {
m := a.Manager
world, err := m.Resolve(ctx, server)
if err != nil {
return backup.Archived{}, err
}
if world.Claim != pvc {
return backup.Archived{}, errors.New("active PVC changed")
}
id := archivetransfer.ID()
name := "reap-" + id[:16]
j, t, err := m.uploadJob(server, pvc, world.Node, name)
if err != nil {
return backup.Archived{}, err
}
if err = m.Client.Create(ctx, j); err != nil {
return backup.Archived{}, err
}
ticker := time.NewTicker(2 * time.Second)
defer ticker.Stop()
for {
rec, err := m.Archive.Receipt(ctx, t.ID)
if err == nil {
return backup.Archived{Ref: backup.ArchiveRef(rec.Ref), Size: rec.Size, SHA256: rec.SHA256}, nil
}
var job batchv1.Job
if err = m.Client.Get(ctx, types.NamespacedName{Namespace: m.Namespace, Name: name}, &job); err != nil {
return backup.Archived{}, err
}
if maintenance.JobFinished(&job) && !jobSucceeded(&job) {
return backup.Archived{}, errors.New("remote archive Job failed")
}
select {
case <-ctx.Done():
return backup.Archived{}, ctx.Err()
case <-ticker.C:
}
}
}
+521
View File
@@ -0,0 +1,521 @@
package distributed
import (
"context"
"encoding/json"
"errors"
"fmt"
"sort"
"time"
"felis.lolicon.best/internal/apis/felis/v1alpha1"
"felis.lolicon.best/internal/archivetransfer"
"felis.lolicon.best/internal/maintenance"
"felis.lolicon.best/internal/placement"
"felis.lolicon.best/internal/restore"
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"
"k8s.io/apimachinery/pkg/api/resource"
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"
)
const MigrationAnnotation = "felis.lolicon.best/migration"
// Admission errors keep HTTP policy in the API layer.
var ErrBusy = errors.New("server must be fully stopped with no maintenance operation")
var ErrNotFound = errors.New("migration not found")
type Node struct {
Name string `json:"name"`
Role string `json:"role"`
Ready bool `json:"ready"`
Approved bool `json:"approved"`
Addresses []string `json:"addresses"`
Architecture string `json:"architecture"`
}
func (m *Manager) Nodes(ctx context.Context) ([]Node, error) {
var list corev1.NodeList
if err := m.Client.List(ctx, &list); err != nil {
return nil, err
}
out := make([]Node, 0, len(list.Items))
for _, n := range list.Items {
info := Node{Name: n.Name, Role: n.Labels[placement.LabelRole], Ready: placement.Online(&n), Approved: n.Labels[placement.LabelApproved] == "true" && !n.Spec.Unschedulable, Addresses: []string{}, Architecture: n.Status.NodeInfo.Architecture}
for _, a := range n.Status.Addresses {
if a.Type == corev1.NodeInternalIP || a.Type == corev1.NodeExternalIP {
info.Addresses = append(info.Addresses, a.Address)
}
}
out = append(out, info)
}
sort.Slice(out, func(i, j int) bool { return out[i].Name < out[j].Name })
return out, nil
}
func (m *Manager) ValidateNode(ctx context.Context, name string) error {
return placement.Worker(ctx, m.Client, name)
}
// Operation is persisted on the CR together with its non-expiring maintenance lock.
// SourcePVC is retained even after success, and retry never changes the committed active world.
type Operation struct {
ID string `json:"id"`
Server string `json:"server"`
State string `json:"state"`
Stage string `json:"stage"`
SourceNode string `json:"sourceNode"`
TargetNode string `json:"targetNode"`
SourcePVC string `json:"sourcePVC"`
TargetPVC string `json:"targetPVC"`
Backup archivetransfer.Receipt `json:"backup"`
Owner string `json:"-"`
Started time.Time `json:"started"`
Updated time.Time `json:"updated"`
Error string `json:"error,omitempty"`
Switched bool `json:"switched"`
Attempt int `json:"attempt"`
}
// owner travels in persistence, but never on the public operation view.
type persistedOperation struct {
Operation
OwnerID string `json:"ownerId"`
}
func readOperation(s *v1alpha1.MinecraftServer) (Operation, error) {
var stored persistedOperation
if s.Annotations[MigrationAnnotation] == "" {
return Operation{}, ErrNotFound
}
if err := json.Unmarshal([]byte(s.Annotations[MigrationAnnotation]), &stored); err != nil {
return Operation{}, err
}
stored.Operation.Owner = stored.OwnerID
return stored.Operation, nil
}
func (m *Manager) Migration(ctx context.Context, server, id string) (Operation, error) {
var s v1alpha1.MinecraftServer
if err := m.Client.Get(ctx, types.NamespacedName{Namespace: m.Namespace, Name: server}, &s); err != nil {
return Operation{}, err
}
op, err := readOperation(&s)
if err == nil && id != "" && op.ID != id {
return Operation{}, ErrNotFound
}
return op, err
}
func (m *Manager) quiet(ctx context.Context, s *v1alpha1.MinecraftServer, own ...string) error {
if s.Spec.DesiredState != v1alpha1.DesiredStopped || s.Status.Phase != v1alpha1.PhaseStopped || s.Status.Ready {
return ErrBusy
}
var pods corev1.PodList
if err := m.Client.List(ctx, &pods, client.InNamespace(m.Namespace), client.MatchingLabels{v1alpha1.LabelServer: s.Name}); err != nil {
return err
}
// Even terminal maintenance Pods must have exited before a new attempt writes its PVC.
for _, p := range pods.Items {
allowed := len(own) > 0 && p.Labels[MigrationAnnotation] == own[0]
if p.Labels[v1alpha1.LabelComponent] == "server" || (!allowed && p.Status.Phase != corev1.PodSucceeded && p.Status.Phase != corev1.PodFailed) {
return ErrBusy
}
}
return nil
}
func (m *Manager) BeginMigration(ctx context.Context, server, target, owner string) (Operation, error) {
if err := m.ValidateNode(ctx, target); err != nil {
return Operation{}, err
}
var op Operation
err := retry.RetryOnConflict(retry.DefaultRetry, func() error {
var s v1alpha1.MinecraftServer
if err := m.Client.Get(ctx, types.NamespacedName{Namespace: m.Namespace, Name: server}, &s); err != nil {
return err
}
if previous, err := readOperation(&s); err == nil && previous.State != "succeeded" {
if previous.TargetNode == target {
op = previous
return nil
}
return ErrBusy
} else if err != nil && !errors.Is(err, ErrNotFound) {
return err
}
if s.Spec.ReaperExempt || s.Labels[v1alpha1.LabelSystemRole] != "" {
return ErrBusy
}
if err := m.quiet(ctx, &s); err != nil {
return err
}
var jobs batchv1.JobList
if err := m.Client.List(ctx, &jobs, client.InNamespace(m.Namespace), client.MatchingLabels{maintenance.LabelServer: server}); err != nil {
return err
}
if _, held := maintenance.Holder(server, s.Annotations, jobs.Items, time.Now()); held {
return ErrBusy
}
world, err := m.Resolve(ctx, server)
if err != nil {
return err
}
source := world.Node
if source == "" {
source = m.Controller
}
if source == target {
return errors.New("source and target nodes are identical")
}
// Bound local-path worlds have a physical node; reject a guessed or mismatched source.
var pvc corev1.PersistentVolumeClaim
if err := m.Client.Get(ctx, types.NamespacedName{Namespace: m.Namespace, Name: world.Claim}, &pvc); err != nil {
return err
}
if pvc.Spec.VolumeName == "" {
return errors.New("source world is not bound")
}
if err := m.volumeOnNode(ctx, pvc.Spec.VolumeName, source); err != nil {
return err
}
id := archivetransfer.ID()
now := time.Now().UTC()
op = Operation{ID: id, Server: server, State: "backing_up", Stage: "backing_up", SourceNode: source, TargetNode: target, SourcePVC: world.Claim, TargetPVC: "world-" + server + "-m" + id[:12], Owner: owner, Started: now, Updated: now}
return m.save(ctx, &s, op, true)
})
return op, err
}
func (m *Manager) volumeOnNode(ctx context.Context, volume, node string) error {
var pv corev1.PersistentVolume
if err := m.Client.Get(ctx, types.NamespacedName{Name: volume}, &pv); err != nil {
return err
}
if pv.Spec.NodeAffinity == nil || pv.Spec.NodeAffinity.Required == nil {
return errors.New("migration requires a node-local volume with node affinity")
}
var n corev1.Node
if err := m.Client.Get(ctx, types.NamespacedName{Name: node}, &n); err != nil {
return err
}
for _, term := range pv.Spec.NodeAffinity.Required.NodeSelectorTerms {
if len(term.MatchFields) > 0 {
continue
}
matched := len(term.MatchExpressions) > 0
for _, e := range term.MatchExpressions {
if e.Operator != corev1.NodeSelectorOpIn {
matched = false
break
}
found := false
for _, v := range e.Values {
if n.Labels[e.Key] == v {
found = true
}
}
if !found {
matched = false
break
}
}
if matched {
return nil
}
}
return errors.New("world volume is not on the recorded execution node")
}
func (m *Manager) save(ctx context.Context, s *v1alpha1.MinecraftServer, op Operation, lock bool) error {
before := s.DeepCopy()
op.Updated = time.Now().UTC()
raw, _ := json.Marshal(persistedOperation{Operation: op, OwnerID: op.Owner})
if s.Annotations == nil {
s.Annotations = map[string]string{}
}
s.Annotations[MigrationAnnotation] = string(raw)
if lock {
s.Annotations[maintenance.Annotation] = maintenance.LockValue(maintenance.KindMigration, op.Started)
} else {
delete(s.Annotations, maintenance.Annotation)
}
return m.Client.Patch(ctx, s, client.MergeFromWithOptions(before, client.MergeFromWithOptimisticLock{}))
}
func (m *Manager) RetryMigration(ctx context.Context, server, id string) (Operation, error) {
var op Operation
err := retry.RetryOnConflict(retry.DefaultRetry, func() error {
var s v1alpha1.MinecraftServer
if err := m.Client.Get(ctx, types.NamespacedName{Namespace: m.Namespace, Name: server}, &s); err != nil {
return err
}
var err error
op, err = readOperation(&s)
if err != nil {
return err
}
if op.ID != id {
return ErrNotFound
}
if op.State != "failed" {
return ErrBusy
}
if err := m.quiet(ctx, &s); err != nil {
return err
}
if err := m.ValidateNode(ctx, op.TargetNode); err != nil {
return err
}
op.State = op.Stage
op.Error = ""
op.Attempt++
return m.save(ctx, &s, op, true)
})
return op, err
}
func (m *Manager) ReconcileMigrations(ctx context.Context) error {
var servers v1alpha1.MinecraftServerList
if err := m.Client.List(ctx, &servers, client.InNamespace(m.Namespace)); err != nil {
return err
}
var errs []error
for i := range servers.Items {
s := &servers.Items[i]
op, err := readOperation(s)
if errors.Is(err, ErrNotFound) {
continue
}
if err != nil {
errs = append(errs, err)
continue
}
if op.State == "succeeded" {
if err := m.expireMigrationJobs(ctx, op.ID); err != nil {
errs = append(errs, err)
}
continue
}
if op.State == "failed" {
continue
}
err = m.advance(ctx, s, &op)
if err != nil {
if apierrors.IsConflict(err) {
continue
}
op.Stage = op.State
op.State = "failed"
op.Error = err.Error()
if saveErr := m.save(ctx, s, op, true); saveErr != nil {
errs = append(errs, saveErr)
}
}
}
return errors.Join(errs...)
}
func (m *Manager) advance(ctx context.Context, s *v1alpha1.MinecraftServer, op *Operation) error {
if err := m.quiet(ctx, s, op.ID); err != nil {
return err
}
if s.Annotations[maintenance.Annotation] != maintenance.LockValue(maintenance.KindMigration, op.Started) {
return errors.New("persistent migration lock was changed")
}
if err := m.ValidateNode(ctx, op.TargetNode); err != nil {
return err
}
if !op.Switched && (s.WorldPVC() != op.SourcePVC || (s.Spec.NodeName != "" && s.Spec.NodeName != op.SourceNode)) {
return errors.New("source placement changed during migration")
}
var source corev1.Node
if !op.Switched {
if err := m.Client.Get(ctx, types.NamespacedName{Name: op.SourceNode}, &source); err != nil {
return err
}
if !placement.Online(&source) {
return errors.New("source node is offline")
}
}
name := fmt.Sprintf("migration-%s-%d", op.ID[:16], op.Attempt)
switch op.State {
case "backing_up":
var j batchv1.Job
err := m.Client.Get(ctx, types.NamespacedName{Namespace: m.Namespace, Name: name + "-backup"}, &j)
if apierrors.IsNotFound(err) {
job, t, err := m.uploadJob(op.Server, op.SourcePVC, op.SourceNode, name+"-backup")
if err != nil {
return err
}
raw, _ := json.Marshal(t)
job.Annotations = map[string]string{archivetransfer.Annotation: string(raw)}
job.Labels[MigrationAnnotation] = op.ID
job.Spec.Template.Labels[MigrationAnnotation] = op.ID
job.Spec.TTLSecondsAfterFinished = nil
return m.Client.Create(ctx, job)
}
if err != nil {
return err
}
var t archivetransfer.Ticket
if err = json.Unmarshal([]byte(j.Annotations[archivetransfer.Annotation]), &t); err != nil {
return err
}
rec, err := m.Archive.Receipt(ctx, t.ID)
if err != nil {
if maintenance.JobFinished(&j) {
return fmt.Errorf("migration backup has no durable receipt: %w", err)
}
return nil
}
if rec.Ref != t.Ref || rec.SHA256 == "" {
return errors.New("migration backup receipt mismatch")
}
if m.Record == nil {
return errors.New("backup recorder unavailable")
}
if err = m.Record(ctx, Backup{ID: t.ID, Server: op.Server, Owner: op.Owner, Reason: "pre_restore", Receipt: rec}); err != nil {
return err
}
op.Backup = rec
op.State = "restoring"
op.Stage = op.State
return m.save(ctx, s, *op, true)
case "restoring":
if err := m.ensureTarget(ctx, s, op); err != nil {
return err
}
var j batchv1.Job
err := m.Client.Get(ctx, types.NamespacedName{Namespace: m.Namespace, Name: name + "-restore"}, &j)
if apierrors.IsNotFound(err) {
p := restore.JobParams{Server: op.Server, WorldPVC: op.TargetPVC, BackupRef: op.Backup.Ref, ArchiveStore: "tarLocal", Namespace: m.Namespace, ServiceAccount: "felis-restore", Image: m.Image, WorldsRoot: "/world", BackupRoot: m.Archive.Root, Deadline: 2 * time.Hour}
p, err = m.restoreParams(ctx, p)
if err != nil {
return err
}
if p.SHA256 != op.Backup.SHA256 {
return errors.New("migration archive digest changed")
}
job, err := restore.RestoreJob(p)
if err != nil {
return err
}
job.Name = name + "-restore"
m.pin(&job.Spec.Template.Spec, op.TargetNode)
job.Labels[LabelTransfer] = "true"
job.Spec.Template.Labels[LabelTransfer] = "true"
job.Labels[MigrationAnnotation] = op.ID
job.Spec.Template.Labels[MigrationAnnotation] = op.ID
job.Spec.TTLSecondsAfterFinished = nil
return m.Client.Create(ctx, job)
}
if err != nil {
return err
}
if !maintenance.JobFinished(&j) {
return nil
}
if !jobSucceeded(&j) {
return errors.New("target restore or read-back verification failed")
}
// Job Complete alone is not enough for a RWO handoff: all its processes must be gone.
if err := m.quiet(ctx, s); err != nil {
return nil
}
var pvc corev1.PersistentVolumeClaim
if err = m.Client.Get(ctx, types.NamespacedName{Namespace: m.Namespace, Name: op.TargetPVC}, &pvc); err != nil {
return err
}
if err = m.volumeOnNode(ctx, pvc.Spec.VolumeName, op.TargetNode); err != nil {
return err
}
op.State = "switching"
op.Stage = op.State
return m.save(ctx, s, *op, true)
case "switching":
var sts appsv1.StatefulSet
err := m.Client.Get(ctx, types.NamespacedName{Namespace: m.Namespace, Name: s.Name}, &sts)
if err == nil {
if sts.Spec.Replicas != nil && *sts.Spec.Replicas != 0 {
return ErrBusy
}
if sts.Status.Replicas != 0 {
return nil
}
policy := metav1.DeletePropagationForeground
return m.Client.Delete(ctx, &sts, &client.DeleteOptions{PropagationPolicy: &policy, Preconditions: &metav1.Preconditions{UID: &sts.UID, ResourceVersion: &sts.ResourceVersion}})
}
if !apierrors.IsNotFound(err) {
return err
}
// One optimistic CR write commits node, claim and progress. No failure can roll this back.
target := s.DeepCopy()
committed := *op
target.Spec.NodeName = op.TargetNode
target.Spec.Storage.ClaimName = op.TargetPVC
committed.Switched = true
committed.State = "succeeded"
committed.Stage = "succeeded"
committed.Updated = time.Now().UTC()
raw, _ := json.Marshal(persistedOperation{Operation: committed, OwnerID: op.Owner})
target.Annotations[MigrationAnnotation] = string(raw)
delete(target.Annotations, maintenance.Annotation)
if err := m.Client.Patch(ctx, target, client.MergeFromWithOptions(s.DeepCopy(), client.MergeFromWithOptimisticLock{})); err != nil {
return err
}
*s, *op = *target, committed
return nil
default:
return fmt.Errorf("unknown migration stage %q", op.State)
}
}
func (m *Manager) ensureTarget(ctx context.Context, s *v1alpha1.MinecraftServer, op *Operation) error {
var pvc corev1.PersistentVolumeClaim
err := m.Client.Get(ctx, types.NamespacedName{Namespace: m.Namespace, Name: op.TargetPVC}, &pvc)
if err == nil {
if pvc.Labels[MigrationAnnotation] != op.ID {
return errors.New("target PVC identity mismatch")
}
return nil
}
if !apierrors.IsNotFound(err) {
return err
}
var source corev1.PersistentVolumeClaim
if err = m.Client.Get(ctx, types.NamespacedName{Namespace: m.Namespace, Name: op.SourcePVC}, &source); err != nil {
return err
}
size := source.Spec.Resources.Requests[corev1.ResourceStorage]
if size.IsZero() {
size = resource.MustParse("8Gi")
}
pvc = corev1.PersistentVolumeClaim{ObjectMeta: metav1.ObjectMeta{Name: op.TargetPVC, Namespace: m.Namespace, Labels: map[string]string{maintenance.LabelServer: s.Name, MigrationAnnotation: op.ID}}, Spec: corev1.PersistentVolumeClaimSpec{AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteOnce}, StorageClassName: source.Spec.StorageClassName, Resources: corev1.VolumeResourceRequirements{Requests: corev1.ResourceList{corev1.ResourceStorage: size}}}}
return m.Client.Create(ctx, &pvc)
}
func (m *Manager) expireMigrationJobs(ctx context.Context, id string) error {
var jobs batchv1.JobList
if err := m.Client.List(ctx, &jobs, client.InNamespace(m.Namespace), client.MatchingLabels{MigrationAnnotation: id}); err != nil {
return err
}
for i := range jobs.Items {
j := &jobs.Items[i]
if j.Spec.TTLSecondsAfterFinished != nil || !maintenance.JobFinished(j) {
continue
}
before := j.DeepCopy()
ttl := int32(600)
j.Spec.TTLSecondsAfterFinished = &ttl
if err := m.Client.Patch(ctx, j, client.MergeFrom(before)); err != nil {
return err
}
}
return nil
}
+352
View File
@@ -0,0 +1,352 @@
package distributed
import (
"context"
"encoding/json"
"errors"
"net/http"
"net/http/httptest"
"strings"
"testing"
"time"
"felis.lolicon.best/internal/apis/felis/v1alpha1"
"felis.lolicon.best/internal/archivetransfer"
"felis.lolicon.best/internal/backup"
"felis.lolicon.best/internal/backupjob"
"felis.lolicon.best/internal/maintenance"
"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"
"k8s.io/apimachinery/pkg/api/resource"
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 fixture(t *testing.T) (*Manager, context.Context) {
t.Helper()
scheme := runtime.NewScheme()
for _, add := range []func(*runtime.Scheme) error{corev1.AddToScheme, batchv1.AddToScheme, appsv1.AddToScheme, v1alpha1.AddToScheme} {
if err := add(scheme); err != nil {
t.Fatal(err)
}
}
node := func(name string) *corev1.Node {
return &corev1.Node{ObjectMeta: metav1.ObjectMeta{Name: name, Labels: map[string]string{placement.LabelIdentity: name, placement.LabelRole: placement.RoleWorker, placement.LabelApproved: "true", "kubernetes.io/hostname": name}}, Status: corev1.NodeStatus{Conditions: []corev1.NodeCondition{{Type: corev1.NodeReady, Status: corev1.ConditionTrue}}}}
}
s := &v1alpha1.MinecraftServer{ObjectMeta: metav1.ObjectMeta{Name: "alice", Namespace: "minecraft"}, Spec: v1alpha1.MinecraftServerSpec{DesiredState: v1alpha1.DesiredStopped, NodeName: "b"}, Status: v1alpha1.MinecraftServerStatus{Phase: v1alpha1.PhaseStopped}}
s.Spec.Storage.Size = "1Gi"
pvc := &corev1.PersistentVolumeClaim{ObjectMeta: metav1.ObjectMeta{Name: s.WorldPVC(), Namespace: "minecraft"}, Spec: corev1.PersistentVolumeClaimSpec{VolumeName: "source", Resources: corev1.VolumeResourceRequirements{Requests: corev1.ResourceList{corev1.ResourceStorage: resource.MustParse("1Gi")}}}}
zero := int32(0)
cl := fake.NewClientBuilder().WithScheme(scheme).WithStatusSubresource(s, &batchv1.Job{}).WithObjects(node("b"), node("c"), s, pvc, localPV("source", "b"), &appsv1.StatefulSet{ObjectMeta: metav1.ObjectMeta{Name: "alice", Namespace: "minecraft"}, Spec: appsv1.StatefulSetSpec{Replicas: &zero}}, &corev1.Service{ObjectMeta: metav1.ObjectMeta{Name: "alice", Namespace: "minecraft"}, Spec: corev1.ServiceSpec{ClusterIP: "10.43.0.80"}}).Build()
var receipt archivetransfer.Receipt
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/inspect" && r.Method == "POST" || strings.HasPrefix(r.URL.Path, "/receipts/") {
var jobs batchv1.JobList
cl.List(r.Context(), &jobs)
for _, j := range jobs.Items {
var ticket archivetransfer.Ticket
if json.Unmarshal([]byte(j.Annotations[archivetransfer.Annotation]), &ticket) == nil && ticket.Ref != "" {
receipt = archivetransfer.Receipt{Ref: ticket.Ref, SHA256: strings.Repeat("a", 64), Size: 50}
}
}
if receipt.Ref == "" {
http.NotFound(w, r)
return
}
json.NewEncoder(w).Encode(receipt)
return
}
http.NotFound(w, r)
}))
t.Cleanup(server.Close)
m := &Manager{Client: cl, Namespace: "minecraft", Controller: "a", Image: "registry.local/felis:1", Archive: archivetransfer.Client{Root: "/backups", URL: server.URL, Key: strings.Repeat("k", 32)}, Record: func(context.Context, Backup) error { return nil }}
m.Resolve = placement.Resolve(cl, "minecraft")
return m, context.Background()
}
func localPV(name, node string) *corev1.PersistentVolume {
return &corev1.PersistentVolume{ObjectMeta: metav1.ObjectMeta{Name: name}, Spec: corev1.PersistentVolumeSpec{NodeAffinity: &corev1.VolumeNodeAffinity{Required: &corev1.NodeSelector{NodeSelectorTerms: []corev1.NodeSelectorTerm{{MatchExpressions: []corev1.NodeSelectorRequirement{{Key: "kubernetes.io/hostname", Operator: corev1.NodeSelectorOpIn, Values: []string{node}}}}}}}}}
}
func getServer(t *testing.T, m *Manager, ctx context.Context) *v1alpha1.MinecraftServer {
t.Helper()
var s v1alpha1.MinecraftServer
if err := m.Client.Get(ctx, types.NamespacedName{Namespace: m.Namespace, Name: "alice"}, &s); err != nil {
t.Fatal(err)
}
return &s
}
func TestMigrationSurvivesRestartAndKeepsIdentity(t *testing.T) {
m, ctx := fixture(t)
op, err := m.BeginMigration(ctx, "alice", "c", "owner")
if err != nil {
t.Fatal(err)
}
if again, err := m.BeginMigration(ctx, "alice", "c", "owner"); err != nil || again.ID != op.ID {
t.Fatalf("duplicate %+v: %v", again, err)
}
s := getServer(t, m, ctx)
if kind, held := maintenance.Holder(s.Name, s.Annotations, nil, time.Now().Add(365*24*time.Hour)); !held || kind != maintenance.KindMigration {
t.Fatal("migration lock expired")
}
if _, err := m.BeginMigration(ctx, "alice", "b", "owner"); !errors.Is(err, ErrBusy) {
t.Fatalf("competing migration: %v", err)
}
// A restart reconstructs the coordinator solely from persisted CR and Job state.
restarted := *m
m = &restarted
for i := 0; i < 3; i++ {
if err := m.ReconcileMigrations(ctx); err != nil {
t.Fatal(err)
}
}
current, err := m.Migration(ctx, "alice", op.ID)
if err != nil || current.State != "restoring" {
t.Fatalf("progress %+v: %v", current, err)
}
var pvc corev1.PersistentVolumeClaim
if err := m.Client.Get(ctx, types.NamespacedName{Namespace: m.Namespace, Name: op.TargetPVC}, &pvc); err != nil {
t.Fatal(err)
}
pvc.Spec.VolumeName = "target"
if err := m.Client.Update(ctx, &pvc); err != nil {
t.Fatal(err)
}
if err := m.Client.Create(ctx, localPV("target", "c")); err != nil {
t.Fatal(err)
}
var jobs batchv1.JobList
if err := m.Client.List(ctx, &jobs); err != nil {
t.Fatal(err)
}
for i := range jobs.Items {
j := &jobs.Items[i]
assertJobIsolation(t, j)
j.Status.Conditions = []batchv1.JobCondition{{Type: batchv1.JobComplete, Status: corev1.ConditionTrue}}
if err := m.Client.Status().Update(ctx, j); err != nil {
t.Fatal(err)
}
}
for i := 0; i < 3; i++ {
if err := m.ReconcileMigrations(ctx); err != nil {
t.Fatal(err)
}
}
s = getServer(t, m, ctx)
if s.Spec.NodeName != "c" || s.WorldPVC() != op.TargetPVC || s.Spec.DesiredState != v1alpha1.DesiredStopped || s.Status.Ready {
t.Fatalf("unsafe switch %+v", s)
}
if _, held := maintenance.Holder(s.Name, s.Annotations, nil, time.Now()); held {
t.Fatal("success did not release lock")
}
var svc corev1.Service
if err := m.Client.Get(ctx, types.NamespacedName{Namespace: m.Namespace, Name: "alice"}, &svc); err != nil || svc.Spec.ClusterIP != "10.43.0.80" {
t.Fatalf("service changed: %v", err)
}
if err := m.Client.Get(ctx, types.NamespacedName{Namespace: m.Namespace, Name: op.SourcePVC}, &pvc); err != nil {
t.Fatal("source PVC removed", err)
}
var sts appsv1.StatefulSet
if err := m.Client.Get(ctx, types.NamespacedName{Namespace: m.Namespace, Name: "alice"}, &sts); !apierrors.IsNotFound(err) {
t.Fatal("stopped immutable StatefulSet not removed", err)
}
}
func assertJobIsolation(t *testing.T, j *batchv1.Job) {
t.Helper()
p := j.Spec.Template.Spec
if p.AutomountServiceAccountToken == nil || *p.AutomountServiceAccountToken || p.HostNetwork || p.HostPID || p.HostIPC {
t.Fatal("maintenance has host credentials/namespaces")
}
claims := 0
for _, v := range p.Volumes {
if v.Secret != nil || v.HostPath != nil {
t.Fatal("maintenance carries secret or host mount")
}
if v.PersistentVolumeClaim != nil {
claims++
}
}
if claims != 1 {
t.Fatalf("cross-node job mounts %d PVCs", claims)
}
for _, c := range p.Containers {
if c.SecurityContext.AllowPrivilegeEscalation == nil || *c.SecurityContext.AllowPrivilegeEscalation {
t.Fatal("maintenance can escalate")
}
for _, e := range c.Env {
if strings.Contains(e.Name, "DATABASE") || e.Name == archivetransfer.KeyEnv {
t.Fatal("full credentials leaked")
}
}
}
}
func TestLostSourceFailsClosedAndRetryKeepsSource(t *testing.T) {
m, ctx := fixture(t)
op, err := m.BeginMigration(ctx, "alice", "c", "owner")
if err != nil {
t.Fatal(err)
}
var n corev1.Node
m.Client.Get(ctx, types.NamespacedName{Name: "b"}, &n)
n.Status.Conditions = nil
if err := m.Client.Status().Update(ctx, &n); err != nil {
t.Fatal(err)
}
if err := m.ReconcileMigrations(ctx); err != nil {
t.Fatal(err)
}
failed, err := m.Migration(ctx, "alice", op.ID)
if err != nil || failed.State != "failed" {
t.Fatalf("offline source %+v: %v", failed, err)
}
s := getServer(t, m, ctx)
if s.WorldPVC() != op.SourcePVC || s.Spec.NodeName != "b" {
t.Fatal("failed migration switched placement")
}
if _, held := maintenance.Holder(s.Name, s.Annotations, nil, time.Now().Add(time.Hour)); !held {
t.Fatal("failed migration lost lock")
}
if _, err := m.RetryMigration(ctx, "alice", op.ID); err != nil {
t.Fatal(err)
}
if err := m.ReconcileMigrations(ctx); err != nil {
t.Fatal(err)
}
if current, _ := m.Migration(ctx, "alice", op.ID); current.State != "failed" {
t.Fatal("retry started tasks on offline source")
}
}
func TestMigrationRejectsLiveProcessAndUnapprovedTarget(t *testing.T) {
m, ctx := fixture(t)
if err := m.Client.Create(ctx, &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Name: "alice-0", Namespace: m.Namespace, Labels: map[string]string{v1alpha1.LabelServer: "alice", v1alpha1.LabelComponent: "server"}}, Status: corev1.PodStatus{Phase: corev1.PodRunning}}); err != nil {
t.Fatal(err)
}
if _, err := m.BeginMigration(ctx, "alice", "c", "owner"); !errors.Is(err, ErrBusy) {
t.Fatal("live source accepted", err)
}
var n corev1.Node
m.Client.Get(ctx, types.NamespacedName{Name: "c"}, &n)
delete(n.Labels, placement.LabelApproved)
m.Client.Update(ctx, &n)
if err := m.ValidateNode(ctx, "c"); err == nil {
t.Fatal("unapproved worker accepted")
}
}
func TestMigrationFailureRetainsSourceAndCanRetry(t *testing.T) {
for _, stage := range []string{"restoring", "switching"} {
t.Run(stage, func(t *testing.T) {
m, ctx := fixture(t)
op, err := m.BeginMigration(ctx, "alice", "c", "owner")
if err != nil {
t.Fatal(err)
}
for i := 0; i < 3; i++ {
if err := m.ReconcileMigrations(ctx); err != nil {
t.Fatal(err)
}
}
var jobs batchv1.JobList
if err := m.Client.List(ctx, &jobs); err != nil {
t.Fatal(err)
}
for i := range jobs.Items {
j := &jobs.Items[i]
condition := batchv1.JobComplete
if stage == "restoring" && strings.HasSuffix(j.Name, "-restore") {
condition = batchv1.JobFailed
}
j.Status.Conditions = []batchv1.JobCondition{{Type: condition, Status: corev1.ConditionTrue}}
if err := m.Client.Status().Update(ctx, j); err != nil {
t.Fatal(err)
}
}
if stage == "switching" {
var pvc corev1.PersistentVolumeClaim
if err := m.Client.Get(ctx, types.NamespacedName{Namespace: m.Namespace, Name: op.TargetPVC}, &pvc); err != nil {
t.Fatal(err)
}
pvc.Spec.VolumeName = "target"
if err := m.Client.Update(ctx, &pvc); err != nil {
t.Fatal(err)
}
if err := m.Client.Create(ctx, localPV("target", "c")); err != nil {
t.Fatal(err)
}
fail := true
m.Client = interceptor.NewClient(m.Client.(client.WithWatch), interceptor.Funcs{
Patch: func(ctx context.Context, c client.WithWatch, obj client.Object, patch client.Patch, opts ...client.PatchOption) error {
if s, ok := obj.(*v1alpha1.MinecraftServer); ok && s.Spec.NodeName == "c" && fail {
fail = false
return errors.New("commit interrupted")
}
return c.Patch(ctx, obj, patch, opts...)
},
})
}
for i := 0; i < 3; i++ {
if err := m.ReconcileMigrations(ctx); err != nil {
t.Fatal(err)
}
}
failed, err := m.Migration(ctx, "alice", op.ID)
if err != nil || failed.State != "failed" || failed.Stage != stage || failed.Switched || failed.Error == "" {
t.Fatalf("failure %+v: %v", failed, err)
}
s := getServer(t, m, ctx)
if s.WorldPVC() != op.SourcePVC || s.Spec.NodeName != "b" || s.Spec.DesiredState != v1alpha1.DesiredStopped {
t.Fatal("failed operation changed active source")
}
if _, held := maintenance.Holder(s.Name, s.Annotations, nil, time.Now().Add(30*24*time.Hour)); !held {
t.Fatal("failed operation released lock")
}
archiver := &RemoteArchiver{TarLocal: &backup.TarLocal{BackupRoot: t.TempDir()}, Manager: m}
if err := archiver.Delete(ctx, backup.ArchiveRef(failed.Backup.Ref)); !errors.Is(err, ErrBusy) {
t.Fatalf("failed migration lost its safety archive: %v", err)
}
if retry, err := m.RetryMigration(ctx, "alice", op.ID); err != nil || retry.State != stage || retry.Attempt != 1 {
t.Fatalf("retry %+v: %v", retry, err)
}
if err := m.ReconcileMigrations(ctx); err != nil {
t.Fatal(err)
}
if stage == "switching" {
if done, _ := m.Migration(ctx, "alice", op.ID); done.State != "succeeded" || !done.Switched {
t.Fatalf("commit retry %+v", done)
}
}
})
}
}
func TestBackupResolvesActiveClaimAndPreservesConflictContract(t *testing.T) {
m, ctx := fixture(t)
p := backupjob.JobParams{Server: "alice", WorldPVC: "stale-claim", JobName: "manual-backup"}
if err := m.CreateBackupJob(ctx, p); err != nil {
t.Fatal(err)
}
var j batchv1.Job
if err := m.Client.Get(ctx, types.NamespacedName{Namespace: m.Namespace, Name: p.JobName}, &j); err != nil {
t.Fatal(err)
}
assertJobIsolation(t, &j)
for _, v := range j.Spec.Template.Spec.Volumes {
if v.PersistentVolumeClaim != nil && v.PersistentVolumeClaim.ClaimName != getServer(t, m, ctx).WorldPVC() {
t.Fatal("backup used a stale world claim")
}
}
if err := m.CreateBackupJob(ctx, p); !errors.Is(err, backupjob.ErrAlreadyExists) {
t.Fatalf("duplicate backup: %v", err)
}
}