oxedyne/fe2o3/fe2o3_o3db_sync/src/dist/engine.rs
30.1 KiB, 227 runs
created by r1870400018:11382, which is this file's identity for as long as the history lasts, whatever it is later renamed to
download · who wrote it · its history
| 1 | //! The distributed-Ozone engine. |
| 2 | //! |
| 3 | //! [`DistOzone`] composes the placement service, the peer set, and a |
| 4 | //! [`Storage`] backend into one cohesive replication engine. It is a pure |
| 5 | //! state machine: every public method either reads state or returns the |
| 6 | //! outbound envelopes the caller should dispatch through its transport |
| 7 | //! adapter. No method calls `send` itself. |
| 8 | //! |
| 9 | //! # Write path |
| 10 | //! |
| 11 | //! [`DistOzone::put`] runs the placement decision, persists locally if the |
| 12 | //! local peer is a holder, and returns [`PutOutcome`] with a |
| 13 | //! [`ReplicatePut`](crate::transport::MsgKind::ReplicatePut) envelope for |
| 14 | //! every remote holder. The caller dispatches the envelopes; recipients |
| 15 | //! re-check their own placement decision, so a put that a sender directed |
| 16 | //! at a peer with a slightly different view of `N` can be silently dropped |
| 17 | //! by the recipient without harm -- the next anti-entropy round fills it |
| 18 | //! in. |
| 19 | //! |
| 20 | //! # Read path |
| 21 | //! |
| 22 | //! [`DistOzone::get`] reads from the local store if the local peer is a |
| 23 | //! holder, returning [`GetOutcome::Local`] or [`GetOutcome::LocalMiss`]. If |
| 24 | //! the local peer is *not* a holder, it returns [`GetOutcome::Remote`] with |
| 25 | //! a request id and the |
| 26 | //! [`GetRequest`](crate::transport::MsgKind::GetRequest) envelopes to |
| 27 | //! dispatch. The caller polls [`DistOzone::poll_get`] to learn when a |
| 28 | //! response has landed. |
| 29 | //! |
| 30 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 31 | //! Anthropic Claude |
| 32 | |
| 33 | use super::cohort; |
| 34 | use super::config::{ |
| 35 | Consistency, |
| 36 | DistOzoneConfig, |
| 37 | TableConfig, |
| 38 | }; |
| 39 | use super::consensus::{ |
| 40 | self, |
| 41 | CohortInstance, |
| 42 | }; |
| 43 | use super::hotstuff::{ |
| 44 | replica::Command as HsCommand, |
| 45 | types::{ |
| 46 | NewView, |
| 47 | Proposal, |
| 48 | Vote, |
| 49 | }, |
| 50 | }; |
| 51 | use super::peer_set::PeerSet; |
| 52 | use super::placement::Placement; |
| 53 | use super::record::{ |
| 54 | Record, |
| 55 | RecordId, |
| 56 | }; |
| 57 | use super::storage::Storage; |
| 58 | use super::transport::{ |
| 59 | Envelope, |
| 60 | MsgKind, |
| 61 | RequestId, |
| 62 | }; |
| 63 | |
| 64 | use oxedyne_fe2o3_core::prelude::*; |
| 65 | use oxedyne_fe2o3_data::iblt::{ |
| 66 | DecodeOutcome, |
| 67 | Iblt, |
| 68 | IbltConfig, |
| 69 | }; |
| 70 | use crate::kademlia::id::{ |
| 71 | ID_LEN, |
| 72 | NodeId, |
| 73 | }; |
| 74 | use crate::oam::config::OamConfig; |
| 75 | |
| 76 | use std::collections::HashMap; |
| 77 | use std::sync::{ |
| 78 | Mutex, |
| 79 | atomic::{ |
| 80 | AtomicU64, |
| 81 | Ordering, |
| 82 | }, |
| 83 | }; |
| 84 | |
| 85 | |
| 86 | // Fixed key and value lengths for the anti-entropy IBLT sketches. The key is a |
| 87 | // 32-byte RecordId and the value is the 32-byte content hash produced by |
| 88 | // Storage::digests. Matching lengths across peers is mandatory for IBLT |
| 89 | // subtraction; callers cannot override. |
| 90 | const ANTI_ENTROPY_KEY_LEN: usize = ID_LEN; |
| 91 | const ANTI_ENTROPY_VALUE_LEN: usize = 32; |
| 92 | |
| 93 | |
| 94 | // Choosing more than one peer protects against a single straggler; choosing |
| 95 | // many wastes bandwidth. Three is the operating-point default referenced in the |
| 96 | // spec ("an OAM holder" -- plural in realistic deployments). |
| 97 | pub const DEFAULT_READ_FANOUT: usize = 3; |
| 98 | |
| 99 | |
| 100 | pub struct DistOzone<S: Storage> { |
| 101 | cfg: DistOzoneConfig, |
| 102 | placement: Placement, |
| 103 | peer_set: PeerSet, |
| 104 | storage: S, |
| 105 | next_rid: AtomicU64, |
| 106 | pending_gets: Mutex<HashMap<RequestId, PendingGet>>, |
| 107 | read_fanout: usize, |
| 108 | // An entry is created lazily on the first message, and kept after Decide so |
| 109 | // duplicate late messages are silently absorbed. |
| 110 | cohorts: Mutex<HashMap<(String, RecordId), CohortInstance>>, |
| 111 | } |
| 112 | |
| 113 | /// The in-flight state of a remote read. |
| 114 | #[derive(Clone, Debug)] |
| 115 | struct PendingGet { |
| 116 | #[allow(dead_code)] // Preserved for diagnostics and future retry logic. |
| 117 | table: String, |
| 118 | #[allow(dead_code)] |
| 119 | id: RecordId, |
| 120 | outstanding: usize, |
| 121 | first_response: Option<Record>, |
| 122 | resolved_empty: bool, |
| 123 | } |
| 124 | |
| 125 | impl<S: Storage> DistOzone<S> { |
| 126 | /// The bootstrap peer list is filtered to exclude the local peer, and the |
| 127 | /// placement service's threshold is precomputed. |
| 128 | pub fn new(cfg: DistOzoneConfig, storage: S) -> Outcome<Self> { |
| 129 | let placement = Placement::new(cfg.local_peer_id, cfg.oam); |
| 130 | let peer_set = PeerSet::from_bootstrap( |
| 131 | &cfg.local_peer_id, |
| 132 | cfg.bootstrap_peers.iter().cloned(), |
| 133 | ); |
| 134 | Ok(Self { |
| 135 | cfg, |
| 136 | placement, |
| 137 | peer_set, |
| 138 | storage, |
| 139 | next_rid: AtomicU64::new(1), |
| 140 | pending_gets: Mutex::new(HashMap::new()), |
| 141 | read_fanout: DEFAULT_READ_FANOUT, |
| 142 | cohorts: Mutex::new(HashMap::new()), |
| 143 | }) |
| 144 | } |
| 145 | |
| 146 | /// One is the minimum; values above the peer-set size are clamped |
| 147 | /// automatically at request time. |
| 148 | pub fn set_read_fanout(&mut self, fanout: usize) { |
| 149 | self.read_fanout = fanout.max(1); |
| 150 | } |
| 151 | |
| 152 | pub fn config(&self) -> &DistOzoneConfig { |
| 153 | &self.cfg |
| 154 | } |
| 155 | |
| 156 | pub fn peer_set(&self) -> &PeerSet { |
| 157 | &self.peer_set |
| 158 | } |
| 159 | |
| 160 | pub fn placement(&self) -> &Placement { |
| 161 | &self.placement |
| 162 | } |
| 163 | |
| 164 | pub fn storage(&self) -> &S { |
| 165 | &self.storage |
| 166 | } |
| 167 | |
| 168 | /// Ignores the local peer; true if the peer was new. |
| 169 | pub fn insert_peer(&mut self, peer: NodeId) -> bool { |
| 170 | if peer == self.cfg.local_peer_id { |
| 171 | return false; |
| 172 | } |
| 173 | self.peer_set.insert(peer) |
| 174 | } |
| 175 | |
| 176 | pub fn remove_peer(&mut self, peer: &NodeId) -> bool { |
| 177 | self.peer_set.remove(peer) |
| 178 | } |
| 179 | |
| 180 | /// Recomputes the cached placement threshold. The estimate typically comes |
| 181 | /// from a HyperLogLog merge. |
| 182 | pub fn update_network_size(&mut self, network_size: u64) -> Outcome<()> { |
| 183 | let oam = res!(OamConfig::new(self.cfg.oam.replication, network_size)); |
| 184 | self.cfg.oam = oam; |
| 185 | self.placement.update_oam(oam); |
| 186 | Ok(()) |
| 187 | } |
| 188 | |
| 189 | fn table_or_err(&self, name: &str) -> Outcome<&TableConfig> { |
| 190 | match self.cfg.table(name) { |
| 191 | Some(t) => Ok(t), |
| 192 | None => Err(err!( |
| 193 | "Unknown table in DistOzone operation: {}.", name; |
| 194 | Invalid, Input, Missing)), |
| 195 | } |
| 196 | } |
| 197 | |
| 198 | /// On an eventual table the record is persisted locally if the local |
| 199 | /// peer is a holder and a [`MsgKind::ReplicatePut`] envelope is emitted |
| 200 | /// for every remote holder. |
| 201 | /// |
| 202 | /// On a cohort-backed table the put enters a HotStuff round. If the |
| 203 | /// local peer is the cohort's initial leader the engine opens a |
| 204 | /// [`MsgKind::CohortPropose`] round and returns the proposal envelopes. |
| 205 | /// Otherwise it emits a [`MsgKind::CohortSubmit`] envelope to the |
| 206 | /// leader, who drives the round on the submitter's behalf. In both |
| 207 | /// cases [`PutOutcome::consensus_pending`] is set; the caller learns |
| 208 | /// when consensus completes through |
| 209 | /// [`InboundOutcome::completed_consensus_put`]. |
| 210 | pub fn put(&self, record: Record) -> Outcome<PutOutcome> { |
| 211 | let tc = res!(self.table_or_err(&record.table)); |
| 212 | match tc.consistency { |
| 213 | Consistency::Eventual => self.put_eventual(record), |
| 214 | Consistency::Cohort { lambda } => |
| 215 | self.put_cohort(record, lambda), |
| 216 | } |
| 217 | } |
| 218 | |
| 219 | fn put_eventual(&self, record: Record) -> Outcome<PutOutcome> { |
| 220 | let decision = self.placement.decide(&record.id, &self.peer_set); |
| 221 | let local_persisted = decision.local_is_holder; |
| 222 | if local_persisted { |
| 223 | res!(self.storage.put(&record)); |
| 224 | } |
| 225 | |
| 226 | let mut outbound = Vec::with_capacity(decision.remote_holders.len()); |
| 227 | for peer in decision.remote_holders { |
| 228 | outbound.push(Envelope::new( |
| 229 | self.cfg.local_peer_id, |
| 230 | *peer, |
| 231 | MsgKind::ReplicatePut { record: record.clone() }, |
| 232 | )); |
| 233 | } |
| 234 | Ok(PutOutcome { |
| 235 | local_persisted, |
| 236 | outbound, |
| 237 | consensus_pending: None, |
| 238 | }) |
| 239 | } |
| 240 | |
| 241 | fn put_cohort(&self, record: Record, lambda: u64) -> Outcome<PutOutcome> { |
| 242 | let table = record.table.clone(); |
| 243 | let id = record.id; |
| 244 | let sel = res!(cohort::select( |
| 245 | &table, |
| 246 | &id, |
| 247 | &self.peer_set, |
| 248 | &self.cfg.local_peer_id, |
| 249 | lambda, |
| 250 | )); |
| 251 | if sel.local_is_leader { |
| 252 | // Drive the round locally. |
| 253 | let outbound = res!(self.leader_open_round(sel, record)); |
| 254 | Ok(PutOutcome { |
| 255 | local_persisted: false, |
| 256 | outbound, |
| 257 | consensus_pending: Some((table, id)), |
| 258 | }) |
| 259 | } else { |
| 260 | // Forward to the leader. |
| 261 | let env = Envelope::new( |
| 262 | self.cfg.local_peer_id, |
| 263 | sel.leader, |
| 264 | MsgKind::CohortSubmit { record }, |
| 265 | ); |
| 266 | Ok(PutOutcome { |
| 267 | local_persisted: false, |
| 268 | outbound: vec![env], |
| 269 | consensus_pending: Some((table, id)), |
| 270 | }) |
| 271 | } |
| 272 | } |
| 273 | |
| 274 | fn leader_open_round( |
| 275 | &self, |
| 276 | sel: cohort::Cohort, |
| 277 | record: Record, |
| 278 | ) |
| 279 | -> Outcome<Vec<Envelope>> |
| 280 | { |
| 281 | let table = record.table.clone(); |
| 282 | let id = record.id; |
| 283 | let block = consensus::encode_record(&record); |
| 284 | let block_hash = consensus::block_hash(&block); |
| 285 | |
| 286 | let lambda = sel.members.len() as u64; |
| 287 | let mut cohorts = lock_mutex!(self.cohorts); |
| 288 | let instance = match cohorts.get_mut(&(table.clone(), id)) { |
| 289 | Some(i) => i, |
| 290 | None => { |
| 291 | let fresh = res!(CohortInstance::new( |
| 292 | sel, |
| 293 | &self.cfg.local_peer_id, |
| 294 | lambda, |
| 295 | )); |
| 296 | cohorts.insert((table.clone(), id), fresh); |
| 297 | cohorts.get_mut(&(table.clone(), id)).expect("just inserted") |
| 298 | }, |
| 299 | }; |
| 300 | if instance.has_decided() { |
| 301 | return Ok(Vec::new()); |
| 302 | } |
| 303 | let cmds = res!(instance.replica.propose(block, block_hash)); |
| 304 | let outbound = res!(self.translate_commands(&table, &id, instance, cmds)); |
| 305 | Ok(outbound) |
| 306 | } |
| 307 | |
| 308 | /// Applies the local side effects -- persistence on Decide -- as it goes. |
| 309 | fn translate_commands( |
| 310 | &self, |
| 311 | table: &str, |
| 312 | id: &RecordId, |
| 313 | instance: &mut CohortInstance, |
| 314 | cmds: Vec<HsCommand>, |
| 315 | ) |
| 316 | -> Outcome<Vec<Envelope>> |
| 317 | { |
| 318 | let mut out = Vec::new(); |
| 319 | for cmd in cmds { |
| 320 | match cmd { |
| 321 | HsCommand::BroadcastProposal(proposal) => { |
| 322 | // Broadcast to every cohort member *except* ourselves. |
| 323 | // HotStuff's on_proposal recipe expects the leader to |
| 324 | // record its own vote via on_proposal too, so we also |
| 325 | // feed the proposal back into our local replica. |
| 326 | let self_id = instance.replica.config().self_id; |
| 327 | for (i, member) in instance.members.iter().enumerate() { |
| 328 | if (i as u16) == self_id { |
| 329 | continue; |
| 330 | } |
| 331 | out.push(Envelope::new( |
| 332 | self.cfg.local_peer_id, |
| 333 | *member, |
| 334 | MsgKind::CohortPropose { |
| 335 | table: table.to_string(), |
| 336 | id: *id, |
| 337 | proposal: proposal.clone(), |
| 338 | }, |
| 339 | )); |
| 340 | } |
| 341 | // Now fold the proposal into the local replica. We |
| 342 | // recurse on the commands produced, which lets a leader |
| 343 | // that is also a voter correctly emit its own SendVote. |
| 344 | let local_cmds = res!(instance.replica.on_proposal(proposal)); |
| 345 | let more = res!(self.translate_commands( |
| 346 | table, id, instance, local_cmds, |
| 347 | )); |
| 348 | out.extend(more); |
| 349 | }, |
| 350 | HsCommand::SendVote { to, vote } => { |
| 351 | let target = match instance.node_id(to) { |
| 352 | Some(n) => n, |
| 353 | None => return Err(err!( |
| 354 | "HotStuff SendVote targets replica id {} \ |
| 355 | which is out of cohort range (size = {}).", |
| 356 | to, instance.members.len(); |
| 357 | Bug, Invalid)), |
| 358 | }; |
| 359 | // The leader of the current view is the usual target; a |
| 360 | // vote sent to ourselves needs to be folded in directly |
| 361 | // rather than sent on the wire. |
| 362 | if target == self.cfg.local_peer_id { |
| 363 | let local_cmds = res!(instance.replica.on_vote(vote)); |
| 364 | let more = res!(self.translate_commands( |
| 365 | table, id, instance, local_cmds, |
| 366 | )); |
| 367 | out.extend(more); |
| 368 | } else { |
| 369 | out.push(Envelope::new( |
| 370 | self.cfg.local_peer_id, |
| 371 | target, |
| 372 | MsgKind::CohortVote { |
| 373 | table: table.to_string(), |
| 374 | id: *id, |
| 375 | vote, |
| 376 | }, |
| 377 | )); |
| 378 | } |
| 379 | }, |
| 380 | HsCommand::SendNewView { to, new_view } => { |
| 381 | let target = match instance.node_id(to) { |
| 382 | Some(n) => n, |
| 383 | None => return Err(err!( |
| 384 | "HotStuff SendNewView targets replica id {} \ |
| 385 | which is out of cohort range (size = {}).", |
| 386 | to, instance.members.len(); |
| 387 | Bug, Invalid)), |
| 388 | }; |
| 389 | if target == self.cfg.local_peer_id { |
| 390 | let local_cmds = res!(instance.replica.on_new_view(new_view)); |
| 391 | let more = res!(self.translate_commands( |
| 392 | table, id, instance, local_cmds, |
| 393 | )); |
| 394 | out.extend(more); |
| 395 | } else { |
| 396 | out.push(Envelope::new( |
| 397 | self.cfg.local_peer_id, |
| 398 | target, |
| 399 | MsgKind::CohortNewView { |
| 400 | table: table.to_string(), |
| 401 | id: *id, |
| 402 | new_view, |
| 403 | }, |
| 404 | )); |
| 405 | } |
| 406 | }, |
| 407 | HsCommand::Decide { view: _, block } => { |
| 408 | let record = res!(consensus::decode_record(&block)); |
| 409 | if record.table != table || &record.id != id { |
| 410 | return Err(err!( |
| 411 | "HotStuff Decide block decoded to a record at \ |
| 412 | ({}, {:?}) that does not match the consensus \ |
| 413 | slot ({}, {:?}).", |
| 414 | record.table, record.id.as_bytes(), |
| 415 | table, id.as_bytes(); |
| 416 | Bug, Invalid, Mismatch)); |
| 417 | } |
| 418 | let hash = consensus::block_hash(&block); |
| 419 | res!(instance.mark_decided(hash)); |
| 420 | res!(self.storage.put(&record)); |
| 421 | }, |
| 422 | } |
| 423 | } |
| 424 | Ok(out) |
| 425 | } |
| 426 | |
| 427 | /// The caller owns the timer; on expiry it calls this and dispatches the |
| 428 | /// returned envelopes. Empty if the instance is absent -- already decided, |
| 429 | /// or never created. |
| 430 | pub fn cohort_timeout( |
| 431 | &self, |
| 432 | table: &str, |
| 433 | id: &RecordId, |
| 434 | ) |
| 435 | -> Outcome<Vec<Envelope>> |
| 436 | { |
| 437 | let mut cohorts = lock_mutex!(self.cohorts); |
| 438 | let instance = match cohorts.get_mut(&(table.to_string(), *id)) { |
| 439 | Some(i) => i, |
| 440 | None => return Ok(Vec::new()), |
| 441 | }; |
| 442 | if instance.has_decided() { |
| 443 | return Ok(Vec::new()); |
| 444 | } |
| 445 | let cmds = res!(instance.replica.on_timeout()); |
| 446 | self.translate_commands(table, id, instance, cmds) |
| 447 | } |
| 448 | |
| 449 | /// Reads from local storage if the local peer is a holder; otherwise |
| 450 | /// dispatches a request to the nearest peers and returns the in-flight |
| 451 | /// request handle. |
| 452 | pub fn get( |
| 453 | &self, |
| 454 | table: &str, |
| 455 | id: &RecordId, |
| 456 | ) |
| 457 | -> Outcome<GetOutcome> |
| 458 | { |
| 459 | res!(self.table_or_err(table)); |
| 460 | |
| 461 | if self.placement.i_am_holder(id) { |
| 462 | return Ok(match res!(self.storage.get(table, id)) { |
| 463 | Some(r) => GetOutcome::Local(r), |
| 464 | None => GetOutcome::LocalMiss, |
| 465 | }); |
| 466 | } |
| 467 | |
| 468 | let targets = self.placement.read_targets( |
| 469 | id, |
| 470 | &self.peer_set, |
| 471 | self.read_fanout, |
| 472 | ); |
| 473 | if targets.is_empty() { |
| 474 | return Ok(GetOutcome::NoTargets); |
| 475 | } |
| 476 | let request_id = self.next_rid.fetch_add(1, Ordering::Relaxed); |
| 477 | let outbound = targets.iter() |
| 478 | .map(|peer| Envelope::new( |
| 479 | self.cfg.local_peer_id, |
| 480 | **peer, |
| 481 | MsgKind::GetRequest { |
| 482 | request_id, |
| 483 | table: table.to_string(), |
| 484 | id: *id, |
| 485 | }, |
| 486 | )) |
| 487 | .collect::<Vec<_>>(); |
| 488 | { |
| 489 | let mut pending = lock_mutex!(self.pending_gets); |
| 490 | pending.insert(request_id, PendingGet { |
| 491 | table: table.to_string(), |
| 492 | id: *id, |
| 493 | outstanding: outbound.len(), |
| 494 | first_response: None, |
| 495 | resolved_empty: false, |
| 496 | }); |
| 497 | } |
| 498 | Ok(GetOutcome::Remote { request_id, outbound }) |
| 499 | } |
| 500 | |
| 501 | pub fn handle_envelope(&self, env: Envelope) -> Outcome<InboundOutcome> { |
| 502 | if env.to != self.cfg.local_peer_id { |
| 503 | // An envelope addressed to somebody else; ignore. This is mostly |
| 504 | // a belt-and-braces guard: the transport adapter should not |
| 505 | // deliver misaddressed envelopes in the first place. |
| 506 | return Ok(InboundOutcome::empty()); |
| 507 | } |
| 508 | match env.body { |
| 509 | MsgKind::ReplicatePut { record } => { |
| 510 | // Re-check placement: sender's view of N may disagree with |
| 511 | // ours near the threshold. Drop silently if we do not |
| 512 | // consider ourselves a holder. |
| 513 | if !self.placement.i_am_holder(&record.id) { |
| 514 | return Ok(InboundOutcome::empty()); |
| 515 | } |
| 516 | // Reject puts for tables we do not know about. |
| 517 | if self.cfg.table(&record.table).is_none() { |
| 518 | return Err(err!( |
| 519 | "ReplicatePut for unknown table '{}'.", record.table; |
| 520 | Invalid, Input, Missing)); |
| 521 | } |
| 522 | res!(self.storage.put(&record)); |
| 523 | Ok(InboundOutcome::empty()) |
| 524 | } |
| 525 | MsgKind::GetRequest { request_id, table, id } => { |
| 526 | if self.cfg.table(&table).is_none() { |
| 527 | return Err(err!( |
| 528 | "GetRequest for unknown table '{}'.", table; |
| 529 | Invalid, Input, Missing)); |
| 530 | } |
| 531 | let record = res!(self.storage.get(&table, &id)); |
| 532 | let reply = Envelope::new( |
| 533 | self.cfg.local_peer_id, |
| 534 | env.from, |
| 535 | MsgKind::GetResponse { request_id, record }, |
| 536 | ); |
| 537 | Ok(InboundOutcome { |
| 538 | outbound: vec![reply], |
| 539 | completed_get: None, |
| 540 | completed_consensus_put: None, |
| 541 | }) |
| 542 | } |
| 543 | MsgKind::GetResponse { request_id, record } => { |
| 544 | let completed = { |
| 545 | let mut pending = lock_mutex!(self.pending_gets); |
| 546 | let Some(slot) = pending.get_mut(&request_id) else { |
| 547 | // Unknown request id; stale or cancelled response. |
| 548 | return Ok(InboundOutcome::empty()); |
| 549 | }; |
| 550 | if slot.outstanding > 0 { |
| 551 | slot.outstanding -= 1; |
| 552 | } |
| 553 | match record { |
| 554 | Some(r) if slot.first_response.is_none() => { |
| 555 | slot.first_response = Some(r); |
| 556 | } |
| 557 | None => { |
| 558 | slot.resolved_empty = true; |
| 559 | } |
| 560 | _ => { /* later non-first response; ignore. */ } |
| 561 | } |
| 562 | // A pending get is "resolved" when either a response with |
| 563 | // a record has landed or every target has replied empty. |
| 564 | slot.first_response.is_some() || slot.outstanding == 0 |
| 565 | }; |
| 566 | Ok(InboundOutcome { |
| 567 | outbound: Vec::new(), |
| 568 | completed_get: if completed { Some(request_id) } else { None }, |
| 569 | completed_consensus_put: None, |
| 570 | }) |
| 571 | } |
| 572 | MsgKind::AntiEntropyDigest { table, sketch } => { |
| 573 | self.handle_anti_entropy_digest(env.from, table, sketch) |
| 574 | } |
| 575 | MsgKind::AntiEntropyReply { table, records, requested_ids, bulk } => { |
| 576 | self.handle_anti_entropy_reply( |
| 577 | env.from, table, records, requested_ids, bulk, |
| 578 | ) |
| 579 | } |
| 580 | MsgKind::AntiEntropyPush { table, records } => { |
| 581 | self.handle_anti_entropy_push(table, records) |
| 582 | } |
| 583 | MsgKind::CohortSubmit { record } => { |
| 584 | self.handle_cohort_submit(record) |
| 585 | } |
| 586 | MsgKind::CohortPropose { table, id, proposal } => { |
| 587 | self.handle_cohort_propose(env.from, table, id, proposal) |
| 588 | } |
| 589 | MsgKind::CohortVote { table, id, vote } => { |
| 590 | self.handle_cohort_vote(table, id, vote) |
| 591 | } |
| 592 | MsgKind::CohortNewView { table, id, new_view } => { |
| 593 | self.handle_cohort_new_view(table, id, new_view) |
| 594 | } |
| 595 | } |
| 596 | } |
| 597 | |
| 598 | pub fn poll_get(&self, request_id: RequestId) -> Outcome<PollOutcome> { |
| 599 | let pending = lock_mutex!(self.pending_gets); |
| 600 | let Some(slot) = pending.get(&request_id) else { |
| 601 | return Ok(PollOutcome::Unknown); |
| 602 | }; |
| 603 | if let Some(r) = &slot.first_response { |
| 604 | return Ok(PollOutcome::Record(r.clone())); |
| 605 | } |
| 606 | if slot.outstanding == 0 && slot.resolved_empty { |
| 607 | return Ok(PollOutcome::NotFound); |
| 608 | } |
| 609 | Ok(PollOutcome::Pending) |
| 610 | } |
| 611 | |
| 612 | /// Late responses that arrive after cancellation are ignored by the next |
| 613 | /// [`handle_envelope`](Self::handle_envelope) call. |
| 614 | pub fn cancel_get(&self, request_id: RequestId) -> Outcome<()> { |
| 615 | let mut pending = lock_mutex!(self.pending_gets); |
| 616 | pending.remove(&request_id); |
| 617 | Ok(()) |
| 618 | } |
| 619 | |
| 620 | /// The envelope carries a serialised IBLT built from the local storage's |
| 621 | /// [`Storage::digests`] enumeration for that table. The recipient |
| 622 | /// subtracts it against its own sketch, decodes the symmetric difference |
| 623 | /// and answers with [`MsgKind::AntiEntropyReply`]. Cohort-backed tables |
| 624 | /// reconcile through consensus rather than anti-entropy, and are rejected. |
| 625 | pub fn build_anti_entropy_request( |
| 626 | &self, |
| 627 | table: &str, |
| 628 | target: NodeId, |
| 629 | ) |
| 630 | -> Outcome<Envelope> |
| 631 | { |
| 632 | let tc = res!(self.table_or_err(table)); |
| 633 | if !matches!(tc.consistency, Consistency::Eventual) { |
| 634 | return Err(err!( |
| 635 | "anti-entropy is defined only for Eventual tables \ |
| 636 | (table '{}').", table; |
| 637 | Invalid, Input, Unimplemented)); |
| 638 | } |
| 639 | let iblt = res!(self.build_table_iblt(tc)); |
| 640 | let sketch = iblt.to_bytes(); |
| 641 | Ok(Envelope::new( |
| 642 | self.cfg.local_peer_id, |
| 643 | target, |
| 644 | MsgKind::AntiEntropyDigest { |
| 645 | table: table.to_string(), |
| 646 | sketch, |
| 647 | }, |
| 648 | )) |
| 649 | } |
| 650 | |
| 651 | /// Factored out so both the request builder and the inbound digest handler |
| 652 | /// use the same sketch shape. |
| 653 | fn build_table_iblt(&self, tc: &TableConfig) -> Outcome<Iblt> { |
| 654 | let cfg = IbltConfig { |
| 655 | num_cells: tc.iblt_cells, |
| 656 | num_hashes: TableConfig::IBLT_NUM_HASHES, |
| 657 | key_len: ANTI_ENTROPY_KEY_LEN, |
| 658 | value_len: ANTI_ENTROPY_VALUE_LEN, |
| 659 | seed: tc.iblt_seed(), |
| 660 | }; |
| 661 | let mut iblt = res!(Iblt::new(cfg)); |
| 662 | let digests = res!(self.storage.digests(&tc.name)); |
| 663 | for d in digests { |
| 664 | res!(iblt.insert(d.id.as_bytes(), &d.content)); |
| 665 | } |
| 666 | Ok(iblt) |
| 667 | } |
| 668 | |
| 669 | /// Decodes the symmetric difference against the local sketch and returns an |
| 670 | /// [`AntiEntropyReply`][ar] envelope carrying records the sender lacks and |
| 671 | /// a list of record identifiers the recipient lacks. On sketch overload it |
| 672 | /// falls back to a bulk reply of every record the recipient holds for the |
| 673 | /// table. |
| 674 | /// |
| 675 | /// [ar]: MsgKind::AntiEntropyReply |
| 676 | fn handle_anti_entropy_digest( |
| 677 | &self, |
| 678 | from: NodeId, |
| 679 | table: String, |
| 680 | sketch: Vec<u8>, |
| 681 | ) |
| 682 | -> Outcome<InboundOutcome> |
| 683 | { |
| 684 | let tc = res!(self.table_or_err(&table)); |
| 685 | if !matches!(tc.consistency, Consistency::Eventual) { |
| 686 | return Err(err!( |
| 687 | "anti-entropy digest received for non-Eventual table '{}'.", |
| 688 | table; |
| 689 | Invalid, Input, Unimplemented)); |
| 690 | } |
| 691 | let expected_cfg = IbltConfig { |
| 692 | num_cells: tc.iblt_cells, |
| 693 | num_hashes: TableConfig::IBLT_NUM_HASHES, |
| 694 | key_len: ANTI_ENTROPY_KEY_LEN, |
| 695 | value_len: ANTI_ENTROPY_VALUE_LEN, |
| 696 | seed: tc.iblt_seed(), |
| 697 | }; |
| 698 | let their_iblt = res!(Iblt::from_bytes(&sketch)); |
| 699 | if their_iblt.config() != expected_cfg { |
| 700 | return Err(err!( |
| 701 | "anti-entropy sketch config mismatch for table '{}'.", |
| 702 | table; |
| 703 | Invalid, Input, Mismatch)); |
| 704 | } |
| 705 | let mut mine = res!(self.build_table_iblt(tc)); |
| 706 | res!(mine.subtract(&their_iblt)); |
| 707 | let decode = res!(mine.decode()); |
| 708 | let (records_for_sender, requested_ids, bulk) = match decode { |
| 709 | DecodeOutcome::Complete { inserted, deleted } => { |
| 710 | // `inserted` = keys in mine not in theirs -> records I |
| 711 | // should send. `deleted` = keys in theirs not in mine -> |
| 712 | // ids I should request. |
| 713 | let mut records_for_sender = Vec::with_capacity(inserted.len()); |
| 714 | for (key_bytes, _value_hash) in inserted { |
| 715 | let rid = res!(RecordId::from_slice(&key_bytes)); |
| 716 | if let Some(r) = res!(self.storage.get(&table, &rid)) { |
| 717 | records_for_sender.push(r); |
| 718 | } |
| 719 | } |
| 720 | let mut requested_ids = Vec::with_capacity(deleted.len()); |
| 721 | for (key_bytes, _value_hash) in deleted { |
| 722 | requested_ids.push(res!(RecordId::from_slice(&key_bytes))); |
| 723 | } |
| 724 | (records_for_sender, requested_ids, false) |
| 725 | } |
| 726 | DecodeOutcome::Incomplete { .. } => { |
| 727 | // Sketch overloaded. Fall back to bulk: send everything I |
| 728 | // have for this table; the sender absorbs what it lacks. |
| 729 | // This is simple and correct; a later optimisation can |
| 730 | // teach the sender to retry with a larger sketch. |
| 731 | let digests = res!(self.storage.digests(&table)); |
| 732 | let mut records = Vec::with_capacity(digests.len()); |
| 733 | for d in digests { |
| 734 | if let Some(r) = res!(self.storage.get(&table, &d.id)) { |
| 735 | records.push(r); |
| 736 | } |
| 737 | } |
| 738 | (records, Vec::new(), true) |
| 739 | } |
| 740 | }; |
| 741 | let reply = Envelope::new( |
| 742 | self.cfg.local_peer_id, |
| 743 | from, |
| 744 | MsgKind::AntiEntropyReply { |
| 745 | table, |
| 746 | records: records_for_sender, |
| 747 | requested_ids, |
| 748 | bulk, |
| 749 | }, |
| 750 | ); |
| 751 | Ok(InboundOutcome { |
| 752 | outbound: vec![reply], |
| 753 | completed_get: None, |
| 754 | completed_consensus_put: None, |
| 755 | }) |
| 756 | } |
| 757 | |
| 758 | /// Applies the records the recipient was missing, and builds an |
| 759 | /// [`AntiEntropyPush`][ap] envelope for any records requested in return. |
| 760 | /// |
| 761 | /// [ap]: MsgKind::AntiEntropyPush |
| 762 | fn handle_anti_entropy_reply( |
| 763 | &self, |
| 764 | from: NodeId, |
| 765 | table: String, |
| 766 | records: Vec<Record>, |
| 767 | requested_ids: Vec<RecordId>, |
| 768 | _bulk: bool, |
| 769 | ) |
| 770 | -> Outcome<InboundOutcome> |
| 771 | { |
| 772 | res!(self.table_or_err(&table)); |
| 773 | |
| 774 | // Apply every record the peer sent us, re-checking placement so a |
| 775 | // stale-N sender cannot push a record to a peer that has since |
| 776 | // stopped considering itself a holder. |
| 777 | for record in records { |
| 778 | if record.table != table { |
| 779 | continue; |
| 780 | } |
| 781 | if self.placement.i_am_holder(&record.id) { |
| 782 | res!(self.storage.put(&record)); |
| 783 | } |
| 784 | } |
| 785 | |
| 786 | // Build a push for every requested id we actually have. |
| 787 | let mut to_push = Vec::with_capacity(requested_ids.len()); |
| 788 | for rid in requested_ids { |
| 789 | if let Some(r) = res!(self.storage.get(&table, &rid)) { |
| 790 | to_push.push(r); |
| 791 | } |
| 792 | } |
| 793 | if to_push.is_empty() { |
| 794 | return Ok(InboundOutcome::empty()); |
| 795 | } |
| 796 | let push = Envelope::new( |
| 797 | self.cfg.local_peer_id, |
| 798 | from, |
| 799 | MsgKind::AntiEntropyPush { |
| 800 | table, |
| 801 | records: to_push, |
| 802 | }, |
| 803 | ); |
| 804 | Ok(InboundOutcome { |
| 805 | outbound: vec![push], |
| 806 | completed_get: None, |
| 807 | completed_consensus_put: None, |
| 808 | }) |
| 809 | } |
| 810 | |
| 811 | /// Opens a HotStuff round if the local peer is the initial leader for the |
| 812 | /// `(table, record_id)` pair. If it is not, the submission is dropped |
| 813 | /// silently -- almost always a stale-cohort race, where the submitter's |
| 814 | /// view of the peer set differed from the leader's. |
| 815 | fn handle_cohort_submit( |
| 816 | &self, |
| 817 | record: Record, |
| 818 | ) |
| 819 | -> Outcome<InboundOutcome> |
| 820 | { |
| 821 | let tc = res!(self.table_or_err(&record.table)); |
| 822 | let lambda = match tc.consistency { |
| 823 | Consistency::Cohort { lambda } => lambda, |
| 824 | Consistency::Eventual => return Err(err!( |
| 825 | "CohortSubmit received for eventual-consistency table '{}'.", |
| 826 | record.table; |
| 827 | Invalid, Input, Mismatch)), |
| 828 | }; |
| 829 | let sel = res!(cohort::select( |
| 830 | &record.table, |
| 831 | &record.id, |
| 832 | &self.peer_set, |
| 833 | &self.cfg.local_peer_id, |
| 834 | lambda, |
| 835 | )); |
| 836 | if !sel.local_is_leader { |
| 837 | // Not our job -- drop silently. |
| 838 | return Ok(InboundOutcome::empty()); |
| 839 | } |
| 840 | let outbound = res!(self.leader_open_round(sel, record)); |
| 841 | Ok(InboundOutcome { |
| 842 | outbound, |
| 843 | completed_get: None, |
| 844 | completed_consensus_put: None, |
| 845 | }) |
| 846 | } |
| 847 | |
| 848 | /// Creates a fresh per-record replica if one does not yet exist. A proposal |
| 849 | /// addressed to a peer that is not a cohort member is dropped silently. |
| 850 | fn handle_cohort_propose( |
| 851 | &self, |
| 852 | from: NodeId, |
| 853 | table: String, |
| 854 | id: RecordId, |
| 855 | proposal: Proposal, |
| 856 | ) |
| 857 | -> Outcome<InboundOutcome> |
| 858 | { |
| 859 | let tc = res!(self.table_or_err(&table)); |
| 860 | let lambda = match tc.consistency { |
| 861 | Consistency::Cohort { lambda } => lambda, |
| 862 | Consistency::Eventual => return Err(err!( |
| 863 | "CohortPropose received for eventual-consistency table '{}'.", |
| 864 | table; |
| 865 | Invalid, Input, Mismatch)), |
| 866 | }; |
| 867 | let sel = res!(cohort::select( |
| 868 | &table, |
| 869 | &id, |
| 870 | &self.peer_set, |
| 871 | &self.cfg.local_peer_id, |
| 872 | lambda, |
| 873 | )); |
| 874 | if !sel.local_is_member { |
| 875 | // Proposal arrived but we are not in the cohort -- drop. This |
| 876 | // can happen if peer-set views disagree; the sender will retry |
| 877 | // after its own set updates. |
| 878 | return Ok(InboundOutcome::empty()); |
| 879 | } |
| 880 | // Verify the proposal came from someone plausibly in the cohort. |
| 881 | // Stronger origin checks (view-aware leader matching) live in the |
| 882 | // HotStuff replica itself. |
| 883 | if !sel.members.iter().any(|m| m == &from) { |
| 884 | return Ok(InboundOutcome::empty()); |
| 885 | } |
| 886 | let mut cohorts = lock_mutex!(self.cohorts); |
| 887 | let key = (table.clone(), id); |
| 888 | let instance = match cohorts.get_mut(&key) { |
| 889 | Some(i) => i, |
| 890 | None => { |
| 891 | let fresh = res!(CohortInstance::new( |
| 892 | sel, |
| 893 | &self.cfg.local_peer_id, |
| 894 | lambda, |
| 895 | )); |
| 896 | cohorts.insert(key.clone(), fresh); |
| 897 | cohorts.get_mut(&key).expect("just inserted") |
| 898 | }, |
| 899 | }; |
| 900 | if instance.has_decided() { |
| 901 | return Ok(InboundOutcome::empty()); |
| 902 | } |
| 903 | let before_decided = instance.has_decided(); |
| 904 | let cmds = res!(instance.replica.on_proposal(proposal)); |
| 905 | let outbound = res!(self.translate_commands(&table, &id, instance, cmds)); |
| 906 | let completed_consensus_put = if !before_decided && instance.has_decided() { |
| 907 | Some((table, id)) |
| 908 | } else { |
| 909 | None |
| 910 | }; |
| 911 | Ok(InboundOutcome { |
| 912 | outbound, |
| 913 | completed_get: None, |
| 914 | completed_consensus_put, |
| 915 | }) |
| 916 | } |
| 917 | |
| 918 | /// Leader-only: a non-leader replica silently ignores a vote, per the |
| 919 | /// HotStuff spec. |
| 920 | fn handle_cohort_vote( |
| 921 | &self, |
| 922 | table: String, |
| 923 | id: RecordId, |
| 924 | vote: Vote, |
| 925 | ) |
| 926 | -> Outcome<InboundOutcome> |
| 927 | { |
| 928 | let mut cohorts = lock_mutex!(self.cohorts); |
| 929 | let key = (table.clone(), id); |
| 930 | let instance = match cohorts.get_mut(&key) { |
| 931 | Some(i) => i, |
| 932 | None => { |
| 933 | // Vote arrived before we created an instance. Drop -- a |
| 934 | // well-behaved cohort member only votes after seeing the |
| 935 | // leader's Propose, so by the time a vote reaches us we |
| 936 | // should already have an instance. This case is most |
| 937 | // likely an adversarial or replayed envelope. |
| 938 | return Ok(InboundOutcome::empty()); |
| 939 | }, |
| 940 | }; |
| 941 | if instance.has_decided() { |
| 942 | return Ok(InboundOutcome::empty()); |
| 943 | } |
| 944 | let before_decided = instance.has_decided(); |
| 945 | let cmds = res!(instance.replica.on_vote(vote)); |
| 946 | let outbound = res!(self.translate_commands(&table, &id, instance, cmds)); |
| 947 | let completed_consensus_put = if !before_decided && instance.has_decided() { |
| 948 | Some((table, id)) |
| 949 | } else { |
| 950 | None |
| 951 | }; |
| 952 | Ok(InboundOutcome { |
| 953 | outbound, |
| 954 | completed_get: None, |
| 955 | completed_consensus_put, |
| 956 | }) |
| 957 | } |
| 958 | |
| 959 | fn handle_cohort_new_view( |
| 960 | &self, |
| 961 | table: String, |
| 962 | id: RecordId, |
| 963 | new_view: NewView, |
| 964 | ) |
| 965 | -> Outcome<InboundOutcome> |
| 966 | { |
| 967 | let mut cohorts = lock_mutex!(self.cohorts); |
| 968 | let key = (table.clone(), id); |
| 969 | let instance = match cohorts.get_mut(&key) { |
| 970 | Some(i) => i, |
| 971 | None => return Ok(InboundOutcome::empty()), |
| 972 | }; |
| 973 | if instance.has_decided() { |
| 974 | return Ok(InboundOutcome::empty()); |
| 975 | } |
| 976 | let cmds = res!(instance.replica.on_new_view(new_view)); |
| 977 | let outbound = res!(self.translate_commands(&table, &id, instance, cmds)); |
| 978 | Ok(InboundOutcome { |
| 979 | outbound, |
| 980 | completed_get: None, |
| 981 | completed_consensus_put: None, |
| 982 | }) |
| 983 | } |
| 984 | |
| 985 | /// Each record is placement-checked before persistence. |
| 986 | fn handle_anti_entropy_push( |
| 987 | &self, |
| 988 | table: String, |
| 989 | records: Vec<Record>, |
| 990 | ) |
| 991 | -> Outcome<InboundOutcome> |
| 992 | { |
| 993 | res!(self.table_or_err(&table)); |
| 994 | for record in records { |
| 995 | if record.table != table { |
| 996 | continue; |
| 997 | } |
| 998 | if self.placement.i_am_holder(&record.id) { |
| 999 | res!(self.storage.put(&record)); |
| 1000 | } |
| 1001 | } |
| 1002 | Ok(InboundOutcome::empty()) |
| 1003 | } |
| 1004 | } |
| 1005 | |
| 1006 | |
| 1007 | /// The result of a [`DistOzone::put`] call. |
| 1008 | #[derive(Clone, Debug)] |
| 1009 | pub struct PutOutcome { |
| 1010 | // Always false for a cohort-backed write, which persists on Decide and is |
| 1011 | // signalled through InboundOutcome::completed_consensus_put. |
| 1012 | pub local_persisted: bool, |
| 1013 | pub outbound: Vec<Envelope>, |
| 1014 | // Set when the put entered a HotStuff consensus round. |
| 1015 | pub consensus_pending: Option<(String, RecordId)>, |
| 1016 | } |
| 1017 | |
| 1018 | |
| 1019 | /// The result of a [`DistOzone::get`] call. |
| 1020 | #[derive(Clone, Debug)] |
| 1021 | pub enum GetOutcome { |
| 1022 | Local(Record), |
| 1023 | LocalMiss, // a holder, but no record at that id |
| 1024 | // A remote read has been initiated; completion is reported through |
| 1025 | // DistOzone::poll_get. |
| 1026 | Remote { |
| 1027 | request_id: RequestId, |
| 1028 | outbound: Vec<Envelope>, |
| 1029 | }, |
| 1030 | NoTargets, // not a holder, and no remote targets are known |
| 1031 | } |
| 1032 | |
| 1033 | |
| 1034 | /// The result of a [`DistOzone::handle_envelope`] call. |
| 1035 | #[derive(Clone, Debug)] |
| 1036 | pub struct InboundOutcome { |
| 1037 | pub outbound: Vec<Envelope>, |
| 1038 | // Set when this envelope completed a pending remote read; the caller polls |
| 1039 | // DistOzone::poll_get to collect the record itself. |
| 1040 | pub completed_get: Option<RequestId>, |
| 1041 | // Set when this envelope drove a cohort round to Decide and the record was |
| 1042 | // persisted locally. |
| 1043 | pub completed_consensus_put: Option<(String, RecordId)>, |
| 1044 | } |
| 1045 | |
| 1046 | impl InboundOutcome { |
| 1047 | fn empty() -> Self { |
| 1048 | Self { |
| 1049 | outbound: Vec::new(), |
| 1050 | completed_get: None, |
| 1051 | completed_consensus_put: None, |
| 1052 | } |
| 1053 | } |
| 1054 | } |
| 1055 | |
| 1056 | |
| 1057 | /// The result of a [`DistOzone::poll_get`] query. |
| 1058 | #[derive(Clone, Debug)] |
| 1059 | pub enum PollOutcome { |
| 1060 | Pending, |
| 1061 | Record(Record), |
| 1062 | NotFound, // every outstanding holder replied that the record is absent |
| 1063 | Unknown, // unknown, or cancelled |
| 1064 | } |