返回 DeepSeek-Reasonix
pipe.go
1 package persistentshell
2
3 import (
4 "os"
5 "os/exec"
6 "sync"
7
8 "reasonix/internal/proc"
9 )
10
11 type pipeProcess struct {
12 input, output *os.File
13 cmd *exec.Cmd
14 job uintptr
15 once sync.Once
16 }
17
18 // Both output streams share one OS pipe so command output and its trailing
19 // completion fence remain ordered. Track the whole process tree before the
20 // shell starts; cancellation must retire its children as well as the shell.
21 func startPipe(argv []string, dir string, env []string) (ptyConn, error) {
22 if len(argv) == 0 {
23 return nil, errEmptyArgv
24 }
25 inputReader, inputWriter, err := os.Pipe()
26 if err != nil {
27 return nil, err
28 }
29 outputReader, outputWriter, err := os.Pipe()
30 if err != nil {
31 _ = inputReader.Close()
32 _ = inputWriter.Close()
33 return nil, err
34 }
35 cmd := proc.Command(argv[0], argv[1:]...)
36 cmd.Dir, cmd.Env = dir, env
37 cmd.Stdin, cmd.Stdout, cmd.Stderr = inputReader, outputWriter, outputWriter
38 job, err := startShellTracked(cmd)
39 _ = inputReader.Close()
40 _ = outputWriter.Close()
41 if err != nil {
42 _ = inputWriter.Close()
43 _ = outputReader.Close()
44 return nil, err
45 }
46 return &pipeProcess{input: inputWriter, output: outputReader, cmd: cmd, job: job}, nil
47 }
48
49 func (p *pipeProcess) Read(b []byte) (int, error) { return p.output.Read(b) }
50 func (p *pipeProcess) Write(b []byte) (int, error) { return p.input.Write(b) }
51 func (p *pipeProcess) Close() error {
52 p.once.Do(func() {
53 proc.KillTracked(p.cmd, p.job)
54 _ = p.input.Close()
55 _ = p.output.Close()
56 _ = p.cmd.Wait()
57 })
58 return nil
59 }
60
60 lines GO