Bases: ABC
A class to represent a raw Redis message.
Source code in faststream/redis/parser/message.py
| def __init__(
self,
data: bytes,
headers: dict[str, Any] | None = None,
) -> None:
self.data = data
self.headers = headers or {}
|
build(
*,
message: Union[
Sequence[SendableMessage], SendableMessage
],
reply_to: str | None,
headers: dict[str, Any] | None,
correlation_id: str,
serializer: Optional[SerializerProto] = None,
codec: Optional[CodecProto] = None,
) -> MessageFormat
Source code in faststream/redis/parser/message.py
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60 | @classmethod
async def build(
cls,
*,
message: Union[Sequence["SendableMessage"], "SendableMessage"],
reply_to: str | None,
headers: dict[str, Any] | None,
correlation_id: str,
serializer: Optional["SerializerProto"] = None,
codec: Optional["CodecProto"] = None,
) -> "MessageFormat":
codec_instance = codec or DefaultCodec()
payload, content_type = await codec_instance.encode(message, serializer) # type: ignore[arg-type]
headers_to_send = {
"correlation_id": correlation_id,
}
if content_type:
headers_to_send["content-type"] = content_type
if reply_to:
headers_to_send["reply_to"] = reply_to
if headers is not None:
headers_to_send.update(headers)
return cls(
data=payload,
headers=headers_to_send,
)
|
encode(
*,
message: Union[
Sequence[SendableMessage], SendableMessage
],
reply_to: str | None,
headers: dict[str, Any] | None,
correlation_id: str,
serializer: Optional[SerializerProto] = None,
codec: Optional[CodecProto] = None,
) -> bytes
Source code in faststream/redis/parser/message.py
62
63
64
65
66
67
68
69
70
71
72
73
74 | @classmethod
@abstractmethod
async def encode(
cls,
*,
message: Union[Sequence["SendableMessage"], "SendableMessage"],
reply_to: str | None,
headers: dict[str, Any] | None,
correlation_id: str,
serializer: Optional["SerializerProto"] = None,
codec: Optional["CodecProto"] = None,
) -> bytes:
raise NotImplementedError
|
parse(data: bytes) -> tuple[bytes, dict[str, Any]]
Source code in faststream/redis/parser/message.py
| @classmethod
@abstractmethod
def parse(cls, data: bytes) -> tuple[bytes, dict[str, Any]]:
raise NotImplementedError
|