diff --git a/.github/actions/build-wework-core-e2e/action.yml b/.github/actions/build-wework-core-e2e/action.yml index cf8831a13d..12c531212e 100644 --- a/.github/actions/build-wework-core-e2e/action.yml +++ b/.github/actions/build-wework-core-e2e/action.yml @@ -8,6 +8,10 @@ inputs: backend-rs-image: description: Content-addressed Rust backend runtime image required: true + runtime-binaries-ready: + description: Whether the caller already provided the runtime binaries from images that could not be pulled + required: false + default: "false" save-cache: description: Whether writable dependency caches may be saved required: false @@ -38,6 +42,7 @@ runs: env: EXECUTOR_IMAGE: ${{ inputs.executor-image }} BACKEND_RS_IMAGE: ${{ inputs.backend-rs-image }} + RUNTIME_BINARIES_READY: ${{ inputs.runtime-binaries-ready }} GITHUB_TOKEN: ${{ github.token }} OCI_RUNTIME_WAIT_SECONDS: "420" CI: true @@ -46,21 +51,34 @@ runs: WEWORK_ELECTRON_DEPENDENCIES_READY: "true" run: | mkdir -p .ci-artifacts - .github/scripts/restore-oci-runtime-binary.sh \ - "$EXECUTOR_IMAGE" \ - /app/executor \ - .ci-artifacts/wegent-executor & - executor_pid=$! - .github/scripts/restore-oci-runtime-binary.sh \ - "$BACKEND_RS_IMAGE" \ - /app/wegent-backend-rs \ - .ci-artifacts/wegent-backend-rs & - backend_pid=$! + executor_pid="" + backend_pid="" + if [[ "$RUNTIME_BINARIES_READY" == "true" ]]; then + for runtime in executor backend-rs; do + if [[ ! -s ".ci-artifacts/wegent-$runtime" ]]; then + echo "Provided runtime binary is missing: .ci-artifacts/wegent-$runtime" >&2 + exit 1 + fi + # Artifact upload does not preserve the executable bit. + chmod 0755 ".ci-artifacts/wegent-$runtime" + done + else + .github/scripts/restore-oci-runtime-binary.sh \ + "$EXECUTOR_IMAGE" \ + /app/executor \ + .ci-artifacts/wegent-executor & + executor_pid=$! + .github/scripts/restore-oci-runtime-binary.sh \ + "$BACKEND_RS_IMAGE" \ + /app/wegent-backend-rs \ + .ci-artifacts/wegent-backend-rs & + backend_pid=$! + fi status=0 pnpm --filter wework ai:verify:electron:build || status=$? - if ! wait "$executor_pid"; then status=1; fi - if ! wait "$backend_pid"; then status=1; fi + if [[ -n "$executor_pid" ]] && ! wait "$executor_pid"; then status=1; fi + if [[ -n "$backend_pid" ]] && ! wait "$backend_pid"; then status=1; fi if ((status == 0)); then test -x wework/electron/release/WeWork-linux-x64/WeWork test -x wework/electron/release/WeWork-linux-x64/resources/bin/wegent-executor diff --git a/.github/workflows/wework-e2e.yml b/.github/workflows/wework-e2e.yml index a58e64722e..bac4651def 100644 --- a/.github/workflows/wework-e2e.yml +++ b/.github/workflows/wework-e2e.yml @@ -51,6 +51,7 @@ jobs: wework_desktop_other_e2e: ${{ steps.classify.outputs.wework_desktop_other_e2e }} wework_desktop_other_e2e_matrix: ${{ steps.classify.outputs.wework_desktop_other_e2e_matrix }} wework_desktop_macos_inspector_e2e: ${{ steps.classify.outputs.wework_desktop_macos_inspector_e2e }} + wework_desktop_rust_runtime_missing: ${{ steps.rust-runtime-check.outputs.missing }} steps: - name: Checkout code @@ -146,6 +147,26 @@ jobs: echo "backend-rs-ref=ghcr.io/${GITHUB_REPOSITORY_OWNER,,}/wegent-backend-rs-e2e:v1-$BACKEND_RS_SOURCE_DIGEST" } >> "$GITHUB_OUTPUT" + # Content-addressed Rust runtimes are only published by main and by same-repo + # pull requests. Pull requests from forks cannot publish them, so they have to + # rebuild the missing runtimes locally instead of waiting for an image that + # will never appear. + - name: Check shared Rust E2E runtimes + if: steps.classify.outputs.wework_desktop_e2e == 'true' && (github.event_name != 'pull_request' || (github.event.pull_request.draft == false && (github.event.action != 'labeled' || github.event.label.name == 'ci:all'))) + id: rust-runtime-check + env: + EXECUTOR_IMAGE: ${{ steps.runtime.outputs.executor-ref }} + BACKEND_RS_IMAGE: ${{ steps.runtime.outputs.backend-rs-ref }} + run: | + missing=false + for image in "$EXECUTOR_IMAGE" "$BACKEND_RS_IMAGE"; do + if ! docker manifest inspect "$image" >/dev/null 2>&1; then + echo "Missing content-addressed runtime: $image" >&2 + missing=true + fi + done + echo "missing=$missing" >> "$GITHUB_OUTPUT" + - name: Build and publish Wework browser E2E image if: steps.browser-image-check.outputs.exists == 'false' uses: docker/build-push-action@ca052bb54ab0790a636c9b5f226502c73d547a25 @@ -218,11 +239,99 @@ jobs: if-no-files-found: ignore retention-days: 7 + prepare-wework-desktop-e2e-runtime: + name: Rebuild Missing Wework Desktop E2E Runtimes + needs: + - changes + if: needs.changes.outputs.wework_desktop_rust_runtime_missing == 'true' && (needs.changes.outputs.wework_desktop_core_e2e == 'true' || needs.changes.outputs.wework_desktop_cloud_e2e == 'true' || needs.changes.outputs.wework_desktop_other_e2e == 'true') && (github.event_name != 'pull_request' || (github.event.pull_request.draft == false && (github.event.action != 'labeled' || github.event.label.name == 'ci:all'))) + runs-on: ubuntu-latest + timeout-minutes: 30 + permissions: + contents: read + packages: read + + steps: + - name: Checkout code + uses: actions/checkout@v4 + with: + filter: blob:none + persist-credentials: false + + - name: Log in to GHCR + uses: docker/login-action@c94ce9fb468520275223c153574b00df6fe4bcc9 + with: + registry: ghcr.io + username: ${{ github.actor }} + password: ${{ secrets.GITHUB_TOKEN }} + + - name: Set up Docker Buildx + uses: docker/setup-buildx-action@8d2750c68a42422c14e847fe6c8ac0403b4cbd6f + + - name: Resolve Executor E2E base image + id: executor-runtime + env: + SOURCE_DIGEST: ${{ hashFiles('frontend/e2e/fixtures/claudecode-executor/Dockerfile', 'executor/Cargo.toml', 'executor/Cargo.lock', 'executor/src/**', 'shared/assets/**') }} + run: .github/scripts/resolve-executor-e2e-runtime.sh + + - name: Build Executor E2E runtime image + uses: docker/build-push-action@ca052bb54ab0790a636c9b5f226502c73d547a25 + with: + context: . + file: frontend/e2e/fixtures/claudecode-executor/Dockerfile + build-args: BASE_IMAGE=${{ steps.executor-runtime.outputs.base-image }} + load: true + push: false + tags: wegent/e2e-claudecode-executor:latest + cache-from: type=registry,ref=ghcr.io/${{ github.repository_owner }}/wegent-executor:buildcache-e2e + + - name: Extract Executor E2E runtime binary + run: | + mkdir -p .ci-artifacts + container_id="$(docker create wegent/e2e-claudecode-executor:latest)" + trap 'docker rm -f "$container_id" >/dev/null 2>&1 || true' EXIT + docker cp "$container_id:/app/executor" .ci-artifacts/wegent-executor + chmod 0755 .ci-artifacts/wegent-executor + + - name: Set up Rust + uses: dtolnay/rust-toolchain@29eef336d9b2848a0b548edc03f92a220660cdb8 + + - name: Set up executor Rust cache + uses: ./.github/actions/setup-executor-rust-cache + + - name: Build Rust backend E2E gateway + env: + CARGO_PROFILE_RELEASE_STRIP: "false" + CARGO_PROFILE_RELEASE_BUILD_OVERRIDE_STRIP: "false" + CARGO_TARGET_DIR: ${{ runner.temp }}/backend-rs-target + run: | + cargo build \ + --release \ + --locked \ + --manifest-path backend-rs/Cargo.toml \ + --bin wegent-backend-rs + mkdir -p .ci-artifacts + cp "$CARGO_TARGET_DIR/release/wegent-backend-rs" .ci-artifacts/wegent-backend-rs + chmod 0755 .ci-artifacts/wegent-backend-rs + + - name: Upload rebuilt Wework desktop E2E runtimes + uses: actions/upload-artifact@v4 + with: + name: wework-desktop-e2e-runtimes + path: | + .ci-artifacts/wegent-executor + .ci-artifacts/wegent-backend-rs + if-no-files-found: error + retention-days: 1 + compression-level: 0 + build-wework-desktop-core-e2e: name: Build Shared Wework Desktop E2E Artifact needs: - changes - if: (needs.changes.outputs.wework_desktop_core_e2e == 'true' || needs.changes.outputs.wework_desktop_cloud_e2e == 'true' || needs.changes.outputs.wework_desktop_other_e2e == 'true') && (github.event_name != 'pull_request' || (github.event.pull_request.draft == false && (github.event.action != 'labeled' || github.event.label.name == 'ci:all'))) + - prepare-wework-desktop-e2e-runtime + # `always()` keeps this job running when the runtime rebuild was skipped, while + # the explicit result checks still stop it when that rebuild failed or was cancelled. + if: always() && needs.prepare-wework-desktop-e2e-runtime.result != 'failure' && needs.prepare-wework-desktop-e2e-runtime.result != 'cancelled' && (needs.changes.outputs.wework_desktop_core_e2e == 'true' || needs.changes.outputs.wework_desktop_cloud_e2e == 'true' || needs.changes.outputs.wework_desktop_other_e2e == 'true') && (github.event_name != 'pull_request' || (github.event.pull_request.draft == false && (github.event.action != 'labeled' || github.event.label.name == 'ci:all'))) runs-on: ubuntu-latest timeout-minutes: 35 env: @@ -240,6 +349,13 @@ jobs: filter: blob:none persist-credentials: false + - name: Download rebuilt Wework desktop E2E runtimes + if: needs.prepare-wework-desktop-e2e-runtime.result == 'success' + uses: actions/download-artifact@v4 + with: + name: wework-desktop-e2e-runtimes + path: .ci-artifacts + - name: Set up ORAS uses: oras-project/setup-oras@1d808f7d7f6995cc68b7bf507bfe5c5446e1dc9d with: @@ -287,6 +403,7 @@ jobs: with: executor-image: ${{ needs.changes.outputs.executor_image }} backend-rs-image: ${{ needs.changes.outputs.backend_rs_image }} + runtime-binaries-ready: ${{ needs.prepare-wework-desktop-e2e-runtime.result == 'success' }} save-cache: ${{ github.ref == 'refs/heads/main' }} - name: Upload Wework desktop Core E2E build diff --git a/executor/src/runtime_work/codex_user_input.rs b/executor/src/runtime_work/codex_user_input.rs new file mode 100644 index 0000000000..6c5924b4e3 --- /dev/null +++ b/executor/src/runtime_work/codex_user_input.rs @@ -0,0 +1,115 @@ +// SPDX-FileCopyrightText: 2026 Weibo, Inc. +// +// SPDX-License-Identifier: Apache-2.0 + +use serde_json::{json, Value}; + +use super::util::string_field; + +/// A `request_user_input_async` question: one prompt plus suggested answers. +/// +/// Codex uses the same shape in the live `agentMessage.questions` payload and in +/// the replayed `request_user_input_async` tool arguments, so both projections of +/// the interactive card normalize here. The answer to such a question is always +/// the next user message, never a runtime response. +pub(crate) fn async_question_render_payload(item_id: &str, questions: &[Value]) -> Option { + let questions = questions + .iter() + .enumerate() + .filter_map(|(index, question)| async_question(index, question)) + .collect::>(); + if questions.is_empty() { + return None; + } + Some(json!({ + "kind": "request_user_input", + "delivery": "async", + "itemId": item_id, + "questions": questions, + })) +} + +fn async_question(index: usize, question: &Value) -> Option { + let prompt = string_field(question, "title") + .or_else(|| string_field(question, "question")) + .unwrap_or_default() + .trim() + .to_owned(); + let options = question + .get("options") + .and_then(Value::as_array) + .into_iter() + .flatten() + .filter_map(option_label) + .collect::>(); + if prompt.is_empty() && options.is_empty() { + return None; + } + Some(json!({ + "id": format!("question_{}", index + 1), + "question": prompt, + "options": options, + // Codex accepts a free-text answer next to the suggested options. + "is_other": true, + })) +} + +fn option_label(option: &Value) -> Option { + let label = option + .as_str() + .map(str::to_owned) + .or_else(|| string_field(option, "label")) + .unwrap_or_default(); + let label = label.trim(); + (!label.is_empty()).then(|| json!({ "label": label })) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn normalizes_title_and_plain_string_options() { + let payload = async_question_render_payload( + "call-1", + &[json!({ + "title": "Which state jitters?", + "options": ["Following", "Reading", "Both"], + })], + ) + .expect("payload"); + + assert_eq!(payload["kind"], "request_user_input"); + assert_eq!(payload["delivery"], "async"); + assert_eq!(payload["itemId"], "call-1"); + assert_eq!(payload["questions"][0]["id"], "question_1"); + assert_eq!(payload["questions"][0]["question"], "Which state jitters?"); + assert_eq!(payload["questions"][0]["is_other"], true); + assert_eq!(payload["questions"][0]["options"][0]["label"], "Following"); + assert_eq!(payload["questions"][0]["options"][2]["label"], "Both"); + } + + #[test] + fn drops_blank_questions_and_blank_options() { + let payload = async_question_render_payload( + "call-2", + &[ + json!({"title": " "}), + json!({"title": "Pick one", "options": [" ", "Second"]}), + ], + ) + .expect("payload"); + + let questions = payload["questions"].as_array().expect("questions"); + assert_eq!(questions.len(), 1); + assert_eq!(questions[0]["question"], "Pick one"); + assert_eq!(questions[0]["options"].as_array().map(Vec::len), Some(1)); + assert_eq!(questions[0]["options"][0]["label"], "Second"); + } + + #[test] + fn rejects_questions_without_prompt_or_options() { + assert!(async_question_render_payload("call-3", &[json!({"title": ""})]).is_none()); + assert!(async_question_render_payload("call-4", &[]).is_none()); + } +} diff --git a/executor/src/runtime_work/events.rs b/executor/src/runtime_work/events.rs index 885916c461..417921842c 100644 --- a/executor/src/runtime_work/events.rs +++ b/executor/src/runtime_work/events.rs @@ -26,6 +26,7 @@ use super::{ codex_notifications::{ codex_notification, debug_ignored_codex_notification, is_root_codex_turn_event, }, + codex_user_input::async_question_render_payload, notification_mapping::{ codex_stream_debug_enabled, log_dropped_notification, log_stream_text_mapping, log_text_mapping, map_text_chunk, map_tool_output_delta, notification_item_id, @@ -537,6 +538,10 @@ impl CodexNotificationEventMapper { if self.emit_applied_guidance(&emit_context, notification.params) { return; } + if emit_async_request_user_input(&emit_context, notification.params) { + self.agent_message_phases.forget_item(notification.params); + return; + } if self.emit_text_chunk( &emit_context, ¬ification.method, @@ -2009,6 +2014,25 @@ fn emit_request_user_input( object.insert("requestId".to_owned(), request_id.clone()); } } + emit_request_user_input_block( + event_tx, + device_id, + local_task_id, + request, + block_id, + render_payload, + ); +} + +/// Emits the interactive question block shared by every request-user-input source. +fn emit_request_user_input_block( + event_tx: &Option>, + device_id: &str, + local_task_id: &str, + request: &ExecutionRequest, + block_id: String, + render_payload: Value, +) { emit_response_event( event_tx, device_id, @@ -2028,6 +2052,34 @@ fn emit_request_user_input( ); } +/// Codex asks non-blocking clarifying questions through `request_user_input_async`. +/// The tool returns immediately, so the app-server delivers the question as an +/// `agentMessage` carrying `questions`, and the answer comes back as the next user +/// message instead of a runtime response. Render the structured choices as the +/// interactive card rather than the plain-text fallback the same item also carries. +fn emit_async_request_user_input(context: &EventEmitContext<'_>, params: &Value) -> bool { + let item = params.get("item").unwrap_or(params); + if item_type(item).as_str() != "agentmessage" { + return false; + } + let Some(questions) = item.get("questions").and_then(Value::as_array) else { + return false; + }; + let item_id = item_id(item, "request-user-input"); + let Some(render_payload) = async_question_render_payload(item_id.as_str(), questions) else { + return false; + }; + emit_request_user_input_block( + context.event_tx, + context.device_id, + context.local_task_id, + context.request, + format!("request-user-input-{item_id}"), + render_payload, + ); + true +} + fn emit_codex_approval_request( event_tx: &Option>, device_id: &str, @@ -2082,22 +2134,13 @@ fn emit_codex_approval_request( object.insert("requestId".to_owned(), request_id.clone()); } } - emit_response_event( + emit_request_user_input_block( event_tx, device_id, - "response.block.created", local_task_id, request, - json!({ - "block": { - "id": block_id, - "type": "tool", - "tool_name": "request_user_input", - "status": "pending", - "timestamp": now_ms(), - "render_payload": render_payload, - } - }), + block_id, + render_payload, ); } @@ -5304,6 +5347,65 @@ mod tests { assert_eq!(block["render_payload"]["questions"][0]["id"], "goal"); } + #[test] + fn maps_codex_async_questions_to_interactive_tool_block() { + let (event_tx, mut event_rx) = broadcast::channel(4); + let request = ExecutionRequest { + task_id: "7".to_owned(), + subtask_id: "8".to_owned(), + ..ExecutionRequest::default() + }; + + map_codex_notification( + &Some(event_tx), + "device-1", + "local-1", + &request, + json!({ + "method": "item/completed", + "params": { + "item": { + "id": "call-question", + "type": "agentMessage", + "phase": "final_answer", + "delivery": "async", + "text": "Which state jitters?\n- Following\n- Reading", + "questions": [ + { + "title": "Which state jitters?", + "options": ["Following", "Reading"] + } + ] + } + } + }), + ); + + let event = event_rx + .try_recv() + .expect("async question should emit an interactive block"); + let block = &event["payload"]["data"]["block"]; + assert_eq!(event["event"], "response.block.created"); + assert_eq!(block["type"], "tool"); + assert_eq!(block["tool_name"], "request_user_input"); + assert_eq!(block["status"], "pending"); + assert_eq!(block["id"], "request-user-input-call-question"); + assert_eq!(block["render_payload"]["kind"], "request_user_input"); + assert_eq!(block["render_payload"]["delivery"], "async"); + assert_eq!(block["render_payload"]["itemId"], "call-question"); + assert_eq!( + block["render_payload"]["questions"][0]["question"], + "Which state jitters?" + ); + assert_eq!( + block["render_payload"]["questions"][0]["options"][1]["label"], + "Reading" + ); + // The plain-text fallback that accompanies the structured question must not + // also land in the transcript as the final answer. + assert!(event_rx.try_recv().is_err()); + } + #[test] fn maps_mcp_form_elicitation_to_interactive_tool_block() { let (event_tx, mut event_rx) = broadcast::channel(4); diff --git a/executor/src/runtime_work/mod.rs b/executor/src/runtime_work/mod.rs index b51364a051..d1a4de4dec 100644 --- a/executor/src/runtime_work/mod.rs +++ b/executor/src/runtime_work/mod.rs @@ -8,6 +8,7 @@ mod codex_global_state; mod codex_notifications; mod codex_rollout; mod codex_transcript_page; +mod codex_user_input; mod collaboration_projects; mod connectors; mod events; diff --git a/executor/src/runtime_work/transcript.rs b/executor/src/runtime_work/transcript.rs index 22597580d5..434cad836f 100644 --- a/executor/src/runtime_work/transcript.rs +++ b/executor/src/runtime_work/transcript.rs @@ -15,6 +15,7 @@ use crate::{ services::turn_file_changes::persist_named_artifact, }; +use super::codex_user_input::async_question_render_payload; use super::util::{ bool_field, codex_wrapped_item_payload, extract_text, id_field, integer_field, is_codex_context_compaction_item_type, is_codex_tool_item_type, is_codex_tool_output_item_type, @@ -221,7 +222,9 @@ impl<'a> TurnTranscriptProjector<'a> { self.assistant.blocks.push(block); } "agentmessage" | "agentmessageevent" if !self.project_subagent_message(item) => { - self.project_assistant_message(item, has_later_process); + if !self.project_async_request_user_input(item) { + self.project_assistant_message(item, has_later_process); + } } "agentmessage" | "agentmessageevent" => {} "message" => self.project_role_message(item, has_later_process), @@ -290,6 +293,31 @@ impl<'a> TurnTranscriptProjector<'a> { } } + /// Codex records a non-blocking clarifying question as an `AgentMessage` that + /// carries `questions`, so history must restore the interactive card instead of + /// folding the question into the assistant's final text. + fn project_async_request_user_input(&mut self, item: &Value) -> bool { + let Some(questions) = item.get("questions").and_then(Value::as_array) else { + return false; + }; + let item_id = item_id(item, "request-user-input"); + let Some(render_payload) = async_question_render_payload(item_id.as_str(), questions) + else { + return false; + }; + self.assistant.blocks.push(json!({ + "id": format!("request-user-input-{item_id}"), + "subtaskId": self.subtask_id, + "type": "tool", + "tool_use_id": item_id, + "tool_name": "request_user_input", + "status": "pending", + "timestamp": item_timestamp(item).unwrap_or(self.created_at), + "render_payload": render_payload, + })); + true + } + fn project_file_change(&mut self, item: &Value, summary: Option) { let Some(summary) = summary else { return; diff --git a/executor/src/runtime_work/transcript/tests_core.rs b/executor/src/runtime_work/transcript/tests_core.rs index 8d165d3ec2..7c0c001cb8 100644 --- a/executor/src/runtime_work/transcript/tests_core.rs +++ b/executor/src/runtime_work/transcript/tests_core.rs @@ -938,6 +938,133 @@ fn transcript_restores_pending_request_user_input_as_interactive_block() { ); } +#[test] +fn transcript_restores_pending_async_request_user_input_as_interactive_block() { + let thread = json!({ + "id": "thread-1", + "turns": [ + { + "id": "turn-1", + "startedAt": 1_780_000_000, + "status": "completed", + "items": [ + { + "type": "response_item", + "payload": { + "type": "function_call", + "call_id": "call-1", + "name": "request_user_input_async", + "arguments": "{\"questions\":[{\"title\":\"Which state jitters?\",\"options\":[\"Following\",\"Reading\"]}]}" + } + }, + { + "type": "response_item", + "payload": { + "type": "function_call_output", + "call_id": "call-1", + "output": "{\"accepted\":true}" + } + } + ] + } + ] + }); + + let messages = transcript_messages(&thread, "device-1"); + let block = messages + .iter() + .flat_map(|message| message["blocks"].as_array().into_iter().flatten().cloned()) + .find(|block| block["tool_name"] == "request_user_input") + .expect("async question block"); + + assert_eq!(block["type"], "tool"); + // The async tool call itself returns immediately, so only the render payload + // marks the interaction; it stays unanswered until the next user message. + assert_eq!(block["render_payload"]["kind"], "request_user_input"); + assert_eq!(block["render_payload"]["delivery"], "async"); + assert_eq!(block["render_payload"]["itemId"], "call-1"); + assert_eq!( + block["render_payload"]["questions"][0]["question"], + "Which state jitters?" + ); + assert_eq!( + block["render_payload"]["questions"][0]["options"][0]["label"], + "Following" + ); + // The `accepted` tool output is not an answer, so the card must stay open. + assert!(block["render_payload"].get("response").is_none()); +} + +#[test] +fn transcript_restores_async_question_recorded_as_agent_message() { + let thread = json!({ + "id": "thread-1", + "turns": [ + { + "id": "turn-1", + "startedAt": 1_780_000_000, + "status": "completed", + "items": [ + { + "type": "AgentMessage", + "id": "call-question", + "content": [{"type": "Text", "text": "Which state jitters?\n- Following\n- Reading"}], + "phase": "final_answer", + "delivery": "async", + "questions": [ + {"title": "Which state jitters?", "options": ["Following", "Reading"]} + ] + } + ] + } + ] + }); + + let messages = transcript_messages(&thread, "device-1"); + let block = messages + .iter() + .flat_map(|message| message["blocks"].as_array().into_iter().flatten().cloned()) + .find(|block| block["tool_name"] == "request_user_input") + .expect("async question block"); + + assert_eq!(block["type"], "tool"); + assert_eq!(block["render_payload"]["kind"], "request_user_input"); + assert_eq!(block["render_payload"]["delivery"], "async"); + assert_eq!(block["render_payload"]["itemId"], "call-question"); + assert_eq!( + block["render_payload"]["questions"][0]["question"], + "Which state jitters?" + ); + assert_eq!( + block["render_payload"]["questions"][0]["options"][1]["label"], + "Reading" + ); + // The plain-text fallback must not also become the assistant's final text. + assert!(messages.iter().all(|message| !message["content"] + .as_str() + .unwrap_or_default() + .contains("Reading"))); + // Turn projection feeds the workspace from `runtimeItems`, so the interactive + // block has to be listed there as well. + let runtime_block = messages + .iter() + .flat_map(|message| { + message["runtimeItems"] + .as_array() + .into_iter() + .flatten() + .cloned() + }) + .find(|item| { + item["type"] == "block" + && item["block"]["tool_name"].as_str() == Some("request_user_input") + }); + assert!( + runtime_block.is_some(), + "runtimeItems must carry the question block" + ); +} + #[test] fn transcript_restores_answered_request_user_input_response() { let thread = json!({ diff --git a/executor/src/runtime_work/transcript/tool_projection.rs b/executor/src/runtime_work/transcript/tool_projection.rs index 7414ef0557..5ad8518541 100644 --- a/executor/src/runtime_work/transcript/tool_projection.rs +++ b/executor/src/runtime_work/transcript/tool_projection.rs @@ -3,6 +3,7 @@ // SPDX-License-Identifier: Apache-2.0 use super::*; +use crate::runtime_work::codex_user_input::async_question_render_payload; use sha2::{Digest, Sha256}; fn command_block(item: &Value, timestamp: i64, options: TranscriptBuildOptions) -> Value { @@ -180,7 +181,12 @@ fn is_meaningful_tool_output(output: &Value) -> bool { } fn insert_request_user_input_render_payload(object: &mut Map, item: &Value) { - if item_type(item) == "functioncall" && tool_name(item) == "request_user_input" { + if item_type(item) == "functioncall" + && matches!( + tool_name(item).as_str(), + "request_user_input" | "request_user_input_async" + ) + { let mut payload = parse_json_object_string(item, "arguments") .and_then(|value| value.as_object().cloned()) .unwrap_or_default(); @@ -188,8 +194,28 @@ fn insert_request_user_input_render_payload(object: &mut Map, ite "kind".to_owned(), Value::String("request_user_input".to_owned()), ); - payload.insert("requestId".to_owned(), Value::String(tool_call_id(item))); - object.insert("render_payload".to_owned(), Value::Object(payload)); + let call_id = tool_call_id(item); + let async_payload = (tool_name(item) == "request_user_input_async") + .then(|| { + payload + .get("questions") + .and_then(Value::as_array) + .and_then(|questions| { + async_question_render_payload(call_id.as_str(), questions) + }) + }) + .flatten(); + let render_payload = async_payload.unwrap_or_else(|| { + payload.insert("requestId".to_owned(), Value::String(call_id)); + Value::Object(payload) + }); + // Every request-user-input source renders the same interactive card; the + // non-blocking variant is carried by `render_payload.delivery`. + object.insert( + "tool_name".to_owned(), + Value::String("request_user_input".to_owned()), + ); + object.insert("render_payload".to_owned(), render_payload); return; } @@ -205,6 +231,11 @@ fn insert_request_user_input_render_payload(object: &mut Map, ite else { return; }; + // Async questions are answered by the next user message, so the tool output + // ("accepted") is never the answer. + if payload.get("delivery").and_then(Value::as_str) == Some("async") { + return; + } let Some(response) = output_payload_text(item) .and_then(|output| serde_json::from_str::(&output).ok()) .filter(Value::is_object) diff --git a/packages/chat-core/src/runtime-user-input.test.ts b/packages/chat-core/src/runtime-user-input.test.ts new file mode 100644 index 0000000000..0899d2f5c6 --- /dev/null +++ b/packages/chat-core/src/runtime-user-input.test.ts @@ -0,0 +1,175 @@ +// SPDX-FileCopyrightText: 2026 Weibo, Inc. +// +// SPDX-License-Identifier: Apache-2.0 + +import { describe, expect, test } from 'vitest' +import type { RequestUserInputPayload } from './runtime' +import type { WorkbenchMessage } from './workbench-message-reducer' +import { + findRequestUserInputPayload, + isAsyncRequestUserInputPayload, + resolveAsyncRequestUserInputAnswers, + resolveAsyncRequestUserInputReplies, +} from './runtime-user-input' + +const ASYNC_PAYLOAD: RequestUserInputPayload = { + kind: 'request_user_input', + delivery: 'async', + itemId: 'call-question', + questions: [ + { + id: 'question_1', + question: 'Which state jitters?', + options: [{ label: 'Following' }, { label: 'Reading' }], + }, + ], +} + +function assistantMessage(payload: RequestUserInputPayload): WorkbenchMessage { + return { + id: 'assistant-1', + role: 'assistant', + content: '', + status: 'done', + createdAt: '2026-09-29T03:01:37.000Z', + blocks: [ + { + id: 'request-user-input-call-question', + subtaskId: 'subtask-1', + type: 'tool', + toolName: 'request_user_input', + status: 'pending', + createdAt: 1, + renderPayload: payload, + }, + ], + } +} + +function userMessage(content: string): WorkbenchMessage { + return { + id: 'user-1', + role: 'user', + content, + status: 'done', + createdAt: '2026-09-29T03:05:00.000Z', + } +} + +describe('async request user input', () => { + test('treats only async delivery payloads as non-blocking', () => { + expect(isAsyncRequestUserInputPayload(ASYNC_PAYLOAD)).toBe(true) + expect(isAsyncRequestUserInputPayload({ kind: 'request_user_input' })).toBe(false) + expect(isAsyncRequestUserInputPayload(null)).toBe(false) + }) + + test('derives the answer from the reply the user typed in the composer', () => { + const messages = [assistantMessage(ASYNC_PAYLOAD), userMessage('Reading')] + + const resolved = resolveAsyncRequestUserInputAnswers(messages) + const block = resolved[0].blocks?.[0] + + expect(resolved).not.toBe(messages) + expect(block?.status).toBe('done') + expect(block && 'renderPayload' in block && block.renderPayload).toMatchObject({ + response: { itemId: 'call-question', answers: { question_1: { answers: ['Reading'] } } }, + }) + }) + + test('leaves an unanswered async question open', () => { + const messages = [assistantMessage(ASYNC_PAYLOAD)] + + const resolved = resolveAsyncRequestUserInputAnswers(messages) + + expect(resolved).toBe(messages) + expect(resolved[0].blocks?.[0].status).toBe('pending') + }) + + test('keeps a question that already carries its own answer', () => { + const answered: RequestUserInputPayload = { + ...ASYNC_PAYLOAD, + response: { itemId: 'call-question', answers: { question_1: { answers: ['Following'] } } }, + } + const messages = [assistantMessage(answered), userMessage('Reading')] + + const resolved = resolveAsyncRequestUserInputAnswers(messages) + + expect(resolved).toBe(messages) + }) + + test('does not resolve blocking questions from later user messages', () => { + const blocking: RequestUserInputPayload = { kind: 'request_user_input', requestId: 42 } + const messages = [assistantMessage(blocking), userMessage('Reading')] + + const resolved = resolveAsyncRequestUserInputAnswers(messages) + + expect(resolved).toBe(messages) + }) + + test('finds the payload behind a runtime answer key', () => { + const messages = [assistantMessage(ASYNC_PAYLOAD)] + + expect(findRequestUserInputPayload(messages, 'item:call-question')).toBe(ASYNC_PAYLOAD) + expect(findRequestUserInputPayload(messages, 'item:missing')).toBeNull() + expect(findRequestUserInputPayload(messages, null)).toBeNull() + }) + + describe('replies attributed to the questions they answered', () => { + test('shows the question above a single answer', () => { + const messages = [assistantMessage(ASYNC_PAYLOAD), userMessage('Reading')] + + const replies = resolveAsyncRequestUserInputReplies(messages) + + expect(replies.get('user-1')).toEqual([ + { question: 'Which state jitters?', answer: 'Reading' }, + ]) + }) + + test('pairs each line of a multi-question reply with its question', () => { + const payload: RequestUserInputPayload = { + ...ASYNC_PAYLOAD, + questions: [ + { id: 'q1', question: '晴天还是雨天?' }, + { id: 'q2', question: '早上还是晚上?' }, + { id: 'q3', question: '猫还是狗?' }, + ], + } + const messages = [assistantMessage(payload), userMessage('晴天\n早上\n猫')] + + const replies = resolveAsyncRequestUserInputReplies(messages) + + expect(replies.get('user-1')).toEqual([ + { question: '晴天还是雨天?', answer: '晴天' }, + { question: '早上还是晚上?', answer: '早上' }, + { question: '猫还是狗?', answer: '猫' }, + ]) + }) + + test('leaves a reply whose lines do not match the questions unattributed', () => { + const payload: RequestUserInputPayload = { + ...ASYNC_PAYLOAD, + questions: [ + { id: 'q1', question: '晴天还是雨天?' }, + { id: 'q2', question: '早上还是晚上?' }, + ], + } + const messages = [assistantMessage(payload), userMessage('晴天')] + + expect(resolveAsyncRequestUserInputReplies(messages).size).toBe(0) + }) + + test('ignores unanswered and blocking questions', () => { + const blocking: RequestUserInputPayload = { kind: 'request_user_input', requestId: 42 } + + expect( + resolveAsyncRequestUserInputReplies([assistantMessage(ASYNC_PAYLOAD)]).size + ).toBe(0) + expect( + resolveAsyncRequestUserInputReplies([ + assistantMessage(blocking), + userMessage('Reading'), + ]).size + ).toBe(0) + }) + }) +}) diff --git a/packages/chat-core/src/runtime-user-input.ts b/packages/chat-core/src/runtime-user-input.ts index dba8b04962..7994343ab8 100644 --- a/packages/chat-core/src/runtime-user-input.ts +++ b/packages/chat-core/src/runtime-user-input.ts @@ -4,6 +4,7 @@ import type { TurnFileChangesSummary, } from './runtime' import type { + WorkbenchMessage, WorkbenchProcessingBlock, WorkbenchToolBlock as ToolBlock, } from './workbench-message-reducer' @@ -12,6 +13,7 @@ type ProcessingBlock = WorkbenchProcessingBlock const EMPTY_HIDDEN_REQUEST_USER_INPUT_IDS = new Set() export const CODEX_IMPLEMENT_PLAN_QUESTION = '执行此计划?' export const CODEX_IMPLEMENT_PLAN_RESPONSE_LABEL = '是的,执行此计划' +export const ASYNC_REQUEST_USER_INPUT_DELIVERY = 'async' const IMPLEMENT_PLAN_TEXT_MARKERS = ['实施此计划', '执行此计划'] export function hasImplementationPlanText(text: string | null | undefined): boolean { @@ -95,6 +97,177 @@ export function isRequestUserInputBlock(block: ProcessingBlock): block is Reques return isRequestUserInputPayload(block.renderPayload) } +/** + * Codex's non-blocking `request_user_input_async` question. The tool returns + * immediately, so the answer arrives as the next user message instead of a + * runtime response. + */ +export function isAsyncRequestUserInputPayload( + payload: RequestUserInputPayload | null | undefined +): boolean { + return payload?.delivery === ASYNC_REQUEST_USER_INPUT_DELIVERY +} + +/** + * Async questions are answered by the next user message, so their response is + * derived from the conversation rather than from a runtime answer. Without this, + * a question the user answered in the composer re-opens as a stale prompt. + */ +export function resolveAsyncRequestUserInputAnswers( + messages: WorkbenchMessage[] +): WorkbenchMessage[] { + const replyByIndex = asyncReplyByMessageIndex(messages) + let changed = false + const resolved = messages.map((message, index) => { + const reply = replyByIndex[index] + if (!reply) return message + const blocks = message.blocks?.map(block => resolveAsyncBlock(block, reply)) + if (!blocks || blocks.every((block, blockIndex) => block === message.blocks![blockIndex])) { + return message + } + changed = true + return { ...message, blocks } + }) + return changed ? resolved : messages +} + +export interface AsyncRequestUserInputReply { + question: string + answer: string +} + +/** + * What each non-blocking answer said, keyed by the reply message id. + * + * An async question is answered by the next user message (that delivery is what + * keeps the model moving), so the transcript recovers the question/answer pairs + * from the conversation instead of from a runtime response. That stays true once + * the question has been answered: sending it only echoes the answer back onto the + * block, so the next user message is still the source of truth. A reply that + * cannot be attributed to its questions is left out, and the message renders as + * typed. + */ +export function resolveAsyncRequestUserInputReplies( + messages: WorkbenchMessage[] +): Map { + const replies = new Map() + messages.forEach((message, index) => { + const questions = asyncQuestionPrompts(message) + if (questions.length === 0) return + const reply = messages.slice(index + 1).find(isUserReply) + if (!reply) return + const rows = pairQuestionsWithReply(questions, reply.content) + if (rows.length > 0) replies.set(reply.id, rows) + }) + return replies +} + +function asyncQuestionPrompts( + message: WorkbenchMessage +): string[] { + return (message.blocks ?? []).flatMap(block => { + if (block.type !== 'tool') return [] + const payload = block.renderPayload + if (!isRequestUserInputPayload(payload) || !isAsyncRequestUserInputPayload(payload)) { + return [] + } + return (payload.questions ?? []).map( + (question, index) => + question.question?.trim() || question.id?.trim() || `question_${index + 1}` + ) + }) +} + +/** + * Answers are composed one per line in question order, so a reply with as many + * lines as questions maps line to question. A single question answers with its + * whole reply; anything else is left unattributed. + */ +function pairQuestionsWithReply( + questions: string[], + reply: string +): AsyncRequestUserInputReply[] { + if (questions.length === 1) { + const answer = reply.trim() + return answer ? [{ question: questions[0], answer }] : [] + } + const lines = reply + .split('\n') + .map(line => line.trim()) + .filter(Boolean) + if (lines.length !== questions.length) return [] + return questions.map((question, index) => ({ question, answer: lines[index] })) +} + +function resolveAsyncBlock( + block: WorkbenchProcessingBlock, + reply: string +): WorkbenchProcessingBlock { + if (block.type !== 'tool') return block + const payload = block.renderPayload + if (!isRequestUserInputPayload(payload)) return block + if (!isAsyncRequestUserInputPayload(payload) || hasRequestUserInputResponse(payload)) { + return block + } + return { + ...block, + status: 'done', + renderPayload: { + ...payload, + response: asyncRequestUserInputResponse(payload, reply), + }, + } +} + +function asyncReplyByMessageIndex(messages: WorkbenchMessage[]): (string | null)[] { + const replies: (string | null)[] = new Array(messages.length).fill(null) + let reply: string | null = null + for (let index = messages.length - 1; index >= 0; index -= 1) { + replies[index] = reply + if (isUserReply(messages[index])) reply = messages[index].content + } + return replies +} + +/** The user message that answers a preceding question, whenever the composer sent it. */ +function isUserReply(message: WorkbenchMessage): boolean { + return message.role === 'user' && Boolean(message.content.trim()) +} + +/** A single free-form reply answers every question the async card asked. */ +function asyncRequestUserInputResponse( + payload: RequestUserInputPayload, + reply: string +): RequestUserInputResponse { + return { + requestId: payload.requestId ?? payload.request_id, + itemId: payload.itemId ?? payload.item_id, + answers: Object.fromEntries( + (payload.questions ?? []).map((question, index) => [ + question.id?.trim() || `question_${index + 1}`, + { answers: [reply] }, + ]) + ), + } +} + +/** Finds the question a runtime answer belongs to so its delivery mode can win. */ +export function findRequestUserInputPayload( + messages: WorkbenchMessage[], + key: string | null +): RequestUserInputPayload | null { + if (!key) return null + for (const message of messages) { + for (const block of message.blocks ?? []) { + if (block.type !== 'tool') continue + const payload = block.renderPayload + if (!isRequestUserInputPayload(payload)) continue + if (requestUserInputPayloadKey(payload) === key) return payload + } + } + return null +} + export function isPendingRequestUserInputBlock( block: ProcessingBlock, hiddenRequestUserInputIds: ReadonlySet = EMPTY_HIDDEN_REQUEST_USER_INPUT_IDS diff --git a/packages/chat-core/src/runtime.ts b/packages/chat-core/src/runtime.ts index 91505fe0f2..7be963df53 100644 --- a/packages/chat-core/src/runtime.ts +++ b/packages/chat-core/src/runtime.ts @@ -339,6 +339,8 @@ export interface RequestUserInputPayload { interaction_kind?: string; approvalKind?: string; approval_kind?: string; + /** `async` marks a Codex question answered by the next user message. */ + delivery?: string; command?: string; cwd?: string; reason?: string; diff --git a/packages/collaboration/src/conversation/MessageList.tsx b/packages/collaboration/src/conversation/MessageList.tsx index 6a0487c4b4..2e86e9a2b1 100644 --- a/packages/collaboration/src/conversation/MessageList.tsx +++ b/packages/collaboration/src/conversation/MessageList.tsx @@ -31,6 +31,10 @@ import { stripPluginWorkspaceResultMarkers } from "@wegent/chat-core/plugin-work import { activityClassNames as cn } from "../issue-detail/activityClassNames"; import { AssistantThinkingIndicator } from "./AssistantThinkingIndicator"; import type { RequestUserInputPayload } from "./RequestUserInputCard"; +import { + resolveAsyncRequestUserInputAnswers, + resolveAsyncRequestUserInputReplies, +} from "@wegent/chat-core/runtime-user-input"; import { getMessagePretextIntrinsicHeight } from "./messagePretextLayout"; import type { AssistantPlanOpenRequest } from "./AssistantPlanCard"; import { @@ -183,10 +187,18 @@ export const MessageList = memo(function MessageList({ const [submittingEditMessageId, setSubmittingEditMessageId] = useState< string | null >(null); - const visibleMessages = useMemo( - () => messages.filter(shouldRenderMessage), + const resolvedMessages = useMemo( + () => resolveAsyncRequestUserInputAnswers(messages), + [messages], + ); + const replyQuestionsByMessageId = useMemo( + () => resolveAsyncRequestUserInputReplies(messages), [messages], ); + const visibleMessages = useMemo( + () => resolvedMessages.filter(shouldRenderMessage), + [resolvedMessages], + ); const runtimeTurnsById = useMemo( () => new Map( @@ -219,7 +231,7 @@ export const MessageList = memo(function MessageList({ : null; const shouldShowWaitingIndicator = isWaitingForAssistant && - !messages.some( + !resolvedMessages.some( (message) => message.role === "assistant" && message.status === "streaming", ); @@ -620,6 +632,7 @@ export const MessageList = memo(function MessageList({ = { export function UserMessage({ services, message, + replyQuestions, onBeforeToggle, onOpenWorkspaceFile, onOpenLocalSkillFile, @@ -102,6 +104,8 @@ export function UserMessage({ }: { services: UserMessageServices message: WorkbenchMessage + /** Questions this message answered, for non-blocking questions answered by the next user message. */ + replyQuestions?: AsyncRequestUserInputReply[] onBeforeToggle?: () => void onOpenWorkspaceFile?: (path: string, options?: WorkspaceFileOpenOptions) => void onOpenLocalSkillFile?: (path: string) => void @@ -277,11 +281,22 @@ export function UserMessage({ shouldCollapse && !isExpanded ? 'max-h-44' : '', ].join(' ')} > - {renderUserContent( - displayContent, - services, - onOpenLocalSkillFile, - onOpenWorkspaceFile + {replyQuestions && replyQuestions.length > 0 ? ( +
+ {replyQuestions.map((row, index) => ( +
+ {row.question} + {row.answer} +
+ ))} +
+ ) : ( + renderUserContent( + displayContent, + services, + onOpenLocalSkillFile, + onOpenWorkspaceFile + ) )} {showGoalRequestBadge && (
diff --git a/packages/collaboration/src/issue-detail/useBrowserConversationActions.test.tsx b/packages/collaboration/src/issue-detail/useBrowserConversationActions.test.tsx index 2e72cf45ff..58b4139989 100644 --- a/packages/collaboration/src/issue-detail/useBrowserConversationActions.test.tsx +++ b/packages/collaboration/src/issue-detail/useBrowserConversationActions.test.tsx @@ -20,6 +20,15 @@ const payload = { { id: 'directory', question: '工作目录?', options: [{ label: '当前目录', value: '/repo' }] }, ], } +const asyncPayload = { + kind: 'request_user_input', + requestId: 'question-async', + itemId: 'tool-async', + delivery: 'async', + questions: [ + { id: 'directory', question: '工作目录?', options: [{ label: '当前目录', value: '/repo' }] }, + ], +} describe('browser runtime question actions', () => { let root: Root @@ -27,19 +36,27 @@ describe('browser runtime question actions', () => { let session: ReturnType let actions: ReturnType let runtime: Parameters[0] - beforeEach(async () => { - globalThis.IS_REACT_ACT_ENVIRONMENT = true + async function mountQuestion( + question: typeof payload = payload, + running = true + ): Promise { + if (root) act(() => root.unmount()) + if (container) container.remove() + if (session) session.stop() container = document.createElement('div') document.body.append(container) root = createRoot(container) runtime = { - work: { sendRuntimeMessage: vi.fn().mockResolvedValue({ accepted: true }) }, + work: { + sendRuntimeMessage: vi.fn().mockResolvedValue({ accepted: true }), + guideRuntimeTask: vi.fn().mockResolvedValue({ accepted: true, turnId: 'turn-1' }), + }, cancel: vi.fn().mockResolvedValue(undefined), } const transcript: RuntimeTranscriptResponse = { runtime: 'codex', workspacePath: '/repo', - running: true, + running, messages: [], turns: [ { @@ -66,7 +83,7 @@ describe('browser runtime question actions', () => { type: 'tool', toolName: 'request_user_input', status: 'pending', - renderPayload: payload, + renderPayload: question, }, }, ], @@ -94,14 +111,18 @@ describe('browser runtime question actions', () => { {actions.error &&
{actions.error}
} void actions.onRequestUserInputIgnore(payload)} + onIgnore={() => void actions.onRequestUserInputIgnore(question)} />
) } await act(async () => root.render()) + } + beforeEach(async () => { + globalThis.IS_REACT_ACT_ENVIRONMENT = true + await mountQuestion() }) afterEach(() => { act(() => root.unmount()) @@ -138,6 +159,16 @@ describe('browser runtime question actions', () => { renderPayload: { response: { requestId: 'question-1' } }, }) }) + it('steers the running turn when a non-blocking question is answered mid-turn', async () => { + await mountQuestion(asyncPayload, true) + await act(async () => button('request-user-input-submit-button').click()) + expect(runtime.work.guideRuntimeTask).toHaveBeenCalledWith({ + address, + message: '/repo', + clientGuidanceId: expect.stringMatching(/^queued-runtime-pane-/), + }) + expect(runtime.work.sendRuntimeMessage).not.toHaveBeenCalled() + }) it('keeps the answer editable and the question pending when runtime rejects it', async () => { vi.mocked(runtime.work.sendRuntimeMessage).mockResolvedValueOnce({ accepted: false, diff --git a/packages/collaboration/src/issue-detail/useBrowserConversationActions.ts b/packages/collaboration/src/issue-detail/useBrowserConversationActions.ts index 52cb31539b..955ae8f008 100644 --- a/packages/collaboration/src/issue-detail/useBrowserConversationActions.ts +++ b/packages/collaboration/src/issue-detail/useBrowserConversationActions.ts @@ -5,7 +5,10 @@ import type { RuntimeTaskAddress, } from '@wegent/chat-core/runtime' import { + findRequestUserInputPayload, + isAsyncRequestUserInputPayload, requestUserInputPayloadKey, + requestUserInputResponseKey, requestUserInputResponseText, } from '@wegent/chat-core/runtime-user-input' import type { createRuntimeConversationSession } from '@wegent/chat-core/runtime-conversation-session' @@ -17,7 +20,7 @@ import { retryRuntimeConversation } from '../execution/retryRuntimeConversation' type Session = ReturnType type QuestionRuntime = Pick & { - work: Pick + work: Pick } /** Bind the PC question cards to the addressed runtime; rejected answers remain editable. */ @@ -92,12 +95,34 @@ export function useBrowserConversationActions( pending.current.add(session) setError(null) try { + const message = requestUserInputResponseText(response) + // Non-blocking Codex questions carry no runtime request to answer, so the + // answer only travels as the next user message. + const appendUserMessage = isAsyncRequestUserInputPayload( + findRequestUserInputPayload( + session.getSnapshot().messages, + requestUserInputResponseKey(response) + ) + ) + // Codex steers the turn that is still running and starts a new one once + // the model has stopped, so the answer is never rejected mid-turn. + if (appendUserMessage && session.getSnapshot().running) { + const guided = await runtime.work.guideRuntimeTask({ + address, + message, + clientGuidanceId: `queued-runtime-pane-${Date.now()}`, + }) + if (guided.accepted === false || guided.success === false) + throw new Error(guided.error || translate('todo.send_failed')) + return true + } const accepted = await runtime.work.sendRuntimeMessage({ address, - message: requestUserInputResponseText(response), - requestUserInputResponse: response, + message, + ...(appendUserMessage ? {} : { requestUserInputResponse: response }), }) if (!accepted.accepted) throw new Error(accepted.error || translate('todo.send_failed')) + if (appendUserMessage) return true session.applyUserInputResponse(response) return true } catch (cause) { diff --git a/wework/e2e/desktop/modules/desktop-server.mjs b/wework/e2e/desktop/modules/desktop-server.mjs index d4b398ac2e..c95249f07f 100644 --- a/wework/e2e/desktop/modules/desktop-server.mjs +++ b/wework/e2e/desktop/modules/desktop-server.mjs @@ -220,6 +220,17 @@ import { REQUEST_USER_INPUT_COMPLETION_TEXT, REQUEST_USER_INPUT_PROMPT, REQUEST_USER_INPUT_QUESTION, + REQUEST_USER_INPUT_ASYNC_ANSWER, + REQUEST_USER_INPUT_ASYNC_CALL_ID, + REQUEST_USER_INPUT_ASYNC_COMPLETION_TEXT, + REQUEST_USER_INPUT_ASYNC_PROMPT, + REQUEST_USER_INPUT_ASYNC_QUESTION, + REQUEST_USER_INPUT_ASYNC_WAITING_TEXT, + REQUEST_USER_INPUT_ASYNC_BUSY_CALL_ID, + REQUEST_USER_INPUT_ASYNC_BUSY_COMPLETION_TEXT, + REQUEST_USER_INPUT_ASYNC_BUSY_PROMPT, + REQUEST_USER_INPUT_ASYNC_BUSY_QUESTION, + REQUEST_USER_INPUT_ASYNC_BUSY_ANSWER, RETRY_COMPLETION_TEXT, RETRY_CONTINUATION_PROMPT, RETRY_FAILURE_TEXT, @@ -362,6 +373,28 @@ function toolOutputText(request, callId) { return findOutput(request.input ?? []) } +/** Text of the user-authored turns codex forwarded to the model, in order. */ +function inputUserMessageTexts(request) { + const texts = [] + const visit = value => { + if (Array.isArray(value)) { + value.forEach(visit) + return + } + if (!value || typeof value !== 'object') return + if (value.type === 'message' && value.role === 'user') { + const content = Array.isArray(value.content) ? value.content : [value.content] + for (const part of content) { + const text = typeof part === 'string' ? part : part?.text + if (typeof text === 'string' && text.trim()) texts.push(text) + } + } + Object.values(value).forEach(visit) + } + visit(request.input) + return texts +} + function pluginWorkspacePublishCommand(body) { const context = findNestedString(body, value => value.includes(PLUGIN_WORKSPACE_PUBLISH_COMMAND_PREFIX) @@ -524,6 +557,9 @@ class DesktopE2EServer { this.requestUserInputResponseWritten = new Promise(resolvePromise => { this.resolveRequestUserInputResponseWritten = resolvePromise }) + this.asyncRequestUserInputKeepsRunningRelease = new Promise(resolvePromise => { + this.releaseAsyncRequestUserInputKeepsRunning = resolvePromise + }) this.taskPlanCompletionRelease = new Promise(resolvePromise => { this.releaseTaskPlanCompletion = resolvePromise }) @@ -812,6 +848,7 @@ class DesktopE2EServer { 'fork_provider_follow_up', 'task_plan', 'request_user_input', + 'request_user_input_async', 'mcp_elicitation', 'window_lifecycle', 'background_completion_restore', @@ -969,6 +1006,10 @@ class DesktopE2EServer { return this.guard(this.requestUserInputResponseWritten) } + releaseAsyncRequestUserInputHold() { + this.releaseAsyncRequestUserInputKeepsRunning() + } + releaseTaskPlanResponse() { this.releaseTaskPlanCompletion() return this.guard(this.taskPlanCompletionWritten) @@ -3744,6 +3785,75 @@ class DesktopE2EServer { return } + if (this.scenario === 'request_user_input_async') { + this.recordScenarioRequest('request_user_input_async', modelRequest) + const serializedBody = JSON.stringify(body) + // Two variants share this scenario: a question answered after the turn + // settles, and one answered while the model is still working. + const keepsRunning = serializedBody.includes(REQUEST_USER_INPUT_ASYNC_BUSY_PROMPT) + const questionText = keepsRunning + ? REQUEST_USER_INPUT_ASYNC_BUSY_QUESTION + : REQUEST_USER_INPUT_ASYNC_QUESTION + const callId = keepsRunning + ? REQUEST_USER_INPUT_ASYNC_BUSY_CALL_ID + : REQUEST_USER_INPUT_ASYNC_CALL_ID + const toolOutput = toolOutputText(body, callId) + if (toolOutput !== null) { + // The non-blocking tool acknowledges immediately and never carries the + // answer, so the chosen option must arrive as the next user message. + assert.deepEqual( + JSON.parse(toolOutput), + { accepted: true }, + 'The async request-user-input tool must acknowledge without an answer' + ) + if (keepsRunning) { + // Keep the turn running until the test has answered the live question, + // then let the model finish with its completion message. + await this.asyncRequestUserInputKeepsRunningRelease + this.writeSse(response, [ + responseCreated(responseId), + assistantMessage(REQUEST_USER_INPUT_ASYNC_BUSY_COMPLETION_TEXT), + responseCompleted(responseId), + ]) + return + } + const answered = inputUserMessageTexts(body).some(text => + text.includes(REQUEST_USER_INPUT_ASYNC_ANSWER) + ) + this.writeSse(response, [ + responseCreated(responseId), + assistantMessage( + answered + ? REQUEST_USER_INPUT_ASYNC_COMPLETION_TEXT + : REQUEST_USER_INPUT_ASYNC_WAITING_TEXT + ), + responseCompleted(responseId), + ]) + return + } + assert.ok( + keepsRunning || serializedBody.includes(REQUEST_USER_INPUT_ASYNC_PROMPT), + 'The real Codex request did not contain the async request-user-input prompt' + ) + const tool = selectTool(body, 'request_user_input_async', { + questions: [ + { + title: questionText, + options: [ + 'Minimal', + keepsRunning ? REQUEST_USER_INPUT_ASYNC_BUSY_ANSWER : REQUEST_USER_INPUT_ASYNC_ANSWER, + ], + }, + ], + }) + this.writeSse(response, [ + responseCreated(responseId), + ...functionCall(callId, tool.name, tool.arguments), + responseCompleted(responseId), + ]) + return + } + if (this.scenario === 'mcp_elicitation') { this.recordScenarioRequest('mcp_elicitation', modelRequest) const requestNumber = this.scenarioRequests.get('mcp_elicitation').length @@ -5885,4 +5995,4 @@ class DesktopE2EServer { } } -export { DesktopE2EServer } +export { DesktopE2EServer, inputUserMessageTexts } diff --git a/wework/e2e/desktop/modules/shared.mjs b/wework/e2e/desktop/modules/shared.mjs index 82f103ae63..fd1d8d594c 100644 --- a/wework/e2e/desktop/modules/shared.mjs +++ b/wework/e2e/desktop/modules/shared.mjs @@ -112,6 +112,21 @@ const REQUEST_USER_INPUT_PROMPT = 'WEWORK_DESKTOP_E2E_REQUEST_INPUT: ask which implementation direction to use.' const REQUEST_USER_INPUT_QUESTION = 'Which implementation direction should be used?' const REQUEST_USER_INPUT_COMPLETION_TEXT = 'WEWORK_DESKTOP_E2E_REQUEST_INPUT_COMPLETE' +const REQUEST_USER_INPUT_ASYNC_PROMPT = + 'WEWORK_DESKTOP_E2E_REQUEST_INPUT_ASYNC: ask which implementation direction to use without blocking the turn.' +const REQUEST_USER_INPUT_ASYNC_QUESTION = 'Which async implementation direction should be used?' +const REQUEST_USER_INPUT_ASYNC_ANSWER = 'Complete' +const REQUEST_USER_INPUT_ASYNC_CALL_ID = 'wework-e2e-request-user-input-async' +const REQUEST_USER_INPUT_ASYNC_WAITING_TEXT = 'WEWORK_DESKTOP_E2E_REQUEST_INPUT_ASYNC_WAITING' +const REQUEST_USER_INPUT_ASYNC_COMPLETION_TEXT = 'WEWORK_DESKTOP_E2E_REQUEST_INPUT_ASYNC_COMPLETE' +const REQUEST_USER_INPUT_ASYNC_BUSY_PROMPT = + 'WEWORK_DESKTOP_E2E_REQUEST_INPUT_ASYNC_BUSY: ask which implementation direction to use while the turn keeps running.' +const REQUEST_USER_INPUT_ASYNC_BUSY_QUESTION = + 'Which implementation direction should be used while the turn keeps running?' +const REQUEST_USER_INPUT_ASYNC_BUSY_ANSWER = 'Proceed' +const REQUEST_USER_INPUT_ASYNC_BUSY_CALL_ID = 'wework-e2e-request-user-input-async-busy' +const REQUEST_USER_INPUT_ASYNC_BUSY_COMPLETION_TEXT = + 'WEWORK_DESKTOP_E2E_REQUEST_INPUT_ASYNC_BUSY_COMPLETE' const MCP_ELICITATION_PROMPT = 'WEWORK_DESKTOP_E2E_MCP_ELICITATION: confirm the inner-site access audience.' const MCP_ELICITATION_COMPLETION_TEXT = 'WEWORK_DESKTOP_E2E_MCP_ELICITATION_COMPLETE' @@ -1583,6 +1598,17 @@ export { REQUEST_USER_INPUT_PROMPT, REQUEST_USER_INPUT_QUESTION, REQUEST_USER_INPUT_COMPLETION_TEXT, + REQUEST_USER_INPUT_ASYNC_PROMPT, + REQUEST_USER_INPUT_ASYNC_QUESTION, + REQUEST_USER_INPUT_ASYNC_ANSWER, + REQUEST_USER_INPUT_ASYNC_CALL_ID, + REQUEST_USER_INPUT_ASYNC_WAITING_TEXT, + REQUEST_USER_INPUT_ASYNC_COMPLETION_TEXT, + REQUEST_USER_INPUT_ASYNC_BUSY_PROMPT, + REQUEST_USER_INPUT_ASYNC_BUSY_QUESTION, + REQUEST_USER_INPUT_ASYNC_BUSY_ANSWER, + REQUEST_USER_INPUT_ASYNC_BUSY_CALL_ID, + REQUEST_USER_INPUT_ASYNC_BUSY_COMPLETION_TEXT, MCP_ELICITATION_PROMPT, MCP_ELICITATION_COMPLETION_TEXT, MCP_ELICITATION_ACCEPTED_MARKER, diff --git a/wework/e2e/desktop/modules/task-flow-main.mjs b/wework/e2e/desktop/modules/task-flow-main.mjs index 92218b538e..65e224fbc3 100644 --- a/wework/e2e/desktop/modules/task-flow-main.mjs +++ b/wework/e2e/desktop/modules/task-flow-main.mjs @@ -277,6 +277,7 @@ import { } from './shared.mjs' import { + verifyAsyncRequestUserInput, verifyBackgroundCompletionRestore, verifyCompletedTurnFork, verifyForkProviderModelPreservation, @@ -2183,6 +2184,12 @@ source = ${JSON.stringify(staleBundledMarketplacePath)}` if (shouldRunDesktopCheckpoint('priority-filter')) { phase = 'priority-filter' await verifyPriorityFilter({ composerSelector: ACTIVE_COMPOSER_SELECTOR, control }) + phase = 'request-user-input-async' + await verifyAsyncRequestUserInput({ + composerSelector: ACTIVE_COMPOSER_SELECTOR, + control, + executorHome, + }) phase = 'runtime-task-order-unread' await verifyRuntimeTaskOrderAndUnreadVisibility({ composerSelector: ACTIVE_COMPOSER_SELECTOR, diff --git a/wework/e2e/desktop/modules/task-state-flows.mjs b/wework/e2e/desktop/modules/task-state-flows.mjs index fe7f340648..94d15e6d7c 100644 --- a/wework/e2e/desktop/modules/task-state-flows.mjs +++ b/wework/e2e/desktop/modules/task-state-flows.mjs @@ -5,6 +5,8 @@ import { waitForSnapshot, } from './conversation-layout.mjs' +import { inputUserMessageTexts } from './desktop-server.mjs' + import { assertConversationTextOrder } from './conversation-navigation.mjs' import { ensureTaskRowVisible } from './memory-tool-flows.mjs' @@ -31,14 +33,26 @@ import { REQUEST_USER_INPUT_COMPLETION_TEXT, REQUEST_USER_INPUT_PROMPT, REQUEST_USER_INPUT_QUESTION, + REQUEST_USER_INPUT_ASYNC_ANSWER, + REQUEST_USER_INPUT_ASYNC_BUSY_ANSWER, + REQUEST_USER_INPUT_ASYNC_BUSY_COMPLETION_TEXT, + REQUEST_USER_INPUT_ASYNC_BUSY_PROMPT, + REQUEST_USER_INPUT_ASYNC_BUSY_QUESTION, + REQUEST_USER_INPUT_ASYNC_COMPLETION_TEXT, + REQUEST_USER_INPUT_ASYNC_PROMPT, + REQUEST_USER_INPUT_ASYNC_QUESTION, + REQUEST_USER_INPUT_ASYNC_WAITING_TEXT, RUNNING_FORK_COMPLETION_TEXT, RUNNING_FORK_FOLLOW_UP_PROMPT, WORKBENCH_READY_TIMEOUT_MS, assert, join, + mkdir, readFile, + rm, selectE2EModel, sendPromptUntilScenarioRequest, + writeFile, withTimeout, } from './shared.mjs' @@ -214,6 +228,175 @@ async function verifyPriorityFilter({ composerSelector, control }) { } } +/** + * Codex's non-blocking `request_user_input_async` question must render as the + * interactive card and be answered with the next user message. + */ +async function setCatalogOverride(executorHome, slug, fields) { + const capabilities = join(executorHome, 'capabilities') + await mkdir(capabilities, { recursive: true }) + await writeFile( + join(capabilities, 'model-catalog-overrides.json'), + `${JSON.stringify([{ slug, fields }], null, 2)}\n` + ) + // Codex caches the served catalog; drop it so the new entry is picked up. + await rm(join(executorHome, 'codex', 'models_cache.json'), { force: true }) +} + +export async function verifyAsyncRequestUserInput({ composerSelector, control, executorHome }) { + // Codex only exposes its non-blocking question tool to models that declare it, + // which the cloud catalog does in production. Declare it for the E2E model so + // the regression exercises the real app-server path. + await setCatalogOverride(executorHome, DEFAULT_MODEL_ID, { + experimental_supported_tools: ['request_user_input_async'], + }) + control.setScenario('request_user_input_async') + await control.command('click', '[data-testid="new-chat-button"]') + await control.command('waitFor', composerSelector, { + timeoutMs: WORKBENCH_READY_TIMEOUT_MS, + }) + await selectE2EModel(control, DEFAULT_MODEL_ID, DEFAULT_MODEL_LABEL) + const composerSnapshot = JSON.parse(await control.command('snapshot', 'body')) + if (composerSnapshot.testIds.includes('plan-mode-pill')) { + await control.command('click', '[data-testid="cancel-plan-mode-button"]') + } + await waitForSnapshot( + control, + snapshot => !snapshot.testIds.includes('plan-mode-pill'), + 'The async request-user-input regression must run in Default mode' + ) + await sendPromptUntilScenarioRequest( + control, + composerSelector, + REQUEST_USER_INPUT_ASYNC_PROMPT, + 'request_user_input_async' + ) + + // The question used to be flattened into the assistant's final text, so the + // interactive card never appeared. + await control.command('waitFor', '[data-testid="request-user-input-card"]', { + text: REQUEST_USER_INPUT_ASYNC_QUESTION, + visible: true, + stableMs: COMPOSER_READY_STABILITY_MS, + timeoutMs: DEFAULT_STEP_TIMEOUT_MS, + }) + await control.command('waitFor', '[data-testid="message-assistant"]', { + text: REQUEST_USER_INPUT_ASYNC_WAITING_TEXT, + visible: true, + stableMs: COMPOSER_READY_STABILITY_MS, + timeoutMs: DEFAULT_STEP_TIMEOUT_MS, + }) + await waitForWorkbenchDebugState( + control, + snapshot => snapshot.pane?.status?.isBusy === false, + 'The async question turn did not settle before it was answered' + ) + await captureVerificationScreenshot(control, 'request-user-input-async-01-question-card.png') + + await control.command('click', '[data-testid="request-user-input-option-question_1-1"]') + await control.command('waitFor', '[data-testid="message-user"]', { + text: REQUEST_USER_INPUT_ASYNC_ANSWER, + visible: true, + stableMs: COMPOSER_READY_STABILITY_MS, + timeoutMs: DEFAULT_STEP_TIMEOUT_MS, + }) + // The answer travels as a normal user message, so the reply has to show the + // question it answered instead of reading like an unprompted message. + await control.command('waitFor', '[data-testid="user-message-question-reply"]', { + text: REQUEST_USER_INPUT_ASYNC_QUESTION, + visible: true, + stableMs: COMPOSER_READY_STABILITY_MS, + timeoutMs: DEFAULT_STEP_TIMEOUT_MS, + }) + await control.command('waitFor', '[data-testid="message-assistant"]', { + text: REQUEST_USER_INPUT_ASYNC_COMPLETION_TEXT, + visible: true, + stableMs: COMPOSER_READY_STABILITY_MS, + timeoutMs: DEFAULT_STEP_TIMEOUT_MS, + }) + assert.equal( + Number(await control.command('getElementCount', '[data-testid="request-user-input-card"]')), + 0, + 'The answered async question must not stay interactive' + ) + await captureVerificationScreenshot(control, 'request-user-input-async-02-answered.png') + + // The non-blocking question must also be answerable while the model keeps + // working, otherwise the "async" contract only holds once the turn settles. + await sendPromptUntilScenarioRequest( + control, + composerSelector, + REQUEST_USER_INPUT_ASYNC_BUSY_PROMPT, + 'request_user_input_async' + ) + await control.command('waitFor', '[data-testid="request-user-input-card"]', { + text: REQUEST_USER_INPUT_ASYNC_BUSY_QUESTION, + visible: true, + stableMs: COMPOSER_READY_STABILITY_MS, + timeoutMs: DEFAULT_STEP_TIMEOUT_MS, + }) + await waitForWorkbenchDebugState( + control, + snapshot => snapshot.pane?.status?.isBusy === true, + 'The model stopped working while the async question waited to be answered' + ) + await captureVerificationScreenshot( + control, + 'request-user-input-async-03-question-while-running.png' + ) + + await control.command('click', '[data-testid="request-user-input-option-question_1-1"]') + // Hold the turn open until the answer has provably landed, so the assertion + // covers delivery while the model is still working rather than after it stops. + await control.command('waitFor', '[data-testid="message-user"]', { + text: REQUEST_USER_INPUT_ASYNC_BUSY_ANSWER, + visible: true, + stableMs: COMPOSER_READY_STABILITY_MS, + timeoutMs: DEFAULT_STEP_TIMEOUT_MS, + }) + await control.command('waitFor', '[data-testid="user-message-question-reply"]', { + text: REQUEST_USER_INPUT_ASYNC_BUSY_QUESTION, + visible: true, + stableMs: COMPOSER_READY_STABILITY_MS, + timeoutMs: DEFAULT_STEP_TIMEOUT_MS, + }) + await waitForWorkbenchDebugState( + control, + snapshot => snapshot.pane?.status?.isBusy === true, + 'The async answer only landed after the model stopped working' + ) + await withTimeout( + control.releaseAsyncRequestUserInputHold(), + DEFAULT_STEP_TIMEOUT_MS, + 'Timed out releasing the running async question turn' + ) + await control.command('waitFor', '[data-testid="message-assistant"]', { + text: REQUEST_USER_INPUT_ASYNC_BUSY_COMPLETION_TEXT, + visible: true, + stableMs: COMPOSER_READY_STABILITY_MS, + timeoutMs: DEFAULT_STEP_TIMEOUT_MS, + }) + assert.ok( + control.modelRequests.some( + request => + request.scenario === 'request_user_input_async' && + inputUserMessageTexts(request.body).some(text => + text.includes(REQUEST_USER_INPUT_ASYNC_BUSY_ANSWER) + ) + ), + 'The async answer never reached the model while the turn kept running' + ) + assert.equal( + Number(await control.command('getElementCount', '[data-testid="request-user-input-card"]')), + 0, + 'The async question answered while running must not stay interactive' + ) + await captureVerificationScreenshot( + control, + 'request-user-input-async-04-answered-while-running.png' + ) +} + async function findRuntimeTaskSortableListSelector(control, taskRowTestIds) { const snapshot = JSON.parse(await control.command('snapshot', 'body')) const candidateTestIds = snapshot.testIds.filter( diff --git a/wework/src/components/chat/SharedMessageList.test.tsx b/wework/src/components/chat/SharedMessageList.test.tsx index ec4b2d7056..ac1ce08e30 100644 --- a/wework/src/components/chat/SharedMessageList.test.tsx +++ b/wework/src/components/chat/SharedMessageList.test.tsx @@ -22,6 +22,40 @@ const assistantMessage: WorkbenchMessage = { status: 'done', createdAt: '2026-09-17T00:00:01Z', } +const asyncQuestion: WorkbenchMessage = { + id: 'assistant-question', + role: 'assistant', + content: '', + status: 'done', + createdAt: '2026-09-17T00:00:01Z', + blocks: [ + { + id: 'request-user-input-call-question', + subtaskId: 'subtask-1', + type: 'tool', + toolName: 'request_user_input', + status: 'pending', + createdAt: 1, + renderPayload: { + kind: 'request_user_input', + delivery: 'async', + itemId: 'call-question', + questions: [ + { id: 'q1', question: '晴天还是雨天?' }, + { id: 'q2', question: '早上还是晚上?' }, + { id: 'q3', question: '猫还是狗?' }, + ], + }, + }, + ], +} +const asyncReply: WorkbenchMessage = { + id: 'user-reply', + role: 'user', + content: '晴天\n早上\n猫', + status: 'done', + createdAt: '2026-09-17T00:00:02Z', +} function renderList(overrides: Partial = {}) { const onEditLastUserMessage = vi.fn().mockResolvedValue(true) @@ -129,4 +163,17 @@ describe('full shared message list in a browser host', () => { expect(reference).toHaveAttribute('tabindex', '-1') expect(reference).toHaveTextContent('spec') }) + + test('shows what a non-blocking answer replied to', () => { + renderList({ messages: [asyncQuestion, asyncReply] }) + const reply = within(screen.getByTestId('user-message-question-reply')) + const answered = ['晴天还是雨天?', '晴天', '早上还是晚上?', '早上', '猫还是狗?', '猫'] + answered.forEach(text => expect(reply.getByText(text)).toBeInTheDocument()) + }) + + test('renders an ordinary user message without a question reference', () => { + renderList() + expect(screen.queryByTestId('user-message-question-reply')).not.toBeInTheDocument() + expect(screen.getByTestId('user-message-content')).toHaveTextContent('Original prompt') + }) }) diff --git a/wework/src/components/layout/useWorkbenchPaneSession.ts b/wework/src/components/layout/useWorkbenchPaneSession.ts index 946cca0845..637309664c 100644 --- a/wework/src/components/layout/useWorkbenchPaneSession.ts +++ b/wework/src/components/layout/useWorkbenchPaneSession.ts @@ -42,6 +42,8 @@ import { persistAttachmentReferences } from '@/lib/attachments' import { localRuntimeAttachments, remoteAttachmentIds } from '@/lib/runtime-attachments' import { applyRequestUserInputResponseToBlock, + findRequestUserInputPayload, + isAsyncRequestUserInputPayload, requestUserInputPayloadKey, requestUserInputResponseKey, requestUserInputResponseText, @@ -1429,72 +1431,6 @@ export function useWorkbenchPaneSession({ ] ) - const sendRequestUserInputResponse = useCallback( - async ( - response: RequestUserInputResponse, - options: SendRequestUserInputResponseOptions = {} - ): Promise => { - if (!currentRuntimeTask) return false - - const message = requestUserInputResponseText(response) - const requestUserInputKey = requestUserInputResponseKey(response) - const runtimeModelOverride = options.forceDefaultCollaborationMode - ? { collaborationMode: 'default' } - : undefined - if (options.forceDefaultCollaborationMode) { - projectChat.setSelectedModelOption('collaborationMode', 'default') - } - const appendedUserMessage = options.appendUserMessage - ? createRuntimeUserMessage(message) - : null - if (appendedUserMessage) { - dispatchMessages({ type: 'user_added', message: appendedUserMessage }) - } - if (requestUserInputKey) { - setAnsweredRequestUserInputIds(current => { - if (current.has(requestUserInputKey)) return current - const next = new Set(current) - next.add(requestUserInputKey) - return next - }) - } - applyLocalRequestUserInputResponse(response) - const runtimeModelFields = options.appendUserMessage - ? getRuntimeModelFields(runtimeModelOverride) - : {} - const additionalContext = readRuntimeTerminalAdditionalContext(currentRuntimeTask) - const sent = await sendRuntimePaneMessage({ - address: currentRuntimeTask, - message, - ...(appendedUserMessage ? { clientUserMessageId: appendedUserMessage.id } : {}), - ...runtimeModelFields, - ...(options.appendUserMessage ? {} : { requestUserInputResponse: response }), - ...(additionalContext ? { additionalContext } : {}), - }) - if (sent) { - markRuntimeTerminalAdditionalContextDelivered(additionalContext) - } else { - if (requestUserInputKey) { - setAnsweredRequestUserInputIds(current => { - if (!current.has(requestUserInputKey)) return current - const next = new Set(current) - next.delete(requestUserInputKey) - return next - }) - } - } - return sent - }, - [ - applyLocalRequestUserInputResponse, - currentRuntimeTask, - dispatchMessages, - getRuntimeModelFields, - projectChat, - sendRuntimePaneMessage, - ] - ) - const editLastUserMessageInPane = useCallback( async (message: WorkbenchMessage, content: string): Promise => { const submittedContent = content.trim() @@ -1839,7 +1775,7 @@ export function useWorkbenchPaneSession({ }, [runtimeTaskLoadTarget]) const sendQueuedMessageAsGuidance = useCallback( - async (queuedMessage: RuntimePaneQueuedMessage, forceActiveTurn = false) => { + async (queuedMessage: RuntimePaneQueuedMessage, forceActiveTurn = false): Promise => { const id = queuedMessage.id queuedMessageBusyBlockSnapshotsRef.current.delete(id) if (!currentRuntimeTask) { @@ -1850,11 +1786,11 @@ export function useWorkbenchPaneSession({ : message ) ) - return + return false } - if (currentRuntimeTask.runtime && currentRuntimeTask.runtime !== 'codex') return + if (currentRuntimeTask.runtime && currentRuntimeTask.runtime !== 'codex') return false - if (queuedMessage.status === 'sending') return + if (queuedMessage.status === 'sending') return false setError(null) if (!readCurrentPaneBusy() && !forceActiveTurn) { @@ -1869,7 +1805,7 @@ export function useWorkbenchPaneSession({ const result = await sendRuntimeMessage(queuedMessage) if (result.queued) { retainRuntimeQueuedMessage(queuedMessage, result.queuePosition) - return + return true } setQueuedMessages(messages => result.accepted @@ -1880,6 +1816,7 @@ export function useWorkbenchPaneSession({ : message ) ) + return result.accepted } catch (error) { console.error('[Wework] Queued runtime message send failed', { id, @@ -1892,8 +1829,8 @@ export function useWorkbenchPaneSession({ : message ) ) + return false } - return } try { @@ -1938,32 +1875,32 @@ export function useWorkbenchPaneSession({ queuedMessage ) } + return true } - if (!result.sent) { - removeOptimisticRuntimeConversationGuidance(currentRuntimeTask, id) - if (takeInterruptedRuntimeConversationGuidance(currentRuntimeTask, id)) { - setQueuedMessages(messages => messages.filter(message => message.id !== id)) - return - } - setQueuedMessages(messages => - messages.map(message => - message.id === id - ? { - ...message, - status: 'failed', - awaitingGuidanceAcceptance: undefined, - notice: undefined, - error: '引导发送失败', - } - : message - ) - ) + removeOptimisticRuntimeConversationGuidance(currentRuntimeTask, id) + if (takeInterruptedRuntimeConversationGuidance(currentRuntimeTask, id)) { + setQueuedMessages(messages => messages.filter(message => message.id !== id)) + return true } + setQueuedMessages(messages => + messages.map(message => + message.id === id + ? { + ...message, + status: 'failed', + awaitingGuidanceAcceptance: undefined, + notice: undefined, + error: '引导发送失败', + } + : message + ) + ) + return false } catch (error) { removeOptimisticRuntimeConversationGuidance(currentRuntimeTask, id) if (takeInterruptedRuntimeConversationGuidance(currentRuntimeTask, id)) { setQueuedMessages(messages => messages.filter(message => message.id !== id)) - return + return true } console.error('[Wework] Queued guidance send failed', { id, @@ -1982,6 +1919,7 @@ export function useWorkbenchPaneSession({ : message ) ) + return false } }, [ @@ -1995,6 +1933,76 @@ export function useWorkbenchPaneSession({ ] ) + const sendRequestUserInputResponse = useCallback( + async ( + response: RequestUserInputResponse, + options: SendRequestUserInputResponseOptions = {} + ): Promise => { + if (!currentRuntimeTask) return false + + const message = requestUserInputResponseText(response) + const requestUserInputKey = requestUserInputResponseKey(response) + const runtimeModelOverride = options.forceDefaultCollaborationMode + ? { collaborationMode: 'default' } + : undefined + if (options.forceDefaultCollaborationMode) { + projectChat.setSelectedModelOption('collaborationMode', 'default') + } + // A question that is answered by the next user message - a non-blocking + // Codex question or an implementation-plan confirmation - is delivered like + // any other user message. Codex then steers the running turn while the model + // is still working and starts a new turn once it has stopped, instead of + // rejecting the answer for arriving mid-turn. + const answerIsUserMessage = + options.appendUserMessage || + isAsyncRequestUserInputPayload( + findRequestUserInputPayload(messagesRef.current, requestUserInputKey) + ) + if (answerIsUserMessage) { + const answerMessage: RuntimePaneQueuedMessage = { + id: `queued-runtime-pane-${Date.now()}-${queuedMessages.length}`, + content: message, + status: 'queued', + createdAt: new Date().toISOString(), + ...getRuntimeModelFields(runtimeModelOverride), + } + setQueuedMessages(messages => [...messages, answerMessage]) + if (!(await sendQueuedMessageAsGuidance(answerMessage))) return false + } else { + const additionalContext = readRuntimeTerminalAdditionalContext(currentRuntimeTask) + const sent = await sendRuntimePaneMessage({ + address: currentRuntimeTask, + message, + requestUserInputResponse: response, + ...(additionalContext ? { additionalContext } : {}), + }) + if (!sent) return false + markRuntimeTerminalAdditionalContextDelivered(additionalContext) + } + if (requestUserInputKey) { + setAnsweredRequestUserInputIds(current => { + if (current.has(requestUserInputKey)) return current + const next = new Set(current) + next.add(requestUserInputKey) + return next + }) + } + applyLocalRequestUserInputResponse(response) + return true + }, + [ + applyLocalRequestUserInputResponse, + currentRuntimeTask, + getRuntimeModelFields, + projectChat, + queuedMessages.length, + sendQueuedMessageAsGuidance, + sendRuntimePaneMessage, + setAnsweredRequestUserInputIds, + setQueuedMessages, + ] + ) + const send: (inputOverride?: string, options?: RuntimePaneSendOptions) => Promise = useCallback( async (inputOverride, options = {}) => { diff --git a/wework/src/components/layout/workspace-panels/TemporaryChatPanel.tsx b/wework/src/components/layout/workspace-panels/TemporaryChatPanel.tsx index d4f570370c..33463babc3 100644 --- a/wework/src/components/layout/workspace-panels/TemporaryChatPanel.tsx +++ b/wework/src/components/layout/workspace-panels/TemporaryChatPanel.tsx @@ -25,7 +25,10 @@ import { ScrollableMessageArea } from '@/components/chat/ScrollableMessageArea' import type { RequestUserInputPayload } from '@/components/chat/RequestUserInputCard' import { applyRequestUserInputResponseToBlock, + findRequestUserInputPayload, + isAsyncRequestUserInputPayload, requestUserInputPayloadKey, + requestUserInputResponseKey, requestUserInputResponseText, } from '@/components/chat/requestUserInputMessages' import type { ChatSubmitOptions, ProjectWorkControls } from '@/components/chat/ChatInput' @@ -778,6 +781,23 @@ export function TemporaryChatPanel({ const submitRequestUserInput = useCallback( async (response: RequestUserInputResponse): Promise => { if (!address) return false + const appendUserMessage = isAsyncRequestUserInputPayload( + findRequestUserInputPayload(messages, requestUserInputResponseKey(response)) + ) + // A non-blocking Codex question carries no runtime request to answer, so it + // is answered by the next user message: Codex steers the running turn while + // the model is still working and starts a new turn once it has stopped. + if (appendUserMessage) { + const answerMessage: RuntimePaneQueuedMessage = { + id: `queued-side-chat-${Date.now()}-${queuedMessages.length}`, + content: requestUserInputResponseText(response), + status: 'queued', + createdAt: new Date().toISOString(), + ...selectedModelFields, + } + conversationQueue.enqueue(answerMessage) + return sendQueuedMessageAsGuidance(answerMessage) + } const sent = await sendRuntimePaneMessage({ address, message: requestUserInputResponseText(response), @@ -792,7 +812,16 @@ export function TemporaryChatPanel({ ) return true }, - [address, runtimeContext, sendRuntimePaneMessage] + [ + address, + conversationQueue, + messages, + queuedMessages.length, + runtimeContext, + selectedModelFields, + sendQueuedMessageAsGuidance, + sendRuntimePaneMessage, + ] ) const ignoreRequestUserInput = useCallback(