Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
62 changes: 48 additions & 14 deletions api/v1alpha1/amaltheasession_children.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,6 @@ import (
"context"
"crypto/rand"
"encoding/base64"
"errors"
"fmt"
"io"
"maps"
Expand Down Expand Up @@ -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{
Expand Down Expand Up @@ -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()
Expand All @@ -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()
Expand All @@ -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()
Expand Down Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -7,3 +7,6 @@ spec:
selector:
matchLabels:
app.kubernetes.io/name: AmaltheaSession
matchExpressions:
- key: batch.kubernetes.io/job-name
operator: DoesNotExist
9 changes: 5 additions & 4 deletions internal/controller/children.go
Original file line number Diff line number Diff line change
Expand Up @@ -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, ""
Expand All @@ -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, ""
Expand All @@ -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
}
Expand Down
23 changes: 16 additions & 7 deletions internal/controller/children_conditions.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading