Repository navigation
Conversation
…ases it (#19) ag2ai/faststream#2933 (Optional[Model] = None resolving correctly for a real tombstone instead of crash-looping on required fields) is open, unmerged, no PyPI release. Producer-side already gets the equivalent treatment via _patch_real_tombstones - do the same here so both halves work today, not just once faststream ships something. Three patches, each guarded for idempotency: - AsyncConfluentParser.parse_message: tags the body with a dedicated TOMBSTONE sentinel (not a byte pattern - can't collide with real content) when the raw value is a genuine null. - AsyncConfluentParser.decode_message: maps TOMBSTONE to None. - faststream._internal.fastapi.route.build_faststream_to_fastapi_parser: a full replacement (it's a per-subscriber closure factory, not a patchable class method) that passes a real None through instead of wrapping it as {param_name: None} - that wrapping is what defeats FastAPI's own no-body shortcut and forces field-level validation on every tombstone. Bumps to 0.4.51. Co-authored-by: Claude Sonnet 5 <noreply@anthropic.com>
…ases it (#19) ag2ai/faststream#2933 (Optional[Model] = None resolving correctly for a real tombstone instead of crash-looping on required fields) is open, unmerged, no PyPI release. Producer-side already gets the equivalent treatment via _patch_real_tombstones - do the same here so both halves work today, not just once faststream ships something. Three patches, each guarded for idempotency: - AsyncConfluentParser.parse_message: tags the body with a dedicated TOMBSTONE sentinel (not a byte pattern - can't collide with real content) when the raw value is a genuine null. - AsyncConfluentParser.decode_message: maps TOMBSTONE to None. - faststream._internal.fastapi.route.build_faststream_to_fastapi_parser: a full replacement (it's a per-subscriber closure factory, not a patchable class method) that passes a real None through instead of wrapping it as {param_name: None} - that wrapping is what defeats FastAPI's own no-body shortcut and forces field-level validation on every tombstone. Bumps to 0.4.51. Co-authored-by: Claude Sonnet 5 <noreply@anthropic.com>
…ases it (#19) (#20) ag2ai/faststream#2933 (Optional[Model] = None resolving correctly for a real tombstone instead of crash-looping on required fields) is open, unmerged, no PyPI release. Producer-side already gets the equivalent treatment via _patch_real_tombstones - do the same here so both halves work today, not just once faststream ships something. Three patches, each guarded for idempotency: - AsyncConfluentParser.parse_message: tags the body with a dedicated TOMBSTONE sentinel (not a byte pattern - can't collide with real content) when the raw value is a genuine null. - AsyncConfluentParser.decode_message: maps TOMBSTONE to None. - faststream._internal.fastapi.route.build_faststream_to_fastapi_parser: a full replacement (it's a per-subscriber closure factory, not a patchable class method) that passes a real None through instead of wrapping it as {param_name: None} - that wrapping is what defeats FastAPI's own no-body shortcut and forces field-level validation on every tombstone. Bumps to 0.4.51. Co-authored-by: Claude Sonnet 5 <noreply@anthropic.com>
Lancetnik
left a comment
There was a problem hiding this comment.
Thanks for the follow-up — the problem is real (consuming a compacted topic through the FastAPI integration crash-loops on every tombstone today), and the FastAPI-side symptom fix works. But the chosen mechanism — a sentinel object traveling through the public message.body — moves the crash instead of removing it, and changes behavior outside the stated scope. Requesting changes on the points below (1–3 are blocking).
1. The sentinel crashes Prometheus / OpenTelemetry middlewares on every tombstone
ConfluentMetricsSettingsProvider.get_consume_attrs_from_message does len(msg.body) (faststream/confluent/prometheus/provider.py:36), and the OTel provider does the same (faststream/confluent/opentelemetry/provider.py:56). Both run in the consume path before the user handler. _Tombstone has no __len__, so any instrumented app gets TypeError: object of type '_Tombstone' has no len() on every tombstone — the exact bug class this PR fixes for FastAPI, reintroduced one layer up. The batch variants are worse: bytearray().join(msg.body) (prometheus :51, otel :87-88) fails the whole batch if it contains a single tombstone. Verified by running the PR code end-to-end.
Whatever representation is chosen, both providers (single + batch) need to be updated in this PR.
2. elif body is None in route.py silently changes behavior for all brokers, not just Kafka tombstones
parsed_consumer is shared cross-broker code. Any message whose decoded body is None — e.g. a plain b"null" JSON payload on RabbitMQ/NATS — now takes the new branch. Reproduced on in-memory Rabbit: a handler async def h(msg: dict) was previously invoked (receiving the leaked {"msg": None} wrapper), and now raises RequestValidationError: Field required before the handler runs. The old behavior was arguably broken too, but this is a user-visible semantic change on every broker and should be deliberate: documented, and covered by a base testcase in tests/brokers/base/fastapi.py so all brokers inherit it — not only by a confluent-local test.
3. The sibling aiokafka broker keeps the bug — and the sentinel is likely unnecessary
faststream/kafka/parser.py:45 / :77 still do body = message.value or b"", so KafkaBroker users get none of this fix while the special cases live in shared layers (decode_message, fastapi route). Both Kafka parsers should implement the same contract.
Related design point: StreamMessage.body is already annotated bytes | Any, so plain None expresses "tombstone" with zero new vocabulary — the parser sets body = message.value(), and decode_message needs the same one-line guard either way. That would (a) avoid a new forever-public TOMBSTONE symbol, (b) make finding #1 visible to mypy as len(Optional[bytes]), and (c) drop the __bool__ = False footgun, which makes if not msg.body: conflate tombstones with b"" again — the very distinction the sentinel exists to preserve. If the sentinel stays, note there's already an established sentinel idiom in the codebase (EMPTY in faststream/_internal/constants.py) it should follow, and the mapping expression is duplicated in parse_message/parse_batch (a third copy will appear with the aiokafka fix) — worth one helper.
4. In-memory TestKafkaBroker now diverges from the real broker
faststream/confluent/testing.py:349 encodes publish(None) through encode_message (→ b"") and MockConfluentMessage.value() can't represent a null value at all. So the same user code sees b"" under TestKafkaBroker and None against a real broker — user test suites will pass while production behaves differently. The testing path needs the same null-value story.
5. Test nits
await asyncio.sleep(self.timeout)burns a fixed 10s per connected run and flakes if delivery is slower; the neighboring tests (test_batch_real, baseFastAPITestcase) use anasyncio.Event+asyncio.wait(..., timeout=...)— same assertion, ~1s, no flake window.br._producer._producer.producerreaches through three levels of private internals; fine as a deliberate encoder bypass, but the exact-order list assert additionally relies on the topic having a single partition — anEventset atlen(received) == 2plus order-insensitive comparison (like the batch-key tests do) would be sturdier.TOMBSTONEis a new public export with no docs; if it survives the redesign discussion it needs a documentation entry.
Also: this PR and #2932 are two halves of one story (produce vs consume) — worth coordinating the merge so tombstones work end-to-end.
🤖 Generated with Claude Code
|
Thanks, fixed all three blocking points, keeping the sentinel.
Also fixed the sleep-based test, uses an event now. On dropping the sentinel for plain None: I'd rather not. Your own point 2 is the reason why, None already means other things on other brokers, and your rabbit repro showed that collision is real. Folding tombstone into plain None would make that worse, not better, since a handler treating None as "optional field" would now silently also fire on "this was deleted" with no error at all. The sentinel is what lets us keep the two apart. Separate question: does #2939 (TestKafkaBroker mocking) need to fold into this PR, or does it stay independent? Happy either way. |
|
Added the TOMBSTONE note to kafka/message.md and confluent/message.md, next to the value() field list. |
f2fb723 to
27faf86
Compare
|
Pushed one more change. publish(None, ...) used to always send a real tombstone, but now None goes through the codec like any other value, so it's symmetric with what consume already does. TOMBSTONE is the explicit way to send a real tombstone now. Two things I wasn't sure about here, and #2939 merged in the meantime assuming the old None = tombstone behavior, so you might not have caught this redesign yet. Happy to go either way:
|
|
Ran a hard self review on the publish side change before you look again, found and fixed a real bug plus a few smaller things. Real bug: Also cleaned up:
All green, mypy and ruff clean. |
|
hey @Lancetnik, would appreciate a re-review! :) |
2d853e1 to
0fcd275
Compare
0fcd275 to
34b34ee
Compare
d908b7a to
19432c6
Compare
Lancetnik
left a comment
There was a problem hiding this comment.
Thanks for the thorough follow-up — points 1–3 from the last review are properly closed: both providers (single + batch, kafka + confluent) no longer assume bytes, the route.py check is scoped to a real tombstone instead of body is None, and aiokafka now has parity with confluent through one shared helper. CI is green. None of what follows is a complaint about the execution.
But I think the mechanism is wrong, and that's what's driving the size and the breakage.
The root cause
A tombstone is a property of the message — "this record has no value". This PR encodes it as a value in the body slot. That slot is read as bytes everywhere in the framework: len() in four telemetry providers, json_loads() in decode_message, bytearray().join() in the batch paths, record sizing in both testing brokers. So one sentinel in that slot leaks into 26 files, needs five new forever-public symbols, and breaks behavior on both sides — msg: bytes handlers on consume, and publish(None, key=...) on produce.
What the issue actually needs
#1967 is a FastAPI crash-loop on a compacted topic. For Optional[Model] = None to resolve to None, exactly one thing has to happen: parsed_consumer needs to know the message has no body, so it can hand body=None to solve_dependencies. That information is already available losslessly — raw_message.value is None. Changing msg.body isn't required for it.
If the tombstone is a flag on StreamMessage and body stays b"":
- the kafka/confluent parsers change by a couple of lines each, additively;
route.pychecks the flag instead of the sentinel — two lines;- the Prometheus/OTel providers,
decode_message,encode_message, both testing brokers and both producers are not touched at all; - existing
msg: byteshandlers keep receivingb""— nothing breaks; - the public surface is one attribute instead of
TOMBSTONE/Tombstone/body_size/batch_body_size/value_or_tombstone.
That's ~5 files instead of 26, with no breaking change — and the tombstone/b"" distinction isn't lost, it becomes available for the first time. Today there's no way to tell them apart at all.
How I'd like to split this
1. The crash fix — mergeable now. Message-level flag plus the route.py branch. Closes #1967, breaks nothing, introduces no new vocabulary. If you cut this out of the current branch as its own PR, I'll merge it.
2. An explicit TOMBSTONE on publish — additive only. publish(None, key=...) shipped in 0.7.2 (#2932) and #2939 followed in 0.7.3, so it can't be revoked the way this branch does: a user who deletes keys that way gets no exception and no warning after upgrading, just b"" on the wire and compaction that silently stops deleting. If we want an explicit spelling, it has to be added alongside, with the None form deprecated by warning and removed in 0.8. And in that case Tombstone has to be part of SendableMessage — right now the documented call doesn't type-check:
error: No overload variant of "publish" of "KafkaBroker" matches argument types "Tombstone", "str", "bytes"
error: No overload variant of "publish_batch" of "KafkaBroker" matches argument types "Tombstone", "str", "bytes"
3. msg.body becoming None/a sentinel for a tombstone — that's 0.8. This part is breaking by definition: anything other than b"" changes behavior for every existing Kafka consumer (your own tests had to move from msg: bytes to msg: bytes | None). There's no breaking-change window open on 0.7.x, so this waits for the major — which is also where the full design belongs, including what a tombstone should look like in body and whether the sentinel is public.
Smaller things, if the current shape survives
Tombstoneis public and constructible, but the guards disagree: the key requirement checkscmd.body is TOMBSTONEwhileencode_or_tombstonechecksisinstance.publish(Tombstone(), topic)without a key skips the friendly error and reaches aiokafka askey=None, value=None.__bool__returnsFalsewith no__eq__, soTOMBSTONE == Tombstone()isFalsewhileisinstancepasses — andif not msg.body:conflates a tombstone withb""again, the exact distinction the sentinel exists to preserve. (EmptyPlaceholderin_internal/constants.pydefines both.)- Requiring a key on confluent is a second, independent breaking change riding along — a null-key/null-value record is legal Kafka and
confluent-kafkaaccepts it. - The batch + custom
BatchCodecProtocombination raises at runtime with nothing about it in the docs.
To be clear about what I'm asking for: not a rewrite of what's here. Just pull piece 1 out into its own PR so the crash is fixed for users now, and let 2 and 3 become a design discussion for 0.8, where TOMBSTONE can land whole instead of in halves.
`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
19432c6 to
80810e8
Compare
|
Rebased onto current The crash fix is out in #3028 and this branch is stacked on it, so once #3028 lands the diff here collapses to the second commit. What's left is only what you scoped to 0.8, plus the smaller items from your review:
On revoking |
80810e8 to
c4c47b4
Compare
|
Correction to my last comment — I said "the smaller items from your review" but only covered two of the four. The other two, straight: Key requirement (your nit 3). Still there, deliberately. Compaction deletes per key, so a keyless tombstone deletes nothing — it's legal Kafka but almost certainly a bug at the call site, and the error is cheaper than the silent no-op. You're right that it's an independent breaking change riding along, though, so it's your call: say the word and I'll drop it and delete Batch + custom |
c4c47b4 to
f6e4c30
Compare
|
Reworked this so it isn't a breaking change anymore.
That drops most of what you flagged. Also folded the duplicated batch encode into one helper next to Since it's additive now, does this still need to wait for 0.8 or can it land on 0.7.x? Would appreciate another look. |
|
would appreciate a re-review @Lancetnik :) |
f6e4c30 to
56ee916
Compare
|
Hi! I’ll join the review this week |
|
#3063 adds
1. Move
2. Move They're wiring for the Kafka parsers and producers. The per-message key helpers already in 3. Drop the #3028 is landing a neutral 4. Replace the batch flag with per-element detection.
5. Move the That's the per-broker decode hook, so 6. Drop the Keep the two stacked; there's nothing left to duplicate here once #3028 is in. Separately, and not yours to fix: 7. Drop the " That leaves this PR as the |
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>
56ee916 to
8f90891
Compare
…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>
8f90891 to
b7f61d4
Compare
…ords A null record value is a delete marker on a compacted topic, and until now nothing could tell one from an empty payload once it was parsed. `Tombstone` is an empty `bytes` subclass, so `msg.body` for a null record still equals `b""`, is falsy, has `len() == 0` and is still `bytes` - every existing consumer behaves exactly as before - while `isinstance(body, Tombstone)` answers the question for the first time. A batch marks its tombstones one element at a time, which a single flag cannot express. On publish, `TOMBSTONE` is the explicit spelling and requires a key, since compaction deletes per key. `publish(None)` is untouched: keyed-only on aiokafka, any `None` on confluent, both byte identical to before. Everything lives in `_internal/kafka/` plus the two Kafka packages, which re-export `Tombstone`/`TOMBSTONE`; the encode helpers stay unexported. `decode_message` for a tombstone is handled by each broker's own decode hook rather than `faststream/message/utils.py`, so `_internal/`, `faststream/message/` and the telemetry providers are untouched. The tombstone tests the two brokers share live in `tests/brokers/base/` behind a `get_message_value` hook, the way `BatchKeysTestcase` already does it for keys, and the telemetry providers get a tombstone through them so the `len()` crash that shaped this design stays impossible by test rather than by argument. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
b7f61d4 to
06088c6
Compare
|
@Lancetnik @borisalekseev ready for another look — rewritten to your 7-item list, and rebased on #3028 (which now carries the All seven are done; the description above walks through them one by one. The short version: Your two open nits. The key requirement now only applies to an explicit Two things worth your judgement:
Test-wise, the cases both brokers share now live in CI green, including both real-broker suites. ruff clean, mypy clean across 519 source files. |
Stacked on #3028 — once that merges the diff here collapses to the last commit. Rewritten to your 7-item list: everything now lives in
_internal/kafka/plus the two Kafka packages, and nothing under_internal/,faststream/message/or_internal/fastapi/is touched. 17 files instead of 25.No breaking changes.
Tombstonesubclassesbytesand is empty, somsg.bodyfor a null record is== b"",len() == 0, falsy, and still an instance ofbytes. Every existing consumer path is unchanged;isinstance(body, Tombstone)is the only thing that tells a tombstone apart from an empty payload, which was previously impossible.Your list
Tombstone/TOMBSTONEmoved to_internal/kafka/tombstone.py, re-exported fromfaststream.kafkaandfaststream.confluent._internal/constants.pyis byte identical tomain, so a Redis or NATS user never sees the symbol.value_or_tombstone,encode_or_tombstone,encode_batch_or_tombstonemoved there too, and out of every__all__— the parsers and producers import them by path.tombstoneflag is gone from this PR. It rides fix(fastapi): resolve Optional body params to None for a Kafka tombstone #3028'sno_bodyplusKafkaMessage.tombstone, so there's one source of truth.msg.bodyis aTombstoneor isn't, which is the one thing a single flag genuinely cannot express.decode_messagespecial case moved intoAioKafkaParser.decode_messageand its Confluent counterpart.faststream/message/utils.pyis untouched.decode_batchnow routes each record through that hook instead of calling the module-level function, so a batch inherits it rather than needing its own copy._internal/fastapi/route.pychange here — fix(fastapi): resolve Optional body params to None for a Kafka tombstone #3028 delivers it.publish(None)deprecation. Both legacy rules are byte identical tomain: keyed-only on aiokafka, anyNoneon confluent. 1.0: encode a None body as JSON null instead of an empty body #3062 keeps the question.The two nits from the earlier review
Key requirement. It now applies only to an explicit
Tombstone, which is new API, so nothing is taken away: a keylesspublish(None)on confluent still produces a tombstone exactly as onmain. The independent breaking change you flagged isn't riding along anymore.Batch + custom
BatchCodecProtostill raises —encode_batch()has no way to express a null value for one record of a batch, so there's nothing correct to encode. Documented in bothmessage.mdpages rather than left as a runtime surprise.Untouched
decode_message/encode_message,SendableMessage,StreamMessage, and the four Prometheus/OpenTelemetry providers are byte identical tomain. A tombstone body is genuinebytes, so none of them need to change, and thelen()crash from the first review cannot recur.Tombstoneneeds noSendableMessageentry either, sincebytesis already in the union.Behavior worth calling out
type(msg.body) is bytesis nowFalsefor a tombstone. Exact-type checks are rare, but it is true.msg: bytes | Nonehandler receives ab""-equal value, notNone. The tombstone signal ismsg.tombstoneorisinstance. This matchesmain, so nothing regresses.Test plan
tests/brokers/base/testclient.py+tests/brokers/base/publish.py: the tombstone cases both brokers share live in one place behind aget_message_valuehook, the wayBatchKeysTestcasealready does it for keys — in-memory (body reads asb"", explicitTOMBSTONEon the wire,request(), the missing-key raise, a tombstone in a batch, a plainNonein a batch staying non-null, theBatchCodecProtoraise, and a custom codec never seeing the tombstone) and connected (TOMBSTONE, the missing-key raise, a batch tombstone).test_publish_none_tombstoneandtest_consume_*_without_valuepass byte-for-byte unmodified — that's the back-compat proof, not a claim.tests/message/test_tombstone.py: the class's own contract —Tombstone(b"data")raises, pickle/copy/deepcopy round-trip, and it prints asTOMBSTONEunderrepr,str, an f-string and%s(plainbytesdoesn't route those through__repr__, so a tombstone would otherwise read asb""in every log line).tests/{prometheus,opentelemetry}/{kafka,confluent}/test_provider.py: a tombstone through all four providers, single and batch. Thelen()crash that shaped this design is now impossible by test rather than by argument.tests/brokers/{kafka,confluent}/test_parser.py: the sentinel for a single record, per-element marking in a batch with the message-level flag pinned toFalseby design,decode_batchon a mixed batch, andapplication/jsonon a tombstone decoding to empty while a genuineb""with the same content type still raises.tests/brokers/confluent/test_test_client.py: confluent's own divergence — a keylesspublish(None)still tombstones there, unlike aiokafka.tests/brokers+tests/message+tests/prometheus+tests/opentelemetry+tests/asyncapi(not connected) pass. ruff clean, mypy clean across 519 source files.docs/docs/en/{kafka,confluent}/message.md.🤖 Generated with Claude Code