Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
b410d91
Initial version of CreateTopics and MetaData implementation
ryerraguntla Sep 22, 2026
dafcba8
updating documentation
ryerraguntla Sep 22, 2026
8a0a0d7
Merge branch 'master' into feat(gateways)/3534-3538-kafka-metadata-cr…
ryerraguntla Sep 22, 2026
cd398c7
Fixed CreateTopics/Metadata review finding
ryerraguntla Sep 23, 2026
b82675e
Merge branch 'master' into feat(gateways)/3534-3538-kafka-metadata-cr…
ryerraguntla Sep 23, 2026
5546dc8
Fixed on review comments
ryerraguntla Sep 24, 2026
e1329fb
Merge branch 'feat(gateways)/3534-3538-kafka-metadata-createtopics' o…
ryerraguntla Sep 24, 2026
941ba08
Merge branch 'master' into feat(gateways)/3534-3538-kafka-metadata-cr…
ryerraguntla Sep 24, 2026
9d5eda8
Fixed review comments
ryerraguntla Sep 24, 2026
78c57e0
fixing the doc test
ryerraguntla Sep 24, 2026
fdc8572
Merge branch 'master' into feat(gateways)/3534-3538-kafka-metadata-cr…
ryerraguntla Sep 24, 2026
0ce3343
fixing doc test and clippy errors
ryerraguntla Sep 24, 2026
19d7b31
Merge branch 'feat(gateways)/3534-3538-kafka-metadata-createtopics' o…
ryerraguntla Sep 24, 2026
6920b66
Merge branch 'master' into feat(gateways)/3534-3538-kafka-metadata-cr…
ryerraguntla Sep 25, 2026
63a6b57
Fixed clippy errors from the merge issues
ryerraguntla Sep 26, 2026
2918e93
Merge branch 'master' into feat(gateways)/3534-3538-kafka-metadata-cr…
ryerraguntla Sep 26, 2026
9b73687
Merge branch 'master' into feat(gateways)/3534-3538-kafka-metadata-cr…
ryerraguntla Sep 27, 2026
912cc34
Resolving merge conflicts
ryerraguntla Sep 27, 2026
ab15f25
Merge branch 'feat(gateways)/3534-3538-kafka-metadata-createtopics' o…
ryerraguntla Sep 27, 2026
0a0f871
Fixing merge issues
ryerraguntla Sep 27, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions gateways/kafka/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,18 @@

