openenv/coding_env
21
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 