diff --git a/cmd/runner/main.go b/cmd/runner/main.go index 09ae98f..b4a1c3d 100644 --- a/cmd/runner/main.go +++ b/cmd/runner/main.go @@ -15,14 +15,15 @@ import ( // LogEvent represents a structured log event emitted by the runner. type LogEvent struct { - Kind string `json:"kind"` // "control" or "data" - Event string `json:"event,omitempty"` // "start", "end" (for control events) - StepID int `json:"step_id"` // 0-based step index - StepName string `json:"step_name,omitempty"` // Step name - Stream string `json:"stream,omitempty"` // "stdout" or "stderr" (for data events) - Content string `json:"content,omitempty"` // Log line content (for data events) - ExitCode *int `json:"exit_code,omitempty"` // Exit code (for control/end events) - Timestamp string `json:"timestamp"` // ISO 8601 timestamp + Kind string `json:"kind"` // "control" or "data" + Event string `json:"event,omitempty"` // "start", "end" (for control events) + StepID int `json:"step_id"` // 0-based step index + StepName string `json:"step_name,omitempty"` // Step name + WorkflowName string `json:"workflow_name"` // Workflow name for log separation + Stream string `json:"stream,omitempty"` // "stdout" or "stderr" (for data events) + Content string `json:"content,omitempty"` // Log line content (for data events) + ExitCode *int `json:"exit_code,omitempty"` // Exit code (for control/end events) + Timestamp string `json:"timestamp"` // ISO 8601 timestamp } func main() { @@ -56,7 +57,7 @@ func run() error { // Execute each step ctx := context.Background() for i, step := range workflow.Steps { - if err := executeStep(ctx, i, step); err != nil { + if err := executeStep(ctx, i, step, workflow.Name); err != nil { return fmt.Errorf("step %d (%s) failed: %w", i, step.Name, err) } } @@ -64,9 +65,9 @@ func run() error { return nil } -func executeStep(ctx context.Context, stepID int, step loomv1alpha1.WorkflowStep) error { +func executeStep(ctx context.Context, stepID int, step loomv1alpha1.WorkflowStep, workflowName string) error { // Emit step start event - emitControlEvent(stepID, step.Name, "start", nil) + emitControlEvent(stepID, step.Name, workflowName, "start", nil) // Set step-specific environment variables if step.Environment != nil { @@ -95,14 +96,14 @@ func executeStep(ctx context.Context, stepID int, step loomv1alpha1.WorkflowStep // Start the command if err := cmd.Start(); err != nil { exitCode := 1 - emitControlEvent(stepID, step.Name, "end", &exitCode) + emitControlEvent(stepID, step.Name, workflowName, "end", &exitCode) return fmt.Errorf("failed to start command: %w", err) } // Stream stdout and stderr concurrently done := make(chan error, 2) - go streamOutput(stdout, stepID, "stdout", done) - go streamOutput(stderr, stepID, "stderr", done) + go streamOutput(stdout, stepID, workflowName, "stdout", done) + go streamOutput(stderr, stepID, workflowName, "stderr", done) // Wait for both streams to complete for i := 0; i < 2; i++ { @@ -124,7 +125,7 @@ func executeStep(ctx context.Context, stepID int, step loomv1alpha1.WorkflowStep } // Emit step end event - emitControlEvent(stepID, step.Name, "end", &exitCode) + emitControlEvent(stepID, step.Name, workflowName, "end", &exitCode) if exitCode != 0 { return fmt.Errorf("command exited with code %d", exitCode) @@ -133,7 +134,7 @@ func executeStep(ctx context.Context, stepID int, step loomv1alpha1.WorkflowStep return nil } -func streamOutput(reader io.Reader, stepID int, stream string, done chan<- error) { +func streamOutput(reader io.Reader, stepID int, workflowName, stream string, done chan<- error) { scanner := bufio.NewScanner(reader) // Increase buffer size for long lines buf := make([]byte, 0, 64*1024) @@ -141,31 +142,33 @@ func streamOutput(reader io.Reader, stepID int, stream string, done chan<- error for scanner.Scan() { line := scanner.Text() - emitDataEvent(stepID, stream, line) + emitDataEvent(stepID, workflowName, stream, line) } done <- scanner.Err() } -func emitControlEvent(stepID int, stepName, event string, exitCode *int) { +func emitControlEvent(stepID int, stepName, workflowName, event string, exitCode *int) { ev := LogEvent{ - Kind: "control", - Event: event, - StepID: stepID, - StepName: stepName, - ExitCode: exitCode, - Timestamp: time.Now().UTC().Format(time.RFC3339Nano), + Kind: "control", + Event: event, + StepID: stepID, + StepName: stepName, + WorkflowName: workflowName, + ExitCode: exitCode, + Timestamp: time.Now().UTC().Format(time.RFC3339Nano), } emitJSON(ev) } -func emitDataEvent(stepID int, stream, content string) { +func emitDataEvent(stepID int, workflowName, stream, content string) { ev := LogEvent{ - Kind: "data", - StepID: stepID, - Stream: stream, - Content: content, - Timestamp: time.Now().UTC().Format(time.RFC3339Nano), + Kind: "data", + StepID: stepID, + WorkflowName: workflowName, + Stream: stream, + Content: content, + Timestamp: time.Now().UTC().Format(time.RFC3339Nano), } emitJSON(ev) } diff --git a/config/manager/kustomization.yaml b/config/manager/kustomization.yaml index 0d68922..ef26823 100644 --- a/config/manager/kustomization.yaml +++ b/config/manager/kustomization.yaml @@ -8,4 +8,4 @@ kind: Kustomization images: - name: controller newName: atcr.io/evan.jarrett.net/loom - newTag: v0.0.8 + newTag: v0.0.9 diff --git a/internal/controller/spindleset_controller.go b/internal/controller/spindleset_controller.go index b710491..bb1462f 100644 --- a/internal/controller/spindleset_controller.go +++ b/internal/controller/spindleset_controller.go @@ -252,6 +252,10 @@ func findCondition(conditions []metav1.Condition, conditionType string) *metav1. // setCondition adds or updates a condition in the list func setCondition(conditions *[]metav1.Condition, newCondition metav1.Condition) { if conditions == nil { + return // Pointer itself is nil, nothing we can do + } + + if *conditions == nil { *conditions = []metav1.Condition{} } @@ -313,6 +317,7 @@ func (r *SpindleSetReconciler) ensurePipelineJobs(ctx context.Context, spindleSe Image: workflowSpec.Image, Architecture: workflowSpec.Architecture, Steps: jobSteps, + WorkflowSpec: workflowSpec, // Pass full workflow spec to runner RepoURL: pipelineRun.RepoURL, CommitSHA: pipelineRun.CommitSHA, Secrets: nil, // TODO: Handle secrets @@ -334,6 +339,11 @@ func (r *SpindleSetReconciler) ensurePipelineJobs(ctx context.Context, spindleSe logger.Info("Creating Job for workflow", "workflow", workflowSpec.Name, "job", job.Name) if err := r.Create(ctx, job); err != nil { + if apierrors.IsAlreadyExists(err) { + // Job already exists (possibly from previous deployment), skip + logger.Info("Job already exists, skipping creation", "workflow", workflowSpec.Name, "job", job.Name) + continue + } return fmt.Errorf("failed to create job for workflow %s: %w", workflowSpec.Name, err) } diff --git a/internal/engine/kubernetes_engine.go b/internal/engine/kubernetes_engine.go index 5c4982d..d54d7c4 100644 --- a/internal/engine/kubernetes_engine.go +++ b/internal/engine/kubernetes_engine.go @@ -412,14 +412,15 @@ func (e *KubernetesEngine) streamJobLogs(ctx context.Context, job *batchv1.Job, // LogEvent represents a structured log event from the runner binary type LogEvent struct { - Kind string `json:"kind"` // "control" or "data" - Event string `json:"event,omitempty"` // "start", "end" (for control events) - StepID int `json:"step_id"` // 0-based step index - StepName string `json:"step_name,omitempty"` // Step name - Stream string `json:"stream,omitempty"` // "stdout" or "stderr" (for data events) - Content string `json:"content,omitempty"` // Log line content (for data events) - ExitCode *int `json:"exit_code,omitempty"` // Exit code (for control/end events) - Timestamp string `json:"timestamp"` // ISO 8601 timestamp + Kind string `json:"kind"` // "control" or "data" + Event string `json:"event,omitempty"` // "start", "end" (for control events) + StepID int `json:"step_id"` // 0-based step index + StepName string `json:"step_name,omitempty"` // Step name + WorkflowName string `json:"workflow_name"` // Workflow name for log separation + Stream string `json:"stream,omitempty"` // "stdout" or "stderr" (for data events) + Content string `json:"content,omitempty"` // Log line content (for data events) + ExitCode *int `json:"exit_code,omitempty"` // Exit code (for control/end events) + Timestamp string `json:"timestamp"` // ISO 8601 timestamp } // parseLogs reads log lines as JSON events and sends them to WorkflowLogger