Team Ai
Apppublic

anu151105/agentic-browser

sourceHugging Facemitupdated 1y agoView on Hugging Face
2likes
agent_communication.py513 linesDownload Raw Back to a2a_protocol
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