use crate::interface::{
cli::{inspect::InspectCommands, server::ServerCommands, x::*},
config::Config,
};
use clap::{Args, Parser, Subcommand};
use eyre::Context;
use futures::StreamExt;
use std::{collections::HashMap, io::Write, path::PathBuf, sync::Arc};
use tokio::{
fs::{self, File, OpenOptions},
io::{AsyncWriteExt, stdout},
sync::RwLock,
};
use xhfs_core::{
device::{ConcreteDevice, disk::Controller, kv_device::*, logical::LogicalDevice},
xhfs::*,
};
mod inspect;
mod server;
mod x;
#[derive(Parser, Debug)]
#[command(
author = "michael-0acf4",
version,
about = "XHFS distributed File System"
)]
pub struct MainCommand {
#[command(subcommand)]
pub command: Commands,
}
#[derive(Subcommand, Debug)]
pub enum Commands {
Upload(UploadCommand),
Download(DownloadCommand),
Read(ReadCommand),
Format(GlobalOptions),
Info(GlobalOptions),
X(FsCommands),
Inspect(InspectCommands),
Server(ServerCommands),
}
#[derive(Args, Debug, Clone)]
pub struct GlobalOptions {
#[arg(long)]
pub config: Option<PathBuf>,
#[arg(long)]
pub password: bool,
#[arg(short, long)]
pub verbose: bool,
#[arg(short, long)]
pub force: bool,
#[arg(long, default_value = "false")]
pub dev: bool,
#[arg(long, default_value = "134217728")]
pub dev_unit_capacity: usize,
#[arg(long, default_value = "1")]
pub dev_replica_count: u8,
#[arg(long, default_value = "1")]
pub dev_logical_count: usize,
}
#[derive(Args, Debug)]
pub struct UploadCommand {
pub src_path: PathBuf,
pub dest_path: Option<PathBuf>,
#[arg(short, long, default_value = "false")]
pub overwrite: bool,
#[arg(short, long, default_value = "1048576")]
pub block_size: usize,
#[arg(long, default_value = "false")]
pub single_block: bool,
#[command(flatten)]
pub global: GlobalOptions,
}
#[derive(Args, Debug)]
pub struct DownloadCommand {
pub src_path: PathBuf,
pub dest_path: Option<PathBuf>,
#[arg(short, long, default_value = "false")]
pub overwrite: bool,
#[arg(short, long, default_value = "1048576")]
pub block_size: usize,
#[arg(short, long, default_value = "false")]
pub single_block: bool,
#[command(flatten)]
pub global: GlobalOptions,
}
#[derive(Args, Debug)]
pub struct ReadCommand {
pub src_path: PathBuf,
#[arg(short, long, default_value = "false")]
pub overwrite: bool,
#[arg(short, long, default_value = "1048576")]
pub block_size: usize,
#[arg(short, long, default_value = "false")]
pub single_block: bool,
#[command(flatten)]
pub global: GlobalOptions,
}
impl MainCommand {
pub async fn run(&self) -> eyre::Result<()> {
match &self.command {
Commands::Format(global_options) => {
if global_options.force || confirm_destructive_action() {
let xhfs = global_options.format_and_get_xhfs().await?;
println!("Formatted {} Bytes", xhfs.total_capacity()?);
} else {
println!("abort");
}
}
Commands::Upload(u) => {
let xhfs = u.global.get_xhfs().await?;
let data = fs::read(&u.src_path).await?;
let dest_path = match &u.dest_path {
Some(path) => path.to_owned(),
None => PathBuf::from(u.src_path.file_name().ok_or_else(|| {
eyre::eyre!(
"Could not derive destination path from source {}",
u.src_path.display()
)
})?),
};
if u.single_block {
xhfs.fwrite(
&dest_path,
data,
WriteOption {
overwrite: u.overwrite,
},
)
.await?;
} else {
let file = File::open(&u.src_path).await?;
xhfs.fwrite_stream(
&dest_path,
file,
u.block_size,
WriteOption {
overwrite: u.overwrite,
},
)
.await?;
}
}
Commands::Download(d) => {
let xhfs = d.global.get_xhfs().await?;
let dest_path = match &d.dest_path {
Some(path) => path.clone(),
None => d.src_path.file_name().expect("Valid filename").into(),
};
if dest_path.exists() && !d.overwrite {
eyre::bail!("File {} already exists", dest_path.display());
}
if d.single_block {
let data = xhfs.fread(&d.src_path).await?;
tokio::fs::write(&dest_path, data).await?;
} else {
let mut data = xhfs.fread_stream(&d.src_path, d.block_size).await?;
let mut file = OpenOptions::new()
.write(true)
.create(true)
.truncate(d.overwrite)
.create_new(!d.overwrite)
.open(&dest_path)
.await?;
while let Some(chunk) = data.next().await {
let chunk = chunk?;
file.write_all(&chunk).await?;
}
file.flush().await?;
}
}
Commands::Read(r) => {
let xhfs = r.global.get_xhfs().await?;
let mut out = stdout();
if r.single_block {
let data = xhfs.fread(&r.src_path).await?;
out.write_all(&data).await?;
out.flush().await?;
} else {
let mut data = xhfs.fread_stream(&r.src_path, r.block_size).await?;
let mut out = stdout();
while let Some(chunk) = data.next().await {
let chunk = chunk?;
out.write_all(&chunk).await?;
}
}
out.flush().await?;
}
Commands::X(x) => x.run().await?,
Commands::Info(global_options) => {
println!(
"Config loaded: {}",
global_options.resolve_config_path().display()
);
let xhfs = global_options.get_xhfs().await?;
println!("{}", xhfs.format_headers_report().await?);
}
Commands::Inspect(i) => i.command.run().await?,
Commands::Server(s) => s.run().await?,
}
Ok(())
}
}
impl GlobalOptions {
fn resolve_config_path(&self) -> PathBuf {
match &self.config {
Some(config) => config.clone(),
None => match std::env::var("XHFS_CONFIG") {
Ok(val) => PathBuf::from(val),
Err(_) => PathBuf::from("xhfs.yaml"),
},
}
}
fn resolve_password(&self) -> Option<String> {
if self.password {
let password = rpassword::prompt_password("Password: ").unwrap();
if password.is_empty() {
return None;
}
Some(password)
} else {
std::env::var("XHFS_PASSWORD").ok()
}
}
async fn load_xhfs(&self, format: bool) -> eyre::Result<XHFS> {
if self.dev {
println!("In-memory mode enabled, the configs will be ignored.");
let xhfs = create_simple_memory_xhfs(
self.dev_unit_capacity,
self.dev_replica_count,
self.dev_logical_count,
)
.await?;
return Ok(xhfs);
}
let config = Config::load(self.resolve_config_path())?;
let password = self.resolve_password();
let xhfs = config.materialize(format, password).await?;
xhfs.get_root_inode().await.with_context(|| {
"Failed to decrypt data. The password may be incorrect or the data may be corrupted.".to_string()
})?;
Ok(xhfs)
}
pub async fn format_and_get_xhfs(&self) -> eyre::Result<XHFS> {
self.load_xhfs(true).await
}
pub async fn get_xhfs(&self) -> eyre::Result<XHFS> {
self.load_xhfs(false).await
}
}
fn confirm_destructive_action() -> bool {
print!("This action is destructive. Proceed? [y/N]: ");
std::io::stdout().flush().unwrap();
let mut input = String::new();
std::io::stdin()
.read_line(&mut input)
.expect("Failed to read line");
match input.trim().to_lowercase().as_str() {
"y" | "yes" => true,
_ => false,
}
}
async fn create_simple_memory_xhfs(
unit_capacity: usize,
replica_count: u8,
logical_count: usize,
) -> eyre::Result<XHFS> {
let slot_capacity = 1024;
if unit_capacity % slot_capacity != 0 {
eyre::bail!("In-memory capacity must be divisible by {slot_capacity}");
}
let mut logical_devices = vec![];
for _ in 0..logical_count {
let devices = (0..replica_count)
.map(|_| {
ConcreteDevice::KVDevice(KVDevice {
store: Arc::new(MemoryKV(RwLock::new(HashMap::new()))),
total_slots: unit_capacity / slot_capacity,
slot_capacity,
})
})
.collect::<Vec<_>>();
logical_devices.push(LogicalDevice::new(2, devices)?);
}
let ctrl = Controller::from(logical_devices).await?;
XHFS::format_new(ctrl, None).await
}