use bytes::Bytes;
use bytesize::ByteSize;
use dragonfly_api::common::v2::Range;
use dragonfly_client_config::dfdaemon::Config;
use dragonfly_client_config::MIN_PIECE_LENGTH;
use dragonfly_client_core::{Error, Result};
use dragonfly_client_util::buffer_pool::BufferPool;
use dragonfly_client_util::fs::fd::{FDCache, DEFAULT_FD_CACHE_CAPACITY};
use dragonfly_client_util::fs::{fadvise_dontneed, fadvise_willneed, fallocate};
use futures::Stream;
use std::cmp::max;
use std::os::unix::fs::MetadataExt;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use tokio::fs;
use tokio::io::AsyncRead;
use tracing::{error, info, instrument, warn};
use walkdir::WalkDir;
pub struct Content {
pub config: Arc<Config>,
pub dir: PathBuf,
fd_cache: FDCache,
buffer_pool: BufferPool,
writeback: super::content::Writeback,
}
impl Content {
pub async fn new(config: Arc<Config>, dir: &Path) -> Result<Content> {
let dir = dir.join(super::content::DEFAULT_CONTENT_DIR);
if !config.storage.keep {
fs::remove_dir_all(&dir).await.unwrap_or_else(|err| {
warn!("remove {:?} failed: {}", dir, err);
});
}
fs::create_dir_all(&dir.join(super::content::DEFAULT_TASK_DIR)).await?;
fs::create_dir_all(&dir.join(super::content::DEFAULT_PERSISTENT_TASK_DIR)).await?;
fs::create_dir_all(&dir.join(super::content::DEFAULT_PERSISTENT_CACHE_TASK_DIR)).await?;
info!("content initialized directory: {:?}", dir);
Ok(Content {
buffer_pool: BufferPool::new(
super::content::MAX_BUFFER_POOL_IDLE_BUFFERS
* max(
config.storage.write_buffer_size,
config.storage.read_buffer_size,
),
),
writeback: super::content::Writeback::new(config.storage.writeback_mode),
config,
dir,
fd_cache: FDCache::new(DEFAULT_FD_CACHE_CAPACITY),
})
}
pub fn available_space(&self) -> Result<u64> {
let disk_threshold = self.config.gc.policy.disk_threshold;
if disk_threshold != ByteSize::default() {
let usage_space = WalkDir::new(&self.dir)
.into_iter()
.filter_map(|entry| entry.ok())
.filter_map(|entry| entry.metadata().ok())
.filter(|metadata| metadata.is_file())
.fold(0, |acc, m| acc + m.len());
if usage_space >= disk_threshold.as_u64() {
warn!(
"usage space {} is greater than disk threshold {}, no need to calculate available space",
usage_space, disk_threshold
);
return Ok(0);
}
return Ok(disk_threshold.as_u64() - usage_space);
}
let stat = fs2::statvfs(&self.dir)?;
Ok(stat.available_space())
}
pub fn total_space(&self) -> Result<u64> {
let disk_threshold = self.config.gc.policy.disk_threshold;
if disk_threshold != ByteSize::default() {
return Ok(disk_threshold.as_u64());
}
let stat = fs2::statvfs(&self.dir)?;
Ok(stat.total_space())
}
pub fn has_enough_space(&self, content_length: u64) -> Result<bool> {
let available_space = self.available_space()?;
if available_space < content_length {
warn!(
"not enough space to store the task: available_space={}, content_length={}",
available_space, content_length
);
return Ok(false);
}
Ok(true)
}
async fn is_same_dev_inode<P: AsRef<Path>, Q: AsRef<Path>>(
&self,
source: P,
target: Q,
) -> Result<bool> {
let source_metadata = fs::metadata(source).await?;
let target_metadata = fs::metadata(target).await?;
Ok(source_metadata.dev() == target_metadata.dev()
&& source_metadata.ino() == target_metadata.ino())
}
pub async fn is_same_dev_inode_as_task(&self, task_id: &str, to: &Path) -> Result<bool> {
let task_path = self.get_task_path(task_id);
self.is_same_dev_inode(&task_path, to).await
}
#[instrument(level = "debug", skip_all)]
pub async fn create_task(&self, task_id: &str, length: u64) -> Result<PathBuf> {
let task_path = self.get_task_path(task_id);
if task_path.exists() {
return Ok(task_path);
}
let task_dir = self
.dir
.join(super::content::DEFAULT_TASK_DIR)
.join(&task_id[..3]);
fs::create_dir_all(&task_dir).await.inspect_err(|err| {
error!("create {:?} failed: {}", task_dir, err);
})?;
let f = fs::File::create(task_dir.join(task_id))
.await
.inspect_err(|err| {
error!("create {:?} failed: {}", task_dir, err);
})?;
fallocate(&f, length).await.inspect_err(|err| {
error!("fallocate {:?} failed: {}", task_dir, err);
})?;
Ok(task_dir.join(task_id))
}
#[instrument(level = "debug", skip_all)]
pub async fn hard_link_task(&self, task_id: &str, to: &Path) -> Result<()> {
let task_path = self.get_task_path(task_id);
if let Err(err) = fs::hard_link(task_path.clone(), to).await {
if err.kind() == std::io::ErrorKind::AlreadyExists {
if let Ok(true) = self.is_same_dev_inode(&task_path, to).await {
info!("hard already exists, no need to operate");
return Ok(());
}
}
warn!("hard link {:?} to {:?} failed: {}", task_path, to, err);
return Err(Error::IO(err));
}
info!("hard link {:?} to {:?} success", task_path, to);
Ok(())
}
#[instrument(level = "debug", skip_all)]
pub async fn hard_link_to_task(&self, from: &Path, task_id: &str) -> Result<()> {
let task_path = self.get_task_path(task_id);
if let Err(err) = fs::hard_link(from, &task_path).await {
if err.kind() == std::io::ErrorKind::AlreadyExists {
if let Ok(true) = self.is_same_dev_inode(from, &task_path).await {
info!("hard already exists, no need to operate");
return Ok(());
}
}
warn!("hard link {:?} to {:?} failed: {}", from, task_path, err);
return Err(Error::IO(err));
}
info!("hard link {:?} to {:?} success", from, task_path);
Ok(())
}
#[instrument(skip_all)]
pub async fn copy_task(&self, task_id: &str, to: &Path) -> Result<()> {
let length = fs::copy(self.get_task_path(task_id), to).await?;
if let Ok(f) = fs::File::open(to).await {
self.writeback
.trigger(&Arc::new(f.into_std().await), 0, length)
.await;
}
info!("copy to {:?} success", to);
Ok(())
}
pub async fn delete_task(&self, task_id: &str) -> Result<()> {
info!("delete task content: {}", task_id);
let task_path = self.get_task_path(task_id);
self.fd_cache.remove(&task_path).unwrap_or_else(|err| {
error!("remove {:?} from fd_cache failed: {}", task_path, err);
});
fs::remove_file(task_path.as_path())
.await
.inspect_err(|err| {
error!("remove {:?} failed: {}", task_path, err);
})?;
Ok(())
}
#[instrument(level = "debug", skip_all)]
pub async fn fadvise_dontneed_task(&self, task_id: &str) -> Result<()> {
let f = fs::File::open(self.get_task_path(task_id)).await?;
fadvise_dontneed(&f).await
}
#[instrument(level = "debug", skip_all)]
pub async fn read_piece(
&self,
task_id: &str,
offset: u64,
length: u64,
range: Option<Range>,
) -> Result<super::io::RangeReader> {
let task_path = self.get_task_path(task_id);
let (target_offset, target_length) =
super::content::calculate_piece_range(offset, length, range);
let fd = self.fd_cache.open(&task_path).await.inspect_err(|err| {
error!("open {:?} failed: {}", task_path, err);
})?;
if target_length >= MIN_PIECE_LENGTH {
fadvise_willneed(&fd, target_offset, target_length)
.await
.unwrap_or_else(|err| warn!("fadvise_willneed failed: {}", err));
}
Ok(super::io::RangeReader::new(
fd,
target_offset,
target_length,
self.config.storage.read_buffer_size,
self.buffer_pool.clone(),
))
}
#[instrument(level = "debug", skip_all)]
pub async fn write_piece_from_stream<S>(
&self,
task_id: &str,
offset: u64,
expected_length: u64,
stream: &mut S,
) -> Result<super::io::WriteRangeResponse>
where
S: Stream<Item = std::io::Result<Bytes>> + Unpin + ?Sized,
{
let task_path = self.get_task_path(task_id);
let fd = self
.fd_cache
.open_write(&task_path)
.await
.inspect_err(|err| {
error!("open {:?} failed: {}", task_path, err);
})?;
let response = super::io::write_range_from_stream(
fd.clone(),
offset,
expected_length,
self.config.storage.write_buffer_size,
stream,
)
.await
.inspect_err(|err| {
error!("write {:?} failed: {}", task_path, err);
})?;
self.writeback.trigger(&fd, offset, response.length).await;
Ok(response)
}
fn get_task_path(&self, task_id: &str) -> PathBuf {
let sub_dir = &task_id[..3];
self.dir
.join(super::content::DEFAULT_TASK_DIR)
.join(sub_dir)
.join(task_id)
}
pub async fn is_same_dev_inode_as_persistent_task(
&self,
task_id: &str,
to: &Path,
) -> Result<bool> {
let task_path = self.get_persistent_task_path(task_id);
self.is_same_dev_inode(&task_path, to).await
}
#[instrument(level = "debug", skip_all)]
pub async fn create_persistent_task(&self, task_id: &str, length: u64) -> Result<PathBuf> {
let task_path = self.get_persistent_task_path(task_id);
if task_path.exists() {
return Ok(task_path);
}
let task_dir = self
.dir
.join(super::content::DEFAULT_PERSISTENT_TASK_DIR)
.join(&task_id[..3]);
fs::create_dir_all(&task_dir).await.inspect_err(|err| {
error!("create {:?} failed: {}", task_dir, err);
})?;
let f = fs::File::create(task_dir.join(task_id))
.await
.inspect_err(|err| {
error!("create {:?} failed: {}", task_dir, err);
})?;
fallocate(&f, length).await.inspect_err(|err| {
error!("fallocate {:?} failed: {}", task_dir, err);
})?;
Ok(task_dir.join(task_id))
}
#[instrument(level = "debug", skip_all)]
pub async fn create_persistent_task_dir(&self, task_id: &str) -> Result<PathBuf> {
let task_path = self.get_persistent_task_path(task_id);
if task_path.exists() {
return Ok(task_path);
}
let task_dir = self
.dir
.join(super::content::DEFAULT_PERSISTENT_TASK_DIR)
.join(&task_id[..3]);
fs::create_dir_all(&task_dir).await.inspect_err(|err| {
error!("create {:?} failed: {}", task_dir, err);
})?;
Ok(task_dir)
}
#[instrument(level = "debug", skip_all)]
pub async fn hard_link_persistent_task(&self, task_id: &str, to: &Path) -> Result<()> {
let task_path = self.get_persistent_task_path(task_id);
if let Err(err) = fs::hard_link(task_path.clone(), to).await {
if err.kind() == std::io::ErrorKind::AlreadyExists {
if let Ok(true) = self.is_same_dev_inode(&task_path, to).await {
info!("hard already exists, no need to operate");
return Ok(());
}
}
warn!("hard link {:?} to {:?} failed: {}", task_path, to, err);
return Err(Error::IO(err));
}
info!("hard link {:?} to {:?} success", task_path, to);
Ok(())
}
#[instrument(level = "debug", skip_all)]
pub async fn hard_link_to_persistent_task(&self, from: &Path, task_id: &str) -> Result<()> {
let task_path = self.get_persistent_task_path(task_id);
if let Err(err) = fs::hard_link(from, &task_path).await {
if err.kind() == std::io::ErrorKind::AlreadyExists {
if let Ok(true) = self.is_same_dev_inode(from, &task_path).await {
info!("hard already exists, no need to operate");
return Ok(());
}
}
warn!("hard link {:?} to {:?} failed: {}", from, task_path, err);
return Err(Error::IO(err));
}
info!("hard link {:?} to {:?} success", from, task_path);
Ok(())
}
#[instrument(level = "debug", skip_all)]
pub async fn copy_persistent_task(&self, task_id: &str, to: &Path) -> Result<()> {
let length = fs::copy(self.get_persistent_task_path(task_id), to).await?;
if let Ok(f) = fs::File::open(to).await {
self.writeback
.trigger(&Arc::new(f.into_std().await), 0, length)
.await;
}
info!("copy to {:?} success", to);
Ok(())
}
#[instrument(level = "debug", skip_all)]
pub async fn read_persistent_piece(
&self,
task_id: &str,
offset: u64,
length: u64,
range: Option<Range>,
) -> Result<super::io::RangeReader> {
let task_path = self.get_persistent_task_path(task_id);
let (target_offset, target_length) =
super::content::calculate_piece_range(offset, length, range);
let fd = self.fd_cache.open(&task_path).await.inspect_err(|err| {
error!("open {:?} failed: {}", task_path, err);
})?;
if target_length >= MIN_PIECE_LENGTH {
fadvise_willneed(&fd, target_offset, target_length)
.await
.unwrap_or_else(|err| warn!("fadvise_willneed failed: {}", err));
}
Ok(super::io::RangeReader::new(
fd,
target_offset,
target_length,
self.config.storage.read_buffer_size,
self.buffer_pool.clone(),
))
}
#[instrument(level = "debug", skip_all)]
pub async fn write_persistent_piece<R: AsyncRead + Unpin + ?Sized>(
&self,
task_id: &str,
offset: u64,
expected_length: u64,
reader: &mut R,
) -> Result<super::io::WriteRangeResponse> {
let task_path = self.get_persistent_task_path(task_id);
let fd = self
.fd_cache
.open_write(&task_path)
.await
.inspect_err(|err| {
error!("open {:?} failed: {}", task_path, err);
})?;
let response = super::io::write_range(
fd.clone(),
offset,
expected_length,
self.config.storage.write_buffer_size,
reader,
&self.buffer_pool,
)
.await
.inspect_err(|err| {
error!("write {:?} failed: {}", task_path, err);
})?;
self.writeback.trigger(&fd, offset, response.length).await;
Ok(response)
}
#[instrument(level = "debug", skip_all)]
pub async fn write_persistent_piece_from_stream<S>(
&self,
task_id: &str,
offset: u64,
expected_length: u64,
stream: &mut S,
) -> Result<super::io::WriteRangeResponse>
where
S: Stream<Item = std::io::Result<Bytes>> + Unpin + ?Sized,
{
let task_path = self.get_persistent_task_path(task_id);
let fd = self
.fd_cache
.open_write(&task_path)
.await
.inspect_err(|err| {
error!("open {:?} failed: {}", task_path, err);
})?;
let response = super::io::write_range_from_stream(
fd.clone(),
offset,
expected_length,
self.config.storage.write_buffer_size,
stream,
)
.await
.inspect_err(|err| {
error!("write {:?} failed: {}", task_path, err);
})?;
self.writeback.trigger(&fd, offset, response.length).await;
Ok(response)
}
pub async fn delete_persistent_task(&self, task_id: &str) -> Result<()> {
info!("delete persistent task content: {}", task_id);
let persistent_task_path = self.get_persistent_task_path(task_id);
self.fd_cache
.remove(&persistent_task_path)
.unwrap_or_else(|err| {
error!(
"remove {:?} from fd_cache failed: {}",
persistent_task_path, err
);
});
fs::remove_file(persistent_task_path.as_path())
.await
.inspect_err(|err| {
error!("remove {:?} failed: {}", persistent_task_path, err);
})?;
Ok(())
}
#[instrument(level = "debug", skip_all)]
pub async fn fadvise_dontneed_persistent_task(&self, task_id: &str) -> Result<()> {
let f = fs::File::open(self.get_persistent_task_path(task_id)).await?;
fadvise_dontneed(&f).await
}
fn get_persistent_task_path(&self, task_id: &str) -> PathBuf {
self.dir
.join(super::content::DEFAULT_PERSISTENT_TASK_DIR)
.join(&task_id[..3])
.join(task_id)
}
pub async fn is_same_dev_inode_as_persistent_cache_task(
&self,
task_id: &str,
to: &Path,
) -> Result<bool> {
let task_path = self.get_persistent_cache_task_path(task_id);
self.is_same_dev_inode(&task_path, to).await
}
#[instrument(level = "debug", skip_all)]
pub async fn create_persistent_cache_task(
&self,
task_id: &str,
length: u64,
) -> Result<PathBuf> {
let task_path = self.get_persistent_cache_task_path(task_id);
if task_path.exists() {
return Ok(task_path);
}
let task_dir = self
.dir
.join(super::content::DEFAULT_PERSISTENT_CACHE_TASK_DIR)
.join(&task_id[..3]);
fs::create_dir_all(&task_dir).await.inspect_err(|err| {
error!("create {:?} failed: {}", task_dir, err);
})?;
let f = fs::File::create(task_dir.join(task_id))
.await
.inspect_err(|err| {
error!("create {:?} failed: {}", task_dir, err);
})?;
fallocate(&f, length).await.inspect_err(|err| {
error!("fallocate {:?} failed: {}", task_dir, err);
})?;
Ok(task_dir.join(task_id))
}
#[instrument(level = "debug", skip_all)]
pub async fn create_persistent_cache_task_dir(&self, task_id: &str) -> Result<PathBuf> {
let task_path = self.get_persistent_cache_task_path(task_id);
if task_path.exists() {
return Ok(task_path);
}
let task_dir = self
.dir
.join(super::content::DEFAULT_PERSISTENT_CACHE_TASK_DIR)
.join(&task_id[..3]);
fs::create_dir_all(&task_dir).await.inspect_err(|err| {
error!("create {:?} failed: {}", task_dir, err);
})?;
Ok(task_dir)
}
#[instrument(level = "debug", skip_all)]
pub async fn hard_link_persistent_cache_task(&self, task_id: &str, to: &Path) -> Result<()> {
let task_path = self.get_persistent_cache_task_path(task_id);
if let Err(err) = fs::hard_link(task_path.clone(), to).await {
if err.kind() == std::io::ErrorKind::AlreadyExists {
if let Ok(true) = self.is_same_dev_inode(&task_path, to).await {
info!("hard already exists, no need to operate");
return Ok(());
}
}
warn!("hard link {:?} to {:?} failed: {}", task_path, to, err);
return Err(Error::IO(err));
}
info!("hard link {:?} to {:?} success", task_path, to);
Ok(())
}
#[instrument(level = "debug", skip_all)]
pub async fn hard_link_to_persistent_cache_task(
&self,
from: &Path,
task_id: &str,
) -> Result<()> {
let task_path = self.get_persistent_cache_task_path(task_id);
if let Err(err) = fs::hard_link(from, &task_path).await {
if err.kind() == std::io::ErrorKind::AlreadyExists {
if let Ok(true) = self.is_same_dev_inode(from, &task_path).await {
info!("hard already exists, no need to operate");
return Ok(());
}
}
warn!("hard link {:?} to {:?} failed: {}", from, task_path, err);
return Err(Error::IO(err));
}
info!("hard link {:?} to {:?} success", from, task_path);
Ok(())
}
#[instrument(level = "debug", skip_all)]
pub async fn copy_persistent_cache_task(&self, task_id: &str, to: &Path) -> Result<()> {
let length = fs::copy(self.get_persistent_cache_task_path(task_id), to).await?;
if let Ok(f) = fs::File::open(to).await {
self.writeback
.trigger(&Arc::new(f.into_std().await), 0, length)
.await;
}
info!("copy to {:?} success", to);
Ok(())
}
#[instrument(level = "debug", skip_all)]
pub async fn read_persistent_cache_piece(
&self,
task_id: &str,
offset: u64,
length: u64,
range: Option<Range>,
) -> Result<super::io::RangeReader> {
let task_path = self.get_persistent_cache_task_path(task_id);
let (target_offset, target_length) =
super::content::calculate_piece_range(offset, length, range);
let fd = self.fd_cache.open(&task_path).await.inspect_err(|err| {
error!("open {:?} failed: {}", task_path, err);
})?;
if target_length >= MIN_PIECE_LENGTH {
fadvise_willneed(&fd, target_offset, target_length)
.await
.unwrap_or_else(|err| warn!("fadvise_willneed failed: {}", err));
}
Ok(super::io::RangeReader::new(
fd,
target_offset,
target_length,
self.config.storage.read_buffer_size,
self.buffer_pool.clone(),
))
}
#[instrument(level = "debug", skip_all)]
pub async fn write_persistent_cache_piece<R: AsyncRead + Unpin + ?Sized>(
&self,
task_id: &str,
offset: u64,
expected_length: u64,
reader: &mut R,
) -> Result<super::io::WriteRangeResponse> {
let task_path = self.get_persistent_cache_task_path(task_id);
let fd = self
.fd_cache
.open_write(&task_path)
.await
.inspect_err(|err| {
error!("open {:?} failed: {}", task_path, err);
})?;
let response = super::io::write_range(
fd.clone(),
offset,
expected_length,
self.config.storage.write_buffer_size,
reader,
&self.buffer_pool,
)
.await
.inspect_err(|err| {
error!("write {:?} failed: {}", task_path, err);
})?;
self.writeback.trigger(&fd, offset, response.length).await;
Ok(response)
}
#[instrument(level = "debug", skip_all)]
pub async fn write_persistent_cache_piece_from_stream<S>(
&self,
task_id: &str,
offset: u64,
expected_length: u64,
stream: &mut S,
) -> Result<super::io::WriteRangeResponse>
where
S: Stream<Item = std::io::Result<Bytes>> + Unpin + ?Sized,
{
let task_path = self.get_persistent_cache_task_path(task_id);
let fd = self
.fd_cache
.open_write(&task_path)
.await
.inspect_err(|err| {
error!("open {:?} failed: {}", task_path, err);
})?;
let response = super::io::write_range_from_stream(
fd.clone(),
offset,
expected_length,
self.config.storage.write_buffer_size,
stream,
)
.await
.inspect_err(|err| {
error!("write {:?} failed: {}", task_path, err);
})?;
self.writeback.trigger(&fd, offset, response.length).await;
Ok(response)
}
pub async fn delete_persistent_cache_task(&self, task_id: &str) -> Result<()> {
info!("delete persistent cache task content: {}", task_id);
let persistent_cache_task_path = self.get_persistent_cache_task_path(task_id);
self.fd_cache
.remove(&persistent_cache_task_path)
.unwrap_or_else(|err| {
error!(
"remove {:?} from fd_cache failed: {}",
persistent_cache_task_path, err
);
});
fs::remove_file(persistent_cache_task_path.as_path())
.await
.inspect_err(|err| {
error!("remove {:?} failed: {}", persistent_cache_task_path, err);
})?;
Ok(())
}
#[instrument(level = "debug", skip_all)]
pub async fn fadvise_dontneed_persistent_cache_task(&self, task_id: &str) -> Result<()> {
let f = fs::File::open(self.get_persistent_cache_task_path(task_id)).await?;
fadvise_dontneed(&f).await
}
fn get_persistent_cache_task_path(&self, task_id: &str) -> PathBuf {
self.dir
.join(super::content::DEFAULT_PERSISTENT_CACHE_TASK_DIR)
.join(&task_id[..3])
.join(task_id)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::content::DEFAULT_TASK_DIR;
use dragonfly_client_config::dfdaemon::WritebackMode;
use std::io::Cursor;
use tempfile::tempdir;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
async fn content(config: Config, dir: &Path) -> Content {
Content::new(Arc::new(config), dir).await.unwrap()
}
#[tokio::test]
async fn new_wipes_the_content_dir_unless_keep_is_set() {
let test_cases = vec![(false, false), (true, true)];
for (keep, expected_exists) in test_cases {
let temp_dir = tempdir().unwrap();
let task_id = "60409bd0ec44160f44c53c39b3fe1c5fdfb23faded0228c68bee83bc15a200e3";
let task_path = content(Config::default(), temp_dir.path())
.await
.create_task(task_id, 0)
.await
.unwrap();
assert!(task_path.exists());
let mut config = Config::default();
config.storage.keep = keep;
content(config, temp_dir.path()).await;
assert_eq!(task_path.exists(), expected_exists);
}
}
#[tokio::test]
async fn create_task_reuses_the_existing_path() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "60409bd0ec44160f44c53c39b3fe1c5fdfb23faded0228c68bee83bc15a200e3";
let task_path = content.create_task(task_id, 0).await.unwrap();
assert!(task_path.exists());
assert_eq!(task_path, temp_dir.path().join("content/tasks/604/60409bd0ec44160f44c53c39b3fe1c5fdfb23faded0228c68bee83bc15a200e3"));
let task_path_exists = content.create_task(task_id, 0).await.unwrap();
assert_eq!(task_path, task_path_exists);
}
#[tokio::test]
async fn hard_link_task_links_and_accepts_the_existing_link() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "c71d239df91726fc519c6eb72d318ec65820627232b2f796219e87dcf35d0ab4";
content.create_task(task_id, 0).await.unwrap();
let to = temp_dir
.path()
.join("c71d239df91726fc519c6eb72d318ec65820627232b2f796219e87dcf35d0ab4");
content.hard_link_task(task_id, &to).await.unwrap();
assert!(to.exists());
assert!(content
.is_same_dev_inode_as_task(task_id, &to)
.await
.unwrap());
content.hard_link_task(task_id, &to).await.unwrap();
}
#[tokio::test]
async fn copy_task_copies_the_content_to_the_destination() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "bfd3c02fb31a7373e25b405fd5fd3082987ccfbaf210889153af9e65bbf13002";
content.create_task(task_id, 64).await.unwrap();
let to = temp_dir
.path()
.join("bfd3c02fb31a7373e25b405fd5fd3082987ccfbaf210889153af9e65bbf13002");
content.copy_task(task_id, &to).await.unwrap();
assert!(to.exists());
}
#[tokio::test]
async fn delete_task_removes_the_content() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "4e19f03b0fceb38f23ff4f657681472a53ef335db3660ae5494912570b7a2bb7";
let task_path = content.create_task(task_id, 0).await.unwrap();
assert!(task_path.exists());
content.delete_task(task_id).await.unwrap();
assert!(!task_path.exists());
}
#[tokio::test]
async fn write_piece_from_stream_reads_back_in_every_writeback_mode() {
let test_cases = vec![
WritebackMode::Sync,
WritebackMode::Async,
WritebackMode::Off,
];
for writeback_mode in test_cases {
let temp_dir = tempdir().unwrap();
let mut config = Config::default();
config.storage.writeback_mode = writeback_mode;
let content = content(config, temp_dir.path()).await;
let task_id = "60409bd0ec44160f44c53c39b3fe1c5fdfb23faded0228c68bee83bc15a200e3";
content.create_task(task_id, 13).await.unwrap();
let data = b"hello, world!";
let response = content
.write_piece_from_stream(
task_id,
0,
13,
&mut futures::stream::iter([Ok(Bytes::from_static(data))]),
)
.await
.unwrap();
assert_eq!(response.length, 13);
if writeback_mode == WritebackMode::Async {
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
}
let mut reader = content.read_piece(task_id, 0, 13, None).await.unwrap();
let mut buffer = Vec::new();
reader.read_to_end(&mut buffer).await.unwrap();
assert_eq!(buffer, data);
}
}
#[tokio::test]
async fn fadvise_dontneed_task_keeps_content_and_rejects_missing_task() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "9f86d081884c7d659a2feaa0c55ad015a3bf4f1b2b0b822cd15d6c15b0f00a08";
content.create_task(task_id, 13).await.unwrap();
let data = b"hello, world!";
content
.write_piece_from_stream(
task_id,
0,
13,
&mut futures::stream::iter([Ok(Bytes::from_static(data))]),
)
.await
.unwrap();
let test_cases = vec![
(task_id, true),
(
"aaa6963dccfd5b4f60b48845606946cea72084f14ed5cce61ec96e69f80a30f8",
false,
),
];
for (task_id, expected_ok) in test_cases {
assert_eq!(
content.fadvise_dontneed_task(task_id).await.is_ok(),
expected_ok
);
}
let mut reader = content.read_piece(task_id, 0, 13, None).await.unwrap();
let mut buffer = Vec::new();
reader.read_to_end(&mut buffer).await.unwrap();
assert_eq!(buffer, data);
}
#[tokio::test]
async fn fadvise_dontneed_persistent_task_requires_the_content() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "ba7816bf8f01cfea414140de5dae2223b00361a396177a9cb410ff61f20015ad";
content.create_persistent_task(task_id, 13).await.unwrap();
let test_cases = vec![
(task_id, true),
(
"aaa6963dccfd5b4f60b48845606946cea72084f14ed5cce61ec96e69f80a30f8",
false,
),
];
for (task_id, expected_ok) in test_cases {
assert_eq!(
content
.fadvise_dontneed_persistent_task(task_id)
.await
.is_ok(),
expected_ok
);
}
}
#[tokio::test]
async fn fadvise_dontneed_persistent_cache_task_requires_the_content() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824";
content
.create_persistent_cache_task(task_id, 13)
.await
.unwrap();
let test_cases = vec![
(task_id, true),
(
"aaa6963dccfd5b4f60b48845606946cea72084f14ed5cce61ec96e69f80a30f8",
false,
),
];
for (task_id, expected_ok) in test_cases {
assert_eq!(
content
.fadvise_dontneed_persistent_cache_task(task_id)
.await
.is_ok(),
expected_ok
);
}
}
#[tokio::test]
async fn read_piece_reads_the_requested_range() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "c794a3bbae81e06d1c8d362509bdd42a7c105b0fb28d80ffe27f94b8f04fc845";
content.create_task(task_id, 13).await.unwrap();
content
.write_piece_from_stream(
task_id,
0,
13,
&mut futures::stream::iter([Ok(Bytes::from_static(b"hello, world!"))]),
)
.await
.unwrap();
let test_cases = vec![
(None, &b"hello, world!"[..]),
(
Some(Range {
start: 0,
length: 5,
}),
&b"hello"[..],
),
(
Some(Range {
start: 7,
length: 6,
}),
&b"world!"[..],
),
];
for (range, expected) in test_cases {
let mut reader = content.read_piece(task_id, 0, 13, range).await.unwrap();
let mut buffer = Vec::new();
reader.read_to_end(&mut buffer).await.unwrap();
assert_eq!(buffer, expected);
}
}
#[tokio::test]
async fn write_piece_from_stream_returns_length_and_crc32() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "60b48845606946cea72084f14ed5cce61ec96e69f80a30f891a6963dccfd5b4f";
content.create_task(task_id, 4).await.unwrap();
let response = content
.write_piece_from_stream(
task_id,
0,
4,
&mut futures::stream::iter([Ok(Bytes::from_static(b"test"))]),
)
.await
.unwrap();
assert_eq!(response.length, 4);
assert_eq!(response.hash, "3632233996");
}
#[tokio::test]
async fn create_persistent_task_reuses_the_existing_path() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "c4f108ab1d2b8cfdffe89ea9676af35123fa02e3c25167d62538f630d5d44745";
let task_path = content.create_persistent_task(task_id, 0).await.unwrap();
assert!(task_path.exists());
assert_eq!(task_path, temp_dir.path().join("content/persistent-tasks/c4f/c4f108ab1d2b8cfdffe89ea9676af35123fa02e3c25167d62538f630d5d44745"));
let task_path_exists = content.create_persistent_task(task_id, 0).await.unwrap();
assert_eq!(task_path, task_path_exists);
}
#[tokio::test]
async fn hard_link_persistent_task_links_and_accepts_the_existing_link() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "5e81970eb2b048910cc84cab026b951f2ceac0a09c72c0717193bb6e466e11cd";
content.create_persistent_task(task_id, 0).await.unwrap();
let to = temp_dir
.path()
.join("5e81970eb2b048910cc84cab026b951f2ceac0a09c72c0717193bb6e466e11cd");
content
.hard_link_persistent_task(task_id, &to)
.await
.unwrap();
assert!(to.exists());
assert!(content
.is_same_dev_inode_as_persistent_task(task_id, &to)
.await
.unwrap());
content
.hard_link_persistent_task(task_id, &to)
.await
.unwrap();
}
#[tokio::test]
async fn hard_link_to_persistent_task_links_a_single_source() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "d7a8fbb307d7809469ca9abcb0082e4f8d5651e46d3cdb762d02d0bf37c9e592";
let task_dir = content.create_persistent_task_dir(task_id).await.unwrap();
assert_eq!(
task_dir,
temp_dir.path().join("content/persistent-tasks/d7a")
);
let from = temp_dir.path().join("source");
fs::write(&from, b"hello, world!").await.unwrap();
content
.hard_link_to_persistent_task(&from, task_id)
.await
.unwrap();
assert!(content
.is_same_dev_inode_as_persistent_task(task_id, &from)
.await
.unwrap());
content
.hard_link_to_persistent_task(&from, task_id)
.await
.unwrap();
let other_source = temp_dir.path().join("other-source");
fs::write(&other_source, b"other").await.unwrap();
assert!(content
.hard_link_to_persistent_task(&other_source, task_id)
.await
.is_err());
}
#[tokio::test]
async fn copy_persistent_task_copies_the_content_to_the_destination() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "194b9c2018429689fb4e596a506c7e9db564c187b9709b55b33b96881dfb6dd5";
content.create_persistent_task(task_id, 64).await.unwrap();
let to = temp_dir
.path()
.join("194b9c2018429689fb4e596a506c7e9db564c187b9709b55b33b96881dfb6dd5");
content.copy_persistent_task(task_id, &to).await.unwrap();
assert!(to.exists());
}
#[tokio::test]
async fn delete_persistent_task_removes_the_content() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "17430ba545c3ce82790e9c9f77e64dca44bb6d6a0c9e18be175037c16c73713d";
let task_path = content.create_persistent_task(task_id, 0).await.unwrap();
assert!(task_path.exists());
content.delete_persistent_task(task_id).await.unwrap();
assert!(!task_path.exists());
}
#[tokio::test]
async fn read_persistent_piece_reads_the_requested_range() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "9cb27a4af09aee4eb9f904170217659683f4a0ea7cd55e1a9fbcb99ddced659a";
content.create_persistent_task(task_id, 13).await.unwrap();
content
.write_persistent_piece(task_id, 0, 13, &mut Cursor::new(b"hello, world!"))
.await
.unwrap();
let test_cases = vec![
(None, &b"hello, world!"[..]),
(
Some(Range {
start: 0,
length: 5,
}),
&b"hello"[..],
),
(
Some(Range {
start: 7,
length: 6,
}),
&b"world!"[..],
),
];
for (range, expected) in test_cases {
let mut reader = content
.read_persistent_piece(task_id, 0, 13, range)
.await
.unwrap();
let mut buffer = Vec::new();
reader.read_to_end(&mut buffer).await.unwrap();
assert_eq!(buffer, expected);
}
}
#[tokio::test]
async fn write_persistent_piece_hashes_reader_and_stream_input() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "ca1afaf856e8a667fbd48093ca3ca1b8eeb4bf735912fbe551676bc5817a720a";
content.create_persistent_task(task_id, 4).await.unwrap();
let response = content
.write_persistent_piece(task_id, 0, 4, &mut Cursor::new(b"test"))
.await
.unwrap();
assert_eq!(response.length, 4);
assert_eq!(response.hash, "3632233996");
let response = content
.write_persistent_piece_from_stream(
task_id,
0,
4,
&mut futures::stream::iter([Ok(Bytes::from_static(b"test"))]),
)
.await
.unwrap();
assert_eq!(response.length, 4);
assert_eq!(response.hash, "3632233996");
}
#[tokio::test]
async fn create_persistent_cache_task_reuses_the_existing_path() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "c4f108ab1d2b8cfdffe89ea9676af35123fa02e3c25167d62538f630d5d44745";
let task_path = content
.create_persistent_cache_task(task_id, 0)
.await
.unwrap();
assert!(task_path.exists());
assert_eq!(task_path, temp_dir.path().join("content/persistent-cache-tasks/c4f/c4f108ab1d2b8cfdffe89ea9676af35123fa02e3c25167d62538f630d5d44745"));
let task_path_exists = content
.create_persistent_cache_task(task_id, 0)
.await
.unwrap();
assert_eq!(task_path, task_path_exists);
}
#[tokio::test]
async fn hard_link_persistent_cache_task_links_and_accepts_the_existing_link() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "5e81970eb2b048910cc84cab026b951f2ceac0a09c72c0717193bb6e466e11cd";
content
.create_persistent_cache_task(task_id, 0)
.await
.unwrap();
let to = temp_dir
.path()
.join("5e81970eb2b048910cc84cab026b951f2ceac0a09c72c0717193bb6e466e11cd");
content
.hard_link_persistent_cache_task(task_id, &to)
.await
.unwrap();
assert!(to.exists());
assert!(content
.is_same_dev_inode_as_persistent_cache_task(task_id, &to)
.await
.unwrap());
content
.hard_link_persistent_cache_task(task_id, &to)
.await
.unwrap();
}
#[tokio::test]
async fn hard_link_to_persistent_cache_task_links_a_single_source() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "ef537f25c895bfa782526529a9b63d97aa631564d5d789c2b765448c8635fb6c";
let task_dir = content
.create_persistent_cache_task_dir(task_id)
.await
.unwrap();
assert_eq!(
task_dir,
temp_dir.path().join("content/persistent-cache-tasks/ef5")
);
let from = temp_dir.path().join("source");
fs::write(&from, b"hello, world!").await.unwrap();
content
.hard_link_to_persistent_cache_task(&from, task_id)
.await
.unwrap();
assert!(content
.is_same_dev_inode_as_persistent_cache_task(task_id, &from)
.await
.unwrap());
content
.hard_link_to_persistent_cache_task(&from, task_id)
.await
.unwrap();
let other_source = temp_dir.path().join("other-source");
fs::write(&other_source, b"other").await.unwrap();
assert!(content
.hard_link_to_persistent_cache_task(&other_source, task_id)
.await
.is_err());
}
#[tokio::test]
async fn copy_persistent_cache_task_copies_the_content_to_the_destination() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "194b9c2018429689fb4e596a506c7e9db564c187b9709b55b33b96881dfb6dd5";
content
.create_persistent_cache_task(task_id, 64)
.await
.unwrap();
let to = temp_dir
.path()
.join("194b9c2018429689fb4e596a506c7e9db564c187b9709b55b33b96881dfb6dd5");
content
.copy_persistent_cache_task(task_id, &to)
.await
.unwrap();
assert!(to.exists());
}
#[tokio::test]
async fn delete_persistent_cache_task_removes_the_content() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "17430ba545c3ce82790e9c9f77e64dca44bb6d6a0c9e18be175037c16c73713d";
let task_path = content
.create_persistent_cache_task(task_id, 0)
.await
.unwrap();
assert!(task_path.exists());
content.delete_persistent_cache_task(task_id).await.unwrap();
assert!(!task_path.exists());
}
#[tokio::test]
async fn read_persistent_cache_piece_reads_the_requested_range() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "9cb27a4af09aee4eb9f904170217659683f4a0ea7cd55e1a9fbcb99ddced659a";
content
.create_persistent_cache_task(task_id, 13)
.await
.unwrap();
content
.write_persistent_cache_piece(task_id, 0, 13, &mut Cursor::new(b"hello, world!"))
.await
.unwrap();
let test_cases = vec![
(None, &b"hello, world!"[..]),
(
Some(Range {
start: 0,
length: 5,
}),
&b"hello"[..],
),
(
Some(Range {
start: 7,
length: 6,
}),
&b"world!"[..],
),
];
for (range, expected) in test_cases {
let mut reader = content
.read_persistent_cache_piece(task_id, 0, 13, range)
.await
.unwrap();
let mut buffer = Vec::new();
reader.read_to_end(&mut buffer).await.unwrap();
assert_eq!(buffer, expected);
}
}
#[tokio::test]
async fn write_persistent_cache_piece_hashes_reader_and_stream_input() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let task_id = "ca1afaf856e8a667fbd48093ca3ca1b8eeb4bf735912fbe551676bc5817a720a";
content
.create_persistent_cache_task(task_id, 4)
.await
.unwrap();
let response = content
.write_persistent_cache_piece(task_id, 0, 4, &mut Cursor::new(b"test"))
.await
.unwrap();
assert_eq!(response.length, 4);
assert_eq!(response.hash, "3632233996");
let response = content
.write_persistent_cache_piece_from_stream(
task_id,
0,
4,
&mut futures::stream::iter([Ok(Bytes::from_static(b"test"))]),
)
.await
.unwrap();
assert_eq!(response.length, 4);
assert_eq!(response.hash, "3632233996");
}
#[tokio::test]
async fn has_enough_space_compares_against_the_free_space() {
let temp_dir = tempdir().unwrap();
let content = content(Config::default(), temp_dir.path()).await;
let test_cases = vec![(1, true), (u64::MAX, false)];
for (content_length, expected) in test_cases {
assert_eq!(content.has_enough_space(content_length).unwrap(), expected);
}
}
#[tokio::test]
async fn space_accounting_uses_disk_threshold_minus_usage() {
let test_cases = vec![
(
ByteSize::mib(10),
ByteSize::mib(9).as_u64() + 1,
ByteSize::mib(9).as_u64(),
false,
),
(
ByteSize::mib(10),
ByteSize::mib(9).as_u64(),
ByteSize::mib(9).as_u64(),
true,
),
(ByteSize::mib(1), 1, 0, false),
];
for (disk_threshold, content_length, expected_available_space, expected_enough) in
test_cases
{
let temp_dir = tempdir().unwrap();
let mut config = Config::default();
config.gc.policy.disk_threshold = disk_threshold;
let content = content(config, temp_dir.path()).await;
let mut file = fs::File::create(content.dir.join(DEFAULT_TASK_DIR).join("1mib"))
.await
.unwrap();
file.write_all(&vec![0u8; ByteSize::mib(1).as_u64() as usize])
.await
.unwrap();
file.flush().await.unwrap();
assert_eq!(content.total_space().unwrap(), disk_threshold.as_u64());
assert_eq!(content.available_space().unwrap(), expected_available_space);
assert_eq!(
content.has_enough_space(content_length).unwrap(),
expected_enough
);
}
}
}