已合并
test(distributed):add test for validate_checkpoint_id, reset, set_up_storage_reader from BroadcastingTorchSaveReader #33733
test(distributed):add test for validate_checkpoint_id, reset, set_up_storage_reader from BroadcastingTorchSaveReader #33733
已合并
xh-zhan创建于 4月14日
共 1 个文件变更+104-0
@@ -0,0 +1,104 @@
1+"""
2+Add validation cases for BroadcastingTorchSaveReader APIs:
3+1. pytorch/test/distributed/checkpoint/test_format_utils.py from PyTorch community lacks sufficient API validations, so this file is added.
4+2. This file validates the following APIs:
5+ torch.distributed.checkpoint.format_utils.BroadcastingTorchSaveReader.reset
6+ torch.distributed.checkpoint.format_utils.BroadcastingTorchSaveReader.set_up_storage_reader
7+ torch.distributed.checkpoint.format_utils.BroadcastingTorchSaveReader.validate_checkpoint_id
8+ (extendable)
9+"""
10+ 
11+import os
12+import tempfile
13+from unittest.mock import patch
14+ 
15+import torch
16+from torch.distributed.checkpoint.format_utils import BroadcastingTorchSaveReader
17+from torch.testing._internal.common_utils import TestCase, run_tests
18+ 
19+ 
20+class TestBroadcastingTorchSaveReader(TestCase):
21+ """Independent unit tests for pure logic APIs of BroadcastingTorchSaveReader.
22+ 
23+ This test uses mocks for distributed functions, allowing full validation
24+ of all APIs without starting a process group.
25+ """
26+ 
27+ def setUp(self):
28+ """Runs before each test method: creates a valid temporary checkpoint file."""
29+ self.temp_file = tempfile.NamedTemporaryFile(delete=False, suffix=".pt")
30+ torch.save({"dummy": torch.tensor([1, 2, 3])}, self.temp_file.name)
31+ self.temp_file.close()
32+ self.valid_checkpoint_id = self.temp_file.name
33+ 
34+ def tearDown(self):
35+ """Runs after each test method: deletes the temporary file to clean up."""
36+ if os.path.exists(self.valid_checkpoint_id):
37+ os.unlink(self.valid_checkpoint_id)
38+ 
39+ def test_validate_checkpoint_id_valid_path(self):
40+ """Passing a path to an existing file should return True."""
41+ result = BroadcastingTorchSaveReader.validate_checkpoint_id(self.valid_checkpoint_id)
42+ self.assertTrue(result, "valid checkpoint file should return True")
43+ 
44+ def test_validate_checkpoint_id_invalid_path(self):
45+ """Passing a path to a non-existent file should return False."""
46+ non_existent = "/tmp/does_not_exist_12345.pt"
47+ result = BroadcastingTorchSaveReader.validate_checkpoint_id(non_existent)
48+ self.assertFalse(result, "non-existent file should return False")
49+ 
50+ def test_reset_updates_checkpoint_id(self):
51+ """Verify that the reset method correctly updates the internal checkpoint_id."""
52+ reader = BroadcastingTorchSaveReader(checkpoint_id="/old/path.pt")
53+ new_path = "/new/path.pt"
54+ reader.reset(new_path)
55+ self.assertEqual(reader.checkpoint_id, new_path, "reset should update checkpoint_id")
56+ 
57+ def test_reset_none(self):
58+ """Verify that reset with None correctly sets checkpoint_id to None."""
59+ reader = BroadcastingTorchSaveReader(checkpoint_id="abc")
60+ reader.reset(None)
61+ self.assertIsNone(reader.checkpoint_id)
62+ 
63+ class DummyMetadata:
64+ """Placeholder Metadata class used for testing."""
65+ pass
66+ 
67+ @patch("torch.distributed.get_rank")
68+ def test_set_up_storage_reader_sets_coordinator_flag(self, mock_get_rank):
69+ """Verify that set_up_storage_reader correctly sets the is_coordinator attribute."""
70+ mock_get_rank.return_value = 0
71+ 
72+ reader = BroadcastingTorchSaveReader(checkpoint_id=self.valid_checkpoint_id)
73+ metadata = self.DummyMetadata()
74+ 
75+ reader.set_up_storage_reader(metadata, is_coordinator=True)
76+ self.assertTrue(reader.is_coordinator, "is_coordinator should be set to True")
77+ 
78+ reader.set_up_storage_reader(metadata, is_coordinator=False)
79+ self.assertFalse(reader.is_coordinator, "is_coordinator should be set to False")
80+ 
81+ @patch("torch.distributed.get_rank")
82+ def test_set_up_storage_reader_coordinator_rank_mismatch_raises(self, mock_get_rank):
83+ """Verify that an AssertionError is raised when is_coordinator=True
84+ but the current rank does not match coordinator_rank."""
85+ mock_get_rank.return_value = 1
86+ 
87+ reader = BroadcastingTorchSaveReader(checkpoint_id=self.valid_checkpoint_id, coordinator_rank=0)
88+ metadata = self.DummyMetadata()
89+ 
90+ with self.assertRaises(AssertionError, msg="Should raise AssertionError when coordinator rank mismatches"):
91+ reader.set_up_storage_reader(metadata, is_coordinator=True)
92+ 
93+ def test_set_up_storage_reader_requires_checkpoint_id(self):
94+ """Verify that an AssertionError is raised when checkpoint_id is None."""
95+ reader = BroadcastingTorchSaveReader(checkpoint_id=None)
96+ metadata = self.DummyMetadata()
97+ 
98+ with patch("torch.distributed.get_rank", return_value=0):
99+ with self.assertRaises(AssertionError, msg="checkpoint_id not set should raise AssertionError"):
100+ reader.set_up_storage_reader(metadata, is_coordinator=True)
101+ 
102+ 
103+if __name__ == "__main__":
104+ run_tests()