feat(gateways): send Kafka records to Iggy from the Produce handler - #4282
Merged
Merged
Conversation
Codecov Report✅ All modified and coverable lines are covered by tests. Additional details and impacted files@@ Coverage Diff @@
## master #4282 +/- ##
============================================
+ Coverage 87.57% 87.65% +0.08%
- Complexity 1575 1576 +1
============================================
Files 1284 1284
Lines 225457 225457
Branches 188821 188821
============================================
+ Hits 197438 197630 +192
+ Misses 23289 23078 -211
- Partials 4730 4749 +19
🚀 New features to boost your workflow:
|
numinnex
reviewed
Sep 24, 2026
7 tasks
ryerraguntla
previously approved these changes
Sep 26, 2026
ryerraguntla
left a comment
Contributor
There was a problem hiding this comment.
Resolve the conflicts in api.rs. Most of the the findings are documented in the PR description as limitations and known issues. Document the known issues and limitations in the form of an issue. I can add this newly created issue as part of parity checking issue. Once the end to end flow is tested along with related API keys this newly created ticket could be validated and resolved
Action items before merging
- resolve first reviewers comments
- Resolve api.rs conflicts and successful prechecks
- create a new ticket for all known issues and add to parity tracking.
The record codec and the bridge were both in place, so Produce was the only thing still discarding every payload it decoded.
IDEMPOTENCE.md chose allocate only, so Produce ignores the producer id and a stock producer works once #4250 lands.
krishvishal
force-pushed
the
kafka-produce
branch
from
September 27, 2026 03:31
e6b7466 to
863c1a7
Compare
mmodzelewski
approved these changes
Sep 27, 2026
spetz
approved these changes
Sep 27, 2026
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Part of #3535. With
IGGY_KAFKA_BRIDGE_ENABLED=true, Produce now writes decoded records to Iggy using the codec from #4229 and bridge from #3533, replacing the record-dropping stub.#3535 remains open: its first acceptance criterion requires a real producer run. That waits for Metadata (#3534, PR #4258), which lets clients discover the topic and reach Produce.
Writes and responses
send_messagescall, routed to the client-selectedPartitioning::partition_id, without balanced partitioning. Every batch is decoded into one Iggy message per record perdocs/BRIDGE_MAPPING.md; batches are not stored as opaque bytes.log_start_offsetis always-1, avoiding an extra round trip for a field producers do not read.ResourceNotFounderrors returnUNKNOWN_TOPIC_OR_PARTITION(3).docs/IDEMPOTENCE.md. feat(gateways): support idempotent Kafka producers, refuse transactions #4250 enables default Java producers to start. Retries are not deduplicated; delivery remains at-least-once.Validation
INVALID_RECORD(87)UNSUPPORTED_VERSION(35), stopping the producerUNSUPPORTED_COMPRESSION_TYPE(76)acks= 0, 1 or -1acksINVALID_REQUIRED_ACKS(21)The README documents every field and code.
ERROR_*constants inapi.rsnow take their values from thekafka-protocolResponseErrortable, preventing numeric typos.Limits and failure handling
timeout_ms, capped at 20 seconds. Expired partitions returnREQUEST_TIMED_OUT(7).max_frame_sizedecompressed bytes andmax_frame_size / 64records, with four requests decoding concurrently. A partition exceeding the budget alone returnsMESSAGE_TOO_LARGE(10); one that fits alone but exceeds the remaining budget returnsNOT_LEADER_OR_FOLLOWER(6) for retry. Large or compressed batches decode on a blocking thread.IGGY_KAFKA_IGGY_MAX_MESSAGE_SIZE(default64MiB) mirrors Iggy'smessage_bus.max_message_size; larger partitions return 10.u32::MAXmicroseconds apart (~71 minutes). Wider spans require multiple sends; if a later send fails, retrying can duplicate earlier records.Disconnected,EmptyResponse,TcpErrorandStaleClientreturn 7 rather than 6. The SDK does not resend; records may already be stored.max.in.flight.requests.per.connection=1for strict ordering.Tests and remaining work
tests/produce_real_bridge_tests.rsadds 22 real-iggy-servertests to thekafka_bridgenextest group. They submit Kafka wire bytes and read records through the Iggy SDK, so matching mapping errors in both directions cannot hide a failure. They callhandle_request_boundeddirectly, bypassing the TCP listener.After #4258 merges, run
kcat -Pand a defaultkafka-console-producer.shthrough the gateway, record the results in row G2 ofdocs/MANUAL_TESTING.md, and close #3535. The console producer also requires InitProducerId from #4250, orenable.idempotence=false.#4250 also changes the Produce stub in
produce.rs; if it merges first, rebase this branch and preserve itstransactional_idguard. Other Produce TODOs, including a client pool, remain indocs/SCOPE.md.