Merge branch 'main' of github.com:rustfs/rustfs into houseme/get-small-file-optimization

This commit is contained in:
houseme
2026-06-28 11:54:44 +08:00
31 changed files with 474 additions and 474 deletions
+1 -1
View File
@@ -61,7 +61,7 @@ pub fn encode_tags(tags: Vec<Tag>) -> String {
for tag in tags.iter() {
if let (Some(k), Some(v)) = (tag.key.as_ref(), tag.value.as_ref()) {
//encoded.append_pair(k.as_ref().unwrap().as_str(), v.as_ref().unwrap().as_str());
//encoded.append_pair(k.as_ref().expect("operation should succeed").as_str(), v.as_ref().expect("operation should succeed").as_str());
encoded.append_pair(k.as_str(), v.as_str());
}
}
+3 -3
View File
@@ -237,10 +237,10 @@ impl PutObjectOptions {
if is_amz_header(k) || is_standard_header(k) || is_storageclass_header(k) || is_rustfs_header(k) || is_minio_header(k)
{
if let Ok(header_name) = HeaderName::from_bytes(k.as_bytes()) {
header.insert(header_name, HeaderValue::from_str(&v).unwrap());
header.insert(header_name, HeaderValue::from_str(&v).expect("operation should succeed"));
}
} else if let Ok(header_name) = HeaderName::from_bytes(format!("x-amz-meta-{}", k).as_bytes()) {
header.insert(header_name, HeaderValue::from_str(&v).unwrap());
header.insert(header_name, HeaderValue::from_str(&v).expect("operation should succeed"));
}
}
@@ -376,7 +376,7 @@ impl TransitionClient {
let mut md5_base64: String = "".to_string();
if opts.send_content_md5 {
if let Some(mut md5_hasher) = self.md5_hasher.lock().unwrap().as_mut() {
if let Some(mut md5_hasher) = self.md5_hasher.lock().expect("operation should succeed").as_mut() {
let hash = md5_hasher.hash_encode(&buf[..length]);
md5_base64 = base64_encode(hash.as_ref());
}
+5 -5
View File
@@ -150,10 +150,10 @@ impl TransitionClient {
) -> Result<ObjectInfo, std::io::Error> {
let mut headers = opts.header();
if opts.internal.replication_delete_marker {
headers.insert("X-Source-DeleteMarker", HeaderValue::from_str("true").unwrap());
headers.insert("X-Source-DeleteMarker", HeaderValue::from_str("true").expect("operation should succeed"));
}
if opts.internal.is_replication_ready_for_delete_marker {
headers.insert("X-Check-Replication-Ready", HeaderValue::from_str("true").unwrap());
headers.insert("X-Check-Replication-Ready", HeaderValue::from_str("true").expect("operation should succeed"));
}
let resp = self
@@ -183,12 +183,12 @@ impl TransitionClient {
Ok(resp) => {
let h = resp.headers();
let delete_marker = if let Some(x_amz_delete_marker) = h.get(X_AMZ_DELETE_MARKER.as_str()) {
x_amz_delete_marker.to_str().unwrap() == "true"
x_amz_delete_marker.to_str().expect("operation should succeed") == "true"
} else {
false
};
let replication_ready = if let Some(x_amz_delete_marker) = h.get("X-Replication-Ready") {
x_amz_delete_marker.to_str().unwrap() == "true"
x_amz_delete_marker.to_str().expect("operation should succeed") == "true"
} else {
false
};
@@ -224,7 +224,7 @@ impl TransitionClient {
//http_resp_to_error_response(resp, bucket_name, object_name)
}
Ok(to_object_info(bucket_name, object_name, h).unwrap())
Ok(to_object_info(bucket_name, object_name, h).expect("operation should succeed"))
}
Err(err) => {
return Err(std::io::Error::other(err));
@@ -121,8 +121,8 @@ pub fn new_getobjectreader<'a>(
let is_compressed = false; //oi.is_compressed_ok();
let rs_;
if rs.is_none() && opts.part_number.is_some() && opts.part_number.unwrap() > 0 {
rs_ = part_number_to_rangespec(oi.clone(), opts.part_number.unwrap());
if rs.is_none() && opts.part_number.is_some() && opts.part_number.expect("operation should succeed") > 0 {
rs_ = part_number_to_rangespec(oi.clone(), opts.part_number.expect("operation should succeed"));
} else {
rs_ = rs.clone();
}
+66 -66
View File
@@ -142,9 +142,9 @@ impl RemoteDisk {
pub(crate) async fn new(ep: &Endpoint, opt: &DiskOption, data_transport: Arc<dyn InternodeDataTransport>) -> Result<Self> {
let addr = if let Some(port) = ep.url.port() {
format!("{}://{}:{}", ep.url.scheme(), ep.url.host_str().unwrap(), port)
format!("{}://{}:{}", ep.url.scheme(), ep.url.host_str().expect("operation should succeed"), port)
} else {
format!("{}://{}", ep.url.scheme(), ep.url.host_str().unwrap())
format!("{}://{}", ep.url.scheme(), ep.url.host_str().expect("operation should succeed"))
};
let env_health_check =
@@ -363,7 +363,7 @@ impl RemoteDisk {
let elapsed = Duration::from_nanos(
(std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.expect("operation should succeed")
.as_nanos() as i64 - last_success_nanos) as u64
);
@@ -587,7 +587,7 @@ impl RemoteDisk {
// Record operation start
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.expect("operation should succeed")
.as_nanos() as i64;
self.health.last_started.store(now, std::sync::atomic::Ordering::Relaxed);
self.health.increment_waiting();
@@ -2550,7 +2550,7 @@ mod tests {
async fn new_remote_disk_with_transport(data_transport: Arc<dyn InternodeDataTransport>) -> RemoteDisk {
let endpoint = Endpoint {
url: url::Url::parse("http://remote-node:9000/data/rustfs0").unwrap(),
url: url::Url::parse("http://remote-node:9000/data/rustfs0").expect("operation should succeed"),
is_local: false,
pool_idx: 0,
set_idx: 0,
@@ -2561,7 +2561,7 @@ mod tests {
health_check: false,
};
RemoteDisk::new(&endpoint, &disk_option, data_transport).await.unwrap()
RemoteDisk::new(&endpoint, &disk_option, data_transport).await.expect("operation should succeed")
}
#[derive(Debug)]
@@ -2603,7 +2603,7 @@ mod tests {
#[tokio::test]
async fn test_remote_disk_creation() {
let url = url::Url::parse("http://example.com:9000/path").unwrap();
let url = url::Url::parse("http://example.com:9000/path").expect("operation should succeed");
let endpoint = Endpoint {
url: url.clone(),
is_local: false,
@@ -2619,7 +2619,7 @@ mod tests {
let remote_disk = RemoteDisk::new(&endpoint, &disk_option, Arc::new(TcpHttpInternodeDataTransport))
.await
.unwrap();
.expect("operation should succeed");
assert!(!remote_disk.is_local());
assert_eq!(remote_disk.endpoint.url, url);
@@ -2631,7 +2631,7 @@ mod tests {
#[tokio::test]
async fn test_remote_disk_basic_properties() {
let url = url::Url::parse("http://remote-server:9000").unwrap();
let url = url::Url::parse("http://remote-server:9000").expect("operation should succeed");
let endpoint = Endpoint {
url: url.clone(),
is_local: false,
@@ -2647,7 +2647,7 @@ mod tests {
let remote_disk = RemoteDisk::new(&endpoint, &disk_option, Arc::new(TcpHttpInternodeDataTransport))
.await
.unwrap();
.expect("operation should succeed");
// Test basic properties
assert!(!remote_disk.is_local());
@@ -2665,7 +2665,7 @@ mod tests {
#[tokio::test]
async fn test_remote_disk_path() {
let url = url::Url::parse("http://remote-server:9000/storage").unwrap();
let url = url::Url::parse("http://remote-server:9000/storage").expect("operation should succeed");
let endpoint = Endpoint {
url: url.clone(),
is_local: false,
@@ -2681,7 +2681,7 @@ mod tests {
let remote_disk = RemoteDisk::new(&endpoint, &disk_option, Arc::new(TcpHttpInternodeDataTransport))
.await
.unwrap();
.expect("operation should succeed");
let path = remote_disk.path();
// Remote disk path should be based on the URL path
@@ -2697,7 +2697,7 @@ mod tests {
};
let addr = listener.local_addr().expect("listener local address should be available");
let url = url::Url::parse(&format!("http://{}:{}/data/rustfs0", addr.ip(), addr.port())).unwrap();
let url = url::Url::parse(&format!("http://{}:{}/data/rustfs0", addr.ip(), addr.port())).expect("operation should succeed");
let endpoint = Endpoint {
url,
is_local: false,
@@ -2713,7 +2713,7 @@ mod tests {
let remote_disk = RemoteDisk::new(&endpoint, &disk_option, Arc::new(TcpHttpInternodeDataTransport))
.await
.unwrap();
.expect("operation should succeed");
assert!(remote_disk.is_online().await);
drop(listener);
@@ -2734,7 +2734,7 @@ mod tests {
drop(listener);
let url = url::Url::parse(&format!("http://{ip}:{port}/data/rustfs0")).unwrap();
let url = url::Url::parse(&format!("http://{ip}:{port}/data/rustfs0")).expect("operation should succeed");
let endpoint = Endpoint {
url,
is_local: false,
@@ -2756,7 +2756,7 @@ mod tests {
let remote_disk = RemoteDisk::new(&endpoint, &disk_option, Arc::new(TcpHttpInternodeDataTransport))
.await
.unwrap();
.expect("operation should succeed");
remote_disk.enable_health_check();
// Wait out the initial success-grace window so the active probe loop
@@ -2796,7 +2796,7 @@ mod tests {
});
let base_addr = format!("http://{}:{}", addr.ip(), addr.port());
let url = url::Url::parse(&format!("{base_addr}/data/rustfs0")).unwrap();
let url = url::Url::parse(&format!("{base_addr}/data/rustfs0")).expect("operation should succeed");
let endpoint = Endpoint {
url,
is_local: false,
@@ -2808,7 +2808,7 @@ mod tests {
health.mark_failure(&endpoint, "test_failure");
health.mark_failure(&endpoint, "test_failure");
assert_eq!(health.runtime_state(), RuntimeDriveHealthState::Offline);
let channel = TonicEndpoint::from_shared(base_addr.clone()).unwrap().connect_lazy();
let channel = TonicEndpoint::from_shared(base_addr.clone()).expect("operation should succeed").connect_lazy();
runtime_sources::cache_test_node_channel(base_addr.clone(), channel).await;
assert!(runtime_sources::test_node_channel_is_cached(&base_addr).await);
@@ -2854,19 +2854,19 @@ mod tests {
let copy_task = tokio::spawn(async move {
let mut cursor = Cursor::new(payload);
copy_stream_with_buffer(&mut cursor, &mut write_half, 4 * 1024).await.unwrap();
copy_stream_with_buffer(&mut cursor, &mut write_half, 4 * 1024).await.expect("operation should succeed");
});
let mut copied = Vec::new();
read_half.read_to_end(&mut copied).await.unwrap();
copy_task.await.unwrap();
read_half.read_to_end(&mut copied).await.expect("operation should succeed");
copy_task.await.expect("operation should succeed");
assert_eq!(copied, expected);
}
#[tokio::test]
async fn test_remote_disk_disk_id() {
let url = url::Url::parse("http://remote-server:9000").unwrap();
let url = url::Url::parse("http://remote-server:9000").expect("operation should succeed");
let endpoint = Endpoint {
url: url.clone(),
is_local: false,
@@ -2882,29 +2882,29 @@ mod tests {
let remote_disk = RemoteDisk::new(&endpoint, &disk_option, Arc::new(TcpHttpInternodeDataTransport))
.await
.unwrap();
.expect("operation should succeed");
// Initially, disk ID should be None
let initial_id = remote_disk.get_disk_id().await.unwrap();
let initial_id = remote_disk.get_disk_id().await.expect("operation should succeed");
assert!(initial_id.is_none());
// Set a disk ID
let test_id = Uuid::new_v4();
remote_disk.set_disk_id(Some(test_id)).await.unwrap();
remote_disk.set_disk_id(Some(test_id)).await.expect("operation should succeed");
// Verify the disk ID was set
let retrieved_id = remote_disk.get_disk_id().await.unwrap();
let retrieved_id = remote_disk.get_disk_id().await.expect("operation should succeed");
assert_eq!(retrieved_id, Some(test_id));
// Clear the disk ID
remote_disk.set_disk_id(None).await.unwrap();
let cleared_id = remote_disk.get_disk_id().await.unwrap();
remote_disk.set_disk_id(None).await.expect("operation should succeed");
let cleared_id = remote_disk.get_disk_id().await.expect("operation should succeed");
assert!(cleared_id.is_none());
}
#[tokio::test]
async fn test_remote_disk_ref_prefers_disk_id() {
let url = url::Url::parse("http://remote-server:9000").unwrap();
let url = url::Url::parse("http://remote-server:9000").expect("operation should succeed");
let endpoint = Endpoint {
url,
is_local: false,
@@ -2919,11 +2919,11 @@ mod tests {
let remote_disk = RemoteDisk::new(&endpoint, &disk_option, Arc::new(TcpHttpInternodeDataTransport))
.await
.unwrap();
.expect("operation should succeed");
assert_eq!(remote_disk.disk_ref().await, endpoint.to_string());
let disk_id = Uuid::new_v4();
remote_disk.set_disk_id(Some(disk_id)).await.unwrap();
remote_disk.set_disk_id(Some(disk_id)).await.expect("operation should succeed");
assert_eq!(remote_disk.disk_ref().await, disk_id.to_string());
}
@@ -2935,7 +2935,7 @@ mod tests {
let remote_disk = new_remote_disk_with_transport(Arc::new(transport.clone())).await;
let expected_disk = remote_disk.disk_ref().await;
let _reader = remote_disk.read_file_stream("bucket", "object/part.1", 7, 11).await.unwrap();
let _reader = remote_disk.read_file_stream("bucket", "object/part.1", 7, 11).await.expect("operation should succeed");
let calls = transport.calls();
assert_eq!(calls.len(), 1);
@@ -2961,7 +2961,7 @@ mod tests {
let transport = RecordingInternodeDataTransport::default();
let remote_disk = new_remote_disk_with_transport(Arc::new(transport.clone())).await;
let _reader = remote_disk.read_file_stream("bucket", "object/part.1", 7, 11).await.unwrap();
let _reader = remote_disk.read_file_stream("bucket", "object/part.1", 7, 11).await.expect("operation should succeed");
let calls = transport.calls();
assert_eq!(calls.len(), 1);
@@ -2982,8 +2982,8 @@ mod tests {
let _created = remote_disk
.create_file("orig-bucket", "bucket", "object/part.1", 4096)
.await
.unwrap();
let _appended = remote_disk.append_file("bucket", "object/part.2").await.unwrap();
.expect("operation should succeed");
let _appended = remote_disk.append_file("bucket", "object/part.2").await.expect("operation should succeed");
let calls = transport.calls();
assert_eq!(calls.len(), 2);
@@ -3070,10 +3070,10 @@ mod tests {
disk_id: String::new(),
..Default::default()
};
let expected_body = serde_json::to_vec(&opts).unwrap();
let expected_body = serde_json::to_vec(&opts).expect("operation should succeed");
let mut writer = Vec::new();
remote_disk.walk_dir(opts, &mut writer).await.unwrap();
remote_disk.walk_dir(opts, &mut writer).await.expect("operation should succeed");
let calls = transport.calls();
assert_eq!(calls.len(), 1);
@@ -3191,7 +3191,7 @@ mod tests {
];
for (url_str, expected_hostname) in test_cases {
let url = url::Url::parse(url_str).unwrap();
let url = url::Url::parse(url_str).expect("operation should succeed");
let endpoint = Endpoint {
url: url.clone(),
is_local: false,
@@ -3207,7 +3207,7 @@ mod tests {
let remote_disk = RemoteDisk::new(&endpoint, &disk_option, Arc::new(TcpHttpInternodeDataTransport))
.await
.unwrap();
.expect("operation should succeed");
assert!(!remote_disk.is_local());
assert_eq!(remote_disk.host_name(), expected_hostname);
@@ -3219,7 +3219,7 @@ mod tests {
#[tokio::test]
async fn test_remote_disk_location_validation() {
// Test valid location
let url = url::Url::parse("http://server:9000").unwrap();
let url = url::Url::parse("http://server:9000").expect("operation should succeed");
let valid_endpoint = Endpoint {
url: url.clone(),
is_local: false,
@@ -3235,7 +3235,7 @@ mod tests {
let remote_disk = RemoteDisk::new(&valid_endpoint, &disk_option, Arc::new(TcpHttpInternodeDataTransport))
.await
.unwrap();
.expect("operation should succeed");
let location = remote_disk.get_disk_location();
assert!(location.valid());
assert_eq!(location.pool_idx, Some(0));
@@ -3253,7 +3253,7 @@ mod tests {
let remote_disk_invalid = RemoteDisk::new(&invalid_endpoint, &disk_option, Arc::new(TcpHttpInternodeDataTransport))
.await
.unwrap();
.expect("operation should succeed");
let invalid_location = remote_disk_invalid.get_disk_location();
assert!(!invalid_location.valid());
assert_eq!(invalid_location.pool_idx, None);
@@ -3263,7 +3263,7 @@ mod tests {
#[tokio::test]
async fn test_remote_disk_close() {
let url = url::Url::parse("http://server:9000").unwrap();
let url = url::Url::parse("http://server:9000").expect("operation should succeed");
let endpoint = Endpoint {
url: url.clone(),
is_local: false,
@@ -3279,7 +3279,7 @@ mod tests {
let remote_disk = RemoteDisk::new(&endpoint, &disk_option, Arc::new(TcpHttpInternodeDataTransport))
.await
.unwrap();
.expect("operation should succeed");
// Test close operation (should succeed)
let result = remote_disk.close().await;
@@ -3288,7 +3288,7 @@ mod tests {
#[tokio::test]
async fn test_execute_with_timeout_marks_remote_disk_faulty() {
let url = url::Url::parse("http://remote-timeout:9000").unwrap();
let url = url::Url::parse("http://remote-timeout:9000").expect("operation should succeed");
let endpoint = Endpoint {
url,
is_local: false,
@@ -3306,7 +3306,7 @@ mod tests {
Arc::new(TcpHttpInternodeDataTransport),
)
.await
.unwrap();
.expect("operation should succeed");
let err = remote_disk
.execute_with_timeout(
@@ -3330,7 +3330,7 @@ mod tests {
#[tokio::test]
async fn test_execute_with_timeout_can_ignore_remote_timeout_failure() {
let url = url::Url::parse("http://remote-timeout-ignored:9000").unwrap();
let url = url::Url::parse("http://remote-timeout-ignored:9000").expect("operation should succeed");
let endpoint = Endpoint {
url,
is_local: false,
@@ -3348,7 +3348,7 @@ mod tests {
Arc::new(TcpHttpInternodeDataTransport),
)
.await
.unwrap();
.expect("operation should succeed");
let err = remote_disk
.execute_with_timeout_for_op_and_health_action(
@@ -3369,7 +3369,7 @@ mod tests {
#[tokio::test]
async fn test_execute_with_timeout_zero_duration_waits_for_operation() {
let url = url::Url::parse("http://remote-no-timeout:9000").unwrap();
let url = url::Url::parse("http://remote-no-timeout:9000").expect("operation should succeed");
let endpoint = Endpoint {
url,
is_local: false,
@@ -3387,7 +3387,7 @@ mod tests {
Arc::new(TcpHttpInternodeDataTransport),
)
.await
.unwrap();
.expect("operation should succeed");
remote_disk
.execute_with_timeout(
@@ -3409,7 +3409,7 @@ mod tests {
#[tokio::test]
async fn test_execute_with_timeout_evicts_cached_connection() {
let addr = "http://127.0.0.1:59991".to_string();
let url = url::Url::parse(&format!("{addr}/data")).unwrap();
let url = url::Url::parse(&format!("{addr}/data")).expect("operation should succeed");
let endpoint = Endpoint {
url,
is_local: false,
@@ -3427,9 +3427,9 @@ mod tests {
Arc::new(TcpHttpInternodeDataTransport),
)
.await
.unwrap();
.expect("operation should succeed");
let channel = TonicEndpoint::from_shared(addr.clone()).unwrap().connect_lazy();
let channel = TonicEndpoint::from_shared(addr.clone()).expect("operation should succeed").connect_lazy();
runtime_sources::cache_test_node_channel(addr.clone(), channel).await;
assert!(runtime_sources::test_node_channel_is_cached(&addr).await);
@@ -3453,7 +3453,7 @@ mod tests {
#[tokio::test]
async fn test_execute_with_timeout_marks_faulty_on_timeout_like_error() {
let addr = "http://127.0.0.1:59992".to_string();
let url = url::Url::parse(&format!("{addr}/data")).unwrap();
let url = url::Url::parse(&format!("{addr}/data")).expect("operation should succeed");
let endpoint = Endpoint {
url,
is_local: false,
@@ -3471,9 +3471,9 @@ mod tests {
Arc::new(TcpHttpInternodeDataTransport),
)
.await
.unwrap();
.expect("operation should succeed");
let channel = TonicEndpoint::from_shared(addr.clone()).unwrap().connect_lazy();
let channel = TonicEndpoint::from_shared(addr.clone()).expect("operation should succeed").connect_lazy();
runtime_sources::cache_test_node_channel(addr.clone(), channel).await;
let err = remote_disk
@@ -3509,7 +3509,7 @@ mod tests {
#[tokio::test]
async fn test_execute_with_timeout_marks_faulty_on_network_like_error() {
let addr = "http://127.0.0.1:59993".to_string();
let url = url::Url::parse(&format!("{addr}/data")).unwrap();
let url = url::Url::parse(&format!("{addr}/data")).expect("operation should succeed");
let endpoint = Endpoint {
url,
is_local: false,
@@ -3527,9 +3527,9 @@ mod tests {
Arc::new(TcpHttpInternodeDataTransport),
)
.await
.unwrap();
.expect("operation should succeed");
let channel = TonicEndpoint::from_shared(addr.clone()).unwrap().connect_lazy();
let channel = TonicEndpoint::from_shared(addr.clone()).expect("operation should succeed").connect_lazy();
runtime_sources::cache_test_node_channel(addr.clone(), channel).await;
let err = remote_disk
@@ -3570,7 +3570,7 @@ mod tests {
#[tokio::test]
async fn test_execute_with_timeout_can_ignore_network_like_error() {
let addr = "http://127.0.0.1:59995".to_string();
let url = url::Url::parse(&format!("{addr}/data")).unwrap();
let url = url::Url::parse(&format!("{addr}/data")).expect("operation should succeed");
let endpoint = Endpoint {
url,
is_local: false,
@@ -3588,9 +3588,9 @@ mod tests {
Arc::new(TcpHttpInternodeDataTransport),
)
.await
.unwrap();
.expect("operation should succeed");
let channel = TonicEndpoint::from_shared(addr.clone()).unwrap().connect_lazy();
let channel = TonicEndpoint::from_shared(addr.clone()).expect("operation should succeed").connect_lazy();
runtime_sources::cache_test_node_channel(addr.clone(), channel).await;
let err = remote_disk
@@ -3623,7 +3623,7 @@ mod tests {
#[tokio::test]
async fn test_execute_with_timeout_keeps_remote_disk_online_for_business_error() {
let addr = "http://127.0.0.1:59994".to_string();
let url = url::Url::parse(&format!("{addr}/data")).unwrap();
let url = url::Url::parse(&format!("{addr}/data")).expect("operation should succeed");
let endpoint = Endpoint {
url,
is_local: false,
@@ -3641,9 +3641,9 @@ mod tests {
Arc::new(TcpHttpInternodeDataTransport),
)
.await
.unwrap();
.expect("operation should succeed");
let channel = TonicEndpoint::from_shared(addr.clone()).unwrap().connect_lazy();
let channel = TonicEndpoint::from_shared(addr.clone()).expect("operation should succeed").connect_lazy();
runtime_sources::cache_test_node_channel(addr.clone(), channel).await;
let err = remote_disk
@@ -3661,7 +3661,7 @@ mod tests {
#[test]
fn test_remote_disk_sync_properties() {
let url = url::Url::parse("https://secure-remote:9000/data").unwrap();
let url = url::Url::parse("https://secure-remote:9000/data").expect("operation should succeed");
let endpoint = Endpoint {
url: url.clone(),
is_local: false,
+165 -165
View File
@@ -1827,7 +1827,7 @@ impl LocalDisk {
while let Some((last_name, _, _)) = dir_stack.last()
&& *last_name < name
{
let (pop, skip_object, dir_to_skip) = dir_stack.pop().unwrap();
let (pop, skip_object, dir_to_skip) = dir_stack.pop().expect("operation should succeed");
out.write_obj(&MetaCacheEntry {
name: pop.clone(),
..Default::default()
@@ -1862,7 +1862,7 @@ impl LocalDisk {
if let Some(_dir) = dir_objes.get(entry) {
is_dir_obj = true;
meta.name
.truncate(meta.name.len() - meta.name.chars().last().unwrap().len_utf8());
.truncate(meta.name.len() - meta.name.chars().last().expect("operation should succeed").len_utf8());
meta.name.push_str(GLOBAL_DIR_SUFFIX_WITH_SLASH);
}
@@ -3758,7 +3758,7 @@ impl DiskAPI for LocalDisk {
let mut info = Cache::get(self.disk_info_cache.clone()).await?;
info.nr_requests = self.nrrequests;
info.rotational = self.rotational;
info.mount_path = self.path().to_str().unwrap().to_string();
info.mount_path = self.path().to_str().expect("operation should succeed").to_string();
info.endpoint = self.endpoint.to_string();
info.scanning = self.scanning.load(Ordering::Acquire) == 1;
@@ -3851,7 +3851,7 @@ mod test {
let paths: Vec<_> = vols.iter().map(|v| path_join(&[Path::new(v), Path::new("test")])).collect();
for p in paths.iter() {
assert!(skip_access_checks(p.to_str().unwrap()));
assert!(skip_access_checks(p.to_str().expect("operation should succeed")));
}
}
@@ -3883,8 +3883,8 @@ mod test {
use crate::disk::format::FormatV3;
use tempfile::tempdir;
let dir = tempdir().unwrap();
let mut endpoint = Endpoint::try_from(dir.path().to_str().unwrap()).unwrap();
let dir = tempdir().expect("operation should succeed");
let mut endpoint = Endpoint::try_from(dir.path().to_str().expect("operation should succeed")).expect("operation should succeed");
endpoint.set_pool_index(0);
endpoint.set_set_index(0);
endpoint.set_disk_index(0);
@@ -3927,15 +3927,15 @@ mod test {
async fn cleanup_tmp_on_startup_moves_existing_tmp_and_recreates_trash() {
use tempfile::tempdir;
let dir = tempdir().unwrap();
let dir = tempdir().expect("operation should succeed");
let tmp = LocalDisk::meta_path(dir.path(), RUSTFS_META_TMP_BUCKET);
let leftover = tmp.join("leftover").join("data");
fs::create_dir_all(leftover.parent().unwrap()).await.unwrap();
fs::write(&leftover, b"temporary").await.unwrap();
fs::create_dir_all(leftover.parent().expect("operation should succeed")).await.expect("operation should succeed");
fs::write(&leftover, b"temporary").await.expect("operation should succeed");
LocalDisk::cleanup_tmp_on_startup(dir.path(), Arc::new(AtomicU32::new(0)), Arc::new(Notify::new()))
.await
.unwrap();
.expect("operation should succeed");
assert!(!tmp.join("leftover").exists());
assert!(LocalDisk::meta_path(dir.path(), RUSTFS_META_TMP_DELETED_BUCKET).exists());
@@ -3945,50 +3945,50 @@ mod test {
async fn cleanup_stale_tmp_objects_moves_expired_tmp_dirs_to_trash() {
use tempfile::tempdir;
let dir = tempdir().unwrap();
let dir = tempdir().expect("operation should succeed");
let tmp = LocalDisk::meta_path(dir.path(), RUSTFS_META_TMP_BUCKET);
let stale = tmp.join("stale").join("data");
let trash = LocalDisk::meta_path(dir.path(), RUSTFS_META_TMP_DELETED_BUCKET);
fs::create_dir_all(stale.parent().unwrap()).await.unwrap();
fs::create_dir_all(&trash).await.unwrap();
fs::write(&stale, b"temporary").await.unwrap();
fs::create_dir_all(stale.parent().expect("operation should succeed")).await.expect("operation should succeed");
fs::create_dir_all(&trash).await.expect("operation should succeed");
fs::write(&stale, b"temporary").await.expect("operation should succeed");
tokio::time::sleep(Duration::from_millis(2)).await;
LocalDisk::cleanup_stale_tmp_objects_with_expiry(dir.path().to_path_buf(), Duration::ZERO)
.await
.unwrap();
.expect("operation should succeed");
assert!(!tmp.join("stale").exists());
assert!(trash.exists());
let mut entries = fs::read_dir(&trash).await.unwrap();
assert!(entries.next_entry().await.unwrap().is_some());
let mut entries = fs::read_dir(&trash).await.expect("operation should succeed");
assert!(entries.next_entry().await.expect("operation should succeed").is_some());
}
#[tokio::test]
async fn cleanup_stale_tmp_objects_keeps_fresh_dirs_and_regular_files() {
use tempfile::tempdir;
let dir = tempdir().unwrap();
let dir = tempdir().expect("operation should succeed");
let tmp = LocalDisk::meta_path(dir.path(), RUSTFS_META_TMP_BUCKET);
let fresh_dir = tmp.join("fresh").join("data");
let regular_file = tmp.join("note.txt");
let trash = LocalDisk::meta_path(dir.path(), RUSTFS_META_TMP_DELETED_BUCKET);
fs::create_dir_all(fresh_dir.parent().unwrap()).await.unwrap();
fs::create_dir_all(&trash).await.unwrap();
fs::write(&fresh_dir, b"temporary").await.unwrap();
fs::write(&regular_file, b"keep").await.unwrap();
fs::create_dir_all(fresh_dir.parent().expect("operation should succeed")).await.expect("operation should succeed");
fs::create_dir_all(&trash).await.expect("operation should succeed");
fs::write(&fresh_dir, b"temporary").await.expect("operation should succeed");
fs::write(&regular_file, b"keep").await.expect("operation should succeed");
LocalDisk::cleanup_stale_tmp_objects_with_expiry(dir.path().to_path_buf(), Duration::from_secs(60))
.await
.unwrap();
.expect("operation should succeed");
assert!(tmp.join("fresh").exists());
assert!(regular_file.exists());
let mut entries = fs::read_dir(&trash).await.unwrap();
assert!(entries.next_entry().await.unwrap().is_none());
let mut entries = fs::read_dir(&trash).await.expect("operation should succeed");
assert!(entries.next_entry().await.expect("operation should succeed").is_none());
}
#[tokio::test(start_paused = true)]
@@ -4019,7 +4019,7 @@ mod test {
ready.store(1, Ordering::Release);
notify.notify_waiters();
assert!(wait.await.unwrap());
assert!(wait.await.expect("operation should succeed"));
}
#[tokio::test(start_paused = true)]
@@ -4036,7 +4036,7 @@ mod test {
tokio::task::yield_now().await;
tokio::time::advance(Duration::from_secs(2)).await;
assert!(!wait.await.unwrap());
assert!(!wait.await.expect("operation should succeed"));
}
#[tokio::test]
@@ -4044,21 +4044,21 @@ mod test {
use rustfs_filemeta::MetacacheReader;
use tempfile::tempdir;
let dir = tempdir().unwrap();
let dir = tempdir().expect("operation should succeed");
let bucket = "test-bucket";
let bucket_dir = dir.path().join(bucket);
fs::create_dir_all(bucket_dir.join("foo/bar/xyzzy")).await.unwrap();
fs::create_dir_all(bucket_dir.join("quux/thud")).await.unwrap();
fs::create_dir_all(bucket_dir.join("asdf")).await.unwrap();
fs::create_dir_all(bucket_dir.join("foo/bar/xyzzy")).await.expect("operation should succeed");
fs::create_dir_all(bucket_dir.join("quux/thud")).await.expect("operation should succeed");
fs::create_dir_all(bucket_dir.join("asdf")).await.expect("operation should succeed");
fs::write(bucket_dir.join("foo/bar/xl.meta"), b"meta").await.unwrap();
fs::write(bucket_dir.join("foo/bar/xyzzy/xl.meta"), b"meta").await.unwrap();
fs::write(bucket_dir.join("quux/thud/xl.meta"), b"meta").await.unwrap();
fs::write(bucket_dir.join("asdf/xl.meta"), b"meta").await.unwrap();
fs::write(bucket_dir.join("foo/bar/xl.meta"), b"meta").await.expect("operation should succeed");
fs::write(bucket_dir.join("foo/bar/xyzzy/xl.meta"), b"meta").await.expect("operation should succeed");
fs::write(bucket_dir.join("quux/thud/xl.meta"), b"meta").await.expect("operation should succeed");
fs::write(bucket_dir.join("asdf/xl.meta"), b"meta").await.expect("operation should succeed");
let endpoint = Endpoint::try_from(dir.path().to_str().unwrap()).unwrap();
let disk = LocalDisk::new(&endpoint, false).await.unwrap();
let endpoint = Endpoint::try_from(dir.path().to_str().expect("operation should succeed")).expect("operation should succeed");
let disk = LocalDisk::new(&endpoint, false).await.expect("operation should succeed");
let (reader, mut writer) = tokio::io::duplex(4096);
let mut out = MetacacheWriter::new(&mut writer);
@@ -4072,11 +4072,11 @@ mod test {
disk.scan_dir("".to_string(), "".to_string(), &opts, &mut out, &mut objs_returned, false, None)
.await
.unwrap();
out.close().await.unwrap();
.expect("operation should succeed");
out.close().await.expect("operation should succeed");
let mut reader = MetacacheReader::new(reader);
let entries = reader.read_all().await.unwrap();
let entries = reader.read_all().await.expect("operation should succeed");
let names: Vec<String> = entries
.into_iter()
.filter(|entry| !entry.metadata.is_empty())
@@ -4094,26 +4094,26 @@ mod test {
use rustfs_filemeta::MetacacheReader;
use tempfile::tempdir;
let dir = tempdir().unwrap();
let dir = tempdir().expect("operation should succeed");
let bucket = "test-bucket";
let bucket_dir = dir.path().join(bucket);
fs::create_dir_all(bucket_dir.join("marker/file.txt")).await.unwrap();
fs::create_dir_all(bucket_dir.join("marker/subdir/file.txt")).await.unwrap();
fs::create_dir_all(bucket_dir.join("marker/file.txt")).await.expect("operation should succeed");
fs::create_dir_all(bucket_dir.join("marker/subdir/file.txt")).await.expect("operation should succeed");
fs::create_dir_all(bucket_dir.join(format!("marker/subdir{GLOBAL_DIR_SUFFIX}")))
.await
.unwrap();
.expect("operation should succeed");
fs::write(bucket_dir.join("marker/file.txt/xl.meta"), b"meta").await.unwrap();
fs::write(bucket_dir.join("marker/file.txt/xl.meta"), b"meta").await.expect("operation should succeed");
fs::write(bucket_dir.join("marker/subdir/file.txt/xl.meta"), b"meta")
.await
.unwrap();
.expect("operation should succeed");
fs::write(bucket_dir.join(format!("marker/subdir{GLOBAL_DIR_SUFFIX}/xl.meta")), b"meta")
.await
.unwrap();
.expect("operation should succeed");
let endpoint = Endpoint::try_from(dir.path().to_str().unwrap()).unwrap();
let disk = LocalDisk::new(&endpoint, false).await.unwrap();
let endpoint = Endpoint::try_from(dir.path().to_str().expect("operation should succeed")).expect("operation should succeed");
let disk = LocalDisk::new(&endpoint, false).await.expect("operation should succeed");
let (reader, mut writer) = tokio::io::duplex(4096);
let mut out = MetacacheWriter::new(&mut writer);
@@ -4127,11 +4127,11 @@ mod test {
disk.scan_dir("marker/".to_string(), "".to_string(), &opts, &mut out, &mut objs_returned, false, None)
.await
.unwrap();
out.close().await.unwrap();
.expect("operation should succeed");
out.close().await.expect("operation should succeed");
let mut reader = MetacacheReader::new(reader);
let entries = reader.read_all().await.unwrap();
let entries = reader.read_all().await.expect("operation should succeed");
let names: Vec<String> = entries
.into_iter()
.filter(|entry| !entry.metadata.is_empty())
@@ -4148,7 +4148,7 @@ mod test {
use rustfs_filemeta::MetacacheReader;
use tempfile::tempdir;
let dir = tempdir().unwrap();
let dir = tempdir().expect("operation should succeed");
let bucket = "test-bucket";
let bucket_dir = dir.path().join(bucket);
@@ -4166,12 +4166,12 @@ mod test {
"unrelated/engineering/repo-0000",
] {
let object_dir = bucket_dir.join(name);
fs::create_dir_all(&object_dir).await.unwrap();
fs::write(object_dir.join(STORAGE_FORMAT_FILE), b"meta").await.unwrap();
fs::create_dir_all(&object_dir).await.expect("operation should succeed");
fs::write(object_dir.join(STORAGE_FORMAT_FILE), b"meta").await.expect("operation should succeed");
}
let endpoint = Endpoint::try_from(dir.path().to_str().unwrap()).unwrap();
let disk = LocalDisk::new(&endpoint, false).await.unwrap();
let endpoint = Endpoint::try_from(dir.path().to_str().expect("operation should succeed")).expect("operation should succeed");
let disk = LocalDisk::new(&endpoint, false).await.expect("operation should succeed");
async fn scan_names(disk: &LocalDisk, bucket: &str, base_dir: &str, forward_to: &str) -> (Vec<String>, i32) {
let (reader, mut writer) = tokio::io::duplex(4096);
@@ -4187,13 +4187,13 @@ mod test {
disk.scan_dir(base_dir.to_string(), "".to_string(), &opts, &mut out, &mut objs_returned, false, None)
.await
.unwrap();
out.close().await.unwrap();
.expect("operation should succeed");
out.close().await.expect("operation should succeed");
drop(out);
drop(writer);
let mut reader = MetacacheReader::new(reader);
let entries = reader.read_all().await.unwrap();
let entries = reader.read_all().await.expect("operation should succeed");
let names: Vec<String> = entries
.into_iter()
.filter(|entry| !entry.metadata.is_empty())
@@ -4291,7 +4291,7 @@ mod test {
fm.marshal_msg().expect("object metadata should encode")
}
let dir = tempdir().unwrap();
let dir = tempdir().expect("operation should succeed");
let bucket = "test-bucket";
let bucket_dir = dir.path().join(bucket);
@@ -4301,14 +4301,14 @@ mod test {
("shard/aaa-trash-0002", "33333333-3333-3333-3333-333333333333"),
] {
let object_dir = bucket_dir.join(name);
fs::create_dir_all(&object_dir).await.unwrap();
fs::create_dir_all(&object_dir).await.expect("operation should succeed");
fs::write(object_dir.join(STORAGE_FORMAT_FILE), delete_marker_metadata(version_id))
.await
.unwrap();
.expect("operation should succeed");
}
let hidden_versioned_dir = bucket_dir.join("shard/aaa-trash-0003");
fs::create_dir_all(&hidden_versioned_dir).await.unwrap();
fs::create_dir_all(&hidden_versioned_dir).await.expect("operation should succeed");
fs::write(
hidden_versioned_dir.join(STORAGE_FORMAT_FILE),
delete_marker_with_old_object_metadata(
@@ -4317,19 +4317,19 @@ mod test {
),
)
.await
.unwrap();
.expect("operation should succeed");
let visible_dir = bucket_dir.join("shard/bbb-visible-0000");
fs::create_dir_all(&visible_dir).await.unwrap();
fs::create_dir_all(&visible_dir).await.expect("operation should succeed");
fs::write(
visible_dir.join(STORAGE_FORMAT_FILE),
object_metadata("66666666-6666-6666-6666-666666666666"),
)
.await
.unwrap();
.expect("operation should succeed");
let endpoint = Endpoint::try_from(dir.path().to_str().unwrap()).unwrap();
let disk = LocalDisk::new(&endpoint, false).await.unwrap();
let endpoint = Endpoint::try_from(dir.path().to_str().expect("operation should succeed")).expect("operation should succeed");
let disk = LocalDisk::new(&endpoint, false).await.expect("operation should succeed");
let (reader, mut writer) = tokio::io::duplex(4096);
let mut out = MetacacheWriter::new(&mut writer);
@@ -4344,8 +4344,8 @@ mod test {
disk.scan_dir("".to_string(), "".to_string(), &opts, &mut out, &mut objs_returned, false, None)
.await
.unwrap();
out.close().await.unwrap();
.expect("operation should succeed");
out.close().await.expect("operation should succeed");
drop(out);
drop(writer);
@@ -4353,7 +4353,7 @@ mod test {
let has_visible_object = reader
.read_all()
.await
.unwrap()
.expect("operation should succeed")
.into_iter()
.any(|entry| !entry.metadata.is_empty() && entry.name == "shard/bbb-visible-0000");
@@ -4368,24 +4368,24 @@ mod test {
use std::os::unix::fs::PermissionsExt;
use tempfile::tempdir;
let dir = tempdir().unwrap();
let dir = tempdir().expect("operation should succeed");
let bucket = "test-bucket";
let bucket_dir = dir.path().join(bucket);
let object_dir = bucket_dir.join("broken");
let meta_path = object_dir.join(STORAGE_FORMAT_FILE);
fs::create_dir_all(&object_dir).await.unwrap();
fs::write(&meta_path, b"meta").await.unwrap();
fs::create_dir_all(&object_dir).await.expect("operation should succeed");
fs::write(&meta_path, b"meta").await.expect("operation should succeed");
let original_permissions = fs::metadata(&meta_path).await.unwrap().permissions();
fs::set_permissions(&meta_path, Permissions::from_mode(0o000)).await.unwrap();
let original_permissions = fs::metadata(&meta_path).await.expect("operation should succeed").permissions();
fs::set_permissions(&meta_path, Permissions::from_mode(0o000)).await.expect("operation should succeed");
if fs::File::open(&meta_path).await.is_ok() {
fs::set_permissions(&meta_path, original_permissions).await.unwrap();
fs::set_permissions(&meta_path, original_permissions).await.expect("operation should succeed");
return;
}
let endpoint = Endpoint::try_from(dir.path().to_str().unwrap()).unwrap();
let disk = LocalDisk::new(&endpoint, false).await.unwrap();
let endpoint = Endpoint::try_from(dir.path().to_str().expect("operation should succeed")).expect("operation should succeed");
let disk = LocalDisk::new(&endpoint, false).await.expect("operation should succeed");
let (_reader, mut writer) = tokio::io::duplex(4096);
let mut out = MetacacheWriter::new(&mut writer);
@@ -4401,7 +4401,7 @@ mod test {
.scan_dir("".to_string(), "".to_string(), &opts, &mut out, &mut objs_returned, false, None)
.await;
fs::set_permissions(&meta_path, original_permissions).await.unwrap();
fs::set_permissions(&meta_path, original_permissions).await.expect("operation should succeed");
assert!(matches!(result, Err(DiskError::FileAccessDenied)));
}
@@ -4422,7 +4422,7 @@ mod test {
const DIR_IN_MULTIPART_DIR: &str = "dir-in-multipart";
const EMPTY_STR: &str = "";
let parse_uuid = |s: &str| Uuid::parse_str(s).unwrap();
let parse_uuid = |s: &str| Uuid::parse_str(s).expect("operation should succeed");
let create_file_info = |version_id: &str, data_dir: &str| FileInfo {
version_id: Some(parse_uuid(version_id)),
data_dir: Some(parse_uuid(data_dir)),
@@ -4430,39 +4430,39 @@ mod test {
..Default::default()
};
let dir = tempdir().unwrap();
let dir = tempdir().expect("operation should succeed");
let obj_base = dir.path().join("test-bucket").join(BASE_DIR);
let multipart_base = obj_base.join(MULTIPART_DIR);
let dir_in_multipart_base = multipart_base.join(DIR_IN_MULTIPART_DIR);
fs::create_dir_all(&multipart_base).await.unwrap();
fs::create_dir_all(&multipart_base).await.expect("operation should succeed");
for uuid in &[UUID_MULTIPART_1, UUID_MULTIPART_2] {
fs::create_dir_all(multipart_base.join(uuid)).await.unwrap();
fs::write(multipart_base.join(uuid).join("part.1"), b"part").await.unwrap();
fs::create_dir_all(multipart_base.join(uuid)).await.expect("operation should succeed");
fs::write(multipart_base.join(uuid).join("part.1"), b"part").await.expect("operation should succeed");
}
fs::create_dir_all(obj_base.join(UUID_OBJ)).await.unwrap();
fs::write(obj_base.join(UUID_OBJ).join("part.1"), b"part").await.unwrap();
fs::create_dir_all(obj_base.join(UUID_OBJ)).await.expect("operation should succeed");
fs::write(obj_base.join(UUID_OBJ).join("part.1"), b"part").await.expect("operation should succeed");
fs::create_dir_all(&dir_in_multipart_base).await.unwrap();
fs::create_dir_all(&dir_in_multipart_base).await.expect("operation should succeed");
fs::write(dir_in_multipart_base.join(STORAGE_FORMAT_FILE), b"meta")
.await
.unwrap();
.expect("operation should succeed");
let mut fm = FileMeta::default();
fm.add_version(create_file_info(VER_ID_1, UUID_MULTIPART_1)).unwrap();
fm.add_version(create_file_info(VER_ID_2, UUID_MULTIPART_2)).unwrap();
fs::write(multipart_base.join(STORAGE_FORMAT_FILE), fm.marshal_msg().unwrap())
fm.add_version(create_file_info(VER_ID_1, UUID_MULTIPART_1)).expect("operation should succeed");
fm.add_version(create_file_info(VER_ID_2, UUID_MULTIPART_2)).expect("operation should succeed");
fs::write(multipart_base.join(STORAGE_FORMAT_FILE), fm.marshal_msg().expect("operation should succeed"))
.await
.unwrap();
.expect("operation should succeed");
let mut fm = FileMeta::default();
fm.add_version(create_file_info(VER_ID_3, UUID_OBJ)).unwrap();
fs::write(obj_base.join(STORAGE_FORMAT_FILE), fm.marshal_msg().unwrap())
fm.add_version(create_file_info(VER_ID_3, UUID_OBJ)).expect("operation should succeed");
fs::write(obj_base.join(STORAGE_FORMAT_FILE), fm.marshal_msg().expect("operation should succeed"))
.await
.unwrap();
.expect("operation should succeed");
let endpoint = Endpoint::try_from(dir.path().to_str().unwrap()).unwrap();
let disk = LocalDisk::new(&endpoint, false).await.unwrap();
let endpoint = Endpoint::try_from(dir.path().to_str().expect("operation should succeed")).expect("operation should succeed");
let disk = LocalDisk::new(&endpoint, false).await.expect("operation should succeed");
let (reader, mut writer) = tokio::io::duplex(4096);
disk.walk_dir(
@@ -4476,11 +4476,11 @@ mod test {
&mut writer,
)
.await
.unwrap();
MetacacheWriter::new(&mut writer).close().await.unwrap();
.expect("operation should succeed");
MetacacheWriter::new(&mut writer).close().await.expect("operation should succeed");
let mut reader = MetacacheReader::new(reader);
let entries = reader.read_all().await.unwrap();
let entries = reader.read_all().await.expect("operation should succeed");
let names: Vec<String> = entries.into_iter().map(|entry| entry.name).collect();
assert_eq!(
@@ -4558,7 +4558,7 @@ mod test {
#[tokio::test]
async fn test_make_volume() {
let p = "./testv0";
fs::create_dir_all(&p).await.unwrap();
fs::create_dir_all(&p).await.expect("operation should succeed");
let ep = match Endpoint::try_from(p) {
Ok(e) => e,
@@ -4568,17 +4568,17 @@ mod test {
}
};
let disk = LocalDisk::new(&ep, false).await.unwrap();
let disk = LocalDisk::new(&ep, false).await.expect("operation should succeed");
let tmpp = disk.resolve_abs_path(Path::new(RUSTFS_META_TMP_DELETED_BUCKET)).unwrap();
let tmpp = disk.resolve_abs_path(Path::new(RUSTFS_META_TMP_DELETED_BUCKET)).expect("operation should succeed");
println!("ppp :{:?}", &tmpp);
let volumes = vec!["a123", "b123", "c123"];
disk.make_volumes(volumes.clone()).await.unwrap();
disk.make_volumes(volumes.clone()).await.expect("operation should succeed");
disk.make_volumes(volumes.clone()).await.unwrap();
disk.make_volumes(volumes.clone()).await.expect("operation should succeed");
let _ = fs::remove_dir_all(&p).await;
}
@@ -4586,7 +4586,7 @@ mod test {
#[tokio::test]
async fn test_delete_volume() {
let p = "./testv1";
fs::create_dir_all(&p).await.unwrap();
fs::create_dir_all(&p).await.expect("operation should succeed");
let ep = match Endpoint::try_from(p) {
Ok(e) => e,
@@ -4596,17 +4596,17 @@ mod test {
}
};
let disk = LocalDisk::new(&ep, false).await.unwrap();
let disk = LocalDisk::new(&ep, false).await.expect("operation should succeed");
let tmpp = disk.resolve_abs_path(Path::new(RUSTFS_META_TMP_DELETED_BUCKET)).unwrap();
let tmpp = disk.resolve_abs_path(Path::new(RUSTFS_META_TMP_DELETED_BUCKET)).expect("operation should succeed");
println!("ppp :{:?}", &tmpp);
let volumes = vec!["a123", "b123", "c123"];
disk.make_volumes(volumes.clone()).await.unwrap();
disk.make_volumes(volumes.clone()).await.expect("operation should succeed");
disk.delete_volume("a").await.unwrap();
disk.delete_volume("a").await.expect("operation should succeed");
let _ = fs::remove_dir_all(&p).await;
}
@@ -4614,10 +4614,10 @@ mod test {
#[tokio::test]
async fn test_local_disk_basic_operations() {
let test_dir = "./test_local_disk_basic";
fs::create_dir_all(&test_dir).await.unwrap();
fs::create_dir_all(&test_dir).await.expect("operation should succeed");
let endpoint = Endpoint::try_from(test_dir).unwrap();
let disk = LocalDisk::new(&endpoint, false).await.unwrap();
let endpoint = Endpoint::try_from(test_dir).expect("operation should succeed");
let disk = LocalDisk::new(&endpoint, false).await.expect("operation should succeed");
// Test basic properties
assert!(disk.is_local());
@@ -4626,15 +4626,15 @@ mod test {
assert!(!disk.to_string().is_empty());
// Test path resolution
let abs_path = disk.resolve_abs_path("test/path").unwrap();
let abs_path = disk.resolve_abs_path("test/path").expect("operation should succeed");
assert!(abs_path.is_absolute());
// Test bucket path
let bucket_path = disk.get_bucket_path("test-bucket").unwrap();
let bucket_path = disk.get_bucket_path("test-bucket").expect("operation should succeed");
assert!(bucket_path.to_string_lossy().contains("test-bucket"));
// Test object path
let object_path = disk.get_object_path("test-bucket", "test-object").unwrap();
let object_path = disk.get_object_path("test-bucket", "test-object").expect("operation should succeed");
assert!(object_path.to_string_lossy().contains("test-bucket"));
assert!(object_path.to_string_lossy().contains("test-object"));
@@ -4648,13 +4648,13 @@ mod test {
use std::os::unix::fs::symlink;
use tempfile::tempdir;
let root_dir = tempdir().unwrap();
let outside_dir = tempdir().unwrap();
let root_dir = tempdir().expect("operation should succeed");
let outside_dir = tempdir().expect("operation should succeed");
let link_path = root_dir.path().join("escape-bucket");
symlink(outside_dir.path(), &link_path).unwrap();
symlink(outside_dir.path(), &link_path).expect("operation should succeed");
let endpoint = Endpoint::try_from(root_dir.path().to_string_lossy().as_ref()).unwrap();
let disk = LocalDisk::new(&endpoint, false).await.unwrap();
let endpoint = Endpoint::try_from(root_dir.path().to_string_lossy().as_ref()).expect("operation should succeed");
let disk = LocalDisk::new(&endpoint, false).await.expect("operation should succeed");
assert!(matches!(disk.get_bucket_path("escape-bucket"), Err(DiskError::InvalidPath)));
}
@@ -4665,15 +4665,15 @@ mod test {
use std::os::unix::fs::symlink;
use tempfile::tempdir;
let root_dir = tempdir().unwrap();
let outside_dir = tempdir().unwrap();
let root_dir = tempdir().expect("operation should succeed");
let outside_dir = tempdir().expect("operation should succeed");
let bucket_dir = root_dir.path().join("bucket");
fs::create_dir_all(&bucket_dir).await.unwrap();
fs::create_dir_all(&bucket_dir).await.expect("operation should succeed");
let link_path = bucket_dir.join("escape");
symlink(outside_dir.path(), &link_path).unwrap();
symlink(outside_dir.path(), &link_path).expect("operation should succeed");
let endpoint = Endpoint::try_from(root_dir.path().to_string_lossy().as_ref()).unwrap();
let disk = LocalDisk::new(&endpoint, false).await.unwrap();
let endpoint = Endpoint::try_from(root_dir.path().to_string_lossy().as_ref()).expect("operation should succeed");
let disk = LocalDisk::new(&endpoint, false).await.expect("operation should succeed");
assert!(matches!(disk.get_object_path("bucket", "escape/object.txt"), Err(DiskError::InvalidPath)));
}
@@ -4681,21 +4681,21 @@ mod test {
#[tokio::test]
async fn test_local_disk_file_operations() {
let test_dir = "./test_local_disk_file_ops";
fs::create_dir_all(&test_dir).await.unwrap();
fs::create_dir_all(&test_dir).await.expect("operation should succeed");
let endpoint = Endpoint::try_from(test_dir).unwrap();
let disk = LocalDisk::new(&endpoint, false).await.unwrap();
let endpoint = Endpoint::try_from(test_dir).expect("operation should succeed");
let disk = LocalDisk::new(&endpoint, false).await.expect("operation should succeed");
// Create test volume
disk.make_volume("test-volume").await.unwrap();
disk.make_volume("test-volume").await.expect("operation should succeed");
// Test write and read operations
let test_data: Vec<u8> = vec![1, 2, 3, 4, 5];
disk.write_all("test-volume", "test-file.txt", test_data.clone().into())
.await
.unwrap();
.expect("operation should succeed");
let read_data = disk.read_all("test-volume", "test-file.txt").await.unwrap();
let read_data = disk.read_all("test-volume", "test-file.txt").await.expect("operation should succeed");
assert_eq!(read_data, test_data);
// Test file deletion
@@ -4705,38 +4705,38 @@ mod test {
undo_write: false,
old_data_dir: None,
};
disk.delete("test-volume", "test-file.txt", delete_opts).await.unwrap();
disk.delete("test-volume", "test-file.txt", delete_opts).await.expect("operation should succeed");
// Clean up
disk.delete_volume("test-volume").await.unwrap();
disk.delete_volume("test-volume").await.expect("operation should succeed");
let _ = fs::remove_dir_all(&test_dir).await;
}
#[tokio::test]
async fn test_local_disk_volume_operations() {
let test_dir = "./test_local_disk_volumes";
fs::create_dir_all(&test_dir).await.unwrap();
fs::create_dir_all(&test_dir).await.expect("operation should succeed");
let endpoint = Endpoint::try_from(test_dir).unwrap();
let disk = LocalDisk::new(&endpoint, false).await.unwrap();
let endpoint = Endpoint::try_from(test_dir).expect("operation should succeed");
let disk = LocalDisk::new(&endpoint, false).await.expect("operation should succeed");
// Test creating multiple volumes
let volumes = vec!["vol1", "vol2", "vol3"];
disk.make_volumes(volumes.clone()).await.unwrap();
disk.make_volumes(volumes.clone()).await.expect("operation should succeed");
// Test listing volumes
let volume_list = disk.list_volumes().await.unwrap();
let volume_list = disk.list_volumes().await.expect("operation should succeed");
assert!(!volume_list.is_empty());
// Test volume stats
for vol in &volumes {
let vol_info = disk.stat_volume(vol).await.unwrap();
let vol_info = disk.stat_volume(vol).await.expect("operation should succeed");
assert_eq!(vol_info.name, *vol);
}
// Test deleting volumes
for vol in &volumes {
disk.delete_volume(vol).await.unwrap();
disk.delete_volume(vol).await.expect("operation should succeed");
}
// Clean up the test directory
@@ -4746,10 +4746,10 @@ mod test {
#[tokio::test]
async fn test_local_disk_disk_info() {
let test_dir = "./test_local_disk_info";
fs::create_dir_all(&test_dir).await.unwrap();
fs::create_dir_all(&test_dir).await.expect("operation should succeed");
let endpoint = Endpoint::try_from(test_dir).unwrap();
let disk = LocalDisk::new(&endpoint, false).await.unwrap();
let endpoint = Endpoint::try_from(test_dir).expect("operation should succeed");
let disk = LocalDisk::new(&endpoint, false).await.expect("operation should succeed");
let disk_info_opts = DiskInfoOptions {
disk_id: "test-disk".to_string(),
@@ -4757,7 +4757,7 @@ mod test {
noop: false,
};
let disk_info = disk.disk_info(&disk_info_opts).await.unwrap();
let disk_info = disk.disk_info(&disk_info_opts).await.expect("operation should succeed");
// Basic checks on disk info
// Note: On macOS, Windows, and some other systems, fs_type may be empty
@@ -4780,14 +4780,14 @@ mod test {
async fn test_read_file_stream_rejects_offset_length_overflow() {
use tempfile::tempdir;
let dir = tempdir().unwrap();
let endpoint = Endpoint::try_from(dir.path().to_str().unwrap()).unwrap();
let disk = LocalDisk::new(&endpoint, false).await.unwrap();
let dir = tempdir().expect("operation should succeed");
let endpoint = Endpoint::try_from(dir.path().to_str().expect("operation should succeed")).expect("operation should succeed");
let disk = LocalDisk::new(&endpoint, false).await.expect("operation should succeed");
disk.make_volume("test-volume").await.unwrap();
disk.make_volume("test-volume").await.expect("operation should succeed");
disk.write_all("test-volume", "test-file.txt", Bytes::from_static(b"test"))
.await
.unwrap();
.expect("operation should succeed");
let result = disk.read_file_stream("test-volume", "test-file.txt", usize::MAX, 1).await;
assert!(matches!(result, Err(DiskError::FileCorrupt)));
@@ -4797,14 +4797,14 @@ mod test {
async fn test_read_file_zero_copy_rejects_offset_length_overflow() {
use tempfile::tempdir;
let dir = tempdir().unwrap();
let endpoint = Endpoint::try_from(dir.path().to_str().unwrap()).unwrap();
let disk = LocalDisk::new(&endpoint, false).await.unwrap();
let dir = tempdir().expect("operation should succeed");
let endpoint = Endpoint::try_from(dir.path().to_str().expect("operation should succeed")).expect("operation should succeed");
let disk = LocalDisk::new(&endpoint, false).await.expect("operation should succeed");
disk.make_volume("test-volume").await.unwrap();
disk.make_volume("test-volume").await.expect("operation should succeed");
disk.write_all("test-volume", "test-file.txt", Bytes::from_static(b"test"))
.await
.unwrap();
.expect("operation should succeed");
let result = disk.read_file_zero_copy("test-volume", "test-file.txt", usize::MAX, 1).await;
assert!(matches!(result, Err(DiskError::FileCorrupt)));
@@ -4857,15 +4857,15 @@ mod test {
let test_file = "./test_read_exists.txt";
// Test non-existent file
let (data, metadata) = read_file_exists(test_file).await.unwrap();
let (data, metadata) = read_file_exists(test_file).await.expect("operation should succeed");
assert!(data.is_empty());
assert!(metadata.is_none());
// Create test file
fs::write(test_file, b"test content").await.unwrap();
fs::write(test_file, b"test content").await.expect("operation should succeed");
// Test existing file
let (data, metadata) = read_file_exists(test_file).await.unwrap();
let (data, metadata) = read_file_exists(test_file).await.expect("operation should succeed");
assert_eq!(data.as_ref(), b"test content");
assert!(metadata.is_some());
@@ -4879,10 +4879,10 @@ mod test {
let test_content = b"test content for read_all";
// Create test file
fs::write(test_file, test_content).await.unwrap();
fs::write(test_file, test_content).await.expect("operation should succeed");
// Test reading file
let (data, metadata) = read_file_all(test_file).await.unwrap();
let (data, metadata) = read_file_all(test_file).await.expect("operation should succeed");
assert_eq!(data.as_ref(), test_content);
assert!(metadata.is_file());
assert_eq!(metadata.len(), test_content.len() as u64);
@@ -4896,10 +4896,10 @@ mod test {
let test_file = "./test_metadata.txt";
// Create test file
fs::write(test_file, b"test").await.unwrap();
fs::write(test_file, b"test").await.expect("operation should succeed");
// Test reading metadata
let metadata = read_file_metadata(test_file).await.unwrap();
let metadata = read_file_metadata(test_file).await.expect("operation should succeed");
assert!(metadata.is_file());
assert_eq!(metadata.len(), 4); // "test" is 4 bytes
+52 -52
View File
@@ -412,7 +412,7 @@ fn recover_empty_payload_data_shards(
/// use rustfs_ecstore::api::erasure::Erasure;
/// let erasure = Erasure::new(4, 2, 8);
/// let data = b"hello world";
/// let shards = erasure.encode_data(data).unwrap();
/// let shards = erasure.encode_data(data).expect("operation should succeed");
/// // Simulate loss and recovery...
/// ```
pub struct Erasure {
@@ -474,13 +474,13 @@ impl Erasure {
/// for decode/reconstruct (for reading and healing old-version files).
pub fn new_with_options(data_shards: usize, parity_shards: usize, block_size: usize, uses_legacy: bool) -> Self {
let encoder = if !uses_legacy && parity_shards > 0 {
Some(ReedSolomonEncoder::new(data_shards, parity_shards).unwrap())
Some(ReedSolomonEncoder::new(data_shards, parity_shards).expect("operation should succeed"))
} else {
None
};
let legacy_encoder = if uses_legacy && parity_shards > 0 {
Some(LegacyReedSolomonEncoder::new(data_shards, parity_shards).unwrap())
Some(LegacyReedSolomonEncoder::new(data_shards, parity_shards).expect("operation should succeed"))
} else {
None
};
@@ -1208,7 +1208,7 @@ mod tests {
let test_data = b"SIMD mode test data for encoding and decoding roundtrip verification with sufficient length to ensure shard size requirements are met for proper SIMD optimization.".repeat(20); // ~3KB for SIMD
let data = &test_data;
let encoded_shards = erasure.encode_data(data).unwrap();
let encoded_shards = erasure.encode_data(data).expect("operation should succeed");
assert_eq!(encoded_shards.len(), data_shards + parity_shards);
// Create decode input with some shards missing, convert to the format expected by decode_data
@@ -1217,12 +1217,12 @@ mod tests {
decode_input[i] = Some(encoded_shards[i].to_vec());
}
erasure.decode_data(&mut decode_input).unwrap();
erasure.decode_data(&mut decode_input).expect("operation should succeed");
// Recover original data
let mut recovered = Vec::new();
for shard in decode_input.iter().take(data_shards) {
recovered.extend_from_slice(shard.as_ref().unwrap());
recovered.extend_from_slice(shard.as_ref().expect("operation should succeed"));
}
recovered.truncate(data.len());
assert_eq!(&recovered, data);
@@ -1238,7 +1238,7 @@ mod tests {
// Generate 1MB test data
let data: Vec<u8> = (0..1048576).map(|i| (i % 256) as u8).collect();
let encoded_shards = erasure.encode_data(&data).unwrap();
let encoded_shards = erasure.encode_data(&data).expect("operation should succeed");
assert_eq!(encoded_shards.len(), data_shards + parity_shards);
// Create decode input with some shards missing, convert to the format expected by decode_data
@@ -1247,12 +1247,12 @@ mod tests {
decode_input[i] = Some(encoded_shards[i].to_vec());
}
erasure.decode_data(&mut decode_input).unwrap();
erasure.decode_data(&mut decode_input).expect("operation should succeed");
// Recover original data
let mut recovered = Vec::new();
for shard in decode_input.iter().take(data_shards) {
recovered.extend_from_slice(shard.as_ref().unwrap());
recovered.extend_from_slice(shard.as_ref().expect("operation should succeed"));
}
recovered.truncate(data.len());
assert_eq!(recovered, data);
@@ -1265,7 +1265,7 @@ mod tests {
let block_size = 6;
let erasure = Erasure::new(data_shards, parity_shards, block_size);
let data = vec![0u8; block_size];
let shards = erasure.encode_data(&data).unwrap();
let shards = erasure.encode_data(&data).expect("operation should succeed");
assert_eq!(shards.len(), data_shards + parity_shards);
let total_len: usize = shards.iter().map(|b| b.len()).sum();
assert_eq!(total_len, erasure.shard_size() * (data_shards + parity_shards));
@@ -1296,7 +1296,7 @@ mod tests {
let erasure = Erasure::new_with_options(data_shards, parity_shards, block_size, true);
let data = b"Legacy encode/decode roundtrip test data with sufficient length.".repeat(20);
let encoded_shards = erasure.encode_data(&data).unwrap();
let encoded_shards = erasure.encode_data(&data).expect("operation should succeed");
assert_eq!(encoded_shards.len(), data_shards + parity_shards);
let mut decode_input: Vec<Option<Vec<u8>>> = vec![None; data_shards + parity_shards];
@@ -1304,11 +1304,11 @@ mod tests {
decode_input[i] = Some(encoded_shards[i].to_vec());
}
erasure.decode_data(&mut decode_input).unwrap();
erasure.decode_data(&mut decode_input).expect("operation should succeed");
let mut recovered = Vec::new();
for shard in decode_input.iter().take(data_shards) {
recovered.extend_from_slice(shard.as_ref().unwrap());
recovered.extend_from_slice(shard.as_ref().expect("operation should succeed"));
}
recovered.truncate(data.len());
assert_eq!(&recovered, &data);
@@ -1322,17 +1322,17 @@ mod tests {
let erasure = Erasure::new_with_options(data_shards, parity_shards, block_size, true);
let data = b"Legacy decode with missing shards test.".repeat(10);
let encoded_shards = erasure.encode_data(&data).unwrap();
let encoded_shards = erasure.encode_data(&data).expect("operation should succeed");
let mut shards_opt: Vec<Option<Vec<u8>>> = encoded_shards.iter().map(|s| Some(s.to_vec())).collect();
shards_opt[1] = None;
shards_opt[5] = None;
erasure.decode_data(&mut shards_opt).unwrap();
erasure.decode_data(&mut shards_opt).expect("operation should succeed");
let mut recovered = Vec::new();
for shard in shards_opt.iter().take(data_shards) {
recovered.extend_from_slice(shard.as_ref().unwrap());
recovered.extend_from_slice(shard.as_ref().expect("operation should succeed"));
}
recovered.truncate(data.len());
assert_eq!(&recovered, &data);
@@ -1373,17 +1373,17 @@ mod tests {
.encode_stream_callback_async::<_, _, (), _>(&mut reader, move |res| {
let tx = tx.clone();
async move {
let shards = res.unwrap();
tx.send(shards).await.unwrap();
let shards = res.expect("operation should succeed");
tx.send(shards).await.expect("operation should succeed");
Ok(())
}
})
.await
.unwrap();
.expect("operation should succeed");
});
let result = handle.await;
assert!(result.is_ok());
let collected_shards = rx.recv().await.unwrap();
let collected_shards = rx.recv().await.expect("operation should succeed");
assert_eq!(collected_shards.len(), data_shards + parity_shards);
}
@@ -1412,17 +1412,17 @@ mod tests {
.encode_stream_callback_async::<_, _, (), _>(&mut reader, move |res| {
let tx = tx.clone();
async move {
let shards = res.unwrap();
tx.send(shards).await.unwrap();
let shards = res.expect("operation should succeed");
tx.send(shards).await.expect("operation should succeed");
Ok(())
}
})
.await
.unwrap();
.expect("operation should succeed");
});
let result = handle.await;
assert!(result.is_ok());
let shards = rx.recv().await.unwrap();
let shards = rx.recv().await.expect("operation should succeed");
assert_eq!(shards.len(), data_shards + parity_shards);
// Test decode using the old API that operates in-place
@@ -1430,12 +1430,12 @@ mod tests {
for i in 0..data_shards {
decode_input[i] = Some(shards[i].to_vec());
}
erasure.decode_data(&mut decode_input).unwrap();
erasure.decode_data(&mut decode_input).expect("operation should succeed");
// Recover original data
let mut recovered = Vec::new();
for shard in decode_input.iter().take(data_shards) {
recovered.extend_from_slice(shard.as_ref().unwrap());
recovered.extend_from_slice(shard.as_ref().expect("operation should succeed"));
}
recovered.truncate(data_clone.len());
assert_eq!(&recovered, &data_clone);
@@ -1456,7 +1456,7 @@ mod tests {
let observed = observed_clone.clone();
async move {
let err = res.expect_err("zero block size should report an error");
*observed.lock().unwrap() = Some((err.kind(), err.to_string()));
*observed.lock().expect("operation should succeed") = Some((err.kind(), err.to_string()));
Ok(())
}
})
@@ -1464,7 +1464,7 @@ mod tests {
.expect("callback should handle the zero block size error");
assert_eq!(total, 0);
let observed = observed.lock().unwrap();
let observed = observed.lock().expect("operation should succeed");
let (kind, message) = observed.as_ref().expect("callback should be invoked once");
assert_eq!(*kind, io::ErrorKind::InvalidInput);
assert!(message.contains("block_size"));
@@ -1485,7 +1485,7 @@ mod tests {
let test_data = b"SIMD mode test data for encoding and decoding roundtrip verification with sufficient length to ensure shard size requirements are met for proper SIMD optimization and validation.";
let data = test_data.repeat(25); // Create much larger data: ~5KB total, ~1.25KB per shard
let encoded_shards = erasure.encode_data(&data).unwrap();
let encoded_shards = erasure.encode_data(&data).expect("operation should succeed");
assert_eq!(encoded_shards.len(), data_shards + parity_shards);
// Create decode input with some shards missing
@@ -1495,12 +1495,12 @@ mod tests {
shards_opt[1] = None; // Lose second data shard
shards_opt[5] = None; // Lose second parity shard
erasure.decode_data(&mut shards_opt).unwrap();
erasure.decode_data(&mut shards_opt).expect("operation should succeed");
// Verify recovered data
let mut recovered = Vec::new();
for shard in shards_opt.iter().take(data_shards) {
recovered.extend_from_slice(shard.as_ref().unwrap());
recovered.extend_from_slice(shard.as_ref().expect("operation should succeed"));
}
recovered.truncate(data.len());
assert_eq!(&recovered, &data);
@@ -1516,7 +1516,7 @@ mod tests {
// Create all-zero data that ensures adequate shard size for SIMD optimization
let data = vec![0u8; 1024]; // 1KB of zeros, each shard will be 256 bytes
let encoded_shards = erasure.encode_data(&data).unwrap();
let encoded_shards = erasure.encode_data(&data).expect("operation should succeed");
assert_eq!(encoded_shards.len(), data_shards + parity_shards);
// Verify that all data shards are zeros
@@ -1531,12 +1531,12 @@ mod tests {
shards_opt[0] = None; // Lose first data shard
shards_opt[4] = None; // Lose first parity shard
erasure.decode_data(&mut shards_opt).unwrap();
erasure.decode_data(&mut shards_opt).expect("operation should succeed");
// Verify recovered data is still all zeros
let mut recovered = Vec::new();
for shard in shards_opt.iter().take(data_shards) {
recovered.extend_from_slice(shard.as_ref().unwrap());
recovered.extend_from_slice(shard.as_ref().expect("operation should succeed"));
}
recovered.truncate(data.len());
assert!(recovered.iter().all(|&x| x == 0), "Recovered data should be all zeros");
@@ -1555,7 +1555,7 @@ mod tests {
data.push((i % 256) as u8);
}
let shards = erasure.encode_data(&data).unwrap();
let shards = erasure.encode_data(&data).expect("operation should succeed");
assert_eq!(shards.len(), data_shards + parity_shards);
// Simulate the loss of multiple shards
@@ -1566,12 +1566,12 @@ mod tests {
shards_opt[11] = None; // Parity shard
// Decode
erasure.decode_data(&mut shards_opt).unwrap();
erasure.decode_data(&mut shards_opt).expect("operation should succeed");
// Recover original data
let mut recovered = Vec::new();
for shard in shards_opt.iter().take(data_shards) {
recovered.extend_from_slice(shard.as_ref().unwrap());
recovered.extend_from_slice(shard.as_ref().expect("operation should succeed"));
}
recovered.truncate(data.len());
assert_eq!(&recovered, &data);
@@ -1603,7 +1603,7 @@ mod tests {
Ok(_) => {
let mut recovered = Vec::new();
for shard in shards_opt.iter().take(data_shards) {
recovered.extend_from_slice(shard.as_ref().unwrap());
recovered.extend_from_slice(shard.as_ref().expect("operation should succeed"));
}
recovered.truncate(data.len());
assert_eq!(&recovered, &data);
@@ -1630,7 +1630,7 @@ mod tests {
let data =
b"Testing maximum erasure capacity with SIMD Reed-Solomon implementation for robustness verification!".repeat(3);
let shards = erasure.encode_data(&data).unwrap();
let shards = erasure.encode_data(&data).expect("operation should succeed");
// Lose exactly the maximum number of shards (equal to parity_shards)
let mut shards_opt: Vec<Option<Vec<u8>>> = shards.iter().map(|b| Some(b.to_vec())).collect();
@@ -1639,11 +1639,11 @@ mod tests {
shards_opt[6] = None; // Parity shard
// Should succeed with maximum erasures
erasure.decode_data(&mut shards_opt).unwrap();
erasure.decode_data(&mut shards_opt).expect("operation should succeed");
let mut recovered = Vec::new();
for shard in shards_opt.iter().take(data_shards) {
recovered.extend_from_slice(shard.as_ref().unwrap());
recovered.extend_from_slice(shard.as_ref().expect("operation should succeed"));
}
recovered.truncate(data.len());
assert_eq!(&recovered, &data);
@@ -1662,7 +1662,7 @@ mod tests {
fn test_reed_solomon_compat() {
let data = generate_compat_test_data(7557);
let erasure = Erasure::new(4, 2, 7557);
let shards = erasure.encode_data(&data).unwrap();
let shards = erasure.encode_data(&data).expect("operation should succeed");
assert_eq!(shards.len(), 6, "expected 6 shards (4 data + 2 parity)");
// Per-shard HighwayHash
@@ -1732,7 +1732,7 @@ mod tests {
// Verify recovered data
let mut recovered = Vec::new();
for shard in shards_opt.iter().take(data_shards) {
recovered.extend_from_slice(shard.as_ref().unwrap());
recovered.extend_from_slice(shard.as_ref().expect("operation should succeed"));
}
recovered.truncate(small_data.len());
println!("recovered: {recovered:?}");
@@ -1776,7 +1776,7 @@ mod tests {
// Encode the data
let start = std::time::Instant::now();
let shards = erasure.encode_data(&data).unwrap();
let shards = erasure.encode_data(&data).expect("operation should succeed");
let encode_duration = start.elapsed();
println!("⏱️ Encoding completed in: {encode_duration:?}");
@@ -1800,7 +1800,7 @@ mod tests {
// Decode and recover data
let start = std::time::Instant::now();
erasure.decode_data(&mut shards_opt).unwrap();
erasure.decode_data(&mut shards_opt).expect("operation should succeed");
let decode_duration = start.elapsed();
println!("⏱️ Decoding completed in: {decode_duration:?}");
@@ -1808,7 +1808,7 @@ mod tests {
// Verify recovered data integrity
let mut recovered = Vec::new();
for shard in shards_opt.iter().take(data_shards) {
recovered.extend_from_slice(shard.as_ref().unwrap());
recovered.extend_from_slice(shard.as_ref().expect("operation should succeed"));
}
recovered.truncate(data.len());
@@ -1853,20 +1853,20 @@ mod tests {
.encode_stream_callback_async::<_, _, (), _>(&mut reader, move |res| {
let tx = tx.clone();
async move {
let shards = res.unwrap();
tx.send(shards).await.unwrap();
let shards = res.expect("operation should succeed");
tx.send(shards).await.expect("operation should succeed");
Ok(())
}
})
.await
.unwrap();
.expect("operation should succeed");
});
let mut all_blocks = Vec::new();
while let Some(block) = rx.recv().await {
all_blocks.push(block);
}
handle.await.unwrap();
handle.await.expect("operation should succeed");
// Verify we got multiple blocks
assert!(all_blocks.len() > 1, "Should have multiple blocks for stream test");
@@ -1879,10 +1879,10 @@ mod tests {
shards_opt[1] = None;
shards_opt[5] = None;
erasure.decode_data(&mut shards_opt).unwrap();
erasure.decode_data(&mut shards_opt).expect("operation should succeed");
for shard in shards_opt.iter().take(data_shards) {
recovered.extend_from_slice(shard.as_ref().unwrap());
recovered.extend_from_slice(shard.as_ref().expect("operation should succeed"));
}
}
+2 -2
View File
@@ -429,7 +429,7 @@ impl SetDisks {
"find_file_info_in_quorum: inspecting meta"
);
let etag_only = mod_time.is_none() && etag.is_some() && meta.get_etag().is_some_and(|v| &v == etag.as_ref().unwrap());
let etag_only = mod_time.is_none() && etag.is_some() && meta.get_etag().is_some_and(|v| &v == etag.as_ref().expect("operation should succeed"));
let mod_valid = mod_time == &meta.mod_time;
if etag_only || mod_valid {
@@ -509,7 +509,7 @@ impl SetDisks {
}
if found {
let mut fi = found_fi.unwrap();
let mut fi = found_fi.expect("operation should succeed");
for (val, &count) in &valid_obj_map {
if count >= quorum {
+2 -2
View File
@@ -112,7 +112,7 @@ pub fn set_global_rustfs_port(value: u16) {
/// * None
///
pub fn set_global_deployment_id(id: Uuid) {
globalDeploymentIDPtr.set(id).unwrap();
globalDeploymentIDPtr.set(id).expect("operation should succeed");
}
/// Get the global deployment id
@@ -274,7 +274,7 @@ pub(crate) type TypeLocalDiskSetDrives = Vec<Vec<Vec<Option<DiskStore>>>>;
/// # Returns
/// * None
pub fn set_global_region(region: s3s::region::Region) {
GLOBAL_REGION.set(region).unwrap();
GLOBAL_REGION.set(region).expect("operation should succeed");
}
/// Get the global region
+5 -5
View File
@@ -352,8 +352,8 @@ impl SetDisks {
// We write at temporary location and then rename to final location.
let tmp_id = Uuid::new_v4().to_string();
let src_data_dir = latest_meta.data_dir.unwrap().to_string();
let dst_data_dir = latest_meta.data_dir.unwrap();
let src_data_dir = latest_meta.data_dir.expect("operation should succeed").to_string();
let dst_data_dir = latest_meta.data_dir.expect("operation should succeed");
if !latest_meta.deleted && !latest_meta.is_remote() {
let erasure_info = latest_meta.erasure.clone();
@@ -549,13 +549,13 @@ impl SetDisks {
} else {
rename_successes += 1;
if parts_metadata[index].is_remote() {
let rm_data_dir = parts_metadata[index].data_dir.unwrap().to_string();
let rm_data_dir = parts_metadata[index].data_dir.expect("operation should succeed").to_string();
let d_path = Path::new(&encode_dir_object(object)).join(rm_data_dir);
disk.delete(
bucket,
d_path.to_str().unwrap(),
d_path.to_str().expect("operation should succeed"),
DeleteOptions {
immediate: true,
@@ -713,7 +713,7 @@ impl SetDisks {
for (index, (err, disk)) in errs.iter().zip(disks.iter()).enumerate() {
if let (Some(DiskError::VolumeNotFound | DiskError::FileNotFound), Some(disk)) = (err, disk) {
let vol_path = Path::new(bucket).join(object);
let drive_state = match disk.make_volume(vol_path.to_str().unwrap()).await {
let drive_state = match disk.make_volume(vol_path.to_str().expect("operation should succeed")).await {
Ok(_) => DriveState::Ok.to_string(),
Err(merr) => match merr {
DiskError::VolumeExists => DriveState::Ok.to_string(),
+2 -2
View File
@@ -286,8 +286,8 @@ impl ECStore {
}
for disk in disks.iter() {
if disk.is_some() && disk.as_ref().unwrap().is_local() {
local_disks.push(disk.as_ref().unwrap().clone());
if disk.is_some() && disk.as_ref().expect("operation should succeed").is_local() {
local_disks.push(disk.as_ref().expect("operation should succeed").clone());
}
}
+6 -6
View File
@@ -36,7 +36,7 @@ impl ECStore {
// if disk.is_none() {
// continue;
// }
// // let disk = disk.as_ref().unwrap().clone();
// // let disk = disk.as_ref().expect("operation should succeed").clone();
// // futures.push(disk.delete(
// // bucket,
// // prefix,
@@ -315,7 +315,7 @@ impl ECStore {
return Ok((pinfo.clone(), self.pools_with_object(&ress, opts).await));
}
let err = pinfo.err.as_ref().unwrap();
let err = pinfo.err.as_ref().expect("operation should succeed");
if err == &Error::ErasureReadQuorum && !opts.metadata_chg {
return Ok((pinfo.clone(), self.pools_with_object(&ress, opts).await));
@@ -460,7 +460,7 @@ impl ECStore {
let mut pool_meta = self.pool_meta.write().await;
*pool_meta = meta;
// *self.pool_meta.write().unwrap() = meta;
// *self.pool_meta.write().expect("operation should succeed") = meta;
Ok(())
}
@@ -638,7 +638,7 @@ mod tests {
fn object_info_with_mod_time(unix_ts: i64, delete_marker: bool) -> ObjectInfo {
ObjectInfo {
mod_time: Some(OffsetDateTime::from_unix_timestamp(unix_ts).unwrap()),
mod_time: Some(OffsetDateTime::from_unix_timestamp(unix_ts).expect("operation should succeed")),
delete_marker,
..Default::default()
}
@@ -660,7 +660,7 @@ mod tests {
];
let (info, idx) =
resolve_latest_object_info_candidates(candidates, "bucket", "object", &ObjectOptions::default()).unwrap();
resolve_latest_object_info_candidates(candidates, "bucket", "object", &ObjectOptions::default()).expect("operation should succeed");
assert_eq!(idx, 1);
assert!(info.delete_marker);
@@ -681,7 +681,7 @@ mod tests {
},
];
let (_, idx) = resolve_latest_object_info_candidates(candidates, "bucket", "object", &ObjectOptions::default()).unwrap();
let (_, idx) = resolve_latest_object_info_candidates(candidates, "bucket", "object", &ObjectOptions::default()).expect("operation should succeed");
assert_eq!(idx, 1);
}
+2 -2
View File
@@ -383,7 +383,7 @@ impl PoolTier {
// Use the pool's shared metrics for recording
let _metrics_lock = self.metrics.lock().unwrap_or_else(|e| e.into_inner());
let _metrics = _metrics_lock.as_ref().unwrap();
let _metrics = _metrics_lock.as_ref().expect("operation should succeed");
// Record acquisition
pool_metrics.total_acquires.fetch_add(1, Ordering::Relaxed);
@@ -406,7 +406,7 @@ impl PoolTier {
// Use the pool's shared metrics for recording
let _metrics_lock = self.metrics.lock().unwrap_or_else(|e| e.into_inner());
let _metrics = _metrics_lock.as_ref().unwrap();
let _metrics = _metrics_lock.as_ref().expect("operation should succeed");
// Record acquisition
pool_metrics.total_acquires.fetch_add(1, Ordering::Relaxed);
+1 -1
View File
@@ -173,7 +173,7 @@ pub fn set_global_guard(guard: OtelGuard) -> Result<(), GlobalError> {
///
/// # async fn trace_operation() -> Result<(), Box<dyn std::error::Error>> {
/// # let guard = get_global_guard()?;
/// # let _lock = guard.lock().unwrap();
/// # let _lock = guard.lock().expect("operation should succeed");
/// # // Perform traced operation
/// # Ok(())
/// # }
@@ -25,7 +25,7 @@
//! use rustfs_obs::metrics::collectors::{GpuCollector, collect_gpu_metrics};
//! use sysinfo::Pid;
//!
//! let pid = sysinfo::get_current_pid().unwrap();
//! let pid = sysinfo::get_current_pid().expect("operation should succeed");
//! let collector = GpuCollector::new(pid)?;
//! let stats = collector.collect()?;
//! let metrics = collect_gpu_metrics(&stats, &labels);
@@ -105,7 +105,7 @@ impl GpuCollector {
/// use rustfs_obs::metrics::collectors::GpuCollector;
/// use sysinfo::Pid;
///
/// let pid = sysinfo::get_current_pid().unwrap();
/// let pid = sysinfo::get_current_pid().expect("operation should succeed");
/// let collector = GpuCollector::new(pid)?;
/// ```
pub fn new(pid: Pid) -> Result<Self, GpuError> {
@@ -342,7 +342,7 @@ impl ExpirationWorker {
metrics: &Arc<RwLock<ExpirationMetrics>>,
) -> SwiftResult<()> {
let start_time = SystemTime::now();
let now = start_time.duration_since(UNIX_EPOCH).unwrap().as_secs();
let now = start_time.duration_since(UNIX_EPOCH).expect("operation should succeed").as_secs();
debug!(
event = EVENT_SWIFT_EXPIRATION_ITERATION_SUMMARY,
@@ -375,7 +375,7 @@ impl ExpirationWorker {
}
// Remove from queue and add to batch
let entry = queue.pop().unwrap().0;
let entry = queue.pop().expect("operation should succeed").0;
drop(queue); // Release lock
batch.push(entry);
@@ -452,7 +452,7 @@ impl ExpirationWorker {
}
// Update metrics
let duration = SystemTime::now().duration_since(start_time).unwrap();
let duration = SystemTime::now().duration_since(start_time).expect("operation should succeed");
let mut m = metrics.write().await;
m.objects_scanned += scanned_count;
m.objects_deleted += deleted_count;
@@ -589,9 +589,9 @@ mod tests {
}));
// Should pop in order: 1000, 2000, 3000
assert_eq!(heap.pop().unwrap().0.expires_at, 1000);
assert_eq!(heap.pop().unwrap().0.expires_at, 2000);
assert_eq!(heap.pop().unwrap().0.expires_at, 3000);
assert_eq!(heap.pop().expect("operation should succeed").0.expires_at, 1000);
assert_eq!(heap.pop().expect("operation should succeed").0.expires_at, 2000);
assert_eq!(heap.pop().expect("operation should succeed").0.expires_at, 3000);
}
#[test]
+6 -6
View File
@@ -215,7 +215,7 @@ impl RateLimiter {
/// Returns (remaining, reset_timestamp) if successful,
/// or SwiftError::TooManyRequests if rate limited
pub fn check_rate_limit(&self, key: &str, rate_limit: &RateLimit) -> SwiftResult<(u32, u64)> {
let mut buckets = self.buckets.lock().unwrap();
let mut buckets = self.buckets.lock().expect("operation should succeed");
// Get or create bucket for this key
let bucket = buckets.entry(key.to_string()).or_insert_with(|| TokenBucket::new(rate_limit));
@@ -241,7 +241,7 @@ impl RateLimiter {
/// Get current rate limit status without consuming quota
pub fn get_status(&self, key: &str, rate_limit: &RateLimit) -> (u32, u64) {
let mut buckets = self.buckets.lock().unwrap();
let mut buckets = self.buckets.lock().expect("operation should succeed");
let bucket = buckets.entry(key.to_string()).or_insert_with(|| TokenBucket::new(rate_limit));
@@ -292,7 +292,7 @@ mod tests {
#[test]
fn test_parse_rate_limit_valid() {
let rate_limit = RateLimit::parse("1000/60").unwrap();
let rate_limit = RateLimit::parse("1000/60").expect("operation should succeed");
assert_eq!(rate_limit.limit, 1000);
assert_eq!(rate_limit.window_seconds, 60);
}
@@ -362,7 +362,7 @@ mod tests {
// Consume 10
for _ in 0..10 {
bucket.try_consume().unwrap();
bucket.try_consume().expect("operation should succeed");
}
assert_eq!(bucket.remaining(), 90);
@@ -395,7 +395,7 @@ mod tests {
let rate_limit = extract_rate_limit(&metadata);
assert!(rate_limit.is_some());
let rate_limit = rate_limit.unwrap();
let rate_limit = rate_limit.expect("operation should succeed");
assert_eq!(rate_limit.limit, 1000);
assert_eq!(rate_limit.window_seconds, 60);
}
@@ -408,7 +408,7 @@ mod tests {
let rate_limit = extract_rate_limit(&metadata);
assert!(rate_limit.is_some());
let rate_limit = rate_limit.unwrap();
let rate_limit = rate_limit.expect("operation should succeed");
assert_eq!(rate_limit.limit, 100);
assert_eq!(rate_limit.window_seconds, 60);
}
+20 -20
View File
@@ -636,7 +636,7 @@ mod tests {
let reader = BufReader::new(Cursor::new(data.to_vec()));
let mut encrypt_reader = EncryptReader::new(reader, key, nonce);
let mut encrypted = Vec::new();
encrypt_reader.read_to_end(&mut encrypted).await.unwrap();
encrypt_reader.read_to_end(&mut encrypted).await.expect("operation should succeed");
encrypted
}
@@ -674,14 +674,14 @@ mod tests {
// Encrypt
let mut encrypt_reader = encrypt_reader;
let mut encrypted = Vec::new();
encrypt_reader.read_to_end(&mut encrypted).await.unwrap();
encrypt_reader.read_to_end(&mut encrypted).await.expect("operation should succeed");
// Decrypt using DecryptReader
let reader = Cursor::new(encrypted.clone());
let decrypt_reader = DecryptReader::new(reader, key, nonce);
let mut decrypt_reader = decrypt_reader;
let mut decrypted = Vec::new();
decrypt_reader.read_to_end(&mut decrypted).await.unwrap();
decrypt_reader.read_to_end(&mut decrypted).await.expect("operation should succeed");
assert_eq!(&decrypted, data);
}
@@ -700,7 +700,7 @@ mod tests {
let encrypt_reader = EncryptReader::new(reader, key, nonce);
let mut encrypt_reader = encrypt_reader;
let mut encrypted = Vec::new();
encrypt_reader.read_to_end(&mut encrypted).await.unwrap();
encrypt_reader.read_to_end(&mut encrypted).await.expect("operation should succeed");
// Now test DecryptReader
@@ -708,7 +708,7 @@ mod tests {
let decrypt_reader = DecryptReader::new(reader, key, nonce);
let mut decrypt_reader = decrypt_reader;
let mut decrypted = Vec::new();
decrypt_reader.read_to_end(&mut decrypted).await.unwrap();
decrypt_reader.read_to_end(&mut decrypted).await.expect("operation should succeed");
assert_eq!(&decrypted, data);
}
@@ -728,13 +728,13 @@ mod tests {
let encrypt_reader = EncryptReader::new(reader, key, nonce);
let mut encrypt_reader = encrypt_reader;
let mut encrypted = Vec::new();
encrypt_reader.read_to_end(&mut encrypted).await.unwrap();
encrypt_reader.read_to_end(&mut encrypted).await.expect("operation should succeed");
let reader = std::io::Cursor::new(encrypted.clone());
let decrypt_reader = DecryptReader::new(reader, key, nonce);
let mut decrypt_reader = decrypt_reader;
let mut decrypted = Vec::new();
decrypt_reader.read_to_end(&mut decrypted).await.unwrap();
decrypt_reader.read_to_end(&mut decrypted).await.expect("operation should succeed");
assert_eq!(&decrypted, &data);
}
@@ -752,12 +752,12 @@ mod tests {
let reader = Cursor::new(data.clone());
let mut encrypt_reader = EncryptReader::new(reader, key, nonce);
let mut encrypted = Vec::new();
encrypt_reader.read_to_end(&mut encrypted).await.unwrap();
encrypt_reader.read_to_end(&mut encrypted).await.expect("operation should succeed");
let reader = ChunkedCursor::new(encrypted, 3);
let mut decrypt_reader = DecryptReader::new(reader, key, nonce);
let mut decrypted = Vec::new();
decrypt_reader.read_to_end(&mut decrypted).await.unwrap();
decrypt_reader.read_to_end(&mut decrypted).await.expect("operation should succeed");
assert_eq!(decrypted, data);
}
@@ -775,12 +775,12 @@ mod tests {
let reader = Cursor::new(data.clone());
let mut encrypt_reader = EncryptReader::new(reader, key, nonce);
let mut encrypted = Vec::new();
encrypt_reader.read_to_end(&mut encrypted).await.unwrap();
encrypt_reader.read_to_end(&mut encrypted).await.expect("operation should succeed");
let reader = PendingChunkedCursor::new(encrypted, 3);
let mut decrypt_reader = DecryptReader::new(reader, key, nonce);
let mut decrypted = Vec::new();
decrypt_reader.read_to_end(&mut decrypted).await.unwrap();
decrypt_reader.read_to_end(&mut decrypted).await.expect("operation should succeed");
assert_eq!(decrypted, data);
}
@@ -798,7 +798,7 @@ mod tests {
let reader = Cursor::new(data.clone());
let mut encrypt_reader = EncryptReader::new(reader, key, nonce);
let mut encrypted = Vec::new();
encrypt_reader.read_to_end(&mut encrypted).await.unwrap();
encrypt_reader.read_to_end(&mut encrypted).await.expect("operation should succeed");
let reader = ChunkedCursor::new(encrypted, 8192);
let decrypt_reader = DecryptReader::new(reader, key, nonce);
@@ -806,7 +806,7 @@ mod tests {
let mut decrypted = Vec::new();
while let Some(chunk) = stream.next().await {
let bytes = chunk.unwrap();
let bytes = chunk.expect("operation should succeed");
decrypted.extend_from_slice(&bytes);
}
@@ -826,7 +826,7 @@ mod tests {
let reader = Cursor::new(data.clone());
let mut encrypt_reader = EncryptReader::new(reader, key, nonce);
let mut encrypted = Vec::new();
encrypt_reader.read_to_end(&mut encrypted).await.unwrap();
encrypt_reader.read_to_end(&mut encrypted).await.expect("operation should succeed");
let reader = ChunkedCursor::new(encrypted, 8192);
let decrypt_reader = DecryptReader::new(reader, key, nonce);
@@ -835,7 +835,7 @@ mod tests {
let mut decrypted = Vec::new();
while let Some(chunk) = stream.next().await {
let bytes = chunk.unwrap();
let bytes = chunk.expect("operation should succeed");
decrypted.extend_from_slice(&bytes);
}
@@ -857,7 +857,7 @@ mod tests {
let reader = BufReader::new(Cursor::new(data.to_vec()));
let mut encrypt_reader = EncryptReader::new(reader, key, nonce);
let mut encrypted = Vec::new();
encrypt_reader.read_to_end(&mut encrypted).await.unwrap();
encrypt_reader.read_to_end(&mut encrypted).await.expect("operation should succeed");
encrypted
}
@@ -871,7 +871,7 @@ mod tests {
let reader = BufReader::new(Cursor::new(combined));
let mut decrypt_reader = DecryptReader::new_multipart(reader, key, base_nonce, vec![1, 2]);
let mut decrypted = Vec::new();
decrypt_reader.read_to_end(&mut decrypted).await.unwrap();
decrypt_reader.read_to_end(&mut decrypted).await.expect("operation should succeed");
let mut expected = Vec::with_capacity(part_one.len() + part_two.len());
expected.extend_from_slice(&part_one);
@@ -891,7 +891,7 @@ mod tests {
let reader = Cursor::new(data);
let mut encrypt_reader = EncryptReader::new(reader, key, nonce);
let mut encrypted = Vec::new();
encrypt_reader.read_to_end(&mut encrypted).await.unwrap();
encrypt_reader.read_to_end(&mut encrypted).await.expect("operation should succeed");
let payloads = extract_encrypted_payloads(&encrypted);
assert!(payloads.len() >= 2);
@@ -920,7 +920,7 @@ mod tests {
let reader = Cursor::new(encrypted);
let mut decrypt_reader = DecryptReader::new(reader, key, nonce);
let mut decrypted = Vec::new();
decrypt_reader.read_to_end(&mut decrypted).await.unwrap();
decrypt_reader.read_to_end(&mut decrypted).await.expect("operation should succeed");
assert_eq!(decrypted, data);
}
@@ -945,7 +945,7 @@ mod tests {
let reader = BufReader::new(Cursor::new(combined));
let mut decrypt_reader = DecryptReader::new_multipart(reader, key, base_nonce, vec![1, 2]);
let mut decrypted = Vec::new();
decrypt_reader.read_to_end(&mut decrypted).await.unwrap();
decrypt_reader.read_to_end(&mut decrypted).await.expect("operation should succeed");
let mut expected = Vec::with_capacity(part_one.len() + part_two.len());
expected.extend_from_slice(&part_one);
+37 -37
View File
@@ -50,13 +50,13 @@
//! let diskable_md5 = false;
//!
//! // Method 1: Simple creation (recommended for most cases)
//! let hash_reader = HashReader::from_stream(reader, size, actual_size, etag.clone(), None, diskable_md5).unwrap();
//! let hash_reader = HashReader::from_stream(reader, size, actual_size, etag.clone(), None, diskable_md5).expect("operation should succeed");
//!
//! // Method 2: With a capability-aware typed wrapper
//! let reader2 = BufReader::new(Cursor::new(&data[..]));
//! let reader2 = HashReader::from_stream(reader2, size, actual_size, etag.clone(), None, diskable_md5).unwrap();
//! let reader2 = HashReader::from_stream(reader2, size, actual_size, etag.clone(), None, diskable_md5).expect("operation should succeed");
//! let wrapped_reader = EtagReader::new(HardLimitReader::new(reader2, size), etag.clone());
//! let hash_reader2 = HashReader::from_reader(wrapped_reader, size, actual_size, etag.clone(), None, diskable_md5).unwrap();
//! let hash_reader2 = HashReader::from_reader(wrapped_reader, size, actual_size, etag.clone(), None, diskable_md5).expect("operation should succeed");
//! # });
//! ```
//!
@@ -72,7 +72,7 @@
//! # tokio_test::block_on(async {
//! let data = b"test";
//! let reader = BufReader::new(Cursor::new(&data[..]));
//! let hash_reader = HashReader::from_stream(reader, 4, 4, None, None,false).unwrap();
//! let hash_reader = HashReader::from_stream(reader, 4, 4, None, None,false).expect("operation should succeed");
//!
//! // Check if a type is a HashReader
//! assert!(hash_reader.is_hash_reader());
@@ -277,7 +277,7 @@ impl HashReader {
let content_hasher = existing_hash_reader
.content_hash()
.clone()
.map(|hash| hash.checksum_type.hasher().unwrap());
.map(|hash| hash.checksum_type.hasher().expect("operation should succeed"));
let content_sha256 = existing_hash_reader.content_sha256().clone();
let content_sha256_hasher = existing_hash_reader.content_sha256().clone().map(|_| Sha256Hasher::new());
let inner = existing_hash_reader.take_inner();
@@ -664,25 +664,25 @@ mod tests {
// Test 1: Simple creation
let reader1 = BufReader::new(Cursor::new(&data[..]));
let hash_reader1 = HashReader::from_stream(reader1, size, actual_size, etag.clone(), None, false).unwrap();
let hash_reader1 = HashReader::from_stream(reader1, size, actual_size, etag.clone(), None, false).expect("operation should succeed");
assert_eq!(hash_reader1.size(), size);
assert_eq!(hash_reader1.actual_size(), actual_size);
// Test 2: With HardLimitReader wrapping
let reader2 =
HashReader::from_stream(BufReader::new(Cursor::new(&data[..])), size, actual_size, etag.clone(), None, false)
.unwrap();
.expect("operation should succeed");
let hard_limit = HardLimitReader::new(reader2, size);
let hash_reader2 = HashReader::from_reader(hard_limit, size, actual_size, etag.clone(), None, false).unwrap();
let hash_reader2 = HashReader::from_reader(hard_limit, size, actual_size, etag.clone(), None, false).expect("operation should succeed");
assert_eq!(hash_reader2.size(), size);
assert_eq!(hash_reader2.actual_size(), actual_size);
// Test 3: With EtagReader wrapping
let reader3 =
HashReader::from_stream(BufReader::new(Cursor::new(&data[..])), size, actual_size, etag.clone(), None, false)
.unwrap();
.expect("operation should succeed");
let etag_reader = EtagReader::new(reader3, etag.clone());
let hash_reader3 = HashReader::from_reader(etag_reader, size, actual_size, etag, None, false).unwrap();
let hash_reader3 = HashReader::from_reader(etag_reader, size, actual_size, etag, None, false).expect("operation should succeed");
assert_eq!(hash_reader3.size(), size);
assert_eq!(hash_reader3.actual_size(), actual_size);
}
@@ -703,7 +703,7 @@ mod tests {
None,
false,
)
.unwrap(),
.expect("operation should succeed"),
);
assert!(boxed_hash_reader.is_hash_reader());
}
@@ -719,7 +719,7 @@ mod tests {
None,
false,
)
.unwrap();
.expect("operation should succeed");
let boxed_encrypt_reader = Box::new(EncryptReader::new(inner, [7u8; 32], [3u8; 12]));
assert!(boxed_encrypt_reader.is_hash_reader());
@@ -732,9 +732,9 @@ mod tests {
None,
false,
)
.unwrap();
.expect("operation should succeed");
let mut encrypted = Vec::new();
hash_reader.read_to_end(&mut encrypted).await.unwrap();
hash_reader.read_to_end(&mut encrypted).await.expect("operation should succeed");
assert!(!encrypted.is_empty());
assert_ne!(encrypted, data);
@@ -745,9 +745,9 @@ mod tests {
async fn test_hashreader_etag_basic() {
let data = b"hello hashreader";
let reader = BufReader::new(Cursor::new(&data[..]));
let mut hash_reader = HashReader::from_stream(reader, data.len() as i64, data.len() as i64, None, None, false).unwrap();
let mut hash_reader = HashReader::from_stream(reader, data.len() as i64, data.len() as i64, None, None, false).expect("operation should succeed");
let mut buf = Vec::new();
let _ = hash_reader.read_to_end(&mut buf).await.unwrap();
let _ = hash_reader.read_to_end(&mut buf).await.expect("operation should succeed");
let etag = hash_reader.try_resolve_etag();
assert!(etag.is_some());
assert_eq!(buf, data);
@@ -757,9 +757,9 @@ mod tests {
async fn test_hashreader_diskable_md5() {
let data = b"no etag";
let reader = BufReader::new(Cursor::new(&data[..]));
let mut hash_reader = HashReader::from_stream(reader, data.len() as i64, data.len() as i64, None, None, true).unwrap();
let mut hash_reader = HashReader::from_stream(reader, data.len() as i64, data.len() as i64, None, None, true).expect("operation should succeed");
let mut buf = Vec::new();
let _ = hash_reader.read_to_end(&mut buf).await.unwrap();
let _ = hash_reader.read_to_end(&mut buf).await.expect("operation should succeed");
// Etag should be None when diskable_md5 is true
let etag = hash_reader.try_resolve_etag();
assert!(etag.is_none());
@@ -770,14 +770,14 @@ mod tests {
async fn test_add_calculated_checksum_records_checksum() {
let data = b"server-side copy checksum";
let reader = BufReader::new(Cursor::new(&data[..]));
let mut hash_reader = HashReader::from_stream(reader, data.len() as i64, data.len() as i64, None, None, false).unwrap();
let mut hash_reader = HashReader::from_stream(reader, data.len() as i64, data.len() as i64, None, None, false).expect("operation should succeed");
hash_reader.add_calculated_checksum(ChecksumType::CRC64_NVME).unwrap();
hash_reader.add_calculated_checksum(ChecksumType::CRC64_NVME).expect("operation should succeed");
let mut buf = Vec::new();
hash_reader.read_to_end(&mut buf).await.unwrap();
hash_reader.read_to_end(&mut buf).await.expect("operation should succeed");
let expected = Checksum::new_from_data(ChecksumType::CRC64_NVME, data).unwrap();
let expected = Checksum::new_from_data(ChecksumType::CRC64_NVME, data).expect("operation should succeed");
let checksums = hash_reader.content_crc();
assert_eq!(buf, data);
@@ -792,7 +792,7 @@ mod tests {
// Create a HashReader first
let hash_reader =
HashReader::from_stream(reader, data.len() as i64, data.len() as i64, Some("test_etag".to_string()), None, false)
.unwrap();
.expect("operation should succeed");
let hash_reader = wrap_reader(hash_reader);
// Now try to create another HashReader from the existing one using new
let result = HashReader::new(
@@ -805,7 +805,7 @@ mod tests {
);
assert!(result.is_ok());
let final_reader = result.unwrap();
let final_reader = result.expect("operation should succeed");
assert_eq!(final_reader.checksum, Some("test_etag".to_string()));
assert_eq!(final_reader.size(), data.len() as i64);
}
@@ -838,14 +838,14 @@ mod tests {
let size = data.len() as i64;
let actual_size = data.len() as i64;
let mut hr = HashReader::from_stream(reader, size, actual_size, Some(expected.clone()), None, false).unwrap();
let mut hr = HashReader::from_stream(reader, size, actual_size, Some(expected.clone()), None, false).expect("operation should succeed");
// If compression is enabled, compress data first
let compressed_data = if is_compress {
let mut compressed_buf = Vec::new();
let compress_reader = CompressReader::new(hr, CompressionAlgorithm::Gzip);
let mut compress_reader = compress_reader;
compress_reader.read_to_end(&mut compressed_buf).await.unwrap();
compress_reader.read_to_end(&mut compressed_buf).await.expect("operation should succeed");
println!("Original size: {}, Compressed size: {}", data.len(), compressed_buf.len());
@@ -853,7 +853,7 @@ mod tests {
} else {
// If not compressing, read original data directly
let mut buf = Vec::new();
hr.read_to_end(&mut buf).await.unwrap();
hr.read_to_end(&mut buf).await.expect("operation should succeed");
buf
};
@@ -869,7 +869,7 @@ mod tests {
let encrypt_reader = encrypt_reader::EncryptReader::new(Cursor::new(compressed_data), key, nonce);
let mut encrypted_data = Vec::new();
let mut encrypt_reader = encrypt_reader;
encrypt_reader.read_to_end(&mut encrypted_data).await.unwrap();
encrypt_reader.read_to_end(&mut encrypted_data).await.expect("operation should succeed");
println!("Encrypted size: {}", encrypted_data.len());
@@ -877,14 +877,14 @@ mod tests {
let decrypt_reader = DecryptReader::new(Cursor::new(encrypted_data), key, nonce);
let mut decrypt_reader = decrypt_reader;
let mut decrypted_data = Vec::new();
decrypt_reader.read_to_end(&mut decrypted_data).await.unwrap();
decrypt_reader.read_to_end(&mut decrypted_data).await.expect("operation should succeed");
if is_compress {
// If compression was used, decompress is needed
let decompress_reader = DecompressReader::new(Cursor::new(decrypted_data), CompressionAlgorithm::Gzip);
let mut decompress_reader = decompress_reader;
let mut final_data = Vec::new();
decompress_reader.read_to_end(&mut final_data).await.unwrap();
decompress_reader.read_to_end(&mut final_data).await.expect("operation should succeed");
println!("Final decompressed size: {}", final_data.len());
assert_eq!(final_data.len() as i64, actual_size);
@@ -902,7 +902,7 @@ mod tests {
let decompress_reader = DecompressReader::new(Cursor::new(compressed_data), CompressionAlgorithm::Gzip);
let mut decompress_reader = decompress_reader;
let mut decompressed = Vec::new();
decompress_reader.read_to_end(&mut decompressed).await.unwrap();
decompress_reader.read_to_end(&mut decompressed).await.expect("operation should succeed");
assert_eq!(decompressed.len() as i64, actual_size);
assert_eq!(&decompressed, &data);
@@ -931,13 +931,13 @@ mod tests {
println!("Original data size: {} bytes", data.len());
let reader = BufReader::new(Cursor::new(data.clone()));
let hash_reader = HashReader::from_stream(reader, data.len() as i64, data.len() as i64, None, None, false).unwrap();
let hash_reader = HashReader::from_stream(reader, data.len() as i64, data.len() as i64, None, None, false).expect("operation should succeed");
// Test compression
let compress_reader = CompressReader::new(hash_reader, CompressionAlgorithm::Gzip);
let mut compressed_data = Vec::new();
let mut compress_reader = compress_reader;
compress_reader.read_to_end(&mut compressed_data).await.unwrap();
compress_reader.read_to_end(&mut compressed_data).await.expect("operation should succeed");
println!("Compressed data size: {} bytes", compressed_data.len());
println!("Compression ratio: {:.2}%", (compressed_data.len() as f64 / data.len() as f64) * 100.0);
@@ -949,7 +949,7 @@ mod tests {
let decompress_reader = DecompressReader::new(Cursor::new(compressed_data), CompressionAlgorithm::Gzip);
let mut decompressed_data = Vec::new();
let mut decompress_reader = decompress_reader;
decompress_reader.read_to_end(&mut decompressed_data).await.unwrap();
decompress_reader.read_to_end(&mut decompressed_data).await.expect("operation should succeed");
// Verify decompressed data matches original
assert_eq!(decompressed_data.len(), data.len());
@@ -976,13 +976,13 @@ mod tests {
println!("\nTesting algorithm: {algorithm:?}");
let reader = BufReader::new(Cursor::new(data.clone()));
let hash_reader = HashReader::from_stream(reader, data.len() as i64, data.len() as i64, None, None, false).unwrap();
let hash_reader = HashReader::from_stream(reader, data.len() as i64, data.len() as i64, None, None, false).expect("operation should succeed");
// Compress
let compress_reader = CompressReader::new(hash_reader, algorithm);
let mut compressed_data = Vec::new();
let mut compress_reader = compress_reader;
compress_reader.read_to_end(&mut compressed_data).await.unwrap();
compress_reader.read_to_end(&mut compressed_data).await.expect("operation should succeed");
println!(
" Compressed size: {} bytes (ratio: {:.2}%)",
@@ -994,7 +994,7 @@ mod tests {
let decompress_reader = DecompressReader::new(Cursor::new(compressed_data), algorithm);
let mut decompressed_data = Vec::new();
let mut decompress_reader = decompress_reader;
decompress_reader.read_to_end(&mut decompressed_data).await.unwrap();
decompress_reader.read_to_end(&mut decompressed_data).await.expect("operation should succeed");
// Verify
assert_eq!(decompressed_data.len(), data.len());
+8 -8
View File
@@ -176,12 +176,12 @@ mod tests {
scan_range: None,
},
};
let db = get_global_db(input.clone(), true).await.unwrap();
let db = get_global_db(input.clone(), true).await.expect("operation should succeed");
let query = Query::new(Context { input: Arc::new(input) }, sql.to_string());
let result = db.execute(&query).await.unwrap();
let result = db.execute(&query).await.expect("operation should succeed");
let results = result.result().chunk_result().await.unwrap().to_vec();
let results = result.result().chunk_result().await.expect("operation should succeed").to_vec();
let expected = [
"+----------------+---------+-----+------------+--------+",
@@ -201,7 +201,7 @@ mod tests {
];
assert_batches_eq!(expected, &results);
pretty::print_batches(&results).unwrap();
pretty::print_batches(&results).expect("operation should succeed");
}
#[tokio::test]
@@ -235,12 +235,12 @@ mod tests {
scan_range: None,
},
};
let db = get_global_db(input.clone(), true).await.unwrap();
let db = get_global_db(input.clone(), true).await.expect("operation should succeed");
let query = Query::new(Context { input: Arc::new(input) }, sql.to_string());
let result = db.execute(&query).await.unwrap();
let result = db.execute(&query).await.expect("operation should succeed");
let results = result.result().chunk_result().await.unwrap().to_vec();
pretty::print_batches(&results).unwrap();
let results = result.result().chunk_result().await.expect("operation should succeed").to_vec();
pretty::print_batches(&results).expect("operation should succeed");
}
}
+14 -14
View File
@@ -23,7 +23,7 @@ use std::sync::LazyLock;
use thiserror::Error;
use url::Url;
static HOST_LABEL_REGEX: LazyLock<Regex> = LazyLock::new(|| Regex::new(r"^[a-zA-Z0-9]([a-zA-Z0-9-]*[a-zA-Z0-9])?$").unwrap());
static HOST_LABEL_REGEX: LazyLock<Regex> = LazyLock::new(|| Regex::new(r"^[a-zA-Z0-9]([a-zA-Z0-9-]*[a-zA-Z0-9])?$").expect("operation should succeed"));
/// NetError represents errors that can occur in network operations.
#[derive(Error, Debug)]
@@ -332,7 +332,7 @@ impl<'de> serde::Deserialize<'de> for ParsedURL {
{
let s: String = serde::Deserialize::deserialize(deserializer)?;
if s.is_empty() {
Ok(ParsedURL(Url::parse("about:blank").unwrap()))
Ok(ParsedURL(Url::parse("about:blank").expect("operation should succeed")))
} else {
parse_url(&s).map_err(serde::de::Error::custom)
}
@@ -363,7 +363,7 @@ pub fn parse_url(s: &str) -> Result<ParsedURL, NetError> {
});
if !port_str.is_empty() {
let host_port = format!("{}:{}", uu.host_str().unwrap(), port_str);
let host_port = format!("{}:{}", uu.host_str().expect("operation should succeed"), port_str);
parse_host(&host_port)?;
}
}
@@ -485,7 +485,7 @@ mod tests {
fn parse_host_with_valid_ipv4() {
let result = parse_host("192.168.1.1:8080");
assert!(result.is_ok());
let host = result.unwrap();
let host = result.expect("operation should succeed");
assert_eq!(host.name, "192.168.1.1");
assert_eq!(host.port, Some(8080));
}
@@ -494,7 +494,7 @@ mod tests {
fn parse_host_with_valid_hostname() {
let result = parse_host("example.com:443");
assert!(result.is_ok());
let host = result.unwrap();
let host = result.expect("operation should succeed");
assert_eq!(host.name, "example.com");
assert_eq!(host.port, Some(443));
}
@@ -503,7 +503,7 @@ mod tests {
fn parse_host_with_ipv6_brackets() {
let result = parse_host("[::1]:8080");
assert!(result.is_ok());
let host = result.unwrap();
let host = result.expect("operation should succeed");
assert_eq!(host.name, "::1");
assert_eq!(host.port, Some(8080));
}
@@ -512,7 +512,7 @@ mod tests {
fn parse_host_with_bare_ipv6_without_port() {
let result = parse_host("::1");
assert!(result.is_ok());
let host = result.unwrap();
let host = result.expect("operation should succeed");
assert_eq!(host.name, "::1");
assert_eq!(host.port, None);
}
@@ -521,7 +521,7 @@ mod tests {
fn parse_host_with_ipv6_zone_without_port() {
let result = parse_host("fe80::1%eth0");
assert!(result.is_ok());
let host = result.unwrap();
let host = result.expect("operation should succeed");
assert_eq!(host.name, "fe80::1%eth0");
assert_eq!(host.port, None);
}
@@ -530,7 +530,7 @@ mod tests {
fn parse_host_with_bracketed_ipv6_zone_and_port() {
let result = parse_host("[fe80::1%eth0]:9000");
assert!(result.is_ok());
let host = result.unwrap();
let host = result.expect("operation should succeed");
assert_eq!(host.name, "fe80::1%eth0");
assert_eq!(host.port, Some(9000));
}
@@ -539,7 +539,7 @@ mod tests {
fn parse_host_with_bracketed_ipv6_without_port() {
let result = parse_host("[::1]");
assert!(result.is_ok());
let host = result.unwrap();
let host = result.expect("operation should succeed");
assert_eq!(host.name, "::1");
assert_eq!(host.port, None);
}
@@ -560,7 +560,7 @@ mod tests {
fn parse_host_without_port() {
let result = parse_host("example.com");
assert!(result.is_ok());
let host = result.unwrap();
let host = result.expect("operation should succeed");
assert_eq!(host.name, "example.com");
assert_eq!(host.port, None);
}
@@ -605,7 +605,7 @@ mod tests {
fn parse_url_with_valid_http_url() {
let result = parse_url("http://example.com/path");
assert!(result.is_ok());
let parsed = result.unwrap();
let parsed = result.expect("operation should succeed");
assert_eq!(parsed.hostname(), "example.com");
assert_eq!(parsed.port(), "80");
assert_eq!(parsed.scheme(), "http");
@@ -616,7 +616,7 @@ mod tests {
fn parse_url_with_explicit_default_https_port() {
let result = parse_url("https://example.com:443/path");
assert!(result.is_ok());
let parsed = result.unwrap();
let parsed = result.expect("operation should succeed");
assert_eq!(parsed.to_string(), "https://example.com/path");
}
@@ -636,7 +636,7 @@ mod tests {
fn parse_url_normalizes_path() {
let result = parse_url("http://example.com//path/../path/");
assert!(result.is_ok());
let parsed = result.unwrap();
let parsed = result.expect("operation should succeed");
assert_eq!(parsed.to_string(), "http://example.com/path/");
}
}
+19 -19
View File
@@ -178,7 +178,7 @@ mod tests {
fn test_compress_decompress_gzip() {
let data = b"hello gzip compress";
let compressed = compress_block(data, CompressionAlgorithm::Gzip);
let decompressed = decompress_block(&compressed, CompressionAlgorithm::Gzip).unwrap();
let decompressed = decompress_block(&compressed, CompressionAlgorithm::Gzip).expect("operation should succeed");
assert_eq!(decompressed, data);
}
@@ -186,7 +186,7 @@ mod tests {
fn test_compress_decompress_deflate() {
let data = b"hello deflate compress";
let compressed = compress_block(data, CompressionAlgorithm::Deflate);
let decompressed = decompress_block(&compressed, CompressionAlgorithm::Deflate).unwrap();
let decompressed = decompress_block(&compressed, CompressionAlgorithm::Deflate).expect("operation should succeed");
assert_eq!(decompressed, data);
}
@@ -194,7 +194,7 @@ mod tests {
fn test_compress_decompress_zstd() {
let data = b"hello zstd compress";
let compressed = compress_block(data, CompressionAlgorithm::Zstd);
let decompressed = decompress_block(&compressed, CompressionAlgorithm::Zstd).unwrap();
let decompressed = decompress_block(&compressed, CompressionAlgorithm::Zstd).expect("operation should succeed");
assert_eq!(decompressed, data);
}
@@ -202,7 +202,7 @@ mod tests {
fn test_compress_decompress_lz4() {
let data = b"hello lz4 compress";
let compressed = compress_block(data, CompressionAlgorithm::Lz4);
let decompressed = decompress_block(&compressed, CompressionAlgorithm::Lz4).unwrap();
let decompressed = decompress_block(&compressed, CompressionAlgorithm::Lz4).expect("operation should succeed");
assert_eq!(decompressed, data);
}
@@ -210,7 +210,7 @@ mod tests {
fn test_compress_decompress_brotli() {
let data = b"hello brotli compress";
let compressed = compress_block(data, CompressionAlgorithm::Brotli);
let decompressed = decompress_block(&compressed, CompressionAlgorithm::Brotli).unwrap();
let decompressed = decompress_block(&compressed, CompressionAlgorithm::Brotli).expect("operation should succeed");
assert_eq!(decompressed, data);
}
@@ -218,18 +218,18 @@ mod tests {
fn test_compress_decompress_snappy() {
let data = b"hello snappy compress";
let compressed = compress_block(data, CompressionAlgorithm::Snappy);
let decompressed = decompress_block(&compressed, CompressionAlgorithm::Snappy).unwrap();
let decompressed = decompress_block(&compressed, CompressionAlgorithm::Snappy).expect("operation should succeed");
assert_eq!(decompressed, data);
}
#[test]
fn test_from_str() {
assert_eq!(CompressionAlgorithm::from_str("gzip").unwrap(), CompressionAlgorithm::Gzip);
assert_eq!(CompressionAlgorithm::from_str("deflate").unwrap(), CompressionAlgorithm::Deflate);
assert_eq!(CompressionAlgorithm::from_str("zstd").unwrap(), CompressionAlgorithm::Zstd);
assert_eq!(CompressionAlgorithm::from_str("lz4").unwrap(), CompressionAlgorithm::Lz4);
assert_eq!(CompressionAlgorithm::from_str("brotli").unwrap(), CompressionAlgorithm::Brotli);
assert_eq!(CompressionAlgorithm::from_str("snappy").unwrap(), CompressionAlgorithm::Snappy);
assert_eq!(CompressionAlgorithm::from_str("gzip").expect("operation should succeed"), CompressionAlgorithm::Gzip);
assert_eq!(CompressionAlgorithm::from_str("deflate").expect("operation should succeed"), CompressionAlgorithm::Deflate);
assert_eq!(CompressionAlgorithm::from_str("zstd").expect("operation should succeed"), CompressionAlgorithm::Zstd);
assert_eq!(CompressionAlgorithm::from_str("lz4").expect("operation should succeed"), CompressionAlgorithm::Lz4);
assert_eq!(CompressionAlgorithm::from_str("brotli").expect("operation should succeed"), CompressionAlgorithm::Brotli);
assert_eq!(CompressionAlgorithm::from_str("snappy").expect("operation should succeed"), CompressionAlgorithm::Snappy);
assert!(CompressionAlgorithm::from_str("unknown").is_err());
}
@@ -278,12 +278,12 @@ mod tests {
println!("{name}: {size} bytes, {dur:?}");
}
// All should decompress to the original
assert_eq!(decompress_block(&gzip, CompressionAlgorithm::Gzip).unwrap(), data);
assert_eq!(decompress_block(&deflate, CompressionAlgorithm::Deflate).unwrap(), data);
assert_eq!(decompress_block(&zstd, CompressionAlgorithm::Zstd).unwrap(), data);
assert_eq!(decompress_block(&lz4, CompressionAlgorithm::Lz4).unwrap(), data);
assert_eq!(decompress_block(&brotli, CompressionAlgorithm::Brotli).unwrap(), data);
assert_eq!(decompress_block(&snappy, CompressionAlgorithm::Snappy).unwrap(), data);
assert_eq!(decompress_block(&gzip, CompressionAlgorithm::Gzip).expect("operation should succeed"), data);
assert_eq!(decompress_block(&deflate, CompressionAlgorithm::Deflate).expect("operation should succeed"), data);
assert_eq!(decompress_block(&zstd, CompressionAlgorithm::Zstd).expect("operation should succeed"), data);
assert_eq!(decompress_block(&lz4, CompressionAlgorithm::Lz4).expect("operation should succeed"), data);
assert_eq!(decompress_block(&brotli, CompressionAlgorithm::Brotli).expect("operation should succeed"), data);
assert_eq!(decompress_block(&snappy, CompressionAlgorithm::Snappy).expect("operation should succeed"), data);
// All compressed results should not be empty
assert!(
!gzip.is_empty()
@@ -326,7 +326,7 @@ mod tests {
// Decompression test
let start = Instant::now();
let _decompressed = decompress_block(&compressed, algo).unwrap();
let _decompressed = decompress_block(&compressed, algo).expect("operation should succeed");
let _decompression_time = start.elapsed();
// Calculate compression ratio
+4 -4
View File
@@ -84,7 +84,7 @@ pub fn is_sha256_checksum(s: &str) -> bool {
/// A 20-byte array containing the HMAC-SHA1 hash of the input data using the provided key
///
pub fn hmac_sha1(key: impl AsRef<[u8]>, data: impl AsRef<[u8]>) -> [u8; 20] {
let mut m = <Hmac<Sha1>>::new_from_slice(key.as_ref()).unwrap();
let mut m = <Hmac<Sha1>>::new_from_slice(key.as_ref()).expect("operation should succeed");
m.update(data.as_ref());
m.finalize().into_bytes().into()
}
@@ -100,7 +100,7 @@ pub fn hmac_sha1(key: impl AsRef<[u8]>, data: impl AsRef<[u8]>) -> [u8; 20] {
/// A 32-byte array containing the HMAC-SHA256 hash of the input data using the provided key
///
pub fn hmac_sha256(key: impl AsRef<[u8]>, data: impl AsRef<[u8]>) -> [u8; 32] {
let mut m = Hmac::<Sha256>::new_from_slice(key.as_ref()).unwrap();
let mut m = Hmac::<Sha256>::new_from_slice(key.as_ref()).expect("operation should succeed");
m.update(data.as_ref());
m.finalize().into_bytes().into()
}
@@ -149,8 +149,8 @@ fn test_base64_encoding_decoding() {
println!("Encoded: {}", &encoded_string);
let decoded_bytes = base64_decode_url_safe_no_pad(encoded_string.as_bytes()).unwrap();
let decoded_string = String::from_utf8(decoded_bytes).unwrap();
let decoded_bytes = base64_decode_url_safe_no_pad(encoded_string.as_bytes()).expect("operation should succeed");
let decoded_string = String::from_utf8(decoded_bytes).expect("operation should succeed");
assert_eq!(decoded_string, original_uuid_timestamp)
}
+2 -2
View File
@@ -30,8 +30,8 @@ pub const X_REAL_IP: &str = "x-real-ip";
/// e.g. Forwarded: for=192.0.2.60;proto=https;by=203.0.113.43
const FORWARDED: &str = "forwarded";
static FOR_REGEX: LazyLock<Regex> = LazyLock::new(|| Regex::new(r"(?i)(?:for=)([^(;|,| )]+)(.*)").unwrap());
static PROTO_REGEX: LazyLock<Regex> = LazyLock::new(|| Regex::new(r"(?i)^(;|,| )+(?:proto=)(https|http)").unwrap());
static FOR_REGEX: LazyLock<Regex> = LazyLock::new(|| Regex::new(r"(?i)(?:for=)([^(;|,| )]+)(.*)").expect("operation should succeed"));
static PROTO_REGEX: LazyLock<Regex> = LazyLock::new(|| Regex::new(r"(?i)^(;|,| )+(?:proto=)(https|http)").expect("operation should succeed"));
/// Used to disable all processing of the X-Forwarded-For header in source IP discovery.
///
+2 -2
View File
@@ -84,7 +84,7 @@ impl Stream for RetryTimer {
self.timer = Some(timer);
}
let mut timer = self.timer.as_mut().unwrap();
let mut timer = self.timer.as_mut().expect("operation should succeed");
match Pin::new(&mut timer).poll_tick(cx) {
Poll::Ready(_) => {
self.rem -= 1;
@@ -151,7 +151,7 @@ pub fn is_request_error_retryable(_err: std::io::Error) -> bool {
}
let uerr = err.(*url.Error);
if uerr.is_ok() {
let e = uerr.unwrap();
let e = uerr.expect("operation should succeed");
return match e.type {
x509.UnknownAuthorityError => {
false
+5 -5
View File
@@ -32,10 +32,10 @@ use std::sync::LazyLock;
/// let false_values = ["0", "f", "F", "false", "FALSE", "False", "off", "OFF", "Off", "disabled"];
///
/// for val in true_values.iter() {
/// assert_eq!(parse_bool(val).unwrap(), true);
/// assert_eq!(parse_bool(val).expect("operation should succeed"), true);
/// }
/// for val in false_values.iter() {
/// assert_eq!(parse_bool(val).unwrap(), false);
/// assert_eq!(parse_bool(val).expect("operation should succeed"), false);
/// }
/// ```
///
@@ -236,7 +236,7 @@ pub fn match_as_pattern_prefix(pattern: &str, text: &str) -> bool {
text.len() <= pattern.len()
}
static ELLIPSES_RE: LazyLock<Regex> = LazyLock::new(|| Regex::new(r"(.*)(\{[0-9A-Fa-f]*\.\.\.[0-9A-Fa-f]*\})(.*)").unwrap());
static ELLIPSES_RE: LazyLock<Regex> = LazyLock::new(|| Regex::new(r"(.*)(\{[0-9A-Fa-f]*\.\.\.[0-9A-Fa-f]*\})(.*)").expect("operation should succeed"));
/// Ellipses constants
const OPEN_BRACES: &str = "{";
@@ -351,7 +351,7 @@ impl ArgPattern {
/// use rustfs_utils::string::find_ellipses_patterns;
///
/// let pattern = "http://rustfs{2...3}/export/set{1...64}";
/// let arg_pattern = find_ellipses_patterns(pattern).unwrap();
/// let arg_pattern = find_ellipses_patterns(pattern).expect("operation should succeed");
/// assert_eq!(arg_pattern.total_sizes(), 128);
/// ```
pub fn find_ellipses_patterns(arg: &str) -> Result<ArgPattern> {
@@ -469,7 +469,7 @@ pub fn has_ellipses<T: AsRef<str>>(s: &[T]) -> bool {
/// ```no_run
/// use rustfs_utils::string::parse_ellipses_range;
///
/// let range = parse_ellipses_range("{1...5}").unwrap();
/// let range = parse_ellipses_range("{1...5}").expect("operation should succeed");
/// assert_eq!(range, vec!["1", "2", "3", "4", "5"]);
/// ```
///
+27 -27
View File
@@ -246,7 +246,7 @@ impl Operation for CreateKeyHandler {
.map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers))
}
@@ -309,7 +309,7 @@ impl Operation for DescribeKeyHandler {
.map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers))
}
@@ -474,7 +474,7 @@ impl Operation for ListKeysHandler {
.map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers))
}
@@ -548,7 +548,7 @@ impl Operation for GenerateDataKeyHandler {
.map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers))
}
@@ -617,7 +617,7 @@ impl Operation for CreateKmsKeyHandler {
let data =
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
return Ok(S3Response::with_headers((StatusCode::SERVICE_UNAVAILABLE, Body::from(data)), headers));
};
@@ -631,7 +631,7 @@ impl Operation for CreateKmsKeyHandler {
let data =
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
return Ok(S3Response::with_headers((StatusCode::SERVICE_UNAVAILABLE, Body::from(data)), headers));
};
@@ -670,7 +670,7 @@ impl Operation for CreateKmsKeyHandler {
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers))
}
@@ -695,7 +695,7 @@ impl Operation for CreateKmsKeyHandler {
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
Ok(S3Response::with_headers((StatusCode::INTERNAL_SERVER_ERROR, Body::from(data)), headers))
}
@@ -760,7 +760,7 @@ impl Operation for DeleteKmsKeyHandler {
let data =
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
return Ok(S3Response::with_headers((StatusCode::BAD_REQUEST, Body::from(data)), headers));
};
@@ -787,7 +787,7 @@ impl Operation for DeleteKmsKeyHandler {
let data =
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
return Ok(S3Response::with_headers((StatusCode::SERVICE_UNAVAILABLE, Body::from(data)), headers));
};
@@ -801,7 +801,7 @@ impl Operation for DeleteKmsKeyHandler {
let data =
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
return Ok(S3Response::with_headers((StatusCode::SERVICE_UNAVAILABLE, Body::from(data)), headers));
};
@@ -833,7 +833,7 @@ impl Operation for DeleteKmsKeyHandler {
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers))
}
@@ -864,7 +864,7 @@ impl Operation for DeleteKmsKeyHandler {
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
Ok(S3Response::with_headers((status, Body::from(data)), headers))
}
@@ -927,7 +927,7 @@ impl Operation for CancelKmsKeyDeletionHandler {
let data =
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
return Ok(S3Response::with_headers((StatusCode::BAD_REQUEST, Body::from(data)), headers));
};
CancelKmsKeyDeletionRequest { key_id: key_id.clone() }
@@ -945,7 +945,7 @@ impl Operation for CancelKmsKeyDeletionHandler {
let data =
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
return Ok(S3Response::with_headers((StatusCode::SERVICE_UNAVAILABLE, Body::from(data)), headers));
};
@@ -959,7 +959,7 @@ impl Operation for CancelKmsKeyDeletionHandler {
let data =
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
return Ok(S3Response::with_headers((StatusCode::SERVICE_UNAVAILABLE, Body::from(data)), headers));
};
@@ -989,7 +989,7 @@ impl Operation for CancelKmsKeyDeletionHandler {
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers))
}
@@ -1015,7 +1015,7 @@ impl Operation for CancelKmsKeyDeletionHandler {
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
Ok(S3Response::with_headers((StatusCode::INTERNAL_SERVER_ERROR, Body::from(data)), headers))
}
@@ -1070,7 +1070,7 @@ impl Operation for ListKmsKeysHandler {
let data =
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
return Ok(S3Response::with_headers((StatusCode::SERVICE_UNAVAILABLE, Body::from(data)), headers));
};
@@ -1085,7 +1085,7 @@ impl Operation for ListKmsKeysHandler {
let data =
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
return Ok(S3Response::with_headers((StatusCode::SERVICE_UNAVAILABLE, Body::from(data)), headers));
};
@@ -1119,7 +1119,7 @@ impl Operation for ListKmsKeysHandler {
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers))
}
@@ -1145,7 +1145,7 @@ impl Operation for ListKmsKeysHandler {
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
Ok(S3Response::with_headers((StatusCode::INTERNAL_SERVER_ERROR, Body::from(data)), headers))
}
@@ -1192,7 +1192,7 @@ impl Operation for DescribeKmsKeyHandler {
let data =
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
return Ok(S3Response::with_headers((StatusCode::BAD_REQUEST, Body::from(data)), headers));
};
@@ -1205,7 +1205,7 @@ impl Operation for DescribeKmsKeyHandler {
let data =
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
return Ok(S3Response::with_headers((StatusCode::SERVICE_UNAVAILABLE, Body::from(data)), headers));
};
@@ -1218,7 +1218,7 @@ impl Operation for DescribeKmsKeyHandler {
let data =
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
return Ok(S3Response::with_headers((StatusCode::SERVICE_UNAVAILABLE, Body::from(data)), headers));
};
@@ -1247,7 +1247,7 @@ impl Operation for DescribeKmsKeyHandler {
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers))
}
@@ -1278,7 +1278,7 @@ impl Operation for DescribeKmsKeyHandler {
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
Ok(S3Response::with_headers((status, Body::from(data)), headers))
}
+3 -3
View File
@@ -208,7 +208,7 @@ impl Operation for KmsStatusHandler {
let data = serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers))
}
@@ -257,7 +257,7 @@ impl Operation for KmsConfigHandler {
let data = serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers))
}
@@ -302,7 +302,7 @@ impl Operation for KmsClearCacheHandler {
serde_json::to_vec(&response).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, "application/json".parse().unwrap());
headers.insert(CONTENT_TYPE, "application/json".parse().expect("operation should succeed"));
Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers))
}
+2 -2
View File
@@ -70,7 +70,7 @@ impl Operation for TriggerProfileCPU {
match crate::profiling::dump_cpu_pprof_for(dur).await {
Ok(path) => {
let mut header = HeaderMap::new();
header.insert(CONTENT_TYPE, "text/html".parse().unwrap());
header.insert(CONTENT_TYPE, "text/html".parse().expect("operation should succeed"));
Ok(S3Response::with_headers((StatusCode::OK, Body::from(path.display().to_string())), header))
}
Err(e) => Err(s3s::s3_error!(InternalError, "{}", format!("Failed to dump CPU profile: {e}"))),
@@ -99,7 +99,7 @@ impl Operation for TriggerProfileMemory {
match crate::profiling::dump_memory_pprof_now().await {
Ok(path) => {
let mut header = HeaderMap::new();
header.insert(CONTENT_TYPE, "text/html".parse().unwrap());
header.insert(CONTENT_TYPE, "text/html".parse().expect("operation should succeed"));
Ok(S3Response::with_headers((StatusCode::OK, Body::from(path.display().to_string())), header))
}
Err(e) => Err(s3s::s3_error!(InternalError, "{}", format!("Failed to dump Memory profile: {e}"))),
+2 -2
View File
@@ -1995,7 +1995,7 @@ async fn put_replication_probe_object(
.customize()
.map_request(move |mut req| {
for (key, value) in headers.clone() {
req.headers_mut().insert(key.unwrap(), value);
req.headers_mut().insert(key.expect("operation should succeed"), value);
}
Result::<_, std::io::Error>::Ok(req)
})
@@ -2039,7 +2039,7 @@ async fn delete_replication_probe_object(
.customize()
.map_request(move |mut req| {
for (key, value) in headers.clone() {
req.headers_mut().insert(key.unwrap(), value);
req.headers_mut().insert(key.expect("operation should succeed"), value);
}
Result::<_, std::io::Error>::Ok(req)
})
+1 -1
View File
@@ -89,7 +89,7 @@ fn compute_default_max_blocking_threads() -> usize {
/// ```no_run
/// // tokio_runtime_builder is pub(crate) - call it from within the rustfs binary:
/// // let builder = tokio_runtime_builder();
/// // let runtime = builder.build().unwrap();
/// // let runtime = builder.build().expect("operation should succeed");
/// ```
pub fn tokio_runtime_builder() -> tokio::runtime::Builder {
let mut builder = tokio::runtime::Builder::new_multi_thread();