Oregami
Repositories/oxedyne/fe2o3

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
23use crate::cas::{
24 Cas,
25 Chunk,
26 ContentId,
27};
28
29use crate::O3db;
30use crate::base::id::usr_kind_id_deleted;
31
32use oxedyne_fe2o3_core::prelude::*;
33use oxedyne_fe2o3_iop_crypto::enc::Encrypter;
34use oxedyne_fe2o3_iop_db::api::ScanOpts;
35use oxedyne_fe2o3_iop_hash::{
36 api::Hasher,
37 csum::Checksummer,
38};
39use oxedyne_fe2o3_jdat::{
40 prelude::*,
41 id::NumIdDat,
42};
43
44use std::sync::Arc;
45
46
47// A scan on this prefix enumerates the store.
48const 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.
57pub 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
69impl<
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
150impl<
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)]
214mod 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}