ringcentral.websocket.web_socket_subscription_test

  1import json
  2import unittest
  3import uuid
  4
  5from observable import Observable
  6
  7from .events import WebSocketEvents
  8from .web_socket_subscription import WebSocketSubscription
  9
 10
 11class RecordingHandler:
 12    def __init__(self):
 13        self.calls = []
 14
 15    def __call__(self, *args):
 16        self.calls.append(args)
 17
 18
 19class FakeWebSocketClient(Observable):
 20    def __init__(self, responder=None, send_error=None):
 21        Observable.__init__(self)
 22        self.sent_messages = []
 23        self.receive_message_listener_count = 0
 24        self._responder = responder
 25        self._send_error = send_error
 26
 27    def on(self, event, *handlers):
 28        if event == WebSocketEvents.receiveMessage:
 29            self.receive_message_listener_count += len(handlers)
 30        return Observable.on(self, event, *handlers)
 31
 32    def off(self, event=None, *handlers):
 33        if event == WebSocketEvents.receiveMessage:
 34            self.receive_message_listener_count -= len(handlers)
 35        return Observable.off(self, event, *handlers)
 36
 37    async def send_message(self, message):
 38        self.sent_messages.append(message)
 39        if self._send_error is not None:
 40            error = self._send_error
 41            self._send_error = None
 42            raise error
 43        if self._responder is not None:
 44            response = self._responder(message)
 45            if response is not None:
 46                self.trigger(WebSocketEvents.receiveMessage, json.dumps(response))
 47
 48
 49def creation_response(request):
 50    return [
 51        {
 52            "type": "ClientRequest",
 53            "messageId": request[0]["messageId"],
 54            "status": 200,
 55            "headers": {
 56                "Server": "nginx",
 57                "Date": "Wed, 20 Aug 2025 22:23:55 GMT",
 58                "Content-Type": "application/json",
 59                "RoutingKey": "SJC01P07",
 60                "RCRequestId": "bedff5ae-9d68-4bc9-8653-7e34603ef562-2686696-1-19",
 61            },
 62        },
 63        {
 64            "uri": "/restapi/v1.0/subscription/1b2a2e6b-2245-4278-b47c-16259ca003a8",
 65            "id": "1b2a2e6b-2245-4278-b47c-16259ca003a8",
 66            "creationTime": "2025-08-20T22:23:55.169Z",
 67            "status": "Active",
 68            "eventFilters": ["/restapi/v1.0/account/809646016/extension/62264425016/presence"],
 69            "expirationTime": "2025-08-21T22:23:55.169Z",
 70            "expiresIn": 86399,
 71            "deliveryMode": {"transportType": "WebSocket", "encryption": False},
 72        },
 73    ]
 74
 75
 76def rejected_creation_response(request, status=403):
 77    response = creation_response(request)
 78    response[0]["status"] = status
 79    return response
 80
 81
 82def update_response(request, status=200, message_type="ClientRequest"):
 83    return [
 84        {
 85            "type": message_type,
 86            "messageId": request[0]["messageId"],
 87            "status": status,
 88            "headers": {
 89                "Server": "nginx",
 90                "Date": "Wed, 20 Aug 2025 22:23:55 GMT",
 91                "Content-Type": "application/json",
 92                "RoutingKey": "SJC01P07",
 93                "RCRequestId": "bedff5ae-9d68-4bc9-8653-7e34603ef562-2686696-1-19",
 94            },
 95        },
 96        {
 97            "uri": "/restapi/v1.0/subscription/9d3b7f10-5c11-4b6e-8a2f-6f8b56f0f1c2",
 98            "id": "9d3b7f10-5c11-4b6e-8a2f-6f8b56f0f1c2",
 99            "creationTime": "2025-08-20T22:23:55.169Z",
100            "status": "Active",
101            "eventFilters": request[1]["eventFilters"],
102            "expirationTime": "2025-08-21T22:23:55.169Z",
103            "expiresIn": 86399,
104            "deliveryMode": {"transportType": "WebSocket", "encryption": False},
105        },
106    ]
107
108
109def removal_response(request, status=200, message_type="ClientRequest"):
110    return [
111        {
112            "type": message_type,
113            "messageId": request[0]["messageId"],
114            "status": status,
115            "headers": {
116                "Server": "nginx",
117                "Date": "Wed, 20 Aug 2025 22:23:55 GMT",
118                "RCRequestId": "bedff5ae-9d68-4bc9-8653-7e34603ef562-2686696-1-19",
119            },
120        }
121    ]
122
123
124EVENT_FILTERS = ["/restapi/v1.0/account/~/extension/~/presence"]
125OTHER_FILTERS = ["/restapi/v1.0/account/~/extension/~/message-store"]
126
127
128def server_notification():
129    return [
130        {
131            "type": "ServerNotification",
132            "messageId": str(uuid.uuid4()),
133            "headers": {"RoutingKey": "SJC01P07"},
134        },
135        {
136            "uri": "/restapi/v1.0/subscription/1b2a2e6b-2245-4278-b47c-16259ca003a8",
137            "event": {"/restapi/v1.0/account/~/extension/~/presence": {"activeCalls": []}},
138        },
139    ]
140
141
142class WebSocketSubscriptionTest(unittest.IsolatedAsyncioTestCase):
143    def setUp(self):
144        self.web_socket_client = FakeWebSocketClient()
145
146    async def test_successful_creation_before_send_returns_stores_response_and_emits_event_once(self):
147        self.web_socket_client._responder = creation_response
148        subscription = WebSocketSubscription(self.web_socket_client)
149        created = RecordingHandler()
150        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
151
152        await subscription.subscribe(events=EVENT_FILTERS)
153
154        self.assertEqual(len(created.calls), 1)
155        self.assertIs(created.calls[0][0], subscription)
156        self.assertEqual(len(self.web_socket_client.sent_messages), 1)
157        stored = subscription.get_subscription_info()
158        self.assertIsNotNone(stored)
159        self.assertEqual(stored[0]["type"], "ClientRequest")
160        self.assertEqual(
161            stored[0]["messageId"], self.web_socket_client.sent_messages[0][0]["messageId"]
162        )
163        self.assertNotIn("WSG-SubscriptionId", stored[0]["headers"])
164        self.assertEqual(
165            stored[1]["id"], "1b2a2e6b-2245-4278-b47c-16259ca003a8"
166        )
167
168    async def test_unrelated_message_id_response_does_not_change_state_or_emit_events(self):
169        def unrelated_response(request):
170            response = creation_response(request)
171            response[0]["messageId"] = "unrelated-" + request[0]["messageId"]
172            return response
173
174        self.web_socket_client._responder = unrelated_response
175        subscription = WebSocketSubscription(self.web_socket_client)
176        created = RecordingHandler()
177        failed = RecordingHandler()
178        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
179        self.web_socket_client.on(WebSocketEvents.createSubscriptionError, failed)
180
181        await subscription.subscribe(events=EVENT_FILTERS)
182
183        self.assertIsNone(subscription.get_subscription_info())
184        self.assertEqual(created.calls, [])
185        self.assertEqual(failed.calls, [])
186
187    async def test_listener_remains_after_creation_so_notifications_are_emitted(self):
188        self.web_socket_client._responder = creation_response
189        subscription = WebSocketSubscription(self.web_socket_client)
190        created = RecordingHandler()
191        notifications = RecordingHandler()
192        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
193        self.web_socket_client.on(WebSocketEvents.receiveSubscriptionNotification, notifications)
194
195        await subscription.subscribe(events=EVENT_FILTERS)
196        self.assertEqual(len(created.calls), 1)
197
198        notification = server_notification()
199        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(notification))
200
201        self.assertEqual(len(notifications.calls), 1)
202        self.assertEqual(notifications.calls[0][0], notification)
203
204    async def test_rejected_creation_emits_single_error_without_state_or_success(self):
205        self.web_socket_client._responder = rejected_creation_response
206        subscription = WebSocketSubscription(self.web_socket_client)
207        created = RecordingHandler()
208        failed = RecordingHandler()
209        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
210        self.web_socket_client.on(WebSocketEvents.createSubscriptionError, failed)
211
212        await subscription.subscribe(events=EVENT_FILTERS)
213
214        self.assertEqual(len(failed.calls), 1)
215        error = failed.calls[0][0]
216        self.assertIsInstance(error, Exception)
217        self.assertEqual(
218            str(error), "WebSocket subscription creation failed with status 403"
219        )
220        self.assertEqual(created.calls, [])
221        self.assertIsNone(subscription.get_subscription_info())
222
223    async def test_retry_after_rejected_creation_sends_new_request_and_can_succeed(self):
224        self.web_socket_client._responder = rejected_creation_response
225        subscription = WebSocketSubscription(self.web_socket_client)
226        created = RecordingHandler()
227        failed = RecordingHandler()
228        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
229        self.web_socket_client.on(WebSocketEvents.createSubscriptionError, failed)
230
231        await subscription.subscribe(events=EVENT_FILTERS)
232        self.assertEqual(len(failed.calls), 1)
233
234        self.web_socket_client._responder = creation_response
235        await subscription.subscribe(events=EVENT_FILTERS)
236
237        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
238        self.assertEqual(len(failed.calls), 1)
239        self.assertEqual(len(created.calls), 1)
240        self.assertIsNotNone(subscription.get_subscription_info())
241
242    async def test_second_creation_while_pending_is_rejected_without_new_request_or_listener(self):
243        self.web_socket_client._responder = None
244        subscription = WebSocketSubscription(self.web_socket_client)
245        created = RecordingHandler()
246        failed = RecordingHandler()
247        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
248        self.web_socket_client.on(WebSocketEvents.createSubscriptionError, failed)
249
250        await subscription.subscribe(events=EVENT_FILTERS)
251        self.assertEqual(len(self.web_socket_client.sent_messages), 1)
252        listeners_before = self.web_socket_client.receive_message_listener_count
253
254        with self.assertRaises(Exception):
255            await subscription.subscribe(events=EVENT_FILTERS)
256
257        self.assertEqual(len(self.web_socket_client.sent_messages), 1)
258        self.assertEqual(
259            self.web_socket_client.receive_message_listener_count, listeners_before
260        )
261
262        self.web_socket_client._responder = creation_response
263        pending_request = self.web_socket_client.sent_messages[0]
264        self.web_socket_client.trigger(
265            WebSocketEvents.receiveMessage, json.dumps(creation_response(pending_request))
266        )
267
268        self.assertEqual(len(created.calls), 1)
269        self.assertEqual(failed.calls, [])
270        self.assertIsNotNone(subscription.get_subscription_info())
271
272    async def test_second_creation_while_pending_leaves_original_filters_for_update(self):
273        self.web_socket_client._responder = None
274        subscription = WebSocketSubscription(self.web_socket_client)
275        created = RecordingHandler()
276        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
277
278        await subscription.subscribe(events=EVENT_FILTERS)
279        self.assertEqual(len(self.web_socket_client.sent_messages), 1)
280
281        with self.assertRaises(Exception):
282            await subscription.subscribe(events=OTHER_FILTERS)
283
284        pending_request = self.web_socket_client.sent_messages[0]
285        self.web_socket_client._responder = creation_response
286        self.web_socket_client.trigger(
287            WebSocketEvents.receiveMessage, json.dumps(creation_response(pending_request))
288        )
289        self.assertEqual(len(created.calls), 1)
290
291        await subscription.update()
292
293        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
294        self.assertEqual(
295            self.web_socket_client.sent_messages[1][1]["eventFilters"], EVENT_FILTERS
296        )
297
298    async def test_failed_send_clears_pending_and_new_listener_and_retry_succeeds(self):
299        self.web_socket_client._responder = creation_response
300        self.web_socket_client._send_error = Exception("connection closed")
301        subscription = WebSocketSubscription(self.web_socket_client)
302        created = RecordingHandler()
303        notifications = RecordingHandler()
304        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
305        self.web_socket_client.on(WebSocketEvents.receiveSubscriptionNotification, notifications)
306
307        with self.assertRaises(Exception):
308            await subscription.subscribe(events=EVENT_FILTERS)
309
310        self.assertEqual(len(self.web_socket_client.sent_messages), 1)
311        self.assertIsNone(subscription.get_subscription_info())
312        self.assertEqual(created.calls, [])
313        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
314
315        await subscription.subscribe(events=EVENT_FILTERS)
316
317        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
318        self.assertEqual(len(created.calls), 1)
319        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
320
321        notification = server_notification()
322        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(notification))
323        self.assertEqual(len(notifications.calls), 1)
324
325    async def test_removal_detaches_listener_and_later_retry_receives_notifications_once(self):
326        self.web_socket_client._responder = creation_response
327        subscription = WebSocketSubscription(self.web_socket_client)
328        notifications = RecordingHandler()
329        self.web_socket_client.on(WebSocketEvents.receiveSubscriptionNotification, notifications)
330
331        await subscription.subscribe(events=EVENT_FILTERS)
332        notification = server_notification()
333        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(notification))
334        self.assertEqual(len(notifications.calls), 1)
335
336        await subscription.remove()
337        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(server_notification()))
338        self.assertEqual(len(notifications.calls), 1)
339        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
340
341        await subscription.subscribe(events=EVENT_FILTERS)
342        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(server_notification()))
343        self.assertEqual(len(notifications.calls), 2)
344        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
345
346    async def test_update_send_without_response_emits_nothing_and_preserves_confirmed_state(self):
347        self.web_socket_client._responder = creation_response
348        subscription = WebSocketSubscription(self.web_socket_client)
349        await subscription.subscribe(events=EVENT_FILTERS)
350        confirmed = subscription.get_subscription_info()
351
352        updated = RecordingHandler()
353        failed = RecordingHandler()
354        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
355        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
356
357        self.web_socket_client._responder = None
358        await subscription.update(events=OTHER_FILTERS)
359
360        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
361        self.assertEqual(self.web_socket_client.sent_messages[1][0]["method"], "PUT")
362        self.assertEqual(self.web_socket_client.sent_messages[1][1]["eventFilters"], OTHER_FILTERS)
363        self.assertEqual(updated.calls, [])
364        self.assertEqual(failed.calls, [])
365        self.assertIs(subscription.get_subscription_info(), confirmed)
366        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
367
368    async def test_confirmed_update_response_commits_envelope_and_filters_before_single_event(self):
369        self.web_socket_client._responder = creation_response
370        subscription = WebSocketSubscription(self.web_socket_client)
371        await subscription.subscribe(events=EVENT_FILTERS)
372
373        observed_at_event = []
374
375        def on_updated(updated_subscription):
376            observed_at_event.append(updated_subscription.get_subscription_info())
377
378        updated = RecordingHandler()
379        failed = RecordingHandler()
380        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, on_updated)
381        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
382        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
383
384        self.web_socket_client._responder = update_response
385        await subscription.update(events=OTHER_FILTERS)
386
387        self.assertEqual(len(updated.calls), 1)
388        self.assertIs(updated.calls[0][0], subscription)
389        self.assertEqual(failed.calls, [])
390        stored = subscription.get_subscription_info()
391        self.assertEqual(
392            stored[0]["messageId"], self.web_socket_client.sent_messages[1][0]["messageId"]
393        )
394        self.assertEqual(stored[1]["id"], "9d3b7f10-5c11-4b6e-8a2f-6f8b56f0f1c2")
395        self.assertEqual(stored[1]["eventFilters"], OTHER_FILTERS)
396        self.assertEqual(observed_at_event, [stored])
397
398        self.web_socket_client._responder = None
399        await subscription.update()
400        self.assertEqual(self.web_socket_client.sent_messages[2][1]["eventFilters"], OTHER_FILTERS)
401
402    async def test_rejected_update_response_emits_single_error_preserves_state_and_permits_retry(self):
403        self.web_socket_client._responder = creation_response
404        subscription = WebSocketSubscription(self.web_socket_client)
405        await subscription.subscribe(events=EVENT_FILTERS)
406        confirmed = subscription.get_subscription_info()
407
408        updated = RecordingHandler()
409        failed = RecordingHandler()
410        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
411        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
412
413        self.web_socket_client._responder = lambda request: update_response(request, status=404)
414        await subscription.update(events=OTHER_FILTERS)
415
416        self.assertEqual(len(failed.calls), 1)
417        error = failed.calls[0][0]
418        self.assertIsInstance(error, Exception)
419        self.assertEqual(
420            str(error), "WebSocket subscription update failed with status 404"
421        )
422        self.assertEqual(updated.calls, [])
423        self.assertIs(subscription.get_subscription_info(), confirmed)
424
425        self.web_socket_client._responder = update_response
426        await subscription.update()
427
428        self.assertEqual(len(updated.calls), 1)
429        self.assertEqual(self.web_socket_client.sent_messages[2][1]["eventFilters"], EVENT_FILTERS)
430
431    async def test_unrelated_update_response_is_ignored_and_pending_update_still_completes(self):
432        def unrelated_update_response(request):
433            response = update_response(request)
434            response[0]["messageId"] = "unrelated-" + request[0]["messageId"]
435            return response
436
437        self.web_socket_client._responder = creation_response
438        subscription = WebSocketSubscription(self.web_socket_client)
439        await subscription.subscribe(events=EVENT_FILTERS)
440        confirmed = subscription.get_subscription_info()
441
442        updated = RecordingHandler()
443        failed = RecordingHandler()
444        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
445        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
446
447        self.web_socket_client._responder = unrelated_update_response
448        await subscription.update(events=OTHER_FILTERS)
449
450        self.assertEqual(updated.calls, [])
451        self.assertEqual(failed.calls, [])
452        self.assertIs(subscription.get_subscription_info(), confirmed)
453
454        pending_request = self.web_socket_client.sent_messages[1]
455        self.web_socket_client.trigger(
456            WebSocketEvents.receiveMessage, json.dumps(update_response(pending_request))
457        )
458
459        self.assertEqual(len(updated.calls), 1)
460        self.assertEqual(failed.calls, [])
461        self.assertEqual(
462            subscription.get_subscription_info()[0]["messageId"], pending_request[0]["messageId"]
463        )
464
465    async def test_update_send_failure_preserves_state_and_permits_retry_after_late_response(self):
466        self.web_socket_client._responder = creation_response
467        subscription = WebSocketSubscription(self.web_socket_client)
468        await subscription.subscribe(events=EVENT_FILTERS)
469        confirmed = subscription.get_subscription_info()
470
471        updated = RecordingHandler()
472        failed = RecordingHandler()
473        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
474        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
475
476        self.web_socket_client._send_error = Exception("connection closed")
477        with self.assertRaises(Exception):
478            await subscription.update(events=OTHER_FILTERS)
479
480        self.assertEqual(updated.calls, [])
481        self.assertEqual(failed.calls, [])
482        self.assertIs(subscription.get_subscription_info(), confirmed)
483
484        pending_request = self.web_socket_client.sent_messages[1]
485        self.web_socket_client.trigger(
486            WebSocketEvents.receiveMessage, json.dumps(update_response(pending_request))
487        )
488        self.assertEqual(updated.calls, [])
489        self.assertEqual(failed.calls, [])
490        self.assertIs(subscription.get_subscription_info(), confirmed)
491
492        self.web_socket_client._responder = update_response
493        await subscription.update()
494
495        self.assertEqual(len(updated.calls), 1)
496        self.assertEqual(self.web_socket_client.sent_messages[2][1]["eventFilters"], EVENT_FILTERS)
497
498    async def test_overlapping_update_and_removal_while_update_pending_are_rejected_without_side_effects(self):
499        self.web_socket_client._responder = creation_response
500        subscription = WebSocketSubscription(self.web_socket_client)
501        await subscription.subscribe(events=EVENT_FILTERS)
502        confirmed = subscription.get_subscription_info()
503
504        updated = RecordingHandler()
505        removed = RecordingHandler()
506        failed = RecordingHandler()
507        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
508        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
509        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
510
511        self.web_socket_client._responder = None
512        await subscription.update(events=OTHER_FILTERS)
513        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
514
515        with self.assertRaises(Exception):
516            await subscription.update(events=OTHER_FILTERS)
517        with self.assertRaises(Exception):
518            await subscription.remove()
519
520        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
521        self.assertIs(subscription.get_subscription_info(), confirmed)
522        self.assertEqual(updated.calls, [])
523        self.assertEqual(removed.calls, [])
524
525        pending_request = self.web_socket_client.sent_messages[1]
526        self.web_socket_client.trigger(
527            WebSocketEvents.receiveMessage, json.dumps(update_response(pending_request, status=403))
528        )
529        self.assertEqual(len(failed.calls), 1)
530        self.assertEqual(updated.calls, [])
531        self.assertIs(subscription.get_subscription_info(), confirmed)
532
533        self.web_socket_client._responder = update_response
534        await subscription.update()
535
536        self.assertEqual(len(self.web_socket_client.sent_messages), 3)
537        self.assertEqual(self.web_socket_client.sent_messages[2][1]["eventFilters"], EVENT_FILTERS)
538
539    async def test_removal_send_without_response_preserves_state_listener_and_emits_nothing(self):
540        self.web_socket_client._responder = creation_response
541        subscription = WebSocketSubscription(self.web_socket_client)
542        await subscription.subscribe(events=EVENT_FILTERS)
543        confirmed = subscription.get_subscription_info()
544
545        removed = RecordingHandler()
546        failed = RecordingHandler()
547        notifications = RecordingHandler()
548        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
549        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, failed)
550        self.web_socket_client.on(WebSocketEvents.receiveSubscriptionNotification, notifications)
551
552        self.web_socket_client._responder = None
553        await subscription.remove()
554
555        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
556        self.assertEqual(self.web_socket_client.sent_messages[1][0]["method"], "DELETE")
557        self.assertEqual(removed.calls, [])
558        self.assertEqual(failed.calls, [])
559        self.assertIs(subscription.get_subscription_info(), confirmed)
560        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
561
562        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(server_notification()))
563        self.assertEqual(len(notifications.calls), 1)
564
565    async def test_confirmed_removal_clears_state_and_detaches_listener_before_single_event(self):
566        self.web_socket_client._responder = creation_response
567        subscription = WebSocketSubscription(self.web_socket_client)
568        await subscription.subscribe(events=EVENT_FILTERS)
569        notifications = RecordingHandler()
570        self.web_socket_client.on(WebSocketEvents.receiveSubscriptionNotification, notifications)
571
572        observed_at_event = []
573
574        def on_removed(*args):
575            observed_at_event.append(
576                (
577                    args,
578                    subscription.get_subscription_info(),
579                    self.web_socket_client.receive_message_listener_count,
580                )
581            )
582
583        removed = RecordingHandler()
584        failed = RecordingHandler()
585        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, on_removed)
586        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
587        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, failed)
588
589        self.web_socket_client._responder = removal_response
590        await subscription.remove()
591
592        self.assertEqual(len(removed.calls), 1)
593        self.assertEqual(removed.calls[0], ())
594        self.assertEqual(failed.calls, [])
595        self.assertEqual(observed_at_event, [((), None, 0)])
596        self.assertIsNone(subscription.get_subscription_info())
597        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
598
599        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(server_notification()))
600        self.assertEqual(len(notifications.calls), 0)
601
602    async def test_rejected_removal_response_emits_single_error_and_preserves_everything(self):
603        self.web_socket_client._responder = creation_response
604        subscription = WebSocketSubscription(self.web_socket_client)
605        await subscription.subscribe(events=EVENT_FILTERS)
606        confirmed = subscription.get_subscription_info()
607        notifications = RecordingHandler()
608        self.web_socket_client.on(WebSocketEvents.receiveSubscriptionNotification, notifications)
609
610        removed = RecordingHandler()
611        failed = RecordingHandler()
612        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
613        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, failed)
614
615        self.web_socket_client._responder = lambda request: removal_response(request, status=404)
616        await subscription.remove()
617
618        self.assertEqual(len(failed.calls), 1)
619        error = failed.calls[0][0]
620        self.assertIsInstance(error, Exception)
621        self.assertEqual(
622            str(error), "WebSocket subscription removal failed with status 404"
623        )
624        self.assertEqual(removed.calls, [])
625        self.assertIs(subscription.get_subscription_info(), confirmed)
626        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
627
628        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(server_notification()))
629        self.assertEqual(len(notifications.calls), 1)
630
631        self.web_socket_client._responder = removal_response
632        await subscription.remove()
633
634        self.assertEqual(len(removed.calls), 1)
635        self.assertIsNone(subscription.get_subscription_info())
636        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
637
638    async def test_unrelated_removal_response_is_ignored_and_pending_removal_still_completes(self):
639        def unrelated_removal_response(request):
640            response = removal_response(request)
641            response[0]["messageId"] = "unrelated-" + request[0]["messageId"]
642            return response
643
644        self.web_socket_client._responder = creation_response
645        subscription = WebSocketSubscription(self.web_socket_client)
646        await subscription.subscribe(events=EVENT_FILTERS)
647        confirmed = subscription.get_subscription_info()
648
649        removed = RecordingHandler()
650        failed = RecordingHandler()
651        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
652        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, failed)
653
654        self.web_socket_client._responder = unrelated_removal_response
655        await subscription.remove()
656
657        self.assertEqual(removed.calls, [])
658        self.assertEqual(failed.calls, [])
659        self.assertIs(subscription.get_subscription_info(), confirmed)
660        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
661
662        pending_request = self.web_socket_client.sent_messages[1]
663        self.web_socket_client.trigger(
664            WebSocketEvents.receiveMessage, json.dumps(removal_response(pending_request))
665        )
666
667        self.assertEqual(len(removed.calls), 1)
668        self.assertEqual(failed.calls, [])
669        self.assertIsNone(subscription.get_subscription_info())
670        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
671
672    async def test_removal_send_failure_preserves_state_and_listener_and_permits_retry(self):
673        self.web_socket_client._responder = creation_response
674        subscription = WebSocketSubscription(self.web_socket_client)
675        await subscription.subscribe(events=EVENT_FILTERS)
676        confirmed = subscription.get_subscription_info()
677
678        removed = RecordingHandler()
679        failed = RecordingHandler()
680        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
681        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, failed)
682
683        self.web_socket_client._send_error = Exception("connection closed")
684        with self.assertRaises(Exception):
685            await subscription.remove()
686
687        self.assertEqual(removed.calls, [])
688        self.assertEqual(failed.calls, [])
689        self.assertIs(subscription.get_subscription_info(), confirmed)
690        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
691
692        pending_request = self.web_socket_client.sent_messages[1]
693        self.web_socket_client.trigger(
694            WebSocketEvents.receiveMessage, json.dumps(removal_response(pending_request))
695        )
696        self.assertEqual(removed.calls, [])
697        self.assertEqual(failed.calls, [])
698        self.assertIs(subscription.get_subscription_info(), confirmed)
699        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
700
701        self.web_socket_client._responder = removal_response
702        await subscription.remove()
703
704        self.assertEqual(len(removed.calls), 1)
705        self.assertIsNone(subscription.get_subscription_info())
706        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
707
708    async def test_overlapping_removal_and_update_while_removal_pending_are_rejected_without_side_effects(self):
709        self.web_socket_client._responder = creation_response
710        subscription = WebSocketSubscription(self.web_socket_client)
711        await subscription.subscribe(events=EVENT_FILTERS)
712        confirmed = subscription.get_subscription_info()
713
714        removed = RecordingHandler()
715        failed = RecordingHandler()
716        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
717        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, failed)
718
719        self.web_socket_client._responder = None
720        await subscription.remove()
721        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
722
723        with self.assertRaises(Exception):
724            await subscription.remove()
725        with self.assertRaises(Exception):
726            await subscription.update(events=OTHER_FILTERS)
727
728        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
729        self.assertIs(subscription.get_subscription_info(), confirmed)
730        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
731        self.assertEqual(removed.calls, [])
732        self.assertEqual(failed.calls, [])
733
734        pending_request = self.web_socket_client.sent_messages[1]
735        self.web_socket_client.trigger(
736            WebSocketEvents.receiveMessage, json.dumps(removal_response(pending_request))
737        )
738        self.assertEqual(len(removed.calls), 1)
739        self.assertEqual(failed.calls, [])
740        self.assertIsNone(subscription.get_subscription_info())
741        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
742
743    async def test_duplicate_and_late_responses_after_completion_produce_no_events_or_changes(self):
744        self.web_socket_client._responder = creation_response
745        subscription = WebSocketSubscription(self.web_socket_client)
746        await subscription.subscribe(events=EVENT_FILTERS)
747        creation_request = self.web_socket_client.sent_messages[0]
748
749        updated = RecordingHandler()
750        removed = RecordingHandler()
751        failed = RecordingHandler()
752        remove_failed = RecordingHandler()
753        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
754        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
755        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
756        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, remove_failed)
757
758        self.web_socket_client._responder = update_response
759        await subscription.update(events=OTHER_FILTERS)
760        self.assertEqual(len(updated.calls), 1)
761        stored = subscription.get_subscription_info()
762
763        update_request = self.web_socket_client.sent_messages[1]
764        self.web_socket_client.trigger(
765            WebSocketEvents.receiveMessage, json.dumps(update_response(update_request))
766        )
767        self.web_socket_client.trigger(
768            WebSocketEvents.receiveMessage, json.dumps(creation_response(creation_request))
769        )
770        self.assertEqual(len(updated.calls), 1)
771        self.assertIs(subscription.get_subscription_info(), stored)
772        self.assertEqual(failed.calls, [])
773        self.assertEqual(removed.calls, [])
774
775        self.web_socket_client._responder = removal_response
776        await subscription.remove()
777        self.assertEqual(len(removed.calls), 1)
778        self.assertIsNone(subscription.get_subscription_info())
779
780        removal_request = self.web_socket_client.sent_messages[2]
781        self.web_socket_client.trigger(
782            WebSocketEvents.receiveMessage, json.dumps(removal_response(removal_request))
783        )
784        self.web_socket_client.trigger(
785            WebSocketEvents.receiveMessage, json.dumps(update_response(update_request))
786        )
787        self.assertEqual(len(removed.calls), 1)
788        self.assertEqual(len(updated.calls), 1)
789        self.assertEqual(failed.calls, [])
790        self.assertEqual(remove_failed.calls, [])
791        self.assertIsNone(subscription.get_subscription_info())
792        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
793
794    async def test_responses_correlated_by_message_id_regardless_of_type(self):
795        self.web_socket_client._responder = creation_response
796        subscription = WebSocketSubscription(self.web_socket_client)
797        await subscription.subscribe(events=EVENT_FILTERS)
798
799        updated = RecordingHandler()
800        removed = RecordingHandler()
801        failed = RecordingHandler()
802        remove_failed = RecordingHandler()
803        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
804        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
805        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
806        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, remove_failed)
807
808        self.web_socket_client._responder = lambda request: update_response(
809            request, message_type="ClientResponse"
810        )
811        await subscription.update(events=OTHER_FILTERS)
812
813        self.assertEqual(len(updated.calls), 1)
814        self.assertEqual(failed.calls, [])
815        self.assertIsNotNone(subscription.get_subscription_info())
816
817        self.web_socket_client._responder = lambda request: removal_response(
818            request, message_type="ClientResponse"
819        )
820        await subscription.remove()
821
822        self.assertEqual(len(removed.calls), 1)
823        self.assertEqual(failed.calls, [])
824        self.assertEqual(remove_failed.calls, [])
825        self.assertIsNone(subscription.get_subscription_info())
826
827
828if __name__ == "__main__":
829    unittest.main()
class RecordingHandler:
12class RecordingHandler:
13    def __init__(self):
14        self.calls = []
15
16    def __call__(self, *args):
17        self.calls.append(args)
calls
class FakeWebSocketClient(observable.core.Observable):
20class FakeWebSocketClient(Observable):
21    def __init__(self, responder=None, send_error=None):
22        Observable.__init__(self)
23        self.sent_messages = []
24        self.receive_message_listener_count = 0
25        self._responder = responder
26        self._send_error = send_error
27
28    def on(self, event, *handlers):
29        if event == WebSocketEvents.receiveMessage:
30            self.receive_message_listener_count += len(handlers)
31        return Observable.on(self, event, *handlers)
32
33    def off(self, event=None, *handlers):
34        if event == WebSocketEvents.receiveMessage:
35            self.receive_message_listener_count -= len(handlers)
36        return Observable.off(self, event, *handlers)
37
38    async def send_message(self, message):
39        self.sent_messages.append(message)
40        if self._send_error is not None:
41            error = self._send_error
42            self._send_error = None
43            raise error
44        if self._responder is not None:
45            response = self._responder(message)
46            if response is not None:
47                self.trigger(WebSocketEvents.receiveMessage, json.dumps(response))

