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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"] }
14 changes: 9 additions & 5 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -17,16 +17,20 @@ 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
```

`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`.
12 changes: 9 additions & 3 deletions src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(&params)
{
let _ = output::write_event(&event);
}

if let Some(id) = value.get("id").cloned() {
Expand Down
45 changes: 45 additions & 0 deletions src/commands.rs
Original file line number Diff line number Diff line change
Expand Up @@ -157,6 +157,11 @@ pub enum Command {
#[arg(trailing_var_arg = true, allow_hyphen_values = true)]
text: Vec<String>,
},
/// 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)]
Expand Down Expand Up @@ -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?),
Expand Down Expand Up @@ -907,8 +916,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 {
Expand All @@ -923,6 +934,7 @@ async fn submit_input(
.map_err(CliError::from_request)?;
Ok(json!({
"sessionId": session_id,
"clientInputId": client_input_id,
"result": extension_payload(&response),
}))
}
Expand Down Expand Up @@ -1000,6 +1012,29 @@ fn split_exec_args(line: &str) -> Vec<String> {
args
}

async fn discard_queued(
client: &AcpClient,
session_id: String,
input_id: String,
) -> Result<Value, CliError> {
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<Value, CliError> {
client.initialize().await.map_err(CliError::from_request)?;
client.cancel(session_id.clone()).map_err(CliError::rpc)?;
Expand Down Expand Up @@ -1181,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!(
Expand Down
80 changes: 80 additions & 0 deletions src/events.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,48 @@ pub fn observe_assistant_text(buffer: &mut String, event: &Value) {
}
}

pub fn compact_input_state(params: &Value) -> Option<Value> {
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<Value> {
let kind = update
.get("sessionUpdate")
Expand Down Expand Up @@ -84,6 +126,44 @@ fn first_string(value: &Value, keys: &[&str]) -> Option<String> {
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_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(
Expand Down
Loading