返回 CodeWhale
lifecycle_tests.rs
根目录 / crates / tui / src / tools / subagent / lifecycle_tests.rs
1 use super::*;
2 use tempfile::tempdir;
3
4 pub(super) fn status_rows(payload: &Value) -> Vec<Value> {
5 let columns = payload["columns"].as_array().expect("roster columns");
6 let unique = columns
7 .iter()
8 .map(|column| column.as_str().expect("column name"))
9 .collect::<HashSet<_>>();
10 assert_eq!(columns.len(), unique.len(), "duplicate roster column");
11 payload["agents"]
12 .as_array()
13 .expect("roster rows")
14 .iter()
15 .map(|row| {
16 let values = row.as_array().expect("roster row array");
17 assert_eq!(columns.len(), values.len(), "row must match its header");
18 Value::Object(
19 columns
20 .iter()
21 .zip(values)
22 .map(|(column, value)| (column.as_str().unwrap().to_string(), value.clone()))
23 .collect(),
24 )
25 })
26 .collect()
27 }
28
29 fn prior_messages() -> Vec<Message> {
30 vec![Message {
31 role: Role::User,
32 content: vec![ContentBlock::Text {
33 text: "retained work".into(),
34 cache_control: None,
35 }],
36 }]
37 }
38
39 #[tokio::test]
40 async fn lifecycle_bulk_followup_preserves_mappings_and_retries_without_duplicate_workers() {
41 let dir = tempdir().unwrap();
42 let manager = new_shared_subagent_manager(dir.path().to_path_buf(), 12);
43 let mut sources = Vec::new();
44 {
45 let mut guard = manager.write().await;
46 for i in 0..6 {
47 let (id, _) = guard.insert_test_interrupted_continuable_agent(
48 &format!("parked-{i}"),
49 dir.path(),
50 prior_messages(),
51 );
52 let agent = guard.agents.get_mut(&id).unwrap();
53 agent.checkpoint.as_mut().unwrap().parked_at_turn_end = true;
54 agent.agent_type = FleetRole::Scout;
55 agent.model = "deepseek-v4-flash".into();
56 agent.allowed_tools = Some(Vec::new());
57 let record = guard.worker_records.get_mut(&id).unwrap();
58 record.parent_run_id = None;
59 record.spec.parent_run_id = None;
60 let spec = &mut record.spec;
61 spec.model = "deepseek-v4-flash".into();
62 spec.agent_type = FleetRole::Scout;
63 spec.runtime_profile = WorkerRuntimeProfile::for_role(FleetRole::Scout);
64 sources.push(id);
65 }
66 }
67 let (client, _, _, fixture_config) =
68 super::tests::delayed_chat_client(Duration::from_secs(30), "fixture result").await;
69 let mut runtime = super::tests::stub_runtime();
70 runtime.manager = Arc::clone(&manager);
71 runtime.client = client;
72 runtime.api_config = Some(Arc::new(fixture_config));
73 runtime.context = ToolContext::new(dir.path());
74 let tool = coord::AgentsFollowupTool::new(Arc::clone(&manager)).with_runtime(runtime);
75 let input = json!({"agent_ids": sources, "message": "Continue the assignment."});
76 let first = tool
77 .execute(input.clone(), &ToolContext::new(dir.path()))
78 .await
79 .unwrap();
80 let first: Value = serde_json::from_str(&first.content).unwrap();
81 assert_eq!(first["results"].as_array().unwrap().len(), 6);
82 assert_eq!(first["errors"], json!([]));
83 let second = tool
84 .execute(input, &ToolContext::new(dir.path()))
85 .await
86 .unwrap();
87 let second: Value = serde_json::from_str(&second.content).unwrap();
88 for (a, b) in first["results"]
89 .as_array()
90 .unwrap()
91 .iter()
92 .zip(second["results"].as_array().unwrap())
93 {
94 assert_eq!(a["from"], b["from"]);
95 assert_eq!(a["to"], b["to"]);
96 assert_ne!(a["from"], a["to"]);
97 }
98 let mut guard = manager.write().await;
99 assert_eq!(guard.agents.len(), 12);
100 for id in sources {
101 let target = guard.continuation_target(&id).unwrap();
102 assert_ne!(target, id);
103 assert_eq!(
104 guard.continuation_source(&target).as_deref(),
105 Some(id.as_str())
106 );
107 let _ = guard.cancel_agent(&target);
108 }
109 }
110
111 #[tokio::test]
112 async fn lifecycle_bulk_followup_reports_unknown_and_foreign_targets_without_hiding_success() {
113 let dir = tempdir().unwrap();
114 let manager = new_shared_subagent_manager(dir.path().to_path_buf(), 4);
115 let (owned, foreign) = {
116 let mut guard = manager.write().await;
117 let owned = guard.insert_test_running_agent("owned", dir.path());
118 let foreign = guard.insert_test_running_agent("foreign", dir.path());
119 guard.assign_test_session_owner(&foreign, "another-session");
120 (owned, foreign)
121 };
122 let tool = coord::AgentsFollowupTool::new(Arc::clone(&manager));
123 let result = tool
124 .execute(
125 json!({"agent_ids": [owned, "missing", foreign], "message": "Check progress"}),
126 &ToolContext::new(dir.path()),
127 )
128 .await
129 .unwrap();
130 let payload: Value = serde_json::from_str(&result.content).unwrap();
131 assert_eq!(payload["results"].as_array().unwrap().len(), 1);
132 assert_eq!(payload["errors"].as_array().unwrap().len(), 2);
133 assert!(!manager.read().await.child_was_woken(&foreign));
134 }
135
136 #[tokio::test]
137 async fn lifecycle_followup_rejects_ambiguous_and_invalid_batch_inputs_before_delivery() {
138 let dir = tempdir().unwrap();
139 let manager = new_shared_subagent_manager(dir.path().to_path_buf(), 2);
140 let id = manager
141 .write()
142 .await
143 .insert_test_running_agent("owned", dir.path());
144 let tool = coord::AgentsFollowupTool::new(Arc::clone(&manager));
145 for input in [
146 json!({"agent_id": id, "agent_ids": [id], "message": "x"}),
147 json!({"agent_ids": [], "message": "x"}),
148 json!({"agent_ids": [id, 7], "message": "x"}),
149 json!({"all_parked": "true", "message": "x"}),
150 json!({"agent_id": id, "message": " "}),
151 ] {
152 assert!(
153 tool.execute(input, &ToolContext::new(dir.path()))
154 .await
155 .is_err()
156 );
157 }
158 assert!(!manager.read().await.child_was_woken(&id));
159 }
160
161 #[tokio::test]
162 async fn lifecycle_resume_lineage_survives_persist_and_rejects_cycles_and_foreign_hops() {
163 let dir = tempdir().unwrap();
164 let base = dir.path().canonicalize().unwrap();
165 let state_path = base.join(".codewhale/subagents/state.json");
166 let mut manager = SubAgentManager::new(base.clone(), 6).with_state_path(state_path.clone());
167 let (a, _) = manager.insert_test_interrupted_continuable_agent("old", &base, prior_messages());
168 let (b, _) =
169 manager.insert_test_interrupted_continuable_agent("continued", &base, prior_messages());
170 let (c, _) =
171 manager.insert_test_interrupted_continuable_agent("latest", &base, prior_messages());
172 manager.resume_targets.insert(a.clone(), b.clone());
173 manager.resume_targets.insert(b.clone(), c.clone());
174 let (path, payload) = manager.build_persist_payload().unwrap().unwrap();
175 write_json_atomic(&base, &path, &payload).unwrap();
176 let mut loaded = SubAgentManager::new(base.clone(), 6).with_state_path(state_path);
177 loaded.load_state().unwrap();
178 assert_eq!(loaded.continuation_target(&a).unwrap(), c);
179 loaded.resume_targets.insert(c.clone(), a.clone());
180 assert!(
181 loaded
182 .continuation_target(&a)
183 .unwrap_err()
184 .to_string()
185 .contains("cycle")
186 );
187 loaded.resume_targets.remove(&c);
188 loaded.assign_test_session_owner(&c, "foreign");
189 assert!(
190 loaded
191 .continuation_target(&a)
192 .unwrap_err()
193 .to_string()
194 .contains("outside")
195 );
196 }
197
198 #[tokio::test]
199 async fn lifecycle_compact_roster_bounds_every_state_and_pages_multibyte_names() {
200 let dir = tempdir().unwrap();
201 let manager = new_shared_subagent_manager(dir.path().to_path_buf(), 50);
202 let mut expected = HashSet::new();
203 {
204 let mut guard = manager.write().await;
205 for i in 0..37 {
206 let id = guard.insert_test_running_agent(&format!("bounded-{i}"), dir.path());
207 expected.insert(id.clone());
208 let agent = guard.agents.get_mut(&id).unwrap();
209 agent.session_name = "🐋\"".repeat(4000);
210 agent.prompt = "archive-only".repeat(10_000);
211 agent.checkpoint = Some(build_subagent_checkpoint(
212 &id,
213 "resume",
214 &prior_messages(),
215 1,
216 true,
217 ));
218 if i % 3 == 1 {
219 agent.status = SubAgentStatus::Interrupted("reason".repeat(10_000));
220 }
221 if i % 3 == 2 {
222 agent.status = SubAgentStatus::Failed("failure".repeat(10_000));
223 }
224 let record = guard.worker_records.get_mut(&id).unwrap();
225 record.usage.total_tokens = Some(100);
226 record.verification.summary = "\"🐋".repeat(10_000);
227 if i > 0 {
228 record.parent_run_id = Some("agent_bounded-0".into());
229 }
230 }
231 }
232 let mut offset = 0;
233 let mut seen = HashSet::new();
234 let mut header = None;
235 loop {
236 let result = inspect_agent_from_input(
237 &json!({"action": "status", "verbose": true, "offset": offset}),
238 Arc::clone(&manager),
239 &ToolContext::new(dir.path()),
240 false,
241 None,
242 )
243 .await
244 .unwrap();
245 assert!(
246 result.content.len() <= lifecycle::COMPACT_STATUS_BYTES,
247 "{}",
248 result.content.len()
249 );
250 let value: Value = serde_json::from_str(&result.content).unwrap();
251 assert_eq!(
252 *header.get_or_insert_with(|| value["columns"].clone()),
253 value["columns"]
254 );
255 assert_eq!(value["total_count"], 37);
256 assert_eq!(
257 value["usage"]["total_tokens"], 3700,
258 "scope totals must not be counted per child"
259 );
260 let rows = status_rows(&value);
261 assert!(!rows.is_empty());
262 for row in rows {
263 assert!(seen.insert(row["agent_id"].as_str().unwrap().to_string()));
264 for key in [
265 "snapshot",
266 "worker_record",
267 "checkpoint",
268 "transcript_handle",
269 ] {
270 assert!(row.get(key).is_none(), "{key}");
271 }
272 }
273 let Some(next) = value["next_offset"].as_u64() else {
274 break;
275 };
276 assert!(next > offset);
277 offset = next;
278 }
279 assert_eq!(seen, expected);
280 }
281
282 #[tokio::test]
283 async fn lifecycle_compact_columns_keep_unknown_usage_and_pending_input_explicit() {
284 let dir = tempdir().unwrap();
285 let mut manager = SubAgentManager::new(dir.path().to_path_buf(), 3);
286 let zero = manager.insert_test_running_agent("zero", dir.path());
287 manager
288 .worker_records
289 .get_mut(&zero)
290 .unwrap()
291 .usage
292 .total_tokens = Some(0);
293 let (pending, _) =
294 manager.insert_test_interrupted_continuable_agent("pending", dir.path(), prior_messages());
295 manager.agents.get_mut(&pending).unwrap().needs_input = Some(SubAgentNeedsInput {
296 question: "Approve the next validation step?".into(),
297 });
298 let payload =
299 lifecycle::compact_roster(&manager, &json!({}), "workspace", false, false).unwrap();
300 let rows = status_rows(&payload);
301 let zero = rows.iter().find(|row| row["agent_id"] == zero).unwrap();
302 let pending = rows.iter().find(|row| row["agent_id"] == pending).unwrap();
303 assert_eq!(zero["total_tokens"], 0);
304 assert!(zero["needs_input"].is_null());
305 assert!(pending["total_tokens"].is_null());
306 assert_eq!(pending["needs_input"], "Approve the next validation step?");
307 assert_eq!(pending["needs_continuation"], true);
308 assert_eq!(payload["usage"]["total_tokens"], 0);
309 assert_eq!(payload["usage"]["reported_workers"], 1);
310 let empty = SubAgentManager::new(dir.path().join("empty"), 1);
311 let empty = lifecycle::compact_roster(&empty, &json!({}), "workspace", false, false).unwrap();
312 assert_eq!(empty["columns"], payload["columns"]);
313 assert!(status_rows(&empty).is_empty());
314 }
315
316 #[tokio::test]
317 async fn lifecycle_addressed_status_follows_lineage_and_detail_is_bounded() {
318 let dir = tempdir().unwrap();
319 let manager = new_shared_subagent_manager(dir.path().to_path_buf(), 3);
320 let (old, latest) = {
321 let mut guard = manager.write().await;
322 let (old, _) =
323 guard.insert_test_interrupted_continuable_agent("old", dir.path(), prior_messages());
324 let latest = guard.insert_test_running_agent("latest", dir.path());
325 guard.resume_targets.insert(old.clone(), latest.clone());
326 guard.agents.get_mut(&latest).unwrap().result = Some("🐋".repeat(100_000));
327 (old, latest)
328 };
329 let context = ToolContext::new(dir.path());
330 let roster = inspect_agent_from_input(
331 &json!({"action": "status"}),
332 Arc::clone(&manager),
333 &context,
334 false,
335 None,
336 )
337 .await
338 .unwrap();
339 let roster: Value = serde_json::from_str(&roster.content).unwrap();
340 let rows = status_rows(&roster);
341 assert_eq!(
342 rows.iter().find(|row| row["agent_id"] == old).unwrap()["resumed_as"],
343 latest
344 );
345 assert_eq!(
346 rows.iter().find(|row| row["agent_id"] == latest).unwrap()["resumed_from"],
347 old
348 );
349 let result = inspect_agent_from_input(
350 &json!({"agent_id": old}),
351 Arc::clone(&manager),
352 &context,
353 false,
354 None,
355 )
356 .await
357 .unwrap();
358 let row: Value = serde_json::from_str(&result.content).unwrap();
359 assert_eq!(row["agent_id"], latest);
360 assert_eq!(row["addressed_agent_id"], old);
361 assert_eq!(row["resumed_from"], old);
362 assert!(row.get("snapshot").is_none());
363 let detail = inspect_agent_from_input(
364 &json!({"agent_id": old, "detail": true}),
365 manager,
366 &context,
367 false,
368 None,
369 )
370 .await
371 .unwrap();
372 assert!(detail.content.len() <= 32 * 1024);
373 let row: Value = serde_json::from_str(&detail.content).unwrap();
374 assert_eq!(row["detail_bounded"], true);
375 assert!(row["transcript_handle"].is_object());
376 }
377
378 #[tokio::test]
379 async fn lifecycle_named_cancel_stops_grandchildren_and_preserves_sibling() {
380 let dir = tempdir().unwrap();
381 let mut manager = SubAgentManager::new(dir.path().to_path_buf(), 4);
382 let parent = manager.insert_test_running_agent("parent", dir.path());
383 let child = manager.insert_test_running_agent("child", dir.path());
384 let sibling = manager.insert_test_running_agent("sibling", dir.path());
385 let record = manager.worker_records.get_mut(&child).unwrap();
386 record.parent_run_id = Some(parent.clone());
387 record.spec.parent_run_id = Some(parent.clone());
388 manager
389 .cancel_agent_for_session("workspace", &parent)
390 .unwrap();
391 assert_eq!(
392 manager.get_result(&child).unwrap().status,
393 SubAgentStatus::Cancelled
394 );
395 assert_eq!(
396 manager.get_result(&parent).unwrap().status,
397 SubAgentStatus::Cancelled
398 );
399 assert_eq!(
400 manager.get_result(&sibling).unwrap().status,
401 SubAgentStatus::Running
402 );
403 }
404
405 #[test]
406 fn lifecycle_recovery_never_forks_and_byte_preview_preserves_utf8() {
407 let instruction = subagent_followup_recovery("agent_parked");
408 assert!(instruction.contains("action=\"followup\""));
409 assert!(!instruction.contains("resume_from"));
410 let preview = lifecycle::text_preview(&"🐋".repeat(10_000), 65);
411 assert!(preview.len() <= 65);
412 assert!(preview.ends_with("..."));
413 }
414
415 #[test]
416 fn lifecycle_park_event_sentinel_and_tool_share_recovery_instruction() {
417 let agent_id = "agent_parked";
418 let instruction = subagent_followup_recovery(agent_id);
419 let parking = Arc::new(std::sync::atomic::AtomicBool::new(true));
420 let (status, output, checkpoint, needs_input, _, _) =
421 subagent_cancellation_projection(agent_id, &prior_messages(), 1, None, Some(&parking));
422 assert!(output.as_deref().unwrap().contains(&instruction));
423 assert_eq!(needs_input.as_ref().unwrap().question, instruction);
424
425 let mut result = super::tests::make_snapshot(status);
426 result.agent_id = agent_id.into();
427 result.result = output;
428 result.checkpoint = checkpoint;
429 result.needs_input = needs_input;
430 let sentinel = subagent_done_sentinel(agent_id, &result, false);
431 let metadata: Value = serde_json::from_str(
432 sentinel
433 .strip_prefix("<codewhale:subagent.done>")
434 .unwrap()
435 .strip_suffix("</codewhale:subagent.done>")
436 .unwrap(),
437 )
438 .unwrap();
439 assert_eq!(metadata["needs_input"]["question"], instruction);
440 assert!(AGENT_TOOL_DESCRIPTION.contains(&subagent_followup_recovery("<agent_id>")));
441 }
442
443 #[tokio::test]
444 async fn lifecycle_followup_rechecks_actual_successor_authority() {
445 let dir = tempdir().unwrap();
446 let manager = new_shared_subagent_manager(dir.path().to_path_buf(), 4);
447 let (caller, old, sibling) = {
448 let mut guard = manager.write().await;
449 let caller = guard.insert_test_running_agent("caller", dir.path());
450 let (old, _) = guard.insert_test_interrupted_continuable_agent(
451 "own-child",
452 dir.path(),
453 prior_messages(),
454 );
455 let sibling = guard.insert_test_running_agent("sibling", dir.path());
456 let record = guard.worker_records.get_mut(&old).unwrap();
457 record.parent_run_id = Some(caller.clone());
458 record.spec.parent_run_id = Some(caller.clone());
459 guard.resume_targets.insert(old.clone(), sibling.clone());
460 (caller, old, sibling)
461 };
462 let tool =
463 coord::AgentsFollowupTool::new(Arc::clone(&manager)).with_optional_caller(Some(caller));
464 assert!(
465 tool.execute(
466 json!({"agent_id": old, "message": "Try to wake sibling"}),
467 &ToolContext::new(dir.path())
468 )
469 .await
470 .is_err()
471 );
472 assert!(!manager.read().await.child_was_woken(&sibling));
473 }
474
475 #[tokio::test]
476 async fn lifecycle_deliverable_preview_reports_omissions_and_detail_pages_the_full_list() {
477 let dir = tempdir().unwrap();
478 let mut manager = SubAgentManager::new(dir.path().to_path_buf(), 2);
479 let id = manager.insert_test_running_agent("outputs", dir.path());
480 manager
481 .worker_records
482 .get_mut(&id)
483 .unwrap()
484 .verification
485 .deliverables = (0..9)
486 .map(|index| DeliverableVerdict {
487 path: format!("report-{index}.md"),
488 status: if index == 8 {
489 "missing".into()
490 } else {
491 "present".into()
492 },
493 bytes: (index != 8).then_some(20),
494 })
495 .collect();
496 let compact = lifecycle::compact_row(&manager, &manager.agents[&id]);
497 assert_eq!(compact["verification"]["deliverables_total"], 9);
498 assert_eq!(compact["verification"]["deliverables_omitted"], 5);
499 assert_eq!(compact["verification"]["deliverable_counts"]["missing"], 1);
500 assert_eq!(
501 compact["verification"]["deliverables"][0]["status"],
502 "missing"
503 );
504 let detail = lifecycle::bounded_detail(
505 json!({"verification": manager.worker_records[&id].verification}),
506 compact,
507 4,
508 2,
509 );
510 assert_eq!(
511 detail["verification"]["deliverables"]
512 .as_array()
513 .unwrap()
514 .len(),
515 2
516 );
517 assert_eq!(
518 detail["verification"]["deliverables"][0]["path"],
519 "report-4.md"
520 );
521 assert_eq!(detail["verification"]["deliverables_next_offset"], 6);
522 }
523
524 #[tokio::test]
525 async fn lifecycle_running_row_reports_declared_vs_observed_writes() {
526 // #6194 item 5: the parent sees the declared/observed write diff while
527 // the child is still alive, not only in the post-mortem receipt.
528 let dir = tempdir().unwrap();
529 let mut manager = SubAgentManager::new(dir.path().to_path_buf(), 2);
530 let id = manager.insert_test_running_agent("writer", dir.path());
531 {
532 let record = manager.worker_records.get_mut(&id).unwrap();
533 record.spec.runtime_profile.permissions.write = true;
534 record.spec.launch_manifest = Some(ChildLaunchManifest {
535 owner_session: "root".to_string(),
536 child_id: id.clone(),
537 profile: record.spec.runtime_profile.clone(),
538 prompt: record.spec.objective.clone(),
539 cwd: None,
540 worktree: false,
541 writable_roots: Vec::new(),
542 writable_files: Vec::new(),
543 coordination_contracts: Vec::new(),
544 expected_artifact: None,
545 deliverables: vec!["a.rs".to_string(), "b.rs".to_string()],
546 resume_identity: None,
547 generation: 1,
548 resume_from_agent_id: None,
549 });
550 record
551 .delivery_evidence
552 .observed_writes
553 .insert("b.rs".to_string());
554 record
555 .delivery_evidence
556 .observed_writes
557 .insert("surprise.rs".to_string());
558 }
559 let row = lifecycle::compact_row(&manager, &manager.agents[&id]);
560 assert_eq!(row["write_progress"]["declared_total"], 2);
561 assert_eq!(row["write_progress"]["observed_total"], 2);
562 assert_eq!(row["write_progress"]["declared"][0], "a.rs");
563 assert_eq!(row["write_progress"]["observed"][1], "surprise.rs");
564 }
565
566 #[tokio::test]
567 async fn lifecycle_read_only_row_omits_write_progress() {
568 let dir = tempdir().unwrap();
569 let mut manager = SubAgentManager::new(dir.path().to_path_buf(), 2);
570 let id = manager.insert_test_running_agent("scout", dir.path());
571 let row = lifecycle::compact_row(&manager, &manager.agents[&id]);
572 assert!(row.get("write_progress").is_none());
573 }
574
575 #[tokio::test]
576 async fn lifecycle_continuation_link_is_durable_before_the_child_can_run() {
577 let dir = tempdir().unwrap();
578 let base = dir.path().canonicalize().unwrap();
579 let path = base.join(".codewhale/subagents/state.json");
580 let manager = Arc::new(RwLock::new(
581 SubAgentManager::new(base.clone(), 4).with_state_path(path.clone()),
582 ));
583 let (client, calls, _, fixture_config) =
584 super::tests::delayed_chat_client(Duration::from_secs(30), "fixture").await;
585 let mut runtime = super::tests::stub_runtime();
586 runtime.client = client;
587 runtime.api_config = Some(Arc::new(fixture_config));
588 runtime.manager = Arc::clone(&manager);
589 runtime.context = ToolContext::new(&base);
590 let mut guard = manager.write().await;
591 let (source, _) =
592 guard.insert_test_interrupted_continuable_agent("durable-source", &base, prior_messages());
593 guard.agents.get_mut(&source).unwrap().model = "deepseek-v4-flash".into();
594 let record = guard.worker_records.get_mut(&source).unwrap();
595 record.parent_run_id = None;
596 record.spec.parent_run_id = None;
597 guard
598 .worker_records
599 .get_mut(&source)
600 .unwrap()
601 .spec
602 .runtime_profile = WorkerRuntimeProfile::for_role(FleetRole::Scout);
603 assert!(
604 !guard.worker_records[&source]
605 .spec
606 .runtime_profile
607 .permissions
608 .write
609 );
610 let successor = guard
611 .resume_from_checkpoint(Arc::clone(&manager), runtime, &source, "Continue")
612 .unwrap();
613 // Holding the manager lock keeps run_subagent_task_inner at its first
614 // await. This is the earliest published snapshot, before any child step.
615 assert_eq!(calls.load(std::sync::atomic::Ordering::SeqCst), 0);
616 let persisted: PersistedSubAgentState =
617 serde_json::from_slice(&std::fs::read(path).unwrap()).unwrap();
618 assert_eq!(
619 persisted.resume_targets.get(&source),
620 Some(&successor.agent_id)
621 );
622 assert!(
623 persisted
624 .agents
625 .iter()
626 .any(|agent| agent.id == successor.agent_id)
627 );
628 let _ = guard.cancel_agent(&successor.agent_id);
629 }
630
631 #[tokio::test]
632 async fn lifecycle_continuation_persist_failure_rolls_back_worker_and_link() {
633 let dir = tempdir().unwrap();
634 let base = dir.path().canonicalize().unwrap();
635 let path = base.join(".codewhale/subagents/state.json");
636 let manager = Arc::new(RwLock::new(
637 SubAgentManager::new(base.clone(), 4).with_state_path(path),
638 ));
639 let mut runtime = super::tests::stub_runtime();
640 runtime.manager = Arc::clone(&manager);
641 runtime.context = ToolContext::new(&base);
642 let mut guard = manager.write().await;
643 let (source, _) =
644 guard.insert_test_interrupted_continuable_agent("failed-source", &base, prior_messages());
645 guard.agents.get_mut(&source).unwrap().model = "deepseek-v4-flash".into();
646 guard
647 .worker_records
648 .get_mut(&source)
649 .unwrap()
650 .spec
651 .runtime_profile = WorkerRuntimeProfile::for_role(FleetRole::Scout);
652 assert!(
653 !guard.worker_records[&source]
654 .spec
655 .runtime_profile
656 .permissions
657 .write
658 );
659 std::fs::create_dir_all(base.join(".codewhale")).unwrap();
660 std::fs::write(base.join(".codewhale/subagents"), "not a directory").unwrap();
661 let result = guard.resume_from_checkpoint(Arc::clone(&manager), runtime, &source, "Continue");
662 assert!(result.is_err());
663 assert_eq!(guard.agents.len(), 1);
664 assert_eq!(guard.worker_records.len(), 1);
665 assert!(guard.resume_targets.is_empty());
666 assert!(matches!(
667 guard.get_result(&source).unwrap().status,
668 SubAgentStatus::Interrupted(_)
669 ));
670 }
671
672 #[tokio::test]
673 async fn lifecycle_cancel_original_stops_its_continuation_and_existing_descendants() {
674 let dir = tempdir().unwrap();
675 let mut manager = SubAgentManager::new(dir.path().to_path_buf(), 5);
676 let (source, _) =
677 manager.insert_test_interrupted_continuable_agent("original", dir.path(), prior_messages());
678 let current = manager.insert_test_running_agent("current", dir.path());
679 let child = manager.insert_test_running_agent("existing-child", dir.path());
680 let sibling = manager.insert_test_running_agent("separate-fork", dir.path());
681 manager
682 .resume_targets
683 .insert(source.clone(), current.clone());
684 let record = manager.worker_records.get_mut(&child).unwrap();
685 record.parent_run_id = Some(source.clone());
686 record.spec.parent_run_id = Some(source.clone());
687
688 let cancelled = manager
689 .cancel_agent_for_session("workspace", &source)
690 .unwrap();
691 assert_eq!(cancelled.agent_id, current);
692 assert_eq!(cancelled.status, SubAgentStatus::Cancelled);
693 assert_eq!(
694 manager.get_result(&child).unwrap().status,
695 SubAgentStatus::Cancelled
696 );
697 assert_eq!(
698 manager.get_result(&sibling).unwrap().status,
699 SubAgentStatus::Running
700 );
701 assert!(matches!(
702 manager.get_result(&source).unwrap().status,
703 SubAgentStatus::Interrupted(_)
704 ));
705 }
706
707 #[tokio::test]
708 async fn lifecycle_resumed_parent_controls_existing_descendants_but_never_itself_or_siblings() {
709 let dir = tempdir().unwrap();
710 let mut manager = SubAgentManager::new(dir.path().to_path_buf(), 5);
711 let (source, _) =
712 manager.insert_test_interrupted_continuable_agent("original", dir.path(), prior_messages());
713 let current = manager.insert_test_running_agent("current", dir.path());
714 let child = manager.insert_test_running_agent("existing-child", dir.path());
715 let sibling = manager.insert_test_running_agent("sibling", dir.path());
716 manager
717 .resume_targets
718 .insert(source.clone(), current.clone());
719 let record = manager.worker_records.get_mut(&child).unwrap();
720 record.parent_run_id = Some(source.clone());
721 record.spec.parent_run_id = Some(source.clone());
722 assert!(
723 manager
724 .continuation_target_for_caller("workspace", &child, Some(&current), "test")
725 .is_ok()
726 );
727 for forbidden in [&source, &current, &sibling] {
728 assert!(
729 manager
730 .continuation_target_for_caller("workspace", forbidden, Some(&current), "test")
731 .is_err()
732 );
733 }
734
735 // A corrupt source→sibling link must not turn source authority into
736 // permission to cancel the unrelated actual successor.
737 manager
738 .resume_targets
739 .insert(child.clone(), sibling.clone());
740 let manager = Arc::new(RwLock::new(manager));
741 assert!(
742 cancel_agent_from_input(
743 &json!({"agent_id": child}),
744 Arc::clone(&manager),
745 &ToolContext::new(dir.path()),
746 Some(&current),
747 )
748 .await
749 .is_err()
750 );
751 assert_eq!(
752 manager.read().await.get_result(&sibling).unwrap().status,
753 SubAgentStatus::Running
754 );
755 }
756
757 #[tokio::test]
758 async fn lifecycle_detail_budget_keeps_the_transcript_handle_retrievable() {
759 let dir = tempdir().unwrap();
760 let context = ToolContext::new(dir.path());
761 let mut manager = SubAgentManager::new(dir.path().to_path_buf(), 2);
762 let id = manager.insert_test_running_agent("large-checkpoint", dir.path());
763 let messages = (0..20)
764 .map(|_| Message {
765 role: Role::User,
766 content: vec![ContentBlock::Text {
767 text: "🐋".repeat(1024),
768 cache_control: None,
769 }],
770 })
771 .collect::<Vec<_>>();
772 manager.agents.get_mut(&id).unwrap().checkpoint = Some(build_subagent_checkpoint(
773 &id, "retained", &messages, 1, true,
774 ));
775 let projection = subagent_session_projection(
776 &new_shared_subagent_manager(dir.path().to_path_buf(), 1),
777 manager.get_result(&id).unwrap(),
778 false,
779 &context,
780 manager.get_worker_record_for_session("workspace", &id),
781 )
782 .await;
783 let expected = projection.transcript_handle.clone();
784 let detail = lifecycle::bounded_detail(
785 serde_json::to_value(projection).unwrap(),
786 lifecycle::compact_row(&manager, &manager.agents[&id]),
787 0,
788 20,
789 );
790 assert!(serde_json::to_vec(&detail).unwrap().len() <= 32 * 1024);
791 let handle: VarHandle = serde_json::from_value(detail["transcript_handle"].clone()).unwrap();
792 assert_eq!(handle.session_id, expected.session_id);
793 assert_eq!(handle.name, expected.name);
794 assert_eq!(handle.sha256, expected.sha256);
795 assert!(
796 context
797 .runtime
798 .handle_store
799 .lock()
800 .await
801 .get(&handle)
802 .is_some()
803 );
804 }
805
806 #[tokio::test]
807 async fn lifecycle_cleanup_keeps_source_identity_while_descendants_still_need_it() {
808 let dir = tempdir().unwrap();
809 let mut manager = SubAgentManager::new(dir.path().to_path_buf(), 5);
810 let source = manager.insert_test_running_agent("original", dir.path());
811 let current = manager.insert_test_running_agent("continuation", dir.path());
812 let child = manager.insert_test_running_agent("existing-child", dir.path());
813 manager
814 .resume_targets
815 .insert(source.clone(), current.clone());
816 let record = manager.worker_records.get_mut(&child).unwrap();
817 record.parent_run_id = Some(source.clone());
818 record.spec.parent_run_id = Some(source.clone());
819 let old = Instant::now() - Duration::from_secs(120);
820 manager.agents.get_mut(&source).unwrap().status = SubAgentStatus::Completed;
821 manager.agents.get_mut(&source).unwrap().started_at = old;
822 let record = manager.worker_records.get_mut(&source).unwrap();
823 record.status = AgentWorkerStatus::Completed;
824 record.completed_at_ms = Some(epoch_millis_now().saturating_sub(120_000));
825
826 manager.cleanup_for_session("workspace", Duration::from_secs(60));
827 assert_eq!(manager.continuation_target(&source).unwrap(), current);
828 assert!(
829 manager
830 .continuation_target_for_caller("workspace", &child, Some(&current), "test")
831 .is_ok()
832 );
833 manager
834 .cancel_agent_for_session("workspace", &source)
835 .unwrap();
836 assert_eq!(
837 manager.get_result(&child).unwrap().status,
838 SubAgentStatus::Cancelled
839 );
840
841 // Once no live work uses either projection, normal expiry reclaims the
842 // archived workers and their continuation edge together.
843 for agent in manager.agents.values_mut() {
844 agent.started_at = old;
845 }
846 for record in manager.worker_records.values_mut() {
847 record.completed_at_ms = Some(epoch_millis_now().saturating_sub(120_000));
848 }
849 manager.cleanup_for_session("workspace", Duration::from_secs(60));
850 assert!(manager.agents.is_empty());
851 assert!(manager.worker_records.is_empty());
852 assert!(manager.resume_targets.is_empty());
853 }
854
854 lines RUST