diff --git a/bluesky/input.py b/bluesky/input.py index 83acf39..d135469 100644 --- a/bluesky/input.py +++ b/bluesky/input.py @@ -303,27 +303,24 @@ class BlueskyJetstreamInputService(BlueskyBaseInputService): url += "&wantedCollections=app.bsky.feed.repost" url += f"&wantedDids={self.did}" - async for ws in websockets.connect( - url, - ping_interval=20, - ping_timeout=10, - close_timeout=5, - ): + while True: try: - self.log.info("Listening to '%s'...", env.JETSTREAM_URL) + async with websockets.connect( + url, + ping_interval=20, + ping_timeout=10, + close_timeout=5, + ) as ws: + self.log.info("Listening to '%s'...", env.JETSTREAM_URL) - async def listen_for_messages(): async for msg in ws: self.submitter(lambda: self._accept_msg(msg)) - listen = asyncio.create_task(listen_for_messages()) - - _ = await asyncio.gather(listen) except websockets.ConnectionClosedError as e: - self.log.error(e, stack_info=True, exc_info=True) - self.log.info("Reconnecting to '%s'...", env.JETSTREAM_URL) - continue + self.log.warning("Connection closed: %s", e) except TimeoutError as e: - self.log.error("Connection timeout: '%s'", e) - self.log.info("Reconnecting to '%s'...", env.JETSTREAM_URL) - continue + self.log.warning("Connection timeout: '%s'", e) + except Exception as e: + self.log.error("Unexpected error: '%s'", e, exc_info=True) + + self.log.info("Reconnecting to '%s'...", env.JETSTREAM_URL) diff --git a/mastodon/input.py b/mastodon/input.py index 8be3490..456b0a3 100644 --- a/mastodon/input.py +++ b/mastodon/input.py @@ -248,28 +248,25 @@ class MastodonInputService(MastodonService, InputService): async def listen(self): url = f"{self.streaming_url}/api/v1/streaming?stream=user" - async for ws in websockets.connect( - url, - additional_headers={"Authorization": f"Bearer {self.options.token}"}, - ping_interval=20, - ping_timeout=10, - close_timeout=5, - ): + while True: try: - self.log.info("Listening to '%s'...", self.streaming_url) + async with websockets.connect( + url, + additional_headers={"Authorization": f"Bearer {self.options.token}"}, + ping_interval=20, + ping_timeout=10, + close_timeout=5, + ) as ws: + self.log.info("Listening to '%s'...", self.streaming_url) - async def listen_for_messages(): async for msg in ws: self.submitter(lambda: self._accept_msg(msg)) - listen = asyncio.create_task(listen_for_messages()) - - _ = await asyncio.gather(listen) except websockets.ConnectionClosedError as e: - self.log.error(e, stack_info=True, exc_info=True) - self.log.info("Reconnecting to '%s'...", self.streaming_url) - continue + self.log.warning("Connection closed: %s", e) except TimeoutError as e: - self.log.error("Connection timeout: '%s'", e) - self.log.info("Reconnecting to '%s'...", self.streaming_url) - continue + self.log.warning("Connection timeout: '%s'", e) + except Exception as e: + self.log.error("Unexpected error: '%s'", e, exc_info=True) + + self.log.info("Reconnecting to '%s'...", self.streaming_url) diff --git a/misskey/input.py b/misskey/input.py index 9416fde..1cd00f5 100644 --- a/misskey/input.py +++ b/misskey/input.py @@ -243,28 +243,25 @@ class MisskeyInputService(MisskeyService, InputService): streaming: str = f"{'wss' if self.url.startswith('https') else 'ws'}://{self.url.split('://', 1)[1]}" url: str = f"{streaming}/streaming?i={self.options.token}" - async for ws in websockets.connect( - url, - ping_interval=20, - ping_timeout=10, - close_timeout=5, - ): + while True: try: - self.log.info("Listening to '%s'...", streaming) - await self._subscribe_to_home(ws) + async with websockets.connect( + url, + ping_interval=20, + ping_timeout=10, + close_timeout=5, + ) as ws: + self.log.info("Listening to '%s'...", streaming) + await self._subscribe_to_home(ws) - async def listen_for_messages(): async for msg in ws: self.submitter(lambda: self._accept_msg(msg)) - listen = asyncio.create_task(listen_for_messages()) - - _ = await asyncio.gather(listen) except websockets.ConnectionClosedError as e: - self.log.error(e, stack_info=True, exc_info=True) - self.log.info("Reconnecting to '%s'...", streaming) - continue + self.log.warning("Connection closed: %s", e) except TimeoutError as e: - self.log.error("Connection timeout: '%s'", e) - self.log.info("Reconnecting to '%s'...", streaming) - continue + self.log.warning("Connection timeout: '%s'", e) + except Exception as e: + self.log.error("Unexpected error: '%s'", e, exc_info=True) + + self.log.info("Reconnecting to '%s'...", streaming)