Skip to main content

cobble_data_structure/
structured_single_db.rs

1#[cfg(feature = "ffi")]
2use crate::list::encode_borrowed_list_for_write;
3use crate::priority_queue::{
4    PriorityQueue, priority_queue_column_family_name, priority_queue_column_family_options,
5    validate_priority_queue_column_family,
6};
7use crate::structured_db::ensure_bytes_column;
8#[cfg(feature = "ffi")]
9use crate::structured_db::list_column_config;
10use crate::structured_db::{
11    StructuredColumnValue, StructuredDbIterator, StructuredReadOptions, StructuredScanOptions,
12    StructuredSchema, StructuredSchemaBuilder, StructuredSchemaOwner, StructuredWriteBatch,
13    StructuredWriteOptions, decode_row, encode_for_write,
14    load_structured_schema_from_cobble_schema, persist_structured_schema_on_db,
15};
16use bytes::Bytes;
17use cobble::{Config, DbIterator, Error, MemtableType, Result, SingleDb};
18use std::ops::Range;
19use std::sync::Arc;
20
21pub struct StructuredSingleDb {
22    db: SingleDb,
23    structured_schema: Arc<StructuredSchema>,
24    default_write_options: StructuredWriteOptions,
25    default_read_options: StructuredReadOptions,
26    default_scan_options: StructuredScanOptions,
27}
28
29impl StructuredSingleDb {
30    pub fn open(config: Config) -> Result<Self> {
31        let db = SingleDb::open(config)?;
32        let structured_schema =
33            load_structured_schema_from_cobble_schema(&db.db().current_schema())?;
34        Ok(Self {
35            db,
36            structured_schema: Arc::new(structured_schema),
37            default_write_options: StructuredWriteOptions::default(),
38            default_read_options: StructuredReadOptions::default(),
39            default_scan_options: StructuredScanOptions::default(),
40        })
41    }
42
43    pub fn db(&self) -> &SingleDb {
44        &self.db
45    }
46
47    #[cfg(feature = "ffi")]
48    pub(crate) fn jni_direct_buffer_pool_config(&self) -> Result<(usize, usize)> {
49        cobble::ffi::db_direct_buffer_pool_config(self.db.db())
50    }
51
52    pub fn current_schema(&self) -> StructuredSchema {
53        self.structured_schema.as_ref().clone()
54    }
55
56    /// Register a factory for raw single-column transform specifications.
57    pub fn register_schema_transform<F, T>(
58        &self,
59        transform_type: impl Into<String>,
60        factory: F,
61    ) -> Result<()>
62    where
63        F: Fn(&[u8]) -> Result<T> + Send + Sync + 'static,
64        T: Fn(Option<Bytes>) -> Result<Option<Bytes>> + Send + Sync + 'static,
65    {
66        self.db
67            .db()
68            .register_schema_transform(transform_type, factory)
69    }
70
71    pub fn update_schema(&mut self) -> StructuredSchemaBuilder<'_, Self> {
72        StructuredSchemaBuilder::new(self)
73    }
74
75    /// Creates a structured priority queue backed by a dedicated column family.
76    ///
77    /// The returned queue uses the same `PriorityQueue` API as sharded `StructuredDb`, and the
78    /// queue itself owns the backend-specific dispatch. This keeps upper layers from having to
79    /// special-case `StructuredSingleDb` when offering, scanning, or advancing queue items.
80    pub fn new_priority_queue<'a>(
81        &'a mut self,
82        name: impl Into<String>,
83    ) -> Result<PriorityQueue<'a>> {
84        let normalized_name = priority_queue_column_family_name(name.into())?;
85        if self
86            .db
87            .db()
88            .current_schema()
89            .column_family_ids()
90            .contains_key(normalized_name.as_str())
91        {
92            return Err(Error::InvalidState(format!(
93                "priority queue '{}' already exists",
94                normalized_name
95            )));
96        }
97
98        let mut builder = self.update_schema();
99        builder.add_bytes_column(Some(normalized_name.clone()), 0);
100        builder.set_column_family_options(
101            Some(normalized_name.clone()),
102            priority_queue_column_family_options(),
103        );
104        builder.commit()?;
105        let column_family_id = self
106            .current_schema()
107            .resolve_column_family_id(Some(normalized_name.as_str()))?;
108
109        Ok(PriorityQueue::from_single_column_family(
110            self,
111            normalized_name,
112            column_family_id,
113        ))
114    }
115
116    /// Opens an existing structured priority queue by name.
117    ///
118    /// The target column family must exist, be marked as a priority queue, and contain exactly one
119    /// bytes column.
120    pub fn get_priority_queue<'a>(&'a self, name: impl Into<String>) -> Result<PriorityQueue<'a>> {
121        let normalized_name = priority_queue_column_family_name(name.into())?;
122        let column_family_id = validate_priority_queue_column_family(
123            self.db.db().current_schema().as_ref(),
124            normalized_name.as_str(),
125        )?;
126        Ok(PriorityQueue::from_single_column_family(
127            self,
128            normalized_name,
129            column_family_id,
130        ))
131    }
132
133    /// Opens an existing structured priority queue or creates it on first use.
134    ///
135    /// If a same-named column family already exists, it must already be a valid priority queue
136    /// family.
137    pub fn get_or_new_priority_queue<'a>(
138        &'a mut self,
139        name: impl Into<String>,
140    ) -> Result<PriorityQueue<'a>> {
141        let normalized_name = priority_queue_column_family_name(name.into())?;
142        if self
143            .db
144            .db()
145            .current_schema()
146            .column_family_ids()
147            .contains_key(normalized_name.as_str())
148        {
149            let column_family_id = validate_priority_queue_column_family(
150                self.db.db().current_schema().as_ref(),
151                normalized_name.as_str(),
152            )?;
153            return Ok(PriorityQueue::from_single_column_family(
154                self,
155                normalized_name,
156                column_family_id,
157            ));
158        }
159
160        let mut builder = self.update_schema();
161        builder.add_bytes_column(Some(normalized_name.clone()), 0);
162        builder.set_column_family_options(
163            Some(normalized_name.clone()),
164            priority_queue_column_family_options(),
165        );
166        builder.commit()?;
167        let column_family_id = self
168            .current_schema()
169            .resolve_column_family_id(Some(normalized_name.as_str()))?;
170
171        Ok(PriorityQueue::from_single_column_family(
172            self,
173            normalized_name,
174            column_family_id,
175        ))
176    }
177}
178
179impl StructuredSingleDb {
180    fn reset_default_options(&mut self) {
181        self.default_write_options = StructuredWriteOptions::default();
182        self.default_read_options = StructuredReadOptions::default();
183        self.default_scan_options = StructuredScanOptions::default();
184    }
185
186    pub fn reload_schema(&mut self) -> Result<()> {
187        let schema = load_structured_schema_from_cobble_schema(&self.db.db().current_schema())?;
188        self.structured_schema = Arc::new(schema);
189        self.reset_default_options();
190        Ok(())
191    }
192
193    pub fn apply_schema(
194        &mut self,
195        structured_schema: StructuredSchema,
196    ) -> Result<StructuredSchema> {
197        persist_structured_schema_on_db(self.db.db(), &structured_schema)?;
198        let reloaded = load_structured_schema_from_cobble_schema(&self.db.db().current_schema())?;
199        self.structured_schema = Arc::new(reloaded.clone());
200        self.reset_default_options();
201        Ok(reloaded)
202    }
203
204    // ── Write operations ────────────────────────────────────────────────
205
206    pub fn put<K, V>(&self, bucket: u16, key: K, column: u16, value: V) -> Result<()>
207    where
208        K: AsRef<[u8]>,
209        V: Into<StructuredColumnValue>,
210    {
211        self.put_with_options(bucket, key, column, value, &self.default_write_options)
212    }
213
214    pub fn put_with_options<K, V>(
215        &self,
216        bucket: u16,
217        key: K,
218        column: u16,
219        value: V,
220        options: &StructuredWriteOptions,
221    ) -> Result<()>
222    where
223        K: AsRef<[u8]>,
224        V: Into<StructuredColumnValue>,
225    {
226        let encoded = encode_for_write(
227            &self.structured_schema,
228            options.column_family(),
229            self.db.db().now_seconds(),
230            column,
231            value.into(),
232            options.ttl_seconds(),
233        )?;
234        self.db
235            .put_with_options(bucket, key, column, encoded, options.as_cobble())
236    }
237
238    #[cfg(feature = "ffi")]
239    pub(crate) fn put_borrowed_bytes_with_options<K>(
240        &self,
241        bucket: u16,
242        key: K,
243        column: u16,
244        value: &[u8],
245        options: &StructuredWriteOptions,
246    ) -> Result<()>
247    where
248        K: AsRef<[u8]>,
249    {
250        ensure_bytes_column(&self.structured_schema, options.column_family(), column)?;
251        self.db
252            .put_with_options(bucket, key, column, value, options.as_cobble())
253    }
254
255    pub fn merge<K, V>(&self, bucket: u16, key: K, column: u16, value: V) -> Result<()>
256    where
257        K: AsRef<[u8]>,
258        V: Into<StructuredColumnValue>,
259    {
260        self.merge_with_options(bucket, key, column, value, &self.default_write_options)
261    }
262
263    pub fn merge_with_options<K, V>(
264        &self,
265        bucket: u16,
266        key: K,
267        column: u16,
268        value: V,
269        options: &StructuredWriteOptions,
270    ) -> Result<()>
271    where
272        K: AsRef<[u8]>,
273        V: Into<StructuredColumnValue>,
274    {
275        let encoded = encode_for_write(
276            &self.structured_schema,
277            options.column_family(),
278            self.db.db().now_seconds(),
279            column,
280            value.into(),
281            options.ttl_seconds(),
282        )?;
283        self.db
284            .merge_with_options(bucket, key, column, encoded, options.as_cobble())
285    }
286
287    pub(crate) fn merge_borrowed_bytes_with_options<K>(
288        &self,
289        bucket: u16,
290        key: K,
291        column: u16,
292        value: &[u8],
293        options: &StructuredWriteOptions,
294    ) -> Result<()>
295    where
296        K: AsRef<[u8]>,
297    {
298        ensure_bytes_column(&self.structured_schema, options.column_family(), column)?;
299        self.db
300            .merge_with_options(bucket, key, column, value, options.as_cobble())
301    }
302
303    #[cfg(feature = "ffi")]
304    pub(crate) fn put_borrowed_list_with_options<K>(
305        &self,
306        bucket: u16,
307        key: K,
308        column: u16,
309        elements: &[&[u8]],
310        options: &StructuredWriteOptions,
311    ) -> Result<()>
312    where
313        K: AsRef<[u8]>,
314    {
315        let config = list_column_config(&self.structured_schema, options.column_family(), column)?;
316        let encoded = encode_borrowed_list_for_write(
317            elements,
318            &config,
319            options.ttl_seconds(),
320            self.db.db().now_seconds(),
321        )?;
322        self.db
323            .put_with_options(bucket, key, column, encoded, options.as_cobble())
324    }
325
326    #[cfg(feature = "ffi")]
327    pub(crate) fn merge_borrowed_list_with_options<K>(
328        &self,
329        bucket: u16,
330        key: K,
331        column: u16,
332        elements: &[&[u8]],
333        options: &StructuredWriteOptions,
334    ) -> Result<()>
335    where
336        K: AsRef<[u8]>,
337    {
338        let config = list_column_config(&self.structured_schema, options.column_family(), column)?;
339        let encoded = encode_borrowed_list_for_write(
340            elements,
341            &config,
342            options.ttl_seconds(),
343            self.db.db().now_seconds(),
344        )?;
345        self.db
346            .merge_with_options(bucket, key, column, encoded, options.as_cobble())
347    }
348
349    pub fn delete<K>(&self, bucket: u16, key: K, column: u16) -> Result<()>
350    where
351        K: AsRef<[u8]>,
352    {
353        self.delete_with_options(bucket, key, column, &self.default_write_options)
354    }
355
356    pub fn delete_with_options<K>(
357        &self,
358        bucket: u16,
359        key: K,
360        column: u16,
361        options: &StructuredWriteOptions,
362    ) -> Result<()>
363    where
364        K: AsRef<[u8]>,
365    {
366        self.db
367            .delete_with_options(bucket, key, column, options.as_cobble())
368    }
369
370    pub fn new_write_batch(&self) -> StructuredWriteBatch {
371        StructuredWriteBatch::new(
372            Arc::clone(&self.structured_schema),
373            self.db.db().now_seconds(),
374        )
375    }
376
377    pub fn write_batch(&self, batch: StructuredWriteBatch) -> Result<()> {
378        self.db.write_batch(batch.into_inner())
379    }
380
381    pub fn write_batch_with_options(
382        &self,
383        batch: StructuredWriteBatch,
384        options: &StructuredWriteOptions,
385    ) -> Result<()> {
386        self.db
387            .write_batch_with_options(batch.into_inner(), options.as_cobble())
388    }
389
390    // ── Read operations ─────────────────────────────────────────────────
391
392    pub fn get<K>(&self, bucket: u16, key: K) -> Result<Option<Vec<Option<StructuredColumnValue>>>>
393    where
394        K: AsRef<[u8]>,
395    {
396        self.get_with_options(bucket, key, &self.default_read_options)
397    }
398
399    pub fn multi_get<K>(
400        &self,
401        keys: &[(u16, K)],
402    ) -> Result<Vec<Option<Vec<Option<StructuredColumnValue>>>>>
403    where
404        K: AsRef<[u8]>,
405    {
406        self.multi_get_with_options(keys, &self.default_read_options)
407    }
408
409    pub fn multi_get_with_options<K>(
410        &self,
411        keys: &[(u16, K)],
412        options: &StructuredReadOptions,
413    ) -> Result<Vec<Option<Vec<Option<StructuredColumnValue>>>>>
414    where
415        K: AsRef<[u8]>,
416    {
417        let raw_keys = keys
418            .iter()
419            .map(|(bucket, key)| (*bucket, key.as_ref()))
420            .collect::<Vec<_>>();
421        let projected_schema = options.resolve_projected_schema_cached(&self.structured_schema)?;
422        self.db
423            .multi_get_with_options(&raw_keys, options.as_cobble())?
424            .into_iter()
425            .map(|raw| {
426                raw.map(|columns| decode_row(&projected_schema, 0, columns))
427                    .transpose()
428            })
429            .collect()
430    }
431
432    pub fn get_with_options<K>(
433        &self,
434        bucket: u16,
435        key: K,
436        options: &StructuredReadOptions,
437    ) -> Result<Option<Vec<Option<StructuredColumnValue>>>>
438    where
439        K: AsRef<[u8]>,
440    {
441        let raw = self
442            .db
443            .get_with_options(bucket, key.as_ref(), options.as_cobble())?;
444        let projected_schema = options.resolve_projected_schema_cached(&self.structured_schema)?;
445        raw.map(|columns| decode_row(&projected_schema, 0, columns))
446            .transpose()
447    }
448
449    pub fn scan(&self, bucket: u16, range: Range<&[u8]>) -> Result<StructuredDbIterator> {
450        self.scan_with_options(bucket, range, &self.default_scan_options)
451    }
452
453    pub fn scan_with_options(
454        &self,
455        bucket: u16,
456        range: Range<&[u8]>,
457        options: &StructuredScanOptions,
458    ) -> Result<StructuredDbIterator> {
459        let inner = self
460            .db
461            .scan_with_options(bucket, range, options.as_cobble())?;
462        let projected_schema = options.resolve_projected_schema_cached(&self.structured_schema)?;
463        Ok(StructuredDbIterator::new(inner, projected_schema, 0))
464    }
465
466    pub(crate) fn scan_raw_bounds(
467        &self,
468        bucket: u16,
469        start_key_inclusive: Option<&[u8]>,
470        end_key_exclusive: Option<&[u8]>,
471        options: &StructuredScanOptions,
472    ) -> Result<DbIterator> {
473        self.db.db().scan_with_options_bounds(
474            bucket,
475            start_key_inclusive,
476            end_key_exclusive,
477            options.as_cobble(),
478        )
479    }
480
481    #[cfg(feature = "ffi")]
482    pub(crate) fn scan_with_options_bounds_for_ffi(
483        &self,
484        bucket: u16,
485        start_key_inclusive: Option<&[u8]>,
486        end_key_exclusive: Option<&[u8]>,
487        options: &StructuredScanOptions,
488    ) -> Result<StructuredDbIterator> {
489        let inner =
490            self.scan_raw_bounds(bucket, start_key_inclusive, end_key_exclusive, options)?;
491        let projected_schema = options.resolve_projected_schema_cached(&self.structured_schema)?;
492        Ok(StructuredDbIterator::new(inner, projected_schema, 0))
493    }
494
495    // ── Snapshot lifecycle ───────────────────────────────────────────────
496
497    pub fn snapshot(&self) -> Result<u64> {
498        self.db.snapshot()
499    }
500
501    /// Change the memtable implementation for future rotations.
502    ///
503    /// `flush_current = false` leaves the active memtable untouched; `true` rotates a non-empty
504    /// active memtable now through the normal flush path.
505    pub fn switch_memtable_type(
506        &self,
507        memtable_type: MemtableType,
508        flush_current: bool,
509    ) -> Result<()> {
510        self.db.switch_memtable_type(memtable_type, flush_current)
511    }
512
513    /// Mark all currently referenced READONLY files for asynchronous loading into primary
514    /// storage.
515    pub fn load_readonly_files_to_primary(&self) -> Result<usize> {
516        self.db.load_readonly_files_to_primary()
517    }
518
519    pub fn snapshot_with_callback<F>(&self, callback: F) -> Result<u64>
520    where
521        F: Fn(Result<cobble::GlobalSnapshotManifest>) + Send + Sync + 'static,
522    {
523        self.db.snapshot_with_callback(callback)
524    }
525
526    pub fn retain_snapshot(&self, global_snapshot_id: u64) -> Result<bool> {
527        self.db.retain_snapshot(global_snapshot_id)
528    }
529
530    pub fn expire_snapshot(&self, global_snapshot_id: u64) -> Result<bool> {
531        self.db.expire_snapshot(global_snapshot_id)
532    }
533
534    pub fn list_snapshots(&self) -> Result<Vec<cobble::GlobalSnapshotManifest>> {
535        self.db.list_snapshots()
536    }
537
538    pub fn set_time(&self, next: u32) {
539        self.db.set_time(next)
540    }
541
542    pub fn close(&self) -> Result<()> {
543        self.db.close()
544    }
545
546    pub(crate) fn advance_column_family_truncation_cursor_by_id(
547        &self,
548        bucket: u16,
549        column_family_id: u8,
550        key: &[u8],
551    ) -> Result<()> {
552        self.db
553            .db()
554            .advance_truncation_cursor_by_id(bucket, column_family_id, key)
555    }
556
557    pub(crate) fn column_family_truncation_cursor_by_id(
558        &self,
559        bucket: u16,
560        column_family_id: u8,
561    ) -> Result<Option<Vec<u8>>> {
562        self.db
563            .db()
564            .truncation_cursor_by_id(bucket, column_family_id)
565    }
566}
567
568impl StructuredSchemaOwner for StructuredSingleDb {
569    fn current_structured_schema(&self) -> StructuredSchema {
570        self.current_schema()
571    }
572
573    fn begin_core_schema_update(&self) -> cobble::SchemaBuilder {
574        self.db.db().update_schema()
575    }
576
577    fn install_committed_structured_schema(
578        &mut self,
579        schema: StructuredSchema,
580    ) -> StructuredSchema {
581        self.structured_schema = Arc::new(schema.clone());
582        self.reset_default_options();
583        schema
584    }
585}
586
587#[cfg(test)]
588#[path = "../tests/unit/structured_single_db.rs"]
589mod tests;