feat: integrate external scheduler for prefetch task registration
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>
This commit is contained in:
@@ -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
|
||||
Reference in New Issue
Block a user