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
class ProtocolAdapter(abc.ABC):
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.

runtime
@abstractmethod
def start(self) -> None:
26    @abstractmethod
27    def start(self) -> None: ...
@abstractmethod
def stop(self) -> None:
29    @abstractmethod
30    def stop(self) -> None: ...
class GRPCAdapter(ProtocolAdapter):
 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.

GRPCAdapter(runtime=None, host: str = '0.0.0.0', port: int = 50051)
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
host
port
def start(self) -> None:
64    def start(self) -> None:
65        self._running = True
66        # In production: grpc.server with generated ZelosServicer
67        # Phase 2 reference: service handlers below
def stop(self) -> None:
69    def stop(self) -> None:
70        self._running = False
def SubmitGoal(self, request: dict) -> dict:
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.

def GetGoalStatus(self, request: dict) -> dict:
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"}
def CancelGoal(self, request: dict) -> dict:
90    def CancelGoal(self, request: dict) -> dict:
91        if not self.runtime:
92            return {}
93        return self.runtime.cancel_goal(request["goal_id"]) or {}
def RegisterAgent(self, request: dict) -> dict:
 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        }
def AgentHeartbeat(self, request: dict) -> dict:
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}
def SubmitResult(self, request: dict) -> dict:
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"}
def GetHealth(self, _=None) -> dict:
124    def GetHealth(self, _=None) -> dict:
125        return self.runtime.get_health() if self.runtime else {"status": "unhealthy"}
def GetMetrics(self, _=None) -> dict:
127    def GetMetrics(self, _=None) -> dict:
128        return self.runtime.get_metrics() if self.runtime else {}
Inherited Members
ProtocolAdapter
runtime
class WebSocketAdapter(ProtocolAdapter):
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.

WebSocketAdapter(runtime=None)
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
def start(self) -> None:
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)
def stop(self) -> None:
159    def stop(self) -> None:
160        self._running = False
161        self._clients.clear()
def register_client(self, client_id: str) -> None:
163    def register_client(self, client_id: str) -> None:
164        self._clients[client_id] = {"subscriptions": [], "queue": []}
def unregister_client(self, client_id: str) -> None:
166    def unregister_client(self, client_id: str) -> None:
167        self._clients.pop(client_id, None)
def watch_goal(self, client_id: str, goal_id: str) -> None:
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}")
def get_events(self, client_id: str) -> list[dict]:
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
ProtocolAdapter
runtime
class MCPAdapter(ProtocolAdapter):
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.

MCPAdapter(runtime=None)
210    def __init__(self, runtime=None):
211        super().__init__(runtime)
212        self._tool_registry: dict[str, dict] = {}  # tool_name → {server, schema}
def start(self) -> None:
214    def start(self) -> None: ...
def stop(self) -> None:
216    def stop(self) -> None: ...
def register_tool(self, name: str, server_endpoint: str, schema: dict) -> None:
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.

def list_tools(self) -> list[dict]:
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).

def call_tool(self, name: str, arguments: dict) -> dict:
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
ProtocolAdapter
runtime
class A2AAdapter(ProtocolAdapter):
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

A2AAdapter(runtime=None)
264    def __init__(self, runtime=None):
265        super().__init__(runtime)
266        self._external_agents: dict[str, dict] = {}
def start(self) -> None:
268    def start(self) -> None: ...
def stop(self) -> None:
270    def stop(self) -> None: ...
def generate_agent_card(self, agent_id: str) -> dict:
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.

def receive_external_task(self, task_data: dict) -> str | None:
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.

def register_external_agent(self, agent_card: dict) -> str:
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.

Inherited Members
ProtocolAdapter
runtime