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
8 changes: 4 additions & 4 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -29,18 +29,18 @@ repository = "https://github.com/taskiq-python/taskiq-nats"
keywords = ["taskiq", "tasks", "distributed", "async", "nats", "result_backend"]
requires-python = ">=3.10,<4"
dependencies = [
"taskiq>=0.11.20,<1",
"taskiq>=0.13.0",
"nats-py>=2.2.0",
]

[dependency-groups]
dev = [
"pre-commit>=4.4.0",
# lint
"ruff>=0.14.5",
"black>=25.11.0",
"ruff>=0.16.9",
"black>=26.5.1",
# type check
"mypy>=1.18.2",
"mypy>=2.3.1",
# tests
"pytest>=9.0.1",
"pytest-cov>=7.0.0",
Expand Down
2 changes: 1 addition & 1 deletion taskiq_nats/broker.py
Original file line number Diff line number Diff line change
Expand Up @@ -105,7 +105,7 @@ class BaseJetStreamBroker(
be sure that messages are delivered to the workers.
"""

def __init__(
def __init__( # noqa: PLR0917
self,
servers: str | list[str],
subject: str = "taskiq_tasks",
Expand Down
8 changes: 4 additions & 4 deletions taskiq_nats/result_backend.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@
from nats.js.object_store import ObjectStore
from taskiq import AsyncResultBackend, ResultGetError
from taskiq.abc.serializer import TaskiqSerializer
from taskiq.compat import model_dump, model_validate
from taskiq.result import TaskiqResult
from taskiq.serializers import PickleSerializer

Expand Down Expand Up @@ -73,7 +72,7 @@ async def set_result(self, task_id: str, result: TaskiqResult[_ReturnType]) -> N
"""
await self.object_store.put(
name=task_id,
data=self.serializer.dumpb(model_dump(result)),
data=self.serializer.dumpb(result.model_dump(mode="json")),
)

async def is_result_ready(self, task_id: str) -> bool:
Expand Down Expand Up @@ -114,8 +113,9 @@ async def get_result(
name=task_id,
)

taskiq_result: TaskiqResult[_ReturnType] = model_validate(
TaskiqResult[_ReturnType],
taskiq_result: TaskiqResult[_ReturnType] = TaskiqResult[
_ReturnType
].model_validate(
self.serializer.loadb(result.data), # type: ignore[arg-type]
)

Expand Down
5 changes: 2 additions & 3 deletions taskiq_nats/schedule_source.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@
from nats.js.kv import KeyValue
from taskiq import ScheduledTask, ScheduleSource
from taskiq.abc.serializer import TaskiqSerializer
from taskiq.compat import model_dump, model_validate
from taskiq.serializers import PickleSerializer

log = logging.getLogger(__name__)
Expand Down Expand Up @@ -85,7 +84,7 @@ async def add_schedule(self, schedule: ScheduledTask) -> None:
"""
await self.kv.put(
f"{self.prefix}.{schedule.schedule_id}",
self.serializer.dumpb(model_dump(schedule)),
self.serializer.dumpb(schedule.model_dump(mode="json")),
)

async def get_schedules(self) -> list[ScheduledTask]:
Expand All @@ -102,7 +101,7 @@ async def get_schedules(self) -> list[ScheduledTask]:
return []

return [
model_validate(ScheduledTask, self.serializer.loadb(schedule.value))
ScheduledTask.model_validate(self.serializer.loadb(schedule.value))
for schedule in schedules
if schedule and schedule.value
]
Expand Down
Loading
Loading