diff --git a/main.py b/main.py index 7cb197b..46ba22d 100644 --- a/main.py +++ b/main.py @@ -38,13 +38,13 @@ def handle_sighup(signum, frame): room_status_log = open("logs/room_status.jsonl", "a", encoding="utf-8") -def on_room_changed(new_room_ids): +async def on_room_changed(new_room_ids): global TEST_ROOM_IDS TEST_ROOM_IDS = list(set(new_room_ids + room_from_config())) for removed in clients.keys() - TEST_ROOM_IDS: - asyncio.create_task(clients[removed].stop_and_close()) + await clients[removed].stop_and_close() del clients[removed] if removed in log_files: @@ -52,7 +52,7 @@ def on_room_changed(new_room_ids): del log_files[removed] for added in TEST_ROOM_IDS - clients.keys(): - asyncio.create_task(connect_room_safety(added)) + await connect_room_safety(added) async def connect_room_safety(room_id: str): diff --git a/redis_util.py b/redis_util.py index 82fe2e4..63ac42a 100644 --- a/redis_util.py +++ b/redis_util.py @@ -1,11 +1,12 @@ +import asyncio import threading -from typing import Callable +from typing import Callable, Awaitable import redis room_set_key = "bilibili:live:danmu:room_set" r: redis.Redis | None = None -on_room_changed: Callable[[list[int]], None] = lambda ignore: None +on_room_changed: Callable[[list[int]], Awaitable[None]] | None = None def init(redis_conf: dict | None): @@ -21,7 +22,7 @@ def init(redis_conf: dict | None): ) # run subscribe_redis on new thread - threading.Thread(target=subscribe_redis, daemon=True).start() + threading.Thread(target=_subscribe_redis, daemon=True).start() def get_room_list() -> list[int]: @@ -31,13 +32,17 @@ def get_room_list() -> list[int]: return [int(room_id) for room_id in r.smembers(room_set_key)] -def subscribe_redis(): +def _subscribe_redis(): + loop = asyncio.new_event_loop() + asyncio.set_event_loop(loop) + loop.run_until_complete(_subscribe_redis_async) + loop.close() + + +async def _subscribe_redis_async(): p = r.pubsub() - p.psubscribe(**{f"__keyspace@0__:{room_set_key}": room_changed_handler}) - for message in p.listen(): - pass - - -# noinspection PyUnusedLocal -def room_changed_handler(ignore): - on_room_changed(get_room_list()) + p.subscribe(f"__keyspace@0__:{room_set_key}") + while p.subscribed: + msg = await asyncio.to_thread(p.get_message, ignore_subscribe_messages=True, timeout=1) + if msg and on_room_changed is not None: + await on_room_changed(get_room_list())