async def build_message(
message: "SendableMessage",
topic: str,
partition: int | None = None,
timestamp_ms: int | None = None,
key: bytes | None = None,
headers: dict[str, str] | None = None,
correlation_id: str | None = None,
*,
reply_to: str = "",
serializer: Optional["SerializerProto"],
codec: Optional["CodecProto"] = None,
id_generator: IdGenerator = gen_cor_id,
) -> "ConsumerRecord":
"""Build a Kafka ConsumerRecord for a sendable message."""
if message is None and key is not None:
# keyed None is a real tombstone, matching publish()'s own rule
# (aiokafka needs a key or value, a keyless None still goes b"")
msg, content_type = None, None
else:
msg, content_type = await (codec or DefaultCodec()).encode(message, serializer)
k = key or b""
headers = {
"content-type": content_type or "",
"correlation_id": correlation_id or id_generator(),
**(headers or {}),
}
if reply_to:
headers["reply_to"] = headers.get("reply_to", reply_to)
return ConsumerRecord(
value=msg,
topic=topic,
partition=partition or 0,
key=k,
serialized_key_size=len(k),
serialized_value_size=0 if msg is None else len(msg),
checksum=0 if msg is None else sum(msg),
offset=0,
headers=[(i, j.encode()) for i, j in headers.items()],
timestamp_type=1,
timestamp=timestamp_ms or int(datetime.now(timezone.utc).timestamp() * 1000),
)