oxedyne/fe2o3/fe2o3_o3db_sync/src/dist/consensus.rs
9.9 KiB, 27 runs
created by r1870400018:11462, 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 | //! Per-record HotStuff consensus state for cohort-backed tables. |
| 2 | //! |
| 3 | //! Cohort-backed tables ([`Consistency::Cohort`][c]) serialise writes through |
| 4 | //! a HotStuff consensus cohort. This module owns the per-`(table, record_id)` |
| 5 | //! state that bridges the pure [`fe2o3_hotstuff::replica::Replica`] state |
| 6 | //! machine to the wire: it maps between [`NodeId`] (the peer-space |
| 7 | //! identifier) and [`ReplicaId`] (HotStuff's cohort-local index), it encodes |
| 8 | //! and decodes the block payload so HotStuff can carry a full [`Record`] as |
| 9 | //! opaque bytes, and it hashes the block deterministically so every peer |
| 10 | //! agrees on the identity of what is being decided. |
| 11 | //! |
| 12 | //! The module is deliberately pure -- it holds no transport, no timers, and |
| 13 | //! no storage handle. The engine ([`crate::dist::DistOzone`]) drives it by |
| 14 | //! feeding inbound HotStuff messages in and translating the emitted |
| 15 | //! [`fe2o3_hotstuff::replica::Command`] list back out into |
| 16 | //! [`crate::transport::Envelope`]s. |
| 17 | //! |
| 18 | //! [c]: crate::dist::config::Consistency::Cohort |
| 19 | //! |
| 20 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 21 | //! Anthropic Claude |
| 22 | |
| 23 | use super::cohort::Cohort; |
| 24 | use super::hotstuff::{ |
| 25 | replica::{ |
| 26 | Config as HsConfig, |
| 27 | Replica, |
| 28 | }, |
| 29 | types::{ |
| 30 | BLOCK_HASH_LEN, |
| 31 | BlockHash, |
| 32 | ReplicaId, |
| 33 | }, |
| 34 | }; |
| 35 | use super::record::{ |
| 36 | Record, |
| 37 | RecordId, |
| 38 | }; |
| 39 | |
| 40 | use oxedyne_fe2o3_core::prelude::*; |
| 41 | use crate::kademlia::id::NodeId; |
| 42 | |
| 43 | |
| 44 | /// `f = floor((lambda - 1) / 3)`, matching the Hematite spec and the |
| 45 | /// requirement that HotStuff sees `cohort_size >= 3f + 1`. |
| 46 | pub fn faults_tolerated(lambda: u64) -> usize { |
| 47 | ((lambda.saturating_sub(1)) / 3) as usize |
| 48 | } |
| 49 | |
| 50 | |
| 51 | /// Per-`(table, record_id)` HotStuff state bundled with the cohort membership. |
| 52 | /// |
| 53 | /// One instance is created per record that goes through consensus. The |
| 54 | /// [`Replica`] inside drives the three-phase protocol; the cohort membership |
| 55 | /// is a snapshot of the deterministic selection at creation time. Members are |
| 56 | /// indexed `0..cohort_size` in the order produced by |
| 57 | /// [`cohort::select`](crate::dist::cohort::select), so every peer's instance uses |
| 58 | /// the same mapping without exchanging it on the wire. |
| 59 | pub struct CohortInstance { |
| 60 | pub replica: Replica, |
| 61 | pub members: Vec<NodeId>, // indexed by ReplicaId, as Cohort::members |
| 62 | // Some only after Decide, which makes a duplicate Decide a no-op. |
| 63 | pub decided_hash: Option<BlockHash>, |
| 64 | } |
| 65 | |
| 66 | impl CohortInstance { |
| 67 | /// Rejects cohorts the local peer does not appear in -- such a peer has |
| 68 | /// no role in the consensus and should never hold an instance. |
| 69 | pub fn new(cohort: Cohort, local_id: &NodeId, lambda: u64) -> Outcome<Self> { |
| 70 | if !cohort.local_is_member { |
| 71 | return Err(err!( |
| 72 | "CohortInstance requires the local peer to be a cohort member."; |
| 73 | Invalid, Input, Missing)); |
| 74 | } |
| 75 | let members = cohort.members; |
| 76 | let self_id = match members.iter().position(|m| m == local_id) { |
| 77 | Some(i) => i as ReplicaId, |
| 78 | None => return Err(err!( |
| 79 | "Local peer missing from cohort members despite \ |
| 80 | local_is_member = true."; |
| 81 | Bug, Invalid)), |
| 82 | }; |
| 83 | let cfg = HsConfig { |
| 84 | cohort_size: members.len(), |
| 85 | f: faults_tolerated(lambda), |
| 86 | self_id, |
| 87 | }; |
| 88 | let replica = res!(Replica::new(cfg)); |
| 89 | Ok(Self { |
| 90 | replica, |
| 91 | members, |
| 92 | decided_hash: None, |
| 93 | }) |
| 94 | } |
| 95 | |
| 96 | pub fn replica_id(&self, node: &NodeId) -> Option<ReplicaId> { |
| 97 | self.members.iter() |
| 98 | .position(|m| m == node) |
| 99 | .map(|i| i as ReplicaId) |
| 100 | } |
| 101 | |
| 102 | pub fn node_id(&self, rid: ReplicaId) -> Option<NodeId> { |
| 103 | self.members.get(rid as usize).copied() |
| 104 | } |
| 105 | |
| 106 | pub fn current_leader(&self) -> Outcome<NodeId> { |
| 107 | let cfg = self.replica.config(); |
| 108 | let leader_rid = cfg.leader_for(self.replica.view()); |
| 109 | match self.node_id(leader_rid) { |
| 110 | Some(n) => Ok(n), |
| 111 | None => Err(err!( |
| 112 | "HotStuff leader id {} out of cohort range (size = {}).", |
| 113 | leader_rid, self.members.len(); |
| 114 | Bug, Invalid)), |
| 115 | } |
| 116 | } |
| 117 | |
| 118 | /// Idempotent: a duplicate Decide on the same hash is accepted silently; a |
| 119 | /// Decide on a different hash is a bug and is rejected. |
| 120 | pub fn mark_decided(&mut self, hash: BlockHash) -> Outcome<()> { |
| 121 | match self.decided_hash { |
| 122 | Some(prev) if prev == hash => Ok(()), |
| 123 | Some(_) => Err(err!( |
| 124 | "CohortInstance decided twice on different blocks -- safety \ |
| 125 | violation in upstream HotStuff."; |
| 126 | Bug, Invalid)), |
| 127 | None => { |
| 128 | self.decided_hash = Some(hash); |
| 129 | Ok(()) |
| 130 | } |
| 131 | } |
| 132 | } |
| 133 | |
| 134 | /// New wire messages for a decided instance are dropped by the engine. |
| 135 | pub fn has_decided(&self) -> bool { |
| 136 | self.decided_hash.is_some() |
| 137 | } |
| 138 | } |
| 139 | |
| 140 | |
| 141 | /// Serialises a [`Record`] into the block payload carried by HotStuff. |
| 142 | /// |
| 143 | /// Wire format (little-endian throughout): |
| 144 | /// ```text |
| 145 | /// [32 bytes record id][4 bytes table_len][table_bytes][4 bytes value_len][value_bytes] |
| 146 | /// ``` |
| 147 | /// |
| 148 | /// Length fields cap at `u32::MAX`; this is enforced at decode time rather |
| 149 | /// than at encode time (callers produce `String` and `Vec<u8>` both of which |
| 150 | /// can exceed that in principle, but in practice are tiny). |
| 151 | pub fn encode_record(record: &Record) -> Vec<u8> { |
| 152 | let table_bytes = record.table.as_bytes(); |
| 153 | let mut out = Vec::with_capacity( |
| 154 | 32 + 4 + table_bytes.len() + 4 + record.value.len(), |
| 155 | ); |
| 156 | out.extend_from_slice(record.id.as_bytes()); |
| 157 | out.extend_from_slice(&(table_bytes.len() as u32).to_le_bytes()); |
| 158 | out.extend_from_slice(table_bytes); |
| 159 | out.extend_from_slice(&(record.value.len() as u32).to_le_bytes()); |
| 160 | out.extend_from_slice(&record.value); |
| 161 | out |
| 162 | } |
| 163 | |
| 164 | |
| 165 | /// Validates that the declared lengths fit the buffer. |
| 166 | pub fn decode_record(bytes: &[u8]) -> Outcome<Record> { |
| 167 | if bytes.len() < 32 + 4 { |
| 168 | return Err(err!( |
| 169 | "Record block too short: {} bytes, need at least 36.", |
| 170 | bytes.len(); |
| 171 | Invalid, Input, Size)); |
| 172 | } |
| 173 | let mut id_bytes = [0u8; 32]; |
| 174 | id_bytes.copy_from_slice(&bytes[..32]); |
| 175 | let id = RecordId::from_bytes(id_bytes); |
| 176 | |
| 177 | let mut cursor = 32; |
| 178 | let table_len = { |
| 179 | let mut buf = [0u8; 4]; |
| 180 | buf.copy_from_slice(&bytes[cursor..cursor + 4]); |
| 181 | cursor += 4; |
| 182 | u32::from_le_bytes(buf) as usize |
| 183 | }; |
| 184 | if bytes.len() < cursor + table_len + 4 { |
| 185 | return Err(err!( |
| 186 | "Record block truncated: table_len {} exceeds remaining bytes.", |
| 187 | table_len; |
| 188 | Invalid, Input, Size)); |
| 189 | } |
| 190 | let table = match std::str::from_utf8(&bytes[cursor..cursor + table_len]) { |
| 191 | Ok(s) => s.to_string(), |
| 192 | Err(_) => return Err(err!( |
| 193 | "Record block table name is not valid UTF-8."; |
| 194 | Invalid, Input)), |
| 195 | }; |
| 196 | cursor += table_len; |
| 197 | |
| 198 | let value_len = { |
| 199 | let mut buf = [0u8; 4]; |
| 200 | buf.copy_from_slice(&bytes[cursor..cursor + 4]); |
| 201 | cursor += 4; |
| 202 | u32::from_le_bytes(buf) as usize |
| 203 | }; |
| 204 | if bytes.len() != cursor + value_len { |
| 205 | return Err(err!( |
| 206 | "Record block length mismatch: expected {} bytes after header, \ |
| 207 | got {}.", |
| 208 | value_len, bytes.len() - cursor; |
| 209 | Invalid, Input, Size)); |
| 210 | } |
| 211 | let value = bytes[cursor..].to_vec(); |
| 212 | Ok(Record { id, table, value }) |
| 213 | } |
| 214 | |
| 215 | |
| 216 | /// Deterministic 32-byte hash of a block payload. |
| 217 | /// |
| 218 | /// Implementation: splitmix64-style mixing, widened to 32 bytes by mixing the |
| 219 | /// accumulated state through four successive multiplicative rounds. Matches |
| 220 | /// the hashing style used elsewhere in this crate (cohort seed derivation, |
| 221 | /// IBLT seed) so that the whole dist-Ozone layer remains dependency-minimal |
| 222 | /// and self-contained. |
| 223 | /// |
| 224 | /// This hash is *deterministic* and *collision-robust against unstructured |
| 225 | /// input*, but is not a cryptographic hash. The HotStuff protocol itself |
| 226 | /// does not require cryptographic-strength block hashes -- only that every |
| 227 | /// honest peer computes the same hash for the same block, which splitmix64 |
| 228 | /// satisfies. A deployment that wants collision-resistance against |
| 229 | /// adversarial leaders can swap this function for SHA3-256 or BLAKE3 without |
| 230 | /// touching the rest of the crate; the [`BlockHash`] alias is already |
| 231 | /// 32 bytes for both. |
| 232 | pub fn block_hash(block: &[u8]) -> BlockHash { |
| 233 | let mut state: u64 = 0x9E3779B97F4A7C15; |
| 234 | for chunk in block.chunks(8) { |
| 235 | let mut buf = [0u8; 8]; |
| 236 | buf[..chunk.len()].copy_from_slice(chunk); |
| 237 | let word = u64::from_le_bytes(buf); |
| 238 | state = state.wrapping_add(word); |
| 239 | state = (state ^ (state >> 30)).wrapping_mul(0xBF58476D1CE4E5B9); |
| 240 | state = (state ^ (state >> 27)).wrapping_mul(0x94D049BB133111EB); |
| 241 | state ^= state >> 31; |
| 242 | } |
| 243 | // Also mix in the length so blocks of different lengths hash differently |
| 244 | // even if the trailing bytes align. |
| 245 | state = state.wrapping_add(block.len() as u64); |
| 246 | state = (state ^ (state >> 30)).wrapping_mul(0xBF58476D1CE4E5B9); |
| 247 | state = (state ^ (state >> 27)).wrapping_mul(0x94D049BB133111EB); |
| 248 | state ^= state >> 31; |
| 249 | let mut out = [0u8; BLOCK_HASH_LEN]; |
| 250 | for i in 0..4 { |
| 251 | let limb = state.wrapping_mul(0x9E3779B97F4A7C15 ^ ((i as u64) + 1)); |
| 252 | out[i * 8..(i + 1) * 8].copy_from_slice(&limb.to_le_bytes()); |
| 253 | } |
| 254 | out |
| 255 | } |
| 256 | |
| 257 | |
| 258 | #[cfg(test)] |
| 259 | mod tests { |
| 260 | use super::*; |
| 261 | |
| 262 | #[test] |
| 263 | fn record_round_trip_through_encode_decode() -> Outcome<()> { |
| 264 | let id = RecordId::from_bytes([0x42; 32]); |
| 265 | let rec = Record::new(id, "identity", b"hello world".to_vec()); |
| 266 | let bytes = encode_record(&rec); |
| 267 | let back = res!(decode_record(&bytes)); |
| 268 | assert_eq!(back, rec); |
| 269 | Ok(()) |
| 270 | } |
| 271 | |
| 272 | #[test] |
| 273 | fn record_with_empty_value_round_trips() -> Outcome<()> { |
| 274 | let id = RecordId::from_bytes([1; 32]); |
| 275 | let rec = Record::new(id, "epoch", Vec::new()); |
| 276 | let bytes = encode_record(&rec); |
| 277 | let back = res!(decode_record(&bytes)); |
| 278 | assert_eq!(back, rec); |
| 279 | Ok(()) |
| 280 | } |
| 281 | |
| 282 | #[test] |
| 283 | fn decode_rejects_truncated_block() { |
| 284 | assert!(decode_record(&[0u8; 10]).is_err()); |
| 285 | } |
| 286 | |
| 287 | #[test] |
| 288 | fn decode_rejects_wrong_trailing_length() { |
| 289 | let id = RecordId::from_bytes([2; 32]); |
| 290 | let rec = Record::new(id, "identity", b"v".to_vec()); |
| 291 | let mut bytes = encode_record(&rec); |
| 292 | bytes.pop(); |
| 293 | assert!(decode_record(&bytes).is_err()); |
| 294 | } |
| 295 | |
| 296 | #[test] |
| 297 | fn block_hash_is_deterministic() { |
| 298 | let a = block_hash(b"some block bytes"); |
| 299 | let b = block_hash(b"some block bytes"); |
| 300 | assert_eq!(a, b); |
| 301 | } |
| 302 | |
| 303 | #[test] |
| 304 | fn block_hash_differs_for_distinct_payloads() { |
| 305 | let a = block_hash(b"some block bytes"); |
| 306 | let b = block_hash(b"some block byteS"); |
| 307 | assert_ne!(a, b); |
| 308 | } |
| 309 | |
| 310 | #[test] |
| 311 | fn faults_tolerated_matches_spec() { |
| 312 | assert_eq!(faults_tolerated(5), 1); |
| 313 | assert_eq!(faults_tolerated(7), 2); |
| 314 | assert_eq!(faults_tolerated(9), 2); |
| 315 | } |
| 316 | } |