Skip to content

feat(kafka): explicit TOMBSTONE for publishing and consuming null records - #2933

Open
aradng wants to merge 6 commits into
ag2ai:mainfrom
aradng:fix/fastapi-body-none-for-empty-message
Open

aradng wants to merge 6 commits into
ag2ai:mainfrom
aradng:fix/fastapi-body-none-for-empty-message

Conversation

@aradng

@aradng aradng commented Jul 13, 2026 •

Copy link
Copy Markdown
Contributor

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. Tombstone subclasses bytes and is empty, so msg.body for a null record is == b"", len() == 0, falsy, and still an instance of bytes. 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

  1. Tombstone / TOMBSTONE moved to _internal/kafka/tombstone.py, re-exported from faststream.kafka and faststream.confluent. _internal/constants.py is byte identical to main, so a Redis or NATS user never sees the symbol.
  2. value_or_tombstone, encode_or_tombstone, encode_batch_or_tombstone moved there too, and out of every __all__ — the parsers and producers import them by path.
  3. The tombstone flag is gone from this PR. It rides fix(fastapi): resolve Optional body params to None for a Kafka tombstone #3028's no_body plus KafkaMessage.tombstone, so there's one source of truth.
  4. The batch flag is replaced by per-element detection. Each element of msg.body is a Tombstone or isn't, which is the one thing a single flag genuinely cannot express.
  5. The decode_message special case moved into AioKafkaParser.decode_message and its Confluent counterpart. faststream/message/utils.py is untouched. decode_batch now routes each record through that hook instead of calling the module-level function, so a batch inherits it rather than needing its own copy.
  6. No _internal/fastapi/route.py change here — fix(fastapi): resolve Optional body params to None for a Kafka tombstone #3028 delivers it.
  7. No publish(None) deprecation. Both legacy rules are byte identical to main: keyed-only on aiokafka, any None on 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 keyless publish(None) on confluent still produces a tombstone exactly as on main. The independent breaking change you flagged isn't riding along anymore.

Batch + custom BatchCodecProto still 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 both message.md pages rather than left as a runtime surprise.

Untouched

decode_message/encode_message, SendableMessage, StreamMessage, and the four Prometheus/OpenTelemetry providers are byte identical to main. A tombstone body is genuine bytes, so none of them need to change, and the len() crash from the first review cannot recur. Tombstone needs no SendableMessage entry either, since bytes is already in the union.

Behavior worth calling out

  • type(msg.body) is bytes is now False for a tombstone. Exact-type checks are rare, but it is true.
  • A msg: bytes | None handler receives a b""-equal value, not None. The tombstone signal is msg.tombstone or isinstance. This matches main, 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 a get_message_value hook, the way BatchKeysTestcase already does it for keys — in-memory (body reads as b"", explicit TOMBSTONE on the wire, request(), the missing-key raise, a tombstone in a batch, a plain None in a batch staying non-null, the BatchCodecProto raise, and a custom codec never seeing the tombstone) and connected (TOMBSTONE, the missing-key raise, a batch tombstone).
  • The pre-existing test_publish_none_tombstone and test_consume_*_without_value pass 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 as TOMBSTONE under repr, str, an f-string and %s (plain bytes doesn't route those through __repr__, so a tombstone would otherwise read as b"" in every log line).
  • tests/{prometheus,opentelemetry}/{kafka,confluent}/test_provider.py: a tombstone through all four providers, single and batch. The len() 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 to False by design, decode_batch on a mixed batch, and application/json on a tombstone decoding to empty while a genuine b"" with the same content type still raises.
  • tests/brokers/confluent/test_test_client.py: confluent's own divergence — a keyless publish(None) still tombstones there, unlike aiokafka.
  • Full tests/brokers + tests/message + tests/prometheus + tests/opentelemetry + tests/asyncapi (not connected) pass. ruff clean, mypy clean across 519 source files.
  • Docs updated in docs/docs/en/{kafka,confluent}/message.md.

🤖 Generated with Claude Code

@aradng
aradng requested a review from Lancetnik as a code owner July 13, 2026 15:56
@github-actions github-actions Bot added the Confluent Issues related to `faststream.confluent` module label Jul 13, 2026
aradng added a commit to aradng/FastLoom that referenced this pull request Jul 13, 2026
…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>
aradng added a commit to aradng/FastLoom that referenced this pull request Jul 13, 2026
…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>
aradng added a commit to aradng/FastLoom that referenced this pull request Jul 13, 2026
…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 Lancetnik left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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, base FastAPITestcase) use an asyncio.Event + asyncio.wait(..., timeout=...) — same assertion, ~1s, no flake window.
  • br._producer._producer.producer reaches 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 — an Event set at len(received) == 2 plus order-insensitive comparison (like the batch-key tests do) would be sturdier.
  • TOMBSTONE is 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

