Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
158 changes: 157 additions & 1 deletion tests/examples/deepswe_mllog_utils_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
import shutil
import tempfile
import types
from typing import Any
from unittest import mock

import numpy as np
Expand All @@ -30,6 +31,15 @@
from tunix.utils import mllog_utils


def _read_mllog_events(path: str) -> list[dict[str, Any]]:
with open(path, "r", encoding="utf-8") as f:
return [
json.loads(line.split(":::MLLOG ", 1)[1])
for line in f
if ":::MLLOG " in line
]


@absltest.skipIf(mllog_utils.mllogger is None, "mlperf_logging is not installed")
class MllogUtilsTest(absltest.TestCase):

Expand Down Expand Up @@ -64,6 +74,7 @@ def test_file_logging_with_metric_logger_dir(self):
tpu_topology="v5p-64",
rollout_engine="vllm",
target_accuracy=0.69,
model_id="",
)

mllog_utils.init_start(args)
Expand Down Expand Up @@ -221,6 +232,7 @@ def test_end_to_end_mlperf_logging_with_train_configs(self):
seed=1,
learning_rate=1e-6,
eval_every_n_steps=5,
model_id="",
)

mock_train_dataset = [None] * 5480
Expand Down Expand Up @@ -602,7 +614,151 @@ def test_train_stop(self):
content = f.read()

self.assertIn('"key": "block_stop"', content)
self.assertIn('"key": "run_stop"', content)
# run_stop is emitted by the offline evaluator, not by train_stop.
self.assertNotIn('"key": "run_stop"', content)

def test_compute_val_start_step(self):
self.assertEqual(mllog_utils.compute_val_start_step(256), 18)
self.assertEqual(mllog_utils.compute_val_start_step(512), 10)
self.assertEqual(mllog_utils.compute_val_start_step(1024), 7)
self.assertEqual(mllog_utils.compute_val_start_step(256, 1), 1)
self.assertEqual(mllog_utils.compute_val_start_step(256, 0), 18)
with self.assertRaises(ValueError):
mllog_utils.compute_val_start_step(0)

def test_append_checkpoint_manifest_upserts_sorted_records(self):
manifest_path = os.path.join(
self.test_dir, "mllog", "eval_checkpoints.jsonl"
)
for step, ts_ms in ((19, 2000), (18, 1000), (18, 1500)):
mllog_utils.append_checkpoint_manifest(
manifest_path,
{
"step": step,
"checkpoint_path": f"gs://ckpt/{step}/model_params",
"timestamp_ms": ts_ms,
"samples_count": step * 256,
"val_start_at": 18,
},
)
with open(manifest_path, "r", encoding="utf-8") as f:
records = [json.loads(line) for line in f]
self.assertEqual([r["step"] for r in records], [18, 19])
self.assertEqual([r["timestamp_ms"] for r in records], [1500, 2000])

def test_configure_logger_downloads_existing_gcs_log_before_config(self):
calls = mock.MagicMock()
fake_mllogger = mock.MagicMock()
fake_mllogger.logger.handlers = []
with (
mock.patch.object(mllog_utils, "mllog", calls.mllog),
mock.patch.object(mllog_utils, "mllogger", fake_mllogger),
mock.patch.object(mllog_utils, "_is_master_process", return_value=True),
mock.patch.object(mllog_utils, "_gcs_target_path", None),
mock.patch.object(mllog_utils, "_local_log_path", None),
mock.patch.object(
mllog_utils, "_download_from_gcs_if_exists", calls.download
),
):
mllog_utils.configure_logger(metric_logger_dir="gs://b/mllog", seed=42)
local_path = mllog_utils._local_log_path # pylint: disable=protected-access

self.assertEqual(
[c[0] for c in calls.mock_calls], ["download", "mllog.config"]
)
calls.download.assert_called_once_with("gs://b/mllog/seed_42.out", local_path)
self.assertEqual(
os.path.abspath(calls.mllog.config.call_args.kwargs["filename"]),
local_path,
)

def test_offline_eval_rcp_sequence_converged(self):
log_dir = self.test_dir
args = types.SimpleNamespace(batch_size=16, num_generations=16)
mllog_utils.configure_logger(metric_logger_dir=log_dir, seed=42)
mllog_utils.train_stop(args, step=19, time_ms=2000)
self.assertFalse(
mllog_utils.log_offline_eval_step(
step=18,
samples_count=4608,
eval_accuracy=0.65,
checkpoint_timestamp_ms=1000,
)
)
self.assertTrue(
mllog_utils.log_offline_eval_step(
step=19,
samples_count=4864,
eval_accuracy=0.70,
checkpoint_timestamp_ms=2000,
is_last_checkpoint=True,
)
)

