Skip to content
Open
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
10 changes: 6 additions & 4 deletions miles/backends/megatron_utils/actor.py
Original file line number Diff line number Diff line change
Expand Up @@ -616,10 +616,12 @@ def save_model(self, rollout_id: int, force_sync: bool = False) -> None:

save_hf_model(self.args, rollout_id, self.model)

if self.args.custom_megatron_post_save_hook_path is not None and dist.get_rank() == 0:
if self.args.async_save:
maybe_finalize_async_save(blocking=True)
post_save_hook_path = self.args.custom_megatron_post_save_hook_path
if post_save_hook_path is not None and self.args.async_save:
# Distributed checkpoint finalization must run on every rank.
maybe_finalize_async_save(blocking=True)

if post_save_hook_path is not None and dist.get_rank() == 0:
from megatron.training.checkpointing import get_checkpoint_name

from miles.utils.misc import load_function
Expand All @@ -630,7 +632,7 @@ def save_model(self, rollout_id: int, force_sync: bool = False) -> None:
if self.args.save_hf is not None and self.role == "actor"
else None
)
post_save_hook = load_function(self.args.custom_megatron_post_save_hook_path)
post_save_hook = load_function(post_save_hook_path)
post_save_hook(self.args, rollout_id, checkpoint_dir, hf_checkpoint_dir)

if self.args.offload_train:
Expand Down
48 changes: 48 additions & 0 deletions tests/fast/backends/megatron_utils/test_actor_checkpoint.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
from argparse import Namespace
from unittest.mock import MagicMock, patch

import pytest

from miles.backends.megatron_utils.actor import MegatronTrainRayActor


@pytest.mark.parametrize(
("rank", "expected_events"),
[
(0, ["finalize", "save", "finalize", "hook"]),
(1, ["finalize", "save", "finalize"]),
],
)
def test_async_post_save_hook_finalizes_on_every_rank(rank: int, expected_events: list[str]) -> None:
events: list[str] = []
actor = MegatronTrainRayActor.__new__(MegatronTrainRayActor)
actor._heartbeat = MagicMock()
actor.args = Namespace(
async_save=True,
custom_megatron_post_save_hook_path="test_checkpoint.post_save_hook",
debug_rollout_only=False,
offload_train=False,
save="/checkpoints",
save_hf=None,
)
actor.model = MagicMock()
actor.optimizer = MagicMock()
actor.opt_param_scheduler = MagicMock()
actor.role = "actor"

def post_save_hook(*_args: object) -> None:
events.append("hook")

with (
patch("miles.backends.megatron_utils.actor.is_multi_lora_enabled", return_value=False),
patch("miles.backends.megatron_utils.actor.save", side_effect=lambda *_args: events.append("save")),
patch("miles.backends.megatron_utils.actor.dist.get_rank", return_value=rank),
patch(
"megatron.training.async_utils.maybe_finalize_async_save",
side_effect=lambda **_kwargs: events.append("finalize"),
),
patch("miles.utils.misc.load_function", return_value=post_save_hook),
):
actor.save_model(rollout_id=7)

assert events == expected_events
Loading