diff --git a/provider_tekton.go b/provider_tekton.go index 7c3201b..cdfdb2c 100644 --- a/provider_tekton.go +++ b/provider_tekton.go @@ -465,26 +465,45 @@ func (p *tektonProvider) Logs( return nil, fmt.Errorf("lookup tekton run mapping: %w", err) } if ref == nil { + // No mapping at all means this provider never spawned a + // PipelineRun for the requested tuple, so a 404 is the + // honest answer. return nil, ErrLogsNotFound } - taskRuns, err := p.taskRunsForPipelineRun(ctx, *ref) - if err != nil { - return nil, err - } - p.log.Debug("Logs: found TaskRuns for PipelineRun", - "pipeline_run", ref.PipelineRunName, "count", len(taskRuns), - ) - if len(taskRuns) == 0 { - return nil, ErrLogsNotFound - } - - terminal := p.isPipelineRunTerminal(ctx, *ref) - p.log.Debug("Logs: pipeline run terminal state", "pipeline_run", ref.PipelineRunName, "terminal", terminal) - + // At this point the workflow *was* spawned — we have a row in the + // store mapping the tuple to a PipelineRun. The TaskRuns the + // PipelineRun fans out to are created asynchronously by Tekton + // once the run is admitted, so a freshly-spawned or still-queueing + // PipelineRun will momentarily report zero TaskRuns. Returning + // ErrLogsNotFound here would mistranslate that into a 404 at the + // HTTP layer (see provider.go contract: ErrLogsNotFound means the + // workflow never ran), making just-spawned runs look nonexistent + // to the appview. Instead, hand back an open channel and poll + // inside the goroutine until TaskRuns materialize, ctx is + // cancelled, or the PipelineRun reaches a terminal state with + // nothing to stream. out := make(chan LogLine, 32) go func() { defer close(out) + + taskRuns := p.waitForTaskRuns(ctx, *ref) + if len(taskRuns) == 0 { + // Either ctx was cancelled, or the PipelineRun + // reached a terminal state without ever scheduling + // any TaskRuns. Both cases close an empty stream; + // there's nothing to send. + return + } + p.log.Debug("Logs: found TaskRuns for PipelineRun", + "pipeline_run", ref.PipelineRunName, "count", len(taskRuns), + ) + + terminal := p.isPipelineRunTerminal(ctx, *ref) + p.log.Debug("Logs: pipeline run terminal state", + "pipeline_run", ref.PipelineRunName, "terminal", terminal, + ) + stepID := 0 for _, tr := range taskRuns { taskName := tr.GetName() @@ -525,6 +544,44 @@ func (p *tektonProvider) Logs( return out, nil } +// waitForTaskRuns blocks until at least one TaskRun belonging to the +// PipelineRun in ref is observable, the PipelineRun reaches a terminal +// state with no TaskRuns, or ctx is cancelled. It exists so Logs can +// keep the channel open across the gap between PipelineRun creation +// and Tekton scheduling its TaskRuns, instead of mistranslating that +// gap into a 404. The first iteration always probes immediately so a +// PipelineRun that already has TaskRuns is returned without waiting. +func (p *tektonProvider) waitForTaskRuns( + ctx context.Context, + ref TektonRunRef, +) []k8s.Object { + const interval = 1 * time.Second + for { + taskRuns, err := p.taskRunsForPipelineRun(ctx, ref) + if err != nil { + p.log.Debug("waitForTaskRuns: list TaskRuns failed", + "err", err, "pipeline_run", ref.PipelineRunName, + ) + } else if len(taskRuns) > 0 { + return taskRuns + } + // If the PipelineRun has already reached a terminal state and + // still has no TaskRuns, there's nothing more to wait for — + // return empty so the caller closes the stream cleanly. + if p.isPipelineRunTerminal(ctx, ref) { + p.log.Debug("waitForTaskRuns: PipelineRun terminal with no TaskRuns", + "pipeline_run", ref.PipelineRunName, + ) + return nil + } + select { + case <-ctx.Done(): + return nil + case <-time.After(interval): + } + } +} + // isPipelineRunTerminal returns true if the PipelineRun is in a terminal state right now. func (p *tektonProvider) isPipelineRunTerminal(ctx context.Context, ref TektonRunRef) bool { obj, err := p.client.GetObject(ctx, pipelineRunsGVR, ref.Namespace, diff --git a/provider_tekton_test.go b/provider_tekton_test.go index f9f3a90..5201666 100644 --- a/provider_tekton_test.go +++ b/provider_tekton_test.go @@ -198,8 +198,42 @@ func TestTektonLogsLookup(t *testing.T) { if err := st.InsertTektonRun(ctx, ref); err != nil { t.Fatalf("insert ref: %v", err) } - if _, err := p.Logs(ctx, "knot.example.com", "rkey-1", "ci.yml"); !errors.Is(err, ErrLogsNotFound) { - t.Fatalf("logs before TaskRuns err = %v; want ErrLogsNotFound", err) + // With the mapping in place but no TaskRuns yet, Logs must NOT + // return ErrLogsNotFound: the workflow has been spawned and is + // just queueing inside Tekton. Surfacing 404 here mistranslates + // "still scheduling" as "doesn't exist" at the HTTP layer (see + // http.go's /logs handler). Verify we get an open channel that + // stays open until ctx is cancelled. + { + waitCtx, cancel := context.WithCancel(ctx) + ch, err := p.Logs(waitCtx, "knot.example.com", "rkey-1", "ci.yml") + if err != nil { + cancel() + t.Fatalf("logs before TaskRuns err = %v; want nil", err) + } + if ch == nil { + cancel() + t.Fatalf("logs before TaskRuns: nil channel") + } + // Channel must not produce any frames or close before we + // cancel. A premature close would mean the goroutine treated + // "no TaskRuns yet" as "stream done", which is what we just + // fixed. + select { + case line, ok := <-ch: + cancel() + if !ok { + t.Fatalf("logs channel closed before TaskRuns appeared") + } + t.Fatalf("unexpected frame before TaskRuns: %+v", line) + case <-time.After(50 * time.Millisecond): + } + cancel() + // After cancellation the producer goroutine must close the + // channel and not strand any frames. + for line := range ch { + t.Fatalf("unexpected frame after cancel: %+v", line) + } } client.SeedObject(taskRunsGVR, "ci", k8s.Object{