From 10c260dfa89adc52ae9e2f217a51ab81b8b17332 Mon Sep 17 00:00:00 2001 From: Dana Powers Date: Sun, 19 Jul 2026 09:13:41 -0700 Subject: [PATCH 1/6] KafkaTCPProxy: sans-IO protocol base class for proxy connections --- kafka/net/manager.py | 20 +++-- kafka/net/proxy.py | 173 +++++++++++++++++++++++++++++++++++++++ test/net/test_manager.py | 19 ----- test/net/test_proxy.py | 71 ++++++++++++++++ 4 files changed, 258 insertions(+), 25 deletions(-) create mode 100644 kafka/net/proxy.py create mode 100644 test/net/test_proxy.py diff --git a/kafka/net/manager.py b/kafka/net/manager.py index 8acc76e4d..4e3c10281 100644 --- a/kafka/net/manager.py +++ b/kafka/net/manager.py @@ -10,6 +10,7 @@ from kafka.net.backend import resolve_backend from kafka.cluster import ClusterMetadata import kafka.errors as Errors +from kafka.net.proxy import KafkaTCPProxy from kafka.net.ssl import KafkaSSLTransport from kafka.net.wakeup_notifier import WakeupNotifier from kafka.protocol.broker_version_data import BrokerVersionData @@ -243,12 +244,19 @@ def ssl_enabled(self): async def _connect(self, node, conn, reset_backoff_on_connect=True, timeout_at=None): try: - await self._net.create_connection( - conn, node.host, node.port, - ssl=self.ssl_context, - proxy_url=self.config['proxy_url'], - socket_options=self.config['socket_options'], - timeout_at=timeout_at) + if self.config['proxy_url']: + proxy = KafkaTCPProxy(self._net, self.config['proxy_url']) + await proxy.create_connection( + conn, node.host, node.port, + ssl=self.ssl_context, + socket_options=self.config['socket_options'], + timeout_at=timeout_at) + else: + await self._net.create_connection( + conn, node.host, node.port, + ssl=self.ssl_context, + socket_options=self.config['socket_options'], + timeout_at=timeout_at) # Note: conn.initialize does not currently raise on error; # errors are pushed to conn.init_future and raised on await conn await conn.initialize(timeout_at=timeout_at) diff --git a/kafka/net/proxy.py b/kafka/net/proxy.py new file mode 100644 index 000000000..c6b06b604 --- /dev/null +++ b/kafka/net/proxy.py @@ -0,0 +1,173 @@ +import logging +import socket +import time +from urllib.parse import urlparse + +import kafka.errors as Errors +from kafka.net.ssl import KafkaSSLTransport + +log = logging.getLogger(__name__) + + +class KafkaTCPProxyStates: + DISCONNECTED = '' + CONNECTING = '' + NEGOTIATE_PROPOSE = '' + NEGOTIATING = '' + AUTHENTICATING = '' + REQUEST_SUBMIT = '' + REQUESTING = '' + READ_ADDRESS = '' + COMPLETE = '' + + +class KafkaTCPProxy: + # scheme => handling class + _registry = {} + SCHEMES = () + + @classmethod + def register_class(cls, klass): + for scheme in klass.SCHEMES: + cls._registry[scheme] = klass + + def __init_subclass__(cls, **kw): + super().__init_subclass__(**kw) + KafkaTCPProxy.register_class(cls) + + def __new__(cls, net, proxy_url): + if proxy_url is None: + return super().__new__(cls) + try: + parsed = urlparse(proxy_url) + except Exception: + raise ValueError('Unable to parse proxy_url: %s' % (proxy_url,)) + if not parsed.scheme: + raise ValueError('proxy_url requires scheme:// (%s)' % (proxy_url,)) + try: + klass = KafkaTCPProxy._registry[parsed.scheme] + except KeyError: + raise ValueError('Unsupported proxy url scheme: %s' % (parsed.scheme)) + return super().__new__(klass) + + def __init__(self, net, proxy_url): + self._net = net + self._proxy_url = urlparse(proxy_url) + if self.proxy_scheme not in self.SCHEMES: + raise ValueError('Unsupported proxy scheme: %s' % (self.proxy_scheme,)) + self._transport = None + self._connect_future = None + self._buf = b'' + self._addrinfo = None + self._state = KafkaTCPProxyStates.DISCONNECTED + self._timeout_at = None + + @property + def proxy_scheme(self): + return self._proxy_url.scheme + + @property + def proxy_host(self): + return self._proxy_url.hostname + + @property + def proxy_port(self): + return self._proxy_url.port + + def set_addrinfo(self, addrinfo): + self._addrinfo = addrinfo + + @property + def _afi(self): + return self._addrinfo[0] if self._addrinfo is not None else None + + @property + def _host(self): + return self._addrinfo[4][0] if self._addrinfo is not None else None + + @property + def _port(self): + return self._addrinfo[4][1] if self._addrinfo is not None else None + + async def create_connection(self, protocol, host, port, *, ssl=None, + socket_options=(), timeout_at=None): + if self.use_remote_lookup(): + self.set_addrinfo((socket.AF_UNSPEC, socket.SOCK_STREAM, socket.IPPROTO_TCP, '', (host, port))) + await self._net.create_connection(self, self.proxy_host, self.proxy_port, + ssl=None, # TODO: support TLS connections to proxy + socket_options=socket_options, + timeout_at=timeout_at) + await self.do_connect(timeout_at=timeout_at) + else: + for addrinfo in await self._net.getaddrinfo(host, port): + self.set_addrinfo(addrinfo) + await self._net.create_connection(self, self.proxy_host, self.proxy_port, + ssl=None, # TODO: support TLS connections to proxy + socket_options=socket_options, + timeout_at=timeout_at) + try: + await self.do_connect(timeout_at=timeout_at) + break + except Exception as exc: + log.debug('Failed to connect to %s: %s', addrinfo, exc) + self._transport.abort(exc) + self.connection_lost(exc) + else: + raise Errors.KafkaConnectionError('Unable to connect to %s:%d via proxy' % (host, port)) + + if ssl: + ssl_wrapper = KafkaSSLTransport(self._net, ssl, host=host) + ssl_wrapper.connection_made(self._transport) + await ssl_wrapper.handshake() + self._transport = ssl_wrapper + + try: + protocol.connection_made(self._transport) + except Exception as e: + self._transport.abort(e) + raise + self._transport = None + + # TODO: what about transport.get_peer() / host_port() ? + + def use_remote_lookup(self): + return True + + def connection_lost(self, exc): + self._state = KafkaTCPProxyStates.DISCONNECTED + self._transport = None + if self._connect_future is not None and not self._connect_future.is_done: + self._connect_future.failure(exc) + + def connection_made(self, transport): + self._transport = transport + self._transport.set_protocol(self) + self._transport.resume_reading() + self._buf = b'' + self._state = KafkaTCPProxyStates.CONNECTING + self._connect_future = self._net.create_future() + self._wrapped_state_machine() + + def data_received(self, data): + """Runs a state machine through connection to authentication to + proxy connection request.""" + self._buf += data + self._wrapped_state_machine() + + def _wrapped_state_machine(self): + try: + if self._timeout_at is not None and self._timeout_at <= time.monotonic(): + raise Errors.KafkaTimeoutError('Proxy connection timeout') + if self._run_state_machine(): + self._connect_future.success(True) + except Exception as exc: + self._state = KafkaTCPProxyStates.DISCONNECTED + self._connect_future.failure(exc) + + async def do_connect(self, timeout_at=None): + self._timeout_at = timeout_at + self._wrapped_state_machine() + await self._connect_future + + def _run_state_machine(self): + raise NotImplementedError() diff --git a/test/net/test_manager.py b/test/net/test_manager.py index ca3816cb4..88a63e7ed 100644 --- a/test/net/test_manager.py +++ b/test/net/test_manager.py @@ -114,25 +114,6 @@ def test_proxy_url_takes_precedence_over_legacy(self, net): ) assert m.config['proxy_url'] == 'socks5://new:1080' - def test_connect_passes_proxy_url(self, net): - """_connect must forward the configured proxy_url to - net.create_connection. Regression guard against the kwarg name - drifting from the create_connection signature.""" - m = KafkaConnectionManager(net, proxy_url='socks5://proxy:1080') - node = MagicMock(host='broker', port=9092, node_id='bootstrap-0') - conn = KafkaConnection(net, node_id='bootstrap-0', **m.config) - - async def fake_create_connection(protocol, host, port, **kwargs): - # Close mid-connect so _connect short-circuits before - # connection_made()/initialize() -- we only assert the forwarded kwarg. - conn.close() - return MagicMock() - - with patch.object(net, 'create_connection', - side_effect=fake_create_connection) as mc: - net.run(m._connect(node, conn)) - assert mc.call_args.kwargs.get('proxy_url') == 'socks5://proxy:1080' - class TestKafkaConnectionManagerBackoff: def test_connection_delay_no_backoff(self, manager): diff --git a/test/net/test_proxy.py b/test/net/test_proxy.py new file mode 100644 index 000000000..3b55c10f2 --- /dev/null +++ b/test/net/test_proxy.py @@ -0,0 +1,71 @@ +from unittest.mock import patch + +import pytest + +from kafka.net.proxy import KafkaTCPProxy + +from kafka.net.http_connect import HttpConnectProxyProtocol +from kafka.net.socks5 import Socks5ProxyProtocol + + +class TestKafkaTCPProxyRegistry: + def test_socks5(self, net): + assert 'socks5' in KafkaTCPProxy._registry + factory = KafkaTCPProxy(net, 'socks5://foo.bar') + assert isinstance(factory, Socks5ProxyProtocol) + + def test_socks5h(self, net): + assert 'socks5h' in KafkaTCPProxy._registry + factory = KafkaTCPProxy(net, 'socks5h://foo.bar') + assert isinstance(factory, Socks5ProxyProtocol) + + def test_http(self, net): + assert 'http' in KafkaTCPProxy._registry + factory = KafkaTCPProxy(net, 'http://proxy:8080') + assert isinstance(factory, HttpConnectProxyProtocol) + + def test_unknown_scheme_raises(self, net): + with pytest.raises(ValueError, match='Unsupported proxy url scheme'): + KafkaTCPProxy(net, 'ftp://proxy:8080') + + def test_no_scheme_raises(self, net): + with pytest.raises(ValueError, match='scheme'): + KafkaTCPProxy(net, 'no-scheme') + + def test_empty_string_raises(self, net): + with pytest.raises(ValueError, match='scheme'): + KafkaTCPProxy(net, '') + + def test_kafka_net_import_registers_socks5(self): + # Importing kafka.net must register Socks5ProxyProtocol. Regression guard + # against accidentally dropping the import from kafka/net/__init__.py. + import kafka.net # noqa: F401 + assert KafkaTCPProxy._registry['socks5'] is Socks5ProxyProtocol + assert KafkaTCPProxy._registry['socks5h'] is Socks5ProxyProtocol + + def test_subclass_auto_registers(self, net): + class _TestProxy(KafkaTCPProxy): + SCHEMES = ('test-autoregister',) + def __init__(self, net, proxy_url): + self.proxy_url = proxy_url + try: + assert KafkaTCPProxy._registry['test-autoregister'] is _TestProxy + proxy = KafkaTCPProxy(net, 'test-autoregister://x') + assert isinstance(proxy, _TestProxy) + assert proxy.proxy_url == 'test-autoregister://x' # pylint: disable=no-member + finally: + KafkaTCPProxy._registry.pop('test-autoregister', None) + + def test_duplicate_scheme_last_wins(self): + prior = KafkaTCPProxy._registry.get('test-dup') + class _First(KafkaTCPProxy): + SCHEMES = ('test-dup',) + class _Second(KafkaTCPProxy): + SCHEMES = ('test-dup',) + try: + assert KafkaTCPProxy._registry['test-dup'] is _Second + finally: + if prior is None: + KafkaTCPProxy._registry.pop('test-dup', None) + else: + KafkaTCPProxy._registry['test-dup'] = prior From 227e2b3b0e7c68653c52d7722da2d41aceb5d799 Mon Sep 17 00:00:00 2001 From: Dana Powers Date: Sun, 19 Jul 2026 10:07:55 -0700 Subject: [PATCH 2/6] Convert Socks5 and HTTP-Connect proxy classes to sans-io state machines --- kafka/net/__init__.py | 6 +- kafka/net/http_connect.py | 142 ++++-------------- kafka/net/socks5.py | 261 ++++++++++------------------------ test/net/test_http_connect.py | 141 ++++++------------ 4 files changed, 152 insertions(+), 398 deletions(-) diff --git a/kafka/net/__init__.py b/kafka/net/__init__.py index 51bf318ca..1f564144f 100644 --- a/kafka/net/__init__.py +++ b/kafka/net/__init__.py @@ -1,8 +1,8 @@ from .connection import KafkaConnection from .manager import KafkaConnectionManager from .metrics import KafkaConnectionMetrics, KafkaManagerMetrics -from .http_connect import HttpConnectProxy -from .socks5 import Socks5Proxy +from .http_connect import HttpConnectProxyProtocol +from .socks5 import Socks5ProxyProtocol from .wakeup_notifier import WakeupNotifier from .compat import KafkaNetClient @@ -11,6 +11,6 @@ __all__ = [ 'KafkaConnection', 'KafkaConnectionManager', 'KafkaConnectionMetrics', 'KafkaManagerMetrics', - 'HttpConnectProxy', 'Socks5Proxy', + 'HttpConnectProxyProtocol', 'Socks5ProxyProtocol', 'WakeupNotifier', 'KafkaNetClient', ] diff --git a/kafka/net/http_connect.py b/kafka/net/http_connect.py index 959b3b71d..e90c94713 100644 --- a/kafka/net/http_connect.py +++ b/kafka/net/http_connect.py @@ -1,29 +1,14 @@ import base64 -import errno import logging -import random -import socket -from urllib.parse import urlparse from kafka.errors import KafkaConnectionError -from kafka.net.backend.inet import KafkaNetSocket +from kafka.net.proxy import KafkaTCPProxy, KafkaTCPProxyStates log = logging.getLogger(__name__) -_WOULD_BLOCK = {errno.EWOULDBLOCK, errno.EAGAIN} -_MAX_RESPONSE_SIZE = 65536 - -class _States: - DISCONNECTED = '' - CONNECTING = '' - SENDING = '' - READING = '' - COMPLETE = '' - - -class HttpConnectProxy(KafkaNetSocket): +class HttpConnectProxyProtocol(KafkaTCPProxy): """Tunnels broker connections through an HTTP CONNECT proxy (RFC 7231 s4.3.6). Registered for the ``http`` scheme -- pass ``proxy_url='http://host:port'`` @@ -35,110 +20,37 @@ class HttpConnectProxy(KafkaNetSocket): SCHEMES = ('http',) - def __init__(self, proxy_url): - self._proxy_url = urlparse(proxy_url) - self._sock = None - self._state = _States.DISCONNECTED - self._send_buf = b'' - self._recv_buf = b'' - self._proxy_addr = self._get_proxy_addr() - - def _get_proxy_addr(self): - addrs = self.dns_lookup(self._proxy_url.hostname, self._proxy_url.port, proxy=True) - if not addrs: - raise KafkaConnectionError('Unable to resolve proxy_url via dns') - return random.choice(addrs) - - def dns_lookup(self, host, port, proxy=False): - if proxy: - return super().dns_lookup(host, port, raise_error=True) - # Always forward broker hostname unresolved; the proxy handles DNS - return [(socket.AF_UNSPEC, socket.SOCK_STREAM, socket.IPPROTO_TCP, '', (host, port))] - - def socket(self, family=socket.AF_UNSPEC, sock_type=socket.SOCK_STREAM, proto=socket.IPPROTO_TCP): - self._target_afi = family - proxy_family, _, _, _, _ = self._proxy_addr - self._sock = socket.socket(proxy_family, sock_type, proto) - return self._sock + def _run_state_machine(self): + if self._state in (KafkaTCPProxyStates.DISCONNECTED, KafkaTCPProxyStates.COMPLETE): + return False - def connect_ex(self, sock, addr): - assert sock is self._sock + if self._state == KafkaTCPProxyStates.CONNECTING: + self._do_connecting() - if self._state == _States.DISCONNECTED: - self._state = _States.CONNECTING + if self._state == KafkaTCPProxyStates.REQUESTING: + self._do_requesting() - if self._state == _States.CONNECTING: - ret = self._do_connecting(addr) - if ret is not None: - return ret + if self._state == KafkaTCPProxyStates.COMPLETE: + return True + else: + return False - if self._state == _States.SENDING: - ret = self._do_sending() - if ret is not None: - return ret - - if self._state == _States.READING: - ret = self._do_reading() - if ret is not None: - return ret - - if self._state == _States.COMPLETE: - return 0 - - return errno.ECONNREFUSED - - def _do_connecting(self, addr): - _, _, _, _, proxy_sockaddr = self._proxy_addr - ret = self._sock.connect_ex(proxy_sockaddr) - if ret and ret != errno.EISCONN: - return ret - host, port = addr[0], addr[1] - headers = 'CONNECT {0}:{1} HTTP/1.1\r\nHost: {0}:{1}\r\n'.format(host, port) + def _do_connecting(self): + headers = 'CONNECT {0}:{1} HTTP/1.1\r\nHost: {0}:{1}\r\n'.format(self._host, self._port) if self._proxy_url.username and self._proxy_url.password: credentials = base64.b64encode( '{0}:{1}'.format(self._proxy_url.username, self._proxy_url.password).encode() ).decode() headers += 'Proxy-Authorization: Basic {}\r\n'.format(credentials) - self._send_buf = (headers + '\r\n').encode() - self._state = _States.SENDING - return None - - def _do_sending(self): - while self._send_buf: - try: - sent = self._sock.send(self._send_buf) - if sent == 0: - log.error('Proxy closed connection while sending CONNECT request') - return errno.ECONNREFUSED - self._send_buf = self._send_buf[sent:] - except OSError as exc: - if exc.errno in _WOULD_BLOCK: - return errno.EWOULDBLOCK - raise - self._state = _States.READING - return None - - def _do_reading(self): - while b'\r\n\r\n' not in self._recv_buf: - try: - chunk = self._sock.recv(4096) - if not chunk: - log.error('Proxy closed connection during CONNECT handshake') - self._sock.close() - return errno.ECONNREFUSED - self._recv_buf += chunk - if len(self._recv_buf) > _MAX_RESPONSE_SIZE: - log.error('Proxy response exceeded %d bytes without end-of-headers', _MAX_RESPONSE_SIZE) - self._sock.close() - return errno.ECONNREFUSED - except OSError as exc: - if exc.errno in _WOULD_BLOCK: - return errno.EWOULDBLOCK - raise - first_line = self._recv_buf.split(b'\r\n')[0] - if b' 200 ' in first_line or first_line.endswith(b' 200'): - self._state = _States.COMPLETE - return None - log.error('HTTP CONNECT to proxy failed: %r', first_line) - self._sock.close() - return errno.ECONNREFUSED + self._transport.write((headers + '\r\n').encode()) + self._state = KafkaTCPProxyStates.REQUESTING + + def _do_requesting(self): + if b'\r\n\r\n' in self._buf: + first_line = self._buf.split(b'\r\n')[0] + if b' 200 ' in first_line or first_line.endswith(b' 200'): + self._state = KafkaTCPProxyStates.COMPLETE + else: + log.error('HTTP CONNECT to proxy failed: %r', first_line) + self._state = KafkaTCPProxyStates.DISCONNECTED + raise KafkaConnectionError('HTTP CONNECT to proxy failed: %r' % (first_line,)) diff --git a/kafka/net/socks5.py b/kafka/net/socks5.py index cb4ab5a17..f92e4bfa9 100644 --- a/kafka/net/socks5.py +++ b/kafka/net/socks5.py @@ -1,262 +1,157 @@ -import errno import logging -import random import socket import struct -from urllib.parse import urlparse from kafka.errors import KafkaConnectionError -from kafka.net.backend.inet import KafkaNetSocket +from kafka.net.proxy import KafkaTCPProxy, KafkaTCPProxyStates log = logging.getLogger(__name__) -class ProxyConnectionStates: - DISCONNECTED = '' - CONNECTING = '' - NEGOTIATE_PROPOSE = '' - NEGOTIATING = '' - AUTHENTICATING = '' - REQUEST_SUBMIT = '' - REQUESTING = '' - READ_ADDRESS = '' - COMPLETE = '' +class Socks5ProxyProtocol(KafkaTCPProxy): + """Socks5 proxy sans-IO protocol handler.""" -class Socks5Proxy(KafkaNetSocket): - """Socks5 proxy - - Manages connection through socks5 proxy with support for username/password - authentication. - """ # socks5h for remote dns SCHEMES = ('socks5', 'socks5h') - def __init__(self, proxy_url): - self._buffer_in = b'' - self._buffer_out = b'' - self._proxy_url = urlparse(proxy_url) - if self._proxy_url.scheme not in self.SCHEMES: - raise ValueError('Unsupported proxy scheme: %s' % (self._proxy_url.scheme,)) - self._sock = None - self._state = ProxyConnectionStates.DISCONNECTED - self._target_afi = socket.AF_UNSPEC - self._proxy_addr = self._get_proxy_addr() - - def _get_proxy_addr(self): - proxy_addrs = self.dns_lookup(self._proxy_url.hostname, self._proxy_url.port, proxy=True) - if not proxy_addrs: - raise KafkaConnectionError('Unable to resolve proxy_url via dns') - return random.choice(proxy_addrs) - - def _use_remote_lookup(self): - return self._proxy_url.scheme == 'socks5h' - - def dns_lookup(self, host, port, proxy=False): - if proxy: - return super().dns_lookup(host, port, raise_error=True) - elif self._use_remote_lookup(): - return [(socket.AF_UNSPEC, socket.SOCK_STREAM, socket.IPPROTO_TCP, '', (host, port))] - else: - return super().dns_lookup(host, port) - - def socket(self, family=socket.AF_UNSPEC, sock_type=socket.SOCK_STREAM, proto=socket.IPPROTO_TCP): - """Open and record a socket. - - Returns the actual underlying socket - object to ensure e.g. selects and ssl wrapping works as expected. - """ - self._target_afi = family # Store the address family of the target - proxy_family, _, _, _, _ = self._proxy_addr - self._sock = socket.socket(proxy_family, sock_type, proto) - return self._sock - - def _flush_buf(self): - """Send out all data that is stored in the outgoing buffer. - - It is expected that the caller handles error handling, including non-blocking - as well as connection failure exceptions. - """ - while self._buffer_out: - sent_bytes = self._sock.send(self._buffer_out) - self._buffer_out = self._buffer_out[sent_bytes:] - - def _peek_buf(self, datalen): - """Ensure local inbound buffer has enough data, and return that data without - consuming the local buffer - - It's expected that the caller handles e.g. blocking exceptions""" - while True: - bytes_remaining = datalen - len(self._buffer_in) - if bytes_remaining <= 0: - break - data = self._sock.recv(bytes_remaining) - if not data: - break - self._buffer_in = self._buffer_in + data + def use_remote_lookup(self): + return self.proxy_scheme == 'socks5h' - return self._buffer_in[:datalen] + def _read_buf(self, num_bytes): + if len(self._buf) < num_bytes: + raise IndexError('Not enought bytes in buffer') + data, self._buf = self._buf[0:num_bytes], self._buf[num_bytes:] + return data - def _read_buf(self, datalen): - """Read and consume bytes from socket connection + def _run_state_machine(self): + if self._state in (KafkaTCPProxyStates.DISCONNECTED, KafkaTCPProxyStates.COMPLETE): + return False - It's expected that the caller handles e.g. blocking exceptions""" - buf = self._peek_buf(datalen) - if buf: - self._buffer_in = self._buffer_in[len(buf):] - return buf - - def connect_ex(self, sock, addr): - """Runs a state machine through connection to authentication to - proxy connection request. - - The somewhat strange setup is to facilitate non-intrusive use from - BrokerConnection state machine. - - This function is called with a socket in non-blocking mode. Both - send and receive calls can return in EWOULDBLOCK/EAGAIN which we - specifically avoid handling here. These are handled in main - BrokerConnection connection loop, which then would retry calls - to this function.""" - assert sock is self._sock - if self._state == ProxyConnectionStates.DISCONNECTED: - self._state = ProxyConnectionStates.CONNECTING - - if self._state == ProxyConnectionStates.CONNECTING: - _, _, _, _, sockaddr = self._proxy_addr - ret = self._sock.connect_ex(sockaddr) - if not ret or ret == errno.EISCONN: - self._state = ProxyConnectionStates.NEGOTIATE_PROPOSE - else: - return ret - - if self._state == ProxyConnectionStates.NEGOTIATE_PROPOSE: + if self._state == KafkaTCPProxyStates.CONNECTING: if self._proxy_url.username and self._proxy_url.password: # Propose username/password - self._buffer_out = b"\x05\x01\x02" + self._transport.write(b"\x05\x01\x02") else: # Propose no auth - self._buffer_out = b"\x05\x01\x00" - self._state = ProxyConnectionStates.NEGOTIATING + self._transport.write(b"\x05\x01\x00") + self._state = KafkaTCPProxyStates.NEGOTIATING + + if self._state == KafkaTCPProxyStates.NEGOTIATING: + try: + buf = self._read_buf(2) + except IndexError: + return False - if self._state == ProxyConnectionStates.NEGOTIATING: - self._flush_buf() - buf = self._read_buf(2) if buf[0:1] != b"\x05": log.error("Unrecognized SOCKS version") - self._state = ProxyConnectionStates.DISCONNECTED - self._sock.close() - return errno.ECONNREFUSED + raise KafkaConnectionError('Unrecognized SOCKS version') if buf[1:2] == b"\x00": # No authentication required - self._state = ProxyConnectionStates.REQUEST_SUBMIT + self._state = KafkaTCPProxyStates.REQUEST_SUBMIT elif buf[1:2] == b"\x02": # Username/password authentication selected userlen = len(self._proxy_url.username) passlen = len(self._proxy_url.password) - self._buffer_out = struct.pack( + self._transport.write(struct.pack( "!bb{}sb{}s".format(userlen, passlen), 1, # version userlen, self._proxy_url.username.encode(), passlen, self._proxy_url.password.encode(), - ) - self._state = ProxyConnectionStates.AUTHENTICATING + )) + self._state = KafkaTCPProxyStates.AUTHENTICATING else: log.error("Unrecognized SOCKS authentication method") - self._state = ProxyConnectionStates.DISCONNECTED - self._sock.close() - return errno.ECONNREFUSED + raise KafkaConnectionError('Unrecognized SOCKS authentication method') - if self._state == ProxyConnectionStates.AUTHENTICATING: - self._flush_buf() - buf = self._read_buf(2) + if self._state == KafkaTCPProxyStates.AUTHENTICATING: + try: + buf = self._read_buf(2) + except IndexError: + return False if buf == b"\x01\x00": # Authentication succesful - self._state = ProxyConnectionStates.REQUEST_SUBMIT + self._state = KafkaTCPProxyStates.REQUEST_SUBMIT else: log.error("Socks5 proxy authentication failure") - self._state = ProxyConnectionStates.DISCONNECTED - self._sock.close() - return errno.ECONNREFUSED + raise KafkaConnectionError('Socks5 proxy authentication failure') - if self._state == ProxyConnectionStates.REQUEST_SUBMIT: - if self._use_remote_lookup(): + if self._state == KafkaTCPProxyStates.REQUEST_SUBMIT: + if self.use_remote_lookup(): addr_type = 3 - addr_len = len(addr[0]) - elif self._target_afi == socket.AF_INET: + addr_len = len(self._host) + elif self._afi == socket.AF_INET: addr_type = 1 addr_len = 4 - elif self._target_afi == socket.AF_INET6: + elif self._afi == socket.AF_INET6: addr_type = 4 addr_len = 16 else: - log.error("Unknown address family, %r", self._target_afi) - self._state = ProxyConnectionStates.DISCONNECTED - self._sock.close() - return errno.ECONNREFUSED + log.error("Unknown address family, %r", self._afi) + raise KafkaConnectionError("Unknown address family, %r" % self._afi) - self._buffer_out = struct.pack( + self._transport.write(struct.pack( "!bbbb", 5, # version 1, # command: connect 0, # reserved addr_type, # 1 for ipv4, 4 for ipv6 address, 3 for domain name - ) + )) # Addr format depends on type if addr_type == 3: # len + domain name (no null terminator) - self._buffer_out += struct.pack( + self._transport.write(struct.pack( "!b{}s".format(addr_len), addr_len, - addr[0].encode('ascii'), - ) + self._host.encode('ascii'), + )) else: # either 4 (type 1) or 16 (type 4) bytes of actual address - self._buffer_out += struct.pack( + self._transport.write(struct.pack( "!{}s".format(addr_len), - socket.inet_pton(self._target_afi, addr[0]), - ) - self._buffer_out += struct.pack("!H", addr[1]) # port - - self._state = ProxyConnectionStates.REQUESTING - - if self._state == ProxyConnectionStates.REQUESTING: - self._flush_buf() - buf = self._read_buf(2) + socket.inet_pton(self._afi, self._host), + )) + self._transport.write(struct.pack("!H", self._port)) # port + self._state = KafkaTCPProxyStates.REQUESTING + + if self._state == KafkaTCPProxyStates.REQUESTING: + try: + buf = self._read_buf(2) + except IndexError: + return False if buf[0:2] == b"\x05\x00": - self._state = ProxyConnectionStates.READ_ADDRESS + self._state = KafkaTCPProxyStates.READ_ADDRESS else: log.error("Proxy request failed: %r", buf[1:2]) - self._state = ProxyConnectionStates.DISCONNECTED - self._sock.close() - return errno.ECONNREFUSED + raise KafkaConnectionError("Proxy request failed: %r" % buf[1:2]) - if self._state == ProxyConnectionStates.READ_ADDRESS: + if self._state == KafkaTCPProxyStates.READ_ADDRESS: # we don't really care about the remote endpoint address, but need to clear the stream - buf = self._peek_buf(2) - if buf[0:2] == b"\x00\x01": - _ = self._read_buf(2 + 4 + 2) # ipv4 address + port - elif buf[0:2] == b"\x00\x05": - _ = self._read_buf(2 + 16 + 2) # ipv6 address + port + if len(self._buf) < 2: + return False + if self._buf[0:2] == b"\x00\x01": + addrbytes = 4 + elif self._buf[0:2] == b"\x00\x05": + addrbytes = 16 else: log.error("Unrecognized remote address type %r", buf[1:2]) - self._state = ProxyConnectionStates.DISCONNECTED - self._sock.close() - return errno.ECONNREFUSED - self._state = ProxyConnectionStates.COMPLETE + raise KafkaConnectionError("Unrecognized remote address type %r", self._buf[1:2]) + try: + self._read_buf(2 + addrbytes + 2) # header + address + port + except IndexError: + return False + self._state = KafkaTCPProxyStates.COMPLETE + assert not self._buf, 'unexpected bytes remaining in buffer' - if self._state == ProxyConnectionStates.COMPLETE: - return 0 + if self._state == KafkaTCPProxyStates.COMPLETE: + return True # not reached; # Send and recv will raise socket error on EWOULDBLOCK/EAGAIN that is assumed to be handled by # the caller. The caller re-enters this state machine from retry logic with timer or via select & family log.error("Internal error, state %r not handled correctly", self._state) - self._state = ProxyConnectionStates.DISCONNECTED - if self._sock: - self._sock.close() - return errno.ECONNREFUSED + raise KafkaConnectionError("Internal Socks5 proxy state error (%r)" % (self._state,)) diff --git a/test/net/test_http_connect.py b/test/net/test_http_connect.py index 19a654090..80f6c913f 100644 --- a/test/net/test_http_connect.py +++ b/test/net/test_http_connect.py @@ -4,113 +4,60 @@ import pytest -from kafka.net.http_connect import HttpConnectProxy -from kafka.net.backend.inet import KafkaNetSocket +from kafka.net.http_connect import HttpConnectProxyProtocol +from kafka.net.proxy import KafkaTCPProxy, KafkaTCPProxyStates _FAKE_PROXY_ADDR = (socket.AF_INET, socket.SOCK_STREAM, socket.IPPROTO_TCP, '', ('1.2.3.4', 8080)) -def _make_proxy(url='http://proxy:8080'): - with patch.object(HttpConnectProxy, 'dns_lookup', return_value=[_FAKE_PROXY_ADDR]): - proxy = HttpConnectProxy(url) - proxy._sock = MagicMock() - return proxy - - class TestHttpConnectProxyRegistry: - def test_registered_for_http_scheme(self): - with patch.object(HttpConnectProxy, 'dns_lookup', return_value=[_FAKE_PROXY_ADDR]): - obj = KafkaNetSocket('http://proxy:8080') - assert isinstance(obj, HttpConnectProxy) + def test_registered_for_http_scheme(self, net): + obj = KafkaTCPProxy(net, 'http://proxy:8080') + assert isinstance(obj, HttpConnectProxyProtocol) - def test_unregistered_scheme_raises(self): + def test_unregistered_scheme_raises(self, net): with pytest.raises(ValueError, match='Unsupported proxy url scheme'): - KafkaNetSocket('socks4://proxy:8080') - - -class TestHttpConnectProxyDnsLookup: - def test_broker_lookup_returns_unresolved(self): - proxy = _make_proxy() - result = proxy.dns_lookup('broker.kafka.internal', 9092) - assert result == [(socket.AF_UNSPEC, socket.SOCK_STREAM, socket.IPPROTO_TCP, '', ('broker.kafka.internal', 9092))] - - def test_proxy_lookup_delegates_to_super(self): - with patch('socket.getaddrinfo', return_value=[_FAKE_PROXY_ADDR]) as mock_gai: - proxy = _make_proxy() - proxy.dns_lookup('proxy', 8080, proxy=True) - mock_gai.assert_called() - - -class TestHttpConnectProxySocket: - def test_uses_proxy_family_not_broker_family(self): - proxy = _make_proxy() - with patch('socket.socket') as mock_ctor: - proxy.socket(socket.AF_INET6, socket.SOCK_STREAM, socket.IPPROTO_TCP) - mock_ctor.assert_called_once_with(socket.AF_INET, socket.SOCK_STREAM, socket.IPPROTO_TCP) - - -class TestHttpConnectProxyConnectEx: - def test_success(self): - proxy = _make_proxy() - proxy._sock.connect_ex.return_value = 0 - proxy._sock.send.side_effect = lambda b: len(b) - proxy._sock.recv.return_value = b'HTTP/1.1 200 Connection Established\r\n\r\n' - assert proxy.connect_ex(proxy._sock, ('broker', 9092)) == 0 - - def test_success_no_reason_phrase(self): - proxy = _make_proxy() - proxy._sock.connect_ex.return_value = 0 - proxy._sock.send.side_effect = lambda b: len(b) - proxy._sock.recv.return_value = b'HTTP/1.1 200\r\n\r\n' - assert proxy.connect_ex(proxy._sock, ('broker', 9092)) == 0 - - def test_basic_auth_header_sent_when_credentials_in_url(self): + KafkaTCPProxy(net, 'socks4://proxy:8080') + + +class TestHttpConnectProxyStateMachine: + def test_success(self, net): + proxy = KafkaTCPProxy(net, 'http://proxy:8080') + proxy.connection_made(MagicMock()) + proxy.data_received(b'HTTP/1.1 200 Connection Established\r\n\r\n') + assert proxy._connect_future.is_done + assert proxy._connect_future.succeeded() + assert proxy._state == KafkaTCPProxyStates.COMPLETE + + def test_success_no_reason_phrase(self, net): + proxy = KafkaTCPProxy(net, 'http://proxy:8080') + proxy.connection_made(MagicMock()) + proxy.data_received(b'HTTP/1.1 200\r\n\r\n') + assert proxy._connect_future.is_done + assert proxy._connect_future.succeeded() + assert proxy._state == KafkaTCPProxyStates.COMPLETE + + def test_basic_auth_header_sent_when_credentials_in_url(self, net): import base64 - proxy = _make_proxy('http://user:pass@proxy:8080') - proxy._sock.connect_ex.return_value = 0 + proxy = KafkaTCPProxy(net, 'http://user:pass@proxy:8080') + transport = MagicMock() sent = [] - proxy._sock.send.side_effect = lambda b: sent.append(b) or len(b) - proxy._sock.recv.return_value = b'HTTP/1.1 200 Connection Established\r\n\r\n' - proxy.connect_ex(proxy._sock, ('broker', 9092)) + transport.write.side_effect = lambda b: sent.append(b) or len(b) + proxy.connection_made(transport) request = b''.join(sent).decode() expected = base64.b64encode(b'user:pass').decode() assert 'Proxy-Authorization: Basic {}'.format(expected) in request - - def test_non_200_response_returns_econnrefused(self): - proxy = _make_proxy() - proxy._sock.connect_ex.return_value = 0 - proxy._sock.send.side_effect = lambda b: len(b) - proxy._sock.recv.return_value = b'HTTP/1.1 407 Proxy Authentication Required\r\n\r\n' - assert proxy.connect_ex(proxy._sock, ('broker', 9092)) == errno.ECONNREFUSED - - def test_ewouldblock_while_sending(self): - proxy = _make_proxy() - proxy._sock.connect_ex.return_value = 0 - proxy._sock.send.side_effect = OSError(errno.EWOULDBLOCK, 'would block') - assert proxy.connect_ex(proxy._sock, ('broker', 9092)) == errno.EWOULDBLOCK - - def test_ewouldblock_while_reading(self): - proxy = _make_proxy() - proxy._sock.connect_ex.return_value = 0 - proxy._sock.send.side_effect = lambda b: len(b) - proxy._sock.recv.side_effect = OSError(errno.EWOULDBLOCK, 'would block') - assert proxy.connect_ex(proxy._sock, ('broker', 9092)) == errno.EWOULDBLOCK - - def test_resumes_after_ewouldblock(self): - proxy = _make_proxy() - proxy._sock.connect_ex.return_value = 0 - proxy._sock.send.side_effect = lambda b: len(b) - proxy._sock.recv.side_effect = [ - OSError(errno.EWOULDBLOCK, 'would block'), - b'HTTP/1.1 200 Connection Established\r\n\r\n', - ] - assert proxy.connect_ex(proxy._sock, ('broker', 9092)) == errno.EWOULDBLOCK - assert proxy.connect_ex(proxy._sock, ('broker', 9092)) == 0 - - def test_eof_during_handshake_returns_econnrefused(self): - proxy = _make_proxy() - proxy._sock.connect_ex.return_value = 0 - proxy._sock.send.side_effect = lambda b: len(b) - proxy._sock.recv.return_value = b'' - assert proxy.connect_ex(proxy._sock, ('broker', 9092)) == errno.ECONNREFUSED + assert not proxy._connect_future.is_done + proxy.data_received(b'HTTP/1.1 200 Connection Established\r\n\r\n') + assert proxy._connect_future.is_done + assert proxy._connect_future.succeeded() + assert proxy._state == KafkaTCPProxyStates.COMPLETE + + def test_non_200_response_disconnects(self, net): + proxy = KafkaTCPProxy(net, 'http://proxy:8080') + proxy.connection_made(MagicMock()) + proxy.data_received(b'HTTP/1.1 407 Proxy Authentication Required\r\n\r\n') + assert proxy._connect_future.is_done + assert proxy._connect_future.failed() + assert proxy._state == KafkaTCPProxyStates.DISCONNECTED From 6bdf982461b129d1df722f38b7e2f4839811e1d6 Mon Sep 17 00:00:00 2001 From: Dana Powers Date: Sun, 19 Jul 2026 11:42:42 -0700 Subject: [PATCH 3/6] Drop proxy_url from backend create_connection interface --- kafka/net/backend/abstract.py | 4 +--- kafka/net/backend/asyncio_backend.py | 6 +----- kafka/net/backend/selector.py | 3 +-- test/net/backend/test_asyncio_backend.py | 7 ------- 4 files changed, 3 insertions(+), 17 deletions(-) diff --git a/kafka/net/backend/abstract.py b/kafka/net/backend/abstract.py index f4fdb8ce6..e6d657b69 100644 --- a/kafka/net/backend/abstract.py +++ b/kafka/net/backend/abstract.py @@ -213,7 +213,6 @@ async def create_connection( port: int, *, ssl: Any = None, - proxy_url: Optional[str] = None, socket_options: Sequence[Any] = (), timeout_at: Optional[float] = None, ) -> None: @@ -228,8 +227,7 @@ async def create_connection( transport handle. ``protocol`` (a ``KafkaConnection``) may *refuse* the transport by raising from ``connection_made`` if it closed mid-connect; on that (or any) failure the backend closes the orphaned transport - before propagating. Backends without native proxy support raise when - ``proxy_url`` is set. + before propagating. """ # --- cross-thread bridge --------------------------------------------- diff --git a/kafka/net/backend/asyncio_backend.py b/kafka/net/backend/asyncio_backend.py index e268ef935..7eeb074ac 100644 --- a/kafka/net/backend/asyncio_backend.py +++ b/kafka/net/backend/asyncio_backend.py @@ -342,11 +342,7 @@ async def getaddrinfo(self, host, port): return await self._loop.getaddrinfo(host, port) async def create_connection(self, protocol, host, port, *, ssl=None, - proxy_url=None, socket_options=(), timeout_at=None): - if proxy_url is not None: - raise NotImplementedError( - 'The asyncio backend does not support proxy_url yet; use the ' - 'default selector backend for SOCKS5/HTTP-CONNECT proxying.') + socket_options=(), timeout_at=None): server_hostname = host.rstrip('.') if ssl is not None else None adapter = _AsyncioProtocolAdapter(protocol, host, port, socket_options) connect = self._loop.create_connection( diff --git a/kafka/net/backend/selector.py b/kafka/net/backend/selector.py index e56626add..c7afff710 100644 --- a/kafka/net/backend/selector.py +++ b/kafka/net/backend/selector.py @@ -624,8 +624,7 @@ async def connect_addrinfo(self, addrinfo, socket_options=(), timeout_at=None): raise Errors.KafkaTimeoutError('Connection timed out') async def create_connection(self, protocol, host, port, *, ssl=None, - proxy_url=None, socket_options=(), - timeout_at=None): + socket_options=(), timeout_at=None): """Establish a connected transport to host:port and wire ``protocol``. The selector owns the raw socket: DNS + non-blocking connect, then diff --git a/test/net/backend/test_asyncio_backend.py b/test/net/backend/test_asyncio_backend.py index d59b916eb..d3cb24872 100644 --- a/test/net/backend/test_asyncio_backend.py +++ b/test/net/backend/test_asyncio_backend.py @@ -319,13 +319,6 @@ async def do(): finally: srv.close() - def test_proxy_url_raises(self, started_backend): - async def do(): - await started_backend.create_connection( - _StubProtocol(), 'h', 1, proxy_url='socks5://proxy:1080') - with pytest.raises(NotImplementedError, match='proxy'): - started_backend.run(do) - class TestLoopConfig: def test_loop_factory_used_and_owned(self): From 75703eb165ebc5f59c8b0a96abe2ff8f571d97b3 Mon Sep 17 00:00:00 2001 From: Dana Powers Date: Mon, 20 Jul 2026 10:04:54 -0700 Subject: [PATCH 4/6] Add Socks5 tests; expand http-connect coverage --- test/net/test_http_connect.py | 50 ++++++++++++ test/net/test_socks5.py | 149 ++++++++++++++++++++++++++++++++++ 2 files changed, 199 insertions(+) create mode 100644 test/net/test_socks5.py diff --git a/test/net/test_http_connect.py b/test/net/test_http_connect.py index 80f6c913f..28a7ab936 100644 --- a/test/net/test_http_connect.py +++ b/test/net/test_http_connect.py @@ -61,3 +61,53 @@ def test_non_200_response_disconnects(self, net): assert proxy._connect_future.is_done assert proxy._connect_future.failed() assert proxy._state == KafkaTCPProxyStates.DISCONNECTED + + +class TestHttpConnectProxyFraming: + TARGET = (socket.AF_UNSPEC, socket.SOCK_STREAM, socket.IPPROTO_TCP, '', ('broker.example.com', 9092)) + + def _proxy(self, net, url='http://proxy:8080'): + proxy = KafkaTCPProxy(net, url) + proxy.set_addrinfo(self.TARGET) + sent = [] + transport = MagicMock() + transport.write.side_effect = lambda b: sent.append(bytes(b)) or len(b) + proxy.connection_made(transport) + return proxy, sent + + def test_connect_line_uses_target_host_port(self, net): + # Regression: the target host/port must be set before the state machine + # sends CONNECT (an earlier ordering bug emitted "CONNECT None:None"). + _, sent = self._proxy(net) + request = b''.join(sent).decode() + assert request.startswith('CONNECT broker.example.com:9092 HTTP/1.1\r\n') + assert 'Host: broker.example.com:9092\r\n' in request + + def test_incomplete_response_waits(self, net): + proxy, _ = self._proxy(net) + # No blank line yet -> keep buffering, stay in REQUESTING. + proxy.data_received(b'HTTP/1.1 200 Connection Established\r\n') + assert not proxy._connect_future.is_done + assert proxy._state == KafkaTCPProxyStates.REQUESTING + + def test_fragmented_response_completes(self, net): + proxy, _ = self._proxy(net) + proxy.data_received(b'HTTP/1.1 200 Connection Established\r\n') + assert not proxy._connect_future.is_done + proxy.data_received(b'\r\n') + assert proxy._connect_future.succeeded() + assert proxy._state == KafkaTCPProxyStates.COMPLETE + + def test_byte_by_byte_delivery(self, net): + proxy, _ = self._proxy(net) + for byte in b'HTTP/1.1 200 OK\r\n\r\n': + assert not proxy._connect_future.is_done + proxy.data_received(bytes([byte])) + assert proxy._connect_future.succeeded() + assert proxy._state == KafkaTCPProxyStates.COMPLETE + + def test_trailing_bytes_after_headers(self, net): + proxy, _ = self._proxy(net) + proxy.data_received(b'HTTP/1.1 200 OK\r\n\r\nleftover') + assert proxy._connect_future.succeeded() + assert proxy._state == KafkaTCPProxyStates.COMPLETE diff --git a/test/net/test_socks5.py b/test/net/test_socks5.py new file mode 100644 index 000000000..6e56640b4 --- /dev/null +++ b/test/net/test_socks5.py @@ -0,0 +1,149 @@ +import socket +import struct + +from unittest.mock import MagicMock + +from kafka.net.proxy import KafkaTCPProxy, KafkaTCPProxyStates +from kafka.net.socks5 import Socks5ProxyProtocol + + +# SOCKS5 reply granting a CONNECT to an IPv4 endpoint: +# VER=05 REP=00 RSV=00 ATYP=01(ipv4) + 4 addr bytes + 2 port bytes. +_GRANT_IPV4 = b'\x05\x00\x00\x01' + b'\x7f\x00\x00\x01' + struct.pack('!H', 0) + + +def _connect(net, url, target): + """Build a proxy, wire a recording transport, and run through connection_made. + + Returns (proxy, sent) where ``sent`` accumulates every ``transport.write``. + Mirrors what manager._connect does: set the target addrinfo before the + transport's connection_made drives the first state-machine step. + """ + proxy = KafkaTCPProxy(net, url) + proxy.set_addrinfo(target) + sent = [] + transport = MagicMock() + transport.write.side_effect = lambda b: sent.append(bytes(b)) or len(b) + proxy.connection_made(transport) + return proxy, sent + + +class TestSocks5Registry: + def test_dispatch(self, net): + assert isinstance(KafkaTCPProxy(net, 'socks5://p:1080'), Socks5ProxyProtocol) + assert isinstance(KafkaTCPProxy(net, 'socks5h://p:1080'), Socks5ProxyProtocol) + + def test_remote_lookup_only_for_socks5h(self, net): + assert KafkaTCPProxy(net, 'socks5://p:1080').use_remote_lookup() is False + assert KafkaTCPProxy(net, 'socks5h://p:1080').use_remote_lookup() is True + + +class TestSocks5NoAuth: + TARGET = (socket.AF_INET, socket.SOCK_STREAM, socket.IPPROTO_TCP, '', ('1.2.3.4', 9092)) + REQUEST = b'\x05\x01\x00\x01' + socket.inet_pton(socket.AF_INET, '1.2.3.4') + struct.pack('!H', 9092) + + def test_proposes_no_auth(self, net): + proxy, sent = _connect(net, 'socks5://proxy:1080', self.TARGET) + assert b''.join(sent) == b'\x05\x01\x00' + assert proxy._state == KafkaTCPProxyStates.NEGOTIATING + assert not proxy._connect_future.is_done + + def test_happy_path(self, net): + proxy, sent = _connect(net, 'socks5://proxy:1080', self.TARGET) + # server selects no-auth -> we submit the CONNECT request + proxy.data_received(b'\x05\x00') + assert proxy._state == KafkaTCPProxyStates.REQUESTING + assert b''.join(sent) == b'\x05\x01\x00' + self.REQUEST + # server grants + proxy.data_received(_GRANT_IPV4) + assert proxy._state == KafkaTCPProxyStates.COMPLETE + assert proxy._connect_future.succeeded() + + def test_request_rejected(self, net): + proxy, _ = _connect(net, 'socks5://proxy:1080', self.TARGET) + proxy.data_received(b'\x05\x00') + proxy.data_received(b'\x05\x01') # REP=01 -> general failure + assert proxy._connect_future.failed() + assert proxy._state == KafkaTCPProxyStates.DISCONNECTED + + def test_bad_socks_version(self, net): + proxy, _ = _connect(net, 'socks5://proxy:1080', self.TARGET) + proxy.data_received(b'\x04\x00') + assert proxy._connect_future.failed() + assert proxy._state == KafkaTCPProxyStates.DISCONNECTED + + def test_unrecognized_auth_method(self, net): + proxy, _ = _connect(net, 'socks5://proxy:1080', self.TARGET) + proxy.data_received(b'\x05\xff') + assert proxy._connect_future.failed() + assert proxy._state == KafkaTCPProxyStates.DISCONNECTED + + +class TestSocks5UserPass: + TARGET = (socket.AF_INET, socket.SOCK_STREAM, socket.IPPROTO_TCP, '', ('1.2.3.4', 9092)) + # subnegotiation: VER=01 ULEN=04 'user' PLEN=04 'pass' + AUTH = struct.pack('!bb4sb4s', 1, 4, b'user', 4, b'pass') + + def test_proposes_user_pass(self, net): + proxy, sent = _connect(net, 'socks5://user:pass@proxy:1080', self.TARGET) + assert b''.join(sent) == b'\x05\x01\x02' + assert proxy._state == KafkaTCPProxyStates.NEGOTIATING + + def test_auth_success(self, net): + proxy, sent = _connect(net, 'socks5://user:pass@proxy:1080', self.TARGET) + # server selects user/pass -> we send credentials + proxy.data_received(b'\x05\x02') + assert proxy._state == KafkaTCPProxyStates.AUTHENTICATING + assert b''.join(sent) == b'\x05\x01\x02' + self.AUTH + # server accepts credentials -> CONNECT request submitted + proxy.data_received(b'\x01\x00') + assert proxy._state == KafkaTCPProxyStates.REQUESTING + proxy.data_received(_GRANT_IPV4) + assert proxy._connect_future.succeeded() + + def test_auth_failure(self, net): + proxy, _ = _connect(net, 'socks5://user:pass@proxy:1080', self.TARGET) + proxy.data_received(b'\x05\x02') + proxy.data_received(b'\x01\x01') # non-zero status -> auth failed + assert proxy._connect_future.failed() + assert proxy._state == KafkaTCPProxyStates.DISCONNECTED + + +class TestSocks5AddressEncoding: + def _request_after_select(self, net, url, target): + proxy, sent = _connect(net, url, target) + proxy.data_received(b'\x05\x00') # no-auth selected -> REQUEST_SUBMIT + # sent[0] is the proposal; the remainder is the CONNECT request + return b''.join(sent[1:]) + + def test_ipv4(self, net): + target = (socket.AF_INET, socket.SOCK_STREAM, socket.IPPROTO_TCP, '', ('1.2.3.4', 9092)) + req = self._request_after_select(net, 'socks5://proxy:1080', target) + assert req == b'\x05\x01\x00\x01' + socket.inet_pton(socket.AF_INET, '1.2.3.4') + struct.pack('!H', 9092) + + def test_ipv6(self, net): + target = (socket.AF_INET6, socket.SOCK_STREAM, socket.IPPROTO_TCP, '', ('::1', 9092)) + req = self._request_after_select(net, 'socks5://proxy:1080', target) + assert req == b'\x05\x01\x00\x04' + socket.inet_pton(socket.AF_INET6, '::1') + struct.pack('!H', 9092) + + def test_remote_lookup_domain(self, net): + # socks5h forwards the hostname unresolved (ATYP=3, domain name). + target = (socket.AF_UNSPEC, socket.SOCK_STREAM, socket.IPPROTO_TCP, '', ('example.com', 9092)) + req = self._request_after_select(net, 'socks5h://proxy:1080', target) + expected = b'\x05\x01\x00\x03' + bytes([len('example.com')]) + b'example.com' + struct.pack('!H', 9092) + assert req == expected + + +class TestSocks5Fragmented: + TARGET = (socket.AF_INET, socket.SOCK_STREAM, socket.IPPROTO_TCP, '', ('1.2.3.4', 9092)) + + def test_byte_by_byte_delivery(self, net): + # The sans-IO invariant: the state machine must buffer across arbitrarily + # fragmented data_received calls, not assume a whole message per read. + proxy, _ = _connect(net, 'socks5://proxy:1080', self.TARGET) + stream = b'\x05\x00' + _GRANT_IPV4 # method-select reply, then CONNECT grant + for byte in stream: + assert not proxy._connect_future.is_done + proxy.data_received(bytes([byte])) + assert proxy._state == KafkaTCPProxyStates.COMPLETE + assert proxy._connect_future.succeeded() From 7b7643c647e03571dc059c840f57eaaa149fdd06 Mon Sep 17 00:00:00 2001 From: Dana Powers Date: Mon, 20 Jul 2026 10:11:15 -0700 Subject: [PATCH 5/6] Add proxy tests --- test/net/test_proxy.py | 212 ++++++++++++++++++++++++++++++++++++++++- 1 file changed, 210 insertions(+), 2 deletions(-) diff --git a/test/net/test_proxy.py b/test/net/test_proxy.py index 3b55c10f2..47f93aee9 100644 --- a/test/net/test_proxy.py +++ b/test/net/test_proxy.py @@ -1,13 +1,74 @@ -from unittest.mock import patch +import asyncio +import socket +from unittest.mock import AsyncMock, MagicMock, patch import pytest -from kafka.net.proxy import KafkaTCPProxy +import kafka.errors as Errors +from kafka.net.proxy import KafkaTCPProxy, KafkaTCPProxyStates from kafka.net.http_connect import HttpConnectProxyProtocol from kafka.net.socks5 import Socks5ProxyProtocol +def _drive(coro): + """Run a create_connection coroutine to completion off any real IO loop. + + The orchestration tests stub every await point (backend + do_connect), so a + throwaway asyncio loop is enough to drive create_connection deterministically. + """ + loop = asyncio.new_event_loop() + try: + return loop.run_until_complete(coro) + finally: + loop.close() + + +class _StubNet: + """Minimal NetBackend stand-in for exercising proxy.create_connection. + + Records getaddrinfo / create_connection calls and, like a real backend, + wires a transport onto the protocol via connection (here: sets _transport). + """ + + def __init__(self, addrs=()): + self.addrs = list(addrs) + self.getaddrinfo_calls = 0 + self.create_connection_calls = [] + self.transports = [] + + async def getaddrinfo(self, host, port): + self.getaddrinfo_calls += 1 + return self.addrs + + async def create_connection(self, protocol, host, port, *, ssl=None, + socket_options=(), timeout_at=None): + transport = MagicMock(name='proxy_transport') + self.create_connection_calls.append((protocol, host, port)) + self.transports.append(transport) + protocol._transport = transport + + +def _stub_do_connect(proxy, outcomes): + """Replace proxy.do_connect with a scripted sequence. + + Each outcome is either None (negotiation succeeded) or an exception to raise + (negotiation failed for that addrinfo). + """ + it = iter(outcomes) + + async def _do_connect(timeout_at=None): + outcome = next(it) + if isinstance(outcome, BaseException): + raise outcome + + proxy.do_connect = _do_connect + + +_IPV4_ADDR = (socket.AF_INET, socket.SOCK_STREAM, socket.IPPROTO_TCP, '', ('10.0.0.1', 9092)) +_IPV4_ADDR2 = (socket.AF_INET, socket.SOCK_STREAM, socket.IPPROTO_TCP, '', ('10.0.0.2', 9092)) + + class TestKafkaTCPProxyRegistry: def test_socks5(self, net): assert 'socks5' in KafkaTCPProxy._registry @@ -69,3 +130,150 @@ class _Second(KafkaTCPProxy): KafkaTCPProxy._registry.pop('test-dup', None) else: KafkaTCPProxy._registry['test-dup'] = prior + + +class TestCreateConnectionOrchestration: + """P3: exercise create_connection's DNS / retry / handoff wiring with the + per-proxy negotiation (do_connect) stubbed out.""" + + def test_remote_lookup_skips_dns_and_connects_to_proxy(self): + # http proxy uses remote lookup: no local getaddrinfo, connect to proxy. + net = _StubNet() + proxy = KafkaTCPProxy(net, 'http://proxy:8080') + _stub_do_connect(proxy, [None]) + protocol = MagicMock() + + _drive(proxy.create_connection(protocol, 'broker', 9092)) + + assert net.getaddrinfo_calls == 0 + assert net.create_connection_calls == [(proxy, 'proxy', 8080)] + # transport handed off to the real protocol, then released by the proxy + protocol.connection_made.assert_called_once_with(net.transports[0]) + assert proxy._transport is None + + def test_local_lookup_resolves_then_connects(self): + # socks5 (no 'h') resolves the broker locally before connecting. + net = _StubNet(addrs=[_IPV4_ADDR]) + proxy = KafkaTCPProxy(net, 'socks5://proxy:1080') + _stub_do_connect(proxy, [None]) + protocol = MagicMock() + + _drive(proxy.create_connection(protocol, 'broker', 9092)) + + assert net.getaddrinfo_calls == 1 + assert net.create_connection_calls == [(proxy, 'proxy', 1080)] + protocol.connection_made.assert_called_once_with(net.transports[0]) + + def test_retries_next_address_on_failure(self): + net = _StubNet(addrs=[_IPV4_ADDR, _IPV4_ADDR2]) + proxy = KafkaTCPProxy(net, 'socks5://proxy:1080') + # first addr fails negotiation, second succeeds + _stub_do_connect(proxy, [Errors.KafkaConnectionError('nope'), None]) + protocol = MagicMock() + + _drive(proxy.create_connection(protocol, 'broker', 9092)) + + assert len(net.create_connection_calls) == 2 + # handed off with the second (successful) transport + protocol.connection_made.assert_called_once_with(net.transports[1]) + + def test_all_addresses_fail_raises(self): + net = _StubNet(addrs=[_IPV4_ADDR, _IPV4_ADDR2]) + proxy = KafkaTCPProxy(net, 'socks5://proxy:1080') + _stub_do_connect(proxy, [Errors.KafkaConnectionError('nope'), + Errors.KafkaConnectionError('nope2')]) + protocol = MagicMock() + + with pytest.raises(Errors.KafkaConnectionError, match='Unable to connect'): + _drive(proxy.create_connection(protocol, 'broker', 9092)) + assert len(net.create_connection_calls) == 2 + protocol.connection_made.assert_not_called() + + def test_protocol_refusing_transport_aborts_and_reraises(self): + # KafkaConnection may refuse the transport (closed mid-connect) by + # raising from connection_made; the proxy must abort it and propagate. + net = _StubNet() + proxy = KafkaTCPProxy(net, 'http://proxy:8080') + _stub_do_connect(proxy, [None]) + boom = Errors.KafkaConnectionError('closed during connect') + protocol = MagicMock() + protocol.connection_made.side_effect = boom + + with pytest.raises(Errors.KafkaConnectionError, match='closed during connect'): + _drive(proxy.create_connection(protocol, 'broker', 9092)) + net.transports[0].abort.assert_called_once_with(boom) + + +class TestCreateConnectionSSL: + """P5: TLS is layered on top of the established proxy tunnel.""" + + @patch('kafka.net.proxy.KafkaSSLTransport') + def test_ssl_wraps_tunnel_before_handoff(self, mock_ssl_cls): + net = _StubNet() + proxy = KafkaTCPProxy(net, 'http://proxy:8080') + _stub_do_connect(proxy, [None]) + ssl_wrapper = mock_ssl_cls.return_value + ssl_wrapper.handshake = AsyncMock() + ssl_ctx = MagicMock(name='ssl_ctx') + protocol = MagicMock() + + _drive(proxy.create_connection(protocol, 'broker', 9092, ssl=ssl_ctx)) + + tunnel = net.transports[0] + # wrapper built for the broker host, layered onto the tunnel, handshaked + mock_ssl_cls.assert_called_once_with(net, ssl_ctx, host='broker') + ssl_wrapper.connection_made.assert_called_once_with(tunnel) + ssl_wrapper.handshake.assert_awaited_once() + # the real protocol receives the SSL transport, not the raw tunnel + protocol.connection_made.assert_called_once_with(ssl_wrapper) + assert proxy._transport is None + + def test_no_ssl_hands_off_raw_tunnel(self): + net = _StubNet() + proxy = KafkaTCPProxy(net, 'http://proxy:8080') + _stub_do_connect(proxy, [None]) + protocol = MagicMock() + + _drive(proxy.create_connection(protocol, 'broker', 9092, ssl=None)) + + protocol.connection_made.assert_called_once_with(net.transports[0]) + + +class TestConnectionLost: + """P4: a transport drop mid-negotiation must fail the pending connect + cleanly (regression for connection_lost calling the .failure setter, not + the .failed predicate).""" + + TARGET = (socket.AF_INET, socket.SOCK_STREAM, socket.IPPROTO_TCP, '', ('1.2.3.4', 9092)) + + def test_drop_mid_negotiation_fails_pending_connect(self, net): + proxy = KafkaTCPProxy(net, 'socks5://proxy:1080') + proxy.set_addrinfo(self.TARGET) + proxy.connection_made(MagicMock()) # negotiation started, future pending + assert not proxy._connect_future.is_done + + exc = Errors.KafkaConnectionError('proxy dropped') + proxy.connection_lost(exc) + + assert proxy._connect_future.failed() + assert proxy._connect_future.exception is exc + assert proxy._state == KafkaTCPProxyStates.DISCONNECTED + + def test_lost_before_connect_is_noop(self, net): + # _connect_future is None before connection_made; must not raise. + proxy = KafkaTCPProxy(net, 'socks5://proxy:1080') + proxy.connection_lost(Errors.KafkaConnectionError('early drop')) + assert proxy._state == KafkaTCPProxyStates.DISCONNECTED + + def test_lost_after_complete_does_not_overwrite(self, net): + proxy = KafkaTCPProxy(net, 'socks5://proxy:1080') + proxy.set_addrinfo(self.TARGET) + proxy.connection_made(MagicMock()) + proxy.data_received(b'\x05\x00') # no-auth selected -> request submitted + # grant: VER REP RSV ATYP=ipv4 + 4 addr + 2 port + proxy.data_received(b'\x05\x00\x00\x01' + b'\x7f\x00\x00\x01' + b'\x00\x00') + assert proxy._connect_future.succeeded() + + proxy.connection_lost(Errors.KafkaConnectionError('late drop')) + # already-resolved future is left untouched + assert proxy._connect_future.succeeded() From 444816fb87e1bdfe79746fb360eea65c5e2ac359 Mon Sep 17 00:00:00 2001 From: Dana Powers Date: Mon, 20 Jul 2026 10:04:32 -0700 Subject: [PATCH 6/6] remove backend/inet.py --- kafka/net/backend/inet.py | 100 -------------------------------------- 1 file changed, 100 deletions(-) delete mode 100644 kafka/net/backend/inet.py diff --git a/kafka/net/backend/inet.py b/kafka/net/backend/inet.py deleted file mode 100644 index 02d6adcbb..000000000 --- a/kafka/net/backend/inet.py +++ /dev/null @@ -1,100 +0,0 @@ -import errno -import logging -import socket -import time -from urllib.parse import urlparse - -import kafka.errors as Errors - - -log = logging.getLogger(__name__) - - -class KafkaNetSocket: - # scheme => handling class - _registry = {} - - @classmethod - def register_class(cls, klass): - for scheme in klass.SCHEMES: - cls._registry[scheme] = klass - - def __init_subclass__(cls, **kw): - super().__init_subclass__(**kw) - KafkaNetSocket.register_class(cls) - - def __new__(cls, proxy_url=None): - if proxy_url is None: - return super().__new__(cls) - try: - parsed = urlparse(proxy_url) - except Exception: - raise ValueError('Unable to parse proxy_url: %s' % (proxy_url,)) - if not parsed.scheme: - raise ValueError('proxy_url requires scheme:// (%s)' % (proxy_url,)) - try: - klass = KafkaNetSocket._registry[parsed.scheme] - except KeyError: - raise ValueError('Unsupported proxy url scheme: %s' % (parsed.scheme)) - return super().__new__(klass) - - def __init__(self, proxy_url=None): - pass - - # simple sockets / no proxy - def dns_lookup(self, host, port, raise_error=False): - # XXX: all DNS functions in Python are blocking. If we really - # want to be non-blocking here, we need to use a 3rd-party - # library like python-adns, or move resolution onto its - # own thread. This will be subject to the default libc - # name resolution timeout (5s on most Linux boxes) - try: - return socket.getaddrinfo(host, port, socket.AF_UNSPEC, socket.SOCK_STREAM) - except socket.gaierror as ex: - err_str = "DNS lookup failed for %s:%d, %r" % (host, port, ex) - if not raise_error: - log.warning(err_str) - return [] - raise Errors.KafkaConnectionError(err_str) - - def socket(self, family=socket.AF_UNSPEC, sock_type=socket.SOCK_STREAM, proto=socket.IPPROTO_TCP): - return socket.socket(family, sock_type, proto) - - async def connect(self, net, addrinfo, socket_options=(), timeout_at=None): - """Create non-blocking socket (with options) and connect to addrinfo tuple""" - family, sock_type, proto, _canonname, sockaddr = addrinfo - sock = self.socket(family, sock_type, proto) - sock.setblocking(False) - for option in socket_options: - sock.setsockopt(*option) - return await self.sock_connect(net, sock, sockaddr, timeout_at=timeout_at) - - async def sock_connect(self, net, sock, sockaddr, timeout_at=None): - while timeout_at is None or time.monotonic() < timeout_at: - ret = None - try: - ret = self.connect_ex(sock, sockaddr) - except BlockingIOError: - ret = errno.EWOULDBLOCK - except socket.error as err: - ret = err.errno - - # Connection succeeded - if not ret or ret == errno.EISCONN: - log.debug('Connected: %s', sock) - return sock - - # Needs retry - # WSAEINVAL == 10022, but errno.WSAEINVAL is not available on non-win systems - elif ret in (errno.EINPROGRESS, errno.EALREADY, errno.EWOULDBLOCK, 10022): - await net.wait_write(sock, timeout_at=timeout_at) - - # Connection failed - else: - errstr = errno.errorcode.get(ret, 'UNKNOWN') - raise Errors.KafkaConnectionError('{} {}'.format(ret, errstr)) - else: - raise Errors.KafkaTimeoutError('Connection timed out') - - def connect_ex(self, sock, sockaddr): - return sock.connect_ex(sockaddr)