Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions cq/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
from ._core.routing.di import DIAdapter
from ._core.routing.dispatchers.abc import Dispatcher
from ._core.routing.dispatchers.bus import Bus
from ._core.routing.dispatchers.functions import dispatch_sequentially
from ._core.routing.dispatchers.pipe import ContextPipeline, Pipe
from ._core.routing.router import Router

Expand Down Expand Up @@ -46,6 +47,7 @@
"Router",
"__router__",
"command_handler",
"dispatch_sequentially",
"event_handler",
"new_command_bus",
"new_event_bus",
Expand Down
11 changes: 11 additions & 0 deletions cq/_core/routing/dispatchers/functions.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
from typing import Any

from cq._core.routing.dispatchers.abc import Dispatcher


async def dispatch_sequentially[T, *Ts](
dispatcher: Dispatcher[T, Any],
/,
*messages: T,
) -> tuple[*Ts]:
return tuple([await dispatcher.dispatch(message) for message in messages])
File renamed without changes.
File renamed without changes.
File renamed without changes.
Empty file added tests/core/routing/__init__.py
Empty file.
Empty file.
21 changes: 21 additions & 0 deletions tests/core/routing/dispatchers/test_functions.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
from typing import Any, Self

from cq import Bus, dispatch_sequentially


async def test_dispatch_sequentially_with_success_return_tuple(
bus: Bus[Any, Any],
) -> None:
class Handler:
async def handle(self, message: str) -> str:
return message

@classmethod
async def async_factory(cls) -> Self:
return cls()

bus.subscribe(str, Handler.async_factory)

messages = ("a", "b")
results: tuple[str, str] = await dispatch_sequentially(bus, *messages)
assert messages == results
File renamed without changes.
476 changes: 238 additions & 238 deletions uv.lock

Large diffs are not rendered by default.

Loading