diff --git a/kubernetes/base/stream/ws_client.py b/kubernetes/base/stream/ws_client.py index ffdd877ef4..7754cfa613 100644 --- a/kubernetes/base/stream/ws_client.py +++ b/kubernetes/base/stream/ws_client.py @@ -214,7 +214,9 @@ def update(self, timeout=0): # efficient as epoll. Will work for fd numbers above 1024. # select.epoll() - newest and most efficient way of polling. # However, only works on linux. - if hasattr(select, "poll"): + if self.sock.is_ssl() and self.sock.sock.pending() > 0: + r = [self.sock.sock] + elif hasattr(select, "poll"): poll = select.poll() poll.register(self.sock.sock, select.POLLIN) if timeout is not None and timeout != float("inf"): diff --git a/kubernetes/base/stream/ws_client_test.py b/kubernetes/base/stream/ws_client_test.py index 80f9c692c5..11ff39652f 100644 --- a/kubernetes/base/stream/ws_client_test.py +++ b/kubernetes/base/stream/ws_client_test.py @@ -18,6 +18,7 @@ from . import ws_client as ws_client_module from .ws_client import get_websocket_url, WSClient, V5_CHANNEL_PROTOCOL, V4_CHANNEL_PROTOCOL, CLOSE_CHANNEL, STDIN_CHANNEL from .ws_client import websocket_proxycare +from .ws_client import STDOUT_CHANNEL from kubernetes.client.configuration import Configuration import os import socket @@ -192,6 +193,7 @@ def test_update_receives_close_v5(self): mock_ws = MagicMock() mock_ws.subprotocol = V5_CHANNEL_PROTOCOL mock_ws.connected = True + mock_ws.is_ssl.return_value = False mock_ws.sock.fileno.return_value = 10 # Setup frame with close signal for channel 0 @@ -216,6 +218,7 @@ def test_update_ignores_close_signal_v4(self): mock_ws = MagicMock() mock_ws.subprotocol = V4_CHANNEL_PROTOCOL mock_ws.connected = True + mock_ws.is_ssl.return_value = False mock_ws.sock.fileno.return_value = 10 # Setup frame that looks like close signal but should be treated as data @@ -315,6 +318,7 @@ def test_peek_channel_closed_with_leftover_data(self): mock_ws = MagicMock() mock_ws.subprotocol = V5_CHANNEL_PROTOCOL mock_ws.connected = True + mock_ws.is_ssl.return_value = False mock_ws.sock.fileno.return_value = 10 mock_create.return_value = mock_ws @@ -350,6 +354,7 @@ def test_update_infinite_timeout_polls_without_overflow(self): mock_ws = MagicMock() mock_ws.subprotocol = V5_CHANNEL_PROTOCOL mock_ws.connected = True + mock_ws.is_ssl.return_value = False mock_ws.sock.fileno.return_value = 10 mock_create.return_value = mock_ws @@ -358,6 +363,46 @@ def test_update_infinite_timeout_polls_without_overflow(self): mock_poll.return_value.poll.assert_called_once_with(None) + def test_update_reads_pending_ssl_data_without_polling(self): + with ( + patch.object( + ws_client_module, + 'create_websocket', + ) as mock_create, + patch('select.poll') as mock_poll, + patch('select.select') as mock_select, + ): + mock_ws = MagicMock() + mock_ws.subprotocol = V5_CHANNEL_PROTOCOL + mock_ws.connected = True + mock_ws.is_ssl.return_value = True + mock_ws.sock.pending.return_value = 1 + + frame = MagicMock() + frame.data = bytes([STDOUT_CHANNEL]) + b"pending" + mock_ws.recv_data_frame.return_value = ( + websocket.ABNF.OPCODE_BINARY, + frame, + ) + mock_create.return_value = mock_ws + + client = WSClient( + self.config_mock, + "wss://test", + headers=None, + capture_all=True, + binary=True, + ) + client.update(timeout=None) + + mock_poll.assert_not_called() + mock_select.assert_not_called() + mock_ws.recv_data_frame.assert_called_once_with(True) + self.assertEqual( + client.read_channel(STDOUT_CHANNEL), + b"pending", + ) + def test_readline_channel_returns_empty_string_on_expired_timeout(self): """Verify readline_channel returns '' (not None) when a finite timeout expires""" with patch.object(ws_client_module, 'create_websocket') as mock_create: