feat(scheduler): add generic REST API executor for universal HTTP task execution
Add rest_api_executor as a universal executor that can call any REST API
endpoint across the system. This provides a standard way to trigger HTTP
operations from scheduled tasks.
Features:
- All HTTP methods: GET, POST, PUT, DELETE, PATCH
- Authentication: Bearer token, Basic auth, API key
- Environment variable substitution: ${VAR_NAME}
- JSONPath response extraction
- Configurable timeouts and SSL verification
- Sensitive data redaction in logs
- Custom headers support
This executor enables scheduler to call any service endpoint (Library Desk,
Core API, external webhooks) without needing service-specific executors.
Example usage:
{
"executor": "rest_api_executor",
"config": {
"url": "http://library-desk:8089/consolidate/knowledge",
"method": "POST",
"payload": {"process_limit": 10},
"auth": {"type": "bearer", "token": "${API_KEY}"}
}
}
This commit is contained in:
@@ -0,0 +1,266 @@
|
||||
"""
|
||||
Generic REST API Executor
|
||||
|
||||
Universal executor for calling any REST API endpoint across the system.
|
||||
Supports GET, POST, PUT, DELETE with configurable payloads, headers, and authentication.
|
||||
|
||||
This executor can be used to trigger any service endpoint:
|
||||
- Library Desk knowledge consolidation
|
||||
- Core API operations
|
||||
- External webhooks
|
||||
- Any HTTP-based task
|
||||
|
||||
Config schema:
|
||||
{
|
||||
"url": "http://service:port/endpoint",
|
||||
"method": "POST", # GET, POST, PUT, DELETE, PATCH
|
||||
"payload": {...}, # Request body (for POST/PUT/PATCH)
|
||||
"headers": {...}, # Additional headers
|
||||
"auth": {
|
||||
"type": "bearer", # bearer, basic, api_key
|
||||
"token": "${ENV_VAR}", # Use ${VAR} for env vars
|
||||
"header": "Authorization" # Optional: header name for API key
|
||||
},
|
||||
"timeout": 300, # Timeout in seconds (default: 300)
|
||||
"verify_ssl": true, # SSL verification (default: true)
|
||||
"success_codes": [200, 201, 202], # Expected success codes
|
||||
"response_path": "result.message" # JSONPath to extract from response
|
||||
}
|
||||
|
||||
Example configs:
|
||||
|
||||
1. Library Desk Knowledge Consolidation:
|
||||
{
|
||||
"url": "http://library-desk:8089/consolidate/knowledge",
|
||||
"method": "POST",
|
||||
"payload": {"process_limit": 10, "lookback_days": 7, "dry_run": false},
|
||||
"auth": {"type": "bearer", "token": "${LIBRARY_DESK_API_KEY}"}
|
||||
}
|
||||
|
||||
2. Core API Container Restart:
|
||||
{
|
||||
"url": "http://core-api:8088/v1/infrastructure/containers/nginx/restart",
|
||||
"method": "POST",
|
||||
"auth": {"type": "bearer", "token": "${CORE_API_KEY}"}
|
||||
}
|
||||
|
||||
3. External Webhook:
|
||||
{
|
||||
"url": "https://hooks.slack.com/services/YOUR/WEBHOOK/URL",
|
||||
"method": "POST",
|
||||
"payload": {"text": "Scheduled task completed"},
|
||||
"verify_ssl": true
|
||||
}
|
||||
"""
|
||||
import logging
|
||||
import os
|
||||
import re
|
||||
import httpx
|
||||
from typing import Any, Dict, Optional
|
||||
|
||||
from src.config import Settings
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
async def execute(config: dict, settings: Settings) -> str:
|
||||
"""
|
||||
Execute REST API call with configured parameters.
|
||||
|
||||
Args:
|
||||
config: REST API call configuration (see module docstring)
|
||||
settings: Global scheduler settings
|
||||
|
||||
Returns:
|
||||
Response summary or extracted result
|
||||
|
||||
Raises:
|
||||
ValueError: On configuration error
|
||||
Exception: On API call failure
|
||||
"""
|
||||
# Required configuration
|
||||
url = config.get('url')
|
||||
if not url:
|
||||
raise ValueError("Missing required config: 'url'")
|
||||
|
||||
method = config.get('method', 'POST').upper()
|
||||
if method not in ['GET', 'POST', 'PUT', 'DELETE', 'PATCH']:
|
||||
raise ValueError(f"Invalid HTTP method: {method}")
|
||||
|
||||
# Optional configuration
|
||||
payload = config.get('payload', {})
|
||||
headers = config.get('headers', {})
|
||||
timeout = config.get('timeout', 300)
|
||||
verify_ssl = config.get('verify_ssl', True)
|
||||
success_codes = config.get('success_codes', [200, 201, 202, 204])
|
||||
response_path = config.get('response_path')
|
||||
|
||||
# Handle authentication
|
||||
auth_config = config.get('auth', {})
|
||||
if auth_config:
|
||||
auth_header = _build_auth_header(auth_config, settings)
|
||||
if auth_header:
|
||||
headers.update(auth_header)
|
||||
|
||||
# Substitute environment variables in URL and payload
|
||||
url = _substitute_env_vars(url)
|
||||
payload = _substitute_env_vars_recursive(payload)
|
||||
|
||||
logger.info(f"Executing REST API call: {method} {url}")
|
||||
if payload:
|
||||
logger.debug(f"Payload: {_redact_sensitive(payload)}")
|
||||
|
||||
# Make HTTP request
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=timeout, verify=verify_ssl) as client:
|
||||
if method == 'GET':
|
||||
response = await client.get(url, headers=headers)
|
||||
elif method == 'POST':
|
||||
response = await client.post(url, json=payload, headers=headers)
|
||||
elif method == 'PUT':
|
||||
response = await client.put(url, json=payload, headers=headers)
|
||||
elif method == 'DELETE':
|
||||
response = await client.delete(url, headers=headers)
|
||||
elif method == 'PATCH':
|
||||
response = await client.patch(url, json=payload, headers=headers)
|
||||
|
||||
# Check status code
|
||||
if response.status_code not in success_codes:
|
||||
error_msg = (
|
||||
f"API call failed with status {response.status_code}: "
|
||||
f"{response.text[:500]}"
|
||||
)
|
||||
logger.error(error_msg)
|
||||
raise Exception(error_msg)
|
||||
|
||||
# Parse response
|
||||
try:
|
||||
response_data = response.json()
|
||||
except:
|
||||
response_data = {"text": response.text}
|
||||
|
||||
# Extract specific field if response_path provided
|
||||
result_text = None
|
||||
if response_path and isinstance(response_data, dict):
|
||||
result_text = _extract_json_path(response_data, response_path)
|
||||
|
||||
if not result_text:
|
||||
# Build summary from response
|
||||
if isinstance(response_data, dict):
|
||||
# Look for common result fields
|
||||
result_text = (
|
||||
response_data.get('message') or
|
||||
response_data.get('result') or
|
||||
response_data.get('summary') or
|
||||
f"Success ({response.status_code})"
|
||||
)
|
||||
else:
|
||||
result_text = f"Success ({response.status_code})"
|
||||
|
||||
logger.info(f"API call succeeded: {result_text}")
|
||||
return str(result_text)
|
||||
|
||||
except httpx.HTTPStatusError as e:
|
||||
error_msg = f"HTTP {e.response.status_code}: {e.response.text[:500]}"
|
||||
logger.error(error_msg)
|
||||
raise Exception(error_msg)
|
||||
except httpx.RequestError as e:
|
||||
error_msg = f"Request failed: {str(e)}"
|
||||
logger.error(error_msg)
|
||||
raise Exception(error_msg)
|
||||
except Exception as e:
|
||||
logger.error(f"REST API call failed: {e}", exc_info=True)
|
||||
raise
|
||||
|
||||
|
||||
def _build_auth_header(auth_config: dict, settings: Settings) -> Optional[Dict[str, str]]:
|
||||
"""Build authentication header from config."""
|
||||
auth_type = auth_config.get('type', '').lower()
|
||||
|
||||
if auth_type == 'bearer':
|
||||
token = auth_config.get('token', '')
|
||||
token = _substitute_env_vars(token)
|
||||
if token:
|
||||
return {"Authorization": f"Bearer {token}"}
|
||||
|
||||
elif auth_type == 'basic':
|
||||
username = _substitute_env_vars(auth_config.get('username', ''))
|
||||
password = _substitute_env_vars(auth_config.get('password', ''))
|
||||
if username and password:
|
||||
import base64
|
||||
credentials = base64.b64encode(f"{username}:{password}".encode()).decode()
|
||||
return {"Authorization": f"Basic {credentials}"}
|
||||
|
||||
elif auth_type == 'api_key':
|
||||
key = _substitute_env_vars(auth_config.get('key', ''))
|
||||
header_name = auth_config.get('header', 'X-API-Key')
|
||||
if key:
|
||||
return {header_name: key}
|
||||
|
||||
return None
|
||||
|
||||
|
||||
def _substitute_env_vars(text: str) -> str:
|
||||
"""Substitute ${ENV_VAR} placeholders with environment variables."""
|
||||
if not isinstance(text, str):
|
||||
return text
|
||||
|
||||
# Find all ${VAR} patterns
|
||||
pattern = r'\$\{([A-Z_][A-Z0-9_]*)\}'
|
||||
matches = re.findall(pattern, text)
|
||||
|
||||
for var_name in matches:
|
||||
env_value = os.getenv(var_name, '')
|
||||
if not env_value:
|
||||
logger.warning(f"Environment variable not found: {var_name}")
|
||||
text = text.replace(f"${{{var_name}}}", env_value)
|
||||
|
||||
return text
|
||||
|
||||
|
||||
def _substitute_env_vars_recursive(data: Any) -> Any:
|
||||
"""Recursively substitute environment variables in nested structures."""
|
||||
if isinstance(data, dict):
|
||||
return {k: _substitute_env_vars_recursive(v) for k, v in data.items()}
|
||||
elif isinstance(data, list):
|
||||
return [_substitute_env_vars_recursive(item) for item in data]
|
||||
elif isinstance(data, str):
|
||||
return _substitute_env_vars(data)
|
||||
else:
|
||||
return data
|
||||
|
||||
|
||||
def _extract_json_path(data: dict, path: str) -> Optional[str]:
|
||||
"""
|
||||
Extract value from nested dict using dot notation.
|
||||
|
||||
Example: "result.message" -> data["result"]["message"]
|
||||
"""
|
||||
try:
|
||||
keys = path.split('.')
|
||||
value = data
|
||||
for key in keys:
|
||||
if isinstance(value, dict):
|
||||
value = value.get(key)
|
||||
else:
|
||||
return None
|
||||
return str(value) if value is not None else None
|
||||
except:
|
||||
return None
|
||||
|
||||
|
||||
def _redact_sensitive(data: Any) -> Any:
|
||||
"""Redact sensitive fields from logs."""
|
||||
if isinstance(data, dict):
|
||||
redacted = {}
|
||||
sensitive_keys = ['password', 'token', 'api_key', 'secret', 'auth']
|
||||
for k, v in data.items():
|
||||
if any(s in k.lower() for s in sensitive_keys):
|
||||
redacted[k] = '***REDACTED***'
|
||||
else:
|
||||
redacted[k] = _redact_sensitive(v)
|
||||
return redacted
|
||||
elif isinstance(data, list):
|
||||
return [_redact_sensitive(item) for item in data]
|
||||
else:
|
||||
return data
|
||||
Reference in New Issue
Block a user