Add SchedulerClient to communicate with external scheduler service for registering volatile prefetch tasks discovered during HybridRAG searches. - Add scheduler_client.py with full REST API for task CRUD operations - Add scheduler_url config setting (default: http://scheduler:8090) - Update consolidation service to use scheduler for prefetch registration - Add scheduler health checks to startup/shutdown lifecycle When HybridRAG classifies web content as prefetch-worthy, it now creates scheduled tasks that periodically refresh the volatile cache via the external scheduler service. 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
316 lines
10 KiB
Python
316 lines
10 KiB
Python
"""
|
|
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
|