Send `AppendEntries` to an observer peer. If the observer's advisory send queue is full (backpressure threshold reached), the send is skipped. The observer will fall behind and recover via snapshot when it reconnects. Source commits are never delayed.
(&mut self, observer: u64)
| 225 | /// reached), the send is skipped. The observer will fall behind and recover |
| 226 | /// via snapshot when it reconnects. Source commits are never delayed. |
| 227 | pub(super) fn send_append_entries_to_observer(&mut self, observer: u64) { |
| 228 | let can_receive = match &self.leader_state { |
| 229 | Some(ls) => ls.observer_can_receive(observer), |
| 230 | None => return, |
| 231 | }; |
| 232 | if !can_receive { |
| 233 | debug!( |
| 234 | node = self.config.node_id, |
| 235 | group = self.config.group_id, |
| 236 | observer, |
| 237 | "observer send queue full; skipping (advisory backpressure)" |
| 238 | ); |
| 239 | return; |
| 240 | } |
| 241 | |
| 242 | let leader = match &self.leader_state { |
| 243 | Some(ls) => ls, |
| 244 | None => return, |
| 245 | }; |
| 246 | |
| 247 | let obs_state = match leader |
| 248 | .observer_states |
| 249 | .iter() |
| 250 | .find(|(id, _)| *id == observer) |
| 251 | { |
| 252 | Some((_, s)) => s.clone(), |
| 253 | None => return, |
| 254 | }; |
| 255 | |
| 256 | let next_index = obs_state.next_index; |
| 257 | let prev_log_index = next_index.saturating_sub(1); |
| 258 | |
| 259 | let prev_log_term = match self.log.term_at(prev_log_index) { |
| 260 | Some(term) => term, |
| 261 | None => { |
| 262 | debug!( |
| 263 | node = self.config.node_id, |
| 264 | group = self.config.group_id, |
| 265 | observer, |
| 266 | next_index, |
| 267 | snapshot_index = self.log.snapshot_index(), |
| 268 | "observer needs snapshot (log compacted)" |
| 269 | ); |
| 270 | self.ready.snapshots_needed.push(observer); |
| 271 | return; |
| 272 | } |
| 273 | }; |
| 274 | |
| 275 | let entries = if next_index <= self.log.last_index() { |
| 276 | match self.log.entries_range(next_index, self.log.last_index()) { |
| 277 | Ok(slice) => slice.to_vec(), |
| 278 | Err(crate::error::RaftError::LogCompacted { .. }) => { |
| 279 | debug!( |
| 280 | node = self.config.node_id, |
| 281 | group = self.config.group_id, |
| 282 | observer, |
| 283 | next_index, |
| 284 | "observer needs snapshot (entries compacted)" |
no test coverage detected