use super::{endpoint::Endpoint, error::DiskError, format::FormatV3}; use super::{ DeleteOptions, DiskAPI, FileReader, FileWriter, ReadMultipleReq, ReadMultipleResp, ReadOptions, RenameDataResp, VolumeInfo, }; use crate::{ error::{Error, Result}, file_meta::FileMeta, store_api::{FileInfo, RawFileInfo}, utils, }; use bytes::Bytes; use path_absolutize::Absolutize; use std::{ fs::Metadata, path::{Path, PathBuf}, }; use time::OffsetDateTime; use tokio::fs::{self, File}; use tokio::io::ErrorKind; use tracing::{debug, warn}; use uuid::Uuid; #[derive(Debug)] pub struct LocalDisk { pub root: PathBuf, pub id: Uuid, pub _format_data: Vec, pub _format_meta: Option, pub _format_path: PathBuf, // pub format_legacy: bool, // drop pub _format_last_check: OffsetDateTime, } impl LocalDisk { pub async fn new(ep: &Endpoint, cleanup: bool) -> Result { let root = fs::canonicalize(ep.url.path()).await?; if cleanup { // TODO: 删除tmp数据 } let format_path = Path::new(super::RUSTFS_META_BUCKET) .join(Path::new(super::FORMAT_CONFIG_FILE)) .absolutize_virtually(&root)? .into_owned(); let (format_data, format_meta) = read_file_exists(&format_path).await?; let mut id = Uuid::nil(); // let mut format_legacy = false; let mut format_last_check = OffsetDateTime::UNIX_EPOCH; if !format_data.is_empty() { let s = format_data.as_slice(); let fm = FormatV3::try_from(s)?; let (set_idx, disk_idx) = fm.find_disk_index_by_disk_id(fm.erasure.this)?; if Some(set_idx) != ep.set_idx || Some(disk_idx) != ep.disk_idx { return Err(Error::from(DiskError::InconsistentDisk)); } id = fm.erasure.this; // format_legacy = fm.erasure.distribution_algo == DistributionAlgoVersion::V1; format_last_check = OffsetDateTime::now_utc(); } let disk = Self { root, id, _format_meta: format_meta, _format_data: format_data, _format_path: format_path, // format_legacy, _format_last_check: format_last_check, }; disk.make_meta_volumes().await?; Ok(disk) } async fn make_meta_volumes(&self) -> Result<()> { let buckets = format!("{}/{}", super::RUSTFS_META_BUCKET, super::BUCKET_META_PREFIX); let multipart = format!("{}/{}", super::RUSTFS_META_BUCKET, "multipart"); let config = format!("{}/{}", super::RUSTFS_META_BUCKET, "config"); let tmp = format!("{}/{}", super::RUSTFS_META_BUCKET, "tmp"); let defaults = vec![buckets.as_str(), multipart.as_str(), config.as_str(), tmp.as_str()]; self.make_volumes(defaults).await } pub fn resolve_abs_path(&self, path: impl AsRef) -> Result { Ok(path.as_ref().absolutize_virtually(&self.root)?.into_owned()) } pub fn get_object_path(&self, bucket: &str, key: &str) -> Result { let dir = Path::new(&bucket); let file_path = Path::new(&key); self.resolve_abs_path(dir.join(file_path)) } pub fn get_bucket_path(&self, bucket: &str) -> Result { let dir = Path::new(&bucket); self.resolve_abs_path(dir) } // /// Write to the filesystem atomically. // /// This is done by first writing to a temporary location and then moving the file. // pub(crate) async fn prepare_file_write<'a>(&self, path: &'a PathBuf) -> Result> { // let tmp_path = self.get_object_path(RUSTFS_META_TMP_BUCKET, Uuid::new_v4().to_string().as_str())?; // debug!("prepare_file_write tmp_path:{:?}, path:{:?}", &tmp_path, &path); // let file = File::create(&tmp_path).await?; // let writer = BufWriter::new(file); // Ok(FileWriter { // tmp_path, // dest_path: path, // writer, // clean_tmp: true, // }) // } pub async fn rename_all(&self, src_data_path: &PathBuf, dst_data_path: &PathBuf, skip: &PathBuf) -> Result<()> { if !skip.starts_with(src_data_path) { fs::create_dir_all(dst_data_path.parent().unwrap_or(Path::new("/"))).await?; } debug!( "rename_all from \n {:?} \n to \n {:?} \n skip:{:?}", &src_data_path, &dst_data_path, &skip ); fs::rename(&src_data_path, &dst_data_path).await?; Ok(()) } pub async fn delete_file(&self, base_path: &PathBuf, delete_path: &PathBuf, recursive: bool, _immediate: bool) -> Result<()> { debug!("delete_file {:?}", &delete_path); if is_root_path(base_path) || is_root_path(delete_path) { return Ok(()); } if !delete_path.starts_with(base_path) || base_path == delete_path { return Ok(()); } if recursive { let trash_path = self.get_object_path(super::RUSTFS_META_TMP_DELETED_BUCKET, Uuid::new_v4().to_string().as_str())?; // fs::create_dir_all(&trash_path).await?; fs::rename(&delete_path, &trash_path).await.map_err(|err| { // 使用文件路径自定义错误信息 std::io::Error::new( err.kind(), format!("Failed to rename file '{:?}' to '{:?}': {}", &delete_path, &trash_path, err), ) })?; // TODO: immediate return Ok(()); } if delete_path.is_dir() { fs::remove_dir(delete_path).await?; } else { fs::remove_file(delete_path).await?; } if let Some(dir_path) = delete_path.parent() { Box::pin(self.delete_file(base_path, &PathBuf::from(dir_path), false, false)).await?; } Ok(()) } /// read xl.meta raw data async fn read_raw( &self, bucket: &str, volume_dir: impl AsRef, path: impl AsRef, read_data: bool, ) -> Result<(Vec, OffsetDateTime)> { let meta_path = path.as_ref().join(Path::new(super::STORAGE_FORMAT_FILE)); if read_data { self.read_all_data(bucket, volume_dir, meta_path).await } else { self.read_metadata_with_dmtime(meta_path).await } } async fn read_metadata_with_dmtime(&self, path: impl AsRef) -> Result<(Vec, OffsetDateTime)> { let (data, meta) = read_file_all(path).await?; let modtime = match meta.modified() { Ok(md) => OffsetDateTime::from(md), Err(_) => return Err(Error::msg("Not supported modified on this platform")), }; Ok((data, modtime)) } async fn read_all_data( &self, _bucket: &str, _volume_dir: impl AsRef, path: impl AsRef, ) -> Result<(Vec, OffsetDateTime)> { let (data, meta) = read_file_all(path).await?; let modtime = match meta.modified() { Ok(md) => OffsetDateTime::from(md), Err(_) => return Err(Error::msg("Not supported modified on this platform")), }; Ok((data, modtime)) } } fn is_root_path(path: impl AsRef) -> bool { path.as_ref().components().count() == 1 && path.as_ref().has_root() } // 过滤 std::io::ErrorKind::NotFound pub async fn read_file_exists(path: impl AsRef) -> Result<(Vec, Option)> { let p = path.as_ref(); let (data, meta) = match read_file_all(&p).await { Ok((data, meta)) => (data, Some(meta)), Err(e) => { if DiskError::FileNotFound.is(&e) { (Vec::new(), None) } else { return Err(e); } } }; // let mut data = Vec::new(); // if meta.is_some() { // data = fs::read(&p).await?; // } Ok((data, meta)) } pub async fn write_all_internal(p: impl AsRef, data: impl AsRef<[u8]>) -> Result<()> { // create top dir if not exists fs::create_dir_all(&p.as_ref().parent().unwrap_or_else(|| Path::new("."))).await?; fs::write(&p, data).await?; Ok(()) } pub async fn read_file_all(path: impl AsRef) -> Result<(Vec, Metadata)> { let p = path.as_ref(); let meta = read_file_metadata(&path).await?; let data = fs::read(&p).await?; Ok((data, meta)) } pub async fn read_file_metadata(p: impl AsRef) -> Result { let meta = fs::metadata(&p).await.map_err(|e| match e.kind() { ErrorKind::NotFound => Error::from(DiskError::FileNotFound), ErrorKind::PermissionDenied => Error::from(DiskError::FileAccessDenied), _ => Error::from(e), })?; Ok(meta) } pub async fn check_volume_exists(p: impl AsRef) -> Result<()> { fs::metadata(&p).await.map_err(|e| match e.kind() { ErrorKind::NotFound => Error::from(DiskError::VolumeNotFound), ErrorKind::PermissionDenied => Error::from(DiskError::FileAccessDenied), _ => Error::from(e), })?; Ok(()) } fn skip_access_checks(p: impl AsRef) -> bool { let vols = [ super::RUSTFS_META_TMP_DELETED_BUCKET, super::RUSTFS_META_TMP_BUCKET, super::RUSTFS_META_MULTIPART_BUCKET, super::RUSTFS_META_BUCKET, ]; for v in vols.iter() { if p.as_ref().starts_with(v) { return true; } } false } #[async_trait::async_trait] impl DiskAPI for LocalDisk { fn is_local(&self) -> bool { true } fn id(&self) -> Uuid { self.id } #[must_use] async fn read_all(&self, volume: &str, path: &str) -> Result { let p = self.get_object_path(volume, path)?; let (data, _) = read_file_all(&p).await?; Ok(Bytes::from(data)) } async fn write_all(&self, volume: &str, path: &str, data: Vec) -> Result<()> { let p = self.get_object_path(volume, path)?; write_all_internal(p, data).await?; Ok(()) } async fn delete(&self, volume: &str, path: &str, opt: DeleteOptions) -> Result<()> { let vol_path = self.get_bucket_path(volume)?; if !skip_access_checks(volume) { check_volume_exists(&vol_path).await?; } let fpath = self.get_object_path(volume, path)?; self.delete_file(&vol_path, &fpath, opt.recursive, opt.immediate).await?; // if opt.recursive { // let trash_path = self.get_object_path(RUSTFS_META_TMP_DELETED_BUCKET, Uuid::new_v4().to_string().as_str())?; // fs::create_dir_all(&trash_path).await?; // fs::rename(&fpath, &trash_path).await?; // // TODO: immediate // return Ok(()); // } // fs::remove_file(fpath).await?; Ok(()) } async fn rename_file(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str) -> Result<()> { let src_volume_path = self.get_bucket_path(src_volume)?; if !skip_access_checks(src_volume) { check_volume_exists(&src_volume_path).await?; } if !skip_access_checks(dst_volume) { let vol_path = self.get_bucket_path(dst_volume)?; check_volume_exists(&vol_path).await?; } let srcp = self.get_object_path(src_volume, src_path)?; let dstp = self.get_object_path(dst_volume, dst_path)?; let src_is_dir = srcp.is_dir(); let dst_is_dir = dstp.is_dir(); if !src_is_dir && dst_is_dir || src_is_dir && !dst_is_dir { return Err(Error::from(DiskError::FileAccessDenied)); } // TODO: check path length if src_is_dir { // TODO: remove dst_dir } fs::create_dir_all(dstp.parent().unwrap_or_else(|| Path::new("."))).await?; let mut idx = 0; loop { if let Err(e) = fs::rename(&srcp, &dstp).await { if e.kind() == ErrorKind::NotFound && idx == 0 { idx += 1; continue; } }; break; } if let Some(dir_path) = srcp.parent() { self.delete_file(&src_volume_path, &PathBuf::from(dir_path), false, false) .await?; } Ok(()) } async fn create_file(&self, _origvolume: &str, volume: &str, path: &str, _file_size: usize) -> Result { let fpath = self.get_object_path(volume, path)?; debug!("CreateFile fpath: {:?}", fpath); if let Some(_dir_path) = fpath.parent() { fs::create_dir_all(&_dir_path).await?; } let file = File::create(&fpath).await?; Ok(FileWriter::new(file)) // let mut writer = BufWriter::new(file); // io::copy(&mut r, &mut writer).await?; // Ok(()) } // async fn append_file(&self, volume: &str, path: &str, mut r: DuplexStream) -> Result { async fn append_file(&self, volume: &str, path: &str) -> Result { let p = self.get_object_path(volume, path)?; if let Some(dir_path) = p.parent() { fs::create_dir_all(&dir_path).await?; } // debug!("append_file open {} {:?}", self.id(), &p); let file = File::options() .read(true) .create(true) .write(true) .append(true) .open(&p) .await?; Ok(FileWriter::new(file)) // let mut writer = BufWriter::new(file); // io::copy(&mut r, &mut writer).await?; // debug!("append_file end {} {}", self.id(), path); // io::copy(&mut r, &mut file).await?; // Ok(()) } async fn read_file(&self, volume: &str, path: &str) -> Result { let p = self.get_object_path(volume, path)?; debug!("read_file {:?}", &p); let file = File::options().read(true).open(&p).await?; Ok(FileReader::new(file)) // file.seek(SeekFrom::Start(offset as u64)).await?; // let mut buffer = vec![0; length]; // let bytes_read = file.read(&mut buffer).await?; // buffer.truncate(bytes_read); // Ok((buffer, bytes_read)) } async fn list_dir(&self, _origvolume: &str, volume: &str, _dir_path: &str, _count: usize) -> Result> { let p = self.get_bucket_path(volume)?; let mut entries = fs::read_dir(&p).await?; let mut volumes = Vec::new(); while let Some(entry) = entries.next_entry().await? { if let Ok(_metadata) = entry.metadata().await { // if !metadata.is_dir() { // continue; // } let name = entry.file_name().to_string_lossy().to_string(); // let created = match metadata.created() { // Ok(md) => OffsetDateTime::from(md), // Err(_) => return Err(Error::msg("Not supported created on this platform")), // }; volumes.push(name); } } Ok(volumes) } async fn walk_dir(&self) -> Result> { unimplemented!() } async fn rename_data( &self, src_volume: &str, src_path: &str, fi: FileInfo, dst_volume: &str, dst_path: &str, ) -> Result { let src_volume_path = self.get_bucket_path(src_volume)?; if !skip_access_checks(src_volume) { check_volume_exists(&src_volume_path).await?; } let dst_volume_path = self.get_bucket_path(dst_volume)?; if !skip_access_checks(dst_volume) { check_volume_exists(&dst_volume_path).await?; } let src_file_path = self.get_object_path(src_volume, format!("{}/{}", &src_path, super::STORAGE_FORMAT_FILE).as_str())?; let dst_file_path = self.get_object_path(dst_volume, format!("{}/{}", &dst_path, super::STORAGE_FORMAT_FILE).as_str())?; let (src_data_path, dst_data_path) = { let mut data_dir = String::new(); if !fi.is_remote() { data_dir = utils::path::retain_slash(fi.data_dir.to_string().as_str()); } if !data_dir.is_empty() { let src_data_path = self.get_object_path( src_volume, utils::path::retain_slash(format!("{}/{}", &src_path, data_dir).as_str()).as_str(), )?; let dst_data_path = self.get_object_path(dst_volume, format!("{}/{}", &dst_path, data_dir).as_str())?; (src_data_path, dst_data_path) } else { (PathBuf::new(), PathBuf::new()) } }; // let curreng_data_path = self.get_object_path(&dst_volume, &dst_path); let mut meta = FileMeta::new(); let (dst_buf, _) = read_file_exists(&dst_file_path).await?; let mut skip_parent = dst_volume_path; if !&dst_buf.is_empty() { skip_parent = PathBuf::from(&dst_file_path.parent().unwrap_or(Path::new("/"))); } if !dst_buf.is_empty() { meta = match FileMeta::unmarshal(&dst_buf) { Ok(m) => m, Err(_) => FileMeta::new(), } // xl.load // meta.from(dst_buf); } meta.add_version(fi.clone())?; let fm_data = meta.marshal_msg()?; fs::write(&src_file_path, fm_data).await?; let no_inline = src_data_path.has_root() && fi.data.is_none() && fi.size > 0; if no_inline { self.rename_all(&src_data_path, &dst_data_path, &skip_parent).await?; } self.rename_all(&src_file_path, &dst_file_path, &skip_parent).await?; if src_volume != super::RUSTFS_META_MULTIPART_BUCKET { fs::remove_dir(&src_file_path.parent().unwrap()).await?; } else { self.delete_file(&src_volume_path, &PathBuf::from(src_file_path.parent().unwrap()), true, false) .await?; } Ok(RenameDataResp { old_data_dir: fi.data_dir.to_string(), }) } async fn make_volumes(&self, volumes: Vec<&str>) -> Result<()> { for vol in volumes { if let Err(e) = self.make_volume(vol).await { match &e.downcast_ref::() { Some(DiskError::VolumeExists) => Ok(()), Some(_) => Err(e), None => Err(e), }?; } // TODO: health check } Ok(()) } async fn make_volume(&self, volume: &str) -> Result<()> { let p = self.get_bucket_path(volume)?; match File::open(&p).await { Ok(_) => (), Err(e) => match e.kind() { ErrorKind::NotFound => { fs::create_dir_all(&p).await?; return Ok(()); } _ => return Err(Error::from(e)), }, } Err(Error::from(DiskError::VolumeExists)) } async fn list_volumes(&self) -> Result> { let mut entries = fs::read_dir(&self.root).await?; let mut volumes = Vec::new(); while let Some(entry) = entries.next_entry().await? { if let Ok(metadata) = entry.metadata().await { // if !metadata.is_dir() { // continue; // } let name = entry.file_name().to_string_lossy().to_string(); let created = match metadata.created() { Ok(md) => OffsetDateTime::from(md), Err(_) => return Err(Error::msg("Not supported created on this platform")), }; volumes.push(VolumeInfo { name, created }); } } Ok(volumes) } async fn stat_volume(&self, volume: &str) -> Result { let p = self.get_bucket_path(volume)?; let m = read_file_metadata(&p).await?; let modtime = match m.modified() { Ok(md) => OffsetDateTime::from(md), Err(_) => return Err(Error::msg("Not supported modified on this platform")), }; Ok(VolumeInfo { name: volume.to_string(), created: modtime, }) } async fn write_metadata(&self, _org_volume: &str, volume: &str, path: &str, fi: FileInfo) -> Result<()> { let p = self.get_object_path(volume, format!("{}/{}", path, super::STORAGE_FORMAT_FILE).as_str())?; warn!("write_metadata {:?} {:?}", &p, &fi); let mut meta = FileMeta::new(); if !fi.fresh { let (buf, _) = read_file_exists(&p).await?; if !buf.is_empty() { meta = FileMeta::unmarshal(&buf)?; } } meta.add_version(fi)?; let fm_data = meta.marshal_msg()?; write_all_internal(p, fm_data).await?; return Ok(()); } async fn read_version( &self, _org_volume: &str, volume: &str, path: &str, version_id: &str, opts: &ReadOptions, ) -> Result { let file_path = self.get_object_path(volume, path)?; let file_dir = self.get_bucket_path(volume)?; let read_data = opts.read_data; let (data, _) = self.read_raw(volume, file_dir, file_path, read_data).await?; let meta = FileMeta::unmarshal(&data)?; let fi = meta.into_fileinfo(volume, path, version_id, false, true)?; Ok(fi) } async fn read_xl(&self, volume: &str, path: &str, read_data: bool) -> Result { let file_path = self.get_object_path(volume, path)?; let file_dir = self.get_bucket_path(volume)?; let (buf, _) = self.read_raw(volume, file_dir, file_path, read_data).await?; Ok(RawFileInfo { buf }) } async fn read_multiple(&self, req: ReadMultipleReq) -> Result> { let mut results = Vec::new(); let mut found = 0; for v in req.files.iter() { let fpath = self.get_object_path(&req.bucket, format!("{}/{}", &req.prefix, v).as_str())?; let mut res = ReadMultipleResp { bucket: req.bucket.clone(), prefix: req.prefix.clone(), file: v.clone(), ..Default::default() }; // if req.metadata_only {} match read_file_all(&fpath).await { Ok((data, meta)) => { found += 1; if req.max_size > 0 && data.len() > req.max_size { res.exists = true; res.error = format!("max size ({}) exceeded: {}", req.max_size, data.len()); results.push(res); break; } res.exists = true; res.data = data; res.mod_time = match meta.modified() { Ok(md) => OffsetDateTime::from(md), Err(_) => return Err(Error::msg("Not supported modified on this platform")), }; results.push(res); if req.max_results > 0 && found >= req.max_results { break; } } Err(e) => { if !(DiskError::FileNotFound.is(&e) || DiskError::VolumeNotFound.is(&e)) { res.exists = true; res.error = e.to_string(); } if req.abort404 && !res.exists { results.push(res); break; } results.push(res); } } } Ok(results) } async fn delete_volume(&self, volume: &str) -> Result<()> { let p = self.get_bucket_path(volume)?; fs::remove_dir_all(&p).await?; Ok(()) } } #[cfg(test)] mod test { use super::*; #[tokio::test] async fn test_skip_access_checks() { // let arr = Vec::new(); let vols = [ super::super::RUSTFS_META_TMP_DELETED_BUCKET, super::super::RUSTFS_META_TMP_BUCKET, super::super::RUSTFS_META_MULTIPART_BUCKET, super::super::RUSTFS_META_BUCKET, ]; let paths: Vec<_> = vols.iter().map(|v| Path::new(v).join("test")).collect(); for p in paths.iter() { assert!(skip_access_checks(p.to_str().unwrap())); } } #[tokio::test] async fn test_make_volume() { let p = "./testv"; fs::create_dir_all(&p).await.unwrap(); let ep = match Endpoint::try_from(p) { Ok(e) => e, Err(e) => { println!("{e}"); return; } }; let disk = LocalDisk::new(&ep, false).await.unwrap(); let tmpp = disk .resolve_abs_path(Path::new(super::super::RUSTFS_META_TMP_DELETED_BUCKET)) .unwrap(); println!("ppp :{:?}", &tmpp); let volumes = vec!["a", "b", "c"]; disk.make_volumes(volumes.clone()).await.unwrap(); disk.make_volumes(volumes.clone()).await.unwrap(); fs::remove_dir_all(&p).await.unwrap(); } #[tokio::test] async fn test_delete_volume() { let p = "./testv"; fs::create_dir_all(&p).await.unwrap(); let ep = match Endpoint::try_from(p) { Ok(e) => e, Err(e) => { println!("{e}"); return; } }; let disk = LocalDisk::new(&ep, false).await.unwrap(); let tmpp = disk .resolve_abs_path(Path::new(super::super::RUSTFS_META_TMP_DELETED_BUCKET)) .unwrap(); println!("ppp :{:?}", &tmpp); let volumes = vec!["a", "b", "c"]; disk.make_volumes(volumes.clone()).await.unwrap(); disk.delete_volume("a").await.unwrap(); fs::remove_dir_all(&p).await.unwrap(); } }