ryerraguntla commented on code in PR #4259:
URL: https://github.com/apache/iggy/pull/4259#discussion_r4077318037


##########
gateways/kafka/src/protocol/handlers/list_offsets.rs:
##########
@@ -17,44 +17,365 @@
 
 //! `ListOffsets` (API key 2).
 
+use std::collections::{HashMap, HashSet};
+use std::time::Duration;
+
 use bytes::Bytes;
+use kafka_protocol::messages::list_offsets_request::{ListOffsetsPartition, 
ListOffsetsTopic};
 use kafka_protocol::messages::list_offsets_response::{
     ListOffsetsPartitionResponse, ListOffsetsTopicResponse,
 };
 use kafka_protocol::messages::{ListOffsetsRequest, ListOffsetsResponse};
+use tokio::time::Instant;
 
+use crate::bridge::{BridgeError, IggyBridge};
 use crate::error::Result;
 use crate::protocol::api::{
-    API_KEY_LIST_OFFSETS, ApiVersionRange, ERROR_NOT_LEADER_OR_FOLLOWER, 
GatewayState,
-    HandleOutcome,
+    API_KEY_LIST_OFFSETS, ApiVersionRange, ERROR_INVALID_REQUEST, 
ERROR_NOT_LEADER_OR_FOLLOWER,
+    ERROR_REQUEST_TIMED_OUT, ERROR_UNKNOWN_TOPIC_OR_PARTITION,
+    ERROR_UNSUPPORTED_FOR_MESSAGE_FORMAT, ERROR_UNSUPPORTED_VERSION, 
GatewayState, HandleOutcome,
 };
 use crate::protocol::bounds_guard::validate_list_offsets_shape;
-use crate::protocol::handlers::{decode_guarded, encode_message, 
handle_versioned_request};
+use crate::protocol::handlers::{
+    decode_guarded, encode_message, handle_versioned_request, 
is_supported_version,
+    respond_or_close, unsupported_version_response,
+};
 
 pub const RANGE: ApiVersionRange = ApiVersionRange {
     api_key: API_KEY_LIST_OFFSETS,
     min_version: 1,
     max_version: 6,
 };
 
