From 2d260393aa2b4cf123e45c6d92aafb67c5b4e8fe Mon Sep 17 00:00:00 2001 From: Wes Johnson Date: Mon, 28 Sep 2026 13:53:24 +0200 Subject: [PATCH 1/2] feat: allow renku jobs to run on spot instances --- api/v1alpha1/amaltheasession_children.go | 62 +++++++++++++++++----- internal/controller/children.go | 9 ++-- internal/controller/children_conditions.go | 23 +++++--- 3 files changed, 69 insertions(+), 25 deletions(-) diff --git a/api/v1alpha1/amaltheasession_children.go b/api/v1alpha1/amaltheasession_children.go index e2e3dbd7..69315b39 100644 --- a/api/v1alpha1/amaltheasession_children.go +++ b/api/v1alpha1/amaltheasession_children.go @@ -4,7 +4,6 @@ import ( "context" "crypto/rand" "encoding/base64" - "errors" "fmt" "io" "maps" @@ -205,9 +204,19 @@ func (cr *AmaltheaSession) Job(cfg config.AmaltheaSessionConfiguration) (batchv1 Annotations: annotations, }, Spec: batchv1.JobSpec{ - Parallelism: ptr.To(int32(1)), - Completions: ptr.To(int32(1)), - BackoffLimit: ptr.To(int32(0)), + Parallelism: ptr.To(int32(1)), + Completions: ptr.To(int32(1)), + BackoffLimit: ptr.To(int32(0)), + PodFailurePolicy: &batchv1.PodFailurePolicy{ + Rules: []batchv1.PodFailurePolicyRule{ + { + Action: batchv1.PodFailurePolicyActionIgnore, + OnPodConditions: []batchv1.PodFailurePolicyOnPodConditionsPattern{ + {Type: v1.DisruptionTarget, Status: v1.ConditionTrue}, + }, + }, + }, + }, ActiveDeadlineSeconds: activeTTL, Suspend: &cr.Spec.Hibernated, Template: v1.PodTemplateSpec{ @@ -510,7 +519,7 @@ func (cr *AmaltheaSession) NeedsDeletion() bool { } } -func (cr *AmaltheaSession) GetPod(ctx context.Context, clnt client.Client) (*v1.Pod, error) { +func (cr *AmaltheaSession) GetPod(ctx context.Context, clnt client.Reader) (*v1.Pod, error) { logger := log.FromContext(ctx) if cr.Spec.SessionType == SessionTypeNonInteractive { selector := labels.Set{"job-name": cr.JobName()}.AsSelector() @@ -520,16 +529,11 @@ func (cr *AmaltheaSession) GetPod(ctx context.Context, clnt client.Client) (*v1. logger.Info("cannot list pods for batch job", "job-name", cr.JobName(), "error", err) return nil, err } - itemLength := len(podList.Items) - if itemLength > 1 { - logger.Info("Too many pods returned for batch job", "job-name", cr.JobName(), "num_pods", itemLength) - return nil, errors.New("more than one pod found for job") - } - if itemLength > 0 { - return &podList.Items[0], nil + pod := currentJobPod(podList.Items) + if pod == nil { + logger.Info("No pod found for batch job", "job-name", cr.JobName()) } - logger.Info("No pod found for batch job", "job-name", cr.JobName()) - return nil, nil + return pod, nil } else { pod := v1.Pod{} podName := cr.PodName() @@ -542,6 +546,26 @@ func (cr *AmaltheaSession) GetPod(ctx context.Context, clnt client.Client) (*v1. } } +func currentJobPod(pods []v1.Pod) *v1.Pod { + var newest, newestActive *v1.Pod + for i := range pods { + pod := &pods[i] + if newest == nil || newest.CreationTimestamp.Before(&pod.CreationTimestamp) { + newest = pod + } + if pod.Status.Phase == v1.PodSucceeded || pod.Status.Phase == v1.PodFailed { + continue + } + if newestActive == nil || newestActive.CreationTimestamp.Before(&pod.CreationTimestamp) { + newestActive = pod + } + } + if newestActive != nil { + return newestActive + } + return newest +} + func (cr *AmaltheaSession) GetJob(ctx context.Context, clnt client.Client) (*batchv1.Job, error) { job := batchv1.Job{} jobName := cr.JobName() @@ -569,6 +593,16 @@ func (as *AmaltheaSession) GetPodEvents(ctx context.Context, c client.Reader) (* logger := log.FromContext(ctx) events := v1.EventList{} podName := as.PodName() + if as.Spec.SessionType == SessionTypeNonInteractive { + pod, err := as.GetPod(ctx, c) + if err != nil { + return nil, err + } + if pod == nil { + return nil, nil + } + podName = pod.Name + } logger.Info("Getting event list for pod", "pod", podName) err := c.List(ctx, &events, diff --git a/internal/controller/children.go b/internal/controller/children.go index 0e0f0d3f..9c52a09e 100644 --- a/internal/controller/children.go +++ b/internal/controller/children.go @@ -526,7 +526,8 @@ func (c ChildResourceUpdates) IsRunning(pod *v1.Pod) bool { } func (c ChildResourceUpdates) State(cr *amaltheadevv1alpha1.AmaltheaSession, pod *v1.Pod, job *batchv1.Job) (amaltheadevv1alpha1.State, string) { - msg := c.failureMessage(pod, job) + isJob := cr.Spec.SessionType == amaltheadevv1alpha1.SessionTypeNonInteractive + msg := c.failureMessage(pod, job, isJob) switch { case cr.GetDeletionTimestamp() != nil: return amaltheadevv1alpha1.NotReady, "" @@ -538,7 +539,7 @@ func (c ChildResourceUpdates) State(cr *amaltheadevv1alpha1.AmaltheaSession, pod return amaltheadevv1alpha1.Failed, msg case c.IsRunning(pod): return amaltheadevv1alpha1.Running, "" - case podIsCompleted(pod): + case !isJob && podIsCompleted(pod): fallthrough case jobIsSuccess(job): return amaltheadevv1alpha1.Succeeded, "" @@ -547,8 +548,8 @@ func (c ChildResourceUpdates) State(cr *amaltheadevv1alpha1.AmaltheaSession, pod } } -func (c ChildResourceUpdates) failureMessage(pod *v1.Pod, job *batchv1.Job) string { - msg := podFailureReason(pod) +func (c ChildResourceUpdates) failureMessage(pod *v1.Pod, job *batchv1.Job, isJob bool) string { + msg := podFailureReason(pod, isJob) if msg != "" { return msg } diff --git a/internal/controller/children_conditions.go b/internal/controller/children_conditions.go index 1e12b876..75507902 100644 --- a/internal/controller/children_conditions.go +++ b/internal/controller/children_conditions.go @@ -32,24 +32,33 @@ func jobFailureReason(job *batchv1.Job) string { return "" } - if job.Status.Failed > 0 { - for _, cond := range job.Status.Conditions { - if (cond.Type == batchv1.JobFailed || cond.Type == batchv1.JobFailureTarget) && cond.Status == v1.ConditionTrue { - return cond.Reason + ": " + cond.Message - } + for _, cond := range job.Status.Conditions { + if (cond.Type == batchv1.JobFailed || cond.Type == batchv1.JobFailureTarget) && cond.Status == v1.ConditionTrue { + return cond.Reason + ": " + cond.Message } - return "The job failed." } return "" } -func podFailureReason(pod *v1.Pod) string { +func podIsDisrupted(pod *v1.Pod) bool { + for _, condition := range pod.Status.Conditions { + if condition.Type == v1.DisruptionTarget && condition.Status == v1.ConditionTrue { + return true + } + } + return false +} + +func podFailureReason(pod *v1.Pod, isJob bool) string { if pod == nil { return "" } if pod.GetDeletionTimestamp() != nil { return "" } + if isJob && podIsDisrupted(pod) { + return "" + } // NOTE: Checking the pod phase is not useful because it will still say "Running" when a container // fails to start for example when the command is referencing an executable that does not exist. // NOTE: The messages from the container statuses list are more descriptive than the conditions From 4a6c21262d427b1c4ec3bc9567776a5ba4b66695 Mon Sep 17 00:00:00 2001 From: Wes Johnson Date: Tue, 29 Sep 2026 11:24:28 +0200 Subject: [PATCH 2/2] feat: allow renku jobs to run on spot instances --- .../amalthea-sessions/templates/poddisruptionbudget.yaml | 3 +++ 1 file changed, 3 insertions(+) diff --git a/helm-chart/amalthea-sessions/templates/poddisruptionbudget.yaml b/helm-chart/amalthea-sessions/templates/poddisruptionbudget.yaml index 170d8b0a..c93527cd 100644 --- a/helm-chart/amalthea-sessions/templates/poddisruptionbudget.yaml +++ b/helm-chart/amalthea-sessions/templates/poddisruptionbudget.yaml @@ -7,3 +7,6 @@ spec: selector: matchLabels: app.kubernetes.io/name: AmaltheaSession + matchExpressions: + - key: batch.kubernetes.io/job-name + operator: DoesNotExist