Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
18 commits
Select commit Hold shift + click to select a range
d1ada70
feat: allow tests to supply the callable which opens websocket connec…
owenpearson Sep 23, 2026
c61872f
feat: read the realtime client's time through the clock
owenpearson Sep 23, 2026
1b15510
test: add the UTS mock websocket and derive the auto-connect spec
owenpearson Sep 23, 2026
9164354
test: derive the realtime client unit specs
owenpearson Sep 23, 2026
a9f3e1c
test: derive the connection state, id and error reason unit specs
owenpearson Sep 23, 2026
de751b2
test: derive the connection failure and realtime auth unit specs
owenpearson Sep 23, 2026
34d3bc6
test: derive the channel attach and detach unit specs
owenpearson Sep 23, 2026
52e4a97
test: derive the channel publish unit specs
owenpearson Sep 23, 2026
1d5695c
test: derive the channel subscribe and message field unit specs
owenpearson Sep 23, 2026
3824e67
test: derive the channel option, property and collection unit specs
owenpearson Sep 23, 2026
bdc33e4
test: derive the channel connection state and event unit specs
owenpearson Sep 23, 2026
e10c6b2
test: derive the presence enter, subscribe and get unit specs
owenpearson Sep 23, 2026
36a168f
test: derive the channel annotation, delta and message version specs
owenpearson Sep 23, 2026
df7afef
test: derive the presence map and sync unit specs
owenpearson Sep 24, 2026
30ccecf
test: derive the presence channel state, reentry and history specs
owenpearson Sep 24, 2026
76ff2ee
test: derive the heartbeat, ping, fallback and recovery unit specs
owenpearson Sep 24, 2026
35e7e56
docs: record the realtime deviations and the specification faults found
owenpearson Sep 24, 2026
af2dd88
docs: record that time is read through the clock
owenpearson Sep 25, 2026
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
385 changes: 375 additions & 10 deletions .claude/skills/uts-to-python/SKILL.md

Large diffs are not rendered by default.

20 changes: 20 additions & 0 deletions CLAUDE.md
Original file line number Diff line number Diff line change
@@ -1,2 +1,22 @@
- after making any code changes, run `uv ruff check` to make sure linting passes
- use `uv` to run any other necessary tasks such as `pytest`

# Time

Production code in `ably/` reads the current time and schedules delayed callbacks through
the clock, never through `time`, `datetime.now` or `asyncio.sleep`:

- `select_clock(options)` in the constructor, from `ably.util.clock`
- `now_ms()` for the time of day
- `monotonic_ms()` for a duration
- `timer(timeout_ms, callback)` for a delayed callback

A test replaces it with `TestOptions(clock=...)`. Converting a caller-supplied value, and
arithmetic on the unix epoch, are not clock readings and stay as they are.

`asyncio.wait_for` in `ConnectionManager.ping` is the one delay that runs on the event
loop's clock rather than this one; a test which waits it out sets a short real
`realtime_request_timeout`.

