nemo_gym.rollout_collection

View as Markdown

Module Contents

Classes

NameDescription
DispatchLatencyTrackerPer-task latency observed by the dispatcher, and the drain margin from it.
E2ERolloutCollectionConfigSpin up all necessary servers and perform a batch of rollout collection using each dataset inside the provided configs.
RolloutAggregationConfigAggregate metrics across rollout shards produced by gym eval run --no-serve +disable_aggregation=true.
RolloutAggregationHelper-
RolloutCollectionConfigPerform a batch of rollout collection.
RolloutCollectionHelper-
SharedRolloutCollectionConfig-
_BoundedCompletionIteratorCompletion-order iterator with a bounded set of resident asyncio tasks.
_CompletedRolloutA finished /run dispatch, with timing carried alongside (not inside) the raw result.

Functions

NameDescription
_agent_request_failure_rowOne sidecar row for a /run call that came back without a result.
_attach_ng_perf-
_attach_trajectory_record-
_build_ng_perfAssemble the per-rollout ng_perf summary from ng_trajectory.
_build_trajectory_record-
_counted_failure_rowThe metric-input copy of a sidecar row counted as zero.
_coverage_reportState how much of the input the score covers, for the runs where it is not all of it.
_dispatch_drained_resultThe result Gym builds for a row the dispatch budget never started.
_drop_truncated_tailRepair a jsonl whose last line a hard kill cut short, so resume can read and append to it.
_environment_server_for_agentReturn the one environment server that fronts an agent.
_environment_server_for_config_rowPick the environment server a row is dispatched to, or None for today’s agent path.
_environment_servers_by_agentMap each agent name to the environment servers whose agent_server names it.
_episode_recordTurn a BaseEpisodeResponse into the rollout record the collector stores.
_expand_input_globExpand a glob-or-comma-separated-globs string into a sorted, deduplicated list of paths.
_failure_rows_counted_as_zeroSidecar rows the caller opted to count in the metrics denominator.
_fill_task_fieldsGive the rows the metric input adds, counted failures and imputed zeros, their task’s dataset fields.
_has_observation_gap-
_identity_keyA rollout’s key across runs, which each number their tasks from 0.
_is_collector_keyKeys rollout collection writes itself; an Environment Server result must not use them.
_is_episode_responseTrue for a BaseEpisodeResponse-shaped reply: object identities plus a result or failure key.
_latest_failure_rowsThe last attempt recorded for each rollout across the failures sidecars.
_masking_step_metricsIn-progress view of what a run is losing to its environment rather than its policy.
_materialized_taskset-
_metrics_groupThe environment server a rollout is scored under, as _call_aggregate_metrics groups it.
_missing_rollout_rows_counted_as_zeroMaterialized rollouts that produced no row anywhere, counted as zeros.
_native_episode_request_body-
_nonnegative_int-
_normalize_health_check_ignored_checks-
_read_jsonl-
_rollout_for_exportReturn an exporter view without the complete trajectory or raw capture payloads.
_rollout_order_keyTask then repeat: metrics that read a task’s repeats positionally need that order.
_rollout_request_debug_summary-
_routing_identityThe agent a rollout ran on, or for a native taskset row, which names none, its environment server.
_strip_capture_payloads-
_trajectory_identity-
_truncated_bodyDecode at most the kept prefix, so a huge error page is never decoded in full.
_turn_contentSplit one captured model call into the turn’s question, answer, reasoning, and tool-call count.
_turns_from_model_callsBuild one turn per captured model call that returned a response, for an agent that sent no trajectory.
get_max_rollout_attemptsRead NEMO_GYM_MAX_ROLLOUT_ATTEMPTS (positive int) or default to 3.
is_terminal_failureWhether a persisted failure must be gated on resume.
loads_jsonl_lineParse one JSONL line, raising a clean ConfigError (naming file + line) on malformed JSON.
migrate_invalid_judge_main_rowsMove legacy invalid-judge rows from the main JSONL into the sidecar.
observed_elapsedBest-effort per-rollout wallclock from a result/failure row.

Data

AGENT_REQUEST_FAILED_FAILURE_CLASS

AGENT_RUN_ERROR_FAILURE_CLASS

ENVIRONMENT_SERVER_FAILURE_CLASS

NG_DISPATCH_DRAINED_KEY

NG_ELAPSED_KEY

NG_ENVIRONMENT_SERVER_KEY

NG_FAILURE_CLASS_KEY

NG_NO_PERSIST_KEY

NG_PERF_KEY

NG_RESULT_TYPE_KEY

NG_TASK_ID_KEY

NG_TERMINAL_KEY

NG_TRAJECTORY_KEY

_DEFAULT_MAX_ROLLOUT_ATTEMPTS

_MAX_FAILURE_BODY_CHARS

_MIGRATED_INVALID_JUDGE_KEY

_MODEL_CALL_PAYLOAD_KEYS

_NO_RESULT_FAILURE_CLASSES

_RUN_FAILURE_ERRORS

_SERVER_DID_NOT_RUN_STATUSES

__getattr__

_failures_path_for

