| import logging
|
| import httpx
|
| import os
|
| import asyncio
|
| from typing import Dict, Any, List, Optional
|
| from datetime import datetime, timezone, timedelta
|
|
|
| logger = logging.getLogger(__name__)
|
|
|
| class NodeBridgeService:
|
| def __init__(self):
|
| self.node_url = os.getenv("NODE_ENGINE_URL", "http://localhost:3003")
|
| self.client = httpx.AsyncClient(base_url=self.node_url, timeout=30.0)
|
|
|
|
|
| self._catalog_cache: List[Dict[str, Any]] = []
|
| self._catalog_last_updated: datetime = datetime.min
|
| self._cache_ttl = timedelta(minutes=5)
|
|
|
| async def get_health(self) -> bool:
|
| try:
|
| resp = await self.client.get("/health")
|
| resp.raise_for_status()
|
| return True
|
| except Exception as e:
|
| logger.error(f"Node engine health check failed: {e}")
|
| return False
|
|
|
| async def get_catalog(self, force_refresh: bool = False) -> List[Dict[str, Any]]:
|
| """Fetches pieces from Node engine."""
|
| if not force_refresh and self._catalog_cache and (datetime.now(timezone.utc) - self._catalog_last_updated < self._cache_ttl):
|
| return self._catalog_cache
|
|
|
| try:
|
| resp = await self.client.get("/pieces")
|
| resp.raise_for_status()
|
| pieces = resp.json()
|
|
|
| self._catalog_cache = pieces
|
| self._catalog_last_updated = datetime.now(timezone.utc)
|
| return pieces
|
| except Exception as e:
|
| logger.error(f"Failed to fetch catalog from Node engine: {e}")
|
| return []
|
|
|
| async def get_piece_details(self, piece_name: str) -> Optional[Dict[str, Any]]:
|
| """Get piece details."""
|
| try:
|
| resp = await self.client.get(f"/pieces/{piece_name}")
|
| resp.raise_for_status()
|
| return resp.json()
|
| except httpx.HTTPStatusError as e:
|
| if e.response.status_code == 404:
|
| return None
|
| logger.error(f"Failed to get piece details for {piece_name}: {e}")
|
| return None
|
| except Exception as e:
|
| return None
|
|
|
| async def execute_action(self,
|
| piece_name: str,
|
| action_name: str,
|
| props: Dict[str, Any],
|
| auth: Optional[Dict[str, Any]] = None) -> Dict[str, Any]:
|
| """Executes a piece action."""
|
| payload = {
|
| "pieceName": piece_name,
|
| "actionName": action_name,
|
| "props": props,
|
| "auth": auth
|
| }
|
|
|
| try:
|
| resp = await self.client.post("/execute/action", json=payload)
|
| resp.raise_for_status()
|
| result = resp.json()
|
|
|
| if not result.get("success", False):
|
| raise Exception(f"Execution failed: {result.get('error', 'Unknown error')}")
|
|
|
| return result.get("output", {})
|
| except Exception as e:
|
| logger.error(f"Execution error for {piece_name}.{action_name}: {e}")
|
| raise
|
|
|
| async def close(self):
|
| await self.client.aclose()
|
|
|
|
|
| node_bridge = NodeBridgeService()
|
|
|