Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 7 additions & 9 deletions .claude/skills/uts-to-python/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
4 changes: 2 additions & 2 deletions .github/workflows/check.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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 }}
Expand Down
6 changes: 3 additions & 3 deletions .github/workflows/lint.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand All @@ -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:
Expand Down
10 changes: 5 additions & 5 deletions .github/workflows/release.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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/
Expand Down Expand Up @@ -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/
Expand Down Expand Up @@ -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/
Expand Down
2 changes: 1 addition & 1 deletion ably/realtime/channel.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
7 changes: 6 additions & 1 deletion ably/realtime/connectionmanager.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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()
Expand Down
6 changes: 3 additions & 3 deletions ably/util/eventemitter.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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)
Expand Down
16 changes: 14 additions & 2 deletions ably/util/helper.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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():
Expand Down Expand Up @@ -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()
Expand Down
11 changes: 9 additions & 2 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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)
Expand Down
8 changes: 6 additions & 2 deletions test/ably/realtime/realtimechannel_publish_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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
Expand Down
30 changes: 19 additions & 11 deletions test/ably/rest/encoders_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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):
Expand All @@ -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ó'
Expand All @@ -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'
Expand All @@ -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']
Expand All @@ -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']
Expand All @@ -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'])
Expand All @@ -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'])
Expand Down Expand Up @@ -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'
Expand All @@ -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'
Expand All @@ -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

Expand All @@ -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'
Expand All @@ -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'
Expand Down
Loading
Loading