ryerraguntla commented on code in PR #4259:
URL: https://github.com/apache/iggy/pull/4259#discussion_r4077011370
##########
gateways/kafka/src/protocol/handlers/list_offsets.rs:
##########
@@ -17,44 +17,272 @@
//! `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 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_SERVER_ERROR,
ERROR_UNKNOWN_TOPIC_OR_PARTITION,
+ 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 may address through the
bridge.
+///
+/// `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
+/// (`bridge/iggy_bridge/mod.rs`'s "Concurrency ceiling"), so a large batch
head-of-line-blocks
+/// every other client's control plane for its duration. 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.
+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.
+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_INVALID_REQUEST)
+ });
+ }
+
+ 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 distinct_names: HashSet<&str> = req.topics.iter().map(|t|
t.name.as_str()).collect();
+ if distinct_names.len() > MAX_BRIDGE_BACKED_TOPICS {
+ tracing::warn!(
+ distinct_topics = distinct_names.len(),
+ max = MAX_BRIDGE_BACKED_TOPICS,
+ "ListOffsets request addresses too many distinct topics; rejecting"
+ );
+ let topics = req
+ .topics
+ .iter()
+ .map(|topic| error_topic_response(topic, ERROR_INVALID_REQUEST))
+ .collect();
+ let resp = ListOffsetsResponse::default().with_topics(topics);
+ return respond_or_close(encode_message(&resp, api_version, 256),
"ListOffsets");
+ }
+
+ let topics =
+ match tokio::time::timeout(REQUEST_DEADLINE,
resolve_all_topics(bridge, &req.topics)).await
+ {
+ Ok(topics) => topics,
+ Err(_elapsed) => {
+ tracing::warn!(
+ distinct_topics = distinct_names.len(),
+ deadline_secs = REQUEST_DEADLINE.as_secs(),
+ "ListOffsets request's aggregate bridge work exceeded its
deadline; \
+ answering retriable instead of blocking further"
+ );
+ req.topics
+ .iter()
+ .map(|topic| error_topic_response(topic,
ERROR_REQUEST_TIMED_OUT))
+ .collect()
+ }
+ };
+ let resp = ListOffsetsResponse::default().with_topics(topics);
+ respond_or_close(encode_message(&resp, api_version, 256), "ListOffsets")
+}
+
+/// Resolves every requested topic entry, deduping by name first so a topic
named in more than
+/// one request entry (or a name repeated verbatim) costs one
`high_watermarks` round trip, not
+/// one per entry - the earlier per-entry loop paid a separate round trip for
each, even for the
+/// same name.
Review Comment:
Dropped from both the doc comment and SCOPE.md:132.
--
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]