|
1 | | -from collections.abc import AsyncGenerator |
2 | | -from typing import NewType, Protocol |
| 1 | +from asyncio import Queue, TaskGroup, gather |
| 2 | +from collections.abc import AsyncGenerator, Awaitable, Callable |
| 3 | +from dataclasses import dataclass, field |
| 4 | +from typing import Any, overload |
3 | 5 |
|
4 | 6 | from nats.aio.msg import Msg as NatsMessage |
5 | 7 | from nats.errors import TimeoutError as NatsTimeoutError |
6 | 8 | from nats.js import JetStreamContext |
7 | 9 | from taskiq import ( |
8 | 10 | AckableMessage, |
9 | 11 | AsyncBroker, |
| 12 | + AsyncTaskiqDecoratedTask, |
10 | 13 | BrokerMessage, |
11 | 14 | ) |
12 | 15 |
|
13 | 16 |
|
14 | | -class PullSubscribe(Protocol): |
| 17 | +@dataclass(frozen=True) |
| 18 | +class PullSubscribe: |
| 19 | + _func: Callable[ |
| 20 | + [JetStreamContext, str], Awaitable[JetStreamContext.PullSubscription], |
| 21 | + ] |
| 22 | + |
15 | 23 | async def __call__( |
16 | | - self, js: JetStreamContext, subject: str, /, |
17 | | - ) -> JetStreamContext.PullSubscription: ... |
| 24 | + self, js: JetStreamContext, subject: str, |
| 25 | + ) -> JetStreamContext.PullSubscription: |
| 26 | + return await self._func(js, subject) |
18 | 27 |
|
19 | 28 |
|
20 | | -class NatsBroker(AsyncBroker): |
21 | | - _consumer: JetStreamContext.PullSubscription |
22 | | - js: JetStreamContext |
| 29 | +@dataclass |
| 30 | +class NatsQueue: |
| 31 | + _subject: str |
| 32 | + _pull_subscribe: PullSubscribe |
| 33 | + _pull_consume_batch: int |
| 34 | + _pull_consume_timeout: float | None |
23 | 35 |
|
24 | | - def __init__( |
25 | | - self, |
26 | | - subject: str, |
27 | | - pull_subscribe: PullSubscribe, |
28 | | - pull_consume_batch: int = 1, |
29 | | - pull_consume_timeout: float | None = 5, |
30 | | - ) -> None: |
31 | | - super().__init__() |
32 | | - self._subject = subject |
33 | | - self._pull_subscribe = pull_subscribe |
34 | | - self._pull_consume_batch = pull_consume_batch |
35 | | - self._pull_consume_timeout = pull_consume_timeout |
| 36 | + _js: JetStreamContext = field(init=False) |
| 37 | + _pull_subscription: JetStreamContext.PullSubscription = field(init=False) |
36 | 38 |
|
37 | | - async def startup(self) -> None: |
38 | | - await super().startup() |
39 | | - self._consumer = await self._pull_subscribe(self.js, self._subject) |
| 39 | + async def startup(self, js: JetStreamContext) -> None: |
| 40 | + self._js = js |
| 41 | + self._pull_subscription = await self._pull_subscribe(js, self._subject) |
40 | 42 |
|
41 | | - async def kick(self, message: BrokerMessage) -> None: |
42 | | - await self.js.publish( |
| 43 | + async def push(self, message: BrokerMessage) -> None: |
| 44 | + await self._js.publish( |
43 | 45 | self._subject, |
44 | 46 | payload=message.message, |
45 | 47 | headers=message.labels, |
46 | 48 | ) |
47 | 49 |
|
48 | | - async def listen(self) -> AsyncGenerator[AckableMessage]: |
| 50 | + async def pull_to(self, output: Queue[AckableMessage]) -> None: |
| 51 | + nats_messages: list[NatsMessage] |
| 52 | + |
49 | 53 | while True: |
50 | 54 | try: |
51 | | - nats_messages: list[NatsMessage] = await self._consumer.fetch( |
| 55 | + nats_messages = await self._pull_subscription.fetch( |
52 | 56 | batch=self._pull_consume_batch, |
53 | 57 | timeout=self._pull_consume_timeout, |
54 | 58 | ) |
55 | | - for nats_message in nats_messages: |
56 | | - yield AckableMessage( |
57 | | - data=nats_message.data, |
58 | | - ack=nats_message.ack, |
59 | | - ) |
60 | 59 | except NatsTimeoutError: |
61 | 60 | continue |
62 | 61 |
|
| 62 | + ackable_messages = ( |
| 63 | + AckableMessage( |
| 64 | + data=nats_message.data, |
| 65 | + ack=nats_message.ack, |
| 66 | + ) |
| 67 | + for nats_message in nats_messages |
| 68 | + ) |
| 69 | + await gather(*( |
| 70 | + output.put(ackable_message) |
| 71 | + for ackable_message in ackable_messages |
| 72 | + )) |
| 73 | + |
| 74 | + |
| 75 | +class NatsBroker(AsyncBroker): |
| 76 | + js: JetStreamContext |
| 77 | + pulling_queue: Queue[AckableMessage] |
| 78 | + |
| 79 | + def __init__( |
| 80 | + self, |
| 81 | + default_subject: str = "taskiq.>", |
| 82 | + default_pull_subscribe: PullSubscribe | None = None, |
| 83 | + default_pull_consume_batch: int = 1, |
| 84 | + default_pull_consume_timeout: float | None = 5, |
| 85 | + ) -> None: |
| 86 | + super().__init__() |
| 87 | + self._default_subject = default_subject |
| 88 | + self._default_pull_subscribe = ( |
| 89 | + default_pull_subscribe |
| 90 | + or (lambda js, sub: js.pull_subscribe( |
| 91 | + sub, |
| 92 | + "taskiq", |
| 93 | + "TASKIQ", |
| 94 | + )) |
| 95 | + ) |
| 96 | + self._default_pull_consume_batch = default_pull_consume_batch |
| 97 | + self._default_pull_consume_timeout = default_pull_consume_timeout |
| 98 | + self._queue_by_task_name = dict[str, NatsQueue]() |
| 99 | + |
| 100 | + @overload |
| 101 | + def task[**PmT, RT]( |
| 102 | + self, |
| 103 | + task_name: Callable[PmT, RT], |
| 104 | + **labels: Any, # noqa: ANN401 |
| 105 | + ) -> AsyncTaskiqDecoratedTask[PmT, RT]: |
| 106 | + ... |
| 107 | + |
| 108 | + @overload |
| 109 | + def task[**PmT, RT]( |
| 110 | + self, |
| 111 | + task_name: str | None = None, |
| 112 | + **labels: Any, # noqa: ANN401 |
| 113 | + ) -> Callable[ |
| 114 | + [Callable[PmT, RT]], |
| 115 | + AsyncTaskiqDecoratedTask[PmT, RT], |
| 116 | + ]: |
| 117 | + ... |
| 118 | + |
| 119 | + def task[**PmT, RT]( |
| 120 | + self, |
| 121 | + task_name: str | Callable[PmT, RT] | None = None, |
| 122 | + **labels: Any, |
| 123 | + ) -> Any: |
| 124 | + subject = labels.pop("subject", self._default_subject) |
| 125 | + pull_subscribe = labels.pop( |
| 126 | + "pull_subscribe", self._default_pull_subscribe, |
| 127 | + ) |
| 128 | + pull_consume_batch = labels.pop( |
| 129 | + "pull_consume_batch", self._default_pull_consume_batch, |
| 130 | + ) |
| 131 | + pull_consume_timeout = labels.pop( |
| 132 | + "pull_consume_timeout", self._default_pull_consume_timeout, |
| 133 | + ) |
| 134 | + queue = NatsQueue( |
| 135 | + subject, |
| 136 | + pull_subscribe, |
| 137 | + pull_consume_batch, |
| 138 | + pull_consume_timeout, |
| 139 | + ) |
| 140 | + |
| 141 | + result = super().task(task_name, **labels) |
| 142 | + |
| 143 | + if isinstance(result, AsyncTaskiqDecoratedTask): |
| 144 | + self._queue_by_task_name[result.task_name] = queue |
| 145 | + return result |
| 146 | + |
| 147 | + def decorator( |
| 148 | + func: Callable[PmT, RT], |
| 149 | + ) -> AsyncTaskiqDecoratedTask[PmT, RT]: |
| 150 | + task = result(func) |
| 151 | + self._queue_by_task_name[task.task_name] = queue |
| 152 | + return task |
| 153 | + |
| 154 | + return decorator |
| 155 | + |
| 156 | + async def startup(self) -> None: |
| 157 | + await super().startup() |
| 158 | + await gather(*( |
| 159 | + queue.startup(self.js) |
| 160 | + for queue in self._queue_by_task_name.values() |
| 161 | + )) |
| 162 | + |
| 163 | + async def kick(self, message: BrokerMessage) -> None: |
| 164 | + await self._queue_by_task_name[message.task_name].push(message) |
| 165 | + |
| 166 | + async def listen(self) -> AsyncGenerator[AckableMessage]: |
| 167 | + async with TaskGroup() as pulling_tasks: |
| 168 | + for nats_queue in self._queue_by_task_name.values(): |
| 169 | + pulling_tasks.create_task( |
| 170 | + nats_queue.pull_to(self.pulling_queue), |
| 171 | + ) |
63 | 172 |
|
64 | | -NatsBrokers = NewType("NatsBrokers", tuple[NatsBroker, ...]) |
| 173 | + while True: |
| 174 | + yield await self.pulling_queue.get() |
0 commit comments