openenv/echo_env
6
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 