63 lines
2.1 KiB
Python
63 lines
2.1 KiB
Python
"""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)
|