oxedyne/fe2o3/fe2o3_o3db_sync/src/dist/storage.rs
4.5 KiB, 22 runs
created by r1870400018:11392, 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 storage abstraction that distributed Ozone sits on top of. |
| 2 | //! |
| 3 | //! [`Storage`] is the trait every concrete backend implements. The canonical |
| 4 | //! production backend is `fe2o3_o3db_sync` driving a per-peer local Ozone, |
| 5 | //! but the distributed-Ozone engine is agnostic -- any type that persists |
| 6 | //! `(table, id) -> Record` with digest enumeration for anti-entropy will do. |
| 7 | //! |
| 8 | //! This keeps the test suite honest: the tests in `tests/` use |
| 9 | //! [`MemoryStorage`], an in-memory `HashMap`-backed adapter, and exercise the |
| 10 | //! full engine without touching disk. A later commit wires the |
| 11 | //! `fe2o3_o3db_sync` adapter in an integration crate. |
| 12 | //! |
| 13 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 14 | //! Anthropic Claude |
| 15 | |
| 16 | use super::record::{ |
| 17 | Record, |
| 18 | RecordDigest, |
| 19 | RecordId, |
| 20 | }; |
| 21 | |
| 22 | use oxedyne_fe2o3_core::prelude::*; |
| 23 | |
| 24 | use std::collections::HashMap; |
| 25 | use std::sync::Mutex; |
| 26 | |
| 27 | |
| 28 | /// Implementations must be internally thread-safe: distributed Ozone holds |
| 29 | /// one `Storage` handle per engine and may access it from multiple threads |
| 30 | /// (the write path, the anti-entropy loop, the consensus cohort driver, and |
| 31 | /// the inbound envelope dispatcher). |
| 32 | pub trait Storage { |
| 33 | /// Overwrites any record already at the same identifier. |
| 34 | fn put(&self, record: &Record) -> Outcome<()>; |
| 35 | |
| 36 | fn get(&self, table: &str, id: &RecordId) -> Outcome<Option<Record>>; |
| 37 | |
| 38 | /// The bool reports whether a record was there to remove. |
| 39 | fn delete(&self, table: &str, id: &RecordId) -> Outcome<bool>; |
| 40 | |
| 41 | /// Used by the IBLT anti-entropy layer to build a symmetric-difference |
| 42 | /// sketch against a peer's view of the same table. The `content` hash |
| 43 | /// on each digest must be deterministic: two peers that hold the same |
| 44 | /// record bytes must produce the same `content`. |
| 45 | fn digests(&self, table: &str) -> Outcome<Vec<RecordDigest>>; |
| 46 | } |
| 47 | |
| 48 | |
| 49 | /// A simple in-memory [`Storage`] backed by a `HashMap`. Intended for tests, |
| 50 | /// two-peer loopback demos, and documentation examples -- not for |
| 51 | /// production. |
| 52 | /// |
| 53 | /// Thread safety is via a single `Mutex` around the inner state; contention |
| 54 | /// is irrelevant at test scale. |
| 55 | pub struct MemoryStorage { |
| 56 | inner: Mutex<HashMap<(String, [u8; 32]), Record>>, |
| 57 | } |
| 58 | |
| 59 | impl MemoryStorage { |
| 60 | pub fn new() -> Self { |
| 61 | Self { inner: Mutex::new(HashMap::new()) } |
| 62 | } |
| 63 | |
| 64 | /// Across all tables. |
| 65 | pub fn len(&self) -> Outcome<usize> { |
| 66 | let guard = lock_mutex!(self.inner); |
| 67 | Ok(guard.len()) |
| 68 | } |
| 69 | } |
| 70 | |
| 71 | impl Default for MemoryStorage { |
| 72 | fn default() -> Self { |
| 73 | Self::new() |
| 74 | } |
| 75 | } |
| 76 | |
| 77 | impl Storage for MemoryStorage { |
| 78 | fn put(&self, record: &Record) -> Outcome<()> { |
| 79 | let mut guard = lock_mutex!(self.inner); |
| 80 | let key = (record.table.clone(), *record.id.as_bytes()); |
| 81 | guard.insert(key, record.clone()); |
| 82 | Ok(()) |
| 83 | } |
| 84 | |
| 85 | fn get(&self, table: &str, id: &RecordId) -> Outcome<Option<Record>> { |
| 86 | let guard = lock_mutex!(self.inner); |
| 87 | let key = (table.to_string(), *id.as_bytes()); |
| 88 | Ok(guard.get(&key).cloned()) |
| 89 | } |
| 90 | |
| 91 | fn delete(&self, table: &str, id: &RecordId) -> Outcome<bool> { |
| 92 | let mut guard = lock_mutex!(self.inner); |
| 93 | let key = (table.to_string(), *id.as_bytes()); |
| 94 | Ok(guard.remove(&key).is_some()) |
| 95 | } |
| 96 | |
| 97 | fn digests(&self, table: &str) -> Outcome<Vec<RecordDigest>> { |
| 98 | let guard = lock_mutex!(self.inner); |
| 99 | let mut out = Vec::new(); |
| 100 | for ((t, _), rec) in guard.iter() { |
| 101 | if t == table { |
| 102 | let content = content_hash(&rec.value); |
| 103 | out.push(RecordDigest { id: rec.id, content }); |
| 104 | } |
| 105 | } |
| 106 | // Sort for deterministic iteration in tests. The digests call site |
| 107 | // already feeds an IBLT, which is order-independent, so sorting here |
| 108 | // is purely for test reproducibility. |
| 109 | out.sort_by(|a, b| a.id.as_bytes().cmp(b.id.as_bytes())); |
| 110 | Ok(out) |
| 111 | } |
| 112 | } |
| 113 | |
| 114 | |
| 115 | /// The deterministic 256-bit hash behind [`RecordDigest::content`]. |
| 116 | /// |
| 117 | /// This is *not* a cryptographic hash -- it is splitmix64-based and intended |
| 118 | /// only for test adapters. Production storage backends should use a proper |
| 119 | /// hash (SHA-3, BLAKE3) that is resistant to adversarial collisions. |
| 120 | fn content_hash(bytes: &[u8]) -> [u8; 32] { |
| 121 | let mut state: u64 = 0x9E3779B97F4A7C15; |
| 122 | for chunk in bytes.chunks(8) { |
| 123 | let mut buf = [0u8; 8]; |
| 124 | buf[..chunk.len()].copy_from_slice(chunk); |
| 125 | let word = u64::from_le_bytes(buf); |
| 126 | state = state.wrapping_add(word); |
| 127 | state = (state ^ (state >> 30)).wrapping_mul(0xBF58476D1CE4E5B9); |
| 128 | state = (state ^ (state >> 27)).wrapping_mul(0x94D049BB133111EB); |
| 129 | state ^= state >> 31; |
| 130 | } |
| 131 | let mut out = [0u8; 32]; |
| 132 | for i in 0..4 { |
| 133 | let limb = state.wrapping_mul(0x9E3779B97F4A7C15 ^ (i as u64)); |
| 134 | out[i * 8..(i + 1) * 8].copy_from_slice(&limb.to_le_bytes()); |
| 135 | } |
| 136 | out |
| 137 | } |