auto heal(1)

Signed-off-by: root <root@DESKTOP-QLJNS6S.localdomain>
This commit is contained in:
root
2024-11-16 20:43:34 +08:00
parent 7eed1a6db3
commit 01e0e4b673
20 changed files with 988 additions and 1791 deletions
+8 -2
View File
@@ -65,7 +65,10 @@ impl LastMinuteLatency {
}
pub fn add(&mut self, t: &Duration) {
let sec = SystemTime::now().duration_since(UNIX_EPOCH).expect("Time went backwards").as_secs();
let sec = SystemTime::now()
.duration_since(UNIX_EPOCH)
.expect("Time went backwards")
.as_secs();
self.forward_to(sec);
let win_idx = sec % 60;
self.totals[win_idx as usize].add(t);
@@ -81,7 +84,10 @@ impl LastMinuteLatency {
pub fn get_total(&mut self) -> AccElem {
let mut res = AccElem::default();
let sec = SystemTime::now().duration_since(UNIX_EPOCH).expect("Time went backwards").as_secs();
let sec = SystemTime::now()
.duration_since(UNIX_EPOCH)
.expect("Time went backwards")
.as_secs();
self.forward_to(sec);
for elem in self.totals.iter() {
res.merge(elem);
@@ -1,10 +1,9 @@
// automatically generated by the FlatBuffers compiler, do not modify
// @generated
use core::mem;
use core::cmp::Ordering;
use core::mem;
extern crate flatbuffers;
use self::flatbuffers::{EndianScalar, Follow};
@@ -12,112 +11,114 @@ use self::flatbuffers::{EndianScalar, Follow};
#[allow(unused_imports, dead_code)]
pub mod models {
use core::mem;
use core::cmp::Ordering;
use core::cmp::Ordering;
use core::mem;
extern crate flatbuffers;
use self::flatbuffers::{EndianScalar, Follow};
extern crate flatbuffers;
use self::flatbuffers::{EndianScalar, Follow};
pub enum PingBodyOffset {}
#[derive(Copy, Clone, PartialEq)]
pub enum PingBodyOffset {}
#[derive(Copy, Clone, PartialEq)]
pub struct PingBody<'a> {
pub _tab: flatbuffers::Table<'a>,
}
impl<'a> flatbuffers::Follow<'a> for PingBody<'a> {
type Inner = PingBody<'a>;
#[inline]
unsafe fn follow(buf: &'a [u8], loc: usize) -> Self::Inner {
Self { _tab: flatbuffers::Table::new(buf, loc) }
}
}
impl<'a> PingBody<'a> {
pub const VT_PAYLOAD: flatbuffers::VOffsetT = 4;
pub const fn get_fully_qualified_name() -> &'static str {
"models.PingBody"
}
#[inline]
pub unsafe fn init_from_table(table: flatbuffers::Table<'a>) -> Self {
PingBody { _tab: table }
}
#[allow(unused_mut)]
pub fn create<'bldr: 'args, 'args: 'mut_bldr, 'mut_bldr, A: flatbuffers::Allocator + 'bldr>(
_fbb: &'mut_bldr mut flatbuffers::FlatBufferBuilder<'bldr, A>,
args: &'args PingBodyArgs<'args>
) -> flatbuffers::WIPOffset<PingBody<'bldr>> {
let mut builder = PingBodyBuilder::new(_fbb);
if let Some(x) = args.payload { builder.add_payload(x); }
builder.finish()
}
#[inline]
pub fn payload(&self) -> Option<flatbuffers::Vector<'a, u8>> {
// Safety:
// Created from valid Table for this object
// which contains a valid value in this slot
unsafe { self._tab.get::<flatbuffers::ForwardsUOffset<flatbuffers::Vector<'a, u8>>>(PingBody::VT_PAYLOAD, None)}
}
}
impl flatbuffers::Verifiable for PingBody<'_> {
#[inline]
fn run_verifier(
v: &mut flatbuffers::Verifier, pos: usize
) -> Result<(), flatbuffers::InvalidFlatbuffer> {
use self::flatbuffers::Verifiable;
v.visit_table(pos)?
.visit_field::<flatbuffers::ForwardsUOffset<flatbuffers::Vector<'_, u8>>>("payload", Self::VT_PAYLOAD, false)?
.finish();
Ok(())
}
}
pub struct PingBodyArgs<'a> {
pub payload: Option<flatbuffers::WIPOffset<flatbuffers::Vector<'a, u8>>>,
}
impl<'a> Default for PingBodyArgs<'a> {
#[inline]
fn default() -> Self {
PingBodyArgs {
payload: None,
pub struct PingBody<'a> {
pub _tab: flatbuffers::Table<'a>,
}
}
}
pub struct PingBodyBuilder<'a: 'b, 'b, A: flatbuffers::Allocator + 'a> {
fbb_: &'b mut flatbuffers::FlatBufferBuilder<'a, A>,
start_: flatbuffers::WIPOffset<flatbuffers::TableUnfinishedWIPOffset>,
}
impl<'a: 'b, 'b, A: flatbuffers::Allocator + 'a> PingBodyBuilder<'a, 'b, A> {
#[inline]
pub fn add_payload(&mut self, payload: flatbuffers::WIPOffset<flatbuffers::Vector<'b , u8>>) {
self.fbb_.push_slot_always::<flatbuffers::WIPOffset<_>>(PingBody::VT_PAYLOAD, payload);
}
#[inline]
pub fn new(_fbb: &'b mut flatbuffers::FlatBufferBuilder<'a, A>) -> PingBodyBuilder<'a, 'b, A> {
let start = _fbb.start_table();
PingBodyBuilder {
fbb_: _fbb,
start_: start,
impl<'a> flatbuffers::Follow<'a> for PingBody<'a> {
type Inner = PingBody<'a>;
#[inline]
unsafe fn follow(buf: &'a [u8], loc: usize) -> Self::Inner {
Self {
_tab: flatbuffers::Table::new(buf, loc),
}
}
}
}
#[inline]
pub fn finish(self) -> flatbuffers::WIPOffset<PingBody<'a>> {
let o = self.fbb_.end_table(self.start_);
flatbuffers::WIPOffset::new(o.value())
}
}
impl core::fmt::Debug for PingBody<'_> {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
let mut ds = f.debug_struct("PingBody");
ds.field("payload", &self.payload());
ds.finish()
}
}
} // pub mod models
impl<'a> PingBody<'a> {
pub const VT_PAYLOAD: flatbuffers::VOffsetT = 4;
pub const fn get_fully_qualified_name() -> &'static str {
"models.PingBody"
}
#[inline]
pub unsafe fn init_from_table(table: flatbuffers::Table<'a>) -> Self {
PingBody { _tab: table }
}
#[allow(unused_mut)]
pub fn create<'bldr: 'args, 'args: 'mut_bldr, 'mut_bldr, A: flatbuffers::Allocator + 'bldr>(
_fbb: &'mut_bldr mut flatbuffers::FlatBufferBuilder<'bldr, A>,
args: &'args PingBodyArgs<'args>,
) -> flatbuffers::WIPOffset<PingBody<'bldr>> {
let mut builder = PingBodyBuilder::new(_fbb);
if let Some(x) = args.payload {
builder.add_payload(x);
}
builder.finish()
}
#[inline]
pub fn payload(&self) -> Option<flatbuffers::Vector<'a, u8>> {
// Safety:
// Created from valid Table for this object
// which contains a valid value in this slot
unsafe {
self._tab
.get::<flatbuffers::ForwardsUOffset<flatbuffers::Vector<'a, u8>>>(PingBody::VT_PAYLOAD, None)
}
}
}
impl flatbuffers::Verifiable for PingBody<'_> {
#[inline]
fn run_verifier(v: &mut flatbuffers::Verifier, pos: usize) -> Result<(), flatbuffers::InvalidFlatbuffer> {
use self::flatbuffers::Verifiable;
v.visit_table(pos)?
.visit_field::<flatbuffers::ForwardsUOffset<flatbuffers::Vector<'_, u8>>>("payload", Self::VT_PAYLOAD, false)?
.finish();
Ok(())
}
}
pub struct PingBodyArgs<'a> {
pub payload: Option<flatbuffers::WIPOffset<flatbuffers::Vector<'a, u8>>>,
}
impl<'a> Default for PingBodyArgs<'a> {
#[inline]
fn default() -> Self {
PingBodyArgs { payload: None }
}
}
pub struct PingBodyBuilder<'a: 'b, 'b, A: flatbuffers::Allocator + 'a> {
fbb_: &'b mut flatbuffers::FlatBufferBuilder<'a, A>,
start_: flatbuffers::WIPOffset<flatbuffers::TableUnfinishedWIPOffset>,
}
impl<'a: 'b, 'b, A: flatbuffers::Allocator + 'a> PingBodyBuilder<'a, 'b, A> {
#[inline]
pub fn add_payload(&mut self, payload: flatbuffers::WIPOffset<flatbuffers::Vector<'b, u8>>) {
self.fbb_
.push_slot_always::<flatbuffers::WIPOffset<_>>(PingBody::VT_PAYLOAD, payload);
}
#[inline]
pub fn new(_fbb: &'b mut flatbuffers::FlatBufferBuilder<'a, A>) -> PingBodyBuilder<'a, 'b, A> {
let start = _fbb.start_table();
PingBodyBuilder {
fbb_: _fbb,
start_: start,
}
}
#[inline]
pub fn finish(self) -> flatbuffers::WIPOffset<PingBody<'a>> {
let o = self.fbb_.end_table(self.start_);
flatbuffers::WIPOffset::new(o.value())
}
}
impl core::fmt::Debug for PingBody<'_> {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
let mut ds = f.debug_struct("PingBody");
ds.field("payload", &self.payload());
ds.finish()
}
}
} // pub mod models
File diff suppressed because it is too large Load Diff
+14
View File
@@ -382,6 +382,19 @@ message DiskInfoResponse {
optional string error_info = 3;
}
message NsScannerRequest {
string disk = 1;
string cache = 2;
uint64 scan_mode = 3;
}
message NsScannerResponse {
bool success = 1;
string update = 2;
string data_usage_cache = 3;
optional string error_info = 4;
}
// lock api have same argument type
message GenerallyLockRequest {
string args = 1;
@@ -432,6 +445,7 @@ service NodeService {
rpc ReadMultiple(ReadMultipleRequest) returns (ReadMultipleResponse) {};
rpc DeleteVolume(DeleteVolumeRequest) returns (DeleteVolumeResponse) {};
rpc DiskInfo(DiskInfoRequest) returns (DiskInfoResponse) {};
rpc NsScanner(stream NsScannerRequest) returns (stream NsScannerResponse) {};
/* -------------------------------lock service-------------------------- */
-1
View File
@@ -9,7 +9,6 @@ use s3s::dto::{
BucketLifecycleConfiguration, NotificationConfiguration, ObjectLockConfiguration, ReplicationConfiguration,
ServerSideEncryptionConfiguration, Tagging, VersioningConfiguration,
};
use s3s::xml;
use serde::Serializer;
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
+5 -5
View File
@@ -14,8 +14,8 @@ use crate::{
error::{Error, Result},
};
type AgreedFn = Box<dyn Fn(MetaCacheEntry) -> Pin<Box<dyn Future<Output = ()>>> + Send + 'static>;
type PartialFn = Box<dyn Fn(MetaCacheEntries, &[Option<Error>]) -> Pin<Box<dyn Future<Output = ()>>> + Send + 'static>;
type AgreedFn = Box<dyn Fn(MetaCacheEntry) -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + 'static>;
type PartialFn = Box<dyn Fn(MetaCacheEntries, &[Option<Error>]) -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + 'static>;
type FinishedFn = Box<dyn Fn(&[Option<Error>]) -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + 'static>;
#[derive(Default)]
@@ -213,7 +213,7 @@ pub async fn list_path_raw(mut rx: B_Receiver<bool>, opts: ListPathRawOptions) -
if at_eof + has_err == readers.len() {
if has_err > 0 {
if let Some(finished_fn) = opts.finished.as_ref() {
finished_fn(&errs);
finished_fn(&errs).await;
}
break;
}
@@ -221,13 +221,13 @@ pub async fn list_path_raw(mut rx: B_Receiver<bool>, opts: ListPathRawOptions) -
if agree == readers.len() {
if let Some(agreed_fn) = opts.agreed.as_ref() {
agreed_fn(current);
agreed_fn(current).await;
}
continue;
}
if let Some(partial_fn) = opts.partial.as_ref() {
partial_fn(MetaCacheEntries(top_entries), &errs);
partial_fn(MetaCacheEntries(top_entries), &errs).await;
}
}
+8 -8
View File
@@ -20,8 +20,8 @@ use crate::disk::{LocalFileReader, LocalFileWriter, STORAGE_FORMAT_FILE};
use crate::error::{Error, Result};
use crate::global::{GLOBAL_IsErasureSD, GLOBAL_RootDiskThreshold};
use crate::heal::data_scanner::{has_active_rules, scan_data_folder, ScannerItem, SizeSummary};
use crate::heal::data_scanner_metric::{globalScannerMetrics, ScannerMetric, ScannerMetrics};
use crate::heal::data_usage_cache::{self, DataUsageCache, DataUsageEntry};
use crate::heal::data_scanner_metric::{ScannerMetric, ScannerMetrics};
use crate::heal::data_usage_cache::{DataUsageCache, DataUsageEntry};
use crate::heal::error::{ERR_IGNORE_FILE_CONTRIB, ERR_SKIP_FILE};
use crate::heal::heal_commands::HealScanMode;
use crate::new_object_layer_fn;
@@ -40,7 +40,6 @@ use crate::{
};
use common::defer;
use path_absolutize::Absolutize;
use s3s::dto::{ReplicationConfiguration, ReplicationRuleStatus};
use std::collections::{HashMap, HashSet};
use std::fmt::Debug;
use std::io::Cursor;
@@ -1880,7 +1879,7 @@ impl DiskAPI for LocalDisk {
}
async fn ns_scanner(
self: Arc<Self>,
&self,
cache: &DataUsageCache,
updates: Sender<DataUsageEntry>,
scan_mode: HealScanMode,
@@ -1908,7 +1907,7 @@ impl DiskAPI for LocalDisk {
};
let loc = self.get_disk_location();
let disks = store.get_disks(loc.pool_idx.unwrap(), loc.disk_idx.unwrap()).await?;
let disk = self.clone();
let disk = Arc::new(LocalDisk::new(&self.endpoint(), false).await?);
let disk_clone = disk.clone();
let mut cache = cache.clone();
cache.info.updates = Some(updates.clone());
@@ -1925,7 +1924,7 @@ impl DiskAPI for LocalDisk {
}
let stop_fn = ScannerMetrics::log(ScannerMetric::ScanObject);
let mut res = HashMap::new();
let done_sz = ScannerMetrics::time_size(ScannerMetric::ReadMetadata);
let done_sz = ScannerMetrics::time_size(ScannerMetric::ReadMetadata).await;
let buf = match disk.read_metadata(item.path.clone()).await {
Ok(buf) => buf,
Err(err) => {
@@ -1965,7 +1964,7 @@ impl DiskAPI for LocalDisk {
let mut obj_deleted = false;
for info in obj_infos.iter() {
let done = ScannerMetrics::time(ScannerMetric::ApplyVersion);
let mut sz = 0;
let sz: usize;
(obj_deleted, sz) = item.apply_actions(&info, &size_s).await;
done().await;
@@ -1994,7 +1993,8 @@ impl DiskAPI for LocalDisk {
}
for frer_version in fivs.free_versions.iter() {
let _obj_info = frer_version.to_object_info(&item.bucket, &item.object_path().to_string_lossy(), versioned);
let _obj_info =
frer_version.to_object_info(&item.bucket, &item.object_path().to_string_lossy(), versioned);
let done = ScannerMetrics::time(ScannerMetric::TierObjSweep);
done().await;
}
+8 -8
View File
@@ -1,7 +1,7 @@
pub mod endpoint;
pub mod error;
pub mod format;
mod local;
pub mod local;
pub mod os;
pub mod remote;
@@ -16,7 +16,7 @@ const STORAGE_FORMAT_FILE: &str = "xl.meta";
use crate::{
erasure::Writer,
error::{Error, Result},
file_meta::{merge_file_meta_versions, FileMeta, FileMetaShallowVersion, FileMetaVersion},
file_meta::{merge_file_meta_versions, FileMeta, FileMetaShallowVersion},
heal::{
data_usage_cache::{DataUsageCache, DataUsageEntry},
heal_commands::HealScanMode,
@@ -344,14 +344,14 @@ impl DiskAPI for Disk {
}
async fn ns_scanner(
self: Arc<Self>,
&self,
cache: &DataUsageCache,
updates: Sender<DataUsageEntry>,
scan_mode: HealScanMode,
) -> Result<DataUsageCache> {
match &*self {
Disk::Local(local_disk) => Arc::new(local_disk).ns_scanner(cache, updates, scan_mode).await,
Disk::Remote(remote_disk) => Arc::new(remote_disk).ns_scanner(cache, updates, scan_mode).await,
Disk::Local(local_disk) => local_disk.ns_scanner(cache, updates, scan_mode).await,
Disk::Remote(remote_disk) => remote_disk.ns_scanner(cache, updates, scan_mode).await,
}
}
}
@@ -453,7 +453,7 @@ pub trait DiskAPI: Debug + Send + Sync + 'static {
async fn read_all(&self, volume: &str, path: &str) -> Result<Vec<u8>>;
async fn disk_info(&self, opts: &DiskInfoOptions) -> Result<DiskInfo>;
async fn ns_scanner(
self: Arc<Self>,
&self,
cache: &DataUsageCache,
updates: Sender<DataUsageEntry>,
scan_mode: HealScanMode,
@@ -1303,10 +1303,10 @@ impl Reader for RemoteFileReader {
Err(Error::from_string(error_info))
}
}
async fn seek(&mut self, offset: usize) -> Result<()> {
async fn seek(&mut self, _offset: usize) -> Result<()> {
unimplemented!()
}
async fn read_exact(&mut self, buf: &mut [u8]) -> Result<usize> {
async fn read_exact(&mut self, _buf: &mut [u8]) -> Result<usize> {
unimplemented!()
}
}
+39 -7
View File
@@ -1,16 +1,17 @@
use std::{path::PathBuf, sync::Arc};
use std::path::PathBuf;
use futures::lock::Mutex;
use protos::{
node_service_time_out_client,
proto_gen::node_service::{
CheckPartsRequest, DeletePathsRequest, DeleteRequest, DeleteVersionRequest, DeleteVersionsRequest, DeleteVolumeRequest,
DiskInfoRequest, ListDirRequest, ListVolumesRequest, MakeVolumeRequest, MakeVolumesRequest, ReadAllRequest,
ReadMultipleRequest, ReadVersionRequest, ReadXlRequest, RenameDataRequest, RenameFileRequst, StatVolumeRequest,
UpdateMetadataRequest, VerifyFileRequest, WalkDirRequest, WriteAllRequest, WriteMetadataRequest,
DiskInfoRequest, ListDirRequest, ListVolumesRequest, MakeVolumeRequest, MakeVolumesRequest, NsScannerRequest,
ReadAllRequest, ReadMultipleRequest, ReadVersionRequest, ReadXlRequest, RenameDataRequest, RenameFileRequst,
StatVolumeRequest, UpdateMetadataRequest, VerifyFileRequest, WalkDirRequest, WriteAllRequest, WriteMetadataRequest,
},
};
use tokio::sync::mpsc::Sender;
use tokio::sync::mpsc::{self, Sender};
use tokio_stream::{wrappers::ReceiverStream, StreamExt};
use tonic::Request;
use tracing::info;
use uuid::Uuid;
@@ -736,11 +737,42 @@ impl DiskAPI for RemoteDisk {
}
async fn ns_scanner(
self: Arc<Self>,
&self,
cache: &DataUsageCache,
updates: Sender<DataUsageEntry>,
scan_mode: HealScanMode,
) -> Result<DataUsageCache> {
todo!()
info!("ns_scanner");
let cache = serde_json::to_string(cache)?;
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?;
let (tx, rx) = mpsc::channel(10);
let in_stream = ReceiverStream::new(rx);
let mut response = client.ns_scanner(in_stream).await?.into_inner();
let request = NsScannerRequest {
disk: self.root.to_string_lossy().to_string(),
cache,
scan_mode: scan_mode as u64,
};
tx.send(request).await?;
loop {
match response.next().await {
Some(Ok(resp)) => {
if !resp.update.is_empty() {
let data_usage_cache = serde_json::from_str::<DataUsageEntry>(&resp.update)?;
let _ = updates.send(data_usage_cache).await;
} else if !resp.data_usage_cache.is_empty() {
let data_usage_cache = serde_json::from_str::<DataUsageCache>(&resp.data_usage_cache)?;
return Ok(data_usage_cache);
} else {
return Err(Error::from_string("scan was interrupted"));
}
}
_ => return Err(Error::from_string("scan was interrupted")),
}
}
}
}
+1 -3
View File
@@ -424,7 +424,7 @@ impl Erasure {
writers: &mut [Option<BitrotWriter>],
readers: Vec<Option<BitrotReader>>,
total_length: usize,
prefer: &[bool],
_prefer: &[bool],
) -> Result<()> {
if writers.len() != self.parity_shards + self.data_shards {
return Err(Error::from_string("invalid argument"));
@@ -437,8 +437,6 @@ impl Erasure {
end_block += 1;
}
let mut bytes_writed = 0;
let mut errs = Vec::new();
for _ in start_block..=end_block {
let mut bufs = reader.read().await?;
+92 -64
View File
@@ -1,17 +1,14 @@
use std::{env, sync::Arc};
use tokio::{
select,
sync::{
broadcast::Receiver as B_Receiver,
mpsc::{self, Receiver, Sender},
RwLock,
},
use tokio::sync::{
mpsc::{self, Receiver, Sender},
RwLock,
};
use crate::{
disk::error::DiskError,
endpoints::Endpoints,
error::{Error, Result},
global::{GLOBAL_BackgroundHealRoutine, GLOBAL_BackgroundHealState, GLOBAL_LOCAL_DISK_MAP},
heal::heal_ops::NOP_HEAL,
new_object_layer_fn,
store_api::StorageAPI,
@@ -20,9 +17,38 @@ use crate::{
use super::{
heal_commands::{HealOpts, HealResultItem},
heal_ops::HealSequence,
heal_ops::{new_bg_heal_sequence, HealSequence},
};
pub async fn init_auto_heal() {
init_background_healing().await;
if let Ok(v) = env::var("_RUSTFS_AUTO_DRIVE_HEALING") {
if v == "on" {
// GLOBAL_BackgroundHealState.write().await.push_heal_local_disks(heal_local_disks).await;
}
}
}
async fn init_background_healing() {
let bg_seq = Arc::new(RwLock::new(new_bg_heal_sequence()));
for _ in 0..GLOBAL_BackgroundHealRoutine.read().await.workers {
let bg_seq_clone = bg_seq.clone();
tokio::spawn(async {
GLOBAL_BackgroundHealRoutine.write().await.add_worker(bg_seq_clone).await;
});
}
let _ = GLOBAL_BackgroundHealState
.write()
.await
.launch_new_heal_sequence(bg_seq)
.await;
}
async fn get_local_disks_to_heal() -> Endpoints {
for (_, disk) in GLOBAL_LOCAL_DISK_MAP.read().await.iter() {}
todo!()
}
#[derive(Clone, Debug)]
pub struct HealTask {
pub bucket: String,
@@ -49,7 +75,7 @@ impl HealTask {
pub struct HealResult {
pub result: HealResultItem,
err: Option<Error>,
_err: Option<Error>,
}
pub struct HealRoutine {
@@ -79,66 +105,68 @@ impl HealRoutine {
}))
}
pub async fn add_worker(&mut self, mut ctx: B_Receiver<bool>, bgseq: &mut HealSequence) {
pub async fn add_worker(&mut self, bgseq: Arc<RwLock<HealSequence>>) {
loop {
select! {
task = self.tasks_rx.recv() => {
let mut d_res = HealResultItem::default();
let d_err: Option<Error>;
match task {
Some(task) => {
if task.bucket == NOP_HEAL {
d_err = Some(Error::from_string("skip file"));
} else if task.bucket == SLASH_SEPARATOR {
match heal_disk_format(task.opts).await {
Ok((res, err)) => {
d_res = res;
d_err = err;
},
Err(err) => {d_err = Some(err)},
}
} else {
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = lock
.as_ref()
.expect("Not init");
if task.object.is_empty() {
match store.heal_object(&task.bucket, &task.object, &task.version_id, &task.opts).await {
Ok((res, err)) => {
d_res = res;
d_err = err;
},
Err(err) => {d_err = Some(err)},
}
} else {
match store.heal_object(&task.bucket, &task.object, &task.version_id, &task.opts).await {
Ok((res, err)) => {
d_res = res;
d_err = err;
},
Err(err) => {d_err = Some(err)},
}
}
let mut d_res = HealResultItem::default();
let d_err: Option<Error>;
match self.tasks_rx.recv().await {
Some(task) => {
if task.bucket == NOP_HEAL {
d_err = Some(Error::from_string("skip file"));
} else if task.bucket == SLASH_SEPARATOR {
match heal_disk_format(task.opts).await {
Ok((res, err)) => {
d_res = res;
d_err = err;
}
if let Some(resp_tx) = task.resp_tx {
let _ = resp_tx.send(HealResult{result: d_res, err: d_err}).await;
} else {
// when respCh is not set caller is not waiting but we
// update the relevant metrics for them
if d_err.is_none() {
bgseq.count_healed(d_res.heal_item_type);
} else {
bgseq.count_failed(d_res.heal_item_type);
Err(err) => d_err = Some(err),
}
} else {
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = lock.as_ref().expect("Not init");
if task.object.is_empty() {
match store
.heal_object(&task.bucket, &task.object, &task.version_id, &task.opts)
.await
{
Ok((res, err)) => {
d_res = res;
d_err = err;
}
Err(err) => d_err = Some(err),
}
},
None => return,
} else {
match store
.heal_object(&task.bucket, &task.object, &task.version_id, &task.opts)
.await
{
Ok((res, err)) => {
d_res = res;
d_err = err;
}
Err(err) => d_err = Some(err),
}
}
}
if let Some(resp_tx) = task.resp_tx {
let _ = resp_tx
.send(HealResult {
result: d_res,
_err: d_err,
})
.await;
} else {
// when respCh is not set caller is not waiting but we
// update the relevant metrics for them
if d_err.is_none() {
bgseq.write().await.count_healed(d_res.heal_item_type);
} else {
bgseq.write().await.count_failed(d_res.heal_item_type);
}
}
}
_ = ctx.recv() => {
return;
}
None => return,
}
}
}
+21 -14
View File
@@ -34,7 +34,6 @@ use super::{
data_usage_cache::{DataUsageCache, DataUsageEntry, DataUsageHash},
heal_commands::{HealScanMode, HEAL_DEEP_SCAN, HEAL_NORMAL_SCAN},
};
use crate::heal::data_scanner_metric::current_path_updater;
use crate::heal::data_usage::DATA_USAGE_ROOT;
use crate::{
cache_value::metacache_set::{list_path_raw, ListPathRawOptions},
@@ -54,15 +53,16 @@ use crate::{
},
new_object_layer_fn,
peer::is_reserved_or_invalid_bucket,
store::{ECStore, ListPathOptions},
store::ECStore,
utils::path::{path_join, path_to_bucket_object, path_to_bucket_object_with_base_path, SLASH_SEPARATOR},
};
use crate::{disk::local::LocalDisk, heal::data_scanner_metric::current_path_updater};
use crate::{
disk::DiskAPI,
store_api::{FileInfo, ObjectInfo},
};
const DATA_SCANNER_SLEEP_PER_FOLDER: Duration = Duration::from_millis(1); // Time to wait between folders.
const _DATA_SCANNER_SLEEP_PER_FOLDER: Duration = Duration::from_millis(1); // Time to wait between folders.
const DATA_USAGE_UPDATE_DIR_CYCLES: u32 = 16; // Visit all folders every n cycles.
const DATA_SCANNER_COMPACT_LEAST_OBJECT: u64 = 500; // Compact when there are less than this many objects in a branch.
const DATA_SCANNER_COMPACT_AT_CHILDREN: u64 = 10000; // Compact when there are this many children in a branch.
@@ -70,12 +70,12 @@ const DATA_SCANNER_COMPACT_AT_FOLDERS: u64 = DATA_SCANNER_COMPACT_AT_CHILDREN /
pub const DATA_SCANNER_FORCE_COMPACT_AT_FOLDERS: u64 = 250_000; // Compact when this many subfolders in a single folder (even top level).
const DATA_SCANNER_START_DELAY: Duration = Duration::from_secs(60); // Time to wait on startup and between cycles.
const HEAL_DELETE_DANGLING: bool = true;
pub const HEAL_DELETE_DANGLING: bool = true;
const HEAL_OBJECT_SELECT_PROB: u64 = 1024; // Overall probability of a file being scanned; one in n.
// static SCANNER_SLEEPER: () = new_dynamic_sleeper(2, Duration::from_secs(1), true); // Keep defaults same as config defaults
static SCANNER_CYCLE: AtomicU64 = AtomicU64::new(DATA_SCANNER_START_DELAY.as_secs());
static SCANNER_IDLE_MODE: AtomicU32 = AtomicU32::new(0); // default is throttled when idle
static _SCANNER_IDLE_MODE: AtomicU32 = AtomicU32::new(0); // default is throttled when idle
static SCANNER_EXCESS_OBJECT_VERSIONS: AtomicU64 = AtomicU64::new(100);
static SCANNER_EXCESS_OBJECT_VERSIONS_TOTAL_SIZE: AtomicU64 = AtomicU64::new(1024 * 1024 * 1024 * 1024); // 1 TB
static SCANNER_EXCESS_FOLDERS: AtomicU64 = AtomicU64::new(50_000);
@@ -167,16 +167,18 @@ async fn run_data_scanner() {
cycle_info.current = 0;
cycle_info.cycle_completed.push(SystemTime::now());
if cycle_info.cycle_completed.len() > DATA_USAGE_UPDATE_DIR_CYCLES as usize {
cycle_info.cycle_completed = cycle_info.cycle_completed[cycle_info.cycle_completed.len() - DATA_USAGE_UPDATE_DIR_CYCLES as usize..].to_vec();
cycle_info.cycle_completed = cycle_info.cycle_completed
[cycle_info.cycle_completed.len() - DATA_USAGE_UPDATE_DIR_CYCLES as usize..]
.to_vec();
}
globalScannerMetrics.write().await.set_cycle(Some(cycle_info.clone())).await;
let mut tmp = Vec::new();
tmp.write_u64::<LittleEndian>(cycle_info.next).unwrap();
let _ = save_config(store, &DATA_USAGE_BLOOM_NAME_PATH, &tmp).await;
},
}
Err(err) => {
res.insert("error".to_string(), err.to_string());
},
}
}
stop_fn(&res).await;
sleep(Duration::from_secs(SCANNER_CYCLE.load(std::sync::atomic::Ordering::SeqCst))).await;
@@ -265,7 +267,7 @@ impl Default for CurrentScannerCycle {
}
impl CurrentScannerCycle {
pub fn marshal_msg(&self, buf: &[u8]) -> Result<Vec<u8>> {
pub fn marshal_msg(&self, next_buf: &[u8]) -> Result<Vec<u8>> {
let len: u32 = 4;
let mut wr = Vec::new();
@@ -291,7 +293,7 @@ impl CurrentScannerCycle {
.serialize(&mut Serializer::new(&mut buf))
.expect("Serialization failed");
rmp::encode::write_bin(&mut wr, &buf)?;
let mut result = buf.to_vec();
let mut result = next_buf.to_vec();
result.extend(wr.iter());
Ok(result)
}
@@ -355,7 +357,7 @@ fn timestamp_to_system_time(timestamp: u64) -> SystemTime {
}
#[derive(Clone, Debug, Default)]
struct Heal {
pub struct Heal {
enabled: bool,
bitrot: bool,
}
@@ -403,7 +405,12 @@ impl ScannerItem {
cumulative_size += obj_info.size;
}
if cumulative_size >= SCANNER_EXCESS_OBJECT_VERSIONS_TOTAL_SIZE.load(Ordering::SeqCst).try_into().unwrap() {
if cumulative_size
>= SCANNER_EXCESS_OBJECT_VERSIONS_TOTAL_SIZE
.load(Ordering::SeqCst)
.try_into()
.unwrap()
{
//todo
}
@@ -421,7 +428,7 @@ impl ScannerItem {
Ok(object_infos)
}
pub async fn apply_actions(&self, oi: &ObjectInfo, size_s: &SizeSummary) -> (bool, usize) {
pub async fn apply_actions(&self, _oi: &ObjectInfo, _size_s: &SizeSummary) -> (bool, usize) {
let done = ScannerMetrics::time(ScannerMetric::Ilm);
//todo: lifecycle
done().await;
@@ -1014,7 +1021,7 @@ pub fn has_active_rules(config: &ReplicationConfiguration, prefix: &str, recursi
false
}
pub type LocalDrive = Arc<dyn DiskAPI>;
pub type LocalDrive = Arc<LocalDisk>;
pub async fn scan_data_folder(
disks: &[Option<DiskStore>],
drive: LocalDrive,
+11 -11
View File
@@ -79,8 +79,8 @@ impl Clone for LockedLastMinuteLatency {
}
impl LockedLastMinuteLatency {
pub fn add(&mut self, value: &Duration) {
self.add_size(value, 0);
pub async fn add(&mut self, value: &Duration) {
self.add_size(value, 0).await;
}
pub async fn add_size(&mut self, value: &Duration, sz: u64) {
@@ -105,16 +105,16 @@ impl LockedLastMinuteLatency {
a.size = old.size;
a.total = old.total;
a.n = old.n;
self.mu.write().await;
let _ = self.mu.write().await;
self.latency.add_all(t - 1, &a);
}
self.cached.n += 1;
self.cached.total += value.as_secs();
self.cached.size != sz;
self.cached.size += sz;
}
pub async fn total(&mut self) -> AccElem {
self.mu.read().await;
let _ = self.mu.read().await;
self.latency.get_total()
}
}
@@ -147,19 +147,19 @@ impl ScannerMetrics {
pub fn log(s: ScannerMetric) -> LogFn {
let start = SystemTime::now();
let s_clone = s as usize;
Arc::new(move |custom: &HashMap<String, String>| {
Arc::new(move |_custom: &HashMap<String, String>| {
Box::pin(async move {
let duration = SystemTime::now().duration_since(start).unwrap_or(Duration::from_secs(0));
let mut sm_w = globalScannerMetrics.write().await;
sm_w.operations[s_clone].fetch_add(1, Ordering::SeqCst);
if s_clone < ScannerMetric::LastRealtime as usize {
sm_w.latency[s_clone].add(&duration);
sm_w.latency[s_clone].add(&duration).await;
}
})
})
}
pub fn time_size(s: ScannerMetric) -> TimeSizeFn {
pub async fn time_size(s: ScannerMetric) -> TimeSizeFn {
let start = SystemTime::now();
let s_clone = s as usize;
Arc::new(move |sz: u64| {
@@ -168,7 +168,7 @@ impl ScannerMetrics {
let mut sm_w = globalScannerMetrics.write().await;
sm_w.operations[s_clone].fetch_add(1, Ordering::SeqCst);
if s_clone < ScannerMetric::LastRealtime as usize {
sm_w.latency[s_clone].add_size(&duration, sz);
sm_w.latency[s_clone].add_size(&duration, sz).await;
}
})
})
@@ -183,7 +183,7 @@ impl ScannerMetrics {
let mut sm_w = globalScannerMetrics.write().await;
sm_w.operations[s_clone].fetch_add(1, Ordering::SeqCst);
if s_clone < ScannerMetric::LastRealtime as usize {
sm_w.latency[s_clone].add(&duration);
sm_w.latency[s_clone].add(&duration).await;
}
})
})
@@ -191,7 +191,7 @@ impl ScannerMetrics {
}
pub type CloseDiskFn = Arc<dyn Fn() -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + Sync + 'static>;
pub fn current_path_updater(disk: &str, initial: &str) -> (UpdateCurrentPathFn, CloseDiskFn) {
pub fn current_path_updater(disk: &str, _initial: &str) -> (UpdateCurrentPathFn, CloseDiskFn) {
let disk_1 = disk.to_string();
let disk_2 = disk.to_string();
(
+12 -11
View File
@@ -5,7 +5,6 @@ use crate::error::{Error, Result};
use crate::new_object_layer_fn;
use crate::set_disk::SetDisks;
use crate::store_api::{BucketInfo, HTTPRangeSpec, ObjectIO, ObjectOptions};
use bytes::Bytes;
use bytesize::ByteSize;
use http::HeaderMap;
use path_clean::PathClean;
@@ -17,7 +16,6 @@ use serde::{Deserialize, Serialize};
use std::collections::{HashMap, HashSet};
use std::hash::{DefaultHasher, Hash, Hasher};
use std::path::Path;
use std::sync::Arc;
use std::time::{Duration, SystemTime};
use std::u64;
use tokio::sync::mpsc::Sender;
@@ -761,15 +759,18 @@ impl DataUsageCache {
bui.replica_count = rs.replica_count;
for (arn, stat) in rs.targets.iter() {
bui.replication_info.insert(arn.clone(), BucketTargetUsageInfo {
replication_pending_size: stat.pending_size,
replicated_size: stat.replicated_size,
replication_failed_size: stat.failed_size,
replication_pending_count: stat.pending_count,
replication_failed_count: stat.failed_count,
replicated_count: stat.replicated_count,
..Default::default()
});
bui.replication_info.insert(
arn.clone(),
BucketTargetUsageInfo {
replication_pending_size: stat.pending_size,
replicated_size: stat.replicated_size,
replication_failed_size: stat.failed_size,
replication_pending_count: stat.pending_count,
replication_failed_count: stat.failed_count,
replicated_count: stat.replicated_count,
..Default::default()
},
);
}
}
dst.insert(bucket.name.clone(), bui);
+12 -2
View File
@@ -75,11 +75,21 @@ pub struct HealResultItem {
pub object_size: usize,
}
#[derive(Debug, Default, Serialize, Deserialize)]
#[derive(Debug, Serialize, Deserialize)]
pub struct HealStartSuccess {
pub client_token: String,
pub client_address: String,
pub start_time: u64,
pub start_time: SystemTime,
}
impl Default for HealStartSuccess {
fn default() -> Self {
Self {
client_token: Default::default(),
client_address: Default::default(),
start_time: SystemTime::now(),
}
}
}
pub type HealStopSuccess = HealStartSuccess;
+69 -32
View File
@@ -3,7 +3,7 @@ use crate::{
endpoints::Endpoints,
error::{Error, Result},
global::GLOBAL_IsDistErasure,
heal::heal_commands::HEAL_UNKNOWN_SCAN,
heal::heal_commands::{HealStartSuccess, HEAL_UNKNOWN_SCAN},
utils::path::has_profix,
};
use lazy_static::lazy_static;
@@ -28,6 +28,7 @@ use uuid::Uuid;
use super::{
background_heal_ops::HealTask,
data_scanner::HEAL_DELETE_DANGLING,
heal_commands::{HealItemType, HealOpts, HealResultItem, HealScanMode, HealStopSuccess, HealingDisk, HealingTracker},
};
@@ -45,6 +46,10 @@ const HEAL_RUNNING_STATUS: &str = "running";
const HEAL_STOPPED_STATUS: &str = "stopped";
const HEAL_FINISHED_STATUS: &str = "finished";
pub const RUESTFS_RESERVED_BUCKET: &str = "rustfs";
pub const RUESTFS_RESERVED_BUCKET_PATH: &str = "/rustfs";
pub const LOGIN_PATH_PREFIX: &str = "/login";
const MAX_UNCONSUMED_HEAL_RESULT_ITEMS: usize = 1000;
const HEAL_UNCONSUMED_TIMEOUT: std::time::Duration = Duration::from_secs(24 * 60 * 60);
pub const NOP_HEAL: &str = "";
@@ -74,8 +79,8 @@ pub struct HealSequence {
pub bucket: String,
pub object: String,
pub report_progress: bool,
pub start_time: u64,
pub end_time: Arc<RwLock<u64>>,
pub start_time: SystemTime,
pub end_time: Arc<RwLock<SystemTime>>,
pub client_token: String,
pub client_address: String,
pub force_started: bool,
@@ -94,6 +99,30 @@ pub struct HealSequence {
rx: Arc<RwLock<Receiver<bool>>>,
}
pub fn new_bg_heal_sequence() -> HealSequence {
let hs = HealOpts {
remove: HEAL_DELETE_DANGLING,
..Default::default()
};
HealSequence {
start_time: SystemTime::now(),
client_token: BG_HEALING_UUID.to_string(),
bucket: RUESTFS_RESERVED_BUCKET.to_string(),
setting: hs,
current_status: Arc::new(RwLock::new(HealSequenceStatus {
summary: HEAL_NOT_STARTED_STATUS.to_string(),
heal_setting: hs,
..Default::default()
})),
report_progress: false,
scanned_items_map: HashMap::new(),
healed_items_map: HashMap::new(),
heal_failed_items_map: HashMap::new(),
..Default::default()
}
}
impl Default for HealSequence {
fn default() -> Self {
let (h_tx, h_rx) = mpsc::channel(1);
@@ -102,8 +131,8 @@ impl Default for HealSequence {
bucket: Default::default(),
object: Default::default(),
report_progress: Default::default(),
start_time: Default::default(),
end_time: Default::default(),
start_time: SystemTime::now(),
end_time: Arc::new(RwLock::new(SystemTime::now())),
client_token: Default::default(),
client_address: Default::default(),
force_started: Default::default(),
@@ -289,9 +318,10 @@ impl HealSequence {
}
}
pub async fn heal_sequence_start(h: Arc<HealSequence>) {
pub async fn heal_sequence_start(h: Arc<RwLock<HealSequence>>) {
let r = h.read().await;
{
let mut current_status_w = h.current_status.write().await;
let mut current_status_w = r.current_status.write().await;
(*current_status_w).summary = HEAL_RUNNING_STATUS.to_string();
(*current_status_w).start_time = SystemTime::now()
.duration_since(UNIX_EPOCH)
@@ -301,22 +331,20 @@ pub async fn heal_sequence_start(h: Arc<HealSequence>) {
let h_clone = h.clone();
spawn(async move {
h_clone.traverse_and_heal().await;
h_clone.read().await.traverse_and_heal().await;
});
let h_clone_1 = h.clone();
let mut x = h.traverse_and_heal_done_rx.write().await;
let mut x = r.traverse_and_heal_done_rx.write().await;
select! {
_ = h.is_done() => {
*(h.end_time.write().await) = SystemTime::now()
.duration_since(UNIX_EPOCH)
.expect("Time went backwards")
.as_secs();
let mut current_status_w = h.current_status.write().await;
_ = r.is_done() => {
*(r.end_time.write().await) = SystemTime::now();
let mut current_status_w = r.current_status.write().await;
(*current_status_w).summary = HEAL_FINISHED_STATUS.to_string();
spawn(async move {
let mut rx_w = h_clone_1.traverse_and_heal_done_rx.write().await;
let binding = h_clone_1.read().await;
let mut rx_w = binding.traverse_and_heal_done_rx.write().await;
rx_w.recv().await;
});
}
@@ -325,12 +353,12 @@ pub async fn heal_sequence_start(h: Arc<HealSequence>) {
Some(err) => {
match err {
Some(err) => {
let mut current_status_w = h.current_status.write().await;
let mut current_status_w = r.current_status.write().await;
(current_status_w).summary = HEAL_STOPPED_STATUS.to_string();
(current_status_w).failure_detail = err.to_string();
},
None => {
let mut current_status_w = h.current_status.write().await;
let mut current_status_w = r.current_status.write().await;
(current_status_w).summary = HEAL_FINISHED_STATUS.to_string();
}
}
@@ -429,13 +457,13 @@ impl AllHealState {
Endpoints::from(endpoints)
}
async fn set_disk_healing_status(&mut self, ep: Endpoint, healing: bool) {
pub async fn set_disk_healing_status(&mut self, ep: Endpoint, healing: bool) {
let _ = self.mu.write().await;
self.heal_local_disks.insert(ep, healing);
}
async fn push_heal_local_disks(&mut self, heal_local_disks: &[Endpoint]) {
pub async fn push_heal_local_disks(&mut self, heal_local_disks: &[Endpoint]) {
let _ = self.mu.write().await;
heal_local_disks.iter().for_each(|heal_local_disk| {
@@ -450,9 +478,7 @@ impl AllHealState {
let mut keys_to_reomve = Vec::new();
for (k, v) in self.heal_seq_map.iter() {
let r = v.read().await;
if r.has_ended().await
&& (UNIX_EPOCH + Duration::from_secs(*(r.end_time.read().await)) + KEEP_HEAL_SEQ_STATE_DURATION) < now
{
if r.has_ended().await && now.duration_since(*(r.end_time.read().await)).unwrap() > KEEP_HEAL_SEQ_STATE_DURATION {
keys_to_reomve.push(k.clone())
}
}
@@ -523,15 +549,16 @@ impl AllHealState {
// `keepHealSeqStateDuration`. This function also launches a
// background routine to clean up heal results after the
// aforementioned duration.
pub async fn launch_new_heal_sequence(&mut self, heal_sequence: &HealSequence) -> Result<Vec<u8>> {
let path = Path::new(&heal_sequence.bucket).join(heal_sequence.object.clone());
pub async fn launch_new_heal_sequence(&mut self, heal_sequence: Arc<RwLock<HealSequence>>) -> Result<Vec<u8>> {
let r = heal_sequence.read().await;
let path = Path::new(&r.bucket).join(r.object.clone());
let path_s = path.to_str().unwrap();
if heal_sequence.force_started {
if r.force_started {
self.stop_heal_sequence(path_s).await?;
} else {
if let Some(hs) = self.get_heal_sequence(path_s).await {
if !hs.read().await.has_ended().await {
return Err(Error::from_string(format!("Heal is already running on the given path (use force-start option to stop and start afresh). The heal was started by IP {} at {}, token is {}", heal_sequence.client_address, heal_sequence.start_time, heal_sequence.client_token)));
return Err(Error::from_string(format!("Heal is already running on the given path (use force-start option to stop and start afresh). The heal was started by IP {} at {:?}, token is {}", r.client_address, r.start_time, r.client_token)));
}
}
}
@@ -547,17 +574,27 @@ impl AllHealState {
}
}
self.heal_seq_map
.insert(path_s.to_string(), Arc::new(RwLock::new(heal_sequence.clone())));
self.heal_seq_map.insert(path_s.to_string(), heal_sequence.clone());
let client_token = heal_sequence.client_token.clone();
let client_token = r.client_token.clone();
if *GLOBAL_IsDistErasure.read().await {
// TODO: proxy
}
if heal_sequence.client_token == BG_HEALING_UUID {
if r.client_token == BG_HEALING_UUID {
// For background heal do nothing, do not spawn an unnecessary goroutine.
} else {
let heal_sequence_clone = heal_sequence.clone();
tokio::spawn(async {
heal_sequence_start(heal_sequence_clone).await;
});
}
todo!()
let b = serde_json::to_vec(&HealStartSuccess {
client_token,
client_address: r.client_address.clone(),
start_time: r.start_time,
})?;
Ok(b)
}
}
+9 -10
View File
@@ -483,7 +483,7 @@ impl StorageAPI for Sets {
.complete_multipart_upload(bucket, object, upload_id, uploaded_parts, opts)
.await
}
async fn get_disks(&self, pool_idx: usize, set_idx: usize) -> Result<Vec<Option<DiskStore>>> {
async fn get_disks(&self, _pool_idx: usize, _set_idx: usize) -> Result<Vec<Option<DiskStore>>> {
unimplemented!()
}
@@ -535,8 +535,7 @@ impl StorageAPI for Sets {
}
let format_op_id = Uuid::new_v4().to_string();
let (new_format_sets, current_disks_info) =
new_heal_format_sets(&ref_format, self.set_count, self.set_drive_count, &formats, &errs);
let (new_format_sets, _) = new_heal_format_sets(&ref_format, self.set_count, self.set_drive_count, &formats, &errs);
if !dry_run {
let mut tmp_new_formats = vec![None; self.set_count * self.set_drive_count];
for (i, set) in new_format_sets.iter().enumerate() {
@@ -578,7 +577,7 @@ impl StorageAPI for Sets {
}
Ok((res, None))
}
async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result<HealResultItem> {
async fn heal_bucket(&self, _bucket: &str, _opts: &HealOpts) -> Result<HealResultItem> {
unimplemented!()
}
async fn heal_object(
@@ -592,18 +591,18 @@ impl StorageAPI for Sets {
.heal_object(bucket, object, version_id, opts)
.await
}
async fn heal_objects(&self, bucket: &str, prefix: &str, opts: &HealOpts, func: HealObjectFn) -> Result<()> {
async fn heal_objects(&self, _bucket: &str, _prefix: &str, _opts: &HealOpts, _func: HealObjectFn) -> Result<()> {
unimplemented!()
}
async fn get_pool_and_set(&self, id: &str) -> Result<(Option<usize>, Option<usize>, Option<usize>)> {
async fn get_pool_and_set(&self, _id: &str) -> Result<(Option<usize>, Option<usize>, Option<usize>)> {
unimplemented!()
}
async fn check_abandoned_parts(&self, bucket: &str, object: &str, opts: &HealOpts) -> Result<()> {
async fn check_abandoned_parts(&self, _bucket: &str, _object: &str, _opts: &HealOpts) -> Result<()> {
unimplemented!()
}
}
async fn close_storage_disks(disks: &[Option<DiskStore>]) {
async fn _close_storage_disks(disks: &[Option<DiskStore>]) {
let mut futures = Vec::with_capacity(disks.len());
for disk in disks.iter() {
if let Some(disk) = disk {
@@ -652,7 +651,7 @@ fn formats_to_drives_info(endpoints: &Endpoints, formats: &[Option<FormatV3>], e
let mut before_drives = Vec::with_capacity(endpoints.as_ref().len());
for (index, format) in formats.iter().enumerate() {
let drive = endpoints.get_string(index);
let mut state = if format.is_some() {
let state = if format.is_some() {
DRIVE_STATE_OK
} else {
if let Some(Some(err)) = errs.get(index) {
@@ -689,7 +688,7 @@ fn new_heal_format_sets(
let mut new_formats = vec![vec![None; set_drive_count]; set_count];
let mut current_disks_info = vec![vec![DiskInfo::default(); set_drive_count]; set_count];
for (i, set) in ref_format.erasure.sets.iter().enumerate() {
for (j, value) in set.iter().enumerate() {
for j in 0..set.len() {
if let Some(Some(err)) = errs.get(i * set_drive_count + j) {
match err.downcast_ref::<DiskError>() {
Some(DiskError::UnformattedDisk) => {
+18 -10
View File
@@ -1,6 +1,5 @@
#![allow(clippy::map_entry)]
use crate::bucket::metadata;
use crate::bucket::metadata_sys::{self, init_bucket_metadata_sys, set_bucket_metadata};
use crate::bucket::utils::{check_valid_bucket_name, check_valid_bucket_name_strict, is_meta_bucketname};
use crate::config::{self, storageclass, GLOBAL_ConfigSys};
@@ -44,7 +43,6 @@ use http::HeaderMap;
use lazy_static::lazy_static;
use rand::Rng;
use s3s::dto::{BucketVersioningStatus, ObjectLockConfiguration, ObjectLockEnabled, VersioningConfiguration};
use tokio::time::interval;
use std::cmp::Ordering;
use std::slice::Iter;
use std::time::SystemTime;
@@ -54,12 +52,13 @@ use std::{
time::Duration,
};
use time::OffsetDateTime;
use tokio::{fs, select};
use tokio::sync::mpsc::Sender;
use tokio::sync::{broadcast, mpsc, RwLock, Semaphore};
use tokio::time::interval;
use tokio::{fs, select};
use crate::heal::data_usage_cache::{DataUsageCache, DataUsageCacheInfo};
use tracing::{debug, info, warn};
use tracing::{debug, info};
use uuid::Uuid;
const MAX_UPLOADS_LIST: usize = 10000;
@@ -510,7 +509,7 @@ impl ECStore {
self.pools.iter().for_each(|pool| {
total_results += pool.disk_set.len();
});
let mut results = Arc::new(RwLock::new(vec![DataUsageCache::default(); total_results]));
let results = Arc::new(RwLock::new(vec![DataUsageCache::default(); total_results]));
let (cancel, _) = broadcast::channel(100);
let first_err = Arc::new(RwLock::new(None));
let mut futures = Vec::new();
@@ -535,13 +534,16 @@ impl ECStore {
}
}
});
if let Err(err) = set.ns_scanner(&all_buckets_clone, want_cycle.try_into().unwrap(), tx, heal_scan_mode).await {
if let Err(err) = set
.ns_scanner(&all_buckets_clone, want_cycle.try_into().unwrap(), tx, heal_scan_mode)
.await
{
let mut f_w = first_err_clone.write().await;
if f_w.is_none() {
*f_w = Some(err);
}
let _ = cancel_clone.send(true);
return ;
return;
}
let _ = task.await;
});
@@ -583,7 +585,7 @@ impl ECStore {
}
_ = ctx_closer.recv() => {
}
}
let _ = task.await;
@@ -688,7 +690,13 @@ impl ECStore {
}
}
async fn update_scan(all_merged: Arc<RwLock<DataUsageCache>>, results: Arc<RwLock<Vec<DataUsageCache>>>, last_update: Option<SystemTime>, all_buckets: Vec<BucketInfo>, updates: Sender<DataUsageInfo>) -> Option<SystemTime> {
async fn update_scan(
all_merged: Arc<RwLock<DataUsageCache>>,
results: Arc<RwLock<Vec<DataUsageCache>>>,
last_update: Option<SystemTime>,
all_buckets: Vec<BucketInfo>,
updates: Sender<DataUsageInfo>,
) -> Option<SystemTime> {
let mut w = all_merged.write().await;
*w = DataUsageCache {
info: DataUsageCacheInfo {
@@ -1527,7 +1535,7 @@ impl StorageAPI for ECStore {
Ok((r, None))
}
async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result<HealResultItem> {
async fn heal_bucket(&self, _bucket: &str, _opts: &HealOpts) -> Result<HealResultItem> {
unimplemented!()
}
async fn heal_object(
+97 -6
View File
@@ -6,6 +6,7 @@ use ecstore::{
UpdateMetadataOpts, WalkDirOptions,
},
erasure::Writer,
heal::data_usage_cache::DataUsageCache,
peer::{LocalPeerS3Client, PeerS3Client},
store::{all_local_disk_path, find_local_disk},
store_api::{BucketOptions, DeleteBucketOptions, FileInfo, MakeBucketOptions},
@@ -22,12 +23,12 @@ use protos::{
DiskInfoRequest, DiskInfoResponse, GenerallyLockRequest, GenerallyLockResponse, GetBucketInfoRequest,
GetBucketInfoResponse, ListBucketRequest, ListBucketResponse, ListDirRequest, ListDirResponse, ListVolumesRequest,
ListVolumesResponse, MakeBucketRequest, MakeBucketResponse, MakeVolumeRequest, MakeVolumeResponse, MakeVolumesRequest,
MakeVolumesResponse, PingRequest, PingResponse, ReadAllRequest, ReadAllResponse, ReadAtRequest, ReadAtResponse,
ReadMultipleRequest, ReadMultipleResponse, ReadVersionRequest, ReadVersionResponse, ReadXlRequest, ReadXlResponse,
RenameDataRequest, RenameDataResponse, RenameFileRequst, RenameFileResponse, RenamePartRequst, RenamePartResponse,
StatVolumeRequest, StatVolumeResponse, UpdateMetadataRequest, UpdateMetadataResponse, VerifyFileRequest,
VerifyFileResponse, WalkDirRequest, WalkDirResponse, WriteAllRequest, WriteAllResponse, WriteMetadataRequest,
WriteMetadataResponse, WriteRequest, WriteResponse,
MakeVolumesResponse, NsScannerRequest, NsScannerResponse, PingRequest, PingResponse, ReadAllRequest, ReadAllResponse,
ReadAtRequest, ReadAtResponse, ReadMultipleRequest, ReadMultipleResponse, ReadVersionRequest, ReadVersionResponse,
ReadXlRequest, ReadXlResponse, RenameDataRequest, RenameDataResponse, RenameFileRequst, RenameFileResponse,
RenamePartRequst, RenamePartResponse, StatVolumeRequest, StatVolumeResponse, UpdateMetadataRequest,
UpdateMetadataResponse, VerifyFileRequest, VerifyFileResponse, WalkDirRequest, WalkDirResponse, WriteAllRequest,
WriteAllResponse, WriteMetadataRequest, WriteMetadataResponse, WriteRequest, WriteResponse,
},
};
use tokio::sync::mpsc;
@@ -1270,6 +1271,96 @@ impl Node for NodeService {
}
}
type NsScannerStream = ResponseStream<NsScannerResponse>;
async fn ns_scanner(&self, request: Request<Streaming<NsScannerRequest>>) -> Result<Response<Self::NsScannerStream>, Status> {
info!("ns_scanner");
let mut in_stream = request.into_inner();
let (tx, rx) = mpsc::channel(10);
tokio::spawn(async move {
match in_stream.next().await {
Some(Ok(request)) => {
if let Some(disk) = find_local_disk(&request.disk).await {
let cache = match serde_json::from_str::<DataUsageCache>(&request.cache) {
Ok(cache) => cache,
Err(_) => {
tx.send(Ok(NsScannerResponse {
success: false,
update: "".to_string(),
data_usage_cache: "".to_string(),
error_info: Some("can not decode DataUsageCache".to_string()),
}))
.await
.expect("working rx");
return;
}
};
let (updates_tx, mut updates_rx) = mpsc::channel(100);
let tx_clone = tx.clone();
let task = tokio::spawn(async move {
loop {
match updates_rx.recv().await {
Some(update) => {
let update = serde_json::to_string(&update).expect("encode failed");
tx_clone
.send(Ok(NsScannerResponse {
success: true,
update,
data_usage_cache: "".to_string(),
error_info: Some("can not decode DataUsageCache".to_string()),
}))
.await
.expect("working rx");
}
None => return,
}
}
});
let data_usage_cache = disk.ns_scanner(&cache, updates_tx, request.scan_mode as usize).await;
let _ = task.await;
match data_usage_cache {
Ok(data_usage_cache) => {
let data_usage_cache = serde_json::to_string(&data_usage_cache).expect("encode failed");
tx.send(Ok(NsScannerResponse {
success: true,
update: "".to_string(),
data_usage_cache,
error_info: Some("can not decode DataUsageCache".to_string()),
}))
.await
.expect("working rx");
}
Err(_) => {
tx.send(Ok(NsScannerResponse {
success: false,
update: "".to_string(),
data_usage_cache: "".to_string(),
error_info: Some("scanner failed".to_string()),
}))
.await
.expect("working rx");
}
}
} else {
tx.send(Ok(NsScannerResponse {
success: false,
update: "".to_string(),
data_usage_cache: "".to_string(),
error_info: Some("can not find disk".to_string()),
}))
.await
.expect("working rx");
}
}
_ => todo!(),
}
});
let out_stream = ReceiverStream::new(rx);
Ok(tonic::Response::new(Box::pin(out_stream)))
}
async fn lock(&self, request: Request<GenerallyLockRequest>) -> Result<Response<GenerallyLockResponse>, Status> {
let request = request.into_inner();
match &serde_json::from_str::<LockArgs>(&request.args) {
-2
View File
@@ -6,12 +6,10 @@ mod storage;
use clap::Parser;
use common::error::{Error, Result};
use ecstore::{
bucket::metadata_sys::init_bucket_metadata_sys,
endpoints::EndpointServerPools,
heal::data_scanner::init_data_scanner,
set_global_endpoints,
store::{init_local_disks, ECStore},
store_api::{BucketOptions, StorageAPI},
update_erasure_type,
};
use grpc::make_server;