Skip to main content

rho_coding_agent/app/
automation.rs

1use std::{
2    fmt,
3    io::{self, Read, Write},
4    num::NonZeroUsize,
5    path::PathBuf,
6    sync::Arc,
7    time::Duration,
8};
9
10use rho_sdk::{SessionOptions, SystemPrompt, UserInput};
11
12use {
13    crate::agent::{PromptPolicy, ToolCapability},
14    crate::cli::{Command, OutputFormat},
15    crate::config::Config,
16    crate::credential_store::AppCredentialStore,
17    crate::diagnostics::RuntimeDiagnostics,
18    crate::herdr::{HerdrReporter, HerdrState},
19    crate::prompt,
20    crate::subagent::{RunState, RunStatus},
21    crate::tools::{
22        agent::BackgroundSubagents,
23        sdk_registry::{AppToolSet, DelegationConfig, ToolSetOptions},
24    },
25    rho_providers::providers::build_automation_provider,
26};
27
28use super::{
29    agent_binding::BoundAgent,
30    automation_protocol::{write_event, JsonlAdapter, TerminalReason, WireEvent},
31    headless_run::{self, HeadlessRunDeps, HostInputResponder},
32    policy::AppPolicy,
33    runtime_builder::{
34        build_runtime_with_max_steps, configured_context_window, RuntimeBuildOptions,
35    },
36    sdk_config::SdkBootstrapOptions,
37};
38
39/// Error returned after an automation run has cleaned up and selected a stable exit code.
40#[derive(Debug)]
41pub struct AutomationExit {
42    code: u8,
43    reason: TerminalReason,
44    message: String,
45}
46
47impl AutomationExit {
48    pub(super) fn new(code: u8, reason: TerminalReason, message: impl Into<String>) -> Self {
49        Self {
50            code,
51            reason,
52            message: message.into(),
53        }
54    }
55
56    /// Returns the documented process exit code for this automation result.
57    pub fn exit_code(&self) -> u8 {
58        self.code
59    }
60
61    fn reason(&self) -> TerminalReason {
62        self.reason
63    }
64}
65
66impl fmt::Display for AutomationExit {
67    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
68        formatter.write_str(&self.message)
69    }
70}
71
72impl std::error::Error for AutomationExit {}
73
74/// Error returned after an automation run handles an interrupt and completes cleanup.
75#[derive(Debug)]
76pub struct AutomationInterrupted {
77    signal: ShutdownSignal,
78}
79
80impl AutomationInterrupted {
81    fn new(signal: ShutdownSignal) -> Self {
82        Self { signal }
83    }
84
85    /// Returns the conventional process exit code for the received signal.
86    pub fn exit_code(&self) -> u8 {
87        self.signal.exit_code()
88    }
89}
90
91impl fmt::Display for AutomationInterrupted {
92    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
93        write!(formatter, "rho run interrupted by {}", self.signal)
94    }
95}
96
97impl std::error::Error for AutomationInterrupted {}
98
99#[derive(Clone, Copy, Debug)]
100enum ShutdownSignal {
101    Interrupt,
102    Terminate,
103}
104
105impl ShutdownSignal {
106    fn exit_code(self) -> u8 {
107        match self {
108            Self::Interrupt => 130,
109            Self::Terminate => 143,
110        }
111    }
112}
113
114impl fmt::Display for ShutdownSignal {
115    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
116        match self {
117            Self::Interrupt => formatter.write_str("SIGINT"),
118            Self::Terminate => formatter.write_str("SIGTERM"),
119        }
120    }
121}
122
123#[derive(Debug)]
124struct SubagentCancelled;
125
126impl fmt::Display for SubagentCancelled {
127    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
128        formatter.write_str("subagent cancellation requested")
129    }
130}
131
132impl std::error::Error for SubagentCancelled {}
133
134pub(super) struct Startup<'a> {
135    pub config: &'a Config,
136    pub config_path: PathBuf,
137    pub cwd: PathBuf,
138    pub no_system_prompt: bool,
139    pub no_tools: bool,
140    pub no_subagents: bool,
141    pub usage_purpose: &'static str,
142    pub parent_session_id: Option<rho_sdk::SessionId>,
143    pub agent: BoundAgent,
144    pub output_file: Option<PathBuf>,
145    pub output: OutputFormat,
146    pub max_steps: Option<NonZeroUsize>,
147    pub timeout: Option<Duration>,
148    pub diagnostics: RuntimeDiagnostics,
149    pub herdr: HerdrReporter,
150    pub host_input: Option<Arc<dyn HostInputResponder>>,
151}
152
153pub(super) fn prompt_for_command(command: &Option<Command>) -> anyhow::Result<Option<String>> {
154    match command {
155        Some(Command::Run { prompt, stdin, .. }) => {
156            prompt_from_stdin(prompt.clone(), *stdin).map(Some)
157        }
158        Some(
159            Command::Attach { .. }
160            | Command::Login { .. }
161            | Command::CredentialStore { .. }
162            | Command::Sessions { .. }
163            | Command::Update,
164        )
165        | None => Ok(None),
166    }
167}
168
169pub(super) fn emit_startup_failure() -> anyhow::Result<()> {
170    let mut adapter = JsonlAdapter::new();
171    let event = adapter.failed(
172        TerminalReason::ConfigurationError,
173        "configuration failed".into(),
174        None,
175    );
176    emit(event)
177}
178
179pub(super) async fn run(prompt_text: String, startup: Startup<'_>) -> anyhow::Result<()> {
180    let mut jsonl = (startup.output == OutputFormat::Jsonl).then(JsonlAdapter::new);
181    let deadline = startup
182        .timeout
183        .map(|timeout| tokio::time::Instant::now() + timeout);
184    // The reporter exists before anything that can fail, so a parent process
185    // watching the output file always sees a terminal state, including startup failures.
186    let reporter_result = startup
187        .output_file
188        .as_ref()
189        .map(|path| {
190            RunReporter::new(
191                path.clone(),
192                RunArtifactIdentity {
193                    agent_id: startup.agent.id().to_string(),
194                    agent_fingerprint: startup.agent.fingerprint().to_string(),
195                    provider: startup.config.provider.clone(),
196                    model: startup.config.model.clone(),
197                },
198                startup.cwd.clone(),
199                &prompt_text,
200                /* stream_output */ startup.output == OutputFormat::Text,
201                None,
202            )
203        })
204        .transpose();
205    let mut reporter = match reporter_result {
206        Ok(reporter) => reporter,
207        Err(error) => {
208            emit_failure(&mut jsonl, TerminalReason::OutputError, &error)?;
209            return Err(
210                AutomationExit::new(1, TerminalReason::OutputError, error.to_string()).into(),
211            );
212        }
213    };
214
215    let cancellation = rho_tools::cancellation::RunCancellation::default();
216    let (result, timed_out) = if let Some(deadline) = deadline {
217        let future = run_session_with_output(
218            prompt_text,
219            &startup,
220            reporter.as_mut(),
221            Some(cancellation.clone()),
222            jsonl.as_mut(),
223        );
224        tokio::pin!(future);
225        tokio::select! {
226            result = &mut future => (result, false),
227            () = tokio::time::sleep_until(deadline) => {
228                cancellation.cancel();
229                (future.await, true)
230            }
231        }
232    } else {
233        (
234            run_session_with_output(
235                prompt_text,
236                &startup,
237                reporter.as_mut(),
238                None,
239                jsonl.as_mut(),
240            )
241            .await,
242            false,
243        )
244    };
245    if let Some(reporter) = reporter.as_mut() {
246        let reached_step_limit = result.as_ref().is_ok_and(|outcome| {
247            outcome.stop_reason() == rho_sdk::StopReason::MaxSteps
248                && (jsonl.is_some() || startup.max_steps.is_some())
249        });
250        if reached_step_limit {
251            let stopped = Err(AutomationExit::new(
252                124,
253                TerminalReason::MaxSteps,
254                "rho run reached its model-step limit",
255            )
256            .into());
257            reporter.finish(&stopped);
258        } else {
259            reporter.finish(&result);
260        }
261    }
262
263    if timed_out {
264        emit_stopped(&mut jsonl, TerminalReason::Timeout)?;
265        return Err(AutomationExit::new(124, TerminalReason::Timeout, "rho run timed out").into());
266    }
267
268    match result {
269        Ok(answer) => {
270            let max_steps = answer.stop_reason() == rho_sdk::StopReason::MaxSteps;
271            if max_steps && (jsonl.is_some() || startup.max_steps.is_some()) {
272                if let Some(adapter) = jsonl.as_mut() {
273                    let text = (!answer.text().is_empty()).then(|| answer.text().into());
274                    let event = adapter.stopped(TerminalReason::MaxSteps, text);
275                    emit(event)?;
276                } else {
277                    write_text_answer(&answer, reporter.is_some())?;
278                }
279                return Err(AutomationExit::new(
280                    124,
281                    TerminalReason::MaxSteps,
282                    "rho run reached its model-step limit",
283                )
284                .into());
285            }
286            if let Some(adapter) = jsonl.as_mut() {
287                let event = adapter.completed(answer.text().into());
288                emit(event)?;
289            } else {
290                write_text_answer(&answer, reporter.is_some())?;
291            }
292            Ok(())
293        }
294        Err(error) => {
295            let (reason, code) = classify_error(&error);
296            if reason == TerminalReason::Interrupted {
297                emit_stopped(&mut jsonl, reason)?;
298            } else if reason != TerminalReason::OutputError {
299                emit_failure(&mut jsonl, reason, &error)?;
300            }
301            let message = terminal_error_message(reason, &error);
302            if error.is::<AutomationInterrupted>() {
303                return Err(error);
304            }
305            Err(AutomationExit::new(code, reason, message).into())
306        }
307    }
308}
309
310fn write_text_answer(answer: &rho_sdk::RunOutcome, has_reporter: bool) -> anyhow::Result<()> {
311    let result = (|| -> io::Result<()> {
312        let mut stdout = io::stdout().lock();
313        if has_reporter {
314            writeln!(stdout, "\n[subagent run complete]")?;
315        } else {
316            writeln!(stdout, "{}", answer.text())?;
317        }
318        stdout.flush()
319    })();
320    result.map_err(|error| {
321        AutomationExit::new(
322            1,
323            TerminalReason::OutputError,
324            format!("could not write output: {error}"),
325        )
326        .into()
327    })
328}
329
330pub(super) fn emit(event: WireEvent) -> anyhow::Result<()> {
331    let mut stdout = io::stdout().lock();
332    write_event(&mut stdout, &event).map_err(|error| {
333        AutomationExit::new(
334            1,
335            TerminalReason::OutputError,
336            format!("could not write JSONL output: {error}"),
337        )
338        .into()
339    })
340}
341
342fn emit_stopped(adapter: &mut Option<JsonlAdapter>, reason: TerminalReason) -> anyhow::Result<()> {
343    if let Some(adapter) = adapter.as_mut() {
344        let text = adapter.partial_text();
345        let event = adapter.stopped(reason, text);
346        emit(event)?;
347    }
348    Ok(())
349}
350
351fn emit_failure(
352    adapter: &mut Option<JsonlAdapter>,
353    reason: TerminalReason,
354    error: &anyhow::Error,
355) -> anyhow::Result<()> {
356    if let Some(adapter) = adapter.as_mut() {
357        let text = adapter.partial_text();
358        let message = terminal_error_message(reason, error);
359        let event = adapter.failed(reason, message, text);
360        emit(event)?;
361    }
362    Ok(())
363}
364
365fn terminal_error_message(reason: TerminalReason, error: &anyhow::Error) -> String {
366    match reason {
367        TerminalReason::Authentication => "authentication failed".to_string(),
368        TerminalReason::ConfigurationError => "configuration failed".to_string(),
369        TerminalReason::OutputError => "output failed".to_string(),
370        TerminalReason::OtherError => "run failed".to_string(),
371        _ => error.to_string(),
372    }
373}
374
375fn classify_error(error: &anyhow::Error) -> (TerminalReason, u8) {
376    if let Some(interrupted) = error.downcast_ref::<AutomationInterrupted>() {
377        return (TerminalReason::Interrupted, interrupted.exit_code());
378    }
379    if let Some(exit) = error.downcast_ref::<AutomationExit>() {
380        return (exit.reason(), exit.exit_code());
381    }
382    for cause in error.chain() {
383        if let Some(error) = cause.downcast_ref::<rho_sdk::Error>() {
384            return match error {
385                rho_sdk::Error::Authentication { .. } => (TerminalReason::Authentication, 1),
386                rho_sdk::Error::Provider(provider)
387                    if provider.kind() == rho_sdk::ProviderErrorKind::Authentication =>
388                {
389                    (TerminalReason::Authentication, 1)
390                }
391                rho_sdk::Error::Provider(_) => (TerminalReason::ProviderError, 1),
392                rho_sdk::Error::Tool(_) => (TerminalReason::ToolHostError, 1),
393                rho_sdk::Error::InvalidConfiguration { .. } => {
394                    (TerminalReason::ConfigurationError, 2)
395                }
396                _ => (TerminalReason::OtherError, 1),
397            };
398        }
399        if let Some(error) = cause.downcast_ref::<rho_providers::model::ModelError>() {
400            use rho_providers::model::ModelError;
401            return match error {
402                ModelError::MissingApiKey
403                | ModelError::MissingCodexAuth
404                | ModelError::MissingAnthropicApiKey
405                | ModelError::MissingGoogleApiKey
406                | ModelError::MissingGithubCopilotAuth
407                | ModelError::MissingMoonshotApiKey
408                | ModelError::MissingPoolsideApiKey
409                | ModelError::MissingOpenRouterApiKey
410                | ModelError::MissingCredentialProfile(_)
411                | ModelError::MissingKimiAuth
412                | ModelError::MissingXaiApiKey
413                | ModelError::MissingXaiAuth
414                | ModelError::Credentials(_) => (TerminalReason::Authentication, 1),
415                ModelError::UnsupportedReasoning { .. } | ModelError::UnsupportedProvider(_) => {
416                    (TerminalReason::ConfigurationError, 2)
417                }
418                _ => (TerminalReason::ProviderError, 1),
419            };
420        }
421    }
422    (TerminalReason::OtherError, 1)
423}
424
425pub(crate) async fn run_session(
426    prompt_text: String,
427    startup: &Startup<'_>,
428    reporter: Option<&mut RunReporter>,
429    cancellation: Option<rho_tools::cancellation::RunCancellation>,
430) -> anyhow::Result<rho_sdk::RunOutcome> {
431    run_session_with_output(prompt_text, startup, reporter, cancellation, None).await
432}
433
434async fn run_session_with_output(
435    prompt_text: String,
436    startup: &Startup<'_>,
437    reporter: Option<&mut RunReporter>,
438    cancellation: Option<rho_tools::cancellation::RunCancellation>,
439    mut jsonl: Option<&mut JsonlAdapter>,
440) -> anyhow::Result<rho_sdk::RunOutcome> {
441    let sdk_options = SdkBootstrapOptions::from_config(startup.config, &startup.cwd)?;
442    let credentials = rho_providers::auth::provider_credentials::ApplicationCredentialSource::new(
443        Arc::new(AppCredentialStore),
444    );
445    let provider = build_automation_provider(sdk_options.provider, &credentials)?;
446    let mut capabilities = startup
447        .agent
448        .rho_capabilities()
449        .cloned()
450        .unwrap_or_default();
451    if startup.no_subagents {
452        capabilities.remove(&ToolCapability::Agent);
453        capabilities.remove(&ToolCapability::Agents);
454    }
455    let launch_delegation_enabled = capabilities.contains(&ToolCapability::Agent);
456    let delegation_enabled =
457        launch_delegation_enabled || capabilities.contains(&ToolCapability::Agents);
458    let tool_set = if startup.no_tools {
459        AppToolSet::disabled()
460    } else {
461        let mut options = ToolSetOptions::new(capabilities);
462        if delegation_enabled {
463            options = options.delegation(DelegationConfig::new(
464                startup.cwd.clone(),
465                startup.config_path.clone(),
466                BackgroundSubagents::Disabled,
467            ));
468        }
469        AppToolSet::new(startup.config, startup.diagnostics.clone(), options)
470    };
471    let tool_specs = tool_set.specs();
472    let system_prompt = if startup.no_system_prompt {
473        startup.diagnostics.update_prompt_sources(Vec::new());
474        SystemPrompt::None
475    } else {
476        let mut text = match startup.agent.prompt() {
477            PromptPolicy::Replace(text) => text.clone(),
478            PromptPolicy::Extend(extra) => {
479                let built = prompt::system_prompt(&tool_specs, &startup.cwd);
480                startup.diagnostics.update_prompt_sources(built.sources);
481                let mut text = built.text;
482                if !launch_delegation_enabled {
483                    prompt::append_subagents_disabled_instruction(&mut text);
484                }
485                if !extra.is_empty() {
486                    text.push_str("\n\n# Agent instructions\n\n");
487                    text.push_str(extra);
488                }
489                text
490            }
491        };
492        if text.is_empty() {
493            text = "You are a coding agent.".into();
494        }
495        SystemPrompt::Custom(text)
496    };
497    startup.diagnostics.update_tools(&tool_specs);
498
499    let workspace_root = sdk_options.workspace.root.clone();
500    let workspace = sdk_options.workspace.build_workspace()?;
501    let context_window = configured_context_window(startup.config);
502    let compaction = sdk_options.runtime.compaction.clone();
503    startup.diagnostics.update_compaction_config(&compaction);
504    let usage_recording = crate::usage::default_recording().await;
505    let runtime = build_runtime_with_max_steps(
506        RuntimeBuildOptions {
507            provider,
508            tools: tool_set.tools(),
509            workspace,
510            workspace_policy: AppPolicy::for_mode(startup.config.permission_mode),
511            approval_handler: None,
512            system_prompt,
513            reasoning: sdk_options.runtime.reasoning,
514            compaction,
515            context_window,
516            usage_purpose: startup.usage_purpose,
517            usage_parent_session_id: startup.parent_session_id.clone(),
518            usage_recording,
519        },
520        startup.max_steps,
521    )?;
522    let session = runtime.session(SessionOptions::default()).await?;
523    if let Some(adapter) = jsonl.as_deref_mut() {
524        adapter.set_run_context(session.id(), &workspace_root);
525    }
526    startup
527        .herdr
528        .report_state(HerdrState::Working, None, None)
529        .await;
530    let result = complete_run(
531        &session,
532        prompt_text,
533        HeadlessRunDeps {
534            reporter,
535            external_cancellation: cancellation,
536            jsonl,
537            host_input: startup.host_input.as_deref(),
538        },
539    )
540    .await;
541
542    runtime.shutdown();
543    tool_set.shutdown().await;
544    startup
545        .herdr
546        .report_state(HerdrState::Idle, None, None)
547        .await;
548    startup.herdr.release().await;
549
550    result
551}
552
553async fn complete_run(
554    session: &rho_sdk::Session,
555    prompt_text: String,
556    dependencies: HeadlessRunDeps<'_>,
557) -> anyhow::Result<rho_sdk::RunOutcome> {
558    let HeadlessRunDeps {
559        reporter,
560        external_cancellation,
561        jsonl,
562        host_input,
563    } = dependencies;
564    let mut run = session.start(UserInput::text(prompt_text)).await?;
565    let cancellation = run.cancellation_handle();
566    let external_cancellation = external_cancellation.unwrap_or_default();
567    tokio::select! {
568        outcome = headless_run::drive(&mut run, reporter, jsonl, host_input) => outcome,
569        signal = shutdown_signal() => {
570            let signal = signal?;
571            cancellation.cancel();
572            let _ = run.outcome().await;
573            Err(AutomationInterrupted::new(signal).into())
574        }
575        () = external_cancellation.cancelled() => {
576            cancellation.cancel();
577            let _ = run.outcome().await;
578            Err(SubagentCancelled.into())
579        }
580    }
581}
582
583pub(crate) use crate::run_artifacts::RunArtifactIdentity;
584
585/// Maintains the `--output-file` status contract for subagent runs and
586/// streams progress to stdout so a watching pane shows live activity.
587pub(crate) struct RunReporter {
588    sink: crate::run_artifacts::RunArtifactSink,
589    adapter: crate::tui::event_adapter::SdkEventAdapter,
590    stream_output: bool,
591}
592
593impl RunReporter {
594    pub(crate) fn new(
595        path: PathBuf,
596        identity: RunArtifactIdentity,
597        cwd: PathBuf,
598        prompt: &str,
599        stream_output: bool,
600        status_tx: Option<tokio::sync::watch::Sender<RunStatus>>,
601    ) -> anyhow::Result<Self> {
602        let sink = crate::run_artifacts::RunArtifactSink::open(path, &identity, prompt, status_tx)?;
603        Ok(Self {
604            sink,
605            adapter: crate::tui::event_adapter::SdkEventAdapter::new(cwd),
606            stream_output,
607        })
608    }
609
610    /// Resume after the executor already wrote the Starting boundary.
611    pub(crate) fn continue_from(
612        path: PathBuf,
613        started_status: RunStatus,
614        cwd: PathBuf,
615        prompt: &str,
616        stream_output: bool,
617        status_tx: Option<tokio::sync::watch::Sender<RunStatus>>,
618    ) -> anyhow::Result<Self> {
619        let sink = crate::run_artifacts::RunArtifactSink::continue_from(
620            path,
621            started_status,
622            prompt,
623            status_tx,
624        )?;
625        Ok(Self {
626            sink,
627            adapter: crate::tui::event_adapter::SdkEventAdapter::new(cwd),
628            stream_output,
629        })
630    }
631
632    pub(super) fn on_event(&mut self, event: &rho_sdk::RunEvent) {
633        use rho_sdk::RunEvent;
634
635        let attachments = crate::tui::translate_run_event(&mut self.adapter, event);
636        if !attachments.is_empty() {
637            let mut saw_text_delta = false;
638            let mut needs_immediate_publish = false;
639            for attachment in attachments {
640                match &attachment {
641                    crate::run_artifacts::AttachmentEvent::AssistantTextDelta(text)
642                        if !text.is_empty() =>
643                    {
644                        self.sink.append_last_text(text);
645                        saw_text_delta = true;
646                        self.sink.write_attachment(attachment);
647                    }
648                    _ => {
649                        needs_immediate_publish = true;
650                        self.sink.write_attachment(attachment);
651                    }
652                }
653            }
654            if needs_immediate_publish {
655                self.sink.publish();
656            } else if saw_text_delta {
657                self.sink.publish_throttled();
658            }
659        }
660        match event {
661            RunEvent::StepStarted { step } => {
662                self.sink.status.state = RunState::Running;
663                self.sink.status.turns = *step as u64;
664                self.sink.publish();
665            }
666            RunEvent::ToolStarted { name, .. } => {
667                self.sink.status.last_activity = Some(format!("tool: {name}"));
668                self.stream(&format!("\n[tool] {name}\n"));
669                self.sink.publish();
670            }
671            RunEvent::HostInputRequested { request }
672            | RunEvent::ToolHostInputRequested { request, .. } => {
673                self.sink.status.last_activity =
674                    Some(format!("waiting for questionnaire: {}", request.title()));
675                self.sink.publish();
676            }
677            RunEvent::AssistantTextDelta { text } => {
678                self.sink.status.last_activity = Some("assistant text".into());
679                self.stream(text);
680                // Attachment path already published throttled when translated.
681            }
682            RunEvent::ProviderStreamReset { .. } => {
683                self.sink.status.last_activity = Some("retrying provider response".into());
684                self.sink.status.last_text = None;
685                self.stream("\n[provider response discarded; retrying]\n");
686                self.sink.publish();
687            }
688            RunEvent::UsageUpdated { usage } => {
689                self.sink.status.input_tokens = usage.total_input_tokens();
690                self.sink.status.output_tokens = usage.output_tokens;
691            }
692            _ => {}
693        }
694    }
695
696    #[cfg(test)]
697    pub(crate) fn status(&self) -> &RunStatus {
698        &self.sink.status
699    }
700
701    pub(super) fn write(&mut self) {
702        self.sink.publish();
703    }
704
705    pub(crate) fn finish(&mut self, result: &anyhow::Result<rho_sdk::RunOutcome>) {
706        match result {
707            Ok(outcome) => {
708                let usage = outcome.usage();
709                self.sink.status.input_tokens = usage.total_input_tokens();
710                self.sink.status.output_tokens = usage.output_tokens;
711                self.sink.finish_ok(Some(outcome.text().to_string()));
712            }
713            Err(error)
714                if error.is::<AutomationInterrupted>()
715                    || error.downcast_ref::<AutomationExit>().is_some_and(|exit| {
716                        matches!(
717                            exit.reason(),
718                            TerminalReason::MaxSteps | TerminalReason::Timeout
719                        )
720                    })
721                    || error.is::<SubagentCancelled>() =>
722            {
723                self.sink.finish_stopped("stopped");
724            }
725            Err(error) => {
726                self.sink.finish_error(format!("{error:#}"));
727            }
728        }
729    }
730
731    fn stream(&self, text: &str) {
732        if !self.stream_output {
733            return;
734        }
735        let mut stdout = io::stdout().lock();
736        let _ = stdout.write_all(text.as_bytes());
737        let _ = stdout.flush();
738    }
739}
740
741#[cfg(unix)]
742async fn shutdown_signal() -> io::Result<ShutdownSignal> {
743    use tokio::signal::unix::{signal, SignalKind};
744
745    let mut interrupt = signal(SignalKind::interrupt())?;
746    let mut terminate = signal(SignalKind::terminate())?;
747    tokio::select! {
748        _ = interrupt.recv() => Ok(ShutdownSignal::Interrupt),
749        _ = terminate.recv() => Ok(ShutdownSignal::Terminate),
750    }
751}
752
753#[cfg(not(unix))]
754async fn shutdown_signal() -> io::Result<ShutdownSignal> {
755    tokio::signal::ctrl_c().await?;
756    Ok(ShutdownSignal::Interrupt)
757}
758
759fn prompt_from_stdin(parts: Vec<String>, read_stdin: bool) -> anyhow::Result<String> {
760    prompt_from_reader(parts, read_stdin, &mut io::stdin())
761}
762
763fn prompt_from_reader(
764    parts: Vec<String>,
765    read_stdin: bool,
766    stdin: &mut impl Read,
767) -> anyhow::Result<String> {
768    let mut chunks = Vec::new();
769    let inline = parts.join(" ").trim().to_string();
770    if !inline.is_empty() {
771        chunks.push(inline);
772    }
773    if read_stdin {
774        let mut buffer = String::new();
775        stdin.read_to_string(&mut buffer)?;
776        let buffer = buffer.trim().to_string();
777        if !buffer.is_empty() {
778            chunks.push(buffer);
779        }
780    }
781
782    let prompt = chunks.join("\n\n");
783    if prompt.is_empty() {
784        anyhow::bail!("rho run requires a prompt argument or --stdin");
785    }
786    Ok(prompt)
787}
788
789#[cfg(test)]
790#[path = "automation_tests.rs"]
791mod tests;