ryerraguntla commented on code in PR #4258:
URL: https://github.com/apache/iggy/pull/4258#discussion_r4093446918
##########
gateways/kafka/src/bridge/iggy_bridge/topics.rs:
##########
@@ -217,4 +255,151 @@ impl IggyBridge {
Err(err) => Err(err),
}
}
+
+ /// Looks up `kafka_topic`, resolved through the configured
+ /// [`TopicMapping`](crate::bridge::topic_map::TopicMapping), without
creating it.
+ ///
+ /// Returns `Ok(None)` when either the mapped stream or the mapped topic
doesn't exist -
+ /// callers (`CreateTopics`' existence check, `Metadata`'s lookup) treat
both the same way:
+ /// nothing answers to this Kafka-side name yet.
+ ///
+ /// # Errors
+ ///
+ /// Returns [`BridgeError::InvalidKafkaTopicName`] if `kafka_topic` fails
Kafka's own
+ /// topic-naming rules. Returns [`BridgeError::Timeout`] if a call takes
longer than
+ /// `REQUEST_TIMEOUT`. Returns [`BridgeError::Iggy`] for connectivity/auth
failures.
+ pub async fn get_kafka_topic(
+ &self,
+ kafka_topic: &str,
+ ) -> Result<Option<TopicDetails>, BridgeError> {
+ validate_kafka_topic_name("kafka_topic", kafka_topic)?;
+ let (stream_name, topic_name) =
self.config.topic_mapping.resolve(kafka_topic);
+ let stream_id =
Identifier::named(stream_name).map_err(BridgeError::Iggy)?;
+ let topic_id =
Identifier::named(topic_name).map_err(BridgeError::Iggy)?;
+ // No separate get_stream probe: get_topic already answers Ok(None)
when the stream
+ // itself is missing (see high_watermarks' own doc on this same fact),
so a probe first
+ // would just pay a second round trip to learn something this one call
already tells us.
+ with_request_timeout(self.client.get_topic(&stream_id,
&topic_id)).await
+ }
+
+ /// Creates the Iggy stream/topic backing `kafka_topic`, or reports that
it already exists.
+ ///
+ /// Atomic from this call's perspective, unlike a separate existence check
+ /// ([`Self::get_kafka_topic`]) followed by
[`Self::ensure_stream_and_topic`]: that sequence
+ /// has a TOCTOU window between the two calls, and
`ensure_stream_and_topic`'s own idempotent
+ /// contract would then absorb a second concurrent caller's create into a
silent `Ok`, so both
+ /// callers see success for a `CreateTopics` request Kafka promises
exactly one `NONE` for.
+ /// Here, the create attempt itself is the existence check: no separate
read precedes it, and
+ /// [`TopicCreationOutcome::AlreadyExists`] comes from the server's own
rejection of the write,
+ /// not from an earlier read that could already be stale by the time this
call's write lands.
+ ///
+ /// # Errors
+ ///
+ /// Same as [`Self::ensure_stream_and_topic`], except an already-existing
topic is reported as
+ /// [`TopicCreationOutcome::AlreadyExists`] rather than
[`BridgeError::PartitionCountMismatch`].
+ /// `CreateTopics` is not an upsert, so a pre-existing topic is never
itself an error here,
+ /// regardless of whether its partition count matches `partition_count`.
+ pub async fn create_kafka_topic(
+ &self,
+ kafka_topic: &str,
+ partition_count: u32,
+ ) -> Result<TopicCreationOutcome, BridgeError> {
+ validate_kafka_topic_name("kafka_topic", kafka_topic)?;
+ if partition_count == 0 {
+ return Err(BridgeError::InvalidPartitionCount {
+ kafka_topic: kafka_topic.to_string(),
+ });
+ }
+ let (stream_name, topic_name) =
self.config.topic_mapping.resolve(kafka_topic);
+ let stream_id = self.ensure_stream(stream_name).await?;
Review Comment:
rather than restructuring ensure_topic/create_kafka_topic to try
create_topic first and only create the stream on StreamIdNotFound, fixed it at
the shared
ensure_stream helper directly — it now attempts create_stream first and
treats StreamNameAlreadyExists as success, never calling get_stream (which
returns StreamDetails embedding the full topics: Vec<Topic>) at
all. Confirmed via iggy_common::StreamDetails's own field list that your
cost claim is accurate. This fixes the cost for both call sites
(ensure_stream_and_topic and create_kafka_topic) in one place rather than
two.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]