Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
4dc570c
fix(process): bound context command cleanup
PierrunoYT Aug 27, 2026
8621fdb
fix(process): retain command tree cancellation identity
ampagent Aug 27, 2026
f743018
fix(process): bind Windows jobs before execution
ampagent Aug 27, 2026
44aff3d
fix(execution): preserve successful Windows descendants
PierrunoYT Aug 28, 2026
e9af688
fix(execution): fall back when tree containment fails
PierrunoYT Aug 29, 2026
91c5895
fix(execution): require Windows job containment
ampagent Sep 2, 2026
803e076
fix(execution): retain process tree lifecycle ownership
PierrunoYT Sep 3, 2026
aac2bd9
fix(process): bound context command cleanup
PierrunoYT Aug 27, 2026
54b457f
fix(process): retain command tree cancellation identity
ampagent Aug 27, 2026
8e69a0b
fix(process): bind Windows jobs before execution
ampagent Aug 27, 2026
20b24c1
fix(execution): preserve successful Windows descendants
PierrunoYT Aug 28, 2026
c6e35f0
fix(execution): fall back when tree containment fails
PierrunoYT Aug 29, 2026
56a3503
fix(execution): require Windows job containment
ampagent Sep 2, 2026
53b3ead
fix(execution): retain process tree lifecycle ownership
PierrunoYT Sep 3, 2026
d6e8f02
fix(execution): bound configured hook process lifecycles
ampagent Sep 7, 2026
109c345
Merge existing PR ancestry after refreshing upstream base
ampagent Sep 7, 2026
6285a65
Merge upstream main into fix/issue-966-process-tree-timeouts
ampagent Sep 12, 2026
aa20bc1
test(plugins): allow race runtime exit delay in origin fixture
ampagent Sep 12, 2026
173ddc3
fix(process): preserve benchmark output drain failures
PierrunoYT Sep 13, 2026
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
8 changes: 4 additions & 4 deletions internal/agenteval/agent_command.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,9 +3,10 @@ package agenteval
import (
"bytes"
"context"
"errors"
"os/exec"
"strings"

"github.com/Gitlawb/zero/internal/execution"
)

