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

package execution

import (
"context"
"errors"
"strings"
"sync"
"testing"
"time"

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

// TestPreferSinkFailure pins the rule: the caller's failing sink is the reason the
// run stopped, so it leads the returned error even when a context error was recorded
// first, and nothing is dropped.
func TestPreferSinkFailure(t *testing.T) {
t.Parallel()

sinkErr := errors.New("sink stopped")
deadline := context.DeadlineExceeded

assert.Equal(t, deadline, preferSinkFailure(deadline, nil))
assert.Equal(t, sinkErr, preferSinkFailure(nil, sinkErr))
assert.Equal(t, sinkErr, preferSinkFailure(sinkErr, sinkErr))

joined := preferSinkFailure(deadline, sinkErr)
require.ErrorIs(t, joined, sinkErr)
require.ErrorIs(t, joined, deadline)
assert.Equal(t, sinkErr.Error(), strings.SplitN(joined.Error(), "\n", 2)[0], "the sink failure leads")
}

// TestRunnerPrefersTheSinkFailureOverCallerCancellation pins the same rule through
// the runner: a caller that cancels while its sink is still running, and a sink that
// then reports its own failure, still get their failure back as the reason the run
// stopped. Which of the two the runner observed first is a schedule, and the answer
// must not depend on it.
func TestRunnerPrefersTheSinkFailureOverCallerCancellation(t *testing.T) {
t.Parallel()

fixture := newExecutorFixture(t)
executor := fixture.executor(t, unavailableBackend{}, systemRunnerDependencies())
operation, authorization := fixture.fullAccessOperation(t, fixture.operationSpec("printf hello; /bin/sleep 10"))

ctx, cancel := context.WithCancel(t.Context())
defer cancel()

sinkErr := errors.New("sink stopped late")
started := make(chan struct{})

var once sync.Once

go func() {
<-started
cancel()
}()

result, err := executor.Execute(ctx, operation, authorization, SinkFunc(func(sinkCtx context.Context, _ OutputChunk) error {
once.Do(func() { close(started) })

// The cancellation reaches the runner first, then this failure lands while it
// is draining the dispatcher.
<-sinkCtx.Done()
time.Sleep(10 * time.Millisecond)

return sinkErr
}))

require.ErrorIs(t, err, sinkErr)
require.ErrorIs(t, err, context.Canceled, "the cancellation is still reported")
assert.Equal(t, StatusCanceled, result.Status)
}
49 changes: 47 additions & 2 deletions internal/coding/execution/runner_unix.go
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,8 @@ func (e *Executor) run(ctx context.Context, plan *Plan, sink Sink) (Result, erro
trigger := make(chan error, 1)
chunks := make(chan OutputChunk, plan.operation.output.QueueDepth)
readersDone := startOutputReaders(pipes, collector, plan.operation.output.ChunkBytes, chunks, trigger)
dispatchDone := startDispatcher(runCtx, sink, chunks, trigger)
failures := &outputFailure{}
dispatchDone := startDispatcher(runCtx, sink, chunks, trigger, failures)
stdinDone := startStdinWriter(pipes.stdinParent, plan.launch.stdin, trigger)

waitDone := make(chan error, 1)
Expand All @@ -66,6 +67,7 @@ func (e *Executor) run(ctx context.Context, plan *Plan, sink Sink) (Result, erro
cleanupErr := cleanupProcessGroup(command.Process.Pid, e.termGrace)
readerErr := e.finishIO(runCtx, ctx, pipes, readersDone, dispatchDone, stdinDone, chunks, trigger, &runErr)
closeErr = errors.Join(closeErr, readerErr, pipes.closeAll())
runErr = preferSinkFailure(runErr, failures.failure())

result.Duration = e.deps.now().Sub(startedAt)
result.Stdout, result.Stderr = collector.results()
Expand Down Expand Up @@ -179,6 +181,7 @@ func startDispatcher(
sink Sink,
chunks <-chan OutputChunk,
trigger chan<- error,
failures *outputFailure,
) <-chan struct{} {
done := make(chan struct{})
go func() {
Expand All @@ -190,7 +193,9 @@ func startDispatcher(
}

if err := sink.WriteOutput(ctx, chunk); err != nil {
notifyTrigger(trigger, fmt.Errorf("coding execution: output sink: %w", err))
failure := fmt.Errorf("coding execution: output sink: %w", err)
failures.record(failure)
notifyTrigger(trigger, failure)

return
}
Expand Down Expand Up @@ -423,6 +428,46 @@ func setFirstError(target *error, candidate error) {
}
}

// outputFailure remembers the sink's own failure. The dispatcher reports it on the
// trigger as well, but the first error to reach the runner wins that race, and the
// caller's failing sink is the more useful answer than an operation deadline that
// expired while the runner was draining.
type outputFailure struct {
mu sync.Mutex
err error
}

func (f *outputFailure) record(err error) {
f.mu.Lock()
defer f.mu.Unlock()

if f.err == nil {
f.err = err
}
}

func (f *outputFailure) failure() error {
f.mu.Lock()
defer f.mu.Unlock()

return f.err
}

// preferSinkFailure returns the sink's failure as the reason the run stopped, with
// anything else that happened joined behind it.
func preferSinkFailure(runErr, sinkErr error) error {
switch {
case sinkErr == nil:
return runErr
case runErr == nil:
return sinkErr
case errors.Is(runErr, sinkErr):
return runErr
default:
return errors.Join(sinkErr, runErr)
}
}

// processGroupSignalUnavailable reports whether a failed process-group signal
// means the group is gone or is no longer ours to signal. Cleanup runs after the
// direct child is reaped, so its pid can be recycled into an unrelated process
Expand Down
Loading