Merge pull request #12 from hannut/fix/proxy-connection-upgrade-headers

Preserve connection-upgrade headers for kubectl streaming
This commit is contained in:
Philip Laine
2026-06-11 14:47:42 +02:00
committed by GitHub
2 changed files with 85 additions and 0 deletions
+31
View File
@@ -34,6 +34,13 @@ const (
AuthorizationHeader = "Authorization"
ImpersonateUserHeader = "Impersonate-User"
ImpersonateGroupHeader = "Impersonate-Group"
ConnectionHeader = "Connection"
UpgradeHeader = "Upgrade"
SecWebsocketKeyHeader = "Sec-Websocket-Key"
SecWebsocketVersionHeader = "Sec-Websocket-Version"
SecWebsocketProtocolHeader = "Sec-Websocket-Protocol"
SecWebsocketExtensionsHeader = "Sec-Websocket-Extensions"
)
type PeerLister interface {
@@ -81,6 +88,15 @@ func proxyHandler(peerLister PeerLister, kubeAPIServerURL *url.URL, certPool *x5
ContentLengthHeader: nil,
ContentTypeHeader: nil,
UserAgentHeader: nil,
// WebSocket negotiation headers for streaming subresources
// (kubectl exec/attach/port-forward/cp). These are not
// hop-by-hop, so httputil.ReverseProxy does not restore them;
// they must survive the allowlist for the upstream API server to
// complete the WebSocket handshake.
SecWebsocketKeyHeader: nil,
SecWebsocketVersionHeader: nil,
SecWebsocketProtocolHeader: nil,
SecWebsocketExtensionsHeader: nil,
}
for k := range pr.Out.Header {
if _, ok := allowedHeaders[k]; !ok {
@@ -88,6 +104,21 @@ func proxyHandler(peerLister PeerLister, kubeAPIServerURL *url.URL, certPool *x5
}
}
// Preserve the connection-upgrade handshake for streaming
// subresources. The allowlist above strips the hop-by-hop
// Connection/Upgrade headers; without them httputil.ReverseProxy
// treats the request as non-upgrade and the API server rejects it
// with "Upgrade request required" (breaking kubectl
// exec/attach/port-forward/cp over both WebSocket and SPDY).
// Reconstruct them from the inbound request rather than allowlisting
// the client-supplied Connection header, so a client cannot name
// proxy-set headers (Authorization, Impersonate-*) as hop-by-hop and
// have them stripped before reaching the API server.
if upgrade := pr.In.Header.Get(UpgradeHeader); upgrade != "" {
pr.Out.Header.Set(ConnectionHeader, "Upgrade")
pr.Out.Header.Set(UpgradeHeader, upgrade)
}
peer, ok := pr.In.Context().Value(peerCtxKey{}).(api.Peer)
if !ok {
return
+54
View File
@@ -135,6 +135,60 @@ func TestProxyHandler(t *testing.T) {
}
}
func TestProxyHandlerPreservesUpgrade(t *testing.T) {
t.Parallel()
peerLister := &mockPeerLister{
peers: map[string]api.Peer{
"192.0.2.1": {
UserId: "foo",
Groups: []api.GroupMinimum{
{
Name: "group1",
},
},
},
},
}
bearerToken := "foobar"
srv := httptest.NewTLSServer(http.HandlerFunc(func(rw http.ResponseWriter, req *http.Request) {
// Echo the headers that reached the API server so the test can
// assert the upgrade handshake survived and identity is the peer's.
body := fmt.Sprintf("%s %s %s %s",
req.Header.Get(ConnectionHeader),
req.Header.Get(UpgradeHeader),
req.Header.Get(SecWebsocketKeyHeader),
req.Header.Get(ImpersonateUserHeader),
)
// nolint: errcheck
rw.Write([]byte(body))
}))
t.Cleanup(func() {
srv.Close()
})
certPool := srv.Client().Transport.(*http.Transport).TLSClientConfig.RootCAs
kubeAPIServerURL, err := url.Parse(srv.URL)
require.NoError(t, err)
req := httptest.NewRequestWithContext(t.Context(), http.MethodGet,
"/api/v1/namespaces/default/pods/example/exec", nil)
req.Header.Set(ConnectionHeader, "Upgrade")
req.Header.Set(UpgradeHeader, "websocket")
req.Header.Set(SecWebsocketKeyHeader, "dGhlIHNhbXBsZSBub25jZQ==")
// A client must not be able to impersonate by sending this directly.
req.Header.Set(ImpersonateUserHeader, "system:admin")
rec := httptest.NewRecorder()
handler := proxyHandler(peerLister, kubeAPIServerURL, certPool, bearerToken)
handler(rec, req)
b, err := io.ReadAll(rec.Result().Body)
require.NoError(t, err)
require.EqualT(t, http.StatusOK, rec.Result().StatusCode)
require.EqualT(t, "Upgrade websocket dGhlIHNhbXBsZSBub25jZQ== foo", string(b))
}
func TestGenerateSelfSignedCert(t *testing.T) {
t.Parallel()