返回 CodeWhale
task_ownership.rs
根目录 / crates / tui / src / runtime_threads / tests / task_ownership.rs
1 use super::*;
2 use crate::task_manager::{TaskManager, TaskManagerConfig};
3
4 fn fixture_config() -> Config {
5 let mut config = Config {
6 ..Config::default()
7 }
8 .with_legacy_root(
9 Some("local-fixture-key".into()),
10 Some("http://127.0.0.1:1/v1".into()),
11 );
12 config.set_feature("mcp", false).unwrap();
13 config.set_feature("subagents", false).unwrap();
14 config
15 }
16
17 #[tokio::test]
18 async fn idle_runtime_engine_cycles_are_drained_before_same_scope_reopen() -> Result<()> {
19 let _env = crate::test_support::lock_test_env();
20 let root = tempfile::tempdir()?;
21 let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", root.path());
22 let cfg = fixture_config();
23 let runtime_config = test_manager_config(root.path().join("runtime"));
24 let runtime = Arc::new(RuntimeThreadManager::open(
25 cfg.clone(),
26 root.path().into(),
27 runtime_config.clone(),
28 )?);
29 let tasks = TaskManager::start_with_runtime_manager(
30 TaskManagerConfig::from_runtime(&cfg, root.path().into(), None, Some(1)),
31 cfg.clone(),
32 runtime.clone(),
33 )
34 .await?;
35 let thread = runtime
36 .create_thread(CreateThreadRequest::default())
37 .await?;
38 // Load the real production Engine with its actual strong Runtime services.
39 // No turn is submitted, so constructing this Engine makes no provider call.
40 let engine = runtime.get_engine(&thread.id).await?;
41 engine.get_session_snapshot().await?;
42 runtime.spawn_goal_continuation(thread.id.clone(), 3_600);
43 let weak = Arc::downgrade(&tasks);
44 tasks.shutdown_and_wait().await?;
45 assert!(runtime.active.lock().await.engines.is_empty());
46 assert!(runtime.get_engine(&thread.id).await.is_err());
47 assert!(
48 RuntimeThreadManager::open(cfg.clone(), root.path().into(), runtime_config.clone())
49 .is_err(),
50 "retained Runtime handle still owns its scope"
51 );
52 drop(engine);
53 drop(tasks);
54 drop(runtime);
55 assert!(
56 weak.upgrade().is_none(),
57 "idle Engine service cycle must have ended"
58 );
59 let reopened = RuntimeThreadManager::open(cfg, root.path().into(), runtime_config)?;
60 reopened.shutdown_and_wait().await?;
61 Ok(())
62 }
63
64 #[tokio::test]
65 async fn shutdown_waits_actual_engine_exit_and_terminal_monitor_before_releasing_scope()
66 -> Result<()> {
67 use crate::core::engine::Engine;
68 use crate::llm_client::mock::{MockLlmClient, canned};
69 let _env = crate::test_support::lock_test_env();
70 let root = tempfile::tempdir()?;
71 let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", root.path());
72 let cfg = fixture_config();
73 let manager_cfg = test_manager_config(root.path().join("runtime"));
74 let runtime = Arc::new(RuntimeThreadManager::open(
75 cfg.clone(),
76 root.path().into(),
77 manager_cfg.clone(),
78 )?);
79 let tasks = TaskManager::start_with_runtime_manager(
80 TaskManagerConfig::from_runtime(&cfg, root.path().into(), None, Some(1)),
81 cfg.clone(),
82 runtime.clone(),
83 )
84 .await?;
85 let thread = runtime
86 .create_thread(CreateThreadRequest::default())
87 .await?;
88 let model = Arc::new(MockLlmClient::new(vec![canned::simple_text_turn(
89 "owned fixture",
90 )]));
91 let (engine, handle) = Engine::new_with_model_client(
92 EngineConfig {
93 workspace: root.path().into(),
94 model: thread.model.clone(),
95 subagents_enabled: false,
96 snapshots_enabled: false,
97 memory_enabled: false,
98 terminal_chrome_enabled: false,
99 runtime_services: crate::tools::spec::RuntimeToolServices {
100 task_manager: Some(tasks.clone()),
101 active_thread_id: Some(thread.id.clone()),
102 dynamic_tool_executor: Some(Arc::new(runtime.as_ref().clone())),
103 ..Default::default()
104 },
105 ..EngineConfig::default()
106 },
107 &cfg,
108 model.clone(),
109 );
110 runtime
111 .install_test_engine(&thread.id, handle.clone())
112 .await?;
113 let (release, blocked) = oneshot::channel();
114 let worker = tokio::spawn(async move {
115 blocked.await.unwrap();
116 engine.run().await;
117 });
118 runtime
119 .engine_workers
120 .lock()
121 .push((handle.clone(), retained_completion(worker)));
122 let turn = runtime
123 .start_turn(
124 &thread.id,
125 StartTurnRequest {
126 prompt: "fixture admission".into(),
127 ..Default::default()
128 },
129 )
130 .await?;
131 let drain = tokio::spawn({
132 let tasks = tasks.clone();
133 async move { tasks.shutdown_and_wait().await }
134 });
135 tokio::time::timeout(Duration::from_secs(5), runtime.cancel_token.cancelled()).await?;
136 assert!(
137 !drain.is_finished(),
138 "a cancellation request cannot prove Engine exit"
139 );
140 assert!(
141 RuntimeThreadManager::open(cfg.clone(), root.path().into(), manager_cfg.clone()).is_err()
142 );
143 assert!(
144 runtime
145 .start_turn(
146 &thread.id,
147 StartTurnRequest {
148 prompt: "late admission".into(),
149 ..Default::default()
150 }
151 )
152 .await
153 .is_err()
154 );
155 let completion = runtime.engine_workers.lock()[0].1.clone();
156 tokio::time::timeout(Duration::from_secs(5), async {
157 while completion.try_lock().is_ok() {
158 sleep(Duration::from_millis(5)).await;
159 }
160 })
161 .await?;
162 drain.abort();
163 assert!(drain.await.unwrap_err().is_cancelled());
164 let retry = tokio::spawn({
165 let tasks = tasks.clone();
166 async move { tasks.shutdown_and_wait().await }
167 });
168 sleep(Duration::from_millis(25)).await;
169 assert!(
170 !retry.is_finished(),
171 "retry must still own and await the original Engine join"
172 );
173 release.send(()).unwrap();
174 tokio::time::timeout(Duration::from_secs(15), retry).await???;
175 drop(completion);
176 assert!(runtime.store.load_turn(&turn.id)?.status != RuntimeTurnStatus::InProgress);
177 drop(handle);
178 drop(tasks);
179 drop(runtime);
180 let reopened = RuntimeThreadManager::open(cfg, root.path().into(), manager_cfg)?;
181 reopened.shutdown_and_wait().await?;
182 Ok(())
183 }
184
185 #[tokio::test]
186 async fn shutdown_drains_accepted_user_input_receipt_after_caller_disconnects() -> Result<()> {
187 let root = tempfile::tempdir()?;
188 let runtime = Arc::new(test_manager(root.path().join("runtime"))?);
189 let thread = runtime
190 .create_thread(CreateThreadRequest::default())
191 .await?;
192 let harness = mock_engine_handle();
193 runtime
194 .install_test_engine(&thread.id, harness.handle.clone())
195 .await?;
196 runtime.register_pending_user_input(
197 &thread.id,
198 PendingUserInputRequest {
199 id: "input_drain".into(),
200 turn_id: "turn_drain".into(),
201 request: crate::tools::user_input::UserInputRequest {
202 questions: Vec::new(),
203 },
204 },
205 );
206 let hold_receipt = runtime.event_emit.lock().await;
207 let submission = tokio::spawn({
208 let runtime = runtime.clone();
209 let thread_id = thread.id.clone();
210 async move {
211 runtime
212 .submit_user_input(
213 &thread_id,
214 "input_drain",
215 crate::tools::user_input::UserInputResponse {
216 answers: Vec::new(),
217 },
218 )
219 .await
220 }
221 });
222 tokio::time::timeout(Duration::from_secs(5), async {
223 while runtime.turn_monitors.lock().is_empty() {
224 sleep(Duration::from_millis(5)).await;
225 }
226 })
227 .await?;
228 submission.abort();
229 assert!(submission.await.unwrap_err().is_cancelled());
230 let drain = tokio::spawn({
231 let runtime = runtime.clone();
232 async move { runtime.shutdown_and_wait().await }
233 });
234 tokio::time::timeout(Duration::from_secs(5), runtime.cancel_token.cancelled()).await?;
235 sleep(Duration::from_millis(25)).await;
236 assert!(
237 !drain.is_finished(),
238 "accepted detached receipt is still part of shutdown"
239 );
240 drop(hold_receipt);
241 tokio::time::timeout(Duration::from_secs(5), drain).await???;
242 let events = runtime.events_since(&thread.id, None)?;
243 assert_eq!(
244 events
245 .iter()
246 .filter(|event| event.event == "user_input.answered")
247 .count(),
248 1
249 );
250 assert!(runtime.pending_user_inputs.lock().is_empty());
251 Ok(())
252 }
253
254 #[tokio::test]
255 async fn shutdown_fences_recovery_readers_without_publishing_late_receipts() -> Result<()> {
256 let root = tempfile::tempdir()?;
257 let runtime = Arc::new(test_manager(root.path().join("runtime"))?);
258 let thread = runtime
259 .create_thread(CreateThreadRequest::default())
260 .await?;
261 let turn = sample_turn(&thread.id, "turn_recovery_drain", RuntimeTurnStatus::Failed);
262 runtime.store.save_turn(&turn)?;
263 runtime.queue_recovery_receipt(RecoveredTurnReceipt {
264 turn,
265 unresolved_dynamic_tools: Vec::new(),
266 });
267 let hold = runtime.recovery_flush.lock().await;
268 let reader = tokio::spawn({
269 let runtime = runtime.clone();
270 let id = thread.id.clone();
271 async move { runtime.get_thread(&id).await }
272 });
273 tokio::task::yield_now().await;
274 let shutdown = tokio::spawn({
275 let runtime = runtime.clone();
276 async move { runtime.shutdown_and_wait().await }
277 });
278 tokio::time::timeout(Duration::from_secs(5), runtime.cancel_token.cancelled()).await?;
279 assert!(
280 !shutdown.is_finished(),
281 "shutdown must fence pending recovery producers"
282 );
283 drop(hold);
284 assert!(
285 tokio::time::timeout(Duration::from_secs(5), reader)
286 .await??
287 .is_err()
288 );
289 tokio::time::timeout(Duration::from_secs(5), shutdown).await???;
290 assert!(
291 runtime.recovery_receipts.lock().contains_key(&thread.id),
292 "unadmitted recovery remains queued for the next owner"
293 );
294 assert!(
295 runtime
296 .events_since(&thread.id, None)?
297 .iter()
298 .all(|event| event.event != "turn.completed")
299 );
300 assert!(runtime.get_thread_detail(&thread.id).await.is_err());
301 Ok(())
302 }
303
304 #[tokio::test]
305 async fn shutdown_drains_an_already_started_recovery_flush() -> Result<()> {
306 let root = tempfile::tempdir()?;
307 let runtime = Arc::new(test_manager(root.path().join("runtime"))?);
308 let thread = runtime
309 .create_thread(CreateThreadRequest::default())
310 .await?;
311 let turn = sample_turn(
312 &thread.id,
313 "turn_started_recovery",
314 RuntimeTurnStatus::Failed,
315 );
316 runtime.store.save_turn(&turn)?;
317 runtime.queue_recovery_receipt(RecoveredTurnReceipt {
318 turn,
319 unresolved_dynamic_tools: Vec::new(),
320 });
321 let hold = runtime.event_emit.lock().await;
322 let reader = tokio::spawn({
323 let runtime = runtime.clone();
324 let id = thread.id.clone();
325 async move { runtime.get_thread(&id).await }
326 });
327 tokio::time::timeout(Duration::from_secs(5), async {
328 while runtime.recovery_flush.try_lock().is_ok() {
329 sleep(Duration::from_millis(5)).await;
330 }
331 })
332 .await?;
333 let shutdown = tokio::spawn({
334 let runtime = runtime.clone();
335 async move { runtime.shutdown_and_wait().await }
336 });
337 tokio::time::timeout(Duration::from_secs(5), runtime.cancel_token.cancelled()).await?;
338 sleep(Duration::from_millis(25)).await;
339 assert!(
340 !shutdown.is_finished(),
341 "accepted recovery write must settle before drain succeeds"
342 );
343 drop(hold);
344 tokio::time::timeout(Duration::from_secs(5), reader).await???;
345 tokio::time::timeout(Duration::from_secs(5), shutdown).await???;
346 assert!(!runtime.recovery_receipts.lock().contains_key(&thread.id));
347 assert_eq!(
348 runtime
349 .events_since(&thread.id, None)?
350 .iter()
351 .filter(|event| event.event == "turn.completed")
352 .count(),
353 1
354 );
355 Ok(())
356 }
357
357 lines RUST