tiphys.ipc.server
IPC Server implementation.
Manages IPC connections and routes messages between host and clients.
1""" 2IPC Server implementation. 3 4Manages IPC connections and routes messages between host and clients. 5""" 6 7import asyncio 8import contextlib 9import json 10import os 11from collections.abc import Callable 12from typing import Any 13 14import structlog 15 16from tiphys.ipc.protocol import IPCRequest, IPCResponse, MessageType 17 18logger = structlog.get_logger(__name__) 19 20 21class IPCServer: 22 """ 23 Server for handling local IPC connections via Unix domain sockets. 24 """ 25 26 def __init__(self, socket_path: str = "/tmp/tiphys.sock"): 27 self.socket_path = socket_path 28 self._server: asyncio.AbstractServer | None = None 29 self._handlers: dict[str, Callable[..., Any]] = {} 30 self._running = False 31 32 async def start(self) -> None: 33 """Start the IPC server.""" 34 if self._running: 35 return 36 37 # Clean up old socket 38 if os.path.exists(self.socket_path): 39 os.remove(self.socket_path) 40 41 self._server = await asyncio.start_unix_server(self._handle_client, path=self.socket_path) 42 self._running = True 43 logger.info("IPC Server started", path=self.socket_path) 44 45 async def stop(self) -> None: 46 """Stop the IPC server.""" 47 self._running = False 48 if self._server: 49 self._server.close() 50 await self._server.wait_closed() 51 self._server = None 52 53 if os.path.exists(self.socket_path): 54 with contextlib.suppress(OSError): 55 os.remove(self.socket_path) 56 logger.info("IPC Server stopped") 57 58 def register_handler(self, method: str, handler: Callable[..., Any]) -> None: 59 """Register a handler for a specific method.""" 60 self._handlers[method] = handler 61 62 async def _handle_client( 63 self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter 64 ) -> None: 65 """Handle individual client connection.""" 66 peer = writer.get_extra_info("peername") 67 logger.debug("IPC Client connected", peer=peer) 68 69 try: 70 while self._running: 71 data = await reader.readline() 72 if not data: 73 break 74 75 try: 76 message_dict = json.loads(data.decode().strip()) 77 await self._process_message(message_dict, writer) 78 except json.JSONDecodeError: 79 logger.warning("Invalid JSON received") 80 except Exception as e: 81 logger.error("Error processing message", error=str(e)) 82 83 except Exception as e: 84 logger.error("IPC Connection error", error=str(e)) 85 finally: 86 writer.close() 87 await writer.wait_closed() 88 89 async def _process_message(self, data: dict[str, Any], writer: asyncio.StreamWriter) -> None: 90 """Process an incoming IPC message.""" 91 try: 92 # Determine message type and validate 93 msg_type = data.get("type") 94 95 if msg_type == MessageType.REQUEST: 96 request = IPCRequest(**data) 97 response = await self._handle_request(request) 98 99 # Send response 100 response_data = response.model_dump_json() + "\n" 101 writer.write(response_data.encode()) 102 await writer.drain() 103 104 elif msg_type == MessageType.EVENT: 105 # Handle event (fire and forget) 106 # TODO: Implement event bus routing 107 pass 108 109 except Exception as e: 110 logger.error("Failed to process IPC message", error=str(e)) 111 # Try to send error response if possible 112 if data.get("type") == MessageType.REQUEST and "id" in data: 113 error_response = IPCResponse( 114 request_id=data["id"], type=MessageType.ERROR, error=str(e) 115 ) 116 writer.write((error_response.model_dump_json() + "\n").encode()) 117 await writer.drain() 118 119 async def _handle_request(self, request: IPCRequest) -> IPCResponse: 120 """Execute the requested method.""" 121 handler = self._handlers.get(request.method) 122 if not handler: 123 return IPCResponse(request_id=request.id, error=f"Method '{request.method}' not found") 124 125 try: 126 # Call handler (support both sync and async) 127 if asyncio.iscoroutinefunction(handler): 128 result = await handler(**request.params) 129 else: 130 result = handler(**request.params) 131 132 return IPCResponse(request_id=request.id, result=result) 133 except Exception as e: 134 logger.error("Handler error", method=request.method, error=str(e)) 135 return IPCResponse(request_id=request.id, error=str(e))
logger =
<BoundLoggerLazyProxy(logger=None, wrapper_class=None, processors=None, context_class=None, initial_values={}, logger_factory_args=('tiphys.ipc.server',))>
class
IPCServer:
22class IPCServer: 23 """ 24 Server for handling local IPC connections via Unix domain sockets. 25 """ 26 27 def __init__(self, socket_path: str = "/tmp/tiphys.sock"): 28 self.socket_path = socket_path 29 self._server: asyncio.AbstractServer | None = None 30 self._handlers: dict[str, Callable[..., Any]] = {} 31 self._running = False 32 33 async def start(self) -> None: 34 """Start the IPC server.""" 35 if self._running: 36 return 37 38 # Clean up old socket 39 if os.path.exists(self.socket_path): 40 os.remove(self.socket_path) 41 42 self._server = await asyncio.start_unix_server(self._handle_client, path=self.socket_path) 43 self._running = True 44 logger.info("IPC Server started", path=self.socket_path) 45 46 async def stop(self) -> None: 47 """Stop the IPC server.""" 48 self._running = False 49 if self._server: 50 self._server.close() 51 await self._server.wait_closed() 52 self._server = None 53 54 if os.path.exists(self.socket_path): 55 with contextlib.suppress(OSError): 56 os.remove(self.socket_path) 57 logger.info("IPC Server stopped") 58 59 def register_handler(self, method: str, handler: Callable[..., Any]) -> None: 60 """Register a handler for a specific method.""" 61 self._handlers[method] = handler 62 63 async def _handle_client( 64 self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter 65 ) -> None: 66 """Handle individual client connection.""" 67 peer = writer.get_extra_info("peername") 68 logger.debug("IPC Client connected", peer=peer) 69 70 try: 71 while self._running: 72 data = await reader.readline() 73 if not data: 74 break 75 76 try: 77 message_dict = json.loads(data.decode().strip()) 78 await self._process_message(message_dict, writer) 79 except json.JSONDecodeError: 80 logger.warning("Invalid JSON received") 81 except Exception as e: 82 logger.error("Error processing message", error=str(e)) 83 84 except Exception as e: 85 logger.error("IPC Connection error", error=str(e)) 86 finally: 87 writer.close() 88 await writer.wait_closed() 89 90 async def _process_message(self, data: dict[str, Any], writer: asyncio.StreamWriter) -> None: 91 """Process an incoming IPC message.""" 92 try: 93 # Determine message type and validate 94 msg_type = data.get("type") 95 96 if msg_type == MessageType.REQUEST: 97 request = IPCRequest(**data) 98 response = await self._handle_request(request) 99 100 # Send response 101 response_data = response.model_dump_json() + "\n" 102 writer.write(response_data.encode()) 103 await writer.drain() 104 105 elif msg_type == MessageType.EVENT: 106 # Handle event (fire and forget) 107 # TODO: Implement event bus routing 108 pass 109 110 except Exception as e: 111 logger.error("Failed to process IPC message", error=str(e)) 112 # Try to send error response if possible 113 if data.get("type") == MessageType.REQUEST and "id" in data: 114 error_response = IPCResponse( 115 request_id=data["id"], type=MessageType.ERROR, error=str(e) 116 ) 117 writer.write((error_response.model_dump_json() + "\n").encode()) 118 await writer.drain() 119 120 async def _handle_request(self, request: IPCRequest) -> IPCResponse: 121 """Execute the requested method.""" 122 handler = self._handlers.get(request.method) 123 if not handler: 124 return IPCResponse(request_id=request.id, error=f"Method '{request.method}' not found") 125 126 try: 127 # Call handler (support both sync and async) 128 if asyncio.iscoroutinefunction(handler): 129 result = await handler(**request.params) 130 else: 131 result = handler(**request.params) 132 133 return IPCResponse(request_id=request.id, result=result) 134 except Exception as e: 135 logger.error("Handler error", method=request.method, error=str(e)) 136 return IPCResponse(request_id=request.id, error=str(e))
Server for handling local IPC connections via Unix domain sockets.
async def
start(self) -> None:
33 async def start(self) -> None: 34 """Start the IPC server.""" 35 if self._running: 36 return 37 38 # Clean up old socket 39 if os.path.exists(self.socket_path): 40 os.remove(self.socket_path) 41 42 self._server = await asyncio.start_unix_server(self._handle_client, path=self.socket_path) 43 self._running = True 44 logger.info("IPC Server started", path=self.socket_path)
Start the IPC server.
async def
stop(self) -> None:
46 async def stop(self) -> None: 47 """Stop the IPC server.""" 48 self._running = False 49 if self._server: 50 self._server.close() 51 await self._server.wait_closed() 52 self._server = None 53 54 if os.path.exists(self.socket_path): 55 with contextlib.suppress(OSError): 56 os.remove(self.socket_path) 57 logger.info("IPC Server stopped")
Stop the IPC server.