_get_max_rollout_attempts

logger

API

class nemo_gym.rollout_collection.DispatchLatencyTracker()

Per-task latency observed by the dispatcher, and the drain margin from it.

Rollouts/hr is the number operators watch, and on its own it is misleading: raising concurrency raises aggregate throughput while making every individual task slower, so the run looks healthier right up until tasks start breaching their per-task timeout. Reporting the latency percentiles next to the rate makes that trade visible while there is still time to react to it.

_drained
= 0
_durations
List[float] = []
drained
int
nemo_gym.rollout_collection.DispatchLatencyTracker.drain_margin(
configured: typing.Optional[float]
) -> typing.Optional[float]

Seconds of headroom a task needs before it is worth starting.

An explicit value wins. Otherwise adapt to this run’s own p75, once enough tasks have finished for that to mean anything.

nemo_gym.rollout_collection.DispatchLatencyTracker.quantile(
q: float
) -> typing.Optional[float]
nemo_gym.rollout_collection.DispatchLatencyTracker.record(
seconds: float
) -> None
nemo_gym.rollout_collection.DispatchLatencyTracker.record_drained() -> None
nemo_gym.rollout_collection.DispatchLatencyTracker.summary() -> str
class nemo_gym.rollout_collection.E2ERolloutCollectionConfig()

Bases: SharedRolloutCollectionConfig

Spin up all necessary servers and perform a batch of rollout collection using each dataset inside the provided configs.

Examples:

gym eval run +output_jsonl_fpath=weather_rollouts.jsonl +num_samples_in_parallel=10
reuse_existing_data_preparation
bool = False
split
Union[Literal['train'], Literal['validation'], Literal['benchmark']]
nemo_gym.rollout_collection.E2ERolloutCollectionConfig._reject_example_split(
data
)
classmethod
nemo_gym.rollout_collection.E2ERolloutCollectionConfig._reject_input_jsonl_fpath(
data
)
classmethod
class nemo_gym.rollout_collection.RolloutAggregationConfig()

Bases: BaseNeMoGymCLIConfig

Aggregate metrics across rollout shards produced by gym eval run --no-serve +disable_aggregation=true.

Reads every JSONL file matching input_glob, computes aggregate metrics by POSTing to each agent server’s /aggregate_metrics endpoint over the global union of records, and writes a single <output_jsonl_fpath stem>_aggregate_metrics.json next to the rollouts. By default also concatenates all shards into output_jsonl_fpath.

Examples:

gym eval aggregate "+config_paths=[benchmarks/aime24/config.yaml,responses_api_models/vllm_model/configs/vllm_model.yaml]" +input_glob='results/rollouts-rs*-chunk*.jsonl' +output_jsonl_fpath=results/rollouts.jsonl
count_failure_classes_as_zero
List[str]
count_missing_rollouts_as_zero
bool
disable_health_check
bool
health_check_ignored_checks
List[str]
health_check_workers
Optional[int]
input_glob
str
merge_shards
bool
output_jsonl_fpath
str
nemo_gym.rollout_collection.RolloutAggregationConfig._validate_health_check_ignored_checks(
value
)
classmethod
class nemo_gym.rollout_collection.RolloutAggregationHelper()

Bases: BaseModel

nemo_gym.rollout_collection.RolloutAggregationHelper.run_from_config(
) -> typing.Optional[pathlib.Path]
async
class nemo_gym.rollout_collection.RolloutCollectionConfig()

Bases: SharedRolloutCollectionConfig

Perform a batch of rollout collection.

Examples:

gym eval run --no-serve +agent_name=example_single_tool_call_simple_agent +input_jsonl_fpath=weather_query.jsonl +output_jsonl_fpath=weather_rollouts.jsonl +limit=100 +num_repeats=4 +num_samples_in_parallel=10
agent_map
Optional[Dict[str, str]]
agent_name
Optional[str]
dispatch_budget_s
Optional[float]
dispatch_longest_first
bool
drain_margin_s
Optional[float]
fan_out
Optional[Dict[str, List[str]]]
input_jsonl_fpath
str
interleave_repeats
bool
limit
Optional[int]
materialized_jsonl_fpath
Path
num_repeats
Union[int, Dict[str, int]]
num_repeats_add_seed
bool
prompt_config
Optional[str]
resume_from_cache
bool
retry_invalid_judge_responses
bool
retry_terminal_timeouts
bool
skills
Optional[SkillsConfig]
nemo_gym.rollout_collection.RolloutCollectionConfig._coerce_null_num_repeats(
v
)
classmethod
class nemo_gym.rollout_collection.RolloutCollectionHelper()

Bases: BaseModel

nemo_gym.rollout_collection.RolloutCollectionHelper._agent_name_for_row(
row: dict[str, typing.Any],
global_config_dict: omegaconf.DictConfig
) -> str | None
staticmethod
nemo_gym.rollout_collection.RolloutCollectionHelper._call_aggregate_metrics(
results: typing.List[typing.Dict],
rows: typing.List[typing.Dict],
output_fpath: pathlib.Path
) -> typing.Optional[pathlib.Path]
async

