From 6f0b30c18a45d9d00e129945de01640f5748670b Mon Sep 17 00:00:00 2001 From: ajianaz Date: Fri, 18 Sep 2026 16:48:55 +0700 Subject: [PATCH 1/5] =?UTF-8?q?feat:=20tole=20acp=20=E2=80=94=20Agent=20Cl?= =?UTF-8?q?ient=20Protocol=20host=20(D2,=20issue=20#95)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit tole as an ACP agent over stdio: editors and ACP-capable clients (Zed et al.) drive durable tole sessions via line-delimited JSON-RPC, with no new dependencies. - initialize / session/new / session/load / session/prompt: ACP sessions map to the durable JSONL store (session jail = the client's cwd; load replays the existing log). - session/prompt runs one full tole turn and delivers the final answer as an agent_message_chunk update before responding (v1: no intra-turn streaming — the synchronous turn loop is untouched). - Approval bridge: Write/Destructive tool calls surface as session/request_permission requests to the EDITOR — the human in the client is the approver, which is why Destructive tools CAN be registered here with genuine per-call consent (unlike MCP server mode), without weakening anything. - Prompt accepts the ACP content-block array form (and plain strings). - run_command/git/gh-detect operate on the SESSION cwd, not the process cwd (CodeCora review finding on this PR). - Provider config is only required when a prompt actually runs. - CI-safe integration test spawns the real binary: initialize, session lifecycle, unknown-method error; live E2E covered permission flow, streaming chunk, and stop reasons. Closes #95 (Phase D2). --- CHANGELOG.md | 8 + README.md | 3 + crates/tole-cli/src/acp.rs | 617 +++++++++++++++++++++++++ crates/tole-cli/src/main.rs | 37 ++ crates/tole-cli/tests/acp_handshake.rs | 155 +++++++ docs/epics.md | 2 +- 6 files changed, 821 insertions(+), 1 deletion(-) create mode 100644 crates/tole-cli/src/acp.rs create mode 100644 crates/tole-cli/tests/acp_handshake.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index 51d8d05..e3078f4 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -31,6 +31,14 @@ 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. - `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..2ec76bf --- /dev/null +++ b/crates/tole-cli/src/acp.rs @@ -0,0 +1,617 @@ +//! 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 { + storage: tole_core::storage::JsonlStorage, + registry: tole_core::tool::ToolRegistry, + system_prompt: Option, + memory: Option, + first_prompt_done: bool, +} + +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>, + 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(); + } + }); + + let mut sessions = 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). + if let (Some(id), Some(result)) = ( + msg.get("id").and_then(Value::as_u64), + msg.get("result").cloned(), + ) { + conn.route_response(id, result); + } + 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 = params + .get("sessionId") + .and_then(Value::as_str) + .map(str::to_string) + .unwrap_or_else(new_session_id); + match open_session( + &session_id, + &cwd, + loading, + allow_patterns, + auto_write, + memory.clone(), + conn.clone(), + ) { + Ok(state) => { + 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) + else { + reply_error(&conn, id, "session/prompt: missing 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_ptr = share_sessions(&mut sessions); + let prompt_clone = prompt_text; + let session_id_clone = session_id.clone(); + std::thread::spawn(move || { + let result = + run_prompt(sessions_ptr, &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(()) +} + +// The sessions map must be shareable with the per-prompt threads; wrap it +// here so the method bodies stay readable. +use std::sync::Arc as StdArc; +type SharedSessions = StdArc>; + +fn share_sessions(sessions: &mut Sessions) -> SharedSessions { + // A fresh Arc per ACP server instance is fine: the reader loop owns + // the map for the process lifetime and hands clones to turn threads. + StdArc::new(Mutex::new(std::mem::replace(sessions, Sessions::new()))) +} + +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(), + ); +} + +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, + 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)?; + reg.register(Box::new(UtekeDocumentTool::new(None))) + .map_err(anyhow::Error::msg)?; + } + 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(JobPollTool::new(workspace_canon.clone()))) + .map_err(anyhow::Error::msg)?; + reg.register(Box::new(ReadFileTool::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)?; + + 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")? + }; + Ok(SessionState { + storage, + registry: reg, + system_prompt, + memory, + first_prompt_done: loading, + }) +} + +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 { + let mut sessions = sessions.lock().expect("sessions lock"); + let Some(state) = sessions.map.get_mut(session_id) else { + anyhow::bail!("unknown session: {session_id}"); + }; + + // Memory loop, pre-turn (first prompt of a fresh session only). + let mut effective = prompt.to_string(); + #[cfg(feature = "shell-tools")] + if let Some(mem) = state.memory.clone() { + if !state.first_prompt_done { + 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(state.registry.specs()); + if let Some(sys) = state.system_prompt.as_deref() { + provider = provider.with_system_prompt(sys); + } + + let outcome = tole_core::turn::run_turn( + &mut state.storage, + &mut provider, + &state.registry, + &effective, + )?; + + #[cfg(feature = "shell-tools")] + if let tole_core::turn::TurnOutcome::Final { text } = &outcome { + if let Some(mem) = state.memory.clone() { + let _ = tole_core::memory::remember_session(&mem, session_id, prompt, text); + } + } + state.first_prompt_done = 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 ed6ceaa..eb42ce6 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}; @@ -152,6 +155,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 (B1). Chat { /// System prompt for a fresh session (ignored when resuming — @@ -234,6 +261,16 @@ fn dispatch(cli: Cli) -> Result<()> { allow_patterns, workspace, } => mcp_server_command(workspace.as_ref(), &allow_patterns), + #[cfg(feature = "shell-tools")] + Command::Acp { + allow_patterns, + yes, + workspace, + memory, + } => { + let memory = resolve_memory(memory.as_ref())?; + crate::acp::run_acp(&allow_patterns, yes, workspace.as_ref(), 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..899bfa1 --- /dev/null +++ b/crates/tole-cli/tests/acp_handshake.rs @@ -0,0 +1,155 @@ +//! 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); +} 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. From 9be00742fc2ad8a720d75c7e3abf99645194cab1 Mon Sep 17 00:00:00 2001 From: ajianaz Date: Fri, 18 Sep 2026 17:00:16 +0700 Subject: [PATCH 2/5] fix: ACP session map survives across prompts; id traversal + poison hardening CodeCora review on the ACP PR caught a real architectural bug: the prompt arm wrapped a FRESH empty map per prompt (share_sessions used mem::replace), so every session vanished after its first turn and session/load could open a divergent second handle mid-turn. - The session map now lives for the whole server lifetime (Arc clone per prompt thread); the map lock is held for the duration of a turn, which serializes a busy session instead of allowing divergent appends. - Poisoning-tolerant locking: one panicking turn no longer bricks the ACP server. - sessionId validation (charset, length, no separators/parent refs) blocks path traversal via session/load before ids touch the filesystem; regression tests cover multi-session persistence and traversal rejection. --- CHANGELOG.md | 5 +- crates/tole-cli/src/acp.rs | 68 ++++++++++++++++++++------ crates/tole-cli/tests/acp_handshake.rs | 46 +++++++++++++++++ 3 files changed, 102 insertions(+), 17 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index e3078f4..66f71e3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -38,7 +38,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 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. + 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). - `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/crates/tole-cli/src/acp.rs b/crates/tole-cli/src/acp.rs index 2ec76bf..e8a67b5 100644 --- a/crates/tole-cli/src/acp.rs +++ b/crates/tole-cli/src/acp.rs @@ -241,7 +241,12 @@ pub fn run_acp( } }); - let mut sessions = Sessions::new(); + // 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 }; @@ -290,11 +295,14 @@ pub fn run_acp( .and_then(Value::as_str) .unwrap_or(".") .to_string(); - let session_id = params - .get("sessionId") - .and_then(Value::as_str) - .map(str::to_string) - .unwrap_or_else(new_session_id); + 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; + }; match open_session( &session_id, &cwd, @@ -305,7 +313,9 @@ pub fn run_acp( conn.clone(), ) { Ok(state) => { - sessions.map.insert(session_id.clone(), state); + lock_sessions(&sessions) + .map + .insert(session_id.clone(), state); reply(&conn, id, json!({ "sessionId": session_id })); } Err(e) => reply_error(&conn, id, &e.to_string()), @@ -316,10 +326,15 @@ pub fn run_acp( .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") { @@ -338,12 +353,12 @@ pub fn run_acp( // error) so this reader loop stays live for permission // requests while the turn runs. let conn = conn.clone(); - let sessions_ptr = share_sessions(&mut sessions); + 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_ptr, &session_id_clone, &prompt_clone, conn.clone()); + 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()), @@ -360,15 +375,15 @@ pub fn run_acp( Ok(()) } -// The sessions map must be shareable with the per-prompt threads; wrap it -// here so the method bodies stay readable. use std::sync::Arc as StdArc; type SharedSessions = StdArc>; -fn share_sessions(sessions: &mut Sessions) -> SharedSessions { - // A fresh Arc per ACP server instance is fine: the reader loop owns - // the map for the process lifetime and hands clones to turn threads. - StdArc::new(Mutex::new(std::mem::replace(sessions, Sessions::new()))) +/// 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) { @@ -384,6 +399,26 @@ fn reply_error(conn: &Conn, id: Option, message: &str) { ); } +/// 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) @@ -529,7 +564,8 @@ fn run_prompt( prompt: &str, conn: Conn, ) -> Result { - let mut sessions = sessions.lock().expect("sessions lock"); + // Poisoning-tolerant: see lock_sessions. + let mut sessions = lock_sessions(&sessions); let Some(state) = sessions.map.get_mut(session_id) else { anyhow::bail!("unknown session: {session_id}"); }; diff --git a/crates/tole-cli/tests/acp_handshake.rs b/crates/tole-cli/tests/acp_handshake.rs index 899bfa1..e388cf9 100644 --- a/crates/tole-cli/tests/acp_handshake.rs +++ b/crates/tole-cli/tests/acp_handshake.rs @@ -153,3 +153,49 @@ fn acp_load_returns_existing_session() { 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"); +} From 13f6506fc31c3b26c93c87546ee13ccd4a44b108 Mon Sep 17 00:00:00 2001 From: ajianaz Date: Mon, 28 Sep 2026 10:23:17 +0700 Subject: [PATCH 3/5] fix: ACP reader loop no longer deadlocks behind a running turn MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit CodeCora review deadlock finding: run_prompt held the session-map lock for the WHOLE turn — including the up-to-600s permission wait — and the reader loop's session/new/load arm takes the same lock. A client opening a session while a permission request was pending froze protocol routing until the timeout denied the tool. - SessionState storage/registry/first_prompt_done are now Arc-wrapped; the map lock is held only to take handles and reject a busy session; the running turn locks its OWN storage mutex. - Busy flag with a panic-safe Drop guard: concurrent turns on one session are refused ("session is busy"), and a panicking turn un-busies via Drop. - Live-verified: three consecutive prompts in one session all end_turn (previously the map swap broke even single-turn persistence); handshake integration tests stay green. --- CHANGELOG.md | 8 +++- crates/tole-cli/src/acp.rs | 89 ++++++++++++++++++++++++++------------ 2 files changed, 69 insertions(+), 28 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 66f71e3..76f9b85 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -41,7 +41,13 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 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). + 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/crates/tole-cli/src/acp.rs b/crates/tole-cli/src/acp.rs index e8a67b5..e2545ef 100644 --- a/crates/tole-cli/src/acp.rs +++ b/crates/tole-cli/src/acp.rs @@ -195,11 +195,24 @@ impl PromptFn for AcpPrompt { // --------------------------------------------------------------------------- struct SessionState { - storage: tole_core::storage::JsonlStorage, - registry: tole_core::tool::ToolRegistry, + // 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: bool, + 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 { @@ -517,11 +530,12 @@ fn open_session( .context("creating session")? }; Ok(SessionState { - storage, - registry: reg, + storage: StdArc::new(Mutex::new(storage)), + registry: StdArc::new(reg), system_prompt, memory, - first_prompt_done: loading, + first_prompt_done: StdArc::new(Mutex::new(loading)), + busy: StdArc::new(Mutex::new(false)), }) } @@ -564,21 +578,46 @@ fn run_prompt( prompt: &str, conn: Conn, ) -> Result { - // Poisoning-tolerant: see lock_sessions. - let mut sessions = lock_sessions(&sessions); - let Some(state) = sessions.map.get_mut(session_id) else { - anyhow::bail!("unknown session: {session_id}"); + // 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")] - if let Some(mem) = state.memory.clone() { - if !state.first_prompt_done { - 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}"); + { + 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}"); + } } } } @@ -591,25 +630,21 @@ fn run_prompt( (or the OPENAI_* equivalents)", )?; let mut provider = - tole_core::openai::OpenAiProvider::new(cfg).with_tool_specs(state.registry.specs()); - if let Some(sys) = state.system_prompt.as_deref() { + 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 outcome = tole_core::turn::run_turn( - &mut state.storage, - &mut provider, - &state.registry, - &effective, - )?; + 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) = state.memory.clone() { - let _ = tole_core::memory::remember_session(&mem, session_id, prompt, text); + if let Some(mem) = &memory { + let _ = tole_core::memory::remember_session(mem, session_id, prompt, text); } } - state.first_prompt_done = true; + *first_prompt_done.lock().expect("fpd lock") = true; let stop = match &outcome { tole_core::turn::TurnOutcome::Final { text } => { From db092930f326444390801baaa697b810ac1f7209 Mon Sep 17 00:00:00 2001 From: ajianaz Date: Mon, 28 Sep 2026 10:34:38 +0700 Subject: [PATCH 4/5] chore: retrigger CI for the D2 PR (previous push predated a stale CI state) From c918c06808447c80212d0a5428a809acafff1f40 Mon Sep 17 00:00:00 2001 From: ajianaz Date: Mon, 28 Sep 2026 11:28:34 +0700 Subject: [PATCH 5/5] =?UTF-8?q?fix:=20CodeCora=20round-2=20on=20the=20ACP?= =?UTF-8?q?=20host=20=E2=80=94=20plan-mode=20filter,=20busy=20load,=20erro?= =?UTF-8?q?r=20routing?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - plan_mode now REALLY filters: open_session registered every tool and the early return still handed back the full registry — under --yes that pre-authorized writes in a supposedly read-only session. Write/ Destructive registrations are now skipped entirely when plan_mode is active (read_file/cora_search/uteke_recall/job_poll remain). - session/load refuses a busy session: the old path opened a second JsonlStorage handle on the same JSONL while a turn was running, allowing divergent concurrent appends (the orphaned busy flag no longer guarded anything after the state was replaced). - client ERROR replies to session/request_permission are routed like results — an errored/cancelled permission now fails closed immediately instead of hanging the turn for the full 600s timeout. - also: duplicated too_many_arguments attribute removed (CI clippy). --- crates/tole-cli/src/acp.rs | 107 ++++++++++++++++++++++++++----------- 1 file changed, 76 insertions(+), 31 deletions(-) diff --git a/crates/tole-cli/src/acp.rs b/crates/tole-cli/src/acp.rs index d04ca11..6a14ffb 100644 --- a/crates/tole-cli/src/acp.rs +++ b/crates/tole-cli/src/acp.rs @@ -276,12 +276,19 @@ pub fn run_acp( .and_then(Value::as_str) .map(str::to_string) else { - // A response to one of OUR requests (permission). - if let (Some(id), Some(result)) = ( - msg.get("id").and_then(Value::as_u64), - msg.get("result").cloned(), - ) { - conn.route_response(id, result); + // 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; }; @@ -317,6 +324,21 @@ pub fn run_acp( 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, @@ -328,9 +350,24 @@ pub fn run_acp( conn.clone(), ) { Ok(state) => { - lock_sessions(&sessions) + // 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 - .insert(session_id.clone(), state); + .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()), @@ -446,7 +483,6 @@ fn new_session_id() -> String { /// the ACP editor (interactive — which is what unlocks Destructive tools /// with genuine human consent). #[allow(clippy::too_many_arguments)] -#[allow(clippy::too_many_arguments)] fn open_session( session_id: &str, cwd: &str, @@ -489,34 +525,43 @@ fn open_session( if tole_cli_binary_available("uteke") { reg.register(Box::new(UtekeRecallTool::new())) .map_err(anyhow::Error::msg)?; - reg.register(Box::new(UtekeDocumentTool::new(None))) + 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(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(JobPollTool::new(workspace_canon.clone()))) .map_err(anyhow::Error::msg)?; reg.register(Box::new(ReadFileTool::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)?; let storage = if loading { let dir = sessions_dir_for(cwd)?;