feat(codec): wire codecs into publishing flow - #2902
ce1ebrimbor wants to merge 28 commits into
Conversation
d20c610 to
6a33f03
Compare
|
|
I will add a schema registry example, I think this will make more sense. |
6a33f03 to
3e0d1d1
Compare
9b45655 to
75cc013
Compare
@Lancetnik I have added a Schema Registry example I have previously used. I won't be available next week, if there are new suggestions I will be able to add them after. |
9a14d12 to
1239f71
Compare
|
@Lancetnik I am not comforable with the fact that Rabbit and Redis still need a destination parameter |
Lancetnik
left a comment
There was a problem hiding this comment.
We should unify behavior and signature for all broker codecs
2b2657e to
efe2d10
Compare
…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.
… inherits CodecProto
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.
2eddff8 to
df5e22f
Compare
|
@Lancetnik if you have time 🙏 |
|
@ce1ebrimbor hi! Sorry for waiting so long. I'll review and merge the PR this week |
Hey, I will ping you on telegram ^^ |
Description
Unified codec interface compatible with all the brokers.
WIP #2837
Type of change
Please delete options that are not relevant.
Checklist
just lintshows no errors)just test-coveragejust static-analysis