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
17pub 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 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 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;