Call /aggregate_metrics on the environment server each rollout ran through.

Rows are grouped by the environment server stamped on them at preprocessing (_ng_environment_server); a row without the stamp is grouped by the environment server that fronts its agent, as before, so the identity decided at dispatch is the one aggregation uses, across shards and resumed runs alike. Writes a single _aggregate_metrics.json with one entry per environment server (same shape as the old _agent_metrics.json, plus the server name). Returns the file path.

nemo_gym.rollout_collection.RolloutCollectionHelper._dispatch_name(
row: dict[str, typing.Any]
) -> str
staticmethod
nemo_gym.rollout_collection.RolloutCollectionHelper._load_from_cache(
retain_results_in_memory: bool = True,
success_keys: typing.Optional[set] = None
) -> typing.Tuple[typing.List[typing.Dict], typing.List[typing.Dict], typing.List[typing.Dict], typing.List[typing.List[bytes]]]
nemo_gym.rollout_collection.RolloutCollectionHelper._preprocess_raw_rows(
raw_rows: typing.List[typing.Tuple[int, str, typing.Dict]],
) -> typing.List[typing.Dict]
staticmethod
nemo_gym.rollout_collection.RolloutCollectionHelper._preprocess_rows_from_config(
) -> typing.List[typing.Dict]
nemo_gym.rollout_collection.RolloutCollectionHelper._run_examples_with_metadata(
examples: typing.List[typing.Dict],
head_server_config: typing.Optional[nemo_gym.config_types.BaseServerConfig] = None,
semaphore: typing.Optional[asyncio.Semaphore] = None,
route_failures_to_sidecar: bool = False,
environment_server_name: str | None = None,
max_resident_tasks: typing.Optional[int] = None,
dispatch_budget_s: typing.Optional[float] = None,
drain_margin_s: typing.Optional[float] = None,
latency_tracker: typing.Optional[nemo_gym.rollout_collection.DispatchLatencyTracker] = None
) -> typing.Iterator[asyncio.Future]

Internal dispatch shared by run_examples and Gym’s own collection paths.

When max_resident_tasks is set, at most that many rollout tasks are admitted at once. When unset, all examples are scheduled as before. The collection owner closes the bounded iterator on cancellation or error.

dispatch_budget_s stops starting rows that many seconds after this call, and drain_margin_s stops sooner for rows that would not have time to finish. A row drained this way resolves to a Gym-built _dispatch_drained_result.

Identical contract to run_examples, but each future resolves to a _CompletedRollout that carries rollout_latency_ms alongside the raw /run result instead of inside it, so internal-only timing never has to be smuggled through (and stripped back out of) a dict that a direct caller of run_examples could also observe.

nemo_gym.rollout_collection.RolloutCollectionHelper._run_from_config(
) -> typing.Tuple[typing.List[typing.Dict]]
async
nemo_gym.rollout_collection.RolloutCollectionHelper._stamp_environment_server_agent_refs(
examples: list[dict[str, typing.Any]],
global_config_dict: omegaconf.DictConfig
) -> None
classmethod

Stamp compatibility-routed rows with the bound agent.

These rows never reach resolve_task_sources, so without this they carry no agent_ref and results, aggregate metrics and reward profiling lose the agent they ran on. A row that already names an agent is left alone; the name is validated against the environment server by _validate_environment_servers. Materialized tasks are skipped: they carry no agent by design, and their result projection is not the legacy shape this key belongs to.

nemo_gym.rollout_collection.RolloutCollectionHelper._validate_agent_names(
examples: typing.List[typing.Dict],
global_config_dict: omegaconf.DictConfig
) -> None
staticmethod

Fail before any dispatch when a row names an agent absent from the running config.

Without this, the first bad row dies mid-collection with a raw omegaconf ConfigKeyError after valid rows have already been dispatched.

nemo_gym.rollout_collection.RolloutCollectionHelper._validate_agent_pairings(
examples: typing.List[typing.Dict],
global_config_dict: omegaconf.DictConfig
) -> None
staticmethod

Fail before dispatch when a row points to agent incompatible with the resources server it runs on.

nemo_gym.rollout_collection.RolloutCollectionHelper._validate_environment_servers(
examples: list[dict],
global_config_dict: omegaconf.DictConfig
) -> None
classmethod
nemo_gym.rollout_collection.RolloutCollectionHelper.preprocess_examples(
examples: typing.List[typing.Dict],
agent_map: typing.Optional[typing.Dict[str, str]] = None,
fan_out: typing.Optional[typing.Dict[str, typing.List[str]]] = None,
num_repeats: typing.Union[int, typing.Dict[str, int]] = 1,
num_repeats_add_seed: bool = False,
global_config_dict: typing.Optional[omegaconf.DictConfig] = None
) -> typing.List[typing.Dict]

Apply run-level routing and repetition to caller-held rows.

Public entry point for direct run_examples callers (e.g. trainer integrations that drive dispatch themselves): run_examples resolves task_sources and validates agent names, but agent_map, fan_out and num_repeats are applied only during preprocessing. Call this first, then pass the returned rows to run_examples.

