aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src/meshbay_node/daemon.py
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/daemon.py')
-rw-r--r--packages/meshbay-node/src/meshbay_node/daemon.py23
1 files changed, 12 insertions, 11 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py
index 9428ddd..518f221 100644
--- a/packages/meshbay-node/src/meshbay_node/daemon.py
+++ b/packages/meshbay-node/src/meshbay_node/daemon.py
@@ -37,6 +37,7 @@ from pathlib import Path
import uvicorn
+from meshbay_common.background import spawn
from meshbay_common.paths import fold
from meshbay_common import MNP_VERSION
from meshbay_common.protocol import MNP
@@ -738,7 +739,7 @@ class NodeDaemon:
try:
loop.add_signal_handler(
signal.SIGHUP,
- lambda: asyncio.ensure_future(self._reload_config()))
+ lambda: spawn(self._reload_config()))
except (NotImplementedError, AttributeError):
pass # no SIGHUP on Windows; `reload` says so there
await stop_event.wait()
@@ -1257,7 +1258,7 @@ class NodeDaemon:
def fire() -> None:
self._pending_broadcasts.pop(group_id, None)
- asyncio.ensure_future(self._broadcast_index_change(indexer))
+ spawn(self._broadcast_index_change(indexer))
self._pending_broadcasts[group_id] = loop.call_later(
self._broadcast_coalesce_secs, fire)
@@ -1309,14 +1310,14 @@ class NodeDaemon:
self._enriched_attempted.discard((group_id, entry.id))
seen = {e.id for e in new_entries}
new_entries = new_entries + [e for e in rebuilt if e.id not in seen]
- asyncio.ensure_future(self._enrich_new_video_entries(indexer, new_entries))
+ spawn(self._enrich_new_video_entries(indexer, new_entries))
# Music app (docs/musicbay.md §6): same shape, gated on audio_root
# exactly like video_root above (added later — musicbay.md's
# original "no root, whole shared tree" call didn't hold up).
- asyncio.ensure_future(self._enrich_new_audio_entries(indexer, new_entries))
+ spawn(self._enrich_new_audio_entries(indexer, new_entries))
# Photos app (docs/photos.md §5): same shape, gated on photo_roots
# (a list, not a single string — §2.1).
- asyncio.ensure_future(self._enrich_new_photo_entries(indexer, new_entries))
+ spawn(self._enrich_new_photo_entries(indexer, new_entries))
# A rename/move changes the very filename (or season folder) that
# §3.3/§3.4's title-parse read display_title/season/episode from,
@@ -1327,11 +1328,11 @@ class NodeDaemon:
# (duration/thumb_hash/... landing via _on_enriched below) leaves
# name/path alone and must not re-trigger itself forever.
if delta is not None and delta.updates and previous is not None:
- asyncio.ensure_future(
+ spawn(
self._reenrich_renamed_video_entries(indexer, delta.updates, previous))
- asyncio.ensure_future(
+ spawn(
self._reenrich_renamed_audio_entries(indexer, delta.updates, previous))
- asyncio.ensure_future(
+ spawn(
self._reenrich_renamed_photo_entries(indexer, delta.updates, previous))
# Videos/Music/Photos apps: a file that leaves the index also loses
@@ -1366,7 +1367,7 @@ class NodeDaemon:
still_referenced = any(
i.index.get_entry(file_id) is not None for i in self._indexers)
if not still_referenced:
- asyncio.ensure_future(self._media_cache.prune_file(file_id))
+ spawn(self._media_cache.prune_file(file_id))
# 11.5 — Push to connected WebRTC peers in this group
if self._webrtc:
@@ -1404,7 +1405,7 @@ class NodeDaemon:
else [e.id for e in idx.entries])
if hashes:
endpoint = f"webrtc:{self._config.node.quic_port}"
- asyncio.ensure_future(self._register_swarm(hashes, endpoint))
+ spawn(self._register_swarm(hashes, endpoint))
async def _enrich_new_video_entries(self, indexer: DirectoryIndexer, entries: list) -> None:
"""
@@ -1666,7 +1667,7 @@ class NodeDaemon:
return
for session in list(self._webrtc._sessions.values()):
if session._group_id == group_id:
- asyncio.ensure_future(session.close())
+ spawn(session.close())
log.info("Dropped session for revoked group %s", group_id[:8])
async def _register_swarm(self, hashes: list[str], endpoint: str) -> None: