/
container_runner.go
123 lines (105 loc) · 2.61 KB
/
container_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
package provider
import (
"context"
"io"
"strings"
"github.com/virtual-kubelet/virtual-kubelet/node/api"
"golang.org/x/crypto/ssh"
"golang.org/x/sync/errgroup"
)
type ContainerRunner struct {
client *ssh.Client
attach api.AttachIO
}
// NewContainerRunner creates a new ContainerRunner object.
func NewContainerRunner(sshClient *ssh.Client, attachIO api.AttachIO) *ContainerRunner {
return &ContainerRunner{
attach: attachIO,
client: sshClient,
}
}
// Exec executes a command on a remote server via SSH protocol.
func (cr *ContainerRunner) Exec(ctx context.Context, cmd []string) error {
session, err := cr.client.NewSession()
if err != nil {
return err
}
defer session.Close()
if cr.attach.TTY() {
// Set up terminal modes
modes := ssh.TerminalModes{
ssh.ECHO: 0,
ssh.ECHOCTL: 0,
}
err = session.RequestPty("Xterm", 120, 60, modes)
if err != nil {
return err
}
}
sessionStdinPipe, err := session.StdinPipe()
if err != nil {
return err
}
defer sessionStdinPipe.Close()
sessionStdoutPipe, err := session.StdoutPipe()
if err != nil {
return err
}
sessionStderrPipe, err := session.StderrPipe()
if err != nil {
return err
}
ctx, cancel := context.WithCancel(ctx)
g, ctx := errgroup.WithContext(ctx)
// Goroutine responsible for listening to the 'resize' channel, updating the session
// window measurements, and handling session closure if the context is canceled
// or the terminal window is closed by the user.
g.Go(func() error {
for {
select {
case size := <-cr.attach.Resize():
// If the height and width are both 0, it likely indicates that the terminal
// window has been closed by the user. In this case, the SSH session is closed
// to terminate any running commands.
if size.Height == 0 && size.Width == 0 {
session.Close()
return nil
}
session.WindowChange(int(size.Height), int(size.Width))
case <-ctx.Done():
// If the context is canceled, the SSH session is closed to terminate any
// running commands.
session.Close()
return ctx.Err()
}
}
})
if aout := cr.attach.Stdout(); aout != nil {
defer aout.Close()
g.Go(func() error {
io.Copy(aout, sessionStdoutPipe)
return nil
})
}
if aerr := cr.attach.Stderr(); aerr != nil {
defer aerr.Close()
g.Go(func() error {
io.Copy(aerr, sessionStderrPipe)
return nil
})
}
if ain := cr.attach.Stdin(); ain != nil {
g.Go(func() error {
io.Copy(sessionStdinPipe, ain)
return nil
})
}
// sending the command
g.Go(func() error {
err := session.Run(strings.Join(cmd, " ") + "\n")
cancel()
return err
})
err = g.Wait()
return err
}