Async in-memory queue implementation (single partition)
| 16 | |
| 17 | |
| 18 | class MemoryQueue(QueueInterface): |
| 19 | """Async in-memory queue implementation (single partition)""" |
| 20 | |
| 21 | def __init__(self): |
| 22 | self._queue: asyncio.Queue[QueueMessage] = asyncio.Queue() |
| 23 | |
| 24 | async def put(self, message: QueueMessage) -> None: |
| 25 | logger.debug("Adding message to queue: agent_id=%s", message.agent_id) |
| 26 | await self._queue.put(message) |
| 27 | |
| 28 | async def get(self, timeout: Optional[float] = None) -> QueueMessage: |
| 29 | """ |
| 30 | Retrieve a message from the queue. |
| 31 | |
| 32 | Args: |
| 33 | timeout: Optional timeout in seconds (None = block indefinitely) |
| 34 | |
| 35 | Returns: |
| 36 | QueueMessage from the queue |
| 37 | |
| 38 | Raises: |
| 39 | asyncio.TimeoutError: If no message available within timeout |
| 40 | """ |
| 41 | if timeout is not None: |
| 42 | message = await asyncio.wait_for(self._queue.get(), timeout=timeout) |
| 43 | else: |
| 44 | message = await self._queue.get() |
| 45 | logger.debug("Retrieved message from queue: agent_id=%s", message.agent_id) |
| 46 | return message |
| 47 | |
| 48 | async def close(self) -> None: |
| 49 | pass |
| 50 | |
| 51 | |
| 52 | class PartitionedMemoryQueue(QueueInterface): |
no outgoing calls