Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_o3db_sync/src/dist/o3db_storage.rs

11.3 KiB, 41 runs

created by r1870400018:11529, 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//! [`Storage`] adapter backed by a local [`O3db`](crate::O3db) instance.
2//!
3//! Plumbs the distributed-Ozone [`Storage`] trait onto the existing local
4//! Ozone engine. Every record is persisted as one key/value pair:
5//!
6//! - Key: `Dat::Str("{table}:{hex_id}")`. The `{table}:` prefix lets a
7//! `scan` with [`ScanOpts::with_str_prefix`] enumerate every record in
8//! a single table, which the anti-entropy loop needs for its digest
9//! round. `{hex_id}` is the lowercase-hex encoding of the 32-byte
10//! [`RecordId`].
11//! - Value: `Dat::BU8(record.value)`. The record's table and id are
12//! recoverable from the key, so the stored value payload carries only
13//! the application-opaque bytes.
14//!
15//! Reads block on the responder channel through the synchronous
16//! [`OzoneApi::get_wait`] wrapper. Writes and deletes block on
17//! `Responder::recv_store_ack`, which waits for each record to be
18//! written and then durable, surfacing any error along the way as an
19//! [`Outcome`] error.
20//!
21//! # Performance notes
22//!
23//! - `digests` performs a prefix scan followed by a per-record fetch.
24//! Scan v1 in `fe2o3_o3db_sync` does not return values, so the
25//! content hash has to come from a second round trip per record.
26//! This is O(n) fetches per anti-entropy round; acceptable for small
27//! cohort-backed tables (identity, peer_set, revocation) but a
28//! bottleneck for large ones. A scan v2 that returns values is the
29//! planned optimisation; the [`Storage`] trait does not need to
30//! change to adopt it.
31//! - The content hash is the splitmix64-based 32-byte hash used by
32//! [`MemoryStorage`](super::storage::MemoryStorage) so that a mixed
33//! cluster (in-memory peer + O3db-backed peer) converges. A
34//! cryptographic replacement can slot in without touching the trait.
35//!
36//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
37//! Anthropic Claude
38
39use super::record::{
40 Record,
41 RecordDigest,
42 RecordId,
43};
44use super::storage::Storage;
45
46use crate::O3db;
47use crate::base::id::usr_kind_id_deleted;
48
49use oxedyne_fe2o3_core::prelude::*;
50use oxedyne_fe2o3_iop_crypto::enc::Encrypter;
51use oxedyne_fe2o3_iop_db::api::ScanOpts;
52use oxedyne_fe2o3_iop_hash::{
53 api::Hasher,
54 csum::Checksummer,
55};
56use oxedyne_fe2o3_jdat::{
57 prelude::*,
58 id::NumIdDat,
59};
60
61use std::sync::Arc;
62
63
64/// Storage adapter persisting records in a local [`O3db`] instance.
65///
66/// The adapter is generic over the six O3db type parameters so a caller
67/// can plug its chosen crypto and hash schemes in directly. In practice
68/// most callers instantiate it once at start-up with concrete types and
69/// then hand it to [`DistOzone::new`](super::engine::DistOzone::new).
70pub struct O3dbStorage<
71 const UIDL: usize,
72 UID: NumIdDat<UIDL> + 'static,
73 ENC: Encrypter + 'static,
74 KH: Hasher + 'static,
75 PR: Hasher + 'static,
76 CS: Checksummer + 'static,
77> {
78 db: Arc<O3db<UIDL, UID, ENC, KH, PR, CS>>,
79 user: UID,
80}
81
82impl<
83 const UIDL: usize,
84 UID: NumIdDat<UIDL> + 'static,
85 ENC: Encrypter + 'static,
86 KH: Hasher + 'static,
87 PR: Hasher + 'static,
88 CS: Checksummer + 'static,
89>
90 O3dbStorage<UIDL, UID, ENC, KH, PR, CS>
91{
92 /// `user` is the caller identity under which every distributed-mode
93 /// write is stamped -- distributed Ozone does not model per-write
94 /// authorship beyond this. Applications that need finer-grained
95 /// multi-user bookkeeping should layer that above the adapter.
96 pub fn new(db: Arc<O3db<UIDL, UID, ENC, KH, PR, CS>>, user: UID) -> Self {
97 Self { db, user }
98 }
99
100 /// Composite key format: `"{table}:{hex_id}"`.
101 fn encode_key(table: &str, id: &RecordId) -> Dat {
102 let mut s = String::with_capacity(table.len() + 1 + 64);
103 s.push_str(table);
104 s.push(':');
105 for b in id.as_bytes() {
106 s.push(hex_char(b >> 4));
107 s.push(hex_char(b & 0x0f));
108 }
109 Dat::Str(s)
110 }
111
112 /// Asserts the expected table prefix along the way.
113 fn parse_key(dat: &Dat, table: &str) -> Outcome<RecordId> {
114 let s = match dat {
115 Dat::Str(s) => s,
116 _ => return Err(err!(
117 "O3dbStorage expected Dat::Str key, got {:?}.", dat;
118 Invalid, Input, Mismatch)),
119 };
120 let expected_prefix_len = table.len() + 1;
121 if s.len() != expected_prefix_len + 64 {
122 return Err(err!(
123 "O3dbStorage key '{}' has unexpected length (table='{}').",
124 s, table;
125 Invalid, Input, Size));
126 }
127 if !s.starts_with(table) || s.as_bytes()[table.len()] != b':' {
128 return Err(err!(
129 "O3dbStorage key '{}' lacks expected '{}:' prefix.",
130 s, table;
131 Invalid, Input, Mismatch));
132 }
133 let hex = &s.as_bytes()[expected_prefix_len..];
134 let mut arr = [0u8; 32];
135 for i in 0..32 {
136 let hi = res!(nibble(hex[i * 2]));
137 let lo = res!(nibble(hex[i * 2 + 1]));
138 arr[i] = (hi << 4) | lo;
139 }
140 Ok(RecordId::from_bytes(arr))
141 }
142
143 /// Waits for the store's acknowledgement: its record count, each record written and
144 /// then durable. Whether the key existed is not needed here.
145 fn drain_store_ack(
146 resp: &crate::comm::response::Responder<UIDL, UID, ENC, KH>,
147 )
148 -> Outcome<()>
149 {
150 res!(resp.recv_store_ack());
151 Ok(())
152 }
153
154 /// Accepts the whole unsigned-bytes family: the store path always writes
155 /// `Dat::BU8`, but the other widths show up where a caller had
156 /// pre-existing data under a different width.
157 fn extract_value(dat: &Dat) -> Outcome<Vec<u8>> {
158 match dat {
159 Dat::BU8(b) | Dat::BU16(b) | Dat::BU32(b) | Dat::BU64(b) =>
160 Ok(b.clone()),
161 other => Err(err!(
162 "O3dbStorage expected a byte-vector Dat value, got {:?}.",
163 other;
164 Decode, Unexpected)),
165 }
166 }
167}
168
169impl<
170 const UIDL: usize,
171 UID: NumIdDat<UIDL> + 'static,
172 ENC: Encrypter + 'static,
173 KH: Hasher + 'static,
174 PR: Hasher + 'static,
175 CS: Checksummer + 'static,
176>
177 Storage for O3dbStorage<UIDL, UID, ENC, KH, PR, CS>
178{
179 fn put(&self, record: &Record) -> Outcome<()> {
180 let key = Self::encode_key(&record.table, &record.id);
181 let value = Dat::BU8(record.value.clone());
182 let resp = res!(self.db.api().store(key, value, self.user));
183 Self::drain_store_ack(&resp)
184 }
185
186 fn get(&self, table: &str, id: &RecordId) -> Outcome<Option<Record>> {
187 let key = Self::encode_key(table, id);
188 match res!(self.db.api().get_wait(&key, None)) {
189 None => Ok(None),
190 Some((dat, _meta)) => {
191 // A deletion tombstone is stored as
192 // `Dat::Usr(usr_kind_id_deleted(), _)`; `get_wait`
193 // returns it verbatim, so we have to recognise it
194 // here and surface "not present" to the caller.
195 if let Dat::Usr(kind, _) = &dat {
196 if *kind == usr_kind_id_deleted() {
197 return Ok(None);
198 }
199 }
200 let value = res!(Self::extract_value(&dat));
201 Ok(Some(Record {
202 id: *id,
203 table: table.to_string(),
204 value,
205 }))
206 }
207 }
208 }
209
210 fn delete(&self, table: &str, id: &RecordId) -> Outcome<bool> {
211 // Fetch first to determine whether a record was present; the
212 // overwrite below cannot distinguish "overwrote an existing
213 // record" from "wrote a fresh tombstone", and [`Storage::delete`]
214 // promises the former answer.
215 let existed = res!(self.get(table, id)).is_some();
216 if !existed {
217 return Ok(false);
218 }
219 // Overwrite the key with a tombstone via the ordinary `store`
220 // path. Going through `store` rather than
221 // `delete_using_responder` means the tombstone flows through
222 // the same encryption and cache pipeline as a regular write,
223 // so subsequent `get` calls see the tombstone on the first
224 // read rather than racing the deletion path's unencrypted
225 // direct-to-disk write.
226 let key = Self::encode_key(table, id);
227 let tombstone = Dat::Usr(
228 usr_kind_id_deleted(),
229 Some(Box::new(Dat::Empty)),
230 );
231 let resp = res!(self.db.api().store(key, tombstone, self.user));
232 res!(Self::drain_store_ack(&resp));
233 Ok(true)
234 }
235
236 fn digests(&self, table: &str) -> Outcome<Vec<RecordDigest>> {
237 let prefix = fmt!("{}:", table);
238 let opts = ScanOpts::with_str_prefix(prefix);
239 let entries = res!(self.db.api().scan(&opts, None));
240 let mut out = Vec::with_capacity(entries.len());
241 for (k, _v, _meta) in &entries {
242 let id = res!(Self::parse_key(k, table));
243 // Fetch each record so we can hash its value bytes. Scan v1
244 // does not return values, hence the per-record round trip.
245 let record = match res!(self.get(table, &id)) {
246 Some(r) => r,
247 None => continue, // Deleted between scan and fetch.
248 };
249 let content = content_hash(&record.value);
250 out.push(RecordDigest { id, content });
251 }
252 out.sort_by(|a, b| a.id.as_bytes().cmp(b.id.as_bytes()));
253 Ok(out)
254 }
255}
256
257
258/// A splitmix64-widened 32-byte hash, matching the one
259/// [`MemoryStorage`](super::storage::MemoryStorage) uses so that a mixed
260/// cluster of in-memory and O3db-backed peers does not disagree on digests.
261/// Not cryptographic; see the module doc comment.
262fn content_hash(bytes: &[u8]) -> [u8; 32] {
263 let mut state: u64 = 0x9E3779B97F4A7C15;
264 for chunk in bytes.chunks(8) {
265 let mut buf = [0u8; 8];
266 buf[..chunk.len()].copy_from_slice(chunk);
267 let word = u64::from_le_bytes(buf);
268 state = state.wrapping_add(word);
269 state = (state ^ (state >> 30)).wrapping_mul(0xBF58476D1CE4E5B9);
270 state = (state ^ (state >> 27)).wrapping_mul(0x94D049BB133111EB);
271 state ^= state >> 31;
272 }
273 let mut out = [0u8; 32];
274 for i in 0..4 {
275 let limb = state.wrapping_mul(0x9E3779B97F4A7C15 ^ (i as u64));
276 out[i * 8..(i + 1) * 8].copy_from_slice(&limb.to_le_bytes());
277 }
278 out
279}
280
281
282fn hex_char(nib: u8) -> char {
283 match nib {
284 0..=9 => (b'0' + nib) as char,
285 10..=15 => (b'a' + nib - 10) as char,
286 _ => '?', // Unreachable: caller masks to 0..=15.
287 }
288}
289
290fn nibble(b: u8) -> Outcome<u8> {
291 match b {
292 b'0'..=b'9' => Ok(b - b'0'),
293 b'a'..=b'f' => Ok(b - b'a' + 10),
294 b'A'..=b'F' => Ok(b - b'A' + 10),
295 _ => Err(err!(
296 "Invalid hex character: 0x{:02x}.", b;
297 Invalid, Input)),
298 }
299}
300
301
302#[cfg(test)]
303mod tests {
304 use super::*;
305
306 #[test]
307 fn key_round_trips() -> Outcome<()> {
308 let id = RecordId::from_bytes([0xab; 32]);
309 let dat = O3dbStorage::<16, oxedyne_fe2o3_jdat::id::IdDat<16, u128>,
310 oxedyne_fe2o3_crypto::enc::EncryptionScheme,
311 oxedyne_fe2o3_hash::hash::HashScheme,
312 oxedyne_fe2o3_hash::hash::HashScheme,
313 oxedyne_fe2o3_hash::csum::ChecksumScheme>::encode_key(
314 "identity", &id,
315 );
316 let back = res!(O3dbStorage::<16, oxedyne_fe2o3_jdat::id::IdDat<16, u128>,
317 oxedyne_fe2o3_crypto::enc::EncryptionScheme,
318 oxedyne_fe2o3_hash::hash::HashScheme,
319 oxedyne_fe2o3_hash::hash::HashScheme,
320 oxedyne_fe2o3_hash::csum::ChecksumScheme>::parse_key(
321 &dat, "identity",
322 ));
323 assert_eq!(back, id);
324 Ok(())
325 }
326
327 #[test]
328 fn content_hash_is_deterministic() {
329 assert_eq!(content_hash(b"hello"), content_hash(b"hello"));
330 assert_ne!(content_hash(b"hello"), content_hash(b"world"));
331 }
332
333 #[test]
334 fn hex_roundtrip() -> Outcome<()> {
335 for byte in 0u8..=255 {
336 let hi = hex_char(byte >> 4);
337 let lo = hex_char(byte & 0x0f);
338 let back = (res!(nibble(hi as u8)) << 4) | res!(nibble(lo as u8));
339 assert_eq!(back, byte);
340 }
341 Ok(())
342 }
343
344 #[test]
345 fn parse_key_rejects_wrong_prefix() {
346 let id = RecordId::from_bytes([0xcc; 32]);
347 let dat = O3dbStorage::<16, oxedyne_fe2o3_jdat::id::IdDat<16, u128>,
348 oxedyne_fe2o3_crypto::enc::EncryptionScheme,
349 oxedyne_fe2o3_hash::hash::HashScheme,
350 oxedyne_fe2o3_hash::hash::HashScheme,
351 oxedyne_fe2o3_hash::csum::ChecksumScheme>::encode_key(
352 "identity", &id,
353 );
354 assert!(O3dbStorage::<16, oxedyne_fe2o3_jdat::id::IdDat<16, u128>,
355 oxedyne_fe2o3_crypto::enc::EncryptionScheme,
356 oxedyne_fe2o3_hash::hash::HashScheme,
357 oxedyne_fe2o3_hash::hash::HashScheme,
358 oxedyne_fe2o3_hash::csum::ChecksumScheme>::parse_key(
359 &dat, "escrow",
360 ).is_err());
361 }
362}