This repository has been archived by the owner on Oct 29, 2021. It is now read-only.
/
k8s_forward.go.go
143 lines (121 loc) · 3.29 KB
/
k8s_forward.go.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
132
133
134
135
136
137
138
139
140
141
142
143
package kubetest
import (
"fmt"
"io/ioutil"
"net"
"net/http"
"github.com/pkg/errors"
v1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/util/httpstream"
spdy2 "k8s.io/apimachinery/pkg/util/httpstream/spdy"
"k8s.io/client-go/rest"
"k8s.io/client-go/tools/portforward"
"k8s.io/client-go/transport/spdy"
)
type PortForward struct {
// The initialized Kubernetes client.
k8s *K8s
pod *v1.Pod
// The port on the pod to forward traffic to.
DestinationPort int
// The port that the port forward should listen to, random if not set.
ListenPort int
stopChan chan struct{}
readyChan chan struct{}
pf *portforward.PortForwarder
}
func (k8s *K8s) NewPortForwarder(pod *v1.Pod, port int) (*PortForward, error) {
pf := &PortForward{
pod: pod,
k8s: k8s,
DestinationPort: port,
}
return pf, nil
}
// Start a port forward to a pod - blocks until the tunnel is ready for use.
func (p *PortForward) Start() error {
p.stopChan = make(chan struct{}, 1)
readyChan := make(chan struct{}, 1)
errChan := make(chan error, 1)
listenPort, err := p.getListenPort()
if err != nil {
return errors.Wrap(err, "Could not find a port to bind to")
}
dialer, err := p.dialer()
if err != nil {
return errors.Wrap(err, "Could not create a dialer")
}
ports := []string{
fmt.Sprintf("%d:%d", listenPort, p.DestinationPort),
}
discard := ioutil.Discard
pf, err := portforward.New(dialer, ports, p.stopChan, readyChan, discard, discard)
if err != nil {
return errors.Wrap(err, "Could not port forward into pod")
}
p.pf = pf
go func() {
errChan <- pf.ForwardPorts()
}()
select {
case err = <-errChan:
return errors.Wrap(err, "Could not create port forward")
case <-readyChan:
return nil
}
}
func (p *PortForward) Stop() {
p.stopChan <- struct{}{}
p.pf.Close()
}
func (p *PortForward) getListenPort() (int, error) {
var err error
if p.ListenPort == 0 {
p.ListenPort, err = p.getFreePort()
}
return p.ListenPort, err
}
func (p *PortForward) getFreePort() (int, error) {
listener, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
return 0, err
}
port := listener.Addr().(*net.TCPAddr).Port
err = listener.Close()
if err != nil {
return 0, err
}
return port, nil
}
// Create an httpstream.Dialer for use with portforward.New
func (p *PortForward) dialer() (httpstream.Dialer, error) {
podname := ""
if p.pod != nil {
podname = p.pod.Name
} else {
return nil, errors.New("Could not do POST request for a non-existing pod")
}
url := p.k8s.clientset.CoreV1().RESTClient().Post().
Resource("pods").
Namespace(p.k8s.GetK8sNamespace()).
Name(podname).
SubResource("portforward").URL()
transport, upgrader, err := roundTripperFor(p.k8s.config)
if err != nil {
return nil, errors.Wrap(err, "Could not create round tripper")
}
dialer := spdy.NewDialer(upgrader, &http.Client{Transport: transport}, "POST", url)
return dialer, nil
}
func roundTripperFor(config *rest.Config) (http.RoundTripper, httpstream.UpgradeRoundTripper, error) {
tlsConfig, err := rest.TLSConfigFor(config)
if err != nil {
return nil, nil, err
}
upgradeRoundTripper := spdy2.NewRoundTripper(tlsConfig, true, true)
wrapper, err := rest.HTTPWrappersForConfig(config, upgradeRoundTripper)
if err != nil {
return nil, nil, err
}
return wrapper, upgradeRoundTripper, nil
}