ringcentral.websocket.web_socket_subscription
1#!/usr/bin/env python 2# encoding: utf-8 3import uuid 4import json 5from .events import WebSocketEvents 6from observable import Observable 7 8# _subscription format example: https://git.ringcentral.com/platform/wsg/-/blob/master/RingCentral_WebSocket_API.md#step-4-subscribing-to-rc-events 9 10 11class WebSocketSubscription(Observable): 12 def __init__(self, web_socket_client): 13 Observable.__init__(self) 14 self._web_socket_client = web_socket_client 15 self._event_filters = [] 16 self._subscription = None 17 self._pending_creation_message_id = None 18 self._pending_update_message_id = None 19 self._pending_update_filters = None 20 self._pending_removal_message_id = None 21 self._receive_message_listener_attached = False 22 23 def _operation_pending(self): 24 return ( 25 self._pending_creation_message_id is not None 26 or self._pending_update_message_id is not None 27 or self._pending_removal_message_id is not None 28 ) 29 30 def on_message(self, message): 31 message_json = json.loads(message) 32 if( 33 self._pending_creation_message_id is not None 34 and message_json[0].get('type') == 'ClientRequest' 35 and message_json[0].get('messageId') == self._pending_creation_message_id 36 ): 37 status = message_json[0].get('status', 0) 38 self._pending_creation_message_id = None 39 if 200 <= status < 300: 40 self.set_subscription(message_json) 41 self._web_socket_client.trigger(WebSocketEvents.subscriptionCreated, self) 42 else: 43 error = Exception(f"WebSocket subscription creation failed with status {status}") 44 self._web_socket_client.trigger(WebSocketEvents.createSubscriptionError, error) 45 elif( 46 self._pending_update_message_id is not None 47 and message_json[0].get('messageId') == self._pending_update_message_id 48 ): 49 status = message_json[0].get('status', 0) 50 proposed_filters = self._pending_update_filters 51 self._pending_update_message_id = None 52 self._pending_update_filters = None 53 if 200 <= status < 300: 54 self.set_subscription(message_json) 55 self.set_events(proposed_filters) 56 self._web_socket_client.trigger(WebSocketEvents.subscriptionUpdated, self) 57 else: 58 error = Exception(f"WebSocket subscription update failed with status {status}") 59 self._web_socket_client.trigger(WebSocketEvents.updateSubscriptionError, error) 60 elif( 61 self._pending_removal_message_id is not None 62 and message_json[0].get('messageId') == self._pending_removal_message_id 63 ): 64 status = message_json[0].get('status', 0) 65 self._pending_removal_message_id = None 66 if 200 <= status < 300: 67 if self._receive_message_listener_attached: 68 self._web_socket_client.off(WebSocketEvents.receiveMessage, self.on_message) 69 self._receive_message_listener_attached = False 70 self.reset() 71 self._web_socket_client.trigger(WebSocketEvents.subscriptionRemoved) 72 else: 73 error = Exception(f"WebSocket subscription removal failed with status {status}") 74 self._web_socket_client.trigger(WebSocketEvents.removeSubscriptionError, error) 75 elif message_json[0].get('type') == 'ServerNotification': 76 self._web_socket_client.trigger(WebSocketEvents.receiveSubscriptionNotification, message_json) 77 78 async def register(self, events=None): 79 if not self._subscription: 80 await self.subscribe(events=events) 81 else: 82 await self.update(events=events) 83 84 def add_events(self, events): 85 self._event_filters += events 86 pass 87 88 def set_events(self, events): 89 self._event_filters = events 90 91 async def subscribe(self, events=None): 92 if self._pending_creation_message_id is not None: 93 raise Exception("Subscription creation is already in progress") 94 95 if events: 96 self.set_events(events) 97 98 if not self._event_filters or len(self._event_filters) == 0: 99 raise Exception("Events are undefined") 100 101 newly_attached = False 102 try: 103 messageId = str(uuid.uuid4()) 104 self._pending_creation_message_id = messageId 105 requestBodyJson = [ 106 { 107 "type": "ClientRequest", 108 "messageId": messageId, 109 "method": "POST", 110 "path": "/restapi/v1.0/subscription/", 111 }, 112 { 113 "eventFilters": self._event_filters, 114 "deliveryMode": {"transportType": "WebSocket"}, 115 }, 116 ] 117 if not self._receive_message_listener_attached: 118 self._web_socket_client.on(WebSocketEvents.receiveMessage, self.on_message) 119 self._receive_message_listener_attached = True 120 newly_attached = True 121 await self._web_socket_client.send_message(requestBodyJson) 122 123 except Exception as e: 124 self._pending_creation_message_id = None 125 if newly_attached: 126 self._web_socket_client.off(WebSocketEvents.receiveMessage, self.on_message) 127 self._receive_message_listener_attached = False 128 self.reset() 129 print(e) 130 raise 131 132 async def update(self, events=None): 133 if self._operation_pending(): 134 raise Exception("Subscription update is already in progress") 135 136 proposed_filters = events if events else self._event_filters 137 if not proposed_filters or len(proposed_filters) == 0: 138 raise Exception("Events are undefined") 139 140 try: 141 subscriptionId = self._subscription[1]["id"] 142 messageId = str(uuid.uuid4()) 143 self._pending_update_message_id = messageId 144 self._pending_update_filters = proposed_filters 145 requestBodyJson = [ 146 { 147 "type": "ClientRequest", 148 "messageId": messageId, 149 "method": "PUT", 150 "path": f"/restapi/v1.0/subscription/{subscriptionId}", 151 }, 152 { 153 "eventFilters": proposed_filters, 154 "deliveryMode": {"transportType": "WebSocket"}, 155 }, 156 ] 157 await self._web_socket_client.send_message(requestBodyJson) 158 159 except Exception as e: 160 self._pending_update_message_id = None 161 self._pending_update_filters = None 162 print(e) 163 raise 164 165 async def remove(self): 166 if self._operation_pending(): 167 raise Exception("Subscription removal is already in progress") 168 169 subscriptionId = self._subscription[1]["id"] 170 if not subscriptionId: 171 raise Exception("Missing subscriptionId") 172 173 try: 174 messageId = str(uuid.uuid4()) 175 self._pending_removal_message_id = messageId 176 requestBodyJson = [ 177 { 178 "type": "ClientRequest", 179 "messageId": messageId, 180 "method": "DELETE", 181 "path": f"/restapi/v1.0/subscription/{subscriptionId}", 182 } 183 ] 184 185 await self._web_socket_client.send_message(requestBodyJson) 186 187 except Exception as e: 188 self._pending_removal_message_id = None 189 print(e) 190 raise 191 192 def set_subscription(self, data): 193 self._subscription = data 194 195 def get_subscription_info(self): 196 return self._subscription 197 198 def reset(self): 199 self._subscription = None 200 201 def destroy(self): 202 self.reset() 203 self.off() 204 205 206if __name__ == "__main__": 207 pass
class
WebSocketSubscription(observable.core.Observable):
12class WebSocketSubscription(Observable): 13 def __init__(self, web_socket_client): 14 Observable.__init__(self) 15 self._web_socket_client = web_socket_client 16 self._event_filters = [] 17 self._subscription = None 18 self._pending_creation_message_id = None 19 self._pending_update_message_id = None 20 self._pending_update_filters = None 21 self._pending_removal_message_id = None 22 self._receive_message_listener_attached = False 23 24 def _operation_pending(self): 25 return ( 26 self._pending_creation_message_id is not None 27 or self._pending_update_message_id is not None 28 or self._pending_removal_message_id is not None 29 ) 30 31 def on_message(self, message): 32 message_json = json.loads(message) 33 if( 34 self._pending_creation_message_id is not None 35 and message_json[0].get('type') == 'ClientRequest' 36 and message_json[0].get('messageId') == self._pending_creation_message_id 37 ): 38 status = message_json[0].get('status', 0) 39 self._pending_creation_message_id = None 40 if 200 <= status < 300: 41 self.set_subscription(message_json) 42 self._web_socket_client.trigger(WebSocketEvents.subscriptionCreated, self) 43 else: 44 error = Exception(f"WebSocket subscription creation failed with status {status}") 45 self._web_socket_client.trigger(WebSocketEvents.createSubscriptionError, error) 46 elif( 47 self._pending_update_message_id is not None 48 and message_json[0].get('messageId') == self._pending_update_message_id 49 ): 50 status = message_json[0].get('status', 0) 51 proposed_filters = self._pending_update_filters 52 self._pending_update_message_id = None 53 self._pending_update_filters = None 54 if 200 <= status < 300: 55 self.set_subscription(message_json) 56 self.set_events(proposed_filters) 57 self._web_socket_client.trigger(WebSocketEvents.subscriptionUpdated, self) 58 else: 59 error = Exception(f"WebSocket subscription update failed with status {status}") 60 self._web_socket_client.trigger(WebSocketEvents.updateSubscriptionError, error) 61 elif( 62 self._pending_removal_message_id is not None 63 and message_json[0].get('messageId') == self._pending_removal_message_id 64 ): 65 status = message_json[0].get('status', 0) 66 self._pending_removal_message_id = None 67 if 200 <= status < 300: 68 if self._receive_message_listener_attached: 69 self._web_socket_client.off(WebSocketEvents.receiveMessage, self.on_message) 70 self._receive_message_listener_attached = False 71 self.reset() 72 self._web_socket_client.trigger(WebSocketEvents.subscriptionRemoved) 73 else: 74 error = Exception(f"WebSocket subscription removal failed with status {status}") 75 self._web_socket_client.trigger(WebSocketEvents.removeSubscriptionError, error) 76 elif message_json[0].get('type') == 'ServerNotification': 77 self._web_socket_client.trigger(WebSocketEvents.receiveSubscriptionNotification, message_json) 78 79 async def register(self, events=None): 80 if not self._subscription: 81 await self.subscribe(events=events) 82 else: 83 await self.update(events=events) 84 85 def add_events(self, events): 86 self._event_filters += events 87 pass 88 89 def set_events(self, events): 90 self._event_filters = events 91 92 async def subscribe(self, events=None): 93 if self._pending_creation_message_id is not None: 94 raise Exception("Subscription creation is already in progress") 95 96 if events: 97 self.set_events(events) 98 99 if not self._event_filters or len(self._event_filters) == 0: 100 raise Exception("Events are undefined") 101 102 newly_attached = False 103 try: 104 messageId = str(uuid.uuid4()) 105 self._pending_creation_message_id = messageId 106 requestBodyJson = [ 107 { 108 "type": "ClientRequest", 109 "messageId": messageId, 110 "method": "POST", 111 "path": "/restapi/v1.0/subscription/", 112 }, 113 { 114 "eventFilters": self._event_filters, 115 "deliveryMode": {"transportType": "WebSocket"}, 116 }, 117 ] 118 if not self._receive_message_listener_attached: 119 self._web_socket_client.on(WebSocketEvents.receiveMessage, self.on_message) 120 self._receive_message_listener_attached = True 121 newly_attached = True 122 await self._web_socket_client.send_message(requestBodyJson) 123 124 except Exception as e: 125 self._pending_creation_message_id = None 126 if newly_attached: 127 self._web_socket_client.off(WebSocketEvents.receiveMessage, self.on_message) 128 self._receive_message_listener_attached = False 129 self.reset() 130 print(e) 131 raise 132 133 async def update(self, events=None): 134 if self._operation_pending(): 135 raise Exception("Subscription update is already in progress") 136 137 proposed_filters = events if events else self._event_filters 138 if not proposed_filters or len(proposed_filters) == 0: 139 raise Exception("Events are undefined") 140 141 try: 142 subscriptionId = self._subscription[1]["id"] 143 messageId = str(uuid.uuid4()) 144 self._pending_update_message_id = messageId 145 self._pending_update_filters = proposed_filters 146 requestBodyJson = [ 147 { 148 "type": "ClientRequest", 149 "messageId": messageId, 150 "method": "PUT", 151 "path": f"/restapi/v1.0/subscription/{subscriptionId}", 152 }, 153 { 154 "eventFilters": proposed_filters, 155 "deliveryMode": {"transportType": "WebSocket"}, 156 }, 157 ] 158 await self._web_socket_client.send_message(requestBodyJson) 159 160 except Exception as e: 161 self._pending_update_message_id = None 162 self._pending_update_filters = None 163 print(e) 164 raise 165 166 async def remove(self): 167 if self._operation_pending(): 168 raise Exception("Subscription removal is already in progress") 169 170 subscriptionId = self._subscription[1]["id"] 171 if not subscriptionId: 172 raise Exception("Missing subscriptionId") 173 174 try: 175 messageId = str(uuid.uuid4()) 176 self._pending_removal_message_id = messageId 177 requestBodyJson = [ 178 { 179 "type": "ClientRequest", 180 "messageId": messageId, 181 "method": "DELETE", 182 "path": f"/restapi/v1.0/subscription/{subscriptionId}", 183 } 184 ] 185 186 await self._web_socket_client.send_message(requestBodyJson) 187 188 except Exception as e: 189 self._pending_removal_message_id = None 190 print(e) 191 raise 192 193 def set_subscription(self, data): 194 self._subscription = data 195 196 def get_subscription_info(self): 197 return self._subscription 198 199 def reset(self): 200 self._subscription = None 201 202 def destroy(self): 203 self.reset() 204 self.off()
Event system for python
WebSocketSubscription(web_socket_client)
13 def __init__(self, web_socket_client): 14 Observable.__init__(self) 15 self._web_socket_client = web_socket_client 16 self._event_filters = [] 17 self._subscription = None 18 self._pending_creation_message_id = None 19 self._pending_update_message_id = None 20 self._pending_update_filters = None 21 self._pending_removal_message_id = None 22 self._receive_message_listener_attached = False
def
on_message(self, message):
31 def on_message(self, message): 32 message_json = json.loads(message) 33 if( 34 self._pending_creation_message_id is not None 35 and message_json[0].get('type') == 'ClientRequest' 36 and message_json[0].get('messageId') == self._pending_creation_message_id 37 ): 38 status = message_json[0].get('status', 0) 39 self._pending_creation_message_id = None 40 if 200 <= status < 300: 41 self.set_subscription(message_json) 42 self._web_socket_client.trigger(WebSocketEvents.subscriptionCreated, self) 43 else: 44 error = Exception(f"WebSocket subscription creation failed with status {status}") 45 self._web_socket_client.trigger(WebSocketEvents.createSubscriptionError, error) 46 elif( 47 self._pending_update_message_id is not None 48 and message_json[0].get('messageId') == self._pending_update_message_id 49 ): 50 status = message_json[0].get('status', 0) 51 proposed_filters = self._pending_update_filters 52 self._pending_update_message_id = None 53 self._pending_update_filters = None 54 if 200 <= status < 300: 55 self.set_subscription(message_json) 56 self.set_events(proposed_filters) 57 self._web_socket_client.trigger(WebSocketEvents.subscriptionUpdated, self) 58 else: 59 error = Exception(f"WebSocket subscription update failed with status {status}") 60 self._web_socket_client.trigger(WebSocketEvents.updateSubscriptionError, error) 61 elif( 62 self._pending_removal_message_id is not None 63 and message_json[0].get('messageId') == self._pending_removal_message_id 64 ): 65 status = message_json[0].get('status', 0) 66 self._pending_removal_message_id = None 67 if 200 <= status < 300: 68 if self._receive_message_listener_attached: 69 self._web_socket_client.off(WebSocketEvents.receiveMessage, self.on_message) 70 self._receive_message_listener_attached = False 71 self.reset() 72 self._web_socket_client.trigger(WebSocketEvents.subscriptionRemoved) 73 else: 74 error = Exception(f"WebSocket subscription removal failed with status {status}") 75 self._web_socket_client.trigger(WebSocketEvents.removeSubscriptionError, error) 76 elif message_json[0].get('type') == 'ServerNotification': 77 self._web_socket_client.trigger(WebSocketEvents.receiveSubscriptionNotification, message_json)
async def
subscribe(self, events=None):
92 async def subscribe(self, events=None): 93 if self._pending_creation_message_id is not None: 94 raise Exception("Subscription creation is already in progress") 95 96 if events: 97 self.set_events(events) 98 99 if not self._event_filters or len(self._event_filters) == 0: 100 raise Exception("Events are undefined") 101 102 newly_attached = False 103 try: 104 messageId = str(uuid.uuid4()) 105 self._pending_creation_message_id = messageId 106 requestBodyJson = [ 107 { 108 "type": "ClientRequest", 109 "messageId": messageId, 110 "method": "POST", 111 "path": "/restapi/v1.0/subscription/", 112 }, 113 { 114 "eventFilters": self._event_filters, 115 "deliveryMode": {"transportType": "WebSocket"}, 116 }, 117 ] 118 if not self._receive_message_listener_attached: 119 self._web_socket_client.on(WebSocketEvents.receiveMessage, self.on_message) 120 self._receive_message_listener_attached = True 121 newly_attached = True 122 await self._web_socket_client.send_message(requestBodyJson) 123 124 except Exception as e: 125 self._pending_creation_message_id = None 126 if newly_attached: 127 self._web_socket_client.off(WebSocketEvents.receiveMessage, self.on_message) 128 self._receive_message_listener_attached = False 129 self.reset() 130 print(e) 131 raise
async def
update(self, events=None):
133 async def update(self, events=None): 134 if self._operation_pending(): 135 raise Exception("Subscription update is already in progress") 136 137 proposed_filters = events if events else self._event_filters 138 if not proposed_filters or len(proposed_filters) == 0: 139 raise Exception("Events are undefined") 140 141 try: 142 subscriptionId = self._subscription[1]["id"] 143 messageId = str(uuid.uuid4()) 144 self._pending_update_message_id = messageId 145 self._pending_update_filters = proposed_filters 146 requestBodyJson = [ 147 { 148 "type": "ClientRequest", 149 "messageId": messageId, 150 "method": "PUT", 151 "path": f"/restapi/v1.0/subscription/{subscriptionId}", 152 }, 153 { 154 "eventFilters": proposed_filters, 155 "deliveryMode": {"transportType": "WebSocket"}, 156 }, 157 ] 158 await self._web_socket_client.send_message(requestBodyJson) 159 160 except Exception as e: 161 self._pending_update_message_id = None 162 self._pending_update_filters = None 163 print(e) 164 raise
async def
remove(self):
166 async def remove(self): 167 if self._operation_pending(): 168 raise Exception("Subscription removal is already in progress") 169 170 subscriptionId = self._subscription[1]["id"] 171 if not subscriptionId: 172 raise Exception("Missing subscriptionId") 173 174 try: 175 messageId = str(uuid.uuid4()) 176 self._pending_removal_message_id = messageId 177 requestBodyJson = [ 178 { 179 "type": "ClientRequest", 180 "messageId": messageId, 181 "method": "DELETE", 182 "path": f"/restapi/v1.0/subscription/{subscriptionId}", 183 } 184 ] 185 186 await self._web_socket_client.send_message(requestBodyJson) 187 188 except Exception as e: 189 self._pending_removal_message_id = None 190 print(e) 191 raise
Inherited Members
- observable.core.Observable
- events
- on
- off
- once
- trigger