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