Skip to main content

cobble_data_structure/
structured_reader.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::{
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
16/// Builder for a structured snapshot reader with runtime schema wiring.
17pub 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    /// Register a factory for raw single-column transform specifications.
43    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    /// Opens a fixed view from an already resolved global snapshot without reloading its manifest.
72    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    // ── Read operations ─────────────────────────────────────────────────
139
140    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    // ── Snapshot management ─────────────────────────────────────────────
214
215    pub fn refresh(&mut self) -> Result<()> {
216        self.reader.refresh()?;
217        self.refresh_structured_schema_if_changed()
218    }
219
220    /// Register a factory for raw single-column transform specifications.
221    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
281/// Load the structured schema from the first shard of the reader's current global snapshot.
282fn 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;