Foundation layer for [apache/iggy#3421](https://github.com/apache/iggy/issues/3421): a TCP listener on the Kafka wire port that decodes requests, validates scoped API keys and versions. With a bridge, Produce writes to Iggy and ListOffsets reads offsets from it. Everything else is a stub.

> **Stub warning:** Produce and Fetch still don't persist or read real data - they return
> retriable `NOT_LEADER_OR_FOLLOWER` (6) so clients keep data locally / retry elsewhere instead of
> trusting a fake success. CreateTopics, Metadata, and ListOffsets are wired to the Iggy bridge:
> with `IGGY_KAFKA_BRIDGE_ENABLED=true`, CreateTopics creates a real Iggy stream/topic, Metadata
> reports real topics and partition counts (a topic not requested by name and not found is
> silently absent from a null-topics "list all" response, and `UNKNOWN_TOPIC_OR_PARTITION` when
> named explicitly), and ListOffsets answers `EARLIEST`/`LATEST` from real partition state; with
> the bridge off (the default), all three stay stubs - CreateTopics answers `NOT_CONTROLLER` (41),
> Metadata reports every requested topic unknown, and ListOffsets answers `NOT_LEADER_OR_FOLLOWER`
> (6). **CreateTopics has no authentication gate yet**: with the bridge on, any client that can
> reach this port can create topics (up to 1000 partitions each) as the bridge's own Iggy user,
> until SASL ([#3549](https://github.com/apache/iggy/issues/3549)) lands. See
> **Stub warning:** When you set `IGGY_KAFKA_BRIDGE_ENABLED=true`, Produce writes to Iggy and
> ListOffsets answers `EARLIEST`/`LATEST` from real partition state. No other API stores or reads
> real data. Produce and ListOffsets without a bridge, and Fetch with or without one, answer
Expand Down
96 changes: 96 additions & 0 deletions gateways/kafka/docs/SCOPE.md
Original file line number Diff line number Diff line change
Expand Up @@ -117,10 +117,106 @@ below it are still open for the issues that build on top of it.
- [ ] Fetch → `poll_messages` ([#3536](https://github.com/apache/iggy/issues/3536)).
- [x] Idempotent `ensure_stream_and_topic()` (create-if-not-exists) - `src/bridge/iggy_bridge.rs`,
exercised end-to-end in `tests/bridge_iggy_integration_tests.rs`.
- [x] Real CreateTopics ([#3538](https://github.com/apache/iggy/issues/3538)): with
`IGGY_KAFKA_BRIDGE_ENABLED=true`, creates the Iggy stream/topic through
`IggyBridge::create_kafka_topic` - an atomic create-or-report-exists call, not a separate
existence read followed by an idempotent create (that sequence has a TOCTOU window: two
concurrent requests for the same new name could both observe "doesn't exist yet" and both
receive `NONE`, when Kafka guarantees exactly one caller does).
`src/protocol/handlers/create_topics.rs`, `tests/create_topics_real_bridge_tests.rs`. With
the bridge off, the stub from #3421 answers `NOT_CONTROLLER` (41) as before.
- Every occurrence of a duplicate topic name within one request is rejected with
`INVALID_REQUEST` (42) and nothing is created for it, matching real Kafka
(`ControllerApis.createTopics`) rather than creating the first occurrence and reporting the
rest as already existing.
- A manual partition `assignments` list combined with an explicit `num_partitions`/
`replication_factor` is rejected with `INVALID_REQUEST` (42) even when the count agrees with
`assignments.len()` - real Kafka's own `ReplicationControlManager` treats the two as mutually
exclusive inputs, not independently-checked values that happen to agree.
- A manual assignment's own partition indices are checked too, not just its length:
`ERROR_INVALID_REPLICA_ASSIGNMENT` (39) for a duplicate index or one that isn't exactly
`0..assignments.len()` - `{5: [...], 7: [...]}` has the right length for a 2-partition topic
but names neither partition `0` nor `1`, matching real Kafka's own
`ReplicationControlManager.createTopic` key-set validation.
- `IggyError::RequestAlreadyApplied` (the SDK's reconnect path replayed a write that already
committed) maps to `NONE`, not the `UNKNOWN_SERVER_ERROR` catch-all - the operation did
succeed, and a Java client treats `UNKNOWN_SERVER_ERROR` as non-retriable. Special-cased
locally in `create_topics.rs`, not in the shared `BridgeError -> Kafka error code` mapping
every handler's error path goes through: "the write already applied" is a write-only fact,
and a read (Metadata, ListOffsets) reaching this variant has no write to have applied.
- `error_message` on a rejected topic never re-embeds the topic name `CreatableTopicResult.name`
already carries - `BridgeError::InvalidKafkaTopicName`'s `Display` does, so using it directly
would roughly double the response cost per invalid name, for free (name validation runs before
any bridge I/O).
- `IggyBridge::get_kafka_topic` no longer probes `get_stream` separately before `get_topic` -
`get_topic` already answers `Ok(None)` when the stream itself is missing, so the probe was a
second round trip to learn something the one call already told it.
- Bridge fan-out is bounded independently of `bounds_guard`'s `MAX_REQUEST_ELEMENTS` (4,096,
still a pre-decode ceiling, not a usability one): a duplicate name never reaches the bridge at
all (rejected up front, see above), a request naming more than 100 distinct non-duplicate
topics is rejected outright (`INVALID_REQUEST`, no bridge call for any of them), and the
request's own wire `timeout_ms` (clamped to `[1s, 30s]`) now bounds the whole handler's
aggregate bridge work, not just decoded and discarded - a deadline that fires answers every
topic `REQUEST_TIMED_OUT` rather than continuing to hold the shared lockstep `IggyClient`.
- [x] Document partition mapping in [`BRIDGE_MAPPING.md`](BRIDGE_MAPPING.md):
- Iggy partitions are **0-based** (same as Kafka) — direct `partition_id` mapping, no offset conversion
- Kafka consumer groups do **not** map onto Iggy consumer groups. Assignment stays client-side, and Iggy's group registry is used as an offset key only ([`OFFSET_STORAGE.md`](OFFSET_STORAGE.md))
- `Partitioning::partition_id(index)` on every Produce. A Kafka producer resolves the partition before it builds the request, so `Partitioning::balanced()` has no trigger there. The `-1` default-partition-count case belongs to CreateTopics
- [x] Real Metadata topic/partition data ([#3534](https://github.com/apache/iggy/issues/3534)):
with `IGGY_KAFKA_BRIDGE_ENABLED=true`, a named lookup answers from
`IggyBridge::get_kafka_topic` (`UNKNOWN_TOPIC_OR_PARTITION` if not found) and a null
topics array ("all topics") answers from `IggyBridge::list_kafka_topics` - the target of
every configured `TopicMapping` override plus every other topic in the default stream, so
an Iggy stream this bridge has no mapping rule pointing at is never listed (nothing a Kafka
client ever named). Every reported partition names this gateway's single broker (node id 1)
as leader/replica/ISR, since there is only ever one. `src/protocol/handlers/metadata.rs`,
`src/bridge/iggy_bridge/topics.rs` (`get_kafka_topic`/`list_kafka_topics`),
`tests/metadata_real_bridge_tests.rs`. With the bridge off, the stub from #3421 reports
every requested topic unknown, as before.
- Fixed alongside: the stub's `decode_topics` collapsed a null topics array ("all topics") and
an explicit empty one (`Some(vec![])`, "these zero topics") to the same `Vec::new()` - the
real path's `decode_requested_topics` keeps `Option<Vec<_>>` so the two aren't conflated.
Also: at `api_version == 0` an explicit empty array is folded into "all topics" too, matching
Kafka's own `isAllTopics()` rule (`topics == null || (topics.isEmpty() && version == 0)`) -
there is no v0 wire shape for "cluster info only, zero topics" (that distinct shape, KIP-4's
`describeCluster()`, starts at v1).
- `get_kafka_topic` shares its implementation with CreateTopics' own existence check above -
both need "does this Kafka-side name resolve to a real Iggy topic," so this bridge exposes one
method returning the SDK's own `TopicDetails`, not two narrower, independently-maintained
lookups.
- The response-size cap above only covers the "all topics" and per-name-expansion cases. A
named lookup is separately bounded on the request side: names repeated in one request are
deduped up front to one response entry, not just one bridge round trip - real Kafka answers a
topic named twice in one request with one response entry, and re-expanding to match the
request would let a handful of repeats of one large topic name amplify a response sized off
the repeat count instead of the distinct count. No cap on distinct names: `bounds_guard`'s
`MAX_REQUEST_ELEMENTS` (4,096) is still the pre-decode ceiling, but `IggyBridge::get_kafka_topics`
batches by the *stream* each name resolves to rather than paying one round trip per name, so
the real bridge cost is bounded by distinct streams involved (config-time-bounded), not by
how many names the client asks about. A per-name cap here previously permanently broke a
long-lived Java producer once its `ProducerMetadata`'s cumulative tracked-topic set - resent
in full on every refresh - crossed the cap: every later request answered every topic
`INVALID_REQUEST`, and the producer had no way to shrink its own tracked set to recover. The
whole lookup's aggregate bridge work still runs under a fixed 20s wall-clock deadline
(`Metadata` carries no `timeout_ms` field in any version, unlike `CreateTopics`, so this cannot
be client-honored) - a deadline that fires answers every name `REQUEST_TIMED_OUT` rather than
continuing to hold the shared lockstep `IggyClient`.
- The "all topics" path is server-driven, not client-count-driven, so it has no distinct-topic
cap - only the response-size projection applies there, and it **truncates** rather than
closes on overflow: `list_kafka_topics()`'s result is trimmed to as many whole topics (in
listing order) as `max_frame_size` allows. Closing instead, as the named-lookup path still
does, would make every all-topics Metadata call - the bootstrap/refresh shape both librdkafka
and the Java client use - fail identically and permanently once the cluster's total partition
count crosses the trip point, since that size is the catalog's own, not anything the
requesting client chose or can shrink.
- [ ] Multi-broker topology (this gateway is, and will stay, a single logical broker - node id 1
always leads every partition it reports; nothing here models an Iggy cluster as multiple
Kafka-visible brokers)
- [ ] A raw/non-conformant client's extra trailing byte on some Metadata request shapes seen
against real `kcat`/`librdkafka` traffic in earlier testing on a since-restructured branch -
not reproduced or root-caused against the current `kafka_protocol`-based decode path in this
session, so not carried forward as a fix here rather than guessed at. Needs fresh
reproduction against a real `kcat` before it's re-closed.
- [x] Real ListOffsets ([#3537](https://github.com/apache/iggy/issues/3537)): with
`IGGY_KAFKA_BRIDGE_ENABLED=true`, `LATEST` answers from `IggyBridge::high_watermarks` and
`EARLIEST` answers `0`. Any other requested timestamp (arbitrary-timestamp offset search,
Expand Down
88 changes: 83 additions & 5 deletions gateways/kafka/src/bridge/error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,9 +19,9 @@ use iggy::prelude::IggyError;
use thiserror::Error;

use crate::protocol::api::{
ERROR_INVALID_PARTITIONS, ERROR_INVALID_TOPIC_EXCEPTION, ERROR_NOT_LEADER_OR_FOLLOWER,
ERROR_REQUEST_TIMED_OUT, ERROR_TOPIC_ALREADY_EXISTS, ERROR_TOPIC_AUTHORIZATION_FAILED,
ERROR_UNKNOWN_SERVER_ERROR, ERROR_UNKNOWN_TOPIC_OR_PARTITION,
ERROR_INVALID_PARTITIONS, ERROR_INVALID_REQUEST, ERROR_INVALID_TOPIC_EXCEPTION,
ERROR_NOT_LEADER_OR_FOLLOWER, ERROR_REQUEST_TIMED_OUT, ERROR_TOPIC_ALREADY_EXISTS,
ERROR_TOPIC_AUTHORIZATION_FAILED, ERROR_UNKNOWN_SERVER_ERROR, ERROR_UNKNOWN_TOPIC_OR_PARTITION,
};

/// Errors from the `IggyBridge`: connection lifecycle, config, and Iggy SDK calls.
Expand Down Expand Up @@ -80,6 +80,12 @@ pub enum BridgeError {
/// raw or non-conformant client, not an expected path.
#[error("invalid Kafka topic name '{kafka_topic}': {reason}")]
InvalidKafkaTopicName { kafka_topic: String, reason: String },
/// `ensure_stream_and_topic` was asked for `partition_count == 0`. Enforced here, not only
/// at the `CreateTopics` wire-validation layer that's this bridge's only caller today - the
/// method is `pub`, so a future caller (a test, or a later Produce auto-create path)
/// bypassing that layer must not be able to provision a topic nothing can produce to.
#[error("partition count must be at least 1, got 0 for topic '{kafka_topic}'")]
InvalidPartitionCount { kafka_topic: String },
}

impl BridgeError {
Expand All @@ -101,6 +107,9 @@ impl BridgeError {
Self::PartitionOutOfRange { .. } => ERROR_UNKNOWN_TOPIC_OR_PARTITION,
Self::PartitionCountMismatch { .. } => ERROR_TOPIC_ALREADY_EXISTS,
Self::InvalidKafkaTopicName { .. } => ERROR_INVALID_TOPIC_EXCEPTION,
// Same code the wire-validation layer already uses for this exact condition - this
// is the same failure reached through a different door, not a new kind of error.
Self::InvalidPartitionCount { .. } => ERROR_INVALID_PARTITIONS,
// Not a wire-response case in practice: an invalid bridge config is caught at
// `IggyBridge::connect` before any handler exists to answer a Kafka request, so this
// is reachable only if a future caller starts constructing configs at request time.
Expand Down Expand Up @@ -180,7 +189,19 @@ const fn iggy_error_to_kafka_code(err: &IggyError) -> i16 {
| IggyError::TcpError
| IggyError::TransientNotAccepted => ERROR_NOT_LEADER_OR_FOLLOWER,
IggyError::TransientNotCommitted => ERROR_REQUEST_TIMED_OUT,
IggyError::TooManyPartitions => ERROR_INVALID_PARTITIONS,
// Not `ERROR_INVALID_PARTITIONS` (37): that code's own text, per `kafka-protocol`'s
// table, is "Number of partitions is below 1" - the opposite condition from "too many"
// (Iggy's server-side cap, above 1000). Reusing 37 for both directions would return a
// client-visible error message that contradicts the actual request it sent.
IggyError::TooManyPartitions => ERROR_INVALID_REQUEST,
// Deliberately NOT special-cased here to ERROR_NONE: this function is shared by every
// handler's error path, but "the operation did commit, so report success" only holds for
// a caller that issued a *write* - the SDK's own reconnect path replayed a write whose
// first attempt already applied, and the server's client-table dedup caught the replay.
// A read that somehow reaches this variant has no write to have "already applied"; a
// write-side caller (CreateTopics) special-cases it locally, close to the write it
// concerns, instead of baking a write-only assumption into a mapping every read also
// goes through.
_ => ERROR_UNKNOWN_SERVER_ERROR,
}
}
Expand Down Expand Up @@ -236,8 +257,28 @@ mod tests {
}

#[test]
fn too_many_partitions_maps_to_invalid_partitions() {
fn too_many_partitions_maps_to_invalid_request_not_invalid_partitions() {
// INVALID_PARTITIONS (37) means "count is below 1" (kafka-protocol's own error table) -
// the opposite condition from "too many": reusing it here would send a client-visible
// message that contradicts the request it just sent.
let err = BridgeError::Iggy(IggyError::TooManyPartitions);
assert_eq!(err.to_kafka_error_code(), ERROR_INVALID_REQUEST);
}

#[test]
fn request_already_applied_falls_to_the_generic_mapping_here() {
// This shared mapping has no write to know "already applied" refers to - CreateTopics
// (the only caller for whom that's a success, not a fault) special-cases it locally
// instead (`create_topics.rs`), close to the write it concerns.
let err = BridgeError::Iggy(IggyError::RequestAlreadyApplied);
assert_eq!(err.to_kafka_error_code(), ERROR_UNKNOWN_SERVER_ERROR);
}

#[test]
fn invalid_partition_count_maps_to_invalid_partitions() {
let err = BridgeError::InvalidPartitionCount {
kafka_topic: "orders".to_string(),
};
assert_eq!(err.to_kafka_error_code(), ERROR_INVALID_PARTITIONS);
}

Expand Down Expand Up @@ -372,6 +413,43 @@ mod tests {
assert_eq!(err.to_kafka_error_code(), ERROR_TOPIC_ALREADY_EXISTS);
}

#[test]
fn every_sent_error_code_matches_kafka_protocols_own_table() {
use kafka_protocol::error::ResponseError;

for (ours, theirs) in [
(
ERROR_UNKNOWN_SERVER_ERROR,
ResponseError::UnknownServerError,
),
(
ERROR_UNKNOWN_TOPIC_OR_PARTITION,
ResponseError::UnknownTopicOrPartition,
),
(
ERROR_NOT_LEADER_OR_FOLLOWER,
ResponseError::NotLeaderOrFollower,
),
(ERROR_REQUEST_TIMED_OUT, ResponseError::RequestTimedOut),
(
ERROR_INVALID_TOPIC_EXCEPTION,
ResponseError::InvalidTopicException,
),
(
ERROR_TOPIC_AUTHORIZATION_FAILED,
ResponseError::TopicAuthorizationFailed,
),
(
ERROR_TOPIC_ALREADY_EXISTS,
ResponseError::TopicAlreadyExists,
),
(ERROR_INVALID_PARTITIONS, ResponseError::InvalidPartitions),
(ERROR_INVALID_REQUEST, ResponseError::InvalidRequest),
] {
assert_eq!(ours, theirs.code());
}
}

#[test]
fn send_lost_maps_to_request_timed_out() {
// The SDK does not replay a send after a lost connection, so it may have landed.
Expand Down
2 changes: 2 additions & 0 deletions gateways/kafka/src/bridge/iggy_bridge/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,8 @@ mod offsets;
mod produce;
mod topics;

pub use topics::{KafkaTopicMetadata, TopicCreationOutcome};

/// Passes attempted, after the first, before [`IggyBridge::connect`] gives up and returns `Err`.
///
/// Not the SDK's own default (`TcpClientReconnectionConfig::default()` is `max_retries: None` -
Expand Down
Loading
Loading