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 register(self, events=None):
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)
def add_events(self, events):
85    def add_events(self, events):
86        self._event_filters += events
87        pass
def set_events(self, events):
89    def set_events(self, events):
90        self._event_filters = events
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
def set_subscription(self, data):
193    def set_subscription(self, data):
194        self._subscription = data
def get_subscription_info(self):
196    def get_subscription_info(self):
197        return self._subscription
def reset(self):
199    def reset(self):
200        self._subscription = None
def destroy(self):
202    def destroy(self):
203        self.reset()
204        self.off()
Inherited Members
observable.core.Observable
events
on
off
once
trigger