Skip to content

feat(codec): wire codecs into publishing flow - #2902

Open
ce1ebrimbor wants to merge 28 commits into
ag2ai:mainfrom
ce1ebrimbor:feat/codec-destination
Open

ce1ebrimbor wants to merge 28 commits into
ag2ai:mainfrom
ce1ebrimbor:feat/codec-destination

Conversation

@ce1ebrimbor

@ce1ebrimbor ce1ebrimbor commented Jun 3, 2026 •

Copy link
Copy Markdown
Collaborator

Description

Unified codec interface compatible with all the brokers.

WIP #2837

Type of change

Please delete options that are not relevant.

  • New feature (a non-breaking change that adds functionality)

Checklist

  • My code adheres to the style guidelines of this project (just lint shows no errors)
  • I have conducted a self-review of my own code
  • I have made the necessary changes to the documentation
  • My changes do not generate any new warnings
  • I have added tests to validate the effectiveness of my fix or the functionality of my new feature
  • Both new and existing unit tests pass successfully on my local environment by running just test-coverage
  • I have ensured that static analysis tests are passing by running just static-analysis
  • I have included code examples to illustrate the modifications

@github-actions github-actions Bot added Confluent Issues related to `faststream.confluent` module AioKafka Issues related to `faststream.kafka` module NATS Issues related to `faststream.nats` module and NATS broker features MQTT Issues related to `faststream.mqtt` module labels Jun 3, 2026
@ce1ebrimbor
ce1ebrimbor force-pushed the feat/codec-destination branch from d20c610 to 6a33f03 Compare June 3, 2026 10:23
@github-actions github-actions Bot added the documentation Improvements or additions to documentation label Jun 3, 2026
@Lancetnik

Copy link
Copy Markdown
Member

GzipCodec is not prefect example because you can achieve the same behavior using middleware - https://github.com/ulbwa/faststream-compressors

Comment thread faststream/confluent/publisher/producer.py Outdated
@ce1ebrimbor

Copy link
Copy Markdown
Collaborator Author

GzipCodec is not prefect example because you can achieve the same behavior using middleware - https://github.com/ulbwa/faststream-compressors

I will add a schema registry example, I think this will make more sense.

@ce1ebrimbor
ce1ebrimbor force-pushed the feat/codec-destination branch from 6a33f03 to 3e0d1d1 Compare June 4, 2026 19:14
@github-actions github-actions Bot added the Redis Issues related to `faststream.redis` module and Redis features label Jun 4, 2026
@ce1ebrimbor
ce1ebrimbor force-pushed the feat/codec-destination branch 4 times, most recently from 9b45655 to 75cc013 Compare June 4, 2026 19:38
@ce1ebrimbor

Copy link
Copy Markdown
Collaborator Author

I will add a schema registry example, I think this will make more sense.

@Lancetnik I have added a Schema Registry example I have previously used.
I will need some help for other brokers.

I won't be available next week, if there are new suggestions I will be able to add them after.

@ce1ebrimbor
ce1ebrimbor force-pushed the feat/codec-destination branch 2 times, most recently from 9a14d12 to 1239f71 Compare June 4, 2026 19:45
@ce1ebrimbor

ce1ebrimbor commented Jun 4, 2026 •

Copy link
Copy Markdown
Collaborator Author

@Lancetnik I am not comforable with the fact that Rabbit and Redis still need a destination parameter
If we want to keep the consistency we should reworkvthose, your call.
Otherwise let's keep this PR as wip. 🙏

Comment thread faststream/confluent/publisher/producer.py Outdated
Comment thread faststream/confluent/testing.py Outdated
Comment thread faststream/confluent/testing.py Outdated
Comment thread faststream/kafka/testing.py Outdated

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

We should unify behavior and signature for all broker codecs

@ce1ebrimbor
ce1ebrimbor marked this pull request as draft August 6, 2026 18:16
@ce1ebrimbor
ce1ebrimbor force-pushed the feat/codec-destination branch 2 times, most recently from 2b2657e to efe2d10 Compare August 6, 2026 19:46
…and Rabbit producers

Unifies the codec API across all brokers. Previously Redis and Rabbit
constructed intermediate PublishCommand objects or decomposed the command
into scalar arguments before calling codec.encode(). Now the actual
RedisPublishCommand/RabbitPublishCommand flows through to the codec,
giving it access to destination and all broker-specific metadata.

Changes:
- redis/parser/message.py: MessageFormat.build() accepts PublishCommand
- redis/parser/binary.py: BinaryMessageFormatV1.encode() accepts cmd
- redis/publisher/producer.py: pass cmd directly to encode()
- rabbit/parser.py: AioPikaParser.encode_message() accepts optional cmd
- rabbit/publisher/producer.py: pass RabbitPublishCommand as cmd
- rabbit/testing.py: construct RabbitPublishCommand for test messages
- Update test callers to use new encode(cmd=...) signature
Move PublishCommand, PublishType imports from function bodies to module
level in confluent/publisher/producer.py and confluent/testing.py.
Rename _BasePublishCommand to _BaseCmd for consistency.
Move PublishCommand, PublishType imports from function bodies to module
level in kafka/publisher/producer.py and kafka/testing.py.
Rename _BasePublishCommand to _BaseCmd for consistency.
Move PublishType, PublishCommand imports from function bodies to module
level in redis/publisher/producer.py, redis/testing.py,
rabbit/parser.py, and rabbit/testing.py.
Since cmd is now required, destination was dead code. Also remove
unused PublishType and _BasePublishCommand imports from rabbit/parser.py.
Introduce EncodedMessage(body, content_type) dataclass as the return
type for CodecProto.encode() and BatchCodecProto.encode_batch().

This replaces the previous tuple[bytes, str | None] contract with a
named, extensible structure. All broker producers, testing files, and
the schema registry example are updated to use encoded.body and
encoded.content_type directly.
Add warn_deprecated_param() to _compat.py. Emit warnings at:
- BrokerConfig.__post_init__ (broker-level parser/decoder)
- SubscriberUsecase.add_call() (subscriber-level parser/decoder)
- SubscriberUsecase.__call__() (handler-level parser/decoder)

Users should migrate to codec with custom encode()/decode().
parser remains valid (raw msg → StreamMessage). Only decoder is
replaced by codec.decode(). Inline warnings directly instead of
using warn_deprecated_param helper.
@ce1ebrimbor
ce1ebrimbor force-pushed the feat/codec-destination branch from 2eddff8 to df5e22f Compare August 18, 2026 16:49
@ce1ebrimbor

Copy link
Copy Markdown
Collaborator Author

@Lancetnik if you have time 🙏

@Lancetnik

Copy link
Copy Markdown
Member

@ce1ebrimbor hi! Sorry for waiting so long. I'll review and merge the PR this week
Incredible job, thank you!
BTW, can we contact some way? Mail / Linkedin / Telegram / WA? What format works for you?

@ce1ebrimbor

Copy link
Copy Markdown
Collaborator Author

@ce1ebrimbor hi! Sorry for waiting so long. I'll review and merge the PR this week
Incredible job, thank you!
BTW, can we contact some way? Mail / Linkedin / Telegram / WA? What format works for you?

Hey, I will ping you on telegram ^^

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 MQTT Issues related to `faststream.mqtt` module NATS Issues related to `faststream.nats` module and NATS broker features Redis Issues related to `faststream.redis` module and Redis features

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants