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


##########
gateways/kafka/src/protocol/handlers/find_coordinator.rs:
##########
@@ -0,0 +1,164 @@
+// 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.
+
+//! `FindCoordinator` (API key 10).
+//!
+//! The answer is always this gateway, matching the single broker entry 
Metadata advertises. No
+//! group state is consulted: where the coordinator lives does not depend on 
which group is asked
+//! about.
+
+use bytes::Bytes;
+use kafka_protocol::messages::find_coordinator_response::Coordinator;
+use kafka_protocol::messages::{BrokerId, FindCoordinatorRequest, 
FindCoordinatorResponse};
+use kafka_protocol::protocol::StrBytes;
+
+use crate::error::Result;
+use crate::protocol::api::{
+    API_KEY_FIND_COORDINATOR, ApiVersionRange, BrokerAdvertise, 
ERROR_INVALID_REQUEST, ERROR_NONE,
+    ERROR_UNSUPPORTED_VERSION, GatewayState, HandleOutcome, 
is_supported_version,
+};
+use crate::protocol::bounds_guard::validate_find_coordinator_shape;
+use crate::protocol::handlers::{
+    decode_guarded, encode_message, respond_or_close, 
unsupported_version_response,
+};
+
+pub const RANGE: ApiVersionRange = ApiVersionRange {
+    api_key: API_KEY_FIND_COORDINATOR,
+    min_version: 0,
+    max_version: 4,
+};
+
+/// `key_type` 0. Types 1 (transaction) and 2 (share) have no coordinator here.
+const COORDINATOR_TYPE_GROUP: i8 = 0;
+
+/// The node id this gateway advertises for itself, in Metadata and here alike.
+const SELF_NODE_ID: i32 = 1;
+
+const UNSUPPORTED_KEY_TYPE_MESSAGE: &str = "only group coordination is 
supported";
+
+#[expect(
+    clippy::unused_async,
+    reason = "the shared handler signature, kept until a handler awaits the 
bridge"
+)]
+pub async fn handle(state: &GatewayState, api_version: i16, body: Bytes) -> 
HandleOutcome {
+    if !is_supported_version(API_KEY_FIND_COORDINATOR, api_version) {
+        return unsupported_version_response(API_KEY_FIND_COORDINATOR, 
api_version, |version| {
+            encode_error_response(version, ERROR_UNSUPPORTED_VERSION)
+        });
+    }
+    match decode_guarded::<FindCoordinatorRequest>(api_version, body, 
|version, body| {
+        validate_find_coordinator_shape(version, body, state.max_frame_size)
+    }) {
+        Ok(request) => respond_or_close(
+            encode_response(api_version, &request, &state.broker),
+            "FindCoordinator",
+        ),
+        Err(error) => {
+            // debug!, not warn!: attacker-controlled, not operator-actionable.
+            tracing::debug!(%error, api_version, "Failed to decode 
FindCoordinator request");
+            if api_version >= 4 {
+                // v4 carries its error per requested key and the keys are 
exactly what failed to
+                // decode, so there is no honest body to send.
+                return HandleOutcome::Close;
+            }
+            respond_or_close(
+                encode_error_response(api_version, ERROR_INVALID_REQUEST),
+                "FindCoordinator",
+            )
+        }
+    }
+}
+
+/// # Errors
+///
+/// Returns an error when `kafka_protocol` cannot encode the response at 
`version`.
+pub fn encode_response(
+    version: i16,
+    request: &FindCoordinatorRequest,
+    broker: &BrokerAdvertise,
+) -> Result<Bytes> {
+    let keys = if version >= 4 {
+        request.coordinator_keys.clone()
+    } else {
+        vec![request.key.clone()]
+    };
+    if request.key_type == COORDINATOR_TYPE_GROUP {
+        encode_inner(version, &keys, ERROR_NONE, None, Some(broker))
+    } else {
+        // Not a retriable code: transactions are out of scope for good, and
+        // COORDINATOR_NOT_AVAILABLE would make a transactional producer retry 
forever.
+        encode_inner(
+            version,
+            &keys,
+            ERROR_INVALID_REQUEST,

Review Comment:
    key_type!=0 → 42. librdkafka rd_kafka_txn_handle_FindCoordinator default → 
txn_coord_set(NULL) → 500ms forever. Java fatal. File not in 4270-vs-4263 
delta; PR vs master still ships. Fix: 53 or 31. (same as in #4263)



##########
gateways/kafka/docs/MANUAL_TESTING.md:
##########
@@ -175,8 +180,13 @@ Requires `kcat` installed. Gateway does **not** implement 
SASL or full broker se
 | ID | Test | Command | Expected (foundation) |
 | ---- | ------ | --------- | --------------------- |
 | G1 | Broker metadata | `kcat -b 127.0.0.1:9093 -L` | ApiVersions + Metadata 
handshake; broker appears in metadata |
-| G2 | Produce (likely fails later) | `echo "hello" \| kcat -b 127.0.0.1:9093 
-t test -P` | May fail at coordinator/group stage — document actual error |
-| G3 | Consumer (likely fails later) | `kcat -b 127.0.0.1:9093 -t test -C -o 
beginning` | May fail without consumer groups — document actual error |
+| G2 | Produce (likely fails later) | `echo "hello" \| kcat -b 127.0.0.1:9093 
-t test -P` | Produce is still a stub: retriable `NOT_LEADER_OR_FOLLOWER` (6), 
so kcat retries — document actual error |
+| G3 | Consumer group rebalance | `kcat -b 127.0.0.1:9093 -G g1 test` in two 
terminals | Each prints its assigned partitions and the two sets are disjoint; 
then both stall, because OffsetFetch (9) is unlisted and closes the connection 
— the client re-runs FindCoordinator and loops. Record the exact librdkafka log 
lines |
+| G4 | Ungraceful consumer exit | `kill -9` one of G3's kcats | Within 
`session.timeout.ms` the survivor logs a rebalance and is assigned every 
partition |
+| G5 | Java console consumer | `kafka-console-consumer.sh --bootstrap-server 
127.0.0.1:9093 --group g2 --topic test` | Exercises JoinGroup v9, SyncGroup v5, 
Heartbeat v4, FindCoordinator v4. A 4.0 client may need `--consumer-property 
group.protocol=classic`, or it sends ConsumerGroupHeartbeat (68) and the 
connection closes |
+| G6 | Graceful kcat exit | G3's two kcats, then Ctrl-C one (kcat closes its 
consumer, which sends LeaveGroup v0/v1) | The survivor rebalances within one 
heartbeat interval, not `session.timeout.ms`, and is assigned every partition |

Review Comment:
   G6/G7 — “assigned partitions”. Metadata stub → RangeAssignor 0. OffsetFetch 
honesty OK. **Fix**: join/leave/sync real; assignment empty. 



##########
gateways/kafka/docs/kafka_api_keys_reference.md:
##########


Review Comment:
   “4.0 default key 68” false. group.protocol default classic.



##########
gateways/kafka/src/group/state.rs:
##########
@@ -0,0 +1,2698 @@
+// 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, HashSet};
+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, 
LeaveRequest, LeaveResult,
+    LeavingMember, LeftMember, SyncRequest, SyncResult,
+};
+use crate::protocol::api::{
+    ERROR_COORDINATOR_NOT_AVAILABLE, ERROR_FENCED_INSTANCE_ID, 
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,
+    },
+}
+
+/// What one `LeaveGroup` removed, which decides how the group reacts 
afterwards.
+#[derive(Default)]
+struct Departures {
+    members: bool,
+    pending: bool,
+}
+
+/// Member ids per `group_instance_id`. Without KIP-345 replacement one 
instance id can be held by
+/// several members, so each maps to a set.
+type InstanceHolders = HashMap<StrBytes, HashSet<StrBytes>>;
+
+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();
+        }
+    }
+
+    fn instance_holders(&self) -> InstanceHolders {
+        let mut holders = InstanceHolders::new();
+        for (member_id, member) in &self.members {
+            if let Some(instance_id) = &member.group_instance_id {
+                holders
+                    .entry(instance_id.clone())
+                    .or_default()
+                    .insert(member_id.clone());
+            }
+        }
+        holders
+    }
+
+    /// Removes one `LeaveGroup` identity, following Kafka's 
`handleLeaveGroup` order.
+    ///
+    /// `holders` must stay exact across calls: a later identity in the same 
request is resolved
+    /// against it.
+    fn leave_one(
+        &mut self,
+        identity: &LeavingMember,
+        holders: &mut InstanceHolders,
+        departed: &mut Departures,
+    ) -> i16 {
+        let holders_of = identity
+            .group_instance_id
+            .as_ref()
+            .and_then(|instance_id| holders.get_mut(instance_id))
+            .filter(|ids| !ids.is_empty());
+
+        if identity.member_id.is_empty() {
+            let Some(ids) = holders_of else {
+                return ERROR_UNKNOWN_MEMBER_ID;
+            };
+            for id in std::mem::take(ids) {
+                self.remove_member(&id);
+            }
+            departed.members = true;
+            return ERROR_NONE;
+        }
+        if self.pending.remove(&identity.member_id).is_some() {
+            departed.pending = true;
+            return ERROR_NONE;
+        }
+        match (identity.group_instance_id.is_some(), holders_of) {
+            (_, Some(ids)) if ids.contains(&identity.member_id) => {}
+            (_, Some(_)) => return ERROR_FENCED_INSTANCE_ID,
+            (true, None) => return ERROR_UNKNOWN_MEMBER_ID,
+            (false, None) if !self.members.contains_key(&identity.member_id) 
=> {
+                return ERROR_UNKNOWN_MEMBER_ID;
+            }
+            (false, None) => {}
+        }
+        if let Some(ids) = self
+            .members
+            .get(&identity.member_id)
+            .and_then(|member| member.group_instance_id.as_ref())
+            .and_then(|instance_id| holders.get_mut(instance_id))
+        {
+            ids.remove(&identity.member_id);
+        }
+        self.remove_member(&identity.member_id);
+        departed.members = true;
+        ERROR_NONE
+    }
+
+    /// React to a `LeaveGroup` the way the session sweep reacts to an expiry.
+    ///
+    /// The rebalance opens only after the removals, so its join window is 
sized from the
+    /// members that remain. The final `bump` is unconditional on any removal: 
it is what answers
+    /// a waiter the departed member still has parked on another connection.
+    fn after_departure(&mut self, departed: &Departures, now: Instant) {
+        if departed.members && !self.members.is_empty() && self.phase != 
Phase::PreparingRebalance {

Review Comment:
    last-member leave while PreparingRebalance skips prepare_rebalance; 
join_deadline stays. Later tick → complete_join pending.clear() (:441) → group 
drop. Pending rejoin 25. Kafka initNextGeneration keeps pendingJoinMembers. 
Test :2425 same Instant only. KIP-394 empty rejoin recovers; no durable loss. 
**Fix**: join_deadline=None when members.empty && pending nonempty. 



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