(self, func: str, leader_only: bool,
async_call: bool, *args, **kwargs)
| 178 | |
| 179 | @unwrap_ray_errors() |
| 180 | def call_all_ray_workers(self, func: str, leader_only: bool, |
| 181 | async_call: bool, *args, **kwargs): |
| 182 | workers = (self.workers[0], ) if leader_only else self.workers |
| 183 | if async_call: |
| 184 | return [ |
| 185 | getattr(worker, func).remote(*args, **kwargs) |
| 186 | for worker in workers |
| 187 | ] |
| 188 | else: |
| 189 | return ray.get([ |
| 190 | getattr(worker, func).remote(*args, **kwargs) |
| 191 | for worker in workers |
| 192 | ]) |
| 193 | |
| 194 | @unwrap_ray_errors() |
| 195 | def collective_rpc(self, |
no test coverage detected