update
This commit is contained in:
6
main.py
6
main.py
@@ -38,13 +38,13 @@ def handle_sighup(signum, frame):
|
|||||||
room_status_log = open("logs/room_status.jsonl", "a", encoding="utf-8")
|
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
|
global TEST_ROOM_IDS
|
||||||
|
|
||||||
TEST_ROOM_IDS = list(set(new_room_ids + room_from_config()))
|
TEST_ROOM_IDS = list(set(new_room_ids + room_from_config()))
|
||||||
|
|
||||||
for removed in clients.keys() - TEST_ROOM_IDS:
|
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]
|
del clients[removed]
|
||||||
|
|
||||||
if removed in log_files:
|
if removed in log_files:
|
||||||
@@ -52,7 +52,7 @@ def on_room_changed(new_room_ids):
|
|||||||
del log_files[removed]
|
del log_files[removed]
|
||||||
|
|
||||||
for added in TEST_ROOM_IDS - clients.keys():
|
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):
|
async def connect_room_safety(room_id: str):
|
||||||
|
|||||||
@@ -1,11 +1,12 @@
|
|||||||
|
import asyncio
|
||||||
import threading
|
import threading
|
||||||
from typing import Callable
|
from typing import Callable, Awaitable
|
||||||
|
|
||||||
import redis
|
import redis
|
||||||
|
|
||||||
room_set_key = "bilibili:live:danmu:room_set"
|
room_set_key = "bilibili:live:danmu:room_set"
|
||||||
r: redis.Redis | None = None
|
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):
|
def init(redis_conf: dict | None):
|
||||||
@@ -21,7 +22,7 @@ def init(redis_conf: dict | None):
|
|||||||
)
|
)
|
||||||
|
|
||||||
# run subscribe_redis on new thread
|
# 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]:
|
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)]
|
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 = r.pubsub()
|
||||||
p.psubscribe(**{f"__keyspace@0__:{room_set_key}": room_changed_handler})
|
p.subscribe(f"__keyspace@0__:{room_set_key}")
|
||||||
for message in p.listen():
|
while p.subscribed:
|
||||||
pass
|
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())
|
||||||
# noinspection PyUnusedLocal
|
|
||||||
def room_changed_handler(ignore):
|
|
||||||
on_room_changed(get_room_list())
|
|
||||||
|
|||||||
Reference in New Issue
Block a user