adds .local/share/millennium and easyeffects
This commit is contained in:
@@ -0,0 +1,7 @@
|
||||
"""
|
||||
WebSocket module for Extendium plugin.
|
||||
Handles communication between frontend and webkit clients.
|
||||
"""
|
||||
from .server import initialize_server, run_server, shutdown_server
|
||||
|
||||
__all__ = ['initialize_server', 'run_server', 'shutdown_server']
|
||||
@@ -0,0 +1,94 @@
|
||||
"""
|
||||
Client manager for WebSocket connections.
|
||||
Handles tracking and managing connected clients.
|
||||
"""
|
||||
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
|
||||
class ClientManager:
|
||||
"""
|
||||
Manages WebSocket clients.
|
||||
Tracks frontend and webkit clients separately.
|
||||
"""
|
||||
|
||||
def __init__(self):
|
||||
"""Initialize the client manager."""
|
||||
self._frontend_client = None
|
||||
self._webkit_clients: List[Dict[str, Any]] = []
|
||||
|
||||
@property
|
||||
def frontend_client(self) -> Optional[Dict[str, Any]]:
|
||||
"""Get the frontend client."""
|
||||
return self._frontend_client
|
||||
|
||||
@frontend_client.setter
|
||||
def frontend_client(self, client: Dict[str, Any]):
|
||||
"""Set the frontend client."""
|
||||
self._frontend_client = client
|
||||
|
||||
@property
|
||||
def webkit_clients(self) -> List[Dict[str, Any]]:
|
||||
"""Get all webkit clients."""
|
||||
return self._webkit_clients
|
||||
|
||||
def add_webkit_client(self, client: Dict[str, Any]):
|
||||
"""
|
||||
Add a webkit client.
|
||||
|
||||
Args:
|
||||
client: The WebSocket client object
|
||||
"""
|
||||
self._webkit_clients.append(client)
|
||||
|
||||
def remove_client(self, client_id: str) -> bool:
|
||||
"""
|
||||
Remove a client by its ID.
|
||||
|
||||
Args:
|
||||
client_id: The ID of the client to remove
|
||||
|
||||
Returns:
|
||||
bool: True if client was removed, False otherwise
|
||||
"""
|
||||
# Check if it's the frontend client
|
||||
if self._frontend_client and self._frontend_client.get('id') == client_id:
|
||||
self._frontend_client = None
|
||||
return True
|
||||
|
||||
# Check if it's a webkit client
|
||||
for i, client in enumerate(self._webkit_clients):
|
||||
if client.get('id') == client_id:
|
||||
del self._webkit_clients[i]
|
||||
return True
|
||||
|
||||
return False
|
||||
|
||||
def has_frontend_client(self) -> bool:
|
||||
"""Check if a frontend client is connected."""
|
||||
return self._frontend_client is not None
|
||||
|
||||
def has_webkit_clients(self) -> bool:
|
||||
"""Check if any webkit clients are connected."""
|
||||
return len(self._webkit_clients) > 0
|
||||
|
||||
def disconnect_all_clients(self):
|
||||
if self._frontend_client:
|
||||
try:
|
||||
websocket = getattr(self._frontend_client, 'websocket', None)
|
||||
if websocket and hasattr(websocket, 'close'):
|
||||
websocket.close()
|
||||
except Exception:
|
||||
pass
|
||||
finally:
|
||||
self._frontend_client = None
|
||||
|
||||
for client in self._webkit_clients:
|
||||
try:
|
||||
websocket = getattr(client, 'websocket', None)
|
||||
if websocket and hasattr(websocket, 'close'):
|
||||
websocket.close()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
self._webkit_clients.clear()
|
||||
@@ -0,0 +1,14 @@
|
||||
"""
|
||||
This file is used to bootup the websocket server without millennium so just from the terminal
|
||||
"""
|
||||
|
||||
from .server import initialize_server, run_server, shutdown_server
|
||||
|
||||
if __name__ == "__main__":
|
||||
initialize_server()
|
||||
run_server()
|
||||
try:
|
||||
input("Press Enter to exit...")
|
||||
except KeyboardInterrupt:
|
||||
pass
|
||||
shutdown_server()
|
||||
@@ -0,0 +1,150 @@
|
||||
"""
|
||||
Message handler for WebSocket communication.
|
||||
Processes different types of messages and routes them appropriately.
|
||||
"""
|
||||
|
||||
import json
|
||||
from typing import Any, Dict, Optional
|
||||
|
||||
from logger.logger import logger # pylint: disable=import-error
|
||||
|
||||
from .client_manager import ClientManager
|
||||
from .message_types import ClientType, MessageType
|
||||
from .request_handler import RequestHandler
|
||||
|
||||
|
||||
class MessageAdapter:
|
||||
def send_message(self, client: Dict[str, Any], message: str):
|
||||
pass
|
||||
|
||||
class MessageHandler:
|
||||
"""
|
||||
Handles processing and routing of WebSocket messages.
|
||||
"""
|
||||
|
||||
def __init__(self, server: MessageAdapter, client_manager: ClientManager, request_handler: RequestHandler):
|
||||
"""
|
||||
Initialize the message handler.
|
||||
|
||||
Args:
|
||||
server: The WebSocket server instance
|
||||
client_manager: The client manager instance
|
||||
request_handler: The request handler instance
|
||||
"""
|
||||
self._server = server
|
||||
self._client_manager = client_manager
|
||||
self._request_handler = request_handler
|
||||
self._message_handlers = {
|
||||
MessageType.IDENTIFY.value: self._handle_identify,
|
||||
MessageType.WEBKIT_MESSAGE.value: self._handle_webkit_message,
|
||||
MessageType.FRONTEND_RESPONSE.value: self._handle_frontend_response,
|
||||
MessageType.ERROR.value: self._handle_error,
|
||||
}
|
||||
|
||||
def process_message(self, client: Dict[str, Any], message: str) -> None:
|
||||
"""
|
||||
Process an incoming WebSocket message.
|
||||
|
||||
Args:
|
||||
client: The client that sent the message
|
||||
message: The message string
|
||||
"""
|
||||
try:
|
||||
data = json.loads(message)
|
||||
message_type = data.get('type')
|
||||
|
||||
if message_type in self._message_handlers:
|
||||
self._message_handlers[message_type](client, data)
|
||||
else:
|
||||
logger.error(f"Unknown message type: {message_type}")
|
||||
self._send_error(client, data.get('requestId'), "Unknown message type")
|
||||
except json.JSONDecodeError:
|
||||
logger.error(f"Invalid JSON message: {message}")
|
||||
self._send_error(client, None, "Invalid JSON message")
|
||||
except Exception as e:
|
||||
logger.error(f"Error processing message: {str(e)}")
|
||||
self._send_error(client, data.get('requestId') if 'data' in locals() else None, f"Error: {str(e)}")
|
||||
|
||||
def _handle_identify(self, client: Dict[str, Any], data: Dict[str, Any]) -> None:
|
||||
client_type = data.get('clientType')
|
||||
|
||||
if client_type == ClientType.FRONTEND.value:
|
||||
self._client_manager.frontend_client = client
|
||||
elif client_type == ClientType.WEBKIT.value:
|
||||
self._client_manager.add_webkit_client(client)
|
||||
else:
|
||||
logger.error(f"Unknown client type: {client_type}")
|
||||
self._send_error(client, None, f"Unknown client type: {client_type}")
|
||||
|
||||
def _handle_webkit_message(self, client: Dict[str, Any], data: Dict[str, Any]) -> None:
|
||||
"""
|
||||
Handle a message from a webkit client to the frontend.
|
||||
|
||||
Args:
|
||||
client: The client that sent the message
|
||||
data: The message data
|
||||
"""
|
||||
if not self._client_manager.has_frontend_client():
|
||||
logger.error("No frontend client connected to receive message")
|
||||
self._send_error(client, data.get('requestId'), "No frontend client connected")
|
||||
return
|
||||
|
||||
request_id = data.get('requestId')
|
||||
|
||||
if request_id:
|
||||
# Store the request for later response matching
|
||||
self._request_handler.add_request(request_id, client)
|
||||
|
||||
# Forward message to frontend
|
||||
self._server.send_message(self._client_manager.frontend_client, json.dumps(data))
|
||||
else:
|
||||
logger.error("Missing requestId in webkit message")
|
||||
self._send_error(client, request_id, "Missing requestId")
|
||||
|
||||
def _handle_frontend_response(self, client: Dict[str, Any], data: Dict[str, Any]) -> None:
|
||||
"""
|
||||
Handle a response from the frontend to a webkit client.
|
||||
|
||||
Args:
|
||||
client: The client that sent the message
|
||||
data: The message data
|
||||
"""
|
||||
request_id = data.get('requestId')
|
||||
|
||||
if not request_id:
|
||||
logger.error("Missing requestId in frontend response")
|
||||
self._send_error(client, None, "Missing requestId in response")
|
||||
return
|
||||
|
||||
request = self._request_handler.get_request(request_id)
|
||||
if request:
|
||||
webkit_client = request['client']
|
||||
# Send response back to the webkit client
|
||||
self._server.send_message(webkit_client, json.dumps(data))
|
||||
# Clean up the pending request
|
||||
self._request_handler.remove_request(request_id)
|
||||
else:
|
||||
# logger.error(f"Received response for unknown request: {request_id}")
|
||||
self._send_error(client, request_id, f"Unknown request: {request_id}")
|
||||
|
||||
def _handle_error(self, client: Dict[str, Any], data: Dict[str, Any]) -> None:
|
||||
logger.error(f"Received error: {data['error']} requestId: {data['requestId']} extensionName: {data['extensionName']}")
|
||||
|
||||
def _send_error(self, client: Dict[str, Any], request_id: Optional[str], error_message: str) -> None:
|
||||
"""
|
||||
Send an error message to a client.
|
||||
|
||||
Args:
|
||||
client: The client to send the error to
|
||||
request_id: The ID of the request that caused the error
|
||||
error_message: The error message
|
||||
"""
|
||||
error_data = {
|
||||
'type': MessageType.ERROR.value,
|
||||
'error': error_message
|
||||
}
|
||||
|
||||
if request_id:
|
||||
error_data['requestId'] = request_id
|
||||
|
||||
self._server.send_message(client, json.dumps(error_data))
|
||||
@@ -0,0 +1,20 @@
|
||||
"""
|
||||
Message type definitions for WebSocket communication.
|
||||
"""
|
||||
|
||||
from enum import Enum
|
||||
from typing import Any, Dict, List, Optional, TypedDict, Union
|
||||
|
||||
|
||||
class MessageType(Enum):
|
||||
"""Message types for WebSocket communication."""
|
||||
IDENTIFY = "identify"
|
||||
WEBKIT_MESSAGE = "webkit_message"
|
||||
FRONTEND_RESPONSE = "frontend_response"
|
||||
ERROR = "error"
|
||||
|
||||
class ClientType(Enum):
|
||||
"""Types of clients that can connect to the WebSocket server."""
|
||||
FRONTEND = "frontend"
|
||||
WEBKIT = "webkit"
|
||||
|
||||
@@ -0,0 +1,49 @@
|
||||
"""
|
||||
Request handler for WebSocket communication.
|
||||
Manages pending requests and their responses.
|
||||
"""
|
||||
|
||||
import time
|
||||
from typing import Any, Dict, Optional
|
||||
|
||||
from logger.logger import logger # pylint: disable=import-error
|
||||
|
||||
|
||||
class RequestHandler:
|
||||
"""
|
||||
Handles tracking and managing WebSocket requests.
|
||||
"""
|
||||
|
||||
def __init__(self):
|
||||
self._pending_requests: Dict[str, Dict[str, Any]] = {}
|
||||
|
||||
def add_request(self, request_id: str, client: Dict[str, Any]) -> None:
|
||||
self._pending_requests[request_id] = {
|
||||
'client': client,
|
||||
'timestamp': time.time()
|
||||
}
|
||||
|
||||
def get_request(self, request_id: str) -> Optional[Dict[str, Any]]:
|
||||
return self._pending_requests.get(request_id)
|
||||
|
||||
def remove_request(self, request_id: str) -> bool:
|
||||
if request_id in self._pending_requests:
|
||||
del self._pending_requests[request_id]
|
||||
return True
|
||||
return False
|
||||
|
||||
def cleanup_old_requests(self, max_age_seconds: int) -> int:
|
||||
current_time = time.time()
|
||||
to_remove = []
|
||||
|
||||
for request_id, request_data in self._pending_requests.items():
|
||||
if current_time - request_data['timestamp'] > max_age_seconds:
|
||||
to_remove.append(request_id)
|
||||
|
||||
for request_id in to_remove:
|
||||
del self._pending_requests[request_id]
|
||||
|
||||
if to_remove:
|
||||
logger.log(f"Cleaned up {len(to_remove)} stale requests")
|
||||
|
||||
return len(to_remove)
|
||||
@@ -0,0 +1,214 @@
|
||||
"""
|
||||
WebSocket server implementation for Extendium plugin.
|
||||
"""
|
||||
|
||||
import threading
|
||||
import time
|
||||
import uuid
|
||||
from typing import Optional
|
||||
|
||||
from logger.logger import logger # pylint: disable=import-error
|
||||
from websockets.exceptions import ConnectionClosed
|
||||
from websockets.sync.server import Server, ServerConnection, serve
|
||||
|
||||
from .client_manager import ClientManager
|
||||
from .message_handler import MessageHandler
|
||||
from .request_handler import RequestHandler
|
||||
|
||||
# Global instances
|
||||
_server: Optional[Server] = None
|
||||
_client_manager: Optional[ClientManager] = None
|
||||
_request_handler: Optional[RequestHandler] = None
|
||||
_message_handler: Optional[MessageHandler] = None
|
||||
_cleanup_thread: Optional[threading.Thread] = None
|
||||
_server_thread: Optional[threading.Thread] = None
|
||||
_running = False
|
||||
|
||||
|
||||
class WebSocketClientAdapter:
|
||||
"""
|
||||
Adapter class to maintain compatibility with the existing code that expects
|
||||
the client structure from the websocket-server package.
|
||||
"""
|
||||
def __init__(self, websocket: ServerConnection):
|
||||
self.websocket = websocket
|
||||
self.id = str(uuid.uuid4())
|
||||
self.data = {'id': self.id}
|
||||
|
||||
def get(self, key, default=None):
|
||||
return self.data.get(key, default)
|
||||
|
||||
def __getitem__(self, key):
|
||||
return self.data[key]
|
||||
|
||||
def __setitem__(self, key, value):
|
||||
self.data[key] = value
|
||||
|
||||
class MessageAdapter:
|
||||
"""
|
||||
Adapter to maintain compatibility with the existing MessageHandler that expects
|
||||
a server with a send_message method.
|
||||
"""
|
||||
def send_message(self, client: WebSocketClientAdapter, message: str) -> None:
|
||||
try:
|
||||
client.websocket.send(message)
|
||||
except ConnectionClosed:
|
||||
# logger.error("Failed to send message to client: connection closed")
|
||||
pass
|
||||
except Exception as e:
|
||||
logger.error(f"Error sending message: {e}")
|
||||
|
||||
def handle_client(websocket: ServerConnection) -> None:
|
||||
"""
|
||||
Main handler for each client connection.
|
||||
|
||||
Args:
|
||||
websocket: The WebSocket connection object
|
||||
"""
|
||||
|
||||
# Create client adapter
|
||||
client = WebSocketClientAdapter(websocket)
|
||||
|
||||
try:
|
||||
# Process messages in a loop
|
||||
for message in websocket:
|
||||
if _message_handler:
|
||||
_message_handler.process_message(client, message)
|
||||
except ConnectionClosed:
|
||||
pass
|
||||
except Exception as e:
|
||||
logger.error(f"Error handling client: {e}")
|
||||
finally:
|
||||
# Handle client disconnection (equivalent to _client_left)
|
||||
if _client_manager:
|
||||
_client_manager.remove_client(client.id)
|
||||
|
||||
def _cleanup_routine() -> None:
|
||||
"""Periodically clean up old pending requests."""
|
||||
global _request_handler, _running
|
||||
|
||||
while _running and _request_handler:
|
||||
# Clean up requests older than 30 seconds
|
||||
_request_handler.cleanup_old_requests(30)
|
||||
|
||||
# split wait into intervals to allow faster shutdown detection
|
||||
for _ in range(20):
|
||||
if not _running:
|
||||
return
|
||||
time.sleep(0.25)
|
||||
|
||||
def initialize_server() -> None:
|
||||
"""
|
||||
Initialize the WebSocket server.
|
||||
|
||||
Args:
|
||||
port: The port to listen on
|
||||
host: The host address to bind to
|
||||
|
||||
Returns:
|
||||
The initialized WebSocket server
|
||||
"""
|
||||
global _server, _client_manager, _request_handler, _message_handler
|
||||
|
||||
if _server:
|
||||
logger.warn("WebSocket server already initialized")
|
||||
return _server
|
||||
|
||||
# Create instances
|
||||
_client_manager = ClientManager()
|
||||
_request_handler = RequestHandler()
|
||||
|
||||
# Create message adapter for compatibility
|
||||
message_adapter = MessageAdapter()
|
||||
|
||||
# Create message handler with the adapter
|
||||
_message_handler = MessageHandler(message_adapter, _client_manager, _request_handler)
|
||||
|
||||
def run_server(port: int, host: str = "localhost") -> None:
|
||||
"""Run the WebSocket server in a separate thread."""
|
||||
global _server, _cleanup_thread, _server_thread, _running
|
||||
|
||||
if _running:
|
||||
logger.log("WebSocket server already running")
|
||||
return
|
||||
|
||||
_running = True
|
||||
|
||||
# Start cleanup thread as non-daemon for proper shutdown
|
||||
_cleanup_thread = threading.Thread(target=_cleanup_routine)
|
||||
_cleanup_thread.daemon = False # non-daemon so we can actually clean it up.
|
||||
_cleanup_thread.start()
|
||||
|
||||
# Start server in a separate thread
|
||||
def start_server():
|
||||
global _server
|
||||
try:
|
||||
_server = serve(handle_client, host, port)
|
||||
_server.serve_forever()
|
||||
except Exception as e:
|
||||
if _running: # Only log if we're not shutting down
|
||||
logger.error(f"Error in WebSocket server: {e}")
|
||||
finally:
|
||||
logger.log("WebSocket server thread ending")
|
||||
|
||||
_server_thread = threading.Thread(target=start_server)
|
||||
_server_thread.daemon = False # Non-daemon so we can wait for proper termination
|
||||
_server_thread.start()
|
||||
|
||||
logger.log(f"WebSocket server is running on {host}:{port}")
|
||||
|
||||
def shutdown_server() -> None:
|
||||
"""Shutdown the WebSocket server."""
|
||||
global _server, _running, _cleanup_thread, _server_thread, _client_manager, _request_handler, _message_handler
|
||||
|
||||
if not _running:
|
||||
logger.log("WebSocket server not running")
|
||||
return
|
||||
|
||||
logger.log("Shutting down WebSocket server...")
|
||||
_running = False
|
||||
|
||||
# disconnect all connected/dangling clients
|
||||
if _client_manager:
|
||||
try:
|
||||
_client_manager.disconnect_all_clients()
|
||||
except Exception as e:
|
||||
logger.error(f"Error disconnecting clients: {e}")
|
||||
|
||||
# shutdown the server
|
||||
if _server:
|
||||
try:
|
||||
_server.shutdown()
|
||||
except Exception as e:
|
||||
logger.error(f"Error shutting down WebSocket server: {e}")
|
||||
|
||||
try:
|
||||
if hasattr(_server, 'close'):
|
||||
_server.close()
|
||||
except Exception as e:
|
||||
logger.error(f"Error closing WebSocket server: {e}")
|
||||
|
||||
# give some time for the server to shutdown, this is somewhat just a safeguard.
|
||||
time.sleep(0.2)
|
||||
|
||||
# wait for cleanup thread to finish
|
||||
if _cleanup_thread and _cleanup_thread.is_alive():
|
||||
_cleanup_thread.join(timeout=5)
|
||||
if _cleanup_thread.is_alive():
|
||||
logger.error("Cleanup thread did not finish within timeout")
|
||||
|
||||
# wait for server thread to finish
|
||||
if _server_thread and _server_thread.is_alive():
|
||||
_server_thread.join(timeout=5)
|
||||
if _server_thread.is_alive():
|
||||
logger.error("Server thread did not finish within timeout")
|
||||
|
||||
# force inline gc cleanup from the gil.
|
||||
_server = None
|
||||
_client_manager = None
|
||||
_request_handler = None
|
||||
_message_handler = None
|
||||
_cleanup_thread = None
|
||||
_server_thread = None
|
||||
|
||||
logger.log("WebSocket server shutdown complete")
|
||||
Reference in New Issue
Block a user