deerflow-code/offline-backend-20260512/backend/app/gateway/skill_knowledge_auto_sync.py
2026-09-07 18:24:55 +08:00

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()