返回 CodeWhale
rlm_host.rs
根目录 / crates / tui / src / core / engine / rlm_host.rs
1 //! Captured, call-local RLM projection onto the canonical Engine.
2 //!
3 //! A Python kernel never owns this receipt or a bridge. The only model/code
4 //! driver is Engine::run_turn; the RPC adapter owns its bounded invocation.
5 use super::turn_loop::usage_has_reported_data;
6 use super::*;
7 use crate::rlm::bridge::{RlmUsageAccumulator, RlmUsageReservation};
8 use crate::rlm::turn::{RlmRoundTrace, RlmTermination, RlmTurnResult};
9 use crate::tui::auto_review::AutoReviewPolicy;
10 use anyhow::anyhow;
11
12 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
13 pub(crate) enum RlmMode {
14 Completion,
15 Recursive { depth_remaining: u32 },
16 }
17
18 pub(crate) struct RlmInvocation {
19 pub prompt: String,
20 pub mode: RlmMode,
21 pub max_tokens: Option<u32>,
22 pub task_instructions: Option<String>,
23 pub deadline: tokio::time::Instant,
24 pub gate: Option<crate::tools::codemode::NestedCallGate>,
25 pub events: Option<mpsc::Sender<Event>>,
26 pub usage: RlmUsageAccumulator,
27 }
28
29 /// Frozen from the actual serving Core turn, never from a client/string alone.
30 /// `context` predates attaching this receipt, so it cannot contain itself.
31 #[derive(Clone)]
32 pub(crate) struct CapturedRlmCaller {
33 #[cfg(test)]
34 _fixture_workspace: Option<Arc<tempfile::TempDir>>,
35 context: ToolContext,
36 authority: TurnAuthority,
37 route: ValidatedRuntimeRoute,
38 client: SharedModelClient,
39 subagent_manager: SharedSubAgentManager,
40 injected: bool,
41 config: EngineConfig,
42 system: SystemPrompt,
43 extension_prompt_block: Option<String>,
44 approval_store: Result<ApprovalReceiptStore, String>,
45 review_policy: Arc<AutoReviewPolicy>,
46 deadline: tokio::time::Instant,
47 origin_scope: crate::cost_status::CostScopeToken,
48 origin_session: String,
49 origin_turn: String,
50 child_accounting: Option<crate::tools::subagent::engine::ChildAccountingProjection>,
51 scheduler: tokio::runtime::Handle,
52 runtime_owner: Option<String>,
53 _runtime_lease: Option<crate::cost_status::RuntimeUsageLease>,
54 }
55 impl std::fmt::Debug for CapturedRlmCaller {
56 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
57 f.debug_struct("CapturedRlmCaller")
58 .field("workspace", &self.context.workspace)
59 .field("session", &self.context.state_namespace)
60 .finish_non_exhaustive()
61 }
62 }
63 impl CapturedRlmCaller {
64 pub(super) fn capture(
65 engine: &Engine,
66 authority: &TurnAuthority,
67 route: &TurnRouteContext,
68 context: &ToolContext,
69 origin_turn: &str,
70 ) -> Result<Self, ToolError> {
71 if context.acp_host.is_some() || origin_turn.trim().is_empty() {
72 return Err(ToolError::not_available(
73 "this host has no RLM model-call authority",
74 ));
75 }
76 let identity = engine
77 .api_provider_identity
78 .as_ref()
79 .ok_or_else(|| ToolError::not_available("the Core route has no admitted identity"))?;
80 route
81 .api_config
82 .verify_provider_identity(identity)
83 .map_err(ToolError::permission_denied)?;
84 let resolved =
85 resolve_runtime_route_for_identity(&route.api_config, identity, Some(&route.model))
86 .map_err(ToolError::permission_denied)?;
87 let client = engine
88 .codewhale_client
89 .as_ref()
90 .ok_or_else(|| ToolError::not_available("the Core route has no captured client"))?;
91 if client.admitted_provider_identity() != identity
92 || engine
93 .active_route_endpoint
94 .as_ref()
95 .map(Engine::endpoint_identity)
96 != Some(Engine::endpoint_identity(resolved.candidate.endpoint()))
97 {
98 return Err(ToolError::permission_denied(
99 "the RLM route differs from the installed Core client",
100 ));
101 }
102 let mut captured_context = context.clone();
103 captured_context.rlm_caller = None;
104 let child_accounting = engine
105 .child_host
106 .as_ref()
107 .map(|child| child.authority.accounting_projection());
108 let scheduler = tokio::runtime::Handle::try_current()
109 .map_err(|_| ToolError::not_available("RLM caller has no held Engine scheduler"))?;
110 let runtime_owner = engine.config.compaction.runtime_cost_owner.clone();
111 let runtime_lease = runtime_owner
112 .as_deref()
113 .and_then(crate::cost_status::acquire_runtime_usage_lease);
114 let parent_deadline = context.turn_deadline;
115 let default_deadline = tokio::time::Instant::now()
116 + engine
117 .config
118 .turn_wall_clock
119 .min(crate::tools::subagent::DEFAULT_CHILD_WALL_TIME);
120 let deadline = parent_deadline.map_or(default_deadline, |d| d.min(default_deadline));
121 let result =
122 Self {
123 #[cfg(test)]
124 _fixture_workspace: None,
125 context: captured_context,
126 authority: authority.clone(),
127 route: ValidatedRuntimeRoute {
128 identity: resolved.identity,
129 candidate: resolved.candidate,
130 config: resolved.config,
131 model: resolved.model,
132 context_window: resolved.context_window,
133 client: client.clone(),
134 },
135 client: engine.model_client.clone().ok_or_else(|| {
136 ToolError::not_available("the Core model client is unavailable")
137 })?,
138 subagent_manager: Arc::clone(&engine.subagent_manager),
139 injected: engine.model_client_injected,
140 config: engine.config.clone(),
141 system: engine.session.system_prompt.clone().ok_or_else(|| {
142 ToolError::not_available("the Core policy prompt is unavailable")
143 })?,
144 extension_prompt_block: engine.extension_prompt_block.clone(),
145 approval_store: engine.approval_receipt_store.clone(),
146 review_policy: Arc::clone(&engine.shared_auto_review_policy),
147 deadline,
148 origin_scope: crate::cost_status::scope_token(),
149 origin_session: context.state_namespace.clone(),
150 origin_turn: origin_turn.to_string(),
151 child_accounting,
152 scheduler,
153 runtime_owner,
154 _runtime_lease: runtime_lease,
155 };
156 result.validate_context(context)?;
157 Ok(result)
158 }
159
160 pub(crate) fn validate_context(&self, context: &ToolContext) -> Result<(), ToolError> {
161 if context.workspace != self.context.workspace
162 || context.state_namespace != self.context.state_namespace
163 || context.owner_agent_id != self.context.owner_agent_id
164 || context
165 .plugin_registry
166 .as_ref()
167 .map(|p| p.caller_selection())
168 != self
169 .context
170 .plugin_registry
171 .as_ref()
172 .map(|p| p.caller_selection())
173 {
174 return Err(ToolError::permission_denied(
175 "RLM caller identity, workspace or composition changed",
176 ));
177 }
178 self.validate_live()
179 }
180
181 pub(super) fn validate_live(&self) -> Result<(), ToolError> {
182 if self
183 .context
184 .cancel_token
185 .as_ref()
186 .is_none_or(CancellationToken::is_cancelled)
187 {
188 return Err(ToolError::cancelled(
189 "RLM originating turn is cancelled or unavailable",
190 ));
191 }
192 if tokio::time::Instant::now() >= self.deadline {
193 return Err(ToolError::execution_failed(
194 "RLM original parent deadline exhausted",
195 ));
196 }
197 crate::extension_host::validate_caller_plugins(self.context.plugin_registry.as_deref())
198 .map_err(ToolError::permission_denied)?;
199 self.route
200 .config
201 .verify_provider_identity(&self.route.identity)
202 .map_err(ToolError::permission_denied)
203 }
204
205 pub(crate) fn deadline(&self) -> tokio::time::Instant {
206 self.deadline
207 }
208
209 /// Boxing breaks recursive future types; the future still borrows this
210 /// caller and never enters the persistent session map or PythonRuntime.
211 pub(crate) fn dispatch<'call>(
212 &'call self,
213 invocation: RlmInvocation,
214 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = RlmTurnResult> + Send + 'call>> {
215 Box::pin(self.dispatch_admitted(invocation))
216 }
217
218 async fn dispatch_admitted(&self, mut invocation: RlmInvocation) -> RlmTurnResult {
219 let started = Instant::now();
220 invocation.deadline = invocation.deadline.min(self.deadline);
221 if let Err(error) = self.validate_live() {
222 return RlmTurnResult::failed(error.to_string(), started.elapsed());
223 }
224 if tokio::time::Instant::now() >= invocation.deadline {
225 return RlmTurnResult::failed(
226 "RLM original parent deadline exhausted".into(),
227 started.elapsed(),
228 );
229 }
230 if invocation.mode != RlmMode::Completion && invocation.gate.is_none() {
231 return RlmTurnResult::failed(
232 "no permission gate is serving this RLM turn".into(),
233 started.elapsed(),
234 );
235 }
236 // Every admitted completion/recursive host owns bounded receipt
237 // capacity, including setup/preparation failures with no provider call.
238 if let Err(error) = invocation.usage.reserve_nested_turn().await {
239 return RlmTurnResult::failed(error, started.elapsed());
240 }
241 let usage = invocation.usage.clone();
242 let events = invocation.events.clone();
243 let depth = match invocation.mode {
244 RlmMode::Completion => 0,
245 RlmMode::Recursive { depth_remaining } => depth_remaining,
246 };
247 let run_id = uuid::Uuid::new_v4().simple().to_string();
248 let (mut engine, handle) =
249 match Engine::new_rlm_admitted(Arc::new(self.clone()), invocation).await {
250 Ok(host) => host,
251 Err(error) => return RlmTurnResult::failed(error.to_string(), started.elapsed()),
252 };
253 let spec = match engine.rlm_turn_spec() {
254 Ok(spec) => spec,
255 Err(error) => return RlmTurnResult::failed(error.to_string(), started.elapsed()),
256 };
257 let cancel = self
258 .context
259 .cancel_token
260 .as_ref()
261 .expect("validated caller cancel")
262 .clone();
263 let deadline = engine.rlm_host.as_ref().expect("RLM host").deadline;
264 engine.rlm_host.as_mut().expect("RLM host").run_id = run_id;
265 let mut receiver = handle.rx_event.write().await;
266 let local_cancel = engine.cancel_token.clone();
267 // Drive the existing admission/turn method while draining its existing
268 // bounded event channel. No authority-bearing task is detached.
269 let outcome = {
270 let mut call = Box::pin(engine.handle_send_message(spec));
271 loop {
272 tokio::select! {
273 biased;
274 () = cancel.cancelled() => { local_cancel.cancel(); break None; }
275 () = tokio::time::sleep_until(deadline) => { local_cancel.cancel(); break None; }
276 result = &mut call => break Some(result),
277 event = receiver.recv() => {
278 let Some(event) = event else { local_cancel.cancel(); break None; };
279 forward_rlm_event(events.as_ref(), event, depth, &cancel, deadline).await;
280 }
281 }
282 }
283 };
284 while let Ok(event) = receiver.try_recv() {
285 forward_rlm_event(events.as_ref(), event, depth, &cancel, deadline).await;
286 }
287 let mut result = engine.rlm_result(outcome, started.elapsed());
288 let terminal = format!(
289 "RLM finished: {:?} after {} iteration(s), {} sub-LLM call(s), answer {} chars{}",
290 result.termination,
291 result.iterations,
292 result.total_rpcs,
293 result.answer.chars().count(),
294 result
295 .error
296 .as_ref()
297 .map_or(String::new(), |error| format!(" — {error}"))
298 );
299 usage.record_nested_event(serde_json::json!({"run_id": engine.rlm_host.as_ref().expect("RLM host").run_id, "depth_remaining": depth, "kind":"status", "content":terminal})).await;
300 // Settlement is already in the canonical nested-event ledger. The
301 // terminal observation must not wait past cancellation/deadline, and
302 // an available channel can still observe it after the deadline.
303 if let Some(events) = events.as_ref()
304 && let Some(message) =
305 crate::rlm::bridge::nested_rlm_status_line(Event::status(terminal), depth)
306 {
307 let _ = events.try_send(Event::status(message));
308 }
309 let snapshot = usage.snapshot().await;
310 result.usage = snapshot.usage;
311 result.routed_usage = snapshot.records;
312 result.routed_usage_drop_records = snapshot.drop_records;
313 result.routed_usage_dropped_records = snapshot.dropped_records;
314 result
315 }
316 }
317
318 pub(super) struct RlmHostSetup {
319 pub caller: Arc<CapturedRlmCaller>,
320 pub invocation: RlmInvocation,
321 pub system_prompt: SystemPrompt,
322 }
323
324 pub(super) struct RlmHostState {
325 pub caller: Arc<CapturedRlmCaller>,
326 pub run_id: String,
327 pub mode: RlmMode,
328 pub system_prompt: SystemPrompt,
329 pub prompt: String,
330 pub deadline: tokio::time::Instant,
331 pub gate: Option<crate::tools::codemode::NestedCallGate>,
332 pub usage: RlmUsageAccumulator,
333 pub max_tokens: Option<u32>,
334 pub trace: Vec<RlmRoundTrace>,
335 pub last_response: String,
336 pub final_answer: Option<String>,
337 pub total_rpcs: u32,
338 pub consecutive_no_code: u32,
339 pub termination: Option<RlmTermination>,
340 pub error: Option<String>,
341 pub model_rounds: u32,
342 }
343 impl From<RlmHostSetup> for RlmHostState {
344 fn from(setup: RlmHostSetup) -> Self {
345 Self {
346 caller: setup.caller,
347 run_id: String::new(),
348 mode: setup.invocation.mode,
349 system_prompt: setup.system_prompt,
350 prompt: setup.invocation.prompt,
351 deadline: setup.invocation.deadline,
352 gate: setup.invocation.gate,
353 usage: setup.invocation.usage,
354 max_tokens: setup.invocation.max_tokens,
355 trace: Vec::new(),
356 last_response: String::new(),
357 final_answer: None,
358 total_rpcs: 0,
359 consecutive_no_code: 0,
360 termination: None,
361 error: None,
362 model_rounds: 0,
363 }
364 }
365 }
366
367 impl Engine {
368 pub(super) async fn new_rlm_admitted(
369 caller: Arc<CapturedRlmCaller>,
370 invocation: RlmInvocation,
371 ) -> Result<(Self, EngineHandle)> {
372 caller.validate_live().map_err(anyhow::Error::new)?;
373 let mut config = caller.config.clone();
374 config.model = caller.route.model.clone();
375 config.workspace = caller.context.workspace.clone();
376 config.session_id = Some(caller.context.state_namespace.clone());
377 config.plugin_registry = caller.context.plugin_registry.clone();
378 config.runtime_services = caller.context.runtime.clone();
379 config.allow_shell = caller.authority.allow_shell;
380 config.trust_mode = caller.authority.trust_mode;
381 config.network_policy = caller.context.network_policy.clone();
382 config.max_steps = if invocation.mode == RlmMode::Completion {
383 1
384 } else {
385 crate::rlm::turn::MAX_RLM_ITERATIONS
386 };
387 config.snapshots_enabled = false;
388 config.memory_enabled = false;
389 config.terminal_chrome_enabled = false;
390 config.subagents_enabled = false;
391 config.advisor_config = crate::tools::subagent::AdvisorConfig::disabled();
392 config.goal_objective = None;
393 config.goal_status = GoalStatus::Paused;
394 config.goal_state = new_shared_goal_state();
395 config.turn_wall_clock = invocation
396 .deadline
397 .saturating_duration_since(tokio::time::Instant::now());
398 config.compaction.runtime_cost_owner = caller.runtime_owner.clone();
399 let system_prompt = crate::rlm::prompt::captured_rlm_prompt(
400 &caller.system,
401 invocation.mode,
402 invocation.task_instructions.as_deref(),
403 )?;
404 let mut kernel = None;
405 if invocation.mode != RlmMode::Completion {
406 let path = tempfile::TempPath::try_from_path(crate::rlm::session::write_context_file(
407 &invocation.prompt,
408 )?)?;
409 let cancel = caller
410 .context
411 .cancel_token
412 .as_ref()
413 .expect("validated originating cancellation");
414 let startup = tokio::select! {
415 biased;
416 () = cancel.cancelled() => Err("RLM Python startup cancelled by its originating turn".to_string()),
417 result = tokio::time::timeout_at(invocation.deadline, crate::repl::PythonRuntime::spawn_with_context(&path)) => {
418 result.map_err(|_| "RLM Python startup exhausted its original wall-clock deadline".to_string()).and_then(|r| r)
419 },
420 };
421 match startup {
422 Ok(runtime) => {
423 // The successfully bootstrapped runtime owns this same
424 // path. Before that handoff, TempPath cleans every Drop.
425 let _context_path = path.keep()?;
426 kernel = Some(runtime);
427 }
428 Err(error) => return Err(anyhow!("RLM Python startup failed: {error}")),
429 }
430 }
431 let api = (*caller.route.config).clone();
432 let setup = RlmHostSetup {
433 caller: Arc::clone(&caller),
434 invocation,
435 system_prompt,
436 };
437 let (mut engine, mut handle) =
438 Self::new_admitted(config, &api, Some(EngineHostSetup::Rlm(setup)));
439 engine.install_validated_runtime_route(caller.route.clone());
440 engine.model_client = Some(Arc::clone(&caller.client));
441 engine.model_client_injected = caller.injected;
442 handle.client_preflight_required = false;
443 engine.repl_kernel = kernel;
444 engine.extension_prompt_block = caller.extension_prompt_block.clone();
445 Ok((engine, handle))
446 }
447
448 fn rlm_turn_spec(&self) -> Result<TurnSpec> {
449 let state = self
450 .rlm_host
451 .as_ref()
452 .ok_or_else(|| anyhow!("not an admitted RLM host"))?;
453 state.caller.validate_live().map_err(anyhow::Error::new)?;
454 let max_output_tokens = state
455 .max_tokens
456 .map(|limit| {
457 std::num::NonZeroU32::new(limit)
458 .ok_or_else(|| anyhow!("RLM max_tokens must be greater than zero"))
459 })
460 .transpose()?;
461 let content = if state.mode == RlmMode::Completion {
462 state.prompt.clone()
463 } else {
464 crate::rlm::turn::metadata_text(&state.prompt, 0, None, None)
465 };
466 Ok(TurnSpec {
467 content,
468 images: Vec::new(),
469 mode: state.caller.authority.mode,
470 route: Box::new(state.caller.route.clone().into_resolved()),
471 compaction: Box::new(self.config.compaction.clone()),
472 initial_routed_usage: Box::default(),
473 goal_objective: None,
474 goal_token_budget: None,
475 goal_status: GoalStatus::Paused,
476 reasoning_effort: self.session.reasoning_effort.clone(),
477 reasoning_effort_auto: false,
478 auto_model: false,
479 allow_shell: state.caller.authority.allow_shell,
480 trust_mode: state.caller.authority.trust_mode,
481 auto_approve: state.caller.authority.auto_approve,
482 approval_mode: state.caller.authority.approval_mode,
483 translation_enabled: false,
484 allowed_tools: Some(if state.mode == RlmMode::Completion {
485 Vec::new()
486 } else {
487 vec![tool_catalog::CODE_EXECUTION_TOOL_NAME.to_string()]
488 }),
489 dynamic_tools: Vec::new(),
490 hook_executor: self.config.hook_executor.clone(),
491 verbosity: None,
492 provenance: UserInputProvenance::Runtime,
493 submission_id: None,
494 max_output_tokens,
495 })
496 }
497
498 fn rlm_result(&self, outcome: Option<SendMessageOutcome>, duration: Duration) -> RlmTurnResult {
499 let state = self.rlm_host.as_ref().expect("admitted RLM state");
500 let error = state.error.clone().or_else(|| match outcome {
501 Some(SendMessageOutcome::Finished {
502 status: TurnOutcomeStatus::Completed,
503 error,
504 }) => error,
505 Some(SendMessageOutcome::Finished { error, .. })
506 | Some(SendMessageOutcome::NotStarted { error }) => {
507 error.or_else(|| Some("RLM Core turn did not complete".into()))
508 }
509 None => Some(
510 "RLM evaluation cancelled or exhausted its original wall-clock deadline".into(),
511 ),
512 });
513 RlmTurnResult {
514 answer: state
515 .final_answer
516 .clone()
517 .unwrap_or_else(|| state.last_response.clone()),
518 iterations: state.model_rounds,
519 duration,
520 error,
521 usage: Usage::default(),
522 routed_usage: Vec::new(),
523 routed_usage_drop_records: Vec::new(),
524 routed_usage_dropped_records: 0,
525 termination: state.termination.unwrap_or(RlmTermination::Error),
526 trace: state.trace.clone(),
527 total_rpcs: state.total_rpcs,
528 }
529 .require_answer_or_error()
530 }
531 }
532
533 impl EngineHostSetup {
534 pub(super) fn initial_authority(&self, api: &Config) -> LiveRuntimeAuthority {
535 self.context()
536 .live_posture
537 .as_ref()
538 .map(LivePosture::read)
539 .unwrap_or_else(|| match self {
540 Self::Child(setup) => {
541 let runtime = &setup.authority.runtime;
542 let context = &runtime.context;
543 LiveRuntimeAuthority::from_fields(
544 runtime.parent_mode,
545 runtime.allow_shell,
546 context.trust_mode,
547 context.auto_approve,
548 context.approval_mode,
549 api.sandbox_mode.clone(),
550 )
551 }
552 Self::Rlm(setup) => LiveRuntimeAuthority::from_fields(
553 setup.caller.authority.mode,
554 setup.caller.authority.allow_shell,
555 setup.caller.authority.trust_mode,
556 setup.caller.authority.auto_approve,
557 setup.caller.authority.approval_mode,
558 api.sandbox_mode.clone(),
559 ),
560 })
561 }
562 pub(super) fn cancel_token(&self) -> &CancellationToken {
563 match self {
564 Self::Child(setup) => &setup.authority.runtime.cancel_token,
565 Self::Rlm(setup) => setup
566 .caller
567 .context
568 .cancel_token
569 .as_ref()
570 .expect("captured RLM cancel"),
571 }
572 }
573 pub(super) fn context(&self) -> &ToolContext {
574 match self {
575 Self::Child(setup) => &setup.authority.runtime.context,
576 Self::Rlm(setup) => &setup.caller.context,
577 }
578 }
579 pub(super) fn subagent_manager(&self) -> &SharedSubAgentManager {
580 match self {
581 Self::Child(setup) => &setup.authority.runtime.manager,
582 Self::Rlm(setup) => &setup.caller.subagent_manager,
583 }
584 }
585 pub(super) fn owner_agent_id(&self) -> Option<String> {
586 match self {
587 Self::Child(setup) => Some(setup.authority.owner_agent_id.clone()),
588 Self::Rlm(setup) => setup.caller.context.owner_agent_id.clone(),
589 }
590 }
591 pub(super) fn client(&self) -> &CodewhaleClient {
592 match self {
593 Self::Child(setup) => &setup.authority.runtime.client,
594 Self::Rlm(setup) => &setup.caller.route.client,
595 }
596 }
597 pub(super) fn system_prompt(&self) -> &SystemPrompt {
598 match self {
599 Self::Child(setup) => &setup.system_prompt,
600 Self::Rlm(setup) => &setup.system_prompt,
601 }
602 }
603 pub(super) fn review_policy(&self) -> Arc<AutoReviewPolicy> {
604 match self {
605 Self::Child(setup) => Arc::clone(&setup.authority.runtime.auto_review_policy),
606 Self::Rlm(setup) => Arc::clone(&setup.caller.review_policy),
607 }
608 }
609 pub(super) fn approval_store(&self) -> Result<ApprovalReceiptStore, String> {
610 match self {
611 Self::Child(setup) => setup.authority.approval_receipt_store(),
612 Self::Rlm(setup) => setup.caller.approval_store.clone(),
613 }
614 }
615 }
616
617 async fn forward_rlm_event(
618 events: Option<&mpsc::Sender<Event>>,
619 event: Event,
620 depth: u32,
621 cancel: &CancellationToken,
622 deadline: tokio::time::Instant,
623 ) {
624 if let Some(events) = events
625 && let Some(message) = crate::rlm::bridge::nested_rlm_status_line(event, depth)
626 {
627 tokio::select! {
628 biased;
629 () = cancel.cancelled() => {},
630 _ = tokio::time::sleep_until(deadline) => {},
631 _ = events.send(Event::status(message)) => {},
632 }
633 }
634 }
635
636 impl Engine {
637 pub(super) fn rlm_tool_build(
638 &self,
639 authority: &TurnAuthority,
640 _route: &TurnRouteContext,
641 ) -> TurnToolBuild {
642 let state = self.rlm_host.as_ref().expect("admitted RLM host");
643 let mut context = state.caller.context.clone();
644 context.rlm_caller = Some(Arc::clone(&state.caller));
645 context.cancel_token = Some(self.cancel_token.clone());
646 context.turn_deadline = Some(state.deadline);
647 let registry = ToolRegistryBuilder::new().build(context);
648 let mut catalog = Vec::new();
649 if state.mode != RlmMode::Completion {
650 tool_catalog::ensure_advanced_tooling(
651 &mut catalog,
652 authority.mode,
653 &HashSet::new(),
654 tool_catalog::ToolMode::Direct,
655 );
656 catalog.retain(|tool| tool.name == tool_catalog::CODE_EXECUTION_TOOL_NAME);
657 }
658 let allowed = Some(catalog.iter().map(|tool| tool.name.clone()).collect());
659 TurnToolBuild {
660 surface: ToolSurfacePolicy::new(
661 registry,
662 Some(catalog),
663 authority.mode,
664 &HashSet::new(),
665 &[],
666 false,
667 allowed,
668 None,
669 self.config.max_tool_calls,
670 tool_catalog::ToolMode::Direct,
671 ),
672 mcp_tool_names: Vec::new(),
673 mcp: McpToolState::Disabled,
674 subagent_runtime_model: None,
675 mailbox: None,
676 plugin_tool_names: HashSet::new(),
677 }
678 }
679 }
680
681 /// One actual producer dispatch, reserved before entering the client. Its Drop
682 /// records unknown execution synchronously; no future is needed to retain it.
683 pub(super) struct RlmDispatchedRequest {
684 caller: Arc<CapturedRlmCaller>,
685 usage: RlmUsageAccumulator,
686 reservation: RlmUsageReservation,
687 source: String,
688 route: crate::cost_status::EffectiveRouteEnvelope,
689 settled: Option<(
690 Usage,
691 Option<u64>,
692 crate::cost_status::RuntimeUsageMissingReason,
693 )>,
694 projected: bool,
695 refused: bool,
696 }
697 impl RlmDispatchedRequest {
698 pub(super) async fn settle_open_error(&mut self, error: &anyhow::Error) {
699 if crate::tools::subagent::engine::provider_request_refusal_proven(error) {
700 self.usage.cancel_sync(
701 self.reservation,
702 false,
703 crate::cost_status::RuntimeUsageMissingReason::RequestOutcomeUnknown,
704 );
705 self.refused = true;
706 } else {
707 self.settle(&Usage::default(), false).await;
708 }
709 }
710 pub(super) async fn settle(&mut self, usage: &Usage, complete: bool) {
711 if self.settled.is_some() || self.refused {
712 return;
713 }
714 let reason = if complete {
715 crate::cost_status::RuntimeUsageMissingReason::SuccessWithoutUsage
716 } else {
717 crate::cost_status::RuntimeUsageMissingReason::RequestOutcomeUnknown
718 };
719 // Publication precedes any cancellable worker projection. Drop retains
720 // this exact settlement and never synthesizes a second charge.
721 let priced = self
722 .caller
723 .publish_response(&self.source, &self.route, usage, reason);
724 self.settled = Some((usage.clone(), priced, reason));
725 if usage_has_reported_data(usage) {
726 self.usage.complete(self.reservation, usage).await;
727 } else {
728 self.usage.cancel_sync(self.reservation, true, reason);
729 }
730 if let Some(child) = &self.caller.child_accounting {
731 child
732 .project_settled(&self.source, &self.route, usage, priced, reason)
733 .await;
734 }
735 self.projected = true;
736 }
737 }
738 impl Drop for RlmDispatchedRequest {
739 fn drop(&mut self) {
740 if self.refused || self.projected {
741 return;
742 }
743 let (usage, priced, reason) = self.settled.clone().unwrap_or_else(|| {
744 let reason = crate::cost_status::RuntimeUsageMissingReason::RequestOutcomeUnknown;
745 let usage = Usage::default();
746 let priced = self
747 .caller
748 .publish_response(&self.source, &self.route, &usage, reason);
749 self.usage.cancel_sync(self.reservation, true, reason);
750 (usage, priced, reason)
751 });
752 if let Some(child) = &self.caller.child_accounting {
753 child.recover_settled(
754 &self.caller.scheduler,
755 self.source.clone(),
756 self.route.clone(),
757 usage,
758 priced,
759 reason,
760 );
761 }
762 }
763 }
764 impl CapturedRlmCaller {
765 fn publish_response(
766 &self,
767 source: &str,
768 route: &crate::cost_status::EffectiveRouteEnvelope,
769 usage: &Usage,
770 reason: crate::cost_status::RuntimeUsageMissingReason,
771 ) -> Option<u64> {
772 if let Some(child) = &self.child_accounting {
773 return child.publish(source, route, usage, reason);
774 }
775 if usage_has_reported_data(usage) {
776 if let Some(owner) = self.runtime_owner.as_deref() {
777 crate::cost_status::report_effective_route_for_runtime(
778 self.origin_scope,
779 Some(owner),
780 source,
781 route,
782 usage,
783 );
784 } else {
785 crate::cost_status::report_effective_route_for_interactive_origin(
786 self.origin_scope,
787 &self.origin_session,
788 &self.origin_turn,
789 source,
790 route,
791 usage,
792 );
793 }
794 } else if let Some(owner) = self.runtime_owner.as_deref() {
795 crate::cost_status::report_missing_runtime_usage(
796 self.origin_scope,
797 Some(owner),
798 source,
799 route,
800 reason,
801 );
802 } else {
803 crate::cost_status::report_missing_usage_for_interactive_origin(
804 self.origin_scope,
805 &self.origin_session,
806 &self.origin_turn,
807 source,
808 route,
809 reason,
810 );
811 }
812 None
813 }
814 }
815 impl Engine {
816 pub(super) async fn rlm_provider_request(
817 &mut self,
818 client: &SharedModelClient,
819 request: &codewhale_models::MessageRequest,
820 ) -> Result<Option<RlmDispatchedRequest>, String> {
821 let Some(state) = self.rlm_host.as_mut() else {
822 return Ok(None);
823 };
824 state
825 .caller
826 .validate_live()
827 .map_err(|error| error.to_string())?;
828 let route = client.effective_route_envelope(&request.model, chrono::Utc::now());
829 let reservation = state.usage.reserve(route.clone()).await?;
830 let source = state
831 .usage
832 .source_id(reservation)
833 .await
834 .ok_or_else(|| "RLM dispatch reservation disappeared".to_string())?;
835 state.model_rounds = state.model_rounds.saturating_add(1);
836 Ok(Some(RlmDispatchedRequest {
837 caller: Arc::clone(&state.caller),
838 usage: state.usage.clone(),
839 reservation,
840 source,
841 route,
842 settled: None,
843 projected: false,
844 refused: false,
845 }))
846 }
847 }
848
849 #[cfg(test)]
850 impl CapturedRlmCaller {
851 pub(crate) fn fixture_origin_cancel(&self) -> CancellationToken {
852 self.context
853 .cancel_token
854 .as_ref()
855 .expect("captured originating token")
856 .clone()
857 }
858 pub(crate) fn hold_fixture_workspace(&mut self, workspace: Arc<tempfile::TempDir>) {
859 self._fixture_workspace = Some(workspace);
860 }
861 }
862
862 lines RUST