返回 CodeWhale
runtime_store_convergence.rs
根目录 / crates / tui / src / runtime_api / tests / runtime_store_convergence.rs
1 //! Real Engine/store through the actual authenticated owner dispatcher. Only
2 //! the model is scripted; this is not a fake HTTP bridge or an independent lease.
3 use super::*;
4 use crate::llm_client::mock::{MockLlmClient, canned};
5 use crate::session_manager::SessionManager;
6
7 struct AbortOnDrop(tokio::task::AbortHandle);
8 impl Drop for AbortOnDrop {
9 fn drop(&mut self) {
10 self.0.abort();
11 }
12 }
13
14 async fn reply(
15 guest: &mut codewhale_app_server::daemon_client::OwnerClient,
16 id: &str,
17 ) -> Result<Value> {
18 tokio::time::timeout(ci_scaled(Duration::from_secs(60)), async {
19 loop {
20 let frame = guest.recv().await?.context("owner closed before reply")?;
21 if frame["id"] == id {
22 return Ok(frame);
23 }
24 }
25 })
26 .await
27 .context("owner reply timed out")?
28 }
29
30 #[tokio::test]
31 async fn authenticated_owner_guest_runs_real_engine_and_reuses_held_store() -> Result<()> {
32 let _env = lock_test_env();
33 let temporary = tempfile::tempdir()?;
34 let root = temporary.path().canonicalize()?;
35 let _home = EnvVarGuard::set("CODEWHALE_HOME", root.join("home"));
36 let config_path = root.join("config.toml");
37 fs::write(&config_path, "model = 'deepseek-v4-pro'\n")?;
38 let config = Config::default().with_legacy_root(
39 Some("owner-fixture-key".into()),
40 Some("http://127.0.0.1:1/v1".into()),
41 );
42 #[cfg(unix)]
43 let socket_path = root.join("private-run").join("owner.sock");
44 #[cfg(windows)]
45 let socket_path = PathBuf::from(format!(r"\\.\pipe\codewhale-test-{}", Uuid::new_v4()));
46 let shutdown = RuntimeServerShutdown::default();
47 let token = format!("owner-fixture-{}", Uuid::new_v4());
48 let (addr, threads, http) = spawn_test_server_with_root_token_mobile_workspace_and_overrides(
49 root.clone(),
50 root.join("sessions"),
51 Some(token.clone()),
52 false,
53 root.join("workspace"),
54 TestServerOverrides {
55 config: Some(config),
56 config_path: Some(config_path.clone()),
57 owner_socket: Some(socket_path.clone()),
58 shutdown: Some(shutdown.clone()),
59 ..TestServerOverrides::default()
60 },
61 )
62 .await?
63 .context("owner acceptance requires a real loopback listener")?;
64 let _http_abort = AbortOnDrop(http.abort_handle());
65 let model = Arc::new(MockLlmClient::new(Vec::new()));
66 threads.set_test_model_client(model.clone());
67 let manager = threads.clone();
68 let (binding, generation) =
69 codewhale_app_server::daemon_socket::owner_work(move || manager.capture_control_owner())
70 .await?;
71 let owner_client = codewhale_app_server::daemon_client::connect(
72 Some(config_path.clone()),
73 Some(socket_path.clone()),
74 )
75 .await?;
76 let captured = owner_client.receipt().clone();
77 assert_eq!(captured.data_dir, binding.data_dir);
78 assert_eq!(captured.execution_scope, binding.execution_scope);
79 assert_eq!(captured.lease_generation, generation);
80 assert_eq!(captured.socket_path, socket_path);
81 assert_eq!(captured.config_path, Some(config_path.clone()));
82 assert_eq!(captured.pid, std::process::id());
83 assert_eq!(
84 captured.process_start,
85 codewhale_app_server::daemon_socket::capture_process_start(std::process::id()).await?
86 );
87 #[cfg(unix)]
88 assert_eq!(
89 captured.principal,
90 codewhale_config::private_directory::PrivateDirectory::current_user_id().to_string()
91 );
92 #[cfg(windows)]
93 assert_eq!(
94 captured.principal,
95 codewhale_app_server::daemon_socket::owner_work(|| {
96 codewhale_config::windows_identity::CurrentWindowsUser::open()?.sid_string()
97 })
98 .await?
99 );
100 drop(owner_client);
101 let acp = codewhale_app_server::daemon_client::connect_acp_if_published(
102 Some(config_path.clone()),
103 Some(socket_path.clone()),
104 )
105 .await?
106 .context("actual Normal owner admits its captured ACP frontend")?;
107 assert_eq!(acp.receipt(), &captured);
108 tokio::time::timeout(
109 Duration::from_secs(5),
110 acp.forward(tokio::io::empty(), tokio::io::sink()),
111 )
112 .await??;
113 let manager = threads.clone();
114 let (after_acp_eof, after_generation) =
115 codewhale_app_server::daemon_socket::owner_work(move || manager.capture_control_owner())
116 .await?;
117 assert_eq!(after_acp_eof.data_dir, captured.data_dir);
118 assert_eq!(after_generation, captured.lease_generation);
119 assert!(!threads.is_acp_host());
120 let unauthenticated = crate::tls::reqwest_client()
121 .get(format!("http://{addr}/v1/threads"))
122 .send()
123 .await?;
124 assert_eq!(unauthenticated.status(), StatusCode::UNAUTHORIZED);
125 // Compatibility prompts address an admitted durable thread, never an
126 // invented alias that could remint an existing conversation.
127 let mut creator = codewhale_app_server::daemon_client::connect(
128 Some(config_path.clone()),
129 Some(socket_path.clone()),
130 )
131 .await?;
132 assert_eq!(creator.receipt(), &captured);
133 creator
134 .send(
135 json!("create"),
136 "thread/create",
137 json!({"metadata":{"operation_key":"same-owner-conversation-create","model":"deepseek-v4-pro"}}),
138 )
139 .await?;
140 let created = reply(&mut creator, "create").await?;
141 assert!(created["error"].is_null(), "{created:#}");
142 let conversation = created["result"]["thread_id"]
143 .as_str()
144 .context("durable creation returns its canonical identity")?
145 .to_owned();
146 drop(creator);
147 for (id, answer) in [
148 ("first", "first owner answer"),
149 ("second", "second owner answer"),
150 ] {
151 model.push_turn(canned::simple_text_turn(answer));
152 let mut guest = codewhale_app_server::daemon_client::connect(
153 Some(config_path.clone()),
154 Some(socket_path.clone()),
155 )
156 .await?;
157 assert_eq!(guest.receipt(), &captured);
158 guest
159 .send(
160 json!(id),
161 "prompt/request",
162 json!({"thread_id":conversation,"prompt":id,"model":"deepseek-v4-pro"}),
163 )
164 .await?;
165 let result = reply(&mut guest, id).await?;
166 assert!(result["error"].is_null(), "{result:#}");
167 assert_eq!(result["result"]["output"], answer);
168 guest
169 .send(json!("guest-shutdown"), "shutdown", json!({}))
170 .await?;
171 assert!(reply(&mut guest, "guest-shutdown").await?["error"].is_object());
172 // Actual per-connection logical EOF on Windows and Unix; host stays up.
173 tokio::time::timeout(
174 Duration::from_secs(5),
175 guest.forward(tokio::io::empty(), tokio::io::sink()),
176 )
177 .await??;
178 }
179 let catalog: Value = crate::tls::reqwest_client()
180 .get(format!("http://{addr}/v1/threads"))
181 .bearer_auth(&token)
182 .send()
183 .await?
184 .error_for_status()?
185 .json()
186 .await?;
187 let rows = catalog.as_array().context("canonical thread catalog")?;
188 assert_eq!(rows.len(), 1, "reattachment must not remint Runtime/store");
189 assert_eq!(rows[0]["id"], conversation);
190 let detail = threads
191 .get_thread_detail(rows[0]["id"].as_str().context("thread id")?)
192 .await?;
193 assert_eq!(detail.turns.len(), 2);
194 assert!(
195 detail
196 .turns
197 .iter()
198 .all(|turn| serde_json::to_value(turn).is_ok_and(|v| v["status"] == "completed"))
199 );
200 assert!(detail.latest_seq > 0);
201 let manager = threads.clone();
202 let (after, generation) =
203 codewhale_app_server::daemon_socket::owner_work(move || manager.capture_control_owner())
204 .await?;
205 assert_eq!(after.data_dir, captured.data_dir);
206 assert_eq!(generation, captured.lease_generation);
207 // Join the actual owned Engines while their hosting runtime is still alive.
208 threads.shutdown_and_wait().await?;
209 assert!(shutdown.drain(Duration::from_secs(5)).await);
210 tokio::time::timeout(Duration::from_secs(5), async {
211 loop {
212 if codewhale_app_server::daemon_client::connect(
213 Some(config_path.clone()),
214 Some(socket_path.clone()),
215 )
216 .await
217 .is_err()
218 {
219 break;
220 }
221 tokio::time::sleep(Duration::from_millis(10)).await;
222 }
223 })
224 .await
225 .context("withdrawal must refuse without child/fresh-store fallback")?;
226 threads.shutdown_and_wait().await?;
227 http.abort();
228 Ok(())
229 }
230
231 struct GatedOwnerModel {
232 scripted: MockLlmClient,
233 first: std::sync::atomic::AtomicBool,
234 entered: tokio::sync::Notify,
235 release: tokio::sync::Semaphore,
236 }
237 impl crate::llm_client::LlmClient for GatedOwnerModel {
238 fn provider_name(&self) -> &'static str {
239 "deepseek"
240 }
241 fn model(&self) -> &str {
242 "deepseek-v4-pro"
243 }
244 async fn create_message(
245 &self,
246 request: codewhale_models::MessageRequest,
247 ) -> Result<codewhale_models::MessageResponse> {
248 crate::llm_client::LlmClient::create_message(&self.scripted, request).await
249 }
250 async fn create_message_stream(
251 &self,
252 request: codewhale_models::MessageRequest,
253 ) -> Result<crate::llm_client::StreamEventBox> {
254 if self.first.swap(false, std::sync::atomic::Ordering::SeqCst) {
255 self.entered.notify_one();
256 self.release.acquire().await?.forget();
257 }
258 crate::llm_client::LlmClient::create_message_stream(&self.scripted, request).await
259 }
260 }
261 fn listener_selection(
262 workspace: PathBuf,
263 port: u16,
264 token: String,
265 ) -> codewhale_app_server::RuntimeListenerSelection {
266 codewhale_app_server::RuntimeListenerSelection {
267 workers: 1,
268 workspace,
269 config_profile: None,
270 config_source: None,
271 host: "127.0.0.1".into(),
272 port,
273 cors_origins: Vec::new(),
274 auth_token: Some(token),
275 insecure_no_auth: false,
276 mobile: false,
277 web: false,
278 }
279 }
280 fn authenticated_http(token: &str) -> Result<reqwest::Client> {
281 let mut headers = HeaderMap::new();
282 headers.insert(
283 header::AUTHORIZATION,
284 HeaderValue::from_str(&format!("Bearer {token}"))?,
285 );
286 Ok(codewhale_release::platform_http_client_builder()
287 .default_headers(headers)
288 .redirect(reqwest::redirect::Policy::none())
289 .build()?)
290 }
291 async fn frontend_ready(
292 guest: &mut codewhale_app_server::daemon_client::OwnerClient,
293 ) -> Result<RuntimeFrontendReady> {
294 let frame = tokio::time::timeout(Duration::from_secs(10), guest.recv())
295 .await??
296 .context("owner closed before frontend readiness")?;
297 anyhow::ensure!(
298 frame["method"] == "daemon/frontend_ready" && frame.get("id").is_none(),
299 "invalid frontend readiness"
300 );
301 Ok(serde_json::from_value(frame["params"].clone())?)
302 }
303
304 #[tokio::test]
305 async fn real_normal_owner_distinct_listener_detach_keeps_turn_compat_acp_and_successor()
306 -> Result<()> {
307 let _env = lock_test_env();
308 let temporary = tempfile::tempdir()?;
309 let root = temporary.path().canonicalize()?;
310 let _home = EnvVarGuard::set("CODEWHALE_HOME", root.join("home"));
311 let workspace = root.join("workspace");
312 #[cfg(unix)]
313 let socket = root.join("private-run").join("normal-owner.sock");
314 #[cfg(windows)]
315 let socket = PathBuf::from(format!(r"\\.\pipe\normal-owner-{}", Uuid::new_v4()));
316 let original_token = "cwrt_original_owner_only";
317 let selected_token = "cwrt_selected_listener_only";
318 let config = Config::default().with_legacy_root(
319 Some("fixture-key".into()),
320 Some("http://127.0.0.1:1/v1".into()),
321 );
322 let (original, manager, server) =
323 spawn_test_server_with_root_token_mobile_workspace_and_overrides(
324 root.clone(),
325 root.join("sessions"),
326 Some(original_token.into()),
327 false,
328 workspace.clone(),
329 TestServerOverrides {
330 config: Some(config),
331 owner_socket: Some(socket.clone()),
332 ..Default::default()
333 },
334 )
335 .await?
336 .context("real listener required")?;
337 let _abort = AbortOnDrop(server.abort_handle());
338 let model = Arc::new(GatedOwnerModel {
339 scripted: MockLlmClient::new(vec![
340 canned::simple_text_turn("pending turn survived detach"),
341 canned::simple_text_turn("compat answer"),
342 canned::simple_text_turn("ACP answer"),
343 canned::simple_text_turn("ordinary successor"),
344 ]),
345 first: std::sync::atomic::AtomicBool::new(true),
346 entered: tokio::sync::Notify::new(),
347 release: tokio::sync::Semaphore::new(0),
348 });
349 manager.set_test_model_client(model.clone());
350 let control = codewhale_app_server::daemon_client::connect(None, Some(socket.clone())).await?;
351 let receipt = control.receipt().clone();
352 validate_selected_owner(&control, receipt.data_dir.clone()).await?;
353 drop(control);
354 let pending_workspace = root.join("pending-workspace");
355 fs::create_dir(&pending_workspace)?;
356 let mut guest = codewhale_app_server::daemon_client::connect_listener_if_published(
357 None,
358 Some(socket.clone()),
359 listener_selection(pending_workspace.clone(), 0, selected_token.into()),
360 receipt.clone(),
361 )
362 .await?
363 .context("published owner")?;
364 let ready = frontend_ready(&mut guest).await?;
365 assert_ne!(ready.endpoint, original);
366 assert!(!ready.reused_owner_listener);
367 let selected = authenticated_http(selected_token)?;
368 let original_http = authenticated_http(original_token)?;
369 assert_eq!(
370 selected
371 .get(format!("http://{original}/v1/threads"))
372 .send()
373 .await?
374 .status(),
375 StatusCode::UNAUTHORIZED
376 );
377 assert_eq!(
378 original_http
379 .get(format!("http://{}/v1/threads", ready.endpoint))
380 .send()
381 .await?
382 .status(),
383 StatusCode::UNAUTHORIZED
384 );
385 let thread: Value = selected
386 .post(format!("http://{}/v1/threads", ready.endpoint))
387 .json(&json!({}))
388 .send()
389 .await?
390 .error_for_status()?
391 .json()
392 .await?;
393 assert_eq!(thread["workspace"].as_str(), pending_workspace.to_str());
394 let id = thread["id"].as_str().context("thread id")?.to_owned();
395 let started: Value = selected
396 .post(format!("http://{}/v1/threads/{id}/turns", ready.endpoint))
397 .json(&json!({"prompt":"persist across listener close"}))
398 .send()
399 .await?
400 .error_for_status()?
401 .json()
402 .await?;
403 let turn = started["turn"]["id"]
404 .as_str()
405 .context("turn id")?
406 .to_owned();
407 tokio::time::timeout(Duration::from_secs(30), model.entered.notified()).await?;
408 drop(guest);
409 tokio::time::timeout(Duration::from_secs(10), async {
410 while tokio::net::TcpStream::connect(ready.endpoint).await.is_ok() {
411 tokio::time::sleep(Duration::from_millis(20)).await;
412 }
413 })
414 .await?;
415 assert_eq!(
416 original_http
417 .get(format!("http://{original}/healthz"))
418 .send()
419 .await?
420 .status(),
421 StatusCode::OK
422 );
423 assert!(!manager.is_acp_host());
424 model.release.add_permits(1);
425 assert_eq!(
426 wait_for_terminal_turn_status(
427 &original_http,
428 original,
429 &id,
430 &turn,
431 Duration::from_secs(30)
432 )
433 .await?,
434 "completed"
435 );
436 // Compatibility uses the same captured original bridge and actual Engine.
437 let selected_workspace = root.join("compat-workspace");
438 fs::create_dir(&selected_workspace)?;
439 let mut compat = codewhale_app_server::daemon_client::connect_scoped_control_if_published(
440 None,
441 Some(socket.clone()),
442 codewhale_app_server::RuntimeFrontendScope {
443 workers: 1,
444 workspace: selected_workspace.clone(),
445 config_profile: None,
446 config_source: None,
447 },
448 receipt.clone(),
449 )
450 .await?
451 .context("captured scoped control")?;
452 compat
453 .send(
454 json!("compat-create"),
455 "thread/create",
456 json!({"metadata":{"operation_key":"scoped-compat-create"}}),
457 )
458 .await?;
459 let created = reply(&mut compat, "compat-create").await?;
460 assert!(created["error"].is_null(), "{created:#}");
461 let compat_id = created["result"]["thread_id"]
462 .as_str()
463 .context("scoped creation returns its canonical identity")?
464 .to_owned();
465 compat
466 .send(
467 json!("compat"),
468 "prompt/request",
469 json!({"thread_id":compat_id,"prompt":"same owner compatibility"}),
470 )
471 .await?;
472 let result = reply(&mut compat, "compat").await?;
473 assert!(result["error"].is_null(), "{result:#}");
474 assert_eq!(result["result"]["output"], "compat answer");
475 drop(compat);
476 let mut acp = codewhale_app_server::daemon_client::connect_selected_acp_if_published(
477 None,
478 Some(socket.clone()),
479 codewhale_app_server::RuntimeFrontendScope {
480 workers: 1,
481 workspace: workspace.clone(),
482 config_profile: None,
483 config_source: None,
484 },
485 "deepseek-v4-pro".into(),
486 receipt.clone(),
487 )
488 .await?
489 .context("same Normal owner ACP")?;
490 acp.send(
491 json!("init"),
492 "initialize",
493 json!({"protocolVersion":1,"clientCapabilities":{}}),
494 )
495 .await?;
496 assert!(reply(&mut acp, "init").await?["error"].is_null());
497 acp.send(
498 json!("new"),
499 "session/new",
500 json!({"cwd":workspace,"mcpServers":[]}),
501 )
502 .await?;
503 let session = reply(&mut acp, "new").await?["result"]["sessionId"]
504 .as_str()
505 .context("ACP session")?
506 .to_owned();
507 acp.send(
508 json!("prompt"),
509 "session/prompt",
510 json!({"sessionId":session,"prompt":"one narrowed turn"}),
511 )
512 .await?;
513 assert_eq!(
514 reply(&mut acp, "prompt").await?["result"]["stopReason"],
515 "end_turn"
516 );
517 drop(acp);
518 let successor: Value = original_http
519 .post(format!("http://{original}/v1/threads/{id}/turns"))
520 .json(&json!({"prompt":"ordinary successor"}))
521 .send()
522 .await?
523 .error_for_status()?
524 .json()
525 .await?;
526 let successor_id = successor["turn"]["id"].as_str().context("successor id")?;
527 assert_eq!(
528 wait_for_terminal_turn_status(
529 &original_http,
530 original,
531 &id,
532 successor_id,
533 Duration::from_secs(30)
534 )
535 .await?,
536 "completed"
537 );
538 assert_eq!(
539 model.scripted.call_count(),
540 4,
541 "no replay after guest detach"
542 );
543 assert_eq!(
544 model.scripted.captured_requests()[2].model,
545 "deepseek-v4-pro"
546 );
547 let catalog: Value = original_http
548 .get(format!("http://{original}/v1/threads"))
549 .send()
550 .await?
551 .error_for_status()?
552 .json()
553 .await?;
554 assert!(
555 catalog
556 .as_array()
557 .context("canonical catalog")?
558 .iter()
559 .any(|row| row["workspace"].as_str() == selected_workspace.to_str()),
560 "scoped control default must reach the same canonical thread store"
561 );
562
563 let captured = manager.clone();
564 let (binding, generation) =
565 codewhale_app_server::daemon_socket::owner_work(move || captured.capture_control_owner())
566 .await?;
567 assert_eq!(generation, receipt.lease_generation);
568 assert_eq!(binding.data_dir, receipt.data_dir);
569 let detail = manager.get_thread_detail(&id).await?;
570 assert_eq!(detail.turns.len(), 2);
571 assert!(detail.latest_seq > 0);
572 manager.shutdown_and_wait().await?;
573 server.abort();
574 Ok(())
575 }
576
577 #[tokio::test]
578 async fn real_owner_exact_listener_reuse_scope_refusal_and_fresh_browser_projection() -> Result<()>
579 {
580 let _env = lock_test_env();
581 let temporary = tempfile::tempdir()?;
582 let root = temporary.path().canonicalize()?;
583 let _home = EnvVarGuard::set("CODEWHALE_HOME", root.join("home"));
584 let workspace = root.join("workspace");
585 #[cfg(unix)]
586 let socket = root.join("private-run").join("browser-owner.sock");
587 #[cfg(windows)]
588 let socket = PathBuf::from(format!(r"\\.\pipe\browser-owner-{}", Uuid::new_v4()));
589 let token = "cwrt_bound_operator_secret";
590 let (original, manager, server) =
591 spawn_test_server_with_root_token_mobile_workspace_and_overrides(
592 root.clone(),
593 root.join("sessions"),
594 Some(token.into()),
595 false,
596 workspace.clone(),
597 TestServerOverrides {
598 owner_socket: Some(socket.clone()),
599 ..Default::default()
600 },
601 )
602 .await?
603 .context("real listener required")?;
604 let _abort = AbortOnDrop(server.abort_handle());
605 let control = codewhale_app_server::daemon_client::connect(None, Some(socket.clone())).await?;
606 let receipt = control.receipt().clone();
607 drop(control);
608 let same = listener_selection(workspace.clone(), original.port(), token.into());
609 assert!(!format!("{same:?}").contains(token));
610 let mut reuse = codewhale_app_server::daemon_client::connect_listener_if_published(
611 None,
612 Some(socket.clone()),
613 same.clone(),
614 receipt.clone(),
615 )
616 .await?
617 .context("published owner")?;
618 let ready = frontend_ready(&mut reuse).await?;
619 assert_eq!(ready.endpoint, original);
620 assert!(ready.reused_owner_listener);
621 drop(reuse);
622 let client = authenticated_http(token)?;
623 assert_eq!(
624 client
625 .get(format!("http://{original}/v1/threads"))
626 .send()
627 .await?
628 .status(),
629 StatusCode::OK
630 );
631 let mut wrong = same;
632 wrong.workers = 2;
633 assert!(
634 codewhale_app_server::daemon_client::connect_listener_if_published(
635 None,
636 Some(socket.clone()),
637 wrong,
638 receipt.clone()
639 )
640 .await
641 .is_err()
642 );
643 assert!(
644 codewhale_app_server::daemon_client::connect_scoped_control_if_published(
645 None,
646 Some(socket.clone()),
647 codewhale_app_server::RuntimeFrontendScope {
648 workers: 1,
649 workspace: workspace.clone(),
650 config_profile: Some("unadmitted-profile".into()),
651 config_source: None,
652 },
653 receipt.clone()
654 )
655 .await
656 .is_err()
657 );
658 let mut stale = receipt.clone();
659 stale.lease_generation.push_str("-stale");
660 assert!(
661 codewhale_app_server::daemon_client::connect_listener_if_published(
662 None,
663 Some(socket.clone()),
664 listener_selection(workspace.clone(), 0, token.into()),
665 stale
666 )
667 .await
668 .is_err()
669 );
670 // A distinct selected root now joins its actual owner-held scope.
671 let other_workspace = root.join("other-workspace");
672 fs::create_dir(&other_workspace)?;
673 let reservation = std::net::TcpListener::bind("127.0.0.1:0")?;
674 let port = reservation.local_addr()?.port();
675 drop(reservation);
676 let mut browser = listener_selection(other_workspace.clone(), port, token.into());
677 browser.web = true;
678 browser.mobile = true;
679 let mut guest = codewhale_app_server::daemon_client::connect_listener_if_published(
680 None,
681 Some(socket.clone()),
682 browser,
683 receipt.clone(),
684 )
685 .await?
686 .context("published owner")?;
687 let selected = frontend_ready(&mut guest).await?;
688 assert!(!selected.reused_owner_listener);
689 let web_url = selected.web_bootstrap_url.context("fresh web bootstrap")?;
690 let mobile_url = selected
691 .mobile_bootstrap_url
692 .context("fresh mobile bootstrap")?;
693 assert!(!web_url.contains(token) && !mobile_url.contains(token));
694 let boot = client.get(&web_url).send().await?;
695 assert_eq!(boot.status(), StatusCode::SEE_OTHER);
696 let cookie = boot.headers()[header::SET_COOKIE]
697 .to_str()?
698 .split(';')
699 .next()
700 .context("web cookie")?
701 .to_owned();
702 let proof = boot.headers()[header::LOCATION]
703 .to_str()?
704 .strip_prefix("/#p=")
705 .context("web proof")?
706 .to_owned();
707 assert_eq!(
708 client.get(&web_url).send().await?.status(),
709 StatusCode::UNAUTHORIZED
710 );
711 // Cookies/proofs are admitted only at this selected origin and cannot
712 // replace the original listener's bearer gate.
713 let unauth = crate::tls::reqwest_client();
714 let selected_origin = format!("http://{}", selected.endpoint);
715 let created: Value = unauth
716 .post(format!("{selected_origin}/v1/threads"))
717 .header(header::COOKIE, &cookie)
718 .header(web::WEB_REQUEST_HEADER, &proof)
719 .header(header::ORIGIN, &selected_origin)
720 .json(&json!({}))
721 .send()
722 .await?
723 .error_for_status()?
724 .json()
725 .await?;
726 assert_eq!(created["workspace"].as_str(), other_workspace.to_str());
727 assert_eq!(
728 unauth
729 .get(format!("http://{original}/v1/threads"))
730 .header(header::COOKIE, &cookie)
731 .header(web::WEB_REQUEST_HEADER, &proof)
732 .send()
733 .await?
734 .status(),
735 StatusCode::UNAUTHORIZED
736 );
737 let mobile = client.get(&mobile_url).send().await?;
738 assert_eq!(mobile.status(), StatusCode::SEE_OTHER);
739 assert!(
740 !mobile.headers()[header::SET_COOKIE]
741 .to_str()?
742 .contains(token)
743 );
744 assert_eq!(
745 client.get(&mobile_url).send().await?.status(),
746 StatusCode::UNAUTHORIZED
747 );
748 drop(guest);
749 assert_eq!(
750 client
751 .get(format!("http://{original}/v1/threads"))
752 .send()
753 .await?
754 .status(),
755 StatusCode::OK
756 );
757 manager.shutdown_and_wait().await?;
758 server.abort();
759 Ok(())
760 }
761
762 #[tokio::test]
763 async fn real_owner_two_workspace_services_share_scope_caches_and_settle_global_mutations()
764 -> Result<()> {
765 let _env = lock_test_env();
766 let temporary = tempfile::tempdir()?;
767 let root = temporary.path().canonicalize()?;
768 let _home = EnvVarGuard::set("CODEWHALE_HOME", root.join("home"));
769 let original_workspace = root.join("original-workspace");
770 let selected_workspace = root.join("selected-workspace");
771 fs::create_dir_all(&original_workspace)?;
772 fs::create_dir_all(&selected_workspace)?;
773 fs::write(original_workspace.join("marker.txt"), "original contents")?;
774 fs::write(selected_workspace.join("marker.txt"), "selected contents")?;
775 #[cfg(unix)]
776 let socket = root.join("private-run").join("scope-owner.sock");
777 #[cfg(windows)]
778 let socket = PathBuf::from(format!(r"\\.\pipe\scope-owner-{}", Uuid::new_v4()));
779 let token = "cwrt_scope_operator_only";
780 let workers = runtime_api_sub_agent_manager(&original_workspace, 2);
781 let (scope_tx, scope_rx) = oneshot::channel();
782 let (original, manager, server) =
783 spawn_test_server_with_root_token_mobile_workspace_and_overrides(
784 root.clone(),
785 root.join("sessions"),
786 Some(token.into()),
787 false,
788 original_workspace.clone(),
789 TestServerOverrides {
790 owner_socket: Some(socket.clone()),
791 sub_agent_manager: Some(workers.clone()),
792 workspace_scopes_handle: Some(scope_tx),
793 ..Default::default()
794 },
795 )
796 .await?
797 .context("real original listener required")?;
798 let _abort = AbortOnDrop(server.abort_handle());
799 let scopes = scope_rx.await?;
800 assert!(
801 Arc::ptr_eq(&scopes.workers, &workers),
802 "workspace service admission must reuse the real owner's worker actor"
803 );
804
805 let control = codewhale_app_server::daemon_client::connect(None, Some(socket.clone())).await?;
806 let receipt = control.receipt().clone();
807 drop(control);
808 let mut first = codewhale_app_server::daemon_client::connect_listener_if_published(
809 None,
810 Some(socket.clone()),
811 listener_selection(selected_workspace.clone(), 0, token.into()),
812 receipt.clone(),
813 )
814 .await?
815 .context("first selected listener")?;
816 let first_ready = frontend_ready(&mut first).await?;
817 let mut second = codewhale_app_server::daemon_client::connect_listener_if_published(
818 None,
819 Some(socket.clone()),
820 listener_selection(selected_workspace.clone(), 0, token.into()),
821 receipt,
822 )
823 .await?
824 .context("matching selected listener")?;
825 let second_ready = frontend_ready(&mut second).await?;
826 assert_ne!(first_ready.endpoint, second_ready.endpoint);
827 let client = authenticated_http(token)?;
828 for (address, workspace, contents) in [
829 (original, &original_workspace, "original contents"),
830 (
831 first_ready.endpoint,
832 &selected_workspace,
833 "selected contents",
834 ),
835 (
836 second_ready.endpoint,
837 &selected_workspace,
838 "selected contents",
839 ),
840 ] {
841 let file: Value = client
842 .get(format!(
843 "http://{address}/v1/workspace/files/read?path=marker.txt"
844 ))
845 .send()
846 .await?
847 .error_for_status()?
848 .json()
849 .await?;
850 assert_eq!(file["content"], contents);
851 let lsp: Value = client
852 .get(format!("http://{address}/v1/lsp"))
853 .send()
854 .await?
855 .error_for_status()?
856 .json()
857 .await?;
858 assert_eq!(lsp["workspace"].as_str(), workspace.to_str());
859 }
860 let (original_scope, selected_scope) = {
861 let held = scopes.scopes.lock();
862 assert_eq!(
863 held.len(),
864 2,
865 "two matching listeners must not create two scopes"
866 );
867 (
868 held[&original_workspace].clone(),
869 held[&selected_workspace].clone(),
870 )
871 };
872 for endpoint in [original, second_ready.endpoint] {
873 let runs: Value = client
874 .get(format!("http://{endpoint}/v1/agent-runs"))
875 .send()
876 .await?
877 .error_for_status()?
878 .json()
879 .await?;
880 assert_eq!(
881 runs["governor"]["max_launch_slots"], 2,
882 "both selected roots must report the same owner-global worker ceiling"
883 );
884 }
885 // Run a real Engine turn in the selected workspace, then retain a live
886 // direct child under that exact owner/session while its execution cwd is
887 // an isolated directory. Endpoint scope must come from admission receipt.
888 let model = Arc::new(MockLlmClient::new(vec![canned::simple_text_turn(
889 "origin-owner-turn",
890 )]));
891 manager.set_test_model_client(model.clone());
892 let selected_thread: Value = client
893 .post(format!("http://{}/v1/threads", second_ready.endpoint))
894 .json(&json!({"model":"deepseek-v4-pro"}))
895 .send()
896 .await?
897 .error_for_status()?
898 .json()
899 .await?;
900 let selected_session = selected_thread["id"]
901 .as_str()
902 .context("selected owner session")?
903 .to_owned();
904 let turn: Value = client
905 .post(format!(
906 "http://{}/v1/threads/{selected_session}/turns",
907 second_ready.endpoint
908 ))
909 .json(&json!({"prompt":"prove selected owner Engine"}))
910 .send()
911 .await?
912 .error_for_status()?
913 .json()
914 .await?;
915 assert_eq!(
916 wait_for_terminal_turn_status(
917 &client,
918 second_ready.endpoint,
919 &selected_session,
920 turn["turn"]["id"].as_str().context("selected turn")?,
921 Duration::from_secs(30)
922 )
923 .await?,
924 "completed"
925 );
926 assert_eq!(model.call_count(), 1);
927 let selected_fleet = FleetManager::open(&selected_workspace)?
928 .with_sub_agent_manager(workers.clone())
929 .with_session_model(DEFAULT_TEXT_MODEL)
930 .with_route_config(test_fleet_route_config());
931 let task: codewhale_protocol::fleet::FleetTaskSpec = serde_json::from_value(json!({
932 "id":"scope-lease", "name":"scope lease", "description":null,
933 "objective":"hold a reviewed local Fleet admission", "instructions":"no process launched",
934 "worker":{"agent_profile":null,"role":"reviewer","loadout":null,"model_class":null,
935 "model":null,"tool_profile":"read-only","tools":[],"capabilities":[]},
936 "workspace":null,"input_files":[],"context":[],"budget":null,"tags":[],
937 "expected_artifacts":[],"scorer":null,"retry_policy":null,"alert_policy":null,
938 "timeout_seconds":null,"metadata":{}
939 }))?;
940 let report = selected_fleet.create_run(
941 crate::fleet::task_spec::FleetTaskSpecDocument {
942 name: Some("selected lease".into()),
943 labels: Default::default(),
944 security_policy: None,
945 workers: Vec::new(),
946 tasks: vec![task],
947 usage_ceiling: None,
948 },
949 1,
950 )?;
951 let ledger = selected_fleet.rebuild_state()?;
952 let fleet_record = workers
953 .read()
954 .await
955 .fleet_worker_records_for_workspace(&selected_workspace)
956 .map_err(anyhow::Error::msg)?
957 .into_iter()
958 .find(|record| record.spec.run_id == report.run_id.0)
959 .context("actual Fleet admission record")?;
960 assert!(fleet_worker_has_selected_lease(&fleet_record, &ledger));
961 let mut wrong_worker = ledger.clone();
962 wrong_worker
963 .tasks
964 .values_mut()
965 .next()
966 .context("actual Fleet task")?
967 .leased_to = Some("different-live-worker".into());
968 assert!(!fleet_worker_has_selected_lease(
969 &fleet_record,
970 &wrong_worker
971 ));
972 let mut wrong_run = ledger;
973 wrong_run
974 .tasks
975 .values_mut()
976 .next()
977 .context("actual Fleet task")?
978 .entry
979 .run_id
980 .0 = "different-run".into();
981 assert!(!fleet_worker_has_selected_lease(&fleet_record, &wrong_run));
982 assert_eq!(
983 client
984 .get(format!(
985 "http://{original}/v1/agent-runs/{}",
986 fleet_record.spec.run_id
987 ))
988 .send()
989 .await?
990 .status(),
991 StatusCode::NOT_FOUND
992 );
993 assert_eq!(
994 client
995 .get(format!(
996 "http://{}/v1/agent-runs/{}",
997 second_ready.endpoint, fleet_record.spec.run_id
998 ))
999 .send()
1000 .await?
1001 .status(),
1002 StatusCode::OK
1003 );
1004 let execution_workspace = root.join("isolated-direct-execution");
1005 fs::create_dir(&execution_workspace)?;
1006 let mut actor = workers.clone().write_owned().await;
1007 let selected_root = selected_workspace.clone();
1008 let direct_id = codewhale_app_server::daemon_socket::owner_work(move || {
1009 actor
1010 .insert_test_running_direct_child_in_origin(
1011 "scope-direct",
1012 &execution_workspace,
1013 &selected_root,
1014 &selected_session,
1015 )
1016 .map_err(anyhow::Error::msg)
1017 })
1018 .await?;
1019 for address in [first_ready.endpoint, second_ready.endpoint] {
1020 let record: Value = client
1021 .get(format!("http://{address}/v1/agent-runs/{direct_id}"))
1022 .send()
1023 .await?
1024 .error_for_status()?
1025 .json()
1026 .await?;
1027 assert_ne!(
1028 record["spec"]["workspace"].as_str(),
1029 selected_workspace.to_str(),
1030 "execution cwd must not become originating scope"
1031 );
1032 }
1033 assert_eq!(
1034 client
1035 .post(format!(
1036 "http://{original}/v1/agent-runs/{direct_id}/cancel"
1037 ))
1038 .send()
1039 .await?
1040 .status(),
1041 StatusCode::NOT_FOUND
1042 );
1043 assert_eq!(
1044 workers.read().await.get_result(&direct_id)?.status,
1045 SubAgentStatus::Running
1046 );
1047 let stopped: Value = client
1048 .post(format!(
1049 "http://{}/v1/agent-runs/{direct_id}/cancel",
1050 second_ready.endpoint
1051 ))
1052 .send()
1053 .await?
1054 .error_for_status()?
1055 .json()
1056 .await?;
1057 let repeated: Value = client
1058 .post(format!(
1059 "http://{}/v1/agent-runs/{direct_id}/cancel",
1060 first_ready.endpoint
1061 ))
1062 .send()
1063 .await?
1064 .error_for_status()?
1065 .json()
1066 .await?;
1067 assert_eq!(
1068 stopped["status"], repeated["status"],
1069 "terminal repeated cancellation stays idempotent"
1070 );
1071 assert_eq!(stopped["spec"]["worker_id"], direct_id);
1072 assert_eq!(
1073 workers
1074 .read()
1075 .await
1076 .rate_limit_governor()
1077 .snapshot(std::time::Instant::now())
1078 .max_capacity,
1079 2
1080 );
1081
1082 let (admitted_a, admitted_b) = tokio::try_join!(
1083 scopes.admit(selected_workspace.clone()),
1084 scopes.admit(selected_workspace.clone())
1085 )?;
1086 assert!(Arc::ptr_eq(&admitted_a, &admitted_b));
1087 assert!(Arc::ptr_eq(&admitted_a, &selected_scope));
1088 assert!(!Arc::ptr_eq(
1089 original_scope.lsp.get().unwrap(),
1090 selected_scope.lsp.get().unwrap()
1091 ));
1092 let a = client
1093 .get(format!(
1094 "http://{}/v1/apps/mcp/tools?connect=true",
1095 first_ready.endpoint
1096 ))
1097 .send();
1098 let b = client
1099 .get(format!(
1100 "http://{}/v1/apps/mcp/tools?connect=true",
1101 second_ready.endpoint
1102 ))
1103 .send();
1104 let (a, b) = tokio::try_join!(a, b)?;
1105 a.error_for_status()?;
1106 b.error_for_status()?;
1107 client
1108 .get(format!("http://{original}/v1/apps/mcp/tools?connect=true"))
1109 .send()
1110 .await?
1111 .error_for_status()?;
1112 let original_pool = original_scope.mcp.lock().await.as_ref().unwrap().1.clone();
1113 let selected_pool = selected_scope.mcp.lock().await.as_ref().unwrap().1.clone();
1114 assert!(!Arc::ptr_eq(&original_pool, &selected_pool));
1115 assert!(Arc::ptr_eq(
1116 &original_pool.lock().await.dynamic_servers,
1117 &selected_pool.lock().await.dynamic_servers
1118 ));
1119 // A pending pool operation keeps its exact live handle while the single
1120 // persisted global mutation settles. No pool replacement or replay occurs.
1121 let pending = selected_pool.lock().await;
1122 let base = format!("http://{original}/v1/apps/mcp/servers");
1123 mcp_test_success(
1124 client
1125 .post(&base)
1126 .header(header::IF_MATCH, mcp_test_revision(&client, &base).await?)
1127 .json(&json!({"name":"scope-proof","command":"scope-proof-never-executed"}))
1128 .send()
1129 .await?,
1130 )
1131 .await?;
1132 assert!(!pending.server_names().contains(&"scope-proof".to_string()));
1133 drop(pending);
1134 for (method, suffix, body) in [
1135 (reqwest::Method::GET, "", None),
1136 (
1137 reqwest::Method::PATCH,
1138 "/scope-proof",
1139 Some(json!({"args":["changed"]})),
1140 ),
1141 (reqwest::Method::POST, "/scope-proof/disable", None),
1142 (reqwest::Method::POST, "/scope-proof/enable", None),
1143 (reqwest::Method::DELETE, "/scope-proof", None),
1144 ] {
1145 if method != reqwest::Method::GET {
1146 let mut request = client
1147 .request(method, format!("{base}{suffix}"))
1148 .header(header::IF_MATCH, mcp_test_revision(&client, &base).await?);
1149 if let Some(body) = body {
1150 request = request.json(&body);
1151 }
1152 mcp_test_success(request.send().await?).await?;
1153 }
1154 for address in [original, first_ready.endpoint, second_ready.endpoint] {
1155 mcp_test_success(
1156 client
1157 .get(format!("http://{address}/v1/apps/mcp/servers"))
1158 .send()
1159 .await?,
1160 )
1161 .await?;
1162 }
1163 let generation = scopes
1164 .mcp_generation
1165 .load(std::sync::atomic::Ordering::SeqCst);
1166 for (scope, pool) in [
1167 (&original_scope, &original_pool),
1168 (&selected_scope, &selected_pool),
1169 ] {
1170 let cache = scope.mcp.lock().await;
1171 let (actual, retained) = cache.as_ref().unwrap();
1172 assert_eq!(
1173 *actual, generation,
1174 "global mutation must invalidate both workspace catalogs"
1175 );
1176 assert!(
1177 Arc::ptr_eq(retained, pool),
1178 "accepted pool identity must survive reload"
1179 );
1180 }
1181 }
1182 assert!(
1183 !selected_pool
1184 .lock()
1185 .await
1186 .server_names()
1187 .contains(&"scope-proof".to_string())
1188 );
1189 drop(first);
1190 client
1191 .get(format!("http://{}/v1/lsp", second_ready.endpoint))
1192 .send()
1193 .await?
1194 .error_for_status()?;
1195 assert!(Arc::ptr_eq(
1196 &scopes.admit(selected_workspace.clone()).await?,
1197 &selected_scope
1198 ));
1199 let retired_workspace = root.join("retired-workspace");
1200 #[cfg(windows)]
1201 {
1202 // This fixture still owns a Fleet ledger whose ancestor pins deny
1203 // deletion. Check that protection before releasing the test owner.
1204 let error = fs::rename(&selected_workspace, &retired_workspace)
1205 .expect_err("a live Fleet ledger must prevent workspace replacement");
1206 assert_eq!(error.raw_os_error(), Some(32), "{error}");
1207 assert!(!retired_workspace.exists());
1208 assert!(Arc::ptr_eq(
1209 &scopes.admit(selected_workspace.clone()).await?,
1210 &selected_scope
1211 ));
1212 }
1213 drop(selected_fleet);
1214 // File identity, not a stable pathname, binds the cache. Replacing the
1215 // selected directory after the Fleet pins close cannot borrow its admitted
1216 // LSP/MCP authority. Keep this check on Windows as well as Unix.
1217 fs::rename(&selected_workspace, &retired_workspace)
1218 .context("replace selected workspace after releasing the test Fleet manager")?;
1219 fs::create_dir(&selected_workspace)?;
1220 assert!(scopes.admit(selected_workspace.clone()).await.is_err());
1221 assert_eq!(
1222 client
1223 .get(format!("http://{}/v1/lsp", second_ready.endpoint))
1224 .send()
1225 .await?
1226 .status(),
1227 StatusCode::CONFLICT
1228 );
1229 client
1230 .get(format!("http://{original}/v1/lsp"))
1231 .send()
1232 .await?
1233 .error_for_status()?;
1234 // Retained scope admission is finite even after frontend detach. The
1235 // current owner holds caches/cleanup instead of creating an idle leak.
1236 for n in 2..MAX_RUNTIME_WORKSPACE_SCOPES {
1237 let path = root.join(format!("scope-{n}"));
1238 fs::create_dir(&path)?;
1239 scopes.admit(path).await?;
1240 }
1241 let excess = root.join("excess-scope");
1242 fs::create_dir(&excess)?;
1243 assert!(scopes.admit(excess).await.is_err());
1244 assert_eq!(scopes.scopes.lock().len(), MAX_RUNTIME_WORKSPACE_SCOPES);
1245 drop(second);
1246 manager.shutdown_and_wait().await?;
1247 drop((
1248 admitted_a,
1249 admitted_b,
1250 original_scope,
1251 selected_scope,
1252 original_pool,
1253 selected_pool,
1254 scopes,
1255 ));
1256 server.abort();
1257 Ok(())
1258 }
1259
1260 #[tokio::test]
1261 async fn real_owner_rejects_different_operator_source_before_attach_and_accepts_exact_alias()
1262 -> Result<()> {
1263 let _env = lock_test_env();
1264 let temporary = tempfile::tempdir()?;
1265 let root = temporary.path().canonicalize()?;
1266 let _home = EnvVarGuard::set("CODEWHALE_HOME", root.join("home"));
1267 let workspace = root.join("workspace");
1268 fs::create_dir_all(&workspace)?;
1269 let config_path = root.join("operator.toml");
1270 fs::write(&config_path, "")?;
1271 let config = Config::load(Some(config_path.clone()), None)?;
1272 let other_config = root.join("other-operator.toml");
1273 fs::write(&other_config, "")?;
1274 let alias_directory = root.join("alias");
1275 fs::create_dir(&alias_directory)?;
1276 let alias = alias_directory.join("..").join("operator.toml");
1277 #[cfg(unix)]
1278 let socket = root.join("private-run").join("config-owner.sock");
1279 #[cfg(windows)]
1280 let socket = PathBuf::from(format!(r"\\.\pipe\config-owner-{}", Uuid::new_v4()));
1281 let token = "cwrt_config_source_private";
1282 let (original, manager, server) =
1283 spawn_test_server_with_root_token_mobile_workspace_and_overrides(
1284 root.clone(),
1285 root.join("sessions"),
1286 Some(token.into()),
1287 false,
1288 workspace.clone(),
1289 TestServerOverrides {
1290 config: Some(config),
1291 config_path: Some(config_path.clone()),
1292 owner_socket: Some(socket.clone()),
1293 ..Default::default()
1294 },
1295 )
1296 .await?
1297 .context("actual config-bound owner required")?;
1298 let _abort = AbortOnDrop(server.abort_handle());
1299 let control = codewhale_app_server::daemon_client::connect(
1300 Some(config_path.clone()),
1301 Some(socket.clone()),
1302 )
1303 .await?;
1304 let receipt = control.receipt().clone();
1305 drop(control);
1306 let mut wrong = listener_selection(workspace.clone(), 0, token.into());
1307 wrong.config_source = Some(other_config.clone());
1308 let error = codewhale_app_server::daemon_client::connect_listener_if_published(
1309 Some(config_path.clone()),
1310 Some(socket.clone()),
1311 wrong,
1312 receipt.clone(),
1313 )
1314 .await
1315 .err()
1316 .context("explicit other operator source must refuse before attach ACK")?;
1317 let message = error.to_string();
1318 assert!(
1319 message == "selected owner refused guest attachment",
1320 "{message}"
1321 );
1322 assert!(!message.contains(token));
1323 assert!(!message.contains(other_config.to_str().unwrap()));
1324 let mut matching = listener_selection(workspace, 0, token.into());
1325 matching.config_source = Some(alias);
1326 let mut selected = codewhale_app_server::daemon_client::connect_listener_if_published(
1327 Some(config_path),
1328 Some(socket),
1329 matching,
1330 receipt,
1331 )
1332 .await?
1333 .context("same canonical captured operator source must be admitted")?;
1334 let ready = frontend_ready(&mut selected).await?;
1335 let client = authenticated_http(token)?;
1336 client
1337 .get(format!("http://{}/v1/threads", ready.endpoint))
1338 .send()
1339 .await?
1340 .error_for_status()?;
1341 drop(selected);
1342 client
1343 .get(format!("http://{original}/v1/threads"))
1344 .send()
1345 .await?
1346 .error_for_status()?;
1347 manager.shutdown_and_wait().await?;
1348 server.abort();
1349 Ok(())
1350 }
1351
1352 #[test]
1353 fn captured_operator_source_preserves_uncreated_nested_selection_without_reminting() -> Result<()> {
1354 let temporary = tempfile::tempdir()?;
1355 let root = temporary.path().canonicalize()?;
1356 let selected = root
1357 .join("not-created")
1358 .join("nested")
1359 .join("operator.toml");
1360 assert_eq!(
1361 canonical_runtime_config_source(Some(selected.clone()))?,
1362 Some(selected.clone())
1363 );
1364 assert!(
1365 !root.join("not-created").exists(),
1366 "selection admission must not create a config or its parents"
1367 );
1368 assert_ne!(
1369 canonical_runtime_config_source(Some(root.join("other.toml")))?,
1370 Some(selected)
1371 );
1372 Ok(())
1373 }
1374
1375 #[tokio::test]
1376 async fn canonical_history_import_runs_real_engine_retains_all_branches_and_refuses_lost_commit()
1377 -> Result<()> {
1378 use codewhale_protocol::{
1379 CanonicalHistoryImportRequest, CanonicalThreadReceipt, LegacyThreadHistory, MessageRecord,
1380 };
1381 let _env = lock_test_env();
1382 let temporary = tempfile::tempdir()?;
1383 let root = temporary.path().canonicalize()?;
1384 let _home = EnvVarGuard::set("CODEWHALE_HOME", root.join("home"));
1385 let workspace = root.join("workspace");
1386 fs::create_dir_all(&workspace)?;
1387 let sessions_dir = root.join("sessions");
1388 let token = format!("history-fixture-{}", Uuid::new_v4());
1389 let config = Config::default().with_legacy_root(
1390 Some("history-fixture-key".into()),
1391 Some("http://127.0.0.1:1/v1".into()),
1392 );
1393 let (addr, runtime, server) = spawn_test_server_with_root_token_mobile_workspace_and_overrides(
1394 root.clone(),
1395 sessions_dir.clone(),
1396 Some(token.clone()),
1397 false,
1398 workspace.clone(),
1399 TestServerOverrides {
1400 config: Some(config),
1401 ..TestServerOverrides::default()
1402 },
1403 )
1404 .await?
1405 .context("canonical history acceptance requires an actual listener")?;
1406 let _server = AbortOnDrop(server.abort_handle());
1407 let model = Arc::new(MockLlmClient::new(vec![canned::simple_text_turn(
1408 "canonical successor",
1409 )]));
1410 runtime.set_test_model_client(model.clone());
1411 let binding = runtime.session_store_binding();
1412 let request = CanonicalHistoryImportRequest {
1413 version: 1,
1414 operation_key: "canonical-history-acceptance".into(),
1415 expected_data_dir: binding.data_dir.clone(),
1416 expected_execution_scope: binding.execution_scope.clone(),
1417 target_runtime_thread_id: None,
1418 workspace,
1419 model: Some("deepseek-v4-pro".into()),
1420 history: LegacyThreadHistory {
1421 goal: None,
1422 version: 1,
1423 state_store_id: "c".repeat(64),
1424 thread_id: "legacy-full-graph".into(),
1425 current_leaf_id: Some(3),
1426 messages: [
1427 (1, None, "system", "retained system context"),
1428 (2, Some(1), "user", "retained user context"),
1429 (3, Some(2), "assistant", "retained prior answer"),
1430 (4, Some(1), "user", "retained inactive branch"),
1431 ]
1432 .into_iter()
1433 .map(|(id, parent_entry_id, role, content)| MessageRecord {
1434 id,
1435 thread_id: "legacy-full-graph".into(),
1436 role: role.into(),
1437 content: content.into(),
1438 item: None,
1439 created_at: 1_700_000_000 + id,
1440 parent_entry_id,
1441 })
1442 .collect(),
1443 },
1444 };
1445 let mut headers = reqwest::header::HeaderMap::new();
1446 headers.insert(
1447 reqwest::header::AUTHORIZATION,
1448 format!("Bearer {token}").parse()?,
1449 );
1450 let client = crate::tls::reqwest_client_builder()
1451 .default_headers(headers)
1452 .build()?;
1453 let endpoint = format!("http://{addr}/v1/thread-history/import");
1454 let receipt: CanonicalThreadReceipt = client
1455 .post(&endpoint)
1456 .json(&request)
1457 .send()
1458 .await?
1459 .error_for_status()?
1460 .json()
1461 .await?;
1462 let repeated: CanonicalThreadReceipt = client
1463 .post(&endpoint)
1464 .json(&request)
1465 .send()
1466 .await?
1467 .error_for_status()?
1468 .json()
1469 .await?;
1470 assert_eq!(
1471 receipt, repeated,
1472 "replay returns the same durable canonical identity"
1473 );
1474 // Inject the precise durable crash boundary after the real owner has
1475 // seeded records/checkpoint but before its operation commit marker. The
1476 // journal retains System while the Runtime item projection omits it.
1477 let operation_path = binding.data_dir.join("turn-operations").join(format!(
1478 "history_{}.json",
1479 crate::hashing::sha256_hex(
1480 format!("codewhale:history-operation:v1\0{}", request.operation_key).as_bytes()
1481 )
1482 ));
1483 let mut operation: Value = serde_json::from_slice(&fs::read(&operation_path)?)?;
1484 operation["committed"] = json!(false);
1485 crate::utils::write_atomic(&operation_path, &serde_json::to_vec(&operation)?)?;
1486 let post_seed_retry: CanonicalThreadReceipt = client
1487 .post(&endpoint)
1488 .json(&request)
1489 .send()
1490 .await?
1491 .error_for_status()?
1492 .json()
1493 .await?;
1494 assert_eq!(
1495 receipt, post_seed_retry,
1496 "a seeded System-containing full graph recovers without comparing lossy item projection"
1497 );
1498 let directory = sessions_dir.clone();
1499 let session_id = receipt.session_id.clone();
1500 let ticket = crate::test_support::env_scope_ticket();
1501 let saved = codewhale_app_server::daemon_socket::owner_work(move || {
1502 let _scope = crate::test_support::join_env_scope(ticket);
1503 Ok(SessionManager::new(directory)?.load_session_snapshot(&session_id)?)
1504 })
1505 .await?;
1506 assert_eq!(
1507 saved
1508 .journal
1509 .as_ref()
1510 .context("full imported journal")?
1511 .entries
1512 .len(),
1513 4
1514 );
1515 assert_eq!(saved.messages.len(), 3);
1516 assert_eq!(saved.metadata.runtime_store.as_ref(), Some(&binding));
1517 let full: codewhale_protocol::CanonicalThreadSnapshot = client
1518 .get(format!(
1519 "http://{addr}/v1/threads/{}/history",
1520 receipt.runtime_thread_id
1521 ))
1522 .send()
1523 .await?
1524 .error_for_status()?
1525 .json()
1526 .await?;
1527 assert_eq!(
1528 full.saved_session_id.as_deref(),
1529 Some(receipt.session_id.as_str())
1530 );
1531 assert_eq!(full.data_dir, binding.data_dir);
1532 let full: crate::session_manager::SavedSession = serde_json::from_value(full.session)?;
1533 assert_eq!(
1534 full.journal
1535 .as_ref()
1536 .context("complete read journal")?
1537 .entries
1538 .len(),
1539 4
1540 );
1541 assert_eq!(full.messages, saved.messages);
1542 let mut changed = request.clone();
1543 changed.history.messages[0].content = "changed operation input".into();
1544 assert_eq!(
1545 client.post(&endpoint).json(&changed).send().await?.status(),
1546 StatusCode::CONFLICT
1547 );
1548 let started: Value = client
1549 .post(format!(
1550 "http://{addr}/v1/threads/{}/turns",
1551 receipt.runtime_thread_id
1552 ))
1553 .json(&json!({"prompt":"continue canonical history"}))
1554 .send()
1555 .await?
1556 .error_for_status()?
1557 .json()
1558 .await?;
1559 let turn_id = started["turn"]["id"].as_str().context("actual turn id")?;
1560 assert_eq!(
1561 wait_for_terminal_turn_status(
1562 &client,
1563 addr,
1564 &receipt.runtime_thread_id,
1565 turn_id,
1566 Duration::from_secs(30)
1567 )
1568 .await?,
1569 "completed"
1570 );
1571 let actual = serde_json::to_string(&model.captured_requests())?;
1572 assert!(actual.contains("retained user context"));
1573 assert!(actual.contains("retained prior answer"));
1574 assert!(
1575 !actual.contains("retained inactive branch"),
1576 "the inactive branch remains saved but is not the active provider transcript"
1577 );
1578 let repeated: CanonicalThreadReceipt = client
1579 .post(&endpoint)
1580 .json(&request)
1581 .send()
1582 .await?
1583 .error_for_status()?
1584 .json()
1585 .await?;
1586 assert_eq!(
1587 receipt, repeated,
1588 "a legitimate successor does not remint the import"
1589 );
1590 let full: codewhale_protocol::CanonicalThreadSnapshot = client
1591 .get(format!(
1592 "http://{addr}/v1/threads/{}/history",
1593 receipt.runtime_thread_id
1594 ))
1595 .send()
1596 .await?
1597 .error_for_status()?
1598 .json()
1599 .await?;
1600 let full: crate::session_manager::SavedSession = serde_json::from_value(full.session)?;
1601 assert!(serde_json::to_string(&full)?.contains("retained inactive branch"));
1602 assert!(serde_json::to_string(&full.messages)?.contains("canonical successor"));
1603 // A pre-existing bare compatibility link adopts into this actual bound
1604 // canonical thread; its current branch and successor remain authoritative.
1605 let mut linked = request.clone();
1606 linked.operation_key = "canonical-existing-link-adoption".into();
1607 linked.target_runtime_thread_id = Some(receipt.runtime_thread_id.clone());
1608 linked.history.state_store_id = "d".repeat(64);
1609 linked.history.messages[1].content = "legacy linked branch evidence".into();
1610 let adopted: CanonicalThreadReceipt = client
1611 .post(&endpoint)
1612 .json(&linked)
1613 .send()
1614 .await?
1615 .error_for_status()?
1616 .json()
1617 .await?;
1618 assert_eq!(adopted.runtime_thread_id, receipt.runtime_thread_id);
1619 assert_eq!(adopted.session_id, receipt.session_id);
1620 let repeated: CanonicalThreadReceipt = client
1621 .post(&endpoint)
1622 .json(&linked)
1623 .send()
1624 .await?
1625 .error_for_status()?
1626 .json()
1627 .await?;
1628 assert_eq!(adopted, repeated);
1629 let full: codewhale_protocol::CanonicalThreadSnapshot = client
1630 .get(format!(
1631 "http://{addr}/v1/threads/{}/history",
1632 receipt.runtime_thread_id
1633 ))
1634 .send()
1635 .await?
1636 .error_for_status()?
1637 .json()
1638 .await?;
1639 let full: crate::session_manager::SavedSession = serde_json::from_value(full.session)?;
1640 assert!(serde_json::to_string(&full.journal)?.contains("legacy linked branch evidence"));
1641 assert!(!serde_json::to_string(&full.messages)?.contains("legacy linked branch evidence"));
1642 assert!(serde_json::to_string(&full.messages)?.contains("canonical successor"));
1643 let directory = sessions_dir;
1644 let session_id = receipt.session_id.clone();
1645 let ticket = crate::test_support::env_scope_ticket();
1646 codewhale_app_server::daemon_socket::owner_work(move || {
1647 let _scope = crate::test_support::join_env_scope(ticket);
1648 let manager = SessionManager::new(directory.clone())?;
1649 let _lease = manager.reserve_session_for_external_write(&session_id)?;
1650 fs::remove_file(directory.join(format!("{session_id}.json")))?;
1651 Ok(())
1652 })
1653 .await?;
1654 assert_eq!(
1655 client.post(&endpoint).json(&request).send().await?.status(),
1656 StatusCode::CONFLICT,
1657 "a missing committed full document is recovery, never a fresh import"
1658 );
1659 assert_eq!(
1660 client
1661 .get(format!(
1662 "http://{addr}/v1/threads/{}/history",
1663 receipt.runtime_thread_id
1664 ))
1665 .send()
1666 .await?
1667 .status(),
1668 StatusCode::CONFLICT,
1669 "full read refuses the missing authoritative document too"
1670 );
1671 assert_eq!(
1672 runtime
1673 .list_threads(
1674 crate::runtime_threads::ThreadListFilter::IncludeArchived,
1675 None
1676 )
1677 .await?
1678 .len(),
1679 1
1680 );
1681 Ok(())
1682 }
1683
1684 /// Actual authenticated owner listener plus real Runtime/Engine; only the
1685 /// provider response is scripted. Each test owns its isolated HOME and store.
1686 struct HistoryOperationFixture {
1687 addr: SocketAddr,
1688 runtime: Arc<RuntimeThreadManager>,
1689 client: reqwest::Client,
1690 sessions_dir: PathBuf,
1691 workspace: PathBuf,
1692 model: Arc<MockLlmClient>,
1693 _server: AbortOnDrop,
1694 }
1695 impl HistoryOperationFixture {
1696 async fn new(root: &FsPath) -> Result<Self> {
1697 let workspace = root.join("workspace");
1698 fs::create_dir_all(&workspace)?;
1699 let sessions_dir = root.join("sessions");
1700 let token = format!("history-operation-fixture-{}", Uuid::new_v4());
1701 let config = Config::default().with_legacy_root(
1702 Some("history-fixture-key".into()),
1703 Some("http://127.0.0.1:1/v1".into()),
1704 );
1705 let (addr, runtime, server) =
1706 spawn_test_server_with_root_token_mobile_workspace_and_overrides(
1707 root.to_path_buf(),
1708 sessions_dir.clone(),
1709 Some(token.clone()),
1710 false,
1711 workspace.clone(),
1712 TestServerOverrides {
1713 config: Some(config),
1714 ..Default::default()
1715 },
1716 )
1717 .await?
1718 .context("history acceptance requires an actual owner listener")?;
1719 let model = Arc::new(MockLlmClient::new(vec![canned::simple_text_turn(
1720 "actual successor",
1721 )]));
1722 runtime.set_test_model_client(model.clone());
1723 let mut headers = reqwest::header::HeaderMap::new();
1724 headers.insert(
1725 reqwest::header::AUTHORIZATION,
1726 format!("Bearer {token}").parse()?,
1727 );
1728 Ok(Self {
1729 addr,
1730 runtime,
1731 client: crate::tls::reqwest_client_builder()
1732 .default_headers(headers)
1733 .build()?,
1734 sessions_dir,
1735 workspace,
1736 model,
1737 _server: AbortOnDrop(server.abort_handle()),
1738 })
1739 }
1740 fn url(&self, suffix: &str) -> String {
1741 format!("http://{}{suffix}", self.addr)
1742 }
1743 fn mutation(
1744 &self,
1745 key: &str,
1746 mutation: codewhale_protocol::CanonicalThreadMutation,
1747 ) -> codewhale_protocol::CanonicalThreadMutationRequest {
1748 let binding = self.runtime.session_store_binding();
1749 codewhale_protocol::CanonicalThreadMutationRequest {
1750 version: 1,
1751 operation_key: key.into(),
1752 expected_data_dir: binding.data_dir,
1753 expected_execution_scope: binding.execution_scope,
1754 workspace: self.workspace.clone(),
1755 mutation,
1756 }
1757 }
1758 async fn submit(
1759 &self,
1760 request: &codewhale_protocol::CanonicalThreadMutationRequest,
1761 ) -> Result<codewhale_protocol::CanonicalThreadReceipt> {
1762 Ok(self
1763 .client
1764 .post(self.url("/v1/thread-history/mutate"))
1765 .json(request)
1766 .send()
1767 .await?
1768 .error_for_status()?
1769 .json()
1770 .await?)
1771 }
1772 async fn snapshot(&self, id: &str) -> Result<codewhale_protocol::CanonicalThreadSnapshot> {
1773 Ok(self
1774 .client
1775 .get(self.url(&format!("/v1/threads/{id}/history")))
1776 .send()
1777 .await?
1778 .error_for_status()?
1779 .json()
1780 .await?)
1781 }
1782 async fn lookup(
1783 &self,
1784 request: &codewhale_protocol::CanonicalThreadMutationRequest,
1785 ) -> Result<codewhale_protocol::CanonicalThreadOperationStatus> {
1786 Ok(self
1787 .client
1788 .post(self.url("/v1/thread-history/operations/lookup"))
1789 .json(&codewhale_protocol::CanonicalThreadOperationLookup {
1790 version: request.version,
1791 operation_key: request.operation_key.clone(),
1792 expected_data_dir: request.expected_data_dir.clone(),
1793 expected_execution_scope: request.expected_execution_scope.clone(),
1794 workspace: request.workspace.clone(),
1795 })
1796 .send()
1797 .await?
1798 .error_for_status()?
1799 .json()
1800 .await?)
1801 }
1802 async fn import_branched_source(&self) -> Result<codewhale_protocol::CanonicalThreadReceipt> {
1803 self.import_branched_source_with_goal(None).await
1804 }
1805 async fn import_branched_source_with_goal(
1806 &self,
1807 goal: Option<codewhale_protocol::ThreadGoal>,
1808 ) -> Result<codewhale_protocol::CanonicalThreadReceipt> {
1809 let binding = self.runtime.session_store_binding();
1810 let request = codewhale_protocol::CanonicalHistoryImportRequest {
1811 version: 1,
1812 operation_key: "operation-source-full-branches".into(),
1813 expected_data_dir: binding.data_dir,
1814 expected_execution_scope: binding.execution_scope,
1815 target_runtime_thread_id: None,
1816 workspace: self.workspace.clone(),
1817 model: Some("deepseek-v4-pro".into()),
1818 history: codewhale_protocol::LegacyThreadHistory {
1819 goal,
1820 version: 1,
1821 state_store_id: "f".repeat(64),
1822 thread_id: "original-operation-source".into(),
1823 current_leaf_id: Some(3),
1824 messages: [
1825 (1, None, "system", "system retained"),
1826 (2, Some(1), "user", "original active prompt"),
1827 (3, Some(2), "assistant", "original active answer"),
1828 (4, Some(1), "user", "inactive branch selected for fork"),
1829 ]
1830 .into_iter()
1831 .map(
1832 |(id, parent_entry_id, role, content)| codewhale_protocol::MessageRecord {
1833 id,
1834 thread_id: "original-operation-source".into(),
1835 role: role.into(),
1836 content: content.into(),
1837 item: None,
1838 created_at: 1_700_000_000 + id,
1839 parent_entry_id,
1840 },
1841 )
1842 .collect(),
1843 },
1844 };
1845 Ok(self
1846 .client
1847 .post(self.url("/v1/thread-history/import"))
1848 .json(&request)
1849 .send()
1850 .await?
1851 .error_for_status()?
1852 .json()
1853 .await?)
1854 }
1855 }
1856
1857 #[tokio::test]
1858 async fn canonical_create_commits_one_identity_and_read_only_lookup_has_no_turn_effect()
1859 -> Result<()> {
1860 use codewhale_protocol::{CanonicalThreadMutation, CanonicalThreadOperationStatus};
1861 let _env = lock_test_env();
1862 let temp = tempfile::tempdir()?;
1863 let root = temp.path().canonicalize()?;
1864 let _home = EnvVarGuard::set("CODEWHALE_HOME", root.join("home"));
1865 let fixture = HistoryOperationFixture::new(&root).await?;
1866 let request = fixture.mutation(
1867 "actual-create-intent",
1868 CanonicalThreadMutation::Create {
1869 config: json!({"model":"deepseek-v4-pro","allow_shell":false}),
1870 },
1871 );
1872 assert!(matches!(
1873 fixture.lookup(&request).await?,
1874 CanonicalThreadOperationStatus::Absent
1875 ));
1876 let receipt = fixture.submit(&request).await?;
1877 assert_eq!(fixture.submit(&request).await?, receipt);
1878 match fixture.lookup(&request).await? {
1879 CanonicalThreadOperationStatus::Committed {
1880 receipt: recovered,
1881 association,
1882 } => {
1883 assert_eq!(receipt, recovered);
1884 assert_eq!(
1885 association.kind,
1886 codewhale_protocol::CanonicalThreadOperationKind::Create
1887 );
1888 assert!(association.source_runtime_thread_id.is_none());
1889 }
1890 other => anyhow::bail!("unexpected settled lookup: {other:?}"),
1891 }
1892 let snapshot = fixture.snapshot(&receipt.runtime_thread_id).await?;
1893 let again = fixture.snapshot(&receipt.runtime_thread_id).await?;
1894 assert_eq!(
1895 snapshot.document_digest, again.document_digest,
1896 "read-only projection is stable"
1897 );
1898 assert_eq!(snapshot.session, again.session);
1899 assert_eq!(
1900 fixture
1901 .runtime
1902 .list_threads(
1903 crate::runtime_threads::ThreadListFilter::IncludeArchived,
1904 None
1905 )
1906 .await?
1907 .len(),
1908 1
1909 );
1910 assert!(
1911 fixture
1912 .runtime
1913 .get_thread_detail(&receipt.runtime_thread_id)
1914 .await?
1915 .turns
1916 .is_empty()
1917 );
1918 assert!(fixture.model.captured_requests().is_empty());
1919 let mut changed = request.clone();
1920 changed.mutation = CanonicalThreadMutation::Create {
1921 config: json!({"model":"deepseek-v4-pro","allow_shell":true}),
1922 };
1923 assert_eq!(
1924 fixture
1925 .client
1926 .post(fixture.url("/v1/thread-history/mutate"))
1927 .json(&changed)
1928 .send()
1929 .await?
1930 .status(),
1931 StatusCode::CONFLICT
1932 );
1933 let mut wrong_scope = codewhale_protocol::CanonicalThreadOperationLookup {
1934 version: 1,
1935 operation_key: request.operation_key,
1936 expected_data_dir: request.expected_data_dir,
1937 expected_execution_scope: request.expected_execution_scope,
1938 workspace: root.join("another-scope"),
1939 };
1940 assert_eq!(
1941 fixture
1942 .client
1943 .post(fixture.url("/v1/thread-history/operations/lookup"))
1944 .json(&wrong_scope)
1945 .send()
1946 .await?
1947 .status(),
1948 StatusCode::CONFLICT
1949 );
1950 wrong_scope.workspace = fixture.workspace.clone();
1951 wrong_scope.expected_execution_scope = "another-owner".into();
1952 assert_eq!(
1953 fixture
1954 .client
1955 .post(fixture.url("/v1/thread-history/operations/lookup"))
1956 .json(&wrong_scope)
1957 .send()
1958 .await?
1959 .status(),
1960 StatusCode::CONFLICT
1961 );
1962 Ok(())
1963 }
1964
1965 #[tokio::test]
1966 async fn canonical_fork_keeps_all_branches_selects_exact_leaf_and_recovers_after_successor()
1967 -> Result<()> {
1968 use codewhale_protocol::{
1969 CanonicalHistoryOptions, CanonicalHistorySource, CanonicalThreadMutation,
1970 CanonicalThreadOperationStatus,
1971 };
1972 let _env = lock_test_env();
1973 let temp = tempfile::tempdir()?;
1974 let root = temp.path().canonicalize()?;
1975 let _home = EnvVarGuard::set("CODEWHALE_HOME", root.join("home"));
1976 let fixture = HistoryOperationFixture::new(&root).await?;
1977 let original = fixture.import_branched_source().await?;
1978 let snapshot = fixture.snapshot(&original.runtime_thread_id).await?;
1979 let source: crate::session_manager::SavedSession =
1980 serde_json::from_value(snapshot.session.clone())?;
1981 let selected = source.journal.as_ref().context("source graph")?.entries[3]
1982 .id
1983 .clone();
1984 let mut request = fixture.mutation(
1985 "actual-fork-intent",
1986 CanonicalThreadMutation::Fork {
1987 source: CanonicalHistorySource::Thread {
1988 runtime_thread_id: original.runtime_thread_id.clone(),
1989 expected_document_digest: snapshot.document_digest,
1990 },
1991 options: CanonicalHistoryOptions::default(),
1992 selected_entry_id: Some(selected.clone()),
1993 },
1994 );
1995 // A selected target root is independent of the source root; the source is unchanged.
1996 request.workspace = root.join("fork-workspace");
1997 fs::create_dir_all(&request.workspace)?;
1998 let receipt = fixture.submit(&request).await?;
1999 assert_ne!(original.session_id, receipt.session_id);
2000 assert_ne!(original.runtime_thread_id, receipt.runtime_thread_id);
2001 // The selected branch has two messages, unlike the three-message source.
2002 // Inspect raw publication before the protected loader can restore its count.
2003 let raw: crate::session_manager::SavedSession = serde_json::from_slice(&fs::read(
2004 fixture
2005 .sessions_dir
2006 .join(format!("{}.json", receipt.session_id)),
2007 )?)?;
2008 assert_eq!(raw.metadata.message_count, 2);
2009 assert_eq!(raw.messages.len(), 2);
2010 assert_ne!(raw.metadata.message_count, source.metadata.message_count);
2011 let prepared_graph = raw.journal.as_ref().context("raw selected full graph")?;
2012 assert_eq!(
2013 prepared_graph.entries,
2014 source.journal.as_ref().unwrap().entries
2015 );
2016 assert_eq!(prepared_graph.leaf_id.as_deref(), Some(selected.as_str()));
2017 assert_eq!(raw.leaf_id, prepared_graph.leaf_id);
2018 assert_eq!(raw.messages, prepared_graph.to_messages());
2019 let operation_path = request
2020 .expected_data_dir
2021 .join("turn-operations")
2022 .join(format!(
2023 "history_{}.json",
2024 crate::hashing::sha256_hex(
2025 format!("codewhale:history-operation:v1\0{}", request.operation_key).as_bytes()
2026 )
2027 ));
2028 let mut prepared: Value = serde_json::from_slice(&fs::read(&operation_path)?)?;
2029 assert_eq!(
2030 prepared["target_document_digest"],
2031 json!(super::super::thread_history::saved_document_digest(&raw)?),
2032 "Fork recovery binds the exact shorter branch count published to disk"
2033 );
2034 // Model interruption after complete publication but before the intent commit.
2035 prepared["committed"] = json!(false);
2036 crate::utils::write_atomic(&operation_path, &serde_json::to_vec(&prepared)?)?;
2037 let association = match fixture.lookup(&request).await? {
2038 CanonicalThreadOperationStatus::Pending { association, .. } => association,
2039 other => anyhow::bail!("expected prepared Fork intent: {other:?}"),
2040 };
2041 let recovery = codewhale_protocol::CanonicalThreadOperationRecovery {
2042 operation: codewhale_protocol::CanonicalThreadOperationLookup {
2043 version: request.version,
2044 operation_key: request.operation_key.clone(),
2045 expected_data_dir: request.expected_data_dir.clone(),
2046 expected_execution_scope: request.expected_execution_scope.clone(),
2047 workspace: request.workspace.clone(),
2048 },
2049 association,
2050 };
2051 let prepared_detail = fixture
2052 .runtime
2053 .get_thread_detail(&receipt.runtime_thread_id)
2054 .await?;
2055 for _ in 0..2 {
2056 let recovered: CanonicalThreadOperationStatus = fixture
2057 .client
2058 .post(fixture.url("/v1/thread-history/operations/recover"))
2059 .json(&recovery)
2060 .send()
2061 .await?
2062 .error_for_status()?
2063 .json()
2064 .await?;
2065 match recovered {
2066 CanonicalThreadOperationStatus::Committed {
2067 receipt: settled,
2068 association,
2069 } => {
2070 assert_eq!(settled, receipt);
2071 assert_eq!(association, recovery.association);
2072 }
2073 other => anyhow::bail!("prepared exact-key Fork recovery did not settle: {other:?}"),
2074 }
2075 }
2076 let recovered_detail = fixture
2077 .runtime
2078 .get_thread_detail(&receipt.runtime_thread_id)
2079 .await?;
2080 assert_eq!(prepared_detail.turns.len(), recovered_detail.turns.len());
2081 assert_eq!(prepared_detail.items.len(), recovered_detail.items.len());
2082 assert!(fixture.model.captured_requests().is_empty());
2083 let fork = fixture.snapshot(&receipt.runtime_thread_id).await?;
2084 let fork: crate::session_manager::SavedSession = serde_json::from_value(fork.session)?;
2085 assert_eq!(
2086 fork.journal.as_ref().context("fork graph")?.entries,
2087 source.journal.as_ref().unwrap().entries,
2088 "inactive and active branches remain complete"
2089 );
2090 assert_eq!(fork.leaf_id.as_deref(), Some(selected.as_str()));
2091 assert_eq!(
2092 fork.metadata.parent_session_id.as_deref(),
2093 Some(original.session_id.as_str())
2094 );
2095 assert_eq!(fork.metadata.workspace, request.workspace);
2096 assert_eq!(
2097 fixture.snapshot(&original.runtime_thread_id).await?.session,
2098 snapshot.session,
2099 "fork never rewrites source"
2100 );
2101 let started: Value = fixture.client.post(fixture.url(&format!("/v1/threads/{}/turns",receipt.runtime_thread_id)))
2102 .json(&json!({"prompt":"continue selected branch","operation_key":"fork-successor","expected_workspace":request.workspace})).send().await?.error_for_status()?.json().await?;
2103 let turn_id = started["turn"]["id"]
2104 .as_str()
2105 .context("actual successor turn")?;
2106 assert_eq!(
2107 wait_for_terminal_turn_status(
2108 &fixture.client,
2109 fixture.addr,
2110 &receipt.runtime_thread_id,
2111 turn_id,
2112 Duration::from_secs(30)
2113 )
2114 .await?,
2115 "completed"
2116 );
2117 let requests = serde_json::to_string(&fixture.model.captured_requests())?;
2118 assert!(requests.contains("inactive branch selected for fork"));
2119 assert!(!requests.contains("original active answer"));
2120 let before = fixture
2121 .runtime
2122 .get_thread_detail(&receipt.runtime_thread_id)
2123 .await?;
2124 let recovered = fixture.lookup(&request).await?;
2125 match recovered {
2126 CanonicalThreadOperationStatus::Committed {
2127 receipt: known,
2128 association,
2129 } => {
2130 assert_eq!(known, receipt);
2131 assert_eq!(
2132 association.kind,
2133 codewhale_protocol::CanonicalThreadOperationKind::Fork
2134 );
2135 assert_eq!(
2136 association.source_session_id.as_deref(),
2137 Some(original.session_id.as_str())
2138 );
2139 }
2140 other => anyhow::bail!("unexpected successor recovery: {other:?}"),
2141 }
2142 let after = fixture
2143 .runtime
2144 .get_thread_detail(&receipt.runtime_thread_id)
2145 .await?;
2146 assert_eq!(before.turns.len(), after.turns.len());
2147 assert_eq!(before.items.len(), after.items.len());
2148 let mut refreshed = request.clone();
2149 if let CanonicalThreadMutation::Fork {
2150 source:
2151 CanonicalHistorySource::Thread {
2152 expected_document_digest,
2153 ..
2154 },
2155 ..
2156 } = &mut refreshed.mutation
2157 {
2158 *expected_document_digest = "0".repeat(64);
2159 }
2160 assert_eq!(
2161 fixture
2162 .client
2163 .post(fixture.url("/v1/thread-history/mutate"))
2164 .json(&refreshed)
2165 .send()
2166 .await?
2167 .status(),
2168 StatusCode::CONFLICT,
2169 "read-only lookup, rather than a reminted source proposal, recovers the retained key"
2170 );
2171 Ok(())
2172 }
2173
2174 #[tokio::test]
2175 async fn canonical_resume_offered_suffix_settles_once_and_full_graph_witness_refuses_tamper()
2176 -> Result<()> {
2177 use codewhale_protocol::{
2178 CanonicalHistoryOptions, CanonicalHistorySource, CanonicalThreadMutation,
2179 CanonicalThreadOperationStatus,
2180 };
2181 let _env = lock_test_env();
2182 let temp = tempfile::tempdir()?;
2183 let root = temp.path().canonicalize()?;
2184 let _home = EnvVarGuard::set("CODEWHALE_HOME", root.join("home"));
2185 let fixture = HistoryOperationFixture::new(&root).await?;
2186 let original = fixture.import_branched_source().await?;
2187 let snapshot = fixture.snapshot(&original.runtime_thread_id).await?;
2188 let source: crate::session_manager::SavedSession = serde_json::from_value(snapshot.session)?;
2189 let mut offered = source
2190 .messages
2191 .iter()
2192 .map(serde_json::to_value)
2193 .collect::<std::result::Result<Vec<_>, _>>()?;
2194 offered.push(serde_json::to_value(codewhale_models::Message {
2195 role: codewhale_models::Role::User,
2196 content: vec![codewhale_models::ContentBlock::Text {
2197 text: "genuine new offered suffix".into(),
2198 cache_control: None,
2199 }],
2200 })?);
2201 let request = fixture.mutation(
2202 "actual-resume-intent",
2203 CanonicalThreadMutation::Resume {
2204 source: CanonicalHistorySource::Thread {
2205 runtime_thread_id: original.runtime_thread_id.clone(),
2206 expected_document_digest: snapshot.document_digest,
2207 },
2208 options: CanonicalHistoryOptions {
2209 expected_session_goal_digest: None,
2210 offered_history: offered,
2211 overrides: json!({"base_instructions":"keep canonical receipt"}),
2212 source_path: Some(
2213 fixture
2214 .sessions_dir
2215 .join(format!("{}.json", original.session_id)),
2216 ),
2217 },
2218 },
2219 );
2220 let receipt = fixture.submit(&request).await?;
2221 assert_eq!(receipt.runtime_thread_id, original.runtime_thread_id);
2222 assert_eq!(receipt.session_id, original.session_id);
2223 // Inspect the exact publication before a protected loader can normalize it.
2224 let raw: crate::session_manager::SavedSession = serde_json::from_slice(&fs::read(
2225 fixture
2226 .sessions_dir
2227 .join(format!("{}.json", receipt.session_id)),
2228 )?)?;
2229 assert_eq!(raw.metadata.message_count, 4);
2230 assert_eq!(raw.messages.len(), 4);
2231 let published_journal = raw.journal.as_ref().context("published full journal")?;
2232 let offered_leaf = format!(
2233 "offered:{}:3",
2234 crate::hashing::sha256_hex(request.operation_key.as_bytes())
2235 );
2236 assert_eq!(
2237 published_journal.leaf_id.as_deref(),
2238 Some(offered_leaf.as_str())
2239 );
2240 assert_eq!(raw.leaf_id, published_journal.leaf_id);
2241 assert_eq!(raw.messages, published_journal.to_messages());
2242 let before = fixture
2243 .runtime
2244 .get_thread_detail(&receipt.runtime_thread_id)
2245 .await?;
2246 assert_eq!(
2247 before.turns.len(),
2248 2,
2249 "only the non-overlapping offered user suffix is seeded"
2250 );
2251 let full = fixture.snapshot(&receipt.runtime_thread_id).await?;
2252 let full: crate::session_manager::SavedSession = serde_json::from_value(full.session)?;
2253 assert_eq!(full.journal.as_ref().unwrap().entries.len(), 5);
2254 assert_eq!(full.messages.len(), 4);
2255 assert!(
2256 full.system_prompt
2257 .as_deref()
2258 .is_some_and(|text| text.contains("keep canonical receipt"))
2259 );
2260 // Actual completed seed/checkpoint, then emulate interruption before the intent marker.
2261 let operation_path = request
2262 .expected_data_dir
2263 .join("turn-operations")
2264 .join(format!(
2265 "history_{}.json",
2266 crate::hashing::sha256_hex(
2267 format!("codewhale:history-operation:v1\0{}", request.operation_key).as_bytes()
2268 )
2269 ));
2270 let mut operation: Value = serde_json::from_slice(&fs::read(&operation_path)?)?;
2271 assert_eq!(
2272 operation["target_document_digest"],
2273 json!(super::super::thread_history::saved_document_digest(&raw)?),
2274 "the retained recovery witness must bind the exact count and offered leaf published to disk"
2275 );
2276 operation["committed"] = json!(false);
2277 crate::utils::write_atomic(&operation_path, &serde_json::to_vec(&operation)?)?;
2278 let association = match fixture.lookup(&request).await? {
2279 CanonicalThreadOperationStatus::Pending { association, .. } => association,
2280 other => anyhow::bail!("expected published pending intent: {other:?}"),
2281 };
2282 let recovery = codewhale_protocol::CanonicalThreadOperationRecovery {
2283 operation: codewhale_protocol::CanonicalThreadOperationLookup {
2284 version: request.version,
2285 operation_key: request.operation_key.clone(),
2286 expected_data_dir: request.expected_data_dir.clone(),
2287 expected_execution_scope: request.expected_execution_scope.clone(),
2288 workspace: request.workspace.clone(),
2289 },
2290 association,
2291 };
2292 // Reconstructing a proposal from the changed current source is not the
2293 // original intent. Exact-key recovery needs neither that proposal nor a
2294 // second source body, and must not remint its already-published session.
2295 let changed = fixture.snapshot(&receipt.runtime_thread_id).await?;
2296 let mut refreshed = request.clone();
2297 let CanonicalThreadMutation::Resume {
2298 source:
2299 CanonicalHistorySource::Thread {
2300 expected_document_digest,
2301 ..
2302 },
2303 ..
2304 } = &mut refreshed.mutation
2305 else {
2306 anyhow::bail!("Resume fixture changed");
2307 };
2308 *expected_document_digest = changed.document_digest;
2309 assert_eq!(
2310 fixture
2311 .client
2312 .post(fixture.url("/v1/thread-history/mutate"))
2313 .json(&refreshed)
2314 .send()
2315 .await?
2316 .status(),
2317 StatusCode::CONFLICT
2318 );
2319 let mut wrong = recovery.clone();
2320 wrong.association.kind = codewhale_protocol::CanonicalThreadOperationKind::Fork;
2321 assert_eq!(
2322 fixture
2323 .client
2324 .post(fixture.url("/v1/thread-history/operations/recover"))
2325 .json(&wrong)
2326 .send()
2327 .await?
2328 .status(),
2329 StatusCode::CONFLICT
2330 );
2331 let mut wrong = recovery.clone();
2332 wrong.operation.workspace = root.join("different-selection");
2333 assert_eq!(
2334 fixture
2335 .client
2336 .post(fixture.url("/v1/thread-history/operations/recover"))
2337 .json(&wrong)
2338 .send()
2339 .await?
2340 .status(),
2341 StatusCode::CONFLICT
2342 );
2343 for _ in 0..2 {
2344 let recovered: CanonicalThreadOperationStatus = fixture
2345 .client
2346 .post(fixture.url("/v1/thread-history/operations/recover"))
2347 .json(&recovery)
2348 .send()
2349 .await?
2350 .error_for_status()?
2351 .json()
2352 .await?;
2353 match recovered {
2354 CanonicalThreadOperationStatus::Committed {
2355 receipt: settled,
2356 association,
2357 } => {
2358 assert_eq!(settled, receipt);
2359 assert_eq!(association, recovery.association);
2360 }
2361 other => anyhow::bail!("prepared exact-key recovery did not settle: {other:?}"),
2362 }
2363 }
2364 assert_eq!(fixture.submit(&request).await?, receipt);
2365 let after = fixture
2366 .runtime
2367 .get_thread_detail(&receipt.runtime_thread_id)
2368 .await?;
2369 assert_eq!(before.turns.len(), after.turns.len());
2370 assert_eq!(before.items.len(), after.items.len());
2371 assert!(
2372 fixture.model.captured_requests().is_empty(),
2373 "history admission and recovery never execute a provider turn"
2374 );
2375 let file = fixture
2376 .sessions_dir
2377 .join(format!("{}.json", receipt.session_id));
2378 let mut tampered: crate::session_manager::SavedSession =
2379 serde_json::from_slice(&fs::read(&file)?)?;
2380 let entry = &mut tampered
2381 .journal
2382 .as_mut()
2383 .context("tampered full graph")?
2384 .entries[3];
2385 let crate::session_tree::SessionEntryKind::Message { message } = &mut entry.kind else {
2386 anyhow::bail!("inactive fixture entry lost its message");
2387 };
2388 let codewhale_models::ContentBlock::Text { text, .. } = &mut message.content[0] else {
2389 anyhow::bail!("inactive fixture text missing");
2390 };
2391 *text = "tampered inactive branch".into();
2392 crate::utils::write_atomic(&file, &serde_json::to_vec(&tampered)?)?;
2393 let lookup = codewhale_protocol::CanonicalThreadOperationLookup {
2394 version: 1,
2395 operation_key: request.operation_key,
2396 expected_data_dir: request.expected_data_dir,
2397 expected_execution_scope: request.expected_execution_scope,
2398 workspace: request.workspace,
2399 };
2400 assert_eq!(
2401 fixture
2402 .client
2403 .post(fixture.url("/v1/thread-history/operations/lookup"))
2404 .json(&lookup)
2405 .send()
2406 .await?
2407 .status(),
2408 StatusCode::CONFLICT,
2409 "matching entry count/selected leaf cannot hide a changed inactive branch"
2410 );
2411 assert_eq!(
2412 fixture
2413 .client
2414 .post(fixture.url("/v1/thread-history/operations/recover"))
2415 .json(&recovery)
2416 .send()
2417 .await?
2418 .status(),
2419 StatusCode::CONFLICT,
2420 "exact-key recovery cannot conceal a changed inactive branch"
2421 );
2422 Ok(())
2423 }
2424
2425 #[tokio::test]
2426 async fn canonical_import_preserves_paused_goal_and_retry_never_overwrites_later_goal() -> Result<()>
2427 {
2428 let _env = lock_test_env();
2429 let temp = tempfile::tempdir()?;
2430 let root = temp.path().canonicalize()?;
2431 let _home = EnvVarGuard::set("CODEWHALE_HOME", root.join("home"));
2432 let fixture = HistoryOperationFixture::new(&root).await?;
2433 let source = codewhale_protocol::ThreadGoal {
2434 thread_id: "original-operation-source".into(),
2435 goal_id: Uuid::new_v4().to_string(),
2436 objective: "preserve the actual prior goal".into(),
2437 status: codewhale_protocol::ThreadGoalStatus::Active,
2438 token_budget: Some(1000),
2439 tokens_used: 41,
2440 time_used_seconds: 9,
2441 continuation_count: 2,
2442 created_at: 1_700_000_000,
2443 updated_at: 1_700_000_010,
2444 last_gap_fingerprint: None,
2445 repeated_gap_count: 0,
2446 last_gap_pass: None,
2447 pause_reason: None,
2448 };
2449 let receipt = fixture
2450 .import_branched_source_with_goal(Some(source.clone()))
2451 .await?;
2452 let goal = fixture
2453 .runtime
2454 .get_goal(&receipt.runtime_thread_id)
2455 .await?
2456 .context("imported actual owner goal")?;
2457 let mut expected = source.clone();
2458 expected.thread_id = receipt.runtime_thread_id.clone();
2459 expected.status = codewhale_protocol::ThreadGoalStatus::Paused;
2460 assert_eq!(
2461 goal, expected,
2462 "source counters/timestamps persist without authorizing provider work"
2463 );
2464 assert!(fixture.model.captured_requests().is_empty());
2465 let mut later = goal;
2466 later.status = codewhale_protocol::ThreadGoalStatus::Complete;
2467 later.updated_at += 7;
2468 later.tokens_used += 4;
2469 fixture.runtime.save_goal(later.clone()).await?;
2470 assert_eq!(
2471 fixture
2472 .import_branched_source_with_goal(Some(source.clone()))
2473 .await?,
2474 receipt
2475 );
2476 assert_eq!(
2477 fixture.runtime.get_goal(&receipt.runtime_thread_id).await?,
2478 Some(later)
2479 );
2480 let mut conflicting = source;
2481 conflicting.thread_id = "wrong-source".into();
2482 assert!(
2483 fixture
2484 .import_branched_source_with_goal(Some(conflicting))
2485 .await
2486 .is_err()
2487 );
2488 assert!(fixture.model.captured_requests().is_empty());
2489 Ok(())
2490 }
2491
2492 #[tokio::test]
2493 async fn canonical_fork_copies_local_goal_under_source_witness_and_never_replays_evolution()
2494 -> Result<()> {
2495 use crate::session_manager::{SessionGoalState, SessionGoalStatus, SessionManager};
2496 use codewhale_protocol::{CanonicalHistorySource, CanonicalThreadMutation};
2497 let _env = lock_test_env();
2498 let temp = tempfile::tempdir()?;
2499 let root = temp.path().canonicalize()?;
2500 let _home = EnvVarGuard::set("CODEWHALE_HOME", root.join("home"));
2501 let fixture = HistoryOperationFixture::new(&root).await?;
2502 let original = fixture.import_branched_source().await?;
2503 let manager = SessionManager::new(fixture.sessions_dir.clone())?;
2504 let goal: SessionGoalState = serde_json::from_value(json!({
2505 "objective":"preserve local goal without launching work", "status":"active",
2506 "token_budget":1000,"tokens_used":13,"time_used_seconds":8,
2507 "continuation_count":2,"elapsed_seconds":9,"goal_id":"source-local-goal",
2508 "last_gap_fingerprint":"a".repeat(64),"repeated_gap_count":1,"last_gap_pass":1
2509 }))?;
2510 manager.save_session_goal(&original.session_id, Some(&goal))?;
2511 let snapshot = fixture.snapshot(&original.runtime_thread_id).await?;
2512 let request = fixture.mutation(
2513 "local-goal-fork",
2514 CanonicalThreadMutation::Fork {
2515 source: CanonicalHistorySource::Thread {
2516 runtime_thread_id: original.runtime_thread_id.clone(),
2517 expected_document_digest: snapshot.document_digest,
2518 },
2519 options: Default::default(),
2520 selected_entry_id: None,
2521 },
2522 );
2523 let receipt = fixture.submit(&request).await?;
2524 let copied = manager
2525 .load_session_goal(&receipt.session_id)?
2526 .context("fork local goal sidecar")?;
2527 let mut paused = goal.clone();
2528 paused.status = SessionGoalStatus::Paused;
2529 paused.pause_reason = None;
2530 assert_eq!(copied, paused);
2531 assert_eq!(
2532 manager.load_session_goal(&original.session_id)?,
2533 Some(goal.clone())
2534 );
2535 let mut evolved = copied;
2536 evolved.status = SessionGoalStatus::Complete;
2537 evolved.tokens_used += 7;
2538 manager.save_session_goal(&receipt.session_id, Some(&evolved))?;
2539 assert_eq!(fixture.submit(&request).await?, receipt);
2540 assert_eq!(
2541 manager.load_session_goal(&receipt.session_id)?,
2542 Some(evolved)
2543 );
2544 assert!(fixture.model.captured_requests().is_empty());
2545
2546 let stale = fixture.snapshot(&original.runtime_thread_id).await?;
2547 let changed_request = fixture.mutation(
2548 "changed-source-goal",
2549 CanonicalThreadMutation::Fork {
2550 source: CanonicalHistorySource::Thread {
2551 runtime_thread_id: original.runtime_thread_id.clone(),
2552 expected_document_digest: stale.document_digest,
2553 },
2554 options: Default::default(),
2555 selected_entry_id: None,
2556 },
2557 );
2558 manager.save_session_goal(&original.session_id, None)?;
2559 assert_eq!(
2560 fixture
2561 .client
2562 .post(fixture.url("/v1/thread-history/mutate"))
2563 .json(&changed_request)
2564 .send()
2565 .await?
2566 .status(),
2567 StatusCode::CONFLICT,
2568 "missing previously captured sidecar is a changed source, never an empty inherited goal"
2569 );
2570 let goal_path = fixture
2571 .sessions_dir
2572 .join(".goals")
2573 .join(format!("{}.json", original.session_id));
2574 crate::utils::write_atomic(&goal_path, b"{corrupt-local-goal")?;
2575 assert_eq!(
2576 fixture
2577 .client
2578 .get(fixture.url(&format!(
2579 "/v1/threads/{}/history",
2580 original.runtime_thread_id
2581 )))
2582 .send()
2583 .await?
2584 .status(),
2585 StatusCode::CONFLICT,
2586 "corrupt sidecar remains visible recovery evidence"
2587 );
2588 Ok(())
2589 }
2590
2591 #[tokio::test]
2592 async fn native_fork_routes_keep_full_graph_local_goal_and_captured_session_directory() -> Result<()>
2593 {
2594 use crate::session_manager::{
2595 SavedSession, SessionGoalState, SessionGoalStatus, SessionManager,
2596 };
2597 let _env = lock_test_env();
2598 let temp = tempfile::tempdir()?;
2599 let root = temp.path().canonicalize()?;
2600 let _home = EnvVarGuard::set("CODEWHALE_HOME", root.join("home"));
2601 let fixture = HistoryOperationFixture::new(&root).await?;
2602 assert_ne!(
2603 fixture.sessions_dir,
2604 crate::session_manager::default_sessions_dir()?
2605 );
2606 let original = fixture.import_branched_source().await?;
2607 let manager = SessionManager::new(fixture.sessions_dir.clone())?;
2608 let mut source = manager.load_session_snapshot_bounded(
2609 &original.session_id,
2610 codewhale_protocol::MAX_CANONICAL_HISTORY_BYTES,
2611 )?;
2612 source.window_title = Some("captured native source window".into());
2613 source.work_state = Some(crate::session_manager::SessionWorkState::default());
2614 source.metadata.cost.session_cost_usd = 1.25;
2615 manager.save_session(&source)?;
2616 let source_bytes = fs::read(
2617 fixture
2618 .sessions_dir
2619 .join(format!("{}.json", original.session_id)),
2620 )?;
2621 let goal: SessionGoalState = serde_json::from_value(json!({
2622 "objective":"native sidecar continuity", "status":"active", "token_budget":900,
2623 "tokens_used":17,"time_used_seconds":4,"continuation_count":2,"elapsed_seconds":5,
2624 "goal_id":"native-goal-revision","last_gap_fingerprint":"b".repeat(64),
2625 "repeated_gap_count":1,"last_gap_pass":1
2626 }))?;
2627 manager.save_session_goal(&original.session_id, Some(&goal))?;
2628 let source_entries = source
2629 .journal
2630 .as_ref()
2631 .context("native source journal")?
2632 .entries
2633 .clone();
2634 let source_detail = fixture
2635 .runtime
2636 .get_thread_detail(&original.runtime_thread_id)
2637 .await?;
2638 let turn_id = source_detail
2639 .turns
2640 .first()
2641 .context("seeded native turn")?
2642 .id
2643 .clone();
2644 for (suffix, body, empty) in [
2645 ("fork", None, false),
2646 ("fork-at-turn", Some(json!({"turn_id":turn_id})), false),
2647 ("undo", Some(json!({"depth":0})), true),
2648 ] {
2649 let request = fixture.client.post(fixture.url(&format!(
2650 "/v1/threads/{}/{suffix}",
2651 original.runtime_thread_id
2652 )));
2653 let request = if let Some(body) = body {
2654 request.json(&body)
2655 } else {
2656 request
2657 };
2658 let response: Value = request.send().await?.error_for_status()?.json().await?;
2659 let fork: crate::runtime_threads::ThreadRecord =
2660 serde_json::from_value(if suffix == "fork" {
2661 response
2662 } else {
2663 response["thread"].clone()
2664 })?;
2665 assert_ne!(fork.session_id, Some(original.session_id.clone()));
2666 let id = fork
2667 .session_id
2668 .as_ref()
2669 .context("native fork owns a session")?;
2670 let saved: SavedSession = manager
2671 .load_session_snapshot_bounded(id, codewhale_protocol::MAX_CANONICAL_HISTORY_BYTES)?;
2672 let graph = saved.journal.as_ref().context("full native fork journal")?;
2673 assert!(
2674 source_entries
2675 .iter()
2676 .all(|entry| graph.entries.contains(entry)),
2677 "every inactive branch and original entry survives {suffix}"
2678 );
2679 assert_eq!(saved.messages.is_empty(), empty);
2680 assert_eq!(saved.window_title, source.window_title);
2681 assert_eq!(saved.work_state, source.work_state);
2682 assert_eq!(
2683 serde_json::to_value(&saved.metadata.cost)?,
2684 serde_json::to_value(&source.metadata.cost)?
2685 );
2686 assert_eq!(
2687 saved.metadata.parent_session_id.as_deref(),
2688 Some(original.session_id.as_str())
2689 );
2690 assert_eq!(
2691 saved.metadata.runtime_store,
2692 Some(fixture.runtime.session_store_binding())
2693 );
2694 let mut paused = goal.clone();
2695 paused.status = SessionGoalStatus::Paused;
2696 paused.pause_reason = None;
2697 assert_eq!(manager.load_session_goal(id)?, Some(paused));
2698 assert!(
2699 fixture.runtime.get_goal(&fork.id).await?.is_none(),
2700 "local goal continuity does not invent public Runtime goal inheritance"
2701 );
2702 }
2703 assert_eq!(
2704 fs::read(
2705 fixture
2706 .sessions_dir
2707 .join(format!("{}.json", original.session_id))
2708 )?,
2709 source_bytes
2710 );
2711 assert_eq!(manager.load_session_goal(&original.session_id)?, Some(goal));
2712 assert!(fixture.model.captured_requests().is_empty());
2713 Ok(())
2714 }
2715
2716 #[tokio::test]
2717 async fn native_fork_copy_refusal_never_publishes_a_shared_session_thread() -> Result<()> {
2718 let _env = lock_test_env();
2719 let temp = tempfile::tempdir()?;
2720 let root = temp.path().canonicalize()?;
2721 let _home = EnvVarGuard::set("CODEWHALE_HOME", root.join("home"));
2722 let fixture = HistoryOperationFixture::new(&root).await?;
2723 let original = fixture.import_branched_source().await?;
2724 let before = fixture
2725 .runtime
2726 .list_threads(
2727 crate::runtime_threads::ThreadListFilter::IncludeArchived,
2728 None,
2729 )
2730 .await?;
2731 let path = fixture
2732 .sessions_dir
2733 .join(format!("{}.json", original.session_id));
2734 let original_bytes = fs::read(&path)?;
2735 crate::utils::write_atomic(&path, b"{invalid-complete-source")?;
2736 let refused = fixture
2737 .client
2738 .post(fixture.url(&format!("/v1/threads/{}/fork", original.runtime_thread_id)))
2739 .send()
2740 .await?;
2741 assert!(!refused.status().is_success());
2742 let after = fixture
2743 .runtime
2744 .list_threads(
2745 crate::runtime_threads::ThreadListFilter::IncludeArchived,
2746 None,
2747 )
2748 .await?;
2749 assert_eq!(before.len(), after.len());
2750 assert_eq!(
2751 fixture
2752 .runtime
2753 .get_thread(&original.runtime_thread_id)
2754 .await?
2755 .session_id,
2756 Some(original.session_id.clone())
2757 );
2758 assert_eq!(fs::read(&path)?, b"{invalid-complete-source");
2759 crate::utils::write_atomic(&path, &original_bytes)?;
2760 assert!(fixture.model.captured_requests().is_empty());
2761 Ok(())
2762 }
2763
2763 lines RUST