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
82pub 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 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 pub fn bucket(mut self, bucket: u16) -> Self {
128 self.bucket = Some(bucket);
129 self
130 }
131
132 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 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 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 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
505pub 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 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 pub fn schema(&self) -> Arc<TableSchema> {
614 self.typed.load().schema_arc()
615 }
616
617 pub fn key_builder(&self) -> TableKeyBuilder {
619 self.typed.load().key_builder()
620 }
621
622 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 pub fn get(&self, key: &TableKey) -> Result<Option<Vec<Value>>> {
629 self.view_for_access()?.get(key)
630 }
631
632 pub fn multi_get(&self, keys: &[TableKey]) -> Result<Vec<Option<Vec<Value>>>> {
634 self.view_for_access()?.multi_get(keys)
635 }
636
637 pub fn scan(&self, bucket: u16) -> Result<TableScan> {
639 self.view_for_access()?.scan(bucket)
640 }
641
642 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 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 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, ¤t.typed) {
740 None
741 } else {
742 Some(crate::ffi::TableReaderView { typed })
743 }
744 }
745}
746
747pub 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 pub fn current_global_snapshot(mut self) -> Self {
794 self.selection = Some(TableReaderSelection::CurrentGlobal);
795 self
796 }
797
798 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
835pub 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 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}