Spaces:
Paused
Paused
| from typing import Dict, Optional, List | |
| from fastapi import WebSocket | |
| from datetime import datetime | |
| from app.core.mesh import MeshRegistry, MeshNode, MeshNodeStatus | |
| from app.core.models import NodeType, WorkerRuntimeType | |
| _mac_node_sockets: Dict[str, WebSocket] = {} | |
| _mac_node_meta: Dict[str, dict] = {} | |
| async def connect_mac_node(node_id: str, websocket: WebSocket) -> bool: | |
| _mac_node_sockets[node_id] = websocket | |
| _mac_node_meta[node_id] = {"capabilities": [], "connected_at": datetime.utcnow(), "last_heartbeat": datetime.utcnow()} | |
| # Register in mesh registry | |
| MeshRegistry.register_node( | |
| MeshNode( | |
| node_id=node_id, | |
| node_type=NodeType.MAC_AGENT, | |
| runtime_type=WorkerRuntimeType.OLLAMA, | |
| capabilities=[], | |
| status=MeshNodeStatus.ONLINE, | |
| supported_models=["llama3.2", "mistral"], | |
| ) | |
| ) | |
| return True | |
| async def disconnect_mac_node(node_id: str) -> None: | |
| _mac_node_sockets.pop(node_id, None) | |
| _mac_node_meta.pop(node_id, None) | |
| node = MeshRegistry.get_node(node_id) | |
| if node: | |
| node.status = MeshNodeStatus.OFFLINE | |
| async def send_to_mac_node(node_id: str, message: dict) -> bool: | |
| ws = _mac_node_sockets.get(node_id) | |
| if ws: | |
| try: | |
| await ws.send_json(message) | |
| return True | |
| except Exception: | |
| return False | |
| return False | |
| def list_mac_nodes() -> list: | |
| return list(_mac_node_sockets.keys()) | |
| def get_mac_nodes_with_capability(capability: str) -> List[str]: | |
| return [nid for nid, meta in _mac_node_meta.items() if capability in meta.get("capabilities", [])] | |
| def update_mac_capabilities(node_id: str, capabilities: List[str]) -> None: | |
| if node_id in _mac_node_meta: | |
| _mac_node_meta[node_id]["capabilities"] = capabilities | |
| def update_mac_heartbeat(node_id: str) -> None: | |
| if node_id in _mac_node_meta: | |
| _mac_node_meta[node_id]["last_heartbeat"] = datetime.utcnow() | |
| async def handle_mac_message(node_id: str, message: dict) -> None: | |
| op = message.get("op") | |
| if op == "node.register": | |
| await send_to_mac_node(node_id, {"op": "node.registered", "node_id": node_id}) | |
| elif op == "node.capabilities": | |
| caps = message.get("capabilities", []) | |
| update_mac_capabilities(node_id, caps) | |
| await send_to_mac_node(node_id, {"op": "node.capabilities_ok", "count": len(caps)}) | |
| elif op == "node.heartbeat": | |
| update_mac_heartbeat(node_id) | |
| elif op == "disconnect": | |
| await disconnect_mac_node(node_id) | |