Skip to content

LogicPublisher

faststream.kafka.publisher.usecase.LogicPublisher #

LogicPublisher(
    config: KafkaPublisherConfig,
    specification: PublisherSpecification[Any, Any],
)

Bases: PublisherUsecase

A class to publish messages to a Kafka topic.

Source code in faststream/kafka/publisher/usecase.py
def __init__(
    self,
    config: "KafkaPublisherConfig",
    specification: "PublisherSpecification[Any, Any]",
) -> None:
    super().__init__(config, specification)

    self._topic = config.topic
    self.partition = config.partition
    self.reply_to = config.reply_to
    self.headers = config.headers or {}

partition instance-attribute #

partition = config.partition

reply_to instance-attribute #

reply_to = config.reply_to

headers instance-attribute #

headers = config.headers or {}

topic property #

topic: str

is_test instance-attribute #

is_test = False

mock property #

mock: MagicMock

The mock recording the endpoint's calls, available under a test broker.

specification instance-attribute #

specification = specification

request async #

request(
    message: SendableMessage,
    topic: str = "",
    *,
    key: bytes | Any | None = None,
    partition: int | None = None,
    timestamp_ms: int | None = None,
    headers: dict[str, str] | None = None,
    correlation_id: str | None = None,
    timeout: float = 0.5,
) -> KafkaMessage

Send a request message to Kafka topic.

PARAMETER DESCRIPTION
message

Message body to send.

TYPE: SendableMessage

topic

Topic where the message will be published.

TYPE: str DEFAULT: ''

key

