diff --git a/src/clients/scheduler_client.py b/src/clients/scheduler_client.py new file mode 100644 index 0000000..90496e5 --- /dev/null +++ b/src/clients/scheduler_client.py @@ -0,0 +1,315 @@ +""" +Client for external Scheduler service. + +Registers and manages scheduled tasks for prefetch operations +(weather, news, etc.) discovered through HybridRAG searches. +""" + +import httpx +import logging +from typing import Optional, Any +from pydantic import BaseModel, Field + +logger = logging.getLogger(__name__) + + +class SchedulerTask(BaseModel): + """Task definition for scheduler registration.""" + + task_name: str = Field(..., description="Unique task identifier") + service: str = Field(default="library-desk", description="Service that owns this task") + executor: str = Field(default="rest_api", description="Executor type") + priority: int = Field(default=50, ge=1, le=100, description="Priority (lower = higher)") + description: Optional[str] = Field(None, description="Human-readable description") + enabled: bool = Field(default=True, description="Whether task is enabled") + max_retries: int = Field(default=3, ge=0, le=10, description="Max retry attempts") + timeout_seconds: int = Field(default=3600, ge=1, description="Execution timeout") + + # Schedule (-1 = every, or specific value) + minute: int = Field(default=-1, ge=-1, le=59, description="Minute (-1=every)") + hour: int = Field(default=-1, ge=-1, le=23, description="Hour (-1=every)") + day_of_month: int = Field(default=-1, ge=-1, le=31, description="Day of month (-1=every)") + month: int = Field(default=-1, ge=-1, le=12, description="Month (-1=every)") + day_of_week: int = Field(default=-1, ge=-1, le=6, description="Day of week (-1=every, 0=Mon)") + + # Executor config (for rest_api executor) + config: Optional[dict[str, Any]] = Field(None, description="Executor-specific config") + + +class SchedulerClient: + """Client for external scheduler service.""" + + def __init__(self, base_url: str, timeout: float = 30.0): + """ + Initialize scheduler client. + + Args: + base_url: Scheduler API base URL (e.g., "http://scheduler:8090") + timeout: HTTP request timeout in seconds + """ + self.base_url = base_url.rstrip("/") + self.timeout = timeout + self._client: Optional[httpx.AsyncClient] = None + + async def _get_client(self) -> httpx.AsyncClient: + """Get or create HTTP client.""" + if self._client is None or self._client.is_closed: + self._client = httpx.AsyncClient( + base_url=self.base_url, + timeout=self.timeout, + ) + return self._client + + async def close(self): + """Close HTTP client.""" + if self._client and not self._client.is_closed: + await self._client.aclose() + self._client = None + logger.info("Scheduler client closed") + + async def health_check(self) -> bool: + """Check scheduler connectivity.""" + try: + client = await self._get_client() + response = await client.get("/health") + return response.status_code == 200 + except Exception as e: + logger.error(f"Scheduler health check failed: {e}") + return False + + async def task_exists(self, task_name: str) -> bool: + """ + Check if a task already exists. + + Args: + task_name: Task identifier to check + + Returns: + True if task exists, False otherwise. + """ + try: + client = await self._get_client() + response = await client.get(f"/tasks/{task_name}") + return response.status_code == 200 + except Exception as e: + logger.error(f"Failed to check task existence: {e}") + return False + + async def get_task(self, task_name: str) -> Optional[dict[str, Any]]: + """ + Get task details. + + Args: + task_name: Task identifier + + Returns: + Task dict or None if not found. + """ + try: + client = await self._get_client() + response = await client.get(f"/tasks/{task_name}") + if response.status_code == 200: + return response.json() + return None + except Exception as e: + logger.error(f"Failed to get task {task_name}: {e}") + return None + + async def list_tasks( + self, + service: Optional[str] = None, + enabled: Optional[bool] = None + ) -> list[dict[str, Any]]: + """ + List scheduled tasks. + + Args: + service: Filter by service name + enabled: Filter by enabled status + + Returns: + List of task dicts. + """ + try: + client = await self._get_client() + params = {} + if service: + params["service"] = service + if enabled is not None: + params["enabled"] = enabled + + response = await client.get("/tasks", params=params) + if response.status_code == 200: + return response.json() + return [] + except Exception as e: + logger.error(f"Failed to list tasks: {e}") + return [] + + async def create_task(self, task: SchedulerTask) -> Optional[dict[str, Any]]: + """ + Create a new scheduled task. + + Args: + task: Task definition + + Returns: + Created task dict or None on failure. + """ + try: + client = await self._get_client() + response = await client.post( + "/tasks", + json=task.model_dump(exclude_none=True) + ) + if response.status_code == 200: + logger.info(f"Created scheduler task: {task.task_name}") + return response.json() + else: + logger.error( + f"Failed to create task {task.task_name}: " + f"{response.status_code} - {response.text}" + ) + return None + except Exception as e: + logger.error(f"Failed to create task {task.task_name}: {e}") + return None + + async def update_task( + self, + task_name: str, + updates: dict[str, Any] + ) -> Optional[dict[str, Any]]: + """ + Update an existing task. + + Args: + task_name: Task identifier + updates: Fields to update + + Returns: + Updated task dict or None on failure. + """ + try: + client = await self._get_client() + response = await client.put(f"/tasks/{task_name}", json=updates) + if response.status_code == 200: + logger.info(f"Updated scheduler task: {task_name}") + return response.json() + else: + logger.error( + f"Failed to update task {task_name}: " + f"{response.status_code} - {response.text}" + ) + return None + except Exception as e: + logger.error(f"Failed to update task {task_name}: {e}") + return None + + async def delete_task(self, task_name: str) -> bool: + """ + Delete a scheduled task. + + Args: + task_name: Task identifier + + Returns: + True if deleted, False otherwise. + """ + try: + client = await self._get_client() + response = await client.delete(f"/tasks/{task_name}") + if response.status_code == 200: + logger.info(f"Deleted scheduler task: {task_name}") + return True + else: + logger.error( + f"Failed to delete task {task_name}: " + f"{response.status_code} - {response.text}" + ) + return False + except Exception as e: + logger.error(f"Failed to delete task {task_name}: {e}") + return False + + async def trigger_task(self, task_name: str) -> bool: + """ + Manually trigger a task to run immediately. + + Args: + task_name: Task identifier + + Returns: + True if triggered, False otherwise. + """ + try: + client = await self._get_client() + response = await client.post(f"/tasks/{task_name}/trigger") + if response.status_code == 200: + logger.info(f"Triggered task: {task_name}") + return True + else: + logger.error( + f"Failed to trigger task {task_name}: " + f"{response.status_code} - {response.text}" + ) + return False + except Exception as e: + logger.error(f"Failed to trigger task {task_name}: {e}") + return False + + async def register_volatile_fetch( + self, + namespace: str, + key: str, + user: str, + schedule: dict[str, int], + description: Optional[str] = None, + ) -> bool: + """ + Register a volatile fetch task for prefetch. + + Convenience method to create tasks that call /volatile/fetch endpoints. + + Args: + namespace: Volatile namespace (e.g., "weather", "news") + key: Volatile key (e.g., "rotterdam", "nos") + user: User for the fetch + schedule: Cron-like schedule dict (minute, hour, etc.) + description: Human-readable description + + Returns: + True if registered (or already exists), False on failure. + """ + task_name = f"volatile_{namespace}_{key}_{user}".replace("-", "_") + + # Check if already exists + if await self.task_exists(task_name): + logger.info(f"Prefetch task already exists: {task_name}") + return True + + task = SchedulerTask( + task_name=task_name, + service="library-desk", + executor="rest_api", + priority=60, # Background maintenance priority + description=description or f"Prefetch {namespace}/{key} for {user}", + minute=schedule.get("minute", -1), + hour=schedule.get("hour", -1), + day_of_month=schedule.get("day_of_month", -1), + month=schedule.get("month", -1), + day_of_week=schedule.get("day_of_week", -1), + config={ + "method": "POST", + "url": f"http://library-desk:8089/volatile/fetch/{namespace}/{key}", + "headers": { + "Content-Type": "application/json" + }, + "body": { + "user": user + } + } + ) + + result = await self.create_task(task) + return result is not None diff --git a/src/config.py b/src/config.py index 53fb2b8..4eac3c5 100644 --- a/src/config.py +++ b/src/config.py @@ -128,6 +128,9 @@ class Settings(BaseSettings): system_settings_user: str = Field(default="settings", description="System settings database user") system_settings_password: str = Field(default="", description="System settings database password") + # Scheduler Service + scheduler_url: str = Field(default="http://scheduler:8090", description="Scheduler service URL") + @property def qdrant_url(self) -> str: """Computed Qdrant URL.""" diff --git a/src/core/dependencies.py b/src/core/dependencies.py index 3ab8943..3df364f 100644 --- a/src/core/dependencies.py +++ b/src/core/dependencies.py @@ -24,6 +24,7 @@ from src.clients.ollama_client import OllamaClient from src.clients.content_extractor import ContentExtractor from src.clients.paperless_client import PaperlessClient from src.clients.settings_client import SettingsClient +from src.clients.scheduler_client import SchedulerClient from src.apis import ( OpenMeteoProvider, AggregatedNewsProvider, @@ -202,6 +203,22 @@ def get_settings_client() -> SettingsClient: return client +@lru_cache +def get_scheduler_client() -> SchedulerClient: + """ + Get scheduler service client singleton. + + Returns: + Initialized SchedulerClient for task management + + Note: Used for registering prefetch tasks discovered during HybridRAG searches + """ + settings = get_settings() + client = SchedulerClient(base_url=settings.scheduler_url) + logger.debug(f"Created Scheduler client: {settings.scheduler_url}") + return client + + # ============================================================================= # External API Providers # ============================================================================= @@ -320,6 +337,7 @@ RedisDep = Annotated[aioredis.Redis, Depends(get_redis_client)] ContentExtractorDep = Annotated[ContentExtractor, Depends(get_content_extractor)] PaperlessDep = Annotated[PaperlessClient, Depends(get_paperless_client)] SettingsClientDep = Annotated[SettingsClient, Depends(get_settings_client)] +SchedulerDep = Annotated[SchedulerClient, Depends(get_scheduler_client)] # External API provider dependencies WeatherProviderDep = Annotated[OpenMeteoProvider, Depends(get_weather_provider)] @@ -392,6 +410,17 @@ async def startup_clients(): else: logger.info("○ System settings not configured") + # Check Scheduler availability + try: + scheduler = get_scheduler_client() + is_healthy = await scheduler.health_check() + if is_healthy: + logger.info(f"✓ Scheduler ready: {settings.scheduler_url}") + else: + logger.warning("✗ Scheduler not responding") + except Exception as e: + logger.error(f"✗ Scheduler health check failed: {e}") + # Qdrant, Wiki.js, SearXNG are lazy-initialized logger.info("Service clients startup complete") @@ -461,6 +490,14 @@ async def shutdown_clients(): except Exception as e: logger.error(f"Error closing settings client: {e}") + # Close scheduler client + try: + scheduler = get_scheduler_client() + await scheduler.close() + logger.info("✓ Scheduler client closed") + except Exception as e: + logger.error(f"Error closing scheduler client: {e}") + logger.info("Service clients shutdown complete") @@ -556,6 +593,14 @@ async def check_service_health() -> dict: else: health["system_settings"] = None # Not configured + # Scheduler + try: + scheduler = get_scheduler_client() + health["scheduler"] = await scheduler.health_check() + except Exception as e: + logger.error(f"Scheduler health check failed: {e}") + health["scheduler"] = False + return health @@ -600,6 +645,7 @@ def get_consolidation_service() -> "ConsolidationService": ingestion_service=get_ingestion_service(), volatile_service=get_volatile_cache_service(), settings_client=get_settings_client(), + scheduler_client=get_scheduler_client(), ) diff --git a/src/services/consolidation_service.py b/src/services/consolidation_service.py index 9fbb1b7..7811ff2 100644 --- a/src/services/consolidation_service.py +++ b/src/services/consolidation_service.py @@ -46,6 +46,7 @@ class ConsolidationService: ingestion_service: Optional["IngestionService"] = None, volatile_service: Optional["VolatileCacheService"] = None, settings_client: Optional["SettingsClient"] = None, + scheduler_client: Optional["SchedulerClient"] = None, ): self.neo4j = neo4j self.ollama = ollama @@ -54,7 +55,8 @@ class ConsolidationService: self.wiki_page_writer = WikiPageWriter(ollama_client=ollama, settings=settings) self.ingestion_service = ingestion_service # Optional to avoid circular dependency self.volatile_service = volatile_service # For ephemeral data caching - self.settings_client = settings_client # For prefetch registration + self.settings_client = settings_client # For prefetch registration (fallback) + self.scheduler_client = scheduler_client # For scheduler-driven prefetch async def consolidate_knowledge( self, @@ -1233,7 +1235,7 @@ JSON:""" user: str, ) -> bool: """ - Register a prefetch pattern with the scheduler. + Register a prefetch pattern with the external scheduler service. Args: classification: The classification with prefetch info @@ -1243,31 +1245,70 @@ JSON:""" Returns: True if successfully registered, False otherwise """ - if not self.settings_client: - logger.warning("Settings client not configured, skipping prefetch registration") + if not self.scheduler_client: + logger.warning("Scheduler client not configured, skipping prefetch registration") return False - cron = classification.prefetch_cron or "0 * * * *" # Default: hourly - endpoint = classification.prefetch_endpoint or "" + # Parse cron pattern into scheduler schedule format + # Format: "minute hour day_of_month month day_of_week" + # Scheduler uses -1 for "every" + cron = classification.prefetch_cron or "0 * * * *" + schedule = self._parse_cron_to_schedule(cron) - if not endpoint: - logger.warning(f"No prefetch endpoint specified for {web_result.get('url')}") + # Determine namespace and key from classification + namespace = classification.volatile_namespace or "custom" + key = classification.volatile_key or web_result.get('url', '').split('/')[-1].split('?')[0] + + if not key: + logger.warning(f"Could not determine prefetch key for {web_result.get('url')}") return False try: - # Store prefetch configuration in settings - prefetch_key = f"prefetch.{user}.{classification.volatile_key or 'auto'}" - prefetch_config = { - "cron": cron, - "endpoint": endpoint, - "source_url": web_result.get('url', ''), - "enabled": True, - } + # Use the scheduler client's convenience method to register volatile fetch + success = await self.scheduler_client.register_volatile_fetch( + namespace=namespace, + key=key, + user=user, + schedule=schedule, + description=f"Auto-prefetch: {classification.title or web_result.get('title', 'Unknown')}", + ) - await self.settings_client.set(prefetch_key, prefetch_config, user=user) - logger.info(f"Registered prefetch: {prefetch_key} ({cron})") - return True + if success: + logger.info(f"Registered scheduler task: volatile_{namespace}_{key}_{user}") + return success except Exception as e: - logger.error(f"Failed to register prefetch: {e}") + logger.error(f"Failed to register prefetch with scheduler: {e}") return False + + def _parse_cron_to_schedule(self, cron: str) -> dict: + """ + Parse cron string to scheduler schedule dict. + + Args: + cron: Cron-style string (e.g., "0 6 * * *" = 6:00 AM daily) + + Returns: + Dict with minute, hour, day_of_month, month, day_of_week + where -1 means "every" + """ + parts = cron.strip().split() + if len(parts) != 5: + # Default to hourly if invalid + return {"minute": 0, "hour": -1} + + def parse_part(part: str) -> int: + if part == "*": + return -1 + try: + return int(part) + except ValueError: + return -1 + + return { + "minute": parse_part(parts[0]), + "hour": parse_part(parts[1]), + "day_of_month": parse_part(parts[2]), + "month": parse_part(parts[3]), + "day_of_week": parse_part(parts[4]), + }