Skip to content

RedisPublishCommand

faststream.redis.response.RedisPublishCommand #

RedisPublishCommand(
    message: SendableMessage,
    /,
    *messages: SendableMessage,
    _publish_type: PublishType,
    correlation_id: str | None = None,
    channel: str | None = None,
    list: str | None = None,
    stream: str | None = None,
    maxlen: int | None = None,
    headers: dict[str, Any] | None = None,
    reply_to: str = "",
    timeout: float | None = 30.0,
    pipeline: Optional[Pipeline[bytes]] = None,
    message_format: type[
        MessageFormat
    ] = BinaryMessageFormatV1,
)

Bases: BatchPublishCommand

Source code in faststream/redis/response.py
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
def __init__(
    self,
    message: "SendableMessage",
    /,
    *messages: "SendableMessage",
    _publish_type: "PublishType",
    correlation_id: str | None = None,
    channel: str | None = None,
    list: str | None = None,
    stream: str | None = None,
    maxlen: int | None = None,
    headers: dict[str, Any] | None = None,
    reply_to: str = "",
    timeout: float | None = 30.0,
    pipeline: Optional["Pipeline[bytes]"] = None,
    message_format: type["MessageFormat"] = BinaryMessageFormatV1,
) -> None:
    super().__init__(
        message,
        *messages,
        _publish_type=_publish_type,
        correlation_id=correlation_id,
        reply_to=reply_to,
        destination="",
        headers=headers,
    )

    self.pipeline = pipeline

    self.message_format = message_format

    self.set_destination(
        channel=channel,
        list=list,
        stream=stream,
    )

    # Stream option
    self.maxlen = maxlen

    # Request option
    self.timeout = timeout

destination_type instance-attribute #

destination_type: DestinationType

pipeline instance-attribute #

pipeline = pipeline

message_format instance-attribute #

message_format = message_format

maxlen instance-attribute #

maxlen = maxlen

timeout instance-attribute #

timeout = timeout

body instance-attribute #

body = body

headers instance-attribute #

headers = headers or {}

correlation_id instance-attribute #

correlation_id = correlation_id

destination instance-attribute #

destination = destination

reply_to instance-attribute #

reply_to = reply_to

publish_type instance-attribute #

publish_type = _publish_type

batch_bodies property writable #

batch_bodies: tuple[Any, ...]

extra_bodies instance-attribute #

extra_bodies = bodies

set_destination #

set_destination(
    *,
    channel: str | None = None,
    list: str | None = None,
    stream: str | None = None,
) -> None
Source code in faststream/redis/response.py
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
def set_destination(
    self,
    *,
    channel: str | None = None,
    list: str | None = None,
    stream: str | None = None,
) -> None:
    if channel is not None:
        self.destination_type = DestinationType.Channel
        self.destination = channel
    elif list is not None:
        self.destination_type = DestinationType.List
        self.destination = list
    elif stream is not None:
        self.destination_type = DestinationType.Stream
        self.destination = stream
    else:
        raise SetupError(INCORRECT_SETUP_MSG)

from_cmd classmethod #

from_cmd(
    cmd: Union[PublishCommand, RedisPublishCommand],
    *,
    batch: bool = False,
    message_format: type[
        MessageFormat
    ] = BinaryMessageFormatV1,
) -> RedisPublishCommand
Source code in faststream/redis/response.py
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
@classmethod
def from_cmd(
    cls,
    cmd: Union["PublishCommand", "RedisPublishCommand"],
    *,
    batch: bool = False,
    message_format: type["MessageFormat"] = BinaryMessageFormatV1,
) -> "RedisPublishCommand":
    if isinstance(cmd, RedisPublishCommand):
        # NOTE: Should return a copy probably.
        return cmd

    body, extra_bodies = cls._parse_bodies(cmd.body, batch=batch)

    return cls(
        body,
        *extra_bodies,
        channel=cmd.destination,
        correlation_id=cmd.correlation_id,
        headers=cmd.headers,
        reply_to=cmd.reply_to,
        message_format=message_format,
        _publish_type=cmd.publish_type,
    )

as_publish_command #

as_publish_command() -> PublishCommand

Method to transform handlers' Response result to DTO for publishers.

Source code in faststream/response/response.py
22
23
24
25
26
27
28
29
def as_publish_command(self) -> "PublishCommand":
    """Method to transform handlers' Response result to DTO for publishers."""
    return PublishCommand(
        body=self.body,
        headers=self.headers,
        correlation_id=self.correlation_id,
        _publish_type=PublishType.PUBLISH,
    )

get_publish_key #

get_publish_key() -> Any | None

Get the key for publishing this message.

Override this method in subclasses to provide broker-specific keys. Default implementation returns None (no key).

RETURNS DESCRIPTION
Any | None

The key for publishing, or None if this Response type doesn't use keys.

Source code in faststream/response/response.py
31
32
33
34
35
36
37
38
39
40
def get_publish_key(self) -> Any | None:  # noqa: PLR6301
    """Get the key for publishing this message.

    Override this method in subclasses to provide broker-specific keys.
    Default implementation returns None (no key).

    Returns:
        The key for publishing, or None if this Response type doesn't use keys.
    """
    return None

add_headers #

add_headers(
    headers: dict[str, Any], *, override: bool = True
) -> None
Source code in faststream/response/response.py
71
72
73
74
75
76
77
78
79
80
def add_headers(
    self,
    headers: dict[str, Any],
    *,
    override: bool = True,
) -> None:
    if override:
        self.headers |= headers
    else:
        self.headers = headers | self.headers