nemo_gym.rollout_collection
nemo_gym.rollout_collection
Module Contents
Classes
Functions
Data
AGENT_REQUEST_FAILED_FAILURE_CLASS
ENVIRONMENT_SERVER_FAILURE_CLASS
API
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.
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.
Bases: SharedRolloutCollectionConfig
Spin up all necessary servers and perform a batch of rollout collection using each dataset inside the provided configs.
Examples:
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:
Bases: BaseModel
Bases: SharedRolloutCollectionConfig
Perform a batch of rollout collection.
Examples:
Bases: BaseModel
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.
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.
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.
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.
Fail before dispatch when a row points to agent incompatible with the resources server it runs on.
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.
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.
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.
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.
Bases: UploadRolloutsConfigMixin, BaseNeMoGymCLIConfig
Reject incomplete submitted runs after saving their partial artifacts.
Completion-order iterator with a bounded set of resident asyncio tasks.
A finished /run dispatch, with timing carried alongside (not inside) the raw result.
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.
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.
The metric-input copy of a sidecar row counted as zero.
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.
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.
Repair a jsonl whose last line a hard kill cut short, so resume can read and append to it.
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.
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.
Map each agent name to the environment servers whose agent_server names it.
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.
Expand a glob-or-comma-separated-globs string into a sorted, deduplicated list of paths.
Examples:
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.
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.
A rollout’s key across runs, which each number their tasks from 0.
Keys rollout collection writes itself; an Environment Server result must not use them.
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.
The last attempt recorded for each rollout across the failures sidecars.
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.
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.
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.
Return an exporter view without the complete trajectory or raw capture payloads.
Task then repeat: metrics that read a task’s repeats positionally need that order.
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.
Decode at most the kept prefix, so a huge error page is never decoded in full.
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.
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.
Read NEMO_GYM_MAX_ROLLOUT_ATTEMPTS (positive int) or default to 3.
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_timeoutin the shared vocabulary.reference_missing,eval_missing,transport_ineligible(GDPVal): environment faults with no shared name (namespaced, they would begdpval:<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.
Parse one JSONL line, raising a clean ConfigError (naming file + line) on malformed JSON.
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).
Best-effort per-rollout wallclock from a result/failure row.