Skip to content
Merged
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
6 changes: 3 additions & 3 deletions kafka/net/__init__.py
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -11,6 +11,6 @@
__all__ = [
'KafkaConnection', 'KafkaConnectionManager',
'KafkaConnectionMetrics', 'KafkaManagerMetrics',
'HttpConnectProxy', 'Socks5Proxy',
'HttpConnectProxyProtocol', 'Socks5ProxyProtocol',
'WakeupNotifier', 'KafkaNetClient',
]
4 changes: 1 addition & 3 deletions kafka/net/backend/abstract.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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 ---------------------------------------------
Expand Down
6 changes: 1 addition & 5 deletions kafka/net/backend/asyncio_backend.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
100 changes: 0 additions & 100 deletions kafka/net/backend/inet.py

This file was deleted.

3 changes: 1 addition & 2 deletions kafka/net/backend/selector.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
142 changes: 27 additions & 115 deletions kafka/net/http_connect.py
Original file line number Diff line number Diff line change
@@ -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 = '<disconnected>'
CONNECTING = '<connecting>'
SENDING = '<sending>'
READING = '<reading>'
COMPLETE = '<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'``
Expand All @@ -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,))
20 changes: 14 additions & 6 deletions kafka/net/manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand Down
Loading