(events_for_one_agent)
| 59 | # Agents are processed in parallel. |
| 60 | # Events for each agent are put on queue sequentially. |
| 61 | async def process_an_agent(events_for_one_agent): |
| 62 | try: |
| 63 | async for event in events_for_one_agent: |
| 64 | resume_signal = asyncio.Event() |
| 65 | await queue.put((event, resume_signal)) |
| 66 | # Wait for upstream to consume event before generating new events. |
| 67 | await resume_signal.wait() |
| 68 | finally: |
| 69 | # Mark agent as finished. |
| 70 | await queue.put((sentinel, None)) |
| 71 | |
| 72 | async with asyncio.TaskGroup() as tg: |
| 73 | for events_for_one_agent in agent_runs: |
no test coverage detected