2022-03-26 21:49:23 +00:00
|
|
|
use crate::{
|
2022-03-27 01:45:12 +00:00
|
|
|
error::Error,
|
2022-03-26 21:49:23 +00:00
|
|
|
file::File,
|
|
|
|
repo::{Repo, SettingsRepo},
|
2022-09-24 22:18:53 +00:00
|
|
|
store::{Store, StoreConfig},
|
2022-03-26 21:49:23 +00:00
|
|
|
};
|
2021-10-23 04:48:56 +00:00
|
|
|
use actix_web::web::Bytes;
|
|
|
|
use futures_util::stream::Stream;
|
|
|
|
use std::{
|
|
|
|
path::{Path, PathBuf},
|
|
|
|
pin::Pin,
|
|
|
|
};
|
|
|
|
use storage_path_generator::Generator;
|
|
|
|
use tokio::io::{AsyncRead, AsyncWrite};
|
2022-09-24 22:18:53 +00:00
|
|
|
use tokio_util::io::StreamReader;
|
2022-04-07 02:40:49 +00:00
|
|
|
use tracing::{debug, error, instrument, Instrument};
|
2021-10-23 04:48:56 +00:00
|
|
|
|
|
|
|
mod file_id;
|
|
|
|
pub(crate) use file_id::FileId;
|
|
|
|
|
|
|
|
// - Settings Tree
|
|
|
|
// - last-path -> last generated path
|
|
|
|
|
2022-04-01 16:51:46 +00:00
|
|
|
const GENERATOR_KEY: &str = "last-path";
|
2021-10-23 04:48:56 +00:00
|
|
|
|
|
|
|
#[derive(Debug, thiserror::Error)]
|
|
|
|
pub(crate) enum FileError {
|
2022-03-26 21:49:23 +00:00
|
|
|
#[error("Failed to read or write file")]
|
2021-10-23 04:48:56 +00:00
|
|
|
Io(#[from] std::io::Error),
|
|
|
|
|
2022-03-26 21:49:23 +00:00
|
|
|
#[error("Failed to generate path")]
|
2021-10-23 04:48:56 +00:00
|
|
|
PathGenerator(#[from] storage_path_generator::PathError),
|
|
|
|
|
|
|
|
#[error("Error formatting file store identifier")]
|
|
|
|
IdError,
|
|
|
|
|
|
|
|
#[error("Mailformed file store identifier")]
|
|
|
|
PrefixError,
|
|
|
|
|
|
|
|
#[error("Tried to save over existing file")]
|
|
|
|
FileExists,
|
|
|
|
}
|
|
|
|
|
|
|
|
#[derive(Clone)]
|
|
|
|
pub(crate) struct FileStore {
|
|
|
|
path_gen: Generator,
|
|
|
|
root_dir: PathBuf,
|
2022-03-26 21:49:23 +00:00
|
|
|
repo: Repo,
|
2021-10-23 04:48:56 +00:00
|
|
|
}
|
|
|
|
|
2022-09-24 22:18:53 +00:00
|
|
|
impl StoreConfig for FileStore {
|
|
|
|
type Store = FileStore;
|
|
|
|
|
|
|
|
fn build(self) -> Self::Store {
|
|
|
|
self
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2021-10-23 04:48:56 +00:00
|
|
|
#[async_trait::async_trait(?Send)]
|
|
|
|
impl Store for FileStore {
|
|
|
|
type Identifier = FileId;
|
|
|
|
type Stream = Pin<Box<dyn Stream<Item = std::io::Result<Bytes>>>>;
|
|
|
|
|
2021-10-29 01:59:11 +00:00
|
|
|
#[tracing::instrument(skip(reader))]
|
2022-09-24 22:18:53 +00:00
|
|
|
async fn save_async_read<Reader>(&self, mut reader: Reader) -> Result<Self::Identifier, Error>
|
2021-10-23 04:48:56 +00:00
|
|
|
where
|
2022-09-24 22:18:53 +00:00
|
|
|
Reader: AsyncRead + Unpin + 'static,
|
2021-10-23 04:48:56 +00:00
|
|
|
{
|
2022-03-26 21:49:23 +00:00
|
|
|
let path = self.next_file().await?;
|
2021-10-23 04:48:56 +00:00
|
|
|
|
2022-09-24 22:18:53 +00:00
|
|
|
if let Err(e) = self.safe_save_reader(&path, &mut reader).await {
|
2021-10-23 04:48:56 +00:00
|
|
|
self.safe_remove_file(&path).await?;
|
2022-03-27 01:45:12 +00:00
|
|
|
return Err(e.into());
|
2021-10-23 04:48:56 +00:00
|
|
|
}
|
|
|
|
|
2022-03-27 01:45:12 +00:00
|
|
|
Ok(self.file_id_from_path(path)?)
|
2021-10-23 04:48:56 +00:00
|
|
|
}
|
|
|
|
|
2022-09-24 22:18:53 +00:00
|
|
|
async fn save_stream<S>(&self, stream: S) -> Result<Self::Identifier, Error>
|
|
|
|
where
|
|
|
|
S: Stream<Item = std::io::Result<Bytes>> + Unpin + 'static,
|
|
|
|
{
|
|
|
|
self.save_async_read(StreamReader::new(stream)).await
|
|
|
|
}
|
|
|
|
|
2021-10-29 01:59:11 +00:00
|
|
|
#[tracing::instrument(skip(bytes))]
|
2022-03-27 01:45:12 +00:00
|
|
|
async fn save_bytes(&self, bytes: Bytes) -> Result<Self::Identifier, Error> {
|
2022-03-26 21:49:23 +00:00
|
|
|
let path = self.next_file().await?;
|
2021-10-23 04:48:56 +00:00
|
|
|
|
|
|
|
if let Err(e) = self.safe_save_bytes(&path, bytes).await {
|
|
|
|
self.safe_remove_file(&path).await?;
|
2022-03-27 01:45:12 +00:00
|
|
|
return Err(e.into());
|
2021-10-23 04:48:56 +00:00
|
|
|
}
|
|
|
|
|
2022-03-27 01:45:12 +00:00
|
|
|
Ok(self.file_id_from_path(path)?)
|
2021-10-23 04:48:56 +00:00
|
|
|
}
|
|
|
|
|
2021-10-29 01:59:11 +00:00
|
|
|
#[tracing::instrument]
|
2021-10-23 04:48:56 +00:00
|
|
|
async fn to_stream(
|
|
|
|
&self,
|
|
|
|
identifier: &Self::Identifier,
|
|
|
|
from_start: Option<u64>,
|
|
|
|
len: Option<u64>,
|
2022-03-27 01:45:12 +00:00
|
|
|
) -> Result<Self::Stream, Error> {
|
2021-10-23 04:48:56 +00:00
|
|
|
let path = self.path_from_file_id(identifier);
|
|
|
|
|
2022-04-07 02:40:49 +00:00
|
|
|
let file_span = tracing::trace_span!(parent: None, "File Stream");
|
|
|
|
let file = file_span
|
|
|
|
.in_scope(|| File::open(path))
|
|
|
|
.instrument(file_span.clone())
|
|
|
|
.await?;
|
|
|
|
|
|
|
|
let stream = file_span
|
|
|
|
.in_scope(|| file.read_to_stream(from_start, len))
|
|
|
|
.instrument(file_span)
|
2021-10-23 04:48:56 +00:00
|
|
|
.await?;
|
|
|
|
|
|
|
|
Ok(Box::pin(stream))
|
|
|
|
}
|
|
|
|
|
2021-10-29 01:59:11 +00:00
|
|
|
#[tracing::instrument(skip(writer))]
|
2021-10-23 04:48:56 +00:00
|
|
|
async fn read_into<Writer>(
|
|
|
|
&self,
|
|
|
|
identifier: &Self::Identifier,
|
|
|
|
writer: &mut Writer,
|
|
|
|
) -> Result<(), std::io::Error>
|
|
|
|
where
|
2022-09-24 22:18:53 +00:00
|
|
|
Writer: AsyncWrite + Unpin,
|
2021-10-23 04:48:56 +00:00
|
|
|
{
|
|
|
|
let path = self.path_from_file_id(identifier);
|
|
|
|
|
|
|
|
File::open(&path).await?.read_to_async_write(writer).await?;
|
|
|
|
|
|
|
|
Ok(())
|
|
|
|
}
|
|
|
|
|
2021-10-29 01:59:11 +00:00
|
|
|
#[tracing::instrument]
|
2022-03-27 01:45:12 +00:00
|
|
|
async fn len(&self, identifier: &Self::Identifier) -> Result<u64, Error> {
|
2021-10-23 04:48:56 +00:00
|
|
|
let path = self.path_from_file_id(identifier);
|
|
|
|
|
|
|
|
let len = tokio::fs::metadata(path).await?.len();
|
|
|
|
|
|
|
|
Ok(len)
|
|
|
|
}
|
|
|
|
|
2021-10-29 01:59:11 +00:00
|
|
|
#[tracing::instrument]
|
2022-03-27 01:45:12 +00:00
|
|
|
async fn remove(&self, identifier: &Self::Identifier) -> Result<(), Error> {
|
2021-10-23 04:48:56 +00:00
|
|
|
let path = self.path_from_file_id(identifier);
|
|
|
|
|
|
|
|
self.safe_remove_file(path).await?;
|
|
|
|
|
|
|
|
Ok(())
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
impl FileStore {
|
2022-03-27 01:45:12 +00:00
|
|
|
pub(crate) async fn build(root_dir: PathBuf, repo: Repo) -> Result<Self, Error> {
|
2022-03-26 21:49:23 +00:00
|
|
|
let path_gen = init_generator(&repo).await?;
|
2021-10-23 04:48:56 +00:00
|
|
|
|
|
|
|
Ok(FileStore {
|
|
|
|
root_dir,
|
|
|
|
path_gen,
|
2022-03-26 21:49:23 +00:00
|
|
|
repo,
|
2021-10-23 04:48:56 +00:00
|
|
|
})
|
|
|
|
}
|
|
|
|
|
2022-03-27 01:45:12 +00:00
|
|
|
async fn next_directory(&self) -> Result<PathBuf, Error> {
|
2021-10-23 04:48:56 +00:00
|
|
|
let path = self.path_gen.next();
|
|
|
|
|
2022-03-26 21:49:23 +00:00
|
|
|
match self.repo {
|
|
|
|
Repo::Sled(ref sled_repo) => {
|
|
|
|
sled_repo
|
|
|
|
.set(GENERATOR_KEY, path.to_be_bytes().into())
|
|
|
|
.await?;
|
|
|
|
}
|
|
|
|
}
|
2021-10-23 04:48:56 +00:00
|
|
|
|
2022-03-25 23:47:50 +00:00
|
|
|
let mut target_path = self.root_dir.clone();
|
2021-10-23 04:48:56 +00:00
|
|
|
for dir in path.to_strings() {
|
|
|
|
target_path.push(dir)
|
|
|
|
}
|
|
|
|
|
|
|
|
Ok(target_path)
|
|
|
|
}
|
|
|
|
|
2022-03-27 01:45:12 +00:00
|
|
|
async fn next_file(&self) -> Result<PathBuf, Error> {
|
2022-03-26 21:49:23 +00:00
|
|
|
let target_path = self.next_directory().await?;
|
|
|
|
let filename = uuid::Uuid::new_v4().to_string();
|
2021-10-23 04:48:56 +00:00
|
|
|
|
|
|
|
Ok(target_path.join(filename))
|
|
|
|
}
|
|
|
|
|
|
|
|
async fn safe_remove_file<P: AsRef<Path>>(&self, path: P) -> Result<(), FileError> {
|
|
|
|
tokio::fs::remove_file(&path).await?;
|
|
|
|
self.try_remove_parents(path.as_ref()).await;
|
|
|
|
Ok(())
|
|
|
|
}
|
|
|
|
|
|
|
|
async fn try_remove_parents(&self, mut path: &Path) {
|
|
|
|
while let Some(parent) = path.parent() {
|
|
|
|
if parent.ends_with(&self.root_dir) {
|
|
|
|
return;
|
|
|
|
}
|
|
|
|
|
|
|
|
if tokio::fs::remove_dir(parent).await.is_err() {
|
|
|
|
return;
|
|
|
|
}
|
|
|
|
|
|
|
|
path = parent;
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
// Try writing to a file
|
|
|
|
#[instrument(name = "Saving file", skip(bytes), fields(path = tracing::field::debug(&path.as_ref())))]
|
|
|
|
async fn safe_save_bytes<P: AsRef<Path>>(
|
|
|
|
&self,
|
|
|
|
path: P,
|
|
|
|
bytes: Bytes,
|
|
|
|
) -> Result<(), FileError> {
|
|
|
|
safe_create_parent(&path).await?;
|
|
|
|
|
|
|
|
// Only write the file if it doesn't already exist
|
|
|
|
debug!("Checking if {:?} already exists", path.as_ref());
|
|
|
|
if let Err(e) = tokio::fs::metadata(&path).await {
|
|
|
|
if e.kind() != std::io::ErrorKind::NotFound {
|
|
|
|
return Err(e.into());
|
|
|
|
}
|
|
|
|
} else {
|
|
|
|
return Ok(());
|
|
|
|
}
|
|
|
|
|
|
|
|
// Open the file for writing
|
|
|
|
debug!("Creating {:?}", path.as_ref());
|
|
|
|
let mut file = File::create(&path).await?;
|
|
|
|
|
|
|
|
// try writing
|
|
|
|
debug!("Writing to {:?}", path.as_ref());
|
|
|
|
if let Err(e) = file.write_from_bytes(bytes).await {
|
|
|
|
error!("Error writing {:?}, {}", path.as_ref(), e);
|
|
|
|
// remove file if writing failed before completion
|
|
|
|
self.safe_remove_file(path).await?;
|
|
|
|
return Err(e.into());
|
|
|
|
}
|
|
|
|
debug!("{:?} written", path.as_ref());
|
|
|
|
|
|
|
|
Ok(())
|
|
|
|
}
|
|
|
|
|
|
|
|
#[instrument(skip(input), fields(to = tracing::field::debug(&to.as_ref())))]
|
|
|
|
async fn safe_save_reader<P: AsRef<Path>>(
|
|
|
|
&self,
|
|
|
|
to: P,
|
|
|
|
input: &mut (impl AsyncRead + Unpin + ?Sized),
|
|
|
|
) -> Result<(), FileError> {
|
|
|
|
safe_create_parent(&to).await?;
|
|
|
|
|
|
|
|
debug!("Checking if {:?} already exists", to.as_ref());
|
|
|
|
if let Err(e) = tokio::fs::metadata(&to).await {
|
|
|
|
if e.kind() != std::io::ErrorKind::NotFound {
|
|
|
|
return Err(e.into());
|
|
|
|
}
|
|
|
|
} else {
|
|
|
|
return Err(FileError::FileExists);
|
|
|
|
}
|
|
|
|
|
|
|
|
debug!("Writing stream to {:?}", to.as_ref());
|
|
|
|
|
|
|
|
let mut file = File::create(to).await?;
|
|
|
|
|
|
|
|
file.write_from_async_read(input).await?;
|
|
|
|
|
|
|
|
Ok(())
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
pub(crate) async fn safe_create_parent<P: AsRef<Path>>(path: P) -> Result<(), FileError> {
|
|
|
|
if let Some(path) = path.as_ref().parent() {
|
|
|
|
debug!("Creating directory {:?}", path);
|
|
|
|
tokio::fs::create_dir_all(path).await?;
|
|
|
|
}
|
|
|
|
|
|
|
|
Ok(())
|
|
|
|
}
|
|
|
|
|
2022-03-27 01:45:12 +00:00
|
|
|
async fn init_generator(repo: &Repo) -> Result<Generator, Error> {
|
2022-03-26 21:49:23 +00:00
|
|
|
match repo {
|
|
|
|
Repo::Sled(sled_repo) => {
|
|
|
|
if let Some(ivec) = sled_repo.get(GENERATOR_KEY).await? {
|
|
|
|
Ok(Generator::from_existing(
|
|
|
|
storage_path_generator::Path::from_be_bytes(ivec.to_vec())?,
|
|
|
|
))
|
|
|
|
} else {
|
|
|
|
Ok(Generator::new())
|
|
|
|
}
|
|
|
|
}
|
2021-10-23 04:48:56 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
impl std::fmt::Debug for FileStore {
|
|
|
|
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
|
|
|
f.debug_struct("FileStore")
|
2022-03-22 02:43:38 +00:00
|
|
|
.field("path_gen", &"generator")
|
2021-10-23 04:48:56 +00:00
|
|
|
.field("root_dir", &self.root_dir)
|
|
|
|
.finish()
|
|
|
|
}
|
|
|
|
}
|