use std::io::{BufRead, IsTerminal, Write};
use std::sync::{Arc, Mutex};
use rpi_agent::events::AgentEvent;
use rpi_ai::types::{AssistantMessage, Content, ImageContent, StopReason};
use rpi_harness::agent_harness::{AgentHarness, AgentLane, HarnessRunOutcome};
use rpi_harness::events::{HarnessEvent, RunEndOutcome};
use crate::args::Args;
pub fn assistant_text(msg: &AssistantMessage) -> String {
msg.content
.iter()
.filter_map(|c| match c {
Content::Text(t) => Some(t.text.clone()),
_ => None,
})
.collect()
}
pub fn outcome_exit_code(outcome: &HarnessRunOutcome) -> i32 {
match outcome {
HarnessRunOutcome::Failed { .. } | HarnessRunOutcome::Aborted { .. } => 1,
_ => 0,
}
}
pub async fn print(
harness: &AgentHarness,
_args: &Args,
initial: Option<String>,
extra_messages: &[String],
initial_images: Vec<ImageContent>,
) -> i32 {
let lane: Arc<dyn AgentLane> = harness.lane("main");
let mut last_exit = 0;
let mut last_msg: Option<AssistantMessage> = None;
let mut prompts: Vec<String> = Vec::new();
if let Some(init) = initial {
prompts.push(init);
}
for m in extra_messages {
prompts.push(m.clone());
}
if prompts.is_empty() {
return 0;
}
let mut images = initial_images;
for prompt in prompts {
match lane.prompt_text(&prompt, std::mem::take(&mut images)).await {
Ok(result) => {
last_exit = outcome_exit_code(&result.outcome);
match &result.outcome {
HarnessRunOutcome::Completed { final_message, .. }
| HarnessRunOutcome::Aborted { final_message, .. } => {
last_msg = Some(final_message.clone());
}
HarnessRunOutcome::Failed {
error,
final_message,
..
} => {
if let Some(m) = final_message {
if m.stop_reason == StopReason::Error {
if let Some(em) = &m.error_message {
eprintln!("{em}");
}
}
}
eprintln!("run failed: {error:?}");
}
HarnessRunOutcome::Suspended { .. } => {
eprintln!("run suspended (deferred) — resume is not supported in v1");
last_exit = 1;
}
}
}
Err(e) => {
eprintln!("prompt rejected: {e}");
return 1;
}
}
}
if let Some(m) = &last_msg {
match m.stop_reason {
StopReason::Error => {
if let Some(em) = &m.error_message {
eprintln!("{em}");
}
last_exit = 1;
}
StopReason::Aborted => {
eprintln!("request aborted");
last_exit = 1;
}
_ => {
let text = assistant_text(m);
let mut out = std::io::stdout();
let _ = out.write_all(text.as_bytes());
if !text.ends_with('\n') {
let _ = out.write_all(b"\n");
}
let _ = out.flush();
}
}
}
last_exit
}
pub async fn json(
harness: &AgentHarness,
_args: &Args,
initial: Option<String>,
extra_messages: &[String],
initial_images: Vec<ImageContent>,
mut agent_events: Option<tokio::sync::broadcast::Receiver<AgentEvent>>,
) -> i32 {
let lane: Arc<dyn AgentLane> = harness.lane("main");
let collected: Arc<Mutex<Vec<HarnessEvent>>> = Arc::new(Mutex::new(Vec::new()));
let collected_for_watch = collected.clone();
let mut watch = harness.events().watch(|| ());
watch.start(Arc::new(move |event: &HarnessEvent| {
emit_json_event(event);
collected_for_watch.lock().unwrap().push(event.clone());
}));
std::mem::forget(watch);
let (agent_done_tx, mut agent_done_rx) = tokio::sync::mpsc::unbounded_channel::<()>();
let agent_event_task = agent_events.take().map(|mut rx| {
tokio::spawn(async move {
let done_tx = agent_done_tx;
loop {
match rx.recv().await {
Ok(event) => {
let terminal = event.is_terminal();
emit_agent_event(&event);
if terminal {
let _ = done_tx.send(());
}
}
Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => continue,
Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
}
}
})
});
let mut prompts: Vec<String> = Vec::new();
if let Some(init) = initial {
prompts.push(init);
}
for m in extra_messages {
prompts.push(m.clone());
}
let mut last_exit = 0;
let mut final_outcome: Option<HarnessRunOutcome> = None;
let mut images = initial_images;
for prompt in prompts {
match lane.prompt_text(&prompt, std::mem::take(&mut images)).await {
Ok(result) => {
let _ = tokio::time::timeout(
std::time::Duration::from_millis(250),
agent_done_rx.recv(),
)
.await;
last_exit = outcome_exit_code(&result.outcome);
final_outcome = Some(result.outcome);
}
Err(e) => {
let line = serde_json::json!({
"type": "error",
"error": e.to_string(),
});
println!("{line}");
if let Some(task) = agent_event_task {
task.abort();
}
return 1;
}
}
}
let (outcome_str, final_text) = match final_outcome {
Some(HarnessRunOutcome::Completed { final_message, .. }) => {
("completed", Some(assistant_text(&final_message)))
}
Some(HarnessRunOutcome::Aborted { final_message, .. }) => {
("aborted", Some(assistant_text(&final_message)))
}
Some(HarnessRunOutcome::Failed { final_message, .. }) => {
let t = final_message.as_ref().map(assistant_text);
("failed", t)
}
Some(HarnessRunOutcome::Suspended { .. }) => ("suspended", None),
None => ("idle", None),
};
let result_line = serde_json::json!({
"type": "result",
"outcome": outcome_str,
"finalText": final_text,
});
println!("{result_line}");
if let Some(task) = agent_event_task {
task.abort();
}
last_exit
}
fn emit_agent_event(event: &AgentEvent) {
println!("{}", agent_event_json(event));
}
fn agent_event_json(event: &AgentEvent) -> serde_json::Value {
use rpi_ai::types::AssistantMessageEvent;
match event {
AgentEvent::AgentStart => serde_json::json!({"type":"agent_start"}),
AgentEvent::AgentEnd { messages } => serde_json::json!({
"type":"agent_end", "messageCount": messages.len(),
"messages": serde_json::to_value(messages).unwrap_or(serde_json::Value::Null)
}),
AgentEvent::RetryScheduled {
attempt,
max_retries,
delay_ms,
error,
} => serde_json::json!({
"type":"retry_scheduled", "attempt":attempt,
"maxRetries":max_retries, "delayMs":delay_ms, "error":error
}),
AgentEvent::TurnStart => serde_json::json!({"type":"turn_start"}),
AgentEvent::TurnEnd {
message,
tool_results,
} => serde_json::json!({
"type":"turn_end", "message": serde_json::to_value(message).ok(),
"toolResultCount": tool_results.len(),
"toolResults": serde_json::to_value(tool_results).unwrap_or(serde_json::Value::Null)
}),
AgentEvent::MessageStart { message } => serde_json::json!({
"type":"message_start", "message": serde_json::to_value(message).ok()
}),
AgentEvent::MessageEnd { message } => serde_json::json!({
"type":"message_end", "message": serde_json::to_value(message).ok()
}),
AgentEvent::MessageUpdate {
message,
assistant_message_event,
} => {
let mut value = serde_json::json!({
"type":"message_update",
"message": serde_json::to_value(message).unwrap_or(serde_json::Value::Null),
"assistantMessageEvent": serde_json::to_value(assistant_message_event)
.unwrap_or(serde_json::Value::Null),
"eventType": assistant_message_event.type_tag(),
});
let object = value.as_object_mut().expect("json object");
match assistant_message_event {
AssistantMessageEvent::TextDelta {
content_index,
delta,
..
}
| AssistantMessageEvent::ThinkingDelta {
content_index,
delta,
..
}
| AssistantMessageEvent::ToolCallDelta {
content_index,
delta,
..
} => {
object.insert("contentIndex".into(), (*content_index).into());
object.insert("delta".into(), delta.clone().into());
}
_ => {}
}
value
}
AgentEvent::ToolExecutionStart {
tool_call_id,
tool_name,
args,
} => serde_json::json!({
"type":"tool_execution_start", "toolCallId":tool_call_id,
"toolName":tool_name, "args":args
}),
AgentEvent::ToolExecutionUpdate {
tool_call_id,
tool_name,
args,
partial_result,
} => serde_json::json!({
"type":"tool_execution_update", "toolCallId":tool_call_id, "toolName":tool_name,
"args": args,
"partialResult": tool_result_json(partial_result)
}),
AgentEvent::ToolExecutionEnd {
tool_call_id,
tool_name,
result,
is_error,
} => serde_json::json!({
"type":"tool_execution_end", "toolCallId":tool_call_id,
"toolName":tool_name, "isError":is_error,
"result": tool_result_json(result)
}),
}
}
fn tool_result_json(result: &rpi_agent::types::AgentToolResult) -> serde_json::Value {
let content: Vec<serde_json::Value> = result
.content
.iter()
.map(|item| match item {
rpi_agent::types::TextContentOrImage::Text(text) => serde_json::json!({
"type": "text",
"text": text.text,
}),
rpi_agent::types::TextContentOrImage::Image(image) => {
serde_json::to_value(image).unwrap_or(serde_json::Value::Null)
}
})
.collect();
serde_json::json!({
"content": content,
"details": result.details,
"usage": result.usage.as_ref().and_then(|usage| serde_json::to_value(usage).ok()),
"addedToolNames": result.added_tool_names,
"terminate": result.terminate,
})
}
fn emit_json_event(event: &HarnessEvent) {
let line = match event {
HarnessEvent::RunStart(e) => serde_json::json!({
"type": "run_start",
"lane": e.lane,
"runId": e.run_id,
}),
HarnessEvent::RunEnd(e) => serde_json::json!({
"type": "run_end",
"lane": e.lane,
"runId": e.run_id,
"outcome": run_end_outcome_str(e.outcome),
"leafId": e.leaf_id,
}),
};
println!("{line}");
}
fn run_end_outcome_str(o: RunEndOutcome) -> &'static str {
match o {
RunEndOutcome::Completed => "completed",
RunEndOutcome::Aborted => "aborted",
RunEndOutcome::Failed => "failed",
}
}
pub async fn interactive(
harness: &AgentHarness,
event_rx: Option<tokio::sync::broadcast::Receiver<rpi_agent::AgentEvent>>,
args: &Args,
model_catalog: Vec<rpi_ai::Model>,
initial: Option<String>,
extra_messages: &[String],
initial_images: Vec<ImageContent>,
theme: Option<&str>,
no_themes: bool,
reload_context: &crate::session::ReloadContext,
) -> i32 {
let force_tui = std::env::var("RPI_FORCE_TUI")
.map(|v| v == "1")
.unwrap_or(false);
if force_tui || crate::interactive_tui::is_tui_supported() {
crate::interactive_tui::interactive_tui(
harness,
event_rx,
args,
model_catalog,
initial,
extra_messages,
initial_images,
theme,
no_themes,
reload_context,
)
.await
} else {
interactive_repl(harness, args, initial, extra_messages, initial_images).await
}
}
pub async fn interactive_repl(
harness: &AgentHarness,
#[allow(unused_variables)] args: &Args,
initial: Option<String>,
extra_messages: &[String],
initial_images: Vec<ImageContent>,
) -> i32 {
let lane: Arc<dyn AgentLane> = harness.lane("main");
let stdin = std::io::stdin();
let is_tty = stdin.is_terminal();
if is_tty {
println!(
"rpi interactive (v1 minimal REPL). Type /exit to quit, /abort to cancel a run.\n"
);
}
let mut prompts: Vec<String> = Vec::new();
if let Some(init) = initial {
prompts.push(init);
}
for m in extra_messages {
prompts.push(m.clone());
}
let mut images = initial_images;
for prompt in prompts {
if let Err(code) = run_one(&lane, &prompt, std::mem::take(&mut images)).await {
return code;
}
}
let mut line = String::new();
loop {
if is_tty {
print!("> ");
let _ = std::io::stdout().flush();
}
line.clear();
match stdin.lock().read_line(&mut line) {
Ok(0) => break, Ok(_) => {}
Err(_) => break,
}
let trimmed = line.trim();
if trimmed.is_empty() {
continue;
}
if trimmed == "/exit" || trimmed == "/quit" {
break;
}
if trimmed == "/abort" {
let _ = lane.abort().await;
eprintln!("(aborted)");
continue;
}
if let Err(code) = run_one(&lane, trimmed, Vec::new()).await {
return code;
}
}
0
}
async fn run_one(
lane: &Arc<dyn AgentLane>,
prompt: &str,
images: Vec<ImageContent>,
) -> Result<(), i32> {
match lane.prompt_text(prompt, images).await {
Ok(result) => {
match &result.outcome {
HarnessRunOutcome::Completed { final_message, .. }
| HarnessRunOutcome::Aborted { final_message, .. } => {
let text = assistant_text(final_message);
if !text.is_empty() {
println!("{text}");
}
}
HarnessRunOutcome::Failed {
error,
final_message,
..
} => {
if let Some(m) = final_message {
if let Some(em) = &m.error_message {
eprintln!("error: {em}");
}
}
eprintln!("run failed: {error:?}");
}
HarnessRunOutcome::Suspended { .. } => {
eprintln!("run suspended (deferred) — resume not supported in v1");
}
}
Ok(())
}
Err(e) => {
eprintln!("prompt rejected: {e}");
Err(1)
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use rpi_agent::events::AgentEvent;
use rpi_agent::types::AgentToolResult;
use rpi_ai::types::{
AssistantMessage, Content, StopReason, TextContent, TextContentType, Usage,
};
use rpi_harness::session::types::OperationError;
fn assistant(text: &str, stop: StopReason) -> AssistantMessage {
AssistantMessage {
role: rpi_ai::types::AssistantRole,
content: vec![Content::Text(TextContent {
kind: TextContentType,
text: text.into(),
text_signature: None,
})],
api: rpi_ai::Api::AnthropicMessages,
provider: "anthropic".into(),
model: "claude-sonnet-5".into(),
response_model: None,
response_id: None,
usage: Usage::zero(),
stop_reason: stop,
deferred: None,
error_message: None,
raw_stop_reason: None,
end_turn: None,
timestamp: 0,
}
}
#[test]
fn assistant_text_concatenates_text_blocks() {
let m = assistant("hello", StopReason::Stop);
assert_eq!(assistant_text(&m), "hello");
}
#[test]
fn outcome_exit_code_maps_failed_aborted_to_1() {
let failed = HarnessRunOutcome::Failed {
leaf_id: "l".into(),
error: OperationError {
code: "boom".into(),
message: "boom".into(),
},
final_entry_id: None,
final_message: None,
};
assert_eq!(outcome_exit_code(&failed), 1);
let completed = HarnessRunOutcome::Completed {
leaf_id: "l".into(),
final_entry_id: "e".into(),
final_message: assistant("ok", StopReason::Stop),
};
assert_eq!(outcome_exit_code(&completed), 0);
}
#[test]
fn run_end_outcome_str_roundtrip() {
assert_eq!(run_end_outcome_str(RunEndOutcome::Completed), "completed");
assert_eq!(run_end_outcome_str(RunEndOutcome::Aborted), "aborted");
assert_eq!(run_end_outcome_str(RunEndOutcome::Failed), "failed");
}
#[test]
fn agent_event_projection_keeps_terminal_and_tool_payloads() {
let end = agent_event_json(&AgentEvent::AgentEnd { messages: vec![] });
assert_eq!(end["type"], "agent_end");
assert_eq!(end["messages"], serde_json::json!([]));
let tool = agent_event_json(&AgentEvent::ToolExecutionEnd {
tool_call_id: "call-1".into(),
tool_name: "read".into(),
result: AgentToolResult::text("hello"),
is_error: false,
});
assert_eq!(tool["type"], "tool_execution_end");
assert_eq!(tool["result"]["content"][0]["text"], "hello");
assert_eq!(tool["result"]["terminate"], false);
let retry = agent_event_json(&AgentEvent::RetryScheduled {
attempt: 3,
max_retries: 10,
delay_ms: 8_000,
error: "503 service unavailable".into(),
});
assert_eq!(retry["type"], "retry_scheduled");
assert_eq!(retry["attempt"], 3);
assert_eq!(retry["maxRetries"], 10);
assert_eq!(retry["delayMs"], 8_000);
assert_eq!(retry["error"], "503 service unavailable");
}
}