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#[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 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#[derive(Debug)]
72pub struct AutomationInterrupted {
73 signal: ShutdownSignal,
74}
75
76impl AutomationInterrupted {
77 fn new(signal: ShutdownSignal) -> Self {
78 Self { signal }
79 }
80
81 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 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 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 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, 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
555pub(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 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 }
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;