use std::path::{Path, PathBuf};
use std::pin::Pin;
use chrono::DateTime;
use futures::{Stream, TryStreamExt};
use tokio::fs;
use tokio::io::AsyncWriteExt;
use tokio_stream::wrappers::ReadDirStream;
use super::{ObjectMeta, StorageBackend, StorageError};
mod rename;
#[derive(Default, Debug)]
pub struct FileStorageBackend {
root: String,
}
impl FileStorageBackend {
pub fn new(root: &str) -> Self {
Self {
root: String::from(root),
}
}
}
#[async_trait::async_trait]
impl StorageBackend for FileStorageBackend {
fn join_path(&self, path: &str, path_to_join: &str) -> String {
let new_path = Path::new(path);
new_path
.join(path_to_join)
.into_os_string()
.into_string()
.unwrap()
}
fn join_paths(&self, paths: &[&str]) -> String {
let mut iter = paths.iter();
let mut path = PathBuf::from(iter.next().unwrap_or(&""));
iter.for_each(|s| path.push(s));
path.into_os_string().into_string().unwrap()
}
async fn head_obj(&self, path: &str) -> Result<ObjectMeta, StorageError> {
let attr = fs::metadata(path).await?;
Ok(ObjectMeta {
path: path.to_string(),
modified: DateTime::from(attr.modified().unwrap()),
})
}
async fn get_obj(&self, path: &str) -> Result<Vec<u8>, StorageError> {
fs::read(path).await.map_err(StorageError::from)
}
async fn list_objs<'a>(
&'a self,
path: &'a str,
) -> Result<
Pin<Box<dyn Stream<Item = Result<ObjectMeta, StorageError>> + Send + 'a>>,
StorageError,
> {
let readdir = ReadDirStream::new(fs::read_dir(path).await?);
Ok(Box::pin(readdir.err_into().and_then(|entry| async move {
Ok(ObjectMeta {
path: String::from(entry.path().to_str().unwrap()),
modified: DateTime::from(entry.metadata().await.unwrap().modified().unwrap()),
})
})))
}
async fn put_obj(&self, path: &str, obj_bytes: &[u8]) -> Result<(), StorageError> {
if let Some(parent) = Path::new(path).parent() {
fs::create_dir_all(parent).await?;
}
let mut f = fs::OpenOptions::new()
.create_new(true)
.write(true)
.open(path)
.await?;
f.write(obj_bytes).await?;
Ok(())
}
async fn rename_obj(&self, src: &str, dst: &str) -> Result<(), StorageError> {
rename::atomic_rename(src, dst)
}
async fn delete_obj(&self, path: &str) -> Result<(), StorageError> {
fs::remove_file(path).await.map_err(StorageError::from)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn put_and_rename() {
let tmp_dir = tempdir::TempDir::new("rename_test").unwrap();
let backend = FileStorageBackend::new(tmp_dir.path().to_str().unwrap());
let tmp_file_path = tmp_dir.path().join("tmp_file");
let new_file_path = tmp_dir.path().join("new_file");
let tmp_file = tmp_file_path.to_str().unwrap();
let new_file = new_file_path.to_str().unwrap();
backend.put_obj(tmp_file, b"hello").await.unwrap();
if let Err(e) = backend.rename_obj(tmp_file, new_file).await {
panic!("Expect put_obj to return Ok, got Err: {:#?}", e)
}
backend.put_obj(tmp_file, b"hello").await.unwrap();
assert!(matches!(
backend.rename_obj(tmp_file, new_file).await,
Err(StorageError::AlreadyExists(s)) if s == new_file_path.to_str().unwrap(),
));
}
#[tokio::test]
async fn delete_obj() {
let tmp_dir = tempdir::TempDir::new("delete_test").unwrap();
let tmp_file_path = tmp_dir.path().join("tmp_file");
let backend = FileStorageBackend::new(tmp_dir.path().to_str().unwrap());
let path = tmp_file_path.to_str().unwrap();
backend.put_obj(path, &[]).await.unwrap();
assert_eq!(fs::metadata(path).await.is_ok(), true);
backend.delete_obj(path).await.unwrap();
assert_eq!(fs::metadata(path).await.is_ok(), false)
}
#[test]
fn join_multiple_paths() {
let backend = FileStorageBackend::new("./");
assert_eq!(
Path::new(&backend.join_paths(&["abc", "efg/", "123"])),
Path::new("abc").join("efg").join("123"),
);
assert_eq!(
&backend.join_paths(&["abc", "efg"]),
&backend.join_path("abc", "efg"),
);
assert_eq!(&backend.join_paths(&["foo"]), "foo",);
assert_eq!(&backend.join_paths(&[]), "",);
}
}