Pass global_config_dict (the merged config) to also resolve task_source-only rows to their agents here; leave it None to defer that to run_examples, which does it against the head server’s config. Input rows are not mutated; the expanded, stamped copies are returned.

nemo_gym.rollout_collection.RolloutCollectionHelper.resolve_task_sources(
examples: typing.List[typing.Dict],
global_config_dict: omegaconf.DictConfig
) -> None
staticmethod

Stamp an agent_ref onto every row that carries only a task_source.

task_source names the config instance that declared the row’s dataset. Resolution is resolve_dataset_agent — the same rules benchmark discovery uses, so dispatch can never disagree with the listing. Conflicting agent: pins across one instance’s datasets are a hard error (rows carry only the instance name), as are unknown/non-routable instances; +agent_map is the disambiguator.

Rows that already have an agent_ref are left untouched, so this is a no-op on legacy datasets and on already-resolved (materialized) rows. Runs before any dispatch.

nemo_gym.rollout_collection.RolloutCollectionHelper.run_examples(
examples: typing.List[typing.Dict],
head_server_config: typing.Optional[nemo_gym.config_types.BaseServerConfig] = None,
semaphore: typing.Optional[asyncio.Semaphore] = None,
route_failures_to_sidecar: bool = False,
environment_server_name: str | None = None,
max_resident_tasks: typing.Optional[int] = None,
dispatch_budget_s: typing.Optional[float] = None,
drain_margin_s: typing.Optional[float] = None,
latency_tracker: typing.Optional[nemo_gym.rollout_collection.DispatchLatencyTracker] = None
) -> typing.Iterator[asyncio.Future]

We provide this function as a lower level interface for running rollout collection.

Rows are dispatched as given: task_sources are resolved and agent names validated here, but run-level knobs (agent_map, fan_out, num_repeats) are NOT applied — call preprocess_examples first if you need them.

route_failures_to_sidecar makes a failed /run a failure row instead of an exception that ends every rollout still in flight. It defaults off because those rollouts then leave the score.

max_resident_tasks limits admitted tasks and therefore concurrent requests, even when semaphore allows more. Admission starts when the first returned awaitable is awaited. None schedules all examples up front, as with asyncio.as_completed. Stopping iteration early leaves up to max_resident_tasks tasks running because this mapped iterator has no aclose().

dispatch_budget_s stops starting rows that many seconds after this call, and drain_margin_s stops sooner for rows that would not have time to finish.

Rows start in the order given, once the returned iterator is first consumed.

A row that ran resolves to exactly the (row, result) pair Gym’s own /run endpoint returned — no Gym-private fields are added to result. A row that produced no /run result resolves to a Gym-built result carrying Gym-private fields instead: a row drained by the dispatch budget gets _ng_failure_class="cancelled" with the _ng_dispatch_drained and _ng_no_persist markers, and a failed /run under route_failures_to_sidecar gets a failure row with _ng_failure_* fields.

nemo_gym.rollout_collection.RolloutCollectionHelper.run_from_config(
) -> typing.Tuple[typing.List[typing.Dict]]
async

Collect rollouts for a whole config. Wrapped in the run-scoped job span.

This is the driver side of an evaluation run and the outermost span Gym produces, so every rollout it dispatches is a descendant of it. job is in the default preset but deliberately not in per_rollout, where each rollout is meant to be its own bounded root trace.

