diff --git a/CHANGELOG.md b/CHANGELOG.md index b33348340..be9484adf 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,14 @@ * Apps expose `NetworkVolume` and `GlobalVolume` resources with explicit path-keyed mounts, lazy create-if-missing provisioning, and execution-bound filesystem access. Global volumes use GraphQL for lookup and creation. +### Bug Fixes + +* Task calls support indefinite execution with bounded interruption cleanup. Failed deletion retains the pod ID for recovery. +* Network volume creation intersects catalog storage support with every sharing resource's hardware stock and datacenter constraints. Unsupported explicit pins fail before creation; existing volumes retain their authoritative IDs and datacenters. +* Global volume attachments require GPU compute for tasks and endpoints. +* Deployment artifacts apply project ignore rules and mandatory secret-file exclusions while preserving vendored dependency assets. +* Synchronous calls wait across polling timeouts on Python 3.10 and propagate operation timeouts. + ## [1.12.0](https://github.com/runpod/runpod-python/compare/v1.11.0...v1.12.0) (2026-08-10) diff --git a/docs/cli/references/projects.md b/docs/cli/references/projects.md index 2480b204e..e1d0fef75 100644 --- a/docs/cli/references/projects.md +++ b/docs/cli/references/projects.md @@ -11,3 +11,27 @@ You may need to update the default configuration within `runpod.toml` to match y ## Ignore Files and Folders Create a `.runpodignore` file in the root of your project to ignore files and folders from being uploaded to the Runpod platform, the same file will also be used to ignore files that should not trigger an API server reload. + +### Flash deployment artifacts + +`rp flash deploy` and `rp flash deploy --build-only` apply Git-style ignore +patterns in this order, with later rules taking precedence: + +1. Default exclusions for local environments, caches, tests, and build archives. +2. Project `.gitignore` files, with deeper files overriding parents in their subtree. +3. The project-root `.runpodignore`. + +Ancestor and global Git ignores are not read. Negation (`!`) can re-include +ordinary files, but excluded parent directories must also be re-included. + +Negation cannot include `.git`, `.runpod`, `.flash`, the root `env/` or +`runpod_manifest.json`, or credential-like source paths such as `.env` variants, +PEM/key files, private SSH keys, cloud credential directories, and credentials, +secrets, or service-account files. These are filename safeguards, not secret +scanning; review the build-only artifact and supply credentials through worker +environment variables or a secret store. + +The manifest and vendored `env/` are added separately. Source ignore rules do +not strip dependency CA bundles. Symlinks are omitted from source and +dependencies; the output artifact and dependency build directory are never +copied back into source. diff --git a/examples/apps/README.md b/examples/apps/README.md index 4fd1fa8ab..c7ca9abb5 100644 --- a/examples/apps/README.md +++ b/examples/apps/README.md @@ -20,6 +20,40 @@ workers pick up the new code automatically. the exhaustive feature-by-feature suite lives in [`tests/e2e/examples`](../../tests/e2e/examples). `await task.spawn.aio(...)` returns a task job that owns its pod. -`await job.wait(timeout=...)` terminates that pod when waiting finishes, including -timeout, task failure, or cancellation. use `await job.cancel()` to abandon a -spawned task explicitly. a failed termination retains `job.pod_id` for cleanup. +`await job.wait(timeout=...)` attempts to terminate that pod when waiting finishes, +including timeout, task failure, or cancellation. use `await job.cancel()` to +abandon a spawned task explicitly. leaving the client without waiting or cancelling +does not cancel intentional detached work. + +cancellation waits for bounded cleanup and retries transient deletion failures. +failed deletion retains `job.pod_id`; check the console after a cleanup warning. +active tasks can run indefinitely, including detached work after client exit. +`remote()` and `job.wait()` impose no execution timeout by default. pod deletion +requires a working control-plane API and pod-scoped credentials. + +## Storage + +`NetworkVolume(name_or_id, size=50, datacenter=None, create=True)` references +datacenter-local storage. New volumes use one catalog capability snapshot per +provisioning run to select a datacenter with network storage support and hardware +stock for every resource sharing the volume, respecting their datacenter pins. +An unsupported explicit volume pin fails before creation instead of relocating; +catalog lookup failures also stop creation. Existing volumes retain their IDs and +datacenters regardless of new-volume eligibility. Their consumers must still be +schedulable in that datacenter. `GlobalVolume(name_or_id, create=True)` +references global storage without a datacenter constraint. Both inherit from +the abstract `Volume` base, resolve by name or ID, and create missing storage +when provisioning remote compute. Set `create=False` to require existing storage. +Global-volume lookup and creation use GraphQL; network volumes use REST. + +Declare attachments with `mounts={"/path": volume}`. Tasks support one network +and one global volume at distinct, non-overlapping paths. Queue and API resources +support one volume at `/runpod-volume`. Global volumes require GPU compute for +both tasks and endpoints. + +Inside worker code, `volume.path` returns the configured mount path. The runtime +binds declared references and resolved IDs before importing user code. Access +raises if the volume is unmounted or has multiple bindings. Evaluate `.path` +inside remote functions, not during local module discovery. Calling `.local()` +does not mount remote storage on the client machine. + diff --git a/requirements.txt b/requirements.txt index a19a076db..21ad7c4b4 100644 --- a/requirements.txt +++ b/requirements.txt @@ -9,6 +9,7 @@ colorama >= 0.4.6, < 0.4.7 cryptography >= 50.0.1 fastapi[all] >= 0.141.1 paramiko >= 5.0.0 +pathspec >= 0.12.1 prettytable >= 3.18.0 psutil >= 7.2.2 py-cpuinfo >= 9.0.0 diff --git a/runpod/apps/api.py b/runpod/apps/api.py index 64ba26b3a..30b6c301d 100644 --- a/runpod/apps/api.py +++ b/runpod/apps/api.py @@ -5,7 +5,7 @@ absent from rest. management verbs for the wider sdk stay in runpod.api.ctl_commands. """ -from typing import Any, Dict, List, Optional +from typing import Any, Dict, List, Optional, Set import aiohttp @@ -387,6 +387,31 @@ async def cpu_stock_status( ) return _stock_in_datacenter(data, data_center_id) + async def network_volume_datacenters(self) -> Set[str]: + """datacenters supporting any network volume tier chosen by the backend.""" + data = await run_rest_request_async( + "GET", "/v2/catalog/datacenters", api_key=self._api_key + ) + datacenters = data.get("dataCenters") if isinstance(data, dict) else None + if not isinstance(datacenters, list): + raise QueryError("datacenter catalog is missing dataCenters") + supported = set() + for dc in datacenters: + if ( + not isinstance(dc, dict) + or not isinstance(dc.get("id"), str) + or not dc["id"] + or not isinstance(dc.get("networkVolumeTypes"), list) + or any( + not isinstance(tier, str) or not tier + for tier in dc["networkVolumeTypes"] + ) + ): + raise QueryError("datacenter catalog has invalid networkVolumeTypes") + if dc["networkVolumeTypes"]: + supported.add(dc["id"]) + return supported + async def list_global_volumes(self) -> List[Dict[str, Any]]: data = await self._execute(app_queries.QUERY_GLOBAL_VOLUMES, retry=True) return data["myself"]["globalStoreBuckets"] diff --git a/runpod/apps/context.py b/runpod/apps/context.py index 3bcffc5a4..b21d38760 100644 --- a/runpod/apps/context.py +++ b/runpod/apps/context.py @@ -1,11 +1,18 @@ """execution context detection and the sync/async bridge.""" import asyncio +import concurrent.futures +import logging import os import threading +import time from enum import Enum from typing import Any, Coroutine +log = logging.getLogger(__name__) +CLEANUP_TIMEOUT = 30.0 +BRIDGE_CLEANUP_TIMEOUT = CLEANUP_TIMEOUT + 5.0 + class Context(Enum): """where the current process is running.""" @@ -68,16 +75,45 @@ def _ensure_loop(cls) -> asyncio.AbstractEventLoop: @classmethod def run(cls, coro: Coroutine[Any, Any, Any]) -> Any: loop = cls._ensure_loop() - future = asyncio.run_coroutine_threadsafe(coro, loop) + completed = threading.Event() + interrupted = threading.Event() + running = [] + + async def run(): + running.append(asyncio.current_task()) + try: + if interrupted.is_set(): + coro.close() + raise asyncio.CancelledError() + return await coro + finally: + completed.set() + + def cancel(): + interrupted.set() + if running: + running[0].cancel() + + future = asyncio.run_coroutine_threadsafe(run(), loop) try: while True: try: return future.result(timeout=0.2) - except TimeoutError: + except concurrent.futures.TimeoutError: if future.done(): raise except BaseException: - future.cancel() + loop.call_soon_threadsafe(cancel) + deadline = time.monotonic() + BRIDGE_CLEANUP_TIMEOUT + while not completed.is_set(): + try: + remaining = deadline - time.monotonic() + if remaining <= 0: + log.warning("interrupted operation cleanup is still pending") + break + completed.wait(min(remaining, 0.2)) + except BaseException: + continue raise diff --git a/runpod/apps/datacenter.py b/runpod/apps/datacenter.py index 04b159b1a..1ba61f469 100644 --- a/runpod/apps/datacenter.py +++ b/runpod/apps/datacenter.py @@ -1,6 +1,6 @@ """datacenter selection for app resources. -only datacenters with storage support and S3 API support are listed""" +network volume creation support is discovered from the datacenter catalog.""" from enum import Enum from typing import List @@ -41,8 +41,8 @@ def all(cls) -> List["DataCenter"]: return list(cls) -# datacenters with high cpu serverless stock, restricted to the -# storage+S3 set above. cpu5c/cpu5g are only stocked in EU-RO-1. +# datacenters with high cpu serverless stock. storage support is checked +# separately; cpu5c/cpu5g are only stocked in EU-RO-1. CPU3_DATACENTERS: List[DataCenter] = [ DataCenter.EU_CZ_1, DataCenter.EU_RO_1, diff --git a/runpod/apps/deploy.py b/runpod/apps/deploy.py index a864099a8..a0ef47ca7 100644 --- a/runpod/apps/deploy.py +++ b/runpod/apps/deploy.py @@ -1,6 +1,5 @@ """deploy pipeline: discovered apps -> manifest -> artifact -> activated build.""" -import fnmatch import json import logging import os @@ -8,7 +7,9 @@ import tempfile from dataclasses import dataclass, field from pathlib import Path -from typing import Any, Dict, List, Optional +from typing import Any, Dict, List, Optional, Tuple + +from pathspec import GitIgnoreSpec import runpod @@ -32,18 +33,62 @@ MANIFEST_VERSION = 1 DEFAULT_IGNORES = [ - ".git", - ".venv", - "venv", - "__pycache__", + ".venv/", + "venv/", + "env/", + "__pycache__/", "*.pyc", - ".runpod", - ".flash", - "node_modules", + "node_modules/", ".DS_Store", "*.tar.gz", + "tests/", + "test/", + "test_*.py", + "*_test.py", + ".pytest_cache/", + ".mypy_cache/", + ".ruff_cache/", + ".tox/", + ".nox/", ] +PROTECTED_IGNORES = GitIgnoreSpec.from_lines( + [ + ".git", + ".runpod", + ".flash", + "/env", + "/runpod_manifest.json", + ".env", + ".env.*", + "*.env", + "*.env.*", + "*.pem", + "*.key", + "id_rsa", + "id_rsa.*", + "id_dsa", + "id_dsa.*", + "id_ecdsa", + "id_ecdsa.*", + "id_ed25519", + "id_ed25519.*", + ".ssh", + ".aws", + ".azure", + ".kube", + ".docker/config.json", + "**/.docker/config.json", + ".netrc", + ".npmrc", + ".pypirc", + ".git-credentials", + ".boto", + "service-account*.json", + "service_account*.json", + ] +) + def _module_path_for(fn, project_root: Path) -> str: """dotted import path for fn's file relative to the project root.""" @@ -117,25 +162,19 @@ def build_manifest( return manifest -def _load_ignores(project_root: Path) -> List[str]: - patterns = list(DEFAULT_IGNORES) - ignore_file = project_root / ".runpodignore" - if ignore_file.exists(): - for line in ignore_file.read_text().splitlines(): - line = line.strip() - if line and not line.startswith("#"): - patterns.append(line) - return patterns +def _load_ignores(path: Path) -> GitIgnoreSpec: + if path.is_symlink() or not path.is_file(): + return GitIgnoreSpec.from_lines([]) + return GitIgnoreSpec.from_lines(path.read_text(encoding="utf-8").splitlines()) -def _is_ignored(rel_path: str, patterns: List[str]) -> bool: - parts = rel_path.split("/") - for pattern in patterns: - if fnmatch.fnmatch(rel_path, pattern): - return True - if any(fnmatch.fnmatch(part, pattern) for part in parts): - return True - return False +def _is_ignored(rel_path: str, patterns: List[Tuple[str, GitIgnoreSpec]]) -> bool: + ignored = False + for prefix, spec in patterns: + result = spec.check_file(rel_path[len(prefix) :]) + if result.include is not None: + ignored = result.include + return ignored ENV_DIR_NAME = "env" @@ -150,35 +189,72 @@ def package_project( """tar source + vendored env + manifest into a build artifact. layout inside the tarball: - {source files} project code, .runpodignore honored + {source files} project code, git-style ignore rules honored env/ vendored site-packages tree runpod_manifest.json """ if output is None: output = Path(tempfile.mkdtemp()) / "artifact.tar.gz" - patterns = _load_ignores(project_root) + project_root = project_root.resolve() + output_resolved = output.resolve() env_resolved = env_dir.resolve() if env_dir is not None else None - - with tarfile.open(output, "w:gz") as tar: - for path in sorted(project_root.rglob("*")): - if not path.is_file(): - continue - if env_resolved is not None and env_resolved in path.resolve().parents: - continue - rel = path.relative_to(project_root).as_posix() - if _is_ignored(rel, patterns): - continue - if rel == ENV_DIR_NAME or rel.startswith(f"{ENV_DIR_NAME}/"): - continue - tar.add(path, arcname=rel) - - if env_dir is not None and env_dir.is_dir(): - for path in sorted(env_dir.rglob("*")): - if not path.is_file(): + runpod_ignores = _load_ignores(project_root / ".runpodignore") + git_ignores = [("", GitIgnoreSpec.from_lines(DEFAULT_IGNORES))] + + with tarfile.open(output, "w:gz", dereference=True) as tar: + for directory, dirs, files in os.walk(project_root, followlinks=False): + root = Path(directory) + relative = root.relative_to(project_root).as_posix() + prefix = "" if relative == "." else relative + "/" + while not prefix.startswith(git_ignores[-1][0]): + git_ignores.pop() + git_ignores.append((prefix, _load_ignores(root / ".gitignore"))) + patterns = [*git_ignores, ("", runpod_ignores)] + kept_dirs = [] + for name in sorted(dirs): + path = root / name + rel = prefix + name + "/" + if ( + path.is_symlink() + or path == env_resolved + or PROTECTED_IGNORES.match_file(rel) + or _is_ignored(rel, patterns) + ): + log.warning("excluded source path %r from deployment artifact", rel) + continue + kept_dirs.append(name) + dirs[:] = kept_dirs + for name in sorted(files): + path = root / name + rel = prefix + name + if ( + path.is_symlink() + or not path.is_file() + or path == output_resolved + or PROTECTED_IGNORES.match_file(rel) + or _is_ignored(rel, patterns) + ): + log.warning("excluded source path %r from deployment artifact", rel) continue - rel = path.relative_to(env_dir).as_posix() - tar.add(path, arcname=f"{ENV_DIR_NAME}/{rel}") + tar.add(path, arcname=rel, recursive=False) + + if env_resolved is not None and env_resolved.is_dir(): + for directory, dirs, files in os.walk(env_resolved, followlinks=False): + root = Path(directory) + dirs[:] = sorted( + name for name in dirs if not (root / name).is_symlink() + ) + for name in sorted(files): + path = root / name + if ( + path.is_symlink() + or not path.is_file() + or path == output_resolved + ): + continue + rel = path.relative_to(env_resolved).as_posix() + tar.add(path, arcname=f"{ENV_DIR_NAME}/{rel}", recursive=False) manifest_bytes = json.dumps(manifest, indent=2).encode() info = tarfile.TarInfo(name="runpod_manifest.json") diff --git a/runpod/apps/init.py b/runpod/apps/init.py index 0bb0e9cd1..d7f212017 100644 --- a/runpod/apps/init.py +++ b/runpod/apps/init.py @@ -42,12 +42,16 @@ def main(): # (also installed locally for rp flash dev) """ -RUNPODIGNORE_TEMPLATE = """# excluded from the deploy artifact -.git -.venv -__pycache__ -*.pyc -.env +RUNPODIGNORE_TEMPLATE = """# git-style patterns, applied after project .gitignore files +# excluded directories must be re-included before their files +# credential safeguards and internal paths cannot be re-included +# add project-specific exclusions here; review the artifact before deploying +.venv/ +__pycache__/ +tests/ +.env* +*.pem +*.key """ PROJECT_FILES: Dict[str, str] = { @@ -59,9 +63,7 @@ def main(): def detect_conflicts(project_dir: Path) -> List[str]: """names of skeleton files that already exist in project_dir.""" - return [ - name for name in PROJECT_FILES if (project_dir / name).exists() - ] + return [name for name in PROJECT_FILES if (project_dir / name).exists()] def create_project( diff --git a/runpod/apps/placement.py b/runpod/apps/placement.py index 931bece77..e9e99add6 100644 --- a/runpod/apps/placement.py +++ b/runpod/apps/placement.py @@ -82,9 +82,7 @@ async def fetch(self, keys: Iterable[StockKey]) -> None: gpu_ids = { (k[1], k[2]) for k in keys - if k[0] == "gpu" - and k[1] != "*" - and (k[1], k[2]) not in self._fetched_gpu + if k[0] == "gpu" and k[1] != "*" and (k[1], k[2]) not in self._fetched_gpu } gpu_pod_ids = { (k[1], k[2]) @@ -133,7 +131,9 @@ async def _fetch_cpu(self, client, instance_id: str, dc: str, pods: bool) -> Non try: status = await client.cpu_stock_status(instance_id, dc, pods=pods) except Exception: # noqa: BLE001 - stock is advisory - log.debug("cpu stock query failed for %s@%s", instance_id, dc, exc_info=True) + log.debug( + "cpu stock query failed for %s@%s", instance_id, dc, exc_info=True + ) status = None self._cpu[(instance_id, dc, pods)] = _score(status) @@ -186,22 +186,22 @@ def solve_placement( stock: StockMap, *, volume_name: str, + volume_datacenters: Set[str], + volume_dc: Optional[str] = None, existing_dc: Optional[str] = None, ) -> str: """pick the datacenter for one volume given every resource using it. an existing volume's DC is a hard constraint (verified schedulable); - a new volume lands in the intersection of every resource's candidate - set, ranked maximin: the DC where the most-constrained resource has - the best stock. + a new volume lands in the intersection of storage capability, its pin, + and every resource's candidate set, ranked maximin: the DC where the + most-constrained resource has the best stock. """ per_resource = {spec.name: candidates(spec, stock) for spec in specs} if existing_dc is not None: existing_dc = DataCenter.from_string(existing_dc).value - blocked = [ - name for name, dcs in per_resource.items() if existing_dc not in dcs - ] + blocked = [name for name, dcs in per_resource.items() if existing_dc not in dcs] if blocked: raise PlacementError( f"volume '{volume_name}' lives in {existing_dc}, but " @@ -210,7 +210,9 @@ def solve_placement( ) return existing_dc - shared = set.intersection(*per_resource.values()) if per_resource else set() + shared = volume_datacenters.intersection(*per_resource.values()) + if volume_dc is not None: + shared &= {volume_dc} if not shared: lines = [ f" {name:<12} schedulable in: {', '.join(sorted(dcs)) or '(nowhere)'}" @@ -218,7 +220,10 @@ def solve_placement( ] raise PlacementError( f"cannot place volume '{volume_name}': no datacenter can host " - f"every resource using it\n" + "\n".join(lines) + "\n" + f"every resource using it with network volume support" + f"{f' in pinned datacenter {volume_dc}' if volume_dc else ''}\n" + + "\n".join(lines) + + "\n" "use separate volumes or compatible hardware" ) diff --git a/runpod/apps/targets.py b/runpod/apps/targets.py index 40769eae0..28e720513 100644 --- a/runpod/apps/targets.py +++ b/runpod/apps/targets.py @@ -806,9 +806,6 @@ class PodTarget(InvocationTarget): runner on the pod (see runpod.apps.tasks). """ - # tasks default to a long window; the pod's terminateAfter is the backstop - TASK_TIMEOUT_SECONDS = 3600.0 - def __init__( self, spec: ResourceSpec, @@ -843,11 +840,11 @@ def build_payload( return request.to_input() async def invoke( - self, payload: Dict[str, Any], *, timeout: float = TASK_TIMEOUT_SECONDS + self, payload: Dict[str, Any], *, timeout: Optional[float] = None ) -> Any: import time - from .tasks import TaskExecution, unwrap_task_response + from .tasks import TaskExecution, finish_cleanup, unwrap_task_response name = self.spec.name hardware = ",".join(self.spec.cpu or self.spec.gpu or ["any"]) @@ -874,6 +871,7 @@ async def invoke( await execution.wait_ready() emit(self.events, "worker_ready", name, execution.pod_id or "") response = await execution.execute(payload, timeout) + result = unwrap_task_response(response) except Exception: emit( self.events, @@ -883,12 +881,15 @@ async def invoke( ) raise finally: - try: - if stream is not None: - await stream.stop() - finally: - await execution.terminate() - result = unwrap_task_response(response) + + async def cleanup(): + try: + if stream is not None: + await stream.stop() + finally: + await execution.terminate() + + await finish_cleanup(cleanup()) emit( self.events, "request_completed", @@ -898,7 +899,7 @@ async def invoke( return result async def submit(self, payload: Dict[str, Any]) -> Any: - from .tasks import TaskExecution, TaskJob + from .tasks import TaskExecution, TaskJob, finish_cleanup execution = TaskExecution(self.spec, specs=self.specs) try: @@ -906,6 +907,6 @@ async def submit(self, payload: Dict[str, Any]) -> Any: await execution.wait_ready() await execution.submit(payload) except BaseException: - await execution.terminate() + await finish_cleanup(execution.terminate()) raise return TaskJob(execution) diff --git a/runpod/apps/tasks.py b/runpod/apps/tasks.py index d1707f300..e6a6419b4 100644 --- a/runpod/apps/tasks.py +++ b/runpod/apps/tasks.py @@ -6,22 +6,21 @@ 2. wait for the runner's /ping via the pod http proxy 3. POST the FunctionRequest to /execute (remote) or /submit (spawn) 4. collect the response, terminate the pod - -`terminateAfter` is set at deploy time as a server-side safety net so a -crashed client cannot leak a running pod indefinitely. """ import asyncio import logging +import math import os as _os import secrets +import sys import time -from datetime import datetime, timedelta, timezone from typing import Any, Dict, List, Optional import aiohttp from .api import AppsApiClient, is_capacity_error +from .context import CLEANUP_TIMEOUT from .errors import RemoteExecutionError from .spec import ResourceSpec @@ -29,12 +28,55 @@ TASK_PORT = 8080 -# safety net: pods self-terminate server-side after this long -DEFAULT_MAX_LIFETIME = timedelta(hours=1) - READY_POLL_INTERVAL = 2.0 READY_TIMEOUT = 600.0 RESULT_POLL_INTERVAL = 2.0 +DELETE_ATTEMPTS = 3 +DELETE_TIMEOUT = 5.0 +_pending_cleanups = set() + + +async def finish_cleanup(coro) -> None: + """finish cleanup despite repeated cancellation, without hiding its cause.""" + original = sys.exc_info()[1] + task = asyncio.create_task(coro) + deadline = asyncio.get_running_loop().time() + CLEANUP_TIMEOUT + cancellation = None + try: + while not task.done(): + remaining = deadline - asyncio.get_running_loop().time() + if remaining <= 0: + raise TimeoutError("task pod cleanup is still pending") + try: + await asyncio.wait({task}, timeout=remaining) + except asyncio.CancelledError as exc: + if cancellation is None: + cancellation = exc + task.result() + except Exception: + if original is None and cancellation is None: + raise + log.warning( + "task cleanup failed while handling %r", + original or cancellation, + exc_info=True, + ) + finally: + if not task.done(): + _pending_cleanups.add(task) + task.add_done_callback(_cleanup_finished) + if original is None and cancellation is not None: + raise cancellation + + +def _cleanup_finished(task) -> None: + _pending_cleanups.discard(task) + if task.cancelled(): + return + try: + task.result() + except Exception: + log.warning("background task pod cleanup failed", exc_info=True) def _runtime_command() -> str: @@ -85,7 +127,6 @@ def _pod_input(spec: ResourceSpec, token: str, task_name: str) -> Dict[str, Any] start the runtime package via dockerArgs. """ spec.validate() - terminate_after = (datetime.now(timezone.utc) + DEFAULT_MAX_LIFETIME).isoformat() from .secret import render_env @@ -102,7 +143,6 @@ def _pod_input(spec: ResourceSpec, token: str, task_name: str) -> Dict[str, Any] "imageName": image_for_spec(spec, python_version=local_python_version()), "ports": f"{TASK_PORT}/http", "containerDiskInGb": spec.container_disk_gb or (10 if spec.is_cpu else 30), - "terminateAfter": terminate_after, "supportPublicIp": True, } @@ -160,6 +200,7 @@ def __init__( self.api = api or AppsApiClient() self.token = secrets.token_urlsafe(32) self.pod_id: Optional[str] = None + self._deployment: Optional[asyncio.Task] = None @property def _headers(self) -> Dict[str, str]: @@ -174,6 +215,10 @@ async def start(self) -> None: self.spec.registry_auth, api=self.api ) pod = await self._attach_mounts(pod) + self._deployment = asyncio.create_task(self._deploy_and_record(pod)) + await asyncio.shield(self._deployment) + + async def _deploy_and_record(self, pod: Dict[str, Any]) -> None: result = await self._deploy_pod(pod) self.pod_id = result["id"] log.info("task pod %s deployed for %s", self.pod_id, self.spec.name) @@ -207,9 +252,7 @@ async def _attach_mounts(self, pod: Dict[str, Any]) -> Dict[str, Any]: """resolve storage mounts and apply their placement constraints.""" from .volume import VolumeResolver, attach_pod_mounts - await attach_pod_mounts( - pod, self.spec, VolumeResolver(self.api), self.specs - ) + await attach_pod_mounts(pod, self.spec, VolumeResolver(self.api), self.specs) return pod async def wait_ready(self, timeout: float = READY_TIMEOUT) -> None: @@ -233,17 +276,19 @@ async def wait_ready(self, timeout: float = READY_TIMEOUT) -> None: f"task pod {self.pod_id} did not become ready within {timeout}s" ) - async def execute(self, request: Dict[str, Any], timeout: float) -> Dict[str, Any]: + async def execute( + self, request: Dict[str, Any], timeout: Optional[float] + ) -> Dict[str, Any]: """run to completion: submit, then poll for the result. the pod proxy caps how long a single request can stay open (long jobs would 524), so completion always goes through the background slot + short /result polls. """ - await self.submit(request) - deadline = time.monotonic() + timeout + await self.submit({**request, "timeout": timeout}) + deadline = time.monotonic() + timeout if timeout is not None else None while True: - if time.monotonic() >= deadline: + if deadline is not None and time.monotonic() >= deadline: raise TimeoutError( f"task on pod {self.pod_id} did not finish in {timeout}s" ) @@ -260,7 +305,21 @@ async def submit(self, request: Dict[str, Any]) -> None: and /ping succeeded), 404s here are propagation races and are retried briefly rather than surfaced. """ - url = f"{_proxy_url(self.pod_id)}/submit" + await self._post("/submit", request) + + async def set_timeout(self, timeout: float) -> None: + await self._post("/timeout", {"timeout": timeout}) + + async def _post(self, path: str, request: Dict[str, Any]) -> None: + timeout = request.get("timeout") + if timeout is not None and ( + isinstance(timeout, bool) + or not isinstance(timeout, (int, float)) + or not math.isfinite(timeout) + or timeout < 0 + ): + raise ValueError("task timeout must be a finite non-negative number") + url = f"{_proxy_url(self.pod_id)}{path}" attempts = 6 async with aiohttp.ClientSession() as session: for attempt in range(attempts): @@ -274,6 +333,16 @@ async def submit(self, request: Dict[str, Any]) -> None: await asyncio.sleep(2 * (attempt + 1)) continue resp.raise_for_status() + if timeout is not None: + response = await resp.json() + if ( + not isinstance(response, dict) + or response.get("timeout") != timeout + ): + raise RuntimeError( + "task runtime did not acknowledge the timeout; " + "update the task runtime image" + ) return async def poll_result(self) -> Optional[Dict[str, Any]]: @@ -305,15 +374,46 @@ async def poll_result(self) -> Optional[Dict[str, Any]]: return None async def terminate(self) -> None: + if self._deployment is not None: + try: + await asyncio.shield(self._deployment) + except Exception: + if self.pod_id is None: + return if self.pod_id is None: return try: - await self.api.terminate_pod(self.pod_id) + for attempt in range(DELETE_ATTEMPTS): + try: + await asyncio.wait_for( + self.api.terminate_pod(self.pod_id), DELETE_TIMEOUT + ) + break + except Exception as exc: + status = getattr(exc, "status_code", None) or getattr( + exc, "status", None + ) + if status == 404: + break + retryable = ( + status == 429 + or (status is not None and 500 <= status < 600) + or isinstance( + exc, + ( + aiohttp.ClientConnectionError, + asyncio.TimeoutError, + OSError, + ), + ) + ) + if not retryable or attempt == DELETE_ATTEMPTS - 1: + raise + await asyncio.sleep(0.5 * (2**attempt)) log.info("task pod %s terminated", self.pod_id) except Exception as exc: log.warning( - "failed to terminate task pod %s (terminateAfter is the " - "backstop): %s", + "failed to terminate task pod %s; terminate it in the console: %s", self.pod_id, exc, ) @@ -348,8 +448,10 @@ def pod_id(self) -> Optional[str]: async def wait(self, timeout: Optional[float] = None) -> Any: """wait for the result, terminating the pod on every exit.""" - deadline = time.monotonic() + timeout if timeout is not None else None try: + deadline = time.monotonic() + timeout if timeout is not None else None + if not self._done and timeout is not None: + await self._execution.set_timeout(timeout) while not self._done: if deadline is not None and time.monotonic() >= deadline: raise TimeoutError( @@ -362,10 +464,10 @@ async def wait(self, timeout: Optional[float] = None) -> Any: break await asyncio.sleep(RESULT_POLL_INTERVAL) finally: - await self._execution.terminate() + await finish_cleanup(self._execution.terminate()) return self._result async def cancel(self) -> None: """terminate the pod, abandoning the task.""" - await self._execution.terminate() + await finish_cleanup(self._execution.terminate()) self._done = True diff --git a/runpod/apps/volume.py b/runpod/apps/volume.py index 6673e09b3..943c80b51 100644 --- a/runpod/apps/volume.py +++ b/runpod/apps/volume.py @@ -6,7 +6,7 @@ from collections.abc import Mapping from pathlib import Path, PurePosixPath from types import MappingProxyType -from typing import Any, Dict, List, Optional, Tuple +from typing import Any, Dict, List, Optional, Set, Tuple from .errors import AppError from .utils.client import default_client @@ -167,6 +167,7 @@ def __init__(self, api=None, events: Optional[object] = None): self.events = events self._resolved: Dict[Tuple[str, str], Dict[str, Any]] = {} self._stock = None + self._volume_datacenters: Optional[Set[str]] = None async def _client(self): self._api = default_client(self._api) @@ -213,6 +214,44 @@ async def resolve( raise VolumeError(f"volume '{volume.name}' not found and create=False") dc = record["dataCenter"] if record is not None else volume.datacenter + if record is None: + from .datacenter import DataCenter + + pins = { + DataCenter.from_string(ref.datacenter).value + for ref in [volume] + + [ + ref + for spec in specs + for ref in spec.mounts.values() + if (ref.kind, ref.reference) == key + ] + if ref.datacenter + } + if len(pins) > 1: + raise VolumeError( + f"volume '{volume.name}' has conflicting datacenter pins: " + + ", ".join(sorted(pins)) + ) + dc = next(iter(pins), None) + if not specs and not dc: + raise VolumeError( + f"creating network volume {volume.name!r} requires a datacenter " + "when no app placement constraints are available" + ) + if self._volume_datacenters is None: + try: + self._volume_datacenters = await client.network_volume_datacenters() + except Exception as exc: + raise VolumeError( + f"cannot determine network volume support for '{volume.name}': " + f"datacenter catalog lookup failed: {exc}" + ) from exc + if dc and dc not in self._volume_datacenters: + raise VolumeError( + f"cannot create network volume '{volume.name}' in {dc}: " + "datacenter does not support network volumes" + ) if specs: from .placement import StockMap, _hardware_keys, solve_placement @@ -220,12 +259,12 @@ async def resolve( self._stock = StockMap(client) await self._stock.fetch([k for spec in specs for k in _hardware_keys(spec)]) dc = solve_placement( - specs, self._stock, volume_name=volume.name, existing_dc=dc - ) - elif not dc: - raise VolumeError( - f"creating network volume {volume.name!r} requires a datacenter " - "when no app placement constraints are available" + specs, + self._stock, + volume_name=volume.name, + volume_datacenters=self._volume_datacenters or set(), + volume_dc=dc if record is None else None, + existing_dc=dc if record is not None else None, ) if record is None: @@ -292,10 +331,8 @@ def validate_mounts(mounts: Mapping[str, Volume], kind: str, is_cpu: bool) -> No if kind in ("queue", "api") and mounts: if len(mounts) != 1 or next(iter(mounts)) != str(ENDPOINT_MOUNT_PATH): raise VolumeError("endpoints support one volume mounted at /runpod-volume") - if is_cpu and any( - isinstance(volume, GlobalVolume) for volume in mounts.values() - ): - raise VolumeError("global volumes require a gpu endpoint") + if is_cpu and any(isinstance(volume, GlobalVolume) for volume in mounts.values()): + raise VolumeError("global volumes require gpu compute") def _bind_worker_mounts( diff --git a/tests/test_apps/test_api_client.py b/tests/test_apps/test_api_client.py index 0a4d42771..0303e84f8 100644 --- a/tests/test_apps/test_api_client.py +++ b/tests/test_apps/test_api_client.py @@ -507,6 +507,23 @@ async def graphql(query, *, api_key, variables, anonymous): rest.assert_not_awaited() +class TestNetworkVolumeCapabilities: + async def test_any_supported_tier_is_eligible(self): + catalog = { + "dataCenters": [ + {"id": "US-IL-1", "networkVolumeTypes": []}, + {"id": "EU-RO-1", "networkVolumeTypes": ["STANDARD"]}, + {"id": "US-KS-2", "networkVolumeTypes": ["PREMIUM"]}, + ] + } + with patch( + "runpod.apps.api.run_rest_request_async", + AsyncMock(return_value=catalog), + ): + supported = await AppsApiClient().network_volume_datacenters() + assert supported == {"EU-RO-1", "US-KS-2"} + + class TestVolumesRegistrySecrets: async def test_secret_crud(self): client, patcher = _client_with( diff --git a/tests/test_apps/test_deploy.py b/tests/test_apps/test_deploy.py index 5af21d320..5445c5b52 100644 --- a/tests/test_apps/test_deploy.py +++ b/tests/test_apps/test_deploy.py @@ -31,9 +31,7 @@ def clean_registry(): def _write_project(tmp_path: Path) -> Path: - (tmp_path / "main.py").write_text( - textwrap.dedent( - """ + (tmp_path / "main.py").write_text(textwrap.dedent(""" import runpod from runpod import App @@ -45,9 +43,7 @@ def que1(x: int): if __name__ == "__main__": raise SystemExit("main guard must not run during discovery") - """ - ) - ) + """)) return tmp_path @@ -128,9 +124,7 @@ class Api: def value(self): return {"value": 1} - payload = _deployed_endpoint_input( - app, Api.spec, "env-1", "build-1", "3.12" - ) + payload = _deployed_endpoint_input(app, Api.spec, "env-1", "build-1", "3.12") assert payload["type"] == "LB" assert payload["template"]["ports"] == "80/http" env = {entry["key"]: entry["value"] for entry in payload["template"]["env"]} @@ -141,6 +135,7 @@ def value(self): class TestPackaging: def test_tarball_contains_source_and_manifest(self, tmp_path): _write_project(tmp_path) + (tmp_path / "runpod_manifest.json").write_text('{"app": "forged"}') manifest = {"version": 1, "app": "demo-app", "resources": []} tar_path = package_project(tmp_path, manifest) @@ -151,25 +146,61 @@ def test_tarball_contains_source_and_manifest(self, tmp_path): extracted = json.load(tar.extractfile("runpod_manifest.json")) assert extracted["app"] == "demo-app" - def test_ignores_applied(self, tmp_path): + def test_ignores_and_non_overridable_credentials(self, tmp_path): _write_project(tmp_path) - (tmp_path / "secret.env").write_text("KEY=1") - (tmp_path / ".runpodignore").write_text("secret.env\n") - pycache = tmp_path / "__pycache__" - pycache.mkdir() - (pycache / "x.pyc").write_text("junk") + for name in ("secret.env", ".env.production", "private.pem"): + (tmp_path / name).write_text("private") + (tmp_path / "local.txt").write_text("local") + (tmp_path / "__pycache__").mkdir() + (tmp_path / "__pycache__/main.pyc").write_text("cache") + (tmp_path / ".runpodignore").write_text( + "!*.env\n!.env.production\n!private.pem\nlocal.txt\n" + ) - tar_path = package_project(tmp_path, {"version": 1, "resources": []}) - with tarfile.open(tar_path) as tar: - names = tar.getnames() - assert "secret.env" not in names - assert not any("__pycache__" in n for n in names) + with tarfile.open(package_project(tmp_path, {})) as tar: + assert set(tar.getnames()) == { + "main.py", + ".runpodignore", + "runpod_manifest.json", + } + + def test_credential_like_source_paths_are_included(self, tmp_path): + names = ( + "credentials.py", + "secrets.py", + "credentials", + "secrets", + "nested/secrets/config.py", + "nested/credentials/config.py", + ) + for name in names: + path = tmp_path / "src" / name + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text("source") + with tarfile.open(package_project(tmp_path, {})) as tar: + for name in names: + assert f"src/{name}" in tar.getnames() + + def test_excluded_paths_are_logged_without_contents(self, tmp_path, caplog): + for name in (".env", "private.key", "private.pem", "local.txt", ".aws/config"): + path = tmp_path / name + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text("sensitive-file-contents") + (tmp_path / ".runpodignore").write_text("!*.key\n!.aws/\nlocal.txt\n") + + with tarfile.open(package_project(tmp_path, {})) as tar: + assert set(tar.getnames()) == {".runpodignore", "runpod_manifest.json"} + for name in (".env", "private.key", "private.pem", "local.txt", ".aws/"): + assert f"excluded source path {name!r} from deployment artifact" in caplog.text + assert "sensitive-file-contents" not in caplog.text def test_vendored_env_included_under_env(self, tmp_path): _write_project(tmp_path) env_dir = tmp_path / "built-env" (env_dir / "numpy").mkdir(parents=True) (env_dir / "numpy" / "__init__.py").write_text("") + (env_dir / "numpy" / "cacert.pem").write_text("dependency CA") + (tmp_path / ".runpodignore").write_text("*.pem\n") tar_path = package_project( tmp_path, {"version": 1, "resources": []}, env_dir=env_dir @@ -180,6 +211,47 @@ def test_vendored_env_included_under_env(self, tmp_path): # env dir under project root must not be double-added as source assert "built-env/numpy/__init__.py" not in names assert "main.py" in names + assert tar.extractfile("env/numpy/cacert.pem").read() == b"dependency CA" + + def test_project_ignore_precedence(self, tmp_path): + for name in ( + "root.log", + "nested/keep.log", + "nested/drop.log", + "other/keep.log", + ): + path = tmp_path / name + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(name) + (tmp_path / ".gitignore").write_text("*.log\n") + (tmp_path / "nested/.gitignore").write_text("!*.log\n") + (tmp_path / ".runpodignore").write_text("!root.log\nnested/drop.log\n") + + with tarfile.open(package_project(tmp_path, {})) as tar: + assert set(tar.getnames()) == { + ".gitignore", + ".runpodignore", + "nested/.gitignore", + "root.log", + "nested/keep.log", + "runpod_manifest.json", + } + + def test_symlinks_cannot_escape_source_or_dependencies(self, tmp_path): + project = tmp_path / "project" + project.mkdir() + (project / "main.py").write_text("source") + outside = tmp_path / "outside" + outside.mkdir() + (outside / "secret.txt").write_text("private") + env_dir = tmp_path / "built-env" + env_dir.mkdir() + for root in (project, env_dir): + (root / "linked.py").symlink_to(outside / "secret.txt") + (root / "linked-dir").symlink_to(outside, target_is_directory=True) + + with tarfile.open(package_project(project, {}, env_dir=env_dir)) as tar: + assert set(tar.getnames()) == {"main.py", "runpod_manifest.json"} def _stub_build(tmp_path): @@ -241,7 +313,9 @@ async def test_deploy_reuses_existing_app_and_env(self, tmp_path): api.create_app.assert_not_awaited() api.create_environment.assert_not_awaited() - async def test_deploy_environment_is_used_for_nested_calls(self, tmp_path, monkeypatch): + async def test_deploy_environment_is_used_for_nested_calls( + self, tmp_path, monkeypatch + ): _write_project(tmp_path) (app,) = discover_apps(tmp_path) handle = next(iter(app.resources.values())) @@ -252,7 +326,8 @@ async def test_deploy_environment_is_used_for_nested_calls(self, tmp_path, monke "flashEnvironments": [{"id": "env-prod", "name": "prod"}], } api.prepare_artifact_upload.return_value = { - "uploadUrl": "https://upload", "objectKey": "key-1", + "uploadUrl": "https://upload", + "objectKey": "key-1", } api.finalize_artifact_upload.return_value = {"id": "build-1"} api.save_endpoint.return_value = {"id": "ep-1"} @@ -300,9 +375,7 @@ def test_no_apps_and_failures_raises_with_causes(self, tmp_path): def test_import_time_invocation_diagnosed(self, tmp_path): _write_project(tmp_path) - (tmp_path / "client.py").write_text( - "from main import que1\nque1.remote(1)\n" - ) + (tmp_path / "client.py").write_text("from main import que1\nque1.remote(1)\n") # directory walk: the client file fails with the precise # diagnosis but the app still discovers apps = discover_apps(tmp_path) diff --git a/tests/test_apps/test_dispatch.py b/tests/test_apps/test_dispatch.py index 5e1080ab8..6a2543a83 100644 --- a/tests/test_apps/test_dispatch.py +++ b/tests/test_apps/test_dispatch.py @@ -1,5 +1,6 @@ """tests for context detection and remote dispatch.""" +import asyncio import os import subprocess import sys @@ -10,6 +11,7 @@ import runpod from runpod.apps import App, Context, current_context, is_local from runpod.apps.app import _clear_registry +from runpod.apps.context import block from runpod.apps.errors import ( EndpointNotFound, InvalidResourceError, @@ -417,6 +419,69 @@ def test_api_stub_http(self): class TestSyncBridge: + def test_waits_across_poll_timeouts(self): + async def slow(): + await asyncio.sleep(0.45) + return "complete" + + assert block(slow()) == "complete" + + @pytest.mark.timeout(5) + def test_propagates_operation_timeout(self): + async def fail(): + raise asyncio.TimeoutError("operation deadline") + + with pytest.raises(asyncio.TimeoutError): + block(fail()) + + @pytest.mark.parametrize("bounded", [False, True]) + def test_interrupt_waits_for_cleanup_and_preserves_first_signal(self, bounded): + script = """ +import asyncio +import os +import signal +import threading +from runpod.apps import context + +started = threading.Event() +cleaning = threading.Event() +finished = threading.Event() +if BOUNDED: + context.BRIDGE_CLEANUP_TIMEOUT = 0.05 + +async def operation(): + started.set() + try: + await asyncio.Future() + finally: + cleaning.set() + await asyncio.sleep(1 if BOUNDED else 0.4) + finished.set() + +def interrupt(): + assert started.wait(5) + os.kill(os.getpid(), signal.SIGINT) + assert cleaning.wait(5) + os.kill(os.getpid(), signal.SIGINT) + +threading.Thread(target=interrupt, daemon=True).start() +try: + context.block(operation()) +except KeyboardInterrupt: + assert cleaning.is_set() + assert finished.is_set() == (not BOUNDED) +else: + raise AssertionError("interruption was swallowed") +""" + result = subprocess.run( + [sys.executable, "-c", f"BOUNDED = {bounded!r}\n" + script], + capture_output=True, + text=True, + timeout=10, + check=False, + ) + assert result.returncode == 0, result.stderr + def test_remote_inside_running_loop(self, monkeypatch): """calling sync .remote() from inside an event loop must not raise.""" import asyncio diff --git a/tests/test_apps/test_placement.py b/tests/test_apps/test_placement.py index 7bfe4824c..99aef01d1 100644 --- a/tests/test_apps/test_placement.py +++ b/tests/test_apps/test_placement.py @@ -112,6 +112,7 @@ def test_intersection_picks_shared_dc(self): ], stock, volume_name="models", + volume_datacenters={dc.value for dc in DataCenter.all()}, ) assert dc == "EU-RO-1" @@ -130,6 +131,7 @@ def test_disjoint_hardware_errors_with_details(self): ], stock, volume_name="models", + volume_datacenters={dc.value for dc in DataCenter.all()}, ) def test_existing_dc_is_hard_constraint(self): @@ -139,6 +141,7 @@ def test_existing_dc_is_hard_constraint(self): stock, volume_name="models", existing_dc="EU-RO-1", + volume_datacenters=set(), ) assert dc == "EU-RO-1" @@ -150,6 +153,7 @@ def test_existing_dc_unschedulable_errors(self): stock, volume_name="models", existing_dc="US-KS-2", + volume_datacenters={dc.value for dc in DataCenter.all()}, ) def test_maximin_prefers_worst_case_stock(self): @@ -170,6 +174,7 @@ def test_maximin_prefers_worst_case_stock(self): ], stock, volume_name="v", + volume_datacenters={dc.value for dc in DataCenter.all()}, ) assert dc == "US-KS-2" @@ -188,6 +193,7 @@ def test_mixed_cpu_gpu_sharing(self): ], stock, volume_name="shared", + volume_datacenters={dc.value for dc in DataCenter.all()}, ) assert dc == "EU-RO-1" @@ -208,11 +214,22 @@ async def gpu_stock(gpu_id, dc, gpu_count=1, pods=False): await stock.fetch(_hardware_keys(single)) await stock.fetch(_hardware_keys(pair)) - assert solve_placement([single], stock, volume_name="one") == "EU-RO-1" - assert solve_placement([single, pair], stock, volume_name="shared") == "US-KS-2" + supported = {"EU-RO-1", "US-KS-2"} + dc = solve_placement( + [single], stock, volume_name="one", volume_datacenters=supported + ) + assert dc == "EU-RO-1" + dc = solve_placement( + [single, pair], stock, volume_name="shared", volume_datacenters=supported + ) + assert dc == "US-KS-2" with pytest.raises(PlacementError): solve_placement( - [pair], stock, volume_name="fixed", existing_dc="EU-RO-1" + [pair], + stock, + volume_name="fixed", + existing_dc="EU-RO-1", + volume_datacenters=supported, ) async def test_cpu_task_and_endpoint_require_shared_product_stock(self): @@ -228,7 +245,20 @@ async def cpu_stock(instance_id, dc, *, pods=False): task = ResourceSpec(kind=ResourceKind.TASK, name="task", cpu="cpu5c-2-4") await stock.fetch(_hardware_keys(endpoint)) await stock.fetch(_hardware_keys(task)) - assert solve_placement([endpoint], stock, volume_name="one") == "EU-RO-1" - assert solve_placement([endpoint, task], stock, volume_name="shared") == "US-KS-2" + supported = {"EU-RO-1", "US-KS-2"} + dc = solve_placement( + [endpoint], stock, volume_name="one", volume_datacenters=supported + ) + assert dc == "EU-RO-1" + dc = solve_placement( + [endpoint, task], stock, volume_name="shared", volume_datacenters=supported + ) + assert dc == "US-KS-2" with pytest.raises(PlacementError): - solve_placement([task], stock, volume_name="fixed", existing_dc="EU-RO-1") + solve_placement( + [task], + stock, + volume_name="fixed", + existing_dc="EU-RO-1", + volume_datacenters=supported, + ) diff --git a/tests/test_apps/test_tasks.py b/tests/test_apps/test_tasks.py index e04b1669f..78b761152 100644 --- a/tests/test_apps/test_tasks.py +++ b/tests/test_apps/test_tasks.py @@ -52,7 +52,6 @@ def test_cpu_pod_input(self): assert pod["imageName"] == f"runpod/task:py{local_python_version()}-latest" assert pod["ports"] == "8080/http" - assert pod["terminateAfter"] env = {e["key"]: e["value"] for e in pod["env"]} assert env["RUNPOD_TASK_TOKEN"] == "tok" assert "RUNPOD_RUNTIME_PACKAGE_SPEC" not in env @@ -142,28 +141,45 @@ def t(x): assert _unb64(payload["args"][0]) == 5 assert payload["dependencies"] == ["numpy"] - def test_task_remote_runs_full_lifecycle(self): + def test_task_remote_runs_past_one_hour_and_cleans_up(self): + from types import SimpleNamespace + + from runpod.apps.tasks import TaskExecution + app = App("a") @app.task(name="t", cpu="cpu3c-1-2") def t(x): return x * 2 - with patch("runpod.apps.tasks.TaskExecution") as MockExec: - instance = MockExec.return_value - instance.start = AsyncMock() - instance.wait_ready = AsyncMock() - instance.execute = AsyncMock( - return_value={"success": True, "result": _b64(10)} - ) - instance.terminate = AsyncMock() + clock = [0.0] + pods = {"pod-1"} - result = t.remote(5) + async def poll_result(): + if clock[0] == 0: + clock[0] = 7200.0 + return None + return {"success": True, "result": _b64(10)} - assert result == 10 - instance.start.assert_awaited_once() - instance.wait_ready.assert_awaited_once() - instance.terminate.assert_awaited_once() + async def delete(pod_id): + pods.remove(pod_id) + + execution = TaskExecution(t.spec, api=MagicMock(terminate_pod=delete)) + execution.pod_id = "pod-1" + execution.start = AsyncMock() + execution.wait_ready = AsyncMock() + execution.submit = AsyncMock() + execution.poll_result = poll_result + with ( + patch("runpod.apps.tasks.TaskExecution", return_value=execution), + patch( + "runpod.apps.tasks.time", + SimpleNamespace(monotonic=lambda: clock[0]), + ), + patch("runpod.apps.tasks.RESULT_POLL_INTERVAL", 0), + ): + assert t.remote(5) == 10 + assert not pods def test_task_terminates_pod_on_failure(self): app = App("a") @@ -399,7 +415,9 @@ async def test_shared_volume_placement_accounts_for_siblings(self): datacenter=["US-IL-1"], ) api = AsyncMock() + api.network_volume_datacenters.return_value = {"EU-RO-1", "US-IL-1"} api.list_network_volumes.return_value = [] + api.network_volume_datacenters.return_value = {"EU-RO-1", "US-IL-1"} api.cpu_stock_status.side_effect = lambda instance, dc, *, pods=False: ( "High" if dc == "US-IL-1" else "Low" ) @@ -459,6 +477,7 @@ async def test_execute_polls_to_done(self): with patch("runpod.apps.tasks.asyncio.sleep", AsyncMock()): response = await execution.execute({"fn": "t"}, timeout=60) assert response == {"success": True, "json_result": 4} + execution.submit.assert_awaited_once_with({"fn": "t", "timeout": 60}) async def test_execute_timeout(self): from runpod.apps.tasks import TaskExecution @@ -478,6 +497,46 @@ async def test_execute_timeout(self): await execution.execute({"fn": "t"}, timeout=10) +class TestTaskTimeoutTransport: + @pytest.mark.parametrize("method", ["submit", "set_timeout"]) + @pytest.mark.parametrize("acknowledged", [False, True]) + async def test_runtime_must_acknowledge_timeout(self, method, acknowledged): + from runpod.apps.tasks import TaskExecution + + execution = TaskExecution( + ResourceSpec(kind=ResourceKind.TASK, name="t"), api=MagicMock() + ) + execution.pod_id = "pod-9" + response = MagicMock(status=200) + response.json = AsyncMock( + return_value={"timeout": 60} if acknowledged else {"status": "RUNNING"} + ) + session = MagicMock() + session.__aenter__.return_value = session + session.post.return_value.__aenter__.return_value = response + with patch("runpod.apps.tasks.aiohttp.ClientSession", return_value=session): + argument = {"timeout": 60} if method == "submit" else 60 + if acknowledged: + await getattr(execution, method)(argument) + else: + with pytest.raises(RuntimeError, match="did not acknowledge"): + await getattr(execution, method)(argument) + assert session.post.call_args.kwargs["json"] == {"timeout": 60} + assert session.post.call_args.kwargs["headers"] == execution._headers + + @pytest.mark.parametrize("timeout", [-1, True, "60", float("nan"), float("inf")]) + async def test_invalid_timeouts_fail_before_http(self, timeout): + from runpod.apps.tasks import TaskExecution + + execution = TaskExecution( + ResourceSpec(kind=ResourceKind.TASK, name="t"), api=MagicMock() + ) + with patch("runpod.apps.tasks.aiohttp.ClientSession") as session: + with pytest.raises(ValueError, match="finite non-negative"): + await execution.set_timeout(timeout) + session.assert_not_called() + + class TestTaskJob: def _job(self): from runpod.apps.tasks import TaskExecution, TaskJob @@ -497,17 +556,43 @@ async def test_wait_returns_result_and_terminates(self): assert result == 9 execution.terminate.assert_awaited_once() - async def test_wait_timeout_terminates_pod(self): + async def test_wait_sends_timeout_to_runtime(self): job, execution = self._job() - execution.poll_result = AsyncMock(return_value=None) - with ( - patch("runpod.apps.tasks.asyncio.sleep", AsyncMock()), - patch("runpod.apps.tasks.time.monotonic", side_effect=[0, 100]), - ): - with pytest.raises(TimeoutError): - await job.wait(timeout=10) + execution.poll_result = AsyncMock( + return_value={"success": True, "json_result": 9} + ) + assert await job.wait(timeout=60) == 9 + execution.set_timeout.assert_awaited_once_with(60) + execution.terminate.assert_awaited_once() + + async def test_timeout_configuration_failure_terminates_pod(self): + job, execution = self._job() + execution.set_timeout = AsyncMock(side_effect=RuntimeError("unsupported")) + with pytest.raises(RuntimeError, match="unsupported"): + await job.wait(timeout=60) + execution.poll_result.assert_not_called() execution.terminate.assert_awaited_once() + async def test_wait_timeout_terminates_pod(self): + from runpod.apps.tasks import TaskExecution, TaskJob + + spec = ResourceSpec(kind=ResourceKind.TASK, name="t", cpu=["cpu3c-1-2"]) + execution = TaskExecution(spec, api=MagicMock()) + execution.pod_id = "pod-9" + pods = {"pod-9"} + + async def delete(pod_id): + pods.remove(pod_id) + + execution.api.terminate_pod = delete + execution.set_timeout = AsyncMock() + job = TaskJob(execution) + with pytest.raises(TimeoutError): + await job.wait(timeout=0) + assert not pods + assert job.pod_id is None + execution.set_timeout.assert_awaited_once_with(0) + @pytest.mark.parametrize( "failure", [RuntimeError("poll failed"), asyncio.CancelledError()], @@ -577,3 +662,63 @@ async def terminate(pod_id): await PodTarget(spec, lambda: None).submit({}) assert not pods assert execution.pod_id is None + + +class TestCancellationCleanup: + @pytest.mark.parametrize("method", ["invoke", "submit"]) + async def test_cancel_during_creation_recovers_and_deletes_pod(self, method): + from runpod.apps.tasks import TaskExecution + from runpod.error import QueryError + + creating, release_create = asyncio.Event(), asyncio.Event() + deleting, release_delete = asyncio.Event(), asyncio.Event() + pods = set() + failures = iter([QueryError("rate limited", status_code=429), None]) + + async def deploy(*args, **kwargs): + creating.set() + await release_create.wait() + pods.add("pod-late") + return {"id": "pod-late"} + + async def delete(pod_id): + deleting.set() + await release_delete.wait() + failure = next(failures) + if failure: + raise failure + pods.remove(pod_id) + + spec = ResourceSpec(kind=ResourceKind.TASK, name="t", cpu=["cpu3c-1-2"]) + execution = TaskExecution( + spec, api=MagicMock(deploy_task_pod=deploy, terminate_pod=delete) + ) + with patch("runpod.apps.tasks.TaskExecution", return_value=execution): + task = asyncio.create_task( + getattr(PodTarget(spec, lambda: None), method)({}) + ) + await creating.wait() + task.cancel() + release_create.set() + await deleting.wait() + task.cancel() + release_delete.set() + await asyncio.wait({task}) + assert task.cancelled() + assert not pods + + async def test_failed_cleanup_preserves_error_and_recoverable_pod(self): + from runpod.apps.tasks import TaskExecution, TaskJob + + failure = ValueError("work failed") + spec = ResourceSpec(kind=ResourceKind.TASK, name="t", cpu=["cpu3c-1-2"]) + api = MagicMock( + terminate_pod=AsyncMock(side_effect=RuntimeError("delete failed")) + ) + execution = TaskExecution(spec, api=api) + execution.pod_id = "pod-recoverable" + execution.poll_result = AsyncMock(side_effect=failure) + with pytest.raises(ValueError) as caught: + await TaskJob(execution).wait() + assert caught.value is failure + assert execution.pod_id == "pod-recoverable" diff --git a/tests/test_apps/test_volume.py b/tests/test_apps/test_volume.py index 6c3cae5b7..70c255db3 100644 --- a/tests/test_apps/test_volume.py +++ b/tests/test_apps/test_volume.py @@ -6,6 +6,8 @@ import pytest +from runpod.apps.datacenter import DataCenter +from runpod.apps.placement import PlacementError from runpod.apps.spec import ResourceKind, ResourceSpec from runpod.apps.volume import ( GlobalVolume, @@ -26,6 +28,7 @@ def _spec(name="r", gpu=None, cpu=None): def _api(volumes=None, created=None, global_volumes=None): api = AsyncMock() api.list_network_volumes.return_value = volumes or [] + api.network_volume_datacenters.return_value = {dc.value for dc in DataCenter.all()} api.list_global_volumes.return_value = global_volumes or [] api.create_global_volume.return_value = {"id": "gv-new", "name": "models"} api.create_network_volume.return_value = created or { @@ -140,6 +143,16 @@ def test_invalid_mount_mapping_raises(self, mounts): with pytest.raises(VolumeError): normalize_mounts(mounts) + @pytest.mark.parametrize("kind", list(ResourceKind)) + def test_global_volume_rejects_cpu_resources(self, kind): + with pytest.raises(VolumeError): + ResourceSpec( + kind=kind, + name="cpu", + cpu="cpu3c-1-2", + mounts={"/runpod-volume": GlobalVolume("models")}, + ) + class TestVolumeResolver: def test_existing_by_name(self): @@ -148,9 +161,10 @@ def test_existing_by_name(self): {"id": "nv-1", "name": "models", "size": 50, "dataCenter": "EU-RO-1"} ] ) + api.network_volume_datacenters.side_effect = RuntimeError("catalog unavailable") resolver = VolumeResolver(api) resolved = asyncio.run( - resolver.resolve(NetworkVolume("models"), [_spec(gpu=None)]) + resolver.resolve(NetworkVolume("models", datacenter="US-IL-1"), [_spec()]) ) assert resolved == {"id": "nv-1", "dataCenterId": "EU-RO-1"} api.create_network_volume.assert_not_awaited() @@ -167,14 +181,17 @@ def test_existing_by_id(self): ) assert resolved["id"] == "nv-1" - def test_missing_creates_with_placement(self): + def test_storage_capability_beats_higher_hardware_stock(self): api = _api() + api.network_volume_datacenters.return_value = {"EU-RO-1"} + api.cpu_stock_status.side_effect = lambda instance, dc, *, pods=False: ( + "HIGH" if dc == "US-IL-1" else "LOW" + ) resolver = VolumeResolver(api) resolved = asyncio.run( - resolver.resolve(NetworkVolume("models"), [_spec(gpu=None)]) + resolver.resolve(NetworkVolume("models"), [_spec(cpu="cpu3c-2-4")]) ) - assert resolved["id"] == "nv-new" - api.create_network_volume.assert_awaited_once() + assert resolved == {"id": "nv-new", "dataCenterId": "EU-RO-1"} @pytest.mark.parametrize("volume_type", [NetworkVolume, GlobalVolume]) def test_missing_no_create_raises(self, volume_type): @@ -262,6 +279,25 @@ def test_creation_without_app_requires_explicit_placement(self): ) assert resolved == {"id": "nv-new", "dataCenterId": "EU-RO-1"} + async def test_unsupported_explicit_pin_never_relocates(self): + api = _api() + api.network_volume_datacenters.return_value = {"EU-RO-1"} + with pytest.raises(VolumeError): + await VolumeResolver(api).resolve( + NetworkVolume("models", datacenter="US-IL-1"), [_spec()] + ) + + async def test_disjoint_storage_and_hardware_cannot_create(self): + api = _api() + api.network_volume_datacenters.return_value = {"EU-RO-1"} + api.cpu_stock_status.side_effect = lambda instance, dc, *, pods=False: ( + "HIGH" if dc == "US-IL-1" else "NONE" + ) + with pytest.raises(PlacementError): + await VolumeResolver(api).resolve( + NetworkVolume("models"), [_spec(cpu="cpu3c-2-4")] + ) + class TestTaskVolume: def test_one_volume_per_backend(self): @@ -285,7 +321,7 @@ def test_pod_pins_to_volume_dc(self): spec = ResourceSpec( kind=ResourceKind.TASK, name="t", - cpu=["cpu3c-1-2"], + gpu="4090", mounts={"/models": NetworkVolume("models"), "/data": GlobalVolume("gv-1")}, ) execution = TaskExecution(spec, api=api) @@ -310,22 +346,18 @@ def test_pod_pins_to_volume_dc(self): class TestEndpointMounts: @pytest.mark.parametrize( - "mounts,cpu", + "mounts", [ - ({"/models": NetworkVolume("models")}, None), - ({"/runpod-volume": GlobalVolume("global")}, "cpu3c-1-2"), - ( - { - "/runpod-volume": NetworkVolume("models"), - "/data": GlobalVolume("global"), - }, - None, - ), + {"/models": NetworkVolume("models")}, + { + "/runpod-volume": NetworkVolume("models"), + "/data": GlobalVolume("global"), + }, ], ) - def test_unsupported_attachment_rejected_before_provisioning(self, mounts, cpu): + def test_unsupported_attachment_rejected_before_provisioning(self, mounts): with pytest.raises(VolumeError): - ResourceSpec(kind=ResourceKind.QUEUE, name="queue", cpu=cpu, mounts=mounts) + ResourceSpec(kind=ResourceKind.QUEUE, name="queue", mounts=mounts) def test_global_attachment_does_not_pin_datacenter(self): from runpod import App