From 28808e200a9f3b2f51a61cc0920893e3591c0bde Mon Sep 17 00:00:00 2001 From: zack <43246297+clickingbuttons@users.noreply.github.com> Date: Tue, 6 Sep 2022 10:20:50 -0400 Subject: [PATCH 1/2] Fixes by at-cf --- examples/websocket/latency.py | 15 +++++++++++++++ polygon/websocket/__init__.py | 22 +++++++++++++--------- 2 files changed, 28 insertions(+), 9 deletions(-) create mode 100644 examples/websocket/latency.py diff --git a/examples/websocket/latency.py b/examples/websocket/latency.py new file mode 100644 index 00000000..dd7cb5fe --- /dev/null +++ b/examples/websocket/latency.py @@ -0,0 +1,15 @@ +from polygon import WebSocketClient +from polygon.websocket.models import WebSocketMessage +from typing import List +import time + +c = WebSocketClient(subscriptions=["Q.SPY"]) + + +def handle_msg(msgs: List[WebSocketMessage]): + for m in msgs: + now = time.time() * 1000 + print(now, m.timestamp, now - m.timestamp) + + +c.run(handle_msg) diff --git a/polygon/websocket/__init__.py b/polygon/websocket/__init__.py index 21f6bfcb..66a9dd9a 100644 --- a/polygon/websocket/__init__.py +++ b/polygon/websocket/__init__.py @@ -118,16 +118,20 @@ async def connect( self.subs = set(self.scheduled_subs) self.schedule_resub = False - cmsg: Union[ - List[WebSocketMessage], Union[str, bytes] - ] = await s.recv() - # we know cmsg is Data - msgJson = json.loads(cmsg) # type: ignore - for m in msgJson: - if m["ev"] == "status": - logger.debug("status: %s", m["message"]) - continue + try: + cmsg: Union[ + List[WebSocketMessage], Union[str, bytes] + ] = await asyncio.wait_for(s.recv(), timeout=1) + except asyncio.TimeoutError: + continue + if not self.raw: + # we know cmsg is Data + msgJson = json.loads(cmsg) # type: ignore + for m in msgJson: + if m["ev"] == "status": + logger.debug("status: %s", m["message"]) + continue cmsg = parse(msgJson, logger) if len(cmsg) > 0: From e8e2eb962dc0141e70136088c310579f40cfe88d Mon Sep 17 00:00:00 2001 From: zack <43246297+clickingbuttons@users.noreply.github.com> Date: Tue, 6 Sep 2022 10:26:28 -0400 Subject: [PATCH 2/2] fix lint --- examples/websocket/latency.py | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/examples/websocket/latency.py b/examples/websocket/latency.py index dd7cb5fe..57a0a183 100644 --- a/examples/websocket/latency.py +++ b/examples/websocket/latency.py @@ -1,6 +1,6 @@ from polygon import WebSocketClient -from polygon.websocket.models import WebSocketMessage -from typing import List +from polygon.websocket.models import WebSocketMessage, EquityQuote +from typing import List, cast import time c = WebSocketClient(subscriptions=["Q.SPY"]) @@ -8,8 +8,10 @@ def handle_msg(msgs: List[WebSocketMessage]): for m in msgs: - now = time.time() * 1000 - print(now, m.timestamp, now - m.timestamp) + q: EquityQuote = cast(EquityQuote, m) + if q.timestamp is not None: + now = time.time() * 1000 + print(now, q.timestamp, now - q.timestamp) c.run(handle_msg)