@aradng

aradng commented Jul 14, 2026

Copy link
Copy Markdown
Contributor Author

Related: #2932 (merged, producer side) and #2939 (test broker side, same gap). This PR covers the consumer/fastapi side. All three together make Model | Tombstone actually usable end to end.

@github-actions github-actions Bot added the AioKafka Issues related to `faststream.kafka` module label Jul 14, 2026
@aradng

aradng commented Jul 14, 2026

Copy link
Copy Markdown
Contributor Author

Thanks, fixed all three blocking points, keeping the sentinel.

  1. Prometheus and otel now check for the sentinel instead of assuming bytes, single and batch, confluent and kafka. Added body_size/batch_body_size helpers so this isn't duplicated per provider.

  2. Scoped the route.py check to message.body is TOMBSTONE instead of body is None. Reproduced your rabbit repro before and after: async def h(msg: dict) receiving a null body still gets {'msg': None} exactly like before, unaffected now. Only a real tombstone takes the new path.

  3. aiokafka parser now does the same TOMBSTONE mapping as confluent (pulled into a shared value_or_tombstone helper so it's one place, not three). Two existing kafka tests asserted the old b""-on-tombstone behavior, updated them to expect None.

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.

@aradng
aradng requested a review from Lancetnik July 14, 2026 16:48
@aradng

aradng commented Jul 14, 2026

Copy link
Copy Markdown
Contributor Author

Added the TOMBSTONE note to kafka/message.md and confluent/message.md, next to the value() field list.

@github-actions github-actions Bot added the documentation Improvements or additions to documentation label Jul 14, 2026
@aradng
aradng force-pushed the fix/fastapi-body-none-for-empty-message branch from f2fb723 to 27faf86 Compare July 14, 2026 17:41
@aradng

aradng commented Jul 14, 2026 •

Copy link
Copy Markdown
Contributor Author

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:

  1. Is this publish side change even in scope for this PR? It's breaking since it flips the existing tombstone on None behavior, I can split it out if you'd rather keep this PR scoped to consume/FastAPI.
  2. Small thing, should plain None keep encoding to b"" like it does now, or match how everything else gets encoded and produce b"null"?

@aradng

aradng commented Jul 14, 2026

Copy link
Copy Markdown
Contributor Author

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: publish_batch never checked for TOMBSTONE, only single message publish() did. A tombstone inside a batch fell through to the codec and crashed with a serialization error instead of producing a real tombstone. Fixed in both confluent and kafka, real producer and the test client.

Also cleaned up:

  • extracted encode_or_tombstone() and ensure_tombstone_key() into message/utils.py so the tombstone check and the key requirement aren't copied across four places
  • confluent now requires a key for a tombstone too, same as aiokafka. A keyless tombstone deletes nothing on either broker so there was no good reason for the difference
  • switched the key check from SetupError to ValueError. Every other SetupError in this repo is a router or broker construction time error, not a per call publish argument problem
  • renamed _Tombstone to Tombstone and exported it, since the test client needs it for an isinstance check at runtime, not just for typing
  • moved the FastAPI tombstone test into a shared testcase so confluent and kafka both run it, added a negative case for a genuinely empty non-null body
  • added tests for a tombstone inside a batch, and gave confluent its own consume tests for a plain tombstone, it had none before

All green, mypy and ruff clean.

@aradng

aradng commented Jul 18, 2026

Copy link
Copy Markdown
Contributor Author

hey @Lancetnik, would appreciate a re-review! :)

@aradng
aradng force-pushed the fix/fastapi-body-none-for-empty-message branch from 2d853e1 to 0fcd275 Compare July 18, 2026 14:42
@aradng
aradng force-pushed the fix/fastapi-body-none-for-empty-message branch from 0fcd275 to 34b34ee Compare August 3, 2026 06:57
@aradng
aradng force-pushed the fix/fastapi-body-none-for-empty-message branch from d908b7a to 19432c6 Compare August 4, 2026 07:06
@aradng

aradng commented Aug 4, 2026

Copy link
Copy Markdown
Contributor Author

@Lancetnik