nemo_gym.rollout_collection.RolloutCollectionHelper.setup_server_client(
head_server_config: typing.Optional[nemo_gym.config_types.BaseServerConfig] = None
class nemo_gym.rollout_collection.SharedRolloutCollectionConfig()

Bases: UploadRolloutsConfigMixin, BaseNeMoGymCLIConfig

count_failure_classes_as_zero
List[str]
count_missing_rollouts_as_zero
bool
disable_aggregation
bool
disable_health_check
bool
environment_routing_mode
Literal['agent', 'legacy', 'taskset']
environment_server_name
str | None
environment_server_routes
dict[str, str]
health_check_ignored_checks
List[str]
health_check_workers
Optional[int]
max_resident_rollout_tasks
Optional[int]
num_samples_in_parallel
Optional[int]
output_jsonl_fpath
str
require_complete
bool
responses_create_params
Dict[str, Any]
retain_results_in_memory
bool
rollout_collection_driver
Optional[str]
route_failures_to_sidecar
bool
nemo_gym.rollout_collection.SharedRolloutCollectionConfig._validate_health_check_ignored_checks(
value
)
classmethod
nemo_gym.rollout_collection.SharedRolloutCollectionConfig.check_completion(
expected: int,
results: typing.List[typing.Dict[str, typing.Any]]
) -> None

Reject incomplete submitted runs after saving their partial artifacts.

class nemo_gym.rollout_collection._BoundedCompletionIterator(
awaitables: typing.Iterator,
max_resident_tasks: int,
total: int
)

Completion-order iterator with a bounded set of resident asyncio tasks.

_awaitables
= iter(awaitables)
_lock
= asyncio.Lock()
_pending
set[Task] = set()
_progress
_ready
list[Task] = []
_resident_task_count
int
nemo_gym.rollout_collection._BoundedCompletionIterator.__iter__()
nemo_gym.rollout_collection._BoundedCompletionIterator.__next__()
nemo_gym.rollout_collection._BoundedCompletionIterator._admit() -> bool
nemo_gym.rollout_collection._BoundedCompletionIterator._fill() -> None
nemo_gym.rollout_collection._BoundedCompletionIterator._next_completed()
async
nemo_gym.rollout_collection._BoundedCompletionIterator.aclose() -> None
async
class nemo_gym.rollout_collection._CompletedRollout(
row: typing.Dict[str, typing.Any],
result: typing.Dict[str, typing.Any],
rollout_latency_ms: typing.Optional[float],
environment_server: typing.Optional[str] = None,
environment_server_type: typing.Optional[str] = None
)
Dataclass

A finished /run dispatch, with timing carried alongside (not inside) the raw result.

environment_server
Optional[str] = None
environment_server_type
Optional[str] = None
result
Dict[str, Any]
rollout_latency_ms
Optional[float]
row
Dict[str, Any]
nemo_gym.rollout_collection._agent_request_failure_row(
exc: BaseException,
status: typing.Optional[int]
) -> typing.Dict[str, typing.Any]

One sidecar row for a /run call that came back without a result.

No reward and no response: an infrastructure failure is not a verifier score of zero, and a placeholder would read as real generation data to token capture, aggregation and trainers. The class says whether the rollout ran. A NeMo Gym server answers 500 when its own handler raises, so any status it answered with means the rollout ran and broke, which is also how a model server rejecting the model’s own output arrives here. A gateway status, or no reply to take a status from, says nothing about the rollout. Neither class carries a reward; an evaluation that wants the first counted names it in count_failure_classes_as_zero.

nemo_gym.rollout_collection._attach_ng_perf(
result: dict[str, typing.Any],
observability_enabled: bool,
rollout_latency_ms: typing.Optional[float] = None
) -> None
nemo_gym.rollout_collection._attach_trajectory_record(
row: dict[str, typing.Any],
result: dict[str, typing.Any]
) -> None
nemo_gym.rollout_collection._build_ng_perf(
result: dict[str, typing.Any],
rollout_latency_ms: typing.Optional[float]
) -> typing.Optional[dict[str, typing.Any]]

Assemble the per-rollout ng_perf summary from ng_trajectory.

Returns None (ng_perf stays absent) unless at least one reasoning turn was observed: per-turn evidence is needed rather than just raw model-call capture, so a rollout collected with observability disabled produces no ng_perf at all.

Token fields are summed over every model call referenced by a reasoning-turn AgentInvocation. This includes compaction calls whenever the harness also lists them in AgentInvocation.model_calls.

num_turns counts reasoning turns summed across all invocations (an AgentInvocation is one root-agent or subagent conversation that may span many turns). Each invocation contributes its explicit TrajectoryTurn count when the harness emits turn records, falling back to its owned model-call count (one assistant response per turn), then to 1 (an invocation that ran had at least one turn) — so hybrid trajectories where only some invocations report turns still count every conversation.

token_observability_coverage reports what fraction of those turns actually resolved to a captured call: a turn whose ModelCallRef was unmatched or ambiguous silently loses its tokens from the sums below, and this is the only signal that it happened.

nemo_gym.rollout_collection._build_trajectory_record(
row: dict[str, typing.Any],
result: dict[str, typing.Any]
nemo_gym.rollout_collection._counted_failure_row(
row: typing.Dict[str, typing.Any]
) -> typing.Dict[str, typing.Any]

The metric-input copy of a sidecar row counted as zero.

nemo_gym.rollout_collection._coverage_report(
expected: int,
scored: int,
failure_counts: collections.Counter,
failures_hint: typing.Any,
imputed: int = 0
) -> str

State how much of the input the score covers, for the runs where it is not all of it.

Silence here is what makes a partial run look complete, so this reports against the materialized input rather than the rollouts one hop happened to dispatch, and names the rollouts that are in the score only as imputed zeros.

nemo_gym.rollout_collection._dispatch_drained_result(
remaining_s: float,
margin_s: typing.Optional[float]
) -> typing.Dict[str, typing.Any]

The result Gym builds for a row the dispatch budget never started.

No row is written anywhere: absence is the resume signal, so the task is re-dispatched intact next allocation instead of being started and killed part-way through. _ng_dispatch_drained keeps it out of capture, token finalization, progress metrics and the rollouts upload, since nothing ran.

nemo_gym.rollout_collection._drop_truncated_tail(
fpath: pathlib.Path
) -> None

Repair a jsonl whose last line a hard kill cut short, so resume can read and append to it.

nemo_gym.rollout_collection._environment_server_for_agent(
agent_name: str,
servers_by_agent: collections.abc.Mapping[str, list[str]]
) -> str

Return the one environment server that fronts an agent.

A row routed by its agent cannot choose between several environment servers. Several servers may still front one agent when every row names its server directly. Native tasksets name their environment server through environment_server_routes.

nemo_gym.rollout_collection._environment_server_for_config_row(
row: collections.abc.Mapping[str, typing.Any],
config: typing.Any
) -> str | None

Pick the environment server a row is dispatched to, or None for today’s agent path.

A materialized task (task_id.taskset plus task_input) always routes by its taskset: it is the native episode request and no agent-server /run accepts it. A flat row follows environment_routing_mode: agent keeps today’s routing (its agent’s environment server is resolved at dispatch), legacy sends every flat row to environment_server_name, and taskset refuses flat rows so a native-only run cannot silently pick up legacy input.

One batch may therefore hold both kinds of rows in agent and legacy mode. The chosen server is stamped on the row as _ng_environment_server and travels with it through the materialized input file, retries, and results.

nemo_gym.rollout_collection._environment_servers_by_agent(
global_config_dict: omegaconf.DictConfig
) -> dict[str, list[str]]

Map each agent name to the environment servers whose agent_server names it.

nemo_gym.rollout_collection._episode_record(
response: typing.Dict[str, typing.Any]
) -> typing.Dict[str, typing.Any]

Turn a BaseEpisodeResponse into the rollout record the collector stores.

A handled failure becomes a failures-sidecar row, so resume retries a non-terminal one and never a terminal one. A result is stored as the Environment Server returned it; the collector adds only its own _ng_* keys, so any Environment Server type can be collected without the collector knowing its result fields.

Every Environment Server type is scored the same way: through the result’s top-level reward, with optional top-level reward_components. A reward nested elsewhere in the result is stored as data, and a result without a top-level reward is unscored.

nemo_gym.rollout_collection._expand_input_glob(
input_glob: str
) -> typing.List[str]

Expand a glob-or-comma-separated-globs string into a sorted, deduplicated list of paths.

Examples:

'results/rollouts.jsonl' -> ['results/rollouts.jsonl'] (if it exists)
'a/*.jsonl, b/*.jsonl' -> matches of both patterns, deduplicated
nemo_gym.rollout_collection._failure_rows_counted_as_zero(
failures_fpaths: typing.List[pathlib.Path],
failure_classes: typing.List[str],
scored_keys: set
) -> typing.List[typing.Dict[str, typing.Any]]

Sidecar rows the caller opted to count in the metrics denominator.

The last attempt of a rollout is the one that stands, so it is selected across every failure class before the wanted classes are picked out. Selecting the other way round would let a stale attempt be counted after a later one landed in a class the caller did not ask for.

A row that already carries a reward is counted as it stands. A row that carries none records that no rollout happened, so it is counted as a zero here and only here: the score enters the metric input, never the sidecar or the rollouts jsonl, which keeps the artifacts free of a verdict no verifier gave. A rollout that also succeeded is never counted.

nemo_gym.rollout_collection._fill_task_fields(
added_rows: typing.List[typing.Dict[str, typing.Any]],
real_rows: typing.List[typing.Dict[str, typing.Any]],
materialized_fpaths: typing.List[pathlib.Path],
group: typing.Callable[[Mapping[str, Any]], typing.Optional[str]]
) -> None

Give the rows the metric input adds, counted failures and imputed zeros, their task’s dataset fields.

Rows are ordered by repeat, so an added row can come first in its task, and metric hooks read task-level fields such as a subset label or a weight from a task’s first rollout. The fields come from the task’s materialized row, and only those that real rows of the same group carry: the aggregator averages every number it is handed, so a field no real row reports would become a metric of its own. A field the verifier computes is not in the materialized row, so the added row lacks it.

nemo_gym.rollout_collection._has_observation_gap(
result: dict[str, typing.Any],
code: str
) -> bool
nemo_gym.rollout_collection._identity_key(
row: collections.abc.Mapping[str, typing.Any],
group: typing.Callable[[Mapping[str, Any]], typing.Optional[str]]
) -> tuple

A rollout’s key across runs, which each number their tasks from 0.

nemo_gym.rollout_collection._is_collector_key(
key: str
) -> bool

Keys rollout collection writes itself; an Environment Server result must not use them.

nemo_gym.rollout_collection._is_episode_response(
result: typing.Any
) -> bool

True for a BaseEpisodeResponse-shaped reply: object identities plus a result or failure key.

The collector only applies this to a row it dispatched as an episode request (_materialized_taskset(row)), so an agent’s verify response that echoes identity fields is left alone.

nemo_gym.rollout_collection._latest_failure_rows(
failures_fpaths: typing.List[pathlib.Path]
) -> typing.Dict[typing.Tuple[typing.Any, typing.Any], typing.Dict[str, typing.Any]]

The last attempt recorded for each rollout across the failures sidecars.

nemo_gym.rollout_collection._masking_step_metrics(
agent_name: str,
scored: collections.Counter,
dropped: collections.Counter
) -> typing.Dict[str, float]

In-progress view of what a run is losing to its environment rather than its policy.

scored covers persisted rollouts only, split into the unmasked ones (count, reward) and the masked ones; dropped counts what never reached the main output at all. reward_unmasked averages over the unmasked rollouts alone, so the gap against the existing reward series is the score lost to infrastructure. Failed and omitted attempts are reported as counts, never folded into a quality average.

Empty until something is actually masked or dropped, so a healthy run exports exactly what it exported before. The final numbers come from /aggregate_metrics; this is the progress view while the run is still going.

nemo_gym.rollout_collection._materialized_taskset(
row: collections.abc.Mapping[str, typing.Any]
) -> str | None
nemo_gym.rollout_collection._metrics_group(
row: collections.abc.Mapping[str, typing.Any],
servers_by_agent: typing.Callable[[], collections.abc.Mapping[str, list[str]]]
) -> typing.Optional[str]

The environment server a rollout is scored under, as _call_aggregate_metrics groups it.

The row’s stamp, else the one server that fronts its agent, looked up only for an unstamped row: a materialized row routed by its agent has no stamp while its result has one, and both must land in the same group. Runs of one agent behind two servers stay apart this way. An agent behind no server or several stays its own group; _call_aggregate_metrics rejects such a row.

nemo_gym.rollout_collection._missing_rollout_rows_counted_as_zero(
materialized_fpaths: typing.List[pathlib.Path],
failures_fpaths: typing.List[pathlib.Path],
scored_keys: set
) -> typing.List[typing.Dict[str, typing.Any]]

Materialized rollouts that produced no row anywhere, counted as zeros.

The failure classes reach a rollout that failed and said so. This reaches the one that never got that far — killed mid-flight, or dispatched and lost — which leaves nothing in the rollouts jsonl and nothing in the sidecar. Without it such a rollout leaves the denominator as well as the numerator, so the score reads higher the more of the run went missing.

The sidecar is read here rather than trusted from the caller. A failure whose class the caller left out of count_failure_classes_as_zero is absent from the scored keys, and counting it here would score the very rollouts that selection excluded — silently turning the selection into a no-op.

The zero carries the rollout’s identity: its agent, and its environment server stamp when the row has one, which a native taskset row needs because it names no agent. _fill_task_fields adds the task’s dataset fields. A row that names neither cannot reach any server’s metrics, so it is warned about rather than counted. The score enters the metric input and nothing else, the same way a counted failure row does.

nemo_gym.rollout_collection._native_episode_request_body(
row: collections.abc.Mapping[str, typing.Any]
) -> dict[str, typing.Any]
nemo_gym.rollout_collection._nonnegative_int(
value: typing.Any
) -> typing.Optional[int]
nemo_gym.rollout_collection._normalize_health_check_ignored_checks(
value
) -> typing.List[str]
nemo_gym.rollout_collection._read_jsonl(
path: pathlib.Path
) -> typing.List[typing.Dict]
nemo_gym.rollout_collection._rollout_for_export(
result: dict[str, typing.Any]
) -> dict[str, typing.Any]

Return an exporter view without the complete trajectory or raw capture payloads.

nemo_gym.rollout_collection._rollout_order_key(
row: typing.Dict[str, typing.Any]
) -> tuple

Task then repeat: metrics that read a task’s repeats positionally need that order.

nemo_gym.rollout_collection._rollout_request_debug_summary(
row: typing.Dict[str, typing.Any]
) -> typing.Dict[str, typing.Any]
nemo_gym.rollout_collection._routing_identity(
row: collections.abc.Mapping[str, typing.Any]
) -> typing.Optional[str]

The agent a rollout ran on, or for a native taskset row, which names none, its environment server.

A materialized row, its result and its sidecar row all name the same one.

nemo_gym.rollout_collection._strip_capture_payloads(
result: dict[str, typing.Any]
) -> None
nemo_gym.rollout_collection._trajectory_identity(
row: dict[str, typing.Any]
) -> tuple[str, str]
nemo_gym.rollout_collection._truncated_body(
body: typing.Optional[bytes]
) -> typing.Optional[str]

Decode at most the kept prefix, so a huge error page is never decoded in full.

nemo_gym.rollout_collection._turn_content(
request: typing.Any,
response: typing.Any
) -> tuple[typing.Any, typing.Any, typing.Any, int]

Split one captured model call into the turn’s question, answer, reasoning, and tool-call count.

Handles Responses API output items and chat-completions messages, the two dialects the Model Server captures.

nemo_gym.rollout_collection._turns_from_model_calls(
task_id: str,
rollout_id: str,
resolved: typing.Any

Build one turn per captured model call that returned a response, for an agent that sent no trajectory.

A call belongs to the invocation whose model_calls reference it. Unreferenced calls stay in the raw evidence but do not become turns: even a single-agent rollout may capture judge or auxiliary calls that do not belong to the agent.

nemo_gym.rollout_collection.get_max_rollout_attempts() -> int

Read NEMO_GYM_MAX_ROLLOUT_ATTEMPTS (positive int) or default to 3.

nemo_gym.rollout_collection.is_terminal_failure(
record: collections.abc.Mapping[str, typing.Any],
retry_terminal_timeouts: bool = False
) -> bool

Whether a persisted failure must be gated on resume.

By default a row is terminal iff it is stamped _ng_failure_terminal, so an agent that marks its timeouts terminal on purpose keeps them gated.

retry_terminal_timeouts is for agents whose older builds incorrectly stamped per-attempt timeouts terminal. A timeout reflects the load/remaining walltime of that attempt, so it remains retryable (up to the normal max-attempt cap). A skipped sample is unusable regardless of which agent version wrote the sidecar and stays terminal.

The class names below are _ng_failure_class labels that agents and resources servers already write to the failures sidecar. They predate nemo_gym.failure_kinds and are not registered there:

  • timeout_exceeded (Stirrup, pinchbench): the per-task timeout; agent_timeout in the shared vocabulary.
  • reference_missing, eval_missing, transport_ineligible (GDPVal): environment faults with no shared name (namespaced, they would be gdpval:<kind>).
  • skipped (Stirrup): the sample cannot be run; no shared name.

As failure_kinds requires, this function decides retryability for the occurrence, under an explicit caller opt-in; the names carry no retry meaning.

nemo_gym.rollout_collection.loads_jsonl_line(
raw,
fpath,
line_no: int
)

Parse one JSONL line, raising a clean ConfigError (naming file + line) on malformed JSON.

nemo_gym.rollout_collection.migrate_invalid_judge_main_rows(
output_fpath: pathlib.Path
) -> int

Move legacy invalid-judge rows from the main JSONL into the sidecar.

Old builds persisted invalid_judge_response=True as a zero-reward main success, which both contaminated aggregation and gated resume. Sidecar-first migration is idempotent via a marker keyed by stage/task/rollout/attempt; the main file is then atomically rewritten so a crash cannot lose retry state.

Run it only for environments that opt in (retry_invalid_judge_responses): other verifiers score an invalid judge response as a zero-reward row on purpose.

Migrated rows are classed judge_invalid, the sidecar label reverification uses for the same condition (judge_unparseable in nemo_gym.failure_kinds).

nemo_gym.rollout_collection.observed_elapsed(
record: typing.Dict[str, typing.Any]
) -> typing.Optional[float]

Best-effort per-rollout wallclock from a result/failure row.

nemo_gym.rollout_collection.AGENT_REQUEST_FAILED_FAILURE_CLASS = 'agent_request_failed'
nemo_gym.rollout_collection.AGENT_RUN_ERROR_FAILURE_CLASS = 'agent_run_error'
nemo_gym.rollout_collection.ENVIRONMENT_SERVER_FAILURE_CLASS = 'environment_server_failed'
nemo_gym.rollout_collection.NG_DISPATCH_DRAINED_KEY = '_ng_dispatch_drained'
nemo_gym.rollout_collection.NG_ELAPSED_KEY = 'elapsed_seconds'
nemo_gym.rollout_collection.NG_ENVIRONMENT_SERVER_KEY = ENVIRONMENT_SERVER_STAMP_KEY_NAME
nemo_gym.rollout_collection.NG_FAILURE_CLASS_KEY = '_ng_failure_class'
nemo_gym.rollout_collection.NG_NO_PERSIST_KEY = '_ng_no_persist'
nemo_gym.rollout_collection.NG_PERF_KEY = 'ng_perf'
nemo_gym.rollout_collection.NG_RESULT_TYPE_KEY = '_ng_result_type'
nemo_gym.rollout_collection.NG_TASK_ID_KEY = '_ng_task_id'
nemo_gym.rollout_collection.NG_TERMINAL_KEY = '_ng_failure_terminal'
nemo_gym.rollout_collection.NG_TRAJECTORY_KEY = 'ng_trajectory'
nemo_gym.rollout_collection._DEFAULT_MAX_ROLLOUT_ATTEMPTS = 3
nemo_gym.rollout_collection._MAX_FAILURE_BODY_CHARS = 2000
nemo_gym.rollout_collection._MIGRATED_INVALID_JUDGE_KEY = '_ng_migrated_invalid_judge_response'
nemo_gym.rollout_collection._MODEL_CALL_PAYLOAD_KEYS = ('request', 'response', 'request_raw', 'response_raw')
nemo_gym.rollout_collection._NO_RESULT_FAILURE_CLASSES = frozenset({AGENT_REQUEST_FAILED_FAILURE_CLASS, AGENT_RUN_ERROR_FAILURE_CLASS, EN...
nemo_gym.rollout_collection._RUN_FAILURE_ERRORS = (ClientError, orjson.JSONDecodeError, TimeoutError)
nemo_gym.rollout_collection._SERVER_DID_NOT_RUN_STATUSES = frozenset({429, 502, 503, 504})
nemo_gym.rollout_collection.__getattr__ = moved_attr_getter(__name__, {'collect_rollouts': 'nemo_gym.cli.eval', 'aggregate...
nemo_gym.rollout_collection._failures_path_for = failures_path_for
nemo_gym.rollout_collection._get_max_rollout_attempts = get_max_rollout_attempts
nemo_gym.rollout_collection.logger = logging.getLogger(__name__)