use std::{
fmt,
io::{self, Read, Write},
num::NonZeroUsize,
path::PathBuf,
sync::Arc,
time::Duration,
};
use rho_sdk::{SessionOptions, SystemPrompt, UserInput};
use {
crate::agent::{PromptPolicy, ToolCapability},
crate::cli::{Command, OutputFormat},
crate::config::Config,
crate::credential_store::AppCredentialStore,
crate::diagnostics::RuntimeDiagnostics,
crate::herdr::{HerdrReporter, HerdrState},
crate::prompt,
crate::subagent::{RunState, RunStatus},
crate::tools::{
agent::BackgroundSubagents,
sdk_registry::{AppToolSet, DelegationConfig, ToolSetOptions},
},
rho_providers::providers::build_automation_provider,
};
use super::{
agent_binding::BoundAgent,
automation_protocol::{write_event, JsonlAdapter, TerminalReason, WireEvent},
headless_run::{self, HeadlessRunDeps, HostInputResponder},
policy::AppPolicy,
runtime_builder::{
build_runtime_with_max_steps, configured_context_window, RuntimeBuildOptions,
},
sdk_config::SdkBootstrapOptions,
};
#[derive(Debug)]
pub struct AutomationExit {
code: u8,
reason: TerminalReason,
message: String,
}
impl AutomationExit {
pub(super) fn new(code: u8, reason: TerminalReason, message: impl Into<String>) -> Self {
Self {
code,
reason,
message: message.into(),
}
}
pub fn exit_code(&self) -> u8 {
self.code
}
fn reason(&self) -> TerminalReason {
self.reason
}
}
impl fmt::Display for AutomationExit {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str(&self.message)
}
}
impl std::error::Error for AutomationExit {}
#[derive(Debug)]
pub struct AutomationInterrupted {
signal: ShutdownSignal,
}
impl AutomationInterrupted {
fn new(signal: ShutdownSignal) -> Self {
Self { signal }
}
pub fn exit_code(&self) -> u8 {
self.signal.exit_code()
}
}
impl fmt::Display for AutomationInterrupted {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(formatter, "rho run interrupted by {}", self.signal)
}
}
impl std::error::Error for AutomationInterrupted {}
#[derive(Clone, Copy, Debug)]
enum ShutdownSignal {
Interrupt,
Terminate,
}
impl ShutdownSignal {
fn exit_code(self) -> u8 {
match self {
Self::Interrupt => 130,
Self::Terminate => 143,
}
}
}
impl fmt::Display for ShutdownSignal {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Interrupt => formatter.write_str("SIGINT"),
Self::Terminate => formatter.write_str("SIGTERM"),
}
}
}
#[derive(Debug)]
struct SubagentCancelled;
impl fmt::Display for SubagentCancelled {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str("subagent cancellation requested")
}
}
impl std::error::Error for SubagentCancelled {}
pub(super) struct Startup<'a> {
pub config: &'a Config,
pub config_path: PathBuf,
pub cwd: PathBuf,
pub no_system_prompt: bool,
pub no_tools: bool,
pub no_subagents: bool,
pub usage_purpose: &'static str,
pub parent_session_id: Option<rho_sdk::SessionId>,
pub agent: BoundAgent,
pub output_file: Option<PathBuf>,
pub output: OutputFormat,
pub max_steps: Option<NonZeroUsize>,
pub timeout: Option<Duration>,
pub diagnostics: RuntimeDiagnostics,
pub herdr: HerdrReporter,
pub host_input: Option<Arc<dyn HostInputResponder>>,
}
pub(super) fn prompt_for_command(command: &Option<Command>) -> anyhow::Result<Option<String>> {
match command {
Some(Command::Run { prompt, stdin, .. }) => {
prompt_from_stdin(prompt.clone(), *stdin).map(Some)
}
Some(
Command::Attach { .. }
| Command::Login { .. }
| Command::CredentialStore { .. }
| Command::Sessions { .. }
| Command::Update,
)
| None => Ok(None),
}
}
pub(super) fn emit_startup_failure() -> anyhow::Result<()> {
let mut adapter = JsonlAdapter::new();
let event = adapter.failed(
TerminalReason::ConfigurationError,
"configuration failed".into(),
None,
);
emit(event)
}
pub(super) async fn run(prompt_text: String, startup: Startup<'_>) -> anyhow::Result<()> {
let mut jsonl = (startup.output == OutputFormat::Jsonl).then(JsonlAdapter::new);
let deadline = startup
.timeout
.map(|timeout| tokio::time::Instant::now() + timeout);
let reporter_result = startup
.output_file
.as_ref()
.map(|path| {
RunReporter::new(
path.clone(),
RunArtifactIdentity {
agent_id: startup.agent.id().to_string(),
agent_fingerprint: startup.agent.fingerprint().to_string(),
provider: startup.config.provider.clone(),
model: startup.config.model.clone(),
},
startup.cwd.clone(),
&prompt_text,
startup.output == OutputFormat::Text,
None,
)
})
.transpose();
let mut reporter = match reporter_result {
Ok(reporter) => reporter,
Err(error) => {
emit_failure(&mut jsonl, TerminalReason::OutputError, &error)?;
return Err(
AutomationExit::new(1, TerminalReason::OutputError, error.to_string()).into(),
);
}
};
let cancellation = rho_tools::cancellation::RunCancellation::default();
let (result, timed_out) = if let Some(deadline) = deadline {
let future = run_session_with_output(
prompt_text,
&startup,
reporter.as_mut(),
Some(cancellation.clone()),
jsonl.as_mut(),
);
tokio::pin!(future);
tokio::select! {
result = &mut future => (result, false),
() = tokio::time::sleep_until(deadline) => {
cancellation.cancel();
(future.await, true)
}
}
} else {
(
run_session_with_output(
prompt_text,
&startup,
reporter.as_mut(),
None,
jsonl.as_mut(),
)
.await,
false,
)
};
if let Some(reporter) = reporter.as_mut() {
let reached_step_limit = result.as_ref().is_ok_and(|outcome| {
outcome.stop_reason() == rho_sdk::StopReason::MaxSteps
&& (jsonl.is_some() || startup.max_steps.is_some())
});
if reached_step_limit {
let stopped = Err(AutomationExit::new(
124,
TerminalReason::MaxSteps,
"rho run reached its model-step limit",
)
.into());
reporter.finish(&stopped);
} else {
reporter.finish(&result);
}
}
if timed_out {
emit_stopped(&mut jsonl, TerminalReason::Timeout)?;
return Err(AutomationExit::new(124, TerminalReason::Timeout, "rho run timed out").into());
}
match result {
Ok(answer) => {
let max_steps = answer.stop_reason() == rho_sdk::StopReason::MaxSteps;
if max_steps && (jsonl.is_some() || startup.max_steps.is_some()) {
if let Some(adapter) = jsonl.as_mut() {
let text = (!answer.text().is_empty()).then(|| answer.text().into());
let event = adapter.stopped(TerminalReason::MaxSteps, text);
emit(event)?;
} else {
write_text_answer(&answer, reporter.is_some())?;
}
return Err(AutomationExit::new(
124,
TerminalReason::MaxSteps,
"rho run reached its model-step limit",
)
.into());
}
if let Some(adapter) = jsonl.as_mut() {
let event = adapter.completed(answer.text().into());
emit(event)?;
} else {
write_text_answer(&answer, reporter.is_some())?;
}
Ok(())
}
Err(error) => {
let (reason, code) = classify_error(&error);
if reason == TerminalReason::Interrupted {
emit_stopped(&mut jsonl, reason)?;
} else if reason != TerminalReason::OutputError {
emit_failure(&mut jsonl, reason, &error)?;
}
let message = terminal_error_message(reason, &error);
if error.is::<AutomationInterrupted>() {
return Err(error);
}
Err(AutomationExit::new(code, reason, message).into())
}
}
}
fn write_text_answer(answer: &rho_sdk::RunOutcome, has_reporter: bool) -> anyhow::Result<()> {
let result = (|| -> io::Result<()> {
let mut stdout = io::stdout().lock();
if has_reporter {
writeln!(stdout, "\n[subagent run complete]")?;
} else {
writeln!(stdout, "{}", answer.text())?;
}
stdout.flush()
})();
result.map_err(|error| {
AutomationExit::new(
1,
TerminalReason::OutputError,
format!("could not write output: {error}"),
)
.into()
})
}
pub(super) fn emit(event: WireEvent) -> anyhow::Result<()> {
let mut stdout = io::stdout().lock();
write_event(&mut stdout, &event).map_err(|error| {
AutomationExit::new(
1,
TerminalReason::OutputError,
format!("could not write JSONL output: {error}"),
)
.into()
})
}
fn emit_stopped(adapter: &mut Option<JsonlAdapter>, reason: TerminalReason) -> anyhow::Result<()> {
if let Some(adapter) = adapter.as_mut() {
let text = adapter.partial_text();
let event = adapter.stopped(reason, text);
emit(event)?;
}
Ok(())
}
fn emit_failure(
adapter: &mut Option<JsonlAdapter>,
reason: TerminalReason,
error: &anyhow::Error,
) -> anyhow::Result<()> {
if let Some(adapter) = adapter.as_mut() {
let text = adapter.partial_text();
let message = terminal_error_message(reason, error);
let event = adapter.failed(reason, message, text);
emit(event)?;
}
Ok(())
}
fn terminal_error_message(reason: TerminalReason, error: &anyhow::Error) -> String {
match reason {
TerminalReason::Authentication => "authentication failed".to_string(),
TerminalReason::ConfigurationError => "configuration failed".to_string(),
TerminalReason::OutputError => "output failed".to_string(),
TerminalReason::OtherError => "run failed".to_string(),
_ => error.to_string(),
}
}
fn classify_error(error: &anyhow::Error) -> (TerminalReason, u8) {
if let Some(interrupted) = error.downcast_ref::<AutomationInterrupted>() {
return (TerminalReason::Interrupted, interrupted.exit_code());
}
if let Some(exit) = error.downcast_ref::<AutomationExit>() {
return (exit.reason(), exit.exit_code());
}
for cause in error.chain() {
if let Some(error) = cause.downcast_ref::<rho_sdk::Error>() {
return match error {
rho_sdk::Error::Authentication { .. } => (TerminalReason::Authentication, 1),
rho_sdk::Error::Provider(provider)
if provider.kind() == rho_sdk::ProviderErrorKind::Authentication =>
{
(TerminalReason::Authentication, 1)
}
rho_sdk::Error::Provider(_) => (TerminalReason::ProviderError, 1),
rho_sdk::Error::Tool(_) => (TerminalReason::ToolHostError, 1),
rho_sdk::Error::InvalidConfiguration { .. } => {
(TerminalReason::ConfigurationError, 2)
}
_ => (TerminalReason::OtherError, 1),
};
}
if let Some(error) = cause.downcast_ref::<rho_providers::model::ModelError>() {
use rho_providers::model::ModelError;
return match error {
ModelError::MissingApiKey
| ModelError::MissingCodexAuth
| ModelError::MissingAnthropicApiKey
| ModelError::MissingGoogleApiKey
| ModelError::MissingGithubCopilotAuth
| ModelError::MissingMoonshotApiKey
| ModelError::MissingPoolsideApiKey
| ModelError::MissingOpenRouterApiKey
| ModelError::MissingCredentialProfile(_)
| ModelError::MissingKimiAuth
| ModelError::MissingXaiApiKey
| ModelError::MissingXaiAuth
| ModelError::Credentials(_) => (TerminalReason::Authentication, 1),
ModelError::UnsupportedReasoning { .. } | ModelError::UnsupportedProvider(_) => {
(TerminalReason::ConfigurationError, 2)
}
_ => (TerminalReason::ProviderError, 1),
};
}
}
(TerminalReason::OtherError, 1)
}
pub(crate) async fn run_session(
prompt_text: String,
startup: &Startup<'_>,
reporter: Option<&mut RunReporter>,
cancellation: Option<rho_tools::cancellation::RunCancellation>,
) -> anyhow::Result<rho_sdk::RunOutcome> {
run_session_with_output(prompt_text, startup, reporter, cancellation, None).await
}
async fn run_session_with_output(
prompt_text: String,
startup: &Startup<'_>,
reporter: Option<&mut RunReporter>,
cancellation: Option<rho_tools::cancellation::RunCancellation>,
mut jsonl: Option<&mut JsonlAdapter>,
) -> anyhow::Result<rho_sdk::RunOutcome> {
let sdk_options = SdkBootstrapOptions::from_config(startup.config, &startup.cwd)?;
let credentials = rho_providers::auth::provider_credentials::ApplicationCredentialSource::new(
Arc::new(AppCredentialStore),
);
let provider = build_automation_provider(sdk_options.provider, &credentials)?;
let mut capabilities = startup
.agent
.rho_capabilities()
.cloned()
.unwrap_or_default();
if startup.no_subagents {
capabilities.remove(&ToolCapability::Agent);
capabilities.remove(&ToolCapability::Agents);
}
let launch_delegation_enabled = capabilities.contains(&ToolCapability::Agent);
let delegation_enabled =
launch_delegation_enabled || capabilities.contains(&ToolCapability::Agents);
let tool_set = if startup.no_tools {
AppToolSet::disabled()
} else {
let mut options = ToolSetOptions::new(capabilities);
if delegation_enabled {
options = options.delegation(DelegationConfig::new(
startup.cwd.clone(),
startup.config_path.clone(),
BackgroundSubagents::Disabled,
));
}
AppToolSet::new(startup.config, startup.diagnostics.clone(), options)
};
let tool_specs = tool_set.specs();
let system_prompt = if startup.no_system_prompt {
startup.diagnostics.update_prompt_sources(Vec::new());
SystemPrompt::None
} else {
let mut text = match startup.agent.prompt() {
PromptPolicy::Replace(text) => text.clone(),
PromptPolicy::Extend(extra) => {
let built = prompt::system_prompt(&tool_specs, &startup.cwd);
startup.diagnostics.update_prompt_sources(built.sources);
let mut text = built.text;
if !launch_delegation_enabled {
prompt::append_subagents_disabled_instruction(&mut text);
}
if !extra.is_empty() {
text.push_str("\n\n# Agent instructions\n\n");
text.push_str(extra);
}
text
}
};
if text.is_empty() {
text = "You are a coding agent.".into();
}
SystemPrompt::Custom(text)
};
startup.diagnostics.update_tools(&tool_specs);
let workspace_root = sdk_options.workspace.root.clone();
let workspace = sdk_options.workspace.build_workspace()?;
let context_window = configured_context_window(startup.config);
let compaction = sdk_options.runtime.compaction.clone();
startup.diagnostics.update_compaction_config(&compaction);
let usage_recording = crate::usage::default_recording().await;
let runtime = build_runtime_with_max_steps(
RuntimeBuildOptions {
provider,
tools: tool_set.tools(),
workspace,
workspace_policy: AppPolicy::for_mode(startup.config.permission_mode),
approval_handler: None,
system_prompt,
reasoning: sdk_options.runtime.reasoning,
compaction,
context_window,
usage_purpose: startup.usage_purpose,
usage_parent_session_id: startup.parent_session_id.clone(),
usage_recording,
},
startup.max_steps,
)?;
let session = runtime.session(SessionOptions::default()).await?;
if let Some(adapter) = jsonl.as_deref_mut() {
adapter.set_run_context(session.id(), &workspace_root);
}
startup
.herdr
.report_state(HerdrState::Working, None, None)
.await;
let result = complete_run(
&session,
prompt_text,
HeadlessRunDeps {
reporter,
external_cancellation: cancellation,
jsonl,
host_input: startup.host_input.as_deref(),
},
)
.await;
runtime.shutdown();
tool_set.shutdown().await;
startup
.herdr
.report_state(HerdrState::Idle, None, None)
.await;
startup.herdr.release().await;
result
}
async fn complete_run(
session: &rho_sdk::Session,
prompt_text: String,
dependencies: HeadlessRunDeps<'_>,
) -> anyhow::Result<rho_sdk::RunOutcome> {
let HeadlessRunDeps {
reporter,
external_cancellation,
jsonl,
host_input,
} = dependencies;
let mut run = session.start(UserInput::text(prompt_text)).await?;
let cancellation = run.cancellation_handle();
let external_cancellation = external_cancellation.unwrap_or_default();
tokio::select! {
outcome = headless_run::drive(&mut run, reporter, jsonl, host_input) => outcome,
signal = shutdown_signal() => {
let signal = signal?;
cancellation.cancel();
let _ = run.outcome().await;
Err(AutomationInterrupted::new(signal).into())
}
() = external_cancellation.cancelled() => {
cancellation.cancel();
let _ = run.outcome().await;
Err(SubagentCancelled.into())
}
}
}
pub(crate) use crate::run_artifacts::RunArtifactIdentity;
pub(crate) struct RunReporter {
sink: crate::run_artifacts::RunArtifactSink,
adapter: crate::tui::event_adapter::SdkEventAdapter,
stream_output: bool,
}
impl RunReporter {
pub(crate) fn new(
path: PathBuf,
identity: RunArtifactIdentity,
cwd: PathBuf,
prompt: &str,
stream_output: bool,
status_tx: Option<tokio::sync::watch::Sender<RunStatus>>,
) -> anyhow::Result<Self> {
let sink = crate::run_artifacts::RunArtifactSink::open(path, &identity, prompt, status_tx)?;
Ok(Self {
sink,
adapter: crate::tui::event_adapter::SdkEventAdapter::new(cwd),
stream_output,
})
}
pub(crate) fn continue_from(
path: PathBuf,
started_status: RunStatus,
cwd: PathBuf,
prompt: &str,
stream_output: bool,
status_tx: Option<tokio::sync::watch::Sender<RunStatus>>,
) -> anyhow::Result<Self> {
let sink = crate::run_artifacts::RunArtifactSink::continue_from(
path,
started_status,
prompt,
status_tx,
)?;
Ok(Self {
sink,
adapter: crate::tui::event_adapter::SdkEventAdapter::new(cwd),
stream_output,
})
}
pub(super) fn on_event(&mut self, event: &rho_sdk::RunEvent) {
use rho_sdk::RunEvent;
let attachments = crate::tui::translate_run_event(&mut self.adapter, event);
if !attachments.is_empty() {
let mut saw_text_delta = false;
let mut needs_immediate_publish = false;
for attachment in attachments {
match &attachment {
crate::run_artifacts::AttachmentEvent::AssistantTextDelta(text)
if !text.is_empty() =>
{
self.sink.append_last_text(text);
saw_text_delta = true;
self.sink.write_attachment(attachment);
}
_ => {
needs_immediate_publish = true;
self.sink.write_attachment(attachment);
}
}
}
if needs_immediate_publish {
self.sink.publish();
} else if saw_text_delta {
self.sink.publish_throttled();
}
}
match event {
RunEvent::StepStarted { step } => {
self.sink.status.state = RunState::Running;
self.sink.status.turns = *step as u64;
self.sink.publish();
}
RunEvent::ToolStarted { name, .. } => {
self.sink.status.last_activity = Some(format!("tool: {name}"));
self.stream(&format!("\n[tool] {name}\n"));
self.sink.publish();
}
RunEvent::HostInputRequested { request }
| RunEvent::ToolHostInputRequested { request, .. } => {
self.sink.status.last_activity =
Some(format!("waiting for questionnaire: {}", request.title()));
self.sink.publish();
}
RunEvent::AssistantTextDelta { text } => {
self.sink.status.last_activity = Some("assistant text".into());
self.stream(text);
}
RunEvent::ProviderStreamReset { .. } => {
self.sink.status.last_activity = Some("retrying provider response".into());
self.sink.status.last_text = None;
self.stream("\n[provider response discarded; retrying]\n");
self.sink.publish();
}
RunEvent::UsageUpdated { usage } => {
self.sink.status.input_tokens = usage.total_input_tokens();
self.sink.status.output_tokens = usage.output_tokens;
}
_ => {}
}
}
#[cfg(test)]
pub(crate) fn status(&self) -> &RunStatus {
&self.sink.status
}
pub(super) fn write(&mut self) {
self.sink.publish();
}
pub(crate) fn finish(&mut self, result: &anyhow::Result<rho_sdk::RunOutcome>) {
match result {
Ok(outcome) => {
let usage = outcome.usage();
self.sink.status.input_tokens = usage.total_input_tokens();
self.sink.status.output_tokens = usage.output_tokens;
self.sink.finish_ok(Some(outcome.text().to_string()));
}
Err(error)
if error.is::<AutomationInterrupted>()
|| error.downcast_ref::<AutomationExit>().is_some_and(|exit| {
matches!(
exit.reason(),
TerminalReason::MaxSteps | TerminalReason::Timeout
)
})
|| error.is::<SubagentCancelled>() =>
{
self.sink.finish_stopped("stopped");
}
Err(error) => {
self.sink.finish_error(format!("{error:#}"));
}
}
}
fn stream(&self, text: &str) {
if !self.stream_output {
return;
}
let mut stdout = io::stdout().lock();
let _ = stdout.write_all(text.as_bytes());
let _ = stdout.flush();
}
}
#[cfg(unix)]
async fn shutdown_signal() -> io::Result<ShutdownSignal> {
use tokio::signal::unix::{signal, SignalKind};
let mut interrupt = signal(SignalKind::interrupt())?;
let mut terminate = signal(SignalKind::terminate())?;
tokio::select! {
_ = interrupt.recv() => Ok(ShutdownSignal::Interrupt),
_ = terminate.recv() => Ok(ShutdownSignal::Terminate),
}
}
#[cfg(not(unix))]
async fn shutdown_signal() -> io::Result<ShutdownSignal> {
tokio::signal::ctrl_c().await?;
Ok(ShutdownSignal::Interrupt)
}
fn prompt_from_stdin(parts: Vec<String>, read_stdin: bool) -> anyhow::Result<String> {
prompt_from_reader(parts, read_stdin, &mut io::stdin())
}
fn prompt_from_reader(
parts: Vec<String>,
read_stdin: bool,
stdin: &mut impl Read,
) -> anyhow::Result<String> {
let mut chunks = Vec::new();
let inline = parts.join(" ").trim().to_string();
if !inline.is_empty() {
chunks.push(inline);
}
if read_stdin {
let mut buffer = String::new();
stdin.read_to_string(&mut buffer)?;
let buffer = buffer.trim().to_string();
if !buffer.is_empty() {
chunks.push(buffer);
}
}
let prompt = chunks.join("\n\n");
if prompt.is_empty() {
anyhow::bail!("rho run requires a prompt argument or --stdin");
}
Ok(prompt)
}
#[cfg(test)]
#[path = "automation_tests.rs"]
mod tests;