1use std::sync::Arc;
21
22use crate::arrow::ArrowReaderBuilder;
23use crate::encryption::EncryptionManager;
24use crate::encryption::kms::KeyManagementClient;
25use crate::error::invalid_data;
26use crate::inspect::MetadataTable;
27use crate::io::FileIO;
28use crate::io::object_cache::ObjectCache;
29use crate::runtime::Runtime;
30use crate::scan::TableScanBuilder;
31use crate::spec::{
32 ManifestListReader, ManifestReader, SchemaRef, SnapshotRef, TableMetadata, TableMetadataRef,
33};
34use crate::{Result, TableIdent};
35
36pub struct TableBuilder {
38 file_io: Option<FileIO>,
39 metadata_location: Option<String>,
40 metadata: Option<TableMetadataRef>,
41 identifier: Option<TableIdent>,
42 kms_client: Option<Arc<dyn KeyManagementClient>>,
43 readonly: bool,
44 disable_cache: bool,
45 cache_size_bytes: Option<u64>,
46 runtime: Option<Runtime>,
47}
48
49impl TableBuilder {
50 pub(crate) fn new() -> Self {
51 Self {
52 file_io: None,
53 metadata_location: None,
54 metadata: None,
55 identifier: None,
56 kms_client: None,
57 readonly: false,
58 disable_cache: false,
59 cache_size_bytes: None,
60 runtime: None,
61 }
62 }
63
64 pub fn file_io(mut self, file_io: FileIO) -> Self {
66 self.file_io = Some(file_io);
67 self
68 }
69
70 pub fn metadata_location<T: Into<String>>(mut self, metadata_location: T) -> Self {
72 self.metadata_location = Some(metadata_location.into());
73 self
74 }
75
76 pub fn metadata<T: Into<TableMetadataRef>>(mut self, metadata: T) -> Self {
78 self.metadata = Some(metadata.into());
79 self
80 }
81
82 pub fn identifier(mut self, identifier: TableIdent) -> Self {
84 self.identifier = Some(identifier);
85 self
86 }
87
88 pub fn readonly(mut self, readonly: bool) -> Self {
90 self.readonly = readonly;
91 self
92 }
93
94 pub fn disable_cache(mut self) -> Self {
98 self.disable_cache = true;
99 self
100 }
101
102 pub fn cache_size_bytes(mut self, cache_size_bytes: u64) -> Self {
104 self.cache_size_bytes = Some(cache_size_bytes);
105 self
106 }
107
108 pub fn runtime(mut self, runtime: Runtime) -> Self {
110 self.runtime = Some(runtime);
111 self
112 }
113
114 pub fn kms_client(mut self, kms_client: Arc<dyn KeyManagementClient>) -> Self {
120 self.kms_client = Some(kms_client);
121 self
122 }
123
124 pub fn build(self) -> Result<Table> {
126 let Self {
127 file_io,
128 metadata_location,
129 metadata,
130 identifier,
131 kms_client,
132 readonly,
133 disable_cache,
134 cache_size_bytes,
135 runtime,
136 } = self;
137
138 let Some(file_io) = file_io else {
139 return Err(invalid_data!(
140 "FileIO must be provided with TableBuilder.file_io()"
141 ));
142 };
143
144 let Some(metadata) = metadata else {
145 return Err(invalid_data!(
146 "TableMetadataRef must be provided with TableBuilder.metadata()"
147 ));
148 };
149
150 let Some(identifier) = identifier else {
151 return Err(invalid_data!(
152 "TableIdent must be provided with TableBuilder.identifier()"
153 ));
154 };
155
156 let Some(runtime) = runtime else {
157 return Err(invalid_data!(
158 "Runtime must be provided with TableBuilder.runtime()"
159 ));
160 };
161
162 let encryption_manager =
163 EncryptionManager::from_table_metadata(kms_client.as_ref(), &metadata)?;
164
165 let object_cache = if disable_cache {
166 Arc::new(ObjectCache::with_disabled_cache(
167 file_io.clone(),
168 encryption_manager.clone(),
169 ))
170 } else if let Some(cache_size_bytes) = cache_size_bytes {
171 Arc::new(ObjectCache::new_with_capacity(
172 file_io.clone(),
173 cache_size_bytes,
174 encryption_manager.clone(),
175 ))
176 } else {
177 Arc::new(ObjectCache::new(
178 file_io.clone(),
179 encryption_manager.clone(),
180 ))
181 };
182
183 Ok(Table {
184 file_io,
185 metadata_location,
186 metadata,
187 identifier,
188 readonly,
189 object_cache,
190 runtime,
191 encryption_manager,
192 })
193 }
194}
195
196#[derive(Debug, Clone)]
198pub struct Table {
199 file_io: FileIO,
200 metadata_location: Option<String>,
201 metadata: TableMetadataRef,
202 identifier: TableIdent,
203 readonly: bool,
204 object_cache: Arc<ObjectCache>,
205 runtime: Runtime,
206 encryption_manager: Option<Arc<EncryptionManager>>,
207}
208
209impl Table {
210 pub(crate) fn with_metadata(mut self, metadata: TableMetadataRef) -> Self {
212 self.metadata = metadata;
213 self
214 }
215
216 pub(crate) fn with_metadata_location(mut self, metadata_location: String) -> Self {
218 self.metadata_location = Some(metadata_location);
219 self
220 }
221
222 pub fn builder() -> TableBuilder {
224 TableBuilder::new()
225 }
226
227 pub fn identifier(&self) -> &TableIdent {
229 &self.identifier
230 }
231 pub fn metadata(&self) -> &TableMetadata {
233 &self.metadata
234 }
235
236 pub fn metadata_ref(&self) -> TableMetadataRef {
238 self.metadata.clone()
239 }
240
241 pub fn metadata_location(&self) -> Option<&str> {
243 self.metadata_location.as_deref()
244 }
245
246 pub fn metadata_location_result(&self) -> Result<&str> {
248 self.metadata_location.as_deref().ok_or(invalid_data!(
249 "Metadata location does not exist for table: {}",
250 self.identifier
251 ))
252 }
253
254 pub fn file_io(&self) -> &FileIO {
256 &self.file_io
257 }
258
259 pub(crate) fn object_cache(&self) -> Arc<ObjectCache> {
261 self.object_cache.clone()
262 }
263
264 pub fn encryption_manager(&self) -> Option<&Arc<EncryptionManager>> {
271 self.encryption_manager.as_ref()
272 }
273
274 pub fn scan(&self) -> TableScanBuilder<'_> {
276 TableScanBuilder::new(self)
277 }
278
279 pub fn inspect(&self) -> MetadataTable<'_> {
282 MetadataTable::new(self)
283 }
284
285 pub(crate) fn runtime(&self) -> &Runtime {
287 &self.runtime
288 }
289
290 pub fn readonly(&self) -> bool {
292 self.readonly
293 }
294
295 pub fn current_schema_ref(&self) -> SchemaRef {
297 self.metadata.current_schema().clone()
298 }
299
300 pub fn manifest_list_reader(&self, snapshot: &SnapshotRef) -> ManifestListReader {
302 ManifestListReader::new(
303 snapshot.clone(),
304 self.file_io.clone(),
305 self.metadata.clone(),
306 self.encryption_manager.clone(),
307 )
308 }
309
310 pub fn manifest_reader(&self) -> ManifestReader {
312 ManifestReader::new(self.file_io.clone())
313 }
314
315 pub fn reader_builder(&self) -> ArrowReaderBuilder {
317 ArrowReaderBuilder::new(self.file_io.clone(), self.runtime().clone())
318 }
319}
320
321#[derive(Debug, Clone)]
345pub struct StaticTable(Table);
346
347impl StaticTable {
348 pub async fn from_metadata(
350 metadata: TableMetadata,
351 table_ident: TableIdent,
352 file_io: FileIO,
353 ) -> Result<Self> {
354 let table = Table::builder()
355 .metadata(metadata)
356 .identifier(table_ident)
357 .file_io(file_io.clone())
358 .runtime(Runtime::try_current()?)
359 .readonly(true)
360 .build();
361
362 Ok(Self(table?))
363 }
364 pub async fn from_metadata_file(
366 metadata_location: &str,
367 table_ident: TableIdent,
368 file_io: FileIO,
369 ) -> Result<Self> {
370 let metadata = TableMetadata::read_from(&file_io, metadata_location).await?;
371
372 let table = Table::builder()
373 .metadata(metadata)
374 .metadata_location(metadata_location)
375 .identifier(table_ident)
376 .file_io(file_io.clone())
377 .runtime(Runtime::try_current()?)
378 .readonly(true)
379 .build();
380
381 Ok(Self(table?))
382 }
383
384 pub fn scan(&self) -> TableScanBuilder<'_> {
386 self.0.scan()
387 }
388
389 pub fn metadata(&self) -> TableMetadataRef {
391 self.0.metadata_ref()
392 }
393
394 pub fn into_table(self) -> Table {
398 self.0
399 }
400
401 pub fn reader_builder(&self) -> ArrowReaderBuilder {
403 self.0.reader_builder()
404 }
405}
406
407#[cfg(test)]
408mod tests {
409 use std::fs;
410
411 use super::*;
412 use crate::ErrorKind;
413 use crate::encryption::SensitiveBytes;
414 use crate::encryption::kms::MemoryKeyManagementClient;
415 use crate::spec::TableProperties;
416
417 fn load_test_metadata(filename: &str) -> TableMetadata {
418 let path = format!(
419 "{}/testdata/table_metadata/{}",
420 env!("CARGO_MANIFEST_DIR"),
421 filename
422 );
423 let json = fs::read_to_string(path).unwrap();
424 serde_json::from_str(&json).unwrap()
425 }
426
427 #[tokio::test]
428 async fn test_static_table_from_file() {
429 let metadata_file_name = "TableMetadataV2Valid.json";
430 let metadata_file_path = format!(
431 "{}/testdata/table_metadata/{}",
432 env!("CARGO_MANIFEST_DIR"),
433 metadata_file_name
434 );
435 let file_io = FileIO::new_with_fs();
436 let static_identifier = TableIdent::from_strs(["static_ns", "static_table"]).unwrap();
437 let static_table =
438 StaticTable::from_metadata_file(&metadata_file_path, static_identifier, file_io)
439 .await
440 .unwrap();
441 let snapshot_id = static_table
442 .metadata()
443 .current_snapshot()
444 .unwrap()
445 .snapshot_id();
446 assert_eq!(
447 snapshot_id, 3055729675574597004,
448 "snapshot id from metadata don't match"
449 );
450 }
451
452 #[tokio::test]
453 async fn test_static_into_table() {
454 let metadata_file_name = "TableMetadataV2Valid.json";
455 let metadata_file_path = format!(
456 "{}/testdata/table_metadata/{}",
457 env!("CARGO_MANIFEST_DIR"),
458 metadata_file_name
459 );
460 let file_io = FileIO::new_with_fs();
461 let static_identifier = TableIdent::from_strs(["static_ns", "static_table"]).unwrap();
462 let static_table =
463 StaticTable::from_metadata_file(&metadata_file_path, static_identifier, file_io)
464 .await
465 .unwrap();
466 let table = static_table.into_table();
467 assert!(table.readonly());
468 assert_eq!(table.identifier.name(), "static_table");
469 assert_eq!(
470 table.metadata_location(),
471 Some(metadata_file_path).as_deref()
472 );
473 }
474
475 #[tokio::test]
476 async fn test_table_readonly_flag() {
477 let metadata_file_name = "TableMetadataV2Valid.json";
478 let metadata_file_path = format!(
479 "{}/testdata/table_metadata/{}",
480 env!("CARGO_MANIFEST_DIR"),
481 metadata_file_name
482 );
483 let file_io = FileIO::new_with_fs();
484 let metadata_file = file_io.new_input(metadata_file_path).unwrap();
485 let metadata_file_content = metadata_file.read().await.unwrap();
486 let table_metadata =
487 serde_json::from_slice::<TableMetadata>(&metadata_file_content).unwrap();
488 let static_identifier = TableIdent::from_strs(["ns", "table"]).unwrap();
489 let table = Table::builder()
490 .metadata(table_metadata)
491 .identifier(static_identifier)
492 .file_io(file_io.clone())
493 .runtime(Runtime::try_current().unwrap())
494 .build()
495 .unwrap();
496 assert!(!table.readonly());
497 assert_eq!(table.identifier.name(), "table");
498 }
499
500 fn make_kms() -> Arc<dyn KeyManagementClient> {
501 let kms = MemoryKeyManagementClient::new();
502 kms.add_master_key("master-1").unwrap();
503 Arc::new(kms)
504 }
505
506 #[tokio::test]
507 async fn table_decrypts_manifest_list_via_object_cache() {
508 let mut metadata: TableMetadata = load_test_metadata("TableMetadataV3ValidEncryption.json");
511
512 let manifest_list_path = format!(
514 "{}/testdata/manifests_lists/manifest-list-v3-encrypted.avro",
515 env!("CARGO_MANIFEST_DIR"),
516 );
517 let snapshot = metadata.snapshots.get_mut(&1).unwrap();
518 let mut patched = snapshot.as_ref().clone();
519 patched.manifest_list = manifest_list_path;
520 *snapshot = Arc::new(patched);
521
522 let kms: Arc<dyn KeyManagementClient> = {
524 let k = MemoryKeyManagementClient::new();
525 k.add_master_key_bytes(
526 "master-1",
527 SensitiveBytes::new([
528 0x00, 0x01, 0x02, 0x03, 0x04, 0x05, 0x06, 0x07, 0x08, 0x09, 0x0a, 0x0b, 0x0c,
529 0x0d, 0x0e, 0x0f,
530 ]),
531 )
532 .unwrap();
533 Arc::new(k)
534 };
535
536 let table = Table::builder()
537 .file_io(FileIO::new_with_fs())
538 .metadata(metadata)
539 .identifier(TableIdent::from_strs(["ns", "enc"]).unwrap())
540 .kms_client(kms)
541 .runtime(Runtime::try_current().unwrap())
542 .build()
543 .unwrap();
544
545 let snapshot_ref = table.metadata().current_snapshot().unwrap();
546 let manifest_list = table
547 .object_cache()
548 .get_manifest_list(snapshot_ref, &table.metadata_ref())
549 .await
550 .unwrap();
551 assert_eq!(manifest_list.entries().len(), 0);
552 }
553
554 #[tokio::test]
555 async fn table_builder_errors_when_encryption_key_id_set_but_no_kms() {
556 let metadata: TableMetadata = load_test_metadata("TableMetadataV3ValidEncryption.json");
557
558 let err = Table::builder()
559 .file_io(FileIO::new_with_memory())
560 .metadata(metadata)
561 .identifier(TableIdent::from_strs(["ns", "enc"]).unwrap())
562 .runtime(Runtime::try_current().unwrap())
563 .build()
564 .unwrap_err();
565 assert_eq!(err.kind(), ErrorKind::PreconditionFailed);
566 }
567
568 #[tokio::test]
569 async fn table_builder_skips_encryption_on_pre_v3_table() {
570 let mut metadata: TableMetadata = load_test_metadata("TableMetadataV2ValidMinimal.json");
573 metadata.properties.insert(
574 TableProperties::PROPERTY_ENCRYPTION_KEY_ID.to_string(),
575 "master-1".to_string(),
576 );
577
578 let table = Table::builder()
579 .file_io(FileIO::new_with_memory())
580 .metadata(metadata)
581 .identifier(TableIdent::from_strs(["ns", "enc"]).unwrap())
582 .kms_client(make_kms())
583 .runtime(Runtime::try_current().unwrap())
584 .build()
585 .unwrap();
586 assert!(table.encryption_manager().is_none());
587 }
588
589 #[tokio::test]
590 async fn table_builder_skips_encryption_when_property_absent() {
591 let metadata: TableMetadata = load_test_metadata("TableMetadataV3ValidMinimal.json");
592 let table = Table::builder()
593 .file_io(FileIO::new_with_memory())
594 .metadata(metadata)
595 .identifier(TableIdent::from_strs(["ns", "plain"]).unwrap())
596 .kms_client(make_kms())
597 .runtime(Runtime::try_current().unwrap())
598 .build()
599 .unwrap();
600 assert!(table.encryption_manager().is_none());
601 }
602}