diff --git a/ably/pubsub/http/channel.py b/ably/pubsub/http/channel.py index 531a17cd..321ea275 100644 --- a/ably/pubsub/http/channel.py +++ b/ably/pubsub/http/channel.py @@ -1,7 +1,5 @@ -import base64 import json import logging -import os from collections import OrderedDict from typing import Iterator, Optional from urllib import parse @@ -15,6 +13,7 @@ Message, MessageAction, MessageVersion, + assign_idempotent_ids, make_message_response_handler, make_single_message_response_handler, ) @@ -58,11 +57,7 @@ def __publish_request_body(self, messages): """ # Idempotent publishing if self.ably.options.idempotent_rest_publishing: - # RSL1k1 - if all(message.id is None for message in messages): - base_id = base64.urlsafe_b64encode(os.urandom(12)).decode() - for serial, message in enumerate(messages): - message.id = f'{base_id}:{serial}' + assign_idempotent_ids(messages) request_body_list = [] for m in messages: diff --git a/ably/pubsub/http/http.py b/ably/pubsub/http/http.py index 457c0635..7d912d0b 100644 --- a/ably/pubsub/http/http.py +++ b/ably/pubsub/http/http.py @@ -1,13 +1,23 @@ +import json import logging -from typing import Optional +from typing import List, Optional, Union from urllib.parse import urlencode +import msgpack + from ably.pubsub.http.auth import Auth from ably.pubsub.http.channel import Channels from ably.pubsub.http.push import Push from ably.pubsub.prototypes import PubSubHttpClient from ably.pubsub.request.http import Http from ably.pubsub.request.paginatedresult import HttpPaginatedResponse, PaginatedResult, format_params +from ably.pubsub.types.batch import ( + BatchPublishSpec, + BatchResult, + batch_presence_result_from_dict, + batch_publish_result_from_dict, +) +from ably.pubsub.types.message import assign_idempotent_ids from ably.pubsub.types.options import Options from ably.pubsub.types.stats import stats_response_processor from ably.pubsub.types.tokendetails import TokenDetails @@ -104,6 +114,63 @@ async def time(self, timeout: Optional[float] = None) -> float: AblyException.raise_for_response(r) return r.to_native()[0] + async def batch_publish(self, specs) -> Union[BatchResult, List[BatchResult]]: + """Publishes messages to several channels in a single request (RSC22) + + :Parameters: + - `specs`: a `BatchPublishSpec`, or a list of them. A dict with + `channels` and `messages` keys may stand in for a `BatchPublishSpec`. + + Returns a `BatchResult` holding a `BatchPublishSuccessResult` or a + `BatchPublishFailureResult` for each channel the spec names, or, given + a list of specs, a list of `BatchResult`s in the same order. + """ + single = not isinstance(specs, (list, tuple)) + specs = [BatchPublishSpec.factory(spec) for spec in ([specs] if single else specs)] + + for spec in specs: + if not spec.channels: + raise AblyException('A batch publish spec must name at least one channel', 400, 40000) + if not spec.messages: + raise AblyException('A batch publish spec must include at least one message', 400, 40000) + + # RSC22d: RSL1k1 applies to each spec separately + if self.options.idempotent_rest_publishing: + for spec in specs: + assign_idempotent_ids(spec.messages) + + binary = self.options.use_binary_protocol + body = [spec.as_dict(binary=binary) for spec in specs] + if single: + body = body[0] + if binary: + body = msgpack.packb(body, use_bin_type=True) + else: + body = json.dumps(body, separators=(',', ':')) + + response = await self.http.post('/messages', body=body) + + # RSC22b: the response is an array with a result for each spec, even when a single spec was sent + results = [BatchResult.from_dict(result, batch_publish_result_from_dict) + for result in response.to_native()] + return results[0] if single else results + + async def batch_presence(self, channels: List[str]) -> BatchResult: + """Retrieves the members present on several channels in a single request (RSC24) + + :Parameters: + - `channels`: the names of the channels + + Returns a `BatchResult` holding a `BatchPresenceSuccessResult` or a + `BatchPresenceFailureResult` for each channel. + """ + if isinstance(channels, str): + raise TypeError('Unexpected str channels, expected a list of channel names') + + path = '/presence?' + urlencode({'channels': ','.join(channels)}) + response = await self.http.get(path) + return BatchResult.from_dict(response.to_native(), batch_presence_result_from_dict) + @property def client_id(self) -> Optional[str]: return self.options.client_id diff --git a/ably/pubsub/prototypes.py b/ably/pubsub/prototypes.py index f269ab7f..4fc2613e 100644 --- a/ably/pubsub/prototypes.py +++ b/ably/pubsub/prototypes.py @@ -28,6 +28,7 @@ from ably.pubsub.realtime.channel import Channels as RealtimeChannels from ably.pubsub.realtime.connection import Connection from ably.pubsub.request.paginatedresult import HttpPaginatedResponse, PaginatedResult + from ably.pubsub.types.batch import BatchPublishSpec, BatchResult from ably.pubsub.types.options import Options @@ -72,6 +73,15 @@ async def time(self, timeout: float | None = None) -> float: """Return the current server time in ms since the unix epoch.""" ... + async def batch_publish(self, specs: BatchPublishSpec | dict | list[BatchPublishSpec | dict] + ) -> BatchResult | list[BatchResult]: + """Publish messages to one or more channels in a single request.""" + ... + + async def batch_presence(self, channels: list[str]) -> BatchResult: + """Return the presence members of several channels in a single request.""" + ... + async def request(self, method: str, path: str, version: str, params: dict | None = None, body=None, headers=None) -> HttpPaginatedResponse: """Make an arbitrary request against the Ably REST API.""" diff --git a/ably/pubsub/server/__init__.py b/ably/pubsub/server/__init__.py index 41efc61e..dcc0f2bf 100644 --- a/ably/pubsub/server/__init__.py +++ b/ably/pubsub/server/__init__.py @@ -28,6 +28,14 @@ from ably.pubsub.realtime.realtime import DefaultPubSubRealtimeClient as _DefaultPubSubRealtimeClient from ably.pubsub.request.paginatedresult import HttpPaginatedResponse, PaginatedResult from ably.pubsub.types.annotation import Annotation, AnnotationAction +from ably.pubsub.types.batch import ( + BatchPresenceFailureResult, + BatchPresenceSuccessResult, + BatchPublishFailureResult, + BatchPublishSpec, + BatchPublishSuccessResult, + BatchResult, +) from ably.pubsub.types.capability import Capability from ably.pubsub.types.channelmode import ChannelMode from ably.pubsub.types.channeloptions import ChannelOptions @@ -248,6 +256,12 @@ def create_realtime_client(**kwargs) -> PubSubRealtimeClient: 'Annotation', 'AnnotationAction', 'Auth', + 'BatchPresenceFailureResult', + 'BatchPresenceSuccessResult', + 'BatchPublishFailureResult', + 'BatchPublishSpec', + 'BatchPublishSuccessResult', + 'BatchResult', 'Capability', 'ChannelMode', 'ChannelOptions', diff --git a/ably/pubsub/server/sync.py b/ably/pubsub/server/sync.py index 5ccab431..d57b5cd0 100644 --- a/ably/pubsub/server/sync.py +++ b/ably/pubsub/server/sync.py @@ -14,6 +14,14 @@ from ably.pubsub.sync.prototypes import PubSubHttpClient from ably.pubsub.sync.request.paginatedresult import HttpPaginatedResponseSync, PaginatedResultSync from ably.pubsub.sync.types.annotation import Annotation, AnnotationAction +from ably.pubsub.sync.types.batch import ( + BatchPresenceFailureResult, + BatchPresenceSuccessResult, + BatchPublishFailureResult, + BatchPublishSpec, + BatchPublishSuccessResult, + BatchResult, +) from ably.pubsub.sync.types.capability import Capability from ably.pubsub.sync.types.channelmode import ChannelMode from ably.pubsub.sync.types.channeloptions import ChannelOptions @@ -61,6 +69,12 @@ def create_http_client(**kwargs) -> PubSubHttpClient: 'Annotation', 'AnnotationAction', 'AuthSync', + 'BatchPresenceFailureResult', + 'BatchPresenceSuccessResult', + 'BatchPublishFailureResult', + 'BatchPublishSpec', + 'BatchPublishSuccessResult', + 'BatchResult', 'Capability', 'ChannelMode', 'ChannelOptions', diff --git a/ably/pubsub/types/batch.py b/ably/pubsub/types/batch.py new file mode 100644 index 00000000..d6d37735 --- /dev/null +++ b/ably/pubsub/types/batch.py @@ -0,0 +1,208 @@ +from typing import Generic, List, TypeVar + +from ably.pubsub.types.presence import PresenceMessage +from ably.pubsub.util.exceptions import AblyException + +T = TypeVar('T') + + +class BatchResult(Generic[T]): + """The results of a batch operation, one per channel (BAR2).""" + + def __init__(self, success_count: int, failure_count: int, results: List[T]): + self.__success_count = success_count + self.__failure_count = failure_count + self.__results = results + + @property + def success_count(self) -> int: + """The number of successful operations (BAR2a).""" + return self.__success_count + + @property + def failure_count(self) -> int: + """The number of unsuccessful operations (BAR2b).""" + return self.__failure_count + + @property + def results(self) -> List[T]: + """The result of each operation in the batch (BAR2c).""" + return self.__results + + @staticmethod + def from_dict(obj, result_from_dict): + """Create a BatchResult from a response envelope, building each result with `result_from_dict`.""" + return BatchResult( + success_count=obj.get('successCount'), + failure_count=obj.get('failureCount'), + results=[result_from_dict(result) for result in obj.get('results') or []], + ) + + +class BatchPublishSpec: + """The messages a batch publish sends, and the channels it sends every one of them to (BSP2).""" + + def __init__(self, channels, messages): + """ + Args: + channels: The names of the channels to publish the messages to (BSP2a). + messages: The `Message` objects to publish to each channel (BSP2b). + """ + self.__channels = channels + self.__messages = messages + + @property + def channels(self): + return self.__channels + + @property + def messages(self): + return self.__messages + + def as_dict(self, binary=False): + """Convert BatchPublishSpec to the wire format, encoding each message per RSL4.""" + return { + 'channels': list(self.channels), + 'messages': [message.as_dict(binary=binary) for message in self.messages], + } + + @classmethod + def factory(cls, spec): + """A BatchPublishSpec from either an instance or a dict with `channels` and `messages` keys.""" + if isinstance(spec, cls): + return spec + if isinstance(spec, dict): + return cls(channels=spec.get('channels'), messages=spec.get('messages')) + raise TypeError(f'Unexpected {type(spec)} batch publish spec, expected a BatchPublishSpec or a dict') + + +class BatchPublishSuccessResult: + """The result of publishing a spec's messages to one channel (BPR2).""" + + def __init__(self, channel, message_id, serials): + """ + Args: + channel: The name of the channel (BPR2a). + message_id: The id prefix shared by the published messages (BPR2b). + serials: The serial of each published message, or None where a conflation rule discarded it (BPR2c). + """ + self.__channel = channel + self.__message_id = message_id + self.__serials = serials + + @property + def channel(self): + return self.__channel + + @property + def message_id(self): + return self.__message_id + + @property + def serials(self): + return self.__serials + + @staticmethod + def from_dict(obj): + return BatchPublishSuccessResult( + channel=obj.get('channel'), + message_id=obj.get('messageId'), + serials=obj.get('serials') or [], + ) + + +class BatchPublishFailureResult: + """The reason a spec's messages could not be published to one channel (BPF2).""" + + def __init__(self, channel, error): + """ + Args: + channel: The name of the channel (BPF2a). + error: An `AblyException` describing why the publish failed (BPF2b). + """ + self.__channel = channel + self.__error = error + + @property + def channel(self): + return self.__channel + + @property + def error(self): + return self.__error + + @staticmethod + def from_dict(obj): + return BatchPublishFailureResult( + channel=obj.get('channel'), + error=AblyException.from_dict(obj['error']), + ) + + +class BatchPresenceSuccessResult: + """The members present on one channel of a batch presence request (BGR2).""" + + def __init__(self, channel, presence): + """ + Args: + channel: The name of the channel (BGR2a). + presence: A `PresenceMessage` for each member present on the channel (BGR2b). + """ + self.__channel = channel + self.__presence = presence + + @property + def channel(self): + return self.__channel + + @property + def presence(self): + return self.__presence + + @staticmethod + def from_dict(obj): + # The server leaves `presence` out for a channel with no members + return BatchPresenceSuccessResult( + channel=obj.get('channel'), + presence=PresenceMessage.from_encoded_array(obj.get('presence') or []), + ) + + +class BatchPresenceFailureResult: + """The reason the members of one channel of a batch presence request could not be retrieved (BGF2).""" + + def __init__(self, channel, error): + """ + Args: + channel: The name of the channel (BGF2a). + error: An `AblyException` describing why the request failed for the channel (BGF2b). + """ + self.__channel = channel + self.__error = error + + @property + def channel(self): + return self.__channel + + @property + def error(self): + return self.__error + + @staticmethod + def from_dict(obj): + return BatchPresenceFailureResult( + channel=obj.get('channel'), + error=AblyException.from_dict(obj['error']), + ) + + +def batch_publish_result_from_dict(obj): + if obj.get('error') is not None: + return BatchPublishFailureResult.from_dict(obj) + return BatchPublishSuccessResult.from_dict(obj) + + +def batch_presence_result_from_dict(obj): + if obj.get('error') is not None: + return BatchPresenceFailureResult.from_dict(obj) + return BatchPresenceSuccessResult.from_dict(obj) diff --git a/ably/pubsub/types/message.py b/ably/pubsub/types/message.py index 93a20f0c..ad574154 100644 --- a/ably/pubsub/types/message.py +++ b/ably/pubsub/types/message.py @@ -1,4 +1,6 @@ +import base64 import logging +import os from enum import IntEnum from ably.pubsub.types.mixins import DeltaExtras, EncodeDataMixin @@ -391,6 +393,14 @@ def update_inner_message_fields(proto_msg: dict): msg_index = msg_index + 1 +def assign_idempotent_ids(messages): + """Gives every message a library-generated id, unless one of them already has an id (RSL1k1).""" + if all(message.id is None for message in messages): + base_id = base64.urlsafe_b64encode(os.urandom(12)).decode() + for serial, message in enumerate(messages): + message.id = f'{base_id}:{serial}' + + def make_message_response_handler(cipher): def encrypted_message_response_handler(response): messages = response.to_native() diff --git a/test/ably/http/httpbatch_test.py b/test/ably/http/httpbatch_test.py new file mode 100644 index 00000000..f362d5cd --- /dev/null +++ b/test/ably/http/httpbatch_test.py @@ -0,0 +1,112 @@ +import logging + +import pytest + +from ably.pubsub.server import AblyException +from ably.pubsub.types.batch import ( + BatchPresenceSuccessResult, + BatchPublishFailureResult, + BatchPublishSpec, + BatchPublishSuccessResult, + BatchResult, +) +from ably.pubsub.types.message import Message +from test.ably.testapp import TestApp +from test.ably.utils import BaseAsyncTestCase, VaryByProtocolTestsMetaclass + +log = logging.getLogger(__name__) + + +class TestHttpBatch(BaseAsyncTestCase, metaclass=VaryByProtocolTestsMetaclass): + + @pytest.fixture(autouse=True) + async def setup(self): + self.test_vars = await TestApp.get_test_vars() + self.ably = await TestApp.get_ably_rest() + # keys[2] may publish to `canpublish:*` and to nothing else + self.ably_restricted = await TestApp.get_ably_rest(key=self.test_vars['keys'][2]['key_str']) + yield + await self.ably.close() + await self.ably_restricted.close() + + def per_protocol_setup(self, use_binary_protocol): + self.ably.options.use_binary_protocol = use_binary_protocol + self.ably_restricted.options.use_binary_protocol = use_binary_protocol + self.use_binary_protocol = use_binary_protocol + + async def test_batch_publish_single_spec(self): + channel_names = [self.get_channel_name('persisted:batch_publish_a'), + self.get_channel_name('persisted:batch_publish_b')] + + result = await self.ably.batch_publish(BatchPublishSpec(channels=channel_names, messages=[ + Message('string', 'This is a string message payload'), + Message('binary', b'This is a byte[] message payload'), + Message('json', {'test': 'This is a JSONObject message payload'}), + ])) + + assert isinstance(result, BatchResult) + assert result.success_count == 2 + assert result.failure_count == 0 + assert [entry.channel for entry in result.results] == channel_names + for entry in result.results: + assert isinstance(entry, BatchPublishSuccessResult) + assert len(entry.serials) == 3 + + for channel_name in channel_names: + history = await self.ably.channels[channel_name].history(direction='forwards') + messages = history.items + assert [message.name for message in messages] == ['string', 'binary', 'json'] + assert messages[0].data == 'This is a string message payload' + assert bytes(messages[1].data) == b'This is a byte[] message payload' + assert messages[2].data == {'test': 'This is a JSONObject message payload'} + # The library-generated ids share the prefix the result reports + assert [message.id for message in messages] == [f'{result.results[0].message_id}:{serial}' + for serial in range(3)] + + async def test_batch_publish_array_of_specs(self): + channel_a = self.get_channel_name('batch_publish_a') + channel_b = self.get_channel_name('batch_publish_b') + + results = await self.ably.batch_publish([ + BatchPublishSpec(channels=[channel_a], messages=[Message('a', 'data-a')]), + {'channels': [channel_a, channel_b], 'messages': [Message('b', 'data-b')]}, + ]) + + assert isinstance(results, list) + assert len(results) == 2 + assert [entry.channel for entry in results[0].results] == [channel_a] + assert [entry.channel for entry in results[1].results] == [channel_a, channel_b] + assert results[0].results[0].message_id != results[1].results[0].message_id + + async def test_batch_publish_partial_failure(self): + allowed_channel = self.get_channel_name('canpublish:batch_publish') + denied_channel = self.get_channel_name('batch_publish') + + result = await self.ably_restricted.batch_publish(BatchPublishSpec( + channels=[allowed_channel, denied_channel], + messages=[Message('event', 'data')], + )) + + assert result.success_count == 1 + assert result.failure_count == 1 + + success, failure = result.results + assert isinstance(success, BatchPublishSuccessResult) + assert success.channel == allowed_channel + assert isinstance(failure, BatchPublishFailureResult) + assert failure.channel == denied_channel + assert isinstance(failure.error, AblyException) + assert failure.error.code == 40160 + assert failure.error.status_code == 401 + + async def test_batch_presence_empty_channels(self): + channel_names = [self.get_channel_name('batch_presence_a'), self.get_channel_name('batch_presence_b')] + + result = await self.ably.batch_presence(channel_names) + + assert result.success_count == 2 + assert result.failure_count == 0 + assert [entry.channel for entry in result.results] == channel_names + for entry in result.results: + assert isinstance(entry, BatchPresenceSuccessResult) + assert entry.presence == [] diff --git a/test/unit/batch_test.py b/test/unit/batch_test.py new file mode 100644 index 00000000..8300c3bd --- /dev/null +++ b/test/unit/batch_test.py @@ -0,0 +1,38 @@ +"""Unit tests for the batch types and the inputs the batch methods accept. + +Tests cover: +- BSP2: a dict standing in for a BatchPublishSpec +- RSC22, RSC24: inputs refused before any request is made +""" + +import pytest + +from ably.pubsub.server import create_http_client +from ably.pubsub.types.batch import BatchPublishSpec +from ably.pubsub.types.message import Message + + +def test_batch_publish_spec_factory_accepts_spec(): + spec = BatchPublishSpec(channels=['a'], messages=[Message('event', 'data')]) + assert BatchPublishSpec.factory(spec) is spec + + +def test_batch_publish_spec_factory_accepts_dict(): + messages = [Message('event', 'data')] + spec = BatchPublishSpec.factory({'channels': ['a', 'b'], 'messages': messages}) + assert spec.channels == ['a', 'b'] + assert spec.messages is messages + + +def test_batch_publish_spec_factory_rejects_other_types(): + with pytest.raises(TypeError): + BatchPublishSpec.factory('a') + + +async def test_batch_presence_rejects_a_single_channel_name(): + client = create_http_client(key='fake.key:secret') + try: + with pytest.raises(TypeError): + await client.batch_presence('channel') + finally: + await client.close() diff --git a/test/uts/deviations.md b/test/uts/deviations.md index 498f0ad3..2eb9daa1 100644 --- a/test/uts/deviations.md +++ b/test/uts/deviations.md @@ -31,8 +31,8 @@ Variants` section and run every one of their tests twice, once per protocol, and `rest/unit` tests are parametrized over a table of fixtures the specification gives inline. That turns 1141 derived tests into 1234 pytest cases. -Of **1132 Test IDs, derived as 1141 tests and run as 1234 pytest cases**: 904 Test IDs -(913 tests, 1001 cases) pass, 213 (213 tests, 218 cases) are gated behind +Of **1132 Test IDs, derived as 1141 tests and run as 1234 pytest cases**: 947 Test IDs +(956 tests, 1047 cases) pass, 170 (170 tests, 172 cases) are gated behind `RUN_DEVIATIONS`, and 15 (15 tests, 15 cases) cannot be run at all. The three groups are disjoint: two Test IDs, and one parametrized test, have a gated part and a passing part, and are counted with the gated. Every gated test has been confirmed to fail when @@ -40,18 +40,18 @@ enabled, so none of them passes under both behaviours. 494 of the Test IDs come `uts/rest/unit` (503 tests, 536 cases), 481 from `uts/realtime/unit` (481, 481), 84 from `uts/rest/integration` (84, 122) and 73 from `uts/realtime/integration` (73, 95); 8 of the REST integration ids (8, 8) and 30 of the realtime ones (30, 30) come from the -`proxy` package within each. Of the gated Test IDs 122 are REST and 91 realtime, which is -126 REST cases and 92 realtime. -A further 122 pytest cases under `helpers/` cover the mock infrastructure itself and are +`proxy` package within each. Of the gated Test IDs 79 are REST and 91 realtime, which is +80 REST cases and 92 realtime. +A further 130 pytest cases under `helpers/` cover the mock infrastructure itself and are not derived from a specification. -The 203 gated Test IDs that record SDK non-compliance — 203 tests, 208 cases — reduce to -**71 distinct root causes**, 27 on the REST side and 44 on the realtime side. Three further +The 159 gated Test IDs that record SDK non-compliance — 159 tests, 161 cases — reduce to +**70 distinct root causes**, 26 on the REST side and 44 on the realtime side. Three further defects are recorded below with no test of their own, because the specification's test cannot discriminate (RTP18a), has nothing to assert against (the timezone split on synthesized LEAVE timestamps), or is worked around in the setup of every test that would otherwise trip over it (`enterClient` on an anonymous connection), so the file -carries **74 SDK root causes** in all. The remaining 10 gated Test IDs are +carries **73 SDK root causes** in all. The remaining 10 gated Test IDs are specification faults, and reduce to 7. Entries closed by a fix are removed rather than kept as history; `git log` holds that. @@ -131,6 +131,7 @@ Raised upstream: | [#552](https://github.com/ably/specification/issues/552) | `proxy/connection_resume.md`: a status code neither SDK returns, a proxy substitution that does not exist, and event-log fields the proxy does not emit | | [#553](https://github.com/ably/specification/issues/553) | A heartbeat-starvation test that closes the socket thirteen seconds inside the idle window | | [#554](https://github.com/ably/specification/issues/554) | Two sections provoking one server response, leaving the revoked-key point uncovered | +| [#559](https://github.com/ably/specification/issues/559) | Batch publish and token revocation fixtures in the response format the server sends below protocol version 3 | `#527` also carries a comment on the realtime wire-format assertions, `#532` one on the same housekeeping categories in `realtime/unit`, and @@ -271,16 +272,36 @@ file, and the protocol, which fix LEAVE at 3 and UPDATE at 4. The closing note o `batch_presence.md` states that with `X-Ably-Version >= 3` the server returns a `BatchResult` envelope "for all batch responses" and calls the plain array legacy. -Every mock in `batch_publish.md` uses the plain array. `features.md` RSC22b backs -`batch_publish.md` — "the response will still be an array" — so `batch_presence.md`'s -claim is the one to revisit. `revoke_tokens.md` has the same internal split: -`RSA17c_1` and `TRS2_1` stub a bare array while asserting envelope fields. +Every mock in `batch_publish.md` uses the plain array, and its RSC22_Headers1 pins +`X-Ably-Version: 2`: `batch_publish.md` is written against the legacy protocol, and +`batch_presence.md` is the one that matches the server ably-python talks to. Measured +against the sandbox, `POST /messages` answers: -`batch_publish.md` RSC22_Headers1 also pins `X-Ably-Version: 2` and -`Content-Type: application/json`; CSV2b templates the version, the sibling spec says -">= 3", and the binary protocol default makes the content type msgpack. - -The revocation half of that split is settled by the server. `POST +| `X-Ably-Version` | Every channel succeeds | Any channel fails | +|---|---|---| +| 5 | 201 with an array holding one `{successCount, failureCount, results}` envelope per spec | the same, 201 | +| 2, or none | 201 with one flat array of `{channel, messageId}` across every spec | 400 with error 40020, and the flat array, failures included, as `batchResponse` | + +A spec sent as a bare object is answered like an array of one: an array holding a single +envelope. That is the array RSC22b means when it says the response "will still be an array" +and the single-spec overload "will have to extract the element" — an array of +`BatchResult`s, one per spec, rather than of per-channel results. `GET /presence` follows the +same split by version, answering a mixed batch with 200 and the envelope at version 5 and +with 400/40020 and `batchResponse` at version 2. The legacy shape carries no `serials` +either; `batch_publish.md`'s mocks add them to it. + +ably-python sends version 5, so every `batch_publish.md` mock is corrected to the envelope, +through `batch_result()` in `rest/unit/batch_publish_test.py`, and the assertions stand as +written: they read `result.results[...]`, the layout BAR2c gives the envelope. RSC22_Headers1's +two pins are corrected beside them — CSV2b templates the version, which ably-python sends as +5, and the binary protocol default (TO3f) makes the content type msgpack, as +[#527](https://github.com/ably/specification/issues/527) records for the unit specs that +read a body without pinning the protocol. The envelope disagreement is filed as +[#559](https://github.com/ably/specification/issues/559), together with the two +`revoke_tokens.md` mocks below. + +`revoke_tokens.md` has the same internal split: `RSA17c_1` and `TRS2_1` stub a bare array +while asserting envelope fields. The revocation half is settled by the server too. `POST /keys/{keyName}/revokeTokens` with `X-Ably-Version: 5` answers 201 with the `{successCount, failureCount, results}` envelope `revoke_tokens.md`'s "Server Response Format" section describes, for a mixed success/failure batch as well as an all-success @@ -773,9 +794,6 @@ Filed as [#554](https://github.com/ably/specification/issues/554). | `message_encoding.md`, `msgpack_interop.md`, `annotations.md` | Six sections carry no Test ID; ids were inferred by sibling convention | | `publish.md`, `rest_presence.md`, `message_encoding.md`, `history.md`, `idempotency.md` | All point at `/Users/paddy/data/worknew/dev/dart-experiments/...` for the mock contract | | `revoke_tokens.md` | The all-success Setup blocks stub HTTP 200 with a plain array, the shape the file's own "Server Response Format" section calls legacy and says no current SDK sees. The `BatchResult` envelope and HTTP 201 that section prescribes appear only in the mixed and all-failure blocks | -| `publish.md` (integration) | The `Spec points:` header reads RSL1d, RSL1l1, RSL1m4, RSL1n, and the file carries a fifth section, `## RSL1k5 - Idempotent publish with client-supplied IDs`, with its own Test ID. The section is sound; only the header is short. Same housekeeping class as [#532](https://github.com/ably/specification/issues/532) | -| `auth.md` (integration) | RSC10's expired-JWT fixture is `generate_jwt(expires_at: now() - 5_seconds)`, naming `exp` and leaving `iat` open. Ably reads a JWT's lifetime as `exp - iat` and rejects a negative one with 400/40003 "Invalid value for ttl" before it considers expiry, so `iat` at now produces a token that fails the wrong way and never reaches the 40140–40149 renewal path the test is about. Backdating `iat` past `exp` gives the already-expired token the test wants, answered 401/40142. An SDK signing its own Ably JWT has to choose, so the fixture should say which | -| `batch_presence.md` | BGR2 says a channel with no members "returns a success result with an empty `presence` array", and the unit tier's mocks all send `'presence': []`. The server sends no `presence` key at all, so an implementation has to default the field for the assertion to hold. The derived test asserts the specification's `length == 0`, with the wire shape in a comment | | `publish.md` (integration) | The `Spec points:` header reads RSL1d, RSL1l1, RSL1m4, RSL1n, and the file carries a fifth section, `## RSL1k5 - Idempotent publish with client-supplied IDs`, with its own Test ID. The section is sound; only the header is short. `auth.md` (integration) has the same shape: its header reads RSA4, RSA8 and it carries `## RSC10` with its own Test ID. Same housekeeping class as [#532](https://github.com/ably/specification/issues/532); filed as [#550](https://github.com/ably/specification/issues/550) | | `auth.md` (integration) | RSC10's expired-JWT fixture is `generate_jwt(expires_at: now() - 5_seconds)`, naming `exp` and leaving `iat` open. Ably reads a JWT's lifetime as `exp - iat` and rejects a negative one with 400/40003 "Invalid value for ttl" before it considers expiry, so `iat` at now produces a token that fails the wrong way and never reaches the 40140–40149 renewal path the test is about. Backdating `iat` past `exp` gives the already-expired token the test wants, answered 401/40142. An SDK signing its own Ably JWT has to choose, so the fixture should say which. Filed as [#550](https://github.com/ably/specification/issues/550) | | `batch_presence.md` | The restricted-key setup's comment reads "only has access to \"batch-allowed\" channel" while the setup fixes `allowed_channel = "channel6"`; `batch-allowed` appears nowhere in the file. Filed with [#548](https://github.com/ably/specification/issues/548), whose fix replaces the same lines | @@ -792,8 +810,8 @@ Nothing to fix here, only something to build. Each row is one feature, and the c the number of gated Test IDs that fall with it, with the pytest case count beside it where the two differ. -Five of these rows are gated at both tiers. `batchPresence`, `Auth#revokeTokens`, the -`PushChannel` surface and the `clientId` filter on `RestPresence#get` each carry +Four of these rows are gated at both tiers. `Auth#revokeTokens`, the `PushChannel` surface +and the `clientId` filter on `RestPresence#get` each carry `uts/rest/integration` tests as well as unit ones, and connection recovery carries two `uts/realtime/integration` ones; all are written against the spelling the unit tier already gates on, so both tiers go green together when the API lands. Those @@ -801,17 +819,15 @@ integration tests do all their real work first — the sandbox app, the channels presence members entered over a realtime connection, the registered device and the issued token are all real, and each test reaches the missing call before it fails, so the assertions either side of it are known to hold against real server responses. The -`batch_presence` and `push_channels` files were additionally run against throwaway shims -— a `batch_presence` forwarding to `GET /presence`, and a `PushChannel` posting and -deleting `/push/channelSubscriptions` with `X-Ably-DeviceToken` — and pass in full -against them. The two recovery tests are a different shape: there is no missing call for +`push_channels` file was additionally run against a throwaway shim — a `PushChannel` +posting and deleting `/push/channelSubscriptions` with `X-Ably-DeviceToken` — and passes in +full against it. The two recovery tests are a different shape: there is no missing call for them to reach, so each runs end to end against the sandbox through `uts-proxy` and fails on the `recover` parameter the connection never sends. | Spec points | Missing | Test IDs | |---|---|---| -| RSC22, RSC24, BSP2, BPR2, BPF2, BAR2, BGR2, BGF2 | `batchPublish` and `batchPresence`, and all six result types. `grep -rn batch ably/` finds nothing | 44 (47 cases) | -| RSA17, RSA17b–g, BAR2, TRS2, TRF2 | `Auth#revokeTokens`, `TokenRevocationTargetSpecifier`, `BatchResult`. Gated against `auth.revoke_tokens(targets, issued_before=, allow_reauth_margin=)` returning `success_count` / `failure_count` / `results`, with `target` / `issued_before` / `applies_at` / `error` per result. RSA17d is the one case that needs no server at all — a token-authenticated client must refuse locally with 40162/401 — so it can be satisfied before any of the wire work | 21 | +| RSA17, RSA17b–g, BAR2, TRS2, TRF2 | `Auth#revokeTokens`, `TokenRevocationTargetSpecifier` and the two token revocation result types; the `BatchResult` they would arrive in is the batch API's. Gated against `auth.revoke_tokens(targets, issued_before=, allow_reauth_margin=)` returning `success_count` / `failure_count` / `results`, with `target` / `issued_before` / `applies_at` / `error` per result. RSA17d is the one case that needs no server at all — a token-authenticated client must refuse locally with 40162/401 — so it can be satisfied before any of the wire work | 21 | | RSH7, RSH7a–e, RSH6, RSH8 | `PushChannel`: `channel.push`, `client.device`, `LocalDevice`. The push *admin* surface (RSH1) does exist | 12 | | RTN16, RTN16f–k, RTC1c (TO3i) | Connection recovery, entire. `recover` is in the `Options` signature, stored, and given a property and a setter (`options.py:30,111,193,196`), and read nowhere. No `Connection#createRecoveryKey`, no `recover` connect parameter, no recovery-key decoding. Measured through the proxy: a client built with a valid `recover=` opened a `ws_connect` whose query parameters were `{'accessToken': …, 'v': '5'}` — no `recover` — and was given a fresh `connectionId`. RTN16l is otherwise fully compliant, taking the proxy's `recovery-failed-new-id`, `recovery-failed-new-key` and error 80008 and staying CONNECTED; only the absent parameter fails it | 8 | | RTL22, RTL22a–d, MFI1, MFI2a–e | `MessageFilter`. `RealtimeChannel.subscribe` (`channel.py:262-273`) accepts only a `str` or a callable, and there is no filter type of any shape to spell. Each test builds its filter through the module's `message_filter()` helper, which is the one place to repoint when the type lands | 5 | @@ -3217,9 +3233,10 @@ The header states how many derived tests there are, how many pass, how many are gated and how many cannot run. Those numbers are the check that the file is still true: in pytest cases, the gated count must equal the number of failures under `RUN_DEVIATIONS=1`, and gated plus unrunnable must equal the number of skips without -it. As of this writing that is 217 failures and 15 skips with the variable set, and -232 skips and 1124 passes without it, the 1124 being 1002 derived cases and 122 -`helpers/` ones. +it. As of this writing that is 172 failures and 15 skips with the variable set, and +187 skips and 1177 passes without it, the 1177 being 1047 derived cases and 130 +`helpers/` ones. Those figures need the `submodules` checkout CI makes: without it, three +`rest/unit/encoding` tests that read the ably-common fixtures skip as well. The other two counts are measured from the source rather than from a run. The number of **derived tests** is the number of `# UTS:` comments, 1141. The number of **Test IDs** is diff --git a/test/uts/rest/integration/batch_presence_test.py b/test/uts/rest/integration/batch_presence_test.py index b9d4e7af..ecf2d93c 100644 --- a/test/uts/rest/integration/batch_presence_test.py +++ b/test/uts/rest/integration/batch_presence_test.py @@ -2,38 +2,26 @@ Spec points: RSC24, BGR2, BGF2 -DEVIATION: ably-python has no batch API. `DefaultPubSubHttpClient` exposes no `batch_presence`, and -the package defines neither `BatchResult` nor `BatchPresenceSuccessResult` / -`BatchPresenceFailureResult`; the word "batch" appears nowhere under `ably/`. Every -test here therefore departs from the specification and is gated behind -RUN_DEVIATIONS, against the same spelling -[test/uts/rest/unit/batch_presence_test.py](../unit/batch_presence_test.py) gates on — -`client.batch_presence([...])` giving a result with `success_count`, `failure_count` -and `results` — so that dropping the marker is the only change either tier needs when -RSC24 lands. - -The setup halves are real, and the responses they assert against were confirmed by -hand through `client.request('GET', '/presence', params={'channels': ...})` against -the sandbox: the server does return the `successCount` / `failureCount` / `results` -envelope the specification describes, with `code` 40160 and `statusCode` 401 for a -channel the key has no capability for. Two details of that confirmation are recorded -beside the assertions they bear on — the `presence` key the server omits for an empty -channel, and the presence members a closed connection takes with it. - -See [deviations.md](../../deviations.md): the gating under *Failing Tests* -> -*Unimplemented features*, the omitted `presence` key and the closed connection under -*UTS Spec Errors*, and `enterClient` on an anonymous connection under *Failing Tests* -> -*Auth*. +The server returns the `successCount` / `failureCount` / `results` envelope the +specification describes, with `code` 40160 and `statusCode` 401 for a channel the key +has no capability for. Two details of the server's responses are recorded beside the +assertions they bear on — the `presence` key the server omits for an empty channel, and +the presence members a closed connection takes with it. + +See [deviations.md](../../deviations.md): the omitted `presence` key and the closed +connection under *UTS Spec Errors*, and `enterClient` on an anonymous connection under +*Failing Tests* -> *Auth*. """ from ably.pubsub.realtime.connection import ConnectionState +from ably.pubsub.types.batch import BatchPresenceFailureResult, BatchPresenceSuccessResult +from ably.pubsub.util.exceptions import AblyException from test.uts.helpers.client import ( await_connection_state, sandbox_realtime_client, sandbox_rest_client, wall_clock_poll_until, ) -from test.uts.helpers.deviations import deviation from test.uts.helpers.sandbox import random_id @@ -45,20 +33,6 @@ def result_for(result, channel_name): raise AssertionError(f'No batch result for channel {channel_name!r}') -def is_success_result(entry): - """The specifications' `entry IS BatchPresenceSuccessResult`. - - Neither result class exists to name, so a success is told apart from a failure by - which attribute it carries, as the unit tier does. - """ - return getattr(entry, 'presence', None) is not None - - -def is_failure_result(entry): - """The specifications' `entry IS BatchPresenceFailureResult`.""" - return getattr(entry, 'error', None) is not None - - def member_for(entry, client_id): """The specifications' `presence.find(m => m.clientId == client_id)`.""" for member in entry.presence: @@ -96,7 +70,6 @@ async def enter_members(realtime, channel_name, members): return channel -@deviation # UTS: rest/integration/RSC24/batch-presence-multiple-channels-0 async def test_rsc24_batch_presence_multiple_channels(sandbox, use_binary_protocol): channel_a_name = 'batch-presence-a-' + random_id() @@ -119,7 +92,7 @@ async def test_rsc24_batch_presence_multiple_channels(sandbox, use_binary_protoc result_a = result_for(result, channel_a_name) result_b = result_for(result, channel_b_name) - assert is_success_result(result_a) + assert isinstance(result_a, BatchPresenceSuccessResult) assert len(result_a.presence) == 2 client_ids_a = [member.client_id for member in result_a.presence] assert 'user-1' in client_ids_a @@ -127,13 +100,12 @@ async def test_rsc24_batch_presence_multiple_channels(sandbox, use_binary_protoc assert member_for(result_a, 'user-1').data == 'data-a1' - assert is_success_result(result_b) + assert isinstance(result_b, BatchPresenceSuccessResult) assert len(result_b.presence) == 1 assert result_b.presence[0].client_id == 'user-3' assert result_b.presence[0].data == 'data-b1' -@deviation # UTS: rest/integration/RSC24/restricted-key-channel-failure-1 async def test_rsc24_restricted_key_channel_failure(sandbox, use_binary_protocol): # The specification hard-codes "channel6" as the channel `keys[2]` is allowed. The @@ -165,7 +137,7 @@ async def test_rsc24_restricted_key_channel_failure(sandbox, use_binary_protocol async def one_member_on_the_allowed_channel(): result = await restricted_rest.batch_presence([allowed_channel, denied_channel]) success = result_for(result, allowed_channel) - return result if is_success_result(success) and len(success.presence) == 1 else None + return result if isinstance(success, BatchPresenceSuccessResult) and len(success.presence) == 1 else None result = await wall_clock_poll_until( one_member_on_the_allowed_channel, @@ -178,16 +150,16 @@ async def one_member_on_the_allowed_channel(): success = result_for(result, allowed_channel) failure = result_for(result, denied_channel) - assert is_success_result(success) + assert isinstance(success, BatchPresenceSuccessResult) assert len(success.presence) == 1 assert success.presence[0].client_id == 'member-1' - assert is_failure_result(failure) + assert isinstance(failure, BatchPresenceFailureResult) + assert isinstance(failure.error, AblyException) assert failure.error.code == 40160 assert failure.error.status_code == 401 -@deviation # UTS: rest/integration/RSC24/empty-channel-presence-2 async def test_rsc24_empty_channel_presence(sandbox, use_binary_protocol): empty_channel = 'batch-empty-' + random_id() @@ -210,9 +182,9 @@ async def test_rsc24_empty_channel_presence(sandbox, use_binary_protocol): # The server counts the empty channel a success and leaves `presence` out of its # result altogether rather than sending `[]`, so an implementation of BGR2 has to # default the field for the specification's assertion to hold. - assert is_success_result(empty_result) + assert isinstance(empty_result, BatchPresenceSuccessResult) assert len(empty_result.presence) == 0 - assert is_success_result(populated_result) + assert isinstance(populated_result, BatchPresenceSuccessResult) assert len(populated_result.presence) == 1 assert populated_result.presence[0].client_id == 'someone' diff --git a/test/uts/rest/unit/batch_presence_test.py b/test/uts/rest/unit/batch_presence_test.py index 599a3892..e1c67fc0 100644 --- a/test/uts/rest/unit/batch_presence_test.py +++ b/test/uts/rest/unit/batch_presence_test.py @@ -2,28 +2,16 @@ Spec points: RSC24, BAR2, BGR2, BGF2 -NOTE: ably-python has no batch API. `DefaultPubSubHttpClient` exposes no `batch_presence`, and the package -defines neither `BatchResult`/`BatchPresenceResponse` nor `BatchPresenceSuccessResult`/ -`BatchPresenceFailureResult`; the word "batch" appears nowhere under `ably/`. Every test in -this file therefore departs from the specification and is gated behind `RUN_DEVIATIONS`. -Each carries the assertion the spec calls for, written against the name ably-python would -use once RSC24 is implemented, so that dropping the `@deviation` marker is the only change -needed when it is. Today a batch presence query has to be hand-rolled by the caller through -`client.request('GET', '/presence', version=..., params={'channels': ...})`, which returns -an `HttpPaginatedResponse` of raw dicts rather than decoded `PresenceMessage`s, and which -does not raise on an error status. - -Counts are read as `success_count` / `failure_count`, the snake_case spelling of BAR2a and -BAR2b, and a success result is told apart from a failure result by which attributes it -carries rather than by `isinstance`, since neither class exists to name. +The specification's `BatchPresenceResponse` is ably-python's `BatchResult`, whose counts are +read as `success_count` / `failure_count`, the snake_case spelling of BAR2a and BAR2b. """ import pytest +from ably.pubsub.types.batch import BatchPresenceFailureResult, BatchPresenceSuccessResult from ably.pubsub.types.presence import PresenceAction from ably.pubsub.util.exceptions import AblyException from test.uts.helpers.client import rest_client -from test.uts.helpers.deviations import deviation from test.uts.helpers.mock_http import MockHttpClient @@ -40,7 +28,6 @@ def respond_with(status, body): # UTS: rest/unit/RSC24/get-presence-channels-param-0 -@deviation async def test_rsc24_batch_presence_get_presence_channels_param(): captured_requests = [] mock_http = MockHttpClient( @@ -65,7 +52,6 @@ async def test_rsc24_batch_presence_get_presence_channels_param(): # UTS: rest/unit/RSC24/single-channel-param-0 -@deviation async def test_rsc24_batch_presence_single_channel_param(): captured_requests = [] mock_http = MockHttpClient( @@ -86,7 +72,6 @@ async def test_rsc24_batch_presence_single_channel_param(): # UTS: rest/unit/RSC24/special-chars-comma-joined-0 -@deviation async def test_rsc24_batch_presence_special_chars_comma_joined(): captured_requests = [] mock_http = MockHttpClient( @@ -108,7 +93,6 @@ async def test_rsc24_batch_presence_special_chars_comma_joined(): # UTS: rest/unit/BAR2/mixed-success-failure-counts-0 -@deviation async def test_bar2_batch_presence_mixed_success_failure_counts(): mock_http = MockHttpClient( on_connection_attempt=lambda conn: conn.respond_with_success(), @@ -136,7 +120,6 @@ async def test_bar2_batch_presence_mixed_success_failure_counts(): # UTS: rest/unit/BAR2/all-success-counts-0 -@deviation async def test_bar2_batch_presence_all_success_counts(): mock_http = MockHttpClient( on_connection_attempt=lambda conn: conn.respond_with_success(), @@ -159,7 +142,6 @@ async def test_bar2_batch_presence_all_success_counts(): # UTS: rest/unit/BAR2/all-failure-counts-0 -@deviation async def test_bar2_batch_presence_all_failure_counts(): error = {'code': 40160, 'statusCode': 401, 'message': 'Not permitted'} mock_http = MockHttpClient( @@ -183,7 +165,6 @@ async def test_bar2_batch_presence_all_failure_counts(): # UTS: rest/unit/BGR2/success-with-members-0 -@deviation async def test_bgr2_batch_presence_success_with_members(): mock_http = MockHttpClient( on_connection_attempt=lambda conn: conn.respond_with_success(), @@ -222,7 +203,7 @@ async def test_bgr2_batch_presence_success_with_members(): assert len(result.results) == 1 success = result.results[0] - assert getattr(success, 'presence', None) is not None + assert isinstance(success, BatchPresenceSuccessResult) assert success.channel == 'my-channel' assert len(success.presence) == 2 @@ -237,7 +218,6 @@ async def test_bgr2_batch_presence_success_with_members(): # UTS: rest/unit/BGR2/success-empty-presence-0 -@deviation async def test_bgr2_batch_presence_success_empty_presence(): mock_http = MockHttpClient( on_connection_attempt=lambda conn: conn.respond_with_success(), @@ -254,13 +234,12 @@ async def test_bgr2_batch_presence_success_empty_presence(): result = await client.batch_presence(['empty-channel']) success = result.results[0] - assert getattr(success, 'presence', None) is not None + assert isinstance(success, BatchPresenceSuccessResult) assert success.channel == 'empty-channel' assert len(success.presence) == 0 # UTS: rest/unit/BGF2/failure-error-details-0 -@deviation async def test_bgf2_batch_presence_failure_error_details(): mock_http = MockHttpClient( on_connection_attempt=lambda conn: conn.respond_with_success(), @@ -286,7 +265,8 @@ async def test_bgf2_batch_presence_failure_error_details(): assert len(result.results) == 1 failure = result.results[0] - assert getattr(failure, 'error', None) is not None + assert isinstance(failure, BatchPresenceFailureResult) + assert isinstance(failure.error, AblyException) assert failure.channel == 'restricted-channel' assert failure.error.code == 40160 assert failure.error.status_code == 401 @@ -294,7 +274,6 @@ async def test_bgf2_batch_presence_failure_error_details(): # UTS: rest/unit/RSC24/mixed-success-failure-results-0 -@deviation async def test_rsc24_batch_presence_mixed_success_failure_results(): mock_http = MockHttpClient( on_connection_attempt=lambda conn: conn.respond_with_success(), @@ -329,18 +308,17 @@ async def test_rsc24_batch_presence_mixed_success_failure_results(): assert result.failure_count == 1 assert len(result.results) == 2 - assert getattr(result.results[0], 'presence', None) is not None + assert isinstance(result.results[0], BatchPresenceSuccessResult) assert result.results[0].channel == 'allowed-channel' assert len(result.results[0].presence) == 1 assert result.results[0].presence[0].client_id == 'user-1' - assert getattr(result.results[1], 'error', None) is not None + assert isinstance(result.results[1], BatchPresenceFailureResult) assert result.results[1].channel == 'restricted-channel' assert result.results[1].error.code == 40160 # UTS: rest/unit/RSC24/server-error-propagated-0 -@deviation async def test_rsc24_batch_presence_server_error_propagated(): mock_http = MockHttpClient( on_connection_attempt=lambda conn: conn.respond_with_success(), @@ -358,7 +336,6 @@ async def test_rsc24_batch_presence_server_error_propagated(): # UTS: rest/unit/RSC24/auth-error-propagated-0 -@deviation async def test_rsc24_batch_presence_auth_error_propagated(): mock_http = MockHttpClient( on_connection_attempt=lambda conn: conn.respond_with_success(), @@ -376,7 +353,6 @@ async def test_rsc24_batch_presence_auth_error_propagated(): # UTS: rest/unit/RSC24/uses-configured-auth-0 -@deviation async def test_rsc24_batch_presence_uses_configured_auth(): captured_requests = [] mock_http = MockHttpClient( diff --git a/test/uts/rest/unit/batch_publish_test.py b/test/uts/rest/unit/batch_publish_test.py index a4765c2f..bf40f99d 100644 --- a/test/uts/rest/unit/batch_publish_test.py +++ b/test/uts/rest/unit/batch_publish_test.py @@ -2,31 +2,30 @@ Spec points: RSC22c, RSC22d, BSP2a, BSP2b, BPR2a, BPR2b, BPR2c, BPF2a, BPF2b -NOTE: ably-python has no batch API. `DefaultPubSubHttpClient` exposes no `batch_publish`, and the package -defines neither `BatchPublishSpec` nor `BatchResult`/`BatchPublishSuccessResult`/ -`BatchPublishFailureResult`; the word "batch" appears nowhere under `ably/`. Every test in -this file therefore departs from the specification and is gated behind `RUN_DEVIATIONS`. -Each carries the assertion the spec calls for, written against the name ably-python would -use once RSC22 is implemented, so that dropping the `@deviation` marker is the only change -needed when it is. Today a batch publish has to be hand-rolled by the caller through -`client.request('POST', '/messages', version=..., body=...)`, which does no spec -construction, no RSL4 message encoding and no RSL1k1 idempotent ID generation. - -Because `BatchPublishSpec` does not exist to construct, a spec is written as a mapping of -the BSP2 attributes. Results are read by attribute, as the spec writes them, and a success -result is told apart from a failure result by which attributes it carries rather than by -`isinstance`, since neither class exists to name. +The specification's mocks answer in the legacy response format, a flat array of per-channel +results, which the server sends to a client sending `X-Ably-Version: 2` or no version at all. +ably-python sends version 5, and to that the server answers every batch publish with HTTP 201 +and an array holding a `BatchResult` envelope for each spec sent, a single spec sent as a bare +object included. Every fixture is corrected to that shape through `batch_result()`, and the +assertions stand as the specification writes them. See *Batch response envelopes disagree +between sibling specs* in [deviations.md](../../deviations.md). """ +import json import uuid import msgpack import pytest +from ably.pubsub.types.batch import ( + BatchPublishFailureResult, + BatchPublishSpec, + BatchPublishSuccessResult, + BatchResult, +) from ably.pubsub.types.message import Message from ably.pubsub.util.exceptions import AblyException from test.uts.helpers.client import rest_client -from test.uts.helpers.deviations import deviation from test.uts.helpers.mock_http import MockHttpClient @@ -37,7 +36,7 @@ def random_id(): def capture_and_respond(captured_requests, status=201, body=None): def on_request(request): captured_requests.append(request) - request.respond_with(status, body if body is not None else []) + request.respond_with(status, body if body is not None else success_response(request)) return on_request @@ -50,8 +49,29 @@ def success_result(channel, message_id='msg', serials=('s1',)): return {'channel': channel, 'messageId': message_id, 'serials': list(serials)} +def batch_result(*results): + """One spec's `BatchResult` envelope, holding the per-channel `results` given. + + UTS SPEC ERROR: batch_publish.md - the specification's mocks send the per-channel results + as a flat array, the legacy format, where the server answers ably-python's + `X-Ably-Version: 5` with an array of these envelopes, one per spec. + """ + failure_count = sum(1 for result in results if 'error' in result) + return { + 'successCount': len(results) - failure_count, + 'failureCount': failure_count, + 'results': list(results), + } + + +def success_response(request): + """What the server answers a batch publish with when every channel succeeds.""" + body = msgpack.unpackb(request.body) + specs = body if isinstance(body, list) else [body] + return [batch_result(*(success_result(channel) for channel in spec['channels'])) for spec in specs] + + # UTS: rest/unit/RSC22c/single-spec-post-messages-0 -@deviation async def test_rsc22c_batch_publish_single_spec_post_messages(): channel_name_1 = f'test-RSC22c1-a-{random_id()}' channel_name_2 = f'test-RSC22c1-b-{random_id()}' @@ -63,10 +83,10 @@ async def test_rsc22c_batch_publish_single_spec_post_messages(): ) client = rest_client(mock_http) - await client.batch_publish({ - 'channels': [channel_name_1, channel_name_2], - 'messages': [Message(name='event', data='hello')], - }) + await client.batch_publish(BatchPublishSpec( + channels=[channel_name_1, channel_name_2], + messages=[Message(name='event', data='hello')], + )) assert len(captured_requests) == 1 request = captured_requests[0] @@ -81,7 +101,6 @@ async def test_rsc22c_batch_publish_single_spec_post_messages(): # UTS: rest/unit/RSC22c/array-specs-post-messages-0 -@deviation async def test_rsc22c_batch_publish_array_specs_post_messages(): channel_name_1 = f'test-RSC22c2-a-{random_id()}' channel_name_2 = f'test-RSC22c2-b-{random_id()}' @@ -94,8 +113,8 @@ async def test_rsc22c_batch_publish_array_specs_post_messages(): client = rest_client(mock_http) await client.batch_publish([ - {'channels': [channel_name_1], 'messages': [Message(name='e1', data='d1')]}, - {'channels': [channel_name_2], 'messages': [Message(name='e2', data='d2')]}, + BatchPublishSpec(channels=[channel_name_1], messages=[Message(name='e1', data='d1')]), + BatchPublishSpec(channels=[channel_name_2], messages=[Message(name='e2', data='d2')]), ]) assert len(captured_requests) == 1 @@ -113,13 +132,13 @@ async def test_rsc22c_batch_publish_array_specs_post_messages(): # UTS: rest/unit/RSC22c/single-spec-single-result-0 -@deviation async def test_rsc22c_batch_publish_single_spec_single_result(): channel_name = f'test-RSC22c3-{random_id()}' # UTS SPEC ERROR: RSC22c3 - the mock body is a bare result object, but RSC22b says the REST # response "will still be an array", from which the single-spec overload extracts one element. - response_body = [{'channel': channel_name, 'messageId': 'msg123', 'serials': ['serial1']}] + # The server sends that element as the spec's BatchResult envelope. + response_body = [batch_result({'channel': channel_name, 'messageId': 'msg123', 'serials': ['serial1']})] mock_http = MockHttpClient( on_connection_attempt=lambda conn: conn.respond_with_success(), @@ -127,26 +146,25 @@ async def test_rsc22c_batch_publish_single_spec_single_result(): ) client = rest_client(mock_http) - result = await client.batch_publish({ - 'channels': [channel_name], - 'messages': [Message(name='event', data='hello')], - }) + result = await client.batch_publish(BatchPublishSpec( + channels=[channel_name], + messages=[Message(name='event', data='hello')], + )) - assert not isinstance(result, list) + assert isinstance(result, BatchResult) assert len(result.results) == 1 assert result.results[0].channel == channel_name assert result.results[0].message_id == 'msg123' # UTS: rest/unit/RSC22c/array-specs-array-results-0 -@deviation async def test_rsc22c_batch_publish_array_specs_array_results(): channel_name_1 = f'test-RSC22c4-a-{random_id()}' channel_name_2 = f'test-RSC22c4-b-{random_id()}' response_body = [ - success_result(channel_name_1, 'msg1', ['s1']), - success_result(channel_name_2, 'msg2', ['s2']), + batch_result(success_result(channel_name_1, 'msg1', ['s1'])), + batch_result(success_result(channel_name_2, 'msg2', ['s2'])), ] mock_http = MockHttpClient( @@ -156,28 +174,28 @@ async def test_rsc22c_batch_publish_array_specs_array_results(): client = rest_client(mock_http) results = await client.batch_publish([ - {'channels': [channel_name_1], 'messages': [Message(name='e1', data='d1')]}, - {'channels': [channel_name_2], 'messages': [Message(name='e2', data='d2')]}, + BatchPublishSpec(channels=[channel_name_1], messages=[Message(name='e1', data='d1')]), + BatchPublishSpec(channels=[channel_name_2], messages=[Message(name='e2', data='d2')]), ]) assert isinstance(results, list) assert len(results) == 2 + assert all(isinstance(result, BatchResult) for result in results) assert results[0].results[0].channel == channel_name_1 assert results[1].results[0].channel == channel_name_2 # UTS: rest/unit/RSC22c/multiple-channels-multiple-results-0 -@deviation async def test_rsc22c_batch_publish_multiple_channels_multiple_results(): channel_name_1 = f'test-RSC22c5-a-{random_id()}' channel_name_2 = f'test-RSC22c5-b-{random_id()}' channel_name_3 = f'test-RSC22c5-c-{random_id()}' - response_body = [ + response_body = [batch_result( success_result(channel_name_1, 'msg1', ['s1']), success_result(channel_name_2, 'msg2', ['s2']), success_result(channel_name_3, 'msg3', ['s3']), - ] + )] mock_http = MockHttpClient( on_connection_attempt=lambda conn: conn.respond_with_success(), @@ -185,17 +203,16 @@ async def test_rsc22c_batch_publish_multiple_channels_multiple_results(): ) client = rest_client(mock_http) - result = await client.batch_publish({ - 'channels': [channel_name_1, channel_name_2, channel_name_3], - 'messages': [Message(name='event', data='hello')], - }) + result = await client.batch_publish(BatchPublishSpec( + channels=[channel_name_1, channel_name_2, channel_name_3], + messages=[Message(name='event', data='hello')], + )) assert len(result.results) == 3 assert [entry.channel for entry in result.results] == [channel_name_1, channel_name_2, channel_name_3] # UTS: rest/unit/RSC22c/messages-encoded-per-rsl4-0 -@deviation async def test_rsc22c_batch_publish_messages_encoded_per_rsl4(): channel_name = f'test-RSC22c6-{random_id()}' @@ -206,14 +223,14 @@ async def test_rsc22c_batch_publish_messages_encoded_per_rsl4(): ) client = rest_client(mock_http) - await client.batch_publish({ - 'channels': [channel_name], - 'messages': [ + await client.batch_publish(BatchPublishSpec( + channels=[channel_name], + messages=[ Message(name='string', data='plain text'), Message(name='binary', data=b'\x01\x02\x03'), Message(name='json', data={'key': 'value'}), ], - }) + )) assert len(captured_requests) == 1 messages = msgpack.unpackb(captured_requests[0].body)['messages'] @@ -228,12 +245,11 @@ async def test_rsc22c_batch_publish_messages_encoded_per_rsl4(): assert 'encoding' not in messages[1] or messages[1]['encoding'] is None # RSL4c3: a JSON payload is stringified and the encoding attribute is set to "json" - assert messages[2]['data'] == '{"key":"value"}' + assert json.loads(messages[2]['data']) == {'key': 'value'} assert messages[2]['encoding'] == 'json' # UTS: rest/unit/RSC22c/uses-configured-auth-0 -@deviation async def test_rsc22c_batch_publish_uses_configured_auth(): channel_name = f'test-RSC22c7-{random_id()}' @@ -244,10 +260,10 @@ async def test_rsc22c_batch_publish_uses_configured_auth(): ) token_client = rest_client(mock_http, token='fake-token') - await token_client.batch_publish({ - 'channels': [channel_name], - 'messages': [Message(name='event', data='hello')], - }) + await token_client.batch_publish(BatchPublishSpec( + channels=[channel_name], + messages=[Message(name='event', data='hello')], + )) assert len(captured_requests) == 1 assert captured_requests[0].headers['Authorization'].startswith('Bearer ') @@ -258,17 +274,16 @@ async def test_rsc22c_batch_publish_uses_configured_auth(): mock_http.on_request = capture_and_respond(basic_requests) basic_client = rest_client(mock_http) - await basic_client.batch_publish({ - 'channels': [basic_channel_name], - 'messages': [Message(name='event', data='hello')], - }) + await basic_client.batch_publish(BatchPublishSpec( + channels=[basic_channel_name], + messages=[Message(name='event', data='hello')], + )) assert len(basic_requests) == 1 assert basic_requests[0].headers['Authorization'].startswith('Basic ') # UTS: rest/unit/RSC22d/idempotent-ids-generated-0 -@deviation async def test_rsc22d_batch_publish_idempotent_ids_generated(): captured_requests = [] mock_http = MockHttpClient( @@ -278,14 +293,14 @@ async def test_rsc22d_batch_publish_idempotent_ids_generated(): client = rest_client(mock_http, idempotent_rest_publishing=True) await client.batch_publish([ - { - 'channels': [f'test-RSC22d-a-{random_id()}'], - 'messages': [Message(name='e1', data='d1'), Message(name='e2', data='d2')], - }, - { - 'channels': [f'test-RSC22d-b-{random_id()}'], - 'messages': [Message(name='e3', data='d3'), Message(name='e4', data='d4')], - }, + BatchPublishSpec( + channels=[f'test-RSC22d-a-{random_id()}'], + messages=[Message(name='e1', data='d1'), Message(name='e2', data='d2')], + ), + BatchPublishSpec( + channels=[f'test-RSC22d-b-{random_id()}'], + messages=[Message(name='e3', data='d3'), Message(name='e4', data='d4')], + ), ]) assert len(captured_requests) == 1 @@ -307,7 +322,6 @@ async def test_rsc22d_batch_publish_idempotent_ids_generated(): # UTS: rest/unit/RSC22d/explicit-ids-preserved-0 -@deviation async def test_rsc22d_batch_publish_explicit_ids_preserved(): channel_name = f'test-RSC22d-explicit-{random_id()}' @@ -318,13 +332,13 @@ async def test_rsc22d_batch_publish_explicit_ids_preserved(): ) client = rest_client(mock_http, idempotent_rest_publishing=True) - await client.batch_publish({ - 'channels': [channel_name], - 'messages': [ + await client.batch_publish(BatchPublishSpec( + channels=[channel_name], + messages=[ Message(name='e1', data='d1', id='explicit-id-1'), Message(name='e2', data='d2', id='explicit-id-2'), ], - }) + )) assert len(captured_requests) == 1 messages = msgpack.unpackb(captured_requests[0].body)['messages'] @@ -334,7 +348,6 @@ async def test_rsc22d_batch_publish_explicit_ids_preserved(): # UTS: rest/unit/RSC22d/ids-not-generated-disabled-0 -@deviation async def test_rsc22d_batch_publish_ids_not_generated_disabled(): channel_name = f'test-RSC22d-disabled-{random_id()}' @@ -345,10 +358,10 @@ async def test_rsc22d_batch_publish_ids_not_generated_disabled(): ) client = rest_client(mock_http, idempotent_rest_publishing=False) - await client.batch_publish({ - 'channels': [channel_name], - 'messages': [Message(name='e1', data='d1'), Message(name='e2', data='d2')], - }) + await client.batch_publish(BatchPublishSpec( + channels=[channel_name], + messages=[Message(name='e1', data='d1'), Message(name='e2', data='d2')], + )) assert len(captured_requests) == 1 messages = msgpack.unpackb(captured_requests[0].body)['messages'] @@ -358,7 +371,6 @@ async def test_rsc22d_batch_publish_ids_not_generated_disabled(): # UTS: rest/unit/BSP2a/channels-array-strings-0 -@deviation async def test_bsp2a_batch_publish_spec_channels_array_strings(): channel_name_1 = f'test-BSP2a-a-{random_id()}' channel_name_2 = f'test-BSP2a-b-{random_id()}' @@ -371,10 +383,10 @@ async def test_bsp2a_batch_publish_spec_channels_array_strings(): ) client = rest_client(mock_http) - await client.batch_publish({ - 'channels': [channel_name_1, channel_name_2, channel_name_3], - 'messages': [Message(name='event', data='hello')], - }) + await client.batch_publish(BatchPublishSpec( + channels=[channel_name_1, channel_name_2, channel_name_3], + messages=[Message(name='event', data='hello')], + )) assert len(captured_requests) == 1 channels = msgpack.unpackb(captured_requests[0].body)['channels'] @@ -385,7 +397,6 @@ async def test_bsp2a_batch_publish_spec_channels_array_strings(): # UTS: rest/unit/BSP2b/messages-array-objects-0 -@deviation async def test_bsp2b_batch_publish_spec_messages_array_objects(): channel_name = f'test-BSP2b-{random_id()}' @@ -396,13 +407,13 @@ async def test_bsp2b_batch_publish_spec_messages_array_objects(): ) client = rest_client(mock_http) - await client.batch_publish({ - 'channels': [channel_name], - 'messages': [ + await client.batch_publish(BatchPublishSpec( + channels=[channel_name], + messages=[ Message(name='event1', data='data1'), Message(name='event2', data={'key': 'value'}), ], - }) + )) assert len(captured_requests) == 1 messages = msgpack.unpackb(captured_requests[0].body)['messages'] @@ -415,118 +426,117 @@ async def test_bsp2b_batch_publish_spec_messages_array_objects(): assert messages[0]['name'] == 'event1' assert messages[0]['data'] == 'data1' assert messages[1]['name'] == 'event2' - assert messages[1]['data'] == '{"key":"value"}' + assert isinstance(messages[1]['data'], str) + assert json.loads(messages[1]['data']) == {'key': 'value'} assert messages[1]['encoding'] == 'json' # UTS: rest/unit/BPR2a/success-channel-name-0 -@deviation async def test_bpr2a_batch_publish_success_channel_name(): channel_name = f'test-BPR2a-{random_id()}' mock_http = MockHttpClient( on_connection_attempt=lambda conn: conn.respond_with_success(), - on_request=respond_with(201, [success_result(channel_name, 'msg123', ['s1'])]), + on_request=respond_with(201, [batch_result(success_result(channel_name, 'msg123', ['s1']))]), ) client = rest_client(mock_http) - result = await client.batch_publish({ - 'channels': [channel_name], - 'messages': [Message(name='event', data='hello')], - }) + result = await client.batch_publish(BatchPublishSpec( + channels=[channel_name], + messages=[Message(name='event', data='hello')], + )) assert result.results[0].channel == channel_name # UTS: rest/unit/BPR2b/success-message-id-prefix-0 -@deviation async def test_bpr2b_batch_publish_success_message_id_prefix(): channel_name = f'test-BPR2b-{random_id()}' mock_http = MockHttpClient( on_connection_attempt=lambda conn: conn.respond_with_success(), - on_request=respond_with(201, [success_result(channel_name, 'unique-id-prefix', ['s1', 's2'])]), + on_request=respond_with(201, [batch_result( + success_result(channel_name, 'unique-id-prefix', ['s1', 's2']), + )]), ) client = rest_client(mock_http) - result = await client.batch_publish({ - 'channels': [channel_name], - 'messages': [Message(name='e1', data='d1'), Message(name='e2', data='d2')], - }) + result = await client.batch_publish(BatchPublishSpec( + channels=[channel_name], + messages=[Message(name='e1', data='d1'), Message(name='e2', data='d2')], + )) assert result.results[0].message_id == 'unique-id-prefix' # UTS: rest/unit/BPR2c/serials-array-0 -@deviation async def test_bpr2c_batch_publish_serials_array(): channel_name = f'test-BPR2c-{random_id()}' mock_http = MockHttpClient( on_connection_attempt=lambda conn: conn.respond_with_success(), - on_request=respond_with(201, [success_result(channel_name, 'msg', ['serial1', 'serial2', 'serial3'])]), + on_request=respond_with(201, [batch_result( + success_result(channel_name, 'msg', ['serial1', 'serial2', 'serial3']), + )]), ) client = rest_client(mock_http) messages = [Message(name=f'e{index}', data=f'd{index}') for index in range(3)] - result = await client.batch_publish({'channels': [channel_name], 'messages': messages}) + result = await client.batch_publish(BatchPublishSpec(channels=[channel_name], messages=messages)) assert result.results[0].serials == ['serial1', 'serial2', 'serial3'] assert len(result.results[0].serials) == len(messages) # UTS: rest/unit/BPR2c/serials-null-conflated-0 -@deviation async def test_bpr2c_batch_publish_serials_null_conflated(): channel_name = f'test-BPR2c1-{random_id()}' mock_http = MockHttpClient( on_connection_attempt=lambda conn: conn.respond_with_success(), - on_request=respond_with(201, [{ + on_request=respond_with(201, [batch_result({ 'channel': channel_name, 'messageId': 'msg', 'serials': ['serial1', None, 'serial3'], - }]), + })]), ) client = rest_client(mock_http) messages = [Message(name=f'e{index}', data=f'd{index}') for index in range(3)] - result = await client.batch_publish({'channels': [channel_name], 'messages': messages}) + result = await client.batch_publish(BatchPublishSpec(channels=[channel_name], messages=messages)) # BPR2c: a null serial marks a message discarded by a conflation rule assert result.results[0].serials == ['serial1', None, 'serial3'] # UTS: rest/unit/BPF2a/failure-channel-name-0 -@deviation async def test_bpf2a_batch_publish_failure_channel_name(): channel_name = f'test-BPF2a-{random_id()}' mock_http = MockHttpClient( on_connection_attempt=lambda conn: conn.respond_with_success(), - on_request=respond_with(201, [{ + on_request=respond_with(201, [batch_result({ 'channel': channel_name, 'error': {'code': 40160, 'statusCode': 401, 'message': 'Not permitted'}, - }]), + })]), ) client = rest_client(mock_http) - result = await client.batch_publish({ - 'channels': [channel_name], - 'messages': [Message(name='event', data='hello')], - }) + result = await client.batch_publish(BatchPublishSpec( + channels=[channel_name], + messages=[Message(name='event', data='hello')], + )) assert result.results[0].channel == channel_name # UTS: rest/unit/BPF2b/failure-error-info-0 -@deviation async def test_bpf2b_batch_publish_failure_error_info(): channel_name = f'test-BPF2b-{random_id()}' mock_http = MockHttpClient( on_connection_attempt=lambda conn: conn.respond_with_success(), - on_request=respond_with(201, [{ + on_request=respond_with(201, [batch_result({ 'channel': channel_name, 'error': { 'code': 40160, @@ -534,78 +544,78 @@ async def test_bpf2b_batch_publish_failure_error_info(): 'message': 'Channel operation not permitted', 'href': 'https://help.ably.io/error/40160', }, - }]), + })]), ) client = rest_client(mock_http) - result = await client.batch_publish({ - 'channels': [channel_name], - 'messages': [Message(name='event', data='hello')], - }) + result = await client.batch_publish(BatchPublishSpec( + channels=[channel_name], + messages=[Message(name='event', data='hello')], + )) error = result.results[0].error + assert isinstance(error, AblyException) assert error.code == 40160 assert error.status_code == 401 assert 'not permitted' in error.message # UTS: rest/unit/RSC22c/partial-success-mixed-results-0 -@deviation async def test_rsc22c_batch_publish_partial_success_mixed_results(): channel_name_allowed = f'test-BatchResult1-allowed-{random_id()}' channel_name_restricted = f'test-BatchResult1-restricted-{random_id()}' mock_http = MockHttpClient( on_connection_attempt=lambda conn: conn.respond_with_success(), - on_request=respond_with(201, [ + on_request=respond_with(201, [batch_result( success_result(channel_name_allowed, 'msg1', ['s1']), { 'channel': channel_name_restricted, 'error': {'code': 40160, 'statusCode': 401, 'message': 'Not permitted'}, }, - ]), + )]), ) client = rest_client(mock_http) - result = await client.batch_publish({ - 'channels': [channel_name_allowed, channel_name_restricted], - 'messages': [Message(name='event', data='hello')], - }) + result = await client.batch_publish(BatchPublishSpec( + channels=[channel_name_allowed, channel_name_restricted], + messages=[Message(name='event', data='hello')], + )) # NOTE: the spec writes result[0] / result[1]; BAR2c puts the per-channel results in # BatchResult.results, so they are read from there. assert len(result.results) == 2 + assert isinstance(result.results[0], BatchPublishSuccessResult) assert result.results[0].channel == channel_name_allowed assert result.results[0].message_id == 'msg1' - assert getattr(result.results[0], 'error', None) is None + assert isinstance(result.results[1], BatchPublishFailureResult) assert result.results[1].channel == channel_name_restricted assert result.results[1].error.code == 40160 # UTS: rest/unit/RSC22c/distinguish-success-failure-0 -@deviation async def test_rsc22c_batch_publish_distinguish_success_failure(): channel_name = f'test-BatchResult2-{random_id()}' failed_channel_name = f'test-BatchResult2-failed-{random_id()}' mock_http = MockHttpClient( on_connection_attempt=lambda conn: conn.respond_with_success(), - on_request=respond_with(201, [ + on_request=respond_with(201, [batch_result( success_result(channel_name, 'msg1', ['s1']), { 'channel': failed_channel_name, 'error': {'code': 40160, 'statusCode': 401, 'message': 'Not permitted'}, }, - ]), + )]), ) client = rest_client(mock_http) - result = await client.batch_publish({ - 'channels': [channel_name, failed_channel_name], - 'messages': [Message(name='event', data='hello')], - }) + result = await client.batch_publish(BatchPublishSpec( + channels=[channel_name, failed_channel_name], + messages=[Message(name='event', data='hello')], + )) for entry in result.results: has_success_fields = getattr(entry, 'message_id', None) is not None and \ @@ -615,7 +625,6 @@ async def test_rsc22c_batch_publish_distinguish_success_failure(): # UTS: rest/unit/RSC22/empty-channels-rejected-0 -@deviation async def test_rsc22_batch_publish_empty_channels_rejected(): mock_http = MockHttpClient( on_connection_attempt=lambda conn: conn.respond_with_success(), @@ -624,14 +633,13 @@ async def test_rsc22_batch_publish_empty_channels_rejected(): client = rest_client(mock_http) with pytest.raises(AblyException) as excinfo: - await client.batch_publish({'channels': [], 'messages': [Message(name='event', data='hello')]}) + await client.batch_publish(BatchPublishSpec(channels=[], messages=[Message(name='event', data='hello')])) # NOTE: the spec says only "the error indicates invalid request"; read as a 400. assert excinfo.value.status_code == 400 # UTS: rest/unit/RSC22/empty-messages-rejected-0 -@deviation async def test_rsc22_batch_publish_empty_messages_rejected(): channel_name = f'test-RSC22-Error2-{random_id()}' @@ -642,14 +650,13 @@ async def test_rsc22_batch_publish_empty_messages_rejected(): client = rest_client(mock_http) with pytest.raises(AblyException) as excinfo: - await client.batch_publish({'channels': [channel_name], 'messages': []}) + await client.batch_publish(BatchPublishSpec(channels=[channel_name], messages=[])) # NOTE: the spec says only "the error indicates invalid request"; read as a 400. assert excinfo.value.status_code == 400 # UTS: rest/unit/RSC22/server-error-propagated-0 -@deviation async def test_rsc22_batch_publish_server_error_propagated(): channel_name = f'test-RSC22-Error3-{random_id()}' @@ -662,17 +669,16 @@ async def test_rsc22_batch_publish_server_error_propagated(): client = rest_client(mock_http) with pytest.raises(AblyException) as excinfo: - await client.batch_publish({ - 'channels': [channel_name], - 'messages': [Message(name='event', data='hello')], - }) + await client.batch_publish(BatchPublishSpec( + channels=[channel_name], + messages=[Message(name='event', data='hello')], + )) assert excinfo.value.code == 50000 assert excinfo.value.status_code == 500 # UTS: rest/unit/RSC22/auth-error-propagated-0 -@deviation async def test_rsc22_batch_publish_auth_error_propagated(): channel_name = f'test-RSC22-Error4-{random_id()}' @@ -685,17 +691,16 @@ async def test_rsc22_batch_publish_auth_error_propagated(): client = rest_client(mock_http) with pytest.raises(AblyException) as excinfo: - await client.batch_publish({ - 'channels': [channel_name], - 'messages': [Message(name='event', data='hello')], - }) + await client.batch_publish(BatchPublishSpec( + channels=[channel_name], + messages=[Message(name='event', data='hello')], + )) assert excinfo.value.code == 40101 assert excinfo.value.status_code == 401 # UTS: rest/unit/RSC22/standard-headers-included-0 -@deviation async def test_rsc22_batch_publish_standard_headers_included(): channel_name = f'test-RSC22-Headers1-{random_id()}' @@ -706,10 +711,10 @@ async def test_rsc22_batch_publish_standard_headers_included(): ) client = rest_client(mock_http) - await client.batch_publish({ - 'channels': [channel_name], - 'messages': [Message(name='event', data='hello')], - }) + await client.batch_publish(BatchPublishSpec( + channels=[channel_name], + messages=[Message(name='event', data='hello')], + )) assert len(captured_requests) == 1 request = captured_requests[0] @@ -726,7 +731,6 @@ async def test_rsc22_batch_publish_standard_headers_included(): # UTS: rest/unit/RSC22/request-id-included-0 -@deviation async def test_rsc22_batch_publish_request_id_included(): channel_name = f'test-RSC22-Headers2-{random_id()}' @@ -737,10 +741,10 @@ async def test_rsc22_batch_publish_request_id_included(): ) client = rest_client(mock_http, add_request_ids=True) - await client.batch_publish({ - 'channels': [channel_name], - 'messages': [Message(name='event', data='hello')], - }) + await client.batch_publish(BatchPublishSpec( + channels=[channel_name], + messages=[Message(name='event', data='hello')], + )) assert len(captured_requests) == 1 request_id = captured_requests[0].url.query_params['request_id'] @@ -751,21 +755,20 @@ async def test_rsc22_batch_publish_request_id_included(): # UTS: rest/unit/RSC22/multiple-messages-per-channel-0 -@deviation async def test_rsc22_batch_publish_multiple_messages_per_channel(): channel_name = f'test-RSC22-Batch1-{random_id()}' captured_requests = [] mock_http = MockHttpClient( on_connection_attempt=lambda conn: conn.respond_with_success(), - on_request=capture_and_respond(captured_requests, body=[ + on_request=capture_and_respond(captured_requests, body=[batch_result( success_result(channel_name, 'msg', [f's{index}' for index in range(100)]), - ]), + )]), ) client = rest_client(mock_http) messages = [Message(name=f'event-{index}', data=f'data-{index}') for index in range(100)] - result = await client.batch_publish({'channels': [channel_name], 'messages': messages}) + result = await client.batch_publish(BatchPublishSpec(channels=[channel_name], messages=messages)) assert len(captured_requests) == 1 sent_messages = msgpack.unpackb(captured_requests[0].body)['messages'] @@ -777,7 +780,6 @@ async def test_rsc22_batch_publish_multiple_messages_per_channel(): # UTS: rest/unit/RSC22/multiple-channels-multiple-messages-0 -@deviation async def test_rsc22_batch_publish_multiple_channels_multiple_messages(): channel_name_1 = f'test-RSC22-Batch2-a-{random_id()}' channel_name_2 = f'test-RSC22-Batch2-b-{random_id()}' @@ -786,22 +788,22 @@ async def test_rsc22_batch_publish_multiple_channels_multiple_messages(): captured_requests = [] mock_http = MockHttpClient( on_connection_attempt=lambda conn: conn.respond_with_success(), - on_request=capture_and_respond(captured_requests, body=[ + on_request=capture_and_respond(captured_requests, body=[batch_result( success_result(channel_name_1, 'msg1', ['s1', 's2', 's3']), success_result(channel_name_2, 'msg2', ['s4', 's5', 's6']), success_result(channel_name_3, 'msg3', ['s7', 's8', 's9']), - ]), + )]), ) client = rest_client(mock_http) - result = await client.batch_publish({ - 'channels': [channel_name_1, channel_name_2, channel_name_3], - 'messages': [ + result = await client.batch_publish(BatchPublishSpec( + channels=[channel_name_1, channel_name_2, channel_name_3], + messages=[ Message(name='msg1', data='d1'), Message(name='msg2', data='d2'), Message(name='msg3', data='d3'), ], - }) + )) assert len(captured_requests) == 1 body = msgpack.unpackb(captured_requests[0].body)