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, ToolsAndPrompt, 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 pub approval_session: Option<rho_sdk::ApprovalSession>,
148 pub hook_host_labels: rho_sdk::hooks::HookHostLabels,
149}
150
151pub(super) fn prompt_for_command(command: &Option<Command>) -> anyhow::Result<Option<String>> {
152 match command {
153 Some(Command::Run { prompt, stdin, .. }) => {
154 prompt_from_stdin(prompt.clone(), *stdin).map(Some)
155 }
156 Some(
157 Command::Attach { .. }
158 | Command::Login { .. }
159 | Command::CredentialStore { .. }
160 | Command::Sessions { .. }
161 | Command::Mcp { .. }
162 | Command::Plugins { .. }
163 | Command::Workflow { .. }
164 | Command::WorkflowPlannerWorker
165 | Command::Update,
166 )
167 | None => Ok(None),
168 }
169}
170
171pub(super) fn emit_startup_failure(message: impl Into<String>) -> anyhow::Result<()> {
172 let mut adapter = JsonlAdapter::new();
173 let event = adapter.failed(TerminalReason::ConfigurationError, message.into(), None);
174 emit(event)
175}
176
177pub(super) async fn run(prompt_text: String, startup: Startup<'_>) -> anyhow::Result<()> {
178 let mut jsonl = (startup.output == OutputFormat::Jsonl).then(JsonlAdapter::new);
179 let deadline = startup
180 .timeout
181 .map(|timeout| tokio::time::Instant::now() + timeout);
182 let reporter_result = startup
185 .output_file
186 .as_ref()
187 .map(|path| {
188 RunReporter::new(
189 path.clone(),
190 RunArtifactIdentity {
191 agent_id: startup.agent.id().to_string(),
192 agent_fingerprint: startup.agent.fingerprint().to_string(),
193 provider: startup.config.provider.clone(),
194 model: startup.config.model.clone(),
195 },
196 startup.cwd.clone(),
197 &prompt_text,
198 startup.output == OutputFormat::Text,
199 None,
200 )
201 })
202 .transpose();
203 let mut reporter = match reporter_result {
204 Ok(reporter) => reporter,
205 Err(error) => {
206 emit_failure(&mut jsonl, TerminalReason::OutputError, &error)?;
207 return Err(
208 AutomationExit::new(1, TerminalReason::OutputError, error.to_string()).into(),
209 );
210 }
211 };
212
213 let cancellation = rho_tools::cancellation::RunCancellation::default();
214 let (result, timed_out) = if let Some(deadline) = deadline {
215 let future = run_session_with_output(
216 prompt_text,
217 &startup,
218 reporter.as_mut(),
219 Some(cancellation.clone()),
220 jsonl.as_mut(),
221 );
222 tokio::pin!(future);
223 tokio::select! {
224 result = &mut future => (result, false),
225 () = tokio::time::sleep_until(deadline) => {
226 cancellation.cancel();
227 (future.await, true)
228 }
229 }
230 } else {
231 (
232 run_session_with_output(
233 prompt_text,
234 &startup,
235 reporter.as_mut(),
236 None,
237 jsonl.as_mut(),
238 )
239 .await,
240 false,
241 )
242 };
243 let terminal = classify_run_terminal(result, timed_out);
244 if let Some(reporter) = reporter.as_mut() {
245 reporter.finish_terminal(&terminal);
246 }
247 emit_and_exit_terminal(terminal, &mut jsonl, reporter.is_some())
248}
249
250fn write_text_answer(answer: &rho_sdk::RunOutcome, has_reporter: bool) -> anyhow::Result<()> {
251 let result = (|| -> io::Result<()> {
252 let mut stdout = io::stdout().lock();
253 if has_reporter {
254 writeln!(stdout, "\n[subagent run complete]")?;
255 } else {
256 writeln!(stdout, "{}", answer.text())?;
257 }
258 stdout.flush()
259 })();
260 result.map_err(|error| {
261 AutomationExit::new(
262 1,
263 TerminalReason::OutputError,
264 format!("could not write output: {error}"),
265 )
266 .into()
267 })
268}
269
270pub(super) fn emit(event: WireEvent) -> anyhow::Result<()> {
271 let mut stdout = io::stdout().lock();
272 write_event(&mut stdout, &event).map_err(|error| {
273 AutomationExit::new(
274 1,
275 TerminalReason::OutputError,
276 format!("could not write JSONL output: {error}"),
277 )
278 .into()
279 })
280}
281
282fn emit_stopped(adapter: &mut Option<JsonlAdapter>, reason: TerminalReason) -> anyhow::Result<()> {
283 if let Some(adapter) = adapter.as_mut() {
284 let text = adapter.partial_text();
285 let event = adapter.stopped(reason, text);
286 emit(event)?;
287 }
288 Ok(())
289}
290
291fn emit_failure(
292 adapter: &mut Option<JsonlAdapter>,
293 reason: TerminalReason,
294 error: &anyhow::Error,
295) -> anyhow::Result<()> {
296 if let Some(adapter) = adapter.as_mut() {
297 let text = adapter.partial_text();
298 let message = terminal_error_message(reason, error);
299 let event = adapter.failed(reason, message, text);
300 emit(event)?;
301 }
302 Ok(())
303}
304
305const MAX_STEPS_MESSAGE: &str = "rho run reached its model-step limit";
306const TIMEOUT_MESSAGE: &str = "rho run timed out";
307
308enum RunTerminal {
311 Completed(rho_sdk::RunOutcome),
312 MaxSteps(rho_sdk::RunOutcome),
313 Timeout,
314 Failed(anyhow::Error),
315}
316
317fn classify_run_terminal(
318 result: anyhow::Result<rho_sdk::RunOutcome>,
319 timed_out: bool,
320) -> RunTerminal {
321 if timed_out {
322 return RunTerminal::Timeout;
323 }
324 match result {
325 Ok(answer) if answer.stop_reason() == rho_sdk::StopReason::MaxSteps => {
326 RunTerminal::MaxSteps(answer)
327 }
328 Ok(answer) => RunTerminal::Completed(answer),
329 Err(error) => RunTerminal::Failed(error),
330 }
331}
332
333fn emit_and_exit_terminal(
334 terminal: RunTerminal,
335 jsonl: &mut Option<JsonlAdapter>,
336 has_reporter: bool,
337) -> anyhow::Result<()> {
338 match terminal {
339 RunTerminal::Timeout => {
340 emit_stopped(jsonl, TerminalReason::Timeout)?;
341 Err(AutomationExit::new(124, TerminalReason::Timeout, TIMEOUT_MESSAGE).into())
342 }
343 RunTerminal::MaxSteps(answer) => {
344 if let Some(adapter) = jsonl.as_mut() {
345 let text = (!answer.text().is_empty()).then(|| answer.text().into());
346 let event = adapter.stopped(TerminalReason::MaxSteps, text);
347 emit(event)?;
348 } else {
349 write_text_answer(&answer, has_reporter)?;
350 }
351 Err(AutomationExit::new(124, TerminalReason::MaxSteps, MAX_STEPS_MESSAGE).into())
352 }
353 RunTerminal::Completed(answer) => {
354 if let Some(adapter) = jsonl.as_mut() {
355 let event = adapter.completed(answer.text().into());
356 emit(event)?;
357 } else {
358 write_text_answer(&answer, has_reporter)?;
359 }
360 Ok(())
361 }
362 RunTerminal::Failed(error) => {
363 let (reason, code) = classify_error(&error);
364 if reason == TerminalReason::Interrupted {
365 emit_stopped(jsonl, reason)?;
366 } else if reason != TerminalReason::OutputError {
367 emit_failure(jsonl, reason, &error)?;
368 }
369 let message = terminal_error_message(reason, &error);
370 if error.is::<AutomationInterrupted>() {
371 return Err(error);
372 }
373 Err(AutomationExit::new(code, reason, message).into())
374 }
375 }
376}
377
378fn terminal_error_message(reason: TerminalReason, error: &anyhow::Error) -> String {
384 match reason {
385 TerminalReason::Authentication => "authentication failed".to_string(),
386 TerminalReason::ProviderError
387 | TerminalReason::ToolHostError
388 | TerminalReason::ConfigurationError
389 | TerminalReason::OutputError
390 | TerminalReason::OtherError
391 | TerminalReason::Interrupted
392 | TerminalReason::MaxSteps
393 | TerminalReason::Timeout
394 | TerminalReason::Completed => error.to_string(),
395 }
396}
397
398fn classify_error(error: &anyhow::Error) -> (TerminalReason, u8) {
399 if let Some(interrupted) = error.downcast_ref::<AutomationInterrupted>() {
400 return (TerminalReason::Interrupted, interrupted.exit_code());
401 }
402 if let Some(exit) = error.downcast_ref::<AutomationExit>() {
403 return (exit.reason(), exit.exit_code());
404 }
405 for cause in error.chain() {
406 if let Some(error) = cause.downcast_ref::<rho_sdk::Error>() {
407 return match error {
408 rho_sdk::Error::Authentication { .. } => (TerminalReason::Authentication, 1),
409 rho_sdk::Error::Provider(provider)
410 if provider.kind() == rho_sdk::ProviderErrorKind::Authentication =>
411 {
412 (TerminalReason::Authentication, 1)
413 }
414 rho_sdk::Error::Provider(_) => (TerminalReason::ProviderError, 1),
415 rho_sdk::Error::Tool(_) => (TerminalReason::ToolHostError, 1),
416 rho_sdk::Error::InvalidConfiguration { .. } => {
417 (TerminalReason::ConfigurationError, 2)
418 }
419 _ => (TerminalReason::OtherError, 1),
420 };
421 }
422 if let Some(error) = cause.downcast_ref::<rho_providers::model::ModelError>() {
423 use rho_providers::model::ModelError;
424 return match error {
425 ModelError::MissingCredentials(_) | ModelError::Credentials(_) => {
426 (TerminalReason::Authentication, 1)
427 }
428 ModelError::UnsupportedReasoning { .. } | ModelError::UnsupportedProvider(_) => {
429 (TerminalReason::ConfigurationError, 2)
430 }
431 _ => (TerminalReason::ProviderError, 1),
432 };
433 }
434 }
435 (TerminalReason::OtherError, 1)
436}
437
438pub(crate) async fn run_session(
439 prompt_text: String,
440 startup: &Startup<'_>,
441 reporter: Option<&mut RunReporter>,
442 cancellation: Option<rho_tools::cancellation::RunCancellation>,
443) -> anyhow::Result<rho_sdk::RunOutcome> {
444 run_session_with_output(prompt_text, startup, reporter, cancellation, None).await
445}
446
447async fn run_session_with_output(
448 prompt_text: String,
449 startup: &Startup<'_>,
450 reporter: Option<&mut RunReporter>,
451 cancellation: Option<rho_tools::cancellation::RunCancellation>,
452 mut jsonl: Option<&mut JsonlAdapter>,
453) -> anyhow::Result<rho_sdk::RunOutcome> {
454 let sdk_options = SdkBootstrapOptions::from_config(startup.config, &startup.cwd)?;
455 let credentials = rho_providers::auth::provider_credentials::ApplicationCredentialSource::new(
456 Arc::new(AppCredentialStore),
457 );
458 let provider = build_automation_provider(sdk_options.provider, &credentials)?;
459 let workspace_root = sdk_options.workspace.root.clone();
460 let workspace = sdk_options.workspace.build_workspace()?;
461 let ToolsAndPrompt {
462 tools: tool_set,
463 system_prompt,
464 ..
465 } = assemble_tools_and_prompt(ToolsAndPromptOptions {
466 config: startup.config,
467 config_path: startup.config_path.clone(),
468 cwd: &startup.cwd,
469 no_system_prompt: startup.no_system_prompt,
470 no_tools: startup.no_tools,
471 no_subagents: startup.no_subagents,
472 questionnaire_enabled: true,
474 background_subagents: BackgroundSubagents::Disabled,
475 diagnostics: &startup.diagnostics,
476 agent: &startup.agent,
477 })
478 .await?;
479
480 let context_window = configured_context_window(startup.config);
481 let compaction = sdk_options.runtime.compaction.clone();
482 startup.diagnostics.update_compaction_config(&compaction);
483 let usage_recording = crate::usage::default_recording().await;
484 let hooks = crate::hooks::start_for_cwd(&workspace_root);
485 if let Some(hooks) = hooks.as_ref() {
486 startup.diagnostics.attach_hooks(hooks);
487 }
488 let startup_result: anyhow::Result<_> = async {
489 let runtime = build_runtime_with_max_steps(
490 RuntimeBuildOptions {
491 provider,
492 tools: tool_set.tools(),
493 workspace,
494 workspace_policy: AppPolicy::for_mode(startup.config.permission_mode),
495 approval_session: startup.approval_session.clone(),
496 system_prompt: system_prompt.for_advisor_mode(tool_set.advisor_registered()),
497 reasoning: sdk_options.runtime.reasoning,
498 service_tier: sdk_options.runtime.service_tier,
499 compaction,
500 context_window,
501 usage_purpose: startup.usage_purpose,
502 usage_parent_session_id: startup.parent_session_id.clone(),
503 usage_recording,
504 hook_host_labels: startup.hook_host_labels.clone(),
505 hooks: hooks.as_ref(),
506 },
507 startup.max_steps,
508 )?;
509 let session = match runtime.session(SessionOptions::default()).await {
510 Ok(session) => session,
511 Err(error) => {
512 runtime.shutdown();
513 return Err(error.into());
514 }
515 };
516 anyhow::Ok((runtime, session))
517 }
518 .await;
519 let (runtime, session) = match startup_result {
520 Ok(startup) => startup,
521 Err(error) => {
522 if let Some(hooks) = hooks {
523 hooks.shutdown(crate::hooks::DRAIN_GRACE).await;
524 }
525 tool_set.shutdown().await;
526 return Err(error);
527 }
528 };
529 if let Some(advisor) = tool_set.advisor() {
530 advisor.bind_session(session.clone());
531 }
532 if let Some(adapter) = jsonl.as_deref_mut() {
533 adapter.set_run_context(session.id(), &workspace_root);
534 }
535 startup
536 .herdr
537 .report_state(HerdrState::Working, None, None)
538 .await;
539 let result = complete_run(
540 &session,
541 prompt_text,
542 HeadlessRunDeps {
543 reporter,
544 external_cancellation: cancellation,
545 jsonl,
546 host_input: startup.host_input.as_deref(),
547 },
548 )
549 .await;
550
551 let session_hooks = runtime.hooks();
552 let session_id = session.id().clone();
553 match &result {
554 Ok(_) => {
555 session_hooks.session_completed(&session_id, 1)
556 }
557 Err(error) => session_hooks.session_failed(
558 &session_id,
559 rho_sdk::hooks::HookSessionFailureKind::RunFailed,
560 &error.to_string(),
561 ),
562 }
563 runtime.shutdown();
564 drop(session);
565 drop(runtime);
566 if let Some(hooks) = hooks {
567 hooks.shutdown(crate::hooks::DRAIN_GRACE).await;
568 }
569 tool_set.shutdown().await;
570 startup
571 .herdr
572 .report_state(HerdrState::Idle, None, None)
573 .await;
574 startup.herdr.release().await;
575
576 result
577}
578
579async fn complete_run(
580 session: &rho_sdk::Session,
581 prompt_text: String,
582 dependencies: HeadlessRunDeps<'_>,
583) -> anyhow::Result<rho_sdk::RunOutcome> {
584 let HeadlessRunDeps {
585 reporter,
586 external_cancellation,
587 jsonl,
588 host_input,
589 } = dependencies;
590 let mut run = session.start(UserInput::text(prompt_text)).await?;
591 let cancellation = run.cancellation_handle();
592 let external_cancellation = external_cancellation.unwrap_or_default();
593 tokio::select! {
594 outcome = headless_run::drive(&mut run, reporter, jsonl, host_input) => outcome,
595 signal = shutdown_signal() => {
596 let signal = signal?;
597 cancellation.cancel();
598 let _ = run.outcome().await;
599 Err(AutomationInterrupted::new(signal).into())
600 }
601 () = external_cancellation.cancelled() => {
602 cancellation.cancel();
603 let _ = run.outcome().await;
604 Err(SubagentCancelled.into())
605 }
606 }
607}
608
609pub(crate) use crate::run_artifacts::RunArtifactIdentity;
610
611pub(crate) struct RunReporter {
614 sink: crate::run_artifacts::RunArtifactSink,
615 adapter: crate::tui::event_adapter::SdkEventAdapter,
616 stream_output: bool,
617}
618
619impl RunReporter {
620 pub(crate) fn new(
621 path: PathBuf,
622 identity: RunArtifactIdentity,
623 cwd: PathBuf,
624 prompt: &str,
625 stream_output: bool,
626 status_tx: Option<tokio::sync::watch::Sender<RunStatus>>,
627 ) -> anyhow::Result<Self> {
628 let sink = crate::run_artifacts::RunArtifactSink::open(path, &identity, prompt, status_tx)?;
629 Ok(Self {
630 sink,
631 adapter: crate::tui::event_adapter::SdkEventAdapter::new(cwd),
632 stream_output,
633 })
634 }
635
636 pub(crate) fn continue_from(
638 path: PathBuf,
639 started_status: RunStatus,
640 cwd: PathBuf,
641 prompt: &str,
642 stream_output: bool,
643 status_tx: Option<tokio::sync::watch::Sender<RunStatus>>,
644 ) -> anyhow::Result<Self> {
645 let sink = crate::run_artifacts::RunArtifactSink::continue_from(
646 path,
647 started_status,
648 prompt,
649 status_tx,
650 )?;
651 Ok(Self {
652 sink,
653 adapter: crate::tui::event_adapter::SdkEventAdapter::new(cwd),
654 stream_output,
655 })
656 }
657
658 pub(super) fn on_event(&mut self, event: &rho_sdk::RunEvent) {
659 use rho_sdk::RunEvent;
660
661 let attachments = crate::tui::translate_run_event(&mut self.adapter, event);
662 if !attachments.is_empty() {
663 let mut saw_text_delta = false;
664 let mut needs_immediate_publish = false;
665 for attachment in attachments {
666 match &attachment {
667 crate::run_artifacts::AttachmentEvent::AssistantTextDelta(text)
668 if !text.is_empty() =>
669 {
670 self.sink.append_last_text(text);
671 saw_text_delta = true;
672 self.sink.write_attachment(attachment);
673 }
674 _ => {
675 needs_immediate_publish = true;
676 self.sink.write_attachment(attachment);
677 }
678 }
679 }
680 if needs_immediate_publish {
681 self.sink.publish();
682 } else if saw_text_delta {
683 self.sink.publish_throttled();
684 }
685 }
686 match event {
687 RunEvent::StepStarted { step, .. } => {
688 self.sink.status.state = RunState::Running;
689 self.sink.status.turns = *step as u64;
690 self.sink.publish();
691 }
692 RunEvent::ToolStarted { name, .. } => {
693 self.sink.status.last_activity = Some(format!("tool: {name}"));
694 self.stream(&format!("\n[tool] {name}\n"));
695 self.sink.publish();
696 }
697 RunEvent::HostInputRequested { request }
698 | RunEvent::ToolHostInputRequested { request, .. } => {
699 self.sink.status.last_activity =
700 Some(format!("waiting for questionnaire: {}", request.title()));
701 self.sink.publish();
702 }
703 RunEvent::AssistantTextDelta { text } => {
704 self.sink.status.last_activity = Some("assistant text".into());
705 self.stream(text);
706 }
708 RunEvent::ProviderStreamReset { .. } => {
709 self.sink.status.last_activity = Some("retrying provider response".into());
710 self.sink.status.last_text = None;
711 self.stream("\n[provider response discarded; retrying]\n");
712 self.sink.publish();
713 }
714 RunEvent::UsageUpdated { usage } => {
715 self.sink.status.input_tokens = usage.total_input_tokens();
716 self.sink.status.output_tokens = usage.output_tokens;
717 }
718 _ => {}
719 }
720 }
721
722 #[cfg(test)]
723 pub(crate) fn status(&self) -> &RunStatus {
724 &self.sink.status
725 }
726
727 pub(super) fn write(&mut self) {
728 self.sink.publish();
729 }
730
731 pub(crate) fn finish(&mut self, result: &anyhow::Result<rho_sdk::RunOutcome>) {
732 match result {
733 Ok(outcome) => {
734 let usage = outcome.usage();
735 self.sink.status.input_tokens = usage.total_input_tokens();
736 self.sink.status.output_tokens = usage.output_tokens;
737 self.sink.finish_ok(Some(outcome.text().to_string()));
738 }
739 Err(error)
740 if error.is::<AutomationInterrupted>()
741 || error.downcast_ref::<AutomationExit>().is_some_and(|exit| {
742 matches!(
743 exit.reason(),
744 TerminalReason::MaxSteps | TerminalReason::Timeout
745 )
746 })
747 || error.is::<SubagentCancelled>() =>
748 {
749 self.sink.finish_stopped("stopped");
750 }
751 Err(error) => {
752 self.sink.finish_error(format!("{error:#}"));
753 }
754 }
755 }
756
757 fn finish_terminal(&mut self, terminal: &RunTerminal) {
758 match terminal {
759 RunTerminal::Completed(outcome) => {
760 let usage = outcome.usage();
761 self.sink.status.input_tokens = usage.total_input_tokens();
762 self.sink.status.output_tokens = usage.output_tokens;
763 self.sink.finish_ok(Some(outcome.text().to_string()));
764 }
765 RunTerminal::MaxSteps(_) | RunTerminal::Timeout => {
766 self.sink.finish_stopped("stopped");
767 }
768 RunTerminal::Failed(error)
769 if error.is::<AutomationInterrupted>() || error.is::<SubagentCancelled>() =>
770 {
771 self.sink.finish_stopped("stopped");
772 }
773 RunTerminal::Failed(error) => {
774 self.sink.finish_error(format!("{error:#}"));
775 }
776 }
777 }
778
779 fn stream(&self, text: &str) {
780 if !self.stream_output {
781 return;
782 }
783 let mut stdout = io::stdout().lock();
784 let _ = stdout.write_all(text.as_bytes());
785 let _ = stdout.flush();
786 }
787}
788
789#[cfg(unix)]
790async fn shutdown_signal() -> io::Result<ShutdownSignal> {
791 use tokio::signal::unix::{signal, SignalKind};
792
793 let mut interrupt = signal(SignalKind::interrupt())?;
794 let mut terminate = signal(SignalKind::terminate())?;
795 tokio::select! {
796 _ = interrupt.recv() => Ok(ShutdownSignal::Interrupt),
797 _ = terminate.recv() => Ok(ShutdownSignal::Terminate),
798 }
799}
800
801#[cfg(not(unix))]
802async fn shutdown_signal() -> io::Result<ShutdownSignal> {
803 tokio::signal::ctrl_c().await?;
804 Ok(ShutdownSignal::Interrupt)
805}
806
807fn prompt_from_stdin(parts: Vec<String>, read_stdin: bool) -> anyhow::Result<String> {
808 if !read_stdin && crate::stdio::stdin_is_redirected() {
809 anyhow::bail!(
810 "stdin is redirected but --stdin was not set; pass --stdin to include piped input"
811 );
812 }
813 prompt_from_reader(parts, read_stdin, &mut io::stdin())
814}
815
816fn prompt_from_reader(
817 parts: Vec<String>,
818 read_stdin: bool,
819 stdin: &mut impl Read,
820) -> anyhow::Result<String> {
821 let mut chunks = Vec::new();
822 let inline = parts.join(" ").trim().to_string();
823 if !inline.is_empty() {
824 chunks.push(inline);
825 }
826 if read_stdin {
827 let mut buffer = String::new();
828 stdin.read_to_string(&mut buffer)?;
829 let buffer = buffer.trim().to_string();
830 if !buffer.is_empty() {
831 chunks.push(buffer);
832 }
833 }
834
835 let prompt = chunks.join("\n\n");
836 if prompt.is_empty() {
837 anyhow::bail!("rho run requires a prompt argument or --stdin");
838 }
839 Ok(prompt)
840}
841
842#[cfg(test)]
843#[path = "automation_tests.rs"]
844mod tests;