events = _read_mllog_events(os.path.join(log_dir, "seed_42.out"))
self.assertEqual(
[e["key"] for e in events],
["block_stop"] + ["eval_start", "eval_accuracy", "eval_stop"] * 2
+ ["run_stop"],
)
self.assertEqual(events[0]["time_ms"], 2000)
self.assertEqual(events[0]["metadata"]["step"], 19)
self.assertEqual([events[2]["value"], events[5]["value"]], [0.65, 0.70])
run_stop = events[-1]
self.assertEqual(run_stop["time_ms"], 2000)
self.assertEqual(run_stop["metadata"]["status"], "success")
self.assertEqual(run_stop["metadata"]["samples_count"], 4864)

def test_offline_eval_rcp_sequence_not_converged(self):
log_dir = self.test_dir
mllog_utils.configure_logger(metric_logger_dir=log_dir, seed=42)
for step, ts_ms in ((18, 1000), (19, 2000), (20, 3000)):
self.assertFalse(
mllog_utils.log_offline_eval_step(
step=step,
samples_count=step * 256,
eval_accuracy=0.5,
checkpoint_timestamp_ms=ts_ms,
is_last_checkpoint=step == 20,
)
)

events = _read_mllog_events(os.path.join(log_dir, "seed_42.out"))
run_stops = [e for e in events if e["key"] == "run_stop"]
self.assertLen(run_stops, 1)
self.assertEqual(events[-1]["key"], "run_stop")
self.assertEqual(run_stops[0]["time_ms"], 3000)
self.assertEqual(run_stops[0]["metadata"]["status"], "aborted")
self.assertEqual(run_stops[0]["metadata"]["samples_count"], 5120)

def test_mlperf_6_1_0_init_print_disclosures(self):
fake_mllogger = mock.MagicMock()
with (
mock.patch.object(mllog_utils, "mllogger", fake_mllogger),
mock.patch.object(mllog_utils, "_is_master_process", return_value=True),
mock.patch.object(mllog_utils, "_flush_to_gcs_if_needed"),
):
args = types.SimpleNamespace(
batch_size=16,
num_generations=16,
max_prompt_length=4096,
max_response_length=61440,
model_id="",
)
mllog_utils.init_print(args)
emitted = {
c.kwargs["key"]: c.kwargs["value"]
for c in fake_mllogger.event.call_args_list
}
self.assertEqual(emitted["eval_samples"], 251)
self.assertEqual(emitted["max_sequence_length"], 65536)
for key in (
"lowest_numerical_precision_in_linear",
"lowest_numerical_precision_in_attn",
"lowest_numerical_precision_in_comm",
):
self.assertEqual(emitted[key], "bfloat16")
self.assertEqual(emitted["config_filename"], "qwen35_397b_grpo")


if __name__ == "__main__":
Expand Down
30 changes: 30 additions & 0 deletions tests/experimental/examples/deepswe_dist/eval_deepswe_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -168,6 +168,21 @@ def test_missing_and_failed_attempts_are_not_dropped(self):
with self.assertRaisesRegex(ValueError, "duplicate"):
eval_lib.summarize(rows + [rows[0]], ["a", "b"], 4)

def test_pass_at_4_is_reported_with_more_attempts(self):
rows = [
dict(
instance_id="a",
attempt=i,
reward=float(i == 5),
resolved=i == 5,
status="SUCCEEDED",
)
for i in range(8)
]
summary = eval_lib.summarize(rows, ["a"], 8)
# Unbiased pass@k = 1 - C(n - c, k) / C(n, k) with n=8, c=1.
self.assertEqual(summary["pass_at_k"], {"1": 1 / 8, "4": 1 / 2, "8": 1.0})

