Skip to content
Merged
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
72 changes: 72 additions & 0 deletions internal/coding/execution/runner_drain_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
//go:build darwin || linux

package execution

import (
"context"
"os"
"testing"
"time"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

// TestFinishIODrainsBeforeClosingTheReadEnds pins the ordering behind a capture that
// described reader scheduling: when the output trigger fires, the read ends stay open
// until the readers report done, so a caller receives what the process wrote rather
// than whatever had been read when the trigger landed.
func TestFinishIODrainsBeforeClosingTheReadEnds(t *testing.T) {
t.Parallel()

stdoutRead, stdoutWrite, err := os.Pipe()
require.NoError(t, err)
stderrRead, stderrWrite, err := os.Pipe()
require.NoError(t, err)
t.Cleanup(func() {
for _, file := range []*os.File{stdoutRead, stdoutWrite, stderrRead, stderrWrite} {
_ = file.Close()
}
})

pipes := &processPipes{
stdoutParent: stdoutRead,
stdoutChild: stdoutWrite,
stderrParent: stderrRead,
stderrChild: stderrWrite,
}

trigger := make(chan error, 1)
trigger <- ErrOutputLimit

// The readers are still running. The grace is far beyond this test, so a timer
// cannot end the drain in its place.
readersDone := make(chan struct{})
ioDone := make(chan struct{})
close(ioDone)

var runErr error

finished := make(chan error, 1)
executor := &Executor{drainGrace: 10 * time.Second}

go func() {
finished <- executor.finishIO(
context.Background(), context.Background(), pipes,
readersDone, ioDone, ioDone, make(chan OutputChunk), trigger, &runErr,
)
}()

// Wait until the runner has taken the trigger, so the branch under test is the one
// running, then let the readers finish.
require.Eventually(t, func() bool { return len(trigger) == 0 },
30*time.Second, 10*time.Millisecond, "the runner did not take the trigger")
close(readersDone)
require.NoError(t, <-finished)
require.ErrorIs(t, runErr, ErrOutputLimit)

// After the drain the read end is still ours: writing succeeds while it is open and
// fails with a broken pipe when the trigger closed it early.
_, writeErr := stdoutWrite.Write([]byte("late"))
assert.NoError(t, writeErr, "the read end was closed before the readers drained")
}
14 changes: 12 additions & 2 deletions internal/coding/execution/runner_unix.go
Original file line number Diff line number Diff line change
Expand Up @@ -249,9 +249,19 @@ func (e *Executor) finishIO(
case err := <-trigger:
setFirstError(runErr, err)

_ = pipes.closeReadEnds()
// The trigger stopped the process in waitForProcess, so the readers reach EOF on
// their own: drain them before touching the pipes. Closing the read ends here
// made the retained capture and the byte count describe how fast the readers
// happened to be scheduled rather than what the process wrote, which is what the
// output-limit fixture hit on CI. The grace still bounds a child that holds the
// pipe open.
select {
case <-readersDone:
case <-drainTimer.C:
_ = pipes.closeReadEnds()

<-readersDone
<-readersDone
}
case <-runCtx.Done():
setFirstError(runErr, executionContextError(runCtx, callerCtx))

Expand Down
Loading