(self, fn, args_kwargs, threads=10, timeout=300, global_kwargs=None)
| 564 | self.log.verbose(f"{self.name}: finished shutting down") |
| 565 | |
| 566 | async def task_pool(self, fn, args_kwargs, threads=10, timeout=300, global_kwargs=None): |
| 567 | if global_kwargs is None: |
| 568 | global_kwargs = {} |
| 569 | |
| 570 | tasks = {} |
| 571 | args_kwargs = list(args_kwargs) |
| 572 | |
| 573 | def new_task(): |
| 574 | if args_kwargs: |
| 575 | kwargs = {} |
| 576 | tracker = None |
| 577 | args = args_kwargs.pop(0) |
| 578 | if isinstance(args, (list, tuple)): |
| 579 | # you can specify a custom tracker value if you want |
| 580 | # this helps with correlating results |
| 581 | with suppress(ValueError): |
| 582 | args, kwargs, tracker = args |
| 583 | # or you can just specify args/kwargs |
| 584 | with suppress(ValueError): |
| 585 | args, kwargs = args |
| 586 | |
| 587 | if not isinstance(kwargs, dict): |
| 588 | raise ValueError(f"kwargs must be dict (got: {kwargs})") |
| 589 | if not isinstance(args, (list, tuple)): |
| 590 | args = [args] |
| 591 | |
| 592 | task = self.new_child_task(fn(*args, **kwargs, **global_kwargs)) |
| 593 | tasks[task] = (args, kwargs, tracker) |
| 594 | |
| 595 | for _ in range(threads): # Start initial batch of tasks |
| 596 | new_task() |
| 597 | |
| 598 | while tasks: # While there are tasks pending |
| 599 | # Wait for the first task to complete |
| 600 | finished = await self.finished_tasks(tasks, timeout=timeout) |
| 601 | for task in finished: |
| 602 | result = task.result() |
| 603 | (args, kwargs, tracker) = tasks.pop(task) |
| 604 | yield (args, kwargs, tracker), result |
| 605 | new_task() |
| 606 | |
| 607 | def new_child_task(self, coro): |
| 608 | """ |
no test coverage detected