返回 CodeWhale
test_cases_13.rs
根目录 / crates / tui / src / core / engine / tests / test_cases_13.rs
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(&registry))
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(&registry))
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 }
1545 lines RUST