mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-02 05:14:35 +08:00
feat(integrity): add inventory, audit, and protected migration (#8065)
* feat(integrity): add inventory, audit, and protected migration * refactor(admin): use gateway facade for integrity handlers * fix(integrity): sort fingerprint metadata explicitly
This commit is contained in:
@@ -504,6 +504,13 @@ pub mod notification {
|
||||
};
|
||||
}
|
||||
|
||||
pub mod integrity {
|
||||
pub use crate::services::integrity::{
|
||||
IntegrityError, InventoryItem, InventoryPage, ItemRequest, ItemResult, ItemState, Job, JobMode, JobRequest, JobState,
|
||||
Protection, Readiness, control_job, create_job, get_job, inventory, readiness, resume_job,
|
||||
};
|
||||
}
|
||||
|
||||
pub mod object {
|
||||
pub use crate::object_api::{
|
||||
BLOCK_SIZE_V2, ERASURE_ALGORITHM, EncryptionResolutionError, EncryptionResolutionErrorKind, GetObjectBodyCacheHook,
|
||||
|
||||
@@ -1291,6 +1291,28 @@ pub(crate) async fn get_object_lock_config_and_incarnation_from_disk_in(
|
||||
}
|
||||
}
|
||||
|
||||
/// Inspect all migration-relevant settings while the caller holds an Object Lock
|
||||
/// snapshot's lifecycle and metadata transaction guards. Cached settings are not
|
||||
/// sufficient to authorize a direct storage writer that cannot apply S3 defaults.
|
||||
pub(crate) async fn integrity_migration_metadata_in(
|
||||
ctx: &crate::runtime::instance::InstanceContext,
|
||||
bucket: &str,
|
||||
) -> Result<Arc<BucketMetadata>> {
|
||||
let sys = bucket_metadata_sys_of(ctx)?.read().await.clone();
|
||||
match sys
|
||||
.read_authoritative_metadata_from_disk_under_transaction_lock(bucket)
|
||||
.await?
|
||||
{
|
||||
BucketMetadataAuthority::Authoritative(metadata)
|
||||
if metadata.bucket_incarnation_sidecar && !metadata.bucket_incarnation_id.is_nil() =>
|
||||
{
|
||||
Ok(metadata)
|
||||
}
|
||||
BucketMetadataAuthority::MissingBucket => Err(Error::BucketNotFound(bucket.to_string())),
|
||||
_ => Err(Error::other("migration requires authoritative bucket metadata")),
|
||||
}
|
||||
}
|
||||
|
||||
/// Re-read the quota configuration and bucket incarnation from the same
|
||||
/// authoritative metadata blob while the caller holds the bucket metadata
|
||||
/// transaction read lock.
|
||||
|
||||
@@ -935,6 +935,8 @@ pub struct ObjectOptions {
|
||||
/// Persisted bucket incarnation observed before authorization.
|
||||
pub expected_bucket_incarnation_id: Option<Uuid>,
|
||||
pub no_lock: bool,
|
||||
/// Internal read-only inspection must not enqueue metadata or payload repairs.
|
||||
pub suppress_read_repair: bool,
|
||||
/// Control-plane writers that immediately read or CAS the same namespace
|
||||
/// key use TailDrained without changing namespace lock ownership.
|
||||
#[doc(hidden)]
|
||||
|
||||
@@ -0,0 +1,290 @@
|
||||
// Copyright 2024 RustFS Team
|
||||
// Licensed under the Apache License, Version 2.0.
|
||||
|
||||
//! Explicit, bounded integrity operations. This service does not participate in ordinary GET or Heal.
|
||||
//! Lock order: cluster worker lock, bucket incarnation fence, object commit lock.
|
||||
//! Control requests use checkpoint CAS without taking the worker lock; a worker that loses CAS stops.
|
||||
|
||||
mod model;
|
||||
mod runner;
|
||||
#[cfg(test)]
|
||||
mod tests;
|
||||
|
||||
pub use model::{
|
||||
InventoryItem, InventoryPage, ItemRequest, ItemResult, ItemState, Job, JobMode, JobRequest, JobState, Protection, Readiness,
|
||||
readiness,
|
||||
};
|
||||
|
||||
use crate::config::com::{read_config_limited_preserve_empty_with_metadata, save_config_with_opts_quiet};
|
||||
use crate::disk::RUSTFS_META_BUCKET;
|
||||
use crate::error::Error as StorageError;
|
||||
use crate::object_api::{ObjectOptions, WriteCompletion};
|
||||
use crate::storage_api_contracts::{
|
||||
list::ListOperations,
|
||||
namespace::NamespaceLocking,
|
||||
object::{HTTPPreconditions, ObjectOperations},
|
||||
};
|
||||
use crate::store::ECStore;
|
||||
use model::InventoryItem as Item;
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
use uuid::Uuid;
|
||||
|
||||
type Result<T> = std::result::Result<T, IntegrityError>;
|
||||
|
||||
#[derive(Debug, thiserror::Error)]
|
||||
pub enum IntegrityError {
|
||||
#[error("invalid integrity request: {0}")]
|
||||
Invalid(&'static str),
|
||||
#[error("integrity job or target changed concurrently")]
|
||||
Conflict,
|
||||
#[error("an integrity worker is already active")]
|
||||
Busy,
|
||||
#[error("integrity job not found")]
|
||||
NotFound,
|
||||
#[error("integrity operation was stopped")]
|
||||
Stopped,
|
||||
#[error("protected writes are not enabled on this coordinator")]
|
||||
NotActivated,
|
||||
#[error("unsupported migration bucket configuration")]
|
||||
UnsupportedBucket,
|
||||
#[error("integrity storage operation failed: {0}")]
|
||||
Storage(#[from] StorageError),
|
||||
#[error("integrity stream operation failed: {0}")]
|
||||
Io(#[from] std::io::Error),
|
||||
#[error("integrity checkpoint is malformed: {0}")]
|
||||
Json(#[from] serde_json::Error),
|
||||
}
|
||||
|
||||
fn check_bucket(bucket: &str) -> Result<()> {
|
||||
crate::bucket::utils::check_valid_bucket_name_strict(bucket).map_err(|_| IntegrityError::Invalid("invalid source bucket"))
|
||||
}
|
||||
|
||||
fn path(bucket: &str, id: Uuid) -> String {
|
||||
format!("buckets/{bucket}/integrity-jobs/{id}.json")
|
||||
}
|
||||
|
||||
pub async fn inventory(
|
||||
api: Arc<ECStore>,
|
||||
bucket: &str,
|
||||
prefix: &str,
|
||||
key_marker: Option<String>,
|
||||
version_marker: Option<String>,
|
||||
limit: i32,
|
||||
) -> Result<InventoryPage> {
|
||||
check_bucket(bucket)?;
|
||||
if !(1..=100).contains(&limit) || prefix.len() > 1024 || (version_marker.is_some() && key_marker.is_none()) {
|
||||
return Err(IntegrityError::Invalid("invalid inventory page"));
|
||||
}
|
||||
let incarnation = api.bucket_incarnation_id_from_disk(bucket).await?;
|
||||
let page = api
|
||||
.clone()
|
||||
.list_object_versions(bucket, prefix, key_marker.clone(), version_marker.clone(), None, limit)
|
||||
.await?;
|
||||
if page.is_truncated && (page.next_marker.clone(), page.next_version_idmarker.clone()) == (key_marker, version_marker) {
|
||||
return Err(IntegrityError::Invalid("inventory cursor did not advance"));
|
||||
}
|
||||
let mut items = Vec::with_capacity(page.objects.len());
|
||||
for listed in page.objects {
|
||||
if listed.delete_marker {
|
||||
items.push(Item::from_info(&listed));
|
||||
continue;
|
||||
}
|
||||
let opts = ObjectOptions {
|
||||
version_id: Some(listed.version_id.unwrap_or_else(Uuid::nil).to_string()),
|
||||
include_part_checksums: true,
|
||||
suppress_read_repair: true,
|
||||
..Default::default()
|
||||
};
|
||||
match api.get_object_info(bucket, &listed.name, &opts).await {
|
||||
Ok(info) => items.push(Item::from_info(&info)),
|
||||
Err(_) => {
|
||||
let mut item = Item::from_info(&listed);
|
||||
item.observation_error = Some("metadata_unavailable_or_changed".into());
|
||||
item.protection = Protection::Unknown;
|
||||
item.audit_unsupported = Some("metadata_unavailable_or_changed".into());
|
||||
items.push(item);
|
||||
}
|
||||
}
|
||||
}
|
||||
if api.bucket_incarnation_id_from_disk(bucket).await? != incarnation {
|
||||
return Err(IntegrityError::Conflict);
|
||||
}
|
||||
Ok(InventoryPage {
|
||||
observed_at: time::OffsetDateTime::now_utc().to_string(),
|
||||
bucket_incarnation: incarnation,
|
||||
items,
|
||||
is_truncated: page.is_truncated,
|
||||
next_key_marker: page.next_marker,
|
||||
next_version_marker: page.next_version_idmarker,
|
||||
})
|
||||
}
|
||||
|
||||
pub async fn create_job(api: Arc<ECStore>, bucket: &str, request: JobRequest) -> Result<Job> {
|
||||
check_bucket(bucket)?;
|
||||
request.validate()?;
|
||||
if request.mode == JobMode::Migrate && !readiness().new_writes_enabled_here {
|
||||
return Err(IntegrityError::NotActivated);
|
||||
}
|
||||
let incarnation = api.bucket_incarnation_id_from_disk(bucket).await?;
|
||||
let job = Job {
|
||||
format_version: 1,
|
||||
id: Uuid::new_v4(),
|
||||
bucket: bucket.into(),
|
||||
bucket_incarnation: incarnation,
|
||||
results: vec![ItemResult::default(); request.items.len()],
|
||||
request,
|
||||
state: JobState::Paused,
|
||||
revision: 0,
|
||||
last_error: None,
|
||||
};
|
||||
save(&api, &job, None, None).await?;
|
||||
Ok(job)
|
||||
}
|
||||
|
||||
pub async fn get_job(api: Arc<ECStore>, bucket: &str, id: Uuid) -> Result<Job> {
|
||||
Ok(load(&api, bucket, id).await?.0)
|
||||
}
|
||||
|
||||
async fn load(api: &Arc<ECStore>, bucket: &str, id: Uuid) -> Result<(Job, String)> {
|
||||
check_bucket(bucket)?;
|
||||
let (bytes, info) = match read_config_limited_preserve_empty_with_metadata(api.clone(), &path(bucket, id), 1024 * 1024).await
|
||||
{
|
||||
Ok(value) => value,
|
||||
Err(StorageError::ConfigNotFound | StorageError::FileNotFound) => return Err(IntegrityError::NotFound),
|
||||
Err(error) => return Err(error.into()),
|
||||
};
|
||||
let job: Job = serde_json::from_slice(&bytes)?;
|
||||
job.request.validate()?;
|
||||
if job.format_version != 1
|
||||
|| job.bucket != bucket
|
||||
|| job.id != id
|
||||
|| job.results.len() != job.request.items.len()
|
||||
|| api.bucket_incarnation_id_from_disk(bucket).await? != job.bucket_incarnation
|
||||
{
|
||||
return Err(IntegrityError::Conflict);
|
||||
}
|
||||
Ok((job, info.etag.ok_or(IntegrityError::Conflict)?))
|
||||
}
|
||||
|
||||
async fn save(
|
||||
api: &Arc<ECStore>,
|
||||
job: &Job,
|
||||
etag: Option<&str>,
|
||||
worker: Option<&rustfs_lock::NamespaceLockGuard>,
|
||||
) -> Result<String> {
|
||||
let fence = api
|
||||
.acquire_bucket_incarnation_fence(&job.bucket, job.bucket_incarnation)
|
||||
.await?;
|
||||
let mut opts = ObjectOptions {
|
||||
max_parity: true,
|
||||
write_completion: WriteCompletion::TailDrained,
|
||||
http_preconditions: Some(match etag {
|
||||
Some(value) => HTTPPreconditions {
|
||||
if_match: Some(value.into()),
|
||||
..Default::default()
|
||||
},
|
||||
None => HTTPPreconditions {
|
||||
if_none_match: Some("*".into()),
|
||||
..Default::default()
|
||||
},
|
||||
}),
|
||||
..Default::default()
|
||||
};
|
||||
fence.attach_to_object_options(&mut opts);
|
||||
if let Some(guard) = worker {
|
||||
opts.add_namespace_lock_guard(guard);
|
||||
}
|
||||
let api = api.clone();
|
||||
let job = job.clone();
|
||||
// Keep the lifecycle guard through the drained commit even if the HTTP waiter is cancelled.
|
||||
tokio::spawn(async move {
|
||||
save_config_with_opts_quiet(api.clone(), &path(&job.bucket, job.id), serde_json::to_vec(&job)?, &opts).await?;
|
||||
let (stored, next) = load(&api, &job.bucket, job.id).await?;
|
||||
if stored.revision != job.revision {
|
||||
return Err(IntegrityError::Conflict);
|
||||
}
|
||||
drop(fence);
|
||||
Ok(next)
|
||||
})
|
||||
.await
|
||||
.map_err(|_| IntegrityError::Invalid("checkpoint commit task failed"))?
|
||||
}
|
||||
|
||||
/// Pause/cancel acknowledge intent. An in-flight conditional publication is allowed to drain.
|
||||
pub async fn control_job(api: Arc<ECStore>, bucket: &str, id: Uuid, cancel: bool) -> Result<Job> {
|
||||
let (mut job, etag) = load(&api, bucket, id).await?;
|
||||
match job.state {
|
||||
JobState::Complete | JobState::Cancelled | JobState::CancelRequested => return Ok(job),
|
||||
JobState::Paused | JobState::Failed => job.state = if cancel { JobState::Cancelled } else { JobState::Paused },
|
||||
_ => {
|
||||
job.state = if cancel {
|
||||
JobState::CancelRequested
|
||||
} else {
|
||||
JobState::PauseRequested
|
||||
}
|
||||
}
|
||||
}
|
||||
job.revision = job.revision.checked_add(1).ok_or(IntegrityError::Conflict)?;
|
||||
save(&api, &job, Some(&etag), None).await?;
|
||||
Ok(job)
|
||||
}
|
||||
|
||||
/// Resume is explicit after restart. The distributed lock excludes another worker, including another node.
|
||||
pub async fn resume_job(api: Arc<ECStore>, bucket: &str, id: Uuid) -> Result<Job> {
|
||||
check_bucket(bucket)?;
|
||||
let lock = api.new_ns_lock(RUSTFS_META_BUCKET, "integrity-jobs/worker.lock").await?;
|
||||
let guard = lock
|
||||
.get_write_lock_quiet(Duration::from_secs(1))
|
||||
.await
|
||||
.map_err(|_| IntegrityError::Busy)?;
|
||||
let (mut job, etag) = load(&api, bucket, id).await?;
|
||||
if matches!(job.state, JobState::Complete | JobState::Cancelled) {
|
||||
return Ok(job);
|
||||
}
|
||||
if job.state == JobState::CancelRequested {
|
||||
job.state = JobState::Cancelled;
|
||||
job.revision = job.revision.checked_add(1).ok_or(IntegrityError::Conflict)?;
|
||||
save(&api, &job, Some(&etag), Some(&guard)).await?;
|
||||
return Ok(job);
|
||||
}
|
||||
if job.request.mode == JobMode::Migrate && !readiness().new_writes_enabled_here {
|
||||
return Err(IntegrityError::NotActivated);
|
||||
}
|
||||
job.state = JobState::Running;
|
||||
job.last_error = None;
|
||||
job.revision = job.revision.checked_add(1).ok_or(IntegrityError::Conflict)?;
|
||||
let etag = save(&api, &job, Some(&etag), Some(&guard)).await?;
|
||||
let response = job.clone();
|
||||
tokio::spawn(async move {
|
||||
let result = runner::run(&api, &mut job, etag, &guard).await;
|
||||
if let Err(error) = result {
|
||||
// A failed write/CAS stops publication. Never overwrite a concurrent control request.
|
||||
if let Ok((mut latest, etag)) = load(&api, &job.bucket, job.id).await {
|
||||
match latest.state {
|
||||
JobState::PauseRequested => latest.state = JobState::Paused,
|
||||
JobState::CancelRequested => latest.state = JobState::Cancelled,
|
||||
JobState::Running if latest.revision <= job.revision => {
|
||||
latest.state = JobState::Failed;
|
||||
latest.last_error = Some(
|
||||
match error {
|
||||
IntegrityError::Conflict => "concurrent_change",
|
||||
IntegrityError::Stopped => "worker_stopped",
|
||||
IntegrityError::UnsupportedBucket => "unsupported_destination_bucket_configuration",
|
||||
IntegrityError::NotActivated => "protected_writes_not_activated",
|
||||
_ => "operation_failed_retry_requires_resume",
|
||||
}
|
||||
.into(),
|
||||
);
|
||||
}
|
||||
_ => return,
|
||||
}
|
||||
if let Some(revision) = latest.revision.checked_add(1) {
|
||||
latest.revision = revision;
|
||||
let _ = save(&api, &latest, Some(&etag), Some(&guard)).await;
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
Ok(response)
|
||||
}
|
||||
@@ -0,0 +1,343 @@
|
||||
// Copyright 2024 RustFS Team
|
||||
// Licensed under the Apache License, Version 2.0.
|
||||
|
||||
use super::{IntegrityError, Result};
|
||||
use crate::object_api::ObjectInfo;
|
||||
use rustfs_rio::{Checksum, read_checksums};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use sha2::{Digest, Sha256};
|
||||
use std::collections::{BTreeMap, HashSet};
|
||||
use uuid::Uuid;
|
||||
|
||||
pub const MAX_ITEMS: usize = 64;
|
||||
pub const MAX_OBJECT_BYTES: u64 = 5 * 1024 * 1024 * 1024;
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct ItemRequest {
|
||||
pub key: String,
|
||||
pub version_id: Option<String>,
|
||||
/// Base64 SHA-256 supplied by the administrator from an independent source.
|
||||
pub expected_sha256: Option<String>,
|
||||
pub target_key: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum JobMode {
|
||||
Audit,
|
||||
Migrate,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct JobRequest {
|
||||
pub mode: JobMode,
|
||||
pub items: Vec<ItemRequest>,
|
||||
/// Bounds each source read, destination PUT, and verification stream; one worker runs per cluster.
|
||||
pub bytes_per_second: u64,
|
||||
pub max_object_bytes: u64,
|
||||
}
|
||||
|
||||
impl JobRequest {
|
||||
pub(crate) fn validate(&self) -> Result<()> {
|
||||
if self.items.is_empty() || self.items.len() > MAX_ITEMS {
|
||||
return Err(IntegrityError::Invalid("a job requires 1..=64 explicit source versions"));
|
||||
}
|
||||
if !(64 * 1024..=1024 * 1024 * 1024).contains(&self.bytes_per_second)
|
||||
|| self.max_object_bytes == 0
|
||||
|| self.max_object_bytes > MAX_OBJECT_BYTES
|
||||
{
|
||||
return Err(IntegrityError::Invalid("invalid byte rate or staging limit"));
|
||||
}
|
||||
let mut targets = HashSet::new();
|
||||
let sources: HashSet<_> = self.items.iter().map(|i| i.key.as_str()).collect();
|
||||
for item in &self.items {
|
||||
validate_key(&item.key)?;
|
||||
if let Some(version) = &item.version_id
|
||||
&& version != "null"
|
||||
&& Uuid::parse_str(version).ok().is_none_or(|id| id.is_nil())
|
||||
{
|
||||
return Err(IntegrityError::Invalid("invalid source version"));
|
||||
}
|
||||
if let Some(value) = &item.expected_sha256 {
|
||||
sha256_bytes(value)?;
|
||||
}
|
||||
match (self.mode, &item.target_key) {
|
||||
(JobMode::Audit, None) => {}
|
||||
(JobMode::Migrate, Some(target)) => {
|
||||
validate_key(target)?;
|
||||
if sources.contains(target.as_str()) || !targets.insert(target) {
|
||||
return Err(IntegrityError::Invalid("targets must be distinct from every source and target"));
|
||||
}
|
||||
}
|
||||
_ => return Err(IntegrityError::Invalid("only migrations require an explicit target key")),
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
fn validate_key(key: &str) -> Result<()> {
|
||||
if key.is_empty() || key.len() > 1024 || key.ends_with('/') || key.contains('\0') {
|
||||
return Err(IntegrityError::Invalid("invalid object key"));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) fn sha256_bytes(value: &str) -> Result<Vec<u8>> {
|
||||
let raw = base64_simd::STANDARD
|
||||
.decode_to_vec(value)
|
||||
.map_err(|_| IntegrityError::Invalid("SHA256 must be canonical Base64"))?;
|
||||
if raw.len() != 32 || base64_simd::STANDARD.encode_to_string(&raw) != value {
|
||||
return Err(IntegrityError::Invalid("SHA256 must be canonical Base64 of 32 bytes"));
|
||||
}
|
||||
Ok(raw)
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum Protection {
|
||||
Unknown,
|
||||
Legacy,
|
||||
IndependentCommitment,
|
||||
InvalidDeclaration,
|
||||
NoPayload,
|
||||
}
|
||||
|
||||
pub(crate) fn protection(info: &ObjectInfo) -> Protection {
|
||||
if info.delete_marker {
|
||||
return Protection::NoPayload;
|
||||
}
|
||||
match rustfs_filemeta::shard_integrity::descriptor_from_metadata(&info.user_defined) {
|
||||
Ok(None) if info.parts.iter().all(|p| p.integrity.is_none()) => Protection::Legacy,
|
||||
Ok(Some(parts))
|
||||
if parts.len() == info.parts.len()
|
||||
&& !parts.is_empty()
|
||||
&& parts
|
||||
.iter()
|
||||
.zip(info.parts.iter())
|
||||
.all(|(declared, part)| part.integrity.as_ref() == Some(declared)) =>
|
||||
{
|
||||
Protection::IndependentCommitment
|
||||
}
|
||||
_ => Protection::InvalidDeclaration,
|
||||
}
|
||||
}
|
||||
|
||||
/// Checksums shown in inventory are not automatically accepted as migration evidence.
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct InventoryItem {
|
||||
pub key: String,
|
||||
/// S3 selector, including the literal "null" for the null version slot.
|
||||
pub version_id: String,
|
||||
pub data_dir: Option<Uuid>,
|
||||
pub size: i64,
|
||||
pub protection: Protection,
|
||||
pub checksums: BTreeMap<String, String>,
|
||||
pub checksum_is_multipart: bool,
|
||||
pub audit_unsupported: Option<String>,
|
||||
pub source_fingerprint: String,
|
||||
pub observation_error: Option<String>,
|
||||
}
|
||||
|
||||
impl InventoryItem {
|
||||
pub(crate) fn from_info(info: &ObjectInfo) -> Self {
|
||||
let (checksums, checksum_is_multipart) = info.checksum.as_ref().map_or_else(Default::default, |b| read_checksums(b, 0));
|
||||
Self {
|
||||
key: info.name.clone(),
|
||||
version_id: info
|
||||
.version_id
|
||||
.filter(|id| !id.is_nil())
|
||||
.map_or_else(|| "null".into(), |id| id.to_string()),
|
||||
data_dir: info.data_dir,
|
||||
size: info.size,
|
||||
protection: protection(info),
|
||||
checksums: checksums.into_iter().collect(),
|
||||
checksum_is_multipart,
|
||||
audit_unsupported: unsupported(info).map(str::to_owned),
|
||||
source_fingerprint: fingerprint(info),
|
||||
observation_error: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn unsupported(info: &ObjectInfo) -> Option<&'static str> {
|
||||
if info.delete_marker {
|
||||
Some("delete_marker")
|
||||
} else if !info.transitioned_object.status.is_empty() {
|
||||
Some("transitioned_object")
|
||||
} else if info.is_encrypted() || info.is_compressed() {
|
||||
Some("transformed_object")
|
||||
} else if info.is_multipart() {
|
||||
Some("multipart_object")
|
||||
} else if protection(info) == Protection::InvalidDeclaration {
|
||||
Some("invalid_protection_declaration")
|
||||
} else {
|
||||
None
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn stored_sha256(info: &ObjectInfo) -> Option<String> {
|
||||
let bytes = info.checksum.as_ref()?;
|
||||
let (sums, multipart) = read_checksums(bytes, 0);
|
||||
let value = sums.get("SHA256")?;
|
||||
let canonical = Checksum::new_from_string("SHA256", value)?.to_bytes(&[]);
|
||||
// The general display decoder tolerates truncated suffixes. Eligibility does not.
|
||||
(!multipart && canonical.as_ref() == bytes.as_ref()).then(|| value.clone())
|
||||
}
|
||||
|
||||
pub(crate) fn fingerprint(info: &ObjectInfo) -> String {
|
||||
let mut hash = Sha256::new();
|
||||
// Ordered metadata includes transformation and protection declarations as well as user metadata.
|
||||
for value in [
|
||||
info.bucket.clone(),
|
||||
info.name.clone(),
|
||||
format!("{:?}", info.version_id),
|
||||
format!("{:?}", info.data_dir),
|
||||
format!("{:?}", info.mod_time),
|
||||
info.size.to_string(),
|
||||
info.actual_size.to_string(),
|
||||
format!("{:?}", info.etag),
|
||||
format!("{:?}", info.checksum),
|
||||
format!("{:?}", info.transitioned_object),
|
||||
info.user_tags.to_string(),
|
||||
format!("{:?}", info.expires),
|
||||
format!("{:?}", info.content_type),
|
||||
format!("{:?}", info.content_encoding),
|
||||
format!("{:?}", info.storage_class),
|
||||
] {
|
||||
hash.update(value.len().to_le_bytes());
|
||||
hash.update(value.as_bytes());
|
||||
}
|
||||
for part in info.parts.iter() {
|
||||
let values = format!(
|
||||
"{}:{}:{}:{:?}:{:?}:{:?}:{:?}",
|
||||
part.number, part.size, part.actual_size, part.etag, part.mod_time, part.index, part.integrity
|
||||
);
|
||||
hash.update(values.len().to_le_bytes());
|
||||
hash.update(values.as_bytes());
|
||||
if let Some(checksums) = &part.checksums {
|
||||
let mut checksums: Vec<_> = checksums.iter().collect();
|
||||
checksums.sort_unstable_by_key(|(key, _)| *key);
|
||||
for (key, value) in checksums {
|
||||
hash.update(key.len().to_le_bytes());
|
||||
hash.update(key.as_bytes());
|
||||
hash.update(value.len().to_le_bytes());
|
||||
hash.update(value.as_bytes());
|
||||
}
|
||||
}
|
||||
}
|
||||
let mut metadata: Vec<_> = info.user_defined.iter().collect();
|
||||
metadata.sort_unstable_by_key(|(key, _)| *key);
|
||||
for (key, value) in metadata {
|
||||
hash.update(key.len().to_le_bytes());
|
||||
hash.update(key.as_bytes());
|
||||
hash.update(value.len().to_le_bytes());
|
||||
hash.update(value.as_bytes());
|
||||
}
|
||||
hex_simd::encode_to_string(hash.finalize(), hex_simd::AsciiCase::Lower)
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum JobState {
|
||||
Paused,
|
||||
Running,
|
||||
PauseRequested,
|
||||
CancelRequested,
|
||||
Cancelled,
|
||||
Complete,
|
||||
Failed,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum ItemState {
|
||||
Pending,
|
||||
Prepared,
|
||||
Verified,
|
||||
Migrated,
|
||||
Unsupported,
|
||||
Unavailable,
|
||||
Mismatch,
|
||||
Stale,
|
||||
Conflict,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct ItemResult {
|
||||
pub state: ItemState,
|
||||
pub source_fingerprint: Option<String>,
|
||||
pub expected_sha256: Option<String>,
|
||||
pub evidence_source: Option<String>,
|
||||
pub source_size: Option<u64>,
|
||||
pub target_version_id: Option<Uuid>,
|
||||
pub detail: Option<String>,
|
||||
}
|
||||
|
||||
impl Default for ItemResult {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
state: ItemState::Pending,
|
||||
source_fingerprint: None,
|
||||
expected_sha256: None,
|
||||
evidence_source: None,
|
||||
source_size: None,
|
||||
target_version_id: None,
|
||||
detail: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct Job {
|
||||
pub format_version: u32,
|
||||
pub id: Uuid,
|
||||
pub bucket: String,
|
||||
pub bucket_incarnation: Uuid,
|
||||
pub request: JobRequest,
|
||||
pub state: JobState,
|
||||
pub revision: u64,
|
||||
pub results: Vec<ItemResult>,
|
||||
pub last_error: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Serialize)]
|
||||
pub struct InventoryPage {
|
||||
pub observed_at: String,
|
||||
pub bucket_incarnation: Uuid,
|
||||
pub items: Vec<InventoryItem>,
|
||||
pub is_truncated: bool,
|
||||
pub next_key_marker: Option<String>,
|
||||
pub next_version_marker: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Serialize)]
|
||||
pub struct Readiness {
|
||||
pub local_reader_supported: bool,
|
||||
pub write_requested: bool,
|
||||
pub fleet_operator_attested: bool,
|
||||
pub new_writes_enabled_here: bool,
|
||||
pub fleet_capability_verified: bool,
|
||||
pub unverified_requirements: [&'static str; 3],
|
||||
}
|
||||
|
||||
pub fn readiness() -> Readiness {
|
||||
let write = rustfs_utils::get_env_bool(rustfs_config::ENV_SHARD_INTEGRITY_WRITE, false);
|
||||
let fleet = rustfs_utils::get_env_bool(rustfs_config::ENV_SHARD_INTEGRITY_FLEET_CONFIRMED, false);
|
||||
Readiness {
|
||||
local_reader_supported: true,
|
||||
write_requested: write,
|
||||
fleet_operator_attested: fleet,
|
||||
new_writes_enabled_here: write && fleet,
|
||||
fleet_capability_verified: false,
|
||||
unverified_requirements: [
|
||||
"peer_reader_writer_and_repair_capabilities",
|
||||
"old_process_storage_admission",
|
||||
"deployment_failure_and_performance_qualification",
|
||||
],
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,468 @@
|
||||
// Copyright 2024 RustFS Team
|
||||
// Licensed under the Apache License, Version 2.0.
|
||||
|
||||
use super::model::{Protection, fingerprint, protection, sha256_bytes, stored_sha256, unsupported};
|
||||
use super::{IntegrityError, ItemState, Job, JobMode, JobState, Result, load, readiness, save};
|
||||
use crate::object_api::{
|
||||
ObjectInfo, ObjectOptions, PutObjReader, ShardIntegrityWriteMode, WriteCompletion, without_get_object_body_cache_hook,
|
||||
};
|
||||
use crate::storage_api_contracts::object::{HTTPPreconditions, ObjectIO, ObjectOperations};
|
||||
use crate::store::ECStore;
|
||||
use rustfs_rio::{Checksum, HashReader};
|
||||
use sha2::{Digest, Sha256};
|
||||
use std::{collections::HashMap, sync::Arc, time::Duration};
|
||||
use tokio::io::{AsyncReadExt, AsyncSeekExt, AsyncWriteExt};
|
||||
|
||||
const MARKER: &str = "integrity-migration-job-v1";
|
||||
|
||||
pub(super) async fn run(
|
||||
api: &Arc<ECStore>,
|
||||
job: &mut Job,
|
||||
mut etag: String,
|
||||
guard: &rustfs_lock::NamespaceLockGuard,
|
||||
) -> Result<()> {
|
||||
without_get_object_body_cache_hook(async {
|
||||
for index in 0..job.results.len() {
|
||||
ensure_running(api, job, guard).await?;
|
||||
if !matches!(
|
||||
job.results[index].state,
|
||||
ItemState::Pending | ItemState::Prepared | ItemState::Unavailable
|
||||
) {
|
||||
continue;
|
||||
}
|
||||
if job.request.mode == JobMode::Migrate {
|
||||
validate_destination(api, job).await?;
|
||||
}
|
||||
if job.results[index].state == ItemState::Prepared
|
||||
&& job.request.mode == JobMode::Migrate
|
||||
&& reconcile_target(api, job, index, guard).await?
|
||||
{
|
||||
etag = checkpoint(api, job, &etag, guard).await?;
|
||||
continue;
|
||||
}
|
||||
let request = job.request.items[index].clone();
|
||||
let opts = ObjectOptions {
|
||||
version_id: request.version_id.as_ref().map(|version| {
|
||||
if version == "null" {
|
||||
uuid::Uuid::nil().to_string()
|
||||
} else {
|
||||
version.clone()
|
||||
}
|
||||
}),
|
||||
include_part_checksums: true,
|
||||
suppress_read_repair: true,
|
||||
..Default::default()
|
||||
};
|
||||
let info = match api.get_object_info(&job.bucket, &request.key, &opts).await {
|
||||
Ok(info) => info,
|
||||
Err(_) => {
|
||||
job.results[index].state = ItemState::Unavailable;
|
||||
job.results[index].detail = Some("source_metadata_unavailable".into());
|
||||
etag = checkpoint(api, job, &etag, guard).await?;
|
||||
continue;
|
||||
}
|
||||
};
|
||||
if let Some(reason) = unsupported(&info) {
|
||||
job.results[index].state = ItemState::Unsupported;
|
||||
job.results[index].detail = Some(reason.into());
|
||||
etag = checkpoint(api, job, &etag, guard).await?;
|
||||
continue;
|
||||
}
|
||||
let identity = fingerprint(&info);
|
||||
if job.results[index]
|
||||
.source_fingerprint
|
||||
.as_ref()
|
||||
.is_some_and(|old| old != &identity)
|
||||
{
|
||||
job.results[index].state = ItemState::Stale;
|
||||
etag = checkpoint(api, job, &etag, guard).await?;
|
||||
continue;
|
||||
}
|
||||
let expected = request.expected_sha256.clone().or_else(|| stored_sha256(&info));
|
||||
let Some(expected) = expected else {
|
||||
job.results[index].state = ItemState::Unsupported;
|
||||
job.results[index].detail = Some("no_canonical_single_part_sha256".into());
|
||||
etag = checkpoint(api, job, &etag, guard).await?;
|
||||
continue;
|
||||
};
|
||||
let size = u64::try_from(info.size).map_err(|_| IntegrityError::Invalid("negative source size"))?;
|
||||
if size > job.request.max_object_bytes {
|
||||
job.results[index].state = ItemState::Unsupported;
|
||||
job.results[index].detail = Some("object_exceeds_staging_budget".into());
|
||||
etag = checkpoint(api, job, &etag, guard).await?;
|
||||
continue;
|
||||
}
|
||||
let result = &mut job.results[index];
|
||||
result.source_fingerprint = Some(identity);
|
||||
result.expected_sha256 = Some(expected.clone());
|
||||
result.evidence_source = Some(
|
||||
if request.expected_sha256.is_some() {
|
||||
"administrator_supplied_sha256"
|
||||
} else {
|
||||
"stored_sha256"
|
||||
}
|
||||
.into(),
|
||||
);
|
||||
result.source_size = Some(size);
|
||||
result.state = ItemState::Prepared;
|
||||
result.detail = None;
|
||||
etag = checkpoint(api, job, &etag, guard).await?;
|
||||
let mut file = match stage(api, job, index, &info, &opts, guard).await {
|
||||
Ok(file) => file,
|
||||
Err(IntegrityError::Conflict) => {
|
||||
job.results[index].state = ItemState::Stale;
|
||||
etag = checkpoint(api, job, &etag, guard).await?;
|
||||
continue;
|
||||
}
|
||||
Err(IntegrityError::Invalid("content_sha256_mismatch")) => {
|
||||
job.results[index].state = ItemState::Mismatch;
|
||||
etag = checkpoint(api, job, &etag, guard).await?;
|
||||
continue;
|
||||
}
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
ensure_running(api, job, guard).await?;
|
||||
if job.request.mode == JobMode::Audit {
|
||||
job.results[index].state = ItemState::Verified;
|
||||
} else {
|
||||
let snapshot = validate_destination(api, job).await?;
|
||||
file.rewind().await?;
|
||||
let raw = sha256_bytes(&expected)?;
|
||||
let expected_hex = hex_simd::encode_to_string(raw, hex_simd::AsciiCase::Lower);
|
||||
let mut reader = HashReader::from_stream(
|
||||
ThrottledFile::new(file, job.request.bytes_per_second),
|
||||
info.size,
|
||||
info.size,
|
||||
None,
|
||||
Some(expected_hex),
|
||||
false,
|
||||
)?;
|
||||
reader.add_non_trailing_checksum(Checksum::new_from_string("SHA256", &expected), false)?;
|
||||
let mut body = PutObjReader::new(reader);
|
||||
let mut metadata: HashMap<String, String> = info
|
||||
.user_defined
|
||||
.iter()
|
||||
.filter(|(key, _)| {
|
||||
key.starts_with("x-amz-meta-")
|
||||
|| matches!(key.as_str(), "cache-control" | "content-disposition" | "content-language" | "expires")
|
||||
})
|
||||
.map(|(key, value)| (key.clone(), value.clone()))
|
||||
.collect();
|
||||
if let Some(content_type) = &info.content_type {
|
||||
metadata.insert("content-type".into(), content_type.clone());
|
||||
}
|
||||
if let Some(content_encoding) = &info.content_encoding {
|
||||
metadata.insert("content-encoding".into(), content_encoding.clone());
|
||||
}
|
||||
if let Some(expires) = info.expires {
|
||||
metadata.insert(
|
||||
"expires".into(),
|
||||
expires
|
||||
.format(&time::format_description::well_known::Rfc3339)
|
||||
.map_err(|_| IntegrityError::Invalid("invalid source expiration"))?,
|
||||
);
|
||||
}
|
||||
if let Some(storage_class) = &info.storage_class {
|
||||
metadata.insert("x-amz-storage-class".into(), storage_class.clone());
|
||||
}
|
||||
if !info.user_tags.is_empty() {
|
||||
metadata.insert(rustfs_utils::http::headers::AMZ_OBJECT_TAGGING.into(), info.user_tags.to_string());
|
||||
}
|
||||
rustfs_utils::http::insert_str(&mut metadata, MARKER, marker(job, index)?);
|
||||
let mut write = ObjectOptions {
|
||||
user_defined: metadata,
|
||||
preserve_delete_marker: true,
|
||||
shard_integrity_write_mode: Some(ShardIntegrityWriteMode::Protected),
|
||||
write_completion: WriteCompletion::TailDrained,
|
||||
http_preconditions: Some(HTTPPreconditions {
|
||||
if_none_match: Some("*".into()),
|
||||
..Default::default()
|
||||
}),
|
||||
..Default::default()
|
||||
};
|
||||
write.add_namespace_lock_guard(guard);
|
||||
write.expected_bucket_incarnation_id = Some(job.bucket_incarnation);
|
||||
snapshot.add_lock_fences(&mut write);
|
||||
write.object_lock_config_snapshot = Some(snapshot.clone());
|
||||
let target = request
|
||||
.target_key
|
||||
.as_deref()
|
||||
.ok_or(IntegrityError::Invalid("missing target"))?;
|
||||
// Even an error may follow a committed PUT. Reconcile before retrying; never overwrite.
|
||||
let put = api.put_object(&job.bucket, target, &mut body, &write).await;
|
||||
drop(write);
|
||||
drop(snapshot);
|
||||
if !reconcile_target(api, job, index, guard).await? {
|
||||
put?;
|
||||
return Err(IntegrityError::Conflict);
|
||||
}
|
||||
}
|
||||
etag = checkpoint(api, job, &etag, guard).await?;
|
||||
}
|
||||
job.state = if job.results.iter().any(|item| item.state == ItemState::Unavailable) {
|
||||
job.last_error = Some("source_unavailable_retry_requires_resume".into());
|
||||
JobState::Failed
|
||||
} else {
|
||||
JobState::Complete
|
||||
};
|
||||
checkpoint(api, job, &etag, guard).await?;
|
||||
Ok(())
|
||||
})
|
||||
.await
|
||||
}
|
||||
|
||||
async fn checkpoint(api: &Arc<ECStore>, job: &mut Job, etag: &str, guard: &rustfs_lock::NamespaceLockGuard) -> Result<String> {
|
||||
ensure_running(api, job, guard).await?;
|
||||
job.revision = job.revision.checked_add(1).ok_or(IntegrityError::Conflict)?;
|
||||
save(api, job, Some(etag), Some(guard)).await
|
||||
}
|
||||
|
||||
async fn ensure_running(api: &Arc<ECStore>, job: &Job, guard: &rustfs_lock::NamespaceLockGuard) -> Result<()> {
|
||||
if guard.lock_lost_signal().is_some_and(|signal| signal.is_lost()) {
|
||||
return Err(IntegrityError::Stopped);
|
||||
}
|
||||
let (current, _) = load(api, &job.bucket, job.id).await?;
|
||||
if current.state != JobState::Running {
|
||||
return Err(IntegrityError::Stopped);
|
||||
}
|
||||
if current.revision != job.revision {
|
||||
return Err(IntegrityError::Conflict);
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn stage(
|
||||
api: &Arc<ECStore>,
|
||||
job: &Job,
|
||||
index: usize,
|
||||
info: &ObjectInfo,
|
||||
opts: &ObjectOptions,
|
||||
guard: &rustfs_lock::NamespaceLockGuard,
|
||||
) -> Result<tokio::fs::File> {
|
||||
let mut reader = api
|
||||
.get_object_reader(&job.bucket, &job.request.items[index].key, None, Default::default(), opts)
|
||||
.await?;
|
||||
if fingerprint(&reader.object_info) != fingerprint(info) {
|
||||
return Err(IntegrityError::Conflict);
|
||||
}
|
||||
let std_file = tokio::task::spawn_blocking(tempfile::tempfile)
|
||||
.await
|
||||
.map_err(|_| IntegrityError::Invalid("staging task failed"))??;
|
||||
let mut file = tokio::fs::File::from_std(std_file);
|
||||
let expected = job.results[index]
|
||||
.expected_sha256
|
||||
.as_deref()
|
||||
.ok_or(IntegrityError::Invalid("missing expected digest"))?;
|
||||
copy_verified(api, job, &mut reader.stream, Some(&mut file), info.size, expected, guard).await?;
|
||||
file.flush().await?;
|
||||
let after = api.get_object_info(&job.bucket, &job.request.items[index].key, opts).await?;
|
||||
if fingerprint(&after) != fingerprint(info) {
|
||||
return Err(IntegrityError::Conflict);
|
||||
}
|
||||
Ok(file)
|
||||
}
|
||||
|
||||
async fn copy_verified(
|
||||
api: &Arc<ECStore>,
|
||||
job: &Job,
|
||||
input: &mut (dyn tokio::io::AsyncRead + Unpin + Send + Sync),
|
||||
mut output: Option<&mut tokio::fs::File>,
|
||||
size: i64,
|
||||
expected: &str,
|
||||
guard: &rustfs_lock::NamespaceLockGuard,
|
||||
) -> Result<()> {
|
||||
let expected = sha256_bytes(expected)?;
|
||||
let expected_size = u64::try_from(size).map_err(|_| IntegrityError::Invalid("negative object size"))?;
|
||||
let mut buffer = vec![0u8; 64 * 1024];
|
||||
let mut count = 0u64;
|
||||
let mut hash = Sha256::new();
|
||||
let started = tokio::time::Instant::now();
|
||||
let mut checked = started;
|
||||
loop {
|
||||
if checked.elapsed() >= Duration::from_secs(1) {
|
||||
ensure_running(api, job, guard).await?;
|
||||
checked = tokio::time::Instant::now();
|
||||
}
|
||||
let n = tokio::time::timeout(Duration::from_secs(30), input.read(&mut buffer))
|
||||
.await
|
||||
.map_err(|_| std::io::Error::new(std::io::ErrorKind::TimedOut, "integrity read timed out"))??;
|
||||
if n == 0 {
|
||||
break;
|
||||
}
|
||||
count = count
|
||||
.checked_add(u64::try_from(n).map_err(|_| IntegrityError::Invalid("read size overflow"))?)
|
||||
.ok_or(IntegrityError::Invalid("read size overflow"))?;
|
||||
if count > expected_size || count > job.request.max_object_bytes {
|
||||
return Err(IntegrityError::Invalid("content_sha256_mismatch"));
|
||||
}
|
||||
hash.update(&buffer[..n]);
|
||||
if let Some(file) = output.as_deref_mut() {
|
||||
file.write_all(&buffer[..n]).await?;
|
||||
}
|
||||
let delay = Duration::from_secs_f64(count as f64 / job.request.bytes_per_second as f64);
|
||||
tokio::time::sleep_until(started + delay).await;
|
||||
}
|
||||
if count != expected_size || hash.finalize().as_slice() != expected {
|
||||
return Err(IntegrityError::Invalid("content_sha256_mismatch"));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn marker(job: &Job, index: usize) -> Result<String> {
|
||||
let result = &job.results[index];
|
||||
Ok(format!(
|
||||
"{}:{}:{}:{}",
|
||||
job.id,
|
||||
index,
|
||||
result
|
||||
.source_fingerprint
|
||||
.as_deref()
|
||||
.ok_or(IntegrityError::Invalid("missing source identity"))?,
|
||||
result
|
||||
.expected_sha256
|
||||
.as_deref()
|
||||
.ok_or(IntegrityError::Invalid("missing source digest"))?
|
||||
))
|
||||
}
|
||||
|
||||
async fn reconcile_target(
|
||||
api: &Arc<ECStore>,
|
||||
job: &mut Job,
|
||||
index: usize,
|
||||
guard: &rustfs_lock::NamespaceLockGuard,
|
||||
) -> Result<bool> {
|
||||
let key = job.request.items[index]
|
||||
.target_key
|
||||
.as_deref()
|
||||
.ok_or(IntegrityError::Invalid("missing target"))?;
|
||||
let opts = ObjectOptions {
|
||||
include_part_checksums: true,
|
||||
suppress_read_repair: true,
|
||||
..Default::default()
|
||||
};
|
||||
let target = match api.get_object_info(&job.bucket, key, &opts).await {
|
||||
Ok(info) => info,
|
||||
Err(crate::error::Error::FileNotFound | crate::error::Error::ObjectNotFound(_, _)) => return Ok(false),
|
||||
Err(error) => return Err(error.into()),
|
||||
};
|
||||
let expected_marker = marker(job, index)?;
|
||||
if rustfs_utils::http::get_consistent_str(&target.user_defined, MARKER) != Some(expected_marker.as_str()) {
|
||||
job.results[index].state = ItemState::Conflict;
|
||||
job.results[index].detail = Some("target_owned_by_another_operation".into());
|
||||
return Ok(true);
|
||||
}
|
||||
if protection(&target) != Protection::IndependentCommitment {
|
||||
return Err(IntegrityError::Invalid("target_not_protected"));
|
||||
}
|
||||
let mut reader = api
|
||||
.get_object_reader(&job.bucket, key, None, Default::default(), &opts)
|
||||
.await?;
|
||||
if fingerprint(&target) != fingerprint(&reader.object_info) {
|
||||
return Err(IntegrityError::Conflict);
|
||||
}
|
||||
let size = job.results[index]
|
||||
.source_size
|
||||
.ok_or(IntegrityError::Invalid("missing source size"))?;
|
||||
let expected = job.results[index]
|
||||
.expected_sha256
|
||||
.clone()
|
||||
.ok_or(IntegrityError::Invalid("missing source digest"))?;
|
||||
copy_verified(
|
||||
api,
|
||||
job,
|
||||
&mut reader.stream,
|
||||
None,
|
||||
i64::try_from(size).map_err(|_| IntegrityError::Invalid("source size overflow"))?,
|
||||
&expected,
|
||||
guard,
|
||||
)
|
||||
.await?;
|
||||
let after = api.get_object_info(&job.bucket, key, &opts).await?;
|
||||
if fingerprint(&after) != fingerprint(&target) {
|
||||
return Err(IntegrityError::Conflict);
|
||||
}
|
||||
job.results[index].state = ItemState::Migrated;
|
||||
job.results[index].target_version_id = target.version_id;
|
||||
Ok(true)
|
||||
}
|
||||
|
||||
async fn validate_destination(api: &Arc<ECStore>, job: &Job) -> Result<Arc<crate::object_api::ObjectLockConfigSnapshot>> {
|
||||
if !readiness().new_writes_enabled_here {
|
||||
return Err(IntegrityError::NotActivated);
|
||||
}
|
||||
let snapshot = api.object_lock_config_snapshot(&job.bucket).await?;
|
||||
if snapshot
|
||||
.metadata_transaction_guard_for(api.id, &job.bucket, Some(job.bucket_incarnation))
|
||||
.is_none()
|
||||
{
|
||||
return Err(IntegrityError::Conflict);
|
||||
}
|
||||
let metadata = crate::bucket::metadata_sys::integrity_migration_metadata_in(&api.ctx, &job.bucket).await?;
|
||||
// The direct storage writer is not the S3 orchestration layer. Reject configurations it cannot preserve.
|
||||
if metadata.bucket_incarnation_id != job.bucket_incarnation {
|
||||
return Err(IntegrityError::Conflict);
|
||||
}
|
||||
if !metadata.versioning_config_xml.is_empty()
|
||||
|| !metadata.encryption_config_xml.is_empty()
|
||||
|| !metadata.object_lock_config_xml.is_empty()
|
||||
|| !metadata.replication_config_xml.is_empty()
|
||||
|| !metadata.notification_config_xml.is_empty()
|
||||
|| !metadata.quota_config_json.is_empty()
|
||||
|| !metadata.on_demand_migration_config_json.is_empty()
|
||||
|| !metadata.lifecycle_config_xml.is_empty()
|
||||
|| !metadata.table_bucket_config_json.is_empty()
|
||||
|| !metadata.bucket_acl_config_json.is_empty()
|
||||
|| !metadata.logging_config_xml.is_empty()
|
||||
{
|
||||
return Err(IntegrityError::UnsupportedBucket);
|
||||
}
|
||||
Ok(snapshot)
|
||||
}
|
||||
|
||||
/// Bound publication throughput as well as the source and destination verification reads.
|
||||
struct ThrottledFile {
|
||||
file: tokio::fs::File,
|
||||
started: tokio::time::Instant,
|
||||
count: u64,
|
||||
rate: u64,
|
||||
sleep: std::pin::Pin<Box<tokio::time::Sleep>>,
|
||||
}
|
||||
|
||||
impl ThrottledFile {
|
||||
fn new(file: tokio::fs::File, rate: u64) -> Self {
|
||||
let started = tokio::time::Instant::now();
|
||||
Self {
|
||||
file,
|
||||
started,
|
||||
count: 0,
|
||||
rate,
|
||||
sleep: Box::pin(tokio::time::sleep_until(started)),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl tokio::io::AsyncRead for ThrottledFile {
|
||||
fn poll_read(
|
||||
mut self: std::pin::Pin<&mut Self>,
|
||||
cx: &mut std::task::Context<'_>,
|
||||
output: &mut tokio::io::ReadBuf<'_>,
|
||||
) -> std::task::Poll<std::io::Result<()>> {
|
||||
use std::future::Future;
|
||||
if self.sleep.as_mut().poll(cx).is_pending() {
|
||||
return std::task::Poll::Pending;
|
||||
}
|
||||
let capacity = output.remaining().min(64 * 1024);
|
||||
let mut buffer = tokio::io::ReadBuf::new(output.initialize_unfilled_to(capacity));
|
||||
match std::pin::Pin::new(&mut self.file).poll_read(cx, &mut buffer) {
|
||||
std::task::Poll::Ready(Ok(())) => {
|
||||
let count = buffer.filled().len();
|
||||
output.advance(count);
|
||||
self.count += count as u64;
|
||||
let deadline = self.started + Duration::from_secs_f64(self.count as f64 / self.rate as f64);
|
||||
self.sleep.as_mut().reset(deadline);
|
||||
std::task::Poll::Ready(Ok(()))
|
||||
}
|
||||
result => result,
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,387 @@
|
||||
// Copyright 2024 RustFS Team
|
||||
// Licensed under the Apache License, Version 2.0.
|
||||
|
||||
use super::*;
|
||||
use crate::object_api::{ObjectOptions, PutObjReader, ShardIntegrityWriteMode};
|
||||
use crate::storage_api_contracts::{
|
||||
bucket::{BucketOperations, MakeBucketOptions},
|
||||
object::{ObjectIO, ObjectOperations},
|
||||
};
|
||||
use tokio::io::AsyncReadExt;
|
||||
|
||||
async fn fixture() -> (Vec<tempfile::TempDir>, Arc<crate::store::ECStore>, String) {
|
||||
let (dirs, store) = crate::services::rebalance::test_store_with_persisted_rebalance_meta(Default::default()).await;
|
||||
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
let bucket = format!("integrity-{}", Uuid::new_v4());
|
||||
store
|
||||
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("fixture bucket");
|
||||
(dirs, store, bucket)
|
||||
}
|
||||
|
||||
fn request(mode: JobMode, digest: String) -> JobRequest {
|
||||
JobRequest {
|
||||
mode,
|
||||
items: vec![ItemRequest {
|
||||
key: "source".into(),
|
||||
version_id: Some("null".into()),
|
||||
expected_sha256: Some(digest),
|
||||
target_key: (mode == JobMode::Migrate).then(|| "target".into()),
|
||||
}],
|
||||
bytes_per_second: 1024 * 1024,
|
||||
max_object_bytes: 16 * 1024 * 1024,
|
||||
}
|
||||
}
|
||||
|
||||
async fn finish(store: Arc<crate::store::ECStore>, bucket: &str, id: Uuid) -> Job {
|
||||
tokio::time::timeout(std::time::Duration::from_secs(60), async {
|
||||
loop {
|
||||
let job = get_job(store.clone(), bucket, id).await.expect("durable job");
|
||||
if job.state != JobState::Running {
|
||||
return job;
|
||||
}
|
||||
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("job must finish")
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial(shard_integrity_rollout)]
|
||||
async fn audit_migration_and_ack_recovery_preserve_source_and_reject_conflicts() {
|
||||
temp_env::async_with_vars(
|
||||
[
|
||||
(rustfs_config::ENV_SHARD_INTEGRITY_WRITE, Some("true")),
|
||||
(rustfs_config::ENV_SHARD_INTEGRITY_FLEET_CONFIRMED, Some("true")),
|
||||
],
|
||||
async {
|
||||
let (_dirs, store, bucket) = fixture().await;
|
||||
let payload = vec![42; 256 * 1024 + 17];
|
||||
let expected = rustfs_rio::Checksum::new_from_data(rustfs_rio::ChecksumType::SHA256, &payload)
|
||||
.expect("digest")
|
||||
.encoded;
|
||||
store
|
||||
.put_object(
|
||||
&bucket,
|
||||
"source",
|
||||
&mut PutObjReader::from_vec(payload.clone()),
|
||||
&ObjectOptions {
|
||||
user_defined: [
|
||||
("cache-control".into(), "max-age=300".into()),
|
||||
("content-disposition".into(), "attachment; filename=sample.bin".into()),
|
||||
("expires".into(), "2030-01-01T00:00:00Z".into()),
|
||||
(rustfs_utils::http::headers::AMZ_OBJECT_TAGGING.into(), "kind=archive".into()),
|
||||
]
|
||||
.into(),
|
||||
shard_integrity_write_mode: Some(ShardIntegrityWriteMode::Legacy),
|
||||
write_completion: crate::object_api::WriteCompletion::TailDrained,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("source");
|
||||
let source = store
|
||||
.get_object_info(&bucket, "source", &ObjectOptions::default())
|
||||
.await
|
||||
.expect("source metadata");
|
||||
assert_eq!(source.user_tags.as_str(), "kind=archive", "source fixture tags");
|
||||
let before = model::fingerprint(&source);
|
||||
let page = inventory(store.clone(), &bucket, "", None, None, 1).await.expect("inventory");
|
||||
assert_eq!(page.items.len(), 1);
|
||||
assert_eq!(page.items[0].protection, Protection::Legacy);
|
||||
assert_eq!(page.items[0].version_id, "null", "inventory must round-trip an explicit null selector");
|
||||
let audit = create_job(store.clone(), &bucket, request(JobMode::Audit, expected.clone()))
|
||||
.await
|
||||
.expect("create audit");
|
||||
resume_job(store.clone(), &bucket, audit.id).await.expect("resume audit");
|
||||
let audited = finish(store.clone(), &bucket, audit.id).await;
|
||||
assert_eq!(audited.state, JobState::Complete, "{audited:?}");
|
||||
assert_eq!(audited.results[0].state, ItemState::Verified);
|
||||
|
||||
let migration = create_job(store.clone(), &bucket, request(JobMode::Migrate, expected.clone()))
|
||||
.await
|
||||
.expect("create migration");
|
||||
resume_job(store.clone(), &bucket, migration.id)
|
||||
.await
|
||||
.expect("resume migration");
|
||||
let migrated = finish(store.clone(), &bucket, migration.id).await;
|
||||
assert_eq!(migrated.state, JobState::Complete, "{migrated:?}");
|
||||
assert_eq!(migrated.results[0].state, ItemState::Migrated);
|
||||
let mut reader = store
|
||||
.get_object_reader(&bucket, "target", None, Default::default(), &ObjectOptions::default())
|
||||
.await
|
||||
.expect("target");
|
||||
assert_eq!(model::protection(&reader.object_info), Protection::IndependentCommitment);
|
||||
assert_eq!(reader.object_info.user_tags.as_str(), "kind=archive");
|
||||
assert_eq!(reader.object_info.expires, source.expires);
|
||||
assert_eq!(
|
||||
reader.object_info.user_defined.get("cache-control"),
|
||||
source.user_defined.get("cache-control")
|
||||
);
|
||||
assert_eq!(
|
||||
reader.object_info.user_defined.get("content-disposition"),
|
||||
source.user_defined.get("content-disposition")
|
||||
);
|
||||
let mut actual = Vec::new();
|
||||
reader.stream.read_to_end(&mut actual).await.expect("target content");
|
||||
assert_eq!(actual, payload);
|
||||
let after = store
|
||||
.get_object_info(&bucket, "source", &ObjectOptions::default())
|
||||
.await
|
||||
.expect("source after");
|
||||
assert_eq!(model::fingerprint(&after), before);
|
||||
|
||||
// Simulate a lost completion checkpoint after the create-only PUT committed.
|
||||
let (mut replay, etag) = load(&store, &bucket, migration.id).await.expect("load receipt");
|
||||
replay.state = JobState::Paused;
|
||||
replay.results[0].state = ItemState::Prepared;
|
||||
replay.revision += 1;
|
||||
save(&store, &replay, Some(&etag), None)
|
||||
.await
|
||||
.expect("lost acknowledgement fixture");
|
||||
resume_job(store.clone(), &bucket, replay.id)
|
||||
.await
|
||||
.expect("reconcile after restart");
|
||||
let reconciled = finish(store.clone(), &bucket, replay.id).await;
|
||||
assert_eq!(reconciled.results[0].state, ItemState::Migrated, "{reconciled:?}");
|
||||
|
||||
let conflict = create_job(store.clone(), &bucket, request(JobMode::Migrate, expected))
|
||||
.await
|
||||
.expect("conflicting job");
|
||||
resume_job(store.clone(), &bucket, conflict.id)
|
||||
.await
|
||||
.expect("resume conflict");
|
||||
let conflict = finish(store.clone(), &bucket, conflict.id).await;
|
||||
assert_eq!(conflict.results[0].state, ItemState::Conflict, "{conflict:?}");
|
||||
let target = store
|
||||
.get_object_info(&bucket, "target", &ObjectOptions::default())
|
||||
.await
|
||||
.expect("unchanged target");
|
||||
assert_eq!(target.data_dir, reader.object_info.data_dir);
|
||||
},
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn mismatch_never_publishes_and_paused_job_survives_reload() {
|
||||
let (_dirs, store, bucket) = fixture().await;
|
||||
store
|
||||
.put_object(
|
||||
&bucket,
|
||||
"source",
|
||||
&mut PutObjReader::from_vec(b"actual".to_vec()),
|
||||
&ObjectOptions {
|
||||
shard_integrity_write_mode: Some(ShardIntegrityWriteMode::Legacy),
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("source");
|
||||
let wrong = rustfs_rio::Checksum::new_from_data(rustfs_rio::ChecksumType::SHA256, b"different")
|
||||
.expect("digest")
|
||||
.encoded;
|
||||
let job = create_job(store.clone(), &bucket, request(JobMode::Audit, wrong))
|
||||
.await
|
||||
.expect("job");
|
||||
assert_eq!(get_job(store.clone(), &bucket, job.id).await.expect("reload").state, JobState::Paused);
|
||||
resume_job(store.clone(), &bucket, job.id).await.expect("resume");
|
||||
let job = finish(store.clone(), &bucket, job.id).await;
|
||||
assert_eq!(job.results[0].state, ItemState::Mismatch, "{job:?}");
|
||||
assert!(
|
||||
store
|
||||
.get_object_info(&bucket, "target", &ObjectOptions::default())
|
||||
.await
|
||||
.is_err()
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn fingerprint_is_independent_of_metadata_insertion_order() {
|
||||
let mut a = crate::object_api::ObjectInfo::default();
|
||||
let mut first = std::collections::HashMap::new();
|
||||
first.insert("a".into(), "one".into());
|
||||
first.insert("b".into(), "two".into());
|
||||
a.user_defined = Arc::new(first);
|
||||
let mut b = a.clone();
|
||||
let mut second = std::collections::HashMap::new();
|
||||
second.insert("b".into(), "two".into());
|
||||
second.insert("a".into(), "one".into());
|
||||
b.user_defined = Arc::new(second);
|
||||
assert_eq!(model::fingerprint(&a), model::fingerprint(&b));
|
||||
b.data_dir = Some(Uuid::new_v4());
|
||||
assert_ne!(model::fingerprint(&a), model::fingerprint(&b));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn fingerprint_is_independent_of_part_checksum_insertion_order() {
|
||||
let mut checksums = [
|
||||
("SHA256", "sha256"),
|
||||
("SHA1", "sha1"),
|
||||
("CRC32", "crc32"),
|
||||
("CRC32C", "crc32c"),
|
||||
];
|
||||
let a = crate::object_api::ObjectInfo {
|
||||
parts: Arc::new(vec![rustfs_filemeta::ObjectPartInfo {
|
||||
checksums: Some(checksums.map(|(key, value)| (key.to_string(), value.to_string())).into()),
|
||||
..Default::default()
|
||||
}]),
|
||||
..Default::default()
|
||||
};
|
||||
let mut b = a.clone();
|
||||
let expected = model::fingerprint(&a);
|
||||
for _ in 0..checksums.len() {
|
||||
checksums.rotate_left(1);
|
||||
Arc::make_mut(&mut b.parts)[0].checksums =
|
||||
Some(checksums.map(|(key, value)| (key.to_string(), value.to_string())).into());
|
||||
assert_eq!(expected, model::fingerprint(&b));
|
||||
}
|
||||
Arc::make_mut(&mut b.parts)[0]
|
||||
.checksums
|
||||
.as_mut()
|
||||
.expect("part checksums")
|
||||
.insert("SHA256".into(), "changed".into());
|
||||
assert_ne!(expected, model::fingerprint(&b));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rejects_migration_into_a_source_key() {
|
||||
let request = JobRequest {
|
||||
mode: JobMode::Migrate,
|
||||
items: vec![ItemRequest {
|
||||
key: "source".into(),
|
||||
version_id: None,
|
||||
expected_sha256: None,
|
||||
target_key: Some("source".into()),
|
||||
}],
|
||||
bytes_per_second: 1024 * 1024,
|
||||
max_object_bytes: 1024,
|
||||
};
|
||||
assert!(request.validate().is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rejects_truncated_checksum_suffix_as_migration_evidence() {
|
||||
let checksum = rustfs_rio::Checksum::new_from_data(rustfs_rio::ChecksumType::SHA256, b"payload").expect("checksum");
|
||||
let mut info = crate::object_api::ObjectInfo {
|
||||
checksum: Some(checksum.to_bytes(&[])),
|
||||
..Default::default()
|
||||
};
|
||||
assert_eq!(model::stored_sha256(&info), Some(checksum.encoded));
|
||||
let mut bytes = info.checksum.take().expect("bytes").to_vec();
|
||||
bytes.push(0x80);
|
||||
info.checksum = Some(bytes.into());
|
||||
assert!(model::stored_sha256(&info).is_none());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn unavailable_source_can_resume_and_cancel_intent_cannot_restart() {
|
||||
let (_dirs, store, bucket) = fixture().await;
|
||||
let payload = b"recovered".to_vec();
|
||||
let digest = rustfs_rio::Checksum::new_from_data(rustfs_rio::ChecksumType::SHA256, &payload)
|
||||
.expect("digest")
|
||||
.encoded;
|
||||
let job = create_job(store.clone(), &bucket, request(JobMode::Audit, digest))
|
||||
.await
|
||||
.expect("create");
|
||||
resume_job(store.clone(), &bucket, job.id).await.expect("start");
|
||||
let failed = finish(store.clone(), &bucket, job.id).await;
|
||||
assert_eq!(failed.state, JobState::Failed);
|
||||
assert_eq!(failed.results[0].state, ItemState::Unavailable);
|
||||
store
|
||||
.put_object(
|
||||
&bucket,
|
||||
"source",
|
||||
&mut PutObjReader::from_vec(payload),
|
||||
&ObjectOptions {
|
||||
shard_integrity_write_mode: Some(ShardIntegrityWriteMode::Legacy),
|
||||
write_completion: WriteCompletion::TailDrained,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("recovered source");
|
||||
resume_job(store.clone(), &bucket, job.id).await.expect("retry");
|
||||
let completed = finish(store.clone(), &bucket, job.id).await;
|
||||
assert_eq!(completed.state, JobState::Complete, "{completed:?}");
|
||||
assert_eq!(completed.results[0].state, ItemState::Verified);
|
||||
|
||||
// A process restart may leave the durable cancel intent without a live worker.
|
||||
let (mut interrupted, etag) = load(&store, &bucket, job.id).await.expect("load");
|
||||
interrupted.state = JobState::CancelRequested;
|
||||
interrupted.revision += 1;
|
||||
save(&store, &interrupted, Some(&etag), None).await.expect("persist intent");
|
||||
assert_eq!(
|
||||
control_job(store.clone(), &bucket, job.id, false).await.expect("pause").state,
|
||||
JobState::CancelRequested
|
||||
);
|
||||
assert_eq!(
|
||||
resume_job(store.clone(), &bucket, job.id)
|
||||
.await
|
||||
.expect("recover cancel")
|
||||
.state,
|
||||
JobState::Cancelled
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn pause_and_cancel_are_durable_during_content_read() {
|
||||
let (_dirs, store, bucket) = fixture().await;
|
||||
let payload = vec![7; 512 * 1024];
|
||||
let digest = rustfs_rio::Checksum::new_from_data(rustfs_rio::ChecksumType::SHA256, &payload)
|
||||
.expect("digest")
|
||||
.encoded;
|
||||
store
|
||||
.put_object(
|
||||
&bucket,
|
||||
"source",
|
||||
&mut PutObjReader::from_vec(payload),
|
||||
&ObjectOptions {
|
||||
shard_integrity_write_mode: Some(ShardIntegrityWriteMode::Legacy),
|
||||
write_completion: WriteCompletion::TailDrained,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("source");
|
||||
let mut req = request(JobMode::Audit, digest);
|
||||
req.bytes_per_second = 64 * 1024;
|
||||
let job = create_job(store.clone(), &bucket, req).await.expect("create");
|
||||
resume_job(store.clone(), &bucket, job.id).await.expect("start");
|
||||
tokio::time::timeout(std::time::Duration::from_secs(20), async {
|
||||
loop {
|
||||
if get_job(store.clone(), &bucket, job.id).await.expect("status").results[0].state == ItemState::Prepared {
|
||||
break;
|
||||
}
|
||||
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("prepared");
|
||||
let intent = control_job(store.clone(), &bucket, job.id, false).await.expect("pause");
|
||||
assert_eq!(intent.state, JobState::PauseRequested);
|
||||
tokio::time::timeout(std::time::Duration::from_secs(20), async {
|
||||
loop {
|
||||
let paused = get_job(store.clone(), &bucket, job.id).await.expect("status");
|
||||
if paused.state == JobState::Paused {
|
||||
break;
|
||||
}
|
||||
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("pause acknowledged");
|
||||
assert_eq!(
|
||||
control_job(store.clone(), &bucket, job.id, true).await.expect("cancel").state,
|
||||
JobState::Cancelled
|
||||
);
|
||||
assert_eq!(
|
||||
resume_job(store.clone(), &bucket, job.id)
|
||||
.await
|
||||
.expect("terminal resume")
|
||||
.state,
|
||||
JobState::Cancelled
|
||||
);
|
||||
}
|
||||
@@ -16,6 +16,7 @@
|
||||
|
||||
pub(crate) mod batch_processor;
|
||||
pub(crate) mod event_notification;
|
||||
pub(crate) mod integrity;
|
||||
pub(crate) mod metrics_realtime;
|
||||
pub(crate) mod notification_sys;
|
||||
pub(crate) mod rebalance;
|
||||
|
||||
@@ -13698,6 +13698,7 @@ mod tests {
|
||||
true,
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
GET_OBJECT_PATH_LEGACY_DUPLEX,
|
||||
GET_CODEC_STREAMING_OBJECT_CLASS_PLAIN_SINGLE_PART,
|
||||
metrics_size_bucket,
|
||||
@@ -13811,6 +13812,7 @@ mod tests {
|
||||
true,
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
GET_OBJECT_PATH_LEGACY_DUPLEX,
|
||||
GET_CODEC_STREAMING_OBJECT_CLASS_PLAIN_SINGLE_PART,
|
||||
metrics_size_bucket,
|
||||
|
||||
@@ -2776,6 +2776,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
|
||||
self.set_index,
|
||||
self.pool_index,
|
||||
opts.skip_verify_bitrot,
|
||||
opts.suppress_read_repair,
|
||||
true,
|
||||
true,
|
||||
GET_OBJECT_PATH_LEGACY_DUPLEX,
|
||||
@@ -2805,6 +2806,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
|
||||
self.set_index,
|
||||
self.pool_index,
|
||||
opts.skip_verify_bitrot,
|
||||
opts.suppress_read_repair,
|
||||
true,
|
||||
false,
|
||||
GET_OBJECT_PATH_LEGACY_DUPLEX,
|
||||
@@ -2893,6 +2895,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
|
||||
self.set_index,
|
||||
self.pool_index,
|
||||
opts.skip_verify_bitrot,
|
||||
opts.suppress_read_repair,
|
||||
true,
|
||||
false,
|
||||
GET_OBJECT_PATH_DIRECT_MEMORY,
|
||||
@@ -3063,6 +3066,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
|
||||
let set_index = self.set_index;
|
||||
let pool_index = self.pool_index;
|
||||
let skip_verify = opts.skip_verify_bitrot;
|
||||
let suppress_read_repair = opts.suppress_read_repair;
|
||||
// The producer runs in a separate Tokio task, so carry the caller's
|
||||
// read policy across the task boundary explicitly. Tokio task-local
|
||||
// values are not inherited by spawned tasks.
|
||||
@@ -3093,6 +3097,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
|
||||
set_index,
|
||||
pool_index,
|
||||
skip_verify,
|
||||
suppress_read_repair,
|
||||
false,
|
||||
false,
|
||||
GET_OBJECT_PATH_LEGACY_DUPLEX,
|
||||
@@ -9493,6 +9498,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
||||
let set_index = self.set_index;
|
||||
let pool_index = self.pool_index;
|
||||
let skip_verify = opts.skip_verify_bitrot;
|
||||
let suppress_read_repair = opts.suppress_read_repair;
|
||||
let metrics_size_bucket = rustfs_io_metrics::get_object_size_bucket(cloned_fi.size);
|
||||
let erasure_cache = Arc::clone(&self.erasure_cache);
|
||||
let producer = async move {
|
||||
@@ -9510,6 +9516,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
||||
set_index,
|
||||
pool_index,
|
||||
skip_verify,
|
||||
suppress_read_repair,
|
||||
false,
|
||||
false,
|
||||
GET_OBJECT_PATH_LEGACY_DUPLEX,
|
||||
|
||||
@@ -579,7 +579,7 @@ impl SetDisks {
|
||||
}
|
||||
}
|
||||
metadata_fanout_diagnostics.record_quorum_candidate_latency(metadata_metrics_path, fileinfo_selection_quorum);
|
||||
if errs.iter().any(|err| err.is_some()) {
|
||||
if !opts.suppress_read_repair && errs.iter().any(|err| err.is_some()) {
|
||||
let version_id = resolved_read_repair_version_id(&fi, opts.version_id.as_deref());
|
||||
submit_read_repair_heal(
|
||||
&fi.volume,
|
||||
@@ -848,6 +848,7 @@ impl SetDisks {
|
||||
set_index: usize,
|
||||
pool_index: usize,
|
||||
skip_verify_bitrot: bool,
|
||||
suppress_read_repair: bool,
|
||||
prefer_data_blocks_first_reader_setup: bool,
|
||||
require_reconstruction_surplus: bool,
|
||||
metrics_path: &'static str,
|
||||
@@ -1192,7 +1193,7 @@ impl SetDisks {
|
||||
"Shard availability check"
|
||||
);
|
||||
|
||||
if missing_shards > 0 && available_shards >= erasure.data_shards {
|
||||
if !suppress_read_repair && missing_shards > 0 && available_shards >= erasure.data_shards {
|
||||
// We have missing shards but enough to read - trigger background heal
|
||||
debug!(
|
||||
bucket,
|
||||
@@ -1315,20 +1316,22 @@ impl SetDisks {
|
||||
// is bound to the read-repair reservation, so only the
|
||||
// first sighting within the dedup TTL books a journal
|
||||
// record instead of one per retried read.
|
||||
submit_read_repair_heal_with_submitter(
|
||||
ReadRepairHealSubmission {
|
||||
bucket,
|
||||
object,
|
||||
version_id: version_id.as_deref(),
|
||||
pool_index,
|
||||
set_index,
|
||||
part_number: Some(part_number),
|
||||
reason: "decode_error",
|
||||
mrf_intent: Some((rustfs_common::mrf_channel::MrfKind::DecodeFailure, fi.version_id)),
|
||||
},
|
||||
send_read_repair_heal_request,
|
||||
)
|
||||
.await;
|
||||
if !suppress_read_repair {
|
||||
submit_read_repair_heal_with_submitter(
|
||||
ReadRepairHealSubmission {
|
||||
bucket,
|
||||
object,
|
||||
version_id: version_id.as_deref(),
|
||||
pool_index,
|
||||
set_index,
|
||||
part_number: Some(part_number),
|
||||
reason: "decode_error",
|
||||
mrf_intent: Some((rustfs_common::mrf_channel::MrfKind::DecodeFailure, fi.version_id)),
|
||||
},
|
||||
send_read_repair_heal_request,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
has_err = false;
|
||||
}
|
||||
}
|
||||
@@ -2515,6 +2518,7 @@ mod metadata_cache_tests {
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
GET_OBJECT_PATH_SET_DISK,
|
||||
"plain",
|
||||
"small",
|
||||
@@ -2547,6 +2551,7 @@ mod metadata_cache_tests {
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
GET_OBJECT_PATH_SET_DISK,
|
||||
"plain",
|
||||
"small",
|
||||
@@ -2572,6 +2577,7 @@ mod metadata_cache_tests {
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
GET_OBJECT_PATH_SET_DISK,
|
||||
"plain",
|
||||
"small",
|
||||
@@ -2595,6 +2601,7 @@ mod metadata_cache_tests {
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
GET_OBJECT_PATH_SET_DISK,
|
||||
"plain",
|
||||
"small",
|
||||
@@ -2620,6 +2627,7 @@ mod metadata_cache_tests {
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
GET_OBJECT_PATH_SET_DISK,
|
||||
"plain",
|
||||
"small",
|
||||
@@ -2659,6 +2667,7 @@ mod metadata_cache_tests {
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
GET_OBJECT_PATH_SET_DISK,
|
||||
"plain",
|
||||
"empty",
|
||||
@@ -2693,6 +2702,7 @@ mod metadata_cache_tests {
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
GET_OBJECT_PATH_SET_DISK,
|
||||
"plain",
|
||||
"small",
|
||||
@@ -2992,6 +3002,74 @@ mod metadata_cache_tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn missing_metadata_read_repair_is_suppressed_only_for_read_only_requests() {
|
||||
use crate::object_api::{PutObjReader, ShardIntegrityWriteMode, WriteCompletion};
|
||||
use crate::storage_api_contracts::object::ObjectIO;
|
||||
|
||||
let runtime = tokio::runtime::Builder::new_current_thread()
|
||||
.enable_all()
|
||||
.build()
|
||||
.expect("runtime");
|
||||
let recorder = crate::test_metrics::CapturingRecorder::default();
|
||||
metrics::with_local_recorder(&recorder, || {
|
||||
runtime.block_on(async {
|
||||
let ctx = Arc::new(crate::runtime::instance::InstanceContext::new());
|
||||
let (dirs, set) = crate::ecstore_validation_blackbox::make_local_set_disks_with_ctx(4, 2, ctx).await;
|
||||
let bucket = format!("read-only-{}", Uuid::new_v4());
|
||||
let object = "missing-metadata";
|
||||
for disk in set.disks.read().await.iter().flatten() {
|
||||
disk.make_volume(&bucket).await.expect("bucket volume");
|
||||
}
|
||||
let written = set
|
||||
.put_object(
|
||||
&bucket,
|
||||
object,
|
||||
&mut PutObjReader::from_vec(vec![7; 1024]),
|
||||
&ObjectOptions {
|
||||
shard_integrity_write_mode: Some(ShardIntegrityWriteMode::Legacy),
|
||||
write_completion: WriteCompletion::TailDrained,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("fully committed fixture");
|
||||
let missing = dirs[0].path().join(&bucket).join(object).join("xl.meta");
|
||||
tokio::fs::remove_file(&missing).await.expect("remove one metadata replica");
|
||||
|
||||
// Intercept the production submitter at its dedup admission boundary.
|
||||
// This avoids initializing the process-global heal channel or allowing
|
||||
// a background worker to repair the physical fault under test.
|
||||
let version = written.version_id.map(|id| id.to_string());
|
||||
let reservation = reserve_read_repair_heal(&bucket, object, version.as_deref(), 0, 0)
|
||||
.await
|
||||
.expect("unique object reservation");
|
||||
for (suppress_read_repair, submissions) in [(true, 0), (false, 1), (true, 1)] {
|
||||
set.get_object_fileinfo(
|
||||
&bucket,
|
||||
object,
|
||||
&ObjectOptions {
|
||||
include_part_checksums: true,
|
||||
suppress_read_repair,
|
||||
..Default::default()
|
||||
},
|
||||
false,
|
||||
false,
|
||||
)
|
||||
.await
|
||||
.expect("remaining metadata replicas retain read quorum");
|
||||
assert_eq!(
|
||||
recorder.counter_value("rustfs_heal_read_repair_dedup_total", &[("reason", "duplicate")]),
|
||||
submissions,
|
||||
"only ordinary reads must reach read-repair admission"
|
||||
);
|
||||
assert!(!missing.exists(), "the read-only probe must not restore metadata");
|
||||
}
|
||||
release_read_repair_heal_reservation(&reservation).await;
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn read_repair_heal_dedupes_same_object_version() {
|
||||
let bucket = format!("bucket-{}", Uuid::new_v4());
|
||||
@@ -5055,6 +5133,7 @@ mod tests {
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
GET_OBJECT_PATH_SET_DISK,
|
||||
"test-object-class",
|
||||
"test-size-bucket",
|
||||
|
||||
@@ -8,6 +8,8 @@ collection:
|
||||
|
||||
## Operations
|
||||
|
||||
For legacy protection assessment and bounded protected copies, see [Object integrity inventory, audit, and migration](operations/shard-integrity-audit.md).
|
||||
|
||||
Operational runbooks live under [`operations/`](operations/). Replication
|
||||
operators should start with:
|
||||
|
||||
|
||||
@@ -12,7 +12,7 @@ The broad `rustfs_ecstore::api` facade is a compatibility boundary, not an archi
|
||||
| `storage`, `layout`, `error`, `runtime`, `cluster`, `rpc` | Compatibility spine for storage, topology, runtime handles, cluster control, and internode calls. | Keep until replacement contracts compile in downstream boundary files. |
|
||||
| `bucket` | Domain facade consumed through owner-local `storage_api` boundaries; explicit submodules and symbol lists, never whole bucket owner modules. | Keep lists aligned with boundary consumers; never restore whole-module passthroughs. |
|
||||
| `config`, `disk`, `tier` | Compatibility paths with explicit nested submodules and symbol lists. | Same as `bucket`. |
|
||||
| `data_usage`, `capacity`, `notification`, `metrics`, `rebalance` | Domain and service facades consumed through owner-local boundaries. | Narrow one group at a time after explicit aliases or wrappers exist. |
|
||||
| `data_usage`, `capacity`, `notification`, `metrics`, `rebalance`, `integrity` | Domain and service facades consumed through owner-local boundaries. | Narrow one group at a time after explicit aliases or wrappers exist. |
|
||||
| `set_disk`, `object`, `object_api_utils`, `rio`, `bitrot`, `erasure`, `compression`, `cache`, `store_list` | Low-level object IO, reader, erasure, cache, and migration helper compatibility. | Keep stable while `SetDisks` remains the shared state carrier. |
|
||||
| `admin`, `event`, `global` | Admin, event hook, and bootstrap-global compatibility. | `global` is limited to bootstrap writes and lifecycle controls; read-only runtime access goes through `runtime`. |
|
||||
|
||||
@@ -28,7 +28,7 @@ External `rustfs_ecstore::api` imports stay in these local boundary files:
|
||||
|
||||
| Boundary file | Facade families consumed |
|
||||
|---|---|
|
||||
| `rustfs/src/storage/storage_api.rs` | Broad storage-owner bridge: admin, bucket submodules, capacity, compression, cluster, config, data usage, disk, error, event, global bootstrap controls, runtime getters, layout, metrics, notification, rebalance, rio, rpc, set disk, storage, tier. Replication pool/stat handles are projected into RustFS-local wrapper types here. |
|
||||
| `rustfs/src/storage/storage_api.rs` | Broad storage-owner bridge: admin, bucket submodules, capacity, compression, cluster, config, data usage, disk, error, event, global bootstrap controls, runtime getters, layout, metrics, notification, rebalance, integrity, rio, rpc, set disk, storage, tier. Replication pool/stat handles are projected into RustFS-local wrapper types here. |
|
||||
| `rustfs/src/storage_api.rs`, `rustfs/src/admin/storage_api.rs`, `rustfs/src/app/storage_api.rs` | Root, admin, and app owner boundaries: explicit aliases only, no `metadata`, `metadata_sys`, `quota`, `com`, or bare `init` module passthroughs; object and error aliases anchor on storage-api associated types and a local `StorageError`. |
|
||||
| `crates/scanner/src/storage_api.rs` | Bucket lifecycle, replication, metadata, capacity, config, data usage, disk, error, runtime, set disk, storage, tier. Replication queue config, admission, and heal object DTOs are projected into scanner-local types. |
|
||||
| `crates/obs/src/metrics/storage_api.rs` | Bucket bandwidth, lifecycle, replication, quota, capacity, data usage, error, runtime, storage; data usage is consumed as a local DTO projection. |
|
||||
|
||||
@@ -0,0 +1,72 @@
|
||||
# Object integrity inventory, audit, and migration
|
||||
|
||||
**Use this when:** assessing legacy object protection, checking a bounded manifest against an existing SHA256, or creating protected copies without overwriting the source.
|
||||
|
||||
## Protection and activation
|
||||
|
||||
Legacy shard-local checksums can detect accidental corruption, but do not independently bind a complete valid donor shard to its original object. A successful legacy read or Heal is not proof of original content. Inventory reports the stored protection declaration; it does not verify payload bytes. `independent_commitment` requires consistent object and part declarations. `unknown` means metadata could not be observed reliably.
|
||||
|
||||
Both `RUSTFS_SHARD_INTEGRITY_WRITE` and `RUSTFS_SHARD_INTEGRITY_FLEET_CONFIRMED` remain false by default. Enabling both permits protected new writes on that process. The second flag is an operator attestation, not a peer capability probe or a storage fencing mechanism. Before enabling, qualify all readers, writers, repair processes, rollback binaries, and replacement nodes against the protected format. Prevent old processes from accessing those drives. Measure full GET, Range GET, PUT, degraded read, and repair latency and throughput on the deployment's erasure layout. The readiness endpoint reports these requirements as unverified; it does not make the deployment safe automatically.
|
||||
|
||||
Disabling the write flags stops protection for new writes without an inherited mode; it does not remove existing commitments or make an old binary safe. Keep compatible readers for already migrated objects. Normal S3 requests and existing legacy reads do not depend on the integrity job worker being available.
|
||||
|
||||
## Administrative API
|
||||
|
||||
Use an authenticated, signed admin client. Paths use `/rustfs/admin/v3/integrity` and normal RustFS admin request authentication and body conventions.
|
||||
|
||||
| Method and suffix | Permission | Behavior |
|
||||
|---|---|---|
|
||||
| `GET /readiness` | `admin:ServerInfo` | Local support, local flags, operator attestation, and unverified fleet requirements |
|
||||
| `GET /{bucket}/inventory` | `admin:InspectData` | One page of version observations, no payload scan |
|
||||
| `POST /{bucket}/jobs` | `admin:StartBatchJob` | Persist a paused manifest and return its ID |
|
||||
| `GET /{bucket}/jobs/{job_id}` | `admin:DescribeBatchJob` | Durable status and per-item outcomes |
|
||||
| `POST /{bucket}/jobs/{job_id}/control` | `admin:StartBatchJob` | Explicit `resume`, `pause`, or `cancel` |
|
||||
|
||||
Inventory accepts `prefix`, `key-marker`, `version-marker`, and `limit` (1–100). Continue with both returned markers only after storing the entire page. Each inventory item returns a reusable string version selector, including `"null"` for the null version. Pages are observations, not a bucket snapshot: concurrent creates and deletes can affect enumeration. A recreated bucket has a different incarnation; do not combine observations across incarnations. Inventory and audit suppress their own read-repair submissions and bypass the body cache for content verification; they do not stop unrelated normal reads or background maintenance.
|
||||
|
||||
Example job body:
|
||||
|
||||
```json
|
||||
{
|
||||
"mode": "audit",
|
||||
"items": [
|
||||
{
|
||||
"key": "archive/object.bin",
|
||||
"version_id": null,
|
||||
"expected_sha256": "ungWv48Bz+pBQUDeXa4iI7ADYaOWF3qctBD/YfIAFa0="
|
||||
}
|
||||
],
|
||||
"bytes_per_second": 10485760,
|
||||
"max_object_bytes": 1073741824
|
||||
}
|
||||
```
|
||||
|
||||
The example digest is illustrative; supply the expected digest of your own object from an independent trusted record. `version_id` omitted or JSON null selects the current version; the string `"null"` selects the null version. A non-null version must be a UUID. The job records and rechecks the observed source identity. For `mode: "migrate"`, each item also requires a distinct `target_key` in the same bucket. No target may equal any source key in the manifest.
|
||||
|
||||
`expected_sha256` is canonical Base64 of 32 bytes, not hex or an ETag. When omitted, only a canonical stored single-part SHA256 is eligible. A stored checksum verifies consistency with stored metadata; it does not establish an external historical provenance. The service never derives a new expected digest from an unchecked read and calls that historical evidence.
|
||||
|
||||
Start the persisted job with:
|
||||
|
||||
```json
|
||||
{"operation":"resume"}
|
||||
```
|
||||
|
||||
Use the same control endpoint with `pause` or `cancel`. Poll status by ID and retain the ID externally; there is no job-list endpoint.
|
||||
|
||||
## Supported scope and resource limits
|
||||
|
||||
Audit accepts local, untransformed single-part objects with an expected SHA256. Multipart, encrypted, compressed, tiered objects, delete markers, and invalid protection declarations are reported as unsupported. Migration additionally requires a plain unversioned bucket without configured encryption, Object Lock, replication, notification, quota, lifecycle, table namespace, ACL, access logging, or on-demand migration. Disk commits continue to use the bucket’s existing durability policy. These checks hold the bucket configuration fence through publication so the direct storage writer cannot bypass those policies during a race.
|
||||
|
||||
Each manifest contains 1–64 items. `max_object_bytes` is 1 byte–5 GiB and limits the anonymous local temporary file used for one object. Ensure that amount of space is available in the process temporary directory. One cluster-wide worker runs at a time. `bytes_per_second` is 64 KiB/s–1 GiB/s and bounds each sequential source read, destination PUT stream, and verification read. It is not a physical disk or network aggregate ceiling: erasure redundancy and internal buffering still apply. Metadata requests are not byte throttled. Lower rates extend the time a source read or destination publication holds its normal object and bucket configuration locks; choose the rate and per-object budget with concurrent traffic in mind.
|
||||
|
||||
Migration uploads the same staged bytes that passed SHA256 verification, rechecks that checksum during PUT, and publishes with a create-only precondition. It preserves user metadata, content type/encoding/language, disposition, cache control, expiration, storage class, and tags. The new key has a new modification time and storage identity; consumers switch keys separately. The source is neither overwritten nor deleted. No original-object metadata is backfilled in place.
|
||||
|
||||
## Recovery and interpretation
|
||||
|
||||
`complete` means every item reached a final outcome, not that every item passed. Inspect `verified`, `migrated`, `mismatch`, `unsupported`, `stale`, and `conflict` individually. A source that is temporarily unavailable leaves the job `failed`; `resume` retries unavailable or prepared items while preserving completed results. A mismatch never publishes a target.
|
||||
|
||||
Pause and cancel first persist intent. An in-flight conditional target publication may finish; terminal acknowledgement follows after the worker stops. Cancellation does not delete a successfully published target. After a coordinator restart, `running` is the last durable state and does not prove a worker is alive. Explicitly resume to reacquire the distributed worker lock and recover. A persisted cancel intent is completed rather than restarted. A busy response means another worker still owns the lock; retry later.
|
||||
|
||||
A prepared checkpoint precedes destination publication. If the PUT acknowledgement or completion checkpoint is lost, resume first checks the destination's internal operation identity, protection declarations, size, and uncached full-content SHA256. Only a matching receipt is accepted as this job's completed copy. An unrelated destination is a conflict and is never overwritten. Source changes before verification finishes produce `stale`; create a new manifest after investigating the change. Bucket deletion/recreation invalidates the old job incarnation.
|
||||
|
||||
This workflow does not automatically repair damaged historical content, certify all historical shards, enable protected writes fleet-wide, clean MRF queues, or authorize rolling back to an incompatible binary.
|
||||
@@ -0,0 +1,207 @@
|
||||
// Copyright 2024 RustFS Team
|
||||
// Licensed under the Apache License, Version 2.0.
|
||||
|
||||
//! Administrator-directed integrity inspection and bounded, recoverable jobs.
|
||||
//! These routes do not change the MinIO batch-job compatibility endpoints.
|
||||
|
||||
use crate::admin::auth::authorize_admin_request;
|
||||
use crate::admin::router::{AdminOperation, Operation, S3Router};
|
||||
use crate::admin::runtime_sources::object_store_from_req;
|
||||
use crate::admin::storage_api::integrity as service;
|
||||
use crate::admin::storage_api::s3::{Body, S3Error, S3ErrorCode, S3Request, S3Response, S3Result, error as admin_s3_error};
|
||||
use crate::admin::utils::{extract_query_params, json_response, read_compatible_admin_body};
|
||||
use crate::server::ADMIN_PREFIX;
|
||||
use hyper::{Method, StatusCode};
|
||||
use matchit::Params;
|
||||
use rustfs_policy::policy::action::{Action, AdminAction};
|
||||
use serde::Deserialize;
|
||||
use uuid::Uuid;
|
||||
|
||||
#[derive(Clone, Copy)]
|
||||
enum Route {
|
||||
Readiness,
|
||||
Inventory,
|
||||
Create,
|
||||
Status,
|
||||
Control,
|
||||
}
|
||||
|
||||
struct Handler(Route);
|
||||
|
||||
fn permission(route: Route) -> AdminAction {
|
||||
match route {
|
||||
Route::Readiness => AdminAction::ServerInfoAdminAction,
|
||||
Route::Inventory => AdminAction::InspectDataAction,
|
||||
Route::Status => AdminAction::DescribeBatchJobAction,
|
||||
Route::Create | Route::Control => AdminAction::StartBatchJobAction,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn register_integrity_routes(router: &mut S3Router<AdminOperation>) -> std::io::Result<()> {
|
||||
for (method, path, handler) in [
|
||||
(Method::GET, "/v3/integrity/readiness", &Handler(Route::Readiness)),
|
||||
(Method::GET, "/v3/integrity/{bucket}/inventory", &Handler(Route::Inventory)),
|
||||
(Method::POST, "/v3/integrity/{bucket}/jobs", &Handler(Route::Create)),
|
||||
(Method::GET, "/v3/integrity/{bucket}/jobs/{job_id}", &Handler(Route::Status)),
|
||||
(Method::POST, "/v3/integrity/{bucket}/jobs/{job_id}/control", &Handler(Route::Control)),
|
||||
] {
|
||||
router.insert(method, &format!("{ADMIN_PREFIX}{path}"), AdminOperation(handler))?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
struct ControlRequest {
|
||||
operation: Control,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
enum Control {
|
||||
Pause,
|
||||
Resume,
|
||||
Cancel,
|
||||
}
|
||||
|
||||
fn map_error(error: service::IntegrityError) -> S3Error {
|
||||
match error {
|
||||
service::IntegrityError::Invalid(message) => admin_s3_error(S3ErrorCode::InvalidArgument, message),
|
||||
service::IntegrityError::Conflict => admin_s3_error(S3ErrorCode::PreconditionFailed, "integrity job or target changed"),
|
||||
service::IntegrityError::Busy => {
|
||||
admin_s3_error(S3ErrorCode::SlowDown, "an integrity worker is already active; retry resume later")
|
||||
}
|
||||
service::IntegrityError::NotFound => admin_s3_error(S3ErrorCode::NoSuchKey, "integrity job not found"),
|
||||
service::IntegrityError::NotActivated => admin_s3_error(
|
||||
S3ErrorCode::InvalidRequest,
|
||||
"protected writes must be explicitly activated before migration",
|
||||
),
|
||||
service::IntegrityError::UnsupportedBucket => admin_s3_error(
|
||||
S3ErrorCode::InvalidRequest,
|
||||
"migration requires a plain unversioned bucket without configured automation or retention",
|
||||
),
|
||||
_ => admin_s3_error(
|
||||
S3ErrorCode::InternalError,
|
||||
"integrity operation failed; inspect the durable job before retrying",
|
||||
),
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl Operation for Handler {
|
||||
async fn call(&self, req: S3Request<Body>, params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
|
||||
let credentials = authorize_admin_request(&req, vec![Action::AdminAction(permission(self.0))]).await?;
|
||||
if matches!(self.0, Route::Readiness) {
|
||||
return json_response(StatusCode::OK, &service::readiness());
|
||||
}
|
||||
let store = object_store_from_req(&req)
|
||||
.ok_or_else(|| admin_s3_error(S3ErrorCode::InternalError, "object store is not initialized"))?;
|
||||
let bucket = params
|
||||
.get("bucket")
|
||||
.ok_or_else(|| admin_s3_error(S3ErrorCode::InvalidArgument, "bucket is required"))?;
|
||||
match self.0 {
|
||||
Route::Inventory => {
|
||||
let query = extract_query_params(&req.uri);
|
||||
if query
|
||||
.keys()
|
||||
.any(|key| !["prefix", "key-marker", "version-marker", "limit"].contains(&key.as_str()))
|
||||
{
|
||||
return Err(admin_s3_error(S3ErrorCode::InvalidArgument, "unknown inventory parameter"));
|
||||
}
|
||||
let limit = query
|
||||
.get("limit")
|
||||
.map_or(Ok(100), |v| v.parse::<i32>())
|
||||
.map_err(|_| admin_s3_error(S3ErrorCode::InvalidArgument, "invalid limit"))?;
|
||||
let page = service::inventory(
|
||||
store,
|
||||
bucket,
|
||||
query.get("prefix").map_or("", String::as_str),
|
||||
query.get("key-marker").cloned(),
|
||||
query.get("version-marker").cloned(),
|
||||
limit,
|
||||
)
|
||||
.await
|
||||
.map_err(map_error)?;
|
||||
json_response(StatusCode::OK, &page)
|
||||
}
|
||||
Route::Create => {
|
||||
let body = read_compatible_admin_body(req.input, 128 * 1024, req.uri.path(), &credentials.secret_key).await?;
|
||||
let request: service::JobRequest = serde_json::from_slice(&body)
|
||||
.map_err(|_| admin_s3_error(S3ErrorCode::InvalidArgument, "invalid integrity job request"))?;
|
||||
let job = service::create_job(store, bucket, request).await.map_err(map_error)?;
|
||||
json_response(StatusCode::CREATED, &job)
|
||||
}
|
||||
Route::Status | Route::Control => {
|
||||
let id = params
|
||||
.get("job_id")
|
||||
.and_then(|s| Uuid::parse_str(s).ok())
|
||||
.filter(|id| !id.is_nil())
|
||||
.ok_or_else(|| admin_s3_error(S3ErrorCode::InvalidArgument, "invalid job id"))?;
|
||||
let job = if matches!(self.0, Route::Status) {
|
||||
service::get_job(store, bucket, id).await.map_err(map_error)?
|
||||
} else {
|
||||
let body = read_compatible_admin_body(req.input, 1024, req.uri.path(), &credentials.secret_key).await?;
|
||||
let request: ControlRequest = serde_json::from_slice(&body)
|
||||
.map_err(|_| admin_s3_error(S3ErrorCode::InvalidArgument, "expected pause, resume or cancel"))?;
|
||||
match request.operation {
|
||||
Control::Resume => service::resume_job(store, bucket, id).await,
|
||||
Control::Pause => service::control_job(store, bucket, id, false).await,
|
||||
Control::Cancel => service::control_job(store, bucket, id, true).await,
|
||||
}
|
||||
.map_err(map_error)?
|
||||
};
|
||||
json_response(StatusCode::OK, &job)
|
||||
}
|
||||
Route::Readiness => json_response(StatusCode::OK, &service::readiness()),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
#[test]
|
||||
fn integrity_read_permissions_cannot_start_or_resume_payload_writes() {
|
||||
assert_eq!(permission(Route::Create), AdminAction::StartBatchJobAction);
|
||||
assert_eq!(permission(Route::Control), AdminAction::StartBatchJobAction);
|
||||
assert_eq!(permission(Route::Inventory), AdminAction::InspectDataAction);
|
||||
assert_eq!(permission(Route::Status), AdminAction::DescribeBatchJobAction);
|
||||
}
|
||||
#[test]
|
||||
fn control_rejects_unknown_fields_and_operations() {
|
||||
assert!(serde_json::from_str::<ControlRequest>(r#"{"operation":"resume","force":true}"#).is_err());
|
||||
assert!(serde_json::from_str::<ControlRequest>(r#"{"operation":"overwrite"}"#).is_err());
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[tokio::test]
|
||||
async fn integrity_routes_reject_unauthenticated_requests_before_storage_access() {
|
||||
let mut router = matchit::Router::new();
|
||||
router.insert("/{bucket}", ()).expect("test route");
|
||||
for route in [
|
||||
Route::Readiness,
|
||||
Route::Inventory,
|
||||
Route::Create,
|
||||
Route::Status,
|
||||
Route::Control,
|
||||
] {
|
||||
let req = S3Request {
|
||||
input: Body::from(String::new()),
|
||||
method: Method::POST,
|
||||
uri: http::Uri::from_static("/rustfs/admin/v3/integrity/example/jobs"),
|
||||
headers: http::HeaderMap::new(),
|
||||
extensions: http::Extensions::new(),
|
||||
credentials: None,
|
||||
region: None,
|
||||
service: None,
|
||||
trailing_headers: None,
|
||||
};
|
||||
let error = Handler(route)
|
||||
.call(req, router.at("/example").expect("params").params)
|
||||
.await
|
||||
.expect_err("credentials required");
|
||||
assert_eq!(error.code(), &S3ErrorCode::InvalidRequest);
|
||||
assert_eq!(error.message(), Some("get cred failed"));
|
||||
}
|
||||
}
|
||||
@@ -33,6 +33,7 @@ pub(crate) mod iam_error;
|
||||
pub mod idp_compat;
|
||||
pub mod ilm_transition;
|
||||
pub mod inspect_archive;
|
||||
pub mod integrity;
|
||||
pub mod is_admin;
|
||||
pub mod kms;
|
||||
pub mod kms_audit;
|
||||
|
||||
@@ -35,10 +35,10 @@ mod route_registration_test;
|
||||
|
||||
use handlers::{
|
||||
account, audit, batch_job, bucket_meta, cluster_snapshot, config_admin, diagnostics, durability as durability_handler,
|
||||
extensions, gateway_key_inventory, heal, health, idp_compat, ilm_transition, inspect_archive, kms, mfa, module_switch,
|
||||
object_data_cache, object_zip_download, oidc, on_demand_migration, plugins_catalog, plugins_instances, pools, profile_admin,
|
||||
quota as quota_handler, rebalance, replication as replication_handler, scanner, site_replication, sts, system, table_catalog,
|
||||
tier, tls_debug, usage_prefix, user,
|
||||
extensions, gateway_key_inventory, heal, health, idp_compat, ilm_transition, inspect_archive, integrity, kms, mfa,
|
||||
module_switch, object_data_cache, object_zip_download, oidc, on_demand_migration, plugins_catalog, plugins_instances, pools,
|
||||
profile_admin, quota as quota_handler, rebalance, replication as replication_handler, scanner, site_replication, sts, system,
|
||||
table_catalog, tier, tls_debug, usage_prefix, user,
|
||||
};
|
||||
use router::{AdminOperation, S3Router};
|
||||
use s3s::route::S3Route;
|
||||
@@ -94,6 +94,7 @@ fn register_admin_routes(r: &mut S3Router<AdminOperation>) -> std::io::Result<()
|
||||
|
||||
replication_handler::register_replication_route(r)?;
|
||||
batch_job::register_batch_job_route(r)?;
|
||||
integrity::register_integrity_routes(r)?;
|
||||
site_replication::register_site_replication_route(r)?;
|
||||
profile_admin::register_profiling_route(r)?;
|
||||
diagnostics::register_diagnostics_route(r)?;
|
||||
|
||||
@@ -1585,6 +1585,36 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[
|
||||
COMMIT_TABLE,
|
||||
RouteRiskLevel::High,
|
||||
),
|
||||
admin(
|
||||
HttpMethod::Get,
|
||||
"/rustfs/admin/v3/integrity/readiness",
|
||||
SERVER_INFO,
|
||||
RouteRiskLevel::Sensitive,
|
||||
),
|
||||
admin(
|
||||
HttpMethod::Get,
|
||||
"/rustfs/admin/v3/integrity/{bucket}/inventory",
|
||||
INSPECT_DATA,
|
||||
RouteRiskLevel::Sensitive,
|
||||
),
|
||||
admin(
|
||||
HttpMethod::Post,
|
||||
"/rustfs/admin/v3/integrity/{bucket}/jobs",
|
||||
START_BATCH_JOB,
|
||||
RouteRiskLevel::High,
|
||||
),
|
||||
admin(
|
||||
HttpMethod::Get,
|
||||
"/rustfs/admin/v3/integrity/{bucket}/jobs/{job_id}",
|
||||
DESCRIBE_BATCH_JOB,
|
||||
RouteRiskLevel::Sensitive,
|
||||
),
|
||||
admin(
|
||||
HttpMethod::Post,
|
||||
"/rustfs/admin/v3/integrity/{bucket}/jobs/{job_id}/control",
|
||||
START_BATCH_JOB,
|
||||
RouteRiskLevel::High,
|
||||
),
|
||||
// MinIO admin compat: batch job lifecycle (backlog#613).
|
||||
admin(HttpMethod::Post, "/rustfs/admin/v3/start-job", START_BATCH_JOB, RouteRiskLevel::High),
|
||||
admin(HttpMethod::Get, "/rustfs/admin/v3/list-jobs", LIST_BATCH_JOBS, RouteRiskLevel::Sensitive),
|
||||
|
||||
@@ -337,6 +337,19 @@ fn expected_admin_route_matrix() -> Vec<RouteMatrixEntry> {
|
||||
admin_route(Method::DELETE, "/v3/remove-remote-target"),
|
||||
admin_route(Method::POST, "/v3/replication/diff"),
|
||||
admin_route(Method::GET, "/v3/replication/mrf"),
|
||||
admin_route(Method::GET, "/v3/integrity/readiness"),
|
||||
admin_route_sample(Method::GET, "/v3/integrity/{bucket}/inventory", "/v3/integrity/example/inventory"),
|
||||
admin_route_sample(Method::POST, "/v3/integrity/{bucket}/jobs", "/v3/integrity/example/jobs"),
|
||||
admin_route_sample(
|
||||
Method::GET,
|
||||
"/v3/integrity/{bucket}/jobs/{job_id}",
|
||||
"/v3/integrity/example/jobs/11111111-1111-4111-8111-111111111111",
|
||||
),
|
||||
admin_route_sample(
|
||||
Method::POST,
|
||||
"/v3/integrity/{bucket}/jobs/{job_id}/control",
|
||||
"/v3/integrity/example/jobs/11111111-1111-4111-8111-111111111111/control",
|
||||
),
|
||||
admin_route(Method::POST, "/v3/start-job"),
|
||||
admin_route(Method::GET, "/v3/list-jobs"),
|
||||
admin_route(Method::GET, "/v3/status-job"),
|
||||
|
||||
@@ -16,6 +16,12 @@ use std::ops::Deref;
|
||||
use std::sync::Arc;
|
||||
|
||||
use rustfs_storage_api as storage_contracts;
|
||||
|
||||
pub(crate) mod integrity {
|
||||
pub(crate) use crate::storage::storage_api::ecstore_integrity::{
|
||||
IntegrityError, JobRequest, control_job, create_job, get_job, inventory, readiness, resume_job,
|
||||
};
|
||||
}
|
||||
use time::OffsetDateTime;
|
||||
|
||||
mod ecstore_bucket {
|
||||
|
||||
@@ -21,6 +21,12 @@ use std::time::Duration;
|
||||
|
||||
use rand::RngExt as _;
|
||||
use rustfs_storage_api as storage_contracts;
|
||||
|
||||
pub(crate) mod ecstore_integrity {
|
||||
pub(crate) use rustfs_ecstore::api::integrity::{
|
||||
IntegrityError, JobRequest, control_job, create_job, get_job, inventory, readiness, resume_job,
|
||||
};
|
||||
}
|
||||
use tokio::sync::{Mutex, OwnedMutexGuard};
|
||||
use tokio_util::sync::CancellationToken;
|
||||
|
||||
|
||||
Reference in New Issue
Block a user