(self, client: Client, message)
| 748 | return await self.un_watch_ohlcv_for_symbols([[symbol, timeframe]], params) |
| 749 | |
| 750 | def handle_ohlcv(self, client: Client, message): |
| 751 | # |
| 752 | # { |
| 753 | # "topic": "kline.5.BTCUSDT", |
| 754 | # "data": [ |
| 755 | # { |
| 756 | # "start": 1672324800000, |
| 757 | # "end": 1672325099999, |
| 758 | # "interval": "5", |
| 759 | # "open": "16649.5", |
| 760 | # "close": "16677", |
| 761 | # "high": "16677", |
| 762 | # "low": "16608", |
| 763 | # "volume": "2.081", |
| 764 | # "turnover": "34666.4005", |
| 765 | # "confirm": False, |
| 766 | # "timestamp": 1672324988882 |
| 767 | # } |
| 768 | # ], |
| 769 | # "ts": 1672324988882, |
| 770 | # "type": "snapshot" |
| 771 | # } |
| 772 | # |
| 773 | data = self.safe_value(message, 'data', {}) |
| 774 | topic = self.safe_string(message, 'topic', '') |
| 775 | topicParts = topic.split('.') |
| 776 | topicLength = len(topicParts) |
| 777 | timeframeId = self.safe_string(topicParts, 1) |
| 778 | timeframe = self.find_timeframe(timeframeId) |
| 779 | if timeframe is None: |
| 780 | return |
| 781 | marketId = self.safe_string(topicParts, topicLength - 1) |
| 782 | isSpot = client.url.find('spot') > -1 |
| 783 | marketType = 'spot' if isSpot else 'contract' |
| 784 | market = self.safe_market(marketId, None, None, marketType) |
| 785 | symbol = market['symbol'] |
| 786 | ohlcvsByTimeframe = self.safe_value(self.ohlcvs, symbol) |
| 787 | if ohlcvsByTimeframe is None: |
| 788 | self.ohlcvs[symbol] = {} |
| 789 | if self.safe_value(ohlcvsByTimeframe, timeframe) is None: |
| 790 | limit = self.safe_integer(self.options, 'OHLCVLimit', 1000) |
| 791 | self.ohlcvs[symbol][timeframe] = ArrayCacheByTimestamp(limit) |
| 792 | stored = self.ohlcvs[symbol][timeframe] |
| 793 | for i in range(0, len(data)): |
| 794 | parsed = self.parse_ws_ohlcv(data[i], market) |
| 795 | stored.append(parsed) |
| 796 | messageHash = 'ohlcv::' + symbol + '::' + timeframe |
| 797 | resolveData = [symbol, timeframe, stored] |
| 798 | client.resolve(resolveData, messageHash) |
| 799 | |
| 800 | def parse_ws_ohlcv(self, ohlcv, market: Market = None) -> list: |
| 801 | # |
nothing calls this directly
no test coverage detected