type AgentRunInput struct {
Expand Down Expand Up @@ -67,7 +68,7 @@ func (runner CommandAgentRunner) Run(ctx context.Context, input AgentRunInput) A
cmd.Stdout = stdout
cmd.Stderr = stderr

err := cmd.Run()
err := execution.RunCommand(ctx, cmd)
result.Stdout = stdout.buf.String()
result.Stderr = stderr.buf.String()
result.Truncated = stdout.truncated || stderr.truncated
Expand All @@ -81,8 +82,7 @@ func (runner CommandAgentRunner) Run(ctx context.Context, input AgentRunInput) A
result.Error = ctxErr.Error()
return result
}
var exitErr *exec.ExitError
if errors.As(err, &exitErr) {
if exitErr, ok := execution.AsPureExitError(err); ok {
result.ExitCode = exitErr.ExitCode()
return result
}
Expand Down
4 changes: 3 additions & 1 deletion internal/agenteval/materialize.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,8 @@ import (
"os/exec"
"path/filepath"
"strings"

"github.com/Gitlawb/zero/internal/execution"
)

type Materializer struct{}
Expand Down Expand Up @@ -198,7 +200,7 @@ func initGitBaseline(ctx context.Context, workspace string) error {
var output bytes.Buffer
cmd.Stdout = &output
cmd.Stderr = &output
if err := cmd.Run(); err != nil {
if err := execution.RunCommand(ctx, cmd); err != nil {
if ctxErr := ctx.Err(); ctxErr != nil {
return ctxErr
}
Expand Down
13 changes: 8 additions & 5 deletions internal/agenteval/run.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,8 @@ import (
"path/filepath"
"strings"
"time"

"github.com/Gitlawb/zero/internal/execution"
)

// defaultCommandTimeout bounds a single verification command so a hung command
Expand Down Expand Up @@ -144,15 +146,14 @@ func execCommand(ctx context.Context, workspace string, command Command) Command
var stderr bytes.Buffer
cmd.Stdout = &stdout
cmd.Stderr = &stderr
err := cmd.Run()
err := execution.RunCommand(ctx, cmd)
result.Stdout = stdout.String()
result.Stderr = stderr.String()
if err == nil {
result.ExitCode = 0
return result
}
var exitErr *exec.ExitError
if errors.As(err, &exitErr) {
if exitErr, ok := execution.AsPureExitError(err); ok {
result.ExitCode = exitErr.ExitCode()
return result
}
Expand All @@ -167,14 +168,16 @@ func execCommand(ctx context.Context, workspace string, command Command) Command
func defaultRunGit(ctx context.Context, workspace string, args ...string) ([]byte, error) {
allArgs := append([]string{"-C", workspace}, args...)
cmd := exec.CommandContext(ctx, "git", allArgs...)
output, err := cmd.Output()
var output bytes.Buffer
cmd.Stdout = &output
err := execution.RunCommand(ctx, cmd)
if err != nil {
if ctxErr := ctx.Err(); ctxErr != nil {
return nil, ctxErr
}
return nil, err
}
return output, nil
return output.Bytes(), nil
}

func parseGitStatusPorcelain(output []byte) []string {
Expand Down
4 changes: 3 additions & 1 deletion internal/dictation/runner.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@ import (
"os"
"os/exec"
"time"

"github.com/Gitlawb/zero/internal/execution"
)

// commandSpec describes one capture-process invocation. Argv is always
Expand Down Expand Up @@ -107,7 +109,7 @@ func runCommandOutput(ctx context.Context, name string, args ...string) ([]byte,
var out bytes.Buffer
cmd.Stdout = &out
cmd.Stderr = &out
err := cmd.Run()
err := execution.RunCommand(ctx, cmd)
return out.Bytes(), err
}

Expand Down
152 changes: 152 additions & 0 deletions internal/execution/command_context.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,152 @@
package execution

import (
"context"
"errors"
"fmt"
"io"
"os"
"os/exec"
"reflect"
"sync"
)

// RunCommand runs a context-bound command in a retained process tree and
// prevents inherited output handles from blocking Wait indefinitely.
func RunCommand(ctx context.Context, command *exec.Cmd) (err error) {
if command == nil {
return errors.New("execution: nil command")
}
if ctx == nil {
ctx = context.Background()
}
if err := ctx.Err(); err != nil {
return err
}
drains := observeOutputDrains(command)
tree, err := prepareCommandTree(command)
if err != nil {
return err
}
defer func() { err = errors.Join(err, tree.close()) }()

command.WaitDelay = processWaitDelay
// Adapters may return exec.Command, which rejects a non-nil Cancel.
// The watcher below also owns cancellation for commands without that hook.
if command.Cancel != nil {
command.Cancel = tree.cancel
}
if err := command.Start(); err != nil {
_ = tree.attach(nil)
return err
}
if err := tree.attach(command.Process); err != nil {
killErr := command.Process.Kill()
waitErr := command.Wait()
return errors.Join(fmt.Errorf("execution: attach process tree: %w", err), killErr, waitErr)
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}
waitComplete := make(chan struct{})
type cancellation struct {
err error
canceled bool
}
cancelResult := make(chan cancellation, 1)
go func() {
select {
case <-ctx.Done():
cancelResult <- cancellation{err: tree.cancel(), canceled: true}
case <-waitComplete:
cancelResult <- cancellation{}
}
}()
waitErr := command.Wait()
close(waitComplete)
canceled := <-cancelResult
if canceled.canceled {
return errors.Join(waitErr, ctx.Err(), canceled.err)
}
if waitErr != nil && drains.err() != nil {
// Cmd.Wait intentionally prefers an ExitError over a copying error. When
// WaitDelay forcibly closes an inherited output pipe after a nonzero root
// exit, retain that cleanup failure so result consumers cannot reconcile it
// as an ordinary command exit.
return errors.Join(waitErr, exec.ErrWaitDelay, tree.cancel())
}
if waitErr != nil {
return errors.Join(waitErr, tree.cancel())
}
return waitErr
}

type outputDrains struct {
stdout *drainObserver
stderr *drainObserver
}

func observeOutputDrains(command *exec.Cmd) outputDrains {
drains := outputDrains{}
sharedOutput := sameWriter(command.Stderr, command.Stdout)
if !isFile(command.Stdout) && command.Stdout != nil {
drains.stdout = &drainObserver{writer: command.Stdout}
command.Stdout = drains.stdout
}
if sharedOutput && drains.stdout != nil {
command.Stderr = drains.stdout
drains.stderr = drains.stdout
} else if !isFile(command.Stderr) && command.Stderr != nil {
drains.stderr = &drainObserver{writer: command.Stderr}
command.Stderr = drains.stderr
}
return drains
}

func (drains outputDrains) err() error {
if drains.stdout != nil && drains.stdout.err() != nil {
return drains.stdout.err()
}
if drains.stderr != nil {
return drains.stderr.err()
}
return nil
}

func isFile(writer io.Writer) bool {
_, ok := writer.(*os.File)
return ok
}

func sameWriter(left io.Writer, right io.Writer) bool {
if left == nil || right == nil {
return false
}
leftType := reflect.TypeOf(left)
return leftType == reflect.TypeOf(right) && leftType.Comparable() && left == right
}

// drainObserver records an error from os/exec's output-copying goroutine.
// Cmd.Wait drops that error when process exit itself fails.
type drainObserver struct {
writer io.Writer
mu sync.Mutex
copyErr error
}

func (observer *drainObserver) Write(data []byte) (int, error) {
return observer.writer.Write(data)
}

func (observer *drainObserver) ReadFrom(reader io.Reader) (int64, error) {
written, err := io.Copy(observer.writer, reader)
if err != nil {
observer.mu.Lock()
observer.copyErr = err
observer.mu.Unlock()
}
return written, err
}

func (observer *drainObserver) err() error {
observer.mu.Lock()
defer observer.mu.Unlock()
return observer.copyErr
}
100 changes: 100 additions & 0 deletions internal/execution/command_context_process_unix_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,100 @@
//go:build !windows

package execution

import (
"errors"
"os"
"strconv"
"strings"
"syscall"
"testing"
"time"
)

type helperProcessOwner struct {
pidFile string
stopFile string
pid int
exited bool
}

func ownHelperProcess(t *testing.T, pidFile, stopFile string) *helperProcessOwner {
t.Helper()
owner := &helperProcessOwner{pidFile: pidFile, stopFile: stopFile}
t.Cleanup(func() { owner.cleanup(t) })
return owner
}

func (owner *helperProcessOwner) waitReady(t *testing.T, timeout time.Duration) int {
t.Helper()
deadline := time.Now().Add(timeout)
for {
data, err := os.ReadFile(owner.pidFile)
if err == nil {
pid, parseErr := strconv.Atoi(strings.TrimSpace(string(data)))
if parseErr == nil && pid > 0 {
owner.pid = pid
return pid
}
}
if time.Now().After(deadline) {
t.Fatalf("helper did not hand off a valid PID within %s", timeout)
}
time.Sleep(10 * time.Millisecond)
}
}

func (owner *helperProcessOwner) cleanup(t *testing.T) {
t.Helper()
if err := os.WriteFile(owner.stopFile, nil, 0o600); err != nil {
t.Errorf("request helper process stop: %v", err)
}
if owner.exited {
return
}
if owner.pid == 0 {
data, err := os.ReadFile(owner.pidFile)
if err == nil {
owner.pid, _ = strconv.Atoi(strings.TrimSpace(string(data)))
}
}
if owner.pid <= 0 {
return
}
deadline := time.Now().Add(2 * time.Second)
for time.Now().Before(deadline) {
err := syscall.Kill(owner.pid, syscall.Signal(0))
if errors.Is(err, syscall.ESRCH) {
owner.pid = 0
owner.exited = true
return
}
if err != nil {
t.Errorf("check helper process %d after cleanup: %v", owner.pid, err)
return
}
time.Sleep(10 * time.Millisecond)
}
t.Errorf("helper process %d survived cleanup", owner.pid)
}

func (owner *helperProcessOwner) awaitExit(t *testing.T) {
t.Helper()
deadline := time.Now().Add(2 * time.Second)
for {
err := syscall.Kill(owner.pid, syscall.Signal(0))
if errors.Is(err, syscall.ESRCH) {
owner.pid = 0
owner.exited = true
return
}
if err != nil {
t.Fatalf("check helper process %d: %v", owner.pid, err)
}
if time.Now().After(deadline) {
t.Fatalf("helper process %d is still running after command cancellation", owner.pid)
}
time.Sleep(10 * time.Millisecond)
}
}
Loading
Loading