diff --git a/.claude/skills/uts-to-python/SKILL.md b/.claude/skills/uts-to-python/SKILL.md index a413d6bb..fcfa8ae3 100644 --- a/.claude/skills/uts-to-python/SKILL.md +++ b/.claude/skills/uts-to-python/SKILL.md @@ -533,11 +533,11 @@ The template has one. **10. A refused connection and a connect timeout are indistinguishable, and skip the fallback loop.** `ws_connect` catches only `WebSocketException` and `socket.gaierror`, so `ConnectionRefusedError` and `asyncio.TimeoutError` never reach `_emit('failed')`: -the attempt hangs until the transition timer fires with a generic 50003/504, and each -one leaks a `connect_base` task. `respond_with_dns_error()` is the only fast, caught -failure — 40000/400 with the real cause — so **prefer it** whenever you just need "the -connect failed", and note the substitution at the site. Keep -`realtime_request_timeout` short in any test that does wait a refusal out. +the attempt hangs until the transition timer fires with a generic 50003/504. +`respond_with_dns_error()` is the only fast, caught failure — 40000/400 with the real +cause — so **prefer it** whenever you just need "the connect failed", and note the +substitution at the site. Keep `realtime_request_timeout` short in any test that does +wait a refusal out. **11. Keep the fallback hosts empty** unless the spec is about them. `check_connection()` is a module-level, **synchronous** `httpx.get` that no seam @@ -585,11 +585,9 @@ side. `client.connection.connection_details.connection_key`. `client.connection.connection_details` and `connection.error_reason` **are** public. -**19. Two noisy-but-harmless teardown messages.** `Task exception was never retrieved` +**19. A noisy-but-harmless teardown message.** `Task exception was never retrieved` for a client whose connect failed — `WebSocketTransport.send` raises a bare -`Exception()` when `self.websocket is None`. And `Task was destroyed but it is -pending!`, one per refused attempt, which is trap 10's leak showing. Neither is a -failure; do not chase them. +`Exception()` when `self.websocket is None`. It is not a failure; do not chase it. **20. A channel needs a SUSPENDED *connection* to reach SUSPENDED.** `_propagate_connection_interruption` fires only for CLOSING/CLOSED/FAILED/SUSPENDED, so diff --git a/.github/workflows/check.yml b/.github/workflows/check.yml index f69bf097..f8ab8cd9 100644 --- a/.github/workflows/check.yml +++ b/.github/workflows/check.yml @@ -23,12 +23,12 @@ jobs: matrix: python-version: ['3.8', '3.9', '3.10', '3.11', '3.12', '3.13', '3.14'] steps: - - uses: actions/checkout@34e114876b0b11c390a56381ad16ebd13914f8d5 # v4 + - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7 with: submodules: 'recursive' persist-credentials: false - name: Set up Python ${{ matrix.python-version }} - uses: actions/setup-python@a26af69be951a213d495a4c3e4e4022e16d87065 # v5 + uses: actions/setup-python@5fda3b95a4ea91299a34e894583c3862153e4b97 # v7 id: setup-python with: python-version: ${{ matrix.python-version }} diff --git a/.github/workflows/lint.yml b/.github/workflows/lint.yml index d1027713..b3a5a113 100644 --- a/.github/workflows/lint.yml +++ b/.github/workflows/lint.yml @@ -14,12 +14,12 @@ jobs: contents: read runs-on: ubuntu-latest steps: - - uses: actions/checkout@34e114876b0b11c390a56381ad16ebd13914f8d5 # v4 + - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7 with: submodules: 'recursive' persist-credentials: false - name: Set up Python 3.9 - uses: actions/setup-python@a26af69be951a213d495a4c3e4e4022e16d87065 # v5 + uses: actions/setup-python@5fda3b95a4ea91299a34e894583c3862153e4b97 # v7 id: setup-python with: python-version: '3.9' @@ -29,7 +29,7 @@ jobs: with: enable-cache: true - - uses: actions/cache@0057852bfaa89a56745cba8c7296529d2fc39830 # v4 + - uses: actions/cache@55cc8345863c7cc4c66a329aec7e433d2d1c52a9 # v6 name: Define a cache for the virtual environment based on the dependencies lock file id: cache with: diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index 2a7a2ae2..ca2e5857 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -16,12 +16,12 @@ jobs: contents: read steps: - - uses: actions/checkout@34e114876b0b11c390a56381ad16ebd13914f8d5 # v4 + - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7 with: submodules: 'recursive' persist-credentials: false - name: Set up Python 3.12 - uses: actions/setup-python@a26af69be951a213d495a4c3e4e4022e16d87065 # v5 + uses: actions/setup-python@5fda3b95a4ea91299a34e894583c3862153e4b97 # v7 id: setup-python with: python-version: 3.12 @@ -38,7 +38,7 @@ jobs: - name: Build a binary wheel and a source tarball run: uv build - name: Store the distribution packages - uses: actions/upload-artifact@ea165f8d65b6e75b540449e92b4886f43607fa02 # v4 + uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a # v7 with: name: python-package-distributions path: dist/ @@ -80,7 +80,7 @@ jobs: steps: - name: Download all the dists - uses: actions/download-artifact@d3f86a106a0bac45b974a628896c90dbdf5c8093 # v4 + uses: actions/download-artifact@3e5f45b2cfb9172054b4087a40e8e0b5a5461e7c # v8 with: name: python-package-distributions path: dist/ @@ -125,7 +125,7 @@ jobs: steps: - name: Download all the dists - uses: actions/download-artifact@d3f86a106a0bac45b974a628896c90dbdf5c8093 # v4 + uses: actions/download-artifact@3e5f45b2cfb9172054b4087a40e8e0b5a5461e7c # v8 with: name: python-package-distributions path: dist/ diff --git a/ably/realtime/channel.py b/ably/realtime/channel.py index fde9bb06..f67a5173 100644 --- a/ably/realtime/channel.py +++ b/ably/realtime/channel.py @@ -728,7 +728,7 @@ def _on_message(self, proto_msg: dict) -> None: elif self.state == ChannelState.ATTACHING: self._notify_state(ChannelState.ATTACHED, resumed=resumed, has_presence=has_presence) else: - log.warn("RealtimeChannel._on_message(): ATTACHED received while not attaching") + log.warning("RealtimeChannel._on_message(): ATTACHED received while not attaching") elif action == ProtocolMessageAction.DETACHED: if self.state == ChannelState.DETACHING: self._notify_state(ChannelState.DETACHED) diff --git a/ably/realtime/connectionmanager.py b/ably/realtime/connectionmanager.py index 9f920b07..dd3f1105 100644 --- a/ably/realtime/connectionmanager.py +++ b/ably/realtime/connectionmanager.py @@ -452,7 +452,7 @@ async def on_disconnected(self, exception: AblyException) -> None: else: self.notify_state(ConnectionState.DISCONNECTED, exception) else: - log.warn("DISCONNECTED message received without error") + log.warning("DISCONNECTED message received without error") async def on_token_error(self, exception: AblyException) -> None: if self.__error_reason is None or not is_token_error(self.__error_reason): @@ -775,6 +775,11 @@ def cancel_retry_timer(self) -> None: def disconnect_transport(self) -> None: log.info('ConnectionManager.disconnect_transport()') + # A connect attempt still in flight is abandoned along with the transport it was + # opening, which reports neither 'connected' nor 'failed' once disposed. connect_base + # reaches here itself when it gives up, and is left to return. + if self.connect_base_task and self.connect_base_task is not asyncio.current_task(): + self.connect_base_task.cancel() if self.transport: # RTN19a: Requeue pending messages before disposing transport self.requeue_pending_messages() diff --git a/ably/util/eventemitter.py b/ably/util/eventemitter.py index 74f0beb6..b3a62428 100644 --- a/ably/util/eventemitter.py +++ b/ably/util/eventemitter.py @@ -3,7 +3,7 @@ from pyee.asyncio import AsyncIOEventEmitter -from ably.util.helper import is_callable_or_coroutine +from ably.util.helper import is_callable_or_coroutine, is_coroutine_function # pyee's event emitter doesn't support attaching a listener to all events # so to patch it, we create a wrapper which uses two event emitters, one @@ -69,7 +69,7 @@ def on(self, *args): else: raise ValueError("EventEmitter.on(): invalid args") - if asyncio.iscoroutinefunction(listener): + if is_coroutine_function(listener): async def wrapped_listener(*args, **kwargs): try: await listener(*args, **kwargs) @@ -114,7 +114,7 @@ def once(self, *args): else: raise ValueError("EventEmitter.on(): invalid args") - if asyncio.iscoroutinefunction(listener): + if is_coroutine_function(listener): async def wrapped_listener(*args, **kwargs): try: await listener(*args, **kwargs) diff --git a/ably/util/helper.py b/ably/util/helper.py index 4d37d4f4..5c0a9992 100644 --- a/ably/util/helper.py +++ b/ably/util/helper.py @@ -11,6 +11,11 @@ from ably.util.exceptions import AblyException +# asyncio.iscoroutinefunction() also recognises callables carrying this marker, +# such as AsyncMock from the `mock` package (and from unittest.mock on older +# Pythons) and callables marked by asgiref before Python 3.12 +_is_coroutine_marker = getattr(asyncio.coroutines, '_is_coroutine', None) + def get_random_id(): # get random string of letters and digits @@ -19,8 +24,15 @@ def get_random_id(): return random_id +def is_coroutine_function(value): + """Whether `value` is a coroutine function, recognising what asyncio.iscoroutinefunction() does.""" + if inspect.iscoroutinefunction(value): + return True + return _is_coroutine_marker is not None and getattr(value, '_is_coroutine', None) is _is_coroutine_marker + + def is_callable_or_coroutine(value): - return asyncio.iscoroutinefunction(value) or inspect.isfunction(value) or inspect.ismethod(value) + return is_coroutine_function(value) or inspect.isfunction(value) or inspect.ismethod(value) def unix_time_ms(): @@ -67,7 +79,7 @@ def __init__(self, timeout: float, callback: Callable): async def _job(self): await asyncio.sleep(self._timeout / 1000) - if asyncio.iscoroutinefunction(self._callback): + if is_coroutine_function(self._callback): await self._callback() else: self._callback() diff --git a/pyproject.toml b/pyproject.toml index 342f92c2..2cb6e597 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -93,6 +93,13 @@ packages = ["ably"] [tool.pytest.ini_options] timeout = 30 asyncio_mode = "auto" +filterwarnings = [ + # pytest-asyncio before 1.0 calls asyncio APIs that Python 3.14 deprecates; + # the releases that avoid them need Python 3.9 and pytest 8.2 or later + "ignore:'asyncio.iscoroutinefunction' is deprecated:DeprecationWarning:pytest_asyncio", + "ignore:'asyncio.get_event_loop_policy' is deprecated:DeprecationWarning:pytest_asyncio", + "ignore:'asyncio.set_event_loop_policy' is deprecated:DeprecationWarning:pytest_asyncio", +] [[tool.uv.index]] name = "experimental" @@ -107,8 +114,8 @@ extend-exclude = [ ] [tool.ruff.lint] -# Enable Pyflakes (F), pycodestyle (E, W), pep8-naming (N), isort (I), pyupgrade (UP), bugbear (B) and comprehensions (C4) -select = ["E", "W", "F", "N", "I", "UP", "B", "C4"] +# Enable Pyflakes (F), pycodestyle (E, W), pep8-naming (N), isort (I), pyupgrade (UP), bugbear (B), comprehensions (C4) and logging-warn (G010) +select = ["E", "W", "F", "N", "I", "UP", "B", "C4", "G010"] ignore = [ "N818", # exception name should end in 'Error' "UP026", # mock -> unittest.mock (need mock package for Python 3.7 AsyncMock support) diff --git a/test/ably/realtime/realtimechannel_publish_test.py b/test/ably/realtime/realtimechannel_publish_test.py index 9ecf10f9..baa67a41 100644 --- a/test/ably/realtime/realtimechannel_publish_test.py +++ b/test/ably/realtime/realtimechannel_publish_test.py @@ -225,7 +225,9 @@ async def check_pending(): return connection_manager.pending_message_queue.count() > 0 await assert_waiter(check_pending, timeout=2) - # Force DISCONNECTED state + # Simulate loss of connection: dispose the transport, then force DISCONNECTED state + assert connection_manager.transport + await connection_manager.transport.dispose() connection_manager.notify_state( ConnectionState.DISCONNECTED, AblyException('Test disconnect', 400, 80003) @@ -266,7 +268,9 @@ async def check_pending(): return connection_manager.pending_message_queue.count() > 0 await assert_waiter(check_pending, timeout=2) - # Force DISCONNECTED state + # Simulate loss of connection: dispose the transport, then force DISCONNECTED state + assert connection_manager.transport + await connection_manager.transport.dispose() connection_manager.notify_state(ConnectionState.DISCONNECTED, None) # Give time for state transition diff --git a/test/ably/rest/encoders_test.py b/test/ably/rest/encoders_test.py index f8023c5d..f6fe1cae 100644 --- a/test/ably/rest/encoders_test.py +++ b/test/ably/rest/encoders_test.py @@ -4,10 +4,12 @@ import sys from unittest import mock +import httpx import msgpack import pytest from ably import CipherParams +from ably.http.http import Response from ably.types.message import Message from ably.util.crypto import get_cipher from test.ably.testapp import TestApp @@ -21,6 +23,12 @@ log = logging.getLogger(__name__) +def patch_http_post(): + # The patched post records the publish request and resolves to an empty 201 response + return mock.patch('ably.rest.rest.Http.post', new_callable=AsyncMock, + return_value=Response(httpx.Response(201))) + + class TestTextEncodersNoEncryption(BaseAsyncTestCase): @pytest.fixture(autouse=True) async def setup(self): @@ -31,7 +39,7 @@ async def setup(self): async def test_text_utf8(self): channel = self.ably.channels["persisted:publish"] - with mock.patch('ably.rest.rest.Http.post', new_callable=AsyncMock) as post_mock: + with patch_http_post() as post_mock: await channel.publish('event', 'foó') _, kwargs = post_mock.call_args assert json.loads(kwargs['body'])['data'] == 'foó' @@ -41,7 +49,7 @@ async def test_str(self): # This test only makes sense for py2 channel = self.ably.channels["persisted:publish"] - with mock.patch('ably.rest.rest.Http.post', new_callable=AsyncMock) as post_mock: + with patch_http_post() as post_mock: await channel.publish('event', 'foo') _, kwargs = post_mock.call_args assert json.loads(kwargs['body'])['data'] == 'foo' @@ -50,7 +58,7 @@ async def test_str(self): async def test_with_binary_type(self): channel = self.ably.channels["persisted:publish"] - with mock.patch('ably.rest.rest.Http.post', new_callable=AsyncMock) as post_mock: + with patch_http_post() as post_mock: await channel.publish('event', bytearray(b'foo')) _, kwargs = post_mock.call_args raw_data = json.loads(kwargs['body'])['data'] @@ -60,7 +68,7 @@ async def test_with_binary_type(self): async def test_with_bytes_type(self): channel = self.ably.channels["persisted:publish"] - with mock.patch('ably.rest.rest.Http.post', new_callable=AsyncMock) as post_mock: + with patch_http_post() as post_mock: await channel.publish('event', b'foo') _, kwargs = post_mock.call_args raw_data = json.loads(kwargs['body'])['data'] @@ -70,7 +78,7 @@ async def test_with_bytes_type(self): async def test_with_json_dict_data(self): channel = self.ably.channels["persisted:publish"] data = {'foó': 'bár'} - with mock.patch('ably.rest.rest.Http.post', new_callable=AsyncMock) as post_mock: + with patch_http_post() as post_mock: await channel.publish('event', data) _, kwargs = post_mock.call_args raw_data = json.loads(json.loads(kwargs['body'])['data']) @@ -80,7 +88,7 @@ async def test_with_json_dict_data(self): async def test_with_json_list_data(self): channel = self.ably.channels["persisted:publish"] data = ['foó', 'bár'] - with mock.patch('ably.rest.rest.Http.post', new_callable=AsyncMock) as post_mock: + with patch_http_post() as post_mock: await channel.publish('event', data) _, kwargs = post_mock.call_args raw_data = json.loads(json.loads(kwargs['body'])['data']) @@ -161,7 +169,7 @@ def decrypt(self, payload, options=None): async def test_text_utf8(self): channel = self.ably.channels.get("persisted:publish_enc", cipher=self.cipher_params) - with mock.patch('ably.rest.rest.Http.post', new_callable=AsyncMock) as post_mock: + with patch_http_post() as post_mock: await channel.publish('event', 'fóo') _, kwargs = post_mock.call_args assert json.loads(kwargs['body'])['encoding'].strip('/') == 'utf-8/cipher+aes-128-cbc/base64' @@ -172,7 +180,7 @@ async def test_str(self): # This test only makes sense for py2 channel = self.ably.channels["persisted:publish"] - with mock.patch('ably.rest.rest.Http.post', new_callable=AsyncMock) as post_mock: + with patch_http_post() as post_mock: await channel.publish('event', 'foo') _, kwargs = post_mock.call_args assert json.loads(kwargs['body'])['data'] == 'foo' @@ -182,7 +190,7 @@ async def test_with_binary_type(self): channel = self.ably.channels.get("persisted:publish_enc", cipher=self.cipher_params) - with mock.patch('ably.rest.rest.Http.post', new_callable=AsyncMock) as post_mock: + with patch_http_post() as post_mock: await channel.publish('event', bytearray(b'foo')) _, kwargs = post_mock.call_args @@ -195,7 +203,7 @@ async def test_with_json_dict_data(self): channel = self.ably.channels.get("persisted:publish_enc", cipher=self.cipher_params) data = {'foó': 'bár'} - with mock.patch('ably.rest.rest.Http.post', new_callable=AsyncMock) as post_mock: + with patch_http_post() as post_mock: await channel.publish('event', data) _, kwargs = post_mock.call_args assert json.loads(kwargs['body'])['encoding'].strip('/') == 'json/utf-8/cipher+aes-128-cbc/base64' @@ -206,7 +214,7 @@ async def test_with_json_list_data(self): channel = self.ably.channels.get("persisted:publish_enc", cipher=self.cipher_params) data = ['foó', 'bár'] - with mock.patch('ably.rest.rest.Http.post', new_callable=AsyncMock) as post_mock: + with patch_http_post() as post_mock: await channel.publish('event', data) _, kwargs = post_mock.call_args assert json.loads(kwargs['body'])['encoding'].strip('/') == 'json/utf-8/cipher+aes-128-cbc/base64' diff --git a/test/unit/connectionmanager_test.py b/test/unit/connectionmanager_test.py new file mode 100644 index 00000000..f85a3f48 --- /dev/null +++ b/test/unit/connectionmanager_test.py @@ -0,0 +1,64 @@ +import asyncio + +import pytest + +from ably.realtime.connection import ConnectionState +from test.uts.helpers.client import close_open_clients, realtime_client +from test.uts.helpers.clock import FakeClock, settle +from test.uts.helpers.mock_websocket import CONNECTED_MESSAGE, MockWebSocket + + +@pytest.fixture(autouse=True) +async def close_clients(): + yield + await close_open_clients() + + +# RTN14c +async def test_an_attempt_the_transition_timer_ends_is_cancelled(): + # The server never answers the handshake, so only the transition timer ends the attempt + mock_ws = MockWebSocket(on_connection_attempt=lambda conn: None) + clock = FakeClock() + client = realtime_client(mock_ws, clock=clock, realtime_request_timeout=1000) + + client.connect() + await settle() + attempt = client.connection.connection_manager.connect_base_task + assert not attempt.done() + + await clock.advance(1100) + await settle() + + assert client.connection.state == ConnectionState.DISCONNECTED + assert attempt.done() + + +# RTN14c +async def test_an_attempt_the_transition_timer_ends_does_not_connect_afterwards(): + authorised = asyncio.Event() + attempts = [] + + async def auth_callback(token_params): + await authorised.wait() + return 'a-token' + + def on_connection_attempt(conn): + attempts.append(conn) + conn.respond_with_success(CONNECTED_MESSAGE) + + mock_ws = MockWebSocket(on_connection_attempt=on_connection_attempt) + clock = FakeClock() + client = realtime_client(mock_ws, clock=clock, auth_callback=auth_callback, + realtime_request_timeout=1000) + + client.connect() + await settle() + await clock.advance(1100) + assert client.connection.state == ConnectionState.DISCONNECTED + + # The auth callback answers after the attempt it was serving has been given up on + authorised.set() + await settle() + + assert client.connection.state == ConnectionState.DISCONNECTED + assert attempts == [] diff --git a/test/uts/deviations.md b/test/uts/deviations.md index 655f9699..487759bd 100644 --- a/test/uts/deviations.md +++ b/test/uts/deviations.md @@ -1161,7 +1161,7 @@ the reason is plumbed through: a missing reason should give an `AblyException`, An ATTACHED arriving while the channel is DETACHING or DETACHED must be answered with a new DETACH, the channel remaining in or returning to DETACHING. `_on_message` handles ATTACHED only for the ATTACHED (RTL12) and ATTACHING cases; every other state falls through to -`log.warn("ATTACHED received while not attaching")` and nothing is sent. While DETACHING +`log.warning("ATTACHED received while not attaching")` and nothing is sent. While DETACHING that leaves the detach to time out, so `detach()` raises "Channel detach timed out" and the channel returns to ATTACHED. @@ -1904,7 +1904,7 @@ transition timer. Measured, with `fallback_hosts=[]` and `realtime_request_timeo | `asyncio.TimeoutError` (`respond_with_timeout`) | still CONNECTING | at t=1000 | 50003 / 504 | | `socket.gaierror` (`respond_with_dns_error`) | already DISCONNECTED | at t=0 | 40000 / 400, naming the cause | -Three consequences: +Two consequences: - A refused connection and a connect timeout are indistinguishable from each other *and* from a server that accepts the socket and says nothing. All three surface as the @@ -1914,11 +1914,6 @@ Three consequences: tried for refused and for timeout, against six attempts — primary plus all five fallbacks — for a DNS error. RTN17d's fallback behaviour therefore cannot happen in practice. -- **Every refused attempt leaks a task and a future.** `try_a_host`'s future - (`connectionmanager.py:646`) is never settled, so each attempt leaves a - `connect_base()` task awaiting it for good, printing `Task was destroyed but it is - pending!` at interpreter shutdown. A long-lived client reconnecting against a refusing - host leaks one per attempt. **Tests affected:** `test_rtn14d_retry_recoverable_failure` is the adapted test that pins it — it asserts that the refusal moves nothing, that DISCONNECTED arrives only when the @@ -1930,14 +1925,13 @@ the fallback loop, each noted at the site: `test_rtn17f_fallback_on_error`, `test_rtn17h_fallback_domains_from_rec2`, `test_rtn17i_prefer_primary_domain`, `test_rtn17j_connectivity_check_before_fallback`, `test_rtn17e_http_uses_same_fallback`, `test_rtn13b_ping_error_suspended`, `test_rtn16g3_recovery_key_null_inactive`, -`test_rtc7_disconnected_retry_timeout`. `test_rtn17g_empty_fallback_set_error` and -`test_rtl6c4_fails_conn_suspended` keep `respond_with_refused()` deliberately — the first -because it asserts that *no* fallback follows, the second because swapping it would silence -the ten `Task was destroyed` lines that are the leak showing. +`test_rtc7_disconnected_retry_timeout`. `test_rtn17g_empty_fallback_set_error` keeps +`respond_with_refused()` deliberately, because it asserts that *no* fallback follows, and +`test_rtl6c4_fails_conn_suspended` keeps it as the specification has it. **Status:** open bug. Widening the `except` to `(WebSocketException, OSError, asyncio.TimeoutError)` — or, better, emitting `failed` from a guard no exception type can -escape — fixes all three consequences. +escape — fixes both consequences. ### The connectivity check bypasses every seam and blocks the event loop @@ -2561,12 +2555,10 @@ vanish.** RTN14d, RTN17d, RTN17e. `websockettransport.py:117` catches only `(WebSocketException, socket.gaierror)`, so a `ConnectionRefusedError` (an `OSError`) and an `asyncio.TimeoutError` never reach `_emit('failed')`, the future `try_host` awaits is never settled, and the attempt is ended only by the transition timer with a generic 50003/504. -Three consequences: refused, timed-out and silently-accepted connections are -indistinguishable; **the fallback loop is unreachable** for refused and timeout (measured: -one attempt and no fallback tried, against six for a DNS error); and each attempt leaks a -`connect_base()` task and its future (`connectionmanager.py:646`), printing `Task was -destroyed but it is pending!` at shutdown. Widening the `except`, or emitting `failed` from a -guard no exception can escape, fixes all three. +Two consequences: refused, timed-out and silently-accepted connections are +indistinguishable; and **the fallback loop is unreachable** for refused and timeout (measured: +one attempt and no fallback tried, against six for a DNS error). Widening the `except`, or +emitting `failed` from a guard no exception can escape, fixes both. `test/uts/realtime/unit/connection/connection_failures_test.py -k rtn14d` (this one is adapted, so it **passes** today and fails when the defect is fixed — read it as the pin, not the proof) diff --git a/test/uts/helpers/clock.py b/test/uts/helpers/clock.py index 6d818715..c58a5b8b 100644 --- a/test/uts/helpers/clock.py +++ b/test/uts/helpers/clock.py @@ -144,12 +144,9 @@ def __next_due(self, target): @staticmethod async def __invoke(callback): - if asyncio.iscoroutinefunction(callback): - await callback() - else: - result = callback() - if inspect.isawaitable(result): - await result + result = callback() + if inspect.isawaitable(result): + await result async def advance_to_connection_state(client, clock, state, step, limit=60):