| 1 | //! Engine state that must stay truthful across a boundary: the live MCP |
| 2 | //! catalog (C02-07/C02-16), a session switch (C02-08/C02-17), a user-input |
| 3 | //! wait (C02-19), an `/edit` whose replacement never starts (C02-02), and a |
| 4 | //! malformed MCP config (C02-18). No fixture reaches a provider or a shell. |
| 5 | |
| 6 | use super::*; |
| 7 | use crate::llm_client::mock::MockLlmClient; |
| 8 | |
| 9 | fn quiet_engine(config: EngineConfig) -> (Engine, EngineHandle) { |
| 10 | Engine::new_with_model_client( |
| 11 | config, |
| 12 | &Config::default(), |
| 13 | Arc::new(MockLlmClient::new(Vec::new())), |
| 14 | ) |
| 15 | } |
| 16 | |
| 17 | async fn session_snapshot(handle: &EngineHandle) -> SessionSnapshot { |
| 18 | let (tx, rx) = tokio::sync::oneshot::channel(); |
| 19 | handle |
| 20 | .send(Op::GetSessionSnapshot { |
| 21 | tx: Arc::new(StdMutex::new(Some(tx))), |
| 22 | }) |
| 23 | .await |
| 24 | .unwrap(); |
| 25 | tokio::time::timeout(model_turn_event_timeout(), rx) |
| 26 | .await |
| 27 | .expect("snapshot response") |
| 28 | .expect("snapshot") |
| 29 | } |
| 30 | |
| 31 | #[test] |
| 32 | fn live_mcp_refresh_drops_tools_the_pool_no_longer_lists() { |
| 33 | // C02-07: the caller's universe is what the pool lists *now*. After a |
| 34 | // live 401 that is only the synthetic login tool — never the dead one. |
| 35 | let mut catalog = vec![api_tool("mcp_alpha_search"), api_tool("exec_shell")]; |
| 36 | let mut active: HashSet<String> = ["mcp_alpha_search", "exec_shell"] |
| 37 | .into_iter() |
| 38 | .map(str::to_string) |
| 39 | .collect(); |
| 40 | let universe: HashSet<String> = ["mcp_alpha_authenticate".to_string()].into(); |
| 41 | let changed = replace_runtime_mcp_tools( |
| 42 | &mut catalog, |
| 43 | &mut active, |
| 44 | &universe, |
| 45 | vec![api_tool("mcp_alpha_authenticate")], |
| 46 | AppMode::Agent, |
| 47 | &HashSet::new(), |
| 48 | crate::model_profile::ToolSurfaceBudget::Standard, |
| 49 | ); |
| 50 | let names: Vec<&str> = catalog.iter().map(|tool| tool.name.as_str()).collect(); |
| 51 | assert_eq!(names, vec!["exec_shell", "mcp_alpha_authenticate"]); |
| 52 | assert!( |
| 53 | !active.contains("mcp_alpha_search"), |
| 54 | "a dead tool is not callable" |
| 55 | ); |
| 56 | assert!(changed); |
| 57 | } |
| 58 | |
| 59 | #[test] |
| 60 | fn live_mcp_refresh_reports_a_schema_only_edit() { |
| 61 | // C02-16: the same names with an edited description is a change the |
| 62 | // prefix check must hear about; an identical refresh is not. |
| 63 | let refresh = |catalog: &mut Vec<Tool>, active: &mut HashSet<String>, tool: Tool| { |
| 64 | replace_runtime_mcp_tools( |
| 65 | catalog, |
| 66 | active, |
| 67 | &["mcp_alpha_read".to_string()].into(), |
| 68 | vec![tool], |
| 69 | AppMode::Agent, |
| 70 | &HashSet::new(), |
| 71 | crate::model_profile::ToolSurfaceBudget::Standard, |
| 72 | ) |
| 73 | }; |
| 74 | let mut catalog = vec![api_tool("exec_shell")]; |
| 75 | let mut active = HashSet::new(); |
| 76 | refresh(&mut catalog, &mut active, api_tool("mcp_alpha_read")); |
| 77 | let mut edited = api_tool("mcp_alpha_read"); |
| 78 | edited.description = "Read a resource (now paginated)".to_string(); |
| 79 | assert!(refresh(&mut catalog, &mut active, edited.clone())); |
| 80 | assert!(!refresh(&mut catalog, &mut active, edited)); |
| 81 | } |
| 82 | |
| 83 | async fn engine_mid_mcp_boot_with_usage() -> (Engine, EngineHandle, tokio::task::JoinHandle<()>) { |
| 84 | let workspace = tempdir().unwrap(); |
| 85 | let (mut engine, handle) = quiet_engine(deterministic_engine_config(workspace.path())); |
| 86 | engine.session.total_usage.input_tokens = 1_000; |
| 87 | engine.session.total_usage.output_tokens = 200; |
| 88 | let boot = tokio::spawn(std::future::pending::<()>()); |
| 89 | engine.mcp_boot_task = Some(boot.abort_handle()); |
| 90 | engine.mcp_boot_in_flight = true; |
| 91 | engine.mcp_boot_generation = Some(7); |
| 92 | engine.mcp_connection_errors.insert( |
| 93 | "previous-server".to_string(), |
| 94 | "previous failure".to_string(), |
| 95 | ); |
| 96 | // A same-conversation re-sync keeps its usage and its boot pass. |
| 97 | let same = engine.session.id.clone(); |
| 98 | assert!(engine.install_synced_session_id(same).is_none()); |
| 99 | assert_eq!(engine.session_snapshot().total_tokens, 1_200); |
| 100 | assert!(engine.mcp_boot_in_flight); |
| 101 | (engine, handle, boot) |
| 102 | } |
| 103 | |
| 104 | #[tokio::test] |
| 105 | async fn a_session_boundary_does_not_inherit_the_previous_usage() { |
| 106 | // C02-17: session A's tokens are not session B's. |
| 107 | let (mut engine, _handle, _boot) = engine_mid_mcp_boot_with_usage().await; |
| 108 | assert!( |
| 109 | engine |
| 110 | .install_synced_session_id("next-session".to_string()) |
| 111 | .is_some() |
| 112 | ); |
| 113 | assert_eq!(engine.session_snapshot().total_tokens, 0); |
| 114 | } |
| 115 | |
| 116 | #[tokio::test] |
| 117 | async fn a_session_boundary_abandons_the_previous_mcp_boot() { |
| 118 | // C02-08: the dropped pool's pass is aborted and its state cleared. |
| 119 | let (mut engine, _handle, boot) = engine_mid_mcp_boot_with_usage().await; |
| 120 | assert!( |
| 121 | engine |
| 122 | .install_synced_session_id("next-session".to_string()) |
| 123 | .is_some() |
| 124 | ); |
| 125 | assert!(!engine.mcp_boot_in_flight); |
| 126 | assert_eq!(engine.mcp_boot_generation, None); |
| 127 | assert!(engine.mcp_boot_rx.is_none()); |
| 128 | assert!(engine.mcp_connection_errors.is_empty()); |
| 129 | let joined = tokio::time::timeout(Duration::from_secs(5), boot) |
| 130 | .await |
| 131 | .expect("the previous boot pass must stop promptly"); |
| 132 | assert!(joined.unwrap_err().is_cancelled(), "the pass was aborted"); |
| 133 | } |
| 134 | |
| 135 | fn empty_user_input_request() -> crate::tools::user_input::UserInputRequest { |
| 136 | crate::tools::user_input::UserInputRequest { |
| 137 | questions: Vec::new(), |
| 138 | } |
| 139 | } |
| 140 | |
| 141 | #[tokio::test] |
| 142 | async fn a_session_boundary_discards_and_stops_the_previous_mcp_supervisor() { |
| 143 | let workspace = tempdir().unwrap(); |
| 144 | let (mut engine, _handle) = quiet_engine(deterministic_engine_config(workspace.path())); |
| 145 | let (tx, rx) = mpsc::channel(1); |
| 146 | tx.send(McpSupervisorUpdate { |
| 147 | died: vec![("previous-server".to_string(), "stale diagnosis".to_string())], |
| 148 | failed: Vec::new(), |
| 149 | recovered: Vec::new(), |
| 150 | parked: Vec::new(), |
| 151 | }) |
| 152 | .await |
| 153 | .unwrap(); |
| 154 | engine.mcp_supervisor_rx = Some(rx); |
| 155 | |
| 156 | // A reconnect future retains its connection resources across await. The |
| 157 | // boundary must abort that future, not merely drop a Weak pool owner. |
| 158 | let held = Arc::new(()); |
| 159 | let weak = Arc::downgrade(&held); |
| 160 | let (started_tx, started_rx) = tokio::sync::oneshot::channel(); |
| 161 | let supervisor = tokio::spawn(async move { |
| 162 | let _held = held; |
| 163 | started_tx.send(()).unwrap(); |
| 164 | std::future::pending::<()>().await; |
| 165 | }); |
| 166 | started_rx.await.unwrap(); |
| 167 | engine.mcp_supervisor_task = Some(supervisor.abort_handle()); |
| 168 | assert!( |
| 169 | engine |
| 170 | .install_synced_session_id("next-session".to_string()) |
| 171 | .is_some() |
| 172 | ); |
| 173 | assert!( |
| 174 | engine.mcp_supervisor_rx.is_none(), |
| 175 | "queued diagnoses retired" |
| 176 | ); |
| 177 | assert!(engine.mcp_supervisor_task.is_none()); |
| 178 | let joined = tokio::time::timeout(Duration::from_secs(1), supervisor) |
| 179 | .await |
| 180 | .expect("the previous supervisor must stop promptly"); |
| 181 | assert!(joined.unwrap_err().is_cancelled()); |
| 182 | assert!(weak.upgrade().is_none(), "reconnect resources released"); |
| 183 | assert!(tx.is_closed(), "old diagnoses cannot cross the boundary"); |
| 184 | |
| 185 | engine.ensure_mcp_pool().await.unwrap(); |
| 186 | assert!(engine.mcp_supervisor_rx.is_some(), "the next pool is armed"); |
| 187 | assert!(engine.mcp_supervisor_task.is_some()); |
| 188 | assert!(engine.mcp_connection_errors.is_empty()); |
| 189 | engine.drop_mcp_pool(); |
| 190 | } |
| 191 | |
| 192 | #[tokio::test] |
| 193 | async fn an_undeliverable_user_input_request_fails_fast() { |
| 194 | // C02-19: the question cannot reach a host (its event channel is closed, |
| 195 | // while the answer channel stays open), so nobody can ever answer it. |
| 196 | let workspace = tempdir().unwrap(); |
| 197 | let (mut engine, handle) = quiet_engine(deterministic_engine_config(workspace.path())); |
| 198 | handle.rx_event.write().await.close(); |
| 199 | let result = tokio::time::timeout( |
| 200 | Duration::from_secs(5), |
| 201 | engine.await_user_input("ask-undeliverable", empty_user_input_request()), |
| 202 | ) |
| 203 | .await |
| 204 | .expect("an undeliverable question must not wait for an answer"); |
| 205 | assert!( |
| 206 | matches!(result, Err(ToolError::ExecutionFailed { .. })), |
| 207 | "{result:?}" |
| 208 | ); |
| 209 | } |
| 210 | |
| 211 | #[tokio::test] |
| 212 | async fn an_open_user_input_question_is_not_charged_to_the_turn() { |
| 213 | // C02-19: a delivered question the person leaves open is human time, not |
| 214 | // the agent's; the per-turn wall clock does not advance across it. |
| 215 | let workspace = tempdir().unwrap(); |
| 216 | let wait = Duration::from_millis(400); |
| 217 | let (mut engine, _handle) = quiet_engine(EngineConfig { |
| 218 | user_input_timeout: Some(wait), |
| 219 | ..deterministic_engine_config(workspace.path()) |
| 220 | }); |
| 221 | let before = engine.turn_wall_clock.spent(); |
| 222 | let result = engine |
| 223 | .await_user_input("ask-unanswered", empty_user_input_request()) |
| 224 | .await; |
| 225 | assert!( |
| 226 | matches!(result, Err(ToolError::Timeout { .. })), |
| 227 | "{result:?}" |
| 228 | ); |
| 229 | let charged = engine.turn_wall_clock.spent().saturating_sub(before); |
| 230 | assert!( |
| 231 | charged < wait / 2, |
| 232 | "the human wait was charged to the turn budget: {charged:?}" |
| 233 | ); |
| 234 | } |
| 235 | |
| 236 | #[tokio::test] |
| 237 | #[allow(clippy::await_holding_lock)] |
| 238 | async fn edit_last_turn_restores_the_exchange_when_the_replacement_never_starts() { |
| 239 | // C02-02: `/edit` cuts the last exchange, then dispatches the new text. |
| 240 | // A dispatch that never starts must hand the cut exchange back. |
| 241 | let _lock = lock_test_env(); |
| 242 | let tmp = tempdir().unwrap(); |
| 243 | let api_config = Config::default().with_legacy_root( |
| 244 | Some("test-key".to_string()), |
| 245 | Some("http://127.0.0.1:9".to_string()), |
| 246 | ); |
| 247 | let (mut engine, handle) = Engine::new_with_model_client( |
| 248 | EngineConfig { |
| 249 | workspace: tmp.path().to_path_buf(), |
| 250 | model: "deepseek-v4-pro".to_string(), |
| 251 | snapshots_enabled: false, |
| 252 | subagents_enabled: false, |
| 253 | ..Default::default() |
| 254 | }, |
| 255 | &api_config, |
| 256 | Arc::new(MockLlmClient::new(Vec::new())), |
| 257 | ); |
| 258 | // No model client: the replacement send returns NotStarted. |
| 259 | engine.model_client = None; |
| 260 | let run = tokio::spawn(engine.run()); |
| 261 | let text = |role: Role, text: &str| Message { |
| 262 | role, |
| 263 | content: vec![ContentBlock::Text { |
| 264 | text: text.to_string(), |
| 265 | cache_control: None, |
| 266 | }], |
| 267 | }; |
| 268 | handle |
| 269 | .send(Op::SyncSession { |
| 270 | session_id: Some("edit-restore".to_string()), |
| 271 | messages: vec![ |
| 272 | text(Role::User, "original prompt"), |
| 273 | text(Role::Assistant, "original answer"), |
| 274 | ], |
| 275 | system_prompt: None, |
| 276 | system_prompt_override: false, |
| 277 | model: "deepseek-v4-pro".to_string(), |
| 278 | workspace: tmp.path().to_path_buf(), |
| 279 | mode: AppMode::Agent, |
| 280 | }) |
| 281 | .await |
| 282 | .unwrap(); |
| 283 | let before = session_snapshot(&handle).await; |
| 284 | assert_eq!(before.messages.len(), 2); |
| 285 | handle |
| 286 | .send(Op::EditLastTurn { |
| 287 | new_message: "edited prompt".to_string(), |
| 288 | submission_id: None, |
| 289 | }) |
| 290 | .await |
| 291 | .unwrap(); |
| 292 | let after = session_snapshot(&handle).await; |
| 293 | assert_eq!( |
| 294 | after.messages, before.messages, |
| 295 | "a replacement that never started must not cost the original exchange" |
| 296 | ); |
| 297 | handle.send(Op::Shutdown).await.unwrap(); |
| 298 | run.await.unwrap(); |
| 299 | } |
| 300 | |
| 301 | #[tokio::test] |
| 302 | async fn malformed_mcp_config_is_reported_instead_of_silently_empty() { |
| 303 | // C02-18: an unreadable config still yields an empty, reloadable pool, |
| 304 | // but the person is told the configured servers are gone. |
| 305 | let workspace = tempdir().unwrap(); |
| 306 | let config_path = workspace.path().join("mcp.json"); |
| 307 | fs::write(&config_path, "{ this is not json").unwrap(); |
| 308 | let (mut engine, handle) = quiet_engine(deterministic_engine_config(workspace.path())); |
| 309 | engine.session.mcp_config_path = config_path; |
| 310 | engine |
| 311 | .ensure_mcp_pool() |
| 312 | .await |
| 313 | .expect("a reloadable empty pool still exists"); |
| 314 | let mut rx = handle.rx_event.write().await; |
| 315 | let reported = std::iter::from_fn(|| rx.try_recv().ok()).any(|event| { |
| 316 | matches!( |
| 317 | event, |
| 318 | Event::Status { ref message } if message.contains("MCP config could not be loaded") |
| 319 | ) |
| 320 | }); |
| 321 | assert!(reported, "a malformed MCP config must not fail silently"); |
| 322 | } |
| 323 |