Repository navigation
Conversation
…ords Builds on the message-level flag from ag2ai#3028. That flag records the fact ("this record has no value") but leaves `body` as `b""`, so a tombstone still can't be told from an empty payload anywhere except the FastAPI route, and there's no way to publish one explicitly. This adds the representation: - `Tombstone` / `TOMBSTONE` live in `_internal/constants.py` next to `EmptyPlaceholder`, with the same `__repr__` / `__bool__` / `__eq__` shape, and `Tombstone` is part of `SendableMessage` so `publish(TOMBSTONE, topic, key=...)` type-checks. - Both Kafka parsers map a null record value to `TOMBSTONE`; `decode_message` maps it back to `None`. `b""` and `b"null"` stay distinct from it. - Both producers and both testing brokers publish a real null value for `TOMBSTONE`, single and batch, and require a key for it. - The Prometheus and OpenTelemetry providers size bodies through `body_size` / `batch_body_size` instead of assuming `bytes`, so a tombstone no longer raises `TypeError` mid-batch. Every guard is `isinstance(..., Tombstone)`, so a user-constructed `Tombstone()` takes the same path as the singleton. Breaking: `msg.body` is `TOMBSTONE` rather than `b""` for a null record, and `publish(None)` encodes normally instead of producing a tombstone. Both are 0.8 changes. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01JoFxdJYipYjPtMdmJaxCFN
…ords Builds on the message-level flag from ag2ai#3028. That flag records the fact ("this record has no value") but leaves `body` as `b""`, so a tombstone still can't be told from an empty payload anywhere except the FastAPI route, and there's no way to publish one explicitly. This adds the representation: - `Tombstone` / `TOMBSTONE` live in `_internal/constants.py` next to `EmptyPlaceholder`, with the same `__repr__` / `__bool__` / `__eq__` shape, and `Tombstone` is part of `SendableMessage` so `publish(TOMBSTONE, topic, key=...)` type-checks. - Both Kafka parsers map a null record value to `TOMBSTONE`; `decode_message` maps it back to `None`. `b""` and `b"null"` stay distinct from it. - Both producers and both testing brokers publish a real null value for `TOMBSTONE`, single and batch, and require a key for it. - The Prometheus and OpenTelemetry providers size bodies through `body_size` / `batch_body_size` instead of assuming `bytes`, so a tombstone no longer raises `TypeError` mid-batch. Every guard is `isinstance(..., Tombstone)`, so a user-constructed `Tombstone()` takes the same path as the singleton. Breaking: `msg.body` is `TOMBSTONE` rather than `b""` for a null record, and `publish(None)` encodes normally instead of producing a tombstone. Both are 0.8 changes. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01JoFxdJYipYjPtMdmJaxCFN
…ords Builds on the message-level flag from ag2ai#3028. That flag records the fact ("this record has no value") but leaves `body` as `b""`, so a tombstone still can't be told from an empty payload anywhere else, and there's no way to publish one explicitly. This adds the representation without changing behavior for anyone: - `Tombstone` subclasses `bytes` and is empty, so `msg.body` for a null record is `== b""`, has `len() == 0`, is falsy, joins into a batch and is an instance of `bytes`. Every existing consumer path is unchanged; `isinstance(body, Tombstone)` is the only way to tell the two apart. - Both Kafka parsers map a null record value to `TOMBSTONE`. `decode_message` is byte-identical to main, so a `msg: bytes` handler still gets an empty body rather than `None`. - Both producers and both testing brokers publish a real null value for `TOMBSTONE`, single and batch, and an explicit `TOMBSTONE` requires a key. - `publish(None)` keeps producing a tombstone exactly as it does today - keyed-only on aiokafka, any `None` on confluent - now with a `DeprecationWarning`; it will encode normally in 0.8. Because a tombstone body is genuine `bytes`, the Prometheus and OpenTelemetry providers, `decode_message`, `encode_message` and `SendableMessage` need no changes at all - those files are identical to main. The public surface is `TOMBSTONE` and `Tombstone`, nothing else. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01JoFxdJYipYjPtMdmJaxCFN
…ords Builds on the message-level flag from ag2ai#3028. That flag records the fact ("this record has no value") but leaves `body` as `b""`, so a tombstone still can't be told from an empty payload anywhere else, and there's no way to publish one explicitly. This adds the representation without changing behavior for anyone: - `Tombstone` subclasses `bytes` and is empty, so `msg.body` for a null record is `== b""`, has `len() == 0`, is falsy, joins into a batch and is an instance of `bytes`. Every existing consumer path is unchanged; `isinstance(body, Tombstone)` is the only way to tell the two apart. - Both Kafka parsers map a null record value to `TOMBSTONE`. `decode_message` is byte-identical to main, so a `msg: bytes` handler still gets an empty body rather than `None`. - Both producers and both testing brokers publish a real null value for `TOMBSTONE`, single and batch, and an explicit `TOMBSTONE` requires a key. - `publish(None)` keeps producing a tombstone exactly as it does today - keyed-only on aiokafka, any `None` on confluent - now with a `DeprecationWarning`; it will encode normally in 0.8. Because a tombstone body is genuine `bytes`, the Prometheus and OpenTelemetry providers, `decode_message`, `encode_message` and `SendableMessage` need no changes at all - those files are identical to main. The public surface is `TOMBSTONE` and `Tombstone`, nothing else. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01JoFxdJYipYjPtMdmJaxCFN
Lancetnik
left a comment
There was a problem hiding this comment.
Two changes and this is good to merge.
1. Rename the flag to no_body — on StreamMessage, both parsers, and route.py.
Core shouldn't carry a Kafka noun; that's the same line we're drawing on #2933. And "tombstone" isn't the concept core needs here — what route.py is deciding is whether the message has a body at all, which is exactly how FastAPI already thinks about a bodyless request. A neutral name says that, and the other brokers get a flag that could mean something to them rather than one that is permanently False and named after someone else's feature.
I checked the rename passes the FastAPI tombstone tests unchanged, so it's mechanical.
2. Expose the domain word one layer down: KafkaMessage.tombstone as a property returning self.no_body, in faststream/kafka/message.py and faststream/confluent/message.py.
Someone consuming a compacted topic should get to write msg.tombstone — that's the right word for them, and it keeps the docs note you added accurate as-is. It just belongs in the Kafka packages rather than the base class. It also gives #2933 something to build on instead of introducing a second flag for the same fact.
Nothing else needed. body staying b"" is the right call for 0.7 — the None-body question is being decided in #3062 for 1.0 — and batches stay out of scope here; we'll settle those in #2933.
|
Both done, rebased on current 1. The flag is 2. The parser tests now assert both the neutral flag and the Kafka-side property, so the property is covered rather than assumed. Reverting only the Batches and the |
|
@Lancetnik ready for another look. Both of your changes are in, plus two gaps I found reviewing my own diff afterwards:
One thing I did not change: CI green, ruff and mypy clean. |
`parsed_consumer()` always wraps the decoded body as `{first_arg: body}`,
so `solve_dependencies()` never sees an absent body - only a present dict
like `{"msg": None}`. It can't take the no-body shortcut it already uses
for bodyless HTTP requests, so a required field on the target model fails
validation before the handler runs. On a compacted topic that's every
tombstone, forever.
A tombstone is a property of the record ("this record has no value"), not
a value in the body slot, so it's carried as a flag on `StreamMessage`:
the kafka/confluent parsers set `tombstone=value is None`, `route.py`
passes `body=None` to `solve_dependencies` when it's set. `body` stays
`b""`, so `decode_message`, the telemetry providers, the testing brokers
and the producers are untouched and `msg: bytes` handlers are unaffected.
Other brokers never set the flag, so a plain `b"null"` payload on
Rabbit/NATS keeps its existing behavior.
Closes ag2ai#1967
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JoFxdJYipYjPtMdmJaxCFN
Review asked for the core flag to stop carrying a Kafka noun. What `route.py` decides is whether the message has a body at all - the same question FastAPI already answers for a bodyless HTTP request - so `StreamMessage` takes `no_body` and the other brokers get a flag that could mean something to them. The domain word lives one layer down: `KafkaMessage.tombstone` is a property over `no_body` in both `faststream/kafka/message.py` and `faststream/confluent/message.py`, so someone consuming a compacted topic still writes `msg.tombstone` and the docs note stays accurate. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
…atch does Confluent already had a connected test driving `value=None` straight through the raw producer, so the "a tombstone from any client, not only `publish(None)`" claim was only proven for one of the two brokers. aiokafka now has the mirror. The docs note also promised more than the code delivers: `parse_batch` never sets the flag, so `msg.tombstone` is always `False` in a `batch=True` subscriber. Said so, rather than leaving a compacted-topic consumer to find out. Per-record tombstones in a batch stay out of scope here. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
9245aa0 to
d704b67
Compare
Piece 1 of #2933, split out per that review — the crash fix on its own, no new vocabulary, no breaking change.
The bug
parsed_consumer()always wraps the decoded body as{first_arg: body}, sosolve_dependencies()never sees a genuinely absent body — only a present dict like{"msg": None}. It has no way to take the no-body shortcut it already uses correctly for bodyless HTTP requests, so a required field on the target model fails validation before the handler runs. On a compacted topic that's every tombstone, forever.The fix
A tombstone is a property of the record — "this record has no value" — not a value in the body slot. So it's carried as a flag on
StreamMessagerather than encoded intobody:StreamMessage.__init__takesno_body: bool = False(additive, defaults preserve today's behavior) — a neutral name, since whatroute.pydecides is whether the message has a body at all, exactly how FastAPI already thinks about a bodyless request;AioKafkaParser.parse_message/AsyncConfluentParser.parse_messagesetno_body=value is None;route.pyhandsbody=Nonetosolve_dependencieswhen the flag is set;KafkaMessage.tombstoneis a property overno_body, in bothfaststream/kafka/message.pyandfaststream/confluent/message.py, so a compacted-topic consumer writesmsg.tombstone.bodystaysb"". That meansdecode_message,encode_message, the four Prometheus/OpenTelemetry providers, both testing brokers and both producers are not touched, and existingmsg: byteshandlers keep receivingb""— nothing breaks. The tombstone/b""distinction isn't lost either; it becomes available for the first time, asmsg.tombstone.Other brokers never set the flag, so a plain
b"null"payload on RabbitMQ/NATS keeps its existing behavior — this also closes point 2 of the first review without a cross-broker semantic change.Scope
Deliberately out of scope here:
parse_batchdoesn't set the flag — Bug: Broker subscribe handler falls in Error when try to handle message without value (but key exists!!) #1967 is a single-message crash, and a per-record flag doesn't fit a batch message's singlebodylist. Worth designing properly rather than guessing.TOMBSTONEspelling on publish, andmsg.bodybecomingNone/a sentinel for a tombstone — both deferred to the 0.8 discussion, per the review.publish(None, key=...)from fix(confluent): publish(None) should send a real Kafka tombstone #2932/fix(testing): TestKafkaBroker should mock a real tombstone too #2939 is untouched and stays the way to produce one.Test plan
tests/brokers/{kafka,confluent}/test_parser.py:no_bodyandmsg.tombstoneare both set for a null value and stayFalseforb""and a real payload;bodyisb""either way.tests/brokers/base/fastapi.py:KafkaTombstoneFastAPILocalTestcase, mixed into both kafka and confluent local routers —Optional[Model] = Noneresolves toNonefor a tombstone, while a genuinely emptyb""body still raisesRequestValidationError.tests/brokers/{confluent,kafka}/test_fastapi.py: connected tests against a real broker, driving the raw producer directly (value=None) so they cover a tombstone from any client, not only FastStream'spublish(None)— both brokers, since the claim shouldn't rest on one.asyncio.Event+wait_for, no fixed sleep.route.pyhunk fails both FastAPI tombstone tests — the fix is load-bearing, not incidentally passing.tests/brokers+tests/message+tests/asyncapi(not connected) pass.docs/docs/en/{kafka,confluent}/message.md, including what the code does not do:parse_batchnever sets the flag, somsg.tombstonestaysFalsein abatch=Truesubscriber.Closes #1967
🤖 Generated with Claude Code
https://claude.ai/code/session_01JoFxdJYipYjPtMdmJaxCFN