use std::{path::PathBuf, process::ExitCode, sync::Arc, time::Duration};
use basis_tasks::{Continuation, RunSpec, TaskHandle, Tasks, WaitOutcome};
use serde_json::json;
use crate::{
cli::{AskArgs, CancelArgs, InboxArgs, RunArgs, SendArgs, WaitArgs, WatchArgs},
duration_arg::DurationArg,
exit::EXIT_OK,
run::prompt_from,
};
use super::{
error::{ClientError, message_timeout, wait_timeout, watch_timeout},
prompt_host::CliPromptHost,
render::{Live, decorate_terminal, print_hint, render_result},
};
const DEFAULT_WAIT: Duration = Duration::from_secs(30 * 60);
const MAX_WAIT: Duration = Duration::from_secs(7 * 24 * 60 * 60);
pub(crate) async fn spawn(args: RunArgs, attach: bool) -> Result<ExitCode, ClientError> {
let workspace = workspace_or_current(args.workspace.clone())?;
let prompt = prompt_from(args.prompt.clone())?;
let approve: basis_tasks::Approve = args.approve.into();
let spec = run_spec(&args, prompt, approve)?;
let tasks = open_tasks(workspace)?;
if !attach {
let handle = tasks.spawn(spec)?;
let payload = json!({
"task": handle,
"state": "resumable",
"next": format!("basis wait {handle}"),
});
return render_result(&payload, args.json);
}
let timeout = bounded_wait(args.timeout);
let live = Live::when(!args.json);
let (handle, outcome) = tasks
.spawn_and_wait(spec, timeout, Some(live_sink(live.clone())))
.await?;
match outcome {
WaitOutcome::Terminal(terminal) => {
live.settled(&decorate_terminal(handle.as_ref(), terminal), args.json)
}
WaitOutcome::TimedOut { attached } => Err(wait_timeout(handle.as_ref(), timeout, attached)),
}
}
fn run_spec(
args: &RunArgs,
prompt: String,
approve: basis_tasks::Approve,
) -> Result<RunSpec, ClientError> {
let mut spec = RunSpec::new(prompt).with_approve(approve);
if let Some(provider) = &args.provider {
spec = spec.with_provider(provider.clone());
}
if let Some(base_url) = &args.base_url {
spec = spec.with_base_url(base_url.clone());
}
if let Some(model) = &args.model {
spec = spec.with_model(model.clone());
}
if args.no_shell {
spec = spec.without_shell();
}
if let Some(system_prompt) = crate::cli::system_prompt(
args.system_prompt.clone(),
args.append_system_prompt.clone(),
) {
spec = spec.with_system_prompt(system_prompt);
}
if let Some(effort) = args.effort {
spec = spec.with_effort(effort.into());
}
if let Some(deadline) = args.deadline {
spec = spec.with_deadline(deadline.duration());
}
if let Some(tool_budget) = args.tool_budget {
spec = spec.with_tool_budget(tool_budget);
}
if let Some(token_budget) = args.token_budget {
spec = spec.with_token_budget(token_budget);
}
if args.detached {
spec = spec.detached();
}
if args.continue_latest {
spec = spec.continuing(Continuation::Latest);
} else if let Some(session) = &args.session {
let handle = TaskHandle::parse(session.clone()).map_err(|_| {
ClientError::usage(format!(
"`{session}` is not a task handle; `basis list` prints them"
))
})?;
spec = spec.continuing(Continuation::Named(handle));
}
Ok(spec)
}
pub(crate) fn has_current_task() -> bool {
basis_tasks::current_task().is_some()
}
pub(crate) async fn send(args: SendArgs) -> Result<ExitCode, ClientError> {
send_message(
args.task,
args.message,
args.await_result,
args.timeout,
args.json,
)
.await
}
pub(crate) async fn ask(args: AskArgs) -> Result<ExitCode, ClientError> {
send_message(args.task, args.message, true, args.timeout, args.json).await
}
async fn send_message(
task: String,
raw_message: String,
await_result: bool,
timeout: Option<DurationArg>,
json: bool,
) -> Result<ExitCode, ClientError> {
let handle = TaskHandle::parse(task.clone())?;
let tasks = open_tasks(current_dir()?)?;
let from_stdin = raw_message == "-";
let message = prompt_from(raw_message)?;
let message = if from_stdin {
message
} else {
let workspace = tasks.workspace_of(&handle)?;
crate::templates::resolve(&message, &workspace)?.unwrap_or(message)
};
let caller = basis_tasks::current_task();
if !await_result {
let message_id = tasks.send(&handle, message)?;
let payload = json!({
"task": task,
"message": message_id,
"state": "accepted",
"next": send_next_hint(&tasks, caller.as_ref(), &handle, &message_id),
});
return render_result(&payload, json);
}
let timeout = bounded_wait(timeout);
let reply = tasks
.ask(&handle, caller.as_ref(), message, timeout)
.await?;
match reply.outcome {
WaitOutcome::Terminal(payload) => render_result(&payload, json),
WaitOutcome::TimedOut { .. } => Err(message_timeout(&task, &reply.message_id, timeout)),
}
}
fn send_next_hint(
tasks: &Tasks,
caller: Option<&TaskHandle>,
target: &TaskHandle,
message: &str,
) -> String {
if tasks.validate_wait_edge(caller, target).is_ok() {
format!("basis wait {target} --message {message}")
} else {
format!("basis inbox {target}")
}
}
pub(crate) async fn wait(args: WaitArgs) -> Result<ExitCode, ClientError> {
let handle = TaskHandle::parse(args.task.clone())?;
let tasks = open_tasks(current_dir()?)?;
let caller = basis_tasks::current_task();
let timeout = bounded_wait(args.timeout);
if let Some(message_id) = args.message {
return match tasks
.wait_message(&handle, caller.as_ref(), &message_id, timeout)
.await?
{
WaitOutcome::Terminal(payload) => render_result(&payload, args.json),
WaitOutcome::TimedOut { .. } => Err(message_timeout(&args.task, &message_id, timeout)),
};
}
if let Some(terminal) = tasks.terminal(&handle)? {
return render_result(&decorate_terminal(&args.task, terminal), args.json);
}
let live = Live::when(!args.json);
match tasks
.wait(
&handle,
caller.as_ref(),
timeout,
Some(live_sink(live.clone())),
)
.await?
{
WaitOutcome::Terminal(terminal) => {
live.settled(&decorate_terminal(&args.task, terminal), args.json)
}
WaitOutcome::TimedOut { attached } => Err(wait_timeout(&args.task, timeout, attached)),
}
}
pub(crate) async fn cancel(args: CancelArgs) -> Result<ExitCode, ClientError> {
let handle = TaskHandle::parse(args.task.clone())?;
let tasks = open_tasks(current_dir()?)?;
let caller = basis_tasks::current_task();
tasks.validate_cancel_target(caller.as_ref(), &handle)?;
if let Some(terminal) = tasks.terminal(&handle)? {
return render_result(&decorate_terminal(&args.task, terminal), args.json);
}
tasks.cancel(&handle, caller.as_ref())?;
let payload = json!({
"task": args.task,
"state": "cancel_requested",
"next": format!("basis wait {}", args.task),
});
render_result(&payload, args.json)
}
pub(crate) async fn watch(args: WatchArgs) -> Result<ExitCode, ClientError> {
let handle = TaskHandle::parse(args.task.clone())?;
let tasks = open_tasks(current_dir()?)?;
if tasks.terminal(&handle)?.is_none() {
tasks.validate_wait_edge(basis_tasks::current_task().as_ref(), &handle)?;
}
let timeout = args
.timeout
.map(DurationArg::duration)
.unwrap_or(DEFAULT_WAIT);
let deadline = tokio::time::Instant::now() + timeout;
let mut cursor = tasks.watch(&handle)?;
let live = Live::when(!args.json);
loop {
let terminal = tasks.terminal(&handle)?;
for record in cursor.poll()? {
if args.json {
println!("{}", record.raw);
} else {
live.show(&record.raw)
.map_err(|error| format!("render task progress: {error}"))?;
}
}
if let Some(terminal) = terminal {
return live.settled(&decorate_terminal(&args.task, terminal), args.json);
}
if tokio::time::Instant::now() >= deadline {
return Err(watch_timeout(&args.task, tasks.is_attached(&handle)?));
}
tokio::time::sleep(basis_tasks::POLL).await;
}
}
pub(crate) async fn inbox(args: InboxArgs) -> Result<ExitCode, ClientError> {
let task = args
.task
.or_else(|| basis_tasks::current_task().map(|handle| handle.to_string()))
.ok_or_else(|| {
"`basis inbox` needs a task id outside a basis task: use `basis inbox <ID>`".to_string()
})?;
let handle = TaskHandle::parse(task.clone())?;
let tasks = open_tasks(current_dir()?)?;
let payload = tasks.inbox(&handle)?;
if args.json {
println!("{payload}");
return Ok(ExitCode::from(EXIT_OK));
}
let messages = payload["messages"].as_array().cloned().unwrap_or_default();
if messages.is_empty() {
println!("inbox is empty");
} else {
for message in messages {
let state = message["state"].as_str().unwrap_or("unknown");
let id = message["id"].as_str().unwrap_or("?");
let body = message["body"].as_str().unwrap_or_default();
println!("[{state}] {id}: {body}");
if let Some(reply) = message["reply"]["result"].as_str()
&& !reply.is_empty()
{
println!(" reply: {reply}");
}
}
}
print_hint(&payload);
Ok(ExitCode::from(EXIT_OK))
}
fn live_sink(live: Live) -> Arc<dyn basis_tasks::LiveSink> {
Arc::new(live)
}
fn open_tasks(workspace: PathBuf) -> Result<Tasks, ClientError> {
Ok(Tasks::open(workspace)
.map_err(ClientError::from)?
.with_prompt_host(Arc::new(CliPromptHost)))
}
fn current_dir() -> Result<PathBuf, ClientError> {
std::env::current_dir().map_err(|error| format!("no working directory: {error}").into())
}
fn workspace_or_current(workspace: Option<PathBuf>) -> Result<PathBuf, ClientError> {
workspace.map_or_else(current_dir, Ok)
}
fn bounded_wait(timeout: Option<DurationArg>) -> Duration {
timeout
.map(|timeout| timeout.duration().min(MAX_WAIT))
.unwrap_or(DEFAULT_WAIT)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn client_waits_are_defaulted_and_clamped() {
assert_eq!(bounded_wait(None), DEFAULT_WAIT);
assert_eq!(
bounded_wait(Some("30s".parse().unwrap())),
Duration::from_secs(30)
);
assert_eq!(bounded_wait(Some("30d".parse().unwrap())), MAX_WAIT);
}
}