(self, storage_type, is_eval, offset)
| 27 | ] |
| 28 | ) |
| 29 | async def test_read_task(self, storage_type, is_eval, offset): |
| 30 | config = get_template_config() |
| 31 | total_samples = 17 |
| 32 | batch_size = 4 |
| 33 | config.buffer.explorer_input.taskset = get_unittest_dataset_config( |
| 34 | "countdown" |
| 35 | ) # 17 samples |
| 36 | config.buffer.explorer_input.taskset.storage_type = storage_type |
| 37 | config.buffer.explorer_input.taskset.is_eval = is_eval |
| 38 | config.buffer.explorer_input.taskset.index = offset |
| 39 | config.buffer.explorer_input.taskset.batch_size = batch_size |
| 40 | if storage_type == StorageType.SQL.value: |
| 41 | dataset = datasets.load_dataset( |
| 42 | config.buffer.explorer_input.taskset.path, split="train" |
| 43 | ) |
| 44 | config.buffer.explorer_input.taskset.path = f"sqlite:///{db_path}" |
| 45 | await SQLTaskStorage.load_from_dataset( |
| 46 | dataset, config.buffer.explorer_input.taskset.to_storage_config() |
| 47 | ) |
| 48 | reader = get_buffer_reader(config.buffer.explorer_input.taskset) |
| 49 | tasks = [] |
| 50 | while True: |
| 51 | try: |
| 52 | cur_tasks = await reader.read() |
| 53 | tasks.extend(cur_tasks) |
| 54 | except StopAsyncIteration: |
| 55 | break |
| 56 | if is_eval: |
| 57 | self.assertEqual(len(tasks), total_samples - offset) |
| 58 | else: |
| 59 | self.assertEqual(len(tasks), (total_samples - offset) // batch_size * batch_size) |
| 60 | |
| 61 | def setUp(self): |
| 62 | if os.path.exists(db_path): |
nothing calls this directly
no test coverage detected