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#[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 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#[derive(Debug)]
75pub struct AutomationInterrupted {
76 signal: ShutdownSignal,
77}
78
79impl AutomationInterrupted {
80 fn new(signal: ShutdownSignal) -> Self {
81 Self { signal }
82 }
83
84 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 pub notice_poster: Option<Arc<dyn super::subagent_messaging::NoticePoster>>,
152 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 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 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
317enum 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
387fn 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 questionnaire_enabled: true,
486 mcp_elicitation: match startup.host_input {
490 Some(_) => crate::tools::mcp::McpElicitationSupport::Available,
491 None => crate::tools::mcp::McpElicitationSupport::Unavailable,
492 },
493 mcp_sampling: crate::app::tools_prompt::McpSamplingSupport::Unavailable,
496 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, 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
632fn 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
716pub(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 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 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 }
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;