fix(table-catalog): prevent object catalog lock reentry (#6144)

* fix(table-catalog): prevent object catalog lock reentry

* test(table-catalog): distinguish unlocked object reads

---------

Co-authored-by: Henry Guo <marshawcoco@users.noreply.github.com>
This commit is contained in:
Henry Guo
2026-08-17 08:06:53 +08:00
committed by GitHub
parent 9f02ca6c36
commit 890ddea94b
5 changed files with 60 additions and 20 deletions
+35 -18
View File
@@ -387,6 +387,7 @@ where
async fn has_active_namespace_descendant(&self, table_bucket: &str, namespace: &Namespace) -> TableCatalogStoreResult<bool> {
let parent = namespace.public_name();
let namespace_path = self.paths.namespace_entry_path(table_bucket, namespace);
let descendant_prefix = format!("{}{}/", self.paths.namespace_entries_prefix(table_bucket), namespace.storage_id());
let scan_limit = NonZeroUsize::new(TABLE_CATALOG_LIST_MAX_KEYS)
.ok_or_else(|| TableCatalogStoreError::Internal("catalog object scan limit must be positive".to_string()))?;
@@ -398,6 +399,9 @@ where
.await?;
let last_scanned = page.objects.last().cloned();
for object in page.objects {
if object == namespace_path {
continue;
}
if self
.read_active_namespace_evidence(&object)
.await?
@@ -1814,7 +1818,6 @@ where
&self,
report: &TableMetadataMaintenanceReport,
) -> TableCatalogStoreResult<()> {
let report = table_maintenance_report_with_recommended_actions(report.clone());
let namespace = parse_namespace_for_store(&report.job.namespace)?;
let table = parse_table_for_store(&report.job.table)?;
let Some((entry, _)) = self
@@ -1828,20 +1831,27 @@ where
table.as_str()
)));
};
validate_table_maintenance_report_owner(&report, &report.job.table_bucket, &namespace, &table, &entry.table_id)?;
let job_path = self.paths.table_maintenance_job_path(
&report.job.table_bucket,
&namespace,
&table,
&report.job.table_id,
&report.job.job_id,
);
self.put_table_metadata_maintenance_report_for_entry(report, &entry).await
}
async fn put_table_metadata_maintenance_report_for_entry(
&self,
report: &TableMetadataMaintenanceReport,
entry: &TableEntry,
) -> TableCatalogStoreResult<()> {
let report = table_maintenance_report_with_recommended_actions(report.clone());
let namespace = parse_namespace_for_store(&entry.namespace)?;
let table = parse_table_for_store(&entry.table)?;
validate_table_maintenance_report_owner(&report, &entry.table_bucket, &namespace, &table, &entry.table_id)?;
let job_path =
self.paths
.table_maintenance_job_path(&entry.table_bucket, &namespace, &table, &entry.table_id, &report.job.job_id);
let latest_job_path =
self.paths
.table_maintenance_latest_job_path(&report.job.table_bucket, &namespace, &table, &report.job.table_id);
.table_maintenance_latest_job_path(&entry.table_bucket, &namespace, &table, &entry.table_id);
let current_job_path =
self.paths
.table_maintenance_current_job_path(&report.job.table_bucket, &namespace, &table, &report.job.table_id);
.table_maintenance_current_job_path(&entry.table_bucket, &namespace, &table, &entry.table_id);
self.write_entry(self.catalog_bucket(), &job_path, &report, TableCatalogPutPrecondition::Any)
.await?;
self.write_entry(self.catalog_bucket(), &latest_job_path, &report, TableCatalogPutPrecondition::Any)
@@ -2157,7 +2167,7 @@ where
before_status,
before_quarantined_object_count,
);
self.put_table_metadata_maintenance_report_unfenced(&report).await?;
self.put_table_metadata_maintenance_report_for_entry(&report, &entry).await?;
report
}
TableMaintenanceSchedulerPreflight::Complete(report) => *report,
@@ -2270,7 +2280,7 @@ where
before_status,
before_quarantined_object_count,
);
self.put_table_metadata_maintenance_report_unfenced(&report).await?;
self.put_table_metadata_maintenance_report_for_entry(&report, &entry).await?;
report
};
@@ -2381,6 +2391,7 @@ where
}
self.expire_table_maintenance_job(
current,
context.entry,
now,
"maintenance worker lease expired",
TableMaintenanceAuditAction::WorkerLeaseExpired,
@@ -2392,6 +2403,7 @@ where
}
self.expire_table_maintenance_job(
current,
context.entry,
now,
"maintenance scheduler lease expired",
TableMaintenanceAuditAction::SchedulerLeaseExpired,
@@ -2555,7 +2567,7 @@ where
before_status,
before_quarantined_object_count,
);
self.put_table_metadata_maintenance_report_unfenced(&report).await?;
self.put_table_metadata_maintenance_report_for_entry(&report, &entry).await?;
let delete = matches!(report.job.operation, TableMetadataMaintenanceOperation::Delete);
(report, effective, delete)
@@ -2621,6 +2633,7 @@ where
}
self.expire_table_maintenance_job(
current,
context.entry,
now,
"maintenance worker lease expired",
TableMaintenanceAuditAction::WorkerLeaseExpired,
@@ -2635,6 +2648,7 @@ where
}
self.expire_table_maintenance_job(
current,
context.entry,
now,
"maintenance scheduler lease expired",
TableMaintenanceAuditAction::SchedulerLeaseExpired,
@@ -2737,13 +2751,14 @@ where
Some(TableMetadataMaintenanceJobStatus::Running),
before_quarantined_object_count,
);
self.put_table_metadata_maintenance_report_unfenced(&report).await?;
self.put_table_metadata_maintenance_report_for_entry(&report, &entry).await?;
Ok(report)
}
async fn expire_table_maintenance_job(
&self,
mut report: TableMetadataMaintenanceReport,
entry: &TableEntry,
now: OffsetDateTime,
reason: &str,
action: TableMaintenanceAuditAction,
@@ -2763,7 +2778,7 @@ where
before_status,
before_quarantined_object_count,
);
self.put_table_metadata_maintenance_report_unfenced(&report).await?;
self.put_table_metadata_maintenance_report_for_entry(&report, entry).await?;
Ok(report)
}
@@ -2840,7 +2855,8 @@ where
None,
None,
);
self.put_table_metadata_maintenance_report_unfenced(&report).await?;
self.put_table_metadata_maintenance_report_for_entry(&report, control.entry)
.await?;
Ok(report)
}
@@ -2917,7 +2933,8 @@ where
None,
None,
);
self.put_table_metadata_maintenance_report_unfenced(&report).await?;
self.put_table_metadata_maintenance_report_for_entry(&report, control.entry)
.await?;
Ok(report)
}
+15
View File
@@ -356,6 +356,7 @@ pub(crate) struct TestCatalogObjectBackend {
pub(crate) missing_read_object_path: Arc<tokio::sync::Mutex<Option<String>>>,
pub(crate) fail_read_object_path: Arc<tokio::sync::Mutex<Option<String>>>,
pub(crate) lock_attempts: Arc<tokio::sync::Mutex<Vec<(String, String)>>>,
pub(crate) reject_reads_while_write_locked: bool,
/// Content-addressed (sha256) etags instead of the store fake's counter.
/// The admin handler tests observe an object's etag and expect rewriting
/// identical bytes to reproduce it, so their fixtures set this.
@@ -689,6 +690,14 @@ impl TableCatalogObjectBackend for TestCatalogObjectBackend {
drop(fail_read_object_path);
let key = (bucket.to_string(), object.to_string());
if self.reject_reads_while_write_locked {
let lock = self.locks.lock().await.get(&key).cloned();
if lock.as_ref().is_some_and(|lock| lock.try_read().is_err()) {
return Err(TableCatalogStoreError::Internal(format!(
"catalog read attempted while its write lock is held: {object}"
)));
}
}
let (attempt, pause_before) = {
let mut state = self.state.lock().await;
state.read_calls += 1;
@@ -737,6 +746,12 @@ impl TableCatalogObjectBackend for TestCatalogObjectBackend {
Ok(result)
}
async fn read_object_unlocked(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<Option<TableCatalogObject>> {
let mut backend = self.clone();
backend.reject_reads_while_write_locked = false;
backend.read_object(bucket, object).await
}
async fn read_object_limited(
&self,
bucket: &str,
+8 -2
View File
@@ -3630,7 +3630,10 @@ async fn configured_table_catalog_store_uses_durable_strong_snapshot() {
#[tokio::test]
async fn object_table_catalog_store_persists_view_entries_and_blocks_non_empty_namespace_drop() {
let backend = TestCatalogObjectBackend::default();
let backend = TestCatalogObjectBackend {
reject_reads_while_write_locked: true,
..Default::default()
};
let store = ObjectTableCatalogStore::new(backend.clone());
let bucket = "analytics";
let namespace = Namespace::parse("sales").unwrap();
@@ -6442,7 +6445,10 @@ async fn maintenance_scheduler_report_marks_disabled_default() {
#[tokio::test]
async fn maintenance_scheduler_run_queues_one_durable_job() {
let backend = TestCatalogObjectBackend::default();
let backend = TestCatalogObjectBackend {
reject_reads_while_write_locked: true,
..Default::default()
};
let store = ObjectTableCatalogStore::new(backend.clone());
let bucket = "analytics";
let namespace = Namespace::parse("sales").expect("namespace should parse");
+1
View File
@@ -1207,6 +1207,7 @@ def smoke_view_request(args: argparse.Namespace, view_name: str, version_id: int
"schema-id": 0,
"timestamp-ms": int(time.time() * 1000),
"summary": {"operation": "replace"},
"default-namespace": [args.namespace],
"representations": [
{
"type": "sql",
@@ -749,6 +749,7 @@ class PyIcebergSmokeConfigTest(unittest.TestCase):
self.assertEqual(request["name"], "orders_view")
self.assertEqual(request["schema"]["type"], "struct")
self.assertEqual(request["view-version"]["version-id"], 7)
self.assertEqual(request["view-version"]["default-namespace"], ["sales"])
self.assertEqual(request["view-version"]["representations"][0]["dialect"], "spark")
self.assertEqual(request["properties"]["rustfs.smoke.table"], "orders")