mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-23 04:39:04 +00:00
Compare commits
15 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 5ff5a7a06e | |||
| b6ba89d9e4 | |||
| 648d5166e2 | |||
| 84eb5aebef | |||
| c58accabe7 | |||
| 4d92835d7b | |||
| 71db3327ad | |||
| 68e95d497e | |||
| 8350d73be9 | |||
| 9a2a33a7eb | |||
| 5df26a13ba | |||
| e3a362989f | |||
| 0cd2ae20e2 | |||
| 55a7fa9f03 | |||
| 9f51f37a0d |
@@ -30,7 +30,8 @@ make build-docker BUILD_OS=ubuntu22.04
|
||||
- Crate membership: `Cargo.toml` `[workspace].members`
|
||||
- Architecture, layering, crate map: [ARCHITECTURE.md](ARCHITECTURE.md)
|
||||
- Migration guardrails & readiness contracts: [docs/architecture/](docs/architecture/README.md)
|
||||
- CI gates: `.github/workflows/ci.yml` (source of truth; never copy its steps into docs)
|
||||
- CI workflow steps: `.github/workflows/`; event, timeout, and required-status
|
||||
matrix: [docs/testing/ci-gates.md](docs/testing/ci-gates.md)
|
||||
- Test-layer taxonomy, per-layer entry commands, serial/nextest rules, flake
|
||||
policy: [docs/testing/README.md](docs/testing/README.md)
|
||||
- Tier/ILM transition debugging (xl.meta inspection, versionId tracing):
|
||||
|
||||
@@ -70,6 +70,8 @@ make pre-pr
|
||||
|
||||
> For the full test-layer taxonomy (unit / ecstore black-box / e2e / s3s-e2e / S3 compatibility / chaos / fuzz / bench), each layer's entry command, the naming conventions the migration gate depends on, and the serial/nextest rules, see [docs/testing/README.md](docs/testing/README.md).
|
||||
|
||||
> For the event, timeout, required-status, and local reproduction matrix, see [docs/testing/ci-gates.md](docs/testing/ci-gates.md).
|
||||
|
||||
### 🔒 Automated Pre-commit Hooks
|
||||
#### What `make pre-commit` and `make pre-pr` actually run
|
||||
|
||||
|
||||
Generated
+23
-27
@@ -1858,9 +1858,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "cc"
|
||||
version = "1.4.3"
|
||||
version = "1.4.4"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "509591b7bcd67f4ef775afad7662703b4935daaa6ec0e5605cfb1090b32a2b6d"
|
||||
checksum = "0ad534f4357a5264cce5019c989cf66a4f0dc4e0d1b1d15f8aacec0ff7360273"
|
||||
dependencies = [
|
||||
"find-msvc-tools",
|
||||
"jobserver",
|
||||
@@ -2522,12 +2522,6 @@ dependencies = [
|
||||
"subtle",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "cty"
|
||||
version = "0.2.2"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "b365fabc795046672053e29c954733ec3b05e4be654ab130fe8f1f94d7051f35"
|
||||
|
||||
[[package]]
|
||||
name = "curve25519-dalek"
|
||||
version = "4.1.3"
|
||||
@@ -5988,15 +5982,6 @@ version = "0.2.16"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "b6d2cec3eae94f9f509c767b45932f1ada8350c4bdb85af2fcab4a3c14807981"
|
||||
|
||||
[[package]]
|
||||
name = "libmimalloc-sys"
|
||||
version = "0.1.49"
|
||||
source = "git+https://github.com/xonatius/mimalloc_rust.git?rev=6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11#6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11"
|
||||
dependencies = [
|
||||
"cc",
|
||||
"cty",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "libredox"
|
||||
version = "0.1.20"
|
||||
@@ -6397,14 +6382,6 @@ dependencies = [
|
||||
"synstructure 0.13.2",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "mimalloc"
|
||||
version = "0.1.52"
|
||||
source = "git+https://github.com/xonatius/mimalloc_rust.git?rev=6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11#6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11"
|
||||
dependencies = [
|
||||
"libmimalloc-sys",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "mime"
|
||||
version = "0.3.17"
|
||||
@@ -9162,13 +9139,11 @@ dependencies = [
|
||||
"insta",
|
||||
"jiff",
|
||||
"libc",
|
||||
"libmimalloc-sys",
|
||||
"libsystemd",
|
||||
"matchit 0.9.2",
|
||||
"md-5 0.11.0",
|
||||
"metrics",
|
||||
"metrics-util",
|
||||
"mimalloc",
|
||||
"mime_guess",
|
||||
"opentelemetry",
|
||||
"opentelemetry_sdk",
|
||||
@@ -9204,6 +9179,8 @@ dependencies = [
|
||||
"rustfs-lock",
|
||||
"rustfs-log-analyzer",
|
||||
"rustfs-madmin",
|
||||
"rustfs-mimalloc",
|
||||
"rustfs-mimalloc-sys",
|
||||
"rustfs-notify",
|
||||
"rustfs-object-capacity",
|
||||
"rustfs-object-data-cache",
|
||||
@@ -9587,6 +9564,7 @@ dependencies = [
|
||||
"serde",
|
||||
"serde_json",
|
||||
"serial_test",
|
||||
"sha2 0.11.0",
|
||||
"temp-env",
|
||||
"tempfile",
|
||||
"thiserror 2.0.20",
|
||||
@@ -9875,6 +9853,24 @@ dependencies = [
|
||||
"tokio",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-mimalloc"
|
||||
version = "0.5.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "a406f4aa07084301d485beec873af6dccc8e3f8762da244743df92038b1db1a6"
|
||||
dependencies = [
|
||||
"rustfs-mimalloc-sys",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-mimalloc-sys"
|
||||
version = "0.5.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "c3051b819175f58445d4c369a72f0ab88149f3885ba8bea2aff3be01f53fe7cd"
|
||||
dependencies = [
|
||||
"cc",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-notify"
|
||||
version = "1.0.0-rc.3"
|
||||
|
||||
+2
-2
@@ -350,8 +350,8 @@ russh-sftp = "2.4.0"
|
||||
dav-server = "0.11.0"
|
||||
|
||||
# Performance Analysis and Memory Profiling
|
||||
mimalloc = { version = "0.1.52", git = "https://github.com/xonatius/mimalloc_rust.git", rev = "6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11" }
|
||||
libmimalloc-sys = { version = "0.1.49", git = "https://github.com/xonatius/mimalloc_rust.git", rev = "6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11", features = ["extended"] }
|
||||
rustfs-mimalloc = { version = "0.5.0" }
|
||||
rustfs-mimalloc-sys = { version = "0.5.0" }
|
||||
hotpath = { version = "0.23.3", default-features = false }
|
||||
# Snapshot testing for output format regression detection
|
||||
insta = { version = "1.48" }
|
||||
|
||||
@@ -8000,10 +8000,15 @@ impl DiskAPI for LocalDisk {
|
||||
use std::io::Write as _;
|
||||
|
||||
let file_path = self.io_get_object_path(volume, path)?;
|
||||
let lock_path = file_path.with_extension("rustfs-cas.lock");
|
||||
let path = path.to_string();
|
||||
let sync_metadata = effective_durability(volume).syncs_commit_metadata();
|
||||
return Ok(tokio::task::spawn_blocking(move || {
|
||||
// A persistent directory lock bounds metadata growth. Removing
|
||||
// per-target lock files can split flock ownership across inodes.
|
||||
let lock_path = file_path
|
||||
.parent()
|
||||
.ok_or_else(|| std::io::Error::new(ErrorKind::InvalidInput, "conditional file has no parent"))?
|
||||
.join(".rustfs-cas.lock");
|
||||
let lock = std::fs::OpenOptions::new()
|
||||
.create(true)
|
||||
.truncate(false)
|
||||
@@ -8068,7 +8073,25 @@ impl DiskAPI for LocalDisk {
|
||||
.map_err(DiskError::from)??);
|
||||
}
|
||||
|
||||
#[cfg(not(unix))]
|
||||
#[cfg(windows)]
|
||||
{
|
||||
let file_path = self.io_get_object_path(volume, path)?;
|
||||
let sync_metadata = effective_durability(volume).syncs_commit_metadata();
|
||||
let publication_root = self.publication_root.clone();
|
||||
return Ok(tokio::task::spawn_blocking(move || {
|
||||
os::compare_and_update_control_file(
|
||||
&file_path,
|
||||
expected.as_deref(),
|
||||
replacement.as_deref(),
|
||||
sync_metadata,
|
||||
&publication_root,
|
||||
)
|
||||
})
|
||||
.await
|
||||
.map_err(DiskError::from)??);
|
||||
}
|
||||
|
||||
#[cfg(not(any(unix, windows)))]
|
||||
{
|
||||
let _ = (volume, path, expected, replacement);
|
||||
Err(DiskError::MethodNotAllowed)
|
||||
@@ -21823,9 +21846,9 @@ mod test {
|
||||
assert!(matches!(results[1].as_ref().unwrap_err(), DiskError::Io(_)));
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[cfg(any(unix, windows))]
|
||||
#[tokio::test]
|
||||
async fn conditional_file_update_never_deletes_a_new_owner() {
|
||||
async fn windows_and_unix_conditional_file_update_never_deletes_a_new_owner() {
|
||||
use tempfile::tempdir;
|
||||
|
||||
let dir = tempdir().expect("temp dir should be created");
|
||||
@@ -21856,8 +21879,18 @@ mod test {
|
||||
disk.read_all(RUSTFS_META_BUCKET, HEALING_MARKER_PATH)
|
||||
.await
|
||||
.expect("new owner marker should remain"),
|
||||
owner_b
|
||||
owner_b.clone()
|
||||
);
|
||||
assert_eq!(
|
||||
disk.compare_and_update_file(RUSTFS_META_BUCKET, HEALING_MARKER_PATH, Some(owner_b), None)
|
||||
.await
|
||||
.expect("current owner should remove marker"),
|
||||
ConditionalFileUpdate::Updated
|
||||
);
|
||||
assert!(matches!(
|
||||
disk.read_all(RUSTFS_META_BUCKET, HEALING_MARKER_PATH).await,
|
||||
Err(DiskError::FileNotFound)
|
||||
));
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
@@ -21872,7 +21905,10 @@ mod test {
|
||||
let marker_path = disk
|
||||
.get_object_path(RUSTFS_META_BUCKET, HEALING_MARKER_PATH)
|
||||
.expect("marker path should resolve");
|
||||
let lock_path = marker_path.with_extension("rustfs-cas.lock");
|
||||
let lock_path = marker_path
|
||||
.parent()
|
||||
.expect("marker path should have a parent")
|
||||
.join(".rustfs-cas.lock");
|
||||
let lock = std::fs::OpenOptions::new()
|
||||
.create(true)
|
||||
.truncate(false)
|
||||
@@ -21893,6 +21929,40 @@ mod test {
|
||||
assert!(matches!(err, DiskError::Io(ref err) if err.kind() == ErrorKind::WouldBlock));
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
#[tokio::test]
|
||||
async fn windows_conditional_file_update_returns_would_block_when_marker_lock_is_contended() {
|
||||
let dir = tempfile::tempdir().expect("temp dir should be created");
|
||||
let endpoint = Endpoint::try_from(dir.path().to_str().expect("temp dir should be utf8")).expect("endpoint should parse");
|
||||
let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created");
|
||||
ensure_test_volume(&disk, RUSTFS_META_BUCKET).await;
|
||||
let marker_path = disk
|
||||
.get_object_path(RUSTFS_META_BUCKET, HEALING_MARKER_PATH)
|
||||
.expect("marker path should resolve");
|
||||
let lock_path = marker_path
|
||||
.parent()
|
||||
.expect("marker path should have a parent")
|
||||
.join(".rustfs-cas.lock");
|
||||
let lock = std::fs::OpenOptions::new()
|
||||
.create(true)
|
||||
.truncate(false)
|
||||
.read(true)
|
||||
.write(true)
|
||||
.open(lock_path)
|
||||
.expect("marker lock should open");
|
||||
lock.try_lock().expect("marker lock should be held");
|
||||
|
||||
let err = tokio::time::timeout(
|
||||
Duration::from_secs(1),
|
||||
disk.compare_and_update_file(RUSTFS_META_BUCKET, HEALING_MARKER_PATH, None, Some(Bytes::from_static(b"owner"))),
|
||||
)
|
||||
.await
|
||||
.expect("contended conditional update must not block")
|
||||
.expect_err("contended conditional update must retry");
|
||||
|
||||
assert!(matches!(err, DiskError::Io(ref err) if err.kind() == ErrorKind::WouldBlock));
|
||||
}
|
||||
|
||||
#[cfg(target_os = "linux")]
|
||||
#[tokio::test]
|
||||
async fn replacement_io_paths_stay_under_the_mount_lease() {
|
||||
|
||||
@@ -12,6 +12,8 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
#[cfg(windows)]
|
||||
use crate::disk::ConditionalFileUpdate;
|
||||
use crate::disk::error::DiskError;
|
||||
use crate::disk::error::Result;
|
||||
use crate::disk::error_conv::to_file_error;
|
||||
@@ -3458,6 +3460,89 @@ fn read_windows_relative_file(file_path: &Path, parent_guard: &ExistingBaseDirec
|
||||
Ok(Some(data))
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
pub(crate) fn compare_and_update_control_file(
|
||||
file_path: &Path,
|
||||
expected: Option<&[u8]>,
|
||||
replacement: Option<&[u8]>,
|
||||
sync_metadata: bool,
|
||||
publication_root: &PublicationRoot,
|
||||
) -> io::Result<ConditionalFileUpdate> {
|
||||
use windows_sys::{
|
||||
Wdk::Storage::FileSystem::{
|
||||
FILE_NON_DIRECTORY_FILE, FILE_OPEN, FILE_OPEN_IF, FILE_OPEN_REPARSE_POINT, FILE_SYNCHRONOUS_IO_NONALERT,
|
||||
},
|
||||
Win32::Storage::FileSystem::{
|
||||
DELETE, FILE_ATTRIBUTE_NORMAL, FILE_READ_ATTRIBUTES, FILE_SHARE_READ, FILE_SHARE_WRITE, FILE_WRITE_DATA, SYNCHRONIZE,
|
||||
},
|
||||
};
|
||||
|
||||
let parent = file_path
|
||||
.parent()
|
||||
.ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "conditional file has no parent"))?;
|
||||
let parent_guard = lock_windows_directory_tree(parent, Some(parent), publication_root)?;
|
||||
let lock = open_windows_relative(
|
||||
parent_guard.last_handle()?,
|
||||
std::ffi::OsStr::new(".rustfs-cas.lock"),
|
||||
SYNCHRONIZE | FILE_READ_ATTRIBUTES | FILE_WRITE_DATA,
|
||||
FILE_SHARE_READ | FILE_SHARE_WRITE,
|
||||
FILE_OPEN_IF,
|
||||
FILE_NON_DIRECTORY_FILE | FILE_OPEN_REPARSE_POINT | FILE_SYNCHRONOUS_IO_NONALERT,
|
||||
FILE_ATTRIBUTE_NORMAL,
|
||||
true,
|
||||
)?;
|
||||
validate_windows_owned_file(&lock)?;
|
||||
match lock.as_file().try_lock() {
|
||||
Ok(()) => {}
|
||||
Err(std::fs::TryLockError::WouldBlock) => return Err(io::Error::from(io::ErrorKind::WouldBlock)),
|
||||
Err(std::fs::TryLockError::Error(err)) => return Err(err),
|
||||
}
|
||||
|
||||
let current = read_windows_relative_file(file_path, &parent_guard)?;
|
||||
let matches = match (¤t, expected) {
|
||||
(None, None) => true,
|
||||
(Some(current), Some(expected)) => current.as_slice() == expected,
|
||||
_ => false,
|
||||
};
|
||||
if !matches {
|
||||
return Ok(match current {
|
||||
None => ConditionalFileUpdate::Missing,
|
||||
Some(_) => ConditionalFileUpdate::Mismatch,
|
||||
});
|
||||
}
|
||||
|
||||
match replacement {
|
||||
Some(replacement) => RenameDestinationPathGuard {
|
||||
directory: parent.to_path_buf(),
|
||||
_directory_guard: parent_guard,
|
||||
}
|
||||
.write_file_for_path_access(file_path, replacement, sync_metadata, sync_metadata)?,
|
||||
None => {
|
||||
let file_name = file_path
|
||||
.file_name()
|
||||
.ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "conditional file must have a name"))?;
|
||||
let file = open_windows_relative(
|
||||
parent_guard.last_handle()?,
|
||||
file_name,
|
||||
DELETE | SYNCHRONIZE | FILE_READ_ATTRIBUTES,
|
||||
FILE_SHARE_READ,
|
||||
FILE_OPEN,
|
||||
FILE_NON_DIRECTORY_FILE | FILE_OPEN_REPARSE_POINT | FILE_SYNCHRONOUS_IO_NONALERT,
|
||||
0,
|
||||
true,
|
||||
)?;
|
||||
validate_windows_owned_file(&file)?;
|
||||
set_windows_file_delete_on_close(&file, true)?;
|
||||
drop(file);
|
||||
if sync_metadata {
|
||||
fsync_dir_std(parent)?;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Ok(ConditionalFileUpdate::Updated)
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
fn open_windows_directory_component(
|
||||
parent: &WindowsDirectoryHandle,
|
||||
|
||||
@@ -2124,26 +2124,13 @@ impl SetDisks {
|
||||
|
||||
let put_object_size = known_put_object_storage_size(data.size());
|
||||
let shard_file_size_raw = erasure.shard_file_size(put_object_size);
|
||||
let is_inline_buffer =
|
||||
storage_class_config.should_inline(shard_file_size_raw, erasure.data_shards, opts.versioned);
|
||||
let is_inline_buffer = storage_class_config.should_inline(shard_file_size_raw, erasure.data_shards, opts.versioned);
|
||||
|
||||
let collect_stage_timing = rustfs_io_metrics::put_stage_metrics_enabled() || issue3031_diag_enabled();
|
||||
let shard_file_size = shard_file_size_raw;
|
||||
let shard_size = erasure.shard_size();
|
||||
let write_path = classify_put_write_path(is_inline_buffer, put_object_size, fi.erasure.block_size);
|
||||
let direct_inline_commit = matches!(write_path, SmallWritePath::Inline);
|
||||
{
|
||||
use std::io::Write;
|
||||
let msg = format!(
|
||||
"INLINE_DEBUG: bucket={} obj={} size={} shard_fs={} ds={} bs={} inline={} direct={} path={} iblock={} ver={}\n",
|
||||
bucket, object, put_object_size, shard_file_size_raw, erasure.data_shards, fi.erasure.block_size,
|
||||
is_inline_buffer, direct_inline_commit, write_path.metric_label(), storage_class_config.inline_block(), opts.versioned
|
||||
);
|
||||
if let Ok(mut f) = std::fs::OpenOptions::new().create(true).append(true).open("/tmp/rustfs_inline_debug.log") {
|
||||
let _ = f.write_all(msg.as_bytes());
|
||||
}
|
||||
let _ = std::io::stderr().write_all(msg.as_bytes());
|
||||
}
|
||||
rustfs_io_metrics::record_put_object_path(write_path.metric_label());
|
||||
let writer_setup_stage_start = collect_stage_timing.then(Instant::now);
|
||||
let (mut writers, errors) = if direct_inline_commit {
|
||||
|
||||
@@ -91,6 +91,7 @@ metrics = { workspace = true }
|
||||
base64 = { workspace = true }
|
||||
bytes = { workspace = true }
|
||||
crc-fast = { workspace = true }
|
||||
sha2 = { workspace = true }
|
||||
|
||||
[dev-dependencies]
|
||||
serde_json = { workspace = true, features = ["raw_value"] }
|
||||
|
||||
@@ -373,6 +373,11 @@ impl ErasureSetHealer {
|
||||
set_disk_id: &str,
|
||||
buckets: &[String],
|
||||
) -> Result<(ResumeManager, CheckpointManager)> {
|
||||
if self.replacement_task_id.is_none() && CheckpointManager::is_blocked(&self.disk, task_id).await {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Resume task {task_id} has a blocked checkpoint"),
|
||||
});
|
||||
}
|
||||
// check if resume state exists
|
||||
let has_resume_state = if self.replacement_task_id.is_some() {
|
||||
ResumeManager::has_replacement_intent(&self.disk, task_id).await
|
||||
|
||||
@@ -51,6 +51,7 @@ const RESUME_STATE_FILE: &str = "ahm_resume_state.json";
|
||||
const REPLACEMENT_INTENT_FILE: &str = "ahm_replacement_intent.json";
|
||||
const RESUME_PROGRESS_FILE: &str = "ahm_progress.json";
|
||||
pub(super) const RESUME_CHECKPOINT_FILE: &str = "ahm_checkpoint.json";
|
||||
pub(super) const RESUME_CHECKPOINT_BLOCKED_FILE: &str = "ahm_checkpoint.blocked";
|
||||
const REPLACEMENT_COMPLETION_PROOF_FILE: &str = "ahm_replacement_completion_proof.json";
|
||||
const REPLACEMENT_RECOVERY_DIR: &str = "ahm-replacement";
|
||||
const REPLACEMENT_INTENT_SEAL_FILE: &str = "ahm_replacement_intent_seal";
|
||||
|
||||
@@ -13,26 +13,31 @@
|
||||
// limitations under the License.
|
||||
|
||||
use crate::{Error, Result};
|
||||
use base64::Engine as _;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use sha2::{Digest, Sha256};
|
||||
use std::collections::HashSet;
|
||||
use std::path::Path;
|
||||
use std::sync::{Arc, Mutex};
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
use tokio::sync::RwLock;
|
||||
use tokio::sync::{Mutex as AsyncMutex, RwLock};
|
||||
use tracing::{debug, warn};
|
||||
|
||||
use super::super::{BUCKET_META_PREFIX, DiskStore, HealDiskExt as _, RUSTFS_META_BUCKET};
|
||||
use super::super::storage_api::owner::{EcstoreConditionalFileUpdate, EcstoreDiskAPI, EcstoreDiskBytes};
|
||||
use super::super::{BUCKET_META_PREFIX, DiskStore, HealDiskExt, RUSTFS_META_BUCKET};
|
||||
use super::{
|
||||
LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, PersistThrottle, RESUME_CHECKPOINT_FILE, delete_resume_file, path_to_str,
|
||||
validate_resume_task_id,
|
||||
LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, PersistThrottle, RESUME_CHECKPOINT_BLOCKED_FILE, RESUME_CHECKPOINT_FILE,
|
||||
delete_resume_file, path_to_str, validate_resume_task_id,
|
||||
};
|
||||
|
||||
const EVENT_HEAL_CHECKPOINT_STATE: &str = "heal_checkpoint_state";
|
||||
const RESUME_CHECKPOINT_DIGEST_FILE: &str = "ahm_checkpoint.sha256";
|
||||
const CHECKPOINT_PER_VERSION_SCHEMA: u32 = 5;
|
||||
|
||||
/// Current on-disk schema version for `ResumeCheckpoint`. Same rationale as
|
||||
/// `CURRENT_RESUME_SCHEMA`: pre-per-version dedup identities are not comparable
|
||||
/// to the new `compose_key` identities, so a stale checkpoint is discarded.
|
||||
pub(super) const CURRENT_CHECKPOINT_SCHEMA: u32 = 5;
|
||||
pub(super) const CURRENT_CHECKPOINT_SCHEMA: u32 = 6;
|
||||
|
||||
/// resume checkpoint
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
@@ -57,6 +62,11 @@ pub struct ResumeCheckpoint {
|
||||
pub failed_objects: HashSet<String>,
|
||||
/// skipped objects
|
||||
pub skipped_objects: HashSet<String>,
|
||||
/// Integrity digest over the checkpoint with this field set to `None`.
|
||||
/// Keeping it in the checkpoint makes the payload and its authentication
|
||||
/// record one CAS generation instead of two independently-written files.
|
||||
#[serde(default)]
|
||||
pub integrity_digest: Option<String>,
|
||||
}
|
||||
|
||||
impl ResumeCheckpoint {
|
||||
@@ -70,6 +80,7 @@ impl ResumeCheckpoint {
|
||||
processed_objects: HashSet::new(),
|
||||
failed_objects: HashSet::new(),
|
||||
skipped_objects: HashSet::new(),
|
||||
integrity_digest: None,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -116,17 +127,111 @@ pub struct CheckpointManager {
|
||||
disk: DiskStore,
|
||||
checkpoint: Arc<RwLock<ResumeCheckpoint>>,
|
||||
throttle: Mutex<PersistThrottle>,
|
||||
save_lock: AsyncMutex<()>,
|
||||
last_saved: Mutex<Option<EcstoreDiskBytes>>,
|
||||
}
|
||||
|
||||
impl CheckpointManager {
|
||||
fn blocked_path(task_id: &str) -> std::path::PathBuf {
|
||||
Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_BLOCKED_FILE}"))
|
||||
}
|
||||
|
||||
/// Return whether a checkpoint was permanently isolated after a malformed
|
||||
/// or unsupported snapshot was observed.
|
||||
pub(crate) async fn is_blocked(disk: &DiskStore, task_id: &str) -> bool {
|
||||
if validate_resume_task_id(task_id).is_err() {
|
||||
return false;
|
||||
}
|
||||
let blocked_path = Self::blocked_path(task_id);
|
||||
let Ok(path) = path_to_str(&blocked_path) else {
|
||||
return false;
|
||||
};
|
||||
match HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, path).await {
|
||||
Ok(_) => true,
|
||||
Err(crate::heal::DiskError::FileNotFound) => false,
|
||||
Err(_) => true,
|
||||
}
|
||||
}
|
||||
|
||||
/// Validate the checkpoint while enumerating resumable state. This reads
|
||||
/// the checkpoint once and also isolates malformed or unsupported data.
|
||||
pub(crate) async fn is_resumable(disk: &DiskStore, task_id: &str) -> Result<bool> {
|
||||
validate_resume_task_id(task_id)?;
|
||||
if Self::is_blocked(disk, task_id).await {
|
||||
return Err(Error::InvalidCheckpoint(format!("Resume task {task_id} has a blocked checkpoint")));
|
||||
}
|
||||
let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"));
|
||||
let Ok(path) = path_to_str(&file_path) else {
|
||||
return Err(Error::InvalidCheckpoint("Resume checkpoint path is not valid UTF-8".to_string()));
|
||||
};
|
||||
match HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, path).await {
|
||||
Ok(bytes) if bytes.is_empty() => Ok(true),
|
||||
Ok(bytes) => Self::load_from_data(disk.clone(), task_id, bytes.to_vec())
|
||||
.await
|
||||
.map(|_| true),
|
||||
Err(crate::heal::DiskError::FileNotFound) => Ok(true),
|
||||
Err(error) => Err(error.into()),
|
||||
}
|
||||
}
|
||||
|
||||
async fn block_invalid_snapshot(disk: &DiskStore, task_id: &str) {
|
||||
// This marker is intentionally version-agnostic: an unsupported reader
|
||||
// must stop selector retries until an operator cleans up the snapshot.
|
||||
let blocked_path = Self::blocked_path(task_id);
|
||||
let Ok(path) = path_to_str(&blocked_path) else {
|
||||
return;
|
||||
};
|
||||
let result = EcstoreDiskAPI::compare_and_update_file(
|
||||
disk.as_ref(),
|
||||
RUSTFS_META_BUCKET,
|
||||
path,
|
||||
None,
|
||||
Some(EcstoreDiskBytes::from_static(b"blocked")),
|
||||
)
|
||||
.await;
|
||||
match result {
|
||||
Ok(EcstoreConditionalFileUpdate::Updated | EcstoreConditionalFileUpdate::Mismatch) => {}
|
||||
Ok(EcstoreConditionalFileUpdate::Missing) => warn!(
|
||||
target: "rustfs::heal::resume",
|
||||
event = EVENT_HEAL_CHECKPOINT_STATE,
|
||||
component = LOG_COMPONENT_HEAL,
|
||||
subsystem = LOG_SUBSYSTEM_RESUME,
|
||||
task_id,
|
||||
state = "blocked_marker_write_failed",
|
||||
error = "marker target disappeared",
|
||||
"Heal checkpoint could not persist its blocked marker"
|
||||
),
|
||||
Err(error) => warn!(
|
||||
target: "rustfs::heal::resume",
|
||||
event = EVENT_HEAL_CHECKPOINT_STATE,
|
||||
component = LOG_COMPONENT_HEAL,
|
||||
subsystem = LOG_SUBSYSTEM_RESUME,
|
||||
task_id,
|
||||
state = "blocked_marker_write_failed",
|
||||
error = %error,
|
||||
"Heal checkpoint could not persist its blocked marker"
|
||||
),
|
||||
}
|
||||
}
|
||||
|
||||
/// create new checkpoint manager
|
||||
pub async fn new(disk: DiskStore, task_id: String) -> Result<Self> {
|
||||
validate_resume_task_id(&task_id)?;
|
||||
let checkpoint_volume = format!("{RUSTFS_META_BUCKET}/{BUCKET_META_PREFIX}");
|
||||
if let Err(error) = EcstoreDiskAPI::make_volume(disk.as_ref(), &checkpoint_volume).await
|
||||
&& error != crate::heal::DiskError::VolumeExists
|
||||
{
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Failed to create checkpoint volume: {error}"),
|
||||
});
|
||||
}
|
||||
let checkpoint = ResumeCheckpoint::new(task_id);
|
||||
let manager = Self {
|
||||
disk,
|
||||
checkpoint: Arc::new(RwLock::new(checkpoint)),
|
||||
throttle: Mutex::new(PersistThrottle::new()),
|
||||
save_lock: AsyncMutex::new(()),
|
||||
last_saved: Mutex::new(None),
|
||||
};
|
||||
|
||||
// save initial checkpoint
|
||||
@@ -140,6 +245,7 @@ impl CheckpointManager {
|
||||
error = %e,
|
||||
"Heal checkpoint persistence failed"
|
||||
);
|
||||
return Err(e);
|
||||
}
|
||||
Ok(manager)
|
||||
}
|
||||
@@ -148,11 +254,22 @@ impl CheckpointManager {
|
||||
pub async fn load_from_disk(disk: DiskStore, task_id: &str) -> Result<Self> {
|
||||
validate_resume_task_id(task_id)?;
|
||||
let checkpoint_data = Self::read_checkpoint_file(&disk, task_id).await?;
|
||||
let mut checkpoint: ResumeCheckpoint =
|
||||
serde_json::from_slice(&checkpoint_data).map_err(|e| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to deserialize checkpoint: {e}"),
|
||||
})?;
|
||||
Self::load_from_data(disk, task_id, checkpoint_data).await
|
||||
}
|
||||
|
||||
async fn load_from_data(disk: DiskStore, task_id: &str, checkpoint_data: Vec<u8>) -> Result<Self> {
|
||||
validate_resume_task_id(task_id)?;
|
||||
let mut checkpoint: ResumeCheckpoint = match serde_json::from_slice(&checkpoint_data) {
|
||||
Ok(checkpoint) => checkpoint,
|
||||
Err(error) => {
|
||||
Self::block_invalid_snapshot(&disk, task_id).await;
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Failed to deserialize checkpoint: {error}"),
|
||||
});
|
||||
}
|
||||
};
|
||||
if checkpoint.task_id != task_id {
|
||||
Self::block_invalid_snapshot(&disk, task_id).await;
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: "Resume checkpoint task id does not match filename".to_string(),
|
||||
});
|
||||
@@ -163,6 +280,7 @@ impl CheckpointManager {
|
||||
// identities. Discard the stale sets and position, then stamp the
|
||||
// current schema so the scan restarts cleanly.
|
||||
if checkpoint.schema_version > CURRENT_CHECKPOINT_SCHEMA {
|
||||
Self::block_invalid_snapshot(&disk, task_id).await;
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!(
|
||||
"Checkpoint schema {} is newer than supported schema {CURRENT_CHECKPOINT_SCHEMA}",
|
||||
@@ -170,7 +288,45 @@ impl CheckpointManager {
|
||||
),
|
||||
});
|
||||
}
|
||||
if checkpoint.schema_version < CURRENT_CHECKPOINT_SCHEMA {
|
||||
|
||||
let integrity_verified = if let Some(expected) = checkpoint.integrity_digest.as_deref() {
|
||||
let actual = Self::checkpoint_digest(&Self::serialize_without_digest(&checkpoint)?);
|
||||
if expected != actual {
|
||||
Self::block_invalid_snapshot(&disk, task_id).await;
|
||||
return Err(Error::InvalidCheckpoint(format!(
|
||||
"Resume checkpoint digest does not match task {task_id}"
|
||||
)));
|
||||
}
|
||||
true
|
||||
} else if checkpoint.schema_version >= CURRENT_CHECKPOINT_SCHEMA {
|
||||
Self::block_invalid_snapshot(&disk, task_id).await;
|
||||
return Err(Error::InvalidCheckpoint(format!(
|
||||
"Resume checkpoint digest is missing for task {task_id}"
|
||||
)));
|
||||
} else {
|
||||
let digest_path = Self::digest_path(task_id);
|
||||
let digest_path = path_to_str(&digest_path)?;
|
||||
match HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, digest_path).await {
|
||||
Ok(expected) => {
|
||||
let actual = Self::checkpoint_digest(&checkpoint_data);
|
||||
if expected.as_ref() != actual.as_bytes() {
|
||||
Self::block_invalid_snapshot(&disk, task_id).await;
|
||||
return Err(Error::InvalidCheckpoint(format!(
|
||||
"Resume checkpoint digest does not match task {task_id}"
|
||||
)));
|
||||
}
|
||||
true
|
||||
}
|
||||
Err(crate::heal::DiskError::FileNotFound) => false,
|
||||
Err(error) => {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Failed to read checkpoint digest: {error}"),
|
||||
});
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
if checkpoint.schema_version < CHECKPOINT_PER_VERSION_SCHEMA || !integrity_verified {
|
||||
warn!(
|
||||
target: "rustfs::heal::resume",
|
||||
event = EVENT_HEAL_CHECKPOINT_STATE,
|
||||
@@ -187,13 +343,15 @@ impl CheckpointManager {
|
||||
checkpoint.skipped_objects.clear();
|
||||
checkpoint.current_bucket_index = 0;
|
||||
checkpoint.current_object_index = 0;
|
||||
checkpoint.schema_version = CURRENT_CHECKPOINT_SCHEMA;
|
||||
}
|
||||
checkpoint.schema_version = CURRENT_CHECKPOINT_SCHEMA;
|
||||
|
||||
Ok(Self {
|
||||
disk,
|
||||
checkpoint: Arc::new(RwLock::new(checkpoint)),
|
||||
throttle: Mutex::new(PersistThrottle::new()),
|
||||
save_lock: AsyncMutex::new(()),
|
||||
last_saved: Mutex::new(Some(EcstoreDiskBytes::from(checkpoint_data))),
|
||||
})
|
||||
}
|
||||
|
||||
@@ -204,7 +362,7 @@ impl CheckpointManager {
|
||||
}
|
||||
let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"));
|
||||
match path_to_str(&file_path) {
|
||||
Ok(path_str) => match disk.read_all(RUSTFS_META_BUCKET, path_str).await {
|
||||
Ok(path_str) => match HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, path_str).await {
|
||||
Ok(data) => !data.is_empty(),
|
||||
Err(_) => false,
|
||||
},
|
||||
@@ -292,6 +450,8 @@ impl CheckpointManager {
|
||||
|
||||
let checkpoint_file = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"));
|
||||
delete_resume_file(&self.disk, &checkpoint_file).await?;
|
||||
delete_resume_file(&self.disk, &Self::digest_path(&task_id)).await?;
|
||||
delete_resume_file(&self.disk, &Self::blocked_path(&task_id)).await?;
|
||||
|
||||
debug!(
|
||||
target: "rustfs::heal::resume",
|
||||
@@ -307,21 +467,130 @@ impl CheckpointManager {
|
||||
|
||||
/// save checkpoint to disk
|
||||
async fn save_checkpoint(&self) -> Result<()> {
|
||||
let checkpoint = self.checkpoint.read().await;
|
||||
// Serialize saves and take the snapshot only after acquiring the lock:
|
||||
// a slower writer must not publish a snapshot taken before a newer one.
|
||||
let _save_guard = self.save_lock.lock().await;
|
||||
let checkpoint = self.checkpoint.read().await.clone();
|
||||
validate_resume_task_id(&checkpoint.task_id)?;
|
||||
let checkpoint_data = serde_json::to_vec(&*checkpoint).map_err(|e| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to serialize checkpoint: {e}"),
|
||||
})?;
|
||||
let unsigned_checkpoint_data = Self::serialize_without_digest(&checkpoint)?;
|
||||
let digest = Self::checkpoint_digest(&unsigned_checkpoint_data);
|
||||
let mut persisted_checkpoint = checkpoint.clone();
|
||||
persisted_checkpoint.integrity_digest = Some(digest);
|
||||
let checkpoint_data =
|
||||
EcstoreDiskBytes::from(serde_json::to_vec(&persisted_checkpoint).map_err(|e| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to serialize checkpoint: {e}"),
|
||||
})?);
|
||||
|
||||
let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{}_{}", checkpoint.task_id, RESUME_CHECKPOINT_FILE));
|
||||
|
||||
let path_str = path_to_str(&file_path)?;
|
||||
self.disk
|
||||
.write_all(RUSTFS_META_BUCKET, path_str, checkpoint_data.into())
|
||||
let last_saved = self
|
||||
.last_saved
|
||||
.lock()
|
||||
.map_err(|_| Error::TaskExecutionFailed {
|
||||
message: "Checkpoint save state lock is poisoned; refusing to save".to_string(),
|
||||
})?
|
||||
.clone();
|
||||
let update = EcstoreDiskAPI::compare_and_update_file(
|
||||
self.disk.as_ref(),
|
||||
RUSTFS_META_BUCKET,
|
||||
path_str,
|
||||
last_saved.clone(),
|
||||
Some(checkpoint_data.clone()),
|
||||
)
|
||||
.await
|
||||
.map_err(|e| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to save checkpoint: {e}"),
|
||||
})?;
|
||||
|
||||
let expected = match update {
|
||||
EcstoreConditionalFileUpdate::Updated => None,
|
||||
EcstoreConditionalFileUpdate::Missing => {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: "Checkpoint was removed after this manager saved it; refusing to recreate it".to_string(),
|
||||
});
|
||||
}
|
||||
EcstoreConditionalFileUpdate::Mismatch => {
|
||||
// A healthy manager normally completes the CAS above without
|
||||
// another read or JSON parse. Inspect only after a mismatch so
|
||||
// corruption and future schemas cannot be overwritten blindly.
|
||||
let existing = match HealDiskExt::read_all(self.disk.as_ref(), RUSTFS_META_BUCKET, path_str).await {
|
||||
Ok(existing) => existing,
|
||||
Err(crate::heal::DiskError::FileNotFound) => {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: "Checkpoint was removed after this manager saved it; refusing to recreate it".to_string(),
|
||||
});
|
||||
}
|
||||
Err(error) => {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Failed to inspect checkpoint after CAS mismatch: {error}"),
|
||||
});
|
||||
}
|
||||
};
|
||||
|
||||
if existing.is_empty() && last_saved.is_none() {
|
||||
Some(existing)
|
||||
} else {
|
||||
let current: ResumeCheckpoint = match serde_json::from_slice(&existing) {
|
||||
Ok(current) => current,
|
||||
Err(error) => {
|
||||
Self::block_invalid_snapshot(&self.disk, &checkpoint.task_id).await;
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Existing checkpoint is corrupt: {error}"),
|
||||
});
|
||||
}
|
||||
};
|
||||
if current.task_id != checkpoint.task_id {
|
||||
Self::block_invalid_snapshot(&self.disk, &checkpoint.task_id).await;
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: "Existing checkpoint task id does not match filename".to_string(),
|
||||
});
|
||||
}
|
||||
if current.schema_version > CURRENT_CHECKPOINT_SCHEMA {
|
||||
Self::block_invalid_snapshot(&self.disk, &checkpoint.task_id).await;
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!(
|
||||
"Existing checkpoint schema {} is newer than supported schema {CURRENT_CHECKPOINT_SCHEMA}",
|
||||
current.schema_version
|
||||
),
|
||||
});
|
||||
}
|
||||
if last_saved.as_ref().is_none_or(|saved| saved.as_ref() != existing.as_ref()) {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: "Checkpoint changed since this manager loaded it; refusing to overwrite newer progress"
|
||||
.to_string(),
|
||||
});
|
||||
}
|
||||
Some(existing)
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
if let Some(expected) = expected {
|
||||
match EcstoreDiskAPI::compare_and_update_file(
|
||||
self.disk.as_ref(),
|
||||
RUSTFS_META_BUCKET,
|
||||
path_str,
|
||||
Some(expected),
|
||||
Some(checkpoint_data.clone()),
|
||||
)
|
||||
.await
|
||||
.map_err(|e| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to save checkpoint: {e}"),
|
||||
})?;
|
||||
message: format!("Failed to save checkpoint after CAS mismatch: {e}"),
|
||||
})? {
|
||||
EcstoreConditionalFileUpdate::Updated => {}
|
||||
EcstoreConditionalFileUpdate::Missing | EcstoreConditionalFileUpdate::Mismatch => {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: "Checkpoint changed while saving; refusing to overwrite newer progress".to_string(),
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let mut last_saved = self.last_saved.lock().map_err(|_| Error::TaskExecutionFailed {
|
||||
message: "Checkpoint save state lock is poisoned after save".to_string(),
|
||||
})?;
|
||||
*last_saved = Some(checkpoint_data);
|
||||
|
||||
debug!(
|
||||
target: "rustfs::heal::resume",
|
||||
@@ -341,11 +610,38 @@ impl CheckpointManager {
|
||||
let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"));
|
||||
|
||||
let path_str = path_to_str(&file_path)?;
|
||||
disk.read_all(RUSTFS_META_BUCKET, path_str)
|
||||
HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, path_str)
|
||||
.await
|
||||
.map(|bytes| bytes.to_vec())
|
||||
.map_err(|e| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to read checkpoint file: {e}"),
|
||||
})
|
||||
}
|
||||
|
||||
fn serialize_without_digest(checkpoint: &ResumeCheckpoint) -> Result<Vec<u8>> {
|
||||
let mut unsigned = checkpoint.clone();
|
||||
unsigned.integrity_digest = None;
|
||||
let mut value = serde_json::to_value(&unsigned).map_err(|e| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to serialize checkpoint: {e}"),
|
||||
})?;
|
||||
for field in ["processed_objects", "failed_objects", "skipped_objects"] {
|
||||
let Some(values) = value.get_mut(field).and_then(serde_json::Value::as_array_mut) else {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Failed to canonicalize checkpoint field: {field}"),
|
||||
});
|
||||
};
|
||||
values.sort_by(|left, right| left.as_str().cmp(&right.as_str()));
|
||||
}
|
||||
serde_json::to_vec(&value).map_err(|e| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to serialize checkpoint: {e}"),
|
||||
})
|
||||
}
|
||||
|
||||
fn checkpoint_digest(checkpoint_data: &[u8]) -> String {
|
||||
base64::engine::general_purpose::STANDARD.encode(Sha256::digest(checkpoint_data))
|
||||
}
|
||||
|
||||
fn digest_path(task_id: &str) -> std::path::PathBuf {
|
||||
Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_DIGEST_FILE}"))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1600,6 +1600,32 @@ async fn test_checkpoint_schema_v4_discarded_on_load() {
|
||||
temp_dir.close().expect("remove schema test directory");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn downgraded_unsigned_checkpoint_resets_untrusted_progress() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone()).await.unwrap();
|
||||
manager.add_processed_object("victim-a".to_string()).await.unwrap();
|
||||
manager.update_position(2, 500).await.unwrap();
|
||||
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
|
||||
let bytes = disk.read_all(RUSTFS_META_BUCKET, &checkpoint_path).await.unwrap();
|
||||
let mut downgraded: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
|
||||
downgraded["schema_version"] = serde_json::json!(CURRENT_CHECKPOINT_SCHEMA - 1);
|
||||
downgraded.as_object_mut().unwrap().remove("integrity_digest");
|
||||
downgraded["processed_objects"] = serde_json::json!(["victim-b"]);
|
||||
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, serde_json::to_vec(&downgraded).unwrap().into())
|
||||
.await
|
||||
.expect("write downgraded checkpoint");
|
||||
|
||||
let manager = CheckpointManager::load_from_disk(disk, &task_id).await.unwrap();
|
||||
let checkpoint = manager.get_checkpoint().await;
|
||||
assert_eq!(checkpoint.schema_version, CURRENT_CHECKPOINT_SCHEMA);
|
||||
assert_eq!(checkpoint.current_bucket_index, 0);
|
||||
assert_eq!(checkpoint.current_object_index, 0);
|
||||
assert!(checkpoint.processed_objects.is_empty());
|
||||
temp_dir.close().unwrap();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn current_normal_resume_schema_preserves_progress() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
@@ -1675,6 +1701,369 @@ async fn future_resume_and_checkpoint_schemas_are_rejected() {
|
||||
temp_dir.close().expect("remove schema test directory");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn checkpoint_save_does_not_replace_a_non_empty_truncated_snapshot() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
|
||||
.await
|
||||
.expect("create checkpoint manager");
|
||||
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
|
||||
let truncated = b"{\"schema_version\":5,\"task_id\":";
|
||||
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, truncated.as_slice().into())
|
||||
.await
|
||||
.expect("write truncated checkpoint fixture");
|
||||
|
||||
let error = manager
|
||||
.update_position(2, 7)
|
||||
.await
|
||||
.expect_err("a truncated checkpoint must fail closed during save");
|
||||
assert!(error.to_string().contains("Existing checkpoint is corrupt"));
|
||||
assert_eq!(
|
||||
disk.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
|
||||
.await
|
||||
.expect("read truncated checkpoint fixture"),
|
||||
truncated.as_slice()
|
||||
);
|
||||
assert!(CheckpointManager::is_blocked(&disk, &task_id).await);
|
||||
temp_dir.close().expect("remove checkpoint save test directory");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn checkpoint_save_does_not_replace_a_future_schema_snapshot() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
|
||||
.await
|
||||
.expect("create checkpoint manager");
|
||||
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
|
||||
let mut future = ResumeCheckpoint::new(task_id.clone());
|
||||
future.schema_version = CURRENT_CHECKPOINT_SCHEMA + 1;
|
||||
let future_bytes = serde_json::to_vec(&future).expect("serialize future checkpoint fixture");
|
||||
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, future_bytes.clone().into())
|
||||
.await
|
||||
.expect("write future checkpoint fixture");
|
||||
|
||||
let error = manager
|
||||
.update_position(2, 7)
|
||||
.await
|
||||
.expect_err("a future schema must fail closed during save");
|
||||
assert!(error.to_string().contains("Existing checkpoint schema"));
|
||||
assert_eq!(
|
||||
disk.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
|
||||
.await
|
||||
.expect("read future checkpoint fixture"),
|
||||
future_bytes
|
||||
);
|
||||
assert!(CheckpointManager::is_blocked(&disk, &task_id).await);
|
||||
temp_dir.close().expect("remove future schema test directory");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn checkpoint_digest_rejects_same_length_progress_tampering() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
|
||||
.await
|
||||
.expect("create checkpoint manager");
|
||||
manager
|
||||
.add_processed_object("victim-a".to_string())
|
||||
.await
|
||||
.expect("persist checkpoint progress");
|
||||
manager.update_position(1, 1).await.expect("flush checkpoint progress");
|
||||
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
|
||||
let original = disk
|
||||
.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
|
||||
.await
|
||||
.expect("read checkpoint fixture");
|
||||
let tampered = original
|
||||
.windows(b"victim-a".len())
|
||||
.position(|window| window == b"victim-a")
|
||||
.map(|index| {
|
||||
let mut bytes = original.to_vec();
|
||||
bytes[index..index + b"victim-a".len()].copy_from_slice(b"victim-b");
|
||||
bytes
|
||||
})
|
||||
.expect("checkpoint should contain the processed object");
|
||||
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, tampered.into())
|
||||
.await
|
||||
.expect("write tampered checkpoint fixture");
|
||||
|
||||
assert!(CheckpointManager::load_from_disk(disk.clone(), &task_id).await.is_err());
|
||||
assert!(CheckpointManager::is_blocked(&disk, &task_id).await);
|
||||
temp_dir.close().expect("remove digest test directory");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn checkpoint_integrity_survives_missing_legacy_sidecar() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone()).await.unwrap();
|
||||
manager.update_position(2, 9).await.unwrap();
|
||||
|
||||
let digest_path = format!("{BUCKET_META_PREFIX}/{task_id}_ahm_checkpoint.sha256");
|
||||
delete_resume_file(&disk, Path::new(&digest_path)).await.unwrap();
|
||||
|
||||
let restored = CheckpointManager::load_from_disk(disk, &task_id).await.unwrap();
|
||||
let checkpoint = restored.get_checkpoint().await;
|
||||
assert_eq!(checkpoint.current_bucket_index, 2);
|
||||
assert_eq!(checkpoint.current_object_index, 9);
|
||||
assert!(checkpoint.integrity_digest.is_some());
|
||||
temp_dir.close().unwrap();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn checkpoint_integrity_survives_multi_object_reload() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone()).await.unwrap();
|
||||
for index in 0..32 {
|
||||
manager.add_processed_object(format!("processed-{index}")).await.unwrap();
|
||||
manager.add_failed_object(format!("failed-{index}")).await.unwrap();
|
||||
manager.add_skipped_object(format!("skipped-{index}")).await.unwrap();
|
||||
}
|
||||
manager.update_position(2, 9).await.unwrap();
|
||||
|
||||
CheckpointManager::load_from_disk(disk, &task_id)
|
||||
.await
|
||||
.expect("a healthy multi-object checkpoint must survive reload");
|
||||
temp_dir.close().unwrap();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn checkpoint_integrity_rejects_a_removed_embedded_digest() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone()).await.unwrap();
|
||||
manager.update_position(2, 9).await.unwrap();
|
||||
|
||||
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
|
||||
let bytes = disk
|
||||
.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
|
||||
.await
|
||||
.expect("read checkpoint fixture");
|
||||
let mut value: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
|
||||
value["current_object_index"] = serde_json::json!(10);
|
||||
value.as_object_mut().unwrap().remove("integrity_digest");
|
||||
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, serde_json::to_vec(&value).unwrap().into())
|
||||
.await
|
||||
.expect("write tampered checkpoint fixture");
|
||||
|
||||
assert!(
|
||||
CheckpointManager::load_from_disk(disk.clone(), &task_id).await.is_err(),
|
||||
"a current checkpoint without its embedded digest must fail closed"
|
||||
);
|
||||
assert!(CheckpointManager::is_blocked(&disk, &task_id).await);
|
||||
temp_dir.close().unwrap();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn new_checkpoint_manager_rebuilds_an_empty_snapshot() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
|
||||
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, EcstoreDiskBytes::new())
|
||||
.await
|
||||
.expect("write empty checkpoint fixture");
|
||||
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
|
||||
.await
|
||||
.expect("a new manager must rebuild an empty checkpoint");
|
||||
manager
|
||||
.update_position(3, 11)
|
||||
.await
|
||||
.expect("rebuilt checkpoint must remain writable");
|
||||
assert!(CheckpointManager::has_checkpoint(&disk, &task_id).await);
|
||||
temp_dir.close().expect("remove empty checkpoint test directory");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn deleted_checkpoint_is_not_recreated_by_an_old_manager() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
|
||||
.await
|
||||
.expect("create checkpoint manager");
|
||||
manager.cleanup().await.expect("delete checkpoint fixture");
|
||||
|
||||
let error = manager
|
||||
.update_position(1, 2)
|
||||
.await
|
||||
.expect_err("an old manager must not resurrect a deleted checkpoint");
|
||||
assert!(error.to_string().contains("removed after this manager saved it"));
|
||||
assert!(!CheckpointManager::has_checkpoint(&disk, &task_id).await);
|
||||
temp_dir.close().expect("remove deleted checkpoint test directory");
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[tokio::test]
|
||||
async fn checkpoint_cleanup_leaves_no_task_specific_lock_artifact() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
|
||||
.await
|
||||
.expect("create checkpoint manager");
|
||||
let lock_path = Path::new(BUCKET_META_PREFIX)
|
||||
.join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"))
|
||||
.with_extension("rustfs-cas.lock");
|
||||
let lock_path = temp_dir.path().join(RUSTFS_META_BUCKET).join(lock_path);
|
||||
|
||||
manager.cleanup().await.expect("delete checkpoint fixture");
|
||||
|
||||
assert!(
|
||||
!lock_path.exists(),
|
||||
"successful checkpoint cleanup must not leave a task-specific lock artifact"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn an_empty_blocked_marker_still_blocks_resume_selection() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
|
||||
.await
|
||||
.expect("create checkpoint manager");
|
||||
let blocked_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_BLOCKED_FILE}");
|
||||
disk.write_all(RUSTFS_META_BUCKET, &blocked_path, EcstoreDiskBytes::new())
|
||||
.await
|
||||
.expect("write empty blocked marker fixture");
|
||||
|
||||
assert!(CheckpointManager::is_blocked(&disk, &task_id).await);
|
||||
assert!(CheckpointManager::is_resumable(&disk, &task_id).await.is_err());
|
||||
// Recovery requires replacing/cleaning the snapshot, then removing the
|
||||
// marker; ordinary selector retries are intentionally not an unlock path.
|
||||
manager.cleanup().await.expect("clean blocked checkpoint");
|
||||
assert!(!CheckpointManager::is_blocked(&disk, &task_id).await);
|
||||
temp_dir.close().expect("remove empty blocked marker test directory");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn resumable_selector_skips_healthy_tasks_with_blocked_markers() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let tasks = [
|
||||
(ResumeUtils::generate_task_id(), EcstoreDiskBytes::new()),
|
||||
(ResumeUtils::generate_task_id(), EcstoreDiskBytes::from_static(b"blocked")),
|
||||
];
|
||||
for (task_id, marker) in &tasks {
|
||||
ResumeManager::new(
|
||||
disk.clone(),
|
||||
task_id.clone(),
|
||||
"erasure_set".to_string(),
|
||||
"pool_0_set_0".to_string(),
|
||||
vec!["bucket".to_string()],
|
||||
)
|
||||
.await
|
||||
.expect("create healthy resume state");
|
||||
CheckpointManager::new(disk.clone(), task_id.clone())
|
||||
.await
|
||||
.expect("create healthy checkpoint");
|
||||
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
|
||||
let checkpoint_bytes = disk
|
||||
.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
|
||||
.await
|
||||
.expect("read healthy checkpoint before blocking");
|
||||
let marker_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_BLOCKED_FILE}");
|
||||
disk.write_all(RUSTFS_META_BUCKET, &marker_path, marker.clone())
|
||||
.await
|
||||
.expect("write blocked marker");
|
||||
|
||||
assert!(ResumeUtils::get_resumable_tasks(&disk).await.is_err());
|
||||
assert_eq!(
|
||||
disk.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
|
||||
.await
|
||||
.expect("read healthy checkpoint after blocking"),
|
||||
checkpoint_bytes
|
||||
);
|
||||
}
|
||||
temp_dir.close().expect("remove blocked selector test directory");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn stale_checkpoint_manager_cannot_overwrite_newer_progress() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let first = CheckpointManager::new(disk.clone(), task_id.clone())
|
||||
.await
|
||||
.expect("create first checkpoint manager");
|
||||
let second = CheckpointManager::load_from_disk(disk.clone(), &task_id)
|
||||
.await
|
||||
.expect("load second checkpoint manager");
|
||||
|
||||
second
|
||||
.update_position(4, 20)
|
||||
.await
|
||||
.expect("persist newer checkpoint progress");
|
||||
let error = first
|
||||
.update_position(1, 3)
|
||||
.await
|
||||
.expect_err("stale checkpoint manager must not overwrite newer progress");
|
||||
assert!(error.to_string().contains("newer progress"));
|
||||
|
||||
let persisted = CheckpointManager::load_from_disk(disk.clone(), &task_id)
|
||||
.await
|
||||
.expect("load newer checkpoint progress")
|
||||
.get_checkpoint()
|
||||
.await;
|
||||
assert_eq!(persisted.current_bucket_index, 4);
|
||||
assert_eq!(persisted.current_object_index, 20);
|
||||
temp_dir.close().expect("remove stale manager test directory");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn resumable_selector_isolates_future_and_corrupt_checkpoints() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let future_task = ResumeUtils::generate_task_id();
|
||||
let corrupt_task = ResumeUtils::generate_task_id();
|
||||
for task_id in [&future_task, &corrupt_task] {
|
||||
ResumeManager::new(
|
||||
disk.clone(),
|
||||
task_id.to_string(),
|
||||
"erasure_set".to_string(),
|
||||
"pool_0_set_0".to_string(),
|
||||
vec!["bucket".to_string()],
|
||||
)
|
||||
.await
|
||||
.expect("create resumable state fixture");
|
||||
}
|
||||
|
||||
let future_path = format!("{BUCKET_META_PREFIX}/{future_task}_{RESUME_CHECKPOINT_FILE}");
|
||||
let mut future = ResumeCheckpoint::new(future_task.clone());
|
||||
future.schema_version = CURRENT_CHECKPOINT_SCHEMA + 1;
|
||||
let future_bytes = serde_json::to_vec(&future).expect("serialize future checkpoint fixture");
|
||||
disk.write_all(RUSTFS_META_BUCKET, &future_path, future_bytes.clone().into())
|
||||
.await
|
||||
.expect("write future checkpoint fixture");
|
||||
let corrupt_path = format!("{BUCKET_META_PREFIX}/{corrupt_task}_{RESUME_CHECKPOINT_FILE}");
|
||||
let corrupt_bytes = b"{truncated";
|
||||
disk.write_all(RUSTFS_META_BUCKET, &corrupt_path, corrupt_bytes.as_slice().into())
|
||||
.await
|
||||
.expect("write corrupt checkpoint fixture");
|
||||
|
||||
assert!(CheckpointManager::is_resumable(&disk, &future_task).await.is_err());
|
||||
assert!(CheckpointManager::is_resumable(&disk, &corrupt_task).await.is_err());
|
||||
assert!(ResumeUtils::get_resumable_tasks(&disk).await.is_err());
|
||||
for (task_id, path, bytes) in [
|
||||
(&future_task, future_path, future_bytes),
|
||||
(&corrupt_task, corrupt_path, corrupt_bytes.to_vec()),
|
||||
] {
|
||||
assert_eq!(
|
||||
disk.read_all(RUSTFS_META_BUCKET, &path)
|
||||
.await
|
||||
.expect("read isolated checkpoint bytes"),
|
||||
bytes
|
||||
);
|
||||
let blocked_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_BLOCKED_FILE}");
|
||||
assert!(
|
||||
!disk
|
||||
.read_all(RUSTFS_META_BUCKET, &blocked_path)
|
||||
.await
|
||||
.expect("read checkpoint blocked marker")
|
||||
.is_empty()
|
||||
);
|
||||
}
|
||||
temp_dir.close().expect("remove selector isolation test directory");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_persist_throttle_batches_until_threshold() {
|
||||
let mut throttle = PersistThrottle::new();
|
||||
|
||||
@@ -21,7 +21,7 @@ use uuid::Uuid;
|
||||
use super::super::{BUCKET_META_PREFIX, DiskError, DiskStore, HealDiskExt as _, RUSTFS_META_BUCKET};
|
||||
use super::replacement::{ReplacementPhase, ReplacementRecoveryRecord};
|
||||
use super::{
|
||||
EVENT_HEAL_RESUME_STATE, LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, REPLACEMENT_COMPLETION_PROOF_FILE,
|
||||
CheckpointManager, EVENT_HEAL_RESUME_STATE, LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, REPLACEMENT_COMPLETION_PROOF_FILE,
|
||||
REPLACEMENT_INTENT_FILE, RESUME_STATE_FILE, ResumeManager, ResumeStateFile, is_replacement_intent, path_to_str,
|
||||
replacement_recovery_corruption_for_state_load, replacement_recovery_dir, validate_resume_task_id,
|
||||
};
|
||||
@@ -67,6 +67,7 @@ impl ResumeUtils {
|
||||
// Extract task ID from filename: {task_id}_ahm_resume_state.json
|
||||
if let Some(task_id) = entry.strip_suffix(&format!("_{RESUME_STATE_FILE}"))
|
||||
&& validate_resume_task_id(task_id).is_ok()
|
||||
&& CheckpointManager::is_resumable(disk, task_id).await?
|
||||
{
|
||||
task_ids.push(task_id.to_string());
|
||||
}
|
||||
|
||||
@@ -43,9 +43,6 @@ allow-git = [
|
||||
# RustFS fork carrying presigned expiry and constant-time authentication fixes.
|
||||
# owner: rustfs-maintainers review: 2026-10
|
||||
"https://github.com/rustfs/s3s.git",
|
||||
# MiMalloc fork pinned for hotpath allocation counting support.
|
||||
# owner: houseme review: 2026-10
|
||||
"https://github.com/xonatius/mimalloc_rust.git",
|
||||
]
|
||||
|
||||
[bans]
|
||||
|
||||
@@ -0,0 +1,149 @@
|
||||
# CI gate matrix
|
||||
|
||||
This file is the source of truth for which validation runs on each event, its
|
||||
configured wall-clock budget, and whether it can block a merge. Test taxonomy,
|
||||
naming, and nextest serialization rules remain in [README.md](README.md); e2e
|
||||
membership and counts remain in
|
||||
[e2e-suite-inventory.md](e2e-suite-inventory.md).
|
||||
|
||||
The distinction between **required** and **report-only** is load-bearing:
|
||||
a failing job blocks a merge only when its exact check name is present in the
|
||||
live `main` ruleset. A workflow name, a `merge_group` trigger, or a red PR check
|
||||
does not make a job required by itself.
|
||||
|
||||
## Required merge checks
|
||||
|
||||
The live `main` ruleset (`6436880`) currently requires exactly these contexts:
|
||||
|
||||
| Required context | Producer | Validation |
|
||||
|---|---|---|
|
||||
| `CLA Check` | `.github/workflows/cla.yml` | Contributor agreement |
|
||||
| `Quick Checks` | `.github/workflows/ci.yml` | Formatting and repository guard scripts |
|
||||
| `Test and Lint` | `.github/workflows/ci.yml` | Clippy, workspace nextest excluding `e2e_test`, doctests, and migration proofs |
|
||||
|
||||
For pull requests limited to the paths excluded by the main CI workflow,
|
||||
`.github/workflows/ci-docs-only.yml` reports `Quick Checks` and
|
||||
`Test and Lint` under the same names. It runs the real quick checks and the
|
||||
planning-document guard; it does not claim that Rust compilation or runtime
|
||||
tests ran. Despite the workflow name, these paths also include selected deploy,
|
||||
workflow, and lock files.
|
||||
|
||||
Verify the live rule rather than trusting this snapshot before changing merge
|
||||
policy:
|
||||
|
||||
```bash
|
||||
gh api repos/rustfs/rustfs/rulesets/6436880 \
|
||||
--jq '.rules[] | select(.type == "required_status_checks") | .parameters'
|
||||
```
|
||||
|
||||
The ruleset currently has `strict_required_status_checks_policy=false`.
|
||||
`Continuous Integration` accepts `merge_group` events and runs `e2e-full` for
|
||||
them, but `End-to-End Tests (full merge gate)` is not currently a required
|
||||
context. Therefore the repository is prepared to test a merge-queue SHA, but
|
||||
the workflow alone does not prove that every merge passed that lane.
|
||||
|
||||
## Pull request and merge matrix
|
||||
|
||||
Budgets below are job `timeout-minutes`, not typical runtimes. “Report-only”
|
||||
means the result is visible and actionable but is not in the live required
|
||||
context list.
|
||||
|
||||
| Event | Validation | Budget | Merge status | Reproduction |
|
||||
|---|---|---:|---|---|
|
||||
| PR, non-doc change | `Quick Checks` | 10 min | Required | `make pre-commit` (broader local umbrella) |
|
||||
| PR, non-doc change | `Test and Lint` | 90 min | Required | `cargo nextest run --profile ci --all --exclude e2e_test` |
|
||||
| PR, non-doc change | `Typos` | 10 min | Report-only | `typos` |
|
||||
| PR, non-doc change | `ILM Integration (serial)` | 90 min | Report-only | Use the exact command in `.github/workflows/ci.yml` |
|
||||
| PR, non-doc change | rio-v2 / swift / sftp test-and-lint variants | 90 min each | Report-only | `cargo nextest run` with the workflow's feature set |
|
||||
| PR, non-doc change | `Build RustFS Debug Binary` | 30 min | Report-only; prerequisite for black-box lanes | `cargo build -p rustfs --bins` |
|
||||
| PR, non-doc change | `io_uring Integration (real)` | 30 min | Report-only | `cargo test -p rustfs-ecstore --lib uring_ -- --test-threads=1 --nocapture` |
|
||||
| PR, non-doc change | `End-to-End Tests` (`e2e-smoke` plus `s3s-e2e`) | 30 min | Report-only | `cargo nextest run --profile e2e-smoke -p e2e_test`; then `./scripts/e2e-run.sh ./target/debug/rustfs <data-dir>` |
|
||||
| PR, non-doc change | `S3 Implemented Tests` | 60 min | Report-only | Build `rustfs`, then run `scripts/s3-tests/run.sh` with `DEPLOY_MODE=binary`, `TEST_MODE=single`, and `MAXFAIL=0` |
|
||||
| PR, non-doc change | `S3 Lifecycle Behavior Tests` | 30 min | Report-only | Use the accelerated scanner environment in `.github/workflows/ci.yml` with `scripts/s3-tests/run.sh` |
|
||||
| PR touching dependency or workflow inputs | Cargo Deny / Workflow Pin Report / Dependency Review | 20 / 5 / 30 min | Report-only | `cargo deny check`; `scripts/security/check_workflow_pins.sh` |
|
||||
| PR touching architecture rules or architecture docs | `Architecture Migration Rules` | 10 min | Report-only | `scripts/check_architecture_migration_rules.sh` |
|
||||
| PR touching Nix or workspace manifests | `Nix Build & Check` | 60 min | Report-only | `nix flake check` |
|
||||
| PR limited to main-CI-excluded paths | companion `Quick Checks` and `Test and Lint` | 10 min each | Required | `git diff --check`; `make doc-paths-check` when documentation paths changed |
|
||||
| `merge_group` | Standard CI plus `e2e-full` | 55 min for `e2e-full` | Standard required contexts only; `e2e-full` report-only | `cargo nextest run --profile e2e-full -p e2e_test` |
|
||||
| Push to `main` | Standard CI plus `e2e-full` | 55 min for `e2e-full` | Post-merge detection | Same as `merge_group` |
|
||||
| PR touching fuzz inputs or harness paths | Build plus five 60-second fuzz smoke targets | 60 min build; 30 min per target | Report-only | `MAX_TOTAL_TIME=60 ./scripts/fuzz/run.sh` |
|
||||
| PR touching selected ecstore disk/format paths | `Rename Safety` on Windows | 60 min | Report-only | Run the four `cargo test -p rustfs-ecstore --lib <filter>` commands in `windows-filesystem.yml` on Windows |
|
||||
|
||||
The authoritative e2e filters live in `.config/nextest.toml`; extend a profile
|
||||
instead of adding a second ad-hoc selector. Before a profile runs,
|
||||
`scripts/check_test_wiring.py` compares its exact membership to the committed
|
||||
digest so a silent test drop fails closed.
|
||||
|
||||
## Scheduled and manual validation
|
||||
|
||||
Scheduled lanes are independent fault domains. They do not block a pull
|
||||
request, but their workflow-local gate can fail the run and scheduled failures
|
||||
are routed to the shared failure-issue action. The scheduled-validation
|
||||
watchdog and freshness workflow separately detect incomplete runs and missing
|
||||
schedules.
|
||||
|
||||
| Cadence (UTC unless noted) | Workflow / validation | Budget | Verdict and artifacts | Reproduction |
|
||||
|---|---|---:|---|---|
|
||||
| Daily 02:17 | Fuzz: five nightly corpus targets | 60 min build; 60 min per target | Gate; corpus/crash artifacts, scheduled failure alert | `MAX_TOTAL_TIME=<seconds> ./scripts/fuzz/run.sh` |
|
||||
| Daily 03:17 | MinIO interop (EC + SSE read parity) | 40 min | Gate; scheduled failure alert | Dispatch `minio-interop.yml` or follow its pinned Docker fixture steps |
|
||||
| Daily 04:29 | Replication / cluster-fault / protocol e2e | 45 / 90 / 90 min | Three independent gates; JUnit, membership, and server logs | `cargo nextest run --profile e2e-repl-nightly -p e2e_test`; `--profile e2e-nightly`; `-j 1 --profile e2e-protocols` |
|
||||
| Daily 06:31 | Warp performance A/B | 180 min | Regression budget gate; A/B summaries and server logs | `bash scripts/run_hotpath_warp_abba.sh --help` |
|
||||
| Daily 00:07 Asia/Shanghai (16:07 UTC previous day) | Nightly GNU build and Vault lanes | 150 / 90 / 60 min | Build, live Vault, and HA failover gates | Use the commands and pinned Vault images in `nightly-gnu.yml` |
|
||||
| Daily 03:23 | Security Audit | 20 / 5 min, plus 30 min on PR dependency review | Cargo Deny and workflow-pin gates; scheduled failure alert | `cargo deny check`; `scripts/security/check_workflow_pins.sh` |
|
||||
| Daily 23:47 | Scheduled Validation Freshness | 10 min | Fails when a critical schedule was never created or is stale | Dispatch `scheduled-validation-freshness.yml` |
|
||||
| Sunday 00:11 | Full `Continuous Integration` matrix | Per-job budgets above | Weekly variant coverage, including dormant rio-v2 binary/e2e lanes | Dispatch `ci.yml` |
|
||||
| Sunday 01:13 | Seven-platform build matrix | 150 min per platform | Build/package integrity; scheduled failure alert | Dispatch `build.yml` with an exact platform set |
|
||||
| Sunday 02:19 | Ceph s3-tests full sweep: single and real four-node, four shards each | 180 min per shard | Compatibility gate; report, JUnit, exact node IDs, and server logs | `scripts/s3-tests/run.sh` against an existing single or distributed target |
|
||||
| Sunday 06:41 | Mint | 120 min | **Report-only by design**; per-suite PASS/FAIL/NA and raw `log.json` | Reproduce the pinned Docker sequence in `mint.yml` or dispatch it |
|
||||
| Sunday 07:43 | Workspace line coverage | 120 min | Report-only trend; lcov and JSON retained 90 days | `make coverage` |
|
||||
| Monthly, day 1 06:37 | Runner Hygiene | 15 min | Validates runner ephemerality; scheduled failure alert | Dispatch `runner-hygiene.yml` |
|
||||
|
||||
Manual `workflow_dispatch` exists for the scheduled workflows above. Manual
|
||||
runs are debugging evidence and intentionally do not open scheduled-failure
|
||||
issues. A manual performance run may explicitly allow a known regression; that
|
||||
override must not be treated as an ordinary passing baseline.
|
||||
|
||||
## Release validation
|
||||
|
||||
Release validation is post-merge and tag-driven; it does not substitute for a
|
||||
pull-request gate.
|
||||
|
||||
| Event | Validation | Budget | Result |
|
||||
|---|---|---:|---|
|
||||
| Push to `main` or weekly schedule | `Build and Release` platform matrix | 150 min per platform | Build artifacts for all selected targets; no release publication on a main push |
|
||||
| Valid release or preview tag | `Build and Release` plus asset checks | 150 min per platform | Draft release, checksummed assets, and publish step |
|
||||
| Successful non-preview release-tag build | Docker image build and image scan | 60 min build; 30 min scan | Multi-architecture images plus vulnerability report |
|
||||
| Successful release-tag build | DEB/RPM packaging | 30 min per architecture | Packages and checksum files uploaded to the release |
|
||||
| Successful non-preview release-tag build | Helm template test and package | 30 min build; 30 min publish | Versioned chart and repository index |
|
||||
|
||||
Use an exact preview tag for end-to-end release rehearsal. Manual dispatches
|
||||
are backfill/debug paths and do not prove the automatic `workflow_run` chain.
|
||||
|
||||
## Evidence requirements
|
||||
|
||||
A green check is useful only when it proves the intended behavior ran:
|
||||
|
||||
- Record the exact commit SHA and run URL.
|
||||
- Separate product failure from runner prerequisites, service readiness, and
|
||||
cancellation. Repair the precondition, then rerun the exact workload.
|
||||
- Preserve membership manifests, JUnit, raw compatibility logs, seeds, and
|
||||
server logs where the workflow provides them.
|
||||
- For a bug fix or a new fault checker, provide sensitivity evidence: the old
|
||||
behavior or an intentional mutation must fail the new oracle, and the fixed
|
||||
behavior must pass it.
|
||||
- Never promote a report-only lane to required from one green run. Require at
|
||||
least 14 days and 30 representative pull requests with at least 99% complete
|
||||
execution, then update the ruleset and this table together.
|
||||
|
||||
## Change checklist
|
||||
|
||||
Update this file in the same pull request when any of these change:
|
||||
|
||||
- workflow triggers, job names, timeouts, or nextest profile ownership;
|
||||
- required status contexts or strict/merge-queue policy;
|
||||
- scheduled cadence, alert routing, artifact contract, or local reproduction;
|
||||
- report-only versus gating semantics.
|
||||
|
||||
Do not copy per-module test counts here. Update
|
||||
[e2e-suite-inventory.md](e2e-suite-inventory.md) and its enforced membership
|
||||
digest instead.
|
||||
+2
-2
@@ -336,13 +336,13 @@ opentelemetry = { workspace = true }
|
||||
tracing-opentelemetry = { workspace = true }
|
||||
# Data structures
|
||||
hashbrown = { workspace = true, features = ["serde", "rayon"] }
|
||||
mimalloc = { workspace = true }
|
||||
rustfs-mimalloc = { workspace = true }
|
||||
|
||||
[target.'cfg(target_os = "linux")'.dependencies]
|
||||
libsystemd.workspace = true
|
||||
|
||||
[target.'cfg(not(target_os = "windows"))'.dependencies]
|
||||
libmimalloc-sys.workspace = true
|
||||
rustfs-mimalloc-sys.workspace = true
|
||||
|
||||
[dev-dependencies]
|
||||
uuid = { workspace = true, features = ["v4", "v5", "fast-rng", "macro-diagnostics"] }
|
||||
|
||||
@@ -369,14 +369,8 @@ pub fn allocator_reclaim_controller_snapshot(ctx: &CancellationToken) -> Allocat
|
||||
}
|
||||
|
||||
#[cfg(not(target_os = "windows"))]
|
||||
#[allow(unsafe_code)]
|
||||
fn collect_allocator_memory(force: bool) -> Result<(), String> {
|
||||
// SAFETY: `mi_collect` is provided by the active global allocator backend
|
||||
// on this target family. It is explicitly intended to reclaim retained
|
||||
// pages/segments and does not require additional invariants from the caller.
|
||||
unsafe {
|
||||
libmimalloc_sys::mi_collect(force);
|
||||
}
|
||||
rustfs_mimalloc::MiMalloc::collect(force);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
||||
+187
-44
@@ -19,19 +19,23 @@ use std::sync::atomic::{AtomicU64, Ordering};
|
||||
use std::time::Duration;
|
||||
|
||||
use chrono::{DateTime, SecondsFormat, Utc};
|
||||
use reqwest::{Client, StatusCode, Url, header};
|
||||
use rustls::RootCertStore;
|
||||
use rustls::pki_types::{CertificateDer, pem::PemObject as _};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use uuid::Uuid;
|
||||
use zeroize::Zeroizing;
|
||||
|
||||
use super::config::HeartbeatConfig;
|
||||
use super::credential_store::CredentialStoreError;
|
||||
use super::credential_store::{CredentialStoreError, DeviceCredential};
|
||||
use super::identity::IdentityError;
|
||||
use super::identity_store::StoreError;
|
||||
use super::registration::CredentialValidationError;
|
||||
use super::telemetry::{TelemetryDelivery, TelemetryError, TelemetryTransport, is_exact_utc_seconds};
|
||||
use super::registration::{CredentialValidationError, validate_stored_credential};
|
||||
|
||||
const PROTOCOL_VERSION: &str = "v1";
|
||||
const AGENT_VERSION: &str = concat!("rustfs-agent/", env!("CARGO_PKG_VERSION"));
|
||||
const MAX_SEQUENCE: u64 = 9_007_199_254_740_991;
|
||||
const MAX_RESPONSE_BYTES: usize = 64 * 1024;
|
||||
#[cfg(unix)]
|
||||
const FILE_MODE: u32 = 0o600;
|
||||
static STAGING_SEQUENCE: AtomicU64 = AtomicU64::new(0);
|
||||
@@ -118,39 +122,133 @@ pub(crate) enum Delivery {
|
||||
}
|
||||
|
||||
pub(crate) struct HeartbeatSender {
|
||||
transport: TelemetryTransport,
|
||||
endpoint: Url,
|
||||
root_store: RootCertStore,
|
||||
roots: Vec<CertificateDer<'static>>,
|
||||
config: HeartbeatConfig,
|
||||
}
|
||||
|
||||
impl HeartbeatSender {
|
||||
pub(crate) fn new(config: HeartbeatConfig) -> Result<Self, HeartbeatError> {
|
||||
let mut endpoint = Url::parse(&config.endpoint).map_err(|_| HeartbeatError::Endpoint)?;
|
||||
if endpoint.scheme() != "https"
|
||||
|| endpoint.cannot_be_a_base()
|
||||
|| !endpoint.username().is_empty()
|
||||
|| endpoint.password().is_some()
|
||||
|| endpoint.query().is_some()
|
||||
|| endpoint.fragment().is_some()
|
||||
{
|
||||
return Err(HeartbeatError::Endpoint);
|
||||
}
|
||||
if !endpoint.path().ends_with('/') {
|
||||
endpoint.set_path(&format!("{}/", endpoint.path()));
|
||||
}
|
||||
let roots = CertificateDer::pem_slice_iter(&config.root_ca_pem)
|
||||
.collect::<Result<Vec<_>, _>>()
|
||||
.map_err(|_| HeartbeatError::RootCertificate)?;
|
||||
if roots.is_empty() {
|
||||
return Err(HeartbeatError::RootCertificate);
|
||||
}
|
||||
let mut root_store = RootCertStore::empty();
|
||||
let (accepted, rejected) = root_store.add_parsable_certificates(roots.clone());
|
||||
if accepted != roots.len() || rejected != 0 {
|
||||
return Err(HeartbeatError::RootCertificate);
|
||||
}
|
||||
let schedule = config.schedule;
|
||||
if schedule.cadence.is_zero() || schedule.jitter > schedule.cadence {
|
||||
if schedule.cadence.is_zero()
|
||||
|| schedule.timeout.is_zero()
|
||||
|| schedule.timeout > Duration::from_secs(5)
|
||||
|| schedule.initial_backoff.is_zero()
|
||||
|| schedule.max_backoff < schedule.initial_backoff
|
||||
|| schedule.max_backoff > Duration::from_secs(5 * 60)
|
||||
|| schedule.jitter > schedule.cadence
|
||||
{
|
||||
return Err(HeartbeatError::Schedule);
|
||||
}
|
||||
Ok(Self {
|
||||
transport: TelemetryTransport::new(config)?,
|
||||
endpoint,
|
||||
root_store,
|
||||
roots,
|
||||
config,
|
||||
})
|
||||
}
|
||||
|
||||
pub(crate) async fn send(&self, heartbeat: &PendingHeartbeat) -> Result<Delivery, HeartbeatError> {
|
||||
match self.transport.post("heartbeats", heartbeat).await? {
|
||||
TelemetryDelivery::Accepted { body, .. } => {
|
||||
let accepted: HeartbeatResponse = serde_json::from_slice(&body).map_err(|_| HeartbeatError::Response)?;
|
||||
if accepted.accepted_version != PROTOCOL_VERSION
|
||||
|| accepted.capability_hints.len() > 32
|
||||
|| accepted.capability_hints.iter().any(|hint| hint.len() > 32)
|
||||
|| !is_exact_utc_seconds(&accepted.server_time)
|
||||
{
|
||||
return Err(HeartbeatError::Response);
|
||||
}
|
||||
Ok(Delivery::Accepted {
|
||||
server_time: accepted.server_time,
|
||||
})
|
||||
let (cluster_uid, client) = {
|
||||
let _lock = self.config.credential_store.lock().await?;
|
||||
let credential = self.config.credential_store.load()?.ok_or(HeartbeatError::NotRegistered)?;
|
||||
let identity = self.config.identity_store.load()?.ok_or(HeartbeatError::IdentityMissing)?;
|
||||
validate_stored_credential(&credential, &identity, &self.root_store, &self.roots)?;
|
||||
let now = Utc::now().timestamp();
|
||||
if now < credential.not_before_unix || now >= credential.not_after_unix {
|
||||
return Err(HeartbeatError::CredentialExpired);
|
||||
}
|
||||
TelemetryDelivery::Retry { retry_after } => Ok(Delivery::Retry { retry_after }),
|
||||
TelemetryDelivery::AuthenticationStopped { status, reason } => Ok(Delivery::AuthenticationStopped { status, reason }),
|
||||
TelemetryDelivery::Rejected { status, reason } => Ok(Delivery::Rejected { status, reason }),
|
||||
let cluster_uid = cluster_uid(&credential)?.to_owned();
|
||||
let client = self.client(&credential, &identity.to_pkcs8_pem()?)?;
|
||||
(cluster_uid, client)
|
||||
};
|
||||
let url = self.endpoint.join(&format!("clusters/{cluster_uid}/heartbeats"))?;
|
||||
let response = match client.post(url).json(heartbeat).send().await {
|
||||
Ok(response) => response,
|
||||
Err(error) if error.is_timeout() || error.is_connect() || error.is_request() => {
|
||||
return Ok(Delivery::Retry { retry_after: None });
|
||||
}
|
||||
Err(error) => return Err(error.into()),
|
||||
};
|
||||
let status = response.status();
|
||||
if status == StatusCode::TOO_MANY_REQUESTS {
|
||||
return Ok(Delivery::Retry {
|
||||
retry_after: retry_after(response.headers(), Utc::now(), self.config.schedule.max_backoff),
|
||||
});
|
||||
}
|
||||
if status == StatusCode::REQUEST_TIMEOUT || status.is_server_error() {
|
||||
return Ok(Delivery::Retry { retry_after: None });
|
||||
}
|
||||
if matches!(status, StatusCode::UNAUTHORIZED | StatusCode::FORBIDDEN) {
|
||||
return Ok(Delivery::AuthenticationStopped {
|
||||
status: status.as_u16(),
|
||||
reason: response_reason(response).await,
|
||||
});
|
||||
}
|
||||
if status != StatusCode::OK {
|
||||
return Ok(Delivery::Rejected {
|
||||
status: status.as_u16(),
|
||||
reason: response_reason(response).await,
|
||||
});
|
||||
}
|
||||
let accepted: HeartbeatResponse =
|
||||
serde_json::from_slice(&bounded_body(response).await?).map_err(|_| HeartbeatError::Response)?;
|
||||
if accepted.accepted_version != PROTOCOL_VERSION
|
||||
|| accepted.capability_hints.len() > 32
|
||||
|| accepted.capability_hints.iter().any(|hint| hint.len() > 32)
|
||||
|| !is_exact_utc_seconds(&accepted.server_time)
|
||||
{
|
||||
return Err(HeartbeatError::Response);
|
||||
}
|
||||
Ok(Delivery::Accepted {
|
||||
server_time: accepted.server_time,
|
||||
})
|
||||
}
|
||||
|
||||
fn client(&self, credential: &DeviceCredential, key: &Zeroizing<String>) -> Result<Client, HeartbeatError> {
|
||||
let mut pem = Zeroizing::new(Vec::with_capacity(credential.certificate_chain.len() + key.len() + 1));
|
||||
pem.extend_from_slice(credential.certificate_chain.as_bytes());
|
||||
pem.push(b'\n');
|
||||
pem.extend_from_slice(key.as_bytes());
|
||||
let identity = reqwest::Identity::from_pem(&pem).map_err(|_| HeartbeatError::IdentityCertificate)?;
|
||||
let roots = self
|
||||
.roots
|
||||
.iter()
|
||||
.map(|root| reqwest::Certificate::from_der(root.as_ref()))
|
||||
.collect::<Result<Vec<_>, _>>()?;
|
||||
Client::builder()
|
||||
.https_only(true)
|
||||
.redirect(reqwest::redirect::Policy::none())
|
||||
.timeout(self.config.schedule.timeout)
|
||||
.tls_certs_only(roots)
|
||||
.identity(identity)
|
||||
.build()
|
||||
.map_err(Into::into)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -280,6 +378,73 @@ impl HeartbeatStateStore {
|
||||
}
|
||||
}
|
||||
|
||||
fn cluster_uid(credential: &DeviceCredential) -> Result<&str, HeartbeatError> {
|
||||
let mut parts = credential.name.split('/');
|
||||
let valid = parts.next() == Some("organizations");
|
||||
let organization_uid = parts.next();
|
||||
let valid = valid && parts.next() == Some("clusters");
|
||||
let cluster_uid = parts.next();
|
||||
let valid = valid && parts.next() == Some("clusterDevices");
|
||||
let device_uid = parts.next();
|
||||
if !valid
|
||||
|| organization_uid.is_none_or(str::is_empty)
|
||||
|| cluster_uid.is_none_or(str::is_empty)
|
||||
|| device_uid != Some(credential.uid.as_str())
|
||||
|| parts.next().is_some()
|
||||
{
|
||||
return Err(HeartbeatError::CredentialName);
|
||||
}
|
||||
cluster_uid.ok_or(HeartbeatError::CredentialName)
|
||||
}
|
||||
|
||||
fn retry_after(headers: &header::HeaderMap, now: DateTime<Utc>, maximum: Duration) -> Option<Duration> {
|
||||
let value = headers.get(header::RETRY_AFTER)?.to_str().ok()?;
|
||||
let delay = value.parse::<u64>().ok().map(Duration::from_secs).or_else(|| {
|
||||
DateTime::parse_from_rfc2822(value)
|
||||
.ok()
|
||||
.and_then(|at| (at.with_timezone(&Utc) - now).to_std().ok())
|
||||
})?;
|
||||
Some(delay.min(maximum))
|
||||
}
|
||||
|
||||
fn is_exact_utc_seconds(value: &str) -> bool {
|
||||
DateTime::parse_from_rfc3339(value).is_ok_and(|time| {
|
||||
time.offset().local_minus_utc() == 0
|
||||
&& value.ends_with('Z')
|
||||
&& time.with_timezone(&Utc).to_rfc3339_opts(SecondsFormat::Secs, true) == value
|
||||
})
|
||||
}
|
||||
|
||||
async fn response_reason(response: reqwest::Response) -> Option<String> {
|
||||
#[derive(Deserialize)]
|
||||
struct Envelope {
|
||||
#[serde(default)]
|
||||
details: Vec<Detail>,
|
||||
}
|
||||
#[derive(Deserialize)]
|
||||
struct Detail {
|
||||
#[serde(default)]
|
||||
reason: String,
|
||||
}
|
||||
|
||||
serde_json::from_slice::<Envelope>(&bounded_body(response).await.ok()?)
|
||||
.ok()?
|
||||
.details
|
||||
.into_iter()
|
||||
.find_map(|detail| (!detail.reason.is_empty()).then_some(detail.reason))
|
||||
}
|
||||
|
||||
async fn bounded_body(mut response: reqwest::Response) -> Result<Vec<u8>, HeartbeatError> {
|
||||
let mut body = Vec::new();
|
||||
while let Some(chunk) = response.chunk().await? {
|
||||
if body.len().saturating_add(chunk.len()) > MAX_RESPONSE_BYTES {
|
||||
return Err(HeartbeatError::ResponseTooLarge);
|
||||
}
|
||||
body.extend_from_slice(&chunk);
|
||||
}
|
||||
Ok(body)
|
||||
}
|
||||
|
||||
fn parent(path: &Path) -> Result<&Path, HeartbeatError> {
|
||||
path.parent()
|
||||
.ok_or_else(|| state_io(path, io::Error::new(io::ErrorKind::InvalidInput, "state path has no parent")))
|
||||
@@ -418,25 +583,3 @@ pub enum HeartbeatError {
|
||||
#[error(transparent)]
|
||||
CredentialValidation(#[from] CredentialValidationError),
|
||||
}
|
||||
|
||||
impl From<TelemetryError> for HeartbeatError {
|
||||
fn from(error: TelemetryError) -> Self {
|
||||
match error {
|
||||
TelemetryError::Endpoint => Self::Endpoint,
|
||||
TelemetryError::RootCertificate => Self::RootCertificate,
|
||||
TelemetryError::Schedule => Self::Schedule,
|
||||
TelemetryError::NotRegistered => Self::NotRegistered,
|
||||
TelemetryError::IdentityMissing => Self::IdentityMissing,
|
||||
TelemetryError::IdentityCertificate => Self::IdentityCertificate,
|
||||
TelemetryError::CredentialName => Self::CredentialName,
|
||||
TelemetryError::CredentialExpired => Self::CredentialExpired,
|
||||
TelemetryError::ResponseTooLarge => Self::ResponseTooLarge,
|
||||
TelemetryError::Url(error) => Self::Url(error),
|
||||
TelemetryError::Transport(error) => Self::Transport(error),
|
||||
TelemetryError::Identity(error) => Self::Identity(error),
|
||||
TelemetryError::IdentityStore(error) => Self::IdentityStore(error),
|
||||
TelemetryError::CredentialStore(error) => Self::CredentialStore(error),
|
||||
TelemetryError::CredentialValidation(error) => Self::CredentialValidation(error),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,594 +0,0 @@
|
||||
// Copyright 2024 RustFS Team
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
use std::collections::BTreeSet;
|
||||
use std::fs;
|
||||
use std::io::{self, Write as _};
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
use std::time::Duration;
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
use sha2::{Digest as _, Sha256};
|
||||
use uuid::Uuid;
|
||||
|
||||
use super::config::HeartbeatConfig;
|
||||
use super::telemetry::{TelemetryDelivery, TelemetryError, TelemetryTransport, is_exact_utc_seconds};
|
||||
|
||||
const PROTOCOL_VERSION: &str = "v1";
|
||||
const RUSTFS_VERSION: &str = concat!(
|
||||
env!("CARGO_PKG_VERSION_MAJOR"),
|
||||
".",
|
||||
env!("CARGO_PKG_VERSION_MINOR"),
|
||||
".",
|
||||
env!("CARGO_PKG_VERSION_PATCH")
|
||||
);
|
||||
const HASH_PREFIX: &[u8] = b"rustfs-connect/agent/v1/inventory-snapshot\n";
|
||||
const MAX_SEQUENCE: u64 = 9_007_199_254_740_991;
|
||||
const MAX_SAFE_INTEGER: u64 = 9_007_199_254_740_991;
|
||||
#[cfg(unix)]
|
||||
const FILE_MODE: u32 = 0o600;
|
||||
static STAGING_SEQUENCE: AtomicU64 = AtomicU64::new(0);
|
||||
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
pub struct InventorySchedule {
|
||||
pub cadence: Duration,
|
||||
pub jitter: Duration,
|
||||
}
|
||||
|
||||
impl Default for InventorySchedule {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
cadence: Duration::from_secs(6 * 60 * 60),
|
||||
jitter: Duration::from_secs(30 * 60),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, PartialEq, Eq)]
|
||||
pub enum InventoryStatus {
|
||||
Starting,
|
||||
Unchanged { content_hash: String },
|
||||
Online { content_hash: String, received_at: String },
|
||||
BackingOff { delay: Duration },
|
||||
AuthenticationStopped { status: u16, reason: Option<String> },
|
||||
Failed { reason: String },
|
||||
Stopped,
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
|
||||
pub enum InventoryFlag {
|
||||
#[serde(rename = "capacity.critical")]
|
||||
CapacityCritical,
|
||||
#[serde(rename = "capacity.warning")]
|
||||
CapacityWarning,
|
||||
#[serde(rename = "clock.skew")]
|
||||
ClockSkew,
|
||||
#[serde(rename = "cluster.degraded")]
|
||||
ClusterDegraded,
|
||||
#[serde(rename = "cluster.healing")]
|
||||
ClusterHealing,
|
||||
#[serde(rename = "cluster.readonly")]
|
||||
ClusterReadonly,
|
||||
#[serde(rename = "drive.offline")]
|
||||
DriveOffline,
|
||||
#[serde(rename = "node.offline")]
|
||||
NodeOffline,
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "lowercase")]
|
||||
pub enum OperatingSystemFamily {
|
||||
Linux,
|
||||
Darwin,
|
||||
Windows,
|
||||
Freebsd,
|
||||
Other,
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields, rename_all = "camelCase")]
|
||||
pub struct InventoryOsVersion {
|
||||
family: OperatingSystemFamily,
|
||||
major: u16,
|
||||
minor: u16,
|
||||
}
|
||||
|
||||
impl InventoryOsVersion {
|
||||
pub fn new(family: OperatingSystemFamily, major: u16, minor: u16) -> Result<Self, InventoryError> {
|
||||
if major > 9999 || minor > 9999 {
|
||||
return Err(InventoryError::OsVersion);
|
||||
}
|
||||
Ok(Self { family, major, minor })
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields, rename_all = "camelCase")]
|
||||
pub struct InventorySnapshot {
|
||||
rustfs_version: String,
|
||||
os_version: Option<InventoryOsVersion>,
|
||||
node_count: u16,
|
||||
drive_count: u32,
|
||||
capacity_total_bytes: u64,
|
||||
capacity_used_bytes: u64,
|
||||
coarse_flags: Vec<InventoryFlag>,
|
||||
}
|
||||
|
||||
impl InventorySnapshot {
|
||||
pub fn current(
|
||||
node_count: usize,
|
||||
drive_count: usize,
|
||||
capacity_total_bytes: u64,
|
||||
capacity_free_bytes: u64,
|
||||
coarse_flags: impl IntoIterator<Item = InventoryFlag>,
|
||||
) -> Result<Self, InventoryError> {
|
||||
let capacity_used_bytes = capacity_total_bytes
|
||||
.checked_sub(capacity_free_bytes)
|
||||
.ok_or(InventoryError::Capacity)?;
|
||||
Self::new(
|
||||
RUSTFS_VERSION,
|
||||
None,
|
||||
node_count,
|
||||
drive_count,
|
||||
capacity_total_bytes,
|
||||
capacity_used_bytes,
|
||||
coarse_flags,
|
||||
)
|
||||
}
|
||||
|
||||
pub fn new(
|
||||
rustfs_version: impl Into<String>,
|
||||
os_version: Option<InventoryOsVersion>,
|
||||
node_count: usize,
|
||||
drive_count: usize,
|
||||
capacity_total_bytes: u64,
|
||||
capacity_used_bytes: u64,
|
||||
coarse_flags: impl IntoIterator<Item = InventoryFlag>,
|
||||
) -> Result<Self, InventoryError> {
|
||||
let snapshot = Self {
|
||||
rustfs_version: rustfs_version.into(),
|
||||
os_version,
|
||||
node_count: u16::try_from(node_count).map_err(|_| InventoryError::NodeCount)?,
|
||||
drive_count: u32::try_from(drive_count).map_err(|_| InventoryError::DriveCount)?,
|
||||
capacity_total_bytes,
|
||||
capacity_used_bytes,
|
||||
coarse_flags: coarse_flags.into_iter().collect::<BTreeSet<_>>().into_iter().collect(),
|
||||
};
|
||||
snapshot.validate()?;
|
||||
Ok(snapshot)
|
||||
}
|
||||
|
||||
pub fn content_hash(&self) -> Result<String, InventoryError> {
|
||||
#[derive(Serialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
struct Canonical<'a> {
|
||||
capacity_total_bytes: u64,
|
||||
capacity_used_bytes: u64,
|
||||
coarse_flags: &'a [InventoryFlag],
|
||||
drive_count: u32,
|
||||
node_count: u16,
|
||||
os_version: Option<InventoryOsVersion>,
|
||||
rustfs_version: &'a str,
|
||||
}
|
||||
|
||||
let canonical = serde_json::to_vec(&Canonical {
|
||||
capacity_total_bytes: self.capacity_total_bytes,
|
||||
capacity_used_bytes: self.capacity_used_bytes,
|
||||
coarse_flags: &self.coarse_flags,
|
||||
drive_count: self.drive_count,
|
||||
node_count: self.node_count,
|
||||
os_version: self.os_version,
|
||||
rustfs_version: &self.rustfs_version,
|
||||
})?;
|
||||
let mut digest = Sha256::new();
|
||||
digest.update(HASH_PREFIX);
|
||||
digest.update(canonical);
|
||||
Ok(hex_simd::encode_to_string(digest.finalize(), hex_simd::AsciiCase::Lower))
|
||||
}
|
||||
|
||||
fn validate(&self) -> Result<(), InventoryError> {
|
||||
if !valid_version(&self.rustfs_version) {
|
||||
return Err(InventoryError::RustfsVersion);
|
||||
}
|
||||
if self.node_count == 0 || self.node_count > 4096 {
|
||||
return Err(InventoryError::NodeCount);
|
||||
}
|
||||
if self.drive_count > 1_048_576 {
|
||||
return Err(InventoryError::DriveCount);
|
||||
}
|
||||
if self.capacity_total_bytes > MAX_SAFE_INTEGER || self.capacity_used_bytes > self.capacity_total_bytes {
|
||||
return Err(InventoryError::Capacity);
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
fn valid_version(version: &str) -> bool {
|
||||
let components = version.split('.').collect::<Vec<_>>();
|
||||
components.len() == 3
|
||||
&& components.iter().all(|component| {
|
||||
!component.is_empty()
|
||||
&& component.len() <= 4
|
||||
&& (component == &"0" || !component.starts_with('0'))
|
||||
&& component.parse::<u16>().is_ok_and(|value| value <= 9999)
|
||||
})
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields, rename_all = "camelCase")]
|
||||
pub(crate) struct PendingInventory {
|
||||
protocol_version: String,
|
||||
request_id: String,
|
||||
sequence: u64,
|
||||
#[serde(flatten)]
|
||||
snapshot: InventorySnapshot,
|
||||
}
|
||||
|
||||
impl PendingInventory {
|
||||
fn new(snapshot: InventorySnapshot, sequence: u64) -> Self {
|
||||
Self {
|
||||
protocol_version: PROTOCOL_VERSION.to_owned(),
|
||||
request_id: Uuid::new_v4().to_string(),
|
||||
sequence,
|
||||
snapshot,
|
||||
}
|
||||
}
|
||||
|
||||
fn is_valid(&self) -> bool {
|
||||
self.protocol_version == PROTOCOL_VERSION
|
||||
&& self.sequence <= MAX_SEQUENCE
|
||||
&& self.snapshot.validate().is_ok()
|
||||
&& Uuid::parse_str(&self.request_id)
|
||||
.is_ok_and(|request_id| request_id.get_version_num() == 4 && request_id.to_string() == self.request_id)
|
||||
}
|
||||
|
||||
fn content_hash(&self) -> Result<String, InventoryError> {
|
||||
self.snapshot.content_hash()
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) enum InventoryDelivery {
|
||||
Accepted { content_hash: String, received_at: String },
|
||||
Retry { retry_after: Option<Duration> },
|
||||
AuthenticationStopped { status: u16, reason: Option<String> },
|
||||
Rejected { status: u16, reason: Option<String> },
|
||||
}
|
||||
|
||||
pub(crate) struct InventorySender {
|
||||
transport: TelemetryTransport,
|
||||
}
|
||||
|
||||
impl InventorySender {
|
||||
pub(crate) fn new(config: HeartbeatConfig) -> Result<Self, InventoryError> {
|
||||
Ok(Self {
|
||||
transport: TelemetryTransport::new(config)?,
|
||||
})
|
||||
}
|
||||
|
||||
pub(crate) async fn send(&self, inventory: &PendingInventory) -> Result<InventoryDelivery, InventoryError> {
|
||||
match self.transport.post("inventorySnapshots", inventory).await? {
|
||||
TelemetryDelivery::Accepted { cluster_name, body } => {
|
||||
#[derive(Deserialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
struct InventoryResponse {
|
||||
name: String,
|
||||
uid: String,
|
||||
content_hash: String,
|
||||
received_at: String,
|
||||
}
|
||||
|
||||
let accepted: InventoryResponse = serde_json::from_slice(&body).map_err(|_| InventoryError::Response)?;
|
||||
let uid = Uuid::parse_str(&accepted.uid).map_err(|_| InventoryError::Response)?;
|
||||
let content_hash = inventory.content_hash()?;
|
||||
if uid.get_version_num() != 7
|
||||
|| uid.to_string() != accepted.uid
|
||||
|| accepted.name != format!("{cluster_name}/inventorySnapshots/{}", accepted.uid)
|
||||
|| accepted.content_hash != content_hash
|
||||
|| !is_exact_utc_seconds(&accepted.received_at)
|
||||
{
|
||||
return Err(InventoryError::Response);
|
||||
}
|
||||
Ok(InventoryDelivery::Accepted {
|
||||
content_hash,
|
||||
received_at: accepted.received_at,
|
||||
})
|
||||
}
|
||||
TelemetryDelivery::Retry { retry_after } => Ok(InventoryDelivery::Retry { retry_after }),
|
||||
TelemetryDelivery::AuthenticationStopped { status, reason } => {
|
||||
Ok(InventoryDelivery::AuthenticationStopped { status, reason })
|
||||
}
|
||||
TelemetryDelivery::Rejected { status, reason } => Ok(InventoryDelivery::Rejected { status, reason }),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
pub(crate) struct InventoryStateStore {
|
||||
path: PathBuf,
|
||||
}
|
||||
|
||||
#[derive(Default, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields, rename_all = "camelCase")]
|
||||
struct InventoryState {
|
||||
next_sequence: u64,
|
||||
pending: Option<PendingInventory>,
|
||||
last_accepted_content_hash: Option<String>,
|
||||
}
|
||||
|
||||
impl InventoryStateStore {
|
||||
pub(crate) fn from_heartbeat_path(path: &Path) -> Result<Self, InventoryError> {
|
||||
let root = path.parent().and_then(Path::parent).ok_or(InventoryError::StatePath)?;
|
||||
Ok(Self {
|
||||
path: root.join("inventory/state.json"),
|
||||
})
|
||||
}
|
||||
|
||||
pub(crate) fn try_runtime_lock(&self) -> Result<fs::File, InventoryError> {
|
||||
let directory = parent(&self.path)?;
|
||||
fs::create_dir_all(directory).map_err(|source| state_io(directory, source))?;
|
||||
let name = filename(&self.path)?;
|
||||
let path = directory.join(format!(".{name}.lock"));
|
||||
let mut options = fs::OpenOptions::new();
|
||||
options.create(true).truncate(false).read(true).write(true);
|
||||
#[cfg(unix)]
|
||||
{
|
||||
use std::os::unix::fs::OpenOptionsExt as _;
|
||||
options.mode(FILE_MODE);
|
||||
}
|
||||
let lock = options.open(&path).map_err(|source| state_io(&path, source))?;
|
||||
check_mode(&path)?;
|
||||
lock.try_lock().map_err(|_| InventoryError::AlreadyRunning)?;
|
||||
Ok(lock)
|
||||
}
|
||||
|
||||
pub(crate) async fn pending(&self) -> Result<Option<PendingInventory>, InventoryError> {
|
||||
let store = self.clone();
|
||||
tokio::task::spawn_blocking(move || {
|
||||
let state = store.read()?;
|
||||
if state.pending.is_none() && state.next_sequence > MAX_SEQUENCE {
|
||||
return Err(InventoryError::SequenceExhausted);
|
||||
}
|
||||
Ok(state.pending)
|
||||
})
|
||||
.await
|
||||
.map_err(|source| state_io(&self.path, io::Error::other(source)))?
|
||||
}
|
||||
|
||||
pub(crate) async fn prepare(&self, snapshot: InventorySnapshot) -> Result<Option<PendingInventory>, InventoryError> {
|
||||
let store = self.clone();
|
||||
tokio::task::spawn_blocking(move || store.prepare_sync(snapshot))
|
||||
.await
|
||||
.map_err(|source| state_io(&self.path, io::Error::other(source)))?
|
||||
}
|
||||
|
||||
pub(crate) async fn mark_accepted(&self, accepted: &PendingInventory) -> Result<(), InventoryError> {
|
||||
let store = self.clone();
|
||||
let accepted = accepted.clone();
|
||||
tokio::task::spawn_blocking(move || store.mark_accepted_sync(&accepted))
|
||||
.await
|
||||
.map_err(|source| state_io(&self.path, io::Error::other(source)))?
|
||||
}
|
||||
|
||||
fn prepare_sync(&self, snapshot: InventorySnapshot) -> Result<Option<PendingInventory>, InventoryError> {
|
||||
let mut state = self.read()?;
|
||||
if state.pending.is_some() {
|
||||
return Ok(state.pending);
|
||||
}
|
||||
let content_hash = snapshot.content_hash()?;
|
||||
if state.last_accepted_content_hash.as_deref() == Some(&content_hash) {
|
||||
return Ok(None);
|
||||
}
|
||||
if state.next_sequence > MAX_SEQUENCE {
|
||||
return Err(InventoryError::SequenceExhausted);
|
||||
}
|
||||
let pending = PendingInventory::new(snapshot, state.next_sequence);
|
||||
state.pending = Some(pending.clone());
|
||||
self.write(&state)?;
|
||||
Ok(Some(pending))
|
||||
}
|
||||
|
||||
fn mark_accepted_sync(&self, accepted: &PendingInventory) -> Result<(), InventoryError> {
|
||||
let mut state = self.read()?;
|
||||
if state.pending.as_ref() != Some(accepted) {
|
||||
return Err(InventoryError::StateConflict);
|
||||
}
|
||||
state.next_sequence = accepted.sequence.checked_add(1).ok_or(InventoryError::SequenceExhausted)?;
|
||||
state.last_accepted_content_hash = Some(accepted.content_hash()?);
|
||||
state.pending = None;
|
||||
self.write(&state)
|
||||
}
|
||||
|
||||
fn read(&self) -> Result<InventoryState, InventoryError> {
|
||||
let bytes = match fs::read(&self.path) {
|
||||
Ok(bytes) => bytes,
|
||||
Err(source) if source.kind() == io::ErrorKind::NotFound => return Ok(InventoryState::default()),
|
||||
Err(source) => return Err(state_io(&self.path, source)),
|
||||
};
|
||||
check_mode(&self.path)?;
|
||||
let state: InventoryState = serde_json::from_slice(&bytes).map_err(|source| InventoryError::StateInvalid {
|
||||
path: self.path.clone(),
|
||||
source,
|
||||
})?;
|
||||
let last_hash_valid = state.last_accepted_content_hash.as_deref().is_none_or(valid_content_hash);
|
||||
let pending_valid = state.pending.as_ref().is_none_or(|pending| {
|
||||
pending.sequence == state.next_sequence
|
||||
&& pending.is_valid()
|
||||
&& pending
|
||||
.content_hash()
|
||||
.is_ok_and(|hash| state.last_accepted_content_hash.as_deref() != Some(&hash))
|
||||
});
|
||||
if state.next_sequence > MAX_SEQUENCE + 1 || !last_hash_valid || !pending_valid {
|
||||
return Err(InventoryError::StateCorrupt { path: self.path.clone() });
|
||||
}
|
||||
Ok(state)
|
||||
}
|
||||
|
||||
fn write(&self, state: &InventoryState) -> Result<(), InventoryError> {
|
||||
let bytes = serde_json::to_vec(state).map_err(|source| InventoryError::StateInvalid {
|
||||
path: self.path.clone(),
|
||||
source,
|
||||
})?;
|
||||
let directory = parent(&self.path)?;
|
||||
fs::create_dir_all(directory).map_err(|source| state_io(directory, source))?;
|
||||
let temp = stage(directory, &self.path, &bytes)?;
|
||||
let result = fs::rename(&temp, &self.path)
|
||||
.map_err(|source| state_io(&self.path, source))
|
||||
.and_then(|()| fsync_dir(directory).map_err(|source| state_io(directory, source)));
|
||||
if result.is_err() {
|
||||
let _ = fs::remove_file(temp);
|
||||
}
|
||||
result
|
||||
}
|
||||
}
|
||||
|
||||
fn valid_content_hash(value: &str) -> bool {
|
||||
value.len() == 64
|
||||
&& value
|
||||
.bytes()
|
||||
.all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
|
||||
}
|
||||
|
||||
fn parent(path: &Path) -> Result<&Path, InventoryError> {
|
||||
path.parent()
|
||||
.ok_or_else(|| state_io(path, io::Error::new(io::ErrorKind::InvalidInput, "state path has no parent")))
|
||||
}
|
||||
|
||||
fn filename(path: &Path) -> Result<&str, InventoryError> {
|
||||
path.file_name()
|
||||
.and_then(|name| name.to_str())
|
||||
.ok_or_else(|| state_io(path, io::Error::new(io::ErrorKind::InvalidInput, "state filename is invalid")))
|
||||
}
|
||||
|
||||
fn stage(directory: &Path, destination: &Path, bytes: &[u8]) -> Result<PathBuf, InventoryError> {
|
||||
let name = filename(destination)?;
|
||||
loop {
|
||||
let path = directory.join(format!(
|
||||
".{name}.{}.{}.tmp",
|
||||
std::process::id(),
|
||||
STAGING_SEQUENCE.fetch_add(1, Ordering::Relaxed)
|
||||
));
|
||||
let mut options = fs::OpenOptions::new();
|
||||
options.write(true).create_new(true);
|
||||
#[cfg(unix)]
|
||||
{
|
||||
use std::os::unix::fs::OpenOptionsExt as _;
|
||||
options.mode(FILE_MODE);
|
||||
}
|
||||
let mut file = match options.open(&path) {
|
||||
Ok(file) => file,
|
||||
Err(source) if source.kind() == io::ErrorKind::AlreadyExists => continue,
|
||||
Err(source) => return Err(state_io(&path, source)),
|
||||
};
|
||||
if let Err(source) = file.write_all(bytes).and_then(|()| file.sync_all()) {
|
||||
let _ = fs::remove_file(&path);
|
||||
return Err(state_io(&path, source));
|
||||
}
|
||||
return Ok(path);
|
||||
}
|
||||
}
|
||||
|
||||
fn state_io(path: &Path, source: io::Error) -> InventoryError {
|
||||
InventoryError::StateIo {
|
||||
path: path.to_path_buf(),
|
||||
source,
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
fn check_mode(path: &Path) -> Result<(), InventoryError> {
|
||||
use std::os::unix::fs::PermissionsExt as _;
|
||||
|
||||
let mode = fs::metadata(path)
|
||||
.map_err(|source| state_io(path, source))?
|
||||
.permissions()
|
||||
.mode()
|
||||
& 0o7777;
|
||||
if mode != FILE_MODE {
|
||||
return Err(InventoryError::StatePermissions {
|
||||
path: path.to_path_buf(),
|
||||
mode,
|
||||
expected: FILE_MODE,
|
||||
});
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(not(unix))]
|
||||
fn check_mode(_path: &Path) -> Result<(), InventoryError> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn fsync_dir(directory: &Path) -> io::Result<()> {
|
||||
#[cfg(unix)]
|
||||
fs::File::open(directory)?.sync_all()?;
|
||||
#[cfg(not(unix))]
|
||||
let _ = directory;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[derive(Debug, thiserror::Error)]
|
||||
pub enum InventoryError {
|
||||
#[error("the RustFS inventory version is outside protocol bounds")]
|
||||
RustfsVersion,
|
||||
#[error("the RustFS inventory operating-system version is outside protocol bounds")]
|
||||
OsVersion,
|
||||
#[error("the RustFS inventory node count is outside protocol bounds")]
|
||||
NodeCount,
|
||||
#[error("the RustFS inventory drive count is outside protocol bounds")]
|
||||
DriveCount,
|
||||
#[error("the RustFS inventory capacity is outside protocol bounds")]
|
||||
Capacity,
|
||||
#[error("the RustFS inventory snapshot is incomplete: observed {observed} of {expected} configured drives")]
|
||||
SnapshotIncomplete { expected: usize, observed: usize },
|
||||
#[error("the Connect inventory schedule is invalid")]
|
||||
Schedule,
|
||||
#[error("the Connect inventory sequence is exhausted")]
|
||||
SequenceExhausted,
|
||||
#[error("a Connect inventory runtime already owns this state")]
|
||||
AlreadyRunning,
|
||||
#[error("the persisted Connect inventory changed while delivery was in flight")]
|
||||
StateConflict,
|
||||
#[error("the Connect inventory state path is invalid")]
|
||||
StatePath,
|
||||
#[error("Connect inventory state I/O failed at {path}: {source}")]
|
||||
StateIo {
|
||||
path: PathBuf,
|
||||
#[source]
|
||||
source: io::Error,
|
||||
},
|
||||
#[error("Connect inventory state at {path} is invalid: {source}")]
|
||||
StateInvalid {
|
||||
path: PathBuf,
|
||||
#[source]
|
||||
source: serde_json::Error,
|
||||
},
|
||||
#[error("Connect inventory state at {path} violates the protocol invariants")]
|
||||
StateCorrupt { path: PathBuf },
|
||||
#[cfg(unix)]
|
||||
#[error("Connect inventory state at {path} has mode {mode:o}, expected {expected:o}")]
|
||||
StatePermissions { path: PathBuf, mode: u32, expected: u32 },
|
||||
#[error("Connect returned an invalid inventory response")]
|
||||
Response,
|
||||
#[error(transparent)]
|
||||
Json(#[from] serde_json::Error),
|
||||
#[error("Connect inventory delivery failed: {0}")]
|
||||
Telemetry(String),
|
||||
}
|
||||
|
||||
impl From<TelemetryError> for InventoryError {
|
||||
fn from(error: TelemetryError) -> Self {
|
||||
Self::Telemetry(error.to_string())
|
||||
}
|
||||
}
|
||||
@@ -31,11 +31,9 @@ pub mod credential_store;
|
||||
pub mod heartbeat;
|
||||
pub mod identity;
|
||||
pub mod identity_store;
|
||||
pub mod inventory;
|
||||
pub mod offline;
|
||||
pub mod registration;
|
||||
pub mod runtime;
|
||||
mod telemetry;
|
||||
|
||||
pub use client::{ClientError, ConnectClient, ConnectConfig};
|
||||
pub use config::{HeartbeatConfig, HeartbeatConfigError, HeartbeatSchedule};
|
||||
@@ -43,10 +41,6 @@ pub use credential_store::{CredentialStore, DeviceCredential};
|
||||
pub use heartbeat::{CoarseNodeSummary, HeartbeatError, HeartbeatStatus};
|
||||
pub use identity::{DeviceIdentity, IdentityError, RegistrationProof, RegistrationTranscript};
|
||||
pub use identity_store::{IdentityStore, StoreError};
|
||||
pub use inventory::{
|
||||
InventoryError, InventoryFlag, InventoryOsVersion, InventorySchedule, InventorySnapshot, InventoryStatus,
|
||||
OperatingSystemFamily,
|
||||
};
|
||||
pub use offline::{EnrollmentError, OfflineEnrollment, OfflineKeyStore, VerifiedChallenge};
|
||||
pub use registration::{RegistrationToken, TokenError};
|
||||
pub use runtime::{HeartbeatRuntime, InventoryRuntime, spawn_heartbeat_runtime, spawn_inventory_runtime};
|
||||
pub use runtime::{HeartbeatRuntime, spawn_heartbeat_runtime};
|
||||
|
||||
@@ -23,16 +23,11 @@ use tokio_util::sync::CancellationToken;
|
||||
|
||||
use super::config::HeartbeatConfig;
|
||||
use super::heartbeat::{CoarseNodeSummary, Delivery, HeartbeatError, HeartbeatSender, HeartbeatStateStore, HeartbeatStatus};
|
||||
use super::inventory::{
|
||||
InventoryDelivery, InventoryError, InventorySchedule, InventorySender, InventorySnapshot, InventoryStateStore,
|
||||
InventoryStatus,
|
||||
};
|
||||
|
||||
pub struct HeartbeatRuntime {
|
||||
shutdown: CancellationToken,
|
||||
status: watch::Receiver<HeartbeatStatus>,
|
||||
task: Option<JoinHandle<()>>,
|
||||
inventory: Option<InventoryRuntime>,
|
||||
}
|
||||
|
||||
impl HeartbeatRuntime {
|
||||
@@ -40,19 +35,11 @@ impl HeartbeatRuntime {
|
||||
self.status.clone()
|
||||
}
|
||||
|
||||
pub(crate) fn with_inventory(mut self, inventory: Option<InventoryRuntime>) -> Self {
|
||||
self.inventory = inventory;
|
||||
self
|
||||
}
|
||||
|
||||
pub async fn shutdown(mut self) {
|
||||
self.shutdown.cancel();
|
||||
if let Some(task) = self.task.take() {
|
||||
let _ = task.await;
|
||||
}
|
||||
if let Some(inventory) = self.inventory.take() {
|
||||
inventory.shutdown().await;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -62,31 +49,6 @@ impl Drop for HeartbeatRuntime {
|
||||
}
|
||||
}
|
||||
|
||||
pub struct InventoryRuntime {
|
||||
shutdown: CancellationToken,
|
||||
status: watch::Receiver<InventoryStatus>,
|
||||
task: Option<JoinHandle<()>>,
|
||||
}
|
||||
|
||||
impl InventoryRuntime {
|
||||
pub fn status(&self) -> watch::Receiver<InventoryStatus> {
|
||||
self.status.clone()
|
||||
}
|
||||
|
||||
pub async fn shutdown(mut self) {
|
||||
self.shutdown.cancel();
|
||||
if let Some(task) = self.task.take() {
|
||||
let _ = task.await;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for InventoryRuntime {
|
||||
fn drop(&mut self) {
|
||||
self.shutdown.cancel();
|
||||
}
|
||||
}
|
||||
|
||||
pub fn spawn_heartbeat_runtime<F>(
|
||||
config: Option<HeartbeatConfig>,
|
||||
parent_shutdown: &CancellationToken,
|
||||
@@ -160,126 +122,6 @@ where
|
||||
shutdown,
|
||||
status: status_rx,
|
||||
task: Some(task),
|
||||
inventory: None,
|
||||
}))
|
||||
}
|
||||
|
||||
pub fn spawn_inventory_runtime<F, Fut>(
|
||||
config: Option<HeartbeatConfig>,
|
||||
schedule: InventorySchedule,
|
||||
parent_shutdown: &CancellationToken,
|
||||
sample: F,
|
||||
) -> Result<Option<InventoryRuntime>, InventoryError>
|
||||
where
|
||||
F: Fn() -> Fut + Send + Sync + 'static,
|
||||
Fut: Future<Output = Result<InventorySnapshot, InventoryError>> + Send + 'static,
|
||||
{
|
||||
let Some(config) = config else {
|
||||
return Ok(None);
|
||||
};
|
||||
if schedule.cadence.is_zero() || schedule.jitter > schedule.cadence {
|
||||
return Err(InventoryError::Schedule);
|
||||
}
|
||||
let retry_schedule = config.schedule;
|
||||
let store = InventoryStateStore::from_heartbeat_path(&config.state_path)?;
|
||||
let lock = store.try_runtime_lock()?;
|
||||
let sender = InventorySender::new(config)?;
|
||||
let shutdown = parent_shutdown.child_token();
|
||||
let task_shutdown = shutdown.clone();
|
||||
let (status_tx, status_rx) = watch::channel(InventoryStatus::Starting);
|
||||
let task = tokio::spawn(async move {
|
||||
let _lock = lock;
|
||||
let mut backoff = retry_schedule.initial_backoff;
|
||||
loop {
|
||||
if task_shutdown.is_cancelled() {
|
||||
break;
|
||||
}
|
||||
let pending = match store.pending().await {
|
||||
Ok(Some(pending)) => pending,
|
||||
Ok(None) => {
|
||||
let snapshot = match cancellable(&task_shutdown, sample()).await {
|
||||
Some(Ok(snapshot)) => snapshot,
|
||||
Some(Err(InventoryError::SnapshotIncomplete { .. })) => {
|
||||
let delay = backoff;
|
||||
backoff = backoff.saturating_mul(2).min(retry_schedule.max_backoff);
|
||||
let _ = status_tx.send(InventoryStatus::BackingOff { delay });
|
||||
if sleep_or_cancel(&task_shutdown, delay).await {
|
||||
break;
|
||||
}
|
||||
continue;
|
||||
}
|
||||
Some(Err(error)) => return failed_inventory(&status_tx, error),
|
||||
None => break,
|
||||
};
|
||||
let content_hash = match snapshot.content_hash() {
|
||||
Ok(content_hash) => content_hash,
|
||||
Err(error) => return failed_inventory(&status_tx, error),
|
||||
};
|
||||
match store.prepare(snapshot).await {
|
||||
Ok(Some(pending)) => pending,
|
||||
Ok(None) => {
|
||||
backoff = retry_schedule.initial_backoff;
|
||||
let _ = status_tx.send(InventoryStatus::Unchanged { content_hash });
|
||||
if sleep_or_cancel(&task_shutdown, schedule.cadence.saturating_add(jitter(schedule.jitter))).await {
|
||||
break;
|
||||
}
|
||||
continue;
|
||||
}
|
||||
Err(error) => return failed_inventory(&status_tx, error),
|
||||
}
|
||||
}
|
||||
Err(error) => return failed_inventory(&status_tx, error),
|
||||
};
|
||||
let delivery = match cancellable(&task_shutdown, sender.send(&pending)).await {
|
||||
Some(Ok(delivery)) => delivery,
|
||||
Some(Err(error)) => return failed_inventory(&status_tx, error),
|
||||
None => break,
|
||||
};
|
||||
let delay = match delivery {
|
||||
InventoryDelivery::Accepted {
|
||||
content_hash,
|
||||
received_at,
|
||||
} => {
|
||||
if let Err(error) = store.mark_accepted(&pending).await {
|
||||
return failed_inventory(&status_tx, error);
|
||||
}
|
||||
backoff = retry_schedule.initial_backoff;
|
||||
let _ = status_tx.send(InventoryStatus::Online {
|
||||
content_hash,
|
||||
received_at,
|
||||
});
|
||||
schedule.cadence.saturating_add(jitter(schedule.jitter))
|
||||
}
|
||||
InventoryDelivery::Retry { retry_after } => {
|
||||
let delay = retry_after
|
||||
.unwrap_or(backoff)
|
||||
.clamp(retry_schedule.initial_backoff, retry_schedule.max_backoff);
|
||||
backoff = backoff.saturating_mul(2).min(retry_schedule.max_backoff);
|
||||
let _ = status_tx.send(InventoryStatus::BackingOff { delay });
|
||||
delay
|
||||
}
|
||||
InventoryDelivery::AuthenticationStopped { status, reason } => {
|
||||
let _ = status_tx.send(InventoryStatus::AuthenticationStopped { status, reason });
|
||||
return;
|
||||
}
|
||||
InventoryDelivery::Rejected { status, reason } => {
|
||||
let suffix = reason.map_or_else(String::new, |reason| format!("; reason={reason}"));
|
||||
let _ = status_tx.send(InventoryStatus::Failed {
|
||||
reason: format!("Connect rejected inventory with HTTP {status}{suffix}"),
|
||||
});
|
||||
return;
|
||||
}
|
||||
};
|
||||
if sleep_or_cancel(&task_shutdown, delay).await {
|
||||
break;
|
||||
}
|
||||
}
|
||||
let _ = status_tx.send(InventoryStatus::Stopped);
|
||||
});
|
||||
Ok(Some(InventoryRuntime {
|
||||
shutdown,
|
||||
status: status_rx,
|
||||
task: Some(task),
|
||||
}))
|
||||
}
|
||||
|
||||
@@ -289,12 +131,6 @@ fn failed(status: &watch::Sender<HeartbeatStatus>, error: HeartbeatError) {
|
||||
});
|
||||
}
|
||||
|
||||
fn failed_inventory(status: &watch::Sender<InventoryStatus>, error: InventoryError) {
|
||||
let _ = status.send(InventoryStatus::Failed {
|
||||
reason: error.to_string(),
|
||||
});
|
||||
}
|
||||
|
||||
fn jitter(maximum: Duration) -> Duration {
|
||||
if maximum.is_zero() {
|
||||
Duration::ZERO
|
||||
|
||||
@@ -1,265 +0,0 @@
|
||||
// Copyright 2024 RustFS Team
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
use std::time::Duration;
|
||||
|
||||
use chrono::{DateTime, SecondsFormat, Utc};
|
||||
use reqwest::{Client, StatusCode, Url, header};
|
||||
use rustls::RootCertStore;
|
||||
use rustls::pki_types::{CertificateDer, pem::PemObject as _};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use zeroize::Zeroizing;
|
||||
|
||||
use super::config::HeartbeatConfig;
|
||||
use super::credential_store::{CredentialStoreError, DeviceCredential};
|
||||
use super::identity::IdentityError;
|
||||
use super::identity_store::StoreError;
|
||||
use super::registration::{CredentialValidationError, validate_stored_credential};
|
||||
|
||||
const MAX_RESPONSE_BYTES: usize = 64 * 1024;
|
||||
|
||||
pub(crate) enum TelemetryDelivery {
|
||||
Accepted { cluster_name: String, body: Vec<u8> },
|
||||
Retry { retry_after: Option<Duration> },
|
||||
AuthenticationStopped { status: u16, reason: Option<String> },
|
||||
Rejected { status: u16, reason: Option<String> },
|
||||
}
|
||||
|
||||
pub(crate) struct TelemetryTransport {
|
||||
endpoint: Url,
|
||||
root_store: RootCertStore,
|
||||
roots: Vec<CertificateDer<'static>>,
|
||||
config: HeartbeatConfig,
|
||||
}
|
||||
|
||||
impl TelemetryTransport {
|
||||
pub(crate) fn new(config: HeartbeatConfig) -> Result<Self, TelemetryError> {
|
||||
let mut endpoint = Url::parse(&config.endpoint).map_err(|_| TelemetryError::Endpoint)?;
|
||||
if endpoint.scheme() != "https"
|
||||
|| endpoint.cannot_be_a_base()
|
||||
|| !endpoint.username().is_empty()
|
||||
|| endpoint.password().is_some()
|
||||
|| endpoint.query().is_some()
|
||||
|| endpoint.fragment().is_some()
|
||||
{
|
||||
return Err(TelemetryError::Endpoint);
|
||||
}
|
||||
if !endpoint.path().ends_with('/') {
|
||||
endpoint.set_path(&format!("{}/", endpoint.path()));
|
||||
}
|
||||
let roots = CertificateDer::pem_slice_iter(&config.root_ca_pem)
|
||||
.collect::<Result<Vec<_>, _>>()
|
||||
.map_err(|_| TelemetryError::RootCertificate)?;
|
||||
if roots.is_empty() {
|
||||
return Err(TelemetryError::RootCertificate);
|
||||
}
|
||||
let mut root_store = RootCertStore::empty();
|
||||
let (accepted, rejected) = root_store.add_parsable_certificates(roots.clone());
|
||||
if accepted != roots.len() || rejected != 0 {
|
||||
return Err(TelemetryError::RootCertificate);
|
||||
}
|
||||
let schedule = config.schedule;
|
||||
if schedule.timeout.is_zero()
|
||||
|| schedule.timeout > Duration::from_secs(5)
|
||||
|| schedule.initial_backoff.is_zero()
|
||||
|| schedule.max_backoff < schedule.initial_backoff
|
||||
|| schedule.max_backoff > Duration::from_secs(5 * 60)
|
||||
{
|
||||
return Err(TelemetryError::Schedule);
|
||||
}
|
||||
Ok(Self {
|
||||
endpoint,
|
||||
root_store,
|
||||
roots,
|
||||
config,
|
||||
})
|
||||
}
|
||||
|
||||
pub(crate) async fn post<T: Serialize>(&self, collection: &str, value: &T) -> Result<TelemetryDelivery, TelemetryError> {
|
||||
let (cluster_name, cluster_uid, client) = self.authenticated_client().await?;
|
||||
let url = self.endpoint.join(&format!("clusters/{cluster_uid}/{collection}"))?;
|
||||
let response = match client.post(url).json(value).send().await {
|
||||
Ok(response) => response,
|
||||
Err(error) if error.is_timeout() || error.is_connect() || error.is_request() => {
|
||||
return Ok(TelemetryDelivery::Retry { retry_after: None });
|
||||
}
|
||||
Err(error) => return Err(error.into()),
|
||||
};
|
||||
let status = response.status();
|
||||
if status == StatusCode::TOO_MANY_REQUESTS {
|
||||
return Ok(TelemetryDelivery::Retry {
|
||||
retry_after: retry_after(response.headers(), Utc::now(), self.config.schedule.max_backoff),
|
||||
});
|
||||
}
|
||||
if status == StatusCode::REQUEST_TIMEOUT || status.is_server_error() {
|
||||
return Ok(TelemetryDelivery::Retry { retry_after: None });
|
||||
}
|
||||
if matches!(status, StatusCode::UNAUTHORIZED | StatusCode::FORBIDDEN) {
|
||||
return Ok(TelemetryDelivery::AuthenticationStopped {
|
||||
status: status.as_u16(),
|
||||
reason: response_reason(response).await,
|
||||
});
|
||||
}
|
||||
if status != StatusCode::OK {
|
||||
return Ok(TelemetryDelivery::Rejected {
|
||||
status: status.as_u16(),
|
||||
reason: response_reason(response).await,
|
||||
});
|
||||
}
|
||||
Ok(TelemetryDelivery::Accepted {
|
||||
cluster_name,
|
||||
body: bounded_body(response).await?,
|
||||
})
|
||||
}
|
||||
|
||||
async fn authenticated_client(&self) -> Result<(String, String, Client), TelemetryError> {
|
||||
let _lock = self.config.credential_store.lock().await?;
|
||||
let credential = self.config.credential_store.load()?.ok_or(TelemetryError::NotRegistered)?;
|
||||
let identity = self.config.identity_store.load()?.ok_or(TelemetryError::IdentityMissing)?;
|
||||
validate_stored_credential(&credential, &identity, &self.root_store, &self.roots)?;
|
||||
let now = Utc::now().timestamp();
|
||||
if now < credential.not_before_unix || now >= credential.not_after_unix {
|
||||
return Err(TelemetryError::CredentialExpired);
|
||||
}
|
||||
let (organization_uid, cluster_uid) = credential_parent(&credential)?;
|
||||
let cluster_name = format!("organizations/{organization_uid}/clusters/{cluster_uid}");
|
||||
let client = self.client(&credential, &identity.to_pkcs8_pem()?)?;
|
||||
Ok((cluster_name, cluster_uid.to_owned(), client))
|
||||
}
|
||||
|
||||
fn client(&self, credential: &DeviceCredential, key: &Zeroizing<String>) -> Result<Client, TelemetryError> {
|
||||
let mut pem = Zeroizing::new(Vec::with_capacity(credential.certificate_chain.len() + key.len() + 1));
|
||||
pem.extend_from_slice(credential.certificate_chain.as_bytes());
|
||||
pem.push(b'\n');
|
||||
pem.extend_from_slice(key.as_bytes());
|
||||
let identity = reqwest::Identity::from_pem(&pem).map_err(|_| TelemetryError::IdentityCertificate)?;
|
||||
let roots = self
|
||||
.roots
|
||||
.iter()
|
||||
.map(|root| reqwest::Certificate::from_der(root.as_ref()))
|
||||
.collect::<Result<Vec<_>, _>>()?;
|
||||
Client::builder()
|
||||
.https_only(true)
|
||||
.redirect(reqwest::redirect::Policy::none())
|
||||
.timeout(self.config.schedule.timeout)
|
||||
.tls_certs_only(roots)
|
||||
.identity(identity)
|
||||
.build()
|
||||
.map_err(Into::into)
|
||||
}
|
||||
}
|
||||
|
||||
fn credential_parent(credential: &DeviceCredential) -> Result<(&str, &str), TelemetryError> {
|
||||
let mut parts = credential.name.split('/');
|
||||
let valid = parts.next() == Some("organizations");
|
||||
let organization_uid = parts.next();
|
||||
let valid = valid && parts.next() == Some("clusters");
|
||||
let cluster_uid = parts.next();
|
||||
let valid = valid && parts.next() == Some("clusterDevices");
|
||||
let device_uid = parts.next();
|
||||
if !valid
|
||||
|| organization_uid.is_none_or(str::is_empty)
|
||||
|| cluster_uid.is_none_or(str::is_empty)
|
||||
|| device_uid != Some(credential.uid.as_str())
|
||||
|| parts.next().is_some()
|
||||
{
|
||||
return Err(TelemetryError::CredentialName);
|
||||
}
|
||||
Ok((
|
||||
organization_uid.ok_or(TelemetryError::CredentialName)?,
|
||||
cluster_uid.ok_or(TelemetryError::CredentialName)?,
|
||||
))
|
||||
}
|
||||
|
||||
fn retry_after(headers: &header::HeaderMap, now: DateTime<Utc>, maximum: Duration) -> Option<Duration> {
|
||||
let value = headers.get(header::RETRY_AFTER)?.to_str().ok()?;
|
||||
let delay = value.parse::<u64>().ok().map(Duration::from_secs).or_else(|| {
|
||||
DateTime::parse_from_rfc2822(value)
|
||||
.ok()
|
||||
.and_then(|at| (at.with_timezone(&Utc) - now).to_std().ok())
|
||||
})?;
|
||||
Some(delay.min(maximum))
|
||||
}
|
||||
|
||||
pub(crate) fn is_exact_utc_seconds(value: &str) -> bool {
|
||||
DateTime::parse_from_rfc3339(value).is_ok_and(|time| {
|
||||
time.offset().local_minus_utc() == 0
|
||||
&& value.ends_with('Z')
|
||||
&& time.with_timezone(&Utc).to_rfc3339_opts(SecondsFormat::Secs, true) == value
|
||||
})
|
||||
}
|
||||
|
||||
async fn response_reason(response: reqwest::Response) -> Option<String> {
|
||||
#[derive(Deserialize)]
|
||||
struct Envelope {
|
||||
#[serde(default)]
|
||||
details: Vec<Detail>,
|
||||
}
|
||||
#[derive(Deserialize)]
|
||||
struct Detail {
|
||||
#[serde(default)]
|
||||
reason: String,
|
||||
}
|
||||
|
||||
serde_json::from_slice::<Envelope>(&bounded_body(response).await.ok()?)
|
||||
.ok()?
|
||||
.details
|
||||
.into_iter()
|
||||
.find_map(|detail| (!detail.reason.is_empty()).then_some(detail.reason))
|
||||
}
|
||||
|
||||
async fn bounded_body(mut response: reqwest::Response) -> Result<Vec<u8>, TelemetryError> {
|
||||
let mut body = Vec::new();
|
||||
while let Some(chunk) = response.chunk().await? {
|
||||
if body.len().saturating_add(chunk.len()) > MAX_RESPONSE_BYTES {
|
||||
return Err(TelemetryError::ResponseTooLarge);
|
||||
}
|
||||
body.extend_from_slice(&chunk);
|
||||
}
|
||||
Ok(body)
|
||||
}
|
||||
|
||||
#[derive(Debug, thiserror::Error)]
|
||||
pub(crate) enum TelemetryError {
|
||||
#[error("Connect telemetry endpoint must be an HTTPS base URL without credentials, query, or fragment")]
|
||||
Endpoint,
|
||||
#[error("Connect telemetry root CA configuration is invalid")]
|
||||
RootCertificate,
|
||||
#[error("Connect telemetry retry schedule is invalid")]
|
||||
Schedule,
|
||||
#[error("RustFS is not registered with Connect")]
|
||||
NotRegistered,
|
||||
#[error("the Connect device private key is missing")]
|
||||
IdentityMissing,
|
||||
#[error("the stored Connect certificate and device private key cannot form a TLS identity")]
|
||||
IdentityCertificate,
|
||||
#[error("the stored Connect credential name is invalid")]
|
||||
CredentialName,
|
||||
#[error("the stored Connect device certificate is not currently valid")]
|
||||
CredentialExpired,
|
||||
#[error("Connect telemetry response exceeded 64 KiB")]
|
||||
ResponseTooLarge,
|
||||
#[error(transparent)]
|
||||
Url(#[from] url::ParseError),
|
||||
#[error(transparent)]
|
||||
Transport(#[from] reqwest::Error),
|
||||
#[error(transparent)]
|
||||
Identity(#[from] IdentityError),
|
||||
#[error(transparent)]
|
||||
IdentityStore(#[from] StoreError),
|
||||
#[error(transparent)]
|
||||
CredentialStore(#[from] CredentialStoreError),
|
||||
#[error(transparent)]
|
||||
CredentialValidation(#[from] CredentialValidationError),
|
||||
}
|
||||
+10
-8
@@ -26,22 +26,22 @@ struct MiMallocAllocator;
|
||||
unsafe impl GlobalAlloc for MiMallocAllocator {
|
||||
unsafe fn alloc(&self, layout: Layout) -> *mut u8 {
|
||||
// SAFETY: the caller upholds GlobalAlloc's contract for layout.
|
||||
unsafe { mimalloc::MiMalloc.alloc(layout) }
|
||||
unsafe { rustfs_mimalloc::MiMalloc.alloc(layout) }
|
||||
}
|
||||
|
||||
unsafe fn alloc_zeroed(&self, layout: Layout) -> *mut u8 {
|
||||
// SAFETY: the caller upholds GlobalAlloc's contract for layout.
|
||||
unsafe { mimalloc::MiMalloc.alloc_zeroed(layout) }
|
||||
unsafe { rustfs_mimalloc::MiMalloc.alloc_zeroed(layout) }
|
||||
}
|
||||
|
||||
unsafe fn dealloc(&self, ptr: *mut u8, layout: Layout) {
|
||||
// SAFETY: ptr and layout came from this allocator and are forwarded unchanged.
|
||||
unsafe { mimalloc::MiMalloc.dealloc(ptr, layout) }
|
||||
unsafe { rustfs_mimalloc::MiMalloc.dealloc(ptr, layout) }
|
||||
}
|
||||
|
||||
unsafe fn realloc(&self, ptr: *mut u8, layout: Layout, new_size: usize) -> *mut u8 {
|
||||
// SAFETY: ptr and layout came from this allocator and are forwarded unchanged.
|
||||
unsafe { mimalloc::MiMalloc.realloc(ptr, layout, new_size) }
|
||||
unsafe { rustfs_mimalloc::MiMalloc.realloc(ptr, layout, new_size) }
|
||||
}
|
||||
}
|
||||
|
||||
@@ -51,7 +51,7 @@ static GLOBAL: hotpath::CountingAllocator<MiMallocAllocator> = hotpath::Counting
|
||||
|
||||
#[cfg(not(all(feature = "hotpath", feature = "hotpath-alloc")))]
|
||||
#[global_allocator]
|
||||
static GLOBAL: mimalloc::MiMalloc = mimalloc::MiMalloc;
|
||||
static GLOBAL: rustfs_mimalloc::MiMalloc = rustfs_mimalloc::MiMalloc;
|
||||
|
||||
fn main() {
|
||||
let _hotpath_guard = hotpath::HotpathGuardBuilder::new("main").build();
|
||||
@@ -71,8 +71,9 @@ mod tests {
|
||||
allocation.extend_from_slice(&[7_u8; 64]);
|
||||
|
||||
assert_eq!(allocation.len(), 64);
|
||||
let heap = rustfs_mimalloc::heap::Heap::main();
|
||||
// SAFETY: the live Vec pointer is valid to inspect for heap ownership.
|
||||
assert!(unsafe { libmimalloc_sys::mi_is_in_heap_region(allocation.as_ptr().cast()) });
|
||||
assert!(unsafe { heap.contains(allocation.as_ptr()) });
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -85,12 +86,13 @@ mod tests {
|
||||
let layout = Layout::from_size_align(32, 8).expect("valid test allocation layout");
|
||||
let grown_layout = Layout::from_size_align(64, 8).expect("valid grown test allocation layout");
|
||||
let allocator = super::MiMallocAllocator;
|
||||
let heap = rustfs_mimalloc::heap::Heap::main();
|
||||
|
||||
// SAFETY: The pointer is checked for null before use and later released
|
||||
// through the same allocator with the corresponding layout.
|
||||
let ptr = unsafe { allocator.alloc_zeroed(layout) };
|
||||
assert!(!ptr.is_null());
|
||||
assert!(unsafe { libmimalloc_sys::mi_is_in_heap_region(ptr.cast()) });
|
||||
assert!(unsafe { heap.contains(ptr) });
|
||||
assert!(unsafe { std::slice::from_raw_parts(ptr, 32).iter().all(|byte| *byte == 0) });
|
||||
|
||||
// SAFETY: `ptr` was allocated by `allocator` with `layout`; on failure
|
||||
@@ -102,7 +104,7 @@ mod tests {
|
||||
panic!("mimalloc realloc failed in allocator smoke test");
|
||||
}
|
||||
|
||||
assert!(unsafe { libmimalloc_sys::mi_is_in_heap_region(grown_ptr.cast()) });
|
||||
assert!(unsafe { heap.contains(grown_ptr) });
|
||||
// SAFETY: `grown_ptr` was reallocated by `allocator` and is released
|
||||
// with the matching grown layout.
|
||||
unsafe { allocator.dealloc(grown_ptr, grown_layout) };
|
||||
|
||||
@@ -17,10 +17,7 @@ use rustfs_io_metrics::{
|
||||
record_cpu_usage, record_memory_usage, record_process_memory_split,
|
||||
};
|
||||
use serde::Serialize;
|
||||
#[cfg(any(test, not(target_os = "windows")))]
|
||||
use serde_json::Value;
|
||||
#[cfg(not(target_os = "windows"))]
|
||||
use std::ffi::CStr;
|
||||
use std::path::Path;
|
||||
use std::sync::{Arc, Mutex, OnceLock};
|
||||
use std::time::Duration;
|
||||
@@ -231,7 +228,18 @@ fn read_cgroup_memory_snapshot() -> Option<CgroupMemorySnapshot> {
|
||||
read_cgroup_v2().or_else(read_cgroup_v1)
|
||||
}
|
||||
|
||||
#[cfg(any(test, not(target_os = "windows")))]
|
||||
fn read_allocator_memory_snapshot() -> Option<AllocatorMemorySnapshot> {
|
||||
let json = rustfs_mimalloc::MiMalloc::stats_json();
|
||||
if json.is_empty() {
|
||||
return None;
|
||||
}
|
||||
let observation = parse_mimalloc_stats_json(&json)?;
|
||||
Some(AllocatorMemorySnapshot {
|
||||
backend: crate::allocator_reclaim::allocator_backend(),
|
||||
observation,
|
||||
})
|
||||
}
|
||||
|
||||
fn numeric_json_value(value: &Value) -> Option<u64> {
|
||||
match value {
|
||||
Value::Number(number) => number
|
||||
@@ -242,7 +250,6 @@ fn numeric_json_value(value: &Value) -> Option<u64> {
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(any(test, not(target_os = "windows")))]
|
||||
fn numeric_json_field(value: &Value, field: &str) -> Option<u64> {
|
||||
match value {
|
||||
Value::Object(fields) => fields
|
||||
@@ -254,7 +261,6 @@ fn numeric_json_field(value: &Value, field: &str) -> Option<u64> {
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(any(test, not(target_os = "windows")))]
|
||||
fn mimalloc_stat_field(value: &Value, metric: &str, field: &str) -> Option<u64> {
|
||||
match value {
|
||||
Value::Object(fields) => {
|
||||
@@ -271,12 +277,10 @@ fn mimalloc_stat_field(value: &Value, metric: &str, field: &str) -> Option<u64>
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(any(test, not(target_os = "windows")))]
|
||||
fn mimalloc_stat_current(value: &Value, metric: &str) -> Option<u64> {
|
||||
mimalloc_stat_field(value, metric, "current")
|
||||
}
|
||||
|
||||
#[cfg(any(test, not(target_os = "windows")))]
|
||||
fn mimalloc_stat_sum(value: &Value, metrics: &[&str], field: &str) -> Option<u64> {
|
||||
metrics
|
||||
.iter()
|
||||
@@ -285,7 +289,6 @@ fn mimalloc_stat_sum(value: &Value, metrics: &[&str], field: &str) -> Option<u64
|
||||
.filter(|value| *value > 0)
|
||||
}
|
||||
|
||||
#[cfg(any(test, not(target_os = "windows")))]
|
||||
fn parse_mimalloc_stats_json(stats_json: &str) -> Option<AllocatorMemoryObservation> {
|
||||
let value = serde_json::from_str::<Value>(stats_json).ok()?;
|
||||
let malloc_metrics = ["malloc_normal", "malloc_huge"];
|
||||
@@ -312,33 +315,6 @@ fn parse_mimalloc_stats_json(stats_json: &str) -> Option<AllocatorMemoryObservat
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(not(target_os = "windows"))]
|
||||
#[allow(unsafe_code)]
|
||||
fn read_allocator_memory_snapshot() -> Option<AllocatorMemorySnapshot> {
|
||||
// SAFETY: `mi_stats_get_json` returns a null-terminated JSON buffer owned by
|
||||
// mimalloc when called with a null input buffer. The mimalloc API requires
|
||||
// freeing that buffer with `mi_free`; parsing finishes before the buffer is freed.
|
||||
let observation = unsafe {
|
||||
let stats_ptr = libmimalloc_sys::mi_stats_get_json(0, std::ptr::null_mut());
|
||||
if stats_ptr.is_null() {
|
||||
return None;
|
||||
}
|
||||
|
||||
let observation = CStr::from_ptr(stats_ptr).to_str().ok().and_then(parse_mimalloc_stats_json);
|
||||
libmimalloc_sys::mi_free(stats_ptr.cast());
|
||||
observation?
|
||||
};
|
||||
Some(AllocatorMemorySnapshot {
|
||||
backend: crate::allocator_reclaim::allocator_backend(),
|
||||
observation,
|
||||
})
|
||||
}
|
||||
|
||||
#[cfg(target_os = "windows")]
|
||||
fn read_allocator_memory_snapshot() -> Option<AllocatorMemorySnapshot> {
|
||||
None
|
||||
}
|
||||
|
||||
fn configured_memory_observability_interval_secs() -> u64 {
|
||||
rustfs_utils::get_env_u64(ENV_MEMORY_OBSERVABILITY_INTERVAL_SECS, DEFAULT_MEMORY_OBSERVABILITY_INTERVAL_SECS).max(1)
|
||||
}
|
||||
@@ -566,6 +542,13 @@ mod tests {
|
||||
assert_eq!(parse_mimalloc_stats_json(r#"{ "allocator": "unknown" }"#), None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_allocator_memory_snapshot_uses_mimalloc_stats_json() {
|
||||
let snapshot = super::read_allocator_memory_snapshot();
|
||||
#[cfg(not(target_os = "windows"))]
|
||||
assert!(snapshot.is_some(), "allocator snapshot should be available on non-Windows");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn memory_observability_snapshot_reports_disabled_when_metrics_are_disabled() {
|
||||
let snapshot = build_memory_observability_status_snapshot(false, 15, false);
|
||||
|
||||
@@ -13,13 +13,10 @@
|
||||
// limitations under the License.
|
||||
|
||||
use crate::site_replication_reconcile::spawn_site_replication_reconcile_task;
|
||||
use crate::storage_api::startup::services::{ECStore, EndpointServerPools, ServerContextSlot, StorageAdminApi};
|
||||
use crate::storage_api::startup::services::{ECStore, EndpointServerPools, ServerContextSlot};
|
||||
use crate::{
|
||||
config::Config,
|
||||
connect::{
|
||||
CoarseNodeSummary, HeartbeatConfig, HeartbeatRuntime, InventoryError, InventoryFlag, InventoryRuntime, InventorySchedule,
|
||||
InventorySnapshot, spawn_heartbeat_runtime, spawn_inventory_runtime,
|
||||
},
|
||||
connect::{CoarseNodeSummary, HeartbeatConfig, HeartbeatRuntime, spawn_heartbeat_runtime},
|
||||
init::{init_buffer_profile_system, init_kms_system},
|
||||
server::ServiceStateManager,
|
||||
startup_audit::init_audit_runtime,
|
||||
@@ -80,9 +77,6 @@ pub(crate) async fn init_startup_runtime_services(
|
||||
let optional_runtimes = init_optional_runtime_services().await?;
|
||||
let heartbeat_config = HeartbeatConfig::from_env().map_err(std::io::Error::other)?;
|
||||
let heartbeat_nodes = heartbeat_config.as_ref().map(|_| endpoint_pools.get_nodes().len());
|
||||
let inventory_drives = heartbeat_config
|
||||
.as_ref()
|
||||
.map(|_| endpoint_pools.as_ref().iter().map(|pool| pool.endpoints.as_ref().len()).sum());
|
||||
|
||||
init_buffer_profile_system(config);
|
||||
init_deadlock_detector_runtime();
|
||||
@@ -102,9 +96,7 @@ pub(crate) async fn init_startup_runtime_services(
|
||||
init_notification_runtime(endpoint_pools, buckets).await?;
|
||||
let enable_scanner = init_background_service_runtime(store.clone()).await?;
|
||||
init_observability_runtime(store.clone(), ctx.clone()).await;
|
||||
let heartbeat = start_heartbeat_runtime(heartbeat_config.clone(), heartbeat_nodes, &ctx)?;
|
||||
let inventory = start_inventory_runtime(heartbeat_config, heartbeat_nodes, inventory_drives, store, &ctx)?;
|
||||
let heartbeat = heartbeat.map(|heartbeat| heartbeat.with_inventory(inventory));
|
||||
let heartbeat = start_heartbeat_runtime(heartbeat_config, heartbeat_nodes, &ctx)?;
|
||||
|
||||
Ok(StartupServiceRuntime {
|
||||
optional_runtimes,
|
||||
@@ -128,69 +120,3 @@ fn start_heartbeat_runtime(
|
||||
.ok_or_else(|| std::io::Error::other("Connect heartbeat node count is outside protocol bounds"))?;
|
||||
spawn_heartbeat_runtime(Some(config), shutdown, move || summary).map_err(std::io::Error::other)
|
||||
}
|
||||
|
||||
fn start_inventory_runtime(
|
||||
config: Option<HeartbeatConfig>,
|
||||
node_count: Option<usize>,
|
||||
expected_drive_count: Option<usize>,
|
||||
store: Arc<ECStore>,
|
||||
shutdown: &CancellationToken,
|
||||
) -> Result<Option<InventoryRuntime>> {
|
||||
let Some(config) = config else {
|
||||
return Ok(None);
|
||||
};
|
||||
let node_count = node_count.unwrap_or_default();
|
||||
let expected_drive_count = expected_drive_count.unwrap_or_default();
|
||||
spawn_inventory_runtime(Some(config), InventorySchedule::default(), shutdown, move || {
|
||||
let store = store.clone();
|
||||
async move {
|
||||
let info = StorageAdminApi::storage_info(store.as_ref()).await;
|
||||
inventory_snapshot(node_count, expected_drive_count, info)
|
||||
}
|
||||
})
|
||||
.map_err(std::io::Error::other)
|
||||
}
|
||||
|
||||
fn inventory_snapshot(
|
||||
node_count: usize,
|
||||
expected_drive_count: usize,
|
||||
info: rustfs_madmin::StorageInfo,
|
||||
) -> std::result::Result<InventorySnapshot, InventoryError> {
|
||||
if info.disks.len() != expected_drive_count {
|
||||
return Err(InventoryError::SnapshotIncomplete {
|
||||
expected: expected_drive_count,
|
||||
observed: info.disks.len(),
|
||||
});
|
||||
}
|
||||
let total = crate::app::storage_api::capacity::get_total_usable_capacity(&info.disks, &info) as u64;
|
||||
let free = crate::app::storage_api::capacity::get_total_usable_capacity_free(&info.disks, &info) as u64;
|
||||
let mut flags = Vec::with_capacity(3);
|
||||
if info.disks.iter().any(|disk| disk.state == rustfs_madmin::ITEM_OFFLINE) {
|
||||
flags.extend([InventoryFlag::ClusterDegraded, InventoryFlag::DriveOffline]);
|
||||
}
|
||||
if info.disks.iter().any(|disk| disk.healing) {
|
||||
flags.push(InventoryFlag::ClusterHealing);
|
||||
}
|
||||
InventorySnapshot::current(node_count, info.disks.len(), total, free, flags)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn inventory_rejects_a_partial_startup_storage_snapshot() {
|
||||
let info = rustfs_madmin::StorageInfo {
|
||||
disks: vec![rustfs_madmin::Disk::default()],
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
assert!(matches!(
|
||||
inventory_snapshot(2, 2, info),
|
||||
Err(InventoryError::SnapshotIncomplete {
|
||||
expected: 2,
|
||||
observed: 1
|
||||
})
|
||||
));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -279,7 +279,6 @@ pub(crate) mod startup {
|
||||
}
|
||||
|
||||
pub(crate) mod services {
|
||||
pub(crate) use super::super::storage_contracts::StorageAdminApi;
|
||||
pub(crate) use crate::storage::storage_api::{ECStore, EndpointServerPools, ServerContextSlot};
|
||||
}
|
||||
|
||||
|
||||
@@ -1,669 +0,0 @@
|
||||
// Copyright 2024 RustFS Team
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
use std::collections::VecDeque;
|
||||
use std::fs;
|
||||
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||
use std::sync::{Arc, Mutex};
|
||||
use std::time::Duration;
|
||||
|
||||
use bytes::Bytes;
|
||||
use http_body_util::{BodyExt as _, Full};
|
||||
use hyper::service::service_fn;
|
||||
use hyper::{Request, Response, StatusCode};
|
||||
use hyper_util::rt::TokioIo;
|
||||
use rcgen::{
|
||||
BasicConstraints, CertificateParams, DistinguishedName, DnType, ExtendedKeyUsagePurpose, IsCa, Issuer, KeyPair,
|
||||
KeyUsagePurpose, SanType,
|
||||
};
|
||||
use rustfs::connect::{
|
||||
CredentialStore, DeviceCredential, HeartbeatConfig, HeartbeatSchedule, IdentityStore, InventoryFlag, InventoryOsVersion,
|
||||
InventorySchedule, InventorySnapshot, InventoryStatus, OperatingSystemFamily, spawn_inventory_runtime,
|
||||
};
|
||||
use rustls::RootCertStore;
|
||||
use rustls::pki_types::{CertificateDer, PrivateKeyDer, PrivatePkcs8KeyDer};
|
||||
use rustls::server::WebPkiClientVerifier;
|
||||
use serde_json::{Value, json};
|
||||
use time::OffsetDateTime;
|
||||
use tokio::net::TcpListener;
|
||||
use tokio::sync::watch;
|
||||
use tokio_rustls::TlsAcceptor;
|
||||
use tokio_util::sync::CancellationToken;
|
||||
|
||||
const ORGANIZATION_UID: &str = "0198f4b0-1a00-7c10-8d21-2e3f4a5b6c70";
|
||||
const CLUSTER_UID: &str = "0198f4b0-2b00-7d20-9e31-3f4a5b6c7d81";
|
||||
const DEVICE_UID: &str = "0198f4b0-3c00-7e30-8f41-4a5b6c7d8e92";
|
||||
const SNAPSHOT_UID: &str = "0198f4b0-4d00-7f40-9051-5b6c7d8e9fa3";
|
||||
|
||||
struct TestPki {
|
||||
root_params: CertificateParams,
|
||||
root_key: KeyPair,
|
||||
root_der: CertificateDer<'static>,
|
||||
root_pem: String,
|
||||
server_der: CertificateDer<'static>,
|
||||
server_key: PrivatePkcs8KeyDer<'static>,
|
||||
}
|
||||
|
||||
impl TestPki {
|
||||
fn new() -> Self {
|
||||
let now = OffsetDateTime::now_utc();
|
||||
let root_key = KeyPair::generate().expect("generate root key");
|
||||
let mut root_params = CertificateParams::default();
|
||||
root_params.is_ca = IsCa::Ca(BasicConstraints::Unconstrained);
|
||||
root_params.not_before = now - time::Duration::days(30);
|
||||
root_params.not_after = now + time::Duration::days(30);
|
||||
root_params.key_usages = vec![KeyUsagePurpose::KeyCertSign, KeyUsagePurpose::DigitalSignature];
|
||||
let root = root_params.self_signed(&root_key).expect("sign root");
|
||||
|
||||
let server_key = KeyPair::generate().expect("generate server key");
|
||||
let mut server_params = CertificateParams::default();
|
||||
server_params.not_before = now - time::Duration::hours(1);
|
||||
server_params.not_after = now + time::Duration::days(2);
|
||||
server_params
|
||||
.subject_alt_names
|
||||
.push(SanType::DnsName("localhost".try_into().expect("valid DNS name")));
|
||||
server_params.extended_key_usages = vec![ExtendedKeyUsagePurpose::ServerAuth];
|
||||
let server = server_params
|
||||
.signed_by(&server_key, &Issuer::from_params(&root_params, &root_key))
|
||||
.expect("sign server certificate");
|
||||
Self {
|
||||
root_params,
|
||||
root_key,
|
||||
root_der: root.der().clone(),
|
||||
root_pem: root.pem(),
|
||||
server_der: server.der().clone(),
|
||||
server_key: PrivatePkcs8KeyDer::from(server_key.serialize_der()),
|
||||
}
|
||||
}
|
||||
|
||||
fn server_config(&self) -> rustls::ServerConfig {
|
||||
let mut roots = RootCertStore::empty();
|
||||
roots.add(self.root_der.clone()).expect("add client root");
|
||||
let verifier = WebPkiClientVerifier::builder(Arc::new(roots))
|
||||
.build()
|
||||
.expect("client verifier");
|
||||
rustls::ServerConfig::builder()
|
||||
.with_client_cert_verifier(verifier)
|
||||
.with_single_cert(vec![self.server_der.clone()], PrivateKeyDer::Pkcs8(self.server_key.clone_key()))
|
||||
.expect("server TLS")
|
||||
}
|
||||
|
||||
fn stores(&self, temp: &tempfile::TempDir) -> (IdentityStore, CredentialStore) {
|
||||
let identity_store = IdentityStore::new(temp.path().join("identity"));
|
||||
let identity = identity_store.load_or_create().expect("create identity");
|
||||
let private_key = PrivatePkcs8KeyDer::from(identity.to_pkcs8_der().expect("serialize key").to_vec());
|
||||
let device_key = KeyPair::from_pkcs8_der_and_sign_algo(&private_key, &rcgen::PKCS_ECDSA_P256_SHA256).expect("device key");
|
||||
let now = OffsetDateTime::now_utc();
|
||||
let mut params = CertificateParams::default();
|
||||
params.not_before = now - time::Duration::hours(1);
|
||||
params.not_after = now + time::Duration::hours(23);
|
||||
params.serial_number = Some(vec![1; 16].into());
|
||||
params.key_usages = vec![KeyUsagePurpose::DigitalSignature];
|
||||
params.extended_key_usages = vec![ExtendedKeyUsagePurpose::ClientAuth];
|
||||
params.distinguished_name = DistinguishedName::new();
|
||||
params.distinguished_name.push(DnType::CommonName, DEVICE_UID);
|
||||
params.subject_alt_names.push(SanType::URI(
|
||||
format!("urn:rustfs:connect:device:{DEVICE_UID}")
|
||||
.try_into()
|
||||
.expect("device URI"),
|
||||
));
|
||||
let certificate = params
|
||||
.signed_by(&device_key, &Issuer::from_params(&self.root_params, &self.root_key))
|
||||
.expect("device certificate");
|
||||
let cluster = format!("organizations/{ORGANIZATION_UID}/clusters/{CLUSTER_UID}");
|
||||
let credential = DeviceCredential {
|
||||
name: format!("{cluster}/clusterDevices/{DEVICE_UID}"),
|
||||
uid: DEVICE_UID.to_owned(),
|
||||
protocol_version: "v1".to_owned(),
|
||||
key_id: format!("x509-{}", "01".repeat(16)),
|
||||
certificate_serial: "01".repeat(16),
|
||||
certificate: certificate.pem(),
|
||||
certificate_chain: certificate.pem(),
|
||||
not_before_unix: (now - time::Duration::hours(1)).unix_timestamp(),
|
||||
not_after_unix: (now + time::Duration::hours(23)).unix_timestamp(),
|
||||
};
|
||||
let directory = temp.path().join("credential");
|
||||
fs::create_dir_all(&directory).expect("credential directory");
|
||||
let path = directory.join("device.crt.json");
|
||||
fs::write(&path, serde_json::to_vec(&credential).expect("credential JSON")).expect("write credential");
|
||||
private_mode(&path);
|
||||
(identity_store, CredentialStore::new(directory))
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
struct Reply {
|
||||
status: StatusCode,
|
||||
body: Value,
|
||||
retry_after: Option<&'static str>,
|
||||
}
|
||||
|
||||
impl Reply {
|
||||
fn ok(content_hash: &str) -> Self {
|
||||
Self {
|
||||
status: StatusCode::OK,
|
||||
body: json!({
|
||||
"name": format!("organizations/{ORGANIZATION_UID}/clusters/{CLUSTER_UID}/inventorySnapshots/{SNAPSHOT_UID}"),
|
||||
"uid": SNAPSHOT_UID,
|
||||
"contentHash": content_hash,
|
||||
"receivedAt": "2026-08-22T01:02:03Z",
|
||||
"futureField": true
|
||||
}),
|
||||
retry_after: None,
|
||||
}
|
||||
}
|
||||
|
||||
fn error(status: StatusCode, reason: &str) -> Self {
|
||||
Self {
|
||||
status,
|
||||
body: json!({"details": [{"reason": reason}]}),
|
||||
retry_after: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
struct TestServer {
|
||||
endpoint: String,
|
||||
seen: Arc<Mutex<Vec<Value>>>,
|
||||
task: tokio::task::JoinHandle<()>,
|
||||
}
|
||||
|
||||
impl Drop for TestServer {
|
||||
fn drop(&mut self) {
|
||||
self.task.abort();
|
||||
}
|
||||
}
|
||||
|
||||
async fn server(pki: &TestPki, replies: Vec<Reply>) -> TestServer {
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind server");
|
||||
let address = listener.local_addr().expect("server address");
|
||||
let acceptor = TlsAcceptor::from(Arc::new(pki.server_config()));
|
||||
let replies = Arc::new(Mutex::new(VecDeque::from(replies)));
|
||||
let seen = Arc::new(Mutex::new(Vec::new()));
|
||||
let captured = seen.clone();
|
||||
let task = tokio::spawn(async move {
|
||||
while let Ok((stream, _)) = listener.accept().await {
|
||||
let acceptor = acceptor.clone();
|
||||
let replies = replies.clone();
|
||||
let seen = captured.clone();
|
||||
tokio::spawn(async move {
|
||||
let Ok(stream) = acceptor.accept(stream).await else { return };
|
||||
let service = service_fn(move |request: Request<hyper::body::Incoming>| {
|
||||
let replies = replies.clone();
|
||||
let seen = seen.clone();
|
||||
async move {
|
||||
assert_eq!(request.uri().path(), format!("/agent/clusters/{CLUSTER_UID}/inventorySnapshots"));
|
||||
let body = request.into_body().collect().await.expect("request body").to_bytes();
|
||||
seen.lock()
|
||||
.expect("seen lock")
|
||||
.push(serde_json::from_slice(&body).expect("request JSON"));
|
||||
let reply = replies
|
||||
.lock()
|
||||
.expect("reply lock")
|
||||
.pop_front()
|
||||
.unwrap_or_else(|| Reply::error(StatusCode::SERVICE_UNAVAILABLE, "UNAVAILABLE"));
|
||||
let mut builder = Response::builder()
|
||||
.status(reply.status)
|
||||
.header("content-type", "application/json");
|
||||
if let Some(value) = reply.retry_after {
|
||||
builder = builder.header("retry-after", value);
|
||||
}
|
||||
Ok::<_, hyper::Error>(
|
||||
builder
|
||||
.body(Full::new(Bytes::from(serde_json::to_vec(&reply.body).expect("reply JSON"))))
|
||||
.expect("reply"),
|
||||
)
|
||||
}
|
||||
});
|
||||
let _ = hyper::server::conn::http1::Builder::new()
|
||||
.serve_connection(TokioIo::new(stream), service)
|
||||
.await;
|
||||
});
|
||||
}
|
||||
});
|
||||
TestServer {
|
||||
endpoint: format!("https://localhost:{}/agent/", address.port()),
|
||||
seen,
|
||||
task,
|
||||
}
|
||||
}
|
||||
|
||||
fn config(temp: &tempfile::TempDir, pki: &TestPki, server: &TestServer) -> HeartbeatConfig {
|
||||
let (identity_store, credential_store) = pki.stores(temp);
|
||||
HeartbeatConfig {
|
||||
endpoint: server.endpoint.clone(),
|
||||
root_ca_pem: pki.root_pem.as_bytes().to_vec(),
|
||||
identity_store,
|
||||
credential_store,
|
||||
state_path: temp.path().join("private-config-secret/heartbeat/state.json"),
|
||||
schedule: HeartbeatSchedule {
|
||||
cadence: Duration::from_secs(30),
|
||||
jitter: Duration::ZERO,
|
||||
timeout: Duration::from_millis(200),
|
||||
initial_backoff: Duration::from_millis(20),
|
||||
max_backoff: Duration::from_millis(80),
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
fn schedule() -> InventorySchedule {
|
||||
InventorySchedule {
|
||||
cadence: Duration::from_secs(60),
|
||||
jitter: Duration::ZERO,
|
||||
}
|
||||
}
|
||||
|
||||
fn snapshot() -> InventorySnapshot {
|
||||
InventorySnapshot::new(
|
||||
"1.4.2",
|
||||
Some(InventoryOsVersion::new(OperatingSystemFamily::Linux, 6, 8).expect("valid operating-system version")),
|
||||
8,
|
||||
96,
|
||||
1_099_511_627_776,
|
||||
412_316_860_416,
|
||||
[InventoryFlag::ClusterDegraded, InventoryFlag::DriveOffline],
|
||||
)
|
||||
.expect("valid inventory")
|
||||
}
|
||||
|
||||
fn collect_strings(value: &Value, strings: &mut Vec<String>) {
|
||||
match value {
|
||||
Value::String(value) => strings.push(value.clone()),
|
||||
Value::Array(values) => values.iter().for_each(|value| collect_strings(value, strings)),
|
||||
Value::Object(values) => values.values().for_each(|value| collect_strings(value, strings)),
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
|
||||
async fn wait_for(
|
||||
status: &mut watch::Receiver<InventoryStatus>,
|
||||
predicate: impl Fn(&InventoryStatus) -> bool,
|
||||
) -> InventoryStatus {
|
||||
tokio::time::timeout(Duration::from_secs(3), async {
|
||||
loop {
|
||||
let current = status.borrow_and_update().clone();
|
||||
if predicate(¤t) {
|
||||
return current;
|
||||
}
|
||||
status.changed().await.expect("status channel");
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("inventory status timeout")
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn connect_inventory_frozen_vector_has_the_exact_canonical_hash_and_no_open_ended_fields() {
|
||||
let fixtures: Value = serde_json::from_str(include_str!("../../protocol/agent/v1/fixtures/inventory/valid-vectors.json"))
|
||||
.expect("valid fixture JSON");
|
||||
let expected = &fixtures["vectors"][0]["expected"];
|
||||
let snapshot = snapshot();
|
||||
|
||||
assert_eq!(snapshot.content_hash().expect("content hash"), expected["contentHash"]);
|
||||
assert_eq!(InventorySchedule::default().cadence, Duration::from_secs(6 * 60 * 60));
|
||||
assert_eq!(InventorySchedule::default().jitter, Duration::from_secs(30 * 60));
|
||||
let encoded = serde_json::to_value(snapshot).expect("snapshot JSON");
|
||||
assert_eq!(
|
||||
encoded,
|
||||
json!({
|
||||
"rustfsVersion": "1.4.2",
|
||||
"osVersion": {"family": "linux", "major": 6, "minor": 8},
|
||||
"nodeCount": 8,
|
||||
"driveCount": 96,
|
||||
"capacityTotalBytes": 1099511627776_u64,
|
||||
"capacityUsedBytes": 412316860416_u64,
|
||||
"coarseFlags": ["cluster.degraded", "drive.offline"]
|
||||
})
|
||||
);
|
||||
|
||||
let fixtures: Value =
|
||||
serde_json::from_str(include_str!("../../protocol/agent/v1/fixtures/inventory/secret-like-vectors.json"))
|
||||
.expect("valid secret-like fixture JSON");
|
||||
let known_fields = [
|
||||
"protocolVersion",
|
||||
"rustfsVersion",
|
||||
"osVersion",
|
||||
"nodeCount",
|
||||
"driveCount",
|
||||
"capacityTotalBytes",
|
||||
"capacityUsedBytes",
|
||||
"coarseFlags",
|
||||
];
|
||||
let known_flags = ["cluster.degraded", "drive.offline"];
|
||||
let mut excluded = Vec::new();
|
||||
for vector in fixtures["vectors"].as_array().expect("fixture vectors") {
|
||||
let input = vector["input"].as_object().expect("fixture input");
|
||||
for (name, value) in input {
|
||||
if !known_fields.contains(&name.as_str()) {
|
||||
collect_strings(value, &mut excluded);
|
||||
}
|
||||
}
|
||||
for (name, value) in input["osVersion"].as_object().expect("fixture OS version") {
|
||||
if !["family", "major", "minor"].contains(&name.as_str()) {
|
||||
collect_strings(value, &mut excluded);
|
||||
}
|
||||
}
|
||||
for flag in input["coarseFlags"].as_array().expect("fixture coarse flags") {
|
||||
let flag = flag.as_str().expect("fixture coarse flag");
|
||||
if !known_flags.contains(&flag) {
|
||||
excluded.push(flag.to_owned());
|
||||
}
|
||||
}
|
||||
}
|
||||
let encoded = serde_json::to_string(&encoded).expect("encoded snapshot");
|
||||
for value in excluded {
|
||||
assert!(!encoded.contains(&value), "snapshot exposed fixture value {value}");
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn connect_inventory_bounds_fail_instead_of_truncating_or_inventing_values() {
|
||||
assert!(matches!(
|
||||
InventorySnapshot::current(0, 0, 0, 0, []),
|
||||
Err(rustfs::connect::InventoryError::NodeCount)
|
||||
));
|
||||
assert!(matches!(
|
||||
InventorySnapshot::current(1, 1_048_577, 0, 0, []),
|
||||
Err(rustfs::connect::InventoryError::DriveCount)
|
||||
));
|
||||
assert!(matches!(
|
||||
InventorySnapshot::current(1, 0, 9_007_199_254_740_992, 0, []),
|
||||
Err(rustfs::connect::InventoryError::Capacity)
|
||||
));
|
||||
assert!(matches!(
|
||||
InventorySnapshot::current(1, 0, 10, 11, []),
|
||||
Err(rustfs::connect::InventoryError::Capacity)
|
||||
));
|
||||
assert!(matches!(
|
||||
InventorySnapshot::new("1.0.0-private.1", None, 1, 0, 0, 0, []),
|
||||
Err(rustfs::connect::InventoryError::RustfsVersion)
|
||||
));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn connect_inventory_restart_replays_the_pending_request_and_then_skips_unchanged_inventory() {
|
||||
let pki = TestPki::new();
|
||||
let content_hash = snapshot().content_hash().expect("content hash");
|
||||
let first_server = server(&pki, vec![Reply::error(StatusCode::SERVICE_UNAVAILABLE, "UNAVAILABLE")]).await;
|
||||
let temp = tempfile::tempdir().expect("tempdir");
|
||||
let shutdown = CancellationToken::new();
|
||||
let samples = Arc::new(AtomicUsize::new(0));
|
||||
let sampled = samples.clone();
|
||||
let runtime = spawn_inventory_runtime(Some(config(&temp, &pki, &first_server)), schedule(), &shutdown, move || {
|
||||
sampled.fetch_add(1, Ordering::Relaxed);
|
||||
std::future::ready(Ok(snapshot()))
|
||||
})
|
||||
.expect("start inventory")
|
||||
.expect("configured inventory");
|
||||
let mut status = runtime.status();
|
||||
|
||||
assert!(matches!(
|
||||
wait_for(&mut status, |status| matches!(status, InventoryStatus::BackingOff { .. })).await,
|
||||
InventoryStatus::BackingOff { delay } if delay == Duration::from_millis(20)
|
||||
));
|
||||
assert_eq!(samples.load(Ordering::Relaxed), 1);
|
||||
let original = first_server.seen.lock().expect("seen lock")[0].clone();
|
||||
runtime.shutdown().await;
|
||||
|
||||
let mut limited = Reply::error(StatusCode::TOO_MANY_REQUESTS, "RATE_LIMITED");
|
||||
limited.retry_after = Some("0");
|
||||
let restart_server = server(&pki, vec![limited, Reply::ok(&content_hash)]).await;
|
||||
let restart_config = config(&temp, &pki, &restart_server);
|
||||
let restart_samples = Arc::new(AtomicUsize::new(0));
|
||||
let sampled = restart_samples.clone();
|
||||
let restart = spawn_inventory_runtime(Some(restart_config.clone()), schedule(), &shutdown, move || {
|
||||
sampled.fetch_add(1, Ordering::Relaxed);
|
||||
std::future::ready(Ok(snapshot()))
|
||||
})
|
||||
.expect("restart inventory")
|
||||
.expect("configured inventory");
|
||||
let mut status = restart.status();
|
||||
|
||||
assert!(matches!(
|
||||
wait_for(&mut status, |status| matches!(status, InventoryStatus::Online { .. })).await,
|
||||
InventoryStatus::Online { content_hash: accepted, received_at }
|
||||
if accepted == content_hash && received_at == "2026-08-22T01:02:03Z"
|
||||
));
|
||||
assert_eq!(restart_samples.load(Ordering::Relaxed), 0);
|
||||
let delivered = restart_server.seen.lock().expect("seen lock").clone();
|
||||
assert_eq!(delivered, vec![original.clone(), original.clone()]);
|
||||
assert_eq!(original["sequence"], 0);
|
||||
let encoded = serde_json::to_string(&original).expect("request JSON");
|
||||
for forbidden in [
|
||||
"private-config-secret",
|
||||
"BEGIN CERTIFICATE",
|
||||
"AKIAIOSFODNN7EXAMPLE",
|
||||
"bucket",
|
||||
"object",
|
||||
"path",
|
||||
] {
|
||||
assert!(!encoded.contains(forbidden), "request exposed {forbidden}");
|
||||
}
|
||||
assert_eq!(original.as_object().expect("request object").len(), 10);
|
||||
restart.shutdown().await;
|
||||
|
||||
let unchanged_samples = Arc::new(AtomicUsize::new(0));
|
||||
let sampled = unchanged_samples.clone();
|
||||
let unchanged = spawn_inventory_runtime(Some(restart_config), schedule(), &shutdown, move || {
|
||||
sampled.fetch_add(1, Ordering::Relaxed);
|
||||
std::future::ready(Ok(snapshot()))
|
||||
})
|
||||
.expect("restart inventory")
|
||||
.expect("configured inventory");
|
||||
let mut unchanged_status = unchanged.status();
|
||||
assert!(matches!(
|
||||
wait_for(&mut unchanged_status, |status| matches!(status, InventoryStatus::Unchanged { .. })).await,
|
||||
InventoryStatus::Unchanged { content_hash: unchanged } if unchanged == content_hash
|
||||
));
|
||||
assert_eq!(unchanged_samples.load(Ordering::Relaxed), 1);
|
||||
assert_eq!(restart_server.seen.lock().expect("seen lock").len(), 2);
|
||||
unchanged.shutdown().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn connect_inventory_disconnect_retries_without_resampling() {
|
||||
let pki = TestPki::new();
|
||||
let unavailable = server(&pki, Vec::new()).await;
|
||||
let temp = tempfile::tempdir().expect("tempdir");
|
||||
let config = config(&temp, &pki, &unavailable);
|
||||
drop(unavailable);
|
||||
let shutdown = CancellationToken::new();
|
||||
let samples = Arc::new(AtomicUsize::new(0));
|
||||
let sampled = samples.clone();
|
||||
let runtime = spawn_inventory_runtime(Some(config), schedule(), &shutdown, move || {
|
||||
sampled.fetch_add(1, Ordering::Relaxed);
|
||||
std::future::ready(Ok(snapshot()))
|
||||
})
|
||||
.expect("start inventory")
|
||||
.expect("configured inventory");
|
||||
let mut status = runtime.status();
|
||||
|
||||
assert!(matches!(
|
||||
wait_for(&mut status, |status| {
|
||||
matches!(status, InventoryStatus::BackingOff { delay } if *delay == Duration::from_millis(40))
|
||||
})
|
||||
.await,
|
||||
InventoryStatus::BackingOff { delay } if delay == Duration::from_millis(40)
|
||||
));
|
||||
assert_eq!(samples.load(Ordering::Relaxed), 1);
|
||||
runtime.shutdown().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn connect_inventory_retries_an_incomplete_sample_before_delivery() {
|
||||
let pki = TestPki::new();
|
||||
let content_hash = snapshot().content_hash().expect("content hash");
|
||||
let server = server(&pki, vec![Reply::ok(&content_hash)]).await;
|
||||
let temp = tempfile::tempdir().expect("tempdir");
|
||||
let shutdown = CancellationToken::new();
|
||||
let samples = Arc::new(AtomicUsize::new(0));
|
||||
let sampled = samples.clone();
|
||||
let runtime = spawn_inventory_runtime(Some(config(&temp, &pki, &server)), schedule(), &shutdown, move || {
|
||||
let attempt = sampled.fetch_add(1, Ordering::Relaxed);
|
||||
std::future::ready(if attempt == 0 {
|
||||
Err(rustfs::connect::InventoryError::SnapshotIncomplete {
|
||||
expected: 96,
|
||||
observed: 12,
|
||||
})
|
||||
} else {
|
||||
Ok(snapshot())
|
||||
})
|
||||
})
|
||||
.expect("start inventory")
|
||||
.expect("configured inventory");
|
||||
let mut status = runtime.status();
|
||||
|
||||
assert!(matches!(
|
||||
wait_for(&mut status, |status| matches!(status, InventoryStatus::Online { .. })).await,
|
||||
InventoryStatus::Online { content_hash: accepted, .. } if accepted == content_hash
|
||||
));
|
||||
assert_eq!(samples.load(Ordering::Relaxed), 2);
|
||||
assert_eq!(server.seen.lock().expect("seen lock").len(), 1);
|
||||
runtime.shutdown().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn connect_inventory_unchanged_sample_resets_incomplete_backoff() {
|
||||
let pki = TestPki::new();
|
||||
let content_hash = snapshot().content_hash().expect("content hash");
|
||||
let server = server(&pki, vec![Reply::ok(&content_hash)]).await;
|
||||
let temp = tempfile::tempdir().expect("tempdir");
|
||||
let shutdown = CancellationToken::new();
|
||||
let config = config(&temp, &pki, &server);
|
||||
let seed = spawn_inventory_runtime(Some(config.clone()), schedule(), &shutdown, || std::future::ready(Ok(snapshot())))
|
||||
.expect("start inventory")
|
||||
.expect("configured inventory");
|
||||
let mut seed_status = seed.status();
|
||||
|
||||
assert!(matches!(
|
||||
wait_for(&mut seed_status, |status| matches!(status, InventoryStatus::Online { .. })).await,
|
||||
InventoryStatus::Online { content_hash: accepted, .. } if accepted == content_hash
|
||||
));
|
||||
seed.shutdown().await;
|
||||
|
||||
let samples = Arc::new(AtomicUsize::new(0));
|
||||
let sampled = samples.clone();
|
||||
let runtime = spawn_inventory_runtime(
|
||||
Some(config),
|
||||
InventorySchedule {
|
||||
cadence: Duration::from_millis(100),
|
||||
jitter: Duration::ZERO,
|
||||
},
|
||||
&shutdown,
|
||||
move || {
|
||||
let attempt = sampled.fetch_add(1, Ordering::Relaxed);
|
||||
std::future::ready(if matches!(attempt, 0 | 1 | 3) {
|
||||
Err(rustfs::connect::InventoryError::SnapshotIncomplete {
|
||||
expected: 96,
|
||||
observed: 12,
|
||||
})
|
||||
} else {
|
||||
Ok(snapshot())
|
||||
})
|
||||
},
|
||||
)
|
||||
.expect("restart inventory")
|
||||
.expect("configured inventory");
|
||||
let mut status = runtime.status();
|
||||
|
||||
assert!(matches!(
|
||||
wait_for(&mut status, |status| {
|
||||
matches!(status, InventoryStatus::BackingOff { delay } if *delay == Duration::from_millis(20))
|
||||
})
|
||||
.await,
|
||||
InventoryStatus::BackingOff { delay } if delay == Duration::from_millis(20)
|
||||
));
|
||||
assert!(matches!(
|
||||
wait_for(&mut status, |status| {
|
||||
matches!(status, InventoryStatus::BackingOff { delay } if *delay == Duration::from_millis(40))
|
||||
})
|
||||
.await,
|
||||
InventoryStatus::BackingOff { delay } if delay == Duration::from_millis(40)
|
||||
));
|
||||
assert!(matches!(
|
||||
wait_for(&mut status, |status| matches!(status, InventoryStatus::Unchanged { .. })).await,
|
||||
InventoryStatus::Unchanged { content_hash: unchanged } if unchanged == content_hash
|
||||
));
|
||||
assert!(matches!(
|
||||
wait_for(&mut status, |status| matches!(status, InventoryStatus::BackingOff { .. })).await,
|
||||
InventoryStatus::BackingOff { delay } if delay == Duration::from_millis(20)
|
||||
));
|
||||
assert_eq!(samples.load(Ordering::Relaxed), 4);
|
||||
assert_eq!(server.seen.lock().expect("seen lock").len(), 1);
|
||||
runtime.shutdown().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn connect_inventory_revoked_device_stops_without_retrying() {
|
||||
let pki = TestPki::new();
|
||||
let server = server(&pki, vec![Reply::error(StatusCode::UNAUTHORIZED, "DEVICE_REVOKED")]).await;
|
||||
let temp = tempfile::tempdir().expect("tempdir");
|
||||
let shutdown = CancellationToken::new();
|
||||
let runtime = spawn_inventory_runtime(Some(config(&temp, &pki, &server)), schedule(), &shutdown, || {
|
||||
std::future::ready(Ok(snapshot()))
|
||||
})
|
||||
.expect("start inventory")
|
||||
.expect("configured inventory");
|
||||
let mut status = runtime.status();
|
||||
|
||||
assert!(matches!(
|
||||
wait_for(&mut status, |status| matches!(status, InventoryStatus::AuthenticationStopped { .. })).await,
|
||||
InventoryStatus::AuthenticationStopped { status: 401, reason: Some(reason) } if reason == "DEVICE_REVOKED"
|
||||
));
|
||||
assert_eq!(server.seen.lock().expect("seen lock").len(), 1);
|
||||
runtime.shutdown().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn connect_inventory_sequence_overflow_fails_before_sampling_or_network_delivery() {
|
||||
let pki = TestPki::new();
|
||||
let server = server(&pki, Vec::new()).await;
|
||||
let temp = tempfile::tempdir().expect("tempdir");
|
||||
let config = config(&temp, &pki, &server);
|
||||
let state = temp.path().join("private-config-secret/inventory/state.json");
|
||||
fs::create_dir_all(state.parent().expect("state directory")).expect("create state directory");
|
||||
fs::write(
|
||||
&state,
|
||||
br#"{"nextSequence":9007199254740992,"pending":null,"lastAcceptedContentHash":null}"#,
|
||||
)
|
||||
.expect("write state");
|
||||
private_mode(&state);
|
||||
let samples = Arc::new(AtomicUsize::new(0));
|
||||
let sampled = samples.clone();
|
||||
let shutdown = CancellationToken::new();
|
||||
let runtime = spawn_inventory_runtime(Some(config), schedule(), &shutdown, move || {
|
||||
sampled.fetch_add(1, Ordering::Relaxed);
|
||||
std::future::ready(Ok(snapshot()))
|
||||
})
|
||||
.expect("start inventory")
|
||||
.expect("configured inventory");
|
||||
let mut status = runtime.status();
|
||||
|
||||
assert!(matches!(
|
||||
wait_for(&mut status, |status| matches!(status, InventoryStatus::Failed { .. })).await,
|
||||
InventoryStatus::Failed { reason } if reason.contains("sequence is exhausted")
|
||||
));
|
||||
assert_eq!(samples.load(Ordering::Relaxed), 0);
|
||||
assert!(server.seen.lock().expect("seen lock").is_empty());
|
||||
runtime.shutdown().await;
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
fn private_mode(path: &std::path::Path) {
|
||||
use std::os::unix::fs::PermissionsExt as _;
|
||||
fs::set_permissions(path, fs::Permissions::from_mode(0o600)).expect("private permissions");
|
||||
}
|
||||
|
||||
#[cfg(not(unix))]
|
||||
fn private_mode(_path: &std::path::Path) {}
|
||||
Reference in New Issue
Block a user