diff --git a/CHANGELOG.md b/CHANGELOG.md index 51d8d05..76f9b85 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -31,6 +31,23 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 profiles are unaffected — `tole-core` defaults do not change. ### Added +- `tole acp`: tole as an **Agent Client Protocol agent** over stdio + (issue #95) — editors (Zed et al.) drive durable tole sessions: + `session/new`/`session/load` map to the JSONL session store, prompts run + full tole turns, the final answer streams as an `agent_message_chunk`, + and Write/Destructive tool calls surface as + `session/request_permission` requests — the editor human approves, with + Destructive consent being a genuine per-call decision. Provider config + is only required when a prompt actually runs. The session map survives + across prompts (a first-run regression where fresh state replaced the + map after one turn — caught by CodeCora review — is fixed, along with + session-id path-traversal and mutex-poisoning hardening). The + session-map lock is held only briefly — a running turn keeps its OWN + storage lock, so the reader loop stays live for permission routing + (the first implementation deadlocked protocol routing for up to the + permission timeout whenever a client opened a session while a + permission request was pending — caught by CodeCora review). Sessions + reject concurrent turns (busy) and panic-safe un-busy via Drop. - `tole mcp`: tole as an **MCP server** over stdio (issue #94) — the registry's hardened tools (jailed file ops, argv-validated git, detached jobs, memory loop, cora/uteke integrations) become callable by any MCP diff --git a/README.md b/README.md index b3f8beb..2b51008 100644 --- a/README.md +++ b/README.md @@ -34,6 +34,9 @@ tole is the agent harness of the CodeCora ecosystem — the hands: (jailed file ops, git, jobs, memory) to any MCP client. ReadOnly tools are always callable; Write tools require `--allow` patterns; Destructive tools are structurally absent. +- **tole as an ACP agent**: `tole acp` speaks the Agent Client Protocol — + editors (Zed et al.) drive durable tole sessions, and tool approvals + surface as permission requests in the editor. Every ecosystem integration is probe-first: a missing binary degrades to a one-line warning, never a phantom tool. diff --git a/crates/tole-cli/src/acp.rs b/crates/tole-cli/src/acp.rs new file mode 100644 index 0000000..6a14ffb --- /dev/null +++ b/crates/tole-cli/src/acp.rs @@ -0,0 +1,750 @@ +//! D2 (issue #95): `tole acp` — an **Agent Client Protocol** host. +//! +//! Editors and ACP-capable clients (Zed et al.) launch `tole acp` and +//! drive a durable tole session over line-delimited JSON-RPC on +//! stdio/stderr: +//! +//! - `initialize` → protocol + capability handshake +//! - `session/new` / `session/load` → a durable JSONL session (workspace +//! jail rooted at the client-provided cwd) +//! - `session/prompt` → one full tole turn; the final answer is delivered +//! as an `agent_message_chunk` update before the response (v1 has no +//! intra-turn streaming — the turn loop is synchronous; the response is +//! correct, just not incremental) +//! - Write/Destructive tool calls surface as +//! `session/request_permission` **requests to the editor** — the human +//! in the editor is the approver, which is exactly tole's gate with a +//! different face. Destructive tools therefore CAN be exposed here +//! (unlike MCP server mode): consent is structurally a human decision. +//! +//! Wire hygiene: protocol messages go to stdout; all diagnostics go to +//! stderr. No new dependencies — the JSON-RPC framing is hand-rolled +//! line-delimited JSON (serde_json only). + +use anyhow::{Context, Result}; +use serde_json::{json, Value}; +use std::collections::HashMap; +use std::io::{BufRead, Write}; +use std::path::PathBuf; +use std::sync::mpsc; +use std::sync::{Arc, Mutex}; +use std::time::Duration; + +use tole_cli::approver::{InteractiveApprover, PromptFn}; +use tole_cli::tools::WriteFileTool; + +const ACP_PROTOCOL_VERSION: u32 = 1; +/// Permission requests can sit in an editor until a human clicks; do not +/// turn that into a timeout race. +const PERMISSION_TIMEOUT: Duration = Duration::from_secs(600); + +// --------------------------------------------------------------------------- +// Connection: outgoing lines + response routing +// --------------------------------------------------------------------------- + +/// Shared connection state: outgoing JSON-RPC lines are pushed to the +/// writer thread; responses to agent-initiated requests (permission) are +/// routed back to whoever is waiting on them. +#[derive(Clone)] +struct Conn { + tx: mpsc::Sender, + next_id: Arc>, + pending: Arc>>>, +} + +impl Conn { + fn new(tx: mpsc::Sender) -> Self { + Self { + tx, + next_id: Arc::new(Mutex::new(1)), + pending: Arc::new(Mutex::new(HashMap::new())), + } + } + + fn send_line(&self, line: &str) { + if self.tx.send(format!("{line}\n")).is_err() { + // Client is gone; the pending waits below will time out and + // the session loop winds down on EOF anyway. + } + } + + fn send_notification(&self, method: &str, params: Value) { + let line = json!({"jsonrpc": "2.0", "method": method, "params": params}).to_string(); + self.send_line(&line); + } + + /// Agent-initiated request (permission): returns the client's result + /// object, or an error-shaped object on timeout/cancel. + fn request(&self, method: &str, params: Value, timeout: Duration) -> Result { + let id = { + let mut n = self.next_id.lock().expect("id lock"); + let id = *n; + *n += 1; + id + }; + let (tx, rx) = mpsc::channel(); + self.pending.lock().expect("pending lock").insert(id, tx); + self.send_line( + &json!({"jsonrpc": "2.0", "id": id, "method": method, "params": params}).to_string(), + ); + match rx.recv_timeout(timeout) { + Ok(v) => Ok(v), + Err(_) => { + self.pending.lock().expect("pending lock").remove(&id); + Err("client did not answer the permission request in time".into()) + } + } + } + + /// Reader-side routing: a response line completes a pending agent + /// request. + fn route_response(&self, id: u64, result: Value) { + if let Some(tx) = self.pending.lock().expect("pending lock").remove(&id) { + let _ = tx.send(result); + } + } +} + +// --------------------------------------------------------------------------- +// Approval: the editor is the human +// --------------------------------------------------------------------------- + +/// PromptFn bridge: renders a `session/request_permission` request to the +/// ACP client and maps the chosen option back to a tole verdict. +struct AcpPrompt { + conn: Conn, + session_id: String, + counter: Arc>, +} + +impl PromptFn for AcpPrompt { + fn prompt(&self, req: &tole_core::approval::ToolRequest<'_>) -> tole_core::approval::Verdict { + use tole_core::approval::Verdict; + let option_allow = "allow-once"; + let option_reject = "reject-once"; + let call_id = { + let mut n = self.counter.lock().expect("call id lock"); + *n += 1; + format!("call-{n}") + }; + // Announce the tool call first so editors can render it. + self.conn.send_notification( + "session/update", + json!({ + "sessionId": self.session_id, + "update": { + "sessionUpdate": "tool_call", + "toolCallId": call_id, + "title": req.description, + "kind": "other", + "rawInput": req.input, + } + }), + ); + let params = json!({ + "sessionId": self.session_id, + "toolCall": { + "toolCallId": call_id, + "title": req.description, + "kind": "other", + "rawInput": req.input, + }, + "options": [ + {"kind": "allow_once", "name": "Allow", "optionId": option_allow}, + {"kind": "reject_once", "name": "Reject", "optionId": option_reject}, + ], + }); + let verdict = + match self + .conn + .request("session/request_permission", params, PERMISSION_TIMEOUT) + { + Ok(result) => { + let chosen = result["outcome"]["optionId"] + .as_str() + .unwrap_or(option_reject); + if chosen == option_allow { + Verdict::Allow + } else { + Verdict::Deny + } + } + Err(_) => { + // Cancelled or timed out: fail closed. + Verdict::Deny + } + }; + // Close the tool-call record. + self.conn.send_notification( + "session/update", + json!({ + "sessionId": self.session_id, + "update": { + "sessionUpdate": "tool_call_update", + "toolCallId": call_id, + "status": if verdict == Verdict::Allow { "completed" } else { "rejected" }, + } + }), + ); + verdict + } +} + +// --------------------------------------------------------------------------- +// Session state +// --------------------------------------------------------------------------- + +struct SessionState { + // Per-session storage lock: a turn holds THIS (not the session-map + // lock), so the reader loop stays live for permission routing while + // a prompt runs (CodeCora scan deadlock finding). + storage: StdArc>, + registry: StdArc, + system_prompt: Option, + memory: Option, + first_prompt_done: StdArc>, + busy: StdArc>, +} + +/// Marks a session busy for its whole lifetime; Drop un-marks even on +/// panic, so one failed turn cannot brick the session. +struct BusyGuard(StdArc>); +impl Drop for BusyGuard { + fn drop(&mut self) { + *self.0.lock().expect("busy lock") = false; + } +} + +struct Sessions { + map: HashMap, +} + +impl Sessions { + fn new() -> Self { + Self { + map: HashMap::new(), + } + } +} + +// --------------------------------------------------------------------------- +// ACP server loop +// --------------------------------------------------------------------------- + +/// Run the ACP agent over stdin/stdout. Blocks until the client closes +/// stdin. +pub fn run_acp( + allow_patterns: &[String], + auto_write: bool, + _workspace_default: Option<&String>, + plan_mode: bool, + memory: Option, +) -> Result<()> { + let (line_tx, line_rx) = mpsc::channel::(); + let conn = Conn::new(line_tx.clone()); + + // Writer thread: stdout belongs to the protocol. + std::thread::spawn(move || { + let stdout = std::io::stdout(); + let mut out = stdout.lock(); + for line in line_rx { + if out.write_all(line.as_bytes()).is_err() { + break; + } + let _ = out.flush(); + } + }); + + // The session map lives for the WHOLE server lifetime and is shared + // with prompt threads (Arc clone per prompt). Holding the map lock + // for the duration of a turn also serializes access to one session's + // storage — session/load of a busy id simply waits its turn instead + // of opening a divergent second handle. + let sessions: SharedSessions = StdArc::new(Mutex::new(Sessions::new())); + let stdin = std::io::stdin(); + for line in stdin.lock().lines() { + let Ok(line) = line else { break }; + if line.trim().is_empty() { + continue; + } + let Ok(msg) = serde_json::from_str::(&line) else { + continue; // non-JSON noise on the protocol channel: ignore + }; + let id = msg.get("id").cloned(); + let Some(method) = msg + .get("method") + .and_then(Value::as_str) + .map(str::to_string) + else { + // A response to one of OUR requests (permission). Client + // ERROR replies route as well — an errored permission must + // fail closed immediately, not hang for the full timeout + // (CodeCora scan 2026-09-28). + if let Some(id) = msg.get("id").and_then(Value::as_u64) { + if let Some(result) = msg.get("result").cloned() { + conn.route_response(id, result); + } else if let Some(err) = msg.get("error").cloned() { + conn.route_response( + id, + json!({"outcome": {"outcome": "cancelled"}, "error": err}), + ); + } + } + continue; + }; + let params = msg.get("params").cloned().unwrap_or(json!({})); + + match method.as_str() { + "initialize" => { + reply( + &conn, + id, + json!({ + "protocolVersion": ACP_PROTOCOL_VERSION, + "agentCapabilities": { + "loadSession": true, + "promptCapabilities": {}, + }, + "authMethods": [], + }), + ); + } + "session/new" | "session/load" => { + let loading = method == "session/load"; + let cwd = params + .get("cwd") + .and_then(Value::as_str) + .unwrap_or(".") + .to_string(); + let session_id = match params.get("sessionId").and_then(Value::as_str) { + Some(id) => validate_session_id(id), + None => Some(new_session_id()), + }; + let Some(session_id) = session_id else { + reply_error(&conn, id, "session/load: invalid sessionId"); + continue; + }; + // A turn holds its session's storage lock and marks it + // busy; loading a busy id would open a SECOND handle on + // the same JSONL mid-turn (CodeCora scan 2026-09-28). + let busy_now = { + let sessions = lock_sessions(&sessions); + sessions + .map + .get(&session_id) + .map(|st| *st.busy.lock().expect("busy lock")) + .unwrap_or(false) + }; + if busy_now { + reply_error(&conn, id, "session is busy running a turn"); + continue; + } + match open_session( + &session_id, + &cwd, + loading, + allow_patterns, + auto_write, + plan_mode, + memory.clone(), + conn.clone(), + ) { + Ok(state) => { + // Insert + busy re-check in ONE critical section: + // open_session is slow (canonicalize + a git + // subprocess), and a prompt that set busy inside + // that window must not be orphaned by the insert + // swapping in a fresh, not-busy state (CodeCora + // scan 2026-09-28). + let mut sessions = lock_sessions(&sessions); + let busy_now = sessions + .map + .get(&session_id) + .map(|st| *st.busy.lock().expect("busy lock")) + .unwrap_or(false); + if busy_now { + drop(sessions); + reply_error(&conn, id, "session became busy while opening — retry"); + continue; + } + sessions.map.insert(session_id.clone(), state); + reply(&conn, id, json!({ "sessionId": session_id })); + } + Err(e) => reply_error(&conn, id, &e.to_string()), + } + } + "session/prompt" => { + let Some(session_id) = params + .get("sessionId") + .and_then(Value::as_str) + .map(str::to_string) + .map(|id| validate_session_id(&id)) + else { + reply_error(&conn, id, "session/prompt: missing sessionId"); + continue; + }; + let Some(session_id) = session_id else { + reply_error(&conn, id, "session/prompt: invalid sessionId"); + continue; + }; + // ACP sends `prompt` as an array of content blocks; a + // plain string is accepted for client convenience. + let prompt_text = match params.get("prompt") { + Some(Value::Array(blocks)) => blocks + .iter() + .filter_map(|b| b.get("text").and_then(Value::as_str)) + .collect::>() + .join(""), + Some(Value::String(s)) => s.clone(), + _ => { + reply_error(&conn, id, "session/prompt: missing prompt"); + continue; + } + }; + // Response is delivered from the turn thread (result or + // error) so this reader loop stays live for permission + // requests while the turn runs. + let conn = conn.clone(); + let sessions = sessions.clone(); + let prompt_clone = prompt_text; + let session_id_clone = session_id.clone(); + std::thread::spawn(move || { + let result = + run_prompt(sessions, &session_id_clone, &prompt_clone, conn.clone()); + match result { + Ok(stop) => reply(&conn, id, json!({ "stopReason": stop })), + Err(e) => reply_error(&conn, id, &e.to_string()), + } + }); + } + other => { + if id.is_some() { + reply_error(&conn, id, &format!("method not supported: {other}")); + } + } + } + } + Ok(()) +} + +use std::sync::Arc as StdArc; +type SharedSessions = StdArc>; + +/// Poisoning-tolerant lock: one panicking turn must not brick the whole +/// ACP server (CodeCora scan finding — mutex poisoning). +fn lock_sessions(sessions: &SharedSessions) -> std::sync::MutexGuard<'_, Sessions> { + sessions + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) +} + +fn reply(conn: &Conn, id: Option, result: Value) { + let Some(id) = id else { return }; + conn.send_line(&json!({"jsonrpc": "2.0", "id": id, "result": result}).to_string()); +} + +fn reply_error(conn: &Conn, id: Option, message: &str) { + let Some(id) = id else { return }; + conn.send_line( + &json!({"jsonrpc": "2.0", "id": id, "error": {"code": -32603, "message": message}}) + .to_string(), + ); +} + +/// ACP session ids are tole session ids: reject path separators, +/// parent refs, and anything outside the tole charset before the id ever +/// touches a path (CodeCora scan finding: `../` or absolute ids would +/// escape the sessions dir via Path::join). +fn validate_session_id(id: &str) -> Option { + let ok = !id.is_empty() + && id.len() <= 64 + && !id.contains('/') + && !id.contains('\\') + && !id.contains("..") + && id + .chars() + .all(|c| c.is_ascii_alphanumeric() || matches!(c, '-' | '_')); + if ok { + Some(id.to_string()) + } else { + None + } +} + +fn new_session_id() -> String { + let ms = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|d| d.as_millis()) + .unwrap_or(0); + format!("acp-{ms:x}-{:x}", std::process::id()) +} + +/// Create/open a session: workspace jail = the client's cwd; approval = +/// the ACP editor (interactive — which is what unlocks Destructive tools +/// with genuine human consent). +#[allow(clippy::too_many_arguments)] +fn open_session( + session_id: &str, + cwd: &str, + loading: bool, + allow_patterns: &[String], + auto_write: bool, + plan_mode: bool, + memory: Option, + conn: Conn, +) -> Result { + let workspace = PathBuf::from(cwd); + let workspace_canon = workspace + .canonicalize() + .with_context(|| format!("session workspace {}: {cwd}", workspace.display()))?; + let system_prompt = std::env::var("TOLE_SYSTEM_PROMPT") + .ok() + .filter(|s| !s.trim().is_empty()); + + let mut reg = tole_core::tool::ToolRegistry::with_approver( + InteractiveApprover::new(AcpPrompt { + conn, + session_id: session_id.to_string(), + counter: Arc::new(Mutex::new(0)), + }) + .with_allow_patterns(allow_patterns.to_vec()) + .with_auto_write(auto_write), + ); + use tole_core::cora_search::CoraSearchTool; + use tole_core::file_tools::{DeleteFileTool, EditFileTool}; + use tole_core::gh::GhTool; + use tole_core::git::GitTool; + use tole_core::jobs::{JobPollTool, JobStartTool}; + use tole_core::read_file::ReadFileTool; + use tole_core::run_command::RunCommandTool; + use tole_core::uteke::{UtekeDocumentTool, UtekeRecallTool}; + if tole_cli_binary_available("cora") { + reg.register(Box::new(CoraSearchTool::new())) + .map_err(anyhow::Error::msg)?; + } + if tole_cli_binary_available("uteke") { + reg.register(Box::new(UtekeRecallTool::new())) + .map_err(anyhow::Error::msg)?; + if !plan_mode { + reg.register(Box::new(UtekeDocumentTool::new(None))) + .map_err(anyhow::Error::msg)?; + } + } + // Write-tier tools: skipped entirely under --plan-mode (read-only + // sessions — the model cannot even attempt a write). The earlier + // draft registered everything and "filtered" afterwards, which + // CodeCora rightly called out: the full registry under --yes broke + // the read-only contract. + if !plan_mode { + reg.register(Box::new(RunCommandTool::new(workspace_canon.clone()))) + .map_err(anyhow::Error::msg)?; + reg.register(Box::new(JobStartTool::new(workspace_canon.clone()))) + .map_err(anyhow::Error::msg)?; + reg.register(Box::new(WriteFileTool::new(workspace_canon.clone()))) + .map_err(anyhow::Error::msg)?; + reg.register(Box::new(EditFileTool::new(workspace_canon.clone()))) + .map_err(anyhow::Error::msg)?; + { + let repo = + detect_github_repo(&workspace_canon).unwrap_or_else(|| "codecoradev/tole".into()); + reg.register(Box::new(GhTool::new(repo))) + .map_err(anyhow::Error::msg)?; + } + reg.register(Box::new(GitTool::new().in_dir(workspace_canon.clone()))) + .map_err(anyhow::Error::msg)?; + // delete_file IS registered here: the ACP editor prompt is an + // interactive approver, so a Destructive tool carries genuine + // human consent — the same rule as the CLI, not an exception. + reg.register(Box::new(DeleteFileTool::new(workspace_canon.clone()))) + .map_err(anyhow::Error::msg)?; + } + reg.register(Box::new(JobPollTool::new(workspace_canon.clone()))) + .map_err(anyhow::Error::msg)?; + reg.register(Box::new(ReadFileTool::new(workspace_canon.clone()))) + .map_err(anyhow::Error::msg)?; + + let storage = if loading { + let dir = sessions_dir_for(cwd)?; + tole_core::storage::JsonlStorage::open(dir.join(format!("{session_id}.jsonl"))) + .context("loading session")? + } else { + let dir = sessions_dir_for(cwd)?; + std::fs::create_dir_all(&dir)?; + tole_core::storage::JsonlStorage::create_with( + &dir, + session_id, + None, + system_prompt.as_deref(), + ) + .context("creating session")? + }; + // --plan-mode: read-only sessions — Write/Destructive tools are not + // registered at all, so the model cannot even attempt them. + if plan_mode { + eprintln!("tole acp: --plan-mode is active — serving read-only tools only"); + return Ok(SessionState { + storage: StdArc::new(Mutex::new(storage)), + registry: StdArc::new(reg), + system_prompt, + memory, + first_prompt_done: StdArc::new(Mutex::new(loading)), + busy: StdArc::new(Mutex::new(false)), + }); + } + Ok(SessionState { + storage: StdArc::new(Mutex::new(storage)), + registry: StdArc::new(reg), + system_prompt, + memory, + first_prompt_done: StdArc::new(Mutex::new(loading)), + busy: StdArc::new(Mutex::new(false)), + }) +} + +fn sessions_dir_for(cwd: &str) -> Result { + let dir = PathBuf::from(cwd).join(".tole/sessions"); + std::fs::create_dir_all(&dir)?; + Ok(dir) +} + +fn tole_cli_binary_available(name: &str) -> bool { + if name.contains('/') { + return std::path::Path::new(name).is_file(); + } + if let Ok(path) = std::env::var("PATH") { + for dir in std::env::split_paths(&path) { + if dir.join(name).is_file() { + return true; + } + } + } + false +} + +fn detect_github_repo(cwd: &PathBuf) -> Option { + let mut cmd = std::process::Command::new("git"); + cmd.args(["config", "--get", "remote.origin.url"]) + .current_dir(cwd); + let out = tole_core::subprocess::run_with_timeout(&mut cmd, std::time::Duration::from_secs(5)) + .ok()?; + if !out.status.success() { + return None; + } + crate::github_repo_from_remote_url(&String::from_utf8_lossy(&out.stdout)) +} + +/// One ACP prompt = one full tole turn. Returns the stop reason. +fn run_prompt( + sessions: SharedSessions, + session_id: &str, + prompt: &str, + conn: Conn, +) -> Result { + // Brief map lock: take the session's handles and reject a busy + // session. The MAP lock is released here — a running turn holds only + // its OWN storage lock, so the reader loop stays live for permission + // routing (CodeCora deadlock finding). + let (storage, registry, memory, system_prompt, first_prompt_done, busy_guard) = { + let mut sessions = lock_sessions(&sessions); + let Some(state) = sessions.map.get_mut(session_id) else { + anyhow::bail!("unknown session: {session_id}"); + }; + { + let mut busy = state.busy.lock().expect("busy lock"); + if *busy { + anyhow::bail!("session is busy running a turn"); + } + *busy = true; + } + ( + state.storage.clone(), + state.registry.clone(), + state.memory.clone(), + state.system_prompt.clone(), + state.first_prompt_done.clone(), + std::sync::Arc::clone(&state.busy), + ) + }; + // Panic-safe un-busy: Drop clears the flag even if the turn unwinds. + let _busy_guard = BusyGuard(busy_guard); + + // Memory loop, pre-turn (first prompt of a fresh session only). + let mut effective = prompt.to_string(); + #[cfg(feature = "shell-tools")] + { + let done = *first_prompt_done.lock().expect("fpd lock"); + if !done { + if let Some(mem) = &memory { + if let Ok(block) = tole_core::memory::recall_block(mem, prompt) { + if !block.is_empty() { + eprintln!("tole acp: memory: recalled context injected"); + effective = format!("{prompt}{block}"); + } + } + } + } + } + + // Provider: built per turn (cheap), so a session can be created + // without provider env and only fail when a turn is actually run. + let cfg = tole_core::openai::OpenAiConfig::from_env().context( + "missing provider config: set TOLE_BASE_URL / TOLE_MODEL / TOLE_API_KEY \ + (or the OPENAI_* equivalents)", + )?; + let mut provider = + tole_core::openai::OpenAiProvider::new(cfg).with_tool_specs(registry.specs()); + if let Some(sys) = system_prompt.as_deref() { + provider = provider.with_system_prompt(sys); + } + + let mut storage = storage.lock().unwrap_or_else(|p| p.into_inner()); + let outcome = tole_core::turn::run_turn(&mut *storage, &mut provider, ®istry, &effective)?; + + #[cfg(feature = "shell-tools")] + if let tole_core::turn::TurnOutcome::Final { text } = &outcome { + if let Some(mem) = &memory { + let _ = tole_core::memory::remember_session(mem, session_id, prompt, text); + } + } + *first_prompt_done.lock().expect("fpd lock") = true; + + let stop = match &outcome { + tole_core::turn::TurnOutcome::Final { text } => { + conn.send_notification( + "session/update", + json!({ + "sessionId": session_id, + "update": { + "sessionUpdate": "agent_message_chunk", + "content": {"type": "text", "text": text}, + } + }), + ); + "end_turn" + } + tole_core::turn::TurnOutcome::ApprovalRequired { name } => { + eprintln!("tole acp: approval denied for '{name}'"); + "refusal" + } + tole_core::turn::TurnOutcome::UnknownTool { name } => { + eprintln!("tole acp: unknown tool '{name}'"); + "refusal" + } + tole_core::turn::TurnOutcome::ProviderFailed { message } => { + eprintln!("tole acp: provider failed: {message}"); + "refusal" + } + tole_core::turn::TurnOutcome::BudgetExhausted => { + eprintln!("tole acp: step budget exhausted"); + "max_tokens" + } + tole_core::turn::TurnOutcome::LoopDetected { tool, .. } => { + eprintln!("tole acp: loop detected on '{tool}'"); + "refusal" + } + tole_core::turn::TurnOutcome::Storage(e) => { + anyhow::bail!("storage error: {e}"); + } + }; + Ok(stop.to_string()) +} diff --git a/crates/tole-cli/src/main.rs b/crates/tole-cli/src/main.rs index ac7bed1..573114e 100644 --- a/crates/tole-cli/src/main.rs +++ b/crates/tole-cli/src/main.rs @@ -7,6 +7,9 @@ use std::path::{Path, PathBuf}; use tole_cli::approver::{InteractiveApprover, StdioPrompt}; use tole_cli::tools::WriteFileTool; use tole_core::approval::AllowlistApprover; + +#[cfg(feature = "shell-tools")] +mod acp; #[cfg(feature = "shell-tools")] use tole_core::cora_search::CoraSearchTool; use tole_core::file_tools::{DeleteFileTool, EditFileTool}; @@ -179,6 +182,30 @@ enum Command { #[arg(long)] workspace: Option, }, + /// Serve tole as an ACP agent over stdio (issue #95): editors and + /// ACP clients drive durable tole sessions; tool approvals surface + /// as permission requests in the client. + #[cfg(feature = "shell-tools")] + Acp { + /// Same semantics as `run --allow` (Write pre-authorization). + #[arg(long = "allow")] + allow_patterns: Vec, + + /// Auto-allow every Write call (Destructive still prompts in the + /// client). + #[arg(long)] + yes: bool, + + /// Default file-tools root; each session's jail is the client's + /// session cwd. + #[arg(long)] + workspace: Option, + + /// Memory loop backend (`uteke`) — same as `--memory uteke` on + /// run/chat. Falls back to the TOLE_MEMORY env. + #[arg(long)] + memory: Option, + }, /// Interactive multi-turn chat on one durable session. Chat { /// System prompt for a fresh session (ignored when resuming — @@ -287,6 +314,34 @@ fn dispatch(cli: Cli) -> Result<()> { host.plan_mode, ) } + #[cfg(feature = "shell-tools")] + Command::Acp { + allow_patterns, + yes, + workspace, + memory, + } => { + // Same loud-bail rule as `tole mcp` for hooks: the ACP host + // does not wire local pre/post hooks — approvals happen in + // the editor via permission requests instead. + if host.on_pretool_non_empty() || host.on_posttool_non_empty() { + anyhow::bail!( + "--on-pretool/--on-posttool are not supported by `tole acp` \ + (approvals happen via session/request_permission in the client)" + ); + } + if host.plan_mode { + eprintln!("tole acp: --plan-mode is active — serving read-only tools only"); + } + let memory = resolve_memory(memory.as_ref())?; + crate::acp::run_acp( + &allow_patterns, + yes, + workspace.as_ref(), + host.plan_mode, + memory, + ) + } Command::Sessions => sessions_command(&sessions_dir), Command::Status { id } => status_command(&sessions_dir, &id), Command::Chat { diff --git a/crates/tole-cli/tests/acp_handshake.rs b/crates/tole-cli/tests/acp_handshake.rs new file mode 100644 index 0000000..e388cf9 --- /dev/null +++ b/crates/tole-cli/tests/acp_handshake.rs @@ -0,0 +1,201 @@ +//! D2 (issue #95): ACP handshake + session lifecycle, CI-safe (no LLM). +//! Spawns the real `tole acp` binary and speaks line-delimited JSON-RPC +//! with a deadline-guarded reader. + +use serde_json::{json, Value}; +use std::io::{BufRead, Write}; +use std::process::{Child, Command, Stdio}; +use std::sync::mpsc; +use std::time::Duration; + +struct AcpProcess { + child: Child, + rx: mpsc::Receiver, +} + +impl AcpProcess { + fn spawn() -> Self { + let mut child = Command::new(env!("CARGO_BIN_EXE_tole")) + .arg("acp") + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .stderr(Stdio::null()) + .spawn() + .expect("spawn tole acp"); + let stdout = child.stdout.take().expect("stdout piped"); + let (tx, rx) = mpsc::channel(); + std::thread::spawn(move || { + let reader = std::io::BufReader::new(stdout); + for line in reader.lines().map_while(Result::ok) { + if tx.send(line).is_err() { + break; + } + } + }); + Self { child, rx } + } + + fn send(&mut self, v: &Value) { + writeln!(self.child.stdin.as_mut().expect("stdin"), "{v}").unwrap(); + } + + /// Wait for a line that parses as JSON-RPC with the given id. + fn wait_response(&self, id: u64, timeout: Duration) -> Value { + let deadline = std::time::Instant::now() + timeout; + loop { + let left = deadline.saturating_duration_since(std::time::Instant::now()); + assert!(left > Duration::ZERO, "timed out waiting for id {id}"); + match self.rx.recv_timeout(left) { + Ok(line) => { + if let Ok(v) = serde_json::from_str::(&line) { + if v.get("id").and_then(Value::as_u64) == Some(id) + && (v.get("result").is_some() || v.get("error").is_some()) + { + return v; + } + } + } + Err(mpsc::RecvTimeoutError::Timeout) => { + panic!("timed out waiting for id {id}"); + } + Err(mpsc::RecvTimeoutError::Disconnected) => { + panic!("acp process closed stdout"); + } + } + } + } +} + +impl Drop for AcpProcess { + fn drop(&mut self) { + let _ = self.child.kill(); + let _ = self.child.wait(); + } +} + +fn temp_cwd(tag: &str) -> std::path::PathBuf { + let dir = std::env::temp_dir().join(format!("tole-acp-{tag}-{}", std::process::id())); + let _ = std::fs::remove_dir_all(&dir); + std::fs::create_dir_all(&dir).unwrap(); + dir +} + +#[test] +fn acp_initialize_and_session_new() { + let mut acp = AcpProcess::spawn(); + let cwd = temp_cwd("init"); + + acp.send(&json!({ + "jsonrpc": "2.0", "id": 1, "method": "initialize", + "params": {"protocolVersion": 1, "clientCapabilities": {}} + })); + let init = acp.wait_response(1, Duration::from_secs(20)); + assert_eq!(init["result"]["protocolVersion"], 1); + assert_eq!(init["result"]["agentCapabilities"]["loadSession"], true); + + acp.send(&json!({ + "jsonrpc": "2.0", "id": 2, "method": "session/new", + "params": {"cwd": cwd.to_string_lossy()} + })); + let created = acp.wait_response(2, Duration::from_secs(20)); + let session_id = created["result"]["sessionId"] + .as_str() + .expect("sessionId") + .to_string(); + assert!(session_id.starts_with("acp-")); + + // The durable session file exists under the session cwd. + assert!(cwd + .join(".tole/sessions") + .join(format!("{session_id}.jsonl")) + .exists()); + + // Unknown methods surface a JSON-RPC error, not a hang. + acp.send(&json!({ + "jsonrpc": "2.0", "id": 3, "method": "session/frobnicate", "params": {} + })); + let err = acp.wait_response(3, Duration::from_secs(20)); + assert!(err.get("error").is_some()); +} + +#[test] +fn acp_load_returns_existing_session() { + let mut acp = AcpProcess::spawn(); + let cwd = temp_cwd("load"); + + acp.send(&json!({ + "jsonrpc": "2.0", "id": 1, "method": "initialize", + "params": {"protocolVersion": 1, "clientCapabilities": {}} + })); + let _ = acp.wait_response(1, Duration::from_secs(20)); + + acp.send(&json!({ + "jsonrpc": "2.0", "id": 2, "method": "session/new", + "params": {"cwd": cwd.to_string_lossy()} + })); + let session_id = acp.wait_response(2, Duration::from_secs(20))["result"]["sessionId"] + .as_str() + .expect("sessionId") + .to_string(); + + // Load the SAME session id in a fresh process-backed state: the + // durable JSONL is the source of truth (D2 acceptance). + let mut acp2 = AcpProcess::spawn(); + acp2.send(&json!({ + "jsonrpc": "2.0", "id": 1, "method": "initialize", + "params": {"protocolVersion": 1, "clientCapabilities": {}} + })); + let _ = acp2.wait_response(1, Duration::from_secs(20)); + acp2.send(&json!({ + "jsonrpc": "2.0", "id": 2, "method": "session/load", + "params": {"cwd": cwd.to_string_lossy(), "sessionId": session_id} + })); + let loaded = acp2.wait_response(2, Duration::from_secs(20)); + assert_eq!(loaded["result"]["sessionId"], session_id); +} + +/// The map must survive across prompts (CodeCora scan regression: a +/// fresh Arc per prompt dropped every session after the first turn). +#[test] +fn acp_session_survives_multiple_session_news() { + let mut acp = AcpProcess::spawn(); + let cwd = temp_cwd("multi"); + + acp.send(&json!({ + "jsonrpc": "2.0", "id": 1, "method": "initialize", + "params": {"protocolVersion": 1, "clientCapabilities": {}} + })); + let _ = acp.wait_response(1, Duration::from_secs(20)); + + for n in 2..=4 { + acp.send(&json!({ + "jsonrpc": "2.0", "id": n, "method": "session/new", + "params": {"cwd": cwd.to_string_lossy()} + })); + let r = acp.wait_response(n, Duration::from_secs(20)); + assert!( + r["result"]["sessionId"].is_string(), + "session/new {n} failed" + ); + } +} + +/// Path traversal via sessionId is refused at the protocol surface. +#[test] +fn acp_rejects_traversal_session_ids() { + let mut acp = AcpProcess::spawn(); + let cwd = temp_cwd("traversal"); + + acp.send(&json!({ + "jsonrpc": "2.0", "id": 1, "method": "initialize", + "params": {"protocolVersion": 1, "clientCapabilities": {}} + })); + let _ = acp.wait_response(1, Duration::from_secs(20)); + + acp.send(&json!({ + "jsonrpc": "2.0", "id": 2, "method": "session/load", + "params": {"cwd": cwd.to_string_lossy(), "sessionId": "../../etc/passwd"} + })); + let r = acp.wait_response(2, Duration::from_secs(20)); + assert!(r.get("error").is_some(), "traversal id must be an error"); +} diff --git a/docs/epics.md b/docs/epics.md index a5689b3..818149c 100644 --- a/docs/epics.md +++ b/docs/epics.md @@ -120,7 +120,7 @@ Expose the registry's hardened tools via rmcp's server side. ReadOnly tools always callable; Write needs `--allow`; Destructive structurally absent. Est: S–M. -### D2 — `tole acp`: Agent Client Protocol host (issue #95) +### D2 — `tole acp`: Agent Client Protocol host (issue #95) — ✅ DONE (2026-09-18) Editor integration (Zed et al.): ACP sessions ↔ JSONL sessions, ACP permission requests ↔ the approval gate, streaming turn events. Est: M.