-
Notifications
You must be signed in to change notification settings - Fork 73
/
runner.go
131 lines (112 loc) · 2.87 KB
/
runner.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
package k8s
import (
"bytes"
"io"
"os"
"os/exec"
"regexp"
"sync"
"syscall"
"github.com/deviceinsight/kafkactl/v5/internal/output"
"github.com/pkg/errors"
)
// Runner interface for shell commands
type Runner interface {
ExecuteAndReturn(cmd string, args []string) ([]byte, error)
Execute(cmd string, args []string) error
}
type ShellRunner struct {
Dir string
}
func (shell *ShellRunner) ExecuteAndReturn(binary string, args []string) ([]byte, error) {
cmd := exec.Command(binary, args...)
cmd.Dir = shell.Dir
cmd.Env = os.Environ()
if cmd.Stdout != nil {
return nil, errors.New("exec: Stdout already set")
}
if cmd.Stderr != nil {
return nil, errors.New("exec: Stderr already set")
}
var stdout bytes.Buffer
var stderr bytes.Buffer
var combined bytes.Buffer
cmd.Stdout = io.MultiWriter(&stdout, &combined)
cmd.Stderr = io.MultiWriter(&stderr, &combined)
err := cmd.Run()
if err != nil {
switch ee := err.(type) {
case *exec.ExitError:
// Propagate any non-zero exit status from the external command
waitStatus := ee.Sys().(syscall.WaitStatus)
exitStatus := waitStatus.ExitStatus()
err = newExitError(cmd.Path, cmd.Args, exitStatus, ee, stderr.String(), combined.String())
default:
output.Fail(errors.Wrap(err, "unexpected error"))
}
}
return stdout.Bytes(), err
}
// Execute a shell command
func (shell *ShellRunner) Execute(binary string, args []string) error {
cmd := exec.Command(binary, args...)
cmd.Dir = shell.Dir
cmd.Env = os.Environ()
// get stdOut of cmd
stdoutIn, _ := cmd.StdoutPipe()
// stdin, stderr directly mapped from outside
cmd.Stdin = os.Stdin
cmd.Stderr = os.Stderr
err := cmd.Start()
if err != nil {
output.Fail(err)
}
var wg sync.WaitGroup
wg.Add(1)
go func() {
if err := filterOutput(output.IoStreams.Out, stdoutIn); err != nil {
output.Fail(errors.Wrap(err, "unable to write std out"))
}
wg.Done()
}()
wg.Wait()
err = cmd.Wait()
if err != nil {
switch ee := err.(type) {
case *exec.ExitError:
// Propagate any non-zero exit status from the external command
waitStatus := ee.Sys().(syscall.WaitStatus)
exitStatus := waitStatus.ExitStatus()
err = newExitError(cmd.Path, cmd.Args, exitStatus, ee, "", "")
default:
output.Fail(errors.Wrap(err, "unexpected error"))
}
}
return err
}
func filterOutput(w io.Writer, r io.Reader) error {
buf := make([]byte, 1024)
for {
n, err := r.Read(buf[:])
if n > 0 {
data := buf[:n]
data = filterData(data)
_, err := w.Write(data)
if err != nil {
return err
}
}
if err != nil {
// Read returns io.EOF at the end of file, which is not an error for us
if err == io.EOF {
err = nil
}
return err
}
}
}
func filterData(data []byte) []byte {
// filter unwanted stuff from kubectl output
var podDeletedRegex = regexp.MustCompile(`pod "kafkactl-.+" deleted\n`)
return podDeletedRegex.ReplaceAll(data, []byte(""))
}