| 1 | //! Actual Core fixtures for the migrated RLM contract. Replies are adapted to |
| 2 | //! the existing stream decoder; there is no fixture model/code loop. |
| 3 | use super::*; |
| 4 | use crate::core::engine::rlm_host::{CapturedRlmCaller, RlmInvocation, RlmMode}; |
| 5 | use crate::llm_client::LlmClient; |
| 6 | use crate::rlm::bridge::{RlmBridge, RlmUsageAccumulator}; |
| 7 | use codewhale_models::{MessageRequest, MessageResponse, StreamEvent}; |
| 8 | use std::future::Future; |
| 9 | use std::pin::Pin; |
| 10 | |
| 11 | pub(crate) trait Replies: Send + Sync { |
| 12 | fn effective_route_envelope( |
| 13 | &self, |
| 14 | model: &str, |
| 15 | at: chrono::DateTime<chrono::Utc>, |
| 16 | ) -> crate::cost_status::EffectiveRouteEnvelope; |
| 17 | fn effective_max_output_tokens(&self, model: &str) -> u32; |
| 18 | fn create_message_boxed( |
| 19 | &self, |
| 20 | request: MessageRequest, |
| 21 | ) -> Pin<Box<dyn Future<Output = anyhow::Result<MessageResponse>> + Send + '_>>; |
| 22 | } |
| 23 | impl<T: LlmClient + Send + Sync> Replies for T { |
| 24 | fn effective_route_envelope( |
| 25 | &self, |
| 26 | model: &str, |
| 27 | at: chrono::DateTime<chrono::Utc>, |
| 28 | ) -> crate::cost_status::EffectiveRouteEnvelope { |
| 29 | LlmClient::effective_route_envelope(self, model, at) |
| 30 | } |
| 31 | fn effective_max_output_tokens(&self, model: &str) -> u32 { |
| 32 | LlmClient::effective_max_output_tokens(self, model) |
| 33 | } |
| 34 | fn create_message_boxed( |
| 35 | &self, |
| 36 | request: MessageRequest, |
| 37 | ) -> Pin<Box<dyn Future<Output = anyhow::Result<MessageResponse>> + Send + '_>> { |
| 38 | Box::pin(LlmClient::create_message(self, request)) |
| 39 | } |
| 40 | } |
| 41 | struct ReplyStream { |
| 42 | replies: Arc<dyn Replies>, |
| 43 | model: String, |
| 44 | } |
| 45 | #[async_trait::async_trait] |
| 46 | impl crate::core::model_client::ModelClient for ReplyStream { |
| 47 | fn provider_name(&self) -> &str { |
| 48 | "custom" |
| 49 | } |
| 50 | fn model(&self) -> &str { |
| 51 | &self.model |
| 52 | } |
| 53 | fn effective_route_envelope( |
| 54 | &self, |
| 55 | model: &str, |
| 56 | at: chrono::DateTime<chrono::Utc>, |
| 57 | ) -> crate::cost_status::EffectiveRouteEnvelope { |
| 58 | self.replies.effective_route_envelope(model, at) |
| 59 | } |
| 60 | fn effective_max_output_tokens(&self, model: &str) -> u32 { |
| 61 | self.replies.effective_max_output_tokens(model) |
| 62 | } |
| 63 | async fn create_message(&self, request: MessageRequest) -> anyhow::Result<MessageResponse> { |
| 64 | self.replies.create_message_boxed(request).await |
| 65 | } |
| 66 | async fn create_message_stream( |
| 67 | &self, |
| 68 | request: MessageRequest, |
| 69 | ) -> anyhow::Result<crate::llm_client::StreamEventBox> { |
| 70 | let response = self.replies.create_message_boxed(request).await?; |
| 71 | let mut start = response.clone(); |
| 72 | start.content.clear(); |
| 73 | let mut events = vec![StreamEvent::MessageStart { message: start }]; |
| 74 | for (index, block) in response.content.iter().enumerate() { |
| 75 | let index = u32::try_from(index).expect("fixture block index"); |
| 76 | match block { |
| 77 | ContentBlock::Text { text, .. } => { |
| 78 | events.push(StreamEvent::ContentBlockStart { |
| 79 | index, |
| 80 | content_block: codewhale_models::ContentBlockStart::Text { |
| 81 | text: String::new(), |
| 82 | }, |
| 83 | }); |
| 84 | events.push(StreamEvent::ContentBlockDelta { |
| 85 | index, |
| 86 | delta: codewhale_models::Delta::TextDelta { text: text.clone() }, |
| 87 | }); |
| 88 | } |
| 89 | _ => { |
| 90 | events.push(StreamEvent::ContentBlockStart { |
| 91 | index, |
| 92 | content_block: serde_json::from_value(serde_json::to_value(block)?)?, |
| 93 | }); |
| 94 | } |
| 95 | } |
| 96 | events.push(StreamEvent::ContentBlockStop { index }); |
| 97 | } |
| 98 | events.push(StreamEvent::MessageDelta { |
| 99 | delta: codewhale_models::MessageDelta { |
| 100 | stop_reason: response.stop_reason, |
| 101 | stop_sequence: response.stop_sequence, |
| 102 | }, |
| 103 | usage: Some(response.usage), |
| 104 | }); |
| 105 | events.push(StreamEvent::MessageStop); |
| 106 | Ok(Box::pin(futures_util::stream::iter( |
| 107 | events.into_iter().map(Ok), |
| 108 | ))) |
| 109 | } |
| 110 | async fn health_check(&self) -> anyhow::Result<bool> { |
| 111 | Ok(true) |
| 112 | } |
| 113 | } |
| 114 | |
| 115 | pub(crate) fn fixture_config(model: &str) -> Config { |
| 116 | let mut config: Config = toml::from_str(&format!("provider = 'rlm-fixture'\n[providers.rlm-fixture]\nkind = 'openai-compatible'\napi_key = 'owned-fixture-key'\nbase_url = 'http://127.0.0.1:9/v1'\nmodel = {}\ncontext_window = 1000000\n", serde_json::to_string(model).unwrap())).expect("fixture Config"); |
| 117 | config.default_text_model = Some(model.to_string()); |
| 118 | config |
| 119 | } |
| 120 | |
| 121 | pub(crate) fn install_fixture_route(engine: &mut Engine) { |
| 122 | let identity = engine |
| 123 | .api_config |
| 124 | .active_provider_identity() |
| 125 | .expect("fixture admitted exact id"); |
| 126 | let resolved = resolve_runtime_route_for_identity( |
| 127 | &engine.api_config, |
| 128 | &identity, |
| 129 | Some(&engine.session.model), |
| 130 | ) |
| 131 | .expect("fixture route"); |
| 132 | let client = CodewhaleClient::new(&resolved.config) |
| 133 | .expect("fixture captured client, never called unless loopback transport was supplied"); |
| 134 | engine.install_validated_runtime_route(ValidatedRuntimeRoute { |
| 135 | identity: resolved.identity, |
| 136 | candidate: resolved.candidate, |
| 137 | config: resolved.config, |
| 138 | model: resolved.model, |
| 139 | context_window: resolved.context_window, |
| 140 | client, |
| 141 | }); |
| 142 | } |
| 143 | |
| 144 | /// Use the actual configured Engine's posture, route, registry and services. |
| 145 | /// No client-only or ambient default configuration grants fixture authority. |
| 146 | pub(crate) fn admitted_context(engine: &Engine, turn_id: &str) -> ToolContext { |
| 147 | let posture = engine.applied_runtime_authority(); |
| 148 | let authority = TurnAuthority::from_effective_fields( |
| 149 | posture.mode, |
| 150 | posture.allow_shell, |
| 151 | posture.trust_mode, |
| 152 | posture.auto_approve, |
| 153 | posture.approval_mode, |
| 154 | ); |
| 155 | let route = TurnRouteContext { |
| 156 | provider: engine.api_provider, |
| 157 | model: engine.session.model.clone(), |
| 158 | capabilities: engine.active_route_capabilities, |
| 159 | limits: engine.active_route_limits, |
| 160 | client: engine.codewhale_client.clone(), |
| 161 | api_config: Box::new(engine.api_config.clone()), |
| 162 | locale_tag: "en".into(), |
| 163 | role_models: engine.subagent_role_models(), |
| 164 | auto_model: false, |
| 165 | reasoning_effort: None, |
| 166 | reasoning_effort_auto: false, |
| 167 | }; |
| 168 | let mut context = engine.build_tool_context_for_turn(&authority, &route); |
| 169 | let caller = CapturedRlmCaller::capture(engine, &authority, &route, &context, turn_id) |
| 170 | .expect("actual Core admission"); |
| 171 | context.rlm_caller = Some(Arc::new(caller)); |
| 172 | context |
| 173 | } |
| 174 | |
| 175 | pub(crate) fn caller_for(replies: Arc<dyn Replies>, model: &str) -> CapturedRlmCaller { |
| 176 | let workspace = Arc::new(tempfile::tempdir().expect("RLM fixture workspace")); |
| 177 | let config = fixture_config(model); |
| 178 | let (mut engine, _handle) = Engine::new_with_model_client( |
| 179 | EngineConfig { |
| 180 | workspace: workspace.path().to_path_buf(), |
| 181 | model: model.to_string(), |
| 182 | snapshots_enabled: false, |
| 183 | memory_enabled: false, |
| 184 | subagents_enabled: false, |
| 185 | terminal_chrome_enabled: false, |
| 186 | turn_wall_clock: Duration::from_secs(60), |
| 187 | ..EngineConfig::default() |
| 188 | }, |
| 189 | &config, |
| 190 | Arc::new(ReplyStream { |
| 191 | replies, |
| 192 | model: model.to_string(), |
| 193 | }), |
| 194 | ); |
| 195 | install_fixture_route(&mut engine); |
| 196 | engine.session.system_prompt = Some(SystemPrompt::Text( |
| 197 | "Captured operator Core policy. Task guidance is additive.".into(), |
| 198 | )); |
| 199 | let context = admitted_context(&engine, "fixture-origin-turn"); |
| 200 | let mut caller = context |
| 201 | .rlm_caller |
| 202 | .as_ref() |
| 203 | .expect("captured receipt") |
| 204 | .as_ref() |
| 205 | .clone(); |
| 206 | caller.hold_fixture_workspace(workspace); |
| 207 | caller |
| 208 | } |
| 209 | |
| 210 | /// Admit a loopback LlmClient through a real Engine with the supplied Config |
| 211 | /// and tool context services, rather than granting authority to its URL alone. |
| 212 | pub(crate) fn context_for_replies( |
| 213 | config: &Config, |
| 214 | model: &str, |
| 215 | context: &ToolContext, |
| 216 | replies: Arc<dyn Replies>, |
| 217 | ) -> ToolContext { |
| 218 | let (mut engine, _handle) = Engine::new_with_model_client( |
| 219 | EngineConfig { |
| 220 | workspace: context.workspace.clone(), |
| 221 | model: model.to_string(), |
| 222 | session_id: Some(context.state_namespace.clone()), |
| 223 | runtime_services: context.runtime.clone(), |
| 224 | plugin_registry: context.plugin_registry.clone(), |
| 225 | snapshots_enabled: false, |
| 226 | memory_enabled: false, |
| 227 | subagents_enabled: false, |
| 228 | terminal_chrome_enabled: false, |
| 229 | turn_wall_clock: Duration::from_secs(60), |
| 230 | ..EngineConfig::default() |
| 231 | }, |
| 232 | config, |
| 233 | Arc::new(ReplyStream { |
| 234 | replies, |
| 235 | model: model.to_string(), |
| 236 | }), |
| 237 | ); |
| 238 | install_fixture_route(&mut engine); |
| 239 | let mut admitted = admitted_context(&engine, "fixture-tool-origin-turn"); |
| 240 | admitted.nested_call_gate = context.nested_call_gate.clone(); |
| 241 | admitted |
| 242 | } |
| 243 | |
| 244 | pub(crate) struct BridgeFixture { |
| 245 | caller: CapturedRlmCaller, |
| 246 | depth: u32, |
| 247 | pub(crate) usage: RlmUsageAccumulator, |
| 248 | deadline: Option<tokio::time::Instant>, |
| 249 | gate: Option<crate::tools::codemode::NestedCallGate>, |
| 250 | events: Option<mpsc::Sender<Event>>, |
| 251 | } |
| 252 | impl BridgeFixture { |
| 253 | pub(crate) fn new(replies: Arc<dyn Replies>, model: String, depth: u32) -> Self { |
| 254 | Self::with_usage_accumulator(replies, model, depth, RlmUsageAccumulator::new()) |
| 255 | } |
| 256 | pub(crate) fn with_usage_accumulator( |
| 257 | replies: Arc<dyn Replies>, |
| 258 | model: String, |
| 259 | depth: u32, |
| 260 | usage: RlmUsageAccumulator, |
| 261 | ) -> Self { |
| 262 | Self { |
| 263 | caller: caller_for(replies, &model), |
| 264 | depth, |
| 265 | usage, |
| 266 | deadline: None, |
| 267 | gate: None, |
| 268 | events: None, |
| 269 | } |
| 270 | } |
| 271 | fn bridge(&self) -> RlmBridge<'_> { |
| 272 | let mut bridge = RlmBridge::with_usage_accumulator( |
| 273 | &self.caller, |
| 274 | self.depth, |
| 275 | Duration::from_secs(600), |
| 276 | self.usage.clone(), |
| 277 | ) |
| 278 | .with_deadline(self.deadline) |
| 279 | .with_gate(self.gate.clone()); |
| 280 | if let Some(events) = &self.events { |
| 281 | bridge = bridge.with_events(events.clone()); |
| 282 | } |
| 283 | bridge |
| 284 | } |
| 285 | pub(crate) fn with_deadline(mut self, deadline: Option<tokio::time::Instant>) -> Self { |
| 286 | self.deadline = deadline; |
| 287 | self |
| 288 | } |
| 289 | pub(crate) fn with_gate( |
| 290 | mut self, |
| 291 | gate: Option<crate::tools::codemode::NestedCallGate>, |
| 292 | ) -> Self { |
| 293 | self.gate = gate; |
| 294 | self |
| 295 | } |
| 296 | pub(crate) fn with_events(mut self, events: mpsc::Sender<Event>) -> Self { |
| 297 | self.events = Some(events); |
| 298 | self |
| 299 | } |
| 300 | pub(crate) async fn usage_snapshot(&self) -> crate::rlm::bridge::RlmUsageSnapshot { |
| 301 | self.usage.snapshot().await |
| 302 | } |
| 303 | pub(crate) async fn dispatch_rlm( |
| 304 | &self, |
| 305 | prompt: String, |
| 306 | model: Option<String>, |
| 307 | ) -> crate::repl::runtime::SingleResp { |
| 308 | self.bridge().dispatch_rlm(prompt, model).await |
| 309 | } |
| 310 | } |
| 311 | impl crate::repl::runtime::RpcDispatcher for BridgeFixture { |
| 312 | fn dispatch<'a>( |
| 313 | &'a self, |
| 314 | request: crate::repl::runtime::RpcRequest, |
| 315 | ) -> Pin<Box<dyn Future<Output = crate::repl::runtime::RpcResponse> + Send + 'a>> { |
| 316 | Box::pin(async move { |
| 317 | crate::repl::runtime::RpcDispatcher::dispatch(&self.bridge(), request).await |
| 318 | }) |
| 319 | } |
| 320 | } |
| 321 | |
| 322 | #[allow(clippy::too_many_arguments)] |
| 323 | pub(crate) async fn run_admitted_fixture( |
| 324 | replies: Arc<dyn Replies>, |
| 325 | model: String, |
| 326 | prompt: String, |
| 327 | task: Option<String>, |
| 328 | _old_child_model: String, |
| 329 | events: mpsc::Sender<Event>, |
| 330 | depth: u32, |
| 331 | usage: RlmUsageAccumulator, |
| 332 | deadline: tokio::time::Instant, |
| 333 | gate: Option<crate::tools::codemode::NestedCallGate>, |
| 334 | ) -> crate::rlm::turn::RlmTurnResult { |
| 335 | let caller = caller_for(replies, &model); |
| 336 | caller |
| 337 | .dispatch(RlmInvocation { |
| 338 | prompt, |
| 339 | mode: RlmMode::Recursive { |
| 340 | depth_remaining: depth, |
| 341 | }, |
| 342 | max_tokens: None, |
| 343 | task_instructions: task, |
| 344 | deadline, |
| 345 | gate, |
| 346 | events: Some(events), |
| 347 | usage, |
| 348 | }) |
| 349 | .await |
| 350 | } |
| 351 | #[allow(clippy::too_many_arguments)] |
| 352 | pub(crate) async fn run_fixture( |
| 353 | replies: Arc<dyn Replies>, |
| 354 | model: String, |
| 355 | prompt: String, |
| 356 | task: Option<String>, |
| 357 | old_child_model: String, |
| 358 | events: mpsc::Sender<Event>, |
| 359 | depth: u32, |
| 360 | ) -> crate::rlm::turn::RlmTurnResult { |
| 361 | run_admitted_fixture( |
| 362 | replies, |
| 363 | model, |
| 364 | prompt, |
| 365 | task, |
| 366 | old_child_model, |
| 367 | events, |
| 368 | depth, |
| 369 | RlmUsageAccumulator::new(), |
| 370 | tokio::time::Instant::now() + Duration::from_secs(60), |
| 371 | Some(crate::tools::codemode::NestedCallGate::admitting_for_test()), |
| 372 | ) |
| 373 | .await |
| 374 | } |
| 375 | |
| 376 | #[tokio::test] |
| 377 | async fn captured_identity_cancel_and_task_bounds_refuse_before_any_provider() { |
| 378 | let workspace = tempfile::tempdir().unwrap(); |
| 379 | let config = fixture_config("captured-model"); |
| 380 | let mock = Arc::new(crate::llm_client::mock::MockLlmClient::new(Vec::new())); |
| 381 | let context = context_for_replies( |
| 382 | &config, |
| 383 | "captured-model", |
| 384 | &ToolContext::new(workspace.path()), |
| 385 | mock.clone(), |
| 386 | ); |
| 387 | let caller = context.rlm_caller.as_ref().unwrap(); |
| 388 | for altered in [ |
| 389 | { |
| 390 | let mut c = context.clone(); |
| 391 | c.workspace = workspace.path().join("other"); |
| 392 | c |
| 393 | }, |
| 394 | { |
| 395 | let mut c = context.clone(); |
| 396 | c.state_namespace.push_str("-other"); |
| 397 | c |
| 398 | }, |
| 399 | { |
| 400 | let mut c = context.clone(); |
| 401 | c.owner_agent_id = Some("foreign-child".into()); |
| 402 | c |
| 403 | }, |
| 404 | ] { |
| 405 | assert!( |
| 406 | caller.validate_context(&altered).is_err(), |
| 407 | "caller mutation cannot grant a nested route" |
| 408 | ); |
| 409 | } |
| 410 | let zero_tokens = caller |
| 411 | .dispatch(RlmInvocation { |
| 412 | prompt: "no effect".into(), |
| 413 | mode: RlmMode::Completion, |
| 414 | max_tokens: Some(0), |
| 415 | task_instructions: None, |
| 416 | deadline: caller.deadline(), |
| 417 | gate: None, |
| 418 | events: None, |
| 419 | usage: RlmUsageAccumulator::new(), |
| 420 | }) |
| 421 | .await; |
| 422 | assert!( |
| 423 | zero_tokens |
| 424 | .error |
| 425 | .unwrap() |
| 426 | .contains("max_tokens must be greater than zero") |
| 427 | ); |
| 428 | assert_eq!(mock.call_count(), 0); |
| 429 | let rejected = caller |
| 430 | .dispatch(RlmInvocation { |
| 431 | prompt: "no effect".into(), |
| 432 | mode: RlmMode::Completion, |
| 433 | max_tokens: None, |
| 434 | task_instructions: Some("x".repeat(4097)), |
| 435 | deadline: caller.deadline(), |
| 436 | gate: None, |
| 437 | events: None, |
| 438 | usage: RlmUsageAccumulator::new(), |
| 439 | }) |
| 440 | .await; |
| 441 | assert!(rejected.error.unwrap().contains("4096 bytes")); |
| 442 | context.cancel_token.as_ref().unwrap().cancel(); |
| 443 | let cancelled = caller |
| 444 | .dispatch(RlmInvocation { |
| 445 | prompt: "no effect".into(), |
| 446 | mode: RlmMode::Completion, |
| 447 | max_tokens: None, |
| 448 | task_instructions: None, |
| 449 | deadline: caller.deadline(), |
| 450 | gate: None, |
| 451 | events: None, |
| 452 | usage: RlmUsageAccumulator::new(), |
| 453 | }) |
| 454 | .await; |
| 455 | assert!( |
| 456 | cancelled |
| 457 | .error |
| 458 | .unwrap() |
| 459 | .contains("originating turn is cancelled") |
| 460 | ); |
| 461 | assert_eq!(mock.call_count(), 0); |
| 462 | } |
| 463 | |
| 464 | #[tokio::test] |
| 465 | async fn dropping_an_actual_dispatched_core_call_keeps_unknown_usage_without_retaining_authority() { |
| 466 | struct Pending; |
| 467 | impl Replies for Pending { |
| 468 | fn effective_route_envelope( |
| 469 | &self, |
| 470 | model: &str, |
| 471 | at: chrono::DateTime<chrono::Utc>, |
| 472 | ) -> crate::cost_status::EffectiveRouteEnvelope { |
| 473 | crate::cost_status::EffectiveRouteEnvelope::capture_observed( |
| 474 | crate::config::ProviderKind::Custom, |
| 475 | "rlm-fixture", |
| 476 | model, |
| 477 | Some("http://127.0.0.1:9/v1"), |
| 478 | at, |
| 479 | ) |
| 480 | } |
| 481 | fn effective_max_output_tokens(&self, _: &str) -> u32 { |
| 482 | 8192 |
| 483 | } |
| 484 | fn create_message_boxed( |
| 485 | &self, |
| 486 | _: MessageRequest, |
| 487 | ) -> Pin<Box<dyn Future<Output = anyhow::Result<MessageResponse>> + Send + '_>> { |
| 488 | Box::pin(std::future::pending()) |
| 489 | } |
| 490 | } |
| 491 | let caller = Arc::new(caller_for(Arc::new(Pending), "captured-model")); |
| 492 | let weak_caller = Arc::downgrade(&caller); |
| 493 | let usage = RlmUsageAccumulator::new(); |
| 494 | let result = caller |
| 495 | .dispatch(RlmInvocation { |
| 496 | prompt: "pending".into(), |
| 497 | mode: RlmMode::Completion, |
| 498 | max_tokens: None, |
| 499 | task_instructions: None, |
| 500 | deadline: tokio::time::Instant::now() + Duration::from_millis(100), |
| 501 | gate: None, |
| 502 | events: None, |
| 503 | usage: usage.clone(), |
| 504 | }) |
| 505 | .await; |
| 506 | assert!( |
| 507 | result |
| 508 | .error |
| 509 | .as_deref() |
| 510 | .unwrap() |
| 511 | .contains("wall-clock deadline") |
| 512 | ); |
| 513 | let snapshot = usage.snapshot().await; |
| 514 | assert_eq!(snapshot.records.len(), 0); |
| 515 | assert_eq!(snapshot.drop_records.len(), 1); |
| 516 | assert_eq!( |
| 517 | snapshot.drop_records[0].reason, |
| 518 | crate::cost_status::RuntimeUsageMissingReason::RequestOutcomeUnknown |
| 519 | ); |
| 520 | assert_eq!(snapshot.drop_records[0].route.model, "captured-model"); |
| 521 | assert_eq!(snapshot.dropped_records, 1); |
| 522 | assert_eq!( |
| 523 | usage.snapshot().await.drop_records, |
| 524 | snapshot.drop_records, |
| 525 | "inspection never republishes or loses the Drop receipt" |
| 526 | ); |
| 527 | drop(caller); |
| 528 | assert!( |
| 529 | weak_caller.upgrade().is_none(), |
| 530 | "dropped nested work retains no caller authority" |
| 531 | ); |
| 532 | } |
| 533 | |
| 534 | #[tokio::test] |
| 535 | async fn full_caller_channel_never_delays_original_cancel_or_loses_usage_and_kernel_cleanup() { |
| 536 | struct FirstThenPending { |
| 537 | mock: crate::llm_client::mock::MockLlmClient, |
| 538 | entered: tokio::sync::Notify, |
| 539 | } |
| 540 | impl Replies for FirstThenPending { |
| 541 | fn effective_route_envelope( |
| 542 | &self, |
| 543 | model: &str, |
| 544 | at: chrono::DateTime<chrono::Utc>, |
| 545 | ) -> crate::cost_status::EffectiveRouteEnvelope { |
| 546 | Replies::effective_route_envelope(&self.mock, model, at) |
| 547 | } |
| 548 | fn effective_max_output_tokens(&self, model: &str) -> u32 { |
| 549 | Replies::effective_max_output_tokens(&self.mock, model) |
| 550 | } |
| 551 | fn create_message_boxed( |
| 552 | &self, |
| 553 | request: MessageRequest, |
| 554 | ) -> Pin<Box<dyn Future<Output = anyhow::Result<MessageResponse>> + Send + '_>> { |
| 555 | Box::pin(async move { |
| 556 | if self.mock.call_count() > 0 { |
| 557 | self.entered.notify_one(); |
| 558 | std::future::pending().await |
| 559 | } else { |
| 560 | Replies::create_message_boxed(&self.mock, request).await |
| 561 | } |
| 562 | }) |
| 563 | } |
| 564 | } |
| 565 | let replies = Arc::new(FirstThenPending { |
| 566 | mock: crate::llm_client::mock::MockLlmClient::new(Vec::new()), |
| 567 | entered: tokio::sync::Notify::new(), |
| 568 | }); |
| 569 | replies.mock.push_message_response(MessageResponse { |
| 570 | id: "completed-python-round".into(), |
| 571 | r#type: "message".into(), |
| 572 | role: "assistant".into(), |
| 573 | content: vec![ContentBlock::Text { |
| 574 | text: "```repl\nprint(_os.environ['RLM_CONTEXT_FILE'])\n```".into(), |
| 575 | cache_control: None, |
| 576 | }], |
| 577 | model: "captured-model".into(), |
| 578 | stop_reason: Some("end_turn".into()), |
| 579 | stop_sequence: None, |
| 580 | container: None, |
| 581 | usage: Usage { |
| 582 | input_tokens: 7, |
| 583 | output_tokens: 11, |
| 584 | ..Usage::default() |
| 585 | }, |
| 586 | }); |
| 587 | let caller = caller_for(replies.clone(), "captured-model"); |
| 588 | let cancel = caller.fixture_origin_cancel(); |
| 589 | let usage = RlmUsageAccumulator::new(); |
| 590 | let (tx, mut rx) = mpsc::channel(1); |
| 591 | let call = caller.dispatch(RlmInvocation { |
| 592 | prompt: "context held only by Python".into(), |
| 593 | mode: RlmMode::Recursive { depth_remaining: 0 }, |
| 594 | max_tokens: None, |
| 595 | task_instructions: None, |
| 596 | deadline: tokio::time::Instant::now() + Duration::from_secs(60), |
| 597 | gate: Some(crate::tools::codemode::NestedCallGate::admitting_for_test()), |
| 598 | events: Some(tx.clone()), |
| 599 | usage: usage.clone(), |
| 600 | }); |
| 601 | tokio::pin!(call); |
| 602 | tokio::time::timeout(Duration::from_secs(10), async { |
| 603 | loop { |
| 604 | tokio::select! { |
| 605 | biased; |
| 606 | _ = replies.entered.notified() => break, |
| 607 | result = &mut call => panic!("second provider request was never reached: {:?}", result.error), |
| 608 | event = rx.recv() => assert!(event.is_some()), |
| 609 | } |
| 610 | } |
| 611 | }).await.expect("completed Python round must reach the second actual Core request"); |
| 612 | // Preserve the receiver but stop observing: it can no longer release a |
| 613 | // blocked forwarding send. Cancellation, not its 60-second deadline, |
| 614 | // must now retire the owned provider and Python work. |
| 615 | let _ = tx.try_send(Event::status("occupied caller channel")); |
| 616 | cancel.cancel(); |
| 617 | let result = tokio::time::timeout(Duration::from_secs(1), &mut call) |
| 618 | .await |
| 619 | .expect("original cancellation must interrupt a full observational channel"); |
| 620 | assert_eq!(result.iterations, 2); |
| 621 | assert_eq!(result.usage.input_tokens, 7); |
| 622 | assert_eq!(result.usage.output_tokens, 11); |
| 623 | assert_eq!(result.routed_usage.len(), 1); |
| 624 | assert_eq!(result.routed_usage_drop_records.len(), 1); |
| 625 | assert_eq!( |
| 626 | result.routed_usage_drop_records[0].reason, |
| 627 | crate::cost_status::RuntimeUsageMissingReason::RequestOutcomeUnknown |
| 628 | ); |
| 629 | let context_path = PathBuf::from( |
| 630 | result |
| 631 | .trace |
| 632 | .first() |
| 633 | .expect("completed round trace") |
| 634 | .stdout_preview |
| 635 | .trim(), |
| 636 | ); |
| 637 | assert!(context_path.is_absolute()); |
| 638 | assert!( |
| 639 | !context_path.exists(), |
| 640 | "cancelled nested Engine releases its owned kernel/context file" |
| 641 | ); |
| 642 | assert_eq!(usage.snapshot().await.dropped_records, 1); |
| 643 | } |
| 644 | |
| 645 | #[tokio::test] |
| 646 | async fn nested_host_never_claims_or_overwrites_parent_posture_revision() { |
| 647 | let workspace = tempfile::tempdir().unwrap(); |
| 648 | let config = fixture_config("captured-model"); |
| 649 | let mock = Arc::new(crate::llm_client::mock::MockLlmClient::new(Vec::new())); |
| 650 | let (mut parent, _handle) = Engine::new_with_model_client( |
| 651 | EngineConfig { |
| 652 | workspace: workspace.path().to_path_buf(), |
| 653 | model: "captured-model".into(), |
| 654 | memory_enabled: false, |
| 655 | snapshots_enabled: false, |
| 656 | subagents_enabled: false, |
| 657 | ..EngineConfig::default() |
| 658 | }, |
| 659 | &config, |
| 660 | mock, |
| 661 | ); |
| 662 | install_fixture_route(&mut parent); |
| 663 | { |
| 664 | let mut state = parent.live_runtime_authority.lock().unwrap(); |
| 665 | state.revision = 17; |
| 666 | state.applied_revision = 8; |
| 667 | } |
| 668 | let context = admitted_context(&parent, "origin-posture-turn"); |
| 669 | let caller = context.rlm_caller.as_ref().unwrap().clone(); |
| 670 | let (nested, _handle) = Engine::new_rlm_admitted( |
| 671 | caller, |
| 672 | RlmInvocation { |
| 673 | prompt: "no provider invocation".into(), |
| 674 | mode: RlmMode::Completion, |
| 675 | max_tokens: None, |
| 676 | task_instructions: None, |
| 677 | deadline: tokio::time::Instant::now() + Duration::from_secs(60), |
| 678 | gate: None, |
| 679 | events: None, |
| 680 | usage: RlmUsageAccumulator::new(), |
| 681 | }, |
| 682 | ) |
| 683 | .await |
| 684 | .unwrap(); |
| 685 | assert!(Arc::ptr_eq( |
| 686 | &nested.subagent_manager, |
| 687 | &parent.subagent_manager |
| 688 | )); |
| 689 | assert!(!Arc::ptr_eq( |
| 690 | &nested.live_runtime_authority, |
| 691 | &parent.live_runtime_authority |
| 692 | )); |
| 693 | nested.record_applied_runtime_authority(&TurnAuthority::from_effective_fields( |
| 694 | AppMode::Plan, |
| 695 | false, |
| 696 | false, |
| 697 | false, |
| 698 | ApprovalMode::Never, |
| 699 | )); |
| 700 | let state = parent.live_runtime_authority.lock().unwrap(); |
| 701 | assert_eq!(state.revision, 17); |
| 702 | assert_eq!(state.applied_revision, 8); |
| 703 | assert_eq!(parent.current_mode, AppMode::Agent); |
| 704 | } |
| 705 | |
| 706 | #[tokio::test] |
| 707 | async fn repeated_real_core_preparation_failures_have_bounded_receipts_without_provider_work() { |
| 708 | let workspace = tempfile::tempdir().unwrap(); |
| 709 | let mut config = fixture_config("captured-model"); |
| 710 | let identity = config.active_provider_identity().unwrap(); |
| 711 | config |
| 712 | .provider_config_for_mut(&identity) |
| 713 | .unwrap() |
| 714 | .context_window = Some(1); |
| 715 | let mock = Arc::new(crate::llm_client::mock::MockLlmClient::new(Vec::new())); |
| 716 | let (mut engine, _handle) = Engine::new_with_model_client( |
| 717 | EngineConfig { |
| 718 | workspace: workspace.path().to_path_buf(), |
| 719 | model: "captured-model".into(), |
| 720 | memory_enabled: false, |
| 721 | snapshots_enabled: false, |
| 722 | subagents_enabled: false, |
| 723 | terminal_chrome_enabled: false, |
| 724 | ..EngineConfig::default() |
| 725 | }, |
| 726 | &config, |
| 727 | mock.clone(), |
| 728 | ); |
| 729 | install_fixture_route(&mut engine); |
| 730 | assert_eq!( |
| 731 | engine |
| 732 | .active_route_limits |
| 733 | .and_then(|limits| limits.context_tokens), |
| 734 | Some(1), |
| 735 | "the installed actual Core route carries the fixture's tiny context budget", |
| 736 | ); |
| 737 | let context = admitted_context(&engine, "bounded-failed-turn"); |
| 738 | let caller = context.rlm_caller.as_ref().unwrap(); |
| 739 | let usage = RlmUsageAccumulator::new(); |
| 740 | let invocation = || RlmInvocation { |
| 741 | prompt: "real request preparation must retain this whole input".into(), |
| 742 | mode: RlmMode::Completion, |
| 743 | max_tokens: None, |
| 744 | task_instructions: None, |
| 745 | deadline: tokio::time::Instant::now() + Duration::from_secs(60), |
| 746 | gate: None, |
| 747 | events: None, |
| 748 | usage: usage.clone(), |
| 749 | }; |
| 750 | for _ in 0..crate::cost_status::MAX_CHILD_USAGE_RECORDS { |
| 751 | let result = caller.dispatch(invocation()).await; |
| 752 | assert!( |
| 753 | result.error.as_deref().unwrap().contains("history exceeds"), |
| 754 | "{result:?}" |
| 755 | ); |
| 756 | assert_eq!(result.iterations, 0); |
| 757 | } |
| 758 | let before = usage.snapshot().await; |
| 759 | assert_eq!( |
| 760 | before.nested_events.len(), |
| 761 | crate::cost_status::MAX_CHILD_USAGE_RECORDS |
| 762 | ); |
| 763 | assert!(before.records.is_empty()); |
| 764 | assert!(before.drop_records.is_empty()); |
| 765 | assert_eq!(before.dropped_records, 0); |
| 766 | assert_eq!(mock.call_count(), 0); |
| 767 | let refused = caller.dispatch(invocation()).await; |
| 768 | assert!( |
| 769 | refused |
| 770 | .error |
| 771 | .as_deref() |
| 772 | .unwrap() |
| 773 | .contains("nested-turn receipt limit") |
| 774 | ); |
| 775 | let after = usage.snapshot().await; |
| 776 | assert_eq!(after.nested_events, before.nested_events); |
| 777 | assert_eq!(mock.call_count(), 0); |
| 778 | } |
| 779 | |
| 780 | #[tokio::test] |
| 781 | async fn actual_expired_python_startup_releases_the_staged_context() { |
| 782 | let mock = Arc::new(crate::llm_client::mock::MockLlmClient::new(Vec::new())); |
| 783 | let caller = Arc::new(caller_for(mock.clone(), "captured-model")); |
| 784 | let prompt = format!("owned startup cleanup fixture {}", uuid::Uuid::new_v4()); |
| 785 | let result = Engine::new_rlm_admitted( |
| 786 | caller, |
| 787 | RlmInvocation { |
| 788 | prompt: prompt.clone(), |
| 789 | mode: RlmMode::Recursive { depth_remaining: 0 }, |
| 790 | max_tokens: None, |
| 791 | task_instructions: None, |
| 792 | deadline: tokio::time::Instant::now(), |
| 793 | gate: Some(crate::tools::codemode::NestedCallGate::admitting_for_test()), |
| 794 | events: None, |
| 795 | usage: RlmUsageAccumulator::new(), |
| 796 | }, |
| 797 | ) |
| 798 | .await; |
| 799 | assert!( |
| 800 | result |
| 801 | .err() |
| 802 | .unwrap() |
| 803 | .to_string() |
| 804 | .contains("Python startup failed") |
| 805 | ); |
| 806 | let prefix = format!("session_{}_", std::process::id()); |
| 807 | for entry in std::fs::read_dir(std::env::temp_dir().join("deepseek_rlm_ctx")).unwrap() { |
| 808 | let entry = entry.unwrap(); |
| 809 | // Only inspect this test process's similarly sized fixture files; |
| 810 | // parallel RLM tests keep their own unrelated kernels untouched. |
| 811 | if !entry.file_name().to_string_lossy().starts_with(&prefix) { |
| 812 | continue; |
| 813 | } |
| 814 | if let Ok(metadata) = entry.metadata() |
| 815 | && metadata.len() == u64::try_from(prompt.len()).unwrap() |
| 816 | && let Ok(content) = std::fs::read_to_string(entry.path()) |
| 817 | { |
| 818 | assert_ne!( |
| 819 | content, prompt, |
| 820 | "failed startup must remove its staged file" |
| 821 | ); |
| 822 | } |
| 823 | } |
| 824 | assert_eq!(mock.call_count(), 0); |
| 825 | } |
| 826 |