Uh oh!
There was an error while loading. Please reload this page.
- Notifications
You must be signed in to change notification settings - Fork 27
Realtime presence#651
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
Merged
Uh oh!
There was an error while loading. Please reload this page.
Merged
Realtime presence #651
Changes from all commits
Commits
Show all changes
11 commits
Select commit
Hold shift + click to select a range
03ec3ff
add presencemap and some presence helpers
owenpearson 140c40a
fix: don't resume connection when explicitly closed
owenpearson 8e39e30
add connection.when_state
owenpearson a2fba41
clear connection state as soon as entering SUSPENDED
owenpearson 25e9257
add support for arbitrary transport params
owenpearson c012993
attempt to send `CLOSE` message on connection close
owenpearson f18382b
implement realtime presence publish/subscribe functionality
owenpearson d3da55b
improve presence encoding and support encryption
owenpearson ca3f388
update presence tests for pytest asyncio
owenpearson d67f987
parameterise realtime presence tests with msgpack and json
owenpearson e1779bb
fix: hanging in wait_for_sync when channel DETACHED/FAILED
owenpearson File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Uh oh!
There was an error while loading. Please reload this page.
Jump to
Jump to file
Failed to load files.
Loading
Uh oh!
There was an error while loading. Please reload this page.
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -135,6 +135,17 @@ def enact_state_change(self, state: ConnectionState, reason: AblyException | Non | ||
| self.__state = state | ||
| if reason: | ||
| self.__error_reason = reason | ||
| # RTN16d: Clear connection state when entering SUSPENDED or terminal states | ||
| if state == ConnectionState.SUSPENDED or state in ( | ||
| ConnectionState.CLOSED, | ||
| ConnectionState.FAILED | ||
| ): | ||
| self.__connection_details = None | ||
| self.connection_id = None | ||
| self.__connection_key = None | ||
| self.msg_serial = 0 | ||
| self._emit('connectionstate', ConnectionStateChange(current_state, state, state, reason)) | ||
| def check_connection(self) -> bool: | ||
| @@ -157,6 +168,10 @@ async def __get_transport_params(self) -> dict: | ||
| # RTN2a: Set format to msgpack if use_binary_protocol is enabled | ||
| if self.options.use_binary_protocol: | ||
| params["format"] = "msgpack" | ||
| # Add any custom transport params from options | ||
| params.update(self.options.transport_params) | ||
| return params | ||
| async def close_impl(self) -> None: | ||
| @@ -165,13 +180,23 @@ async def close_impl(self) -> None: | ||
| self.cancel_suspend_timer() | ||
| self.start_transition_timer(ConnectionState.CLOSING, fail_state=ConnectionState.CLOSED) | ||
| if self.transport: | ||
| await self.transport.dispose() | ||
| # Try to send protocol CLOSE message in the background | ||
| asyncio.create_task(self.transport.close()) | ||
| # Yield to event loop to give the close message a chance to send | ||
| await asyncio.sleep(0) | ||
| await self.transport.dispose() # Dispose transport resources | ||
| if self.connect_base_task: | ||
| self.connect_base_task.cancel() | ||
| if self.disconnect_transport_task: | ||
| await self.disconnect_transport_task | ||
| self.cancel_retry_timer() | ||
| # Clear connection details to prevent resume on next connect | ||
| # When explicitly closed, we want a fresh connection, not a resume | ||
| self.__connection_details = None | ||
| self.connection_id = None | ||
| self.msg_serial = 0 | ||
coderabbitai[bot] marked this conversation as resolved.
Uh oh!There was an error while loading. Please reload this page. | ||
| self.notify_state(ConnectionState.CLOSED) | ||
| async def send_protocol_message(self, protocol_message: dict) -> None: | ||
| @@ -648,7 +673,6 @@ def on_suspend_timer_expire() -> None: | ||
| AblyException("Connection to server unavailable", 400, 80002) | ||
| ) | ||
| self.__fail_state = ConnectionState.SUSPENDED | ||
| self.__connection_details = None | ||
| self.suspend_timer = Timer(Defaults.connection_state_ttl, on_suspend_timer_expire) | ||
Oops, something went wrong.
Uh oh!
There was an error while loading. Please reload this page.
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.