Skip to content

feat(gateways): send Kafka records to Iggy from the Produce handler - #4282

Merged
krishvishal merged 4 commits into
masterfrom
kafka-produce
Sep 27, 2026
Merged

krishvishal merged 4 commits into
masterfrom
kafka-produce

Conversation

@krishvishal

Copy link
Copy Markdown
Member

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

  • Each partition entry normally becomes one send_messages call, routed to the client-selected Partitioning::partition_id, without balanced partitioning. Every batch is decoded into one Iggy message per record per docs/BRIDGE_MAPPING.md; batches are not stored as opaque bytes.
  • Partitions respond independently: a failure gets its own error code without losing other partitions' offsets. Base offsets come from send confirmations. log_start_offset is always -1, avoiding an extra round trip for a field producers do not read.
  • Produce never creates topics. Missing topics and send-side ResourceNotFound errors return UNKNOWN_TOPIC_OR_PARTITION (3).
  • Idempotent batches are accepted, but producer id, epoch and sequence are ignored, following the “allocate only” decision in 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

Condition Behavior
Multiple batches in one partition, or more records than declared INVALID_RECORD (87)
Transactional or control batch UNSUPPORTED_VERSION (35), stopping the producer
zstd compression Supported from Produce v7; older versions return UNSUPPORTED_COMPRESSION_TYPE (76)
acks = 0, 1 or -1 Same write path
Any other acks INVALID_REQUIRED_ACKS (21)

The README documents every field and code. ERROR_* constants in api.rs now take their values from the kafka-protocol ResponseError table, preventing numeric typos.

Limits and failure handling

  • Deadline: honors timeout_ms, capped at 20 seconds. Expired partitions return REQUEST_TIMED_OUT (7).
  • Decode budget: each request allows max_frame_size decompressed bytes and max_frame_size / 64 records, with four requests decoding concurrently. A partition exceeding the budget alone returns MESSAGE_TOO_LARGE (10); one that fits alone but exceeds the remaining budget returns NOT_LEADER_OR_FOLLOWER (6) for retry. Large or compressed batches decode on a blocking thread.
  • SDK backpressure: only one send runs inside the SDK at a time because an active send cannot be cancelled. A send outliving its deadline retains the slot for up to 45 seconds. Others wait in the gateway until their own deadlines, then return 7, rather than queuing inside the SDK.
  • Send size: IGGY_KAFKA_IGGY_MAX_MESSAGE_SIZE (default 64MiB) mirrors Iggy's message_bus.max_message_size; larger partitions return 10.
  • Timestamp span: an Iggy send supports timestamps at most u32::MAX microseconds apart (~71 minutes). Wider spans require multiple sends; if a later send fails, retrying can duplicate earlier records.
  • Unknown write outcome: send-side Disconnected, EmptyResponse, TcpError and StaleClient return 7 rather than 6. The SDK does not resend; records may already be stored.
  • Ordering: each Kafka connection handles one request at a time, but retries with multiple requests in flight can place a failed batch after later batches. The README requires max.in.flight.requests.per.connection=1 for strict ordering.

Tests and remaining work

tests/produce_real_bridge_tests.rs adds 22 real-iggy-server tests to the kafka_bridge nextest 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 call handle_request_bounded directly, bypassing the TCP listener.

After #4258 merges, run kcat -P and a default kafka-console-producer.sh through the gateway, record the results in row G2 of docs/MANUAL_TESTING.md, and close #3535. The console producer also requires InitProducerId from #4250, or enable.idempotence=false.

#4250 also changes the Produce stub in produce.rs; if it merges first, rebase this branch and preserve its transactional_id guard. Other Produce TODOs, including a client pool, remain in docs/SCOPE.md.

@github-actions github-actions Bot added the S-waiting-on-review PR is waiting on a reviewer label Sep 24, 2026
@codecov

codecov Bot commented Sep 24, 2026 •

Copy link
Copy Markdown

Codecov Report

✅ All modified and coverable lines are covered by tests.
✅ Project coverage is 87.65%. Comparing base (b8bf85f) to head (863c1a7).

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     
Components Coverage Δ
Rust Core 88.75% <ø> (-0.02%) ⬇️
Java SDK 68.70% <ø> (+0.01%) ⬆️
C# SDK 77.41% <ø> (-0.02%) ⬇️
Python SDK 90.97% <ø> (ø)
PHP SDK 85.67% <ø> (ø)
Node SDK 96.43% <ø> (+1.68%) ⬆️
Go SDK 70.11% <ø> (+0.05%) ⬆️
see 47 files with indirect coverage changes
🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

Comment thread gateways/kafka/src/records.rs
Comment thread gateways/kafka/src/records.rs Outdated
Comment thread gateways/kafka/src/records.rs Outdated
Comment thread gateways/kafka/src/protocol/handlers/produce.rs Outdated
ryerraguntla
ryerraguntla previously approved these changes Sep 26, 2026

@ryerraguntla ryerraguntla left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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

  1. resolve first reviewers comments
  2. Resolve api.rs conflicts and successful prechecks
  3. 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
krishvishal merged commit 8d0e28d into master Sep 27, 2026
182 of 186 checks passed
@krishvishal
krishvishal deleted the kafka-produce branch September 27, 2026 05:37
@github-actions github-actions Bot removed the S-waiting-on-review PR is waiting on a reviewer label Sep 27, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Kafka gateway: Kafka Produce API (key 0) — bridge to Iggy send_messages

5 participants