Event system for python

FakeWebSocketClient(responder=None, send_error=None)
21    def __init__(self, responder=None, send_error=None):
22        Observable.__init__(self)
23        self.sent_messages = []
24        self.receive_message_listener_count = 0
25        self._responder = responder
26        self._send_error = send_error
sent_messages
receive_message_listener_count
def on(self, event, *handlers):
28    def on(self, event, *handlers):
29        if event == WebSocketEvents.receiveMessage:
30            self.receive_message_listener_count += len(handlers)
31        return Observable.on(self, event, *handlers)

Register a handler to a specified event

def off(self, event=None, *handlers):
33    def off(self, event=None, *handlers):
34        if event == WebSocketEvents.receiveMessage:
35            self.receive_message_listener_count -= len(handlers)
36        return Observable.off(self, event, *handlers)

Unregister an event or handler from an event

async def send_message(self, message):
38    async def send_message(self, message):
39        self.sent_messages.append(message)
40        if self._send_error is not None:
41            error = self._send_error
42            self._send_error = None
43            raise error
44        if self._responder is not None:
45            response = self._responder(message)
46            if response is not None:
47                self.trigger(WebSocketEvents.receiveMessage, json.dumps(response))
Inherited Members
observable.core.Observable
events
once
trigger
def creation_response(request):
50def creation_response(request):
51    return [
52        {
53            "type": "ClientRequest",
54            "messageId": request[0]["messageId"],
55            "status": 200,
56            "headers": {
57                "Server": "nginx",
58                "Date": "Wed, 20 Aug 2025 22:23:55 GMT",
59                "Content-Type": "application/json",
60                "RoutingKey": "SJC01P07",
61                "RCRequestId": "bedff5ae-9d68-4bc9-8653-7e34603ef562-2686696-1-19",
62            },
63        },
64        {
65            "uri": "/restapi/v1.0/subscription/1b2a2e6b-2245-4278-b47c-16259ca003a8",
66            "id": "1b2a2e6b-2245-4278-b47c-16259ca003a8",
67            "creationTime": "2025-08-20T22:23:55.169Z",
68            "status": "Active",
69            "eventFilters": ["/restapi/v1.0/account/809646016/extension/62264425016/presence"],
70            "expirationTime": "2025-08-21T22:23:55.169Z",
71            "expiresIn": 86399,
72            "deliveryMode": {"transportType": "WebSocket", "encryption": False},
73        },
74    ]
def rejected_creation_response(request, status=403):
77def rejected_creation_response(request, status=403):
78    response = creation_response(request)
79    response[0]["status"] = status
80    return response
def update_response(request, status=200, message_type='ClientRequest'):
 83def update_response(request, status=200, message_type="ClientRequest"):
 84    return [
 85        {
 86            "type": message_type,
 87            "messageId": request[0]["messageId"],
 88            "status": status,
 89            "headers": {
 90                "Server": "nginx",
 91                "Date": "Wed, 20 Aug 2025 22:23:55 GMT",
 92                "Content-Type": "application/json",
 93                "RoutingKey": "SJC01P07",
 94                "RCRequestId": "bedff5ae-9d68-4bc9-8653-7e34603ef562-2686696-1-19",
 95            },
 96        },
 97        {
 98            "uri": "/restapi/v1.0/subscription/9d3b7f10-5c11-4b6e-8a2f-6f8b56f0f1c2",
 99            "id": "9d3b7f10-5c11-4b6e-8a2f-6f8b56f0f1c2",
100            "creationTime": "2025-08-20T22:23:55.169Z",
101            "status": "Active",
102            "eventFilters": request[1]["eventFilters"],
103            "expirationTime": "2025-08-21T22:23:55.169Z",
104            "expiresIn": 86399,
105            "deliveryMode": {"transportType": "WebSocket", "encryption": False},
106        },
107    ]
def removal_response(request, status=200, message_type='ClientRequest'):
110def removal_response(request, status=200, message_type="ClientRequest"):
111    return [
112        {
113            "type": message_type,
114            "messageId": request[0]["messageId"],
115            "status": status,
116            "headers": {
117                "Server": "nginx",
118                "Date": "Wed, 20 Aug 2025 22:23:55 GMT",
119                "RCRequestId": "bedff5ae-9d68-4bc9-8653-7e34603ef562-2686696-1-19",
120            },
121        }
122    ]
EVENT_FILTERS = ['/restapi/v1.0/account/~/extension/~/presence']
OTHER_FILTERS = ['/restapi/v1.0/account/~/extension/~/message-store']
def server_notification():
129def server_notification():
130    return [
131        {
132            "type": "ServerNotification",
133            "messageId": str(uuid.uuid4()),
134            "headers": {"RoutingKey": "SJC01P07"},
135        },
136        {
137            "uri": "/restapi/v1.0/subscription/1b2a2e6b-2245-4278-b47c-16259ca003a8",
138            "event": {"/restapi/v1.0/account/~/extension/~/presence": {"activeCalls": []}},
139        },
140    ]
class WebSocketSubscriptionTest(unittest.async_case.IsolatedAsyncioTestCase):
143class WebSocketSubscriptionTest(unittest.IsolatedAsyncioTestCase):
144    def setUp(self):
145        self.web_socket_client = FakeWebSocketClient()
146
147    async def test_successful_creation_before_send_returns_stores_response_and_emits_event_once(self):
148        self.web_socket_client._responder = creation_response
149        subscription = WebSocketSubscription(self.web_socket_client)
150        created = RecordingHandler()
151        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
152
153        await subscription.subscribe(events=EVENT_FILTERS)
154
155        self.assertEqual(len(created.calls), 1)
156        self.assertIs(created.calls[0][0], subscription)
157        self.assertEqual(len(self.web_socket_client.sent_messages), 1)
158        stored = subscription.get_subscription_info()
159        self.assertIsNotNone(stored)
160        self.assertEqual(stored[0]["type"], "ClientRequest")
161        self.assertEqual(
162            stored[0]["messageId"], self.web_socket_client.sent_messages[0][0]["messageId"]
163        )
164        self.assertNotIn("WSG-SubscriptionId", stored[0]["headers"])
165        self.assertEqual(
166            stored[1]["id"], "1b2a2e6b-2245-4278-b47c-16259ca003a8"
167        )
168
169    async def test_unrelated_message_id_response_does_not_change_state_or_emit_events(self):
170        def unrelated_response(request):
171            response = creation_response(request)
172            response[0]["messageId"] = "unrelated-" + request[0]["messageId"]
173            return response
174
175        self.web_socket_client._responder = unrelated_response
176        subscription = WebSocketSubscription(self.web_socket_client)
177        created = RecordingHandler()
178        failed = RecordingHandler()
179        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
180        self.web_socket_client.on(WebSocketEvents.createSubscriptionError, failed)
181
182        await subscription.subscribe(events=EVENT_FILTERS)
183
184        self.assertIsNone(subscription.get_subscription_info())
185        self.assertEqual(created.calls, [])
186        self.assertEqual(failed.calls, [])
187
188    async def test_listener_remains_after_creation_so_notifications_are_emitted(self):
189        self.web_socket_client._responder = creation_response
190        subscription = WebSocketSubscription(self.web_socket_client)
191        created = RecordingHandler()
192        notifications = RecordingHandler()
193        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
194        self.web_socket_client.on(WebSocketEvents.receiveSubscriptionNotification, notifications)
195
196        await subscription.subscribe(events=EVENT_FILTERS)
197        self.assertEqual(len(created.calls), 1)
198
199        notification = server_notification()
200        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(notification))
201
202        self.assertEqual(len(notifications.calls), 1)
203        self.assertEqual(notifications.calls[0][0], notification)
204
205    async def test_rejected_creation_emits_single_error_without_state_or_success(self):
206        self.web_socket_client._responder = rejected_creation_response
207        subscription = WebSocketSubscription(self.web_socket_client)
208        created = RecordingHandler()
209        failed = RecordingHandler()
210        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
211        self.web_socket_client.on(WebSocketEvents.createSubscriptionError, failed)
212
213        await subscription.subscribe(events=EVENT_FILTERS)
214
215        self.assertEqual(len(failed.calls), 1)
216        error = failed.calls[0][0]
217        self.assertIsInstance(error, Exception)
218        self.assertEqual(
219            str(error), "WebSocket subscription creation failed with status 403"
220        )
221        self.assertEqual(created.calls, [])
222        self.assertIsNone(subscription.get_subscription_info())
223
224    async def test_retry_after_rejected_creation_sends_new_request_and_can_succeed(self):
225        self.web_socket_client._responder = rejected_creation_response
226        subscription = WebSocketSubscription(self.web_socket_client)
227        created = RecordingHandler()
228        failed = RecordingHandler()
229        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
230        self.web_socket_client.on(WebSocketEvents.createSubscriptionError, failed)
231
232        await subscription.subscribe(events=EVENT_FILTERS)
233        self.assertEqual(len(failed.calls), 1)
234
235        self.web_socket_client._responder = creation_response
236        await subscription.subscribe(events=EVENT_FILTERS)
237
238        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
239        self.assertEqual(len(failed.calls), 1)
240        self.assertEqual(len(created.calls), 1)
241        self.assertIsNotNone(subscription.get_subscription_info())
242
243    async def test_second_creation_while_pending_is_rejected_without_new_request_or_listener(self):
244        self.web_socket_client._responder = None
245        subscription = WebSocketSubscription(self.web_socket_client)
246        created = RecordingHandler()
247        failed = RecordingHandler()
248        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
249        self.web_socket_client.on(WebSocketEvents.createSubscriptionError, failed)
250
251        await subscription.subscribe(events=EVENT_FILTERS)
252        self.assertEqual(len(self.web_socket_client.sent_messages), 1)
253        listeners_before = self.web_socket_client.receive_message_listener_count
254
255        with self.assertRaises(Exception):
256            await subscription.subscribe(events=EVENT_FILTERS)
257
258        self.assertEqual(len(self.web_socket_client.sent_messages), 1)
259        self.assertEqual(
260            self.web_socket_client.receive_message_listener_count, listeners_before
261        )
262
263        self.web_socket_client._responder = creation_response
264        pending_request = self.web_socket_client.sent_messages[0]
265        self.web_socket_client.trigger(
266            WebSocketEvents.receiveMessage, json.dumps(creation_response(pending_request))
267        )
268
269        self.assertEqual(len(created.calls), 1)
270        self.assertEqual(failed.calls, [])
271        self.assertIsNotNone(subscription.get_subscription_info())
272
273    async def test_second_creation_while_pending_leaves_original_filters_for_update(self):
274        self.web_socket_client._responder = None
275        subscription = WebSocketSubscription(self.web_socket_client)
276        created = RecordingHandler()
277        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
278
279        await subscription.subscribe(events=EVENT_FILTERS)
280        self.assertEqual(len(self.web_socket_client.sent_messages), 1)
281
282        with self.assertRaises(Exception):
283            await subscription.subscribe(events=OTHER_FILTERS)
284
285        pending_request = self.web_socket_client.sent_messages[0]
286        self.web_socket_client._responder = creation_response
287        self.web_socket_client.trigger(
288            WebSocketEvents.receiveMessage, json.dumps(creation_response(pending_request))
289        )
290        self.assertEqual(len(created.calls), 1)
291
292        await subscription.update()
293
294        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
295        self.assertEqual(
296            self.web_socket_client.sent_messages[1][1]["eventFilters"], EVENT_FILTERS
297        )
298
299    async def test_failed_send_clears_pending_and_new_listener_and_retry_succeeds(self):
300        self.web_socket_client._responder = creation_response
301        self.web_socket_client._send_error = Exception("connection closed")
302        subscription = WebSocketSubscription(self.web_socket_client)
303        created = RecordingHandler()
304        notifications = RecordingHandler()
305        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
306        self.web_socket_client.on(WebSocketEvents.receiveSubscriptionNotification, notifications)
307
308        with self.assertRaises(Exception):
309            await subscription.subscribe(events=EVENT_FILTERS)
310
311        self.assertEqual(len(self.web_socket_client.sent_messages), 1)
312        self.assertIsNone(subscription.get_subscription_info())
313        self.assertEqual(created.calls, [])
314        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
315
316        await subscription.subscribe(events=EVENT_FILTERS)
317
318        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
319        self.assertEqual(len(created.calls), 1)
320        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
321
322        notification = server_notification()
323        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(notification))
324        self.assertEqual(len(notifications.calls), 1)
325
326    async def test_removal_detaches_listener_and_later_retry_receives_notifications_once(self):
327        self.web_socket_client._responder = creation_response
328        subscription = WebSocketSubscription(self.web_socket_client)
329        notifications = RecordingHandler()
330        self.web_socket_client.on(WebSocketEvents.receiveSubscriptionNotification, notifications)
331
332        await subscription.subscribe(events=EVENT_FILTERS)
333        notification = server_notification()
334        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(notification))
335        self.assertEqual(len(notifications.calls), 1)
336
337        await subscription.remove()
338        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(server_notification()))
339        self.assertEqual(len(notifications.calls), 1)
340        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
341
342        await subscription.subscribe(events=EVENT_FILTERS)
343        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(server_notification()))
344        self.assertEqual(len(notifications.calls), 2)
345        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
346
347    async def test_update_send_without_response_emits_nothing_and_preserves_confirmed_state(self):
348        self.web_socket_client._responder = creation_response
349        subscription = WebSocketSubscription(self.web_socket_client)
350        await subscription.subscribe(events=EVENT_FILTERS)
351        confirmed = subscription.get_subscription_info()
352
353        updated = RecordingHandler()
354        failed = RecordingHandler()
355        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
356        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
357
358        self.web_socket_client._responder = None
359        await subscription.update(events=OTHER_FILTERS)
360
361        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
362        self.assertEqual(self.web_socket_client.sent_messages[1][0]["method"], "PUT")
363        self.assertEqual(self.web_socket_client.sent_messages[1][1]["eventFilters"], OTHER_FILTERS)
364        self.assertEqual(updated.calls, [])
365        self.assertEqual(failed.calls, [])
366        self.assertIs(subscription.get_subscription_info(), confirmed)
367        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
368
369    async def test_confirmed_update_response_commits_envelope_and_filters_before_single_event(self):
370        self.web_socket_client._responder = creation_response
371        subscription = WebSocketSubscription(self.web_socket_client)
372        await subscription.subscribe(events=EVENT_FILTERS)
373
374        observed_at_event = []
375
376        def on_updated(updated_subscription):
377            observed_at_event.append(updated_subscription.get_subscription_info())
378
379        updated = RecordingHandler()
380        failed = RecordingHandler()
381        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, on_updated)
382        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
383        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
384
385        self.web_socket_client._responder = update_response
386        await subscription.update(events=OTHER_FILTERS)
387
388        self.assertEqual(len(updated.calls), 1)
389        self.assertIs(updated.calls[0][0], subscription)
390        self.assertEqual(failed.calls, [])
391        stored = subscription.get_subscription_info()
392        self.assertEqual(
393            stored[0]["messageId"], self.web_socket_client.sent_messages[1][0]["messageId"]
394        )
395        self.assertEqual(stored[1]["id"], "9d3b7f10-5c11-4b6e-8a2f-6f8b56f0f1c2")
396        self.assertEqual(stored[1]["eventFilters"], OTHER_FILTERS)
397        self.assertEqual(observed_at_event, [stored])
398
399        self.web_socket_client._responder = None
400        await subscription.update()
401        self.assertEqual(self.web_socket_client.sent_messages[2][1]["eventFilters"], OTHER_FILTERS)
402
403    async def test_rejected_update_response_emits_single_error_preserves_state_and_permits_retry(self):
404        self.web_socket_client._responder = creation_response
405        subscription = WebSocketSubscription(self.web_socket_client)
406        await subscription.subscribe(events=EVENT_FILTERS)
407        confirmed = subscription.get_subscription_info()
408
409        updated = RecordingHandler()
410        failed = RecordingHandler()
411        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
412        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
413
414        self.web_socket_client._responder = lambda request: update_response(request, status=404)
415        await subscription.update(events=OTHER_FILTERS)
416
417        self.assertEqual(len(failed.calls), 1)
418        error = failed.calls[0][0]
419        self.assertIsInstance(error, Exception)
420        self.assertEqual(
421            str(error), "WebSocket subscription update failed with status 404"
422        )
423        self.assertEqual(updated.calls, [])
424        self.assertIs(subscription.get_subscription_info(), confirmed)
425
426        self.web_socket_client._responder = update_response
427        await subscription.update()
428
429        self.assertEqual(len(updated.calls), 1)
430        self.assertEqual(self.web_socket_client.sent_messages[2][1]["eventFilters"], EVENT_FILTERS)
431
432    async def test_unrelated_update_response_is_ignored_and_pending_update_still_completes(self):
433        def unrelated_update_response(request):
434            response = update_response(request)
435            response[0]["messageId"] = "unrelated-" + request[0]["messageId"]
436            return response
437
438        self.web_socket_client._responder = creation_response
439        subscription = WebSocketSubscription(self.web_socket_client)
440        await subscription.subscribe(events=EVENT_FILTERS)
441        confirmed = subscription.get_subscription_info()
442
443        updated = RecordingHandler()
444        failed = RecordingHandler()
445        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
446        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
447
448        self.web_socket_client._responder = unrelated_update_response
449        await subscription.update(events=OTHER_FILTERS)
450
451        self.assertEqual(updated.calls, [])
452        self.assertEqual(failed.calls, [])
453        self.assertIs(subscription.get_subscription_info(), confirmed)
454
455        pending_request = self.web_socket_client.sent_messages[1]
456        self.web_socket_client.trigger(
457            WebSocketEvents.receiveMessage, json.dumps(update_response(pending_request))
458        )
459
460        self.assertEqual(len(updated.calls), 1)
461        self.assertEqual(failed.calls, [])
462        self.assertEqual(
463            subscription.get_subscription_info()[0]["messageId"], pending_request[0]["messageId"]
464        )
465
466    async def test_update_send_failure_preserves_state_and_permits_retry_after_late_response(self):
467        self.web_socket_client._responder = creation_response
468        subscription = WebSocketSubscription(self.web_socket_client)
469        await subscription.subscribe(events=EVENT_FILTERS)
470        confirmed = subscription.get_subscription_info()
471
472        updated = RecordingHandler()
473        failed = RecordingHandler()
474        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
475        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
476
477        self.web_socket_client._send_error = Exception("connection closed")
478        with self.assertRaises(Exception):
479            await subscription.update(events=OTHER_FILTERS)
480
481        self.assertEqual(updated.calls, [])
482        self.assertEqual(failed.calls, [])
483        self.assertIs(subscription.get_subscription_info(), confirmed)
484
485        pending_request = self.web_socket_client.sent_messages[1]
486        self.web_socket_client.trigger(
487            WebSocketEvents.receiveMessage, json.dumps(update_response(pending_request))
488        )
489        self.assertEqual(updated.calls, [])
490        self.assertEqual(failed.calls, [])
491        self.assertIs(subscription.get_subscription_info(), confirmed)
492
493        self.web_socket_client._responder = update_response
494        await subscription.update()
495
496        self.assertEqual(len(updated.calls), 1)
497        self.assertEqual(self.web_socket_client.sent_messages[2][1]["eventFilters"], EVENT_FILTERS)
498
499    async def test_overlapping_update_and_removal_while_update_pending_are_rejected_without_side_effects(self):
500        self.web_socket_client._responder = creation_response
501        subscription = WebSocketSubscription(self.web_socket_client)
502        await subscription.subscribe(events=EVENT_FILTERS)
503        confirmed = subscription.get_subscription_info()
504
505        updated = RecordingHandler()
506        removed = RecordingHandler()
507        failed = RecordingHandler()
508        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
509        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
510        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
511
512        self.web_socket_client._responder = None
513        await subscription.update(events=OTHER_FILTERS)
514        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
515
516        with self.assertRaises(Exception):
517            await subscription.update(events=OTHER_FILTERS)
518        with self.assertRaises(Exception):
519            await subscription.remove()
520
521        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
522        self.assertIs(subscription.get_subscription_info(), confirmed)
523        self.assertEqual(updated.calls, [])
524        self.assertEqual(removed.calls, [])
525
526        pending_request = self.web_socket_client.sent_messages[1]
527        self.web_socket_client.trigger(
528            WebSocketEvents.receiveMessage, json.dumps(update_response(pending_request, status=403))
529        )
530        self.assertEqual(len(failed.calls), 1)
531        self.assertEqual(updated.calls, [])
532        self.assertIs(subscription.get_subscription_info(), confirmed)
533
534        self.web_socket_client._responder = update_response
535        await subscription.update()
536
537        self.assertEqual(len(self.web_socket_client.sent_messages), 3)
538        self.assertEqual(self.web_socket_client.sent_messages[2][1]["eventFilters"], EVENT_FILTERS)
539
540    async def test_removal_send_without_response_preserves_state_listener_and_emits_nothing(self):
541        self.web_socket_client._responder = creation_response
542        subscription = WebSocketSubscription(self.web_socket_client)
543        await subscription.subscribe(events=EVENT_FILTERS)
544        confirmed = subscription.get_subscription_info()
545
546        removed = RecordingHandler()
547        failed = RecordingHandler()
548        notifications = RecordingHandler()
549        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
550        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, failed)
551        self.web_socket_client.on(WebSocketEvents.receiveSubscriptionNotification, notifications)
552
553        self.web_socket_client._responder = None
554        await subscription.remove()
555
556        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
557        self.assertEqual(self.web_socket_client.sent_messages[1][0]["method"], "DELETE")
558        self.assertEqual(removed.calls, [])
559        self.assertEqual(failed.calls, [])
560        self.assertIs(subscription.get_subscription_info(), confirmed)
561        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
562
563        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(server_notification()))
564        self.assertEqual(len(notifications.calls), 1)
565
566    async def test_confirmed_removal_clears_state_and_detaches_listener_before_single_event(self):
567        self.web_socket_client._responder = creation_response
568        subscription = WebSocketSubscription(self.web_socket_client)
569        await subscription.subscribe(events=EVENT_FILTERS)
570        notifications = RecordingHandler()
571        self.web_socket_client.on(WebSocketEvents.receiveSubscriptionNotification, notifications)
572
573        observed_at_event = []
574
575        def on_removed(*args):
576            observed_at_event.append(
577                (
578                    args,
579                    subscription.get_subscription_info(),
580                    self.web_socket_client.receive_message_listener_count,
581                )
582            )
583
584        removed = RecordingHandler()
585        failed = RecordingHandler()
586        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, on_removed)
587        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
588        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, failed)
589
590        self.web_socket_client._responder = removal_response
591        await subscription.remove()
592
593        self.assertEqual(len(removed.calls), 1)
594        self.assertEqual(removed.calls[0], ())
595        self.assertEqual(failed.calls, [])
596        self.assertEqual(observed_at_event, [((), None, 0)])
597        self.assertIsNone(subscription.get_subscription_info())
598        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
599
600        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(server_notification()))
601        self.assertEqual(len(notifications.calls), 0)
602
603    async def test_rejected_removal_response_emits_single_error_and_preserves_everything(self):
604        self.web_socket_client._responder = creation_response
605        subscription = WebSocketSubscription(self.web_socket_client)
606        await subscription.subscribe(events=EVENT_FILTERS)
607        confirmed = subscription.get_subscription_info()
608        notifications = RecordingHandler()
609        self.web_socket_client.on(WebSocketEvents.receiveSubscriptionNotification, notifications)
610
611        removed = RecordingHandler()
612        failed = RecordingHandler()
613        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
614        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, failed)
615
616        self.web_socket_client._responder = lambda request: removal_response(request, status=404)
617        await subscription.remove()
618
619        self.assertEqual(len(failed.calls), 1)
620        error = failed.calls[0][0]
621        self.assertIsInstance(error, Exception)
622        self.assertEqual(
623            str(error), "WebSocket subscription removal failed with status 404"
624        )
625        self.assertEqual(removed.calls, [])
626        self.assertIs(subscription.get_subscription_info(), confirmed)
627        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
628
629        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(server_notification()))
630        self.assertEqual(len(notifications.calls), 1)
631
632        self.web_socket_client._responder = removal_response
633        await subscription.remove()
634
635        self.assertEqual(len(removed.calls), 1)
636        self.assertIsNone(subscription.get_subscription_info())
637        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
638
639    async def test_unrelated_removal_response_is_ignored_and_pending_removal_still_completes(self):
640        def unrelated_removal_response(request):
641            response = removal_response(request)
642            response[0]["messageId"] = "unrelated-" + request[0]["messageId"]
643            return response
644
645        self.web_socket_client._responder = creation_response
646        subscription = WebSocketSubscription(self.web_socket_client)
647        await subscription.subscribe(events=EVENT_FILTERS)
648        confirmed = subscription.get_subscription_info()
649
650        removed = RecordingHandler()
651        failed = RecordingHandler()
652        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
653        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, failed)
654
655        self.web_socket_client._responder = unrelated_removal_response
656        await subscription.remove()
657
658        self.assertEqual(removed.calls, [])
659        self.assertEqual(failed.calls, [])
660        self.assertIs(subscription.get_subscription_info(), confirmed)
661        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
662
663        pending_request = self.web_socket_client.sent_messages[1]
664        self.web_socket_client.trigger(
665            WebSocketEvents.receiveMessage, json.dumps(removal_response(pending_request))
666        )
667
668        self.assertEqual(len(removed.calls), 1)
669        self.assertEqual(failed.calls, [])
670        self.assertIsNone(subscription.get_subscription_info())
671        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
672
673    async def test_removal_send_failure_preserves_state_and_listener_and_permits_retry(self):
674        self.web_socket_client._responder = creation_response
675        subscription = WebSocketSubscription(self.web_socket_client)
676        await subscription.subscribe(events=EVENT_FILTERS)
677        confirmed = subscription.get_subscription_info()
678
679        removed = RecordingHandler()
680        failed = RecordingHandler()
681        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
682        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, failed)
683
684        self.web_socket_client._send_error = Exception("connection closed")
685        with self.assertRaises(Exception):
686            await subscription.remove()
687
688        self.assertEqual(removed.calls, [])
689        self.assertEqual(failed.calls, [])
690        self.assertIs(subscription.get_subscription_info(), confirmed)
691        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
692
693        pending_request = self.web_socket_client.sent_messages[1]
694        self.web_socket_client.trigger(
695            WebSocketEvents.receiveMessage, json.dumps(removal_response(pending_request))
696        )
697        self.assertEqual(removed.calls, [])
698        self.assertEqual(failed.calls, [])
699        self.assertIs(subscription.get_subscription_info(), confirmed)
700        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
701
702        self.web_socket_client._responder = removal_response
703        await subscription.remove()
704
705        self.assertEqual(len(removed.calls), 1)
706        self.assertIsNone(subscription.get_subscription_info())
707        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
708
709    async def test_overlapping_removal_and_update_while_removal_pending_are_rejected_without_side_effects(self):
710        self.web_socket_client._responder = creation_response
711        subscription = WebSocketSubscription(self.web_socket_client)
712        await subscription.subscribe(events=EVENT_FILTERS)
713        confirmed = subscription.get_subscription_info()
714
715        removed = RecordingHandler()
716        failed = RecordingHandler()
717        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
718        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, failed)
719
720        self.web_socket_client._responder = None
721        await subscription.remove()
722        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
723
724        with self.assertRaises(Exception):
725            await subscription.remove()
726        with self.assertRaises(Exception):
727            await subscription.update(events=OTHER_FILTERS)
728
729        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
730        self.assertIs(subscription.get_subscription_info(), confirmed)
731        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
732        self.assertEqual(removed.calls, [])
733        self.assertEqual(failed.calls, [])
734
735        pending_request = self.web_socket_client.sent_messages[1]
736        self.web_socket_client.trigger(
737            WebSocketEvents.receiveMessage, json.dumps(removal_response(pending_request))
738        )
739        self.assertEqual(len(removed.calls), 1)
740        self.assertEqual(failed.calls, [])
741        self.assertIsNone(subscription.get_subscription_info())
742        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
743
744    async def test_duplicate_and_late_responses_after_completion_produce_no_events_or_changes(self):
745        self.web_socket_client._responder = creation_response
746        subscription = WebSocketSubscription(self.web_socket_client)
747        await subscription.subscribe(events=EVENT_FILTERS)
748        creation_request = self.web_socket_client.sent_messages[0]
749
750        updated = RecordingHandler()
751        removed = RecordingHandler()
752        failed = RecordingHandler()
753        remove_failed = RecordingHandler()
754        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
755        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
756        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
757        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, remove_failed)
758
759        self.web_socket_client._responder = update_response
760        await subscription.update(events=OTHER_FILTERS)
761        self.assertEqual(len(updated.calls), 1)
762        stored = subscription.get_subscription_info()
763
764        update_request = self.web_socket_client.sent_messages[1]
765        self.web_socket_client.trigger(
766            WebSocketEvents.receiveMessage, json.dumps(update_response(update_request))
767        )
768        self.web_socket_client.trigger(
769            WebSocketEvents.receiveMessage, json.dumps(creation_response(creation_request))
770        )
771        self.assertEqual(len(updated.calls), 1)
772        self.assertIs(subscription.get_subscription_info(), stored)
773        self.assertEqual(failed.calls, [])
774        self.assertEqual(removed.calls, [])
775
776        self.web_socket_client._responder = removal_response
777        await subscription.remove()
778        self.assertEqual(len(removed.calls), 1)
779        self.assertIsNone(subscription.get_subscription_info())
780
781        removal_request = self.web_socket_client.sent_messages[2]
782        self.web_socket_client.trigger(
783            WebSocketEvents.receiveMessage, json.dumps(removal_response(removal_request))
784        )
785        self.web_socket_client.trigger(
786            WebSocketEvents.receiveMessage, json.dumps(update_response(update_request))
787        )
788        self.assertEqual(len(removed.calls), 1)
789        self.assertEqual(len(updated.calls), 1)
790        self.assertEqual(failed.calls, [])
791        self.assertEqual(remove_failed.calls, [])
792        self.assertIsNone(subscription.get_subscription_info())
793        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
794
795    async def test_responses_correlated_by_message_id_regardless_of_type(self):
796        self.web_socket_client._responder = creation_response
797        subscription = WebSocketSubscription(self.web_socket_client)
798        await subscription.subscribe(events=EVENT_FILTERS)
799
800        updated = RecordingHandler()
801        removed = RecordingHandler()
802        failed = RecordingHandler()
803        remove_failed = RecordingHandler()
804        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
805        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
806        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
807        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, remove_failed)
808
809        self.web_socket_client._responder = lambda request: update_response(
810            request, message_type="ClientResponse"
811        )
812        await subscription.update(events=OTHER_FILTERS)
813
814        self.assertEqual(len(updated.calls), 1)
815        self.assertEqual(failed.calls, [])
816        self.assertIsNotNone(subscription.get_subscription_info())
817
818        self.web_socket_client._responder = lambda request: removal_response(
819            request, message_type="ClientResponse"
820        )
821        await subscription.remove()
822
823        self.assertEqual(len(removed.calls), 1)
824        self.assertEqual(failed.calls, [])
825        self.assertEqual(remove_failed.calls, [])
826        self.assertIsNone(subscription.get_subscription_info())

