diff --git a/cryptofeed/exchanges/poloniex.py b/cryptofeed/exchanges/poloniex.py index c382062f4..a34217ec9 100644 --- a/cryptofeed/exchanges/poloniex.py +++ b/cryptofeed/exchanges/poloniex.py @@ -78,19 +78,18 @@ async def _trade(self, msg: dict, timestamp: float): }] } """ - price = Decimal(msg['data'][0]['price']) - amount = Decimal(msg['data'][0]['amount']) - t = Trade( - self.id, - self.exchange_symbol_to_std_symbol(msg['data'][0]['symbol']), - SELL if msg['data'][0]['takerSide'] == 'sell' else BUY, - amount, - price, - self.timestamp_normalize(msg['data'][0]['ts']), - id=str(msg['data'][0]['id']), - raw=msg - ) - await self.callback(TRADES, t, timestamp) + for entry in msg['data']: + t = Trade( + self.id, + self.exchange_symbol_to_std_symbol(entry['symbol']), + SELL if entry['takerSide'] == 'sell' else BUY, + Decimal(entry['quantity']), + Decimal(entry['price']), + self.timestamp_normalize(entry['ts']), + id=str(entry['id']), + raw=msg + ) + await self.callback(TRADES, t, timestamp) async def _book(self, msg: dict, timestamp: float): data = msg['data'][0] diff --git a/tests/unit/test_poloniex.py b/tests/unit/test_poloniex.py new file mode 100644 index 000000000..57c598689 --- /dev/null +++ b/tests/unit/test_poloniex.py @@ -0,0 +1,57 @@ +from decimal import Decimal +from unittest.mock import AsyncMock + +import pytest + +from cryptofeed.defines import BUY, POLONIEX, SELL, TRADES +from cryptofeed.exchanges.poloniex import Poloniex + + +pytestmark = pytest.mark.unit + + +@pytest.mark.parametrize('batch_size', [1, 2]) +@pytest.mark.asyncio +async def test_trade_uses_base_quantity_and_emits_each_batch_entry(batch_size): + # Avoid symbol discovery or any live exchange access. + feed = object.__new__(Poloniex) + feed.exchange_symbol_to_std_symbol = lambda symbol: symbol.replace('_', '-') + feed.callback = AsyncMock() + entries = [ + { + 'symbol': 'BTC_USDT', + 'amount': '364.89973', + 'quantity': '0.017', + 'takerSide': 'sell', + 'price': '21464.69', + 'id': '60183607', + 'ts': 1661120814823, + }, + { + 'symbol': 'ETH_USDT', + 'amount': '3200', + 'quantity': '2', + 'takerSide': 'buy', + 'price': '1600', + 'id': '60183608', + 'ts': 1661120814824, + }, + ][:batch_size] + message = {'channel': 'trades', 'data': entries} + receipt_timestamp = 1661120815.0 + + await feed._trade(message, receipt_timestamp) + + assert feed.callback.await_count == batch_size + for call, entry in zip(feed.callback.await_args_list, entries): + channel, trade, received_at = call.args + assert channel == TRADES + assert trade.exchange == POLONIEX + assert trade.symbol == entry['symbol'].replace('_', '-') + assert trade.side == (SELL if entry['takerSide'] == 'sell' else BUY) + assert trade.amount == Decimal(entry['quantity']) + assert trade.price == Decimal(entry['price']) + assert trade.id == entry['id'] + assert trade.timestamp == entry['ts'] / 1000.0 + assert trade.raw == message + assert received_at == receipt_timestamp