Skip to main content

cobble_table/
runtime.rs

1use crate::catalog::{TableId, physical_table_name};
2use crate::metadata::TableMetadata;
3use crate::table::{
4    TypedRead, load_table_metadata, load_table_metadata_for_shard,
5    load_table_metadata_from_snapshot, validate_name,
6};
7use crate::transform::TABLE_TRANSFORM_TYPE;
8use crate::{
9    ReadOnlyTable, Result, Table, TableError, TableKey, TableKeyBuilder, TableProjection,
10    TableScan, TableSchema, TableWritePlan, Value, register_schema_transforms,
11};
12use arc_swap::{ArcSwap, Guard};
13use bytes::Bytes;
14use cobble::{
15    Config, Db, DbBuilder, DbGovernance, DbIterator, FileSystemDbGovernance, GovernanceMode,
16    NoopDbGovernance, ReadOnlyDbBuilder, ReadOptions, Reader, ReaderBuilder, ReaderConfig,
17    ScanOptions, SchemaTransformRegistrar, VolumeUsageKind, bucket_snapshot_manifest_path,
18};
19use std::fs::{File, OpenOptions};
20use std::io;
21use std::ops::RangeInclusive;
22use std::path::PathBuf;
23use std::sync::atomic::{AtomicU64, Ordering};
24use std::sync::{Arc, Mutex};
25use std::time::{Duration, Instant};
26use url::Url;
27
28type SchemaTransformCallback =
29    Box<dyn Fn(Option<Bytes>) -> cobble::Result<Option<Bytes>> + Send + Sync>;
30type SchemaTransformFactory =
31    Arc<dyn Fn(&[u8]) -> cobble::Result<SchemaTransformCallback> + Send + Sync>;
32
33#[derive(Default)]
34pub(crate) struct TableSchemaTransformFactories(Vec<(String, SchemaTransformFactory)>);
35
36impl TableSchemaTransformFactories {
37    pub(crate) fn register<F, T>(
38        &mut self,
39        transform_type: impl Into<String>,
40        factory: F,
41    ) -> Result<()>
42    where
43        F: Fn(&[u8]) -> cobble::Result<T> + Send + Sync + 'static,
44        T: Fn(Option<Bytes>) -> cobble::Result<Option<Bytes>> + Send + Sync + 'static,
45    {
46        let transform_type = transform_type.into();
47        if transform_type == TABLE_TRANSFORM_TYPE {
48            return Err(TableError::Storage(cobble::Error::InvalidState(format!(
49                "Schema transform '{TABLE_TRANSFORM_TYPE}' is reserved by cobble-table"
50            ))));
51        }
52        if transform_type.trim().is_empty() {
53            return Err(TableError::Storage(cobble::Error::InvalidState(
54                "Schema transform type must not be empty".to_string(),
55            )));
56        }
57        if self
58            .0
59            .iter()
60            .any(|(existing, _)| existing == &transform_type)
61        {
62            return Err(TableError::Storage(cobble::Error::InvalidState(format!(
63                "Schema transform '{}' is already registered",
64                transform_type
65            ))));
66        }
67        let factory: SchemaTransformFactory = Arc::new(move |spec| Ok(Box::new(factory(spec)?)));
68        self.0.push((transform_type, factory));
69        Ok(())
70    }
71
72    pub(crate) fn apply_to<B: SchemaTransformRegistrar>(&self, builder: B) -> Result<B> {
73        register_schema_transforms(&builder)?;
74        for (transform_type, factory) in &self.0 {
75            let factory = Arc::clone(factory);
76            builder.register_schema_transform(transform_type.clone(), move |spec| factory(spec))?;
77        }
78        Ok(builder)
79    }
80}
81
82/// Builder for a standalone writable typed table shard.
83pub struct TableWriterBuilder {
84    config: Config,
85    table_name: Option<String>,
86    bucket: Option<u16>,
87    catalog_binding: Option<(TableWritePlan, Config)>,
88    transforms: TableSchemaTransformFactories,
89}
90
91impl TableWriterBuilder {
92    pub fn new(config: Config) -> Self {
93        Self {
94            config,
95            table_name: None,
96            bucket: None,
97            catalog_binding: None,
98            transforms: TableSchemaTransformFactories::default(),
99        }
100    }
101
102    pub(crate) fn from_write_plan(
103        config: Config,
104        table_name: String,
105        plan: TableWritePlan,
106        catalog_store_config: Config,
107    ) -> Self {
108        Self {
109            config,
110            table_name: Some(table_name),
111            bucket: None,
112            catalog_binding: Some((plan, catalog_store_config)),
113            transforms: TableSchemaTransformFactories::default(),
114        }
115    }
116
117    /// Select the physical column-family name for this typed table.
118    pub fn table_name(mut self, table_name: impl Into<String>) -> Self {
119        self.table_name = Some(table_name.into());
120        self
121    }
122
123    /// Own exactly one physical bucket using the stable database identity `bucket-N`.
124    ///
125    /// The writer keeps a local-filesystem advisory lock for that database until its underlying
126    /// [`Db`] is dropped.
127    pub fn bucket(mut self, bucket: u16) -> Self {
128        self.bucket = Some(bucket);
129        self
130    }
131
132    /// Register a factory for persisted schema transform specifications before opening.
133    pub fn register_schema_transform<F, T>(
134        mut self,
135        transform_type: impl Into<String>,
136        factory: F,
137    ) -> Result<Self>
138    where
139        F: Fn(&[u8]) -> cobble::Result<T> + Send + Sync + 'static,
140        T: Fn(Option<Bytes>) -> cobble::Result<Option<Bytes>> + Send + Sync + 'static,
141    {
142        self.transforms.register(transform_type, factory)?;
143        Ok(self)
144    }
145
146    /// Initialize the bucket's empty baseline, then create or reopen this standalone table.
147    ///
148    /// A subsequent call starts again from that empty baseline, for overwrite semantics.
149    pub fn create(self, schema: TableSchema) -> Result<Table> {
150        if self.catalog_binding.is_some() {
151            return Err(TableError::InvalidSchema(
152                "catalog-bound TableWriterBuilder requires open()".into(),
153            ));
154        }
155        let name = self.required_table_name()?;
156        let bucket = self.required_bucket()?;
157        let (db, needs_baseline) =
158            self.open_or_resume_bucket_baseline(&name, Some(&schema), None)?;
159        let table = Table::create(db, name, schema)?;
160        if needs_baseline {
161            ensure_empty_baseline(&table)?;
162        } else {
163            ensure_empty_bucket(&table, bucket)?;
164        }
165        Ok(table)
166    }
167
168    /// Initialize the bucket's empty baseline, then open a writable catalog-bound table.
169    ///
170    /// A subsequent call starts again from that empty baseline, for overwrite semantics.
171    pub fn open(self) -> Result<Table> {
172        let (plan, store_config) = self.catalog_binding.as_ref().ok_or_else(|| {
173            TableError::InvalidSchema("TableWriterBuilder::open requires a catalog table".into())
174        })?;
175        let name = self.required_table_name()?;
176        let bucket = self.required_bucket()?;
177        let (db, needs_baseline) =
178            self.open_or_resume_bucket_baseline(&name, None, Some(plan.table_id()))?;
179        let table = open_materialized_catalog_table(db, plan, store_config)?;
180        if needs_baseline {
181            ensure_empty_baseline(&table)?;
182        } else {
183            ensure_empty_bucket(&table, bucket)?;
184        }
185        Ok(table)
186    }
187
188    /// Resume a writable shard from a selected snapshot boundary.
189    pub fn resume_from_snapshot(self, snapshot_id: u64) -> Result<Table> {
190        let name = self.required_table_name()?;
191        let db = self.resume_bucket_snapshot(
192            snapshot_id,
193            &name,
194            self.catalog_binding
195                .as_ref()
196                .map(|(plan, _)| plan.table_id()),
197        )?;
198        match self.catalog_binding {
199            Some((plan, store_config)) => open_materialized_catalog_table(db, &plan, &store_config),
200            None => Table::open(db, name),
201        }
202    }
203
204    fn required_table_name(&self) -> Result<String> {
205        let name = self
206            .table_name
207            .clone()
208            .ok_or_else(|| {
209                TableError::InvalidSchema("TableWriterBuilder requires table_name".into())
210            })
211            .and_then(validate_name)?;
212        if let Some((plan, _)) = &self.catalog_binding
213            && name != physical_table_name(plan.table_id())
214        {
215            return Err(TableError::InvalidSchema(
216                "catalog-bound TableWriterBuilder cannot change table_name".into(),
217            ));
218        }
219        Ok(name)
220    }
221
222    fn required_bucket(&self) -> Result<u16> {
223        self.bucket
224            .ok_or_else(|| TableError::InvalidSchema("TableWriterBuilder requires bucket".into()))
225    }
226
227    fn bucket_config(&self) -> Result<Config> {
228        let bucket = self.required_bucket()?;
229        if !(1..=u16::MAX as u32 + 1).contains(&self.config.total_buckets)
230            || u32::from(bucket) >= self.config.total_buckets
231        {
232            return Err(TableError::InvalidSchema(
233                "single bucket must be in [0, total_buckets), with total_buckets in 1..=65536"
234                    .into(),
235            ));
236        }
237        let mut config = self.config.clone();
238        config.wal_enabled = false;
239        if config.snapshot_retention.is_some() && !config.snapshot_only_track {
240            return Err(TableError::InvalidSchema(
241                "single-bucket writers require snapshot retention to be disabled".into(),
242            ));
243        }
244        Ok(config)
245    }
246
247    fn open_or_resume_bucket_baseline(
248        &self,
249        table_name: &str,
250        standalone_schema: Option<&TableSchema>,
251        catalog_table_id: Option<crate::catalog::TableId>,
252    ) -> Result<(Arc<Db>, bool)> {
253        let bucket = self.required_bucket()?;
254        let governance = self.bucket_governance(bucket)?;
255        if bucket_snapshot_exists(&self.config, bucket, 0)? {
256            self.validate_bucket_snapshot(0, table_name, standalone_schema, catalog_table_id)?;
257            let db = Arc::new(
258                self.bucket_db_builder(governance)?
259                    .resume_from_snapshot(0)?,
260            );
261            return Ok((db, false));
262        }
263        if bucket_has_persisted_state(&self.config, bucket)? {
264            return Err(TableError::Storage(cobble::Error::InvalidState(format!(
265                "{} has persisted state but no empty baseline snapshot 0",
266                bucket_db_id(bucket)
267            ))));
268        }
269        let db = Arc::new(self.bucket_db_builder(governance)?.open()?);
270        Ok((db, true))
271    }
272
273    fn resume_bucket_snapshot(
274        &self,
275        snapshot_id: u64,
276        table_name: &str,
277        catalog_table_id: Option<crate::catalog::TableId>,
278    ) -> Result<Arc<Db>> {
279        let bucket = self.required_bucket()?;
280        let governance = self.bucket_governance(bucket)?;
281        self.validate_bucket_snapshot(snapshot_id, table_name, None, catalog_table_id)?;
282        Ok(Arc::new(
283            self.bucket_db_builder(governance)?
284                .resume_from_snapshot(snapshot_id)?,
285        ))
286    }
287
288    fn bucket_governance(&self, bucket: u16) -> Result<Arc<dyn DbGovernance>> {
289        self.bucket_config()?;
290        Ok(Arc::new(LockedDbGovernance::new(
291            &self.config,
292            bucket,
293            default_governance(&self.config)?,
294        )?))
295    }
296
297    fn bucket_db_builder(&self, governance: Arc<dyn DbGovernance>) -> Result<DbBuilder> {
298        let bucket = self.required_bucket()?;
299        let config = self.bucket_config()?;
300        self.transforms.apply_to(
301            DbBuilder::new(config)
302                .bucket_ranges(vec![bucket..=bucket])
303                .db_id(bucket_db_id(bucket))
304                .governance(governance),
305        )
306    }
307
308    fn validate_bucket_snapshot(
309        &self,
310        snapshot_id: u64,
311        table_name: &str,
312        standalone_schema: Option<&TableSchema>,
313        catalog_table_id: Option<crate::catalog::TableId>,
314    ) -> Result<()> {
315        let bucket = self.required_bucket()?;
316        let config = self.bucket_config()?;
317        let db_id = bucket_db_id(bucket);
318        let manifest_path = bucket_snapshot_manifest_path(&db_id, snapshot_id);
319        let absolute_manifest = configured_local_meta_root(&config)?.join(manifest_path);
320        if !absolute_manifest.exists() {
321            return Err(TableError::InvalidSchema(format!(
322                "single-bucket writer is missing snapshot {snapshot_id} for {}",
323                bucket_db_id(bucket)
324            )));
325        }
326        let snapshot = cobble::load_shard_snapshot_metadata(
327            &config,
328            &db_id,
329            &absolute_manifest.to_string_lossy(),
330        )?;
331        if snapshot.snapshot_id != snapshot_id || snapshot.ranges.as_slice() != [bucket..=bucket] {
332            return Err(TableError::InvalidSchema(format!(
333                "single-bucket manifest must cover exactly bucket {bucket} at snapshot {snapshot_id}"
334            )));
335        }
336        if snapshot_id == 0 && snapshot.data_size_bytes != 0 {
337            return Err(TableError::InvalidSchema(
338                "single-bucket baseline snapshot 0 must not contain data".into(),
339            ));
340        }
341        let metadata = load_table_metadata_from_snapshot(&snapshot, table_name)?;
342        if let Some(schema) = standalone_schema
343            && (metadata.catalog_binding.is_some() || metadata.schema != *schema)
344        {
345            return Err(TableError::InvalidSchema(format!(
346                "column family '{table_name}' is not this standalone table"
347            )));
348        }
349        validate_catalog_binding(&metadata, catalog_table_id)
350    }
351}
352
353fn bucket_db_id(bucket: u16) -> String {
354    format!("bucket-{bucket}")
355}
356
357fn ensure_empty_baseline(table: &Table) -> Result<()> {
358    let snapshot = table.snapshot_and_wait()?;
359    if snapshot.snapshot_id != 0 {
360        return Err(TableError::Storage(cobble::Error::InvalidState(format!(
361            "new single-bucket writer must create empty baseline snapshot 0, got {}",
362            snapshot.snapshot_id
363        ))));
364    }
365    Ok(())
366}
367
368fn ensure_empty_bucket(table: &Table, bucket: u16) -> Result<()> {
369    if table.scan(bucket)?.next().transpose()?.is_some() {
370        return Err(TableError::InvalidSchema(
371            "single-bucket baseline snapshot 0 must not contain table rows".into(),
372        ));
373    }
374    Ok(())
375}
376
377fn default_governance(config: &Config) -> cobble::Result<Arc<dyn DbGovernance>> {
378    match config.governance_mode {
379        GovernanceMode::Filesystem => Ok(Arc::new(FileSystemDbGovernance::from_config(config)?)),
380        GovernanceMode::Noop => Ok(Arc::new(NoopDbGovernance)),
381    }
382}
383
384struct LockedDbGovernance {
385    inner: Arc<dyn DbGovernance>,
386    _lock: LocalBucketLock,
387}
388
389impl LockedDbGovernance {
390    fn new(config: &Config, bucket: u16, inner: Arc<dyn DbGovernance>) -> cobble::Result<Self> {
391        Ok(Self {
392            inner,
393            _lock: LocalBucketLock::acquire(config, bucket)?,
394        })
395    }
396}
397
398impl DbGovernance for LockedDbGovernance {
399    fn register_db(
400        &self,
401        db_id: &str,
402        ranges: &[RangeInclusive<u16>],
403        total_buckets: u32,
404    ) -> cobble::Result<()> {
405        self.inner.register_db(db_id, ranges, total_buckets)
406    }
407
408    fn unregister_db(&self, db_id: &str) -> cobble::Result<()> {
409        self.inner.unregister_db(db_id)
410    }
411}
412
413struct LocalBucketLock {
414    _file: File,
415}
416
417impl LocalBucketLock {
418    fn acquire(config: &Config, bucket: u16) -> cobble::Result<Self> {
419        let root = local_meta_root(config)?;
420        let bucket_dir = root.join(bucket_db_id(bucket));
421        std::fs::create_dir_all(&bucket_dir).map_err(lock_io_error)?;
422        let bucket_dir = std::fs::canonicalize(&bucket_dir).map_err(lock_io_error)?;
423        let file = OpenOptions::new()
424            .create(true)
425            .read(true)
426            .write(true)
427            .truncate(false)
428            .open(bucket_dir.join(".cobble-table-writer.lock"))
429            .map_err(lock_io_error)?;
430        file.try_lock().map_err(|error| {
431            cobble::Error::InvalidState(format!("single-bucket writer is already active: {error}"))
432        })?;
433        Ok(Self { _file: file })
434    }
435}
436
437fn local_meta_root(config: &Config) -> cobble::Result<PathBuf> {
438    let root = configured_local_meta_root(config)?;
439    std::fs::create_dir_all(&root).map_err(lock_io_error)?;
440    std::fs::canonicalize(root).map_err(lock_io_error)
441}
442
443fn configured_local_meta_root(config: &Config) -> cobble::Result<PathBuf> {
444    let volume = config
445        .volumes
446        .iter()
447        .find(|volume| volume.supports(VolumeUsageKind::Meta))
448        .ok_or_else(|| {
449            cobble::Error::ConfigError("single-bucket writer requires a META volume".into())
450        })?;
451    let root = match Url::parse(&volume.base_dir) {
452        Ok(url) if url.scheme() == "file" => url.to_file_path().map_err(|_| {
453            cobble::Error::ConfigError("single-bucket META volume is not a local path".into())
454        })?,
455        Ok(_) => {
456            return Err(cobble::Error::ConfigError(
457                "single-bucket table writers require a local file:// META volume".into(),
458            ));
459        }
460        Err(_) if volume.base_dir.contains("://") => {
461            return Err(cobble::Error::ConfigError(format!(
462                "invalid single-bucket META volume URL: {}",
463                volume.base_dir
464            )));
465        }
466        Err(_) => PathBuf::from(&volume.base_dir),
467    };
468    Ok(root)
469}
470
471fn bucket_has_persisted_state(config: &Config, bucket: u16) -> Result<bool> {
472    let bucket_dir = local_meta_root(config)?.join(bucket_db_id(bucket));
473    if !bucket_dir.exists() {
474        return Ok(false);
475    }
476    for entry in std::fs::read_dir(bucket_dir).map_err(lock_io_error)? {
477        let entry = entry.map_err(lock_io_error)?;
478        if entry.file_name() != ".cobble-table-writer.lock" {
479            return Ok(true);
480        }
481    }
482    Ok(false)
483}
484
485fn bucket_snapshot_exists(config: &Config, bucket: u16, snapshot_id: u64) -> Result<bool> {
486    let db_id = bucket_db_id(bucket);
487    Ok(configured_local_meta_root(config)?
488        .join(bucket_snapshot_manifest_path(&db_id, snapshot_id))
489        .exists())
490}
491
492fn lock_io_error(error: io::Error) -> cobble::Error {
493    cobble::Error::IoError(format!("single-bucket writer lock: {error}"))
494}
495
496fn open_materialized_catalog_table(
497    db: Arc<Db>,
498    plan: &TableWritePlan,
499    store_config: &Config,
500) -> Result<Table> {
501    let (name, metadata) = crate::catalog::materialize_write_plan(store_config, db.as_ref(), plan)?;
502    Table::from_metadata(db, name, metadata)
503}
504
505/// A typed global snapshot read proxy.
506pub struct TableReader {
507    typed: ArcSwap<TypedRead>,
508    refresh: AutoRefreshController,
509}
510
511struct AutoRefreshController {
512    interval: Option<Duration>,
513    started_at: Instant,
514    next_check_at: AtomicU64,
515    refreshing: Mutex<()>,
516}
517
518impl AutoRefreshController {
519    fn new(interval: Option<Duration>) -> Self {
520        let started_at = Instant::now();
521        let next_check_at = interval.map_or(0, duration_nanos);
522        Self {
523            interval,
524            started_at,
525            next_check_at: AtomicU64::new(next_check_at),
526            refreshing: Mutex::new(()),
527        }
528    }
529
530    fn due(&self) -> bool {
531        self.interval.is_some()
532            && self.elapsed_nanos() >= self.next_check_at.load(Ordering::Acquire)
533    }
534
535    fn schedule_next_check(&self) {
536        let Some(interval) = self.interval else {
537            return;
538        };
539        self.next_check_at.store(
540            self.elapsed_nanos()
541                .saturating_add(duration_nanos(interval)),
542            Ordering::Release,
543        );
544    }
545
546    fn elapsed_nanos(&self) -> u64 {
547        self.started_at
548            .elapsed()
549            .as_nanos()
550            .min(u128::from(u64::MAX)) as u64
551    }
552
553    fn lock(&self) -> Result<std::sync::MutexGuard<'_, ()>> {
554        self.refreshing
555            .lock()
556            .map_err(|_| TableError::internal("table reader refresh lock poisoned"))
557    }
558
559    fn try_lock(&self) -> Result<Option<std::sync::MutexGuard<'_, ()>>> {
560        match self.refreshing.try_lock() {
561            Ok(guard) => Ok(Some(guard)),
562            Err(std::sync::TryLockError::WouldBlock) => Ok(None),
563            Err(std::sync::TryLockError::Poisoned(_)) => {
564                Err(TableError::internal("table reader refresh lock poisoned"))
565            }
566        }
567    }
568}
569
570fn duration_nanos(duration: Duration) -> u64 {
571    duration.as_nanos().min(u128::from(u64::MAX)) as u64
572}
573
574#[cfg(test)]
575#[path = "../tests/unit/runtime.rs"]
576mod tests;
577
578impl TableReader {
579    /// Open the snapshot carried by this core reader.
580    ///
581    /// Readers opened from a core current-pointer reader check for later committed snapshots on
582    /// its configured interval. Core fixed-snapshot readers remain fixed.
583    pub fn open(mut reader: Reader, name: impl Into<String>) -> Result<Self> {
584        let refresh_interval = reader.auto_refresh_interval();
585        reader.pin_current_snapshot();
586        let name = validate_name(name.into())?;
587        let metadata = global_metadata(&reader, &name)?;
588        Self::from_metadata(reader, name, metadata, refresh_interval)
589    }
590
591    fn from_metadata(
592        reader: Reader,
593        name: String,
594        metadata: TableMetadata,
595        refresh_interval: Option<Duration>,
596    ) -> Result<Self> {
597        let state = Arc::new(GlobalReaderState::new(
598            reader,
599            name.clone(),
600            metadata.clone(),
601        ));
602        Ok(Self {
603            typed: ArcSwap::from_pointee(TypedRead::from_global_metadata(state, name, metadata)?),
604            refresh: AutoRefreshController::new(refresh_interval),
605        })
606    }
607
608    /// Return the schema of this reader's currently loaded global snapshot.
609    ///
610    /// This accessor never checks `CURRENT`. A later fallible read may refresh a current-pointer
611    /// reader; call [`Self::refresh`] before schema inspection when the latest committed snapshot
612    /// is required.
613    pub fn schema(&self) -> Arc<TableSchema> {
614        self.typed.load().schema_arc()
615    }
616
617    /// Start building one primary key in schema order.
618    pub fn key_builder(&self) -> TableKeyBuilder {
619        self.typed.load().key_builder()
620    }
621
622    /// Compile a reusable read projection from top-level field names.
623    pub fn project_by_names<S: AsRef<str>>(&self, field_names: &[S]) -> Result<TableProjection> {
624        self.view_for_access()?.project_by_names(field_names)
625    }
626
627    /// Read one row by primary key.
628    pub fn get(&self, key: &TableKey) -> Result<Option<Vec<Value>>> {
629        self.view_for_access()?.get(key)
630    }
631
632    /// Read many primary keys while preserving order and duplicates.
633    pub fn multi_get(&self, keys: &[TableKey]) -> Result<Vec<Option<Vec<Value>>>> {
634        self.view_for_access()?.multi_get(keys)
635    }
636
637    /// Scan all rows in one bucket.
638    pub fn scan(&self, bucket: u16) -> Result<TableScan> {
639        self.view_for_access()?.scan(bucket)
640    }
641
642    /// Scan one bucket from an inclusive primary-key bound to an exclusive bound.
643    pub fn scan_bounds(
644        &self,
645        bucket: u16,
646        start_key_inclusive: Option<&TableKey>,
647        end_key_exclusive: Option<&TableKey>,
648    ) -> Result<TableScan> {
649        self.view_for_access()?
650            .scan_bounds(bucket, start_key_inclusive, end_key_exclusive)
651    }
652
653    /// Build a portable full-scan plan pinned to this reader's current snapshot.
654    pub fn scan_plan(&self) -> Result<crate::TableScanPlan> {
655        let typed = self.view_for_access()?;
656        let state = typed
657            .global_state()
658            .ok_or_else(|| TableError::internal("table reader is missing its global read state"))?;
659        state.scan_plan()
660    }
661
662    /// Refresh this current-pointer reader to the latest committed global snapshot.
663    ///
664    /// Existing projections, scans, and scan plans retain their current fixed view. If the
665    /// current pointer is unchanged, this returns `false`; an error leaves this reader unchanged.
666    /// Fixed snapshot readers always return `false`.
667    pub fn refresh(&self) -> Result<bool> {
668        if self.refresh.interval.is_none() {
669            return Ok(false);
670        }
671        let _guard = self.refresh.lock()?;
672        let result = self.refresh_loaded_view();
673        self.refresh.schedule_next_check();
674        result
675    }
676
677    #[cfg(feature = "ffi")]
678    #[doc(hidden)]
679    pub fn auto_refresh_interval_nanos(&self) -> Option<u64> {
680        self.refresh
681            .interval
682            .map(|interval| interval.as_nanos().min(u128::from(u64::MAX)) as u64)
683    }
684
685    fn view_for_access(&self) -> Result<Guard<Arc<TypedRead>>> {
686        let view = self.typed.load();
687        if !self.refresh.due() {
688            return Ok(view);
689        }
690        let Some(_guard) = self.refresh.try_lock()? else {
691            return Ok(view);
692        };
693        drop(view);
694        if self.refresh.due() {
695            let refreshed = self.refresh_loaded_view();
696            self.refresh.schedule_next_check();
697            refreshed?;
698        }
699        Ok(self.typed.load())
700    }
701
702    fn refresh_loaded_view(&self) -> Result<bool> {
703        let typed = self.typed.load_full();
704        let state = typed
705            .global_state()
706            .ok_or_else(|| TableError::internal("table reader is missing its global read state"))?;
707        let Some((reader, name, previous_table_id)) = state.refreshed_snapshot()? else {
708            return Ok(false);
709        };
710        let metadata = global_metadata(&reader, &name)?;
711        if previous_table_id != metadata.catalog_binding.map(|binding| binding.table_id) {
712            return Err(TableError::InvalidSchema(
713                "refreshed global snapshot belongs to a different catalog table".into(),
714            ));
715        }
716        let state = Arc::new(GlobalReaderState::new(
717            reader,
718            name.clone(),
719            metadata.clone(),
720        ));
721        let typed = TypedRead::from_global_metadata(state, name, metadata)?;
722        self.typed.store(Arc::new(typed));
723        Ok(true)
724    }
725
726    #[cfg(feature = "ffi")]
727    pub(crate) fn ffi_acquire_view(&self) -> crate::ffi::TableReaderView {
728        crate::ffi::TableReaderView {
729            typed: self.typed.load_full(),
730        }
731    }
732
733    #[cfg(feature = "ffi")]
734    pub(crate) fn ffi_acquire_view_if_changed(
735        &self,
736        current: &crate::ffi::TableReaderView,
737    ) -> Option<crate::ffi::TableReaderView> {
738        let typed = self.typed.load_full();
739        if Arc::ptr_eq(&typed, &current.typed) {
740            None
741        } else {
742            Some(crate::ffi::TableReaderView { typed })
743        }
744    }
745}
746
747/// Builder for a global typed reader from a fixed or current snapshot selection.
748pub struct TableReaderBuilder {
749    config: Config,
750    table_name: Option<String>,
751    selection: Option<TableReaderSelection>,
752    catalog_table_id: Option<TableId>,
753    transforms: TableSchemaTransformFactories,
754}
755
756enum TableReaderSelection {
757    Global { snapshot_id: u64 },
758    CurrentGlobal,
759}
760
761impl TableReaderBuilder {
762    pub fn new(config: Config) -> Self {
763        Self {
764            config,
765            table_name: None,
766            selection: None,
767            catalog_table_id: None,
768            transforms: TableSchemaTransformFactories::default(),
769        }
770    }
771
772    pub(crate) fn from_catalog(config: Config, table_name: String, table_id: TableId) -> Self {
773        Self {
774            config,
775            table_name: Some(table_name),
776            selection: None,
777            catalog_table_id: Some(table_id),
778            transforms: TableSchemaTransformFactories::default(),
779        }
780    }
781
782    pub fn table_name(mut self, table_name: impl Into<String>) -> Self {
783        self.table_name = Some(table_name.into());
784        self
785    }
786
787    pub fn global_snapshot(mut self, snapshot_id: u64) -> Self {
788        self.selection = Some(TableReaderSelection::Global { snapshot_id });
789        self
790    }
791
792    /// Open the current committed snapshot and check for later commits on access intervals.
793    pub fn current_global_snapshot(mut self) -> Self {
794        self.selection = Some(TableReaderSelection::CurrentGlobal);
795        self
796    }
797
798    /// Register a factory for persisted schema transform specifications before opening.
799    pub fn register_schema_transform<F, T>(
800        mut self,
801        transform_type: impl Into<String>,
802        factory: F,
803    ) -> Result<Self>
804    where
805        F: Fn(&[u8]) -> cobble::Result<T> + Send + Sync + 'static,
806        T: Fn(Option<Bytes>) -> cobble::Result<Option<Bytes>> + Send + Sync + 'static,
807    {
808        self.transforms.register(transform_type, factory)?;
809        Ok(self)
810    }
811
812    pub fn open(self) -> Result<TableReader> {
813        let name = self.required_table_name()?;
814        let builder = self
815            .transforms
816            .apply_to(ReaderBuilder::new(ReaderConfig::from_config(&self.config)))?;
817        let mut reader = match self.selection.ok_or_else(|| {
818            TableError::InvalidSchema("TableReaderBuilder requires a snapshot selection".into())
819        })? {
820            TableReaderSelection::Global { snapshot_id } => builder.open(snapshot_id)?,
821            TableReaderSelection::CurrentGlobal => builder.open_current()?,
822        };
823        let refresh_interval = reader.auto_refresh_interval();
824        reader.pin_current_snapshot();
825        let metadata = global_metadata(&reader, &name)?;
826        validate_catalog_binding(&metadata, self.catalog_table_id)?;
827        TableReader::from_metadata(reader, name, metadata, refresh_interval)
828    }
829
830    fn required_table_name(&self) -> Result<String> {
831        required_reader_table_name(&self.table_name, self.catalog_table_id)
832    }
833}
834
835/// Builder for a typed table over one shard snapshot.
836pub struct ReadOnlyTableBuilder {
837    config: Config,
838    table_name: Option<String>,
839    snapshot: Option<(String, u64)>,
840    catalog_table_id: Option<TableId>,
841    transforms: TableSchemaTransformFactories,
842}
843
844impl ReadOnlyTableBuilder {
845    pub fn new(config: Config) -> Self {
846        Self {
847            config,
848            table_name: None,
849            snapshot: None,
850            catalog_table_id: None,
851            transforms: TableSchemaTransformFactories::default(),
852        }
853    }
854
855    pub(crate) fn from_catalog(config: Config, table_name: String, table_id: TableId) -> Self {
856        Self {
857            config,
858            table_name: Some(table_name),
859            snapshot: None,
860            catalog_table_id: Some(table_id),
861            transforms: TableSchemaTransformFactories::default(),
862        }
863    }
864
865    pub fn table_name(mut self, table_name: impl Into<String>) -> Self {
866        self.table_name = Some(table_name.into());
867        self
868    }
869
870    pub fn shard_snapshot(mut self, db_id: impl Into<String>, snapshot_id: u64) -> Self {
871        self.snapshot = Some((db_id.into(), snapshot_id));
872        self
873    }
874
875    /// Register a factory for persisted schema transform specifications before opening.
876    pub fn register_schema_transform<F, T>(
877        mut self,
878        transform_type: impl Into<String>,
879        factory: F,
880    ) -> Result<Self>
881    where
882        F: Fn(&[u8]) -> cobble::Result<T> + Send + Sync + 'static,
883        T: Fn(Option<Bytes>) -> cobble::Result<Option<Bytes>> + Send + Sync + 'static,
884    {
885        self.transforms.register(transform_type, factory)?;
886        Ok(self)
887    }
888
889    pub fn open(self) -> Result<ReadOnlyTable> {
890        let name = self.required_table_name()?;
891        let (db_id, snapshot_id) = self.snapshot.ok_or_else(|| {
892            TableError::InvalidSchema(
893                "ReadOnlyTableBuilder requires a shard snapshot selection".into(),
894            )
895        })?;
896        let db = Arc::new(
897            self.transforms
898                .apply_to(ReadOnlyDbBuilder::new(self.config).db_id(db_id))?
899                .open(snapshot_id)?,
900        );
901        let metadata = load_table_metadata(&db.current_schema(), &name)?;
902        validate_catalog_binding(&metadata, self.catalog_table_id)?;
903        ReadOnlyTable::from_shard_metadata(db, name, metadata)
904    }
905
906    fn required_table_name(&self) -> Result<String> {
907        required_reader_table_name(&self.table_name, self.catalog_table_id)
908    }
909}
910
911fn required_reader_table_name(
912    table_name: &Option<String>,
913    catalog_table_id: Option<TableId>,
914) -> Result<String> {
915    let name = table_name
916        .clone()
917        .ok_or_else(|| TableError::InvalidSchema("table reader requires table_name".into()))
918        .and_then(validate_name)?;
919    if let Some(table_id) = catalog_table_id
920        && name != physical_table_name(table_id)
921    {
922        return Err(TableError::InvalidSchema(
923            "catalog-bound table reader cannot change table_name".into(),
924        ));
925    }
926    Ok(name)
927}
928
929pub(crate) struct GlobalReaderState {
930    total_buckets: u32,
931    reader: Mutex<Reader>,
932    name: String,
933    metadata: TableMetadata,
934}
935
936impl GlobalReaderState {
937    fn new(reader: Reader, name: String, metadata: TableMetadata) -> Self {
938        let total_buckets = reader.current_global_snapshot().total_buckets;
939        Self {
940            total_buckets,
941            reader: Mutex::new(reader),
942            name,
943            metadata,
944        }
945    }
946
947    pub(crate) fn total_buckets(&self) -> u32 {
948        self.total_buckets
949    }
950
951    pub(crate) fn get(
952        &self,
953        bucket: u16,
954        key: &[u8],
955        options: &ReadOptions,
956    ) -> Result<Option<Vec<Option<bytes::Bytes>>>> {
957        let mut reader = lock_physical_reader(&self.reader)?;
958        Ok(reader.get_with_options(bucket, key, options)?)
959    }
960
961    pub(crate) fn multi_get(
962        &self,
963        keys: &[(u16, &[u8])],
964        options: &ReadOptions,
965    ) -> Result<Vec<Option<Vec<Option<bytes::Bytes>>>>> {
966        let mut reader = lock_physical_reader(&self.reader)?;
967        Ok(reader.multi_get_with_options(keys, options)?)
968    }
969
970    pub(crate) fn scan(
971        &self,
972        bucket: u16,
973        start: Option<&[u8]>,
974        end: Option<&[u8]>,
975        options: &ScanOptions,
976    ) -> Result<DbIterator> {
977        let mut reader = lock_physical_reader(&self.reader)?;
978        Ok(reader.scan_with_options_bounds(bucket, start, end, options)?)
979    }
980
981    pub(crate) fn scan_plan(&self) -> Result<crate::TableScanPlan> {
982        let reader = lock_physical_reader(&self.reader)?;
983        let snapshot = reader.current_global_snapshot();
984        crate::TableScanPlan::from_global_reader(
985            self.name.clone(),
986            self.metadata.clone(),
987            snapshot.id,
988            snapshot.total_buckets,
989            snapshot.shard_snapshots.clone(),
990            reader.config().clone(),
991        )
992    }
993
994    fn refreshed_snapshot(&self) -> Result<Option<(Reader, String, Option<TableId>)>> {
995        Ok(lock_physical_reader(&self.reader)?
996            .refreshed_snapshot()?
997            .map(|reader| {
998                (
999                    reader,
1000                    self.name.clone(),
1001                    self.metadata
1002                        .catalog_binding
1003                        .map(|binding| binding.table_id),
1004                )
1005            }))
1006    }
1007}
1008
1009fn global_metadata(reader: &Reader, name: &str) -> Result<TableMetadata> {
1010    let shard = reader
1011        .current_global_snapshot()
1012        .shard_snapshots
1013        .iter()
1014        .find(|shard| !shard.ranges.is_empty())
1015        .ok_or_else(|| TableError::InvalidSchema("global snapshot has no shard buckets".into()))?;
1016    load_table_metadata_for_shard(reader.config(), shard, name)
1017}
1018
1019fn lock_physical_reader(reader: &Mutex<Reader>) -> Result<std::sync::MutexGuard<'_, Reader>> {
1020    reader
1021        .lock()
1022        .map_err(|_| TableError::internal("table reader core lock poisoned"))
1023}
1024
1025fn validate_catalog_binding(
1026    metadata: &crate::metadata::TableMetadata,
1027    table_id: Option<TableId>,
1028) -> Result<()> {
1029    let Some(table_id) = table_id else {
1030        return Ok(());
1031    };
1032    let binding = metadata.catalog_binding.ok_or_else(|| {
1033        TableError::InvalidSchema("table metadata is not bound to a catalog table".into())
1034    })?;
1035    if binding.table_id != table_id {
1036        return Err(TableError::InvalidSchema(
1037            "table metadata belongs to another catalog table".into(),
1038        ));
1039    }
1040    Ok(())
1041}