Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions docs/docs/en/confluent/message.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,19 @@ This object serves as a unified **FastStream** wrapper around the native broker
* `#!python topic(): str`
* `#!python value(): Optional[Union[str, bytes]]`

!!! note
A record with a `None value()` is a Kafka tombstone, the delete marker on a compacted topic. `#!python msg.tombstone` is `#!python True` for it, keeping it distinct from an empty payload (`#!python b""`), and a `#!python None`-able body parameter of a **FastAPI** subscriber resolves to `#!python None` instead of failing validation.

`#!python msg.body` is `#!python TOMBSTONE` for such a record: an empty `#!python bytes` subclass, so it still equals `#!python b""` and existing handlers are unaffected. `#!python isinstance(body, Tombstone)` is what tells a tombstone from an empty payload - and it is the only thing that does in a `#!python batch=True` subscriber, where each record is marked on its own and `#!python msg.tombstone` stays `#!python False` for the batch as a whole.

```python
from faststream.confluent import TOMBSTONE, Tombstone

await broker.publish(TOMBSTONE, "topic", key=b"user-1")
```

An explicit `#!python TOMBSTONE` requires a key, since compaction deletes per key. It can also ride in a `#!python publish_batch()`, but not together with a custom `BatchCodecProto`, whose `#!python encode_batch()` has no way to express a null value for a single record.

For example, if you would like to access the headers of an incoming message, you would do so like this:

```python hl_lines="1 6"
Expand Down
13 changes: 13 additions & 0 deletions docs/docs/en/kafka/message.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,19 @@ This object serves as a unified **FastStream** wrapper around the native broker
* `#!python topic: str`
* `#!python value: Optional[aiokafka.structs.VT]`

!!! note
A record with a `None value` is a Kafka tombstone, the delete marker on a compacted topic. `#!python msg.tombstone` is `#!python True` for it, keeping it distinct from an empty payload (`#!python b""`), and a `#!python None`-able body parameter of a **FastAPI** subscriber resolves to `#!python None` instead of failing validation.

`#!python msg.body` is `#!python TOMBSTONE` for such a record: an empty `#!python bytes` subclass, so it still equals `#!python b""` and existing handlers are unaffected. `#!python isinstance(body, Tombstone)` is what tells a tombstone from an empty payload - and it is the only thing that does in a `#!python batch=True` subscriber, where each record is marked on its own and `#!python msg.tombstone` stays `#!python False` for the batch as a whole.

```python
from faststream.kafka import TOMBSTONE, Tombstone

await broker.publish(TOMBSTONE, "topic", key=b"user-1")
```

An explicit `#!python TOMBSTONE` requires a key, since compaction deletes per key. It can also ride in a `#!python publish_batch()`, but not together with a custom `BatchCodecProto`, whose `#!python encode_batch()` has no way to express a null value for a single record.

For example, if you would like to access the headers of an incoming message, you would do so like this:

