From 2672ac2b792b85f1705cda186f317a17a1bda8c9 Mon Sep 17 00:00:00 2001 From: Dana Powers Date: Wed, 15 Jul 2026 09:35:04 -0700 Subject: [PATCH 01/10] net.backend -> net.backends.abstract --- kafka/admin/client.py | 2 +- kafka/consumer/group.py | 2 +- kafka/future.py | 2 +- kafka/net/asyncio_backend.py | 2 +- kafka/net/backends/__init__.py | 7 +++++++ kafka/net/{backend.py => backends/abstract.py} | 4 ---- kafka/net/compat.py | 2 +- kafka/net/manager.py | 2 +- kafka/net/selector.py | 4 ++-- kafka/producer/kafka.py | 2 +- test/net/test_asyncio_backend.py | 2 +- test/net/test_backend.py | 6 +++--- test/net/test_net_backend_future.py | 2 +- 13 files changed, 21 insertions(+), 18 deletions(-) create mode 100644 kafka/net/backends/__init__.py rename kafka/net/{backend.py => backends/abstract.py} (98%) diff --git a/kafka/admin/client.py b/kafka/admin/client.py index f51b7682c..8bda276c8 100644 --- a/kafka/admin/client.py +++ b/kafka/admin/client.py @@ -147,7 +147,7 @@ class KafkaAdminClient( metadata before the configured timeout. Note that bootstrap is called eagerly from __init__(). Default: 30000 - net (str or kafka.net.backend.NetBackend): The async backend that runs + net (str or kafka.net.backends.NetBackend): The async backend that runs this client's network I/O event loop. One of: a NetBackend instance; a registered name -- 'selector' (the built-in NetworkSelector) or 'asyncio' (runs I/O on an asyncio loop); or None diff --git a/kafka/consumer/group.py b/kafka/consumer/group.py index 3804147ab..687b0bf15 100644 --- a/kafka/consumer/group.py +++ b/kafka/consumer/group.py @@ -270,7 +270,7 @@ class KafkaConsumer: metrics. Default: 2 metrics_sample_window_ms (int): The maximum age in milliseconds of samples used to compute metrics. Default: 30000 - net (str or kafka.net.backend.NetBackend): The async backend that runs + net (str or kafka.net.backends.NetBackend): The async backend that runs this client's network I/O event loop. One of: a NetBackend instance; a registered name -- 'selector' (the built-in NetworkSelector) or 'asyncio' (runs I/O on an asyncio loop); or None diff --git a/kafka/future.py b/kafka/future.py index f09a00526..a673c2fcb 100644 --- a/kafka/future.py +++ b/kafka/future.py @@ -19,7 +19,7 @@ class Future: 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 - one backend and breaking on another. See ``kafka.net.backend.NetBackendFuture`` + one backend and breaking on another. See ``kafka.net.backends.NetBackendFuture`` for the awaitable contract. """ __slots__ = ('is_done', 'value', 'exception', '_callbacks', '_errbacks', '_lock') diff --git a/kafka/net/asyncio_backend.py b/kafka/net/asyncio_backend.py index 15ee91de8..ee63b5681 100644 --- a/kafka/net/asyncio_backend.py +++ b/kafka/net/asyncio_backend.py @@ -4,7 +4,7 @@ ``NetworkSelector``'s threading model, and implements the ``NetBackend`` contract on top of asyncio primitives. Selected via ``net='asyncio'`` or auto-detected when constructed inside a running asyncio loop (see -``kafka.net.backend.resolve_backend``). +``kafka.net.backends.resolve_backend``). Phase 1 preserves the synchronous public API: ``run()`` blocks the calling thread on the loop thread; it does not run on the caller's own loop. diff --git a/kafka/net/backends/__init__.py b/kafka/net/backends/__init__.py new file mode 100644 index 000000000..bf39e3b3c --- /dev/null +++ b/kafka/net/backends/__init__.py @@ -0,0 +1,7 @@ +from .abstract import ( + NetBackend, NetTransport, NetProtocol, NetBackendFuture, + resolve_backend, register_backend_lazy, +) + +register_backend_lazy('selector', 'kafka.net.selector', 'NetworkSelector') +register_backend_lazy('asyncio', 'kafka.net.asyncio_backend', 'AsyncioBackend') diff --git a/kafka/net/backend.py b/kafka/net/backends/abstract.py similarity index 98% rename from kafka/net/backend.py rename to kafka/net/backends/abstract.py index 3843d46a5..a95eadaba 100644 --- a/kafka/net/backend.py +++ b/kafka/net/backends/abstract.py @@ -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/compat.py b/kafka/net/compat.py index 19c03e223..6f4c7ba52 100644 --- a/kafka/net/compat.py +++ b/kafka/net/compat.py @@ -4,7 +4,7 @@ import time import kafka.errors as Errors -from kafka.net.backend import resolve_backend +from kafka.net.backends import resolve_backend from kafka.net.manager import KafkaConnectionManager from kafka.util import Timer diff --git a/kafka/net/manager.py b/kafka/net/manager.py index 2dd824465..6dcbf0298 100644 --- a/kafka/net/manager.py +++ b/kafka/net/manager.py @@ -7,7 +7,7 @@ from .connection import KafkaConnection from .metrics import KafkaManagerMetrics -from kafka.net.backend import resolve_backend +from kafka.net.backends import resolve_backend from kafka.cluster import ClusterMetadata import kafka.errors as Errors from kafka.net.transport import KafkaSSLTransport diff --git a/kafka/net/selector.py b/kafka/net/selector.py index 6520a7c6e..f39576557 100644 --- a/kafka/net/selector.py +++ b/kafka/net/selector.py @@ -43,7 +43,7 @@ def _initialize_coro(maybe_coro): class SelectorFuture(Future): - """The NetworkSelector's loop-awaitable future (see backend.NetBackendFuture). + """The NetworkSelector's loop-awaitable future (see backends.NetBackendFuture). ``kafka.future.Future`` is the thread-safe callback/handoff core with no ``__await__``; ``SelectorFuture`` adds it: ``yield self`` suspends the @@ -550,7 +550,7 @@ def reschedule(self, when, task): return task def create_future(self): - """Create a loop-awaitable future (see backend.NetBackendFuture). + """Create a loop-awaitable future (see backends.NetBackendFuture). Portability seam for pluggable backends: core coroutines call this instead of constructing ``Future`` directly, so an alternate backend diff --git a/kafka/producer/kafka.py b/kafka/producer/kafka.py index 1a3bee1a9..69728c259 100644 --- a/kafka/producer/kafka.py +++ b/kafka/producer/kafka.py @@ -380,7 +380,7 @@ class KafkaProducer: metrics. Default: 2 metrics_sample_window_ms (int): The maximum age in milliseconds of samples used to compute metrics. Default: 30000 - net (str or kafka.net.backend.NetBackend): The async backend that runs + net (str or kafka.net.backends.NetBackend): The async backend that runs this client's network I/O event loop. One of: a NetBackend instance; a registered name -- 'selector' (the built-in NetworkSelector) or 'asyncio' (runs I/O on an asyncio loop); or None diff --git a/test/net/test_asyncio_backend.py b/test/net/test_asyncio_backend.py index b86205832..3bcadfab9 100644 --- a/test/net/test_asyncio_backend.py +++ b/test/net/test_asyncio_backend.py @@ -14,7 +14,7 @@ import kafka.errors as Errors from kafka.net.asyncio_backend import AsyncioBackend, AsyncioFuture -from kafka.net.backend import NetBackend +from kafka.net.backends import NetBackend from kafka.net.manager import KafkaConnectionManager from kafka.protocol.metadata import MetadataRequest from test.mock_broker import MockBroker diff --git a/test/net/test_backend.py b/test/net/test_backend.py index d7a92f19d..0b4474fb4 100644 --- a/test/net/test_backend.py +++ b/test/net/test_backend.py @@ -1,4 +1,4 @@ -"""Conformance tests for the NetBackend contract (kafka/net/backend.py). +"""Conformance tests for the NetBackend contract (kafka/net/backends/abstract.py). NetworkSelector is the reference implementation; these pin that it satisfies the NetBackend Protocol structurally and that the shared lifecycle helper @@ -10,7 +10,7 @@ import pytest -from kafka.net.backend import ( +from kafka.net.backends.abstract import ( NetBackend, NetTransport, resolve_backend, register_backend, _BACKENDS, ) from kafka.net.selector import NetworkSelector @@ -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.backends.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_net_backend_future.py b/test/net/test_net_backend_future.py index faa2dc463..9548930e7 100644 --- a/test/net/test_net_backend_future.py +++ b/test/net/test_net_backend_future.py @@ -11,7 +11,7 @@ import pytest from kafka.future import Future -from kafka.net.backend import NetBackendFuture +from kafka.net.backends import NetBackendFuture from kafka.net.selector import NetworkSelector From 2ddf849e86a1f60369cb4e9861e22b99659100bb Mon Sep 17 00:00:00 2001 From: Dana Powers Date: Wed, 15 Jul 2026 10:05:36 -0700 Subject: [PATCH 02/10] asyncio_backend -> net/backends/ --- kafka/net/backends/__init__.py | 2 +- kafka/net/{ => backends}/asyncio_backend.py | 0 test/net/{ => backends}/test_asyncio_backend.py | 4 ++-- test/net/test_backend.py | 4 ++-- 4 files changed, 5 insertions(+), 5 deletions(-) rename kafka/net/{ => backends}/asyncio_backend.py (100%) rename test/net/{ => backends}/test_asyncio_backend.py (98%) diff --git a/kafka/net/backends/__init__.py b/kafka/net/backends/__init__.py index bf39e3b3c..5df3ec545 100644 --- a/kafka/net/backends/__init__.py +++ b/kafka/net/backends/__init__.py @@ -4,4 +4,4 @@ ) register_backend_lazy('selector', 'kafka.net.selector', 'NetworkSelector') -register_backend_lazy('asyncio', 'kafka.net.asyncio_backend', 'AsyncioBackend') +register_backend_lazy('asyncio', 'kafka.net.backends.asyncio_backend', 'AsyncioBackend') diff --git a/kafka/net/asyncio_backend.py b/kafka/net/backends/asyncio_backend.py similarity index 100% rename from kafka/net/asyncio_backend.py rename to kafka/net/backends/asyncio_backend.py diff --git a/test/net/test_asyncio_backend.py b/test/net/backends/test_asyncio_backend.py similarity index 98% rename from test/net/test_asyncio_backend.py rename to test/net/backends/test_asyncio_backend.py index 3bcadfab9..cbf1c031d 100644 --- a/test/net/test_asyncio_backend.py +++ b/test/net/backends/test_asyncio_backend.py @@ -1,4 +1,4 @@ -"""Tests for the asyncio NetBackend (kafka/net/asyncio_backend.py). +"""Tests for the asyncio NetBackend (kafka/net/backends/asyncio_backend.py). Covers backend-specific behavior (lifecycle, timers, cross-thread run), reuses the shared NetBackendFuture conformance suite against the asyncio-backed future, @@ -13,8 +13,8 @@ import pytest import kafka.errors as Errors -from kafka.net.asyncio_backend import AsyncioBackend, AsyncioFuture from kafka.net.backends import NetBackend +from kafka.net.backends.asyncio_backend import AsyncioBackend, AsyncioFuture from kafka.net.manager import KafkaConnectionManager from kafka.protocol.metadata import MetadataRequest from test.mock_broker import MockBroker diff --git a/test/net/test_backend.py b/test/net/test_backend.py index 0b4474fb4..03f0a18ec 100644 --- a/test/net/test_backend.py +++ b/test/net/test_backend.py @@ -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.backends.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.backends.asyncio_backend import AsyncioBackend async def main(): return resolve_backend(None, {'client_id': 'auto'}) From 26c54626830b597d5e7e0399d4b810956f4fdfcc Mon Sep 17 00:00:00 2001 From: Dana Powers Date: Wed, 15 Jul 2026 10:15:09 -0700 Subject: [PATCH 03/10] selector -> net/backends/ --- kafka/future.py | 2 +- kafka/net/__init__.py | 3 +- kafka/net/backends/__init__.py | 2 +- kafka/net/backends/abstract.py | 6 +-- kafka/net/{ => backends}/selector.py | 0 test/conftest.py | 2 +- test/net/{ => backends}/test_selector.py | 2 +- test/net/test_backend.py | 2 +- test/net/test_connection.py | 6 --- test/net/test_inet.py | 52 ++++++++---------------- test/net/test_manager.py | 11 ++--- test/net/test_net_backend_future.py | 2 +- test/net/test_sasl_reauthentication.py | 10 ----- test/net/test_transport.py | 2 +- test/net/test_wakeup_notifier.py | 9 +--- 15 files changed, 32 insertions(+), 79 deletions(-) rename kafka/net/{ => backends}/selector.py (100%) rename test/net/{ => backends}/test_selector.py (99%) diff --git a/kafka/future.py b/kafka/future.py index a673c2fcb..15c75ec5a 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.backends.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..6c16c0fce 100644 --- a/kafka/net/__init__.py +++ b/kafka/net/__init__.py @@ -2,7 +2,6 @@ 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 @@ -14,6 +13,6 @@ __all__ = [ 'KafkaConnection', 'create_connection', 'KafkaNetSocket', 'KafkaConnectionManager', 'KafkaConnectionMetrics', 'KafkaManagerMetrics', - 'NetworkSelector', 'HttpConnectProxy', 'Socks5Proxy', 'KafkaTCPTransport', 'KafkaSSLTransport', + 'HttpConnectProxy', 'Socks5Proxy', 'KafkaTCPTransport', 'KafkaSSLTransport', 'WakeupNotifier', 'KafkaNetClient', ] diff --git a/kafka/net/backends/__init__.py b/kafka/net/backends/__init__.py index 5df3ec545..4dc51aacf 100644 --- a/kafka/net/backends/__init__.py +++ b/kafka/net/backends/__init__.py @@ -3,5 +3,5 @@ resolve_backend, register_backend_lazy, ) -register_backend_lazy('selector', 'kafka.net.selector', 'NetworkSelector') +register_backend_lazy('selector', 'kafka.net.backends.selector', 'NetworkSelector') register_backend_lazy('asyncio', 'kafka.net.backends.asyncio_backend', 'AsyncioBackend') diff --git a/kafka/net/backends/abstract.py b/kafka/net/backends/abstract.py index a95eadaba..e3f2dff13 100644 --- a/kafka/net/backends/abstract.py +++ b/kafka/net/backends/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/backends/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.backends.selector.SelectorFuture``; an asyncio (and eventually Twisted) backend supplies its own. Networking is a **connection seam**, not fd-readiness. asyncio and Twisted own @@ -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.backends 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: diff --git a/kafka/net/selector.py b/kafka/net/backends/selector.py similarity index 100% rename from kafka/net/selector.py rename to kafka/net/backends/selector.py diff --git a/test/conftest.py b/test/conftest.py index a0040543d..590271e3d 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.backends.selector import NetworkSelector from kafka.protocol.metadata import MetadataResponse diff --git a/test/net/test_selector.py b/test/net/backends/test_selector.py similarity index 99% rename from test/net/test_selector.py rename to test/net/backends/test_selector.py index dd52e62e5..0b249eced 100644 --- a/test/net/test_selector.py +++ b/test/net/backends/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.backends.selector import ( KernelEvent, NetworkSelector, Task, diff --git a/test/net/test_backend.py b/test/net/test_backend.py index 03f0a18ec..d682f6914 100644 --- a/test/net/test_backend.py +++ b/test/net/test_backend.py @@ -13,7 +13,7 @@ from kafka.net.backends.abstract import ( NetBackend, NetTransport, resolve_backend, register_backend, _BACKENDS, ) -from kafka.net.selector import NetworkSelector +from kafka.net.backends.selector import NetworkSelector from kafka.net.transport import KafkaTCPTransport diff --git a/test/net/test_connection.py b/test/net/test_connection.py index 0ab54103b..84f45c831 100644 --- a/test/net/test_connection.py +++ b/test/net/test_connection.py @@ -6,7 +6,6 @@ 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.protocol.broker_version_data import BrokerVersionData @@ -15,11 +14,6 @@ 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_inet.py b/test/net/test_inet.py index 87cad1682..438564a99 100644 --- a/test/net/test_inet.py +++ b/test/net/test_inet.py @@ -4,7 +4,6 @@ import pytest -from kafka.net.selector import NetworkSelector from kafka.net.inet import create_connection, KafkaNetSocket from kafka.net.socks5 import Socks5Proxy from kafka.net.http_connect import HttpConnectProxy @@ -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,22 +80,19 @@ def test_error_after_wait_write(self): class TestCreateConnection: - def test_dns_failure(self): - net = NetworkSelector() + def test_dns_failure(self, net): with patch('kafka.net.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 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 @@ -112,8 +103,7 @@ def test_successful_connection(self): 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() @@ -127,8 +117,7 @@ def test_tries_multiple_addresses(self): 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 @@ -147,8 +136,7 @@ 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)) @@ -160,8 +148,7 @@ def test_proxy_creates_socket(self): 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'), \ @@ -172,8 +159,7 @@ def test_proxy_remote_dns_skips_local_lookup(self): 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 @@ -185,11 +171,10 @@ def test_no_proxy_uses_direct_socket(self): 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, \ @@ -210,10 +195,9 @@ def test_socks5_proxy_dns_empty_raises(self): 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'), \ @@ -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,7 +303,6 @@ 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]), \ @@ -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,7 +330,6 @@ 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: diff --git a/test/net/test_manager.py b/test/net/test_manager.py index adaf71f2c..650e4dede 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.backends.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.backends.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_net_backend_future.py b/test/net/test_net_backend_future.py index 9548930e7..0cf48632a 100644 --- a/test/net/test_net_backend_future.py +++ b/test/net/test_net_backend_future.py @@ -12,7 +12,7 @@ from kafka.future import Future from kafka.net.backends import NetBackendFuture -from kafka.net.selector import NetworkSelector +from kafka.net.backends.selector import NetworkSelector class NetBackendFutureContract: 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_transport.py b/test/net/test_transport.py index 849eacfc6..855d79b93 100644 --- a/test/net/test_transport.py +++ b/test/net/test_transport.py @@ -7,7 +7,7 @@ import kafka.errors as Errors from kafka.future import Future -from kafka.net.selector import NetworkSelector, TaskState +from kafka.net.backends.selector import NetworkSelector, TaskState from kafka.net.transport import KafkaSSLTransport, KafkaTCPTransport 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) From 3556f8cafb16f1ce2bf04c52a359417771159711 Mon Sep 17 00:00:00 2001 From: Dana Powers Date: Wed, 15 Jul 2026 10:16:57 -0700 Subject: [PATCH 04/10] test_backend -> test/net/backends/ --- test/net/{test_backend.py => backends/test_abstract.py} | 0 test/net/backends/test_asyncio_backend.py | 2 +- test/net/{ => backends}/test_net_backend_future.py | 0 3 files changed, 1 insertion(+), 1 deletion(-) rename test/net/{test_backend.py => backends/test_abstract.py} (100%) rename test/net/{ => backends}/test_net_backend_future.py (100%) diff --git a/test/net/test_backend.py b/test/net/backends/test_abstract.py similarity index 100% rename from test/net/test_backend.py rename to test/net/backends/test_abstract.py diff --git a/test/net/backends/test_asyncio_backend.py b/test/net/backends/test_asyncio_backend.py index cbf1c031d..b0c1dc110 100644 --- a/test/net/backends/test_asyncio_backend.py +++ b/test/net/backends/test_asyncio_backend.py @@ -18,7 +18,7 @@ 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.backends.test_net_backend_future import NetBackendFutureContract @pytest.fixture diff --git a/test/net/test_net_backend_future.py b/test/net/backends/test_net_backend_future.py similarity index 100% rename from test/net/test_net_backend_future.py rename to test/net/backends/test_net_backend_future.py From 4fc0a088440878f94b70b3dcc7f02a27dde1ea92 Mon Sep 17 00:00:00 2001 From: Dana Powers Date: Wed, 15 Jul 2026 10:24:36 -0700 Subject: [PATCH 05/10] inet -> net/backends/ --- kafka/net/__init__.py | 5 ++- kafka/net/{ => backends}/inet.py | 0 kafka/net/backends/selector.py | 2 +- kafka/net/http_connect.py | 2 +- kafka/net/socks5.py | 2 +- test/net/{ => backends}/test_inet.py | 46 ++++++++++++++-------------- test/net/test_http_connect.py | 2 +- 7 files changed, 29 insertions(+), 30 deletions(-) rename kafka/net/{ => backends}/inet.py (100%) rename test/net/{ => backends}/test_inet.py (86%) diff --git a/kafka/net/__init__.py b/kafka/net/__init__.py index 6c16c0fce..5bedc53f9 100644 --- a/kafka/net/__init__.py +++ b/kafka/net/__init__.py @@ -1,5 +1,4 @@ from .connection import KafkaConnection -from .inet import create_connection, KafkaNetSocket from .manager import KafkaConnectionManager from .metrics import KafkaConnectionMetrics, KafkaManagerMetrics from .http_connect import HttpConnectProxy @@ -11,8 +10,8 @@ __all__ = [ - 'KafkaConnection', 'create_connection', 'KafkaNetSocket', - 'KafkaConnectionManager', 'KafkaConnectionMetrics', 'KafkaManagerMetrics', + 'KafkaConnection', 'KafkaConnectionManager', + 'KafkaConnectionMetrics', 'KafkaManagerMetrics', 'HttpConnectProxy', 'Socks5Proxy', 'KafkaTCPTransport', 'KafkaSSLTransport', 'WakeupNotifier', 'KafkaNetClient', ] diff --git a/kafka/net/inet.py b/kafka/net/backends/inet.py similarity index 100% rename from kafka/net/inet.py rename to kafka/net/backends/inet.py diff --git a/kafka/net/backends/selector.py b/kafka/net/backends/selector.py index f39576557..9dc366f9c 100644 --- a/kafka/net/backends/selector.py +++ b/kafka/net/backends/selector.py @@ -11,7 +11,7 @@ import kafka.errors as Errors from kafka.future import Future -from kafka.net.inet import create_connection as _inet_create_connection +from kafka.net.backends.inet import create_connection as _inet_create_connection from kafka.net.transport import KafkaSSLTransport, KafkaTCPTransport from kafka.version import __version__ diff --git a/kafka/net/http_connect.py b/kafka/net/http_connect.py index ec646730a..ff9b65a15 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.backends.inet import KafkaNetSocket log = logging.getLogger(__name__) diff --git a/kafka/net/socks5.py b/kafka/net/socks5.py index 20ee590c3..1794bbb00 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.backends.inet import KafkaNetSocket log = logging.getLogger(__name__) diff --git a/test/net/test_inet.py b/test/net/backends/test_inet.py similarity index 86% rename from test/net/test_inet.py rename to test/net/backends/test_inet.py index 438564a99..d8913d9e3 100644 --- a/test/net/test_inet.py +++ b/test/net/backends/test_inet.py @@ -4,7 +4,7 @@ import pytest -from kafka.net.inet import create_connection, KafkaNetSocket +from kafka.net.backends.inet import create_connection, KafkaNetSocket from kafka.net.socks5 import Socks5Proxy from kafka.net.http_connect import HttpConnectProxy import kafka.errors as Errors @@ -18,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.backends.inet.socket.getaddrinfo', side_effect=socket.gaierror): results = KafkaNetSocket().dns_lookup('invalid.host', 9092) assert results == [] @@ -81,14 +81,14 @@ def test_error_after_wait_write(self, net): class TestCreateConnection: def test_dns_failure(self, net): - with patch('kafka.net.inet.KafkaNetSocket.dns_lookup', return_value=[]): + with patch('kafka.net.backends.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): 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.backends.inet.KafkaNetSocket.dns_lookup', return_value=fake_addr), \ + patch('kafka.net.backends.inet.socket.socket', side_effect=OSError('no socket')): with pytest.raises(Errors.KafkaConnectionError): net.run(create_connection(net, 'host', 9092)) @@ -96,8 +96,8 @@ 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.backends.inet.KafkaNetSocket.dns_lookup', return_value=fake_addr), \ + patch('kafka.net.backends.inet.socket.socket', return_value=mock_sock): result = net.run( create_connection(net, 'host', 9092)) assert result is mock_sock @@ -111,8 +111,8 @@ def test_tries_multiple_addresses(self, net): 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.backends.inet.KafkaNetSocket.dns_lookup', return_value=[addr1, addr2]), \ + patch('kafka.net.backends.inet.socket.socket', side_effect=lambda *a: next(sockets)): result = net.run( create_connection(net, 'host', 9092)) assert result is mock_sock2 @@ -125,8 +125,8 @@ def test_socket_options_applied(self, net): (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.backends.inet.KafkaNetSocket.dns_lookup', return_value=fake_addr), \ + patch('kafka.net.backends.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), @@ -141,7 +141,7 @@ def test_proxy_creates_socket(self, net): 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.backends.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')) @@ -154,7 +154,7 @@ def test_proxy_remote_dns_skips_local_lookup(self, net): 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.backends.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() @@ -163,8 +163,8 @@ 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.backends.inet.KafkaNetSocket.dns_lookup', return_value=[fake_addr]), \ + patch('kafka.net.backends.inet.socket.socket', return_value=mock_sock), \ patch('kafka.net.socks5.Socks5Proxy.connect') as mock_connect: result = net.run( create_connection(net, 'host', 9092)) @@ -177,7 +177,7 @@ def test_socks5h_does_dns_for_proxy_not_target(self, net): it is for the proxy hostname, not the target.""" 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.backends.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( @@ -186,12 +186,12 @@ def test_socks5h_does_dns_for_proxy_not_target(self, net): 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.backends.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.backends.inet.KafkaNetSocket.dns_lookup', return_value=[]): with pytest.raises(Errors.KafkaConnectionError): KafkaNetSocket('socks5://proxy:1080') @@ -201,7 +201,7 @@ def test_proxy_connect_dispatches_through_inherited_connect(self, net): 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.backends.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( @@ -305,8 +305,8 @@ def connect_ex(self, sock, sockaddr): try: 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.backends.inet.KafkaNetSocket.dns_lookup', return_value=[fake_addr]), \ + patch('kafka.net.backends.inet.socket.socket', return_value=mock_sock): result = net.run( create_connection(net, 'broker', 9092, proxy_url='test-httpconnect://proxy:8080')) @@ -331,8 +331,8 @@ async def connect(self, net, addrinfo, socket_options=(), timeout_at=None): try: 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.backends.inet.KafkaNetSocket.dns_lookup', return_value=[fake_addr]), \ + patch('kafka.net.backends.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_http_connect.py b/test/net/test_http_connect.py index a1df08595..cc415b2e4 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.backends.inet import KafkaNetSocket _FAKE_PROXY_ADDR = (socket.AF_INET, socket.SOCK_STREAM, socket.IPPROTO_TCP, '', ('1.2.3.4', 8080)) From bc9834f49b06438cf375c01e0fad77b72e3b2bf7 Mon Sep 17 00:00:00 2001 From: Dana Powers Date: Wed, 15 Jul 2026 10:25:23 -0700 Subject: [PATCH 06/10] drop kafka.net.inet from apidocs --- docs/apidoc/modules.rst | 3 --- docs/apidoc/net/inet.rst | 10 ---------- 2 files changed, 13 deletions(-) delete mode 100644 docs/apidoc/net/inet.rst diff --git a/docs/apidoc/modules.rst b/docs/apidoc/modules.rst index f65b55647..d5640d821 100644 --- a/docs/apidoc/modules.rst +++ b/docs/apidoc/modules.rst @@ -95,8 +95,6 @@ driving the protocol layer directly from the REPL. handshake. - :mod:`~kafka.net.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 @@ -109,7 +107,6 @@ driving the protocol layer directly from the REPL. manager connection transport - inet http_connect socks5 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 From 0030d3d755dbb2677cffb0f795b8862631ba4d29 Mon Sep 17 00:00:00 2001 From: Dana Powers Date: Wed, 15 Jul 2026 10:28:36 -0700 Subject: [PATCH 07/10] net.backends -> net.backend --- kafka/admin/client.py | 2 +- kafka/consumer/group.py | 2 +- kafka/future.py | 4 +- kafka/net/backend/__init__.py | 7 +++ kafka/net/{backends => backend}/abstract.py | 12 ++--- .../{backends => backend}/asyncio_backend.py | 2 +- kafka/net/{backends => backend}/inet.py | 0 kafka/net/{backends => backend}/selector.py | 10 ++-- kafka/net/backends/__init__.py | 7 --- kafka/net/compat.py | 2 +- kafka/net/http_connect.py | 2 +- kafka/net/manager.py | 2 +- kafka/net/socks5.py | 2 +- kafka/producer/kafka.py | 2 +- test/conftest.py | 2 +- .../{backends => backend}/test_abstract.py | 12 ++--- .../test_asyncio_backend.py | 10 ++-- test/net/{backends => backend}/test_inet.py | 46 +++++++++---------- .../test_net_backend_future.py | 4 +- .../{backends => backend}/test_selector.py | 2 +- test/net/test_http_connect.py | 2 +- test/net/test_manager.py | 4 +- test/net/test_transport.py | 2 +- 23 files changed, 70 insertions(+), 70 deletions(-) create mode 100644 kafka/net/backend/__init__.py rename kafka/net/{backends => backend}/abstract.py (96%) rename kafka/net/{backends => backend}/asyncio_backend.py (99%) rename kafka/net/{backends => backend}/inet.py (100%) rename kafka/net/{backends => backend}/selector.py (98%) delete mode 100644 kafka/net/backends/__init__.py rename test/net/{backends => backend}/test_abstract.py (95%) rename test/net/{backends => backend}/test_asyncio_backend.py (97%) rename test/net/{backends => backend}/test_inet.py (86%) rename test/net/{backends => backend}/test_net_backend_future.py (98%) rename test/net/{backends => backend}/test_selector.py (99%) diff --git a/kafka/admin/client.py b/kafka/admin/client.py index 8bda276c8..f51b7682c 100644 --- a/kafka/admin/client.py +++ b/kafka/admin/client.py @@ -147,7 +147,7 @@ class KafkaAdminClient( metadata before the configured timeout. Note that bootstrap is called eagerly from __init__(). Default: 30000 - net (str or kafka.net.backends.NetBackend): The async backend that runs + net (str or kafka.net.backend.NetBackend): The async backend that runs this client's network I/O event loop. One of: a NetBackend instance; a registered name -- 'selector' (the built-in NetworkSelector) or 'asyncio' (runs I/O on an asyncio loop); or None diff --git a/kafka/consumer/group.py b/kafka/consumer/group.py index 687b0bf15..3804147ab 100644 --- a/kafka/consumer/group.py +++ b/kafka/consumer/group.py @@ -270,7 +270,7 @@ class KafkaConsumer: metrics. Default: 2 metrics_sample_window_ms (int): The maximum age in milliseconds of samples used to compute metrics. Default: 30000 - net (str or kafka.net.backends.NetBackend): The async backend that runs + net (str or kafka.net.backend.NetBackend): The async backend that runs this client's network I/O event loop. One of: a NetBackend instance; a registered name -- 'selector' (the built-in NetworkSelector) or 'asyncio' (runs I/O on an asyncio loop); or None diff --git a/kafka/future.py b/kafka/future.py index 15c75ec5a..e3fe807aa 100644 --- a/kafka/future.py +++ b/kafka/future.py @@ -15,11 +15,11 @@ 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.backends.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 - one backend and breaking on another. See ``kafka.net.backends.NetBackendFuture`` + one backend and breaking on another. See ``kafka.net.backend.NetBackendFuture`` for the awaitable contract. """ __slots__ = ('is_done', 'value', 'exception', '_callbacks', '_errbacks', '_lock') 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/backends/abstract.py b/kafka/net/backend/abstract.py similarity index 96% rename from kafka/net/backends/abstract.py rename to kafka/net/backend/abstract.py index e3f2dff13..197604e04 100644 --- a/kafka/net/backends/abstract.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/backends/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.backends.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 @@ -54,16 +54,16 @@ class NetBackendFuture(Protocol): """Contract for the awaitable futures returned by ``net.create_future()``. - A pluggable async backend (the kafka.net.backends 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: + backend. The selector's ``SelectorFuture`` is the reference implementation: it subclasses the thread-safe ``kafka.future.Future`` (the portable callback core) and adds ``__await__``. A plain ``Future`` is deliberately NOT a NetBackendFuture -- it has no ``__await__`` -- so awaiting a cross-thread handoff future fails loudly instead of silently working on one backend. - Pinned semantics -- the three axes where backends could otherwise diverge: + Pinned semantics -- the three axes where backend could otherwise diverge: 1. **Resolution thread.** A future from ``create_future()`` is created and resolved (``success`` / ``failure``) on the loop/IO thread only. @@ -302,7 +302,7 @@ def resolve_backend(net, config): 'or None; got %r' % (net,)) return net # net is None: auto-detect, else default. Auto-detected-but-unregistered - # backends fall back silently (an explicit name would have raised above). + # backend fall back silently (an explicit name would have raised above). name = _detect_async_library() if name is None or name not in _BACKENDS: name = 'selector' diff --git a/kafka/net/backends/asyncio_backend.py b/kafka/net/backend/asyncio_backend.py similarity index 99% rename from kafka/net/backends/asyncio_backend.py rename to kafka/net/backend/asyncio_backend.py index ee63b5681..15ee91de8 100644 --- a/kafka/net/backends/asyncio_backend.py +++ b/kafka/net/backend/asyncio_backend.py @@ -4,7 +4,7 @@ ``NetworkSelector``'s threading model, and implements the ``NetBackend`` contract on top of asyncio primitives. Selected via ``net='asyncio'`` or auto-detected when constructed inside a running asyncio loop (see -``kafka.net.backends.resolve_backend``). +``kafka.net.backend.resolve_backend``). Phase 1 preserves the synchronous public API: ``run()`` blocks the calling thread on the loop thread; it does not run on the caller's own loop. diff --git a/kafka/net/backends/inet.py b/kafka/net/backend/inet.py similarity index 100% rename from kafka/net/backends/inet.py rename to kafka/net/backend/inet.py diff --git a/kafka/net/backends/selector.py b/kafka/net/backend/selector.py similarity index 98% rename from kafka/net/backends/selector.py rename to kafka/net/backend/selector.py index 9dc366f9c..e8a42cac5 100644 --- a/kafka/net/backends/selector.py +++ b/kafka/net/backend/selector.py @@ -11,7 +11,7 @@ import kafka.errors as Errors from kafka.future import Future -from kafka.net.backends.inet import create_connection as _inet_create_connection +from kafka.net.backend.inet import create_connection as _inet_create_connection from kafka.net.transport import KafkaSSLTransport, KafkaTCPTransport from kafka.version import __version__ @@ -43,7 +43,7 @@ def _initialize_coro(maybe_coro): class SelectorFuture(Future): - """The NetworkSelector's loop-awaitable future (see backends.NetBackendFuture). + """The NetworkSelector's loop-awaitable future (see backend.NetBackendFuture). ``kafka.future.Future`` is the thread-safe callback/handoff core with no ``__await__``; ``SelectorFuture`` adds it: ``yield self`` suspends the @@ -297,7 +297,7 @@ def on_io_thread(self): The clean form of the ``current_thread() is _io_thread`` identity check; callers use it to avoid blocking the loop on itself (e.g. a producer ``close()`` invoked from a produce callback). Part of the - NetBackend contract so alternate backends can answer it their own way. + NetBackend contract so alternate backend can answer it their own way. """ return self._io_thread is not None and threading.current_thread() is self._io_thread @@ -550,9 +550,9 @@ def reschedule(self, when, task): return task def create_future(self): - """Create a loop-awaitable future (see backends.NetBackendFuture). + """Create a loop-awaitable future (see backend.NetBackendFuture). - Portability seam for pluggable backends: core coroutines call this + Portability seam for pluggable backend: core coroutines call this instead of constructing ``Future`` directly, so an alternate backend (asyncio, Twisted) can return its own awaitable type. The selector's native awaitable is ``SelectorFuture`` (a ``Future`` with ``__await__``). diff --git a/kafka/net/backends/__init__.py b/kafka/net/backends/__init__.py deleted file mode 100644 index 4dc51aacf..000000000 --- a/kafka/net/backends/__init__.py +++ /dev/null @@ -1,7 +0,0 @@ -from .abstract import ( - NetBackend, NetTransport, NetProtocol, NetBackendFuture, - resolve_backend, register_backend_lazy, -) - -register_backend_lazy('selector', 'kafka.net.backends.selector', 'NetworkSelector') -register_backend_lazy('asyncio', 'kafka.net.backends.asyncio_backend', 'AsyncioBackend') diff --git a/kafka/net/compat.py b/kafka/net/compat.py index 6f4c7ba52..19c03e223 100644 --- a/kafka/net/compat.py +++ b/kafka/net/compat.py @@ -4,7 +4,7 @@ import time import kafka.errors as Errors -from kafka.net.backends import resolve_backend +from kafka.net.backend import resolve_backend from kafka.net.manager import KafkaConnectionManager from kafka.util import Timer diff --git a/kafka/net/http_connect.py b/kafka/net/http_connect.py index ff9b65a15..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.backends.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 6dcbf0298..2dd824465 100644 --- a/kafka/net/manager.py +++ b/kafka/net/manager.py @@ -7,7 +7,7 @@ from .connection import KafkaConnection from .metrics import KafkaManagerMetrics -from kafka.net.backends import resolve_backend +from kafka.net.backend import resolve_backend from kafka.cluster import ClusterMetadata import kafka.errors as Errors from kafka.net.transport import KafkaSSLTransport diff --git a/kafka/net/socks5.py b/kafka/net/socks5.py index 1794bbb00..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.backends.inet import KafkaNetSocket +from kafka.net.backend.inet import KafkaNetSocket log = logging.getLogger(__name__) diff --git a/kafka/producer/kafka.py b/kafka/producer/kafka.py index 69728c259..1a3bee1a9 100644 --- a/kafka/producer/kafka.py +++ b/kafka/producer/kafka.py @@ -380,7 +380,7 @@ class KafkaProducer: metrics. Default: 2 metrics_sample_window_ms (int): The maximum age in milliseconds of samples used to compute metrics. Default: 30000 - net (str or kafka.net.backends.NetBackend): The async backend that runs + net (str or kafka.net.backend.NetBackend): The async backend that runs this client's network I/O event loop. One of: a NetBackend instance; a registered name -- 'selector' (the built-in NetworkSelector) or 'asyncio' (runs I/O on an asyncio loop); or None diff --git a/test/conftest.py b/test/conftest.py index 590271e3d..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.backends.selector import NetworkSelector +from kafka.net.backend.selector import NetworkSelector from kafka.protocol.metadata import MetadataResponse diff --git a/test/net/backends/test_abstract.py b/test/net/backend/test_abstract.py similarity index 95% rename from test/net/backends/test_abstract.py rename to test/net/backend/test_abstract.py index d682f6914..b16e1ef67 100644 --- a/test/net/backends/test_abstract.py +++ b/test/net/backend/test_abstract.py @@ -1,4 +1,4 @@ -"""Conformance tests for the NetBackend contract (kafka/net/backends/abstract.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,10 +10,10 @@ import pytest -from kafka.net.backends.abstract import ( +from kafka.net.backend.abstract import ( NetBackend, NetTransport, resolve_backend, register_backend, _BACKENDS, ) -from kafka.net.backends.selector import NetworkSelector +from kafka.net.backend.selector import NetworkSelector from kafka.net.transport import KafkaTCPTransport @@ -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.backends.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.backends.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.backends.abstract 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/backends/test_asyncio_backend.py b/test/net/backend/test_asyncio_backend.py similarity index 97% rename from test/net/backends/test_asyncio_backend.py rename to test/net/backend/test_asyncio_backend.py index b0c1dc110..d59b916eb 100644 --- a/test/net/backends/test_asyncio_backend.py +++ b/test/net/backend/test_asyncio_backend.py @@ -1,9 +1,9 @@ -"""Tests for the asyncio NetBackend (kafka/net/backends/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.backends import NetBackend -from kafka.net.backends.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.backends.test_net_backend_future import NetBackendFutureContract +from test.net.backend.test_net_backend_future import NetBackendFutureContract @pytest.fixture diff --git a/test/net/backends/test_inet.py b/test/net/backend/test_inet.py similarity index 86% rename from test/net/backends/test_inet.py rename to test/net/backend/test_inet.py index d8913d9e3..12dd3f569 100644 --- a/test/net/backends/test_inet.py +++ b/test/net/backend/test_inet.py @@ -4,7 +4,7 @@ import pytest -from kafka.net.backends.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 @@ -18,7 +18,7 @@ def test_valid_host(self): assert len(res) == 5 def test_invalid_host(self): - with patch('kafka.net.backends.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 == [] @@ -81,14 +81,14 @@ def test_error_after_wait_write(self, net): class TestCreateConnection: def test_dns_failure(self, net): - with patch('kafka.net.backends.inet.KafkaNetSocket.dns_lookup', return_value=[]): + 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): fake_addr = [(socket.AF_INET, socket.SOCK_STREAM, 6, '', ('127.0.0.1', 9092))] - with patch('kafka.net.backends.inet.KafkaNetSocket.dns_lookup', return_value=fake_addr), \ - patch('kafka.net.backends.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)) @@ -96,8 +96,8 @@ 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.backends.inet.KafkaNetSocket.dns_lookup', return_value=fake_addr), \ - patch('kafka.net.backends.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 @@ -111,8 +111,8 @@ def test_tries_multiple_addresses(self, net): mock_sock2 = MagicMock() mock_sock2.connect_ex.return_value = 0 sockets = iter([mock_sock1, mock_sock2]) - with patch('kafka.net.backends.inet.KafkaNetSocket.dns_lookup', return_value=[addr1, addr2]), \ - patch('kafka.net.backends.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 @@ -125,8 +125,8 @@ def test_socket_options_applied(self, net): (socket.IPPROTO_TCP, socket.TCP_NODELAY, 1), (socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1), ] - with patch('kafka.net.backends.inet.KafkaNetSocket.dns_lookup', return_value=fake_addr), \ - patch('kafka.net.backends.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), @@ -141,7 +141,7 @@ def test_proxy_creates_socket(self, net): 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.backends.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')) @@ -154,7 +154,7 @@ def test_proxy_remote_dns_skips_local_lookup(self, net): 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.backends.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() @@ -163,8 +163,8 @@ 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.backends.inet.KafkaNetSocket.dns_lookup', return_value=[fake_addr]), \ - patch('kafka.net.backends.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)) @@ -177,7 +177,7 @@ def test_socks5h_does_dns_for_proxy_not_target(self, net): it is for the proxy hostname, not the target.""" proxy_addr = (socket.AF_INET, socket.SOCK_STREAM, 6, '', ('1.2.3.4', 1080)) mock_sock = MagicMock() - with patch('kafka.net.backends.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( @@ -186,12 +186,12 @@ def test_socks5h_does_dns_for_proxy_not_target(self, net): assert mock_dns.call_args.args[:2] == ('proxy', 1080) def test_socks5_proxy_dns_gaierror_raises(self): - with patch('kafka.net.backends.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.backends.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') @@ -201,7 +201,7 @@ def test_proxy_connect_dispatches_through_inherited_connect(self, net): 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.backends.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( @@ -305,8 +305,8 @@ def connect_ex(self, sock, sockaddr): try: fake_addr = (socket.AF_INET, socket.SOCK_STREAM, 6, '', ('10.0.0.1', 9092)) mock_sock = MagicMock() - with patch('kafka.net.backends.inet.KafkaNetSocket.dns_lookup', return_value=[fake_addr]), \ - patch('kafka.net.backends.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')) @@ -331,8 +331,8 @@ async def connect(self, net, addrinfo, socket_options=(), timeout_at=None): try: fake_addr = (socket.AF_INET, socket.SOCK_STREAM, 6, '', ('10.0.0.1', 9092)) - with patch('kafka.net.backends.inet.KafkaNetSocket.dns_lookup', return_value=[fake_addr]), \ - patch('kafka.net.backends.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/backends/test_net_backend_future.py b/test/net/backend/test_net_backend_future.py similarity index 98% rename from test/net/backends/test_net_backend_future.py rename to test/net/backend/test_net_backend_future.py index 0cf48632a..8cdc1529c 100644 --- a/test/net/backends/test_net_backend_future.py +++ b/test/net/backend/test_net_backend_future.py @@ -11,8 +11,8 @@ import pytest from kafka.future import Future -from kafka.net.backends import NetBackendFuture -from kafka.net.backends.selector import NetworkSelector +from kafka.net.backend import NetBackendFuture +from kafka.net.backend.selector import NetworkSelector class NetBackendFutureContract: diff --git a/test/net/backends/test_selector.py b/test/net/backend/test_selector.py similarity index 99% rename from test/net/backends/test_selector.py rename to test/net/backend/test_selector.py index 0b249eced..0c1542a09 100644 --- a/test/net/backends/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.backends.selector import ( +from kafka.net.backend.selector import ( KernelEvent, NetworkSelector, Task, diff --git a/test/net/test_http_connect.py b/test/net/test_http_connect.py index cc415b2e4..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.backends.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 650e4dede..ca3816cb4 100644 --- a/test/net/test_manager.py +++ b/test/net/test_manager.py @@ -7,7 +7,7 @@ from kafka.cluster import ClusterMetadata from kafka.future import Future -from kafka.net.backends.selector import NetworkSelector +from kafka.net.backend.selector import NetworkSelector from kafka.net.connection import KafkaConnection from kafka.net.manager import KafkaConnectionManager import kafka.errors as Errors @@ -536,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.backends.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_transport.py b/test/net/test_transport.py index 855d79b93..3d042d408 100644 --- a/test/net/test_transport.py +++ b/test/net/test_transport.py @@ -7,7 +7,7 @@ import kafka.errors as Errors from kafka.future import Future -from kafka.net.backends.selector import NetworkSelector, TaskState +from kafka.net.backend.selector import NetworkSelector, TaskState from kafka.net.transport import KafkaSSLTransport, KafkaTCPTransport From 9ecc18baec98072b8262386b62021388e2b600cb Mon Sep 17 00:00:00 2001 From: Dana Powers Date: Wed, 15 Jul 2026 10:32:47 -0700 Subject: [PATCH 08/10] transport -> net.backend/ --- kafka/net/__init__.py | 3 +-- kafka/net/backend/abstract.py | 2 +- kafka/net/backend/selector.py | 2 +- kafka/net/{ => backend}/transport.py | 0 kafka/net/manager.py | 2 +- test/net/backend/test_abstract.py | 2 +- test/net/{ => backend}/test_transport.py | 2 +- test/net/test_connection.py | 2 +- 8 files changed, 7 insertions(+), 8 deletions(-) rename kafka/net/{ => backend}/transport.py (100%) rename test/net/{ => backend}/test_transport.py (99%) diff --git a/kafka/net/__init__.py b/kafka/net/__init__.py index 5bedc53f9..51bf318ca 100644 --- a/kafka/net/__init__.py +++ b/kafka/net/__init__.py @@ -3,7 +3,6 @@ from .metrics import KafkaConnectionMetrics, KafkaManagerMetrics from .http_connect import HttpConnectProxy from .socks5 import Socks5Proxy -from .transport import KafkaTCPTransport, KafkaSSLTransport from .wakeup_notifier import WakeupNotifier from .compat import KafkaNetClient @@ -12,6 +11,6 @@ __all__ = [ 'KafkaConnection', 'KafkaConnectionManager', 'KafkaConnectionMetrics', 'KafkaManagerMetrics', - 'HttpConnectProxy', 'Socks5Proxy', 'KafkaTCPTransport', 'KafkaSSLTransport', + 'HttpConnectProxy', 'Socks5Proxy', 'WakeupNotifier', 'KafkaNetClient', ] diff --git a/kafka/net/backend/abstract.py b/kafka/net/backend/abstract.py index 197604e04..3664d15e6 100644 --- a/kafka/net/backend/abstract.py +++ b/kafka/net/backend/abstract.py @@ -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 diff --git a/kafka/net/backend/selector.py b/kafka/net/backend/selector.py index e8a42cac5..f22ef7e69 100644 --- a/kafka/net/backend/selector.py +++ b/kafka/net/backend/selector.py @@ -12,7 +12,7 @@ import kafka.errors as Errors from kafka.future import Future from kafka.net.backend.inet import create_connection as _inet_create_connection -from kafka.net.transport import KafkaSSLTransport, KafkaTCPTransport +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/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/test/net/backend/test_abstract.py b/test/net/backend/test_abstract.py index b16e1ef67..9c15cb0f6 100644 --- a/test/net/backend/test_abstract.py +++ b/test/net/backend/test_abstract.py @@ -14,7 +14,7 @@ NetBackend, NetTransport, resolve_backend, register_backend, _BACKENDS, ) from kafka.net.backend.selector import NetworkSelector -from kafka.net.transport import KafkaTCPTransport +from kafka.net.backend.transport import KafkaTCPTransport # The full contract surface, kept here so a missing/renamed method fails loudly. 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 3d042d408..b746a6f8e 100644 --- a/test/net/test_transport.py +++ b/test/net/backend/test_transport.py @@ -8,7 +8,7 @@ import kafka.errors as Errors from kafka.future import Future from kafka.net.backend.selector import NetworkSelector, TaskState -from kafka.net.transport import KafkaSSLTransport, KafkaTCPTransport +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 84f45c831..fe185686d 100644 --- a/test/net/test_connection.py +++ b/test/net/test_connection.py @@ -7,7 +7,7 @@ from kafka.future import Future 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 From 980f812d755a66cc4ef746950ca30ef135233251 Mon Sep 17 00:00:00 2001 From: Dana Powers Date: Wed, 15 Jul 2026 10:33:43 -0700 Subject: [PATCH 09/10] update transport docs module --- docs/apidoc/modules.rst | 4 ++-- docs/apidoc/net/{ => backend}/transport.rst | 0 2 files changed, 2 insertions(+), 2 deletions(-) rename docs/apidoc/net/{ => backend}/transport.rst (100%) diff --git a/docs/apidoc/modules.rst b/docs/apidoc/modules.rst index d5640d821..b9c410439 100644 --- a/docs/apidoc/modules.rst +++ b/docs/apidoc/modules.rst @@ -93,7 +93,7 @@ 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.http_connect` - Tunnels broker connections through an HTTP CONNECT proxy (RFC 7231). @@ -106,7 +106,7 @@ driving the protocol layer directly from the REPL. manager connection - transport + 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 From c85ac92922ced9155e3b8d4d4dda6fe31430eec4 Mon Sep 17 00:00:00 2001 From: Dana Powers Date: Wed, 15 Jul 2026 13:03:51 -0700 Subject: [PATCH 10/10] fixups --- kafka/net/backend/abstract.py | 8 ++++---- kafka/net/backend/selector.py | 4 ++-- test/net/backend/test_selector.py | 10 +++++----- 3 files changed, 11 insertions(+), 11 deletions(-) diff --git a/kafka/net/backend/abstract.py b/kafka/net/backend/abstract.py index 3664d15e6..04521bb94 100644 --- a/kafka/net/backend/abstract.py +++ b/kafka/net/backend/abstract.py @@ -57,13 +57,13 @@ class NetBackendFuture(Protocol): 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 - backend. The selector's ``SelectorFuture`` is the reference implementation: + backends. The selector's ``SelectorFuture`` is the reference implementation: it subclasses the thread-safe ``kafka.future.Future`` (the portable callback core) and adds ``__await__``. A plain ``Future`` is deliberately NOT a NetBackendFuture -- it has no ``__await__`` -- so awaiting a cross-thread handoff future fails loudly instead of silently working on one backend. - Pinned semantics -- the three axes where backend could otherwise diverge: + Pinned semantics -- the three axes where backends could otherwise diverge: 1. **Resolution thread.** A future from ``create_future()`` is created and resolved (``success`` / ``failure``) on the loop/IO thread only. @@ -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 = {} @@ -302,7 +302,7 @@ def resolve_backend(net, config): 'or None; got %r' % (net,)) return net # net is None: auto-detect, else default. Auto-detected-but-unregistered - # backend fall back silently (an explicit name would have raised above). + # backends fall back silently (an explicit name would have raised above). name = _detect_async_library() if name is None or name not in _BACKENDS: name = 'selector' diff --git a/kafka/net/backend/selector.py b/kafka/net/backend/selector.py index f22ef7e69..ecb7b3f58 100644 --- a/kafka/net/backend/selector.py +++ b/kafka/net/backend/selector.py @@ -297,7 +297,7 @@ def on_io_thread(self): The clean form of the ``current_thread() is _io_thread`` identity check; callers use it to avoid blocking the loop on itself (e.g. a producer ``close()`` invoked from a produce callback). Part of the - NetBackend contract so alternate backend can answer it their own way. + NetBackend contract so alternate backends can answer it their own way. """ return self._io_thread is not None and threading.current_thread() is self._io_thread @@ -552,7 +552,7 @@ def reschedule(self, when, task): def create_future(self): """Create a loop-awaitable future (see backend.NetBackendFuture). - Portability seam for pluggable backend: core coroutines call this + Portability seam for pluggable backends: core coroutines call this instead of constructing ``Future`` directly, so an alternate backend (asyncio, Twisted) can return its own awaitable type. The selector's native awaitable is ``SelectorFuture`` (a ``Future`` with ``__await__``). diff --git a/test/net/backend/test_selector.py b/test/net/backend/test_selector.py index 0c1542a09..bc1c61a4b 100644 --- a/test/net/backend/test_selector.py +++ b/test/net/backend/test_selector.py @@ -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.