mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-27 16:48:58 +00:00
feat(table-catalog): harden credential and maintenance ops (#3631)
Co-authored-by: Henry Guo <marshawcoco@users.noreply.github.com>
This commit is contained in:
@@ -59,6 +59,7 @@ const CREDENTIAL_EXPIRATION_CONFIG_KEY: &str = "rustfs.credential-expiration-uni
|
||||
const CREDENTIAL_VENDING_UNSUPPORTED: &str = "unsupported";
|
||||
const CREDENTIAL_VENDING_SUPPORTED: &str = "supported";
|
||||
const CREDENTIAL_VENDING_UNSUPPORTED_REASON: &str = "temporary-credentials-not-implemented";
|
||||
const CREDENTIAL_VENDING_DISABLED_REASON: &str = "credential-vending-disabled";
|
||||
const CREDENTIAL_SCOPE_WAREHOUSE_PREFIX: &str = "warehouse-prefix";
|
||||
const CREDENTIAL_SCOPE_TABLE_PREFIX: &str = "table-prefix";
|
||||
const CREDENTIAL_MODE_CLIENT_PROVIDED: &str = "client-provided-s3-credentials-required";
|
||||
@@ -731,6 +732,7 @@ struct RestLoadTableResponse {
|
||||
|
||||
#[derive(Debug, Serialize)]
|
||||
struct RestLoadCredentialsResponse {
|
||||
config: BTreeMap<String, String>,
|
||||
#[serde(rename = "storage-credentials")]
|
||||
storage_credentials: Vec<RestStorageCredential>,
|
||||
}
|
||||
@@ -1568,6 +1570,21 @@ fn load_view_response_from_entry(entry: crate::table_catalog::ViewEntry, metadat
|
||||
}
|
||||
}
|
||||
|
||||
fn load_credentials_response_config(vending: &str, mode: &str, reason: Option<&str>) -> BTreeMap<String, String> {
|
||||
let mut config = BTreeMap::new();
|
||||
config.insert(CREDENTIAL_VENDING_CONFIG_KEY.to_string(), vending.to_string());
|
||||
config.insert(CREDENTIAL_MODE_CONFIG_KEY.to_string(), mode.to_string());
|
||||
if let Some(reason) = reason {
|
||||
config.insert(CREDENTIAL_VENDING_REASON_CONFIG_KEY.to_string(), reason.to_string());
|
||||
}
|
||||
config
|
||||
}
|
||||
|
||||
fn add_table_credential_scope_config(config: &mut BTreeMap<String, String>, scope_prefix: &str) {
|
||||
config.insert(CREDENTIAL_SCOPE_CONFIG_KEY.to_string(), CREDENTIAL_SCOPE_TABLE_PREFIX.to_string());
|
||||
config.insert(CREDENTIAL_SCOPE_PREFIX_CONFIG_KEY.to_string(), scope_prefix.to_string());
|
||||
}
|
||||
|
||||
async fn load_credentials_response_from_entry(
|
||||
entry: &crate::table_catalog::TableEntry,
|
||||
issuer: &dyn TableCredentialIssuer,
|
||||
@@ -1575,6 +1592,11 @@ async fn load_credentials_response_from_entry(
|
||||
) -> S3Result<RestLoadCredentialsResponse> {
|
||||
if !issuer.enabled() {
|
||||
return Ok(RestLoadCredentialsResponse {
|
||||
config: load_credentials_response_config(
|
||||
CREDENTIAL_VENDING_UNSUPPORTED,
|
||||
CREDENTIAL_MODE_CLIENT_PROVIDED,
|
||||
Some(CREDENTIAL_VENDING_DISABLED_REASON),
|
||||
),
|
||||
storage_credentials: Vec::new(),
|
||||
});
|
||||
}
|
||||
@@ -1585,11 +1607,28 @@ async fn load_credentials_response_from_entry(
|
||||
scope_prefix: scope.scope_prefix.clone(),
|
||||
object_prefix: scope.object_prefix.clone(),
|
||||
};
|
||||
let scope_prefix = scope.scope_prefix.clone();
|
||||
let storage_credentials = match issuer.issue_table_credentials(request).await? {
|
||||
Some(issued) => vec![storage_credential_from_issued(scope, issued)],
|
||||
None => Vec::new(),
|
||||
None => {
|
||||
let mut config = load_credentials_response_config(
|
||||
CREDENTIAL_VENDING_UNSUPPORTED,
|
||||
CREDENTIAL_MODE_CLIENT_PROVIDED,
|
||||
Some(CREDENTIAL_VENDING_UNSUPPORTED_REASON),
|
||||
);
|
||||
add_table_credential_scope_config(&mut config, &scope_prefix);
|
||||
return Ok(RestLoadCredentialsResponse {
|
||||
config,
|
||||
storage_credentials: Vec::new(),
|
||||
});
|
||||
}
|
||||
};
|
||||
Ok(RestLoadCredentialsResponse { storage_credentials })
|
||||
let mut config = load_credentials_response_config(CREDENTIAL_VENDING_SUPPORTED, CREDENTIAL_MODE_CATALOG_VENDED, None);
|
||||
add_table_credential_scope_config(&mut config, &scope_prefix);
|
||||
Ok(RestLoadCredentialsResponse {
|
||||
config,
|
||||
storage_credentials,
|
||||
})
|
||||
}
|
||||
|
||||
fn commit_table_response_from_result(
|
||||
@@ -8540,6 +8579,18 @@ mod tests {
|
||||
.expect("disabled issuer should build an empty response");
|
||||
|
||||
assert!(response.storage_credentials.is_empty());
|
||||
assert_eq!(
|
||||
response.config.get(CREDENTIAL_VENDING_CONFIG_KEY),
|
||||
Some(&CREDENTIAL_VENDING_UNSUPPORTED.to_string())
|
||||
);
|
||||
assert_eq!(
|
||||
response.config.get(CREDENTIAL_MODE_CONFIG_KEY),
|
||||
Some(&CREDENTIAL_MODE_CLIENT_PROVIDED.to_string())
|
||||
);
|
||||
assert_eq!(
|
||||
response.config.get(CREDENTIAL_VENDING_REASON_CONFIG_KEY),
|
||||
Some(&"credential-vending-disabled".to_string())
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
@@ -8553,6 +8604,54 @@ mod tests {
|
||||
.expect("disabled issuer should not validate credential scopes");
|
||||
|
||||
assert!(response.storage_credentials.is_empty());
|
||||
assert_eq!(
|
||||
response.config.get(CREDENTIAL_VENDING_CONFIG_KEY),
|
||||
Some(&CREDENTIAL_VENDING_UNSUPPORTED.to_string())
|
||||
);
|
||||
assert!(!response.config.contains_key(CREDENTIAL_SCOPE_PREFIX_CONFIG_KEY));
|
||||
}
|
||||
|
||||
struct UnavailableTableCredentialIssuer;
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl TableCredentialIssuer for UnavailableTableCredentialIssuer {
|
||||
async fn issue_table_credentials(
|
||||
&self,
|
||||
request: TableCredentialIssueRequest<'_>,
|
||||
) -> S3Result<Option<IssuedTableCredentials>> {
|
||||
assert_eq!(request.scope_prefix, "s3://warehouse/tables/table-id/");
|
||||
Ok(None)
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn unavailable_table_credential_issuer_reports_fallback_scope() {
|
||||
let issuer = UnavailableTableCredentialIssuer;
|
||||
let response = load_credentials_response_from_entry(&table_entry_for_credentials(), &issuer, None)
|
||||
.await
|
||||
.expect("unavailable issuer should build a fallback response");
|
||||
|
||||
assert!(response.storage_credentials.is_empty());
|
||||
assert_eq!(
|
||||
response.config.get(CREDENTIAL_VENDING_CONFIG_KEY),
|
||||
Some(&CREDENTIAL_VENDING_UNSUPPORTED.to_string())
|
||||
);
|
||||
assert_eq!(
|
||||
response.config.get(CREDENTIAL_MODE_CONFIG_KEY),
|
||||
Some(&CREDENTIAL_MODE_CLIENT_PROVIDED.to_string())
|
||||
);
|
||||
assert_eq!(
|
||||
response.config.get(CREDENTIAL_VENDING_REASON_CONFIG_KEY),
|
||||
Some(&CREDENTIAL_VENDING_UNSUPPORTED_REASON.to_string())
|
||||
);
|
||||
assert_eq!(
|
||||
response.config.get(CREDENTIAL_SCOPE_CONFIG_KEY),
|
||||
Some(&CREDENTIAL_SCOPE_TABLE_PREFIX.to_string())
|
||||
);
|
||||
assert_eq!(
|
||||
response.config.get(CREDENTIAL_SCOPE_PREFIX_CONFIG_KEY),
|
||||
Some(&"s3://warehouse/tables/table-id/".to_string())
|
||||
);
|
||||
}
|
||||
|
||||
struct TestTableCredentialIssuer;
|
||||
@@ -8588,6 +8687,17 @@ mod tests {
|
||||
.await
|
||||
.expect("issuer should build a scoped credential response");
|
||||
|
||||
assert_eq!(
|
||||
response.config.get(CREDENTIAL_VENDING_CONFIG_KEY),
|
||||
Some(&CREDENTIAL_VENDING_SUPPORTED.to_string())
|
||||
);
|
||||
assert_eq!(
|
||||
response.config.get(CREDENTIAL_MODE_CONFIG_KEY),
|
||||
Some(&CREDENTIAL_MODE_CATALOG_VENDED.to_string())
|
||||
);
|
||||
assert!(!response.config.contains_key(S3_ACCESS_KEY_ID_CONFIG_KEY));
|
||||
assert!(!response.config.contains_key(S3_SECRET_ACCESS_KEY_CONFIG_KEY));
|
||||
assert!(!response.config.contains_key(S3_SESSION_TOKEN_CONFIG_KEY));
|
||||
assert_eq!(response.storage_credentials.len(), 1);
|
||||
let credential = &response.storage_credentials[0];
|
||||
assert_eq!(credential.prefix, "s3://warehouse/tables/table-id/");
|
||||
@@ -8609,6 +8719,28 @@ mod tests {
|
||||
assert!(!credential.config.contains_key("rustfs.credential-vending-reason"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn credential_response_serializes_sensitive_config_only_inside_storage_credentials() {
|
||||
let issuer = TestTableCredentialIssuer;
|
||||
let response = load_credentials_response_from_entry(&table_entry_for_credentials(), &issuer, None)
|
||||
.await
|
||||
.expect("issuer should build a scoped credential response");
|
||||
|
||||
let value = serde_json::to_value(&response).expect("credential response should serialize");
|
||||
|
||||
assert_eq!(
|
||||
value["config"][CREDENTIAL_VENDING_CONFIG_KEY],
|
||||
serde_json::Value::String(CREDENTIAL_VENDING_SUPPORTED.to_string())
|
||||
);
|
||||
assert!(value["config"].get(S3_ACCESS_KEY_ID_CONFIG_KEY).is_none());
|
||||
assert!(value["config"].get(S3_SECRET_ACCESS_KEY_CONFIG_KEY).is_none());
|
||||
assert!(value["config"].get(S3_SESSION_TOKEN_CONFIG_KEY).is_none());
|
||||
assert_eq!(
|
||||
value["storage-credentials"][0]["config"][S3_ACCESS_KEY_ID_CONFIG_KEY],
|
||||
serde_json::Value::String("temporary-access-key".to_string())
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn table_credentials_do_not_snapshot_parent_groups() {
|
||||
let principal = rustfs_credentials::Credentials {
|
||||
|
||||
+168
-9
@@ -108,6 +108,7 @@ const TABLE_METADATA_CLEANUP_SAFETY_WINDOW_SECONDS: i64 = 15 * 60;
|
||||
const TABLE_MAINTENANCE_RETRY_BACKOFF_MAX_SECONDS: u64 = 24 * 60 * 60;
|
||||
const TABLE_MAINTENANCE_WORKER_LEASE_TIMEOUT_DEFAULT_SECONDS: u64 = 15 * 60;
|
||||
const TABLE_MAINTENANCE_WORKER_LEASE_TIMEOUT_MAX_SECONDS: u64 = 24 * 60 * 60;
|
||||
const TABLE_MAINTENANCE_DELETE_DISABLED_REASON: &str = "metadata delete is disabled by maintenance config";
|
||||
const TABLE_COMMIT_SLOW_LOG_THRESHOLD: StdDuration = StdDuration::from_secs(2);
|
||||
const ICEBERG_MAIN_REF: &str = "main";
|
||||
const ICEBERG_MIN_SNAPSHOTS_TO_KEEP_PROPERTY: &str = "history.expire.min-snapshots-to-keep";
|
||||
@@ -490,6 +491,8 @@ pub(crate) struct TableMetadataMaintenanceJob {
|
||||
pub status: TableMetadataMaintenanceJobStatus,
|
||||
#[serde(default)]
|
||||
pub failure_reason: Option<String>,
|
||||
#[serde(default, rename = "recommended-actions")]
|
||||
pub recommended_actions: Vec<TableMaintenanceRecommendedAction>,
|
||||
#[serde(default)]
|
||||
pub config_source: TableMaintenanceConfigSource,
|
||||
#[serde(default)]
|
||||
@@ -755,6 +758,19 @@ pub(crate) enum TableMetadataMaintenanceJobStatus {
|
||||
Paused,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
|
||||
pub(crate) enum TableMaintenanceRecommendedAction {
|
||||
NoActionRequired,
|
||||
ReviewAndRunDelete,
|
||||
EnableDelete,
|
||||
EnableBackgroundMaintenance,
|
||||
ResumeMaintenanceWorker,
|
||||
WaitForRetryBackoff,
|
||||
WaitForActiveWorker,
|
||||
InvestigateFailure,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
|
||||
pub(crate) enum TableMetadataMaintenanceObjectState {
|
||||
@@ -2153,6 +2169,7 @@ 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 job_path = self.paths.table_maintenance_job_path(
|
||||
@@ -2168,11 +2185,11 @@ where
|
||||
let current_job_path =
|
||||
self.paths
|
||||
.table_maintenance_current_job_path(&report.job.table_bucket, &namespace, &table, &report.job.table_id);
|
||||
self.write_entry(self.catalog_bucket(), &job_path, report, TableCatalogPutPrecondition::Any)
|
||||
self.write_entry(self.catalog_bucket(), &job_path, &report, TableCatalogPutPrecondition::Any)
|
||||
.await?;
|
||||
self.write_entry(self.catalog_bucket(), &latest_job_path, report, TableCatalogPutPrecondition::Any)
|
||||
self.write_entry(self.catalog_bucket(), &latest_job_path, &report, TableCatalogPutPrecondition::Any)
|
||||
.await?;
|
||||
self.write_entry(self.catalog_bucket(), ¤t_job_path, report, TableCatalogPutPrecondition::Any)
|
||||
self.write_entry(self.catalog_bucket(), ¤t_job_path, &report, TableCatalogPutPrecondition::Any)
|
||||
.await
|
||||
}
|
||||
|
||||
@@ -2209,7 +2226,7 @@ where
|
||||
};
|
||||
self.read_entry::<TableMetadataMaintenanceReport>(self.catalog_bucket(), &job_path)
|
||||
.await
|
||||
.map(|entry| entry.map(|(report, _)| report))
|
||||
.map(|entry| entry.map(|(report, _)| table_maintenance_report_with_recommended_actions(report)))
|
||||
}
|
||||
|
||||
pub(crate) async fn run_table_metadata_maintenance_worker_once(
|
||||
@@ -2388,6 +2405,7 @@ where
|
||||
operation: TableMetadataMaintenanceOperation::DryRun,
|
||||
status: control.status,
|
||||
failure_reason: Some(control.reason.to_string()),
|
||||
recommended_actions: Vec::new(),
|
||||
config_source: control.effective.source,
|
||||
worker_id: Some(control.worker_id),
|
||||
lease_id: String::new(),
|
||||
@@ -2428,6 +2446,7 @@ where
|
||||
snapshot_expiration: None,
|
||||
compaction: None,
|
||||
};
|
||||
let report = table_maintenance_report_with_recommended_actions(report);
|
||||
self.put_table_metadata_maintenance_report(&report).await?;
|
||||
Ok(report)
|
||||
}
|
||||
@@ -2994,7 +3013,7 @@ where
|
||||
)
|
||||
.await?;
|
||||
|
||||
Ok(TableMetadataMaintenanceReport {
|
||||
Ok(table_maintenance_report_with_recommended_actions(TableMetadataMaintenanceReport {
|
||||
job: TableMetadataMaintenanceJob {
|
||||
job_id: Uuid::new_v4().to_string(),
|
||||
table_bucket: table_bucket.to_string(),
|
||||
@@ -3004,6 +3023,7 @@ where
|
||||
operation: TableMetadataMaintenanceOperation::DryRun,
|
||||
status: TableMetadataMaintenanceJobStatus::Successful,
|
||||
failure_reason: None,
|
||||
recommended_actions: Vec::new(),
|
||||
config_source: TableMaintenanceConfigSource::Default,
|
||||
worker_id: None,
|
||||
lease_id: String::new(),
|
||||
@@ -3044,7 +3064,7 @@ where
|
||||
reachability_graph,
|
||||
snapshot_expiration: None,
|
||||
compaction: None,
|
||||
})
|
||||
}))
|
||||
}
|
||||
|
||||
pub(crate) async fn delete_table_metadata_maintenance_candidates(
|
||||
@@ -3125,14 +3145,16 @@ where
|
||||
report.job.heartbeat_at = Some(started_at.clone());
|
||||
report.job.started_at = Some(started_at);
|
||||
report.job.finished_at = None;
|
||||
refresh_table_maintenance_report_recommended_actions(&mut report);
|
||||
self.put_table_metadata_maintenance_report(&report).await?;
|
||||
|
||||
if delete && !effective.config.delete_enabled {
|
||||
let mut failed = report;
|
||||
failed.job.status = TableMetadataMaintenanceJobStatus::Failed;
|
||||
failed.job.failure_reason = Some("metadata delete is disabled by maintenance config".to_string());
|
||||
failed.job.failure_reason = Some(TABLE_MAINTENANCE_DELETE_DISABLED_REASON.to_string());
|
||||
apply_maintenance_retry_after(&mut failed.job, &effective.config, OffsetDateTime::now_utc());
|
||||
failed.job.finished_at = Some(maintenance_timestamp(OffsetDateTime::now_utc()));
|
||||
refresh_table_maintenance_report_recommended_actions(&mut failed);
|
||||
self.put_table_metadata_maintenance_report(&failed).await?;
|
||||
return Ok(failed);
|
||||
}
|
||||
@@ -3150,17 +3172,20 @@ where
|
||||
failed.job.failure_reason = Some(err.to_string());
|
||||
apply_maintenance_retry_after(&mut failed.job, &effective.config, OffsetDateTime::now_utc());
|
||||
failed.job.finished_at = Some(maintenance_timestamp(OffsetDateTime::now_utc()));
|
||||
refresh_table_maintenance_report_recommended_actions(&mut failed);
|
||||
self.put_table_metadata_maintenance_report(&failed).await?;
|
||||
return Err(err);
|
||||
}
|
||||
};
|
||||
deleted.job.finished_at = Some(maintenance_timestamp(OffsetDateTime::now_utc()));
|
||||
refresh_table_maintenance_report_recommended_actions(&mut deleted);
|
||||
self.put_table_metadata_maintenance_report(&deleted).await?;
|
||||
return Ok(deleted);
|
||||
}
|
||||
|
||||
report.job.status = TableMetadataMaintenanceJobStatus::Successful;
|
||||
report.job.finished_at = Some(maintenance_timestamp(OffsetDateTime::now_utc()));
|
||||
refresh_table_maintenance_report_recommended_actions(&mut report);
|
||||
self.put_table_metadata_maintenance_report(&report).await?;
|
||||
Ok(report)
|
||||
}
|
||||
@@ -3337,7 +3362,7 @@ where
|
||||
let mut object_cleanup_reports = report.object_cleanup_reports;
|
||||
mark_deleted_object_cleanup_reports(&mut object_cleanup_reports, &deleted_object_locations);
|
||||
|
||||
Ok(TableMetadataMaintenanceReport {
|
||||
Ok(table_maintenance_report_with_recommended_actions(TableMetadataMaintenanceReport {
|
||||
job,
|
||||
current_metadata_location: entry.metadata_location,
|
||||
retained_metadata_locations,
|
||||
@@ -3351,7 +3376,7 @@ where
|
||||
reachability_graph: report.reachability_graph,
|
||||
snapshot_expiration: report.snapshot_expiration,
|
||||
compaction: report.compaction,
|
||||
})
|
||||
}))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -6123,6 +6148,58 @@ fn parse_maintenance_timestamp(timestamp: &str) -> Option<OffsetDateTime> {
|
||||
OffsetDateTime::parse(timestamp, &time::format_description::well_known::Rfc3339).ok()
|
||||
}
|
||||
|
||||
fn table_maintenance_recommended_actions(job: &TableMetadataMaintenanceJob) -> Vec<TableMaintenanceRecommendedAction> {
|
||||
let mut actions = Vec::new();
|
||||
match job.status {
|
||||
TableMetadataMaintenanceJobStatus::NotYetRun => {}
|
||||
TableMetadataMaintenanceJobStatus::Running => {
|
||||
actions.push(TableMaintenanceRecommendedAction::WaitForActiveWorker);
|
||||
}
|
||||
TableMetadataMaintenanceJobStatus::Successful => {
|
||||
if matches!(job.operation, TableMetadataMaintenanceOperation::DryRun)
|
||||
&& (job.deletable_metadata_file_count > 0 || job.deletable_object_count > 0)
|
||||
{
|
||||
actions.push(TableMaintenanceRecommendedAction::ReviewAndRunDelete);
|
||||
} else {
|
||||
actions.push(TableMaintenanceRecommendedAction::NoActionRequired);
|
||||
}
|
||||
}
|
||||
TableMetadataMaintenanceJobStatus::Failed => {
|
||||
if job
|
||||
.failure_reason
|
||||
.as_deref()
|
||||
.is_some_and(|reason| reason == TABLE_MAINTENANCE_DELETE_DISABLED_REASON)
|
||||
{
|
||||
actions.push(TableMaintenanceRecommendedAction::EnableDelete);
|
||||
}
|
||||
if job.next_retry_after.is_some() {
|
||||
actions.push(TableMaintenanceRecommendedAction::WaitForRetryBackoff);
|
||||
}
|
||||
if actions.is_empty() {
|
||||
actions.push(TableMaintenanceRecommendedAction::InvestigateFailure);
|
||||
}
|
||||
}
|
||||
TableMetadataMaintenanceJobStatus::Disabled => {
|
||||
actions.push(TableMaintenanceRecommendedAction::EnableBackgroundMaintenance);
|
||||
}
|
||||
TableMetadataMaintenanceJobStatus::Paused => {
|
||||
actions.push(TableMaintenanceRecommendedAction::ResumeMaintenanceWorker);
|
||||
}
|
||||
}
|
||||
actions
|
||||
}
|
||||
|
||||
fn refresh_table_maintenance_report_recommended_actions(report: &mut TableMetadataMaintenanceReport) {
|
||||
report.job.recommended_actions = table_maintenance_recommended_actions(&report.job);
|
||||
}
|
||||
|
||||
fn table_maintenance_report_with_recommended_actions(
|
||||
mut report: TableMetadataMaintenanceReport,
|
||||
) -> TableMetadataMaintenanceReport {
|
||||
refresh_table_maintenance_report_recommended_actions(&mut report);
|
||||
report
|
||||
}
|
||||
|
||||
fn table_maintenance_job_lease_is_active(
|
||||
job: &TableMetadataMaintenanceJob,
|
||||
worker_lease_timeout_seconds: u64,
|
||||
@@ -8107,6 +8184,10 @@ mod tests {
|
||||
assert_eq!(report.job.retained_metadata_file_count, 3);
|
||||
assert_eq!(report.job.cleanup_candidate_count, 2);
|
||||
assert_eq!(report.job.deletable_metadata_file_count, 1);
|
||||
assert_eq!(
|
||||
report.job.recommended_actions,
|
||||
vec![TableMaintenanceRecommendedAction::ReviewAndRunDelete]
|
||||
);
|
||||
|
||||
let current_report = maintenance_object_report(&report, ¤t);
|
||||
assert_eq!(current_report.state, TableMetadataMaintenanceObjectState::Retained);
|
||||
@@ -8141,6 +8222,54 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn maintenance_report_read_back_derives_actions_for_legacy_records() {
|
||||
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 = default_table_metadata_file_path(&namespace, &table, "00001.metadata.json");
|
||||
|
||||
seed_table_for_metadata_maintenance(&store, bucket, &namespace, &table, current.clone()).await;
|
||||
backend
|
||||
.seed_object(bucket, ¤t, br#"{"metadata-log":[]}"#.to_vec())
|
||||
.await;
|
||||
let mut report = store
|
||||
.plan_table_metadata_maintenance(bucket, "sales", "orders", 0)
|
||||
.await
|
||||
.expect("maintenance report should be planned");
|
||||
report.job.status = TableMetadataMaintenanceJobStatus::Running;
|
||||
report.job.worker_id = Some("worker-a".to_string());
|
||||
report.job.lease_id = "lease-a".to_string();
|
||||
report.job.heartbeat_at = Some(maintenance_timestamp(OffsetDateTime::UNIX_EPOCH + Duration::seconds(10)));
|
||||
|
||||
let job_path = store
|
||||
.paths
|
||||
.table_maintenance_job_path(bucket, &namespace, &table, "table-id", &report.job.job_id);
|
||||
let mut legacy_report = serde_json::to_value(&report).expect("legacy report should serialize");
|
||||
legacy_report
|
||||
.get_mut("job")
|
||||
.and_then(serde_json::Value::as_object_mut)
|
||||
.expect("legacy report job should be an object")
|
||||
.remove("recommended-actions");
|
||||
store
|
||||
.write_entry(store.catalog_bucket(), &job_path, &legacy_report, TableCatalogPutPrecondition::Any)
|
||||
.await
|
||||
.expect("legacy maintenance report should be seeded");
|
||||
|
||||
let loaded = store
|
||||
.get_table_metadata_maintenance_report(bucket, "sales", "orders", &report.job.job_id)
|
||||
.await
|
||||
.expect("legacy maintenance report lookup should succeed")
|
||||
.expect("legacy maintenance report should be returned");
|
||||
|
||||
assert_eq!(
|
||||
loaded.job.recommended_actions,
|
||||
vec![TableMaintenanceRecommendedAction::WaitForActiveWorker]
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn maintenance_state_is_scoped_to_current_table_identity() {
|
||||
let backend = TestCatalogObjectBackend::default();
|
||||
@@ -8598,6 +8727,13 @@ mod tests {
|
||||
.as_deref()
|
||||
.is_some_and(|reason| reason.contains("disabled"))
|
||||
);
|
||||
assert_eq!(
|
||||
report.job.recommended_actions,
|
||||
vec![
|
||||
TableMaintenanceRecommendedAction::EnableDelete,
|
||||
TableMaintenanceRecommendedAction::WaitForRetryBackoff,
|
||||
]
|
||||
);
|
||||
assert!(backend.object_exists(bucket, &old).await.unwrap());
|
||||
|
||||
let latest = store
|
||||
@@ -8607,6 +8743,7 @@ mod tests {
|
||||
.expect("failed maintenance job should be stored");
|
||||
assert_eq!(latest.job.job_id, report.job.job_id);
|
||||
assert_eq!(latest.job.status, TableMetadataMaintenanceJobStatus::Failed);
|
||||
assert_eq!(latest.job.recommended_actions, report.job.recommended_actions);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
@@ -8632,6 +8769,10 @@ mod tests {
|
||||
|
||||
assert_eq!(report.job.status, TableMetadataMaintenanceJobStatus::Disabled);
|
||||
assert_eq!(report.job.worker_id.as_deref(), Some("worker-a"));
|
||||
assert_eq!(
|
||||
report.job.recommended_actions,
|
||||
vec![TableMaintenanceRecommendedAction::EnableBackgroundMaintenance]
|
||||
);
|
||||
assert_eq!(report.job.deleted_metadata_file_count, 0);
|
||||
assert!(backend.object_exists(bucket, &old).await.unwrap());
|
||||
}
|
||||
@@ -8675,6 +8816,10 @@ mod tests {
|
||||
|
||||
assert_eq!(report.job.status, TableMetadataMaintenanceJobStatus::Paused);
|
||||
assert_eq!(report.job.operation, TableMetadataMaintenanceOperation::DryRun);
|
||||
assert_eq!(
|
||||
report.job.recommended_actions,
|
||||
vec![TableMaintenanceRecommendedAction::ResumeMaintenanceWorker]
|
||||
);
|
||||
assert_eq!(report.job.deleted_metadata_file_count, 0);
|
||||
assert!(backend.object_exists(bucket, &old).await.unwrap());
|
||||
}
|
||||
@@ -8731,6 +8876,12 @@ mod tests {
|
||||
assert_eq!(deferred.job.job_id, failed.job.job_id);
|
||||
assert_eq!(deferred.job.status, TableMetadataMaintenanceJobStatus::Failed);
|
||||
assert_eq!(deferred.job.worker_id.as_deref(), Some("worker-a"));
|
||||
assert!(
|
||||
deferred
|
||||
.job
|
||||
.recommended_actions
|
||||
.contains(&TableMaintenanceRecommendedAction::WaitForRetryBackoff)
|
||||
);
|
||||
assert!(backend.object_exists(bucket, &old).await.unwrap());
|
||||
}
|
||||
|
||||
@@ -8785,6 +8936,10 @@ mod tests {
|
||||
assert_eq!(report.job.job_id, running.job.job_id);
|
||||
assert_eq!(report.job.status, TableMetadataMaintenanceJobStatus::Running);
|
||||
assert_eq!(report.job.worker_id.as_deref(), Some("worker-a"));
|
||||
assert_eq!(
|
||||
report.job.recommended_actions,
|
||||
vec![TableMaintenanceRecommendedAction::WaitForActiveWorker]
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
@@ -8855,6 +9010,10 @@ mod tests {
|
||||
.as_deref()
|
||||
.is_some_and(|reason| reason.contains("lease expired"))
|
||||
);
|
||||
assert_eq!(
|
||||
expired.job.recommended_actions,
|
||||
vec![TableMaintenanceRecommendedAction::InvestigateFailure]
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
|
||||
Reference in New Issue
Block a user