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


##########
gateways/kafka/src/protocol/handlers/join_group.rs:
##########
@@ -0,0 +1,157 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! `JoinGroup` (API key 11).
+//!
+//! The handler parks for as long as the group's join barrier takes, so it 
owns no state of its
+//! own: everything lives in `crate::group`. Subscription metadata is relayed 
verbatim to the
+//! elected leader and never decoded - which assignor a group runs is the 
client's business.
+
+use bytes::Bytes;
+use kafka_protocol::messages::join_group_response::JoinGroupResponseMember;
+use kafka_protocol::messages::{JoinGroupRequest, JoinGroupResponse};
+use kafka_protocol::protocol::StrBytes;
+
+use crate::error::Result;
+use crate::group::{JoinRequest, JoinResult};
+use crate::protocol::api::{
+    API_KEY_JOIN_GROUP, ApiVersionRange, ERROR_INVALID_REQUEST, 
ERROR_UNSUPPORTED_VERSION,
+    GatewayState, HandleOutcome, is_supported_version,
+};
+use crate::protocol::bounds_guard::validate_join_group_shape;
+use crate::protocol::handlers::{
+    decode_guarded, encode_message, respond_or_close, 
unsupported_version_response,
+};
+
+pub const RANGE: ApiVersionRange = ApiVersionRange {
+    api_key: API_KEY_JOIN_GROUP,
+    min_version: 0,
+    max_version: 9,
+};
+
+/// `protocol_name` is an `Option` in every version's struct but only nullable 
on the wire from
+/// v7. Below that a `None` encodes as a `-1` length a client reading a 
non-nullable string
+/// rejects, so an error response carries an empty name instead.
+const FIRST_NULLABLE_PROTOCOL_NAME_VERSION: i16 = 7;
+
+pub async fn handle(state: &GatewayState, api_version: i16, body: Bytes) -> 
HandleOutcome {
+    if !is_supported_version(API_KEY_JOIN_GROUP, api_version) {
+        return unsupported_version_response(API_KEY_JOIN_GROUP, api_version, 
|version| {
+            encode_error_response(version, ERROR_UNSUPPORTED_VERSION)
+        });
+    }
+    let request = match decode_guarded::<JoinGroupRequest>(api_version, body, 
|version, body| {
+        validate_join_group_shape(version, body, state.max_frame_size)
+    }) {
+        Ok(request) => request,
+        Err(error) => {
+            // debug!, not warn!: attacker-controlled, not operator-actionable.
+            tracing::debug!(%error, api_version, "Failed to decode JoinGroup 
request");
+            return respond_or_close(
+                encode_error_response(api_version, ERROR_INVALID_REQUEST),
+                "JoinGroup",
+            );
+        }
+    };
+    if let Some(reason) = request.reason.as_ref() {
+        tracing::debug!(%reason, "JoinGroup rejoin reason");
+    }
+
+    let result = state

Review Comment:
   decoded JoinGroupRequest + JoinRequest live across join().await (up to 
30min). Same pin. **Fix**: owned copy in From; drop request before await. 



##########
gateways/kafka/src/group/mod.rs:
##########
@@ -0,0 +1,399 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! In-memory coordinator for Kafka's classic consumer group protocol.
+//!
+//! [`GroupCoordinator`] owns every group this gateway instance coordinates 
and is the only
+//! module that awaits: `FindCoordinator`/`JoinGroup`/`Heartbeat`/`SyncGroup` 
handlers translate
+//! wire messages into the request types here, and `state` holds the 
synchronous state machine
+//! those requests drive.
+//!
+//! Membership is process memory, not Iggy state. Two gateway instances 
fronting one Iggy cluster
+//! therefore coordinate two independent groups under one name; see 
`docs/CONSUMER_GROUPS.md`.
+
+mod state;
+
+use std::collections::HashMap;
+use std::time::Duration;
+
+use bytes::Bytes;
+use kafka_protocol::messages::{JoinGroupRequest, SyncGroupRequest};
+use kafka_protocol::protocol::StrBytes;
+use tokio::sync::{Mutex, watch};
+use tokio::time::Instant;
+use tokio_util::sync::CancellationToken;
+
+use crate::group::state::{GroupState, Step};
+use crate::protocol::api::{ERROR_NOT_COORDINATOR, ERROR_UNKNOWN_MEMBER_ID};
+
+/// Kafka's own `group.min.session.timeout.ms` default.
+const DEFAULT_MIN_SESSION_TIMEOUT: Duration = Duration::from_secs(6);
+/// Kafka's own `group.max.session.timeout.ms` default.
+const DEFAULT_MAX_SESSION_TIMEOUT: Duration = Duration::from_mins(30);
+
+/// Kafka does not cap `rebalance.timeout.ms`, so this is a gateway resource 
bound rather than a
+/// protocol rule. It matches `DEFAULT_MAX_SESSION_TIMEOUT` because a barrier 
deadline is what
+/// bounds a park, and a larger value here would lengthen how long one request 
holds a connection
+/// and its `max_connections` permit. A client asking for more is clamped, 
never refused.
+const DEFAULT_MAX_REBALANCE_TIMEOUT: Duration = Duration::from_mins(30);
+/// Kafka's own `group.initial.rebalance.delay.ms` default.
+const DEFAULT_INITIAL_REBALANCE_DELAY: Duration = Duration::from_secs(3);
+
+/// Bounds on the state an unauthenticated client can make this gateway retain.
+///
+/// Only `max_members_per_group` has a Kafka analogue (`group.max.size`, 
unlimited by default).
+/// The rest exist because group state outlives the connection that created it 
- a member stays
+/// until its session expires, up to `max_session_timeout` - so nothing in the 
frame-level bounds
+/// guard bounds the total.
+///
+/// Worst-case retained opaque bytes is `max_total_members * 2 * 
max_member_blob_bytes` (one
+/// subscription plus one assignment per member): ~1.3 GiB at the defaults 
below.
+#[derive(Debug, Clone)]
+pub struct GroupCoordinatorConfig {
+    pub min_session_timeout: Duration,
+    pub max_session_timeout: Duration,
+    /// Ceiling on what a member's `rebalance_timeout` may contribute to a 
barrier deadline.
+    ///
+    /// A client derives this from `max.poll.interval.ms`, which no broker 
range-checks, so it is
+    /// clamped rather than rejected: a value above the ceiling is honoured up 
to it instead of
+    /// failing the join. Without a bound it would be the only limit on how 
long a parked waiter
+    /// holds its connection, since a parked member's session is refreshed 
rather than expiring.
+    pub max_rebalance_timeout: Duration,
+    /// How long a brand-new group waits for more members before completing 
its first join.
+    pub initial_rebalance_delay: Duration,
+    pub max_groups: usize,
+    pub max_members_per_group: usize,
+    /// Cap across every group, checked before a new member id is handed out.
+    pub max_total_members: usize,
+    /// Cap on the opaque bytes retained per member: the sum of one 
`JoinGroup`'s
+    /// `protocols[].metadata`, and the size of one `SyncGroup` assignment 
blob.
+    pub max_member_blob_bytes: usize,
+}
+
+impl Default for GroupCoordinatorConfig {
+    fn default() -> Self {
+        Self {
+            min_session_timeout: DEFAULT_MIN_SESSION_TIMEOUT,
+            max_session_timeout: DEFAULT_MAX_SESSION_TIMEOUT,
+            max_rebalance_timeout: DEFAULT_MAX_REBALANCE_TIMEOUT,
+            initial_rebalance_delay: DEFAULT_INITIAL_REBALANCE_DELAY,
+            max_groups: 1_000,

Review Comment:
    1000×64KiB roster = 64MiB. send_response caps i32::MAX, not 8MiB 
max_frame_size. ~125 fat members → client close → rejoin. **Fix**: bound 
product to frame. 



##########
gateways/kafka/src/group/mod.rs:
##########
@@ -0,0 +1,399 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! In-memory coordinator for Kafka's classic consumer group protocol.
+//!
+//! [`GroupCoordinator`] owns every group this gateway instance coordinates 
and is the only
+//! module that awaits: `FindCoordinator`/`JoinGroup`/`Heartbeat`/`SyncGroup` 
handlers translate
+//! wire messages into the request types here, and `state` holds the 
synchronous state machine
+//! those requests drive.
+//!
+//! Membership is process memory, not Iggy state. Two gateway instances 
fronting one Iggy cluster
+//! therefore coordinate two independent groups under one name; see 
`docs/CONSUMER_GROUPS.md`.
+
+mod state;
+
+use std::collections::HashMap;
+use std::time::Duration;
+
+use bytes::Bytes;
+use kafka_protocol::messages::{JoinGroupRequest, SyncGroupRequest};
+use kafka_protocol::protocol::StrBytes;
+use tokio::sync::{Mutex, watch};
+use tokio::time::Instant;
+use tokio_util::sync::CancellationToken;
+
+use crate::group::state::{GroupState, Step};
+use crate::protocol::api::{ERROR_NOT_COORDINATOR, ERROR_UNKNOWN_MEMBER_ID};
+
+/// Kafka's own `group.min.session.timeout.ms` default.
+const DEFAULT_MIN_SESSION_TIMEOUT: Duration = Duration::from_secs(6);
+/// Kafka's own `group.max.session.timeout.ms` default.
+const DEFAULT_MAX_SESSION_TIMEOUT: Duration = Duration::from_mins(30);
+
+/// Kafka does not cap `rebalance.timeout.ms`, so this is a gateway resource 
bound rather than a
+/// protocol rule. It matches `DEFAULT_MAX_SESSION_TIMEOUT` because a barrier 
deadline is what
+/// bounds a park, and a larger value here would lengthen how long one request 
holds a connection
+/// and its `max_connections` permit. A client asking for more is clamped, 
never refused.
+const DEFAULT_MAX_REBALANCE_TIMEOUT: Duration = Duration::from_mins(30);
+/// Kafka's own `group.initial.rebalance.delay.ms` default.
+const DEFAULT_INITIAL_REBALANCE_DELAY: Duration = Duration::from_secs(3);
+
+/// Bounds on the state an unauthenticated client can make this gateway retain.
+///
+/// Only `max_members_per_group` has a Kafka analogue (`group.max.size`, 
unlimited by default).
+/// The rest exist because group state outlives the connection that created it 
- a member stays
+/// until its session expires, up to `max_session_timeout` - so nothing in the 
frame-level bounds
+/// guard bounds the total.
+///
+/// Worst-case retained opaque bytes is `max_total_members * 2 * 
max_member_blob_bytes` (one
+/// subscription plus one assignment per member): ~1.3 GiB at the defaults 
below.
+#[derive(Debug, Clone)]
+pub struct GroupCoordinatorConfig {
+    pub min_session_timeout: Duration,
+    pub max_session_timeout: Duration,
+    /// Ceiling on what a member's `rebalance_timeout` may contribute to a 
barrier deadline.
+    ///
+    /// A client derives this from `max.poll.interval.ms`, which no broker 
range-checks, so it is
+    /// clamped rather than rejected: a value above the ceiling is honoured up 
to it instead of
+    /// failing the join. Without a bound it would be the only limit on how 
long a parked waiter
+    /// holds its connection, since a parked member's session is refreshed 
rather than expiring.
+    pub max_rebalance_timeout: Duration,
+    /// How long a brand-new group waits for more members before completing 
its first join.
+    pub initial_rebalance_delay: Duration,
+    pub max_groups: usize,
+    pub max_members_per_group: usize,
+    /// Cap across every group, checked before a new member id is handed out.
+    pub max_total_members: usize,
+    /// Cap on the opaque bytes retained per member: the sum of one 
`JoinGroup`'s
+    /// `protocols[].metadata`, and the size of one `SyncGroup` assignment 
blob.
+    pub max_member_blob_bytes: usize,
+}
+
+impl Default for GroupCoordinatorConfig {
+    fn default() -> Self {
+        Self {
+            min_session_timeout: DEFAULT_MIN_SESSION_TIMEOUT,
+            max_session_timeout: DEFAULT_MAX_SESSION_TIMEOUT,
+            max_rebalance_timeout: DEFAULT_MAX_REBALANCE_TIMEOUT,
+            initial_rebalance_delay: DEFAULT_INITIAL_REBALANCE_DELAY,
+            max_groups: 1_000,
+            max_members_per_group: 1_000,
+            max_total_members: 10_000,
+            max_member_blob_bytes: 64 * 1024,
+        }
+    }
+}
+
+/// One `JoinGroup` request, normalized across wire versions.
+#[derive(Debug, Clone)]
+pub struct JoinRequest {
+    pub group_id: StrBytes,
+    pub session_timeout: Duration,
+    pub rebalance_timeout: Duration,
+    pub member_id: StrBytes,
+    pub group_instance_id: Option<StrBytes>,
+    pub protocol_type: StrBytes,
+    /// `(name, metadata)` in request order; the order is this member's 
preference vote.
+    pub protocols: Vec<(StrBytes, Bytes)>,
+    /// KIP-394: from v4 a member must claim an id the coordinator handed it 
first.
+    pub require_known_member_id: bool,
+}
+
+impl From<(i16, &JoinGroupRequest)> for JoinRequest {
+    fn from((api_version, request): (i16, &JoinGroupRequest)) -> Self {
+        let session_timeout = millis_to_duration(request.session_timeout_ms);
+        Self {
+            group_id: request.group_id.0.clone(),
+            session_timeout,
+            // v0 has no rebalance timeout and decodes to -1. A non-positive 
value from any other
+            // version is treated the same way rather than admitting a 
zero-length rebalance
+            // window, which would expire the moment it was set.
+            rebalance_timeout: if request.rebalance_timeout_ms > 0 {
+                millis_to_duration(request.rebalance_timeout_ms)
+            } else {
+                session_timeout
+            },
+            member_id: request.member_id.clone(),
+            group_instance_id: request.group_instance_id.clone(),
+            protocol_type: request.protocol_type.clone(),
+            protocols: request
+                .protocols
+                .iter()
+                .map(|protocol| (protocol.name.clone(), 
protocol.metadata.clone()))
+                .collect(),
+            require_known_member_id: api_version >= 4,
+        }
+    }
+}
+
+/// One `SyncGroup` request, normalized across wire versions.
+#[derive(Debug, Clone)]
+pub struct SyncRequest {
+    pub group_id: StrBytes,
+    pub generation_id: i32,
+    pub member_id: StrBytes,
+    /// Present from v5 only; validated against the group when set.
+    pub protocol_type: Option<StrBytes>,
+    pub protocol_name: Option<StrBytes>,
+    pub assignments: Vec<(StrBytes, Bytes)>,
+}
+
+impl From<&SyncGroupRequest> for SyncRequest {
+    fn from(request: &SyncGroupRequest) -> Self {
+        Self {
+            group_id: request.group_id.0.clone(),
+            generation_id: request.generation_id,
+            member_id: request.member_id.clone(),
+            protocol_type: request.protocol_type.clone(),
+            protocol_name: request.protocol_name.clone(),
+            assignments: request
+                .assignments
+                .iter()
+                .map(|assignment| (assignment.member_id.clone(), 
assignment.assignment.clone()))
+                .collect(),
+        }
+    }
+}
+
+/// One member as the group leader sees it in its `JoinGroup` response.
+#[derive(Debug, Clone)]
+pub struct JoinedMember {
+    pub member_id: StrBytes,
+    pub group_instance_id: Option<StrBytes>,
+    pub metadata: Bytes,
+}
+
+/// Everything a `JoinGroup` response carries, before the handler shapes it 
for a wire version.
+#[derive(Debug, Clone)]
+pub struct JoinResult {
+    pub error: i16,
+    pub generation_id: i32,
+    pub protocol_type: Option<StrBytes>,
+    pub protocol_name: Option<StrBytes>,
+    pub leader: StrBytes,
+    /// The id the member must use next. Set on `MEMBER_ID_REQUIRED` too.
+    pub member_id: StrBytes,
+    /// Non-empty only for the leader.
+    pub members: Vec<JoinedMember>,
+}
+
+impl JoinResult {
+    #[must_use]
+    pub const fn error(error: i16, member_id: StrBytes) -> Self {
+        Self {
+            error,
+            generation_id: -1,
+            protocol_type: None,
+            protocol_name: None,
+            leader: StrBytes::new(),
+            member_id,
+            members: Vec::new(),
+        }
+    }
+}
+
+/// Everything a `SyncGroup` response carries.
+#[derive(Debug, Clone)]
+pub struct SyncResult {
+    pub error: i16,
+    pub protocol_type: Option<StrBytes>,
+    pub protocol_name: Option<StrBytes>,
+    pub assignment: Bytes,
+}
+
+impl SyncResult {
+    #[must_use]
+    pub const fn error(error: i16) -> Self {
+        Self {
+            error,
+            protocol_type: None,
+            protocol_name: None,
+            assignment: Bytes::new(),
+        }
+    }
+}
+
+/// Every consumer group this gateway instance coordinates.
+///
+/// There is no timer task. A request that touches a group first expires 
whatever is overdue in
+/// it, and a parked `JoinGroup`/`SyncGroup` waiter sleeps until that group's 
next deadline, so
+/// the coroutine waiting on a barrier is also the timer that fires it.
+pub struct GroupCoordinator {
+    config: GroupCoordinatorConfig,
+    groups: Mutex<HashMap<StrBytes, GroupState>>,
+    /// Resolves parked waiters on shutdown drain instead of holding it open 
for a full rebalance
+    /// timeout.
+    shutdown: CancellationToken,
+}
+
+impl GroupCoordinator {
+    #[must_use]
+    pub fn new(config: GroupCoordinatorConfig, shutdown: CancellationToken) -> 
Self {
+        Self {
+            config,
+            groups: Mutex::new(HashMap::new()),
+            shutdown,
+        }
+    }
+
+    /// Joins `request`'s member, parking until the group's join barrier 
completes.
+    pub async fn join(&self, request: &JoinRequest) -> JoinResult {
+        let mut parked: Option<StrBytes> = None;
+        loop {
+            let outcome = {
+                let mut groups = self.groups.lock().await;
+                let now = Instant::now();
+                let step = match parked.as_ref() {
+                    None => state::join_step(&mut groups, &self.config, 
request, now),
+                    Some(member_id) => {
+                        state::join_resume_step(&mut groups, 
&request.group_id, member_id, now)
+                    }
+                };
+                let outcome = park_outcome(&groups, &request.group_id, step, 
|member_id| {
+                    JoinResult::error(ERROR_UNKNOWN_MEMBER_ID, member_id)
+                });
+                drop(groups);
+                outcome
+            };
+            let (member_id, wake_at, receiver) = match outcome {
+                Parked::Done(result) => return result,
+                Parked::Wait(member_id, wake_at, receiver) => (member_id, 
wake_at, receiver),
+            };
+            if !self.wait_until(receiver, wake_at).await {
+                return JoinResult::error(ERROR_NOT_COORDINATOR, member_id);
+            }
+            parked = Some(member_id);
+        }
+    }
+
+    /// Delivers `request`'s member its assignment, parking a follower until 
the leader syncs.
+    pub async fn sync(&self, request: &SyncRequest) -> SyncResult {
+        let mut parked = false;
+        loop {
+            let outcome = {
+                let mut groups = self.groups.lock().await;
+                let now = Instant::now();
+                let step = if parked {
+                    state::sync_resume_step(
+                        &mut groups,
+                        &request.group_id,
+                        &request.member_id,
+                        request.generation_id,
+                        now,
+                    )
+                } else {
+                    state::sync_step(&mut groups, &self.config, request, now)
+                };
+                let outcome = park_outcome(&groups, &request.group_id, step, 
|_| {
+                    SyncResult::error(ERROR_UNKNOWN_MEMBER_ID)
+                });
+                drop(groups);
+                outcome
+            };
+            let (wake_at, receiver) = match outcome {
+                Parked::Done(result) => return result,
+                Parked::Wait(_, wake_at, receiver) => (wake_at, receiver),
+            };
+            if !self.wait_until(receiver, wake_at).await {
+                return SyncResult::error(ERROR_NOT_COORDINATOR);
+            }
+            parked = true;
+        }
+    }
+
+    /// Refreshes a member's session and reports whether it must rejoin. Never 
parks.
+    pub async fn heartbeat(
+        &self,
+        group_id: &StrBytes,
+        generation_id: i32,
+        member_id: &StrBytes,
+    ) -> i16 {
+        let mut groups = self.groups.lock().await;

Review Comment:
    Heartbeat takes process-wide Mutex + tick() every member. Cap Join 
reclaim_expired full map (state.rs:512). Park holds TCP permit vs 
max_connections 1024 / 1000 members (server.rs:266). 



##########
gateways/kafka/src/group/state.rs:
##########
@@ -0,0 +1,1923 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Synchronous state machine behind [`crate::group::GroupCoordinator`].
+//!
+//! Every entry point takes the whole group map, a config, and `now`, and 
returns either a
+//! response or a deadline to park until. Nothing here allocates a future or 
touches a clock, so
+//! the protocol rules are testable without a runtime and the coordinator's 
lock is never held
+//! across an `.await`.
+//!
+//! Kafka's `Empty` group state is "absent from the map": offsets live in 
Iggy, so an empty group
+//! holds nothing worth keeping and retaining it would be an unbounded-memory 
vector.
+
+use std::collections::{BTreeMap, HashMap};
+use std::time::Duration;
+
+use bytes::Bytes;
+use kafka_protocol::protocol::StrBytes;
+use tokio::sync::watch;
+use tokio::time::Instant;
+use uuid::Uuid;
+
+use crate::group::{
+    GroupCoordinatorConfig, JoinRequest, JoinResult, JoinedMember, 
SyncRequest, SyncResult,
+};
+use crate::protocol::api::{
+    ERROR_COORDINATOR_NOT_AVAILABLE, ERROR_GROUP_MAX_SIZE_REACHED, 
ERROR_ILLEGAL_GENERATION,
+    ERROR_INCONSISTENT_GROUP_PROTOCOL, ERROR_INVALID_GROUP_ID, 
ERROR_INVALID_REQUEST,
+    ERROR_INVALID_SESSION_TIMEOUT, ERROR_MEMBER_ID_REQUIRED, ERROR_NONE,
+    ERROR_REBALANCE_IN_PROGRESS, ERROR_UNKNOWN_MEMBER_ID,
+};
+
+/// An Iggy name caps at 255 bytes and a Kafka group's offset key is 
`kafka.cg.<group>`, so a
+/// group id this gateway admits must leave room for that prefix 
(`docs/OFFSET_STORAGE.md`).
+pub const MAX_GROUP_ID_BYTES: usize = 246;
+
+/// Prefix for a generated member id when the client sent no 
`group_instance_id`. Kafka uses the
+/// header's `client_id`, which handlers do not receive.
+const DEFAULT_MEMBER_PREFIX: &str = "member";
+
+/// Floor on how long a parked waiter sleeps. Every deadline a tick leaves 
behind is strictly in
+/// the future, so this only guards against a future rule that forgets to 
maintain that and turns
+/// a park into a spin.
+const MIN_PARK: Duration = Duration::from_millis(1);
+
+pub type Groups = HashMap<StrBytes, GroupState>;
+
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub enum Phase {
+    PreparingRebalance,
+    CompletingRebalance,
+    Stable,
+}
+
+/// Outcome of one state-machine step: answer now, or park until `wake_at`.
+pub enum Step<T> {
+    Respond(T),
+    Wait {
+        member_id: StrBytes,
+        wake_at: Instant,
+    },
+}
+
+pub struct Member {
+    group_instance_id: Option<StrBytes>,
+    session_timeout: Duration,
+    rebalance_timeout: Duration,
+    protocols: Vec<(StrBytes, Bytes)>,
+    assignment: Bytes,
+    session_deadline: Instant,
+    /// Has this member rejoined in the rebalance currently being prepared?
+    rejoined: bool,
+    /// Has this member sent `SyncGroup` in the current generation?
+    synced: bool,
+    /// Taken by the member's own parked `JoinGroup` handler. Snapshotting the 
answer at join
+    /// completion means a rebalance that starts before the waiter wakes 
cannot change what the
+    /// member is told about the generation it just completed.
+    join_response: Option<JoinResult>,
+}
+
+/// Copy what a member keeps out of the request frame.
+///
+/// `StrBytes` and `Bytes` are refcounted views, so retaining them verbatim 
keeps the whole frame
+/// alive: a member whose counted metadata is a few hundred bytes can pin 
megabytes. Copying costs
+/// one allocation per protocol on a control-plane path, and bounds what a 
member actually retains
+/// to what the caps actually measure.
+fn retained_protocols(protocols: &[(StrBytes, Bytes)]) -> Vec<(StrBytes, 
Bytes)> {
+    protocols
+        .iter()
+        .map(|(name, metadata)| {
+            (
+                StrBytes::from_string(name.as_str().to_owned()),
+                Bytes::copy_from_slice(metadata),
+            )
+        })
+        .collect()
+}
+
+impl Member {
+    fn new(request: &JoinRequest, now: Instant, max_rebalance_timeout: 
Duration) -> Self {
+        Self {
+            group_instance_id: request
+                .group_instance_id
+                .as_ref()
+                .map(|id| StrBytes::from_string(id.as_str().to_owned())),
+            session_timeout: request.session_timeout,
+            rebalance_timeout: 
request.rebalance_timeout.min(max_rebalance_timeout),
+            protocols: retained_protocols(&request.protocols),
+            assignment: Bytes::new(),
+            session_deadline: now + request.session_timeout,
+            rejoined: true,
+            synced: false,
+            join_response: None,
+        }
+    }
+
+    fn rejoin(&mut self, request: &JoinRequest, now: Instant, 
max_rebalance_timeout: Duration) {
+        self.group_instance_id = request
+            .group_instance_id
+            .as_ref()
+            .map(|id| StrBytes::from_string(id.as_str().to_owned()));
+        self.session_timeout = request.session_timeout;
+        self.rebalance_timeout = 
request.rebalance_timeout.min(max_rebalance_timeout);
+        self.protocols = retained_protocols(&request.protocols);
+        self.session_deadline = now + request.session_timeout;
+        self.rejoined = true;
+        // Cleared here rather than when a rebalance opens: a member that has 
rejoined is asking
+        // for the next generation's answer, while one still parked on the 
previous generation
+        // must keep the snapshot it is waiting to collect.
+        self.join_response = None;
+    }
+
+    fn supports(&self, name: &StrBytes) -> bool {
+        self.protocols
+            .iter()
+            .any(|(candidate, _)| candidate == name)
+    }
+
+    /// This member's vote: the first protocol it listed that the whole group 
can speak.
+    fn vote<'a>(&self, candidates: &'a [StrBytes]) -> Option<&'a StrBytes> {
+        self.protocols
+            .iter()
+            .find_map(|(name, _)| candidates.iter().find(|candidate| 
*candidate == name))
+    }
+
+    fn metadata_for(&self, name: &StrBytes) -> Bytes {
+        self.protocols
+            .iter()
+            .find(|(candidate, _)| candidate == name)
+            .map_or_else(Bytes::new, |(_, metadata)| metadata.clone())
+    }
+}
+
+pub struct GroupState {
+    phase: Phase,
+    /// A group that has completed one rebalance is at generation 1.
+    generation_id: i32,
+    protocol_type: StrBytes,
+    protocol_name: Option<StrBytes>,
+    leader: Option<StrBytes>,
+    /// Ordered, not hashed: iteration decides leader fallback and protocol 
tie-breaks.
+    members: BTreeMap<StrBytes, Member>,
+    /// Ids handed out with `MEMBER_ID_REQUIRED` that have not rejoined yet, 
and their expiry.
+    pending: BTreeMap<StrBytes, Instant>,
+    join_deadline: Option<Instant>,
+    sync_deadline: Option<Instant>,
+    rebalance_started: Instant,
+    /// First join of a new group: the barrier waits out the full delay even 
once every known
+    /// member has joined, so a second consumer starting a moment later lands 
in generation 1.
+    initial: bool,
+    changed: watch::Sender<u64>,
+}
+
+impl GroupState {
+    /// A new group's join window does not open here: it opens when the first 
member is admitted
+    /// (`admit`). A `MEMBER_ID_REQUIRED` reply creates the group but adds no 
member, and a window
+    /// opened at that moment would expire while the client was still on its 
way back with the id
+    /// it was just given.
+    fn new(protocol_type: StrBytes, now: Instant) -> Self {
+        Self {
+            phase: Phase::PreparingRebalance,
+            generation_id: 0,
+            protocol_type,
+            protocol_name: None,
+            leader: None,
+            members: BTreeMap::new(),
+            pending: BTreeMap::new(),
+            join_deadline: None,
+            sync_deadline: None,
+            rebalance_started: now,
+            initial: true,
+            changed: watch::channel(0).0,
+        }
+    }
+
+    pub fn subscribe(&self) -> watch::Receiver<u64> {
+        self.changed.subscribe()
+    }
+
+    /// `send_modify`, never `send`: the latter errors once the last receiver 
is gone, which is
+    /// the normal state of a group whose members are all between requests.
+    fn bump(&self) {
+        self.changed
+            .send_modify(|version| *version = version.wrapping_add(1));
+    }
+
+    fn is_empty(&self) -> bool {
+        self.members.is_empty() && self.pending.is_empty()
+    }
+
+    fn max_rebalance_timeout(&self) -> Duration {
+        self.members
+            .values()
+            .map(|member| member.rebalance_timeout)
+            .max()
+            .unwrap_or(Duration::ZERO)
+    }
+
+    /// Expire whatever is overdue and complete whichever phase that unblocks.
+    ///
+    /// Expiry is judged against the recorded deadline, not against when this 
happens to run, so
+    /// a request arriving after its own member's deadline finds that member 
already gone. That
+    /// is stricter than a broker, whose timer thread may not have fired yet, 
and it is what
+    /// makes eviction deterministic here.
+    fn tick(&mut self, now: Instant) {
+        let mut changed = false;
+
+        let expired_pending: Vec<StrBytes> = self
+            .pending
+            .iter()
+            .filter(|(_, deadline)| **deadline <= now)
+            .map(|(id, _)| id.clone())
+            .collect();
+        for id in &expired_pending {
+            self.pending.remove(id);
+        }
+        changed |= !expired_pending.is_empty();
+
+        let expired: Vec<StrBytes> = self
+            .members
+            .iter()
+            .filter(|(_, member)| member.session_deadline <= now)
+            .map(|(id, _)| id.clone())
+            .collect();
+        for id in &expired {
+            self.remove_member(id);
+        }
+        if !expired.is_empty() {
+            changed = true;
+            if !self.members.is_empty() && self.phase != 
Phase::PreparingRebalance {
+                self.prepare_rebalance(now, None);
+            }
+        }
+
+        if self.phase == Phase::PreparingRebalance {
+            self.maybe_complete_join(now);
+        }
+
+        if self.phase == Phase::CompletingRebalance
+            && self.sync_deadline.is_some_and(|deadline| deadline <= now)
+        {
+            let unsynced: Vec<StrBytes> = self
+                .members
+                .iter()
+                .filter(|(_, member)| !member.synced)
+                .map(|(id, _)| id.clone())
+                .collect();
+            for id in &unsynced {
+                self.remove_member(id);
+            }
+            self.sync_deadline = None;
+            if !self.members.is_empty() {
+                self.prepare_rebalance(now, None);
+            }
+            changed = true;
+        }
+
+        if changed {
+            self.bump();
+        }
+    }
+
+    fn remove_member(&mut self, member_id: &StrBytes) {
+        self.members.remove(member_id);
+        if self.leader.as_ref() == Some(member_id) {
+            self.leader = self.members.keys().next().cloned();
+        }
+    }
+
+    /// The earliest moment any rule in this group could fire.
+    fn next_deadline(&self) -> Option<Instant> {
+        let phase_deadline = match self.phase {
+            Phase::PreparingRebalance => self.join_deadline,
+            Phase::CompletingRebalance => self.sync_deadline,
+            Phase::Stable => None,
+        };
+        phase_deadline
+            .into_iter()
+            .chain(self.members.values().map(|member| member.session_deadline))
+            .chain(self.pending.values().copied())
+            .min()
+    }
+
+    fn wake_at(&self, now: Instant) -> Instant {
+        let floor = now + MIN_PARK;
+        self.next_deadline().unwrap_or(floor).max(floor)
+    }
+
+    fn prepare_rebalance(&mut self, now: Instant, trigger: Option<&StrBytes>) {
+        self.phase = Phase::PreparingRebalance;
+        self.initial = false;
+        self.rebalance_started = now;
+        self.sync_deadline = None;
+        self.join_deadline = Some(now + self.max_rebalance_timeout());
+        for (member_id, member) in &mut self.members {
+            member.rejoined = trigger == Some(member_id);
+        }
+        self.bump();
+    }
+
+    fn maybe_complete_join(&mut self, now: Instant) {
+        let all_joined = !self.initial
+            && self.pending.is_empty()
+            && !self.members.is_empty()
+            && self.members.values().all(|member| member.rejoined);
+        if all_joined || self.join_deadline.is_some_and(|deadline| now >= 
deadline) {
+            self.complete_join(now);
+        }
+    }
+
+    fn complete_join(&mut self, now: Instant) {
+        self.members.retain(|_, member| member.rejoined);
+        self.pending.clear();
+        self.join_deadline = None;
+        self.initial = false;
+        if self.members.is_empty() {
+            self.leader = None;
+            self.bump();
+            return;
+        }
+
+        let leader = match self.leader.clone() {
+            Some(leader) if self.members.contains_key(&leader) => leader,
+            _ => self.members.keys().next().cloned().unwrap_or_default(),
+        };
+        let protocol = self.select_protocol(&leader);
+        self.leader = Some(leader.clone());
+        self.protocol_name.clone_from(&protocol);
+        self.generation_id = self.generation_id.wrapping_add(1);
+        self.phase = Phase::CompletingRebalance;
+        self.sync_deadline = Some(now + self.max_rebalance_timeout());
+
+        let selected = protocol.unwrap_or_default();
+        let roster: Vec<JoinedMember> = self
+            .members
+            .iter()
+            .map(|(member_id, member)| JoinedMember {
+                member_id: member_id.clone(),
+                group_instance_id: member.group_instance_id.clone(),
+                metadata: member.metadata_for(&selected),
+            })
+            .collect();
+
+        let generation_id = self.generation_id;
+        let protocol_type = self.protocol_type.clone();
+        for (member_id, member) in &mut self.members {
+            member.session_deadline = now + member.session_timeout;
+            member.synced = false;
+            member.rejoined = false;
+            member.assignment = Bytes::new();
+            member.join_response = Some(JoinResult {
+                error: ERROR_NONE,
+                generation_id,
+                protocol_type: Some(protocol_type.clone()),
+                protocol_name: Some(selected.clone()),
+                leader: leader.clone(),
+                member_id: member_id.clone(),
+                members: if *member_id == leader {
+                    roster.clone()
+                } else {
+                    Vec::new()
+                },
+            });
+        }
+        self.bump();
+    }
+
+    /// The protocol every member speaks, most first-preference votes winning.
+    ///
+    /// Candidates are walked in the leader's own listed order, which breaks a 
tie deterministically
+    /// where Kafka breaks it by set iteration order.
+    fn select_protocol(&self, leader_id: &StrBytes) -> Option<StrBytes> {
+        let leader = self.members.get(leader_id)?;
+        let candidates: Vec<StrBytes> = leader
+            .protocols
+            .iter()
+            .map(|(name, _)| name.clone())
+            .filter(|name| self.members.values().all(|member| 
member.supports(name)))
+            .collect();
+
+        let mut best: Option<(StrBytes, usize)> = None;
+        for name in &candidates {
+            let votes = self
+                .members
+                .values()
+                .filter(|member| member.vote(&candidates) == Some(name))
+                .count();
+            if best.as_ref().is_none_or(|(_, most)| votes > *most) {
+                best = Some((name.clone(), votes));
+            }
+        }
+        best.map(|(name, _)| name)
+    }
+
+    /// Can `protocols` still leave one name every other member speaks?
+    fn protocols_compatible(&self, member_id: &StrBytes, protocols: 
&[(StrBytes, Bytes)]) -> bool {
+        protocols.iter().any(|(name, _)| {
+            self.members
+                .iter()
+                .all(|(id, member)| id == member_id || member.supports(name))
+        })
+    }
+
+    fn current_generation_result(&self, member_id: &StrBytes) -> JoinResult {
+        // Only the leader runs an assignor, so only the leader is given the 
roster. Handing it an
+        // empty one would have it assign nothing to everybody and leave the 
group consuming no
+        // partitions, silently.
+        let members = if self.leader.as_deref() == Some(member_id.as_ref()) {
+            self.roster()
+        } else {
+            Vec::new()
+        };
+        JoinResult {
+            error: ERROR_NONE,
+            generation_id: self.generation_id,
+            protocol_type: Some(self.protocol_type.clone()),
+            protocol_name: self.protocol_name.clone(),
+            leader: self.leader.clone().unwrap_or_default(),
+            member_id: member_id.clone(),
+            members,
+        }
+    }
+
+    /// The member list an assignor needs, in the group's selected protocol.
+    fn roster(&self) -> Vec<JoinedMember> {
+        let selected = self.protocol_name.clone().unwrap_or_default();
+        self.members
+            .iter()
+            .map(|(member_id, member)| JoinedMember {
+                member_id: member_id.clone(),
+                group_instance_id: member.group_instance_id.clone(),
+                metadata: member.metadata_for(&selected),
+            })
+            .collect()
+    }
+
+    fn sync_result(&self, member_id: &StrBytes) -> SyncResult {
+        SyncResult {
+            error: ERROR_NONE,
+            protocol_type: Some(self.protocol_type.clone()),
+            protocol_name: self.protocol_name.clone(),
+            assignment: self
+                .members
+                .get(member_id)
+                .map_or_else(Bytes::new, |member| member.assignment.clone()),
+        }
+    }
+
+    /// Fan the leader's blobs out to every member. A member the leader left 
out gets empty bytes,
+    /// which is what a broker stores for it too.
+    fn apply_assignments(&mut self, assignments: &[(StrBytes, Bytes)], now: 
Instant) {
+        for (member_id, member) in &mut self.members {
+            member.assignment = assignments
+                .iter()
+                .find(|(target, _)| target == member_id)
+                .map_or_else(Bytes::new, |(_, blob)| 
Bytes::copy_from_slice(blob));
+            member.session_deadline = now + member.session_timeout;
+        }
+        if let Some(leader) = self.leader.clone()
+            && let Some(member) = self.members.get_mut(&leader)
+        {
+            member.synced = true;
+        }
+        self.phase = Phase::Stable;
+        self.sync_deadline = None;
+        self.bump();
+    }
+}
+
+/// Tick every group and drop the ones that emptied, returning how many were 
reclaimed.
+///
+/// A group is only ever ticked by a request naming it, so a group whose 
consumers all died is
+/// never reclaimed on its own: it holds its slot against `max_groups` and its 
members against
+/// `max_total_members` forever. Clients that mint a fresh group id per run, 
which
+/// `kafka-console-consumer` does, then walk the gateway into permanent 
rejection. This runs only
+/// where a cap is about to reject, so the cost lands on the path that would 
otherwise wedge and
+/// never on the hot path.
+fn reclaim_expired(groups: &mut Groups, now: Instant, except: &StrBytes) -> 
usize {
+    let before = groups.len();
+    // `except` is the group the caller is mid-way through admitting. It is 
legitimately empty
+    // until its first member lands, so sweeping it here would delete the 
group out from under
+    // the request that just created it.
+    groups.retain(|id, group| {
+        if id == except {
+            return true;
+        }
+        group.tick(now);
+        !group.is_empty()
+    });
+    before - groups.len()
+}
+
+/// Expire what is overdue in `group_id`. `false` means the group is gone: 
either it never
+/// existed, or ticking emptied it.
+fn tick_group(groups: &mut Groups, group_id: &StrBytes, now: Instant) -> bool {
+    let Some(group) = groups.get_mut(group_id) else {
+        return false;
+    };
+    group.tick(now);
+    if group.is_empty() {
+        groups.remove(group_id);
+        return false;
+    }
+    true
+}
+
+/// Everything a `JoinGroup` can be rejected for before any group is touched.
+fn join_request_error(config: &GroupCoordinatorConfig, request: &JoinRequest) 
-> Option<i16> {
+    if request.group_id.is_empty() || request.group_id.len() > 
MAX_GROUP_ID_BYTES {
+        return Some(ERROR_INVALID_GROUP_ID);
+    }
+    if request.session_timeout < config.min_session_timeout
+        || request.session_timeout > config.max_session_timeout
+    {
+        return Some(ERROR_INVALID_SESSION_TIMEOUT);
+    }
+    if request.protocol_type.is_empty() || request.protocols.is_empty() {
+        return Some(ERROR_INCONSISTENT_GROUP_PROTOCOL);
+    }
+    // Names count too: they are retained alongside the metadata, and a 
request can carry many
+    // long ones while declaring almost no metadata at all.
+    let retained_bytes: usize = request
+        .protocols
+        .iter()
+        .map(|(name, metadata)| name.len() + metadata.len())
+        .sum::<usize>()
+        + request
+            .group_instance_id
+            .as_ref()
+            .map_or(0, |id| id.as_str().len());
+    if retained_bytes > config.max_member_blob_bytes {
+        return Some(ERROR_INVALID_REQUEST);
+    }
+    None
+}
+
+/// Make sure `request.group_id` exists and is ticked, returning whether this 
call created it.
+///
+/// Reclamation runs only when the group cap is about to reject: walking every 
group on every
+/// request would multiply the rebalance wakeup storm by the size of the whole 
map, while a group
+/// that emptied only matters at the moment its slot is needed.
+fn ensure_group(
+    groups: &mut Groups,
+    config: &GroupCoordinatorConfig,
+    request: &JoinRequest,
+    now: Instant,
+) -> Result<bool, i16> {
+    if groups.contains_key(&request.group_id) && tick_group(groups, 
&request.group_id, now) {
+        return Ok(false);
+    }
+    if !request.member_id.is_empty() {
+        return Err(ERROR_UNKNOWN_MEMBER_ID);
+    }
+    if groups.len() >= config.max_groups {
+        reclaim_expired(groups, now, &request.group_id);
+    }
+    if groups.len() >= config.max_groups {
+        tracing::warn!(
+            max_groups = config.max_groups,
+            "consumer group limit reached; rejecting JoinGroup"
+        );
+        return Err(ERROR_COORDINATOR_NOT_AVAILABLE);
+    }
+    groups.insert(
+        request.group_id.clone(),

Review Comment:
   group_id / protocol_type / member_id kept as StrBytes clones. kafka-protocol 
get_bytes = split_to (shares inbound Bytes). v8 reason up to default 8MiB not 
in retained_bytes. Unauth OOM: 1000 groups × frame. **Fix**: 
StrBytes::from_string / copy_from_slice at retain.



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