`Auth._timestamp` signs a token request the server validates, so a client which reads a
fake clock talks to a mock rather than to the sandbox.
7 changes: 5 additions & 2 deletions ably/realtime/channel.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
from ably.types.mixins import DecodingContext
from ably.types.operations import MessageOperation, PublishResult, UpdateDeleteResult
from ably.types.presence import PresenceMessage
from ably.util.clock import select_clock
from ably.util.eventemitter import EventEmitter
from ably.util.exceptions import AblyException, IncompatibleClientIdException
from ably.util.helper import Timer, is_callable_or_coroutine, validate_message_size
Expand Down Expand Up @@ -58,6 +59,7 @@ def __init__(self, realtime: AblyRealtime, name: str, channel_options: ChannelOp
EventEmitter.__init__(self)
self.__name = name
self.__realtime = realtime
self.__clock = select_clock(realtime.options)
self.__state = ChannelState.INITIALIZED
self.__message_emitter = EventEmitter()
self.__state_timer: Timer | None = None
Expand Down Expand Up @@ -846,7 +848,7 @@ def on_timeout() -> None:
self.__state_timer = None
self.__timeout_pending_state()

self.__state_timer = Timer(self.__realtime.options.realtime_request_timeout, on_timeout)
self.__state_timer = self.__clock.timer(self.__realtime.options.realtime_request_timeout, on_timeout)

def __clear_state_timer(self) -> None:
if self.__state_timer:
Expand All @@ -866,7 +868,8 @@ def __start_retry_timer(self) -> None:
if self.__retry_timer:
return

self.__retry_timer = Timer(self.ably.options.channel_retry_timeout, self.__on_retry_timer_expire)
self.__retry_timer = self.__clock.timer(
self.ably.options.channel_retry_timeout, self.__on_retry_timer_expire)

def __cancel_retry_timer(self) -> None:
if self.__retry_timer:
Expand Down
20 changes: 11 additions & 9 deletions ably/realtime/connectionmanager.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,6 @@

import asyncio
import logging
import time
from collections import deque
from itertools import zip_longest
from typing import TYPE_CHECKING
Expand All @@ -16,6 +15,7 @@
from ably.types.connectionstate import ConnectionEvent, ConnectionState, ConnectionStateChange
from ably.types.operations import PublishResult
from ably.types.tokendetails import TokenDetails
from ably.util.clock import Clock, select_clock
from ably.util.eventemitter import EventEmitter
from ably.util.exceptions import AblyException, IncompatibleClientIdException
from ably.util.helper import Timer, get_random_id, is_token_error
Expand Down Expand Up @@ -123,23 +123,25 @@ def clear(self) -> None:
class PendingPing:
"""Represents a ping awaiting its heartbeat echo from the server"""

def __init__(self, id: str):
def __init__(self, id: str, clock: Clock):
self.id = id
self.clock = clock
self.future: asyncio.Future[None] = asyncio.get_running_loop().create_future()
self.start_time: float = time.monotonic()
self.start_time: float = clock.monotonic_ms()

@property
def message(self) -> dict:
return {"action": ProtocolMessageAction.HEARTBEAT, "id": self.id}

@property
def response_time_ms(self) -> float:
return round((time.monotonic() - self.start_time) * 1000, 2)
return round(self.clock.monotonic_ms() - self.start_time, 2)


class ConnectionManager(EventEmitter):
def __init__(self, realtime: AblyRealtime, initial_state):
self.options = realtime.options
self.clock = select_clock(self.options)
self.__ably = realtime
self.__state: ConnectionState = initial_state
self.__pending_pings: dict[str, PendingPing] = {}
Expand Down Expand Up @@ -364,7 +366,7 @@ async def ping(self) -> float:
raise AblyException("Cannot send ping request. Calling ping in invalid state", 400, 40000)

# RTN13e: the id tells this ping's echo apart from server heartbeats and other pings
pending_ping = PendingPing(get_random_id())
pending_ping = PendingPing(get_random_id(), self.clock)
self.__pending_pings[pending_ping.id] = pending_ping
try:
# RTN13d: while connecting, the heartbeat goes out once the connection is established
Expand All @@ -379,7 +381,7 @@ async def ping(self) -> float:
return pending_ping.response_time_ms

async def __send_heartbeat(self, pending_ping: PendingPing) -> None:
pending_ping.start_time = time.monotonic()
pending_ping.start_time = self.clock.monotonic_ms()
try:
await self.send_protocol_message(pending_ping.message)
except Exception as error:
Expand Down Expand Up @@ -718,7 +720,7 @@ def on_transition_timer_expire():

log.debug(f'ConnectionManager.start_transition_timer(): setting timer for {timeout}ms')

self.transition_timer = Timer(timeout, on_transition_timer_expire)
self.transition_timer = self.clock.timer(timeout, on_transition_timer_expire)

def cancel_transition_timer(self):
log.debug('ConnectionManager.cancel_transition_timer()')
Expand All @@ -741,7 +743,7 @@ def on_suspend_timer_expire() -> None:
)
self.__fail_state = ConnectionState.SUSPENDED

