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]
