From ea8dd6838066b165875ad1123aa221f9ac05f48f Mon Sep 17 00:00:00 2001 From: Ivan Zatevakhin Date: Sat, 19 Sep 2026 16:34:53 +0100 Subject: [PATCH 1/3] fix: expose steering input events --- Cargo.lock | 1 + Cargo.toml | 1 + src/client.rs | 12 +++++++--- src/commands.rs | 3 +++ src/events.rs | 64 +++++++++++++++++++++++++++++++++++++++++++++++++ 5 files changed, 78 insertions(+), 3 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index c41b7b1..7f40d79 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1240,6 +1240,7 @@ dependencies = [ "tokio", "tokio-tungstenite", "url", + "uuid", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index 17351cb..4d4a2d0 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -22,3 +22,4 @@ serde_json = "1.0" tokio = { version = "1.0", features = ["full"] } tokio-tungstenite = "0.29" url = "2" +uuid = { version = "1", features = ["v4"] } diff --git a/src/client.rs b/src/client.rs index 438c91a..22ee613 100644 --- a/src/client.rs +++ b/src/client.rs @@ -360,10 +360,16 @@ async fn handle_inbound( let update = params.get("update").cloned().unwrap_or(Value::Null); if let Some(event) = events::compact_session_update(&session_id, &update) { events::observe_assistant_text(&mut *assistant_text.lock().await, &event); - if stream_events { - let _ = output::write_event(&event); - } + let _ = output::write_event(&event); } + } else if stream_events + && matches!( + method.as_str(), + "querymt/session/inputState" | "_querymt/session/inputState" + ) + && let Some(event) = events::compact_input_state(¶ms) + { + let _ = output::write_event(&event); } if let Some(id) = value.get("id").cloned() { diff --git a/src/commands.rs b/src/commands.rs index 395eecb..5cb5c8c 100644 --- a/src/commands.rs +++ b/src/commands.rs @@ -907,8 +907,10 @@ async fn submit_input( } else { run_id }; + let client_input_id = uuid::Uuid::new_v4().to_string(); let mut params = json!({ "session_id": session_id, + "client_input_id": client_input_id, "prompt": [{ "type": "text", "text": text }], }); if let Some(run_id) = run_id { @@ -923,6 +925,7 @@ async fn submit_input( .map_err(CliError::from_request)?; Ok(json!({ "sessionId": session_id, + "clientInputId": client_input_id, "result": extension_payload(&response), })) } diff --git a/src/events.rs b/src/events.rs index b6d1486..f0246da 100644 --- a/src/events.rs +++ b/src/events.rs @@ -12,6 +12,48 @@ pub fn observe_assistant_text(buffer: &mut String, event: &Value) { } } +pub fn compact_input_state(params: &Value) -> Option { + let session_id = params + .get("session_id") + .or_else(|| params.get("sessionId"))? + .clone(); + let input_id = params + .get("input_id") + .or_else(|| params.get("inputId"))? + .clone(); + let mut event = serde_json::Map::from_iter([ + ("type".to_string(), json!("input_state")), + ( + "version".to_string(), + params.get("version").cloned().unwrap_or_else(|| json!(1)), + ), + ("sessionId".to_string(), session_id), + ("inputId".to_string(), input_id), + ( + "delivery".to_string(), + params.get("delivery").cloned().unwrap_or(Value::Null), + ), + ( + "state".to_string(), + params.get("state").cloned().unwrap_or(Value::Null), + ), + ]); + for (output_key, input_keys) in [ + ("runId", ["run_id", "runId"]), + ("latencyMs", ["latency_ms", "latencyMs"]), + ] { + if let Some(value) = input_keys.iter().find_map(|key| params.get(*key)).cloned() { + event.insert(output_key.to_string(), value); + } + } + for key in ["position", "boundary", "reason"] { + if let Some(value) = params.get(key).cloned() { + event.insert(key.to_string(), value); + } + } + Some(Value::Object(event)) +} + pub fn compact_session_update(session_id: &Value, update: &Value) -> Option { let kind = update .get("sessionUpdate") @@ -84,6 +126,28 @@ fn first_string(value: &Value, keys: &[&str]) -> Option { mod tests { use super::*; + #[test] + fn compact_input_state_notification() { + let event = compact_input_state(&json!({ + "version": 1, + "session_id": "s1", + "input_id": "i1", + "delivery": "steer", + "state": "applied", + "run_id": "r1", + "boundary": "after_tools", + "latency_ms": 42 + })) + .unwrap(); + assert_eq!(event["type"], "input_state"); + assert_eq!(event["sessionId"], "s1"); + assert_eq!(event["inputId"], "i1"); + assert_eq!(event["state"], "applied"); + assert_eq!(event["runId"], "r1"); + assert_eq!(event["latencyMs"], 42); + assert!(event.get("position").is_none()); + } + #[test] fn compact_agent_text_chunk() { let event = compact_session_update( From 0f6459e16166b045839481ced9f289a79f593de1 Mon Sep 17 00:00:00 2001 From: Ivan Zatevakhin Date: Sun, 20 Sep 2026 01:44:55 +0100 Subject: [PATCH 2/3] feat: discard queued session input --- README.md | 1 + src/commands.rs | 42 ++++++++++++++++++++++++++++++++++++++++++ src/events.rs | 16 ++++++++++++++++ 3 files changed, 59 insertions(+) diff --git a/README.md b/README.md index 9cba5b7..a09fc9c 100644 --- a/README.md +++ b/README.md @@ -17,6 +17,7 @@ cargo run -- runtime SESSION_ID cargo run -- follow SESSION_ID cargo run -- steer SESSION_ID --run-id RUN "stop and summarize" cargo run -- queue SESSION_ID "next: run tests" +cargo run -- discard-queued SESSION_ID INPUT_ID cargo run -- inspect SESSION_ID --messages 20 cargo run -- set-mode SESSION_ID plan cargo run -- cancel SESSION_ID diff --git a/src/commands.rs b/src/commands.rs index 5cb5c8c..300cb35 100644 --- a/src/commands.rs +++ b/src/commands.rs @@ -157,6 +157,11 @@ pub enum Command { #[arg(trailing_var_arg = true, allow_hyphen_values = true)] text: Vec, }, + /// Remove an input that has not started from the session queue. + DiscardQueued { + session_id: String, + input_id: String, + }, /// Run multiple commands on one connection. Lines from stdin or remaining args. Exec { #[arg(trailing_var_arg = true, allow_hyphen_values = true)] @@ -283,6 +288,10 @@ pub async fn run( pretty, submit_input(client, "querymt/session/queue", session_id, None, text).await?, ), + Command::DiscardQueued { + session_id, + input_id, + } => write_ok(pretty, discard_queued(client, session_id, input_id).await?), Command::Exec { commands } => exec(client, commands, pretty).await, Command::Cancel { session_id } => write_ok(pretty, cancel(client, session_id).await?), Command::Close { session_id } => write_ok(pretty, close(client, session_id).await?), @@ -1003,6 +1012,29 @@ fn split_exec_args(line: &str) -> Vec { args } +async fn discard_queued( + client: &AcpClient, + session_id: String, + input_id: String, +) -> Result { + client.initialize().await.map_err(CliError::from_request)?; + let response = client + .extension( + "querymt/session/discardQueuedInput", + json!({ + "session_id": session_id, + "input_id": input_id, + }), + ) + .await + .map_err(CliError::from_request)?; + Ok(json!({ + "sessionId": session_id, + "inputId": input_id, + "result": extension_payload(&response), + })) +} + async fn cancel(client: &AcpClient, session_id: String) -> Result { client.initialize().await.map_err(CliError::from_request)?; client.cancel(session_id.clone()).map_err(CliError::rpc)?; @@ -1184,6 +1216,16 @@ mod tests { } } + #[test] + fn exec_line_parses_discard_queued_command() { + let command = parse_exec_line("discard-queued session-1 input-1").unwrap(); + assert!(matches!( + command, + Command::DiscardQueued { session_id, input_id } + if session_id == "session-1" && input_id == "input-1" + )); + } + #[test] fn split_exec_args_keeps_quoted_text() { assert_eq!( diff --git a/src/events.rs b/src/events.rs index f0246da..fb3ba0b 100644 --- a/src/events.rs +++ b/src/events.rs @@ -148,6 +148,22 @@ mod tests { assert!(event.get("position").is_none()); } + #[test] + fn compact_discarded_queue_notification() { + let event = compact_input_state(&json!({ + "version": 1, + "session_id": "s1", + "input_id": "i2", + "delivery": "queue", + "state": "discarded", + "reason": "removed_by_user" + })) + .unwrap(); + assert_eq!(event["delivery"], "queue"); + assert_eq!(event["state"], "discarded"); + assert_eq!(event["reason"], "removed_by_user"); + } + #[test] fn compact_agent_text_chunk() { let event = compact_session_update( From 35c83e8fb314dfed37d58128c7e00ed3cfb26186 Mon Sep 17 00:00:00 2001 From: Ivan Zatevakhin Date: Sun, 20 Sep 2026 16:14:37 +0100 Subject: [PATCH 3/3] docs: clarify queued input identifiers --- README.md | 13 ++++++++----- 1 file changed, 8 insertions(+), 5 deletions(-) diff --git a/README.md b/README.md index a09fc9c..ab20720 100644 --- a/README.md +++ b/README.md @@ -23,11 +23,14 @@ cargo run -- set-mode SESSION_ID plan cargo run -- cancel SESSION_ID ``` -`prompt` streams compact NDJSON (`text`, `tool`, `mode`, `plan`, `done`) and waits until -the session is idle, including queued turns. `done` includes assembled assistant -`text`. For an existing busy session, `--delivery auto` (default) steers if -possible, otherwise queues. `follow SESSION` streams updates until idle. `exec` -runs multiple commands on one connection. `--quiet` hides connection logs. +`prompt` streams compact NDJSON (`text`, `tool`, `mode`, `plan`, `input_state`, +`done`) and waits until the session is idle, including queued turns. `done` +includes assembled assistant `text`. For an existing busy session, +`--delivery auto` (default) steers if possible, otherwise queues. The +`INPUT_ID` passed to `discard-queued` is the `inputId` from an `input_state` +event, which matches the `clientInputId` returned when the input is submitted. +`follow SESSION` streams updates until idle. `exec` runs multiple commands on +one connection. `--quiet` hides connection logs. `--permission allow-once` is the default so coder tools can run. Elicitation requests are cancelled and emitted as NDJSON events during `prompt`.