mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
Delete worker pods in terminal states (#2051)
Fixes #1935 - [ x ] Tests pass - [ x ] Appropriate changes to documentation are included in the PR
This commit is contained in:
@@ -16,10 +16,11 @@ package controllers
|
||||
|
||||
// RBAC needed by atecontroller components outside this package, which
|
||||
// controller-gen (paths="./...") does not scan:
|
||||
// - internal/workersync's pod informer lists and watches worker pods.
|
||||
// - internal/workersync's pod informer lists and watches worker pods, and
|
||||
// the syncer deletes worker pods that reach a terminal phase.
|
||||
// - internal/k8sresolver watches ateapi's EndpointSlices to dial it.
|
||||
//
|
||||
//+kubebuilder:rbac:groups=core,resources=pods,verbs=get;list;watch
|
||||
//+kubebuilder:rbac:groups=core,resources=pods,verbs=get;list;watch;delete
|
||||
//+kubebuilder:rbac:groups=discovery.k8s.io,resources=endpointslices,verbs=get;list;watch,namespace=ate-system
|
||||
|
||||
//go:generate bash ../../../../hack/run-tool.sh controller-gen rbac:headerFile=../../../../hack/boilerplate/sh.txt,roleName=ate-controller paths="./..." output:rbac:artifacts:config=../../../../manifests/ate-install/generated/
|
||||
|
||||
@@ -28,7 +28,10 @@ import (
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/status"
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
apierrors "k8s.io/apimachinery/pkg/api/errors"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/util/wait"
|
||||
corev1client "k8s.io/client-go/kubernetes/typed/core/v1"
|
||||
"k8s.io/client-go/tools/cache"
|
||||
"k8s.io/client-go/util/workqueue"
|
||||
)
|
||||
@@ -95,6 +98,7 @@ func (k workerKey) logAttrs() []any {
|
||||
// backoff on transient failures such as a lost version precondition.
|
||||
type WorkerPoolSyncer struct {
|
||||
client ateapipb.ControlClient
|
||||
pods corev1client.PodsGetter
|
||||
workerInformer cache.SharedIndexInformer
|
||||
workerPoolInformer cache.SharedIndexInformer
|
||||
queue workqueue.TypedRateLimitingInterface[workerKey]
|
||||
@@ -106,10 +110,12 @@ type WorkerPoolSyncer struct {
|
||||
listCap time.Duration
|
||||
}
|
||||
|
||||
// NewWorkerPoolSyncer creates a new WorkerPoolSyncer.
|
||||
func NewWorkerPoolSyncer(client ateapipb.ControlClient, workerInformer, workerPoolInformer cache.SharedIndexInformer) *WorkerPoolSyncer {
|
||||
// NewWorkerPoolSyncer creates a new WorkerPoolSyncer. pods is used to delete
|
||||
// worker pods that have reached a terminal phase.
|
||||
func NewWorkerPoolSyncer(client ateapipb.ControlClient, pods corev1client.PodsGetter, workerInformer, workerPoolInformer cache.SharedIndexInformer) *WorkerPoolSyncer {
|
||||
return &WorkerPoolSyncer{
|
||||
client: client,
|
||||
pods: pods,
|
||||
workerInformer: workerInformer,
|
||||
workerPoolInformer: workerPoolInformer,
|
||||
queue: workqueue.NewTypedRateLimitingQueue(workqueue.DefaultTypedControllerRateLimiter[workerKey]()),
|
||||
@@ -269,6 +275,12 @@ func (s *WorkerPoolSyncer) reconcile(ctx context.Context, key workerKey) error {
|
||||
// Deleted event.
|
||||
return s.markWorkerDraining(ctx, key)
|
||||
}
|
||||
// Checked before eligibility for the same reason: a terminal pod is never
|
||||
// Ready, so the eligibility gate would leave its Worker, and the Actors bound
|
||||
// to it, registered for as long as the pod object lingers.
|
||||
if isPodTerminal(pod) {
|
||||
return s.deleteTerminalPod(ctx, key, pod)
|
||||
}
|
||||
if !isWorkerEligible(pod) {
|
||||
// The pod has no IP or is not Ready yet; a later update event re-enqueues
|
||||
// it. A registered Worker still takes a raised epoch: an ateom that is
|
||||
@@ -423,6 +435,40 @@ func isWorkerEligible(pod *corev1.Pod) bool {
|
||||
return false
|
||||
}
|
||||
|
||||
// isPodTerminal reports whether every container in the pod has stopped for
|
||||
// good: a terminal phase is never left, so the pod will not serve again.
|
||||
func isPodTerminal(pod *corev1.Pod) bool {
|
||||
return pod.Status.Phase == corev1.PodSucceeded || pod.Status.Phase == corev1.PodFailed
|
||||
}
|
||||
|
||||
// deleteTerminalPod deletes a worker pod that has reached a terminal phase.
|
||||
// The API server deletes a terminal pod without a grace period, and the
|
||||
// resulting Pod Deleted event deregisters the Worker and releases its Actors
|
||||
// through reconcileDeadWorker. The Worker is marked DRAINING first so the
|
||||
// scheduler stops routing to it even while a failed delete is being retried.
|
||||
//
|
||||
// The delete is preconditioned on the key's UID so it can never remove a
|
||||
// same-named replacement. A pod already gone, or replaced, is the state this
|
||||
// drives towards, so NotFound and Conflict are success.
|
||||
func (s *WorkerPoolSyncer) deleteTerminalPod(ctx context.Context, key workerKey, pod *corev1.Pod) error {
|
||||
if err := s.markWorkerDraining(ctx, key); err != nil {
|
||||
return err
|
||||
}
|
||||
slog.InfoContext(ctx, "Syncer: deleting worker pod (terminal phase)",
|
||||
append(key.logAttrs(), slog.String("phase", string(pod.Status.Phase)))...)
|
||||
uid := pod.UID
|
||||
err := s.pods.Pods(key.namespace).Delete(ctx, key.name, metav1.DeleteOptions{
|
||||
Preconditions: &metav1.Preconditions{UID: &uid},
|
||||
})
|
||||
if apierrors.IsNotFound(err) || apierrors.IsConflict(err) {
|
||||
return nil
|
||||
}
|
||||
if err != nil {
|
||||
return fmt.Errorf("deleting terminal pod: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// markWorkerDraining transitions a worker to STATE_DRAINING so the scheduler
|
||||
// stops routing new actors to it while its pod is Terminating. DrainWorker is
|
||||
// idempotent, so a worker already draining costs nothing. If the worker is
|
||||
|
||||
@@ -29,11 +29,14 @@ import (
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/status"
|
||||
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"
|
||||
"k8s.io/apimachinery/pkg/util/wait"
|
||||
"k8s.io/client-go/kubernetes/fake"
|
||||
k8stesting "k8s.io/client-go/testing"
|
||||
"k8s.io/client-go/tools/cache"
|
||||
)
|
||||
|
||||
@@ -126,7 +129,7 @@ func setupSyncerTest(t *testing.T, ctx context.Context, api *fakeControl, initPo
|
||||
|
||||
// Start before the factory: the informer's initial list is what seeds the
|
||||
// queue with the pods that already exist.
|
||||
NewWorkerPoolSyncer(api, workerInformer, workerPoolInformer).Start(ctx)
|
||||
NewWorkerPoolSyncer(api, fakeK8s.CoreV1(), workerInformer, workerPoolInformer).Start(ctx)
|
||||
workerFactory.Start(ctx.Done())
|
||||
workerFactory.WaitForCacheSync(ctx.Done())
|
||||
|
||||
@@ -136,14 +139,16 @@ func setupSyncerTest(t *testing.T, ctx context.Context, api *fakeControl, initPo
|
||||
// setupReconcileTest builds a syncer whose pod and pool caches can be seeded
|
||||
// directly, for tests that drive reconcile synchronously without starting
|
||||
// factories or worker goroutines. It returns those caches alongside the syncer.
|
||||
// The syncer's pod client is a fake clientset that the caches do not watch.
|
||||
func setupReconcileTest(t *testing.T, api *fakeControl, initPools ...*atev1alpha1.WorkerPool) (*WorkerPoolSyncer, cache.Indexer, cache.Indexer) {
|
||||
t.Helper()
|
||||
|
||||
//nolint:staticcheck // NewSimpleClientset is what the informer machinery takes.
|
||||
_, workerInformer := WorkerPodInformer(fake.NewSimpleClientset())
|
||||
fakeK8s := fake.NewSimpleClientset()
|
||||
_, workerInformer := WorkerPodInformer(fakeK8s)
|
||||
workerPoolInformer, poolIndexer := newWorkerPoolInformer(t, initPools...)
|
||||
|
||||
return NewWorkerPoolSyncer(api, workerInformer, workerPoolInformer), workerInformer.GetIndexer(), poolIndexer
|
||||
return NewWorkerPoolSyncer(api, fakeK8s.CoreV1(), workerInformer, workerPoolInformer), workerInformer.GetIndexer(), poolIndexer
|
||||
}
|
||||
|
||||
// seedPod puts a pod in the syncer's cache as though the informer had delivered
|
||||
@@ -716,6 +721,112 @@ func TestSyncer_SoftDelete_ViaInformer(t *testing.T) {
|
||||
})
|
||||
}
|
||||
|
||||
// TestSyncer_TerminalPod_DeletesPod verifies that a worker pod in a terminal
|
||||
// phase is deleted, preconditioned on its UID, and its worker marked DRAINING
|
||||
// until the Pod Deleted event removes the record.
|
||||
func TestSyncer_TerminalPod_DeletesPod(t *testing.T) {
|
||||
ns, poolName, podName, ip := "ns-terminal", "pool1", "worker-terminal", "10.0.0.20"
|
||||
|
||||
for _, phase := range []corev1.PodPhase{corev1.PodFailed, corev1.PodSucceeded} {
|
||||
t.Run(string(phase), func(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
api := newFakeControl()
|
||||
api.put(registeredWorker(ns, poolName, podName, testPodUID, ip))
|
||||
s, pods, _ := setupReconcileTest(t, api, workerPool(ns, poolName, "gvisor", nil))
|
||||
|
||||
pod := workerPod(ns, podName, poolName, testPodUID, ip)
|
||||
pod.Status.Phase = phase
|
||||
//nolint:staticcheck // NewSimpleClientset is the fake the syncer's pod client takes.
|
||||
fakeK8s := fake.NewSimpleClientset(pod.DeepCopy())
|
||||
s.pods = fakeK8s.CoreV1()
|
||||
mustReconcile(t, ctx, s, seedPod(t, pods, pod))
|
||||
|
||||
if _, err := fakeK8s.CoreV1().Pods(ns).Get(ctx, podName, metav1.GetOptions{}); !apierrors.IsNotFound(err) {
|
||||
t.Errorf("get pod after reconcile = %v, want NotFound", err)
|
||||
}
|
||||
var deletes []k8stesting.DeleteAction
|
||||
for _, a := range fakeK8s.Actions() {
|
||||
if d, ok := a.(k8stesting.DeleteAction); ok {
|
||||
deletes = append(deletes, d)
|
||||
}
|
||||
}
|
||||
if len(deletes) != 1 {
|
||||
t.Fatalf("pod deletes = %d, want 1", len(deletes))
|
||||
}
|
||||
if got := deletes[0].GetDeleteOptions().Preconditions; got == nil || got.UID == nil || *got.UID != testPodUID {
|
||||
t.Errorf("delete preconditions = %+v, want UID %s", got, testPodUID)
|
||||
}
|
||||
if got := api.get(testPodUID).GetStatus().GetState(); got != ateapipb.WorkerState_WORKER_STATE_DRAINING {
|
||||
t.Errorf("worker state = %v, want DRAINING", got)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestSyncer_TerminalPod_AlreadyGoneOrReplaced verifies that a terminal pod
|
||||
// whose delete finds it gone (NotFound) or replaced under its name (the UID
|
||||
// precondition fails with Conflict) reconciles cleanly rather than requeueing.
|
||||
func TestSyncer_TerminalPod_AlreadyGoneOrReplaced(t *testing.T) {
|
||||
ns, poolName, podName := "ns-terminal-gone", "pool1", "worker-terminal-gone"
|
||||
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
err error
|
||||
}{
|
||||
{"not found", apierrors.NewNotFound(corev1.Resource("pods"), podName)},
|
||||
{"conflict", apierrors.NewConflict(corev1.Resource("pods"), podName, errors.New("precondition failed: UID"))},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
s, pods, _ := setupReconcileTest(t, newFakeControl())
|
||||
|
||||
//nolint:staticcheck // NewSimpleClientset is the fake the syncer's pod client takes.
|
||||
fakeK8s := fake.NewSimpleClientset()
|
||||
fakeK8s.PrependReactor("delete", "pods", func(k8stesting.Action) (bool, runtime.Object, error) {
|
||||
return true, nil, tc.err
|
||||
})
|
||||
s.pods = fakeK8s.CoreV1()
|
||||
|
||||
pod := workerPod(ns, podName, poolName, testPodUID, "")
|
||||
pod.Status.Phase = corev1.PodFailed
|
||||
mustReconcile(t, ctx, s, seedPod(t, pods, pod))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestSyncer_TerminalPod_ViaInformer walks a registered worker through its pod
|
||||
// failing: the syncer deletes the pod, and the resulting Pod Deleted event
|
||||
// removes the worker record, which is what releases its Actors.
|
||||
func TestSyncer_TerminalPod_ViaInformer(t *testing.T) {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||
defer cancel()
|
||||
|
||||
ns, podName, poolName := "ns-syncer-terminal", "worker-terminal-1", "pool1"
|
||||
|
||||
api := newFakeControl()
|
||||
fakeK8s := setupSyncerTest(t, ctx, api, workerPool(ns, poolName, "gvisor", nil))
|
||||
|
||||
pod := workerPod(ns, podName, poolName, testPodUID, "10.0.0.21")
|
||||
if _, err := fakeK8s.CoreV1().Pods(ns).Create(ctx, pod, metav1.CreateOptions{}); err != nil {
|
||||
t.Fatalf("create pod: %v", err)
|
||||
}
|
||||
waitForWorker(t, ctx, api, testPodUID, func(w *ateapipb.Worker) bool {
|
||||
return w.GetStatus().GetState() == ateapipb.WorkerState_WORKER_STATE_ACTIVE
|
||||
})
|
||||
|
||||
failed := pod.DeepCopy()
|
||||
failed.Status.Phase = corev1.PodFailed
|
||||
failed.Status.Conditions = []corev1.PodCondition{{Type: corev1.PodReady, Status: corev1.ConditionFalse}}
|
||||
if _, err := fakeK8s.CoreV1().Pods(ns).UpdateStatus(ctx, failed, metav1.UpdateOptions{}); err != nil {
|
||||
t.Fatalf("update pod status: %v", err)
|
||||
}
|
||||
|
||||
waitForWorker(t, ctx, api, testPodUID, func(w *ateapipb.Worker) bool { return w == nil })
|
||||
if _, err := fakeK8s.CoreV1().Pods(ns).Get(ctx, podName, metav1.GetOptions{}); !apierrors.IsNotFound(err) {
|
||||
t.Errorf("get pod = %v, want NotFound", err)
|
||||
}
|
||||
}
|
||||
|
||||
// TestSyncer_PodRecreatedWithNewUID verifies that when a pod is deleted and
|
||||
// recreated under the same name, the registry converges to the new pod's UID
|
||||
// and IP. A coalesced Delete+Create surfaces as a single update event, which
|
||||
|
||||
@@ -253,7 +253,7 @@ func main() {
|
||||
// Start registers the informer event handlers, so it has to run before the
|
||||
// factory does: the initial list then synthesizes an Add for every pod that
|
||||
// already exists, and no explicit startup re-list is needed.
|
||||
workersync.NewWorkerPoolSyncer(ateapiClient, workerPodInformer, workerPoolInformer.Informer()).Start(runCtx)
|
||||
workersync.NewWorkerPoolSyncer(ateapiClient, k8sClient.CoreV1(), workerPodInformer, workerPoolInformer.Informer()).Start(runCtx)
|
||||
|
||||
workerPodInformerFactory.Start(runCtx.Done())
|
||||
ateFactory.Start(runCtx.Done())
|
||||
|
||||
@@ -22,6 +22,14 @@ rules:
|
||||
- ""
|
||||
resources:
|
||||
- pods
|
||||
verbs:
|
||||
- delete
|
||||
- get
|
||||
- list
|
||||
- watch
|
||||
- apiGroups:
|
||||
- ""
|
||||
resources:
|
||||
- secrets
|
||||
verbs:
|
||||
- get
|
||||
|
||||
Reference in New Issue
Block a user