Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
154 changes: 77 additions & 77 deletions Cargo.lock

Large diffs are not rendered by default.

2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -87,7 +87,7 @@ members = [
exclude = ["services/rc-server"]

[workspace.package]
version = "1.8.6"
version = "1.8.7"
edition = "2021"
license = "Apache-2.0"
rust-version = "1.80"
Expand Down
188 changes: 181 additions & 7 deletions crates/rebon-cli/src/exec.rs
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,10 @@ use rebon_types::{AgentCapabilityMode, ContentBlock as AcpContentBlock, TextCont

use crate::session::commands::effort::{resolve_thinking_from_effort, ThinkingOverrides};

use std::path::PathBuf;
use std::io::Read;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::time::Duration;

/// Parsed inputs for `rebon exec`, assembled from the subcommand args plus the
Expand Down Expand Up @@ -84,6 +87,52 @@ pub struct ExecArgs {
pub max_duration_sec: Option<u64>,
}

pub fn prepare_prompt(
words: &[String],
prompt_file: Option<&Path>,
json: bool,
) -> anyhow::Result<String> {
resolve_prompt(words, prompt_file, std::io::stdin().lock()).inspect_err(|error| {
if json {
emit_error_once(&AtomicBool::new(false), &error.to_string());
}
})
}

fn resolve_prompt(
words: &[String],
prompt_file: Option<&Path>,
mut stdin: impl Read,
) -> anyhow::Result<String> {
let prompt = match prompt_file {
Some(path) => {
let source = if path == Path::new("-") {
"stdin".to_string()
} else {
format!("prompt file {}", path.display())
};
let mut bytes = Vec::new();
if path == Path::new("-") {
stdin
.read_to_end(&mut bytes)
.with_context(|| format!("rebon exec: failed to read {source}"))?;
} else {
bytes = std::fs::read(path)
.with_context(|| format!("rebon exec: failed to read {source}"))?;
}
let bytes = bytes.strip_prefix(&[0xef, 0xbb, 0xbf]).unwrap_or(&bytes);
std::str::from_utf8(bytes)
.with_context(|| format!("rebon exec: {source} is not valid UTF-8"))?
.to_owned()
}
None => words.join(" "),
};
if prompt.trim().is_empty() {
anyhow::bail!("rebon exec: empty prompt");
}
Ok(prompt)
}

fn resolve_exec_thinking(
effort: Option<ReasoningEffort>,
provider_format: ProviderFormat,
Expand Down Expand Up @@ -118,10 +167,29 @@ fn select_unattended_allow_option(options: &[PermissionQueryOption]) -> Option<S

/// Run one headless turn and stream its events to stdout.
pub async fn run(args: ExecArgs) -> anyhow::Result<()> {
if args.prompt.trim().is_empty() {
anyhow::bail!("rebon exec: empty prompt");
let json = args.json;
let error_emitted = Arc::new(AtomicBool::new(false));
let outcome = run_inner(args, &error_emitted).await;
if json {
if let Err(error) = &outcome {
emit_error_once(&error_emitted, &error.to_string());
}
}
outcome
}

fn emit_error_once(error_emitted: &AtomicBool, message: &str) {
if !error_emitted.swap(true, Ordering::Relaxed) {
let message = if message.trim().is_empty() {
"rebon exec failed"
} else {
message
};
emit(&serde_json::json!({ "type": "error", "message": message }));
}
}

async fn run_inner(args: ExecArgs, error_emitted: &Arc<AtomicBool>) -> anyhow::Result<()> {
// Say out loud what this process is, before anything builds a tool list
// from it. `exec` has no one to approve a plan, answer a question, or
// retry a call the classifier stopped, and the tools that need one of
Expand Down Expand Up @@ -161,10 +229,20 @@ pub async fn run(args: ExecArgs) -> anyhow::Result<()> {
// Observe every raw QueryEvent and project it onto stdout. The callback runs
// synchronously inside the executor's single consume loop, so lines are
// emitted in event order without interleaving.
let observer_error_emitted = Arc::clone(error_emitted);
let observer = QueryEventObserver::new(move |event| {
if json {
for value in project_json(event) {
emit(&value);
if value["type"] == "error" {
emit_error_once(
&observer_error_emitted,
value["message"]
.as_str()
.expect("projected error has a message"),
);
} else {
emit(&value);
}
}
} else {
project_text(event);
Expand Down Expand Up @@ -379,6 +457,15 @@ pub async fn run(args: ExecArgs) -> anyhow::Result<()> {
_server_state.close_session(&session_id);
rebon_core::system_prompt::remove_scratchpad_for(&cwd, &session_id);

report_outcome(json, &session_id, deadline_reached, result)
}

fn report_outcome(
json: bool,
session_id: &str,
deadline_reached: bool,
result: Result<rebon_agent_core::PromptOutcome, rebon_agent_core::PromptExecutorError>,
) -> anyhow::Result<()> {
match result {
Ok(outcome) => {
// A caller cannot tell "the harness stopped this" from "the user
Expand Down Expand Up @@ -426,9 +513,6 @@ pub async fn run(args: ExecArgs) -> anyhow::Result<()> {
}
Err(err) => {
let message = format!("{err:?}");
if json {
emit(&serde_json::json!({ "type": "error", "message": message }));
}
Err(anyhow::anyhow!("rebon exec turn failed: {message}"))
}
}
Expand Down Expand Up @@ -600,6 +684,96 @@ mod tests {
use super::*;
use serde_json::json;

#[test]
fn prompt_input_preserves_positional_words() {
let words = vec![" 检查\n项目".into(), "'quoted' --flag ".into()];
assert_eq!(
resolve_prompt(&words, None, std::io::empty()).unwrap(),
" 检查\n项目 'quoted' --flag "
);
}

#[test]
fn prompt_input_reads_utf8_files_and_bom_without_trimming() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("提示词.txt");
let content = " 第一行\r\n第二行 café\n\"quoted\"\t\n";
for prefix in ["", "\u{feff}"] {
std::fs::write(&path, format!("{prefix}{content}")).unwrap();
assert_eq!(
resolve_prompt(&[], Some(&path), std::io::empty()).unwrap(),
content
);
}
}

#[test]
fn prompt_input_reads_stdin_and_bom_without_trimming() {
let content = " 多行提示\n非 ASCII café\r\n";
for prefix in ["", "\u{feff}"] {
let input = format!("{prefix}{content}");
assert_eq!(
resolve_prompt(&[], Some(Path::new("-")), input.as_bytes()).unwrap(),
content
);
}
}

#[test]
fn prompt_input_rejects_empty_and_whitespace_from_all_sources() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("empty.txt");
assert!(resolve_prompt(&[], None, std::io::empty()).is_err());
for input in ["", " \t\r\n\u{2003}"] {
assert!(resolve_prompt(&[input.into()], None, std::io::empty()).is_err());
for prefix in ["", "\u{feff}"] {
let input = format!("{prefix}{input}");
std::fs::write(&path, &input).unwrap();
for error in [
resolve_prompt(&[], Some(&path), std::io::empty()).unwrap_err(),
resolve_prompt(&[], Some(Path::new("-")), input.as_bytes()).unwrap_err(),
] {
assert!(error.to_string().contains("empty prompt"));
}
}
}
}

#[test]
fn prompt_input_rejects_missing_unreadable_and_invalid_utf8_files() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("input.txt");
for unreadable in [path.as_path(), dir.path()] {
let error = resolve_prompt(&[], Some(unreadable), std::io::empty()).unwrap_err();
assert!(error.to_string().contains("failed to read prompt file"));
assert!(error
.to_string()
.contains(&unreadable.display().to_string()));
}
for bytes in [vec![0xff], vec![0xef, 0xbb, 0xbf, 0xc3, 0x28]] {
std::fs::write(&path, &bytes).unwrap();
let error = resolve_prompt(&[], Some(&path), std::io::empty()).unwrap_err();
assert!(error.to_string().contains("not valid UTF-8"));
let error = resolve_prompt(&[], Some(Path::new("-")), bytes.as_slice()).unwrap_err();
assert!(error.to_string().contains("stdin is not valid UTF-8"));
}
}

#[test]
fn prompt_input_reports_stdin_read_errors() {
struct Unreadable;
impl Read for Unreadable {
fn read(&mut self, _: &mut [u8]) -> std::io::Result<usize> {
Err(std::io::Error::new(
std::io::ErrorKind::Other,
"read failed",
))
}
}
let error = resolve_prompt(&[], Some(Path::new("-")), Unreadable).unwrap_err();
assert!(error.to_string().contains("failed to read stdin"));
}

fn permission_option(option_id: &str, kind: PermissionOptionKind) -> PermissionQueryOption {
PermissionQueryOption {
option_id: option_id.into(),
Expand Down
Loading
Loading