返回 CodeWhale
bridge.rs
根目录 / crates / tui / src / rlm / bridge.rs
1 //! RPC bridge that services `llm_query` / `rlm_query` calls coming back
2 //! from the long-lived Python REPL during an RLM turn.
3 //!
4 //! This is the spiritual successor to the HTTP sidecar from earlier
5 //! versions — except instead of binding a localhost port and routing
6 //! through `urllib`, requests come in through stdin/stdout and we just
7 //! submit an admitted call to the canonical Engine.
8 //!
9 //! The bridge tracks cumulative token usage and the recursion budget. For
10 //! `Rlm` / `RlmBatch` requests it submits an admitted Core turn at depth-1.
11 //! Python and its persistent session never retain this borrowed dispatcher.
12
13 use crate::repl::runtime::{BatchResp, RpcDispatcher, RpcRequest, RpcResponse, SingleResp};
14 use codewhale_models::Usage;
15 use futures_util::future::join_all;
16 use std::sync::Arc;
17 use std::sync::Mutex;
18 use std::time::Duration;
19 use uuid::Uuid;
20
21 pub const MAX_BATCH: usize = 16;
22
23 /// One pre-dispatch reservation in the shared routed-usage ledger.
24 #[derive(Debug, Clone, Copy)]
25 pub(crate) struct RlmUsageReservation {
26 index: usize,
27 }
28
29 #[derive(Debug, Default)]
30 struct RlmUsageState {
31 ledger_id: String,
32 usage: Usage,
33 records: Vec<Option<RlmUsageSlot>>,
34 drop_records: Vec<crate::cost_status::RuntimeUsageDropRecord>,
35 dropped_records: u64,
36 nested_events: Vec<serde_json::Value>,
37 nested_runs: usize,
38 }
39
40 #[derive(Debug)]
41 struct RlmUsageSlot {
42 record: crate::cost_status::RuntimeUsageRecord,
43 completed: bool,
44 }
45
46 /// Shared, bounded provider-call ledger for one complete RLM tree.
47 ///
48 /// Every root, child, batch member, and recursive call reserves one slot
49 /// before invoking a provider. A distinct call is never coalesced merely
50 /// because it used the same route: its dispatch instant and frozen quote are
51 /// independent accounting evidence. Sharing one accumulator across recursion
52 /// makes the bound global instead of allowing every nested bridge to reset it.
53 #[derive(Debug, Clone)]
54 pub(crate) struct RlmUsageAccumulator {
55 state: Arc<Mutex<RlmUsageState>>,
56 }
57
58 /// Atomic snapshot returned after all RPC work for a round has settled.
59 #[derive(Debug, Clone, Default)]
60 pub(crate) struct RlmUsageSnapshot {
61 pub usage: Usage,
62 pub records: Vec<crate::cost_status::RuntimeUsageRecord>,
63 /// Exact frozen routes for provider-success responses that did not carry
64 /// authoritative usage. Keeping these separate prevents a missing payload
65 /// from becoming a priced zero-usage receipt.
66 pub drop_records: Vec<crate::cost_status::RuntimeUsageDropRecord>,
67 /// Calls whose execution/usage became ambiguous (for example a timeout).
68 /// They are never represented as authoritative zero-usage responses.
69 pub dropped_records: u64,
70 /// Producer-owned nested loop events, retained in the enclosing tool body
71 /// so every session host persists them through its existing result path.
72 pub nested_events: Vec<serde_json::Value>,
73 }
74
75 impl RlmUsageAccumulator {
76 #[must_use]
77 pub(crate) fn new() -> Self {
78 Self {
79 state: Arc::new(Mutex::new(RlmUsageState {
80 ledger_id: Uuid::new_v4().simple().to_string(),
81 ..RlmUsageState::default()
82 })),
83 }
84 }
85
86 /// Reserve durable accounting capacity before a provider request.
87 /// Definite transport failure cancels the slot; ambiguous cancellation is
88 /// explicit incomplete coverage. Reaching the cap rejects before any
89 /// unreceipted provider work can occur.
90 pub(crate) async fn reserve(
91 &self,
92 route: crate::cost_status::EffectiveRouteEnvelope,
93 ) -> std::result::Result<RlmUsageReservation, String> {
94 let mut state = self
95 .state
96 .lock()
97 .unwrap_or_else(std::sync::PoisonError::into_inner);
98 if state.records.len() == crate::cost_status::MAX_CHILD_USAGE_RECORDS {
99 return Err(format!(
100 "RLM provider-call receipt limit reached ({}); request rejected before dispatch",
101 crate::cost_status::MAX_CHILD_USAGE_RECORDS
102 ));
103 }
104 let index = state.records.len();
105 let source_id = format!("rlm:{}:request:{index}", state.ledger_id);
106 state.records.push(Some(RlmUsageSlot {
107 record: crate::cost_status::RuntimeUsageRecord {
108 source_id,
109 usage: crate::cost_status::EffectiveRouteUsage {
110 route: route.sanitized_for_persistence(),
111 usage: Usage::default(),
112 },
113 },
114 completed: false,
115 }));
116 Ok(RlmUsageReservation { index })
117 }
118
119 pub(crate) async fn source_id(&self, reservation: RlmUsageReservation) -> Option<String> {
120 self.state
121 .lock()
122 .unwrap_or_else(std::sync::PoisonError::into_inner)
123 .records
124 .get(reservation.index)
125 .and_then(Option::as_ref)
126 .map(|slot| slot.record.source_id.clone())
127 }
128
129 /// Attach a provider's reported usage to its already-reserved exact route.
130 pub(crate) async fn complete(&self, reservation: RlmUsageReservation, usage: &Usage) {
131 let mut state = self
132 .state
133 .lock()
134 .unwrap_or_else(std::sync::PoisonError::into_inner);
135 let completed = if let Some(Some(slot)) = state.records.get_mut(reservation.index)
136 && !slot.completed
137 {
138 super::add_usage_with_prompt_cache(&mut slot.record.usage.usage, usage);
139 slot.completed = true;
140 true
141 } else {
142 false
143 };
144 if completed {
145 super::add_usage_with_prompt_cache(&mut state.usage, usage);
146 }
147 }
148
149 /// Remove a reservation that never produced provider-reported usage.
150 /// Ambiguous execution increments explicit incomplete coverage instead of
151 /// being persisted as a priced-zero response.
152 #[cfg(test)]
153 pub(crate) async fn cancel(&self, reservation: RlmUsageReservation, coverage_unknown: bool) {
154 self.cancel_sync(
155 reservation,
156 coverage_unknown,
157 crate::cost_status::RuntimeUsageMissingReason::RequestOutcomeUnknown,
158 );
159 }
160
161 pub(crate) fn cancel_sync(
162 &self,
163 reservation: RlmUsageReservation,
164 coverage_unknown: bool,
165 reason: crate::cost_status::RuntimeUsageMissingReason,
166 ) {
167 let mut state = self
168 .state
169 .lock()
170 .unwrap_or_else(std::sync::PoisonError::into_inner);
171 let cancelled = state.records.get_mut(reservation.index).and_then(|slot| {
172 if slot.as_ref().is_some_and(|slot| !slot.completed) {
173 slot.take()
174 } else {
175 None
176 }
177 });
178 if let Some(slot) = cancelled
179 && coverage_unknown
180 {
181 state
182 .drop_records
183 .push(crate::cost_status::RuntimeUsageDropRecord {
184 reason,
185 source_id: slot.record.source_id,
186 route: slot.record.usage.route,
187 });
188 state.dropped_records = state.dropped_records.saturating_add(1);
189 }
190 }
191
192 /// Bound even nested runs that fail before their first provider request.
193 /// Without this reservation repeated setup failures could grow event
194 /// receipts while never consuming the existing provider-call bound.
195 pub(crate) async fn reserve_nested_turn(&self) -> std::result::Result<(), String> {
196 let mut state = self
197 .state
198 .lock()
199 .unwrap_or_else(std::sync::PoisonError::into_inner);
200 if state.nested_runs == crate::cost_status::MAX_CHILD_USAGE_RECORDS {
201 return Err("RLM nested-turn receipt limit reached before dispatch".into());
202 }
203 state.nested_runs += 1;
204 Ok(())
205 }
206
207 pub(crate) async fn record_nested_event(&self, event: serde_json::Value) {
208 self.state
209 .lock()
210 .unwrap_or_else(std::sync::PoisonError::into_inner)
211 .nested_events
212 .push(event);
213 }
214
215 pub(crate) async fn snapshot(&self) -> RlmUsageSnapshot {
216 let state = self
217 .state
218 .lock()
219 .unwrap_or_else(std::sync::PoisonError::into_inner);
220 let pending = state
221 .records
222 .iter()
223 .flatten()
224 .filter(|slot| !slot.completed)
225 .map(|slot| crate::cost_status::RuntimeUsageDropRecord {
226 reason: crate::cost_status::RuntimeUsageMissingReason::RequestOutcomeUnknown,
227 source_id: slot.record.source_id.clone(),
228 route: slot.record.usage.route.clone(),
229 })
230 .collect::<Vec<_>>();
231 let mut drop_records = state.drop_records.clone();
232 drop_records.extend(pending.iter().cloned());
233 RlmUsageSnapshot {
234 nested_events: state.nested_events.clone(),
235 usage: state.usage.clone(),
236 records: state
237 .records
238 .iter()
239 .flatten()
240 .filter(|slot| slot.completed)
241 .map(|slot| slot.record.clone())
242 .collect(),
243 drop_records,
244 dropped_records: state
245 .dropped_records
246 .saturating_add(u64::try_from(pending.len()).unwrap_or(u64::MAX)),
247 }
248 }
249 }
250
251 /// A dispatcher borrowed by one Python round. Persistent kernels never own
252 /// this caller, its services or an Engine. All nested work ends with this borrow.
253 pub(crate) struct RlmBridge<'call> {
254 caller: &'call crate::core::engine::rlm_host::CapturedRlmCaller,
255 depth_remaining: u32,
256 usage: RlmUsageAccumulator,
257 events: Option<tokio::sync::mpsc::Sender<crate::core::events::Event>>,
258 deadline: tokio::time::Instant,
259 query_timeout: Duration,
260 gate: Option<crate::tools::codemode::NestedCallGate>,
261 }
262
263 impl<'call> RlmBridge<'call> {
264 pub(crate) fn new(
265 caller: &'call crate::core::engine::rlm_host::CapturedRlmCaller,
266 depth_remaining: u32,
267 query_timeout: Duration,
268 ) -> Self {
269 Self::with_usage_accumulator(
270 caller,
271 depth_remaining,
272 query_timeout,
273 RlmUsageAccumulator::new(),
274 )
275 }
276
277 pub(crate) fn with_usage_accumulator(
278 caller: &'call crate::core::engine::rlm_host::CapturedRlmCaller,
279 depth_remaining: u32,
280 query_timeout: Duration,
281 usage: RlmUsageAccumulator,
282 ) -> Self {
283 Self {
284 caller,
285 depth_remaining,
286 usage,
287 events: None,
288 deadline: caller.deadline(),
289 query_timeout: query_timeout.clamp(Duration::from_secs(1), Duration::from_secs(600)),
290 gate: None,
291 }
292 }
293
294 pub(crate) fn with_gate(
295 mut self,
296 gate: Option<crate::tools::codemode::NestedCallGate>,
297 ) -> Self {
298 self.gate = gate;
299 self
300 }
301
302 pub(crate) fn with_deadline(mut self, deadline: Option<tokio::time::Instant>) -> Self {
303 if let Some(deadline) = deadline {
304 self.deadline = self.deadline.min(deadline);
305 }
306 self
307 }
308
309 pub(crate) fn deadline(&self) -> tokio::time::Instant {
310 self.deadline
311 }
312
313 pub(crate) fn with_events(
314 mut self,
315 events: tokio::sync::mpsc::Sender<crate::core::events::Event>,
316 ) -> Self {
317 self.events = Some(events);
318 self
319 }
320
321 pub(crate) async fn usage_snapshot(&self) -> RlmUsageSnapshot {
322 self.usage.snapshot().await
323 }
324
325 async fn invoke(
326 &self,
327 prompt: String,
328 mode: crate::core::engine::rlm_host::RlmMode,
329 max_tokens: Option<u32>,
330 system: Option<String>,
331 ) -> SingleResp {
332 let result = self
333 .caller
334 .dispatch(crate::core::engine::rlm_host::RlmInvocation {
335 prompt,
336 mode,
337 max_tokens,
338 task_instructions: system,
339 deadline: self
340 .deadline
341 .min(tokio::time::Instant::now() + self.query_timeout),
342 gate: self.gate.clone(),
343 events: self.events.clone(),
344 usage: self.usage.clone(),
345 })
346 .await;
347 SingleResp {
348 text: result.answer,
349 error: result.error,
350 }
351 }
352
353 async fn dispatch_llm(
354 &self,
355 prompt: String,
356 _model: Option<String>,
357 max_tokens: Option<u32>,
358 system: Option<String>,
359 ) -> SingleResp {
360 self.invoke(
361 prompt,
362 crate::core::engine::rlm_host::RlmMode::Completion,
363 max_tokens,
364 system,
365 )
366 .await
367 }
368
369 async fn dispatch_llm_batch(
370 &self,
371 prompts: Vec<String>,
372 _model: Option<String>,
373 dependency_mode: Option<String>,
374 ) -> BatchResp {
375 if let Some(resp) = batch_guard(prompts.len(), dependency_mode.as_deref()) {
376 return resp;
377 }
378 BatchResp {
379 results: join_all(
380 prompts
381 .into_iter()
382 .map(|prompt| self.dispatch_llm(prompt, None, None, None)),
383 )
384 .await,
385 }
386 }
387
388 pub(crate) async fn dispatch_rlm(&self, prompt: String, _model: Option<String>) -> SingleResp {
389 if self.depth_remaining == 0 {
390 return self.dispatch_llm(prompt, None, None, None).await;
391 }
392 self.invoke(
393 prompt,
394 crate::core::engine::rlm_host::RlmMode::Recursive {
395 depth_remaining: self.depth_remaining.saturating_sub(1),
396 },
397 None,
398 None,
399 )
400 .await
401 }
402
403 async fn dispatch_rlm_batch(
404 &self,
405 prompts: Vec<String>,
406 _model: Option<String>,
407 dependency_mode: Option<String>,
408 ) -> BatchResp {
409 if let Some(resp) = batch_guard(prompts.len(), dependency_mode.as_deref()) {
410 return resp;
411 }
412 BatchResp {
413 results: join_all(
414 prompts
415 .into_iter()
416 .map(|prompt| self.dispatch_rlm(prompt, None)),
417 )
418 .await,
419 }
420 }
421 }
422
423 /// One nested sub-RLM event as a parent-stream status line, or `None` for an
424 /// event kind the nested loop does not produce.
425 pub(crate) fn nested_rlm_status_line(
426 event: crate::core::events::Event,
427 depth: u32,
428 ) -> Option<String> {
429 use crate::core::events::Event;
430 let body = match event {
431 Event::Status { message } => message,
432 Event::MessageDelta { content, .. } => content.trim().to_string(),
433 _ => return None,
434 };
435 Some(format!("{NESTED_RLM_STATUS_PREFIX}{depth}): {body}"))
436 }
437
438 /// Status-line prefix for forwarded nested sub-RLM events;
439 /// `core::events::status_visibility` classifies these as internal receipts.
440 pub(crate) const NESTED_RLM_STATUS_PREFIX: &str = "sub-RLM (depth ";
441
442 fn batch_guard(prompt_count: usize, dependency_mode: Option<&str>) -> Option<BatchResp> {
443 if prompt_count == 0 {
444 return Some(BatchResp { results: vec![] });
445 }
446 if prompt_count > MAX_BATCH {
447 return Some(BatchResp {
448 results: (0..prompt_count)
449 .map(|_| SingleResp {
450 text: String::new(),
451 error: Some(format!("batch too large: {prompt_count} > {MAX_BATCH}")),
452 })
453 .collect(),
454 });
455 }
456 let mode = dependency_mode
457 .unwrap_or_default()
458 .trim()
459 .to_ascii_lowercase()
460 .replace(['-', ' '], "_");
461 if !matches!(
462 mode.as_str(),
463 "independent" | "parallel_safe" | "map_reduce"
464 ) {
465 return Some(BatchResp {
466 results: (0..prompt_count)
467 .map(|_| SingleResp {
468 text: String::new(),
469 error: Some(
470 "batch requires dependency_mode='independent'; use sub_query_sequence or sequential sub_query calls for dependent work"
471 .to_string(),
472 ),
473 })
474 .collect(),
475 });
476 }
477 None
478 }
479
480 impl RpcDispatcher for RlmBridge<'_> {
481 fn dispatch<'a>(
482 &'a self,
483 req: RpcRequest,
484 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = RpcResponse> + Send + 'a>> {
485 Box::pin(async move {
486 match req {
487 RpcRequest::Llm {
488 prompt,
489 model,
490 max_tokens,
491 system,
492 } => {
493 RpcResponse::Single(self.dispatch_llm(prompt, model, max_tokens, system).await)
494 }
495 RpcRequest::LlmBatch {
496 prompts,
497 model,
498 dependency_mode,
499 safety_note: _,
500 } => RpcResponse::Batch(
501 self.dispatch_llm_batch(prompts, model, dependency_mode)
502 .await,
503 ),
504 RpcRequest::Rlm { prompt, model } => {
505 RpcResponse::Single(self.dispatch_rlm(prompt, model).await)
506 }
507 RpcRequest::RlmBatch {
508 prompts,
509 model,
510 dependency_mode,
511 safety_note: _,
512 } => RpcResponse::Batch(
513 self.dispatch_rlm_batch(prompts, model, dependency_mode)
514 .await,
515 ),
516 }
517 })
518 }
519 }
520
521 #[cfg(test)]
522 mod tests {
523 use super::*;
524 use crate::core::engine::tests::rlm_host::{BridgeFixture, Replies};
525 use crate::llm_client::mock::MockLlmClient;
526 use anyhow::Result;
527 use codewhale_models::{ContentBlock, MessageRequest, MessageResponse};
528 use std::future::Future;
529 use std::pin::Pin;
530
531 fn mock_response_with_usage(text: &str, usage: Usage) -> MessageResponse {
532 MessageResponse {
533 id: "mock_msg".to_string(),
534 r#type: "message".to_string(),
535 role: "assistant".to_string(),
536 content: vec![ContentBlock::Text {
537 text: text.to_string(),
538 cache_control: None,
539 }],
540 model: "mock-model".to_string(),
541 stop_reason: Some("end_turn".to_string()),
542 stop_sequence: None,
543 container: None,
544 usage,
545 }
546 }
547
548 fn mock_response(text: &str, input_tokens: u32, output_tokens: u32) -> MessageResponse {
549 mock_response_with_usage(
550 text,
551 Usage {
552 input_tokens,
553 output_tokens,
554 ..Usage::default()
555 },
556 )
557 }
558
559 fn bridge_for(mock: Arc<MockLlmClient>, depth_remaining: u32) -> BridgeFixture {
560 let client: Arc<dyn Replies> = mock;
561 BridgeFixture::new(client, "child-model".to_string(), depth_remaining).with_gate(Some(
562 crate::tools::codemode::NestedCallGate::admitting_for_test(),
563 ))
564 }
565
566 #[tokio::test]
567 async fn expired_parent_deadline_and_nested_receipt_bound_refuse_before_work() {
568 let mock = Arc::new(MockLlmClient::new(Vec::new()));
569 for depth in [0, 1] {
570 let bridge =
571 bridge_for(mock.clone(), depth).with_deadline(Some(tokio::time::Instant::now()));
572 let response = bridge.dispatch_rlm("must not dispatch".into(), None).await;
573 assert!(response.error.unwrap().contains("deadline exhausted"));
574 assert_eq!(mock.call_count(), 0);
575 assert!(bridge.usage_snapshot().await.nested_events.is_empty());
576 }
577 let bridge = bridge_for(mock.clone(), 1);
578 for _ in 0..crate::cost_status::MAX_CHILD_USAGE_RECORDS {
579 bridge.usage.reserve_nested_turn().await.unwrap();
580 }
581 let response = bridge
582 .dispatch_rlm("must not spawn Python".into(), None)
583 .await;
584 assert!(response.error.unwrap().contains("receipt limit"));
585 assert_eq!(mock.call_count(), 0);
586 assert!(bridge.usage_snapshot().await.nested_events.is_empty());
587 }
588
589 struct PendingClient(MockLlmClient);
590
591 impl Replies for PendingClient {
592 fn effective_route_envelope(
593 &self,
594 model: &str,
595 at: chrono::DateTime<chrono::Utc>,
596 ) -> crate::cost_status::EffectiveRouteEnvelope {
597 Replies::effective_route_envelope(&self.0, model, at)
598 }
599 fn effective_max_output_tokens(&self, model: &str) -> u32 {
600 Replies::effective_max_output_tokens(&self.0, model)
601 }
602 fn create_message_boxed(
603 &self,
604 _: MessageRequest,
605 ) -> Pin<Box<dyn Future<Output = Result<MessageResponse>> + Send + '_>> {
606 Box::pin(std::future::pending())
607 }
608 }
609
610 #[tokio::test]
611 async fn parent_deadline_interrupts_plain_and_recursive_model_calls() {
612 for depth in [0, 1] {
613 let bridge = BridgeFixture::new(
614 Arc::new(PendingClient(MockLlmClient::new(Vec::new()))),
615 "child-model".into(),
616 depth,
617 )
618 .with_deadline(Some(
619 tokio::time::Instant::now() + Duration::from_millis(500),
620 ))
621 .with_gate(Some(
622 crate::tools::codemode::NestedCallGate::admitting_for_test(),
623 ));
624 let response = tokio::time::timeout(
625 Duration::from_secs(5),
626 bridge.dispatch_rlm("bounded nested context".into(), None),
627 )
628 .await
629 .expect("a child must not replace the inherited budget");
630 assert!(
631 response
632 .error
633 .as_deref()
634 .is_some_and(|error| error.contains("deadline")),
635 "{:?}",
636 response.error
637 );
638 if depth == 0 {
639 assert_eq!(
640 bridge.usage_snapshot().await.dropped_records,
641 1,
642 "canceled provider work has unknown usage, never priced zero"
643 );
644 }
645 }
646 }
647
648 #[test]
649 fn batch_guard_allows_non_empty_batches_at_the_cap() {
650 assert!(batch_guard(MAX_BATCH, Some("independent")).is_none());
651 }
652
653 #[test]
654 fn batch_guard_returns_empty_response_for_empty_batches() {
655 let response = batch_guard(0, None).expect("empty batch should be handled");
656 assert!(response.results.is_empty());
657 }
658
659 #[test]
660 fn batch_guard_returns_one_error_per_oversized_prompt() {
661 let response = batch_guard(MAX_BATCH + 2, Some("independent"))
662 .expect("oversized batch should be handled");
663 assert_eq!(response.results.len(), MAX_BATCH + 2);
664 assert!(response.results.iter().all(|result| {
665 result.text.is_empty()
666 && result
667 .error
668 .as_deref()
669 .is_some_and(|err| err.contains("batch too large"))
670 }));
671 }
672
673 #[test]
674 fn batch_guard_requires_explicit_independence_for_parallel_work() {
675 let response = batch_guard(2, None).expect("missing dependency mode should be handled");
676 assert_eq!(response.results.len(), 2);
677 assert!(response.results.iter().all(|result| {
678 result.text.is_empty()
679 && result
680 .error
681 .as_deref()
682 .is_some_and(|err| err.contains("dependency_mode='independent'"))
683 }));
684
685 let response = batch_guard(2, Some("sequential"))
686 .expect("dependent dependency mode should be handled");
687 assert!(response.results.iter().all(|result| {
688 result
689 .error
690 .as_deref()
691 .is_some_and(|err| err.contains("sub_query_sequence"))
692 }));
693 }
694
695 #[tokio::test]
696 async fn llm_dispatch_pins_configured_child_model() {
697 let mock = Arc::new(MockLlmClient::new(Vec::new()));
698 mock.push_message_response(mock_response("child answer", 7, 11));
699 let bridge = bridge_for(Arc::clone(&mock), 1);
700
701 let response = bridge
702 .dispatch(RpcRequest::Llm {
703 prompt: "child prompt".to_string(),
704 model: Some("override-model".to_string()),
705 max_tokens: Some(123),
706 system: Some("child system".to_string()),
707 })
708 .await;
709
710 match response {
711 RpcResponse::Single(single) => {
712 assert_eq!(single.text, "child answer");
713 assert!(single.error.is_none());
714 }
715 other => panic!("expected single response, got {other:?}"),
716 }
717
718 let captured = mock.captured_requests();
719 assert_eq!(captured.len(), 1);
720 assert_eq!(captured[0].model, "child-model");
721 assert_eq!(captured[0].max_tokens, 123);
722 let system = serde_json::to_string(&captured[0].system).unwrap();
723 assert!(system.contains("Captured operator Core policy"));
724 assert!(system.contains("child system"));
725
726 let snapshot = bridge.usage_snapshot().await;
727 assert_eq!(snapshot.usage.input_tokens, 7);
728 assert_eq!(snapshot.usage.output_tokens, 11);
729 assert_eq!(snapshot.records.len(), 1);
730 assert_eq!(snapshot.records[0].usage.usage, snapshot.usage);
731 assert!(snapshot.drop_records.is_empty());
732 assert_eq!(snapshot.dropped_records, 0);
733 }
734
735 #[tokio::test]
736 async fn llm_dispatch_keeps_semantic_success_but_marks_missing_usage_once() {
737 let mock = Arc::new(MockLlmClient::new(Vec::new()));
738 mock.push_message_response(mock_response_with_usage(
739 "usable child answer",
740 Usage::default(),
741 ));
742 let bridge = bridge_for(Arc::clone(&mock), 1);
743
744 let response = bridge
745 .dispatch(RpcRequest::Llm {
746 prompt: "child prompt".to_string(),
747 model: None,
748 max_tokens: None,
749 system: None,
750 })
751 .await;
752
753 let RpcResponse::Single(response) = response else {
754 panic!("expected single response");
755 };
756 assert_eq!(response.text, "usable child answer");
757 assert!(response.error.is_none());
758
759 let first = bridge.usage_snapshot().await;
760 let replay = bridge.usage_snapshot().await;
761 assert_eq!(first.usage, Usage::default());
762 assert!(first.records.is_empty());
763 assert_eq!(first.drop_records.len(), 1);
764 assert_eq!(first.dropped_records, 1);
765 assert_eq!(replay.drop_records, first.drop_records);
766 assert_eq!(replay.dropped_records, 1);
767 assert_eq!(first.drop_records[0].route.model, "child-model");
768 assert!(first.drop_records[0].source_id.starts_with("rlm:"));
769 }
770
771 #[tokio::test]
772 async fn repeated_reservation_settlement_cannot_duplicate_usage_or_missing_coverage() {
773 let mock = Arc::new(MockLlmClient::new(Vec::new()));
774 let bridge = bridge_for(Arc::clone(&mock), 1);
775 let route =
776 Replies::effective_route_envelope(mock.as_ref(), "child-model", chrono::Utc::now());
777
778 let usage_reservation = bridge
779 .usage
780 .reserve(route.clone())
781 .await
782 .expect("usage reservation");
783 let reported = Usage {
784 input_tokens: 3,
785 output_tokens: 5,
786 ..Usage::default()
787 };
788 bridge.usage.complete(usage_reservation, &reported).await;
789 bridge.usage.complete(usage_reservation, &reported).await;
790
791 let missing_reservation = bridge
792 .usage
793 .reserve(route)
794 .await
795 .expect("missing reservation");
796 bridge.usage.cancel(missing_reservation, true).await;
797 bridge.usage.cancel(missing_reservation, true).await;
798
799 let snapshot = bridge.usage_snapshot().await;
800 assert_eq!(snapshot.usage, reported);
801 assert_eq!(snapshot.records.len(), 1);
802 assert_eq!(snapshot.drop_records.len(), 1);
803 assert_eq!(snapshot.dropped_records, 1);
804 }
805
806 #[tokio::test]
807 async fn llm_dispatch_preserves_prompt_cache_usage() {
808 let mock = Arc::new(MockLlmClient::new(Vec::new()));
809 mock.push_message_response(mock_response_with_usage(
810 "cached child answer",
811 Usage {
812 input_tokens: 1000,
813 output_tokens: 100,
814 prompt_cache_hit_tokens: Some(800),
815 prompt_cache_miss_tokens: Some(200),
816 ..Usage::default()
817 },
818 ));
819 let bridge = bridge_for(Arc::clone(&mock), 1);
820
821 let response = bridge
822 .dispatch(RpcRequest::Llm {
823 prompt: "child prompt".to_string(),
824 model: None,
825 max_tokens: None,
826 system: None,
827 })
828 .await;
829
830 match response {
831 RpcResponse::Single(single) => {
832 assert_eq!(single.text, "cached child answer");
833 assert!(single.error.is_none());
834 }
835 other => panic!("expected single response, got {other:?}"),
836 }
837
838 let usage = bridge.usage_snapshot().await.usage;
839 assert_eq!(usage.input_tokens, 1000);
840 assert_eq!(usage.output_tokens, 100);
841 assert_eq!(usage.prompt_cache_hit_tokens, Some(800));
842 assert_eq!(usage.prompt_cache_miss_tokens, Some(200));
843 }
844
845 #[tokio::test]
846 async fn llm_dispatch_rejects_max_tokens_partial_output_after_charging_usage() {
847 let mock = Arc::new(MockLlmClient::new(Vec::new()));
848 let usage = Usage {
849 input_tokens: 23,
850 output_tokens: 4096,
851 reasoning_tokens: Some(4000),
852 ..Usage::default()
853 };
854 let mut response = mock_response_with_usage(
855 "FINAL('partial answer')\n```repl\nFINAL('also partial')\n```",
856 usage.clone(),
857 );
858 response.stop_reason = Some("max_tokens".to_string());
859 mock.push_message_response(response);
860 let bridge = bridge_for(Arc::clone(&mock), 1);
861
862 let response = bridge
863 .dispatch(RpcRequest::Llm {
864 prompt: "child prompt".to_string(),
865 model: None,
866 max_tokens: None,
867 system: None,
868 })
869 .await;
870
871 match response {
872 RpcResponse::Single(single) => {
873 assert!(
874 single.text.is_empty(),
875 "partial output must not be accepted"
876 );
877 let error = single.error.expect("truncation must surface as an error");
878 assert!(error.contains("incomplete"), "{error}");
879 assert!(error.contains("max_tokens"), "{error}");
880 }
881 other => panic!("expected single response, got {other:?}"),
882 }
883
884 let snapshot = bridge.usage_snapshot().await;
885 assert_eq!(snapshot.usage, usage);
886 assert_eq!(snapshot.records.len(), 1);
887 assert_eq!(snapshot.records[0].usage.usage, usage);
888 let rejected_fragments: Vec<_> = snapshot
889 .nested_events
890 .iter()
891 .filter(|event| {
892 event["kind"] == "code"
893 && event["content"]
894 == "FINAL('partial answer')\n```repl\nFINAL('also partial')\n```"
895 })
896 .collect();
897 assert_eq!(
898 rejected_fragments.len(),
899 1,
900 "retain the real rejected diagnostic fragment"
901 );
902 assert_eq!(mock.call_count(), 1, "truncation must not retry");
903 }
904
905 #[tokio::test]
906 async fn llm_batch_dispatch_pins_configured_child_model() {
907 let mock = Arc::new(MockLlmClient::new(Vec::new()));
908 mock.push_message_response(mock_response("one", 1, 2));
909 mock.push_message_response(mock_response("two", 3, 4));
910 mock.push_message_response(mock_response("three", 5, 6));
911 let bridge = bridge_for(Arc::clone(&mock), 1);
912
913 let response = bridge
914 .dispatch(RpcRequest::LlmBatch {
915 prompts: vec!["a".to_string(), "b".to_string(), "c".to_string()],
916 model: Some("batch-model".to_string()),
917 dependency_mode: Some("independent".to_string()),
918 safety_note: Some("test prompts are independent".to_string()),
919 })
920 .await;
921
922 match response {
923 RpcResponse::Batch(batch) => {
924 let texts: Vec<_> = batch
925 .results
926 .iter()
927 .map(|result| result.text.as_str())
928 .collect();
929 assert_eq!(texts, ["one", "two", "three"]);
930 assert!(batch.results.iter().all(|result| result.error.is_none()));
931 }
932 other => panic!("expected batch response, got {other:?}"),
933 }
934
935 let captured = mock.captured_requests();
936 assert_eq!(captured.len(), 3);
937 assert!(
938 captured
939 .iter()
940 .all(|request| request.model == "child-model")
941 );
942
943 let snapshot = bridge.usage_snapshot().await;
944 assert_eq!(snapshot.usage.input_tokens, 9);
945 assert_eq!(snapshot.usage.output_tokens, 12);
946 assert_eq!(snapshot.records.len(), 3);
947 assert_ne!(
948 snapshot.records[0].source_id, snapshot.records[1].source_id,
949 "distinct provider calls must keep distinct stable identities"
950 );
951 }
952
953 #[tokio::test]
954 async fn shared_accumulator_rejects_the_first_unreceipted_request_before_provider_work() {
955 let mock = Arc::new(MockLlmClient::new(Vec::new()));
956 let client: Arc<dyn Replies> = mock.clone();
957 let usage = RlmUsageAccumulator::new();
958 let bridge = BridgeFixture::with_usage_accumulator(
959 Arc::clone(&client),
960 "child-model".to_string(),
961 1,
962 usage.clone(),
963 );
964 let nested_bridge =
965 BridgeFixture::with_usage_accumulator(client, "child-model".to_string(), 1, usage);
966 let route =
967 Replies::effective_route_envelope(mock.as_ref(), "child-model", chrono::Utc::now());
968 for _ in 0..crate::cost_status::MAX_CHILD_USAGE_RECORDS {
969 let reservation = bridge
970 .usage
971 .reserve(route.clone())
972 .await
973 .expect("receipt slot below cap");
974 bridge
975 .usage
976 .complete(
977 reservation,
978 &Usage {
979 input_tokens: 1,
980 ..Usage::default()
981 },
982 )
983 .await;
984 }
985
986 let response = nested_bridge
987 .dispatch(RpcRequest::Llm {
988 prompt: "must not reach provider".to_string(),
989 model: None,
990 max_tokens: None,
991 system: None,
992 })
993 .await;
994 let RpcResponse::Single(response) = response else {
995 panic!("expected single response");
996 };
997 assert!(
998 response
999 .error
1000 .as_deref()
1001 .is_some_and(|error| error.contains("rejected before dispatch"))
1002 );
1003 assert_eq!(mock.call_count(), 0);
1004 let snapshot = bridge.usage_snapshot().await;
1005 assert_eq!(
1006 snapshot.records.len(),
1007 crate::cost_status::MAX_CHILD_USAGE_RECORDS
1008 );
1009 assert_eq!(snapshot.dropped_records, 0);
1010 assert!(snapshot.drop_records.is_empty());
1011 }
1012
1013 /// #6511: a nested sub-RLM's events used to go to a drain task, so its
1014 /// model calls never reached the parent's record.
1015 #[tokio::test]
1016 async fn nested_rlm_events_reach_the_parent_stream() {
1017 let mock = Arc::new(MockLlmClient::new(Vec::new()));
1018 mock.push_message_response(mock_response("```repl\nFINAL('nested answer')\n```", 3, 4));
1019 let (tx, mut rx) = tokio::sync::mpsc::channel(256);
1020 let bridge = bridge_for(Arc::clone(&mock), 1).with_events(tx);
1021
1022 let response = bridge
1023 .dispatch_rlm("nested context".to_string(), None)
1024 .await;
1025 assert_eq!(response.text, "nested answer");
1026 assert!(response.error.is_none(), "{:?}", response.error);
1027
1028 let mut lines = Vec::new();
1029 while let Ok(event) = rx.try_recv() {
1030 match event {
1031 crate::core::events::Event::Status { message } => lines.push(message),
1032 other => panic!("nested events must arrive as status lines, got {other:?}"),
1033 }
1034 }
1035 assert!(
1036 lines
1037 .iter()
1038 .all(|line| line.starts_with(NESTED_RLM_STATUS_PREFIX)),
1039 "{lines:#?}"
1040 );
1041 assert!(
1042 lines
1043 .iter()
1044 .any(|line| line.contains("FINAL('nested answer')")),
1045 "the nested code round is part of the record: {lines:#?}"
1046 );
1047 assert!(
1048 lines
1049 .iter()
1050 .any(|line| line.contains("RLM finished: Final")),
1051 "{lines:#?}"
1052 );
1053 let receipts = bridge.usage_snapshot().await.nested_events;
1054 assert_eq!(
1055 receipts.len(),
1056 2,
1057 "one complete model/code reply and one terminal receipt"
1058 );
1059 assert!(receipts.iter().all(|entry| {
1060 let content = entry["content"]
1061 .as_str()
1062 .expect("canonical receipt content");
1063 lines.iter().any(|line| line.contains(content))
1064 }));
1065 assert_eq!(
1066 receipts
1067 .iter()
1068 .filter(|entry| entry["kind"] == "code")
1069 .count(),
1070 1
1071 );
1072 assert!(receipts.iter().any(|entry| {
1073 entry["content"]
1074 .as_str()
1075 .is_some_and(|text| text.contains("RLM finished: Final"))
1076 }));
1077 assert_eq!(
1078 crate::core::events::status_visibility(&lines[0]),
1079 crate::core::events::StatusVisibility::Internal
1080 );
1081 }
1082
1083 #[tokio::test]
1084 async fn rlm_dispatch_at_depth_zero_pins_configured_child_model() {
1085 let mock = Arc::new(MockLlmClient::new(Vec::new()));
1086 mock.push_message_response(mock_response("fallback answer", 3, 5));
1087 let bridge = bridge_for(Arc::clone(&mock), 0);
1088
1089 let response = bridge
1090 .dispatch(RpcRequest::Rlm {
1091 prompt: "nested prompt".to_string(),
1092 model: Some("override-model".to_string()),
1093 })
1094 .await;
1095
1096 match response {
1097 RpcResponse::Single(single) => {
1098 assert_eq!(single.text, "fallback answer");
1099 assert!(single.error.is_none());
1100 }
1101 other => panic!("expected single response, got {other:?}"),
1102 }
1103
1104 let usage = bridge.usage_snapshot().await.usage;
1105 assert_eq!(usage.input_tokens, 3);
1106 assert_eq!(usage.output_tokens, 5);
1107
1108 let captured = mock.captured_requests();
1109 assert_eq!(captured.len(), 1);
1110 assert_eq!(captured[0].model, "child-model");
1111 }
1112 }
1113
1113 lines RUST