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]

Reply via email to