"""Debounced file watcher and cron scheduler for opted-in knowledge bindings.""" from __future__ import annotations import asyncio import logging from datetime import UTC, datetime from typing import Any from croniter import croniter from app.gateway.skill_knowledge_job_executor import _skill_directory from deerflow.config.system_settings import load_system_settings from deerflow.skill_knowledge import scan_skill logger = logging.getLogger(__name__) class SkillKnowledgeAutoSyncService: def __init__(self, app: Any, *, poll_seconds: float = 10.0) -> None: self.app = app self.store = app.state.skill_knowledge_store self.poll_seconds = poll_seconds self._task: asyncio.Task | None = None self._wake = asyncio.Event() self._observed: dict[str, tuple[str, datetime]] = {} self._schedule_expression = "" self._next_schedule: datetime | None = None def start(self) -> None: if self._task is None or self._task.done(): self._task = asyncio.create_task(self._loop()) def nudge(self) -> None: self._wake.set() async def close(self) -> None: if self._task is None: return self._task.cancel() try: await self._task except asyncio.CancelledError: pass async def _queue_skill( self, skill_name: str, bindings: list[dict[str, Any]], *, trigger: str, review_mode: str, ) -> None: targets = [ { "target_type": binding["target_type"], "target_mode": binding["target_mode"], "target_id": binding["target_id"], "target_name": binding["target_name"], "remote_business_key": binding.get("remote_business_key"), } for binding in bindings ] job = await self.store.create_job( skill_names=[skill_name], targets=targets, created_by=None, trigger=trigger, review_mode=review_mode, ) if job["total_count"]: self.app.state.skill_knowledge_dispatcher.nudge() async def _watch( self, bindings_by_skill: dict[str, list[dict[str, Any]]], *, debounce_seconds: int, write_remote_enabled: bool, review_mode: str, ) -> None: now = datetime.now(UTC) for skill_name, bindings in bindings_by_skill.items(): try: root = _skill_directory(self.app.state.config, skill_name) scan = await asyncio.to_thread(scan_skill, skill_name, root) except Exception: logger.exception("Failed to fingerprint skill %s", skill_name) continue synced = {str(binding.get("synced_digest") or "") for binding in bindings} if synced == {scan.source_digest}: self._observed.pop(skill_name, None) continue await self.store.mark_skill_stale(skill_name) previous = self._observed.get(skill_name) if previous is None or previous[0] != scan.source_digest: self._observed[skill_name] = (scan.source_digest, now) continue if (now - previous[1]).total_seconds() < debounce_seconds: continue if write_remote_enabled: current = await self.store.list_bindings(skill_names=[skill_name]) stale = [binding for binding in current if binding.get("auto_sync_enabled") and binding.get("status") == "stale"] if stale: await self._queue_skill( skill_name, stale, trigger="watcher", review_mode=review_mode, ) self._observed.pop(skill_name, None) def _schedule_due(self, expression: str, now: datetime) -> bool: if expression != self._schedule_expression or self._next_schedule is None: self._schedule_expression = expression self._next_schedule = croniter(expression, now).get_next(datetime) return False if now < self._next_schedule: return False self._next_schedule = croniter(expression, now).get_next(datetime) return True async def _tick(self) -> None: settings = (await asyncio.to_thread(load_system_settings)).skill_knowledge_auto_sync bindings = [row for row in await self.store.list_bindings() if row.get("auto_sync_enabled")] by_skill: dict[str, list[dict[str, Any]]] = {} for binding in bindings: by_skill.setdefault(str(binding["skill_name"]), []).append(binding) if settings.watch_enabled: await self._watch( by_skill, debounce_seconds=settings.debounce_seconds, write_remote_enabled=settings.write_remote_enabled, review_mode=settings.review_mode, ) if settings.schedule_enabled and settings.write_remote_enabled and by_skill and self._schedule_due(settings.schedule_cron, datetime.now(UTC)): for skill_name, skill_bindings in by_skill.items(): await self._queue_skill( skill_name, skill_bindings, trigger="scheduled", review_mode=settings.review_mode, ) async def _loop(self) -> None: while True: try: await self._tick() except Exception: logger.exception("Skill knowledge automatic sync tick failed") try: await asyncio.wait_for(self._wake.wait(), timeout=self.poll_seconds) except TimeoutError: pass self._wake.clear()