返回 CodeWhale
turn.rs
根目录 / crates / tui / src / rlm / turn.rs
1 //! Pure RLM result and feedback facts. The canonical Engine owns every model
2 //! request, code round, cancellation, session message and permission settlement.
3 use codewhale_models::{ContentBlock, Message, Role, Usage};
4 use std::time::Duration;
5
6 pub(crate) const MAX_RLM_ITERATIONS: u32 = 25;
7 pub(crate) const MAX_CONSECUTIVE_NO_CODE: u32 = 3;
8 pub(crate) const STDOUT_METADATA_PREVIEW_LEN: usize = 800;
9 const PROMPT_PREVIEW_LEN: usize = 500;
10
11 /// How an RLM turn ended.
12 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
13 pub enum RlmTermination {
14 /// `FINAL(value)` was called inside the REPL or `FINAL(...)` appeared
15 /// at the top of the model's response on its own line.
16 Final,
17 /// The model failed to emit a `repl` block for too many rounds in a
18 /// row. The accumulated last response text is surfaced as the answer
19 /// rather than being thrown away.
20 NoCode,
21 /// Iteration cap reached without `FINAL`. The last root response is
22 /// surfaced as the answer alongside the error.
23 Exhausted,
24 /// Hard error — LLM call failed, REPL crashed, timeout.
25 Error,
26 }
27
28 /// Per-round trace entry. Surfaced in the tool result so the user can see
29 /// exactly what the sub-agent did.
30 #[derive(Debug, Clone)]
31 pub struct RlmRoundTrace {
32 pub round: u32,
33 pub code_summary: String,
34 pub stdout_preview: String,
35 pub had_error: bool,
36 pub rpc_count: u32,
37 pub elapsed_ms: u64,
38 }
39
40 /// Result of an RLM turn.
41 #[derive(Debug, Clone)]
42 pub struct RlmTurnResult {
43 pub answer: String,
44 pub iterations: u32,
45 pub duration: Duration,
46 pub error: Option<String>,
47 pub usage: Usage,
48 /// One exact frozen route/quote receipt per admitted provider request.
49 /// Distinct calls are never coalesced, even when they share a route.
50 pub routed_usage: Vec<crate::cost_status::RuntimeUsageRecord>,
51 /// Exact routes for provider-success responses that omitted authoritative
52 /// usage metadata.
53 pub routed_usage_drop_records: Vec<crate::cost_status::RuntimeUsageDropRecord>,
54 pub routed_usage_dropped_records: u64,
55 pub termination: RlmTermination,
56 /// Per-round trace. Empty when the loop never reached the REPL.
57 pub trace: Vec<RlmRoundTrace>,
58 /// Total sub-LLM RPCs made by the sub-agent (sum of `rpc_count` across
59 /// rounds). Useful for verifying that the model engaged with `context`
60 /// rather than answering directly.
61 pub total_rpcs: u32,
62 }
63
64 impl RlmTurnResult {
65 pub(crate) fn failed(error: String, duration: Duration) -> Self {
66 Self {
67 answer: String::new(),
68 iterations: 0,
69 duration,
70 error: Some(error),
71 usage: Usage::default(),
72 routed_usage: Vec::new(),
73 routed_usage_drop_records: Vec::new(),
74 routed_usage_dropped_records: 0,
75 termination: RlmTermination::Error,
76 trace: Vec::new(),
77 total_rpcs: 0,
78 }
79 }
80 pub(crate) fn require_answer_or_error(mut self) -> Self {
81 if self.termination != RlmTermination::Final
82 && self.answer.trim().is_empty()
83 && self.error.is_none()
84 {
85 self.error = Some(format!(
86 "RLM ended ({:?}) after {} iteration(s) with an empty answer",
87 self.termination, self.iterations
88 ));
89 }
90 self
91 }
92 }
93
94 pub(crate) fn metadata_text(
95 prompt: &str,
96 iteration: u32,
97 code: Option<&str>,
98 stdout: Option<&str>,
99 ) -> String {
100 extract_text_blocks(&build_metadata_message(prompt, None, iteration, code, stdout).content)
101 }
102
103 fn build_metadata_message(
104 prompt: &str,
105 root_prompt: Option<&str>,
106 iteration: u32,
107 previous_code: Option<&str>,
108 previous_stdout: Option<&str>,
109 ) -> Message {
110 let prompt_len = prompt.chars().count();
111 let prompt_preview = truncate_text(prompt, PROMPT_PREVIEW_LEN);
112
113 let mut parts = Vec::new();
114 parts.push(format!("## REPL state (round {iteration})"));
115 parts.push(String::new());
116 if let Some(rp) = root_prompt
117 && !rp.trim().is_empty()
118 {
119 parts.push("**Original task** (re-shown every round)".to_string());
120 parts.push(format!("> {}", truncate_text(rp.trim(), 600)));
121 parts.push(String::new());
122 }
123 parts.push("**`_context`** — the long input lives in the REPL only".to_string());
124 parts.push(format!("- Length: {prompt_len} chars"));
125 parts.push(format!("- Preview: \"{prompt_preview}\""));
126 parts.push(String::new());
127
128 parts.push("**REPL helpers** (use inside ```repl blocks)".to_string());
129 parts.push(
130 "- `_context` / `_ctx` / `content` — the full input string"
131 .to_string(),
132 );
133 parts.push(
134 "- `len(_context)` / `_context[a:b]` / `_context.splitlines()` — slice it".to_string(),
135 );
136 parts.push(
137 "- `chunk(max_chars=20000, overlap=0)` — full-coverage chunks with index/start/end/text"
138 .to_string(),
139 );
140 parts.push(
141 "- `chunk_coverage(chunks)` — coverage report for chunk output".to_string(),
142 );
143 parts.push(
144 "- `llm_query(prompt, model=None)` — one-shot child LLM; `model` is ignored and child calls keep the captured Core route"
145 .to_string(),
146 );
147 parts.push(
148 "- `llm_query_batched([p1, p2, ...], dependency_mode=\"independent\")` — concurrent fan-out for independent prompts only; `model` is ignored"
149 .to_string(),
150 );
151 parts.push(
152 "- `rlm_query(prompt, model=None)` — recursive sub-RLM; `model` is ignored"
153 .to_string(),
154 );
155 parts.push(
156 "- `rlm_query_batched([p1, p2, ...], dependency_mode=\"independent\")` — concurrent recursive sub-RLMs for independent prompts only; `model` is ignored"
157 .to_string(),
158 );
159 parts.push(
160 "- `sub_query_sequence(prompt, slices)` — sequential child calls for A->B dependencies and rollback-sensitive work"
161 .to_string(),
162 );
163 parts.push(
164 "- Batch safety: never batch dependent steps, global-state refactors, schema migrations, or rollback-sensitive tasks"
165 .to_string(),
166 );
167 parts.push("- `SHOW_VARS()` — list user variables".to_string());
168 parts.push("- `repl_set(name, value)` / `repl_get(name)` — explicit store".to_string());
169 parts.push(
170 "- `FINAL(value)` — end the loop with this answer".to_string(),
171 );
172 parts.push(
173 "- `FINAL_VAR(name)` — end the loop with a variable's value"
174 .to_string(),
175 );
176 parts.push(String::new());
177
178 if iteration > 0 {
179 parts.push("**Previous round**".to_string());
180 if let Some(code) = previous_code {
181 parts.push(format!("- Code: {}", summarize_code(code)));
182 }
183 if let Some(stdout) = previous_stdout {
184 let stdout_clean = stdout.trim();
185 if !stdout_clean.is_empty() {
186 parts.push(format!("- Stdout preview: \"{stdout_clean}\""));
187 } else {
188 parts.push("- Stdout: (empty)".to_string());
189 }
190 }
191 }
192
193 let text = parts.join("\n");
194
195 Message {
196 role: Role::User,
197 content: vec![ContentBlock::Text {
198 text,
199 cache_control: None,
200 }],
201 }
202 }
203
204 pub(crate) fn summarize_code(code: &str) -> String {
205 let lines: Vec<&str> = code.lines().collect();
206 if lines.len() <= 8 {
207 return code.to_string();
208 }
209 let head = lines[..4].join("\n");
210 let tail = lines[lines.len() - 4..].join("\n");
211 format!("{} lines:\n{head}\n…\n{tail}", lines.len())
212 }
213
214 fn extract_text_blocks(blocks: &[ContentBlock]) -> String {
215 blocks
216 .iter()
217 .filter_map(|b| match b {
218 ContentBlock::Text { text, .. } => Some(text.as_str()),
219 _ => None,
220 })
221 .collect::<Vec<_>>()
222 .join("\n")
223 }
224
225 /// Extract the first ` ```repl ` block from `text`. Falls back to
226 /// ` ```python `/`` ```py `` for compatibility with prompts that learned
227 /// the older fence style.
228 ///
229 /// The opening fence must start its own line (up to three spaces of indent,
230 /// as in Markdown) with no other info string, so prose that mentions a fence
231 /// mid-line is never code to run.
232 pub(crate) fn extract_repl_code(text: &str) -> Option<String> {
233 let mut lines = text.split_inclusive('\n');
234 let mut offset = 0;
235 let code_start = loop {
236 let line = lines.next()?;
237 offset += line.len();
238 let indent = line.len() - line.trim_start_matches(' ').len();
239 if indent <= 3
240 && line[indent..].strip_prefix("```").is_some_and(|info| {
241 matches!(info.trim(), "repl" | "python" | "py") && line.ends_with('\n')
242 })
243 {
244 break offset;
245 }
246 };
247 let after_fence = &text[code_start..];
248 let end_idx = if after_fence.starts_with("```") {
249 0
250 } else {
251 after_fence.find("\n```")?
252 };
253
254 let code = after_fence[..end_idx].trim().to_string();
255 if code.is_empty() {
256 return None;
257 }
258 Some(code)
259 }
260
261 /// Parse a top-level `FINAL(...)` directive from the model's raw text.
262 /// Mirrors the reference RLM's `find_final_answer`: directive must appear
263 /// at the start of a line, *outside* any code fence.
264 pub(crate) fn parse_text_final(text: &str) -> Option<String> {
265 let outside_fence = strip_code_fences(text);
266
267 for line in outside_fence.lines() {
268 let trimmed = line.trim_start();
269 if trimmed.starts_with("FINAL_VAR(") {
270 // FINAL_VAR can't be resolved from text alone — defer to REPL.
271 continue;
272 }
273 if let Some(rest) = trimmed.strip_prefix("FINAL(") {
274 let inner = rest.trim_end();
275 if let Some(end) = inner.rfind(')') {
276 let value = inner[..end].trim();
277 if !value.is_empty() {
278 return Some(strip_quotes(value));
279 }
280 }
281 }
282 }
283 None
284 }
285
286 fn strip_code_fences(text: &str) -> String {
287 let mut out = String::with_capacity(text.len());
288 let mut in_fence = false;
289 for line in text.lines() {
290 if line.trim_start().starts_with("```") {
291 in_fence = !in_fence;
292 continue;
293 }
294 if !in_fence {
295 out.push_str(line);
296 out.push('\n');
297 }
298 }
299 out
300 }
301
302 fn strip_quotes(s: &str) -> String {
303 let bytes = s.as_bytes();
304 if bytes.len() >= 2
305 && ((bytes[0] == b'"' && bytes[bytes.len() - 1] == b'"')
306 || (bytes[0] == b'\'' && bytes[bytes.len() - 1] == b'\''))
307 {
308 return s[1..s.len() - 1].to_string();
309 }
310 s.to_string()
311 }
312
313 pub(crate) fn truncate_text(text: &str, max_chars: usize) -> String {
314 let count = text.chars().count();
315 if count <= max_chars {
316 return text.to_string();
317 }
318 let take = max_chars.saturating_sub(3);
319 let mut result: String = text.chars().take(take).collect();
320 result.push_str("...");
321 result
322 }
323
324 #[cfg(test)]
325 mod tests {
326 use super::*;
327 use crate::core::engine::tests::rlm_host::{Replies, run_admitted_fixture, run_fixture};
328 use crate::core::events::Event;
329 use crate::llm_client::mock::MockLlmClient;
330 use crate::rlm::bridge::RlmUsageAccumulator;
331 use codewhale_models::{MessageRequest, MessageResponse};
332 use std::path::PathBuf;
333 use std::sync::Arc;
334 use tokio::sync::mpsc;
335
336 /// One model round whose code writes a marker file, run under `gate`.
337 async fn marker_round(
338 gate: Option<crate::tools::codemode::NestedCallGate>,
339 ) -> (RlmTurnResult, bool) {
340 let workspace = tempfile::tempdir().expect("tempdir");
341 let marker = workspace.path().join("rlm-round-executed.txt");
342 let marker_literal = serde_json::to_string(&marker.to_string_lossy()).expect("literal");
343 let mock = Arc::new(MockLlmClient::new(Vec::new()));
344 mock.push_message_response(MessageResponse {
345 id: "mock_gated_rlm".to_string(),
346 r#type: "message".to_string(),
347 role: "assistant".to_string(),
348 content: vec![ContentBlock::Text {
349 text: format!(
350 "```repl\nfrom pathlib import Path\nPath({marker_literal}).write_text('executed')\nFINAL('ran')\n```"
351 ),
352 cache_control: None,
353 }],
354 model: "mock-model".to_string(),
355 stop_reason: Some("end_turn".to_string()),
356 stop_sequence: None,
357 container: None,
358 usage: Usage::default(),
359 });
360 let client: Arc<dyn Replies> = mock;
361 let (tx, _rx) = mpsc::channel(8);
362 let result = run_admitted_fixture(
363 client,
364 "root-model".to_string(),
365 "long context".to_string(),
366 None,
367 "child-model".to_string(),
368 tx,
369 0,
370 RlmUsageAccumulator::new(),
371 tokio::time::Instant::now() + Duration::from_secs(60),
372 gate,
373 )
374 .await;
375 (result, marker.exists())
376 }
377
378 #[tokio::test]
379 async fn code_round_runs_only_when_the_gate_admits_exactly_that_code() {
380 use crate::tools::codemode::{NestedCallGate, NestedCallVerdict, NestedDecision};
381
382 let (result, ran) = marker_round(Some(NestedCallGate::admitting_for_test())).await;
383 assert!(ran, "an admitted round runs: {:?}", result.error);
384 assert_eq!(result.termination, RlmTermination::Final);
385
386 let (result, ran) = marker_round(None).await;
387 assert!(!ran, "no gate, no code");
388 assert_eq!(result.termination, RlmTermination::Error);
389 let error = result.error.expect("refusal is reported");
390 assert!(error.contains("no permission gate"), "{error}");
391
392 let refusing = NestedCallGate::answering_for_test(|_, _| NestedCallVerdict::Refused {
393 error: crate::tools::spec::ToolError::permission_denied("not approved"),
394 decision: NestedDecision::Denied,
395 });
396 let (result, ran) = marker_round(Some(refusing)).await;
397 assert!(!ran, "a refused round does not run");
398 let error = result.error.expect("refusal is reported");
399 assert!(error.contains("not approved"), "{error}");
400
401 // A hook-rewritten or otherwise different admitted input is not
402 // permission for the code the model wrote.
403 let rewriting = NestedCallGate::answering_for_test(|name, _| NestedCallVerdict::Run {
404 name: name.to_string(),
405 input: serde_json::json!({ "code": "print('something else')" }),
406 supports_parallel: false,
407 decision: NestedDecision::Auto,
408 hook_context: None,
409 });
410 let (result, ran) = marker_round(Some(rewriting)).await;
411 assert!(!ran, "a different admitted input does not run this code");
412 assert!(
413 result
414 .error
415 .is_some_and(|error| error.contains("was not this code"))
416 );
417 }
418
419 #[tokio::test]
420 async fn max_tokens_complete_repl_is_not_executed_or_accepted() {
421 let workspace = tempfile::tempdir().expect("tempdir");
422 let marker = workspace.path().join("truncated-repl-executed.txt");
423 let marker_literal = serde_json::to_string(&marker.to_string_lossy())
424 .expect("marker path should serialize as a Python string literal");
425 let partial = format!(
426 "```repl\nfrom pathlib import Path\nPath({marker_literal}).write_text('executed')\nFINAL('partial answer')\n```"
427 );
428 let usage = Usage {
429 input_tokens: 17,
430 output_tokens: 4096,
431 ..Usage::default()
432 };
433 let mock = Arc::new(MockLlmClient::new(Vec::new()));
434 mock.push_message_response(MessageResponse {
435 id: "mock_truncated_rlm".to_string(),
436 r#type: "message".to_string(),
437 role: "assistant".to_string(),
438 content: vec![ContentBlock::Text {
439 text: partial,
440 cache_control: None,
441 }],
442 model: "mock-model".to_string(),
443 stop_reason: Some("max_tokens".to_string()),
444 stop_sequence: None,
445 container: None,
446 usage: usage.clone(),
447 });
448 let client: Arc<dyn Replies> = mock.clone();
449 let (tx, _rx) = mpsc::channel(8);
450
451 let result = run_fixture(
452 client,
453 "root-model".to_string(),
454 "long context".to_string(),
455 None,
456 "child-model".to_string(),
457 tx,
458 0,
459 )
460 .await;
461
462 assert_eq!(result.termination, RlmTermination::Error);
463 assert!(
464 result.answer.is_empty(),
465 "partial FINAL must not be accepted"
466 );
467 let error = result.error.expect("truncation must fail the RLM turn");
468 assert!(error.contains("incomplete"), "{error}");
469 assert!(error.contains("max_tokens"), "{error}");
470 assert_eq!(result.usage, usage, "billed usage must still be charged");
471 assert_eq!(result.routed_usage.len(), 1);
472 assert_eq!(result.routed_usage[0].usage.usage, usage);
473 assert_eq!(result.routed_usage_dropped_records, 0);
474 assert_eq!(mock.call_count(), 1, "truncation must not retry");
475 assert!(
476 !marker.exists(),
477 "complete-looking code from a truncated response must not execute"
478 );
479 }
480
481 #[tokio::test]
482 async fn root_provider_success_without_usage_retains_exact_missing_receipt() {
483 let mock = Arc::new(MockLlmClient::new(Vec::new()));
484 mock.push_message_response(MessageResponse {
485 id: "mock_missing_rlm_usage".to_string(),
486 r#type: "message".to_string(),
487 role: "assistant".to_string(),
488 content: vec![ContentBlock::Text {
489 text: "partial".to_string(),
490 cache_control: None,
491 }],
492 model: "mock-model".to_string(),
493 stop_reason: Some("max_tokens".to_string()),
494 stop_sequence: None,
495 container: None,
496 usage: Usage::default(),
497 });
498 let client: Arc<dyn Replies> = mock;
499 let (tx, _rx) = mpsc::channel(8);
500
501 let result = run_fixture(
502 client,
503 "root-model".to_string(),
504 "long context".to_string(),
505 None,
506 "child-model".to_string(),
507 tx,
508 0,
509 )
510 .await;
511
512 assert_eq!(result.termination, RlmTermination::Error);
513 assert_eq!(result.usage, Usage::default());
514 assert!(result.routed_usage.is_empty());
515 assert_eq!(result.routed_usage_drop_records.len(), 1);
516 assert_eq!(result.routed_usage_dropped_records, 1);
517 assert_eq!(
518 result.routed_usage_drop_records[0].route.model,
519 "root-model"
520 );
521 }
522
523 fn text_response(text: &str) -> MessageResponse {
524 MessageResponse {
525 id: "mock_rlm_round".to_string(),
526 r#type: "message".to_string(),
527 role: "assistant".to_string(),
528 content: vec![ContentBlock::Text {
529 text: text.to_string(),
530 cache_control: None,
531 }],
532 model: "mock-model".to_string(),
533 stop_reason: Some("end_turn".to_string()),
534 stop_sequence: None,
535 container: None,
536 usage: Usage {
537 input_tokens: 5,
538 output_tokens: 5,
539 ..Usage::default()
540 },
541 }
542 }
543
544 /// #6511: an exhausted loop used to return `answer: String::new()` and
545 /// drop the middle of its history after 20 messages.
546 #[tokio::test]
547 async fn exhausted_loop_returns_last_response_and_keeps_whole_history() {
548 let mock = Arc::new(MockLlmClient::new(Vec::new()));
549 for round in 0..MAX_RLM_ITERATIONS {
550 mock.push_message_response(text_response(&format!(
551 "```repl\nprint('round {round}')\n```"
552 )));
553 }
554 let client: Arc<dyn Replies> = mock.clone();
555 let (tx, mut rx) = mpsc::channel(1024);
556
557 let result = run_fixture(
558 client,
559 "root-model".to_string(),
560 "long context".to_string(),
561 None,
562 "child-model".to_string(),
563 tx,
564 0,
565 )
566 .await;
567
568 assert_eq!(result.termination, RlmTermination::Exhausted);
569 assert_eq!(result.iterations, MAX_RLM_ITERATIONS);
570 let last_round = MAX_RLM_ITERATIONS - 1;
571 assert!(
572 result.answer.contains(&format!("round {last_round}")),
573 "exhaustion must surface the last root response, got {:?}",
574 result.answer
575 );
576 assert!(
577 result
578 .error
579 .as_deref()
580 .is_some_and(|e| e.contains("exhausted")),
581 "{:?}",
582 result.error
583 );
584
585 let requests = mock.captured_requests();
586 assert_eq!(requests.len(), MAX_RLM_ITERATIONS as usize);
587 // Initial metadata plus (code, result) per completed round: nothing
588 // from the middle is dropped.
589 let last = requests.last().expect("last root request");
590 assert_eq!(
591 last.messages.len(),
592 1 + 2 * (MAX_RLM_ITERATIONS as usize - 1)
593 );
594
595 let mut statuses = Vec::new();
596 while let Ok(event) = rx.try_recv() {
597 if let Event::Status { message } = event {
598 statuses.push(message);
599 }
600 }
601 assert!(
602 statuses
603 .iter()
604 .any(|line| line.contains("RLM finished: Exhausted")),
605 "the loop must log how it ended: {statuses:#?}"
606 );
607 }
608
609 struct PendingAfterResponses(MockLlmClient, usize);
610
611 impl Replies for PendingAfterResponses {
612 fn effective_route_envelope(
613 &self,
614 model: &str,
615 dispatched_at: chrono::DateTime<chrono::Utc>,
616 ) -> crate::cost_status::EffectiveRouteEnvelope {
617 self.0.effective_route_envelope(model, dispatched_at)
618 }
619
620 fn effective_max_output_tokens(&self, model: &str) -> u32 {
621 self.0.effective_max_output_tokens(model)
622 }
623
624 fn create_message_boxed(
625 &self,
626 request: MessageRequest,
627 ) -> std::pin::Pin<
628 Box<dyn std::future::Future<Output = anyhow::Result<MessageResponse>> + Send + '_>,
629 > {
630 if self.0.call_count() < self.1 {
631 self.0.create_message_boxed(request)
632 } else {
633 Box::pin(std::future::pending())
634 }
635 }
636 }
637
638 /// Checking time only between rounds cannot stop a pending root
639 /// request, a Python block, or an event send within the current round.
640 #[tokio::test]
641 async fn wall_clock_deadline_interrupts_pending_work_and_keeps_partial_result() {
642 for pending_model in [true, false] {
643 let partial = if pending_model {
644 "```repl\nprint(_os.environ['RLM_CONTEXT_FILE'])\n```"
645 } else {
646 "```repl\nimport time\ntime.sleep(60)\nFINAL('too late')\n```"
647 };
648 let mock = MockLlmClient::new(Vec::new());
649 mock.push_message_response(text_response(partial));
650 let client = Arc::new(PendingAfterResponses(mock, 1));
651 let (tx, mut rx) = mpsc::channel(32);
652 let usage = RlmUsageAccumulator::new();
653
654 let result = tokio::time::timeout(
655 Duration::from_secs(5),
656 run_admitted_fixture(
657 client.clone(),
658 "root-model".to_string(),
659 "long context".to_string(),
660 None,
661 "child-model".to_string(),
662 tx,
663 0,
664 usage.clone(),
665 tokio::time::Instant::now() + Duration::from_secs(1),
666 Some(crate::tools::codemode::NestedCallGate::admitting_for_test()),
667 ),
668 )
669 .await
670 .expect("the turn deadline must interrupt in-flight work");
671
672 assert_eq!(result.termination, RlmTermination::Error);
673 assert_eq!(result.answer, partial);
674 assert_eq!(result.iterations, if pending_model { 2 } else { 1 });
675 assert!(
676 result
677 .error
678 .as_deref()
679 .unwrap()
680 .contains("wall-clock deadline")
681 );
682 assert_eq!(result.usage, text_response(partial).usage);
683 assert_eq!(client.0.call_count(), 1);
684 let snapshot = usage.snapshot().await;
685 assert_eq!(snapshot.records.len(), 1, "keep completed usage");
686 assert_eq!(snapshot.dropped_records, u64::from(pending_model));
687 if pending_model {
688 assert_eq!(result.trace.len(), 1, "keep completed rounds");
689 assert!(!result.trace[0].had_error);
690 let context_path = PathBuf::from(result.trace[0].stdout_preview.trim());
691 assert!(
692 context_path.is_absolute(),
693 "REPL must report its context path"
694 );
695 assert!(!context_path.exists(), "timeout must clean up context");
696 }
697 let mut saw_timeout = false;
698 while let Ok(event) = rx.try_recv() {
699 if let Event::Status { message } = event {
700 saw_timeout |= message.contains("RLM finished: Error")
701 && message.contains("wall-clock deadline");
702 }
703 }
704 assert!(saw_timeout, "timeout must record how the turn ended");
705 let terminal_receipts: Vec<_> = snapshot
706 .nested_events
707 .iter()
708 .filter(|event| {
709 event["kind"] == "status"
710 && event["content"].as_str().is_some_and(|content| {
711 content.contains("RLM finished: Error")
712 && content.contains("wall-clock deadline")
713 })
714 })
715 .collect();
716 assert_eq!(terminal_receipts.len(), 1, "one canonical terminal receipt");
717 }
718 }
719
720 #[tokio::test]
721 async fn wall_clock_deadline_returns_when_event_stream_is_full() {
722 // An empty response queue fails immediately and cannot test the deadline.
723 // Keep the admitted provider future pending while the host channel is full.
724 let client = Arc::new(PendingAfterResponses(MockLlmClient::new(Vec::new()), 0));
725 let usage = RlmUsageAccumulator::new();
726 let (tx, _rx) = mpsc::channel(1);
727 tx.try_send(Event::status("fixture occupies the caller event channel"))
728 .unwrap();
729 let result = tokio::time::timeout(
730 Duration::from_secs(5),
731 run_admitted_fixture(
732 client.clone(),
733 "root-model".to_string(),
734 "long context".to_string(),
735 None,
736 "child-model".to_string(),
737 tx,
738 0,
739 usage.clone(),
740 tokio::time::Instant::now() + Duration::from_secs(1),
741 Some(crate::tools::codemode::NestedCallGate::admitting_for_test()),
742 ),
743 )
744 .await
745 .expect("deadline hand-back must not wait on a full event stream");
746
747 assert_eq!(result.termination, RlmTermination::Error);
748 assert!(
749 result
750 .error
751 .as_deref()
752 .unwrap()
753 .contains("wall-clock deadline")
754 );
755 assert_eq!(client.0.call_count(), 0, "no scripted response completed");
756 let snapshot = usage.snapshot().await;
757 assert_eq!(snapshot.dropped_records, 1, "one pending admitted request");
758 assert_eq!(
759 snapshot
760 .nested_events
761 .iter()
762 .filter(|event| {
763 event["kind"] == "status"
764 && event["content"].as_str().is_some_and(|content| {
765 content.contains("RLM finished: Error")
766 && content.contains("wall-clock deadline")
767 })
768 })
769 .count(),
770 1,
771 "full host channel must retain one canonical terminal receipt"
772 );
773 }
774
775 #[test]
776 fn empty_answer_is_never_returned_without_an_error() {
777 let empty = |termination| RlmTurnResult {
778 answer: " ".to_string(),
779 iterations: 2,
780 duration: Duration::ZERO,
781 error: None,
782 usage: Usage::default(),
783 routed_usage: Vec::new(),
784 routed_usage_drop_records: Vec::new(),
785 routed_usage_dropped_records: 0,
786 termination,
787 trace: Vec::new(),
788 total_rpcs: 0,
789 };
790 let error = empty(RlmTermination::NoCode)
791 .require_answer_or_error()
792 .error
793 .expect("empty answer needs a reason");
794 assert!(error.contains("empty answer"), "{error}");
795 assert!(error.contains("NoCode"), "{error}");
796
797 // `FINAL("")` is an answer the model chose; callers keep getting "".
798 let deliberate = empty(RlmTermination::Final).require_answer_or_error();
799 assert!(deliberate.error.is_none(), "{:?}", deliberate.error);
800 }
801
802 #[test]
803 fn extract_repl_code_finds_simple_block() {
804 let text = "Here:\n```repl\nprint('hi')\n```\nEnd.";
805 let code = extract_repl_code(text).unwrap();
806 assert_eq!(code, "print('hi')");
807 }
808
809 #[test]
810 fn extract_repl_code_falls_back_to_python_marker() {
811 let text = "Code:\n```python\nx = 1 + 2\n```";
812 let code = extract_repl_code(text).unwrap();
813 assert_eq!(code, "x = 1 + 2");
814 }
815
816 #[test]
817 fn extract_repl_code_returns_none_when_missing() {
818 assert!(extract_repl_code("Just text.").is_none());
819 }
820
821 #[test]
822 fn extract_repl_code_returns_none_on_empty_block() {
823 assert!(extract_repl_code("```repl\n\n```").is_none());
824 }
825
826 #[test]
827 fn extract_repl_code_handles_multiple_blocks() {
828 let text = "```repl\na=1\n```\n```repl\nb=2\n```";
829 let code = extract_repl_code(text).unwrap();
830 assert_eq!(code, "a=1");
831 }
832
833 #[test]
834 fn extract_repl_code_ignores_other_fences() {
835 let text = "```\nfoo\n```\n```repl\nreal_code()\n```";
836 let code = extract_repl_code(text).unwrap();
837 assert_eq!(code, "real_code()");
838 }
839
840 #[test]
841 fn extract_repl_code_requires_a_line_anchored_fence() {
842 assert!(extract_repl_code("see ```python\nx = 1\n```").is_none());
843 assert!(extract_repl_code("a ```repl mention\nx = 1\n```").is_none());
844 assert!(extract_repl_code("```python-output\nx = 1\n```").is_none());
845 let text = "prose ```py inline\n```py\nreal()\n```";
846 assert_eq!(extract_repl_code(text).as_deref(), Some("real()"));
847 }
848
849 #[test]
850 fn parse_text_final_extracts_simple_value() {
851 let text = "OK.\nFINAL(42)\nThanks.";
852 assert_eq!(parse_text_final(text).as_deref(), Some("42"));
853 }
854
855 #[test]
856 fn parse_text_final_strips_quotes() {
857 let text = "FINAL(\"the answer is yes\")";
858 assert_eq!(parse_text_final(text).as_deref(), Some("the answer is yes"));
859 }
860
861 #[test]
862 fn parse_text_final_ignores_inside_code_fence() {
863 let text =
864 "Some prose.\n```repl\n# Note: when ready, call FINAL(value)\nx = 1\n```\nMore prose.";
865 assert!(parse_text_final(text).is_none());
866 }
867
868 #[test]
869 fn parse_text_final_returns_none_when_absent() {
870 assert!(parse_text_final("just talking, no final.").is_none());
871 }
872
873 #[test]
874 fn build_metadata_contains_key_information() {
875 let msg = build_metadata_message("Hello, world!", None, 0, None, None);
876 let text = extract_text_blocks(&msg.content);
877 assert!(text.contains("context"));
878 assert!(text.contains("Hello, world!"));
879 assert!(text.contains("round 0"));
880 assert!(text.contains("llm_query"));
881 assert!(text.contains("rlm_query"));
882 assert!(text.contains("FINAL"));
883 }
884
885 #[test]
886 fn build_metadata_truncates_long_context_without_leaking_tail() {
887 let secret_tail = "DO_NOT_LEAK_CONTEXT_TAIL";
888 let prompt = format!("{}{}", "a".repeat(PROMPT_PREVIEW_LEN + 100), secret_tail);
889 let msg = build_metadata_message(&prompt, None, 0, None, None);
890 let text = extract_text_blocks(&msg.content);
891
892 assert!(text.contains(&format!("- Length: {} chars", prompt.chars().count())));
893 assert!(text.contains("- Preview: \""));
894 assert!(text.contains("..."));
895 assert!(
896 !text.contains(secret_tail),
897 "metadata leaked the non-preview tail of context"
898 );
899 }
900
901 #[tokio::test]
902 async fn build_root_request_keeps_context_tail_out_of_root_payload() {
903 let secret_tail = "DO_NOT_LEAK_ROOT_REQUEST";
904 let prompt = format!("{}{}", "a".repeat(PROMPT_PREVIEW_LEN + 100), secret_tail);
905 let mock = Arc::new(MockLlmClient::new(Vec::new()));
906 mock.push_message_response(MessageResponse {
907 id: "context-metadata".into(),
908 r#type: "message".into(),
909 role: "assistant".into(),
910 content: vec![ContentBlock::Text {
911 text: "```repl\nFINAL('metadata only')\n```".into(),
912 cache_control: None,
913 }],
914 model: "root-model".into(),
915 stop_reason: Some("end_turn".into()),
916 stop_sequence: None,
917 container: None,
918 usage: Usage::default(),
919 });
920 let (tx, _rx) = mpsc::channel(64);
921 let result = run_fixture(
922 mock.clone(),
923 "root-model".into(),
924 prompt.clone(),
925 Some("answer from the long context".into()),
926 "child-model".into(),
927 tx,
928 0,
929 )
930 .await;
931 assert_eq!(
932 result.termination,
933 RlmTermination::Final,
934 "{:?}",
935 result.error
936 );
937 let requests = mock.captured_requests();
938 assert_eq!(requests.len(), 1);
939 let request = &requests[0];
940 let payload = serde_json::to_string(request).expect("actual Core request serializes");
941 assert_eq!(request.temperature, None);
942 assert_eq!(request.top_p, None);
943 assert!(payload.contains(&format!("- Length: {} chars", prompt.chars().count())));
944 assert!(
945 !payload.contains(secret_tail),
946 "Core request leaked the non-preview context tail"
947 );
948 assert!(payload.contains("Captured operator Core policy"));
949 assert!(payload.contains("answer from the long context"));
950 }
951
952 #[test]
953 fn build_metadata_with_iteration_shows_previous_code() {
954 let msg = build_metadata_message("Test prompt", None, 3, Some("print('hi')"), Some("hi"));
955 let text = extract_text_blocks(&msg.content);
956 assert!(text.contains("round 3"));
957 assert!(text.contains("print('hi')"));
958 assert!(text.contains("hi"));
959 }
960
961 #[test]
962 fn build_metadata_includes_root_prompt() {
963 let msg = build_metadata_message(
964 "long context",
965 Some("Summarize the security model"),
966 1,
967 Some("# noop"),
968 Some("ok"),
969 );
970 let text = extract_text_blocks(&msg.content);
971 assert!(text.contains("Original task"));
972 assert!(text.contains("Summarize the security model"));
973 }
974
975 #[test]
976 fn truncate_text_leaves_short_alone() {
977 assert_eq!(truncate_text("hello", 100), "hello");
978 }
979
980 #[test]
981 fn truncate_text_shortens_long_text() {
982 let long = "a".repeat(1000);
983 let truncated = truncate_text(&long, 10);
984 assert_eq!(truncated.chars().count(), 10);
985 assert!(truncated.ends_with("..."));
986 }
987
988 #[test]
989 fn truncate_text_is_unicode_safe() {
990 let s = "日本語テスト";
991 let out = truncate_text(s, 4);
992 assert_eq!(out.chars().count(), 4);
993 assert!(out.ends_with("..."));
994 assert!(std::str::from_utf8(out.as_bytes()).is_ok());
995 }
996
997 #[test]
998 fn extract_text_blocks_joins_text() {
999 let blocks = vec![
1000 ContentBlock::Text {
1001 text: "first".to_string(),
1002 cache_control: None,
1003 },
1004 ContentBlock::Thinking {
1005 signature: None,
1006 state: None,
1007 thinking: "skip".to_string(),
1008 },
1009 ContentBlock::Text {
1010 text: "second".to_string(),
1011 cache_control: None,
1012 },
1013 ];
1014 assert_eq!(extract_text_blocks(&blocks), "first\nsecond");
1015 }
1016
1017 #[test]
1018 fn metadata_msg_role_is_user() {
1019 let msg = build_metadata_message("test", None, 0, None, None);
1020 assert_eq!(msg.role, "user");
1021 }
1022
1023 #[test]
1024 fn summarize_code_keeps_short() {
1025 assert_eq!(summarize_code("a\nb\nc"), "a\nb\nc");
1026 }
1027
1028 #[test]
1029 fn summarize_code_compresses_long() {
1030 let lines: Vec<String> = (0..20).map(|i| format!("line{i}")).collect();
1031 let code = lines.join("\n");
1032 let s = summarize_code(&code);
1033 assert!(s.starts_with("20 lines:"));
1034 assert!(s.contains("line0"));
1035 assert!(s.contains("line19"));
1036 assert!(s.contains("…"));
1037 }
1038 }
1039
1039 lines RUST