diff --git a/pkg/collect/copy.go b/pkg/collect/copy.go index 541e3181..df69fc1c 100644 --- a/pkg/collect/copy.go +++ b/pkg/collect/copy.go @@ -11,6 +11,7 @@ import ( "github.com/pkg/errors" troubleshootv1beta2 "github.com/replicatedhq/troubleshoot/pkg/apis/troubleshoot/v1beta2" + "github.com/replicatedhq/troubleshoot/pkg/k8sutil" corev1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/runtime" "k8s.io/client-go/kubernetes" @@ -102,14 +103,14 @@ func copyFilesFromPod(ctx context.Context, dstPath string, clientConfig *restcli Command: command, Container: containerName, Stdin: true, - Stdout: false, + Stdout: true, Stderr: true, TTY: false, }, parameterCodec) - exec, err := remotecommand.NewSPDYExecutor(clientConfig, "POST", req.URL()) + exec, err := k8sutil.NewFallbackExecutor(clientConfig, req.URL()) if err != nil { - return nil, nil, errors.Wrap(err, "failed to create SPDY executor") + return nil, nil, errors.Wrap(err, "failed to create executor") } result := NewResult() diff --git a/pkg/collect/copy_from_host.go b/pkg/collect/copy_from_host.go index 0a9ea574..1d1e87ec 100644 --- a/pkg/collect/copy_from_host.go +++ b/pkg/collect/copy_from_host.go @@ -304,14 +304,14 @@ func copyFilesFromHost(ctx context.Context, dstPath string, clientConfig *restcl Command: command, Container: containerName, Stdin: true, - Stdout: false, + Stdout: true, Stderr: true, TTY: false, }, parameterCodec) - exec, err := remotecommand.NewSPDYExecutor(clientConfig, "POST", req.URL()) + exec, err := k8sutil.NewFallbackExecutor(clientConfig, req.URL()) if err != nil { - return nil, nil, errors.Wrap(err, "failed to create SPDY executor") + return nil, nil, errors.Wrap(err, "failed to create executor") } result := NewResult() diff --git a/pkg/collect/etcd.go b/pkg/collect/etcd.go index d73f15de..6ae9f47b 100644 --- a/pkg/collect/etcd.go +++ b/pkg/collect/etcd.go @@ -10,6 +10,7 @@ import ( "github.com/pkg/errors" troubleshootv1beta2 "github.com/replicatedhq/troubleshoot/pkg/apis/troubleshoot/v1beta2" + "github.com/replicatedhq/troubleshoot/pkg/k8sutil" corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/client-go/kubernetes" @@ -325,7 +326,7 @@ func (c *etcdDebug) executeCommand(command string) ([]byte, []byte, error) { TTY: false, }, scheme.ParameterCodec) - exec, err := remotecommand.NewSPDYExecutor(c.clientConfig, "POST", req.URL()) + exec, err := k8sutil.NewFallbackExecutor(c.clientConfig, req.URL()) if err != nil { return nil, nil, err } diff --git a/pkg/collect/exec.go b/pkg/collect/exec.go index 1aaeb853..d5a2cdc4 100644 --- a/pkg/collect/exec.go +++ b/pkg/collect/exec.go @@ -9,6 +9,7 @@ import ( "github.com/pkg/errors" troubleshootv1beta2 "github.com/replicatedhq/troubleshoot/pkg/apis/troubleshoot/v1beta2" + "github.com/replicatedhq/troubleshoot/pkg/k8sutil" corev1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/runtime" "k8s.io/client-go/kubernetes" @@ -132,12 +133,12 @@ func getExecOutputs( Command: append(execCollector.Command, execCollector.Args...), Container: container, Stdin: true, - Stdout: false, + Stdout: true, Stderr: true, TTY: false, }, parameterCodec) - exec, err := remotecommand.NewSPDYExecutor(clientConfig, "POST", req.URL()) + exec, err := k8sutil.NewFallbackExecutor(clientConfig, req.URL()) if err != nil { return nil, nil, []string{err.Error()} } diff --git a/pkg/collect/longhorn.go b/pkg/collect/longhorn.go index 0a00c5c3..a2a6cbe2 100644 --- a/pkg/collect/longhorn.go +++ b/pkg/collect/longhorn.go @@ -12,6 +12,7 @@ import ( "github.com/pkg/errors" troubleshootv1beta2 "github.com/replicatedhq/troubleshoot/pkg/apis/troubleshoot/v1beta2" + "github.com/replicatedhq/troubleshoot/pkg/k8sutil" longhornv1beta1types "github.com/replicatedhq/troubleshoot/pkg/longhorn/apis/longhorn/v1beta1" longhornv1beta1 "github.com/replicatedhq/troubleshoot/pkg/longhorn/client/clientset/versioned/typed/longhorn/v1beta1" longhorntypes "github.com/replicatedhq/troubleshoot/pkg/longhorn/types" @@ -391,7 +392,7 @@ func GetLonghornReplicaChecksum(clientConfig *rest.Config, replica longhornv1bet Param("command", "-c"). Param("command", fmt.Sprintf("if [ -d %s ]; then md5sum %s/*; fi", dir, dir)) - executor, err := remotecommand.NewSPDYExecutor(clientConfig, "POST", req.URL()) + executor, err := k8sutil.NewFallbackExecutor(clientConfig, req.URL()) if err != nil { return "", errors.Wrapf(err, "create remote exec") } diff --git a/pkg/collect/sonobuoy_results.go b/pkg/collect/sonobuoy_results.go index deb763b6..f31550fe 100644 --- a/pkg/collect/sonobuoy_results.go +++ b/pkg/collect/sonobuoy_results.go @@ -12,6 +12,7 @@ import ( "github.com/pkg/errors" troubleshootv1beta2 "github.com/replicatedhq/troubleshoot/pkg/apis/troubleshoot/v1beta2" "github.com/replicatedhq/troubleshoot/pkg/client/troubleshootclientset/scheme" + "github.com/replicatedhq/troubleshoot/pkg/k8sutil" corev1 "k8s.io/api/core/v1" kerrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -147,7 +148,7 @@ func sonobuoyRetrieveResults( Stdout: true, Stderr: false, }, scheme.ParameterCodec) - executor, err := remotecommand.NewSPDYExecutor(restConfig, "POST", req.URL()) + executor, err := k8sutil.NewFallbackExecutor(restConfig, req.URL()) if err != nil { return nil, ec, err } diff --git a/pkg/k8sutil/exec.go b/pkg/k8sutil/exec.go new file mode 100644 index 00000000..79344a87 --- /dev/null +++ b/pkg/k8sutil/exec.go @@ -0,0 +1,28 @@ +package k8sutil + +import ( + "net/url" + + "k8s.io/apimachinery/pkg/util/httpstream" + restclient "k8s.io/client-go/rest" + "k8s.io/client-go/tools/remotecommand" +) + +// NewFallbackExecutor creates an executor that tries WebSocket first and falls +// back to SPDY if the server does not support it. Use this in place of +// remotecommand.NewSPDYExecutor everywhere. +func NewFallbackExecutor(config *restclient.Config, u *url.URL) (remotecommand.Executor, error) { + // WebSocket upgrade requires GET per RFC 6455; SPDY uses POST. + wsExec, err := remotecommand.NewWebSocketExecutor(config, "GET", u.String()) + if err != nil { + return nil, err + } + spdyExec, err := remotecommand.NewSPDYExecutor(config, "POST", u) + if err != nil { + return nil, err + } + shouldFallback := func(err error) bool { + return httpstream.IsUpgradeFailure(err) || httpstream.IsHTTPSProxyError(err) + } + return remotecommand.NewFallbackExecutor(wsExec, spdyExec, shouldFallback) +} diff --git a/pkg/k8sutil/exec_test.go b/pkg/k8sutil/exec_test.go new file mode 100644 index 00000000..34d9eb2e --- /dev/null +++ b/pkg/k8sutil/exec_test.go @@ -0,0 +1,19 @@ +package k8sutil + +import ( + "net/url" + "testing" + + "github.com/stretchr/testify/require" + restclient "k8s.io/client-go/rest" +) + +func TestNewFallbackExecutor(t *testing.T) { + config := &restclient.Config{Host: "http://localhost:8080"} + u, err := url.Parse("http://localhost:8080/api/v1/namespaces/default/pods/foo/exec") + require.NoError(t, err) + + exec, err := NewFallbackExecutor(config, u) + require.NoError(t, err) + require.NotNil(t, exec) +} diff --git a/pkg/k8sutil/portforward.go b/pkg/k8sutil/portforward.go deleted file mode 100644 index bbab72f9..00000000 --- a/pkg/k8sutil/portforward.go +++ /dev/null @@ -1,72 +0,0 @@ -package k8sutil - -import ( - "bytes" - "fmt" - "net/http" - "net/url" - "strings" - "time" - - restclient "k8s.io/client-go/rest" - "k8s.io/client-go/tools/portforward" - "k8s.io/client-go/transport/spdy" -) - -func PortForward(config *restclient.Config, localPort int, remotePort int, namespace string, podName string) (chan struct{}, error) { - roundTripper, upgrader, err := spdy.RoundTripperFor(config) - if err != nil { - return nil, err - } - - path := fmt.Sprintf("/api/v1/namespaces/%s/pods/%s/portforward", namespace, podName) - hostIP := strings.TrimLeft(config.Host, "htps:/") - serverURL := url.URL{Scheme: "http", Path: path, Host: hostIP} - dialer := spdy.NewDialer(upgrader, &http.Client{Transport: roundTripper}, http.MethodPost, &serverURL) - - stopChan, readyChan := make(chan struct{}, 1), make(chan struct{}, 1) - out, errOut := new(bytes.Buffer), new(bytes.Buffer) - - forwarder, err := portforward.New(dialer, []string{fmt.Sprintf("%d:%d", localPort, remotePort)}, stopChan, readyChan, out, errOut) - if err != nil { - return nil, err - } - - go func() { - for range readyChan { // Kubernetes will close this channel when it has something to tell us. - } - if errOut.String() != "" { - panic(errOut.String()) - } else if out.String() != "" { - // fmt.Println(out.String()) - } - }() - - go func() error { - if err = forwarder.ForwardPorts(); err != nil { // Locks until stopChan is closed. - panic(err) - } - - return nil - }() - - // Block until the new service is responding, limited to (math) seconds - quickClient := &http.Client{ - Timeout: time.Millisecond * 200, - } - - start := time.Now() - for { - response, err := quickClient.Get(fmt.Sprintf("http://localhost:%d", localPort)) - if err == nil && response.StatusCode == http.StatusOK { - break - } - if time.Now().Sub(start) > time.Duration(time.Second*5) { - return nil, err - } - - time.Sleep(time.Millisecond * 100) - } - - return stopChan, nil -} diff --git a/pkg/supportbundle/collect.go b/pkg/supportbundle/collect.go index 0867b348..e2f55d71 100644 --- a/pkg/supportbundle/collect.go +++ b/pkg/supportbundle/collect.go @@ -19,6 +19,7 @@ import ( "github.com/replicatedhq/troubleshoot/pkg/collect" "github.com/replicatedhq/troubleshoot/pkg/constants" "github.com/replicatedhq/troubleshoot/pkg/convert" + "github.com/replicatedhq/troubleshoot/pkg/k8sutil" "github.com/replicatedhq/troubleshoot/pkg/redact" "github.com/replicatedhq/troubleshoot/pkg/version" "go.opentelemetry.io/otel" @@ -367,7 +368,7 @@ func getExecOutputs( TTY: false, }, parameterCodec) - exec, err := remotecommand.NewSPDYExecutor(clientConfig, "POST", req.URL()) + exec, err := k8sutil.NewFallbackExecutor(clientConfig, req.URL()) if err != nil { return nil, nil, err }