返回 CodeWhale
handle.rs
根目录 / crates / tui / src / core / engine / handle.rs
1 //! Public `EngineHandle` methods.
2 //!
3 //! The struct itself lives next door in `engine.rs` because two
4 //! construction sites (`Engine::new` and the test-only
5 //! `mock_engine_handle`) need access to its private mpsc channels.
6 //! The method surface — `send`, `cancel*`, `is_cancelled`,
7 //! `approve_tool_call` / `deny_tool_call` / `retry_tool_with_policy`,
8 //! `submit_user_input` / `cancel_user_input`, and `steer` — moves here
9 //! so the agent loop's mailbox API is reviewable on its own.
10
11 use anyhow::Result;
12 use std::collections::VecDeque;
13 use std::sync::{Arc, Mutex as StdMutex};
14 use tokio::sync::{mpsc, oneshot};
15 use tokio_util::sync::CancellationToken;
16
17 use codewhale_config::AppMode;
18 use codewhale_execpolicy::ApprovalMode;
19
20 use super::approval::{ApprovalDecision, UserInputDecision};
21 use super::{
22 CancelReason, EngineHandle, LiveRuntimeAuthority, Op, RuntimePermissionAuthority,
23 UserInputResponse,
24 };
25 use crate::approval_log::ApprovalDecider;
26
27 #[derive(Clone)]
28 pub(super) struct TurnControl {
29 pub id: u64,
30 pub narrowing: super::host_profile::TurnNarrowing,
31 pub cancel: CancellationToken,
32 pub reason: Arc<StdMutex<Option<CancelReason>>>,
33 }
34
35 #[derive(Default)]
36 pub(super) struct TurnControls {
37 next_id: u64,
38 pub active: Option<TurnControl>,
39 pub pending: VecDeque<TurnControl>,
40 }
41
42 impl TurnControls {
43 pub fn fresh(&mut self) -> TurnControl {
44 self.next_id = self
45 .next_id
46 .checked_add(1)
47 .expect("turn control id exhausted");
48 TurnControl {
49 id: self.next_id,
50 narrowing: super::host_profile::TurnNarrowing::Inherit,
51 cancel: CancellationToken::new(),
52 reason: Arc::new(StdMutex::new(None)),
53 }
54 }
55
56 fn target(&self) -> Option<&TurnControl> {
57 self.active.as_ref().or_else(|| self.pending.front())
58 }
59 }
60
61 pub(super) struct TurnControlGuard {
62 pub controls: Arc<StdMutex<TurnControls>>,
63 pub id: u64,
64 pub narrowing: super::host_profile::TurnNarrowing,
65 }
66
67 impl Drop for TurnControlGuard {
68 fn drop(&mut self) {
69 let mut controls = self
70 .controls
71 .lock()
72 .unwrap_or_else(std::sync::PoisonError::into_inner);
73 if controls
74 .active
75 .as_ref()
76 .is_some_and(|active| active.id == self.id)
77 {
78 controls.active = None;
79 }
80 }
81 }
82
83 #[derive(Debug)]
84 pub(crate) struct SteerInput {
85 pub(super) turn_id: Option<u64>,
86 pub(super) replace_pending: bool,
87 pub(crate) content: String,
88 pub(super) outcome: Option<oneshot::Sender<SteerOutcome>>,
89 }
90
91 /// The engine's verdict on one steer. A steer whose turn had already moved
92 /// on is discarded by `next_turn_steer`; the verdict tells the sender which
93 /// happened, because "the channel accepted the text" is not "the model saw
94 /// it" (#6276).
95 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
96 pub(crate) enum SteerOutcome {
97 /// The steer's text was committed into the session record inside the
98 /// turn it was sent to; the model received it.
99 Accepted,
100 /// The turn had already moved on (or ended) before the steer reached a
101 /// commit boundary. The model never saw the text.
102 Dropped,
103 }
104
105 /// A steer the engine has taken ownership of. Committing it reports
106 /// [`SteerOutcome::Accepted`]; any other exit — interrupt, failure, early
107 /// return, silent drop of the pending queue — reports `Dropped` from `Drop`,
108 /// so no path can lose a verdict.
109 pub(crate) struct PendingSteer {
110 pub(crate) replace_pending: bool,
111 pub(crate) content: String,
112 outcome: Option<oneshot::Sender<SteerOutcome>>,
113 }
114
115 impl PendingSteer {
116 pub(crate) fn new(content: String, outcome: Option<oneshot::Sender<SteerOutcome>>) -> Self {
117 Self {
118 content,
119 outcome,
120 replace_pending: false,
121 }
122 }
123
124 /// Commit the steer into the turn's record: report `Accepted`, then hand
125 /// back the text. Consuming `self` without calling this reports
126 /// `Dropped` via `Drop`.
127 pub(crate) fn commit(mut self) -> String {
128 if let Some(outcome) = self.outcome.take() {
129 let _ = outcome.send(SteerOutcome::Accepted);
130 }
131 // `Drop` runs after this returns and finds `outcome` already taken,
132 // so the verdict stays exactly one `Accepted`.
133 std::mem::take(&mut self.content)
134 }
135 }
136
137 impl Drop for PendingSteer {
138 fn drop(&mut self) {
139 if let Some(outcome) = self.outcome.take() {
140 let _ = outcome.send(SteerOutcome::Dropped);
141 }
142 }
143 }
144
145 impl SteerInput {
146 /// Take ownership of this steer as an unsettled [`PendingSteer`].
147 ///
148 /// This is the only way to claim a steer off the channel. Whatever the
149 /// claimant then does — `commit()` or drop — settles it exactly once, so
150 /// there is one settlement mechanism rather than two (#6276).
151 pub(crate) fn into_pending(mut self) -> PendingSteer {
152 // Both fields are taken, so the `Drop` below finds nothing left to
153 // settle and the verdict travels with the `PendingSteer`.
154 let mut pending = PendingSteer::new(std::mem::take(&mut self.content), self.outcome.take());
155 pending.replace_pending = self.replace_pending;
156 pending
157 }
158 }
159
160 impl Drop for SteerInput {
161 fn drop(&mut self) {
162 if let Some(outcome) = self.outcome.take() {
163 let _ = outcome.send(SteerOutcome::Dropped);
164 }
165 }
166 }
167
168 impl std::ops::Deref for SteerInput {
169 type Target = str;
170 fn deref(&self) -> &str {
171 &self.content
172 }
173 }
174
175 impl std::fmt::Display for SteerInput {
176 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
177 self.content.fmt(f)
178 }
179 }
180
181 pub(crate) struct SteerPermit {
182 permit: mpsc::OwnedPermit<SteerInput>,
183 turn_id: Option<u64>,
184 }
185
186 impl SteerPermit {
187 pub(crate) fn send(self, content: String) {
188 self.permit.send(SteerInput {
189 turn_id: self.turn_id,
190 replace_pending: false,
191 content,
192 outcome: None,
193 });
194 }
195
196 /// Send a steer and receive the engine's verdict on it. The receiver
197 /// resolves to [`SteerOutcome::Accepted`] when the turn commits the text
198 /// into its record, [`SteerOutcome::Dropped`] when the turn moved on
199 /// first, and closes without a verdict only if the engine itself is gone
200 /// (#6276).
201 pub(crate) fn send_with_outcome(self, content: String) -> oneshot::Receiver<SteerOutcome> {
202 self.send_with_replacement_outcome(content, false)
203 }
204
205 pub(crate) fn send_replacing_with_outcome(
206 self,
207 content: String,
208 ) -> oneshot::Receiver<SteerOutcome> {
209 self.send_with_replacement_outcome(content, true)
210 }
211
212 fn send_with_replacement_outcome(
213 self,
214 content: String,
215 replace_pending: bool,
216 ) -> oneshot::Receiver<SteerOutcome> {
217 let (outcome_tx, outcome_rx) = oneshot::channel();
218 self.permit.send(SteerInput {
219 turn_id: self.turn_id,
220 replace_pending,
221 content,
222 outcome: Some(outcome_tx),
223 });
224 outcome_rx
225 }
226 }
227
228 impl EngineHandle {
229 /// The engine's turn-phase heartbeat (#6184). Hosts read it to tell a
230 /// bounded model wait from a wedged turn without inferring liveness from
231 /// the stream-chunk timeout.
232 #[must_use]
233 pub(crate) fn turn_heartbeat(&self) -> &Arc<super::turn_heartbeat::TurnHeartbeat> {
234 &self.turn_heartbeat
235 }
236
237 /// Called only while Runtime holds the idle turn admission claim. The
238 /// following SendMessage refreshes the existing prompt/config projection.
239 pub(crate) fn restore_runtime_goal(
240 &self,
241 goal: Option<&codewhale_protocol::ThreadGoal>,
242 ) -> Result<()> {
243 let mut state = self
244 .goal_state
245 .lock()
246 .map_err(|_| anyhow::anyhow!("goal state lock poisoned"))?;
247 let current = state.snapshot();
248 if current.goal_id.as_deref() != goal.map(|goal| goal.goal_id.as_str()) {
249 *state = goal.map_or_else(crate::tools::goal::GoalState::default, |goal| {
250 crate::tools::goal::GoalState::from_snapshot(
251 &crate::tools::goal::GoalSnapshot::from_thread_goal(goal),
252 )
253 });
254 }
255 Ok(())
256 }
257
258 /// Apply the host's latest durable goal control to the live continuation
259 /// gate. The mailbox drains only between turns, so it cannot carry stop
260 /// controls. A replacement parks the old goal until normal turn admission
261 /// restores the new revision; it never starts a second turn here.
262 pub(crate) fn sync_runtime_goal_control(
263 &self,
264 goal: Option<&codewhale_protocol::ThreadGoal>,
265 ) -> Result<()> {
266 let mut state = self
267 .goal_state
268 .lock()
269 .map_err(|_| anyhow::anyhow!("goal state lock poisoned"))?;
270 let current = state.snapshot();
271 let Some(goal) = goal else {
272 state.clear();
273 return Ok(());
274 };
275 if current.goal_id.as_deref() != Some(goal.goal_id.as_str()) {
276 if state.is_active() {
277 state.sync_from_host_status(
278 current.objective.as_deref(),
279 current.token_budget,
280 crate::tools::goal::GoalStatus::Paused,
281 );
282 }
283 } else {
284 let (status, _) =
285 crate::tools::goal::thread_goal_status_projection(goal.status.clone());
286 if status != crate::tools::goal::GoalStatus::Active {
287 state.sync_from_host_status(
288 current.objective.as_deref(),
289 current.token_budget,
290 status,
291 );
292 }
293 }
294 Ok(())
295 }
296
297 /// True when the caller must preflight a concrete provider client before
298 /// committing UI/runtime turn state. Test and embedding handles with an
299 /// injected model client return false because that client owns model I/O.
300 #[must_use]
301 pub(crate) fn client_preflight_required(&self) -> bool {
302 self.client_preflight_required
303 }
304
305 /// Send an operation to the engine
306 ///
307 /// This awaits channel capacity, and the engine drains `rx_op` only
308 /// between turns — so on the UI event loop an awaited send into a
309 /// saturated mailbox freezes input for the rest of the turn (#6150).
310 /// Input-path callers instead either `try_send` a droppable op (report
311 /// the rejection) or `try_reserve_owned` before committing UI state and
312 /// hand off with `send_reserved_op`. An awaited `send` remains correct
313 /// only where the operation is part of a committed, ordered transition
314 /// (session/provider reload) whose drop would desync engine and UI.
315 pub async fn send(&self, op: Op) -> Result<()> {
316 let authority = Self::change_mode_authority(&op);
317 let permit = self.tx_op.clone().reserve_owned().await?;
318 if let Some(authority) = authority {
319 self.publish_runtime_authority(authority);
320 }
321 self.send_reserved_op(permit, op);
322 Ok(())
323 }
324
325 /// Try to send an operation without blocking.
326 ///
327 /// Returns `Err` if the channel is full or closed. Use this for
328 /// non-critical, refresh-type ops (e.g. `Op::ListSubAgents`) that can
329 /// safely be dropped and re-requested on the next drain cycle.
330 pub fn try_send(&self, op: Op) -> Result<()> {
331 let authority = Self::change_mode_authority(&op);
332 let result = self.tx_op.clone().try_reserve_owned();
333 // A full channel already guarantees that the engine will wake and
334 // drain an operation. Publish the typed authority anyway: the drain
335 // applies pending authority before handling that queued operation, so
336 // a posture edit never blocks behind refresh traffic. A closed
337 // channel has no engine left to observe the update.
338 if !matches!(&result, Err(mpsc::error::TrySendError::Closed(_)))
339 && let Some(authority) = authority
340 {
341 self.publish_runtime_authority(authority);
342 }
343 // Keep the public error bound to the rejected operation. Callers use
344 // TrySendError<Op> to distinguish a retryable full mailbox from a
345 // stopped engine; reservation errors otherwise carry a Sender<Op>.
346 match result {
347 Ok(permit) => {
348 self.send_reserved_op(permit, op);
349 Ok(())
350 }
351 Err(mpsc::error::TrySendError::Full(_)) => {
352 Err(mpsc::error::TrySendError::Full(op).into())
353 }
354 Err(mpsc::error::TrySendError::Closed(_)) => {
355 Err(mpsc::error::TrySendError::Closed(op).into())
356 }
357 }
358 }
359
360 /// Bind controls and enqueue under one lock, preserving the same FIFO as
361 /// the operation mailbox even when several senders hold reserved slots.
362 pub(crate) fn send_reserved_op(&self, permit: mpsc::OwnedPermit<Op>, op: Op) {
363 self.send_reserved_narrowed_op(permit, op, super::host_profile::TurnNarrowing::Inherit);
364 }
365
366 /// Called only by the captured ACP Runtime admission. One lock retains the
367 /// exact control/mailbox FIFO; no profile is added to the public operation.
368 pub(crate) fn send_reserved_acp_op(&self, permit: mpsc::OwnedPermit<Op>, op: Op) {
369 self.send_reserved_narrowed_op(permit, op, super::host_profile::TurnNarrowing::Acp);
370 }
371
372 fn send_reserved_narrowed_op(
373 &self,
374 permit: mpsc::OwnedPermit<Op>,
375 op: Op,
376 narrowing: super::host_profile::TurnNarrowing,
377 ) {
378 let mut controls = self
379 .turn_controls
380 .lock()
381 .unwrap_or_else(std::sync::PoisonError::into_inner);
382 if matches!(&op, Op::SendMessage(_)) {
383 let mut control = controls.fresh();
384 control.narrowing = narrowing;
385 controls.pending.push_back(control);
386 }
387 permit.send(op);
388 }
389
390 fn change_mode_authority(op: &Op) -> Option<LiveRuntimeAuthority> {
391 let Op::ChangeMode {
392 mode,
393 allow_shell,
394 trust_mode,
395 auto_approve,
396 approval_mode,
397 configured_sandbox_mode,
398 } = op
399 else {
400 return None;
401 };
402 Some(LiveRuntimeAuthority::from_fields(
403 *mode,
404 *allow_shell,
405 *trust_mode,
406 *auto_approve,
407 *approval_mode,
408 configured_sandbox_mode.clone(),
409 ))
410 }
411
412 fn publish_runtime_authority(&self, authority: LiveRuntimeAuthority) {
413 let mut state = self
414 .live_runtime_authority
415 .lock()
416 .unwrap_or_else(std::sync::PoisonError::into_inner);
417 state.revision = state.revision.wrapping_add(1).max(1);
418 state.authority = authority;
419 }
420
421 pub(crate) fn publish_turn_authority(
422 &self,
423 mode: AppMode,
424 allow_shell: bool,
425 trust_mode: bool,
426 auto_approve: bool,
427 approval_mode: ApprovalMode,
428 configured_sandbox_mode: Option<String>,
429 ) {
430 self.publish_runtime_authority(LiveRuntimeAuthority::from_fields(
431 mode,
432 allow_shell,
433 trust_mode,
434 auto_approve,
435 approval_mode,
436 configured_sandbox_mode,
437 ));
438 }
439
440 /// Exact live permission authority for runtime approval and elevation
441 /// gates. This is the same typed state the active engine turn drains.
442 #[must_use]
443 pub(crate) fn runtime_permission_authority(&self) -> RuntimePermissionAuthority {
444 self.live_runtime_authority
445 .lock()
446 .unwrap_or_else(std::sync::PoisonError::into_inner)
447 .authority
448 .permission_snapshot()
449 }
450
451 /// Reserve capacity for a runtime steer before it mutates durable state.
452 /// The owned permit lets the caller persist and dispatch synchronously,
453 /// without a cancellation point between those two operations.
454 pub(crate) async fn reserve_steer(&self) -> Result<SteerPermit> {
455 let permit = self.tx_steer.clone().reserve_owned().await?;
456 let turn_id = self
457 .turn_controls
458 .lock()
459 .unwrap_or_else(std::sync::PoisonError::into_inner)
460 .target()
461 .map(|control| control.id);
462 Ok(SteerPermit { permit, turn_id })
463 }
464
465 /// Cancel the current request (user-initiated path — keeps the
466 /// public `cancel()` signature stable). Equivalent to
467 /// `cancel_with_reason(CancelReason::User)`.
468 pub fn cancel(&self) {
469 self.cancel_with_reason(CancelReason::User);
470 }
471
472 /// Cancel the current request and latch the reason so downstream
473 /// "request cancelled" error messages can name a cause.
474 pub fn cancel_with_reason(&self, reason: CancelReason) {
475 // Keep turn activation excluded until both the admitted control and
476 // the legacy shared token have been canceled.
477 let controls = self
478 .turn_controls
479 .lock()
480 .unwrap_or_else(std::sync::PoisonError::into_inner);
481 if let Some(control) = controls.target() {
482 *control
483 .reason
484 .lock()
485 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(reason);
486 control.cancel.cancel();
487 }
488 match self.cancel_reason.lock() {
489 Ok(mut slot) => *slot = Some(reason),
490 Err(poisoned) => *poisoned.into_inner() = Some(reason),
491 }
492 match self.cancel_token.lock() {
493 Ok(token) => token.cancel(),
494 Err(poisoned) => poisoned.into_inner().cancel(),
495 }
496 crate::retry_status::clear();
497 }
498
499 /// Snapshot the exact existing target control for one captured transport
500 /// wait. It cannot follow a later turn or reset the admitted token.
501 pub(crate) fn captured_turn_cancel(&self) -> CancellationToken {
502 self.turn_controls
503 .lock()
504 .unwrap_or_else(std::sync::PoisonError::into_inner)
505 .target()
506 .map(|control| control.cancel.clone())
507 .unwrap_or_else(|| {
508 self.cancel_token
509 .lock()
510 .unwrap_or_else(std::sync::PoisonError::into_inner)
511 .clone()
512 })
513 }
514
515 /// Check if a request is currently cancelled
516 #[must_use]
517 pub fn is_cancelled(&self) -> bool {
518 if let Some(control) = self
519 .turn_controls
520 .lock()
521 .unwrap_or_else(std::sync::PoisonError::into_inner)
522 .target()
523 {
524 return control.cancel.is_cancelled();
525 }
526 match self.cancel_token.lock() {
527 Ok(token) => token.is_cancelled(),
528 Err(poisoned) => poisoned.into_inner().is_cancelled(),
529 }
530 }
531
532 /// Pause or resume the current pausable command.
533 pub fn set_paused(&self, paused: bool) {
534 match self.shared_paused.lock() {
535 Ok(mut slot) => *slot = paused,
536 Err(poisoned) => *poisoned.into_inner() = paused,
537 }
538 }
539
540 /// Check whether the engine pause gate is set.
541 #[cfg(test)]
542 #[must_use]
543 pub fn is_paused(&self) -> bool {
544 match self.shared_paused.lock() {
545 Ok(slot) => *slot,
546 Err(poisoned) => *poisoned.into_inner(),
547 }
548 }
549
550 /// Deliver one approval answer. An agent's pending call gets it directly
551 /// — the engine may be streaming or running tools and not reading
552 /// approvals — and everything else goes to the engine's own waiting call.
553 async fn send_approval(&self, decision: ApprovalDecision) -> Result<()> {
554 use crate::tools::subagent::ChildApprovalOutcome;
555 let child = match &decision {
556 ApprovalDecision::Approved { id, .. } => Some((id, ChildApprovalOutcome::Approved)),
557 ApprovalDecision::Denied { id, .. } => Some((id, ChildApprovalOutcome::Denied)),
558 // An agent has no timeout outcome of its own (#6101): an expired
559 // card is a deny for whichever call it was answering.
560 ApprovalDecision::TimedOut { id } => Some((id, ChildApprovalOutcome::Denied)),
561 ApprovalDecision::Unavailable { id } => Some((id, ChildApprovalOutcome::Unavailable)),
562 // A sandbox retry only exists for the parent's own tool call.
563 ApprovalDecision::RetryWithPolicy { .. } => None,
564 };
565 if let Some((id, outcome)) = child
566 && crate::tools::subagent::SubAgentManager::is_child_approval_id(id)
567 && self
568 .subagent_manager
569 .write()
570 .await
571 .resolve_child_approval(id, outcome)
572 {
573 return Ok(());
574 }
575 self.tx_approval.send(decision).await?;
576 Ok(())
577 }
578
579 /// Approve a pending tool call because a person said yes.
580 pub async fn approve_tool_call(&self, id: impl Into<String>) -> Result<()> {
581 self.approve_tool_call_by(id, ApprovalDecider::User).await
582 }
583
584 /// Approve a pending tool call, recording who answered: a person, a
585 /// session rule, or the active posture. The approval receipt keeps it.
586 pub async fn approve_tool_call_by(
587 &self,
588 id: impl Into<String>,
589 by: ApprovalDecider,
590 ) -> Result<()> {
591 self.send_approval(ApprovalDecision::Approved { id: id.into(), by })
592 .await
593 }
594
595 /// Deny a pending tool call because a person said no.
596 pub async fn deny_tool_call(&self, id: impl Into<String>) -> Result<()> {
597 self.deny_tool_call_by(id, ApprovalDecider::User).await
598 }
599
600 /// Deny a pending tool call, recording who answered.
601 pub async fn deny_tool_call_by(
602 &self,
603 id: impl Into<String>,
604 by: ApprovalDecider,
605 ) -> Result<()> {
606 self.send_approval(ApprovalDecision::Denied { id: id.into(), by })
607 .await
608 }
609
610 /// Deny a pending tool call because its interactive approval card
611 /// expired (#6101). Kept distinct from [`Self::deny_tool_call`] so the
612 /// receipt records a timeout instead of an operator denial.
613 pub async fn deny_tool_call_timed_out(&self, id: impl Into<String>) -> Result<()> {
614 self.send_approval(ApprovalDecision::TimedOut { id: id.into() })
615 .await
616 }
617
618 /// Resolve a pending tool call whose request could not be put in front
619 /// of a person (a stale turn's request, or a child from another
620 /// conversation). Kept distinct from [`Self::deny_tool_call`] so neither
621 /// the receipt nor the model message claims the person denied it.
622 pub async fn deny_tool_call_unavailable(&self, id: impl Into<String>) -> Result<()> {
623 self.send_approval(ApprovalDecision::Unavailable { id: id.into() })
624 .await
625 }
626
627 /// Retry a tool call with an elevated sandbox policy a person chose.
628 pub async fn retry_tool_with_policy(
629 &self,
630 id: impl Into<String>,
631 policy: crate::sandbox::SandboxPolicy,
632 ) -> Result<()> {
633 self.retry_tool_with_policy_by(id, policy, ApprovalDecider::User)
634 .await
635 }
636
637 /// Retry a tool call with an elevated sandbox policy, recording who chose it.
638 pub async fn retry_tool_with_policy_by(
639 &self,
640 id: impl Into<String>,
641 policy: crate::sandbox::SandboxPolicy,
642 by: ApprovalDecider,
643 ) -> Result<()> {
644 self.tx_approval
645 .send(ApprovalDecision::RetryWithPolicy {
646 id: id.into(),
647 policy,
648 by,
649 })
650 .await?;
651 Ok(())
652 }
653
654 /// Submit a response for request_user_input.
655 pub async fn submit_user_input(
656 &self,
657 id: impl Into<String>,
658 response: UserInputResponse,
659 ) -> Result<()> {
660 self.tx_user_input
661 .send(UserInputDecision::Submitted {
662 id: id.into(),
663 response,
664 })
665 .await?;
666 Ok(())
667 }
668
669 /// Cancel a request_user_input prompt.
670 pub async fn cancel_user_input(&self, id: impl Into<String>) -> Result<()> {
671 self.tx_user_input
672 .send(UserInputDecision::Cancelled { id: id.into() })
673 .await?;
674 Ok(())
675 }
676
677 /// Steer an in-flight turn with additional user input.
678 pub async fn steer(&self, content: impl Into<String>) -> Result<()> {
679 self.reserve_steer().await?.send(content.into());
680 Ok(())
681 }
682
683 /// Request the live context-window budget for this session's route.
684 /// `None` means the route cannot express a bounded window (e.g. an
685 /// unknown model with no catalog or configured limits) — callers should
686 /// surface "unavailable" rather than inventing a number.
687 pub async fn get_context_budget(
688 &self,
689 ) -> Result<Option<crate::core::ops::SessionContextBudget>> {
690 let (tx, rx) = tokio::sync::oneshot::channel();
691 let tx = std::sync::Arc::new(std::sync::Mutex::new(Some(tx)));
692 self.send(Op::GetContextBudget { tx }).await?;
693 rx.await
694 .map_err(|_| anyhow::anyhow!("Engine dropped context budget oneshot"))
695 }
696
697 /// Request a snapshot of the current session state.
698 /// Returns the snapshot directly via a oneshot channel, avoiding
699 /// competition with the SSE event stream on the mpsc receiver.
700 pub async fn get_session_snapshot(&self) -> Result<crate::core::ops::SessionSnapshot> {
701 let (tx, rx) = tokio::sync::oneshot::channel();
702 let tx = std::sync::Arc::new(std::sync::Mutex::new(Some(tx)));
703 self.send(Op::GetSessionSnapshot { tx }).await?;
704 rx.await
705 .map_err(|_| anyhow::anyhow!("Engine dropped session snapshot oneshot"))
706 }
707
708 /// Query after the active turn settles, without competing with events.
709 /// The caller must keep draining events and bound this future: an active
710 /// turn can be awaiting provider/tool work or a full event channel.
711 pub(crate) async fn get_subagent_settlement(
712 &self,
713 ) -> Result<crate::core::ops::SubAgentSettlement> {
714 let (tx, rx) = tokio::sync::oneshot::channel();
715 let tx = Arc::new(StdMutex::new(Some(tx)));
716 self.send(Op::GetSubAgentSettlement { tx }).await?;
717 rx.await
718 .map_err(|_| anyhow::anyhow!("Engine dropped child settlement receipt"))
719 }
720
721 /// Request active provider request concurrency state.
722 pub async fn get_provider_runtime_status(
723 &self,
724 ) -> Result<crate::core::ops::ProviderRuntimeStatus> {
725 let (tx, rx) = tokio::sync::oneshot::channel();
726 let tx = std::sync::Arc::new(std::sync::Mutex::new(Some(tx)));
727 self.send(Op::GetProviderRuntimeStatus { tx }).await?;
728 rx.await
729 .map_err(|_| anyhow::anyhow!("Engine dropped provider runtime status oneshot"))
730 }
731
732 /// Run the bounded initial connection pass on the engine-owned MCP pool.
733 ///
734 /// The returned manager snapshot and every later tool call therefore see
735 /// the same connections and catalog generation. Unlike `reload_mcp`, this
736 /// does not force a config re-read or drop ready transports. Optional
737 /// servers are connected in the background at engine spawn; this waits
738 /// only if the caller explicitly asked for the settled receipt.
739 pub async fn bootstrap_mcp(&self) -> Result<crate::core::ops::McpManagerUpdate> {
740 let (tx, rx) = tokio::sync::oneshot::channel();
741 let tx = std::sync::Arc::new(std::sync::Mutex::new(Some(tx)));
742 self.send(Op::BootstrapMcp { tx }).await?;
743 rx.await
744 .map_err(|_| anyhow::anyhow!("Engine dropped MCP bootstrap oneshot"))?
745 .map_err(anyhow::Error::msg)
746 }
747
748 /// Retry one failed server through the existing engine-owned pool.
749 pub async fn retry_mcp_server(
750 &self,
751 name: impl Into<String>,
752 ) -> Result<crate::core::ops::McpManagerUpdate> {
753 let (tx, rx) = tokio::sync::oneshot::channel();
754 let tx = std::sync::Arc::new(std::sync::Mutex::new(Some(tx)));
755 self.send(Op::RetryMcpServer {
756 name: name.into(),
757 tx,
758 })
759 .await?;
760 rx.await
761 .map_err(|_| anyhow::anyhow!("Engine dropped MCP retry oneshot"))?
762 .map_err(anyhow::Error::msg)
763 }
764
765 /// Force the engine-owned MCP pool to reload and reconnect, returning a
766 /// snapshot from the exact live pool that supplies the next model turn.
767 pub async fn reload_mcp(
768 &self,
769 config_path: std::path::PathBuf,
770 ) -> Result<crate::core::ops::McpManagerUpdate> {
771 let (tx, rx) = tokio::sync::oneshot::channel();
772 let tx = std::sync::Arc::new(std::sync::Mutex::new(Some(tx)));
773 self.send(Op::ReloadMcp { config_path, tx }).await?;
774 rx.await
775 .map_err(|_| anyhow::anyhow!("Engine dropped MCP reload oneshot"))?
776 .map_err(anyhow::Error::msg)
777 }
778 }
779
779 lines RUST