Oregami
Repositories/oxedyne/fe2o3

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
16use super::record::{
17 Record,
18 RecordDigest,
19 RecordId,
20};
21
22use oxedyne_fe2o3_core::prelude::*;
23
24use std::collections::HashMap;
25use 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).
32pub 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.
55pub struct MemoryStorage {
56 inner: Mutex<HashMap<(String, [u8; 32]), Record>>,
57}
58
59impl 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
71impl Default for MemoryStorage {
72 fn default() -> Self {
73 Self::new()
74 }
75}
76
77impl 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.
120fn 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}