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::{
7 Config, GlobalSnapshotManifest, GlobalSnapshotSummary, MergeOperatorResolver,
8 ReadOnlyDbBuilder, Reader, ReaderBuilder, ReaderConfig, Result, VolumeDescriptor,
9};
10use std::ops::Range;
11use std::sync::Arc;
12
13type SchemaTransformCallback = Arc<dyn Fn(Option<Bytes>) -> Result<Option<Bytes>> + Send + Sync>;
14type SchemaTransformFactory = Arc<dyn Fn(&[u8]) -> Result<SchemaTransformCallback> + Send + Sync>;
15
16pub struct StructuredReaderBuilder {
18 inner: ReaderBuilder,
19 volumes: Vec<VolumeDescriptor>,
20 resolver: Option<Arc<dyn MergeOperatorResolver>>,
21 factories: Vec<(String, SchemaTransformFactory)>,
22}
23
24impl StructuredReaderBuilder {
25 pub fn new(config: ReaderConfig) -> Self {
26 Self {
27 volumes: config.volumes.clone(),
28 inner: ReaderBuilder::new(config).merge_operator_resolver(combined_resolver(None)),
29 resolver: None,
30 factories: Vec::new(),
31 }
32 }
33
34 pub fn merge_operator_resolver(mut self, resolver: Arc<dyn MergeOperatorResolver>) -> Self {
35 self.resolver = Some(resolver);
36 self.inner = self
37 .inner
38 .merge_operator_resolver(combined_resolver(self.resolver.clone()));
39 self
40 }
41
42 pub fn register_schema_transform<F, T>(
44 mut self,
45 transform_type: impl Into<String>,
46 factory: F,
47 ) -> Result<Self>
48 where
49 F: Fn(&[u8]) -> Result<T> + Send + Sync + 'static,
50 T: Fn(Option<Bytes>) -> Result<Option<Bytes>> + Send + Sync + 'static,
51 {
52 let transform_type = transform_type.into();
53 let factory: SchemaTransformFactory = Arc::new(move |spec| Ok(Arc::new(factory(spec)?)));
54 self.inner = self
55 .inner
56 .register_schema_transform(transform_type.clone(), {
57 let factory = Arc::clone(&factory);
58 move |spec| {
59 let callback = factory(spec)?;
60 Ok(move |value| callback(value))
61 }
62 })?;
63 self.factories.push((transform_type, factory));
64 Ok(self)
65 }
66
67 pub fn open(self, global_snapshot_id: u64) -> Result<StructuredReader> {
68 self.open_inner(Some(global_snapshot_id))
69 }
70
71 pub fn open_from_global_snapshot(
73 self,
74 global_snapshot: GlobalSnapshotManifest,
75 ) -> Result<StructuredReader> {
76 let reader = self.inner.open_from_global_snapshot(global_snapshot)?;
77 StructuredReader::from_reader(reader, self.volumes, self.resolver, self.factories)
78 }
79
80 pub fn open_current(self) -> Result<StructuredReader> {
81 self.open_inner(None)
82 }
83
84 fn open_inner(self, global_snapshot_id: Option<u64>) -> Result<StructuredReader> {
85 let reader = match global_snapshot_id {
86 Some(snapshot_id) => self.inner.open(snapshot_id)?,
87 None => self.inner.open_current()?,
88 };
89 StructuredReader::from_reader(reader, self.volumes, self.resolver, self.factories)
90 }
91}
92
93pub struct StructuredReader {
94 reader: Reader,
95 structured_schema: Arc<StructuredSchema>,
96 schema_snapshot_id: u64,
97 volumes: Vec<VolumeDescriptor>,
98 resolver: Option<Arc<dyn MergeOperatorResolver>>,
99 factories: Vec<(String, SchemaTransformFactory)>,
100 default_read_options: StructuredReadOptions,
101 default_scan_options: StructuredScanOptions,
102}
103
104impl StructuredReader {
105 pub fn open(read_config: ReaderConfig, global_snapshot_id: u64) -> Result<Self> {
106 StructuredReaderBuilder::new(read_config).open(global_snapshot_id)
107 }
108
109 pub fn open_current(read_config: ReaderConfig) -> Result<Self> {
110 StructuredReaderBuilder::new(read_config).open_current()
111 }
112
113 fn from_reader(
114 reader: Reader,
115 volumes: Vec<VolumeDescriptor>,
116 resolver: Option<Arc<dyn MergeOperatorResolver>>,
117 factories: Vec<(String, SchemaTransformFactory)>,
118 ) -> Result<Self> {
119 let structured_schema =
120 load_schema_from_reader(&reader, &volumes, resolver.as_ref(), &factories)?;
121 let schema_snapshot_id = reader.current_global_snapshot().id;
122 Ok(Self {
123 reader,
124 structured_schema: Arc::new(structured_schema),
125 schema_snapshot_id,
126 volumes,
127 resolver,
128 factories,
129 default_read_options: StructuredReadOptions::default(),
130 default_scan_options: StructuredScanOptions::default(),
131 })
132 }
133
134 pub fn current_schema(&self) -> StructuredSchema {
135 self.structured_schema.as_ref().clone()
136 }
137
138 pub fn get(
141 &mut self,
142 bucket_id: u16,
143 key: &[u8],
144 ) -> Result<Option<Vec<Option<StructuredColumnValue>>>> {
145 let default_options = self.default_read_options.clone();
146 self.get_with_options(bucket_id, key, &default_options)
147 }
148
149 pub fn multi_get<K: AsRef<[u8]>>(
150 &mut self,
151 keys: &[(u16, K)],
152 ) -> Result<Vec<Option<Vec<Option<StructuredColumnValue>>>>> {
153 let default_options = self.default_read_options.clone();
154 self.multi_get_with_options(keys, &default_options)
155 }
156
157 pub fn multi_get_with_options<K: AsRef<[u8]>>(
158 &mut self,
159 keys: &[(u16, K)],
160 options: &StructuredReadOptions,
161 ) -> Result<Vec<Option<Vec<Option<StructuredColumnValue>>>>> {
162 let raw_keys = keys
163 .iter()
164 .map(|(bucket, key)| (*bucket, key.as_ref()))
165 .collect::<Vec<_>>();
166 let raw = self
167 .reader
168 .multi_get_with_options(&raw_keys, options.as_cobble())?;
169 self.refresh_structured_schema_if_changed()?;
170 let projected_schema = options.resolve_projected_schema_cached(&self.structured_schema)?;
171 raw.into_iter()
172 .map(|raw| {
173 raw.map(|columns| decode_row(&projected_schema, 0, columns))
174 .transpose()
175 })
176 .collect()
177 }
178
179 pub fn get_with_options(
180 &mut self,
181 bucket_id: u16,
182 key: &[u8],
183 options: &StructuredReadOptions,
184 ) -> Result<Option<Vec<Option<StructuredColumnValue>>>> {
185 let raw = self
186 .reader
187 .get_with_options(bucket_id, key, options.as_cobble())?;
188 self.refresh_structured_schema_if_changed()?;
189 let projected_schema = options.resolve_projected_schema_cached(&self.structured_schema)?;
190 raw.map(|columns| decode_row(&projected_schema, 0, columns))
191 .transpose()
192 }
193
194 pub fn scan(&mut self, bucket_id: u16, range: Range<&[u8]>) -> Result<StructuredDbIterator> {
195 let default_options = self.default_scan_options.clone();
196 self.scan_with_options(bucket_id, range, &default_options)
197 }
198
199 pub fn scan_with_options(
200 &mut self,
201 bucket_id: u16,
202 range: Range<&[u8]>,
203 options: &StructuredScanOptions,
204 ) -> Result<StructuredDbIterator> {
205 let inner = self
206 .reader
207 .scan_with_options(bucket_id, range, options.as_cobble())?;
208 self.refresh_structured_schema_if_changed()?;
209 let projected_schema = options.resolve_projected_schema_cached(&self.structured_schema)?;
210 Ok(StructuredDbIterator::new(inner, projected_schema, 0))
211 }
212
213 pub fn refresh(&mut self) -> Result<()> {
216 self.reader.refresh()?;
217 self.refresh_structured_schema_if_changed()
218 }
219
220 pub fn register_schema_transform<F, T>(
222 &mut self,
223 transform_type: impl Into<String>,
224 factory: F,
225 ) -> Result<()>
226 where
227 F: Fn(&[u8]) -> Result<T> + Send + Sync + 'static,
228 T: Fn(Option<Bytes>) -> Result<Option<Bytes>> + Send + Sync + 'static,
229 {
230 let transform_type = transform_type.into();
231 let factory: SchemaTransformFactory = Arc::new(move |spec| Ok(Arc::new(factory(spec)?)));
232 self.reader
233 .register_schema_transform(transform_type.clone(), {
234 let factory = Arc::clone(&factory);
235 move |spec| {
236 let callback = factory(spec)?;
237 Ok(move |value| callback(value))
238 }
239 })?;
240 self.factories.push((transform_type, factory));
241 Ok(())
242 }
243
244 pub fn read_mode(&self) -> &'static str {
245 self.reader.read_mode()
246 }
247
248 pub fn configured_snapshot_id(&self) -> Option<u64> {
249 self.reader.configured_snapshot_id()
250 }
251
252 pub fn current_global_snapshot(&self) -> &GlobalSnapshotManifest {
253 self.reader.current_global_snapshot()
254 }
255
256 pub fn list_global_snapshots(&self) -> Result<Vec<GlobalSnapshotSummary>> {
257 self.reader.list_global_snapshots()
258 }
259
260 pub fn list_global_snapshot_manifests(&self) -> Result<Vec<GlobalSnapshotManifest>> {
261 self.reader.list_global_snapshot_manifests()
262 }
263
264 fn refresh_structured_schema_if_changed(&mut self) -> Result<()> {
265 let snapshot_id = self.reader.current_global_snapshot().id;
266 if snapshot_id == self.schema_snapshot_id {
267 return Ok(());
268 }
269 let schema = load_schema_from_reader(
270 &self.reader,
271 &self.volumes,
272 self.resolver.as_ref(),
273 &self.factories,
274 )?;
275 self.structured_schema = Arc::new(schema);
276 self.schema_snapshot_id = snapshot_id;
277 Ok(())
278 }
279}
280
281fn load_schema_from_reader(
283 reader: &Reader,
284 volumes: &[VolumeDescriptor],
285 resolver: Option<&Arc<dyn cobble::MergeOperatorResolver>>,
286 factories: &[(String, SchemaTransformFactory)],
287) -> Result<StructuredSchema> {
288 let manifest = reader.current_global_snapshot();
289 let shard = manifest.shard_snapshots.first().ok_or_else(|| {
290 cobble::Error::ConfigError("global snapshot has no shard snapshots".to_string())
291 })?;
292 let config = Config {
293 volumes: volumes.to_vec(),
294 total_buckets: manifest.total_buckets,
295 ..cobble::Config::default()
296 };
297 let mut builder = ReadOnlyDbBuilder::new(config)
298 .db_id(shard.db_id.clone())
299 .merge_operator_resolver(combined_resolver(resolver.cloned()));
300 for (transform_id, factory) in factories {
301 let transform_id = transform_id.clone();
302 let factory = Arc::clone(factory);
303 builder = builder.register_schema_transform(transform_id, move |spec| {
304 let callback = factory(spec)?;
305 Ok(move |value| callback(value))
306 })?;
307 }
308 let read_only = builder.open(shard.snapshot_id)?;
309 load_structured_schema_from_cobble_schema(&read_only.current_schema())
310}
311
312#[cfg(test)]
313#[path = "../tests/unit/structured_reader.rs"]
314mod tests;