-#[expect(
-    clippy::unused_async,
-    reason = "the shared handler signature, kept until a handler awaits the 
bridge"
-)]
+/// Cap on distinct topics one `ListOffsets` request resolves through the 
bridge in one pass.
+///
+/// `bounds_guard`'s `MAX_REQUEST_ELEMENTS` (4,096) is a pre-decode `DoS` 
ceiling, not a usability
+/// recommendation: each distinct topic here costs one `high_watermarks` round 
trip against the
+/// single lockstep `IggyClient` every Kafka connection on this gateway shares 
(`README.md`'s
+/// "Concurrency ceiling"). Each call takes its own turn on that shared client 
and releases it
+/// before the next, so a large batch does not hold other connections off for 
its whole duration -
+/// only for whichever single call is in flight at a time. 100 keeps a 
worst-case batch's aggregate
+/// bridge cost small relative to that shared resource while remaining 
generous for any real
+/// consumer's offset lookup. A request naming more than this many distinct 
topics gets the first
+/// 100 resolved and the rest answered [`ERROR_REQUEST_TIMED_OUT`] with no 
bridge call at all - a
+/// client that retries only its still-erroring topics (the common case) 
narrows below the cap on
+/// its own within a couple of retries, rather than resending the same 
oversized request forever.
+const MAX_BRIDGE_BACKED_TOPICS: usize = 100;
+
+/// Wall-clock ceiling for one request's aggregate bridge work.
+///
+/// `ListOffsets` carries no `timeout_ms` field in any version this gateway 
supports (that field
+/// is v10+; [`RANGE`] tops out at v6) - unlike `CreateTopics`, there is no 
client-supplied value
+/// to honor here, so this is a fixed ceiling instead. Sized well above one 
`high_watermarks`
+/// call's own `REQUEST_TIMEOUT` (15s, bridge-internal) so a single 
slow-but-alive call is not the
+/// common trigger, while still bounding the sum across up to 
[`MAX_BRIDGE_BACKED_TOPICS`] calls -
+/// without this, a large batch against a struggling bridge could hold the 
shared client for
+/// `MAX_BRIDGE_BACKED_TOPICS * 15s`, not just one call's worth.
+///
+/// Applied per call, not once around the whole batch: [`resolve_all_topics`] 
checks it before
+/// starting each topic's `high_watermarks` call and wraps the call itself in
+/// [`tokio::time::timeout_at`] against the same instant, so a topic already 
resolved when the
+/// deadline arrives keeps its real answer and only the not-yet-started ones 
fall back to
+/// [`ERROR_REQUEST_TIMED_OUT`].
+const REQUEST_DEADLINE: Duration = Duration::from_secs(20);
+
+/// KIP-79 sentinel: the offset of the next message that would be produced.
+const LATEST_TIMESTAMP: i64 = -1;
+/// KIP-79 sentinel: the offset of the first message still retained.
+const EARLIEST_TIMESTAMP: i64 = -2;
+/// Placeholder offset/timestamp for a partition result that carries an error 
- matches real
+/// Kafka's own convention on the error path.
+const NO_OFFSET: i64 = -1;
+
+/// [`IggyBridge::high_watermarks`]'s return type, spelled once for 
[`resolve_one_partition`].
+type HighWatermarksResult =
+    core::result::Result<Vec<(u32, core::result::Result<i64, BridgeError>)>, 
BridgeError>;
+
 pub async fn handle(state: &GatewayState, api_version: i16, body: Bytes) -> 
HandleOutcome {
-    handle_versioned_request(
-        API_KEY_LIST_OFFSETS,
-        api_version,
-        body,
-        |v, b| {
-            decode_guarded::<ListOffsetsRequest>(v, b, |v, b| {
-                validate_list_offsets_shape(v, b, state.max_frame_size)
-            })
-        },
-        encode_response,
-        encode_error_response,
-        "ListOffsets",
-    )
+    let Some(bridge) = &state.bridge else {
+        return handle_versioned_request(
+            API_KEY_LIST_OFFSETS,
+            api_version,
+            body,
+            |v, b| {
+                decode_guarded::<ListOffsetsRequest>(v, b, |v, b| {
+                    validate_list_offsets_shape(v, b, state.max_frame_size)
+                })
+            },
+            encode_response,
+            encode_error_response,
+            "ListOffsets",
+        );
+    };
+
+    if !is_supported_version(API_KEY_LIST_OFFSETS, api_version) {
+        return unsupported_version_response(API_KEY_LIST_OFFSETS, api_version, 
|version| {
+            encode_error_response(version, ERROR_UNSUPPORTED_VERSION)
+        });
+    }
+
+    let req = match decode_guarded::<ListOffsetsRequest>(api_version, body, 
|v, b| {
+        validate_list_offsets_shape(v, b, state.max_frame_size)
+    }) {
+        Ok(req) => req,
+        Err(error) => {
+            // debug!, not warn!: attacker-controlled, not operator-actionable.
+            tracing::debug!(%error, "Failed to decode ListOffsets request");
+            return respond_or_close(
+                encode_error_response(api_version, ERROR_INVALID_REQUEST),
+                "ListOffsets",
+            );
+        }
+    };
+
+    let deadline = Instant::now() + REQUEST_DEADLINE;
+    let topics = resolve_all_topics(bridge, &req.topics, deadline).await;
+    let resp = ListOffsetsResponse::default().with_topics(topics);
+    respond_or_close(encode_message(&resp, api_version, 256), "ListOffsets")
+}
+
+/// One topic's bridge-lookup outcome, decided once per distinct name in 
[`resolve_all_topics`]
+/// and reused for every partition of that topic in [`resolve_one_partition`].
+enum TopicLookup {
+    /// A real `high_watermarks` call was made and returned.
+    Watermarks(HighWatermarksResult),
+    /// No partition of this topic asked for [`LATEST_TIMESTAMP`] or 
[`EARLIEST_TIMESTAMP`], so no
+    /// call was made at all - every partition here answers
+    /// [`ERROR_UNSUPPORTED_FOR_MESSAGE_FORMAT`] regardless of what a lookup 
would have returned.
+    NoLookupNeeded,
+    /// Beyond [`MAX_BRIDGE_BACKED_TOPICS`], or the deadline elapsed before 
this topic's turn - no
+    /// call was made. Answered [`ERROR_REQUEST_TIMED_OUT`] (retriable) rather 
than
+    /// [`ERROR_INVALID_REQUEST`] so a client's own per-topic retry narrows 
the batch on its own.
+    NotAttempted,
+}
+
+/// Dedupes `requested` by topic name, merging every entry's partitions and 
noting - per name -
+/// whether any partition asked for [`LATEST_TIMESTAMP`] or 
[`EARLIEST_TIMESTAMP`].
+///
+/// `order` preserves first-seen order so the topic cap in 
[`resolve_topic_lookups`] keeps a
+/// deterministic prefix of the request rather than an arbitrary hash-order 
subset.
+fn group_requested_topics(
+    requested: &[ListOffsetsTopic],
+) -> (Vec<&str>, HashMap<&str, Vec<u32>>, HashSet<&str>) {
+    let mut order: Vec<&str> = Vec::new();
+    let mut partitions_by_name: HashMap<&str, Vec<u32>> = HashMap::new();
+    let mut needs_lookup: HashSet<&str> = HashSet::new();
+    for topic in requested {
+        let name = topic.name.as_str();
+        if !partitions_by_name.contains_key(name) {
+            order.push(name);
+        }
+        let entry = partitions_by_name.entry(name).or_default();
+        for p in &topic.partitions {
+            if let Ok(index) = u32::try_from(p.partition_index) {
+                entry.push(index);
+            }
+            if matches!(p.timestamp, LATEST_TIMESTAMP | EARLIEST_TIMESTAMP) {
+                needs_lookup.insert(name);
+            }
+        }
+    }
+    for partitions in partitions_by_name.values_mut() {
+        partitions.sort_unstable();
+        partitions.dedup();
+    }
+    (order, partitions_by_name, needs_lookup)
+}
+
+/// Resolves one [`TopicLookup`] per name in `order`: 
[`TopicLookup::NotAttempted`] beyond
+/// [`MAX_BRIDGE_BACKED_TOPICS`] or once `deadline` has passed, 
[`TopicLookup::NoLookupNeeded`]
+/// when `needs_lookup` excludes the name, otherwise a real `high_watermarks` 
call wrapped in
+/// [`tokio::time::timeout_at`] against `deadline` so a topic already resolved 
when time runs out
+/// keeps its real answer.
+async fn resolve_topic_lookups<'a>(
+    bridge: &IggyBridge,
+    order: &[&'a str],
+    partitions_by_name: &HashMap<&'a str, Vec<u32>>,
+    needs_lookup: &HashSet<&'a str>,
+    deadline: Instant,
+) -> HashMap<&'a str, TopicLookup> {
+    let accepted: HashSet<&str> = order
+        .iter()
+        .take(MAX_BRIDGE_BACKED_TOPICS)
+        .copied()
+        .collect();
+    if order.len() > MAX_BRIDGE_BACKED_TOPICS {
+        // debug!, not warn!: the client controls how many topics it batches 
into one request and
+        // the connection stays open, so a consumer stuck above the cap logs 
this every retry.
+        tracing::debug!(
+            distinct_topics = order.len(),
+            max = MAX_BRIDGE_BACKED_TOPICS,
+            "ListOffsets request exceeds the per-request topic cap; resolving 
the first {} and \
+             answering the rest retriable",
+            MAX_BRIDGE_BACKED_TOPICS
+        );
+    }
+
+    let mut lookups: HashMap<&str, TopicLookup> = HashMap::new();
+    let mut deadline_exceeded = false;
+    for &name in order {
+        if !accepted.contains(name) {
+            lookups.insert(name, TopicLookup::NotAttempted);
+            continue;
+        }
+        if !needs_lookup.contains(name) {
+            lookups.insert(name, TopicLookup::NoLookupNeeded);
+            continue;
+        }
+        if deadline_exceeded || Instant::now() >= deadline {
+            if !deadline_exceeded {
+                deadline_exceeded = true;
+                tracing::warn!(
+                    deadline_secs = REQUEST_DEADLINE.as_secs(),
+                    "ListOffsets request's aggregate bridge work exceeded its 
deadline; \
+                     answering remaining topics retriable instead of starting 
new bridge calls"
+                );
+            }
+            lookups.insert(name, TopicLookup::NotAttempted);
+            continue;
+        }
+
+        let partitions = partitions_by_name[name].as_slice();
+        let result =
+            match tokio::time::timeout_at(deadline, 
bridge.high_watermarks(name, partitions)).await
+            {
+                Ok(result) => result,
+                Err(_elapsed) => {
+                    deadline_exceeded = true;
+                    tracing::warn!(
+                        topic = name,
+                        deadline_secs = REQUEST_DEADLINE.as_secs(),
+                        "ListOffsets bridge call for this topic exceeded the 
request's aggregate \
+                     deadline; answering retriable instead of blocking further"
+                    );
+                    lookups.insert(name, TopicLookup::NotAttempted);
+                    continue;
+                }
+            };
+
+        if let Err(call_err) = &result {
+            let kafka_code = call_err.to_kafka_error_code();
+            if kafka_code == ERROR_UNKNOWN_TOPIC_OR_PARTITION {
+                // Client-caused (topic doesn't exist / isn't mapped): 
expected traffic, not
+                // operator-actionable.
+                tracing::debug!(topic = name, %call_err, "ListOffsets bridge 
lookup: topic not found");
+            } else {
+                tracing::error!(topic = name, %call_err, "ListOffsets bridge 
lookup failed");
+            }
+        }
+        lookups.insert(name, TopicLookup::Watermarks(result));
+    }
+    lookups
+}
+
+/// Resolves every requested topic entry via [`group_requested_topics`] +
+/// [`resolve_topic_lookups`], then stamps each requested partition with its 
topic's
+/// [`TopicLookup`] outcome. `EARLIEST`/`LATEST` are the only timestamps this 
bridge resolves -
+/// Iggy exposes no per-message timestamp index - so every other requested 
timestamp gets
+/// [`ERROR_UNSUPPORTED_FOR_MESSAGE_FORMAT`] rather than a fabricated offset.
+async fn resolve_all_topics(
+    bridge: &IggyBridge,
+    requested: &[ListOffsetsTopic],
+    deadline: Instant,
+) -> Vec<ListOffsetsTopicResponse> {
+    let (order, partitions_by_name, needs_lookup) = 
group_requested_topics(requested);
+    let lookups =
+        resolve_topic_lookups(bridge, &order, &partitions_by_name, 
&needs_lookup, deadline).await;
+
+    requested
+        .iter()
+        .map(|topic| {
+            // Always present: `order` (and so `lookups`) was built from 
exactly these same
+            // requested topic names, just above.
+            let lookup = lookups
+                .get(topic.name.as_str())
+                .expect("every requested topic name was resolved above");
+            let partitions = topic
+                .partitions
+                .iter()
+                .map(|requested| resolve_one_partition(requested, lookup))
+                .collect();
+            ListOffsetsTopicResponse::default()
+                .with_name(topic.name.clone())
+                .with_partitions(partitions)
+        })
+        .collect()
+}
+
+/// `lookup` is the whole topic's resolution outcome: an errored 
[`TopicLookup::Watermarks`] is a
+/// call-level failure (e.g. the mapped stream doesn't exist) applying to 
every partition alike;
+/// the inner per-partition `Result` inside its `Ok` is 
[`BridgeError::PartitionOutOfRange`] for one
+/// bad index among otherwise resolvable ones.
+fn resolve_one_partition(
+    requested: &ListOffsetsPartition,
+    lookup: &TopicLookup,
+) -> ListOffsetsPartitionResponse {
+    let Ok(partition_index) = u32::try_from(requested.partition_index) else {
+        return error_response(requested.partition_index, 
ERROR_UNKNOWN_TOPIC_OR_PARTITION);
+    };
+
+    let results = match lookup {
+        TopicLookup::NotAttempted => {
+            return error_response(requested.partition_index, 
ERROR_REQUEST_TIMED_OUT);
+        }
+        // By construction (`needs_lookup` in `resolve_all_topics`) every 
partition of a topic in
+        // this state has a non-sentinel timestamp, so there is nothing to 
branch on here.
+        TopicLookup::NoLookupNeeded => {

Review Comment:
   Nonexistent topics now correctly report UNKNOWN_TOPIC_OR_PARTITION 
regardless of what timestamp was asked. 



-- 
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