From 50e4f3617499c3fe1241d43581f775020f1e52c6 Mon Sep 17 00:00:00 2001 From: Bob Lee Date: Mon, 21 Sep 2026 22:55:53 +0800 Subject: [PATCH 1/2] feat(remote): add host-owned mobile message queues --- docs/interactive-capabilities/README.md | 4 +- .../technical/product-control-open-audit.json | 1 + .../technical/tauri-command-map.json | 22 +- src/apps/cli/src/peer_host/commands/dialog.rs | 14 + src/apps/cli/src/peer_host/commands/mod.rs | 2 + src/apps/desktop/src/api/agentic_api.rs | 12 + src/apps/desktop/src/lib.rs | 1 + src/crates/assembly/core/AGENTS.md | 6 + .../coordination/host_message_queue.rs | 578 ++++++++++++++++++ .../coordination/host_message_queue_tests.rs | 573 +++++++++++++++++ .../src/agentic/coordination/scheduler.rs | 202 +++++- .../service/remote_connect/remote_server.rs | 9 + .../core/src/service_agent_runtime.rs | 31 + .../generated/remote-surface-registry.json | 16 +- .../src/remote_surface/capabilities.rs | 4 + .../src/remote_surface/table.rs | 1 + .../contracts/runtime-ports/src/agent_api.rs | 11 + .../runtime-ports/src/dialog_queue.rs | 135 ++++ src/crates/contracts/runtime-ports/src/lib.rs | 5 + .../execution/agent-runtime/src/runtime.rs | 20 + src/crates/execution/agent-runtime/src/sdk.rs | 7 + .../src/remote_connect.rs | 20 + .../tests/remote_connect_contracts.rs | 18 + src/mobile-web/AGENTS.md | 1 + src/mobile-web/README.md | 31 + src/mobile-web/package.json | 1 + .../src/components/ChatComposerBar.tsx | 6 +- .../src/components/MobileHostQueue.tsx | 48 ++ src/mobile-web/src/i18n/messages.ts | 54 ++ src/mobile-web/src/pages/ChatPage.tsx | 12 +- .../src/services/RemoteSessionManager.ts | 32 + src/mobile-web/src/styles/host-queue.scss | 12 + src/mobile-web/tests/fixtures/host-queue.tsx | 28 + .../tests/host-dialog-queue-browser.test.mjs | 54 ++ .../tests/host-dialog-queue.test.mjs | 103 ++++ src/shared/dialog-queue/HostDialogQueue.ts | 246 ++++++++ .../interactive-capabilities/catalog.json | 1 + .../components/HostPendingQueuePanel.tsx | 85 +++ .../components/PendingQueuePanel.test.tsx | 3 + .../components/PendingQueuePanel.tsx | 19 +- .../services/SessionRollbackService.ts | 5 + .../flow-chat-manager/MessageModule.ts | 15 + .../flow-chat-manager/PendingQueueModule.ts | 5 +- .../src/flow_chat/services/hostDialogQueue.ts | 52 ++ .../local/LocalSessionDriver.test.ts | 29 + .../local/LocalSessionDriver.ts | 103 ++-- .../api/generated/remoteSurface.ts | 5 +- .../peer-device/PeerConnectionManager.ts | 3 + src/web-ui/src/locales/en-US/flow-chat.json | 18 + src/web-ui/src/locales/zh-CN/flow-chat.json | 18 + src/web-ui/src/locales/zh-TW/flow-chat.json | 18 + 51 files changed, 2627 insertions(+), 72 deletions(-) create mode 100644 src/crates/assembly/core/src/agentic/coordination/host_message_queue.rs create mode 100644 src/crates/assembly/core/src/agentic/coordination/host_message_queue_tests.rs create mode 100644 src/crates/contracts/runtime-ports/src/dialog_queue.rs create mode 100644 src/mobile-web/src/components/MobileHostQueue.tsx create mode 100644 src/mobile-web/src/styles/host-queue.scss create mode 100644 src/mobile-web/tests/fixtures/host-queue.tsx create mode 100644 src/mobile-web/tests/host-dialog-queue-browser.test.mjs create mode 100644 src/mobile-web/tests/host-dialog-queue.test.mjs create mode 100644 src/shared/dialog-queue/HostDialogQueue.ts create mode 100644 src/web-ui/src/flow_chat/components/HostPendingQueuePanel.tsx create mode 100644 src/web-ui/src/flow_chat/services/hostDialogQueue.ts diff --git a/docs/interactive-capabilities/README.md b/docs/interactive-capabilities/README.md index c280bb8829..89c1843520 100644 --- a/docs/interactive-capabilities/README.md +++ b/docs/interactive-capabilities/README.md @@ -27,9 +27,9 @@ OpenBitFun Playbook currently contains **22 features**, **21 settings pages**, a - Generated per-item interaction audit: `docs/interactive-capabilities/technical/product-control-open-audit.json` - Generated low-level audit map: `docs/interactive-capabilities/technical/tauri-command-map.json` -说明书、网站、搜索和智能体只看“功能 + 设置 + 子能力”。每项子能力都必须引用已注册 Tauri Command 或可解析的源码标记;这些证据不会进入公开目录。当前 **664** 个 Tauri 命令只用于实现覆盖审计。产品 UI 交互源码会在生成和检查时扫描并校验,但不会保存成随普通 UI 改动频繁变化的版本化快照。 +说明书、网站、搜索和智能体只看“功能 + 设置 + 子能力”。每项子能力都必须引用已注册 Tauri Command 或可解析的源码标记;这些证据不会进入公开目录。当前 **665** 个 Tauri 命令只用于实现覆盖审计。产品 UI 交互源码会在生成和检查时扫描并校验,但不会保存成随普通 UI 改动频繁变化的版本化快照。 -Docs, website, search, and agents see only features, settings, and documented sub-capabilities. Every sub-capability must reference a registered Tauri command or a resolvable source marker; evidence is stripped from public projections. The **664** Tauri commands remain implementation-audit evidence only. Product UI interaction sources are scanned and validated during generation and checks, but are not stored as a versioned snapshot that churns with ordinary UI changes. +Docs, website, search, and agents see only features, settings, and documented sub-capabilities. Every sub-capability must reference a registered Tauri command or a resolvable source marker; evidence is stripped from public projections. The **665** Tauri commands remain implementation-audit evidence only. Product UI interaction sources are scanned and validated during generation and checks, but are not stored as a versioned snapshot that churns with ordinary UI changes. ## 控制边界 / Control boundary diff --git a/docs/interactive-capabilities/technical/product-control-open-audit.json b/docs/interactive-capabilities/technical/product-control-open-audit.json index a97fa23eb9..1f985ec111 100644 --- a/docs/interactive-capabilities/technical/product-control-open-audit.json +++ b/docs/interactive-capabilities/technical/product-control-open-audit.json @@ -69,6 +69,7 @@ }, "evidence": [ "command:start_dialog_turn", + "command:manage_dialog_queue", "command:steer_dialog_turn", "command:interrupt_dialog_turn", "command:cancel_dialog_turn", diff --git a/docs/interactive-capabilities/technical/tauri-command-map.json b/docs/interactive-capabilities/technical/tauri-command-map.json index 08386ae0b9..105e46145e 100644 --- a/docs/interactive-capabilities/technical/tauri-command-map.json +++ b/docs/interactive-capabilities/technical/tauri-command-map.json @@ -2,10 +2,10 @@ "schemaVersion": 2, "generatedFrom": "src/shared/interactive-capabilities/catalog.json", "catalogDigest": "6587344a6b5e75a80457ea4ba2670bf37cf3038bdb436089ae402eb7b5a5a025", - "commandCount": 664, + "commandCount": 665, "coverage": { - "commandCount": 664, - "documentedCommandCount": 616, + "commandCount": 665, + "documentedCommandCount": 617, "implementationCommandCount": 48, "implementationDigest": "f6d38a24a70708988cb47ada81d07eccf0668684e8734c9ee5155b9ffa7e3db8" }, @@ -5424,6 +5424,22 @@ "signature": "fn logout_subscription_account( request: SubscriptionProviderRequest, ) -> Result", "remoteWorkspacePolicy": "LocalOnly" }, + { + "id": "manage_dialog_queue", + "moduleId": "agentic", + "capabilityId": "feature.ai-assistant", + "capabilityIds": [ + "feature.ai-assistant" + ], + "documentedItemIds": [ + "feature.ai-assistant:turn-control" + ], + "visibility": "documented", + "rustPath": "api::agentic_api::manage_dialog_queue", + "sourceFile": "src/apps/desktop/src/api/agentic_api.rs", + "signature": "fn manage_dialog_queue( runtime: State<'_, DesktopRuntimeContext>, request: openbitfun_runtime_ports::DialogQueueRequest, ) -> Result", + "remoteWorkspacePolicy": "RemoteRouted" + }, { "id": "mark_announcement_seen", "moduleId": "announcement", diff --git a/src/apps/cli/src/peer_host/commands/dialog.rs b/src/apps/cli/src/peer_host/commands/dialog.rs index 2657318f89..0f1ee7de23 100644 --- a/src/apps/cli/src/peer_host/commands/dialog.rs +++ b/src/apps/cli/src/peer_host/commands/dialog.rs @@ -502,3 +502,17 @@ mod image_attachment_tests { assert!(peer_image_attachments(&json!({"imageContexts": "bad"})).is_err()); } } + +pub(crate) async fn manage_dialog_queue( + state: &PeerHostState, + args: &Value, +) -> Result { + let request = serde_json::from_value(request_value(args).clone()) + .map_err(|e| format!("Invalid queue request: {e}"))?; + let snapshot = state + .agent_runtime + .manage_dialog_queue(request) + .await + .map_err(|e| e.into_message())?; + serde_json::to_value(snapshot).map_err(|e| e.to_string()) +} diff --git a/src/apps/cli/src/peer_host/commands/mod.rs b/src/apps/cli/src/peer_host/commands/mod.rs index 62456c38ea..750205488c 100644 --- a/src/apps/cli/src/peer_host/commands/mod.rs +++ b/src/apps/cli/src/peer_host/commands/mod.rs @@ -161,6 +161,7 @@ pub(crate) fn dispatch<'a>( "get_session_files" => Box::pin(snapshot::get_session_files(state, args)), // Dialog / tools + "manage_dialog_queue" => Box::pin(dialog::manage_dialog_queue(state, args)), "start_dialog_turn" => Box::pin(dialog::start_dialog_turn(state, args)), "cancel_dialog_turn" => Box::pin(dialog::cancel_dialog_turn(state, args)), "start_user_question_interaction" => Box::pin(dialog::start_user_question_interaction(state, args)), @@ -336,6 +337,7 @@ pub(crate) const HANDLED_COMMANDS: &[&str] = &[ "set_external_tool_targets_enabled_command", "set_active_workspace", "start_dialog_turn", + "manage_dialog_queue", "submit_user_answers", "start_user_question_interaction", "subscribe_permission_requests", diff --git a/src/apps/desktop/src/api/agentic_api.rs b/src/apps/desktop/src/api/agentic_api.rs index ef4c9258f7..5e23546f5f 100644 --- a/src/apps/desktop/src/api/agentic_api.rs +++ b/src/apps/desktop/src/api/agentic_api.rs @@ -2472,6 +2472,18 @@ pub async fn ensure_coordinator_session( .map_err(|error| error.to_string()) } +#[tauri::command] +pub async fn manage_dialog_queue( + runtime: State<'_, DesktopRuntimeContext>, + request: openbitfun_runtime_ports::DialogQueueRequest, +) -> Result { + runtime + .agent_runtime() + .manage_dialog_queue(request) + .await + .map_err(|e| e.into_message()) +} + #[tauri::command] pub async fn start_dialog_turn( _app: AppHandle, diff --git a/src/apps/desktop/src/lib.rs b/src/apps/desktop/src/lib.rs index 203b26dbe1..7a1e74bb5d 100644 --- a/src/apps/desktop/src/lib.rs +++ b/src/apps/desktop/src/lib.rs @@ -1286,6 +1286,7 @@ pub async fn run() { api::agentic_api::interrupt_dialog_turn, api::agentic_api::recover_interrupted_dialog_turn, api::agentic_api::steer_dialog_turn, + api::agentic_api::manage_dialog_queue, api::agentic_api::control_deep_review_queue, api::agentic_api::cancel_session, api::agentic_api::set_subagent_timeout, diff --git a/src/crates/assembly/core/AGENTS.md b/src/crates/assembly/core/AGENTS.md index 68d57bbbdb..c14ebb6507 100644 --- a/src/crates/assembly/core/AGENTS.md +++ b/src/crates/assembly/core/AGENTS.md @@ -309,3 +309,9 @@ For remote search ID binding without a live SSH connection: ```bash cargo test --locked -p openbitfun-core --no-default-features --features agent-runtime,git,ssh-remote --lib service::search::remote::identity_tests ``` + +For host-owned user queue admission, cancellation, steering receipts and client disconnects: + +```bash +cargo test --locked -p openbitfun-core --no-default-features --features remote-connect,git --lib host_queue_ +``` diff --git a/src/crates/assembly/core/src/agentic/coordination/host_message_queue.rs b/src/crates/assembly/core/src/agentic/coordination/host_message_queue.rs new file mode 100644 index 0000000000..929953f548 --- /dev/null +++ b/src/crates/assembly/core/src/agentic/coordination/host_message_queue.rs @@ -0,0 +1,578 @@ +//! User queue receipts and management. Execution stays in DialogScheduler. +use super::*; +use openbitfun_runtime_ports::{ + DialogQueueAction, DialogQueueItem, DialogQueueMessage, DialogQueueRequest, + DialogQueueSnapshot, DialogQueueStatus, +}; +use sha2::{Digest, Sha256}; +use std::collections::BTreeMap; + +// Compact receipts must not be silently evicted: that would allow a delayed +// duplicate to execute twice. At the budget boundary fail admission explicitly. +const RECEIPT_LIMIT: usize = 16_384; +const PREVIEW_CHARS: usize = 2_000; + +#[derive(Default)] +pub(super) struct HostQueueState { + sessions: BTreeMap, +} +struct QueueSession { + epoch: String, + revision: u64, + active_turn_id: Option, + entries: BTreeMap, + order: Vec, + operations: BTreeMap, +} +struct Entry { + fingerprint: String, + view: DialogQueueItem, + held: Option, +} +impl Default for QueueSession { + fn default() -> Self { + Self { + epoch: Uuid::new_v4().to_string(), + revision: 0, + active_turn_id: None, + entries: BTreeMap::new(), + order: Vec::new(), + operations: BTreeMap::new(), + } + } +} +fn error(message: impl Into) -> PortError { + PortError::new(PortErrorKind::InvalidRequest, message) +} +fn fingerprint(value: &impl serde::Serialize) -> PortResult { + let bytes = serde_json::to_vec(value).map_err(|e| error(e.to_string()))?; + Ok(format!("{:x}", Sha256::digest(bytes))) +} +impl HostQueueState { + pub(super) fn contains(&self, session: &str, turn: &str) -> bool { + self.sessions + .get(session) + .is_some_and(|s| s.entries.contains_key(turn)) + } + pub(super) fn pending_held(&self, session: &str) -> usize { + self.sessions.get(session).map_or(0, |s| { + s.entries + .values() + .filter(|e| e.held.is_some() && e.view.status.is_pending()) + .count() + }) + } + pub(super) fn admission_valid(&self, session: &str, turn: &str) -> bool { + self.sessions + .get(session) + .and_then(|s| s.entries.get(turn)) + .is_none_or(|e| e.view.status == DialogQueueStatus::Queued) + } + pub(super) fn retire(&mut self, session: &str) -> Vec { + let Some(s) = self.sessions.get_mut(session) else { + return Vec::new(); + }; + s.epoch = Uuid::new_v4().to_string(); + s.revision += 1; + let mut held = Vec::new(); + for e in s.entries.values_mut() { + if e.view.status.is_pending() { + e.view.status = DialogQueueStatus::Cancelled; + if let Some(turn) = e.held.take() { + held.push(turn); + } + } + } + held + } + pub(super) fn cancelled(&mut self, session: &str, turn: &str) { + if let Some(s) = self.sessions.get_mut(session) { + if let Some(e) = s.entries.get_mut(turn) { + e.view.status = DialogQueueStatus::Cancelled; + e.view.reason = None; + e.held = None; + s.revision += 1; + } + } + } + pub(super) fn started(&mut self, session: &str, turn: &str) { + if let Some(s) = self.sessions.get_mut(session) { + if let Some(e) = s.entries.get_mut(turn) { + e.view.status = DialogQueueStatus::Started; + e.view.reason = None; + e.held = None; + s.revision += 1; + } + } + } + pub(super) fn hold(&mut self, session: &str, turn: &QueuedTurn, reason: &str) -> bool { + let Some(s) = self.sessions.get_mut(session) else { + return false; + }; + let Some(e) = turn.turn_id.as_ref().and_then(|id| s.entries.get_mut(id)) else { + return false; + }; + e.view.status = DialogQueueStatus::Blocked; + e.view.reason = Some(reason.to_string()); + e.held = Some(turn.clone()); + s.revision += 1; + true + } + pub(super) fn consumed(&mut self, session: &str, target: &str, injection: &str) { + let Some(s) = self.sessions.get_mut(session) else { + return; + }; + for e in s.entries.values_mut() { + if e.view.status == DialogQueueStatus::SteeringPending + && e.view.target_turn_id.as_deref() == Some(target) + && e.view.steering_id.as_deref() == Some(injection) + { + e.view.status = DialogQueueStatus::Steered; + e.held = None; + s.revision += 1; + break; + } + } + } + pub(super) fn outcome( + &mut self, + session: &str, + turn: &str, + status: TurnOutcomeStatus, + ) -> Vec { + let Some(s) = self.sessions.get_mut(session) else { + return Vec::new(); + }; + let mut retired_injections = Vec::new(); + for (id, e) in &mut s.entries { + if id == turn + && matches!( + e.view.status, + DialogQueueStatus::Started | DialogQueueStatus::Interrupted + ) + { + e.view.status = match status { + TurnOutcomeStatus::Completed => DialogQueueStatus::Completed, + TurnOutcomeStatus::Cancelled => DialogQueueStatus::Cancelled, + TurnOutcomeStatus::Interrupted => DialogQueueStatus::Interrupted, + _ => DialogQueueStatus::Failed, + }; + s.revision += 1; + } + if e.view.status == DialogQueueStatus::SteeringPending + && e.view.target_turn_id.as_deref() == Some(turn) + { + // Outcome is emitted after the execution future has retired; + // no outstanding consumer can inject this entry afterwards. + if let Some(id) = &e.view.steering_id { + retired_injections.push(id.clone()); + } + e.view.status = DialogQueueStatus::Blocked; + e.view.reason = + Some("Steering was not consumed before the target turn ended".into()); + s.revision += 1; + } + } + retired_injections + } +} +impl DialogScheduler { + pub(super) fn hold_managed_queue(&self, session: &str, reason: &str) { + let ids: Vec = self + .queue_state() + .sessions + .get(session) + .map(|s| { + s.entries + .iter() + .filter(|(_, e)| e.view.status == DialogQueueStatus::Queued) + .map(|(id, _)| id.clone()) + .collect() + }) + .unwrap_or_default(); + for id in ids { + if let Some(turn) = remove_queued_turn_by_id(&self.queues, session, &id) { + self.queue_state().hold(session, &turn, reason); + } + } + } + fn queue_state(&self) -> std::sync::MutexGuard<'_, HostQueueState> { + self.host_queue.lock().unwrap_or_else(|e| e.into_inner()) + } + fn queue_snapshot(&self, session: &str, receipt_id: Option<&str>) -> DialogQueueSnapshot { + let mut state = self.queue_state(); + let s = state.sessions.entry(session.to_string()).or_default(); + let active_turn_id = self + .session_manager + .get_session(session) + .and_then(|session| match &session.state { + SessionState::Processing { + current_turn_id, .. + } => Some(current_turn_id.clone()), + _ => None, + }); + if s.active_turn_id != active_turn_id { + s.active_turn_id = active_turn_id.clone(); + s.revision += 1; + } + let items = s + .order + .iter() + .filter_map(|id| s.entries.get(id)) + .filter(|e| e.view.status.is_pending()) + .map(|e| e.view.clone()) + .collect(); + let held = s + .entries + .values() + .filter(|e| e.held.is_some() && e.view.status.is_pending()) + .count(); + DialogQueueSnapshot { + session_id: session.to_string(), + queue_epoch: s.epoch.clone(), + revision: s.revision, + active_turn_id, + items, + capacity: self.queues.max_depth(), + used: self.queues.depth(session) + held, + receipt: receipt_id + .and_then(|id| s.entries.get(id)) + .map(|e| e.view.clone()), + } + } + pub(super) async fn manage_host_queue( + &self, + request: DialogQueueRequest, + ) -> PortResult { + // Once sent to the host, a mutation must outlive a disconnected RPC. + let scheduler = self + .self_ref + .upgrade() + .ok_or_else(|| error("Queue owner unavailable"))?; + tokio::spawn(async move { scheduler.execute_queue_request(request).await }) + .await + .map_err(|e| PortError::new(PortErrorKind::Backend, e.to_string()))? + } + async fn execute_queue_request( + &self, + request: DialogQueueRequest, + ) -> PortResult { + let session = &request.session_id; + openbitfun_core_types::validate_session_id(session).map_err(error)?; + let _admission = self.host_queue_locks.lock(session).await; + if self.session_manager.get_session(session).is_none() { + return Err(PortError::new( + PortErrorKind::NotFound, + "Session is not loaded on the execution host", + )); + } + let snapshot = self.queue_snapshot(session, None); + if !matches!(request.action, DialogQueueAction::List) + && request.queue_epoch.as_deref() != Some(snapshot.queue_epoch.as_str()) + { + return Err(error( + "queue_scope_expired: refresh the host queue; do not automatically resend", + )); + } + match request.action { + DialogQueueAction::List => { + let _guard = self.lock_session_operation(session).await; + Ok(self.queue_snapshot(session, None)) + } + DialogQueueAction::Get { turn_id } => { + let _guard = self.lock_session_operation(session).await; + Ok(self.queue_snapshot(session, Some(&turn_id))) + } + DialogQueueAction::Submit { message } => { + self.submit_host_message(session, message).await + } + action => self.mutate_host_queue(session, action).await, + } + } + async fn submit_host_message( + &self, + session: &str, + message: DialogQueueMessage, + ) -> PortResult { + if message.turn_id.trim().is_empty() + || message.turn_id.len() > 200 + || (message.content.trim().is_empty() && message.attachments.is_empty()) + { + return Err(error( + "A stable turn ID and message content or attachments are required", + )); + } + let digest = fingerprint(&message)?; + { + let mut state = self.queue_state(); + let s = state.sessions.get_mut(session).expect("queue initialized"); + if let Some(e) = s.entries.get(&message.turn_id) { + if e.fingerprint != digest { + return Err(error("idempotency_conflict: message payload changed")); + } + drop(state); + return Ok(self.queue_snapshot(session, Some(&message.turn_id))); + } + if s.entries.len() + s.operations.len() >= RECEIPT_LIMIT { + return Err(error("Queue receipt budget exhausted; start a new session")); + } + if self.queues.depth(session) + s.entries.values().filter(|e| e.held.is_some()).count() + >= self.queues.max_depth() + { + return Err(error("Message queue is full")); + } + let display = message + .display_content + .as_deref() + .unwrap_or(&message.content); + let view = DialogQueueItem { + turn_id: message.turn_id.clone(), + display_content: display.chars().take(PREVIEW_CHARS).collect(), + preview_truncated: display.chars().count() > PREVIEW_CHARS, + attachment_count: message.attachments.len(), + agent_type: message.agent_type.clone(), + created_at_ms: SystemTime::now() + .duration_since(SystemTime::UNIX_EPOCH) + .unwrap_or_default() + .as_millis() as u64, + status: DialogQueueStatus::Queued, + reason: None, + target_turn_id: None, + steering_id: None, + }; + s.entries.insert( + message.turn_id.clone(), + Entry { + fingerprint: digest, + view, + held: None, + }, + ); + s.order.push(message.turn_id.clone()); + } + let binding = self + .session_manager + .get_session(session) + .ok_or_else(|| error("Session was removed"))?; + let mut metadata = message.metadata; + for key in [ + "acp_transport", + "backgroundTaskId", + "parentSessionId", + "parentDialogTurnId", + "subagentSessionId", + "subagentDialogTurnId", + "require_tool_confirmation", + ] { + metadata.remove(key); + } + let request = AgentDialogTurnRequest { + session_id: session.to_string(), + message: message.content, + original_message: message.display_content, + turn_id: Some(message.turn_id.clone()), + agent_type: message.agent_type, + workspace_path: binding.config.workspace_path.clone(), + workspace_id: None, + remote_connection_id: binding.config.remote_connection_id.clone(), + remote_ssh_host: binding.config.remote_ssh_host.clone(), + policy: DialogSubmissionPolicy::for_source(DialogTriggerSource::DesktopUi), + execution: Default::default(), + output_schema: None, + reply_route: None, + prepended_reminders: Vec::new(), + attachments: message.attachments, + metadata, + }; + let result = AgentDialogTurnPort::submit_dialog_turn(self, request).await; + let mut state = self.queue_state(); + let s = state.sessions.get_mut(session).expect("queue initialized"); + if let Err(err) = result { + s.entries.remove(&message.turn_id); + s.order.retain(|id| id != &message.turn_id); + return Err(err); + } + s.revision += 1; + drop(state); + Ok(self.queue_snapshot(session, Some(&message.turn_id))) + } + async fn mutate_host_queue( + &self, + session: &str, + action: DialogQueueAction, + ) -> PortResult { + let digest = fingerprint(&action)?; + let (id, operation) = match &action { + DialogQueueAction::Cancel { + turn_id, + operation_id, + } + | DialogQueueAction::Promote { + turn_id, + operation_id, + .. + } => (turn_id.clone(), operation_id.clone()), + _ => return Err(error("Invalid queue operation")), + }; + if operation.trim().is_empty() || operation.len() > 200 { + return Err(error("A stable operation ID is required")); + } + let _guard = self.lock_session_operation(session).await; + { + let state = self.queue_state(); + let s = state.sessions.get(session).expect("queue initialized"); + if let Some((previous, previous_id)) = s.operations.get(&operation) { + if previous != &digest { + return Err(error("idempotency_conflict: operation changed")); + } + let previous_id = previous_id.clone(); + drop(state); + return Ok(self.queue_snapshot(session, Some(&previous_id))); + } + if s.entries.len() + s.operations.len() >= RECEIPT_LIMIT { + return Err(error("Queue receipt budget exhausted")); + } + let entry = s + .entries + .get(&id) + .ok_or_else(|| error("Queue entry not found"))?; + if !matches!( + entry.view.status, + DialogQueueStatus::Queued | DialogQueueStatus::Blocked + ) { + return Err(error( + "too_late: message is no longer available for this queue operation", + )); + } + } + if let DialogQueueAction::Promote { + expected_active_turn_id, + .. + } = &action + { + let actual = self.queue_snapshot(session, None).active_turn_id; + if &actual != expected_active_turn_id { + return Err(error("queue_conflict: active turn changed")); + } + if expected_active_turn_id + .as_deref() + .is_some_and(|target| !self.active_turns.matches_turn(session, target)) + { + return Err(error("queue_conflict: target is no longer active")); + } + if self + .session_manager + .latest_dialog_turn_holds_dispatch(session) + .await + .map_err(|e| error(e.to_string()))? + { + return Err(error("Queue is blocked by interrupted turn recovery")); + } + } + let held = { + let mut state = self.queue_state(); + state + .sessions + .get_mut(session) + .and_then(|s| s.entries.get_mut(&id)) + .and_then(|e| e.held.take()) + }; + let turn = held + .or_else(|| remove_queued_turn_by_id(&self.queues, session, &id)) + .ok_or_else(|| error("too_late: message has already started"))?; + match &action { + DialogQueueAction::Cancel { .. } => { + self.finish_removed_queued_turn(session, turn).await; + self.queue_state() + .sessions + .get_mut(session) + .unwrap() + .entries + .get_mut(&id) + .unwrap() + .view + .status = DialogQueueStatus::Cancelled; + } + DialogQueueAction::Promote { + expected_active_turn_id: Some(target), + .. + } => { + let attachments = turn + .image_contexts + .clone() + .unwrap_or_default() + .into_iter() + .map(|image| { + AgentInputAttachment::image_context( + image.id, + image.image_path, + image.data_url, + image.mime_type, + image.metadata, + ) + }) + .collect(); + let steering_id = Uuid::new_v4().to_string(); + let metadata = turn + .user_message_metadata + .as_ref() + .and_then(|v| v.as_object()) + .cloned() + .unwrap_or_default(); + let decision = resolve_dialog_steering_action( + Some(target), + session, + target, + turn.user_input.clone(), + turn.original_user_input.clone(), + attachments, + metadata, + steering_id.clone(), + SystemTime::now(), + ); + if let DialogSteeringAction::Buffer { injection, .. } = decision { + { + let mut state = self.queue_state(); + let e = state + .sessions + .get_mut(session) + .unwrap() + .entries + .get_mut(&id) + .unwrap(); + e.view.status = DialogQueueStatus::SteeringPending; + e.view.reason = None; + e.view.target_turn_id = Some(target.clone()); + e.view.steering_id = Some(steering_id); + e.held = Some(turn); + } + self.round_injection_buffer.push(session, injection); + } else { + self.queue_state() + .hold(session, &turn, "Unable to steer queued message"); + return Err(error("Unable to steer queued message")); + } + } + DialogQueueAction::Promote { + expected_active_turn_id: None, + .. + } => { + if let Err(e) = self.start_turn(session, &turn).await { + self.queue_state().hold(session, &turn, &e.to_string()); + } + } + _ => unreachable!(), + } + { + let mut state = self.queue_state(); + let s = state.sessions.get_mut(session).unwrap(); + s.operations.insert(operation, (digest, id.clone())); + s.revision += 1; + } + if matches!(action, DialogQueueAction::Cancel { .. }) { + // Removing the final held message releases otherwise healthy work. + let _ = self.try_start_next_queued_locked(session).await; + } + Ok(self.queue_snapshot(session, Some(&id))) + } +} diff --git a/src/crates/assembly/core/src/agentic/coordination/host_message_queue_tests.rs b/src/crates/assembly/core/src/agentic/coordination/host_message_queue_tests.rs new file mode 100644 index 0000000000..168185dd0d --- /dev/null +++ b/src/crates/assembly/core/src/agentic/coordination/host_message_queue_tests.rs @@ -0,0 +1,573 @@ +use openbitfun_runtime_ports::{ + DialogQueueAction as Action, DialogQueueMessage, DialogQueueRequest, + DialogQueueStatus as Status, +}; + +fn request(epoch: Option<&str>, action: Action) -> DialogQueueRequest { + DialogQueueRequest { + session_id: "host-queue-session".into(), + queue_epoch: epoch.map(str::to_owned), + action, + } +} +fn message(id: &str) -> DialogQueueMessage { + DialogQueueMessage { + turn_id: id.into(), + content: "follow up while offline".into(), + display_content: None, + agent_type: "Standard".into(), + attachments: Vec::new(), + metadata: Default::default(), + } +} +async fn fixture() -> ( + Arc, + Arc, + tempfile::TempDir, + String, +) { + let (scheduler, sessions, _, root) = test_scheduler(); + mark_session_processing(&sessions, &root, "host-queue-session", "active-turn").await; + scheduler.active_turns.insert( + "host-queue-session".into(), + desktop_active_turn("active-turn"), + ); + let snapshot = scheduler + .manage_host_queue(request(None, Action::List)) + .await + .unwrap(); + (scheduler, sessions, root, snapshot.queue_epoch) +} + +#[tokio::test] +async fn host_queue_duplicate_and_conflicting_submissions() { + let (scheduler, _, _root, epoch) = fixture().await; + let submit = request( + Some(&epoch), + Action::Submit { + message: message("queued-a"), + }, + ); + let (a, b) = tokio::join!( + scheduler.manage_host_queue(submit.clone()), + scheduler.manage_host_queue(submit) + ); + assert_eq!(a.unwrap().receipt.unwrap().status, Status::Queued); + assert_eq!(b.unwrap().items.len(), 1); + assert_eq!(scheduler.queue_depth("host-queue-session"), 1); + let mut changed = message("queued-a"); + changed.content = "different".into(); + assert!(scheduler + .manage_host_queue(request(Some(&epoch), Action::Submit { message: changed })) + .await + .unwrap_err() + .message + .contains("idempotency_conflict")); +} + +#[tokio::test] +async fn host_queue_cancel_is_idempotent_and_never_cancels_active_turn() { + let (scheduler, _, _root, epoch) = fixture().await; + scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Submit { + message: message("queued-a"), + }, + )) + .await + .unwrap(); + let cancel = request( + Some(&epoch), + Action::Cancel { + turn_id: "queued-a".into(), + operation_id: "cancel-a".into(), + }, + ); + for _ in 0..2 { + assert_eq!( + scheduler + .manage_host_queue(cancel.clone()) + .await + .unwrap() + .receipt + .unwrap() + .status, + Status::Cancelled + ); + } + assert!(scheduler + .active_turns + .matches_turn("host-queue-session", "active-turn")); + assert_eq!(scheduler.queue_depth("host-queue-session"), 0); +} + +#[tokio::test] +async fn host_queue_steering_retains_payload_until_consumption_and_rejects_cancel() { + let (scheduler, _, _root, epoch) = fixture().await; + scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Submit { + message: message("queued-a"), + }, + )) + .await + .unwrap(); + let promote = request( + Some(&epoch), + Action::Promote { + turn_id: "queued-a".into(), + operation_id: "promote-a".into(), + expected_active_turn_id: Some("active-turn".into()), + }, + ); + let snapshot = scheduler.manage_host_queue(promote.clone()).await.unwrap(); + assert_eq!(snapshot.receipt.unwrap().status, Status::SteeringPending); + scheduler.manage_host_queue(promote).await.unwrap(); + let injections = scheduler + .round_injection_source + .take_pending("host-queue-session", "active-turn"); + assert_eq!(injections.len(), 1); + assert_eq!( + scheduler.queue_depth("host-queue-session"), + 1, + "draining the buffer is not consumption" + ); + assert!(scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Cancel { + turn_id: "queued-a".into(), + operation_id: "cancel-a".into() + } + )) + .await + .unwrap_err() + .message + .contains("too_late")); + let injection = &injections[0]; + scheduler.round_injection_source.acknowledge_consumed( + "host-queue-session", + "active-turn", + &injection.id, + injection.kind, + ); + let result = scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Get { + turn_id: "queued-a".into(), + }, + )) + .await + .unwrap(); + assert_eq!(result.receipt.unwrap().status, Status::Steered); + assert_eq!(result.used, 0); +} + +#[tokio::test] +async fn host_queue_unconsumed_steering_and_failed_queue_remain_recoverable() { + let (scheduler, _, _root, epoch) = fixture().await; + for id in ["queued-a", "queued-b"] { + scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Submit { + message: message(id), + }, + )) + .await + .unwrap(); + } + scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Promote { + turn_id: "queued-a".into(), + operation_id: "promote-a".into(), + expected_active_turn_id: Some("active-turn".into()), + }, + )) + .await + .unwrap(); + let _ = scheduler + .round_injection_source + .take_pending("host-queue-session", "active-turn"); + scheduler + .outcome_sender() + .send(( + "host-queue-session".into(), + TurnOutcome::Failed { + turn_id: "active-turn".into(), + error: "provider unavailable".into(), + }, + )) + .unwrap(); + tokio::time::timeout(Duration::from_secs(3), async { + loop { + let snapshot = scheduler + .manage_host_queue(request(None, Action::List)) + .await + .unwrap(); + if snapshot.items.len() == 2 + && snapshot + .items + .iter() + .all(|item| item.status == Status::Blocked) + { + break; + } + tokio::task::yield_now().await; + } + }) + .await + .unwrap(); + assert_eq!(scheduler.queue_depth("host-queue-session"), 2); +} + +#[tokio::test] +async fn host_queue_promote_fences_target_and_epoch() { + let (scheduler, _, _root, epoch) = fixture().await; + scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Submit { + message: message("queued-a"), + }, + )) + .await + .unwrap(); + assert!(scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Promote { + turn_id: "queued-a".into(), + operation_id: "promote-a".into(), + expected_active_turn_id: None + } + )) + .await + .unwrap_err() + .message + .contains("queue_conflict")); + scheduler + .host_queue + .lock() + .unwrap() + .retire("host-queue-session"); + assert!(scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Submit { + message: message("queued-b") + } + )) + .await + .unwrap_err() + .message + .contains("queue_scope_expired")); +} + +#[tokio::test] +async fn host_queue_request_survives_disconnected_caller() { + let (scheduler, _, _root, epoch) = fixture().await; + let guard = scheduler.lock_session_operation("host-queue-session").await; + let owner = scheduler.clone(); + let task = tokio::spawn(async move { + owner + .manage_host_queue(request( + Some(&epoch), + Action::Submit { + message: message("queued-offline"), + }, + )) + .await + }); + // Wait for host admission, then drop the caller while the host is locked. + tokio::time::timeout(Duration::from_secs(3), async { + while !scheduler + .host_queue + .lock() + .unwrap() + .contains("host-queue-session", "queued-offline") + { + tokio::task::yield_now().await; + } + }) + .await + .unwrap(); + task.abort(); + drop(guard); + tokio::time::timeout(Duration::from_secs(3), async { + while scheduler.queue_depth("host-queue-session") != 1 { + tokio::task::yield_now().await; + } + }) + .await + .unwrap(); + assert!(scheduler + .active_turns + .matches_turn("host-queue-session", "active-turn")); +} + +#[tokio::test] +async fn host_queue_host_outcomes_start_followups_without_any_client() { + let (scheduler, sessions, _root, epoch) = fixture().await; + let id = "host-queue-session"; + let ai_config = AIConfig { + models: vec![AIModelConfig { + id: "queue-test-model".into(), + name: "Queue test".into(), + provider: "openai".into(), + model_name: "test-model".into(), + base_url: "http://127.0.0.1:1".into(), + enabled: true, + ..Default::default() + }], + ..Default::default() + }; + TEST_MODEL_RESOLUTION_AI_CONFIG + .scope( + ai_config.clone(), + sessions.update_session_model_id(id, "queue-test-model"), + ) + .await + .unwrap(); + for turn in ["offline-b", "offline-c"] { + scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Submit { + message: message(turn), + }, + )) + .await + .unwrap(); + } + // No query, RPC or controller drives the following transitions. Feed real + // scheduler outcomes; the real coordinator must create both follow-up turns. + let (tx, rx) = mpsc::unbounded_channel(); + let owner = scheduler.clone(); + let handler = tokio::spawn(async move { + TEST_MODEL_RESOLUTION_AI_CONFIG + .scope( + AIConfig { + models: vec![AIModelConfig { + id: "queue-test-model".into(), + name: "Queue test".into(), + provider: "openai".into(), + model_name: "test-model".into(), + base_url: "http://127.0.0.1:1".into(), + enabled: true, + ..Default::default() + }], + ..Default::default() + }, + owner.run_outcome_handler(rx), + ) + .await; + }); + for (finished, started, count) in [ + ("active-turn", "offline-b", 1), + ("offline-b", "offline-c", 2), + ] { + sessions + .update_session_state(id, SessionState::Idle) + .await + .unwrap(); + tx.send(( + id.into(), + TurnOutcome::Completed { + turn_id: finished.into(), + final_response: "done".into(), + }, + )) + .unwrap(); + tokio::time::timeout(Duration::from_secs(5), async { + loop { + if scheduler.active_turns.matches_turn(id, started) { + break; + } + tokio::task::yield_now().await; + } + }) + .await + .expect("host dispatch must start the follow-up"); + assert_eq!(sessions.get_turn_count(id), count); + let _ = scheduler.coordinator.cancel_dialog_turn(id, started).await; + } + handler.abort(); +} + +#[tokio::test] +async fn host_queue_capacity_counts_unconsumed_steering() { + let (scheduler, _, _root, epoch) = fixture().await; + for index in 0..scheduler.queues.max_depth() { + scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Submit { + message: message(&format!("queued-{index}")), + }, + )) + .await + .unwrap(); + } + scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Promote { + turn_id: "queued-0".into(), + operation_id: "promote-first".into(), + expected_active_turn_id: Some("active-turn".into()), + }, + )) + .await + .unwrap(); + assert!(scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Submit { + message: message("overflow") + } + )) + .await + .unwrap_err() + .message + .contains("queue is full")); + assert_eq!( + scheduler.queue_depth("host-queue-session"), + scheduler.queues.max_depth() + ); +} + +#[tokio::test] +async fn host_queue_concurrent_cancel_and_promote_has_one_winner() { + let (scheduler, _, _root, epoch) = fixture().await; + scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Submit { + message: message("queued-a"), + }, + )) + .await + .unwrap(); + let (cancel, promote) = tokio::join!( + scheduler.manage_host_queue(request( + Some(&epoch), + Action::Cancel { + turn_id: "queued-a".into(), + operation_id: "cancel-a".into() + } + )), + scheduler.manage_host_queue(request( + Some(&epoch), + Action::Promote { + turn_id: "queued-a".into(), + operation_id: "promote-a".into(), + expected_active_turn_id: Some("active-turn".into()) + } + )), + ); + assert_ne!(cancel.is_ok(), promote.is_ok()); + assert!(scheduler + .active_turns + .matches_turn("host-queue-session", "active-turn")); +} + +#[tokio::test] +async fn host_queue_cancel_during_terminal_transition_does_not_replace_active_owner() { + let (scheduler, sessions, _root, epoch) = fixture().await; + for id in ["queued-a", "queued-b"] { + scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Submit { + message: message(id), + }, + )) + .await + .unwrap(); + } + sessions + .update_session_state("host-queue-session", SessionState::Idle) + .await + .unwrap(); + scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Cancel { + turn_id: "queued-a".into(), + operation_id: "cancel-a".into(), + }, + )) + .await + .unwrap(); + assert!(scheduler + .active_turns + .matches_turn("host-queue-session", "active-turn")); + assert_eq!(scheduler.queue_depth("host-queue-session"), 1); +} + +#[tokio::test] +async fn host_queue_interrupted_target_cannot_consume_a_blocked_injection_on_resume() { + let (scheduler, _, _root, epoch) = fixture().await; + for id in ["queued-a", "queued-b"] { + scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Submit { + message: message(id), + }, + )) + .await + .unwrap(); + } + scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Promote { + turn_id: "queued-a".into(), + operation_id: "promote-a".into(), + expected_active_turn_id: Some("active-turn".into()), + }, + )) + .await + .unwrap(); + scheduler + .outcome_sender() + .send(( + "host-queue-session".into(), + TurnOutcome::Interrupted { + turn_id: "active-turn".into(), + execution_generation: 0, + }, + )) + .unwrap(); + tokio::time::timeout(Duration::from_secs(3), async { + loop { + let snapshot = scheduler + .manage_host_queue(request(None, Action::List)) + .await + .unwrap(); + if snapshot.items.len() == 2 + && snapshot + .items + .iter() + .all(|item| item.status == Status::Blocked) + { + break; + } + tokio::task::yield_now().await; + } + }) + .await + .unwrap(); + assert!(scheduler + .round_injection_source + .take_pending("host-queue-session", "active-turn") + .is_empty()); + assert_eq!(scheduler.queue_depth("host-queue-session"), 2); +} diff --git a/src/crates/assembly/core/src/agentic/coordination/scheduler.rs b/src/crates/assembly/core/src/agentic/coordination/scheduler.rs index f8da52705b..6a74e1a3eb 100644 --- a/src/crates/assembly/core/src/agentic/coordination/scheduler.rs +++ b/src/crates/assembly/core/src/agentic/coordination/scheduler.rs @@ -10,6 +10,10 @@ //! - FIFO ordering within the same priority level //! - Queue cleared on unrecoverable failure +#[path = "host_message_queue.rs"] +mod host_message_queue; +use host_message_queue::HostQueueState; + use super::coordinator::{ session_storage_workspace_locator, ConversationCoordinator, DialogTriggerSource, DialogTurnStopDisposition, HiddenSubagentExecutionRequest, SubagentResult, @@ -282,6 +286,7 @@ struct BackgroundResultDelivery { } struct SchedulerRoundInjectionSource { + host_queue: Arc>, buffer: Arc, } @@ -305,11 +310,15 @@ impl DialogRoundInjectionSource for SchedulerRoundInjectionSource { fn acknowledge_consumed( &self, - _session_id: &str, - _turn_id: &str, - _injection_id: &str, + session_id: &str, + turn_id: &str, + injection_id: &str, _kind: RoundInjectionKind, ) { + self.host_queue + .lock() + .unwrap_or_else(|e| e.into_inner()) + .consumed(session_id, turn_id, injection_id); } } @@ -319,6 +328,9 @@ impl DialogRoundInjectionSource for SchedulerRoundInjectionSource { /// should submit messages through this scheduler instead of calling /// ConversationCoordinator directly. pub struct DialogScheduler { + self_ref: std::sync::Weak, + host_queue: Arc>, + host_queue_locks: KeyedAsyncLock, coordinator: Arc, session_manager: Arc, /// Per-session priority message queues. @@ -439,11 +451,16 @@ impl DialogScheduler { // retirement of the active-turn owner depends on their delivery. let (outcome_tx, outcome_rx) = mpsc::unbounded_channel(); let round_injection_buffer = Arc::new(SessionRoundInjectionBuffer::default()); + let host_queue = Arc::new(std::sync::Mutex::new(HostQueueState::default())); let round_injection_source = Arc::new(SchedulerRoundInjectionSource { buffer: round_injection_buffer.clone(), + host_queue: host_queue.clone(), }); - let scheduler = Arc::new(Self { + let scheduler = Arc::new_cyclic(|weak| Self { + self_ref: weak.clone(), + host_queue, + host_queue_locks: KeyedAsyncLock::default(), coordinator, session_manager, queues: Arc::new(DialogTurnQueue::default()), @@ -1173,6 +1190,16 @@ impl DialogScheduler { mut queued_turn: QueuedTurn, reject_if_busy: bool, ) -> Result { + if !self + .host_queue + .lock() + .unwrap_or_else(|e| e.into_inner()) + .admission_valid(&session_id, &resolved_turn_id) + { + return Err(SchedulerSubmitError::Message( + "queue_scope_expired: session maintenance retired this submission".into(), + )); + } if let Some(session) = self.session_manager.get_session(&session_id) { queued_turn.workspace_path = session_storage_workspace_locator( queued_turn.workspace_path.as_deref(), @@ -1186,7 +1213,25 @@ impl DialogScheduler { .map(str::trim) .filter(|id| !id.is_empty()) .map(ToOwned::to_owned); - let requested_storage_path = if let Some(workspace_id) = requested_workspace_id.as_deref() { + let host_owned = self + .host_queue + .lock() + .unwrap_or_else(|e| e.into_inner()) + .contains(&session_id, &resolved_turn_id); + let requested_storage_path = if host_owned { + // Queue commands carry no controller filesystem locator. Preserve + // the loaded session's authoritative local/remote storage binding. + Some( + self.session_manager + .effective_session_storage_path(&session_id) + .await + .ok_or_else(|| { + SchedulerSubmitError::Message( + "Host session storage binding unavailable".into(), + ) + })?, + ) + } else if let Some(workspace_id) = requested_workspace_id.as_deref() { // ID-aware callers locate the session by its owning workspace; the // path on the request is only an execution-root projection. Some( @@ -1222,7 +1267,7 @@ impl DialogScheduler { self.session_manager .validate_session_storage_path_binding(&session_id, &requested_storage_path) .map_err(SchedulerSubmitError::Core)?; - if requested_workspace_id.is_some() { + if host_owned || requested_workspace_id.is_some() { // The session is loaded and bound by ID; an omitted locator makes // the coordinator reuse that binding instead of re-resolving a path. queued_turn.workspace_path = None; @@ -1238,6 +1283,18 @@ impl DialogScheduler { .latest_dialog_turn_holds_dispatch(&session_id) .await .map_err(SchedulerSubmitError::Core)?; + if interrupted_hold + && queued_turn.turn_id.as_ref().is_some_and(|id| { + self.host_queue + .lock() + .unwrap_or_else(|e| e.into_inner()) + .contains(&session_id, id) + }) + { + return Err(SchedulerSubmitError::Message( + "Queue is blocked by interrupted turn recovery".into(), + )); + } let interrupted_turn_to_abandon = if interrupted_hold && !matches!( queued_turn.policy.trigger_source, @@ -1250,11 +1307,18 @@ impl DialogScheduler { } else { None }; - let state_fact = if self.active_turns.contains(&session_id) || interrupted_hold { - DialogSessionStateFact::Processing - } else { - Self::session_state_fact(state.as_ref()) - }; + let held_user_messages = self + .host_queue + .lock() + .unwrap_or_else(|e| e.into_inner()) + .pending_held(&session_id) + > 0; + let state_fact = + if self.active_turns.contains(&session_id) || interrupted_hold || held_user_messages { + DialogSessionStateFact::Processing + } else { + Self::session_state_fact(state.as_ref()) + }; let queue_has_items = self.queues.has_items(&session_id); if matches!( @@ -1440,6 +1504,11 @@ impl DialogScheduler { /// Number of messages currently queued for a session. pub fn queue_depth(&self, session_id: &str) -> usize { self.queues.depth(session_id) + + self + .host_queue + .lock() + .unwrap_or_else(|e| e.into_inner()) + .pending_held(session_id) } /// Whether a session has a running or queued turn. This is intentionally a @@ -1447,7 +1516,7 @@ impl DialogScheduler { /// depending on scheduler internals. pub fn is_session_busy_or_queued(&self, session_id: &str) -> bool { self.active_turns.contains(session_id) - || self.queues.has_items(session_id) + || self.queue_depth(session_id) > 0 || self .session_manager .get_session_state(session_id) @@ -1455,6 +1524,12 @@ impl DialogScheduler { } async fn finish_removed_queued_turn(&self, session_id: &str, removed_turn: QueuedTurn) { + if let Some(id) = removed_turn.turn_id.as_deref() { + self.host_queue + .lock() + .unwrap_or_else(|e| e.into_inner()) + .cancelled(session_id, id); + } match removed_turn.execution { QueuedTurnExecution::Standard | QueuedTurnExecution::FreshExternalSubagent(_) => { if let Some(turn_id) = removed_turn.turn_id { @@ -1715,14 +1790,21 @@ impl DialogScheduler { ) -> OpenBitFunResult { openbitfun_core_types::validate_session_id(session_id) .map_err(OpenBitFunError::Validation)?; + let _queue_admission_guard = self.host_queue_locks.lock(session_id).await; let operation_guard = self.lock_session_operation(session_id).await; self.session_manager .validate_session_storage_path_binding(session_id, requested_storage_path)?; - let mut retired_turn_ids = if self.queue_depth(session_id) > 0 { - self.clear_queue(session_id).await - } else { - Vec::new() - }; + let held_turns = self + .host_queue + .lock() + .unwrap_or_else(|e| e.into_inner()) + .retire(session_id); + let mut retired_turn_ids = Vec::new(); + for turn in held_turns { + retired_turn_ids.extend(turn.turn_id.iter().cloned()); + self.finish_removed_queued_turn(session_id, turn).await; + } + retired_turn_ids.extend(self.clear_queue_with_policy(session_id, false).await); abort_thread_goal_continuation_for_session(session_id); let deadline = Instant::now() + wait_timeout; let cancelled_before_parent = self @@ -1816,6 +1898,11 @@ impl DialogScheduler { // ── Private helpers ────────────────────────────────────────────────────── fn enqueue(&self, session_id: &str, queued_turn: QueuedTurn) -> Result<(), String> { + // Called under the session operation lock: held user messages consume + // the same capacity as physical queue entries, including for old producers. + if self.queue_depth(session_id) >= self.queues.max_depth() { + return Err("Message queue is full".into()); + } let priority = queued_turn.policy.queue_priority; let new_len = match self.queues.enqueue(session_id, queued_turn, priority) { Ok(new_len) => new_len, @@ -1837,10 +1924,31 @@ impl DialogScheduler { } async fn clear_queue(&self, session_id: &str) -> Vec { + self.clear_queue_with_policy(session_id, true).await + } + + async fn clear_queue_with_policy( + &self, + session_id: &str, + preserve_user_messages: bool, + ) -> Vec { let cleared_turns = self.queues.clear(session_id); let count = cleared_turns.len(); let mut retired_turn_ids = Vec::new(); for queued_turn in cleared_turns { + if preserve_user_messages + && self + .host_queue + .lock() + .unwrap_or_else(|e| e.into_inner()) + .hold( + session_id, + &queued_turn, + "Previous turn failed; retry or cancel this message", + ) + { + continue; + } match queued_turn.execution { QueuedTurnExecution::Standard | QueuedTurnExecution::FreshExternalSubagent(_) => { if let Some(turn_id) = queued_turn.turn_id { @@ -1907,7 +2015,9 @@ impl DialogScheduler { .session_manager .get_session(session_id) .map(|s| s.state.clone()); - if matches!(state, Some(SessionState::Processing { .. })) { + if self.active_turns.contains(session_id) + || matches!(state, Some(SessionState::Processing { .. })) + { return Ok(None); } if self @@ -1918,6 +2028,15 @@ impl DialogScheduler { return Ok(None); } + if self + .host_queue + .lock() + .unwrap_or_else(|e| e.into_inner()) + .pending_held(session_id) + > 0 + { + return Ok(None); + } let Some(next_turn) = self.dequeue_next(session_id) else { return Ok(None); }; @@ -1931,7 +2050,14 @@ impl DialogScheduler { match self.start_turn(session_id, &next_turn).await { Ok(tid) => Ok(Some(tid)), Err(err) => { - self.requeue_front(session_id, next_turn); + if !self + .host_queue + .lock() + .unwrap_or_else(|e| e.into_inner()) + .hold(session_id, &next_turn, &err.to_string()) + { + self.requeue_front(session_id, next_turn); + } Err(err) } } @@ -1941,6 +2067,21 @@ impl DialogScheduler { &self, session_id: &str, queued_turn: &QueuedTurn, + ) -> Result { + let result = self.start_turn_inner(session_id, queued_turn).await; + if let Ok(id) = &result { + self.host_queue + .lock() + .unwrap_or_else(|e| e.into_inner()) + .started(session_id, id); + } + result + } + + async fn start_turn_inner( + &self, + session_id: &str, + queued_turn: &QueuedTurn, ) -> Result { match &queued_turn.execution { QueuedTurnExecution::HiddenSubagent(execution) => { @@ -2339,6 +2480,21 @@ impl DialogScheduler { }); let lifecycle_plan = resolve_turn_outcome_lifecycle_plan(&outcome, active_turn.is_some()); + let retired_injections = self + .host_queue + .lock() + .unwrap_or_else(|e| e.into_inner()) + .outcome(&session_id, outcome.turn_id(), lifecycle_plan.status); + for injection_id in retired_injections { + self.round_injection_buffer + .remove_by_id(&session_id, &injection_id); + } + if lifecycle_plan.status == TurnOutcomeStatus::Interrupted { + self.hold_managed_queue( + &session_id, + "Turn interrupted; recover it before retrying queued messages", + ); + } if lifecycle_plan.queue_action == TurnOutcomeQueueAction::ClearQueue { debug!( "Turn {}, clearing queue: session_id={}", @@ -2809,6 +2965,13 @@ impl DialogScheduler { #[async_trait::async_trait] impl AgentDialogTurnPort for DialogScheduler { + async fn manage_dialog_queue( + &self, + request: openbitfun_runtime_ports::DialogQueueRequest, + ) -> PortResult { + self.manage_host_queue(request).await + } + async fn submit_dialog_turn( &self, request: AgentDialogTurnRequest, @@ -3180,6 +3343,7 @@ pub fn clear_thread_goal_continuation_abort(session_id: &str) { #[cfg(test)] mod tests { use super::*; + include!("host_message_queue_tests.rs"); use crate::agentic::core::{ProcessingPhase, SessionConfig}; use crate::agentic::events::{EventQueue, EventQueueConfig, EventRouter}; use crate::agentic::execution::{ diff --git a/src/crates/assembly/core/src/service/remote_connect/remote_server.rs b/src/crates/assembly/core/src/service/remote_connect/remote_server.rs index 4c0788356b..b91f2b8cd3 100644 --- a/src/crates/assembly/core/src/service/remote_connect/remote_server.rs +++ b/src/crates/assembly/core/src/service/remote_connect/remote_server.rs @@ -271,6 +271,15 @@ impl RemoteCommandRuntimeHost for CoreRemoteCommandRuntimeHost<'_> { } } + async fn manage_dialog_queue( + &self, + request: openbitfun_runtime_ports::DialogQueueRequest, + ) -> std::result::Result { + CoreServiceAgentRuntime::remote_dialog_host(self.dispatcher)? + .manage_dialog_queue(request) + .await + } + async fn submit_dialog( &self, request: RemoteDialogSubmissionRequest, diff --git a/src/crates/assembly/core/src/service_agent_runtime.rs b/src/crates/assembly/core/src/service_agent_runtime.rs index 6eb2db285a..a98b723aad 100644 --- a/src/crates/assembly/core/src/service_agent_runtime.rs +++ b/src/crates/assembly/core/src/service_agent_runtime.rs @@ -324,6 +324,20 @@ struct ConfiguredPluginDialogTurnPort { #[cfg(feature = "opencode-plugin-host")] #[async_trait::async_trait] impl AgentDialogTurnPort for ConfiguredPluginDialogTurnPort { + async fn manage_dialog_queue( + &self, + request: openbitfun_runtime_ports::DialogQueueRequest, + ) -> PortResult { + if matches!( + &request.action, + openbitfun_runtime_ports::DialogQueueAction::Submit { .. } + | openbitfun_runtime_ports::DialogQueueAction::Promote { .. } + ) { + self.submission.ensure_session(&request.session_id).await; + } + self.inner.manage_dialog_queue(request).await + } + async fn submit_dialog_turn( &self, request: AgentDialogTurnRequest, @@ -2755,6 +2769,23 @@ impl<'a> CoreRemoteDialogRuntimeHost<'a> { }) } + pub(crate) async fn manage_dialog_queue( + &self, + request: openbitfun_runtime_ports::DialogQueueRequest, + ) -> Result { + let binding = self.resolve_binding_workspace(&request.session_id).await; + if !self.remote_session_exists(&request.session_id).await? { + let binding = binding + .ok_or_else(|| "Session workspace is unavailable on this host".to_string())?; + self.restore_remote_session(&request.session_id, binding) + .await?; + } + self.runtime + .manage_dialog_queue(request) + .await + .map_err(CoreServiceAgentRuntime::runtime_error_message) + } + pub(crate) async fn steer_dialog( &self, request: RemoteDialogSteerRequest, diff --git a/src/crates/contracts/product-domains/src/generated/remote-surface-registry.json b/src/crates/contracts/product-domains/src/generated/remote-surface-registry.json index 528c381d86..ec3d8d53eb 100644 --- a/src/crates/contracts/product-domains/src/generated/remote-surface-registry.json +++ b/src/crates/contracts/product-domains/src/generated/remote-surface-registry.json @@ -1,6 +1,6 @@ { "schemaVersion": 1, - "digest": "fnv1a64:01a7a1c7756afe23", + "digest": "fnv1a64:4a41f1e11a6651e1", "retiredCommandPrefixes": [ { "prefix": "lsp_", @@ -9,6 +9,7 @@ ], "capabilities": { "ids": [ + "dialog_queue_v1", "idempotent_dialog_submit", "inline_image_attachments_v1", "btw_initial_model_selection_v1", @@ -29,6 +30,7 @@ "user_question_interaction_v1" ], "desktop": [ + "dialog_queue_v1", "idempotent_dialog_submit", "inline_image_attachments_v1", "btw_initial_model_selection_v1", @@ -49,6 +51,7 @@ "user_question_interaction_v1" ], "cli": [ + "dialog_queue_v1", "control_conversation_v1", "control_conversation_reset_v1", "idempotent_dialog_submit", @@ -4259,6 +4262,17 @@ "reason": "the CLI peer host has no handler for this command" } }, + { + "id": "manage_dialog_queue", + "surface": "tauri_command", + "remoteWorkspace": "RemoteRouted", + "peer": { + "kind": "proxied" + }, + "cliPeer": { + "kind": "handled" + } + }, { "id": "mark_announcement_seen", "surface": "tauri_command", diff --git a/src/crates/contracts/product-domains/src/remote_surface/capabilities.rs b/src/crates/contracts/product-domains/src/remote_surface/capabilities.rs index 1dab574540..3f20b1ef5d 100644 --- a/src/crates/contracts/product-domains/src/remote_surface/capabilities.rs +++ b/src/crates/contracts/product-domains/src/remote_surface/capabilities.rs @@ -22,6 +22,7 @@ use super::PeerHostKind; #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum PeerHostCapability { + DialogQueueV1, /// `start_dialog_turn` / `start_acp_dialog_turn` may be retried with the /// same `(sessionId, turnId)`; the host coalesces duplicate attempts. IdempotentDialogSubmit, @@ -64,6 +65,7 @@ pub enum PeerHostCapability { impl PeerHostCapability { /// Every capability id, in wire order. pub const ALL: &'static [PeerHostCapability] = &[ + Self::DialogQueueV1, Self::IdempotentDialogSubmit, Self::InlineImageAttachmentsV1, Self::BtwInitialModelSelectionV1, @@ -87,6 +89,7 @@ impl PeerHostCapability { /// The key used in the `peer_mode_ping` `capabilities` object. pub const fn key(self) -> &'static str { match self { + Self::DialogQueueV1 => "dialog_queue_v1", Self::IdempotentDialogSubmit => "idempotent_dialog_submit", Self::InlineImageAttachmentsV1 => "inline_image_attachments_v1", Self::ControlConversationV1 => "control_conversation_v1", @@ -119,6 +122,7 @@ const DESKTOP_CAPABILITIES: &[PeerHostCapability] = PeerHostCapability::ALL; /// The CLI peer host has no MiniApp runtime, WSL connection setup, host-native /// ProductControl providers, or presentation surface. const CLI_CAPABILITIES: &[PeerHostCapability] = &[ + PeerHostCapability::DialogQueueV1, PeerHostCapability::ControlConversationV1, PeerHostCapability::ControlConversationResetV1, PeerHostCapability::IdempotentDialogSubmit, diff --git a/src/crates/contracts/product-domains/src/remote_surface/table.rs b/src/crates/contracts/product-domains/src/remote_surface/table.rs index 0040bdcdd4..1736f2b3ce 100644 --- a/src/crates/contracts/product-domains/src/remote_surface/table.rs +++ b/src/crates/contracts/product-domains/src/remote_surface/table.rs @@ -444,6 +444,7 @@ pub(super) const OPERATIONS: &[OperationDefinition] = &[ op("load_session_turns", Unaudited, Proxied, HANDLED), op("local_file_download", Agnostic, ControllerLocal, REFUSED), op("logout_subscription_account", LocalOnly, Proxied, CLI_NOT_IMPLEMENTED), + op("manage_dialog_queue", Routed, Proxied, HANDLED), op("mark_announcement_seen", Agnostic, ControllerLocal, REFUSED), op("mark_openbitfun_control_surface_ready", Agnostic, ControllerLocal, REFUSED), op("mark_openbitfun_control_surface_unready", Agnostic, ControllerLocal, REFUSED), diff --git a/src/crates/contracts/runtime-ports/src/agent_api.rs b/src/crates/contracts/runtime-ports/src/agent_api.rs index 45b50bc523..43ae7b8c3f 100644 --- a/src/crates/contracts/runtime-ports/src/agent_api.rs +++ b/src/crates/contracts/runtime-ports/src/agent_api.rs @@ -1885,6 +1885,16 @@ pub trait AgentTurnSettlementPort: Send + Sync { #[async_trait::async_trait] pub trait AgentDialogTurnPort: Send + Sync { + async fn manage_dialog_queue( + &self, + _request: crate::DialogQueueRequest, + ) -> PortResult { + Err(PortError::new( + PortErrorKind::NotAvailable, + "dialog_queue_v1 is not supported", + )) + } + async fn submit_dialog_turn( &self, request: AgentDialogTurnRequest, @@ -3398,6 +3408,7 @@ mod tests { let snapshot = AgentSessionLineageSnapshot { root_session_id: "root_1".to_string(), sessions: vec![AgentSessionLineageEntry { + workspace_id: None, session_id: "child_1".to_string(), session_name: "Research".to_string(), agent_type: "explore".to_string(), diff --git a/src/crates/contracts/runtime-ports/src/dialog_queue.rs b/src/crates/contracts/runtime-ports/src/dialog_queue.rs new file mode 100644 index 0000000000..bf25d0c95c --- /dev/null +++ b/src/crates/contracts/runtime-ports/src/dialog_queue.rs @@ -0,0 +1,135 @@ +//! Host-owned user message queue. The epoch fences retries across owner restarts. +use crate::AgentInputAttachment; +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct DialogQueueMessage { + pub turn_id: String, + pub content: String, + #[serde(default)] + pub display_content: Option, + pub agent_type: String, + #[serde(default)] + pub attachments: Vec, + #[serde(default)] + pub metadata: serde_json::Map, +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde( + tag = "action", + rename_all = "snake_case", + rename_all_fields = "camelCase" +)] +pub enum DialogQueueAction { + List, + Get { + turn_id: String, + }, + Submit { + message: DialogQueueMessage, + }, + Cancel { + turn_id: String, + operation_id: String, + }, + Promote { + turn_id: String, + operation_id: String, + expected_active_turn_id: Option, + }, +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct DialogQueueRequest { + pub session_id: String, + #[serde(default)] + pub queue_epoch: Option, + #[serde(flatten)] + pub action: DialogQueueAction, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum DialogQueueStatus { + Queued, + Blocked, + SteeringPending, + Steered, + Started, + Interrupted, + Completed, + Failed, + Cancelled, +} +impl DialogQueueStatus { + pub fn is_pending(self) -> bool { + matches!(self, Self::Queued | Self::Blocked | Self::SteeringPending) + } +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct DialogQueueItem { + pub turn_id: String, + pub display_content: String, + pub preview_truncated: bool, + pub attachment_count: usize, + pub agent_type: String, + pub created_at_ms: u64, + pub status: DialogQueueStatus, + pub reason: Option, + pub target_turn_id: Option, + pub steering_id: Option, +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct DialogQueueSnapshot { + pub session_id: String, + pub queue_epoch: String, + pub revision: u64, + pub active_turn_id: Option, + pub items: Vec, + pub capacity: usize, + pub used: usize, + /// Present for submission queries and mutations; terminal receipts remain queryable. + pub receipt: Option, +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn queue_request_wire_defaults_and_extensions_are_compatible() { + let legacy = serde_json::json!({"sessionId":"session", "action":"list"}); + let request: DialogQueueRequest = serde_json::from_value(legacy).unwrap(); + assert_eq!(request.queue_epoch, None); + assert_eq!(request.action, DialogQueueAction::List); + let input = serde_json::json!({"sessionId":"session", "queueEpoch":"owner", "action":"submit", + "message":{"turnId":"turn", "content":"follow up", "agentType":"Standard", "futureField":true}, + "futureExtension": 1}); + let request: DialogQueueRequest = serde_json::from_value(input).unwrap(); + let DialogQueueAction::Submit { message } = &request.action else { + panic!("submit") + }; + assert!(message.attachments.is_empty()); + assert!(message.metadata.is_empty()); + assert_eq!( + serde_json::from_value::(serde_json::to_value(&request).unwrap()) + .unwrap(), + request + ); + } + + #[test] + fn promote_round_trip_preserves_operation_and_expected_turn() { + let input = serde_json::json!({"sessionId":"session", "queueEpoch":"owner", "action":"promote", + "turnId":"queued", "operationId":"op", "expectedActiveTurnId":"active"}); + let request: DialogQueueRequest = serde_json::from_value(input.clone()).unwrap(); + assert_eq!(serde_json::to_value(request).unwrap(), input); + } +} diff --git a/src/crates/contracts/runtime-ports/src/lib.rs b/src/crates/contracts/runtime-ports/src/lib.rs index d73077be9d..5c54fed612 100644 --- a/src/crates/contracts/runtime-ports/src/lib.rs +++ b/src/crates/contracts/runtime-ports/src/lib.rs @@ -461,3 +461,8 @@ mod tests { assert!(json.get("providerId").is_none()); } } + +#[cfg(feature = "agent-api")] +mod dialog_queue; +#[cfg(feature = "agent-api")] +pub use dialog_queue::*; diff --git a/src/crates/execution/agent-runtime/src/runtime.rs b/src/crates/execution/agent-runtime/src/runtime.rs index 08d825de72..f6d35074cb 100644 --- a/src/crates/execution/agent-runtime/src/runtime.rs +++ b/src/crates/execution/agent-runtime/src/runtime.rs @@ -1605,6 +1605,26 @@ impl AgentRuntime { .map_err(RuntimeError::from) } + pub async fn manage_dialog_queue( + &self, + request: openbitfun_runtime_ports::DialogQueueRequest, + ) -> Result { + let session_id = request.session_id.clone(); + let result = self + .dialog_turn + .as_ref() + .ok_or(RuntimeError::MissingDialogTurnPort)? + .manage_dialog_queue(request) + .await + .map_err(RuntimeError::from)?; + if result.session_id != session_id { + return Err( + PortError::new(PortErrorKind::Backend, "Queue session identity mismatch").into(), + ); + } + Ok(result) + } + pub async fn submit_dialog_turn( &self, request: AgentDialogTurnRequest, diff --git a/src/crates/execution/agent-runtime/src/sdk.rs b/src/crates/execution/agent-runtime/src/sdk.rs index b1029b717f..8b37008490 100644 --- a/src/crates/execution/agent-runtime/src/sdk.rs +++ b/src/crates/execution/agent-runtime/src/sdk.rs @@ -678,6 +678,13 @@ impl AgentRuntime { self.inner.submit_turn(request).await } + pub async fn manage_dialog_queue( + &self, + request: openbitfun_runtime_ports::DialogQueueRequest, + ) -> Result { + self.inner.manage_dialog_queue(request).await + } + pub async fn submit_dialog_turn( &self, request: AgentDialogTurnRequest, diff --git a/src/crates/services/services-integrations/src/remote_connect.rs b/src/crates/services/services-integrations/src/remote_connect.rs index af9e166285..a49b6de172 100644 --- a/src/crates/services/services-integrations/src/remote_connect.rs +++ b/src/crates/services/services-integrations/src/remote_connect.rs @@ -565,6 +565,7 @@ fn remote_host_capabilities() -> Vec { "workspace_id_references_v1".to_string(), REMOTE_CAPABILITY_HARNESS_PROFILES_V1.to_string(), REMOTE_CAPABILITY_DIALOG_STEER_V1.to_string(), + "dialog_queue_v1".to_string(), REMOTE_CAPABILITY_PLAN_BUILD_V1.to_string(), REMOTE_CAPABILITY_USER_QUESTION_INTERACTION_V1.to_string(), REMOTE_CAPABILITY_HOST_STREAM_V1.to_string(), @@ -2564,6 +2565,9 @@ pub struct RemoteControlClient { #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] #[serde(tag = "cmd", rename_all = "snake_case")] pub enum RemoteCommand { + DialogQueue { + request: openbitfun_runtime_ports::DialogQueueRequest, + }, /// Retired: relay-stored session history. Kept so older controllers get an /// explicit upgrade message instead of an unknown-command failure. GetSessionKey { @@ -2803,6 +2807,9 @@ pub enum RemoteCommand { #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] #[serde(tag = "resp", rename_all = "snake_case")] pub enum RemoteResponse { + DialogQueue { + snapshot: openbitfun_runtime_ports::DialogQueueSnapshot, + }, /// Retired shape; new hosts never produce it but older peers may still send it. SessionKey { session_id: String, @@ -3056,6 +3063,13 @@ pub trait RemoteCommandRuntimeHost: Send + Sync { } } + async fn manage_dialog_queue( + &self, + _request: openbitfun_runtime_ports::DialogQueueRequest, + ) -> Result { + Err("dialog_queue_v1 is not supported".into()) + } + async fn submit_dialog( &self, request: RemoteDialogSubmissionRequest, @@ -3083,6 +3097,12 @@ where H: RemoteCommandRuntimeHost + ?Sized, { match command { + RemoteCommand::DialogQueue { request } => { + match host.manage_dialog_queue(request.clone()).await { + Ok(snapshot) => RemoteResponse::DialogQueue { snapshot }, + Err(message) => RemoteResponse::Error { message }, + } + } RemoteCommand::Ping { .. } => RemoteResponse::Pong, RemoteCommand::GetWorkspaceInfo diff --git a/src/crates/services/services-integrations/tests/remote_connect_contracts.rs b/src/crates/services/services-integrations/tests/remote_connect_contracts.rs index fb2403d108..c6a7d50f60 100644 --- a/src/crates/services/services-integrations/tests/remote_connect_contracts.rs +++ b/src/crates/services/services-integrations/tests/remote_connect_contracts.rs @@ -3562,3 +3562,21 @@ fn file_chunk_revision_is_additive_for_legacy_peers() { let response: RemoteResponse = serde_json::from_value(current.clone()).unwrap(); assert_eq!(serde_json::to_value(response).unwrap(), current); } + +#[test] +fn dialog_queue_wire_keeps_legacy_send_messages_readable() { + let legacy: RemoteCommand = serde_json::from_value(serde_json::json!({ + "cmd": "send_message", "session_id": "session", "content": "hello" + })) + .unwrap(); + assert!(matches!( + legacy, + RemoteCommand::SendMessage { turn_id: None, .. } + )); + let queue = serde_json::json!({"cmd":"dialog_queue", "request":{ + "sessionId":"session", "queueEpoch":"epoch", "action":"cancel", + "turnId":"queued", "operationId":"cancel-1" + }}); + let command: RemoteCommand = serde_json::from_value(queue.clone()).unwrap(); + assert_eq!(serde_json::to_value(command).unwrap(), queue); +} diff --git a/src/mobile-web/AGENTS.md b/src/mobile-web/AGENTS.md index 18edeb7adb..8e178621d3 100644 --- a/src/mobile-web/AGENTS.md +++ b/src/mobile-web/AGENTS.md @@ -54,6 +54,7 @@ pnpm --dir src/mobile-web run test:interaction-mailbox # independent question/pe pnpm --dir src/mobile-web run test:terminal-browser # real xterm keyboard, ANSI and native bridge pnpm --dir src/mobile-web run test:ui-components pnpm --dir src/mobile-web run test:session-stream-browser # host-driven latest/backward pages, hints, restart gaps, no relay history routes +pnpm --dir src/mobile-web run test:host-queue # idempotent outbox and real browser close/reopen pnpm --dir src/mobile-web run test:host-stream # HostStream reader: hints, reconnects, host restart, unsupported hosts; catalog subscription pnpm --dir src/mobile-web run test:account-login # account login without an online desktop pnpm --dir src/mobile-web run test:account-browser # real Chrome tabs, persistence, migration, races; simulated Relay diff --git a/src/mobile-web/README.md b/src/mobile-web/README.md index 004fbeb239..df0d423a55 100644 --- a/src/mobile-web/README.md +++ b/src/mobile-web/README.md @@ -72,3 +72,34 @@ File uploads stream 3 MiB chunks through the runtime transfer owner, use a whole Account sign-in offers independent GitHub and email-code accounts through the shared OpenBitFun authorization page. No password registration is needed. Use the same login method and account on the desktop/CLI and phone to see its devices. + +## Host-owned message queue + +On hosts advertising `dialog_queue_v1`, the composer remains available while a +turn runs. Accepted follow-ups appear in the host message queue, shared with the +desktop and supported Peer Device controllers. Closing this page, disconnecting +the phone, or leaving the session does not stop host-side dispatch. + +- **Send now** starts the selected message when idle or steers it into the + observed active turn. Waiting for steering is distinct from being consumed. +- **Remove from queue** only removes an unstarted message. It never stops the + active turn. An operation that lost a race with dispatch is rejected. +- A failed turn or unconsumed steering retains the message as blocked on the + host. Resolve the cause, then explicitly send now or remove it. +- A lost response leaves an unconfirmed local record in IndexedDB. **Check / + retry** queries the original message ID before retransmission; it does not + allocate a second submission. Local storage must succeed before sending. +- The guarantee starts when the execution host accepts the message. A request + that never reached the host is not guaranteed to run. Pending messages are + held in host memory: quitting or restarting the execution host can lose them. + A changed queue epoch prevents automatic replay; the submitting browser keeps + its cached text for an explicit recovery decision. + +Older hosts retain the legacy send path and do not expose this queue management +UI. ACP and Detached Dispatch retain their own driver behavior. Permissions and +questions still use the existing remote interaction mailbox. + +Verification: `pnpm --dir src/mobile-web run test:host-queue` covers ambiguous +retries and a real Chromium page close/reopen with IndexedDB. The browser tests +use simulated host/relay data and disposable profiles; they are not evidence of +a physical phone or an SSH workspace test. diff --git a/src/mobile-web/package.json b/src/mobile-web/package.json index dd6e8293a1..6aa4f166d0 100644 --- a/src/mobile-web/package.json +++ b/src/mobile-web/package.json @@ -19,6 +19,7 @@ "type-check": "tsc --noEmit", "build": "vite build", "preview": "vite preview", + "test:host-queue": "node --test tests/host-dialog-queue.test.mjs tests/host-dialog-queue-browser.test.mjs", "test:host-stream": "node --test tests/host-stream.test.mjs tests/host-catalog-subscription.test.mjs", "test:session-stream-browser": "node --test tests/session-stream-browser.test.mjs", "test:runtime-files": "node --test tests/runtime-file-upload.test.mjs tests/runtime-file-download.test.mjs", diff --git a/src/mobile-web/src/components/ChatComposerBar.tsx b/src/mobile-web/src/components/ChatComposerBar.tsx index 8f62c7d71d..13912b0fcf 100644 --- a/src/mobile-web/src/components/ChatComposerBar.tsx +++ b/src/mobile-web/src/components/ChatComposerBar.tsx @@ -78,7 +78,8 @@ export default function ChatComposerBar({ )} size="sm" /> - ) : streaming ? ( + ) : null} + {streaming && ( - ) : expanded ? ( + )} + {!imageAnalyzing && (expanded || input.trim() || pendingImages.length > 0) ? ( void }) { + const { t } = useI18n(); + const view = useSyncExternalStore(queue.subscribe, queue.getSnapshot, queue.getSnapshot); + const [busy, setBusy] = useState(false); + const [error, setError] = useState(null); + useEffect(() => observeHostQueue(queue), [queue]); + const run = async (action: () => Promise) => { + setBusy(true); setError(null); + try { await action(); } catch (e) { setError(String(e)); } + finally { setBusy(false); } + }; + if (!view.snapshot?.items.length && !view.pending.length && !view.error && !error) return null; + return
+ {t('queue.title')} +

