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.

IPCServer(socket_path: str = '/tmp/tiphys.sock')
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
socket_path
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.

def register_handler(self, method: str, handler: Callable[..., typing.Any]) -> None:
59    def register_handler(self, method: str, handler: Callable[..., Any]) -> None:
60        """Register a handler for a specific method."""
61        self._handlers[method] = handler

Register a handler for a specific method.