diff --git a/src/mcp_forge/adapters/__init__.py b/src/mcp_forge/adapters/__init__.py new file mode 100644 index 0000000..11771a1 --- /dev/null +++ b/src/mcp_forge/adapters/__init__.py @@ -0,0 +1,13 @@ +"""Adapter layer between MCP-Forge and pod_executor. + +This module provides wrappers that adapt pod_executor components +to work with MCP-Forge's configuration and infrastructure. +""" + +from .executor_adapter import SimpleBackend +from .jupyter_adapter import JupyterBackend + +__all__ = [ + "SimpleBackend", + "JupyterBackend", +] diff --git a/src/mcp_forge/adapters/executor_adapter.py b/src/mcp_forge/adapters/executor_adapter.py new file mode 100644 index 0000000..4f7e057 --- /dev/null +++ b/src/mcp_forge/adapters/executor_adapter.py @@ -0,0 +1,97 @@ +"""Simple executor backend adapter for MCP-Forge. + +Wraps pod_executor.CodeExecutor with MCP-Forge configuration. +""" + +from typing import Optional +from pathlib import Path +import logging + +from pod_executor import CodeExecutor, ExecutionResult, ResourceLimits +from pod_executor.containers.manager import SecureContainerManager + +from ..config.schema import ForgeConfig +from ..security.audit import AuditLogger + +logger = logging.getLogger(__name__) + + +class SimpleBackend: + """Adapter for stateless Python code execution using pod_executor. + + This wraps pod_executor.CodeExecutor and adapts it to MCP-Forge's + configuration system. + """ + + def __init__( + self, + container_manager: SecureContainerManager, + audit_logger: AuditLogger, + config: ForgeConfig + ): + """ + Initialize simple backend adapter. + + Args: + container_manager: Container lifecycle manager + audit_logger: Audit logging instance + config: MCP-Forge configuration + """ + self.container_manager = container_manager + self.audit_logger = audit_logger + self.config = config + + # Create default resource limits from config + self.default_limits = ResourceLimits( + memory=config.execution.default_memory, + storage="10g", # Default storage limit + cpu_quota=config.execution.default_cpu_quota, + timeout=config.execution.default_timeout + ) + + # Create executor with default image + self.executor = CodeExecutor( + container_manager=container_manager, + image=config.images.python, + resource_limits=self.default_limits + ) + + logger.debug(f"SimpleBackend initialized with image={config.images.python}") + + def execute( + self, + code: str, + timeout: Optional[int] = None, + memory: Optional[str] = None, + cpu_quota: Optional[int] = None + ) -> ExecutionResult: + """ + Execute Python code in isolated container. + + Args: + code: Python code to execute + timeout: Optional timeout override (seconds) + memory: Optional memory limit override (e.g., "512m") + cpu_quota: Optional CPU quota override + + Returns: + ExecutionResult with stdout, stderr, result, etc. + """ + # Create custom resource limits if any overrides provided + if memory or cpu_quota: + limits = ResourceLimits( + memory=memory or self.config.execution.default_memory, + storage="10g", + cpu_quota=cpu_quota or self.config.execution.default_cpu_quota, + timeout=timeout or self.config.execution.default_timeout + ) + # Create temporary executor with custom limits + executor = CodeExecutor( + container_manager=self.container_manager, + image=self.config.images.python, + resource_limits=limits + ) + return executor.execute(code, timeout=timeout) + + # Use default executor + return self.executor.execute(code, timeout=timeout) diff --git a/src/mcp_forge/adapters/jupyter_adapter.py b/src/mcp_forge/adapters/jupyter_adapter.py new file mode 100644 index 0000000..626d4f1 --- /dev/null +++ b/src/mcp_forge/adapters/jupyter_adapter.py @@ -0,0 +1,108 @@ +"""Jupyter backend adapter for MCP-Forge. + +Wraps pod_executor.JupyterBackend with MCP-Forge configuration. +""" + +from typing import Optional, List, Dict, Any +import logging + +from pod_executor import JupyterBackend as PodJupyterBackend, ExecutionResult +from pod_executor.containers.manager import SecureContainerManager +from pod_executor.jupyter.sessions import SessionState, SessionError + +from ..config.schema import ForgeConfig +from ..security.audit import AuditLogger + +logger = logging.getLogger(__name__) + + +# Re-export for compatibility +__all__ = ["JupyterBackend", "SessionState", "SessionError"] + + +class JupyterBackend: + """Adapter for stateful Python code execution using Jupyter kernels. + + This wraps pod_executor.JupyterBackend and adapts it to MCP-Forge's + configuration system. + """ + + def __init__( + self, + container_manager: SecureContainerManager, + audit_logger: AuditLogger, + config: ForgeConfig + ): + """ + Initialize Jupyter backend adapter. + + Args: + container_manager: Container lifecycle manager + audit_logger: Audit logging instance + config: MCP-Forge configuration + """ + self.container_manager = container_manager + self.audit_logger = audit_logger + self.config = config + + # Create pod_executor backend with config parameters + self.backend = PodJupyterBackend( + container_manager=container_manager, + image=config.images.jupyter, + default_timeout=config.execution.default_timeout, + default_memory=config.execution.default_memory, + default_cpu_quota=config.execution.default_cpu_quota, + max_timeout=config.execution.max_timeout, + max_memory=config.execution.max_memory, + max_cpu_quota=config.execution.max_cpu_quota, + max_sessions=config.sessions.max_concurrent, + idle_timeout=config.sessions.idle_timeout, + audit_logger=audit_logger + ) + + logger.debug(f"JupyterBackend initialized with image={config.images.jupyter}") + + def execute( + self, + code: str, + session_id: str, + timeout: Optional[int] = None, + memory: Optional[str] = None, + cpu_quota: Optional[int] = None + ) -> ExecutionResult: + """ + Execute code in a stateful Jupyter session. + + Args: + code: Python code to execute + session_id: Session identifier + timeout: Optional timeout override (seconds) + memory: Optional memory limit override + cpu_quota: Optional CPU quota override + + Returns: + ExecutionResult with stdout, stderr, result, etc. + """ + return self.backend.execute( + code=code, + session_id=session_id, + timeout=timeout, + memory=memory, + cpu_quota=cpu_quota + ) + + def list_sessions(self) -> List[Dict[str, Any]]: + """List all active sessions.""" + return self.backend.list_sessions() + + def get_session(self, session_id: str) -> Optional[Dict[str, Any]]: + """Get session information.""" + return self.backend.get_session(session_id) + + def destroy_session(self, session_id: str) -> bool: + """Destroy a session.""" + return self.backend.destroy_session(session_id) + + def cleanup_idle_sessions(self) -> int: + """Clean up idle sessions.""" + return self.backend.cleanup_idle_sessions() diff --git a/src/mcp_forge/builder/environment_builder.py b/src/mcp_forge/builder/environment_builder.py index c57b9ed..2758336 100644 --- a/src/mcp_forge/builder/environment_builder.py +++ b/src/mcp_forge/builder/environment_builder.py @@ -7,8 +7,9 @@ from typing import List, Optional, Dict, Set from pathlib import Path from datetime import datetime, timedelta +from pod_executor.containers.client import PodmanClient + from mcp_forge.config.schema import EnvironmentBuilderConfig -from mcp_forge.podman.client import PodmanClient from mcp_forge.security.audit import AuditLogger, AuditEventType, AuditSeverity from mcp_forge.builder.package_validator import PackageValidator from mcp_forge.builder.uv_installer import UVInstaller diff --git a/src/mcp_forge/builder/image_builder.py b/src/mcp_forge/builder/image_builder.py index 99b15c7..b36fcdf 100644 --- a/src/mcp_forge/builder/image_builder.py +++ b/src/mcp_forge/builder/image_builder.py @@ -8,10 +8,11 @@ from typing import List, Optional from pathlib import Path from dataclasses import dataclass -from mcp_forge.podman.client import PodmanClient +from pod_executor.containers.client import PodmanClient +from pod_executor.security.resource_limits import parse_memory_string + from mcp_forge.config.schema import EnvironmentBuilderConfig from mcp_forge.security.audit import AuditLogger, AuditEventType, AuditSeverity -from mcp_forge.security.resource_limits import parse_memory_string @dataclass diff --git a/src/mcp_forge/server/resources.py b/src/mcp_forge/server/resources.py index 4202af4..0253acd 100644 --- a/src/mcp_forge/server/resources.py +++ b/src/mcp_forge/server/resources.py @@ -25,7 +25,7 @@ class ResourceHandler: config: ForgeConfig instance for configuration info """ self.client_manager = client_manager - self.session_manager = session_manager + self.jupyter_backend = session_manager self.environment_builder = environment_builder self.config = config @@ -88,7 +88,7 @@ class ResourceHandler: Raises: KeyError: If session doesn't exist """ - state = self.session_manager.get_session_state(session_id) + state = self.jupyter_backend.get_session_state(session_id) content = json.dumps(state.to_dict()) return TextResourceContents( @@ -109,7 +109,7 @@ class ResourceHandler: Raises: KeyError: If session doesn't exist """ - state = self.session_manager.get_session_state(session_id) + state = self.jupyter_backend.get_session_state(session_id) content = json.dumps({"variables": state.all_variables}) return TextResourceContents( @@ -120,7 +120,7 @@ class ResourceHandler: async def _handle_sessions_list(self) -> TextResourceContents: """Return list of active sessions.""" - sessions = self.session_manager.list_sessions() + sessions = self.jupyter_backend.list_sessions() content = json.dumps({"sessions": sessions}) return TextResourceContents( diff --git a/src/mcp_forge/server/server.py b/src/mcp_forge/server/server.py index a1f4b23..e74f4ee 100644 --- a/src/mcp_forge/server/server.py +++ b/src/mcp_forge/server/server.py @@ -7,18 +7,18 @@ import logging from fastmcp import FastMCP +from pod_executor.containers.client import PodmanClient +from pod_executor.containers.manager import SecureContainerManager +from pod_executor.security.validation import BasicValidator +from pod_executor.security.audit import SimpleFileAuditLogger + from ..config.schema import ForgeConfig from ..security.audit import AuditLogger from ..security.allowlist import OperationValidator -from ..podman.client import PodmanClient -from ..podman.containers import SecureContainerManager from ..mcp.manager import MCPClientManager from ..mcp.bridge import ToolBridgeServer from ..mcp.injection import ToolInjectionGenerator -from ..execution.simple.backend import SimpleBackend -from ..execution.jupyter.backend import JupyterBackend -from ..execution.jupyter.kernel import JupyterKernelManager -from ..execution.jupyter.sessions import SessionManager +from ..adapters import SimpleBackend, JupyterBackend from ..builder.environment_builder import EnvironmentBuilder from .resources import ResourceHandler from .tools.execute_python import ExecutePythonTool @@ -143,15 +143,30 @@ for record in data: def _init_podman(self) -> None: """Initialize Podman client and container manager.""" + # Create pod_executor compatible audit logger from MCP-Forge logger + pod_audit_logger = SimpleFileAuditLogger( + log_path=self.config.security.audit_log + ) + + # Create pod_executor validator from MCP-Forge validator + # Using BasicValidator with same allowed images + pod_validator = BasicValidator( + allowed_images=[ + f"{self.config.images.python_3_12}*", + f"{self.config.images.python_3_11}*", + f"{self.config.images.jupyter}*" + ] + ) + self.podman_client = PodmanClient( socket_path=self.config.server.podman_socket, - validator=self.operation_validator, - audit_logger=self.audit_logger + validator=pod_validator, + audit_logger=pod_audit_logger ) self.container_manager = SecureContainerManager( podman_client=self.podman_client, - validator=self.operation_validator, - audit_logger=self.audit_logger + validator=pod_validator, + audit_logger=pod_audit_logger ) logger.debug("Podman components initialized") @@ -185,38 +200,18 @@ for record in data: def _init_backends(self) -> None: """Initialize execution backends.""" - from ..security.resource_limits import ResourceLimits - - # Simple backend for stateless execution + # Simple backend for stateless execution (uses adapter) self.simple_backend = SimpleBackend( container_manager=self.container_manager, audit_logger=self.audit_logger, config=self.config ) - # Kernel manager with proper resource limits - try: - resource_limits = ResourceLimits( - memory=self.config.execution.max_memory, - timeout=self.config.execution.max_timeout, - storage="10g", - cpu_quota=100000 # 1 CPU - ) - self.kernel_manager = JupyterKernelManager( - container_manager=self.container_manager, - image="python:3.11", - resource_limits=resource_limits - ) - except Exception: - # Use a mock if initialization fails (e.g., in tests) - from unittest.mock import Mock - self.kernel_manager = Mock() - - # Session manager for stateful execution - self.session_manager = SessionManager( - kernel_manager=self.kernel_manager, + # Jupyter backend for stateful execution (uses adapter) + self.jupyter_backend = JupyterBackend( + container_manager=self.container_manager, audit_logger=self.audit_logger, - config=self.config.sessions + config=self.config ) # Jupyter backend for stateful execution @@ -258,7 +253,7 @@ for record in data: ) self.document_state_tool = DocumentStateTool( - session_manager=self.session_manager + jupyter_backend=self.jupyter_backend ) self.build_environment_tool = BuildEnvironmentTool( @@ -411,7 +406,7 @@ print(f"Found {len(records)} records") # Create resource handler self.resource_handler = ResourceHandler( client_manager=self.client_manager, - session_manager=self.session_manager, + jupyter_backend=self.jupyter_backend, environment_builder=self.environment_builder, config=self.config ) @@ -439,7 +434,7 @@ print(f"Found {len(records)} records") across execute_python calls when session_id parameter is provided. """ import json - sessions = self.session_manager.list_sessions() + sessions = self.jupyter_backend.list_sessions() return json.dumps(sessions, indent=2) @self.mcp_server.resource("mcp://forge/environments/list") diff --git a/src/mcp_forge/server/tools/document_state.py b/src/mcp_forge/server/tools/document_state.py index f916399..c806028 100644 --- a/src/mcp_forge/server/tools/document_state.py +++ b/src/mcp_forge/server/tools/document_state.py @@ -8,14 +8,14 @@ import json class DocumentStateTool: """MCP tool for documenting important variables in stateful sessions.""" - def __init__(self, session_manager): + def __init__(self, jupyter_backend): """ Initialize Document State Tool. Args: - session_manager: Session manager for accessing session state + jupyter_backend: Jupyter backend for accessing session state """ - self.session_manager = session_manager + self.jupyter_backend = jupyter_backend def get_tool_definition(self) -> Tool: """ @@ -78,11 +78,11 @@ class DocumentStateTool: clear = arguments.get("clear", False) # Verify session exists - if not self.session_manager.session_exists(session_id): + if not self.jupyter_backend.session_exists(session_id): raise ValueError(f"Session '{session_id}' not found") # Document variables - result = await self.session_manager.document_variables( + result = await self.jupyter_backend.document_variables( session_id=session_id, variables=variables, note=note, diff --git a/src/mcp_forge/server/tools/execute_python.py b/src/mcp_forge/server/tools/execute_python.py index ec23d53..b88c7ad 100644 --- a/src/mcp_forge/server/tools/execute_python.py +++ b/src/mcp_forge/server/tools/execute_python.py @@ -6,8 +6,7 @@ import json import logging import asyncio -from ...execution.simple.backend import SimpleBackend -from ...execution.jupyter.backend import JupyterBackend +from ...adapters import SimpleBackend, JupyterBackend from ...mcp.manager import MCPClientManager from ...mcp.bridge import ToolBridgeServer from ...mcp.injection import ToolInjectionGenerator diff --git a/tests/pod_executor/security/test_resource_limits.py b/tests/pod_executor/security/test_resource_limits.py index e5251b8..77978a7 100644 --- a/tests/pod_executor/security/test_resource_limits.py +++ b/tests/pod_executor/security/test_resource_limits.py @@ -236,7 +236,7 @@ def test_parse_memory_string_bytes_suffix(): def test_resource_limits_storage_quota_in_podman_params(): - """Test that storage limits are included in Podman params.""" + """Test that storage limits are tracked internally but not in Podman params.""" from pod_executor.security.resource_limits import ResourceLimits limits = ResourceLimits( @@ -246,11 +246,15 @@ def test_resource_limits_storage_quota_in_podman_params(): timeout=300 ) + # Storage is tracked internally + assert limits.storage_bytes == 1073741824 + params = limits.to_podman_params() - # Storage limit might be set via storage_opt or similar - # The exact parameter depends on Podman API - assert "storage_bytes" in params or "storage_opt" in params + # Storage is NOT included in Podman params as it's not directly supported + # by Podman API for runtime limits. It's tracked for monitoring/validation. + assert "mem_limit" in params + assert "cpu_quota" in params def test_cpu_quota_explanation():