diff --git a/docs/environments.md b/docs/environments.md index f02987c9a4..a9bb10edad 100644 --- a/docs/environments.md +++ b/docs/environments.md @@ -789,6 +789,7 @@ Supported third-party environment integrations include: - **`ReasoningGymEnv`** — wraps [reasoning-gym](https://github.com/open-thought/reasoning-gym) procedural datasets - **`BrowserEnv`** — unified browser automation via [Browserbase](https://browserbase.com) with DOM and CUA modes - **`OpenEnvEnv`** — wraps OpenEnv gym and MCP contracts using Prime Sandboxes with prebuilt images referenced from `.build.json` +- **`NemoGymEnv`** — wraps [NVIDIA NeMo Gym](https://github.com/NVIDIA-NeMo/Gym) environments. These require additional dependencies installed via extras (e.g., `uv add 'verifiers[ta]'` for TextArena, `uv add 'verifiers[browser]'` for BrowserEnv, `uv add 'verifiers[openenv]'` for OpenEnvEnv). For OpenEnv environments, build the bundled project image with `prime env build ` before evaluation or training. diff --git a/environments/README.md b/environments/README.md index 61261d7bc6..619b695696 100644 --- a/environments/README.md +++ b/environments/README.md @@ -49,6 +49,15 @@ This folder contains installable example environments that showcase common usage - **opencode_harbor**: Runs the OpenCode CLI agent on Harbor tasks with API interception via Prime Tunnel. - **terminus_harbor**: Runs the Terminus agent on Harbor tasks with API interception via Prime Tunnel. +- **NeMo Gym**: wraps NVIDIA NeMo Gym environments. A few examples to get started: + - **nemo-gym-workplace-assistant**: Multi-step tool use workplace assistant + - **nemo-gym-reasoning-gym**: Reasoning Gym tasks with a simple agent + - **nemo-gym-reasoning-gym-reflection**: Reasoning Gym with a LangGraph reflection agent + - **nemo-gym-structured-outputs**: Structured output tasks + - **nemo-gym-code-gen**: Code generation tasks + - **nemo-gym-mcqa**: Science multiple choice question answering + - **nemo-gym-arc-agi**: ARC-AGI 1 and 2 + ### Composition - **EnvGroup** - **math_group**: Groups two `SingleTurnEnv` tasks (GSM8K + Math) into one environment with shared interface. @@ -73,6 +82,7 @@ This folder contains installable example environments that showcase common usage - **GymEnv integration**: `gem_wordle` - **OpenEnv integration (gym + MCP)**: `openenv_textarena`, `openenv_echo` - **CLI agent sandboxes**: `opencode_harbor`, `terminus_harbor` +- **NeMo Gym integration**: `nemo-gym-reasoning-gym`, `nemo-gym-workplace-assistant` - **MCP integration**: `mcp_search_env` - **RLM (recursive LLM)**: `rlm_secrets` - **Environment and rubric composition**: `math_group`, `math_python`, `wiki_search` diff --git a/environments/nemo_gym/nemo_gym_arc_agi/nemo_gym_arc_agi.py b/environments/nemo_gym/nemo_gym_arc_agi/nemo_gym_arc_agi.py new file mode 100644 index 0000000000..d501d750b0 --- /dev/null +++ b/environments/nemo_gym/nemo_gym_arc_agi/nemo_gym_arc_agi.py @@ -0,0 +1,20 @@ +from typing import Any + +import verifiers as vf +from verifiers.envs.integrations.nemo_gym import ( + NemoGymEnv, + _build_dataset, + _resolve_gym_config, +) + + +def load_environment( + dataset_path: str | None = None, + **kwargs: Any, +) -> vf.Environment: + dataset, _ = _build_dataset("arc_agi", "example", dataset_path=dataset_path) + return NemoGymEnv( + gym_configs=[_resolve_gym_config("arc_agi")], + dataset=dataset, + **kwargs, + ) diff --git a/environments/nemo_gym/nemo_gym_arc_agi/pyproject.toml b/environments/nemo_gym/nemo_gym_arc_agi/pyproject.toml new file mode 100644 index 0000000000..a2bda46f92 --- /dev/null +++ b/environments/nemo_gym/nemo_gym_arc_agi/pyproject.toml @@ -0,0 +1,21 @@ +[project] +name = "nemo-gym-arc-agi" +description = "NeMo Gym ARC AGI 1 and 2 environment" +tags = ["nemo-gym", "knowledge", "single-turn"] +version = "0.1.0" +requires-python = ">=3.12" +dependencies = [ + "verifiers>=0.1.11.dev1", + "nemo-gym>=0.2.0", +] + +[build-system] +requires = ["hatchling"] +build-backend = "hatchling.build" + +[tool.hatch.build] +include = ["nemo_gym_arc_agi.py", "pyproject.toml"] + +[tool.verifiers.eval] +num_examples = 5 +rollouts_per_example = 1 diff --git a/environments/nemo_gym/nemo_gym_code_gen/nemo_gym_code_gen.py b/environments/nemo_gym/nemo_gym_code_gen/nemo_gym_code_gen.py new file mode 100644 index 0000000000..550079e373 --- /dev/null +++ b/environments/nemo_gym/nemo_gym_code_gen/nemo_gym_code_gen.py @@ -0,0 +1,20 @@ +from typing import Any + +import verifiers as vf +from verifiers.envs.integrations.nemo_gym import ( + NemoGymEnv, + _build_dataset, + _resolve_gym_config, +) + + +def load_environment( + dataset_path: str | None = None, + **kwargs: Any, +) -> vf.Environment: + dataset, _ = _build_dataset("code_gen", "example", dataset_path=dataset_path) + return NemoGymEnv( + gym_configs=[_resolve_gym_config("code_gen")], + dataset=dataset, + **kwargs, + ) diff --git a/environments/nemo_gym/nemo_gym_code_gen/pyproject.toml b/environments/nemo_gym/nemo_gym_code_gen/pyproject.toml new file mode 100644 index 0000000000..0da06d2b5f --- /dev/null +++ b/environments/nemo_gym/nemo_gym_code_gen/pyproject.toml @@ -0,0 +1,21 @@ +[project] +name = "nemo-gym-code-gen" +description = "NeMo Gym code generation environment" +tags = ["nemo-gym"] +version = "0.1.0" +requires-python = ">=3.12" +dependencies = [ + "verifiers>=0.1.11.dev1", + "nemo-gym>=0.2.0", +] + +[build-system] +requires = ["hatchling"] +build-backend = "hatchling.build" + +[tool.hatch.build] +include = ["nemo_gym_code_gen.py", "pyproject.toml"] + +[tool.verifiers.eval] +num_examples = 5 +rollouts_per_example = 1 diff --git a/environments/nemo_gym/nemo_gym_mcqa/nemo_gym_mcqa.py b/environments/nemo_gym/nemo_gym_mcqa/nemo_gym_mcqa.py new file mode 100644 index 0000000000..dc9db405c9 --- /dev/null +++ b/environments/nemo_gym/nemo_gym_mcqa/nemo_gym_mcqa.py @@ -0,0 +1,20 @@ +from typing import Any + +import verifiers as vf +from verifiers.envs.integrations.nemo_gym import ( + NemoGymEnv, + _build_dataset, + _resolve_gym_config, +) + + +def load_environment( + dataset_path: str | None = None, + **kwargs: Any, +) -> vf.Environment: + dataset, _ = _build_dataset("mcqa", "example", dataset_path=dataset_path) + return NemoGymEnv( + gym_configs=[_resolve_gym_config("mcqa")], + dataset=dataset, + **kwargs, + ) diff --git a/environments/nemo_gym/nemo_gym_mcqa/pyproject.toml b/environments/nemo_gym/nemo_gym_mcqa/pyproject.toml new file mode 100644 index 0000000000..2faa781b4c --- /dev/null +++ b/environments/nemo_gym/nemo_gym_mcqa/pyproject.toml @@ -0,0 +1,21 @@ +[project] +name = "nemo-gym-mcqa" +description = "NeMo Gym mcqa environment" +tags = ["nemo-gym"] +version = "0.1.0" +requires-python = ">=3.12" +dependencies = [ + "verifiers>=0.1.11.dev1", + "nemo-gym>=0.2.0", +] + +[build-system] +requires = ["hatchling"] +build-backend = "hatchling.build" + +[tool.hatch.build] +include = ["nemo_gym_mcqa.py", "pyproject.toml"] + +[tool.verifiers.eval] +num_examples = 5 +rollouts_per_example = 1 diff --git a/environments/nemo_gym/nemo_gym_reasoning_gym/nemo_gym_reasoning_gym.py b/environments/nemo_gym/nemo_gym_reasoning_gym/nemo_gym_reasoning_gym.py new file mode 100644 index 0000000000..c9ca71a80e --- /dev/null +++ b/environments/nemo_gym/nemo_gym_reasoning_gym/nemo_gym_reasoning_gym.py @@ -0,0 +1,20 @@ +from typing import Any + +import verifiers as vf +from verifiers.envs.integrations.nemo_gym import ( + NemoGymEnv, + _build_dataset, + _resolve_gym_config, +) + + +def load_environment( + dataset_path: str | None = None, + **kwargs: Any, +) -> vf.Environment: + dataset, _ = _build_dataset("reasoning_gym", "example", dataset_path=dataset_path) + return NemoGymEnv( + gym_configs=[_resolve_gym_config("reasoning_gym")], + dataset=dataset, + **kwargs, + ) diff --git a/environments/nemo_gym/nemo_gym_reasoning_gym/pyproject.toml b/environments/nemo_gym/nemo_gym_reasoning_gym/pyproject.toml new file mode 100644 index 0000000000..5ea6cbffdf --- /dev/null +++ b/environments/nemo_gym/nemo_gym_reasoning_gym/pyproject.toml @@ -0,0 +1,21 @@ +[project] +name = "nemo-gym-reasoning-gym" +description = "NeMo Gym Reasoning Gym simple agent environment" +tags = ["nemo-gym"] +version = "0.1.0" +requires-python = ">=3.12" +dependencies = [ + "verifiers>=0.1.11.dev1", + "nemo-gym>=0.2.0", +] + +[build-system] +requires = ["hatchling"] +build-backend = "hatchling.build" + +[tool.hatch.build] +include = ["nemo_gym_reasoning_gym.py", "pyproject.toml"] + +[tool.verifiers.eval] +num_examples = 5 +rollouts_per_example = 1 diff --git a/environments/nemo_gym/nemo_gym_reasoning_gym_parallel_thinking/nemo_gym_reasoning_gym_parallel_thinking.py b/environments/nemo_gym/nemo_gym_reasoning_gym_parallel_thinking/nemo_gym_reasoning_gym_parallel_thinking.py new file mode 100644 index 0000000000..79bc290a8f --- /dev/null +++ b/environments/nemo_gym/nemo_gym_reasoning_gym_parallel_thinking/nemo_gym_reasoning_gym_parallel_thinking.py @@ -0,0 +1,20 @@ +from typing import Any + +import verifiers as vf +from verifiers.envs.integrations.nemo_gym import ( + NemoGymEnv, + _build_dataset, + _resolve_gym_config, +) + + +def load_environment( + dataset_path: str | None = None, + **kwargs: Any, +) -> vf.Environment: + dataset, _ = _build_dataset("reasoning_gym", "example", dataset_path=dataset_path) + return NemoGymEnv( + gym_configs=[_resolve_gym_config("reasoning_gym", "parallel_thinking_agent")], + dataset=dataset, + **kwargs, + ) diff --git a/environments/nemo_gym/nemo_gym_reasoning_gym_parallel_thinking/pyproject.toml b/environments/nemo_gym/nemo_gym_reasoning_gym_parallel_thinking/pyproject.toml new file mode 100644 index 0000000000..79722b460f --- /dev/null +++ b/environments/nemo_gym/nemo_gym_reasoning_gym_parallel_thinking/pyproject.toml @@ -0,0 +1,21 @@ +[project] +name = "nemo-gym-reasoning-gym-parallel-thinking" +description = "NeMo Gym Reasoning Gym environment with LangGraph parallel thinking agent" +tags = ["nemo-gym"] +version = "0.1.0" +requires-python = ">=3.12" +dependencies = [ + "verifiers>=0.1.11.dev1", + "nemo-gym>=0.2.0", +] + +[build-system] +requires = ["hatchling"] +build-backend = "hatchling.build" + +[tool.hatch.build] +include = ["nemo_gym_reasoning_gym_parallel_thinking.py", "pyproject.toml"] + +[tool.verifiers.eval] +num_examples = 5 +rollouts_per_example = 1 diff --git a/environments/nemo_gym/nemo_gym_reasoning_gym_reflection/nemo_gym_reasoning_gym_reflection.py b/environments/nemo_gym/nemo_gym_reasoning_gym_reflection/nemo_gym_reasoning_gym_reflection.py new file mode 100644 index 0000000000..3be8252706 --- /dev/null +++ b/environments/nemo_gym/nemo_gym_reasoning_gym_reflection/nemo_gym_reasoning_gym_reflection.py @@ -0,0 +1,20 @@ +from typing import Any + +import verifiers as vf +from verifiers.envs.integrations.nemo_gym import ( + NemoGymEnv, + _build_dataset, + _resolve_gym_config, +) + + +def load_environment( + dataset_path: str | None = None, + **kwargs: Any, +) -> vf.Environment: + dataset, _ = _build_dataset("reasoning_gym", "example", dataset_path=dataset_path) + return NemoGymEnv( + gym_configs=[_resolve_gym_config("reasoning_gym", "reflection_agent")], + dataset=dataset, + **kwargs, + ) diff --git a/environments/nemo_gym/nemo_gym_reasoning_gym_reflection/pyproject.toml b/environments/nemo_gym/nemo_gym_reasoning_gym_reflection/pyproject.toml new file mode 100644 index 0000000000..66941dbf2b --- /dev/null +++ b/environments/nemo_gym/nemo_gym_reasoning_gym_reflection/pyproject.toml @@ -0,0 +1,21 @@ +[project] +name = "nemo-gym-reasoning-gym-reflection" +description = "NeMo Gym Reasoning Gym environment with LangGraph reflection agent" +tags = ["nemo-gym"] +version = "0.1.0" +requires-python = ">=3.12" +dependencies = [ + "verifiers>=0.1.11.dev1", + "nemo-gym>=0.2.0", +] + +[build-system] +requires = ["hatchling"] +build-backend = "hatchling.build" + +[tool.hatch.build] +include = ["nemo_gym_reasoning_gym_reflection.py", "pyproject.toml"] + +[tool.verifiers.eval] +num_examples = 5 +rollouts_per_example = 1 diff --git a/environments/nemo_gym/nemo_gym_structured_outputs/nemo_gym_structured_outputs.py b/environments/nemo_gym/nemo_gym_structured_outputs/nemo_gym_structured_outputs.py new file mode 100644 index 0000000000..e6b61bb207 --- /dev/null +++ b/environments/nemo_gym/nemo_gym_structured_outputs/nemo_gym_structured_outputs.py @@ -0,0 +1,24 @@ +from typing import Any + +import verifiers as vf +from verifiers.envs.integrations.nemo_gym import ( + NemoGymEnv, + _build_dataset, + _resolve_gym_config, +) + + +def load_environment( + dataset_path: str | None = None, + **kwargs: Any, +) -> vf.Environment: + dataset, _ = _build_dataset( + "structured_outputs", "example", dataset_path=dataset_path + ) + return NemoGymEnv( + gym_configs=[ + _resolve_gym_config("structured_outputs", "structured_outputs_json") + ], + dataset=dataset, + **kwargs, + ) diff --git a/environments/nemo_gym/nemo_gym_structured_outputs/pyproject.toml b/environments/nemo_gym/nemo_gym_structured_outputs/pyproject.toml new file mode 100644 index 0000000000..c18db17fca --- /dev/null +++ b/environments/nemo_gym/nemo_gym_structured_outputs/pyproject.toml @@ -0,0 +1,21 @@ +[project] +name = "nemo-gym-structured-outputs" +description = "NeMo Gym structured outputs environment" +tags = ["nemo-gym"] +version = "0.1.0" +requires-python = ">=3.12" +dependencies = [ + "verifiers>=0.1.11.dev1", + "nemo-gym>=0.2.0", +] + +[build-system] +requires = ["hatchling"] +build-backend = "hatchling.build" + +[tool.hatch.build] +include = ["nemo_gym_structured_outputs.py", "pyproject.toml"] + +[tool.verifiers.eval] +num_examples = 5 +rollouts_per_example = 1 diff --git a/environments/nemo_gym/nemo_gym_workplace_assistant/nemo_gym_workplace_assistant.py b/environments/nemo_gym/nemo_gym_workplace_assistant/nemo_gym_workplace_assistant.py new file mode 100644 index 0000000000..8cd4732c0f --- /dev/null +++ b/environments/nemo_gym/nemo_gym_workplace_assistant/nemo_gym_workplace_assistant.py @@ -0,0 +1,22 @@ +from typing import Any + +import verifiers as vf +from verifiers.envs.integrations.nemo_gym import ( + NemoGymEnv, + _build_dataset, + _resolve_gym_config, +) + + +def load_environment( + dataset_path: str | None = None, + **kwargs: Any, +) -> vf.Environment: + dataset, _ = _build_dataset( + "workplace_assistant", "example", dataset_path=dataset_path + ) + return NemoGymEnv( + gym_configs=[_resolve_gym_config("workplace_assistant")], + dataset=dataset, + **kwargs, + ) diff --git a/environments/nemo_gym/nemo_gym_workplace_assistant/pyproject.toml b/environments/nemo_gym/nemo_gym_workplace_assistant/pyproject.toml new file mode 100644 index 0000000000..609f19b544 --- /dev/null +++ b/environments/nemo_gym/nemo_gym_workplace_assistant/pyproject.toml @@ -0,0 +1,21 @@ +[project] +name = "nemo-gym-workplace-assistant" +description = "NeMo Gym workplace assistant environment" +tags = ["nemo-gym", "agent", "tools", "multi-turn"] +version = "0.1.0" +requires-python = ">=3.12" +dependencies = [ + "verifiers>=0.1.11.dev1", + "nemo-gym>=0.2.0", +] + +[build-system] +requires = ["hatchling"] +build-backend = "hatchling.build" + +[tool.hatch.build] +include = ["nemo_gym_workplace_assistant.py", "pyproject.toml"] + +[tool.verifiers.eval] +num_examples = 5 +rollouts_per_example = 1 diff --git a/environments/nemo_gym/nemo_gym_xlam_fc/nemo_gym_xlam_fc.py b/environments/nemo_gym/nemo_gym_xlam_fc/nemo_gym_xlam_fc.py new file mode 100644 index 0000000000..c425a260e7 --- /dev/null +++ b/environments/nemo_gym/nemo_gym_xlam_fc/nemo_gym_xlam_fc.py @@ -0,0 +1,20 @@ +from typing import Any + +import verifiers as vf +from verifiers.envs.integrations.nemo_gym import ( + NemoGymEnv, + _build_dataset, + _resolve_gym_config, +) + + +def load_environment( + dataset_path: str | None = None, + **kwargs: Any, +) -> vf.Environment: + dataset, _ = _build_dataset("xlam_fc", "example", dataset_path=dataset_path) + return NemoGymEnv( + gym_configs=[_resolve_gym_config("xlam_fc")], + dataset=dataset, + **kwargs, + ) diff --git a/environments/nemo_gym/nemo_gym_xlam_fc/pyproject.toml b/environments/nemo_gym/nemo_gym_xlam_fc/pyproject.toml new file mode 100644 index 0000000000..8b266c61e4 --- /dev/null +++ b/environments/nemo_gym/nemo_gym_xlam_fc/pyproject.toml @@ -0,0 +1,21 @@ +[project] +name = "nemo-gym-xlam-fc" +description = "NeMo Gym xlam function calling environment" +tags = ["nemo-gym"] +version = "0.1.0" +requires-python = ">=3.12" +dependencies = [ + "verifiers>=0.1.11.dev1", + "nemo-gym>=0.2.0", +] + +[build-system] +requires = ["hatchling"] +build-backend = "hatchling.build" + +[tool.hatch.build] +include = ["nemo_gym_xlam_fc.py", "pyproject.toml"] + +[tool.verifiers.eval] +num_examples = 5 +rollouts_per_example = 1 diff --git a/tests/test_envs.py b/tests/test_envs.py index ffbc90a384..7199b54ef0 100644 --- a/tests/test_envs.py +++ b/tests/test_envs.py @@ -24,6 +24,8 @@ # Uses prime-tunnel which is still experimental and has low usage limits "terminus_harbor", "opencode_harbor", + # Contains nested sub-environments (nemo_gym_*/), not a top-level environment itself + "nemo_gym", ] SKIPPED_ENV_LOADING_ENVS = [ diff --git a/verifiers/envs/integrations/nemo_gym/__init__.py b/verifiers/envs/integrations/nemo_gym/__init__.py new file mode 100644 index 0000000000..ccc4cbc4d6 --- /dev/null +++ b/verifiers/envs/integrations/nemo_gym/__init__.py @@ -0,0 +1,9 @@ +from .env import NemoGymEnv +from .utils import _build_dataset, _resolve_gym_config, _reward_from_nemo_gym + +__all__ = [ + "NemoGymEnv", + "_build_dataset", + "_resolve_gym_config", + "_reward_from_nemo_gym", +] diff --git a/verifiers/envs/integrations/nemo_gym/env.py b/verifiers/envs/integrations/nemo_gym/env.py new file mode 100644 index 0000000000..8f7923a6dc --- /dev/null +++ b/verifiers/envs/integrations/nemo_gym/env.py @@ -0,0 +1,225 @@ +from __future__ import annotations + +import asyncio +import importlib.util +import json +import os +import threading +import time +from pathlib import Path +from typing import Any + +from datasets import Dataset + +import verifiers as vf +from verifiers.clients import Client +from verifiers.types import RolloutInput, SamplingArgs, State + +from .utils import ( + _map_nemo_gym_result_to_state, + _resolve_agent_name, + _reward_from_nemo_gym, +) + + +class NemoGymEnv(vf.Environment): + def __init__( + self, + *, + gym_configs: list[str], + dataset: Dataset, + rubric: vf.Rubric | None = None, + vllm_server_host: str = "127.0.0.1", + vllm_server_port: int = 8000, + head_server_host: str = "0.0.0.0", + head_server_port: int = 11000, + head_server_client_host: str = "127.0.0.1", + policy_base_url: str | None = None, + policy_api_key: str | None = None, + policy_model_config: str | None = None, + system_prompt: str | None = None, + **kwargs: Any, + ): + self.gym_configs = gym_configs + self.vllm_server_host = vllm_server_host + self.vllm_server_port = vllm_server_port + self.head_server_host = head_server_host + self.head_server_port = head_server_port + self.head_server_client_host = head_server_client_host + self.policy_base_url = policy_base_url + self.policy_api_key = policy_api_key + self.policy_model_config = policy_model_config + + self._run_helper: Any | None = None + self._rch: Any | None = None + self._head_server_config: Any | None = None + self._server_lock = asyncio.Lock() + self._agent_name: str | None = None + + self._bg_loop: asyncio.AbstractEventLoop = asyncio.new_event_loop() + self._bg_thread = threading.Thread( + target=self._bg_loop.run_forever, daemon=True, name="nemo-gym-loop" + ) + self._bg_thread.start() + + super().__init__( + dataset=dataset, + rubric=rubric or vf.Rubric(funcs=[_reward_from_nemo_gym], weights=[1.0]), + system_prompt=system_prompt, + message_type="chat", + **kwargs, + ) + + def _start_run_helper(self, model: str) -> None: + try: + from nemo_gym.cli import GlobalConfigDictParserConfig, RunHelper + from nemo_gym.rollout_collection import RolloutCollectionHelper + from nemo_gym.server_utils import HEAD_SERVER_KEY_NAME, BaseServerConfig + from omegaconf import DictConfig + except ImportError as exc: + raise ImportError( + "NemoGymEnv currently requires nemo-gym installed as an editable local clone (this will update to support PyPI soon):\n" + " git clone https://github.com/NVIDIA-NeMo/Gym /path/to/Gym\n" + " pip install -e /path/to/Gym" + ) from exc + + if self.policy_model_config: + policy_model_config = self.policy_model_config + else: + responses_spec = importlib.util.find_spec("responses_api_models") + if responses_spec and responses_spec.submodule_search_locations: + responses_root = Path(next(iter(responses_spec.submodule_search_locations))) + policy_model_config = str( + responses_root / "vllm_model" / "configs" / "vllm_model_for_training.yaml" + ) + else: + raise RuntimeError( + "Could not locate responses_api_models. " + "nemo-gym must be installed as an editable local clone: pip install -e /path/to/Gym (this will update to support PyPI soon)." + ) + + config = { + HEAD_SERVER_KEY_NAME: { + "host": self.head_server_host, + "port": self.head_server_port, + }, + "config_paths": [ + policy_model_config, + *self.gym_configs, + ], + "policy_base_url": self.policy_base_url + or f"http://{self.vllm_server_host}:{self.vllm_server_port}/v1", + "policy_api_key": self.policy_api_key + or os.environ.get("POLICY_API_KEY", "EMPTY"), + "policy_model_name": model, + "global_aiohttp_connector_limit_per_host": 16_384, + "global_aiohttp_connector_limit": 65_536, + "skip_venv_if_present": True, + } + + hf_token = os.environ.get("HF_TOKEN") or os.environ.get( + "HUGGING_FACE_HUB_TOKEN" + ) + if hf_token: + config["hf_token"] = hf_token + + rh = RunHelper() + rh.start( + global_config_dict_parser_config=GlobalConfigDictParserConfig( + initial_global_config_dict=DictConfig(config), + skip_load_from_cli=True, + ) + ) + + self._run_helper = rh + self._head_server_config = BaseServerConfig( + host=self.head_server_client_host, + port=self.head_server_port, + ) + self._rch = RolloutCollectionHelper() + + async def _ensure_server(self, model: str) -> tuple[Any, Any]: + if self._rch is not None: + return self._rch, self._head_server_config + async with self._server_lock: + if self._rch is not None: + return self._rch, self._head_server_config + loop = asyncio.get_running_loop() + await loop.run_in_executor(None, self._start_run_helper, model) + return self._rch, self._head_server_config + + @vf.teardown + async def teardown_server(self) -> None: + if self._run_helper is not None: + try: + loop = asyncio.get_running_loop() + await loop.run_in_executor(None, self._run_helper.shutdown) + except Exception: + pass + self._run_helper = None + self._rch = None + self._head_server_config = None + self._bg_loop.call_soon_threadsafe(self._bg_loop.stop) + + async def rollout( + self, + input: RolloutInput, + client: Client, + model: str, + sampling_args: SamplingArgs | None = None, + ) -> State: + state = await self.init_state(input, client, model, sampling_args) + start_time: float = state["timing"]["start_time"] + + try: + rch, head_server_config = await self._ensure_server(model) + + dataset_row: dict[str, Any] = json.loads(state["info"]["dataset_row_json"]) + dataset_row["_rowidx"] = 0 + if "agent_ref" not in dataset_row: + if self._agent_name is None: + self._agent_name = _resolve_agent_name(self.gym_configs[0]) + dataset_row["agent_ref"] = {"name": self._agent_name} + + rcp: dict[str, Any] = dataset_row.setdefault("responses_create_params", {}) + if sampling_args: + for key in ("temperature", "top_p"): + if sampling_args.get(key) is not None: + rcp[key] = sampling_args[key] + + # reuse the zmq loop, kinda ugly + async def _run() -> Any: + for task in rch.run_examples( + examples=[dataset_row], + head_server_config=head_server_config, + ): + _row, result = await task + return result + return None + + future = asyncio.run_coroutine_threadsafe(_run(), self._bg_loop) + nemo_gym_result = await asyncio.get_running_loop().run_in_executor( + None, future.result + ) + + except Exception as exc: + state["error"] = vf.InfraError( + f"NemoGymEnv rollout failed: {type(exc).__name__}: {exc}" + ) + state["completion"] = [] + state["is_completed"] = True + state["stop_condition"] = "has_error" + _fill_timing(state, start_time) + return state + + _map_nemo_gym_result_to_state(state, nemo_gym_result, model) + state["stop_condition"] = "has_error" if state.get("error") else "completed" + state["is_completed"] = True + _fill_timing(state, start_time) + return state + + +def _fill_timing(state: State, start_time: float) -> None: + elapsed_ms = (time.time() - start_time) * 1000.0 + state["timing"]["generation_ms"] = elapsed_ms + state["timing"]["total_ms"] = elapsed_ms diff --git a/verifiers/envs/integrations/nemo_gym/utils.py b/verifiers/envs/integrations/nemo_gym/utils.py new file mode 100644 index 0000000000..f920112896 --- /dev/null +++ b/verifiers/envs/integrations/nemo_gym/utils.py @@ -0,0 +1,306 @@ +from __future__ import annotations + +import importlib.util +import json +import time +import uuid + +import yaml +from pathlib import Path +from typing import Any + +from datasets import Dataset + +from verifiers.types import ( + AssistantMessage, + Messages, + Response, + ResponseMessage, + ResponseTokens, + State, + ToolCall, + ToolMessage, + TrajectoryStep, +) + + +def _json_dumps(value: Any) -> str: + return json.dumps(value, ensure_ascii=False) + + +def _stringify(value: Any) -> str: + if value is None: + return "" + if isinstance(value, str): + return value + try: + return _json_dumps(value) + except (TypeError, ValueError): + return str(value) + + +def _resolve_resources_servers_root() -> Path: + resources_spec = importlib.util.find_spec("resources_servers") + if resources_spec and resources_spec.submodule_search_locations: + root = Path(next(iter(resources_spec.submodule_search_locations))).resolve() + if root.exists(): + return root + + nemo_spec = importlib.util.find_spec("nemo_gym") + if nemo_spec and nemo_spec.origin: + nemo_root = Path(nemo_spec.origin).resolve().parent + sibling = nemo_root.parent / "resources_servers" + if sibling.exists(): + return sibling + + raise RuntimeError( + "Unable to locate NeMo Gym resources_servers package. " + "Install `nemo-gym` or pass `dataset_path` explicitly." + ) + + +def _build_dataset( + resources_server: str, + dataset_split: str, + dataset_path: str | None = None, + dataset_limit: int | None = None, +) -> tuple[Dataset, Path]: + if dataset_path is not None: + path = Path(dataset_path).expanduser().resolve() + if not path.exists(): + raise FileNotFoundError(f"dataset_path does not exist: {path}") + else: + root = _resolve_resources_servers_root() + path = root / resources_server / "data" / f"{dataset_split}.jsonl" + if not path.exists(): + raise FileNotFoundError( + f"Could not find dataset for '{resources_server}' split '{dataset_split}': {path}" + ) + + rows: list[dict[str, Any]] = [] + with path.open("r", encoding="utf-8") as f: + for line_no, line in enumerate(f, start=1): + line = line.strip() + if not line: + continue + try: + row = json.loads(line) + except json.JSONDecodeError as exc: + raise ValueError( + f"Invalid JSON in {path} line {line_no}: {exc}" + ) from exc + if not isinstance(row, dict): + raise ValueError(f"Row {line_no} in {path} is not an object") + if "responses_create_params" not in row: + raise ValueError( + f"Row {line_no} in {path} is missing 'responses_create_params'" + ) + rows.append(row) + if not rows: + raise ValueError(f"Dataset file {path} contains no rows") + + if dataset_limit is not None: + if dataset_limit <= 0: + raise ValueError("dataset_limit must be > 0") + rows = rows[:dataset_limit] + + dataset_rows: list[dict[str, Any]] = [] + for row in rows: + rcp = row["responses_create_params"] + raw_input = rcp.get("input", []) + if isinstance(raw_input, str): + prompt = [{"role": "user", "content": raw_input}] + elif isinstance(raw_input, list): + prompt = raw_input + else: + prompt = [{"role": "user", "content": _stringify(raw_input)}] + dataset_rows.append( + { + "prompt": prompt, + "answer": _stringify(row.get("answer", "")), + "task": resources_server, + "info": {"dataset_row_json": _json_dumps(row)}, + } + ) + + return Dataset.from_list(dataset_rows), path + + +def _resolve_gym_config(resources_server: str, config_name: str | None = None) -> str: + root = _resolve_resources_servers_root() + name = config_name or resources_server + path = root / resources_server / "configs" / f"{name}.yaml" + if not path.exists(): + raise FileNotFoundError( + f"Could not find NeMo Gym config for '{resources_server}/{name}': {path}" + ) + return str(path) + + +def _resolve_agent_name(gym_config_path: str) -> str: + with open(gym_config_path) as f: + config = yaml.safe_load(f) + for key, value in config.items(): + if isinstance(value, dict) and "responses_api_agents" in value: + return key + raise RuntimeError( + f"Could not find a responses_api_agents entry in {gym_config_path}" + ) + + +def _reward_from_nemo_gym(state: State, **kwargs: Any) -> float: + return float(state.get("nemo_gym_reward", 0.0) or 0.0) + + +def _nemo_item_to_assistant_message(item: dict[str, Any]) -> AssistantMessage: + item_type = item.get("type") + + if item_type == "message": + content_blocks = item.get("content") or [] + text = "\n".join( + c.get("text", "") + for c in content_blocks + if isinstance(c, dict) and c.get("type") == "output_text" + ) + return AssistantMessage(role="assistant", content=text or None) + + if item_type == "function_call": + tool_call = ToolCall( + id=str(item.get("call_id") or item.get("id") or uuid.uuid4().hex[:8]), + name=str(item.get("name", "")), + arguments=str(item.get("arguments", "{}")), + ) + return AssistantMessage(role="assistant", content=None, tool_calls=[tool_call]) + + return AssistantMessage(role="assistant", content=str(item)) + + +def _make_response( + msg: AssistantMessage, + model: str, + gen_ids: list[int], + logprobs: list[float], + prompt_ids: list[int], +) -> Response: + tokens = ResponseTokens( + prompt_ids=prompt_ids, + prompt_mask=[1] * len(prompt_ids), + completion_ids=gen_ids, + completion_mask=[1] * len(gen_ids), + completion_logprobs=logprobs, + routed_experts=None, + ) + return Response( + id=f"nemo_gym-{uuid.uuid4().hex[:8]}", + created=int(time.time()), + model=model, + usage=None, + message=ResponseMessage( + role="assistant", + content=msg.content, + tool_calls=msg.tool_calls, + finish_reason="tool_calls" if msg.tool_calls else "stop", + is_truncated=False, + tokens=tokens, + ), + ) + + +def _build_trajectory_from_nemo( + output_items: list[dict[str, Any]], + initial_prompt: Messages, + model: str, + trajectory_id: str, +) -> tuple[list[TrajectoryStep], Messages]: + trajectory: list[TrajectoryStep] = [] + completion_messages: list = [] + all_messages: list = list(initial_prompt) + + for item in output_items: + if item.get("type") == "function_call_output": + tool_msg = ToolMessage( + role="tool", + tool_call_id=str(item.get("call_id", "")), + content=str(item.get("output", "")), + ) + all_messages.append(tool_msg) + completion_messages.append(tool_msg) + continue + + if "generation_token_ids" not in item: + continue + + prompt_ids: list[int] = list(item.get("prompt_token_ids") or []) + gen_ids: list[int] = list(item.get("generation_token_ids") or []) + logprobs: list[float] = list( + item.get("generation_log_probs") or [0.0] * len(gen_ids) + ) + + step_prompt: Messages = list(all_messages) + assistant_msg = _nemo_item_to_assistant_message(item) + all_messages.append(assistant_msg) + completion_messages.append(assistant_msg) + + trajectory.append( + { + "prompt": step_prompt, + "completion": [assistant_msg], + "response": _make_response( + assistant_msg, model, gen_ids, logprobs, prompt_ids + ), + "tokens": { + "prompt_ids": prompt_ids, + "prompt_mask": [1] * len(prompt_ids), + "completion_ids": gen_ids, + "completion_mask": [1] * len(gen_ids), + "completion_logprobs": logprobs, + "overlong_prompt": False, + "is_truncated": False, + "routed_experts": None, + }, + "reward": None, + "advantage": None, + "is_truncated": False, + "trajectory_id": trajectory_id, + "extras": {}, + } + ) + + return trajectory, completion_messages + + +def _map_nemo_gym_result_to_state( + state: State, nemo_gym_result: Any, model: str +) -> None: + import verifiers as vf + + if not isinstance(nemo_gym_result, dict) or nemo_gym_result.get("error"): + error_detail = ( + nemo_gym_result.get("error", "unknown error") + if isinstance(nemo_gym_result, dict) + else repr(nemo_gym_result) + ) + state["error"] = vf.InfraError( + f"NeMo Gym agent server rollout failed: {error_detail}" + ) + state["nemo_gym_reward"] = 0.0 + state["completion"] = [] + return + + state["nemo_gym_reward"] = float(nemo_gym_result.get("reward", 0.0) or 0.0) + state["nemo_gym_result"] = nemo_gym_result + + output_items: list[dict[str, Any]] = (nemo_gym_result.get("response") or {}).get( + "output" + ) or [] + + trajectory, completion_messages = _build_trajectory_from_nemo( + output_items=output_items, + initial_prompt=state["prompt"], + model=model, + trajectory_id=state["trajectory_id"], + ) + + state["trajectory"] = trajectory + state["completion"] = completion_messages + state["is_truncated"] = False diff --git a/verifiers/utils/install_utils.py b/verifiers/utils/install_utils.py index b15b34edbd..8561ef304e 100644 --- a/verifiers/utils/install_utils.py +++ b/verifiers/utils/install_utils.py @@ -200,7 +200,15 @@ def install_from_local(env_name: str, env_dir: str = "./environments") -> bool: True if installation succeeded """ env_folder = normalize_package_name(env_name) - env_path = Path(env_dir) / env_folder + base = Path(env_dir) + env_path = base / env_folder + + # Also check one level of subdirectories (e.g. environments/nemo_gym/nemo_mcqa/). + if not env_path.exists() and base.is_dir(): + for subdir in base.iterdir(): + if subdir.is_dir() and (subdir / env_folder).is_dir(): + env_path = subdir / env_folder + break if not env_path.exists(): logger.error(f"Local environment not found: {env_path}") diff --git a/verifiers/utils/message_utils.py b/verifiers/utils/message_utils.py index da8994f6f7..1334508902 100644 --- a/verifiers/utils/message_utils.py +++ b/verifiers/utils/message_utils.py @@ -107,6 +107,8 @@ def from_raw_message(message: dict) -> Message: return TextMessage.model_validate(message) elif message["role"] == "system": return SystemMessage.model_validate(message) + elif message["role"] == "developer": + return SystemMessage.model_validate({**message, "role": "system"}) elif message["role"] == "user": return UserMessage.model_validate(message) elif message["role"] == "assistant":