zelos.protocol_adapters
Phase 2 Protocol Adapters — gRPC, MCP, A2A, WebSocket.
Each adapter translates an external protocol → Runtime API. Adapters contain ZERO business logic. They are stateless, replaceable plugins.
Architecture:
Protocol Adapter (translate) → Runtime API → Kernel
1""" 2Phase 2 Protocol Adapters — gRPC, MCP, A2A, WebSocket. 3 4Each adapter translates an external protocol → Runtime API. 5Adapters contain ZERO business logic. They are stateless, replaceable plugins. 6 7Architecture: 8 Protocol Adapter (translate) → Runtime API → Kernel 9""" 10 11import threading 12import time 13from abc import ABC, abstractmethod 14from typing import Any 15 16# ═══════════════════ Adapter Base ═══════════════════ 17 18 19class ProtocolAdapter(ABC): 20 """Base for all protocol adapters.""" 21 22 def __init__(self, runtime=None): 23 self.runtime = runtime 24 25 @abstractmethod 26 def start(self) -> None: ... 27 28 @abstractmethod 29 def stop(self) -> None: ... 30 31 32# ═══════════════════ gRPC Adapter ═══════════════════ 33 34 35class GRPCAdapter(ProtocolAdapter): 36 """ 37 gRPC Adapter — translates gRPC calls to Runtime API. 38 39 Proto service definition (zelos.proto): 40 41 service Zelos { 42 rpc SubmitGoal(SubmitGoalRequest) returns (GoalResponse); 43 rpc GetGoalStatus(GetGoalStatusRequest) returns (GoalStatusResponse); 44 rpc CancelGoal(CancelGoalRequest) returns (GoalStatusResponse); 45 rpc WatchGoal(WatchGoalRequest) returns (stream Event); 46 rpc RegisterAgent(RegisterAgentRequest) returns (AgentResponse); 47 rpc AgentHeartbeat(HeartbeatRequest) returns (HeartbeatResponse); 48 rpc SubmitResult(SubmitResultRequest) returns (SubmitResultResponse); 49 rpc GetHealth(Empty) returns (HealthResponse); 50 rpc GetMetrics(Empty) returns (MetricsResponse); 51 } 52 53 Phase 2: Full service handler implementation. For production, 54 use grpcio to create the server with generated stubs. 55 """ 56 57 def __init__(self, runtime=None, host: str = "0.0.0.0", port: int = 50051): 58 super().__init__(runtime) 59 self.host = host 60 self.port = port 61 self._running = False 62 63 def start(self) -> None: 64 self._running = True 65 # In production: grpc.server with generated ZelosServicer 66 # Phase 2 reference: service handlers below 67 68 def stop(self) -> None: 69 self._running = False 70 71 # ── Service Handlers ── 72 73 def SubmitGoal(self, request: dict) -> dict: 74 """gRPC → Runtime API: SubmitGoal.""" 75 if not self.runtime: 76 return {"goal_id": "", "status": "rejected", "reason": "Runtime not connected"} 77 return self.runtime.submit_goal( 78 description=request.get("description", ""), 79 budget=request.get("budget"), 80 deadline=request.get("deadline"), 81 priority=request.get("priority", "medium"), 82 ) 83 84 def GetGoalStatus(self, request: dict) -> dict: 85 if not self.runtime: 86 return {"goal_id": "", "status": "not_found"} 87 return self.runtime.get_goal_status(request["goal_id"]) or {"status": "not_found"} 88 89 def CancelGoal(self, request: dict) -> dict: 90 if not self.runtime: 91 return {} 92 return self.runtime.cancel_goal(request["goal_id"]) or {} 93 94 def RegisterAgent(self, request: dict) -> dict: 95 if not self.runtime: 96 return {"agent_id": "", "status": "rejected"} 97 agent_id = self.runtime.add_agent( 98 name=request.get("name", "grpc-agent"), 99 entrypoint=request.get("entrypoint", ""), 100 capabilities=request.get("capabilities", []), 101 ) 102 return { 103 "agent_id": agent_id, 104 "status": "registered", 105 "heartbeat_interval_ms": 30000, 106 "runtime_version": "0.7.0", 107 } 108 109 def AgentHeartbeat(self, request: dict) -> dict: 110 if not self.runtime: 111 return {"status": "re-register"} 112 ok = self.runtime._execution_engine.heartbeat(request["agent_id"]) 113 return {"status": "ok" if ok else "re-register", "pending_tasks": 0} 114 115 def SubmitResult(self, request: dict) -> dict: 116 if not self.runtime: 117 return {"status": "rejected"} 118 ok = self.runtime._execution_engine.submit_result( 119 request.get("task_id", ""), request.get("agent_id", ""), request.get("result", {}) 120 ) 121 return {"status": "accepted" if ok else "rejected"} 122 123 def GetHealth(self, _=None) -> dict: 124 return self.runtime.get_health() if self.runtime else {"status": "unhealthy"} 125 126 def GetMetrics(self, _=None) -> dict: 127 return self.runtime.get_metrics() if self.runtime else {} 128 129 130# ═══════════════════ WebSocket Adapter ═══════════════════ 131 132 133class WebSocketAdapter(ProtocolAdapter): 134 """ 135 WebSocket Adapter — streams Events to connected clients. 136 137 Supports: 138 - watch_goal(goal_id): stream all events for a Goal 139 - watch_tasks(pattern): stream task.* events 140 - health streaming 141 142 Phase 2: Reference implementation using a simple event subscription model. 143 """ 144 145 def __init__(self, runtime=None): 146 super().__init__(runtime) 147 self._clients: dict[str, dict] = {} # client_id → {subscriptions, queue} 148 self._running = False 149 self._thread: threading.Thread | None = None 150 151 def start(self) -> None: 152 self._running = True 153 if self.runtime: 154 # Subscribe to all major event domains 155 for domain in ["goal", "task", "plan", "agent", "artifact", "verification", "plugin"]: 156 self.runtime._event_bus.subscribe_pattern(f"{domain}.*", self._on_event) 157 158 def stop(self) -> None: 159 self._running = False 160 self._clients.clear() 161 162 def register_client(self, client_id: str) -> None: 163 self._clients[client_id] = {"subscriptions": [], "queue": []} 164 165 def unregister_client(self, client_id: str) -> None: 166 self._clients.pop(client_id, None) 167 168 def watch_goal(self, client_id: str, goal_id: str) -> None: 169 if client_id in self._clients: 170 self._clients[client_id]["subscriptions"].append(f"goal.{goal_id}") 171 172 def get_events(self, client_id: str) -> list[dict]: 173 """Drain event queue for a client.""" 174 if client_id not in self._clients: 175 return [] 176 queue = self._clients[client_id]["queue"] 177 events = list(queue) 178 queue.clear() 179 return events 180 181 def _on_event(self, event) -> Any: 182 """Fan out event to interested clients.""" 183 for _cid, client in self._clients.items(): 184 for sub in client["subscriptions"]: 185 # Simple match: "goal.g-001" matches events with correlation_id 186 parts = sub.split(".", 1) 187 if len(parts) == 2 and event.correlation_id == parts[1]: 188 client["queue"].append(event.to_dict()) 189 elif event.event_type.startswith(sub): 190 client["queue"].append(event.to_dict()) 191 from zelos.event_bus import HandlerResult 192 193 return HandlerResult.ACK 194 195 196# ═══════════════════ MCP Adapter ═══════════════════ 197 198 199class MCPAdapter(ProtocolAdapter): 200 """ 201 MCP (Model Context Protocol) Adapter. 202 203 MCP is how Agents access TOOLS during Task execution — NOT how they 204 communicate with the Runtime. The Runtime is unaware of MCP internals. 205 206 Phase 2: MCP Server lifecycle management + Tool registry. 207 """ 208 209 def __init__(self, runtime=None): 210 super().__init__(runtime) 211 self._tool_registry: dict[str, dict] = {} # tool_name → {server, schema} 212 213 def start(self) -> None: ... 214 215 def stop(self) -> None: ... 216 217 def register_tool(self, name: str, server_endpoint: str, schema: dict) -> None: 218 """Register an MCP tool available to Agents.""" 219 self._tool_registry[name] = {"endpoint": server_endpoint, "schema": schema} 220 221 def list_tools(self) -> list[dict]: 222 """Return all registered MCP tools (MCP tools/list format).""" 223 return [ 224 { 225 "name": name, 226 "description": info["schema"].get("description", ""), 227 "inputSchema": info["schema"].get("input_schema", {}), 228 } 229 for name, info in self._tool_registry.items() 230 ] 231 232 def call_tool(self, name: str, arguments: dict) -> dict: 233 """Invoke an MCP tool (MCP tools/call format).""" 234 tool = self._tool_registry.get(name) 235 if not tool: 236 return {"isError": True, "content": [{"type": "text", "text": f"Tool not found: {name}"}]} 237 # Phase 2: In production, forward to the tool server 238 return {"content": [{"type": "text", "text": f"Tool '{name}' invoked with {arguments}"}]} 239 240 241# ═══════════════════ A2A Adapter ═══════════════════ 242 243 244class A2AAdapter(ProtocolAdapter): 245 """ 246 A2A (Agent-to-Agent Protocol) Adapter. 247 248 A2A is an EXTERNAL interoperability protocol for: 249 - Runtime ↔ Runtime coordination 250 - External Agent integration 251 - Third-party Agent services 252 253 Key translation: A2A Agent Card ↔ Zelos Capability list 254 A2A Task ↔ Zelos Task 255 A2A Message ↔ Zelos Event 256 257 A2A assumes direct agent-to-agent communication. Zelos PROHIBITS this. 258 The A2A Adapter routes ALL A2A messages through the Runtime: 259 A2A: Agent A → Agent B 260 → Zelos: Agent A → Runtime → Agent B 261 """ 262 263 def __init__(self, runtime=None): 264 super().__init__(runtime) 265 self._external_agents: dict[str, dict] = {} 266 267 def start(self) -> None: ... 268 269 def stop(self) -> None: ... 270 271 def generate_agent_card(self, agent_id: str) -> dict: 272 """Generate A2A Agent Card from Zelos Agent Capabilities.""" 273 if not self.runtime: 274 return {} 275 agent = self.runtime.get_agent(agent_id) 276 if not agent: 277 return {} 278 return { 279 "agentId": agent.get("agent_id", ""), 280 "name": agent.get("name", ""), 281 "description": f"Zelos Agent providing {len(agent.get('capabilities', []))} capabilities", 282 "skills": [ 283 {"name": c["name"], "description": c.get("description", "")} for c in agent.get("capabilities", []) 284 ], 285 "endpoint": f"zelos://agents/{agent.get('agent_id', '')}", 286 } 287 288 def receive_external_task(self, task_data: dict) -> str | None: 289 """Receive an A2A Task from an external system → create Zelos Task.""" 290 if not self.runtime: 291 return None 292 goal = self.runtime.submit_goal( 293 description=task_data.get("description", "External A2A task"), 294 priority=task_data.get("priority", "medium"), 295 ) 296 return goal.get("goal_id") 297 298 def register_external_agent(self, agent_card: dict) -> str: 299 """Register an external (non-Zelos) Agent via its A2A Card.""" 300 import uuid 301 302 aid = str(uuid.uuid4()) 303 self._external_agents[aid] = { 304 "card": agent_card, 305 "registered_at": time.time(), 306 } 307 return aid
20class ProtocolAdapter(ABC): 21 """Base for all protocol adapters.""" 22 23 def __init__(self, runtime=None): 24 self.runtime = runtime 25 26 @abstractmethod 27 def start(self) -> None: ... 28 29 @abstractmethod 30 def stop(self) -> None: ...
Base for all protocol adapters.
36class GRPCAdapter(ProtocolAdapter): 37 """ 38 gRPC Adapter — translates gRPC calls to Runtime API. 39 40 Proto service definition (zelos.proto): 41 42 service Zelos { 43 rpc SubmitGoal(SubmitGoalRequest) returns (GoalResponse); 44 rpc GetGoalStatus(GetGoalStatusRequest) returns (GoalStatusResponse); 45 rpc CancelGoal(CancelGoalRequest) returns (GoalStatusResponse); 46 rpc WatchGoal(WatchGoalRequest) returns (stream Event); 47 rpc RegisterAgent(RegisterAgentRequest) returns (AgentResponse); 48 rpc AgentHeartbeat(HeartbeatRequest) returns (HeartbeatResponse); 49 rpc SubmitResult(SubmitResultRequest) returns (SubmitResultResponse); 50 rpc GetHealth(Empty) returns (HealthResponse); 51 rpc GetMetrics(Empty) returns (MetricsResponse); 52 } 53 54 Phase 2: Full service handler implementation. For production, 55 use grpcio to create the server with generated stubs. 56 """ 57 58 def __init__(self, runtime=None, host: str = "0.0.0.0", port: int = 50051): 59 super().__init__(runtime) 60 self.host = host 61 self.port = port 62 self._running = False 63 64 def start(self) -> None: 65 self._running = True 66 # In production: grpc.server with generated ZelosServicer 67 # Phase 2 reference: service handlers below 68 69 def stop(self) -> None: 70 self._running = False 71 72 # ── Service Handlers ── 73 74 def SubmitGoal(self, request: dict) -> dict: 75 """gRPC → Runtime API: SubmitGoal.""" 76 if not self.runtime: 77 return {"goal_id": "", "status": "rejected", "reason": "Runtime not connected"} 78 return self.runtime.submit_goal( 79 description=request.get("description", ""), 80 budget=request.get("budget"), 81 deadline=request.get("deadline"), 82 priority=request.get("priority", "medium"), 83 ) 84 85 def GetGoalStatus(self, request: dict) -> dict: 86 if not self.runtime: 87 return {"goal_id": "", "status": "not_found"} 88 return self.runtime.get_goal_status(request["goal_id"]) or {"status": "not_found"} 89 90 def CancelGoal(self, request: dict) -> dict: 91 if not self.runtime: 92 return {} 93 return self.runtime.cancel_goal(request["goal_id"]) or {} 94 95 def RegisterAgent(self, request: dict) -> dict: 96 if not self.runtime: 97 return {"agent_id": "", "status": "rejected"} 98 agent_id = self.runtime.add_agent( 99 name=request.get("name", "grpc-agent"), 100 entrypoint=request.get("entrypoint", ""), 101 capabilities=request.get("capabilities", []), 102 ) 103 return { 104 "agent_id": agent_id, 105 "status": "registered", 106 "heartbeat_interval_ms": 30000, 107 "runtime_version": "0.7.0", 108 } 109 110 def AgentHeartbeat(self, request: dict) -> dict: 111 if not self.runtime: 112 return {"status": "re-register"} 113 ok = self.runtime._execution_engine.heartbeat(request["agent_id"]) 114 return {"status": "ok" if ok else "re-register", "pending_tasks": 0} 115 116 def SubmitResult(self, request: dict) -> dict: 117 if not self.runtime: 118 return {"status": "rejected"} 119 ok = self.runtime._execution_engine.submit_result( 120 request.get("task_id", ""), request.get("agent_id", ""), request.get("result", {}) 121 ) 122 return {"status": "accepted" if ok else "rejected"} 123 124 def GetHealth(self, _=None) -> dict: 125 return self.runtime.get_health() if self.runtime else {"status": "unhealthy"} 126 127 def GetMetrics(self, _=None) -> dict: 128 return self.runtime.get_metrics() if self.runtime else {}
gRPC Adapter — translates gRPC calls to Runtime API.
Proto service definition (zelos.proto):
service Zelos { rpc SubmitGoal(SubmitGoalRequest) returns (GoalResponse); rpc GetGoalStatus(GetGoalStatusRequest) returns (GoalStatusResponse); rpc CancelGoal(CancelGoalRequest) returns (GoalStatusResponse); rpc WatchGoal(WatchGoalRequest) returns (stream Event); rpc RegisterAgent(RegisterAgentRequest) returns (AgentResponse); rpc AgentHeartbeat(HeartbeatRequest) returns (HeartbeatResponse); rpc SubmitResult(SubmitResultRequest) returns (SubmitResultResponse); rpc GetHealth(Empty) returns (HealthResponse); rpc GetMetrics(Empty) returns (MetricsResponse); }
Phase 2: Full service handler implementation. For production, use grpcio to create the server with generated stubs.
74 def SubmitGoal(self, request: dict) -> dict: 75 """gRPC → Runtime API: SubmitGoal.""" 76 if not self.runtime: 77 return {"goal_id": "", "status": "rejected", "reason": "Runtime not connected"} 78 return self.runtime.submit_goal( 79 description=request.get("description", ""), 80 budget=request.get("budget"), 81 deadline=request.get("deadline"), 82 priority=request.get("priority", "medium"), 83 )
gRPC → Runtime API: SubmitGoal.
95 def RegisterAgent(self, request: dict) -> dict: 96 if not self.runtime: 97 return {"agent_id": "", "status": "rejected"} 98 agent_id = self.runtime.add_agent( 99 name=request.get("name", "grpc-agent"), 100 entrypoint=request.get("entrypoint", ""), 101 capabilities=request.get("capabilities", []), 102 ) 103 return { 104 "agent_id": agent_id, 105 "status": "registered", 106 "heartbeat_interval_ms": 30000, 107 "runtime_version": "0.7.0", 108 }
116 def SubmitResult(self, request: dict) -> dict: 117 if not self.runtime: 118 return {"status": "rejected"} 119 ok = self.runtime._execution_engine.submit_result( 120 request.get("task_id", ""), request.get("agent_id", ""), request.get("result", {}) 121 ) 122 return {"status": "accepted" if ok else "rejected"}
Inherited Members
134class WebSocketAdapter(ProtocolAdapter): 135 """ 136 WebSocket Adapter — streams Events to connected clients. 137 138 Supports: 139 - watch_goal(goal_id): stream all events for a Goal 140 - watch_tasks(pattern): stream task.* events 141 - health streaming 142 143 Phase 2: Reference implementation using a simple event subscription model. 144 """ 145 146 def __init__(self, runtime=None): 147 super().__init__(runtime) 148 self._clients: dict[str, dict] = {} # client_id → {subscriptions, queue} 149 self._running = False 150 self._thread: threading.Thread | None = None 151 152 def start(self) -> None: 153 self._running = True 154 if self.runtime: 155 # Subscribe to all major event domains 156 for domain in ["goal", "task", "plan", "agent", "artifact", "verification", "plugin"]: 157 self.runtime._event_bus.subscribe_pattern(f"{domain}.*", self._on_event) 158 159 def stop(self) -> None: 160 self._running = False 161 self._clients.clear() 162 163 def register_client(self, client_id: str) -> None: 164 self._clients[client_id] = {"subscriptions": [], "queue": []} 165 166 def unregister_client(self, client_id: str) -> None: 167 self._clients.pop(client_id, None) 168 169 def watch_goal(self, client_id: str, goal_id: str) -> None: 170 if client_id in self._clients: 171 self._clients[client_id]["subscriptions"].append(f"goal.{goal_id}") 172 173 def get_events(self, client_id: str) -> list[dict]: 174 """Drain event queue for a client.""" 175 if client_id not in self._clients: 176 return [] 177 queue = self._clients[client_id]["queue"] 178 events = list(queue) 179 queue.clear() 180 return events 181 182 def _on_event(self, event) -> Any: 183 """Fan out event to interested clients.""" 184 for _cid, client in self._clients.items(): 185 for sub in client["subscriptions"]: 186 # Simple match: "goal.g-001" matches events with correlation_id 187 parts = sub.split(".", 1) 188 if len(parts) == 2 and event.correlation_id == parts[1]: 189 client["queue"].append(event.to_dict()) 190 elif event.event_type.startswith(sub): 191 client["queue"].append(event.to_dict()) 192 from zelos.event_bus import HandlerResult 193 194 return HandlerResult.ACK
WebSocket Adapter — streams Events to connected clients.
Supports:
- watch_goal(goal_id): stream all events for a Goal
- watch_tasks(pattern): stream task.* events
- health streaming
Phase 2: Reference implementation using a simple event subscription model.
173 def get_events(self, client_id: str) -> list[dict]: 174 """Drain event queue for a client.""" 175 if client_id not in self._clients: 176 return [] 177 queue = self._clients[client_id]["queue"] 178 events = list(queue) 179 queue.clear() 180 return events
Drain event queue for a client.
Inherited Members
200class MCPAdapter(ProtocolAdapter): 201 """ 202 MCP (Model Context Protocol) Adapter. 203 204 MCP is how Agents access TOOLS during Task execution — NOT how they 205 communicate with the Runtime. The Runtime is unaware of MCP internals. 206 207 Phase 2: MCP Server lifecycle management + Tool registry. 208 """ 209 210 def __init__(self, runtime=None): 211 super().__init__(runtime) 212 self._tool_registry: dict[str, dict] = {} # tool_name → {server, schema} 213 214 def start(self) -> None: ... 215 216 def stop(self) -> None: ... 217 218 def register_tool(self, name: str, server_endpoint: str, schema: dict) -> None: 219 """Register an MCP tool available to Agents.""" 220 self._tool_registry[name] = {"endpoint": server_endpoint, "schema": schema} 221 222 def list_tools(self) -> list[dict]: 223 """Return all registered MCP tools (MCP tools/list format).""" 224 return [ 225 { 226 "name": name, 227 "description": info["schema"].get("description", ""), 228 "inputSchema": info["schema"].get("input_schema", {}), 229 } 230 for name, info in self._tool_registry.items() 231 ] 232 233 def call_tool(self, name: str, arguments: dict) -> dict: 234 """Invoke an MCP tool (MCP tools/call format).""" 235 tool = self._tool_registry.get(name) 236 if not tool: 237 return {"isError": True, "content": [{"type": "text", "text": f"Tool not found: {name}"}]} 238 # Phase 2: In production, forward to the tool server 239 return {"content": [{"type": "text", "text": f"Tool '{name}' invoked with {arguments}"}]}
MCP (Model Context Protocol) Adapter.
MCP is how Agents access TOOLS during Task execution — NOT how they communicate with the Runtime. The Runtime is unaware of MCP internals.
Phase 2: MCP Server lifecycle management + Tool registry.
218 def register_tool(self, name: str, server_endpoint: str, schema: dict) -> None: 219 """Register an MCP tool available to Agents.""" 220 self._tool_registry[name] = {"endpoint": server_endpoint, "schema": schema}
Register an MCP tool available to Agents.
222 def list_tools(self) -> list[dict]: 223 """Return all registered MCP tools (MCP tools/list format).""" 224 return [ 225 { 226 "name": name, 227 "description": info["schema"].get("description", ""), 228 "inputSchema": info["schema"].get("input_schema", {}), 229 } 230 for name, info in self._tool_registry.items() 231 ]
Return all registered MCP tools (MCP tools/list format).
233 def call_tool(self, name: str, arguments: dict) -> dict: 234 """Invoke an MCP tool (MCP tools/call format).""" 235 tool = self._tool_registry.get(name) 236 if not tool: 237 return {"isError": True, "content": [{"type": "text", "text": f"Tool not found: {name}"}]} 238 # Phase 2: In production, forward to the tool server 239 return {"content": [{"type": "text", "text": f"Tool '{name}' invoked with {arguments}"}]}
Invoke an MCP tool (MCP tools/call format).
Inherited Members
245class A2AAdapter(ProtocolAdapter): 246 """ 247 A2A (Agent-to-Agent Protocol) Adapter. 248 249 A2A is an EXTERNAL interoperability protocol for: 250 - Runtime ↔ Runtime coordination 251 - External Agent integration 252 - Third-party Agent services 253 254 Key translation: A2A Agent Card ↔ Zelos Capability list 255 A2A Task ↔ Zelos Task 256 A2A Message ↔ Zelos Event 257 258 A2A assumes direct agent-to-agent communication. Zelos PROHIBITS this. 259 The A2A Adapter routes ALL A2A messages through the Runtime: 260 A2A: Agent A → Agent B 261 → Zelos: Agent A → Runtime → Agent B 262 """ 263 264 def __init__(self, runtime=None): 265 super().__init__(runtime) 266 self._external_agents: dict[str, dict] = {} 267 268 def start(self) -> None: ... 269 270 def stop(self) -> None: ... 271 272 def generate_agent_card(self, agent_id: str) -> dict: 273 """Generate A2A Agent Card from Zelos Agent Capabilities.""" 274 if not self.runtime: 275 return {} 276 agent = self.runtime.get_agent(agent_id) 277 if not agent: 278 return {} 279 return { 280 "agentId": agent.get("agent_id", ""), 281 "name": agent.get("name", ""), 282 "description": f"Zelos Agent providing {len(agent.get('capabilities', []))} capabilities", 283 "skills": [ 284 {"name": c["name"], "description": c.get("description", "")} for c in agent.get("capabilities", []) 285 ], 286 "endpoint": f"zelos://agents/{agent.get('agent_id', '')}", 287 } 288 289 def receive_external_task(self, task_data: dict) -> str | None: 290 """Receive an A2A Task from an external system → create Zelos Task.""" 291 if not self.runtime: 292 return None 293 goal = self.runtime.submit_goal( 294 description=task_data.get("description", "External A2A task"), 295 priority=task_data.get("priority", "medium"), 296 ) 297 return goal.get("goal_id") 298 299 def register_external_agent(self, agent_card: dict) -> str: 300 """Register an external (non-Zelos) Agent via its A2A Card.""" 301 import uuid 302 303 aid = str(uuid.uuid4()) 304 self._external_agents[aid] = { 305 "card": agent_card, 306 "registered_at": time.time(), 307 } 308 return aid
A2A (Agent-to-Agent Protocol) Adapter.
A2A is an EXTERNAL interoperability protocol for:
- Runtime ↔ Runtime coordination
- External Agent integration
- Third-party Agent services
Key translation: A2A Agent Card ↔ Zelos Capability list A2A Task ↔ Zelos Task A2A Message ↔ Zelos Event
A2A assumes direct agent-to-agent communication. Zelos PROHIBITS this. The A2A Adapter routes ALL A2A messages through the Runtime: A2A: Agent A → Agent B → Zelos: Agent A → Runtime → Agent B
272 def generate_agent_card(self, agent_id: str) -> dict: 273 """Generate A2A Agent Card from Zelos Agent Capabilities.""" 274 if not self.runtime: 275 return {} 276 agent = self.runtime.get_agent(agent_id) 277 if not agent: 278 return {} 279 return { 280 "agentId": agent.get("agent_id", ""), 281 "name": agent.get("name", ""), 282 "description": f"Zelos Agent providing {len(agent.get('capabilities', []))} capabilities", 283 "skills": [ 284 {"name": c["name"], "description": c.get("description", "")} for c in agent.get("capabilities", []) 285 ], 286 "endpoint": f"zelos://agents/{agent.get('agent_id', '')}", 287 }
Generate A2A Agent Card from Zelos Agent Capabilities.
289 def receive_external_task(self, task_data: dict) -> str | None: 290 """Receive an A2A Task from an external system → create Zelos Task.""" 291 if not self.runtime: 292 return None 293 goal = self.runtime.submit_goal( 294 description=task_data.get("description", "External A2A task"), 295 priority=task_data.get("priority", "medium"), 296 ) 297 return goal.get("goal_id")
Receive an A2A Task from an external system → create Zelos Task.
299 def register_external_agent(self, agent_card: dict) -> str: 300 """Register an external (non-Zelos) Agent via its A2A Card.""" 301 import uuid 302 303 aid = str(uuid.uuid4()) 304 self._external_agents[aid] = { 305 "card": agent_card, 306 "registered_at": time.time(), 307 } 308 return aid
Register an external (non-Zelos) Agent via its A2A Card.