use super::*;
pub(crate) fn dispatch_hotbar_slot(
app: &mut App,
config: &Config,
slot: u8,
) -> Result<Option<HotbarDispatch>> {
let known_action_ids = app
.hotbar_actions
.iter()
.map(|action| action.id())
.collect::<Vec<_>>();
let bindings = config.resolve_hotbar_bindings(&known_action_ids).bindings;
let Some(action_id) = bindings
.iter()
.find(|binding| binding.slot == slot)
.map(|binding| binding.action.clone())
else {
return Ok(None);
};
let Some(action) = app.hotbar_actions.get(&action_id) else {
app.status_message = Some(format!(
"Hotbar slot {slot} action is not available: {action_id}"
));
app.needs_redraw = true;
return Ok(Some(HotbarDispatch::Handled));
};
if let Some(reason) = action.disabled_reason(app) {
app.status_message = Some(format!(
"Hotbar slot {slot} action is not available: {reason}"
));
app.needs_redraw = true;
return Ok(Some(HotbarDispatch::Handled));
}
action.dispatch(app).map(Some)
}
pub(crate) fn queued_ui_to_session(msg: &QueuedMessage) -> QueuedSessionMessage {
QueuedSessionMessage {
display: msg.display.clone(),
skill_instruction: msg.skill_instruction.clone(),
skill_provenance: msg.skill_provenance.clone(),
}
}
pub(crate) fn queued_session_to_ui(msg: QueuedSessionMessage) -> QueuedMessage {
QueuedMessage {
display: msg.display,
skill_instruction: msg.skill_instruction,
skill_provenance: msg.skill_provenance,
}
}
pub(crate) fn enqueue_offline_message(app: &mut App, message: QueuedMessage) {
app.queue_message(message);
persist_offline_queue_state(app);
}
pub(crate) fn push_assistant_message(
app: &mut App,
text: String,
thinking: Option<String>,
tool_uses: PendingToolUses,
) {
let mut blocks = Vec::new();
if let Some(thinking) = thinking {
blocks.push(ContentBlock::Thinking {
thinking,
signature: None,
state: None,
});
}
if !text.is_empty() {
blocks.push(ContentBlock::Text {
text,
cache_control: None,
});
}
for (id, name, input) in tool_uses {
blocks.push(ContentBlock::ToolUse {
id,
name,
input,
caller: None,
thought_signature: None,
});
}
let has_sendable_content = blocks.iter().any(|block| {
matches!(
block,
ContentBlock::Text { .. } | ContentBlock::ToolUse { .. }
)
});
if has_sendable_content {
app.api_messages.push(Message {
role: "assistant".to_string(),
content: blocks,
});
}
}
pub(crate) fn replace_matching_assistant_text(
app: &mut App,
original_text: &str,
translated_text: String,
) -> bool {
for message in app.api_messages.iter_mut().rev() {
if message.role != "assistant" && message.role != crate::models::INTERRUPTED_ASSISTANT_ROLE
{
continue;
}
for block in &mut message.content {
if let ContentBlock::Text { text, .. } = block
&& text == original_text
{
*text = translated_text;
return true;
}
}
}
false
}
pub(crate) fn build_queued_message(app: &mut App, input: String) -> QueuedMessage {
let skill_instruction = app.active_skill.take();
let skill_provenance = app.active_skill_provenance.take();
QueuedMessage::new(input, skill_instruction).with_skill_provenance(skill_provenance)
}
pub(crate) async fn submit_initial_input_if_ready(
app: &mut App,
config: &Config,
engine_handle: &EngineHandle,
) -> Result<()> {
if !app.auto_submit_initial_input {
return Ok(());
}
if app
.view_stack
.contains_kind(crate::tui::views::ModalKind::TelemetryNotice)
{
return Ok(());
}
if app.onboarding != OnboardingState::None {
if app.status_message.is_none() && !app.input.trim().is_empty() {
app.status_message = Some(INITIAL_PROMPT_DEFERRED_STATUS.to_string());
}
return Ok(());
}
app.auto_submit_initial_input = false;
if let Some(input) = app.submit_input() {
if app.status_message.as_deref() == Some(INITIAL_PROMPT_DEFERRED_STATUS) {
app.status_message = None;
}
let queued = build_queued_message(app, input);
dispatch_user_message_with_recovery(
app,
config,
engine_handle,
queued,
DispatchRecovery::Initial,
)
.await?;
}
Ok(())
}
pub(crate) fn message_from_submitted_input(
app: &mut App,
input: String,
) -> (QueuedMessage, DispatchRecovery) {
if let Some(mut draft) = app.queued_draft.take() {
draft.display = input;
(draft, DispatchRecovery::Draft)
} else {
(
build_queued_message(app, input),
DispatchRecovery::Immediate,
)
}
}
pub(crate) fn take_next_queued_message(app: &mut App) -> Option<(QueuedMessage, DispatchRecovery)> {
if app.input.is_empty() {
return app.remove_queued_message(0).map(|message| {
(
message,
DispatchRecovery::Queued {
restore_index: Some(0),
},
)
});
}
None
}
pub(crate) async fn send_next_queued_message_now(
app: &mut App,
config: &Config,
engine_handle: &EngineHandle,
) -> Result<bool> {
let Some((message, recovery)) = take_next_queued_message(app) else {
return Ok(false);
};
send_taken_queued_message_now(app, config, engine_handle, message, recovery).await?;
Ok(true)
}
pub(crate) async fn send_queued_message_at_index_now(
app: &mut App,
config: &Config,
engine_handle: &EngineHandle,
index: usize,
) -> Result<bool> {
let Some(message) = app.remove_queued_message(index) else {
app.status_message = Some("Queued message not found".to_string());
return Ok(true);
};
send_taken_queued_message_now(
app,
config,
engine_handle,
message,
DispatchRecovery::Queued {
restore_index: Some(index),
},
)
.await?;
Ok(true)
}
pub(crate) async fn send_taken_queued_message_now(
app: &mut App,
config: &Config,
engine_handle: &EngineHandle,
message: QueuedMessage,
recovery: DispatchRecovery,
) -> Result<()> {
if app.offline_mode {
restore_queued_or_draft_message(app, recovery, message);
app.status_message = Some(format!(
"Offline: {} queued follow-up(s) — /queue send <n>, /queue clear",
app.queued_message_count()
));
return Ok(());
}
let display = message.display.clone();
if app.dispatch_in_flight {
restore_queued_or_draft_message(app, recovery, message);
app.status_message = Some(format!(
"{} queued follow-up(s) — sends after current dispatch starts",
app.queued_message_count()
));
return Ok(());
}
if app.is_loading {
match steer_user_message(app, config, engine_handle, message.clone()).await {
Ok(true) => app.push_status_toast(
"Sent queued follow-up into current turn",
StatusToastLevel::Info,
Some(1_500),
),
Ok(false) => {
restore_queued_or_draft_message(app, recovery, message);
app.push_status_toast(
"message_submit hook blocked the follow-up; original queue/draft restored",
StatusToastLevel::Warning,
Some(4_000),
);
}
Err(err) => {
restore_queued_or_draft_message(app, recovery, message);
app.status_message = Some(format!(
"Steer failed ({err}); {} queued follow-up(s) — /queue send <n>, /queue clear",
app.queued_message_count()
));
}
}
} else if let Err(_err) =
dispatch_user_message_with_recovery(app, config, engine_handle, message, recovery).await
{
} else {
app.status_message = Some(format!("Sent queued follow-up: {display}"));
}
Ok(())
}
pub(crate) fn queued_message_content_for_app(
app: &App,
message: &QueuedMessage,
cwd: Option<PathBuf>,
git_cache: &mut crate::tui::git_mention::GitMentionCache,
) -> Result<String> {
if let Some(authority) = message.skill_provenance.as_ref() {
if authority.workspace != app.workspace {
anyhow::bail!("Queued plugin skill belongs to a different workspace and was denied");
}
crate::plugins::registry::verify_plugin_component_authority(
authority,
crate::plugins::activation::PluginActivationCapability::Skills,
)
.map_err(anyhow::Error::msg)?;
}
let completion_index = app.composer.mention_discovery.fuzzy_candidates(
&app.workspace,
&app.composer.mention_cwd,
app.mention_walk_depth,
app.workspace_follow_symlinks,
);
let stabilization_dir = crate::tui::file_mention::screenshot_stabilization_dir(&app.workspace);
let display = crate::tui::file_mention::stabilize_screenshot_references(
&message.display,
&stabilization_dir,
);
let user_request = crate::tui::file_mention::user_request_with_file_mentions_cached(
&display,
&app.workspace,
cwd,
git_cache,
completion_index,
);
if let Some(skill_instruction) = message.skill_instruction.as_ref() {
Ok(format!(
"{skill_instruction}\n\n---\n\nUser request: {user_request}"
))
} else {
Ok(user_request)
}
}
pub(crate) fn dispatch_completion_permit(
app: &App,
) -> std::result::Result<
tokio::sync::mpsc::OwnedPermit<crate::tui::app::DispatchApplyFn>,
&'static str,
> {
let sender = app
.dispatch_completion_tx
.clone()
.ok_or("dispatch completion mailbox is unavailable")?;
sender.try_reserve_owned().map_err(|error| match error {
tokio::sync::mpsc::error::TrySendError::Full(_) => "dispatch completion mailbox is full",
tokio::sync::mpsc::error::TrySendError::Closed(_) => {
"dispatch completion mailbox is closed"
}
})
}
#[cfg(test)]
pub(crate) async fn dispatch_user_message(
app: &mut App,
config: &Config,
engine_handle: &EngineHandle,
message: QueuedMessage,
) -> Result<()> {
dispatch_user_message_with_recovery(
app,
config,
engine_handle,
message,
DispatchRecovery::Immediate,
)
.await
}
pub(crate) async fn dispatch_user_message_with_recovery(
app: &mut App,
config: &Config,
engine_handle: &EngineHandle,
mut message: QueuedMessage,
recovery: DispatchRecovery,
) -> Result<()> {
let stop_words = config.stop_words();
if is_stop_word(&message.display, &stop_words).is_some() {
engine_handle.cancel();
app.stopped_turn = true;
app.status_message = Some("Turn stopped. Tool calls blocked for this turn.".to_string());
return Ok(());
}
app.stopped_turn = false;
if app
.hooks
.has_hooks_for_event(crate::hooks::HookEvent::MessageSubmit)
{
let context = app.base_hook_context().with_message(&message.display);
let strict_gates = app
.hooks
.matched_strict_gate_labels(crate::hooks::HookEvent::MessageSubmit, &context);
let hooks = app.hooks.clone();
let original_text = message.display.clone();
if app.dispatch_completion_tx.is_some() {
let completion_permit = match dispatch_completion_permit(app) {
Ok(permit) => permit,
Err(error) => {
recover_unstarted_external_message(app, message, recovery, error);
return Err(anyhow::Error::msg(error));
}
};
app.dispatch_in_flight = true;
tokio::spawn(async move {
let outcome = match tokio::task::spawn_blocking(move || {
hooks.execute_message_submit_transform_for_dispatch(&context, &original_text)
})
.await
{
Ok(outcome) => outcome,
Err(error) => {
tracing::error!(target: "hooks", %error, "message_submit executor task was lost");
lost_message_submit_outcome(&strict_gates)
}
};
let apply: crate::tui::app::DispatchApplyFn = Box::new(
move |app: &mut App,
engine_handle: &EngineHandle,
config: &Config|
-> anyhow::Result<()> {
if !apply_message_submit_outcome(app, &mut message, outcome) {
app.dispatch_in_flight = false;
restore_message_submit_denial(app, message, recovery);
return Ok(());
}
let _ = start_user_dispatch(app, config, engine_handle, message, recovery);
Ok(())
},
);
completion_permit.send(apply);
});
return Ok(());
}
let outcome = match tokio::task::spawn_blocking(move || {
hooks.execute_message_submit_transform_for_dispatch(&context, &original_text)
})
.await
{
Ok(outcome) => outcome,
Err(error) => {
tracing::error!(target: "hooks", %error, "message_submit executor task was lost");
lost_message_submit_outcome(&strict_gates)
}
};
if !apply_message_submit_outcome(app, &mut message, outcome) {
restore_message_submit_denial(app, message, recovery);
return Ok(());
}
}
if app.dispatch_completion_tx.is_some() {
return start_user_dispatch(app, config, engine_handle, message, recovery);
}
let prepare = match prepare_user_dispatch(app, config, message.clone()) {
Ok(prepare) => prepare,
Err(error) => {
recover_unstarted_external_message(app, message, recovery, &error.to_string());
return Err(error);
}
};
run_prepared_dispatch(app, config, engine_handle, prepare, recovery).await
}
pub(crate) fn lost_message_submit_outcome(
strict_gates: &[String],
) -> crate::hooks::MessageSubmitOutcome {
if strict_gates.is_empty() {
crate::hooks::MessageSubmitOutcome::Unchanged {
warning: Some(
"message_submit hook executor did not run; submission continued because no strict gate matched"
.to_string(),
),
}
} else {
crate::hooks::MessageSubmitOutcome::Blocked {
reason: "message_submit hook executor did not run; a strict gate blocked submission"
.to_string(),
}
}
}
pub(crate) fn prepare_user_dispatch(
app: &mut App,
config: &Config,
message: QueuedMessage,
) -> Result<UserDispatchPrepare> {
let _ = app.maybe_nudge_for_planning_prompt(&message.display);
let paused_dispatch = plan_paused_command_message(app, &message.display);
let cwd = std::env::current_dir().ok();
let mut git_cache = crate::tui::git_mention::GitMentionCache::default();
let completion_index = app.composer.mention_discovery.fuzzy_candidates(
&app.workspace,
&app.composer.mention_cwd,
app.mention_walk_depth,
app.workspace_follow_symlinks,
);
let references = crate::tui::file_mention::context_references_from_input_cached(
&message.display,
&app.workspace,
cwd.clone(),
&mut git_cache,
completion_index,
);
let mut content = queued_message_content_for_app(app, &message, cwd, &mut git_cache)?;
if let Some(note) = paused_dispatch.note() {
content.push_str(note);
}
let (app_route_identity, route_config) = app_scoped_runtime_config(app, config);
let should_auto_resolve = auto_router::should_resolve_auto_model_selection(app);
let auto_router_context = auto_router::recent_auto_router_context(&app.api_messages);
let snapshot = UserDispatchSnapshot {
is_loading: app.is_loading,
runtime_turn_status: app.runtime_turn_status.clone(),
receipt_text: app.receipt_text.clone(),
receipt_started_at: app.receipt_started_at,
tool_evidence: app.tool_evidence.clone(),
history_len: app.history.len(),
history_revisions_len: app.history_revisions.len(),
history_version: app.history_version,
api_messages_len: app.api_messages.len(),
last_send_at: app.last_send_at,
};
app.is_loading = true;
app.runtime_turn_status = None;
app.clear_receipt();
app.tool_evidence.clear();
app.needs_redraw = true;
let message_index = app.api_messages.len();
app.add_message(HistoryCell::User {
content: message.display.clone(),
});
let history_cell = app.history.len().saturating_sub(1);
app.scroll_to_bottom();
app.last_send_at = Some(Instant::now());
app.api_messages.push(Message {
role: "user".to_string(),
content: vec![ContentBlock::Text {
text: content.clone(),
cache_control: None,
}],
});
let goal_objective = paused_dispatch.goal_objective(app);
Ok(UserDispatchPrepare {
message,
content,
references,
paused_dispatch,
app_route_identity,
route_config,
goal_objective,
goal_status: app.hunt.verdict.goal_status(),
goal_token_budget: app.hunt.token_budget,
mode: app.mode,
api_provider: app.api_provider,
app_model: app.model.clone(),
auto_model: app.auto_model,
reasoning_effort: app.reasoning_effort,
allow_shell: app.allow_shell,
trust_mode: app.trust_mode,
auto_approve: app_auto_approve_enabled(app),
approval_mode: app.approval_mode,
translation_enabled: app.translation_enabled,
allowed_tools: app.active_allowed_tools.clone(),
hook_executor: app.runtime_services.hook_executor.clone(),
verbosity: app.verbosity.clone(),
provenance: UserInputProvenance::ExternalUser,
auto_router_context,
should_auto_resolve,
auto_compact_user_configured: app.auto_compact_user_configured,
auto_compact: app.auto_compact,
auto_compact_threshold_percent: app.auto_compact_threshold_percent,
snapshot,
message_index,
history_cell,
})
}
pub(crate) fn start_user_dispatch(
app: &mut App,
config: &Config,
engine_handle: &EngineHandle,
message: QueuedMessage,
recovery: DispatchRecovery,
) -> Result<()> {
let completion_permit = match dispatch_completion_permit(app) {
Ok(permit) => permit,
Err(error) => {
recover_unstarted_external_message(app, message, recovery, error);
return Err(anyhow::Error::msg(error));
}
};
let recovery_message = message.clone();
let prepare = match prepare_user_dispatch(app, config, message) {
Ok(prepare) => prepare,
Err(error) => {
recover_unstarted_external_message(app, recovery_message, recovery, &error.to_string());
return Err(error);
}
};
app.dispatch_in_flight = true;
tokio::spawn(spawned_dispatch_execute(
prepare,
recovery,
engine_handle.clone(),
completion_permit,
));
Ok(())
}
pub(crate) async fn spawned_dispatch_execute(
prepare: UserDispatchPrepare,
recovery: DispatchRecovery,
engine_handle: EngineHandle,
completion_permit: tokio::sync::mpsc::OwnedPermit<crate::tui::app::DispatchApplyFn>,
) {
let apply = spawned_dispatch_inner(prepare, recovery, engine_handle).await;
completion_permit.send(apply);
}
pub(crate) async fn spawned_dispatch_inner(
prepare: UserDispatchPrepare,
recovery: DispatchRecovery,
engine_handle: EngineHandle,
) -> crate::tui::app::DispatchApplyFn {
let plan_result = plan_turn_route(TurnRoutePlanRequest {
route_config: &prepare.route_config,
app_route_identity: &prepare.app_route_identity,
api_provider: prepare.api_provider,
app_model: &prepare.app_model,
auto_model: prepare.auto_model,
reasoning_effort: prepare.reasoning_effort,
mode: prepare.mode,
content: &prepare.content,
display_text: &prepare.message.display,
auto_router_context: &prepare.auto_router_context,
should_auto_resolve: prepare.should_auto_resolve,
allow_auto_router_response_cache: true,
preflight_required: engine_handle.client_preflight_required(),
auto_compact_user_configured: prepare.auto_compact_user_configured,
auto_compact: prepare.auto_compact,
auto_compact_threshold_percent: prepare.auto_compact_threshold_percent,
})
.await;
let planned = match plan_result {
Ok(planned) => planned,
Err(err) => return build_dispatch_error_closure(prepare, recovery, err),
};
let PlannedTurnRoute {
route: turn_route,
compaction: turn_compaction,
effective_provider,
effective_model,
effective_provider_identity,
effective_provider_label,
selected_reasoning_effort,
effective_reasoning_effort,
auto_controls_reasoning,
auto_selection,
routing_source: _,
} = planned;
let effective_reasoning_tier = selected_reasoning_effort
.unwrap_or(prepare.reasoning_effort)
.normalize_for_route(
effective_provider,
&turn_route.candidate.endpoint().base_url,
&turn_route.model,
);
let effective_reasoning_receipt = reasoning_effort_receipt_for_route(
effective_reasoning_tier,
effective_provider,
&turn_route.candidate.endpoint().base_url,
&turn_route.model,
);
if let Err(err) = engine_handle
.send(Op::SendMessage {
content: prepare.content.clone(),
mode: prepare.mode,
route: Box::new(turn_route),
compaction: Box::new(turn_compaction.clone()),
goal_objective: prepare.goal_objective.clone(),
goal_token_budget: prepare.goal_token_budget,
goal_status: prepare.goal_status,
reasoning_effort: effective_reasoning_effort,
reasoning_effort_auto: auto_controls_reasoning,
auto_model: prepare.auto_model,
allow_shell: prepare.allow_shell,
trust_mode: prepare.trust_mode,
auto_approve: prepare.auto_approve,
approval_mode: prepare.approval_mode,
translation_enabled: prepare.translation_enabled,
allowed_tools: prepare.allowed_tools.clone(),
dynamic_tools: Vec::new(),
hook_executor: prepare.hook_executor.clone(),
verbosity: prepare.verbosity.clone(),
provenance: prepare.provenance,
})
.await
{
return build_dispatch_error_closure(prepare, recovery, err.to_string());
}
build_dispatch_success_closure(
prepare,
UserDispatchOutcome {
turn_compaction,
effective_provider,
effective_model,
effective_provider_identity,
effective_provider_label,
effective_reasoning_effort: effective_reasoning_receipt,
auto_selection,
},
)
}
pub(crate) fn build_dispatch_success_closure(
prepare: UserDispatchPrepare,
outcome: UserDispatchOutcome,
) -> crate::tui::app::DispatchApplyFn {
Box::new(
move |app: &mut App, engine_handle: &EngineHandle, config: &Config| -> anyhow::Result<()> {
app.dispatch_in_flight = false;
prepare.paused_dispatch.apply(app, engine_handle);
let dispatch_started_at = Instant::now();
app.is_loading = true;
app.dispatch_started_at = Some(dispatch_started_at);
app.runtime_turn_status = None;
app.last_submitted_prompt = Some(prepare.message.display.clone());
app.clear_receipt();
app.tool_evidence.clear();
app.system_prompt = Some(build_app_system_prompt_with_goal(
app,
config,
app.hunt.quarry.as_deref(),
));
app.record_context_references(
prepare.history_cell,
prepare.message_index,
prepare.references,
);
app.scroll_to_bottom();
app.last_effective_reasoning_effort = Some(outcome.effective_reasoning_effort);
if prepare.auto_model {
app.last_effective_model = Some(outcome.effective_model.clone());
app.last_effective_provider = Some(outcome.effective_provider);
app.last_effective_provider_identity =
Some(outcome.effective_provider_identity.clone());
if let Some(selection) = outcome.auto_selection.as_ref() {
app.last_auto_route_receipt = selection.receipt.clone();
let status = app
.tr(MessageId::AutoRouteSelectedToast)
.replace("{provider}", &outcome.effective_provider_label)
.replace("{model}", &outcome.effective_model)
.replace("{source}", selection.source.label());
app.push_status_toast(status, StatusToastLevel::Info, Some(6_000));
}
} else {
app.last_effective_model = None;
app.last_effective_provider = None;
app.last_effective_provider_identity = None;
app.last_auto_route_receipt = None;
}
app.pending_auto_route_receipt = outcome
.auto_selection
.as_ref()
.and_then(|selection| selection.receipt.clone());
app.pending_turn_route = Some((
outcome.effective_provider,
outcome.effective_model,
prepare.auto_model,
));
maybe_warn_context_pressure_for_config(app, &outcome.turn_compaction);
app.session.last_prompt_tokens = None;
app.session.last_completion_tokens = None;
app.session.last_output_throughput = None;
app.session.last_prompt_cache_hit_tokens = None;
app.session.last_prompt_cache_miss_tokens = None;
app.session.last_reasoning_replay_tokens = None;
if let Ok(manager) = SessionManager::default_location()
&& let Ok(session) = build_session_snapshot(app, &manager)
{
if app.current_session_id.is_none() {
app.current_session_id = Some(session.metadata.id.clone());
}
if let Err(err) = persist_with_pending_work_boundary(
app,
PersistRequest::SaveCheckpoint { session },
) {
app.status_message = Some(format!(
"To-do list update pending: turn checkpoint could not be queued ({err})"
));
}
}
Ok(())
},
)
}
pub(crate) fn build_dispatch_error_closure(
prepare: UserDispatchPrepare,
recovery: DispatchRecovery,
error: String,
) -> crate::tui::app::DispatchApplyFn {
Box::new(
move |app: &mut App,
_engine_handle: &EngineHandle,
_config: &Config|
-> anyhow::Result<()> {
app.remote_control.fail_active_dispatch(&error);
app.dispatch_in_flight = false;
app.is_loading = prepare.snapshot.is_loading;
app.runtime_turn_status = prepare.snapshot.runtime_turn_status.clone();
app.receipt_text = prepare.snapshot.receipt_text.clone();
app.receipt_started_at = prepare.snapshot.receipt_started_at;
app.tool_evidence = prepare.snapshot.tool_evidence.clone();
app.history.truncate(prepare.snapshot.history_len);
app.prune_transcript_index_state(prepare.snapshot.history_len);
app.history_revisions
.truncate(prepare.snapshot.history_revisions_len);
app.history_version = prepare.snapshot.history_version;
app.api_messages.truncate(prepare.snapshot.api_messages_len);
app.last_send_at = prepare.snapshot.last_send_at;
app.needs_redraw = true;
match recovery {
DispatchRecovery::Immediate => {
restore_failed_immediate_submit(
app,
prepare.message,
&anyhow::Error::msg(error.clone()),
);
}
DispatchRecovery::Queued { restore_index } => {
restore_queued_message(app, restore_index, prepare.message);
app.status_message = Some(
app.tr(MessageId::DispatchFailedQueued)
.replace("{error}", &error)
.replace("{count}", &app.queued_message_count().to_string()),
);
}
DispatchRecovery::Draft => {
restore_queued_or_draft_message(app, DispatchRecovery::Draft, prepare.message);
app.status_message = Some(format!(
"Message dispatch failed ({error}); queued draft restored"
));
}
DispatchRecovery::Initial => {
let initial_error = app
.tr(MessageId::DispatchFailedInitial)
.replace("{error}", &error);
restore_failed_immediate_submit(
app,
prepare.message,
&anyhow::Error::msg(initial_error),
);
}
}
Err(anyhow::Error::msg(error))
},
)
}
pub(crate) fn parse_queue_send_command(input: &str) -> Option<Result<usize, String>> {
let rest = strip_queue_command_prefix(input.trim())?;
let mut parts = rest.split_whitespace();
let action = parts.next()?;
if !action.eq_ignore_ascii_case("send") && !action.eq_ignore_ascii_case("now") {
return None;
}
let Some(raw_index) = parts.next() else {
return Some(Err("Usage: /queue send <n>".to_string()));
};
if parts.next().is_some() {
return Some(Err("Usage: /queue send <n>".to_string()));
}
let Ok(index) = raw_index.parse::<usize>() else {
return Some(Err("Index must be a positive number".to_string()));
};
if index == 0 {
return Some(Err("Index must be >= 1".to_string()));
}
Some(Ok(index - 1))
}
pub(crate) fn strip_queue_command_prefix(input: &str) -> Option<&str> {
for prefix in ["/queue", "/queued"] {
if let Some(rest) = input.strip_prefix(prefix)
&& (rest.is_empty() || rest.chars().next().is_some_and(char::is_whitespace))
{
return Some(rest);
}
}
None
}
pub(crate) async fn steer_user_message(
app: &mut App,
config: &Config,
engine_handle: &EngineHandle,
mut message: QueuedMessage,
) -> Result<bool> {
let stop_words = config.stop_words();
if is_stop_word(&message.display, &stop_words).is_some() {
engine_handle.cancel();
app.stopped_turn = true;
app.status_message = Some("Turn stopped. Tool calls blocked for this turn.".to_string());
return Ok(false);
}
app.stopped_turn = false;
if app
.hooks
.has_hooks_for_event(crate::hooks::HookEvent::MessageSubmit)
{
let context = app.base_hook_context().with_message(&message.display);
let strict_gates = app
.hooks
.matched_strict_gate_labels(crate::hooks::HookEvent::MessageSubmit, &context);
let hooks = app.hooks.clone();
let original_text = message.display.clone();
let outcome = match tokio::task::spawn_blocking(move || {
hooks.execute_message_submit_transform_for_dispatch(&context, &original_text)
})
.await
{
Ok(outcome) => outcome,
Err(error) => {
tracing::error!(target: "hooks", %error, "steer message_submit executor task was lost");
lost_message_submit_outcome(&strict_gates)
}
};
if !apply_message_submit_outcome(app, &mut message, outcome) {
return Ok(false);
}
}
let paused_snapshot = snapshot_steer_paused_state(app);
let paused_dispatch = plan_paused_command_message(app, &message.display);
let paused_note = paused_dispatch.note().map(str::to_string);
paused_dispatch.apply(app, engine_handle);
let cwd = std::env::current_dir().ok();
let mut git_cache = crate::tui::git_mention::GitMentionCache::default();
let completion_index = app.composer.mention_discovery.fuzzy_candidates(
&app.workspace,
&app.composer.mention_cwd,
app.mention_walk_depth,
app.workspace_follow_symlinks,
);
let references = crate::tui::file_mention::context_references_from_input_cached(
&message.display,
&app.workspace,
cwd.clone(),
&mut git_cache,
completion_index,
);
let mut content = queued_message_content_for_app(app, &message, cwd, &mut git_cache)?;
if let Some(note) = paused_note.as_deref() {
content.push_str(note);
}
let message_index = app.api_messages.len();
if active_foreground_shell_running(app)
&& let Err(err) = request_active_foreground_shell_background(app)
{
restore_steer_paused_state(app, &paused_snapshot);
engine_handle.set_paused(paused_snapshot.paused);
return Err(err.context("could not move foreground shell to /jobs before steering"));
}
if let Err(err) = engine_handle.steer(content.clone()).await {
restore_steer_paused_state(app, &paused_snapshot);
engine_handle.set_paused(paused_snapshot.paused);
return Err(err);
}
app.last_submitted_prompt = Some(message.display.clone());
app.flush_active_cell();
app.add_message(HistoryCell::User {
content: format!("+ {}", message.display),
});
let history_cell = app.history.len().saturating_sub(1);
app.record_context_references(history_cell, message_index, references);
app.api_messages.push(Message {
role: "user".to_string(),
content: vec![ContentBlock::Text {
text: content.clone(),
cache_control: None,
}],
});
app.status_message = Some("Steering current turn...".to_string());
Ok(true)
}
pub(crate) fn snapshot_steer_paused_state(app: &App) -> SteerPausedSnapshot {
SteerPausedSnapshot {
paused: app.paused,
pausable: app.pausable,
paused_quarry: app.paused_quarry.clone(),
quarry: app.hunt.quarry.clone(),
tokens_used: app.hunt.tokens_used,
time_used_seconds: app.hunt.time_used_seconds,
continuation_count: app.hunt.continuation_count,
}
}
pub(crate) fn restore_steer_paused_state(app: &mut App, snapshot: &SteerPausedSnapshot) {
app.paused = snapshot.paused;
app.pausable = snapshot.pausable;
app.paused_quarry = snapshot.paused_quarry.clone();
app.hunt.quarry = snapshot.quarry.clone();
app.hunt.tokens_used = snapshot.tokens_used;
app.hunt.time_used_seconds = snapshot.time_used_seconds;
app.hunt.continuation_count = snapshot.continuation_count;
}
pub(crate) async fn attempt_steer_with_queue_fallback(
app: &mut App,
config: &Config,
engine_handle: &EngineHandle,
message: QueuedMessage,
recovery: DispatchRecovery,
) {
match steer_user_message(app, config, engine_handle, message.clone()).await {
Ok(true) => {
app.push_status_toast(
"Steering into current turn",
StatusToastLevel::Info,
Some(1_500),
);
}
Ok(false) => {
restore_queued_or_draft_message(app, recovery, message);
app.push_status_toast(
"message_submit hook blocked the steer; original queue/draft restored",
StatusToastLevel::Warning,
Some(4_000),
);
}
Err(err) => {
restore_queued_or_draft_message(app, recovery, message);
let status = format!(
"Steer failed ({err}); {} queued follow-up(s) — /queue send <n>",
app.queued_message_count()
);
app.status_message = Some(status.clone());
app.push_status_toast(status, StatusToastLevel::Warning, Some(4_000));
}
}
}
pub(crate) async fn queue_follow_up(app: &mut App, message: QueuedMessage) -> Result<()> {
let display = message.display.clone();
enqueue_offline_message(app, message);
let toast = if app.mode == AppMode::Operate {
format!(
"Queued task: {display} ({} total) — dispatches next while workers continue; ↑ to edit",
app.queued_message_count()
)
} else {
format!(
"Queued: {display} ({} total) — sends after current output; ↑ to edit",
app.queued_message_count()
)
};
app.status_message = Some(toast.clone());
app.push_status_toast(toast, StatusToastLevel::Info, Some(3_000));
Ok(())
}
pub(crate) async fn dispatch_composer_message(
app: &mut App,
config: &Config,
engine_handle: &EngineHandle,
message: QueuedMessage,
recovery: DispatchRecovery,
action: ComposerSubmitAction,
) -> Result<()> {
if app.remote_control.blocks_local_input() {
app.input = message.display;
app.cursor_position = app.input.chars().count();
let status =
"Web remote control owns prompts. Use /rc stop to return input to this terminal."
.to_string();
app.status_message = Some(status.clone());
app.push_status_toast(status, StatusToastLevel::Warning, Some(6_000));
return Ok(());
}
if let Some(focus) = app.agent_focus.as_ref() {
let agent_id = focus.agent_id.clone();
let label = focus.label.clone();
let text = message.display.clone();
crate::tui::agent_focus::echo_user_follow_up(app, &text);
let receipt = app
.tr(crate::localization::MessageId::AgentFocusFollowUpQueued)
.replace("{agent}", &label);
app.push_history_cell(crate::tui::history::HistoryCell::System { content: receipt });
if engine_handle
.send(crate::core::ops::Op::FollowUpSubAgent {
agent_id: agent_id.clone(),
text,
})
.await
.is_err()
{
let failed = app
.tr(crate::localization::MessageId::AgentFocusFollowUpFailed)
.replace("{agent}", &label)
.replace("{reason}", "engine unavailable");
app.status_message = Some(failed.clone());
app.push_status_toast(failed, StatusToastLevel::Warning, Some(5_000));
}
return Ok(());
}
let disposition = match action {
ComposerSubmitAction::Submit(disposition) => disposition,
ComposerSubmitAction::SendQueuedNow | ComposerSubmitAction::Noop => {
SubmitDisposition::Queue
}
};
match disposition {
SubmitDisposition::Immediate => {
let _ =
dispatch_user_message_with_recovery(app, config, engine_handle, message, recovery)
.await;
Ok(())
}
SubmitDisposition::Queue => {
let count = app.queued_message_count().saturating_add(1);
enqueue_offline_message(app, message);
let (status, toast) = if app.offline_mode {
(
format!("Offline: {count} queued follow-up(s) — ↑ edit last, /queue send <n>"),
format!("Offline: queued follow-up ({count} total)"),
)
} else if app.mode == AppMode::Operate {
(
format!(
"{count} queued task(s) — dispatches next while workers continue; ↑ edit last, /queue send <n>"
),
format!("Queued task ({count} total) — dispatches next"),
)
} else {
(
format!(
"{count} queued follow-up(s) — sends after current output; ↑ edit last, /queue send <n>"
),
format!("Queued follow-up ({count} total) — sends after current output"),
)
};
app.status_message = Some(status);
app.push_status_toast(toast, StatusToastLevel::Info, Some(3_000));
Ok(())
}
SubmitDisposition::Steer => {
attempt_steer_with_queue_fallback(app, config, engine_handle, message, recovery).await;
Ok(())
}
SubmitDisposition::QueueFollowUp => queue_follow_up(app, message).await,
}
}
#[cfg(test)]
pub(crate) async fn submit_or_steer_message(
app: &mut App,
config: &Config,
engine_handle: &EngineHandle,
message: QueuedMessage,
recovery: DispatchRecovery,
) -> Result<()> {
let action = ComposerSubmitAction::Submit(app.decide_submit_disposition());
dispatch_composer_message(app, config, engine_handle, message, recovery, action).await
}
pub(crate) fn merge_pending_steers(app: &mut App) -> Option<QueuedMessage> {
let drained = app.drain_pending_steers();
if drained.is_empty() {
return None;
}
if drained.len() == 1 {
return drained.into_iter().next();
}
let mut skill_instruction: Option<String> = None;
let mut skill_provenance = None;
let mut bodies: Vec<String> = Vec::with_capacity(drained.len());
for msg in drained {
if skill_instruction.is_none() {
skill_instruction = msg.skill_instruction;
skill_provenance = msg.skill_provenance;
}
bodies.push(msg.display);
}
Some(
QueuedMessage::new(bodies.join("\n\n"), skill_instruction)
.with_skill_provenance(skill_provenance),
)
}