A key to associate with the message. Can be used to determine which partition to send the message to. If partition is None (and producer's partitioner config is left as default), then messages with the same key will be delivered to the same partition (but if key is None, partition is chosen randomly). Must be type bytes, or be serializable to bytes via configured key_serializer.

TYPE: bytes | Any | None DEFAULT: None

partition

Specify a partition. If not set, the partition will be selected using the configured partitioner.

TYPE: int | None DEFAULT: None

timestamp_ms

Epoch milliseconds (from Jan 1 1970 UTC) to use as the message timestamp. Defaults to current time.

TYPE: int | None DEFAULT: None

headers

Message headers to store metainformation.

TYPE: dict[str, str] | None DEFAULT: None

correlation_id

Manual message correlation_id setter. correlation_id is a useful option to trace messages.

TYPE: str | None DEFAULT: None

timeout

Timeout to send RPC request.

TYPE: float DEFAULT: 0.5

RETURNS DESCRIPTION
KafkaMessage

The response message.

TYPE: KafkaMessage

Source code in faststream/kafka/publisher/usecase.py
@override
async def request(
    self,
    message: "SendableMessage",
    topic: str = "",
    *,
    key: bytes | Any | None = None,
    partition: int | None = None,
    timestamp_ms: int | None = None,
    headers: dict[str, str] | None = None,
    correlation_id: str | None = None,
    timeout: float = 0.5,
) -> "KafkaMessage":
    """Send a request message to Kafka topic.

    Args:
        message: Message body to send.
        topic: Topic where the message will be published.
        key: A key to associate with the message. Can be used to
            determine which partition to send the message to. If partition
            is `None` (and producer's partitioner config is left as default),
            then messages with the same key will be delivered to the same
            partition (but if key is `None`, partition is chosen randomly).
            Must be type `bytes`, or be serializable to bytes via configured
            `key_serializer`.
        partition: Specify a partition. If not set, the partition will be
            selected using the configured `partitioner`.
        timestamp_ms: Epoch milliseconds (from Jan 1 1970 UTC) to use as
            the message timestamp. Defaults to current time.
        headers: Message headers to store metainformation.
        correlation_id: Manual message **correlation_id** setter.
            **correlation_id** is a useful option to trace messages.
        timeout: Timeout to send RPC request.

    Returns:
        KafkaMessage: The response message.
    """
    cmd = KafkaPublishCommand(
        message,
        topic=topic or self.topic,
        key=key,
        partition=partition if partition is not None else self.partition,
        headers=self.headers | (headers or {}),
        correlation_id=correlation_id or self._outer_config.id_generator(),
        timestamp_ms=timestamp_ms,
        timeout=timeout,
        _publish_type=PublishType.REQUEST,
    )

    msg: KafkaMessage = await self._basic_request(
        cmd,
        producer=self._outer_config.producer,
    )
    return msg

flush async #

flush() -> None
Source code in faststream/kafka/publisher/usecase.py
async def flush(self) -> None:
    producer = cast("AioKafkaFastProducer", self._outer_config.producer)
    await producer.flush()

publish abstractmethod async #

publish(
    message: SendableMessage,
    /,
    *,
    correlation_id: str | None = None,
) -> Any | None

Public method to publish a message.

Should be called by user only broker.publisher(...).publish(...).

Source code in faststream/_internal/endpoint/publisher/proto.py
@abstractmethod
async def publish(
    self,
    message: "SendableMessage",
    /,
    *,
    correlation_id: str | None = None,
) -> Any | None:
    """Public method to publish a message.

    Should be called by user only `broker.publisher(...).publish(...)`.
    """
    ...

assert_called_once_with async #

assert_called_once_with(
    body: Any = EMPTY,
    /,
    *,
    headers: Any = EMPTY,
    correlation_id: Any = EMPTY,
    reply_to: Any = EMPTY,
    content_type: Any = EMPTY,
    path: Any = EMPTY,
    context: Mapping[str, Any] = EMPTY,
) -> None

Assert the endpoint was called once, with the message described here.

PARAMETER DESCRIPTION
body

The body as a dict, a model or a matcher; it goes through the codec.

TYPE: Any DEFAULT: EMPTY

headers

Headers the message must carry; the rest may carry more.

TYPE: Any DEFAULT: EMPTY

correlation_id

The exact correlation id.

TYPE: Any DEFAULT: EMPTY

reply_to

The exact reply-to destination.

TYPE: Any DEFAULT: EMPTY

content_type

The exact content type.

TYPE: Any DEFAULT: EMPTY

path

The exact path parameters the subject template matched.

TYPE: Any DEFAULT: EMPTY

context

Context paths, as given to Context(), mapped to their values.

TYPE: Mapping[str, Any] DEFAULT: EMPTY

Source code in faststream/_internal/testing/calls.py
async def assert_called_once_with(
    self,
    body: Any = EMPTY,
    /,
    *,
    headers: Any = EMPTY,
    correlation_id: Any = EMPTY,
    reply_to: Any = EMPTY,
    content_type: Any = EMPTY,
    path: Any = EMPTY,
    context: Mapping[str, Any] = EMPTY,
) -> None:
    """Assert the endpoint was called once, with the message described here.

    Args:
        body: The body as a dict, a model or a matcher; it goes through the codec.
        headers: Headers the message must carry; the rest may carry more.
        correlation_id: The exact correlation id.
        reply_to: The exact reply-to destination.
        content_type: The exact content type.
        path: The exact path parameters the subject template matched.
        context: Context paths, as given to `Context()`, mapped to their values.
    """
    recorder = self._recorder_with_calls()
    recorder.mock.assert_called_once()
    await recorder.assert_last_call(
        ExpectedCall(
            body=body,
            headers=headers,
            correlation_id=correlation_id,
            reply_to=reply_to,
            content_type=content_type,
            path=path,
            context=context,
        )
    )

assert_called_with async #

assert_called_with(
    body: Any = EMPTY,
    /,
    *,
    headers: Any = EMPTY,
    correlation_id: Any = EMPTY,
    reply_to: Any = EMPTY,
    content_type: Any = EMPTY,
    path: Any = EMPTY,
    context: Mapping[str, Any] = EMPTY,
) -> None

Assert the last message the endpoint saw is the one described here.

PARAMETER DESCRIPTION
body

The body as a dict, a model or a matcher; it goes through the codec.

TYPE: Any DEFAULT: EMPTY

headers

Headers the message must carry; the rest may carry more.

TYPE: Any DEFAULT: EMPTY

correlation_id

The exact correlation id.

TYPE: Any DEFAULT: EMPTY

reply_to

The exact reply-to destination.

TYPE: Any DEFAULT: EMPTY

content_type

The exact content type.

TYPE: Any DEFAULT: EMPTY

path

The exact path parameters the subject template matched.

TYPE: Any DEFAULT: EMPTY

context

Context paths, as given to Context(), mapped to their values.

TYPE: Mapping[str, Any] DEFAULT: EMPTY

Source code in faststream/_internal/testing/calls.py
async def assert_called_with(
    self,
    body: Any = EMPTY,
    /,
    *,
    headers: Any = EMPTY,
    correlation_id: Any = EMPTY,
    reply_to: Any = EMPTY,
    content_type: Any = EMPTY,
    path: Any = EMPTY,
    context: Mapping[str, Any] = EMPTY,
) -> None:
    """Assert the last message the endpoint saw is the one described here.

    Args:
        body: The body as a dict, a model or a matcher; it goes through the codec.
        headers: Headers the message must carry; the rest may carry more.
        correlation_id: The exact correlation id.
        reply_to: The exact reply-to destination.
        content_type: The exact content type.
        path: The exact path parameters the subject template matched.
        context: Context paths, as given to `Context()`, mapped to their values.
    """
    recorder = self._recorder_with_calls()
    await recorder.assert_last_call(
        ExpectedCall(
            body=body,
            headers=headers,
            correlation_id=correlation_id,
            reply_to=reply_to,
            content_type=content_type,
            path=path,
            context=context,
        )
    )

assert_any_call async #

assert_any_call(
    body: Any = EMPTY,
    /,
    *,
    headers: Any = EMPTY,
    correlation_id: Any = EMPTY,
    reply_to: Any = EMPTY,
    content_type: Any = EMPTY,
    path: Any = EMPTY,
    context: Mapping[str, Any] = EMPTY,
) -> None

Assert one of the messages the endpoint saw is the one described here.

PARAMETER DESCRIPTION
body

The body as a dict, a model or a matcher; it goes through the codec.

TYPE: Any DEFAULT: EMPTY

headers

Headers the message must carry; the rest may carry more.

TYPE: Any DEFAULT: EMPTY

correlation_id

The exact correlation id.

TYPE: Any DEFAULT: EMPTY

reply_to

The exact reply-to destination.

TYPE: Any DEFAULT: EMPTY

content_type

The exact content type.

TYPE: Any DEFAULT: EMPTY

path

The exact path parameters the subject template matched.

TYPE: Any DEFAULT: EMPTY

context

Context paths, as given to Context(), mapped to their values.

TYPE: Mapping[str, Any] DEFAULT: EMPTY

Source code in faststream/_internal/testing/calls.py
async def assert_any_call(
    self,
    body: Any = EMPTY,
    /,
    *,
    headers: Any = EMPTY,
    correlation_id: Any = EMPTY,
    reply_to: Any = EMPTY,
    content_type: Any = EMPTY,
    path: Any = EMPTY,
    context: Mapping[str, Any] = EMPTY,
) -> None:
    """Assert one of the messages the endpoint saw is the one described here.

    Args:
        body: The body as a dict, a model or a matcher; it goes through the codec.
        headers: Headers the message must carry; the rest may carry more.
        correlation_id: The exact correlation id.
        reply_to: The exact reply-to destination.
        content_type: The exact content type.
        path: The exact path parameters the subject template matched.
        context: Context paths, as given to `Context()`, mapped to their values.
    """
    recorder = self._recorder_with_calls()
    await recorder.assert_any_call(
        ExpectedCall(
            body=body,
            headers=headers,
            correlation_id=correlation_id,
            reply_to=reply_to,
            content_type=content_type,
            path=path,
            context=context,
        )
    )

start async #

start() -> None
Source code in faststream/_internal/endpoint/publisher/usecase.py
async def start(self) -> None:
    pass

set_test #

set_test(
    *, recorder: CallRecorder | None = None, with_fake: bool
) -> None

Turn publisher to testing mode, sharing recorder when one is given.

Source code in faststream/_internal/endpoint/publisher/usecase.py
def set_test(
    self,
    *,
    recorder: CallRecorder | None = None,
    with_fake: bool,
) -> None:
    """Turn publisher to testing mode, sharing `recorder` when one is given."""
    self.is_test = True
    if recorder is None:
        self._recorder.reset()
    else:
        self._recorder = recorder
    self._fake_handler = with_fake

reset_test #

reset_test() -> None

Turn off publisher's testing mode.

Source code in faststream/_internal/endpoint/publisher/usecase.py
def reset_test(self) -> None:
    """Turn off publisher's testing mode."""
    self.is_test = False
    self._recorder.reset()
    # A shared recorder goes back to the handler it belongs to
    self._recorder = CallRecorder(self.specification.name, self._outer_config)
    self._fake_handler = False

schema #

schema() -> dict[str, PublisherSpec]
Source code in faststream/_internal/endpoint/publisher/usecase.py
def schema(self) -> dict[str, "PublisherSpec"]:
    return self.specification.get_schema()