(self)
| 440 | self.kwargs = kwargs |
| 441 | |
| 442 | async def producer(self): |
| 443 | # Data generation logic here |
| 444 | for step_result in _generate_tokens( |
| 445 | self.process_func, *self.args, **self.kwargs |
| 446 | ): |
| 447 | await self.queue.put(step_result) |
| 448 | await asyncio.sleep(0.0001) |
| 449 | # asyc sleep otherwise this doesn't yield any result |
| 450 | if self.shutdown_event.is_set(): |
| 451 | break |
| 452 | await self.queue.put(None) |
| 453 | |
| 454 | def __aiter__(self): |
| 455 | self.iterator_task = asyncio.create_task(self.producer()) |
no test coverage detected