A class whose instances are single test cases.

By default, the test code itself should be placed in a method named 'runTest'.

If the fixture may be used for many test cases, create as many test methods as are needed. When instantiating such a TestCase subclass, specify in the constructor arguments the name of the test method that the instance is to execute.

Test authors should subclass TestCase for their own tests. Construction and deconstruction of the test's environment ('fixture') can be implemented by overriding the 'setUp' and 'tearDown' methods respectively.

If it is necessary to override the __init__ method, the base class __init__ method must always be called. It is important that subclasses should not change the signature of their __init__ method, since instances of the classes are instantiated automatically by parts of the framework in order to be run.

When subclassing TestCase, you can set these attributes:

  • failureException: determines which exception will be raised when the instance's assertion methods fail; test methods raising this exception will be deemed to have 'failed' rather than 'errored'.
  • longMessage: determines whether long messages (including repr of objects used in assert methods) will be printed on failure in addition to any explicit message passed.
  • maxDiff: sets the maximum length of a diff in failure messages by assert methods using difflib. It is looked up as an instance attribute so can be configured by individual tests if required.
def setUp(self):
144    def setUp(self):
145        self.web_socket_client = FakeWebSocketClient()

Hook method for setting up the test fixture before exercising it.

async def test_successful_creation_before_send_returns_stores_response_and_emits_event_once(self):
147    async def test_successful_creation_before_send_returns_stores_response_and_emits_event_once(self):
148        self.web_socket_client._responder = creation_response
149        subscription = WebSocketSubscription(self.web_socket_client)
150        created = RecordingHandler()
151        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
152
153        await subscription.subscribe(events=EVENT_FILTERS)
154
155        self.assertEqual(len(created.calls), 1)
156        self.assertIs(created.calls[0][0], subscription)
157        self.assertEqual(len(self.web_socket_client.sent_messages), 1)
158        stored = subscription.get_subscription_info()
159        self.assertIsNotNone(stored)
160        self.assertEqual(stored[0]["type"], "ClientRequest")
161        self.assertEqual(
162            stored[0]["messageId"], self.web_socket_client.sent_messages[0][0]["messageId"]
163        )
164        self.assertNotIn("WSG-SubscriptionId", stored[0]["headers"])
165        self.assertEqual(
166            stored[1]["id"], "1b2a2e6b-2245-4278-b47c-16259ca003a8"
167        )
async def test_unrelated_message_id_response_does_not_change_state_or_emit_events(self):
169    async def test_unrelated_message_id_response_does_not_change_state_or_emit_events(self):
170        def unrelated_response(request):
171            response = creation_response(request)
172            response[0]["messageId"] = "unrelated-" + request[0]["messageId"]
173            return response
174
175        self.web_socket_client._responder = unrelated_response
176        subscription = WebSocketSubscription(self.web_socket_client)
177        created = RecordingHandler()
178        failed = RecordingHandler()
179        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
180        self.web_socket_client.on(WebSocketEvents.createSubscriptionError, failed)
181
182        await subscription.subscribe(events=EVENT_FILTERS)
183
184        self.assertIsNone(subscription.get_subscription_info())
185        self.assertEqual(created.calls, [])
186        self.assertEqual(failed.calls, [])
async def test_listener_remains_after_creation_so_notifications_are_emitted(self):
188    async def test_listener_remains_after_creation_so_notifications_are_emitted(self):
189        self.web_socket_client._responder = creation_response
190        subscription = WebSocketSubscription(self.web_socket_client)
191        created = RecordingHandler()
192        notifications = RecordingHandler()
193        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
194        self.web_socket_client.on(WebSocketEvents.receiveSubscriptionNotification, notifications)
195
196        await subscription.subscribe(events=EVENT_FILTERS)
197        self.assertEqual(len(created.calls), 1)
198
199        notification = server_notification()
200        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(notification))
201
202        self.assertEqual(len(notifications.calls), 1)
203        self.assertEqual(notifications.calls[0][0], notification)
async def test_rejected_creation_emits_single_error_without_state_or_success(self):
205    async def test_rejected_creation_emits_single_error_without_state_or_success(self):
206        self.web_socket_client._responder = rejected_creation_response
207        subscription = WebSocketSubscription(self.web_socket_client)
208        created = RecordingHandler()
209        failed = RecordingHandler()
210        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
211        self.web_socket_client.on(WebSocketEvents.createSubscriptionError, failed)
212
213        await subscription.subscribe(events=EVENT_FILTERS)
214
215        self.assertEqual(len(failed.calls), 1)
216        error = failed.calls[0][0]
217        self.assertIsInstance(error, Exception)
218        self.assertEqual(
219            str(error), "WebSocket subscription creation failed with status 403"
220        )
221        self.assertEqual(created.calls, [])
222        self.assertIsNone(subscription.get_subscription_info())
async def test_retry_after_rejected_creation_sends_new_request_and_can_succeed(self):
224    async def test_retry_after_rejected_creation_sends_new_request_and_can_succeed(self):
225        self.web_socket_client._responder = rejected_creation_response
226        subscription = WebSocketSubscription(self.web_socket_client)
227        created = RecordingHandler()
228        failed = RecordingHandler()
229        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
230        self.web_socket_client.on(WebSocketEvents.createSubscriptionError, failed)
231
232        await subscription.subscribe(events=EVENT_FILTERS)
233        self.assertEqual(len(failed.calls), 1)
234
235        self.web_socket_client._responder = creation_response
236        await subscription.subscribe(events=EVENT_FILTERS)
237
238        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
239        self.assertEqual(len(failed.calls), 1)
240        self.assertEqual(len(created.calls), 1)
241        self.assertIsNotNone(subscription.get_subscription_info())
async def test_second_creation_while_pending_is_rejected_without_new_request_or_listener(self):
243    async def test_second_creation_while_pending_is_rejected_without_new_request_or_listener(self):
244        self.web_socket_client._responder = None
245        subscription = WebSocketSubscription(self.web_socket_client)
246        created = RecordingHandler()
247        failed = RecordingHandler()
248        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
249        self.web_socket_client.on(WebSocketEvents.createSubscriptionError, failed)
250
251        await subscription.subscribe(events=EVENT_FILTERS)
252        self.assertEqual(len(self.web_socket_client.sent_messages), 1)
253        listeners_before = self.web_socket_client.receive_message_listener_count
254
255        with self.assertRaises(Exception):
256            await subscription.subscribe(events=EVENT_FILTERS)
257
258        self.assertEqual(len(self.web_socket_client.sent_messages), 1)
259        self.assertEqual(
260            self.web_socket_client.receive_message_listener_count, listeners_before
261        )
262
263        self.web_socket_client._responder = creation_response
264        pending_request = self.web_socket_client.sent_messages[0]
265        self.web_socket_client.trigger(
266            WebSocketEvents.receiveMessage, json.dumps(creation_response(pending_request))
267        )
268
269        self.assertEqual(len(created.calls), 1)
270        self.assertEqual(failed.calls, [])
271        self.assertIsNotNone(subscription.get_subscription_info())
async def test_second_creation_while_pending_leaves_original_filters_for_update(self):
273    async def test_second_creation_while_pending_leaves_original_filters_for_update(self):
274        self.web_socket_client._responder = None
275        subscription = WebSocketSubscription(self.web_socket_client)
276        created = RecordingHandler()
277        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
278
279        await subscription.subscribe(events=EVENT_FILTERS)
280        self.assertEqual(len(self.web_socket_client.sent_messages), 1)
281
282        with self.assertRaises(Exception):
283            await subscription.subscribe(events=OTHER_FILTERS)
284
285        pending_request = self.web_socket_client.sent_messages[0]
286        self.web_socket_client._responder = creation_response
287        self.web_socket_client.trigger(
288            WebSocketEvents.receiveMessage, json.dumps(creation_response(pending_request))
289        )
290        self.assertEqual(len(created.calls), 1)
291
292        await subscription.update()
293
294        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
295        self.assertEqual(
296            self.web_socket_client.sent_messages[1][1]["eventFilters"], EVENT_FILTERS
297        )
async def test_failed_send_clears_pending_and_new_listener_and_retry_succeeds(self):
299    async def test_failed_send_clears_pending_and_new_listener_and_retry_succeeds(self):
300        self.web_socket_client._responder = creation_response
301        self.web_socket_client._send_error = Exception("connection closed")
302        subscription = WebSocketSubscription(self.web_socket_client)
303        created = RecordingHandler()
304        notifications = RecordingHandler()
305        self.web_socket_client.on(WebSocketEvents.subscriptionCreated, created)
306        self.web_socket_client.on(WebSocketEvents.receiveSubscriptionNotification, notifications)
307
308        with self.assertRaises(Exception):
309            await subscription.subscribe(events=EVENT_FILTERS)
310
311        self.assertEqual(len(self.web_socket_client.sent_messages), 1)
312        self.assertIsNone(subscription.get_subscription_info())
313        self.assertEqual(created.calls, [])
314        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
315
316        await subscription.subscribe(events=EVENT_FILTERS)
317
318        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
319        self.assertEqual(len(created.calls), 1)
320        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
321
322        notification = server_notification()
323        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(notification))
324        self.assertEqual(len(notifications.calls), 1)
async def test_removal_detaches_listener_and_later_retry_receives_notifications_once(self):
326    async def test_removal_detaches_listener_and_later_retry_receives_notifications_once(self):
327        self.web_socket_client._responder = creation_response
328        subscription = WebSocketSubscription(self.web_socket_client)
329        notifications = RecordingHandler()
330        self.web_socket_client.on(WebSocketEvents.receiveSubscriptionNotification, notifications)
331
332        await subscription.subscribe(events=EVENT_FILTERS)
333        notification = server_notification()
334        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(notification))
335        self.assertEqual(len(notifications.calls), 1)
336
337        await subscription.remove()
338        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(server_notification()))
339        self.assertEqual(len(notifications.calls), 1)
340        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
341
342        await subscription.subscribe(events=EVENT_FILTERS)
343        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(server_notification()))
344        self.assertEqual(len(notifications.calls), 2)
345        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
async def test_update_send_without_response_emits_nothing_and_preserves_confirmed_state(self):
347    async def test_update_send_without_response_emits_nothing_and_preserves_confirmed_state(self):
348        self.web_socket_client._responder = creation_response
349        subscription = WebSocketSubscription(self.web_socket_client)
350        await subscription.subscribe(events=EVENT_FILTERS)
351        confirmed = subscription.get_subscription_info()
352
353        updated = RecordingHandler()
354        failed = RecordingHandler()
355        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
356        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
357
358        self.web_socket_client._responder = None
359        await subscription.update(events=OTHER_FILTERS)
360
361        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
362        self.assertEqual(self.web_socket_client.sent_messages[1][0]["method"], "PUT")
363        self.assertEqual(self.web_socket_client.sent_messages[1][1]["eventFilters"], OTHER_FILTERS)
364        self.assertEqual(updated.calls, [])
365        self.assertEqual(failed.calls, [])
366        self.assertIs(subscription.get_subscription_info(), confirmed)
367        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
async def test_confirmed_update_response_commits_envelope_and_filters_before_single_event(self):
369    async def test_confirmed_update_response_commits_envelope_and_filters_before_single_event(self):
370        self.web_socket_client._responder = creation_response
371        subscription = WebSocketSubscription(self.web_socket_client)
372        await subscription.subscribe(events=EVENT_FILTERS)
373
374        observed_at_event = []
375
376        def on_updated(updated_subscription):
377            observed_at_event.append(updated_subscription.get_subscription_info())
378
379        updated = RecordingHandler()
380        failed = RecordingHandler()
381        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, on_updated)
382        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
383        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
384
385        self.web_socket_client._responder = update_response
386        await subscription.update(events=OTHER_FILTERS)
387
388        self.assertEqual(len(updated.calls), 1)
389        self.assertIs(updated.calls[0][0], subscription)
390        self.assertEqual(failed.calls, [])
391        stored = subscription.get_subscription_info()
392        self.assertEqual(
393            stored[0]["messageId"], self.web_socket_client.sent_messages[1][0]["messageId"]
394        )
395        self.assertEqual(stored[1]["id"], "9d3b7f10-5c11-4b6e-8a2f-6f8b56f0f1c2")
396        self.assertEqual(stored[1]["eventFilters"], OTHER_FILTERS)
397        self.assertEqual(observed_at_event, [stored])
398
399        self.web_socket_client._responder = None
400        await subscription.update()
401        self.assertEqual(self.web_socket_client.sent_messages[2][1]["eventFilters"], OTHER_FILTERS)
async def test_rejected_update_response_emits_single_error_preserves_state_and_permits_retry(self):
403    async def test_rejected_update_response_emits_single_error_preserves_state_and_permits_retry(self):
404        self.web_socket_client._responder = creation_response
405        subscription = WebSocketSubscription(self.web_socket_client)
406        await subscription.subscribe(events=EVENT_FILTERS)
407        confirmed = subscription.get_subscription_info()
408
409        updated = RecordingHandler()
410        failed = RecordingHandler()
411        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
412        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
413
414        self.web_socket_client._responder = lambda request: update_response(request, status=404)
415        await subscription.update(events=OTHER_FILTERS)
416
417        self.assertEqual(len(failed.calls), 1)
418        error = failed.calls[0][0]
419        self.assertIsInstance(error, Exception)
420        self.assertEqual(
421            str(error), "WebSocket subscription update failed with status 404"
422        )
423        self.assertEqual(updated.calls, [])
424        self.assertIs(subscription.get_subscription_info(), confirmed)
425
426        self.web_socket_client._responder = update_response
427        await subscription.update()
428
429        self.assertEqual(len(updated.calls), 1)
430        self.assertEqual(self.web_socket_client.sent_messages[2][1]["eventFilters"], EVENT_FILTERS)
async def test_unrelated_update_response_is_ignored_and_pending_update_still_completes(self):
432    async def test_unrelated_update_response_is_ignored_and_pending_update_still_completes(self):
433        def unrelated_update_response(request):
434            response = update_response(request)
435            response[0]["messageId"] = "unrelated-" + request[0]["messageId"]
436            return response
437
438        self.web_socket_client._responder = creation_response
439        subscription = WebSocketSubscription(self.web_socket_client)
440        await subscription.subscribe(events=EVENT_FILTERS)
441        confirmed = subscription.get_subscription_info()
442
443        updated = RecordingHandler()
444        failed = RecordingHandler()
445        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
446        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
447
448        self.web_socket_client._responder = unrelated_update_response
449        await subscription.update(events=OTHER_FILTERS)
450
451        self.assertEqual(updated.calls, [])
452        self.assertEqual(failed.calls, [])
453        self.assertIs(subscription.get_subscription_info(), confirmed)
454
455        pending_request = self.web_socket_client.sent_messages[1]
456        self.web_socket_client.trigger(
457            WebSocketEvents.receiveMessage, json.dumps(update_response(pending_request))
458        )
459
460        self.assertEqual(len(updated.calls), 1)
461        self.assertEqual(failed.calls, [])
462        self.assertEqual(
463            subscription.get_subscription_info()[0]["messageId"], pending_request[0]["messageId"]
464        )
async def test_update_send_failure_preserves_state_and_permits_retry_after_late_response(self):
466    async def test_update_send_failure_preserves_state_and_permits_retry_after_late_response(self):
467        self.web_socket_client._responder = creation_response
468        subscription = WebSocketSubscription(self.web_socket_client)
469        await subscription.subscribe(events=EVENT_FILTERS)
470        confirmed = subscription.get_subscription_info()
471
472        updated = RecordingHandler()
473        failed = RecordingHandler()
474        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
475        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
476
477        self.web_socket_client._send_error = Exception("connection closed")
478        with self.assertRaises(Exception):
479            await subscription.update(events=OTHER_FILTERS)
480
481        self.assertEqual(updated.calls, [])
482        self.assertEqual(failed.calls, [])
483        self.assertIs(subscription.get_subscription_info(), confirmed)
484
485        pending_request = self.web_socket_client.sent_messages[1]
486        self.web_socket_client.trigger(
487            WebSocketEvents.receiveMessage, json.dumps(update_response(pending_request))
488        )
489        self.assertEqual(updated.calls, [])
490        self.assertEqual(failed.calls, [])
491        self.assertIs(subscription.get_subscription_info(), confirmed)
492
493        self.web_socket_client._responder = update_response
494        await subscription.update()
495
496        self.assertEqual(len(updated.calls), 1)
497        self.assertEqual(self.web_socket_client.sent_messages[2][1]["eventFilters"], EVENT_FILTERS)
async def test_overlapping_update_and_removal_while_update_pending_are_rejected_without_side_effects(self):
499    async def test_overlapping_update_and_removal_while_update_pending_are_rejected_without_side_effects(self):
500        self.web_socket_client._responder = creation_response
501        subscription = WebSocketSubscription(self.web_socket_client)
502        await subscription.subscribe(events=EVENT_FILTERS)
503        confirmed = subscription.get_subscription_info()
504
505        updated = RecordingHandler()
506        removed = RecordingHandler()
507        failed = RecordingHandler()
508        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
509        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
510        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
511
512        self.web_socket_client._responder = None
513        await subscription.update(events=OTHER_FILTERS)
514        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
515
516        with self.assertRaises(Exception):
517            await subscription.update(events=OTHER_FILTERS)
518        with self.assertRaises(Exception):
519            await subscription.remove()
520
521        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
522        self.assertIs(subscription.get_subscription_info(), confirmed)
523        self.assertEqual(updated.calls, [])
524        self.assertEqual(removed.calls, [])
525
526        pending_request = self.web_socket_client.sent_messages[1]
527        self.web_socket_client.trigger(
528            WebSocketEvents.receiveMessage, json.dumps(update_response(pending_request, status=403))
529        )
530        self.assertEqual(len(failed.calls), 1)
531        self.assertEqual(updated.calls, [])
532        self.assertIs(subscription.get_subscription_info(), confirmed)
533
534        self.web_socket_client._responder = update_response
535        await subscription.update()
536
537        self.assertEqual(len(self.web_socket_client.sent_messages), 3)
538        self.assertEqual(self.web_socket_client.sent_messages[2][1]["eventFilters"], EVENT_FILTERS)
async def test_removal_send_without_response_preserves_state_listener_and_emits_nothing(self):
540    async def test_removal_send_without_response_preserves_state_listener_and_emits_nothing(self):
541        self.web_socket_client._responder = creation_response
542        subscription = WebSocketSubscription(self.web_socket_client)
543        await subscription.subscribe(events=EVENT_FILTERS)
544        confirmed = subscription.get_subscription_info()
545
546        removed = RecordingHandler()
547        failed = RecordingHandler()
548        notifications = RecordingHandler()
549        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
550        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, failed)
551        self.web_socket_client.on(WebSocketEvents.receiveSubscriptionNotification, notifications)
552
553        self.web_socket_client._responder = None
554        await subscription.remove()
555
556        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
557        self.assertEqual(self.web_socket_client.sent_messages[1][0]["method"], "DELETE")
558        self.assertEqual(removed.calls, [])
559        self.assertEqual(failed.calls, [])
560        self.assertIs(subscription.get_subscription_info(), confirmed)
561        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
562
563        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(server_notification()))
564        self.assertEqual(len(notifications.calls), 1)
async def test_confirmed_removal_clears_state_and_detaches_listener_before_single_event(self):
566    async def test_confirmed_removal_clears_state_and_detaches_listener_before_single_event(self):
567        self.web_socket_client._responder = creation_response
568        subscription = WebSocketSubscription(self.web_socket_client)
569        await subscription.subscribe(events=EVENT_FILTERS)
570        notifications = RecordingHandler()
571        self.web_socket_client.on(WebSocketEvents.receiveSubscriptionNotification, notifications)
572
573        observed_at_event = []
574
575        def on_removed(*args):
576            observed_at_event.append(
577                (
578                    args,
579                    subscription.get_subscription_info(),
580                    self.web_socket_client.receive_message_listener_count,
581                )
582            )
583
584        removed = RecordingHandler()
585        failed = RecordingHandler()
586        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, on_removed)
587        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
588        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, failed)
589
590        self.web_socket_client._responder = removal_response
591        await subscription.remove()
592
593        self.assertEqual(len(removed.calls), 1)
594        self.assertEqual(removed.calls[0], ())
595        self.assertEqual(failed.calls, [])
596        self.assertEqual(observed_at_event, [((), None, 0)])
597        self.assertIsNone(subscription.get_subscription_info())
598        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
599
600        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(server_notification()))
601        self.assertEqual(len(notifications.calls), 0)
async def test_rejected_removal_response_emits_single_error_and_preserves_everything(self):
603    async def test_rejected_removal_response_emits_single_error_and_preserves_everything(self):
604        self.web_socket_client._responder = creation_response
605        subscription = WebSocketSubscription(self.web_socket_client)
606        await subscription.subscribe(events=EVENT_FILTERS)
607        confirmed = subscription.get_subscription_info()
608        notifications = RecordingHandler()
609        self.web_socket_client.on(WebSocketEvents.receiveSubscriptionNotification, notifications)
610
611        removed = RecordingHandler()
612        failed = RecordingHandler()
613        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
614        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, failed)
615
616        self.web_socket_client._responder = lambda request: removal_response(request, status=404)
617        await subscription.remove()
618
619        self.assertEqual(len(failed.calls), 1)
620        error = failed.calls[0][0]
621        self.assertIsInstance(error, Exception)
622        self.assertEqual(
623            str(error), "WebSocket subscription removal failed with status 404"
624        )
625        self.assertEqual(removed.calls, [])
626        self.assertIs(subscription.get_subscription_info(), confirmed)
627        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
628
629        self.web_socket_client.trigger(WebSocketEvents.receiveMessage, json.dumps(server_notification()))
630        self.assertEqual(len(notifications.calls), 1)
631
632        self.web_socket_client._responder = removal_response
633        await subscription.remove()
634
635        self.assertEqual(len(removed.calls), 1)
636        self.assertIsNone(subscription.get_subscription_info())
637        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
async def test_unrelated_removal_response_is_ignored_and_pending_removal_still_completes(self):
639    async def test_unrelated_removal_response_is_ignored_and_pending_removal_still_completes(self):
640        def unrelated_removal_response(request):
641            response = removal_response(request)
642            response[0]["messageId"] = "unrelated-" + request[0]["messageId"]
643            return response
644
645        self.web_socket_client._responder = creation_response
646        subscription = WebSocketSubscription(self.web_socket_client)
647        await subscription.subscribe(events=EVENT_FILTERS)
648        confirmed = subscription.get_subscription_info()
649
650        removed = RecordingHandler()
651        failed = RecordingHandler()
652        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
653        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, failed)
654
655        self.web_socket_client._responder = unrelated_removal_response
656        await subscription.remove()
657
658        self.assertEqual(removed.calls, [])
659        self.assertEqual(failed.calls, [])
660        self.assertIs(subscription.get_subscription_info(), confirmed)
661        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
662
663        pending_request = self.web_socket_client.sent_messages[1]
664        self.web_socket_client.trigger(
665            WebSocketEvents.receiveMessage, json.dumps(removal_response(pending_request))
666        )
667
668        self.assertEqual(len(removed.calls), 1)
669        self.assertEqual(failed.calls, [])
670        self.assertIsNone(subscription.get_subscription_info())
671        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
async def test_removal_send_failure_preserves_state_and_listener_and_permits_retry(self):
673    async def test_removal_send_failure_preserves_state_and_listener_and_permits_retry(self):
674        self.web_socket_client._responder = creation_response
675        subscription = WebSocketSubscription(self.web_socket_client)
676        await subscription.subscribe(events=EVENT_FILTERS)
677        confirmed = subscription.get_subscription_info()
678
679        removed = RecordingHandler()
680        failed = RecordingHandler()
681        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
682        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, failed)
683
684        self.web_socket_client._send_error = Exception("connection closed")
685        with self.assertRaises(Exception):
686            await subscription.remove()
687
688        self.assertEqual(removed.calls, [])
689        self.assertEqual(failed.calls, [])
690        self.assertIs(subscription.get_subscription_info(), confirmed)
691        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
692
693        pending_request = self.web_socket_client.sent_messages[1]
694        self.web_socket_client.trigger(
695            WebSocketEvents.receiveMessage, json.dumps(removal_response(pending_request))
696        )
697        self.assertEqual(removed.calls, [])
698        self.assertEqual(failed.calls, [])
699        self.assertIs(subscription.get_subscription_info(), confirmed)
700        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
701
702        self.web_socket_client._responder = removal_response
703        await subscription.remove()
704
705        self.assertEqual(len(removed.calls), 1)
706        self.assertIsNone(subscription.get_subscription_info())
707        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
async def test_overlapping_removal_and_update_while_removal_pending_are_rejected_without_side_effects(self):
709    async def test_overlapping_removal_and_update_while_removal_pending_are_rejected_without_side_effects(self):
710        self.web_socket_client._responder = creation_response
711        subscription = WebSocketSubscription(self.web_socket_client)
712        await subscription.subscribe(events=EVENT_FILTERS)
713        confirmed = subscription.get_subscription_info()
714
715        removed = RecordingHandler()
716        failed = RecordingHandler()
717        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
718        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, failed)
719
720        self.web_socket_client._responder = None
721        await subscription.remove()
722        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
723
724        with self.assertRaises(Exception):
725            await subscription.remove()
726        with self.assertRaises(Exception):
727            await subscription.update(events=OTHER_FILTERS)
728
729        self.assertEqual(len(self.web_socket_client.sent_messages), 2)
730        self.assertIs(subscription.get_subscription_info(), confirmed)
731        self.assertEqual(self.web_socket_client.receive_message_listener_count, 1)
732        self.assertEqual(removed.calls, [])
733        self.assertEqual(failed.calls, [])
734
735        pending_request = self.web_socket_client.sent_messages[1]
736        self.web_socket_client.trigger(
737            WebSocketEvents.receiveMessage, json.dumps(removal_response(pending_request))
738        )
739        self.assertEqual(len(removed.calls), 1)
740        self.assertEqual(failed.calls, [])
741        self.assertIsNone(subscription.get_subscription_info())
742        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
async def test_duplicate_and_late_responses_after_completion_produce_no_events_or_changes(self):
744    async def test_duplicate_and_late_responses_after_completion_produce_no_events_or_changes(self):
745        self.web_socket_client._responder = creation_response
746        subscription = WebSocketSubscription(self.web_socket_client)
747        await subscription.subscribe(events=EVENT_FILTERS)
748        creation_request = self.web_socket_client.sent_messages[0]
749
750        updated = RecordingHandler()
751        removed = RecordingHandler()
752        failed = RecordingHandler()
753        remove_failed = RecordingHandler()
754        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
755        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
756        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
757        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, remove_failed)
758
759        self.web_socket_client._responder = update_response
760        await subscription.update(events=OTHER_FILTERS)
761        self.assertEqual(len(updated.calls), 1)
762        stored = subscription.get_subscription_info()
763
764        update_request = self.web_socket_client.sent_messages[1]
765        self.web_socket_client.trigger(
766            WebSocketEvents.receiveMessage, json.dumps(update_response(update_request))
767        )
768        self.web_socket_client.trigger(
769            WebSocketEvents.receiveMessage, json.dumps(creation_response(creation_request))
770        )
771        self.assertEqual(len(updated.calls), 1)
772        self.assertIs(subscription.get_subscription_info(), stored)
773        self.assertEqual(failed.calls, [])
774        self.assertEqual(removed.calls, [])
775
776        self.web_socket_client._responder = removal_response
777        await subscription.remove()
778        self.assertEqual(len(removed.calls), 1)
779        self.assertIsNone(subscription.get_subscription_info())
780
781        removal_request = self.web_socket_client.sent_messages[2]
782        self.web_socket_client.trigger(
783            WebSocketEvents.receiveMessage, json.dumps(removal_response(removal_request))
784        )
785        self.web_socket_client.trigger(
786            WebSocketEvents.receiveMessage, json.dumps(update_response(update_request))
787        )
788        self.assertEqual(len(removed.calls), 1)
789        self.assertEqual(len(updated.calls), 1)
790        self.assertEqual(failed.calls, [])
791        self.assertEqual(remove_failed.calls, [])
792        self.assertIsNone(subscription.get_subscription_info())
793        self.assertEqual(self.web_socket_client.receive_message_listener_count, 0)
async def test_responses_correlated_by_message_id_regardless_of_type(self):
795    async def test_responses_correlated_by_message_id_regardless_of_type(self):
796        self.web_socket_client._responder = creation_response
797        subscription = WebSocketSubscription(self.web_socket_client)
798        await subscription.subscribe(events=EVENT_FILTERS)
799
800        updated = RecordingHandler()
801        removed = RecordingHandler()
802        failed = RecordingHandler()
803        remove_failed = RecordingHandler()
804        self.web_socket_client.on(WebSocketEvents.subscriptionUpdated, updated)
805        self.web_socket_client.on(WebSocketEvents.subscriptionRemoved, removed)
806        self.web_socket_client.on(WebSocketEvents.updateSubscriptionError, failed)
807        self.web_socket_client.on(WebSocketEvents.removeSubscriptionError, remove_failed)
808
809        self.web_socket_client._responder = lambda request: update_response(
810            request, message_type="ClientResponse"
811        )
812        await subscription.update(events=OTHER_FILTERS)
813
814        self.assertEqual(len(updated.calls), 1)
815        self.assertEqual(failed.calls, [])
816        self.assertIsNotNone(subscription.get_subscription_info())
817
818        self.web_socket_client._responder = lambda request: removal_response(
819            request, message_type="ClientResponse"
820        )
821        await subscription.remove()
822
823        self.assertEqual(len(removed.calls), 1)
824        self.assertEqual(failed.calls, [])
825        self.assertEqual(remove_failed.calls, [])
826        self.assertIsNone(subscription.get_subscription_info())
Inherited Members
unittest.async_case.IsolatedAsyncioTestCase
IsolatedAsyncioTestCase
asyncSetUp
asyncTearDown
addAsyncCleanup
enterAsyncContext
run
debug
unittest.case.TestCase
failureException
longMessage
maxDiff
addTypeEqualityFunc
addCleanup
enterContext
addClassCleanup
enterClassContext
tearDown
setUpClass
tearDownClass
countTestCases
defaultTestResult
shortDescription
id
subTest
doCleanups
doClassCleanups
skipTest
fail
assertFalse
assertTrue
assertRaises
assertWarns
assertLogs
assertNoLogs
assertEqual
assertNotEqual
assertAlmostEqual
assertNotAlmostEqual
assertSequenceEqual
assertListEqual
assertTupleEqual
assertSetEqual
assertIn
assertNotIn
assertIs
assertIsNot
assertDictEqual
assertDictContainsSubset
assertCountEqual
assertMultiLineEqual
assertLess
assertLessEqual
assertGreater
assertGreaterEqual
assertIsNone
assertIsNotNone
assertIsInstance
assertNotIsInstance
assertRaisesRegex
assertWarnsRegex
assertRegex
assertNotRegex
failUnlessRaises
failIf
assertRaisesRegexp
assertRegexpMatches
assertNotRegexpMatches
failUnlessEqual
assertEquals
failIfEqual
assertNotEquals
failUnlessAlmostEqual
assertAlmostEquals
failIfAlmostEqual
assertNotAlmostEquals
failUnless
assert_