From eb4fb8ea7696f72ea7323fcb4cee6293f82cac44 Mon Sep 17 00:00:00 2001 From: Bob Lee Date: Wed, 23 Sep 2026 15:22:15 +0800 Subject: [PATCH 1/2] fix(queue): accept new prompts after recoverable interruption --- .../coordination/host_message_queue.rs | 2 +- .../coordination/host_message_queue_tests.rs | 171 ++++++++++++++++++ .../src/agentic/coordination/scheduler.rs | 29 ++- 3 files changed, 183 insertions(+), 19 deletions(-) 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 index 929953f548..352a64326f 100644 --- a/src/crates/assembly/core/src/agentic/coordination/host_message_queue.rs +++ b/src/crates/assembly/core/src/agentic/coordination/host_message_queue.rs @@ -253,7 +253,7 @@ impl DialogScheduler { .await .map_err(|e| PortError::new(PortErrorKind::Backend, e.to_string()))? } - async fn execute_queue_request( + pub(super) async fn execute_queue_request( &self, request: DialogQueueRequest, ) -> PortResult { 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 index 3cadbe6c0a..2264fa9c1c 100644 --- 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 @@ -600,3 +600,174 @@ async fn host_queue_interrupted_target_cannot_consume_a_blocked_injection_on_res .is_empty()); assert_eq!(scheduler.queue_depth("host-queue-session"), 2); } + +#[test] +fn host_queue_new_prompt_after_stop_supersedes_interruption() { + // This fixture polls coordinator admission inline to retain the task-local + // model configuration. Give its large debug future a dedicated test stack. + std::thread::Builder::new() + .stack_size(16 * 1024 * 1024) + .spawn(|| { + tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .unwrap() + .block_on(async { + for has_held_message in [false, true] { + let (scheduler, sessions, _, root) = test_scheduler_with_persistence(true); + let id = "host-queue-session"; + let workspace = + fixture_workspace_dir(root.path().join("stopped-workspace")); + sessions + .create_session_with_id( + Some(id.into()), + "Stopped".into(), + "Standard".into(), + SessionConfig { + workspace_path: Some(workspace.to_string_lossy().into_owned()), + ..Default::default() + }, + ) + .await + .unwrap(); + sessions + .start_dialog_turn( + id, + "Standard".into(), + "original work".into(), + Some("stopped-turn".into()), + None, + None, + ) + .await + .unwrap(); + let epoch = scheduler + .manage_host_queue(request(None, Action::List)) + .await + .unwrap() + .queue_epoch; + if has_held_message { + scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Submit { + message: message("previously-queued"), + }, + )) + .await + .unwrap(); + scheduler.hold_managed_queue(id, "Turn interrupted"); + } + sessions + .mark_dialog_turn_interrupted(id, "stopped-turn") + .await + .unwrap(); + sessions + .update_session_state_for_turn_if_processing( + id, + "stopped-turn", + SessionState::Idle, + ) + .await + .unwrap(); + assert!(sessions + .latest_dialog_turn_holds_dispatch(id) + .await + .unwrap()); + assert!(scheduler.try_start_next_queued(id).await.unwrap().is_none()); + + let failed = TEST_MODEL_RESOLUTION_AI_CONFIG + .scope( + AIConfig::default(), + Box::pin(scheduler.execute_queue_request(request( + Some(&epoch), + Action::Submit { + message: message("new-prompt"), + }, + ))), + ) + .await; + assert!(failed.is_err(), "missing model must reject admission"); + assert!( + sessions + .latest_dialog_turn_holds_dispatch(id) + .await + .unwrap(), + "failed submission must preserve the stopped turn for recovery" + ); + assert_eq!(sessions.get_turn_count(id), 1); + + let ai_config = AIConfig { + models: vec![AIModelConfig { + id: "queue-stop-test-model".into(), + name: "Queue stop 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-stop-test-model"), + ) + .await + .unwrap(); + // Execute the same host-side request body inline so the model fixture's + // task-local scope covers admission; no model response is needed. + let submit = request( + Some(&epoch), + Action::Submit { + message: message("new-prompt"), + }, + ); + let snapshot = TEST_MODEL_RESOLUTION_AI_CONFIG + .scope( + ai_config, + Box::pin(scheduler.execute_queue_request(submit.clone())), + ) + .await + .expect("a fresh user prompt must supersede the stopped turn"); + assert_eq!(snapshot.receipt.unwrap().status, Status::Started); + assert!(!sessions + .latest_dialog_turn_holds_dispatch(id) + .await + .unwrap()); + assert_eq!(sessions.get_turn_count(id), 2); + let storage = sessions.effective_session_storage_path(id).await.unwrap(); + let stopped = sessions + .persistence_manager() + .load_dialog_turn(&storage, id, 0) + .await + .unwrap() + .unwrap(); + assert_eq!(stopped.turn_id, "stopped-turn"); + assert!( + stopped.recovery.is_none(), + "accepted prompt must retire old recovery" + ); + + if has_held_message { + assert!(snapshot + .items + .iter() + .any(|item| item.turn_id == "previously-queued" + && item.status == Status::Blocked)); + } + // Reconnect/retry must return the existing receipt, not start again. + scheduler.manage_host_queue(submit).await.unwrap(); + assert_eq!(sessions.get_turn_count(id), 2); + let _ = scheduler + .coordinator + .cancel_dialog_turn(id, "new-prompt") + .await; + } + }); + }) + .unwrap() + .join() + .unwrap(); +} diff --git a/src/crates/assembly/core/src/agentic/coordination/scheduler.rs b/src/crates/assembly/core/src/agentic/coordination/scheduler.rs index 2a302e3e04..54c77581ad 100644 --- a/src/crates/assembly/core/src/agentic/coordination/scheduler.rs +++ b/src/crates/assembly/core/src/agentic/coordination/scheduler.rs @@ -1292,18 +1292,9 @@ 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(), - )); - } + // A newly submitted user prompt supersedes recoverable interruption even + // when it arrives through the host queue. Existing queued work still + // stays parked in try_start_next_queued_locked until that user decision. let interrupted_turn_to_abandon = if interrupted_hold && !matches!( queued_turn.policy.trigger_source, @@ -1322,12 +1313,14 @@ impl DialogScheduler { .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 state_fact = if self.active_turns.contains(&session_id) + || interrupted_hold + || (held_user_messages && interrupted_turn_to_abandon.is_none()) + { + DialogSessionStateFact::Processing + } else { + Self::session_state_fact(state.as_ref()) + }; let queue_has_items = self.queues.has_items(&session_id); if matches!( From 40e03484aeed9f3e2973f6caa879bf1edcd1fc32 Mon Sep 17 00:00:00 2001 From: Bob Lee Date: Wed, 23 Sep 2026 15:48:27 +0800 Subject: [PATCH 2/2] fix(queue): protect new work from retiring turns and preserve draft recovery --- .../coordination/host_message_queue.rs | 81 +++- .../coordination/host_message_queue_tests.rs | 350 ++++++++++++++++++ .../src/agentic/coordination/scheduler.rs | 82 +++- .../components/HostPendingQueuePanel.test.tsx | 33 ++ .../components/HostPendingQueuePanel.tsx | 14 +- src/web-ui/src/locales/en-US/flow-chat.json | 2 +- src/web-ui/src/locales/zh-CN/flow-chat.json | 2 +- src/web-ui/src/locales/zh-TW/flow-chat.json | 2 +- 8 files changed, 538 insertions(+), 28 deletions(-) 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 index 352a64326f..728f10c77f 100644 --- a/src/crates/assembly/core/src/agentic/coordination/host_message_queue.rs +++ b/src/crates/assembly/core/src/agentic/coordination/host_message_queue.rs @@ -28,6 +28,9 @@ struct Entry { fingerprint: String, view: DialogQueueItem, held: Option, + // User work accepted after this turn settled must outlive its delayed + // outcome cleanup. This is host-local admission bookkeeping, not wire data. + after_terminal_turn: Option, } impl Default for QueueSession { fn default() -> Self { @@ -54,6 +57,31 @@ impl HostQueueState { .get(session) .is_some_and(|s| s.entries.contains_key(turn)) } + pub(super) fn mark_after_terminal_turn(&mut self, session: &str, turn: &str, previous: String) { + if let Some(e) = self + .sessions + .get_mut(session) + .and_then(|s| s.entries.get_mut(turn)) + { + e.after_terminal_turn = Some(previous); + } + } + pub(super) fn admitted_after( + &self, + session: &str, + previous: &str, + ) -> std::collections::HashSet { + self.sessions + .get(session) + .map(|s| { + s.entries + .iter() + .filter(|(_, e)| e.after_terminal_turn.as_deref() == Some(previous)) + .map(|(id, _)| id.clone()) + .collect() + }) + .unwrap_or_default() + } pub(super) fn pending_held(&self, session: &str) -> usize { self.sessions.get(session).map_or(0, |s| { s.entries @@ -178,6 +206,14 @@ impl HostQueueState { } impl DialogScheduler { pub(super) fn hold_managed_queue(&self, session: &str, reason: &str) { + self.hold_managed_queue_for_outcome(session, reason, None); + } + pub(super) fn hold_managed_queue_for_outcome( + &self, + session: &str, + reason: &str, + previous: Option<&str>, + ) { let ids: Vec = self .queue_state() .sessions @@ -185,7 +221,11 @@ impl DialogScheduler { .map(|s| { s.entries .iter() - .filter(|(_, e)| e.view.status == DialogQueueStatus::Queued) + .filter(|(_, e)| { + e.view.status == DialogQueueStatus::Queued + && previous + .is_none_or(|id| e.after_terminal_turn.as_deref() != Some(id)) + }) .map(|(id, _)| id.clone()) .collect() }) @@ -346,6 +386,7 @@ impl DialogScheduler { fingerprint: digest, view, held: None, + after_terminal_turn: None, }, ); s.order.push(message.turn_id.clone()); @@ -445,6 +486,7 @@ impl DialogScheduler { )); } } + let mut interrupted_turn_to_abandon = None; if let DialogQueueAction::Promote { expected_active_turn_id, .. @@ -460,13 +502,21 @@ impl DialogScheduler { { 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()))? + if expected_active_turn_id.is_none() && self.active_turns.contains(session) { + return Err(error("queue_conflict: previous turn is still retiring")); + } + if expected_active_turn_id.is_none() + && 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")); + interrupted_turn_to_abandon = self + .session_manager + .get_session(session) + .and_then(|s| s.dialog_turn_ids.last().cloned()); + self.hold_managed_queue(session, "Turn interrupted; retry this message explicitly"); } } let held = { @@ -556,11 +606,22 @@ impl DialogScheduler { DialogQueueAction::Promote { expected_active_turn_id: None, .. - } => { - if let Err(e) = self.start_turn(session, &turn).await { + } => match self.start_turn(session, &turn).await { + Err(e) => { self.queue_state().hold(session, &turn, &e.to_string()); } - } + Ok(_) => { + if let Err(e) = self + .abandon_superseded_interrupted_turn( + session, + interrupted_turn_to_abandon.as_deref(), + ) + .await + { + warn!("Failed to retire interrupted recovery after explicit queue promotion: session_id={}, error={}", session, e); + } + } + }, _ => unreachable!(), } { 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 index 2264fa9c1c..23c6480fb0 100644 --- 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 @@ -771,3 +771,353 @@ fn host_queue_new_prompt_after_stop_supersedes_interruption() { .join() .unwrap(); } + +fn run_host_queue_lifecycle_test>( + test: impl FnOnce() -> F + Send + 'static, +) { + std::thread::Builder::new() + .stack_size(16 * 1024 * 1024) + .spawn(move || { + tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .unwrap() + .block_on(test()); + }) + .unwrap() + .join() + .unwrap(); +} + +fn host_queue_lifecycle_model() -> AIConfig { + AIConfig { + models: vec![AIModelConfig { + id: "queue-lifecycle-model".into(), + name: "Queue lifecycle".into(), + provider: "openai".into(), + model_name: "test-model".into(), + base_url: "http://127.0.0.1:1".into(), + enabled: true, + ..Default::default() + }], + ..Default::default() + } +} + +async fn stopped_host_queue_fixture( + retiring: bool, +) -> ( + Arc, + Arc, + tempfile::TempDir, + String, +) { + let (scheduler, sessions, _, root) = test_scheduler_with_persistence(true); + let id = "host-queue-session"; + let workspace = fixture_workspace_dir(root.path().join("stopped-lifecycle")); + sessions + .create_session_with_id( + Some(id.into()), + "Stopped".into(), + "Standard".into(), + SessionConfig { + workspace_path: Some(workspace.to_string_lossy().into_owned()), + ..Default::default() + }, + ) + .await + .unwrap(); + sessions + .start_dialog_turn( + id, + "Standard".into(), + "original".into(), + Some("stopped-turn".into()), + None, + None, + ) + .await + .unwrap(); + let epoch = scheduler + .manage_host_queue(request(None, Action::List)) + .await + .unwrap() + .queue_epoch; + for turn_id in ["old-a", "old-b"] { + scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Submit { + message: message(turn_id), + }, + )) + .await + .unwrap(); + } + sessions + .mark_dialog_turn_interrupted(id, "stopped-turn") + .await + .unwrap(); + sessions + .update_session_state_for_turn_if_processing(id, "stopped-turn", SessionState::Idle) + .await + .unwrap(); + if retiring { + scheduler + .active_turns + .insert(id.into(), desktop_active_turn("stopped-turn")); + } else { + scheduler.hold_managed_queue(id, "Turn interrupted"); + } + TEST_MODEL_RESOLUTION_AI_CONFIG + .scope( + host_queue_lifecycle_model(), + sessions.update_session_model_id(id, "queue-lifecycle-model"), + ) + .await + .unwrap(); + (scheduler, sessions, root, epoch) +} + +async fn process_host_queue_outcome(scheduler: &DialogScheduler, outcome: TurnOutcome) { + let (tx, rx) = mpsc::unbounded_channel(); + tx.send(("host-queue-session".into(), outcome)).unwrap(); + drop(tx); + TEST_MODEL_RESOLUTION_AI_CONFIG + .scope( + host_queue_lifecycle_model(), + scheduler.run_outcome_handler(rx), + ) + .await; +} + +#[test] +fn host_queue_prompt_during_stop_retirement_is_not_parked_by_old_outcome() { + run_host_queue_lifecycle_test(|| async { + let (scheduler, sessions, _root, epoch) = stopped_host_queue_fixture(true).await; + let id = "host-queue-session"; + let mut background = standard_queued_turn("background-result"); + background.policy = DialogSubmissionPolicy::for_source(DialogTriggerSource::AgentSession); + scheduler + .queues + .enqueue(id, background, DialogQueuePriority::Low) + .unwrap(); + let submitted = scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Submit { + message: message("after-stop"), + }, + )) + .await + .unwrap(); + assert_eq!(submitted.receipt.unwrap().status, Status::Queued); + assert!(!sessions + .latest_dialog_turn_holds_dispatch(id) + .await + .unwrap()); + process_host_queue_outcome( + &scheduler, + TurnOutcome::Interrupted { + turn_id: "stopped-turn".into(), + execution_generation: 0, + }, + ) + .await; + assert_eq!( + sessions.get_turn_count(id), + 2, + "new prompt must start without another client request" + ); + let snapshot = scheduler + .manage_host_queue(request(None, Action::List)) + .await + .unwrap(); + assert!(snapshot + .items + .iter() + .all(|item| item.turn_id != "after-stop")); + assert_eq!( + snapshot + .items + .iter() + .filter(|item| item.status == Status::Blocked) + .count(), + 2 + ); + // A second explicit prompt must also run while the older held messages + // remain intact, including when submitted before the current turn ends. + scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Submit { + message: message("another-prompt"), + }, + )) + .await + .unwrap(); + sessions + .update_session_state(id, SessionState::Idle) + .await + .unwrap(); + process_host_queue_outcome( + &scheduler, + TurnOutcome::Completed { + turn_id: "after-stop".into(), + final_response: "done".into(), + }, + ) + .await; + assert_eq!(sessions.get_turn_count(id), 3); + assert_eq!( + scheduler.queues.depth(id), + 1, + "background work must remain parked behind held messages" + ); + let _ = scheduler + .coordinator + .cancel_dialog_turn(id, "another-prompt") + .await; + }); +} + +#[test] +fn host_queue_promote_after_stop_preserves_recovery_on_failure_and_starts_once() { + run_host_queue_lifecycle_test(|| async { + let (scheduler, sessions, _root, epoch) = stopped_host_queue_fixture(false).await; + let id = "host-queue-session"; + let promote = |op: &str| { + request( + Some(&epoch), + Action::Promote { + turn_id: "old-a".into(), + operation_id: op.into(), + expected_active_turn_id: None, + }, + ) + }; + let failed = TEST_MODEL_RESOLUTION_AI_CONFIG + .scope( + AIConfig::default(), + scheduler.execute_queue_request(promote("failed-promotion")), + ) + .await + .unwrap(); + assert_eq!(failed.receipt.unwrap().status, Status::Blocked); + assert!(sessions + .latest_dialog_turn_holds_dispatch(id) + .await + .unwrap()); + assert_eq!(sessions.get_turn_count(id), 1); + let accepted = TEST_MODEL_RESOLUTION_AI_CONFIG + .scope( + host_queue_lifecycle_model(), + scheduler.execute_queue_request(promote("retry-promotion")), + ) + .await + .unwrap(); + assert_eq!(accepted.receipt.unwrap().status, Status::Started); + assert!(!sessions + .latest_dialog_turn_holds_dispatch(id) + .await + .unwrap()); + scheduler + .manage_host_queue(promote("retry-promotion")) + .await + .unwrap(); + assert_eq!(sessions.get_turn_count(id), 2); + assert!(accepted + .items + .iter() + .any(|item| item.turn_id == "old-b" && item.status == Status::Blocked)); + let _ = scheduler.coordinator.cancel_dialog_turn(id, "old-a").await; + }); +} + +#[test] +fn host_queue_cancel_after_stop_does_not_resume_interrupted_or_background_work() { + run_host_queue_lifecycle_test(|| async { + let (scheduler, sessions, _root, epoch) = stopped_host_queue_fixture(false).await; + let id = "host-queue-session"; + let mut background = standard_queued_turn("background-result"); + background.policy = DialogSubmissionPolicy::for_source(DialogTriggerSource::AgentSession); + scheduler + .queues + .enqueue(id, background, DialogQueuePriority::Low) + .unwrap(); + for turn in ["old-a", "old-b"] { + let cancel = request( + Some(&epoch), + Action::Cancel { + turn_id: turn.into(), + operation_id: format!("cancel-{turn}"), + }, + ); + for _ in 0..2 { + let result = scheduler.manage_host_queue(cancel.clone()).await.unwrap(); + assert_eq!(result.receipt.unwrap().status, Status::Cancelled); + } + } + assert!(sessions + .latest_dialog_turn_holds_dispatch(id) + .await + .unwrap()); + assert!(scheduler.try_start_next_queued(id).await.unwrap().is_none()); + assert_eq!(scheduler.queues.depth(id), 1); + assert_eq!(sessions.get_turn_count(id), 1); + }); +} + +#[test] +fn host_queue_new_prompt_after_error_survives_delayed_failure_cleanup() { + run_host_queue_lifecycle_test(|| async { + let (scheduler, sessions, _root, epoch) = stopped_host_queue_fixture(true).await; + let id = "host-queue-session"; + sessions + .abandon_interrupted_dialog_turn(id, Some("stopped-turn")) + .await + .unwrap(); + sessions + .update_session_state( + id, + SessionState::Error { + error: "provider failed".into(), + recoverable: true, + }, + ) + .await + .unwrap(); + let submitted = scheduler + .manage_host_queue(request( + Some(&epoch), + Action::Submit { + message: message("after-error"), + }, + )) + .await + .unwrap(); + assert_eq!(submitted.receipt.unwrap().status, Status::Queued); + process_host_queue_outcome( + &scheduler, + TurnOutcome::Failed { + turn_id: "stopped-turn".into(), + error: "provider failed".into(), + }, + ) + .await; + assert_eq!(sessions.get_turn_count(id), 2); + let snapshot = scheduler + .manage_host_queue(request(None, Action::List)) + .await + .unwrap(); + assert_eq!(snapshot.items.len(), 2); + assert!(snapshot + .items + .iter() + .all(|item| item.status == Status::Blocked && item.turn_id != "after-error")); + let _ = scheduler + .coordinator + .cancel_dialog_turn(id, "after-error") + .await; + }); +} diff --git a/src/crates/assembly/core/src/agentic/coordination/scheduler.rs b/src/crates/assembly/core/src/agentic/coordination/scheduler.rs index 54c77581ad..cf3aaefd0f 100644 --- a/src/crates/assembly/core/src/agentic/coordination/scheduler.rs +++ b/src/crates/assembly/core/src/agentic/coordination/scheduler.rs @@ -109,6 +109,13 @@ pub struct QueuedTurn { } impl QueuedTurn { + fn is_user_submission(&self) -> bool { + !matches!( + self.policy.trigger_source, + DialogTriggerSource::AgentSession | DialogTriggerSource::ScheduledJob + ) + } + fn accept_settlement(&self) { if let Some(registration) = self._settlement_registration.as_ref() { registration.accept(); @@ -1286,6 +1293,21 @@ impl DialogScheduler { .session_manager .get_session(&session_id) .map(|s| s.state.clone()); + if queued_turn.is_user_submission() + && matches!(state, Some(SessionState::Idle | SessionState::Error { .. })) + { + if let Some(previous) = self + .session_manager + .get_session(&session_id) + .and_then(|s| s.dialog_turn_ids.last().cloned()) + .filter(|id| self.active_turns.matches_turn(&session_id, id)) + { + self.host_queue + .lock() + .unwrap_or_else(|e| e.into_inner()) + .mark_after_terminal_turn(&session_id, &resolved_turn_id, previous); + } + } let mut interrupted_hold = matches!(state, Some(SessionState::Idle)) && self .session_manager @@ -1295,11 +1317,13 @@ impl DialogScheduler { // A newly submitted user prompt supersedes recoverable interruption even // when it arrives through the host queue. Existing queued work still // stays parked in try_start_next_queued_locked until that user decision. - let interrupted_turn_to_abandon = if interrupted_hold - && !matches!( - queued_turn.policy.trigger_source, - DialogTriggerSource::AgentSession | DialogTriggerSource::ScheduledJob - ) { + let interrupted_turn_to_abandon = if interrupted_hold && queued_turn.is_user_submission() { + // Park work accepted before this explicit user decision while the + // same session lock still excludes the retiring outcome handler. + self.hold_managed_queue( + &session_id, + "Turn interrupted; retry this message explicitly", + ); interrupted_hold = false; self.session_manager .get_session(&session_id) @@ -1315,7 +1339,7 @@ impl DialogScheduler { > 0; let state_fact = if self.active_turns.contains(&session_id) || interrupted_hold - || (held_user_messages && interrupted_turn_to_abandon.is_none()) + || (held_user_messages && !queued_turn.is_user_submission()) { DialogSessionStateFact::Processing } else { @@ -2057,16 +2081,21 @@ impl DialogScheduler { return Ok(None); } - if self + let held_user_messages = 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 { + > 0; + // Blocked entries are stored outside the runnable queue. They must not + // deadlock newly accepted user work; background work still waits. + let next = if held_user_messages { + self.queues + .remove_first_matching(session_id, QueuedTurn::is_user_submission) + } else { + self.dequeue_next(session_id) + }; + let Some(next_turn) = next else { return Ok(None); }; @@ -2519,9 +2548,10 @@ impl DialogScheduler { .remove_by_id(&session_id, &injection_id); } if lifecycle_plan.status == TurnOutcomeStatus::Interrupted { - self.hold_managed_queue( + self.hold_managed_queue_for_outcome( &session_id, "Turn interrupted; recover it before retrying queued messages", + Some(outcome.turn_id()), ); } if lifecycle_plan.queue_action == TurnOutcomeQueueAction::ClearQueue { @@ -2529,7 +2559,25 @@ impl DialogScheduler { "Turn {}, clearing queue: session_id={}", lifecycle_plan.status, session_id ); + // Snapshot receipt IDs before taking the physical queue + // lock; admission takes these locks in the opposite order. + let newer_ids = self + .host_queue + .lock() + .unwrap_or_else(|e| e.into_inner()) + .admitted_after(&session_id, outcome.turn_id()); + let mut newer_user_work = Vec::new(); + while let Some(turn) = self.queues.remove_first_matching(&session_id, |turn| { + turn.turn_id + .as_ref() + .is_some_and(|id| newer_ids.contains(id)) + }) { + newer_user_work.push(turn); + } let _ = self.clear_queue(&session_id).await; + for turn in newer_user_work.into_iter().rev() { + self.requeue_front(&session_id, turn); + } } (active_turn, active_internal_turn, lifecycle_plan) }; @@ -2750,7 +2798,13 @@ impl DialogScheduler { ), } } - TurnOutcomeQueueAction::ClearQueue => {} + TurnOutcomeQueueAction::ClearQueue => { + // Only user work admitted after the failed turn settled was + // retained above. Previously queued work remains blocked. + if let Err(error) = self.dispatch_next_if_idle(&session_id).await { + warn!("Failed to dispatch newly admitted work after failed turn cleanup: session_id={}, error={}", session_id, error); + } + } } } } diff --git a/src/web-ui/src/flow_chat/components/HostPendingQueuePanel.test.tsx b/src/web-ui/src/flow_chat/components/HostPendingQueuePanel.test.tsx index ecc2c13852..60680a97d3 100644 --- a/src/web-ui/src/flow_chat/components/HostPendingQueuePanel.test.tsx +++ b/src/web-ui/src/flow_chat/components/HostPendingQueuePanel.test.tsx @@ -93,3 +93,36 @@ it('keeps four queued messages compact, exposes actions, and toggles help indepe container.remove(); } }); + +it('copies the complete unresolved draft after an owner restart without resending or deleting it', async () => { + const saved: QueueOutboxRecord = { key: 'draft', scope: 'scope', accepted: true, + request: { sessionId: 'session', queueEpoch: 'old-owner', action: 'submit', + message: { turnId: 'turn', content: 'full prompt', displayContent: 'display prompt', + agentType: 'Standard', attachments: [], metadata: { context: 'original' } } }, + draft: { composerDraft: { text: 'full prompt' }, imageContexts: [{ id: 'attachment' }] } }; + const records = new Map([[saved.key, saved]]); + const invoke = vi.fn(async () => ({ sessionId: 'session', queueEpoch: 'new-owner', revision: 0, + activeTurnId: null, items: [], capacity: 20, used: 0, receipt: null })); + const queue = new HostDialogQueue('scope', 'session', invoke, { + list: async () => [...records.values()], put: async record => { records.set(record.key, record); }, + remove: async key => { records.delete(key); }, + }); + const onRestore = vi.fn(() => true); + const container = document.createElement('div'); + document.body.append(container); + const root = createRoot(container); + try { + await act(async () => root.render()); + await act(async () => container.querySelector('[aria-label="hostQueue.copyDraft"]')!.click()); + expect(onRestore).toHaveBeenCalledWith(expect.objectContaining({ + content: 'full prompt', displayMessage: 'display prompt', composerDraft: { text: 'full prompt' }, + imageContexts: [{ id: 'attachment' }], userMessageMetadata: { context: 'original' }, + })); + expect(invoke.mock.calls).toHaveLength(1); + expect(records.has(saved.key)).toBe(true); + expect(queue.getSnapshot().pending).toHaveLength(1); + } finally { + await act(async () => root.unmount()); + container.remove(); + } +}); diff --git a/src/web-ui/src/flow_chat/components/HostPendingQueuePanel.tsx b/src/web-ui/src/flow_chat/components/HostPendingQueuePanel.tsx index 874eb6f147..5fcbc7475f 100644 --- a/src/web-ui/src/flow_chat/components/HostPendingQueuePanel.tsx +++ b/src/web-ui/src/flow_chat/components/HostPendingQueuePanel.tsx @@ -3,7 +3,7 @@ import type { QueuedMessage } from '../types/flow-chat'; import { useEffect, useId, useState, useSyncExternalStore } from 'react'; import { useI18n } from '@/infrastructure/i18n'; import { Button, IconButton, OverflowText, Tooltip } from '@openbitfun/ui'; -import { ArrowUp, ChevronDown, ChevronUp, Info, Pencil, RotateCw, X } from 'lucide-react'; +import { ArrowUp, ChevronDown, ChevronUp, Copy, Info, Pencil, RotateCw, X } from 'lucide-react'; import './HostPendingQueuePanel.scss'; import { HostDialogQueue, observeHostQueue, type QueueOutboxRecord, type QueueItem } from '../../../../shared/dialog-queue/HostDialogQueue'; import { @@ -89,6 +89,18 @@ export function HostPendingQueuePanel({ queue, onRestore }: { queue: HostDialogQ {record.restoreIntent && record.request.action === 'submit' ? } disabled={busy} onClick={() => void run(() => restoreDraft(record))} /> : } disabled={busy} onClick={() => void run(() => queue.retry(record))} />} + {record.request.action === 'submit' && } disabled={busy} onClick={() => { + if (record.request.action !== 'submit') return; + const message = record.request.message; + const cache = record.draft as Partial | undefined; + // Copy locally even after an owner restart. Keep the unresolved + // receipt: copying is neither a cancellation nor a resend. + onRestore({ ...cache, id: message.turnId, sessionId: queue.sessionId, + content: message.content, displayMessage: message.displayContent, + agentType: message.agentType, timestamp: cache?.timestamp ?? Date.now(), + status: 'queued', retryCount: cache?.retryCount ?? 0, + userMessageMetadata: message.metadata }); + }} />} } disabled={busy} onClick={() => void run(() => queue.dismiss(record))} /> )} 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 415b6c9782..9441d23a44 100644 --- a/src/web-ui/src/locales/en-US/flow-chat.json +++ b/src/web-ui/src/locales/en-US/flow-chat.json @@ -2892,7 +2892,7 @@ "refresh": "Refresh", "unknown": "Delivery is unconfirmed. Check before sending again.", "checkRetry": "Check / retry", - "copyDraft": "Copy text to composer", + "copyDraft": "Copy draft to composer", "dismiss": "Dismiss local reminder", "attachments": "{{count}} attachments", "edit": "Restore draft", 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 49077384fc..6400a70515 100644 --- a/src/web-ui/src/locales/zh-CN/flow-chat.json +++ b/src/web-ui/src/locales/zh-CN/flow-chat.json @@ -2892,7 +2892,7 @@ "refresh": "刷新", "unknown": "尚未确认是否送达,请先检查再重发。", "checkRetry": "检查 / 重试", - "copyDraft": "复制文字到输入框", + "copyDraft": "复制草稿到输入框", "dismiss": "关闭本地提醒", "attachments": "{{count}} 个附件", "edit": "恢复草稿", 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 ffef18db20..a3e9c417ed 100644 --- a/src/web-ui/src/locales/zh-TW/flow-chat.json +++ b/src/web-ui/src/locales/zh-TW/flow-chat.json @@ -2892,7 +2892,7 @@ "refresh": "重新整理", "unknown": "尚未確認是否送達,請先檢查再重新傳送。", "checkRetry": "檢查 / 重試", - "copyDraft": "複製文字到輸入框", + "copyDraft": "複製草稿到輸入框", "dismiss": "關閉本機提醒", "attachments": "{{count}} 個附件", "edit": "恢復草稿",