self.suspend_timer = Timer(Defaults.connection_state_ttl, on_suspend_timer_expire)
self.suspend_timer = self.clock.timer(Defaults.connection_state_ttl, on_suspend_timer_expire)

def check_suspend_timer(self, state: ConnectionState) -> None:
if state not in (
Expand All @@ -764,7 +766,7 @@ def on_retry_timeout():
self.retry_timer = None
self.request_state(ConnectionState.CONNECTING)

self.retry_timer = Timer(interval, on_retry_timeout)
self.retry_timer = self.clock.timer(interval, on_retry_timeout)

def cancel_retry_timer(self) -> None:
if self.retry_timer:
Expand Down
21 changes: 15 additions & 6 deletions ably/transport/websockettransport.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,9 +14,9 @@
from ably.transport.defaults import Defaults
from ably.types.connectiondetails import ConnectionDetails
from ably.types.operations import PublishResult
from ably.util.clock import select_clock
from ably.util.eventemitter import EventEmitter
from ably.util.exceptions import AblyException
from ably.util.helper import Timer, unix_time_ms

try:
# websockets 15+ preferred imports
Expand Down Expand Up @@ -68,6 +68,8 @@ def __init__(self, connection_manager: ConnectionManager, host: str, params: dic
self.ws_connect_task: asyncio.Task | None = None
self.connection_manager = connection_manager
self.options = self.connection_manager.options
self.connect_func = self.__select_connect_func(self.options)
self.clock = select_clock(self.options)
self.is_connected = False
self.idle_timer = None
self.last_activity = None
Expand All @@ -78,6 +80,13 @@ def __init__(self, connection_manager: ConnectionManager, host: str, params: dic
self.format = params.get('format', 'json')
super().__init__()

@staticmethod
def __select_connect_func(options):
test_options = getattr(options, '_test_options', None)
if test_options is not None and test_options.websocket_connect is not None:
return test_options.websocket_connect
return ws_connect

def connect(self):
headers = HttpUtils.default_headers()
query_params = urllib.parse.urlencode(self.params)
Expand All @@ -103,11 +112,11 @@ async def ws_connect(self, ws_url, headers):
try:
# Use additional_headers for websockets 15+, fallback to extra_headers for older versions
try:
async with ws_connect(ws_url, additional_headers=headers) as websocket:
async with self.connect_func(ws_url, additional_headers=headers) as websocket:
await self._handle_websocket_connection(ws_url, websocket)
except TypeError:
# Fallback for websockets 14 and earlier
async with ws_connect(ws_url, extra_headers=headers) as websocket:
async with self.connect_func(ws_url, extra_headers=headers) as websocket:
await self._handle_websocket_connection(ws_url, websocket)
except (WebSocketException, socket.gaierror) as e:
exception = AblyException(f'Error opening websocket connection: {e}', 400, 40000)
Expand Down Expand Up @@ -294,11 +303,11 @@ async def send(self, message: dict):
def set_idle_timer(self, timeout: float):
if self.idle_timer:
self.idle_timer.cancel()
self.idle_timer = Timer(timeout, self.on_idle_timer_expire)
self.idle_timer = self.clock.timer(timeout, self.on_idle_timer_expire)

async def on_idle_timer_expire(self):
self.idle_timer = None
since_last = unix_time_ms() - self.last_activity
since_last = self.clock.now_ms() - self.last_activity
time_remaining = self.max_idle_interval - since_last
msg = f"No activity seen from realtime in {since_last} ms; assuming connection has dropped"
if time_remaining <= 0:
Expand All @@ -310,7 +319,7 @@ async def on_idle_timer_expire(self):
def on_activity(self):
if not self.max_idle_interval:
return
self.last_activity = unix_time_ms()
self.last_activity = self.clock.now_ms()
self.set_idle_timer(self.max_idle_interval + 100)

async def disconnect(self, reason=None):
Expand Down
9 changes: 8 additions & 1 deletion ably/types/testoptions.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,12 @@ class TestOptions:
:Parameters:
- `http_transport`: an `httpx.AsyncBaseTransport` which handles every
HTTP request the client makes, in place of the network.
- `websocket_connect`: a callable which opens every realtime websocket
connection the client makes, in place of `websockets.connect`. It is
called as `websocket_connect(url, additional_headers=headers)`, or with
`extra_headers=headers` if that raises `TypeError`, and returns an async
context manager yielding an object supporting `__aiter__`, `send` and
`close`.
- `clock`: the source of time the client reads, in place of
`ably.util.clock.Clock`. It supplies `now_ms()`, `monotonic_ms()` and
`timer(timeout_ms, callback)`, where `callback` is either a coroutine
Expand All @@ -17,6 +23,7 @@ class TestOptions:
# any module importing it as declaring a test suite.
__test__ = False

def __init__(self, http_transport=None, clock=None):
def __init__(self, http_transport=None, websocket_connect=None, clock=None):
self.http_transport = http_transport
self.websocket_connect = websocket_connect
self.clock = clock
6 changes: 3 additions & 3 deletions ably/util/clock.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,9 +7,9 @@
class Clock:
"""The source of time for every decision a client makes from it.

Token expiry, the server-time offset and the fallback-host cache all read
the clock, and delayed callbacks are scheduled through it, so replacing one
moves all of them together.
Token expiry, the server-time offset, the fallback-host cache and the
transport's idle detection all read the clock, and every delayed callback
is scheduled through it, so replacing one moves all of them together.

`now_ms` is the time of day and may step; `monotonic_ms` only ever moves
forward and is what a duration is measured with.
Expand Down
1 change: 1 addition & 0 deletions ably/util/helper.py
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,7 @@ async def _job(self):
def cancel(self):
self._task.cancel()


def validate_message_size(encoded_messages: list, use_binary_protocol: bool, max_message_size: int) -> None:
"""Validate that encoded messages don't exceed the maximum size limit.

Expand Down
98 changes: 87 additions & 11 deletions test/uts/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,25 +7,32 @@ carries a `# UTS: <id>` comment identifying the specification it came from.

Read `uts/docs/writing-derived-tests.md` in the specification repository before
adding or changing tests here, alongside `.claude/skills/uts-to-python/SKILL.md`,
which covers what is particular to this SDK. Record anything that departs from a
specification in [deviations.md](deviations.md), which also covers how the
specifications are adopted here and why.
which covers what is particular to this SDK and lists every helper below.

Record anything that departs from a specification in [deviations.md](deviations.md),
which also covers the faults found in the specifications themselves and raised
upstream, and the choices behind how the specifications are adopted here.

## Layout

```
helpers/ shared infrastructure the specifications assume
helpers/ shared infrastructure the specifications assume, and its own tests
rest/ specifications under uts/rest
realtime/ specifications under uts/realtime
```

Unit tests serve every request from a mock and reach no network. Integration
tests run against a sandbox app.
Every directory needs an `__init__.py`, because `test` is a package.

Unit tests serve every request from a mock and reach no network — neither the REST
suite nor the realtime one. The seams are installed per client, so a test holds to
that by installing them; one that omits a seam, or that lets the host fallback loop
run, reaches the real internet. Integration tests run against a sandbox app.

## Installing the mock
## Installing the mocks

The specifications express mock installation as a global `install_mock(mock_http)`.
Here a mock is passed to the client it serves:
Here a mock is passed to the client it serves, through `TestOptions`. There are three
seams, all client-scoped; [deviations.md](deviations.md) says why.

```python
mock_http = MockHttpClient(
Expand All @@ -35,11 +42,80 @@ mock_http = MockHttpClient(
ably = AblyRest(key=key, _test_options=TestOptions(http_transport=mock_http.as_transport()))
```

The client builds its HTTP client once, so construct the mock first. Teardown is
`await ably.close()`, which stands in for `uninstall_mock()`.
A realtime client takes its websocket mock the same way, through
`TestOptions(websocket_connect=...)`, and the time it reads through
`TestOptions(clock=...)`:

```python
mock_ws = MockWebSocket(
on_connection_attempt=lambda conn: conn.respond_with_success(CONNECTED_MESSAGE),
)
ably = AblyRealtime(key=key, auto_connect=False,
_test_options=TestOptions(websocket_connect=mock_ws.as_connect(),
clock=FakeClock()))
```

`rest_client(mock_http, ...)` and `realtime_client(mock_ws, ...)` in
[helpers/client.py](helpers/client.py) wrap all three, defaulting the credentials
and registering the client for teardown:

```python
realtime_client(mock_websocket=None, mock_http=None, clock=None, **kwargs)
```

`mock_http=` gives a realtime client an HTTP mock for the specifications that drive
REST over a realtime client; `clock=` installs a `FakeClock`. `realtime_client`
defaults `auto_connect` to **false** and `fallback_hosts` to **empty** — both
deliberate, and both explained in [deviations.md](deviations.md).

The client builds its HTTP client once and reads its websocket hook once, so
construct the mocks first. Teardown is automatic: `conftest.py` closes every
registered client after each test, which stands in for `uninstall_mock()`. **Do not
close clients in a test** — the fixture survives every connection state, and a test
that closes its own leaves nothing to clean up if it fails first.

## What the helpers offer

| Module | Holds |
|---|---|
| [helpers/mock_http.py](helpers/mock_http.py) | `MockHttpClient`, matching `uts/rest/unit/helpers/mock_http.md`, including the superseded `queue_*` family |
| [helpers/mock_websocket.py](helpers/mock_websocket.py) | `MockWebSocket`, matching `uts/realtime/unit/helpers/mock_websocket.md`, plus the protocol-message templates and builders the specifications assume |
| [helpers/client.py](helpers/client.py) | client constructors, and the `AWAIT_STATE` / `AWAIT UNTIL` equivalents |
| [helpers/clock.py](helpers/clock.py) | `FakeClock`, `settle()` and `advance_to_connection_state()` — `enable_fake_timers()` and `ADVANCE_TIME(ms)` |
| [helpers/presence.py](helpers/presence.py) | the presence-map stubs and wire-message builders the presence specifications share |
| [helpers/deviations.py](helpers/deviations.py) | the `@deviation` and `@spec_error` gates |

`SKILL.md` lists every name in each. The helpers have their own tests
(`helpers/*_test.py`), which are not derived from a specification and are not counted
in the derived-test totals.

## Running

```
uv run --extra crypto pytest test/uts
uv run --frozen --extra crypto --extra dev pytest test/uts -q
```

`--frozen` is required: without it dependency resolution reaches past the
environment's cutoff. `--extra dev` carries pytest.

Tests that record a deviation or a specification fault are skipped by default and run
under an environment variable:

```
RUN_DEVIATIONS=1 uv run --frozen --extra crypto --extra dev pytest test/uts -q
```

Every gated test is confirmed to fail when enabled — none passes under both
behaviours — so the two runs are the check that the record in
[deviations.md](deviations.md) is still true. The counts either run should produce are
in that file's header.

Both seams are installed per client, so a test that forgets one, or that lets the host
fallback loop run, reaches the real internet; see the fallback host note in
[deviations.md](deviations.md).

Linting is `ruff`, line length 115:

```
uv run --frozen --extra crypto --extra dev ruff check ably/ test/
```
Loading
Loading