oxedyne/fe2o3/fe2o3_o3db_sync/src/cas_o3db.rs
6.9 KiB, 31 runs
created by r1870400018:16779, 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 | //! [`Cas`] adapter backed by a local [`O3db`](crate::O3db) instance. |
| 2 | //! |
| 3 | //! Plumbs the content-addressed [`Cas`] trait onto the local Ozone engine, so |
| 4 | //! a caller that already runs an O3db (a gateway, say) gains a chunk store |
| 5 | //! without a second storage system. Every chunk is one key/value pair: |
| 6 | //! |
| 7 | //! - Key: `Dat::Str("chunk:{hex_addr}")`. The `chunk:` prefix lets a `scan` |
| 8 | //! with [`ScanOpts::with_str_prefix`] enumerate every chunk for |
| 9 | //! [`Cas::ids`], which garbage collection needs. `{hex_addr}` is the |
| 10 | //! lowercase-hex [`ContentId`]. |
| 11 | //! - Value: `Dat::BU8(chunk.bytes)`. The address is recoverable from the key, |
| 12 | //! so the value carries only the opaque chunk bytes -- ciphertext, when the |
| 13 | //! caller has encrypted before storing. |
| 14 | //! |
| 15 | //! Deletes are tombstones flowing through the ordinary store path, matching |
| 16 | //! `O3dbStorage`, behind the `dist` feature, so a subsequent `get` sees the |
| 17 | //! tombstone on its first read rather than racing a direct-to-disk delete. The |
| 18 | //! log-structured engine reclaims the space on compaction. |
| 19 | //! |
| 20 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 21 | //! Anthropic Claude |
| 22 | |
| 23 | use crate::cas::{ |
| 24 | Cas, |
| 25 | Chunk, |
| 26 | ContentId, |
| 27 | }; |
| 28 | |
| 29 | use crate::O3db; |
| 30 | use crate::base::id::usr_kind_id_deleted; |
| 31 | |
| 32 | use oxedyne_fe2o3_core::prelude::*; |
| 33 | use oxedyne_fe2o3_iop_crypto::enc::Encrypter; |
| 34 | use oxedyne_fe2o3_iop_db::api::ScanOpts; |
| 35 | use oxedyne_fe2o3_iop_hash::{ |
| 36 | api::Hasher, |
| 37 | csum::Checksummer, |
| 38 | }; |
| 39 | use oxedyne_fe2o3_jdat::{ |
| 40 | prelude::*, |
| 41 | id::NumIdDat, |
| 42 | }; |
| 43 | |
| 44 | use std::sync::Arc; |
| 45 | |
| 46 | |
| 47 | // A scan on this prefix enumerates the store. |
| 48 | const CHUNK_PREFIX: &str = "chunk:"; |
| 49 | |
| 50 | |
| 51 | /// Content-addressed chunk store persisting chunks in a local [`O3db`]. |
| 52 | /// |
| 53 | /// Generic over the six O3db type parameters so a caller plugs its chosen |
| 54 | /// crypto, hash and checksum schemes in directly, exactly as |
| 55 | /// `O3dbStorage`, behind the `dist` feature, does. Most callers |
| 56 | /// construct one at start-up with concrete types. |
| 57 | pub struct O3dbCas< |
| 58 | const UIDL: usize, |
| 59 | UID: NumIdDat<UIDL> + 'static, |
| 60 | ENC: Encrypter + 'static, |
| 61 | KH: Hasher + 'static, |
| 62 | PR: Hasher + 'static, |
| 63 | CS: Checksummer + 'static, |
| 64 | > { |
| 65 | db: Arc<O3db<UIDL, UID, ENC, KH, PR, CS>>, |
| 66 | user: UID, |
| 67 | } |
| 68 | |
| 69 | impl< |
| 70 | const UIDL: usize, |
| 71 | UID: NumIdDat<UIDL> + 'static, |
| 72 | ENC: Encrypter + 'static, |
| 73 | KH: Hasher + 'static, |
| 74 | PR: Hasher + 'static, |
| 75 | CS: Checksummer + 'static, |
| 76 | > |
| 77 | O3dbCas<UIDL, UID, ENC, KH, PR, CS> |
| 78 | { |
| 79 | /// `user` is the caller identity every write is stamped with; the chunk |
| 80 | /// store does not model per-chunk authorship beyond this. |
| 81 | pub fn new(db: Arc<O3db<UIDL, UID, ENC, KH, PR, CS>>, user: UID) -> Self { |
| 82 | Self { db, user } |
| 83 | } |
| 84 | |
| 85 | /// Composite key format: `"chunk:{hex_addr}"`. |
| 86 | fn encode_key(id: &ContentId) -> Dat { |
| 87 | let mut s = String::with_capacity(CHUNK_PREFIX.len() + 64); |
| 88 | s.push_str(CHUNK_PREFIX); |
| 89 | s.push_str(&id.to_hex()); |
| 90 | Dat::Str(s) |
| 91 | } |
| 92 | |
| 93 | fn parse_key(dat: &Dat) -> Outcome<ContentId> { |
| 94 | let s = match dat { |
| 95 | Dat::Str(s) => s, |
| 96 | _ => return Err(err!( |
| 97 | "O3dbCas expected a Dat::Str key, got {:?}.", dat; |
| 98 | Invalid, Input, Mismatch)), |
| 99 | }; |
| 100 | if !s.starts_with(CHUNK_PREFIX) { |
| 101 | return Err(err!( |
| 102 | "O3dbCas key '{}' lacks the '{}' prefix.", s, CHUNK_PREFIX; |
| 103 | Invalid, Input, Mismatch)); |
| 104 | } |
| 105 | ContentId::from_hex(&s[CHUNK_PREFIX.len()..]) |
| 106 | } |
| 107 | |
| 108 | /// Waits for the store's acknowledgement: its record count, each record written and |
| 109 | /// then durable. Whether the key existed is not needed here. |
| 110 | fn drain_store_ack( |
| 111 | resp: &crate::comm::response::Responder<UIDL, UID, ENC, KH>, |
| 112 | ) |
| 113 | -> Outcome<()> |
| 114 | { |
| 115 | res!(resp.recv_store_ack()); |
| 116 | Ok(()) |
| 117 | } |
| 118 | |
| 119 | /// Accepts the whole unsigned-bytes family; the store path always writes |
| 120 | /// `Dat::BU8`. |
| 121 | fn extract_bytes(dat: &Dat) -> Outcome<Vec<u8>> { |
| 122 | match dat { |
| 123 | Dat::BU8(b) | Dat::BU16(b) | Dat::BU32(b) | Dat::BU64(b) => |
| 124 | Ok(b.clone()), |
| 125 | other => Err(err!( |
| 126 | "O3dbCas expected a byte-vector Dat value, got {:?}.", other; |
| 127 | Decode, Unexpected)), |
| 128 | } |
| 129 | } |
| 130 | |
| 131 | /// `None` for an absent key or a deletion tombstone. |
| 132 | fn read(&self, id: &ContentId) |
| 133 | -> Outcome<Option<Vec<u8>>> |
| 134 | { |
| 135 | let key = Self::encode_key(id); |
| 136 | match res!(self.db.api().get_wait(&key, None)) { |
| 137 | None => Ok(None), |
| 138 | Some((dat, _meta)) => { |
| 139 | if let Dat::Usr(kind, _) = &dat { |
| 140 | if *kind == usr_kind_id_deleted() { |
| 141 | return Ok(None); |
| 142 | } |
| 143 | } |
| 144 | Ok(Some(res!(Self::extract_bytes(&dat)))) |
| 145 | } |
| 146 | } |
| 147 | } |
| 148 | } |
| 149 | |
| 150 | impl< |
| 151 | const UIDL: usize, |
| 152 | UID: NumIdDat<UIDL> + 'static, |
| 153 | ENC: Encrypter + 'static, |
| 154 | KH: Hasher + 'static, |
| 155 | PR: Hasher + 'static, |
| 156 | CS: Checksummer + 'static, |
| 157 | > |
| 158 | Cas for O3dbCas<UIDL, UID, ENC, KH, PR, CS> |
| 159 | { |
| 160 | fn put(&self, chunk: &Chunk) -> Outcome<()> { |
| 161 | if !chunk.id.verifies(&chunk.bytes) { |
| 162 | return Err(err!( |
| 163 | "Refusing chunk whose bytes do not hash to its address {}.", |
| 164 | chunk.id; |
| 165 | Invalid, Input, Mismatch)); |
| 166 | } |
| 167 | let key = Self::encode_key(&chunk.id); |
| 168 | let value = Dat::BU8(chunk.bytes.clone()); |
| 169 | let resp = res!(self.db.api().store(key, value, self.user)); |
| 170 | Self::drain_store_ack(&resp) |
| 171 | } |
| 172 | |
| 173 | fn get(&self, id: &ContentId) -> Outcome<Option<Vec<u8>>> { |
| 174 | self.read(id) |
| 175 | } |
| 176 | |
| 177 | fn has(&self, id: &ContentId) -> Outcome<bool> { |
| 178 | Ok(res!(self.read(id)).is_some()) |
| 179 | } |
| 180 | |
| 181 | fn delete(&self, id: &ContentId) -> Outcome<bool> { |
| 182 | let existed = res!(self.read(id)).is_some(); |
| 183 | if !existed { |
| 184 | return Ok(false); |
| 185 | } |
| 186 | let key = Self::encode_key(id); |
| 187 | let tombstone = Dat::Usr( |
| 188 | usr_kind_id_deleted(), |
| 189 | Some(Box::new(Dat::Empty)), |
| 190 | ); |
| 191 | let resp = res!(self.db.api().store(key, tombstone, self.user)); |
| 192 | res!(Self::drain_store_ack(&resp)); |
| 193 | Ok(true) |
| 194 | } |
| 195 | |
| 196 | fn ids(&self) -> Outcome<Vec<ContentId>> { |
| 197 | let opts = ScanOpts::with_str_prefix(CHUNK_PREFIX.to_string()); |
| 198 | let entries = res!(self.db.api().scan(&opts, None)); |
| 199 | let mut out = Vec::with_capacity(entries.len()); |
| 200 | for (k, _v, _meta) in &entries { |
| 201 | let id = res!(Self::parse_key(k)); |
| 202 | // Skip a key whose value is a tombstone: scan enumerates the key |
| 203 | // space, and a deleted chunk lingers until compaction. |
| 204 | if res!(self.read(&id)).is_some() { |
| 205 | out.push(id); |
| 206 | } |
| 207 | } |
| 208 | Ok(out) |
| 209 | } |
| 210 | } |
| 211 | |
| 212 | |
| 213 | #[cfg(test)] |
| 214 | mod tests { |
| 215 | use super::*; |
| 216 | |
| 217 | #[test] |
| 218 | fn key_round_trips() -> Outcome<()> { |
| 219 | let id = ContentId::of(b"a chunk of bytes"); |
| 220 | let dat = O3dbCas::<16, oxedyne_fe2o3_jdat::id::IdDat<16, u128>, |
| 221 | oxedyne_fe2o3_crypto::enc::EncryptionScheme, |
| 222 | oxedyne_fe2o3_hash::hash::HashScheme, |
| 223 | oxedyne_fe2o3_hash::hash::HashScheme, |
| 224 | oxedyne_fe2o3_hash::csum::ChecksumScheme>::encode_key(&id); |
| 225 | let back = res!(O3dbCas::<16, oxedyne_fe2o3_jdat::id::IdDat<16, u128>, |
| 226 | oxedyne_fe2o3_crypto::enc::EncryptionScheme, |
| 227 | oxedyne_fe2o3_hash::hash::HashScheme, |
| 228 | oxedyne_fe2o3_hash::hash::HashScheme, |
| 229 | oxedyne_fe2o3_hash::csum::ChecksumScheme>::parse_key(&dat)); |
| 230 | assert_eq!(back, id); |
| 231 | Ok(()) |
| 232 | } |
| 233 | |
| 234 | #[test] |
| 235 | fn parse_key_rejects_wrong_prefix() { |
| 236 | let dat = Dat::Str("other:deadbeef".to_string()); |
| 237 | assert!(O3dbCas::<16, oxedyne_fe2o3_jdat::id::IdDat<16, u128>, |
| 238 | oxedyne_fe2o3_crypto::enc::EncryptionScheme, |
| 239 | oxedyne_fe2o3_hash::hash::HashScheme, |
| 240 | oxedyne_fe2o3_hash::hash::HashScheme, |
| 241 | oxedyne_fe2o3_hash::csum::ChecksumScheme>::parse_key(&dat).is_err()); |
| 242 | } |
| 243 | } |