"""Best-effort periodic full scans for the DeerFlow-local Wiki index.""" from __future__ import annotations import asyncio import logging from collections.abc import Coroutine from typing import Any from fastapi import FastAPI logger = logging.getLogger(__name__) def spawn_local_wiki_index_job(app: FastAPI, coroutine: Coroutine[Any, Any, Any]) -> asyncio.Task: """Keep fire-and-forget index jobs strongly referenced until completion.""" jobs: set[asyncio.Task] = getattr(app.state, "llmwiki_index_jobs", None) if jobs is None: jobs = set() app.state.llmwiki_index_jobs = jobs task = asyncio.create_task(coroutine) jobs.add(task) task.add_done_callback(jobs.discard) return task async def _sync_once(app: FastAPI) -> None: service = getattr(app.state, "llmwiki_sync_service", None) mapping_store = getattr(app.state, "llmwiki_store", None) if service is None or mapping_store is None: return config = app.state.config.llmwiki.local_wiki_index mappings = [row for row in await mapping_store.list_visible("system", scope="all", is_admin=True) if row.get("wiki_index_enabled", True)] semaphore = asyncio.Semaphore(config.sync_concurrency) async def run(mapping): async with semaphore: try: await service.sync_mapping(mapping) except asyncio.CancelledError: raise except Exception: logger.warning("Periodic local Wiki sync failed mapping_id=%s", mapping.get("id"), exc_info=True) await asyncio.gather(*(run(mapping) for mapping in mappings)) async def local_wiki_index_scheduler_loop(app: FastAPI) -> None: interval = app.state.config.llmwiki.local_wiki_index.sync_interval_seconds logger.info( "Local Wiki auto-sync scheduler started; first scan is immediate, interval=%s seconds", interval, ) while True: try: await _sync_once(app) except asyncio.CancelledError: raise except Exception: logger.exception("Local Wiki auto-sync scheduler tick failed (non-fatal)") await asyncio.sleep(interval)