From de19853aaf86ad70b77ad152d6c84907ab043943 Mon Sep 17 00:00:00 2001 From: zenfyr Date: Thu, 19 Mar 2026 10:54:06 +0700 Subject: [PATCH] i have no idea what i'm doing --- bluesky/input.py | 31 ++++++++++++++----------------- mastodon/input.py | 33 +++++++++++++++------------------ misskey/input.py | 33 +++++++++++++++------------------ 3 files changed, 44 insertions(+), 53 deletions(-) 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) -- 2.51.2