Metadata-Version: 2.4
Name: xtr-messenger
Version: 3.0.0
Summary: A message bus for Python: envelopes, stamps, a middleware chain, and pluggable transports.
Keywords: message-bus,message-queue,middleware,envelope,taskiq,amqp
License-Expression: MIT
License-File: LICENSE
Classifier: Development Status :: 3 - Alpha
Classifier: Intended Audience :: Developers
Classifier: Programming Language :: Python :: 3.11
Classifier: Programming Language :: Python :: 3.12
Classifier: Programming Language :: Python :: 3.13
Classifier: Programming Language :: Python :: 3.14
Classifier: Typing :: Typed
Requires-Dist: msgspec>=0.18
Requires-Dist: typing-extensions>=4.4
Requires-Dist: xtr-event-dispatcher-contracts>=3.0,<4
Requires-Dist: xtr-logging-contracts>=3.0,<4
Requires-Dist: aio-pika>=9.4 ; extra == 'amqp'
Requires-Dist: taskiq>=0.12.4 ; extra == 'amqp'
Requires-Dist: taskiq-aio-pika>=0.6.0 ; extra == 'amqp'
Requires-Dist: anyio>=4.0 ; extra == 'amqp'
Requires-Dist: typing-extensions>=4.4 ; extra == 'amqp'
Requires-Dist: xtr-console>=3.0,<4 ; extra == 'console'
Requires-Dist: xtr-dependency-injection>=3.0,<4 ; extra == 'di'
Requires-Dist: xtr-service-contracts>=3.0,<4 ; extra == 'di'
Requires-Dist: pydantic>=2.0 ; extra == 'pydantic'
Requires-Dist: taskiq>=0.12.4 ; extra == 'taskiq'
Requires-Dist: anyio>=4.0 ; extra == 'taskiq'
Requires-Dist: typing-extensions>=4.4 ; extra == 'taskiq'
Requires-Python: >=3.11
Provides-Extra: amqp
Provides-Extra: console
Provides-Extra: di
Provides-Extra: pydantic
Provides-Extra: taskiq
Description-Content-Type: text/markdown

<div align="center">

# xtr-messenger

**A message bus for Python — envelopes, stamps, a middleware chain, and pluggable transports.**

<img alt="python 3.11+" src="https://img.shields.io/badge/python-%E2%89%A5%203.11-3776AB?logo=python&logoColor=white">
<img alt="core dependencies: 4" src="https://img.shields.io/badge/core%20deps-4-3FB950">
<img alt="typed" src="https://img.shields.io/badge/typed-ty%20%2B%20basedpyright-1f6feb">
<img alt="license MIT" src="https://img.shields.io/badge/license-MIT-blue">

</div>

---

## Why?

Publishing a message should not drag a broker into your import graph, and moving a message
onto a queue should not mean rewriting the code that sends it.

Dispatch wraps a message in an **envelope** and walks it through a **middleware chain**. Each
middleware records what it did by appending a **stamp**. A **routing table** maps a message
type to one or more named **transports**. Nothing at the dispatch site knows which.

- 🪶 **Four core dependencies** — `msgspec`, `typing-extensions`, and two interface-only
  packages, `xtr-logging-contracts` and `xtr-event-dispatcher-contracts`. A broker is an extra.
- 🔌 **Transports are discovered** — one entry point adds a scheme; no fork, no registry to edit.
- 💤 **Lazily loaded** — an app speaking `sync://` never imports a broker library.
- 🎯 **Handlers run in one place** — the same way in-process, in a worker, or under a broker's own loop.
- 🧩 **Protocol-based** — every collaborator is a constructor argument, so a DI container can own the graph.
- 📨 **Dataclasses or pydantic** — both on one bus; dataclasses are checked for shape, models by their own rules.

```python
await bus.dispatch(IngestDocument(document_id=doc.id))
```

Where that goes is configuration. Whether it happens in-process or on RabbitMQ is a DSN.

## Install

```sh
uv add xtr-messenger                    # sync:// and in-memory://
uv add "xtr-messenger[amqp]"            # + RabbitMQ
uv add "xtr-messenger[taskiq]"          # + any other taskiq broker
uv add "xtr-messenger[pydantic]"        # + pydantic messages
uv add "xtr-messenger[di]"              # + a MessengerBundle for xtr-dependency-injection
uv add "xtr-messenger[console]"         # + the messenger:consume console command
```

