|
| 1 | +import asyncio |
| 2 | +import uuid |
| 3 | + |
| 4 | +import pytest |
| 5 | + |
| 6 | +from taskiq import InMemoryBroker |
| 7 | +from taskiq.events import TaskiqEvents |
| 8 | +from taskiq.state import TaskiqState |
| 9 | + |
| 10 | + |
| 11 | +@pytest.mark.anyio |
| 12 | +async def test_inmemory_success() -> None: |
| 13 | + broker = InMemoryBroker() |
| 14 | + test_val = uuid.uuid4().hex |
| 15 | + |
| 16 | + @broker.task |
| 17 | + async def task() -> str: |
| 18 | + return test_val |
| 19 | + |
| 20 | + kicked = await task.kiq() |
| 21 | + result = await kicked.wait_result() |
| 22 | + assert result.return_value == test_val |
| 23 | + assert not broker._running_tasks |
| 24 | + |
| 25 | + |
| 26 | +@pytest.mark.anyio |
| 27 | +async def test_cannot_listen() -> None: |
| 28 | + broker = InMemoryBroker() |
| 29 | + |
| 30 | + with pytest.raises(RuntimeError): |
| 31 | + async for _ in broker.listen(): |
| 32 | + pass |
| 33 | + |
| 34 | + |
| 35 | +@pytest.mark.anyio |
| 36 | +async def test_startup() -> None: |
| 37 | + broker = InMemoryBroker() |
| 38 | + test_value = uuid.uuid4().hex |
| 39 | + |
| 40 | + @broker.on_event(TaskiqEvents.WORKER_STARTUP) |
| 41 | + async def _w_startup(state: TaskiqState) -> None: |
| 42 | + state.from_worker = test_value |
| 43 | + |
| 44 | + @broker.on_event(TaskiqEvents.CLIENT_STARTUP) |
| 45 | + async def _c_startup(state: TaskiqState) -> None: |
| 46 | + state.from_client = test_value |
| 47 | + |
| 48 | + await broker.startup() |
| 49 | + |
| 50 | + assert broker.state.from_worker == test_value |
| 51 | + assert broker.state.from_client == test_value |
| 52 | + |
| 53 | + |
| 54 | +@pytest.mark.anyio |
| 55 | +async def test_shutdown() -> None: |
| 56 | + broker = InMemoryBroker() |
| 57 | + test_value = uuid.uuid4().hex |
| 58 | + |
| 59 | + @broker.on_event(TaskiqEvents.WORKER_SHUTDOWN) |
| 60 | + async def _w_startup(state: TaskiqState) -> None: |
| 61 | + state.from_worker = test_value |
| 62 | + |
| 63 | + @broker.on_event(TaskiqEvents.CLIENT_SHUTDOWN) |
| 64 | + async def _c_startup(state: TaskiqState) -> None: |
| 65 | + state.from_client = test_value |
| 66 | + |
| 67 | + await broker.shutdown() |
| 68 | + |
| 69 | + assert broker.state.from_worker == test_value |
| 70 | + assert broker.state.from_client == test_value |
| 71 | + |
| 72 | + |
| 73 | +@pytest.mark.anyio |
| 74 | +async def test_execution() -> None: |
| 75 | + broker = InMemoryBroker() |
| 76 | + test_value = uuid.uuid4().hex |
| 77 | + |
| 78 | + @broker.task |
| 79 | + async def test_task() -> str: |
| 80 | + await asyncio.sleep(0.5) |
| 81 | + return test_value |
| 82 | + |
| 83 | + task = await test_task.kiq() |
| 84 | + assert not await task.is_ready() |
| 85 | + |
| 86 | + result = await task.wait_result() |
| 87 | + assert result.return_value == test_value |
0 commit comments