oxedyne/fe2o3/fe2o3_data/src/iblt/imp.rs
13.8 KiB, 5 runs
created by r1870400018:11214, 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 Invertible Bloom Lookup Table. |
| 2 | //! |
| 3 | //! An [`Iblt`] holds `num_cells` cells, each consisting of four accumulators: |
| 4 | //! a key XOR, a value XOR, a 64-bit key-hash XOR used as a purity check and a |
| 5 | //! signed insert count. Every key is inserted into `num_hashes` cells chosen |
| 6 | //! by double-hashing. Symmetric difference between two IBLTs with the same |
| 7 | //! shape is computed by cellwise XOR on the byte accumulators and cellwise |
| 8 | //! subtraction on the counts; the peeling decoder then extracts keys from |
| 9 | //! "pure" cells -- those with `|count| == 1` and a key-hash fingerprint that |
| 10 | //! matches the recomputed hash of the extracted key -- and removes their |
| 11 | //! contributions from all cells, iterating until no pure cells remain. |
| 12 | |
| 13 | use super::hash::{ |
| 14 | hash_bytes, |
| 15 | hash_pair, |
| 16 | }; |
| 17 | |
| 18 | use oxedyne_fe2o3_core::prelude::*; |
| 19 | |
| 20 | |
| 21 | /// The fixed number of bytes used for the purity-check fingerprint. |
| 22 | pub const FINGERPRINT_LEN: usize = 8; |
| 23 | |
| 24 | /// The fixed number of bytes used for the signed count accumulator. |
| 25 | pub const COUNT_LEN: usize = 4; |
| 26 | |
| 27 | |
| 28 | /// Parameters shared by all cells of an [`Iblt`]. Returned by |
| 29 | /// [`Iblt::config`] and consumed by [`Iblt::from_bytes`] to restore an IBLT |
| 30 | /// from a serialised form. |
| 31 | #[derive(Clone, Copy, Debug, Eq, PartialEq)] |
| 32 | pub struct IbltConfig { |
| 33 | /// Number of cells in the table. |
| 34 | pub num_cells: usize, |
| 35 | /// Number of hash-derived cell indices each key maps to. |
| 36 | pub num_hashes: usize, |
| 37 | /// Fixed key length in bytes. |
| 38 | pub key_len: usize, |
| 39 | /// Fixed value length in bytes. Zero for key-only IBLTs. |
| 40 | pub value_len: usize, |
| 41 | /// Seed shared by peers that intend to reconcile with one another. |
| 42 | pub seed: u64, |
| 43 | } |
| 44 | |
| 45 | |
| 46 | /// The outcome of a peeling decode. |
| 47 | #[derive(Clone, Debug)] |
| 48 | pub enum DecodeOutcome { |
| 49 | /// Every cell was drained. `inserted` and `deleted` are the two halves of |
| 50 | /// the symmetric difference between the IBLTs that produced the decoded |
| 51 | /// table. |
| 52 | Complete { |
| 53 | /// Keys that had `count > 0` at extraction time, paired with their |
| 54 | /// recovered values. |
| 55 | inserted: Vec<(Vec<u8>, Vec<u8>)>, |
| 56 | /// Keys that had `count < 0` at extraction time, paired with their |
| 57 | /// recovered values. |
| 58 | deleted: Vec<(Vec<u8>, Vec<u8>)>, |
| 59 | }, |
| 60 | /// Decoding halted with non-empty cells remaining -- the IBLT was |
| 61 | /// overloaded relative to the true symmetric-difference size. The partial |
| 62 | /// results before the halt are preserved; callers can either fall back to |
| 63 | /// a bulk transfer or retry with a larger IBLT. |
| 64 | Incomplete { |
| 65 | /// Keys extracted before decoding stalled. |
| 66 | inserted: Vec<(Vec<u8>, Vec<u8>)>, |
| 67 | /// Keys extracted before decoding stalled. |
| 68 | deleted: Vec<(Vec<u8>, Vec<u8>)>, |
| 69 | /// Number of cells still containing non-trivial state. |
| 70 | remaining_cells: usize, |
| 71 | }, |
| 72 | } |
| 73 | |
| 74 | |
| 75 | /// An Invertible Bloom Lookup Table over fixed-length keys and values. |
| 76 | #[derive(Clone, Debug)] |
| 77 | pub struct Iblt { |
| 78 | cfg: IbltConfig, |
| 79 | /// XOR accumulator of inserted keys, `key_len` bytes per cell, flattened. |
| 80 | key_xor: Vec<u8>, |
| 81 | /// XOR accumulator of inserted values, `value_len` bytes per cell, |
| 82 | /// flattened. Empty when `value_len == 0`. |
| 83 | value_xor: Vec<u8>, |
| 84 | /// 64-bit key-hash XOR for the purity check, one entry per cell. |
| 85 | fp_xor: Vec<u64>, |
| 86 | /// Signed insert count per cell. |
| 87 | count: Vec<i32>, |
| 88 | } |
| 89 | |
| 90 | impl Iblt { |
| 91 | /// Constructs an empty IBLT with the given configuration. |
| 92 | pub fn new(cfg: IbltConfig) -> Outcome<Self> { |
| 93 | if cfg.num_cells == 0 { |
| 94 | return Err(err!( |
| 95 | "IBLT num_cells must be greater than zero."; |
| 96 | Invalid, Input)); |
| 97 | } |
| 98 | if cfg.num_hashes == 0 { |
| 99 | return Err(err!( |
| 100 | "IBLT num_hashes must be greater than zero."; |
| 101 | Invalid, Input)); |
| 102 | } |
| 103 | if cfg.num_hashes > cfg.num_cells { |
| 104 | return Err(err!( |
| 105 | "IBLT num_hashes ({}) cannot exceed num_cells ({}).", |
| 106 | cfg.num_hashes, cfg.num_cells; |
| 107 | Invalid, Input)); |
| 108 | } |
| 109 | if cfg.key_len == 0 { |
| 110 | return Err(err!( |
| 111 | "IBLT key_len must be greater than zero."; |
| 112 | Invalid, Input)); |
| 113 | } |
| 114 | Ok(Self { |
| 115 | cfg, |
| 116 | key_xor: vec![0u8; cfg.num_cells * cfg.key_len], |
| 117 | value_xor: vec![0u8; cfg.num_cells * cfg.value_len], |
| 118 | fp_xor: vec![0u64; cfg.num_cells], |
| 119 | count: vec![0i32; cfg.num_cells], |
| 120 | }) |
| 121 | } |
| 122 | |
| 123 | /// Returns the shared configuration. |
| 124 | pub fn config(&self) -> IbltConfig { |
| 125 | self.cfg |
| 126 | } |
| 127 | |
| 128 | /// Inserts `(key, value)`, or re-inserts with the same sign if already |
| 129 | /// present. Length mismatches are errors. |
| 130 | pub fn insert(&mut self, key: &[u8], value: &[u8]) -> Outcome<()> { |
| 131 | self.apply(key, value, 1) |
| 132 | } |
| 133 | |
| 134 | /// Records a deletion of `(key, value)`. After subtraction with another |
| 135 | /// IBLT this is what distinguishes "B has an extra copy" from "A has an |
| 136 | /// extra copy" during decoding. |
| 137 | pub fn delete(&mut self, key: &[u8], value: &[u8]) -> Outcome<()> { |
| 138 | self.apply(key, value, -1) |
| 139 | } |
| 140 | |
| 141 | /// Subtracts another IBLT in place. Both IBLTs must share the same |
| 142 | /// [`IbltConfig`]; mismatches are errors. |
| 143 | pub fn subtract(&mut self, other: &Self) -> Outcome<()> { |
| 144 | if self.cfg != other.cfg { |
| 145 | return Err(err!( |
| 146 | "IBLT subtract requires matching configuration."; |
| 147 | Invalid, Input, Mismatch)); |
| 148 | } |
| 149 | for (dst, src) in self.key_xor.iter_mut().zip(other.key_xor.iter()) { |
| 150 | *dst ^= *src; |
| 151 | } |
| 152 | for (dst, src) in self.value_xor.iter_mut().zip(other.value_xor.iter()) { |
| 153 | *dst ^= *src; |
| 154 | } |
| 155 | for (dst, src) in self.fp_xor.iter_mut().zip(other.fp_xor.iter()) { |
| 156 | *dst ^= *src; |
| 157 | } |
| 158 | for (dst, src) in self.count.iter_mut().zip(other.count.iter()) { |
| 159 | *dst = dst.wrapping_sub(*src); |
| 160 | } |
| 161 | Ok(()) |
| 162 | } |
| 163 | |
| 164 | /// Runs the peeling decoder, draining the IBLT of every entry it can |
| 165 | /// recover. |
| 166 | /// |
| 167 | /// Mutates the IBLT in place: on return the cells that contributed to a |
| 168 | /// recovered entry have been reduced, and the remaining cells (if any) |
| 169 | /// are those that could not be peeled. |
| 170 | pub fn decode(&mut self) -> Outcome<DecodeOutcome> { |
| 171 | let mut inserted: Vec<(Vec<u8>, Vec<u8>)> = Vec::new(); |
| 172 | let mut deleted: Vec<(Vec<u8>, Vec<u8>)> = Vec::new(); |
| 173 | |
| 174 | // Queue of cell indices that may currently be pure. We re-enqueue |
| 175 | // the `num_hashes` cells touched by each extraction. |
| 176 | let mut queue: Vec<usize> = (0..self.cfg.num_cells).collect(); |
| 177 | |
| 178 | while let Some(idx) = queue.pop() { |
| 179 | if !self.cell_is_pure(idx) { |
| 180 | continue; |
| 181 | } |
| 182 | let (sign, key, value) = res!(self.extract_cell(idx)); |
| 183 | // Remove this entry from every cell it affects. |
| 184 | let cell_idxs = self.cells_for(&key); |
| 185 | for &ci in &cell_idxs { |
| 186 | self.apply_at(ci, &key, &value, -sign); |
| 187 | } |
| 188 | // Any of those cells might now be pure. |
| 189 | for ci in cell_idxs { |
| 190 | queue.push(ci); |
| 191 | } |
| 192 | if sign > 0 { |
| 193 | inserted.push((key, value)); |
| 194 | } else { |
| 195 | deleted.push((key, value)); |
| 196 | } |
| 197 | } |
| 198 | |
| 199 | // Count residual cells -- any cell still holding state is a failure. |
| 200 | let mut remaining = 0usize; |
| 201 | for i in 0..self.cfg.num_cells { |
| 202 | if !self.cell_is_empty(i) { |
| 203 | remaining += 1; |
| 204 | } |
| 205 | } |
| 206 | Ok(if remaining == 0 { |
| 207 | DecodeOutcome::Complete { inserted, deleted } |
| 208 | } else { |
| 209 | DecodeOutcome::Incomplete { |
| 210 | inserted, |
| 211 | deleted, |
| 212 | remaining_cells: remaining, |
| 213 | } |
| 214 | }) |
| 215 | } |
| 216 | |
| 217 | /// Returns `true` if every cell is at its identity state (no keys, no |
| 218 | /// values, zero fingerprint, zero count). |
| 219 | pub fn is_empty(&self) -> bool { |
| 220 | self.count.iter().all(|&c| c == 0) |
| 221 | && self.fp_xor.iter().all(|&f| f == 0) |
| 222 | && self.key_xor.iter().all(|&b| b == 0) |
| 223 | && self.value_xor.iter().all(|&b| b == 0) |
| 224 | } |
| 225 | |
| 226 | /// Serialises the IBLT into a compact byte buffer. |
| 227 | /// |
| 228 | /// Format: |
| 229 | /// |
| 230 | /// - 5 × u64 little-endian: `num_cells`, `num_hashes`, `key_len`, |
| 231 | /// `value_len`, `seed`. |
| 232 | /// - `num_cells × (key_len + value_len + 8 + 4)` bytes: per-cell |
| 233 | /// `key_xor || value_xor || fp_xor_le || count_le`. |
| 234 | pub fn to_bytes(&self) -> Vec<u8> { |
| 235 | let per_cell = self.cfg.key_len + self.cfg.value_len |
| 236 | + FINGERPRINT_LEN + COUNT_LEN; |
| 237 | let mut out = Vec::with_capacity(8 * 5 + per_cell * self.cfg.num_cells); |
| 238 | out.extend_from_slice(&(self.cfg.num_cells as u64).to_le_bytes()); |
| 239 | out.extend_from_slice(&(self.cfg.num_hashes as u64).to_le_bytes()); |
| 240 | out.extend_from_slice(&(self.cfg.key_len as u64).to_le_bytes()); |
| 241 | out.extend_from_slice(&(self.cfg.value_len as u64).to_le_bytes()); |
| 242 | out.extend_from_slice(&self.cfg.seed.to_le_bytes()); |
| 243 | for i in 0..self.cfg.num_cells { |
| 244 | let kr = self.key_range(i); |
| 245 | out.extend_from_slice(&self.key_xor[kr]); |
| 246 | let vr = self.value_range(i); |
| 247 | out.extend_from_slice(&self.value_xor[vr]); |
| 248 | out.extend_from_slice(&self.fp_xor[i].to_le_bytes()); |
| 249 | out.extend_from_slice(&self.count[i].to_le_bytes()); |
| 250 | } |
| 251 | out |
| 252 | } |
| 253 | |
| 254 | /// How many cells the serialised form declares, read out of its header |
| 255 | /// without the table being built. |
| 256 | /// |
| 257 | /// A reader sizes its answer from this, and sizing an answer is exactly what |
| 258 | /// must not cost the allocation the number asks for. |
| 259 | pub fn cells_in(bytes: &[u8]) -> Outcome<usize> { |
| 260 | if bytes.len() < 8 * 5 { |
| 261 | return Err(err!( |
| 262 | "IBLT serialised form too short: {} bytes.", bytes.len(); |
| 263 | Invalid, Input, Size)); |
| 264 | } |
| 265 | let mut buf = [0u8; 8]; |
| 266 | buf.copy_from_slice(&bytes[..8]); |
| 267 | Ok(u64::from_le_bytes(buf) as usize) |
| 268 | } |
| 269 | |
| 270 | /// Parses the serialised form produced by [`Iblt::to_bytes`]. |
| 271 | pub fn from_bytes(bytes: &[u8]) -> Outcome<Self> { |
| 272 | if bytes.len() < 8 * 5 { |
| 273 | return Err(err!( |
| 274 | "IBLT serialised form too short: {} bytes.", bytes.len(); |
| 275 | Invalid, Input, Size)); |
| 276 | } |
| 277 | let read_u64 = |off: usize| -> u64 { |
| 278 | let mut buf = [0u8; 8]; |
| 279 | buf.copy_from_slice(&bytes[off..off + 8]); |
| 280 | u64::from_le_bytes(buf) |
| 281 | }; |
| 282 | let num_cells = read_u64(0) as usize; |
| 283 | let num_hashes = read_u64(8) as usize; |
| 284 | let key_len = read_u64(16) as usize; |
| 285 | let value_len = read_u64(24) as usize; |
| 286 | let seed = read_u64(32); |
| 287 | let cfg = IbltConfig { num_cells, num_hashes, key_len, value_len, seed }; |
| 288 | |
| 289 | let per_cell = key_len + value_len + FINGERPRINT_LEN + COUNT_LEN; |
| 290 | let body_len = res!(num_cells.checked_mul(per_cell).ok_or_else(|| err!( |
| 291 | "IBLT dimensions overflow: num_cells * per_cell."; |
| 292 | Invalid, Input, Size))); |
| 293 | let expected = 40 + body_len; |
| 294 | if bytes.len() != expected { |
| 295 | return Err(err!( |
| 296 | "IBLT serialised form length mismatch: got {}, expected {}.", |
| 297 | bytes.len(), expected; |
| 298 | Invalid, Input, Size)); |
| 299 | } |
| 300 | |
| 301 | let mut iblt = res!(Self::new(cfg)); |
| 302 | let mut off = 40; |
| 303 | for i in 0..num_cells { |
| 304 | let kr = iblt.key_range(i); |
| 305 | iblt.key_xor[kr].copy_from_slice(&bytes[off..off + key_len]); |
| 306 | off += key_len; |
| 307 | let vr = iblt.value_range(i); |
| 308 | iblt.value_xor[vr].copy_from_slice(&bytes[off..off + value_len]); |
| 309 | off += value_len; |
| 310 | let mut fp_buf = [0u8; 8]; |
| 311 | fp_buf.copy_from_slice(&bytes[off..off + FINGERPRINT_LEN]); |
| 312 | iblt.fp_xor[i] = u64::from_le_bytes(fp_buf); |
| 313 | off += FINGERPRINT_LEN; |
| 314 | let mut c_buf = [0u8; 4]; |
| 315 | c_buf.copy_from_slice(&bytes[off..off + COUNT_LEN]); |
| 316 | iblt.count[i] = i32::from_le_bytes(c_buf); |
| 317 | off += COUNT_LEN; |
| 318 | } |
| 319 | Ok(iblt) |
| 320 | } |
| 321 | |
| 322 | // --- internals ---------------------------------------------------- |
| 323 | |
| 324 | fn apply(&mut self, key: &[u8], value: &[u8], sign: i32) -> Outcome<()> { |
| 325 | if key.len() != self.cfg.key_len { |
| 326 | return Err(err!( |
| 327 | "IBLT key length mismatch: got {}, expected {}.", |
| 328 | key.len(), self.cfg.key_len; |
| 329 | Invalid, Input, Size)); |
| 330 | } |
| 331 | if value.len() != self.cfg.value_len { |
| 332 | return Err(err!( |
| 333 | "IBLT value length mismatch: got {}, expected {}.", |
| 334 | value.len(), self.cfg.value_len; |
| 335 | Invalid, Input, Size)); |
| 336 | } |
| 337 | let cells = self.cells_for(key); |
| 338 | for ci in cells { |
| 339 | self.apply_at(ci, key, value, sign); |
| 340 | } |
| 341 | Ok(()) |
| 342 | } |
| 343 | |
| 344 | /// Applies `(sign × (key, value))` to a specific cell without re-deriving |
| 345 | /// the target cells. Used both by the public insert/delete paths (which |
| 346 | /// walk every target cell) and by the decoder (which removes a recovered |
| 347 | /// entry from the cells it affected). |
| 348 | fn apply_at(&mut self, ci: usize, key: &[u8], value: &[u8], sign: i32) { |
| 349 | let kr = self.key_range(ci); |
| 350 | for (dst, src) in self.key_xor[kr].iter_mut().zip(key.iter()) { |
| 351 | *dst ^= *src; |
| 352 | } |
| 353 | let vr = self.value_range(ci); |
| 354 | for (dst, src) in self.value_xor[vr].iter_mut().zip(value.iter()) { |
| 355 | *dst ^= *src; |
| 356 | } |
| 357 | self.fp_xor[ci] ^= self.fingerprint(key); |
| 358 | self.count[ci] = self.count[ci].wrapping_add(sign); |
| 359 | } |
| 360 | |
| 361 | fn cells_for(&self, key: &[u8]) -> Vec<usize> { |
| 362 | let (h1, h2) = hash_pair(key, self.cfg.seed); |
| 363 | let m = self.cfg.num_cells as u64; |
| 364 | let mut out = Vec::with_capacity(self.cfg.num_hashes); |
| 365 | // Double hashing. To keep the k hashes truly distinct even when h2 is |
| 366 | // a factor of m (unlikely for typical seeds but possible), guard with |
| 367 | // a linear-probe fallback that advances by one cell until a fresh |
| 368 | // index is found. This preserves correctness for pathological seeds |
| 369 | // without distorting typical behaviour. |
| 370 | for i in 0..self.cfg.num_hashes { |
| 371 | let base = h1.wrapping_add((i as u64).wrapping_mul(h2)); |
| 372 | let mut idx = (base % m) as usize; |
| 373 | while out.contains(&idx) { |
| 374 | idx = (idx + 1) % self.cfg.num_cells; |
| 375 | } |
| 376 | out.push(idx); |
| 377 | } |
| 378 | out |
| 379 | } |
| 380 | |
| 381 | fn fingerprint(&self, key: &[u8]) -> u64 { |
| 382 | hash_bytes(key, self.cfg.seed ^ 0xc2b2_ae3d_27d4_eb4f) |
| 383 | } |
| 384 | |
| 385 | fn cell_is_pure(&self, ci: usize) -> bool { |
| 386 | let c = self.count[ci]; |
| 387 | if c != 1 && c != -1 { |
| 388 | return false; |
| 389 | } |
| 390 | let kr = self.key_range(ci); |
| 391 | let key = &self.key_xor[kr]; |
| 392 | self.fingerprint(key) == self.fp_xor[ci] |
| 393 | } |
| 394 | |
| 395 | fn cell_is_empty(&self, ci: usize) -> bool { |
| 396 | if self.count[ci] != 0 { |
| 397 | return false; |
| 398 | } |
| 399 | if self.fp_xor[ci] != 0 { |
| 400 | return false; |
| 401 | } |
| 402 | let kr = self.key_range(ci); |
| 403 | if self.key_xor[kr].iter().any(|&b| b != 0) { |
| 404 | return false; |
| 405 | } |
| 406 | let vr = self.value_range(ci); |
| 407 | if self.value_xor[vr].iter().any(|&b| b != 0) { |
| 408 | return false; |
| 409 | } |
| 410 | true |
| 411 | } |
| 412 | |
| 413 | fn extract_cell(&self, ci: usize) -> Outcome<(i32, Vec<u8>, Vec<u8>)> { |
| 414 | let sign = self.count[ci]; |
| 415 | if sign != 1 && sign != -1 { |
| 416 | return Err(err!( |
| 417 | "IBLT cell {} is not pure (count = {}).", ci, sign; |
| 418 | Invalid, Input, Bug)); |
| 419 | } |
| 420 | let kr = self.key_range(ci); |
| 421 | let key = self.key_xor[kr].to_vec(); |
| 422 | let vr = self.value_range(ci); |
| 423 | let value = self.value_xor[vr].to_vec(); |
| 424 | Ok((sign, key, value)) |
| 425 | } |
| 426 | |
| 427 | fn key_range(&self, ci: usize) -> std::ops::Range<usize> { |
| 428 | let start = ci * self.cfg.key_len; |
| 429 | start..start + self.cfg.key_len |
| 430 | } |
| 431 | |
| 432 | fn value_range(&self, ci: usize) -> std::ops::Range<usize> { |
| 433 | let start = ci * self.cfg.value_len; |
| 434 | start..start + self.cfg.value_len |
| 435 | } |
| 436 | } |