diff options
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/daemon.py')
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/daemon.py | 23 |
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: |