返回 CodeWhale
daemon_socket.rs
根目录 / crates / app-server / tests / daemon_socket.rs
1 //! Desktop Phase 0 acceptance for the daemon socket: spawn the daemon, connect
2 //! over the unix socket, complete the attach/claim handshake, round-trip
3 //! requests through the same JSON-RPC dispatcher the stdio transport uses,
4 //! and shut down cleanly (socket file removed, listener gone).
5
6 #![cfg(unix)]
7
8 use std::os::unix::fs::PermissionsExt;
9 use std::path::{Path, PathBuf};
10 use std::sync::atomic::{AtomicU64, Ordering};
11 use std::time::{Duration, SystemTime, UNIX_EPOCH};
12
13 use codewhale_app_server::daemon_socket::{
14 DaemonSocketError, DaemonSocketOptions, bind_daemon_socket,
15 };
16 use serde_json::{Value, json};
17 use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
18 use tokio::net::UnixStream;
19 use tokio::net::unix::{OwnedReadHalf, OwnedWriteHalf};
20 use tokio::task::JoinHandle;
21
22 static NONCE: AtomicU64 = AtomicU64::new(0);
23
24 /// A short, unique socket path: unix socket paths are capped near 100 bytes,
25 /// so `std::env::temp_dir()` (deep under `/var/folders` on macOS) is too long.
26 /// `/tmp` is the same choice the hooks crate's socket test makes.
27 fn short_socket_root(label: &str) -> PathBuf {
28 let millis = SystemTime::now()
29 .duration_since(UNIX_EPOCH)
30 .expect("clock")
31 .as_millis()
32 % 1_000_000;
33 let nonce = NONCE.fetch_add(1, Ordering::Relaxed);
34 let pid = std::process::id();
35 let root = PathBuf::from("/tmp")
36 .canonicalize()
37 .expect("selected temporary parent")
38 .join(format!("cw-ds-{label}-{pid}-{nonce}-{millis}"));
39 assert!(
40 root.as_os_str().len() < 60,
41 "socket root too long for a unix socket test: {}",
42 root.display()
43 );
44 root
45 }
46
47 struct Harness {
48 root: PathBuf,
49 socket_path: PathBuf,
50 _config_dir: tempfile::TempDir,
51 }
52
53 impl Harness {
54 fn new(label: &str) -> Self {
55 let root = short_socket_root(label);
56 let config_dir = tempfile::tempdir().expect("tempdir");
57 std::fs::write(config_dir.path().join("config.toml"), "").expect("config");
58 Self {
59 socket_path: root.join("run").join("daemon.sock"),
60 root,
61 _config_dir: config_dir,
62 }
63 }
64
65 fn options(&self) -> DaemonSocketOptions {
66 DaemonSocketOptions {
67 socket_path: Some(self.socket_path.clone()),
68 config_path: Some(self._config_dir.path().join("config.toml")),
69 }
70 }
71
72 /// Bind and serve on a background task; returns the serve join handle.
73 async fn spawn_daemon(&self) -> JoinHandle<Result<(), DaemonSocketError>> {
74 let daemon = bind_daemon_socket(self.options())
75 .await
76 .expect("bind daemon socket");
77 assert_eq!(daemon.local_path(), self.socket_path.as_path());
78 tokio::spawn(daemon.serve())
79 }
80 }
81
82 impl Drop for Harness {
83 fn drop(&mut self) {
84 let _ = std::fs::remove_dir_all(&self.root);
85 }
86 }
87
88 struct Client {
89 reader: BufReader<OwnedReadHalf>,
90 writer: OwnedWriteHalf,
91 }
92
93 impl Client {
94 async fn connect(path: &Path) -> Self {
95 let stream = tokio::time::timeout(Duration::from_secs(5), UnixStream::connect(path))
96 .await
97 .expect("connect timeout")
98 .expect("connect");
99 let (rx, writer) = stream.into_split();
100 Self {
101 reader: BufReader::new(rx),
102 writer,
103 }
104 }
105
106 async fn call(&mut self, id: u64, method: &str, params: Value) -> Value {
107 let line = serde_json::to_string(&json!({
108 "jsonrpc": "2.0",
109 "id": id,
110 "method": method,
111 "params": params,
112 }))
113 .expect("encode");
114 self.writer
115 .write_all(format!("{line}\n").as_bytes())
116 .await
117 .expect("write");
118 let mut response = String::new();
119 let read = tokio::time::timeout(
120 Duration::from_secs(10),
121 self.reader.read_line(&mut response),
122 )
123 .await
124 .expect("response timeout")
125 .expect("read");
126 assert!(
127 read > 0,
128 "daemon closed the connection before answering `{method}`"
129 );
130 let value: Value = serde_json::from_str(&response).expect("json response");
131 assert_eq!(value["id"], json!(id), "response id mismatch: {value}");
132 value
133 }
134
135 async fn attach(&mut self, id: u64, name: &str, mode: &str) -> Value {
136 self.call(
137 id,
138 "daemon/attach",
139 json!({ "client": { "name": name, "version": "0.0.0-test", "pid": std::process::id() }, "mode": mode }),
140 )
141 .await
142 }
143
144 /// Read until EOF; proves the daemon closed the socket.
145 async fn wait_for_close(mut self) {
146 let mut sink = String::new();
147 let read = tokio::time::timeout(Duration::from_secs(10), self.reader.read_line(&mut sink))
148 .await
149 .expect("close timeout")
150 .expect("read");
151 assert_eq!(read, 0, "expected EOF, got: {sink}");
152 }
153 }
154
155 async fn wait_for_socket_removed(path: &Path) {
156 tokio::time::timeout(Duration::from_secs(10), async {
157 while path.exists() {
158 tokio::time::sleep(Duration::from_millis(20)).await;
159 }
160 })
161 .await
162 .expect("socket file must be removed on shutdown");
163 }
164
165 #[tokio::test]
166 async fn owner_attaches_round_trips_and_shuts_down_cleanly() {
167 let harness = Harness::new("owner");
168 let server = harness.spawn_daemon().await;
169
170 let socket_mode = std::fs::metadata(&harness.socket_path)
171 .expect("socket metadata")
172 .permissions()
173 .mode()
174 & 0o777;
175 assert_eq!(socket_mode, 0o600, "socket must be private to the user");
176 let dir_mode = std::fs::metadata(harness.socket_path.parent().expect("parent"))
177 .expect("dir metadata")
178 .permissions()
179 .mode()
180 & 0o777;
181 assert_eq!(dir_mode, 0o700, "runtime dir must be private to the user");
182
183 let mut client = Client::connect(&harness.socket_path).await;
184
185 // Anything but healthz before attaching is refused with a typed error:
186 // a read-only probe, a thread/* read, and a prompt run alike.
187 for (id, method, params) in [
188 (1, "capabilities", json!({})),
189 (10, "thread/list", json!({})),
190 (11, "prompt/run", json!({ "prompt": "hi" })),
191 ] {
192 let early = client.call(id, method, params).await;
193 assert_eq!(early["error"]["code"], json!(-32010), "{method}: {early}");
194 assert_eq!(early["error"]["data"]["error"], json!("attach_required"));
195 assert_eq!(early["error"]["data"]["method"], json!(method));
196 }
197
198 // healthz is allowed pre-attach so a shell can probe liveness first.
199 let health = client.call(2, "healthz", json!({})).await;
200 assert_eq!(health["result"]["status"], json!("ok"), "{health}");
201 assert_eq!(health["result"]["transport"], json!("unix-socket"));
202
203 let attached = client.attach(3, "codewhale-desktop", "claim").await;
204 assert_eq!(attached["result"]["attached"], json!(true), "{attached}");
205 assert_eq!(attached["result"]["role"], json!("owner"));
206 assert_eq!(attached["result"]["transport"], json!("unix-socket"));
207 assert_eq!(
208 attached["result"]["daemon"]["pid"],
209 json!(std::process::id())
210 );
211 assert_eq!(
212 attached["result"]["daemon"]["version"],
213 json!(env!("CARGO_PKG_VERSION"))
214 );
215 assert_eq!(
216 attached["result"]["owner"]["name"],
217 json!("codewhale-desktop")
218 );
219 assert_eq!(attached["result"]["connections"], json!(1));
220
221 // Post-attach, the socket transport advertises its own handshake next to
222 // the stdio method set.
223 let advertised = client.call(8, "capabilities", json!({})).await;
224 let methods = advertised["result"]["methods"]
225 .as_array()
226 .expect("methods array");
227 assert_eq!(methods[0], json!("healthz"), "{advertised}");
228 assert_eq!(methods[1], json!("daemon/attach"), "{advertised}");
229 assert!(methods.contains(&json!("shutdown")));
230
231 // Round-trip JSON-RPC requests through the shared dispatcher: `app/*`
232 // methods in, their JSON results out — byte-for-byte the shapes the
233 // stdio transport emits. (No protocol-crate Op/EventMsg envelope is on
234 // this wire; the framing is the stdio transport's newline-delimited
235 // JSON-RPC.)
236 let caps = client.call(4, "app/capabilities", json!({})).await;
237 assert_eq!(caps["result"]["ok"], json!(true), "{caps}");
238 assert!(caps["result"]["data"]["routes"].is_array());
239 let config = client
240 .call(5, "app/config/get", json!({ "key": "model" }))
241 .await;
242 assert_eq!(config["result"]["ok"], json!(true), "{config}");
243 assert_eq!(config["result"]["data"]["key"], json!("model"));
244
245 // A second attach on an attached connection is a typed refusal, not
246 // method_not_found.
247 let again = client.attach(6, "codewhale-desktop", "attach").await;
248 assert_eq!(again["error"]["code"], json!(-32014), "{again}");
249
250 let stopped = client.call(7, "shutdown", json!({})).await;
251 assert_eq!(stopped["result"]["status"], json!("stopped"), "{stopped}");
252
253 let outcome = tokio::time::timeout(Duration::from_secs(10), server)
254 .await
255 .expect("daemon must exit after the owner's shutdown")
256 .expect("join");
257 outcome.expect("serve result");
258 wait_for_socket_removed(&harness.socket_path).await;
259 client.wait_for_close().await;
260 }
261
262 #[tokio::test]
263 async fn guests_share_the_daemon_but_cannot_stop_it() {
264 let harness = Harness::new("guest");
265 let server = harness.spawn_daemon().await;
266
267 let mut owner = Client::connect(&harness.socket_path).await;
268 let claimed = owner.attach(1, "desktop-window-1", "claim").await;
269 assert_eq!(claimed["result"]["role"], json!("owner"), "{claimed}");
270
271 let mut guest = Client::connect(&harness.socket_path).await;
272 let lost = guest.attach(1, "desktop-window-2", "claim").await;
273 assert_eq!(lost["error"]["code"], json!(-32011), "{lost}");
274 assert_eq!(
275 lost["error"]["data"]["owner"]["name"],
276 json!("desktop-window-1")
277 );
278
279 let attached = guest.attach(2, "desktop-window-2", "attach").await;
280 assert_eq!(attached["result"]["role"], json!("attached"), "{attached}");
281 assert_eq!(
282 attached["result"]["owner"]["name"],
283 json!("desktop-window-1")
284 );
285 assert_eq!(attached["result"]["connections"], json!(2));
286
287 let health = guest.call(3, "healthz", json!({})).await;
288 assert_eq!(health["result"]["status"], json!("ok"));
289
290 let refused = guest.call(4, "shutdown", json!({})).await;
291 assert_eq!(refused["error"]["code"], json!(-32012), "{refused}");
292 assert_eq!(refused["error"]["data"]["error"], json!("not_daemon_owner"));
293 assert!(
294 !server.is_finished(),
295 "a guest's shutdown must not stop the daemon"
296 );
297 assert!(harness.socket_path.exists());
298
299 // Once the owner leaves, the slot frees and a relaunched shell can claim.
300 drop(owner);
301 let mut relaunched = Client::connect(&harness.socket_path).await;
302 let reclaimed = tokio::time::timeout(Duration::from_secs(10), async {
303 loop {
304 let response = relaunched.attach(1, "desktop-relaunch", "claim").await;
305 if response.get("result").is_some() {
306 return response;
307 }
308 tokio::time::sleep(Duration::from_millis(20)).await;
309 }
310 })
311 .await
312 .expect("owner slot must free when the owner disconnects");
313 assert_eq!(reclaimed["result"]["role"], json!("owner"), "{reclaimed}");
314
315 // The guest is still attached and served while the new owner is in.
316 let health = guest.call(5, "healthz", json!({})).await;
317 assert_eq!(health["result"]["status"], json!("ok"));
318
319 let stopped = relaunched.call(2, "shutdown", json!({})).await;
320 assert_eq!(stopped["result"]["status"], json!("stopped"));
321 tokio::time::timeout(Duration::from_secs(10), server)
322 .await
323 .expect("daemon exits")
324 .expect("join")
325 .expect("serve result");
326 wait_for_socket_removed(&harness.socket_path).await;
327 // The owner's shutdown closes every other connection, not just its own.
328 guest.wait_for_close().await;
329 relaunched.wait_for_close().await;
330 }
331
332 #[tokio::test]
333 async fn version_skew_is_refused_at_attach() {
334 let harness = Harness::new("skew");
335 let server = harness.spawn_daemon().await;
336 let mut client = Client::connect(&harness.socket_path).await;
337 let refused = client
338 .call(
339 1,
340 "daemon/attach",
341 json!({ "client": { "name": "old-desktop" }, "expect_daemon_version": "0.0.1-other" }),
342 )
343 .await;
344 assert_eq!(refused["error"]["code"], json!(-32013), "{refused}");
345 assert_eq!(
346 refused["error"]["data"]["actual"],
347 json!(env!("CARGO_PKG_VERSION"))
348 );
349
350 let daemon = bind_daemon_socket(harness.options()).await;
351 // Meanwhile the original daemon is live, so a second bind must refuse.
352 match daemon {
353 Err(DaemonSocketError::AlreadyRunning { path }) => {
354 assert_eq!(path, harness.socket_path);
355 }
356 Err(other) => panic!("unexpected error: {other}"),
357 Ok(_) => panic!("second daemon must not replace a live socket"),
358 }
359 server.abort();
360 let _ = server.await;
361 }
362
363 #[tokio::test]
364 async fn stale_socket_is_cleaned_up_and_foreign_files_are_refused() {
365 let harness = Harness::new("stale");
366 std::fs::create_dir_all(harness.socket_path.parent().expect("parent")).expect("mkdir");
367 std::fs::set_permissions(
368 harness.socket_path.parent().unwrap(),
369 std::fs::Permissions::from_mode(0o700),
370 )
371 .expect("private stale fixture parent");
372
373 // A socket file whose listener is gone: bind must reclaim it.
374 {
375 let dead = tokio::net::UnixListener::bind(&harness.socket_path).expect("bind dead");
376 drop(dead);
377 }
378 assert!(
379 harness.socket_path.exists(),
380 "dropping a listener leaves the file"
381 );
382 let daemon = bind_daemon_socket(harness.options())
383 .await
384 .expect("stale socket must be reclaimed");
385 let handle = daemon.shutdown_handle();
386 let server = tokio::spawn(daemon.serve());
387 let mut client = Client::connect(&harness.socket_path).await;
388 let health = client.call(1, "healthz", json!({})).await;
389 assert_eq!(health["result"]["status"], json!("ok"));
390 handle.trigger();
391 tokio::time::timeout(Duration::from_secs(10), server)
392 .await
393 .expect("daemon exits on handle")
394 .expect("join")
395 .expect("serve result");
396 wait_for_socket_removed(&harness.socket_path).await;
397
398 // A regular file at the path is never deleted.
399 std::fs::write(&harness.socket_path, b"not a socket").expect("write file");
400 match bind_daemon_socket(harness.options()).await {
401 Err(DaemonSocketError::NotASocket { path }) => assert_eq!(path, harness.socket_path),
402 Err(other) => panic!("unexpected error: {other}"),
403 Ok(_) => panic!("must refuse to replace a non-socket"),
404 }
405 assert_eq!(
406 std::fs::read(&harness.socket_path).expect("file intact"),
407 b"not a socket"
408 );
409 }
410
411 #[tokio::test]
412 async fn binding_refuses_public_parent_and_drop_before_serve_retires_exact_socket() {
413 let harness = Harness::new("private");
414 let parent = harness.socket_path.parent().unwrap();
415 std::fs::create_dir_all(parent).unwrap();
416 std::fs::set_permissions(parent, std::fs::Permissions::from_mode(0o755)).unwrap();
417 assert!(bind_daemon_socket(harness.options()).await.is_err());
418 assert_eq!(
419 std::fs::metadata(parent).unwrap().permissions().mode() & 0o777,
420 0o755
421 );
422 assert!(!harness.socket_path.exists());
423 std::fs::set_permissions(parent, std::fs::Permissions::from_mode(0o700)).unwrap();
424 let daemon = bind_daemon_socket(harness.options()).await.unwrap();
425 assert!(harness.socket_path.exists());
426 drop(daemon);
427 tokio::time::timeout(Duration::from_secs(5), async {
428 while harness.socket_path.exists() {
429 tokio::task::yield_now().await;
430 }
431 })
432 .await
433 .expect("cancelled binder retirement finishes off-runtime");
434 assert!(!harness.socket_path.exists());
435 }
436
437 #[tokio::test]
438 async fn dropping_bound_daemon_preserves_a_replacement_at_the_selected_path() {
439 let harness = Harness::new("replaced");
440 let daemon = bind_daemon_socket(harness.options()).await.unwrap();
441 let captured = harness.socket_path.with_extension("captured");
442 std::fs::rename(&harness.socket_path, &captured).unwrap();
443 std::fs::write(&harness.socket_path, b"operator replacement").unwrap();
444 drop(daemon);
445 assert_eq!(
446 std::fs::read(&harness.socket_path).unwrap(),
447 b"operator replacement"
448 );
449 assert!(captured.exists());
450 }
451
452 #[tokio::test]
453 async fn captured_owner_frontend_authenticates_generation_without_publishing_bearer() {
454 let harness = Harness::new("owner");
455 let owner = codewhale_protocol::RuntimeOwnerReceipt {
456 version: 1,
457 data_dir: harness.root.join("runtime"),
458 execution_scope: "captured-test-store".into(),
459 lease_generation: "captured-test-generation".into(),
460 pid: std::process::id(),
461 process_start: codewhale_app_server::daemon_socket::capture_process_start(
462 std::process::id(),
463 )
464 .await
465 .unwrap(),
466 principal: codewhale_config::private_directory::PrivateDirectory::current_user_id()
467 .to_string(),
468 socket_path: harness.socket_path.clone(),
469 config_path: harness.options().config_path,
470 };
471 // Transport acceptance: the captured owner is a fixture, not an Engine/store proof.
472 let daemon = codewhale_app_server::bind_runtime_owner(
473 owner.config_path.clone(),
474 "127.0.0.1:1".parse().unwrap(),
475 Some("private-fixture-bearer".into()),
476 owner.clone(),
477 )
478 .await
479 .unwrap();
480 let receipt_path = harness.socket_path.with_file_name("daemon.sock.owner.json");
481 let bytes = std::fs::read(&receipt_path).unwrap();
482 assert_eq!(
483 serde_json::from_slice::<codewhale_protocol::RuntimeOwnerReceipt>(&bytes).unwrap(),
484 owner
485 );
486 assert!(
487 !String::from_utf8(bytes)
488 .unwrap()
489 .contains("private-fixture-bearer")
490 );
491 assert_eq!(
492 std::fs::metadata(&receipt_path)
493 .unwrap()
494 .permissions()
495 .mode()
496 & 0o777,
497 0o600
498 );
499 let handle = daemon.shutdown_handle();
500 let server = tokio::spawn(daemon.serve());
501 let mut guest = Client::connect(&harness.socket_path).await;
502 let missing = guest.attach(1, "guest", "attach").await;
503 assert!(missing["error"].is_object());
504 let mut stale = owner.clone();
505 stale.lease_generation = "other-generation".into();
506 let rejected = guest
507 .call(
508 2,
509 "daemon/attach",
510 json!({"client":{"name":"guest"},"mode":"attach","expect_owner":stale}),
511 )
512 .await;
513 assert!(rejected["error"].is_object());
514 let claim = guest
515 .call(
516 3,
517 "daemon/attach",
518 json!({"client":{"name":"guest"},"mode":"claim","expect_owner":owner}),
519 )
520 .await;
521 assert!(claim["error"].is_object());
522 let attached = guest
523 .call(
524 4,
525 "daemon/attach",
526 json!({"client":{"name":"guest","pid":1},"mode":"attach","expect_owner":owner}),
527 )
528 .await;
529 assert_eq!(attached["result"]["role"], "attached");
530 assert_eq!(attached["result"]["owner_receipt"], json!(owner));
531 // The display PID is deliberately false: authorization uses kernel peer credentials.
532 let denied = guest.call(5, "shutdown", json!({})).await;
533 assert!(denied["error"].is_object());
534 assert!(guest.call(6, "healthz", json!({})).await["result"].is_object());
535 handle.trigger();
536 tokio::time::timeout(Duration::from_secs(5), server)
537 .await
538 .unwrap()
539 .unwrap()
540 .unwrap();
541 assert!(!harness.socket_path.exists());
542 assert!(
543 !receipt_path.exists(),
544 "normal shutdown awaits exact receipt retirement"
545 );
546 }
547
548 /// Native IPC/HTTP transport proof with a fake owned Runtime; no Engine or
549 /// provider acceptance is inferred from this source fixture.
550 #[tokio::test]
551 async fn canonical_cli_control_uses_authenticated_owner_and_same_dispatcher() {
552 use axum::{Json, Router, routing::get};
553 let harness = Harness::new("thread-control");
554 let workspace = harness._config_dir.path().to_path_buf();
555 let thread = serde_json::json!({"id":"canonical-transport","created_at":"2026-10-02T00:00:00Z",
556 "updated_at":"2026-10-02T00:00:00Z","model":"fixture-model","model_provider":"custom",
557 "model_provider_id":"fixture-owner","workspace":workspace,"archived":false});
558 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
559 let endpoint = listener.local_addr().unwrap();
560 let app = Router::new()
561 .route(
562 "/v1/threads",
563 get(move || {
564 let thread = thread.clone();
565 async move { Json(serde_json::json!([thread])) }
566 }),
567 )
568 .route(
569 "/v1/threads/running",
570 get(|| async { Json(serde_json::json!([])) }),
571 );
572 let fake_http = tokio::spawn(async move {
573 axum::serve(listener, app).await.unwrap();
574 });
575 let owner = codewhale_protocol::RuntimeOwnerReceipt {
576 version: 1,
577 data_dir: harness.root.join("runtime"),
578 execution_scope: "control-fixture".into(),
579 lease_generation: "control-generation".into(),
580 pid: std::process::id(),
581 process_start: codewhale_app_server::daemon_socket::capture_process_start(
582 std::process::id(),
583 )
584 .await
585 .unwrap(),
586 principal: codewhale_config::private_directory::PrivateDirectory::current_user_id()
587 .to_string(),
588 socket_path: harness.socket_path.clone(),
589 config_path: harness.options().config_path,
590 };
591 let daemon = codewhale_app_server::bind_runtime_owner(
592 owner.config_path.clone(),
593 endpoint,
594 Some("private-fixture-control".into()),
595 owner.clone(),
596 )
597 .await
598 .unwrap();
599 let shutdown = daemon.shutdown_handle();
600 let socket = tokio::spawn(daemon.serve());
601 let result = codewhale_app_server::request_thread_control(
602 owner.config_path.clone(),
603 Some(harness.socket_path.clone()),
604 None,
605 codewhale_protocol::ThreadRequest::List(codewhale_protocol::ThreadListParams {
606 include_archived: false,
607 limit: None,
608 }),
609 )
610 .await
611 .unwrap();
612 assert_eq!(result.threads.len(), 1);
613 assert_eq!(result.threads[0].id, "canonical-transport");
614 assert_eq!(result.threads[0].model_provider, "fixture-owner");
615 assert!(
616 !serde_json::to_string(&result)
617 .unwrap()
618 .contains("private-fixture-control")
619 );
620 let error = codewhale_app_server::request_thread_control(
621 owner.config_path.clone(),
622 Some(harness.socket_path.clone()),
623 Some(codewhale_app_server::ThreadControlSelection {
624 workspace: Some(workspace),
625 config_profile: None,
626 config_source: None,
627 }),
628 codewhale_protocol::ThreadRequest::List(codewhale_protocol::ThreadListParams {
629 include_archived: false,
630 limit: None,
631 }),
632 )
633 .await
634 .unwrap_err();
635 assert!(format!("{error:#}").contains("no captured worker setting"));
636 shutdown.trigger();
637 tokio::time::timeout(Duration::from_secs(5), socket)
638 .await
639 .unwrap()
640 .unwrap()
641 .unwrap();
642 fake_http.abort();
643 }
644
645 /// This only validates the adapter's captured scope transport. Actual Runtime
646 /// policy/Config admission is covered by the owning Runtime acceptance suite.
647 struct ScopedControlFixture {
648 workers: usize,
649 config: PathBuf,
650 scopes: std::sync::Arc<tokio::sync::Mutex<Vec<codewhale_app_server::RuntimeFrontendScope>>>,
651 }
652 impl codewhale_app_server::RuntimeOwnerFrontend for ScopedControlFixture {
653 fn validate_selection<'a>(
654 &'a self,
655 selection: &'a codewhale_app_server::RuntimeOwnerFrontendSelection,
656 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = anyhow::Result<()>> + Send + 'a>> {
657 Box::pin(async move {
658 let codewhale_app_server::RuntimeOwnerFrontendSelection::Control(scope) = selection
659 else {
660 anyhow::bail!("fixture expects the existing control frontend");
661 };
662 anyhow::ensure!(
663 scope.workers == self.workers,
664 "captured scheduler setting changed"
665 );
666 anyhow::ensure!(
667 scope.workspace.is_absolute() && scope.workspace.is_dir(),
668 "fixture scope is not selected"
669 );
670 anyhow::ensure!(
671 scope
672 .config_profile
673 .as_deref()
674 .is_none_or(|profile| profile == "reviewed"),
675 "profile is not admitted by held owner"
676 );
677 anyhow::ensure!(
678 scope
679 .config_source
680 .as_ref()
681 .is_none_or(|source| source == &self.config),
682 "config is not admitted by held owner"
683 );
684 self.scopes.lock().await.push(scope.clone());
685 Ok(())
686 })
687 }
688 fn serve(
689 &self,
690 selection: codewhale_app_server::RuntimeOwnerFrontendSelection,
691 compatibility: codewhale_app_server::AppState,
692 input: Box<dyn tokio::io::AsyncBufRead + Send + Unpin>,
693 output: Box<dyn tokio::io::AsyncWrite + Send + Unpin>,
694 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = anyhow::Result<()>> + Send + '_>> {
695 Box::pin(async move {
696 let codewhale_app_server::RuntimeOwnerFrontendSelection::Control(scope) = selection
697 else {
698 anyhow::bail!("fixture expects control");
699 };
700 codewhale_app_server::run_guest_control(compatibility, scope.workspace, input, output)
701 .await
702 })
703 }
704 }
705
706 #[tokio::test]
707 async fn canonical_cli_scoped_resume_fork_keep_exact_owner_workers_and_refuse_wrong_selection() {
708 use axum::{Json, Router, response::IntoResponse as _};
709 use std::sync::Arc;
710 let harness = Harness::new("scoped-control");
711 let workspace = harness._config_dir.path().to_path_buf();
712 let selected = workspace.join("explicit-target");
713 std::fs::create_dir(&selected).unwrap();
714 let owner = codewhale_protocol::RuntimeOwnerReceipt {
715 version: 1,
716 data_dir: harness.root.join("runtime"),
717 execution_scope: "scoped-control-store".into(),
718 lease_generation: "scoped-control-generation".into(),
719 pid: std::process::id(),
720 process_start: codewhale_app_server::daemon_socket::capture_process_start(
721 std::process::id(),
722 )
723 .await
724 .unwrap(),
725 principal: codewhale_config::private_directory::PrivateDirectory::current_user_id()
726 .to_string(),
727 socket_path: harness.socket_path.clone(),
728 config_path: harness.options().config_path,
729 };
730 let calls = Arc::new(tokio::sync::Mutex::new(Vec::<(String, Value)>::new()));
731 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
732 let endpoint = listener.local_addr().unwrap();
733 let captured = owner.clone();
734 let http_calls = calls.clone();
735 let target = selected.clone();
736 let app = Router::new().fallback(move |request: axum::extract::Request| {
737 let owner = captured.clone(); let calls = http_calls.clone(); let target = target.clone();
738 async move {
739 let path = request.uri().path().to_string();
740 let bytes = axum::body::to_bytes(request.into_body(), 8 * 1024 * 1024).await.unwrap();
741 let body = if bytes.is_empty() { Value::Null } else { serde_json::from_slice(&bytes).unwrap() };
742 calls.lock().await.push((path.clone(), body.clone()));
743 let record = |id: &str, root: &Path| json!({"id":id,"created_at":"2026-10-02T00:00:00Z",
744 "updated_at":"2026-10-02T00:00:00Z","model":"fixture-model","model_provider":"custom",
745 "model_provider_id":"fixture-owner","workspace":root,"archived":false});
746 let value = if path.ends_with("/operations/lookup") { json!({"state":"absent"}) }
747 else if path.ends_with("/mutate") {
748 assert_eq!(body["expected_data_dir"], json!(owner.data_dir));
749 assert_eq!(body["expected_execution_scope"], owner.execution_scope);
750 assert_eq!(body["workspace"], json!(target));
751 assert_eq!(body["mutation"]["options"]["overrides"], json!({"model":"explicit-model", "model_provider":"fixture-owner", "approval_policy":"on-request", "sandbox":"workspace-write"}));
752 let fork = body["mutation"]["action"] == "fork";
753 json!({"version":1,"data_dir":owner.data_dir,"execution_scope":owner.execution_scope,
754 "operation_key":body["operation_key"],"request_digest":"1".repeat(64),"history_digest":"2".repeat(64),
755 "runtime_thread_id":if fork {"canonical-scoped-fork"} else {"canonical-source"},
756 "session_id":if fork {"session-scoped-fork"} else {"session-source"}})
757 } else if path.ends_with("/history") {
758 json!({"version":1,"data_dir":owner.data_dir,"execution_scope":owner.execution_scope,
759 "runtime_thread_id":"canonical-source","saved_session_id":"session-source",
760 "saved_document_digest":"a".repeat(64),"session_goal_digest":"74234e98afe7498fb5daf1f36ac2d78acc339464f950703b8c019892f982b90b",
761 "document_digest":"b".repeat(64),"session":{"metadata":{"id":"session-source"},"messages":[],"journal":{"entries":[]}}})
762 } else if path == "/v1/threads/running" { json!([]) }
763 else if path == "/v1/threads" { json!([record("canonical-source", &target), record("other-workspace", Path::new("/recorded-other-workspace"))]) }
764 else if path == "/v1/threads/canonical-scoped-fork" { record("canonical-scoped-fork", &target) }
765 else if path == "/v1/threads/canonical-source" { record("canonical-source", &target) }
766 else { return (axum::http::StatusCode::NOT_FOUND, Json(json!({"error":"fixture route missing"}))).into_response(); };
767 Json(value).into_response()
768 }
769 });
770 let fake_http = tokio::spawn(async move {
771 axum::serve(listener, app).await.unwrap();
772 });
773 let scopes = Arc::new(tokio::sync::Mutex::new(Vec::new()));
774 let frontend = Arc::new(ScopedControlFixture {
775 workers: 7,
776 config: owner.config_path.clone().unwrap(),
777 scopes: scopes.clone(),
778 });
779 let (daemon, _) = codewhale_app_server::bind_runtime_frontends(
780 owner.config_path.clone(),
781 Some("fixture-private-scope-token".into()),
782 owner.clone(),
783 codewhale_app_server::RuntimeOwnerRouting {
784 endpoint,
785 workspace: Some(workspace),
786 workers: Some(7),
787 mobile: false,
788 web: false,
789 acp: true,
790 acp_only: false,
791 },
792 Some(frontend),
793 )
794 .await
795 .unwrap();
796 let shutdown = daemon.shutdown_handle();
797 let socket = tokio::spawn(daemon.serve());
798 let selection = codewhale_app_server::ThreadControlSelection {
799 workspace: Some(selected.clone()),
800 config_profile: Some("reviewed".into()),
801 config_source: owner.config_path.clone(),
802 };
803 for (kind, operation, expected_id) in [
804 ("resume", "scoped-resume", "canonical-source"),
805 ("fork", "scoped-fork", "canonical-scoped-fork"),
806 ] {
807 let request = serde_json::from_value(
808 json!({"kind":kind,"thread_id":"canonical-source","operation_key":operation,
809 "model":"explicit-model", "model_provider":"fixture-owner", "approval_policy":"on-request", "sandbox":"workspace-write"}),
810 )
811 .unwrap();
812 let response = codewhale_app_server::request_thread_control(
813 owner.config_path.clone(),
814 Some(harness.socket_path.clone()),
815 Some(selection.clone()),
816 request,
817 )
818 .await
819 .unwrap();
820 assert_eq!(response.data["receipt"]["operation_key"], operation);
821 assert_eq!(response.data["receipt"]["runtime_thread_id"], expected_id);
822 assert_eq!(response.thread.unwrap().id, expected_id);
823 }
824 let list = codewhale_app_server::request_thread_control(
825 owner.config_path.clone(),
826 Some(harness.socket_path.clone()),
827 Some(selection.clone()),
828 codewhale_protocol::ThreadRequest::List(codewhale_protocol::ThreadListParams {
829 include_archived: false,
830 limit: None,
831 }),
832 )
833 .await
834 .unwrap();
835 assert!(
836 list.threads
837 .iter()
838 .any(|thread| thread.id == "other-workspace"),
839 "scoped execution does not filter owner-store-wide read-only list"
840 );
841 assert!(scopes.lock().await.iter().all(|scope| scope.workers == 7
842 && scope.workspace == selected
843 && scope.config_profile.as_deref() == Some("reviewed")));
844 let before = calls.lock().await.len();
845 for wrong in [
846 codewhale_app_server::ThreadControlSelection {
847 config_profile: Some("wrong".into()),
848 ..selection.clone()
849 },
850 codewhale_app_server::ThreadControlSelection {
851 config_source: Some(selected.join("other.toml")),
852 ..selection.clone()
853 },
854 ] {
855 let request = serde_json::from_value(json!({"kind":"resume","thread_id":"canonical-source","operation_key":"never-dispatched"})).unwrap();
856 assert!(
857 codewhale_app_server::request_thread_control(
858 owner.config_path.clone(),
859 Some(harness.socket_path.clone()),
860 Some(wrong),
861 request
862 )
863 .await
864 .is_err()
865 );
866 }
867 assert_eq!(
868 calls.lock().await.len(),
869 before,
870 "selection refusal happens before HTTP effect dispatch"
871 );
872 shutdown.trigger();
873 tokio::time::timeout(Duration::from_secs(5), socket)
874 .await
875 .unwrap()
876 .unwrap()
877 .unwrap();
878 fake_http.abort();
879 }
880
880 lines RUST