From 359ec68985b93674675b72ef390fbfebd0f13d4f Mon Sep 17 00:00:00 2001 From: zeke <40004347+KAJdev@users.noreply.github.com> Date: Wed, 7 Oct 2026 12:10:46 -0700 Subject: [PATCH 1/8] fix: CON-1755 handle Python 3.10 future timeouts --- runpod/apps/context.py | 3 ++- tests/test_apps/test_dispatch.py | 24 ++++++++++++++++++++++++ 2 files changed, 26 insertions(+), 1 deletion(-) diff --git a/runpod/apps/context.py b/runpod/apps/context.py index 3bcffc5a..b9ee470d 100644 --- a/runpod/apps/context.py +++ b/runpod/apps/context.py @@ -1,6 +1,7 @@ """execution context detection and the sync/async bridge.""" import asyncio +import concurrent.futures import os import threading from enum import Enum @@ -73,7 +74,7 @@ def run(cls, coro: Coroutine[Any, Any, Any]) -> Any: while True: try: return future.result(timeout=0.2) - except TimeoutError: + except concurrent.futures.TimeoutError: if future.done(): raise except BaseException: diff --git a/tests/test_apps/test_dispatch.py b/tests/test_apps/test_dispatch.py index 5e1080ab..51426b6d 100644 --- a/tests/test_apps/test_dispatch.py +++ b/tests/test_apps/test_dispatch.py @@ -1,5 +1,7 @@ """tests for context detection and remote dispatch.""" +import asyncio +import concurrent.futures import os import subprocess import sys @@ -10,6 +12,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, @@ -54,6 +57,27 @@ def test_dev_session(self, monkeypatch): assert is_local() is True +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): + failure = concurrent.futures.TimeoutError("operation deadline") + + async def fail(): + raise failure + + with pytest.raises(concurrent.futures.TimeoutError) as caught: + block(fail()) + + assert caught.value is failure + + class TestArgsToInput: def test_positional_mapped_to_names(self): def fn(prompt, temp): From 38624ad2d982a0573bd2aa2fc42771778e701f0c Mon Sep 17 00:00:00 2001 From: zeke <40004347+KAJdev@users.noreply.github.com> Date: Wed, 7 Oct 2026 12:37:28 -0700 Subject: [PATCH 2/8] test: cover CON-1755 sync bridge timeout propagation --- tests/test_apps/test_dispatch.py | 37 +++++++++++++------------------- 1 file changed, 15 insertions(+), 22 deletions(-) diff --git a/tests/test_apps/test_dispatch.py b/tests/test_apps/test_dispatch.py index 51426b6d..ce129594 100644 --- a/tests/test_apps/test_dispatch.py +++ b/tests/test_apps/test_dispatch.py @@ -1,7 +1,6 @@ """tests for context detection and remote dispatch.""" import asyncio -import concurrent.futures import os import subprocess import sys @@ -57,27 +56,6 @@ def test_dev_session(self, monkeypatch): assert is_local() is True -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): - failure = concurrent.futures.TimeoutError("operation deadline") - - async def fail(): - raise failure - - with pytest.raises(concurrent.futures.TimeoutError) as caught: - block(fail()) - - assert caught.value is failure - - class TestArgsToInput: def test_positional_mapped_to_names(self): def fn(prompt, temp): @@ -441,6 +419,21 @@ 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()) + def test_remote_inside_running_loop(self, monkeypatch): """calling sync .remote() from inside an event loop must not raise.""" import asyncio From d117d02299b2c38dffefe14f5bdc6aecd93fa13a Mon Sep 17 00:00:00 2001 From: zeke <40004347+KAJdev@users.noreply.github.com> Date: Wed, 7 Oct 2026 12:38:36 -0700 Subject: [PATCH 3/8] fix: protect CON-1753 deployment artifacts from secret files --- docs/cli/references/projects.md | 46 +++++ requirements.txt | 1 + runpod/apps/deploy.py | 172 +++++++++++++------ runpod/apps/init.py | 20 ++- tests/test_apps/test_deploy.py | 249 ++++++++++++++++++++++++++-- tests/test_cli/test_init_command.py | 22 +++ 6 files changed, 440 insertions(+), 70 deletions(-) diff --git a/docs/cli/references/projects.md b/docs/cli/references/projects.md index 2480b204..90fc05bc 100644 --- a/docs/cli/references/projects.md +++ b/docs/cli/references/projects.md @@ -11,3 +11,49 @@ 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 patterns to +source files, using `pathspec`. Precedence, from lowest to highest, is: + +1. Default local-file exclusions: virtual environments, Python caches, test + directories and test modules, `node_modules`, `.DS_Store`, and `*.tar.gz`. +2. `.gitignore` files within the project, with deeper files overriding parent + rules for their subtree. +3. The project-root `.runpodignore`. + +Ancestor and global Git ignore files are not read. Within each ignore file, the +last matching rule wins. `/` anchors a pattern to that ignore file's directory; +trailing `/` matches directories; `**`, comments, escaped characters and `!` +negation follow Git ignore syntax. Excluded directories are not traversed, so +re-include the parent before its files. For example, to deploy a fixture from +the otherwise excluded `tests` directory: + +```gitignore +!tests/ +tests/* +!tests/fixture.json +``` + +Some safeguards cannot be overridden by negation: `.git`, `.runpod`, `.flash`, +the root `env/` and `runpod_manifest.json` paths, and credential-like source +paths. Credential safeguards include `.env` and its variants, `*.env` and its +variants, `*.pem`, `*.key`, SSH private-key names, `.ssh`, `.aws`, `.azure`, +`.kube`, `.docker/config.json`, `.netrc`, `.npmrc`, `.pypirc`, `.git-credentials`, +`.boto`, `credentials`/`credentials.*`, `secrets`/`secrets.*`, and +`service-account*.json`/`service_account*.json`. Use worker environment variables +or a secret store for credentials instead of packaging them with source. + +These are filename safeguards, not secret scanning: credentials embedded in +ordinary source or differently named files can still be uploaded. Review the +build-only artifact before deployment. + +The generated manifest and vendored `env/` are added separately and cannot be +replaced by source files. Source ignore rules, including the PEM/key safeguards, +do not apply to vendored dependencies, so dependency CA bundles are preserved. +The output artifact itself and the dependency build directory are never copied +back into source. Source and dependency symlinks are omitted, including +project-contained links and linked ignore files; directory links are not +traversed. Hardlinked regular files are stored as independent regular members, +so the archive requires no link extraction support. diff --git a/requirements.txt b/requirements.txt index a19a076d..21ad7c4b 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/deploy.py b/runpod/apps/deploy.py index a864099a..7d2741c7 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,66 @@ 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", + "credentials", + "credentials.*", + "secrets", + "secrets.*", + "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 +166,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 +193,70 @@ 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) + ): + 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) + ): 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 0bb0e9cd..d7f21201 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/tests/test_apps/test_deploy.py b/tests/test_apps/test_deploy.py index 5af21d32..41152777 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"]} @@ -181,6 +175,232 @@ def test_vendored_env_included_under_env(self, tmp_path): assert "built-env/numpy/__init__.py" not in names assert "main.py" in names + def test_source_credentials_and_local_files_excluded_by_default(self, tmp_path): + excluded = { + ".env", + ".env.local", + "settings/dev.env", + "settings/dev.env.backup", + "keys/server.pem", + "keys/server.key", + ".ssh/id_rsa", + "keys/id_ed25519", + ".aws/credentials", + ".azure/token", + ".kube/config", + ".docker/config.json", + ".netrc", + ".npmrc", + ".pypirc", + ".git-credentials", + ".boto", + "credentials.json", + "secrets.yaml", + "service-account-prod.json", + "service_account.json", + "tests/test_app.py", + "test/unit.py", + "test_app.py", + "app_test.py", + ".venv/lib/local.py", + "venv/lib/local.py", + "env/lib/local.py", + ".pytest_cache/state", + "__pycache__/main.pyc", + } + for name in excluded | {"main.py", "package/module.py", "data/input.json"}: + path = tmp_path / name + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(name) + + with tarfile.open(package_project(tmp_path, {})) as tar: + assert set(tar.getnames()) == { + "main.py", + "package/module.py", + "data/input.json", + "runpod_manifest.json", + } + + def test_broad_negation_cannot_override_protected_paths(self, tmp_path): + protected = { + ".env.production", + "keys/private.pem", + "keys/private.key", + ".git/config", + ".runpod/cache", + ".flash/cache", + ".aws/credentials", + "env/local.py", + "runpod_manifest.json/forged.json", + } + for name in protected | {"tests/fixture.py", ".venv/local.py", "main.py"}: + path = tmp_path / name + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(name) + (tmp_path / ".runpodignore").write_text("!**\n") + + with tarfile.open(package_project(tmp_path, {"app": "generated"})) as tar: + assert set(tar.getnames()) == { + ".runpodignore", + "tests/fixture.py", + ".venv/local.py", + "main.py", + "runpod_manifest.json", + } + assert json.load(tar.extractfile("runpod_manifest.json")) == { + "app": "generated" + } + + def test_gitignore_scopes_and_runpod_precedence(self, tmp_path): + for name in [ + "main.py", + "root-only.txt", + "nested/root-only.txt", + "omit.log", + "nested/omit.log", + "nested/keep.log", + "nested/hidden.txt", + "nested/override.txt", + "sibling/hidden.txt", + ]: + path = tmp_path / name + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(name) + (tmp_path / ".gitignore").write_text("/root-only.txt\n*.log\n") + (tmp_path / "nested/.gitignore").write_text( + "hidden.txt\noverride.txt\n!keep.log\n" + ) + (tmp_path / ".runpodignore").write_text( + "!/root-only.txt\n!nested/override.txt\nnested/keep.log\n" + ) + + with tarfile.open(package_project(tmp_path, {})) as tar: + assert set(tar.getnames()) == { + ".gitignore", + ".runpodignore", + "nested/.gitignore", + "main.py", + "root-only.txt", + "nested/root-only.txt", + "nested/override.txt", + "sibling/hidden.txt", + "runpod_manifest.json", + } + + def test_directory_patterns_require_parent_reinclusion(self, tmp_path): + for name in [ + "data/keep.txt", + "cache/keep.txt", + "nested/data/drop.txt", + "tests/fixture.py", + ]: + path = tmp_path / name + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(name) + (tmp_path / "ordinary").write_text("not a directory") + (tmp_path / ".gitignore").write_text("data/\ncache/\nordinary/\n") + (tmp_path / ".runpodignore").write_text( + "!data/keep.txt\n!cache/\ncache/*\n!cache/keep.txt\n!tests/\n" + ) + + with tarfile.open(package_project(tmp_path, {})) as tar: + assert set(tar.getnames()) == { + ".gitignore", + ".runpodignore", + "cache/keep.txt", + "ordinary", + "tests/fixture.py", + "runpod_manifest.json", + } + + def test_ignore_escaping_and_double_star(self, tmp_path): + for name in [ + "#notes", + "!notes", + "trailing ", + "assets/a/cache/x", + "assets/cache/y", + "assets/a/keep", + ]: + path = tmp_path / name + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(name) + (tmp_path / ".runpodignore").write_text( + "\\#notes\n\\!notes\ntrailing\\ \nassets/**/cache/\n" + ) + + with tarfile.open(package_project(tmp_path, {})) as tar: + assert set(tar.getnames()) == { + ".runpodignore", + "assets/a/keep", + "runpod_manifest.json", + } + + def test_generated_entries_and_dependency_certificates_preserved(self, tmp_path): + (tmp_path / "main.py").write_text("source") + (tmp_path / "runpod_manifest.json").write_text('{"app": "forged"}') + (tmp_path / "env").mkdir() + (tmp_path / "env/local.py").write_text("local environment") + env_dir = tmp_path / "built-env" + for name in ["certifi/cacert.pem", "library/public.key", "package.py"]: + path = env_dir / name + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(name) + (tmp_path / ".runpodignore").write_text("!**\n*.pem\n*.key\n") + output = tmp_path / "build.bundle" + + package_project(tmp_path, {"app": "generated"}, output, env_dir) + + with tarfile.open(output) as tar: + assert set(tar.getnames()) == { + ".runpodignore", + "main.py", + "env/certifi/cacert.pem", + "env/library/public.key", + "env/package.py", + "runpod_manifest.json", + } + assert tar.getnames().count("runpod_manifest.json") == 1 + assert json.load(tar.extractfile("runpod_manifest.json")) == { + "app": "generated" + } + assert ( + tar.extractfile("env/certifi/cacert.pem").read() + == b"certifi/cacert.pem" + ) + + def test_links_are_not_followed_and_hardlinks_are_regular_members(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") + (project / "external.py").symlink_to(outside / "secret.txt") + (project / "external-dir").symlink_to(outside, target_is_directory=True) + (project / "internal.py").symlink_to("main.py") + (project / "broken.py").symlink_to("missing.py") + (project / "hardlink.py").hardlink_to(project / "main.py") + (outside / "ignore").write_text("*\n") + (project / ".gitignore").symlink_to(outside / "ignore") + (project / ".runpodignore").symlink_to(outside / "ignore") + (tmp_path / ".gitignore").write_text("*\n") + env_dir = tmp_path / "built-env" + env_dir.mkdir() + (env_dir / "package.py").write_text("dependency") + (env_dir / "external.py").symlink_to(outside / "secret.txt") + (env_dir / "external-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", + "hardlink.py", + "env/package.py", + "runpod_manifest.json", + } + assert all(member.isfile() for member in tar.getmembers()) + assert tar.extractfile("hardlink.py").read() == b"source" + def _stub_build(tmp_path): """patch environment vendoring with a tiny fake env tree.""" @@ -241,7 +461,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 +474,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 +523,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_cli/test_init_command.py b/tests/test_cli/test_init_command.py index 42be0898..fe0b5915 100644 --- a/tests/test_cli/test_init_command.py +++ b/tests/test_cli/test_init_command.py @@ -1,7 +1,10 @@ """rp flash init: project scaffolding.""" +import tarfile + from click.testing import CliRunner +from runpod.apps.deploy import package_project from runpod.apps.init import create_project, detect_conflicts from runpod.rp_cli.main import cli @@ -29,6 +32,25 @@ def test_creates_directory(self, tmp_path): create_project(target, "new-project") assert (target / "main.py").exists() + def test_scaffold_packages_source_without_local_secrets(self, tmp_path): + create_project(tmp_path, "my-app") + (tmp_path / ".env.production").write_text("TOKEN=private") + (tmp_path / "private.pem").write_text("private") + (tmp_path / "tests").mkdir() + (tmp_path / "tests/test_app.py").write_text("local test") + (tmp_path / ".gitignore").write_text("local-data/\n") + (tmp_path / "local-data").mkdir() + (tmp_path / "local-data/input.json").write_text("{}") + + with tarfile.open(package_project(tmp_path, {})) as tar: + assert set(tar.getnames()) == { + "main.py", + "requirements.txt", + ".runpodignore", + ".gitignore", + "runpod_manifest.json", + } + class TestDetectConflicts: def test_empty_dir_no_conflicts(self, tmp_path): From dc20914ffbf8350b8e84f4287d02674c16cbb199 Mon Sep 17 00:00:00 2001 From: zeke <40004347+KAJdev@users.noreply.github.com> Date: Wed, 7 Oct 2026 12:39:09 -0700 Subject: [PATCH 4/8] fix: complete CON-1751 task cleanup after cancellation --- CHANGELOG.md | 4 + examples/apps/README.md | 16 ++- runpod/apps/context.py | 39 +++++++- runpod/apps/targets.py | 22 ++-- runpod/apps/tasks.py | 100 ++++++++++++++++-- tests/test_apps/test_dispatch.py | 87 ++++++++++++---- tests/test_apps/test_tasks.py | 167 +++++++++++++++++++++++++++++-- 7 files changed, 381 insertions(+), 54 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index b3334834..916fa6d6 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,10 @@ * 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 + +* Wait for task cleanup after interruption, preserve cancellation and task errors, recover in-flight pod creation, and retry transient task deletion failures without cancelling detached work. + ## [1.12.0](https://github.com/runpod/runpod-python/compare/v1.11.0...v1.12.0) (2026-08-10) diff --git a/examples/apps/README.md b/examples/apps/README.md index 4fd1fa8a..1894cd0c 100644 --- a/examples/apps/README.md +++ b/examples/apps/README.md @@ -20,6 +20,16 @@ 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 gives cleanup up to 30 seconds to recover an in-flight create response +and delete the pod; synchronous calls allow up to 35 seconds for the coroutine to +acknowledge cleanup after ctrl-c. repeated interrupts do not interrupt deletion. +transient delete failures are retried up to three times. cleanup errors do not +replace an existing task error or cancellation, and failed deletion retains +`job.pod_id` for recovery. cleanup that exceeds the grace period continues only +while the client's event loop remains alive. check the pod in the console after a +cleanup warning; client process exit alone does not stop billing. diff --git a/runpod/apps/context.py b/runpod/apps/context.py index b9ee470d..b21d3876 100644 --- a/runpod/apps/context.py +++ b/runpod/apps/context.py @@ -2,11 +2,17 @@ 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.""" @@ -69,7 +75,26 @@ 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: @@ -78,7 +103,17 @@ def run(cls, coro: Coroutine[Any, Any, Any]) -> Any: 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/targets.py b/runpod/apps/targets.py index 40769eae..763b39c3 100644 --- a/runpod/apps/targets.py +++ b/runpod/apps/targets.py @@ -847,7 +847,7 @@ async def invoke( ) -> 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 +874,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 +884,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 +902,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 +910,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 d1707f30..1f4124d5 100644 --- a/runpod/apps/tasks.py +++ b/runpod/apps/tasks.py @@ -15,6 +15,7 @@ import logging import os as _os import secrets +import sys import time from datetime import datetime, timedelta, timezone from typing import Any, Dict, List, Optional @@ -22,6 +23,7 @@ import aiohttp from .api import AppsApiClient, is_capacity_error +from .context import CLEANUP_TIMEOUT from .errors import RemoteExecutionError from .spec import ResourceSpec @@ -29,12 +31,58 @@ TASK_PORT = 8080 -# safety net: pods self-terminate server-side after this long +# requested server-side deadline for abandoned tasks 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) + try: + task.result() + except asyncio.CancelledError: + pass + except Exception: + log.warning("background task pod cleanup failed", exc_info=True) def _runtime_command() -> str: @@ -160,6 +208,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 +223,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 +260,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: @@ -305,15 +356,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, ) @@ -362,10 +444,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/tests/test_apps/test_dispatch.py b/tests/test_apps/test_dispatch.py index 51426b6d..4bab6a45 100644 --- a/tests/test_apps/test_dispatch.py +++ b/tests/test_apps/test_dispatch.py @@ -1,7 +1,6 @@ """tests for context detection and remote dispatch.""" import asyncio -import concurrent.futures import os import subprocess import sys @@ -57,27 +56,6 @@ def test_dev_session(self, monkeypatch): assert is_local() is True -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): - failure = concurrent.futures.TimeoutError("operation deadline") - - async def fail(): - raise failure - - with pytest.raises(concurrent.futures.TimeoutError) as caught: - block(fail()) - - assert caught.value is failure - - class TestArgsToInput: def test_positional_mapped_to_names(self): def fn(prompt, temp): @@ -441,6 +419,71 @@ 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): + failure = asyncio.TimeoutError("operation deadline") + + async def fail(): + raise failure + + 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_tasks.py b/tests/test_apps/test_tasks.py index e04b1669..c220f660 100644 --- a/tests/test_apps/test_tasks.py +++ b/tests/test_apps/test_tasks.py @@ -399,6 +399,7 @@ 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.cpu_stock_status.side_effect = lambda instance, dc, *, pods=False: ( "High" if dc == "US-IL-1" else "Low" @@ -498,15 +499,22 @@ async def test_wait_returns_result_and_terminates(self): execution.terminate.assert_awaited_once() async def test_wait_timeout_terminates_pod(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.terminate.assert_awaited_once() + 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 + job = TaskJob(execution) + with pytest.raises(TimeoutError): + await job.wait(timeout=0) + assert not pods + assert job.pod_id is None @pytest.mark.parametrize( "failure", @@ -577,3 +585,144 @@ 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 + + spec = ResourceSpec(kind=ResourceKind.TASK, name="t", cpu=["cpu3c-1-2"]) + creating = asyncio.Event() + release_create = asyncio.Event() + deleting = asyncio.Event() + release_delete = asyncio.Event() + pods = set() + + 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() + pods.remove(pod_id) + + api = MagicMock() + api.deploy_task_pod = deploy + api.terminate_pod = delete + execution = TaskExecution(spec, api=api) + with patch("runpod.apps.tasks.TaskExecution", return_value=execution): + task = asyncio.create_task( + getattr(PodTarget(spec, lambda: None), method)({}) + ) + await creating.wait() + task.cancel("original interruption") + await asyncio.sleep(0) + task.cancel("second interruption") + release_create.set() + await deleting.wait() + task.cancel("third interruption") + release_delete.set() + with pytest.raises(asyncio.CancelledError): + await task + assert not pods + assert execution.pod_id is None + + @pytest.mark.parametrize("method", ["invoke", "submit", "wait"]) + @pytest.mark.parametrize("cancelled", [False, True]) + async def test_cleanup_failure_preserves_original(self, method, cancelled): + from runpod.apps.tasks import TaskExecution, TaskJob + + spec = ResourceSpec(kind=ResourceKind.TASK, name="t", cpu=["cpu3c-1-2"]) + failure = ( + asyncio.CancelledError("cancelled") + if cancelled + else ValueError("work failed") + ) + api = MagicMock() + api.terminate_pod = AsyncMock(side_effect=RuntimeError("delete failed")) + execution = TaskExecution(spec, api=api) + execution.pod_id = "pod-recoverable" + execution.start = AsyncMock(side_effect=failure) + execution.poll_result = AsyncMock(side_effect=failure) + with patch("runpod.apps.tasks.TaskExecution", return_value=execution): + with pytest.raises(type(failure)) as caught: + if method == "wait": + await TaskJob(execution).wait() + else: + await getattr(PodTarget(spec, lambda: None), method)({}) + assert caught.value is failure + assert execution.pod_id == "pod-recoverable" + + async def test_cleanup_deadline_does_not_abandon_deletion(self, monkeypatch): + from runpod.apps import tasks + + monkeypatch.setattr(tasks, "CLEANUP_TIMEOUT", 0.02) + entered = asyncio.Event() + release = asyncio.Event() + deleted = asyncio.Event() + + async def delete(): + entered.set() + await release.wait() + deleted.set() + + async def operation(): + try: + await asyncio.Future() + finally: + await tasks.finish_cleanup(delete()) + + task = asyncio.create_task(operation()) + await asyncio.sleep(0) + task.cancel("original") + await entered.wait() + task.cancel("again") + with pytest.raises(asyncio.CancelledError): + await asyncio.wait_for(task, 0.5) + assert not deleted.is_set() + release.set() + await asyncio.wait_for(deleted.wait(), 0.5) + + @pytest.mark.parametrize("status", [429, 503, None]) + async def test_transient_delete_recovers(self, status): + from runpod.apps.tasks import TaskExecution + from runpod.error import QueryError + + spec = ResourceSpec(kind=ResourceKind.TASK, name="t", cpu=["cpu3c-1-2"]) + pods = {"pod-retry"} + attempts = 0 + + async def delete(pod_id): + nonlocal attempts + attempts += 1 + if attempts == 1: + if status is None: + raise OSError("connection reset") + raise QueryError("retry later", status_code=status) + pods.remove(pod_id) + + api = MagicMock() + api.terminate_pod = delete + execution = TaskExecution(spec, api=api) + execution.pod_id = "pod-retry" + await execution.terminate() + assert not pods + assert execution.pod_id is None + + async def test_retry_exhaustion_keeps_recoverable_id(self): + from runpod.apps import tasks + from runpod.error import QueryError + + spec = ResourceSpec(kind=ResourceKind.TASK, name="t", cpu=["cpu3c-1-2"]) + api = MagicMock() + api.terminate_pod = AsyncMock(side_effect=QueryError("busy", status_code=503)) + execution = tasks.TaskExecution(spec, api=api) + execution.pod_id = "pod-retry" + with pytest.raises(QueryError): + await execution.terminate() + assert api.terminate_pod.await_count == tasks.DELETE_ATTEMPTS + assert execution.pod_id == "pod-retry" From 21a2d823351f58c5f182778b1f20d330641788d9 Mon Sep 17 00:00:00 2001 From: zeke <40004347+KAJdev@users.noreply.github.com> Date: Wed, 7 Oct 2026 12:39:36 -0700 Subject: [PATCH 5/8] fix: use storage-capable placement for CON-1752 volumes --- CHANGELOG.md | 4 + README.md | 9 +- runpod/apps/api.py | 27 +++++- runpod/apps/datacenter.py | 6 +- runpod/apps/placement.py | 29 +++--- runpod/apps/volume.py | 53 +++++++++-- tests/test_apps/test_api_client.py | 40 ++++++++ tests/test_apps/test_placement.py | 60 ++++++++++-- tests/test_apps/test_tasks.py | 1 + tests/test_apps/test_volume.py | 141 +++++++++++++++++++++++++++++ 10 files changed, 339 insertions(+), 31 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index b3334834..eae152aa 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,10 @@ * 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 + +* 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. + ## [1.12.0](https://github.com/runpod/runpod-python/compare/v1.11.0...v1.12.0) (2026-08-10) diff --git a/README.md b/README.md index be098910..8b6f7f1c 100644 --- a/README.md +++ b/README.md @@ -76,8 +76,13 @@ Queue `.remote()` calls request a synchronous result and poll the same job if th ## Storage `NetworkVolume(name_or_id, size=50, datacenter=None, create=True)` references -datacenter-local storage. Apps resolve names and choose a datacenter compatible -with every resource sharing the volume. `GlobalVolume(name_or_id, create=True)` +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. diff --git a/runpod/apps/api.py b/runpod/apps/api.py index 64ba26b3..30b6c301 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/datacenter.py b/runpod/apps/datacenter.py index 04b159b1..1ba61f46 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/placement.py b/runpod/apps/placement.py index 931bece7..e9e99add 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/volume.py b/runpod/apps/volume.py index 6673e09b..f3a2523e 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: diff --git a/tests/test_apps/test_api_client.py b/tests/test_apps/test_api_client.py index 0a4d4277..d5b17b4e 100644 --- a/tests/test_apps/test_api_client.py +++ b/tests/test_apps/test_api_client.py @@ -507,6 +507,46 @@ 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"]}, + {"id": "US-MO-2", "networkVolumeTypes": ["STANDARD", "PREMIUM"]}, + ] + } + with patch( + "runpod.apps.api.run_rest_request_async", + AsyncMock(return_value=catalog), + ) as rest: + supported = await AppsApiClient( + api_key="test-key" + ).network_volume_datacenters() + assert supported == {"EU-RO-1", "US-KS-2", "US-MO-2"} + rest.assert_awaited_once_with( + "GET", "/v2/catalog/datacenters", api_key="test-key" + ) + + @pytest.mark.parametrize( + "catalog", + [ + {}, + {"dataCenters": None}, + {"dataCenters": [{"id": "EU-RO-1"}]}, + {"dataCenters": [{"id": "EU-RO-1", "networkVolumeTypes": "STANDARD"}]}, + ], + ) + async def test_missing_capability_is_not_assumed_available(self, catalog): + with patch( + "runpod.apps.api.run_rest_request_async", + AsyncMock(return_value=catalog), + ): + with pytest.raises(QueryError): + await AppsApiClient().network_volume_datacenters() + + class TestVolumesRegistrySecrets: async def test_secret_crud(self): client, patcher = _client_with( diff --git a/tests/test_apps/test_placement.py b/tests/test_apps/test_placement.py index 7bfe4824..de3d641a 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={dc.value for dc in DataCenter.all()}, ) 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,31 @@ 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" + assert ( + solve_placement( + [single], + stock, + volume_name="one", + volume_datacenters={dc.value for dc in DataCenter.all()}, + ) + == "EU-RO-1" + ) + assert ( + solve_placement( + [single, pair], + stock, + volume_name="shared", + volume_datacenters={dc.value for dc in DataCenter.all()}, + ) + == "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={dc.value for dc in DataCenter.all()}, ) async def test_cpu_task_and_endpoint_require_shared_product_stock(self): @@ -228,7 +254,29 @@ 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" + assert ( + solve_placement( + [endpoint], + stock, + volume_name="one", + volume_datacenters={dc.value for dc in DataCenter.all()}, + ) + == "EU-RO-1" + ) + assert ( + solve_placement( + [endpoint, task], + stock, + volume_name="shared", + volume_datacenters={dc.value for dc in DataCenter.all()}, + ) + == "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={dc.value for dc in DataCenter.all()}, + ) diff --git a/tests/test_apps/test_tasks.py b/tests/test_apps/test_tasks.py index e04b1669..57d4aa70 100644 --- a/tests/test_apps/test_tasks.py +++ b/tests/test_apps/test_tasks.py @@ -400,6 +400,7 @@ async def test_shared_volume_placement_accounts_for_siblings(self): ) api = AsyncMock() 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" ) diff --git a/tests/test_apps/test_volume.py b/tests/test_apps/test_volume.py index 6c3cae5b..fdf5fd48 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 { @@ -262,6 +265,144 @@ def test_creation_without_app_requires_explicit_placement(self): ) assert resolved == {"id": "nv-new", "dataCenterId": "EU-RO-1"} + async 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" + ) + resolved = await VolumeResolver(api).resolve( + NetworkVolume("models"), [_spec(cpu="cpu3c-2-4")] + ) + assert resolved["dataCenterId"] == "EU-RO-1" + assert api.create_network_volume.call_args.kwargs["data_center_id"] == "EU-RO-1" + + async def test_reported_datacenters_remain_eligible_when_catalog_supports_them( + self, + ): + api = _api() + api.network_volume_datacenters.return_value = {"US-IL-1"} + resolved = await VolumeResolver(api).resolve( + NetworkVolume("models"), [_spec(cpu="cpu3c-2-4")] + ) + assert resolved["dataCenterId"] == "US-IL-1" + + @pytest.mark.parametrize("specs", [[], [_spec(cpu="cpu3c-2-4")]]) + async def test_unsupported_explicit_pin_never_relocates(self, specs): + api = _api() + api.network_volume_datacenters.return_value = {"EU-RO-1"} + with pytest.raises(VolumeError, match="does not support network volumes"): + await VolumeResolver(api).resolve( + NetworkVolume("models", datacenter="US-IL-1"), specs + ) + api.create_network_volume.assert_not_awaited() + + @pytest.mark.parametrize("supported", [set(), {"EU-RO-1"}]) + async def test_empty_storage_hardware_intersection_never_creates(self, supported): + api = _api() + api.network_volume_datacenters.return_value = supported + 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")] + ) + api.create_network_volume.assert_not_awaited() + + @pytest.mark.parametrize("reference", ["models", "nv-1"]) + async def test_existing_volume_ignores_new_creation_capability(self, reference): + api = _api(volumes=[{"id": "nv-1", "name": "models", "dataCenter": "US-IL-1"}]) + api.network_volume_datacenters.side_effect = RuntimeError("catalog unavailable") + resolved = await VolumeResolver(api).resolve( + NetworkVolume(reference, datacenter="EU-RO-1"), + [_spec(cpu="cpu3c-2-4")], + ) + assert resolved == {"id": "nv-1", "dataCenterId": "US-IL-1"} + api.network_volume_datacenters.assert_not_awaited() + api.create_network_volume.assert_not_awaited() + + async def test_catalog_failure_does_not_create(self): + api = _api() + api.network_volume_datacenters.side_effect = RuntimeError("catalog unavailable") + with pytest.raises(VolumeError, match="datacenter catalog lookup failed"): + await VolumeResolver(api).resolve( + NetworkVolume("models", datacenter="EU-RO-1") + ) + api.create_network_volume.assert_not_awaited() + + async def test_creation_uses_one_catalog_snapshot_per_run(self): + api = _api() + api.network_volume_datacenters.side_effect = [{"EU-RO-1"}, {"US-IL-1"}] + resolver = VolumeResolver(api) + first = await resolver.resolve(NetworkVolume("first"), [_spec()]) + second = await resolver.resolve(NetworkVolume("second"), [_spec()]) + assert first["dataCenterId"] == second["dataCenterId"] == "EU-RO-1" + assert api.network_volume_datacenters.await_count == 1 + + async def test_all_shared_consumers_constrain_creation(self): + api = _api() + api.network_volume_datacenters.return_value = {"EU-RO-1", "US-IL-1"} + volume = NetworkVolume("models") + producer = ResourceSpec( + kind=ResourceKind.TASK, + name="producer", + cpu="cpu3c-2-4", + mounts={"/models": volume}, + ) + consumer = ResourceSpec( + kind=ResourceKind.QUEUE, + name="consumer", + gpu="4090", + gpu_count=2, + datacenter="EU-RO-1", + mounts={"/runpod-volume": NetworkVolume("models")}, + ) + unrelated = ResourceSpec( + kind=ResourceKind.TASK, name="unrelated", datacenter="US-IL-1" + ) + api.cpu_stock_status.side_effect = lambda instance, dc, *, pods=False: ( + "HIGH" if dc == "US-IL-1" else "LOW" + ) + api.gpu_stock_status.side_effect = lambda gpu, dc, gpu_count=1, pods=False: ( + "LOW" if dc == "EU-RO-1" and gpu_count == 2 and not pods else "NONE" + ) + bindings = await VolumeResolver(api).resolve_mounts( + producer.mounts, [producer, consumer, unrelated] + ) + assert bindings[0]["dataCenterId"] == "EU-RO-1" + assert api.create_network_volume.call_args.kwargs["data_center_id"] == "EU-RO-1" + + async def test_supported_volume_pin_conflicting_with_consumer_never_relocates(self): + api = _api() + spec = ResourceSpec(kind=ResourceKind.TASK, name="worker", datacenter="EU-RO-1") + with pytest.raises(PlacementError): + await VolumeResolver(api).resolve( + NetworkVolume("models", datacenter="US-IL-1"), [spec] + ) + api.create_network_volume.assert_not_awaited() + + async def test_shared_volume_reference_pin_is_not_ignored(self): + api = _api() + api.network_volume_datacenters.return_value = {"EU-RO-1"} + volume = NetworkVolume("models") + sibling = ResourceSpec( + kind=ResourceKind.TASK, + name="sibling", + mounts={"/models": NetworkVolume("models", datacenter="US-IL-1")}, + ) + with pytest.raises(VolumeError): + await VolumeResolver(api).resolve(volume, [sibling]) + api.create_network_volume.assert_not_awaited() + + async def test_global_and_storage_free_workloads_do_not_need_catalog(self): + api = _api() + api.network_volume_datacenters.side_effect = RuntimeError("catalog unavailable") + resolver = VolumeResolver(api) + assert await resolver.resolve_mounts({}) == [] + assert await resolver.resolve(GlobalVolume("models")) == {"id": "gv-new"} + api.network_volume_datacenters.assert_not_awaited() + class TestTaskVolume: def test_one_volume_per_backend(self): From 3e1a2c130fac786d2cc72c96776559b2b2961d52 Mon Sep 17 00:00:00 2001 From: zeke <40004347+KAJdev@users.noreply.github.com> Date: Wed, 7 Oct 2026 13:22:36 -0700 Subject: [PATCH 6/8] fix: enforce CON-1751 runtime deadlines with focused regressions --- CHANGELOG.md | 2 +- docs/cli/references/projects.md | 58 ++----- examples/apps/README.md | 13 +- runpod/apps/targets.py | 1 - runpod/apps/tasks.py | 10 +- tests/test_apps/test_api_client.py | 27 +-- tests/test_apps/test_deploy.py | 248 ++++------------------------ tests/test_apps/test_dispatch.py | 1 + tests/test_apps/test_placement.py | 52 ++---- tests/test_apps/test_tasks.py | 120 +++----------- tests/test_apps/test_volume.py | 143 ++-------------- tests/test_cli/test_init_command.py | 22 --- 12 files changed, 116 insertions(+), 581 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 9e95ec7d..b86ad246 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,7 +8,7 @@ ### Bug Fixes -* Wait for task cleanup after interruption, preserve cancellation and task errors, recover in-flight pod creation, and retry transient task deletion failures without cancelling detached work. +* Task calls complete bounded interruption cleanup and carry absolute runtime-enforced lifetime deadlines. 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. * 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. diff --git a/docs/cli/references/projects.md b/docs/cli/references/projects.md index 90fc05bc..e1d0fef7 100644 --- a/docs/cli/references/projects.md +++ b/docs/cli/references/projects.md @@ -14,46 +14,24 @@ Create a `.runpodignore` file in the root of your project to ignore files and fo ### Flash deployment artifacts -`rp flash deploy` and `rp flash deploy --build-only` apply git-style patterns to -source files, using `pathspec`. Precedence, from lowest to highest, is: +`rp flash deploy` and `rp flash deploy --build-only` apply Git-style ignore +patterns in this order, with later rules taking precedence: -1. Default local-file exclusions: virtual environments, Python caches, test - directories and test modules, `node_modules`, `.DS_Store`, and `*.tar.gz`. -2. `.gitignore` files within the project, with deeper files overriding parent - rules for their subtree. +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 ignore files are not read. Within each ignore file, the -last matching rule wins. `/` anchors a pattern to that ignore file's directory; -trailing `/` matches directories; `**`, comments, escaped characters and `!` -negation follow Git ignore syntax. Excluded directories are not traversed, so -re-include the parent before its files. For example, to deploy a fixture from -the otherwise excluded `tests` directory: - -```gitignore -!tests/ -tests/* -!tests/fixture.json -``` - -Some safeguards cannot be overridden by negation: `.git`, `.runpod`, `.flash`, -the root `env/` and `runpod_manifest.json` paths, and credential-like source -paths. Credential safeguards include `.env` and its variants, `*.env` and its -variants, `*.pem`, `*.key`, SSH private-key names, `.ssh`, `.aws`, `.azure`, -`.kube`, `.docker/config.json`, `.netrc`, `.npmrc`, `.pypirc`, `.git-credentials`, -`.boto`, `credentials`/`credentials.*`, `secrets`/`secrets.*`, and -`service-account*.json`/`service_account*.json`. Use worker environment variables -or a secret store for credentials instead of packaging them with source. - -These are filename safeguards, not secret scanning: credentials embedded in -ordinary source or differently named files can still be uploaded. Review the -build-only artifact before deployment. - -The generated manifest and vendored `env/` are added separately and cannot be -replaced by source files. Source ignore rules, including the PEM/key safeguards, -do not apply to vendored dependencies, so dependency CA bundles are preserved. -The output artifact itself and the dependency build directory are never copied -back into source. Source and dependency symlinks are omitted, including -project-contained links and linked ignore files; directory links are not -traversed. Hardlinked regular files are stored as independent regular members, -so the archive requires no link extraction support. +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 1894cd0c..b381f222 100644 --- a/examples/apps/README.md +++ b/examples/apps/README.md @@ -25,11 +25,8 @@ 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 gives cleanup up to 30 seconds to recover an in-flight create response -and delete the pod; synchronous calls allow up to 35 seconds for the coroutine to -acknowledge cleanup after ctrl-c. repeated interrupts do not interrupt deletion. -transient delete failures are retried up to three times. cleanup errors do not -replace an existing task error or cancellation, and failed deletion retains -`job.pod_id` for recovery. cleanup that exceeds the grace period continues only -while the client's event loop remains alive. check the pod in the console after a -cleanup warning; client process exit alone does not stop billing. +cancellation waits for bounded cleanup and retries transient deletion failures. +failed deletion retains `job.pod_id`; check the console after a cleanup warning. +each task carries an absolute one-hour deadline that the runtime checks even during +active work. polling and container restarts do not extend it. pod deletion still +requires a working control-plane API and pod-scoped credentials. diff --git a/runpod/apps/targets.py b/runpod/apps/targets.py index 763b39c3..b943424e 100644 --- a/runpod/apps/targets.py +++ b/runpod/apps/targets.py @@ -806,7 +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__( diff --git a/runpod/apps/tasks.py b/runpod/apps/tasks.py index 1f4124d5..8f0c10ea 100644 --- a/runpod/apps/tasks.py +++ b/runpod/apps/tasks.py @@ -7,8 +7,8 @@ 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. +the runtime watchdog requests pod deletion at the absolute task deadline, +including when the client has exited or the function is still running. """ import asyncio @@ -17,7 +17,7 @@ import secrets import sys import time -from datetime import datetime, timedelta, timezone +from datetime import timedelta from typing import Any, Dict, List, Optional import aiohttp @@ -31,7 +31,6 @@ TASK_PORT = 8080 -# requested server-side deadline for abandoned tasks DEFAULT_MAX_LIFETIME = timedelta(hours=1) READY_POLL_INTERVAL = 2.0 @@ -133,7 +132,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 @@ -141,6 +139,7 @@ def _pod_input(spec: ResourceSpec, token: str, task_name: str) -> Dict[str, Any] "RUNPOD_TASK_TOKEN": token, "RUNPOD_TASK_PORT": str(TASK_PORT), **render_env(spec.env), + "RUNPOD_TASK_DEADLINE": str(time.time() + DEFAULT_MAX_LIFETIME.total_seconds()), } from .images import image_for_spec, local_python_version @@ -150,7 +149,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, } diff --git a/tests/test_apps/test_api_client.py b/tests/test_apps/test_api_client.py index d5b17b4e..0303e84f 100644 --- a/tests/test_apps/test_api_client.py +++ b/tests/test_apps/test_api_client.py @@ -514,37 +514,14 @@ async def test_any_supported_tier_is_eligible(self): {"id": "US-IL-1", "networkVolumeTypes": []}, {"id": "EU-RO-1", "networkVolumeTypes": ["STANDARD"]}, {"id": "US-KS-2", "networkVolumeTypes": ["PREMIUM"]}, - {"id": "US-MO-2", "networkVolumeTypes": ["STANDARD", "PREMIUM"]}, ] } - with patch( - "runpod.apps.api.run_rest_request_async", - AsyncMock(return_value=catalog), - ) as rest: - supported = await AppsApiClient( - api_key="test-key" - ).network_volume_datacenters() - assert supported == {"EU-RO-1", "US-KS-2", "US-MO-2"} - rest.assert_awaited_once_with( - "GET", "/v2/catalog/datacenters", api_key="test-key" - ) - - @pytest.mark.parametrize( - "catalog", - [ - {}, - {"dataCenters": None}, - {"dataCenters": [{"id": "EU-RO-1"}]}, - {"dataCenters": [{"id": "EU-RO-1", "networkVolumeTypes": "STANDARD"}]}, - ], - ) - async def test_missing_capability_is_not_assumed_available(self, catalog): with patch( "runpod.apps.api.run_rest_request_async", AsyncMock(return_value=catalog), ): - with pytest.raises(QueryError): - await AppsApiClient().network_volume_datacenters() + supported = await AppsApiClient().network_volume_datacenters() + assert supported == {"EU-RO-1", "US-KS-2"} class TestVolumesRegistrySecrets: diff --git a/tests/test_apps/test_deploy.py b/tests/test_apps/test_deploy.py index 41152777..34eedda4 100644 --- a/tests/test_apps/test_deploy.py +++ b/tests/test_apps/test_deploy.py @@ -135,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) @@ -145,25 +146,31 @@ 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_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 @@ -174,232 +181,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_source_credentials_and_local_files_excluded_by_default(self, tmp_path): - excluded = { - ".env", - ".env.local", - "settings/dev.env", - "settings/dev.env.backup", - "keys/server.pem", - "keys/server.key", - ".ssh/id_rsa", - "keys/id_ed25519", - ".aws/credentials", - ".azure/token", - ".kube/config", - ".docker/config.json", - ".netrc", - ".npmrc", - ".pypirc", - ".git-credentials", - ".boto", - "credentials.json", - "secrets.yaml", - "service-account-prod.json", - "service_account.json", - "tests/test_app.py", - "test/unit.py", - "test_app.py", - "app_test.py", - ".venv/lib/local.py", - "venv/lib/local.py", - "env/lib/local.py", - ".pytest_cache/state", - "__pycache__/main.pyc", - } - for name in excluded | {"main.py", "package/module.py", "data/input.json"}: - path = tmp_path / name - path.parent.mkdir(parents=True, exist_ok=True) - path.write_text(name) - - with tarfile.open(package_project(tmp_path, {})) as tar: - assert set(tar.getnames()) == { - "main.py", - "package/module.py", - "data/input.json", - "runpod_manifest.json", - } - - def test_broad_negation_cannot_override_protected_paths(self, tmp_path): - protected = { - ".env.production", - "keys/private.pem", - "keys/private.key", - ".git/config", - ".runpod/cache", - ".flash/cache", - ".aws/credentials", - "env/local.py", - "runpod_manifest.json/forged.json", - } - for name in protected | {"tests/fixture.py", ".venv/local.py", "main.py"}: - path = tmp_path / name - path.parent.mkdir(parents=True, exist_ok=True) - path.write_text(name) - (tmp_path / ".runpodignore").write_text("!**\n") - - with tarfile.open(package_project(tmp_path, {"app": "generated"})) as tar: - assert set(tar.getnames()) == { - ".runpodignore", - "tests/fixture.py", - ".venv/local.py", - "main.py", - "runpod_manifest.json", - } - assert json.load(tar.extractfile("runpod_manifest.json")) == { - "app": "generated" - } - - def test_gitignore_scopes_and_runpod_precedence(self, tmp_path): - for name in [ - "main.py", - "root-only.txt", - "nested/root-only.txt", - "omit.log", - "nested/omit.log", + def test_project_ignore_precedence(self, tmp_path): + for name in ( + "root.log", "nested/keep.log", - "nested/hidden.txt", - "nested/override.txt", - "sibling/hidden.txt", - ]: + "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("/root-only.txt\n*.log\n") - (tmp_path / "nested/.gitignore").write_text( - "hidden.txt\noverride.txt\n!keep.log\n" - ) - (tmp_path / ".runpodignore").write_text( - "!/root-only.txt\n!nested/override.txt\nnested/keep.log\n" - ) + (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", - "main.py", - "root-only.txt", - "nested/root-only.txt", - "nested/override.txt", - "sibling/hidden.txt", - "runpod_manifest.json", - } - - def test_directory_patterns_require_parent_reinclusion(self, tmp_path): - for name in [ - "data/keep.txt", - "cache/keep.txt", - "nested/data/drop.txt", - "tests/fixture.py", - ]: - path = tmp_path / name - path.parent.mkdir(parents=True, exist_ok=True) - path.write_text(name) - (tmp_path / "ordinary").write_text("not a directory") - (tmp_path / ".gitignore").write_text("data/\ncache/\nordinary/\n") - (tmp_path / ".runpodignore").write_text( - "!data/keep.txt\n!cache/\ncache/*\n!cache/keep.txt\n!tests/\n" - ) - - with tarfile.open(package_project(tmp_path, {})) as tar: - assert set(tar.getnames()) == { - ".gitignore", - ".runpodignore", - "cache/keep.txt", - "ordinary", - "tests/fixture.py", - "runpod_manifest.json", - } - - def test_ignore_escaping_and_double_star(self, tmp_path): - for name in [ - "#notes", - "!notes", - "trailing ", - "assets/a/cache/x", - "assets/cache/y", - "assets/a/keep", - ]: - path = tmp_path / name - path.parent.mkdir(parents=True, exist_ok=True) - path.write_text(name) - (tmp_path / ".runpodignore").write_text( - "\\#notes\n\\!notes\ntrailing\\ \nassets/**/cache/\n" - ) - - with tarfile.open(package_project(tmp_path, {})) as tar: - assert set(tar.getnames()) == { - ".runpodignore", - "assets/a/keep", - "runpod_manifest.json", - } - - def test_generated_entries_and_dependency_certificates_preserved(self, tmp_path): - (tmp_path / "main.py").write_text("source") - (tmp_path / "runpod_manifest.json").write_text('{"app": "forged"}') - (tmp_path / "env").mkdir() - (tmp_path / "env/local.py").write_text("local environment") - env_dir = tmp_path / "built-env" - for name in ["certifi/cacert.pem", "library/public.key", "package.py"]: - path = env_dir / name - path.parent.mkdir(parents=True, exist_ok=True) - path.write_text(name) - (tmp_path / ".runpodignore").write_text("!**\n*.pem\n*.key\n") - output = tmp_path / "build.bundle" - - package_project(tmp_path, {"app": "generated"}, output, env_dir) - - with tarfile.open(output) as tar: - assert set(tar.getnames()) == { - ".runpodignore", - "main.py", - "env/certifi/cacert.pem", - "env/library/public.key", - "env/package.py", + "root.log", + "nested/keep.log", "runpod_manifest.json", } - assert tar.getnames().count("runpod_manifest.json") == 1 - assert json.load(tar.extractfile("runpod_manifest.json")) == { - "app": "generated" - } - assert ( - tar.extractfile("env/certifi/cacert.pem").read() - == b"certifi/cacert.pem" - ) - def test_links_are_not_followed_and_hardlinks_are_regular_members(self, tmp_path): + 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") - (project / "external.py").symlink_to(outside / "secret.txt") - (project / "external-dir").symlink_to(outside, target_is_directory=True) - (project / "internal.py").symlink_to("main.py") - (project / "broken.py").symlink_to("missing.py") - (project / "hardlink.py").hardlink_to(project / "main.py") - (outside / "ignore").write_text("*\n") - (project / ".gitignore").symlink_to(outside / "ignore") - (project / ".runpodignore").symlink_to(outside / "ignore") - (tmp_path / ".gitignore").write_text("*\n") env_dir = tmp_path / "built-env" env_dir.mkdir() - (env_dir / "package.py").write_text("dependency") - (env_dir / "external.py").symlink_to(outside / "secret.txt") - (env_dir / "external-dir").symlink_to(outside, target_is_directory=True) + 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", - "hardlink.py", - "env/package.py", - "runpod_manifest.json", - } - assert all(member.isfile() for member in tar.getmembers()) - assert tar.extractfile("hardlink.py").read() == b"source" + assert set(tar.getnames()) == {"main.py", "runpod_manifest.json"} def _stub_build(tmp_path): diff --git a/tests/test_apps/test_dispatch.py b/tests/test_apps/test_dispatch.py index 83bed7ed..6a2543a8 100644 --- a/tests/test_apps/test_dispatch.py +++ b/tests/test_apps/test_dispatch.py @@ -481,6 +481,7 @@ def interrupt(): 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 de3d641a..99aef01d 100644 --- a/tests/test_apps/test_placement.py +++ b/tests/test_apps/test_placement.py @@ -141,7 +141,7 @@ def test_existing_dc_is_hard_constraint(self): stock, volume_name="models", existing_dc="EU-RO-1", - volume_datacenters={dc.value for dc in DataCenter.all()}, + volume_datacenters=set(), ) assert dc == "EU-RO-1" @@ -214,31 +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", - volume_datacenters={dc.value for dc in DataCenter.all()}, - ) - == "EU-RO-1" + supported = {"EU-RO-1", "US-KS-2"} + dc = solve_placement( + [single], stock, volume_name="one", volume_datacenters=supported ) - assert ( - solve_placement( - [single, pair], - stock, - volume_name="shared", - volume_datacenters={dc.value for dc in DataCenter.all()}, - ) - == "US-KS-2" + 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", - volume_datacenters={dc.value for dc in DataCenter.all()}, + volume_datacenters=supported, ) async def test_cpu_task_and_endpoint_require_shared_product_stock(self): @@ -254,29 +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", - volume_datacenters={dc.value for dc in DataCenter.all()}, - ) - == "EU-RO-1" + supported = {"EU-RO-1", "US-KS-2"} + dc = solve_placement( + [endpoint], stock, volume_name="one", volume_datacenters=supported ) - assert ( - solve_placement( - [endpoint, task], - stock, - volume_name="shared", - volume_datacenters={dc.value for dc in DataCenter.all()}, - ) - == "US-KS-2" + 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", - volume_datacenters={dc.value for dc in DataCenter.all()}, + volume_datacenters=supported, ) diff --git a/tests/test_apps/test_tasks.py b/tests/test_apps/test_tasks.py index cadd756d..a971c602 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 @@ -592,13 +591,12 @@ 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 - spec = ResourceSpec(kind=ResourceKind.TASK, name="t", cpu=["cpu3c-1-2"]) - creating = asyncio.Event() - release_create = asyncio.Event() - deleting = asyncio.Event() - release_delete = asyncio.Event() + 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() @@ -609,121 +607,41 @@ async def deploy(*args, **kwargs): async def delete(pod_id): deleting.set() await release_delete.wait() + failure = next(failures) + if failure: + raise failure pods.remove(pod_id) - api = MagicMock() - api.deploy_task_pod = deploy - api.terminate_pod = delete - execution = TaskExecution(spec, api=api) + 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("original interruption") - await asyncio.sleep(0) - task.cancel("second interruption") + task.cancel() release_create.set() await deleting.wait() - task.cancel("third interruption") + task.cancel() release_delete.set() with pytest.raises(asyncio.CancelledError): await task assert not pods - assert execution.pod_id is None - @pytest.mark.parametrize("method", ["invoke", "submit", "wait"]) - @pytest.mark.parametrize("cancelled", [False, True]) - async def test_cleanup_failure_preserves_original(self, method, cancelled): + 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"]) - failure = ( - asyncio.CancelledError("cancelled") - if cancelled - else ValueError("work failed") + api = MagicMock( + terminate_pod=AsyncMock(side_effect=RuntimeError("delete failed")) ) - api = MagicMock() - api.terminate_pod = AsyncMock(side_effect=RuntimeError("delete failed")) execution = TaskExecution(spec, api=api) execution.pod_id = "pod-recoverable" - execution.start = AsyncMock(side_effect=failure) execution.poll_result = AsyncMock(side_effect=failure) - with patch("runpod.apps.tasks.TaskExecution", return_value=execution): - with pytest.raises(type(failure)) as caught: - if method == "wait": - await TaskJob(execution).wait() - else: - await getattr(PodTarget(spec, lambda: None), method)({}) + with pytest.raises(ValueError) as caught: + await TaskJob(execution).wait() assert caught.value is failure assert execution.pod_id == "pod-recoverable" - - async def test_cleanup_deadline_does_not_abandon_deletion(self, monkeypatch): - from runpod.apps import tasks - - monkeypatch.setattr(tasks, "CLEANUP_TIMEOUT", 0.02) - entered = asyncio.Event() - release = asyncio.Event() - deleted = asyncio.Event() - - async def delete(): - entered.set() - await release.wait() - deleted.set() - - async def operation(): - try: - await asyncio.Future() - finally: - await tasks.finish_cleanup(delete()) - - task = asyncio.create_task(operation()) - await asyncio.sleep(0) - task.cancel("original") - await entered.wait() - task.cancel("again") - with pytest.raises(asyncio.CancelledError): - await asyncio.wait_for(task, 0.5) - assert not deleted.is_set() - release.set() - await asyncio.wait_for(deleted.wait(), 0.5) - - @pytest.mark.parametrize("status", [429, 503, None]) - async def test_transient_delete_recovers(self, status): - from runpod.apps.tasks import TaskExecution - from runpod.error import QueryError - - spec = ResourceSpec(kind=ResourceKind.TASK, name="t", cpu=["cpu3c-1-2"]) - pods = {"pod-retry"} - attempts = 0 - - async def delete(pod_id): - nonlocal attempts - attempts += 1 - if attempts == 1: - if status is None: - raise OSError("connection reset") - raise QueryError("retry later", status_code=status) - pods.remove(pod_id) - - api = MagicMock() - api.terminate_pod = delete - execution = TaskExecution(spec, api=api) - execution.pod_id = "pod-retry" - await execution.terminate() - assert not pods - assert execution.pod_id is None - - async def test_retry_exhaustion_keeps_recoverable_id(self): - from runpod.apps import tasks - from runpod.error import QueryError - - spec = ResourceSpec(kind=ResourceKind.TASK, name="t", cpu=["cpu3c-1-2"]) - api = MagicMock() - api.terminate_pod = AsyncMock(side_effect=QueryError("busy", status_code=503)) - execution = tasks.TaskExecution(spec, api=api) - execution.pod_id = "pod-retry" - with pytest.raises(QueryError): - await execution.terminate() - assert api.terminate_pod.await_count == tasks.DELETE_ATTEMPTS - assert execution.pod_id == "pod-retry" diff --git a/tests/test_apps/test_volume.py b/tests/test_apps/test_volume.py index fdf5fd48..484ab31d 100644 --- a/tests/test_apps/test_volume.py +++ b/tests/test_apps/test_volume.py @@ -151,9 +151,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() @@ -170,14 +171,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): @@ -265,42 +269,17 @@ def test_creation_without_app_requires_explicit_placement(self): ) assert resolved == {"id": "nv-new", "dataCenterId": "EU-RO-1"} - async def test_storage_capability_beats_higher_hardware_stock(self): + async def test_unsupported_explicit_pin_never_relocates(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" - ) - resolved = await VolumeResolver(api).resolve( - NetworkVolume("models"), [_spec(cpu="cpu3c-2-4")] - ) - assert resolved["dataCenterId"] == "EU-RO-1" - assert api.create_network_volume.call_args.kwargs["data_center_id"] == "EU-RO-1" - - async def test_reported_datacenters_remain_eligible_when_catalog_supports_them( - self, - ): - api = _api() - api.network_volume_datacenters.return_value = {"US-IL-1"} - resolved = await VolumeResolver(api).resolve( - NetworkVolume("models"), [_spec(cpu="cpu3c-2-4")] - ) - assert resolved["dataCenterId"] == "US-IL-1" - - @pytest.mark.parametrize("specs", [[], [_spec(cpu="cpu3c-2-4")]]) - async def test_unsupported_explicit_pin_never_relocates(self, specs): - api = _api() - api.network_volume_datacenters.return_value = {"EU-RO-1"} - with pytest.raises(VolumeError, match="does not support network volumes"): + with pytest.raises(VolumeError): await VolumeResolver(api).resolve( - NetworkVolume("models", datacenter="US-IL-1"), specs + NetworkVolume("models", datacenter="US-IL-1"), [_spec()] ) - api.create_network_volume.assert_not_awaited() - @pytest.mark.parametrize("supported", [set(), {"EU-RO-1"}]) - async def test_empty_storage_hardware_intersection_never_creates(self, supported): + async def test_disjoint_storage_and_hardware_cannot_create(self): api = _api() - api.network_volume_datacenters.return_value = supported + 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" ) @@ -308,100 +287,6 @@ async def test_empty_storage_hardware_intersection_never_creates(self, supported await VolumeResolver(api).resolve( NetworkVolume("models"), [_spec(cpu="cpu3c-2-4")] ) - api.create_network_volume.assert_not_awaited() - - @pytest.mark.parametrize("reference", ["models", "nv-1"]) - async def test_existing_volume_ignores_new_creation_capability(self, reference): - api = _api(volumes=[{"id": "nv-1", "name": "models", "dataCenter": "US-IL-1"}]) - api.network_volume_datacenters.side_effect = RuntimeError("catalog unavailable") - resolved = await VolumeResolver(api).resolve( - NetworkVolume(reference, datacenter="EU-RO-1"), - [_spec(cpu="cpu3c-2-4")], - ) - assert resolved == {"id": "nv-1", "dataCenterId": "US-IL-1"} - api.network_volume_datacenters.assert_not_awaited() - api.create_network_volume.assert_not_awaited() - - async def test_catalog_failure_does_not_create(self): - api = _api() - api.network_volume_datacenters.side_effect = RuntimeError("catalog unavailable") - with pytest.raises(VolumeError, match="datacenter catalog lookup failed"): - await VolumeResolver(api).resolve( - NetworkVolume("models", datacenter="EU-RO-1") - ) - api.create_network_volume.assert_not_awaited() - - async def test_creation_uses_one_catalog_snapshot_per_run(self): - api = _api() - api.network_volume_datacenters.side_effect = [{"EU-RO-1"}, {"US-IL-1"}] - resolver = VolumeResolver(api) - first = await resolver.resolve(NetworkVolume("first"), [_spec()]) - second = await resolver.resolve(NetworkVolume("second"), [_spec()]) - assert first["dataCenterId"] == second["dataCenterId"] == "EU-RO-1" - assert api.network_volume_datacenters.await_count == 1 - - async def test_all_shared_consumers_constrain_creation(self): - api = _api() - api.network_volume_datacenters.return_value = {"EU-RO-1", "US-IL-1"} - volume = NetworkVolume("models") - producer = ResourceSpec( - kind=ResourceKind.TASK, - name="producer", - cpu="cpu3c-2-4", - mounts={"/models": volume}, - ) - consumer = ResourceSpec( - kind=ResourceKind.QUEUE, - name="consumer", - gpu="4090", - gpu_count=2, - datacenter="EU-RO-1", - mounts={"/runpod-volume": NetworkVolume("models")}, - ) - unrelated = ResourceSpec( - kind=ResourceKind.TASK, name="unrelated", datacenter="US-IL-1" - ) - api.cpu_stock_status.side_effect = lambda instance, dc, *, pods=False: ( - "HIGH" if dc == "US-IL-1" else "LOW" - ) - api.gpu_stock_status.side_effect = lambda gpu, dc, gpu_count=1, pods=False: ( - "LOW" if dc == "EU-RO-1" and gpu_count == 2 and not pods else "NONE" - ) - bindings = await VolumeResolver(api).resolve_mounts( - producer.mounts, [producer, consumer, unrelated] - ) - assert bindings[0]["dataCenterId"] == "EU-RO-1" - assert api.create_network_volume.call_args.kwargs["data_center_id"] == "EU-RO-1" - - async def test_supported_volume_pin_conflicting_with_consumer_never_relocates(self): - api = _api() - spec = ResourceSpec(kind=ResourceKind.TASK, name="worker", datacenter="EU-RO-1") - with pytest.raises(PlacementError): - await VolumeResolver(api).resolve( - NetworkVolume("models", datacenter="US-IL-1"), [spec] - ) - api.create_network_volume.assert_not_awaited() - - async def test_shared_volume_reference_pin_is_not_ignored(self): - api = _api() - api.network_volume_datacenters.return_value = {"EU-RO-1"} - volume = NetworkVolume("models") - sibling = ResourceSpec( - kind=ResourceKind.TASK, - name="sibling", - mounts={"/models": NetworkVolume("models", datacenter="US-IL-1")}, - ) - with pytest.raises(VolumeError): - await VolumeResolver(api).resolve(volume, [sibling]) - api.create_network_volume.assert_not_awaited() - - async def test_global_and_storage_free_workloads_do_not_need_catalog(self): - api = _api() - api.network_volume_datacenters.side_effect = RuntimeError("catalog unavailable") - resolver = VolumeResolver(api) - assert await resolver.resolve_mounts({}) == [] - assert await resolver.resolve(GlobalVolume("models")) == {"id": "gv-new"} - api.network_volume_datacenters.assert_not_awaited() class TestTaskVolume: diff --git a/tests/test_cli/test_init_command.py b/tests/test_cli/test_init_command.py index fe0b5915..42be0898 100644 --- a/tests/test_cli/test_init_command.py +++ b/tests/test_cli/test_init_command.py @@ -1,10 +1,7 @@ """rp flash init: project scaffolding.""" -import tarfile - from click.testing import CliRunner -from runpod.apps.deploy import package_project from runpod.apps.init import create_project, detect_conflicts from runpod.rp_cli.main import cli @@ -32,25 +29,6 @@ def test_creates_directory(self, tmp_path): create_project(target, "new-project") assert (target / "main.py").exists() - def test_scaffold_packages_source_without_local_secrets(self, tmp_path): - create_project(tmp_path, "my-app") - (tmp_path / ".env.production").write_text("TOKEN=private") - (tmp_path / "private.pem").write_text("private") - (tmp_path / "tests").mkdir() - (tmp_path / "tests/test_app.py").write_text("local test") - (tmp_path / ".gitignore").write_text("local-data/\n") - (tmp_path / "local-data").mkdir() - (tmp_path / "local-data/input.json").write_text("{}") - - with tarfile.open(package_project(tmp_path, {})) as tar: - assert set(tar.getnames()) == { - "main.py", - "requirements.txt", - ".runpodignore", - ".gitignore", - "runpod_manifest.json", - } - class TestDetectConflicts: def test_empty_dir_no_conflicts(self, tmp_path): From e716258164122aa1f6f13f3ac174269057aeab56 Mon Sep 17 00:00:00 2001 From: zeke <40004347+KAJdev@users.noreply.github.com> Date: Wed, 7 Oct 2026 22:58:21 -0700 Subject: [PATCH 7/8] fix: reject global volumes on cpu compute --- CHANGELOG.md | 1 + examples/apps/README.md | 3 ++- runpod/apps/volume.py | 6 ++---- tests/test_apps/test_volume.py | 32 +++++++++++++++++++------------- 4 files changed, 24 insertions(+), 18 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index e596f07f..be9484ad 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -10,6 +10,7 @@ * 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. diff --git a/examples/apps/README.md b/examples/apps/README.md index 8c9e0440..c7ca9abb 100644 --- a/examples/apps/README.md +++ b/examples/apps/README.md @@ -48,7 +48,8 @@ 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 storage requires a GPU endpoint. +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 diff --git a/runpod/apps/volume.py b/runpod/apps/volume.py index f3a2523e..943c80b5 100644 --- a/runpod/apps/volume.py +++ b/runpod/apps/volume.py @@ -331,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_volume.py b/tests/test_apps/test_volume.py index 484ab31d..70c255db 100644 --- a/tests/test_apps/test_volume.py +++ b/tests/test_apps/test_volume.py @@ -143,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): @@ -311,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) @@ -336,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 From 15999425af2b51ba5d92f83d17dc2b3ce96d06eb Mon Sep 17 00:00:00 2001 From: zeke <40004347+KAJdev@users.noreply.github.com> Date: Thu, 8 Oct 2026 17:14:33 -0700 Subject: [PATCH 8/8] fix: enforce task timeouts remotely and refine artifact exclusions --- runpod/apps/deploy.py | 6 ++-- runpod/apps/tasks.py | 36 ++++++++++++++++---- tests/test_apps/test_deploy.py | 30 +++++++++++++++++ tests/test_apps/test_tasks.py | 60 ++++++++++++++++++++++++++++++++++ 4 files changed, 122 insertions(+), 10 deletions(-) diff --git a/runpod/apps/deploy.py b/runpod/apps/deploy.py index 7d2741c7..a0ef47ca 100644 --- a/runpod/apps/deploy.py +++ b/runpod/apps/deploy.py @@ -84,10 +84,6 @@ ".pypirc", ".git-credentials", ".boto", - "credentials", - "credentials.*", - "secrets", - "secrets.*", "service-account*.json", "service_account*.json", ] @@ -225,6 +221,7 @@ def package_project( 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 @@ -238,6 +235,7 @@ def package_project( or PROTECTED_IGNORES.match_file(rel) or _is_ignored(rel, patterns) ): + log.warning("excluded source path %r from deployment artifact", rel) continue tar.add(path, arcname=rel, recursive=False) diff --git a/runpod/apps/tasks.py b/runpod/apps/tasks.py index 6355f19f..e6a6419b 100644 --- a/runpod/apps/tasks.py +++ b/runpod/apps/tasks.py @@ -6,13 +6,11 @@ 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 - -active functions have no runtime lifetime limit. the watchdog only reclaims -idle pods without active background or inline work. """ import asyncio import logging +import math import os as _os import secrets import sys @@ -287,7 +285,7 @@ async def execute( (long jobs would 524), so completion always goes through the background slot + short /result polls. """ - await self.submit(request) + await self.submit({**request, "timeout": timeout}) deadline = time.monotonic() + timeout if timeout is not None else None while True: if deadline is not None and time.monotonic() >= deadline: @@ -307,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): @@ -321,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]]: @@ -426,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( diff --git a/tests/test_apps/test_deploy.py b/tests/test_apps/test_deploy.py index 34eedda4..5445c5b5 100644 --- a/tests/test_apps/test_deploy.py +++ b/tests/test_apps/test_deploy.py @@ -164,6 +164,36 @@ def test_ignores_and_non_overridable_credentials(self, tmp_path): "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" diff --git a/tests/test_apps/test_tasks.py b/tests/test_apps/test_tasks.py index 67921683..78b76115 100644 --- a/tests/test_apps/test_tasks.py +++ b/tests/test_apps/test_tasks.py @@ -477,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 @@ -496,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 @@ -515,6 +556,23 @@ async def test_wait_returns_result_and_terminates(self): assert result == 9 execution.terminate.assert_awaited_once() + async def test_wait_sends_timeout_to_runtime(self): + job, execution = self._job() + 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 @@ -527,11 +585,13 @@ 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",