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