diff --git a/internal/executor/agent.go b/internal/executor/agent.go index b410710..18f4590 100644 --- a/internal/executor/agent.go +++ b/internal/executor/agent.go @@ -23,8 +23,16 @@ type AgentExecutor struct { Cfg model.ExecutorConfig LookPath func(string) (string, error) Runner CommandRunner + // PipeGrace bounds how long Run waits for output after the process has + // exited, in case a descendant inherited the pipe and lingers. Defaults + // to agentPipeGrace; tests shorten it. + PipeGrace time.Duration } +// agentPipeGrace is how long Run keeps reading output after the process exits, +// in case a descendant inherited the output pipe and lingers. +const agentPipeGrace = 10 * time.Second + func NewAgentExecutor(cfg model.ExecutorConfig) *AgentExecutor { return &AgentExecutor{Cfg: cfg.Normalize()} } @@ -73,17 +81,31 @@ func (e *AgentExecutor) Run(ctx context.Context, req CodingRequest, emit Emit) ( emit("agent", fmt.Sprintf("Starting %s in %s", resolved.DisplayName, req.RepoName), resolved.Command+" "+strings.Join(args, " ")) cmd := e.startCmd(runCtx, resolved.Command, args, req.RepoPath, stdin) - stdout, err := cmd.StdoutPipe() + // Own the output pipes instead of using StdoutPipe/StderrPipe: Cmd.Wait + // closes those parent ends as soon as the process exits, which races the + // scanners below and truncates output they have not read yet. + stdoutR, stdoutW, err := os.Pipe() if err != nil { return Result{}, err } - stderr, err := cmd.StderrPipe() + stderrR, stderrW, err := os.Pipe() if err != nil { + stdoutR.Close() + stdoutW.Close() return Result{}, err } + cmd.Stdout = stdoutW + cmd.Stderr = stderrW if err := cmd.Start(); err != nil { + stdoutR.Close() + stdoutW.Close() + stderrR.Close() + stderrW.Close() return Result{}, fmt.Errorf("start %s: %w", resolved.Command, err) } + // Drop the parent's write ends so the scanners can observe EOF. + stdoutW.Close() + stderrW.Close() var stderrBuf bytes.Buffer doneOut := make(chan struct{}) @@ -91,11 +113,11 @@ func (e *AgentExecutor) Run(ctx context.Context, req CodingRequest, emit Emit) ( sessionCh := make(chan string, 1) go func() { defer close(doneOut) - sessionCh <- scanAgentOutput(stdout, emit, resolved.DisplayName) + sessionCh <- scanAgentOutput(stdoutR, emit, resolved.DisplayName) }() go func() { defer close(doneErr) - sc := bufio.NewScanner(io.TeeReader(stderr, &stderrBuf)) + sc := bufio.NewScanner(io.TeeReader(stderrR, &stderrBuf)) sc.Buffer(make([]byte, 0, 64*1024), 1024*1024) for sc.Scan() { line := sc.Text() @@ -105,9 +127,41 @@ func (e *AgentExecutor) Run(ctx context.Context, req CodingRequest, emit Emit) ( } }() - waitErr := cmd.Wait() - <-doneOut - <-doneErr + // Wait for the process and for the scanners to drain. These pipes are + // ours, so Cmd.Wait does not close them; the scanners see EOF only once + // every writer — including a descendant that inherited the pipe — is + // gone. A lingering descendant is bounded by PipeGrace. + waitCh := make(chan error, 1) + go func() { waitCh <- cmd.Wait() }() + drained := make(chan struct{}) + go func() { + <-doneOut + <-doneErr + close(drained) + }() + + grace := e.PipeGrace + if grace <= 0 { + grace = agentPipeGrace + } + var waitErr error + select { + case waitErr = <-waitCh: + select { + case <-drained: + case <-time.After(grace): + // A descendant still holds the pipes open; keep the output read + // so far and stop waiting (bounded instead of hanging forever). + stdoutR.Close() + stderrR.Close() + <-drained + emit("agent", fmt.Sprintf("%s output truncated: a descendant still holds the output pipe after %s", resolved.DisplayName, grace), "") + } + case <-drained: + waitErr = <-waitCh + } + stdoutR.Close() + stderrR.Close() sessionID := <-sessionCh if runCtx.Err() == context.DeadlineExceeded { diff --git a/internal/executor/executor_test.go b/internal/executor/executor_test.go index d6e0d84..8f3d5d5 100644 --- a/internal/executor/executor_test.go +++ b/internal/executor/executor_test.go @@ -10,6 +10,7 @@ import ( "path/filepath" "strings" "testing" + "time" "github.com/ymhhh/vibecoding/internal/llm" "github.com/ymhhh/vibecoding/internal/model" @@ -324,6 +325,97 @@ func TestAgentExecutorFakeCommand(t *testing.T) { } } +func TestAgentExecutorDrainsAfterParentExit(t *testing.T) { + t.Parallel() + dir := t.TempDir() + script := filepath.Join(dir, "lingering.sh") + // The parent exits at once; a grandchild inherits stdout, emits more + // output, and exits 300ms later. Run must capture both. + body := "#!/bin/sh\n" + + "echo '{\"type\":\"system\",\"session_id\":\"sess-lingering\"}'\n" + + "sh -c 'sleep 0.3; echo late-grandchild' &\n" + + "exit 0\n" + if err := os.WriteFile(script, []byte(body), 0o755); err != nil { + t.Fatal(err) + } + ex := NewAgentExecutor(model.ExecutorConfig{ + Type: "agent", Preset: "custom", Command: script, Args: []string{"{prompt}"}, TimeoutSec: 30, + }) + ex.Runner = func(ctx context.Context, name string, args []string, cwd string, stdin io.Reader) *exec.Cmd { + cmd := exec.CommandContext(ctx, name, args...) + cmd.Dir = cwd + return cmd + } + var logs []string + emit := func(phase, msg, details string) { + logs = append(logs, phase+":"+msg) + } + res, err := ex.Run(context.Background(), CodingRequest{ + RepoPath: dir, + RepoName: "demo", + Title: "t", + Spec: &model.DevSpec{RawMarkdown: "do it"}, + }, emit) + if err != nil { + t.Fatal(err) + } + if res.SessionID != "sess-lingering" { + t.Fatalf("session=%q", res.SessionID) + } + joined := strings.Join(logs, "\n") + if !strings.Contains(joined, "late-grandchild") { + t.Fatalf("missing grandchild output; logs=%v", logs) + } +} + +func TestAgentExecutorPipeGraceBoundsLingeringWriters(t *testing.T) { + t.Parallel() + dir := t.TempDir() + script := filepath.Join(dir, "linger-forever.sh") + // The parent exits at once; the grandchild keeps the pipe open far past + // the short grace. Run must return promptly, keep the session line, and + // note the truncation. + body := "#!/bin/sh\n" + + "echo '{\"type\":\"system\",\"session_id\":\"sess-grace\"}'\n" + + "sh -c 'sleep 30' &\n" + + "exit 0\n" + if err := os.WriteFile(script, []byte(body), 0o755); err != nil { + t.Fatal(err) + } + ex := NewAgentExecutor(model.ExecutorConfig{ + Type: "agent", Preset: "custom", Command: script, Args: []string{"{prompt}"}, TimeoutSec: 30, + }) + ex.PipeGrace = 200 * time.Millisecond + ex.Runner = func(ctx context.Context, name string, args []string, cwd string, stdin io.Reader) *exec.Cmd { + cmd := exec.CommandContext(ctx, name, args...) + cmd.Dir = cwd + return cmd + } + var logs []string + emit := func(phase, msg, details string) { + logs = append(logs, phase+":"+msg) + } + start := time.Now() + res, err := ex.Run(context.Background(), CodingRequest{ + RepoPath: dir, + RepoName: "demo", + Title: "t", + Spec: &model.DevSpec{RawMarkdown: "do it"}, + }, emit) + if err != nil { + t.Fatal(err) + } + if elapsed := time.Since(start); elapsed > 5*time.Second { + t.Fatalf("Run blocked on a lingering descendant for %s", elapsed) + } + if res.SessionID != "sess-grace" { + t.Fatalf("session=%q", res.SessionID) + } + if !strings.Contains(strings.Join(logs, "\n"), "truncated") { + t.Fatalf("missing truncation note; logs=%v", logs) + } +} + func TestAgentExecutorStdinPrompt(t *testing.T) { t.Parallel() dir := t.TempDir()