Register a consumer joining a group. Triggers rebalance.
(&self, tenant_id: u64, stream: &str, group: &str, consumer_id: &str)
| 72 | |
| 73 | /// Register a consumer joining a group. Triggers rebalance. |
| 74 | pub fn join(&self, tenant_id: u64, stream: &str, group: &str, consumer_id: &str) { |
| 75 | let key = (tenant_id, stream.to_string(), group.to_string()); |
| 76 | let mut groups = self.groups.write().unwrap_or_else(|p| p.into_inner()); |
| 77 | let state = groups.entry(key).or_insert_with(GroupAssignment::new); |
| 78 | state.consumers.insert(consumer_id.to_string()); |
| 79 | state.rebalance(); |
| 80 | |
| 81 | tracing::debug!( |
| 82 | stream, |
| 83 | group, |
| 84 | consumer_id, |
| 85 | total_consumers = state.consumers.len(), |
| 86 | "consumer joined, rebalanced" |
| 87 | ); |
| 88 | } |
| 89 | |
| 90 | /// Deregister a consumer leaving a group. Triggers rebalance. |
| 91 | pub fn leave(&self, tenant_id: u64, stream: &str, group: &str, consumer_id: &str) { |