Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion kubernetes/base/stream/ws_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"):
Expand Down
45 changes: 45 additions & 0 deletions kubernetes/base/stream/ws_client_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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

Expand All @@ -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:
Expand Down