156 lines
5.8 KiB
Python
156 lines
5.8 KiB
Python
"""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()
|