mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-06 13:27:43 +00:00
init metacache io
This commit is contained in:
@@ -12,6 +12,7 @@ use tokio::{
|
||||
use crate::{
|
||||
disk::{DiskAPI, DiskStore, MetaCacheEntries, MetaCacheEntry, WalkDirOptions},
|
||||
error::{Error, Result},
|
||||
io::Writer,
|
||||
};
|
||||
|
||||
type AgreedFn = Box<dyn Fn(MetaCacheEntry) -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + 'static>;
|
||||
@@ -77,16 +78,19 @@ pub async fn list_path_raw(mut rx: B_Receiver<bool>, opts: ListPathRawOptions) -
|
||||
match disk
|
||||
.as_ref()
|
||||
.unwrap()
|
||||
.walk_dir(WalkDirOptions {
|
||||
bucket: opts_clone.bucket.clone(),
|
||||
base_dir: opts_clone.path.clone(),
|
||||
recursive: opts_clone.recursice,
|
||||
report_notfound: opts_clone.report_not_found,
|
||||
filter_prefix: opts_clone.filter_prefix.clone(),
|
||||
forward_to: opts_clone.forward_to.clone(),
|
||||
limit: opts_clone.per_disk_limit,
|
||||
..Default::default()
|
||||
})
|
||||
.walk_dir(
|
||||
WalkDirOptions {
|
||||
bucket: opts_clone.bucket.clone(),
|
||||
base_dir: opts_clone.path.clone(),
|
||||
recursive: opts_clone.recursice,
|
||||
report_notfound: opts_clone.report_not_found,
|
||||
filter_prefix: opts_clone.filter_prefix.clone(),
|
||||
forward_to: opts_clone.forward_to.clone(),
|
||||
limit: opts_clone.per_disk_limit,
|
||||
..Default::default()
|
||||
},
|
||||
Writer::NotUse,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(r) => {
|
||||
@@ -115,16 +119,19 @@ pub async fn list_path_raw(mut rx: B_Receiver<bool>, opts: ListPathRawOptions) -
|
||||
match disk
|
||||
.as_ref()
|
||||
.unwrap()
|
||||
.walk_dir(WalkDirOptions {
|
||||
bucket: opts_clone.bucket.clone(),
|
||||
base_dir: opts_clone.path.clone(),
|
||||
recursive: opts_clone.recursice,
|
||||
report_notfound: opts_clone.report_not_found,
|
||||
filter_prefix: opts_clone.filter_prefix.clone(),
|
||||
forward_to: opts_clone.forward_to.clone(),
|
||||
limit: opts_clone.per_disk_limit,
|
||||
..Default::default()
|
||||
})
|
||||
.walk_dir(
|
||||
WalkDirOptions {
|
||||
bucket: opts_clone.bucket.clone(),
|
||||
base_dir: opts_clone.path.clone(),
|
||||
recursive: opts_clone.recursice,
|
||||
report_notfound: opts_clone.report_not_found,
|
||||
filter_prefix: opts_clone.filter_prefix.clone(),
|
||||
forward_to: opts_clone.forward_to.clone(),
|
||||
limit: opts_clone.per_disk_limit,
|
||||
..Default::default()
|
||||
},
|
||||
Writer::NotUse,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(r) => {
|
||||
|
||||
@@ -1291,10 +1291,18 @@ impl DiskAPI for LocalDisk {
|
||||
Ok(entries)
|
||||
}
|
||||
|
||||
// TODO: io.writer
|
||||
async fn walk_dir(&self, opts: WalkDirOptions) -> Result<Vec<MetaCacheEntry>> {
|
||||
// FIXME: TODO: io.writer
|
||||
async fn walk_dir(&self, opts: WalkDirOptions, wr: crate::io::Writer) -> Result<Vec<MetaCacheEntry>> {
|
||||
// warn!("walk_dir opts {:?}", &opts);
|
||||
|
||||
let volume_dir = self.get_bucket_path(&opts.bucket)?;
|
||||
|
||||
if !skip_access_checks(&opts.bucket) {
|
||||
if let Err(e) = access(&volume_dir).await {
|
||||
return Err(convert_access_error(e, DiskError::VolumeAccessDenied));
|
||||
}
|
||||
}
|
||||
|
||||
let mut metas = Vec::new();
|
||||
|
||||
if opts.base_dir.ends_with(SLASH_SEPARATOR) {
|
||||
|
||||
@@ -22,6 +22,7 @@ use crate::{
|
||||
data_usage_cache::{DataUsageCache, DataUsageEntry},
|
||||
heal_commands::{HealScanMode, HealingTracker},
|
||||
},
|
||||
io,
|
||||
store_api::{FileInfo, RawFileInfo},
|
||||
};
|
||||
use endpoint::Endpoint;
|
||||
@@ -208,10 +209,10 @@ impl DiskAPI for Disk {
|
||||
}
|
||||
}
|
||||
|
||||
async fn walk_dir(&self, opts: WalkDirOptions) -> Result<Vec<MetaCacheEntry>> {
|
||||
async fn walk_dir(&self, opts: WalkDirOptions, wr: crate::io::Writer) -> Result<Vec<MetaCacheEntry>> {
|
||||
match self {
|
||||
Disk::Local(local_disk) => local_disk.walk_dir(opts).await,
|
||||
Disk::Remote(remote_disk) => remote_disk.walk_dir(opts).await,
|
||||
Disk::Local(local_disk) => local_disk.walk_dir(opts, wr).await,
|
||||
Disk::Remote(remote_disk) => remote_disk.walk_dir(opts, wr).await,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -403,7 +404,7 @@ pub trait DiskAPI: Debug + Send + Sync + 'static {
|
||||
async fn delete_volume(&self, volume: &str) -> Result<()>;
|
||||
|
||||
// 并发边读边写 TODO: wr io.Writer
|
||||
async fn walk_dir(&self, opts: WalkDirOptions) -> Result<Vec<MetaCacheEntry>>;
|
||||
async fn walk_dir(&self, opts: WalkDirOptions, wr: crate::io::Writer) -> Result<Vec<MetaCacheEntry>>;
|
||||
|
||||
// Metadata operations
|
||||
async fn delete_version(
|
||||
@@ -599,10 +600,10 @@ pub struct MetaCacheEntry {
|
||||
pub metadata: Vec<u8>,
|
||||
|
||||
// cached contains the metadata if decoded.
|
||||
cached: Option<FileMeta>,
|
||||
pub cached: Option<FileMeta>,
|
||||
|
||||
// Indicates the entry can be reused and only one reference to metadata is expected.
|
||||
_reusable: bool,
|
||||
pub reusable: bool,
|
||||
}
|
||||
|
||||
impl MetaCacheEntry {
|
||||
@@ -616,6 +617,7 @@ impl MetaCacheEntry {
|
||||
|
||||
Ok(wr)
|
||||
}
|
||||
|
||||
pub fn is_dir(&self) -> bool {
|
||||
self.metadata.is_empty() && self.name.ends_with('/')
|
||||
}
|
||||
@@ -838,7 +840,7 @@ impl MetaCacheEntries {
|
||||
meta_ver: selected.as_ref().unwrap().cached.as_ref().unwrap().meta_ver,
|
||||
..Default::default()
|
||||
}),
|
||||
_reusable: true,
|
||||
reusable: true,
|
||||
..Default::default()
|
||||
});
|
||||
|
||||
|
||||
@@ -346,7 +346,7 @@ impl DiskAPI for RemoteDisk {
|
||||
Ok(response.volumes)
|
||||
}
|
||||
|
||||
async fn walk_dir(&self, opts: WalkDirOptions) -> Result<Vec<MetaCacheEntry>> {
|
||||
async fn walk_dir(&self, opts: WalkDirOptions, wr: crate::io::Writer) -> Result<Vec<MetaCacheEntry>> {
|
||||
info!("walk_dir");
|
||||
let walk_dir_options = serde_json::to_string(&opts)?;
|
||||
let mut client = node_service_time_out_client(&self.addr)
|
||||
|
||||
@@ -0,0 +1,68 @@
|
||||
use std::pin::Pin;
|
||||
use std::task::{Context, Poll};
|
||||
use tokio::fs::File;
|
||||
use tokio::io::{self, AsyncRead, AsyncWrite, ReadBuf};
|
||||
|
||||
#[derive(Default)]
|
||||
pub enum Reader {
|
||||
#[default]
|
||||
NotUse,
|
||||
File(File),
|
||||
}
|
||||
|
||||
impl AsyncRead for Reader {
|
||||
fn poll_read(self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll<io::Result<()>> {
|
||||
match self.get_mut() {
|
||||
Reader::File(file) => {
|
||||
let file = Pin::new(file);
|
||||
file.poll_read(cx, buf)
|
||||
}
|
||||
Reader::NotUse => Poll::Ready(Ok(())),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
pub enum Writer {
|
||||
#[default]
|
||||
NotUse,
|
||||
File(File),
|
||||
}
|
||||
|
||||
impl AsyncWrite for Writer {
|
||||
fn poll_write(self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &[u8]) -> Poll<io::Result<usize>> {
|
||||
match self.get_mut() {
|
||||
Writer::File(file) => {
|
||||
// Create a pinned reference from the file
|
||||
let file = Pin::new(file);
|
||||
file.poll_write(cx, buf)
|
||||
}
|
||||
Writer::NotUse => Poll::Ready(Ok(0)),
|
||||
}
|
||||
}
|
||||
|
||||
fn poll_flush(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<io::Result<()>> {
|
||||
match self.get_mut() {
|
||||
Writer::File(file) => {
|
||||
let file = Pin::new(file);
|
||||
file.poll_flush(cx)
|
||||
}
|
||||
Writer::NotUse => Poll::Ready(Ok(())),
|
||||
}
|
||||
}
|
||||
|
||||
fn poll_shutdown(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<io::Result<()>> {
|
||||
match self.get_mut() {
|
||||
Writer::File(file) => {
|
||||
let file = Pin::new(file);
|
||||
file.poll_shutdown(cx)
|
||||
}
|
||||
Writer::NotUse => Poll::Ready(Ok(())),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// #[tokio::test]
|
||||
// async fn test_reader{
|
||||
|
||||
// }
|
||||
@@ -13,6 +13,8 @@ mod file_meta;
|
||||
pub mod file_meta_inline;
|
||||
pub mod global;
|
||||
pub mod heal;
|
||||
pub mod io;
|
||||
pub mod metacache;
|
||||
pub mod metrics_realtime;
|
||||
pub mod notification_sys;
|
||||
pub mod peer;
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
pub mod writer;
|
||||
@@ -0,0 +1,206 @@
|
||||
use std::io::ErrorKind;
|
||||
use std::io::Read;
|
||||
use std::io::Write;
|
||||
use std::str::from_utf8;
|
||||
|
||||
use crate::disk::MetaCacheEntry;
|
||||
use crate::error::Error;
|
||||
use crate::error::Result;
|
||||
|
||||
const METACACHE_STREAM_VERSION: u8 = 2;
|
||||
|
||||
pub struct MetacacheWriter<W> {
|
||||
wr: W,
|
||||
buf: Vec<u8>,
|
||||
created: bool,
|
||||
}
|
||||
|
||||
impl<W: Write + Unpin> MetacacheWriter<W> {
|
||||
pub fn new(wr: W, block_size: usize) -> Self {
|
||||
Self {
|
||||
wr,
|
||||
buf: Vec::with_capacity(block_size),
|
||||
created: false,
|
||||
}
|
||||
}
|
||||
|
||||
async fn write(&mut self, objs: &[MetaCacheEntry]) -> Result<()> {
|
||||
if objs.is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
if !self.created {
|
||||
rmp::encode::write_u8(&mut self.wr, METACACHE_STREAM_VERSION)?;
|
||||
self.created = false;
|
||||
}
|
||||
|
||||
for obj in objs.iter() {
|
||||
if obj.name.is_empty() {
|
||||
return Err(Error::msg("metacacheWriter: no name"));
|
||||
}
|
||||
|
||||
rmp::encode::write_bool(&mut self.wr, true)?;
|
||||
|
||||
rmp::encode::write_str(&mut self.wr, &obj.name)?;
|
||||
|
||||
rmp::encode::write_bin(&mut self.wr, &obj.metadata)?;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn close(&mut self) -> Result<()> {
|
||||
rmp::encode::write_bool(&mut self.wr, false)?;
|
||||
|
||||
self.wr.flush()?;
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
pub struct MetacacheReader<R> {
|
||||
rd: R,
|
||||
init: bool,
|
||||
err: Option<Error>,
|
||||
buf: Vec<u8>,
|
||||
}
|
||||
|
||||
impl<R: Read> MetacacheReader<R> {
|
||||
pub fn new(rd: R) -> Self {
|
||||
Self {
|
||||
rd,
|
||||
init: false,
|
||||
err: None,
|
||||
buf: Vec::new(),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn check_init(&mut self) {
|
||||
if !self.init {
|
||||
let ver = match rmp::decode::read_u8(&mut self.rd) {
|
||||
Ok(res) => res,
|
||||
Err(err) => {
|
||||
self.err = Some(Error::msg(err.to_string()));
|
||||
0
|
||||
}
|
||||
};
|
||||
match ver {
|
||||
1 | 2 => (),
|
||||
_ => {
|
||||
self.err = Some(Error::msg("invalid version"));
|
||||
}
|
||||
}
|
||||
|
||||
self.init = true;
|
||||
}
|
||||
}
|
||||
|
||||
pub fn peek(&mut self) -> Result<MetaCacheEntry> {
|
||||
self.check_init();
|
||||
|
||||
if let Some(err) = &self.err {
|
||||
return Err(err.clone());
|
||||
}
|
||||
|
||||
match rmp::decode::read_bool(&mut self.rd) {
|
||||
Ok(res) => {
|
||||
if !res {
|
||||
self.err = Some(Error::new(std::io::Error::from(ErrorKind::UnexpectedEof)));
|
||||
return Err(Error::new(std::io::Error::from(ErrorKind::UnexpectedEof)));
|
||||
}
|
||||
}
|
||||
Err(err) => {
|
||||
self.err = Some(Error::msg(err.to_string()));
|
||||
return Err(Error::new(err));
|
||||
}
|
||||
};
|
||||
|
||||
let l = match rmp::decode::read_str_len(&mut self.rd) {
|
||||
Ok(res) => res,
|
||||
Err(err) => {
|
||||
self.err = Some(Error::msg(err.to_string()));
|
||||
return Err(Error::new(err));
|
||||
}
|
||||
};
|
||||
|
||||
self.buf.resize(l as usize, 0);
|
||||
let name = match self.rd.read_exact(&mut self.buf) {
|
||||
Ok(()) => {
|
||||
let name_buf = self.buf.to_vec();
|
||||
match from_utf8(&name_buf) {
|
||||
Ok(decoded) => Ok(decoded.to_owned()),
|
||||
Err(err) => {
|
||||
self.err = Some(Error::msg(err.to_string()));
|
||||
Err(Error::msg(err.to_string()))
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(err) => {
|
||||
self.err = Some(Error::msg(err.to_string()));
|
||||
Err(Error::msg(err.to_string()))
|
||||
}
|
||||
}?;
|
||||
|
||||
let l = match rmp::decode::read_bin_len(&mut self.rd) {
|
||||
Ok(res) => res,
|
||||
Err(err) => {
|
||||
self.err = Some(Error::msg(err.to_string()));
|
||||
return Err(Error::new(err));
|
||||
}
|
||||
};
|
||||
self.buf.resize(l as usize, 0);
|
||||
match self.rd.read_exact(&mut self.buf) {
|
||||
Ok(res) => res,
|
||||
Err(err) => {
|
||||
self.err = Some(Error::msg(err.to_string()));
|
||||
return Err(Error::new(err));
|
||||
}
|
||||
};
|
||||
|
||||
let metadata = self.buf.clone();
|
||||
|
||||
Ok(MetaCacheEntry {
|
||||
name,
|
||||
metadata,
|
||||
cached: None,
|
||||
reusable: false,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_writer() {
|
||||
use std::fs::File;
|
||||
use std::fs::OpenOptions;
|
||||
|
||||
let file_path = "./test_writer.txt";
|
||||
let f = OpenOptions::new()
|
||||
.create(true)
|
||||
.read(true)
|
||||
.write(true)
|
||||
.truncate(true)
|
||||
.open(file_path)
|
||||
.unwrap();
|
||||
|
||||
// let wr = Writer::File(f);
|
||||
|
||||
let mut w = MetacacheWriter::new(f, 1024);
|
||||
|
||||
let mut objs = Vec::new();
|
||||
for i in 0..10 {
|
||||
objs.push(MetaCacheEntry {
|
||||
name: format!("item{}", i),
|
||||
metadata: vec![0u8, 10],
|
||||
cached: None,
|
||||
reusable: false,
|
||||
});
|
||||
}
|
||||
|
||||
w.write(&objs).await.unwrap();
|
||||
w.close().await.unwrap();
|
||||
|
||||
let nf = File::open(file_path).unwrap();
|
||||
|
||||
let meta = nf.metadata().unwrap();
|
||||
|
||||
println!("{}", meta.len());
|
||||
}
|
||||
@@ -6,7 +6,6 @@ use std::{
|
||||
time::Duration,
|
||||
};
|
||||
|
||||
use crate::heal::heal_ops::{HealEntryFn, HealSequence};
|
||||
use crate::{
|
||||
bitrot::{bitrot_verify, close_bitrot_writers, new_bitrot_filereader, new_bitrot_filewriter, BitrotFileWriter},
|
||||
cache_value::metacache_set::{list_path_raw, ListPathRawOptions},
|
||||
@@ -56,6 +55,10 @@ use crate::{
|
||||
heal::data_scanner::{globalHealConfig, HEAL_DELETE_DANGLING},
|
||||
store_api::ListObjectVersionsInfo,
|
||||
};
|
||||
use crate::{
|
||||
heal::heal_ops::{HealEntryFn, HealSequence},
|
||||
io::Writer,
|
||||
};
|
||||
use futures::future::join_all;
|
||||
use glob::Pattern;
|
||||
use http::HeaderMap;
|
||||
@@ -1333,7 +1336,7 @@ impl SetDisks {
|
||||
let disk = disk.as_ref().unwrap();
|
||||
let opts = opts.clone();
|
||||
// let mut wr = &mut wr;
|
||||
futures.push(disk.walk_dir(opts));
|
||||
futures.push(disk.walk_dir(opts, Writer::NotUse));
|
||||
// tokio::spawn(async move { disk.walk_dir(opts, wr).await });
|
||||
}
|
||||
|
||||
|
||||
+1
-1
@@ -743,7 +743,7 @@ impl Node for NodeService {
|
||||
}));
|
||||
}
|
||||
};
|
||||
match disk.walk_dir(opts).await {
|
||||
match disk.walk_dir(opts, ecstore::io::Writer::NotUse).await {
|
||||
Ok(entries) => {
|
||||
let entries = entries
|
||||
.into_iter()
|
||||
|
||||
Reference in New Issue
Block a user