| Extra | Brings | For |
| --- | --- | --- |
| *(none)* | `msgspec`, `typing-extensions`, `xtr-logging-contracts`, `xtr-event-dispatcher-contracts` | `sync://`, `in-memory://`, the whole core |
| `pydantic` | `pydantic` | Messages validated by a model, not just a shape |
| `taskiq` | `taskiq` | Publishing and consuming over **any** taskiq broker |
| `amqp` | `taskiq-aio-pika` | RabbitMQ, with retries and real dead-lettering |
| `di` | `xtr-dependency-injection` | A `MessengerBundle` for [xtr-dependency-injection](../xtr-dependency-injection) |
| `console` | `xtr-console` | `messenger:consume`, on an [xtr-console](https://github.com/xterr/python-xtr-console) application |

Requires Python 3.11+. An application that wants records written somewhere adds `xtr-logging`
itself.

## Quick start

A message is a frozen dataclass carrying identifiers. `@as_message` declares it, and its
properties live on the class:

```python
from dataclasses import dataclass
from uuid import UUID

from xtr_messenger import as_message


@as_message(name="ingest.document.v1")
@dataclass(frozen=True, slots=True)
class IngestDocument:
    document_id: UUID
    tenant_id: UUID
```

`name` is the contract between producer and consumer, and the only thing they share. It
defaults to `module:QualName`, which changes if the class moves — pin an explicit, versioned
name for anything that outlives a deploy. `@as_message(transport="jobs")` gives the message a
default transport, used when the routing table says nothing about it.

A consumer decodes only the messages it declared: the name on the wire is looked up among the
declared ones and never used to import a module, so whoever produces a message cannot choose
what the consumer imports. A message that crosses a serializer — any transport but `sync://` —
must therefore be declared on the consuming side; an unknown name fails with
`MessageDecodingFailedError`.

A handler is a function that imports nothing but the message:

```python
from xtr_messenger import as_message_handler


@as_message_handler(IngestDocument)
async def ingest(message: IngestDocument) -> None: ...
```

Describe the transports and where messages go, then build a bus:

```python
from xtr_messenger import MessageBusConfig, MessageBusFactory, TransportConfig

CONFIG = MessageBusConfig(
    transports={"sync": TransportConfig("sync://")},
    routing={IngestDocument: "sync"},
)

bus = MessageBusFactory(CONFIG).bus()
await bus.dispatch(IngestDocument(document_id=doc_id, tenant_id=tenant_id))
```

Moving that message onto RabbitMQ later is one line of configuration. The dispatch site does
not change.

## Handlers

The bus calls a handler with the message, and with the envelope too when a second parameter is
annotated `Envelope`. Anything else is refused with `HandlerSignatureError` where the handler is
declared, not on the first message that reaches it.

```python
from xtr_messenger import Envelope, RedeliveryStamp, as_message_handler


@as_message_handler(IngestDocument)
async def ingest(message: IngestDocument, envelope: Envelope) -> None:
    stamp = envelope.last(RedeliveryStamp)  # which delivery attempt this is


@as_message_handler(IngestDocument)
class AuditIngest:  # built once, on its first message
    async def __call__(self, message: IngestDocument) -> None: ...
```

A function, a callable object or a class all work. A class is built **once**, on its first
message, and shared by every message after — with no arguments on its own, or by the container
when [one is wired](#wiring-with-a-container). A handler is a service, not a value: building one
per message would cost a construction each, thousands of times a second. Messages may be handled
concurrently, so keep per-message state off `self`.

**Every handler of a message runs**, in the order declared, and each leaves a `HandledStamp`
behind, carrying what it returned in `result` — `None` when it returns nothing. The value
stays in the process; it is never sent anywhere, so it need not be serializable. Lookup walks the message's bases, so a handler on a marker class still fires for a
subclass that has handlers of its own.

A handler that raises does not stop the others: every one runs, and then the dispatch fails with
one `HandlersFailedError` carrying what each failed handler raised, by name, in `errors`, and the
envelope with a `HandledStamp` for each that succeeded. A worker records the first handler's
error on the rejected message. A retry runs **every** handler again, those that succeeded
included — so keep handlers idempotent.

Declaring writes to a process-wide registry, which is what lets a handler module import nothing
but the message. Where one registry per process is too coarse — two applications in one test
run, say — declare into a `HandlersLocator` of your own and hand it to the bus or the worker:

```python
from xtr_messenger import HandlersLocator

orders = HandlersLocator()


@as_message_handler(PlaceOrder, orders)
async def place_order(message: PlaceOrder) -> None: ...


bus = MessageBusFactory(CONFIG, handlers=orders).bus()
```

## Configuring transports

Configuration is inert data: which transports exist and where messages go. It builds nothing —
a factory does — so it can come from a settings module, an environment variable or a parsed
file without dragging a broker along.

```python
CONFIG = MessageBusConfig(
    transports={
        "high": TransportConfig(AMQP_URL, queue="jobs_high"),
        "low": TransportConfig(f"{AMQP_URL}?queue=jobs_low"),
        "sync": TransportConfig("sync://"),
        "test": TransportConfig("in-memory://?serialize=true"),
    },
    routing={
        UrgentJob: "high",
        AuditRecorded: ["low", "test"],  # fan out
        "*": "low",  # catch-all
    },
)
```

A DSN selects the adapter and carries its settings. It is read when the bus or a worker is built,
not when the `TransportConfig` is made — a configuration may be written before the environment it
reads is set — so one without a scheme fails with `InvalidDsnError` there. Transports differing only in their query string address the same server, so they share
one connection.

| DSN | Transport |
| --- | --- |
| `sync://` | Handles in the calling process, during the dispatch |
| `in-memory://` | Records what was dispatched, for a worker or a test to drain; `?serialize=true` round-trips it |
| `amqp://…`, `amqps://…` | RabbitMQ via taskiq |

Settings go in the query string or in `options`, whichever suits — a DSN travels in one
environment variable, `options` survives review once there are several. `options` wins.

```python
TransportConfig(
    "amqp://user:pass@rabbit:5672/?queue=jobs&max_attempts=10",
    options={"prefetch_count": "50", "dead_letter_queue": "jobs.dlq"},
)
```

<details>
<summary><b>All 26 settings an <code>amqp://</code> transport accepts</b></summary>

| Group | Settings |
| --- | --- |
| Retries | `max_attempts`, `base_delay_seconds`, `dead_letter_queue` |
| Exchange | `exchange`, `exchange_type`, `exchange_durable`, `exchange_auto_delete` |
| Queue | `queue`, `queue_type`, `queue_durable`, `queue_auto_delete`, `queue_exclusive`, `queue_max_priority`, `routing_key` |
| Connection | `heartbeat`, `connect_timeout`, `connection_name`, `frame_max`, `channel_max` |
| TLS | `cacert`, `cert`, `key`, `verify` |
| Consumption | `prefetch_count`, `max_async_tasks`, `auto_setup` |

`max_async_tasks` is how many messages a worker handles at once: ten per CPU, capped at 100,
unless set. A message is acked once handled, so `prefetch_count` (10 by default) caps it too —
a worker handles the smaller of the two at once. Raise both together to go wider.

An unrecognised setting is **refused** with `UnknownTransportOptionError`, naming what the scheme
does accept, and a value it cannot use — `prefetch_count=lots` — with
`InvalidTransportOptionError`. A typo in configuration is always a mistake, and one that is
ignored leaves a transport running on defaults nobody chose.

Credentials are not settings — a URL already expresses them, so they stay in the DSN.

</details>

## Running a worker

A worker entrypoint is an ordinary Python program that names the transports it serves:

```python
# app/worker_high.py
import asyncio
import signal

import app.handlers.ingest  # noqa: F401 — importing declares the message and its handler

from app.bus import CONFIG
from xtr_messenger import WorkerFactory


async def main() -> None:
    worker = WorkerFactory(CONFIG).worker(["high"])
    asyncio.get_running_loop().add_signal_handler(signal.SIGTERM, worker.stop)
    await worker.run()


asyncio.run(main())
```

```sh
uv run python -m app.worker_high
```

What comes back is a `WorkerInterface`, never the broker underneath. Every message it hands the
handlers carries a `ReceivedStamp` naming the transport it was consumed from, as configured —
`high`, not the DSN's scheme. Point one process at
`["high"]` and another at `["low"]` and each consumes its own workload. `stop()` lets it finish
the message in hand and return — at once, if it is waiting for one.

Most transports are driven by the library's own loop. A broker that brings its own worker —
taskiq does — supplies it instead, and the entrypoint above cannot tell. Either way what it
receives is dispatched into a bus, and the bus calls the handlers.

> On AMQP a worker registers one task per message declared with `@as_message`, so import the
> modules declaring your messages before building it. A message that arrives with no handler
> fails with `NoHandlerForMessageError` and is retried, then dead-lettered — never acknowledged
> and lost.

### From the console

With the `console` extra, `messenger:consume` does the same from an
[xtr-console](https://github.com/xterr/python-xtr-console) application. Importing
`xtr_messenger.command` declares it; it needs only a `WorkerFactory` to build workers with.

```python
# app/console.py
import app.handlers  # noqa: F401

from app.bus import CONFIG
from xtr_console import Application
from xtr_messenger import WorkerFactory
from xtr_messenger.command import ConsumeMessagesCommand

ConsumeMessagesCommand.use_workers(WorkerFactory(CONFIG))
raise SystemExit(Application("app").run())
```

```sh
uv run python -m app.console messenger:consume high low
uv run python -m app.console messenger:consume high --time-limit 3600
```

SIGTERM, or the time limit running out, stops the worker once the message in hand is settled;
Ctrl-C cancels it. With a [kernel](#kernel--bundle) the `MessengerBundle` builds the command
from the `WorkerFactory` it provides — handlers wired to the container, and the console
bundle picks the command up automatically when both bundles are active.

### Who retries

The loop makes **exactly one attempt per message** and never retries. Redelivery is something
only a transport can do correctly: it owns the delivery count, the backoff state and the
dead-letter destination. A message that fails is rejected, carrying an `ErrorDetailsStamp`
saying why, and the transport decides what happens next. One undeliverable message does not
stop the worker.

On AMQP that means a retry ladder and real dead-lettering. When attempts run out the message
is republished to `taskiq.dlq` — or the `dead_letter_queue` you name — in its original wire
format, so it can be inspected and replayed. A handler that asks for the envelope sees which
attempt it is on in its `RedeliveryStamp`.

### Worker events

Give a worker an event dispatcher — any `EventDispatcherInterface`, such as
[xtr-event-dispatcher](../xtr-event-dispatcher)'s — and it announces itself and every message:

| Event | When | A listener can |
| --- | --- | --- |
| `WorkerStartedEvent` | once, as the worker starts | |
| `WorkerMessageReceivedEvent` | a message was collected | add stamps; `should_handle(value=False)` to skip it |
| `WorkerMessageHandledEvent` | handled, before it is acknowledged | add stamps |
| `WorkerMessageFailedEvent` | handling raised, before it is rejected | read `error`, `will_retry`; add stamps |
| `WorkerRunningEvent` | after each message is settled | `event.worker.stop()` |
| `WorkerStoppedEvent` | once, however the worker stopped | |

```python
from xtr_event_dispatcher import EventDispatcher
from xtr_messenger.event import WorkerMessageFailedEvent

events = EventDispatcher()
events.add_listener(WorkerMessageFailedEvent, lambda e: alert(e.receiver_name, e.error))

worker = WorkerFactory(CONFIG, event_dispatcher=events).worker(["high"])
```

A skipped message is **acknowledged**, not rejected: skipping is a decision, and a rejection
sends a message wherever failed ones go. A listener raising while a message is received or
handled fails that message like a handler would; one raising on a failure is a bug in the
listener, so the message is rejected first and the exception stops the worker.

Each message event names the transport it came from in `receiver_name`. The events are the
same under taskiq's worker, where they come from the task running the message: a failure
there is announced, then raised for taskiq's retry and dead-letter middleware, and
`will_retry` says whether another attempt follows. taskiq runs messages concurrently and
reports none back to the worker, so no `WorkerRunningEvent` is dispatched there.

## Shaping the bus

Every dispatch-side concern is middleware, and which middleware runs is configuration — the
bus and every worker read the same list:

```python
from xtr_messenger import Envelope, MiddlewareInterface, StackInterface


class RejectOutOfHours(MiddlewareInterface):
    async def handle(self, envelope: Envelope, stack: StackInterface, /) -> Envelope:
        if not within_business_hours():
            return envelope  # short-circuit: nothing downstream runs
        return await stack.next().handle(envelope, stack)


CONFIG = MessageBusConfig(
    transports={...},
    routing={...},
    middleware=["logging", RejectOutOfHours()],  # in order, ahead of routing and handling
)
bus = MessageBusFactory(CONFIG, logger=logger).bus()
```

| Setting | Default | Means |
| --- | --- | --- |
| `middleware` | `()` | What runs, in order: a name, middleware already built, or `{name: {argument: value}}` |
| `default_middleware` | `True` | `False` leaves out [holding messages back](#dispatching-after-the-current-message), routing and handling, on the bus and in workers |
| `require_sender` | `False` | Refuse a message routed nowhere, with `NoSenderForMessageError` |
| `handle_unrouted` | `False` | Handle a message routed nowhere in this process instead |
| `require_handler` | `True` | Refuse a message handled here that no handler takes; `False` where processes each handle some of a shared bus's messages |

`"logging"` is the one name the library ships. Add your own with `@as_middleware("audit")` on the
class — built with no argument where a chain names it — or with `named=`, mapping a name to what
builds it: `MessageBusFactory(CONFIG, named={"audit": Audit})`, which wins over a declared name
— or, with a container, with `injectables(middleware=...)`. A name nothing is registered for is refused with
`UnknownMiddlewareError` when the bus or a worker is built.

A name can carry arguments for its middleware's constructor, by parameter name:

```python
middleware = ["audit", {"audit": {"channel": "billing", "sampled": True}}]
```

```python
@as_middleware("audit")
class Audit(MiddlewareInterface):
    def __init__(self, logger: LoggerInterface, channel: str = "app", sampled: bool = False): ...
```

Each entry with arguments is a middleware of its own, beside the bare name. An argument the
constructor does not take — or an entry in another shape — is refused with
`InvalidMiddlewareArgumentsError`; under a kernel, the container's own check of definition
arguments fails the build.

Middleware keeps no per-message state. One instance serves every dispatch, concurrently, and an
instance in the configuration — or a container's one instance of a class — serves the bus and
the workers alike.

`LoggingMiddleware` reports each dispatch through an
[xtr-logging-contracts](https://github.com/xterr/python-xtr-logging-contracts) `LoggerInterface`,
with everything it has to say travelling as context rather than baked into the text. Where those
records go is the application's choice — below, [xtr-logging](https://github.com/xterr/python-xtr-logging):

```python
from xtr_logging import ConsoleHandler, Logger

logger = Logger("messenger", [ConsoleHandler()])
# [2026-09-24T12:30:45+03:00] messenger.NOTICE: message dispatched
#   {"message_type":"IngestDocument","transport":"high","message_id":"id-1"} []
```

**What it says is graded by severity, so a command's `-v` flags decide how much of a dispatch
it shows.** A `ConsoleHandler` prints notices at `-v`, info at `-vv` and everything at `-vvv`:

| Level | Shown at | Records |
| --- | --- | --- |
| `NOTICE` | `-v` | One per dispatch: the message type, the transport it went to, the id the broker gave it |
| `INFO` | `-vv` | One per handler that ran, and what it returned |
| `DEBUG` | `-vvv` | One per dispatch, carrying every stamp the envelope came back with |

Nothing appears at normal verbosity — a dispatch is routine, and a bus running thousands a
second should not say so unasked. Each tier adds what the one above it left out rather than
repeating it, so `-vvv` on a worker reads as a trace and a plain run stays silent.

It runs only where `middleware` names it. The `logger` given to `MessageBusFactory` and
`WorkerFactory` is what `"logging"` writes through; without one it writes to a `NullLogger`.

With a [kernel](#kernel--bundle), name middleware with `@as_middleware("audit")` and the
container supplies each named entry; the logging middleware writes through the
``"messenger"`` channel that the messenger bundle prepends to logging's config. Under
[xtr-console](https://github.com/xterr/python-xtr-console) the same container makes its
console handlers follow every command, so `messenger:consume -vv` reads out every handler
that ran.

What `middleware` names runs in that order, after the middleware that
[holds messages back](#dispatching-after-the-current-message). Routing and handling always come
last, in that order, and two rules carry the producer/consumer split:

- An envelope that arrived **from** a transport carries a `ReceivedStamp` and is never routed
  again, so a consumer cannot re-publish what it consumes.
- Once a sender accepts an envelope the chain **short-circuits**. A message with a transport
  configured is handed off, not also handled locally — unless the sender hands it back
  received, which is all `sync://` does, and then the bus handles it here.

Handlers are called in exactly one place, at the end of that chain. No transport calls one
itself, so what runs is the same whether a message is handled in-process, by the library's
worker, or by a broker's own; and a message routed to be handled with nothing to handle it
fails loudly everywhere, with `NoHandlerForMessageError`.

A message routed nowhere passes quietly by default. `require_sender=True` in the configuration
refuses it with `NoSenderForMessageError`; `handle_unrouted=True` handles it in-process instead.

Routing resolves most specific first: a `TransportNamesStamp` on the envelope, then the table
walking the message's bases, then `"*"`, then whatever the message declared.

### Stamps of your own

A stamp is a frozen dataclass deriving from `StampInterface`. One that should reach the
consumer must be declared, because decoding is an allow-list — a stamp header names a class
the consumer has to build, so an unknown one is dropped rather than resolved:

```python
from dataclasses import dataclass

from xtr_messenger import StampInterface, as_stamp


@as_stamp
@dataclass(frozen=True, slots=True)
class TenantStamp(StampInterface):
    tenant_id: str
```

Every `JsonSerializer` built without `stamp_types` restores the stamps the library ships plus
every declared one — including one declared after the serializer was built. A stamp travels
under its class name, so two declared stamps may not share one. A
`NonSendableStampInterface` stamp never leaves the process and cannot be declared.
`JsonSerializer(stamp_types=[...])` restores exactly the list given, declarations ignored.

### Redispatching a message

`RedispatchMessage` asks for a message to be dispatched again, **through routing**. Something
that produces messages on a worker — a scheduler, a fan-out handler — decides only that a
message goes out; routing decides where:

```python
from xtr_messenger import RedispatchMessage

RedispatchMessage(BuildReport(report_id))  # wherever routing sends BuildReport
RedispatchMessage(BuildReport(report_id), "urgent")  # to the "urgent" transport only
```

Handling it dispatches the carried envelope — or bare message — through a bus that routes,
with the stamps that describe the current process (`ReceivedStamp` above all) removed, and a
`TransportNamesStamp` when transports are named. The handler returns what the carried
message's handler returned, when it was handled rather than sent.

Nothing needs declaring: every bus from `MessageBusFactory` handles a redispatch through
itself, and every worker from `WorkerFactory` through a publishing bus it builds from the same
configuration on the first redispatch — a worker that never redispatches opens no extra
connection. With a [kernel](#kernel--bundle) the container's bus does it. A `RedispatchMessage`
holds an envelope, which no codec carries, so handle it in the process that creates it rather
than routing it to a remote transport.

### Dispatching after the current message

A handler asking for a follow-up usually means "once what I did is final". Stamp it with
`DispatchAfterCurrentBusStamp` and it waits until the message being handled was handled —
every handler, and every middleware around them, a transaction's commit included — then goes
out through the rest of its chain. If that message fails, the follow-up is never dispatched:

```python
from xtr_messenger import DispatchAfterCurrentBusStamp


@as_message_handler(PlaceOrder)
async def place(message: PlaceOrder, bus: Injected[MessageBusInterface]) -> None:
    ...  # record the order
    await bus.dispatch(SendReceipt(message.order_id), DispatchAfterCurrentBusStamp())
```

- Held-back messages go out in the order dispatched; one they dispatch with the stamp joins the
  end of the line. Each is tried even when one before it failed, and then
  `DelayedMessageHandlingError` carries every failure, and the envelope of the message that
  succeeded.
- The message being handled is the first dispatched in the task — on the bus or by a worker,
  whichever bus the handler dispatches through. A stamped message dispatched with nothing
  being handled goes out at once.
- `DispatchAfterCurrentBusMiddleware` does it, first in every chain with `default_middleware`,
  so every configured middleware has finished with the current message first. A worker whose
  message succeeded but whose held-back message failed rejects its message, as any failure.

## Validating messages

Dataclass messages are checked for **shape** by msgspec: a missing or wrongly typed field
raises `MessageDecodingFailedError` at the boundary rather than arriving half-built. Decoding
never coerces — `"3"` is not accepted where an `int` is declared.

A field the message does not declare is **ignored**, which is what lets a producer add one
without redeploying every consumer first. Where both sides ship together and a stray field
means a typo, ask for the other behaviour:

```python
from xtr_messenger import DataclassCodec, JsonSerializer

strict = JsonSerializer(codecs=[DataclassCodec(forbid_unknown_fields=True)])
```

For rules a type cannot express, model the message with pydantic. Install the extra and it is
picked up automatically; both styles work on the same bus. Strictness is the model's own
choice, so a validator written to normalise its input keeps working.

```python
from decimal import Decimal
from typing import Annotated, ClassVar
from uuid import UUID

from pydantic import BaseModel, ConfigDict, Field


@as_message(name="billing.issue_invoice.v1")
class IssueInvoice(BaseModel):
    model_config: ClassVar[ConfigDict] = ConfigDict(frozen=True, extra="forbid")

    invoice_id: UUID
    amount: Annotated[Decimal, Field(gt=0)]
```

> Message annotations are resolved at runtime. If you lint with ruff, set
> `runtime-evaluated-decorators = ["dataclasses.dataclass"]` so field imports are not moved
> into `TYPE_CHECKING` blocks.

### Errors

Everything the library raises derives from `MessageBusError`, and carries what went wrong as
typed attributes rather than only a message.

| Error | Raised when |
| --- | --- |
| `HandlerSignatureError` | A handler is declared with a shape the bus cannot call |
| `InvalidDsnError` | A DSN has no scheme |
| `UnknownTransportOptionError`, `InvalidTransportOptionError` | A setting is not accepted, or its value is not usable |
| `UnsupportedDsnError` | No installed transport serves a scheme |
| `UnknownTransportError` | A route or a worker names a transport neither configured nor registered |
| `UnknownMiddlewareError` | The configuration names middleware nothing is registered for |
| `InvalidMiddlewareArgumentsError` | The configuration gives middleware arguments it does not take |
| `MixedDsnError` | One AMQP worker is asked to serve two servers |
| `NotConsumableError` | A worker is asked to consume a transport that can only send |
| `IncompatibleReceiversError` | A worker is asked to drain a registered receiver beside a transport with its own worker |
| `NoSenderForMessageError` | A message is routed nowhere and the bus requires a sender |
| `NoHandlerForMessageError` | A message is to be handled and nothing handles it |
| `HandlersFailedError` | One or more handlers raised, once every handler has run |
| `DelayedMessageHandlingError` | A message held back with `DispatchAfterCurrentBusStamp` failed once the current one succeeded |
| `MessageEncodingFailedError` | A message cannot be put on the wire |
| `MessageDecodingFailedError`, `UnknownMessageNameError` | A payload cannot be turned back into its message |

## Testing your application

`in-memory://` records instead of sending, and can be drained by a worker, so a test runs the
whole path without a broker. Build the bus and the worker from the same factory instance so
they share the recorder:

```python
from xtr_messenger import InMemoryTransport, InMemoryTransportFactory, WorkerFactory

CONFIG = MessageBusConfig(
    transports={"jobs": TransportConfig("in-memory://?serialize=true")},
    routing={IngestDocument: "jobs"},
)
factories = [InMemoryTransportFactory()]

bus = MessageBusFactory(CONFIG, factories).bus()
_ = await bus.dispatch(IngestDocument(document_id=doc_id, tenant_id=tenant_id))

jobs = factories[0].create(CONFIG.transports)["jobs"]
assert isinstance(jobs, InMemoryTransport)
assert jobs.messages == (IngestDocument(document_id=doc_id, tenant_id=tenant_id),)

await WorkerFactory(CONFIG, factories).worker(["jobs"]).run()  # returns once drained
assert jobs.rejected == ()
```

`?serialize=true` round-trips every message through the serializer on the way in, so a field
that would only fail on a real broker fails in the test instead. `sent` keeps the whole history
while a worker drains the queue, and `clear()` resets it between tests.

## Transports

| Transport | Lives in | Needs | Use for |
| --- | --- | --- | --- |
| `SyncTransport` | `transport/sync/` | — | Handling in-process; local development |
| `InMemoryTransport` | `transport/in_memory/` | — | Tests: records what was dispatched, and can be consumed |
| `TaskiqSender` | `bridge/taskiq/` | `[taskiq]` | Publishing via **any** taskiq broker |
| `TaskiqWorker` | `bridge/taskiq/` | `[taskiq]` | Consuming via taskiq's own worker |
| `create_amqp_broker` | `bridge/amqp/` | `[amqp]` | RabbitMQ, with retries and dead-lettering |

The tree encodes what each piece costs to import:

```
transport/          nothing    — sync, in_memory, the contracts
bridge/taskiq/      [taskiq]   — broker-agnostic: publish, consume, bind a bus
bridge/amqp/        [amqp]     — RabbitMQ only: connection, retry ladder, dead-lettering
```

Everything under `transport/` imports with no extra installed; a transport needing a driver is
a bridge, reachable only once its extra is present. Nothing in `bridge/taskiq/` names a broker
driver either, so `[taskiq]` is usable on its own for Redis, NATS or an in-memory broker. Tests
enforce all three claims.

### Writing your own

Implement two methods and advertise one entry point:

```python
class TransportFactoryInterface(Protocol):
    def supports(self, dsn: Dsn) -> bool: ...
    def create(self, group: Mapping[str, TransportConfig]) -> Mapping[str, SenderInterface]: ...
```

```toml
[project.entry-points."xtr_messenger.transport_factories"]
kafka = "my_package.kafka:KafkaTransportFactory"
```

The entry point **name is the DSN scheme**, which is what keeps discovery lazy: only the module
serving a scheme in use is imported. A sender that can also be consumed implements
`TransportInterface` — `send`, plus `get`, `ack` and `reject` — and the library's `Worker`
drives it.

Settings reach a factory through its method arguments — each `TransportConfig` carries them —
because settings belong to a transport, not to a factory. One discovered `amqp://` factory
serves two transports pointing at different queues with different retry policies.

A *collaborator* cannot arrive that way. A serializer is an object, not a string, so supplying
one means passing the factory yourself:

```python
MessageBusFactory(CONFIG, [AmqpTransportFactory(serializer=mine)]).bus()
```

Handlers never reach a transport at all. Only the bus calls them, and every transport hands
messages to a bus instead — `sync://` hands the envelope straight back marked received, and a
worker dispatches what it receives. So a private registry given to the bus or the worker
reaches every transport from there.

If your broker owns its own consume loop, also implement `WorkerProvidingInterface`, and dispatch
what it receives into the `bus` its `worker()` is given. Most transports should not: a whole
transport is driven by the library's `Worker`.

## Use in an application

Everything adding this package to an application on
[xtr-dependency-injection](../xtr-dependency-injection) takes — and, read backwards, what removing it undoes.

- **Install** — `uv add "xtr-messenger[di,console]"`; add `amqp`, `taskiq` or `pydantic` for
  what you use.
- **Recipe** — `uv run xtr-recipes recipes:sync` does the *Activate*, *Configure* and
  *Environment* steps below: it lists `MessengerBundle`, writes a starting `<app>/config/messenger.py`
  and `MESSENGER_DSN` (commented out) in `.env`. It prints the step to name and route your
  transports, which a recipe cannot make for you.
- **Activate** — `MessengerBundle: {"all": True}` in `BUNDLES` in `<app>/bundles.py`, imported
  from `xtr_messenger.bundle`.
- **Brings along** — the logging, console and event dispatcher bundles, when those packages
  are installed.
- **Configure** — needed to move a message: with no configuration there are no transports, and
  a message routed nowhere is neither sent nor handled. Transports and routing go in
  `<app>/config/messenger.py`, a `@configure` function returning `MessageBusConfig` — see
  [Kernel / bundle](#kernel--bundle).
- **Environment** — nothing required; a broker DSN is usually `env("MESSENGER_DSN")`, read
  only when the bus is built.
- **Ignore** — nothing.
- **Run** — a worker is `<script> messenger:consume <transport>`.
- **Remove** — drop the `BUNDLES` entry, delete `<app>/config/messenger.py`, then
  `uv remove xtr-messenger` — unless xtr-scheduler is installed, which depends on it.
- **Check** — `debug:bundles` shows `messenger` as `listed` and `active`; `debug:config
  messenger` shows the resolved transports and routing.

## Kernel / bundle

An application using [xtr-dependency-injection](../xtr-dependency-injection) lists
`MessengerBundle` in its `app/bundles.py` and configures it with `@configure`. Handlers ask
for what they need the way any container-injected service does; import nothing from this
library's integration:

```sh
uv add "xtr-messenger[di]"
```

```python
# app/bundles.py
from xtr_messenger.bundle import MessengerBundle

BUNDLES = {MessengerBundle: {"all": True}}
```

```python
# app/config/messenger.py
from xtr_dependency_injection import configure

from xtr_messenger import MessageBusConfig, TransportConfig

from app.messages import IngestDocument


@configure
def messenger() -> MessageBusConfig:
    return MessageBusConfig(
        transports={"jobs": TransportConfig("amqp://queue")},
        routing={IngestDocument: "jobs"},
    )
```

```python
# app/handlers.py
from xtr_dependency_injection import Injected, as_service

from xtr_messenger import as_message_handler


@as_service
class InvoiceRepository: ...


@as_message_handler(IngestDocument)
async def ingest(message: IngestDocument, db: Injected[Session]) -> None:
    await db.record(message.document_id)


@as_message_handler(IssueInvoice)
class IssueInvoiceHandler:
    def __init__(self, invoices: InvoiceRepository) -> None: ...

    async def __call__(self, message: IssueInvoice, db: Injected[Session]) -> None: ...
```

The bundle registers a `MessageBusInterface`, a `WorkerFactory` (and the `messenger:consume`
command when the console bundle is active), and a per-kernel `HandlersLocator`. Every
handler its scan finds is bound with `bind_callable` at boot, so a handler asking for
something the container cannot provide fails at boot, not on its first message.

Container dependencies reach a handler only through a container marker. A **function
handler** takes the message — and the `Envelope`, if it asks for one — as ordinary
parameters; every other parameter it needs from the container **must** be annotated
`Injected[T]`, `Annotated[T, Target("name")]` for a qualified service, or
`Annotated[T, Autowire(param=... | env=...)]`, because a bare `T` is not injected and would
arrive unfilled:

```python
from xtr_dependency_injection import Injected

from xtr_messenger import Envelope, as_message_handler


@as_message_handler(IngestDocument)
async def ingest(
    message: IngestDocument,
    envelope: Envelope,  # the message and its envelope arrive as-is
    db: Injected[Session],  # everything from the container is Injected[...]
    metrics: Injected[Metrics],
) -> None:
    await db.record(message.document_id)
```

A **handler class** is a singleton: its constructor takes what lives as long as the handler
(plain parameters), and anything a single message needs goes on `__call__` as `Injected[T]`.

`@required_bundle` pulls in the logging, console and event dispatcher bundles when installed;
the logging middleware writes through a ``"messenger"`` channel added to logging's config
automatically, and with the event dispatcher bundle active every worker announces
[its events](#worker-events) through the container's dispatcher — a subscriber in the
application hears them with no wiring.

Middleware referred to by name in the config is resolved from the container: give the class
a name with `@as_middleware("audit")`, or register it under `(MiddlewareInterface, name)`
manually. The bundle uses a `ServiceLocator` for that lookup, so a class is built only when
the configuration names it. Another package can ship middleware this way — its bundle loads
the module declaring it — and an application names it in its chain, without this package
knowing it exists. An entry giving a middleware arguments is registered as a middleware of its
own, built by the container with those arguments set; one built by a factory takes none.

**Every message is a unit of work.** The bundle puts `UnitOfWorkMiddleware` ahead of the
configured middleware on the bus and in every worker — right after the one
[holding messages back](#dispatching-after-the-current-message), so a message held back is a
unit of its own — and the middleware after it and every handler of a message share one
[unit of work](../xtr-dependency-injection#units-of-work): a `lifetime="scoped"` service — a
database session — is built once per message, handed to each handler asking for it, and
released when the message is done with, even when a handler raised. A message dispatched while
another is handled — or while a request or a command runs — joins that unit; a message a worker
received is always a unit of its own, even when a command runs the worker. A middleware needing the message's instance resolves it
from `current_unit_of_work()`.

A class in the application — or another bundle — that implements `TransportFactoryInterface`
is registered automatically and consulted **ahead of** the factories discovery finds by entry
point, in registration order. That is how an app serves a private DSN scheme, or overrides a
discovered one, without passing a factory list by hand — the class needs no decorator beyond
living in a scanned module and being buildable with no arguments. With none registered the
bundle falls back to entry-point discovery, so a zero-config app keeps working:

```python
from collections.abc import Mapping

from typing_extensions import override

from xtr_messenger import Dsn, TransportConfig, TransportFactoryInterface
from xtr_messenger.transport.sender import SenderInterface


class KafkaTransportFactory(TransportFactoryInterface):
    @override
    def supports(self, dsn: Dsn) -> bool:
        return dsn.scheme == "kafka"

    @override
    def create(
        self, group: Mapping[str, TransportConfig]
    ) -> Mapping[str, SenderInterface]: ...  # build a sender per name in `group`
```

A service implementing `ReceiverInterface` and tagged `RECEIVER_TAG` (`"messenger.receiver"`)
with an `alias` is consumable under that alias — `messenger:consume <alias>` — with no entry in
`MessageBusConfig.transports`. That is for something that only receives: messages a process
generates rather than reads from a broker, which have no DSN to configure. Another bundle can
register one this way for the application; a configured transport of the same name wins, which
is how an application overrides it. Every tagged receiver is built with the `WorkerFactory`, so
building one must do no I/O. Two receivers under one alias fail the build.

```python
from xtr_dependency_injection import as_service, autoconfigure

from xtr_messenger.bundle import RECEIVER_TAG


@autoconfigure(tags=[(RECEIVER_TAG, {"alias": "ticks"})])
@as_service
class TickReceiver(ReceiverInterface): ...
```

Without a container, pass them to the factory — `WorkerFactory(CONFIG, receivers={"ticks":
TickReceiver()})`. One worker may drain registered receivers and configured transports
together, except a transport that brings its own worker, which cannot share it
(`IncompatibleReceiversError`).

Between messages a worker calls `ServicesResetter.reset()`, so services opting in with
`ResetInterface` — or explicitly tagged `kernel.reset` — are cleared per unit of work.

## Layout

```
xtr_messenger/
├── envelope.py              the message plus its stamps
├── message_registry.py      what @as_message declared: names, default transports
├── stamp_registry.py        what @as_stamp declared: stamps serializers restore
├── message_bus_config.py    which transports exist, and where messages go
├── message_bus_factory.py   builds the bus a process publishes through
├── xtr_messenger.py           the loop — wrap, then walk the middleware chain
├── worker_factory.py        builds what a worker process runs
├── worker.py                the other loop — collect, dispatch, ack or reject
├── dsn.py                   reading a transport's DSN
├── decorator/               @as_message, @as_message_handler, @as_middleware, @as_stamp
├── event/                   what a worker announces about itself and each message
├── message/                 messages the library handles itself: RedispatchMessage
├── handler/                 which function handles which message, and how to call it
├── middleware/              holding back, routing, handling, logging, units of work, the cursor
├── stamp/                   one class per module
├── exception/               one error per module, all a MessageBusError
├── transport/
│   ├── sender/              SenderInterface, SendersLocator — the routing table
│   ├── receiver/            ReceiverInterface, ChainedReceiver
│   ├── sync/  in_memory/    a transport and its factory, each
│   └── serialization/       the wire format, and codec/ for message shapes
├── bridge/
│   ├── taskiq/              broker-agnostic publish and consume
│   └── amqp/                RabbitMQ on top of it
├── command/
│   └── consume.py           messenger:consume, on xtr-console
└── bundle/                  MessengerBundle for xtr-dependency-injection
```

## Development

Developed in the [python-xtr](https://github.com/xterr/python-xtr) monorepo, under
`packages/xtr-messenger`; run the commands below from there. The `python-xtr-messenger` repository is a
read-only copy, so send issues and pull requests to the monorepo.

```sh
uv sync --all-extras
uv run ruff check src tests
uv run ruff format --check src tests
uv run ty check
uv run basedpyright src tests
uv run pytest
```

Two type checkers on purpose — they disagree often enough to be worth both, and `ty` has
already caught a crash `basedpyright` accepted. Both run strict on the tests too, with nothing
suppressed.

The suite mirrors the source tree. `tests/unit/` holds a `test_<module>.py` for each module,
testing it alone against small fakes; `tests/integration/` holds what needs several real
components or a fresh interpreter — end-to-end flows, and the checks that a module never
imports a library it should not. Shared message classes and fakes live in `tests/support/`.

## License

[MIT](LICENSE) © xterr
