mod cli;
mod render;
use std::io::{self, Write};
use std::process::ExitCode;
use std::str::FromStr as _;
use clap::{CommandFactory as _, Parser};
use onetaskgraph_core::config::Layer;
use onetaskgraph_core::{
DependencyRequest, Engine, Environment, Filters, GlobalId, LabelRequest, Loaded, OutputFormat,
PageToken, Paging, ProjectRequest, ProjectSelector, QueryResponse, SearchRequest,
SourceFailure, TaskRequest,
};
use onetaskgraph_plugin_api::{LabelFilter, NativeId, SourceName, TextQuery};
use serde::Serialize;
use crate::cli::{
Cli, Command, ConfigCommand, DependencyArgs, FilterArgs, LabelCommand, PageArgs,
ProjectCommand, SelectionArgs, ShowArgs, SourcesCommand, TaskCommand,
};
const EXIT_OK: u8 = 0;
const EXIT_FAILURE: u8 = 1;
const EXIT_USAGE: u8 = 2;
const EXIT_PARTIAL: u8 = 4;
#[tokio::main(flavor = "current_thread")]
async fn main() -> ExitCode {
let cli = Cli::parse();
if let Command::PluginServe { source } = &cli.command {
let input = io::BufReader::new(io::stdin().lock());
return match onetaskgraph_core::serve_plugin(input, io::stdout().lock(), *source).await {
Ok(()) => ExitCode::SUCCESS,
Err(error) => fail(&error.to_string(), EXIT_FAILURE),
};
}
let flags = match cli.overrides.layer() {
Ok(flags) => flags,
Err(message) => return fail(&message, EXIT_USAGE),
};
match run(&cli.command, &flags, &mut io::stdout().lock()).await {
Ok(code) => ExitCode::from(code),
Err(message) => fail(&message, EXIT_FAILURE),
}
}
fn fail(message: &str, code: u8) -> ExitCode {
eprintln!("onetaskgraph: {message}");
ExitCode::from(code)
}
async fn run(command: &Command, flags: &Layer, out: &mut impl Write) -> Result<u8, String> {
let loaded = load(flags)?;
let loaded = &loaded;
match command {
Command::PluginServe { .. } => unreachable!("plugin serving is dispatched before config"),
Command::Schema => {
emit(out, schema_bundle()?.trim_end(), "the schema bundle")?;
Ok(EXIT_OK)
}
Command::Config {
command: ConfigCommand::Show,
} => {
emit(
out,
effective_config(loaded)?.trim_end(),
"the configuration",
)?;
Ok(EXIT_OK)
}
Command::Sources {
command: SourcesCommand::List,
} => {
let listings = engine(loaded).listing();
let rendered = match loaded.config.output() {
OutputFormat::Text => render::sources(&listings),
OutputFormat::Json => json(&listings, "the sources")?,
};
emit(out, rendered.trim_end(), "the sources")?;
Ok(EXIT_OK)
}
Command::Task {
command: TaskCommand::List(args),
} => {
let engine = engine(loaded);
let request = TaskRequest {
sources: selection(&args.selection)?,
filters: filters(&args.filters)?,
project: selector(&engine, args.project.as_deref(), args.no_project),
paging: paging(loaded, &args.paging)?,
};
let response = engine
.tasks(&request)
.await
.map_err(|error| error.to_string())?;
respond(out, loaded, response, render::tasks, &args.paging, "tasks")
}
Command::Task {
command: TaskCommand::Show(args),
} => {
let response = engine(loaded)
.task(&qualified(&args.id)?)
.await
.map_err(|error| error.to_string())?;
show(out, loaded, response, render::task_detail, args, "task")
}
Command::Task {
command: TaskCommand::Deps(args),
} => {
let request = dependency_request(loaded, args)?;
let response = engine(loaded)
.task_dependencies(&request)
.await
.map_err(|error| error.to_string())?;
respond(
out,
loaded,
response,
render::edges,
&args.paging,
"dependencies",
)
}
Command::Project {
command: ProjectCommand::List(args),
} => {
let request = ProjectRequest {
sources: selection(&args.selection)?,
filters: filters(&args.filters)?,
paging: paging(loaded, &args.paging)?,
};
let response = engine(loaded)
.projects(&request)
.await
.map_err(|error| error.to_string())?;
respond(
out,
loaded,
response,
render::projects,
&args.paging,
"projects",
)
}
Command::Project {
command: ProjectCommand::Show(args),
} => {
let response = engine(loaded)
.project(&qualified(&args.id)?)
.await
.map_err(|error| error.to_string())?;
show(
out,
loaded,
response,
render::project_detail,
args,
"project",
)
}
Command::Project {
command: ProjectCommand::Deps(args),
} => {
let request = dependency_request(loaded, args)?;
let response = engine(loaded)
.project_dependencies(&request)
.await
.map_err(|error| error.to_string())?;
respond(
out,
loaded,
response,
render::edges,
&args.paging,
"dependencies",
)
}
Command::Label {
command: LabelCommand::List(args),
} => {
let request = LabelRequest {
sources: selection(&args.selection)?,
paging: paging(loaded, &args.paging)?,
};
let response = engine(loaded)
.labels(&request)
.await
.map_err(|error| error.to_string())?;
respond(
out,
loaded,
response,
render::labels,
&args.paging,
"labels",
)
}
Command::Search(args) => {
let request = SearchRequest {
sources: selection(&args.selection)?,
text: TextQuery {
terms: args.text.clone(),
fields: args.fields.fields(),
},
kind: args.kind.kind(),
paging: paging(loaded, &args.paging)?,
};
let response = engine(loaded)
.search(&request)
.await
.map_err(|error| error.to_string())?;
respond(
out,
loaded,
response,
render::hits,
&args.paging,
"the search",
)
}
}
}
fn engine(loaded: &Loaded) -> Engine {
Engine::build(&loaded.config, &loaded.secrets)
}
fn respond<T: Serialize>(
out: &mut impl Write,
loaded: &Loaded,
response: QueryResponse<T>,
text: impl FnOnce(&[T]) -> String,
paging: &PageArgs,
what: &str,
) -> Result<u8, String> {
let rendered = match loaded.config.output() {
OutputFormat::Text => {
let mut rendered = text(&response.items);
if paging.explain {
rendered.push('\n');
rendered.push_str(&render::plan(&response.plan));
}
if let Some(next) = &response.next {
rendered.push_str(&format!("\nnext page: --page {next}\n"));
}
rendered
}
OutputFormat::Json => json(&response, what)?,
};
emit(out, rendered.trim_end(), what)?;
Ok(report(&response.errors, paging.allow_partial))
}
fn show<T: Serialize>(
out: &mut impl Write,
loaded: &Loaded,
response: QueryResponse<T>,
text: impl FnOnce(&T) -> String,
args: &ShowArgs,
what: &str,
) -> Result<u8, String> {
match (response.items.first(), response.errors.is_empty()) {
(None, true) => Err(format!(
"no {what} with that id\n\
next: check the id, or list what is there — `onetaskgraph {what} list` \
reports every {what} the configured sources hold."
)),
_ => {
let rendered = match loaded.config.output() {
OutputFormat::Text => {
let mut rendered = response.items.first().map(text).unwrap_or_default();
if args.explain {
rendered.push('\n');
rendered.push_str(&render::plan(&response.plan));
}
rendered
}
OutputFormat::Json => json(&response, what)?,
};
emit(out, rendered.trim_end(), what)?;
Ok(report(&response.errors, args.allow_partial))
}
}
}
fn report(errors: &[SourceFailure], allow_partial: bool) -> u8 {
if errors.is_empty() {
return EXIT_OK;
}
for failure in errors {
eprintln!(
"onetaskgraph: source {} could not answer: {}",
failure.source, failure.error
);
}
if allow_partial {
eprintln!(
"onetaskgraph: the answer above is partial, and --allow-partial says that is \
acceptable."
);
return EXIT_OK;
}
eprintln!(
"onetaskgraph: next: fix the source(s) named above — `onetaskgraph sources list` \
reports each one's state — or re-run with --allow-partial to accept an answer \
without them."
);
EXIT_PARTIAL
}
fn selection(args: &SelectionArgs) -> Result<Vec<SourceName>, String> {
args.source
.iter()
.map(|name| {
SourceName::new(name.clone()).map_err(|error| {
format!(
"--source {name}: {error}\n\
next: name a configured source — `onetaskgraph sources list` reports \
them."
)
})
})
.collect()
}
fn filters(args: &FilterArgs) -> Result<Filters, String> {
Ok(Filters {
text: args.search.as_ref().map(|terms| TextQuery {
terms: terms.clone(),
fields: args.fields.fields(),
}),
labels: LabelFilter {
all_of: args.label.clone(),
none_of: args.not_label.clone(),
any_of: Vec::new(),
},
statuses: args.status.iter().map(|status| status.category()).collect(),
})
}
fn selector(engine: &Engine, project: Option<&str>, orphans: bool) -> ProjectSelector {
if orphans {
return ProjectSelector::Orphans;
}
let Some(project) = project else {
return ProjectSelector::Any;
};
match GlobalId::from_str(project) {
Ok(id) if engine.has(&id.source) => ProjectSelector::Qualified(id),
Ok(_) | Err(_) => ProjectSelector::Native(NativeId::from(project)),
}
}
fn qualified(id: &str) -> Result<GlobalId, String> {
GlobalId::from_str(id).map_err(|error| {
format!(
"{error}\n\
next: qualify the id with the source it belongs to — `onetaskgraph sources \
list` reports the configured names."
)
})
}
fn paging(loaded: &Loaded, args: &PageArgs) -> Result<Paging, String> {
let limit = args.limit.unwrap_or_else(|| loaded.config.page_size());
let token = args
.page
.as_ref()
.map(|raw| {
PageToken::parse(raw.clone()).map_err(|error| {
format!(
"--page: {error}\n\
next: pass a token exactly as a previous page reported it, or drop \
--page to start the walk again."
)
})
})
.transpose()?;
Ok(Paging { limit, token })
}
fn dependency_request(loaded: &Loaded, args: &DependencyArgs) -> Result<DependencyRequest, String> {
Ok(DependencyRequest {
id: qualified(&args.id)?,
direction: args.direction.direction(),
paging: paging(loaded, &args.paging)?,
})
}
fn json(value: &impl Serialize, what: &str) -> Result<String, String> {
serde_json::to_string_pretty(value).map_err(|error| format!("could not render {what}: {error}"))
}
fn schema_bundle() -> Result<String, String> {
let mut bundle = onetaskgraph_core::schema_bundle();
bundle["commands"] = serde_json::to_value(public_commands()?)
.map_err(|error| format!("could not render the command surface: {error}"))?;
serde_json::to_string_pretty(&bundle)
.map_err(|error| format!("could not render the schema bundle: {error}"))
}
#[derive(Serialize)]
#[serde(transparent)]
struct PublicCommand(String);
impl PublicCommand {
fn try_new(path: String) -> Result<Self, String> {
if path.is_empty() || path.split(' ').any(|part| part.is_empty()) {
return Err(format!("invalid public command path {path:?}"));
}
Ok(Self(path))
}
}
fn public_commands() -> Result<Vec<PublicCommand>, String> {
fn leaves(
command: &clap::Command,
prefix: &str,
commands: &mut Vec<PublicCommand>,
) -> Result<(), String> {
let visible: Vec<_> = command
.get_subcommands()
.filter(|child| !child.is_hide_set())
.collect();
if visible.is_empty() {
if !prefix.is_empty() {
commands.push(PublicCommand::try_new(prefix.to_owned())?);
}
return Ok(());
}
for child in visible {
let path = if prefix.is_empty() {
child.get_name().to_owned()
} else {
format!("{prefix} {}", child.get_name())
};
leaves(child, &path, commands)?;
}
Ok(())
}
let mut commands = Vec::new();
leaves(&Cli::command(), "", &mut commands)?;
Ok(commands)
}
fn effective_config(loaded: &Loaded) -> Result<String, String> {
match loaded.config.output() {
OutputFormat::Text => Ok(loaded.effective.render_text()),
OutputFormat::Json => serde_json::to_string_pretty(&loaded.effective)
.map_err(|error| format!("could not render the configuration: {error}")),
}
}
fn load(flags: &Layer) -> Result<Loaded, String> {
let working_directory = std::env::current_dir().map_err(|error| {
format!(
"could not read the working directory: {error}\n\
next: run this from a directory that still exists."
)
})?;
onetaskgraph_core::config::load(&working_directory, &Environment::from_process(), flags)
.map_err(|error| error.to_string())
}
fn emit(out: &mut impl Write, rendered: &str, what: &str) -> Result<(), String> {
writeln!(out, "{rendered}").map_err(|error| format!("could not write {what}: {error}"))?;
out.flush()
.map_err(|error| format!("could not write {what}: {error}"))
}
#[cfg(test)]
mod tests {
use super::*;
fn write_schema_bundle(out: &mut impl Write) -> Result<(), String> {
emit(out, schema_bundle()?.trim_end(), "the schema bundle")
}
#[test]
fn the_schema_verb_writes_a_bundle_with_every_contract_root() {
let mut out = Vec::new();
write_schema_bundle(&mut out).expect("the bundle renders");
let bundle: serde_json::Value =
serde_json::from_slice(&out).expect("the bundle is valid JSON");
assert!(bundle["roots"]["Task"].is_object());
assert!(bundle["plugin_config"]["in-memory"].is_object());
assert_eq!(
bundle["commands"],
serde_json::json!([
"schema",
"config show",
"sources list",
"task list",
"task show",
"task deps",
"project list",
"project show",
"project deps",
"label list",
"search"
])
);
}
struct Failing {
fail_on_write: bool,
}
impl Write for Failing {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
if self.fail_on_write {
Err(io::Error::other("the pipe is closed"))
} else {
Ok(buf.len())
}
}
fn flush(&mut self) -> io::Result<()> {
Err(io::Error::other("the pipe closed before the flush"))
}
}
#[test]
fn a_verb_reports_a_failed_write_rather_than_panicking() {
let mut sink = Failing {
fail_on_write: true,
};
let message = write_schema_bundle(&mut sink).expect_err("writes refused");
assert!(
message.contains("could not write the schema bundle"),
"{message}"
);
}
#[test]
fn a_verb_reports_a_failed_flush_rather_than_exiting_zero_on_a_truncated_document() {
let mut sink = Failing {
fail_on_write: false,
};
let message = write_schema_bundle(&mut sink).expect_err("flushes refused");
assert!(
message.contains("could not write the schema bundle"),
"{message}"
);
}
}