```python hl_lines="1 6"
Expand Down
10 changes: 6 additions & 4 deletions faststream/_internal/fastapi/route.py
Original file line number Diff line number Diff line change
Expand Up @@ -51,13 +51,13 @@ class StreamMessage(Request):
scope: "dict[str, Any]"
_cookies: "dict[str, Any]"
_headers: "dict[str, Any]" # type: ignore[assignment]
_body: Union["dict[str, Any]", list[Any]] # type: ignore[assignment]
_body: Union["dict[str, Any]", list[Any], None] # type: ignore[assignment]
_query_params: "dict[str, Any]" # type: ignore[assignment]

def __init__(
self,
*,
body: Union["dict[str, Any]", list[Any]],
body: Union["dict[str, Any]", list[Any], None],
headers: "dict[str, Any]",
path: "dict[str, Any]",
) -> None:
Expand Down Expand Up @@ -172,9 +172,11 @@ async def parsed_consumer(message: "NativeMessage[Any]") -> Any:
"""Wrapper, that parser FastStream message to FastAPI compatible one."""
body = await message.decode()

fastapi_body: dict[str, Any] | list[Any]
fastapi_body: dict[str, Any] | list[Any] | None
if first_arg is not None:
if isinstance(body, dict):
if message.no_body:
fastapi_body, path = None, {}
elif isinstance(body, dict):
path = fastapi_body = body or {}
elif isinstance(body, list):
fastapi_body, path = body, {}
Expand Down
3 changes: 3 additions & 0 deletions faststream/_internal/kafka/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,8 +10,11 @@
key_for_index,
realign_keys,
)
from .tombstone import TOMBSTONE, Tombstone

__all__ = (
"TOMBSTONE",
"Tombstone",
"extract_per_message_keys_and_bodies",
"key_for_index",
"realign_keys",
Expand Down
108 changes: 108 additions & 0 deletions faststream/_internal/kafka/tombstone.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,108 @@
"""Tombstones: Kafka records whose value is null.

A tombstone is the delete marker of a compacted topic, so it has to survive
the round trip as something other than an empty payload. `Tombstone` is an
empty `bytes` subclass: it compares equal to `b""` and is `bytes` everywhere
the framework expects bytes, and `isinstance(body, Tombstone)` is what tells
the two apart.

Note that any operation producing a new `bytes` from a tombstone - slicing,
concatenation, `bytes(...)`, `join(...)` - returns a plain `bytes`, so an
`isinstance` check belongs before that kind of work, not after it.
"""

from collections.abc import Callable, Sequence
from typing import TYPE_CHECKING

from typing_extensions import Self

from faststream._internal.parser import BatchCodecProto

if TYPE_CHECKING:
from fast_depends.library.serializer import SerializerProto

from faststream._internal.basic_types import SendableMessage
from faststream._internal.parser import CodecProto


class Tombstone(bytes):
"""An empty `bytes` marking a record published with a null value."""

__slots__ = ()

def __new__(cls, value: bytes = b"") -> Self:
"""Build the marker, rejecting any attempt to give it a payload.

`bytes.__getnewargs__` hands `b""` back on unpickle, hence the
argument.
"""
if value:
msg = f"{cls.__name__} carries no data, got {value!r}"
raise ValueError(msg)
return super().__new__(cls, b"")

def __repr__(self) -> str:
return "TOMBSTONE"

# NOTE: bytes doesn't route str()/f-strings through __repr__, so without
# this a tombstone reads as b"" in every log line that doesn't use !r
__str__ = __repr__


TOMBSTONE: Tombstone = Tombstone()


def value_or_tombstone(value: bytes | None) -> bytes:
"""Map a raw record value to a body, marking a null one."""
return TOMBSTONE if value is None else value


async def encode_or_tombstone(
message: "SendableMessage",
codec: "CodecProto",
serializer: "SerializerProto | None",
*,
key: bytes | str | None = None,
none_is_tombstone: bool = False,
) -> tuple[bytes | None, str | None]:
"""Encode a body, or answer with a null value for a tombstone.

An explicit tombstone requires a key, since compaction deletes per key.
`none_is_tombstone` carries each broker's own legacy rule for a `None`
body and never implies the key requirement.
"""
if isinstance(message, Tombstone):
if key is None:
msg = "a Kafka tombstone requires a key"
raise ValueError(msg)
return None, None

if none_is_tombstone and message is None:
return None, None

return await codec.encode(message, serializer)


async def encode_batch_or_tombstone(
bodies: Sequence["SendableMessage"],
codec: "CodecProto",
serializer: "SerializerProto | None",
key_for: Callable[[int], bytes | str | None],
) -> Sequence[tuple[bytes | None, str | None]]:
"""Encode a batch, keeping each element's own key for its tombstone check.

A custom `BatchCodecProto` encodes the batch as a whole and has no way to
express a null value for one record, so a tombstone cannot ride along.
"""
if isinstance(codec, BatchCodecProto):
if any(isinstance(body, Tombstone) for body in bodies):
msg = "a tombstone in a batch isn't supported with a custom BatchCodecProto"
raise ValueError(msg)
return await codec.encode_batch(bodies, serializer)

# NOTE: no `none_is_tombstone` here - a bare None in a batch encodes
# normally, exactly as it did before tombstones existed
return [
await encode_or_tombstone(body, codec, serializer, key=key_for(position))
for position, body in enumerate(bodies)
]
3 changes: 3 additions & 0 deletions faststream/confluent/__init__.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
from typing import TYPE_CHECKING, TypeAlias

from faststream._internal.kafka import TOMBSTONE, Tombstone
from faststream._internal.parser import ParserProto
from faststream._internal.testing.app import TestApp

Expand All @@ -24,6 +25,7 @@
raise ImportError(INSTALL_FASTSTREAM_CONFLUENT) from e

__all__ = (
"TOMBSTONE",
"ConfluentParserType",
"KafkaBroker",
"KafkaMessage",
Expand All @@ -35,6 +37,7 @@
"KafkaRouter",
"TestApp",
"TestKafkaBroker",
"Tombstone",
"Topic",
"TopicPartition",
)
4 changes: 4 additions & 0 deletions faststream/confluent/message.py
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,10 @@ def __init__(
if not is_manual:
self.committed = AckStatus.ACKED

@property
def tombstone(self) -> bool:
return self.no_body

async def ack(self) -> None:
"""Acknowledge the Kafka message."""
if self.is_manual and not self.committed:
Expand Down
18 changes: 14 additions & 4 deletions faststream/confluent/parser.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
from typing import TYPE_CHECKING, Any, cast

from faststream._internal.kafka.tombstone import Tombstone, value_or_tombstone
from faststream.message import StreamMessage, decode_message

from .message import FAKE_CONSUMER, KafkaMessage
Expand Down Expand Up @@ -38,12 +39,14 @@ async def parse_message(
"""Parses a Kafka message."""
headers = _parse_msg_headers(cast("_HeadersInput", message.headers() or ()))

body = message.value() or b""
value = message.value()
body = value_or_tombstone(value)
offset = message.offset()
_, timestamp = message.timestamp()

return KafkaMessage(
body=body,
no_body=value is None,
headers=headers,
reply_to=headers.get("reply_to", ""),
content_type=headers.get("content-type"),
Expand All @@ -66,7 +69,7 @@ async def parse_batch(
last = message[-1]

for m in message:
body.append(m.value() or b"")
body.append(value_or_tombstone(m.value()))
batch_headers.append(
_parse_msg_headers(cast("_HeadersInput", m.headers() or ()))
)
Expand All @@ -90,17 +93,24 @@ async def parse_batch(

async def decode_message(
self,
msg: "StreamMessage[Message]",
msg: "StreamMessage[Any]",
) -> "DecodedMessage":
"""Decodes a message."""
# NOTE: a tombstone carries no payload, so any content-type on it is a lie
if isinstance(msg.body, Tombstone):
return msg.body

return decode_message(msg)

async def decode_batch(
self,
msg: "StreamMessage[tuple[Message, ...]]",
) -> "DecodedMessage":
"""Decode a batch of messages."""
return [decode_message(await self.parse_message(m)) for m in msg.raw_message]
return [
await self.decode_message(await self.parse_message(m))
for m in msg.raw_message
]


def _parse_msg_headers(headers: "_HeadersInput") -> dict[str, str]:
Expand Down
29 changes: 16 additions & 13 deletions faststream/confluent/publisher/producer.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,11 @@
from typing_extensions import override

from faststream._internal.endpoint.utils import ParserComposition
from faststream._internal.parser import BatchCodecProto, DefaultCodec
from faststream._internal.kafka.tombstone import (
encode_batch_or_tombstone,
encode_or_tombstone,
)
from faststream._internal.parser import DefaultCodec
from faststream._internal.producer import ProducerProto
from faststream.confluent.parser import AsyncConfluentParser
from faststream.confluent.response import KafkaPublishCommand
Expand Down Expand Up @@ -139,10 +143,14 @@ async def publish(
cmd: "KafkaPublishCommand",
) -> "asyncio.Future[Message | None] | Message | None":
"""Publish a message to a topic."""
if cmd.body is None:
message, content_type = None, None
else:
message, content_type = await self.codec.encode(cmd.body, self.serializer)
message, content_type = await encode_or_tombstone(
cmd.body,
self.codec,
self.serializer,
key=cmd.key,
# confluent tombstones any None body, keyed or not, unlike aiokafka
none_is_tombstone=True,
)

headers_to_send = {
"content-type": content_type or "",
Expand All @@ -166,14 +174,9 @@ async def publish_batch(self, cmd: "KafkaPublishCommand") -> None:

headers_to_send = cmd.headers_to_publish()

if isinstance(self.codec, BatchCodecProto):
encoded_batch = await self.codec.encode_batch(
cmd.batch_bodies, self.serializer
)
else:
encoded_batch = [
await self.codec.encode(msg, self.serializer) for msg in cmd.batch_bodies
]
encoded_batch = await encode_batch_or_tombstone(
cmd.batch_bodies, self.codec, self.serializer, cmd.key_for
)

for message_position, (message, content_type) in enumerate(encoded_batch):
if content_type:
Expand Down
31 changes: 17 additions & 14 deletions faststream/confluent/testing.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,11 @@
from typing_extensions import override

from faststream._internal.endpoint.utils import ParserComposition
from faststream._internal.parser import BatchCodecProto, DefaultCodec
from faststream._internal.kafka.tombstone import (
encode_batch_or_tombstone,
encode_or_tombstone,
)
from faststream._internal.parser import DefaultCodec
from faststream._internal.testing.broker import (
EnterType,
TestBroker,
Expand Down Expand Up @@ -196,12 +200,9 @@ async def publish_batch(self, cmd: "KafkaPublishCommand") -> None:
"""Publish a batch of messages to the Kafka broker."""
serializer = self.broker.config.fd_config._serializer

if isinstance(self.codec, BatchCodecProto):
encoded = await self.codec.encode_batch(cmd.batch_bodies, serializer)
else:
encoded = [
await self.codec.encode(body, serializer) for body in cmd.batch_bodies
]
encoded = await encode_batch_or_tombstone(
cmd.batch_bodies, self.codec, serializer, cmd.key_for
)

for handler in _find_handler(
self.subscribers,
Expand Down Expand Up @@ -351,12 +352,14 @@ async def build_message(
id_generator: IdGenerator = gen_cor_id,
) -> 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)
msg, content_type = await encode_or_tombstone(
message,
codec or DefaultCodec(),
serializer,
key=key,
# confluent tombstones any None body, keyed or not, unlike aiokafka
none_is_tombstone=True,
)
k = key or b""
headers = {
"content-type": content_type or "",
Expand All @@ -379,7 +382,7 @@ async def build_message(


def _build_mock_message(
body: bytes,
body: bytes | None,
content_type: str | None,
topic: str,
partition: int | None = None,
Expand Down
Loading
Loading