Team Ai
Apppublic

openenv/echo_env

sourceHugging Faceupdated 1d agoView on Hugging Face
6likes
providers.py728 linesDownload Raw Back to runtime
1# SPDX-License-Identifier: BSD-3-Clause2 3"""4Container provider abstractions for running environment servers.5 6This module provides a pluggable architecture for different container providers7(local Docker, Kubernetes, cloud providers, etc.) to be used with EnvClient.8"""9 10from __future__ import annotations11 12from abc import ABC, abstractmethod13from typing import Any, Dict, Optional, Sequence, TypeVar14 15_ContainerProviderT = TypeVar("_ContainerProviderT", bound="ContainerProvider")16 17 18class ContainerProvider(ABC):19    """20    Abstract base class for container providers.21 22    Providers implement this interface to support different container platforms:23    - LocalDockerProvider: Runs containers on local Docker daemon24    - KubernetesProvider: Runs containers in Kubernetes cluster25    - FargateProvider: Runs containers on AWS Fargate26    - CloudRunProvider: Runs containers on Google Cloud Run27 28    The provider manages a single container lifecycle and provides the base URL29    for connecting to it.30 31    Examples:32 33        ```python34        provider = LocalDockerProvider()35        base_url = provider.start_container("echo-env:latest")36        print(base_url)  # http://localhost:800037        # Use the environment via base_url38        provider.stop_container()39        ```40    """41 42    @abstractmethod43    def start_container(44        self,45        image: str,46        port: Optional[int] = None,47        env_vars: Optional[Dict[str, str]] = None,48        **kwargs: Any,49    ) -> str:50        """51        Start a container from the specified image.52 53        Args:54            image (`str`):55                Provider-specific container *source* identifier. For56                container-based providers this is a registry image name (e.g.57                `"echo-env:latest"`); other providers may map it to a58                provider-specific source (see the provider's documentation).59            port (`int`, *optional*):60                Port to expose. If `None`, the provider chooses.61            env_vars (`dict`, *optional*):62                Environment variables to pass to container.63            **kwargs:64                Provider-specific options.65 66        Returns:67            `str`: Base URL to connect to the container (e.g., `"http://localhost:8000"`).68 69        Raises:70            RuntimeError: If container fails to start.71        """72        pass73 74    @abstractmethod75    def stop_container(self) -> None:76        """77        Stop and remove the running container.78 79        This cleans up the container that was started by start_container().80        """81        pass82 83    @abstractmethod84    def wait_for_ready(self, base_url: str, timeout_s: float = 30.0) -> None:85        """86        Wait for the container to be ready to accept requests.87 88        This typically polls the /health endpoint until it returns 200.89 90        Args:91            base_url (`str`):92                Base URL of the container.93            timeout_s (`float`, *optional*, defaults to `30.0`):94                Maximum time to wait in seconds.95 96        Raises:97            TimeoutError: If container doesn't become ready in time.98        """99        pass100 101    def close(self) -> None:102        """103        Release provider-held resources (e.g. SDK clients, connections).104 105        Defaults to a no-op so existing providers are unaffected. Providers that106        hold external resources beyond the container itself (such as a cloud SDK107        client) override this to release them; it is also invoked on context-108        manager exit. Lightweight providers need not override it.109        """110        pass111 112    def __enter__(self: _ContainerProviderT) -> _ContainerProviderT:113        return self114 115    def __exit__(self, exc_type, exc, tb) -> Optional[bool]:116        self.close()117        return None118 119 120class LocalDockerProvider(ContainerProvider):121    """122    Container provider for local Docker daemon.123 124    This provider runs containers on the local machine using Docker.125    Useful for development and testing.126 127    Examples:128 129        ```python130        provider = LocalDockerProvider()131        base_url = provider.start_container("echo-env:latest")132        # Container running on http://localhost:<random-port>133        provider.stop_container()134        ```135    """136 137    def __init__(self):138        """Initialize the local Docker provider."""139        self._container_id: Optional[str] = None140        self._container_name: Optional[str] = None141 142        # Check if Docker is available143        import subprocess144 145        try:146            subprocess.run(147                ["docker", "version"],148                check=True,149                capture_output=True,150                timeout=5,151            )152        except (153            subprocess.CalledProcessError,154            FileNotFoundError,155            subprocess.TimeoutExpired,156        ):157            raise RuntimeError(158                "Docker is not available. Please install Docker Desktop or Docker Engine."159            )160 161    def start_container(162        self,163        image: str,164        port: Optional[int] = None,165        env_vars: Optional[Dict[str, str]] = None,166        **kwargs: Any,167    ) -> str:168        """169        Start a Docker container locally.170 171        Args:172            image (`str`):173                Docker image name.174            port (`int`, *optional*):175                Port to expose. If `None`, finds an available port.176            env_vars (`dict`, *optional*):177                Environment variables for the container.178            **kwargs:179                Additional Docker run options.180 181        Returns:182            `str`: Base URL to connect to the container.183        """184        import subprocess185        import time186 187        # Find available port if not specified188        if port is None:189            port = self._find_available_port()190 191        # Generate container name192        self._container_name = self._generate_container_name(image)193 194        # Build docker run command195        cmd = [196            "docker",197            "run",198            "-d",  # Detached199            "--name",200            self._container_name,201            "-p",202            f"{port}:8000",  # Map port203        ]204 205        # Add environment variables206        if env_vars:207            for key, value in env_vars.items():208                cmd.extend(["-e", f"{key}={value}"])209 210        # Add image211        cmd.append(image)212 213        # Run container214        try:215            result = subprocess.run(cmd, capture_output=True, text=True, check=True)216            self._container_id = result.stdout.strip()217        except subprocess.CalledProcessError as e:218            error_msg = f"Failed to start Docker container.\nCommand: {' '.join(cmd)}\nExit code: {e.returncode}\nStderr: {e.stderr}\nStdout: {e.stdout}"219            raise RuntimeError(error_msg) from e220 221        # Wait a moment for container to start222        time.sleep(1)223 224        base_url = f"http://localhost:{port}"225        return base_url226 227    def stop_container(self) -> None:228        """229        Stop and remove the Docker container.230        """231        if self._container_id is None:232            return233 234        import subprocess235 236        try:237            # Stop container238            subprocess.run(239                ["docker", "stop", self._container_id],240                capture_output=True,241                check=True,242                timeout=10,243            )244 245            # Remove container246            subprocess.run(247                ["docker", "rm", self._container_id],248                capture_output=True,249                check=True,250                timeout=10,251            )252        except subprocess.CalledProcessError:253            # Container might already be stopped/removed254            pass255        finally:256            self._container_id = None257            self._container_name = None258 259    def wait_for_ready(self, base_url: str, timeout_s: float = 30.0) -> None:260        """261        Wait for container to be ready by polling /health endpoint.262 263        Args:264            base_url (`str`):265                Base URL of the container.266            timeout_s (`float`, *optional*, defaults to `30.0`):267                Maximum time to wait in seconds.268 269        Raises:270            TimeoutError: If container doesn't become ready.271        """272        import time273 274        import requests275 276        start_time = time.time()277        health_url = f"{base_url}/health"278 279        # Bypass proxy for localhost to avoid proxy issues280        proxies = {"http": None, "https": None}281 282        while time.time() - start_time < timeout_s:283            try:284                response = requests.get(health_url, timeout=2.0, proxies=proxies)285                if response.status_code == 200:286                    return287            except requests.RequestException:288                pass289 290            time.sleep(0.5)291 292        raise TimeoutError(293            f"Container at {base_url} did not become ready within {timeout_s}s"294        )295 296    def _find_available_port(self) -> int:297        """298        Find an available port on localhost.299 300        Returns:301            `int`: An available port number.302        """303        import socket304 305        with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:306            s.bind(("", 0))307            s.listen(1)308            port = s.getsockname()[1]309        return port310 311    def _generate_container_name(self, image: str) -> str:312        """313        Generate a unique container name based on image name and timestamp.314 315        Args:316            image (`str`):317                Docker image name.318 319        Returns:320            `str`: A unique container name.321        """322        import time323 324        clean_image = image.split("/")[-1].split(":")[0]325        timestamp = int(time.time() * 1000)326        return f"{clean_image}-{timestamp}"327 328 329class DockerSwarmProvider(ContainerProvider):330    """331    Container provider that uses Docker Swarm services for local concurrency.332 333    This provider creates a replicated Swarm service backed by the local Docker334    engine. The built-in load-balancer fans requests across the replicas,335    allowing multiple container instances to run concurrently on the developer336    workstation (mirroring the workflow described in the Docker stack docs).337    """338 339    def __init__(340        self,341        *,342        auto_init_swarm: bool = True,343        overlay_network: Optional[str] = None,344    ):345        """346        Args:347            auto_init_swarm (`bool`, *optional*, defaults to `True`):348                Whether to call `docker swarm init` when Swarm is not active.349                Otherwise, the user must manually initialize Swarm.350            overlay_network (`str`, *optional*):351                Overlay network name for the service. When provided, the network352                is created with `docker network create --driver overlay --attachable`353                if it does not already exist.354        """355        self._service_name: Optional[str] = None356        self._service_id: Optional[str] = None357        self._published_port: Optional[int] = None358        self._overlay_network = overlay_network359        self._auto_init_swarm = auto_init_swarm360 361        self._ensure_docker_available()362        self._ensure_swarm_initialized()363        if self._overlay_network:364            self._ensure_overlay_network(self._overlay_network)365 366    def start_container(367        self,368        image: str,369        port: Optional[int] = None,370        env_vars: Optional[Dict[str, str]] = None,371        **kwargs: Any,372    ) -> str:373        """374        Start (or scale) a Swarm service for the given image.375 376        Args:377            image (`str`):378                Docker image name.379            port (`int`, *optional*):380                Port to expose. If `None`, finds an available port.381            env_vars (`dict`, *optional*):382                Environment variables for the container.383            replicas (`int`, *optional*, defaults to `2`):384                Number of container replicas.385            cpu_limit (`float` or `str`, *optional*):386                CPU limit passed to `--limit-cpu`.387            memory_limit (`str`, *optional*):388                Memory limit passed to `--limit-memory`.389            constraints (`Sequence[str]`, *optional*):390                Placement constraints.391            labels (`dict`, *optional*):392                Service labels.393            command (`Sequence[str]` or `str`, *optional*):394                Override container command.395 396        Returns:397            `str`: Base URL to connect to the service.398        """399        import shlex400        import subprocess401        import time402 403        allowed_kwargs = {404            "replicas",405            "cpu_limit",406            "memory_limit",407            "constraints",408            "labels",409            "command",410        }411        unknown = set(kwargs) - allowed_kwargs412        if unknown:413            raise ValueError(f"Unsupported kwargs for DockerSwarmProvider: {unknown}")414 415        replicas = int(kwargs.get("replicas", 2))416        cpu_limit = kwargs.get("cpu_limit")417        memory_limit = kwargs.get("memory_limit")418        constraints: Optional[Sequence[str]] = kwargs.get("constraints")419        labels: Optional[Dict[str, str]] = kwargs.get("labels")420        command_override = kwargs.get("command")421 422        if port is None:423            port = self._find_available_port()424 425        self._service_name = self._generate_service_name(image)426        self._published_port = port427 428        cmd = [429            "docker",430            "service",431            "create",432            "--detach",433            "--name",434            self._service_name,435            "--replicas",436            str(max(1, replicas)),437            "--publish",438            f"{port}:8000",439        ]440 441        if self._overlay_network:442            cmd.extend(["--network", self._overlay_network])443 444        if env_vars:445            for key, value in env_vars.items():446                cmd.extend(["--env", f"{key}={value}"])447 448        if cpu_limit is not None:449            cmd.extend(["--limit-cpu", str(cpu_limit)])450 451        if memory_limit is not None:452            cmd.extend(["--limit-memory", str(memory_limit)])453 454        if constraints:455            for constraint in constraints:456                cmd.extend(["--constraint", constraint])457 458        if labels:459            for key, value in labels.items():460                cmd.extend(["--label", f"{key}={value}"])461 462        cmd.append(image)463 464        if command_override:465            if isinstance(command_override, str):466                cmd.extend(shlex.split(command_override))467            else:468                cmd.extend(command_override)469 470        try:471            result = subprocess.run(472                cmd,473                capture_output=True,474                text=True,475                check=True,476            )477            self._service_id = result.stdout.strip()478        except subprocess.CalledProcessError as e:479            error_msg = (480                "Failed to start Docker Swarm service.\n"481                f"Command: {' '.join(cmd)}\n"482                f"Exit code: {e.returncode}\n"483                f"Stdout: {e.stdout}\n"484                f"Stderr: {e.stderr}"485            )486            raise RuntimeError(error_msg) from e487 488        # Give Swarm a brief moment to schedule the tasks.489        time.sleep(1.0)490 491        return f"http://localhost:{port}"492 493    def stop_container(self) -> None:494        """495        Remove the Swarm service (and keep the Swarm manager running).496        """497        if not self._service_name:498            return499 500        import subprocess501 502        try:503            subprocess.run(504                ["docker", "service", "rm", self._service_name],505                capture_output=True,506                check=True,507                timeout=10,508            )509        except subprocess.CalledProcessError:510            # Service may already be gone; ignore.511            pass512        finally:513            self._service_name = None514            self._service_id = None515            self._published_port = None516 517    def wait_for_ready(self, base_url: str, timeout_s: float = 30.0) -> None:518        """519        Wait for at least one replica to become healthy by polling /health.520 521        With Swarm's load balancer, requests round-robin across replicas,522        so this only verifies that at least one replica is responding. Some523        replicas may still be starting when this returns.524        """525        import time526 527        import requests528 529        deadline = time.time() + timeout_s530        health_url = f"{base_url}/health"531 532        # Bypass proxy for localhost to avoid proxy issues533        proxies = {"http": None, "https": None}534 535        while time.time() < deadline:536            try:537                response = requests.get(health_url, timeout=2.0, proxies=proxies)538                if response.status_code == 200:539                    return540            except requests.RequestException:541                pass542 543            time.sleep(0.5)544 545        raise TimeoutError(546            f"Swarm service at {base_url} did not become ready within {timeout_s}s"547        )548 549    def _ensure_docker_available(self) -> None:550        import subprocess551 552        try:553            subprocess.run(554                ["docker", "version"],555                check=True,556                capture_output=True,557                timeout=5,558            )559        except (560            subprocess.CalledProcessError,561            FileNotFoundError,562            subprocess.TimeoutExpired,563        ) as exc:564            raise RuntimeError(565                "Docker is not available. Please install Docker Desktop or Docker Engine."566            ) from exc567 568    def _ensure_swarm_initialized(self) -> None:569        import subprocess570 571        try:572            result = subprocess.run(573                ["docker", "info", "--format", "{{.Swarm.LocalNodeState}}"],574                capture_output=True,575                text=True,576                check=True,577                timeout=5,578            )579            state = result.stdout.strip().lower()580            if state == "active":581                return582        except subprocess.CalledProcessError:583            state = "unknown"584 585        if not self._auto_init_swarm:586            raise RuntimeError(587                f"Docker Swarm is not active (state={state}). Enable Swarm manually or pass auto_init_swarm=True."588            )589 590        try:591            subprocess.run(592                ["docker", "swarm", "init"],593                check=True,594                capture_output=True,595                timeout=10,596            )597        except subprocess.CalledProcessError as e:598            raise RuntimeError("Failed to initialize Docker Swarm") from e599 600    def _ensure_overlay_network(self, network: str) -> None:601        import subprocess602 603        inspect = subprocess.run(604            ["docker", "network", "inspect", network],605            capture_output=True,606            text=True,607            check=False,608        )609        if inspect.returncode == 0:610            return611 612        try:613            subprocess.run(614                [615                    "docker",616                    "network",617                    "create",618                    "--driver",619                    "overlay",620                    "--attachable",621                    network,622                ],623                check=True,624                capture_output=True,625                timeout=10,626            )627        except subprocess.CalledProcessError as e:628            raise RuntimeError(f"Failed to create overlay network '{network}'") from e629 630    def _find_available_port(self) -> int:631        import socket632 633        with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:634            s.bind(("", 0))635            s.listen(1)636            port = s.getsockname()[1]637        return port638 639    def _generate_service_name(self, image: str) -> str:640        import time641 642        clean_image = image.split("/")[-1].split(":")[0]643        timestamp = int(time.time() * 1000)644        return f"{clean_image}-swarm-{timestamp}"645 646 647class KubernetesProvider(ContainerProvider):648    """649    Planned container provider for Kubernetes clusters.650 651    Not yet implemented: this is a placeholder for the planned Kubernetes652    backend and does not implement the abstract `ContainerProvider` methods, so653    it cannot be instantiated. Use `LocalDockerProvider`, `DockerSwarmProvider`,654    `DaytonaProvider`, or `ACASandboxProvider` instead.655    """656 657    pass658 659 660class RuntimeProvider(ABC):661    """662    Abstract base class for runtime providers that are not container providers.663    Providers implement this interface to support different runtime platforms:664    - UVProvider: Runs environments via `uv run`665 666    The provider manages a single runtime lifecycle and provides the base URL667    for connecting to it.668 669    Examples:670 671        ```python672        provider = UVProvider(project_path="/path/to/env")673        base_url = provider.start()674        print(base_url)  # http://localhost:8000675        provider.stop()676        ```677    """678 679    @abstractmethod680    def start(681        self,682        port: Optional[int] = None,683        env_vars: Optional[Dict[str, str]] = None,684        **kwargs: Any,685    ) -> str:686        """687        Start the runtime.688 689        Args:690            port (`int`, *optional*):691                Port to expose. If `None`, the provider chooses.692            env_vars (`dict`, *optional*):693                Environment variables for the runtime.694            **kwargs:695                Additional runtime options.696 697        Returns:698            `str`: Base URL to connect to the runtime.699        """700 701    @abstractmethod702    def stop(self) -> None:703        """704        Stop the runtime.705        """706        pass707 708    @abstractmethod709    def wait_for_ready(self, timeout_s: float = 30.0) -> None:710        """711        Wait for the runtime to be ready to accept requests.712        """713        pass714 715    def __enter__(self) -> "RuntimeProvider":716        """717        Enter the runtime provider.718        """719        self.start()720        return self721 722    def __exit__(self, exc_type, exc, tb) -> None:723        """724        Exit the runtime provider.725        """726        self.stop()727        return False728