Team Ai
Apppublic

openenv/coding_env

sourceHugging Faceupdated 3mo agoView on Hugging Face
21likes
providers.py670 linesDownload Raw Back to runtime
1# Copyright (c) Meta Platforms, Inc. and affiliates.2# All rights reserved.3#4# This source code is licensed under the BSD-style license found in the5# LICENSE file in the root directory of this source tree.6 7"""8Container provider abstractions for running environment servers.9 10This module provides a pluggable architecture for different container providers11(local Docker, Kubernetes, cloud providers, etc.) to be used with EnvClient.12"""13 14from __future__ import annotations15 16from abc import ABC, abstractmethod17from typing import Any, Dict, Optional, Sequence18 19 20class ContainerProvider(ABC):21    """22    Abstract base class for container providers.23 24    Providers implement this interface to support different container platforms:25    - LocalDockerProvider: Runs containers on local Docker daemon26    - KubernetesProvider: Runs containers in Kubernetes cluster27    - FargateProvider: Runs containers on AWS Fargate28    - CloudRunProvider: Runs containers on Google Cloud Run29 30    The provider manages a single container lifecycle and provides the base URL31    for connecting to it.32 33    Example:34        >>> 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    @abstractmethod42    def start_container(43        self,44        image: str,45        port: Optional[int] = None,46        env_vars: Optional[Dict[str, str]] = None,47        **kwargs: Any,48    ) -> str:49        """50        Start a container from the specified image.51 52        Args:53            image: Container image name (e.g., "echo-env:latest")54            port: Port to expose (if None, provider chooses)55            env_vars: Environment variables to pass to container56            **kwargs: Provider-specific options57 58        Returns:59            Base URL to connect to the container (e.g., "http://localhost:8000")60 61        Raises:62            RuntimeError: If container fails to start63        """64        pass65 66    @abstractmethod67    def stop_container(self) -> None:68        """69        Stop and remove the running container.70 71        This cleans up the container that was started by start_container().72        """73        pass74 75    @abstractmethod76    def wait_for_ready(self, base_url: str, timeout_s: float = 30.0) -> None:77        """78        Wait for the container to be ready to accept requests.79 80        This typically polls the /health endpoint until it returns 200.81 82        Args:83            base_url: Base URL of the container84            timeout_s: Maximum time to wait85 86        Raises:87            TimeoutError: If container doesn't become ready in time88        """89        pass90 91 92class LocalDockerProvider(ContainerProvider):93    """94    Container provider for local Docker daemon.95 96    This provider runs containers on the local machine using Docker.97    Useful for development and testing.98 99    Example:100        >>> provider = LocalDockerProvider()101        >>> base_url = provider.start_container("echo-env:latest")102        >>> # Container running on http://localhost:<random-port>103        >>> provider.stop_container()104    """105 106    def __init__(self):107        """Initialize the local Docker provider."""108        self._container_id: Optional[str] = None109        self._container_name: Optional[str] = None110 111        # Check if Docker is available112        import subprocess113 114        try:115            subprocess.run(116                ["docker", "version"],117                check=True,118                capture_output=True,119                timeout=5,120            )121        except (122            subprocess.CalledProcessError,123            FileNotFoundError,124            subprocess.TimeoutExpired,125        ):126            raise RuntimeError(127                "Docker is not available. Please install Docker Desktop or Docker Engine."128            )129 130    def start_container(131        self,132        image: str,133        port: Optional[int] = None,134        env_vars: Optional[Dict[str, str]] = None,135        **kwargs: Any,136    ) -> str:137        """138        Start a Docker container locally.139 140        Args:141            image: Docker image name142            port: Port to expose (if None, finds available port)143            env_vars: Environment variables for the container144            **kwargs: Additional Docker run options145 146        Returns:147            Base URL to connect to the container148        """149        import subprocess150        import time151 152        # Find available port if not specified153        if port is None:154            port = self._find_available_port()155 156        # Generate container name157        self._container_name = self._generate_container_name(image)158 159        # Build docker run command160        cmd = [161            "docker",162            "run",163            "-d",  # Detached164            "--name",165            self._container_name,166            "-p",167            f"{port}:8000",  # Map port168        ]169 170        # Add environment variables171        if env_vars:172            for key, value in env_vars.items():173                cmd.extend(["-e", f"{key}={value}"])174 175        # Add image176        cmd.append(image)177 178        # Run container179        try:180            result = subprocess.run(cmd, capture_output=True, text=True, check=True)181            self._container_id = result.stdout.strip()182        except subprocess.CalledProcessError as e:183            error_msg = f"Failed to start Docker container.\nCommand: {' '.join(cmd)}\nExit code: {e.returncode}\nStderr: {e.stderr}\nStdout: {e.stdout}"184            raise RuntimeError(error_msg) from e185 186        # Wait a moment for container to start187        time.sleep(1)188 189        base_url = f"http://localhost:{port}"190        return base_url191 192    def stop_container(self) -> None:193        """194        Stop and remove the Docker container.195        """196        if self._container_id is None:197            return198 199        import subprocess200 201        try:202            # Stop container203            subprocess.run(204                ["docker", "stop", self._container_id],205                capture_output=True,206                check=True,207                timeout=10,208            )209 210            # Remove container211            subprocess.run(212                ["docker", "rm", self._container_id],213                capture_output=True,214                check=True,215                timeout=10,216            )217        except subprocess.CalledProcessError:218            # Container might already be stopped/removed219            pass220        finally:221            self._container_id = None222            self._container_name = None223 224    def wait_for_ready(self, base_url: str, timeout_s: float = 30.0) -> None:225        """226        Wait for container to be ready by polling /health endpoint.227 228        Args:229            base_url: Base URL of the container230            timeout_s: Maximum time to wait231 232        Raises:233            TimeoutError: If container doesn't become ready234        """235        import time236 237        import requests238 239        start_time = time.time()240        health_url = f"{base_url}/health"241 242        # Bypass proxy for localhost to avoid proxy issues243        proxies = {"http": None, "https": None}244 245        while time.time() - start_time < timeout_s:246            try:247                response = requests.get(health_url, timeout=2.0, proxies=proxies)248                if response.status_code == 200:249                    return250            except requests.RequestException:251                pass252 253            time.sleep(0.5)254 255        raise TimeoutError(256            f"Container at {base_url} did not become ready within {timeout_s}s"257        )258 259    def _find_available_port(self) -> int:260        """261        Find an available port on localhost.262 263        Returns:264            An available port number265        """266        import socket267 268        with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:269            s.bind(("", 0))270            s.listen(1)271            port = s.getsockname()[1]272        return port273 274    def _generate_container_name(self, image: str) -> str:275        """276        Generate a unique container name based on image name and timestamp.277 278        Args:279            image: Docker image name280 281        Returns:282            A unique container name283        """284        import time285 286        clean_image = image.split("/")[-1].split(":")[0]287        timestamp = int(time.time() * 1000)288        return f"{clean_image}-{timestamp}"289 290 291class DockerSwarmProvider(ContainerProvider):292    """293    Container provider that uses Docker Swarm services for local concurrency.294 295    This provider creates a replicated Swarm service backed by the local Docker296    engine. The built-in load-balancer fans requests across the replicas,297    allowing multiple container instances to run concurrently on the developer298    workstation (mirroring the workflow described in the Docker stack docs).299    """300 301    def __init__(302        self,303        *,304        auto_init_swarm: bool = True,305        overlay_network: Optional[str] = None,306    ):307        """308        Args:309            auto_init_swarm: Whether to call ``docker swarm init`` when Swarm310                is not active. Otherwise, user must manually initialize Swarm.311            overlay_network: Optional overlay network name for the service.312                When provided, the network is created with313                ``docker network create --driver overlay --attachable`` if it314                does not already exist.315        """316        self._service_name: Optional[str] = None317        self._service_id: Optional[str] = None318        self._published_port: Optional[int] = None319        self._overlay_network = overlay_network320        self._auto_init_swarm = auto_init_swarm321 322        self._ensure_docker_available()323        self._ensure_swarm_initialized()324        if self._overlay_network:325            self._ensure_overlay_network(self._overlay_network)326 327    def start_container(328        self,329        image: str,330        port: Optional[int] = None,331        env_vars: Optional[Dict[str, str]] = None,332        **kwargs: Any,333    ) -> str:334        """335        Start (or scale) a Swarm service for the given image.336 337        Supported kwargs:338            replicas (int): Number of container replicas (default: 2).339            cpu_limit (float | str): CPU limit passed to ``--limit-cpu``.340            memory_limit (str): Memory limit passed to ``--limit-memory``.341            constraints (Sequence[str]): Placement constraints.342            labels (Dict[str, str]): Service labels.343            command (Sequence[str] | str): Override container command.344        """345        import shlex346        import subprocess347        import time348 349        allowed_kwargs = {350            "replicas",351            "cpu_limit",352            "memory_limit",353            "constraints",354            "labels",355            "command",356        }357        unknown = set(kwargs) - allowed_kwargs358        if unknown:359            raise ValueError(f"Unsupported kwargs for DockerSwarmProvider: {unknown}")360 361        replicas = int(kwargs.get("replicas", 2))362        cpu_limit = kwargs.get("cpu_limit")363        memory_limit = kwargs.get("memory_limit")364        constraints: Optional[Sequence[str]] = kwargs.get("constraints")365        labels: Optional[Dict[str, str]] = kwargs.get("labels")366        command_override = kwargs.get("command")367 368        if port is None:369            port = self._find_available_port()370 371        self._service_name = self._generate_service_name(image)372        self._published_port = port373 374        cmd = [375            "docker",376            "service",377            "create",378            "--detach",379            "--name",380            self._service_name,381            "--replicas",382            str(max(1, replicas)),383            "--publish",384            f"{port}:8000",385        ]386 387        if self._overlay_network:388            cmd.extend(["--network", self._overlay_network])389 390        if env_vars:391            for key, value in env_vars.items():392                cmd.extend(["--env", f"{key}={value}"])393 394        if cpu_limit is not None:395            cmd.extend(["--limit-cpu", str(cpu_limit)])396 397        if memory_limit is not None:398            cmd.extend(["--limit-memory", str(memory_limit)])399 400        if constraints:401            for constraint in constraints:402                cmd.extend(["--constraint", constraint])403 404        if labels:405            for key, value in labels.items():406                cmd.extend(["--label", f"{key}={value}"])407 408        cmd.append(image)409 410        if command_override:411            if isinstance(command_override, str):412                cmd.extend(shlex.split(command_override))413            else:414                cmd.extend(command_override)415 416        try:417            result = subprocess.run(418                cmd,419                capture_output=True,420                text=True,421                check=True,422            )423            self._service_id = result.stdout.strip()424        except subprocess.CalledProcessError as e:425            error_msg = (426                "Failed to start Docker Swarm service.\n"427                f"Command: {' '.join(cmd)}\n"428                f"Exit code: {e.returncode}\n"429                f"Stdout: {e.stdout}\n"430                f"Stderr: {e.stderr}"431            )432            raise RuntimeError(error_msg) from e433 434        # Give Swarm a brief moment to schedule the tasks.435        time.sleep(1.0)436 437        return f"http://localhost:{port}"438 439    def stop_container(self) -> None:440        """441        Remove the Swarm service (and keep the Swarm manager running).442        """443        if not self._service_name:444            return445 446        import subprocess447 448        try:449            subprocess.run(450                ["docker", "service", "rm", self._service_name],451                capture_output=True,452                check=True,453                timeout=10,454            )455        except subprocess.CalledProcessError:456            # Service may already be gone; ignore.457            pass458        finally:459            self._service_name = None460            self._service_id = None461            self._published_port = None462 463    def wait_for_ready(self, base_url: str, timeout_s: float = 30.0) -> None:464        """465        Wait for at least one replica to become healthy by polling /health.466 467        Note: With Swarm's load balancer, requests round-robin across replicas,468        so this only verifies that at least one replica is responding. Some469        replicas may still be starting when this returns.470        """471        import time472 473        import requests474 475        deadline = time.time() + timeout_s476        health_url = f"{base_url}/health"477 478        # Bypass proxy for localhost to avoid proxy issues479        proxies = {"http": None, "https": None}480 481        while time.time() < deadline:482            try:483                response = requests.get(health_url, timeout=2.0, proxies=proxies)484                if response.status_code == 200:485                    return486            except requests.RequestException:487                pass488 489            time.sleep(0.5)490 491        raise TimeoutError(492            f"Swarm service at {base_url} did not become ready within {timeout_s}s"493        )494 495    def _ensure_docker_available(self) -> None:496        import subprocess497 498        try:499            subprocess.run(500                ["docker", "version"],501                check=True,502                capture_output=True,503                timeout=5,504            )505        except (506            subprocess.CalledProcessError,507            FileNotFoundError,508            subprocess.TimeoutExpired,509        ) as exc:510            raise RuntimeError(511                "Docker is not available. Please install Docker Desktop or Docker Engine."512            ) from exc513 514    def _ensure_swarm_initialized(self) -> None:515        import subprocess516 517        try:518            result = subprocess.run(519                ["docker", "info", "--format", "{{.Swarm.LocalNodeState}}"],520                capture_output=True,521                text=True,522                check=True,523                timeout=5,524            )525            state = result.stdout.strip().lower()526            if state == "active":527                return528        except subprocess.CalledProcessError:529            state = "unknown"530 531        if not self._auto_init_swarm:532            raise RuntimeError(533                f"Docker Swarm is not active (state={state}). Enable Swarm manually or pass auto_init_swarm=True."534            )535 536        try:537            subprocess.run(538                ["docker", "swarm", "init"],539                check=True,540                capture_output=True,541                timeout=10,542            )543        except subprocess.CalledProcessError as e:544            raise RuntimeError("Failed to initialize Docker Swarm") from e545 546    def _ensure_overlay_network(self, network: str) -> None:547        import subprocess548 549        inspect = subprocess.run(550            ["docker", "network", "inspect", network],551            capture_output=True,552            text=True,553            check=False,554        )555        if inspect.returncode == 0:556            return557 558        try:559            subprocess.run(560                [561                    "docker",562                    "network",563                    "create",564                    "--driver",565                    "overlay",566                    "--attachable",567                    network,568                ],569                check=True,570                capture_output=True,571                timeout=10,572            )573        except subprocess.CalledProcessError as e:574            raise RuntimeError(f"Failed to create overlay network '{network}'") from e575 576    def _find_available_port(self) -> int:577        import socket578 579        with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:580            s.bind(("", 0))581            s.listen(1)582            port = s.getsockname()[1]583        return port584 585    def _generate_service_name(self, image: str) -> str:586        import time587 588        clean_image = image.split("/")[-1].split(":")[0]589        timestamp = int(time.time() * 1000)590        return f"{clean_image}-swarm-{timestamp}"591 592 593class KubernetesProvider(ContainerProvider):594    """595    Container provider for Kubernetes clusters.596 597    This provider creates pods in a Kubernetes cluster and exposes them598    via services or port-forwarding.599 600    Example:601        >>> provider = KubernetesProvider(namespace="envtorch-dev")602        >>> base_url = provider.start_container("echo-env:latest")603        >>> # Pod running in k8s, accessible via service or port-forward604        >>> provider.stop_container()605    """606 607    pass608 609 610class RuntimeProvider(ABC):611    """612    Abstract base class for runtime providers that are not container providers.613    Providers implement this interface to support different runtime platforms:614    - UVProvider: Runs environments via `uv run`615 616    The provider manages a single runtime lifecycle and provides the base URL617    for connecting to it.618 619    Example:620        >>> provider = UVProvider(project_path="/path/to/env")621        >>> base_url = provider.start()622        >>> print(base_url)  # http://localhost:8000623        >>> provider.stop()624    """625 626    @abstractmethod627    def start(628        self,629        port: Optional[int] = None,630        env_vars: Optional[Dict[str, str]] = None,631        **kwargs: Any,632    ) -> str:633        """634        Start a runtime from the specified image.635 636        Args:637            image: Runtime image name638            port: Port to expose (if None, provider chooses)639            env_vars: Environment variables for the runtime640            **kwargs: Additional runtime options641        """642 643    @abstractmethod644    def stop(self) -> None:645        """646        Stop the runtime.647        """648        pass649 650    @abstractmethod651    def wait_for_ready(self, timeout_s: float = 30.0) -> None:652        """653        Wait for the runtime to be ready to accept requests.654        """655        pass656 657    def __enter__(self) -> "RuntimeProvider":658        """659        Enter the runtime provider.660        """661        self.start()662        return self663 664    def __exit__(self, exc_type, exc, tb) -> None:665        """666        Exit the runtime provider.667        """668        self.stop()669        return False670