Override message handling to support request-response
(self, msg)
| 53 | self._log.error("Subscription queue full, dropping error message") |
| 54 | |
| 55 | def _handle_message(self, msg): |
| 56 | """Override message handling to support request-response""" |
| 57 | parsed_msg = super()._handle_message(msg) |
| 58 | self._log.debug(f"Received message: {parsed_msg}") |
| 59 | if parsed_msg is None: |
| 60 | return None |
| 61 | |
| 62 | # Check if this is a subscription event (user data stream, etc.) |
| 63 | # These have 'subscriptionId' and 'event' fields instead of 'id' |
| 64 | if "subscriptionId" in parsed_msg and "event" in parsed_msg: |
| 65 | subscription_id = parsed_msg["subscriptionId"] |
| 66 | event = parsed_msg["event"] |
| 67 | # Route to the registered subscription queue if one exists |
| 68 | if subscription_id in self._subscription_queues: |
| 69 | queue = self._subscription_queues[subscription_id] |
| 70 | try: |
| 71 | queue.put_nowait(event) |
| 72 | except asyncio.QueueFull: |
| 73 | self._log.error( |
| 74 | f"Subscription queue full for {subscription_id}, dropping event" |
| 75 | ) |
| 76 | except Exception as e: |
| 77 | self._log.error( |
| 78 | f"Error putting event in subscription queue for {subscription_id}: {e}" |
| 79 | ) |
| 80 | return None # Don't put in main queue |
| 81 | else: |
| 82 | # No registered queue, return event for main queue (backward compat) |
| 83 | return event |
| 84 | |
| 85 | req_id, exception = None, None |
| 86 | if "id" in parsed_msg: |
| 87 | req_id = parsed_msg["id"] |
| 88 | if "status" in parsed_msg: |
| 89 | if parsed_msg["status"] != 200: |
| 90 | exception = BinanceAPIException( |
| 91 | parsed_msg, |
| 92 | parsed_msg["status"], |
| 93 | self.json_dumps(parsed_msg["error"]), |
| 94 | ) |
| 95 | if req_id is not None and req_id in self._responses: |
| 96 | if exception is not None: |
| 97 | self._responses[req_id].set_exception(exception) |
| 98 | else: |
| 99 | self._responses[req_id].set_result(parsed_msg) |
| 100 | return None # Don't queue request-response messages |
| 101 | elif exception is not None: |
| 102 | raise exception |
| 103 | else: |
| 104 | self._log.warning(f"WS api receieved unknown message: {parsed_msg}") |
| 105 | return None |
| 106 | |
| 107 | async def _ensure_ws_connection(self) -> None: |
| 108 | """Ensure WebSocket connection is established and ready |