From e42b257a3d463400a448c5728f7d666bd43af746 Mon Sep 17 00:00:00 2001 From: Evan Jarrett Date: Mon, 24 Nov 2025 21:28:45 -0600 Subject: [PATCH] various fixes --- internal/engine/kubernetes_engine.go | 48 +++++++++++++++++++++------- 1 file changed, 37 insertions(+), 11 deletions(-) diff --git a/internal/engine/kubernetes_engine.go b/internal/engine/kubernetes_engine.go index 53e954f..146ad85 100644 --- a/internal/engine/kubernetes_engine.go +++ b/internal/engine/kubernetes_engine.go @@ -30,9 +30,10 @@ import ( // workflowLogStream holds the state for streaming logs from a workflow's pod type workflowLogStream struct { - scanner *bufio.Scanner - stream io.ReadCloser - pod *corev1.Pod + scanner *bufio.Scanner + stream io.ReadCloser + pod *corev1.Pod + podPhase corev1.PodPhase // Track pod phase at stream creation time } // KubernetesEngine implements the spindle Engine interface for Kubernetes Jobs. @@ -278,7 +279,13 @@ func (e *KubernetesEngine) DestroyWorkflow(ctx context.Context, wid models.Workf PropagationPolicy: &deletePolicy, } - if err := e.client.Delete(ctx, spindleSet, deleteOptions); err != nil { + // Use a fresh context for cleanup to ensure deletion succeeds even if the + // original context was canceled (e.g., by errgroup when another workflow completes). + // This prevents orphaned SpindleSets when running multiple workflows in parallel. + cleanupCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + + if err := e.client.Delete(cleanupCtx, spindleSet, deleteOptions); err != nil { // Ignore not found errors (SpindleSet may have already been deleted) if client.IgnoreNotFound(err) != nil { return fmt.Errorf("failed to delete SpindleSet: %w", err) @@ -436,12 +443,19 @@ func (e *KubernetesEngine) getOrCreateLogStream(ctx context.Context, wid models. time.Sleep(1 * time.Second) } - logger.Info("Pod is ready, streaming logs", "podName", pod.Name, "phase", pod.Status.Phase) + // Only use Follow mode for running pods. For completed pods, we need to read + // existing logs (Follow:true only streams NEW logs after connection). + shouldFollow := pod.Status.Phase == corev1.PodRunning + if !shouldFollow { + logger.Info("Pod already completed, reading existing logs", "podName", pod.Name, "phase", pod.Status.Phase) + } else { + logger.Info("Pod is running, streaming logs", "podName", pod.Name, "phase", pod.Status.Phase) + } // Stream logs from the main container (not init containers) req := clientset.CoreV1().Pods(pod.Namespace).GetLogs(pod.Name, &corev1.PodLogOptions{ Container: "runner", - Follow: true, + Follow: shouldFollow, }) logStream, err := req.Stream(ctx) @@ -456,9 +470,10 @@ func (e *KubernetesEngine) getOrCreateLogStream(ctx context.Context, wid models. // Create and store stream stream = &workflowLogStream{ - scanner: scanner, - stream: logStream, - pod: pod, + scanner: scanner, + stream: logStream, + pod: pod, + podPhase: pod.Status.Phase, } e.streamMutex.Lock() @@ -534,8 +549,19 @@ func (e *KubernetesEngine) readUntilStepEnd(stream *workflowLogStream, stepID in } } - // If we get here, scanner ended without seeing step end event - // This could mean pod terminated early + // Scanner ended without seeing step end event. + // Check pod phase to determine if this is an error or expected behavior. + if stream.podPhase == corev1.PodSucceeded { + // Pod succeeded but we didn't find the control event - treat as success. + // This can happen if logs were truncated or runner didn't emit events. + return nil + } + + if stream.podPhase == corev1.PodFailed { + return fmt.Errorf("pod failed before step %d completed", stepID) + } + + // Pod was running when we started but stream ended unexpectedly return fmt.Errorf("log stream ended before step %d completed", stepID) } -- 2.51.2