mod daemon;
mod publish;
mod queues;
mod request;
mod run;
mod status;
mod store;
mod table;
mod validate;
use std::path::PathBuf;
use std::process::ExitCode;
use std::time::Duration;
use clap::{Parser, Subcommand};
use swale::partition::InvalidPartition;
use swale::{DefinitionError, Error, Partition, Request};
use self::request::RequestArgs;
use self::store::StoreArg;
pub(crate) type CommandResult = Result<ExitCode, Box<dyn std::error::Error>>;
#[derive(Parser)]
#[command(name = "swale", version, about)]
struct Cli {
#[command(subcommand)]
command: Command,
}
#[derive(Subcommand)]
enum Command {
Validate {
file: PathBuf,
},
Run {
file: PathBuf,
#[command(flatten)]
store: StoreArg,
#[arg(long, value_parser = parse_partition_arg)]
partition: Option<Partition>,
#[arg(long, default_value_t = 4)]
concurrency: usize,
},
Publish {
file: PathBuf,
#[command(flatten)]
store: StoreArg,
},
Daemon {
#[command(flatten)]
store: StoreArg,
#[arg(long, default_value_t = 4)]
concurrency: usize,
#[arg(long = "pool", value_parser = daemon::parse_pool_arg)]
pools: Vec<(String, usize)>,
#[arg(long, default_value_t = 30)]
sync_interval: u64,
#[arg(long, default_value = "90d", value_parser = swale::duration::parse)]
retention: Duration,
},
Start {
graph: String,
#[arg(required = true, value_parser = parse_partition_arg)]
partitions: Vec<Partition>,
#[command(flatten)]
request: RequestArgs,
},
Rerun {
graph: String,
#[arg(value_parser = parse_partition_arg)]
partition: Partition,
node: String,
#[command(flatten)]
request: RequestArgs,
},
Cancel {
graph: String,
#[arg(value_parser = parse_partition_arg)]
partition: Partition,
#[command(flatten)]
request: RequestArgs,
},
Status {
graph: Option<String>,
#[arg(value_parser = parse_partition_arg)]
partition: Option<Partition>,
#[command(flatten)]
store: StoreArg,
#[arg(long, default_value_t = 20)]
limit: usize,
},
Queues {
queue: Option<String>,
#[command(flatten)]
store: StoreArg,
#[arg(long, default_value_t = 20)]
limit: usize,
},
}
pub(crate) async fn main() -> ExitCode {
let result = match Cli::parse().command {
Command::Validate { file } => validate::validate(&file),
Command::Run {
file,
store,
partition,
concurrency,
} => run::run(&file, store, partition, concurrency).await,
Command::Publish { file, store } => publish::publish(&file, store).await,
Command::Daemon {
store,
concurrency,
pools,
sync_interval,
retention,
} => {
daemon::daemon(
store,
concurrency,
pools,
Duration::from_secs(sync_interval),
retention,
)
.await
}
Command::Start {
graph,
partitions,
request: args,
} => request::send_request(Request::Start { graph, partitions }, args).await,
Command::Rerun {
graph,
partition,
node,
request: args,
} => {
request::send_request(
Request::Rerun {
graph,
partition,
node,
},
args,
)
.await
}
Command::Cancel {
graph,
partition,
request: args,
} => request::send_request(Request::Cancel { graph, partition }, args).await,
Command::Status {
graph,
partition,
store,
limit,
} => status::status(store, graph, partition, limit).await,
Command::Queues {
queue,
store,
limit,
} => queues::queues(store, queue, limit).await,
};
match result {
Ok(code) => code,
Err(e) => {
print_error(e.as_ref());
ExitCode::FAILURE
}
}
}
fn print_error(error: &(dyn std::error::Error + 'static)) {
let problems = match error.downcast_ref::<Error>() {
Some(Error::Invalid(problems)) => Some(problems),
_ => match error.downcast_ref::<DefinitionError>() {
Some(DefinitionError::Invalid(Error::Invalid(problems))) => Some(problems),
_ => None,
},
};
match problems {
Some(problems) => {
for problem in problems {
eprintln!("error: {problem}");
}
}
None => eprintln!("error: {error}"),
}
}
pub(crate) fn parse_partition_arg(raw: &str) -> Result<Partition, InvalidPartition> {
Partition::new(raw)
}