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, ToolsAndPrompt, 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    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    // The reporter exists before anything that can fail, so a parent process
183    // watching the output file always sees a terminal state, including startup failures.
184    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                /* stream_output */ 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
308/// Single classification of a finished automation session before reporter,
309/// JSONL/text emission, and process exit share one decision.
310enum 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
378/// Builds the human-readable terminal message for a classified automation error.
379///
380/// Reason codes stay stable machine labels. The message carries actionable detail
381/// except for authentication failures, which stay generic so credentials and
382/// token material never leave the process through stdout, stderr, or JSONL.
383fn 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        // Automation keeps questionnaire capability when the agent exposes it.
473        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, /* completed_runs */ 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
611/// Maintains the `--output-file` status contract for subagent runs and
612/// streams progress to stdout so a watching pane shows live activity.
613pub(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    /// Resume after the executor already wrote the Starting boundary.
637    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                // Attachment path already published throttled when translated.
707            }
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;