anu151105/agentic-browser
2
1"""2Agent Communication module for the A2A Protocol Layer.3 4This module implements the Agent-to-Agent (A2A) protocol for enabling5communication and collaboration between multiple AI agents.6"""7 8import asyncio9import json10import logging11import time12import uuid13from typing import Dict, List, Any, Optional, Union, Callable14 15import aiohttp16import httpx17 18# Configure logging19logging.basicConfig(level=logging.INFO)20logger = logging.getLogger(__name__)21 22class A2AProtocol:23 """24 Implements the Agent-to-Agent (A2A) protocol.25 26 This class manages communication between AI agents, enabling task delegation,27 coordination, and information exchange for complex multi-agent workflows.28 """29 30 def __init__(self):31 """Initialize the A2AProtocol."""32 self.agent_id = str(uuid.uuid4())33 self.agent_name = "Browser Agent"34 self.agent_capabilities = [35 "web_browsing",36 "form_filling",37 "data_extraction",38 "api_calling",39 "content_analysis"40 ]41 42 self.known_agents = {} # agent_id -> agent_info43 self.active_conversations = {} # conversation_id -> conversation_data44 self.pending_requests = {} # request_id -> future45 46 self.session = None47 self.message_handlers = []48 49 logger.info("A2AProtocol instance created")50 51 async def initialize(self, agent_name: str = None, capabilities: List[str] = None):52 """53 Initialize the A2A protocol.54 55 Args:56 agent_name: Optional custom name for this agent57 capabilities: Optional list of agent capabilities58 59 Returns:60 bool: True if initialization was successful61 """62 try:63 self.session = aiohttp.ClientSession(64 headers={65 "Content-Type": "application/json",66 "User-Agent": f"A2A-Protocol/2025 {self.agent_name}"67 }68 )69 70 if agent_name:71 self.agent_name = agent_name72 73 if capabilities:74 self.agent_capabilities = capabilities75 76 # Register default message handlers77 self._register_default_handlers()78 79 logger.info(f"A2A Protocol initialized for agent {self.agent_name} ({self.agent_id})")80 return True81 82 except Exception as e:83 logger.error(f"Error initializing A2A Protocol: {str(e)}")84 return False85 86 async def register_with_directory(self, directory_url: str) -> bool:87 """88 Register this agent with a central agent directory.89 90 Args:91 directory_url: URL of the agent directory service92 93 Returns:94 bool: True if registration was successful95 """96 try:97 registration_data = {98 "agent_id": self.agent_id,99 "agent_name": self.agent_name,100 "capabilities": self.agent_capabilities,101 "status": "online",102 "version": "1.0.0",103 "protocol_version": "A2A-2025",104 "timestamp": time.time()105 }106 107 async with self.session.post(108 f"{directory_url}/register",109 json=registration_data110 ) as response:111 if response.status == 200:112 data = await response.json()113 logger.info(f"Successfully registered with agent directory: {data.get('message', 'OK')}")114 return True115 else:116 logger.error(f"Failed to register with agent directory. Status: {response.status}")117 return False118 119 except Exception as e:120 logger.error(f"Error registering with agent directory: {str(e)}")121 return False122 123 async def discover_agents(self, directory_url: str, capabilities: List[str] = None) -> List[Dict]:124 """125 Discover other agents with specific capabilities.126 127 Args:128 directory_url: URL of the agent directory service129 capabilities: Optional list of required capabilities to filter by130 131 Returns:132 List[Dict]: List of discovered agents133 """134 try:135 params = {}136 if capabilities:137 params["capabilities"] = ",".join(capabilities)138 139 async with self.session.get(140 f"{directory_url}/discover",141 params=params142 ) as response:143 if response.status == 200:144 data = await response.json()145 agents = data.get("agents", [])146 147 # Update known agents148 for agent in agents:149 agent_id = agent.get("agent_id")150 if agent_id and agent_id != self.agent_id:151 self.known_agents[agent_id] = agent152 153 logger.info(f"Discovered {len(agents)} agents")154 return agents155 else:156 logger.error(f"Failed to discover agents. Status: {response.status}")157 return []158 159 except Exception as e:160 logger.error(f"Error discovering agents: {str(e)}")161 return []162 163 async def send_message(self, recipient_id: str, message_type: str, content: Dict, conversation_id: str = None) -> Dict:164 """165 Send a message to another agent.166 167 Args:168 recipient_id: ID of the recipient agent169 message_type: Type of message (e.g., 'request', 'response', 'update')170 content: Message content171 conversation_id: Optional ID for an ongoing conversation172 173 Returns:174 Dict: Message receipt confirmation175 """176 try:177 # Create a new conversation if not provided178 if not conversation_id:179 conversation_id = str(uuid.uuid4())180 self.active_conversations[conversation_id] = {181 "start_time": time.time(),182 "participants": [self.agent_id, recipient_id],183 "messages": []184 }185 186 # Create message187 message = {188 "message_id": str(uuid.uuid4()),189 "conversation_id": conversation_id,190 "sender_id": self.agent_id,191 "sender_name": self.agent_name,192 "recipient_id": recipient_id,193 "message_type": message_type,194 "content": content,195 "timestamp": time.time(),196 "protocol_version": "A2A-2025"197 }198 199 # Store in conversation history200 if conversation_id in self.active_conversations:201 self.active_conversations[conversation_id]["messages"].append(message)202 203 # Get recipient endpoint204 recipient = self.known_agents.get(recipient_id)205 if not recipient or "endpoint" not in recipient:206 raise ValueError(f"Unknown recipient or missing endpoint: {recipient_id}")207 208 # Send message to recipient209 async with self.session.post(210 f"{recipient['endpoint']}/receive",211 json=message212 ) as response:213 if response.status == 200:214 data = await response.json()215 logger.info(f"Message sent successfully to {recipient_id}: {message_type}")216 return {217 "success": True,218 "message_id": message["message_id"],219 "conversation_id": conversation_id,220 "receipt": data221 }222 else:223 error_text = await response.text()224 logger.error(f"Failed to send message to {recipient_id}. Status: {response.status}, Error: {error_text}")225 return {226 "success": False,227 "error": f"Failed to send message. Status: {response.status}"228 }229 230 except Exception as e:231 logger.error(f"Error sending message to {recipient_id}: {str(e)}")232 return {233 "success": False,234 "error": str(e)235 }236 237 async def request_task(self, recipient_id: str, task_description: str, parameters: Dict = None, timeout: int = 60) -> Dict:238 """239 Request another agent to perform a task.240 241 Args:242 recipient_id: ID of the recipient agent243 task_description: Description of the requested task244 parameters: Optional parameters for the task245 timeout: Timeout in seconds for the request246 247 Returns:248 Dict: Task result or error249 """250 try:251 # Create a unique request ID252 request_id = str(uuid.uuid4())253 254 # Create a future for the response255 future = asyncio.get_event_loop().create_future()256 self.pending_requests[request_id] = future257 258 # Create the task request message259 content = {260 "request_id": request_id,261 "task_description": task_description,262 "parameters": parameters or {},263 "response_required": True,264 "timeout_seconds": timeout265 }266 267 # Send the request268 result = await self.send_message(recipient_id, "task_request", content)269 270 if not result.get("success", False):271 # Failed to send request272 if request_id in self.pending_requests:273 del self.pending_requests[request_id]274 return {275 "success": False,276 "error": result.get("error", "Failed to send task request")277 }278 279 # Wait for response with timeout280 try:281 response = await asyncio.wait_for(future, timeout=timeout)282 return {283 "success": True,284 "request_id": request_id,285 "response": response286 }287 288 except asyncio.TimeoutError:289 # Request timed out290 if request_id in self.pending_requests:291 del self.pending_requests[request_id]292 return {293 "success": False,294 "error": f"Request timed out after {timeout} seconds"295 }296 297 except Exception as e:298 logger.error(f"Error requesting task from {recipient_id}: {str(e)}")299 return {300 "success": False,301 "error": str(e)302 }303 304 async def handle_incoming_message(self, message: Dict) -> Dict:305 """306 Handle an incoming message from another agent.307 308 Args:309 message: Incoming message data310 311 Returns:312 Dict: Response data313 """314 try:315 # Validate message format316 if not self._validate_message(message):317 return {318 "success": False,319 "error": "Invalid message format"320 }321 322 # Store in conversation history323 conversation_id = message.get("conversation_id")324 if conversation_id not in self.active_conversations:325 self.active_conversations[conversation_id] = {326 "start_time": time.time(),327 "participants": [message.get("sender_id"), self.agent_id],328 "messages": []329 }330 331 self.active_conversations[conversation_id]["messages"].append(message)332 333 # Process message based on type334 message_type = message.get("message_type", "").lower()335 content = message.get("content", {})336 337 # Find handler for message type338 for handler in self.message_handlers:339 if handler["type"] == message_type:340 return await handler["handler"](message, content)341 342 # Default handler if no specific handler found343 logger.warning(f"No handler for message type: {message_type}")344 return {345 "success": True,346 "message": "Message received but not processed",347 "handled": False348 }349 350 except Exception as e:351 logger.error(f"Error handling incoming message: {str(e)}")352 return {353 "success": False,354 "error": str(e)355 }356 357 async def respond_to_task(self, request_id: str, result: Dict, conversation_id: str, recipient_id: str) -> Dict:358 """359 Respond to a task request from another agent.360 361 Args:362 request_id: Original request ID363 result: Task execution result364 conversation_id: Ongoing conversation ID365 recipient_id: ID of the requesting agent366 367 Returns:368 Dict: Response status369 """370 content = {371 "request_id": request_id,372 "result": result,373 "timestamp": time.time()374 }375 376 return await self.send_message(recipient_id, "task_response", content, conversation_id)377 378 def register_message_handler(self, message_type: str, handler: Callable[[Dict, Dict], Dict]) -> bool:379 """380 Register a handler for a specific message type.381 382 Args:383 message_type: Type of messages to handle384 handler: Async function to handle messages of this type385 386 Returns:387 bool: True if registration was successful388 """389 # Check for existing handler of the same type390 for existing in self.message_handlers:391 if existing["type"] == message_type:392 # Update existing handler393 existing["handler"] = handler394 logger.info(f"Updated handler for message type: {message_type}")395 return True396 397 # Add new handler398 self.message_handlers.append({399 "type": message_type,400 "handler": handler401 })402 403 logger.info(f"Registered new handler for message type: {message_type}")404 return True405 406 def _validate_message(self, message: Dict) -> bool:407 """408 Validate the format of an incoming message.409 410 Args:411 message: Message to validate412 413 Returns:414 bool: True if valid, False otherwise415 """416 required_fields = ["message_id", "conversation_id", "sender_id", "message_type", "content"]417 418 for field in required_fields:419 if field not in message:420 logger.error(f"Missing required field in message: {field}")421 return False422 423 return True424 425 async def _handle_task_request(self, message: Dict, content: Dict) -> Dict:426 """427 Handle a task request from another agent.428 429 Args:430 message: Full message data431 content: Message content432 433 Returns:434 Dict: Response data435 """436 # This would typically delegate to the appropriate handler based on task type437 # For now, we'll just acknowledge receipt438 return {439 "success": True,440 "message": "Task request acknowledged",441 "request_id": content.get("request_id")442 }443 444 async def _handle_task_response(self, message: Dict, content: Dict) -> Dict:445 """446 Handle a task response from another agent.447 448 Args:449 message: Full message data450 content: Message content451 452 Returns:453 Dict: Response data454 """455 request_id = content.get("request_id")456 457 # Check if we have a pending request with this ID458 if request_id in self.pending_requests:459 future = self.pending_requests[request_id]460 if not future.done():461 # Resolve the future with the response462 future.set_result(content.get("result"))463 464 # Clean up465 del self.pending_requests[request_id]466 467 return {468 "success": True,469 "message": "Task response processed"470 }471 else:472 logger.warning(f"Received response for unknown request: {request_id}")473 return {474 "success": True,475 "message": "No pending request found for this ID",476 "request_id": request_id477 }478 479 def _register_default_handlers(self):480 """Register default message handlers."""481 self.register_message_handler("task_request", self._handle_task_request)482 self.register_message_handler("task_response", self._handle_task_response)483 484 # Other default handlers could be registered here485 486 async def get_conversation_history(self, conversation_id: str) -> Dict:487 """488 Get the history of a conversation.489 490 Args:491 conversation_id: ID of the conversation492 493 Returns:494 Dict: Conversation data495 """496 if conversation_id in self.active_conversations:497 return {498 "success": True,499 "conversation": self.active_conversations[conversation_id]500 }501 else:502 return {503 "success": False,504 "error": f"Conversation not found: {conversation_id}"505 }506 507 async def shutdown(self):508 """Clean up resources."""509 if self.session:510 await self.session.close()511 512 logger.info("A2A Protocol resources cleaned up")513 