diff --git a/crates/agent/examples/probe.rs b/crates/agent/examples/probe.rs index 252ee317..308ab284 100644 --- a/crates/agent/examples/probe.rs +++ b/crates/agent/examples/probe.rs @@ -363,6 +363,7 @@ async fn run_probe( AgentEvent::UserInputRequested { request_id, questions, + .. } => { let answers = questions .iter() diff --git a/crates/agent/src/claude.rs b/crates/agent/src/claude.rs index 3472674e..91bf4e5a 100644 --- a/crates/agent/src/claude.rs +++ b/crates/agent/src/claude.rs @@ -42,7 +42,8 @@ use crate::{ InteractionMode, ItemContent, ItemStatus, LaunchEnv, ModelSpec, OptionDescriptor, OptionSelection, PlanStep, PlanStepStatus, ProviderCommand, ProviderCommandKind, ProviderKind, ResumeCursor, RewindMode, SessionCommand, SessionHandle, SessionOptions, ThreadItem, - TokenUsage, TurnStatus, UserInputOption, UserInputQuestion, selection_bool, selection_str, + TokenUsage, TurnStatus, UserInputDelivery, UserInputOption, UserInputQuestion, selection_bool, + selection_str, }; /// Denial returned to `ExitPlanMode` after the client captures the plan. @@ -2715,6 +2716,7 @@ impl Mapper { return vec![AgentEvent::UserInputRequested { request_id, questions, + delivery: UserInputDelivery::Blocking, }]; } @@ -4984,6 +4986,7 @@ mod tests { AgentEvent::UserInputRequested { request_id, questions, + .. } => { assert_eq!(request_id, "ctrl-9"); assert_eq!(questions.len(), 2); diff --git a/crates/agent/src/codex.rs b/crates/agent/src/codex.rs index 59c7a2e6..403d443f 100644 --- a/crates/agent/src/codex.rs +++ b/crates/agent/src/codex.rs @@ -19,7 +19,8 @@ use crate::{ ItemContent, ItemStatus, LaunchEnv, ModelSpec, OptionDescriptor, OptionSelection, PlanStep, PlanStepStatus, ProviderCommand, ProviderCommandKind, ProviderKind, ResumeCursor, SelectOption, SessionCommand, SessionHandle, SessionOptions, ThreadItem, TokenUsage, TurnOptions, TurnStatus, - UserInputOption, UserInputQuestion, file_changes_from_unified_diff, selection_str, + UserInputDelivery, UserInputOption, UserInputQuestion, file_changes_from_unified_diff, + selection_str, }; mod developer_instructions; @@ -409,6 +410,11 @@ struct Actor { /// Pending `mcpServer/elicitation/request`s: canonical request_id → the /// JSON-RPC id and field typing needed to rebuild a typed response. elicitations: HashMap, + /// Open `request_user_input_async` questions of the running turn, under + /// the request id of the latest asking message. The model keeps working; + /// answers go back as a `turn/steer` carrying Codex's reply envelope, and + /// the turn's end discards whatever is still unanswered. + async_questions: Option<(String, Vec)>, items: HashMap, subagents: HashMap, /// Stable parent capsule for each provider-native child thread. Codex 0.150 @@ -485,6 +491,7 @@ async fn run_actor( approvals: HashMap::new(), user_inputs: HashMap::new(), elicitations: HashMap::new(), + async_questions: None, items: HashMap::new(), subagents: HashMap::new(), subagent_parent_by_thread: HashMap::new(), @@ -638,6 +645,14 @@ impl SessionActor for Actor { &json!({ "id": pending.rpc_id, "result": result }), ) .map_err(|e| e.to_string())?; + } else if self + .async_questions + .as_ref() + .is_some_and(|(id, _)| *id == request_id) + { + let (_, questions) = self.async_questions.take().unwrap_or_default(); + self.steer_async_question_reply(&request_id, &questions, &answers) + .await?; } else { self.events .emit(AgentEvent::Warning { @@ -1340,6 +1355,7 @@ impl Actor { .emit(AgentEvent::UserInputRequested { request_id: key, questions, + delivery: UserInputDelivery::Blocking, }) .await; return; @@ -1451,6 +1467,7 @@ impl Actor { .emit(AgentEvent::UserInputRequested { request_id, questions, + delivery: UserInputDelivery::Blocking, }) .await; } @@ -1510,6 +1527,7 @@ impl Actor { _ => TurnStatus::Completed, }; self.active_turn = None; + self.async_questions = None; let usage = self.usage_by_turn.remove(&id); self.events .emit(AgentEvent::TurnCompleted { @@ -1624,6 +1642,11 @@ impl Actor { }; self.events.emit(event).await; } + if method == "item/completed" + && let Some(item) = item_value + { + self.request_async_questions(item).await; + } } "turn/plan/updated" => { let turn_id = notification_turn_id(params, &self.active_turn); @@ -1922,6 +1945,82 @@ impl Actor { self.events.emit(AgentEvent::ItemUpdated(child)).await; } + /// Deliver `request_user_input_async` answers into the running turn. The + /// model reads them like any steer, so the reply is a `turn/steer` whose + /// text is Codex's own question-reply envelope; the timeline shows the + /// same plain rendering the Codex TUI uses for that envelope. + async fn steer_async_question_reply( + &mut self, + request_id: &str, + questions: &[UserInputQuestion], + answers: &serde_json::Map, + ) -> Result<(), String> { + let Some((wire, display)) = async_question_reply(questions, answers) else { + return Ok(()); + }; + let Some(turn_id) = self.active_turn.clone() else { + self.events + .emit(AgentEvent::Warning { + message: "cannot answer: the Codex turn that asked has ended".into(), + }) + .await; + return Ok(()); + }; + let steer_id = format!("codex-question-reply-{request_id}"); + self.events + .emit(AgentEvent::SteerRequested { + request_id: steer_id.clone(), + text: display, + attachments: Vec::new(), + }) + .await; + let thread_id = self.thread_id.clone(); + self.request( + "turn/steer", + json!({ + "threadId": thread_id, + "expectedTurnId": turn_id, + "input": user_input(&wire, &[]), + }), + PendingRequest::Steer { + request_id: steer_id, + text: wire, + }, + ) + } + + /// A completed `agentMessage` carrying `questions` is the + /// `request_user_input_async` tool's output: the message text already + /// lists the questions, and the turn keeps running. Open questions of the + /// turn accumulate under one request so a later ask does not hide an + /// earlier unanswered one. + async fn request_async_questions(&mut self, item: &Value) { + if item.get("type").and_then(Value::as_str) != Some("agentMessage") { + return; + } + let Some(message_id) = item.get("id").and_then(Value::as_str) else { + return; + }; + let questions = parse_codex_async_questions(message_id, item); + if questions.is_empty() { + return; + } + let mut open = self + .async_questions + .take() + .map(|(_, questions)| questions) + .unwrap_or_default(); + open.extend(questions); + self.async_questions = Some((message_id.to_owned(), open.clone())); + self.events + .emit(AgentEvent::UserInputRequested { + request_id: message_id.to_owned(), + questions: open, + delivery: UserInputDelivery::Async, + }) + .await; + } + /// Settle every outstanding native user-input request and MCP elicitation /// on teardown, replying with the protocol's empty/cancel outcome. async fn settle_pending_user_inputs_on_shutdown(&mut self) { @@ -2298,6 +2397,95 @@ fn parse_codex_user_input(params: &Value) -> Vec { .collect() } +/// Map a completed `agentMessage`'s `questions` (the `request_user_input_async` +/// tool's `{title, options?}` list) into canonical questions. The id is the +/// desktop app's `JSON.stringify([tool, message id, index])`, which is what +/// the reply envelope must echo back; the tool has no header. Questions with +/// an empty title and options with an empty label are dropped. +const ASYNC_QUESTION_TOOL: &str = "request_user_input_async"; + +fn has_async_questions(item: &Value) -> bool { + item.get("questions") + .and_then(Value::as_array) + .is_some_and(|questions| !questions.is_empty()) +} + +fn parse_codex_async_questions(message_id: &str, item: &Value) -> Vec { + let questions = match item.get("questions").and_then(Value::as_array) { + Some(questions) => questions, + None => return Vec::new(), + }; + questions + .iter() + .enumerate() + .filter_map(|(index, q)| { + let title = q + .get("title") + .and_then(Value::as_str) + .map(str::trim) + .filter(|t| !t.is_empty())?; + let options = strings(q.get("options")) + .into_iter() + .map(|label| label.trim().to_owned()) + .filter(|label| !label.is_empty()) + .map(|label| UserInputOption { + label, + description: String::new(), + }) + .collect(); + Some(UserInputQuestion { + id: json!([ASYNC_QUESTION_TOOL, message_id, index]).to_string(), + header: String::new(), + question: title.to_owned(), + options, + multi_select: false, + prefill: None, + }) + }) + .collect() +} + +/// Codex's `request_user_input_async` reply: the `` +/// envelope the desktop app and TUI send, listing each answered question with +/// its `questionItemId` so the model correlates the answer. Returns the wire +/// text and the plain `> question` / answer rendering the timeline shows, or +/// `None` when no question was answered. Question text is flattened to one +/// line and capped at 512 characters, as the TUI does. +fn async_question_reply( + questions: &[UserInputQuestion], + answers: &serde_json::Map, +) -> Option<(String, String)> { + let mut replies = Vec::new(); + let mut display = Vec::new(); + for question in questions { + let answer = strings(answers.get(&question.id)).join(", "); + let answer = answer.trim(); + if answer.is_empty() { + continue; + } + let title = question + .question + .chars() + .take(512) + .collect::() + .replace(['\n', '\r'], " "); + replies.push(json!({ + "answer": answer, + "question": title, + "questionItemId": question.id, + })); + display.push(format!("> {title}\n\n{answer}")); + } + if replies.is_empty() { + return None; + } + let wire = format!( + "\n{}\n", + Value::Array(replies) + ); + Some((wire, display.join("\n\n"))) +} + fn map_item(item: &Value) -> Option { let id = item.get("id").and_then(Value::as_str)?.to_owned(); let provider_kind = item @@ -2313,6 +2501,14 @@ fn map_item(item: &Value) -> Option { return None; } let content = match provider_kind { + // `request_user_input_async` output: the text only lists the questions + // the input panel shows, so the transcript records the tool call. + "agentMessage" if has_async_questions(item) => ItemContent::ToolCall { + name: ASYNC_QUESTION_TOOL.into(), + input: json!({ "questions": item.get("questions").cloned().unwrap_or_default() }), + output: None, + status: ItemStatus::Completed, + }, "agentMessage" => ItemContent::AssistantMessage { text: string_field(item, "text"), }, @@ -2648,6 +2844,7 @@ mod tests { approvals: HashMap::new(), user_inputs: HashMap::new(), elicitations: HashMap::new(), + async_questions: None, items: HashMap::new(), subagents: HashMap::new(), subagent_parent_by_thread: HashMap::new(), @@ -2867,6 +3064,148 @@ mod tests { }); } + #[test] + fn async_question_answers_steer_the_running_turn_with_codex_reply_envelope() { + smol::block_on(async { + let (mut actor, events) = test_actor(); + actor.active_turn = Some("turn-1".into()); + // Shape recorded from codex-cli 0.156 `request_user_input_async`. + actor + .handle_line( + &json!({"method":"item/completed","params":{"threadId":"thread-1","turnId":"turn-1","item":{ + "type":"agentMessage","id":"call_q","text":"Which DB?\n- Postgres\n- SQLite", + "phase":"final_answer","delivery":"async", + "questions":[{"title":"Which DB?","options":["Postgres","SQLite"]},{"title":" "}] + }}}) + .to_string(), + ) + .await; + // The listing text is not an answer; the transcript keeps the call. + assert!(matches!( + events.recv().await.unwrap(), + AgentEvent::ItemCompleted(ThreadItem { + content: ItemContent::ToolCall { ref name, ref input, .. }, + .. + }) if name == "request_user_input_async" + && input["questions"][0]["title"] == "Which DB?" + )); + let AgentEvent::UserInputRequested { + request_id, + questions, + delivery, + } = events.recv().await.unwrap() + else { + panic!("expected UserInputRequested") + }; + assert_eq!(request_id, "call_q"); + assert_eq!(delivery, UserInputDelivery::Async); + assert_eq!(questions.len(), 1, "blank titles are dropped"); + assert_eq!( + questions[0].id, + r#"["request_user_input_async","call_q",0]"# + ); + assert_eq!(questions[0].question, "Which DB?"); + assert_eq!( + questions[0] + .options + .iter() + .map(|option| option.label.as_str()) + .collect::>(), + ["Postgres", "SQLite"] + ); + + let mut answers = serde_json::Map::new(); + answers.insert(questions[0].id.clone(), json!("SQLite")); + actor + .handle_command(SessionCommand::RespondUserInput { + request_id: "call_q".into(), + answers, + }) + .await + .unwrap(); + assert!(matches!( + events.recv().await.unwrap(), + AgentEvent::SteerRequested { ref request_id, ref text, .. } + if request_id == "codex-question-reply-call_q" + && text == "> Which DB?\n\nSQLite" + )); + assert!(matches!( + events.recv().await.unwrap(), + AgentEvent::UserInputResolved { ref request_id, .. } if request_id == "call_q" + )); + let ChildOutput::Line(request) = actor.lines.recv().await.unwrap() else { + panic!("expected echoed request") + }; + let request: Value = serde_json::from_str(&request).unwrap(); + assert_eq!(request["method"], "turn/steer"); + assert_eq!(request["params"]["expectedTurnId"], "turn-1"); + let envelope = "\n\ + [{\"answer\":\"SQLite\",\"question\":\"Which DB?\",\ + \"questionItemId\":\"[\\\"request_user_input_async\\\",\\\"call_q\\\",0]\"}]\n\ + "; + assert_eq!(request["params"]["input"][0]["text"], envelope); + + let id = request["id"].as_i64().unwrap(); + actor + .handle_line(&json!({"id": id, "result": {"turnId": "turn-1"}}).to_string()) + .await; + actor + .handle_line( + &json!({"method":"item/completed","params":{"threadId":"thread-1","item":{ + "type":"userMessage","id":"u1","content":[{"type":"text","text":envelope}] + }}}) + .to_string(), + ) + .await; + assert!(matches!( + events.recv().await.unwrap(), + AgentEvent::SteerAccepted { ref request_id } if request_id == "codex-question-reply-call_q" + )); + + // Questions still open when the turn ends can no longer be answered. + actor + .handle_line(&json!({"method":"item/completed","params":{"threadId":"thread-1","item":{ + "type":"agentMessage","id":"call_r","text":"Later?","questions":[{"title":"Later?"}] + }}}).to_string()) + .await; + assert!(matches!( + events.recv().await.unwrap(), + AgentEvent::ItemCompleted(_) + )); + assert!(matches!( + events.recv().await.unwrap(), + AgentEvent::UserInputRequested { ref request_id, .. } if request_id == "call_r" + )); + actor + .handle_line(&json!({"method":"turn/completed","params":{"threadId":"thread-1","turn":{"id":"turn-1","status":"completed"}}}).to_string()) + .await; + assert!(matches!( + events.recv().await.unwrap(), + AgentEvent::TurnCompleted { .. } + )); + let mut answers = serde_json::Map::new(); + answers.insert( + r#"["request_user_input_async","call_r",0]"#.into(), + json!("yes"), + ); + actor + .handle_command(SessionCommand::RespondUserInput { + request_id: "call_r".into(), + answers, + }) + .await + .unwrap(); + assert!(matches!( + events.recv().await.unwrap(), + AgentEvent::Warning { .. } + )); + assert!(events.try_recv().is_err()); + + let _ = actor.child.kill(); + let _ = actor.child.wait(); + }); + } + #[test] fn steer_acceptance_waits_for_the_consumed_user_message_echo() { smol::block_on(async { @@ -3936,6 +4275,7 @@ mod tests { AgentEvent::UserInputRequested { request_id, questions, + .. } => { assert_eq!(request_id, "55"); assert_eq!(questions.len(), 2, "question missing id is dropped"); @@ -4032,6 +4372,7 @@ mod tests { let AgentEvent::UserInputRequested { request_id, questions, + .. } = events.recv().await.unwrap() else { panic!("expected UserInputRequested") diff --git a/crates/agent/src/lib.rs b/crates/agent/src/lib.rs index a06c2bcf..456c8f04 100644 --- a/crates/agent/src/lib.rs +++ b/crates/agent/src/lib.rs @@ -614,6 +614,27 @@ fn selection_bool(selections: &[OptionSelection], id: &str) -> Option { .and_then(|selection| selection.value.as_bool()) } +/// Whether a [`AgentEvent::UserInputRequested`] holds the agent until it is +/// answered. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum UserInputDelivery { + /// The agent is blocked until a matching + /// [`SessionCommand::RespondUserInput`] arrives. + #[default] + Blocking, + /// The agent keeps working on the turn and reads the answer when it is + /// delivered (Codex `request_user_input_async`). The request stays open + /// until answered or the turn ends; the composer stays usable meanwhile. + Async, +} + +impl UserInputDelivery { + pub fn is_blocking(&self) -> bool { + matches!(self, Self::Blocking) + } +} + /// A structured question the agent asks the user (Claude `AskUserQuestion`, /// Codex `item/tool/requestUserInput`). Rendered as a multiple-choice (or /// free-text) prompt; answers ride back through [`SessionCommand::RespondUserInput`]. @@ -805,9 +826,9 @@ pub enum SessionCommand { decision: ApprovalDecision, }, /// Answer a pending user-input request (Claude `AskUserQuestion`, Codex - /// `item/tool/requestUserInput`). Each value is a string (single-select / - /// free text) or an array of strings (multi-select), keyed by the matching - /// [`UserInputQuestion::id`]. + /// `item/tool/requestUserInput` and `request_user_input_async`). Each value + /// is a string (single-select / free text) or an array of strings + /// (multi-select), keyed by the matching [`UserInputQuestion::id`]. RespondUserInput { request_id: String, answers: serde_json::Map, @@ -1077,11 +1098,14 @@ pub enum AgentEvent { request_id: String, decision: ApprovalDecision, }, - /// The agent is asking the user one or more structured questions and is - /// blocked until a matching [`SessionCommand::RespondUserInput`] arrives. + /// The agent is asking the user one or more structured questions. + /// `delivery` says whether it waits for the matching + /// [`SessionCommand::RespondUserInput`] or keeps working meanwhile. UserInputRequested { request_id: String, questions: Vec, + #[serde(default, skip_serializing_if = "UserInputDelivery::is_blocking")] + delivery: UserInputDelivery, }, /// A pending user-input request has been settled (answered, or cancelled on /// teardown — in which case `answers` is empty). @@ -1858,6 +1882,27 @@ mod thread_item_serde_tests { ); } + #[test] + fn user_input_delivery_is_async_only_when_recorded() { + let legacy = r#"{"type":"user_input_requested","request_id":"q","questions":[]}"#; + assert!(matches!( + serde_json::from_str::(legacy).unwrap(), + AgentEvent::UserInputRequested { + delivery: UserInputDelivery::Blocking, + .. + } + )); + let event = AgentEvent::UserInputRequested { + request_id: "q".into(), + questions: Vec::new(), + delivery: UserInputDelivery::Async, + }; + assert_eq!( + serde_json::to_string(&event).unwrap(), + r#"{"type":"user_input_requested","request_id":"q","questions":[],"delivery":"async"}"# + ); + } + #[test] fn stored_items_preserve_parent_links_and_read_legacy_records() { let legacy = r#"{"type":"item_completed","id":"old","content":{"kind":"user_message","text":"hello"}}"#; diff --git a/crates/agent/src/opencode.rs b/crates/agent/src/opencode.rs index 4ca61120..30422ed8 100644 --- a/crates/agent/src/opencode.rs +++ b/crates/agent/src/opencode.rs @@ -21,7 +21,8 @@ use crate::{ Attachment, ChangeCompleteness, DeltaKind, FileChange, FileChangeKind, InteractionMode, ItemContent, ItemStatus, LaunchEnv, ModelSpec, OptionDescriptor, OptionSelection, ProviderCommand, ProviderCommandKind, ProviderKind, ResumeCursor, SelectOption, SessionCommand, - SessionHandle, SessionOptions, ThreadItem, TokenUsage, UserInputOption, UserInputQuestion, + SessionHandle, SessionOptions, ThreadItem, TokenUsage, UserInputDelivery, UserInputOption, + UserInputQuestion, }; pub async fn start(opts: SessionOptions) -> Result { @@ -707,6 +708,7 @@ impl OpenCodeMapper { mapped.events.push(AgentEvent::UserInputRequested { request_id, questions, + delivery: UserInputDelivery::Blocking, }); } } @@ -1978,7 +1980,11 @@ mod tests { })); assert!(matches!( mapped.events.as_slice(), - [AgentEvent::UserInputRequested { request_id, questions }] + [AgentEvent::UserInputRequested { + request_id, + questions, + .. + }] if request_id == "que_1" && questions.len() == 2 && questions[0].id == "que_1:0" diff --git a/crates/agent/src/pi.rs b/crates/agent/src/pi.rs index 9bae2524..8963b4d5 100644 --- a/crates/agent/src/pi.rs +++ b/crates/agent/src/pi.rs @@ -21,7 +21,7 @@ use crate::{ Attachment, DeltaKind, FileChange, FileChangeKind, InteractionMode, ItemContent, ItemStatus, LaunchEnv, ModelSpec, OptionDescriptor, OptionSelection, ProviderCommand, ProviderCommandKind, ProviderKind, ResumeCursor, SelectOption, SessionCommand, SessionHandle, SessionOptions, - ThreadItem, TokenUsage, UserInputOption, UserInputQuestion, + ThreadItem, TokenUsage, UserInputDelivery, UserInputOption, UserInputQuestion, }; const PERMISSION_EXTENSION: &str = include_str!("../assets/pi/tcode-permissions.ts"); @@ -715,6 +715,7 @@ impl PiActor { .emit(AgentEvent::UserInputRequested { request_id: id, questions: vec![question], + delivery: UserInputDelivery::Blocking, }) .await; return; diff --git a/crates/core/src/session.rs b/crates/core/src/session.rs index 7bab9cc2..06130e21 100644 --- a/crates/core/src/session.rs +++ b/crates/core/src/session.rs @@ -8,7 +8,8 @@ use std::sync::Arc; use agent::{ AgentEvent, ApprovalRequest, ChangeCompleteness, DeltaKind, FileChange, ItemContent, - ItemStatus, PlanStep, ResumeCursor, ThreadItem, TokenUsage, TurnStatus, UserInputQuestion, + ItemStatus, PlanStep, ResumeCursor, ThreadItem, TokenUsage, TurnStatus, UserInputDelivery, + UserInputQuestion, }; use serde::{Deserialize, Serialize}; @@ -482,6 +483,15 @@ pub struct ProposedPlan { pub turn: usize, } +/// A structured question set the agent is waiting on, or working past, from +/// [`AgentEvent::UserInputRequested`]. +#[derive(Debug, Clone, PartialEq)] +pub struct PendingUserInput { + pub request_id: String, + pub questions: Vec, + pub delivery: UserInputDelivery, +} + /// Folded view of a session's event history. #[derive(Debug, Clone, Default)] pub struct Timeline { @@ -505,10 +515,10 @@ pub struct Timeline { pub plan_steps: Vec, /// The explanation string from the latest `PlanUpdated`, if any. pub plan_explanation: Option, - /// The active user-input request (Claude `AskUserQuestion` / Codex - /// `requestUserInput`), if the agent is currently blocked on one. Cleared - /// when it resolves or the turn ends. - pub pending_user_input: Option<(String, Vec)>, + /// The open user-input request (Claude `AskUserQuestion`, Codex + /// `requestUserInput` / `request_user_input_async`), if any. Cleared when + /// it resolves or the turn ends. + pub pending_user_input: Option, pub usage: Option, pub resume: Option, pub provider_session_id: Option, @@ -821,14 +831,19 @@ impl Timeline { AgentEvent::UserInputRequested { request_id, questions, + delivery, } => { - self.pending_user_input = Some((request_id.clone(), questions.clone())); + self.pending_user_input = Some(PendingUserInput { + request_id: request_id.clone(), + questions: questions.clone(), + delivery: *delivery, + }); } AgentEvent::UserInputResolved { request_id, .. } => { if self .pending_user_input .as_ref() - .is_some_and(|(id, _)| id == request_id) + .is_some_and(|pending| pending.request_id == *request_id) { self.pending_user_input = None; } @@ -1638,6 +1653,7 @@ mod tests { multi_select: false, prefill: None, }], + delivery: UserInputDelivery::Blocking, }; let mut timeline = Timeline::default(); timeline.apply_at(None, &request); @@ -1645,7 +1661,7 @@ mod tests { timeline .pending_user_input .as_ref() - .map(|(request_id, _)| request_id.as_str()), + .map(|pending| pending.request_id.as_str()), Some("que_1") ); diff --git a/crates/runtime/src/app/command_validation.rs b/crates/runtime/src/app/command_validation.rs index 43f242fc..19e3c78f 100644 --- a/crates/runtime/src/app/command_validation.rs +++ b/crates/runtime/src/app/command_validation.rs @@ -202,7 +202,7 @@ impl AppState { .timeline .pending_user_input .as_ref() - .is_none_or(|request| &request.0 != request_id) => + .is_none_or(|pending| pending.request_id != *request_id) => { return Err(error( "unknown_user_input", @@ -411,7 +411,11 @@ mod tests { let active = state.resident_mut(&id).unwrap(); active.runtime = Runtime::Live(commands); active.turn_in_flight = true; - active.timeline.pending_user_input = Some(("question".into(), Vec::new())); + active.timeline.pending_user_input = Some(tcode_core::session::PendingUserInput { + request_id: "question".into(), + questions: Vec::new(), + delivery: agent::UserInputDelivery::Blocking, + }); state.record_approval_event( &id, &AgentEvent::ApprovalRequested(agent::ApprovalRequest { diff --git a/crates/runtime/src/app/tests.rs b/crates/runtime/src/app/tests.rs index 6b21b832..563769e6 100644 --- a/crates/runtime/src/app/tests.rs +++ b/crates/runtime/src/app/tests.rs @@ -7634,7 +7634,11 @@ fn settling_rejects_busy_descendants_and_accepted_input_reactivates_ancestors() child.turn_in_flight = busy == "turn"; child.background_task_count = usize::from(busy == "background"); child.timeline.pending_user_input = - (busy == "input").then(|| ("input".into(), Vec::new())); + (busy == "input").then(|| tcode_core::session::PendingUserInput { + request_id: "input".into(), + questions: Vec::new(), + delivery: agent::UserInputDelivery::Blocking, + }); child.timeline.pending_approvals.clear(); if busy == "queued" { child.push_queued("queued".into(), Vec::new()); diff --git a/crates/ui/src/chat/model.rs b/crates/ui/src/chat/model.rs index b725d6b4..7c4bfbcf 100644 --- a/crates/ui/src/chat/model.rs +++ b/crates/ui/src/chat/model.rs @@ -352,6 +352,12 @@ pub(crate) fn tool_brief(input: &serde_json::Value) -> String { .or_else(|| map.get("path")) .or_else(|| map.get("command")) .or_else(|| map.get("summary")) + // Question tools: Claude `AskUserQuestion` items carry `question`, + // Codex `request_user_input_async` items carry `title`. + .or_else(|| { + let first = map.get("questions")?.get(0)?; + first.get("question").or_else(|| first.get("title")) + }) .and_then(|v| v.as_str()) .map(one_line) .unwrap_or_default(), diff --git a/crates/ui/src/composer/components/user_input.rs b/crates/ui/src/composer/components/user_input.rs index 731979c8..e7d4ac16 100644 --- a/crates/ui/src/composer/components/user_input.rs +++ b/crates/ui/src/composer/components/user_input.rs @@ -1,12 +1,10 @@ use super::super::*; use crate::scroll::ScrollableElement as _; +use tcode_core::session::PendingUserInput; impl Composer { /// The active session's pending user-input request, if any. - pub(in super::super) fn pending_user_input( - &self, - cx: &App, - ) -> Option<(String, Vec)> { + pub(in super::super) fn pending_user_input(&self, cx: &App) -> Option { self.workspace_store .read(cx) .composer_state() @@ -25,14 +23,15 @@ impl Composer { .read(cx) .composer_state() .pending_user_input; - let current_id = current.as_ref().map(|(id, _)| id.clone()); + let current_id = current.as_ref().map(|pending| pending.request_id.clone()); if current_id != self.ui_request_id { self.ui_request_id = current_id; self.ui_question_index = 0; self.ui_selections.clear(); + self.ui_dismissed_request_id = None; let prefill = current .as_ref() - .and_then(|(_, questions)| questions.first()) + .and_then(|pending| pending.questions.first()) .and_then(|question| question.prefill.as_deref()) .unwrap_or_default(); self.user_input_custom.update(cx, |state, cx| { @@ -43,10 +42,12 @@ impl Composer { pub(in super::super) fn render_user_input_panel( &self, - request_id: String, - questions: Vec, + pending: &PendingUserInput, cx: &mut Context, ) -> AnyElement { + let request_id = pending.request_id.clone(); + let questions = pending.questions.clone(); + let blocking = pending.delivery.is_blocking(); let muted = cx.theme().muted_foreground; let primary = cx.theme().primary; let total = questions.len(); @@ -61,6 +62,12 @@ impl Composer { .cloned() .unwrap_or_default(); + // Codex's non-blocking questions carry no header. + let header_text = if question.header.is_empty() { + crate::tr!("userinput.async_header").into_owned() + } else { + question.header.clone() + }; let header = h_flex() .w_full() .gap_2() @@ -70,7 +77,7 @@ impl Composer { .flex_1() .text_size(px(13.)) .font_medium() - .child(question.header.clone()), + .child(header_text), ) .when(total > 1, |this| { this.child(div().text_size(px(11.)).text_color(muted).child(crate::tr!( @@ -78,6 +85,20 @@ impl Composer { index = index + 1, total = total ))) + }) + .when(!blocking, |this| { + let request_dismiss = request_id.clone(); + this.child( + Button::new("ui-dismiss") + .ghost() + .xsmall() + .icon(IconName::Close) + .tooltip(crate::tr!("userinput.dismiss")) + .on_click(cx.listener(move |this, _, _, cx| { + this.ui_dismissed_request_id = Some(request_dismiss.clone()); + cx.notify(); + })), + ) }); let mut options_content = v_flex().w_full().gap_1(); @@ -310,6 +331,14 @@ impl Composer { .shadow_sm() .child(header) .child(div().text_size(px(13.)).child(question.question.clone())) + .when(!blocking, |this| { + this.child( + div() + .text_size(px(11.)) + .text_color(muted) + .child(crate::tr!("userinput.async_hint")), + ) + }) .child(options) .child(custom_answer) .when(multi, |this| { @@ -326,9 +355,11 @@ impl Composer { } /// Number keys 1-9 pressed in the (empty) main composer input select the - /// matching option of the pending question. Returns true when consumed. - /// Deliberately NOT wired to the panel itself: the only focusable child - /// there is the custom-answer textarea, where digits must stay literal. + /// matching option of the pending blocking question. Returns true when + /// consumed. Deliberately NOT wired to the panel itself: the only focusable + /// child there is the custom-answer textarea, where digits must stay + /// literal. While a non-blocking question is open the composer still + /// writes ordinary messages, so digits stay literal there too. pub(in super::super) fn handle_user_input_digit( &mut self, ev: &gpui::KeyDownEvent, @@ -338,9 +369,17 @@ impl Composer { if ev.keystroke.modifiers.modified() || !self.input.read(cx).value().is_empty() { return false; } - let Some((request_id, questions)) = self.pending_user_input(cx) else { + let Some(PendingUserInput { + request_id, + questions, + delivery, + }) = self.pending_user_input(cx) + else { return false; }; + if !delivery.is_blocking() { + return false; + } let index = self .ui_question_index .min(questions.len().saturating_sub(1)); @@ -388,7 +427,12 @@ impl Composer { window: &mut Window, cx: &mut Context, ) { - let Some((request_id, questions)) = self.pending_user_input(cx) else { + let Some(PendingUserInput { + request_id, + questions, + .. + }) = self.pending_user_input(cx) + else { return; }; let text = input.read(cx).value().trim().to_string(); diff --git a/crates/ui/src/composer/mod.rs b/crates/ui/src/composer/mod.rs index fa7179a3..de9cda70 100644 --- a/crates/ui/src/composer/mod.rs +++ b/crates/ui/src/composer/mod.rs @@ -137,6 +137,10 @@ pub struct Composer { ui_request_id: Option, ui_question_index: usize, ui_selections: std::collections::HashMap>, + /// A non-blocking request the user closed on this client; its panel stays + /// hidden until a newer request replaces it. The agent keeps working and + /// the question text remains in the transcript. + ui_dismissed_request_id: Option, /// The placeholder text last applied to the input (so it is only re-set — /// which notifies — when it actually changes). applied_placeholder: String, @@ -406,6 +410,7 @@ impl Composer { ui_request_id: None, ui_question_index: 0, ui_selections: std::collections::HashMap::new(), + ui_dismissed_request_id: None, applied_placeholder: crate::tr!("composer.placeholder").into_owned(), model_picker_token: 0, control_width: Rc::new(Cell::new(None)), @@ -541,9 +546,13 @@ impl Composer { { return; } - // Pending questions use this text as the current custom answer, - // advancing through the same path as an option click. - if self.pending_user_input(cx).is_some() { + // A blocking question uses this text as the current custom answer, + // advancing through the same path as an option click. A non-blocking + // one leaves the composer to ordinary sends and steers. + if self + .pending_user_input(cx) + .is_some_and(|pending| pending.delivery.is_blocking()) + { self.submit_custom_user_input(input, window, cx); return; } @@ -1104,7 +1113,9 @@ impl Render for Composer { }); } - let user_input = self.pending_user_input(cx); + let user_input = self + .pending_user_input(cx) + .filter(|pending| self.ui_dismissed_request_id.as_ref() != Some(&pending.request_id)); let fallback_block = self .workspace_store .read(cx) @@ -1306,8 +1317,8 @@ impl Render for Composer { .when_some(approval, |this, request| { this.child(self.render_approval_panel(&request, approval_count, cx)) }) - .when_some(user_input, |this, (request_id, questions)| { - this.child(self.render_user_input_panel(request_id, questions, cx)) + .when_some(user_input, |this, pending| { + this.child(self.render_user_input_panel(&pending, cx)) }) .when_some(fallback_block, |this, block| { this.child(self.render_fallback_panel(&block, cx)) diff --git a/crates/ui/src/store/snapshots.rs b/crates/ui/src/store/snapshots.rs index 0db476f8..3ccd8885 100644 --- a/crates/ui/src/store/snapshots.rs +++ b/crates/ui/src/store/snapshots.rs @@ -2,7 +2,7 @@ use std::path::PathBuf; use tcode_core::{ project::WorktreeInfo, - session::Timeline, + session::{PendingUserInput, Timeline}, settings::Settings, ui::{RightTab, WorkspaceMode}, }; @@ -44,7 +44,7 @@ pub(crate) struct ComposerState { pub active_cwd: Option, pub provider_commands: Vec, pub attachments_dir: Option, - pub pending_user_input: Option<(String, Vec)>, + pub pending_user_input: Option, pub active_model: Option, pub model_pending_restart: bool, pub active_model_spec: Option, diff --git a/locales/en.yml b/locales/en.yml index d0a5efab..a2ddb856 100644 --- a/locales/en.yml +++ b/locales/en.yml @@ -764,6 +764,9 @@ userinput: previous: "Previous" next_question: "Next question" done: "Done" + async_header: "Question from the agent" + async_hint: "The agent keeps working. Answer whenever you like, or keep chatting." + dismiss: "Hide" fallback: model_changed_title: "Turn stopped after a model change" reason_model_changed: "The responding model differs from the selected model, %{model}. Tcode stopped the turn." diff --git a/locales/zh-CN.yml b/locales/zh-CN.yml index 1accd1ed..508de3b9 100644 --- a/locales/zh-CN.yml +++ b/locales/zh-CN.yml @@ -758,6 +758,9 @@ userinput: previous: "上一个" next_question: "下一个问题" done: "完成" + async_header: "来自代理的问题" + async_hint: "代理会继续工作。你可以随时回答,也可以继续对话。" + dismiss: "隐藏" fallback: model_changed_title: "本轮因模型更改而中止" reason_model_changed: "实际响应的模型与所选模型 %{model} 不同,Tcode 已中止本轮。"