From dc898993ba5c6537d2e26ff7721c68fc4e38ec4c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Marc=20Schottst=C3=A4dt?= Date: Sun, 30 Aug 2026 01:30:39 +0200 Subject: [PATCH] fix(kubernetes): pin oldest pending consumer --- internal/runtime/kubernetes/resources_test.go | 26 +++++++++++++++++++ internal/runtime/kubernetes/scheduling.go | 14 ++++++++-- 2 files changed, 38 insertions(+), 2 deletions(-) diff --git a/internal/runtime/kubernetes/resources_test.go b/internal/runtime/kubernetes/resources_test.go index a383383..ae58d81 100644 --- a/internal/runtime/kubernetes/resources_test.go +++ b/internal/runtime/kubernetes/resources_test.go @@ -162,6 +162,32 @@ func TestPinPodToRuntimeNodeUsesScheduledPendingPVCConsumer(t *testing.T) { } } +func TestPinPodToRuntimeNodeUsesOldestScheduledPendingPVCConsumer(t *testing.T) { + pvc := "druid-static-web-data" + root := ref("druid", pvc) + older := runningProcedurePod("druid", root, "start", "coldstart", 1, "runtime", "start-job") + older.Name = "runtime-pod" + older.Status.Phase = corev1.PodPending + older.Spec.Volumes = []corev1.Volume{pvcVolume("data", pvc)} + older.CreationTimestamp = metav1.NewTime(time.Unix(1, 0)) + newer := runningProcedurePod("druid", root, "watch", "watch", 1, "watcher", "watch-job") + newer.Name = "watcher-pod" + newer.Status.Phase = corev1.PodPending + newer.Spec.NodeName = "node-b" + newer.Spec.Volumes = []corev1.Volume{pvcVolume("data", pvc)} + newer.CreationTimestamp = metav1.NewTime(time.Unix(2, 0)) + client := fake.NewSimpleClientset(newer, older) + backend := NewWithClient(Config{Namespace: "druid"}, client) + podSpec := corev1.PodSpec{} + + if err := backend.pinPodToRuntimeNode(context.Background(), "druid", pvc, &podSpec); err != nil { + t.Fatal(err) + } + if got := podSpec.NodeSelector[corev1.LabelHostname]; got != "node-a" { + t.Fatalf("expected oldest pending consumer node-a, got %q", got) + } +} + func TestPinPodToRuntimeNodeIgnoresInactiveOrUnrelatedPods(t *testing.T) { root := ref("druid", "druid-static-web-data") inactive := runningProcedurePod("druid", root, "start", "start", 1, "inactive", "start-job") diff --git a/internal/runtime/kubernetes/scheduling.go b/internal/runtime/kubernetes/scheduling.go index abebcfc..00af039 100644 --- a/internal/runtime/kubernetes/scheduling.go +++ b/internal/runtime/kubernetes/scheduling.go @@ -2,6 +2,7 @@ package kubernetes import ( "context" + "time" corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -31,6 +32,8 @@ func (b *Backend) runtimePVCNode(ctx context.Context, namespace string, pvc stri if err != nil { return "", err } + var pendingNode string + var pendingCreated time.Time for idx := range pods.Items { pod := &pods.Items[idx] if pod.Status.Phase != corev1.PodPending && pod.Status.Phase != corev1.PodRunning { @@ -39,9 +42,16 @@ func (b *Backend) runtimePVCNode(ctx context.Context, namespace string, pvc stri if pod.Spec.NodeName == "" || !podUsesPVC(pod, pvc) { continue } - return pod.Spec.NodeName, nil + if pod.Status.Phase == corev1.PodRunning { + return pod.Spec.NodeName, nil + } + created := pod.CreationTimestamp.Time + if pendingNode == "" || created.Before(pendingCreated) { + pendingNode = pod.Spec.NodeName + pendingCreated = created + } } - return "", nil + return pendingNode, nil } func podUsesPVC(pod *corev1.Pod, pvc string) bool {