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::agent::PERMISSION_CLASSIFIER_AGENT_ID,
14    crate::cli::{Command, OutputFormat},
15    crate::config::Config,
16    crate::credential_store::AppCredentialStore,
17    crate::diagnostics::RuntimeDiagnostics,
18    crate::herdr::{HerdrReporter, HerdrState},
19    crate::permission::{PermissionMode, SessionWriteLog},
20    crate::permission_classifier_handler::ClassifierApprovalHandler,
21    crate::subagent::{RunState, RunStatus},
22    crate::tools::agent::BackgroundSubagents,
23    rho_providers::providers::build_automation_provider,
24};
25
26use super::{
27    agent_binding::BoundAgent,
28    automation_protocol::{write_event, JsonlAdapter, TerminalReason, WireEvent},
29    headless_run::{self, HeadlessRunDeps, HostInputResponder},
30    policy::AppPolicy,
31    runtime_builder::{
32        build_runtime_with_max_steps, configured_context_window, RuntimeBuildOptions,
33    },
34    sdk_config::SdkBootstrapOptions,
35    tools_prompt::{assemble_tools_and_prompt, ToolsAndPrompt, ToolsAndPromptOptions},
36};
37
38/// Error returned after an automation run has cleaned up and selected a stable exit code.
39#[derive(Debug)]
40pub struct AutomationExit {
41    code: u8,
42    reason: TerminalReason,
43    message: String,
44}
45
46impl AutomationExit {
47    pub(super) fn new(code: u8, reason: TerminalReason, message: impl Into<String>) -> Self {
48        Self {
49            code,
50            reason,
51            message: message.into(),
52        }
53    }
54
55    /// Returns the documented process exit code for this automation result.
56    pub fn exit_code(&self) -> u8 {
57        self.code
58    }
59
60    fn reason(&self) -> TerminalReason {
61        self.reason
62    }
63}
64
65impl fmt::Display for AutomationExit {
66    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
67        formatter.write_str(&self.message)
68    }
69}
70
71impl std::error::Error for AutomationExit {}
72
73/// Error returned after an automation run handles an interrupt and completes cleanup.
74#[derive(Debug)]
75pub struct AutomationInterrupted {
76    signal: ShutdownSignal,
77}
78
79impl AutomationInterrupted {
80    fn new(signal: ShutdownSignal) -> Self {
81        Self { signal }
82    }
83
84    /// Returns the conventional process exit code for the received signal.
85    pub fn exit_code(&self) -> u8 {
86        self.signal.exit_code()
87    }
88}
89
90impl fmt::Display for AutomationInterrupted {
91    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
92        write!(formatter, "rho run interrupted by {}", self.signal)
93    }
94}
95
96impl std::error::Error for AutomationInterrupted {}
97
98#[derive(Clone, Copy, Debug)]
99enum ShutdownSignal {
100    Interrupt,
101    Terminate,
102}
103
104impl ShutdownSignal {
105    fn exit_code(self) -> u8 {
106        match self {
107            Self::Interrupt => 130,
108            Self::Terminate => 143,
109        }
110    }
111}
112
113impl fmt::Display for ShutdownSignal {
114    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
115        match self {
116            Self::Interrupt => formatter.write_str("SIGINT"),
117            Self::Terminate => formatter.write_str("SIGTERM"),
118        }
119    }
120}
121
122#[derive(Debug)]
123struct SubagentCancelled;
124
125impl fmt::Display for SubagentCancelled {
126    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
127        formatter.write_str("subagent cancellation requested")
128    }
129}
130
131impl std::error::Error for SubagentCancelled {}
132
133pub(super) struct Startup<'a> {
134    pub config: &'a Config,
135    pub config_path: PathBuf,
136    pub cwd: PathBuf,
137    pub no_system_prompt: bool,
138    pub no_tools: bool,
139    pub no_subagents: bool,
140    pub usage_purpose: &'static str,
141    pub parent_session_id: Option<rho_sdk::SessionId>,
142    pub agent: BoundAgent,
143    pub output_file: Option<PathBuf>,
144    pub output: OutputFormat,
145    pub max_steps: Option<NonZeroUsize>,
146    pub timeout: Option<Duration>,
147    pub diagnostics: RuntimeDiagnostics,
148    pub herdr: HerdrReporter,
149    pub host_input: Option<Arc<dyn HostInputResponder>>,
150    /// Non-blocking parent notices for background delegated Rho agents.
151    pub notice_poster: Option<Arc<dyn super::subagent_messaging::NoticePoster>>,
152    /// Receives the live steering port once the Rho session starts.
153    pub steering_slot: Option<super::subagent_messaging::SteeringSlot>,
154    pub approval_session: Option<rho_sdk::ApprovalSession>,
155    pub approval_classifier: Option<Arc<ClassifierApprovalHandler>>,
156    pub hook_host_labels: rho_sdk::hooks::HookHostLabels,
157}
158
159pub(super) fn prompt_for_command(command: &Option<Command>) -> anyhow::Result<Option<String>> {
160    match command {
161        Some(Command::Run { prompt, stdin, .. }) => {
162            prompt_from_stdin(prompt.clone(), *stdin).map(Some)
163        }
164        Some(
165            Command::Attach { .. }
166            | Command::Login { .. }
167            | Command::CredentialStore { .. }
168            | Command::Sessions { .. }
169            | Command::Mcp { .. }
170            | Command::Plugins { .. }
171            | Command::Workflow { .. }
172            | Command::WorkflowPlannerWorker
173            | Command::Update,
174        )
175        | None => Ok(None),
176    }
177}
178
179pub(super) fn emit_startup_failure(message: impl Into<String>) -> anyhow::Result<()> {
180    let mut adapter = JsonlAdapter::new();
181    let event = adapter.failed(TerminalReason::ConfigurationError, message.into(), None);
182    emit(event)
183}
184
185pub(super) async fn run(prompt_text: String, startup: Startup<'_>) -> anyhow::Result<()> {
186    let mut jsonl = (startup.output == OutputFormat::Jsonl).then(JsonlAdapter::new);
187    let deadline = startup
188        .timeout
189        .map(|timeout| tokio::time::Instant::now() + timeout);
190    // The reporter exists before anything that can fail, so a parent process
191    // watching the output file always sees a terminal state, including startup failures.
192    let reporter_result = startup
193        .output_file
194        .as_ref()
195        .map(|path| {
196            RunReporter::new(
197                path.clone(),
198                RunArtifactIdentity {
199                    agent_id: startup.agent.id().to_string(),
200                    agent_fingerprint: startup.agent.fingerprint().to_string(),
201                    provider: startup.config.provider.clone(),
202                    model: startup.config.model.clone(),
203                    runtime: crate::agent::AgentRuntime::Rho,
204                },
205                startup.cwd.clone(),
206                &prompt_text,
207                /* stream_output */ startup.output == OutputFormat::Text,
208                None,
209            )
210        })
211        .transpose();
212    let mut reporter = match reporter_result {
213        Ok(reporter) => reporter,
214        Err(error) => {
215            emit_failure(&mut jsonl, TerminalReason::OutputError, &error)?;
216            return Err(
217                AutomationExit::new(1, TerminalReason::OutputError, error.to_string()).into(),
218            );
219        }
220    };
221
222    let cancellation = rho_tools::cancellation::RunCancellation::default();
223    let (result, timed_out) = if let Some(deadline) = deadline {
224        let future = run_session_with_output(
225            prompt_text,
226            &startup,
227            reporter.as_mut(),
228            Some(cancellation.clone()),
229            jsonl.as_mut(),
230        );
231        tokio::pin!(future);
232        tokio::select! {
233            result = &mut future => (result, false),
234            () = tokio::time::sleep_until(deadline) => {
235                cancellation.cancel();
236                (future.await, true)
237            }
238        }
239    } else {
240        (
241            run_session_with_output(
242                prompt_text,
243                &startup,
244                reporter.as_mut(),
245                None,
246                jsonl.as_mut(),
247            )
248            .await,
249            false,
250        )
251    };
252    let terminal = classify_run_terminal(result, timed_out);
253    if let Some(reporter) = reporter.as_mut() {
254        reporter.finish_terminal(&terminal);
255    }
256    emit_and_exit_terminal(terminal, &mut jsonl, reporter.is_some())
257}
258
259fn write_text_answer(answer: &rho_sdk::RunOutcome, has_reporter: bool) -> anyhow::Result<()> {
260    let result = (|| -> io::Result<()> {
261        let mut stdout = io::stdout().lock();
262        if has_reporter {
263            writeln!(stdout, "\n[subagent run complete]")?;
264        } else {
265            writeln!(stdout, "{}", answer.text())?;
266        }
267        stdout.flush()
268    })();
269    result.map_err(|error| {
270        AutomationExit::new(
271            1,
272            TerminalReason::OutputError,
273            format!("could not write output: {error}"),
274        )
275        .into()
276    })
277}
278
279pub(super) fn emit(event: WireEvent) -> anyhow::Result<()> {
280    let mut stdout = io::stdout().lock();
281    write_event(&mut stdout, &event).map_err(|error| {
282        AutomationExit::new(
283            1,
284            TerminalReason::OutputError,
285            format!("could not write JSONL output: {error}"),
286        )
287        .into()
288    })
289}
290
291fn emit_stopped(adapter: &mut Option<JsonlAdapter>, reason: TerminalReason) -> anyhow::Result<()> {
292    if let Some(adapter) = adapter.as_mut() {
293        let text = adapter.partial_text();
294        let event = adapter.stopped(reason, text);
295        emit(event)?;
296    }
297    Ok(())
298}
299
300fn emit_failure(
301    adapter: &mut Option<JsonlAdapter>,
302    reason: TerminalReason,
303    error: &anyhow::Error,
304) -> anyhow::Result<()> {
305    if let Some(adapter) = adapter.as_mut() {
306        let text = adapter.partial_text();
307        let message = terminal_error_message(reason, error);
308        let event = adapter.failed(reason, message, text);
309        emit(event)?;
310    }
311    Ok(())
312}
313
314const MAX_STEPS_MESSAGE: &str = "rho run reached its model-step limit";
315const TIMEOUT_MESSAGE: &str = "rho run timed out";
316
317/// Single classification of a finished automation session before reporter,
318/// JSONL/text emission, and process exit share one decision.
319enum RunTerminal {
320    Completed(rho_sdk::RunOutcome),
321    MaxSteps(rho_sdk::RunOutcome),
322    Timeout,
323    Failed(anyhow::Error),
324}
325
326fn classify_run_terminal(
327    result: anyhow::Result<rho_sdk::RunOutcome>,
328    timed_out: bool,
329) -> RunTerminal {
330    if timed_out {
331        return RunTerminal::Timeout;
332    }
333    match result {
334        Ok(answer) if answer.stop_reason() == rho_sdk::StopReason::MaxSteps => {
335            RunTerminal::MaxSteps(answer)
336        }
337        Ok(answer) => RunTerminal::Completed(answer),
338        Err(error) => RunTerminal::Failed(error),
339    }
340}
341
342fn emit_and_exit_terminal(
343    terminal: RunTerminal,
344    jsonl: &mut Option<JsonlAdapter>,
345    has_reporter: bool,
346) -> anyhow::Result<()> {
347    match terminal {
348        RunTerminal::Timeout => {
349            emit_stopped(jsonl, TerminalReason::Timeout)?;
350            Err(AutomationExit::new(124, TerminalReason::Timeout, TIMEOUT_MESSAGE).into())
351        }
352        RunTerminal::MaxSteps(answer) => {
353            if let Some(adapter) = jsonl.as_mut() {
354                let text = (!answer.text().is_empty()).then(|| answer.text().into());
355                let event = adapter.stopped(TerminalReason::MaxSteps, text);
356                emit(event)?;
357            } else {
358                write_text_answer(&answer, has_reporter)?;
359            }
360            Err(AutomationExit::new(124, TerminalReason::MaxSteps, MAX_STEPS_MESSAGE).into())
361        }
362        RunTerminal::Completed(answer) => {
363            if let Some(adapter) = jsonl.as_mut() {
364                let event = adapter.completed(answer.text().into());
365                emit(event)?;
366            } else {
367                write_text_answer(&answer, has_reporter)?;
368            }
369            Ok(())
370        }
371        RunTerminal::Failed(error) => {
372            let (reason, code) = classify_error(&error);
373            if reason == TerminalReason::Interrupted {
374                emit_stopped(jsonl, reason)?;
375            } else if reason != TerminalReason::OutputError {
376                emit_failure(jsonl, reason, &error)?;
377            }
378            let message = terminal_error_message(reason, &error);
379            if error.is::<AutomationInterrupted>() {
380                return Err(error);
381            }
382            Err(AutomationExit::new(code, reason, message).into())
383        }
384    }
385}
386
387/// Builds the human-readable terminal message for a classified automation error.
388///
389/// Reason codes stay stable machine labels. The message carries actionable detail
390/// except for authentication failures, which stay generic so credentials and
391/// token material never leave the process through stdout, stderr, or JSONL.
392fn terminal_error_message(reason: TerminalReason, error: &anyhow::Error) -> String {
393    match reason {
394        TerminalReason::Authentication => "authentication failed".to_string(),
395        TerminalReason::ProviderError
396        | TerminalReason::ToolHostError
397        | TerminalReason::ConfigurationError
398        | TerminalReason::OutputError
399        | TerminalReason::OtherError
400        | TerminalReason::Interrupted
401        | TerminalReason::MaxSteps
402        | TerminalReason::Timeout
403        | TerminalReason::Completed => error.to_string(),
404    }
405}
406
407fn classify_error(error: &anyhow::Error) -> (TerminalReason, u8) {
408    if let Some(interrupted) = error.downcast_ref::<AutomationInterrupted>() {
409        return (TerminalReason::Interrupted, interrupted.exit_code());
410    }
411    if let Some(exit) = error.downcast_ref::<AutomationExit>() {
412        return (exit.reason(), exit.exit_code());
413    }
414    for cause in error.chain() {
415        if let Some(error) = cause.downcast_ref::<rho_sdk::Error>() {
416            return match error {
417                rho_sdk::Error::Authentication { .. } => (TerminalReason::Authentication, 1),
418                rho_sdk::Error::Provider(provider)
419                    if provider.kind() == rho_sdk::ProviderErrorKind::Authentication =>
420                {
421                    (TerminalReason::Authentication, 1)
422                }
423                rho_sdk::Error::Provider(_) => (TerminalReason::ProviderError, 1),
424                rho_sdk::Error::Tool(_) => (TerminalReason::ToolHostError, 1),
425                rho_sdk::Error::InvalidConfiguration { .. } => {
426                    (TerminalReason::ConfigurationError, 2)
427                }
428                _ => (TerminalReason::OtherError, 1),
429            };
430        }
431        if let Some(error) = cause.downcast_ref::<rho_providers::model::ModelError>() {
432            use rho_providers::model::ModelError;
433            return match error {
434                ModelError::MissingCredentials(_) | ModelError::Credentials(_) => {
435                    (TerminalReason::Authentication, 1)
436                }
437                ModelError::UnsupportedReasoning { .. } | ModelError::UnsupportedProvider(_) => {
438                    (TerminalReason::ConfigurationError, 2)
439                }
440                _ => (TerminalReason::ProviderError, 1),
441            };
442        }
443    }
444    (TerminalReason::OtherError, 1)
445}
446
447pub(crate) async fn run_session(
448    prompt_text: String,
449    startup: &Startup<'_>,
450    reporter: Option<&mut RunReporter>,
451    cancellation: Option<rho_tools::cancellation::RunCancellation>,
452) -> anyhow::Result<rho_sdk::RunOutcome> {
453    ensure_headless_auto_classifier_model(startup.config)?;
454    run_session_with_output(prompt_text, startup, reporter, cancellation, None).await
455}
456
457async fn run_session_with_output(
458    prompt_text: String,
459    startup: &Startup<'_>,
460    reporter: Option<&mut RunReporter>,
461    cancellation: Option<rho_tools::cancellation::RunCancellation>,
462    mut jsonl: Option<&mut JsonlAdapter>,
463) -> anyhow::Result<rho_sdk::RunOutcome> {
464    ensure_headless_auto_classifier_model(startup.config)?;
465    let _scope = startup.config.providers.thread_scope()?;
466    let sdk_options = SdkBootstrapOptions::from_config(startup.config, &startup.cwd)?;
467    let credentials = rho_providers::auth::provider_credentials::ApplicationCredentialSource::new(
468        Arc::new(AppCredentialStore),
469    );
470    let provider = build_automation_provider(sdk_options.provider, &credentials)?;
471    let workspace_root = sdk_options.workspace.root.clone();
472    let workspace = sdk_options.workspace.build_workspace()?;
473    let ToolsAndPrompt {
474        tools: mut tool_set,
475        system_prompt,
476        ..
477    } = assemble_tools_and_prompt(ToolsAndPromptOptions {
478        config: startup.config,
479        config_path: startup.config_path.clone(),
480        cwd: &startup.cwd,
481        no_system_prompt: startup.no_system_prompt,
482        no_tools: startup.no_tools,
483        no_subagents: startup.no_subagents,
484        // Automation keeps questionnaire capability when the agent exposes it.
485        questionnaire_enabled: true,
486        // An automation run can only show a server's question when the caller
487        // supplied a responder for host input; without one the run would fail
488        // on the first question instead of declining it.
489        mcp_elicitation: match startup.host_input {
490            Some(_) => crate::tools::mcp::McpElicitationSupport::Available,
491            None => crate::tools::mcp::McpElicitationSupport::Unavailable,
492        },
493        // Automation binds no model for sampling, so it never declares the
494        // capability and rejects any request that arrives anyway.
495        mcp_sampling: crate::app::tools_prompt::McpSamplingSupport::Unavailable,
496        // Do not block automation startup on models.dev; use whatever cache is warm.
497        await_catalog_names: false,
498        background_subagents: BackgroundSubagents::Disabled,
499        diagnostics: &startup.diagnostics,
500        agent: &startup.agent,
501    })
502    .await?;
503    if let Some(poster) = startup.notice_poster.clone() {
504        tool_set.add_bundle(crate::tools::message_parent_bundle(poster));
505    }
506
507    let context_window = configured_context_window(startup.config);
508    let compaction = sdk_options.runtime.compaction.clone();
509    startup.diagnostics.update_compaction_config(&compaction);
510    let usage_recording = crate::usage::default_recording().await;
511    let session_writes = SessionWriteLog::default();
512    let approval_session = headless_approval_session(
513        startup.config,
514        startup.approval_session.clone(),
515        startup.approval_classifier.clone(),
516        workspace_root.clone(),
517        usage_recording.clone(),
518        session_writes.clone(),
519    )?;
520    let hooks = crate::hooks::start_for_cwd(&workspace_root);
521    if let Some(hooks) = hooks.as_ref() {
522        startup.diagnostics.attach_hooks(hooks);
523    }
524    let startup_result: anyhow::Result<_> = async {
525        let runtime = build_runtime_with_max_steps(
526            RuntimeBuildOptions {
527                provider,
528                tools: tool_set.tools(),
529                workspace,
530                workspace_policy: AppPolicy::for_mode(
531                    startup.config.permission_mode,
532                    session_writes,
533                ),
534                approval_session,
535                system_prompt,
536                reasoning: sdk_options.runtime.reasoning,
537                service_tier: sdk_options.runtime.service_tier,
538                compaction,
539                context_window,
540                usage_purpose: startup.usage_purpose,
541                usage_parent_session_id: startup.parent_session_id.clone(),
542                usage_recording,
543                hook_host_labels: startup.hook_host_labels.clone(),
544                hooks: hooks.as_ref(),
545            },
546            startup.max_steps,
547        )?;
548        let session = match runtime.session(SessionOptions::default()).await {
549            Ok(session) => session,
550            Err(error) => {
551                runtime.shutdown();
552                return Err(error.into());
553            }
554        };
555        anyhow::Ok((runtime, session))
556    }
557    .await;
558    let (runtime, session) = match startup_result {
559        Ok(startup) => startup,
560        Err(error) => {
561            if let Some(hooks) = hooks {
562                hooks.shutdown(crate::hooks::DRAIN_GRACE).await;
563            }
564            tool_set.shutdown().await;
565            return Err(error);
566        }
567    };
568    if let Some(advisor) = tool_set.advisor() {
569        advisor.bind_session(session.clone());
570    }
571    if let Some(adapter) = jsonl.as_deref_mut() {
572        adapter.set_run_context(session.id(), &workspace_root);
573    }
574    startup
575        .herdr
576        .report_state(HerdrState::Working, None, None)
577        .await;
578    let result = complete_run(
579        &session,
580        prompt_text,
581        HeadlessRunDeps {
582            reporter,
583            external_cancellation: cancellation,
584            jsonl,
585            host_input: startup.host_input.as_deref(),
586        },
587        startup.steering_slot.clone(),
588    )
589    .await;
590
591    let session_hooks = runtime.hooks();
592    let session_id = session.id().clone();
593    match &result {
594        Ok(_) => {
595            session_hooks.session_completed(&session_id, /* completed_runs */ 1)
596        }
597        Err(error) => session_hooks.session_failed(
598            &session_id,
599            rho_sdk::hooks::HookSessionFailureKind::RunFailed,
600            &error.to_string(),
601        ),
602    }
603    runtime.shutdown();
604    drop(session);
605    drop(runtime);
606    if let Some(hooks) = hooks {
607        hooks.shutdown(crate::hooks::DRAIN_GRACE).await;
608    }
609    tool_set.shutdown().await;
610    startup
611        .herdr
612        .report_state(HerdrState::Idle, None, None)
613        .await;
614    startup.herdr.release().await;
615
616    result
617}
618
619pub(crate) fn ensure_headless_auto_classifier_model(config: &Config) -> anyhow::Result<()> {
620    if config.permission_mode == PermissionMode::Auto
621        && config
622            .internal_agent_model(PERMISSION_CLASSIFIER_AGENT_ID)
623            .is_none()
624    {
625        anyhow::bail!(
626            "permission mode auto requires a configured permission-classifier model (set via /config or config.toml [internal_agents.permission-classifier])"
627        );
628    }
629    Ok(())
630}
631
632/// Resolves the approval session for one headless run.
633///
634/// Non-Auto keeps the inherited session. Auto always installs a classifier:
635/// isolate a workflow/subagent template onto this run's write log when present,
636/// otherwise build a fresh headless classifier. A stray non-classifier
637/// `approval_session` is ignored in Auto so callers do not juggle paired Option
638/// knobs.
639fn headless_approval_session(
640    config: &Config,
641    approval_session: Option<rho_sdk::ApprovalSession>,
642    approval_classifier: Option<Arc<ClassifierApprovalHandler>>,
643    workspace_root: PathBuf,
644    usage_recording: rho_sdk::ProviderRequestUsageRecording,
645    session_writes: SessionWriteLog,
646) -> anyhow::Result<Option<rho_sdk::ApprovalSession>> {
647    if config.permission_mode != PermissionMode::Auto {
648        return Ok(approval_session);
649    }
650    Ok(Some(rho_sdk::ApprovalSession::from_shared(
651        headless_auto_classifier(
652            config,
653            approval_classifier,
654            workspace_root,
655            usage_recording,
656            session_writes,
657        ),
658    )))
659}
660
661fn headless_auto_classifier(
662    config: &Config,
663    approval_classifier: Option<Arc<ClassifierApprovalHandler>>,
664    workspace_root: PathBuf,
665    usage_recording: rho_sdk::ProviderRequestUsageRecording,
666    session_writes: SessionWriteLog,
667) -> Arc<ClassifierApprovalHandler> {
668    match approval_classifier {
669        Some(template) => template.isolate_for_run(session_writes),
670        None => ClassifierApprovalHandler::shared(
671            config.clone(),
672            workspace_root,
673            usage_recording,
674            None,
675            Some(session_writes),
676        ),
677    }
678}
679
680async fn complete_run(
681    session: &rho_sdk::Session,
682    prompt_text: String,
683    dependencies: HeadlessRunDeps<'_>,
684    steering_slot: Option<super::subagent_messaging::SteeringSlot>,
685) -> anyhow::Result<rho_sdk::RunOutcome> {
686    let HeadlessRunDeps {
687        reporter,
688        external_cancellation,
689        jsonl,
690        host_input,
691    } = dependencies;
692    let mut run = session.start(UserInput::text(prompt_text)).await?;
693    if let Some(slot) = steering_slot {
694        slot.publish(run.steering_handle());
695    }
696    let cancellation = run.cancellation_handle();
697    let external_cancellation = external_cancellation.unwrap_or_default();
698    tokio::select! {
699        outcome = headless_run::drive(&mut run, reporter, jsonl, host_input) => outcome,
700        signal = shutdown_signal() => {
701            let signal = signal?;
702            cancellation.cancel();
703            let _ = run.outcome().await;
704            Err(AutomationInterrupted::new(signal).into())
705        }
706        () = external_cancellation.cancelled() => {
707            cancellation.cancel();
708            let _ = run.outcome().await;
709            Err(SubagentCancelled.into())
710        }
711    }
712}
713
714pub(crate) use crate::run_artifacts::RunArtifactIdentity;
715
716/// Maintains the `--output-file` status contract for subagent runs and
717/// streams progress to stdout so a watching pane shows live activity.
718pub(crate) struct RunReporter {
719    sink: crate::run_artifacts::RunArtifactSink,
720    adapter: crate::tui::event_adapter::SdkEventAdapter,
721    stream_output: bool,
722}
723
724impl RunReporter {
725    pub(crate) fn new(
726        path: PathBuf,
727        identity: RunArtifactIdentity,
728        cwd: PathBuf,
729        prompt: &str,
730        stream_output: bool,
731        status_tx: Option<tokio::sync::watch::Sender<RunStatus>>,
732    ) -> anyhow::Result<Self> {
733        let sink = crate::run_artifacts::RunArtifactSink::open(path, &identity, prompt, status_tx)?;
734        Ok(Self {
735            sink,
736            adapter: crate::tui::event_adapter::SdkEventAdapter::new(cwd),
737            stream_output,
738        })
739    }
740
741    /// Resume after the executor already wrote the Starting boundary.
742    pub(crate) fn continue_from(
743        path: PathBuf,
744        started_status: RunStatus,
745        cwd: PathBuf,
746        prompt: &str,
747        stream_output: bool,
748        status_tx: Option<tokio::sync::watch::Sender<RunStatus>>,
749    ) -> anyhow::Result<Self> {
750        let sink = crate::run_artifacts::RunArtifactSink::continue_from(
751            path,
752            started_status,
753            prompt,
754            status_tx,
755        )?;
756        Ok(Self {
757            sink,
758            adapter: crate::tui::event_adapter::SdkEventAdapter::new(cwd),
759            stream_output,
760        })
761    }
762
763    pub(super) fn on_event(&mut self, event: &rho_sdk::RunEvent) {
764        use rho_sdk::RunEvent;
765
766        let attachments = crate::tui::translate_run_event(&mut self.adapter, event);
767        for attachment in attachments {
768            // Reasoning is deliberately kept out of `last_text`: the status file
769            // carries the answer, not the thinking.
770            if let crate::run_artifacts::AttachmentEvent::AssistantTextDelta(text) = &attachment {
771                if !text.is_empty() {
772                    self.sink.append_last_text(text);
773                }
774            }
775            self.sink.record_attachment(attachment);
776        }
777        match event {
778            RunEvent::StepStarted { step, .. } => {
779                self.sink.status.state = RunState::Running;
780                self.sink.status.turns = *step as u64;
781                self.sink.publish();
782            }
783            RunEvent::ToolStarted { name, .. } => {
784                self.sink.status.last_activity = Some(format!("tool: {name}"));
785                self.stream(&format!("\n[tool] {name}\n"));
786                self.sink.publish();
787            }
788            RunEvent::HostInputRequested { request }
789            | RunEvent::ToolHostInputRequested { request, .. } => {
790                self.sink.status.last_activity =
791                    Some(format!("waiting for questionnaire: {}", request.title()));
792                self.sink.publish();
793            }
794            RunEvent::AssistantTextDelta { text } => {
795                self.sink.status.last_activity = Some("assistant text".into());
796                self.stream(text);
797                // Attachment path already published throttled when translated.
798            }
799            RunEvent::ProviderStreamReset { .. } => {
800                self.sink.status.last_activity = Some("retrying provider response".into());
801                self.sink.status.last_text = None;
802                self.stream("\n[provider response discarded; retrying]\n");
803                self.sink.publish();
804            }
805            RunEvent::UsageUpdated { usage } => {
806                self.sink.status.input_tokens = usage.total_input_tokens();
807                self.sink.status.output_tokens = usage.output_tokens;
808            }
809            _ => {}
810        }
811    }
812
813    #[cfg(test)]
814    pub(crate) fn status(&self) -> &RunStatus {
815        &self.sink.status
816    }
817
818    pub(super) fn write(&mut self) {
819        self.sink.publish();
820    }
821
822    pub(crate) fn finish(&mut self, result: &anyhow::Result<rho_sdk::RunOutcome>) {
823        match result {
824            Ok(outcome) => {
825                let usage = outcome.usage();
826                self.sink.status.input_tokens = usage.total_input_tokens();
827                self.sink.status.output_tokens = usage.output_tokens;
828                self.sink.finish_ok(Some(outcome.text().to_string()));
829            }
830            Err(error)
831                if error.is::<AutomationInterrupted>()
832                    || error.downcast_ref::<AutomationExit>().is_some_and(|exit| {
833                        matches!(
834                            exit.reason(),
835                            TerminalReason::MaxSteps | TerminalReason::Timeout
836                        )
837                    })
838                    || error.is::<SubagentCancelled>() =>
839            {
840                self.sink.finish_stopped("stopped");
841            }
842            Err(error) => {
843                self.sink.finish_error(format!("{error:#}"));
844            }
845        }
846    }
847
848    fn finish_terminal(&mut self, terminal: &RunTerminal) {
849        match terminal {
850            RunTerminal::Completed(outcome) => {
851                let usage = outcome.usage();
852                self.sink.status.input_tokens = usage.total_input_tokens();
853                self.sink.status.output_tokens = usage.output_tokens;
854                self.sink.finish_ok(Some(outcome.text().to_string()));
855            }
856            RunTerminal::MaxSteps(_) | RunTerminal::Timeout => {
857                self.sink.finish_stopped("stopped");
858            }
859            RunTerminal::Failed(error)
860                if error.is::<AutomationInterrupted>() || error.is::<SubagentCancelled>() =>
861            {
862                self.sink.finish_stopped("stopped");
863            }
864            RunTerminal::Failed(error) => {
865                self.sink.finish_error(format!("{error:#}"));
866            }
867        }
868    }
869
870    fn stream(&self, text: &str) {
871        if !self.stream_output {
872            return;
873        }
874        let mut stdout = io::stdout().lock();
875        let _ = stdout.write_all(text.as_bytes());
876        let _ = stdout.flush();
877    }
878}
879
880#[cfg(unix)]
881async fn shutdown_signal() -> io::Result<ShutdownSignal> {
882    use tokio::signal::unix::{signal, SignalKind};
883
884    let mut interrupt = signal(SignalKind::interrupt())?;
885    let mut terminate = signal(SignalKind::terminate())?;
886    tokio::select! {
887        _ = interrupt.recv() => Ok(ShutdownSignal::Interrupt),
888        _ = terminate.recv() => Ok(ShutdownSignal::Terminate),
889    }
890}
891
892#[cfg(not(unix))]
893async fn shutdown_signal() -> io::Result<ShutdownSignal> {
894    tokio::signal::ctrl_c().await?;
895    Ok(ShutdownSignal::Interrupt)
896}
897
898fn prompt_from_stdin(parts: Vec<String>, read_stdin: bool) -> anyhow::Result<String> {
899    if !read_stdin && crate::stdio::stdin_is_redirected() {
900        anyhow::bail!(
901            "stdin is redirected but --stdin was not set; pass --stdin to include piped input"
902        );
903    }
904    prompt_from_reader(parts, read_stdin, &mut io::stdin())
905}
906
907fn prompt_from_reader(
908    parts: Vec<String>,
909    read_stdin: bool,
910    stdin: &mut impl Read,
911) -> anyhow::Result<String> {
912    let mut chunks = Vec::new();
913    let inline = parts.join(" ").trim().to_string();
914    if !inline.is_empty() {
915        chunks.push(inline);
916    }
917    if read_stdin {
918        let mut buffer = String::new();
919        stdin.read_to_string(&mut buffer)?;
920        let buffer = buffer.trim().to_string();
921        if !buffer.is_empty() {
922            chunks.push(buffer);
923        }
924    }
925
926    let prompt = chunks.join("\n\n");
927    if prompt.is_empty() {
928        anyhow::bail!("rho run requires a prompt argument or --stdin");
929    }
930    Ok(prompt)
931}
932
933#[cfg(test)]
934#[path = "automation_tests.rs"]
935mod tests;