mirror of
https://github.com/jesseduffield/lazygit.git
synced 2026-09-28 02:07:09 -04:00
Extract the reusable parts of a command pipeline
PipeCommands runs a chain of commands to completion and gathers what they wrote to stderr. Rendering a diff through a renderer needs the same chain, but has to read the last command's output as it arrives, so it can't use PipeCommands as it stands. Pull out what both need: naming the chain for a log, wiring each command's output to the next one's input, and starting them all with the cleanup a failure to start requires. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5
parent
5fcbdadc67
commit
e13c09f961
@@ -4,12 +4,10 @@ import (
|
||||
"bytes"
|
||||
"io"
|
||||
"os"
|
||||
"os/exec"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
|
||||
"github.com/go-errors/errors"
|
||||
"github.com/samber/lo"
|
||||
|
||||
"github.com/atotto/clipboard"
|
||||
"github.com/jesseduffield/lazygit/pkg/common"
|
||||
@@ -203,61 +201,25 @@ func (c *OSCommand) FileExists(path string) (bool, error) {
|
||||
|
||||
// PipeCommands runs a heap of commands and pipes their inputs/outputs together like A | B | C
|
||||
func (c *OSCommand) PipeCommands(cmdObjs ...*CmdObj) error {
|
||||
cmds := lo.Map(cmdObjs, func(cmdObj *CmdObj, _ int) *exec.Cmd {
|
||||
return cmdObj.GetCmd()
|
||||
})
|
||||
c.LogCommand(pipelineString(cmdObjs), true)
|
||||
|
||||
logCmdStr := strings.Join(
|
||||
lo.Map(cmdObjs, func(cmdObj *CmdObj, _ int) string {
|
||||
return cmdObj.ToString()
|
||||
}),
|
||||
" | ",
|
||||
)
|
||||
|
||||
c.LogCommand(logCmdStr, true)
|
||||
|
||||
for i := range len(cmds) - 1 {
|
||||
stdout, err := cmds[i].StdoutPipe()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
cmds[i+1].Stdin = stdout
|
||||
cmds, err := wirePipeline(cmdObjs)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// keeping this here in case I adapt this code for some other purpose in the future
|
||||
// cmds[len(cmds)-1].Stdout = os.Stdout
|
||||
|
||||
stderrs := make([]bytes.Buffer, len(cmds))
|
||||
for i := range cmds {
|
||||
cmds[i].Stderr = &stderrs[i]
|
||||
}
|
||||
|
||||
// Start every command before waiting for any of them: waiting for a command
|
||||
// closes our end of the pipe that feeds the next one, and a command that
|
||||
// hasn't been started by then would inherit a closed stdin.
|
||||
started := 0
|
||||
var startErr error
|
||||
for _, cmd := range cmds {
|
||||
if err := cmd.Start(); err != nil {
|
||||
startErr = err
|
||||
break
|
||||
}
|
||||
|
||||
started++
|
||||
}
|
||||
started, startErr := startPipeline(cmds)
|
||||
|
||||
finalErrors := []string{}
|
||||
|
||||
if startErr != nil {
|
||||
c.Log.Error(startErr)
|
||||
finalErrors = append(finalErrors, startErr.Error())
|
||||
|
||||
// Without the rest of the pipeline to drain them, the commands we did
|
||||
// start could block forever writing to a full pipe.
|
||||
for _, cmd := range cmds[:started] {
|
||||
_ = cmd.Process.Kill()
|
||||
}
|
||||
}
|
||||
|
||||
for i, cmd := range cmds[:started] {
|
||||
|
||||
@@ -0,0 +1,60 @@
|
||||
package oscommands
|
||||
|
||||
import (
|
||||
"os/exec"
|
||||
"strings"
|
||||
|
||||
"github.com/samber/lo"
|
||||
)
|
||||
|
||||
// pipelineString names a chain of commands the way a shell would write it.
|
||||
func pipelineString(cmdObjs []*CmdObj) string {
|
||||
return strings.Join(
|
||||
lo.Map(cmdObjs, func(cmdObj *CmdObj, _ int) string {
|
||||
return cmdObj.ToString()
|
||||
}),
|
||||
" | ",
|
||||
)
|
||||
}
|
||||
|
||||
// wirePipeline connects each command's output to the next one's input, like
|
||||
// A | B | C. The last command's output is left for the caller to direct.
|
||||
func wirePipeline(cmdObjs []*CmdObj) ([]*exec.Cmd, error) {
|
||||
cmds := lo.Map(cmdObjs, func(cmdObj *CmdObj, _ int) *exec.Cmd {
|
||||
return cmdObj.GetCmd()
|
||||
})
|
||||
|
||||
for i := range len(cmds) - 1 {
|
||||
stdout, err := cmds[i].StdoutPipe()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
cmds[i+1].Stdin = stdout
|
||||
}
|
||||
|
||||
return cmds, nil
|
||||
}
|
||||
|
||||
// startPipeline starts every command and reports how many it got going. Every
|
||||
// one is started before any of them is waited for: waiting closes our end of
|
||||
// the pipe that feeds the next command, and one that hasn't been started by
|
||||
// then would inherit a closed stdin.
|
||||
//
|
||||
// When a command fails to start, the ones already running are killed, since
|
||||
// without the rest of the pipeline to drain them they could block forever
|
||||
// writing to a full pipe. They still have to be reaped, so the count covers
|
||||
// them too.
|
||||
func startPipeline(cmds []*exec.Cmd) (int, error) {
|
||||
for i, cmd := range cmds {
|
||||
if err := cmd.Start(); err != nil {
|
||||
for _, started := range cmds[:i] {
|
||||
_ = started.Process.Kill()
|
||||
}
|
||||
|
||||
return i, err
|
||||
}
|
||||
}
|
||||
|
||||
return len(cmds), nil
|
||||
}
|
||||
Reference in New Issue
Block a user