返回 CodeWhale
dispatch.rs
根目录 / crates / tui / src / tui / ui / dispatch.rs
1 //! Getting a composed user message into a turn: dispatch, steering, and the
2 //! offline/queued message paths.
3 //!
4 //! Moved verbatim out of `ui.rs`.
5
6 use super::*;
7 use crate::core::ops::TurnSpec;
8 use codewhale_models::Role;
9
10 pub(crate) fn dispatch_hotbar_slot(
11 app: &mut App,
12 config: &Config,
13 slot: u8,
14 ) -> Result<Option<HotbarDispatch>> {
15 let known_action_ids = app
16 .hotbar_actions
17 .iter()
18 .map(|action| action.id())
19 .collect::<Vec<_>>();
20 let bindings = config.resolve_hotbar_bindings(&known_action_ids).bindings;
21 let Some(action_id) = bindings
22 .iter()
23 .find(|binding| binding.slot == slot)
24 .map(|binding| binding.action.clone())
25 else {
26 return Ok(None);
27 };
28
29 let Some(action) = app.hotbar_actions.get(&action_id) else {
30 app.status_message = Some(format!(
31 "Hotbar slot {slot} action is not available: {action_id}"
32 ));
33 app.needs_redraw = true;
34 return Ok(Some(HotbarDispatch::Handled));
35 };
36
37 if let Some(reason) = action.disabled_reason(app) {
38 app.status_message = Some(format!(
39 "Hotbar slot {slot} action is not available: {reason}"
40 ));
41 app.needs_redraw = true;
42 return Ok(Some(HotbarDispatch::Handled));
43 }
44
45 action.dispatch(app, config).map(Some)
46 }
47
48 pub(crate) fn queued_ui_to_session(msg: &QueuedMessage) -> QueuedSessionMessage {
49 QueuedSessionMessage {
50 display: msg.display.clone(),
51 skill_instruction: msg.skill_instruction.clone(),
52 skill_provenance: msg.skill_provenance.clone(),
53 }
54 }
55
56 pub(crate) fn queued_session_to_ui(msg: QueuedSessionMessage) -> QueuedMessage {
57 QueuedMessage {
58 display: msg.display,
59 skill_instruction: msg.skill_instruction,
60 skill_provenance: msg.skill_provenance,
61 // Persistence does not carry this flag; restore may re-echo or rely on
62 // pending preview. Live Queue path sets true after painting.
63 history_echoed: false,
64 }
65 }
66
67 /// Echo a freshly submitted turn into the transcript before it waits on the
68 /// queue / offline bucket. Ops contract: every submit paints `HistoryCell::User`
69 /// before the model runs (queued included).
70 pub(crate) fn echo_queued_user_turn(app: &mut App, message: &mut QueuedMessage) {
71 if message.history_echoed {
72 return;
73 }
74 app.add_message(HistoryCell::User {
75 content: message.display.clone(),
76 });
77 message.history_echoed = true;
78 app.needs_redraw = true;
79 app.scroll_to_bottom();
80 }
81
82 /// Paint the transcript cell for a submitted user turn, reusing the cell that
83 /// queue-time echo already painted when there is one. Exactly one
84 /// `HistoryCell::User` must represent a message across queue -> steer ->
85 /// dispatch; returns that cell's index.
86 pub(crate) fn paint_user_turn_cell(
87 app: &mut App,
88 message: &QueuedMessage,
89 content: String,
90 ) -> usize {
91 if message.history_echoed
92 && let Some(idx) = echoed_user_turn_cell(app, &message.display)
93 {
94 app.history[idx] = HistoryCell::User { content };
95 app.bump_history_cell(idx);
96 return idx;
97 }
98 app.add_message(HistoryCell::User { content });
99 app.history.len().saturating_sub(1)
100 }
101
102 /// The newest transcript cell that queue-time echo painted for `display`.
103 pub(crate) fn echoed_user_turn_cell(app: &App, display: &str) -> Option<usize> {
104 app.history
105 .iter()
106 .enumerate()
107 .rev()
108 .find_map(|(idx, cell)| match cell {
109 HistoryCell::User { content } if content == display => Some(idx),
110 _ => None,
111 })
112 }
113
114 pub(crate) fn enqueue_offline_message(app: &mut App, message: QueuedMessage) {
115 app.queue_message(message);
116 persist_offline_queue_state(app);
117 }
118
119 pub(crate) fn push_assistant_message(
120 app: &mut App,
121 text: String,
122 thinking: Option<String>,
123 tool_uses: PendingToolUses,
124 ) {
125 let mut blocks = Vec::new();
126 if let Some(thinking) = thinking {
127 blocks.push(ContentBlock::Thinking {
128 thinking,
129 signature: None,
130 state: None,
131 });
132 }
133 if !text.is_empty() {
134 blocks.push(ContentBlock::Text {
135 text,
136 cache_control: None,
137 });
138 }
139 blocks.extend(tool_uses);
140
141 let has_sendable_content = blocks.iter().any(|block| {
142 matches!(
143 block,
144 ContentBlock::Text { .. } | ContentBlock::ToolUse { .. }
145 )
146 });
147 if has_sendable_content {
148 app.push_api_message(Message {
149 role: Role::Assistant,
150 content: blocks,
151 });
152 }
153 }
154
155 pub(crate) fn replace_matching_assistant_text(
156 app: &mut App,
157 original_text: &str,
158 translated_text: String,
159 ) -> bool {
160 for message in app.api_messages_mut().iter_mut().rev() {
161 if message.role != "assistant"
162 && message.role != codewhale_models::INTERRUPTED_ASSISTANT_ROLE
163 {
164 continue;
165 }
166 for block in &mut message.content {
167 if let ContentBlock::Text { text, .. } = block
168 && text == original_text
169 {
170 *text = translated_text;
171 return true;
172 }
173 }
174 }
175 false
176 }
177
178 pub(crate) fn build_queued_message(app: &mut App, input: String) -> QueuedMessage {
179 let skill_instruction = app.active_skill.take();
180 let skill_provenance = app.active_skill_provenance.take();
181 QueuedMessage::new(input, skill_instruction).with_skill_provenance(skill_provenance)
182 }
183
184 pub(crate) fn allowed_tools_for_message(
185 configured: Option<Vec<String>>,
186 message: &QueuedMessage,
187 ) -> Option<Vec<String>> {
188 if message.is_workflow_draft() {
189 // `/workflow <objective>` is review-first. The model may draft and ask
190 // for confirmation, but the host makes execution impossible in the
191 // same turn even if the provider ignores that instruction.
192 Some(Vec::new())
193 } else {
194 configured
195 }
196 }
197
198 pub(crate) async fn submit_initial_input_if_ready(
199 app: &mut App,
200 config: &Config,
201 engine_handle: &EngineHandle,
202 ) -> Result<()> {
203 if !app.auto_submit_initial_input {
204 return Ok(());
205 }
206
207 if app.onboarding != OnboardingState::None || app.redaction_gate {
208 if app.status_message.is_none() && !app.input.trim().is_empty() {
209 app.status_message = Some(INITIAL_PROMPT_DEFERRED_STATUS.to_string());
210 }
211 return Ok(());
212 }
213
214 app.auto_submit_initial_input = false;
215 if let Some(input) = app.submit_input() {
216 if app.status_message.as_deref() == Some(INITIAL_PROMPT_DEFERRED_STATUS) {
217 app.status_message = None;
218 }
219 let queued = build_queued_message(app, input);
220 dispatch_user_message_with_recovery(
221 app,
222 config,
223 engine_handle,
224 queued,
225 DispatchRecovery::Initial,
226 )
227 .await?;
228 }
229 Ok(())
230 }
231
232 pub(crate) fn message_from_submitted_input(
233 app: &mut App,
234 input: String,
235 ) -> (QueuedMessage, DispatchRecovery) {
236 if let Some(mut draft) = app.queued_draft.take() {
237 draft.display = input;
238 (draft, DispatchRecovery::Draft)
239 } else {
240 (
241 build_queued_message(app, input),
242 DispatchRecovery::Immediate,
243 )
244 }
245 }
246
247 pub(crate) fn take_next_queued_message(app: &mut App) -> Option<(QueuedMessage, DispatchRecovery)> {
248 if app.input.is_empty() {
249 return app.remove_queued_message(0).map(|message| {
250 (
251 message,
252 DispatchRecovery::Queued {
253 restore_index: Some(0),
254 },
255 )
256 });
257 }
258 None
259 }
260
261 pub(crate) async fn send_next_queued_message_now(
262 app: &mut App,
263 config: &Config,
264 engine_handle: &EngineHandle,
265 ) -> Result<bool> {
266 let Some((message, recovery)) = take_next_queued_message(app) else {
267 return Ok(false);
268 };
269 send_taken_queued_message_now(app, config, engine_handle, message, recovery).await?;
270 Ok(true)
271 }
272
273 pub(crate) async fn send_queued_message_at_index_now(
274 app: &mut App,
275 config: &Config,
276 engine_handle: &EngineHandle,
277 index: usize,
278 ) -> Result<bool> {
279 let Some(message) = app.remove_queued_message(index) else {
280 app.status_message = Some("Queued message not found".to_string());
281 return Ok(true);
282 };
283 send_taken_queued_message_now(
284 app,
285 config,
286 engine_handle,
287 message,
288 DispatchRecovery::Queued {
289 restore_index: Some(index),
290 },
291 )
292 .await?;
293 Ok(true)
294 }
295
296 pub(crate) async fn send_taken_queued_message_now(
297 app: &mut App,
298 config: &Config,
299 engine_handle: &EngineHandle,
300 message: QueuedMessage,
301 recovery: DispatchRecovery,
302 ) -> Result<()> {
303 if app.offline_mode {
304 restore_queued_or_draft_message(app, recovery, message);
305 app.status_message = Some(
306 app.tr(MessageId::ToastOfflineQueuedCount)
307 .replace("{count}", &app.queued_message_count().to_string()),
308 );
309 return Ok(());
310 }
311
312 if app.dispatch_in_flight {
313 // A spawned dispatch is still resolving route/sending its op (#4605):
314 // there is no turn to steer into yet. Re-queue; the completion/turn
315 // lifecycle will drive the next drain.
316 restore_queued_or_draft_message(app, recovery, message);
317 app.status_message = Some(queued_follow_up_toast(app));
318 return Ok(());
319 }
320 if app.is_loading {
321 match steer_user_message(app, config, engine_handle, message.clone()).await {
322 Ok(true) => app.push_status_toast(
323 app.tr(MessageId::ToastSentIntoTurn).into_owned(),
324 StatusToastLevel::Info,
325 Some(1_500),
326 ),
327 Ok(false) => {
328 restore_queued_or_draft_message(app, recovery, message);
329 app.push_status_toast(
330 app.tr(MessageId::ToastHookBlockedFollowUp).into_owned(),
331 StatusToastLevel::Warning,
332 Some(4_000),
333 );
334 }
335 Err(err) => {
336 restore_queued_or_draft_message(app, recovery, message);
337 app.status_message = Some(format!(
338 "{} ({err})",
339 app.tr(MessageId::ToastCouldNotSendIntoTurn)
340 ));
341 }
342 }
343 } else if let Err(_err) =
344 dispatch_user_message_with_recovery(app, config, engine_handle, message, recovery).await
345 {
346 // The completion closure re-queued the message and set the status.
347 } else {
348 app.status_message = Some(app.tr(MessageId::ToastSentIntoTurn).into_owned());
349 }
350 Ok(())
351 }
352
353 pub(crate) fn queued_message_content_for_app(
354 app: &App,
355 message: &QueuedMessage,
356 cwd: Option<PathBuf>,
357 git_cache: &mut crate::tui::git_mention::GitMentionCache,
358 ) -> Result<String> {
359 if let Some(provenance) = message.skill_provenance.as_ref() {
360 provenance
361 .verify_for(&app.workspace, Some(app.extension_plugin_view().as_ref()))
362 .map_err(anyhow::Error::msg)?;
363 }
364 // Pass the process CWD explicitly so the resolver's two-pass logic can
365 // honor the user's launch directory when it differs from `--workspace`
366 // (issue #101 — file mentions silently routing to the wrong root).
367 // The completion index is the composer's already-built fuzzy scan: a
368 // bounded fallback for exact misses, with no submit-time tree walk (#4365).
369 let completion_index = app.composer.mention_discovery.fuzzy_candidates(
370 &app.workspace,
371 &app.composer.mention_cwd,
372 app.mention_walk_depth,
373 app.workspace_follow_symlinks,
374 );
375 // Stabilize macOS screencapture temp references before anything else sees
376 // the text: macOS deletes those Temporary Items dirs minutes after capture.
377 let stabilization_dir = crate::tui::file_mention::screenshot_stabilization_dir(&app.workspace);
378 let display = crate::tui::file_mention::stabilize_screenshot_references(
379 &message.display,
380 &stabilization_dir,
381 );
382 let user_request = crate::tui::file_mention::user_request_with_file_mentions_cached(
383 &display,
384 &app.workspace,
385 cwd,
386 git_cache,
387 completion_index,
388 );
389 if let Some(skill_instruction) = message.skill_instruction.as_ref() {
390 Ok(format!(
391 "{skill_instruction}\n\n---\n\nUser request: {user_request}"
392 ))
393 } else {
394 Ok(user_request)
395 }
396 }
397
398 pub(crate) fn dispatch_completion_permit(
399 app: &App,
400 ) -> std::result::Result<
401 tokio::sync::mpsc::OwnedPermit<crate::tui::app::DispatchApplyFn>,
402 &'static str,
403 > {
404 let sender = app
405 .dispatch_completion_tx
406 .clone()
407 .ok_or("dispatch completion mailbox is unavailable")?;
408 sender.try_reserve_owned().map_err(|error| match error {
409 tokio::sync::mpsc::error::TrySendError::Full(_) => "dispatch completion mailbox is full",
410 tokio::sync::mpsc::error::TrySendError::Closed(_) => {
411 "dispatch completion mailbox is closed"
412 }
413 })
414 }
415
416 #[cfg(test)]
417 pub(crate) async fn dispatch_user_message(
418 app: &mut App,
419 config: &Config,
420 engine_handle: &EngineHandle,
421 message: QueuedMessage,
422 ) -> Result<()> {
423 dispatch_user_message_with_recovery(
424 app,
425 config,
426 engine_handle,
427 message,
428 DispatchRecovery::Immediate,
429 )
430 .await
431 }
432
433 pub(crate) async fn dispatch_user_message_with_recovery(
434 app: &mut App,
435 config: &Config,
436 engine_handle: &EngineHandle,
437 mut message: QueuedMessage,
438 recovery: DispatchRecovery,
439 ) -> Result<()> {
440 if app.redaction_gate {
441 recover_unstarted_external_message(app, message, recovery, INITIAL_PROMPT_DEFERRED_STATUS);
442 return Ok(());
443 }
444 let stop_words = config.stop_words();
445 if is_stop_word(&message.display, &stop_words).is_some() {
446 engine_handle.cancel();
447 app.stopped_turn = true;
448 app.status_message = Some("Turn stopped. Tool calls blocked for this turn.".to_string());
449 return Ok(());
450 }
451 app.stopped_turn = false;
452
453 // #1364: run mutable `message_submit` hooks before dispatch. Hooks see the
454 // user's display text and may replace or block it before file mentions,
455 // skill wrapping, history, and model input are resolved.
456 // Fast-path skip when no hooks configured.
457 if app
458 .hooks
459 .has_hooks_for_event(crate::hooks::HookEvent::MessageSubmit)
460 {
461 let context = app.base_hook_context().with_message(&message.display);
462 let strict_gates = app
463 .hooks
464 .matched_strict_gate_labels(crate::hooks::HookEvent::MessageSubmit, &context);
465 let hooks = app.hooks.clone();
466 let original_text = message.display.clone();
467
468 if app.dispatch_completion_tx.is_some() {
469 // The foreground transform is a gate, but its child wait belongs
470 // on the blocking pool, never on the terminal event loop. Result
471 // delivery reserves bounded mailbox capacity before any work or
472 // state mutation, so the recovery closure cannot be dropped.
473 let completion_permit = match dispatch_completion_permit(app) {
474 Ok(permit) => permit,
475 Err(error) => {
476 recover_unstarted_external_message(app, message, recovery, error);
477 return Err(anyhow::Error::msg(error));
478 }
479 };
480 app.dispatch_in_flight = true;
481 tokio::spawn(async move {
482 let outcome = match tokio::task::spawn_blocking(move || {
483 hooks.execute_message_submit_transform_for_dispatch(&context, &original_text)
484 })
485 .await
486 {
487 Ok(outcome) => outcome,
488 Err(error) => {
489 tracing::error!(target: "hooks", %error, "message_submit executor task was lost");
490 lost_message_submit_outcome(&strict_gates)
491 }
492 };
493 let apply: crate::tui::app::DispatchApplyFn = Box::new(
494 move |app: &mut App,
495 engine_handle: &EngineHandle,
496 config: &Config|
497 -> anyhow::Result<()> {
498 if !apply_message_submit_outcome(app, &mut message, outcome) {
499 app.dispatch_in_flight = false;
500 restore_message_submit_denial(app, message, recovery);
501 return Ok(());
502 }
503 let _ = start_user_dispatch(app, config, engine_handle, message, recovery);
504 Ok(())
505 },
506 );
507 completion_permit.send(apply);
508 });
509 return Ok(());
510 }
511
512 // Unit tests intentionally omit the event-loop completion channel.
513 // Keep those synchronous from the test's perspective while still
514 // running the blocking child wait off the async runtime worker.
515 let outcome = match tokio::task::spawn_blocking(move || {
516 hooks.execute_message_submit_transform_for_dispatch(&context, &original_text)
517 })
518 .await
519 {
520 Ok(outcome) => outcome,
521 Err(error) => {
522 tracing::error!(target: "hooks", %error, "message_submit executor task was lost");
523 lost_message_submit_outcome(&strict_gates)
524 }
525 };
526 if !apply_message_submit_outcome(app, &mut message, outcome) {
527 restore_message_submit_denial(app, message, recovery);
528 return Ok(());
529 }
530 }
531
532 if app.dispatch_completion_tx.is_some() {
533 return start_user_dispatch(app, config, engine_handle, message, recovery);
534 }
535
536 let prepare = match prepare_user_dispatch(app, config, message.clone()) {
537 Ok(prepare) => prepare,
538 Err(error) => {
539 recover_unstarted_external_message(app, message, recovery, &error.to_string());
540 return Err(error);
541 }
542 };
543 run_prepared_dispatch(app, config, engine_handle, prepare, recovery).await
544 }
545
546 pub(crate) fn lost_message_submit_outcome(
547 strict_gates: &[String],
548 ) -> crate::hooks::MessageSubmitOutcome {
549 if strict_gates.is_empty() {
550 crate::hooks::MessageSubmitOutcome::Unchanged {
551 warning: Some(
552 "message_submit hook executor did not run; submission continued because no strict gate matched"
553 .to_string(),
554 ),
555 }
556 } else {
557 crate::hooks::MessageSubmitOutcome::Blocked {
558 reason: "message_submit hook executor did not run; a strict gate blocked submission"
559 .to_string(),
560 }
561 }
562 }
563
564 pub(crate) fn prepare_user_dispatch(
565 app: &mut App,
566 config: &Config,
567 message: QueuedMessage,
568 ) -> Result<UserDispatchPrepare> {
569 anyhow::ensure!(!app.redaction_gate, "{INITIAL_PROMPT_DEFERRED_STATUS}");
570 let _ = app.maybe_nudge_plugin_for_prompt(&message.display);
571
572 // Plan paused-command changes without touching App or the engine pause
573 // gate. Route selection can await and client preflight can fail; neither
574 // may resume or discard a paused command unless a turn is ready to send.
575 let paused_dispatch = plan_paused_command_message(app, &message.display);
576
577 let cwd = std::env::current_dir().ok();
578 // One cache for this submit: the references pass and the payload pass
579 // otherwise each shell out for `@git`/`@diff`, making git compute a large
580 // working-tree diff twice to attach it once (#4067 review follow-up).
581 let mut git_cache = crate::tui::git_mention::GitMentionCache::default();
582 let completion_index = app.composer.mention_discovery.fuzzy_candidates(
583 &app.workspace,
584 &app.composer.mention_cwd,
585 app.mention_walk_depth,
586 app.workspace_follow_symlinks,
587 );
588 let references = crate::tui::file_mention::context_references_from_input_cached(
589 &message.display,
590 &app.workspace,
591 cwd.clone(),
592 &mut git_cache,
593 completion_index,
594 );
595 let mut content = queued_message_content_for_app(app, &message, cwd, &mut git_cache)?;
596 if let Some(note) = paused_dispatch.note() {
597 content.push_str(note);
598 }
599 let (app_route_identity, route_config) =
600 app_scoped_runtime_config(app, config).map_err(anyhow::Error::msg)?;
601
602 let should_auto_resolve = auto_router::should_resolve_auto_model_selection(app);
603 let auto_router_context = auto_router::recent_auto_router_context(&app.api_messages);
604
605 // Capture the App state before any optimistic mutation so a failure can
606 // roll back cleanly.
607 let snapshot = UserDispatchSnapshot {
608 is_loading: app.is_loading,
609 suppress_stream_events_until_turn_complete: app.suppress_stream_events_until_turn_complete,
610 runtime_turn_status: app.runtime_turn_status.clone(),
611 receipt_text: app.receipt_text.clone(),
612 receipt_started_at: app.receipt_started_at,
613 tool_evidence: app.tool_evidence.clone(),
614 history_len: app.history.len(),
615 history_revisions_len: app.history_revisions.len(),
616 api_messages_len: app.api_messages.len(),
617 last_send_at: app.last_send_at,
618 };
619
620 // --- Sync prepare: show the user message and spinner immediately so the
621 // event loop can repaint before network I/O (#4605). The async phase runs
622 // the auto-model route, compaction, and engine send off the render thread.
623 app.is_loading = true;
624 app.runtime_turn_status = None;
625 app.clear_receipt();
626 app.tool_evidence.clear();
627 app.needs_redraw = true;
628
629 let message_index = app.api_messages.len();
630 // Already painted at Queue time — reuse that cell for reference
631 // recording instead of duplicating the bubble.
632 let history_cell = paint_user_turn_cell(app, &message, message.display.clone());
633 app.scroll_to_bottom();
634 // Anchor the tail-flash to the moment the user message appears, not to
635 // the async dispatch completion (which can lag by a route plan). The
636 // failure path restores the pre-send timestamp from the snapshot.
637 app.last_send_at = Some(Instant::now());
638 app.push_api_message(Message {
639 role: Role::User,
640 content: vec![ContentBlock::Text {
641 text: content.clone(),
642 cache_control: None,
643 }],
644 });
645
646 let goal_objective = paused_dispatch.goal_objective(app);
647 let allowed_tools = allowed_tools_for_message(app.active_allowed_tools.clone(), &message);
648
649 Ok(UserDispatchPrepare {
650 message,
651 content,
652 references,
653 paused_dispatch,
654 app_route_identity,
655 route_config,
656 goal_objective,
657 goal_status: app.goal.status,
658 goal_token_budget: app.goal.token_budget,
659 mode: app.mode,
660 api_provider: app.api_provider,
661 app_model: app.model.clone(),
662 auto_model: app.auto_model,
663 reasoning_effort: app.reasoning_effort,
664 allow_shell: app.allow_shell,
665 trust_mode: app.trust_mode,
666 auto_approve: app_auto_approve_enabled(app),
667 approval_mode: app.approval_mode,
668 translation_enabled: app.translation_enabled,
669 allowed_tools,
670 hook_executor: app.runtime_services.hook_executor.clone(),
671 verbosity: app.verbosity.clone(),
672 provenance: UserInputProvenance::ExternalUser,
673 auto_router_context,
674 should_auto_resolve,
675 auto_compact_user_configured: app.auto_compact_user_configured,
676 auto_compact: app.auto_compact,
677 auto_compact_threshold_percent: app.auto_compact_threshold_percent,
678 snapshot,
679 cost_scope: crate::cost_status::scope_token(),
680 message_index,
681 history_cell,
682 })
683 }
684
685 pub(crate) fn start_user_dispatch(
686 app: &mut App,
687 config: &Config,
688 engine_handle: &EngineHandle,
689 message: QueuedMessage,
690 recovery: DispatchRecovery,
691 ) -> Result<()> {
692 let completion_permit = match dispatch_completion_permit(app) {
693 Ok(permit) => permit,
694 Err(error) => {
695 recover_unstarted_external_message(app, message, recovery, error);
696 return Err(anyhow::Error::msg(error));
697 }
698 };
699 let recovery_message = message.clone();
700 let prepare = match prepare_user_dispatch(app, config, message) {
701 Ok(prepare) => prepare,
702 Err(error) => {
703 recover_unstarted_external_message(app, recovery_message, recovery, &error.to_string());
704 return Err(error);
705 }
706 };
707 app.dispatch_in_flight = true;
708 // #6800: a local cancel (Esc, stall recovery) trips this token so the
709 // dispatch fails back to the composer at once instead of holding
710 // `dispatch_in_flight` — and queueing every new send — for the full
711 // `DISPATCH_TASK_BOUND` while it waits on engine admission.
712 let cancel = tokio_util::sync::CancellationToken::new();
713 app.dispatch_cancel = Some(cancel.clone());
714 // Supervised: `spawned_dispatch_execute` owns the whole dispatch future,
715 // so its completion callback always arrives — on success, on a panic, on
716 // a local cancel, or when the dispatch exceeds its bound (#6184).
717 tokio::spawn(spawned_dispatch_execute(
718 prepare,
719 recovery,
720 engine_handle.clone(),
721 completion_permit,
722 cancel,
723 ));
724 Ok(())
725 }
726
727 /// Longest a dispatch may spend routing and waiting for engine admission
728 /// before it is failed back to the composer (#6184). An engine whose op
729 /// mailbox never frees (a wedged turn) used to hold the dispatch — and the
730 /// user's message — forever.
731 pub(crate) const DISPATCH_TASK_BOUND: std::time::Duration = std::time::Duration::from_secs(60);
732
733 pub(crate) async fn spawned_dispatch_execute(
734 prepare: UserDispatchPrepare,
735 recovery: DispatchRecovery,
736 engine_handle: EngineHandle,
737 completion_permit: tokio::sync::mpsc::OwnedPermit<crate::tui::app::DispatchApplyFn>,
738 cancel: tokio_util::sync::CancellationToken,
739 ) {
740 let apply = supervised_dispatch(
741 prepare,
742 recovery,
743 DISPATCH_TASK_BOUND,
744 cancel,
745 |prepare, recovery| spawned_dispatch_inner(prepare, recovery, engine_handle),
746 )
747 .await;
748 completion_permit.send(apply);
749 }
750
751 /// Run one dispatch future under supervision. A dropped JoinHandle used to
752 /// turn a panic or a hang into a dispatch that never reported back: the
753 /// completion permit was dropped, `dispatch_in_flight` stayed set and the
754 /// message sat in limbo. Every outcome now yields a callback; a panic or an
755 /// overrun also leaves a log line and a `crashes/` record. A tripped `cancel`
756 /// token abandons the dispatch and fails it back through the same error
757 /// closure, so the unsent message returns to the composer exactly as on any
758 /// other dispatch failure (#6800).
759 pub(crate) async fn supervised_dispatch<F, Fut>(
760 prepare: UserDispatchPrepare,
761 recovery: DispatchRecovery,
762 bound: std::time::Duration,
763 cancel: tokio_util::sync::CancellationToken,
764 run: F,
765 ) -> crate::tui::app::DispatchApplyFn
766 where
767 F: FnOnce(UserDispatchPrepare, DispatchRecovery) -> Fut,
768 Fut: std::future::Future<Output = crate::tui::app::DispatchApplyFn>,
769 {
770 use futures_util::FutureExt as _;
771 let fallback = prepare.clone();
772 let started = std::time::Instant::now();
773 let supervised = std::panic::AssertUnwindSafe(run(prepare, recovery)).catch_unwind();
774 let outcome = tokio::select! {
775 biased;
776 () = cancel.cancelled() => {
777 return build_dispatch_error_closure(
778 fallback,
779 recovery,
780 "Message dispatch was cancelled before it reached the engine".to_string(),
781 );
782 }
783 outcome = tokio::time::timeout(bound, supervised) => outcome,
784 };
785 match outcome {
786 Ok(Ok(apply)) => apply,
787 Ok(Err(panic)) => {
788 let detail = crate::utils::panic_message(&*panic);
789 crate::utils::record_caught_panic("user-dispatch", &detail);
790 build_dispatch_error_closure(
791 fallback,
792 recovery,
793 format!("Message dispatch hit an internal error: {detail}"),
794 )
795 }
796 Err(_elapsed) => {
797 crate::core::engine::turn_heartbeat::report_stall(
798 &crate::core::engine::turn_heartbeat::StallReport {
799 source: "ui",
800 phase: "while dispatching the message (route planning / engine admission)"
801 .to_string(),
802 detail: Some(format!(
803 "{} / {}",
804 fallback
805 .app_route_identity
806 .compatibility()
807 .map_or(fallback.app_route_identity.key.as_str(), |row| row.label),
808 fallback.app_model
809 )),
810 turn_id: None,
811 provider_request: None,
812 since_progress: started.elapsed(),
813 bound: Some(bound),
814 },
815 );
816 build_dispatch_error_closure(
817 fallback,
818 recovery,
819 format!(
820 "Message dispatch stalled for {}s before the engine accepted it; your message was restored. Press Esc to cancel the running turn, then retry.",
821 bound.as_secs()
822 ),
823 )
824 }
825 }
826 }
827
828 /// Keep classifier receipts owned until the UI admits the operation to Engine.
829 /// Dropping a reserved dispatch (including a closed completion mailbox) must
830 /// settle its already-incurred usage in the original session scope.
831 pub(super) struct UnacceptedDispatchUsage {
832 pub(super) scope: crate::cost_status::CostScopeToken,
833 pub(super) batch: Option<crate::cost_status::RuntimeUsageBatch>,
834 }
835
836 impl Drop for UnacceptedDispatchUsage {
837 fn drop(&mut self) {
838 if let Some(batch) = self.batch.as_ref() {
839 crate::cost_status::report_runtime_usage_batch(self.scope, None, batch);
840 }
841 }
842 }
843
844 pub(crate) async fn spawned_dispatch_inner(
845 prepare: UserDispatchPrepare,
846 recovery: DispatchRecovery,
847 engine_handle: EngineHandle,
848 ) -> crate::tui::app::DispatchApplyFn {
849 // Bound in its own statement: the planner borrows `prepare`, and the error
850 // arm moves it into the failure closure.
851 let plan_result = plan_turn_route(TurnRoutePlanRequest {
852 route_config: &prepare.route_config,
853 app_route_identity: &prepare.app_route_identity,
854 api_provider: prepare.api_provider,
855 app_model: &prepare.app_model,
856 auto_model: prepare.auto_model,
857 reasoning_effort: prepare.reasoning_effort,
858 mode: prepare.mode,
859 content: &prepare.content,
860 auto_router_context: &prepare.auto_router_context,
861 should_auto_resolve: prepare.should_auto_resolve,
862 allow_auto_router_response_cache: true,
863 preflight_required: engine_handle.client_preflight_required(),
864 auto_compact_user_configured: prepare.auto_compact_user_configured,
865 auto_compact: prepare.auto_compact,
866 auto_compact_threshold_percent: prepare.auto_compact_threshold_percent,
867 })
868 .await;
869 let planned = match plan_result {
870 Ok(planned) => planned,
871 Err(err) => return build_dispatch_error_closure(prepare, recovery, err),
872 };
873
874 let PlannedTurnRoute {
875 route: turn_route,
876 compaction: turn_compaction,
877 effective_provider,
878 effective_model,
879 effective_provider_identity,
880 effective_provider_label,
881 selected_reasoning_effort,
882 effective_reasoning_effort,
883 auto_controls_reasoning,
884 auto_selection,
885 initial_routed_usage,
886 routing_source: _,
887 } = planned;
888 let effective_reasoning_tier = selected_reasoning_effort
889 .unwrap_or(prepare.reasoning_effort)
890 .normalize_for_route(
891 effective_provider,
892 &turn_route.candidate.endpoint().base_url,
893 &turn_route.model,
894 );
895 let effective_reasoning_receipt = reasoning_effort_receipt_for_route(
896 effective_reasoning_tier,
897 effective_provider,
898 &turn_route.candidate.endpoint().base_url,
899 &turn_route.model,
900 );
901
902 let mut usage = UnacceptedDispatchUsage {
903 scope: prepare.cost_scope,
904 batch: Some(initial_routed_usage.clone()),
905 };
906 let op = Op::SendMessage(TurnSpec {
907 max_output_tokens: None,
908 content: prepare.content.clone(),
909 images: Vec::new(),
910 mode: prepare.mode,
911 route: Box::new(turn_route),
912 compaction: Box::new(turn_compaction.clone()),
913 initial_routed_usage: Box::new(initial_routed_usage),
914 goal_objective: prepare.goal_objective.clone(),
915 goal_token_budget: prepare.goal_token_budget,
916 goal_status: prepare.goal_status,
917 reasoning_effort: effective_reasoning_effort,
918 reasoning_effort_auto: auto_controls_reasoning,
919 auto_model: prepare.auto_model,
920 allow_shell: prepare.allow_shell,
921 trust_mode: prepare.trust_mode,
922 auto_approve: prepare.auto_approve,
923 approval_mode: prepare.approval_mode,
924 translation_enabled: prepare.translation_enabled,
925 allowed_tools: prepare.allowed_tools.clone(),
926 dynamic_tools: Vec::new(),
927 hook_executor: prepare.hook_executor.clone(),
928 verbosity: prepare.verbosity.clone(),
929 provenance: prepare.provenance,
930 // Interactive TUI submissions do not correlate submissions.
931 submission_id: None,
932 });
933 // Reserve capacity off the render thread, but do not let Engine start
934 // until the completion callback has installed the UI's acceptance state.
935 // Separate completion/event mailboxes otherwise allow TurnStarted (or
936 // TurnComplete) to arrive before a callback that resets those newer facts.
937 let permit = match engine_handle.tx_op.clone().reserve_owned().await {
938 Ok(permit) => permit,
939 Err(err) => return build_dispatch_error_closure(prepare, recovery, err.to_string()),
940 };
941 let outcome = UserDispatchOutcome {
942 turn_compaction,
943 effective_provider,
944 effective_model,
945 effective_provider_identity,
946 effective_provider_label,
947 effective_reasoning_effort: effective_reasoning_receipt,
948 auto_selection,
949 };
950 Box::new(move |app, current_engine, config| {
951 // Admission stays serialized by this flag until its callback retires,
952 // even if the user replaced the Engine/session while routing waited.
953 app.dispatch_in_flight = false;
954 // This request has no admitted Op and cannot emit TurnComplete. Retire
955 // its local cancellation even after replacement, but leave a previous
956 // admitted turn's suppression for that turn's terminal event to retire.
957 if !prepare.snapshot.suppress_stream_events_until_turn_complete {
958 app.suppress_stream_events_until_turn_complete = false;
959 }
960 if !engine_handle.tx_op.same_channel(&current_engine.tx_op)
961 || prepare.cost_scope != crate::cost_status::scope_token()
962 {
963 anyhow::bail!("Message dispatch belongs to a previous engine or session");
964 }
965 if !app.is_loading || engine_handle.tx_op.is_closed() {
966 let error = if engine_handle.tx_op.is_closed() {
967 "Engine stopped before accepting the message"
968 } else {
969 "Message dispatch was cancelled before it reached the engine"
970 };
971 return build_dispatch_error_closure(prepare, recovery, error.to_string())(
972 app,
973 &engine_handle,
974 config,
975 );
976 }
977 build_dispatch_success_closure(prepare, outcome)(app, &engine_handle, config)?;
978 // Existing Engine admission binds cancellation controls and the Op in
979 // one FIFO. No await separates the UI checkpoint from this handoff.
980 engine_handle.send_reserved_op(permit, op);
981 drop(usage.batch.take());
982 Ok(())
983 })
984 }
985
986 pub(crate) fn build_dispatch_success_closure(
987 prepare: UserDispatchPrepare,
988 outcome: UserDispatchOutcome,
989 ) -> crate::tui::app::DispatchApplyFn {
990 Box::new(
991 move |app: &mut App, engine_handle: &EngineHandle, config: &Config| -> anyhow::Result<()> {
992 app.dispatch_in_flight = false;
993 prepare.paused_dispatch.apply(app, engine_handle);
994
995 let dispatch_started_at = Instant::now();
996 app.is_loading = true;
997 app.dispatch_started_at = Some(dispatch_started_at);
998 app.runtime_turn_status = None;
999 // last_send_at was already anchored in the sync prepare phase so
1000 // the tail-flash starts together with the visible user cell.
1001 app.last_submitted_prompt = Some(prepare.message.display.clone());
1002 app.unanswered_submission = Some(crate::tui::app::UnansweredSubmission {
1003 message: prepare.message.clone(),
1004 history_cell: prepare.history_cell,
1005 });
1006 app.clear_receipt();
1007 app.tool_evidence.clear();
1008
1009 app.system_prompt = Some(build_app_system_prompt_with_goal(
1010 app,
1011 config,
1012 app.goal.objective.as_deref(),
1013 ));
1014 // History and api_messages were already appended in the sync prepare
1015 // phase; record references now that the turn is accepted.
1016 app.record_context_references(
1017 prepare.history_cell,
1018 prepare.message_index,
1019 prepare.references,
1020 );
1021 app.scroll_to_bottom();
1022
1023 app.last_effective_reasoning_effort = Some(outcome.effective_reasoning_effort);
1024 if prepare.auto_model {
1025 app.last_effective_model = Some(outcome.effective_model.clone());
1026 app.last_effective_provider = Some(outcome.effective_provider);
1027 app.last_effective_provider_identity =
1028 Some(outcome.effective_provider_identity.clone());
1029 if let Some(selection) = outcome.auto_selection.as_ref() {
1030 app.last_auto_route_receipt = selection.receipt.clone();
1031 let status = app
1032 .tr(MessageId::AutoRouteSelectedToast)
1033 .replace("{provider}", &outcome.effective_provider_label)
1034 .replace("{model}", &outcome.effective_model)
1035 .replace("{source}", selection.source.label());
1036 app.push_status_toast(status, StatusToastLevel::Info, Some(6_000));
1037 }
1038 } else {
1039 app.last_effective_model = None;
1040 app.last_effective_provider = None;
1041 app.last_effective_provider_identity = None;
1042 app.last_auto_route_receipt = None;
1043 }
1044 app.pending_auto_route_receipt = outcome
1045 .auto_selection
1046 .as_ref()
1047 .and_then(|selection| selection.receipt.clone());
1048 app.pending_turn_route = Some((
1049 outcome.effective_provider,
1050 outcome.effective_model,
1051 prepare.auto_model,
1052 ));
1053
1054 maybe_warn_context_pressure_for_config(app, &outcome.turn_compaction);
1055 if let Some(message) = crate::plugins::plugin_reload_nudge(
1056 app.plugin_registry.as_ref(),
1057 &mut app.plugin_reload_nudge_stamp,
1058 ) {
1059 app.push_status_toast(message, StatusToastLevel::Warning, Some(8_000));
1060 }
1061 app.session.last_prompt_tokens = None;
1062 app.session.last_completion_tokens = None;
1063 app.session.last_prompt_cache_hit_tokens = None;
1064 app.session.last_prompt_cache_miss_tokens = None;
1065 app.session.last_reasoning_replay_tokens = None;
1066
1067 if let Ok(manager) = SessionManager::default_location()
1068 && let Ok(session) = build_session_snapshot(app, &manager)
1069 {
1070 if app.current_session_id.is_none() {
1071 app.current_session_id = Some(session.metadata.id.clone());
1072 }
1073 if let Err(err) = persist_with_pending_work_boundary(
1074 app,
1075 PersistRequest::SaveCheckpoint { session },
1076 ) {
1077 app.status_message = Some(format!(
1078 "To-do list update pending: turn checkpoint could not be queued ({err})"
1079 ));
1080 }
1081 }
1082
1083 Ok(())
1084 },
1085 )
1086 }
1087
1088 /// Missing-credential / auth preflight failures. The message never reached a
1089 /// model, so it returns to the composer with a transcript line saying why
1090 /// (`keep_unsent_message_for_connect`) rather than as an echo that has to be
1091 /// typed again — which lost the text and then showed it twice (#6566).
1092 pub(crate) fn is_missing_credential_dispatch_error(error: &str) -> bool {
1093 let lower = error.to_ascii_lowercase();
1094 lower.contains("api key not found")
1095 || lower.contains("access token")
1096 || (lower.contains("credential")
1097 && (lower.contains("not found")
1098 || lower.contains("missing")
1099 || lower.contains("unavailable")
1100 || lower.contains("unsupported")))
1101 }
1102
1103 pub(crate) fn build_dispatch_error_closure(
1104 prepare: UserDispatchPrepare,
1105 recovery: DispatchRecovery,
1106 error: String,
1107 ) -> crate::tui::app::DispatchApplyFn {
1108 Box::new(
1109 move |app: &mut App,
1110 _engine_handle: &EngineHandle,
1111 _config: &Config|
1112 -> anyhow::Result<()> {
1113 app.dispatch_in_flight = false;
1114 // No operation was admitted, including route/reservation failures:
1115 // retire only cancellation introduced by this dispatch. A previous
1116 // admitted turn may still need to suppress its queued events.
1117 if !prepare.snapshot.suppress_stream_events_until_turn_complete {
1118 app.suppress_stream_events_until_turn_complete = false;
1119 }
1120 if prepare.cost_scope != crate::cost_status::scope_token() {
1121 anyhow::bail!("Message dispatch belongs to a previous session");
1122 }
1123 app.remote_control.fail_active_dispatch(&error);
1124 // Roll back the optimistic sync prepare mutations.
1125 app.is_loading = prepare.snapshot.is_loading;
1126 app.runtime_turn_status = prepare.snapshot.runtime_turn_status.clone();
1127 app.receipt_text = prepare.snapshot.receipt_text.clone();
1128 app.receipt_started_at = prepare.snapshot.receipt_started_at;
1129 app.tool_evidence = prepare.snapshot.tool_evidence.clone();
1130 let missing_credential = is_missing_credential_dispatch_error(&error);
1131 // Nothing was sent, so the optimistic echo goes too. A message
1132 // that is not in the transcript can be sent again without
1133 // appearing twice (#6566).
1134 app.history.truncate(prepare.snapshot.history_len);
1135 app.prune_transcript_index_state(prepare.snapshot.history_len);
1136 app.history_revisions
1137 .truncate(prepare.snapshot.history_revisions_len);
1138 // Never rewind the version: a cache keyed by (version, len) that
1139 // saw the rolled-back echo would match a different cell that
1140 // later lands at the same version and length.
1141 app.history_version = app.history_version.wrapping_add(1);
1142 app.truncate_api_messages(prepare.snapshot.api_messages_len);
1143 app.last_send_at = prepare.snapshot.last_send_at;
1144 app.needs_redraw = true;
1145
1146 match recovery {
1147 DispatchRecovery::Immediate => {
1148 if missing_credential {
1149 keep_unsent_message_for_connect(app, prepare.message, &error);
1150 } else {
1151 restore_failed_immediate_submit(
1152 app,
1153 prepare.message,
1154 &anyhow::Error::msg(error.clone()),
1155 );
1156 }
1157 }
1158 DispatchRecovery::Queued { restore_index } => {
1159 restore_queued_message(app, restore_index, prepare.message);
1160 app.status_message = Some(
1161 app.tr(MessageId::DispatchFailedQueued)
1162 .replace("{error}", &error)
1163 .replace("{count}", &app.queued_message_count().to_string()),
1164 );
1165 }
1166 DispatchRecovery::Draft => {
1167 restore_queued_or_draft_message(app, DispatchRecovery::Draft, prepare.message);
1168 app.status_message = Some(format!(
1169 "Message dispatch failed ({error}); queued draft restored"
1170 ));
1171 }
1172 DispatchRecovery::Initial => {
1173 if missing_credential {
1174 keep_unsent_message_for_connect(app, prepare.message, &error);
1175 } else {
1176 let initial_error = app
1177 .tr(MessageId::DispatchFailedInitial)
1178 .replace("{error}", &error);
1179 restore_failed_immediate_submit(
1180 app,
1181 prepare.message,
1182 &anyhow::Error::msg(initial_error),
1183 );
1184 }
1185 }
1186 }
1187
1188 Err(anyhow::Error::msg(error))
1189 },
1190 )
1191 }
1192
1193 pub(crate) fn parse_queue_send_command(input: &str) -> Option<Result<usize, String>> {
1194 let rest = strip_queue_command_prefix(input.trim())?;
1195 let mut parts = rest.split_whitespace();
1196 let action = parts.next()?;
1197 if !action.eq_ignore_ascii_case("send") && !action.eq_ignore_ascii_case("now") {
1198 return None;
1199 }
1200 let Some(raw_index) = parts.next() else {
1201 return Some(Err("Usage: /queue send <n>".to_string()));
1202 };
1203 if parts.next().is_some() {
1204 return Some(Err("Usage: /queue send <n>".to_string()));
1205 }
1206 let Ok(index) = raw_index.parse::<usize>() else {
1207 return Some(Err("Use a positive number".to_string()));
1208 };
1209 if index == 0 {
1210 return Some(Err("Use 1 or more".to_string()));
1211 }
1212 Some(Ok(index - 1))
1213 }
1214
1215 pub(crate) fn strip_queue_command_prefix(input: &str) -> Option<&str> {
1216 for prefix in ["/queue", "/queued"] {
1217 if let Some(rest) = input.strip_prefix(prefix)
1218 && (rest.is_empty() || rest.chars().next().is_some_and(char::is_whitespace))
1219 {
1220 return Some(rest);
1221 }
1222 }
1223 None
1224 }
1225
1226 pub(crate) async fn steer_user_message(
1227 app: &mut App,
1228 config: &Config,
1229 engine_handle: &EngineHandle,
1230 mut message: QueuedMessage,
1231 ) -> Result<bool> {
1232 let stop_words = config.stop_words();
1233 if is_stop_word(&message.display, &stop_words).is_some() {
1234 engine_handle.cancel();
1235 app.stopped_turn = true;
1236 app.status_message = Some("Turn stopped. Tool calls blocked for this turn.".to_string());
1237 return Ok(false);
1238 }
1239 app.stopped_turn = false;
1240 // Same-turn steering is an engine-bound external-user path just like a
1241 // fresh dispatch. Run the mutable gate exactly once on the blocking pool
1242 // before pause state, history, references, or engine input are touched.
1243 if app
1244 .hooks
1245 .has_hooks_for_event(crate::hooks::HookEvent::MessageSubmit)
1246 {
1247 let context = app.base_hook_context().with_message(&message.display);
1248 let strict_gates = app
1249 .hooks
1250 .matched_strict_gate_labels(crate::hooks::HookEvent::MessageSubmit, &context);
1251 let hooks = app.hooks.clone();
1252 let original_text = message.display.clone();
1253 let outcome = match tokio::task::spawn_blocking(move || {
1254 hooks.execute_message_submit_transform_for_dispatch(&context, &original_text)
1255 })
1256 .await
1257 {
1258 Ok(outcome) => outcome,
1259 Err(error) => {
1260 tracing::error!(target: "hooks", %error, "steer message_submit executor task was lost");
1261 lost_message_submit_outcome(&strict_gates)
1262 }
1263 };
1264 if !apply_message_submit_outcome(app, &mut message, outcome) {
1265 return Ok(false);
1266 }
1267 }
1268
1269 let paused_snapshot = snapshot_steer_paused_state(app);
1270 let paused_dispatch = plan_paused_command_message(app, &message.display);
1271 let paused_note = paused_dispatch.note().map(str::to_string);
1272 paused_dispatch.apply(app, engine_handle);
1273 let cwd = std::env::current_dir().ok();
1274 // Same single-submit cache as the other send path — see #4067 follow-up.
1275 let mut git_cache = crate::tui::git_mention::GitMentionCache::default();
1276 let completion_index = app.composer.mention_discovery.fuzzy_candidates(
1277 &app.workspace,
1278 &app.composer.mention_cwd,
1279 app.mention_walk_depth,
1280 app.workspace_follow_symlinks,
1281 );
1282 let references = crate::tui::file_mention::context_references_from_input_cached(
1283 &message.display,
1284 &app.workspace,
1285 cwd.clone(),
1286 &mut git_cache,
1287 completion_index,
1288 );
1289 let mut content = queued_message_content_for_app(app, &message, cwd, &mut git_cache)?;
1290 if let Some(note) = paused_note.as_deref() {
1291 content.push_str(note);
1292 }
1293 // Send exactly what the engine will store. `turn_loop` commits a steer as
1294 // `pending.commit().trim()`, so a composer newline or an appended note
1295 // left the held copy differing from the record by whitespace alone --
1296 // `accepted_steer_index` then never matched, the steer was never
1297 // promoted, and the "sending into this turn" card kept showing a message
1298 // the transcript had already delivered. Trimming here keeps that match an
1299 // exact comparison, which is the stronger invariant.
1300 let content = content.trim().to_string();
1301 let message_index = app.api_messages.len();
1302
1303 // A foreground shell blocks the turn loop that consumes steer input.
1304 // Ask the shared shell manager to detach it before enqueueing the steer so
1305 // the loop can leave the foreground wait and process this message (#4930).
1306 if active_foreground_shell_running(app)
1307 && let Err(err) = request_active_foreground_shell_background(app)
1308 {
1309 restore_steer_paused_state(app, &paused_snapshot);
1310 engine_handle.set_paused(paused_snapshot.paused);
1311 return Err(err.context("could not move foreground shell to /jobs before steering"));
1312 }
1313
1314 if let Err(err) = engine_handle.steer(content.clone()).await {
1315 restore_steer_paused_state(app, &paused_snapshot);
1316 engine_handle.set_paused(paused_snapshot.paused);
1317 return Err(err);
1318 }
1319 app.last_submitted_prompt = Some(message.display.clone());
1320
1321 // #6190: the steer channel accepting the text is not the turn accepting
1322 // it. The engine commits a steer at the next step boundary and discards
1323 // one whose turn has already moved on, so painting a settled cell and
1324 // pushing `api_messages` here produced two defects at once: the cell sat
1325 // above the assistant content the record places before it, and a dropped
1326 // steer left a transcript entry the model never saw. Hold it as in-flight
1327 // instead — it renders in the "sending into turn" preview until the
1328 // engine's own `SessionUpdated` shows it, which is also where it learns
1329 // its real message index.
1330 app.inflight_steers
1331 .push_back(crate::tui::app::InflightSteer {
1332 message,
1333 content,
1334 sent_after_index: message_index,
1335 references,
1336 });
1337 app.needs_redraw = true;
1338
1339 app.status_message = Some("Steering current turn...".to_string());
1340 Ok(true)
1341 }
1342
1343 /// Promote every in-flight steer the engine's record now contains.
1344 ///
1345 /// Called from `apply_engine_session_projection` after the projection lands,
1346 /// so the transcript cell is appended in the position the record gives it:
1347 /// below the assistant work that preceded the steer, as the newest entry.
1348 /// Matching is on the exact text handed to `EngineHandle::steer`, which the
1349 /// engine stores as the accepted user message's first text block, searched
1350 /// from the index the steer was sent after so an identical earlier message
1351 /// cannot claim it.
1352 pub(crate) fn settle_accepted_steers(app: &mut App) {
1353 if app.inflight_steers.is_empty() {
1354 return;
1355 }
1356 let mut claimed: Vec<usize> = Vec::new();
1357 let mut unsettled = VecDeque::new();
1358 for steer in std::mem::take(&mut app.inflight_steers) {
1359 let Some(index) = accepted_steer_index(app, &steer, &claimed) else {
1360 unsettled.push_back(steer);
1361 continue;
1362 };
1363 claimed.push(index);
1364 // Settle the streaming thinking/tool content that chronologically
1365 // preceded the steer before the steer's own cell is appended.
1366 app.flush_active_cell();
1367 let display = format!("+ {}", steer.message.display);
1368 let history_cell = paint_user_turn_cell(app, &steer.message, display);
1369 app.record_context_references(history_cell, index, steer.references);
1370 app.needs_redraw = true;
1371 }
1372 app.inflight_steers = unsettled;
1373 }
1374
1375 fn accepted_steer_index(
1376 app: &App,
1377 steer: &crate::tui::app::InflightSteer,
1378 claimed: &[usize],
1379 ) -> Option<usize> {
1380 let start = steer.sent_after_index.min(app.api_messages.len());
1381 app.api_messages
1382 .iter()
1383 .enumerate()
1384 .skip(start)
1385 .find(|(index, message)| {
1386 !claimed.contains(index)
1387 && message.role == Role::User
1388 && matches!(
1389 message.content.first(),
1390 Some(ContentBlock::Text { text, .. }) if text == &steer.content
1391 )
1392 })
1393 .map(|(index, _)| index)
1394 }
1395
1396 /// A turn that ended without accepting a steer must not swallow it (#6190,
1397 /// #6297). The message becomes a queued follow-up — the queue is the one path
1398 /// that actually drains into a turn, where the engine's context-pressure gate
1399 /// sees it like any other send — instead of a display-only "rejected" string
1400 /// that nothing ever dispatches.
1401 pub(crate) fn settle_unaccepted_steers_at_turn_end(app: &mut App) {
1402 if app.inflight_steers.is_empty() {
1403 return;
1404 }
1405 let deferred = std::mem::take(&mut app.inflight_steers);
1406 app.queued_messages
1407 .extend(deferred.into_iter().map(|steer| steer.message));
1408 app.needs_redraw = true;
1409 }
1410
1411 pub(crate) fn snapshot_steer_paused_state(app: &App) -> SteerPausedSnapshot {
1412 SteerPausedSnapshot {
1413 paused: app.paused,
1414 pausable: app.pausable,
1415 paused_goal_objective: app.paused_goal_objective.clone(),
1416 objective: app.goal.objective.clone(),
1417 tokens_used: app.goal.tokens_used,
1418 time_used_seconds: app.goal.time_used_seconds,
1419 continuation_count: app.goal.continuation_count,
1420 }
1421 }
1422
1423 pub(crate) fn restore_steer_paused_state(app: &mut App, snapshot: &SteerPausedSnapshot) {
1424 app.paused = snapshot.paused;
1425 app.pausable = snapshot.pausable;
1426 app.paused_goal_objective = snapshot.paused_goal_objective.clone();
1427 app.goal.objective = snapshot.objective.clone();
1428 app.goal.tokens_used = snapshot.tokens_used;
1429 app.goal.time_used_seconds = snapshot.time_used_seconds;
1430 app.goal.continuation_count = snapshot.continuation_count;
1431 }
1432
1433 pub(crate) async fn attempt_steer_with_queue_fallback(
1434 app: &mut App,
1435 config: &Config,
1436 engine_handle: &EngineHandle,
1437 message: QueuedMessage,
1438 recovery: DispatchRecovery,
1439 ) -> bool {
1440 match steer_user_message(app, config, engine_handle, message.clone()).await {
1441 Ok(true) => {
1442 app.push_status_toast(
1443 app.tr(MessageId::ToastSentIntoTurn).into_owned(),
1444 StatusToastLevel::Info,
1445 Some(1_500),
1446 );
1447 true
1448 }
1449 Ok(false) => {
1450 restore_queued_or_draft_message(app, recovery, message);
1451 app.push_status_toast(
1452 app.tr(MessageId::ToastHookBlockedFollowUp).into_owned(),
1453 StatusToastLevel::Warning,
1454 Some(4_000),
1455 );
1456 false
1457 }
1458 Err(err) => {
1459 restore_queued_or_draft_message(app, recovery, message);
1460 let status = format!("{} ({err})", app.tr(MessageId::ToastCouldNotSendIntoTurn));
1461 app.status_message = Some(status.clone());
1462 app.push_status_toast(status, StatusToastLevel::Warning, Some(4_000));
1463 false
1464 }
1465 }
1466 }
1467
1468 /// Park a draft on the queued-messages bucket for dispatch after TurnComplete.
1469 /// Unlike a steer, the message is NOT forwarded immediately — it waits for
1470 /// the current turn to finish, then dispatches as a normal user message.
1471 pub(crate) async fn queue_follow_up(app: &mut App, message: QueuedMessage) -> Result<()> {
1472 let mut message = message;
1473 echo_queued_user_turn(app, &mut message);
1474 enqueue_offline_message(app, message);
1475 let toast = queued_follow_up_toast(app);
1476 app.status_message = Some(toast.clone());
1477 app.push_status_toast(toast, StatusToastLevel::Info, Some(3_000));
1478 Ok(())
1479 }
1480
1481 fn queued_follow_up_toast(app: &App) -> String {
1482 if app.offline_mode {
1483 return app.tr(MessageId::ToastQueuedOffline).into_owned();
1484 }
1485 let count = app.queued_message_count();
1486 if count <= 1 {
1487 app.tr(MessageId::ToastQueuedFollowUp).into_owned()
1488 } else {
1489 app.tr(MessageId::ToastQueuedFollowUpCount)
1490 .replace("{count}", &count.to_string())
1491 }
1492 }
1493
1494 pub(crate) async fn dispatch_composer_message(
1495 app: &mut App,
1496 config: &Config,
1497 engine_handle: &EngineHandle,
1498 message: QueuedMessage,
1499 recovery: DispatchRecovery,
1500 action: ComposerSubmitAction,
1501 ) -> Result<()> {
1502 // The ordinary web mirror does not block local prompts. Isolated Runtime
1503 // Chat is different: it owns a provider-backed native turn that is not
1504 // reflected by the interactive engine's `is_loading` bit. Preserve the
1505 // local message in the queue until that exact turn and its terminal relay
1506 // receipt have settled, so one attached run never has two inference loops.
1507 if app.remote_control.runtime_chat_blocks_local_dispatch() {
1508 let mut message = message;
1509 echo_queued_user_turn(app, &mut message);
1510 enqueue_offline_message(app, message);
1511 let queued = app
1512 .tr(MessageId::AgentRailQueuedCount)
1513 .replace("{count}", &app.queued_message_count().to_string());
1514 let status = format!(
1515 "Runtime Chat · {} · {queued}",
1516 app.tr(MessageId::PhaseFinishing)
1517 );
1518 app.push_status_toast(status, StatusToastLevel::Info, Some(4_000));
1519 return Ok(());
1520 }
1521
1522 // Agent focus: the composer addresses one child's fork, not the main
1523 // session. The follow-up is real runtime work (Op::FollowUpSubAgent); the
1524 // main transcript keeps a receipt line and the focused view echoes the
1525 // message until the child's own transcript carries it.
1526 if let Some(focus) = app.agent_focus.as_ref() {
1527 let agent_id = focus.agent_id.clone();
1528 let label = focus.label.clone();
1529 let text = message.display.clone();
1530 crate::tui::agent_focus::echo_user_follow_up(app, &text);
1531 let receipt = app
1532 .tr(codewhale_localization::MessageId::AgentFocusFollowUpQueued)
1533 .replace("{agent}", &label);
1534 app.push_history_cell(crate::tui::history::HistoryCell::System { content: receipt });
1535 // #6150: the input path never awaits a full op channel. The follow-up
1536 // is retryable; a rejected send surfaces immediately.
1537 if let Err(err) = engine_handle.try_send(crate::core::ops::Op::FollowUpSubAgent {
1538 agent_id: agent_id.clone(),
1539 text,
1540 }) {
1541 let reason = if err
1542 .downcast_ref::<tokio::sync::mpsc::error::TrySendError<crate::core::ops::Op>>()
1543 .is_some_and(|e| matches!(e, tokio::sync::mpsc::error::TrySendError::Full(_)))
1544 {
1545 "engine busy"
1546 } else {
1547 "engine unavailable"
1548 };
1549 let failed = app
1550 .tr(codewhale_localization::MessageId::AgentFocusFollowUpFailed)
1551 .replace("{agent}", &label)
1552 .replace("{reason}", reason);
1553 app.status_message = Some(failed.clone());
1554 app.push_status_toast(failed, StatusToastLevel::Warning, Some(5_000));
1555 }
1556 return Ok(());
1557 }
1558 let disposition = match action {
1559 ComposerSubmitAction::Submit(disposition) => disposition,
1560 ComposerSubmitAction::SendQueuedNow | ComposerSubmitAction::Noop => {
1561 // The caller extracted a non-empty input, so these can only arise
1562 // if state changed between key resolution and dispatch. Queueing
1563 // is lossless and preserves ordering in that narrow race.
1564 SubmitDisposition::Queue
1565 }
1566 };
1567 match disposition {
1568 SubmitDisposition::Immediate => {
1569 let _ =
1570 dispatch_user_message_with_recovery(app, config, engine_handle, message, recovery)
1571 .await;
1572 Ok(())
1573 }
1574 SubmitDisposition::Queue => {
1575 let mut message = message;
1576 echo_queued_user_turn(app, &mut message);
1577 enqueue_offline_message(app, message);
1578 // A second, empty Enter inside the window sends this now.
1579 app.arm_double_tap_window();
1580 let toast = queued_follow_up_toast(app);
1581 app.status_message = Some(toast.clone());
1582 app.push_status_toast(toast, StatusToastLevel::Info, Some(3_000));
1583 Ok(())
1584 }
1585 SubmitDisposition::Steer => {
1586 attempt_steer_with_queue_fallback(app, config, engine_handle, message, recovery).await;
1587 Ok(())
1588 }
1589 SubmitDisposition::QueueFollowUp => queue_follow_up(app, message).await,
1590 }
1591 }
1592
1593 #[cfg(test)]
1594 pub(crate) async fn submit_or_steer_message(
1595 app: &mut App,
1596 config: &Config,
1597 engine_handle: &EngineHandle,
1598 message: QueuedMessage,
1599 recovery: DispatchRecovery,
1600 ) -> Result<()> {
1601 let action = ComposerSubmitAction::Submit(app.decide_submit_disposition());
1602 dispatch_composer_message(app, config, engine_handle, message, recovery, action).await
1603 }
1604
1605 /// Drain `app.pending_steers` into a single `QueuedMessage` ready for
1606 /// `dispatch_user_message`. Returns `None` if the queue was empty (caller
1607 /// then falls back to `app.queued_messages`). Skill instruction is taken
1608 /// from the first message that supplies one — multiple steers shouldn't
1609 /// double-up the system framing.
1610 pub(crate) fn merge_pending_steers(app: &mut App) -> Option<QueuedMessage> {
1611 let drained = app.drain_pending_steers();
1612 if drained.is_empty() {
1613 return None;
1614 }
1615 if drained.len() == 1 {
1616 return drained.into_iter().next();
1617 }
1618 let mut skill_instruction: Option<String> = None;
1619 let mut skill_provenance = None;
1620 let mut bodies: Vec<String> = Vec::with_capacity(drained.len());
1621 for msg in drained {
1622 if skill_instruction.is_none() {
1623 skill_instruction = msg.skill_instruction;
1624 skill_provenance = msg.skill_provenance;
1625 }
1626 bodies.push(msg.display);
1627 }
1628 Some(
1629 QueuedMessage::new(bodies.join("\n\n"), skill_instruction)
1630 .with_skill_provenance(skill_provenance),
1631 )
1632 }
1633
1633 lines RUST