mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-16 01:48:21 +00:00
Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 398d16f58c | |||
| 9e48c05493 |
Generated
-4
@@ -278,7 +278,6 @@ checksum = "312c1ea69e5fe9966e0029fb95aca8790100b85aff4f0d3b00a9337c74069a9c"
|
||||
dependencies = [
|
||||
"bigdecimal",
|
||||
"bon",
|
||||
"crc32fast",
|
||||
"digest 0.11.3",
|
||||
"log",
|
||||
"miniz_oxide 0.9.1",
|
||||
@@ -290,11 +289,9 @@ dependencies = [
|
||||
"serde",
|
||||
"serde_bytes",
|
||||
"serde_json",
|
||||
"snap",
|
||||
"strum",
|
||||
"thiserror 2.0.20",
|
||||
"uuid",
|
||||
"zstd",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
@@ -9203,7 +9200,6 @@ dependencies = [
|
||||
"serial_test",
|
||||
"sha2 0.11.0",
|
||||
"shadow-rs",
|
||||
"snap",
|
||||
"socket2",
|
||||
"subtle",
|
||||
"sysinfo",
|
||||
|
||||
+1
-1
@@ -171,7 +171,7 @@ tower = { version = "0.5.3" }
|
||||
tower-http = { version = "0.7.0" }
|
||||
|
||||
# Serialization and Data Formats
|
||||
apache-avro = { version = "0.22.0", features = ["snappy", "zstandard"] }
|
||||
apache-avro = "0.22.0"
|
||||
bytes = { version = "1.12.1" }
|
||||
bytesize = "2.7.0"
|
||||
byteorder = "1.5.0"
|
||||
|
||||
@@ -135,8 +135,7 @@ pub mod bucket {
|
||||
pub use crate::bucket::metadata_sys::ConfigWriteLockProbe;
|
||||
pub use crate::bucket::metadata_sys::{
|
||||
BucketMetadataMutationGuard, BucketMetadataSys, ObjectLockConfigState, acquire_bucket_metadata_transaction_lock,
|
||||
acquire_bucket_metadata_transaction_lock_for_incarnation, capture_bucket_metadata_incarnation, delete,
|
||||
delete_if_incarnation, delete_under_transaction_lock, get, get_accelerate_config, get_bucket_policy,
|
||||
capture_bucket_metadata_incarnation, delete, delete_if_incarnation, get, get_accelerate_config, get_bucket_policy,
|
||||
get_bucket_policy_raw, get_bucket_targets_config, get_config_from_disk, get_cors_config, get_durability_config,
|
||||
get_global_bucket_metadata_sys, get_lifecycle_config, get_logging_config, get_notification_config,
|
||||
get_object_lock_config, get_object_lock_config_state, get_public_access_block_config, get_quota_config,
|
||||
@@ -185,18 +184,17 @@ pub mod bucket {
|
||||
mrf_backlog_observability_snapshot,
|
||||
};
|
||||
pub use crate::bucket::replication::{
|
||||
BucketReplicationResyncStatus, BucketReplicationStat, BucketReplicationStats, BucketStats,
|
||||
DeleteReplicationConfigSnapshot, DeletedObjectReplicationInfo, DurableMrfBacklog, DynReplicationPool, InQueueMetric,
|
||||
MrfOpKind, MrfReplicateEntry, MustReplicateOptions, ObjectOpts, REMOTE_TARGET_CAPABILITY_CONTRACT_VERSION,
|
||||
REMOTE_TARGET_UNSUPPORTED_FIELDS, REMOTE_TARGET_WRITABLE_FIELDS, REPLICATE_INCOMING_DELETE,
|
||||
REPLICATION_CAPABILITY_CONTRACT_VERSION, REPLICATION_READ_ONLY_HISTORICAL_FIELDS, REPLICATION_WRITABLE_FIELDS,
|
||||
ReplicateDecision, ReplicateObjectInfo, ReplicationBatchAdmission, ReplicationConfig,
|
||||
ReplicationConfigStructureError, ReplicationConfigurationExt, ReplicationDeleteScheduleInput,
|
||||
ReplicationDeleteStateSource, ReplicationHealQueueResult, ReplicationObjectBridge, ReplicationObjectIO,
|
||||
ReplicationOperation, ReplicationPoolTrait, ReplicationPriority, ReplicationQueueAdmission, ReplicationScannerBridge,
|
||||
ReplicationState, ReplicationStats, ReplicationStatusType, ReplicationStorage, ReplicationTargetValidationError,
|
||||
ReplicationType, ResyncOpts, ResyncStatusType, RuntimeReplicationTargetBacklog, TargetReplicationResyncStatus,
|
||||
VersionPurgeStatusType, XferStats, commit_force_delete_intent, complete_force_delete_intent,
|
||||
BucketReplicationResyncStatus, BucketReplicationStats, BucketStats, DeleteReplicationConfigSnapshot,
|
||||
DeletedObjectReplicationInfo, DurableMrfBacklog, DynReplicationPool, MrfOpKind, MrfReplicateEntry,
|
||||
MustReplicateOptions, ObjectOpts, REMOTE_TARGET_CAPABILITY_CONTRACT_VERSION, REMOTE_TARGET_UNSUPPORTED_FIELDS,
|
||||
REMOTE_TARGET_WRITABLE_FIELDS, REPLICATE_INCOMING_DELETE, REPLICATION_CAPABILITY_CONTRACT_VERSION,
|
||||
REPLICATION_READ_ONLY_HISTORICAL_FIELDS, REPLICATION_WRITABLE_FIELDS, ReplicateDecision, ReplicateObjectInfo,
|
||||
ReplicationBatchAdmission, ReplicationConfig, ReplicationConfigStructureError, ReplicationConfigurationExt,
|
||||
ReplicationDeleteScheduleInput, ReplicationDeleteStateSource, ReplicationHealQueueResult, ReplicationObjectBridge,
|
||||
ReplicationObjectIO, ReplicationOperation, ReplicationPoolTrait, ReplicationPriority, ReplicationQueueAdmission,
|
||||
ReplicationScannerBridge, ReplicationState, ReplicationStats, ReplicationStatusType, ReplicationStorage,
|
||||
ReplicationTargetValidationError, ReplicationType, ResyncOpts, ResyncStatusType, RuntimeReplicationTargetBacklog,
|
||||
TargetReplicationResyncStatus, VersionPurgeStatusType, commit_force_delete_intent, complete_force_delete_intent,
|
||||
delete_replication_state_from_config, delete_replication_version_id, get_global_replication_pool,
|
||||
get_global_replication_stats, init_background_replication, invalid_replication_config_status_field,
|
||||
persist_force_delete_intent, read_durable_mrf_backlog, replication_state_to_filemeta, replication_status_to_filemeta,
|
||||
|
||||
@@ -656,16 +656,6 @@ pub async fn update_under_transaction_lock(
|
||||
update_under_config_write_guard(get_bucket_metadata_sys()?, guard, config_file, data).await
|
||||
}
|
||||
|
||||
/// Clear one config file while the caller holds this bucket's transaction lock.
|
||||
pub async fn delete_under_transaction_lock(
|
||||
guard: &BucketMetadataMutationGuard,
|
||||
bucket: &str,
|
||||
config_file: &str,
|
||||
) -> Result<OffsetDateTime> {
|
||||
guard.ensure_valid(bucket)?;
|
||||
delete_under_config_write_guard(get_bucket_metadata_sys()?, guard, config_file).await
|
||||
}
|
||||
|
||||
pub async fn update_quota_if_incarnation(
|
||||
bucket: &str,
|
||||
data: Vec<u8>,
|
||||
@@ -805,14 +795,6 @@ pub async fn acquire_bucket_metadata_transaction_lock(bucket: &str) -> Result<Bu
|
||||
acquire_config_write_guard(get_bucket_metadata_sys()?, bucket).await
|
||||
}
|
||||
|
||||
/// Acquire the bucket transaction lock only if its incarnation still matches.
|
||||
pub async fn acquire_bucket_metadata_transaction_lock_for_incarnation(
|
||||
bucket: &str,
|
||||
expected_incarnation_id: Uuid,
|
||||
) -> Result<BucketMetadataMutationGuard> {
|
||||
acquire_config_write_guard_for_incarnation(get_bucket_metadata_sys()?, bucket, Some(expected_incarnation_id)).await
|
||||
}
|
||||
|
||||
pub(crate) async fn acquire_bucket_metadata_transaction_lock_in(
|
||||
ctx: &crate::runtime::instance::InstanceContext,
|
||||
bucket: &str,
|
||||
|
||||
@@ -81,6 +81,6 @@ pub use replication_queue_boundary::{
|
||||
pub use replication_resync_boundary::{BucketReplicationResyncStatus, ResyncOpts, TargetReplicationResyncStatus};
|
||||
pub use replication_scanner_bridge::ReplicationScannerBridge;
|
||||
pub use replication_state::{ReplicationStats, RuntimeReplicationTargetBacklog};
|
||||
pub use replication_stats_boundary::{BucketReplicationStat, BucketReplicationStats, BucketStats, InQueueMetric, XferStats};
|
||||
pub use replication_stats_boundary::{BucketReplicationStats, BucketStats};
|
||||
pub use replication_storage_boundary::{ReplicationObjectIO, ReplicationStorage};
|
||||
pub(crate) use replication_target_config_bridge::ReplicationTargetConfigBridge;
|
||||
|
||||
@@ -704,12 +704,6 @@ impl ReplicationStats {
|
||||
} else {
|
||||
BucketReplicationStats::new()
|
||||
};
|
||||
// Stamp the serializable failure windows from the live samples: the
|
||||
// samples themselves do not cross the peer-RPC wire, so this snapshot
|
||||
// is what cluster aggregation and the metrics endpoints see.
|
||||
for stat in replication_stats.stats.values_mut() {
|
||||
stat.fail_stats.refresh_windows();
|
||||
}
|
||||
let uptime = if cache.contains_key(bucket) {
|
||||
SystemTime::now()
|
||||
.duration_since(SystemTime::UNIX_EPOCH)
|
||||
|
||||
@@ -15,9 +15,7 @@
|
||||
#[cfg(test)]
|
||||
pub(crate) use rustfs_replication::FailStats;
|
||||
pub(crate) use rustfs_replication::{
|
||||
ActiveWorkerStat, ProxyMetric, ProxyStatsCache, QueueCache, ReplicationMetricScope, SRMetricsSummary,
|
||||
ActiveWorkerStat, BucketReplicationStat, InQueueMetric, ProxyMetric, ProxyStatsCache, QueueCache, ReplicationMetricScope,
|
||||
SRMetricsSummary, XferStats,
|
||||
};
|
||||
// Public so the admin wire DTOs (rustfs/src/admin/replication_metrics_wire.rs)
|
||||
// can project the internal stats onto the minio-go response shapes through
|
||||
// the storage_api facade chain.
|
||||
pub use rustfs_replication::{BucketReplicationStat, BucketReplicationStats, BucketStats, InQueueMetric, XferStats};
|
||||
pub use rustfs_replication::{BucketReplicationStats, BucketStats};
|
||||
|
||||
@@ -40,7 +40,14 @@ impl ARN {
|
||||
|
||||
impl Display for ARN {
|
||||
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||
write!(f, "arn:rustfs:{}:{}:{}:{}", self.arn_type, self.region, self.id, self.bucket)
|
||||
// The `minio` partition is deliberate: madmin-go's ParseARN
|
||||
// hard-rejects any other partition, so native mc/madmin tooling can
|
||||
// only decode remote-target ARNs minted in this form (backlog#1675
|
||||
// P1-7). Legacy `arn:rustfs:` ARNs persisted by older releases stay
|
||||
// readable via the FromStr whitelist below; runtime matching between
|
||||
// targets and replication rules is by full-string equality, so mixed
|
||||
// partitions coexist safely.
|
||||
write!(f, "arn:minio:{}:{}:{}:{}", self.arn_type, self.region, self.id, self.bucket)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -48,7 +55,12 @@ impl FromStr for ARN {
|
||||
type Err = std::io::Error;
|
||||
|
||||
fn from_str(s: &str) -> Result<Self, Self::Err> {
|
||||
if !s.starts_with("arn:rustfs:") {
|
||||
// Partition whitelist, not just an `arn:` check: `BucketTargetType::
|
||||
// from_str(...).unwrap_or_default()` below never fails, so this is
|
||||
// the only structural gate rejecting foreign ARNs. `arn:rustfs:` is
|
||||
// the legacy partition and must stay accepted forever (persisted
|
||||
// bucket-targets.json / replication configs from older releases).
|
||||
if !s.starts_with("arn:minio:") && !s.starts_with("arn:rustfs:") {
|
||||
return Err(std::io::Error::new(std::io::ErrorKind::InvalidInput, "Invalid ARN format"));
|
||||
}
|
||||
|
||||
@@ -101,14 +113,50 @@ mod tests {
|
||||
}
|
||||
|
||||
/// RustFS commonly generates ARNs with an empty region:
|
||||
/// `arn:rustfs:replication::<deployment_id>:<bucket>`.
|
||||
/// `arn:minio:replication::<deployment_id>:<bucket>`.
|
||||
#[test]
|
||||
fn from_str_handles_empty_region_segment() {
|
||||
let parsed = ARN::from_str("arn:rustfs:replication::depl-123:bucket-a").expect("valid ARN must parse");
|
||||
let parsed = ARN::from_str("arn:minio:replication::depl-123:bucket-a").expect("valid ARN must parse");
|
||||
|
||||
assert_eq!(parsed.arn_type, BucketTargetType::ReplicationService);
|
||||
assert_eq!(parsed.region, "", "region segment is empty in this form");
|
||||
assert_eq!(parsed.id, "depl-123");
|
||||
assert_eq!(parsed.bucket, "bucket-a");
|
||||
}
|
||||
|
||||
/// madmin-go's `ParseARN` hard-rejects anything that does not start with
|
||||
/// `arn:minio:`, so generated ARNs must use the `minio` partition or the
|
||||
/// native mc/madmin tooling cannot decode remote-target listings.
|
||||
#[test]
|
||||
fn display_emits_minio_partition() {
|
||||
let arn = ARN::new(
|
||||
BucketTargetType::ReplicationService,
|
||||
"depl-123".to_string(),
|
||||
String::new(),
|
||||
"bucket-a".to_string(),
|
||||
);
|
||||
|
||||
assert_eq!(arn.to_string(), "arn:minio:replication::depl-123:bucket-a");
|
||||
}
|
||||
|
||||
/// Persisted bucket-targets.json files from older RustFS releases carry
|
||||
/// `arn:rustfs:` ARNs; the legacy partition must stay parseable forever.
|
||||
#[test]
|
||||
fn from_str_accepts_legacy_rustfs_partition() {
|
||||
let parsed = ARN::from_str("arn:rustfs:replication:us-east-1:depl-123:bucket-a").expect("legacy ARN must parse");
|
||||
|
||||
assert_eq!(parsed.arn_type, BucketTargetType::ReplicationService);
|
||||
assert_eq!(parsed.region, "us-east-1");
|
||||
assert_eq!(parsed.id, "depl-123");
|
||||
assert_eq!(parsed.bucket, "bucket-a");
|
||||
}
|
||||
|
||||
/// The partition whitelist is the only structural gate: `BucketTargetType::
|
||||
/// from_str(...).unwrap_or_default()` never fails, so any 6-segment string
|
||||
/// would otherwise parse as `type=None`.
|
||||
#[test]
|
||||
fn from_str_rejects_unknown_partition() {
|
||||
assert!(ARN::from_str("arn:aws:replication::depl-123:bucket-a").is_err());
|
||||
assert!(ARN::from_str("not-an-arn").is_err());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -190,17 +190,6 @@ pub(crate) const GET_METADATA_CACHE_REASON_VERSION_SUSPENDED: &str = "version_su
|
||||
pub(crate) const GET_METADATA_CACHE_REASON_VERSIONED: &str = "versioned";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_CONFLICTING_METADATA: &str = "conflicting_metadata";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DELETE_MARKER: &str = "delete_marker";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY: &str = "data_read_inline_body_verify";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_DELETED: &str = "data_read_inline_deleted";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY: &str = "data_read_inline_geometry";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_IDENTITY_MISMATCH: &str = "data_read_inline_identity_mismatch";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_PAYLOAD: &str = "data_read_inline_missing_payload";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_SHARD: &str = "data_read_inline_missing_shard";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_NOT_INLINE: &str = "data_read_inline_not_inline";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE: &str = "data_read_inline_part_shape";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_REMOTE: &str = "data_read_inline_remote";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE: &str = "data_read_inline_size";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_TRANSFORMED: &str = "data_read_inline_transformed";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_ERROR: &str = "error";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM: &str = "insufficient_quorum";
|
||||
pub(crate) const GET_METADATA_EARLY_STOP_REASON_NOT_FOUND: &str = "not_found";
|
||||
@@ -562,32 +551,6 @@ mod tests {
|
||||
assert_eq!(GET_METADATA_CACHE_REASON_VERSIONED, "versioned");
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_CONFLICTING_METADATA, "conflicting_metadata");
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_DELETE_MARKER, "delete_marker");
|
||||
assert_eq!(
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY,
|
||||
"data_read_inline_body_verify"
|
||||
);
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_DELETED, "data_read_inline_deleted");
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY, "data_read_inline_geometry");
|
||||
assert_eq!(
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_IDENTITY_MISMATCH,
|
||||
"data_read_inline_identity_mismatch"
|
||||
);
|
||||
assert_eq!(
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_PAYLOAD,
|
||||
"data_read_inline_missing_payload"
|
||||
);
|
||||
assert_eq!(
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_SHARD,
|
||||
"data_read_inline_missing_shard"
|
||||
);
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_NOT_INLINE, "data_read_inline_not_inline");
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE, "data_read_inline_part_shape");
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_REMOTE, "data_read_inline_remote");
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE, "data_read_inline_size");
|
||||
assert_eq!(
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_TRANSFORMED,
|
||||
"data_read_inline_transformed"
|
||||
);
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_ERROR, "error");
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM, "insufficient_quorum");
|
||||
assert_eq!(GET_METADATA_EARLY_STOP_REASON_NOT_FOUND, "not_found");
|
||||
|
||||
@@ -32,22 +32,15 @@ use crate::diagnostics::get::{
|
||||
GET_METADATA_CACHE_REASON_NOT_READ_DATA, GET_METADATA_CACHE_REASON_PART_NUMBER,
|
||||
GET_METADATA_CACHE_REASON_RAW_DATA_MOVEMENT_READ, GET_METADATA_CACHE_REASON_USABLE, GET_METADATA_CACHE_REASON_VERSION_ID,
|
||||
GET_METADATA_CACHE_REASON_VERSION_SUSPENDED, GET_METADATA_CACHE_REASON_VERSIONED,
|
||||
GET_METADATA_EARLY_STOP_REASON_CONFLICTING_METADATA, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY,
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_DELETED, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY,
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_IDENTITY_MISMATCH,
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_PAYLOAD,
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_SHARD, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_NOT_INLINE,
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_REMOTE,
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_TRANSFORMED,
|
||||
GET_METADATA_EARLY_STOP_REASON_DELETE_MARKER, GET_METADATA_EARLY_STOP_REASON_ERROR,
|
||||
GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM, GET_METADATA_EARLY_STOP_REASON_NOT_FOUND,
|
||||
GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST, GET_METADATA_EARLY_STOP_REASON_VALID_QUORUM,
|
||||
GET_METADATA_EARLY_STOP_REASON_VERSION_MATCH_QUORUM, GET_METADATA_EARLY_STOP_REASON_VERSION_NOT_FOUND,
|
||||
GET_METADATA_RESPONSE_CORRUPT, GET_METADATA_RESPONSE_DISK_NOT_FOUND, GET_METADATA_RESPONSE_ERROR,
|
||||
GET_METADATA_RESPONSE_IGNORED, GET_METADATA_RESPONSE_NOT_FOUND, GET_METADATA_RESPONSE_TIMEOUT, GET_METADATA_RESPONSE_VALID,
|
||||
GET_METADATA_RESPONSE_VERSION_NOT_FOUND, GET_OBJECT_PATH_CODEC_STREAMING, GET_OBJECT_PATH_DIRECT_MEMORY,
|
||||
GET_OBJECT_PATH_INTERNAL_META, GET_OBJECT_PATH_LEGACY_DUPLEX, GET_OBJECT_PATH_SET_DISK, GET_STAGE_DECODE,
|
||||
GET_STAGE_METADATA_CACHE_LOOKUP, GET_STAGE_METADATA_RESOLVE, GET_STAGE_RANGE, GET_STAGE_READER_SETUP,
|
||||
GET_METADATA_EARLY_STOP_REASON_CONFLICTING_METADATA, GET_METADATA_EARLY_STOP_REASON_DELETE_MARKER,
|
||||
GET_METADATA_EARLY_STOP_REASON_ERROR, GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM,
|
||||
GET_METADATA_EARLY_STOP_REASON_NOT_FOUND, GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST,
|
||||
GET_METADATA_EARLY_STOP_REASON_VALID_QUORUM, GET_METADATA_EARLY_STOP_REASON_VERSION_MATCH_QUORUM,
|
||||
GET_METADATA_EARLY_STOP_REASON_VERSION_NOT_FOUND, GET_METADATA_RESPONSE_CORRUPT, GET_METADATA_RESPONSE_DISK_NOT_FOUND,
|
||||
GET_METADATA_RESPONSE_ERROR, GET_METADATA_RESPONSE_IGNORED, GET_METADATA_RESPONSE_NOT_FOUND, GET_METADATA_RESPONSE_TIMEOUT,
|
||||
GET_METADATA_RESPONSE_VALID, GET_METADATA_RESPONSE_VERSION_NOT_FOUND, GET_OBJECT_PATH_CODEC_STREAMING,
|
||||
GET_OBJECT_PATH_DIRECT_MEMORY, GET_OBJECT_PATH_INTERNAL_META, GET_OBJECT_PATH_LEGACY_DUPLEX, GET_OBJECT_PATH_SET_DISK,
|
||||
GET_STAGE_DECODE, GET_STAGE_METADATA_CACHE_LOOKUP, GET_STAGE_METADATA_RESOLVE, GET_STAGE_RANGE, GET_STAGE_READER_SETUP,
|
||||
GET_STAGE_READER_SETUP_DROP_PENDING, GET_STAGE_READER_SETUP_SCHEDULE, GET_STAGE_READER_SETUP_WAIT_QUORUM,
|
||||
GET_STAGE_READER_TASK_BITROT_READER_INIT, GET_STAGE_READER_TASK_FILE_OPEN, GET_STAGE_READER_TASK_READER_CONSTRUCTION,
|
||||
GetObjectFailureReason, classify_disk_error, get_stage_timer_if_enabled, record_get_object_pipeline_failure,
|
||||
@@ -659,49 +652,36 @@ pub(in crate::set_disk) fn metadata_early_stop_candidate_matches(left: &FileInfo
|
||||
&& left.erasure.distribution == right.erasure.distribution
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) async fn data_read_early_stop_inline_body_miss_reason(
|
||||
pub(in crate::set_disk) async fn data_read_early_stop_inline_body_verified(
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
candidate: &FileInfo,
|
||||
parts_metadata: &[FileInfo],
|
||||
disks: &[Option<DiskStore>],
|
||||
) -> Option<&'static str> {
|
||||
// `inline_data` excludes remote objects; this diagnostic reports them separately.
|
||||
if !rustfs_utils::http::contains_key_str(&candidate.metadata, rustfs_utils::http::SUFFIX_INLINE_DATA) {
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_NOT_INLINE);
|
||||
}
|
||||
if candidate.is_compressed()
|
||||
) -> bool {
|
||||
if !candidate.inline_data()
|
||||
|| candidate.is_compressed()
|
||||
|| candidate
|
||||
.metadata
|
||||
.keys()
|
||||
.any(|key| rustfs_utils::http::is_object_encryption_marker(key))
|
||||
|| candidate.is_remote()
|
||||
|| candidate.deleted
|
||||
|| candidate.size <= 0
|
||||
|| candidate.parts.len() != 1
|
||||
|| !candidate.has_valid_erasure_geometry()
|
||||
{
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_TRANSFORMED);
|
||||
}
|
||||
if candidate.is_remote() {
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_REMOTE);
|
||||
}
|
||||
if candidate.deleted {
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_DELETED);
|
||||
}
|
||||
if candidate.size <= 0 {
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE);
|
||||
}
|
||||
if candidate.parts.len() != 1 {
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE);
|
||||
}
|
||||
if !candidate.has_valid_erasure_geometry() {
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY);
|
||||
return false;
|
||||
}
|
||||
|
||||
let Ok(object_size) = usize::try_from(candidate.size) else {
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE);
|
||||
return false;
|
||||
};
|
||||
if candidate.parts.first().is_none_or(|part| part.size != object_size) {
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE);
|
||||
return false;
|
||||
}
|
||||
if !can_try_inline_data_shards_direct(object_size, candidate.erasure.block_size) {
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE);
|
||||
return false;
|
||||
}
|
||||
|
||||
let Ok(erasure) = coding::Erasure::try_new_with_options(
|
||||
@@ -710,18 +690,18 @@ pub(in crate::set_disk) async fn data_read_early_stop_inline_body_miss_reason(
|
||||
candidate.erasure.block_size,
|
||||
candidate.uses_legacy_checksum,
|
||||
) else {
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY);
|
||||
return false;
|
||||
};
|
||||
let data_files =
|
||||
match collect_inline_data_shard_fileinfos_by_index_or_reason(parts_metadata, candidate, erasure.data_shards, |index| {
|
||||
let Some(data_files) =
|
||||
collect_inline_data_shard_fileinfos_by_index(parts_metadata, candidate, erasure.data_shards, |index| {
|
||||
disks.get(index).is_some_and(Option::is_some)
|
||||
}) {
|
||||
Ok(data_files) => data_files,
|
||||
Err(reason) => return Some(reason),
|
||||
};
|
||||
})
|
||||
else {
|
||||
return false;
|
||||
};
|
||||
|
||||
let Some(part) = candidate.parts.first() else {
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE);
|
||||
return false;
|
||||
};
|
||||
let checksum_info = candidate.erasure.get_checksum_info(part.number);
|
||||
let checksum_algo = if candidate.uses_legacy_checksum && checksum_info.algorithm == HashAlgorithm::HighwayHash256S {
|
||||
@@ -741,13 +721,12 @@ pub(in crate::set_disk) async fn data_read_early_stop_inline_body_miss_reason(
|
||||
let Ok(mut readers) =
|
||||
build_inline_bitrot_readers_from_refs(&data_files, bucket, object, read_length, shard_size, &checksum_algo, false).await
|
||||
else {
|
||||
return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY);
|
||||
return false;
|
||||
};
|
||||
|
||||
match try_read_inline_data_shards_direct(&mut readers, erasure.data_shards, read_length, object_size).await {
|
||||
Some(body) if body.len() == object_size => None,
|
||||
_ => Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY),
|
||||
}
|
||||
try_read_inline_data_shards_direct(&mut readers, erasure.data_shards, read_length, object_size)
|
||||
.await
|
||||
.is_some_and(|body| body.len() == object_size)
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn classify_metadata_response_error(err: &DiskError) -> &'static str {
|
||||
@@ -2490,7 +2469,6 @@ impl SetDisks {
|
||||
let mut next_fanout_index = 0usize;
|
||||
let mut scheduled_count = 0usize;
|
||||
let mut force_full_wait = false;
|
||||
let mut final_miss_reason_override = None;
|
||||
let spawn_read_version =
|
||||
|join_set: &mut JoinSet<(usize, disk::error::Result<FileInfo>, Duration)>, index: usize, disk: Option<DiskStore>| {
|
||||
let task_opts = opts;
|
||||
@@ -2563,29 +2541,17 @@ impl SetDisks {
|
||||
.or_else(|| accumulator.version_early_stop_decision())
|
||||
{
|
||||
let should_return_early = if read_data {
|
||||
match accumulator.candidate.as_ref() {
|
||||
Some(candidate) => match data_read_early_stop_inline_body_miss_reason(
|
||||
bucket.as_ref(),
|
||||
object.as_ref(),
|
||||
candidate,
|
||||
&ress,
|
||||
disks,
|
||||
)
|
||||
.await
|
||||
{
|
||||
None => true,
|
||||
Some(reason) => {
|
||||
force_full_wait = true;
|
||||
final_miss_reason_override = Some(reason);
|
||||
false
|
||||
}
|
||||
},
|
||||
None => {
|
||||
force_full_wait = true;
|
||||
final_miss_reason_override = Some(GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM);
|
||||
false
|
||||
let allow_data_read_early_stop = match accumulator.candidate.as_ref() {
|
||||
Some(candidate) => {
|
||||
data_read_early_stop_inline_body_verified(bucket.as_ref(), object.as_ref(), candidate, &ress, disks)
|
||||
.await
|
||||
}
|
||||
None => false,
|
||||
};
|
||||
if !allow_data_read_early_stop {
|
||||
force_full_wait = true;
|
||||
}
|
||||
allow_data_read_early_stop
|
||||
} else {
|
||||
true
|
||||
};
|
||||
@@ -2647,12 +2613,7 @@ impl SetDisks {
|
||||
}
|
||||
}
|
||||
|
||||
let accumulator_miss_reason = accumulator.final_miss_reason();
|
||||
let final_miss_reason = match (final_miss_reason_override, accumulator_miss_reason) {
|
||||
(Some(reason), GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM) => reason,
|
||||
_ => accumulator_miss_reason,
|
||||
};
|
||||
rustfs_io_metrics::record_get_object_metadata_early_stop_miss(metrics_path, final_miss_reason);
|
||||
rustfs_io_metrics::record_get_object_metadata_early_stop_miss(metrics_path, accumulator.final_miss_reason());
|
||||
rustfs_io_metrics::record_get_object_metadata_early_stop_saved_responses(metrics_path, 0);
|
||||
rustfs_io_metrics::record_get_object_metadata_fanout_lifecycle(metrics_path, scheduled_count, scheduled_count, 0);
|
||||
let diagnostics = MetadataFanoutDiagnostics::new(fanout_start.elapsed(), observations);
|
||||
@@ -6106,133 +6067,11 @@ mod tests {
|
||||
.clone();
|
||||
|
||||
assert!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &candidate, &parts_metadata, &disks)
|
||||
.await
|
||||
.is_none(),
|
||||
data_read_early_stop_inline_body_verified(bucket, object, &candidate, &parts_metadata, &disks).await,
|
||||
"legacy inline metadata must use the legacy bitrot shard sizing and checksum algorithm"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn data_read_early_stop_reports_inline_miss_reasons() {
|
||||
let bucket = "inline-data-get-miss-reason-bucket";
|
||||
let object = "inline-data-get-miss-reason-object";
|
||||
let payload = b"verified inline payload";
|
||||
let (_dirs, disks) = call_counter_local_disks(bucket, 4).await;
|
||||
let files = inline_metadata_fanout_fileinfos_with_mode(bucket, object, payload, false).await;
|
||||
let distribution = files
|
||||
.first()
|
||||
.map(|file| file.erasure.distribution.clone())
|
||||
.expect("fixture should include metadata");
|
||||
let order = bounded_metadata_fanout_order(bucket, object, 4, 2);
|
||||
let mut parts_metadata = vec![FileInfo::default(); 4];
|
||||
for disk_index in order.into_iter().take(3) {
|
||||
let block_index = distribution
|
||||
.get(disk_index)
|
||||
.copied()
|
||||
.expect("fixture distribution should cover every disk");
|
||||
parts_metadata[disk_index] = files
|
||||
.get(block_index.checked_sub(1).expect("erasure block indexes are one-based"))
|
||||
.expect("fixture should include every distributed shard")
|
||||
.clone();
|
||||
}
|
||||
let candidate = parts_metadata
|
||||
.iter()
|
||||
.find(|file| file.name == object)
|
||||
.expect("fixture should include observed metadata")
|
||||
.clone();
|
||||
let data_disk = distribution
|
||||
.iter()
|
||||
.position(|block_index| *block_index == 1)
|
||||
.expect("fixture distribution should include first data shard");
|
||||
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &candidate, &parts_metadata, &disks).await,
|
||||
None
|
||||
);
|
||||
|
||||
let mut not_inline = candidate.clone();
|
||||
rustfs_utils::http::remove_str(&mut not_inline.metadata, rustfs_utils::http::SUFFIX_INLINE_DATA);
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, ¬_inline, &parts_metadata, &disks).await,
|
||||
Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_NOT_INLINE)
|
||||
);
|
||||
|
||||
let mut remote = candidate.clone();
|
||||
remote.transition_status = TRANSITION_COMPLETE.to_string();
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &remote, &parts_metadata, &disks).await,
|
||||
Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_REMOTE)
|
||||
);
|
||||
|
||||
let mut transformed = candidate.clone();
|
||||
rustfs_utils::http::insert_str(&mut transformed.metadata, rustfs_utils::http::SUFFIX_COMPRESSION, "zstd".to_string());
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &transformed, &parts_metadata, &disks).await,
|
||||
Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_TRANSFORMED)
|
||||
);
|
||||
|
||||
let mut deleted = candidate.clone();
|
||||
deleted.deleted = true;
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &deleted, &parts_metadata, &disks).await,
|
||||
Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_DELETED)
|
||||
);
|
||||
|
||||
let mut zero_size = candidate.clone();
|
||||
zero_size.size = 0;
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &zero_size, &parts_metadata, &disks).await,
|
||||
Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE)
|
||||
);
|
||||
|
||||
let mut multipart = candidate.clone();
|
||||
multipart.parts.push(multipart.parts[0].clone());
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &multipart, &parts_metadata, &disks).await,
|
||||
Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE)
|
||||
);
|
||||
|
||||
let mut invalid_geometry = candidate.clone();
|
||||
invalid_geometry.erasure.data_blocks = 0;
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &invalid_geometry, &parts_metadata, &disks).await,
|
||||
Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY)
|
||||
);
|
||||
|
||||
let mut missing_shard = parts_metadata.clone();
|
||||
missing_shard[data_disk] = FileInfo::default();
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &candidate, &missing_shard, &disks).await,
|
||||
Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_SHARD)
|
||||
);
|
||||
|
||||
let mut missing_payload = parts_metadata.clone();
|
||||
missing_payload[data_disk].data = None;
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &candidate, &missing_payload, &disks).await,
|
||||
Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_PAYLOAD)
|
||||
);
|
||||
|
||||
let mut identity_mismatch = parts_metadata.clone();
|
||||
identity_mismatch[data_disk].version_id = Some(Uuid::new_v4());
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &candidate, &identity_mismatch, &disks).await,
|
||||
Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_IDENTITY_MISMATCH)
|
||||
);
|
||||
|
||||
let mut corrupt = parts_metadata.clone();
|
||||
if let Some(data) = corrupt[data_disk].data.as_mut() {
|
||||
let mut corrupt_data = data.to_vec();
|
||||
corrupt_data[0] ^= 0x01;
|
||||
*data = Bytes::from(corrupt_data);
|
||||
}
|
||||
assert_eq!(
|
||||
data_read_early_stop_inline_body_miss_reason(bucket, object, &candidate, &corrupt, &disks).await,
|
||||
Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
#[serial_test::serial]
|
||||
fn metadata_fanout_lifecycle_records_real_early_stop_abort() {
|
||||
@@ -6322,7 +6161,7 @@ mod tests {
|
||||
&[
|
||||
("path", GET_OBJECT_PATH_INTERNAL_META),
|
||||
("decision", "miss"),
|
||||
("reason", GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY),
|
||||
("reason", GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM),
|
||||
],
|
||||
),
|
||||
1,
|
||||
@@ -6334,7 +6173,7 @@ mod tests {
|
||||
&[
|
||||
("path", GET_OBJECT_PATH_LEGACY_DUPLEX),
|
||||
("decision", "miss"),
|
||||
("reason", GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY),
|
||||
("reason", GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM),
|
||||
],
|
||||
),
|
||||
0,
|
||||
|
||||
@@ -59,10 +59,7 @@ use crate::client::{object_api_utils::get_raw_etag, transition_api::ReaderImpl};
|
||||
use crate::cluster::rpc::heal_bucket_local_on_disks;
|
||||
use crate::data_usage::record_compression_total_memory;
|
||||
use crate::diagnostics::get::{
|
||||
GET_CODEC_STREAMING_OBJECT_CLASS_PLAIN_SINGLE_PART, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY,
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_IDENTITY_MISMATCH,
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_PAYLOAD,
|
||||
GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_SHARD, GET_OBJECT_PATH_BODY_CACHE, GET_OBJECT_PATH_CODEC_STREAMING,
|
||||
GET_CODEC_STREAMING_OBJECT_CLASS_PLAIN_SINGLE_PART, GET_OBJECT_PATH_BODY_CACHE, GET_OBJECT_PATH_CODEC_STREAMING,
|
||||
GET_OBJECT_PATH_CODEC_STREAMING_LEGACY_ENGINE, GET_OBJECT_PATH_CODEC_STREAMING_RUSTFS_ENGINE, GET_OBJECT_PATH_DIRECT_MEMORY,
|
||||
GET_OBJECT_PATH_EMPTY, GET_OBJECT_PATH_INLINE_DIRECT, GET_OBJECT_PATH_INTERNAL_META, GET_OBJECT_PATH_LEGACY_DUPLEX,
|
||||
GET_OBJECT_PATH_REMOTE_TRANSITION, GET_OBJECT_PATH_SET_DISK, GET_STAGE_DECODE, GET_STAGE_EMIT, GET_STAGE_INLINE_PREPARE,
|
||||
@@ -3866,20 +3863,11 @@ fn inline_erasure_shard_file_offset(
|
||||
}
|
||||
|
||||
fn collect_inline_data_shard_fileinfos_by_index<'a>(
|
||||
parts_metadata: &'a [FileInfo],
|
||||
fi: &FileInfo,
|
||||
data_shards: usize,
|
||||
disk_is_online: impl FnMut(usize) -> bool,
|
||||
) -> Option<Vec<&'a FileInfo>> {
|
||||
collect_inline_data_shard_fileinfos_by_index_or_reason(parts_metadata, fi, data_shards, disk_is_online).ok()
|
||||
}
|
||||
|
||||
fn collect_inline_data_shard_fileinfos_by_index_or_reason<'a>(
|
||||
parts_metadata: &'a [FileInfo],
|
||||
fi: &FileInfo,
|
||||
data_shards: usize,
|
||||
mut disk_is_online: impl FnMut(usize) -> bool,
|
||||
) -> std::result::Result<Vec<&'a FileInfo>, &'static str> {
|
||||
) -> Option<Vec<&'a FileInfo>> {
|
||||
let distribution = &fi.erasure.distribution;
|
||||
let mut data_files = vec![None; data_shards];
|
||||
|
||||
@@ -3887,35 +3875,27 @@ fn collect_inline_data_shard_fileinfos_by_index_or_reason<'a>(
|
||||
if !disk_is_online(disk_index) {
|
||||
continue;
|
||||
}
|
||||
let Some(&block_index) = distribution.get(disk_index) else {
|
||||
return Err(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY);
|
||||
};
|
||||
let block_index = *distribution.get(disk_index)?;
|
||||
if block_index == 0 || block_index > data_shards {
|
||||
continue;
|
||||
}
|
||||
if file_info.name.is_empty() {
|
||||
return Err(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_SHARD);
|
||||
}
|
||||
if file_info.erasure.index != block_index {
|
||||
return Err(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_IDENTITY_MISMATCH);
|
||||
continue;
|
||||
}
|
||||
if !file_info.has_valid_erasure_geometry() {
|
||||
return Err(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY);
|
||||
continue;
|
||||
}
|
||||
if !core::io_primitives::metadata_early_stop_candidate_matches(file_info, fi) {
|
||||
return Err(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_IDENTITY_MISMATCH);
|
||||
continue;
|
||||
}
|
||||
if file_info.data.as_ref().is_none_or(|data| data.is_empty()) {
|
||||
return Err(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_PAYLOAD);
|
||||
continue;
|
||||
}
|
||||
|
||||
data_files[block_index - 1] = Some(file_info);
|
||||
}
|
||||
|
||||
data_files
|
||||
.into_iter()
|
||||
.collect::<Option<Vec<_>>>()
|
||||
.ok_or(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_SHARD)
|
||||
data_files.into_iter().collect()
|
||||
}
|
||||
|
||||
impl SetDisks {
|
||||
|
||||
@@ -135,7 +135,10 @@ impl FileMeta {
|
||||
let i = buf.len() as u64;
|
||||
|
||||
// check version, buf = buf[8..]
|
||||
let (buf, _, _) = Self::check_xl2_v1(buf)?;
|
||||
let (buf, _, _) = Self::check_xl2_v1(buf).map_err(|e| {
|
||||
error!("failed to check XL2 v1 format: {}", e);
|
||||
e
|
||||
})?;
|
||||
|
||||
if buf.len() < 5 {
|
||||
error!(
|
||||
|
||||
@@ -520,14 +520,6 @@ struct FailureSample {
|
||||
pub struct FailStats {
|
||||
pub count: i64,
|
||||
pub size: i64,
|
||||
/// Rolling-window snapshots refreshed at collection time
|
||||
/// ([`Self::refresh_windows`]). The raw samples (`recent`) are process
|
||||
/// local (serde-skipped), so these fields are what survives the peer-RPC
|
||||
/// wire and [`Self::merge`]-based cluster aggregation.
|
||||
#[serde(default)]
|
||||
pub last_minute: FailedMetric,
|
||||
#[serde(default)]
|
||||
pub last_hour: FailedMetric,
|
||||
#[serde(skip)]
|
||||
recent: VecDeque<FailureSample>,
|
||||
}
|
||||
@@ -545,17 +537,6 @@ impl FailStats {
|
||||
self.prune(observed_at);
|
||||
}
|
||||
|
||||
/// Recompute the serializable rolling-window snapshots from the local
|
||||
/// samples. Called at the collection point (per-node stats snapshot),
|
||||
/// never on the failure hot path — the two deque scans are O(window) and
|
||||
/// `add_size` runs under the bucket-stats write lock. Only meaningful on
|
||||
/// the live per-node struct: a deserialized or merged struct has no
|
||||
/// samples, and refreshing it would wipe the aggregated windows.
|
||||
pub fn refresh_windows(&mut self) {
|
||||
self.last_minute = self.recent_since(Duration::from_secs(60));
|
||||
self.last_hour = self.recent_since(Duration::from_secs(3600));
|
||||
}
|
||||
|
||||
fn prune(&mut self, observed_at: Instant) {
|
||||
while self
|
||||
.recent
|
||||
@@ -584,16 +565,6 @@ impl FailStats {
|
||||
Self {
|
||||
count: self.count.saturating_add(other.count),
|
||||
size: self.size.saturating_add(other.size),
|
||||
// The window snapshots sum across nodes; the raw samples do not
|
||||
// travel and stay empty on aggregated structs.
|
||||
last_minute: FailedMetric {
|
||||
count: self.last_minute.count.saturating_add(other.last_minute.count),
|
||||
size: self.last_minute.size.saturating_add(other.last_minute.size),
|
||||
},
|
||||
last_hour: FailedMetric {
|
||||
count: self.last_hour.count.saturating_add(other.last_hour.count),
|
||||
size: self.last_hour.size.saturating_add(other.last_hour.size),
|
||||
},
|
||||
recent: VecDeque::new(),
|
||||
}
|
||||
}
|
||||
@@ -665,9 +636,7 @@ impl BucketReplicationStat {
|
||||
}
|
||||
|
||||
pub fn update_xfer_rate(&mut self, size: i64, duration: Duration) {
|
||||
// Same boundary as the worker-pool split and minio-go's
|
||||
// Large/Small transfer-summary labels: >= 128 MiB is "large".
|
||||
if size >= crate::runtime::MIN_LARGE_OBJ_SIZE {
|
||||
if size > 1024 * 1024 {
|
||||
self.xfer_rate_lrg.add_size(size, duration);
|
||||
} else {
|
||||
self.xfer_rate_sml.add_size(size, duration);
|
||||
|
||||
@@ -102,7 +102,7 @@ bytes.workspace = true
|
||||
hex-simd.workspace = true
|
||||
|
||||
[dev-dependencies]
|
||||
tracing-subscriber = { workspace = true, features = ["json", "env-filter", "time"] }
|
||||
tracing-subscriber = { workspace = true, features = ["env-filter", "time"] }
|
||||
serial_test = { workspace = true }
|
||||
temp-env = { workspace = true }
|
||||
tempfile = { workspace = true }
|
||||
|
||||
@@ -65,7 +65,6 @@ const LOG_SUBSYSTEM_FOLDER: &str = "folder";
|
||||
const LOG_SUBSYSTEM_LIFECYCLE: &str = "lifecycle";
|
||||
const LOG_SUBSYSTEM_HEAL: &str = "heal";
|
||||
const EVENT_SCANNER_FOLDER_STATE: &str = "scanner_folder_state";
|
||||
const EVENT_SCANNER_METADATA_CORRUPT: &str = "scanner_metadata_corrupt";
|
||||
const EVENT_SCANNER_LIFECYCLE_ACTION: &str = "scanner_lifecycle_action";
|
||||
const EVENT_SCANNER_HEAL_ADMISSION: &str = "scanner_heal_admission";
|
||||
const EVENT_SCANNER_ALERT_STATE: &str = "scanner_alert_state";
|
||||
@@ -2155,34 +2154,17 @@ impl FolderScanner {
|
||||
self.record_failed(&item.path);
|
||||
|
||||
if should_log_failed_object(into.failed_objects) {
|
||||
if let GetSizeFailureAction::HealMetadata { object } = &failure_action {
|
||||
error!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_METADATA_CORRUPT,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_FOLDER,
|
||||
drive = %self.local_disk.path().display(),
|
||||
bucket = %item.bucket,
|
||||
object = %object,
|
||||
metadata_path = %item.path,
|
||||
failed_objects = into.failed_objects,
|
||||
state = "metadata_corrupt",
|
||||
error = %e,
|
||||
"Scanner detected corrupt object metadata"
|
||||
);
|
||||
} else {
|
||||
warn!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_FOLDER_STATE,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_FOLDER,
|
||||
path = %item.path,
|
||||
failed_objects = into.failed_objects,
|
||||
state = "get_size_failed",
|
||||
error = %e,
|
||||
"Scanner folder failed to get object size"
|
||||
);
|
||||
}
|
||||
warn!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_FOLDER_STATE,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_FOLDER,
|
||||
path = %item.path,
|
||||
failed_objects = into.failed_objects,
|
||||
state = "get_size_failed",
|
||||
error = %e,
|
||||
"Scanner folder failed to get object size"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -3072,59 +3054,12 @@ mod tests {
|
||||
use crate::{DiskOption, Endpoint, STORAGE_FORMAT_FILE, TierStats, new_disk, storageclass};
|
||||
use rustfs_filemeta::{FileInfo, FileMeta};
|
||||
use serial_test::serial;
|
||||
use std::io::Write;
|
||||
#[cfg(unix)]
|
||||
use std::os::unix::fs::{PermissionsExt, symlink};
|
||||
use std::sync::Mutex;
|
||||
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
||||
use temp_env::{with_var, with_var_unset};
|
||||
use tracing_subscriber::fmt::MakeWriter;
|
||||
use uuid::Uuid;
|
||||
|
||||
#[derive(Clone, Default)]
|
||||
struct CapturedLogs {
|
||||
buffer: Arc<Mutex<Vec<u8>>>,
|
||||
}
|
||||
|
||||
struct CapturedLogWriter {
|
||||
buffer: Arc<Mutex<Vec<u8>>>,
|
||||
}
|
||||
|
||||
impl CapturedLogs {
|
||||
fn contents(&self) -> String {
|
||||
let buffer = self
|
||||
.buffer
|
||||
.lock()
|
||||
.expect("captured logs mutex should not be poisoned")
|
||||
.clone();
|
||||
String::from_utf8(buffer).expect("captured logs should be valid UTF-8")
|
||||
}
|
||||
}
|
||||
|
||||
impl Write for CapturedLogWriter {
|
||||
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
|
||||
self.buffer
|
||||
.lock()
|
||||
.expect("captured logs mutex should not be poisoned")
|
||||
.extend_from_slice(buf);
|
||||
Ok(buf.len())
|
||||
}
|
||||
|
||||
fn flush(&mut self) -> std::io::Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
impl<'a> MakeWriter<'a> for CapturedLogs {
|
||||
type Writer = CapturedLogWriter;
|
||||
|
||||
fn make_writer(&'a self) -> Self::Writer {
|
||||
CapturedLogWriter {
|
||||
buffer: Arc::clone(&self.buffer),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn scanner_size_summary_application_saturates_usage_counters() {
|
||||
let target = "arn:minio:replication::target".to_string();
|
||||
@@ -4607,19 +4542,9 @@ mod tests {
|
||||
assert!(budget.entries_visited() >= 1);
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_scan_folder_corrupt_xl_meta_stops_erasure_data_dir_descent() {
|
||||
let logs = CapturedLogs::default();
|
||||
let subscriber = tracing_subscriber::fmt()
|
||||
.json()
|
||||
.with_max_level(tracing::Level::ERROR)
|
||||
.with_writer(logs.clone())
|
||||
.with_ansi(false)
|
||||
.without_time()
|
||||
.finish();
|
||||
let _subscriber_guard = tracing::subscriber::set_default(subscriber);
|
||||
|
||||
let (mut scanner, temp_dir) = build_test_scanner().await;
|
||||
let _guard = TestGuard::new(60, 100, &mut scanner, temp_dir.clone());
|
||||
|
||||
@@ -4671,30 +4596,6 @@ mod tests {
|
||||
assert!(!budget.budget_elapsed());
|
||||
assert_eq!(budget.reason(), None);
|
||||
|
||||
let captured = logs.contents();
|
||||
assert!(
|
||||
!captured.contains("failed to check XL2 v1 format"),
|
||||
"the context-free filemeta parser error must not be emitted"
|
||||
);
|
||||
let events = captured
|
||||
.lines()
|
||||
.map(|line| serde_json::from_str::<serde_json::Value>(line).expect("captured scanner log should be valid JSON"))
|
||||
.filter(|line| line["fields"]["event"] == EVENT_SCANNER_METADATA_CORRUPT)
|
||||
.collect::<Vec<_>>();
|
||||
assert_eq!(
|
||||
events.len(),
|
||||
1,
|
||||
"one corrupt metadata observation must emit one scanner-owned diagnostic event"
|
||||
);
|
||||
let fields = &events[0]["fields"];
|
||||
assert_eq!(fields["component"], LOG_COMPONENT_SCANNER);
|
||||
assert_eq!(fields["subsystem"], LOG_SUBSYSTEM_FOLDER);
|
||||
assert_eq!(fields["drive"], temp_dir.to_string_lossy().as_ref());
|
||||
assert_eq!(fields["bucket"], "bucket");
|
||||
assert_eq!(fields["object"], "object");
|
||||
assert_eq!(fields["metadata_path"], metadata_path.to_string_lossy().as_ref());
|
||||
assert_eq!(fields["state"], "metadata_corrupt");
|
||||
|
||||
let retry_budget = ScannerCycleBudget::new_with_progress_tracking(
|
||||
&parent,
|
||||
crate::scanner_budget::ScannerCycleBudgetConfig {
|
||||
|
||||
@@ -3849,6 +3849,17 @@ impl ScannerIODisk for Disk {
|
||||
let fivs = match meta.get_file_info_versions(item.bucket.as_str(), item.object_path().as_str(), false) {
|
||||
Ok(versions) => versions,
|
||||
Err(e) => {
|
||||
error!(
|
||||
target: "rustfs::scanner::io",
|
||||
event = EVENT_SCANNER_DISK_BUCKET_STATE,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_IO,
|
||||
bucket = %item.bucket,
|
||||
object = %item.object_path(),
|
||||
state = "file_info_versions_failed",
|
||||
error = %e,
|
||||
"Scanner disk bucket failed to resolve file info versions"
|
||||
);
|
||||
return Err(scanner_metadata_corrupt_error(
|
||||
format!("failed to resolve file info versions: {e}"),
|
||||
&item.bucket,
|
||||
|
||||
@@ -61,17 +61,17 @@ catalog extension.
|
||||
|
||||
| Area | Status | Covered behavior |
|
||||
|---|---|---|
|
||||
| Catalog config | Supported | `GET /v1/config` advertises RustFS catalog defaults and only the supported OpenAPI REST paths in `endpoints`. RustFS administration, maintenance, migration, diagnostics, refs, and metadata-location extensions remain available but are not presented as standard Iceberg REST endpoints. |
|
||||
| Catalog config | Supported | `GET /v1/config` advertises RustFS catalog defaults and route capabilities. |
|
||||
| Table bucket discovery | Supported | `PUT` and `GET /v1/buckets/{warehouse}` enable and inspect table bucket state. |
|
||||
| Namespaces | Supported | Create, list, load, existence check, and drop namespace routes are registered on both catalog prefixes. List responses support Iceberg REST `pageSize`/`pageToken` pagination with context-bound tokens and bounded catalog-store reads. Namespace identifiers are limited to 512 ASCII characters so persisted paths and stateless continuation tokens remain bounded. |
|
||||
| Tables | Supported | Create, register, list, load, existence check, commit, metadata-location get/update, and drop table routes are registered on both catalog prefixes. Table and view listings support Iceberg REST `pageSize`/`pageToken` pagination with context-bound tokens and bounded catalog-store reads. Commit identifiers must match the URL resource; unknown requirements, updates, and snapshot operations fail as bad requests; staged create, register overwrite, purge-on-drop, and v3-only encryption-key updates return an explicit unsupported-operation response. Standard statistics, partition statistics, and schema/spec cleanup updates are accepted. |
|
||||
| Commit CAS | Supported | Single-table commits validate base metadata, expected version token, referenced object existence, warehouse scope, and Iceberg commit requirements before advancing the current metadata pointer. Externally supplied metadata transitions preserve monotonic column, partition, and sequence assignment watermarks and immutable definitions for retained schemas, partition specs, sort orders, and snapshots. Standard commits preserve the normal commit-token file name and use an immutable-table-scoped fallback when rename followed by source-name reuse would otherwise collide at the same generation and commit ID. The catalog does not advertise `idempotency-key-lifetime`; clients must treat standard mutation-wide `Idempotency-Key` semantics as unsupported. |
|
||||
| Tables | Supported | Create, register, list, load, existence check, commit, metadata-location get/update, and drop table routes are registered on both catalog prefixes. Table and view listings support Iceberg REST `pageSize`/`pageToken` pagination with context-bound tokens and bounded catalog-store reads. |
|
||||
| Commit CAS | Supported | Single-table commits validate base metadata, expected version token, referenced object existence, warehouse scope, and Iceberg commit requirements before advancing the current metadata pointer. Standard commits preserve the normal commit-token file name and use an immutable-table-scoped fallback when rename followed by source-name reuse would otherwise collide at the same generation and commit ID. |
|
||||
| Commit recovery | Supported | Commit log, idempotency lookup, diagnostics, and recovery routes expose staged/finalization gaps and repair safe idempotency gaps without moving the table pointer. |
|
||||
| Snapshot refs | Supported | Refs can be listed, created or replaced, and deleted through catalog commits. `main` is protected and refs with explicit retention require forced delete. |
|
||||
| Iceberg views | Supported | Basic create, list, load, replace, existence check, and drop routes persist view metadata with view-scoped authorization. Replace identifiers must match the URL resource, `schema-id: -1` resolves to the last added schema, one commit timestamp is used consistently, and only Iceberg view format version 1 is accepted. |
|
||||
| Table credentials endpoint | Supported | Returns an empty `storage-credentials` list by default. Returns table-scoped temporary credentials only when credential vending is enabled. Credential responses set `Cache-Control: no-store, private`, `Pragma: no-cache`, and `Expires: 0`. |
|
||||
| Iceberg views | Supported | Basic create, list, load, replace, existence check, and drop routes persist view metadata with view-scoped authorization. |
|
||||
| Table credentials endpoint | Supported | Returns an empty `storage-credentials` list by default. Returns table-scoped temporary credentials only when credential vending is enabled. |
|
||||
| Catalog diagnostics and export | Supported | Exposes recovery state, consistency state, backing manifest, recoverable commit-log WAL state, strong backing migration target, single-active-writer policy, and scale validation matrix. |
|
||||
| Catalog import and rollback | Supported | Import/register and online rollback use catalog validation and commit paths rather than direct pointer mutation. Online rollback accepts only a forward-safe metadata target that preserves assignment watermarks and retained definitions. Restoring an older target that lowers those watermarks is an offline disaster-recovery operation and requires every writer to be stopped. |
|
||||
| Catalog import and rollback | Supported | Import/register and rollback use catalog validation and commit paths rather than direct pointer mutation. |
|
||||
| External catalog bridge | Supported operator path | Operator-supplied metadata pointer sync/import is supported for external catalog identity boundaries. Online vendor SDK polling and policy mirroring are not claimed. |
|
||||
| Multi-table transactions | Not claimed | RustFS currently claims single-table commit atomicity only. |
|
||||
|
||||
|
||||
+1
-2
@@ -278,8 +278,6 @@ rustfs-signer.workspace = true
|
||||
serde = { workspace = true, features = ["derive"] }
|
||||
serde_json = { workspace = true, features = ["raw_value"] }
|
||||
serde_urlencoded = { workspace = true }
|
||||
snap.workspace = true
|
||||
zstd.workspace = true
|
||||
|
||||
# Cryptography and Security
|
||||
rustls = { workspace = true, default-features = false, features = ["aws-lc-rs", "logging", "tls12", "prefer-post-quantum", "std"] }
|
||||
@@ -357,6 +355,7 @@ rcgen = { workspace = true }
|
||||
rustfs-test-utils.workspace = true
|
||||
# diagnose_e2e fixtures (archives are generated in-test, never checked in)
|
||||
zip = { workspace = true }
|
||||
zstd = { workspace = true }
|
||||
# Enables the shared MockWarmBackend / xl.meta assertion helpers exposed via
|
||||
# the ecstore `api::tier::test_util` facade module (rustfs/backlog#1148 ilm-6).
|
||||
rustfs-ecstore = { workspace = true, features = ["test-util"] }
|
||||
|
||||
@@ -535,13 +535,8 @@ impl Operation for GetReplicationMetricsHandler {
|
||||
|
||||
let bucket_stats = cluster_replication_stats(bucket, app_context_from_req(&req)).await;
|
||||
|
||||
// Same minio-go `replication.Metrics` wire shape as
|
||||
// `?replication-metrics` — the internal snake_case stats are the peer
|
||||
// RPC wire format and must not leak here.
|
||||
let data = serde_json::to_vec(&crate::admin::replication_metrics_wire::MetricsWire::from(
|
||||
&bucket_stats.replication_stats,
|
||||
))
|
||||
.map_err(|_| S3Error::with_message(S3ErrorCode::InternalError, "serialize failed"))?;
|
||||
let data = serde_json::to_vec(&bucket_stats.replication_stats)
|
||||
.map_err(|_| S3Error::with_message(S3ErrorCode::InternalError, "serialize failed"))?;
|
||||
let mut headers = HeaderMap::new();
|
||||
headers.insert(CONTENT_TYPE, HeaderValue::from_static("application/json"));
|
||||
Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers))
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -29,6 +29,6 @@ impl Operation for RestLoadCredentialsHandler {
|
||||
let issuer = IamTableCredentialIssuer::from_request(&req)?;
|
||||
let response =
|
||||
load_credentials_response(&store, &warehouse, &namespace, &table, &issuer, Some(&principal.credentials)).await?;
|
||||
build_sensitive_json_response(StatusCode::OK, &response)
|
||||
build_json_response(StatusCode::OK, &response)
|
||||
}
|
||||
}
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -125,9 +125,7 @@ impl Operation for RestLoadTableHandler {
|
||||
ensure_table_bucket_enabled_from_extensions(&req.extensions, &warehouse).await?;
|
||||
let metadata_backend = table_catalog_backend_from_extensions(&req.extensions)?;
|
||||
let store = table_catalog_store_from_backend(metadata_backend.clone())?;
|
||||
let snapshot_selection = rest_table_snapshot_selection_from_query(&req.uri)?;
|
||||
let mut response = load_table_response(&store, &metadata_backend, &warehouse, &namespace, &table).await?;
|
||||
apply_rest_table_snapshot_selection(&mut response.metadata, snapshot_selection);
|
||||
let response = load_table_response(&store, &metadata_backend, &warehouse, &namespace, &table).await?;
|
||||
build_json_response(StatusCode::OK, &response)
|
||||
}
|
||||
}
|
||||
@@ -160,7 +158,7 @@ impl Operation for RestCommitTableHandler {
|
||||
let principal = authorize_table_catalog_resource_request(&req, &resource, AdminAction::CommitTableAction).await?;
|
||||
install_table_catalog_s3_request_info(&mut req, &principal)?;
|
||||
ensure_table_bucket_enabled_from_extensions(&req.extensions, &warehouse).await?;
|
||||
let request = read_rest_commit_table_request(std::mem::take(&mut req.input)).await?;
|
||||
let request = read_json_body::<RestCommitTableRequest>(std::mem::take(&mut req.input)).await?;
|
||||
let metadata_backend = table_catalog_backend_from_extensions(&req.extensions)?;
|
||||
let store = table_catalog_store_from_backend(metadata_backend.clone())?;
|
||||
let commit_backend = TableCommitObjectBackend::for_request(metadata_backend, req);
|
||||
@@ -180,14 +178,6 @@ impl Operation for RestDropTableHandler {
|
||||
let table = table_name_from_params(¶ms)?;
|
||||
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
|
||||
authorize_table_catalog_resource_request(&req, &resource, AdminAction::DeleteTableAction).await?;
|
||||
let purge_requested = rest_purge_requested_from_query(&req.uri)?;
|
||||
if purge_requested {
|
||||
return Err(iceberg_rest_error(
|
||||
ICEBERG_ERROR_UNSUPPORTED_OPERATION,
|
||||
StatusCode::NOT_ACCEPTABLE,
|
||||
"purgeRequested=true is not supported",
|
||||
));
|
||||
}
|
||||
ensure_table_bucket_enabled_from_extensions(&req.extensions, &warehouse).await?;
|
||||
let store = table_catalog_store_from_extensions(&req.extensions)?;
|
||||
drop_table_in_store(&store, &warehouse, &namespace, &table).await?;
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -43,9 +43,8 @@ impl Operation for RestCreateViewHandler {
|
||||
let metadata_backend = table_catalog_backend_from_extensions(&req.extensions)?;
|
||||
let store = table_catalog_store_from_backend(metadata_backend.clone())?;
|
||||
let table_bucket_enabled = table_bucket_enabled_from_extensions(&req.extensions, &warehouse).await?;
|
||||
let publication_backend = TableCommitObjectBackend::preauthorized(metadata_backend);
|
||||
let response =
|
||||
create_view_response(&store, &publication_backend, &warehouse, &namespace, request, table_bucket_enabled).await?;
|
||||
create_view_response(&store, &metadata_backend, &warehouse, &namespace, request, table_bucket_enabled).await?;
|
||||
build_json_response(StatusCode::OK, &response)
|
||||
}
|
||||
}
|
||||
@@ -88,20 +87,17 @@ pub struct RestReplaceViewHandler {}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl Operation for RestReplaceViewHandler {
|
||||
async fn call(&self, mut req: S3Request<Body>, params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
|
||||
async fn call(&self, req: S3Request<Body>, params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
|
||||
let warehouse = warehouse_from_params(¶ms)?;
|
||||
let namespace = namespace_from_params(¶ms)?;
|
||||
let view = view_name_from_params(¶ms)?;
|
||||
let resource = TableCatalogResource::view(&warehouse, &namespace, &view);
|
||||
let principal = authorize_table_catalog_resource_request(&req, &resource, AdminAction::CommitTableAction).await?;
|
||||
install_table_catalog_s3_request_info(&mut req, &principal)?;
|
||||
authorize_table_catalog_resource_request(&req, &resource, AdminAction::CommitTableAction).await?;
|
||||
ensure_table_bucket_enabled_from_extensions(&req.extensions, &warehouse).await?;
|
||||
let request = read_rest_commit_view_request(std::mem::take(&mut req.input)).await?;
|
||||
let request = read_json_body::<RestCommitViewRequest>(req.input).await?;
|
||||
let metadata_backend = table_catalog_backend_from_extensions(&req.extensions)?;
|
||||
let store = table_catalog_store_from_backend(metadata_backend.clone())?;
|
||||
let commit_backend = TableCommitObjectBackend::for_request(metadata_backend, req);
|
||||
let result = replace_view_response(&store, &commit_backend, &warehouse, &namespace, &view, request).await;
|
||||
let response = commit_backend.finish(result).await?;
|
||||
let response = replace_view_response(&store, &metadata_backend, &warehouse, &namespace, &view, request).await?;
|
||||
build_json_response(StatusCode::OK, &response)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -17,7 +17,6 @@ mod auth;
|
||||
pub mod console;
|
||||
pub mod handlers;
|
||||
mod plugin_contract;
|
||||
pub(crate) mod replication_metrics_wire;
|
||||
// Contract inventory is validated by tests before later runtime integration.
|
||||
#[allow(dead_code)]
|
||||
pub(crate) mod route_policy;
|
||||
|
||||
@@ -1,673 +0,0 @@
|
||||
// Copyright 2024 RustFS Team
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
//! Serialize-only wire projections of the internal replication statistics
|
||||
//! onto the minio-go `replication.Metrics` / `replication.MetricsV2` json
|
||||
//! shapes consumed by `mc replicate status` (`?replication-metrics[=2]` and
|
||||
//! the admin `replicationmetrics` endpoint).
|
||||
//!
|
||||
//! Red line: the internal `BucketStats` family in
|
||||
//! `crates/replication/src/stats.rs` is ALSO the intra-cluster peer-RPC wire
|
||||
//! format — `node_service.rs` encodes it with `rmp_serde::to_vec_named`, so
|
||||
//! its Rust field names travel between nodes as msgpack map keys. Renaming
|
||||
//! those serde names would break mixed-version clusters mid rolling upgrade.
|
||||
//! All madmin/minio-go interop therefore happens in these DTOs; never add
|
||||
//! `#[serde(rename)]` to the internal structs instead.
|
||||
//!
|
||||
//! Field names below are the exact json tags of minio-go
|
||||
//! `pkg/replication/replication.go` (v7.0.91). Keys minio-go does not know
|
||||
//! are RustFS extensions; Go decoders ignore unknown keys. `max`/`peak` are
|
||||
//! both emitted for the queue peak because the MinIO server writes `max`
|
||||
//! while minio-go reads `peak` (an upstream drift); emitting both keeps every
|
||||
//! decoder working.
|
||||
|
||||
use serde::Serialize;
|
||||
use std::collections::HashMap;
|
||||
use std::time::Duration;
|
||||
|
||||
use crate::admin::storage_api::replication::{
|
||||
BucketReplicationStat as InternalReplicationStat, BucketReplicationStats as InternalReplicationStats, BucketStats,
|
||||
InQueueMetric as InternalInQueueMetric, XferStats as InternalXferStats,
|
||||
};
|
||||
|
||||
/// minio-go `replication.RStat`.
|
||||
#[derive(Debug, Default, Clone, Copy, Serialize)]
|
||||
pub(crate) struct RStatWire {
|
||||
#[serde(rename = "count")]
|
||||
pub count: f64,
|
||||
#[serde(rename = "bytes")]
|
||||
pub bytes: i64,
|
||||
}
|
||||
|
||||
/// minio-go `replication.TimedErrStats`.
|
||||
#[derive(Debug, Default, Clone, Copy, Serialize)]
|
||||
pub(crate) struct TimedErrStatsWire {
|
||||
#[serde(rename = "lastMinute")]
|
||||
pub last_minute: RStatWire,
|
||||
#[serde(rename = "lastHour")]
|
||||
pub last_hour: RStatWire,
|
||||
#[serde(rename = "totals")]
|
||||
pub totals: RStatWire,
|
||||
}
|
||||
|
||||
impl TimedErrStatsWire {
|
||||
fn add(self, other: TimedErrStatsWire) -> TimedErrStatsWire {
|
||||
fn add(a: RStatWire, b: RStatWire) -> RStatWire {
|
||||
RStatWire {
|
||||
count: a.count + b.count,
|
||||
bytes: a.bytes.saturating_add(b.bytes),
|
||||
}
|
||||
}
|
||||
TimedErrStatsWire {
|
||||
last_minute: add(self.last_minute, other.last_minute),
|
||||
last_hour: add(self.last_hour, other.last_hour),
|
||||
totals: add(self.totals, other.totals),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// minio-go `replication.QStat`.
|
||||
#[derive(Debug, Default, Clone, Copy, Serialize)]
|
||||
pub(crate) struct QStatWire {
|
||||
#[serde(rename = "count")]
|
||||
pub count: f64,
|
||||
#[serde(rename = "bytes")]
|
||||
pub bytes: f64,
|
||||
}
|
||||
|
||||
/// minio-go `replication.InQueueMetric`, with the queue peak emitted under
|
||||
/// both `peak` (minio-go tag) and `max` (MinIO server tag).
|
||||
#[derive(Debug, Default, Clone, Copy, Serialize)]
|
||||
pub(crate) struct InQueueMetricWire {
|
||||
#[serde(rename = "curr")]
|
||||
pub curr: QStatWire,
|
||||
#[serde(rename = "avg")]
|
||||
pub avg: QStatWire,
|
||||
#[serde(rename = "max")]
|
||||
pub max: QStatWire,
|
||||
#[serde(rename = "peak")]
|
||||
pub peak: QStatWire,
|
||||
}
|
||||
|
||||
impl From<&InternalInQueueMetric> for InQueueMetricWire {
|
||||
fn from(metric: &InternalInQueueMetric) -> Self {
|
||||
fn qstat(bytes: i64, count: i64) -> QStatWire {
|
||||
QStatWire {
|
||||
count: count as f64,
|
||||
bytes: bytes as f64,
|
||||
}
|
||||
}
|
||||
let peak = qstat(metric.max.bytes, metric.max.count);
|
||||
InQueueMetricWire {
|
||||
curr: qstat(metric.curr.bytes, metric.curr.count),
|
||||
avg: qstat(metric.avg.bytes, metric.avg.count),
|
||||
max: peak,
|
||||
peak,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// minio-go `replication.XferStats`.
|
||||
#[derive(Debug, Default, Clone, Copy, Serialize)]
|
||||
pub(crate) struct XferStatsWire {
|
||||
#[serde(rename = "avgRate")]
|
||||
pub avg_rate: f64,
|
||||
#[serde(rename = "peakRate")]
|
||||
pub peak_rate: f64,
|
||||
#[serde(rename = "currRate")]
|
||||
pub curr_rate: f64,
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
struct XferStatsAverage {
|
||||
sum: XferStatsWire,
|
||||
active: u32,
|
||||
}
|
||||
|
||||
impl XferStatsAverage {
|
||||
fn add_active(&mut self, stats: XferStatsWire) {
|
||||
if stats.peak_rate <= 0.0 {
|
||||
return;
|
||||
}
|
||||
self.add_raw(stats);
|
||||
self.active += 1;
|
||||
}
|
||||
|
||||
fn add_raw(&mut self, stats: XferStatsWire) {
|
||||
self.sum.avg_rate += stats.avg_rate;
|
||||
self.sum.curr_rate += stats.curr_rate;
|
||||
self.sum.peak_rate = self.sum.peak_rate.max(stats.peak_rate);
|
||||
}
|
||||
|
||||
fn finish(self) -> XferStatsWire {
|
||||
let active = self.active;
|
||||
self.finish_with_divisor(active)
|
||||
}
|
||||
|
||||
fn finish_with_divisor(self, divisor: u32) -> XferStatsWire {
|
||||
if divisor == 0 {
|
||||
return self.sum;
|
||||
}
|
||||
XferStatsWire {
|
||||
avg_rate: self.sum.avg_rate / f64::from(divisor),
|
||||
peak_rate: self.sum.peak_rate,
|
||||
curr_rate: self.sum.curr_rate / f64::from(divisor),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<&InternalXferStats> for XferStatsWire {
|
||||
fn from(stats: &InternalXferStats) -> Self {
|
||||
XferStatsWire {
|
||||
avg_rate: stats.avg,
|
||||
peak_rate: stats.peak,
|
||||
curr_rate: stats.curr,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// minio-go `replication.WorkerStat`. RustFS does not track per-bucket worker
|
||||
/// occupancy yet, so this always reports zeros.
|
||||
#[derive(Debug, Default, Clone, Copy, Serialize)]
|
||||
pub(crate) struct WorkerStatWire {
|
||||
#[serde(rename = "curr")]
|
||||
pub curr: i32,
|
||||
#[serde(rename = "avg")]
|
||||
pub avg: f32,
|
||||
#[serde(rename = "max")]
|
||||
pub max: i32,
|
||||
}
|
||||
|
||||
/// minio-go `replication.ReplMRFStats`. RustFS does not track the 5-minute /
|
||||
/// dropped MRF windows, so this always reports zeros; the durable backlog is
|
||||
/// enumerable via `/v3/replication/mrf` instead.
|
||||
#[derive(Debug, Default, Clone, Copy, Serialize)]
|
||||
pub(crate) struct ReplMrfStatsWire {
|
||||
#[serde(rename = "failedCount_last5min")]
|
||||
pub last_failed_count: u64,
|
||||
#[serde(rename = "droppedCount_since_uptime")]
|
||||
pub total_dropped_count: u64,
|
||||
#[serde(rename = "droppedBytes_since_uptime")]
|
||||
pub total_dropped_bytes: u64,
|
||||
}
|
||||
|
||||
/// minio-go `replication.CounterSummary`.
|
||||
#[derive(Debug, Default, Clone, Copy, Serialize)]
|
||||
pub(crate) struct CounterSummaryWire {
|
||||
#[serde(rename = "last1hr")]
|
||||
pub last1hr: u64,
|
||||
#[serde(rename = "last1m")]
|
||||
pub last1m: u64,
|
||||
#[serde(rename = "total")]
|
||||
pub total: u64,
|
||||
}
|
||||
|
||||
/// minio-go `replication.TargetMetrics` (one remote target / ARN).
|
||||
#[derive(Debug, Default, Serialize)]
|
||||
pub(crate) struct TargetMetricsWire {
|
||||
#[serde(rename = "replicationCount")]
|
||||
pub replicated_count: i64,
|
||||
#[serde(rename = "completedReplicationSize")]
|
||||
pub replicated_size: i64,
|
||||
/// Bandwidth limit for this target. The tag says "bits" but both MinIO
|
||||
/// and minio-go treat the value as bytes/sec; keep bytes/sec.
|
||||
#[serde(rename = "limitInBits")]
|
||||
pub bandwidth_limit_bytes_per_sec: i64,
|
||||
#[serde(rename = "currentBandwidth")]
|
||||
pub current_bandwidth_bytes_per_sec: f64,
|
||||
#[serde(rename = "failed")]
|
||||
pub failed: TimedErrStatsWire,
|
||||
#[serde(rename = "failedReplicationSize")]
|
||||
pub failed_size: i64,
|
||||
#[serde(rename = "failedReplicationCount")]
|
||||
pub failed_count: i64,
|
||||
}
|
||||
|
||||
fn target_timed_err_stats(stat: &InternalReplicationStat) -> TimedErrStatsWire {
|
||||
// Cluster aggregation merges FailStats without the process-local samples,
|
||||
// so the serializable window snapshots (refreshed at each node's
|
||||
// collection point, summed by merge) are authoritative here; the live
|
||||
// samples only ever agree with or lag them, so take the larger.
|
||||
let sampled_minute = stat.fail_stats.recent_since(Duration::from_secs(60));
|
||||
let sampled_hour = stat.fail_stats.recent_since(Duration::from_secs(3600));
|
||||
let window = |sampled_count: i64, sampled_size: i64, snapshot_count: i64, snapshot_size: i64| RStatWire {
|
||||
count: sampled_count.max(snapshot_count) as f64,
|
||||
bytes: sampled_size.max(snapshot_size),
|
||||
};
|
||||
TimedErrStatsWire {
|
||||
last_minute: window(
|
||||
sampled_minute.count,
|
||||
sampled_minute.size,
|
||||
stat.fail_stats.last_minute.count,
|
||||
stat.fail_stats.last_minute.size,
|
||||
),
|
||||
last_hour: window(
|
||||
sampled_hour.count,
|
||||
sampled_hour.size,
|
||||
stat.fail_stats.last_hour.count,
|
||||
stat.fail_stats.last_hour.size,
|
||||
),
|
||||
totals: RStatWire {
|
||||
count: stat.failed.count as f64,
|
||||
bytes: stat.failed.size,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
impl From<&InternalReplicationStat> for TargetMetricsWire {
|
||||
fn from(stat: &InternalReplicationStat) -> Self {
|
||||
TargetMetricsWire {
|
||||
replicated_count: stat.replicated_count,
|
||||
replicated_size: stat.replicated_size,
|
||||
bandwidth_limit_bytes_per_sec: stat.bandwidth_limit_bytes_per_sec,
|
||||
current_bandwidth_bytes_per_sec: stat.current_bandwidth_bytes_per_sec,
|
||||
failed: target_timed_err_stats(stat),
|
||||
failed_size: stat.failed.size,
|
||||
failed_count: stat.failed.count,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// minio-go `replication.Metrics` — the `currStats` member of `MetricsV2` and
|
||||
/// the whole v1 response body. The trailing snake_case fields are RustFS
|
||||
/// source-health extension keys (ignored by Go decoders) carried over from
|
||||
/// the previous response shape.
|
||||
#[derive(Debug, Default, Serialize)]
|
||||
pub(crate) struct MetricsWire {
|
||||
#[serde(rename = "Stats")]
|
||||
pub stats: HashMap<String, TargetMetricsWire>,
|
||||
#[serde(rename = "completedReplicationSize")]
|
||||
pub replicated_size: i64,
|
||||
#[serde(rename = "replicaSize")]
|
||||
pub replica_size: i64,
|
||||
#[serde(rename = "replicaCount")]
|
||||
pub replica_count: i64,
|
||||
#[serde(rename = "replicationCount")]
|
||||
pub replicated_count: i64,
|
||||
#[serde(rename = "failed")]
|
||||
pub failed: TimedErrStatsWire,
|
||||
#[serde(rename = "queued")]
|
||||
pub queued: InQueueMetricWire,
|
||||
// RustFS extension keys (source health of the aggregation).
|
||||
pub provider_available: bool,
|
||||
pub cluster_complete: bool,
|
||||
pub observed_node_count: u32,
|
||||
pub expected_node_count: u32,
|
||||
}
|
||||
|
||||
impl From<&InternalReplicationStats> for MetricsWire {
|
||||
fn from(stats: &InternalReplicationStats) -> Self {
|
||||
let mut failed = TimedErrStatsWire::default();
|
||||
let mut targets = HashMap::with_capacity(stats.stats.len());
|
||||
for (arn, stat) in &stats.stats {
|
||||
let target = TargetMetricsWire::from(stat);
|
||||
failed = failed.add(target.failed);
|
||||
targets.insert(arn.clone(), target);
|
||||
}
|
||||
MetricsWire {
|
||||
stats: targets,
|
||||
replicated_size: stats.replicated_size,
|
||||
replica_size: stats.replica_size,
|
||||
replica_count: stats.replica_count,
|
||||
replicated_count: stats.replicated_count,
|
||||
failed,
|
||||
queued: InQueueMetricWire::from(&stats.q_stat),
|
||||
provider_available: stats.provider_available,
|
||||
cluster_complete: stats.cluster_complete,
|
||||
observed_node_count: stats.observed_node_count,
|
||||
expected_node_count: stats.expected_node_count,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// minio-go `replication.ReplQNodeStats`.
|
||||
#[derive(Debug, Default, Serialize)]
|
||||
pub(crate) struct ReplQNodeStatsWire {
|
||||
#[serde(rename = "nodeName")]
|
||||
pub node_name: String,
|
||||
#[serde(rename = "uptime")]
|
||||
pub uptime: i64,
|
||||
#[serde(rename = "activeWorkers")]
|
||||
pub workers: WorkerStatWire,
|
||||
#[serde(rename = "transferSummary")]
|
||||
pub xfer_stats: XferSummaryWire,
|
||||
#[serde(rename = "tgtTransferStats")]
|
||||
pub tgt_xfer_stats: TargetXferSummaryWire,
|
||||
#[serde(rename = "queueStats")]
|
||||
pub q_stats: InQueueMetricWire,
|
||||
#[serde(rename = "mrfStats")]
|
||||
pub mrf_stats: ReplMrfStatsWire,
|
||||
#[serde(rename = "retries")]
|
||||
pub retries: CounterSummaryWire,
|
||||
#[serde(rename = "errors")]
|
||||
pub errors: CounterSummaryWire,
|
||||
}
|
||||
|
||||
/// minio-go `replication.ReplQueueStats`.
|
||||
#[derive(Debug, Default, Serialize)]
|
||||
pub(crate) struct ReplQueueStatsWire {
|
||||
#[serde(rename = "nodes")]
|
||||
pub nodes: Vec<ReplQNodeStatsWire>,
|
||||
}
|
||||
|
||||
/// minio-go `replication.MetricsV2` — the `?replication-metrics=2` body.
|
||||
#[derive(Debug, Default, Serialize)]
|
||||
pub(crate) struct MetricsV2Wire {
|
||||
#[serde(rename = "uptime")]
|
||||
pub uptime: i64,
|
||||
#[serde(rename = "currStats")]
|
||||
pub current_stats: MetricsWire,
|
||||
#[serde(rename = "queueStats")]
|
||||
pub queue_stats: ReplQueueStatsWire,
|
||||
#[serde(rename = "downtimeInfo")]
|
||||
pub downtime_info: HashMap<String, serde_json::Value>,
|
||||
}
|
||||
|
||||
/// `transferSummary` map keyed by minio-go `MetricName` (Large/Small/Total).
|
||||
type XferSummaryWire = HashMap<&'static str, XferStatsWire>;
|
||||
/// `tgtTransferStats` map keyed by target ARN.
|
||||
type TargetXferSummaryWire = HashMap<String, XferSummaryWire>;
|
||||
|
||||
fn transfer_summaries(stats: &InternalReplicationStats) -> (XferSummaryWire, TargetXferSummaryWire) {
|
||||
let mut per_target: TargetXferSummaryWire = HashMap::new();
|
||||
let mut large_summary = XferStatsAverage::default();
|
||||
let mut small_summary = XferStatsAverage::default();
|
||||
let mut total_summary = XferStatsAverage::default();
|
||||
let mut active_targets = 0;
|
||||
for (arn, stat) in &stats.stats {
|
||||
let large = XferStatsWire::from(&stat.xfer_rate_lrg);
|
||||
let small = XferStatsWire::from(&stat.xfer_rate_sml);
|
||||
let mut target_total = XferStatsAverage::default();
|
||||
target_total.add_active(large);
|
||||
target_total.add_active(small);
|
||||
let total = target_total.finish();
|
||||
per_target.insert(arn.clone(), HashMap::from([("Large", large), ("Small", small), ("Total", total)]));
|
||||
if large.peak_rate > 0.0 || small.peak_rate > 0.0 {
|
||||
active_targets += 1;
|
||||
large_summary.add_raw(large);
|
||||
small_summary.add_raw(small);
|
||||
total_summary.add_raw(large);
|
||||
total_summary.add_raw(small);
|
||||
}
|
||||
}
|
||||
let summary = HashMap::from([
|
||||
("Large", large_summary.finish_with_divisor(active_targets)),
|
||||
("Small", small_summary.finish_with_divisor(active_targets)),
|
||||
("Total", total_summary.finish_with_divisor(active_targets)),
|
||||
]);
|
||||
(summary, per_target)
|
||||
}
|
||||
|
||||
impl MetricsV2Wire {
|
||||
/// Project the aggregated internal stats onto the `MetricsV2` shape.
|
||||
///
|
||||
/// The aggregation path leaves `queue_stats.nodes` empty today, so a
|
||||
/// single node entry is synthesized from the bucket queue snapshot —
|
||||
/// `mc replicate status` derives its queue/worker panels from
|
||||
/// `queueStats.nodes` and treats an empty list as "no data".
|
||||
pub(crate) fn from_stats(bucket_stats: &BucketStats, node_name: &str) -> Self {
|
||||
let (xfer_stats, tgt_xfer_stats) = transfer_summaries(&bucket_stats.replication_stats);
|
||||
let mut nodes: Vec<ReplQNodeStatsWire> = bucket_stats
|
||||
.queue_stats
|
||||
.nodes
|
||||
.iter()
|
||||
.map(|node| ReplQNodeStatsWire {
|
||||
node_name: node_name.to_string(),
|
||||
uptime: bucket_stats.uptime,
|
||||
q_stats: InQueueMetricWire::from(&node.q_stats),
|
||||
..Default::default()
|
||||
})
|
||||
.collect();
|
||||
if nodes.is_empty() {
|
||||
nodes.push(ReplQNodeStatsWire {
|
||||
node_name: node_name.to_string(),
|
||||
uptime: bucket_stats.uptime,
|
||||
q_stats: InQueueMetricWire::from(&bucket_stats.replication_stats.q_stat),
|
||||
xfer_stats: xfer_stats.clone(),
|
||||
tgt_xfer_stats: tgt_xfer_stats.clone(),
|
||||
..Default::default()
|
||||
});
|
||||
} else {
|
||||
// Attach the transfer summaries to the first node; the internal
|
||||
// snapshot does not attribute transfer rates per node.
|
||||
if let Some(first) = nodes.first_mut() {
|
||||
first.xfer_stats = xfer_stats.clone();
|
||||
first.tgt_xfer_stats = tgt_xfer_stats.clone();
|
||||
}
|
||||
}
|
||||
|
||||
MetricsV2Wire {
|
||||
uptime: bucket_stats.uptime,
|
||||
current_stats: MetricsWire::from(&bucket_stats.replication_stats),
|
||||
queue_stats: ReplQueueStatsWire { nodes },
|
||||
downtime_info: HashMap::new(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn sample_bucket_stats() -> BucketStats {
|
||||
let mut stats = BucketStats {
|
||||
uptime: 42,
|
||||
..Default::default()
|
||||
};
|
||||
stats.replication_stats.replica_count = 2;
|
||||
stats.replication_stats.replica_size = 128;
|
||||
stats.replication_stats.replicated_count = 9;
|
||||
stats.replication_stats.replicated_size = 4096;
|
||||
let target = stats
|
||||
.replication_stats
|
||||
.stats
|
||||
.entry("arn:minio:replication::t:b".to_string())
|
||||
.or_default();
|
||||
target.replicated_count = 9;
|
||||
target.replicated_size = 4096;
|
||||
target.failed.count = 3;
|
||||
target.failed.size = 900;
|
||||
target.bandwidth_limit_bytes_per_sec = 1024;
|
||||
target.current_bandwidth_bytes_per_sec = 512.5;
|
||||
stats
|
||||
.replication_stats
|
||||
.q_stat
|
||||
.curr
|
||||
.now_count
|
||||
.store(4, std::sync::atomic::Ordering::Relaxed);
|
||||
stats
|
||||
.replication_stats
|
||||
.q_stat
|
||||
.curr
|
||||
.now_bytes
|
||||
.store(1200, std::sync::atomic::Ordering::Relaxed);
|
||||
stats.replication_stats.q_stat = stats.replication_stats.q_stat.snapshot();
|
||||
stats
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn metrics_wire_matches_minio_go_tags() {
|
||||
let stats = sample_bucket_stats();
|
||||
let json = serde_json::to_value(MetricsWire::from(&stats.replication_stats)).expect("v1 wire should serialize");
|
||||
|
||||
assert_eq!(json["replicaCount"], 2);
|
||||
assert_eq!(json["replicaSize"], 128);
|
||||
assert_eq!(json["replicationCount"], 9);
|
||||
assert_eq!(json["completedReplicationSize"], 4096);
|
||||
assert_eq!(json["queued"]["curr"]["count"], 4.0);
|
||||
assert_eq!(json["queued"]["curr"]["bytes"], 1200.0);
|
||||
let target = &json["Stats"]["arn:minio:replication::t:b"];
|
||||
assert_eq!(target["replicationCount"], 9);
|
||||
assert_eq!(target["completedReplicationSize"], 4096);
|
||||
assert_eq!(target["limitInBits"], 1024);
|
||||
assert_eq!(target["currentBandwidth"], 512.5);
|
||||
// failed is the madmin TimedErrStats envelope, not the internal
|
||||
// {count,size} pair.
|
||||
assert_eq!(target["failed"]["totals"]["count"], 3.0);
|
||||
assert_eq!(target["failed"]["totals"]["bytes"], 900);
|
||||
assert!(target["failed"].get("count").is_none());
|
||||
// Aggregate failed mirrors the per-target totals.
|
||||
assert_eq!(json["failed"]["totals"]["count"], 3.0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn metrics_v2_wire_synthesizes_queue_node() {
|
||||
let stats = sample_bucket_stats();
|
||||
let json = serde_json::to_value(MetricsV2Wire::from_stats(&stats, "node-1:9000")).expect("v2 wire should serialize");
|
||||
|
||||
assert_eq!(json["uptime"], 42);
|
||||
assert_eq!(json["currStats"]["replicaCount"], 2);
|
||||
let node = &json["queueStats"]["nodes"][0];
|
||||
assert_eq!(node["nodeName"], "node-1:9000");
|
||||
assert_eq!(node["uptime"], 42);
|
||||
assert_eq!(node["queueStats"]["curr"]["count"], 4.0);
|
||||
// The queue peak is emitted under both the minio-go tag (`peak`) and
|
||||
// the MinIO server tag (`max`).
|
||||
assert_eq!(node["queueStats"]["peak"], node["queueStats"]["max"]);
|
||||
assert!(node["activeWorkers"].get("curr").is_some());
|
||||
assert!(node["transferSummary"].get("Total").is_some());
|
||||
assert_eq!(json["downtimeInfo"], serde_json::json!({}));
|
||||
}
|
||||
|
||||
/// minio-go's transferSummary labels mean >= 128 MiB for Large; the
|
||||
/// producer must bin on the same boundary (MIN_LARGE_OBJ_SIZE, shared
|
||||
/// with the worker-pool split), or a 2 MiB replication shows under Large
|
||||
/// while Small stays zero.
|
||||
#[test]
|
||||
fn transfer_summary_bins_on_the_128_mib_boundary() {
|
||||
const MIB: i64 = 1024 * 1024;
|
||||
let mut stats = BucketStats::default();
|
||||
let stat = stats
|
||||
.replication_stats
|
||||
.stats
|
||||
.entry("arn:minio:replication::t:b".to_string())
|
||||
.or_default();
|
||||
stat.update_xfer_rate(2 * MIB, std::time::Duration::from_secs(1));
|
||||
stat.update_xfer_rate(127 * MIB, std::time::Duration::from_secs(1));
|
||||
stat.update_xfer_rate(128 * MIB, std::time::Duration::from_secs(1));
|
||||
|
||||
let json = serde_json::to_value(MetricsV2Wire::from_stats(&stats, "node-1")).expect("v2 wire should serialize");
|
||||
let summary = &json["queueStats"]["nodes"][0]["tgtTransferStats"]["arn:minio:replication::t:b"];
|
||||
let small_peak = summary["Small"]["peakRate"].as_f64().expect("Small peakRate");
|
||||
let large_peak = summary["Large"]["peakRate"].as_f64().expect("Large peakRate");
|
||||
assert!(
|
||||
(small_peak - (127 * MIB) as f64).abs() < 1.0,
|
||||
"2 MiB and 127 MiB transfers must bin as Small (peak {small_peak})"
|
||||
);
|
||||
assert!(
|
||||
(large_peak - (128 * MIB) as f64).abs() < 1.0,
|
||||
"exactly 128 MiB must bin as Large (peak {large_peak})"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn transfer_summaries_average_active_bins_and_targets() {
|
||||
let mut stats = BucketStats::default();
|
||||
let first = stats.replication_stats.stats.entry("target-a".to_string()).or_default();
|
||||
first.xfer_rate_sml.avg = 50.0;
|
||||
first.xfer_rate_sml.curr = 40.0;
|
||||
first.xfer_rate_sml.peak = 60.0;
|
||||
first.xfer_rate_lrg.avg = 100.0;
|
||||
first.xfer_rate_lrg.curr = 80.0;
|
||||
first.xfer_rate_lrg.peak = 120.0;
|
||||
|
||||
let second = stats.replication_stats.stats.entry("target-b".to_string()).or_default();
|
||||
second.xfer_rate_sml.avg = 30.0;
|
||||
second.xfer_rate_sml.curr = 20.0;
|
||||
second.xfer_rate_sml.peak = 40.0;
|
||||
|
||||
let json = serde_json::to_value(MetricsV2Wire::from_stats(&stats, "node-1")).expect("v2 wire should serialize");
|
||||
let node = &json["queueStats"]["nodes"][0];
|
||||
let target_a = &node["tgtTransferStats"]["target-a"]["Total"];
|
||||
assert_eq!(target_a["avgRate"], 75.0);
|
||||
assert_eq!(target_a["currRate"], 60.0);
|
||||
assert_eq!(target_a["peakRate"], 120.0);
|
||||
|
||||
let summary = &node["transferSummary"];
|
||||
assert_eq!(summary["Small"]["avgRate"], 40.0);
|
||||
assert_eq!(summary["Small"]["currRate"], 30.0);
|
||||
assert_eq!(summary["Large"]["avgRate"], 50.0);
|
||||
assert_eq!(summary["Total"]["avgRate"], 90.0);
|
||||
assert_eq!(summary["Total"]["currRate"], 70.0);
|
||||
assert_eq!(summary["Total"]["peakRate"], 120.0);
|
||||
}
|
||||
|
||||
/// Review regression: both metrics endpoints aggregate first, and the
|
||||
/// FailStats merge drops the process-local samples — the rolling windows
|
||||
/// must survive a peer-RPC round trip plus aggregation and still reach
|
||||
/// the wire body.
|
||||
#[test]
|
||||
fn failure_windows_survive_aggregation_before_serialization() {
|
||||
// Node A: live failure; the windows are stamped at the collection
|
||||
// point (get_latest_replication_stats calls refresh_windows before
|
||||
// the stats cross the wire), never on the failure hot path.
|
||||
let mut node_a = crate::admin::storage_api::replication::BucketReplicationStat::default();
|
||||
node_a.fail_stats.add_size(512, None::<&std::io::Error>);
|
||||
node_a.fail_stats.refresh_windows();
|
||||
node_a.failed = node_a.fail_stats.to_metric();
|
||||
|
||||
// Node A's stats cross the peer RPC wire: the samples are dropped,
|
||||
// the window snapshots travel.
|
||||
let encoded = rmp_serde::to_vec_named(&node_a).expect("stat should encode");
|
||||
let remote: crate::admin::storage_api::replication::BucketReplicationStat =
|
||||
rmp_serde::from_slice(&encoded).expect("stat should decode");
|
||||
|
||||
// Aggregation merges the remote stat with an empty local one.
|
||||
let merged_fail = remote.fail_stats.merge(&Default::default());
|
||||
let aggregated = crate::admin::storage_api::replication::BucketReplicationStat {
|
||||
failed: merged_fail.to_metric(),
|
||||
fail_stats: merged_fail,
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let mut stats = BucketStats::default();
|
||||
stats
|
||||
.replication_stats
|
||||
.stats
|
||||
.insert("arn:minio:replication::t:b".to_string(), aggregated);
|
||||
|
||||
let json = serde_json::to_value(MetricsWire::from(&stats.replication_stats)).expect("wire should serialize");
|
||||
let failed = &json["Stats"]["arn:minio:replication::t:b"]["failed"];
|
||||
assert_eq!(failed["totals"]["count"], 1.0);
|
||||
assert_eq!(
|
||||
failed["lastMinute"]["count"], 1.0,
|
||||
"the rolling minute window must survive RPC + aggregation"
|
||||
);
|
||||
assert_eq!(failed["lastMinute"]["bytes"], 512);
|
||||
assert_eq!(failed["lastHour"]["count"], 1.0);
|
||||
}
|
||||
|
||||
/// Pin the intra-cluster peer-RPC wire format of the internal stats: it
|
||||
/// is msgpack with the Rust field names as map keys
|
||||
/// (`rmp_serde::to_vec_named` in node_service.rs). If someone "fixes"
|
||||
/// the interop bug by renaming the internal serde fields instead of using
|
||||
/// these DTOs, this test fails and points them here.
|
||||
#[test]
|
||||
fn internal_bucket_stats_rpc_wire_stays_snake_case() {
|
||||
let stats = sample_bucket_stats();
|
||||
let encoded = rmp_serde::to_vec_named(&stats).expect("internal stats should encode");
|
||||
let value: serde_json::Value = rmp_serde::from_slice(&encoded).expect("named msgpack should decode generically");
|
||||
|
||||
assert!(
|
||||
value.get("replication_stats").is_some(),
|
||||
"peer RPC key replication_stats must not be renamed"
|
||||
);
|
||||
assert!(value["replication_stats"].get("q_stat").is_some());
|
||||
assert!(value.get("queue_stats").is_some());
|
||||
assert!(value.get("proxy_stats").is_some());
|
||||
|
||||
let decoded: BucketStats = rmp_serde::from_slice(&encoded).expect("round-trip through the peer RPC wire");
|
||||
assert_eq!(decoded.replication_stats.replica_count, 2);
|
||||
}
|
||||
}
|
||||
+13
-67
@@ -1548,8 +1548,7 @@ async fn build_replication_metrics_response(
|
||||
let bucket_stats = apply_replication_metrics_bandwidth_report(bucket_stats, collect_replication_metrics_bandwidth(bucket));
|
||||
let bucket_stats = apply_replication_metrics_runtime_fields(bucket_stats, route, replication_metrics_uptime_seconds());
|
||||
|
||||
let node_name = crate::runtime_sources::current_local_node_name().await.unwrap_or_default();
|
||||
let body = serialize_replication_metrics_body(&bucket_stats, route, &node_name)?;
|
||||
let body = serialize_replication_metrics_body(&bucket_stats, route)?;
|
||||
|
||||
let mut resp = S3Response::with_status(Body::from(body), StatusCode::OK);
|
||||
resp.headers
|
||||
@@ -1609,24 +1608,12 @@ fn apply_replication_metrics_runtime_fields(
|
||||
bucket_stats
|
||||
}
|
||||
|
||||
/// Serialize the metrics body in the minio-go wire shapes
|
||||
/// (`replication.Metrics` for v1, `replication.MetricsV2` for v2). The
|
||||
/// internal `BucketStats` serde names are the intra-cluster peer-RPC wire
|
||||
/// format and must never appear here — see
|
||||
/// `crate::admin::replication_metrics_wire`.
|
||||
fn serialize_replication_metrics_body(
|
||||
bucket_stats: &BucketStats,
|
||||
route: ReplicationExtRoute,
|
||||
node_name: &str,
|
||||
) -> S3Result<Vec<u8>> {
|
||||
use crate::admin::replication_metrics_wire::{MetricsV2Wire, MetricsWire};
|
||||
fn serialize_replication_metrics_body(bucket_stats: &BucketStats, route: ReplicationExtRoute) -> S3Result<Vec<u8>> {
|
||||
match route {
|
||||
ReplicationExtRoute::MetricsV1 => {
|
||||
serde_json::to_vec(&MetricsWire::from(&bucket_stats.replication_stats)).map_err(|e| s3_error!(InternalError, "{e}"))
|
||||
}
|
||||
ReplicationExtRoute::MetricsV2 => {
|
||||
serde_json::to_vec(&MetricsV2Wire::from_stats(bucket_stats, node_name)).map_err(|e| s3_error!(InternalError, "{e}"))
|
||||
serde_json::to_vec(&bucket_stats.replication_stats).map_err(|e| s3_error!(InternalError, "{e}"))
|
||||
}
|
||||
ReplicationExtRoute::MetricsV2 => serde_json::to_vec(bucket_stats).map_err(|e| s3_error!(InternalError, "{e}")),
|
||||
ReplicationExtRoute::Check | ReplicationExtRoute::ResetStart | ReplicationExtRoute::ResetStatus => {
|
||||
Err(s3_error!(InternalError, "invalid route for metrics response"))
|
||||
}
|
||||
@@ -4160,37 +4147,22 @@ mod tests {
|
||||
assert!(err.message().unwrap_or_default().contains("rule-stale"));
|
||||
}
|
||||
|
||||
/// The v1 body must decode into minio-go `replication.Metrics` (exact
|
||||
/// json tags); Go's decoder matches case-insensitively but does not
|
||||
/// ignore underscores, so the internal snake_case names read as all-zero.
|
||||
#[test]
|
||||
fn serialize_replication_metrics_body_v1_returns_minio_go_metrics_shape() {
|
||||
fn serialize_replication_metrics_body_v1_returns_replication_stats_only() {
|
||||
let mut stats = BucketStats {
|
||||
uptime: 99,
|
||||
..Default::default()
|
||||
};
|
||||
stats.replication_stats.replica_count = 7;
|
||||
stats.replication_stats.replicated_size = 2048;
|
||||
stats
|
||||
.replication_stats
|
||||
.stats
|
||||
.entry("arn:minio:replication::t:b".to_string())
|
||||
.or_default()
|
||||
.replicated_count = 5;
|
||||
stats.proxy_stats.put_total = 3;
|
||||
|
||||
let body = serialize_replication_metrics_body(&stats, ReplicationExtRoute::MetricsV1, "node-1:9000")
|
||||
.expect("metrics v1 body should serialize");
|
||||
let body =
|
||||
serialize_replication_metrics_body(&stats, ReplicationExtRoute::MetricsV1).expect("metrics v1 body should serialize");
|
||||
let payload: serde_json::Value = serde_json::from_slice(&body).expect("body should be json");
|
||||
|
||||
assert_eq!(payload["replicaCount"], 7);
|
||||
assert_eq!(payload["completedReplicationSize"], 2048);
|
||||
assert_eq!(payload["Stats"]["arn:minio:replication::t:b"]["replicationCount"], 5);
|
||||
assert_eq!(payload["replica_count"], 7);
|
||||
assert!(payload.get("uptime").is_none());
|
||||
assert!(payload.get("proxy_stats").is_none());
|
||||
// The internal snake_case names must not leak into the wire body.
|
||||
assert!(payload.get("replica_count").is_none());
|
||||
assert!(payload.get("q_stat").is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -4276,48 +4248,22 @@ mod tests {
|
||||
assert_eq!(target.current_bandwidth_bytes_per_sec, 3000.0);
|
||||
}
|
||||
|
||||
/// The v2 body must decode into minio-go `replication.MetricsV2`
|
||||
/// (`uptime`/`currStats`/`queueStats`); `mc replicate status` reads
|
||||
/// `currStats` and `queueStats.nodes` and silently shows zeros when the
|
||||
/// keys do not match.
|
||||
#[test]
|
||||
fn serialize_replication_metrics_body_v2_returns_minio_go_metrics_v2_shape() {
|
||||
fn serialize_replication_metrics_body_v2_returns_full_bucket_stats() {
|
||||
let mut stats = BucketStats {
|
||||
uptime: 99,
|
||||
..Default::default()
|
||||
};
|
||||
stats.replication_stats.replica_count = 7;
|
||||
stats
|
||||
.replication_stats
|
||||
.q_stat
|
||||
.curr
|
||||
.now_count
|
||||
.store(4, std::sync::atomic::Ordering::Relaxed);
|
||||
stats
|
||||
.replication_stats
|
||||
.q_stat
|
||||
.curr
|
||||
.now_bytes
|
||||
.store(1200, std::sync::atomic::Ordering::Relaxed);
|
||||
stats.replication_stats.q_stat = stats.replication_stats.q_stat.snapshot();
|
||||
stats.proxy_stats.put_total = 3;
|
||||
|
||||
let body = serialize_replication_metrics_body(&stats, ReplicationExtRoute::MetricsV2, "node-1:9000")
|
||||
.expect("metrics v2 body should serialize");
|
||||
let body =
|
||||
serialize_replication_metrics_body(&stats, ReplicationExtRoute::MetricsV2).expect("metrics v2 body should serialize");
|
||||
let payload: serde_json::Value = serde_json::from_slice(&body).expect("body should be json");
|
||||
|
||||
assert_eq!(payload["uptime"], 99);
|
||||
assert_eq!(payload["currStats"]["replicaCount"], 7);
|
||||
assert_eq!(payload["currStats"]["queued"]["curr"]["count"], 4.0);
|
||||
// The queue snapshot must surface at least one node: mc derives the
|
||||
// worker/queue panels from queueStats.nodes and treats an empty list
|
||||
// as "no data".
|
||||
assert_eq!(payload["queueStats"]["nodes"][0]["queueStats"]["curr"]["count"], 4.0);
|
||||
assert_eq!(payload["queueStats"]["nodes"][0]["uptime"], 99);
|
||||
// The internal snake_case names must not leak into the wire body.
|
||||
assert!(payload.get("replication_stats").is_none());
|
||||
assert!(payload.get("queue_stats").is_none());
|
||||
assert!(payload.get("proxy_stats").is_none());
|
||||
assert_eq!(payload["replication_stats"]["replica_count"], 7);
|
||||
assert_eq!(payload["proxy_stats"]["put_total"], 3);
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
@@ -321,34 +321,6 @@ pub(crate) mod metadata_sys {
|
||||
crate::storage::storage_api::acquire_bucket_metadata_transaction_lock(bucket).await
|
||||
}
|
||||
|
||||
pub(crate) async fn acquire_bucket_metadata_transaction_lock_for_incarnation(
|
||||
bucket: &str,
|
||||
expected_incarnation_id: uuid::Uuid,
|
||||
) -> Result<super::ecstore_bucket::metadata_sys::BucketMetadataMutationGuard> {
|
||||
super::ecstore_bucket::metadata_sys::acquire_bucket_metadata_transaction_lock_for_incarnation(
|
||||
bucket,
|
||||
expected_incarnation_id,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
pub(crate) async fn update_under_transaction_lock(
|
||||
guard: &super::ecstore_bucket::metadata_sys::BucketMetadataMutationGuard,
|
||||
bucket: &str,
|
||||
config_file: &str,
|
||||
data: Vec<u8>,
|
||||
) -> Result<OffsetDateTime> {
|
||||
super::ecstore_bucket::metadata_sys::update_under_transaction_lock(guard, bucket, config_file, data).await
|
||||
}
|
||||
|
||||
pub(crate) async fn delete_under_transaction_lock(
|
||||
guard: &super::ecstore_bucket::metadata_sys::BucketMetadataMutationGuard,
|
||||
bucket: &str,
|
||||
config_file: &str,
|
||||
) -> Result<OffsetDateTime> {
|
||||
super::ecstore_bucket::metadata_sys::delete_under_transaction_lock(guard, bucket, config_file).await
|
||||
}
|
||||
|
||||
pub(crate) async fn update_bucket_targets_under_transaction_lock(
|
||||
guard: &super::ecstore_bucket::metadata_sys::BucketMetadataMutationGuard,
|
||||
bucket: &str,
|
||||
@@ -445,10 +417,6 @@ pub(crate) mod replication {
|
||||
};
|
||||
pub(crate) type BucketReplicationResyncStatus = super::ecstore_bucket::replication::BucketReplicationResyncStatus;
|
||||
pub(crate) type BucketStats = super::ecstore_bucket::replication::BucketStats;
|
||||
pub(crate) type BucketReplicationStats = super::ecstore_bucket::replication::BucketReplicationStats;
|
||||
pub(crate) type BucketReplicationStat = super::ecstore_bucket::replication::BucketReplicationStat;
|
||||
pub(crate) type InQueueMetric = super::ecstore_bucket::replication::InQueueMetric;
|
||||
pub(crate) type XferStats = super::ecstore_bucket::replication::XferStats;
|
||||
pub(crate) type ReplicationStatusType = super::ecstore_bucket::replication::ReplicationStatusType;
|
||||
pub(crate) type ResyncOpts = super::ecstore_bucket::replication::ResyncOpts;
|
||||
pub(crate) type ResyncStatusType = super::ecstore_bucket::replication::ResyncStatusType;
|
||||
|
||||
@@ -1158,19 +1158,6 @@ fn lifecycle_has_expiry_rules(config: &BucketLifecycleConfiguration) -> bool {
|
||||
})
|
||||
}
|
||||
|
||||
/// Status-independent presence of the expiry subset that site replication
|
||||
/// propagates (`replicateILMExpiry`): expiration / noncurrent-version
|
||||
/// expiration only. Distinct from [`lifecycle_has_expiry_rules`], which
|
||||
/// filters on ENABLED for scanner scheduling — editing a Disabled expiry rule
|
||||
/// must still advance the replication axis. Del-marker expiration and
|
||||
/// abort-multipart are site-local and never travel.
|
||||
fn lifecycle_rules_have_expiry(config: &BucketLifecycleConfiguration) -> bool {
|
||||
config
|
||||
.rules
|
||||
.iter()
|
||||
.any(|rule| rule.expiration.is_some() || rule.noncurrent_version_expiration.is_some())
|
||||
}
|
||||
|
||||
fn lifecycle_has_abort_multipart_rules(config: &BucketLifecycleConfiguration) -> bool {
|
||||
config.rules.iter().any(|rule| {
|
||||
rule.status == ExpirationStatus::from_static(ExpirationStatus::ENABLED)
|
||||
@@ -2199,24 +2186,7 @@ impl DefaultBucketUsecase {
|
||||
return Err(s3_error!(InvalidArgument, "{err}"));
|
||||
}
|
||||
|
||||
// Stamp the expiry axis only when the expiry subset can have changed
|
||||
// (MinIO: HasExpiry() || expiryRuleRemoved). Site-replication peers
|
||||
// judge lc-config staleness on this axis; a transition-only edit that
|
||||
// advanced it would let this site's stale expiry subset shadow — and
|
||||
// roll back — a newer peer expiry edit fleet-wide.
|
||||
let previous_expiry_updated_at = match metadata_sys::get_lifecycle_config(&bucket).await {
|
||||
Ok((previous, _)) => {
|
||||
if lifecycle_rules_have_expiry(&input_cfg) || lifecycle_rules_have_expiry(&previous) {
|
||||
Some(Timestamp::from(time::OffsetDateTime::now_utc()))
|
||||
} else {
|
||||
previous.expiry_updated_at
|
||||
}
|
||||
}
|
||||
// No previous config (or unreadable): stamping is the
|
||||
// conservative pre-existing behavior.
|
||||
Err(_) => lifecycle_rules_have_expiry(&input_cfg).then(|| Timestamp::from(time::OffsetDateTime::now_utc())),
|
||||
};
|
||||
input_cfg.expiry_updated_at = previous_expiry_updated_at;
|
||||
input_cfg.expiry_updated_at = Some(Timestamp::from(time::OffsetDateTime::now_utc()));
|
||||
let data = serialize_config(&input_cfg)?;
|
||||
update_bucket_config_for_incarnation(&bucket, BUCKET_LIFECYCLE_CONFIG, data, expected_incarnation_id)
|
||||
.await
|
||||
@@ -2227,14 +2197,7 @@ impl DefaultBucketUsecase {
|
||||
let mut item = sr_bucket_meta_item(bucket.clone(), "lc-config");
|
||||
item.expiry_lc_config =
|
||||
Some(serialize_config(&input_cfg).and_then(|bytes| String::from_utf8(bytes).map_err(to_internal_error))?);
|
||||
// The item travels with the expiry axis, not the wall clock: a site
|
||||
// whose expiry knowledge is old (or absent — UNIX_EPOCH) must not
|
||||
// out-rank newer peer expiry state at the receivers.
|
||||
item.expiry_updated_at = input_cfg
|
||||
.expiry_updated_at
|
||||
.clone()
|
||||
.map(time::OffsetDateTime::from)
|
||||
.or(Some(time::OffsetDateTime::UNIX_EPOCH));
|
||||
item.expiry_updated_at = item.updated_at;
|
||||
if let Err(err) = site_replication_bucket_meta_hook(item).await {
|
||||
warn!(bucket = %bucket, error = ?err, "site replication bucket lifecycle hook failed");
|
||||
}
|
||||
|
||||
@@ -1478,7 +1478,7 @@ async fn retain_table_data_plane_publication_guard<T>(
|
||||
.map_err(|err| s3_error!(InternalError, "failed to acquire table publication guard: {}", err))?;
|
||||
let mut state = retained.state.lock();
|
||||
state.keys.insert(key);
|
||||
state.guards.push(Box::new(guard));
|
||||
state.guards.push(guard);
|
||||
drop(state);
|
||||
req.extensions.insert(retained);
|
||||
Ok(())
|
||||
|
||||
@@ -16,8 +16,6 @@ use std::io::Read;
|
||||
|
||||
use super::super::*;
|
||||
|
||||
const AVRO_ZSTANDARD_MAX_WINDOW_LOG: u32 = 27;
|
||||
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
pub(crate) struct ManifestDataFileReference {
|
||||
pub location: String,
|
||||
@@ -68,7 +66,6 @@ pub(crate) struct DecodedManifestList {
|
||||
pub(crate) struct DecodedManifest {
|
||||
pub references: Vec<ManifestDataFileReference>,
|
||||
pub decoded_size: usize,
|
||||
pub partition_spec_id: Option<i32>,
|
||||
}
|
||||
|
||||
pub(crate) fn manifest_paths_from_manifest_list_avro(data: &[u8]) -> TableCatalogStoreResult<Vec<String>> {
|
||||
@@ -95,25 +92,6 @@ pub(crate) fn decode_manifest_list_avro(data: &[u8]) -> TableCatalogStoreResult<
|
||||
.map_err(|err| TableCatalogStoreError::Invalid(format!("failed to read manifest list Avro: {err}")))?;
|
||||
let format_version =
|
||||
avro_record_format_version(reader.writer_schema(), &["sequence_number", "min_sequence_number"], "manifest list")?;
|
||||
if format_version == 2 {
|
||||
let apache_avro::Schema::Record(record) = reader.writer_schema() else {
|
||||
return Err(TableCatalogStoreError::Invalid("manifest list Avro schema must be a record".to_string()));
|
||||
};
|
||||
for field in [
|
||||
"added_files_count",
|
||||
"existing_files_count",
|
||||
"deleted_files_count",
|
||||
"added_rows_count",
|
||||
"existing_rows_count",
|
||||
"deleted_rows_count",
|
||||
] {
|
||||
if !record.lookup.contains_key(field) {
|
||||
return Err(TableCatalogStoreError::Invalid(format!(
|
||||
"Iceberg v2 manifest list Avro schema is missing {field}"
|
||||
)));
|
||||
}
|
||||
}
|
||||
}
|
||||
let mut manifest_paths = Vec::new();
|
||||
for value in reader {
|
||||
if manifest_paths.len() >= TABLE_MANIFEST_AVRO_MAX_RECORDS {
|
||||
@@ -137,12 +115,24 @@ pub(crate) fn decode_manifest_list_avro(data: &[u8]) -> TableCatalogStoreResult<
|
||||
sequence_number: avro_record_field(&value, "sequence_number").and_then(avro_i64_value),
|
||||
min_sequence_number: avro_record_field(&value, "min_sequence_number").and_then(avro_i64_value),
|
||||
added_snapshot_id: avro_record_field(&value, "added_snapshot_id").and_then(avro_i64_value),
|
||||
added_files_count: avro_nullable_non_negative_i32(&value, "added_files_count", "manifest list")?,
|
||||
existing_files_count: avro_nullable_non_negative_i32(&value, "existing_files_count", "manifest list")?,
|
||||
deleted_files_count: avro_nullable_non_negative_i32(&value, "deleted_files_count", "manifest list")?,
|
||||
added_rows_count: avro_nullable_non_negative_i64(&value, "added_rows_count", "manifest list")?,
|
||||
existing_rows_count: avro_nullable_non_negative_i64(&value, "existing_rows_count", "manifest list")?,
|
||||
deleted_rows_count: avro_nullable_non_negative_i64(&value, "deleted_rows_count", "manifest list")?,
|
||||
added_files_count: avro_record_field(&value, "added_files_count")
|
||||
.and_then(avro_i32_value)
|
||||
.and_then(|value| u64::try_from(value).ok()),
|
||||
existing_files_count: avro_record_field(&value, "existing_files_count")
|
||||
.and_then(avro_i32_value)
|
||||
.and_then(|value| u64::try_from(value).ok()),
|
||||
deleted_files_count: avro_record_field(&value, "deleted_files_count")
|
||||
.and_then(avro_i32_value)
|
||||
.and_then(|value| u64::try_from(value).ok()),
|
||||
added_rows_count: avro_record_field(&value, "added_rows_count")
|
||||
.and_then(avro_i64_value)
|
||||
.and_then(|value| u64::try_from(value).ok()),
|
||||
existing_rows_count: avro_record_field(&value, "existing_rows_count")
|
||||
.and_then(avro_i64_value)
|
||||
.and_then(|value| u64::try_from(value).ok()),
|
||||
deleted_rows_count: avro_record_field(&value, "deleted_rows_count")
|
||||
.and_then(avro_i64_value)
|
||||
.and_then(|value| u64::try_from(value).ok()),
|
||||
});
|
||||
}
|
||||
Ok(DecodedManifestList {
|
||||
@@ -181,17 +171,6 @@ pub(crate) fn decode_manifest_avro(data: &[u8]) -> TableCatalogStoreResult<Decod
|
||||
.map_err(|err| TableCatalogStoreError::Invalid(format!("failed to read manifest Avro: {err}")))?;
|
||||
let format_version =
|
||||
avro_record_format_version(reader.writer_schema(), &["sequence_number", "file_sequence_number"], "manifest")?;
|
||||
let partition_spec_id = reader
|
||||
.user_metadata()
|
||||
.get("partition-spec-id")
|
||||
.map(|value| {
|
||||
std::str::from_utf8(value)
|
||||
.ok()
|
||||
.and_then(|value| value.parse::<i32>().ok())
|
||||
.filter(|value| *value >= 0)
|
||||
.ok_or_else(|| TableCatalogStoreError::Invalid("manifest partition-spec-id metadata is invalid".to_string()))
|
||||
})
|
||||
.transpose()?;
|
||||
let mut files = Vec::new();
|
||||
for value in reader {
|
||||
if files.len() >= TABLE_MANIFEST_AVRO_MAX_RECORDS {
|
||||
@@ -223,9 +202,6 @@ pub(crate) fn decode_manifest_avro(data: &[u8]) -> TableCatalogStoreResult<Decod
|
||||
)));
|
||||
}
|
||||
};
|
||||
let partition = avro_record_field(data_file, "partition")
|
||||
.and_then(avro_record_value_fields)
|
||||
.ok_or_else(|| TableCatalogStoreError::Invalid("manifest data file partition must be a record".to_string()))?;
|
||||
files.push(ManifestDataFileReference {
|
||||
location: file_path.to_string(),
|
||||
format_version,
|
||||
@@ -242,14 +218,15 @@ pub(crate) fn decode_manifest_avro(data: &[u8]) -> TableCatalogStoreResult<Decod
|
||||
file_size_bytes: avro_record_field(data_file, "file_size_in_bytes")
|
||||
.and_then(avro_i64_value)
|
||||
.and_then(|value| u64::try_from(value).ok()),
|
||||
partition,
|
||||
partition: avro_record_field(data_file, "partition")
|
||||
.and_then(avro_record_value_fields)
|
||||
.unwrap_or_default(),
|
||||
sort_order_id: avro_record_field(data_file, "sort_order_id").and_then(avro_i32_value),
|
||||
});
|
||||
}
|
||||
Ok(DecodedManifest {
|
||||
references: files,
|
||||
decoded_size,
|
||||
partition_spec_id,
|
||||
})
|
||||
}
|
||||
|
||||
@@ -263,8 +240,6 @@ pub(crate) async fn decode_manifest_avro_async(data: Vec<u8>) -> TableCatalogSto
|
||||
enum AvroContainerCodec {
|
||||
Null,
|
||||
Deflate,
|
||||
Snappy,
|
||||
Zstandard,
|
||||
}
|
||||
|
||||
fn validate_avro_container(data: &[u8]) -> TableCatalogStoreResult<usize> {
|
||||
@@ -325,8 +300,6 @@ fn validate_avro_container(data: &[u8]) -> TableCatalogStoreResult<usize> {
|
||||
let codec = match codec.unwrap_or(b"null") {
|
||||
b"null" => AvroContainerCodec::Null,
|
||||
b"deflate" => AvroContainerCodec::Deflate,
|
||||
b"snappy" => AvroContainerCodec::Snappy,
|
||||
b"zstandard" => AvroContainerCodec::Zstandard,
|
||||
codec => {
|
||||
return Err(TableCatalogStoreError::Unsupported(format!(
|
||||
"Avro codec {} is not supported for table commit validation",
|
||||
@@ -396,34 +369,6 @@ fn avro_block_decoded_size(codec: AvroContainerCodec, block: &[u8], remaining_si
|
||||
usize::try_from(decoded_size)
|
||||
.map_err(|_| TableCatalogStoreError::Invalid("Avro decoded data size is invalid".to_string()))
|
||||
}
|
||||
AvroContainerCodec::Snappy => {
|
||||
let data_end = block
|
||||
.len()
|
||||
.checked_sub(4)
|
||||
.ok_or_else(|| TableCatalogStoreError::Invalid("Avro snappy block is missing its checksum".to_string()))?;
|
||||
let decoded_size = snap::raw::decompress_len(&block[..data_end])
|
||||
.map_err(|err| TableCatalogStoreError::Invalid(format!("failed to inspect Avro snappy block: {err}")))?;
|
||||
if decoded_size > remaining_size {
|
||||
return Err(TableCatalogStoreError::Invalid("Avro decoded data exceeds the commit limit".to_string()));
|
||||
}
|
||||
Ok(decoded_size)
|
||||
}
|
||||
AvroContainerCodec::Zstandard => {
|
||||
let limit = remaining_size
|
||||
.checked_add(1)
|
||||
.ok_or_else(|| TableCatalogStoreError::Invalid("Avro decoded data size limit overflowed".to_string()))?;
|
||||
let limit = u64::try_from(limit)
|
||||
.map_err(|_| TableCatalogStoreError::Invalid("Avro decoded data size limit is invalid".to_string()))?;
|
||||
let mut decoder = zstd::stream::read::Decoder::new(block)
|
||||
.map_err(|err| TableCatalogStoreError::Invalid(format!("failed to decompress Avro zstandard block: {err}")))?;
|
||||
decoder
|
||||
.window_log_max(AVRO_ZSTANDARD_MAX_WINDOW_LOG)
|
||||
.map_err(|err| TableCatalogStoreError::Invalid(format!("failed to bound Avro zstandard window: {err}")))?;
|
||||
let decoded_size = std::io::copy(&mut decoder.take(limit), &mut std::io::sink())
|
||||
.map_err(|err| TableCatalogStoreError::Invalid(format!("failed to decompress Avro zstandard block: {err}")))?;
|
||||
usize::try_from(decoded_size)
|
||||
.map_err(|_| TableCatalogStoreError::Invalid("Avro decoded data size is invalid".to_string()))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -500,40 +445,6 @@ fn avro_record_value_fields(value: &apache_avro::types::Value) -> Option<Vec<(St
|
||||
)
|
||||
}
|
||||
|
||||
fn avro_nullable_non_negative_i32(
|
||||
value: &apache_avro::types::Value,
|
||||
field: &str,
|
||||
label: &str,
|
||||
) -> TableCatalogStoreResult<Option<u64>> {
|
||||
let Some(value) = avro_record_field(value, field) else {
|
||||
return Ok(None);
|
||||
};
|
||||
match avro_non_union_value(value) {
|
||||
apache_avro::types::Value::Null => Ok(None),
|
||||
apache_avro::types::Value::Int(value) => u64::try_from(*value)
|
||||
.map(Some)
|
||||
.map_err(|_| TableCatalogStoreError::Invalid(format!("{label} field {field} must be a non-negative int"))),
|
||||
_ => Err(TableCatalogStoreError::Invalid(format!("{label} field {field} must be a nullable int"))),
|
||||
}
|
||||
}
|
||||
|
||||
fn avro_nullable_non_negative_i64(
|
||||
value: &apache_avro::types::Value,
|
||||
field: &str,
|
||||
label: &str,
|
||||
) -> TableCatalogStoreResult<Option<u64>> {
|
||||
let Some(value) = avro_record_field(value, field) else {
|
||||
return Ok(None);
|
||||
};
|
||||
match avro_non_union_value(value) {
|
||||
apache_avro::types::Value::Null => Ok(None),
|
||||
apache_avro::types::Value::Long(value) => u64::try_from(*value)
|
||||
.map(Some)
|
||||
.map_err(|_| TableCatalogStoreError::Invalid(format!("{label} field {field} must be a non-negative long"))),
|
||||
_ => Err(TableCatalogStoreError::Invalid(format!("{label} field {field} must be a nullable long"))),
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn avro_non_union_value(value: &apache_avro::types::Value) -> &apache_avro::types::Value {
|
||||
match value {
|
||||
apache_avro::types::Value::Union(_, inner) => avro_non_union_value(inner),
|
||||
@@ -561,127 +472,3 @@ fn avro_i64_value(value: &apache_avro::types::Value) -> Option<i64> {
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn rejects_v2_manifest_lists_without_required_count_fields() {
|
||||
let schema = apache_avro::Schema::parse_str(
|
||||
r#"{
|
||||
"type": "record",
|
||||
"name": "manifest_file",
|
||||
"fields": [
|
||||
{"name": "manifest_path", "type": "string"},
|
||||
{"name": "manifest_length", "type": "long"},
|
||||
{"name": "partition_spec_id", "type": "int"},
|
||||
{"name": "content", "type": "int"},
|
||||
{"name": "sequence_number", "type": "long"},
|
||||
{"name": "min_sequence_number", "type": "long"},
|
||||
{"name": "added_snapshot_id", "type": "long"}
|
||||
]
|
||||
}"#,
|
||||
)
|
||||
.expect("incomplete manifest-list schema should parse");
|
||||
let data = apache_avro::Writer::new(&schema, Vec::new())
|
||||
.expect("manifest-list writer should initialize")
|
||||
.into_inner()
|
||||
.expect("manifest-list bytes should flush");
|
||||
|
||||
let error = match decode_manifest_list_avro(&data) {
|
||||
Ok(_) => panic!("v2 count fields must be declared in the writer schema"),
|
||||
Err(error) => error,
|
||||
};
|
||||
assert_eq!(
|
||||
error,
|
||||
TableCatalogStoreError::Invalid("Iceberg v2 manifest list Avro schema is missing added_files_count".to_string())
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rejects_negative_nullable_manifest_list_counts() {
|
||||
let value = apache_avro::types::Value::Record(vec![
|
||||
("added_files_count".to_string(), apache_avro::types::Value::Int(-1)),
|
||||
("added_rows_count".to_string(), apache_avro::types::Value::Long(-1)),
|
||||
]);
|
||||
|
||||
assert_eq!(
|
||||
avro_nullable_non_negative_i32(&value, "added_files_count", "manifest list")
|
||||
.expect_err("negative file counts must be rejected"),
|
||||
TableCatalogStoreError::Invalid("manifest list field added_files_count must be a non-negative int".to_string())
|
||||
);
|
||||
assert_eq!(
|
||||
avro_nullable_non_negative_i64(&value, "added_rows_count", "manifest list")
|
||||
.expect_err("negative row counts must be rejected"),
|
||||
TableCatalogStoreError::Invalid("manifest list field added_rows_count must be a non-negative long".to_string())
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rejects_manifest_partition_with_non_record_schema() {
|
||||
let schema = apache_avro::Schema::parse_str(
|
||||
r#"{
|
||||
"type": "record",
|
||||
"name": "manifest_entry",
|
||||
"fields": [
|
||||
{"name": "status", "type": "int"},
|
||||
{"name": "snapshot_id", "type": "long"},
|
||||
{
|
||||
"name": "data_file",
|
||||
"type": {
|
||||
"type": "record",
|
||||
"name": "data_file",
|
||||
"fields": [
|
||||
{"name": "file_path", "type": "string"},
|
||||
{"name": "record_count", "type": "long"},
|
||||
{"name": "file_size_in_bytes", "type": "long"},
|
||||
{"name": "partition", "type": "string"}
|
||||
]
|
||||
}
|
||||
}
|
||||
]
|
||||
}"#,
|
||||
)
|
||||
.expect("manifest schema should parse");
|
||||
let mut writer = apache_avro::Writer::new(&schema, Vec::new()).expect("manifest writer should initialize");
|
||||
writer
|
||||
.append_value(apache_avro::types::Value::Record(vec![
|
||||
("status".to_string(), apache_avro::types::Value::Int(1)),
|
||||
("snapshot_id".to_string(), apache_avro::types::Value::Long(1)),
|
||||
(
|
||||
"data_file".to_string(),
|
||||
apache_avro::types::Value::Record(vec![
|
||||
(
|
||||
"file_path".to_string(),
|
||||
apache_avro::types::Value::String("s3://warehouse/tables/table-id/data/file.parquet".to_string()),
|
||||
),
|
||||
("record_count".to_string(), apache_avro::types::Value::Long(1)),
|
||||
("file_size_in_bytes".to_string(), apache_avro::types::Value::Long(1)),
|
||||
("partition".to_string(), apache_avro::types::Value::String("not-a-record".to_string())),
|
||||
]),
|
||||
),
|
||||
]))
|
||||
.expect("manifest record should append");
|
||||
let data = writer.into_inner().expect("manifest bytes should flush");
|
||||
|
||||
let error = match decode_manifest_avro(&data) {
|
||||
Ok(_) => panic!("manifest partitions must preserve their record shape"),
|
||||
Err(error) => error,
|
||||
};
|
||||
assert_eq!(
|
||||
error,
|
||||
TableCatalogStoreError::Invalid("manifest data file partition must be a record".to_string())
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rejects_oversized_zstandard_windows() {
|
||||
// Non-single-segment frame with a 2^28-byte window and one empty final block.
|
||||
let compressed = [0x28, 0xb5, 0x2f, 0xfd, 0x00, 0x90, 0x01, 0x00, 0x00];
|
||||
|
||||
let error = avro_block_decoded_size(AvroContainerCodec::Zstandard, &compressed, TABLE_MANIFEST_AVRO_MAX_DECODED_SIZE)
|
||||
.expect_err("zstandard windows larger than the manifest decode budget must be rejected");
|
||||
assert!(matches!(error, TableCatalogStoreError::Invalid(_)));
|
||||
}
|
||||
}
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -111,12 +111,8 @@ const TABLE_MANIFEST_AVRO_MAX_DECODED_SIZE: usize = 128 * 1024 * 1024;
|
||||
const TABLE_MANIFEST_AVRO_MAX_RECORDS: usize = 1_000_000;
|
||||
const TABLE_MANIFEST_AVRO_MAX_HEADER_ENTRIES: usize = 1_024;
|
||||
const TABLE_COMMIT_MAX_MANIFESTS: usize = 10_000;
|
||||
const TABLE_COMMIT_MAX_MANIFEST_TRAVERSALS: usize = 20_000;
|
||||
const TABLE_COMMIT_MAX_AVRO_BYTES: usize = 512 * 1024 * 1024;
|
||||
const TABLE_COMMIT_MAX_FILE_REFERENCES: usize = 1_000_000;
|
||||
const TABLE_COMMIT_MAX_STATISTICS_OBJECTS: usize = 1_024;
|
||||
const TABLE_COMMIT_MAX_STATISTICS_BYTES: usize = 512 * 1024 * 1024;
|
||||
const TABLE_STATISTICS_FILE_MAX_SIZE: usize = 128 * 1024 * 1024;
|
||||
pub(crate) const TABLE_COMMIT_OBJECT_VALIDATION_CONCURRENCY: usize = 16;
|
||||
pub const TABLE_RESERVED_PREFIX: &str = BUCKET_TABLE_RESERVED_PREFIX;
|
||||
const WAREHOUSE_ROOT: &str = "warehouses";
|
||||
|
||||
@@ -216,7 +216,7 @@ where
|
||||
Ok(fence)
|
||||
}
|
||||
|
||||
pub(super) async fn acquire_table_bucket_registry_write_permit(&self) -> TableCatalogStoreResult<TableCatalogLockGuard> {
|
||||
pub(super) async fn acquire_table_bucket_registry_write_permit(&self) -> TableCatalogStoreResult<Box<dyn Send>> {
|
||||
let fence_path = self.paths.backing_migration_global_fence_path();
|
||||
let lock_path = self.paths.backing_migration_global_fence_lock_path();
|
||||
let guard = self.backend.acquire_read_lock(self.catalog_bucket(), &lock_path).await?;
|
||||
@@ -235,7 +235,7 @@ where
|
||||
pub(super) async fn acquire_object_backed_catalog_write_permit(
|
||||
&self,
|
||||
table_bucket: &str,
|
||||
) -> TableCatalogStoreResult<TableCatalogLockGuard> {
|
||||
) -> TableCatalogStoreResult<Box<dyn Send>> {
|
||||
let lock_path = self.paths.backing_migration_fence_lock_path(table_bucket);
|
||||
let guard = self.backend.acquire_read_lock(self.catalog_bucket(), &lock_path).await?;
|
||||
if self.read_backing_migration_fence(table_bucket).await?.is_some() {
|
||||
@@ -290,7 +290,7 @@ where
|
||||
async fn collect_bucket_snapshot_with_locks(
|
||||
&self,
|
||||
table_bucket: &str,
|
||||
guards: &mut Vec<TableCatalogLockGuard>,
|
||||
guards: &mut Vec<Box<dyn Send>>,
|
||||
) -> TableCatalogStoreResult<StrongTableCatalogBucketSnapshot> {
|
||||
let bucket_path = self.paths.table_bucket_entry_path(table_bucket);
|
||||
guards.push(self.backend.acquire_write_lock(self.catalog_bucket(), &bucket_path).await?);
|
||||
|
||||
@@ -262,29 +262,6 @@ pub(crate) trait TableCatalogStore: Send + Sync {
|
||||
|
||||
async fn create_view(&self, entry: ViewEntry) -> TableCatalogStoreResult<()>;
|
||||
|
||||
async fn create_view_with_publication(
|
||||
&self,
|
||||
entry: ViewEntry,
|
||||
publication: &(dyn TableCommitPublication + Sync),
|
||||
) -> TableCatalogStoreResult<()> {
|
||||
publication.begin_table_bucket(&entry.table_bucket).await?;
|
||||
if !publication.holds_table_bucket(&entry.table_bucket) {
|
||||
return Err(TableCatalogStoreError::Internal(
|
||||
"view creation requires a table-bucket publication fence".to_string(),
|
||||
));
|
||||
}
|
||||
publication
|
||||
.prepare(&entry.table_bucket, &entry.namespace, &entry.view)
|
||||
.await?;
|
||||
if !publication.holds_table(&entry.table_bucket, &entry.namespace, &entry.view) {
|
||||
return Err(TableCatalogStoreError::Internal(
|
||||
"view creation requires a view publication fence".to_string(),
|
||||
));
|
||||
}
|
||||
let _publication_completion = TableCommitPublicationCompletion::new(publication);
|
||||
self.create_view(entry).await
|
||||
}
|
||||
|
||||
async fn list_views(&self, table_bucket: &str, namespace: &str) -> TableCatalogStoreResult<Vec<ViewEntry>>;
|
||||
|
||||
async fn list_views_page(
|
||||
@@ -306,32 +283,6 @@ pub(crate) trait TableCatalogStore: Send + Sync {
|
||||
|
||||
async fn replace_view(&self, request: ViewCommitRequest) -> TableCatalogStoreResult<ViewCommitResult>;
|
||||
|
||||
async fn replace_view_with_publication(
|
||||
&self,
|
||||
request: ViewCommitRequest,
|
||||
table_bucket_fence_required: bool,
|
||||
publication: &(dyn TableCommitPublication + Sync),
|
||||
) -> TableCatalogStoreResult<ViewCommitResult> {
|
||||
if table_bucket_fence_required {
|
||||
publication.begin_table_bucket(&request.table_bucket).await?;
|
||||
if !publication.holds_table_bucket(&request.table_bucket) {
|
||||
return Err(TableCatalogStoreError::Internal(
|
||||
"view replacement requires a table-bucket publication fence".to_string(),
|
||||
));
|
||||
}
|
||||
}
|
||||
publication
|
||||
.prepare(&request.table_bucket, &request.namespace, &request.view)
|
||||
.await?;
|
||||
if !publication.holds_table(&request.table_bucket, &request.namespace, &request.view) {
|
||||
return Err(TableCatalogStoreError::Internal(
|
||||
"view replacement requires a view publication fence".to_string(),
|
||||
));
|
||||
}
|
||||
let _publication_completion = TableCommitPublicationCompletion::new(publication);
|
||||
self.replace_view(request).await
|
||||
}
|
||||
|
||||
async fn drop_view(&self, table_bucket: &str, namespace: &str, view: &str) -> TableCatalogStoreResult<()>;
|
||||
|
||||
async fn get_commit_by_id(
|
||||
@@ -387,7 +338,7 @@ struct TableCommitLockPublication<'a, B> {
|
||||
struct TableCommitLockPublicationState {
|
||||
table_bucket: Option<String>,
|
||||
table: Option<(String, String, String)>,
|
||||
guards: Vec<TableCatalogLockGuard>,
|
||||
guards: Vec<Box<dyn Send>>,
|
||||
}
|
||||
|
||||
impl<'a, B> TableCommitLockPublication<'a, B> {
|
||||
@@ -458,17 +409,15 @@ where
|
||||
}
|
||||
|
||||
fn holds_table_bucket(&self, table_bucket: &str) -> bool {
|
||||
let state = self.state.lock();
|
||||
state.table_bucket.as_deref() == Some(table_bucket) && state.guards.iter().all(|guard| !guard.is_lock_lost())
|
||||
self.state.lock().table_bucket.as_deref() == Some(table_bucket)
|
||||
}
|
||||
|
||||
fn holds_table(&self, table_bucket: &str, namespace: &str, table: &str) -> bool {
|
||||
let state = self.state.lock();
|
||||
state
|
||||
self.state
|
||||
.lock()
|
||||
.table
|
||||
.as_ref()
|
||||
.is_some_and(|held| held.0 == table_bucket && held.1 == namespace && held.2 == table)
|
||||
&& state.guards.iter().all(|guard| !guard.is_lock_lost())
|
||||
}
|
||||
|
||||
fn complete(&self) {
|
||||
@@ -489,32 +438,6 @@ pub(crate) struct TableCatalogObjectMetadata {
|
||||
pub mod_time: Option<OffsetDateTime>,
|
||||
}
|
||||
|
||||
pub(crate) struct TableCatalogLockGuard {
|
||||
_guard: Box<dyn Send>,
|
||||
lock_lost: Option<Arc<rustfs_lock::distributed_lock::LockLostSignal>>,
|
||||
}
|
||||
|
||||
impl TableCatalogLockGuard {
|
||||
pub(crate) fn stable(guard: impl Send + 'static) -> Self {
|
||||
Self {
|
||||
_guard: Box::new(guard),
|
||||
lock_lost: None,
|
||||
}
|
||||
}
|
||||
|
||||
fn namespace(guard: rustfs_lock::NamespaceLockGuard) -> Self {
|
||||
let lock_lost = guard.lock_lost_signal();
|
||||
Self {
|
||||
_guard: Box::new(guard),
|
||||
lock_lost,
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn is_lock_lost(&self) -> bool {
|
||||
self.lock_lost.as_ref().is_some_and(|signal| signal.is_lost())
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub(crate) struct TableCatalogObjectListPage {
|
||||
pub objects: Vec<String>,
|
||||
@@ -665,11 +588,11 @@ pub(crate) trait TableCatalogObjectBackend: Clone + Send + Sync + 'static {
|
||||
Ok(TableCatalogObjectListPage { objects, is_truncated })
|
||||
}
|
||||
|
||||
async fn acquire_read_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<TableCatalogLockGuard> {
|
||||
async fn acquire_read_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<Box<dyn Send>> {
|
||||
self.acquire_write_lock(bucket, object).await
|
||||
}
|
||||
|
||||
async fn acquire_write_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<TableCatalogLockGuard>;
|
||||
async fn acquire_write_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<Box<dyn Send>>;
|
||||
|
||||
async fn begin_table_bucket_commit_publication(&self, _table_bucket: &str) -> TableCatalogStoreResult<()> {
|
||||
Ok(())
|
||||
@@ -1246,17 +1169,6 @@ where
|
||||
}
|
||||
}
|
||||
|
||||
async fn create_view_with_publication(
|
||||
&self,
|
||||
entry: ViewEntry,
|
||||
publication: &(dyn TableCommitPublication + Sync),
|
||||
) -> TableCatalogStoreResult<()> {
|
||||
match self {
|
||||
Self::ObjectBacked(store) => store.create_view_with_publication(entry, publication).await,
|
||||
Self::DurableStrong(store) => store.create_view_with_publication(entry, publication).await,
|
||||
}
|
||||
}
|
||||
|
||||
async fn list_views(&self, table_bucket: &str, namespace: &str) -> TableCatalogStoreResult<Vec<ViewEntry>> {
|
||||
match self {
|
||||
Self::ObjectBacked(store) => store.list_views(table_bucket, namespace).await,
|
||||
@@ -1291,26 +1203,6 @@ where
|
||||
}
|
||||
}
|
||||
|
||||
async fn replace_view_with_publication(
|
||||
&self,
|
||||
request: ViewCommitRequest,
|
||||
table_bucket_fence_required: bool,
|
||||
publication: &(dyn TableCommitPublication + Sync),
|
||||
) -> TableCatalogStoreResult<ViewCommitResult> {
|
||||
match self {
|
||||
Self::ObjectBacked(store) => {
|
||||
store
|
||||
.replace_view_with_publication(request, table_bucket_fence_required, publication)
|
||||
.await
|
||||
}
|
||||
Self::DurableStrong(store) => {
|
||||
store
|
||||
.replace_view_with_publication(request, table_bucket_fence_required, publication)
|
||||
.await
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn drop_view(&self, table_bucket: &str, namespace: &str, view: &str) -> TableCatalogStoreResult<()> {
|
||||
match self {
|
||||
Self::ObjectBacked(store) => store.drop_view(table_bucket, namespace, view).await,
|
||||
@@ -1794,7 +1686,7 @@ where
|
||||
})
|
||||
}
|
||||
|
||||
async fn acquire_write_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<TableCatalogLockGuard> {
|
||||
async fn acquire_write_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<Box<dyn Send>> {
|
||||
let lock = self
|
||||
.store
|
||||
.new_ns_lock(bucket, object)
|
||||
@@ -1804,10 +1696,10 @@ where
|
||||
.get_write_lock(get_lock_acquire_timeout())
|
||||
.await
|
||||
.map_err(|err| TableCatalogStoreError::Internal(format!("failed to acquire catalog table lock: {err}")))?;
|
||||
Ok(TableCatalogLockGuard::namespace(guard))
|
||||
Ok(Box::new(guard))
|
||||
}
|
||||
|
||||
async fn acquire_read_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<TableCatalogLockGuard> {
|
||||
async fn acquire_read_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<Box<dyn Send>> {
|
||||
let lock = self
|
||||
.store
|
||||
.new_ns_lock(bucket, object)
|
||||
@@ -1817,7 +1709,7 @@ where
|
||||
.get_read_lock(get_lock_acquire_timeout())
|
||||
.await
|
||||
.map_err(|err| TableCatalogStoreError::Internal(format!("failed to acquire catalog migration lock: {err}")))?;
|
||||
Ok(TableCatalogLockGuard::namespace(guard))
|
||||
Ok(Box::new(guard))
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1036,20 +1036,6 @@ where
|
||||
.await
|
||||
}
|
||||
|
||||
async fn restore_table_warehouse_index_after_failed_drop(&self, entry: &TableEntry, reason: &'static str) {
|
||||
if let Err(err) = self.reserve_table_warehouse_index(entry).await {
|
||||
tracing::warn!(
|
||||
table_bucket = %entry.table_bucket,
|
||||
namespace = %entry.namespace,
|
||||
table = %entry.table,
|
||||
table_id = %entry.table_id,
|
||||
reason,
|
||||
error = %err,
|
||||
"failed to restore table warehouse index after table drop stopped"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
async fn delete_table_warehouse_index_if_changed(&self, current: &TableEntry, next: &TableEntry) {
|
||||
let Ok(current_index) = table_warehouse_index_entry(current) else {
|
||||
return;
|
||||
@@ -1330,15 +1316,6 @@ where
|
||||
}
|
||||
self.ensure_table_warehouse_prefix_available(&entry).await?;
|
||||
let reservation = self.reserve_table_warehouse_index(&entry).await?;
|
||||
if !publication.holds_table_bucket(&entry.table_bucket)
|
||||
|| !publication.holds_table(&entry.table_bucket, &entry.namespace, &entry.table)
|
||||
{
|
||||
self.delete_created_table_warehouse_index(&entry, reservation, "table publication fence lost")
|
||||
.await;
|
||||
return Err(TableCatalogStoreError::Internal(
|
||||
"table registration publication fence was lost before catalog update".to_string(),
|
||||
));
|
||||
}
|
||||
let result = self
|
||||
.write_entry_unlocked(self.catalog_bucket(), &table_path, &entry, precondition)
|
||||
.await;
|
||||
@@ -1350,25 +1327,7 @@ where
|
||||
}
|
||||
|
||||
async fn write_view_entry(&self, entry: ViewEntry, precondition: TableCatalogPutPrecondition) -> TableCatalogStoreResult<()> {
|
||||
let publication = TableCommitLockPublication::new(&self.backend);
|
||||
self.write_view_entry_with_publication(entry, precondition, &publication)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn write_view_entry_with_publication(
|
||||
&self,
|
||||
entry: ViewEntry,
|
||||
precondition: TableCatalogPutPrecondition,
|
||||
publication: &(dyn TableCommitPublication + Sync),
|
||||
) -> TableCatalogStoreResult<()> {
|
||||
validate_view_entry_version_and_id(&entry)?;
|
||||
publication.begin_table_bucket(&entry.table_bucket).await?;
|
||||
if !publication.holds_table_bucket(&entry.table_bucket) {
|
||||
return Err(TableCatalogStoreError::Internal(
|
||||
"view creation requires a table-bucket publication fence".to_string(),
|
||||
));
|
||||
}
|
||||
let _publication_completion = TableCommitPublicationCompletion::new(publication);
|
||||
self.require_table_bucket(&entry.table_bucket).await?;
|
||||
let namespace = parse_namespace_for_store(&entry.namespace)?;
|
||||
let view = parse_table_for_store(&entry.view)?;
|
||||
@@ -1394,17 +1353,6 @@ where
|
||||
entry.table_bucket, entry.namespace, entry.view
|
||||
)));
|
||||
}
|
||||
// Preserve catalog -> publication -> object lock order across rolling upgrades.
|
||||
publication
|
||||
.prepare(&entry.table_bucket, &entry.namespace, &entry.view)
|
||||
.await?;
|
||||
if !publication.holds_table_bucket(&entry.table_bucket)
|
||||
|| !publication.holds_table(&entry.table_bucket, &entry.namespace, &entry.view)
|
||||
{
|
||||
return Err(TableCatalogStoreError::Internal(
|
||||
"view creation publication fence was lost before catalog update".to_string(),
|
||||
));
|
||||
}
|
||||
self.write_entry_unlocked(self.catalog_bucket(), &view_path, &entry, precondition)
|
||||
.await
|
||||
}
|
||||
@@ -4545,16 +4493,16 @@ where
|
||||
validate_commit_metadata_digest(&request, &new_metadata_object)?;
|
||||
let table_bucket = request.table_bucket.clone();
|
||||
let metadata_location = request.new_metadata_location.clone();
|
||||
let next_metadata_state = tokio::task::spawn_blocking(move || {
|
||||
table_metadata_commit_state(&table_bucket, &metadata_location, &new_metadata_object)
|
||||
let next_warehouse_location = tokio::task::spawn_blocking(move || {
|
||||
table_metadata_warehouse_location(&table_bucket, &metadata_location, &new_metadata_object)
|
||||
})
|
||||
.await
|
||||
.map_err(|err| TableCatalogStoreError::Internal(format!("table metadata parser task failed: {err}")))??;
|
||||
let warehouse_relocation = next_metadata_state
|
||||
.warehouse_location
|
||||
if next_warehouse_location
|
||||
.as_ref()
|
||||
.is_some_and(|warehouse_location| warehouse_location != ¤t.warehouse_location);
|
||||
if warehouse_relocation && !publication.holds_table_bucket(&request.table_bucket) {
|
||||
.is_some_and(|warehouse_location| warehouse_location != ¤t.warehouse_location)
|
||||
&& !publication.holds_table_bucket(&request.table_bucket)
|
||||
{
|
||||
return table_commit_result(
|
||||
&request.table_bucket,
|
||||
&request.namespace,
|
||||
@@ -4589,12 +4537,9 @@ where
|
||||
|
||||
let mut next = current.clone();
|
||||
next.metadata_location = staged_commit_log.new_metadata_location.clone();
|
||||
if let Some(warehouse_location) = next_metadata_state.warehouse_location {
|
||||
if let Some(warehouse_location) = next_warehouse_location {
|
||||
next.warehouse_location = warehouse_location;
|
||||
}
|
||||
if let Some(format_version) = next_metadata_state.format_version {
|
||||
next.format_version = format_version;
|
||||
}
|
||||
next.version_token = staged_commit_log.new_version_token.clone();
|
||||
next.generation = current.generation.saturating_add(1);
|
||||
if next.warehouse_location != current.warehouse_location {
|
||||
@@ -4640,24 +4585,6 @@ where
|
||||
);
|
||||
}
|
||||
|
||||
if !publication.holds_table(&request.table_bucket, &request.namespace, &request.table)
|
||||
|| (warehouse_relocation && !publication.holds_table_bucket(&request.table_bucket))
|
||||
{
|
||||
self.delete_created_table_warehouse_index(&next, reservation, "table publication fence lost")
|
||||
.await;
|
||||
return table_commit_result(
|
||||
&request.table_bucket,
|
||||
&request.namespace,
|
||||
&request.table,
|
||||
&request.commit_id,
|
||||
&request.operation,
|
||||
commit_started,
|
||||
Err(TableCatalogStoreError::Internal(
|
||||
"table commit publication fence was lost before pointer update".to_string(),
|
||||
)),
|
||||
);
|
||||
}
|
||||
|
||||
let cas_started = Instant::now();
|
||||
let cas_result = self
|
||||
.write_entry_unlocked(
|
||||
@@ -4735,21 +4662,20 @@ where
|
||||
)));
|
||||
};
|
||||
self.delete_owned_table_warehouse_index_for_drop(&entry).await?;
|
||||
if !publication.holds_table_bucket(table_bucket)
|
||||
|| !publication.holds_table(table_bucket, &namespace.public_name(), table.as_str())
|
||||
{
|
||||
self.restore_table_warehouse_index_after_failed_drop(&entry, "table publication fence lost")
|
||||
.await;
|
||||
return Err(TableCatalogStoreError::Internal(
|
||||
"table drop publication fence was lost before catalog update".to_string(),
|
||||
));
|
||||
}
|
||||
if let Err(err) = self.backend.delete_object_unlocked(self.catalog_bucket(), &object).await {
|
||||
match self.read_table_with_etag_unlocked(table_bucket, &namespace, &table).await {
|
||||
Ok(None) => return Ok(()),
|
||||
Ok(Some((current, _))) if current == entry => {
|
||||
self.restore_table_warehouse_index_after_failed_drop(&entry, "table entry delete failed")
|
||||
.await;
|
||||
if let Err(restore_err) = self.reserve_table_warehouse_index(&entry).await {
|
||||
tracing::warn!(
|
||||
table_bucket = %entry.table_bucket,
|
||||
namespace = %entry.namespace,
|
||||
table = %entry.table,
|
||||
table_id = %entry.table_id,
|
||||
error = %restore_err,
|
||||
"failed to restore table warehouse index after table entry delete failure"
|
||||
);
|
||||
}
|
||||
}
|
||||
Ok(Some(_)) => {
|
||||
return Err(TableCatalogStoreError::Internal(format!(
|
||||
@@ -4777,15 +4703,6 @@ where
|
||||
self.write_view_entry(entry, TableCatalogPutPrecondition::IfAbsent).await
|
||||
}
|
||||
|
||||
async fn create_view_with_publication(
|
||||
&self,
|
||||
entry: ViewEntry,
|
||||
publication: &(dyn TableCommitPublication + Sync),
|
||||
) -> TableCatalogStoreResult<()> {
|
||||
self.write_view_entry_with_publication(entry, TableCatalogPutPrecondition::IfAbsent, publication)
|
||||
.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();
|
||||
@@ -4840,26 +4757,8 @@ where
|
||||
}
|
||||
|
||||
async fn replace_view(&self, request: ViewCommitRequest) -> TableCatalogStoreResult<ViewCommitResult> {
|
||||
let publication = TableCommitLockPublication::new(&self.backend);
|
||||
self.replace_view_with_publication(request, true, &publication).await
|
||||
}
|
||||
|
||||
async fn replace_view_with_publication(
|
||||
&self,
|
||||
request: ViewCommitRequest,
|
||||
table_bucket_fence_required: bool,
|
||||
publication: &(dyn TableCommitPublication + Sync),
|
||||
) -> TableCatalogStoreResult<ViewCommitResult> {
|
||||
let namespace = parse_namespace_for_store(&request.namespace)?;
|
||||
let view = parse_table_for_store(&request.view)?;
|
||||
if table_bucket_fence_required {
|
||||
publication.begin_table_bucket(&request.table_bucket).await?;
|
||||
if !publication.holds_table_bucket(&request.table_bucket) {
|
||||
return Err(TableCatalogStoreError::Internal(
|
||||
"view replacement requires a table-bucket publication fence".to_string(),
|
||||
));
|
||||
}
|
||||
}
|
||||
let _migration_guard = self.acquire_object_backed_catalog_write_permit(&request.table_bucket).await?;
|
||||
let namespace_path = self.paths.namespace_entry_path(&request.table_bucket, &namespace);
|
||||
let _namespace_guard = self
|
||||
@@ -4868,16 +4767,6 @@ where
|
||||
.await?;
|
||||
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?;
|
||||
// Preserve catalog -> publication -> object lock order across rolling upgrades.
|
||||
publication
|
||||
.prepare(&request.table_bucket, &request.namespace, &request.view)
|
||||
.await?;
|
||||
if !publication.holds_table(&request.table_bucket, &request.namespace, &request.view) {
|
||||
return Err(TableCatalogStoreError::Internal(
|
||||
"view replacement requires a table publication fence".to_string(),
|
||||
));
|
||||
}
|
||||
let _publication_completion = TableCommitPublicationCompletion::new(publication);
|
||||
let Some((current, current_etag)) = self
|
||||
.read_view_with_etag_unlocked(&request.table_bucket, &namespace, &view)
|
||||
.await?
|
||||
@@ -4925,14 +4814,6 @@ where
|
||||
})
|
||||
.await
|
||||
.map_err(|err| TableCatalogStoreError::Internal(format!("view metadata parser task failed: {err}")))??;
|
||||
let warehouse_relocation = next_warehouse_location
|
||||
.as_deref()
|
||||
.is_some_and(|location| location != current.warehouse_location);
|
||||
if warehouse_relocation && !publication.holds_table_bucket(&request.table_bucket) {
|
||||
return Err(TableCatalogStoreError::Internal(
|
||||
"view warehouse relocation requires a table-bucket publication fence".to_string(),
|
||||
));
|
||||
}
|
||||
|
||||
let mut next = current;
|
||||
next.metadata_location = request.new_metadata_location;
|
||||
@@ -4941,40 +4822,13 @@ where
|
||||
}
|
||||
next.version_token = format!("token-{}", Uuid::new_v4());
|
||||
next.generation = next.generation.saturating_add(1);
|
||||
if !publication.holds_table(&request.table_bucket, &request.namespace, &request.view)
|
||||
|| ((table_bucket_fence_required || warehouse_relocation) && !publication.holds_table_bucket(&request.table_bucket))
|
||||
{
|
||||
return Err(TableCatalogStoreError::Internal(
|
||||
"view replacement publication fence was lost before catalog update".to_string(),
|
||||
));
|
||||
}
|
||||
let write_result = self
|
||||
.write_entry_unlocked(
|
||||
self.catalog_bucket(),
|
||||
&view_path,
|
||||
&next,
|
||||
TableCatalogPutPrecondition::IfMatch(current_etag),
|
||||
)
|
||||
.await;
|
||||
if let Err(err) = write_result {
|
||||
match self
|
||||
.read_view_with_etag_unlocked(&request.table_bucket, &namespace, &view)
|
||||
.await
|
||||
{
|
||||
Ok(Some((persisted, _))) if persisted == next => {}
|
||||
Ok(_) => return Err(err),
|
||||
Err(read_err) => {
|
||||
tracing::warn!(
|
||||
table_bucket = %request.table_bucket,
|
||||
namespace = %request.namespace,
|
||||
view = %request.view,
|
||||
error = %read_err,
|
||||
"failed to verify view state after an ambiguous catalog update"
|
||||
);
|
||||
return Err(err);
|
||||
}
|
||||
}
|
||||
}
|
||||
self.write_entry_unlocked(
|
||||
self.catalog_bucket(),
|
||||
&view_path,
|
||||
&next,
|
||||
TableCatalogPutPrecondition::IfMatch(current_etag),
|
||||
)
|
||||
.await?;
|
||||
Ok(ViewCommitResult { view: next })
|
||||
}
|
||||
|
||||
|
||||
@@ -554,7 +554,7 @@ where
|
||||
|
||||
// Ordinary mutations hold the global migration read lock before the local write lock; migration takes the
|
||||
// write side before invoking its dedicated snapshot mutation methods.
|
||||
async fn acquire_snapshot_write_permit(&self) -> TableCatalogStoreResult<TableCatalogLockGuard> {
|
||||
async fn acquire_snapshot_write_permit(&self) -> TableCatalogStoreResult<Box<dyn Send>> {
|
||||
let lock_path = TableCatalogObjectPaths::default().backing_migration_global_fence_lock_path();
|
||||
self.object_backend.acquire_read_lock(RUSTFS_META_BUCKET, &lock_path).await
|
||||
}
|
||||
@@ -1775,7 +1775,7 @@ where
|
||||
request: &TableCommitRequest,
|
||||
namespace: &Namespace,
|
||||
table: &IdentifierSegment,
|
||||
next_metadata_state: TableMetadataCommitState,
|
||||
next_warehouse_location: Option<String>,
|
||||
) -> TableCatalogStoreResult<TableCommitResult> {
|
||||
let key = Self::table_key(&request.table_bucket, namespace, table);
|
||||
let current = Self::validate_new_table_commit_locked(state, &key, request)?;
|
||||
@@ -1802,12 +1802,9 @@ where
|
||||
|
||||
let mut next = current;
|
||||
next.metadata_location = commit_log.new_metadata_location.clone();
|
||||
if let Some(warehouse_location) = next_metadata_state.warehouse_location {
|
||||
if let Some(warehouse_location) = next_warehouse_location {
|
||||
next.warehouse_location = warehouse_location;
|
||||
}
|
||||
if let Some(format_version) = next_metadata_state.format_version {
|
||||
next.format_version = format_version;
|
||||
}
|
||||
Self::ensure_table_warehouse_prefix_available_locked(state, &next, &key)?;
|
||||
next.version_token = commit_log.new_version_token.clone();
|
||||
next.generation = next.generation.saturating_add(1);
|
||||
@@ -2173,7 +2170,6 @@ where
|
||||
let _write_guard = self.write_lock.lock().await;
|
||||
self.hydrate_state().await?;
|
||||
let key = Self::table_key(&entry.table_bucket, &namespace, &table);
|
||||
let publication_identity = (entry.table_bucket.clone(), entry.namespace.clone(), entry.table.clone());
|
||||
let (snapshot, precondition, postcondition) = {
|
||||
let state = self.state.lock().await;
|
||||
Self::require_table_bucket_in_state(&state, &entry.table_bucket)?;
|
||||
@@ -2202,13 +2198,6 @@ where
|
||||
StrongSnapshotWritePostcondition::TablePresent(entry),
|
||||
)
|
||||
};
|
||||
if !publication.holds_table_bucket(&publication_identity.0)
|
||||
|| !publication.holds_table(&publication_identity.0, &publication_identity.1, &publication_identity.2)
|
||||
{
|
||||
return Err(TableCatalogStoreError::Internal(
|
||||
"table registration publication fence was lost before snapshot update".to_string(),
|
||||
));
|
||||
}
|
||||
self.finalize_snapshot_write(snapshot, precondition, postcondition).await
|
||||
}
|
||||
|
||||
@@ -2490,15 +2479,9 @@ where
|
||||
let result = match prepared_result {
|
||||
Ok((result, Some((snapshot, precondition)))) => {
|
||||
let postcondition = Self::commit_write_postcondition(&request.table_bucket, &result.commit_log);
|
||||
if !publication.holds_table(&request.table_bucket, &request.namespace, &request.table) {
|
||||
Err(TableCatalogStoreError::Internal(
|
||||
"table commit publication fence was lost before snapshot update".to_string(),
|
||||
))
|
||||
} else {
|
||||
self.finalize_snapshot_write(snapshot, precondition, postcondition)
|
||||
.await
|
||||
.map(|_| result)
|
||||
}
|
||||
self.finalize_snapshot_write(snapshot, precondition, postcondition)
|
||||
.await
|
||||
.map(|_| result)
|
||||
}
|
||||
Ok((result, None)) => Ok(result),
|
||||
Err(err) => Err(err),
|
||||
@@ -2536,8 +2519,8 @@ where
|
||||
validate_commit_metadata_digest(&request, &new_metadata_object)?;
|
||||
let table_bucket = request.table_bucket.clone();
|
||||
let metadata_location = request.new_metadata_location.clone();
|
||||
let next_metadata_state = tokio::task::spawn_blocking(move || {
|
||||
table_metadata_commit_state(&table_bucket, &metadata_location, &new_metadata_object)
|
||||
let next_warehouse_location = tokio::task::spawn_blocking(move || {
|
||||
table_metadata_warehouse_location(&table_bucket, &metadata_location, &new_metadata_object)
|
||||
})
|
||||
.await
|
||||
.map_err(|err| TableCatalogStoreError::Internal(format!("table metadata parser task failed: {err}")))??;
|
||||
@@ -2556,11 +2539,11 @@ where
|
||||
))
|
||||
})?
|
||||
};
|
||||
let warehouse_relocation = next_metadata_state
|
||||
.warehouse_location
|
||||
if next_warehouse_location
|
||||
.as_ref()
|
||||
.is_some_and(|warehouse_location| warehouse_location != ¤t_warehouse_location);
|
||||
if warehouse_relocation && !publication.holds_table_bucket(&request.table_bucket) {
|
||||
.is_some_and(|warehouse_location| warehouse_location != ¤t_warehouse_location)
|
||||
&& !publication.holds_table_bucket(&request.table_bucket)
|
||||
{
|
||||
return table_commit_result(
|
||||
&request.table_bucket,
|
||||
&request.namespace,
|
||||
@@ -2578,7 +2561,7 @@ where
|
||||
let prepared_result = {
|
||||
let state = self.state.lock().await;
|
||||
let (precondition, mut draft_state) = Self::snapshot_draft_context_locked(&state);
|
||||
match Self::apply_commit_locked(&mut draft_state, &request, &namespace, &table, next_metadata_state) {
|
||||
match Self::apply_commit_locked(&mut draft_state, &request, &namespace, &table, next_warehouse_location) {
|
||||
Ok(result) => Self::snapshot_from_mutated_state_locked(&mut draft_state, self.snapshot_write_version)
|
||||
.map(|snapshot| (result, snapshot, precondition)),
|
||||
Err(err) => Err(err),
|
||||
@@ -2587,16 +2570,7 @@ where
|
||||
let result = match prepared_result {
|
||||
Ok((result, snapshot, precondition)) => {
|
||||
let postcondition = Self::commit_write_postcondition(&request.table_bucket, &result.commit_log);
|
||||
let snapshot_result = if publication.holds_table(&request.table_bucket, &request.namespace, &request.table)
|
||||
&& (!warehouse_relocation || publication.holds_table_bucket(&request.table_bucket))
|
||||
{
|
||||
self.finalize_snapshot_write(snapshot, precondition, postcondition).await
|
||||
} else {
|
||||
Err(TableCatalogStoreError::Internal(
|
||||
"table commit publication fence was lost before snapshot update".to_string(),
|
||||
))
|
||||
};
|
||||
match snapshot_result {
|
||||
match self.finalize_snapshot_write(snapshot, precondition, postcondition).await {
|
||||
Ok(()) => Ok(result),
|
||||
Err(err) => {
|
||||
let replay = {
|
||||
@@ -2673,26 +2647,13 @@ where
|
||||
},
|
||||
)
|
||||
};
|
||||
if !publication.holds_table_bucket(table_bucket)
|
||||
|| !publication.holds_table(table_bucket, &namespace.public_name(), table.as_str())
|
||||
{
|
||||
return Err(TableCatalogStoreError::Internal(
|
||||
"table drop publication fence was lost before snapshot update".to_string(),
|
||||
));
|
||||
}
|
||||
self.finalize_snapshot_write(snapshot, precondition, postcondition).await
|
||||
}
|
||||
|
||||
async fn create_view(&self, entry: ViewEntry) -> TableCatalogStoreResult<()> {
|
||||
let publication = TableCommitLockPublication::new(&self.object_backend);
|
||||
self.create_view_with_publication(entry, &publication).await
|
||||
}
|
||||
|
||||
async fn create_view_with_publication(
|
||||
&self,
|
||||
entry: ViewEntry,
|
||||
publication: &(dyn TableCommitPublication + Sync),
|
||||
) -> TableCatalogStoreResult<()> {
|
||||
let _migration_guard = self.acquire_snapshot_write_permit().await?;
|
||||
let _write_guard = self.write_lock.lock().await;
|
||||
self.hydrate_state().await?;
|
||||
validate_view_entry_version_and_id(&entry)?;
|
||||
validate_view_warehouse_location(&entry.table_bucket, &entry.warehouse_location)?;
|
||||
let namespace = parse_namespace_for_store(&entry.namespace)?;
|
||||
@@ -2702,26 +2663,7 @@ where
|
||||
"view metadata location must be inside the view metadata directory".to_string(),
|
||||
));
|
||||
}
|
||||
publication.begin_table_bucket(&entry.table_bucket).await?;
|
||||
if !publication.holds_table_bucket(&entry.table_bucket) {
|
||||
return Err(TableCatalogStoreError::Internal(
|
||||
"view creation requires a table-bucket publication fence".to_string(),
|
||||
));
|
||||
}
|
||||
let _publication_completion = TableCommitPublicationCompletion::new(publication);
|
||||
let _migration_guard = self.acquire_snapshot_write_permit().await?;
|
||||
publication
|
||||
.prepare(&entry.table_bucket, &entry.namespace, &entry.view)
|
||||
.await?;
|
||||
if !publication.holds_table(&entry.table_bucket, &entry.namespace, &entry.view) {
|
||||
return Err(TableCatalogStoreError::Internal(
|
||||
"view creation requires a table publication fence".to_string(),
|
||||
));
|
||||
}
|
||||
let _write_guard = self.write_lock.lock().await;
|
||||
self.hydrate_state().await?;
|
||||
let key = Self::table_key(&entry.table_bucket, &namespace, &view);
|
||||
let publication_identity = (entry.table_bucket.clone(), entry.namespace.clone(), entry.view.clone());
|
||||
let (snapshot, precondition, postcondition) = {
|
||||
let state = self.state.lock().await;
|
||||
Self::require_table_bucket_in_state(&state, &entry.table_bucket)?;
|
||||
@@ -2740,13 +2682,6 @@ where
|
||||
StrongSnapshotWritePostcondition::ViewPresent(entry),
|
||||
)
|
||||
};
|
||||
if !publication.holds_table_bucket(&publication_identity.0)
|
||||
|| !publication.holds_table(&publication_identity.0, &publication_identity.1, &publication_identity.2)
|
||||
{
|
||||
return Err(TableCatalogStoreError::Internal(
|
||||
"view creation publication fence was lost before snapshot update".to_string(),
|
||||
));
|
||||
}
|
||||
self.finalize_snapshot_write(snapshot, precondition, postcondition).await
|
||||
}
|
||||
|
||||
@@ -2816,38 +2751,11 @@ where
|
||||
}
|
||||
|
||||
async fn replace_view(&self, request: ViewCommitRequest) -> TableCatalogStoreResult<ViewCommitResult> {
|
||||
let publication = TableCommitLockPublication::new(&self.object_backend);
|
||||
self.replace_view_with_publication(request, true, &publication).await
|
||||
}
|
||||
|
||||
async fn replace_view_with_publication(
|
||||
&self,
|
||||
request: ViewCommitRequest,
|
||||
table_bucket_fence_required: bool,
|
||||
publication: &(dyn TableCommitPublication + Sync),
|
||||
) -> TableCatalogStoreResult<ViewCommitResult> {
|
||||
if table_bucket_fence_required {
|
||||
publication.begin_table_bucket(&request.table_bucket).await?;
|
||||
if !publication.holds_table_bucket(&request.table_bucket) {
|
||||
return Err(TableCatalogStoreError::Internal(
|
||||
"view replacement requires a table-bucket publication fence".to_string(),
|
||||
));
|
||||
}
|
||||
}
|
||||
let _publication_completion = TableCommitPublicationCompletion::new(publication);
|
||||
let _migration_guard = self.acquire_snapshot_write_permit().await?;
|
||||
let namespace = parse_namespace_for_store(&request.namespace)?;
|
||||
let view = parse_table_for_store(&request.view)?;
|
||||
publication
|
||||
.prepare(&request.table_bucket, &request.namespace, &request.view)
|
||||
.await?;
|
||||
if !publication.holds_table(&request.table_bucket, &request.namespace, &request.view) {
|
||||
return Err(TableCatalogStoreError::Internal(
|
||||
"view replacement requires a table publication fence".to_string(),
|
||||
));
|
||||
}
|
||||
let write_guard = self.write_lock.lock().await;
|
||||
self.hydrate_state().await?;
|
||||
let namespace = parse_namespace_for_store(&request.namespace)?;
|
||||
let view = parse_table_for_store(&request.view)?;
|
||||
let key = Self::table_key(&request.table_bucket, &namespace, &view);
|
||||
let expected_view_id = {
|
||||
let state = self.state.lock().await;
|
||||
@@ -2890,7 +2798,7 @@ where
|
||||
|
||||
let _write_guard = self.write_lock.lock().await;
|
||||
self.hydrate_state().await?;
|
||||
let (snapshot, precondition, next, postcondition, warehouse_relocation) = {
|
||||
let (snapshot, precondition, next, postcondition) = {
|
||||
let state = self.state.lock().await;
|
||||
Self::ensure_identifier_is_unambiguous_locked(&state, &key)?;
|
||||
let Some(current) = state.views.get(&key).cloned() else {
|
||||
@@ -2920,14 +2828,6 @@ where
|
||||
"current view metadata location does not match expected location".to_string(),
|
||||
));
|
||||
}
|
||||
let warehouse_relocation = next_warehouse_location
|
||||
.as_deref()
|
||||
.is_some_and(|location| location != current.warehouse_location);
|
||||
if warehouse_relocation && !publication.holds_table_bucket(&request.table_bucket) {
|
||||
return Err(TableCatalogStoreError::Internal(
|
||||
"view warehouse relocation requires a table-bucket publication fence".to_string(),
|
||||
));
|
||||
}
|
||||
|
||||
let mut next = current;
|
||||
next.metadata_location = request.new_metadata_location;
|
||||
@@ -2943,16 +2843,8 @@ where
|
||||
precondition,
|
||||
next.clone(),
|
||||
StrongSnapshotWritePostcondition::ViewPresent(next),
|
||||
warehouse_relocation,
|
||||
)
|
||||
};
|
||||
if !publication.holds_table(&request.table_bucket, &request.namespace, &request.view)
|
||||
|| ((table_bucket_fence_required || warehouse_relocation) && !publication.holds_table_bucket(&request.table_bucket))
|
||||
{
|
||||
return Err(TableCatalogStoreError::Internal(
|
||||
"view replacement publication fence was lost before snapshot update".to_string(),
|
||||
));
|
||||
}
|
||||
self.finalize_snapshot_write(snapshot, precondition, postcondition).await?;
|
||||
Ok(ViewCommitResult { view: next })
|
||||
}
|
||||
|
||||
@@ -55,35 +55,23 @@ pub(crate) fn table_metadata_json(table_uuid: &str, location: &str) -> serde_jso
|
||||
})
|
||||
}
|
||||
|
||||
pub(crate) fn manifest_list_avro_bytes(manifests: &[(&str, usize)], sequence_number: i64, snapshot_id: i64) -> Vec<u8> {
|
||||
let manifests = manifests
|
||||
pub(crate) fn manifest_list_avro_bytes(manifest_paths: &[&str], sequence_number: i64, snapshot_id: i64) -> Vec<u8> {
|
||||
let manifests = manifest_paths
|
||||
.iter()
|
||||
.map(|(manifest_path, manifest_length)| (*manifest_path, *manifest_length, 0, sequence_number, snapshot_id))
|
||||
.map(|manifest_path| (*manifest_path, 0, sequence_number, snapshot_id))
|
||||
.collect::<Vec<_>>();
|
||||
manifest_list_avro_entries_with_partition_specs(&manifests)
|
||||
}
|
||||
|
||||
pub(crate) fn manifest_list_avro_entries(manifests: &[(&str, usize, i64, i64)]) -> Vec<u8> {
|
||||
pub(crate) fn manifest_list_avro_entries(manifests: &[(&str, i64, i64)]) -> Vec<u8> {
|
||||
let manifests = manifests
|
||||
.iter()
|
||||
.map(|(manifest_path, manifest_length, sequence_number, snapshot_id)| {
|
||||
(*manifest_path, *manifest_length, 0, *sequence_number, *snapshot_id)
|
||||
})
|
||||
.map(|(manifest_path, sequence_number, snapshot_id)| (*manifest_path, 0, *sequence_number, *snapshot_id))
|
||||
.collect::<Vec<_>>();
|
||||
manifest_list_avro_entries_with_partition_specs(&manifests)
|
||||
}
|
||||
|
||||
pub(crate) fn manifest_list_avro_entries_with_partition_specs(manifests: &[(&str, usize, i32, i64, i64)]) -> Vec<u8> {
|
||||
let manifests = manifests
|
||||
.iter()
|
||||
.map(|(path, length, spec_id, sequence_number, snapshot_id)| {
|
||||
(*path, *length, *spec_id, 0, *sequence_number, *snapshot_id)
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
manifest_list_avro_entries_with_content(&manifests)
|
||||
}
|
||||
|
||||
pub(crate) fn manifest_list_avro_entries_with_content(manifests: &[(&str, usize, i32, i32, i64, i64)]) -> Vec<u8> {
|
||||
pub(crate) fn manifest_list_avro_entries_with_partition_specs(manifests: &[(&str, i32, i64, i64)]) -> Vec<u8> {
|
||||
let schema = apache_avro::Schema::parse_str(
|
||||
r#"
|
||||
{
|
||||
@@ -109,19 +97,16 @@ pub(crate) fn manifest_list_avro_entries_with_content(manifests: &[(&str, usize,
|
||||
)
|
||||
.expect("manifest list avro schema should parse");
|
||||
let mut writer = apache_avro::Writer::new(&schema, Vec::new()).expect("manifest list writer should initialize");
|
||||
for (manifest_path, manifest_length, partition_spec_id, content, sequence_number, snapshot_id) in manifests {
|
||||
for (manifest_path, partition_spec_id, sequence_number, snapshot_id) in manifests {
|
||||
writer
|
||||
.append_value(apache_avro::types::Value::Record(vec![
|
||||
(
|
||||
"manifest_path".to_string(),
|
||||
apache_avro::types::Value::String((*manifest_path).to_string()),
|
||||
),
|
||||
(
|
||||
"manifest_length".to_string(),
|
||||
apache_avro::types::Value::Long(i64::try_from(*manifest_length).expect("test manifest length should fit")),
|
||||
),
|
||||
("manifest_length".to_string(), apache_avro::types::Value::Long(1)),
|
||||
("partition_spec_id".to_string(), apache_avro::types::Value::Int(*partition_spec_id)),
|
||||
("content".to_string(), apache_avro::types::Value::Int(*content)),
|
||||
("content".to_string(), apache_avro::types::Value::Int(0)),
|
||||
("sequence_number".to_string(), apache_avro::types::Value::Long(*sequence_number)),
|
||||
("min_sequence_number".to_string(), apache_avro::types::Value::Long(*sequence_number)),
|
||||
("added_snapshot_id".to_string(), apache_avro::types::Value::Long(*snapshot_id)),
|
||||
@@ -137,86 +122,7 @@ pub(crate) fn manifest_list_avro_entries_with_content(manifests: &[(&str, usize,
|
||||
writer.into_inner().expect("manifest list avro bytes should flush")
|
||||
}
|
||||
|
||||
pub(crate) fn manifest_list_avro_entries_with_nullable_counts(manifests: &[(&str, usize, i32, i64, i64)]) -> Vec<u8> {
|
||||
let schema = apache_avro::Schema::parse_str(
|
||||
r#"
|
||||
{
|
||||
"type": "record",
|
||||
"name": "manifest_file",
|
||||
"fields": [
|
||||
{"name": "manifest_path", "type": "string"},
|
||||
{"name": "manifest_length", "type": "long"},
|
||||
{"name": "partition_spec_id", "type": "int"},
|
||||
{"name": "content", "type": "int"},
|
||||
{"name": "sequence_number", "type": "long"},
|
||||
{"name": "min_sequence_number", "type": "long"},
|
||||
{"name": "added_snapshot_id", "type": "long"},
|
||||
{"name": "added_files_count", "type": ["null", "int"], "default": null},
|
||||
{"name": "existing_files_count", "type": ["null", "int"], "default": null},
|
||||
{"name": "deleted_files_count", "type": ["null", "int"], "default": null},
|
||||
{"name": "added_rows_count", "type": ["null", "long"], "default": null},
|
||||
{"name": "existing_rows_count", "type": ["null", "long"], "default": null},
|
||||
{"name": "deleted_rows_count", "type": ["null", "long"], "default": null}
|
||||
]
|
||||
}
|
||||
"#,
|
||||
)
|
||||
.expect("manifest list avro schema should parse");
|
||||
let mut writer = apache_avro::Writer::new(&schema, Vec::new()).expect("manifest list writer should initialize");
|
||||
for (manifest_path, manifest_length, partition_spec_id, sequence_number, snapshot_id) in manifests {
|
||||
writer
|
||||
.append_value(apache_avro::types::Value::Record(vec![
|
||||
(
|
||||
"manifest_path".to_string(),
|
||||
apache_avro::types::Value::String((*manifest_path).to_string()),
|
||||
),
|
||||
(
|
||||
"manifest_length".to_string(),
|
||||
apache_avro::types::Value::Long(i64::try_from(*manifest_length).expect("test manifest length should fit")),
|
||||
),
|
||||
("partition_spec_id".to_string(), apache_avro::types::Value::Int(*partition_spec_id)),
|
||||
("content".to_string(), apache_avro::types::Value::Int(0)),
|
||||
("sequence_number".to_string(), apache_avro::types::Value::Long(*sequence_number)),
|
||||
("min_sequence_number".to_string(), apache_avro::types::Value::Long(*sequence_number)),
|
||||
("added_snapshot_id".to_string(), apache_avro::types::Value::Long(*snapshot_id)),
|
||||
(
|
||||
"added_files_count".to_string(),
|
||||
apache_avro::types::Value::Union(0, Box::new(apache_avro::types::Value::Null)),
|
||||
),
|
||||
(
|
||||
"existing_files_count".to_string(),
|
||||
apache_avro::types::Value::Union(0, Box::new(apache_avro::types::Value::Null)),
|
||||
),
|
||||
(
|
||||
"deleted_files_count".to_string(),
|
||||
apache_avro::types::Value::Union(0, Box::new(apache_avro::types::Value::Null)),
|
||||
),
|
||||
(
|
||||
"added_rows_count".to_string(),
|
||||
apache_avro::types::Value::Union(0, Box::new(apache_avro::types::Value::Null)),
|
||||
),
|
||||
(
|
||||
"existing_rows_count".to_string(),
|
||||
apache_avro::types::Value::Union(0, Box::new(apache_avro::types::Value::Null)),
|
||||
),
|
||||
(
|
||||
"deleted_rows_count".to_string(),
|
||||
apache_avro::types::Value::Union(0, Box::new(apache_avro::types::Value::Null)),
|
||||
),
|
||||
]))
|
||||
.expect("manifest list record should append");
|
||||
}
|
||||
writer.into_inner().expect("manifest list avro bytes should flush")
|
||||
}
|
||||
|
||||
pub(crate) fn manifest_avro_bytes(files: &[(&str, i32, i32, i64, i64)]) -> Vec<u8> {
|
||||
manifest_avro_bytes_with_partition_spec(files, None)
|
||||
}
|
||||
|
||||
pub(crate) fn manifest_avro_bytes_with_partition_spec(
|
||||
files: &[(&str, i32, i32, i64, i64)],
|
||||
partition_spec_id: Option<i32>,
|
||||
) -> Vec<u8> {
|
||||
let schema = apache_avro::Schema::parse_str(
|
||||
r#"
|
||||
{
|
||||
@@ -235,7 +141,6 @@ pub(crate) fn manifest_avro_bytes_with_partition_spec(
|
||||
"fields": [
|
||||
{"name": "content", "type": "int"},
|
||||
{"name": "file_path", "type": "string"},
|
||||
{"name": "partition", "type": {"type": "record", "name": "partition", "fields": []}},
|
||||
{"name": "record_count", "type": "long"},
|
||||
{"name": "file_size_in_bytes", "type": "long"}
|
||||
]
|
||||
@@ -247,11 +152,6 @@ pub(crate) fn manifest_avro_bytes_with_partition_spec(
|
||||
)
|
||||
.expect("manifest avro schema should parse");
|
||||
let mut writer = apache_avro::Writer::new(&schema, Vec::new()).expect("manifest writer should initialize");
|
||||
if let Some(partition_spec_id) = partition_spec_id {
|
||||
writer
|
||||
.add_user_metadata("partition-spec-id".to_string(), partition_spec_id.to_string())
|
||||
.expect("manifest partition spec metadata should write");
|
||||
}
|
||||
for (file_path, content, status, snapshot_id, sequence_number) in files {
|
||||
writer
|
||||
.append_value(apache_avro::types::Value::Record(vec![
|
||||
@@ -264,7 +164,6 @@ pub(crate) fn manifest_avro_bytes_with_partition_spec(
|
||||
apache_avro::types::Value::Record(vec![
|
||||
("content".to_string(), apache_avro::types::Value::Int(*content)),
|
||||
("file_path".to_string(), apache_avro::types::Value::String((*file_path).to_string())),
|
||||
("partition".to_string(), apache_avro::types::Value::Record(Vec::new())),
|
||||
("record_count".to_string(), apache_avro::types::Value::Long(1)),
|
||||
("file_size_in_bytes".to_string(), apache_avro::types::Value::Long(1)),
|
||||
]),
|
||||
@@ -301,7 +200,6 @@ pub(crate) fn manifest_avro_bytes_with_nullable_sequences(files: &[(&str, i32, i
|
||||
"fields": [
|
||||
{"name": "content", "type": "int"},
|
||||
{"name": "file_path", "type": "string"},
|
||||
{"name": "partition", "type": {"type": "record", "name": "partition", "fields": []}},
|
||||
{"name": "record_count", "type": "long"},
|
||||
{"name": "file_size_in_bytes", "type": "long"}
|
||||
]
|
||||
@@ -325,7 +223,6 @@ pub(crate) fn manifest_avro_bytes_with_nullable_sequences(files: &[(&str, i32, i
|
||||
apache_avro::types::Value::Record(vec![
|
||||
("content".to_string(), apache_avro::types::Value::Int(*content)),
|
||||
("file_path".to_string(), apache_avro::types::Value::String((*file_path).to_string())),
|
||||
("partition".to_string(), apache_avro::types::Value::Record(Vec::new())),
|
||||
("record_count".to_string(), apache_avro::types::Value::Long(1)),
|
||||
("file_size_in_bytes".to_string(), apache_avro::types::Value::Long(1)),
|
||||
]),
|
||||
@@ -387,7 +284,7 @@ pub(crate) struct BlockingObjectPublication {
|
||||
backend: TestCatalogObjectBackend,
|
||||
object: String,
|
||||
started: Arc<tokio::sync::Notify>,
|
||||
guard: Arc<parking_lot::Mutex<Option<TableCatalogLockGuard>>>,
|
||||
guard: Arc<parking_lot::Mutex<Option<Box<dyn Send>>>>,
|
||||
}
|
||||
|
||||
impl BlockingObjectPublication {
|
||||
@@ -915,7 +812,7 @@ impl TableCatalogObjectBackend for TestCatalogObjectBackend {
|
||||
.collect())
|
||||
}
|
||||
|
||||
async fn acquire_write_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<TableCatalogLockGuard> {
|
||||
async fn acquire_write_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<Box<dyn Send>> {
|
||||
self.lock_attempts.lock().await.push((bucket.to_string(), object.to_string()));
|
||||
{
|
||||
let mut state = self.state.lock().await;
|
||||
@@ -931,10 +828,10 @@ impl TableCatalogObjectBackend for TestCatalogObjectBackend {
|
||||
.or_insert_with(|| std::sync::Arc::new(tokio::sync::RwLock::new(())))
|
||||
.clone()
|
||||
};
|
||||
Ok(TableCatalogLockGuard::stable(lock.write_owned().await))
|
||||
Ok(Box::new(lock.write_owned().await))
|
||||
}
|
||||
|
||||
async fn acquire_read_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<TableCatalogLockGuard> {
|
||||
async fn acquire_read_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult<Box<dyn Send>> {
|
||||
// The admin fake implemented only acquire_write_lock, so the trait's
|
||||
// default read->write delegation made read acquisitions observable in
|
||||
// lock_attempts as well; keep that (backlog#1837 PR2).
|
||||
@@ -953,7 +850,7 @@ impl TableCatalogObjectBackend for TestCatalogObjectBackend {
|
||||
.or_insert_with(|| std::sync::Arc::new(tokio::sync::RwLock::new(())))
|
||||
.clone()
|
||||
};
|
||||
Ok(TableCatalogLockGuard::stable(lock.read_owned().await))
|
||||
Ok(Box::new(lock.read_owned().await))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1268,8 +1165,6 @@ pub(crate) struct TestTableCatalogStore {
|
||||
pub(crate) fail_put_table_bucket: tokio::sync::Mutex<bool>,
|
||||
pub(crate) register_table_pause: Option<TestCatalogPublishPause>,
|
||||
pub(crate) commit_table_pause: Option<TestCatalogPublishPause>,
|
||||
pub(crate) create_view_pause: Option<TestCatalogPublishPause>,
|
||||
pub(crate) replace_view_pause: Option<TestCatalogPublishPause>,
|
||||
}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
@@ -1578,10 +1473,6 @@ impl crate::table_catalog::TableCatalogStore for TestTableCatalogStore {
|
||||
entry.table_bucket, entry.namespace
|
||||
)));
|
||||
}
|
||||
if let Some(pause) = &self.create_view_pause {
|
||||
pause.started.notify_one();
|
||||
pause.release.notified().await;
|
||||
}
|
||||
self.views.lock().await.push(entry);
|
||||
Ok(())
|
||||
}
|
||||
@@ -1640,10 +1531,6 @@ impl crate::table_catalog::TableCatalogStore for TestTableCatalogStore {
|
||||
"current view metadata location does not match expected location".to_string(),
|
||||
));
|
||||
}
|
||||
if let Some(pause) = &self.replace_view_pause {
|
||||
pause.started.notify_one();
|
||||
pause.release.notified().await;
|
||||
}
|
||||
let mut next = current;
|
||||
next.metadata_location = request.new_metadata_location;
|
||||
next.version_token = "token-view-committed".to_string();
|
||||
|
||||
+166
-1673
File diff suppressed because it is too large
Load Diff
@@ -132,8 +132,18 @@ def failure_probe_plan(warehouse: str, namespace: str, table: str, rest_path: st
|
||||
"expected-version-token": "stale-token-from-previous-load",
|
||||
"expected-metadata-location": "current-metadata-location-from-load-table",
|
||||
"new-metadata-location": f"s3://{warehouse}/tables/table-id/metadata/conflict_probe.metadata.json",
|
||||
"requirements": [],
|
||||
"updates": [],
|
||||
"requirements": [
|
||||
{
|
||||
"type": "assert-current-snapshot-id",
|
||||
"snapshot-id": 0,
|
||||
}
|
||||
],
|
||||
"updates": [
|
||||
{
|
||||
"action": "set-current-schema",
|
||||
"schema-id": 0,
|
||||
}
|
||||
],
|
||||
},
|
||||
),
|
||||
probe_step(
|
||||
@@ -147,8 +157,6 @@ def failure_probe_plan(warehouse: str, namespace: str, table: str, rest_path: st
|
||||
"expected-version-token": "current-version-token-from-load-table",
|
||||
"expected-metadata-location": "current-metadata-location-from-load-table",
|
||||
"new-metadata-location": f"s3://{warehouse}/tables/table-id/metadata/does_not_exist.metadata.json",
|
||||
"requirements": [],
|
||||
"updates": [],
|
||||
},
|
||||
),
|
||||
probe_step(
|
||||
|
||||
@@ -48,8 +48,6 @@ class FailureCoverageTest(unittest.TestCase):
|
||||
self.assertIn("expected-version-token", by_name["stale-token-commit-conflict"]["body"])
|
||||
self.assertIn("expected-metadata-location", by_name["stale-token-commit-conflict"]["body"])
|
||||
self.assertIn("new-metadata-location", by_name["stale-token-commit-conflict"]["body"])
|
||||
self.assertEqual(by_name["stale-token-commit-conflict"]["body"]["requirements"], [])
|
||||
self.assertEqual(by_name["stale-token-commit-conflict"]["body"]["updates"], [])
|
||||
self.assertNotIn("base", by_name["stale-token-commit-conflict"]["body"])
|
||||
self.assertEqual(
|
||||
by_name["diagnostics-after-finalization-gap"]["path"],
|
||||
@@ -58,8 +56,6 @@ class FailureCoverageTest(unittest.TestCase):
|
||||
self.assertEqual(by_name["diagnostics-after-finalization-gap"]["method"], "GET")
|
||||
self.assertEqual(by_name["recovery-repairs-idempotency-index"]["method"], "POST")
|
||||
self.assertIn("does_not_exist.metadata.json", json.dumps(by_name["missing-metadata-object-rejected"]))
|
||||
self.assertEqual(by_name["missing-metadata-object-rejected"]["body"]["requirements"], [])
|
||||
self.assertEqual(by_name["missing-metadata-object-rejected"]["body"]["updates"], [])
|
||||
self.assertNotIn("base", by_name["missing-metadata-object-rejected"]["body"])
|
||||
|
||||
def test_cli_prints_failure_matrix_and_probe_plan(self) -> None:
|
||||
|
||||
Reference in New Issue
Block a user