nemo_gym.batch_status

View as Markdown

Validated Eval Factory batch manifests and compact, atomic Gym status artifacts.

Module Contents

Classes

NameDescription
BatchManifestVersioned, secret-free handoff from Eval Factory to Gym.
BatchManifestMemberEval Factory’s expectations for one logical benchmark in a shared Gym run.
BatchRepeatPolicyThe repeat shape Gym must observe for every task belonging to one member.
BatchStatusTrackerValidate one batch manifest and materialize its current per-agent status.
ObservedBatchMemberDataset and repeat facts derived from materialized Gym input rows.

Functions

NameDescription
_atomic_write_json-
_canonical_task_bytesCanonicalize task content independently of routing, repeats, retries, and skills.
_read_aggregate_entries-
_safe_aggregation_errorCopy only non-secret diagnostics into batch_status.json.
_seed_from_row-
load_batch_manifestLoad a manifest and return it with the SHA-256 of the exact input bytes.
observe_materialized_rowsDerive per-agent dataset fingerprints and repeat policies from full materialized inputs.
validate_batch_manifestReject any manifest expectation that disagrees with Gym’s materialized inputs.

Data

AGGREGATION_ERROR_KEY

BATCH_STATUS_FNAME

BATCH_STATUS_WRITE_INTERVAL_SECONDS

_DATASET_DIGEST_DOMAIN

_FAILURE_CLASS_KEY

_HEX_SHA256

_MISSING

API

class nemo_gym.batch_status.BatchManifest()

Bases: BaseModel

Versioned, secret-free handoff from Eval Factory to Gym.

members
Dict[str, BatchManifestMember] = Field(min_length=1)
model_config
= ConfigDict(extra='forbid')
schema_version
Literal['1'] = '1'
class nemo_gym.batch_status.BatchManifestMember()

Bases: BaseModel

Eval Factory’s expectations for one logical benchmark in a shared Gym run.

agent_name
str
dataset_sha256
str
expected_rollout_count
int = Field(ge=1)
expected_task_count
int = Field(ge=1)
metric_keys
List[str]
model_config
= ConfigDict(extra='forbid')
repeat_policy
BatchRepeatPolicy
resolved_recipe_sha256
str
task_sources
List[str] = Field(min_length=1)
nemo_gym.batch_status.BatchManifestMember._validate_agent_name(
value: str
) -> str
classmethod
nemo_gym.batch_status.BatchManifestMember._validate_sha256(
value: str
) -> str
classmethod
nemo_gym.batch_status.BatchManifestMember._validate_unique_names(
values: typing.List[str]
) -> typing.List[str]
classmethod
class nemo_gym.batch_status.BatchRepeatPolicy()

Bases: BaseModel

The repeat shape Gym must observe for every task belonging to one member.

model_config
= ConfigDict(extra='forbid')
num_repeats
int = Field(ge=1)
seeded
bool
class nemo_gym.batch_status.BatchStatusTracker(
manifest_fpath: pathlib.Path,
materialized_rows: typing.Sequence[typing.Mapping[str, typing.Any]],
write_interval_seconds: typing.Optional[float] = None
)

Validate one batch manifest and materialize its current per-agent status.

_expected_rollout_identities
set[tuple[str, int, int]]
_last_progress_count
= 0
_last_write
= 0.0
observations
= observe_materialized_rows(materialized_rows)
status_fpath
= manifest_fpath.with_name(BATCH_STATUS_FNAME)
write_interval_seconds
nemo_gym.batch_status.BatchStatusTracker._validated_rollout_identity(
row: typing.Mapping[str, object],
record_kind: str,
row_index: int
) -> tuple[str, int, int]

Reject malformed IDs and records belonging to a different materialized batch.

nemo_gym.batch_status.BatchStatusTracker.build_status(
completed_results: typing.Sequence[typing.Mapping[str, typing.Any]],
failure_rows: typing.Sequence[typing.Mapping[str, typing.Any]],
aggregate_metrics_fpath: typing.Optional[pathlib.Path] = None,
aggregation_deferred: bool = False
) -> typing.Dict[str, typing.Any]

Build status from expected identities, raising ConfigError for mismatched result records.

nemo_gym.batch_status.BatchStatusTracker.write_status(
completed_results: typing.Sequence[typing.Mapping[str, typing.Any]],
failure_rows: typing.Sequence[typing.Mapping[str, typing.Any]],
aggregate_metrics_fpath: typing.Optional[pathlib.Path] = None,
aggregation_deferred: bool = False,
force: bool = False
) -> bool

Atomically refresh status, throttling progress writes unless force is true.

class nemo_gym.batch_status.ObservedBatchMember(
task_sources: typing.List[str],
dataset_sha256: str,
task_count: int,
rollout_count: int,
)
Dataclass

Dataset and repeat facts derived from materialized Gym input rows.

dataset_sha256
str
repeat_policy
BatchRepeatPolicy
rollout_count
int
task_count
int
task_sources
List[str]
nemo_gym.batch_status._atomic_write_json(
path: pathlib.Path,
payload: typing.Mapping[str, typing.Any]
) -> None
nemo_gym.batch_status._canonical_task_bytes(
row: typing.Mapping[str, typing.Any]
) -> bytes

Canonicalize task content independently of routing, repeats, retries, and skills.

nemo_gym.batch_status._read_aggregate_entries(
metrics_fpath: typing.Optional[pathlib.Path]
) -> tuple[typing.Dict[str, typing.Dict[str, typing.Any]], typing.Optional[str]]
nemo_gym.batch_status._safe_aggregation_error(
error: typing.Any
) -> typing.Dict[str, typing.Any]

Copy only non-secret diagnostics into batch_status.json.

nemo_gym.batch_status._seed_from_row(
row: typing.Mapping[str, typing.Any]
) -> typing.Any
nemo_gym.batch_status.load_batch_manifest(
manifest_fpath: pathlib.Path

Load a manifest and return it with the SHA-256 of the exact input bytes.

nemo_gym.batch_status.observe_materialized_rows(
rows: typing.Sequence[typing.Mapping[str, typing.Any]]

Derive per-agent dataset fingerprints and repeat policies from full materialized inputs.

The dataset digest hashes one canonical row per task in its materialized order. Runtime routing and identity fields, repeat seeds, retry attempts, and run-level skills are excluded; prompt-expanded inputs and task_source remain part of the dataset identity.

nemo_gym.batch_status.validate_batch_manifest(
observations: typing.Mapping[str, nemo_gym.batch_status.ObservedBatchMember]
) -> None

Reject any manifest expectation that disagrees with Gym’s materialized inputs.

nemo_gym.batch_status.AGGREGATION_ERROR_KEY = 'aggregation_error'
nemo_gym.batch_status.BATCH_STATUS_FNAME = 'batch_status.json'
nemo_gym.batch_status.BATCH_STATUS_WRITE_INTERVAL_SECONDS = 5.0
nemo_gym.batch_status._DATASET_DIGEST_DOMAIN = b'nemo-gym-batch-dataset-v1\x00'
nemo_gym.batch_status._FAILURE_CLASS_KEY = '_ng_failure_class'
nemo_gym.batch_status._HEX_SHA256 = re.compile('^[0-9a-f]{64}$')
nemo_gym.batch_status._MISSING = object()