"""MCP Server for Code Execution.""" import asyncio import base64 import json import logging import mimetypes import os import shutil import subprocess import tempfile import time import uuid from datetime import datetime from typing import Any import uvicorn from mcp.server.fastmcp import FastMCP from starlette.applications import Starlette from starlette.middleware.base import BaseHTTPMiddleware from starlette.requests import Request from starlette.responses import JSONResponse from starlette.routing import Route from .auth import ( OAUTH_ENABLED, extract_bearer_token, get_protected_resource_metadata, get_www_authenticate_header, verify_access_token, ) from .languages import ( BaseLanguageHandler, JavaHandler, JavaScriptHandler, PythonHandler, RHandler, RustHandler, SwiftHandler, ) logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) # Server configuration MCP_PORT = int(os.environ.get("MCP_PORT", "8000")) class AcceptHeaderMiddleware(BaseHTTPMiddleware): """ Middleware to fix Accept headers for MCP Streamable HTTP compatibility. MCP SDK strictly requires both application/json and text/event-stream in the Accept header. Some clients (e.g., wildcard */* or missing text/event-stream) get rejected with 406 Not Acceptable. This middleware normalizes the Accept header before it reaches the MCP handler. """ async def dispatch(self, request: Request, call_next): if request.method == "POST" and request.url.path.startswith("/mcp"): accept = request.headers.get("accept", "") needs_fix = ( "text/event-stream" not in accept or "application/json" not in accept ) if needs_fix: fixed = "application/json, text/event-stream" # Replace Accept header directly in ASGI scope request.scope["headers"] = [ (k, v) for k, v in request.scope["headers"] if k != b"accept" ] + [(b"accept", fixed.encode())] return await call_next(request) class OAuthMiddleware(BaseHTTPMiddleware): """ Middleware for OAuth2 Bearer token authentication. - Skips authentication for health check and metadata endpoints - Returns 401 with WWW-Authenticate header for missing/invalid tokens - Stores verified user info in request.state.user """ PUBLIC_PATHS = { "/health", "/.well-known/oauth-protected-resource", "/.well-known/oauth-authorization-server", "/.well-known/openid-configuration", } PUBLIC_PREFIXES = ( "/.well-known/", "/mcp/.well-known/", ) async def dispatch(self, request: Request, call_next): # Skip auth if OAuth is disabled if not OAUTH_ENABLED: return await call_next(request) # Skip auth for public endpoints path = request.url.path if path in self.PUBLIC_PATHS or path.startswith(self.PUBLIC_PREFIXES): return await call_next(request) # Extract Bearer token authorization = request.headers.get("authorization") token = extract_bearer_token(authorization) if not token: logger.warning(f"Missing or invalid Authorization header for {request.url.path}") return JSONResponse( status_code=401, content={"error": "unauthorized", "message": "Missing or invalid Bearer token"}, headers={"WWW-Authenticate": get_www_authenticate_header()} ) # Verify token try: user_info = await verify_access_token(token) request.state.user = user_info logger.debug(f"Authenticated: {user_info}") except ValueError as e: logger.warning(f"Token verification failed: {e}") return JSONResponse( status_code=401, content={"error": "invalid_token", "message": str(e)}, headers={"WWW-Authenticate": get_www_authenticate_header()} ) return await call_next(request) # Create FastMCP server mcp = FastMCP("code-executor", host="0.0.0.0", port=MCP_PORT) # Language handlers registry _handlers: dict[str, BaseLanguageHandler] = {} # Executions storage for async code execution # Key: execution_id, Value: execution info dict _executions: dict[str, dict[str, Any]] = {} # Environments storage for pre-installed dependencies # Key: environment_id, Value: environment info dict _environments: dict[str, dict[str, Any]] = {} def _register_handlers(): """Register all language handlers.""" global _handlers # Python versions _handlers["python"] = PythonHandler("3.12") _handlers["python3"] = PythonHandler("3.12") _handlers["python3.11"] = PythonHandler("3.11") _handlers["python3.12"] = PythonHandler("3.12") # Other languages _handlers["r"] = RHandler() _handlers["rust"] = RustHandler() _handlers["java"] = JavaHandler() _handlers["swift"] = SwiftHandler() _handlers["javascript"] = JavaScriptHandler() _handlers["js"] = JavaScriptHandler() _handlers["node"] = JavaScriptHandler() # Initialize handlers _register_handlers() def _get_handler(language: str) -> BaseLanguageHandler | None: """Get handler for the specified language.""" return _handlers.get(language.lower()) def _list_languages() -> list[str]: """List all supported languages.""" return list({h.name for h in _handlers.values()}) def _generate_execution_id() -> str: """Generate a unique execution ID.""" return uuid.uuid4().hex[:8] def _safe_b64decode(data: str) -> bytes: """Decode base64 with automatic padding correction.""" # Remove whitespace and newlines data = data.strip().replace('\n', '').replace('\r', '').replace(' ', '') # Add padding if needed padding = 4 - (len(data) % 4) if padding != 4: data += '=' * padding return base64.b64decode(data) def _get_container_status(container_name: str) -> dict[str, Any] | None: """Get Docker container status.""" try: result = subprocess.run( ["docker", "inspect", "--format", '{"running": {{.State.Running}}, "exit_code": {{.State.ExitCode}}, "status": "{{.State.Status}}"}', container_name], capture_output=True, text=True, timeout=5 ) if result.returncode == 0: return json.loads(result.stdout.strip()) except Exception as e: logger.debug(f"Failed to get container status for {container_name}: {e}") return None def _collect_output_files(output_dir: str) -> list[dict[str, Any]]: """Collect output files from directory.""" output_files = [] if os.path.exists(output_dir): for filename in os.listdir(output_dir): file_path = os.path.join(output_dir, filename) if os.path.isfile(file_path): try: with open(file_path, 'rb') as f: file_data = f.read() mime_type, _ = mimetypes.guess_type(filename) output_files.append({ "filename": filename, "mime_type": mime_type or "application/octet-stream", "data": base64.b64encode(file_data).decode('utf-8'), "size": len(file_data) }) except Exception as e: logger.warning(f"Failed to read output file {filename}: {e}") return output_files async def _start_code_async( language: str, code: str, input_files: list[dict[str, str]] | None = None, environment_id: str | None = None ) -> dict[str, Any]: """Start code execution in a detached Docker container.""" execution_id = _generate_execution_id() start_time = time.time() # Check if using pre-configured environment env = None use_env = False if environment_id: env = _environments.get(environment_id) if not env: return { "success": False, "execution_id": None, "error": f"Environment {environment_id} not found" } # Check environment status if env["status"] == "setting_up": return { "success": False, "execution_id": None, "error": f"Environment {environment_id} is still setting up. Use get_environment_status to check progress." } if env["status"] == "failed": return { "success": False, "execution_id": None, "error": f"Environment {environment_id} setup failed: {env.get('error', 'Unknown error')}" } language = env["language"] use_env = True handler = _get_handler(language) if not handler: return { "success": False, "execution_id": None, "error": f"Unsupported language: {language}. Supported: {_list_languages()}" } # Use environment directory or create temporary directory if use_env: temp_dir = env["env_dir"] output_dir = env["output_dir"] input_dir = env["input_dir"] # Clear previous output files for f in os.listdir(output_dir): os.remove(os.path.join(output_dir, f)) else: temp_dir = tempfile.mkdtemp(prefix=f"code_exec_{execution_id}_") output_dir = os.path.join(temp_dir, "output") input_dir = os.path.join(temp_dir, "input") os.makedirs(output_dir, exist_ok=True) os.makedirs(input_dir, exist_ok=True) try: # Write input files if provided if input_files: for input_file in input_files: filename = input_file.get("filename", "") data = input_file.get("data", "") if filename and data: safe_filename = os.path.basename(filename) file_path = os.path.join(input_dir, safe_filename) try: file_data = _safe_b64decode(data) with open(file_path, 'wb') as f: f.write(file_data) logger.info(f"Wrote input file: {safe_filename} ({len(file_data)} bytes)") except Exception as e: logger.warning(f"Failed to write input file {filename}: {e}") # Prepare files (no requirements if using environment) files = handler.prepare_files(code, "") for filename, content in files.items(): file_path = os.path.join(temp_dir, filename) os.makedirs(os.path.dirname(file_path), exist_ok=True) with open(file_path, 'w') as f: f.write(content) # Create execution script (no dependency installation) script = handler.get_execution_script("") script_path = os.path.join(temp_dir, "run.sh") with open(script_path, 'w') as f: f.write(script) os.chmod(script_path, 0o755) # Container name for tracking container_name = f"code_exec_{execution_id}" # Build Docker command (detached mode) docker_cmd = [ "docker", "run", "-d", # Detached mode "--name", container_name, "--network", "none", "--memory", "512m", "--cpus", "1", "--pids-limit", "100", "-v", f"{temp_dir}:/app", "-w", "/app", handler.docker_image, "/bin/bash", "/app/run.sh" ] logger.info(f"Starting detached execution {execution_id}: {' '.join(docker_cmd)}") # Start container result = subprocess.run(docker_cmd, capture_output=True, text=True, timeout=30) if result.returncode != 0: if not use_env: shutil.rmtree(temp_dir, ignore_errors=True) return { "success": False, "execution_id": None, "error": f"Failed to start container: {result.stderr}" } container_id = result.stdout.strip() # Store execution info _executions[execution_id] = { "execution_id": execution_id, "container_id": container_id, "container_name": container_name, "language": language, "status": "running", "temp_dir": temp_dir, "output_dir": output_dir, "start_time": start_time, "created_at": datetime.now().isoformat(), "exit_code": None, "output": None, "error": None, "files": None, "use_env": use_env } logger.info(f"Execution {execution_id} started in container {container_name}") return { "success": True, "execution_id": execution_id, "container_name": container_name, "status": "running", "message": "Code execution started. Use get_execution_status to check progress." } except Exception as e: logger.exception(f"Failed to start execution {execution_id}") if not use_env: shutil.rmtree(temp_dir, ignore_errors=True) return { "success": False, "execution_id": None, "error": str(e) } def _update_execution_status(execution_id: str) -> dict[str, Any] | None: """Update and return execution status from Docker.""" if execution_id not in _executions: return None execution = _executions[execution_id] # Already completed, return cached result if execution["status"] in ("completed", "failed", "stopped"): return execution # Check container status container_status = _get_container_status(execution["container_name"]) if container_status is None: # Container doesn't exist anymore execution["status"] = "unknown" execution["error"] = "Container not found" return execution if container_status["running"]: execution["status"] = "running" else: # Container finished exit_code = container_status["exit_code"] execution["exit_code"] = exit_code execution["status"] = "completed" if exit_code == 0 else "failed" execution["execution_time"] = time.time() - execution["start_time"] # Get logs try: logs_result = subprocess.run( ["docker", "logs", execution["container_name"]], capture_output=True, text=True, timeout=10 ) output = logs_result.stdout error = logs_result.stderr # Filter output output_lines = output.split('\n') filtered_lines = [line for line in output_lines if not line.startswith('=== ')] execution["output"] = '\n'.join(filtered_lines).strip() execution["error"] = error except Exception as e: logger.warning(f"Failed to get logs for {execution_id}: {e}") execution["output"] = "" execution["error"] = str(e) # Collect output files execution["files"] = _collect_output_files(execution["output_dir"]) # Remove container try: subprocess.run( ["docker", "rm", "-f", execution["container_name"]], capture_output=True, timeout=10 ) except Exception as e: logger.warning(f"Failed to remove container {execution['container_name']}: {e}") return execution def _cleanup_execution(execution_id: str): """Cleanup execution resources.""" if execution_id not in _executions: return execution = _executions[execution_id] # Remove container if still exists try: subprocess.run( ["docker", "rm", "-f", execution["container_name"]], capture_output=True, timeout=10 ) except Exception: pass # Remove temp directory (only if not using environment) if not execution.get("use_env"): temp_dir = execution.get("temp_dir") if temp_dir and os.path.exists(temp_dir): try: shutil.rmtree(temp_dir) except Exception as e: logger.warning(f"Failed to cleanup temp dir for {execution_id}: {e}") # Remove from executions del _executions[execution_id] logger.info(f"Cleaned up execution {execution_id}") async def _setup_environment( language: str, requirements: str ) -> dict[str, Any]: """Start environment setup in a detached Docker container.""" environment_id = _generate_execution_id() start_time = time.time() handler = _get_handler(language) if not handler: return { "success": False, "environment_id": None, "error": f"Unsupported language: {language}. Supported: {_list_languages()}" } if not requirements.strip(): return { "success": False, "environment_id": None, "error": "No requirements provided. Use run_code directly for code without dependencies." } # Create persistent directory for this environment env_dir = tempfile.mkdtemp(prefix=f"code_env_{environment_id}_") output_dir = os.path.join(env_dir, "output") input_dir = os.path.join(env_dir, "input") os.makedirs(output_dir, exist_ok=True) os.makedirs(input_dir, exist_ok=True) try: # Prepare files for dependency installation files = handler.prepare_files("", requirements) for filename, content in files.items(): file_path = os.path.join(env_dir, filename) os.makedirs(os.path.dirname(file_path), exist_ok=True) with open(file_path, 'w') as f: f.write(content) # Create installation script (only install dependencies, no code execution) install_script = handler.get_execution_script(requirements) # Modify script to only install dependencies (remove code execution part) script_lines = install_script.split('\n') install_only_lines = [] for line in script_lines: install_only_lines.append(line) # Stop after dependency installation commands if 'pip install' in line or 'npm install' in line or 'cargo build' in line: install_only_lines.append('echo "Dependencies installed successfully"') break if 'install.packages' in line or 'Rscript' in line.lower(): install_only_lines.append('echo "Dependencies installed successfully"') break # If no install command found, use the full script but add success message if len(install_only_lines) <= 2: install_only_lines = script_lines[:len(script_lines)//2] install_only_lines.append('echo "Environment setup completed"') install_script = '\n'.join(install_only_lines) script_path = os.path.join(env_dir, "setup.sh") with open(script_path, 'w') as f: f.write(install_script) os.chmod(script_path, 0o755) # Container name for tracking container_name = f"code_env_setup_{environment_id}" # Build Docker command (detached mode) docker_cmd = [ "docker", "run", "-d", # Detached mode "--name", container_name, "--network", "host", # Allow network for package downloads "--memory", "1g", "--cpus", "2", "-v", f"{env_dir}:/app", "-w", "/app", handler.docker_image, "/bin/bash", "/app/setup.sh" ] logger.info(f"Setting up environment {environment_id}: {' '.join(docker_cmd)}") # Start container in detached mode result = subprocess.run(docker_cmd, capture_output=True, text=True, timeout=30) if result.returncode != 0: shutil.rmtree(env_dir, ignore_errors=True) return { "success": False, "environment_id": None, "error": f"Failed to start setup container: {result.stderr}" } container_id = result.stdout.strip() # Store environment info with "setting_up" status _environments[environment_id] = { "environment_id": environment_id, "container_id": container_id, "container_name": container_name, "language": language, "requirements": requirements, "status": "setting_up", "env_dir": env_dir, "output_dir": output_dir, "input_dir": input_dir, "docker_image": handler.docker_image, "start_time": start_time, "created_at": datetime.now().isoformat(), "setup_time": None, "error": None } logger.info(f"Environment {environment_id} setup started in container {container_name}") return { "success": True, "environment_id": environment_id, "status": "setting_up", "message": "Environment setup started. Use get_environment_status to check progress." } except Exception as e: logger.exception(f"Failed to start environment setup {environment_id}") shutil.rmtree(env_dir, ignore_errors=True) return { "success": False, "environment_id": None, "error": str(e) } def _update_environment_status(environment_id: str) -> dict[str, Any] | None: """Update and return environment status from Docker.""" if environment_id not in _environments: return None env = _environments[environment_id] # Already ready or failed, return cached result if env["status"] in ("ready", "failed"): return env # Check container status container_status = _get_container_status(env["container_name"]) if container_status is None: # Container doesn't exist anymore - check if it completed env["status"] = "failed" env["error"] = "Setup container not found" return env if container_status["running"]: env["status"] = "setting_up" else: # Container finished exit_code = container_status["exit_code"] env["setup_time"] = time.time() - env["start_time"] if exit_code == 0: env["status"] = "ready" logger.info(f"Environment {environment_id} is ready (setup time: {env['setup_time']:.2f}s)") else: env["status"] = "failed" # Get error logs try: logs_result = subprocess.run( ["docker", "logs", env["container_name"]], capture_output=True, text=True, timeout=10 ) env["error"] = logs_result.stderr or logs_result.stdout except Exception as e: env["error"] = str(e) logger.warning(f"Environment {environment_id} setup failed: {env['error']}") # Remove setup container try: subprocess.run( ["docker", "rm", "-f", env["container_name"]], capture_output=True, timeout=10 ) except Exception: pass return env def _cleanup_environment(environment_id: str): """Cleanup environment resources.""" if environment_id not in _environments: return env = _environments[environment_id] env_dir = env.get("env_dir") if env_dir and os.path.exists(env_dir): try: shutil.rmtree(env_dir) except Exception as e: logger.warning(f"Failed to cleanup env dir for {environment_id}: {e}") del _environments[environment_id] logger.info(f"Cleaned up environment {environment_id}") async def _execute_code( language: str, code: str, requirements: str = "", timeout: int = 60, input_files: list[dict[str, str]] | None = None, environment_id: str | None = None, network: str = "" ) -> dict[str, Any]: """Execute code in a Docker container.""" start_time = time.time() # Check if using pre-configured environment env = None if environment_id: env = _environments.get(environment_id) if not env: return { "success": False, "output": "", "error": f"Environment {environment_id} not found", "execution_time": time.time() - start_time, "exit_code": -1, "files": [] } # Check environment status if env["status"] == "setting_up": return { "success": False, "output": "", "error": f"Environment {environment_id} is still setting up. Use get_environment_status to check progress.", "execution_time": time.time() - start_time, "exit_code": -1, "files": [] } if env["status"] == "failed": return { "success": False, "output": "", "error": f"Environment {environment_id} setup failed: {env.get('error', 'Unknown error')}", "execution_time": time.time() - start_time, "exit_code": -1, "files": [] } # Use environment's language language = env["language"] handler = _get_handler(language) if not handler: return { "success": False, "output": "", "error": f"Unsupported language: {language}. Supported: {_list_languages()}", "execution_time": time.time() - start_time, "exit_code": -1, "files": [] } # Use environment directory or create temporary directory if env: temp_dir = env["env_dir"] output_dir = env["output_dir"] input_dir = env["input_dir"] # Clear previous output files for f in os.listdir(output_dir): os.remove(os.path.join(output_dir, f)) use_env = True else: temp_dir = tempfile.mkdtemp(prefix="code_exec_") output_dir = os.path.join(temp_dir, "output") input_dir = os.path.join(temp_dir, "input") os.makedirs(output_dir, exist_ok=True) os.makedirs(input_dir, exist_ok=True) use_env = False try: # Write input files if provided if input_files: for input_file in input_files: filename = input_file.get("filename", "") data = input_file.get("data", "") if filename and data: # Sanitize filename to prevent path traversal safe_filename = os.path.basename(filename) file_path = os.path.join(input_dir, safe_filename) try: file_data = _safe_b64decode(data) with open(file_path, 'wb') as f: f.write(file_data) logger.info(f"Wrote input file: {safe_filename} ({len(file_data)} bytes)") except Exception as e: logger.warning(f"Failed to write input file {filename}: {e}") # Prepare files (skip requirements if using environment) effective_requirements = "" if use_env else requirements files = handler.prepare_files(code, effective_requirements) for filename, content in files.items(): file_path = os.path.join(temp_dir, filename) os.makedirs(os.path.dirname(file_path), exist_ok=True) with open(file_path, 'w') as f: f.write(content) # Create execution script (skip dependency installation if using environment) script = handler.get_execution_script(effective_requirements) script_path = os.path.join(temp_dir, "run.sh") with open(script_path, 'w') as f: f.write(script) os.chmod(script_path, 0o755) # Build Docker command docker_cmd = [ "docker", "run", "--rm", "--network", network if network else "none", "--memory", "512m", "--cpus", "1", "--pids-limit", "100", "-v", f"{temp_dir}:/app", "-w", "/app", handler.docker_image, "/bin/bash", "/app/run.sh" ] logger.info(f"Executing: {' '.join(docker_cmd)}") # Run Docker container process = await asyncio.create_subprocess_exec( *docker_cmd, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE ) try: stdout, stderr = await asyncio.wait_for( process.communicate(), timeout=timeout ) except TimeoutError: process.kill() await process.wait() return { "success": False, "output": "", "error": f"Execution timed out after {timeout} seconds", "execution_time": time.time() - start_time, "exit_code": -1, "files": [] } execution_time = time.time() - start_time exit_code = process.returncode # Decode output output = stdout.decode('utf-8', errors='replace') if stdout else "" error = stderr.decode('utf-8', errors='replace') if stderr else "" # Remove the echo headers from output output_lines = output.split('\n') filtered_lines = [ line for line in output_lines if not line.startswith('=== ') ] output = '\n'.join(filtered_lines).strip() # Collect output files (images, etc.) output_files = [] if os.path.exists(output_dir): for filename in os.listdir(output_dir): file_path = os.path.join(output_dir, filename) if os.path.isfile(file_path): try: with open(file_path, 'rb') as f: file_data = f.read() mime_type, _ = mimetypes.guess_type(filename) output_files.append({ "filename": filename, "mime_type": mime_type or "application/octet-stream", "data": base64.b64encode(file_data).decode('utf-8'), "size": len(file_data) }) except Exception as e: logger.warning(f"Failed to read output file {filename}: {e}") return { "success": (exit_code == 0), "output": output, "error": error, "execution_time": execution_time, "exit_code": exit_code, "files": output_files } except Exception as e: logger.exception("Execution failed") return { "success": False, "output": "", "error": str(e), "execution_time": time.time() - start_time, "exit_code": -1, "files": [] } finally: # Cleanup temporary directory (only if not using environment) if not use_env: try: shutil.rmtree(temp_dir) except Exception as e: logger.warning(f"Failed to cleanup temp dir: {e}") @mcp.tool() async def run_code( language: str, code: str, timeout: int = 60, input_files: list[dict[str, str]] | None = None, environment_id: str | None = None, requirements: str = "", network: str = "" ) -> str: """ Execute code in an isolated Docker container. If requirements are provided without an environment_id, an environment is automatically created, dependencies are installed, code is executed, and the environment is cleaned up afterward. Supported languages are python, python3.11, python3.12, r, rust, java, swift, javascript. FILE INPUT - To upload files for code to process, use the input_files parameter. Each file should be a dict with "filename" and "data" (base64-encoded). Uploaded files are saved to /app/input/ directory. In your code, read files from /app/input/. For example if you upload "data.csv", read it with open('/app/input/data.csv'). FILE OUTPUT - To return files (images, data, etc.), save them to /app/output/ directory. Files saved to /app/output/ will be returned as base64-encoded data in the response. For example plt.savefig('/app/output/chart.png') Parameters - language is the programming language (e.g. python3.12, r, rust, java, swift, javascript) - code is the source code to execute - timeout is the execution timeout in seconds (default 60, max 300) - input_files is a list of input files, each with "filename" and "data" (base64) keys - environment_id is a pre-configured environment ID from setup_environment (optional) - requirements is the dependencies in language-specific format (e.g. "requests>=2.31.0\\nnumpy") - network is the Docker network name for the execution container (e.g. "host", "agent-network"). Empty string means no network access Returns a JSON string with execution result containing success (boolean), output (stdout), error (stderr or error message), execution_time (seconds), exit_code, and files array (output files with filename, mime_type, data as base64, size). """ # Validate timeout if timeout < 1: timeout = 1 elif timeout > 300: timeout = 300 logger.info(f"Executing {language} code (timeout: {timeout}s, env: {environment_id}, requirements: {bool(requirements)}, network: {network})") # Auto setup environment when requirements are provided without environment_id auto_env_id = None if requirements.strip() and not environment_id: logger.info(f"Auto-setting up environment for {language} with requirements") setup_result = await _setup_environment(language, requirements) if not setup_result.get("success"): return json.dumps({ "success": False, "output": "", "error": f"Auto environment setup failed: {setup_result.get('error', 'Unknown error')}", "execution_time": 0, "exit_code": -1, "files": [] }, indent=2, ensure_ascii=False) auto_env_id = setup_result["environment_id"] # Poll until environment is ready (use remaining timeout) poll_start = time.time() poll_timeout = min(timeout, 300) while time.time() - poll_start < poll_timeout: env_status = _update_environment_status(auto_env_id) if env_status and env_status["status"] == "ready": break if env_status and env_status["status"] == "failed": error_msg = env_status.get("error", "Unknown error") _cleanup_environment(auto_env_id) return json.dumps({ "success": False, "output": "", "error": f"Dependency installation failed: {error_msg}", "execution_time": time.time() - poll_start, "exit_code": -1, "files": [] }, indent=2, ensure_ascii=False) await asyncio.sleep(1) else: _cleanup_environment(auto_env_id) return json.dumps({ "success": False, "output": "", "error": f"Dependency installation timed out after {poll_timeout}s", "execution_time": time.time() - poll_start, "exit_code": -1, "files": [] }, indent=2, ensure_ascii=False) environment_id = auto_env_id try: result = await _execute_code( language, code, "", timeout, input_files, environment_id, network ) logger.info(f"Execution completed: success={result['success']}, time={result['execution_time']:.2f}s") return json.dumps(result, indent=2, ensure_ascii=False) finally: # Cleanup auto-created environment if auto_env_id: _cleanup_environment(auto_env_id) @mcp.tool() async def start_code( language: str, code: str, input_files: list[dict[str, str]] | None = None, environment_id: str | None = None ) -> str: """ Start code execution asynchronously in a Docker container. Unlike run_code which waits for completion, this starts execution in the background and returns immediately with an execution_id. Use this for long-running code that might exceed timeout limits. For code with dependencies, first call setup_environment to create an environment, then use the returned environment_id with this function. WORKFLOW 1. Call setup_environment(language, requirements) to get environment_id (if needed) 2. Call start_code to begin execution and get execution_id 3. Call get_execution_status(execution_id) to check progress 4. When status is "completed" or "failed", call get_execution_result(execution_id) 5. Optionally call stop_execution(execution_id) to terminate early Supported languages are python, python3.11, python3.12, r, rust, java, swift, javascript. File Input/Output is same as run_code. Uploaded files are saved to /app/input/ (read with open('/app/input/')). Save output files to /app/output/ to return them. Parameters - language is the programming language - code is the source code to execute - input_files is a list of input files, each with "filename" and "data" (base64) keys - environment_id is a pre-configured environment ID from setup_environment (optional) Returns a JSON string with success (boolean), execution_id (to track this execution), status ("running" if started), and message (instructions for next steps). """ logger.info(f"Starting async {language} code execution (env: {environment_id})") result = await _start_code_async(language, code, input_files, environment_id) if result["success"]: logger.info(f"Async execution started: {result['execution_id']}") else: logger.warning(f"Failed to start async execution: {result.get('error')}") return json.dumps(result, indent=2, ensure_ascii=False) @mcp.tool() async def get_execution_status(execution_id: str) -> str: """ Get the current status of a code execution. Use this to monitor progress of code started with start_code. The execution_id parameter is the execution ID returned by start_code. Returns a JSON string with execution_id, status ("running", "completed", "failed", "stopped", or "unknown"), execution_time (elapsed time if completed), and exit_code (process exit code if completed). """ execution = _update_execution_status(execution_id) if execution is None: return json.dumps({ "success": False, "error": f"Execution {execution_id} not found" }, indent=2, ensure_ascii=False) # Return status without full output/files (use get_execution_result for that) return json.dumps({ "success": True, "execution_id": execution["execution_id"], "status": execution["status"], "language": execution["language"], "created_at": execution["created_at"], "execution_time": execution.get("execution_time"), "exit_code": execution.get("exit_code") }, indent=2, ensure_ascii=False) @mcp.tool() async def get_execution_result(execution_id: str, cleanup: bool = True) -> str: """ Get the full result of a completed code execution. Call this after get_execution_status shows "completed" or "failed". By default, this cleans up resources after returning results. Parameters - execution_id is the execution ID returned by start_code - cleanup indicates whether to cleanup resources after getting result (default True) Returns a JSON string with full execution result (same format as run_code) including success (boolean), output (stdout), error (stderr or error message), execution_time (seconds), exit_code, and files array (output files with filename, mime_type, data as base64, size). """ execution = _update_execution_status(execution_id) if execution is None: return json.dumps({ "success": False, "error": f"Execution {execution_id} not found" }, indent=2, ensure_ascii=False) if execution["status"] == "running": return json.dumps({ "success": False, "error": "Execution still running. Use get_execution_status to check progress.", "status": "running" }, indent=2, ensure_ascii=False) result = { "success": execution["status"] == "completed", "execution_id": execution["execution_id"], "status": execution["status"], "output": execution.get("output", ""), "error": execution.get("error", ""), "execution_time": execution.get("execution_time", 0), "exit_code": execution.get("exit_code", -1), "files": execution.get("files", []) } if cleanup: _cleanup_execution(execution_id) return json.dumps(result, indent=2, ensure_ascii=False) @mcp.tool() async def stop_execution(execution_id: str) -> str: """ Stop a running code execution. Use this to terminate a long-running execution early. The execution_id parameter is the execution ID returned by start_code. Returns a JSON string with success (whether the execution was stopped) and status (final status "stopped"). """ if execution_id not in _executions: return json.dumps({ "success": False, "error": f"Execution {execution_id} not found" }, indent=2, ensure_ascii=False) execution = _executions[execution_id] if execution["status"] != "running": return json.dumps({ "success": False, "error": f"Execution is not running (status: {execution['status']})" }, indent=2, ensure_ascii=False) # Kill the container try: subprocess.run( ["docker", "kill", execution["container_name"]], capture_output=True, timeout=10 ) execution["status"] = "stopped" execution["execution_time"] = time.time() - execution["start_time"] logger.info(f"Stopped execution {execution_id}") # Cleanup _cleanup_execution(execution_id) return json.dumps({ "success": True, "execution_id": execution_id, "status": "stopped", "message": "Execution stopped and cleaned up" }, indent=2, ensure_ascii=False) except Exception as e: logger.warning(f"Failed to stop execution {execution_id}: {e}") return json.dumps({ "success": False, "error": f"Failed to stop execution: {e}" }, indent=2, ensure_ascii=False) @mcp.tool() async def list_executions() -> str: """ List all tracked code executions. Returns a JSON string with list of executions and their current status. """ # Update all execution statuses for exec_id in list(_executions.keys()): _update_execution_status(exec_id) executions_list = [] for execution in _executions.values(): executions_list.append({ "execution_id": execution["execution_id"], "language": execution["language"], "status": execution["status"], "created_at": execution["created_at"], "execution_time": execution.get("execution_time") }) return json.dumps({ "success": True, "count": len(executions_list), "executions": executions_list }, indent=2, ensure_ascii=False) @mcp.tool() async def setup_environment( language: str, requirements: str ) -> str: """ Start setting up a reusable environment with dependencies. This starts dependency installation in the background and returns immediately. Use get_environment_status to check when setup is complete. WORKFLOW 1. Call setup_environment(language, requirements) to get environment_id 2. Call get_environment_status(environment_id) until status is "ready" 3. Call run_code(language, code, environment_id=environment_id) 4. Optionally call delete_environment(environment_id) when done Parameters - language is the programming language (python, python3.12, r, rust, java, swift, javascript) - requirements is the dependencies in language-specific format. For Python use pip format (requests>=2.31.0). For R use package names (ggplot2). For Rust use Cargo.toml format (serde = "1.0"). For Java use Maven coordinates (com.google.guava/guava/32.1.2-jre). For JavaScript use package@version (lodash@4.17.21). Returns a JSON string with success (whether setup started), environment_id (to track this environment), and status ("setting_up" if started). """ logger.info(f"Setting up {language} environment with requirements") result = await _setup_environment(language, requirements) if result["success"]: logger.info(f"Environment setup started: {result['environment_id']}") else: logger.warning(f"Failed to start environment setup: {result.get('error')}") return json.dumps(result, indent=2, ensure_ascii=False) @mcp.tool() async def get_environment_status(environment_id: str) -> str: """ Get the current status of an environment setup. Use this to check if setup_environment has completed. The environment_id parameter is the environment ID returned by setup_environment. Returns a JSON string with environment_id, status ("setting_up", "ready", or "failed"), setup_time (time taken if ready), and error (error message if failed). """ env = _update_environment_status(environment_id) if env is None: return json.dumps({ "success": False, "error": f"Environment {environment_id} not found" }, indent=2, ensure_ascii=False) return json.dumps({ "success": True, "environment_id": env["environment_id"], "status": env["status"], "language": env["language"], "created_at": env["created_at"], "setup_time": env.get("setup_time"), "error": env.get("error") }, indent=2, ensure_ascii=False) @mcp.tool() async def list_environments() -> str: """ List all environments and their current status. Returns a JSON string with list of environments and their details. """ # Update all environment statuses for env_id in list(_environments.keys()): _update_environment_status(env_id) environments_list = [] for env in _environments.values(): environments_list.append({ "environment_id": env["environment_id"], "language": env["language"], "status": env["status"], "requirements": env["requirements"][:100] + "..." if len(env["requirements"]) > 100 else env["requirements"], "created_at": env["created_at"], "setup_time": env.get("setup_time") }) return json.dumps({ "success": True, "count": len(environments_list), "environments": environments_list }, indent=2, ensure_ascii=False) @mcp.tool() async def delete_environment(environment_id: str) -> str: """ Delete a pre-configured environment. Use this to clean up environments that are no longer needed. The environment_id parameter is the environment ID returned by setup_environment. Returns a JSON string with success (whether the environment was deleted). """ if environment_id not in _environments: return json.dumps({ "success": False, "error": f"Environment {environment_id} not found" }, indent=2, ensure_ascii=False) _cleanup_environment(environment_id) return json.dumps({ "success": True, "message": f"Environment {environment_id} deleted" }, indent=2, ensure_ascii=False) @mcp.tool() async def list_supported_languages() -> str: """ List all supported programming languages and their requirements format. Returns a JSON string with supported languages and details. """ return json.dumps({ "languages": _list_languages(), "details": { "python": { "aliases": ["python", "python3", "python3.11", "python3.12"], "requirements_format": "pip requirements.txt format (package>=version)", "example_requirements": "requests>=2.31.0\nnumpy", "example_code": 'import requests\nprint(requests.__version__)' }, "r": { "aliases": ["r"], "requirements_format": "Package names, one per line", "example_requirements": "ggplot2\ndplyr", "example_code": 'print("Hello from R!")\nx <- c(1,2,3)\nprint(mean(x))' }, "rust": { "aliases": ["rust"], "requirements_format": 'Cargo.toml [dependencies] format (crate = "version")', "example_requirements": 'serde = "1.0"\nrand = "0.8"', "example_code": 'fn main() {\n println!("Hello from Rust!");\n}' }, "java": { "aliases": ["java"], "requirements_format": "Maven coordinates (groupId:artifactId:version)", "example_requirements": "com.google.guava:guava:32.1.2-jre", "example_code": 'public class Main {\n public static void main(String[] args) {\n System.out.println("Hello from Java!");\n }\n}' }, "swift": { "aliases": ["swift"], "requirements_format": "name,url,version per line (for Swift Package Manager)", "example_requirements": "", "example_code": 'print("Hello from Swift!")' }, "javascript": { "aliases": ["javascript", "js", "node"], "requirements_format": "package@version or package:version per line", "example_requirements": "lodash@4.17.21\naxios@1.6.0", "example_code": 'console.log("Hello from JavaScript!");' } } }, indent=2, ensure_ascii=False) # Health check endpoint async def health_check(request): """Health check endpoint.""" return JSONResponse({ "status": "healthy", "service": "code-executor", "languages": _list_languages() }) # OAuth protected resource metadata endpoint async def oauth_protected_resource(request): """Return OAuth 2.0 Protected Resource Metadata.""" return JSONResponse(get_protected_resource_metadata()) async def _monitor_executions(): """Background task to monitor executions and environments.""" while True: try: await asyncio.sleep(5) # Check every 5 seconds # Update status for all running executions for exec_id in list(_executions.keys()): try: execution = _executions.get(exec_id) if execution and execution["status"] == "running": _update_execution_status(exec_id) except Exception as e: logger.warning(f"Error monitoring execution {exec_id}: {e}") # Update status for all setting_up environments for env_id in list(_environments.keys()): try: env = _environments.get(env_id) if env and env["status"] == "setting_up": _update_environment_status(env_id) except Exception as e: logger.warning(f"Error monitoring environment {env_id}: {e}") # Cleanup old completed executions (older than 1 hour) now = time.time() for exec_id in list(_executions.keys()): try: execution = _executions.get(exec_id) if execution and execution["status"] in ("completed", "failed", "stopped"): age = now - execution["start_time"] if age > 3600: # 1 hour logger.info(f"Cleaning up old execution {exec_id} (age: {age:.0f}s)") _cleanup_execution(exec_id) except Exception as e: logger.warning(f"Error cleaning up execution {exec_id}: {e}") # Cleanup old failed environments (older than 1 hour) for env_id in list(_environments.keys()): try: env = _environments.get(env_id) if env and env["status"] == "failed": age = now - env["start_time"] if age > 3600: # 1 hour logger.info(f"Cleaning up failed environment {env_id} (age: {age:.0f}s)") _cleanup_environment(env_id) except Exception as e: logger.warning(f"Error cleaning up environment {env_id}: {e}") except asyncio.CancelledError: break except Exception as e: logger.error(f"Error in monitor: {e}") def create_app(): """Create the combined ASGI app with health check and MCP.""" from contextlib import asynccontextmanager from starlette.middleware import Middleware from starlette.middleware.trustedhost import TrustedHostMiddleware @asynccontextmanager async def lifespan(app): """Manage MCP session manager lifecycle.""" # Start background monitor monitor_task = asyncio.create_task(_monitor_executions()) async with mcp.session_manager.run(): logger.info(f"Code Executor MCP Server started on port {MCP_PORT}") logger.info(f"Supported languages: {_list_languages()}") yield # Stop background monitor monitor_task.cancel() try: await monitor_task except asyncio.CancelledError: pass # Cleanup all remaining executions for exec_id in list(_executions.keys()): _cleanup_execution(exec_id) # Get the MCP ASGI app mcp_app = mcp.streamable_http_app() # Create Starlette app with health endpoint and MCP app = Starlette( routes=[ Route("/health", health_check, methods=["GET"]), Route("/.well-known/oauth-protected-resource", oauth_protected_resource, methods=["GET"]), ], lifespan=lifespan, middleware=[ Middleware(TrustedHostMiddleware, allowed_hosts=["*"]), Middleware(AcceptHeaderMiddleware), Middleware(OAuthMiddleware), ] ) # Mount MCP app at root app.mount("/", mcp_app) return app def main(): """Main entry point.""" logger.info(f"Starting Code Executor MCP Server on port {MCP_PORT}") app = create_app() uvicorn.run(app, host="0.0.0.0", port=MCP_PORT) if __name__ == "__main__": main()