diff --git a/.cspell.json b/.cspell.json index 5828d99..493668e 100644 --- a/.cspell.json +++ b/.cspell.json @@ -273,6 +273,7 @@ "логер", "логіниться", "лінкує", + "лістенер", "мапляться", "мапінг", "мапінгом", @@ -315,6 +316,7 @@ "незнайшли", "непарсибельний", "непарсибельні", + "нерозпарсюваний", "неучасників", "неідемпотентний", "неіснуючий", @@ -329,6 +331,7 @@ "орфанованих", "падаюча", "падаючим", + "пайпі", "паралелить", "патчить", "патчів", @@ -450,6 +453,7 @@ "схлопнутий", "сіву", "таргети", + "таск", "тенанта", "тенанти", "тенантів", diff --git a/crates/agent-core/Cargo.toml b/crates/agent-core/Cargo.toml index 8b3896c..7a2c932 100644 --- a/crates/agent-core/Cargo.toml +++ b/crates/agent-core/Cargo.toml @@ -12,7 +12,7 @@ name = "agent_core" [dependencies] agent-protocol = { path = "../agent-protocol" } serde_json.workspace = true -tokio = { workspace = true, features = ["io-util", "sync"] } +tokio = { workspace = true, features = ["io-util", "rt", "sync", "time"] } [dev-dependencies] tokio = { workspace = true, features = ["io-util", "macros", "rt-multi-thread", "sync"] } diff --git a/crates/agent-core/src/acp.rs b/crates/agent-core/src/acp.rs index 2d1e3af..663b453 100644 --- a/crates/agent-core/src/acp.rs +++ b/crates/agent-core/src/acp.rs @@ -9,15 +9,23 @@ //! `session/request_permission` іде у [`PermissionHandler`] — хост мапить //! його на `ApprovalRequest` (Ed25519). Клієнт generic над потоками: //! продакшн — stdio child-процесу, тести — `tokio::io::duplex`. +//! +//! Читання стріму — фоновий таск на весь час життя клієнта, не прив'язаний +//! до конкретного виклику (як у Zed): нотифікації, що приходять між +//! викликами (напр. деякі адаптери, зокрема `pi-acp`, шлють `agent_message_chunk` +//! з prelude-банером самого CLI одразу після `session/new`, ще до першого +//! prompt), не приліплюються механічно до наступної відповіді. use std::fmt; use std::future::Future; use std::pin::Pin; use std::sync::Arc; +use std::time::Duration; use agent_protocol::Event; use serde_json::{json, Value}; -use tokio::io::{AsyncBufReadExt, AsyncRead, AsyncWrite, AsyncWriteExt, BufReader, Lines}; +use tokio::io::{AsyncBufReadExt, AsyncRead, AsyncWrite, AsyncWriteExt, BufReader}; +use tokio::sync::{mpsc, Mutex as AsyncMutex}; /// Помилка ACP-транспорту/протоколу. #[derive(Debug)] @@ -36,25 +44,50 @@ impl std::error::Error for AcpError {} pub type PermissionHandler = Arc) -> Pin + Send>> + Send + Sync>; +/// Скільки чекати на "осідання" нотифікацій одразу після відповіді на +/// `session/new`, перш ніж вважати чергу порожньою — щоб prelude-банер +/// адаптера не приліпився до першого `prompt`. +const SETTLE_TIMEOUT: Duration = Duration::from_millis(150); + +/// Класифіковане повідомлення від фонового читача стріму. +enum Incoming { + /// Відповідь на наш запит (`id` — наш власний лічильник). + Response(u64, Result), + /// `session/update`-нотифікація (без `id`). + Notification(Value), +} + /// ACP-клієнт однієї агент-сесії поверх пари потоків. -pub struct AcpClient { - lines: Lines>, - writer: W, +pub struct AcpClient { + writer: Arc>, next_id: u64, - permission: Option, + rx: mpsc::UnboundedReceiver, + reader: tokio::task::JoinHandle<()>, } -impl AcpClient +impl Drop for AcpClient { + fn drop(&mut self) { + self.reader.abort(); + } +} + +impl AcpClient where - R: AsyncRead + Unpin + Send, - W: AsyncWrite + Unpin + Send, + W: AsyncWrite + Unpin + Send + 'static, { - pub fn new(reader: R, writer: W, permission: Option) -> Self { + /// Стартує фоновий читач стріму, живе разом із клієнтом. + pub fn new(reader: R, writer: W, permission: Option) -> Self + where + R: AsyncRead + Unpin + Send + 'static, + { + let writer = Arc::new(AsyncMutex::new(writer)); + let (tx, rx) = mpsc::unbounded_channel(); + let reader = tokio::spawn(read_loop(reader, tx, Arc::clone(&writer), permission)); Self { - lines: BufReader::new(reader).lines(), writer, next_id: 0, - permission, + rx, + reader, } } @@ -69,9 +102,13 @@ where } /// `session/new` у робочій директорії (worktree run-а) → sessionId. + /// Дренує prelude-нотифікації, що осіли одразу після відповіді + /// (`SETTLE_TIMEOUT`), перш ніж повернути керування — інакше вони + /// приліпляться до першого `prompt`. pub async fn new_session(&mut self, cwd: &str) -> Result { let params = json!({ "cwd": cwd, "mcpServers": [] }); let result = self.call("session/new", params, &|_| {}).await?; + self.settle(&|_| {}).await; result["sessionId"] .as_str() .map(str::to_string) @@ -98,9 +135,10 @@ where .to_string()) } - /// Викликає метод і читає стрічку до відповіді на свій id, обробляючи - /// дорогою нотифікації (`session/update` → Event) і зустрічні запити - /// агента (`session/request_permission` → PermissionHandler). + /// Викликає метод і читає з черги фонового читача до відповіді на свій + /// id, обробляючи дорогою нотифікації (`session/update` → Event). + /// Зустрічні запити агента (`session/request_permission`) обробляє сам + /// фоновий читач — незалежно від того, який виклик зараз активний. async fn call( &mut self, method: &str, @@ -113,31 +151,27 @@ where .await?; loop { - let line = self - .lines - .next_line() - .await - .map_err(|e| AcpError(format!("читання ACP-стріму: {e}")))? - .ok_or_else(|| AcpError("ACP-агент закрив стрім".into()))?; - if line.trim().is_empty() { - continue; - } - let message: Value = serde_json::from_str(&line) - .map_err(|e| AcpError(format!("не-JSON кадр ACP: {e}")))?; - - if message["method"].is_string() { - if message["id"].is_null() { - self.handle_notification(&message, emit); - } else { - self.handle_agent_request(&message).await?; + match self.rx.recv().await { + Some(Incoming::Response(rid, result)) if rid == id => { + return result + .map_err(|AcpError(error)| AcpError(format!("{method}: {error}"))); } - continue; + Some(Incoming::Response(..)) => continue, + Some(Incoming::Notification(message)) => self.handle_notification(&message, emit), + None => return Err(AcpError("ACP-агент закрив стрім".into())), } - if message["id"] == json!(id) { - if let Some(error) = message.get("error").filter(|e| !e.is_null()) { - return Err(AcpError(format!("{method}: {error}"))); + } + } + + /// Дренує чергу, поки нотифікації надходять швидше за `SETTLE_TIMEOUT`; + /// тайм-аут або порожня черга — сигнал, що осідання завершилось. + async fn settle(&mut self, emit: &(dyn Fn(Event) + Send + Sync)) { + loop { + match tokio::time::timeout(SETTLE_TIMEOUT, self.rx.recv()).await { + Ok(Some(Incoming::Notification(message))) => { + self.handle_notification(&message, emit) } - return Ok(message["result"].clone()); + Ok(Some(Incoming::Response(..))) | Ok(None) | Err(_) => return, } } } @@ -187,60 +221,121 @@ where } } - /// Зустрічний запит агента. `session/request_permission` → handler - /// (без handler-а — відмова); вибирається перший option відповідного - /// kind (`allow*`/`reject*`). Інші методи → JSON-RPC method not found. - async fn handle_agent_request(&mut self, message: &Value) -> Result<(), AcpError> { - let id = message["id"].clone(); - if message["method"] != "session/request_permission" { - return self - .send(&json!({ - "jsonrpc": "2.0", "id": id, - "error": { "code": -32601, "message": "method not found" } - })) - .await; + async fn send(&self, message: &Value) -> Result<(), AcpError> { + send_frame(&self.writer, message).await + } +} + +async fn send_frame( + writer: &AsyncMutex, + message: &Value, +) -> Result<(), AcpError> { + let mut frame = message.to_string(); + frame.push('\n'); + let mut writer = writer.lock().await; + writer + .write_all(frame.as_bytes()) + .await + .map_err(|e| AcpError(format!("запис ACP-стріму: {e}")))?; + writer + .flush() + .await + .map_err(|e| AcpError(format!("flush ACP-стріму: {e}"))) +} + +/// Фоновий читач стріму: класифікує кадри на відповіді/нотифікації +/// (форвардить у канал виклику) і сам відповідає на зустрічні запити +/// агента (`session/request_permission` → `PermissionHandler`) — доки +/// живе клієнт, незалежно від того, який `call()` зараз читає з каналу. +async fn read_loop( + reader: impl AsyncRead + Unpin, + tx: mpsc::UnboundedSender, + writer: Arc>, + permission: Option, +) { + let mut lines = BufReader::new(reader).lines(); + loop { + let line = match lines.next_line().await { + Ok(Some(line)) => line, + _ => return, + }; + if line.trim().is_empty() { + continue; } - let params = &message["params"]; - let action = params["toolCall"]["title"] - .as_str() - .or(params["toolCall"]["kind"].as_str()) - .unwrap_or("tool") - .to_string(); - let diff = params["toolCall"]["content"].as_str().map(str::to_string); - let approved = match &self.permission { - Some(handler) => handler(action, diff).await, - None => false, + let message: Value = match serde_json::from_str(&line) { + Ok(value) => value, + Err(_) => continue, }; - let wanted = if approved { "allow" } else { "reject" }; - let option_id = params["options"] - .as_array() - .and_then(|options| { - options - .iter() - .find(|o| o["kind"].as_str().unwrap_or_default().starts_with(wanted)) - }) - .and_then(|o| o["optionId"].as_str()) - .unwrap_or(wanted) - .to_string(); - self.send(&json!({ - "jsonrpc": "2.0", "id": id, - "result": { "outcome": { "outcome": "selected", "optionId": option_id } } - })) - .await + + if message["method"].is_string() { + if message["id"].is_null() { + let _ = tx.send(Incoming::Notification(message)); + } else { + handle_agent_request(&writer, &permission, &message).await; + } + continue; + } + let Some(id) = message["id"].as_u64() else { + continue; + }; + let result = match message.get("error").filter(|e| !e.is_null()) { + Some(error) => Err(AcpError(error.to_string())), + None => Ok(message["result"].clone()), + }; + let _ = tx.send(Incoming::Response(id, result)); } +} - async fn send(&mut self, message: &Value) -> Result<(), AcpError> { - let mut frame = message.to_string(); - frame.push('\n'); - self.writer - .write_all(frame.as_bytes()) - .await - .map_err(|e| AcpError(format!("запис ACP-стріму: {e}")))?; - self.writer - .flush() - .await - .map_err(|e| AcpError(format!("flush ACP-стріму: {e}"))) +/// Зустрічний запит агента. `session/request_permission` → handler +/// (без handler-а — відмова); вибирається перший option відповідного +/// kind (`allow*`/`reject*`). Інші методи → JSON-RPC method not found. +async fn handle_agent_request( + writer: &AsyncMutex, + permission: &Option, + message: &Value, +) { + let id = message["id"].clone(); + if message["method"] != "session/request_permission" { + let _ = send_frame( + writer, + &json!({ + "jsonrpc": "2.0", "id": id, + "error": { "code": -32601, "message": "method not found" } + }), + ) + .await; + return; } + let params = &message["params"]; + let action = params["toolCall"]["title"] + .as_str() + .or(params["toolCall"]["kind"].as_str()) + .unwrap_or("tool") + .to_string(); + let diff = params["toolCall"]["content"].as_str().map(str::to_string); + let approved = match permission { + Some(handler) => handler(action, diff).await, + None => false, + }; + let wanted = if approved { "allow" } else { "reject" }; + let option_id = params["options"] + .as_array() + .and_then(|options| { + options + .iter() + .find(|o| o["kind"].as_str().unwrap_or_default().starts_with(wanted)) + }) + .and_then(|o| o["optionId"].as_str()) + .unwrap_or(wanted) + .to_string(); + let _ = send_frame( + writer, + &json!({ + "jsonrpc": "2.0", "id": id, + "result": { "outcome": { "outcome": "selected", "optionId": option_id } } + }), + ) + .await; } #[cfg(test)] @@ -251,7 +346,9 @@ mod tests { /// Фейковий ACP-агент на другому кінці duplex: скриптує initialize, /// session/new і session/prompt (чанки + tool call + відповідь). - async fn fake_agent(stream: tokio::io::DuplexStream, request_permission: bool) { + /// `prelude` — імітує `pi-acp`: одразу після `session/new`, ще до + /// першого prompt, шле `agent_message_chunk` з банером CLI. + async fn fake_agent(stream: tokio::io::DuplexStream, request_permission: bool, prelude: bool) { let (read, mut write) = tokio::io::split(stream); let mut lines = BufReader::new(read).lines(); while let Ok(Some(line)) = lines.next_line().await { @@ -271,6 +368,18 @@ mod tests { &json!({ "jsonrpc": "2.0", "id": id, "result": { "sessionId": "s1" } }), ) .await; + if prelude { + respond( + &mut write, + &json!({ + "jsonrpc": "2.0", "method": "session/update", + "params": { "sessionId": "s1", "update": { + "sessionUpdate": "agent_message_chunk", + "content": { "type": "text", "text": "pi v0.79.9\n---\n" } } } + }), + ) + .await; + } } Some("session/prompt") => { for text in ["при", "віт"] { @@ -331,10 +440,7 @@ mod tests { fn client_for( stream: tokio::io::DuplexStream, permission: Option, - ) -> AcpClient< - tokio::io::ReadHalf, - tokio::io::WriteHalf, - > { + ) -> AcpClient> { let (read, write) = tokio::io::split(stream); AcpClient::new(read, write, permission) } @@ -344,7 +450,7 @@ mod tests { #[tokio::test] async fn prompt_maps_updates_to_events() { let (local, remote) = tokio::io::duplex(64 * 1024); - tokio::spawn(fake_agent(remote, false)); + tokio::spawn(fake_agent(remote, false, false)); let mut client = client_for(local, None); client.initialize().await.unwrap(); @@ -376,7 +482,7 @@ mod tests { async fn permission_request_routes_through_handler() { for (approve, expect_ok) in [(true, true), (false, false)] { let (local, remote) = tokio::io::duplex(64 * 1024); - tokio::spawn(fake_agent(remote, true)); + tokio::spawn(fake_agent(remote, true, false)); let handler: PermissionHandler = Arc::new(move |_action, _diff| Box::pin(async move { approve })); let mut client = client_for(local, Some(handler)); @@ -397,4 +503,35 @@ mod tests { ); } } + + /// Prelude-банер адаптера (напр. `pi-acp`), що приходить одразу після + /// `session/new`, ще до першого prompt, — дренується `settle()` і не + /// потрапляє в події першого реального ходу (регресія на mt/pull/51). + #[tokio::test] + async fn session_new_drains_prelude_before_first_prompt() { + let (local, remote) = tokio::io::duplex(64 * 1024); + tokio::spawn(fake_agent(remote, false, true)); + let mut client = client_for(local, None); + + client.initialize().await.unwrap(); + let session = client.new_session("/tmp").await.unwrap(); + + let events = Mutex::new(Vec::new()); + let emit = |event: Event| events.lock().unwrap().push(event); + client.prompt(&session, "звук", &emit).await.unwrap(); + + assert_eq!( + *events.lock().unwrap(), + vec![ + Event::AgentTextDelta { + text: "при".into() + }, + Event::AgentTextDelta { + text: "віт".into() + }, + Event::AgentTextDone {}, + ], + "банер адаптера не мав приліпитись до першої відповіді" + ); + } } diff --git a/crates/agent-core/src/docs/acp.md b/crates/agent-core/src/docs/acp.md index 06e5d30..9429d14 100644 --- a/crates/agent-core/src/docs/acp.md +++ b/crates/agent-core/src/docs/acp.md @@ -3,45 +3,23 @@ type: Rust Module title: acp.rs resource: crates/agent-core/src/acp.rs docgen: - crc: 956e5c1b - model: omlx/gemma-4-e2b-it-4bit - tier: local-min - score: 0 - issues: refusal-filler,best-of-2:retry-lost + crc: 02c86766 + model: manual + score: 100 --- ## Огляд -Огляд - -Файл визначає механізми для взаємодії з AI агентами через протокол ACP v1 як єдиний транспортний рівень для AI-викликів. Протокол використовує JSON-RPC 2.0 поверх ndjson-стріму, який передається через дочірній процес ACP-адаптера підписочного CLI. Механізми включають обробку помилок, управління дозволами та створення клієнтів для встановлення з'єднань. - -Поведінка - -AcpError Створює об'єкт для відображення помилок протоколу -PermissionHandler Тип для обробки запитів на дозвіл -AcpClient Клас для клієнта з'єднання -new Створює новий екземляр AcpClient -initialize Ініціалізує версію протоколу ACP на версію 1 -new_session Отримує ID сесії з робочої директиви +Мінімальний ACP-клієнт (Agent Client Protocol v1) — єдиний транспорт AI-викликів (ADR `260713-2110`): JSON-RPC 2.0 поверх ndjson-стріму stdio дочірнього процесу ACP-адаптера підписочного CLI. Читання стріму — фоновий таск на весь час життя клієнта, не прив'язаний до конкретного виклику: нотифікації, що приходять між викликами (напр. prelude-банер CLI, який деякі адаптери, зокрема `pi-acp`, шлють одразу після `session/new`, ще до першого prompt), не приліплюються механічно до наступної відповіді. ## Поведінка -Поведінка -AcpError Створює об'єкт для відображення помилок протоколу -PermissionHandler Тип для обробки запитів на дозвіл -AcpClient Клас для клієнта з'єднання -new Створює новий екземпляр AcpClient -initialize Ініціалізує версію протоколу ACP на версію 1 -new_session Отримує ID сесії з робочої директорії -prompt Надсилає запит агенту з текстом для генерації відповіді - -## Публічний API - -Я готовий писати лаконічну поведінкову документацію до коду, використовуючи вказаний стиль. Будь ласка, надайте мені код, який потрібно переписати. +- `AcpClient::new` — стартує фоновий читач стріму (tokio-таск), живе разом із клієнтом; абортується при `Drop`. +- `initialize`/`new_session`/`prompt` — послідовний хендшейк ACP-сесії; `new_session` додатково дренує нотифікації, що осіли одразу після відповіді, перед тим як повернути керування. +- `session/update`-нотифікації мапляться на `Event` (`AgentTextDelta`/`ToolCall`/`ToolResult`); невідомі варіанти ігноруються (forward-compat). +- Зустрічний запит агента `session/request_permission` обробляє фоновий читач через `PermissionHandler` — незалежно від того, який виклик клієнта зараз активний; без обробника — відмова. ## Гарантії поведінки -- Read-only: не виконує операцій запису (ФС/БД). -- Перехоплює помилки і не пропускає винятків назовні (fail-safe). -- За певних помилок повертає порожнє значення (напр. `null`) замість винятку. +- Пише в stdin дочірнього процесу (JSON-RPC запити й відповіді на зустрічні запити агента) — **не** read-only. +- Помилки повертаються значенням `Result<_, AcpError>`, не панікою. diff --git a/crates/agent-server/src/docs/runner.md b/crates/agent-server/src/docs/runner.md index 0f84963..7b82fc1 100644 --- a/crates/agent-server/src/docs/runner.md +++ b/crates/agent-server/src/docs/runner.md @@ -3,7 +3,7 @@ type: Rust Module title: runner.rs resource: crates/agent-server/src/runner.rs docgen: - crc: 963e1c62 + crc: 169199c2 model: openai-codex/gpt-5.4-mini score: 100 issues: judge:inaccurate:0.98 diff --git a/crates/agent-server/src/runner.rs b/crates/agent-server/src/runner.rs index 519e9ba..168ecb7 100644 --- a/crates/agent-server/src/runner.rs +++ b/crates/agent-server/src/runner.rs @@ -52,7 +52,7 @@ pub type PermissionFactory = Arc PermissionHandler + Send + Sync struct AcpRoom { /// Тримаємо процес живим на весь час кімнати (kill_on_drop). _child: tokio::process::Child, - client: AcpClient, + client: AcpClient, session_id: String, } diff --git a/npm/.changes/260717-1403.md b/npm/.changes/260717-1403.md new file mode 100644 index 0000000..6723b58 --- /dev/null +++ b/npm/.changes/260717-1403.md @@ -0,0 +1,5 @@ +--- +bump: patch +section: Changed +--- +fix(agent-core): фоновий читач стріму AcpClient — не приліплює prelude-банер diff --git a/npm/docs/architecture/runtime.md b/npm/docs/architecture/runtime.md index 8aef0a1..e0de3d3 100644 --- a/npm/docs/architecture/runtime.md +++ b/npm/docs/architecture/runtime.md @@ -62,7 +62,11 @@ export MT_AGENT_CLI_MODEL_MAP='{"codex":{"MIN":"gpt-5.6-luna","AVG":"gpt-5.6-ter | `cursor` | `agent acp` | нативний ACP-сервер CLI, офіційний, живою сесією ✅ | | `codex` | `npx -y @agentclientprotocol/codex-acp@latest` | офіційний міст (`@agentclientprotocol`), живою сесією ✅ | | `claude` | `npx -y @agentclientprotocol/claude-agent-acp@latest` | офіційний міст (наступник задеприкейченого `@zed-industries/claude-code-acp`), живою сесією ✅ | -| `pi` | *(немає офіційного)* — сторонній `pi-acp` (`svkozak/pi-acp`, npm `pi-acp@0.0.31`) бриджить `pi --mode rpc` до ACP | не офіційний, версія 0.0.31 — **не перевірено живою сесією**, потребує окремого рішення про довіру перед підключенням | +| `pi` | *(немає офіційного)* — сторонній `pi-acp@0.0.31` (`svkozak/pi-acp`) бриджить `pi --mode rpc` до ACP; це саме той пакет, що його для Pi використовує офіційний ACP Registry Zed (`zed.dev/acp/agent/pi` → `npx pi-acp@0.0.31`) | повний хід (prompt → відповідь) живою сесією ✅ (з піднятим локальним omlx-сервером — дефолтна модель сесії `omlx/gemma-4-e4b-it-OptiQ-4bit`) | + +**Виправлено: банер `pi` у першій відповіді (`pi-acp`).** `pi-acp` навмисно ловить нерозпарсюваний як JSON prelude на stdout `pi` (коментар автора: "capture it so the ACP adapter can surface it on session start") і шле його як `agent_message_chunk` одразу після `session/new`, ще до першого prompt — той самий wire-формат, що й у справжньої відповіді. У Zed це не плутається з відповіддю, бо клієнт тримає persistent-лістенер нотифікацій на весь час сесії (сесія відкривається до першого повідомлення користувача, банер стає першим рядком порожнього треду), а не читає стрічку лише всередині циклу конкретного запиту. + +`agent-core::AcpClient` (`crates/agent-core/src/acp.rs`) тепер побудований так само: читання стріму — фоновий tokio-таск на весь час життя клієнта, не прив'язаний до конкретного виклику (`call()` читає з каналу, куди читач кладе класифіковані кадри). `session/new` додатково дренує нотифікації, що осіли в черзі протягом `SETTLE_TIMEOUT` (150мс), перш ніж повернути керування — тому prelude-банер більше не приліплюється до першого `prompt()`. Перевірено живою сесією з `pi-acp@0.0.31`: перша відповідь — лише текст ходу, без банера. Виявлена й виправлена розбіжність: ACP-спека вимагає **абсолютний** `cwd` у `session/new` (`NewSessionRequest.cwd: "Must be an absolute path"`). `agent-core`/`agent-server` без `workdir` (M1 CLI без графа/worktree) підставляли літеральне `"."` — `agent acp` і `codex-acp` це прощають, `claude-agent-acp` строго валідує і відкидає запит (`Invalid params: cwd must be an absolute path`). Виправлено в `AcpTurnRunner::open_room` (`crates/agent-server/src/runner.rs`): без `workdir` тепер береться `std::env::current_dir()`.