async def build_message(
message: "SendableMessage",
topic: str,
*,
correlation_id: str | None = None,
partition: int | None = None,
timestamp_ms: int | None = None,
key: bytes | str | None = None,
headers: dict[str, str] | None = None,
reply_to: str = "",
serializer: Optional["SerializerProto"] = None,
codec: Optional["CodecProto"] = None,
) -> MockConfluentMessage:
"""Build a mock confluent_kafka.Message for a sendable message."""
if message is None:
# keep a real tombstone (message.value() is None) distinct from b""
msg, content_type = None, None
else:
codec_instance = codec or DefaultCodec()
msg, content_type = await codec_instance.encode(message, serializer)
k = key or b""
headers = {
"content-type": content_type or "",
"correlation_id": correlation_id or gen_cor_id(),
"reply_to": reply_to,
**(headers or {}),
}
# https://docs.confluent.io/platform/current/clients/confluent-kafka-python/html/index.html#confluent_kafka.Message.timestamp
return MockConfluentMessage(
raw_msg=msg,
topic=topic,
key=k,
headers=[(i, j.encode()) for i, j in headers.items()],
offset=0,
partition=partition or 0,
timestamp_type=1,
timestamp_ms=timestamp_ms or int(datetime.now(timezone.utc).timestamp() * 1000),
)