use crate::group::{Assignment, Subscription, TopicPartition};
pub mod codes {
pub const NONE: i16 = 0;
pub const COORDINATOR_LOAD_IN_PROGRESS: i16 = 14;
pub const COORDINATOR_NOT_AVAILABLE: i16 = 15;
pub const NOT_COORDINATOR: i16 = 16;
pub const ILLEGAL_GENERATION: i16 = 22;
pub const UNKNOWN_MEMBER_ID: i16 = 25;
pub const REBALANCE_IN_PROGRESS: i16 = 27;
pub const FENCED_INSTANCE_ID: i16 = 82;
pub const MEMBER_ID_REQUIRED: i16 = 79;
}
pub const NO_GENERATION: i32 = -1;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum RebalanceProtocol {
#[default]
Eager,
Cooperative,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MemberState {
Unjoined,
Joining,
Syncing,
Stable,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Step {
Join { member_id: String },
AssignAndSync { members: Vec<Subscription> },
Sync,
Heartbeat,
FindCoordinator,
}
#[derive(Debug, Clone)]
pub struct GroupMember {
group_id: String,
member_id: String,
generation: i32,
state: MemberState,
topics: Vec<String>,
assignment: Vec<TopicPartition>,
leader: bool,
protocol: RebalanceProtocol,
lost: Vec<TopicPartition>,
pending_members: Vec<Subscription>,
}
impl GroupMember {
#[must_use]
pub fn new(group_id: impl Into<String>, topics: Vec<String>) -> Self {
Self {
group_id: group_id.into(),
member_id: String::new(),
generation: NO_GENERATION,
state: MemberState::Unjoined,
topics,
assignment: Vec::new(),
leader: false,
protocol: RebalanceProtocol::Eager,
lost: Vec::new(),
pending_members: Vec::new(),
}
}
#[must_use]
pub fn with_protocol(mut self, protocol: RebalanceProtocol) -> Self {
self.protocol = protocol;
self
}
#[must_use]
pub fn protocol(&self) -> RebalanceProtocol {
self.protocol
}
#[must_use]
pub fn lost(&self) -> &[TopicPartition] {
&self.lost
}
#[must_use]
pub fn group_id(&self) -> &str {
&self.group_id
}
#[must_use]
pub fn topics(&self) -> &[String] {
&self.topics
}
#[must_use]
pub fn state(&self) -> MemberState {
self.state
}
#[must_use]
pub fn generation(&self) -> i32 {
self.generation
}
#[must_use]
pub fn member_id(&self) -> &str {
&self.member_id
}
#[must_use]
pub fn is_leader(&self) -> bool {
self.leader
}
#[must_use]
pub fn assignment(&self) -> &[TopicPartition] {
&self.assignment
}
#[must_use]
pub fn can_commit(&self) -> bool {
self.state == MemberState::Stable && self.generation != NO_GENERATION
}
#[must_use]
pub fn step(&self) -> Step {
match self.state {
MemberState::Unjoined | MemberState::Joining => Step::Join {
member_id: self.member_id.clone(),
},
MemberState::Syncing if self.leader => Step::AssignAndSync {
members: self.pending_members.clone(),
},
MemberState::Syncing => Step::Sync,
MemberState::Stable => Step::Heartbeat,
}
}
#[must_use]
pub fn subscription(&self) -> Subscription {
Subscription {
member_id: self.member_id.clone(),
topics: self.topics.clone(),
owned: self.assignment.clone(),
generation: self.generation,
}
}
pub fn set_topics(&mut self, topics: Vec<String>) {
if topics != self.topics {
self.topics = topics;
self.revoke_and_rejoin();
}
}
pub fn request_rejoin(&mut self) {
if self.state != MemberState::Unjoined {
self.revoke_and_rejoin();
}
}
pub fn on_join(
&mut self,
error: i16,
generation: i32,
member_id: &str,
leader_id: &str,
members: Vec<Subscription>,
) -> Step {
match error {
codes::NONE => {
self.member_id = member_id.to_owned();
self.generation = generation;
self.leader = leader_id == member_id;
self.state = MemberState::Syncing;
if self.leader {
self.pending_members = members.clone();
Step::AssignAndSync { members }
} else {
self.pending_members.clear();
Step::Sync
}
}
codes::MEMBER_ID_REQUIRED => {
self.member_id = member_id.to_owned();
self.state = MemberState::Unjoined;
Step::Join {
member_id: self.member_id.clone(),
}
}
codes::UNKNOWN_MEMBER_ID | codes::FENCED_INSTANCE_ID => {
self.member_id = String::new();
self.revoke_and_rejoin();
Step::Join {
member_id: String::new(),
}
}
codes::COORDINATOR_NOT_AVAILABLE
| codes::NOT_COORDINATOR
| codes::COORDINATOR_LOAD_IN_PROGRESS => {
self.state = MemberState::Unjoined;
Step::FindCoordinator
}
_ => {
self.revoke_and_rejoin();
Step::Join {
member_id: self.member_id.clone(),
}
}
}
}
pub fn on_sync(&mut self, error: i16, assigned: Vec<TopicPartition>) -> Step {
match error {
codes::NONE => {
self.lost.clear();
if self.protocol == RebalanceProtocol::Cooperative {
let mut assigned_sorted = assigned.clone();
assigned_sorted.sort();
self.lost = self
.assignment
.iter()
.filter(|tp| !assigned_sorted.contains(tp))
.cloned()
.collect();
}
self.assignment = assigned;
self.assignment.sort();
self.state = MemberState::Stable;
if !self.lost.is_empty() {
self.state = MemberState::Unjoined;
self.leader = false;
self.pending_members.clear();
return Step::Join {
member_id: self.member_id.clone(),
};
}
Step::Heartbeat
}
codes::REBALANCE_IN_PROGRESS | codes::ILLEGAL_GENERATION => {
self.revoke_and_rejoin();
Step::Join {
member_id: self.member_id.clone(),
}
}
codes::UNKNOWN_MEMBER_ID | codes::FENCED_INSTANCE_ID => {
self.member_id = String::new();
self.revoke_and_rejoin();
Step::Join {
member_id: String::new(),
}
}
codes::COORDINATOR_NOT_AVAILABLE | codes::NOT_COORDINATOR => {
self.revoke_and_rejoin();
Step::FindCoordinator
}
_ => {
self.revoke_and_rejoin();
Step::Join {
member_id: self.member_id.clone(),
}
}
}
}
pub fn on_heartbeat(&mut self, error: i16) -> Step {
match error {
codes::NONE => Step::Heartbeat,
codes::REBALANCE_IN_PROGRESS | codes::ILLEGAL_GENERATION => {
self.revoke_and_rejoin();
Step::Join {
member_id: self.member_id.clone(),
}
}
codes::UNKNOWN_MEMBER_ID | codes::FENCED_INSTANCE_ID => {
self.member_id = String::new();
self.revoke_and_rejoin();
Step::Join {
member_id: String::new(),
}
}
codes::COORDINATOR_NOT_AVAILABLE | codes::NOT_COORDINATOR => {
self.revoke_and_rejoin();
Step::FindCoordinator
}
_ => {
self.revoke_and_rejoin();
Step::Join {
member_id: self.member_id.clone(),
}
}
}
}
pub fn on_leave(&mut self) {
self.member_id = String::new();
self.revoke_and_rejoin();
}
#[must_use]
pub fn my_share(&self, assignment: &Assignment) -> Vec<TopicPartition> {
assignment.get(&self.member_id).cloned().unwrap_or_default()
}
fn revoke_and_rejoin(&mut self) {
if self.protocol == RebalanceProtocol::Eager {
self.lost = std::mem::take(&mut self.assignment);
} else {
self.lost.clear();
}
self.pending_members.clear();
self.leader = false;
self.state = MemberState::Unjoined;
}
}
#[cfg(test)]
mod tests {
use super::*;
fn member() -> GroupMember {
GroupMember::new("g", vec!["t".to_owned()])
}
fn tp(partition: i32) -> TopicPartition {
TopicPartition::new("t", partition)
}
#[test]
fn a_first_join_is_refused_and_then_accepted() {
let mut m = member();
assert_eq!(
m.step(),
Step::Join {
member_id: String::new()
}
);
let step = m.on_join(codes::MEMBER_ID_REQUIRED, NO_GENERATION, "m-1", "", vec![]);
assert_eq!(
step,
Step::Join {
member_id: "m-1".to_owned()
}
);
assert_eq!(m.member_id(), "m-1");
assert_eq!(m.state(), MemberState::Unjoined);
let step = m.on_join(codes::NONE, 7, "m-1", "m-2", vec![]);
assert_eq!(step, Step::Sync, "a follower syncs with no assignment");
assert_eq!(m.state(), MemberState::Syncing);
assert!(!m.is_leader());
assert_eq!(m.on_sync(codes::NONE, vec![tp(0), tp(1)]), Step::Heartbeat);
assert_eq!(m.state(), MemberState::Stable);
assert_eq!(m.assignment(), &[tp(0), tp(1)]);
assert_eq!(m.generation(), 7);
}
#[test]
fn the_leader_is_asked_to_assign() {
let mut m = member();
let members = vec![Subscription {
member_id: "m-1".to_owned(),
topics: vec!["t".to_owned()],
owned: vec![],
generation: 1,
}];
let step = m.on_join(codes::NONE, 3, "m-1", "m-1", members.clone());
assert_eq!(step, Step::AssignAndSync { members });
assert!(m.is_leader());
}
#[test]
fn the_leader_is_still_the_leader_on_the_next_step() {
let mut m = member();
let members = vec![Subscription {
member_id: "m-1".to_owned(),
topics: vec!["t".to_owned()],
owned: vec![],
generation: 1,
}];
m.on_join(codes::NONE, 3, "m-1", "m-1", members.clone());
assert_eq!(m.step(), Step::AssignAndSync { members });
}
#[test]
fn a_follower_is_still_a_follower_on_the_next_step() {
let mut m = member();
m.on_join(codes::NONE, 3, "m-1", "m-2", vec![]);
assert_eq!(m.step(), Step::Sync);
}
#[test]
fn commits_are_refused_outside_stable() {
let mut m = member();
assert!(!m.can_commit(), "not a member yet");
m.on_join(codes::NONE, 1, "m-1", "m-1", vec![]);
assert!(!m.can_commit(), "syncing is not stable");
m.on_sync(codes::NONE, vec![tp(0)]);
assert!(m.can_commit());
m.on_heartbeat(codes::REBALANCE_IN_PROGRESS);
assert!(!m.can_commit(), "a rebalance suspends commits");
}
#[test]
fn a_rebalance_revokes_the_assignment_immediately() {
let mut m = member();
m.on_join(codes::NONE, 1, "m-1", "m-2", vec![]);
m.on_sync(codes::NONE, vec![tp(0), tp(1)]);
assert_eq!(m.assignment().len(), 2);
let step = m.on_heartbeat(codes::REBALANCE_IN_PROGRESS);
assert_eq!(
step,
Step::Join {
member_id: "m-1".to_owned()
}
);
assert!(m.assignment().is_empty(), "partitions must be given up");
assert_eq!(m.state(), MemberState::Unjoined);
assert_eq!(m.generation(), 1);
}
#[test]
fn an_illegal_generation_revokes_too() {
let mut m = member();
m.on_join(codes::NONE, 4, "m-1", "m-2", vec![]);
m.on_sync(codes::NONE, vec![tp(0)]);
m.on_heartbeat(codes::ILLEGAL_GENERATION);
assert!(m.assignment().is_empty());
assert!(!m.can_commit());
assert_eq!(
m.member_id(),
"m-1",
"the member id survives a generation bump"
);
}
#[test]
fn an_unknown_member_id_forgets_the_identity() {
let mut m = member();
m.on_join(codes::NONE, 4, "m-1", "m-2", vec![]);
m.on_sync(codes::NONE, vec![tp(0)]);
let step = m.on_heartbeat(codes::UNKNOWN_MEMBER_ID);
assert_eq!(
step,
Step::Join {
member_id: String::new()
}
);
assert_eq!(m.member_id(), "", "the id is no longer ours to use");
assert!(m.assignment().is_empty());
}
#[test]
fn a_lost_coordinator_is_rediscovered() {
let mut m = member();
m.on_join(codes::NONE, 2, "m-1", "m-2", vec![]);
m.on_sync(codes::NONE, vec![tp(0)]);
assert_eq!(
m.on_heartbeat(codes::NOT_COORDINATOR),
Step::FindCoordinator
);
assert!(
m.assignment().is_empty(),
"still revoked: we cannot heartbeat"
);
}
#[test]
fn a_rebalance_during_sync_restarts_the_join() {
let mut m = member();
m.on_join(codes::NONE, 5, "m-1", "m-1", vec![]);
assert_eq!(m.state(), MemberState::Syncing);
let step = m.on_sync(codes::REBALANCE_IN_PROGRESS, vec![]);
assert_eq!(
step,
Step::Join {
member_id: "m-1".to_owned()
}
);
assert!(
!m.is_leader(),
"leadership is not carried across a rebalance"
);
}
#[test]
fn changing_topics_forces_a_rejoin() {
let mut m = member();
m.on_join(codes::NONE, 1, "m-1", "m-2", vec![]);
m.on_sync(codes::NONE, vec![tp(0)]);
m.set_topics(vec!["t".to_owned(), "u".to_owned()]);
assert_eq!(m.state(), MemberState::Unjoined);
assert!(m.assignment().is_empty());
m.on_join(codes::NONE, 2, "m-1", "m-2", vec![]);
m.on_sync(codes::NONE, vec![tp(0)]);
m.set_topics(vec!["t".to_owned(), "u".to_owned()]);
assert_eq!(m.state(), MemberState::Stable);
}
#[test]
fn the_subscription_carries_what_is_owned() {
let mut m = member();
m.on_join(codes::NONE, 1, "m-1", "m-2", vec![]);
m.on_sync(codes::NONE, vec![tp(3)]);
let s = m.subscription();
assert_eq!(s.member_id, "m-1");
assert_eq!(s.topics, vec!["t".to_owned()]);
assert_eq!(s.owned, vec![tp(3)]);
}
#[test]
fn leaving_gives_everything_up() {
let mut m = member();
m.on_join(codes::NONE, 1, "m-1", "m-2", vec![]);
m.on_sync(codes::NONE, vec![tp(0)]);
m.on_leave();
assert_eq!(m.member_id(), "");
assert!(m.assignment().is_empty());
assert!(!m.can_commit());
}
fn cooperative() -> GroupMember {
GroupMember::new("g", vec!["t".to_owned()]).with_protocol(RebalanceProtocol::Cooperative)
}
#[test]
fn a_cooperative_rebalance_keeps_the_partitions() {
let mut m = cooperative();
m.on_join(codes::NONE, 1, "m-1", "m-2", vec![]);
m.on_sync(codes::NONE, vec![tp(0), tp(1)]);
m.on_heartbeat(codes::REBALANCE_IN_PROGRESS);
assert_eq!(
m.assignment(),
&[tp(0), tp(1)],
"cooperative members keep what nobody has taken yet"
);
assert!(
!m.can_commit(),
"but they are not stable, so they do not commit"
);
}
#[test]
fn an_eager_rebalance_gives_everything_up() {
let mut m = member();
m.on_join(codes::NONE, 1, "m-1", "m-2", vec![]);
m.on_sync(codes::NONE, vec![tp(0), tp(1)]);
m.on_heartbeat(codes::REBALANCE_IN_PROGRESS);
assert!(m.assignment().is_empty());
}
#[test]
fn a_smaller_assignment_is_a_handover() {
let mut m = cooperative();
m.on_join(codes::NONE, 1, "m-1", "m-2", vec![]);
m.on_sync(codes::NONE, vec![tp(0), tp(1), tp(2)]);
assert!(m.lost().is_empty());
m.on_heartbeat(codes::REBALANCE_IN_PROGRESS);
m.on_join(codes::NONE, 2, "m-1", "m-2", vec![]);
let step = m.on_sync(codes::NONE, vec![tp(0)]);
assert_eq!(m.lost(), &[tp(1), tp(2)], "the two that moved");
assert_eq!(m.assignment(), &[tp(0)], "and the one that did not");
assert_eq!(
step,
Step::Join {
member_id: "m-1".to_owned()
},
"rejoin at once: the released partitions have no owner until we do"
);
}
#[test]
fn a_handover_rejoin_keeps_its_generation() {
let mut m = cooperative();
m.on_join(codes::NONE, 4, "m-1", "m-2", vec![]);
m.on_sync(codes::NONE, vec![tp(0), tp(1)]);
m.on_heartbeat(codes::REBALANCE_IN_PROGRESS);
m.on_join(codes::NONE, 5, "m-1", "m-2", vec![]);
m.on_sync(codes::NONE, vec![tp(0)]);
assert_eq!(
m.generation(),
5,
"still a member of the generation just assigned"
);
assert_eq!(
m.subscription().generation,
5,
"and it says so, so its claim on tp(0) is believed"
);
}
#[test]
fn an_unchanged_assignment_does_not_rejoin() {
let mut m = cooperative();
m.on_join(codes::NONE, 1, "m-1", "m-2", vec![]);
assert_eq!(m.on_sync(codes::NONE, vec![tp(0)]), Step::Heartbeat);
assert!(m.lost().is_empty());
assert_eq!(m.state(), MemberState::Stable);
}
#[test]
fn outside_stable_it_owns_nothing_and_commits_nothing() {
let errors = [
codes::REBALANCE_IN_PROGRESS,
codes::ILLEGAL_GENERATION,
codes::UNKNOWN_MEMBER_ID,
codes::FENCED_INSTANCE_ID,
codes::NOT_COORDINATOR,
codes::COORDINATOR_NOT_AVAILABLE,
9_999, ];
for error in errors {
let mut m = member();
let _ = RebalanceProtocol::default();
m.on_join(codes::NONE, 1, "m-1", "m-2", vec![]);
m.on_sync(codes::NONE, vec![tp(0), tp(1)]);
assert!(m.can_commit());
m.on_heartbeat(error);
assert_ne!(m.state(), MemberState::Stable, "error {error}");
assert!(m.assignment().is_empty(), "error {error} kept partitions");
assert!(!m.can_commit(), "error {error} still allowed a commit");
}
}
}