@Lancetnik Lancetnik left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.py checks 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: bytes handlers keep receiving b"" — 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

  • Tombstone is public and constructible, but the guards disagree: the key requirement checks cmd.body is TOMBSTONE while encode_or_tombstone checks isinstance. publish(Tombstone(), topic) without a key skips the friendly error and reaches aiokafka as key=None, value=None.
  • __bool__ returns False with no __eq__, so TOMBSTONE == Tombstone() is False while isinstance passes — and if not msg.body: conflates a tombstone with b"" again, the exact distinction the sentinel exists to preserve. (EmptyPlaceholder in _internal/constants.py defines 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-kafka accepts it.
  • The batch + custom BatchCodecProto combination 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
@aradng
aradng force-pushed the fix/fastapi-body-none-for-empty-message branch from 19432c6 to 80810e8 Compare August 20, 2026 16:50
@aradng aradng changed the title fix(fastapi): resolve Optional body params to None for a real tombstone feat(kafka): explicit TOMBSTONE for publishing and consuming null records Aug 20, 2026
@aradng

aradng commented Aug 20, 2026

Copy link
Copy Markdown
Contributor Author

Rebased onto current main (was 23 behind) and restructured — description rewritten above.

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:

  • Tombstone/TOMBSTONE moved next to EmptyPlaceholder in _internal/constants.py, same __repr__/__bool__/__eq__ shape;
  • Tombstone added to SendableMessage — both mypy errors you pasted are gone (the leftover KafkaResponse-in-publish_batch one reproduces on main, so it's not from this branch);
  • one guard spelling everywhere (isinstance), so publish(Tombstone(), topic) without a key hits the friendly error instead of reaching aiokafka as key=None, value=None;
  • providers size through body_size/batch_body_size, testing brokers publish a real null value.

On revoking publish(None) — I'd rather not keep None overloaded, but you're right that 0.7.x can't take that away after #2932/#2939. Happy to add it back as a deprecated alias that warns and still produces a tombstone, removed in 0.8. That makes the publish side additive and leaves msg.body as the only breaking change. Say the word and I'll push it.

@aradng
aradng force-pushed the fix/fastapi-body-none-for-empty-message branch from 80810e8 to c4c47b4 Compare August 20, 2026 16:53
@aradng

aradng commented Aug 20, 2026

Copy link
Copy Markdown
Contributor Author

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 test_publish_tombstone_without_key_raises.

Batch + custom BatchCodecProto (your nit 4). Still raises — encode_batch() has no way to express a null value for one record of a batch, so there's no correct thing to encode. Now documented in both message.md pages alongside the key requirement, rather than being a runtime surprise. Just pushed.

@aradng
aradng force-pushed the fix/fastapi-body-none-for-empty-message branch from c4c47b4 to f6e4c30 Compare August 20, 2026 18:41
@aradng

aradng commented Aug 20, 2026

Copy link
Copy Markdown
Contributor Author

Reworked this so it isn't a breaking change anymore.

Tombstone now subclasses bytes and is empty. So msg.body for a null record is == b"", len() == 0, falsy, still bytes. Existing consumers see exactly what they see today, and isinstance(body, Tombstone) is what tells the two apart.

That drops most of what you flagged. publish(None) keeps working exactly as it does now, just with a DeprecationWarning, so nothing shipped in 0.7.2/0.7.3 gets revoked. decode_message, encode_message, SendableMessage and all four Prometheus/OTel providers are byte identical to main, so the len() crash from your first review can't come back. Public surface is down to TOMBSTONE and Tombstone.

Also folded the duplicated batch encode into one helper next to BatchCodecProto, and both mypy errors you pasted are gone.

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.

@aradng

aradng commented Aug 26, 2026

Copy link
Copy Markdown
Contributor Author

would appreciate a re-review @Lancetnik :)

@aradng
aradng force-pushed the fix/fastapi-body-none-for-empty-message branch from f6e4c30 to 56ee916 Compare August 26, 2026 07:53
@borisalekseev
borisalekseev self-requested a review August 27, 2026 12:31
@borisalekseev

Copy link
Copy Markdown
Collaborator

Hi! I’ll join the review this week

@Lancetnik

Lancetnik commented Aug 27, 2026 •

Copy link
Copy Markdown
Member

#3063 adds faststream/_internal/kafka/ — the home for machinery both Kafka packages need, since they're independent by design (faststream/confluent/ imports nothing from faststream/kafka/). We'd like this feature to live entirely in the Kafka packages plus that module, with no changes under faststream/_internal/ or faststream/message/. Concretely:

Edited. Items 3, 4 and 6 have changed since I first posted this — I wrote them before settling the review on #3028, and two of them contradicted it. Items 1, 2, 5 and 7 are unchanged. Sorry for the churn.

1. Move Tombstone / TOMBSTONE into _internal/kafka/, re-exported from faststream.kafka and faststream.confluent.

_internal/constants.py is for framework sentinels like EMPTY. Keeping the type out of faststream.message also means a Redis or NATS user never sees TOMBSTONE in their public API.

2. Move value_or_tombstone, encode_or_tombstone, and encode_batch_or_tombstone there too, and leave them out of __all__.

They're wiring for the Kafka parsers and producers. The per-message key helpers already in _internal/kafka/ are unexported for the same reason.

3. Drop the tombstone flag from this PR entirely, and rebase on #3028.

#3028 is landing a neutral no_body flag on StreamMessage plus KafkaMessage.tombstone as a property over it, which already answers "is this single record a tombstone". A second flag here would be a second source of truth for the same fact — and the one in this PR is derived from the body anyway, so the two could disagree.

4. Replace the batch flag with per-element detection.

any(...) reports that some element is a tombstone but not which, and reads as "this batch is a tombstone" — the wrong inference to offer. Per-element isinstance(body[i], Tombstone) is the right shape, and it's the one thing a single flag genuinely cannot express, so this is where the sentinel earns its keep. Batches are this PR's to settle; #3028 deliberately leaves them alone.

5. Move the decode_message special case into AioKafkaParser.decode_message and its Confluent counterpart.

That's the per-broker decode hook, so faststream/message/utils.py stays out of it.

6. Drop the _internal/fastapi/route.py change — #3028 delivers it.

Keep the two stacked; there's nothing left to duplicate here once #3028 is in.

Separately, and not yours to fix: route.py raises RequestValidationError for any decoder returning None, because {first_arg: None} gets resolved against the model. We'll fix that on its own.

7. Drop the publish(None) deprecation.

"None will be encoded normally in 0.8" decides the question #3062 was opened for, and it changes None semantics for every broker. That call belongs there.

That leaves this PR as the Tombstone sentinel, the batch story, and an explicit TOMBSTONE on publish — no core changes, built on top of #3028 rather than duplicating it.

Lancetnik and others added 3 commits August 27, 2026 18:26
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>
@aradng
aradng force-pushed the fix/fastapi-body-none-for-empty-message branch from 56ee916 to 8f90891 Compare September 17, 2026 18:19
…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>
@aradng
aradng force-pushed the fix/fastapi-body-none-for-empty-message branch from 8f90891 to b7f61d4 Compare September 17, 2026 18:44
@aradng
aradng requested a review from draincoder as a code owner September 17, 2026 18:44
…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>
@aradng
aradng force-pushed the fix/fastapi-body-none-for-empty-message branch from b7f61d4 to 06088c6 Compare September 17, 2026 18:44
@aradng

aradng commented Sep 17, 2026

Copy link
Copy Markdown
Contributor Author

@Lancetnik @borisalekseev ready for another look — rewritten to your 7-item list, and rebased on #3028 (which now carries the no_body rename and KafkaMessage.tombstone).

All seven are done; the description above walks through them one by one. The short version: Tombstone/TOMBSTONE and the three encode helpers moved into _internal/kafka/, re-exported from the two Kafka packages only and the helpers out of every __all__; the second flag is gone in favour of #3028's; batches mark per element; the decode_message case moved into each broker's parser; no route.py change; no publish(None) deprecation. _internal/constants.py, _internal/parser.py, faststream/message/ and the four telemetry providers are byte identical to main. 17 files instead of 25.

Your two open nits. The key requirement now only applies to an explicit Tombstone, which is new API — a keyless publish(None) on confluent still tombstones exactly as on main, so the independent breaking change you flagged isn't riding along anymore, and there's now a test pinning that divergence. The BatchCodecProto combination still raises, documented in both message.md pages, since encode_batch() has no way to express a null value for one record.

Two things worth your judgement:

  1. Tombstone overrides __str__ as well as __repr__. bytes doesn't route str(), f-strings or %s through __repr__, so without it a tombstone reads as b"" in every log line that doesn't use !r — which is the one place a human, rather than an isinstance, needs to tell it from an empty payload. Happy to drop it if you'd rather the sentinel carried no behavior beyond the type.

  2. Not this PR's to fix, but this PR makes it reachable: realign_keys in _internal/kafka/keys.py locates a body with current_bodies.index(body), and Tombstone() == b"". A batch holding both a real b"" and a tombstone can resolve to the wrong index there. Pre-existing, no test covers it either way — flagging rather than fixing it here.

Test-wise, the cases both brokers share now live in tests/brokers/base/ behind a get_message_value hook, the way BatchKeysTestcase already does it for keys, instead of being copy-pasted per broker. The class's own contract (constructor guard, pickle/copy round-trip, how it prints) has its own unit test, and a tombstone now goes through all four Prometheus/OTel providers single and batch — the len() crash that shaped this whole design is impossible by test now, not by argument. The pre-existing test_publish_none_tombstone and test_consume_*_without_value pass unmodified, which is the back-compat proof.

CI green, including both real-broker suites. ruff clean, mypy clean across 519 source files.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

AioKafka Issues related to `faststream.kafka` module Confluent Issues related to `faststream.confluent` module documentation Improvements or additions to documentation

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants