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 | |
| 39 | use super::record::{ |
| 40 | Record, |
| 41 | RecordDigest, |
| 42 | RecordId, |
| 43 | }; |
| 44 | use super::storage::Storage; |
| 45 | |
| 46 | use crate::O3db; |
| 47 | use crate::base::id::usr_kind_id_deleted; |
| 48 | |
| 49 | use oxedyne_fe2o3_core::prelude::*; |
| 50 | use oxedyne_fe2o3_iop_crypto::enc::Encrypter; |
| 51 | use oxedyne_fe2o3_iop_db::api::ScanOpts; |
| 52 | use oxedyne_fe2o3_iop_hash::{ |
| 53 | api::Hasher, |
| 54 | csum::Checksummer, |
| 55 | }; |
| 56 | use oxedyne_fe2o3_jdat::{ |
| 57 | prelude::*, |
| 58 | id::NumIdDat, |
| 59 | }; |
| 60 | |
| 61 | use 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). |
| 70 | pub 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 | |
| 82 | impl< |
| 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 | |
| 169 | impl< |
| 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. |
| 262 | fn 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 | |
| 282 | fn 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 | |
| 290 | fn 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)] |
| 303 | mod 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 | } |