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()
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
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
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
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
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 ]
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 ]
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 ]
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 ]
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.
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 )
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)
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())
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())
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())
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 )
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)
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)
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)
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)
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)
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)
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)
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)
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)
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)
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)
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)
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)
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_