diff --git a/internal/k8s/fake.go b/internal/k8s/fake.go index 160e221..9d3646f 100644 --- a/internal/k8s/fake.go +++ b/internal/k8s/fake.go @@ -228,6 +228,7 @@ func (c *FakeClient) StreamPodLogs( namespace string, podName string, container string, + opts LogOptions, ) (io.ReadCloser, error) { select { case <-ctx.Done(): @@ -236,21 +237,57 @@ func (c *FakeClient) StreamPodLogs( } c.mu.Lock() - defer c.mu.Unlock() raw, ok := c.podLogs[podLogKey{ namespace: namespace, podName: podName, container: container, }] + c.mu.Unlock() if !ok { return nil, fmt.Errorf( "%w: pod log %s/%s[%s]", ErrNotFound, namespace, podName, container, ) } - return io.NopCloser(strings.NewReader(raw)), nil + if !opts.Follow { + return io.NopCloser(strings.NewReader(raw)), nil + } + // In Follow mode, mirror the kube apiserver: yield the seeded bytes, + // then keep the stream open until ctx is cancelled. Tests that exercise + // the live-streaming path drive termination by cancelling the context; + // tests that want a one-shot snapshot should leave Follow=false. + return &fakeFollowReader{rest: raw, ctx: ctx}, nil } +// fakeFollowReader returns the seeded log bytes, then blocks on ctx.Done() +// before EOFing. This matches the apiserver behaviour for follow=true: +// the connection only closes once the container terminates (or the caller +// hangs up), so callers using bufio.Scanner do not see a premature EOF and +// treat a still-running container as completed. +type fakeFollowReader struct { + rest string + ctx context.Context + done bool +} + +var _ io.ReadCloser = (*fakeFollowReader)(nil) + +func (r *fakeFollowReader) Read(p []byte) (int, error) { + if len(r.rest) > 0 { + n := copy(p, r.rest) + r.rest = r.rest[n:] + return n, nil + } + if r.done { + return 0, io.EOF + } + r.done = true + <-r.ctx.Done() + return 0, io.EOF +} + +func (r *fakeFollowReader) Close() error { return nil } + func (c *FakeClient) createObject( gvr GVR, namespace string, diff --git a/internal/k8s/incluster.go b/internal/k8s/incluster.go index 5641f28..d899ec6 100644 --- a/internal/k8s/incluster.go +++ b/internal/k8s/incluster.go @@ -242,11 +242,20 @@ func (c *InClusterClient) StreamPodLogs( namespace string, podName string, container string, + opts LogOptions, ) (io.ReadCloser, error) { query := url.Values{} if container != "" { query.Set("container", container) } + // follow=true makes the API server hold the connection open and push + // new log bytes as the container writes them, only closing when the + // container terminates. Without it, the response is a snapshot of + // whatever has been written so far, which silently truncates live + // runs. + if opts.Follow { + query.Set("follow", "true") + } resp, err := c.do(ctx, http.MethodGet, coreResourcePath(namespace, podsResource, podName, podLogsSubresource), query, nil, "") diff --git a/internal/k8s/k8s.go b/internal/k8s/k8s.go index 5fd46b6..6a0bafd 100644 --- a/internal/k8s/k8s.go +++ b/internal/k8s/k8s.go @@ -61,9 +61,23 @@ type Client interface { namespace string, podName string, container string, + opts LogOptions, ) (io.ReadCloser, error) } +// LogOptions controls how StreamPodLogs reads container logs. +type LogOptions struct { + // Follow tells the API server to keep the connection open and stream + // new log bytes as the container writes them, only EOFing when the + // container terminates (or ctx is cancelled). Without Follow, the + // returned reader is a snapshot: it yields the bytes that exist at + // request time and then EOFs immediately, even if the container is + // still running. Live-streaming callers MUST set Follow=true; using + // the snapshot mode for a still-running container makes the caller + // look at a frozen view and treat the container as if it had finished. + Follow bool +} + // GVR identifies a Kubernetes resource by the path segments the API server // routes on. type GVR struct { diff --git a/provider_tekton.go b/provider_tekton.go index cdfdb2c..1b8ced1 100644 --- a/provider_tekton.go +++ b/provider_tekton.go @@ -486,31 +486,83 @@ func (p *tektonProvider) Logs( out := make(chan LogLine, 32) go func() { defer close(out) + p.streamPipelineRunLogs(ctx, out, *ref) + }() + return out, nil +} - 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. +// streamPipelineRunLogs drives the Logs channel for a single PipelineRun. +// It polls for TaskRuns until the PipelineRun is terminal (or ctx is +// cancelled), streaming each TaskRun's logs exactly once. +// +// The loop is deliberately a poll rather than a one-shot pass for two +// reasons: +// +// 1. Live runs need follow semantics. Streaming a still-running TaskRun +// uses Follow=true on the pod log stream, so streamTaskRunLogs only +// returns once the container actually terminates. The outer loop +// then re-checks for new TaskRuns and the PipelineRun's terminal +// state, instead of closing the channel mid-run. +// +// 2. PipelineRuns can spawn additional TaskRuns over time (sequential +// `runAfter` tasks, finally blocks, retries). A single snapshot +// would silently drop any TaskRun that appears after the snapshot, +// even though the workflow is still running. +// +// The loop terminates only when the PipelineRun is terminal AND the +// most recent listing produced no new TaskRuns, which together mean no +// further TaskRuns will ever appear. +func (p *tektonProvider) streamPipelineRunLogs( + ctx context.Context, + out chan<- LogLine, + ref TektonRunRef, +) { + const pollInterval = 1 * time.Second + seen := map[string]bool{} + stepID := 0 + for { + if err := ctx.Err(); err != nil { 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, + taskRuns, err := p.taskRunsForPipelineRun(ctx, ref) + if err != nil { + p.log.Debug("Logs: list TaskRuns failed", + "err", err, "pipeline_run", ref.PipelineRunName, + ) + } + // Snapshot the terminal state *after* the listing so we never + // observe terminal=true while still missing a TaskRun that + // existed at list time. The reverse race (terminal observed + // before a TaskRun spawn) is handled by the next iteration: we + // only exit when terminal is true AND no new TaskRuns showed up + // in this pass. + terminal := p.isPipelineRunTerminal(ctx, ref) + + var fresh []k8s.Object + for _, tr := range taskRuns { + name := tr.GetName() + if name == "" || seen[name] { + continue + } + seen[name] = true + fresh = append(fresh, tr) + } + p.log.Debug("Logs: poll iteration", + "pipeline_run", ref.PipelineRunName, + "task_runs_total", len(taskRuns), + "task_runs_new", len(fresh), + "terminal", terminal, ) - stepID := 0 - for _, tr := range taskRuns { + for _, tr := range fresh { taskName := tr.GetName() if taskName == "" { taskName = fmt.Sprintf("task %d", stepID) } - p.log.Debug("Logs: streaming TaskRun", "task_run", taskName, "step_id", stepID, "terminal", terminal) + p.log.Debug("Logs: streaming TaskRun", + "task_run", taskName, "step_id", stepID, "terminal", terminal, + ) if !sendLine(ctx, out, LogLine{ Kind: LogKindControl, Time: time.Now(), @@ -521,10 +573,16 @@ func (p *tektonProvider) Logs( return } + // terminal is the snapshot taken at the top of this + // iteration. If the pipeline was terminal then, all + // TaskRuns we see are guaranteed complete and we can + // take the snapshot fast-path; otherwise we follow the + // pod logs live so StepStatusEnd is only emitted once + // the container actually exits. if terminal { - p.fetchCompletedTaskRunLogs(ctx, out, *ref, tr, stepID) + p.fetchCompletedTaskRunLogs(ctx, out, ref, tr, stepID) } else { - p.streamTaskRunLogs(ctx, out, *ref, tr, stepID) + p.streamTaskRunLogs(ctx, out, ref, tr, stepID) } if !sendLine(ctx, out, LogLine{ @@ -536,48 +594,35 @@ func (p *tektonProvider) Logs( }) { return } - p.log.Debug("Logs: finished TaskRun", "task_run", taskName, "step_id", stepID) + p.log.Debug("Logs: finished TaskRun", + "task_run", taskName, "step_id", stepID, + ) stepID++ } - p.log.Debug("Logs: all TaskRuns streamed", "pipeline_run", ref.PipelineRunName) - }() - 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", + // Done condition: the PipelineRun is terminal AND we found no + // new TaskRuns this iteration. Both are required because a + // terminal PipelineRun can still expose a freshly-listed + // TaskRun whose pod we haven't drained yet. + if terminal && len(fresh) == 0 { + p.log.Debug("Logs: pipeline run terminal, no new TaskRuns", "pipeline_run", ref.PipelineRunName, ) - return nil + return + } + + // When we processed new TaskRuns this round, loop again + // immediately to re-list — Follow=true streaming may have + // blocked us long enough for additional TaskRuns or terminal + // transitions to have happened. + if len(fresh) > 0 { + continue } + select { case <-ctx.Done(): - return nil - case <-time.After(interval): + return + case <-time.After(pollInterval): } } } @@ -661,7 +706,12 @@ func (p *tektonProvider) fetchCompletedTaskRunLogs( p.log.Debug("fetchCompletedTaskRunLogs: reading container logs", "pod", pod.Name, "container", c.Name, ) - rc, err := p.client.StreamPodLogs(ctx, ref.Namespace, pod.Name, c.Name) + // Snapshot read (Follow=false) is correct here: the TaskRun + // has already terminated, so the API server has the full log + // available and there is nothing more to wait for. + rc, err := p.client.StreamPodLogs(ctx, ref.Namespace, pod.Name, c.Name, + k8s.LogOptions{Follow: false}, + ) if err != nil { p.log.Debug("fetchCompletedTaskRunLogs: stream failed", "err", err, "pod", pod.Name, "container", c.Name, @@ -720,7 +770,15 @@ func (p *tektonProvider) streamTaskRunLogs( p.log.Debug("streamTaskRunLogs: streaming container", "pod", pod.Name, "container", c.Name, "step_id", stepID, ) - rc, err := p.client.StreamPodLogs(ctx, ref.Namespace, pod.Name, c.Name) + // Live tail with Follow=true: the apiserver holds the + // connection open until the container terminates (or ctx is + // cancelled), so sendReaderLines only returns once the + // container is actually done. Without this, the read EOFs at + // the current tail position and the caller emits a spurious + // StepStatusEnd while the step is still running. + rc, err := p.client.StreamPodLogs(ctx, ref.Namespace, pod.Name, c.Name, + k8s.LogOptions{Follow: true}, + ) if err != nil { p.log.Debug("streamTaskRunLogs: stream pod logs failed", "err", err, "pod", pod.Name, "container", c.Name, diff --git a/provider_tekton_test.go b/provider_tekton_test.go index 5201666..580d0e5 100644 --- a/provider_tekton_test.go +++ b/provider_tekton_test.go @@ -236,6 +236,25 @@ func TestTektonLogsLookup(t *testing.T) { } } + // Seed a terminal PipelineRun so Logs takes the snapshot path + // (fetchCompletedTaskRunLogs, Follow=false) and closes the + // channel after draining the seeded TaskRun. The non-terminal + // path follows pod logs live and would only EOF on ctx cancel, + // which is not what this assertion is exercising. + client.SeedObject(pipelineRunsGVR, "ci", k8s.Object{ + "apiVersion": "tekton.dev/v1", + "kind": "PipelineRun", + "metadata": map[string]any{ + "name": "run-1", + "namespace": "ci", + }, + "status": map[string]any{ + "conditions": []any{map[string]any{ + "type": "Succeeded", + "status": "True", + }}, + }, + }) client.SeedObject(taskRunsGVR, "ci", k8s.Object{ "apiVersion": "tekton.dev/v1", "kind": "TaskRun",