From a7be7c558d99b9bcf2f73b230577fc4d260abc9d Mon Sep 17 00:00:00 2001 From: Henry Guo Date: Thu, 4 Jun 2026 15:50:13 +0800 Subject: [PATCH] feat(table-catalog): add object-backed catalog store (#3206) Co-authored-by: Henry Guo --- rustfs/src/table_catalog.rs | 1196 ++++++++++++++++++++++++++++++++++- 1 file changed, 1190 insertions(+), 6 deletions(-) diff --git a/rustfs/src/table_catalog.rs b/rustfs/src/table_catalog.rs index c4a6ea240..873ee86ce 100644 --- a/rustfs/src/table_catalog.rs +++ b/rustfs/src/table_catalog.rs @@ -20,13 +20,26 @@ #![allow(dead_code)] -use std::{collections::BTreeMap, fmt}; +use std::{ + collections::{BTreeMap, BTreeSet}, + fmt, + sync::Arc, +}; +use http::HeaderMap; use rustfs_ecstore::bucket::{ metadata::{BUCKET_TABLE_CONFIG, BUCKET_TABLE_RESERVED_PREFIX}, metadata_sys, }; -use serde::{Deserialize, Serialize}; +use rustfs_ecstore::error::StorageError; +use rustfs_ecstore::{ + set_disk::get_lock_acquire_timeout, + store_api::{HTTPPreconditions, ObjectOptions, PutObjReader, StorageAPI}, +}; +use serde::{Deserialize, Serialize, de::DeserializeOwned}; +use sha2::{Digest, Sha256}; +use tokio::io::AsyncReadExt; +use uuid::Uuid; pub(crate) const TABLE_BUCKET_MARKER_CONFIG: &str = BUCKET_TABLE_CONFIG; pub(crate) const RESERVED_CATALOG_OBJECT_MESSAGE: &str = "Object key is reserved for the table catalog"; @@ -47,6 +60,13 @@ const TABLE_MARKER_FILE: &str = "table.json"; const CURRENT_POINTER_FILE: &str = "current.json"; const LIFECYCLE_FILE: &str = "lifecycle.json"; const METADATA_DIR: &str = "metadata"; +const CATALOG_ROOT: &str = "catalog"; +const TABLE_BUCKET_ENTRY_FILE: &str = "table-bucket.json"; +const NAMESPACE_ENTRY_FILE: &str = "namespace-entry.json"; +const TABLE_ENTRY_FILE: &str = "table-entry.json"; +const COMMIT_LOG_ROOT: &str = "commits"; +const COMMIT_IDEMPOTENCY_ROOT: &str = "commit-idempotency"; +const TABLE_CATALOG_LIST_MAX_KEYS: i32 = 1000; #[derive(Debug, Clone, PartialEq, Eq)] pub enum CatalogIdentifierError { @@ -170,6 +190,7 @@ pub(crate) struct TableEntry { #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "SCREAMING_SNAKE_CASE")] pub(crate) enum CommitLogStatus { + Staged, Committed, Failed, } @@ -267,15 +288,743 @@ pub(crate) trait TableCatalogStore: Send + Sync { async fn drop_table(&self, table_bucket: &str, namespace: &str, table: &str) -> TableCatalogStoreResult<()>; - async fn get_commit_by_id(&self, table_id: &str, commit_id: &str) -> TableCatalogStoreResult>; + async fn get_commit_by_id( + &self, + table_bucket: &str, + table_id: &str, + commit_id: &str, + ) -> TableCatalogStoreResult>; async fn get_commit_by_idempotency_key( &self, + table_bucket: &str, table_id: &str, idempotency_key: &str, ) -> TableCatalogStoreResult>; } +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct TableCatalogObject { + pub data: Vec, + pub etag: Option, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) enum TableCatalogPutPrecondition { + Any, + IfAbsent, + IfMatch(String), +} + +#[async_trait::async_trait] +pub(crate) trait TableCatalogObjectBackend: Clone + Send + Sync + 'static { + async fn read_object(&self, bucket: &str, object: &str) -> TableCatalogStoreResult>; + + async fn object_exists(&self, bucket: &str, object: &str) -> TableCatalogStoreResult; + + async fn put_object( + &self, + bucket: &str, + object: &str, + data: Vec, + precondition: TableCatalogPutPrecondition, + ) -> TableCatalogStoreResult<()>; + + async fn delete_object(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<()>; + + async fn list_objects(&self, bucket: &str, prefix: &str) -> TableCatalogStoreResult>; + + async fn acquire_write_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult>; +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct TableCatalogObjectPaths { + reserved_prefix: &'static str, +} + +impl Default for TableCatalogObjectPaths { + fn default() -> Self { + Self { + reserved_prefix: TABLE_RESERVED_PREFIX, + } + } +} + +impl TableCatalogObjectPaths { + pub fn table_bucket_entry_path(&self) -> String { + format!("{}{}", self.catalog_root_prefix(), TABLE_BUCKET_ENTRY_FILE) + } + + pub fn namespace_entries_prefix(&self) -> String { + format!("{}{}/", self.catalog_root_prefix(), NAMESPACE_ROOT) + } + + pub fn namespace_entry_path(&self, namespace: &Namespace) -> String { + format!("{}{}/{}", self.namespace_entries_prefix(), namespace.storage_id(), NAMESPACE_ENTRY_FILE) + } + + pub fn table_entries_prefix(&self, namespace: &Namespace) -> String { + format!("{}{}/{}/", self.namespace_entries_prefix(), namespace.storage_id(), TABLE_ROOT) + } + + pub fn table_entry_path(&self, namespace: &Namespace, table: &IdentifierSegment) -> String { + format!("{}{}/{}", self.table_entries_prefix(namespace), table.as_str(), TABLE_ENTRY_FILE) + } + + pub fn commit_log_entry_path(&self, table_id: &str, commit_id: &str) -> String { + format!( + "{}{}/{}/{}.json", + self.catalog_root_prefix(), + COMMIT_LOG_ROOT, + catalog_path_hash(table_id), + catalog_path_hash(commit_id) + ) + } + + pub fn commit_idempotency_entry_path(&self, table_id: &str, idempotency_key: &str) -> String { + format!( + "{}{}/{}/{}.json", + self.catalog_root_prefix(), + COMMIT_IDEMPOTENCY_ROOT, + catalog_path_hash(table_id), + catalog_path_hash(idempotency_key) + ) + } + + fn catalog_root_prefix(&self) -> String { + format!("{}/{}/{}/{}/", self.reserved_prefix, WAREHOUSE_ROOT, DEFAULT_WAREHOUSE_ID, CATALOG_ROOT) + } +} + +#[derive(Clone)] +pub(crate) struct ObjectTableCatalogStore { + backend: B, + paths: TableCatalogObjectPaths, +} + +impl ObjectTableCatalogStore +where + B: TableCatalogObjectBackend, +{ + pub fn new(backend: B) -> Self { + Self { + backend, + paths: TableCatalogObjectPaths::default(), + } + } + + async fn read_entry(&self, bucket: &str, object: &str) -> TableCatalogStoreResult)>> + where + T: DeserializeOwned, + { + let Some(object_data) = self.backend.read_object(bucket, object).await? else { + return Ok(None); + }; + + let entry = serde_json::from_slice(&object_data.data) + .map_err(|err| TableCatalogStoreError::Invalid(format!("failed to parse catalog entry {object}: {err}")))?; + Ok(Some((entry, object_data.etag))) + } + + async fn write_entry( + &self, + bucket: &str, + object: &str, + entry: &T, + precondition: TableCatalogPutPrecondition, + ) -> TableCatalogStoreResult<()> + where + T: Serialize, + { + let data = serde_json::to_vec(entry) + .map_err(|err| TableCatalogStoreError::Internal(format!("failed to serialize catalog entry {object}: {err}")))?; + self.backend.put_object(bucket, object, data, precondition).await + } + + async fn require_table_bucket(&self, table_bucket: &str) -> TableCatalogStoreResult<()> { + if self.get_table_bucket(table_bucket).await?.is_none() { + return Err(TableCatalogStoreError::NotFound(format!("table bucket {table_bucket}"))); + } + Ok(()) + } + + async fn read_table_with_etag( + &self, + table_bucket: &str, + namespace: &Namespace, + table: &IdentifierSegment, + ) -> TableCatalogStoreResult> { + let table_path = self.paths.table_entry_path(namespace, table); + let Some((entry, etag)) = self.read_entry::(table_bucket, &table_path).await? else { + return Ok(None); + }; + let Some(etag) = etag else { + return Err(TableCatalogStoreError::Internal(format!("catalog table entry has no etag: {table_path}"))); + }; + Ok(Some((entry, etag))) + } + + async fn write_table_entry( + &self, + entry: TableEntry, + precondition: TableCatalogPutPrecondition, + ) -> TableCatalogStoreResult<()> { + validate_catalog_entry_version("table", entry.version)?; + self.require_table_bucket(&entry.table_bucket).await?; + let namespace = parse_namespace_for_store(&entry.namespace)?; + let table = parse_table_for_store(&entry.table)?; + if self.get_namespace(&entry.table_bucket, &entry.namespace).await?.is_none() { + return Err(TableCatalogStoreError::NotFound(format!( + "namespace {}/{}", + entry.table_bucket, entry.namespace + ))); + } + let table_path = self.paths.table_entry_path(&namespace, &table); + self.write_entry(&entry.table_bucket, &table_path, &entry, precondition).await + } + + async fn read_commit_by_path(&self, table_bucket: &str, object: &str) -> TableCatalogStoreResult> { + self.read_entry::(table_bucket, object) + .await + .map(|entry| entry.map(|(commit, _)| commit)) + } + + async fn finalize_commit_log( + &self, + table_bucket: &str, + commit_path: &str, + idempotency_path: Option<&str>, + commit_log: &CommitLogEntry, + ) -> TableCatalogStoreResult<()> { + self.write_entry(table_bucket, commit_path, commit_log, TableCatalogPutPrecondition::Any) + .await?; + if let Some(idempotency_path) = idempotency_path { + self.write_entry(table_bucket, idempotency_path, commit_log, TableCatalogPutPrecondition::Any) + .await?; + } + Ok(()) + } +} + +#[async_trait::async_trait] +impl TableCatalogStore for ObjectTableCatalogStore +where + B: TableCatalogObjectBackend, +{ + async fn get_table_bucket(&self, table_bucket: &str) -> TableCatalogStoreResult> { + self.read_entry::(table_bucket, &self.paths.table_bucket_entry_path()) + .await + .map(|entry| entry.map(|(bucket, _)| bucket)) + } + + async fn put_table_bucket(&self, entry: TableBucketEntry) -> TableCatalogStoreResult<()> { + validate_catalog_entry_version("table bucket", entry.version)?; + if entry.table_bucket.is_empty() { + return Err(TableCatalogStoreError::Invalid("table bucket name cannot be empty".to_string())); + } + if entry.catalog_type != TABLE_BUCKET_CATALOG_TYPE { + return Err(TableCatalogStoreError::Invalid("unsupported table bucket catalog type".to_string())); + } + + self.write_entry( + &entry.table_bucket, + &self.paths.table_bucket_entry_path(), + &entry, + TableCatalogPutPrecondition::Any, + ) + .await + } + + async fn create_namespace(&self, entry: NamespaceEntry) -> TableCatalogStoreResult<()> { + validate_catalog_entry_version("namespace", entry.version)?; + self.require_table_bucket(&entry.table_bucket).await?; + let namespace = parse_namespace_for_store(&entry.namespace)?; + let object = self.paths.namespace_entry_path(&namespace); + self.write_entry(&entry.table_bucket, &object, &entry, TableCatalogPutPrecondition::IfAbsent) + .await + } + + async fn list_namespaces(&self, table_bucket: &str) -> TableCatalogStoreResult> { + let mut entries = Vec::new(); + for object in self + .backend + .list_objects(table_bucket, &self.paths.namespace_entries_prefix()) + .await? + { + if !object.ends_with(NAMESPACE_ENTRY_FILE) { + continue; + } + if let Some((entry, _)) = self.read_entry::(table_bucket, &object).await? { + entries.push(entry); + } + } + entries.sort_by(|left, right| left.namespace.cmp(&right.namespace)); + Ok(entries) + } + + async fn get_namespace(&self, table_bucket: &str, namespace: &str) -> TableCatalogStoreResult> { + let namespace = parse_namespace_for_store(namespace)?; + self.read_entry::(table_bucket, &self.paths.namespace_entry_path(&namespace)) + .await + .map(|entry| entry.map(|(namespace, _)| namespace)) + } + + async fn drop_namespace(&self, table_bucket: &str, namespace: &str) -> TableCatalogStoreResult<()> { + let namespace = parse_namespace_for_store(namespace)?; + if self.get_namespace(table_bucket, &namespace.public_name()).await?.is_none() { + return Err(TableCatalogStoreError::NotFound(format!( + "namespace {}/{}", + table_bucket, + namespace.public_name() + ))); + } + if !self.list_tables(table_bucket, &namespace.public_name()).await?.is_empty() { + return Err(TableCatalogStoreError::Conflict(format!( + "namespace {}/{} is not empty", + table_bucket, + namespace.public_name() + ))); + } + self.backend + .delete_object(table_bucket, &self.paths.namespace_entry_path(&namespace)) + .await + } + + async fn create_table(&self, entry: TableEntry) -> TableCatalogStoreResult<()> { + self.write_table_entry(entry, TableCatalogPutPrecondition::IfAbsent).await + } + + async fn register_table(&self, entry: TableEntry) -> TableCatalogStoreResult<()> { + self.write_table_entry(entry, TableCatalogPutPrecondition::IfAbsent).await + } + + async fn list_tables(&self, table_bucket: &str, namespace: &str) -> TableCatalogStoreResult> { + let namespace = parse_namespace_for_store(namespace)?; + let mut entries = Vec::new(); + for object in self + .backend + .list_objects(table_bucket, &self.paths.table_entries_prefix(&namespace)) + .await? + { + if !object.ends_with(TABLE_ENTRY_FILE) { + continue; + } + if let Some((entry, _)) = self.read_entry::(table_bucket, &object).await? { + entries.push(entry); + } + } + entries.sort_by(|left, right| left.table.cmp(&right.table)); + Ok(entries) + } + + async fn load_table(&self, table_bucket: &str, namespace: &str, table: &str) -> TableCatalogStoreResult> { + let namespace = parse_namespace_for_store(namespace)?; + let table = parse_table_for_store(table)?; + self.read_entry::(table_bucket, &self.paths.table_entry_path(&namespace, &table)) + .await + .map(|entry| entry.map(|(table, _)| table)) + } + + async fn commit_table(&self, request: TableCommitRequest) -> TableCatalogStoreResult { + let namespace = parse_namespace_for_store(&request.namespace)?; + let table = parse_table_for_store(&request.table)?; + let table_path = self.paths.table_entry_path(&namespace, &table); + let _guard = self.backend.acquire_write_lock(&request.table_bucket, &table_path).await?; + + let Some((current, current_etag)) = self.read_table_with_etag(&request.table_bucket, &namespace, &table).await? else { + return Err(TableCatalogStoreError::NotFound(format!( + "table {}/{}/{}", + request.table_bucket, request.namespace, request.table + ))); + }; + + let commit_path = self.paths.commit_log_entry_path(¤t.table_id, &request.commit_id); + let existing_commit = self.read_commit_by_path(&request.table_bucket, &commit_path).await?; + let idempotency_path = request + .idempotency_key + .as_deref() + .map(|idempotency_key| self.paths.commit_idempotency_entry_path(¤t.table_id, idempotency_key)); + let existing_idempotency_commit = match idempotency_path.as_deref() { + Some(idempotency_path) => self.read_commit_by_path(&request.table_bucket, idempotency_path).await?, + None => None, + }; + + if let Some(existing) = existing_commit.as_ref() { + if !commit_log_matches_request(existing, &request, ¤t.table_id) { + return Err(TableCatalogStoreError::Conflict(format!( + "commit id already exists: {}", + request.commit_id + ))); + } + if matches!(existing.status, CommitLogStatus::Committed) || table_matches_committed_log(¤t, existing) { + let mut committed = existing.clone(); + committed.status = CommitLogStatus::Committed; + let _ = self + .finalize_commit_log(&request.table_bucket, &commit_path, idempotency_path.as_deref(), &committed) + .await; + return Ok(TableCommitResult { + table: current, + commit_log: committed, + }); + } + if !matches!(existing.status, CommitLogStatus::Staged) || !table_matches_staged_base(¤t, existing) { + return Err(TableCatalogStoreError::Conflict( + "existing commit record does not match current table state".to_string(), + )); + } + } + if let Some(existing) = existing_idempotency_commit.as_ref() + && !commit_log_matches_request(existing, &request, ¤t.table_id) + { + return Err(TableCatalogStoreError::Conflict("idempotency key already exists".to_string())); + } + if existing_commit.is_none() && existing_idempotency_commit.is_some() { + return Err(TableCatalogStoreError::Conflict( + "idempotency key exists without a recoverable commit record".to_string(), + )); + } + + if current.version_token != request.expected_version_token { + return Err(TableCatalogStoreError::Conflict( + "current table version token does not match expected token".to_string(), + )); + } + if current.metadata_location != request.expected_metadata_location { + return Err(TableCatalogStoreError::Conflict( + "current table metadata location does not match expected location".to_string(), + )); + } + if !is_valid_table_metadata_location(&namespace, &table, &request.new_metadata_location) { + return Err(TableCatalogStoreError::Invalid( + "new metadata location must be inside the table metadata directory".to_string(), + )); + } + if !self + .backend + .object_exists(&request.table_bucket, &request.new_metadata_location) + .await? + { + return Err(TableCatalogStoreError::NotFound(format!( + "new metadata object {}", + request.new_metadata_location + ))); + } + + let has_existing_commit = existing_commit.is_some(); + let mut staged_commit_log = existing_commit.unwrap_or_else(|| CommitLogEntry { + version: TABLE_CATALOG_ENTRY_VERSION, + commit_id: request.commit_id.clone(), + idempotency_key: request.idempotency_key.clone(), + table_id: current.table_id.clone(), + operation: request.operation.clone(), + expected_version_token: request.expected_version_token.clone(), + new_version_token: format!("token-{}", Uuid::new_v4()), + previous_metadata_location: current.metadata_location.clone(), + new_metadata_location: request.new_metadata_location.clone(), + requirements: request.requirements.clone(), + status: CommitLogStatus::Staged, + writer: request.writer.clone(), + created_at: None, + updated_at: None, + }); + staged_commit_log.status = CommitLogStatus::Staged; + + let mut next = current.clone(); + next.metadata_location = staged_commit_log.new_metadata_location.clone(); + next.version_token = staged_commit_log.new_version_token.clone(); + next.generation = current.generation.saturating_add(1); + + if !has_existing_commit { + self.write_entry( + &request.table_bucket, + &commit_path, + &staged_commit_log, + TableCatalogPutPrecondition::IfAbsent, + ) + .await?; + } + if let Some(idempotency_path) = idempotency_path.as_deref() + && existing_idempotency_commit.is_none() + { + self.write_entry( + &request.table_bucket, + idempotency_path, + &staged_commit_log, + TableCatalogPutPrecondition::IfAbsent, + ) + .await?; + } + + self.write_entry( + &request.table_bucket, + &table_path, + &next, + TableCatalogPutPrecondition::IfMatch(current_etag), + ) + .await?; + + let mut commit_log = staged_commit_log; + commit_log.status = CommitLogStatus::Committed; + // After the table CAS succeeds, the staged record is the durable recovery source. + // A finalization failure must not turn an externally committed pointer into a failed commit response. + let _ = self + .finalize_commit_log(&request.table_bucket, &commit_path, idempotency_path.as_deref(), &commit_log) + .await; + + Ok(TableCommitResult { table: next, commit_log }) + } + + async fn drop_table(&self, table_bucket: &str, namespace: &str, table: &str) -> TableCatalogStoreResult<()> { + let namespace = parse_namespace_for_store(namespace)?; + let table = parse_table_for_store(table)?; + let object = self.paths.table_entry_path(&namespace, &table); + if self + .load_table(table_bucket, &namespace.public_name(), table.as_str()) + .await? + .is_none() + { + return Err(TableCatalogStoreError::NotFound(format!( + "table {}/{}/{}", + table_bucket, + namespace.public_name(), + table.as_str() + ))); + } + self.backend.delete_object(table_bucket, &object).await + } + + async fn get_commit_by_id( + &self, + table_bucket: &str, + table_id: &str, + commit_id: &str, + ) -> TableCatalogStoreResult> { + let object = self.paths.commit_log_entry_path(table_id, commit_id); + self.read_commit_by_path(table_bucket, &object).await + } + + async fn get_commit_by_idempotency_key( + &self, + table_bucket: &str, + table_id: &str, + idempotency_key: &str, + ) -> TableCatalogStoreResult> { + let object = self.paths.commit_idempotency_entry_path(table_id, idempotency_key); + self.read_commit_by_path(table_bucket, &object).await + } +} + +pub(crate) struct EcStoreTableCatalogObjectBackend { + store: Arc, +} + +impl Clone for EcStoreTableCatalogObjectBackend { + fn clone(&self) -> Self { + Self { + store: self.store.clone(), + } + } +} + +impl EcStoreTableCatalogObjectBackend +where + S: StorageAPI, +{ + pub fn new(store: Arc) -> Self { + Self { store } + } +} + +pub(crate) type EcStoreTableCatalogStore = ObjectTableCatalogStore>; + +#[async_trait::async_trait] +impl TableCatalogObjectBackend for EcStoreTableCatalogObjectBackend +where + S: StorageAPI, +{ + async fn read_object(&self, bucket: &str, object: &str) -> TableCatalogStoreResult> { + let info = match self.store.get_object_info(bucket, object, &ObjectOptions::default()).await { + Ok(info) => info, + Err(err) if is_missing_storage_error(&err) => return Ok(None), + Err(err) => return Err(storage_error_to_catalog("read catalog object info", err)), + }; + let mut reader = match self + .store + .get_object_reader(bucket, object, None, HeaderMap::new(), &ObjectOptions::default()) + .await + { + Ok(reader) => reader, + Err(err) if is_missing_storage_error(&err) => return Ok(None), + Err(err) => return Err(storage_error_to_catalog("read catalog object", err)), + }; + let mut data = Vec::new(); + reader + .stream + .read_to_end(&mut data) + .await + .map_err(|err| TableCatalogStoreError::Internal(format!("failed to read catalog object {bucket}/{object}: {err}")))?; + Ok(Some(TableCatalogObject { data, etag: info.etag })) + } + + async fn object_exists(&self, bucket: &str, object: &str) -> TableCatalogStoreResult { + match self.store.get_object_info(bucket, object, &ObjectOptions::default()).await { + Ok(_) => Ok(true), + Err(err) if is_missing_storage_error(&err) => Ok(false), + Err(err) => Err(storage_error_to_catalog("check catalog object", err)), + } + } + + async fn put_object( + &self, + bucket: &str, + object: &str, + data: Vec, + precondition: TableCatalogPutPrecondition, + ) -> TableCatalogStoreResult<()> { + let mut reader = PutObjReader::from_vec(data); + let opts = ObjectOptions { + http_preconditions: http_preconditions_for_catalog_put(precondition), + ..Default::default() + }; + self.store + .put_object(bucket, object, &mut reader, &opts) + .await + .map(|_| ()) + .map_err(|err| storage_error_to_catalog("write catalog object", err)) + } + + async fn delete_object(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<()> { + match self.store.delete_object(bucket, object, ObjectOptions::default()).await { + Ok(_) => Ok(()), + Err(err) if is_missing_storage_error(&err) => Ok(()), + Err(err) => Err(storage_error_to_catalog("delete catalog object", err)), + } + } + + async fn list_objects(&self, bucket: &str, prefix: &str) -> TableCatalogStoreResult> { + let mut continuation = None; + let mut objects = BTreeSet::new(); + + loop { + let result = self + .store + .clone() + .list_objects_v2(bucket, prefix, continuation, None, TABLE_CATALOG_LIST_MAX_KEYS, false, None, false) + .await + .map_err(|err| storage_error_to_catalog("list catalog objects", err))?; + + for object in result.objects { + objects.insert(object.name); + } + + if !result.is_truncated { + break; + } + + let Some(next) = result.next_continuation_token else { + break; + }; + continuation = Some(next); + } + + Ok(objects.into_iter().collect()) + } + + async fn acquire_write_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult> { + let lock = self + .store + .new_ns_lock(bucket, object) + .await + .map_err(|err| storage_error_to_catalog("create catalog table lock", err))?; + let guard = lock + .get_write_lock(get_lock_acquire_timeout()) + .await + .map_err(|err| TableCatalogStoreError::Internal(format!("failed to acquire catalog table lock: {err}")))?; + Ok(Box::new(guard)) + } +} + +fn parse_namespace_for_store(namespace: &str) -> TableCatalogStoreResult { + Namespace::parse(namespace).map_err(|err| TableCatalogStoreError::Invalid(format!("invalid namespace: {err}"))) +} + +fn parse_table_for_store(table: &str) -> TableCatalogStoreResult { + IdentifierSegment::parse(table).map_err(|err| TableCatalogStoreError::Invalid(format!("invalid table name: {err}"))) +} + +fn catalog_path_hash(value: &str) -> String { + let digest = Sha256::digest(value.as_bytes()); + let mut output = String::with_capacity(digest.len() * 2); + const HEX: &[u8; 16] = b"0123456789abcdef"; + for byte in digest { + output.push(char::from(HEX[usize::from(byte >> 4)])); + output.push(char::from(HEX[usize::from(byte & 0x0f)])); + } + output +} + +fn validate_catalog_entry_version(kind: &str, version: u16) -> TableCatalogStoreResult<()> { + if version != TABLE_CATALOG_ENTRY_VERSION { + return Err(TableCatalogStoreError::Invalid(format!("unsupported {kind} entry version"))); + } + Ok(()) +} + +fn commit_log_matches_request(commit_log: &CommitLogEntry, request: &TableCommitRequest, table_id: &str) -> bool { + commit_log.version == TABLE_CATALOG_ENTRY_VERSION + && commit_log.commit_id == request.commit_id + && commit_log.idempotency_key == request.idempotency_key + && commit_log.table_id == table_id + && commit_log.operation == request.operation + && commit_log.expected_version_token == request.expected_version_token + && commit_log.previous_metadata_location == request.expected_metadata_location + && commit_log.new_metadata_location == request.new_metadata_location + && commit_log.requirements == request.requirements + && commit_log.writer == request.writer +} + +fn table_matches_committed_log(table: &TableEntry, commit_log: &CommitLogEntry) -> bool { + table.table_id == commit_log.table_id + && table.metadata_location == commit_log.new_metadata_location + && table.version_token == commit_log.new_version_token +} + +fn table_matches_staged_base(table: &TableEntry, commit_log: &CommitLogEntry) -> bool { + table.table_id == commit_log.table_id + && table.metadata_location == commit_log.previous_metadata_location + && table.version_token == commit_log.expected_version_token +} + +fn http_preconditions_for_catalog_put(precondition: TableCatalogPutPrecondition) -> Option { + match precondition { + TableCatalogPutPrecondition::Any => None, + TableCatalogPutPrecondition::IfAbsent => Some(HTTPPreconditions { + if_none_match: Some("*".to_string()), + ..Default::default() + }), + TableCatalogPutPrecondition::IfMatch(etag) => Some(HTTPPreconditions { + if_match: Some(etag), + ..Default::default() + }), + } +} + +fn is_missing_storage_error(err: &StorageError) -> bool { + matches!( + err, + StorageError::ObjectNotFound(_, _) | StorageError::FileNotFound | StorageError::ConfigNotFound + ) +} + +fn storage_error_to_catalog(action: &str, err: StorageError) -> TableCatalogStoreError { + match err { + StorageError::ObjectNotFound(bucket, object) => TableCatalogStoreError::NotFound(format!("{action}: {bucket}/{object}")), + StorageError::BucketNotFound(bucket) => TableCatalogStoreError::NotFound(format!("{action}: bucket {bucket}")), + StorageError::PreconditionFailed => TableCatalogStoreError::Conflict(format!("{action}: precondition failed")), + other => TableCatalogStoreError::Internal(format!("{action}: {other}")), + } +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize)] pub(crate) struct NamespaceMarker { pub version: u16, @@ -878,12 +1627,18 @@ mod tests { Ok(()) } - async fn get_commit_by_id(&self, _table_id: &str, _commit_id: &str) -> TableCatalogStoreResult> { + async fn get_commit_by_id( + &self, + _table_bucket: &str, + _table_id: &str, + _commit_id: &str, + ) -> TableCatalogStoreResult> { Ok(None) } async fn get_commit_by_idempotency_key( &self, + _table_bucket: &str, _table_id: &str, _idempotency_key: &str, ) -> TableCatalogStoreResult> { @@ -897,10 +1652,16 @@ mod tests { assert!(store.get_table_bucket("analytics").await.unwrap().is_none()); assert!(store.list_namespaces("analytics").await.unwrap().is_empty()); - assert!(store.get_commit_by_id("table-id", "commit-id").await.unwrap().is_none()); assert!( store - .get_commit_by_idempotency_key("table-id", "client-request-id") + .get_commit_by_id("analytics", "table-id", "commit-id") + .await + .unwrap() + .is_none() + ); + assert!( + store + .get_commit_by_idempotency_key("analytics", "table-id", "client-request-id") .await .unwrap() .is_none() @@ -931,6 +1692,429 @@ mod tests { assert_eq!(result.commit_log.status, CommitLogStatus::Committed); } + #[test] + fn catalog_object_entry_paths_use_reserved_prefix_and_hashed_untrusted_ids() { + let paths = TableCatalogObjectPaths::default(); + let namespace = Namespace::parse("analytics.daily_events").unwrap(); + let table = IdentifierSegment::parse("events").unwrap(); + + assert_eq!( + paths.table_bucket_entry_path(), + ".rustfs-table/warehouses/default/catalog/table-bucket.json" + ); + assert_eq!( + paths.namespace_entry_path(&namespace), + ".rustfs-table/warehouses/default/catalog/namespaces/analytics/daily_events/namespace-entry.json" + ); + assert_eq!( + paths.table_entry_path(&namespace, &table), + ".rustfs-table/warehouses/default/catalog/namespaces/analytics/daily_events/tables/events/table-entry.json" + ); + + let commit_path = paths.commit_log_entry_path("table/../id", "commit/%2f\nid"); + let idempotency_path = paths.commit_idempotency_entry_path("table/../id", "client/%2f\nrequest"); + + for path in [commit_path, idempotency_path] { + assert!(path.starts_with(".rustfs-table/warehouses/default/catalog/")); + assert!(path.ends_with(".json")); + assert!(!path.contains("..")); + assert!(!path.contains('%')); + assert!(!path.contains('\n')); + assert!(!path.contains("table/../id")); + assert!(!path.contains("client/%2f")); + } + } + + #[derive(Clone, Default)] + struct TestCatalogObjectBackend { + state: std::sync::Arc>, + write_lock: std::sync::Arc>, + } + + #[derive(Default)] + struct TestCatalogObjectState { + objects: BTreeMap<(String, String), TestCatalogObjectRecord>, + fail_put_attempts: BTreeMap<(String, String), BTreeSet>, + put_attempts: BTreeMap<(String, String), usize>, + next_etag: u64, + } + + #[derive(Clone)] + struct TestCatalogObjectRecord { + data: Vec, + etag: String, + } + + impl TestCatalogObjectBackend { + async fn seed_object(&self, bucket: &str, object: &str, data: Vec) { + let mut state = self.state.lock().await; + let etag = state.next_etag(); + state + .objects + .insert((bucket.to_string(), object.to_string()), TestCatalogObjectRecord { data, etag }); + } + + async fn fail_put_attempt(&self, bucket: &str, object: &str, attempt: usize) { + let mut state = self.state.lock().await; + state + .fail_put_attempts + .entry((bucket.to_string(), object.to_string())) + .or_default() + .insert(attempt); + } + } + + impl TestCatalogObjectState { + fn next_etag(&mut self) -> String { + self.next_etag += 1; + format!("etag-{}", self.next_etag) + } + } + + #[async_trait::async_trait] + impl TableCatalogObjectBackend for TestCatalogObjectBackend { + async fn read_object(&self, bucket: &str, object: &str) -> TableCatalogStoreResult> { + let state = self.state.lock().await; + Ok(state + .objects + .get(&(bucket.to_string(), object.to_string())) + .map(|record| TableCatalogObject { + data: record.data.clone(), + etag: Some(record.etag.clone()), + })) + } + + async fn object_exists(&self, bucket: &str, object: &str) -> TableCatalogStoreResult { + let state = self.state.lock().await; + Ok(state.objects.contains_key(&(bucket.to_string(), object.to_string()))) + } + + async fn put_object( + &self, + bucket: &str, + object: &str, + data: Vec, + precondition: TableCatalogPutPrecondition, + ) -> TableCatalogStoreResult<()> { + let mut state = self.state.lock().await; + let key = (bucket.to_string(), object.to_string()); + let attempt = { + let attempts = state.put_attempts.entry(key.clone()).or_default(); + *attempts += 1; + *attempts + }; + if state + .fail_put_attempts + .get(&key) + .is_some_and(|attempts| attempts.contains(&attempt)) + { + return Err(TableCatalogStoreError::Internal(format!( + "injected put failure for {object} attempt {attempt}" + ))); + } + match precondition { + TableCatalogPutPrecondition::IfAbsent if state.objects.contains_key(&key) => { + return Err(TableCatalogStoreError::Conflict(format!("object already exists: {object}"))); + } + TableCatalogPutPrecondition::IfMatch(expected) => { + let Some(current) = state.objects.get(&key) else { + return Err(TableCatalogStoreError::Conflict(format!("object is missing: {object}"))); + }; + if current.etag != expected { + return Err(TableCatalogStoreError::Conflict(format!("object changed: {object}"))); + } + } + _ => {} + } + + let etag = state.next_etag(); + state.objects.insert(key, TestCatalogObjectRecord { data, etag }); + Ok(()) + } + + async fn delete_object(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<()> { + let mut state = self.state.lock().await; + state.objects.remove(&(bucket.to_string(), object.to_string())); + Ok(()) + } + + async fn list_objects(&self, bucket: &str, prefix: &str) -> TableCatalogStoreResult> { + let state = self.state.lock().await; + Ok(state + .objects + .keys() + .filter(|(entry_bucket, object)| entry_bucket == bucket && object.starts_with(prefix)) + .map(|(_, object)| object.clone()) + .collect()) + } + + async fn acquire_write_lock(&self, _bucket: &str, _object: &str) -> TableCatalogStoreResult> { + Ok(Box::new(self.write_lock.clone().lock_owned().await)) + } + } + + fn test_bucket_entry(bucket: &str) -> TableBucketEntry { + TableBucketEntry { + version: TABLE_CATALOG_ENTRY_VERSION, + table_bucket: bucket.to_string(), + catalog_type: TABLE_BUCKET_CATALOG_TYPE.to_string(), + warehouse_root: format!("s3://{bucket}/"), + state: TableCatalogEntryState::Active, + properties: BTreeMap::new(), + created_at: None, + updated_at: None, + } + } + + fn test_namespace_entry(bucket: &str, namespace: &Namespace) -> NamespaceEntry { + NamespaceEntry { + version: TABLE_CATALOG_ENTRY_VERSION, + table_bucket: bucket.to_string(), + namespace: namespace.public_name(), + namespace_id: namespace.storage_id(), + state: TableCatalogEntryState::Active, + properties: BTreeMap::new(), + created_at: None, + updated_at: None, + } + } + + fn test_table_entry(bucket: &str, namespace: &Namespace, table: &IdentifierSegment, metadata_location: String) -> TableEntry { + TableEntry { + version: TABLE_CATALOG_ENTRY_VERSION, + table_bucket: bucket.to_string(), + namespace: namespace.public_name(), + table: table.as_str().to_string(), + table_id: "table-id".to_string(), + table_uuid: "table-uuid".to_string(), + format: "ICEBERG".to_string(), + format_version: 2, + warehouse_location: format!("s3://{bucket}/tables/table-id"), + metadata_location, + version_token: "token-v1".to_string(), + generation: 1, + state: TableCatalogEntryState::Active, + properties: BTreeMap::new(), + created_at: None, + updated_at: None, + } + } + + #[tokio::test] + async fn object_table_catalog_store_commits_with_token_match_and_writes_log() { + let backend = TestCatalogObjectBackend::default(); + let store = ObjectTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + let namespace = Namespace::parse("sales").unwrap(); + let table = IdentifierSegment::parse("orders").unwrap(); + let current_metadata = default_table_metadata_file_path(&namespace, &table, "00001.metadata.json"); + let new_metadata = default_table_metadata_file_path(&namespace, &table, "00002.metadata.json"); + + store.put_table_bucket(test_bucket_entry(bucket)).await.unwrap(); + store + .create_namespace(test_namespace_entry(bucket, &namespace)) + .await + .unwrap(); + store + .create_table(test_table_entry(bucket, &namespace, &table, current_metadata.clone())) + .await + .unwrap(); + backend.seed_object(bucket, &new_metadata, b"{}".to_vec()).await; + + let result = store + .commit_table(TableCommitRequest { + table_bucket: bucket.to_string(), + namespace: namespace.public_name(), + table: table.as_str().to_string(), + commit_id: "commit-1".to_string(), + idempotency_key: Some("client/%2f\nrequest".to_string()), + operation: "append".to_string(), + expected_version_token: "token-v1".to_string(), + expected_metadata_location: current_metadata, + new_metadata_location: new_metadata.clone(), + requirements: vec![serde_json::json!({"type": "assert-table-uuid", "uuid": "table-uuid"})], + writer: Some("pyiceberg/test".to_string()), + }) + .await + .unwrap(); + + assert_eq!(result.table.metadata_location, new_metadata); + assert_ne!(result.table.version_token, "token-v1"); + assert_eq!(result.table.generation, 2); + assert_eq!(result.commit_log.status, CommitLogStatus::Committed); + + let loaded = store.load_table(bucket, "sales", "orders").await.unwrap().unwrap(); + assert_eq!(loaded.metadata_location, result.table.metadata_location); + assert_eq!(loaded.version_token, result.table.version_token); + assert!( + store + .get_commit_by_id(bucket, "table-id", "commit-1") + .await + .unwrap() + .is_some() + ); + assert!( + store + .get_commit_by_idempotency_key(bucket, "table-id", "client/%2f\nrequest") + .await + .unwrap() + .is_some() + ); + } + + #[tokio::test] + async fn object_table_catalog_store_does_not_advance_table_when_idempotency_staging_fails() { + let backend = TestCatalogObjectBackend::default(); + let store = ObjectTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + let namespace = Namespace::parse("sales").unwrap(); + let table = IdentifierSegment::parse("orders").unwrap(); + let current_metadata = default_table_metadata_file_path(&namespace, &table, "00001.metadata.json"); + let new_metadata = default_table_metadata_file_path(&namespace, &table, "00002.metadata.json"); + let idempotency_key = "client-request"; + let idempotency_path = TableCatalogObjectPaths::default().commit_idempotency_entry_path("table-id", idempotency_key); + + store.put_table_bucket(test_bucket_entry(bucket)).await.unwrap(); + store + .create_namespace(test_namespace_entry(bucket, &namespace)) + .await + .unwrap(); + store + .create_table(test_table_entry(bucket, &namespace, &table, current_metadata.clone())) + .await + .unwrap(); + backend.seed_object(bucket, &new_metadata, b"{}".to_vec()).await; + backend.fail_put_attempt(bucket, &idempotency_path, 1).await; + + let err = store + .commit_table(TableCommitRequest { + table_bucket: bucket.to_string(), + namespace: namespace.public_name(), + table: table.as_str().to_string(), + commit_id: "commit-1".to_string(), + idempotency_key: Some(idempotency_key.to_string()), + operation: "append".to_string(), + expected_version_token: "token-v1".to_string(), + expected_metadata_location: current_metadata.clone(), + new_metadata_location: new_metadata, + requirements: Vec::new(), + writer: None, + }) + .await + .unwrap_err(); + + assert!(matches!(err, TableCatalogStoreError::Internal(_))); + let loaded = store.load_table(bucket, "sales", "orders").await.unwrap().unwrap(); + assert_eq!(loaded.metadata_location, current_metadata); + assert_eq!(loaded.version_token, "token-v1"); + let staged = store.get_commit_by_id(bucket, "table-id", "commit-1").await.unwrap().unwrap(); + assert_eq!(staged.status, CommitLogStatus::Staged); + } + + #[tokio::test] + async fn object_table_catalog_store_recovers_staged_commit_after_post_cas_finalization_failure() { + let backend = TestCatalogObjectBackend::default(); + let store = ObjectTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + let namespace = Namespace::parse("sales").unwrap(); + let table = IdentifierSegment::parse("orders").unwrap(); + let current_metadata = default_table_metadata_file_path(&namespace, &table, "00001.metadata.json"); + let new_metadata = default_table_metadata_file_path(&namespace, &table, "00002.metadata.json"); + let commit_path = TableCatalogObjectPaths::default().commit_log_entry_path("table-id", "commit-1"); + + store.put_table_bucket(test_bucket_entry(bucket)).await.unwrap(); + store + .create_namespace(test_namespace_entry(bucket, &namespace)) + .await + .unwrap(); + store + .create_table(test_table_entry(bucket, &namespace, &table, current_metadata.clone())) + .await + .unwrap(); + backend.seed_object(bucket, &new_metadata, b"{}".to_vec()).await; + backend.fail_put_attempt(bucket, &commit_path, 2).await; + + let request = TableCommitRequest { + table_bucket: bucket.to_string(), + namespace: namespace.public_name(), + table: table.as_str().to_string(), + commit_id: "commit-1".to_string(), + idempotency_key: None, + operation: "append".to_string(), + expected_version_token: "token-v1".to_string(), + expected_metadata_location: current_metadata, + new_metadata_location: new_metadata.clone(), + requirements: Vec::new(), + writer: None, + }; + + let result = store.commit_table(request.clone()).await.unwrap(); + + assert_eq!(result.table.metadata_location, new_metadata); + assert_eq!(result.commit_log.status, CommitLogStatus::Committed); + let loaded = store.load_table(bucket, "sales", "orders").await.unwrap().unwrap(); + assert_eq!(loaded.version_token, result.table.version_token); + let staged = store.get_commit_by_id(bucket, "table-id", "commit-1").await.unwrap().unwrap(); + assert_eq!(staged.status, CommitLogStatus::Staged); + + let retry = store.commit_table(request).await.unwrap(); + assert_eq!(retry.table.version_token, result.table.version_token); + assert_eq!(retry.commit_log.status, CommitLogStatus::Committed); + let committed = store.get_commit_by_id(bucket, "table-id", "commit-1").await.unwrap().unwrap(); + assert_eq!(committed.status, CommitLogStatus::Committed); + } + + #[tokio::test] + async fn object_table_catalog_store_rejects_stale_commit_token() { + let backend = TestCatalogObjectBackend::default(); + let store = ObjectTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + let namespace = Namespace::parse("sales").unwrap(); + let table = IdentifierSegment::parse("orders").unwrap(); + let current_metadata = default_table_metadata_file_path(&namespace, &table, "00001.metadata.json"); + let new_metadata = default_table_metadata_file_path(&namespace, &table, "00002.metadata.json"); + + store.put_table_bucket(test_bucket_entry(bucket)).await.unwrap(); + store + .create_namespace(test_namespace_entry(bucket, &namespace)) + .await + .unwrap(); + store + .create_table(test_table_entry(bucket, &namespace, &table, current_metadata.clone())) + .await + .unwrap(); + backend.seed_object(bucket, &new_metadata, b"{}".to_vec()).await; + + let err = store + .commit_table(TableCommitRequest { + table_bucket: bucket.to_string(), + namespace: namespace.public_name(), + table: table.as_str().to_string(), + commit_id: "commit-1".to_string(), + idempotency_key: None, + operation: "append".to_string(), + expected_version_token: "stale-token".to_string(), + expected_metadata_location: current_metadata.clone(), + new_metadata_location: new_metadata, + requirements: Vec::new(), + writer: None, + }) + .await + .unwrap_err(); + + assert!(matches!(err, TableCatalogStoreError::Conflict(_))); + let loaded = store.load_table(bucket, "sales", "orders").await.unwrap().unwrap(); + assert_eq!(loaded.metadata_location, current_metadata); + assert_eq!(loaded.version_token, "token-v1"); + assert!( + store + .get_commit_by_id(bucket, "table-id", "commit-1") + .await + .unwrap() + .is_none() + ); + } + #[test] fn namespace_marker_path_stays_under_default_reserved_boundary() { let namespace = Namespace::parse("analytics.daily_events").unwrap();