| 221 | namespace=message.get('namespace')) |
| 222 | |
| 223 | def _thread(self): |
| 224 | while True: |
| 225 | try: |
| 226 | for message in self._listen(): |
| 227 | data = None |
| 228 | if isinstance(message, dict): |
| 229 | data = message |
| 230 | else: |
| 231 | try: |
| 232 | data = self.json.loads(message) |
| 233 | except: |
| 234 | pass |
| 235 | if data and 'method' in data: |
| 236 | self._get_logger().debug('pubsub message: {}'.format( |
| 237 | data['method'])) |
| 238 | try: |
| 239 | if data['method'] == 'callback': |
| 240 | self._handle_callback(data) |
| 241 | elif data.get('host_id') != self.host_id: |
| 242 | if data['method'] == 'emit': |
| 243 | self._handle_emit(data) |
| 244 | elif data['method'] == 'disconnect': |
| 245 | self._handle_disconnect(data) |
| 246 | elif data['method'] == 'enter_room': |
| 247 | self._handle_enter_room(data) |
| 248 | elif data['method'] == 'leave_room': |
| 249 | self._handle_leave_room(data) |
| 250 | elif data['method'] == 'close_room': |
| 251 | self._handle_close_room(data) |
| 252 | except Exception: |
| 253 | self.server.logger.exception( |
| 254 | 'Handler error in pubsub listening thread') |
| 255 | self.server.logger.error('pubsub listen() exited unexpectedly') |
| 256 | break # loop should never exit except in unit tests! |
| 257 | except Exception: # pragma: no cover |
| 258 | self.server.logger.exception('Unexpected Error in pubsub ' |
| 259 | 'listening thread') |