Skip to main content

cobble_data_structure/
structured_read_only_db.rs

1use crate::structured_db::{
2    StructuredColumnValue, StructuredDbIterator, StructuredReadOptions, StructuredScanOptions,
3    StructuredSchema, combined_resolver, decode_row, load_structured_schema_from_cobble_schema,
4};
5use bytes::Bytes;
6use cobble::{Config, MergeOperatorResolver, ReadOnlyDb, ReadOnlyDbBuilder, Result};
7use std::ops::Range;
8use std::sync::Arc;
9
10pub struct StructuredReadOnlyDb {
11    db: ReadOnlyDb,
12    structured_schema: Arc<StructuredSchema>,
13    default_read_options: StructuredReadOptions,
14    default_scan_options: StructuredScanOptions,
15}
16
17/// Builder for a snapshot-backed structured database with runtime schema wiring.
18pub struct StructuredReadOnlyDbBuilder {
19    inner: ReadOnlyDbBuilder,
20}
21
22impl StructuredReadOnlyDbBuilder {
23    pub fn new(config: Config) -> Self {
24        Self {
25            inner: ReadOnlyDbBuilder::new(config).merge_operator_resolver(combined_resolver(None)),
26        }
27    }
28
29    pub fn db_id(mut self, db_id: impl Into<String>) -> Self {
30        self.inner = self.inner.db_id(db_id);
31        self
32    }
33
34    pub fn merge_operator_resolver(mut self, resolver: Arc<dyn MergeOperatorResolver>) -> Self {
35        self.inner = self
36            .inner
37            .merge_operator_resolver(combined_resolver(Some(resolver)));
38        self
39    }
40
41    /// Register a factory for raw single-column transform specifications.
42    pub fn register_schema_transform<F, T>(
43        mut self,
44        transform_type: impl Into<String>,
45        factory: F,
46    ) -> Result<Self>
47    where
48        F: Fn(&[u8]) -> Result<T> + Send + Sync + 'static,
49        T: Fn(Option<Bytes>) -> Result<Option<Bytes>> + Send + Sync + 'static,
50    {
51        self.inner = self
52            .inner
53            .register_schema_transform(transform_type, factory)?;
54        Ok(self)
55    }
56
57    pub fn open(self, snapshot_id: u64) -> Result<StructuredReadOnlyDb> {
58        StructuredReadOnlyDb::from_db(self.inner.open(snapshot_id)?)
59    }
60}
61
62impl StructuredReadOnlyDb {
63    fn from_db(db: ReadOnlyDb) -> Result<Self> {
64        let structured_schema = load_structured_schema_from_cobble_schema(&db.current_schema())?;
65        Ok(Self {
66            db,
67            structured_schema: Arc::new(structured_schema),
68            default_read_options: StructuredReadOptions::default(),
69            default_scan_options: StructuredScanOptions::default(),
70        })
71    }
72
73    pub fn open(config: Config, snapshot_id: u64, db_id: impl Into<String>) -> Result<Self> {
74        Self::open_with_resolver(config, snapshot_id, db_id, None)
75    }
76
77    pub fn open_with_resolver(
78        config: Config,
79        snapshot_id: u64,
80        db_id: impl Into<String>,
81        resolver: Option<Arc<dyn MergeOperatorResolver>>,
82    ) -> Result<Self> {
83        let db = ReadOnlyDb::open_with_db_id_and_resolver(
84            config,
85            snapshot_id,
86            db_id,
87            combined_resolver(resolver),
88        )?;
89        Self::from_db(db)
90    }
91
92    pub fn id(&self) -> &str {
93        self.db.id()
94    }
95
96    pub fn current_schema(&self) -> StructuredSchema {
97        self.structured_schema.as_ref().clone()
98    }
99
100    /// Register a factory for raw single-column transform specifications.
101    pub fn register_schema_transform<F, T>(
102        &self,
103        transform_type: impl Into<String>,
104        factory: F,
105    ) -> Result<()>
106    where
107        F: Fn(&[u8]) -> Result<T> + Send + Sync + 'static,
108        T: Fn(Option<Bytes>) -> Result<Option<Bytes>> + Send + Sync + 'static,
109    {
110        self.db.register_schema_transform(transform_type, factory)
111    }
112
113    pub fn get(
114        &self,
115        bucket: u16,
116        key: &[u8],
117    ) -> Result<Option<Vec<Option<StructuredColumnValue>>>> {
118        self.get_with_options(bucket, key, &self.default_read_options)
119    }
120
121    pub fn multi_get<K: AsRef<[u8]>>(
122        &self,
123        keys: &[(u16, K)],
124    ) -> Result<Vec<Option<Vec<Option<StructuredColumnValue>>>>> {
125        self.multi_get_with_options(keys, &self.default_read_options)
126    }
127
128    pub fn multi_get_with_options<K: AsRef<[u8]>>(
129        &self,
130        keys: &[(u16, K)],
131        options: &StructuredReadOptions,
132    ) -> Result<Vec<Option<Vec<Option<StructuredColumnValue>>>>> {
133        let raw_keys = keys
134            .iter()
135            .map(|(bucket, key)| (*bucket, key.as_ref()))
136            .collect::<Vec<_>>();
137        let projected_schema = options.resolve_projected_schema_cached(&self.structured_schema)?;
138        self.db
139            .multi_get_with_options(&raw_keys, options.as_cobble())?
140            .into_iter()
141            .map(|raw| {
142                raw.map(|columns| decode_row(&projected_schema, 0, columns))
143                    .transpose()
144            })
145            .collect()
146    }
147
148    pub fn get_with_options(
149        &self,
150        bucket: u16,
151        key: &[u8],
152        options: &StructuredReadOptions,
153    ) -> Result<Option<Vec<Option<StructuredColumnValue>>>> {
154        let raw = self.db.get_with_options(bucket, key, options.as_cobble())?;
155        let projected_schema = options.resolve_projected_schema_cached(&self.structured_schema)?;
156        raw.map(|columns| decode_row(&projected_schema, 0, columns))
157            .transpose()
158    }
159
160    pub fn scan(&self, bucket: u16, range: Range<&[u8]>) -> Result<StructuredDbIterator> {
161        self.scan_with_options(bucket, range, &self.default_scan_options)
162    }
163
164    pub fn scan_with_options(
165        &self,
166        bucket: u16,
167        range: Range<&[u8]>,
168        options: &StructuredScanOptions,
169    ) -> Result<StructuredDbIterator> {
170        let inner = self
171            .db
172            .scan_with_options(bucket, range, options.as_cobble())?;
173        let projected_schema = options.resolve_projected_schema_cached(&self.structured_schema)?;
174        Ok(StructuredDbIterator::new(inner, projected_schema, 0))
175    }
176}
177
178#[cfg(test)]
179#[path = "../tests/unit/structured_read_only_db.rs"]
180mod tests;