From b4c83b43c84181950af472f7b2193dd5c7e6b16b Mon Sep 17 00:00:00 2001 From: pandyamarut Date: Thu, 10 Jul 2025 17:57:05 -0700 Subject: [PATCH 1/4] feat: Lock the network volume region Signed-off-by: pandyamarut --- src/tetra_rp/client.py | 14 +--- src/tetra_rp/core/resources/network_volume.py | 8 +- src/tetra_rp/core/resources/serverless.py | 84 ++++++++++++++++++- 3 files changed, 87 insertions(+), 19 deletions(-) diff --git a/src/tetra_rp/client.py b/src/tetra_rp/client.py index 12c5d14c..a61c4636 100644 --- a/src/tetra_rp/client.py +++ b/src/tetra_rp/client.py @@ -1,7 +1,7 @@ import logging from functools import wraps -from typing import List, Optional -from .core.resources import ServerlessResource, ResourceManager, NetworkVolume +from typing import List +from .core.resources import ServerlessResource, ResourceManager from .stubs import stub_resource @@ -12,7 +12,6 @@ def remote( resource_config: ServerlessResource, dependencies: List[str] = None, system_dependencies: List[str] = None, - mount_volume: Optional[NetworkVolume] = None, **extra, ): """ @@ -49,15 +48,6 @@ async def my_function(data): def decorator(func): @wraps(func) async def wrapper(*args, **kwargs): - # Create netowrk volume if mount_volume is provided - if mount_volume: - try: - network_volume = await mount_volume.deploy() - resource_config.networkVolumeId = network_volume.id - except Exception as e: - log.error(f"Failed to create or mount network volume: {e}") - raise - resource_manager = ResourceManager() remote_resource = await resource_manager.get_or_deploy_resource( resource_config diff --git a/src/tetra_rp/core/resources/network_volume.py b/src/tetra_rp/core/resources/network_volume.py index c4f6e095..fd863d8b 100644 --- a/src/tetra_rp/core/resources/network_volume.py +++ b/src/tetra_rp/core/resources/network_volume.py @@ -20,8 +20,6 @@ class DataCenter(str, Enum): """ EU_RO_1 = "EU-RO-1" - US_WA_1 = "US-WA-1" - US_CA_1 = "US-CA-1" class NetworkVolume(DeployableResource): @@ -33,10 +31,12 @@ class NetworkVolume(DeployableResource): """ - dataCenterId: Optional[DataCenter] = None + # Internal fixed value + dataCenterId: DataCenter = Field(default=DataCenter.EU_RO_1, frozen=True) + id: Optional[str] = Field(default=None) name: Optional[str] = None - size: Optional[int] = None # Size in GB + size: Optional[int] = Field(default=10, gt=0) # Size in GB @property def is_created(self) -> bool: diff --git a/src/tetra_rp/core/resources/serverless.py b/src/tetra_rp/core/resources/serverless.py index 51267fab..f97bd7ed 100644 --- a/src/tetra_rp/core/resources/serverless.py +++ b/src/tetra_rp/core/resources/serverless.py @@ -1,6 +1,6 @@ import asyncio import logging -from typing import Any, Dict, List, Optional +from typing import Any, Dict, List, Optional, Union from enum import Enum from pydantic import ( field_serializer, @@ -22,6 +22,7 @@ from .cpu import CpuInstanceType from .environment import EnvironmentVars from .constants import CONSOLE_URL +from .network_volume import NetworkVolume # Environment variables are loaded from the .env file @@ -62,7 +63,15 @@ class ServerlessResource(DeployableResource): Base class for GPU serverless resource """ - _input_only = {"id", "cudaVersions", "env", "gpus", "flashboot", "imageName"} + _input_only = { + "id", + "cudaVersions", + "env", + "gpus", + "flashboot", + "imageName", + "networkVolume", + } # === Input-only Fields === cudaVersions: Optional[List[CudaVersion]] = [] # for allowedCudaVersions @@ -71,6 +80,11 @@ class ServerlessResource(DeployableResource): gpus: Optional[List[GpuGroup]] = [GpuGroup.ANY] # for gpuIds imageName: Optional[str] = "" # for template.imageName + # Input-only field that accepts NetworkVolume object or string ID + networkVolume: Optional[Union[NetworkVolume, str]] = Field( + default=None, exclude=True + ) + # === Input Fields === executionTimeoutMs: Optional[int] = None gpuCount: Optional[int] = 1 @@ -78,7 +92,7 @@ class ServerlessResource(DeployableResource): instanceIds: Optional[List[CpuInstanceType]] = None locations: Optional[str] = None name: str - networkVolumeId: Optional[str] = None + networkVolumeId: Optional[str] = None # This gets set from networkVolume scalerType: Optional[ServerlessScalerType] = ServerlessScalerType.QUEUE_DELAY scalerValue: Optional[int] = 4 templateId: Optional[str] = None @@ -116,6 +130,26 @@ def endpoint(self) -> runpod.Endpoint: raise ValueError("Missing self.id") return runpod.Endpoint(self.id) + @field_validator("networkVolume") + @classmethod + def validate_network_volume( + cls, value: Optional[Union[NetworkVolume, str]] + ) -> Optional[Union[NetworkVolume, str]]: + """Validate networkVolume input""" + if value is None: + return None + + if isinstance(value, str): + # If it's a string, assume it's a volume ID + return value + elif isinstance(value, NetworkVolume): + # If it's a NetworkVolume object, validate it + return value + else: + raise ValueError( + "networkVolume must be either a NetworkVolume object or a string ID" + ) + @field_serializer("scalerType") def serialize_scaler_type( self, value: Optional[ServerlessScalerType] @@ -142,6 +176,16 @@ def sync_input_fields(self): if self.flashboot: self.name += "-fb" + if self.networkVolume: + if isinstance(self.networkVolume, str): + # It's already an ID + self.networkVolumeId = self.networkVolume + elif isinstance(self.networkVolume, NetworkVolume): + # It's a NetworkVolume object + if self.networkVolume.is_created: + # Volume already exists, use its ID + self.networkVolumeId = self.networkVolume.id + if self.instanceIds: return self._sync_input_fields_cpu() else: @@ -177,6 +221,37 @@ def _sync_input_fields_cpu(self): return self + async def _ensure_network_volume_deployed(self) -> None: + """ + Ensures network volume is deployed and ready. + Updates networkVolumeId with the deployed volume ID. + """ + if not self.networkVolume: + log.info( + f"No network volume provided for {self.name}, creating default network volume" + ) + default_volume = NetworkVolume( + name=f"{self.name}-volume", + ) + self.networkVolume = default_volume + + if isinstance(self.networkVolume, str): + # It's already an ID, set it + self.networkVolumeId = self.networkVolume + return + + if isinstance(self.networkVolume, NetworkVolume): + if not self.networkVolume.is_created: + # Deploy the network volume + log.info(f"Deploying network volume for {self.name}") + deployed_volume = await self.networkVolume.deploy() + self.networkVolume = deployed_volume + self.networkVolumeId = deployed_volume.id + log.info(f"Network volume deployed with ID: {deployed_volume.id}") + else: + # Already deployed, just set the ID + self.networkVolumeId = self.networkVolume.id + def is_deployed(self) -> bool: """ Checks if the serverless resource is deployed and available. @@ -202,6 +277,9 @@ async def deploy(self) -> "DeployableResource": log.debug(f"{self} exists") return self + # NEW: Ensure network volume is deployed first + await self._ensure_network_volume_deployed() + async with RunpodGraphQLClient() as client: payload = self.model_dump(exclude=self._input_only, exclude_none=True) result = await client.create_endpoint(payload) From 96c5a442543641f8190ae8ed549621fba9cab09d Mon Sep 17 00:00:00 2001 From: pandyamarut Date: Sun, 20 Jul 2025 23:03:11 -0700 Subject: [PATCH 2/4] chore: address comments and reformat Signed-off-by: pandyamarut --- src/tetra_rp/core/api/runpod.py | 36 ++------ src/tetra_rp/core/resources/network_volume.py | 23 +++-- src/tetra_rp/core/resources/serverless.py | 83 ++++--------------- 3 files changed, 42 insertions(+), 100 deletions(-) diff --git a/src/tetra_rp/core/api/runpod.py b/src/tetra_rp/core/api/runpod.py index c5d94347..7524e623 100644 --- a/src/tetra_rp/core/api/runpod.py +++ b/src/tetra_rp/core/api/runpod.py @@ -3,11 +3,12 @@ Bypasses the outdated runpod-python SDK limitations. """ -import os import json -import aiohttp -from typing import Dict, Any, Optional import logging +import os +from typing import Any, Dict, Optional + +import aiohttp log = logging.getLogger(__name__) @@ -267,31 +268,12 @@ async def _execute_rest( raise Exception(f"HTTP request failed: {e}") async def create_network_volume(self, payload: Dict[str, Any]) -> Dict[str, Any]: - """ - Create a network volume in Runpod. + """Create a network volume in Runpod.""" + log.debug(f"Creating network volume: {payload.get('name', 'unnamed')}") - Args: - datacenter_id (str): The ID of the datacenter where the volume will be created. - name (str): The name of the network volume. - size_gb (int): The size of the volume in GB. - - Returns: - Dict[str, Any]: The created network volume details. - """ - datacenter_id = payload.get("dataCenterId") - if hasattr(datacenter_id, "value"): - # If datacenter_id is an enum, get its value - datacenter_id = datacenter_id.value - data = { - "dataCenterId": datacenter_id, - "name": payload.get("name"), - "size": payload.get("size"), - } - url = f"{RUNPOD_REST_API_URL}/networkvolumes" - - log.debug(f"Creating network volume: {data.get('name', 'unnamed')}") - - result = await self._execute_rest("POST", url, data) + result = await self._execute_rest( + "POST", f"{RUNPOD_REST_API_URL}/networkvolumes", payload + ) log.info( f"Created network volume: {result.get('id', 'unknown')} - {result.get('name', 'unnamed')}" diff --git a/src/tetra_rp/core/resources/network_volume.py b/src/tetra_rp/core/resources/network_volume.py index fd863d8b..1b2f7285 100644 --- a/src/tetra_rp/core/resources/network_volume.py +++ b/src/tetra_rp/core/resources/network_volume.py @@ -4,6 +4,7 @@ from pydantic import ( Field, + field_serializer, ) from ..api.runpod import RunpodRestClient @@ -38,6 +39,14 @@ class NetworkVolume(DeployableResource): name: Optional[str] = None size: Optional[int] = Field(default=10, gt=0) # Size in GB + def __str__(self) -> str: + return f"{self.__class__.__name__}:{self.id}" + + @field_serializer("dataCenterId") + def serialize_data_center_id(self, value: Optional[DataCenter]) -> Optional[str]: + """Convert DataCenter enum to string.""" + return value.value if value is not None else None + @property def is_created(self) -> bool: "Returns True if the network volume already exists." @@ -79,17 +88,17 @@ async def deploy(self) -> "DeployableResource": try: # If the resource is already deployed, return it if self.is_deployed(): - log.debug( - f"Network volume {self.id} is already deployed. Mounting existing volume." - ) - log.info(f"Mounted existing network volume: {self.id}") + log.debug(f"{self} exists") return self # Create the network volume - self = await self.create_network_volume() + async with RunpodRestClient() as client: + # Create the network volume + payload = self.model_dump(exclude_none=True) + result = await client.create_network_volume(payload) - if self.is_deployed(): - return self + if volume := self.__class__(**result): + return volume raise ValueError("Deployment failed, no volume was created.") diff --git a/src/tetra_rp/core/resources/serverless.py b/src/tetra_rp/core/resources/serverless.py index f97bd7ed..05beee2f 100644 --- a/src/tetra_rp/core/resources/serverless.py +++ b/src/tetra_rp/core/resources/serverless.py @@ -1,28 +1,27 @@ import asyncio import logging -from typing import Any, Dict, List, Optional, Union from enum import Enum +from typing import Any, Dict, List, Optional + from pydantic import ( + BaseModel, + Field, field_serializer, field_validator, model_validator, - BaseModel, - Field, ) - from runpod.endpoint.runner import Job from ..api.runpod import RunpodGraphQLClient from ..utils.backoff import get_backoff_delay - -from .cloud import runpod from .base import DeployableResource -from .template import PodTemplate, KeyValuePair -from .gpu import GpuGroup +from .cloud import runpod +from .constants import CONSOLE_URL from .cpu import CpuInstanceType from .environment import EnvironmentVars -from .constants import CONSOLE_URL +from .gpu import GpuGroup from .network_volume import NetworkVolume +from .template import KeyValuePair, PodTemplate # Environment variables are loaded from the .env file @@ -80,10 +79,7 @@ class ServerlessResource(DeployableResource): gpus: Optional[List[GpuGroup]] = [GpuGroup.ANY] # for gpuIds imageName: Optional[str] = "" # for template.imageName - # Input-only field that accepts NetworkVolume object or string ID - networkVolume: Optional[Union[NetworkVolume, str]] = Field( - default=None, exclude=True - ) + networkVolume: Optional[NetworkVolume] = None # sets the networkVolumeId # === Input Fields === executionTimeoutMs: Optional[int] = None @@ -130,26 +126,6 @@ def endpoint(self) -> runpod.Endpoint: raise ValueError("Missing self.id") return runpod.Endpoint(self.id) - @field_validator("networkVolume") - @classmethod - def validate_network_volume( - cls, value: Optional[Union[NetworkVolume, str]] - ) -> Optional[Union[NetworkVolume, str]]: - """Validate networkVolume input""" - if value is None: - return None - - if isinstance(value, str): - # If it's a string, assume it's a volume ID - return value - elif isinstance(value, NetworkVolume): - # If it's a NetworkVolume object, validate it - return value - else: - raise ValueError( - "networkVolume must be either a NetworkVolume object or a string ID" - ) - @field_serializer("scalerType") def serialize_scaler_type( self, value: Optional[ServerlessScalerType] @@ -176,15 +152,9 @@ def sync_input_fields(self): if self.flashboot: self.name += "-fb" - if self.networkVolume: - if isinstance(self.networkVolume, str): - # It's already an ID - self.networkVolumeId = self.networkVolume - elif isinstance(self.networkVolume, NetworkVolume): - # It's a NetworkVolume object - if self.networkVolume.is_created: - # Volume already exists, use its ID - self.networkVolumeId = self.networkVolume.id + if self.networkVolume and self.networkVolume.is_created: + # Volume already exists, use its ID + self.networkVolumeId = self.networkVolume.id if self.instanceIds: return self._sync_input_fields_cpu() @@ -227,30 +197,11 @@ async def _ensure_network_volume_deployed(self) -> None: Updates networkVolumeId with the deployed volume ID. """ if not self.networkVolume: - log.info( - f"No network volume provided for {self.name}, creating default network volume" - ) - default_volume = NetworkVolume( - name=f"{self.name}-volume", - ) - self.networkVolume = default_volume - - if isinstance(self.networkVolume, str): - # It's already an ID, set it - self.networkVolumeId = self.networkVolume - return - - if isinstance(self.networkVolume, NetworkVolume): - if not self.networkVolume.is_created: - # Deploy the network volume - log.info(f"Deploying network volume for {self.name}") - deployed_volume = await self.networkVolume.deploy() - self.networkVolume = deployed_volume - self.networkVolumeId = deployed_volume.id - log.info(f"Network volume deployed with ID: {deployed_volume.id}") - else: - # Already deployed, just set the ID - self.networkVolumeId = self.networkVolume.id + log.info(f"{self.name} requires a default network volume") + self.networkVolume = NetworkVolume(name=f"{self.name}-volume") + + if deployedNetworkVolume := await self.networkVolume.deploy(): + self.networkVolumeId = deployedNetworkVolume.id def is_deployed(self) -> bool: """ From 02addd0232fb02b91af139b96edb960597255de8 Mon Sep 17 00:00:00 2001 From: pandyamarut Date: Sun, 20 Jul 2025 23:06:09 -0700 Subject: [PATCH 3/4] chore: remove comment Signed-off-by: pandyamarut --- src/tetra_rp/core/resources/serverless.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/tetra_rp/core/resources/serverless.py b/src/tetra_rp/core/resources/serverless.py index 05beee2f..80c563ae 100644 --- a/src/tetra_rp/core/resources/serverless.py +++ b/src/tetra_rp/core/resources/serverless.py @@ -79,7 +79,7 @@ class ServerlessResource(DeployableResource): gpus: Optional[List[GpuGroup]] = [GpuGroup.ANY] # for gpuIds imageName: Optional[str] = "" # for template.imageName - networkVolume: Optional[NetworkVolume] = None # sets the networkVolumeId + networkVolume: Optional[NetworkVolume] = None # === Input Fields === executionTimeoutMs: Optional[int] = None @@ -88,7 +88,7 @@ class ServerlessResource(DeployableResource): instanceIds: Optional[List[CpuInstanceType]] = None locations: Optional[str] = None name: str - networkVolumeId: Optional[str] = None # This gets set from networkVolume + networkVolumeId: Optional[str] = None scalerType: Optional[ServerlessScalerType] = ServerlessScalerType.QUEUE_DELAY scalerValue: Optional[int] = 4 templateId: Optional[str] = None From 6e918a30344ff494d05e008434c59605933c1d7e Mon Sep 17 00:00:00 2001 From: pandyamarut Date: Mon, 21 Jul 2025 12:02:42 -0700 Subject: [PATCH 4/4] add check Signed-off-by: pandyamarut --- src/tetra_rp/core/resources/serverless.py | 3 +++ 1 file changed, 3 insertions(+) diff --git a/src/tetra_rp/core/resources/serverless.py b/src/tetra_rp/core/resources/serverless.py index 80c563ae..75c28684 100644 --- a/src/tetra_rp/core/resources/serverless.py +++ b/src/tetra_rp/core/resources/serverless.py @@ -196,6 +196,9 @@ async def _ensure_network_volume_deployed(self) -> None: Ensures network volume is deployed and ready. Updates networkVolumeId with the deployed volume ID. """ + if self.networkVolumeId: + return + if not self.networkVolume: log.info(f"{self.name} requires a default network volume") self.networkVolume = NetworkVolume(name=f"{self.name}-volume")