| 1 | |
| 2 | |
| 3 | #[test] |
| 4 | fn agent_mode_elevates_writes_without_granting_network() { |
| 5 | // #273 elevated Agent mode's sandbox so `curl`, package managers, and |
| 6 | // similar shell commands worked, and justified it by saying the |
| 7 | // application-level NetworkPolicy would remain "the only outbound |
| 8 | // boundary". That premise did not hold: NetworkPolicy governs |
| 9 | // fetch_url/web_search/MCP HTTP and never constrained shell subprocesses, |
| 10 | // so workspace-write turns had unrestricted egress with no boundary at |
| 11 | // all. Writing to the workspace is now decoupled from reaching the |
| 12 | // network: the write elevation stays, the network grant does not. |
| 13 | // Network comes from `sandbox_network_access`, from a danger-full-access |
| 14 | // posture, or from the post-denial elevation prompt. |
| 15 | let (engine, _handle) = Engine::new(EngineConfig::default(), &Config::default()); |
| 16 | |
| 17 | let agent_ctx = engine.build_tool_context(AppMode::Agent, false); |
| 18 | let agent_policy = agent_ctx |
| 19 | .elevated_sandbox_policy |
| 20 | .as_ref() |
| 21 | .expect("Agent mode should elevate the sandbox policy"); |
| 22 | assert!( |
| 23 | !agent_policy.has_network_access(), |
| 24 | "Agent mode must not grant shell network access by default; got {agent_policy:?}", |
| 25 | ); |
| 26 | assert!( |
| 27 | !agent_policy |
| 28 | .get_writable_roots(&std::env::current_dir().unwrap_or_else(|_| PathBuf::from("."))) |
| 29 | .is_empty(), |
| 30 | "Agent mode must still elevate workspace writes; got {agent_policy:?}", |
| 31 | ); |
| 32 | |
| 33 | let full_access_ctx = engine.build_tool_context(AppMode::Agent, true); |
| 34 | let full_access_policy = full_access_ctx |
| 35 | .elevated_sandbox_policy |
| 36 | .as_ref() |
| 37 | .expect("Full Access should elevate the sandbox policy"); |
| 38 | assert!(full_access_policy.has_network_access()); |
| 39 | // v0.8.11: Full Access drops to DangerFullAccess (no sandbox) so the |
| 40 | // user is not bounced through approval round-trips for legitimate |
| 41 | // outside-workspace writes (package installs, sub-agent |
| 42 | // workspaces, ~/.cache mutations, etc.). Full Access is opt-in and |
| 43 | // already enables trust mode + auto-approve; the sandbox was the |
| 44 | // last guardrail and contradicts the contract. |
| 45 | assert!( |
| 46 | matches!( |
| 47 | full_access_policy, |
| 48 | crate::sandbox::SandboxPolicy::DangerFullAccess |
| 49 | ), |
| 50 | "Full Access must use DangerFullAccess (no sandbox); got {full_access_policy:?}", |
| 51 | ); |
| 52 | |
| 53 | // Plan mode (#1077): the sandbox must actually deny workspace writes. |
| 54 | // The previous WorkspaceWrite-with-empty-network policy whitelisted the |
| 55 | // workspace as writable, so `python -c "open('f','w').write('x')"` |
| 56 | // mutated files inside the workspace despite Plan-mode's intent. Lock |
| 57 | // it to ReadOnly: no writes anywhere, no network. The shell tool stays |
| 58 | // exposed for read-only inspection (`ls`, `git log`, `grep`, …) and |
| 59 | // the per-platform sandbox enforces the rest. |
| 60 | let plan_ctx = engine.build_tool_context(AppMode::Plan, false); |
| 61 | let plan_policy = plan_ctx |
| 62 | .elevated_sandbox_policy |
| 63 | .as_ref() |
| 64 | .expect("Plan mode should make the shell sandbox policy explicit"); |
| 65 | assert!( |
| 66 | matches!(plan_policy, crate::sandbox::SandboxPolicy::ReadOnly), |
| 67 | "Plan mode must use ReadOnly sandbox to deny workspace writes (#1077); got {plan_policy:?}", |
| 68 | ); |
| 69 | assert!(!plan_policy.has_network_access()); |
| 70 | assert!(!plan_policy.has_full_disk_write_access()); |
| 71 | assert!( |
| 72 | plan_policy |
| 73 | .get_writable_roots(&std::env::current_dir().unwrap_or_else(|_| PathBuf::from("."))) |
| 74 | .is_empty(), |
| 75 | "ReadOnly policy must enumerate zero writable roots; got {plan_policy:?}", |
| 76 | ); |
| 77 | assert!( |
| 78 | plan_ctx |
| 79 | .shell_network_denied_hint |
| 80 | .as_deref() |
| 81 | .is_some_and(|hint| hint.contains("Plan mode") && hint.contains("read-only")), |
| 82 | ); |
| 83 | } |
| 84 | |
| 85 | #[test] |
| 86 | fn sandbox_policy_for_turn_returns_correct_default_policy_per_mode() { |
| 87 | use crate::core::authority::{SandboxNetworkAccess, sandbox_policy_for_turn}; |
| 88 | use crate::sandbox::SandboxPolicy; |
| 89 | use ApprovalMode; |
| 90 | |
| 91 | let workspace = PathBuf::from("/tmp/example-workspace"); |
| 92 | |
| 93 | // Plan: ReadOnly. The whole point of #1077. |
| 94 | assert!(matches!( |
| 95 | sandbox_policy_for_turn( |
| 96 | AppMode::Plan, |
| 97 | ApprovalMode::Suggest, |
| 98 | None, |
| 99 | &workspace, |
| 100 | SandboxNetworkAccess::Restricted, |
| 101 | ), |
| 102 | SandboxPolicy::ReadOnly |
| 103 | )); |
| 104 | |
| 105 | // Agent: WorkspaceWrite with workspace as writable root, network OFF. |
| 106 | match sandbox_policy_for_turn( |
| 107 | AppMode::Agent, |
| 108 | ApprovalMode::Suggest, |
| 109 | None, |
| 110 | &workspace, |
| 111 | SandboxNetworkAccess::Restricted, |
| 112 | ) { |
| 113 | SandboxPolicy::WorkspaceWrite { |
| 114 | writable_roots, |
| 115 | network_access, |
| 116 | .. |
| 117 | } => { |
| 118 | assert_eq!(writable_roots, vec![workspace.clone()]); |
| 119 | assert!( |
| 120 | !network_access, |
| 121 | "workspace-write must not imply shell network access" |
| 122 | ); |
| 123 | } |
| 124 | other => panic!("Agent mode should be WorkspaceWrite; got {other:?}"), |
| 125 | } |
| 126 | |
| 127 | // Agent with the explicit opt-in: same posture, network on. |
| 128 | match sandbox_policy_for_turn( |
| 129 | AppMode::Agent, |
| 130 | ApprovalMode::Suggest, |
| 131 | None, |
| 132 | &workspace, |
| 133 | SandboxNetworkAccess::Allowed, |
| 134 | ) { |
| 135 | SandboxPolicy::WorkspaceWrite { network_access, .. } => { |
| 136 | assert!( |
| 137 | network_access, |
| 138 | "sandbox_network_access = true must grant shell network access" |
| 139 | ); |
| 140 | } |
| 141 | other => panic!("Agent mode should be WorkspaceWrite; got {other:?}"), |
| 142 | } |
| 143 | |
| 144 | // Bypass posture: DangerFullAccess. |
| 145 | assert!(matches!( |
| 146 | sandbox_policy_for_turn( |
| 147 | AppMode::Agent, |
| 148 | ApprovalMode::Bypass, |
| 149 | None, |
| 150 | &workspace, |
| 151 | SandboxNetworkAccess::Restricted, |
| 152 | ), |
| 153 | SandboxPolicy::DangerFullAccess |
| 154 | )); |
| 155 | } |
| 156 | |
| 157 | #[tokio::test] |
| 158 | async fn session_update_preserves_reasoning_tool_only_turn() { |
| 159 | let (mut engine, handle) = Engine::new(EngineConfig::default(), &Config::default()); |
| 160 | let assistant = Message { |
| 161 | role: Role::Assistant, |
| 162 | content: vec![ |
| 163 | ContentBlock::Thinking { |
| 164 | signature: None, |
| 165 | state: None, |
| 166 | thinking: "Need a tool before answering.".to_string(), |
| 167 | }, |
| 168 | ContentBlock::ToolUse { |
| 169 | execution_id: None, |
| 170 | id: "tool-1".to_string(), |
| 171 | name: "read_file".to_string(), |
| 172 | input: json!({"path": "Cargo.toml"}), |
| 173 | caller: None, |
| 174 | thought_signature: None, |
| 175 | }, |
| 176 | ], |
| 177 | }; |
| 178 | |
| 179 | engine.add_session_message(assistant.clone()).await; |
| 180 | |
| 181 | let event = { |
| 182 | let mut rx = handle.rx_event.write().await; |
| 183 | rx.recv().await.expect("session update event") |
| 184 | }; |
| 185 | let Event::SessionUpdated { messages, .. } = event else { |
| 186 | panic!("expected session update event"); |
| 187 | }; |
| 188 | |
| 189 | assert_eq!(*messages, vec![assistant]); |
| 190 | } |
| 191 | |
| 192 | #[tokio::test] |
| 193 | async fn set_model_reloads_instruction_sources_and_updates_session_prompt() { |
| 194 | let tmp = tempdir().expect("tempdir"); |
| 195 | let instructions = tmp.path().join("instructions.md"); |
| 196 | fs::write(&instructions, "FLASH_INSTRUCTIONS_MARKER").expect("write instructions"); |
| 197 | let config = EngineConfig { |
| 198 | workspace: tmp.path().to_path_buf(), |
| 199 | model: "deepseek-v4-flash".to_string(), |
| 200 | instructions: vec![instructions.clone().into()], |
| 201 | ..Default::default() |
| 202 | }; |
| 203 | let (engine, handle) = Engine::new(config, &Config::default()); |
| 204 | fs::write(&instructions, "PRO_INSTRUCTIONS_MARKER").expect("rewrite instructions"); |
| 205 | |
| 206 | let run = tokio::spawn(engine.run()); |
| 207 | handle |
| 208 | .send(Op::SetModel { |
| 209 | model: "deepseek-v4-pro".to_string(), |
| 210 | mode: AppMode::Agent, |
| 211 | route_limits: None, |
| 212 | }) |
| 213 | .await |
| 214 | .expect("send set model"); |
| 215 | |
| 216 | let (model, prompt) = { |
| 217 | let mut rx = handle.rx_event.write().await; |
| 218 | loop { |
| 219 | let event = tokio::time::timeout(std::time::Duration::from_secs(1), rx.recv()) |
| 220 | .await |
| 221 | .expect("session update after model switch") |
| 222 | .expect("event"); |
| 223 | if let Event::SessionUpdated { |
| 224 | model, |
| 225 | system_prompt, |
| 226 | .. |
| 227 | } = event |
| 228 | { |
| 229 | let prompt = match system_prompt.expect("system prompt") { |
| 230 | SystemPrompt::Text(text) => text, |
| 231 | SystemPrompt::Blocks(blocks) => blocks |
| 232 | .into_iter() |
| 233 | .map(|block| block.text) |
| 234 | .collect::<Vec<_>>() |
| 235 | .join("\n"), |
| 236 | }; |
| 237 | break (model, prompt); |
| 238 | } |
| 239 | } |
| 240 | }; |
| 241 | run.abort(); |
| 242 | |
| 243 | assert_eq!(model, "deepseek-v4-pro"); |
| 244 | assert!(prompt.contains("PRO_INSTRUCTIONS_MARKER")); |
| 245 | assert!(!prompt.contains("FLASH_INSTRUCTIONS_MARKER")); |
| 246 | } |
| 247 | |
| 248 | #[tokio::test] |
| 249 | async fn change_mode_refreshes_session_prompt_and_updates_session() { |
| 250 | let tmp = tempdir().expect("tempdir"); |
| 251 | let config = EngineConfig { |
| 252 | workspace: tmp.path().to_path_buf(), |
| 253 | model: "deepseek-v4-pro".to_string(), |
| 254 | ..Default::default() |
| 255 | }; |
| 256 | let (engine, handle) = Engine::new(config, &Config::default()); |
| 257 | |
| 258 | let run = tokio::spawn(engine.run()); |
| 259 | handle |
| 260 | .send(Op::ChangeMode { |
| 261 | mode: AppMode::Agent, |
| 262 | allow_shell: true, |
| 263 | trust_mode: true, |
| 264 | auto_approve: true, |
| 265 | approval_mode: ApprovalMode::Bypass, |
| 266 | configured_sandbox_mode: None, |
| 267 | }) |
| 268 | .await |
| 269 | .expect("send change mode"); |
| 270 | |
| 271 | let (_prompt, messages) = { |
| 272 | let mut rx = handle.rx_event.write().await; |
| 273 | loop { |
| 274 | let event = tokio::time::timeout(std::time::Duration::from_secs(1), rx.recv()) |
| 275 | .await |
| 276 | .expect("session update after mode switch") |
| 277 | .expect("event"); |
| 278 | if let Event::SessionUpdated { |
| 279 | system_prompt, |
| 280 | messages, |
| 281 | .. |
| 282 | } = event |
| 283 | { |
| 284 | let prompt = match system_prompt.expect("system prompt") { |
| 285 | SystemPrompt::Text(text) => text, |
| 286 | SystemPrompt::Blocks(blocks) => blocks |
| 287 | .into_iter() |
| 288 | .map(|block| block.text) |
| 289 | .collect::<Vec<_>>() |
| 290 | .join("\n"), |
| 291 | }; |
| 292 | break (prompt, messages); |
| 293 | } |
| 294 | } |
| 295 | }; |
| 296 | run.abort(); |
| 297 | |
| 298 | assert!( |
| 299 | messages.iter().all(|message| message.role != "system"), |
| 300 | "mode switch must not persist appended system messages: {messages:?}" |
| 301 | ); |
| 302 | } |
| 303 | |
| 304 | /// A posture change announces itself in product words (§19): Permissions, |
| 305 | /// then Plan / Work / Operate. A republished identical posture says nothing. |
| 306 | #[tokio::test] |
| 307 | async fn posture_change_status_uses_permissions_and_work() { |
| 308 | let tmp = tempdir().expect("tempdir"); |
| 309 | let config = EngineConfig { |
| 310 | workspace: tmp.path().to_path_buf(), |
| 311 | ..Default::default() |
| 312 | }; |
| 313 | let (mut engine, handle) = Engine::new(config, &Config::default()); |
| 314 | let publish = |handle: &EngineHandle| { |
| 315 | handle |
| 316 | .try_send(Op::ChangeMode { |
| 317 | mode: AppMode::Agent, |
| 318 | allow_shell: true, |
| 319 | trust_mode: false, |
| 320 | auto_approve: true, |
| 321 | approval_mode: ApprovalMode::Bypass, |
| 322 | configured_sandbox_mode: None, |
| 323 | }) |
| 324 | .expect("publish live runtime authority"); |
| 325 | }; |
| 326 | publish(&handle); |
| 327 | assert!(engine.apply_pending_runtime_authority().await); |
| 328 | publish(&handle); |
| 329 | assert!(!engine.apply_pending_runtime_authority().await); |
| 330 | |
| 331 | let mut statuses = Vec::new(); |
| 332 | let mut rx = handle.rx_event.write().await; |
| 333 | while let Ok(event) = rx.try_recv() { |
| 334 | if let Event::Status { message } = event { |
| 335 | statuses.push(message); |
| 336 | } |
| 337 | } |
| 338 | assert_eq!( |
| 339 | statuses, |
| 340 | vec!["Permissions: Full Access · Work".to_string()] |
| 341 | ); |
| 342 | } |
| 343 | |
| 344 | #[tokio::test] |
| 345 | async fn live_runtime_authority_applies_latest_posture_and_sandbox_before_tools() { |
| 346 | use crate::sandbox::SandboxPolicy; |
| 347 | use ApprovalMode; |
| 348 | |
| 349 | let tmp = tempdir().expect("tempdir"); |
| 350 | let config = EngineConfig { |
| 351 | workspace: tmp.path().to_path_buf(), |
| 352 | ..Default::default() |
| 353 | }; |
| 354 | let (mut engine, handle) = Engine::new(config, &Config::default()); |
| 355 | let registry = ToolRegistryBuilder::new() |
| 356 | .build(engine.build_tool_context(engine.current_mode, engine.session.auto_approve)); |
| 357 | |
| 358 | for (mode, posture, auto_approve, sandbox_mode, expected_sandbox) in [ |
| 359 | ( |
| 360 | AppMode::Operate, |
| 361 | ApprovalMode::Auto, |
| 362 | false, |
| 363 | Some("read-only".to_string()), |
| 364 | SandboxPolicy::ReadOnly, |
| 365 | ), |
| 366 | ( |
| 367 | AppMode::Agent, |
| 368 | ApprovalMode::Bypass, |
| 369 | true, |
| 370 | None, |
| 371 | SandboxPolicy::DangerFullAccess, |
| 372 | ), |
| 373 | ( |
| 374 | AppMode::Agent, |
| 375 | ApprovalMode::Suggest, |
| 376 | false, |
| 377 | None, |
| 378 | SandboxPolicy::WorkspaceWrite { |
| 379 | writable_roots: vec![tmp.path().to_path_buf()], |
| 380 | network_access: false, |
| 381 | exclude_tmpdir: false, |
| 382 | exclude_slash_tmp: false, |
| 383 | }, |
| 384 | ), |
| 385 | ] { |
| 386 | handle |
| 387 | .try_send(Op::ChangeMode { |
| 388 | mode, |
| 389 | allow_shell: true, |
| 390 | trust_mode: false, |
| 391 | auto_approve, |
| 392 | approval_mode: posture, |
| 393 | configured_sandbox_mode: sandbox_mode, |
| 394 | }) |
| 395 | .expect("publish live runtime authority"); |
| 396 | |
| 397 | let published = handle.runtime_permission_authority(); |
| 398 | assert_eq!(published.approval_mode, posture); |
| 399 | assert_eq!(published.auto_approve, auto_approve); |
| 400 | assert!(engine.apply_pending_runtime_authority().await); |
| 401 | assert_eq!(engine.current_mode, mode); |
| 402 | assert_eq!(engine.session.approval_mode, posture); |
| 403 | assert_eq!(engine.session.auto_approve, auto_approve); |
| 404 | assert_eq!( |
| 405 | engine |
| 406 | .live_tool_context(Some(®istry)) |
| 407 | .expect("live registry context") |
| 408 | .elevated_sandbox_policy, |
| 409 | Some(expected_sandbox), |
| 410 | ); |
| 411 | // Tools carry the posture the turn resolved, so a task they create can |
| 412 | // pin the authority it was actually granted. |
| 413 | assert_eq!( |
| 414 | engine |
| 415 | .live_tool_context(Some(®istry)) |
| 416 | .expect("live registry context") |
| 417 | .approval_mode, |
| 418 | posture, |
| 419 | ); |
| 420 | } |
| 421 | } |
| 422 | |
| 423 | #[test] |
| 424 | fn turn_approval_mode_prefers_auto_approve_flag() { |
| 425 | use ApprovalMode; |
| 426 | |
| 427 | assert_eq!( |
| 428 | agent_approval_mode_for_turn(true, ApprovalMode::Suggest), |
| 429 | ApprovalMode::Bypass |
| 430 | ); |
| 431 | assert_eq!( |
| 432 | agent_approval_mode_for_turn(true, ApprovalMode::Never), |
| 433 | ApprovalMode::Bypass |
| 434 | ); |
| 435 | } |
| 436 | |
| 437 | #[test] |
| 438 | fn messages_with_turn_metadata_returns_stored_session_messages() { |
| 439 | use ApprovalMode; |
| 440 | |
| 441 | let tmp = tempdir().expect("tempdir"); |
| 442 | let config = EngineConfig { |
| 443 | workspace: tmp.path().to_path_buf(), |
| 444 | ..Default::default() |
| 445 | }; |
| 446 | let (mut engine, _handle) = Engine::new(config, &Config::default()); |
| 447 | engine.current_mode = AppMode::Plan; |
| 448 | engine.session.approval_mode = ApprovalMode::Suggest; |
| 449 | engine.session.messages = vec![Message { |
| 450 | role: Role::User, |
| 451 | content: vec![ContentBlock::Text { |
| 452 | text: "summary after compaction".to_string(), |
| 453 | cache_control: None, |
| 454 | }], |
| 455 | }] |
| 456 | .into(); |
| 457 | let stored = engine.session.messages.clone(); |
| 458 | |
| 459 | let request_messages = engine.messages_with_turn_metadata(); |
| 460 | |
| 461 | assert_eq!(&*engine.session.messages, &*stored); |
| 462 | assert_eq!(request_messages.len(), stored.len()); |
| 463 | assert!( |
| 464 | request_messages |
| 465 | .iter() |
| 466 | .all(|message| message.role != "system"), |
| 467 | "model request projection must not create appended system messages" |
| 468 | ); |
| 469 | } |
| 470 | |
| 471 | // === To-do state reaches the model through its own tool results === |
| 472 | // |
| 473 | // Codewhale has one To-do list. The model learns what is on it the way it |
| 474 | // learns anything else: from the result its own `work_update` call returned, |
| 475 | // which is ordinary persisted history. No request re-states the list, on any |
| 476 | // step. The complete list stays visible in the UI, which is a different |
| 477 | // surface from the request. |
| 478 | |
| 479 | fn todo_engine() -> ( |
| 480 | Engine, |
| 481 | EngineHandle, |
| 482 | crate::tools::todo::SharedTodoList, |
| 483 | crate::work_graph::SharedWorkRuntime, |
| 484 | tempfile::TempDir, |
| 485 | ) { |
| 486 | let tmp = tempdir().expect("tempdir"); |
| 487 | let todos = crate::tools::todo::new_shared_todo_list(); |
| 488 | let plan = crate::tools::plan::new_shared_plan_state(); |
| 489 | let work = crate::work_graph::new_shared_work_runtime(todos.clone(), plan.clone()); |
| 490 | let config = EngineConfig { |
| 491 | workspace: tmp.path().to_path_buf(), |
| 492 | todos: todos.clone(), |
| 493 | plan_state: plan, |
| 494 | runtime_services: crate::tools::spec::RuntimeToolServices { |
| 495 | work: Some(work.clone()), |
| 496 | ..Default::default() |
| 497 | }, |
| 498 | ..Default::default() |
| 499 | }; |
| 500 | let (mut engine, handle) = Engine::new(config, &Config::default()); |
| 501 | engine.session.messages = vec![Message { |
| 502 | role: Role::User, |
| 503 | content: vec![ContentBlock::Text { |
| 504 | text: "land the To-do seam".to_string(), |
| 505 | cache_control: None, |
| 506 | }], |
| 507 | }] |
| 508 | .into(); |
| 509 | (engine, handle, todos, work, tmp) |
| 510 | } |
| 511 | |
| 512 | fn message_text_of(message: &Message) -> String { |
| 513 | message |
| 514 | .content |
| 515 | .iter() |
| 516 | .filter_map(|block| match block { |
| 517 | ContentBlock::Text { text, .. } => Some(text.as_str()), |
| 518 | _ => None, |
| 519 | }) |
| 520 | .collect::<Vec<_>>() |
| 521 | .join("\n") |
| 522 | } |
| 523 | |
| 524 | /// Run the real `work_update` tool against the attached graph — not a direct |
| 525 | /// mutation of the legacy list, which is exactly the state the fork seam must |
| 526 | /// stop trusting. |
| 527 | async fn run_graph_backed_work_update( |
| 528 | todos: &crate::tools::todo::SharedTodoList, |
| 529 | work: &crate::work_graph::SharedWorkRuntime, |
| 530 | items: serde_json::Value, |
| 531 | ) { |
| 532 | use crate::tools::spec::ToolSpec as _; |
| 533 | let mut context = crate::tools::spec::ToolContext::new(std::env::temp_dir()); |
| 534 | context.runtime.work = Some(work.clone()); |
| 535 | crate::tools::todo::TodoWriteTool::new(todos.clone()) |
| 536 | .execute(json!({ "todos": items }), &context) |
| 537 | .await |
| 538 | .expect("graph-backed todo_write"); |
| 539 | } |
| 540 | |
| 541 | /// A non-empty To-do adds nothing to the messages a request is built from. |
| 542 | #[tokio::test] |
| 543 | async fn a_non_empty_todo_adds_nothing_to_the_request_messages() { |
| 544 | let (engine, _handle, todos, work, _tmp) = todo_engine(); |
| 545 | let stored = engine.session.messages.clone(); |
| 546 | |
| 547 | run_graph_backed_work_update( |
| 548 | &todos, |
| 549 | &work, |
| 550 | json!([ |
| 551 | { "content": "read the runtime seam", "status": "completed" }, |
| 552 | { "content": "write the renderer", "status": "in_progress" } |
| 553 | ]), |
| 554 | ) |
| 555 | .await; |
| 556 | |
| 557 | let request = engine.messages_with_turn_metadata(); |
| 558 | |
| 559 | assert_eq!(request.len(), stored.len(), "nothing may be appended"); |
| 560 | assert_eq!(&*engine.session.messages, &*stored, "history is untouched"); |
| 561 | for message in &request { |
| 562 | let text = message_text_of(message); |
| 563 | assert!( |
| 564 | !text.contains("To-do ("), |
| 565 | "request re-stated the list: {text}" |
| 566 | ); |
| 567 | assert!( |
| 568 | !text.contains("write the renderer"), |
| 569 | "request re-stated an item: {text}" |
| 570 | ); |
| 571 | assert!(!text.contains("codewhale:work"), "{text}"); |
| 572 | } |
| 573 | } |
| 574 | |
| 575 | /// The outbound payload itself, across a whole turn and the turn after it: |
| 576 | /// with work on the list, no provider request body mentions it. |
| 577 | #[tokio::test] |
| 578 | #[allow(clippy::await_holding_lock)] |
| 579 | async fn provider_request_bodies_never_carry_the_todo_list() { |
| 580 | use wiremock::matchers::{method, path}; |
| 581 | use wiremock::{Mock, MockServer, ResponseTemplate}; |
| 582 | |
| 583 | let _lock = lock_test_env(); |
| 584 | let workspace = tempdir().expect("tempdir"); |
| 585 | let server = MockServer::start().await; |
| 586 | let done_sse = concat!( |
| 587 | "data: {\"id\":\"chatcmpl-todo\",\"choices\":[{\"index\":0,", |
| 588 | "\"delta\":{\"content\":\"noted.\"},\"finish_reason\":null}]}\n\n", |
| 589 | "data: {\"id\":\"chatcmpl-todo\",\"choices\":[{\"index\":0,", |
| 590 | "\"delta\":{},\"finish_reason\":\"stop\"}]}\n\n", |
| 591 | "data: [DONE]\n\n", |
| 592 | ); |
| 593 | Mock::given(method("POST")) |
| 594 | .and(path("/v1/chat/completions")) |
| 595 | .respond_with( |
| 596 | ResponseTemplate::new(200) |
| 597 | .insert_header("content-type", "text/event-stream") |
| 598 | .set_body_string(done_sse), |
| 599 | ) |
| 600 | .mount(&server) |
| 601 | .await; |
| 602 | |
| 603 | let api_config = Config { |
| 604 | ..Config::default() |
| 605 | } |
| 606 | .with_legacy_root(Some("test-key".to_string()), Some(server.uri())); |
| 607 | let todos = crate::tools::todo::new_shared_todo_list(); |
| 608 | let plan = crate::tools::plan::new_shared_plan_state(); |
| 609 | let work = crate::work_graph::new_shared_work_runtime(todos.clone(), plan.clone()); |
| 610 | let engine_config = EngineConfig { |
| 611 | workspace: workspace.path().to_path_buf(), |
| 612 | snapshots_enabled: false, |
| 613 | subagents_enabled: false, |
| 614 | todos: todos.clone(), |
| 615 | plan_state: plan, |
| 616 | runtime_services: crate::tools::spec::RuntimeToolServices { |
| 617 | work: Some(work.clone()), |
| 618 | ..Default::default() |
| 619 | }, |
| 620 | ..EngineConfig::default() |
| 621 | }; |
| 622 | run_graph_backed_work_update( |
| 623 | &todos, |
| 624 | &work, |
| 625 | json!([{ "content": "ship the outbound payload test", "status": "in_progress" }]), |
| 626 | ) |
| 627 | .await; |
| 628 | |
| 629 | let (engine, handle) = Engine::new(engine_config, &api_config); |
| 630 | let task = tokio::spawn(engine.run()); |
| 631 | |
| 632 | for prompt in ["first turn", "second turn"] { |
| 633 | handle |
| 634 | .send(external_user_message_op( |
| 635 | prompt, |
| 636 | AppMode::Agent, |
| 637 | &api_config, |
| 638 | )) |
| 639 | .await |
| 640 | .expect("send turn"); |
| 641 | let mut rx = handle.rx_event.write().await; |
| 642 | while let Some(event) = tokio::time::timeout(model_turn_event_timeout(), rx.recv()) |
| 643 | .await |
| 644 | .expect("timed out waiting for turn completion") |
| 645 | { |
| 646 | match event { |
| 647 | Event::Error { envelope, .. } => panic!("turn errored: {envelope:?}"), |
| 648 | Event::TurnComplete { status, error, .. } => { |
| 649 | assert_eq!(status, TurnOutcomeStatus::Completed, "{error:?}"); |
| 650 | break; |
| 651 | } |
| 652 | _ => {} |
| 653 | } |
| 654 | } |
| 655 | } |
| 656 | |
| 657 | let requests = server |
| 658 | .received_requests() |
| 659 | .await |
| 660 | .expect("recorded provider requests"); |
| 661 | assert!(requests.len() >= 2, "expected one request per turn"); |
| 662 | for request in &requests { |
| 663 | let body = String::from_utf8_lossy(&request.body); |
| 664 | assert!( |
| 665 | !body.contains("ship the outbound payload test"), |
| 666 | "a provider request restated the To-do list: {body}" |
| 667 | ); |
| 668 | assert!(!body.contains("To-do ("), "{body}"); |
| 669 | assert!(!body.contains("codewhale:work"), "{body}"); |
| 670 | } |
| 671 | |
| 672 | handle.send(Op::Shutdown).await.expect("shutdown engine"); |
| 673 | task.await.expect("engine task"); |
| 674 | } |
| 675 | |
| 676 | /// The turn-start structured state is deliberately To-do-free: the list moves |
| 677 | /// during a turn, so it is resolved at the fork seam instead. |
| 678 | #[test] |
| 679 | fn turn_start_structured_state_carries_no_todo_section() { |
| 680 | let state = StructuredState { |
| 681 | mode_label: "Agent".to_string(), |
| 682 | workspace: PathBuf::from("/workspace/codewhale"), |
| 683 | cwd: None, |
| 684 | working_set_summary: None, |
| 685 | subagent_snapshots: Vec::new(), |
| 686 | }; |
| 687 | |
| 688 | let block = state.to_system_block().expect("fork state block"); |
| 689 | |
| 690 | assert!( |
| 691 | !block.contains(crate::todo_snapshot::FORK_TODO_SECTION_HEADING), |
| 692 | "stable capture must not pin a To-do section: {block}" |
| 693 | ); |
| 694 | assert!(!block.contains("To-do (")); |
| 695 | } |
| 696 | |
| 697 | /// The fork handoff and `/relay` show the same To-do snapshot body. Relay |
| 698 | /// parity is asserted in `commands::tests`. |
| 699 | #[tokio::test] |
| 700 | async fn fork_state_block_reuses_the_snapshot_body() { |
| 701 | let (engine, _handle, todos, work, _tmp) = todo_engine(); |
| 702 | run_graph_backed_work_update( |
| 703 | &todos, |
| 704 | &work, |
| 705 | json!([ |
| 706 | { "content": "Wire Fleet progress projection", "status": "in_progress" }, |
| 707 | { "content": "Run focused gates", "status": "pending" } |
| 708 | ]), |
| 709 | ) |
| 710 | .await; |
| 711 | let snapshot = engine.todo_source().snapshot().await; |
| 712 | let body = crate::todo_snapshot::todo_snapshot_body(&snapshot).expect("body"); |
| 713 | |
| 714 | let state = StructuredState { |
| 715 | mode_label: "Agent".to_string(), |
| 716 | workspace: PathBuf::from("/workspace/codewhale"), |
| 717 | cwd: None, |
| 718 | working_set_summary: None, |
| 719 | subagent_snapshots: Vec::new(), |
| 720 | }; |
| 721 | let fork_context = crate::tools::subagent::SubAgentForkContext { |
| 722 | messages: engine.messages_with_turn_metadata(), |
| 723 | structured_state_block: state.to_system_block(), |
| 724 | work_source: Some(engine.todo_source()), |
| 725 | }; |
| 726 | |
| 727 | let resolved = fork_context |
| 728 | .with_resolved_state_block() |
| 729 | .await |
| 730 | .structured_state_block |
| 731 | .expect("resolved fork state block"); |
| 732 | |
| 733 | assert!(resolved.contains(&body), "fork body drifted: {resolved}"); |
| 734 | } |
| 735 | |
| 736 | /// `update_plan` is conversational reasoning, not a second list: plan-only |
| 737 | /// state must not produce a To-do snapshot. |
| 738 | #[tokio::test] |
| 739 | async fn plan_only_state_produces_no_todo_snapshot() { |
| 740 | let (engine, _handle, _todos, _work, _tmp) = todo_engine(); |
| 741 | { |
| 742 | let mut plan = engine.config.plan_state.lock().await; |
| 743 | plan.update(crate::tools::plan::UpdatePlanArgs { |
| 744 | objective: Some("Ship the To-do seam".to_string()), |
| 745 | plan: vec![crate::tools::plan::PlanItemArg { |
| 746 | step: "draft the renderer".to_string(), |
| 747 | status: crate::tools::plan::StepStatus::InProgress, |
| 748 | }], |
| 749 | ..crate::tools::plan::UpdatePlanArgs::default() |
| 750 | }); |
| 751 | assert!(!plan.snapshot().is_empty()); |
| 752 | } |
| 753 | |
| 754 | assert!( |
| 755 | engine.todo_source().body().await.is_none(), |
| 756 | "legacy plan-only state must not become a To-do snapshot" |
| 757 | ); |
| 758 | } |
| 759 | |
| 760 | /// A real graph-backed `work_update` stages the new projection in the |
| 761 | /// `WorkRuntime` and publishes into `config.todos` only later, asynchronously, |
| 762 | /// from the UI. The fork seam must read the staged projection, not the |
| 763 | /// pre-write legacy view. |
| 764 | #[tokio::test] |
| 765 | async fn fork_seam_reflects_a_graph_backed_work_update() { |
| 766 | let (engine, _handle, todos, work, _tmp) = todo_engine(); |
| 767 | assert!( |
| 768 | engine.todo_source().is_graph_backed(), |
| 769 | "this engine must read the graph, not the legacy view" |
| 770 | ); |
| 771 | |
| 772 | run_graph_backed_work_update( |
| 773 | &todos, |
| 774 | &work, |
| 775 | json!([ |
| 776 | { "content": "read the runtime seam", "status": "completed" }, |
| 777 | { "content": "hand the child the live list", "status": "in_progress" } |
| 778 | ]), |
| 779 | ) |
| 780 | .await; |
| 781 | |
| 782 | // The staleness this test exists for: the legacy view is still empty here. |
| 783 | assert!( |
| 784 | todos.lock().await.snapshot().is_empty(), |
| 785 | "precondition: work_update stages in the graph and publishes later" |
| 786 | ); |
| 787 | |
| 788 | let body = engine.todo_source().body().await.expect("body"); |
| 789 | assert!(body.contains("[x] #1 read the runtime seam"), "{body}"); |
| 790 | assert!( |
| 791 | body.contains("[~] #2 hand the child the live list"), |
| 792 | "{body}" |
| 793 | ); |
| 794 | } |
| 795 | |
| 796 | /// A compaction checkpoint is history, never a second To-do surface or a |
| 797 | /// stable-system-prefix mutation. |
| 798 | #[tokio::test] |
| 799 | async fn compaction_keeps_todos_out_of_the_prefix() { |
| 800 | let (mut engine, _handle, todos, work, _tmp) = todo_engine(); |
| 801 | run_graph_backed_work_update( |
| 802 | &todos, |
| 803 | &work, |
| 804 | json!([{ "content": "staged graph-only todo", "status": "in_progress" }]), |
| 805 | ) |
| 806 | .await; |
| 807 | assert!( |
| 808 | todos.lock().await.snapshot().is_empty(), |
| 809 | "precondition: the compatibility list is still stale" |
| 810 | ); |
| 811 | |
| 812 | let stable_before = engine.session.system_prompt.clone(); |
| 813 | engine.commit_compaction_checkpoint(Some(SystemPrompt::Text(format!( |
| 814 | "{COMPACTION_SUMMARY_MARKER}\nsummary" |
| 815 | )))); |
| 816 | assert_eq!(engine.session.system_prompt, stable_before); |
| 817 | let checkpoint = engine.rendered_compaction_summary().expect("checkpoint"); |
| 818 | assert!( |
| 819 | !checkpoint.contains("staged graph-only todo"), |
| 820 | "{checkpoint}" |
| 821 | ); |
| 822 | assert!(!checkpoint.contains("### Todos"), "{checkpoint}"); |
| 823 | } |
| 824 | |
| 825 | #[tokio::test] |
| 826 | async fn compaction_completed_reports_complete_post_input_tokens() { |
| 827 | let _env = crate::test_support::lock_test_env(); |
| 828 | let tmp = tempdir().expect("tempdir"); |
| 829 | let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", tmp.path()); |
| 830 | let config = EngineConfig { |
| 831 | workspace: tmp.path().to_path_buf(), |
| 832 | ..Default::default() |
| 833 | }; |
| 834 | let (mut engine, handle) = Engine::new(config, &Config::default()); |
| 835 | engine.session.system_prompt = Some(SystemPrompt::Text("stable system context ".repeat(400))); |
| 836 | engine.session.replace_messages(vec![Message { |
| 837 | role: Role::User, |
| 838 | content: vec![ContentBlock::Text { |
| 839 | text: "post-compaction message".to_string(), |
| 840 | cache_control: None, |
| 841 | }], |
| 842 | }]); |
| 843 | engine.commit_compaction_checkpoint(Some(SystemPrompt::Text(format!( |
| 844 | "{COMPACTION_SUMMARY_MARKER}\npost-compaction summary" |
| 845 | )))); |
| 846 | |
| 847 | let messages_only = |
| 848 | crate::compaction::estimate_input_tokens_for_pressure(&engine.session.messages, None); |
| 849 | let expected = engine.estimated_input_tokens(); |
| 850 | assert!(expected > messages_only); |
| 851 | |
| 852 | engine |
| 853 | .emit_compaction_completed( |
| 854 | "compact_test".to_string(), |
| 855 | false, |
| 856 | "Made room".to_string(), |
| 857 | Some(4), |
| 858 | Some(1), |
| 859 | super::compaction::CompactionPass { |
| 860 | trigger: "manual", |
| 861 | path: crate::compaction::CompactionPath::Summary, |
| 862 | tokens_before: 9000, |
| 863 | threshold_tokens: 8000, |
| 864 | usage: Usage { |
| 865 | input_tokens: 120, |
| 866 | output_tokens: 15, |
| 867 | ..Default::default() |
| 868 | }, |
| 869 | }, |
| 870 | ) |
| 871 | .await; |
| 872 | |
| 873 | let event = handle |
| 874 | .rx_event |
| 875 | .write() |
| 876 | .await |
| 877 | .recv() |
| 878 | .await |
| 879 | .expect("compaction completed event"); |
| 880 | let Event::CompactionCompleted { |
| 881 | post_input_tokens, .. |
| 882 | } = event |
| 883 | else { |
| 884 | panic!("expected CompactionCompleted, got {event:?}"); |
| 885 | }; |
| 886 | assert_eq!(post_input_tokens, Some(expected as u64)); |
| 887 | let log = std::fs::read_to_string(tmp.path().join("audit.log")).unwrap(); |
| 888 | let records = log |
| 889 | .lines() |
| 890 | .map(|line| serde_json::from_str::<serde_json::Value>(line).unwrap()) |
| 891 | .filter(|row| row["event"] == "compaction.completed") |
| 892 | .collect::<Vec<_>>(); |
| 893 | assert_eq!(records.len(), 1); |
| 894 | let details = &records[0]["details"]; |
| 895 | assert_eq!(details["messages_before"], 4); |
| 896 | assert_eq!(details["messages_after"], 1); |
| 897 | assert_eq!(details["reduction_ratio"], 0.75); |
| 898 | assert_eq!(details["estimated_tokens_before"], 9000); |
| 899 | assert_eq!(details["estimated_tokens_after"], expected); |
| 900 | assert_eq!(details["threshold_tokens"], 8000); |
| 901 | assert_eq!(details["summarizer_usage"]["input_tokens"], 120); |
| 902 | assert_eq!(details["trigger"], "manual"); |
| 903 | assert_eq!(details["path"], "summary"); |
| 904 | assert!(!log.contains("post-compaction message")); |
| 905 | engine |
| 906 | .record_compaction_event( |
| 907 | "compaction.refused", |
| 908 | serde_json::json!({ |
| 909 | "trigger": "auto", "reason": "retained_floor", "threshold_tokens": 8000, |
| 910 | }), |
| 911 | ) |
| 912 | .await; |
| 913 | assert!( |
| 914 | std::fs::read_to_string(tmp.path().join("audit.log")) |
| 915 | .unwrap() |
| 916 | .contains("compaction.refused") |
| 917 | ); |
| 918 | } |
| 919 | |
| 920 | /// `fork_context` is captured once at turn start, so a `work_update` followed |
| 921 | /// by an `agent` spawn *in the same turn* must still hand the child the |
| 922 | /// current snapshot. Only the To-do portion is refreshed; the inherited |
| 923 | /// transcript and stable state text are unchanged. |
| 924 | #[tokio::test] |
| 925 | async fn same_turn_fork_carries_the_updated_todo() { |
| 926 | let (engine, _handle, todos, work, _tmp) = todo_engine(); |
| 927 | |
| 928 | // Turn start: capture the fork context, before any work exists. |
| 929 | let stable_block = StructuredState { |
| 930 | mode_label: "Agent".to_string(), |
| 931 | workspace: engine.config.workspace.clone(), |
| 932 | cwd: None, |
| 933 | working_set_summary: None, |
| 934 | subagent_snapshots: Vec::new(), |
| 935 | } |
| 936 | .to_system_block(); |
| 937 | let fork_context = crate::tools::subagent::SubAgentForkContext { |
| 938 | messages: engine.messages_with_turn_metadata(), |
| 939 | structured_state_block: stable_block.clone(), |
| 940 | work_source: Some(engine.todo_source()), |
| 941 | }; |
| 942 | let captured_messages = fork_context.messages.clone(); |
| 943 | assert!( |
| 944 | !fork_context |
| 945 | .with_resolved_state_block() |
| 946 | .await |
| 947 | .structured_state_block |
| 948 | .expect("stable block") |
| 949 | .contains("To-do ("), |
| 950 | "no work yet, so no To-do section" |
| 951 | ); |
| 952 | |
| 953 | // Mid-turn: the model calls work_update, then spawns an agent. |
| 954 | run_graph_backed_work_update( |
| 955 | &todos, |
| 956 | &work, |
| 957 | json!([{ "content": "hand the child the live list", "status": "in_progress" }]), |
| 958 | ) |
| 959 | .await; |
| 960 | |
| 961 | let resolved = fork_context.with_resolved_state_block().await; |
| 962 | let block = resolved |
| 963 | .structured_state_block |
| 964 | .as_deref() |
| 965 | .expect("resolved block"); |
| 966 | let snapshot = engine.todo_source().snapshot().await; |
| 967 | let body = crate::todo_snapshot::todo_snapshot_body(&snapshot).expect("body"); |
| 968 | |
| 969 | assert!( |
| 970 | block.contains(&body), |
| 971 | "same-turn fork must carry the current body: {block}" |
| 972 | ); |
| 973 | assert!( |
| 974 | block.contains("[~] #1 hand the child the live list"), |
| 975 | "{block}" |
| 976 | ); |
| 977 | // Stable history semantics are untouched. |
| 978 | assert_eq!(resolved.messages, captured_messages); |
| 979 | assert!( |
| 980 | block.starts_with(stable_block.as_deref().expect("stable").trim()), |
| 981 | "the stable capture must stay a byte-identical prefix: {block}" |
| 982 | ); |
| 983 | } |
| 984 | |
| 985 | /// U1: hosts resend the compaction config on every model, route or session |
| 986 | /// sync. Neither a session sync nor a config whose switch did not move should |
| 987 | /// produce a status line: it would overwrite a real error (the missing-key |
| 988 | /// notice) or the host's confirmed "Resumed:" receipt in the footer. |
| 989 | #[tokio::test] |
| 990 | async fn unchanged_compaction_config_is_acknowledged_silently() { |
| 991 | let tmp = tempdir().expect("tempdir"); |
| 992 | let (engine, handle) = Engine::new( |
| 993 | EngineConfig { |
| 994 | workspace: tmp.path().to_path_buf(), |
| 995 | ..Default::default() |
| 996 | }, |
| 997 | &Config::default(), |
| 998 | ); |
| 999 | let current = engine.config.compaction.clone(); |
| 1000 | let restored_messages = vec![Message { |
| 1001 | role: Role::User, |
| 1002 | content: vec![ContentBlock::Text { |
| 1003 | text: "Restored conversation proof".to_string(), |
| 1004 | cache_control: None, |
| 1005 | }], |
| 1006 | }]; |
| 1007 | let run = tokio::spawn(engine.run()); |
| 1008 | handle |
| 1009 | .send(Op::SyncSession { |
| 1010 | session_id: Some("resumed-session".to_string()), |
| 1011 | messages: restored_messages.clone(), |
| 1012 | system_prompt: None, |
| 1013 | system_prompt_override: false, |
| 1014 | model: current.model.clone(), |
| 1015 | workspace: tmp.path().to_path_buf(), |
| 1016 | mode: AppMode::Agent, |
| 1017 | }) |
| 1018 | .await |
| 1019 | .expect("sync restored session"); |
| 1020 | handle |
| 1021 | .send(Op::SetCompaction { |
| 1022 | config: current.clone(), |
| 1023 | }) |
| 1024 | .await |
| 1025 | .expect("send unchanged config"); |
| 1026 | // A session restore resyncs the model and window with the switch as it |
| 1027 | // was: applied, but not news. |
| 1028 | let mut resynced = current.clone(); |
| 1029 | resynced.model = format!("{}-resynced", current.model); |
| 1030 | resynced.effective_context_window = Some(64_000); |
| 1031 | handle |
| 1032 | .send(Op::SetCompaction { |
| 1033 | config: resynced.clone(), |
| 1034 | }) |
| 1035 | .await |
| 1036 | .expect("send resynced config"); |
| 1037 | let mut changed = resynced; |
| 1038 | changed.enabled = !changed.enabled; |
| 1039 | let expected = if changed.enabled { |
| 1040 | "Make room automatically: on" |
| 1041 | } else { |
| 1042 | "Make room automatically: off" |
| 1043 | }; |
| 1044 | handle |
| 1045 | .send(Op::SetCompaction { config: changed }) |
| 1046 | .await |
| 1047 | .expect("send changed config"); |
| 1048 | |
| 1049 | let mut rx = handle.rx_event.write().await; |
| 1050 | let mut session_updated = false; |
| 1051 | let first_status = loop { |
| 1052 | let event = tokio::time::timeout(Duration::from_secs(2), rx.recv()) |
| 1053 | .await |
| 1054 | .expect("status after a real change") |
| 1055 | .expect("event"); |
| 1056 | match event { |
| 1057 | Event::SessionUpdated { |
| 1058 | session_id, |
| 1059 | messages, |
| 1060 | model, |
| 1061 | workspace, |
| 1062 | .. |
| 1063 | } => { |
| 1064 | assert_eq!(session_id, "resumed-session"); |
| 1065 | assert_eq!(*messages, restored_messages); |
| 1066 | assert_eq!(model, current.model); |
| 1067 | assert_eq!(workspace, tmp.path()); |
| 1068 | session_updated = true; |
| 1069 | } |
| 1070 | Event::Status { message } => break message, |
| 1071 | _ => {} |
| 1072 | } |
| 1073 | }; |
| 1074 | assert!( |
| 1075 | session_updated, |
| 1076 | "session sync still publishes its authoritative update" |
| 1077 | ); |
| 1078 | assert_eq!( |
| 1079 | first_status, expected, |
| 1080 | "session sync and unchanged/resynced configs produced no status; only the switch did" |
| 1081 | ); |
| 1082 | drop(rx); |
| 1083 | run.abort(); |
| 1084 | } |
| 1085 | |
| 1086 | #[tokio::test] |
| 1087 | async fn change_mode_op_updates_current_mode_and_emits_status() { |
| 1088 | let tmp = tempdir().expect("tempdir"); |
| 1089 | let config = EngineConfig { |
| 1090 | workspace: tmp.path().to_path_buf(), |
| 1091 | model: "deepseek-v4-pro".to_string(), |
| 1092 | ..Default::default() |
| 1093 | }; |
| 1094 | let (engine, handle) = Engine::new(config, &Config::default()); |
| 1095 | |
| 1096 | let run = tokio::spawn(engine.run()); |
| 1097 | handle |
| 1098 | .send(Op::ChangeMode { |
| 1099 | mode: AppMode::Agent, |
| 1100 | allow_shell: true, |
| 1101 | trust_mode: true, |
| 1102 | auto_approve: true, |
| 1103 | approval_mode: ApprovalMode::Bypass, |
| 1104 | configured_sandbox_mode: None, |
| 1105 | }) |
| 1106 | .await |
| 1107 | .expect("send change mode"); |
| 1108 | |
| 1109 | // Expect a SessionUpdated event confirming the mode change. |
| 1110 | let mut rx = handle.rx_event.write().await; |
| 1111 | let session_updated = tokio::time::timeout(std::time::Duration::from_secs(2), rx.recv()) |
| 1112 | .await |
| 1113 | .expect("session update after mode switch") |
| 1114 | .expect("event"); |
| 1115 | let Event::SessionUpdated { messages, .. } = session_updated else { |
| 1116 | panic!("should emit SessionUpdated after mode change, got: {session_updated:?}"); |
| 1117 | }; |
| 1118 | assert!( |
| 1119 | messages.iter().all(|message| message.role != "system"), |
| 1120 | "mode switch must not persist synthetic system messages: {messages:?}" |
| 1121 | ); |
| 1122 | |
| 1123 | // Also expect a status event |
| 1124 | let status = tokio::time::timeout(std::time::Duration::from_secs(2), rx.recv()) |
| 1125 | .await |
| 1126 | .expect("status after mode switch") |
| 1127 | .expect("event"); |
| 1128 | assert!( |
| 1129 | matches!(status, Event::Status { .. }), |
| 1130 | "should emit Status after mode change, got: {status:?}" |
| 1131 | ); |
| 1132 | |
| 1133 | run.abort(); |
| 1134 | } |
| 1135 | |
| 1136 | #[test] |
| 1137 | fn runtime_mode_policy_updates_engine_session_mirrors() { |
| 1138 | let tmp = tempdir().expect("tempdir"); |
| 1139 | let config = EngineConfig { |
| 1140 | workspace: tmp.path().to_path_buf(), |
| 1141 | model: "deepseek-v4-pro".to_string(), |
| 1142 | allow_shell: false, |
| 1143 | trust_mode: false, |
| 1144 | ..Default::default() |
| 1145 | }; |
| 1146 | let (mut engine, _handle) = Engine::new(config, &Config::default()); |
| 1147 | engine.current_mode = AppMode::Plan; |
| 1148 | engine.session.allow_shell = false; |
| 1149 | engine.session.trust_mode = false; |
| 1150 | engine.session.auto_approve = false; |
| 1151 | engine.session.approval_mode = ApprovalMode::Suggest; |
| 1152 | |
| 1153 | let agent_authority = crate::core::authority::TurnAuthority::from_effective_fields( |
| 1154 | AppMode::Agent, |
| 1155 | true, |
| 1156 | false, |
| 1157 | false, |
| 1158 | ApprovalMode::Never, |
| 1159 | ); |
| 1160 | engine.apply_runtime_mode_policy(&agent_authority); |
| 1161 | |
| 1162 | assert_eq!(engine.current_mode, AppMode::Agent); |
| 1163 | assert!(engine.session.allow_shell); |
| 1164 | assert!(engine.config.allow_shell); |
| 1165 | assert!(!engine.session.trust_mode); |
| 1166 | assert!(!engine.config.trust_mode); |
| 1167 | assert!(!engine.session.auto_approve); |
| 1168 | assert_eq!(engine.session.approval_mode, ApprovalMode::Never); |
| 1169 | |
| 1170 | let full_access_authority = crate::core::authority::TurnAuthority::from_effective_fields( |
| 1171 | AppMode::Agent, |
| 1172 | true, |
| 1173 | true, |
| 1174 | true, |
| 1175 | ApprovalMode::Bypass, |
| 1176 | ); |
| 1177 | engine.apply_runtime_mode_policy(&full_access_authority); |
| 1178 | |
| 1179 | assert_eq!(engine.current_mode, AppMode::Agent); |
| 1180 | assert!(engine.session.allow_shell); |
| 1181 | assert!(engine.session.trust_mode); |
| 1182 | assert!(engine.config.trust_mode); |
| 1183 | assert!(engine.session.auto_approve); |
| 1184 | assert_eq!(engine.session.approval_mode, ApprovalMode::Bypass); |
| 1185 | } |
| 1186 | |
| 1187 | #[tokio::test] |
| 1188 | async fn sync_session_restores_current_mode() { |
| 1189 | let tmp = tempdir().expect("tempdir"); |
| 1190 | let config = EngineConfig { |
| 1191 | workspace: tmp.path().to_path_buf(), |
| 1192 | model: "deepseek-v4-pro".to_string(), |
| 1193 | ..Default::default() |
| 1194 | }; |
| 1195 | let (engine, handle) = Engine::new(config, &Config::default()); |
| 1196 | |
| 1197 | let run = tokio::spawn(engine.run()); |
| 1198 | handle |
| 1199 | .send(Op::SyncSession { |
| 1200 | session_id: Some("plan-session".to_string()), |
| 1201 | messages: Vec::new(), |
| 1202 | system_prompt: None, |
| 1203 | system_prompt_override: false, |
| 1204 | model: "deepseek-v4-pro".to_string(), |
| 1205 | workspace: tmp.path().to_path_buf(), |
| 1206 | mode: AppMode::Plan, |
| 1207 | }) |
| 1208 | .await |
| 1209 | .expect("sync session"); |
| 1210 | |
| 1211 | let (tx, rx) = tokio::sync::oneshot::channel(); |
| 1212 | handle |
| 1213 | .send(Op::GetSessionSnapshot { |
| 1214 | tx: std::sync::Arc::new(std::sync::Mutex::new(Some(tx))), |
| 1215 | }) |
| 1216 | .await |
| 1217 | .expect("request snapshot"); |
| 1218 | let snapshot = tokio::time::timeout(Duration::from_secs(2), rx) |
| 1219 | .await |
| 1220 | .expect("snapshot response") |
| 1221 | .expect("snapshot"); |
| 1222 | |
| 1223 | assert_eq!(snapshot.mode, "plan"); |
| 1224 | |
| 1225 | run.abort(); |
| 1226 | } |
| 1227 | |
| 1228 | #[tokio::test] |
| 1229 | async fn sync_session_without_prompt_repins_full_system_prompt_on_next_turn() { |
| 1230 | use crate::llm_client::mock::{MockLlmClient, canned}; |
| 1231 | |
| 1232 | const WORKSPACE_RULE: &str = "SYNC_SESSION_FULL_PROMPT_PROOF"; |
| 1233 | |
| 1234 | async fn wait_for_completed_turn(handle: &EngineHandle) { |
| 1235 | let mut rx = handle.rx_event.write().await; |
| 1236 | loop { |
| 1237 | let event = tokio::time::timeout(model_turn_event_timeout(), rx.recv()) |
| 1238 | .await |
| 1239 | .expect("timed out waiting for turn completion") |
| 1240 | .expect("engine event channel closed before turn completion"); |
| 1241 | if let Event::TurnComplete { status, error, .. } = event { |
| 1242 | assert_eq!(status, TurnOutcomeStatus::Completed, "{error:?}"); |
| 1243 | break; |
| 1244 | } |
| 1245 | } |
| 1246 | } |
| 1247 | |
| 1248 | let workspace = tempdir().expect("tempdir"); |
| 1249 | fs::write( |
| 1250 | workspace.path().join("AGENTS.md"), |
| 1251 | format!("# Rules\n\nAlways preserve {WORKSPACE_RULE}.\n"), |
| 1252 | ) |
| 1253 | .expect("write AGENTS.md fixture"); |
| 1254 | let config = Config::default(); |
| 1255 | let mock = std::sync::Arc::new(MockLlmClient::new(vec![canned::simple_text_turn( |
| 1256 | "Turn complete.", |
| 1257 | )])); |
| 1258 | let client: crate::core::model_client::SharedModelClient = mock.clone(); |
| 1259 | let (mut engine, handle) = Engine::new_with_model_client( |
| 1260 | deterministic_engine_config(workspace.path()), |
| 1261 | &config, |
| 1262 | client, |
| 1263 | ); |
| 1264 | let established_context = engine.installed_next_turn_prompt_context(); |
| 1265 | assert_eq!( |
| 1266 | engine.refresh_pinned_header_for_turn(&established_context), |
| 1267 | None |
| 1268 | ); |
| 1269 | assert_eq!( |
| 1270 | engine.session.pinned_prompt_context.as_ref(), |
| 1271 | Some(&established_context), |
| 1272 | "precondition: the outgoing conversation has an established prompt pin" |
| 1273 | ); |
| 1274 | let task = tokio::spawn(engine.run()); |
| 1275 | |
| 1276 | handle |
| 1277 | .send(Op::SyncSession { |
| 1278 | session_id: Some("fresh-session".to_string()), |
| 1279 | messages: Vec::new(), |
| 1280 | system_prompt: None, |
| 1281 | system_prompt_override: false, |
| 1282 | model: crate::config::DEFAULT_TEXT_MODEL.to_string(), |
| 1283 | workspace: workspace.path().to_path_buf(), |
| 1284 | mode: AppMode::Agent, |
| 1285 | }) |
| 1286 | .await |
| 1287 | .expect("sync fresh session without a persisted prompt"); |
| 1288 | let synced = handle |
| 1289 | .get_session_snapshot() |
| 1290 | .await |
| 1291 | .expect("drain session sync"); |
| 1292 | assert!( |
| 1293 | synced.system_prompt.is_none(), |
| 1294 | "SyncSession must install the persisted prompt exactly before the next turn" |
| 1295 | ); |
| 1296 | |
| 1297 | handle |
| 1298 | .send(external_user_message_op( |
| 1299 | "Start the newly synchronized conversation.", |
| 1300 | AppMode::Agent, |
| 1301 | &config, |
| 1302 | )) |
| 1303 | .await |
| 1304 | .expect("send first turn after sync"); |
| 1305 | wait_for_completed_turn(&handle).await; |
| 1306 | |
| 1307 | let requests = mock.captured_requests(); |
| 1308 | assert_eq!(requests.len(), 1); |
| 1309 | let repinned = requests[0] |
| 1310 | .system |
| 1311 | .clone() |
| 1312 | .map(system_prompt_text) |
| 1313 | .expect("the first turn after SyncSession must send a full system prompt"); |
| 1314 | assert!(repinned.contains(WORKSPACE_RULE), "{repinned}"); |
| 1315 | assert!( |
| 1316 | requests[0].messages.iter().all(|message| { |
| 1317 | message.content.iter().all(|block| { |
| 1318 | !matches!( |
| 1319 | block, |
| 1320 | ContentBlock::Text { text, .. } if text.starts_with("<context_update>") |
| 1321 | ) |
| 1322 | }) |
| 1323 | }), |
| 1324 | "the fresh session must not inherit a context-update delta: {:?}", |
| 1325 | requests[0].messages |
| 1326 | ); |
| 1327 | |
| 1328 | let refreshed = handle |
| 1329 | .get_session_snapshot() |
| 1330 | .await |
| 1331 | .expect("snapshot refreshed session"); |
| 1332 | let refreshed_prompt = refreshed |
| 1333 | .system_prompt |
| 1334 | .map(system_prompt_text) |
| 1335 | .expect("refreshed session prompt"); |
| 1336 | assert!(refreshed_prompt.contains(WORKSPACE_RULE)); |
| 1337 | handle.send(Op::Shutdown).await.expect("shutdown engine"); |
| 1338 | task.await.expect("engine task"); |
| 1339 | } |
| 1340 | |
| 1341 | #[tokio::test] |
| 1342 | async fn sync_session_same_id_does_not_finalize_live_worker() { |
| 1343 | let tmp = tempdir().expect("tempdir"); |
| 1344 | let workspace = tmp.path().to_path_buf(); |
| 1345 | let config = EngineConfig { |
| 1346 | workspace: workspace.clone(), |
| 1347 | model: "deepseek-v4-pro".to_string(), |
| 1348 | ..Default::default() |
| 1349 | }; |
| 1350 | let (engine, handle) = Engine::new(config, &Config::default()); |
| 1351 | let manager = engine.subagent_manager.clone(); |
| 1352 | |
| 1353 | let run = tokio::spawn(engine.run()); |
| 1354 | // Install the conversation identity first so the manager is not |
| 1355 | // finalized by the very first identity transition away from the |
| 1356 | // construction-time UUID. |
| 1357 | handle |
| 1358 | .send(Op::SyncSession { |
| 1359 | session_id: Some("session-keep".to_string()), |
| 1360 | messages: Vec::new(), |
| 1361 | system_prompt: None, |
| 1362 | system_prompt_override: false, |
| 1363 | model: "deepseek-v4-pro".to_string(), |
| 1364 | workspace: workspace.clone(), |
| 1365 | mode: AppMode::Agent, |
| 1366 | }) |
| 1367 | .await |
| 1368 | .expect("install session"); |
| 1369 | handle.get_session_snapshot().await.expect("drain install"); |
| 1370 | |
| 1371 | let agent_id = { |
| 1372 | let mut manager = manager.write().await; |
| 1373 | manager.insert_test_running_agent("keep", &workspace) |
| 1374 | }; |
| 1375 | |
| 1376 | // A same-id re-sync is a reload, not a conversation boundary. |
| 1377 | handle |
| 1378 | .send(Op::SyncSession { |
| 1379 | session_id: Some("session-keep".to_string()), |
| 1380 | messages: Vec::new(), |
| 1381 | system_prompt: None, |
| 1382 | system_prompt_override: false, |
| 1383 | model: "deepseek-v4-pro".to_string(), |
| 1384 | workspace: workspace.clone(), |
| 1385 | mode: AppMode::Agent, |
| 1386 | }) |
| 1387 | .await |
| 1388 | .expect("re-sync same session"); |
| 1389 | handle.get_session_snapshot().await.expect("drain re-sync"); |
| 1390 | |
| 1391 | let record = manager |
| 1392 | .read() |
| 1393 | .await |
| 1394 | .get_worker_record(&agent_id) |
| 1395 | .expect("live worker record"); |
| 1396 | assert!( |
| 1397 | !record.status.is_terminal(), |
| 1398 | "same-id re-sync must not finalize the worker: {:?}", |
| 1399 | record.status |
| 1400 | ); |
| 1401 | |
| 1402 | run.abort(); |
| 1403 | } |
| 1404 | |
| 1405 | #[tokio::test] |
| 1406 | async fn sync_session_different_id_finalizes_live_worker() { |
| 1407 | let tmp = tempdir().expect("tempdir"); |
| 1408 | let workspace = tmp.path().to_path_buf(); |
| 1409 | let config = EngineConfig { |
| 1410 | workspace: workspace.clone(), |
| 1411 | model: "deepseek-v4-pro".to_string(), |
| 1412 | ..Default::default() |
| 1413 | }; |
| 1414 | let (engine, handle) = Engine::new(config, &Config::default()); |
| 1415 | let manager = engine.subagent_manager.clone(); |
| 1416 | |
| 1417 | let run = tokio::spawn(engine.run()); |
| 1418 | handle |
| 1419 | .send(Op::SyncSession { |
| 1420 | session_id: Some("session-a".to_string()), |
| 1421 | messages: Vec::new(), |
| 1422 | system_prompt: None, |
| 1423 | system_prompt_override: false, |
| 1424 | model: "deepseek-v4-pro".to_string(), |
| 1425 | workspace: workspace.clone(), |
| 1426 | mode: AppMode::Agent, |
| 1427 | }) |
| 1428 | .await |
| 1429 | .expect("install session-a"); |
| 1430 | handle.get_session_snapshot().await.expect("drain install"); |
| 1431 | |
| 1432 | let agent_id = { |
| 1433 | let mut manager = manager.write().await; |
| 1434 | let agent_id = manager.insert_test_running_agent("close", &workspace); |
| 1435 | manager.assign_test_session_owner(&agent_id, "session-a"); |
| 1436 | agent_id |
| 1437 | }; |
| 1438 | |
| 1439 | // A different id is a conversation boundary: the live worker must be |
| 1440 | // finalized with the session-closed reason. |
| 1441 | handle |
| 1442 | .send(Op::SyncSession { |
| 1443 | session_id: Some("session-b".to_string()), |
| 1444 | messages: Vec::new(), |
| 1445 | system_prompt: None, |
| 1446 | system_prompt_override: false, |
| 1447 | model: "deepseek-v4-pro".to_string(), |
| 1448 | workspace: workspace.clone(), |
| 1449 | mode: AppMode::Agent, |
| 1450 | }) |
| 1451 | .await |
| 1452 | .expect("switch to session-b"); |
| 1453 | handle.get_session_snapshot().await.expect("drain switch"); |
| 1454 | |
| 1455 | let record = manager |
| 1456 | .read() |
| 1457 | .await |
| 1458 | .get_worker_record(&agent_id) |
| 1459 | .expect("worker record after close"); |
| 1460 | assert!( |
| 1461 | record.status.is_terminal(), |
| 1462 | "different-id switch must finalize the worker: {:?}", |
| 1463 | record.status |
| 1464 | ); |
| 1465 | assert_eq!( |
| 1466 | record.status, |
| 1467 | crate::tools::subagent::AgentWorkerStatus::Interrupted |
| 1468 | ); |
| 1469 | let reason = record.latest_message.as_deref().unwrap_or(""); |
| 1470 | assert!( |
| 1471 | reason.contains("parent session closed"), |
| 1472 | "session-closed reason missing: {reason}" |
| 1473 | ); |
| 1474 | |
| 1475 | run.abort(); |
| 1476 | } |
| 1477 | |
| 1478 | #[tokio::test] |
| 1479 | async fn sync_session_migrates_one_checkpoint_and_strips_its_system_carrier() { |
| 1480 | let tmp = tempdir().expect("tempdir"); |
| 1481 | let config = EngineConfig { |
| 1482 | workspace: tmp.path().to_path_buf(), |
| 1483 | model: "deepseek-v4-pro".to_string(), |
| 1484 | ..Default::default() |
| 1485 | }; |
| 1486 | let (engine, handle) = Engine::new(config, &Config::default()); |
| 1487 | let carrier = SystemPrompt::Text(format!( |
| 1488 | "stable host prompt\n\n<!-- compaction-summary:begin -->\n{COMPACTION_SUMMARY_MARKER}\nnew checkpoint\n<!-- compaction-summary:end -->" |
| 1489 | )); |
| 1490 | let old_checkpoint = crate::compaction::compaction_checkpoint_message(&SystemPrompt::Text( |
| 1491 | format!("{COMPACTION_SUMMARY_MARKER}\nold checkpoint"), |
| 1492 | )); |
| 1493 | |
| 1494 | let run = tokio::spawn(engine.run()); |
| 1495 | let mut messages = vec![old_checkpoint]; |
| 1496 | for round in 0..2 { |
| 1497 | handle |
| 1498 | .send(Op::SyncSession { |
| 1499 | session_id: Some("compacted-session".to_string()), |
| 1500 | messages, |
| 1501 | system_prompt: Some(carrier.clone()), |
| 1502 | system_prompt_override: true, |
| 1503 | model: "deepseek-v4-pro".to_string(), |
| 1504 | workspace: tmp.path().to_path_buf(), |
| 1505 | mode: AppMode::Agent, |
| 1506 | }) |
| 1507 | .await |
| 1508 | .expect("sync compacted session"); |
| 1509 | |
| 1510 | let (tx, rx) = tokio::sync::oneshot::channel(); |
| 1511 | handle |
| 1512 | .send(Op::GetSessionSnapshot { |
| 1513 | tx: std::sync::Arc::new(std::sync::Mutex::new(Some(tx))), |
| 1514 | }) |
| 1515 | .await |
| 1516 | .expect("request snapshot"); |
| 1517 | let snapshot = tokio::time::timeout(Duration::from_secs(2), rx) |
| 1518 | .await |
| 1519 | .expect("snapshot response") |
| 1520 | .expect("snapshot"); |
| 1521 | |
| 1522 | let checkpoints = snapshot |
| 1523 | .messages |
| 1524 | .iter() |
| 1525 | .filter(|message| crate::compaction::is_compaction_checkpoint_message(message)) |
| 1526 | .collect::<Vec<_>>(); |
| 1527 | assert_eq!(checkpoints.len(), 1, "round {round}: {checkpoints:?}"); |
| 1528 | let checkpoint_text = message_text_of(checkpoints[0]); |
| 1529 | assert!( |
| 1530 | checkpoint_text.contains("new checkpoint"), |
| 1531 | "{checkpoint_text}" |
| 1532 | ); |
| 1533 | assert!( |
| 1534 | !checkpoint_text.contains("old checkpoint"), |
| 1535 | "{checkpoint_text}" |
| 1536 | ); |
| 1537 | assert_eq!( |
| 1538 | snapshot.system_prompt, |
| 1539 | Some(SystemPrompt::Text("stable host prompt".to_string())) |
| 1540 | ); |
| 1541 | messages = snapshot.messages; |
| 1542 | } |
| 1543 | |
| 1544 | run.abort(); |
| 1545 | } |