(self)
| 163 | port=port) |
| 164 | |
| 165 | async def init_workers_async(self): |
| 166 | self.create_workers(RayGPUWorker, self.worker_kwargs) |
| 167 | try: |
| 168 | await asyncio.gather(*self._get_worker_ready_futures()) |
| 169 | except ray.exceptions.ActorDiedError as e: |
| 170 | raise RuntimeError("RayGPUWorker died during initialization") from e |
| 171 | port = (await asyncio.gather(*self.call_all_ray_workers( |
| 172 | "setup_tcp_store", leader_only=True, async_call=True)))[0] |
| 173 | await asyncio.gather( |
| 174 | *self.call_all_ray_workers("setup_distributed_env_and_worker", |
| 175 | leader_only=False, |
| 176 | async_call=True, |
| 177 | port=port)) |
| 178 | |
| 179 | @unwrap_ray_errors() |
| 180 | def call_all_ray_workers(self, func: str, leader_only: bool, |
no test coverage detected