Skip to content

Usage with taskiq

How to use

1. Install modern-di-taskiq

uv add modern-di-taskiq
pip install modern-di-taskiq
poetry add modern-di-taskiq

2. Apply to your application

import dataclasses
import typing

from modern_di import Container, Group, Scope, providers
from modern_di_taskiq import FromDI, setup_di
from taskiq import InMemoryBroker


@dataclasses.dataclass(kw_only=True, slots=True, frozen=True)
class Settings:
    service_name: str = "catalog"


@dataclasses.dataclass(kw_only=True, slots=True)
class Report:
    settings: Settings   # APP-scoped, injected by type

    def as_dict(self) -> dict[str, str]:
        return {"service": self.settings.service_name}


class AppGroup(Group):
    settings = providers.Factory(Settings, scope=Scope.APP, cache=True)
    report = providers.Factory(Report, scope=Scope.REQUEST)


broker = InMemoryBroker()
container = Container(groups=[AppGroup])
setup_di(broker, container)
container.validate()  # after setup_di: its context provider is now registered


@broker.task
async def get_report(
    report: typing.Annotated[Report, FromDI(Report)],
) -> dict[str, str]:
    return report.as_dict()

setup_di(broker, container) stores the container on broker.state, registers the message context provider, and hooks the root container's lifecycle to taskiq's worker events. Each task that uses FromDI gets its own Scope.REQUEST child container.

Scopes

The integration creates a Scope.REQUEST child container for each task that uses FromDI, built through TaskiqDepends when the task's dependencies are resolved. A task with no FromDI parameter gets no child container. REQUEST-scoped providers and their finalizers live for that one task, and the child is closed after the task returns or raises.

Resolution is synchronous, as everywhere in modern-di, while the child is closed with close_async(), so REQUEST-scoped finalizers may be async or sync.

There is no Scope.SESSION for taskiq: a task queue doesn't have a session concept comparable to websockets.

Root container lifecycle

setup_di opens the root container on TaskiqEvents.WORKER_STARTUP and closes it with close_async() on WORKER_SHUTDOWN, which runs APP-scoped finalizers. taskiq worker fires both. With an InMemoryBroker, async with broker: fires both too.

run_receiver_task never closes the root

A receiver embedded with taskiq.api.run_receiver_task(...) never calls broker.shutdown(), so WORKER_SHUTDOWN never fires and APP-scoped finalizers never run, whatever run_startup is set to. Tasks still resolve, since a constructed container is already open. After the receiver stops, call await broker.shutdown() or await container.close_async() yourself.

Framework context objects

taskiq.TaskiqMessage is automatically made available by the integration, so factories can declare it as a parameter and get the message that triggered the current task. See Framework context objects for how implicit and explicit resolution work.

The following context provider is also available for explicit import:

  • taskiq_message_provider provides the current taskiq.TaskiqMessage object.

Implicit (type-based) usage

import taskiq
from modern_di import Group, Scope, providers


def create_task_info(message: taskiq.TaskiqMessage) -> dict[str, str]:
    return {
        "task_id": message.task_id,
        "task_name": message.task_name,
    }


class AppGroup(Group):
    # The message dependency is resolved by type annotation
    task_info = providers.Factory(
        create_task_info,
        scope=Scope.REQUEST,
    )

Explicit (provider-based) usage

import taskiq
import modern_di_taskiq
from modern_di import Group, Scope, providers


def create_task_info(message: taskiq.TaskiqMessage) -> dict[str, str]:
    return {"task_id": message.task_id}


class AppGroup(Group):
    task_info = providers.Factory(
        create_task_info,
        scope=Scope.REQUEST,
        kwargs={"message": modern_di_taskiq.taskiq_message_provider},
    )

Testing

Enter an InMemoryBroker with async with so the worker events fire, then send a task and wait for its result:

async with broker:
    task = await get_report.kiq()
    result = await task.wait_result()
assert result.return_value == {"service": "catalog"}

See also

API

Symbol Description
setup_di(broker, container) Wire the APP-scope container into taskiq: a REQUEST child container is then built through TaskiqDepends for each task that uses FromDI; opens/closes the APP container on worker startup/shutdown.
FromDI(provider_or_type, *, use_cache=True) Marker for Annotated[T, FromDI(...)] in task signatures; accepts a provider instance or a plain type. use_cache is passed through to TaskiqDepends. Raises RuntimeError naming setup_di when a task reaches it without setup_di called.
fetch_di_container(broker) Returns the APP-scope container registered with the taskiq broker.
taskiq_message_provider ContextProvider for the current taskiq.TaskiqMessage.