Repository navigation
feat: add batch publish and batch presence to the HTTP client #727
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: integration/v4
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -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)}) | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🗄️ Data Integrity & Integration | 🟠 Major | ⚡ Quick win Choose a delimiter that is absent from the channel names. If a channel is named 🤖 Prompt for AI Agents |
||
| 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 | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
🗄️ Data Integrity & Integration | 🟠 Major | 🏗️ Heavy lift
Generate IDs per spec when Message objects are shared.
If two specs reuse the same ID-less
Messageand target the same channel, the first iteration sets its ID. The second iteration keeps that ID, and both specs serialize it unchanged. Ably can then discard the second publication as a duplicate, although the caller supplied two specs. Build a separate wire representation for each spec before assigning IDs, while preserving IDs supplied by the caller. Separate batch specs support separate publications, and Ably deduplicates repeated message IDs on a channel. (ably.com)🤖 Prompt for AI Agents