feat(table-catalog): add refs and views semantics (#3502)

Co-authored-by: Henry Guo <marshawcoco@users.noreply.github.com>
This commit is contained in:
Henry Guo
2026-06-16 21:19:27 +08:00
committed by GitHub
parent 966756d788
commit cd96491b39
6 changed files with 1624 additions and 47 deletions
File diff suppressed because it is too large Load Diff
+52 -1
View File
@@ -718,6 +718,12 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[
GET_TABLE_METADATA,
RouteRiskLevel::Sensitive,
),
admin(
HttpMethod::Head,
"/iceberg/v1/{warehouse}/namespaces/{namespace}/views/{view}",
GET_TABLE,
RouteRiskLevel::Sensitive,
),
admin(
HttpMethod::Post,
"/iceberg/v1/{warehouse}/namespaces/{namespace}/views/{view}",
@@ -736,6 +742,18 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[
GET_TABLE_METADATA,
RouteRiskLevel::Sensitive,
),
admin(
HttpMethod::Put,
"/iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}",
COMMIT_TABLE,
RouteRiskLevel::High,
),
admin(
HttpMethod::Delete,
"/iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}",
COMMIT_TABLE,
RouteRiskLevel::High,
),
admin(
HttpMethod::Get,
"/iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/metadata-location",
@@ -929,6 +947,12 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[
GET_TABLE_METADATA,
RouteRiskLevel::Sensitive,
),
admin(
HttpMethod::Head,
"/_iceberg/v1/{warehouse}/namespaces/{namespace}/views/{view}",
GET_TABLE,
RouteRiskLevel::Sensitive,
),
admin(
HttpMethod::Post,
"/_iceberg/v1/{warehouse}/namespaces/{namespace}/views/{view}",
@@ -947,6 +971,18 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[
GET_TABLE_METADATA,
RouteRiskLevel::Sensitive,
),
admin(
HttpMethod::Put,
"/_iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}",
COMMIT_TABLE,
RouteRiskLevel::High,
),
admin(
HttpMethod::Delete,
"/_iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}",
COMMIT_TABLE,
RouteRiskLevel::High,
),
admin(
HttpMethod::Get,
"/_iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/metadata-location",
@@ -1214,7 +1250,7 @@ mod tests {
let table_specs = ADMIN_ROUTE_POLICY_SPECS
.iter()
.filter(|spec| spec.path().starts_with("/iceberg/v1") || spec.path().starts_with("/_iceberg/v1"));
assert_eq!(table_specs.count(), 72);
assert_eq!(table_specs.count(), 78);
assert_action(HttpMethod::Put, "/iceberg/v1/buckets/{warehouse}", SET_TABLE_BUCKET);
assert_action(HttpMethod::Get, "/_iceberg/v1/buckets/{warehouse}", GET_TABLE_BUCKET);
assert_action(HttpMethod::Get, "/iceberg/v1/{warehouse}/namespaces", GET_TABLE_NAMESPACE);
@@ -1244,6 +1280,11 @@ mod tests {
"/iceberg/v1/{warehouse}/namespaces/{namespace}/views/{view}",
GET_TABLE_METADATA,
);
assert_action(
HttpMethod::Head,
"/_iceberg/v1/{warehouse}/namespaces/{namespace}/views/{view}",
GET_TABLE,
);
assert_action(
HttpMethod::Post,
"/_iceberg/v1/{warehouse}/namespaces/{namespace}/views/{view}",
@@ -1259,6 +1300,16 @@ mod tests {
"/iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/refs",
GET_TABLE_METADATA,
);
assert_action(
HttpMethod::Put,
"/iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}",
COMMIT_TABLE,
);
assert_action(
HttpMethod::Delete,
"/_iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}",
COMMIT_TABLE,
);
assert_action(
HttpMethod::Get,
"/_iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}",
@@ -352,6 +352,11 @@ fn expected_admin_route_matrix() -> Vec<RouteMatrixEntry> {
"/{warehouse}/namespaces/{namespace}/views/{view}",
"/analytics/namespaces/sales/views/recent_orders",
),
table_route_sample(
Method::HEAD,
"/{warehouse}/namespaces/{namespace}/views/{view}",
"/analytics/namespaces/sales/views/recent_orders",
),
table_route_sample(
Method::POST,
"/{warehouse}/namespaces/{namespace}/views/{view}",
@@ -367,6 +372,16 @@ fn expected_admin_route_matrix() -> Vec<RouteMatrixEntry> {
"/{warehouse}/namespaces/{namespace}/tables/{table}/refs",
"/analytics/namespaces/sales/tables/orders/refs",
),
table_route_sample(
Method::PUT,
"/{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}",
"/analytics/namespaces/sales/tables/orders/refs/audit",
),
table_route_sample(
Method::DELETE,
"/{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}",
"/analytics/namespaces/sales/tables/orders/refs/audit",
),
table_route_sample(
Method::POST,
"/{warehouse}/namespaces/{namespace}/tables/{table}/maintenance/metadata",
@@ -500,6 +515,11 @@ fn expected_admin_route_matrix() -> Vec<RouteMatrixEntry> {
"/{warehouse}/namespaces/{namespace}/views/{view}",
"/analytics/namespaces/sales/views/recent_orders",
),
compat_table_route_sample(
Method::HEAD,
"/{warehouse}/namespaces/{namespace}/views/{view}",
"/analytics/namespaces/sales/views/recent_orders",
),
compat_table_route_sample(
Method::POST,
"/{warehouse}/namespaces/{namespace}/views/{view}",
@@ -515,6 +535,16 @@ fn expected_admin_route_matrix() -> Vec<RouteMatrixEntry> {
"/{warehouse}/namespaces/{namespace}/tables/{table}/refs",
"/analytics/namespaces/sales/tables/orders/refs",
),
compat_table_route_sample(
Method::PUT,
"/{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}",
"/analytics/namespaces/sales/tables/orders/refs/audit",
),
compat_table_route_sample(
Method::DELETE,
"/{warehouse}/namespaces/{namespace}/tables/{table}/refs/{ref}",
"/analytics/namespaces/sales/tables/orders/refs/audit",
),
compat_table_route_sample(
Method::POST,
"/{warehouse}/namespaces/{namespace}/tables/{table}/maintenance/metadata",
@@ -719,6 +749,11 @@ fn test_register_routes_cover_representative_admin_paths() {
Method::GET,
&table_catalog_path("/analytics/namespaces/sales/views/recent_orders"),
);
assert_route(
&router,
Method::HEAD,
&table_catalog_path("/analytics/namespaces/sales/views/recent_orders"),
);
assert_route(
&router,
Method::POST,
@@ -734,6 +769,16 @@ fn test_register_routes_cover_representative_admin_paths() {
Method::GET,
&table_catalog_path("/analytics/namespaces/sales/tables/orders/refs"),
);
assert_route(
&router,
Method::PUT,
&table_catalog_path("/analytics/namespaces/sales/tables/orders/refs/audit"),
);
assert_route(
&router,
Method::DELETE,
&table_catalog_path("/analytics/namespaces/sales/tables/orders/refs/audit"),
);
assert_route(
&router,
Method::POST,
@@ -841,6 +886,11 @@ fn test_register_routes_cover_representative_admin_paths() {
Method::GET,
&compat_table_catalog_path("/analytics/namespaces/sales/views/recent_orders"),
);
assert_route(
&router,
Method::HEAD,
&compat_table_catalog_path("/analytics/namespaces/sales/views/recent_orders"),
);
assert_route(
&router,
Method::POST,
@@ -856,6 +906,16 @@ fn test_register_routes_cover_representative_admin_paths() {
Method::GET,
&compat_table_catalog_path("/analytics/namespaces/sales/tables/orders/refs"),
);
assert_route(
&router,
Method::PUT,
&compat_table_catalog_path("/analytics/namespaces/sales/tables/orders/refs/audit"),
);
assert_route(
&router,
Method::DELETE,
&compat_table_catalog_path("/analytics/namespaces/sales/tables/orders/refs/audit"),
);
assert_route(
&router,
Method::POST,
+380
View File
@@ -67,6 +67,7 @@ pub const TABLE_RESERVED_PREFIX: &str = BUCKET_TABLE_RESERVED_PREFIX;
const WAREHOUSE_ROOT: &str = "warehouses";
const NAMESPACE_ROOT: &str = "namespaces";
const TABLE_ROOT: &str = "tables";
const VIEW_ROOT: &str = "views";
const NAMESPACE_MARKER_FILE: &str = "namespace.json";
const TABLE_MARKER_FILE: &str = "table.json";
const CURRENT_POINTER_FILE: &str = "current.json";
@@ -77,6 +78,7 @@ const DELETE_DIR: &str = "delete";
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 VIEW_ENTRY_FILE: &str = "view-entry.json";
const INTERNAL_CATALOG_ROOT: &str = BUCKET_TABLE_CATALOG_META_PREFIX;
const TABLE_BUCKET_ROOT: &str = BUCKET_TABLE_CATALOG_TABLE_BUCKETS_PREFIX;
const COMMIT_LOG_ROOT: &str = "commits";
@@ -221,6 +223,28 @@ pub(crate) struct TableEntry {
pub updated_at: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct ViewEntry {
pub version: u16,
pub table_bucket: String,
pub namespace: String,
pub view: String,
pub view_id: String,
pub view_uuid: String,
pub format: String,
pub format_version: u16,
pub warehouse_location: String,
pub metadata_location: String,
pub version_token: String,
pub generation: u64,
pub state: TableCatalogEntryState,
#[serde(default)]
pub properties: BTreeMap<String, String>,
pub created_at: Option<String>,
pub updated_at: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub(crate) enum CommitLogStatus {
@@ -273,6 +297,23 @@ pub(crate) struct TableCommitResult {
pub commit_log: CommitLogEntry,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct ViewCommitRequest {
pub table_bucket: String,
pub namespace: String,
pub view: String,
pub expected_version_token: String,
pub expected_metadata_location: String,
pub new_metadata_location: String,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct ViewCommitResult {
pub view: ViewEntry,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct TableMaintenanceConfig {
@@ -984,6 +1025,16 @@ pub(crate) trait TableCatalogStore: Send + Sync {
async fn drop_table(&self, table_bucket: &str, namespace: &str, table: &str) -> TableCatalogStoreResult<()>;
async fn create_view(&self, entry: ViewEntry) -> TableCatalogStoreResult<()>;
async fn list_views(&self, table_bucket: &str, namespace: &str) -> TableCatalogStoreResult<Vec<ViewEntry>>;
async fn load_view(&self, table_bucket: &str, namespace: &str, view: &str) -> TableCatalogStoreResult<Option<ViewEntry>>;
async fn replace_view(&self, request: ViewCommitRequest) -> TableCatalogStoreResult<ViewCommitResult>;
async fn drop_view(&self, table_bucket: &str, namespace: &str, view: &str) -> TableCatalogStoreResult<()>;
async fn get_commit_by_id(
&self,
table_bucket: &str,
@@ -1104,6 +1155,19 @@ impl TableCatalogObjectPaths {
)
}
pub fn view_entries_prefix(&self, table_bucket: &str, namespace: &Namespace) -> String {
format!("{}{}/{}/", self.namespace_entries_prefix(table_bucket), namespace.storage_id(), VIEW_ROOT)
}
pub fn view_entry_path(&self, table_bucket: &str, namespace: &Namespace, view: &IdentifierSegment) -> String {
format!(
"{}{}/{}",
self.view_entries_prefix(table_bucket, namespace),
view.as_str(),
VIEW_ENTRY_FILE
)
}
pub fn table_maintenance_config_path(
&self,
table_bucket: &str,
@@ -1338,6 +1402,25 @@ where
Ok(Some((entry, etag)))
}
async fn read_view_with_etag_unlocked(
&self,
table_bucket: &str,
namespace: &Namespace,
view: &IdentifierSegment,
) -> TableCatalogStoreResult<Option<(ViewEntry, String)>> {
let view_path = self.paths.view_entry_path(table_bucket, namespace, view);
let Some((entry, etag)) = self
.read_entry_unlocked::<ViewEntry>(self.catalog_bucket(), &view_path)
.await?
else {
return Ok(None);
};
let Some(etag) = etag else {
return Err(TableCatalogStoreError::Internal(format!("catalog view entry has no etag: {view_path}")));
};
Ok(Some((entry, etag)))
}
async fn write_table_entry(
&self,
entry: TableEntry,
@@ -1358,6 +1441,22 @@ where
.await
}
async fn write_view_entry(&self, entry: ViewEntry, precondition: TableCatalogPutPrecondition) -> TableCatalogStoreResult<()> {
validate_catalog_entry_version("view", entry.version)?;
self.require_table_bucket(&entry.table_bucket).await?;
let namespace = parse_namespace_for_store(&entry.namespace)?;
let view = parse_table_for_store(&entry.view)?;
if self.get_namespace(&entry.table_bucket, &entry.namespace).await?.is_none() {
return Err(TableCatalogStoreError::NotFound(format!(
"namespace {}/{}",
entry.table_bucket, entry.namespace
)));
}
let view_path = self.paths.view_entry_path(&entry.table_bucket, &namespace, &view);
self.write_entry(self.catalog_bucket(), &view_path, &entry, precondition)
.await
}
async fn read_commit_by_path(&self, object: &str) -> TableCatalogStoreResult<Option<CommitLogEntry>> {
self.read_entry::<CommitLogEntry>(self.catalog_bucket(), object)
.await
@@ -2912,6 +3011,13 @@ where
namespace.public_name()
)));
}
if !self.list_views(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(self.catalog_bucket(), &self.paths.namespace_entry_path(table_bucket, &namespace))
.await
@@ -3234,6 +3340,115 @@ where
self.backend.delete_object(self.catalog_bucket(), &object).await
}
async fn create_view(&self, entry: ViewEntry) -> TableCatalogStoreResult<()> {
self.write_view_entry(entry, TableCatalogPutPrecondition::IfAbsent).await
}
async fn list_views(&self, table_bucket: &str, namespace: &str) -> TableCatalogStoreResult<Vec<ViewEntry>> {
let namespace = parse_namespace_for_store(namespace)?;
let mut entries = Vec::new();
for object in self
.backend
.list_objects(self.catalog_bucket(), &self.paths.view_entries_prefix(table_bucket, &namespace))
.await?
{
if !object.ends_with(VIEW_ENTRY_FILE) {
continue;
}
if let Some((entry, _)) = self.read_entry::<ViewEntry>(self.catalog_bucket(), &object).await? {
entries.push(entry);
}
}
entries.sort_by(|left, right| left.view.cmp(&right.view));
Ok(entries)
}
async fn load_view(&self, table_bucket: &str, namespace: &str, view: &str) -> TableCatalogStoreResult<Option<ViewEntry>> {
let namespace = parse_namespace_for_store(namespace)?;
let view = parse_table_for_store(view)?;
self.read_entry::<ViewEntry>(self.catalog_bucket(), &self.paths.view_entry_path(table_bucket, &namespace, &view))
.await
.map(|entry| entry.map(|(view, _)| view))
}
async fn replace_view(&self, request: ViewCommitRequest) -> TableCatalogStoreResult<ViewCommitResult> {
let namespace = parse_namespace_for_store(&request.namespace)?;
let view = parse_table_for_store(&request.view)?;
let view_path = self.paths.view_entry_path(&request.table_bucket, &namespace, &view);
let _guard = self.backend.acquire_write_lock(self.catalog_bucket(), &view_path).await?;
let Some((current, current_etag)) = self
.read_view_with_etag_unlocked(&request.table_bucket, &namespace, &view)
.await?
else {
return Err(TableCatalogStoreError::NotFound(format!(
"view {}/{}/{}",
request.table_bucket, request.namespace, request.view
)));
};
if current.version_token != request.expected_version_token {
return Err(TableCatalogStoreError::Conflict(
"current view version token does not match expected token".to_string(),
));
}
if current.metadata_location != request.expected_metadata_location {
return Err(TableCatalogStoreError::Conflict(
"current view metadata location does not match expected location".to_string(),
));
}
if !is_valid_view_metadata_location(&namespace, &view, &request.new_metadata_location) {
return Err(TableCatalogStoreError::Invalid(
"new metadata location must be inside the view metadata directory".to_string(),
));
}
let Some(new_metadata_object) = self
.backend
.read_object(&request.table_bucket, &request.new_metadata_location)
.await?
else {
return Err(TableCatalogStoreError::NotFound(format!(
"new view metadata object {}",
request.new_metadata_location
)));
};
let next_warehouse_location =
metadata_warehouse_location(&request.table_bucket, &request.new_metadata_location, &new_metadata_object)?;
let mut next = current;
next.metadata_location = request.new_metadata_location;
if let Some(warehouse_location) = next_warehouse_location {
next.warehouse_location = warehouse_location;
}
next.version_token = format!("token-{}", Uuid::new_v4());
next.generation = next.generation.saturating_add(1);
self.write_entry_unlocked(
self.catalog_bucket(),
&view_path,
&next,
TableCatalogPutPrecondition::IfMatch(current_etag),
)
.await?;
Ok(ViewCommitResult { view: next })
}
async fn drop_view(&self, table_bucket: &str, namespace: &str, view: &str) -> TableCatalogStoreResult<()> {
let namespace = parse_namespace_for_store(namespace)?;
let view = parse_table_for_store(view)?;
let object = self.paths.view_entry_path(table_bucket, &namespace, &view);
if self
.load_view(table_bucket, &namespace.public_name(), view.as_str())
.await?
.is_none()
{
return Err(TableCatalogStoreError::NotFound(format!(
"view {}/{}/{}",
table_bucket,
namespace.public_name(),
view.as_str()
)));
}
self.backend.delete_object(self.catalog_bucket(), &object).await
}
async fn get_commit_by_id(
&self,
table_bucket: &str,
@@ -6163,6 +6378,22 @@ pub(crate) fn default_table_delete_dir_path(namespace: &Namespace, table: &Ident
format!("{}{}/{}", default_table_root_prefix(namespace), table.as_str(), DELETE_DIR)
}
pub(crate) fn default_view_root_prefix(namespace: &Namespace) -> String {
format!("{}{}/{}/", default_namespace_root_prefix(), namespace.storage_id(), VIEW_ROOT)
}
pub(crate) fn default_view_metadata_dir_path(namespace: &Namespace, view: &IdentifierSegment) -> String {
format!("{}{}/{}", default_view_root_prefix(namespace), view.as_str(), METADATA_DIR)
}
pub(crate) fn default_view_metadata_file_path(
namespace: &Namespace,
view: &IdentifierSegment,
metadata_file_name: &str,
) -> String {
format!("{}/{}", default_view_metadata_dir_path(namespace, view), metadata_file_name)
}
pub(crate) fn default_table_metadata_file_path(
namespace: &Namespace,
table: &IdentifierSegment,
@@ -6229,6 +6460,17 @@ pub(crate) fn is_valid_table_metadata_location(
.is_some_and(is_valid_table_metadata_file_name)
}
pub(crate) fn is_valid_view_metadata_location(namespace: &Namespace, view: &IdentifierSegment, metadata_location: &str) -> bool {
if metadata_location.is_empty() {
return false;
}
let metadata_prefix = format!("{}/", default_view_metadata_dir_path(namespace, view));
metadata_location
.strip_prefix(&metadata_prefix)
.is_some_and(is_valid_table_metadata_file_name)
}
pub(crate) fn is_valid_table_metadata_file_name(metadata_file_name: &str) -> bool {
if metadata_file_name.is_empty()
|| metadata_file_name.len() > TABLE_METADATA_FILE_NAME_MAX_LEN
@@ -6554,6 +6796,50 @@ mod tests {
Ok(())
}
async fn create_view(&self, _entry: ViewEntry) -> TableCatalogStoreResult<()> {
Ok(())
}
async fn list_views(&self, _table_bucket: &str, _namespace: &str) -> TableCatalogStoreResult<Vec<ViewEntry>> {
Ok(Vec::new())
}
async fn load_view(
&self,
_table_bucket: &str,
_namespace: &str,
_view: &str,
) -> TableCatalogStoreResult<Option<ViewEntry>> {
Ok(None)
}
async fn replace_view(&self, request: ViewCommitRequest) -> TableCatalogStoreResult<ViewCommitResult> {
Ok(ViewCommitResult {
view: ViewEntry {
version: TABLE_CATALOG_ENTRY_VERSION,
table_bucket: request.table_bucket,
namespace: request.namespace,
view: request.view,
view_id: "view-id".to_string(),
view_uuid: "view-uuid".to_string(),
format: "ICEBERG_VIEW".to_string(),
format_version: 1,
warehouse_location: "s3://analytics/views/view-id".to_string(),
metadata_location: request.new_metadata_location,
version_token: "token-v2".to_string(),
generation: 2,
state: TableCatalogEntryState::Active,
properties: BTreeMap::new(),
created_at: None,
updated_at: None,
},
})
}
async fn drop_view(&self, _table_bucket: &str, _namespace: &str, _view: &str) -> TableCatalogStoreResult<()> {
Ok(())
}
async fn get_commit_by_id(
&self,
_table_bucket: &str,
@@ -6636,6 +6922,10 @@ mod tests {
paths.table_entry_path(bucket, &namespace, &table),
format!("{bucket_root}namespaces/analytics/daily_events/tables/events/table-entry.json")
);
assert_eq!(
paths.view_entry_path(bucket, &namespace, &table),
format!("{bucket_root}namespaces/analytics/daily_events/views/events/view-entry.json")
);
let commit_path = paths.commit_log_entry_path("table/../bucket", "table/../id", "commit/%2f\nid");
let idempotency_path = paths.commit_idempotency_entry_path("table/../bucket", "table/../id", "client/%2f\nrequest");
@@ -7045,6 +7335,27 @@ mod tests {
}
}
fn test_view_entry(bucket: &str, namespace: &Namespace, view: &IdentifierSegment, metadata_location: String) -> ViewEntry {
ViewEntry {
version: TABLE_CATALOG_ENTRY_VERSION,
table_bucket: bucket.to_string(),
namespace: namespace.public_name(),
view: view.as_str().to_string(),
view_id: "view-id".to_string(),
view_uuid: "view-uuid".to_string(),
format: "ICEBERG_VIEW".to_string(),
format_version: 1,
warehouse_location: format!("s3://{bucket}/views/view-id"),
metadata_location,
version_token: "token-v1".to_string(),
generation: 1,
state: TableCatalogEntryState::Active,
properties: BTreeMap::new(),
created_at: None,
updated_at: None,
}
}
async fn seed_table_for_metadata_maintenance(
store: &ObjectTableCatalogStore<TestCatalogObjectBackend>,
bucket: &str,
@@ -7078,6 +7389,75 @@ mod tests {
assert_eq!(object_buckets, BTreeSet::from([rustfs_ecstore::disk::RUSTFS_META_BUCKET]));
}
#[tokio::test]
async fn object_table_catalog_store_persists_view_entries_and_blocks_non_empty_namespace_drop() {
let backend = TestCatalogObjectBackend::default();
let store = ObjectTableCatalogStore::new(backend.clone());
let bucket = "analytics";
let namespace = Namespace::parse("sales").unwrap();
let view = IdentifierSegment::parse("recent_orders").unwrap();
let current_metadata = default_view_metadata_file_path(&namespace, &view, "00001.metadata.json");
let next_metadata = default_view_metadata_file_path(&namespace, &view, "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_view(test_view_entry(bucket, &namespace, &view, current_metadata.clone()))
.await
.unwrap();
assert_eq!(store.list_views(bucket, &namespace.public_name()).await.unwrap()[0].view, "recent_orders");
assert!(
store
.load_view(bucket, &namespace.public_name(), view.as_str())
.await
.unwrap()
.is_some()
);
assert!(matches!(
store.drop_namespace(bucket, &namespace.public_name()).await,
Err(TableCatalogStoreError::Conflict(_))
));
backend
.seed_object(
bucket,
&next_metadata,
serde_json::to_vec(&serde_json::json!({
"format-version": 1,
"view-uuid": "view-uuid",
"location": format!("s3://{bucket}/views/view-id")
}))
.unwrap(),
)
.await;
let result = store
.replace_view(ViewCommitRequest {
table_bucket: bucket.to_string(),
namespace: namespace.public_name(),
view: view.as_str().to_string(),
expected_version_token: "token-v1".to_string(),
expected_metadata_location: current_metadata,
new_metadata_location: next_metadata.clone(),
})
.await
.unwrap();
assert_eq!(result.view.metadata_location, next_metadata);
assert_eq!(result.view.generation, 2);
assert_ne!(result.view.version_token, "token-v1");
store
.drop_view(bucket, &namespace.public_name(), view.as_str())
.await
.unwrap();
assert!(store.list_views(bucket, &namespace.public_name()).await.unwrap().is_empty());
store.drop_namespace(bucket, &namespace.public_name()).await.unwrap();
}
#[tokio::test]
async fn maintenance_dry_run_keeps_current_metadata() {
let backend = TestCatalogObjectBackend::default();
+9 -1
View File
@@ -120,6 +120,15 @@ work items. They are intentionally conservative: only PyIceberg is automated by
this script today; other engines are documented until a repeatable harness is
added.
RustFS also exposes catalog-backed advanced Iceberg surfaces that are not part
of the PyIceberg append smoke path yet:
- table refs can be listed, created or replaced, and deleted through catalog
commits; refs with explicit retention policy require a forced delete, and
`main` cannot be deleted
- Iceberg views support basic create, list, load, replace, existence check, and
drop routes with persisted view metadata and view-scoped authorization
## Client Matrix
| Client | Current status | Claim |
@@ -154,7 +163,6 @@ current unsupported inventory is:
- snapshot expiration dry-run planning and manual catalog commit: supported through metadata maintenance reports
- automatic maintenance scheduling: external scheduler hook supported through the worker run endpoint; built-in periodic scheduling is not claimed
- compaction rewrite: controlled run-once support for unpartitioned Parquet binpack through metadata maintenance; built-in periodic scheduling, sort compaction, delete-file rewrite, and row-level compaction are not claimed
- Iceberg views: stable unsupported routes are registered and return explicit unsupported JSON
- external catalog bridges: metadata import/register is supported, but Polaris/Glue/DLF/Hive synchronization is unsupported
- multi-table transactions: not a short-term production claim
-6
View File
@@ -156,12 +156,6 @@ UNSUPPORTED_INVENTORY: list[dict[str, str]] = [
"roadmap_area": "snapshot-maintenance",
"expected_behavior": "metadata maintenance can plan binpack candidates and commit a safe unpartitioned Parquet rewrite through the catalog; built-in periodic scheduling, sort compaction, delete-file rewrite, and row-level compaction are not claimed",
},
{
"capability": "iceberg-views",
"status": "stable-unsupported-routes",
"roadmap_area": "view-api",
"expected_behavior": "view routes are registered and should return a stable unsupported JSON response until implemented",
},
{
"capability": "external-catalog-bridge",
"status": "metadata-import-only",