def test_reward_comes_from_trajectory_not_completion_status(self):
response = types.SimpleNamespace(
error=None,
Expand All @@ -182,6 +197,21 @@ def test_reward_comes_from_trajectory_not_completion_status(self):
response.error = "infrastructure failure"
self.assertFalse(eval_lib.compact_result(response)["resolved"])

def test_rcp_logging_fails_fast_without_mlperf_logging(self):
from tunix.utils import mllog_utils # pylint: disable=g-import-not-at-top

log_file = "gs://bucket/mllog/seed_42.out"
a = self.args("--rcp_logging=true", "--metric_logger_dir", log_file)
with mock.patch.object(mllog_utils, "configure_logger") as configure:
with mock.patch.object(mllog_utils, "mllogger", None):
with self.assertRaisesRegex(RuntimeError, "mlperf_logging"):
eval_lib.setup_rcp_logging(a)
configure.assert_not_called()
with mock.patch.object(mllog_utils, "mllogger", object()):
eval_lib.setup_rcp_logging(a)
eval_lib.setup_rcp_logging(self.args("--rcp_logging=false"))
configure.assert_called_once_with(metric_logger_dir=log_file, seed=42)

def test_individual_results_are_persisted(self):
with tempfile.TemporaryDirectory() as directory:
writer = eval_lib.ResultWriter(directory)
Expand Down
42 changes: 42 additions & 0 deletions tests/experimental/orchestrator/rl_program_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -1480,6 +1480,48 @@ async def _run():

asyncio.run(_run())

def test_train_stage_on_checkpoint_saved_gets_path_from_response(self):
async def _run():
self.mock_algo.num_generations = 1
self.mock_algo.mini_batch_size = 1
self.mock_engine.save_checkpoint.return_value = datatypes.Response(
metadata={
"checkpoint_saved": True,
"checkpoint_path": "gs://ckpt/1/model_params",
}
)
saved = []
program = self._create_program(
batch_size=1, on_checkpoint_saved=saved.append
)
program.engine = self.mock_engine

payload = datatypes.RLTrainerPayload(
prompt_ids=np.array([1, 2], dtype=np.int32),
prompt_mask=np.array([1.0, 1.0], dtype=np.float32),
completion_ids=np.array([3, 4], dtype=np.int32),
completion_mask=np.array([1.0, 1.0], dtype=np.float32),
advantages=np.array([1.0, 1.0], dtype=np.float32),
)
item = datatypes.TrajectoryItem(
group_index=0,
prompt_id="prompt_0",
start_step=0,
traj={"trajectory_reward": 1.0},
)
item.payload = payload
await program.scored_q.put(item)
await program.scored_q.close()

await program.train_stage()

self.assertLen(saved, 1)
self.assertEqual(saved[0]["step"], 1)
self.assertEqual(saved[0]["checkpoint_path"], "gs://ckpt/1/model_params")
self.assertEqual(saved[0]["timestamp_ms"], program.last_step_timestamp_ms)

asyncio.run(_run())

def test_train_stage_sequence_packed_final_batch_broken_down_into_multiple_microbatches(
self,
):
Expand Down
18 changes: 18 additions & 0 deletions tests/experimental/worker/trainer_worker_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,11 @@ def __init__(self):
self.step_count = 10
self.target_state = None
self.gen_model_input_fn = None
self._checkpoint_dir = None

@property
def checkpoint_dir(self) -> str | None:
return self._checkpoint_dir

def compile(self, dummy_data=None):
pass
Expand Down Expand Up @@ -152,6 +157,19 @@ def test_update_returns_step_count(self):
step = self.worker.update()
self.assertEqual(step, 11)

def test_save_checkpoint_returns_checkpoint_path_from_checkpoint_dir(self):
resp_empty = self.worker.save_checkpoint(metadata={"step": 5})
self.assertTrue(resp_empty.metadata["checkpoint_saved"])
self.assertEqual(resp_empty.metadata["checkpoint_path"], "")

self.fake_trainer._checkpoint_dir = "gs://bucket/checkpoints"
resp = self.worker.save_checkpoint(metadata={"step": 5})
self.assertTrue(resp.metadata["checkpoint_saved"])
self.assertEqual(
resp.metadata["checkpoint_path"],
"gs://bucket/checkpoints/5/model_params",
)

def test_set_target_state_configures_trainer(self):
target_state = {"params": np.zeros((4, 4))}
resp = self.worker.set_target_state(target_state=target_state)
Expand Down
2 changes: 2 additions & 0 deletions tests/utils/maxtext_utils_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -913,6 +913,7 @@ def _init_state(self):
mock_cfg = mock.MagicMock()
mock_cfg.weight_dtype = "bfloat16"
mock_cfg.float32_gate_logits = True
mock_cfg.checkpoint_dir = "/tmp/ckpts"
mock_mesh = mock.MagicMock()

with mock.patch.object(
Expand All @@ -927,6 +928,7 @@ def _init_state(self):
self.assertIsNot(captured["optimizer_cls_during_init"], FakeOptimizer)
self.assertEqual(FakeNNX.Optimizer, FakeOptimizer)
self.assertIsNone(engine._weight_converter._direct.target_dtype)
self.assertEqual(engine.checkpoint_dir, "/tmp/ckpts")


if __name__ == "__main__":
Expand Down
Loading
Loading