diff --git a/docs/apidoc/modules.rst b/docs/apidoc/modules.rst index f65b55647..b9c410439 100644 --- a/docs/apidoc/modules.rst +++ b/docs/apidoc/modules.rst @@ -93,10 +93,8 @@ driving the protocol layer directly from the REPL. - :mod:`~kafka.net.connection` - per-broker async connection: state machine, request/response correlation, and SASL handshake. -- :mod:`~kafka.net.transport` - Async socket I/O with write buffering, +- :mod:`~kafka.net.backend.transport` - Async socket I/O with write buffering, pause/resume hooks, and the asyncio-shaped protocol callback surface. -- :mod:`~kafka.net.inet` - DNS lookup + non-blocking connect, plus a - URL-scheme registry that resolves ``proxy_url`` to socket factories. - :mod:`~kafka.net.http_connect` - Tunnels broker connections through an HTTP CONNECT proxy (RFC 7231). - :mod:`~kafka.net.socks5` - SOCKS5 client with optional username/password @@ -108,8 +106,7 @@ driving the protocol layer directly from the REPL. manager connection - transport - inet + transport http_connect socks5 diff --git a/docs/apidoc/net/transport.rst b/docs/apidoc/net/backend/transport.rst similarity index 100% rename from docs/apidoc/net/transport.rst rename to docs/apidoc/net/backend/transport.rst diff --git a/docs/apidoc/net/inet.rst b/docs/apidoc/net/inet.rst deleted file mode 100644 index 9c051645f..000000000 --- a/docs/apidoc/net/inet.rst +++ /dev/null @@ -1,10 +0,0 @@ -kafka.net.inet -============== - -.. module:: kafka.net.inet - -.. autoclass:: kafka.net.inet.KafkaNetSocket - :members: - :undoc-members: - -.. autofunction:: kafka.net.inet.create_connection diff --git a/kafka/future.py b/kafka/future.py index f09a00526..e3fe807aa 100644 --- a/kafka/future.py +++ b/kafka/future.py @@ -15,7 +15,7 @@ class Future: ``is_done`` / ``value`` / ``exception`` state) and is safe to resolve from any thread. It is deliberately **not** awaitable: awaiting happens only on the event loop, via the backend's loop-awaitable future from - ``net.create_future()`` (``kafka.net.selector.SelectorFuture`` for the + ``net.create_future()`` (``kafka.net.backend.selector.SelectorFuture`` for the selector), which subclasses ``Future`` and adds ``__await__``. Keeping ``__await__`` off the base makes the invariant type-enforced -- awaiting a plain handoff ``Future`` raises immediately rather than silently working on diff --git a/kafka/net/__init__.py b/kafka/net/__init__.py index f7cc11384..51bf318ca 100644 --- a/kafka/net/__init__.py +++ b/kafka/net/__init__.py @@ -1,19 +1,16 @@ from .connection import KafkaConnection -from .inet import create_connection, KafkaNetSocket from .manager import KafkaConnectionManager from .metrics import KafkaConnectionMetrics, KafkaManagerMetrics -from .selector import NetworkSelector from .http_connect import HttpConnectProxy from .socks5 import Socks5Proxy -from .transport import KafkaTCPTransport, KafkaSSLTransport from .wakeup_notifier import WakeupNotifier from .compat import KafkaNetClient __all__ = [ - 'KafkaConnection', 'create_connection', 'KafkaNetSocket', - 'KafkaConnectionManager', 'KafkaConnectionMetrics', 'KafkaManagerMetrics', - 'NetworkSelector', 'HttpConnectProxy', 'Socks5Proxy', 'KafkaTCPTransport', 'KafkaSSLTransport', + 'KafkaConnection', 'KafkaConnectionManager', + 'KafkaConnectionMetrics', 'KafkaManagerMetrics', + 'HttpConnectProxy', 'Socks5Proxy', 'WakeupNotifier', 'KafkaNetClient', ] diff --git a/kafka/net/backend/__init__.py b/kafka/net/backend/__init__.py new file mode 100644 index 000000000..fedff3730 --- /dev/null +++ b/kafka/net/backend/__init__.py @@ -0,0 +1,7 @@ +from .abstract import ( + NetBackend, NetTransport, NetProtocol, NetBackendFuture, + resolve_backend, register_backend_lazy, +) + +register_backend_lazy('selector', 'kafka.net.backend.selector', 'NetworkSelector') +register_backend_lazy('asyncio', 'kafka.net.backend.asyncio_backend', 'AsyncioBackend') diff --git a/kafka/net/backend.py b/kafka/net/backend/abstract.py similarity index 96% rename from kafka/net/backend.py rename to kafka/net/backend/abstract.py index 3843d46a5..04521bb94 100644 --- a/kafka/net/backend.py +++ b/kafka/net/backend/abstract.py @@ -5,13 +5,13 @@ ``KafkaAdminClient`` (and the manager, cluster, connection, coordinator, fetcher, sender) reach for through ``self._net`` / ``manager._net``. -``NetworkSelector`` (``kafka/net/selector.py``) is the reference +``NetworkSelector`` (``kafka/net/backend/selector.py``) is the reference implementation; an asyncio backend (and eventually Twisted) implements the same surface so it can be swapped in via ``net=`` without touching core code. The :class:`NetBackendFuture` contract is the surface of the loop-awaitable futures a backend hands out from ``net.create_future()``. The -selector's implementation is ``kafka.net.selector.SelectorFuture``; an asyncio +selector's implementation is ``kafka.net.backend.selector.SelectorFuture``; an asyncio (and eventually Twisted) backend supplies its own. Networking is a **connection seam**, not fd-readiness. asyncio and Twisted own @@ -27,7 +27,7 @@ * ``wait_read`` / ``wait_write`` / ``unregister_event`` -- the low-level fd-readiness primitives. They are the *selector's* private mechanism (used - only inside ``kafka/net/transport.py`` + ``inet.py``, zero core callers) and + only inside ``kafka/net/backend/transport.py`` + ``inet.py``, zero core callers) and do not port to asyncio/Twisted. The connection seam replaces them. * ``poll(timeout_ms, future=...)`` -- the legacy single-tick driver. Its only remaining caller is the ``KafkaNetClient`` compat shim @@ -54,7 +54,7 @@ class NetBackendFuture(Protocol): """Contract for the awaitable futures returned by ``net.create_future()``. - A pluggable async backend (the kafka.net selector, asyncio, Twisted, ...) + A pluggable async backend (the kafka.net.backend selector, asyncio, Twisted, ...) returns its own future type from ``create_future()``. Core loop coroutines touch it only through this surface, so the type is interchangeable across backends. The selector's ``SelectorFuture`` is the reference implementation: @@ -235,7 +235,7 @@ def wakeup(self) -> None: # --- backend selection ---------------------------------------------------- # name -> factory(**config) -> NetBackend. Populated by register_backend(); -# 'selector' is always available, 'asyncio' registers itself in Step 4. +# 'selector' and 'asyncio' are lazily registered in kafka/net/backend/__init__.py. _BACKENDS = {} @@ -307,7 +307,3 @@ def resolve_backend(net, config): if name is None or name not in _BACKENDS: name = 'selector' return _BACKENDS[name](**config) - - -register_backend_lazy('selector', 'kafka.net.selector', 'NetworkSelector') -register_backend_lazy('asyncio', 'kafka.net.asyncio_backend', 'AsyncioBackend') diff --git a/kafka/net/asyncio_backend.py b/kafka/net/backend/asyncio_backend.py similarity index 100% rename from kafka/net/asyncio_backend.py rename to kafka/net/backend/asyncio_backend.py diff --git a/kafka/net/inet.py b/kafka/net/backend/inet.py similarity index 100% rename from kafka/net/inet.py rename to kafka/net/backend/inet.py diff --git a/kafka/net/selector.py b/kafka/net/backend/selector.py similarity index 99% rename from kafka/net/selector.py rename to kafka/net/backend/selector.py index 6520a7c6e..ecb7b3f58 100644 --- a/kafka/net/selector.py +++ b/kafka/net/backend/selector.py @@ -11,8 +11,8 @@ import kafka.errors as Errors from kafka.future import Future -from kafka.net.inet import create_connection as _inet_create_connection -from kafka.net.transport import KafkaSSLTransport, KafkaTCPTransport +from kafka.net.backend.inet import create_connection as _inet_create_connection +from kafka.net.backend.transport import KafkaSSLTransport, KafkaTCPTransport from kafka.version import __version__ diff --git a/kafka/net/transport.py b/kafka/net/backend/transport.py similarity index 100% rename from kafka/net/transport.py rename to kafka/net/backend/transport.py diff --git a/kafka/net/http_connect.py b/kafka/net/http_connect.py index ec646730a..959b3b71d 100644 --- a/kafka/net/http_connect.py +++ b/kafka/net/http_connect.py @@ -6,7 +6,7 @@ from urllib.parse import urlparse from kafka.errors import KafkaConnectionError -from kafka.net.inet import KafkaNetSocket +from kafka.net.backend.inet import KafkaNetSocket log = logging.getLogger(__name__) diff --git a/kafka/net/manager.py b/kafka/net/manager.py index 2dd824465..64eec5a02 100644 --- a/kafka/net/manager.py +++ b/kafka/net/manager.py @@ -10,7 +10,7 @@ from kafka.net.backend import resolve_backend from kafka.cluster import ClusterMetadata import kafka.errors as Errors -from kafka.net.transport import KafkaSSLTransport +from kafka.net.backend.transport import KafkaSSLTransport from kafka.net.wakeup_notifier import WakeupNotifier from kafka.protocol.broker_version_data import BrokerVersionData from kafka.version import __version__ diff --git a/kafka/net/socks5.py b/kafka/net/socks5.py index 20ee590c3..cb4ab5a17 100644 --- a/kafka/net/socks5.py +++ b/kafka/net/socks5.py @@ -6,7 +6,7 @@ from urllib.parse import urlparse from kafka.errors import KafkaConnectionError -from kafka.net.inet import KafkaNetSocket +from kafka.net.backend.inet import KafkaNetSocket log = logging.getLogger(__name__) diff --git a/test/conftest.py b/test/conftest.py index a0040543d..111af8d5b 100644 --- a/test/conftest.py +++ b/test/conftest.py @@ -6,7 +6,7 @@ from kafka.cluster import ClusterMetadata from kafka.net.compat import KafkaNetClient from kafka.net.manager import KafkaConnectionManager -from kafka.net.selector import NetworkSelector +from kafka.net.backend.selector import NetworkSelector from kafka.protocol.metadata import MetadataResponse diff --git a/test/net/test_backend.py b/test/net/backend/test_abstract.py similarity index 95% rename from test/net/test_backend.py rename to test/net/backend/test_abstract.py index d7a92f19d..9c15cb0f6 100644 --- a/test/net/test_backend.py +++ b/test/net/backend/test_abstract.py @@ -1,4 +1,4 @@ -"""Conformance tests for the NetBackend contract (kafka/net/backend.py). +"""Conformance tests for the NetBackend contract (kafka/net/backend/abstract.py). NetworkSelector is the reference implementation; these pin that it satisfies the NetBackend Protocol structurally and that the shared lifecycle helper @@ -10,11 +10,11 @@ import pytest -from kafka.net.backend import ( +from kafka.net.backend.abstract import ( NetBackend, NetTransport, resolve_backend, register_backend, _BACKENDS, ) -from kafka.net.selector import NetworkSelector -from kafka.net.transport import KafkaTCPTransport +from kafka.net.backend.selector import NetworkSelector +from kafka.net.backend.transport import KafkaTCPTransport # The full contract surface, kept here so a missing/renamed method fails loudly. @@ -120,7 +120,7 @@ def test_unknown_name_raises(self): def test_asyncio_name_resolves(self): # net='asyncio' lazily imports + registers the asyncio backend. - from kafka.net.asyncio_backend import AsyncioBackend + from kafka.net.backend.asyncio_backend import AsyncioBackend b = resolve_backend('asyncio', {'client_id': 'x'}) assert isinstance(b, AsyncioBackend) b.close() @@ -156,7 +156,7 @@ async def main(): def test_autodetect_asyncio_in_loop_returns_asyncio_backend(self): # In a running asyncio loop with no explicit net, auto-detect lazily # registers + selects the asyncio backend (Phase-1: still own thread). - from kafka.net.asyncio_backend import AsyncioBackend + from kafka.net.backend.asyncio_backend import AsyncioBackend async def main(): return resolve_backend(None, {'client_id': 'auto'}) @@ -168,7 +168,7 @@ async def main(): def test_autodetect_falls_back_for_unknown_framework(self, monkeypatch): # A detected-but-unregistered framework (e.g. trio, no backend) falls # back to the default selector rather than erroring. - import kafka.net.backend as backend_mod + import kafka.net.backend.abstract as backend_mod monkeypatch.setattr(backend_mod, '_detect_async_library', lambda: 'trio') assert isinstance(resolve_backend(None, {}), NetworkSelector) diff --git a/test/net/test_asyncio_backend.py b/test/net/backend/test_asyncio_backend.py similarity index 97% rename from test/net/test_asyncio_backend.py rename to test/net/backend/test_asyncio_backend.py index b86205832..d59b916eb 100644 --- a/test/net/test_asyncio_backend.py +++ b/test/net/backend/test_asyncio_backend.py @@ -1,9 +1,9 @@ -"""Tests for the asyncio NetBackend (kafka/net/asyncio_backend.py). +"""Tests for the asyncio NetBackend (kafka/net/backend/asyncio_backend.py). Covers backend-specific behavior (lifecycle, timers, cross-thread run), reuses the shared NetBackendFuture conformance suite against the asyncio-backed future, and drives a real protocol round-trip through a MockBroker on a started -AsyncioBackend -- the both-backends coverage for the async paths. +AsyncioBackend -- the both-backend coverage for the async paths. """ import asyncio import socket @@ -13,12 +13,12 @@ import pytest import kafka.errors as Errors -from kafka.net.asyncio_backend import AsyncioBackend, AsyncioFuture from kafka.net.backend import NetBackend +from kafka.net.backend.asyncio_backend import AsyncioBackend, AsyncioFuture from kafka.net.manager import KafkaConnectionManager from kafka.protocol.metadata import MetadataRequest from test.mock_broker import MockBroker -from test.net.test_net_backend_future import NetBackendFutureContract +from test.net.backend.test_net_backend_future import NetBackendFutureContract @pytest.fixture diff --git a/test/net/test_inet.py b/test/net/backend/test_inet.py similarity index 80% rename from test/net/test_inet.py rename to test/net/backend/test_inet.py index 87cad1682..12dd3f569 100644 --- a/test/net/test_inet.py +++ b/test/net/backend/test_inet.py @@ -4,8 +4,7 @@ import pytest -from kafka.net.selector import NetworkSelector -from kafka.net.inet import create_connection, KafkaNetSocket +from kafka.net.backend.inet import create_connection, KafkaNetSocket from kafka.net.socks5 import Socks5Proxy from kafka.net.http_connect import HttpConnectProxy import kafka.errors as Errors @@ -19,7 +18,7 @@ def test_valid_host(self): assert len(res) == 5 def test_invalid_host(self): - with patch('kafka.net.inet.socket.getaddrinfo', side_effect=socket.gaierror): + with patch('kafka.net.backend.inet.socket.getaddrinfo', side_effect=socket.gaierror): results = KafkaNetSocket().dns_lookup('invalid.host', 9092) assert results == [] @@ -30,8 +29,7 @@ def test_numeric_host(self): class TestSockConnect: - def test_immediate_connect(self): - net = NetworkSelector() + def test_immediate_connect(self, net): factory = KafkaNetSocket() sock = MagicMock() sock.connect_ex.return_value = 0 @@ -39,34 +37,30 @@ def test_immediate_connect(self): assert result is sock sock.connect_ex.assert_called_once_with(('127.0.0.1', 9092)) - def test_eisconn(self): - net = NetworkSelector() + def test_eisconn(self, net): factory = KafkaNetSocket() sock = MagicMock() sock.connect_ex.return_value = errno.EISCONN result = net.run(factory.sock_connect(net, sock, ('127.0.0.1', 9092))) assert result is sock - def test_connection_refused(self): - net = NetworkSelector() + def test_connection_refused(self, net): factory = KafkaNetSocket() sock = MagicMock() sock.connect_ex.return_value = errno.ECONNREFUSED with pytest.raises(Errors.KafkaConnectionError): net.run(factory.sock_connect(net, sock, ('127.0.0.1', 9092))) - def test_socket_error_uses_errno(self): - net = NetworkSelector() + def test_socket_error_uses_errno(self, net): factory = KafkaNetSocket() sock = MagicMock() sock.connect_ex.side_effect = socket.error(errno.ECONNREFUSED, 'refused') with pytest.raises(Errors.KafkaConnectionError): net.run(factory.sock_connect(net, sock, ('127.0.0.1', 9092))) - def test_error_after_wait_write(self): + def test_error_after_wait_write(self, net): """connect_ex returns EINPROGRESS, then after wait_write fires the second connect_ex returns the real error.""" - net = NetworkSelector() factory = KafkaNetSocket() # socketpair endpoints are always immediately writable, so wait_write # fires on the first poll and we re-enter the loop. @@ -86,34 +80,30 @@ def test_error_after_wait_write(self): class TestCreateConnection: - def test_dns_failure(self): - net = NetworkSelector() - with patch('kafka.net.inet.KafkaNetSocket.dns_lookup', return_value=[]): + def test_dns_failure(self, net): + with patch('kafka.net.backend.inet.KafkaNetSocket.dns_lookup', return_value=[]): with pytest.raises(Errors.KafkaConnectionError, match='DNS'): net.run(create_connection(net, 'badhost', 9092)) - def test_socket_init_failure(self): - net = NetworkSelector() + def test_socket_init_failure(self, net): fake_addr = [(socket.AF_INET, socket.SOCK_STREAM, 6, '', ('127.0.0.1', 9092))] - with patch('kafka.net.inet.KafkaNetSocket.dns_lookup', return_value=fake_addr), \ - patch('kafka.net.inet.socket.socket', side_effect=OSError('no socket')): + with patch('kafka.net.backend.inet.KafkaNetSocket.dns_lookup', return_value=fake_addr), \ + patch('kafka.net.backend.inet.socket.socket', side_effect=OSError('no socket')): with pytest.raises(Errors.KafkaConnectionError): net.run(create_connection(net, 'host', 9092)) - def test_successful_connection(self): - net = NetworkSelector() + def test_successful_connection(self, net): fake_addr = [(socket.AF_INET, socket.SOCK_STREAM, 6, '', ('127.0.0.1', 9092))] mock_sock = MagicMock() mock_sock.connect_ex.return_value = 0 - with patch('kafka.net.inet.KafkaNetSocket.dns_lookup', return_value=fake_addr), \ - patch('kafka.net.inet.socket.socket', return_value=mock_sock): + with patch('kafka.net.backend.inet.KafkaNetSocket.dns_lookup', return_value=fake_addr), \ + patch('kafka.net.backend.inet.socket.socket', return_value=mock_sock): result = net.run( create_connection(net, 'host', 9092)) assert result is mock_sock mock_sock.setblocking.assert_called_with(False) - def test_tries_multiple_addresses(self): - net = NetworkSelector() + def test_tries_multiple_addresses(self, net): addr1 = (socket.AF_INET, socket.SOCK_STREAM, 6, '', ('10.0.0.1', 9092)) addr2 = (socket.AF_INET, socket.SOCK_STREAM, 6, '', ('10.0.0.2', 9092)) mock_sock1 = MagicMock() @@ -121,14 +111,13 @@ def test_tries_multiple_addresses(self): mock_sock2 = MagicMock() mock_sock2.connect_ex.return_value = 0 sockets = iter([mock_sock1, mock_sock2]) - with patch('kafka.net.inet.KafkaNetSocket.dns_lookup', return_value=[addr1, addr2]), \ - patch('kafka.net.inet.socket.socket', side_effect=lambda *a: next(sockets)): + with patch('kafka.net.backend.inet.KafkaNetSocket.dns_lookup', return_value=[addr1, addr2]), \ + patch('kafka.net.backend.inet.socket.socket', side_effect=lambda *a: next(sockets)): result = net.run( create_connection(net, 'host', 9092)) assert result is mock_sock2 - def test_socket_options_applied(self): - net = NetworkSelector() + def test_socket_options_applied(self, net): fake_addr = [(socket.AF_INET, socket.SOCK_STREAM, 6, '', ('127.0.0.1', 9092))] mock_sock = MagicMock() mock_sock.connect_ex.return_value = 0 @@ -136,8 +125,8 @@ def test_socket_options_applied(self): (socket.IPPROTO_TCP, socket.TCP_NODELAY, 1), (socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1), ] - with patch('kafka.net.inet.KafkaNetSocket.dns_lookup', return_value=fake_addr), \ - patch('kafka.net.inet.socket.socket', return_value=mock_sock): + with patch('kafka.net.backend.inet.KafkaNetSocket.dns_lookup', return_value=fake_addr), \ + patch('kafka.net.backend.inet.socket.socket', return_value=mock_sock): net.run(create_connection(net, 'host', 9092, socket_options=opts)) mock_sock.setsockopt.assert_has_calls([ call(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1), @@ -147,52 +136,48 @@ def test_socket_options_applied(self): class TestCreateConnectionWithProxy: - def test_proxy_creates_socket(self): - net = NetworkSelector() + def test_proxy_creates_socket(self, net): mock_sock = MagicMock() mock_sock.connect_ex.return_value = 0 fake_addr = (socket.AF_INET, socket.SOCK_STREAM, 6, '', ('127.0.0.1', 9092)) with patch('kafka.net.socks5.Socks5Proxy._get_proxy_addr'), \ - patch('kafka.net.inet.KafkaNetSocket.dns_lookup', return_value=[fake_addr]), \ + patch('kafka.net.backend.inet.KafkaNetSocket.dns_lookup', return_value=[fake_addr]), \ patch('kafka.net.socks5.Socks5Proxy.connect', return_value=mock_sock) as mock_connect: result = net.run( create_connection(net, 'broker', 9092, proxy_url='socks5://proxy:1080')) mock_connect.assert_called_once_with(net, fake_addr, (), timeout_at=None) assert result is mock_sock - def test_proxy_remote_dns_skips_local_lookup(self): - net = NetworkSelector() + def test_proxy_remote_dns_skips_local_lookup(self, net): mock_sock = MagicMock() mock_sock.connect_ex.return_value = 0 with patch('kafka.net.socks5.Socks5Proxy._get_proxy_addr'), \ patch('kafka.net.socks5.Socks5Proxy.socket', return_value=mock_sock), \ patch('kafka.net.socks5.Socks5Proxy.connect_ex', return_value=0), \ - patch('kafka.net.inet.KafkaNetSocket.dns_lookup') as mock_dns: + patch('kafka.net.backend.inet.KafkaNetSocket.dns_lookup') as mock_dns: result = net.run( create_connection(net, 'broker', 9092, proxy_url='socks5h://proxy:1080')) mock_dns.assert_not_called() - def test_no_proxy_uses_direct_socket(self): - net = NetworkSelector() + def test_no_proxy_uses_direct_socket(self, net): fake_addr = (socket.AF_INET, socket.SOCK_STREAM, 6, '', ('127.0.0.1', 9092)) mock_sock = MagicMock() mock_sock.connect_ex.return_value = 0 - with patch('kafka.net.inet.KafkaNetSocket.dns_lookup', return_value=[fake_addr]), \ - patch('kafka.net.inet.socket.socket', return_value=mock_sock), \ + with patch('kafka.net.backend.inet.KafkaNetSocket.dns_lookup', return_value=[fake_addr]), \ + patch('kafka.net.backend.inet.socket.socket', return_value=mock_sock), \ patch('kafka.net.socks5.Socks5Proxy.connect') as mock_connect: result = net.run( create_connection(net, 'host', 9092)) mock_connect.assert_not_called() assert result is mock_sock - def test_socks5h_does_dns_for_proxy_not_target(self): + def test_socks5h_does_dns_for_proxy_not_target(self, net): """Companion to test_proxy_remote_dns_skips_local_lookup: with _get_proxy_addr running normally, exactly one dns_lookup is made and it is for the proxy hostname, not the target.""" - net = NetworkSelector() proxy_addr = (socket.AF_INET, socket.SOCK_STREAM, 6, '', ('1.2.3.4', 1080)) mock_sock = MagicMock() - with patch('kafka.net.inet.KafkaNetSocket.dns_lookup', return_value=[proxy_addr]) as mock_dns, \ + with patch('kafka.net.backend.inet.KafkaNetSocket.dns_lookup', return_value=[proxy_addr]) as mock_dns, \ patch('kafka.net.socks5.Socks5Proxy.socket', return_value=mock_sock), \ patch('kafka.net.socks5.Socks5Proxy.connect_ex', return_value=0): net.run( @@ -201,23 +186,22 @@ def test_socks5h_does_dns_for_proxy_not_target(self): assert mock_dns.call_args.args[:2] == ('proxy', 1080) def test_socks5_proxy_dns_gaierror_raises(self): - with patch('kafka.net.inet.socket.getaddrinfo', side_effect=socket.gaierror): + with patch('kafka.net.backend.inet.socket.getaddrinfo', side_effect=socket.gaierror): with pytest.raises(Errors.KafkaConnectionError): KafkaNetSocket('socks5://bogus.proxy:1080') def test_socks5_proxy_dns_empty_raises(self): - with patch('kafka.net.inet.KafkaNetSocket.dns_lookup', return_value=[]): + with patch('kafka.net.backend.inet.KafkaNetSocket.dns_lookup', return_value=[]): with pytest.raises(Errors.KafkaConnectionError): KafkaNetSocket('socks5://proxy:1080') - def test_proxy_connect_dispatches_through_inherited_connect(self): + def test_proxy_connect_dispatches_through_inherited_connect(self, net): """create_connection -> Socks5Proxy.connect (inherited from KafkaNetSocket) -> Socks5Proxy.socket + Socks5Proxy.connect_ex.""" - net = NetworkSelector() fake_addr = (socket.AF_INET, socket.SOCK_STREAM, 6, '', ('127.0.0.1', 9092)) mock_sock = MagicMock() with patch('kafka.net.socks5.Socks5Proxy._get_proxy_addr'), \ - patch('kafka.net.inet.KafkaNetSocket.dns_lookup', return_value=[fake_addr]), \ + patch('kafka.net.backend.inet.KafkaNetSocket.dns_lookup', return_value=[fake_addr]), \ patch('kafka.net.socks5.Socks5Proxy.socket', return_value=mock_sock) as mock_socket, \ patch('kafka.net.socks5.Socks5Proxy.connect_ex', return_value=0) as mock_connect_ex: result = net.run( @@ -306,7 +290,7 @@ class TestKafkaNetSocketExtensionPattern: to an existing asyncio socket implementation). """ - def test_connect_ex_only_subclass(self): + def test_connect_ex_only_subclass(self, net): """An HTTP CONNECT-style handler that only overrides connect_ex.""" ex_calls = [] @@ -319,11 +303,10 @@ def connect_ex(self, sock, sockaddr): return 0 try: - net = NetworkSelector() fake_addr = (socket.AF_INET, socket.SOCK_STREAM, 6, '', ('10.0.0.1', 9092)) mock_sock = MagicMock() - with patch('kafka.net.inet.KafkaNetSocket.dns_lookup', return_value=[fake_addr]), \ - patch('kafka.net.inet.socket.socket', return_value=mock_sock): + with patch('kafka.net.backend.inet.KafkaNetSocket.dns_lookup', return_value=[fake_addr]), \ + patch('kafka.net.backend.inet.socket.socket', return_value=mock_sock): result = net.run( create_connection(net, 'broker', 9092, proxy_url='test-httpconnect://proxy:8080')) @@ -333,7 +316,7 @@ def connect_ex(self, sock, sockaddr): finally: KafkaNetSocket._registry.pop('test-httpconnect', None) - def test_connect_override_subclass(self): + def test_connect_override_subclass(self, net): """An asyncio-style handler that overrides connect() entirely; the default socket()/sock_connect()/connect_ex() flow is bypassed.""" connect_calls = [] @@ -347,10 +330,9 @@ async def connect(self, net, addrinfo, socket_options=(), timeout_at=None): return 'asyncio-stream-handle' try: - net = NetworkSelector() fake_addr = (socket.AF_INET, socket.SOCK_STREAM, 6, '', ('10.0.0.1', 9092)) - with patch('kafka.net.inet.KafkaNetSocket.dns_lookup', return_value=[fake_addr]), \ - patch('kafka.net.inet.socket.socket') as mock_sock_cls: + with patch('kafka.net.backend.inet.KafkaNetSocket.dns_lookup', return_value=[fake_addr]), \ + patch('kafka.net.backend.inet.socket.socket') as mock_sock_cls: result = net.run( create_connection(net, 'broker', 9092, proxy_url='test-asyncio://x', diff --git a/test/net/test_net_backend_future.py b/test/net/backend/test_net_backend_future.py similarity index 99% rename from test/net/test_net_backend_future.py rename to test/net/backend/test_net_backend_future.py index faa2dc463..8cdc1529c 100644 --- a/test/net/test_net_backend_future.py +++ b/test/net/backend/test_net_backend_future.py @@ -12,7 +12,7 @@ from kafka.future import Future from kafka.net.backend import NetBackendFuture -from kafka.net.selector import NetworkSelector +from kafka.net.backend.selector import NetworkSelector class NetBackendFutureContract: diff --git a/test/net/test_selector.py b/test/net/backend/test_selector.py similarity index 98% rename from test/net/test_selector.py rename to test/net/backend/test_selector.py index dd52e62e5..bc1c61a4b 100644 --- a/test/net/test_selector.py +++ b/test/net/backend/test_selector.py @@ -7,7 +7,7 @@ from kafka.errors import KafkaTimeoutError from kafka.future import Future -from kafka.net.selector import ( +from kafka.net.backend.selector import ( KernelEvent, NetworkSelector, Task, @@ -813,7 +813,7 @@ async def hog(): done.success(True) net.call_soon(hog) - with caplog.at_level('WARNING', logger='kafka.net.selector'): + with caplog.at_level('WARNING', logger='kafka.net.backend.selector'): net.poll(timeout_ms=1000, future=done) assert any('blocking the event loop' in rec.message for rec in caplog.records), ( 'expected slow-task warning, got: %r' @@ -828,7 +828,7 @@ async def quick(): done.success(True) net.call_soon(quick) - with caplog.at_level('WARNING', logger='kafka.net.selector'): + with caplog.at_level('WARNING', logger='kafka.net.backend.selector'): net.poll(timeout_ms=1000, future=done) assert not any('blocking the event loop' in rec.message for rec in caplog.records) @@ -841,7 +841,7 @@ async def hog(): done.success(True) net.call_soon(hog) - with caplog.at_level('WARNING', logger='kafka.net.selector'): + with caplog.at_level('WARNING', logger='kafka.net.backend.selector'): net.poll(timeout_ms=1000, future=done) assert not any('blocking the event loop' in rec.message for rec in caplog.records) @@ -1073,7 +1073,7 @@ async def work(): release = threading.Event() try: self._wedge(net, release) - with caplog.at_level('WARNING', logger='kafka.net.selector'): + with caplog.at_level('WARNING', logger='kafka.net.backend.selector'): th, outcome = self._run_in_thread(net, work) assert isinstance(outcome.get('exc'), KafkaTimeoutError) assert any('did not complete within' in r.message and 'work' in r.message @@ -1097,7 +1097,7 @@ async def work(): release = threading.Event() try: self._wedge(net, release) - with caplog.at_level('WARNING', logger='kafka.net.selector'): + with caplog.at_level('WARNING', logger='kafka.net.backend.selector'): th, outcome = self._run_in_thread(net, work) # caller times out assert isinstance(outcome.get('exc'), KafkaTimeoutError) # Release the wedge so the abandoned coroutine now completes. diff --git a/test/net/test_transport.py b/test/net/backend/test_transport.py similarity index 99% rename from test/net/test_transport.py rename to test/net/backend/test_transport.py index 849eacfc6..b746a6f8e 100644 --- a/test/net/test_transport.py +++ b/test/net/backend/test_transport.py @@ -7,8 +7,8 @@ import kafka.errors as Errors from kafka.future import Future -from kafka.net.selector import NetworkSelector, TaskState -from kafka.net.transport import KafkaSSLTransport, KafkaTCPTransport +from kafka.net.backend.selector import NetworkSelector, TaskState +from kafka.net.backend.transport import KafkaSSLTransport, KafkaTCPTransport @pytest.fixture diff --git a/test/net/test_connection.py b/test/net/test_connection.py index 0ab54103b..fe185686d 100644 --- a/test/net/test_connection.py +++ b/test/net/test_connection.py @@ -6,20 +6,14 @@ import pytest from kafka.future import Future -from kafka.net.selector import NetworkSelector from kafka.net.connection import KafkaConnection -from kafka.net.transport import KafkaTCPTransport +from kafka.net.backend.transport import KafkaTCPTransport from kafka.protocol.broker_version_data import BrokerVersionData from kafka.protocol.metadata import ApiVersionsRequest from kafka.protocol.parser import KafkaProtocol import kafka.errors as Errors -@pytest.fixture -def net(): - return NetworkSelector() - - @pytest.fixture def connection(net): return KafkaConnection(net, node_id='test-0') diff --git a/test/net/test_http_connect.py b/test/net/test_http_connect.py index a1df08595..19a654090 100644 --- a/test/net/test_http_connect.py +++ b/test/net/test_http_connect.py @@ -5,7 +5,7 @@ import pytest from kafka.net.http_connect import HttpConnectProxy -from kafka.net.inet import KafkaNetSocket +from kafka.net.backend.inet import KafkaNetSocket _FAKE_PROXY_ADDR = (socket.AF_INET, socket.SOCK_STREAM, socket.IPPROTO_TCP, '', ('1.2.3.4', 8080)) diff --git a/test/net/test_manager.py b/test/net/test_manager.py index adaf71f2c..ca3816cb4 100644 --- a/test/net/test_manager.py +++ b/test/net/test_manager.py @@ -7,18 +7,13 @@ from kafka.cluster import ClusterMetadata from kafka.future import Future -from kafka.net.selector import NetworkSelector -from kafka.net.manager import KafkaConnectionManager +from kafka.net.backend.selector import NetworkSelector from kafka.net.connection import KafkaConnection +from kafka.net.manager import KafkaConnectionManager import kafka.errors as Errors from kafka.protocol.broker_version_data import BrokerVersionData -@pytest.fixture -def net(): - return NetworkSelector() - - class TestKafkaConnectionManagerNetResolution: """Backend selection lives on the manager (not the compat shim).""" @@ -541,7 +536,7 @@ def test_run_survives_gc_during_poll(self, manager, monkeypatch): for as long as the wrapper Future is pending. """ import gc - from kafka.net.selector import NetworkSelector + from kafka.net.backend.selector import NetworkSelector # Force a GC cycle on every _poll_once entry to deterministically # trigger the orphan-collection race that was masking timeouts in CI. diff --git a/test/net/test_sasl_reauthentication.py b/test/net/test_sasl_reauthentication.py index e0c6f4615..a9a93d84a 100644 --- a/test/net/test_sasl_reauthentication.py +++ b/test/net/test_sasl_reauthentication.py @@ -6,7 +6,6 @@ import kafka.errors as Errors from kafka.net.connection import KafkaConnection from kafka.net.manager import KafkaConnectionManager -from kafka.net.selector import NetworkSelector from kafka.protocol.sasl import ( SaslAuthenticateRequest, SaslAuthenticateResponse, @@ -25,15 +24,6 @@ } -@pytest.fixture -def net(): - sel = NetworkSelector() - try: - yield sel - finally: - sel.close() - - @pytest.fixture def sasl_broker(): return MockBroker(broker_version=(2, 5)) # supports SaslAuthenticate v0-2 diff --git a/test/net/test_wakeup_notifier.py b/test/net/test_wakeup_notifier.py index 7554e5fd6..134d0e00c 100644 --- a/test/net/test_wakeup_notifier.py +++ b/test/net/test_wakeup_notifier.py @@ -15,15 +15,9 @@ import pytest from kafka.future import Future -from kafka.net.selector import NetworkSelector from kafka.net.wakeup_notifier import WakeupNotifier -@pytest.fixture -def net(): - return NetworkSelector() - - @pytest.fixture def notifier(net): return WakeupNotifier(net) @@ -168,7 +162,7 @@ async def task(): assert elapsed < 0.5, ( 'second cycle should wake immediately; took %.3fs' % elapsed) - def test_no_lost_wakeup_under_concurrent_notify_stress(self): + def test_no_lost_wakeup_under_concurrent_notify_stress(self, net): """Probabilistic regression guard for the coalescing path: a consumer coroutine awaits the notifier every iteration (max race exposure) and drains a shared queue; many cross-thread producers append work and @@ -181,7 +175,6 @@ def test_no_lost_wakeup_under_concurrent_notify_stress(self): that. With the latch + coalescing correct, it finishes in milliseconds. """ import collections - net = NetworkSelector() net.start() try: notifier = WakeupNotifier(net)