{t('queue.memoryNotice')}

+ {(error || view.error) && {error || view.error}} + {(error || view.error) && void run(() => queue.refresh())}>{t('queue.refresh')}} +
    + {view.snapshot?.items.map(item =>
  • +

    {item.displayContent}

    + {item.status === 'steering_pending' ? t('queue.steeringPending') : item.status === 'blocked' ? t('queue.blocked') : t('queue.queued')} + {item.attachmentCount > 0 && {t('queue.attachments', { count: item.attachmentCount })}} + {item.reason &&

    {item.reason}

    } +
    + void run(() => queue.act(item, 'promote'))}>{t('queue.sendNow')} + void run(() => queue.act(item, 'cancel'))}>{t('queue.cancel')} +
    +
  • )} + {view.pending.map(record =>
  • +

    {t('queue.unknown')}

    + {record.request.action === 'submit' &&

    {record.request.message.displayContent ?? record.request.message.content}

    } +
    + void run(() => queue.retry(record))}>{t('queue.checkRetry')} + {record.request.action === 'submit' && { + if (record.request.action === 'submit') onRestore(record.request.message.content); + }}>{t('queue.copyDraft')}} + void run(() => queue.dismiss(record))}>{t('queue.dismiss')} +
    +
  • )} +
+
; +} diff --git a/src/mobile-web/src/i18n/messages.ts b/src/mobile-web/src/i18n/messages.ts index 4e45e896e8..5a7762b997 100644 --- a/src/mobile-web/src/i18n/messages.ts +++ b/src/mobile-web/src/i18n/messages.ts @@ -8,6 +8,24 @@ type MessageTree = { readonly [key: string]: MessageLeaf | MessageTree }; export const messages: Record = { 'en-US': { + queue: { + "title": "Host message queue", + "memoryNotice": "Accepted messages run while this page is closed. Restarting the execution device clears pending messages.", + "queued": "Queued", + "blocked": "Waiting for recovery", + "steeringPending": "Waiting to steer", + "sendNow": "Send now", + "cancel": "Remove from queue", + "refresh": "Refresh", + "unknown": "Delivery is unconfirmed. Check before sending again.", + "checkRetry": "Check / retry", + "copyDraft": "Copy text to composer", + "dismiss": "Dismiss local reminder", + "attachments": "{count} attachments", + "edit": "Restore draft", + "noDraft": "The complete draft is only available on the device that submitted it.", + "legacy": "Saved drafts from the previous queue. Send or restore each explicitly." +}, shared: SHARED_TERMS_BY_LOCALE['en-US'], common: { questionTimeoutActive: 'Could not stop the question timeout. Submit promptly or update the execution device.', @@ -341,6 +359,24 @@ export const messages: Record = { }, }, 'zh-CN': { + queue: { + "title": "宿主消息队列", + "memoryNotice": "消息接受后,关闭此页面仍会执行。执行设备重启会清空待执行消息。", + "queued": "排队中", + "blocked": "等待恢复", + "steeringPending": "等待注入", + "sendNow": "立即发送", + "cancel": "移出队列", + "refresh": "刷新", + "unknown": "尚未确认是否送达,请先检查再重发。", + "checkRetry": "检查 / 重试", + "copyDraft": "复制文字到输入框", + "dismiss": "关闭本地提醒", + "attachments": "{count} 个附件", + "edit": "恢复草稿", + "noDraft": "完整草稿仅在提交消息的设备上可用。", + "legacy": "旧队列中的草稿,请逐条发送或恢复。" +}, shared: SHARED_TERMS_BY_LOCALE['zh-CN'], common: { questionTimeoutActive: '未能停止提问计时,请尽快提交或更新执行设备。', @@ -674,6 +710,24 @@ export const messages: Record = { }, }, 'zh-TW': { + queue: { + "title": "主機訊息佇列", + "memoryNotice": "訊息接受後,關閉此頁面仍會執行。執行裝置重新啟動會清空待執行訊息。", + "queued": "排隊中", + "blocked": "等待恢復", + "steeringPending": "等待注入", + "sendNow": "立即傳送", + "cancel": "移出佇列", + "refresh": "重新整理", + "unknown": "尚未確認是否送達,請先檢查再重新傳送。", + "checkRetry": "檢查 / 重試", + "copyDraft": "複製文字到輸入框", + "dismiss": "關閉本機提醒", + "attachments": "{count} 個附件", + "edit": "恢復草稿", + "noDraft": "完整草稿僅在提交訊息的裝置上可用。", + "legacy": "舊佇列中的草稿,請逐條傳送或恢復。" +}, shared: SHARED_TERMS_BY_LOCALE['zh-TW'], common: { questionTimeoutActive: '無法停止提問計時,請儘快提交或更新執行裝置。', diff --git a/src/mobile-web/src/pages/ChatPage.tsx b/src/mobile-web/src/pages/ChatPage.tsx index 22c7846d14..b1ffe269e2 100644 --- a/src/mobile-web/src/pages/ChatPage.tsx +++ b/src/mobile-web/src/pages/ChatPage.tsx @@ -1,3 +1,4 @@ +import { MobileHostQueue } from '../components/MobileHostQueue'; import { downloadRuntimeFile } from '../services/RuntimeFileDownload'; import { PermissionMailbox } from '../components/PermissionMailbox'; import { QuestionInteractionContext } from "../components/ChatAskQuestionCard"; @@ -142,6 +143,9 @@ const ChatPage: React.FC = ({ const isLoadingMoreRef = useRef(false); const hasMoreRef = useRef(true); const controlTargetEpoch = useControlTargetEpoch(sessionMgr); + const queueSupported = sessionMgr.supportsHostCapability('dialog_queue_v1'); + const hostQueue = useMemo(() => queueSupported ? sessionMgr.dialogQueue(sessionId) : null, + [sessionMgr, sessionId, controlTargetEpoch, queueSupported]); const cacheScope = useMemo(() => createRemoteCacheScope( authenticatedUserId, controlTarget?.deviceId ?? sessionMgr.controlTargetDeviceId, @@ -904,8 +908,10 @@ const ChatPage: React.FC = ({ setPendingImages(current => current.filter(image => !imgs.includes(image))); if (!wasStreaming && draftUnchanged) setInputExpanded(false); streamRef.current?.nudge(); - if (wasStreaming) { + if (hostQueue?.getSnapshot().snapshot?.receipt?.status === 'queued') { setInfoToast(t('chat.messageQueued')); + } else if (!hostQueue && wasStreaming) { + setInfoToast(t('common.submitted')); } } catch (e: any) { if (!isChatTargetCurrent(targetEpoch)) return; @@ -918,7 +924,7 @@ const ChatPage: React.FC = ({ setOptimisticMsg(null); } } - }, [captureChatTargetEpoch, imageAnalyzing, input, isChatTargetCurrent, isStreaming, pendingImages, sessionAgentType, sessionId, sessionMgr, setError, t]); + }, [captureChatTargetEpoch, hostQueue, imageAnalyzing, input, isChatTargetCurrent, isStreaming, pendingImages, sessionAgentType, sessionId, sessionMgr, setError, t]); const handleImageSelect = useCallback(() => { fileInputRef.current?.click(); @@ -1111,6 +1117,8 @@ const ChatPage: React.FC = ({ style={{ display: 'none' }} onChange={handleFileChange} /> + {hostQueue && { setInput(current => current ? `${current}\n\n${content}` : content); setInputExpanded(true); }} />} (); + + dialogQueue(sessionId: string): HostDialogQueue { + const target = this.client.getControlTargetSnapshot(); + const account = this.client.accountUserId; + const epoch = this.controlTargetEpoch; + const scope = JSON.stringify([account, target.deviceId, sessionId]); + const key = JSON.stringify([scope, this.controlTargetEpoch]); + let queue = this.queueClients.get(key); + if (!queue) { + queue = new HostDialogQueue(scope, sessionId, async request => { + if (this.client.accountUserId !== account || this.controlTargetEpoch !== epoch) throw new Error('Queue target changed'); + if (!this.supportsHostCapability('dialog_queue_v1')) throw new Error('Host message queue is unsupported'); + const response = await this.request<{ snapshot: QueueSnapshot }>({ cmd: 'dialog_queue', request }, target); + if (this.client.accountUserId !== account || this.controlTargetEpoch !== epoch) throw new Error('Queue target changed'); + return response.snapshot; + }); + this.queueClients.set(key, queue); + } + return queue; + } + async sendMessage( sessionId: string, content: string, @@ -756,6 +779,15 @@ export class RemoteSessionManager { metadata?: Record; }>, ): Promise { + if (this.supportsHostCapability('dialog_queue_v1')) { + const result = await this.dialogQueue(sessionId).submit({ content, agentType: agentType || 'Standard', + attachments: (imageContexts ?? []).map(image => ({ kind: 'remote_image', id: image.id, + metadata: { ...(image.data_url ? { dataUrl: image.data_url } : {}), + ...(image.image_path ? { imagePath: image.image_path } : {}), mimeType: image.mime_type, + metadata: image.metadata } })), metadata: {} }); + if (!result.receipt) throw new Error('Host did not acknowledge the submitted message'); + return result.receipt.turnId; + } const resp = await this.request<{ resp: string; turn_id: string }>({ cmd: 'send_message', session_id: sessionId, diff --git a/src/mobile-web/src/styles/host-queue.scss b/src/mobile-web/src/styles/host-queue.scss new file mode 100644 index 0000000000..0a4a1fe008 --- /dev/null +++ b/src/mobile-web/src/styles/host-queue.scss @@ -0,0 +1,12 @@ +.host-message-queue { + padding: 8px 12px; + max-height: 28vh; + overflow: auto; + flex-shrink: 0; + font-size: 12px; + ul { list-style: none; margin: 0; padding: 0; } + li { padding-block: 8px; } + p { margin-block: 4px; } + &__preview { white-space: pre-wrap; overflow-wrap: anywhere; max-height: 5em; overflow: auto; } + &__actions { display: flex; flex-wrap: wrap; gap: 8px; margin-top: 4px; } +} diff --git a/src/mobile-web/tests/fixtures/host-queue.tsx b/src/mobile-web/tests/fixtures/host-queue.tsx new file mode 100644 index 0000000000..c068567434 --- /dev/null +++ b/src/mobile-web/tests/fixtures/host-queue.tsx @@ -0,0 +1,28 @@ +import React from 'react'; +import { createRoot } from 'react-dom/client'; +import ChatComposerBar from '../../src/components/ChatComposerBar'; +import { MobileHostQueue } from '../../src/components/MobileHostQueue'; +import { I18nProvider } from '../../src/i18n'; +import { HostDialogQueue } from '../../../shared/dialog-queue/HostDialogQueue'; + +export function mountHostQueueFixture() { + const element = document.createElement('main'); + document.body.replaceChildren(element); + const calls: string[] = []; + const queue = new HostDialogQueue('ui-fixture', 'session', async request => { + calls.push(request.action); + return { sessionId: 'session', queueEpoch: 'epoch', revision: calls.length, activeTurnId: 'active', + items: [{ turnId: 'queued', content: '', displayContent: '接着检查错误处理和测试覆盖', previewTruncated: false, + attachmentCount: 0, agentType: 'Standard', createdAtMs: 1, status: 'queued', reason: null, targetTurnId: null, steeringId: null }], + capacity: 20, used: 1, receipt: null }; + }); + const noop = () => {}; + const root = createRoot(element); + root.render( + calls.push('stop')} + onChange={noop} onCompositionEnd={noop} onCompositionStart={noop} onKeyDown={noop} onRemoveImage={noop} + onSend={() => calls.push('send')} pendingImages={[]} remoteUnavailable={false} streaming /> + ); + return { calls, dispose: () => root.unmount() }; +} diff --git a/src/mobile-web/tests/host-dialog-queue-browser.test.mjs b/src/mobile-web/tests/host-dialog-queue-browser.test.mjs new file mode 100644 index 0000000000..7e2599ccdb --- /dev/null +++ b/src/mobile-web/tests/host-dialog-queue-browser.test.mjs @@ -0,0 +1,54 @@ +import assert from 'node:assert/strict'; +import { test } from 'node:test'; +import { fileURLToPath } from 'node:url'; +import { launchBrowser, startSourceServer } from './helpers/browser-account-harness.mjs'; +const modulePath = '/@fs' + fileURLToPath(new URL('../../shared/dialog-queue/HostDialogQueue.ts', import.meta.url)); + +test('real IndexedDB retains an ambiguous submission after closing the browser page', {timeout:60000}, async () => { + const server=await startSourceServer();const browser=await launchBrowser(); + try { + const context=await browser.createIncognitoBrowserContext();let page=await context.newPage(); + await page.goto(server.origin); + const id=await page.evaluate(async path=>{ + const {HostDialogQueue}=await import(path); + const q=new HostDialogQueue('browser-account/host/session','session',async request=>{ + if(request.action==='submit')throw new Error('Lost acknowledgement'); + return {sessionId:'session',queueEpoch:'host-epoch',revision:0,activeTurnId:'running',items:[],capacity:20,used:0,receipt:null}; + }); + try {await q.submit({content:'offline follow up',agentType:'Standard',attachments:[],metadata:{}});}catch{} + return q.getSnapshot().pending[0].request.message.turnId; + },modulePath); + await page.close();page=await context.newPage();await page.goto(server.origin); + const result=await page.evaluate(async ({path,id})=>{ + const {HostDialogQueue}=await import(path);const mutations=[]; + const q=new HostDialogQueue('browser-account/host/session','session',async request=>{ + if(request.action==='submit')mutations.push(request); + return {sessionId:'session',queueEpoch:'host-epoch',revision:2,activeTurnId:null,items:[],capacity:20,used:0, + receipt:request.action==='get'?{turnId:id,status:'completed',displayContent:'offline follow up'}:null}; + }); + await q.refresh();const pending=q.getSnapshot().pending; + await q.retry(pending[0]);return {restoredId:pending[0].request.message.turnId,pending:q.getSnapshot().pending.length,mutations:mutations.length}; + },{path:modulePath,id}); + assert.deepEqual(result,{restoredId:id,pending:0,mutations:0}); + } finally {await browser.close();await server.close();} +}); + +test('mobile running composer keeps send and stop independently available alongside the queue', {timeout:60000}, async () => { + const server=await startSourceServer();const browser=await launchBrowser(); + try { + const page=await browser.newPage();await page.setViewport({width:390,height:844,isMobile:true,hasTouch:true}); + await page.goto(server.origin+'/?lang=zh-CN'); + await page.evaluate(async()=>{ + const {mountHostQueueFixture}=await import('/tests/fixtures/host-queue.tsx'); + window.queueFixture=mountHostQueueFixture(); + }); + await page.waitForSelector('.host-message-queue li'); + const actions=await page.$$eval('.chat-page__send-btn', buttons=>buttons.map(button=>({disabled:button.disabled,stop:button.classList.contains('is-stop')}))); + assert.deepEqual(actions,[{disabled:false,stop:true},{disabled:false,stop:false}]); + await page.click('.chat-page__send-btn:not(.is-stop)'); + assert.ok((await page.evaluate(()=>window.queueFixture.calls)).includes('send')); + assert.ok(!(await page.evaluate(()=>window.queueFixture.calls)).includes('stop')); + assert.equal(await page.$eval('body', body=>body.scrollWidth<=window.innerWidth),true); + await page.screenshot({path:'/tmp/mobile-host-message-queue.png',fullPage:true}); + }finally{await browser.close();await server.close();} +}); diff --git a/src/mobile-web/tests/host-dialog-queue.test.mjs b/src/mobile-web/tests/host-dialog-queue.test.mjs new file mode 100644 index 0000000000..417311b626 --- /dev/null +++ b/src/mobile-web/tests/host-dialog-queue.test.mjs @@ -0,0 +1,103 @@ +import assert from 'node:assert/strict'; +import { readFile } from 'node:fs/promises'; +import test from 'node:test'; +import ts from 'typescript'; +const source = await readFile(new URL('../../shared/dialog-queue/HostDialogQueue.ts', import.meta.url), 'utf8'); +const code = ts.transpileModule(source, { compilerOptions: { module: ts.ModuleKind.ESNext, target: ts.ScriptTarget.ES2022 } }).outputText; +const { HostDialogQueue } = await import(`data:text/javascript;base64,${Buffer.from(code).toString('base64')}`); +function fixture() { + const records = new Map(); const calls = []; const receipts = new Map(); + let epoch = 'owner-1'; let loseResponse = false; let executions = 0; + const storage = { + list: async scope => [...records.values()].filter(record => record.scope === scope), + put: async record => { records.set(record.key, structuredClone(record)); }, + remove: async key => { records.delete(key); }, + }; + const snapshot = receipt => ({sessionId:'session-a',queueEpoch:epoch,revision:executions,activeTurnId:'running',items:[...receipts.values()].filter(item=>item.status==='queued'),capacity:20,used:receipts.size,receipt}); + const invoke = async request => { + calls.push(structuredClone(request)); + if(request.action==='submit') { + assert.ok([...records.values()].some(record=>record.request.message?.turnId===request.message.turnId),'outbox must commit before RPC'); + if(!receipts.has(request.message.turnId)) { executions++; receipts.set(request.message.turnId,{turnId:request.message.turnId,status:'queued',displayContent:request.message.content}); } + if(loseResponse) {loseResponse=false;throw new Error('Disconnected after host accepted');} + return snapshot(receipts.get(request.message.turnId)); + } + if(request.action==='get')return snapshot(receipts.get(request.turnId)??null); + if(request.action==='cancel'||request.action==='promote') { + if(loseResponse) {loseResponse=false;throw new Error('Connection lost');} + return snapshot(null); + } + return snapshot(null); + }; + return {storage,calls,records,receipts,invoke,create:()=>new HostDialogQueue('account/host/session-a','session-a',invoke,storage), + lose:()=>{loseResponse=true;},restart:()=>{epoch='owner-2';receipts.clear();},get executions(){return executions;}}; +} +const message={content:'continue',agentType:'Standard',attachments:[],metadata:{}}; +test('lost acceptance survives page recreation and queries original ID without a duplicate send',async()=>{ + const f=fixture();const first=f.create();f.lose();await assert.rejects(first.submit(message)); + const stored=[...f.records.values()][0];assert.equal(stored.accepted,undefined); + const reopened=f.create();await reopened.refresh();await reopened.retry(reopened.getSnapshot().pending[0]); + assert.equal(f.executions,1);assert.equal(f.calls.filter(call=>call.action==='submit').length,1); + assert.equal(reopened.getSnapshot().pending.length,0);assert.ok(await reopened.savedDraft(stored.request.message.turnId)); +}); +test('an unreceived request retries with the identical payload and stable ID',async()=>{ + const f=fixture();const q=f.create();f.lose();await assert.rejects(q.submit(message));f.receipts.clear(); + const before=f.calls.find(call=>call.action==='submit');await q.submit(message); + assert.deepEqual(f.calls.filter(call=>call.action==='submit')[1],before); +}); +test('host restart never automatically replays an ambiguous message',async()=>{ + const f=fixture();const q=f.create();f.lose();await assert.rejects(q.submit(message));f.restart(); + await assert.rejects(q.retry(q.getSnapshot().pending[0]),/queue_scope_expired/); + assert.equal(f.calls.filter(call=>call.action==='submit').length,1);assert.equal(q.getSnapshot().pending.length,1); +}); +test('failed local storage prevents any submission',async()=>{ + const f=fixture();const q=new HostDialogQueue('scope','session-a',f.invoke,{...f.storage,put:async()=>{throw new Error('storage full');}}); + await assert.rejects(q.submit(message),/storage full/);assert.equal(f.executions,0); +}); +test('a stale revision cannot roll back the replica',async()=>{ + let revision=4;const f=fixture();const q=new HostDialogQueue('scope','session-a',async request=>({...await f.invoke(request),revision}),f.storage); + await q.refresh();revision=2;await q.refresh();assert.equal(q.getSnapshot().snapshot.revision,4); +}); +test('promotion retry preserves operation identity and observed active turn',async()=>{ + const f=fixture();const q=f.create();await q.refresh();f.lose();await assert.rejects(q.act({turnId:'queued'},'promote')); + const before=f.calls.find(call=>call.action==='promote');await q.retry(q.getSnapshot().pending[0]); + assert.deepEqual(f.calls.filter(call=>call.action==='promote')[1],before);assert.equal(before.expectedActiveTurnId,'running'); +}); +test('outbox retry rejects a different target scope',async()=>{ + const f=fixture();const q=f.create();f.lose();await assert.rejects(q.submit(message)); + const other=new HostDialogQueue('other','session-a',f.invoke,f.storage); + await assert.rejects(other.retry(q.getSnapshot().pending[0]),/Queue target changed/);assert.equal(f.executions,1); +}); + +test('accepted pending drafts become explicit recovery records after owner restart',async()=>{ + const f=fixture();const q=f.create();await q.submit(message);assert.equal(q.getSnapshot().pending.length,0); + f.restart();await q.refresh();assert.equal(q.getSnapshot().pending.length,1); + await assert.rejects(q.retry(q.getSnapshot().pending[0]),/queue_scope_expired/);assert.equal(f.executions,1); +}); +test('terminal accepted drafts do not accumulate attachment payloads indefinitely',async()=>{ + const f=fixture();const q=f.create();await q.submit(message);assert.equal(f.records.size,1); + for(const item of f.receipts.values())item.status='completed';await q.refresh();assert.equal(f.records.size,0); +}); + +test('an edit intent preserves the complete draft when cancellation succeeds but its response is lost',async()=>{ + const f=fixture();const q=f.create();const accepted=await q.submit(message,{attachments:['original attachment context']}); + const saved=await q.savedDraft(accepted.receipt.turnId);await q.prepareRestore(saved); + f.receipts.get(accepted.receipt.turnId).status='cancelled'; + // An observer learns of cancellation before the edit caller gets a response. + const reopened=f.create();await reopened.refresh(); + assert.equal(reopened.getSnapshot().pending.length,1); + assert.deepEqual(reopened.getSnapshot().pending[0].draft,{attachments:['original attachment context']}); + assert.equal((await reopened.receipt(accepted.receipt.turnId)).status,'cancelled'); + assert.equal(f.calls.filter(call=>call.action==='submit').length,1); +}); + +test('a serialized promote keeps the turn observed at click time',async()=>{ + const f=fixture();const q=new HostDialogQueue('scope','session-a',async request=>{ + const snapshot=await f.invoke(request); + return {...snapshot,activeTurnId:request.action==='cancel'?'new-turn':'running'}; + },f.storage); + await q.refresh(); + const first=q.act({turnId:'a'},'cancel');const second=q.act({turnId:'b'},'promote'); + await Promise.all([first,second]); + assert.equal(f.calls.find(call=>call.action==='promote').expectedActiveTurnId,'running'); +}); diff --git a/src/shared/dialog-queue/HostDialogQueue.ts b/src/shared/dialog-queue/HostDialogQueue.ts new file mode 100644 index 0000000000..3a21b40a6f --- /dev/null +++ b/src/shared/dialog-queue/HostDialogQueue.ts @@ -0,0 +1,246 @@ +/** One host's queue replica. Timers observe; only the host dispatches work. */ +export type QueueStatus = 'queued' | 'blocked' | 'steering_pending' | 'steered' + | 'started' | 'interrupted' | 'completed' | 'failed' | 'cancelled'; +export interface QueueAttachment { kind: string; id: string; metadata: Record } +export interface QueueMessage { + turnId: string; + content: string; + displayContent?: string; + agentType: string; + attachments: QueueAttachment[]; + metadata: Record; +} +export interface QueueItem { + turnId: string; displayContent: string; previewTruncated: boolean; attachmentCount: number; + agentType: string; createdAtMs: number; status: QueueStatus; reason: string | null; + targetTurnId: string | null; steeringId: string | null; +} +export interface QueueSnapshot { + sessionId: string; queueEpoch: string; revision: number; activeTurnId: string | null; + items: QueueItem[]; capacity: number; used: number; receipt: QueueItem | null; +} +export type QueueAction = { action: 'list' } | { action: 'get'; turnId: string } + | { action: 'submit'; message: QueueMessage } + | { action: 'cancel'; turnId: string; operationId: string } + | { action: 'promote'; turnId: string; operationId: string; expectedActiveTurnId: string | null }; +export type QueueRequest = QueueAction & { sessionId: string; queueEpoch?: string }; +export interface QueueOutboxRecord { key: string; scope: string; request: QueueRequest; accepted?: boolean; restoreIntent?: boolean; draft?: unknown } +export interface QueueStorage { + list(scope: string): Promise; + put(record: QueueOutboxRecord): Promise; + remove(key: string): Promise; +} +const DB_NAME = 'openbitfun.host-dialog-queue.v1'; +const STORE = 'outbox'; +function openDatabase(): Promise { + return new Promise((resolve, reject) => { + const request = indexedDB.open(DB_NAME, 1); + request.onupgradeneeded = () => request.result.createObjectStore(STORE, { keyPath: 'key' }); + request.onsuccess = () => resolve(request.result); + request.onerror = () => reject(request.error ?? new Error('Queue storage unavailable')); + request.onblocked = () => reject(new Error('Queue storage is blocked by another browser tab')); + }); +} +async function transaction(mode: IDBTransactionMode, run: (store: IDBObjectStore) => IDBRequest): Promise { + const db = await openDatabase(); + try { + return await new Promise((resolve, reject) => { + const tx = db.transaction(STORE, mode); + const request = run(tx.objectStore(STORE)); + tx.oncomplete = () => resolve(request.result); + tx.onerror = () => reject(tx.error ?? request.error ?? new Error('Queue storage failed')); + tx.onabort = () => reject(tx.error ?? new Error('Queue storage transaction aborted')); + }); + } finally { db.close(); } +} +export const queueStorage: QueueStorage = { + async list(scope) { + const records = await transaction('readonly', store => store.getAll()); + return records.filter(record => record.scope === scope); + }, + async put(record) { await transaction('readwrite', store => store.put(record)); }, + async remove(key) { await transaction('readwrite', store => store.delete(key)); }, +}; +export interface QueueView { snapshot: QueueSnapshot | null; error: string | null; pending: QueueOutboxRecord[] } +const newId = () => crypto.randomUUID(); +function messageIdentity(message: Omit): string { + // Reopening the composer can allocate fresh image IDs. An ambiguous retry + // must send the stored payload, including its original IDs, unchanged. + return JSON.stringify({ ...message, attachments: message.attachments.map(({ kind, metadata }) => ({ kind, metadata })) }); +} +export class HostDialogQueue { + private view: QueueView = { snapshot: null, error: null, pending: [] }; + private listeners = new Set<() => void>(); + private refreshing: Promise | null = null; + private mutation: Promise = Promise.resolve(); + constructor(readonly scope: string, readonly sessionId: string, + private invoke: (request: QueueRequest) => Promise, private storage: QueueStorage = queueStorage) {} + getSnapshot = (): QueueView => this.view; + subscribe = (listener: () => void): (() => void) => { + this.listeners.add(listener); return () => { this.listeners.delete(listener); }; + }; + private publish(patch: Partial): void { + this.view = { ...this.view, ...patch }; for (const listener of this.listeners) listener(); + } + private accept(snapshot: QueueSnapshot, authoritative: boolean): void { + if (snapshot.sessionId !== this.sessionId || !snapshot.queueEpoch || !Number.isSafeInteger(snapshot.revision)) { + throw new Error('Invalid host queue snapshot'); + } + const previous = this.view.snapshot; + if (previous && snapshot.queueEpoch !== previous.queueEpoch && !authoritative) return; + if (previous?.queueEpoch === snapshot.queueEpoch && snapshot.revision < previous.revision) return; + this.publish({ snapshot, error: null }); + } + async refresh(): Promise { + if (this.refreshing) return this.refreshing; + this.refreshing = (async () => { + try { + const snapshot = await this.invoke({ sessionId: this.sessionId, action: 'list' }); + this.accept(snapshot, true); + const records = await this.storage.list(this.scope); + const current = this.view.snapshot!; + for (const record of records) { + if (!record.accepted || record.restoreIntent || record.request.action !== 'submit') continue; + if (record.request.queueEpoch !== current.queueEpoch) { + // A restarted owner cannot vouch for an accepted in-memory message. + // Preserve the original draft and require an explicit user decision. + await this.storage.put({ ...record, accepted: false }); + } else if (!current.items.some(item => item.turnId === (record.request as Extract).message.turnId)) { + await this.storage.remove(record.key); + } + } + this.publish({ pending: (await this.storage.list(this.scope)).filter(record => !record.accepted || record.restoreIntent) }); + return current; + } catch (error) { this.publish({ error: String(error) }); throw error; } + finally { this.refreshing = null; } + })(); + return this.refreshing; + } + private serialize(run: () => Promise): Promise { + const result = this.mutation.then(run, run); + this.mutation = result.catch(() => undefined); + return result; + } + private async transmit(record: QueueOutboxRecord): Promise { + // The transaction commits before an RPC is permitted to leave this client. + await this.storage.put(record); + this.publish({ pending: (await this.storage.list(this.scope)).filter(record => !record.accepted || record.restoreIntent) }); + try { + const snapshot = await this.invoke(record.request); + if (snapshot.queueEpoch !== record.request.queueEpoch) throw new Error('queue_scope_expired: host queue owner changed'); + if (record.request.action === 'submit' && snapshot.receipt?.turnId !== record.request.message.turnId) { + throw new Error('Host did not acknowledge the submitted message'); + } + this.accept(snapshot, false); + if (record.request.action === 'submit') await this.storage.put({ ...record, accepted: true }); + else await this.storage.remove(record.key); + this.publish({ pending: (await this.storage.list(this.scope)).filter(record => !record.accepted || record.restoreIntent) }); + return snapshot; + } catch (error) { + const text = String(error); + if (/queue_conflict:|too_late:|idempotency_conflict:|Message queue is full|Queue is blocked by interrupted/.test(text)) { + await this.storage.remove(record.key); + this.publish({ pending: (await this.storage.list(this.scope)).filter(record => !record.accepted || record.restoreIntent) }); + } + this.publish({ error: text }); throw error; + } + } + submit(message: Omit, draft?: unknown, requestedTurnId?: string): Promise { + return this.serialize(async () => { + const snapshot = await this.refresh(); + const identity = messageIdentity(message); + const previous = this.view.pending.find(record => { + if (record.restoreIntent || record.request.action !== 'submit') return false; + const { turnId: _id, ...payload } = record.request.message; + return messageIdentity(payload) === identity; + }); + if (previous) return this.retryRecord(previous, snapshot); + const turnId = requestedTurnId?.trim() || newId(); + return this.transmit({ scope: this.scope, key: JSON.stringify([this.scope, turnId]), draft, + request: { sessionId: this.sessionId, queueEpoch: snapshot.queueEpoch, action: 'submit', message: { ...message, turnId } } }); + }); + } + act(item: QueueItem, action: 'cancel' | 'promote'): Promise { + const observed = this.view; + return this.serialize(async () => { + const snapshot = observed.snapshot; + if (!snapshot || observed.error) throw new Error('Refresh the host queue before changing it'); + const previous = this.view.pending.find(record => record.request.action === action + && 'turnId' in record.request && record.request.turnId === item.turnId); + if (previous) return this.retryRecord(previous, snapshot); + const operationId = newId(); + const request: QueueRequest = action === 'cancel' + ? { sessionId: this.sessionId, queueEpoch: snapshot.queueEpoch, action, turnId: item.turnId, operationId } + : { sessionId: this.sessionId, queueEpoch: snapshot.queueEpoch, action, turnId: item.turnId, operationId, + expectedActiveTurnId: snapshot.activeTurnId }; + return this.transmit({ scope: this.scope, key: JSON.stringify([this.scope, operationId]), request }); + }); + } + retry(record: QueueOutboxRecord): Promise { + return this.serialize(async () => this.retryRecord(record, await this.refresh())); + } + private async retryRecord(record: QueueOutboxRecord, snapshot: QueueSnapshot): Promise { + if (record.scope !== this.scope || record.request.sessionId !== this.sessionId) throw new Error('Queue target changed'); + if (record.request.queueEpoch !== snapshot.queueEpoch) { + throw new Error('queue_scope_expired: host restarted; delivery is unknown. Restore the draft before submitting again.'); + } + if (record.request.action === 'submit') { + const result = await this.invoke({ sessionId: this.sessionId, queueEpoch: snapshot.queueEpoch, + action: 'get', turnId: record.request.message.turnId }); + if (result.receipt) { + this.accept(result, false); + await this.storage.put({ ...record, accepted: true }); + this.publish({ pending: (await this.storage.list(this.scope)).filter(record => !record.accepted || record.restoreIntent) }); + return result; + } + } + return this.transmit(record); + } + /** Commit edit intent before cancelling remotely, so a lost reply cannot + * let background refresh garbage-collect the only complete draft. */ + async prepareRestore(record: QueueOutboxRecord): Promise { + if (record.scope !== this.scope) throw new Error('Queue target changed'); + await this.storage.put({ ...record, restoreIntent: true }); + this.publish({ pending: (await this.storage.list(this.scope)).filter(record => !record.accepted || record.restoreIntent) }); + } + async receipt(turnId: string, expectedEpoch?: string): Promise { + const snapshot = await this.refresh(); + if (expectedEpoch && snapshot.queueEpoch !== expectedEpoch) throw new Error('queue_scope_expired: host queue owner changed'); + const result = await this.invoke({ sessionId: this.sessionId, queueEpoch: snapshot.queueEpoch, action: 'get', turnId }); + this.accept(result, false); + return result.receipt; + } + async savedDraft(turnId: string): Promise { + return (await this.storage.list(this.scope)).find(record => record.request.action === 'submit' && record.request.message.turnId === turnId); + } + + /** Explicitly dismiss an unresolved record; never sends or cancels host work. */ + async dismiss(record: QueueOutboxRecord): Promise { + if (record.scope !== this.scope) throw new Error('Queue target changed'); + await this.storage.remove(record.key); + this.publish({ pending: (await this.storage.list(this.scope)).filter(record => !record.accepted || record.restoreIntent) }); + } +} + +/** Refresh only while observed. Closing the view cancels no host work. */ +export function observeHostQueue(queue: HostDialogQueue): () => void { + let closed = false; + let timer: ReturnType | undefined; + let running = false; + const refresh = async () => { + if (closed || running) return; + if (timer) clearTimeout(timer); + running = true; + let delay = 2_000; + try { await queue.refresh(); } catch { delay = 10_000; } + finally { running = false; if (!closed) timer = setTimeout(() => void refresh(), delay); } + }; + const wake = () => { if (document.visibilityState !== 'hidden') void refresh(); }; + window.addEventListener('online', wake); + window.addEventListener('focus', wake); + document.addEventListener('visibilitychange', wake); + void refresh(); + return () => { closed = true; if (timer) clearTimeout(timer); + window.removeEventListener('online', wake); window.removeEventListener('focus', wake); + document.removeEventListener('visibilitychange', wake); }; +} diff --git a/src/shared/interactive-capabilities/catalog.json b/src/shared/interactive-capabilities/catalog.json index 95aa0609b3..14407d9c2e 100644 --- a/src/shared/interactive-capabilities/catalog.json +++ b/src/shared/interactive-capabilities/catalog.json @@ -402,6 +402,7 @@ }, "evidence": [ "command:start_dialog_turn", + "command:manage_dialog_queue", "command:steer_dialog_turn", "command:interrupt_dialog_turn", "command:cancel_dialog_turn", diff --git a/src/web-ui/src/flow_chat/components/HostPendingQueuePanel.tsx b/src/web-ui/src/flow_chat/components/HostPendingQueuePanel.tsx new file mode 100644 index 0000000000..fb99a31448 --- /dev/null +++ b/src/web-ui/src/flow_chat/components/HostPendingQueuePanel.tsx @@ -0,0 +1,85 @@ +import { pendingQueueManager } from '../services/flow-chat-manager/PendingQueueModule'; +import type { QueuedMessage } from '../types/flow-chat'; +import { useEffect, useState, useSyncExternalStore } from 'react'; +import { useI18n } from '@/infrastructure/i18n'; +import { Button } from '@openbitfun/ui'; +import { HostDialogQueue, observeHostQueue, type QueueOutboxRecord, type QueueItem } from '../../../../shared/dialog-queue/HostDialogQueue'; +import { + ChatComposerQueue, ChatComposerQueueHeader, ChatComposerQueueTitle, + ChatComposerQueueList, ChatComposerQueueItem, ChatComposerQueueItemContent, ChatComposerQueueItemActions, +} from '@openbitfun/ui/flow-chat'; + +export function HostPendingQueuePanel({ queue, onRestore }: { queue: HostDialogQueue; onRestore: (item: QueuedMessage) => boolean }) { + const { t } = useI18n('flow-chat'); + const view = useSyncExternalStore(queue.subscribe, queue.getSnapshot, queue.getSnapshot); + const [busy, setBusy] = useState(false); + const [error, setError] = useState(null); + useEffect(() => observeHostQueue(queue), [queue]); + const run = async (operation: () => Promise) => { + setBusy(true); setError(null); + try { await operation(); } catch (e) { setError(String(e)); } + finally { setBusy(false); } + }; + const restoreDraft = async (saved: QueueOutboxRecord, knownItem?: QueueItem) => { + if (saved.request.action !== 'submit') return; + await queue.prepareRestore(saved); + const message = saved.request.message; + const item = knownItem ?? await queue.receipt(message.turnId, saved.request.queueEpoch); + if (!item) throw new Error(t('hostQueue.unknown')); + if (item.status !== 'cancelled') { + const result = await queue.act(item, 'cancel'); + if (result.receipt?.status !== 'cancelled') throw new Error(t('hostQueue.unknown')); + } + const cache = saved.draft as Partial | undefined; + const restored = onRestore({ ...cache, id: item.turnId, sessionId: queue.sessionId, + content: message.content, displayMessage: message.displayContent, agentType: message.agentType, + timestamp: item.createdAtMs, status: 'queued', retryCount: cache?.retryCount ?? 0, + userMessageMetadata: message.metadata }); + if (!restored) { + // Preserve the complete draft if the composer changed during cancellation. + pendingQueueManager.enqueue({ ...cache, sessionId: queue.sessionId, + content: message.content, displayMessage: message.displayContent, + agentType: message.agentType, userMessageMetadata: message.metadata, + retryCount: 1, initialStatus: 'failed' }); + } + await queue.dismiss(saved); + }; + const items = view.snapshot?.items ?? []; + const visibleError = error || (view.error?.includes('Session is not loaded') ? null : view.error); + if (!items.length && !view.pending.length && !visibleError) return null; + return + {t('hostQueue.title')} +

{t('hostQueue.memoryNotice')}

+ {visibleError &&
{visibleError}
} + + {items.map(item => + +
{item.displayContent}
+ {item.status === 'blocked' ? t('hostQueue.blocked') : item.status === 'steering_pending' ? t('hostQueue.steeringPending') : t('hostQueue.queued')} + {item.attachmentCount > 0 && {t('hostQueue.attachments', { count: item.attachmentCount })}} + {item.reason &&

{item.reason}

} +
+ + + + + +
)} + {view.pending.map(record => + {t('hostQueue.unknown')} + {record.request.action === 'submit' &&
{record.request.message.displayContent ?? record.request.message.content}
} +
+ + {record.restoreIntent && record.request.action === 'submit' + ? + : } + + +
)} +
+
; +} diff --git a/src/web-ui/src/flow_chat/components/PendingQueuePanel.test.tsx b/src/web-ui/src/flow_chat/components/PendingQueuePanel.test.tsx index 711f7b9d11..ae93cc91f4 100644 --- a/src/web-ui/src/flow_chat/components/PendingQueuePanel.test.tsx +++ b/src/web-ui/src/flow_chat/components/PendingQueuePanel.test.tsx @@ -17,6 +17,9 @@ const mocks = vi.hoisted(() => ({ steerDialogTurn: vi.fn(), })); +vi.mock('../services/hostDialogQueue', () => ({ hostQueueSupported: () => false })); +vi.mock('./HostPendingQueuePanel', () => ({ HostPendingQueuePanel: () => null })); + vi.mock('react-i18next', () => ({ useTranslation: () => ({ t: (key: string) => key }), })); diff --git a/src/web-ui/src/flow_chat/components/PendingQueuePanel.tsx b/src/web-ui/src/flow_chat/components/PendingQueuePanel.tsx index 633150fe31..b2e0b87bfb 100644 --- a/src/web-ui/src/flow_chat/components/PendingQueuePanel.tsx +++ b/src/web-ui/src/flow_chat/components/PendingQueuePanel.tsx @@ -1,3 +1,6 @@ +import { HostPendingQueuePanel } from './HostPendingQueuePanel'; +import { hostQueueSupported, hostDialogQueue } from '../services/hostDialogQueue'; +import { getActiveSurfaceScope, onSurfaceActivated } from '@/infrastructure/peer-device/deviceSurface'; /** * Pending queue panel * @@ -52,7 +55,7 @@ interface PendingQueuePanelProps { onRestoreToComposer: (item: QueuedMessage) => boolean; } -export function PendingQueuePanel({ +function LegacyPendingQueuePanel({ sessionId, className, onRestoreToComposer, @@ -374,3 +377,17 @@ export function PendingQueuePanel({ } export default PendingQueuePanel; + + +export function PendingQueuePanel(props: PendingQueuePanelProps): JSX.Element | null { + const [, setRevision] = useState(0); + useEffect(() => onSurfaceActivated(() => setRevision(value => value + 1)), []); + const supported = props.sessionId && hostQueueSupported(props.sessionId); + const scope = getActiveSurfaceScope(); + const { t } = useTranslation('flow-chat'); + return <> + {supported && props.sessionId && pendingQueueManager.list(props.sessionId).length > 0 &&

{t('hostQueue.legacy')}

} + {supported && props.sessionId && } + + ; +} diff --git a/src/web-ui/src/flow_chat/services/SessionRollbackService.ts b/src/web-ui/src/flow_chat/services/SessionRollbackService.ts index 3cb9135f8b..8e061b9c56 100644 --- a/src/web-ui/src/flow_chat/services/SessionRollbackService.ts +++ b/src/web-ui/src/flow_chat/services/SessionRollbackService.ts @@ -1,3 +1,4 @@ +import { hostQueueSupported, hostDialogQueue } from './hostDialogQueue'; import { agentAPI } from '@/infrastructure/api'; import { globalEventBus } from '@/infrastructure/event-bus'; import { createLogger } from '@/shared/utils/logger'; @@ -80,6 +81,10 @@ export async function rollbackSessionToTurn( if (resolveSessionDriverId(request.sessionId, flowChatStore.getState().sessions.get(request.sessionId)) === 'dispatch') { throw new Error('History rollback is unavailable for a detached remote session.'); } + if (hostQueueSupported(request.sessionId)) { + const queue = await hostDialogQueue(request.sessionId).refresh(); + if (queue.items.length) throw new Error('Clear the host message queue before changing Session history'); + } const lease = request.lease ?? tryBeginSessionMutation(request.sessionId, request.kind, request.targetTurnId); if (!lease) { diff --git a/src/web-ui/src/flow_chat/services/flow-chat-manager/MessageModule.ts b/src/web-ui/src/flow_chat/services/flow-chat-manager/MessageModule.ts index 8df4105111..07c49ea668 100644 --- a/src/web-ui/src/flow_chat/services/flow-chat-manager/MessageModule.ts +++ b/src/web-ui/src/flow_chat/services/flow-chat-manager/MessageModule.ts @@ -1,3 +1,4 @@ +import { hostQueueSupported, hostDialogQueue, queueImageAttachments } from '../hostDialogQueue'; /** * Message handling module * Shared submission choreography: busy-gate planning, queueing, mode @@ -268,6 +269,19 @@ export async function sendMessage( await driverForSession(sessionId, session).steer(context, sessionId, draft); return; } + if (hostQueueSupported(sessionId)) { + beginSubmission(); + try { + const queue = hostDialogQueue(sessionId); + await queue.submit({ content: message, displayContent: displayMessage, + agentType: agentType?.trim() || session.mode || 'Standard', + attachments: queueImageAttachments(options?.imageContexts), metadata: options?.userMessageMetadata ?? {} }, + { composerDraft: options?.pendingQueueDraft, imageContexts: options?.imageContexts, imageDisplayData: options?.imageDisplayData }); + surfaceScopeAtSend.assertCurrent('accept queued message'); + completeSessionSend(sendCoordinationKey, sendAttempt); + return; + } finally { endSubmission(); } + } try { const item = pendingQueueManager.enqueue({ sessionId, @@ -540,6 +554,7 @@ export async function drainPendingQueue( sessionId: string, options?: { allowInterruptedRecoveryAbandon?: boolean }, ): Promise { + if (hostQueueSupported(sessionId) && !options?.allowInterruptedRecoveryAbandon) return; if (isRuntimeSessionAttachmentInFlight(getActiveSurfaceId(), sessionId) || isRuntimeSessionProjectionStale(getActiveSurfaceId(), sessionId)) { return; diff --git a/src/web-ui/src/flow_chat/services/flow-chat-manager/PendingQueueModule.ts b/src/web-ui/src/flow_chat/services/flow-chat-manager/PendingQueueModule.ts index a19184b64d..e527c623a4 100644 --- a/src/web-ui/src/flow_chat/services/flow-chat-manager/PendingQueueModule.ts +++ b/src/web-ui/src/flow_chat/services/flow-chat-manager/PendingQueueModule.ts @@ -1,7 +1,10 @@ /** * Pending queue module * - * Frontend-side message queue used while a session's current dialog turn is + * Legacy/fallback frontend queue and explicit recovery drafts. Native hosts + * advertising dialog_queue_v1 use HostDialogQueue for new submissions. + * + * This frontend queue is used while a session's current dialog turn is * still running. Items are kept here (NOT submitted to the backend scheduler) * until the session returns to IDLE, at which point the head item is drained * via the regular `sendMessage` path. Users may also pop an item early through diff --git a/src/web-ui/src/flow_chat/services/hostDialogQueue.ts b/src/web-ui/src/flow_chat/services/hostDialogQueue.ts new file mode 100644 index 0000000000..b95eda805c --- /dev/null +++ b/src/web-ui/src/flow_chat/services/hostDialogQueue.ts @@ -0,0 +1,52 @@ +import { HostDialogQueue, type QueueAttachment } from '../../../../shared/dialog-queue/HostDialogQueue'; +import { api } from '@/infrastructure/api/service-api/ApiClient'; +import { getActiveSurfaceScope, isLocalSurface } from '@/infrastructure/peer-device/deviceSurface'; +import { peerConnectionManager } from '@/infrastructure/peer-device/PeerConnectionManager'; +import { accountIdentityService } from '@/infrastructure/account-identity/AccountIdentityService'; +import { FlowChatStore } from '../store/FlowChatStore'; +import { isAcpFlowSession } from '../utils/acpSession'; +import { resolveSessionDriverId } from '../session-drivers/resolve'; +import { translateAgentIdentityFields } from '../../../../shared/agent-harness/wire'; + +const clients = new Map(); +function currentAccount(): string { + const user = accountIdentityService.getSnapshot().me?.user; + return String(user?.accountId || user?.githubId || 'local'); +} +export function hostQueueSupported(sessionId: string): boolean { + const session = FlowChatStore.getInstance().getState().sessions.get(sessionId); + if (!session || isAcpFlowSession(session) || resolveSessionDriverId(sessionId, session) !== 'local') return false; + const scope = getActiveSurfaceScope(); + return isLocalSurface(scope.surfaceId) + || peerConnectionManager.get(scope.surfaceId)?.getState().capabilities.dialogQueueV1 === true; +} +export function hostDialogQueue(sessionId: string): HostDialogQueue { + const scope = getActiveSurfaceScope(); + const account = currentAccount(); + const owner = JSON.stringify([account, scope.surfaceId, sessionId]); + const key = JSON.stringify([owner, scope.epoch]); + let client = clients.get(key); + if (!client) { + client = new HostDialogQueue(owner, sessionId, async request => { + scope.assertCurrent('send queue operation'); + if (currentAccount() !== account) throw new Error('Queue account changed'); + const result = await api.invoke( + 'manage_dialog_queue', { request: translateAgentIdentityFields(request, 'legacy') }); + scope.assertCurrent('apply queue operation'); + if (currentAccount() !== account) throw new Error('Queue account changed'); + return result; + }); + clients.set(key, client); + } + return client; +} +export function queueImageAttachments(images?: unknown[]): QueueAttachment[] { + return (images ?? []).map(value => { + const image = value as { id: string; data_url?: string; image_path?: string; mime_type?: string; metadata?: unknown }; + return { kind: 'remote_image', id: image.id, metadata: { + ...(image.data_url ? { dataUrl: image.data_url } : {}), + ...(image.image_path ? { imagePath: image.image_path } : {}), + mimeType: image.mime_type, metadata: image.metadata, + } }; + }); +} diff --git a/src/web-ui/src/flow_chat/session-drivers/local/LocalSessionDriver.test.ts b/src/web-ui/src/flow_chat/session-drivers/local/LocalSessionDriver.test.ts index 00f4f67dbb..6189b768da 100644 --- a/src/web-ui/src/flow_chat/session-drivers/local/LocalSessionDriver.test.ts +++ b/src/web-ui/src/flow_chat/session-drivers/local/LocalSessionDriver.test.ts @@ -211,3 +211,32 @@ describe('localSessionDriver review repair permissions', () => { expect(mockStartAgenticDialogTurn).not.toHaveBeenCalled(); }); }); + + +const queueMocks = vi.hoisted(() => ({ supported: vi.fn(() => false), submit: vi.fn() })); +vi.mock('../../services/hostDialogQueue', () => ({ + hostQueueSupported: queueMocks.supported, + hostDialogQueue: () => ({ submit: queueMocks.submit }), + queueImageAttachments: () => [], +})); + +describe('host queue submissions', () => { + beforeEach(() => { vi.clearAllMocks(); queueMocks.supported.mockReturnValue(true); queueMocks.submit.mockResolvedValue({ receipt: { status: 'queued' } }); }); + it('lets authoritative host events start a turn without clearing an existing active presentation', async () => { + const { context, session, addedTurns } = createHarness([]); + session.mode = 'Standard'; + context.contentBuffers.set(SESSION_ID, 'active output'); + context.activeTextItems.set(SESSION_ID, 'active item'); + const input = { ...startTurnInput(session), acpClientId: undefined, currentAgentType: 'Standard', isFirstMessage: false, + options: { turnId: 'stable-request-id' } }; + const tracker = { createdLocalTurnId: null, hostAcceptedTurn: false }; + await localSessionDriver.startTurn(context, input, tracker); + expect(queueMocks.submit).toHaveBeenCalledWith(expect.objectContaining({ content: 'hello' }), expect.any(Object), 'stable-request-id'); + expect(mockTransition).not.toHaveBeenCalled(); + expect(mockStartAgenticDialogTurn).not.toHaveBeenCalled(); + expect(context.contentBuffers.get(SESSION_ID)).toBe('active output'); + expect(context.activeTextItems.get(SESSION_ID)).toBe('active item'); + expect(addedTurns).toHaveLength(0); + expect(tracker.hostAcceptedTurn).toBe(true); + }); +}); diff --git a/src/web-ui/src/flow_chat/session-drivers/local/LocalSessionDriver.ts b/src/web-ui/src/flow_chat/session-drivers/local/LocalSessionDriver.ts index 5dd6eb12a9..cf58141a23 100644 --- a/src/web-ui/src/flow_chat/session-drivers/local/LocalSessionDriver.ts +++ b/src/web-ui/src/flow_chat/session-drivers/local/LocalSessionDriver.ts @@ -1,3 +1,4 @@ +import { hostQueueSupported, hostDialogQueue, queueImageAttachments } from '../../services/hostDialogQueue'; import { requireSessionOwningWorkspaceId, sessionOwningWorkspaceId } from '../../utils/sessionOrdering'; /** * Local session driver: the default flavor backed by this machine's (or the @@ -321,6 +322,67 @@ export const localSessionDriver: SessionDriver = { options, } = input; + const prepareSubmission = async () => { + if (readySession.config.worktreeIsolationRequested !== undefined) { + const materialization = sessionWorktreeMaterializationPlan(readySession); + if (materialization) { + log.info('Materializing requested worktree after prompt submission', { + sessionId, + enabled: materialization.enabled, + projectWorkspaceId: materialization.projectWorkspaceId, + projectWorkspacePath: materialization.projectWorkspacePath, + }); + const result = await worktreeAPI.bindSession( + sessionId, + materialization.enabled, + globalThis.crypto?.randomUUID?.() ?? `worktree-first-turn-${Date.now()}`, + materialization, + ); + surfaceScope.assertCurrent('bind session worktree'); + context.flowChatStore.updateSessionExecutionTarget(sessionId, { + workspacePath: result.workspacePath, + projectWorkspacePath: result.projectWorkspacePath, + workspaceId: result.workspaceId, + projectWorkspaceId: result.projectWorkspaceId, + executionTarget: result.executionTarget, + }); + if (result.retainedWorktreePath) { + log.warn('Released worktree retained because it contains local work', { + sessionId, + retainedWorktreePath: result.retainedWorktreePath, + }); + } + } + context.flowChatStore.setSessionWorktreeIsolationRequested(sessionId, undefined); + } + + if (isFirstMessage) { + applyGeneratingTitlePlaceholder(context, sessionId, message); + } + + if (!acpClientId) { + await syncSessionModelSelection(context, sessionId, currentAgentType, surfaceScope); + } + }; + if (!acpClientId && hostQueueSupported(sessionId) && (!options?.execution || options.execution.kind === 'standard')) { + if (readySession.isHistorical || context.pendingHistoryLoads.has(surfaceScope.key('history-load', surfaceScope.epoch, sessionId))) { + throw new Error('Session history is still restoring, please retry once loading finishes'); + } + await prepareSubmission(); + await inheritReviewPermissionMode(readySession, context.flowChatStore.getState().sessions, + () => surfaceScope.assertCurrent('inherit review session permission mode')); + tracker.hostSubmitStarted = true; + await hostDialogQueue(sessionId).submit({ content: message, displayContent: displayMessage, + agentType: currentAgentType, attachments: queueImageAttachments(options?.imageContexts), + metadata: options?.userMessageMetadata ?? {} }, + { composerDraft: options?.pendingQueueDraft, imageContexts: options?.imageContexts, imageDisplayData: options?.imageDisplayData }, options?.turnId); + tracker.hostAcceptedTurn = true; + surfaceScope.assertCurrent('accept host message'); + context.flowChatStore.updateSessionLastSubmittedMode(sessionId, currentAgentType); + if (isFirstMessage) await updateSessionMetadata(context, sessionId, ['titleMetadata']); + return 'completed'; + } + const dialogTurnId = options?.turnId?.trim() || `dialog_${Date.now()}_${Math.random().toString(36).substr(2, 9)}`; const hasImages = (options?.imageContexts?.length ?? 0) > 0; @@ -389,46 +451,7 @@ export const localSessionDriver: SessionDriver = { metadata: { sessionId: sessionId, dialogTurnId } }); - if (readySession.config.worktreeIsolationRequested !== undefined) { - const materialization = sessionWorktreeMaterializationPlan(readySession); - if (materialization) { - log.info('Materializing requested worktree after prompt submission', { - sessionId, - enabled: materialization.enabled, - projectWorkspaceId: materialization.projectWorkspaceId, - projectWorkspacePath: materialization.projectWorkspacePath, - }); - const result = await worktreeAPI.bindSession( - sessionId, - materialization.enabled, - globalThis.crypto?.randomUUID?.() ?? `worktree-first-turn-${Date.now()}`, - materialization, - ); - surfaceScope.assertCurrent('bind session worktree'); - context.flowChatStore.updateSessionExecutionTarget(sessionId, { - workspacePath: result.workspacePath, - projectWorkspacePath: result.projectWorkspacePath, - workspaceId: result.workspaceId, - projectWorkspaceId: result.projectWorkspaceId, - executionTarget: result.executionTarget, - }); - if (result.retainedWorktreePath) { - log.warn('Released worktree retained because it contains local work', { - sessionId, - retainedWorktreePath: result.retainedWorktreePath, - }); - } - } - context.flowChatStore.setSessionWorktreeIsolationRequested(sessionId, undefined); - } - - if (isFirstMessage) { - applyGeneratingTitlePlaceholder(context, sessionId, message); - } - - if (!acpClientId) { - await syncSessionModelSelection(context, sessionId, currentAgentType, surfaceScope); - } + await prepareSubmission(); const updatedSession = context.flowChatStore.getState().sessions.get(sessionId); if (!updatedSession) { diff --git a/src/web-ui/src/infrastructure/api/generated/remoteSurface.ts b/src/web-ui/src/infrastructure/api/generated/remoteSurface.ts index f0ea6434c2..697acfb1a1 100644 --- a/src/web-ui/src/infrastructure/api/generated/remoteSurface.ts +++ b/src/web-ui/src/infrastructure/api/generated/remoteSurface.ts @@ -1,6 +1,6 @@ // Generated by scripts/generate-interactive-capabilities.mjs; do not edit. // Source: openbitfun_product_domains::remote_surface (Product Operation Registry). -export const REMOTE_SURFACE_REGISTRY_DIGEST = "fnv1a64:01a7a1c7756afe23" as const; +export const REMOTE_SURFACE_REGISTRY_DIGEST = "fnv1a64:4a41f1e11a6651e1" as const; /** * Registered Tauri commands the Peer Device controller keeps on the controller @@ -163,6 +163,7 @@ export const PEER_CONTROLLER_LOCAL_COMMANDS: ReadonlySet = new Set([ /** Every capability id a peer host may advertise in `peer_mode_ping`. */ export const PEER_HOST_CAPABILITY_IDS = [ + "dialog_queue_v1", "idempotent_dialog_submit", "inline_image_attachments_v1", "btw_initial_model_selection_v1", @@ -190,6 +191,7 @@ export const PEER_HOST_ADVERTISED_CAPABILITIES: Readonly< Record<'desktop' | 'cli', readonly PeerHostCapabilityId[]> > = { desktop: [ + "dialog_queue_v1", "idempotent_dialog_submit", "inline_image_attachments_v1", "btw_initial_model_selection_v1", @@ -210,6 +212,7 @@ export const PEER_HOST_ADVERTISED_CAPABILITIES: Readonly< "user_question_interaction_v1", ], cli: [ + "dialog_queue_v1", "control_conversation_v1", "control_conversation_reset_v1", "idempotent_dialog_submit", diff --git a/src/web-ui/src/infrastructure/peer-device/PeerConnectionManager.ts b/src/web-ui/src/infrastructure/peer-device/PeerConnectionManager.ts index 49446952b3..e2603c6b20 100644 --- a/src/web-ui/src/infrastructure/peer-device/PeerConnectionManager.ts +++ b/src/web-ui/src/infrastructure/peer-device/PeerConnectionManager.ts @@ -87,6 +87,7 @@ export interface PeerHostCapabilities { */ readonly userQuestionResponse: boolean | null; readonly userQuestionInteraction?: boolean; + readonly dialogQueueV1?: boolean; /** * Which kind of host answered `peer_mode_ping` (`"desktop"` | `"cli"`). * `null` = the host did not advertise `host_type` (even older host, or the @@ -521,6 +522,7 @@ export class PeerConnectionManager { workspaceIdReferencesV1: caps?.workspace_id_references_v1 === true, userQuestionResponse, userQuestionInteraction: caps?.user_question_interaction_v1 === true, + dialogQueueV1: caps?.dialog_queue_v1 === true, hostKind, }; } @@ -754,6 +756,7 @@ function capabilitiesEqual( a.workspaceIdReferencesV1 === b.workspaceIdReferencesV1 && a.userQuestionResponse === b.userQuestionResponse && a.userQuestionInteraction === b.userQuestionInteraction && + a.dialogQueueV1 === b.dialogQueueV1 && a.hostKind === b.hostKind; } diff --git a/src/web-ui/src/locales/en-US/flow-chat.json b/src/web-ui/src/locales/en-US/flow-chat.json index f11ca073d8..15b919ab8f 100644 --- a/src/web-ui/src/locales/en-US/flow-chat.json +++ b/src/web-ui/src/locales/en-US/flow-chat.json @@ -2879,5 +2879,23 @@ "loadFailed": "Unable to open the assistant", "retry": "Retry", "loading": "Opening assistant…" + }, + "hostQueue": { + "title": "Host message queue", + "memoryNotice": "Accepted messages run while this page is closed. Restarting the execution device clears pending messages.", + "queued": "Queued", + "blocked": "Waiting for recovery", + "steeringPending": "Waiting to steer", + "sendNow": "Send now", + "cancel": "Remove from queue", + "refresh": "Refresh", + "unknown": "Delivery is unconfirmed. Check before sending again.", + "checkRetry": "Check / retry", + "copyDraft": "Copy text to composer", + "dismiss": "Dismiss local reminder", + "attachments": "{{count}} attachments", + "edit": "Restore draft", + "noDraft": "The complete draft is only available on the device that submitted it.", + "legacy": "Saved drafts from the previous queue. Send or restore each explicitly." } } diff --git a/src/web-ui/src/locales/zh-CN/flow-chat.json b/src/web-ui/src/locales/zh-CN/flow-chat.json index 4b29026a10..3980762c96 100644 --- a/src/web-ui/src/locales/zh-CN/flow-chat.json +++ b/src/web-ui/src/locales/zh-CN/flow-chat.json @@ -2879,5 +2879,23 @@ "loadFailed": "暂时无法打开助手", "retry": "重试", "loading": "正在打开助手…" + }, + "hostQueue": { + "title": "宿主消息队列", + "memoryNotice": "消息接受后,关闭此页面仍会执行。执行设备重启会清空待执行消息。", + "queued": "排队中", + "blocked": "等待恢复", + "steeringPending": "等待注入", + "sendNow": "立即发送", + "cancel": "移出队列", + "refresh": "刷新", + "unknown": "尚未确认是否送达,请先检查再重发。", + "checkRetry": "检查 / 重试", + "copyDraft": "复制文字到输入框", + "dismiss": "关闭本地提醒", + "attachments": "{{count}} 个附件", + "edit": "恢复草稿", + "noDraft": "完整草稿仅在提交消息的设备上可用。", + "legacy": "旧队列中的草稿,请逐条发送或恢复。" } } diff --git a/src/web-ui/src/locales/zh-TW/flow-chat.json b/src/web-ui/src/locales/zh-TW/flow-chat.json index 8ac22da03d..714fd03da9 100644 --- a/src/web-ui/src/locales/zh-TW/flow-chat.json +++ b/src/web-ui/src/locales/zh-TW/flow-chat.json @@ -2879,5 +2879,23 @@ "loadFailed": "暫時無法開啟助手", "retry": "重試", "loading": "正在開啟助手…" + }, + "hostQueue": { + "title": "主機訊息佇列", + "memoryNotice": "訊息接受後,關閉此頁面仍會執行。執行裝置重新啟動會清空待執行訊息。", + "queued": "排隊中", + "blocked": "等待恢復", + "steeringPending": "等待注入", + "sendNow": "立即傳送", + "cancel": "移出佇列", + "refresh": "重新整理", + "unknown": "尚未確認是否送達,請先檢查再重新傳送。", + "checkRetry": "檢查 / 重試", + "copyDraft": "複製文字到輸入框", + "dismiss": "關閉本機提醒", + "attachments": "{{count}} 個附件", + "edit": "恢復草稿", + "noDraft": "完整草稿僅在提交訊息的裝置上可用。", + "legacy": "舊佇列中的草稿,請逐條傳送或恢復。" } } From 133a437d53983ad57888e3370fa3791b3002eb04 Mon Sep 17 00:00:00 2001 From: Bob Lee Date: Mon, 21 Sep 2026 23:38:36 +0800 Subject: [PATCH 2/2] test(remote): cover host queue rollback guard and update mobile mocks --- .../MobileChatSubmission.test.tsx | 1 + .../services/SessionRollbackService.test.ts | 16 ++++++++++++++++ 2 files changed, 17 insertions(+) diff --git a/src/web-ui/src/app/components/RemoteConnectDialog/MobileChatSubmission.test.tsx b/src/web-ui/src/app/components/RemoteConnectDialog/MobileChatSubmission.test.tsx index 96d77faf69..f0197d6b4f 100644 --- a/src/web-ui/src/app/components/RemoteConnectDialog/MobileChatSubmission.test.tsx +++ b/src/web-ui/src/app/components/RemoteConnectDialog/MobileChatSubmission.test.tsx @@ -67,6 +67,7 @@ describe('mobile chat submission acknowledgement', () => { subscribeSessionStream: vi.fn().mockResolvedValue({ close() {}, wake() {}, async loadOlder() {} }), getSessionMessages: vi.fn().mockResolvedValue({ messages: [], has_more: false }), getModelCatalog: vi.fn().mockResolvedValue({ version: 1, models: [], default_models: {}, session_model_id: 'auto' }), + supportsHostCapability: vi.fn().mockReturnValue(false), sendMessage, } as unknown as RemoteSessionManager; }); diff --git a/src/web-ui/src/flow_chat/services/SessionRollbackService.test.ts b/src/web-ui/src/flow_chat/services/SessionRollbackService.test.ts index df5e020b00..6c4a752c48 100644 --- a/src/web-ui/src/flow_chat/services/SessionRollbackService.test.ts +++ b/src/web-ui/src/flow_chat/services/SessionRollbackService.test.ts @@ -1,5 +1,11 @@ import { beforeEach, describe, expect, it, vi } from 'vitest'; +const hostQueueMock = vi.hoisted(() => ({ supported: vi.fn(() => false), refresh: vi.fn() })); +vi.mock('./hostDialogQueue', () => ({ + hostQueueSupported: hostQueueMock.supported, + hostDialogQueue: () => ({ refresh: hostQueueMock.refresh }), +})); + const agentApiMock = vi.hoisted(() => ({ rollbackSessionToTurn: vi.fn() })); const eventBusMock = vi.hoisted(() => ({ emit: vi.fn() })); const loadSessionHistory = vi.hoisted(() => vi.fn(async () => undefined)); @@ -37,6 +43,8 @@ describe('SessionRollbackService', () => { }); beforeEach(() => { vi.clearAllMocks(); + hostQueueMock.supported.mockReturnValue(false); + hostQueueMock.refresh.mockResolvedValue({ items: [] }); loadSessionHistory.mockResolvedValue(undefined); sessions.clear(); historyViews.clear(); @@ -59,6 +67,14 @@ describe('SessionRollbackService', () => { }); }); + it('refuses host-owned queued work before invoking rollback', async () => { + hostQueueMock.supported.mockReturnValue(true); + hostQueueMock.refresh.mockResolvedValue({ items: [{ turnId: 'queued-turn', status: 'blocked' }] }); + await expect(rollbackSessionToTurn({ sessionId: 'session-1', targetTurnId: 'turn-7', kind: 'rollback' })) + .rejects.toThrow('Clear the host message queue'); + expect(agentApiMock.rollbackSessionToTurn).not.toHaveBeenCalled(); + }); + it('refuses path-only mutations until legacy identity has been hydrated', async () => { const session = sessions.get('session-1'); delete session.workspaceId;