Skip to main content

iceberg/
table.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18//! Table API for Apache Iceberg
19
20use 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
36/// Builder to create table scan.
37pub 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    /// required - sets the necessary FileIO to use for the table
65    pub fn file_io(mut self, file_io: FileIO) -> Self {
66        self.file_io = Some(file_io);
67        self
68    }
69
70    /// optional - sets the tables metadata location
71    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    /// required - passes in the TableMetadata to use for the Table
77    pub fn metadata<T: Into<TableMetadataRef>>(mut self, metadata: T) -> Self {
78        self.metadata = Some(metadata.into());
79        self
80    }
81
82    /// required - passes in the TableIdent to use for the Table
83    pub fn identifier(mut self, identifier: TableIdent) -> Self {
84        self.identifier = Some(identifier);
85        self
86    }
87
88    /// specifies if the Table is readonly or not (default not)
89    pub fn readonly(mut self, readonly: bool) -> Self {
90        self.readonly = readonly;
91        self
92    }
93
94    /// specifies if the Table's metadata cache will be disabled,
95    /// so that reads of Manifests and ManifestLists will never
96    /// get cached.
97    pub fn disable_cache(mut self) -> Self {
98        self.disable_cache = true;
99        self
100    }
101
102    /// optionally set a non-default metadata cache size
103    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    /// Set the Runtime for this table to use when spawning tasks.
109    pub fn runtime(mut self, runtime: Runtime) -> Self {
110        self.runtime = Some(runtime);
111        self
112    }
113
114    /// optional - sets the KMS client used to unwrap keys for table encryption.
115    ///
116    /// If the table metadata has the `encryption.key-id` property set, a
117    /// [`KeyManagementClient`] must be provided here so the table can build
118    /// an [`EncryptionManager`]; otherwise [`Self::build`] will return an error.
119    pub fn kms_client(mut self, kms_client: Arc<dyn KeyManagementClient>) -> Self {
120        self.kms_client = Some(kms_client);
121        self
122    }
123
124    /// build the Table
125    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/// Table represents a table in the catalog.
197#[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    /// Sets the [`Table`] metadata and returns an updated instance with the new metadata applied.
211    pub(crate) fn with_metadata(mut self, metadata: TableMetadataRef) -> Self {
212        self.metadata = metadata;
213        self
214    }
215
216    /// Sets the [`Table`] metadata location and returns an updated instance.
217    pub(crate) fn with_metadata_location(mut self, metadata_location: String) -> Self {
218        self.metadata_location = Some(metadata_location);
219        self
220    }
221
222    /// Returns a TableBuilder to build a table
223    pub fn builder() -> TableBuilder {
224        TableBuilder::new()
225    }
226
227    /// Returns table identifier.
228    pub fn identifier(&self) -> &TableIdent {
229        &self.identifier
230    }
231    /// Returns current metadata.
232    pub fn metadata(&self) -> &TableMetadata {
233        &self.metadata
234    }
235
236    /// Returns current metadata ref.
237    pub fn metadata_ref(&self) -> TableMetadataRef {
238        self.metadata.clone()
239    }
240
241    /// Returns current metadata location.
242    pub fn metadata_location(&self) -> Option<&str> {
243        self.metadata_location.as_deref()
244    }
245
246    /// Returns current metadata location in a result.
247    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    /// Returns file io used in this table.
255    pub fn file_io(&self) -> &FileIO {
256        &self.file_io
257    }
258
259    /// Returns this table's object cache
260    pub(crate) fn object_cache(&self) -> Arc<ObjectCache> {
261        self.object_cache.clone()
262    }
263
264    /// Returns the [`EncryptionManager`] for this table, if encryption is
265    /// configured.
266    ///
267    /// A manager is present iff the table metadata has the
268    /// `encryption.key-id` property set and a [`KeyManagementClient`] was
269    /// supplied to the [`TableBuilder`].
270    pub fn encryption_manager(&self) -> Option<&Arc<EncryptionManager>> {
271        self.encryption_manager.as_ref()
272    }
273
274    /// Creates a table scan.
275    pub fn scan(&self) -> TableScanBuilder<'_> {
276        TableScanBuilder::new(self)
277    }
278
279    /// Creates a metadata table which provides table-like APIs for inspecting metadata.
280    /// See [`MetadataTable`] for more details.
281    pub fn inspect(&self) -> MetadataTable<'_> {
282        MetadataTable::new(self)
283    }
284
285    /// Returns the [`Runtime`] for this table.
286    pub(crate) fn runtime(&self) -> &Runtime {
287        &self.runtime
288    }
289
290    /// Returns the flag indicating whether the `Table` is readonly or not
291    pub fn readonly(&self) -> bool {
292        self.readonly
293    }
294
295    /// Returns the current schema as a shared reference.
296    pub fn current_schema_ref(&self) -> SchemaRef {
297        self.metadata.current_schema().clone()
298    }
299
300    /// Creates a [`ManifestListReader`] for the given snapshot.
301    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    /// Creates a [`ManifestReader`] for loading manifests referenced by this table.
311    pub fn manifest_reader(&self) -> ManifestReader {
312        ManifestReader::new(self.file_io.clone())
313    }
314
315    /// Create a reader for the table.
316    pub fn reader_builder(&self) -> ArrowReaderBuilder {
317        ArrowReaderBuilder::new(self.file_io.clone(), self.runtime().clone())
318    }
319}
320
321/// `StaticTable` is a read-only table struct that can be created from a metadata file or from `TableMetaData` without a catalog.
322/// It can only be used to read metadata and for table scan.
323/// # Examples
324///
325/// ```rust, no_run
326/// # use iceberg::io::FileIO;
327/// # use iceberg::table::StaticTable;
328/// # use iceberg::TableIdent;
329/// # async fn example() {
330/// let metadata_file_location = "s3://bucket_name/path/to/metadata.json";
331/// let file_io = FileIO::new_with_fs();
332/// let static_identifier = TableIdent::from_strs(["static_ns", "static_table"]).unwrap();
333/// let static_table =
334///     StaticTable::from_metadata_file(&metadata_file_location, static_identifier, file_io)
335///         .await
336///         .unwrap();
337/// let snapshot_id = static_table
338///     .metadata()
339///     .current_snapshot()
340///     .unwrap()
341///     .snapshot_id();
342/// # }
343/// ```
344#[derive(Debug, Clone)]
345pub struct StaticTable(Table);
346
347impl StaticTable {
348    /// Creates a static table from a given `TableMetadata` and `FileIO`
349    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    /// Creates a static table directly from metadata file and `FileIO`
365    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    /// Create a TableScanBuilder for the static table.
385    pub fn scan(&self) -> TableScanBuilder<'_> {
386        self.0.scan()
387    }
388
389    /// Get TableMetadataRef for the static table
390    pub fn metadata(&self) -> TableMetadataRef {
391        self.0.metadata_ref()
392    }
393
394    /// Consumes the `StaticTable` and return it as a `Table`
395    /// Please use this method carefully as the Table it returns remains detached from a catalog
396    /// and can't be used to perform modifications on the table.
397    pub fn into_table(self) -> Table {
398        self.0
399    }
400
401    /// Create a reader for the table.
402    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        // The fixture contains a snapshot with key-id, encryption-keys (KEK + wrapped DEK),
509        // all generated with the master key bytes below.
510        let mut metadata: TableMetadata = load_test_metadata("TableMetadataV3ValidEncryption.json");
511
512        // Point the snapshot's manifest-list at the testdata file on disk.
513        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        // Seed the KMS with the same master key bytes used to generate the fixture.
523        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        // Encryption is a v3 spec feature; pre-v3 tables silently skip
571        // encryption even if encryption.key-id is set.
572        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}