返回 CodeWhale
lib.rs
根目录 / crates / app-server / src / lib.rs
1 use codewhale_protocol::runtime::{MAX_RUNTIME_IMAGE_BODY_BYTES, RuntimeImageInput};
2 use std::collections::{HashMap, VecDeque};
3 use std::net::SocketAddr;
4 use std::path::{Path, PathBuf};
5 #[cfg(test)]
6 use std::process::{Child, Command, Stdio};
7 use std::sync::Arc;
8 use std::time::Duration;
9
10 use anyhow::{Context, Result, anyhow, bail};
11 use axum::extract::{DefaultBodyLimit, Request, State};
12 use axum::http::{HeaderValue, Method, StatusCode, header};
13 use axum::middleware::{self, Next};
14 use axum::response::{IntoResponse, Response};
15 use axum::routing::{get, post};
16 use axum::{Json, Router};
17 use codewhale_agent::ModelRegistry;
18 use codewhale_config::ConfigStore;
19 use codewhale_core::Runtime;
20 use codewhale_hooks::{HookDispatcher, JsonlHookSink, StdoutHookSink, UnixSocketHookSink};
21 use codewhale_protocol::{
22 AppRequest, AppResponse, EventFrame, PromptRequest, PromptResponse, ResponseChannel,
23 ThreadGoalClearParams, ThreadGoalGetParams, ThreadGoalSetParams,
24 };
25 use codewhale_state::StateStore;
26 use serde::Deserialize;
27 use serde::de::DeserializeOwned;
28 use serde_json::{Value, json};
29 use tokio::io::{AsyncBufRead, AsyncBufReadExt, AsyncWrite, AsyncWriteExt, BufReader};
30 use tokio::sync::{Mutex, RwLock};
31 use tower_http::cors::CorsLayer;
32 use uuid::Uuid;
33
34 mod chat_completions;
35 mod thread_control;
36 pub use codewhale_protocol::{
37 ThreadListParams, ThreadReadParams, ThreadRequest, ThreadResponse, ThreadSetNameParams,
38 };
39 pub use thread_control::{ThreadControlSelection, request_thread_control};
40
41 /// Capture once in a client before sending a durable control. Servers require
42 /// the caller's retained key and never manufacture a replacement on retry.
43 pub fn capture_thread_operation_key() -> String {
44 Uuid::new_v4().to_string()
45 }
46 pub mod daemon_client;
47 pub mod daemon_socket;
48 #[cfg(windows)]
49 mod daemon_windows;
50
51 /// Legacy DeepSeek-era naming kept for external compatibility.
52 ///
53 /// CodeWhale began life as DeepSeek-TUI; existing health probes, SDK
54 /// harnesses, and on-disk layouts still key off these names. Every remaining
55 /// legacy reference in this crate routes through this shim so a future
56 /// coordinated migration touches exactly one place (repo policy: preserve
57 /// legacy migration care).
58 mod legacy_deepseek_compat {
59 use std::path::PathBuf;
60
61 /// Service name advertised by the HTTP and stdio health probes.
62 pub(crate) const SERVICE_NAME: &str = "deepseek-app-server";
63
64 /// Fallback hook-event log location used when no config path is
65 /// provided (legacy `.deepseek/` dot-directory layout).
66 pub(crate) fn default_events_log_path() -> PathBuf {
67 PathBuf::from(".deepseek/events.jsonl")
68 }
69 }
70
71 /// Upper bound on JSON request bodies accepted by the HTTP app-server.
72 const MAX_HTTP_BODY_BYTES: usize = 16 * 1024 * 1024;
73 const MAX_SSE_FRAME_BYTES: usize = 16 * 1024 * 1024;
74
75 const DEFAULT_CORS_ORIGINS: &[&str] = &[
76 "http://localhost",
77 "http://localhost:1420",
78 "http://localhost:3000",
79 "http://localhost:5173",
80 "http://127.0.0.1",
81 "http://127.0.0.1:1420",
82 "tauri://localhost",
83 ];
84
85 #[derive(Clone)]
86 pub struct AppServerOptions {
87 pub listen: SocketAddr,
88 pub config_path: Option<PathBuf>,
89 pub auth_token: Option<String>,
90 pub insecure_no_auth: bool,
91 pub cors_origins: Vec<String>,
92 }
93
94 impl std::fmt::Debug for AppServerOptions {
95 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
96 f.debug_struct("AppServerOptions")
97 .field("listen", &self.listen)
98 .field("config_path", &self.config_path)
99 .field(
100 "auth_token",
101 &self.auth_token.as_ref().map(|_| "<redacted>"),
102 )
103 .field("insecure_no_auth", &self.insecure_no_auth)
104 .field("cors_origins", &self.cors_origins)
105 .finish()
106 }
107 }
108
109 /// Selected frontend facts passed by the canonical CLI in memory.
110 #[derive(Debug, Clone, PartialEq, Eq)]
111 pub enum RuntimeControlFrontend {
112 Stdio,
113 Socket { path: Option<PathBuf> },
114 LegacyHttp,
115 Acp { model: String },
116 }
117
118 /// A bounded listener request admitted only over the authenticated owner channel.
119 /// The bearer is explicit operator input, never discovery or response data.
120 #[derive(Clone, serde::Serialize, serde::Deserialize)]
121 #[serde(deny_unknown_fields)]
122 pub struct RuntimeListenerSelection {
123 pub workers: usize,
124 pub workspace: PathBuf,
125 pub config_profile: Option<String>,
126 #[serde(default)]
127 pub config_source: Option<PathBuf>,
128 pub host: String,
129 pub port: u16,
130 pub cors_origins: Vec<String>,
131 pub auth_token: Option<String>,
132 pub insecure_no_auth: bool,
133 pub mobile: bool,
134 pub web: bool,
135 }
136 impl std::fmt::Debug for RuntimeListenerSelection {
137 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
138 f.debug_struct("RuntimeListenerSelection")
139 .field("host", &self.host)
140 .field("port", &self.port)
141 .field(
142 "auth_token",
143 &self.auth_token.as_ref().map(|_| "<redacted>"),
144 )
145 .field("mobile", &self.mobile)
146 .field("web", &self.web)
147 .finish_non_exhaustive()
148 }
149 }
150 impl RuntimeListenerSelection {
151 pub fn validate_bounds(&self) -> Result<()> {
152 anyhow::ensure!(
153 self.workspace.as_os_str().len() <= 32768
154 && self
155 .config_source
156 .as_ref()
157 .is_none_or(|path| path.as_os_str().len() <= 32768)
158 && self
159 .config_profile
160 .as_ref()
161 .is_none_or(|profile| profile.len() <= 1024)
162 && self.host.len() <= 128
163 && self.cors_origins.len() <= 64
164 && self.cors_origins.iter().map(String::len).sum::<usize>() <= 32768
165 && self
166 .auth_token
167 .as_ref()
168 .is_none_or(|token| token.len() <= 8192),
169 "selected frontend input exceeds its bounds"
170 );
171 Ok(())
172 }
173 }
174 #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
175 #[serde(deny_unknown_fields)]
176 pub struct RuntimeFrontendScope {
177 pub workers: usize,
178 pub workspace: PathBuf,
179 pub config_profile: Option<String>,
180 #[serde(default)]
181 pub config_source: Option<PathBuf>,
182 }
183 impl RuntimeFrontendScope {
184 pub fn validate_bounds(&self) -> Result<()> {
185 anyhow::ensure!(
186 self.workspace.as_os_str().len() <= 32768
187 && self
188 .config_source
189 .as_ref()
190 .is_none_or(|path| path.as_os_str().len() <= 32768)
191 && self
192 .config_profile
193 .as_ref()
194 .is_none_or(|profile| profile.len() <= 1024),
195 "selected frontend scope exceeds its bounds"
196 );
197 Ok(())
198 }
199 }
200 #[derive(Clone)]
201 pub enum RuntimeOwnerFrontendSelection {
202 Control(RuntimeFrontendScope),
203 Acp {
204 scope: Option<RuntimeFrontendScope>,
205 model: Option<String>,
206 },
207 Listener(RuntimeListenerSelection),
208 }
209 /// Captured Runtime-owned IO projection. Implementations use the held manager
210 /// and existing routers; this port owns neither a store nor a turn controller.
211 pub trait RuntimeOwnerFrontend: Send + Sync {
212 fn validate_selection<'a>(
213 &'a self,
214 selection: &'a RuntimeOwnerFrontendSelection,
215 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = Result<()>> + Send + 'a>>;
216 fn serve(
217 &self,
218 selection: RuntimeOwnerFrontendSelection,
219 compatibility: AppState,
220 input: Box<dyn tokio::io::AsyncBufRead + Send + Unpin>,
221 output: Box<dyn AsyncWrite + Send + Unpin>,
222 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = Result<()>> + Send + '_>>;
223 }
224
225 #[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
226 #[serde(deny_unknown_fields)]
227 pub struct RuntimeOwnerRouting {
228 pub endpoint: SocketAddr,
229 /// The owning manager's acknowledged workspace. A legacy cold wrapper
230 /// without one cannot supply a default mutation scope.
231 #[serde(default, skip_serializing_if = "Option::is_none")]
232 pub workspace: Option<PathBuf>,
233 /// Actual held scheduler setting; absent on historical cold wrappers.
234 #[serde(default, skip_serializing_if = "Option::is_none")]
235 pub workers: Option<usize>,
236 pub mobile: bool,
237 pub web: bool,
238 pub acp: bool,
239 #[serde(default)]
240 pub acp_only: bool,
241 }
242
243 /// Cached app-server→runtime bridge handle.
244 ///
245 /// The outer [`AppState::runtime_bridge`] mutex guards only the cache slot;
246 /// this inner mutex serializes traffic and event cursors for the captured
247 /// canonical owner. It never creates or replaces an execution process.
248 type SharedRuntimeBridge = Arc<Mutex<RuntimeBridge>>;
249
250 #[derive(Clone)]
251 pub struct AppState {
252 captured_owner: Option<codewhale_protocol::RuntimeOwnerReceipt>,
253 captured_routing: Option<RuntimeOwnerRouting>,
254 frontend_workspace: Option<PathBuf>,
255 owner_frontend: Option<Arc<dyn RuntimeOwnerFrontend>>,
256 config_path: Option<PathBuf>,
257 config: Arc<RwLock<codewhale_config::ConfigToml>>,
258 /// Bookkeeping config/jobs and read-only historical archive access.
259 /// Actual turns and transcript writes belong to the captured owner.
260 runtime: Arc<RwLock<Runtime>>,
261 registry: ModelRegistry,
262 auth_token: Option<String>,
263 /// Cached bridge to the real runtime API. Shared by every surface that
264 /// executes a turn — stdio `thread/message`, HTTP `/thread` messages, and
265 /// both `/prompt` transports — because there is exactly one turn engine.
266 runtime_bridge: Arc<Mutex<Option<SharedRuntimeBridge>>>,
267 /// Client-facing thread key → durable runtime thread id.
268 ///
269 /// Runtime threads are persisted by the captured owner's store. Durable
270 /// aliases keep their exact target across frontend detach/reconnect;
271 /// withdrawing a bridge never remints a replacement thread.
272 /// Callers serialize traffic on the bridge mutex.
273 runtime_thread_map: Arc<Mutex<HashMap<String, String>>>,
274 stdio_thread_hints: Arc<Mutex<HashMap<String, RuntimeThreadHint>>>,
275 /// Turns currently streaming over stdio, keyed by stdio thread id.
276 ///
277 /// Deliberately kept *outside* the bridge mutex: a streaming turn holds
278 /// that mutex for its entire duration, so anything reachable only through
279 /// it cannot be used to stop the turn. This holds its own copy of what an
280 /// interrupt needs, so a cancel never waits on the turn it is cancelling.
281 in_flight_turns: Arc<Mutex<HashMap<String, InFlightTurn>>>,
282 }
283
284 /// Everything needed to interrupt a running turn without the bridge lock.
285 #[derive(Debug, Clone)]
286 struct InFlightTurn {
287 base_url: String,
288 auth_token: Option<String>,
289 /// Thread id as the *runtime* knows it, not the stdio-facing id.
290 runtime_thread_id: String,
291 turn_id: String,
292 }
293
294 type TurnRegistry = Arc<Mutex<HashMap<String, InFlightTurn>>>;
295
296 #[derive(Debug, Deserialize)]
297 struct JsonRpcRequest {
298 #[serde(default)]
299 jsonrpc: Option<String>,
300 #[serde(default)]
301 id: Option<Value>,
302 method: String,
303 #[serde(default)]
304 params: Value,
305 }
306
307 /// Server error: the app-server could not reach the runtime that executes
308 /// turns. Kept in the JSON-RPC implementation-defined server range
309 /// (-32000..-32099) alongside `thread_not_found` (-32004).
310 const RUNTIME_UNAVAILABLE_CODE: i64 = -32005;
311 /// Server error: the named thread does not exist.
312 const THREAD_NOT_FOUND_CODE: i64 = -32004;
313 /// Server error: a daemon-socket client tried to act before `daemon/attach`.
314 /// Only the unix listener raises it; gated so the Windows build (where the
315 /// listener is a typed-unsupported stub) does not fail `warnings = "deny"`
316 /// on dead code.
317 #[cfg(any(unix, windows))]
318 const ATTACH_REQUIRED_CODE: i64 = -32010;
319 /// Server error: a `daemon/attach` claim lost to a live owner.
320 #[cfg(any(unix, windows))]
321 const DAEMON_ALREADY_CLAIMED_CODE: i64 = -32011;
322 /// Server error: only the owning client may `shutdown` the daemon.
323 const NOT_DAEMON_OWNER_CODE: i64 = -32012;
324 /// Server error: the client refused the daemon's version at attach time.
325 #[cfg(any(unix, windows))]
326 const DAEMON_VERSION_SKEW_CODE: i64 = -32013;
327 /// Server error: `daemon/attach` sent twice on one connection.
328 const ALREADY_ATTACHED_CODE: i64 = -32014;
329
330 #[derive(Debug)]
331 struct JsonRpcError {
332 code: i64,
333 message: String,
334 data: Option<Value>,
335 }
336
337 #[derive(Debug)]
338 struct StdioDispatchResult {
339 result: Value,
340 should_exit: bool,
341 }
342
343 struct RuntimeBridge {
344 base_url: String,
345 client: reqwest::Client,
346 auth_token: Option<String>,
347 #[cfg(test)]
348 child: Option<Child>,
349 /// Captured owner SSE cursors, keyed by Runtime thread identity.
350 last_seq_by_thread: HashMap<String, u64>,
351 }
352
353 impl std::fmt::Debug for RuntimeBridge {
354 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
355 formatter
356 .debug_struct("RuntimeBridge")
357 .field("base_url", &self.base_url)
358 .field("authenticated", &self.auth_token.is_some())
359 .field("last_seq_by_thread", &self.last_seq_by_thread)
360 .finish_non_exhaustive()
361 }
362 }
363
364 #[derive(Debug, Clone, Default)]
365 struct RuntimeThreadHint {
366 model: Option<String>,
367 workspace: Option<PathBuf>,
368 }
369
370 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
371 enum TurnTerminalStatus {
372 Completed,
373 Failed,
374 Interrupted,
375 Canceled,
376 }
377
378 /// Structured capture of one bridged turn, for callers that must *return*
379 /// the turn instead of streaming it (HTTP `/prompt`, HTTP `/thread` messages).
380 ///
381 /// The stdio path streams the same events to its writer and needs none of
382 /// this, so it passes `None` and pays nothing.
383 #[derive(Debug, Default)]
384 struct TurnTranscript {
385 /// Concatenated `agent_message` deltas — the model's actual output.
386 text: String,
387 /// The model the runtime reports for the thread that ran the turn.
388 model: Option<String>,
389 /// The same frames the stdio path writes, in order.
390 events: Vec<EventFrame>,
391 }
392
393 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
394 enum AppTransport {
395 Http,
396 Stdio,
397 /// Unix-domain-socket daemon transport (`daemon_socket`). Speaks the
398 /// stdio JSON-RPC protocol verbatim after a `daemon/attach` handshake.
399 Socket,
400 }
401
402 impl AppTransport {
403 /// Wire label reported by `healthz` / `capabilities`.
404 fn label(self) -> &'static str {
405 match self {
406 Self::Http => "http",
407 Self::Stdio => "stdio",
408 Self::Socket => {
409 #[cfg(windows)]
410 {
411 "named-pipe"
412 }
413 #[cfg(not(windows))]
414 {
415 "unix-socket"
416 }
417 }
418 }
419 }
420 }
421
422 /// Whether the peer driving a JSON-RPC loop may stop the whole server.
423 ///
424 /// The process-owned stdio loop always may (its peer *is* the supervisor).
425 /// On the daemon socket only the client that claimed the daemon may; every
426 /// other attached client is refused with `not_daemon_owner` — the brief's
427 /// "never terminate a daemon the app did not spawn", enforced server-side.
428 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
429 enum ShutdownAuthority {
430 Granted,
431 Denied,
432 }
433
434 /// Per-connection policy for [`run_stdio_loop`].
435 #[derive(Debug, Clone, Copy)]
436 struct StdioLoopPolicy {
437 transport: AppTransport,
438 shutdown: ShutdownAuthority,
439 }
440
441 impl StdioLoopPolicy {
442 /// Legacy stdio qualification comparator; production stdio attaches its owner.
443 const fn process_stdio() -> Self {
444 Self {
445 transport: AppTransport::Stdio,
446 shutdown: ShutdownAuthority::Granted,
447 }
448 }
449 }
450
451 /// Why [`run_stdio_loop`] returned.
452 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
453 enum StdioLoopExit {
454 /// The peer closed its write side; nothing asked the server to stop.
455 InputClosed,
456 /// The peer sent an honoured `shutdown`.
457 Shutdown,
458 }
459
460 #[derive(Debug, Deserialize)]
461 struct ConfigGetParams {
462 key: String,
463 }
464
465 #[derive(Debug, Deserialize)]
466 struct ConfigSetParams {
467 key: String,
468 value: String,
469 }
470
471 #[derive(Debug, Deserialize)]
472 struct ThreadIdParams {
473 thread_id: String,
474 }
475
476 #[derive(Debug, Deserialize)]
477 struct ThreadMessageParams {
478 #[serde(default, rename = "maxOutputTokens", alias = "max_output_tokens")]
479 max_output_tokens: Option<std::num::NonZeroU32>,
480 thread_id: String,
481 input: String,
482 #[serde(default)]
483 images: Vec<RuntimeImageInput>,
484 }
485
486 #[derive(Debug, Deserialize)]
487 struct ThreadInterruptParams {
488 thread_id: String,
489 }
490
491 pub async fn run(options: AppServerOptions) -> Result<()> {
492 let auth_token = resolve_auth_token(&options)?;
493 let state =
494 build_state_off_runtime(options.config_path.clone(), auth_token, AppTransport::Http)
495 .await?;
496 let app = app_router(state, &options.cors_origins);
497
498 let listener = tokio::net::TcpListener::bind(options.listen).await?;
499 axum::serve(listener, app)
500 .with_graceful_shutdown(shutdown_signal())
501 .await?;
502 Ok(())
503 }
504
505 async fn shutdown_signal() {
506 let ctrl_c = async {
507 let _ = tokio::signal::ctrl_c().await;
508 };
509
510 #[cfg(unix)]
511 let terminate = async {
512 match tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate()) {
513 Ok(mut signal) => {
514 signal.recv().await;
515 }
516 Err(_) => std::future::pending::<()>().await,
517 }
518 };
519
520 #[cfg(not(unix))]
521 let terminate = std::future::pending::<()>();
522
523 tokio::select! {
524 _ = ctrl_c => {}
525 _ = terminate => {}
526 }
527 }
528
529 /// Protected routes the `Capabilities` response advertises. A test sends a
530 /// request to each one through [`app_router`], so an entry here without a
531 /// handler fails the tests instead of shipping a dead route.
532 ///
533 /// There is no `/tool`: a direct tool call outside a turn would need its own
534 /// tool catalog and approval decision, and the Engine behind the runtime
535 /// bridge is the only tool and approval authority. Tools run inside turns
536 /// (`/prompt`, `/thread` messages). This server does not surface approvals:
537 /// `RuntimeBridge::stream_turn_events` forwards only `item.delta` and the
538 /// turn's completion, and there is no decision route, so approval-gated work
539 /// belongs on the Runtime API (`/v1/threads/*`, `POST /v1/approvals/{id}`).
540 const ADVERTISED_ROUTES: &[&str] = &["/thread", "/app", "/prompt", "/jobs"];
541
542 /// Existing compatibility routes adopted by the canonical host listener.
543 /// Both callers share the exact existing dispatcher/auth/body-limit routes.
544 pub fn runtime_compatibility_router(
545 mut state: AppState,
546 cors_origins: &[String],
547 auth_token: Option<String>,
548 workspace: Option<PathBuf>,
549 ) -> Router {
550 // Only the listener's gate changes. The captured bridge keeps the original
551 // owner's in-memory credential and all shared dispatcher state.
552 state.auth_token = auth_token;
553 state.frontend_workspace = workspace;
554 app_router(state, cors_origins)
555 }
556
557 fn app_router(state: AppState, cors_origins: &[String]) -> Router {
558 let protected_routes = Router::new()
559 .route(
560 "/thread",
561 post(thread_handler).layer(axum::extract::DefaultBodyLimit::max(
562 MAX_RUNTIME_IMAGE_BODY_BYTES,
563 )),
564 )
565 .route("/app", post(app_handler))
566 .route(
567 "/prompt",
568 post(prompt_handler).layer(axum::extract::DefaultBodyLimit::max(
569 MAX_RUNTIME_IMAGE_BODY_BYTES,
570 )),
571 )
572 .route("/jobs", get(jobs_handler))
573 .route(
574 "/v1/chat/completions",
575 post(chat_completions::chat_completions_handler),
576 )
577 .route_layer(middleware::from_fn_with_state(
578 state.clone(),
579 require_app_server_token,
580 ));
581
582 Router::new()
583 .route("/healthz", get(healthz))
584 .merge(protected_routes)
585 .layer(DefaultBodyLimit::max(MAX_HTTP_BODY_BYTES))
586 .layer(cors_layer(cors_origins))
587 .with_state(state)
588 }
589
590 /// Attach the existing compatibility dispatcher to the actual held Runtime.
591 /// Auth remains process-private; it is never part of the owner receipt or IPC.
592 #[cfg(any(unix, windows))]
593 pub async fn bind_runtime_owner(
594 config_path: Option<PathBuf>,
595 endpoint: SocketAddr,
596 auth_token: Option<String>,
597 owner: codewhale_protocol::RuntimeOwnerReceipt,
598 ) -> Result<daemon_socket::DaemonSocket> {
599 let (daemon, _) = bind_runtime_frontends(
600 config_path,
601 auth_token,
602 owner,
603 RuntimeOwnerRouting {
604 endpoint,
605 workspace: None,
606 workers: None,
607 mobile: false,
608 web: false,
609 acp: false,
610 acp_only: false,
611 },
612 None,
613 )
614 .await?;
615 Ok(daemon)
616 }
617
618 /// Build the actual compatibility state once for all owner frontends.
619 #[cfg(any(unix, windows))]
620 pub async fn bind_runtime_frontends(
621 config_path: Option<PathBuf>,
622 auth_token: Option<String>,
623 owner: codewhale_protocol::RuntimeOwnerReceipt,
624 routing: RuntimeOwnerRouting,
625 owner_frontend: Option<Arc<dyn RuntimeOwnerFrontend>>,
626 ) -> Result<(daemon_socket::DaemonSocket, AppState)> {
627 anyhow::ensure!(
628 owner.version == 1 && owner.pid == std::process::id() && !owner.lease_generation.is_empty(),
629 "invalid captured Runtime owner"
630 );
631 let mut state =
632 build_state_off_runtime(config_path, auth_token.clone(), AppTransport::Socket).await?;
633 install_rustls_crypto_provider();
634 let endpoint = routing.endpoint;
635 let address = match endpoint.ip() {
636 std::net::IpAddr::V4(ip) if ip.is_unspecified() => {
637 SocketAddr::from((std::net::Ipv4Addr::LOCALHOST, endpoint.port()))
638 }
639 std::net::IpAddr::V6(ip) if ip.is_unspecified() => {
640 SocketAddr::from((std::net::Ipv6Addr::LOCALHOST, endpoint.port()))
641 }
642 _ => endpoint,
643 };
644 let bridge = RuntimeBridge {
645 base_url: format!("http://{address}"),
646 client: codewhale_release::platform_http_client_builder()
647 .redirect(reqwest::redirect::Policy::none())
648 .build()?,
649 auth_token,
650 #[cfg(test)]
651 child: None,
652 last_seq_by_thread: HashMap::new(),
653 };
654 state.captured_owner = Some(owner.clone());
655 anyhow::ensure!(
656 routing.acp == owner_frontend.is_some() && (!routing.acp_only || routing.acp),
657 "ACP projection must belong to its actual captured owner"
658 );
659 state.frontend_workspace = routing.workspace.clone();
660 state.captured_routing = Some(routing);
661 state.owner_frontend = owner_frontend;
662 *state.runtime_bridge.lock().await = Some(Arc::new(Mutex::new(bridge)));
663 let daemon = daemon_socket::bind_captured_owner(state.clone(), owner).await?;
664 Ok((daemon, state))
665 }
666
667 pub async fn run_owned_acp(state: AppState) -> Result<()> {
668 anyhow::ensure!(
669 state
670 .captured_owner
671 .as_ref()
672 .is_some_and(|owner| owner.pid == std::process::id()),
673 "ACP stdio must belong to its actual captured host"
674 );
675 state
676 .owner_frontend
677 .as_ref()
678 .context("captured owner has no ACP projection")?
679 .serve(
680 RuntimeOwnerFrontendSelection::Acp {
681 scope: None,
682 model: None,
683 },
684 state.clone(),
685 Box::new(BufReader::new(tokio::io::stdin())),
686 Box::new(tokio::io::BufWriter::new(tokio::io::stdout())),
687 )
688 .await
689 }
690
691 /// The existing charged socket connection invokes the same priority-aware
692 /// dispatcher after its captured Runtime has admitted the selected scope.
693 pub async fn run_guest_control(
694 mut state: AppState,
695 workspace: PathBuf,
696 input: Box<dyn AsyncBufRead + Send + Unpin>,
697 output: Box<dyn AsyncWrite + Send + Unpin>,
698 ) -> Result<()> {
699 anyhow::ensure!(
700 state
701 .captured_owner
702 .as_ref()
703 .is_some_and(|owner| owner.pid == std::process::id()),
704 "control projection requires its actual captured host"
705 );
706 state.frontend_workspace = Some(workspace);
707 let policy = StdioLoopPolicy {
708 transport: AppTransport::Socket,
709 shutdown: ShutdownAuthority::Denied,
710 };
711 run_stdio_loop(&state, BoundedLines::new(input), output, policy, None::<()>).await?;
712 Ok(())
713 }
714
715 /// Logical write EOF for duplex transports. This is only a per-connection
716 /// lifecycle notification; it carries no tool, turn or host-shutdown authority.
717 pub fn is_control_input_closed(message: &Value) -> bool {
718 message["jsonrpc"] == "2.0"
719 && message.get("id").is_none()
720 && message["method"] == "daemon/input_closed"
721 && message["params"]
722 .as_object()
723 .is_some_and(|params| params.is_empty())
724 }
725
726 pub async fn run_owned_stdio(state: AppState) -> Result<()> {
727 anyhow::ensure!(
728 state
729 .captured_owner
730 .as_ref()
731 .is_some_and(|owner| owner.pid == std::process::id()),
732 "stdio must belong to its actual captured Runtime host"
733 );
734 let lines = BoundedLines::new(BufReader::new(tokio::io::stdin()));
735 let writer = tokio::io::BufWriter::new(tokio::io::stdout());
736 run_stdio_loop(
737 &state,
738 lines,
739 writer,
740 StdioLoopPolicy::process_stdio(),
741 None::<()>,
742 )
743 .await?;
744 Ok(())
745 }
746
747 pub async fn run_stdio(config_path: Option<PathBuf>) -> Result<()> {
748 daemon_client::forward_stdio(config_path).await
749 }
750
751 /// Cancellation-safe, allocation-bounded newline framing shared by every
752 /// local control transport. Consumed partial bytes remain owned by the reader.
753 pub struct BoundedLines<R> {
754 reader: R,
755 pending: Vec<u8>,
756 }
757
758 impl<R: AsyncBufRead + Unpin> BoundedLines<R> {
759 pub fn new(reader: R) -> Self {
760 Self {
761 reader,
762 pending: Vec::new(),
763 }
764 }
765
766 pub fn into_inner(self) -> R {
767 self.reader
768 }
769 pub async fn next_line(&mut self) -> std::io::Result<Option<String>> {
770 loop {
771 let available = self.reader.fill_buf().await?;
772 if available.is_empty() {
773 if self.pending.is_empty() {
774 return Ok(None);
775 }
776 if self.pending.len() > MAX_RUNTIME_IMAGE_BODY_BYTES {
777 return Err(std::io::Error::new(
778 std::io::ErrorKind::InvalidData,
779 "control request exceeds the 8 MiB transport limit",
780 ));
781 }
782 let bytes = std::mem::take(&mut self.pending);
783 return String::from_utf8(bytes).map(Some).map_err(|_| {
784 std::io::Error::new(
785 std::io::ErrorKind::InvalidData,
786 "control frame is not UTF-8",
787 )
788 });
789 }
790 let newline = available.iter().position(|&byte| byte == b'\n');
791 let count = newline.unwrap_or(available.len());
792 if count > (MAX_RUNTIME_IMAGE_BODY_BYTES + 1).saturating_sub(self.pending.len()) {
793 return Err(std::io::Error::new(
794 std::io::ErrorKind::InvalidData,
795 "control request exceeds the 8 MiB transport limit",
796 ));
797 }
798 // Do not let Vec's geometric growth double an almost-full frame.
799 self.pending.reserve_exact(count);
800 self.pending.extend_from_slice(&available[..count]);
801 self.reader.consume(count + usize::from(newline.is_some()));
802 if self.pending.len() > MAX_RUNTIME_IMAGE_BODY_BYTES
803 && self.pending.last() != Some(&b'\r')
804 {
805 return Err(std::io::Error::new(
806 std::io::ErrorKind::InvalidData,
807 "control request exceeds the 8 MiB transport limit",
808 ));
809 }
810 if newline.is_some() {
811 if self.pending.last() == Some(&b'\r') {
812 self.pending.pop();
813 }
814 let bytes = std::mem::take(&mut self.pending);
815 return String::from_utf8(bytes).map(Some).map_err(|_| {
816 std::io::Error::new(
817 std::io::ErrorKind::InvalidData,
818 "control frame is not UTF-8",
819 )
820 });
821 }
822 }
823 }
824 }
825
826 /// The stdio JSON-RPC loop, generic over its transport so it can be driven by
827 /// a duplex pipe in tests rather than the process's real stdin/stdout.
828 async fn run_stdio_loop<R, W, C>(
829 state: &AppState,
830 mut reader: BoundedLines<R>,
831 mut writer: W,
832 policy: StdioLoopPolicy,
833 // Dropped the moment input closes, not when the in-flight turn ends. The
834 // socket transport passes its owner claim here: a `thread/message` can run
835 // for minutes, and an owner who disconnects mid-turn must not keep the
836 // daemon claimed for the rest of it, or a relaunched client is locked out
837 // with `daemon_already_claimed` and cannot even shut the daemon down.
838 // Process stdio has no claim and passes `None`.
839 mut input_claim: Option<C>,
840 ) -> Result<StdioLoopExit>
841 where
842 R: AsyncBufRead + Unpin,
843 W: AsyncWrite + Unpin,
844 C: Send,
845 {
846 // Work that arrived while a turn was streaming. The turn owns the writer
847 // for its whole duration, so these wait for it rather than interleaving
848 // into the middle of a response.
849 let mut pending: VecDeque<(PendingStdioWork, usize)> = VecDeque::new();
850 let mut pending_bytes = 0usize;
851 let mut stdin_open = true;
852
853 loop {
854 let next = pending.pop_front().map(|(work, bytes)| {
855 pending_bytes -= bytes;
856 work
857 });
858 let request = match next {
859 Some(PendingStdioWork::Response(response)) => {
860 write_stdio_line(&mut writer, &response).await?;
861 continue;
862 }
863 Some(PendingStdioWork::Request(request)) => request,
864 None => {
865 if !stdin_open {
866 return Ok(StdioLoopExit::InputClosed);
867 }
868 let Some(line) = reader.next_line().await? else {
869 return Ok(StdioLoopExit::InputClosed);
870 };
871 match parse_stdio_line(&line) {
872 ParsedStdioLine::Blank => continue,
873 ParsedStdioLine::Rejected(response) => {
874 write_stdio_line(&mut writer, &response).await?;
875 continue;
876 }
877 ParsedStdioLine::Request(request) => request,
878 }
879 }
880 };
881
882 if is_control_detach(&request) {
883 drop(input_claim.take());
884 return Ok(StdioLoopExit::InputClosed);
885 }
886 let id = request.id.clone();
887 if request.method == "shutdown" && policy.shutdown == ShutdownAuthority::Denied {
888 write_stdio_line(
889 &mut writer,
890 &jsonrpc_error(id, JsonRpcError::not_daemon_owner()),
891 )
892 .await?;
893 continue;
894 }
895 let dispatched = if request.method == "thread/message" {
896 // A turn can run for minutes. Keep reading stdin while it streams
897 // so an interrupt (or a shutdown) can actually reach it — with a
898 // plain `await` here, nothing could be read until it finished.
899 let dispatch = dispatch_stdio_request_with_writer(
900 state,
901 &mut writer,
902 &request.method,
903 request.params,
904 policy.transport,
905 );
906 tokio::pin!(dispatch);
907 loop {
908 tokio::select! {
909 outcome = &mut dispatch => break outcome,
910 line = reader.next_line(), if stdin_open => {
911 match line? {
912 None => {
913 stdin_open = false;
914 // Release the claim here, not after `dispatch`
915 // resolves.
916 drop(input_claim.take());
917 }
918 Some(line) => {
919 if matches!(parse_stdio_line(&line), ParsedStdioLine::Request(ref request) if is_control_detach(request)) {
920 stdin_open=false; drop(input_claim.take());
921 } else {handle_line_during_turn(state, &line, &mut pending, &mut pending_bytes, policy).await?;}
922 }
923 }
924 }
925 }
926 }
927 } else {
928 dispatch_stdio_request_with_writer(
929 state,
930 &mut writer,
931 &request.method,
932 request.params,
933 policy.transport,
934 )
935 .await
936 };
937
938 match dispatched {
939 Ok(dispatch) => {
940 write_stdio_line(&mut writer, &jsonrpc_result(id, dispatch.result)).await?;
941 if dispatch.should_exit {
942 return Ok(StdioLoopExit::Shutdown);
943 }
944 }
945 Err(err) => {
946 write_stdio_line(&mut writer, &jsonrpc_error(id, err)).await?;
947 }
948 }
949 }
950 }
951
952 /// Work deferred until a streaming turn releases the writer.
953 enum PendingStdioWork {
954 /// Already answered (an interrupt acted immediately); just needs writing.
955 Response(Value),
956 /// Not started yet; runs normally once the turn is done.
957 Request(JsonRpcRequest),
958 }
959
960 /// Fixed retained queue limits; exhaustion is a visible connection refusal
961 /// before deferred work is admitted. Already-running effects are never replayed.
962 fn queue_stdio_work(
963 pending: &mut VecDeque<(PendingStdioWork, usize)>,
964 bytes: &mut usize,
965 work: PendingStdioWork,
966 input_bytes: usize,
967 ) -> Result<()> {
968 let retained = input_bytes
969 .checked_add(1024)
970 .context("control queue size overflow")?;
971 anyhow::ensure!(
972 pending.len() < 64 && retained <= MAX_RUNTIME_IMAGE_BODY_BYTES.saturating_sub(*bytes),
973 "control queue limit exceeded; deferred request not admitted, in-flight outcomes may be uncertain"
974 );
975 *bytes += retained;
976 pending.push_back((work, retained));
977 Ok(())
978 }
979
980 enum ParsedStdioLine {
981 Blank,
982 Request(JsonRpcRequest),
983 Rejected(Value),
984 }
985
986 fn parse_stdio_line(line: &str) -> ParsedStdioLine {
987 if line.len() > MAX_RUNTIME_IMAGE_BODY_BYTES {
988 return ParsedStdioLine::Rejected(jsonrpc_error(
989 None,
990 JsonRpcError::invalid_params("request exceeds the 8 MiB transport limit"),
991 ));
992 }
993 if line.trim().is_empty() {
994 return ParsedStdioLine::Blank;
995 }
996 let request: JsonRpcRequest = match serde_json::from_str(line) {
997 Ok(value) => value,
998 Err(err) => {
999 return ParsedStdioLine::Rejected(jsonrpc_error(
1000 None,
1001 JsonRpcError::parse_error(format!("invalid json: {err}")),
1002 ));
1003 }
1004 };
1005 if request
1006 .jsonrpc
1007 .as_deref()
1008 .is_some_and(|version| version != "2.0")
1009 {
1010 return ParsedStdioLine::Rejected(jsonrpc_error(
1011 request.id,
1012 JsonRpcError::invalid_request("jsonrpc version must be 2.0"),
1013 ));
1014 }
1015 ParsedStdioLine::Request(request)
1016 }
1017
1018 /// Triage a request that arrived mid-turn.
1019 ///
1020 /// Cancellation is the whole point of reading here, so `thread/interrupt`
1021 /// runs immediately and only its reply waits for the writer. `shutdown` also
1022 /// interrupts immediately — otherwise it would block on the bridge mutex the
1023 /// turn is holding — and then queues so the turn can unwind first. Everything
1024 /// else simply queues: it was never urgent, and running it now would race the
1025 /// turn for the writer.
1026 async fn handle_line_during_turn(
1027 state: &AppState,
1028 line: &str,
1029 pending: &mut VecDeque<(PendingStdioWork, usize)>,
1030 pending_bytes: &mut usize,
1031 policy: StdioLoopPolicy,
1032 ) -> Result<()> {
1033 let request = match parse_stdio_line(line) {
1034 ParsedStdioLine::Blank => return Ok(()),
1035 ParsedStdioLine::Rejected(response) => {
1036 return queue_stdio_work(
1037 pending,
1038 pending_bytes,
1039 PendingStdioWork::Response(response),
1040 line.len(),
1041 );
1042 }
1043 ParsedStdioLine::Request(request) => request,
1044 };
1045
1046 match request.method.as_str() {
1047 "thread/interrupt" => {
1048 let id = request.id.clone();
1049 let response = match parse_params::<ThreadInterruptParams>(params_or_object(
1050 request.params.clone(),
1051 )) {
1052 Ok(parsed) => match interrupt_stdio_turn(state, &parsed.thread_id).await {
1053 Ok(interrupted) => jsonrpc_result(
1054 id,
1055 json!({ "thread_id": parsed.thread_id, "interrupted": interrupted }),
1056 ),
1057 Err(err) => jsonrpc_error(id, err),
1058 },
1059 Err(err) => jsonrpc_error(id, err),
1060 };
1061 queue_stdio_work(
1062 pending,
1063 pending_bytes,
1064 PendingStdioWork::Response(response),
1065 line.len(),
1066 )?;
1067 }
1068 "shutdown" if policy.shutdown == ShutdownAuthority::Denied => {
1069 // A non-owner may not even interrupt the live turns: that is the
1070 // first half of what shutdown does.
1071 queue_stdio_work(
1072 pending,
1073 pending_bytes,
1074 PendingStdioWork::Response(jsonrpc_error(
1075 request.id,
1076 JsonRpcError::not_daemon_owner(),
1077 )),
1078 line.len(),
1079 )?;
1080 }
1081 "shutdown" => {
1082 let _ = interrupt_all_stdio_turns(state).await;
1083 queue_stdio_work(
1084 pending,
1085 pending_bytes,
1086 PendingStdioWork::Request(request),
1087 line.len(),
1088 )?;
1089 }
1090 _ => queue_stdio_work(
1091 pending,
1092 pending_bytes,
1093 PendingStdioWork::Request(request),
1094 line.len(),
1095 )?,
1096 }
1097 Ok(())
1098 }
1099
1100 fn is_control_detach(request: &JsonRpcRequest) -> bool {
1101 request.jsonrpc.as_deref() == Some("2.0")
1102 && request.id.is_none()
1103 && request.method == "daemon/input_closed"
1104 && request
1105 .params
1106 .as_object()
1107 .is_some_and(serde_json::Map::is_empty)
1108 }
1109
1110 async fn write_stdio_line<W: AsyncWrite + Unpin>(writer: &mut W, response: &Value) -> Result<()> {
1111 writer.write_all(&serde_json::to_vec(response)?).await?;
1112 writer.write_all(b"\n").await?;
1113 writer.flush().await?;
1114 Ok(())
1115 }
1116
1117 async fn healthz() -> Json<Value> {
1118 Json(json!({
1119 "status": "ok",
1120 "protocol": "v2",
1121 "service": legacy_deepseek_compat::SERVICE_NAME
1122 }))
1123 }
1124
1125 /// Render a routing failure as a typed HTTP error body.
1126 ///
1127 /// Deliberately *not* a success-shaped payload with the error stuffed into a
1128 /// content field: a client must be able to tell "the model said this" from
1129 /// "nothing ran".
1130 fn http_error_from_jsonrpc(err: JsonRpcError) -> (StatusCode, Json<Value>) {
1131 let (status, code) = match err.code {
1132 -32600 | -32602 => (StatusCode::BAD_REQUEST, "invalid_request"),
1133 THREAD_NOT_FOUND_CODE => (StatusCode::NOT_FOUND, "thread_not_found"),
1134 RUNTIME_UNAVAILABLE_CODE => (StatusCode::SERVICE_UNAVAILABLE, "runtime_unavailable"),
1135 _ => (StatusCode::INTERNAL_SERVER_ERROR, "internal_error"),
1136 };
1137 (
1138 status,
1139 Json(json!({
1140 "error": {
1141 "code": code,
1142 "jsonrpc_code": err.code,
1143 "message": err.message,
1144 }
1145 })),
1146 )
1147 }
1148
1149 async fn thread_handler(State(state): State<AppState>, Json(req): Json<ThreadRequest>) -> Response {
1150 // A message is a turn, and turns belong to the runtime — not to the
1151 // bookkeeping `Runtime` behind the other thread operations. This mirrors
1152 // the interception stdio `thread/message` has always done.
1153 if let ThreadRequest::Message {
1154 thread_id,
1155 input,
1156 images,
1157 max_output_tokens,
1158 } = req
1159 {
1160 return match run_http_thread_message(&state, thread_id, input, images, max_output_tokens)
1161 .await
1162 {
1163 Ok(res) => (StatusCode::OK, Json(res)).into_response(),
1164 Err(err) => http_error_from_jsonrpc(err).into_response(),
1165 };
1166 }
1167 match handle_thread_request(&state, req).await {
1168 Ok(res) => (StatusCode::OK, Json(res)).into_response(),
1169 Err(err) => http_error_from_jsonrpc(err).into_response(),
1170 }
1171 }
1172
1173 /// `POST /prompt` — runs a genuine model turn through the runtime bridge.
1174 ///
1175 /// Note what this handler does *not* do: it never takes the `Runtime` write
1176 /// lock. The old implementation held it across the whole request while doing
1177 /// no model work at all.
1178 async fn prompt_handler(State(state): State<AppState>, Json(req): Json<PromptRequest>) -> Response {
1179 let mut sink = tokio::io::sink();
1180 match run_prompt_turn(&state, &mut sink, req).await {
1181 Ok(res) => (StatusCode::OK, Json(res)).into_response(),
1182 Err(err) => http_error_from_jsonrpc(err).into_response(),
1183 }
1184 }
1185
1186 async fn jobs_handler(State(state): State<AppState>) -> Json<AppResponse> {
1187 let runtime = state.runtime.read().await;
1188 Json(runtime.app_status())
1189 }
1190
1191 async fn app_handler(
1192 State(state): State<AppState>,
1193 Json(req): Json<AppRequest>,
1194 ) -> (StatusCode, Json<AppResponse>) {
1195 let response = process_app_request(&state, req, AppTransport::Http).await;
1196 (app_response_status(&response), Json(response))
1197 }
1198
1199 fn app_response_status(response: &AppResponse) -> StatusCode {
1200 if response.ok {
1201 return StatusCode::OK;
1202 }
1203 if response.data.get("request_id").is_some() {
1204 StatusCode::CONFLICT
1205 } else if response
1206 .data
1207 .get("error")
1208 .and_then(Value::as_str)
1209 .is_some_and(|err| err.starts_with(CONFIG_LOAD_ERROR) || err.starts_with(CONFIG_SAVE_ERROR))
1210 {
1211 StatusCode::INTERNAL_SERVER_ERROR
1212 } else {
1213 StatusCode::BAD_REQUEST
1214 }
1215 }
1216
1217 #[cfg(test)]
1218 fn build_state(config_path: Option<PathBuf>, auth_token: Option<String>) -> Result<AppState> {
1219 build_state_with_transport(config_path, auth_token, AppTransport::Http)
1220 }
1221
1222 /// [`build_state_with_transport`] on the blocking pool. Server startup is
1223 /// async, but building state is not: it reads and parses the config file,
1224 /// creates directories, and opens SQLite (a schema migration that may wait
1225 /// out the 5s busy timeout behind another process). Run inline, that parked
1226 /// a Tokio worker.
1227 async fn build_state_off_runtime(
1228 config_path: Option<PathBuf>,
1229 auth_token: Option<String>,
1230 transport: AppTransport,
1231 ) -> Result<AppState> {
1232 daemon_socket::owner_work(move || {
1233 build_state_with_transport(config_path, auth_token, transport)
1234 })
1235 .await
1236 .context("app-server state setup task failed")
1237 }
1238
1239 fn build_state_with_transport(
1240 config_path: Option<PathBuf>,
1241 auth_token: Option<String>,
1242 transport: AppTransport,
1243 ) -> Result<AppState> {
1244 let has_explicit_config_path = config_path.is_some();
1245 let store = ConfigStore::load(config_path)?;
1246 let config_path = has_explicit_config_path.then(|| store.path().to_path_buf());
1247 let config = store.config.clone();
1248 let registry = ModelRegistry::default();
1249
1250 let state_db_path = config_path
1251 .as_ref()
1252 .and_then(|p| p.parent().map(|parent| parent.join("state.db")));
1253 let state_store = StateStore::open(state_db_path)?;
1254
1255 let mut hooks = HookDispatcher::default();
1256 // Stdio carries JSON-RPC on stdout: printing raw hook events there
1257 // corrupts the protocol stream (#5165). HTTP mode keeps the stdout
1258 // sink for local development visibility.
1259 if transport == AppTransport::Http {
1260 hooks.add_sink(Arc::new(StdoutHookSink));
1261 }
1262 let hook_log_path = config_path
1263 .as_ref()
1264 .and_then(|p| p.parent().map(|parent| parent.join("events.jsonl")))
1265 .unwrap_or_else(legacy_deepseek_compat::default_events_log_path);
1266 hooks.add_sink(Arc::new(JsonlHookSink::new(hook_log_path)));
1267
1268 if let Some(socket_path) = config
1269 .hook_sinks
1270 .as_ref()
1271 .and_then(|sinks| sinks.unix_socket_path.as_ref())
1272 .filter(|path| !path.as_os_str().is_empty())
1273 {
1274 hooks.add_sink(Arc::new(UnixSocketHookSink::new(socket_path.clone())));
1275 }
1276
1277 let runtime = Runtime::new(config.clone(), state_store, hooks);
1278
1279 Ok(AppState {
1280 captured_owner: None,
1281 captured_routing: None,
1282 frontend_workspace: None,
1283 owner_frontend: None,
1284 config_path,
1285 config: Arc::new(RwLock::new(config)),
1286 runtime: Arc::new(RwLock::new(runtime)),
1287 registry,
1288 auth_token,
1289 runtime_bridge: Arc::new(Mutex::new(None)),
1290 runtime_thread_map: Arc::new(Mutex::new(HashMap::new())),
1291 stdio_thread_hints: Arc::new(Mutex::new(HashMap::new())),
1292 in_flight_turns: Arc::new(Mutex::new(HashMap::new())),
1293 })
1294 }
1295
1296 fn resolve_auth_token(options: &AppServerOptions) -> Result<Option<String>> {
1297 let configured = options.auth_token.as_ref().map(|token| token.trim());
1298 if let Some(token) = configured
1299 && token.is_empty()
1300 {
1301 bail!("app-server auth token cannot be empty");
1302 }
1303 let has_explicit_token = configured.is_some();
1304
1305 if options.insecure_no_auth {
1306 if !options.listen.ip().is_loopback() {
1307 bail!("refusing unauthenticated app-server bind on non-loopback address");
1308 }
1309 eprintln!("warning: app-server HTTP auth disabled by --insecure-no-auth");
1310 return Ok(None);
1311 }
1312
1313 if !has_explicit_token && !options.listen.ip().is_loopback() {
1314 bail!(
1315 "refusing non-loopback app-server bind without explicit auth token; pass --auth-token or set CODEWHALE_APP_SERVER_TOKEN"
1316 );
1317 }
1318
1319 let token = configured
1320 .map(str::to_string)
1321 .unwrap_or_else(|| format!("cwapp_{}", Uuid::new_v4().simple()));
1322 for line in app_server_auth_status_lines(has_explicit_token) {
1323 eprintln!("{line}");
1324 }
1325 Ok(Some(token))
1326 }
1327
1328 fn app_server_auth_status_lines(has_explicit_token: bool) -> Vec<&'static str> {
1329 if has_explicit_token {
1330 return vec!["app-server auth: bearer token required for HTTP routes."];
1331 }
1332 vec![
1333 "app-server auth: generated bearer token for this process (not printed).",
1334 " Pass --auth-token or set CODEWHALE_APP_SERVER_TOKEN when another client needs to connect.",
1335 ]
1336 }
1337
1338 fn cors_layer(extra_origins: &[String]) -> CorsLayer {
1339 let mut origins: Vec<HeaderValue> = DEFAULT_CORS_ORIGINS
1340 .iter()
1341 .filter_map(|origin| HeaderValue::from_str(origin).ok())
1342 .collect();
1343 for raw in extra_origins {
1344 let trimmed = raw.trim();
1345 if trimmed.is_empty() {
1346 continue;
1347 }
1348 match HeaderValue::from_str(trimmed) {
1349 Ok(value) if !origins.contains(&value) => origins.push(value),
1350 Ok(_) => {}
1351 Err(err) => {
1352 eprintln!("warning: ignoring invalid app-server CORS origin `{trimmed}`: {err}")
1353 }
1354 }
1355 }
1356
1357 CorsLayer::new()
1358 .allow_origin(origins)
1359 .allow_methods([Method::GET, Method::POST, Method::OPTIONS])
1360 .allow_headers([header::AUTHORIZATION, header::CONTENT_TYPE])
1361 }
1362
1363 async fn require_app_server_token(
1364 State(state): State<AppState>,
1365 req: Request,
1366 next: Next,
1367 ) -> Response {
1368 let Some(expected) = state.auth_token.as_deref() else {
1369 return next.run(req).await;
1370 };
1371 let authorized = req
1372 .headers()
1373 .get(header::AUTHORIZATION)
1374 .and_then(|value| value.to_str().ok())
1375 .and_then(|raw| raw.strip_prefix("Bearer "))
1376 .is_some_and(|token| {
1377 codewhale_core::secret_eq::constant_time_eq(token.as_bytes(), expected.as_bytes())
1378 });
1379
1380 if authorized {
1381 next.run(req).await
1382 } else {
1383 (
1384 StatusCode::UNAUTHORIZED,
1385 Json(json!({
1386 "error": {
1387 "message": "app-server bearer token required",
1388 "status": StatusCode::UNAUTHORIZED.as_u16(),
1389 }
1390 })),
1391 )
1392 .into_response()
1393 }
1394 }
1395
1396 fn params_or_object(params: Value) -> Value {
1397 if params.is_null() { json!({}) } else { params }
1398 }
1399
1400 fn parse_params<T: DeserializeOwned>(params: Value) -> std::result::Result<T, JsonRpcError> {
1401 serde_json::from_value(params).map_err(|err| JsonRpcError::invalid_params(err.to_string()))
1402 }
1403
1404 fn jsonrpc_result(id: Option<Value>, result: Value) -> Value {
1405 json!({
1406 "jsonrpc": "2.0",
1407 "id": id.unwrap_or(Value::Null),
1408 "result": result
1409 })
1410 }
1411
1412 fn jsonrpc_error(id: Option<Value>, err: JsonRpcError) -> Value {
1413 json!({
1414 "jsonrpc": "2.0",
1415 "id": id.unwrap_or(Value::Null),
1416 "error": {
1417 "code": err.code,
1418 "message": err.message,
1419 "data": err.data
1420 }
1421 })
1422 }
1423
1424 impl JsonRpcError {
1425 fn parse_error(message: impl Into<String>) -> Self {
1426 Self {
1427 code: -32700,
1428 message: message.into(),
1429 data: None,
1430 }
1431 }
1432
1433 fn invalid_request(message: impl Into<String>) -> Self {
1434 Self {
1435 code: -32600,
1436 message: message.into(),
1437 data: None,
1438 }
1439 }
1440
1441 fn method_not_found(method: &str) -> Self {
1442 Self {
1443 code: -32601,
1444 message: format!("unsupported method: {method}"),
1445 data: None,
1446 }
1447 }
1448
1449 fn invalid_params(message: impl Into<String>) -> Self {
1450 Self {
1451 code: -32602,
1452 message: message.into(),
1453 data: None,
1454 }
1455 }
1456
1457 /// Server error (-32000..-32099): the turn engine could not be reached,
1458 /// or refused to start the turn — either way nothing ran. Distinct from
1459 /// `internal` because the caller can retry this one once a runtime is up.
1460 fn runtime_unavailable(message: impl Into<String>) -> Self {
1461 let message = message.into();
1462 Self {
1463 code: RUNTIME_UNAVAILABLE_CODE,
1464 message: message.clone(),
1465 data: Some(json!({
1466 "error": "runtime_unavailable",
1467 "detail": message,
1468 })),
1469 }
1470 }
1471
1472 /// Server error (-32000..-32099): the named thread does not exist.
1473 fn thread_not_found(thread_id: &str) -> Self {
1474 Self {
1475 code: THREAD_NOT_FOUND_CODE,
1476 message: format!("thread not found: {thread_id}"),
1477 data: Some(json!({
1478 "error": "thread_not_found",
1479 "thread_id": thread_id,
1480 })),
1481 }
1482 }
1483
1484 fn internal(message: impl Into<String>) -> Self {
1485 Self {
1486 code: -32603,
1487 message: message.into(),
1488 data: None,
1489 }
1490 }
1491
1492 /// Server error (-32000..-32099): the daemon-socket connection has not
1493 /// completed `daemon/attach`, so nothing but `healthz` is allowed yet.
1494 #[cfg(any(unix, windows))]
1495 fn attach_required(method: &str) -> Self {
1496 Self {
1497 code: ATTACH_REQUIRED_CODE,
1498 message: format!("send daemon/attach before `{method}`"),
1499 data: Some(json!({
1500 "error": "attach_required",
1501 "method": method,
1502 "attach_method": daemon_socket::ATTACH_METHOD,
1503 })),
1504 }
1505 }
1506
1507 /// Server error (-32000..-32099): a `claim` attach lost to a live owner.
1508 #[cfg(any(unix, windows))]
1509 fn daemon_already_claimed(owner: &Value) -> Self {
1510 Self {
1511 code: DAEMON_ALREADY_CLAIMED_CODE,
1512 message: "daemon already claimed by another client; attach with mode=attach"
1513 .to_string(),
1514 data: Some(json!({
1515 "error": "daemon_already_claimed",
1516 "owner": owner,
1517 })),
1518 }
1519 }
1520
1521 /// Server error (-32000..-32099): only the owner may stop the daemon.
1522 fn not_daemon_owner() -> Self {
1523 Self {
1524 code: NOT_DAEMON_OWNER_CODE,
1525 message: "only the client that claimed this daemon may shut it down".to_string(),
1526 data: Some(json!({ "error": "not_daemon_owner" })),
1527 }
1528 }
1529
1530 /// Server error (-32000..-32099): the client's expected daemon version
1531 /// does not match the running binary (bundle skew).
1532 #[cfg(any(unix, windows))]
1533 fn daemon_version_skew(expected: &str, actual: &str) -> Self {
1534 Self {
1535 code: DAEMON_VERSION_SKEW_CODE,
1536 message: format!(
1537 "daemon version {actual} does not match the client's expected {expected}"
1538 ),
1539 data: Some(json!({
1540 "error": "daemon_version_skew",
1541 "expected": expected,
1542 "actual": actual,
1543 })),
1544 }
1545 }
1546
1547 /// Server error (-32000..-32099): `daemon/attach` after attaching.
1548 fn already_attached() -> Self {
1549 Self {
1550 code: ALREADY_ATTACHED_CODE,
1551 message: "this connection is already attached".to_string(),
1552 data: Some(json!({ "error": "already_attached" })),
1553 }
1554 }
1555 }
1556
1557 async fn handle_thread_request(
1558 state: &AppState,
1559 req: ThreadRequest,
1560 ) -> std::result::Result<ThreadResponse, JsonRpcError> {
1561 thread_control::handle(state, req).await
1562 }
1563
1564 /// One turn's worth of routing decisions, shared by every surface that runs
1565 /// a turn through the bridge.
1566 struct RuntimeTurnInput<'a> {
1567 input: &'a str,
1568 images: &'a [RuntimeImageInput],
1569 max_output_tokens: Option<std::num::NonZeroU32>,
1570 expected_workspace: Option<&'a Path>,
1571 }
1572
1573 struct BridgedTurn<'a> {
1574 max_output_tokens: Option<std::num::NonZeroU32>,
1575 /// Client-facing thread id; the bridge maps it to a runtime thread.
1576 thread_key: &'a str,
1577 input: &'a str,
1578 images: &'a [RuntimeImageInput],
1579 /// Model for the runtime thread when this call is the one that creates
1580 /// it. An existing thread keeps the model it was created with.
1581 model_override: Option<String>,
1582 /// Publish the live turn so a concurrent `thread/interrupt` can cancel
1583 /// it. Only stdio has a mid-turn channel, so only stdio sets this.
1584 interruptible: bool,
1585 /// Forget the thread mapping once the turn ends. Set for one-shot
1586 /// prompts, whose synthetic thread key no client can name again.
1587 ephemeral: bool,
1588 /// Refuse a `thread_key` that is neither mapped, durably linked, nor a
1589 /// persisted thread, instead of minting an empty runtime thread for it.
1590 /// Set by `thread/message`, whose ids come from `thread/create`.
1591 require_known_thread: bool,
1592 }
1593
1594 /// Execute exactly one turn on the real runtime.
1595 ///
1596 /// This is the only way any app-server surface runs a model: `/prompt`,
1597 /// `prompt/request`, `prompt/run`, stdio `thread/message`, and HTTP `/thread`
1598 /// messages all land here. There is no local fallback that fabricates a
1599 /// response — if the runtime cannot be reached the caller gets
1600 /// [`JsonRpcError::runtime_unavailable`] and nothing is written to history.
1601 async fn run_bridged_turn<W: AsyncWrite + Unpin>(
1602 state: &AppState,
1603 writer: &mut W,
1604 turn: BridgedTurn<'_>,
1605 transcript: Option<&mut TurnTranscript>,
1606 ) -> std::result::Result<Value, JsonRpcError> {
1607 let mut hint = {
1608 let hints = state.stdio_thread_hints.lock().await;
1609 hints.get(turn.thread_key).cloned()
1610 };
1611 if let Some(model) = turn.model_override {
1612 hint.get_or_insert_with(RuntimeThreadHint::default).model = Some(model);
1613 }
1614 // Durable compatibility IDs resolve through full canonical history and
1615 // exact owner receipts before the existing turn transport can begin.
1616 let resolved = if turn.ephemeral {
1617 None
1618 } else {
1619 Some(
1620 thread_control::resolve(state, turn.thread_key, true)
1621 .await
1622 .map_err(|error| thread_control::rpc_error(error, Some(turn.thread_key)))?,
1623 )
1624 };
1625 if turn.require_known_thread && resolved.is_none() {
1626 return Err(JsonRpcError::thread_not_found(turn.thread_key));
1627 }
1628 if let Some((_, workspace)) = resolved.as_ref() {
1629 hint.get_or_insert_with(RuntimeThreadHint::default)
1630 .workspace = Some(workspace.clone());
1631 } else if let Some(workspace) = state.frontend_workspace.as_ref() {
1632 hint.get_or_insert_with(RuntimeThreadHint::default)
1633 .workspace
1634 .get_or_insert_with(|| workspace.clone());
1635 }
1636 // Retain the acknowledged/observed scope across all transport awaits.
1637 // The owning Runtime rechecks this assertion at its final turn admission.
1638 let expected_workspace = hint.as_ref().and_then(|hint| hint.workspace.clone());
1639 let mut bridge = acquire_live_runtime_bridge(state).await?;
1640 let mut thread_map = state.runtime_thread_map.clone().lock_owned().await;
1641 if let Some((id, _)) = resolved {
1642 thread_map.insert(turn.thread_key.to_string(), id);
1643 }
1644 if turn.max_output_tokens.is_some() {
1645 let info = bridge
1646 .request_json(
1647 bridge.authed(
1648 bridge
1649 .client
1650 .get(format!("{}/v1/runtime/info", bridge.base_url)),
1651 ),
1652 )
1653 .await
1654 .map_err(|error| JsonRpcError::runtime_unavailable(error.to_string()))?;
1655 if info
1656 .pointer("/capabilities/turn_output_token_limit")
1657 .and_then(Value::as_bool)
1658 != Some(true)
1659 {
1660 return Err(JsonRpcError::invalid_params(
1661 "Runtime does not support maxOutputTokens",
1662 ));
1663 }
1664 if !thread_map.contains_key(turn.thread_key) {
1665 bridge
1666 .require_output_limited_model(hint.as_ref().and_then(|hint| hint.model.as_deref()))
1667 .await
1668 .map_err(|error| JsonRpcError::invalid_params(error.to_string()))?;
1669 }
1670 }
1671 let runtime_thread_id = bridge
1672 .ensure_runtime_thread(&mut thread_map, turn.thread_key, hint)
1673 .await
1674 .map_err(|error| JsonRpcError::runtime_unavailable(error.to_string()))?;
1675 drop(thread_map);
1676 let registration = turn
1677 .interruptible
1678 .then(|| (state.in_flight_turns.clone(), turn.thread_key.to_string()));
1679 let result = bridge
1680 .message_thread(
1681 &runtime_thread_id,
1682 RuntimeTurnInput {
1683 input: turn.input,
1684 images: turn.images,
1685 max_output_tokens: turn.max_output_tokens,
1686 expected_workspace: expected_workspace.as_deref(),
1687 },
1688 writer,
1689 registration,
1690 transcript,
1691 )
1692 .await;
1693 if turn.ephemeral {
1694 // Drop the mapping while we still hold the bridge lock, so a
1695 // long-lived app-server does not accumulate one entry per one-shot
1696 // prompt.
1697 let mut thread_map = state.runtime_thread_map.lock().await;
1698 bridge.forget_thread(&mut thread_map, turn.thread_key);
1699 }
1700 result.map_err(|err| JsonRpcError::internal(err.to_string()))
1701 }
1702
1703 /// Run a prompt as a genuine model turn and return what the model actually
1704 /// said.
1705 ///
1706 /// `writer` receives the same streaming frames stdio `thread/message` emits;
1707 /// HTTP callers pass a sink and read the frames back out of
1708 /// [`PromptResponse::events`].
1709 async fn run_prompt_turn<W: AsyncWrite + Unpin>(
1710 state: &AppState,
1711 writer: &mut W,
1712 req: PromptRequest,
1713 ) -> std::result::Result<PromptResponse, JsonRpcError> {
1714 if req.prompt.trim().is_empty() {
1715 return Err(JsonRpcError::invalid_params("prompt must not be empty"));
1716 }
1717 // The turn engine has no threadless mode, so a prompt without a thread
1718 // gets a fresh one. Keying it on a uuid keeps a one-shot prompt out of
1719 // any caller's history and out of the way of concurrent prompts.
1720 let ephemeral = req.thread_id.is_none();
1721 let thread_key = req
1722 .thread_id
1723 .clone()
1724 .unwrap_or_else(|| format!("prompt-{}", Uuid::new_v4()));
1725
1726 let mut transcript = TurnTranscript::default();
1727 run_bridged_turn(
1728 state,
1729 writer,
1730 BridgedTurn {
1731 max_output_tokens: req.max_output_tokens,
1732 thread_key: &thread_key,
1733 input: &req.prompt,
1734 images: &req.images,
1735 model_override: req.model.clone(),
1736 // `thread/interrupt` addresses client-facing thread ids. A
1737 // one-shot prompt has none to hand back, and a caller-supplied
1738 // thread id is already interruptible through `thread/message`.
1739 interruptible: false,
1740 ephemeral,
1741 // A caller-chosen `/prompt` thread key keeps its conversation
1742 // but need not name a `thread/create` thread.
1743 require_known_thread: false,
1744 },
1745 Some(&mut transcript),
1746 )
1747 .await?;
1748
1749 // Report the model the runtime actually ran, never a locally resolved
1750 // guess. The fallbacks only matter for a runtime that omits the field.
1751 let model = match transcript.model {
1752 Some(model) => model,
1753 None => match req.model {
1754 Some(model) => model,
1755 None => state
1756 .config
1757 .read()
1758 .await
1759 .model
1760 .clone()
1761 .unwrap_or_else(|| "unknown".to_string()),
1762 },
1763 };
1764
1765 Ok(PromptResponse {
1766 output: transcript.text,
1767 model,
1768 events: transcript.events,
1769 })
1770 }
1771
1772 async fn handle_prompt_request<W: AsyncWrite + Unpin>(
1773 state: &AppState,
1774 writer: &mut W,
1775 req: PromptRequest,
1776 ) -> std::result::Result<PromptResponse, JsonRpcError> {
1777 run_prompt_turn(state, writer, req).await
1778 }
1779
1780 /// HTTP `/thread` with a `Message` body: same engine as stdio
1781 /// `thread/message`, but the turn is collected rather than streamed because
1782 /// this transport is request/response.
1783 async fn run_http_thread_message(
1784 state: &AppState,
1785 thread_id: String,
1786 input: String,
1787 images: Vec<RuntimeImageInput>,
1788 max_output_tokens: Option<std::num::NonZeroU32>,
1789 ) -> std::result::Result<ThreadResponse, JsonRpcError> {
1790 let mut transcript = TurnTranscript::default();
1791 let mut sink = tokio::io::sink();
1792 let result = run_bridged_turn(
1793 state,
1794 &mut sink,
1795 BridgedTurn {
1796 max_output_tokens,
1797 thread_key: &thread_id,
1798 input: &input,
1799 images: &images,
1800 model_override: None,
1801 interruptible: false,
1802 ephemeral: false,
1803 require_known_thread: true,
1804 },
1805 Some(&mut transcript),
1806 )
1807 .await?;
1808
1809 Ok(ThreadResponse {
1810 thread_id,
1811 // The turn ran to a terminal state before this response was built,
1812 // which is exactly what the old `accepted` did not mean.
1813 status: "completed".to_string(),
1814 thread: None,
1815 threads: Vec::new(),
1816 goal: None,
1817 model: transcript.model,
1818 model_provider: None,
1819 cwd: None,
1820 approval_policy: None,
1821 sandbox: None,
1822 events: transcript.events,
1823 data: result.get("data").cloned().unwrap_or_else(|| json!({})),
1824 })
1825 }
1826
1827 async fn handle_stdio_thread_message<W: AsyncWrite + Unpin>(
1828 state: &AppState,
1829 writer: &mut W,
1830 parsed: ThreadMessageParams,
1831 ) -> std::result::Result<Value, JsonRpcError> {
1832 let mut result = run_bridged_turn(
1833 state,
1834 writer,
1835 BridgedTurn {
1836 max_output_tokens: parsed.max_output_tokens,
1837 thread_key: &parsed.thread_id,
1838 input: &parsed.input,
1839 images: &parsed.images,
1840 model_override: None,
1841 interruptible: true,
1842 ephemeral: false,
1843 require_known_thread: true,
1844 },
1845 None,
1846 )
1847 .await?;
1848 if let Some(object) = result.as_object_mut() {
1849 object.insert("thread_id".to_string(), Value::String(parsed.thread_id));
1850 }
1851 Ok(result)
1852 }
1853
1854 /// Resuming, forking, archiving or unarchiving a thread the runtime reports
1855 /// as `missing` must fail with a named not-found error. Recording the null model/workspace of that
1856 /// response as a stdio hint would clobber any previously cached hint for
1857 /// the same thread id (#5171).
1858 fn ensure_thread_found(response: &ThreadResponse) -> std::result::Result<(), JsonRpcError> {
1859 if response.status == "missing" {
1860 return Err(JsonRpcError::thread_not_found(&response.thread_id));
1861 }
1862 Ok(())
1863 }
1864
1865 async fn record_stdio_thread_hint(state: &AppState, response: &ThreadResponse) {
1866 let mut hints = state.stdio_thread_hints.lock().await;
1867 hints.insert(
1868 response.thread_id.clone(),
1869 RuntimeThreadHint {
1870 model: response.model.clone(),
1871 workspace: response.cwd.clone(),
1872 },
1873 );
1874 }
1875
1876 /// Historical cold-child cache comparator. It has no production producer.
1877 #[cfg(all(test, unix))]
1878 async fn acquire_historical_runtime_bridge<F, Fut>(
1879 state: &AppState,
1880 start: F,
1881 ) -> std::result::Result<SharedRuntimeBridge, JsonRpcError>
1882 where
1883 F: FnOnce() -> Fut,
1884 Fut: std::future::Future<Output = Result<RuntimeBridge>>,
1885 {
1886 if let Some(bridge) = state.runtime_bridge.lock().await.as_ref() {
1887 return Ok(bridge.clone());
1888 }
1889 if state.captured_owner.is_some() {
1890 return Err(JsonRpcError::runtime_unavailable(
1891 "the captured Runtime owner is unavailable; refusing another owner",
1892 ));
1893 }
1894 let bridge =
1895 Arc::new(Mutex::new(start().await.map_err(|err| {
1896 JsonRpcError::runtime_unavailable(err.to_string())
1897 })?));
1898 let mut slot = state.runtime_bridge.lock().await;
1899 // Prefer a bridge cached by a concurrent caller while we were spawning;
1900 // dropping our unused one kills the extra child via `Drop`.
1901 Ok(slot.get_or_insert_with(|| bridge.clone()).clone())
1902 }
1903
1904 /// Use the exact bridge installed by the held canonical owner. No process,
1905 /// store or transcript can be created when this captured connection is absent.
1906 async fn acquire_runtime_bridge(
1907 state: &AppState,
1908 ) -> std::result::Result<SharedRuntimeBridge, JsonRpcError> {
1909 state
1910 .runtime_bridge
1911 .lock()
1912 .await
1913 .as_ref()
1914 .cloned()
1915 .ok_or_else(|| {
1916 JsonRpcError::runtime_unavailable(
1917 "the captured Runtime owner is unavailable; refusing another owner",
1918 )
1919 })
1920 }
1921
1922 async fn acquire_live_runtime_bridge(
1923 state: &AppState,
1924 ) -> std::result::Result<tokio::sync::OwnedMutexGuard<RuntimeBridge>, JsonRpcError> {
1925 Ok(acquire_runtime_bridge(state).await?.lock_owned().await)
1926 }
1927
1928 /// Historical respawn comparison; all actual frontends use the held owner.
1929 #[cfg(all(test, unix))]
1930 async fn acquire_live_runtime_bridge_with<F, Fut>(
1931 state: &AppState,
1932 start: F,
1933 ) -> std::result::Result<tokio::sync::OwnedMutexGuard<RuntimeBridge>, JsonRpcError>
1934 where
1935 F: Fn() -> Fut,
1936 Fut: std::future::Future<Output = Result<RuntimeBridge>>,
1937 {
1938 for _ in 0..2 {
1939 let shared = acquire_historical_runtime_bridge(state, &start).await?;
1940 let mut bridge = shared.clone().lock_owned().await;
1941 if !bridge.child_exited() {
1942 return Ok(bridge);
1943 }
1944 drop(bridge);
1945 let mut slot = state.runtime_bridge.lock().await;
1946 // Evict only the dead bridge: a concurrent caller may already have
1947 // replaced it with a live one.
1948 if slot
1949 .as_ref()
1950 .is_some_and(|cached| Arc::ptr_eq(cached, &shared))
1951 {
1952 *slot = None;
1953 }
1954 }
1955 Err(JsonRpcError::runtime_unavailable(
1956 "runtime API bridge exited immediately after starting",
1957 ))
1958 }
1959
1960 /// Ask the runtime to interrupt a turn that is streaming right now.
1961 ///
1962 /// Everything this needs was copied out of the bridge when the turn started,
1963 /// so it never touches the bridge mutex the turn is holding. Returns whether
1964 /// a live turn was found for `thread_id`.
1965 /// Interrupt one in-flight turn over HTTP, from an owned snapshot.
1966 ///
1967 /// Split from [`interrupt_stdio_turn`] so teardown paths can run many
1968 /// concurrently (#6211 R8b) — each future owns its snapshot and never holds
1969 /// the turn registry across the request.
1970 async fn interrupt_turn_request(turn: &InFlightTurn) -> std::result::Result<bool, JsonRpcError> {
1971 let mut request = codewhale_release::platform_http_client_builder()
1972 .redirect(reqwest::redirect::Policy::none())
1973 .timeout(Duration::from_secs(10))
1974 .build()
1975 .map_err(|err| JsonRpcError::internal(err.to_string()))?
1976 .post(format!(
1977 "{}/v1/threads/{}/turns/{}/interrupt",
1978 turn.base_url, turn.runtime_thread_id, turn.turn_id
1979 ));
1980 if let Some(token) = turn.auth_token.as_deref() {
1981 request = request.bearer_auth(token);
1982 }
1983 request
1984 .send()
1985 .await
1986 .and_then(reqwest::Response::error_for_status)
1987 .map_err(|err| JsonRpcError::internal(format!("interrupt failed: {err}")))?;
1988 Ok(true)
1989 }
1990
1991 async fn interrupt_stdio_turn(
1992 state: &AppState,
1993 thread_id: &str,
1994 ) -> std::result::Result<bool, JsonRpcError> {
1995 let Some(turn) = state.in_flight_turns.lock().await.get(thread_id).cloned() else {
1996 return Ok(false);
1997 };
1998 interrupt_turn_request(&turn).await
1999 }
2000
2001 /// Interrupt every in-flight turn concurrently, reporting how many were
2002 /// reached (#6211 R8b).
2003 ///
2004 /// The teardown paths used to await each turn's interrupt in sequence, so
2005 /// their latency grew with the number of live turns — up to the 10s
2006 /// per-request timeout apiece. Each interrupt now owns its snapshot and runs
2007 /// as an independent task; individual failures are ignored exactly as the
2008 /// sequential loop ignored them, and the registry lock is never held across
2009 /// the requests.
2010 async fn interrupt_all_stdio_turns(state: &AppState) -> usize {
2011 let turns: Vec<InFlightTurn> = {
2012 let map = state.in_flight_turns.lock().await;
2013 map.values().cloned().collect()
2014 };
2015 let mut set = tokio::task::JoinSet::new();
2016 for turn in turns {
2017 set.spawn(async move { interrupt_turn_request(&turn).await });
2018 }
2019 let mut interrupted = 0usize;
2020 while let Some(joined) = set.join_next().await {
2021 if matches!(joined, Ok(Ok(true))) {
2022 interrupted += 1;
2023 }
2024 }
2025 interrupted
2026 }
2027
2028 /// Historical config/cache comparator. Captured owners retain their bridge.
2029 #[cfg(test)]
2030 async fn invalidate_runtime_bridge(state: &AppState) {
2031 if state.captured_owner.is_some() {
2032 return;
2033 }
2034 let mut bridge = state.runtime_bridge.lock().await;
2035 *bridge = None;
2036 }
2037
2038 impl RuntimeBridge {
2039 /// The child binds an ephemeral loopback port itself (`--port 0`) and
2040 /// reports it on stdout; the parent never reserves a port for it to race.
2041 #[cfg(test)]
2042 fn runtime_command(config_path: Option<&Path>, auth_token: &str) -> Result<Command> {
2043 let current_exe = std::env::current_exe().ok();
2044 let mut command = if let Some(path) = current_exe {
2045 Command::new(path)
2046 } else {
2047 Command::new("codewhale")
2048 };
2049 // Pass the runtime auth token out-of-band via env (not argv) so local
2050 // `ps` cannot read credential material from the child command line.
2051 // The TUI/runtime server already accepts CODEWHALE_RUNTIME_TOKEN /
2052 // DEEPSEEK_RUNTIME_TOKEN when --auth-token is absent.
2053 command
2054 .arg("app-server")
2055 .arg("--http")
2056 .arg("--host")
2057 .arg("127.0.0.1")
2058 .arg("--port")
2059 .arg("0")
2060 .env("CODEWHALE_RUNTIME_TOKEN", auth_token)
2061 .env("DEEPSEEK_RUNTIME_TOKEN", auth_token)
2062 .stdin(Stdio::null())
2063 .stdout(Stdio::piped())
2064 .stderr(Stdio::null());
2065 if let Some(config_path) = config_path {
2066 command.arg("--config").arg(config_path);
2067 }
2068 Ok(command)
2069 }
2070
2071 fn authed(&self, builder: reqwest::RequestBuilder) -> reqwest::RequestBuilder {
2072 match self.auth_token.as_deref() {
2073 Some(token) => builder.bearer_auth(token),
2074 None => builder,
2075 }
2076 }
2077
2078 async fn request_json(&self, builder: reqwest::RequestBuilder) -> Result<Value> {
2079 thread_control::read_json_response(builder.send().await?).await
2080 }
2081
2082 /// Read the existing Runtime catalog before a one-shot prompt can create
2083 /// its thread. Existing threads are checked by canonical turn admission.
2084 async fn require_output_limited_model(&self, requested_model: Option<&str>) -> Result<()> {
2085 let providers = self
2086 .request_json(self.authed(self.client.get(format!("{}/v1/providers", self.base_url))))
2087 .await?;
2088 let current = providers
2089 .get("current")
2090 .and_then(Value::as_str)
2091 .context("Runtime provider is unavailable")?;
2092 let provider = providers
2093 .get("providers")
2094 .and_then(Value::as_array)
2095 .and_then(|providers| {
2096 providers
2097 .iter()
2098 .find(|provider| provider.get("id").and_then(Value::as_str) == Some(current))
2099 })
2100 .context("Runtime provider is unavailable")?;
2101 let model = requested_model
2102 .or_else(|| provider.get("default_model").and_then(Value::as_str))
2103 .context("maxOutputTokens requires an exact model")?;
2104 if model.trim().is_empty() || model.eq_ignore_ascii_case("auto") {
2105 bail!("maxOutputTokens requires an exact model");
2106 }
2107 if !current
2108 .bytes()
2109 .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'-')
2110 {
2111 bail!("Runtime provider identity is invalid");
2112 }
2113 let mut cursor = None;
2114 let mut seen = std::collections::HashSet::new();
2115 loop {
2116 let mut url =
2117 reqwest::Url::parse(&format!("{}/v1/providers/{current}/models", self.base_url))?;
2118 url.query_pairs_mut().append_pair("limit", "250");
2119 if let Some(cursor) = cursor.as_deref() {
2120 url.query_pairs_mut().append_pair("cursor", cursor);
2121 }
2122 let catalog = self.request_json(self.authed(self.client.get(url))).await?;
2123 if let Some(entry) =
2124 catalog
2125 .get("models")
2126 .and_then(Value::as_array)
2127 .and_then(|models| {
2128 models
2129 .iter()
2130 .find(|entry| entry.get("id").and_then(Value::as_str) == Some(model))
2131 })
2132 {
2133 if entry.get("output_token_limit").and_then(Value::as_str) == Some("supported") {
2134 return Ok(());
2135 }
2136 bail!("The selected Runtime model does not support maxOutputTokens");
2137 }
2138 let next = catalog
2139 .get("nextCursor")
2140 .and_then(Value::as_str)
2141 .context("Output-limit support is unknown for the selected Runtime model")?;
2142 if !seen.insert(next.to_string()) {
2143 bail!("Runtime model catalog cursor repeated");
2144 }
2145 cursor = Some(next.to_string());
2146 }
2147 }
2148
2149 /// Resolve `stdio_thread_id` to a runtime thread, minting one only when
2150 /// `thread_map` has no entry. The map lives on [`AppState`] and outlives
2151 /// this bridge, so a thread created under a previous child keeps its id
2152 /// here as long as the store it was persisted to is shared (#6246).
2153 async fn ensure_runtime_thread(
2154 &mut self,
2155 thread_map: &mut HashMap<String, String>,
2156 stdio_thread_id: &str,
2157 hint: Option<RuntimeThreadHint>,
2158 ) -> Result<String> {
2159 if let Some(runtime_thread_id) = thread_map.get(stdio_thread_id) {
2160 return Ok(runtime_thread_id.clone());
2161 }
2162 let hint = hint.unwrap_or_default();
2163 let runtime_thread_id = self
2164 .create_runtime_thread(hint.model, hint.workspace)
2165 .await?;
2166 thread_map.insert(stdio_thread_id.to_string(), runtime_thread_id.clone());
2167 Ok(runtime_thread_id)
2168 }
2169
2170 /// Drop a thread mapping (and its seq cursor) once no caller can name
2171 /// the client-facing key again.
2172 fn forget_thread(&mut self, thread_map: &mut HashMap<String, String>, stdio_thread_id: &str) {
2173 if let Some(runtime_thread_id) = thread_map.remove(stdio_thread_id) {
2174 self.last_seq_by_thread.remove(&runtime_thread_id);
2175 }
2176 }
2177
2178 async fn create_runtime_thread(
2179 &mut self,
2180 model: Option<String>,
2181 workspace: Option<PathBuf>,
2182 ) -> Result<String> {
2183 let record = self
2184 .request_json(
2185 self.authed(self.client.post(format!("{}/v1/threads", self.base_url)))
2186 .json(&json!({
2187 "model": model,
2188 "workspace": workspace,
2189 "mode": "agent",
2190 "archived": false,
2191 })),
2192 )
2193 .await?;
2194 let thread_id = extract_runtime_thread_id(&record)?.to_string();
2195 self.last_seq_by_thread
2196 .entry(thread_id.clone())
2197 .or_insert(0);
2198 Ok(thread_id)
2199 }
2200
2201 /// Run one turn to completion, streaming its events to `writer`.
2202 ///
2203 /// `registration` is `Some` on the stdio path: it publishes the live turn
2204 /// so an `thread/interrupt` arriving mid-stream can reach the runtime
2205 /// without waiting on the bridge mutex this call holds.
2206 async fn message_thread<W: AsyncWrite + Unpin>(
2207 &mut self,
2208 thread_id: &str,
2209 input: RuntimeTurnInput<'_>,
2210 writer: &mut W,
2211 registration: Option<(TurnRegistry, String)>,
2212 mut transcript: Option<&mut TurnTranscript>,
2213 ) -> Result<Value> {
2214 let RuntimeTurnInput {
2215 input,
2216 images,
2217 max_output_tokens,
2218 expected_workspace,
2219 } = input;
2220 let mut request = json!({ "prompt": input });
2221 if let Some(workspace) = expected_workspace {
2222 request["expected_workspace"] = json!(workspace);
2223 }
2224 if !images.is_empty() {
2225 let info = self
2226 .request_json(
2227 self.authed(
2228 self.client
2229 .get(format!("{}/v1/runtime/info", self.base_url)),
2230 ),
2231 )
2232 .await?;
2233 if info
2234 .pointer("/capabilities/turn_image_inputs")
2235 .and_then(Value::as_bool)
2236 != Some(true)
2237 {
2238 bail!(
2239 "Runtime image input is unavailable; update the Runtime before sending attachments"
2240 );
2241 }
2242 request["images"] = json!(images);
2243 }
2244 if let Some(limit) = max_output_tokens {
2245 request["maxOutputTokens"] = json!(limit);
2246 }
2247 let turn = self
2248 .request_json(
2249 self.authed(
2250 self.client
2251 .post(format!("{}/v1/threads/{thread_id}/turns", self.base_url)),
2252 )
2253 .json(&request),
2254 )
2255 .await?;
2256 let turn_id = turn
2257 .pointer("/turn/id")
2258 .and_then(Value::as_str)
2259 .ok_or_else(|| anyhow!("runtime API turn response missing turn.id"))?
2260 .to_string();
2261 let response_id = format!("{thread_id}:{turn_id}");
2262
2263 if let Some(transcript) = transcript.as_deref_mut() {
2264 transcript.model = turn
2265 .pointer("/thread/model")
2266 .and_then(Value::as_str)
2267 .map(str::to_string);
2268 transcript.events.push(EventFrame::ResponseStart {
2269 response_id: response_id.clone(),
2270 });
2271 }
2272
2273 emit_stdio_event(
2274 writer,
2275 json!({
2276 "type": "response_start",
2277 "response_id": response_id,
2278 }),
2279 )
2280 .await?;
2281
2282 // Publish the turn only for the streaming window, and take it back
2283 // before any `?` below: a turn that has already finished must never
2284 // look cancellable.
2285 if let Some((registry, key)) = registration.as_ref() {
2286 registry.lock().await.insert(
2287 key.clone(),
2288 InFlightTurn {
2289 base_url: self.base_url.clone(),
2290 auth_token: self.auth_token.clone(),
2291 runtime_thread_id: thread_id.to_string(),
2292 turn_id: turn_id.clone(),
2293 },
2294 );
2295 }
2296
2297 let since_seq = self.last_seq_by_thread.get(thread_id).copied().unwrap_or(0);
2298 let stream_result = self
2299 .stream_turn_events(
2300 thread_id,
2301 &turn_id,
2302 &response_id,
2303 writer,
2304 since_seq,
2305 transcript.as_deref_mut(),
2306 )
2307 .await;
2308
2309 if let Some((registry, key)) = registration.as_ref() {
2310 registry.lock().await.remove(key);
2311 }
2312
2313 if stream_result.is_ok() {
2314 let _ = emit_stdio_event(
2315 writer,
2316 json!({
2317 "type": "response_end",
2318 "response_id": response_id,
2319 }),
2320 )
2321 .await;
2322 if let Some(transcript) = transcript {
2323 transcript.events.push(EventFrame::ResponseEnd {
2324 response_id: response_id.clone(),
2325 });
2326 }
2327 } else {
2328 // The stream broke before `turn.completed` (transport error,
2329 // oversized or invalid frame, a writer that went away). Ending
2330 // the response here reported success ahead of the error, and the
2331 // runtime turn kept running with nothing left able to interrupt
2332 // it (the registry entry is gone), so the thread refused its next
2333 // message until the orphan finished. Stop it, best effort, and
2334 // let the error below be the only outcome the client sees.
2335 let orphan = InFlightTurn {
2336 base_url: self.base_url.clone(),
2337 auth_token: self.auth_token.clone(),
2338 runtime_thread_id: thread_id.to_string(),
2339 turn_id: turn_id.clone(),
2340 };
2341 if let Err(error) = interrupt_turn_request(&orphan).await {
2342 tracing::warn!(
2343 thread_id,
2344 turn_id = %turn_id,
2345 "failed to interrupt a turn whose event stream broke: {}",
2346 error.message
2347 );
2348 }
2349 }
2350
2351 let (last_seq, status, error) = stream_result?;
2352 self.last_seq_by_thread
2353 .insert(thread_id.to_string(), last_seq);
2354
2355 match status {
2356 TurnTerminalStatus::Completed => Ok(json!({
2357 "thread_id": thread_id,
2358 "status": "accepted",
2359 "thread": Value::Null,
2360 "threads": [],
2361 "model": Value::Null,
2362 "model_provider": Value::Null,
2363 "cwd": Value::Null,
2364 "approval_policy": Value::Null,
2365 "sandbox": Value::Null,
2366 "events": [],
2367 "data": { "turn_id": turn_id },
2368 })),
2369 TurnTerminalStatus::Failed => Err(anyhow!(
2370 "{}",
2371 error.unwrap_or_else(|| "turn failed".to_string())
2372 )),
2373 TurnTerminalStatus::Interrupted => Err(anyhow!(
2374 "{}",
2375 error.unwrap_or_else(|| "turn interrupted".to_string())
2376 )),
2377 TurnTerminalStatus::Canceled => Err(anyhow!(
2378 "{}",
2379 error.unwrap_or_else(|| "turn canceled".to_string())
2380 )),
2381 }
2382 }
2383
2384 async fn stream_turn_events<W: AsyncWrite + Unpin>(
2385 &self,
2386 thread_id: &str,
2387 turn_id: &str,
2388 response_id: &str,
2389 writer: &mut W,
2390 since_seq: u64,
2391 mut transcript: Option<&mut TurnTranscript>,
2392 ) -> Result<(u64, TurnTerminalStatus, Option<String>)> {
2393 let mut response = self
2394 .authed(self.client.get(format!(
2395 "{}/v1/threads/{thread_id}/events?since_seq={since_seq}",
2396 self.base_url
2397 )))
2398 .send()
2399 .await?
2400 .error_for_status()?;
2401
2402 let mut buffer = Vec::new();
2403 let mut last_seq = since_seq;
2404
2405 while let Some(chunk) = response.chunk().await? {
2406 buffer.extend_from_slice(&chunk);
2407 if buffer.len() > MAX_SSE_FRAME_BYTES {
2408 bail!(
2409 "runtime SSE frame exceeded {MAX_SSE_FRAME_BYTES} bytes without a frame delimiter"
2410 );
2411 }
2412 while let Some(frame_bytes) = take_sse_frame(&mut buffer) {
2413 let Some((event_name, frame_data)) = parse_sse_frame(&frame_bytes) else {
2414 continue;
2415 };
2416 let envelope: Value = serde_json::from_str(&frame_data)
2417 .with_context(|| format!("invalid SSE json for {event_name}: {frame_data}"))?;
2418 if let Some(seq) = envelope.get("seq").and_then(Value::as_u64) {
2419 last_seq = last_seq.max(seq);
2420 }
2421 if envelope.get("turn_id").and_then(Value::as_str) != Some(turn_id) {
2422 continue;
2423 }
2424 let payload = envelope.get("payload").cloned().unwrap_or(Value::Null);
2425 match event_name.as_str() {
2426 "item.delta" => {
2427 let kind = payload
2428 .get("kind")
2429 .and_then(Value::as_str)
2430 .unwrap_or_default();
2431 if kind == "agent_message"
2432 && let Some(delta) = payload.get("delta").and_then(Value::as_str)
2433 && !delta.is_empty()
2434 {
2435 emit_stdio_event(
2436 writer,
2437 json!({
2438 "type": "response_delta",
2439 "response_id": response_id,
2440 "delta": delta,
2441 }),
2442 )
2443 .await?;
2444 if let Some(transcript) = transcript.as_deref_mut() {
2445 transcript.text.push_str(delta);
2446 transcript.events.push(EventFrame::ResponseDelta {
2447 response_id: response_id.to_string(),
2448 delta: delta.to_string(),
2449 channel: ResponseChannel::Text,
2450 });
2451 }
2452 }
2453 }
2454 "turn.completed" => {
2455 let status = turn_terminal_status(&payload);
2456 let error = payload
2457 .pointer("/turn/error")
2458 .and_then(Value::as_str)
2459 .map(str::to_string);
2460 return Ok((last_seq, status, error));
2461 }
2462 _ => {}
2463 }
2464 }
2465 }
2466
2467 bail!("runtime event stream ended before turn.completed")
2468 }
2469
2470 #[cfg(test)]
2471 fn from_base_url_for_test(base_url: String) -> Self {
2472 install_rustls_crypto_provider();
2473 Self {
2474 base_url,
2475 client: codewhale_release::platform_http_client_builder()
2476 .timeout(Duration::from_secs(5))
2477 .build()
2478 .expect("build reqwest test client"),
2479 auth_token: None,
2480 child: None,
2481 last_seq_by_thread: HashMap::new(),
2482 }
2483 }
2484 }
2485
2486 #[cfg(test)]
2487 impl RuntimeBridge {
2488 /// Whether the historical fixture child has exited. A bridge without a child (tests,
2489 /// or an externally managed runtime) never reports exited. A `try_wait`
2490 /// error counts as exited: with `WNOHANG` it only fails when the pid is no
2491 /// longer this process's child (already reaped elsewhere), and such a
2492 /// child can neither be tracked nor killed on drop.
2493 #[cfg(unix)]
2494 fn child_exited(&mut self) -> bool {
2495 self.child
2496 .as_mut()
2497 .is_some_and(|child| !matches!(child.try_wait(), Ok(None)))
2498 }
2499
2500 /// Kills the managed runtime child and reaps it on a detached thread so
2501 /// neither an explicit shutdown nor Drop blocks a Tokio runtime thread.
2502 fn shutdown_child(&mut self) {
2503 if let Some(mut child) = self.child.take() {
2504 let _ = child.kill();
2505 std::thread::spawn(move || {
2506 let _ = child.wait();
2507 });
2508 }
2509 }
2510 }
2511
2512 #[cfg(test)]
2513 impl Drop for RuntimeBridge {
2514 fn drop(&mut self) {
2515 self.shutdown_child();
2516 }
2517 }
2518
2519 /// The line the Runtime prints once it holds its listener (the Runtime's
2520 /// `RUNTIME_LISTENING_PREFIX`; this crate does not depend on that one).
2521 #[cfg(test)]
2522 const RUNTIME_LISTENING_PREFIX: &str = "Runtime API listening on http://";
2523 /// The endpoint line is short and comes first; anything larger is not it.
2524 #[cfg(test)]
2525 const RUNTIME_READY_MAX_BYTES: usize = 1024;
2526
2527 /// The endpoint a Runtime child reports for itself: exactly a nonzero port on
2528 /// `127.0.0.1`, the host the parent asked it to bind.
2529 #[cfg(test)]
2530 fn parse_runtime_endpoint(line: &str) -> Result<std::net::SocketAddr> {
2531 let address = line
2532 .trim_end()
2533 .strip_prefix(RUNTIME_LISTENING_PREFIX)
2534 .context("runtime API bridge did not report its endpoint")?;
2535 let endpoint: std::net::SocketAddr = address
2536 .parse()
2537 .context("runtime API bridge reported an invalid endpoint")?;
2538 if endpoint.ip() != std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST) || endpoint.port() == 0
2539 {
2540 bail!("runtime API bridge reported an endpoint that is not a loopback port");
2541 }
2542 Ok(endpoint)
2543 }
2544
2545 /// Read the first stdout line, bounded.
2546 #[cfg(test)]
2547 fn read_runtime_ready_line(stdout: &mut impl std::io::Read) -> Result<String> {
2548 let mut line = Vec::new();
2549 let mut byte = [0u8; 1];
2550 loop {
2551 if stdout.read(&mut byte)? == 0 {
2552 bail!("runtime API bridge closed stdout before reporting its endpoint");
2553 }
2554 if byte[0] == b'\n' {
2555 break;
2556 }
2557 if line.len() >= RUNTIME_READY_MAX_BYTES {
2558 bail!("runtime API bridge sent an oversized readiness line");
2559 }
2560 line.push(byte[0]);
2561 }
2562 String::from_utf8(line).context("runtime API bridge readiness line is not UTF-8")
2563 }
2564
2565 fn install_rustls_crypto_provider() {
2566 let _ = rustls::crypto::ring::default_provider().install_default();
2567 }
2568
2569 fn extract_runtime_thread_id(record: &Value) -> Result<&str> {
2570 record
2571 .get("id")
2572 .and_then(Value::as_str)
2573 .ok_or_else(|| anyhow!("runtime API thread response missing id"))
2574 }
2575
2576 fn turn_terminal_status(payload: &Value) -> TurnTerminalStatus {
2577 match payload
2578 .pointer("/turn/status")
2579 .and_then(Value::as_str)
2580 .unwrap_or("completed")
2581 .to_ascii_lowercase()
2582 .as_str()
2583 {
2584 "failed" => TurnTerminalStatus::Failed,
2585 "interrupted" => TurnTerminalStatus::Interrupted,
2586 "canceled" | "cancelled" => TurnTerminalStatus::Canceled,
2587 _ => TurnTerminalStatus::Completed,
2588 }
2589 }
2590
2591 async fn emit_stdio_event<W: AsyncWrite + Unpin>(writer: &mut W, event: Value) -> Result<()> {
2592 writer.write_all(&serde_json::to_vec(&event)?).await?;
2593 writer.write_all(b"\n").await?;
2594 writer.flush().await?;
2595 Ok(())
2596 }
2597
2598 fn take_sse_frame(buffer: &mut Vec<u8>) -> Option<Vec<u8>> {
2599 if let Some(pos) = buffer.windows(4).position(|window| window == b"\r\n\r\n") {
2600 return Some(buffer.drain(..pos + 4).collect());
2601 }
2602 buffer
2603 .windows(2)
2604 .position(|window| window == b"\n\n")
2605 .map(|pos| buffer.drain(..pos + 2).collect())
2606 }
2607
2608 fn parse_sse_frame(frame_bytes: &[u8]) -> Option<(String, String)> {
2609 let text = String::from_utf8(frame_bytes.to_vec()).ok()?;
2610 let mut event_name = None;
2611 let mut data_lines = Vec::new();
2612 for raw_line in text.lines() {
2613 let line = raw_line.trim_end_matches('\r');
2614 if let Some(value) = line.strip_prefix("event:") {
2615 event_name = Some(value.trim().to_string());
2616 } else if let Some(value) = line.strip_prefix("data:") {
2617 data_lines.push(value.trim_start().to_string());
2618 }
2619 }
2620 match (event_name, data_lines.is_empty()) {
2621 (Some(event), false) => Some((event, data_lines.join("\n"))),
2622 _ => None,
2623 }
2624 }
2625
2626 #[cfg(test)]
2627 async fn dispatch_stdio_request(
2628 state: &AppState,
2629 method: &str,
2630 params: Value,
2631 ) -> std::result::Result<StdioDispatchResult, JsonRpcError> {
2632 let mut sink = tokio::io::sink();
2633 dispatch_stdio_request_with_writer(state, &mut sink, method, params, AppTransport::Stdio).await
2634 }
2635
2636 async fn dispatch_stdio_app_request(
2637 state: &AppState,
2638 request: AppRequest,
2639 transport: AppTransport,
2640 ) -> std::result::Result<StdioDispatchResult, JsonRpcError> {
2641 let response = Box::pin(process_app_request(state, request, transport)).await;
2642 Ok(StdioDispatchResult {
2643 result: serde_json::to_value(response)
2644 .map_err(|err| JsonRpcError::internal(err.to_string()))?,
2645 should_exit: false,
2646 })
2647 }
2648
2649 async fn dispatch_stdio_request_with_writer<W: AsyncWrite + Unpin>(
2650 state: &AppState,
2651 writer: &mut W,
2652 method: &str,
2653 params: Value,
2654 transport: AppTransport,
2655 ) -> std::result::Result<StdioDispatchResult, JsonRpcError> {
2656 let outcome = match method {
2657 "healthz" | "app/healthz" => StdioDispatchResult {
2658 result: json!({
2659 "status": "ok",
2660 "service": legacy_deepseek_compat::SERVICE_NAME,
2661 "transport": transport.label()
2662 }),
2663 should_exit: false,
2664 },
2665 "capabilities" => {
2666 let mut methods = vec![
2667 "healthz",
2668 "thread/capabilities",
2669 "thread/request",
2670 "thread/create",
2671 "thread/start",
2672 "thread/resume",
2673 "thread/fork",
2674 "thread/list",
2675 "thread/read",
2676 "thread/set_name",
2677 "thread/goal/set",
2678 "thread/goal/get",
2679 "thread/goal/clear",
2680 "thread/archive",
2681 "thread/unarchive",
2682 "thread/message",
2683 "thread/interrupt",
2684 "app/capabilities",
2685 "app/request",
2686 "app/config/get",
2687 "app/config/set",
2688 "app/config/unset",
2689 "app/config/list",
2690 "app/config/reload",
2691 "app/models",
2692 "app/thread_loaded_list",
2693 "prompt/capabilities",
2694 "prompt/request",
2695 "prompt/run",
2696 "shutdown",
2697 ];
2698 if transport == AppTransport::Socket {
2699 // The daemon handshake exists only on the socket transport;
2700 // stdio/HTTP clients never see it, so the stdio pin is unchanged.
2701 methods.insert(1, daemon_socket::ATTACH_METHOD);
2702 }
2703 StdioDispatchResult {
2704 result: json!({
2705 "transport": transport.label(),
2706 "families": ["thread/*", "app/*", "prompt/*"],
2707 "turn_image_inputs": true,
2708 "methods": methods,
2709 }),
2710 should_exit: false,
2711 }
2712 }
2713 "thread/capabilities" => StdioDispatchResult {
2714 result: json!({
2715 "turn_image_inputs": true,
2716 "methods": [
2717 "thread/request",
2718 "thread/create",
2719 "thread/start",
2720 "thread/resume",
2721 "thread/fork",
2722 "thread/list",
2723 "thread/read",
2724 "thread/set_name",
2725 "thread/goal/set",
2726 "thread/goal/get",
2727 "thread/goal/clear",
2728 "thread/archive",
2729 "thread/unarchive",
2730 "thread/message",
2731 "thread/interrupt"
2732 ]
2733 }),
2734 should_exit: false,
2735 },
2736 "thread/request" => {
2737 let request: ThreadRequest = parse_params(params)?;
2738 if let ThreadRequest::Message {
2739 thread_id,
2740 input,
2741 images,
2742 max_output_tokens,
2743 } = request
2744 {
2745 let response = handle_stdio_thread_message(
2746 state,
2747 writer,
2748 ThreadMessageParams {
2749 thread_id,
2750 input,
2751 images,
2752 max_output_tokens,
2753 },
2754 )
2755 .await?;
2756 return Ok(StdioDispatchResult {
2757 result: response,
2758 should_exit: false,
2759 });
2760 }
2761 let should_record_hint = matches!(
2762 &request,
2763 ThreadRequest::Create { .. }
2764 | ThreadRequest::Start(_)
2765 | ThreadRequest::Resume(_)
2766 | ThreadRequest::Fork(_)
2767 );
2768 let response = handle_thread_request(state, request).await?;
2769 if should_record_hint {
2770 record_stdio_thread_hint(state, &response).await;
2771 }
2772 StdioDispatchResult {
2773 result: serde_json::to_value(response)
2774 .map_err(|err| JsonRpcError::internal(err.to_string()))?,
2775 should_exit: false,
2776 }
2777 }
2778 "thread/create" => {
2779 #[derive(Debug, Deserialize)]
2780 struct CreateParams {
2781 #[serde(default)]
2782 metadata: Value,
2783 }
2784 let parsed: CreateParams = parse_params(params_or_object(params))?;
2785 let response = handle_thread_request(
2786 state,
2787 ThreadRequest::Create {
2788 metadata: parsed.metadata,
2789 },
2790 )
2791 .await?;
2792 record_stdio_thread_hint(state, &response).await;
2793 StdioDispatchResult {
2794 result: serde_json::to_value(response)
2795 .map_err(|err| JsonRpcError::internal(err.to_string()))?,
2796 should_exit: false,
2797 }
2798 }
2799 "thread/start" => {
2800 let request = ThreadRequest::Start(parse_params(params_or_object(params))?);
2801 let response = handle_thread_request(state, request).await?;
2802 record_stdio_thread_hint(state, &response).await;
2803 StdioDispatchResult {
2804 result: serde_json::to_value(response)
2805 .map_err(|err| JsonRpcError::internal(err.to_string()))?,
2806 should_exit: false,
2807 }
2808 }
2809 "thread/resume" => {
2810 let request = ThreadRequest::Resume(parse_params(params_or_object(params))?);
2811 let response = handle_thread_request(state, request).await?;
2812 ensure_thread_found(&response)?;
2813 record_stdio_thread_hint(state, &response).await;
2814 StdioDispatchResult {
2815 result: serde_json::to_value(response)
2816 .map_err(|err| JsonRpcError::internal(err.to_string()))?,
2817 should_exit: false,
2818 }
2819 }
2820 "thread/fork" => {
2821 let request = ThreadRequest::Fork(parse_params(params_or_object(params))?);
2822 let response = handle_thread_request(state, request).await?;
2823 ensure_thread_found(&response)?;
2824 record_stdio_thread_hint(state, &response).await;
2825 StdioDispatchResult {
2826 result: serde_json::to_value(response)
2827 .map_err(|err| JsonRpcError::internal(err.to_string()))?,
2828 should_exit: false,
2829 }
2830 }
2831 "thread/list" => {
2832 let request = ThreadRequest::List(parse_params(params_or_object(params))?);
2833 let response = handle_thread_request(state, request).await?;
2834 StdioDispatchResult {
2835 result: serde_json::to_value(response)
2836 .map_err(|err| JsonRpcError::internal(err.to_string()))?,
2837 should_exit: false,
2838 }
2839 }
2840 "thread/read" => {
2841 let request = ThreadRequest::Read(parse_params(params_or_object(params))?);
2842 let response = handle_thread_request(state, request).await?;
2843 StdioDispatchResult {
2844 result: serde_json::to_value(response)
2845 .map_err(|err| JsonRpcError::internal(err.to_string()))?,
2846 should_exit: false,
2847 }
2848 }
2849 "thread/set_name" | "thread/set-name" => {
2850 let request = ThreadRequest::SetName(parse_params(params_or_object(params))?);
2851 let response = handle_thread_request(state, request).await?;
2852 StdioDispatchResult {
2853 result: serde_json::to_value(response)
2854 .map_err(|err| JsonRpcError::internal(err.to_string()))?,
2855 should_exit: false,
2856 }
2857 }
2858 "thread/goal/set" | "thread/goal_set" | "thread/goal-set" => {
2859 let request = ThreadRequest::GoalSet(parse_params::<ThreadGoalSetParams>(
2860 params_or_object(params),
2861 )?);
2862 let response = handle_thread_request(state, request).await?;
2863 StdioDispatchResult {
2864 result: serde_json::to_value(response)
2865 .map_err(|err| JsonRpcError::internal(err.to_string()))?,
2866 should_exit: false,
2867 }
2868 }
2869 "thread/goal/get" | "thread/goal_get" | "thread/goal-get" => {
2870 let request = ThreadRequest::GoalGet(parse_params::<ThreadGoalGetParams>(
2871 params_or_object(params),
2872 )?);
2873 let response = handle_thread_request(state, request).await?;
2874 StdioDispatchResult {
2875 result: serde_json::to_value(response)
2876 .map_err(|err| JsonRpcError::internal(err.to_string()))?,
2877 should_exit: false,
2878 }
2879 }
2880 "thread/goal/clear" | "thread/goal_clear" | "thread/goal-clear" => {
2881 let request = ThreadRequest::GoalClear(parse_params::<ThreadGoalClearParams>(
2882 params_or_object(params),
2883 )?);
2884 let response = handle_thread_request(state, request).await?;
2885 StdioDispatchResult {
2886 result: serde_json::to_value(response)
2887 .map_err(|err| JsonRpcError::internal(err.to_string()))?,
2888 should_exit: false,
2889 }
2890 }
2891 "thread/archive" => {
2892 let parsed: ThreadIdParams = parse_params(params_or_object(params))?;
2893 let response = handle_thread_request(
2894 state,
2895 ThreadRequest::Archive {
2896 thread_id: parsed.thread_id,
2897 },
2898 )
2899 .await?;
2900 ensure_thread_found(&response)?;
2901 StdioDispatchResult {
2902 result: serde_json::to_value(response)
2903 .map_err(|err| JsonRpcError::internal(err.to_string()))?,
2904 should_exit: false,
2905 }
2906 }
2907 "thread/unarchive" => {
2908 let parsed: ThreadIdParams = parse_params(params_or_object(params))?;
2909 let response = handle_thread_request(
2910 state,
2911 ThreadRequest::Unarchive {
2912 thread_id: parsed.thread_id,
2913 },
2914 )
2915 .await?;
2916 ensure_thread_found(&response)?;
2917 StdioDispatchResult {
2918 result: serde_json::to_value(response)
2919 .map_err(|err| JsonRpcError::internal(err.to_string()))?,
2920 should_exit: false,
2921 }
2922 }
2923 "thread/message" => {
2924 let parsed: ThreadMessageParams = parse_params(params_or_object(params))?;
2925 let response = handle_stdio_thread_message(state, writer, parsed).await?;
2926 StdioDispatchResult {
2927 result: response,
2928 should_exit: false,
2929 }
2930 }
2931 "app/capabilities" => {
2932 dispatch_stdio_app_request(state, AppRequest::Capabilities, transport).await?
2933 }
2934 "app/request" => {
2935 let request: AppRequest = parse_params(params)?;
2936 dispatch_stdio_app_request(state, request, transport).await?
2937 }
2938 "app/config/get" => {
2939 let parsed: ConfigGetParams = parse_params(params_or_object(params))?;
2940 dispatch_stdio_app_request(state, AppRequest::ConfigGet { key: parsed.key }, transport)
2941 .await?
2942 }
2943 "app/config/set" => {
2944 let parsed: ConfigSetParams = parse_params(params_or_object(params))?;
2945 dispatch_stdio_app_request(
2946 state,
2947 AppRequest::ConfigSet {
2948 key: parsed.key,
2949 value: parsed.value,
2950 },
2951 transport,
2952 )
2953 .await?
2954 }
2955 "app/config/unset" => {
2956 let parsed: ConfigGetParams = parse_params(params_or_object(params))?;
2957 dispatch_stdio_app_request(
2958 state,
2959 AppRequest::ConfigUnset { key: parsed.key },
2960 transport,
2961 )
2962 .await?
2963 }
2964 "app/config/list" => {
2965 dispatch_stdio_app_request(state, AppRequest::ConfigList, transport).await?
2966 }
2967 "app/config/reload" => {
2968 dispatch_stdio_app_request(state, AppRequest::ConfigReload, transport).await?
2969 }
2970 "app/models" => dispatch_stdio_app_request(state, AppRequest::Models, transport).await?,
2971 "app/thread_loaded_list" | "app/thread-loaded-list" => {
2972 dispatch_stdio_app_request(state, AppRequest::ThreadLoadedList, transport).await?
2973 }
2974 "prompt/capabilities" => StdioDispatchResult {
2975 result: json!({
2976 "methods": ["prompt/request", "prompt/run"]
2977 }),
2978 should_exit: false,
2979 },
2980 "prompt/request" | "prompt/run" => {
2981 let request: PromptRequest = parse_params(params)?;
2982 let response = handle_prompt_request(state, writer, request).await?;
2983 StdioDispatchResult {
2984 result: serde_json::to_value(response)
2985 .map_err(|err| JsonRpcError::internal(err.to_string()))?,
2986 should_exit: false,
2987 }
2988 }
2989 "thread/interrupt" => {
2990 let parsed: ThreadInterruptParams = parse_params(params_or_object(params))?;
2991 let interrupted = interrupt_stdio_turn(state, &parsed.thread_id).await?;
2992 StdioDispatchResult {
2993 result: json!({
2994 "thread_id": parsed.thread_id,
2995 "interrupted": interrupted,
2996 }),
2997 should_exit: false,
2998 }
2999 }
3000 "shutdown" => {
3001 // The transport checks shutdown authority before dispatch.
3002 // Interrupt live turns before returning the flushed shutdown
3003 // result to its captured host. Keep the shared bridge bound while
3004 // that owner drains its manager and listeners.
3005 let _ = interrupt_all_stdio_turns(state).await;
3006 StdioDispatchResult {
3007 result: json!({"ok": true, "status": "stopped"}),
3008 should_exit: true,
3009 }
3010 }
3011 daemon_socket::ATTACH_METHOD if transport == AppTransport::Socket => {
3012 return Err(JsonRpcError::already_attached());
3013 }
3014 _ => return Err(JsonRpcError::method_not_found(method)),
3015 };
3016 Ok(outcome)
3017 }
3018
3019 async fn process_app_request(
3020 state: &AppState,
3021 req: AppRequest,
3022 _transport: AppTransport,
3023 ) -> AppResponse {
3024 match req {
3025 AppRequest::Capabilities => AppResponse {
3026 ok: true,
3027 data: json!({
3028 "routes": ADVERTISED_ROUTES,
3029 "config": ["get", "set", "unset", "list", "reload"],
3030 "events": ["response_start", "response_delta", "response_end", "tool_call_start", "tool_call_result"],
3031 "transport": "stdio+http",
3032 "config_path": state.config_path.as_ref().map(|p| p.display().to_string()),
3033 }),
3034 events: Vec::new(),
3035 },
3036 AppRequest::ConfigGet { key } => {
3037 let cfg = state.config.read().await;
3038 let value = cfg.get_display_value(&key);
3039 AppResponse {
3040 ok: true,
3041 data: json!({ "key": key, "value": value }),
3042 events: Vec::new(),
3043 }
3044 }
3045 AppRequest::ConfigSet { key, value } => {
3046 // Only propagate a mutation that actually happened. `set_value`
3047 // leaves the config untouched on an unknown key or invalid value,
3048 // so this is a no-op from the caller's point of view — but
3049 // `propagate_config` invalidates the cached stdio bridge
3050 // regardless, and dropping the last reference kills the running
3051 // child runtime along with its thread map. A single typo'd key
3052 // would orphan every in-flight thread on that bridge.
3053 let result = {
3054 let (key, value) = (key.clone(), value.clone());
3055 persist_config_mutation(state, move |cfg| cfg.set_value(&key, &value)).await
3056 };
3057 let ok = result.is_ok();
3058 let message = result.err().map(|e| e.to_string());
3059 let value = if ok {
3060 value
3061 } else {
3062 codewhale_config::persistence::redact_secrets(&value)
3063 };
3064 AppResponse {
3065 ok,
3066 data: json!({ "key": key, "value": value, "error": message }),
3067 events: Vec::new(),
3068 }
3069 }
3070 AppRequest::ConfigUnset { key } => {
3071 // See ConfigSet: a failed unset changed nothing and must not tear
3072 // down the runtime bridge.
3073 let result = {
3074 let key = key.clone();
3075 persist_config_mutation(state, move |cfg| cfg.unset_value(&key)).await
3076 };
3077 let ok = result.is_ok();
3078 let message = result.err().map(|e| e.to_string());
3079 AppResponse {
3080 ok,
3081 data: json!({ "key": key, "error": message }),
3082 events: Vec::new(),
3083 }
3084 }
3085 AppRequest::ConfigList => {
3086 let cfg = state.config.read().await;
3087 AppResponse {
3088 ok: true,
3089 data: json!({ "values": cfg.list_values() }),
3090 events: Vec::new(),
3091 }
3092 }
3093 AppRequest::ConfigReload => {
3094 // Re-read both `config.toml` and the sibling `permissions.toml`
3095 // from disk (the headless equivalent of the TUI
3096 // `reload_runtime_config` codepath) and push the fresh
3097 // snapshots into `state.config` and the live `Runtime`.
3098 //
3099 // `ConfigStore::load` resolves the same default config path
3100 // that `build_state` used at startup when `config_path` is
3101 // `None`, so a `None` here reloads from the same on-disk file
3102 // the server booted from.
3103 // Disk is already the source of truth here, so nothing to
3104 // persist. External `permissions.toml` edits reach the Engine
3105 // because the update invalidates the runtime bridge; the next
3106 // turn's child loads both files fresh.
3107 if let Err(error) = update_config_store(state, |_| Ok(())).await {
3108 return AppResponse {
3109 ok: false,
3110 data: json!({ "error": error.to_string() }),
3111 events: Vec::new(),
3112 };
3113 }
3114
3115 AppResponse {
3116 ok: true,
3117 data: json!({ "reloaded": true }),
3118 events: Vec::new(),
3119 }
3120 }
3121 AppRequest::Models => AppResponse {
3122 ok: true,
3123 data: json!({ "models": state.registry.list() }),
3124 events: Vec::new(),
3125 },
3126 AppRequest::ThreadLoadedList => {
3127 let response = handle_thread_request(
3128 state,
3129 ThreadRequest::List(ThreadListParams {
3130 include_archived: false,
3131 limit: Some(50),
3132 }),
3133 )
3134 .await;
3135 match response {
3136 Ok(thread_resp) => AppResponse {
3137 ok: true,
3138 data: json!({ "threads": thread_resp.threads }),
3139 events: thread_resp.events,
3140 },
3141 Err(err) => AppResponse {
3142 ok: false,
3143 data: json!({ "error": err.message }),
3144 events: Vec::new(),
3145 },
3146 }
3147 }
3148 AppRequest::SubmitUserInput { request_id, .. } => {
3149 // This transport cannot deliver a clarification answer, and
3150 // saying otherwise was the bug: the previous implementation
3151 // reported `resolved: true` and filed the answers in a map with
3152 // no reader anywhere in this crate.
3153 //
3154 // It cannot be made to work here. `handle_line_during_turn`
3155 // executes exactly one method while a turn is streaming —
3156 // `thread/interrupt`. Everything else, `app/request` included,
3157 // queues until the turn ends, so an answer sent over this
3158 // transport would wait on the very turn that is waiting for it.
3159 // The runtime API owns the pending request and can resume the
3160 // turn, so that is where the reply belongs.
3161 AppResponse {
3162 ok: false,
3163 data: json!({
3164 "error": "user_input_reply_unsupported",
3165 "request_id": request_id,
3166 "message": concat!(
3167 "the app-server control transport cannot deliver clarification answers: ",
3168 "only `thread/interrupt` runs while a turn is streaming, so an answer sent ",
3169 "here would queue behind the turn waiting for it. Reply on the runtime API ",
3170 "instead: POST /v1/user-input/{thread_id}/{request_id}."
3171 ),
3172 }),
3173 events: Vec::new(),
3174 }
3175 }
3176 }
3177 }
3178
3179 /// Install the saved config in the bookkeeping Runtime, then apply it through
3180 /// the already captured canonical owner. Unbound production frontends refuse
3181 /// application while retaining the saved config for explicit recovery.
3182 /// A saved config whose owner reload fails is reported as retained but unapplied.
3183 async fn propagate_config(state: &AppState) -> Result<()> {
3184 {
3185 let mut runtime = state.runtime.write().await;
3186 let snapshot = state.config.read().await.clone();
3187 runtime.update_config(snapshot);
3188 }
3189 if state.captured_owner.is_some() {
3190 let shared = acquire_runtime_bridge(state)
3191 .await
3192 .map_err(|_| anyhow!("config saved but captured Runtime owner is unavailable"))?;
3193 let request = {
3194 let bridge = shared.lock().await;
3195 bridge
3196 .authed(
3197 bridge
3198 .client
3199 .post(format!("{}/v1/config/reload", bridge.base_url)),
3200 )
3201 .timeout(Duration::from_secs(10))
3202 };
3203 let response = request
3204 .send()
3205 .await
3206 .context("config saved; captured owner reload outcome is uncertain, not replayed")?;
3207 anyhow::ensure!(
3208 response.status().is_success(),
3209 "config saved but captured owner rejected reload (status {})",
3210 response.status()
3211 );
3212 } else {
3213 #[cfg(test)]
3214 invalidate_runtime_bridge(state).await;
3215 #[cfg(not(test))]
3216 bail!("config saved but no captured Runtime owner can apply it");
3217 }
3218 Ok(())
3219 }
3220
3221 /// Prefix of a config error that is the server's fault (the file could not
3222 /// be read, parsed, or written), reported over HTTP `/app` as a 500 rather
3223 /// than the 400 a rejected key or value gets.
3224 const CONFIG_LOAD_ERROR: &str = "failed to load config";
3225 /// See [`CONFIG_LOAD_ERROR`].
3226 const CONFIG_SAVE_ERROR: &str = "failed to save config";
3227
3228 /// Apply `mutate` to the config on disk, then propagate the saved result.
3229 ///
3230 /// The mutation runs against a freshly loaded store, not the in-memory
3231 /// snapshot, so edits another process (TUI, `codewhale login`) saved since
3232 /// startup survive. With no explicit `--config` the store resolves the same
3233 /// default path the runtime child reads, so the change reaches turns and
3234 /// survives a restart. Any load, mutation, or save failure is returned and
3235 /// nothing is propagated, so the caller never reports `ok` for a change
3236 /// that was not kept.
3237 ///
3238 /// The work runs on its own task: once the file is written, a caller that
3239 /// goes away (an HTTP client disconnect) must not leave disk ahead of
3240 /// `state.config`, the live runtime, and the cached bridge.
3241 async fn persist_config_mutation(
3242 state: &AppState,
3243 mutate: impl FnOnce(&mut codewhale_config::ConfigToml) -> Result<()> + Send + 'static,
3244 ) -> Result<()> {
3245 update_config_store(state, move |store| {
3246 mutate(&mut store.config)?;
3247 store
3248 .save()
3249 .map_err(|err| anyhow!("{CONFIG_SAVE_ERROR}: {err}"))
3250 })
3251 .await
3252 }
3253
3254 /// Serialize reloads and writes before reading disk, and finish propagation
3255 /// even if the requesting connection disappears. A reload mutates nothing.
3256 async fn update_config_store(
3257 state: &AppState,
3258 update: impl FnOnce(&mut ConfigStore) -> Result<()> + Send + 'static,
3259 ) -> Result<()> {
3260 let state = state.clone();
3261 tokio::spawn(async move {
3262 // Own the write guard across load→mutate→save→install so two
3263 // concurrent mutations cannot overwrite one another. All disk work
3264 // runs off-runtime; the owned operation still finishes propagation
3265 // when the requesting connection goes away.
3266 let mut config = state.config.clone().write_owned().await;
3267 let config_path = state.config_path.clone();
3268 tokio::task::spawn_blocking(move || -> Result<()> {
3269 let mut store = ConfigStore::load(config_path)
3270 .map_err(|err| anyhow!("{CONFIG_LOAD_ERROR}: {err}"))?;
3271 update(&mut store)?;
3272 *config = store.config;
3273 Ok(())
3274 })
3275 .await
3276 .map_err(|err| anyhow!("config store task failed: {err}"))??;
3277 propagate_config(&state).await?;
3278 Ok(())
3279 })
3280 .await
3281 .map_err(|err| anyhow!("config update task failed: {err}"))?
3282 }
3283
3284 /// Install the process-wide rustls crypto provider once for tests that build
3285 /// an HTTP client. Production installs it at startup; each test must do the
3286 /// same instead of relying on another test in the process having run first
3287 /// (nextest runs every test in its own process).
3288 #[cfg(test)]
3289 pub(crate) fn install_test_crypto_provider() {
3290 static INIT: std::sync::OnceLock<()> = std::sync::OnceLock::new();
3291 INIT.get_or_init(|| {
3292 let _ = rustls::crypto::ring::default_provider().install_default();
3293 });
3294 }
3295
3296 #[cfg(test)]
3297 mod tests {
3298 use super::*;
3299 use axum::body::{Body, to_bytes};
3300 use axum::extract::{Path as AxumPath, Query};
3301 use axum::http::header;
3302 use codewhale_protocol::AppRequest;
3303 use std::collections::HashMap;
3304 use std::fs;
3305 use tokio::io::AsyncReadExt;
3306 use tower::ServiceExt;
3307
3308 #[tokio::test]
3309 async fn control_framing_accepts_exact_limit_crlf_and_preserves_next_frame() {
3310 let mut input = vec![b'x'; MAX_RUNTIME_IMAGE_BODY_BYTES];
3311 input.extend_from_slice(b"\r\nnext\n");
3312 let mut reader = BoundedLines::new(BufReader::with_capacity(8192, input.as_slice()));
3313 assert_eq!(
3314 reader.next_line().await.unwrap().unwrap().len(),
3315 MAX_RUNTIME_IMAGE_BODY_BYTES
3316 );
3317 assert_eq!(reader.next_line().await.unwrap().as_deref(), Some("next"));
3318 assert!(reader.next_line().await.unwrap().is_none());
3319 }
3320
3321 #[tokio::test]
3322 async fn control_framing_refuses_oversize_before_retaining_extra_bytes() {
3323 let input = vec![b'x'; MAX_RUNTIME_IMAGE_BODY_BYTES + 4096];
3324 let mut reader = BoundedLines::new(BufReader::with_capacity(8192, input.as_slice()));
3325 assert_eq!(
3326 reader.next_line().await.unwrap_err().kind(),
3327 std::io::ErrorKind::InvalidData
3328 );
3329 assert!(reader.pending.len() <= MAX_RUNTIME_IMAGE_BODY_BYTES + 1);
3330 assert!(reader.pending.capacity() <= MAX_RUNTIME_IMAGE_BODY_BYTES + 1);
3331 }
3332
3333 #[tokio::test]
3334 async fn control_framing_keeps_consumed_prefix_when_read_future_is_cancelled() {
3335 let (mut input, output) = tokio::io::duplex(64);
3336 input.write_all(b"partial").await.unwrap();
3337 let mut reader = BoundedLines::new(BufReader::new(output));
3338 std::future::poll_fn(|cx| {
3339 let mut read = std::pin::pin!(reader.next_line());
3340 assert!(std::future::Future::poll(read.as_mut(), cx).is_pending());
3341 std::task::Poll::Ready(())
3342 })
3343 .await;
3344 assert_eq!(reader.pending, b"partial");
3345 input.write_all(b"-continued\n").await.unwrap();
3346 assert_eq!(
3347 reader.next_line().await.unwrap().as_deref(),
3348 Some("partial-continued")
3349 );
3350 }
3351
3352 #[tokio::test]
3353 async fn control_framing_rejects_non_utf8_and_keeps_eof_line_contract() {
3354 let mut invalid = BoundedLines::new(BufReader::new(&b"\xff\n"[..]));
3355 assert_eq!(
3356 invalid.next_line().await.unwrap_err().kind(),
3357 std::io::ErrorKind::InvalidData
3358 );
3359 let mut eof = BoundedLines::new(BufReader::new(&b"last\r"[..]));
3360 assert_eq!(eof.next_line().await.unwrap().as_deref(), Some("last\r"));
3361 assert!(eof.next_line().await.unwrap().is_none());
3362 }
3363
3364 #[test]
3365 fn retained_control_queue_refuses_exhaustion_without_admitting_more_work() {
3366 let mut pending = VecDeque::new();
3367 let mut bytes = 0;
3368 for _ in 0..64 {
3369 queue_stdio_work(
3370 &mut pending,
3371 &mut bytes,
3372 PendingStdioWork::Response(json!({})),
3373 2,
3374 )
3375 .unwrap();
3376 }
3377 let before = bytes;
3378 assert!(
3379 queue_stdio_work(
3380 &mut pending,
3381 &mut bytes,
3382 PendingStdioWork::Response(json!({})),
3383 2
3384 )
3385 .is_err()
3386 );
3387 assert_eq!(pending.len(), 64);
3388 assert_eq!(bytes, before);
3389 let mut pending = VecDeque::new();
3390 let mut bytes = 0;
3391 assert!(
3392 queue_stdio_work(
3393 &mut pending,
3394 &mut bytes,
3395 PendingStdioWork::Response(json!({})),
3396 MAX_RUNTIME_IMAGE_BODY_BYTES
3397 )
3398 .is_err()
3399 );
3400 assert!(pending.is_empty());
3401 assert_eq!(bytes, 0);
3402 }
3403
3404 #[tokio::test]
3405 async fn captured_owner_bridge_is_retained_and_never_spawns_a_replacement() {
3406 install_test_crypto_provider();
3407 let tmp = tempfile::tempdir().unwrap();
3408 let config = tmp.path().join("config.toml");
3409 fs::write(&config, "").unwrap();
3410 let mut state = build_state(Some(config), None).unwrap();
3411 state.captured_owner = Some(codewhale_protocol::RuntimeOwnerReceipt {
3412 version: 1,
3413 data_dir: tmp.path().join("runtime"),
3414 execution_scope: "fixture".into(),
3415 lease_generation: "fixture".into(),
3416 pid: std::process::id(),
3417 process_start: "fixture".into(),
3418 principal: "fixture".into(),
3419 socket_path: tmp.path().join("owner.sock"),
3420 config_path: state.config_path.clone(),
3421 });
3422 let captured = Arc::new(Mutex::new(RuntimeBridge {
3423 base_url: "http://127.0.0.1:1".into(),
3424 client: codewhale_release::tls::reqwest_client(),
3425 auth_token: None,
3426 child: None,
3427 last_seq_by_thread: HashMap::new(),
3428 }));
3429 *state.runtime_bridge.lock().await = Some(captured.clone());
3430 invalidate_runtime_bridge(&state).await;
3431 let live = acquire_runtime_bridge(&state).await.unwrap();
3432 assert!(Arc::ptr_eq(&live, &captured));
3433 *state.runtime_bridge.lock().await = None;
3434 let refused = acquire_runtime_bridge(&state).await;
3435 assert!(refused.is_err());
3436 }
3437
3438 #[tokio::test]
3439 async fn unbound_frontend_refuses_turn_without_creating_a_runtime_owner() {
3440 install_test_crypto_provider();
3441 let (state, _temporary) = capability_test_state();
3442 assert!(state.runtime_bridge.lock().await.is_none());
3443 let refused = dispatch_stdio_request(
3444 &state,
3445 "prompt/request",
3446 json!({"prompt":"must not dispatch"}),
3447 )
3448 .await
3449 .expect_err("no captured owner can execute this prompt");
3450 assert_eq!(refused.code, RUNTIME_UNAVAILABLE_CODE);
3451 assert!(refused.message.contains("refusing another owner"));
3452 assert!(state.runtime_bridge.lock().await.is_none());
3453 assert!(state.runtime_thread_map.lock().await.is_empty());
3454 assert!(state.in_flight_turns.lock().await.is_empty());
3455 }
3456
3457 #[test]
3458 fn captured_bridge_diagnostics_redact_private_authentication() {
3459 let mut bridge = RuntimeBridge::from_base_url_for_test("http://127.0.0.1:1".into());
3460 bridge.auth_token = Some("private-bridge-token-must-stay-in-memory".into());
3461 let diagnostic = format!("{bridge:?}");
3462 assert!(!diagnostic.contains("private-bridge-token"));
3463 assert!(diagnostic.contains("authenticated: true"));
3464 }
3465
3466 #[tokio::test]
3467 async fn full_control_queue_keeps_interrupt_priority_and_guest_shutdown_refusal() {
3468 use std::sync::atomic::{AtomicUsize, Ordering};
3469 async fn interrupted(State(count): State<Arc<AtomicUsize>>) -> StatusCode {
3470 count.fetch_add(1, Ordering::SeqCst);
3471 StatusCode::OK
3472 }
3473 install_test_crypto_provider();
3474 let count = Arc::new(AtomicUsize::new(0));
3475 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
3476 let address = listener.local_addr().unwrap();
3477 let app = Router::new()
3478 .route(
3479 "/v1/threads/runtime/turns/turn/interrupt",
3480 post(interrupted),
3481 )
3482 .with_state(count.clone());
3483 let server = tokio::spawn(async move {
3484 axum::serve(listener, app).await.unwrap();
3485 });
3486 let (state, _root) = capability_test_state();
3487 state.in_flight_turns.lock().await.insert(
3488 "client".into(),
3489 InFlightTurn {
3490 base_url: format!("http://{address}"),
3491 auth_token: None,
3492 runtime_thread_id: "runtime".into(),
3493 turn_id: "turn".into(),
3494 },
3495 );
3496 let mut pending = VecDeque::new();
3497 let mut bytes = 0;
3498 for _ in 0..64 {
3499 queue_stdio_work(
3500 &mut pending,
3501 &mut bytes,
3502 PendingStdioWork::Response(json!({})),
3503 2,
3504 )
3505 .unwrap();
3506 }
3507 let guest = StdioLoopPolicy {
3508 transport: AppTransport::Socket,
3509 shutdown: ShutdownAuthority::Denied,
3510 };
3511 assert!(
3512 handle_line_during_turn(
3513 &state,
3514 r#"{"id":1,"method":"shutdown","params":{}}"#,
3515 &mut pending,
3516 &mut bytes,
3517 guest
3518 )
3519 .await
3520 .is_err()
3521 );
3522 assert_eq!(
3523 count.load(Ordering::SeqCst),
3524 0,
3525 "guest must not interrupt before its denied shutdown"
3526 );
3527 assert!(
3528 handle_line_during_turn(
3529 &state,
3530 r#"{"id":2,"method":"thread/interrupt","params":{"thread_id":"client"}}"#,
3531 &mut pending,
3532 &mut bytes,
3533 guest
3534 )
3535 .await
3536 .is_err()
3537 );
3538 assert_eq!(
3539 count.load(Ordering::SeqCst),
3540 1,
3541 "interrupt acts before its overloaded reply queue"
3542 );
3543 let owner = StdioLoopPolicy {
3544 transport: AppTransport::Socket,
3545 shutdown: ShutdownAuthority::Granted,
3546 };
3547 assert!(
3548 handle_line_during_turn(
3549 &state,
3550 r#"{"id":3,"method":"shutdown","params":{}}"#,
3551 &mut pending,
3552 &mut bytes,
3553 owner
3554 )
3555 .await
3556 .is_err()
3557 );
3558 assert_eq!(
3559 count.load(Ordering::SeqCst),
3560 2,
3561 "authorized shutdown retains immediate interrupt priority"
3562 );
3563 assert_eq!(pending.len(), 64);
3564 server.abort();
3565 }
3566
3567 #[tokio::test]
3568 async fn bound_owner_config_reload_uses_private_auth_and_preserves_rejected_owner() {
3569 use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
3570 #[derive(Clone)]
3571 struct ReloadFixture {
3572 count: Arc<AtomicUsize>,
3573 refuse: Arc<AtomicBool>,
3574 }
3575 async fn reload(
3576 State(fixture): State<ReloadFixture>,
3577 headers: axum::http::HeaderMap,
3578 ) -> StatusCode {
3579 assert_eq!(
3580 headers.get(header::AUTHORIZATION).unwrap(),
3581 "Bearer private-reload-fixture"
3582 );
3583 fixture.count.fetch_add(1, Ordering::SeqCst);
3584 if fixture.refuse.load(Ordering::SeqCst) {
3585 StatusCode::BAD_REQUEST
3586 } else {
3587 StatusCode::OK
3588 }
3589 }
3590 install_test_crypto_provider();
3591 let fixture = ReloadFixture {
3592 count: Arc::new(AtomicUsize::new(0)),
3593 refuse: Arc::new(AtomicBool::new(false)),
3594 };
3595 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
3596 let address = listener.local_addr().unwrap();
3597 let app = Router::new()
3598 .route("/v1/config/reload", post(reload))
3599 .with_state(fixture.clone());
3600 let server = tokio::spawn(async move {
3601 axum::serve(listener, app).await.unwrap();
3602 });
3603 let (mut state, root) = capability_test_state();
3604 state.captured_owner = Some(codewhale_protocol::RuntimeOwnerReceipt {
3605 version: 1,
3606 data_dir: root.path().join("runtime"),
3607 execution_scope: "fixture".into(),
3608 lease_generation: "fixture".into(),
3609 pid: std::process::id(),
3610 process_start: "fixture".into(),
3611 principal: "fixture".into(),
3612 socket_path: root.path().join("owner.sock"),
3613 config_path: state.config_path.clone(),
3614 });
3615 let mut bridge = RuntimeBridge::from_base_url_for_test(format!("http://{address}"));
3616 bridge.auth_token = Some("private-reload-fixture".into());
3617 let captured = Arc::new(Mutex::new(bridge));
3618 *state.runtime_bridge.lock().await = Some(captured.clone());
3619 propagate_config(&state).await.unwrap();
3620 assert_eq!(fixture.count.load(Ordering::SeqCst), 1);
3621 fixture.refuse.store(true, Ordering::SeqCst);
3622 assert!(
3623 propagate_config(&state)
3624 .await
3625 .unwrap_err()
3626 .to_string()
3627 .contains("config saved")
3628 );
3629 assert_eq!(fixture.count.load(Ordering::SeqCst), 2);
3630 assert!(Arc::ptr_eq(
3631 state.runtime_bridge.lock().await.as_ref().unwrap(),
3632 &captured
3633 ));
3634 server.abort();
3635 }
3636
3637 fn app_with_config(auth_token: Option<&str>) -> (Router, tempfile::TempDir) {
3638 let tmp = tempfile::tempdir().expect("tempdir");
3639 let config_path = tmp.path().join("config.toml");
3640 fs::write(&config_path, "api_key = \"sk-deepseek-secret\"\n").expect("write config");
3641 let state = build_state(
3642 Some(config_path),
3643 auth_token.map(std::string::ToString::to_string),
3644 )
3645 .expect("state");
3646 (app_router(state, &[]), tmp)
3647 }
3648
3649 #[test]
3650 fn build_state_keeps_resolved_explicit_config_path() {
3651 let tmp = tempfile::tempdir().expect("tempdir");
3652 let config_dir = tmp.path().join("config-dir");
3653 fs::create_dir_all(&config_dir).expect("config dir");
3654 let config_path = config_dir.join("config.toml");
3655 fs::write(&config_path, "api_key = \"sk-deepseek-secret\"\n").expect("write config");
3656
3657 let state = build_state(Some(config_path.clone()), None).expect("state");
3658
3659 assert_eq!(
3660 state.config_path.as_deref(),
3661 Some(
3662 config_path
3663 .canonicalize()
3664 .expect("canonical config")
3665 .as_path()
3666 )
3667 );
3668 }
3669
3670 #[tokio::test]
3671 async fn stdio_transport_never_registers_the_stdout_hook_sink() {
3672 let tmp = tempfile::tempdir().expect("tempdir");
3673 let config_path = tmp.path().join("config.toml");
3674 fs::write(&config_path, "api_key = \"sk-deepseek-secret\"\n").expect("write config");
3675
3676 let http_state =
3677 build_state_with_transport(Some(config_path.clone()), None, AppTransport::Http)
3678 .expect("http state");
3679 let stdio_state = build_state_with_transport(Some(config_path), None, AppTransport::Stdio)
3680 .expect("stdio state");
3681
3682 let http_sinks = http_state.runtime.read().await.hooks.sink_count();
3683 let stdio_sinks = stdio_state.runtime.read().await.hooks.sink_count();
3684 assert_eq!(
3685 http_sinks,
3686 stdio_sinks + 1,
3687 "HTTP mode keeps StdoutHookSink + JsonlHookSink; stdio must drop the stdout sink (#5165)"
3688 );
3689 }
3690
3691 async fn response_body_json(response: Response) -> Value {
3692 let bytes = to_bytes(response.into_body(), usize::MAX)
3693 .await
3694 .expect("body bytes");
3695 serde_json::from_slice(&bytes).expect("json response")
3696 }
3697
3698 #[tokio::test]
3699 async fn http_app_routes_require_bearer_token_when_auth_enabled() {
3700 let (app, _tmp) = app_with_config(Some("test-token"));
3701 let response = app
3702 .oneshot(
3703 Request::builder()
3704 .method(Method::POST)
3705 .uri("/app")
3706 .header(header::CONTENT_TYPE, "application/json")
3707 .body(Body::from(
3708 serde_json::to_vec(&AppRequest::ConfigGet {
3709 key: "api_key".to_string(),
3710 })
3711 .expect("request json"),
3712 ))
3713 .expect("request"),
3714 )
3715 .await
3716 .expect("response");
3717
3718 assert_eq!(response.status(), StatusCode::UNAUTHORIZED);
3719 }
3720
3721 #[tokio::test]
3722 async fn http_config_get_redacts_sensitive_values_after_auth() {
3723 let (app, _tmp) = app_with_config(Some("test-token"));
3724 let response = app
3725 .oneshot(
3726 Request::builder()
3727 .method(Method::POST)
3728 .uri("/app")
3729 .header(header::AUTHORIZATION, "Bearer test-token")
3730 .header(header::CONTENT_TYPE, "application/json")
3731 .body(Body::from(
3732 serde_json::to_vec(&AppRequest::ConfigGet {
3733 key: "api_key".to_string(),
3734 })
3735 .expect("request json"),
3736 ))
3737 .expect("request"),
3738 )
3739 .await
3740 .expect("response");
3741
3742 assert_eq!(response.status(), StatusCode::OK);
3743 let body = response_body_json(response).await;
3744 assert_eq!(body["data"]["value"], "sk-d***cret");
3745 }
3746
3747 /// Every route `Capabilities` advertises must reach a handler. A POST
3748 /// with an empty JSON body is enough: a registered handler answers with
3749 /// its own status (200, or 422 for a body it rejects), while a path with
3750 /// no handler answers 404 and a wrong method 405.
3751 #[tokio::test]
3752 async fn every_advertised_route_has_a_handler() {
3753 let (_app, tmp) = app_with_config(Some("test-token"));
3754 let state = build_state(Some(tmp.path().join("config.toml")), None).expect("state");
3755 let caps = process_app_request(&state, AppRequest::Capabilities, AppTransport::Http).await;
3756 let advertised: Vec<String> =
3757 serde_json::from_value(caps.data["routes"].clone()).expect("routes list");
3758 assert_eq!(advertised, ADVERTISED_ROUTES);
3759
3760 for route in &advertised {
3761 let mut status = StatusCode::METHOD_NOT_ALLOWED;
3762 for method in [Method::POST, Method::GET] {
3763 let (app, _tmp) = app_with_config(Some("test-token"));
3764 let body = if method == Method::POST {
3765 Body::from("{}")
3766 } else {
3767 Body::empty()
3768 };
3769 status = app
3770 .oneshot(
3771 Request::builder()
3772 .method(method)
3773 .uri(route.as_str())
3774 .header(header::AUTHORIZATION, "Bearer test-token")
3775 .header(header::CONTENT_TYPE, "application/json")
3776 .body(body)
3777 .expect("request"),
3778 )
3779 .await
3780 .expect("response")
3781 .status();
3782 if status != StatusCode::METHOD_NOT_ALLOWED {
3783 break;
3784 }
3785 }
3786 assert!(
3787 status != StatusCode::NOT_FOUND
3788 && status != StatusCode::METHOD_NOT_ALLOWED
3789 && !status.is_server_error(),
3790 "advertised route {route} has no working handler: {status}"
3791 );
3792 }
3793 }
3794
3795 /// `/tool` ran calls against an empty tool registry under an approval
3796 /// mapping of its own. It is gone from both the router and the
3797 /// advertised list; tools run only inside Engine turns.
3798 #[tokio::test]
3799 async fn tool_route_is_not_served_or_advertised() {
3800 let (app, tmp) = app_with_config(Some("test-token"));
3801 let response = app
3802 .oneshot(
3803 Request::builder()
3804 .method(Method::POST)
3805 .uri("/tool")
3806 .header(header::AUTHORIZATION, "Bearer test-token")
3807 .header(header::CONTENT_TYPE, "application/json")
3808 .body(Body::from(
3809 r#"{"call":{"name":"exec_shell","payload":{"type":"local_shell","command":["true"]}}}"#,
3810 ))
3811 .expect("request"),
3812 )
3813 .await
3814 .expect("response");
3815 assert_eq!(response.status(), StatusCode::NOT_FOUND);
3816
3817 let state = build_state(Some(tmp.path().join("config.toml")), None).expect("state");
3818 let caps = process_app_request(&state, AppRequest::Capabilities, AppTransport::Http).await;
3819 let advertised = caps.data["routes"].as_array().expect("routes list");
3820 assert!(!advertised.iter().any(|route| route == "/tool"));
3821 }
3822
3823 #[tokio::test]
3824 async fn mcp_startup_route_cannot_start_a_parallel_pool() {
3825 let (app, tmp) = app_with_config(Some("test-token"));
3826 let response = app
3827 .oneshot(
3828 Request::builder()
3829 .method(Method::POST)
3830 .uri("/mcp/startup")
3831 .header(header::AUTHORIZATION, "Bearer test-token")
3832 .body(Body::empty())
3833 .expect("request"),
3834 )
3835 .await
3836 .expect("response");
3837 assert_eq!(response.status(), StatusCode::NOT_FOUND);
3838
3839 let state = build_state(Some(tmp.path().join("config.toml")), None).expect("state");
3840 let caps = process_app_request(&state, AppRequest::Capabilities, AppTransport::Http).await;
3841 assert!(
3842 !caps.data["routes"]
3843 .as_array()
3844 .expect("routes")
3845 .iter()
3846 .any(|route| route == "/mcp/startup")
3847 );
3848 assert!(
3849 !caps.data["events"]
3850 .as_array()
3851 .expect("events")
3852 .iter()
3853 .any(|event| event == "mcp_startup_update" || event == "mcp_startup_complete")
3854 );
3855 }
3856
3857 #[tokio::test]
3858 async fn cors_does_not_allow_arbitrary_origins() {
3859 let (app, _tmp) = app_with_config(Some("test-token"));
3860 let response = app
3861 .oneshot(
3862 Request::builder()
3863 .method(Method::GET)
3864 .uri("/healthz")
3865 .header(header::ORIGIN, "https://attacker.example")
3866 .body(Body::empty())
3867 .expect("request"),
3868 )
3869 .await
3870 .expect("response");
3871
3872 assert_eq!(response.status(), StatusCode::OK);
3873 assert!(
3874 response
3875 .headers()
3876 .get(header::ACCESS_CONTROL_ALLOW_ORIGIN)
3877 .is_none()
3878 );
3879 }
3880
3881 #[tokio::test]
3882 async fn config_reload_refreshes_runtime_config_from_disk() {
3883 let tmp = tempfile::tempdir().expect("tempdir");
3884 let config_path = tmp.path().join("config.toml");
3885 fs::write(
3886 &config_path,
3887 "api_key = \"sk-deepseek-secret\"\nmodel = \"deepseek-chat\"\n",
3888 )
3889 .expect("write config");
3890 let state = build_state(Some(config_path.clone()), None).expect("state");
3891 {
3892 let runtime = state.runtime.read().await;
3893 assert_eq!(runtime.config.model.as_deref(), Some("deepseek-chat"));
3894 }
3895
3896 fs::write(
3897 &config_path,
3898 "api_key = \"sk-deepseek-secret\"\nmodel = \"deepseek-reasoner\"\n",
3899 )
3900 .expect("rewrite config");
3901
3902 // ConfigReload must re-read the file and push it into the live
3903 // Runtime without a restart.
3904 let response =
3905 process_app_request(&state, AppRequest::ConfigReload, AppTransport::Stdio).await;
3906 assert!(response.ok, "reload should succeed");
3907 assert_eq!(response.data["reloaded"], true);
3908
3909 // The shared config lock reflects the new model.
3910 {
3911 let cfg = state.config.read().await;
3912 assert_eq!(cfg.model.as_deref(), Some("deepseek-reasoner"));
3913 }
3914 // The live Runtime reflects the new model.
3915 {
3916 let runtime = state.runtime.read().await;
3917 assert_eq!(runtime.config.model.as_deref(), Some("deepseek-reasoner"));
3918 }
3919 }
3920
3921 #[tokio::test]
3922 async fn config_set_propagates_to_runtime_config() {
3923 let tmp = tempfile::tempdir().expect("tempdir");
3924 let config_path = tmp.path().join("config.toml");
3925 fs::write(
3926 &config_path,
3927 "api_key = \"sk-deepseek-secret\"\nmodel = \"deepseek-chat\"\n",
3928 )
3929 .expect("write config");
3930 let state = build_state(Some(config_path.clone()), None).expect("state");
3931
3932 // Set a new model via the API.
3933 let response = process_app_request(
3934 &state,
3935 AppRequest::ConfigSet {
3936 key: "model".to_string(),
3937 value: "deepseek-reasoner".to_string(),
3938 },
3939 AppTransport::Stdio,
3940 )
3941 .await;
3942 assert!(response.ok, "set should succeed");
3943
3944 // Live runtime sees the new model.
3945 {
3946 let runtime = state.runtime.read().await;
3947 assert_eq!(runtime.config.model.as_deref(), Some("deepseek-reasoner"));
3948 }
3949 // The on-disk file was persisted.
3950 let persisted = fs::read_to_string(&config_path).expect("read config");
3951 assert!(persisted.contains("deepseek-reasoner"));
3952 }
3953
3954 /// A bridge stand-in with no child process: this test only cares about
3955 /// whether the cache slot survives, not about talking to a runtime.
3956 fn sentinel_bridge() -> SharedRuntimeBridge {
3957 Arc::new(Mutex::new(RuntimeBridge {
3958 base_url: "http://127.0.0.1:0".to_string(),
3959 client: codewhale_release::tls::reqwest_client(),
3960 auth_token: None,
3961 child: None,
3962 last_seq_by_thread: HashMap::new(),
3963 }))
3964 }
3965
3966 #[tokio::test]
3967 async fn failed_config_set_keeps_the_stdio_bridge() {
3968 crate::install_test_crypto_provider();
3969 // #4737: `set_value` rejects an invalid value before assigning, so the
3970 // request is a no-op — but `apply_config_update` ran anyway and
3971 // invalidated the cached bridge, dropping the child runtime along with
3972 // its thread map. A single bad value orphaned every in-flight stdio
3973 // thread, behind a response that correctly reported `ok: false`.
3974 //
3975 // Only `set_value` is exercised: an unknown key lands in `extras` and
3976 // succeeds, and `unset_value` has no failing input today, so its
3977 // identical guard has nothing to assert against.
3978 let tmp = tempfile::tempdir().expect("tempdir");
3979 let config_path = tmp.path().join("config.toml");
3980 fs::write(&config_path, "model = \"deepseek-chat\"\n").expect("write config");
3981 let state = build_state(Some(config_path.clone()), None).expect("state");
3982 let bridge = sentinel_bridge();
3983 *state.runtime_bridge.lock().await = Some(bridge.clone());
3984 state
3985 .runtime_thread_map
3986 .lock()
3987 .await
3988 .insert("stdio-1".to_string(), "runtime-1".to_string());
3989
3990 let disk_before = fs::read(&config_path).unwrap();
3991 let config_before = serde_json::to_value(&*state.config.read().await).unwrap();
3992 let runtime_before = serde_json::to_value(&state.runtime.read().await.config).unwrap();
3993 let token = ["sk-live-", "Z7qX4mNb2Vc9Lk3PwR8t"].concat();
3994 for (key, value) in [
3995 ("telemetry", "not-a-bool"),
3996 ("approval_policy", "ask"),
3997 ("sandbox_mode", "full"),
3998 ("verbosity", "quiet"),
3999 ("approval_policy", token.as_str()),
4000 ("sandbox_mode", token.as_str()),
4001 ("verbosity", token.as_str()),
4002 ] {
4003 let response = process_app_request(
4004 &state,
4005 AppRequest::ConfigSet {
4006 key: key.into(),
4007 value: value.into(),
4008 },
4009 AppTransport::Stdio,
4010 )
4011 .await;
4012 assert!(!response.ok, "invalid {key} must fail");
4013 assert!(response.data["error"].is_string(), "refusal detail");
4014 let rendered = serde_json::to_string(&response).unwrap();
4015 assert!(
4016 !rendered.contains(&token),
4017 "credential must not enter diagnostics or the echoed value"
4018 );
4019 assert_eq!(
4020 response.data["value"].as_str(),
4021 Some(if value == token { "[redacted]" } else { value })
4022 );
4023 assert_eq!(fs::read(&config_path).unwrap(), disk_before);
4024 assert_eq!(
4025 serde_json::to_value(&*state.config.read().await).unwrap(),
4026 config_before
4027 );
4028 assert_eq!(
4029 serde_json::to_value(&state.runtime.read().await.config).unwrap(),
4030 runtime_before
4031 );
4032 assert!(
4033 state
4034 .runtime_bridge
4035 .lock()
4036 .await
4037 .as_ref()
4038 .is_some_and(|cached| Arc::ptr_eq(cached, &bridge)),
4039 "the same bridge must survive a failed config/set",
4040 );
4041 assert_eq!(
4042 state
4043 .runtime_thread_map
4044 .lock()
4045 .await
4046 .get("stdio-1")
4047 .map(String::as_str),
4048 Some("runtime-1"),
4049 "the live thread map must be intact",
4050 );
4051 }
4052 }
4053
4054 #[tokio::test]
4055 async fn successful_config_set_still_invalidates_the_stdio_bridge() {
4056 crate::install_test_crypto_provider();
4057 // The other half of #4737: a mutation that *did* happen must still
4058 // rebuild the bridge, or the runtime keeps serving the old config.
4059 let tmp = tempfile::tempdir().expect("tempdir");
4060 let config_path = tmp.path().join("config.toml");
4061 fs::write(&config_path, "model = \"deepseek-chat\"\n").expect("write config");
4062 let state = build_state(Some(config_path.clone()), None).expect("state");
4063 *state.runtime_bridge.lock().await = Some(sentinel_bridge());
4064
4065 let response = process_app_request(
4066 &state,
4067 AppRequest::ConfigSet {
4068 key: "model".to_string(),
4069 value: "deepseek-reasoner".to_string(),
4070 },
4071 AppTransport::Stdio,
4072 )
4073 .await;
4074 assert!(response.ok, "valid set should succeed: {response:?}");
4075 assert_eq!(response.data["value"], "deepseek-reasoner");
4076 assert!(
4077 state.runtime_bridge.lock().await.is_none(),
4078 "a successful config change must invalidate the cached bridge",
4079 );
4080 }
4081
4082 #[tokio::test]
4083 async fn config_update_keeps_the_stdio_thread_mapping() {
4084 let (state, _tmp, server) = thread_control::compatibility_fixture().await;
4085 thread_control::resolve(&state, "canonical-1", true)
4086 .await
4087 .unwrap();
4088 let bridge = state.runtime_bridge.lock().await.clone().unwrap();
4089 let response =
4090 process_app_request(&state, AppRequest::ConfigReload, AppTransport::Stdio).await;
4091 assert!(response.ok, "captured owner reload failed: {response:?}");
4092 assert!(
4093 state
4094 .runtime_bridge
4095 .lock()
4096 .await
4097 .as_ref()
4098 .is_some_and(|captured| Arc::ptr_eq(captured, &bridge)),
4099 "reload keeps the already captured owner"
4100 );
4101 assert_eq!(
4102 state
4103 .runtime_thread_map
4104 .lock()
4105 .await
4106 .get("canonical-1")
4107 .map(String::as_str),
4108 Some("canonical-1")
4109 );
4110 let result = dispatch_stdio_request(
4111 &state,
4112 "thread/message",
4113 json!({"thread_id":"canonical-1","input":"next"}),
4114 )
4115 .await
4116 .unwrap();
4117 assert_eq!(result.result["status"], "accepted");
4118 server.abort();
4119 }
4120
4121 #[tokio::test]
4122 async fn thread_message_on_an_unknown_thread_is_not_found() {
4123 let (state, _tmp, server) = thread_control::compatibility_fixture().await;
4124 let error = dispatch_stdio_request(
4125 &state,
4126 "thread/message",
4127 json!({"thread_id":"never-created","input":"hello"}),
4128 )
4129 .await
4130 .unwrap_err();
4131 assert_eq!(error.code, THREAD_NOT_FOUND_CODE);
4132 let error = run_http_thread_message(
4133 &state,
4134 "never-created".into(),
4135 "hello".into(),
4136 Vec::new(),
4137 None,
4138 )
4139 .await
4140 .unwrap_err();
4141 assert_eq!(error.code, THREAD_NOT_FOUND_CODE);
4142 assert!(state.runtime_thread_map.lock().await.is_empty());
4143 assert!(state.in_flight_turns.lock().await.is_empty());
4144 server.abort();
4145 }
4146
4147 #[tokio::test]
4148 async fn thread_runtime_link_survives_an_app_server_restart() {
4149 let (state, tmp, server) = thread_control::compatibility_fixture().await;
4150 let workspace = state.frontend_workspace.clone().unwrap();
4151 let owner = state.captured_owner.clone();
4152 let store = state.runtime.read().await.state_store().clone();
4153 let mut metadata = codewhale_state::ThreadMetadata {
4154 cwd: workspace.clone(),
4155 ..test_client_metadata("client-1")
4156 };
4157 metadata.model_provider = "fixture-account".into();
4158 restore_test_thread_archive(&store, &metadata);
4159 dispatch_stdio_request(
4160 &state,
4161 "thread/message",
4162 json!({"thread_id":"client-1","input":"first"}),
4163 )
4164 .await
4165 .unwrap();
4166 let receipt = store
4167 .get_canonical_runtime_link("client-1", owner.as_ref().unwrap())
4168 .unwrap()
4169 .unwrap();
4170 let bridge = state.runtime_bridge.lock().await.clone();
4171 drop(state);
4172 let mut restarted = build_state(Some(tmp.path().join("config.toml")), None).unwrap();
4173 restarted.captured_owner = owner;
4174 restarted.frontend_workspace = Some(workspace);
4175 *restarted.runtime_bridge.lock().await = bridge;
4176 assert!(restarted.runtime_thread_map.lock().await.is_empty());
4177 dispatch_stdio_request(
4178 &restarted,
4179 "thread/message",
4180 json!({"thread_id":"client-1","input":"second"}),
4181 )
4182 .await
4183 .unwrap();
4184 assert_eq!(
4185 restarted
4186 .runtime_thread_map
4187 .lock()
4188 .await
4189 .get("client-1")
4190 .map(String::as_str),
4191 Some("canonical-1")
4192 );
4193 assert_eq!(
4194 store
4195 .get_canonical_runtime_link("client-1", restarted.captured_owner.as_ref().unwrap())
4196 .unwrap()
4197 .unwrap(),
4198 receipt
4199 );
4200 server.abort();
4201 }
4202
4203 #[tokio::test]
4204 async fn a_mapped_thread_without_a_saved_link_is_linked_on_its_next_turn() {
4205 let (state, _tmp, server) = thread_control::compatibility_fixture().await;
4206 let store = state.runtime.read().await.state_store().clone();
4207 restore_test_thread_archive(
4208 &store,
4209 &codewhale_state::ThreadMetadata {
4210 cwd: state.frontend_workspace.clone().unwrap(),
4211 ..test_client_metadata("client-1")
4212 },
4213 );
4214 state
4215 .runtime_thread_map
4216 .lock()
4217 .await
4218 .insert("client-1".into(), "unreceipted-cache".into());
4219 dispatch_stdio_request(
4220 &state,
4221 "thread/message",
4222 json!({"thread_id":"client-1","input":"hello"}),
4223 )
4224 .await
4225 .unwrap();
4226 let receipt = store
4227 .get_canonical_runtime_link("client-1", state.captured_owner.as_ref().unwrap())
4228 .unwrap()
4229 .unwrap();
4230 assert_eq!(receipt.runtime_thread_id, "canonical-1");
4231 assert_eq!(
4232 store
4233 .get_runtime_thread_link("client-1")
4234 .unwrap()
4235 .as_deref(),
4236 Some("canonical-1")
4237 );
4238 assert_eq!(
4239 state
4240 .runtime_thread_map
4241 .lock()
4242 .await
4243 .get("client-1")
4244 .map(String::as_str),
4245 Some("canonical-1")
4246 );
4247 server.abort();
4248 }
4249
4250 #[tokio::test]
4251 async fn a_thread_messaged_after_a_restart_starts_in_its_stored_workspace() {
4252 let (state, _tmp, server) = thread_control::compatibility_fixture().await;
4253 let workspace = state.frontend_workspace.clone().unwrap();
4254 let store = state.runtime.read().await.state_store().clone();
4255 restore_test_thread_archive(
4256 &store,
4257 &codewhale_state::ThreadMetadata {
4258 cwd: workspace.clone(),
4259 ..test_client_metadata("client-1")
4260 },
4261 );
4262 assert!(state.stdio_thread_hints.lock().await.is_empty());
4263 dispatch_stdio_request(
4264 &state,
4265 "thread/message",
4266 json!({"thread_id":"client-1","input":"hello"}),
4267 )
4268 .await
4269 .unwrap();
4270 assert_eq!(
4271 thread_control::resolve(&state, "client-1", true)
4272 .await
4273 .unwrap()
4274 .1,
4275 workspace
4276 );
4277 assert_eq!(
4278 store.get_thread("client-1").unwrap().unwrap().cwd,
4279 workspace
4280 );
4281 server.abort();
4282 }
4283
4284 #[tokio::test]
4285 async fn a_link_to_a_runtime_thread_that_is_gone_is_replaced() {
4286 // Keep the historical case identity. Missing owner checkpoints now
4287 // refuse; replacing them with an empty conversation loses history.
4288 let (state, _tmp, server) = thread_control::compatibility_fixture().await;
4289 let store = state.runtime.read().await.state_store().clone();
4290 restore_test_thread_archive(
4291 &store,
4292 &codewhale_state::ThreadMetadata {
4293 cwd: state.frontend_workspace.clone().unwrap(),
4294 ..test_client_metadata("client-1")
4295 },
4296 );
4297 rusqlite::Connection::open(store.db_path()).unwrap().execute(
4298 "INSERT INTO thread_runtime_links(thread_id,runtime_thread_id,created_at) VALUES('client-1','thr_gone',1)", []
4299 ).unwrap();
4300 let before = store.snapshot_legacy_thread_history("client-1").unwrap();
4301 let error = dispatch_stdio_request(
4302 &state,
4303 "thread/message",
4304 json!({"thread_id":"client-1","input":"hello"}),
4305 )
4306 .await
4307 .unwrap_err();
4308 assert_eq!(error.code, THREAD_NOT_FOUND_CODE);
4309 assert_eq!(
4310 store
4311 .get_runtime_thread_link("client-1")
4312 .unwrap()
4313 .as_deref(),
4314 Some("thr_gone")
4315 );
4316 assert!(
4317 store
4318 .get_canonical_runtime_link("client-1", state.captured_owner.as_ref().unwrap())
4319 .unwrap()
4320 .is_none()
4321 );
4322 assert_eq!(
4323 serde_json::to_value(store.snapshot_legacy_thread_history("client-1").unwrap())
4324 .unwrap(),
4325 serde_json::to_value(before).unwrap()
4326 );
4327 assert!(state.in_flight_turns.lock().await.is_empty());
4328 server.abort();
4329 }
4330
4331 #[tokio::test]
4332 async fn restored_thread_link_retries_validation_after_a_transient_failure() {
4333 let (state, _tmp, server) = thread_control::compatibility_fixture().await;
4334 let store = state.runtime.read().await.state_store().clone();
4335 restore_test_thread_archive(
4336 &store,
4337 &codewhale_state::ThreadMetadata {
4338 cwd: state.frontend_workspace.clone().unwrap(),
4339 ..test_client_metadata("client-1")
4340 },
4341 );
4342 thread_control::resolve(&state, "client-1", true)
4343 .await
4344 .unwrap();
4345 let receipt = store
4346 .get_canonical_runtime_link("client-1", state.captured_owner.as_ref().unwrap())
4347 .unwrap();
4348 let bridge = state.runtime_bridge.lock().await.take();
4349 let first = acquire_live_runtime_bridge(&state)
4350 .await
4351 .expect_err("transient connection refused");
4352 assert_eq!(first.code, RUNTIME_UNAVAILABLE_CODE);
4353 *state.runtime_bridge.lock().await = bridge;
4354 dispatch_stdio_request(
4355 &state,
4356 "thread/message",
4357 json!({"thread_id":"client-1","input":"retry"}),
4358 )
4359 .await
4360 .unwrap();
4361 assert_eq!(
4362 store
4363 .get_canonical_runtime_link("client-1", state.captured_owner.as_ref().unwrap())
4364 .unwrap(),
4365 receipt
4366 );
4367 server.abort();
4368 }
4369
4370 #[cfg(unix)]
4371 fn dead_bridge(base_url: &str) -> RuntimeBridge {
4372 let mut child = Command::new("true").spawn().expect("spawn fixture child");
4373 child.wait().expect("fixture child exits");
4374 let mut bridge = RuntimeBridge::from_base_url_for_test(base_url.to_string());
4375 bridge.child = Some(child);
4376 bridge
4377 }
4378
4379 #[cfg(unix)]
4380 #[tokio::test]
4381 async fn a_dead_runtime_child_is_replaced_by_a_live_one() {
4382 crate::install_test_crypto_provider();
4383 // A crashed or killed child stayed cached, so every later turn failed
4384 // against it until the app-server restarted.
4385 let (state, _tmp) = capability_test_state();
4386 let dead = Arc::new(Mutex::new(dead_bridge("http://dead.invalid")));
4387 *state.runtime_bridge.lock().await = Some(dead.clone());
4388
4389 let starts = std::sync::atomic::AtomicUsize::new(0);
4390 let bridge = acquire_live_runtime_bridge_with(&state, || {
4391 starts.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
4392 async {
4393 Ok(RuntimeBridge::from_base_url_for_test(
4394 "http://live.invalid".to_string(),
4395 ))
4396 }
4397 })
4398 .await
4399 .expect("a live replacement is returned");
4400 assert_eq!(bridge.base_url, "http://live.invalid");
4401 drop(bridge);
4402 assert_eq!(starts.load(std::sync::atomic::Ordering::SeqCst), 1);
4403 let slot = state.runtime_bridge.lock().await;
4404 let cached = slot.as_ref().expect("the replacement is cached");
4405 assert!(!Arc::ptr_eq(cached, &dead), "the dead bridge is evicted");
4406 assert_eq!(cached.lock().await.base_url, "http://live.invalid");
4407 }
4408
4409 #[cfg(unix)]
4410 #[tokio::test]
4411 async fn a_runtime_child_that_dies_on_respawn_is_reported_after_one_retry() {
4412 crate::install_test_crypto_provider();
4413 let (state, _tmp) = capability_test_state();
4414 *state.runtime_bridge.lock().await =
4415 Some(Arc::new(Mutex::new(dead_bridge("http://dead.invalid"))));
4416
4417 let starts = std::sync::atomic::AtomicUsize::new(0);
4418 let err = acquire_live_runtime_bridge_with(&state, || {
4419 starts.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
4420 async { Ok(dead_bridge("http://dead-again.invalid")) }
4421 })
4422 .await
4423 .expect_err("a child that exits at once is not handed out");
4424 assert_eq!(err.code, RUNTIME_UNAVAILABLE_CODE);
4425 assert!(
4426 err.message.contains("exited immediately"),
4427 "{}",
4428 err.message
4429 );
4430 assert_eq!(
4431 starts.load(std::sync::atomic::Ordering::SeqCst),
4432 1,
4433 "respawned exactly once",
4434 );
4435 assert!(
4436 state.runtime_bridge.lock().await.is_none(),
4437 "no dead bridge stays cached",
4438 );
4439 }
4440
4441 #[tokio::test]
4442 async fn config_set_keeps_edits_saved_by_other_processes() {
4443 // Persisting used to write the startup snapshot back over a freshly
4444 // loaded store, erasing edits another process (the TUI, `login`)
4445 // had saved since.
4446 let tmp = tempfile::tempdir().expect("tempdir");
4447 let config_path = tmp.path().join("config.toml");
4448 fs::write(&config_path, "model = \"deepseek-chat\"\n").expect("write config");
4449 let state = build_state(Some(config_path.clone()), None).expect("state");
4450 fs::write(
4451 &config_path,
4452 "model = \"deepseek-chat\"\ntelemetry = true\n",
4453 )
4454 .expect("external edit");
4455
4456 let response = process_app_request(
4457 &state,
4458 AppRequest::ConfigSet {
4459 key: "model".to_string(),
4460 value: "deepseek-reasoner".to_string(),
4461 },
4462 AppTransport::Stdio,
4463 )
4464 .await;
4465 assert!(response.ok, "set should succeed: {response:?}");
4466 let persisted = fs::read_to_string(&config_path).expect("read config");
4467 assert!(persisted.contains("deepseek-reasoner"), "{persisted}");
4468 assert!(
4469 persisted.contains("telemetry = true"),
4470 "the external edit must survive: {persisted}"
4471 );
4472 assert_eq!(
4473 state.config.read().await.telemetry,
4474 Some(true),
4475 "the live config reflects what was saved",
4476 );
4477 }
4478
4479 #[tokio::test(flavor = "current_thread")]
4480 async fn config_store_work_does_not_block_the_runtime() {
4481 let (state, _tmp) = capability_test_state();
4482 let (started_tx, started_rx) = tokio::sync::oneshot::channel();
4483 let (resume_tx, resume_rx) = std::sync::mpsc::channel();
4484 let update_state = state.clone();
4485 let update = tokio::spawn(async move {
4486 update_config_store(&update_state, move |store| {
4487 started_tx
4488 .send(())
4489 .map_err(|_| anyhow!("configuration test caller disappeared"))?;
4490 // Simulate a stalled file/store operation. The only Tokio
4491 // thread must remain free to send the resume signal.
4492 resume_rx
4493 .recv_timeout(Duration::from_secs(3))
4494 .context("configuration store work blocked the runtime")?;
4495 store.config.model = Some("deepseek-reasoner".to_string());
4496 Ok(())
4497 })
4498 .await
4499 });
4500 tokio::time::timeout(Duration::from_secs(5), started_rx)
4501 .await
4502 .expect("the store operation started")
4503 .expect("start signal");
4504 resume_tx
4505 .send(())
4506 .expect("Tokio must run while the store operation is waiting");
4507 update
4508 .await
4509 .expect("configuration task joined")
4510 .expect("configuration update finished");
4511 assert_eq!(
4512 state.config.read().await.model.as_deref(),
4513 Some("deepseek-reasoner")
4514 );
4515 }
4516
4517 #[tokio::test]
4518 async fn config_set_reports_a_failed_load() {
4519 crate::install_test_crypto_provider();
4520 // Persisting failures were logged and swallowed: the reply said ok
4521 // while disk (and so every future turn) kept the old value.
4522 let tmp = tempfile::tempdir().expect("tempdir");
4523 let config_path = tmp.path().join("config.toml");
4524 fs::write(&config_path, "model = \"deepseek-chat\"\n").expect("write config");
4525 let state = build_state(Some(config_path.clone()), None).expect("state");
4526 *state.runtime_bridge.lock().await = Some(sentinel_bridge());
4527 fs::write(&config_path, "model = [unterminated\n").expect("corrupt config");
4528
4529 let response = process_app_request(
4530 &state,
4531 AppRequest::ConfigSet {
4532 key: "model".to_string(),
4533 value: "deepseek-reasoner".to_string(),
4534 },
4535 AppTransport::Stdio,
4536 )
4537 .await;
4538 assert!(!response.ok, "an unsaved change must not report ok");
4539 assert!(response.data["error"].is_string());
4540 assert_eq!(
4541 app_response_status(&response),
4542 StatusCode::INTERNAL_SERVER_ERROR,
4543 "a config file the server cannot read is a server fault",
4544 );
4545 assert_eq!(
4546 state.config.read().await.model.as_deref(),
4547 Some("deepseek-chat"),
4548 "nothing is propagated when the load fails",
4549 );
4550 assert!(
4551 state.runtime_bridge.lock().await.is_some(),
4552 "a failed load must not tear down the bridge",
4553 );
4554 }
4555
4556 #[cfg(unix)]
4557 #[tokio::test]
4558 async fn config_set_reports_a_failed_save() {
4559 use std::os::unix::fs::PermissionsExt as _;
4560 crate::install_test_crypto_provider();
4561 let tmp = tempfile::tempdir().expect("tempdir");
4562 let config_dir = tmp.path().join("cfg");
4563 fs::create_dir(&config_dir).expect("config dir");
4564 let config_path = config_dir.join("config.toml");
4565 fs::write(&config_path, "model = \"deepseek-chat\"\n").expect("write config");
4566 let state = build_state(Some(config_path.clone()), None).expect("state");
4567 *state.runtime_bridge.lock().await = Some(sentinel_bridge());
4568 fs::set_permissions(&config_dir, fs::Permissions::from_mode(0o500)).expect("read-only");
4569 if fs::write(config_dir.join("probe"), "").is_ok() {
4570 // Running as a user that ignores directory permissions (root);
4571 // the save cannot be made to fail this way.
4572 fs::set_permissions(&config_dir, fs::Permissions::from_mode(0o700)).ok();
4573 return;
4574 }
4575
4576 let set = |key: &str, value: &str| AppRequest::ConfigSet {
4577 key: key.to_string(),
4578 value: value.to_string(),
4579 };
4580 let response = process_app_request(
4581 &state,
4582 set("model", "deepseek-reasoner"),
4583 AppTransport::Stdio,
4584 )
4585 .await;
4586 let rejected =
4587 process_app_request(&state, set("telemetry", "not-a-bool"), AppTransport::Stdio).await;
4588 fs::set_permissions(&config_dir, fs::Permissions::from_mode(0o700)).expect("restore");
4589
4590 assert!(!response.ok, "an unsaved change must not report ok");
4591 let error = response.data["error"].as_str().expect("error message");
4592 assert!(error.starts_with("failed to save config"), "{error}");
4593 assert_eq!(
4594 app_response_status(&response),
4595 StatusCode::INTERNAL_SERVER_ERROR,
4596 "a failed save is a server fault, not a bad request",
4597 );
4598 assert!(!rejected.ok);
4599 assert_eq!(
4600 app_response_status(&rejected),
4601 StatusCode::BAD_REQUEST,
4602 "an invalid value is still the caller's mistake",
4603 );
4604 assert_eq!(
4605 fs::read_to_string(&config_path).expect("read config"),
4606 "model = \"deepseek-chat\"\n",
4607 );
4608 assert_eq!(
4609 state.config.read().await.model.as_deref(),
4610 Some("deepseek-chat"),
4611 "nothing is propagated when the save fails",
4612 );
4613 assert!(
4614 state.runtime_bridge.lock().await.is_some(),
4615 "a failed save must not tear down the bridge",
4616 );
4617 }
4618
4619 #[tokio::test]
4620 async fn a_saved_config_change_still_propagates_when_the_caller_goes_away() {
4621 crate::install_test_crypto_provider();
4622 // Once the file was written, a request future dropped before the
4623 // propagation (an HTTP client disconnect) left disk ahead of the live
4624 // runtime and the cached bridge until a reload or restart.
4625 let tmp = tempfile::tempdir().expect("tempdir");
4626 let config_path = tmp.path().join("config.toml");
4627 fs::write(&config_path, "model = \"deepseek-chat\"\n").expect("write config");
4628 let state = build_state(Some(config_path.clone()), None).expect("state");
4629 *state.runtime_bridge.lock().await = Some(sentinel_bridge());
4630
4631 // A concurrent reader holds the runtime, so the set
4632 // saves to disk and then waits to propagate.
4633 let runtime_reader = state.runtime.read().await;
4634 let request = process_app_request(
4635 &state,
4636 AppRequest::ConfigSet {
4637 key: "model".to_string(),
4638 value: "deepseek-reasoner".to_string(),
4639 },
4640 AppTransport::Http,
4641 );
4642 // Long enough for the save to reach disk on a slow runner: 200 ms was
4643 // not, once, on hosted Windows.
4644 assert!(
4645 tokio::time::timeout(Duration::from_secs(2), request)
4646 .await
4647 .is_err(),
4648 "the set waits for the runtime",
4649 );
4650 assert!(
4651 fs::read_to_string(&config_path)
4652 .expect("read config")
4653 .contains("deepseek-reasoner"),
4654 "the change was saved before the caller went away",
4655 );
4656 drop(runtime_reader);
4657
4658 tokio::time::timeout(Duration::from_secs(5), async {
4659 while state.runtime_bridge.lock().await.is_some() {
4660 tokio::time::sleep(Duration::from_millis(10)).await;
4661 }
4662 })
4663 .await
4664 .expect("the saved change still invalidates the bridge");
4665 assert_eq!(
4666 state.runtime.read().await.config.model.as_deref(),
4667 Some("deepseek-reasoner"),
4668 );
4669 assert_eq!(
4670 state.config.read().await.model.as_deref(),
4671 Some("deepseek-reasoner"),
4672 );
4673 }
4674
4675 #[tokio::test]
4676 async fn config_reload_reads_disk_after_the_earlier_queued_write() {
4677 let tmp = tempfile::tempdir().unwrap();
4678 let config_path = tmp.path().join("config.toml");
4679 fs::write(&config_path, "model = \"deepseek-chat\"\n").unwrap();
4680 let state = build_state(Some(config_path.clone()), None).unwrap();
4681 let reader = state.config.read().await;
4682 let set_state = state.clone();
4683 let set = tokio::spawn(async move {
4684 process_app_request(
4685 &set_state,
4686 AppRequest::ConfigSet {
4687 key: "model".into(),
4688 value: "deepseek-reasoner".into(),
4689 },
4690 AppTransport::Http,
4691 )
4692 .await
4693 });
4694 // Tokio's writer-preferring lock rejects new readers only once the
4695 // first writer is queued. Keep the original reader until reload has
4696 // also been polled, deterministically reproducing the stale-load race.
4697 tokio::time::timeout(Duration::from_secs(5), async {
4698 while state.config.try_read().is_ok() {
4699 tokio::task::yield_now().await;
4700 }
4701 })
4702 .await
4703 .unwrap();
4704 let mut reload = Box::pin(process_app_request(
4705 &state,
4706 AppRequest::ConfigReload,
4707 AppTransport::Http,
4708 ));
4709 std::future::poll_fn(|cx| {
4710 assert!(std::future::Future::poll(reload.as_mut(), cx).is_pending());
4711 std::task::Poll::Ready(())
4712 })
4713 .await;
4714 drop(reader);
4715 let (set, reload) =
4716 tokio::time::timeout(Duration::from_secs(5), async { tokio::join!(set, reload) })
4717 .await
4718 .expect("both queued operations finish");
4719 assert!(set.unwrap().ok);
4720 assert!(reload.ok);
4721 assert_eq!(
4722 state.config.read().await.model.as_deref(),
4723 Some("deepseek-reasoner")
4724 );
4725 assert_eq!(
4726 state.runtime.read().await.config.model.as_deref(),
4727 Some("deepseek-reasoner")
4728 );
4729 assert_eq!(
4730 ConfigStore::load(Some(config_path))
4731 .unwrap()
4732 .config
4733 .model
4734 .as_deref(),
4735 Some("deepseek-reasoner")
4736 );
4737 }
4738
4739 #[tokio::test]
4740 async fn config_reload_finishes_propagation_when_the_caller_goes_away() {
4741 crate::install_test_crypto_provider();
4742 let tmp = tempfile::tempdir().unwrap();
4743 let config_path = tmp.path().join("config.toml");
4744 fs::write(&config_path, "model = \"deepseek-chat\"\n").unwrap();
4745 let state = build_state(Some(config_path.clone()), None).unwrap();
4746 *state.runtime_bridge.lock().await = Some(sentinel_bridge());
4747 fs::write(&config_path, "model = \"deepseek-reasoner\"\n").unwrap();
4748 let reader = state.runtime.read().await;
4749 let mut reload = Box::pin(process_app_request(
4750 &state,
4751 AppRequest::ConfigReload,
4752 AppTransport::Http,
4753 ));
4754 std::future::poll_fn(|cx| {
4755 assert!(std::future::Future::poll(reload.as_mut(), cx).is_pending());
4756 std::task::Poll::Ready(())
4757 })
4758 .await;
4759 tokio::time::timeout(Duration::from_secs(5), async {
4760 while state.config.read().await.model.as_deref() != Some("deepseek-reasoner") {
4761 tokio::task::yield_now().await;
4762 }
4763 })
4764 .await
4765 .unwrap();
4766 drop(reload);
4767 drop(reader);
4768 tokio::time::timeout(Duration::from_secs(5), async {
4769 while state.runtime_bridge.lock().await.is_some() {
4770 tokio::task::yield_now().await;
4771 }
4772 })
4773 .await
4774 .expect("reload propagation outlives its caller");
4775 assert_eq!(
4776 state.runtime.read().await.config.model.as_deref(),
4777 Some("deepseek-reasoner")
4778 );
4779 assert_eq!(
4780 fs::read_to_string(config_path).unwrap(),
4781 "model = \"deepseek-reasoner\"\n",
4782 "reload must not rewrite the file"
4783 );
4784 }
4785
4786 #[tokio::test]
4787 async fn config_unset_propagates_to_runtime_config() {
4788 let tmp = tempfile::tempdir().expect("tempdir");
4789 let config_path = tmp.path().join("config.toml");
4790 fs::write(
4791 &config_path,
4792 "api_key = \"sk-deepseek-secret\"\nmodel = \"deepseek-chat\"\n",
4793 )
4794 .expect("write config");
4795 let state = build_state(Some(config_path.clone()), None).expect("state");
4796
4797 // Sanity: runtime starts with the on-disk model.
4798 {
4799 let runtime = state.runtime.read().await;
4800 assert_eq!(runtime.config.model.as_deref(), Some("deepseek-chat"));
4801 }
4802
4803 // Unset the model via the API. This walks a separate code path
4804 // from ConfigSet (unset_value + update_config), so it needs its
4805 // own regression coverage.
4806 let response = process_app_request(
4807 &state,
4808 AppRequest::ConfigUnset {
4809 key: "model".to_string(),
4810 },
4811 AppTransport::Stdio,
4812 )
4813 .await;
4814 assert!(response.ok, "unset should succeed");
4815
4816 // Live runtime sees the cleared model.
4817 {
4818 let runtime = state.runtime.read().await;
4819 assert!(runtime.config.model.is_none());
4820 }
4821 // Shared config lock agrees.
4822 {
4823 let cfg = state.config.read().await;
4824 assert!(cfg.model.is_none());
4825 }
4826 // The on-disk file no longer carries the model value.
4827 let persisted = fs::read_to_string(&config_path).expect("read config");
4828 assert!(!persisted.contains("deepseek-chat"));
4829 }
4830
4831 #[tokio::test]
4832 async fn config_reload_returns_error_when_disk_config_is_invalid() {
4833 let tmp = tempfile::tempdir().expect("tempdir");
4834 let config_path = tmp.path().join("config.toml");
4835 fs::write(
4836 &config_path,
4837 "api_key = \"sk-deepseek-secret\"\nmodel = \"deepseek-chat\"\n",
4838 )
4839 .expect("write config");
4840 let state = build_state(Some(config_path.clone()), None).expect("state");
4841
4842 // Corrupt the on-disk config so ConfigStore::load fails to parse.
4843 fs::write(&config_path, "api_key = \"unterminated\n").expect("corrupt config");
4844
4845 let response =
4846 process_app_request(&state, AppRequest::ConfigReload, AppTransport::Stdio).await;
4847 assert!(!response.ok, "reload of corrupt config must fail");
4848 let err = response.data["error"]
4849 .as_str()
4850 .expect("error message present")
4851 .to_string();
4852 assert!(
4853 err.contains("failed to load config"),
4854 "error should mention load failure, got: {err}"
4855 );
4856
4857 // Live state is untouched: the early-return on load error must
4858 // not have clobbered runtime.config or state.config.
4859 {
4860 let runtime = state.runtime.read().await;
4861 assert_eq!(runtime.config.model.as_deref(), Some("deepseek-chat"));
4862 }
4863 {
4864 let cfg = state.config.read().await;
4865 assert_eq!(cfg.model.as_deref(), Some("deepseek-chat"));
4866 }
4867 }
4868
4869 async fn seed_test_bridge(state: &AppState) -> SharedRuntimeBridge {
4870 let bridge = Arc::new(Mutex::new(RuntimeBridge::from_base_url_for_test(
4871 "http://127.0.0.1:9".to_string(),
4872 )));
4873 *state.runtime_bridge.lock().await = Some(bridge.clone());
4874 bridge
4875 }
4876
4877 #[tokio::test]
4878 async fn config_set_invalidates_cached_stdio_bridge() {
4879 let tmp = tempfile::tempdir().expect("tempdir");
4880 let config_path = tmp.path().join("config.toml");
4881 fs::write(&config_path, "model = \"deepseek-chat\"\n").expect("write config");
4882 let state = build_state(Some(config_path), None).expect("state");
4883 seed_test_bridge(&state).await;
4884
4885 let response = process_app_request(
4886 &state,
4887 AppRequest::ConfigSet {
4888 key: "model".to_string(),
4889 value: "deepseek-reasoner".to_string(),
4890 },
4891 AppTransport::Stdio,
4892 )
4893 .await;
4894 assert!(response.ok, "set should succeed");
4895
4896 // The cached bridge child must be dropped so the next stdio request
4897 // spawns a fresh runtime that reads the persisted config.
4898 assert!(state.runtime_bridge.lock().await.is_none());
4899 }
4900
4901 #[tokio::test]
4902 async fn config_reload_invalidates_cached_stdio_bridge() {
4903 let tmp = tempfile::tempdir().expect("tempdir");
4904 let config_path = tmp.path().join("config.toml");
4905 fs::write(&config_path, "model = \"deepseek-chat\"\n").expect("write config");
4906 let state = build_state(Some(config_path), None).expect("state");
4907 seed_test_bridge(&state).await;
4908
4909 let response =
4910 process_app_request(&state, AppRequest::ConfigReload, AppTransport::Stdio).await;
4911 assert!(response.ok, "reload should succeed");
4912
4913 assert!(state.runtime_bridge.lock().await.is_none());
4914 }
4915
4916 #[tokio::test]
4917 async fn stdio_bridge_invalidation_not_blocked_by_in_flight_turn() {
4918 let (state, _tmp) = capability_test_state();
4919 let bridge = seed_test_bridge(&state).await;
4920
4921 // Simulate a long streaming turn holding the inner bridge lock.
4922 let _in_flight = bridge.lock().await;
4923
4924 // Invalidation only touches the cache slot, so it must complete
4925 // without waiting for the in-flight turn to release the bridge.
4926 tokio::time::timeout(Duration::from_secs(1), invalidate_runtime_bridge(&state))
4927 .await
4928 .expect("invalidation must not wait on bridge traffic");
4929 assert!(state.runtime_bridge.lock().await.is_none());
4930 }
4931
4932 #[tokio::test]
4933 async fn runtime_read_paths_run_concurrently() {
4934 // Tool/status/mcp handlers take read guards; two must coexist so a
4935 // long-running tool call cannot serialize unrelated requests. With
4936 // the old `Mutex<Runtime>` this pattern would deadlock.
4937 let (state, _tmp) = capability_test_state();
4938 let first = state.runtime.read().await;
4939 let second = state.runtime.read().await;
4940 assert!(first.app_status().ok);
4941 assert!(second.app_status().ok);
4942 }
4943
4944 #[tokio::test]
4945 async fn health_probes_advertise_legacy_deepseek_service_name() {
4946 // External probes still key off the DeepSeek-era service name; both
4947 // transports must serve it from the single compat shim.
4948 let (app, _tmp) = app_with_config(None);
4949 let response = app
4950 .oneshot(
4951 Request::builder()
4952 .method(Method::GET)
4953 .uri("/healthz")
4954 .body(Body::empty())
4955 .expect("request"),
4956 )
4957 .await
4958 .expect("response");
4959 let body = response_body_json(response).await;
4960 assert_eq!(body["service"], legacy_deepseek_compat::SERVICE_NAME);
4961 assert_eq!(body["service"], "deepseek-app-server");
4962
4963 let (state, _tmp) = capability_test_state();
4964 let stdio = dispatch_stdio_request(&state, "healthz", json!({}))
4965 .await
4966 .expect("stdio healthz");
4967 assert_eq!(
4968 stdio.result["service"],
4969 legacy_deepseek_compat::SERVICE_NAME
4970 );
4971 }
4972
4973 #[test]
4974 fn non_loopback_bind_without_auth_fails_fast() {
4975 let options = AppServerOptions {
4976 listen: "0.0.0.0:8787".parse().expect("socket addr"),
4977 config_path: None,
4978 auth_token: None,
4979 insecure_no_auth: false,
4980 cors_origins: Vec::new(),
4981 };
4982
4983 let err =
4984 resolve_auth_token(&options).expect_err("non-loopback generated auth should fail");
4985 assert!(err.to_string().contains("without explicit auth token"));
4986 }
4987
4988 #[tokio::test]
4989 async fn stdio_transport_redacts_config_get_secrets() {
4990 let tmp = tempfile::tempdir().expect("tempdir");
4991 let config_path = tmp.path().join("config.toml");
4992 fs::write(&config_path, "").expect("write config");
4993 let state = build_state(Some(config_path), None).expect("state");
4994 {
4995 let mut cfg = state.config.write().await;
4996 cfg.providers.deepseek.api_key = Some("sk-deepseek-secret".to_string());
4997 }
4998
4999 let response = process_app_request(
5000 &state,
5001 AppRequest::ConfigGet {
5002 key: "api_key".to_string(),
5003 },
5004 AppTransport::Stdio,
5005 )
5006 .await;
5007
5008 assert_eq!(response.data["value"], "sk-d***cret");
5009 }
5010
5011 #[tokio::test]
5012 async fn stdio_thread_goal_methods_round_trip_persisted_goal() {
5013 let (state, _tmp, server) = thread_control::compatibility_fixture().await;
5014 let capabilities = dispatch_stdio_request(&state, "thread/capabilities", json!({}))
5015 .await
5016 .unwrap();
5017 assert!(
5018 capabilities.result["methods"]
5019 .as_array()
5020 .unwrap()
5021 .iter()
5022 .any(|method| method == "thread/goal/set")
5023 );
5024 let started = dispatch_stdio_request(
5025 &state,
5026 "thread/start",
5027 json!({"operation_key":"goal-start"}),
5028 )
5029 .await
5030 .unwrap();
5031 let id = started.result["thread_id"].as_str().unwrap();
5032 let set = dispatch_stdio_request(
5033 &state,
5034 "thread/goal/set",
5035 json!({"thread_id":id,"objective":"Release 0.10.1","token_budget":59000}),
5036 )
5037 .await
5038 .unwrap();
5039 assert_eq!(set.result["goal"]["objective"], "Release 0.10.1");
5040 assert_eq!(set.result["goal"]["status"], "active");
5041 let got = dispatch_stdio_request(&state, "thread/goal/get", json!({"thread_id":id}))
5042 .await
5043 .unwrap();
5044 assert_eq!(got.result["goal"]["token_budget"], 59000);
5045 let cleared = dispatch_stdio_request(&state, "thread/goal/clear", json!({"thread_id":id}))
5046 .await
5047 .unwrap();
5048 assert_eq!(cleared.result["status"], "cleared");
5049 assert_eq!(cleared.result["data"]["cleared"], true);
5050 server.abort();
5051 }
5052
5053 #[tokio::test]
5054 async fn stdio_resume_of_missing_thread_fails_without_clobbering_the_hint() {
5055 let (state, tmp, server) = thread_control::compatibility_fixture().await;
5056 let workspace = tmp.path().join("ws");
5057 state.stdio_thread_hints.lock().await.insert(
5058 "ghost-thread".into(),
5059 RuntimeThreadHint {
5060 model: Some("fixture-model".into()),
5061 workspace: Some(workspace.clone()),
5062 },
5063 );
5064 for method in ["thread/resume", "thread/fork"] {
5065 let error = dispatch_stdio_request(
5066 &state,
5067 method,
5068 json!({"thread_id":"ghost-thread","operation_key":format!("{method}-missing")}),
5069 )
5070 .await
5071 .unwrap_err();
5072 assert_eq!(error.code, THREAD_NOT_FOUND_CODE);
5073 assert!(error.message.contains("ghost-thread"));
5074 }
5075 let hints = state.stdio_thread_hints.lock().await;
5076 assert_eq!(
5077 hints["ghost-thread"].model.as_deref(),
5078 Some("fixture-model")
5079 );
5080 assert_eq!(hints["ghost-thread"].workspace.as_ref(), Some(&workspace));
5081 server.abort();
5082 }
5083
5084 #[tokio::test]
5085 async fn stdio_archive_of_missing_thread_fails_instead_of_reporting_success() {
5086 let (state, _tmp, server) = thread_control::compatibility_fixture().await;
5087 for method in ["thread/archive", "thread/unarchive"] {
5088 let error = dispatch_stdio_request(&state, method, json!({"thread_id":"ghost-thread"}))
5089 .await
5090 .unwrap_err();
5091 assert_eq!(error.code, THREAD_NOT_FOUND_CODE);
5092 assert!(error.message.contains("ghost-thread"));
5093 }
5094 server.abort();
5095 }
5096
5097 fn sse_frame(event: &str, payload: Value) -> String {
5098 format!("event: {event}\ndata: {payload}\n\n")
5099 }
5100 /// A runtime whose turn never ends on its own — only an interrupt stops
5101 /// it. That is the shape of the runaway turn this protects against.
5102 async fn spawn_uninterruptible_until_asked_runtime(
5103 owner_router: Router,
5104 ) -> (
5105 String,
5106 Arc<tokio::sync::Notify>,
5107 tokio::task::JoinHandle<()>,
5108 ) {
5109 use axum::body::Body;
5110 use axum::extract::Path as AxumPath;
5111
5112 let interrupted = Arc::new(tokio::sync::Notify::new());
5113
5114 async fn create_turn(AxumPath(_thread_id): AxumPath<String>) -> Json<Value> {
5115 Json(json!({ "turn": { "id": "turn_runaway" } }))
5116 }
5117 async fn create_thread() -> Json<Value> {
5118 Json(json!({ "id": "thr_runaway" }))
5119 }
5120 async fn interrupt(
5121 State(notify): State<Arc<tokio::sync::Notify>>,
5122 AxumPath((_thread_id, _turn_id)): AxumPath<(String, String)>,
5123 ) -> Json<Value> {
5124 notify.notify_waiters();
5125 Json(json!({ "ok": true }))
5126 }
5127 async fn thread_events(
5128 State(notify): State<Arc<tokio::sync::Notify>>,
5129 AxumPath(_thread_id): AxumPath<String>,
5130 ) -> ([(header::HeaderName, &'static str); 1], Body) {
5131 // Hold the event response open until something interrupts the
5132 // turn. Nothing else can end it, which is the point.
5133 notify.notified().await;
5134 let body = [
5135 sse_frame(
5136 "item.delta",
5137 json!({
5138 "seq": 1,
5139 "turn_id": "turn_runaway",
5140 "payload": { "kind": "agent_message", "delta": "thinking" }
5141 }),
5142 ),
5143 sse_frame(
5144 "turn.completed",
5145 json!({
5146 "seq": 2,
5147 "turn_id": "turn_runaway",
5148 "payload": { "turn": { "status": "interrupted" } }
5149 }),
5150 ),
5151 ]
5152 .concat();
5153 (
5154 [(header::CONTENT_TYPE, "text/event-stream")],
5155 Body::from(body),
5156 )
5157 }
5158
5159 let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
5160 .await
5161 .expect("bind test listener");
5162 let addr = listener.local_addr().expect("listener addr");
5163 let app = Router::new()
5164 .route("/v1/threads", post(create_thread))
5165 .route("/v1/threads/{thread_id}/turns", post(create_turn))
5166 .route(
5167 "/v1/threads/{thread_id}/turns/{turn_id}/interrupt",
5168 post(interrupt),
5169 )
5170 .route("/v1/threads/{thread_id}/events", get(thread_events))
5171 .with_state(interrupted.clone())
5172 .fallback_service(owner_router);
5173 let server = tokio::spawn(async move {
5174 axum::serve(listener, app)
5175 .await
5176 .expect("serve test runtime");
5177 });
5178 (format!("http://{addr}"), interrupted, server)
5179 }
5180
5181 #[tokio::test]
5182 async fn interrupt_stops_a_turn_that_would_otherwise_stream_forever() {
5183 let (state, _tmp, owner_router) = thread_control::compatibility_router();
5184 let (base_url, _notify, server) =
5185 spawn_uninterruptible_until_asked_runtime(owner_router).await;
5186 seed_client_thread(&state, "thr_a").await;
5187 *state.runtime_bridge.lock().await = Some(Arc::new(Mutex::new(
5188 RuntimeBridge::from_base_url_for_test(base_url),
5189 )));
5190
5191 let (client, server_side) = tokio::io::duplex(16 * 1024);
5192 let (client_reader, mut client_writer) = tokio::io::split(client);
5193
5194 let loop_state = state.clone();
5195 let loop_handle = tokio::spawn(async move {
5196 let (rx, tx) = tokio::io::split(server_side);
5197 run_stdio_loop(
5198 &loop_state,
5199 BoundedLines::new(BufReader::new(rx)),
5200 tx,
5201 StdioLoopPolicy::process_stdio(),
5202 None::<()>,
5203 )
5204 .await
5205 });
5206
5207 // Start the runaway turn.
5208 client_writer
5209 .write_all(
5210 b"{\"jsonrpc\":\"2.0\",\"id\":1,\"method\":\"thread/message\",\
5211 \"params\":{\"thread_id\":\"thr_a\",\"input\":\"go\"}}\n",
5212 )
5213 .await
5214 .expect("send thread/message");
5215
5216 // Wait until the turn is genuinely in flight before cancelling, so the
5217 // test exercises mid-stream cancellation rather than a race.
5218 tokio::time::timeout(Duration::from_secs(10), async {
5219 loop {
5220 if state.in_flight_turns.lock().await.contains_key("thr_a") {
5221 return;
5222 }
5223 tokio::time::sleep(Duration::from_millis(10)).await;
5224 }
5225 })
5226 .await
5227 .expect("turn should register itself as in flight");
5228
5229 // The read loop must accept this while the turn holds the bridge.
5230 client_writer
5231 .write_all(
5232 b"{\"jsonrpc\":\"2.0\",\"id\":2,\"method\":\"thread/interrupt\",\
5233 \"params\":{\"thread_id\":\"thr_a\"}}\n",
5234 )
5235 .await
5236 .expect("send thread/interrupt");
5237 client_writer
5238 .write_all(b"{\"jsonrpc\":\"2.0\",\"id\":3,\"method\":\"shutdown\"}\n")
5239 .await
5240 .expect("send shutdown");
5241
5242 let finished = tokio::time::timeout(Duration::from_secs(20), loop_handle)
5243 .await
5244 .expect("the loop must exit rather than hang on the runaway turn");
5245 finished.expect("join loop").expect("loop result");
5246
5247 let mut output = String::new();
5248 let mut lines = BufReader::new(client_reader);
5249 lines
5250 .read_to_string(&mut output)
5251 .await
5252 .expect("read stdio output");
5253
5254 let responses: Vec<Value> = output
5255 .lines()
5256 .filter_map(|line| serde_json::from_str::<Value>(line).ok())
5257 .collect();
5258 let by_id = |id: u64| {
5259 responses
5260 .iter()
5261 .find(|value| value["id"] == json!(id))
5262 .unwrap_or_else(|| panic!("no response for id {id} in {output}"))
5263 .clone()
5264 };
5265
5266 // The turn ended as interrupted rather than running to completion.
5267 assert!(
5268 by_id(1)["error"].is_object(),
5269 "the interrupted turn should report an error, got: {}",
5270 by_id(1)
5271 );
5272 assert_eq!(by_id(2)["result"]["interrupted"], json!(true));
5273 assert_eq!(by_id(3)["result"]["status"], json!("stopped"));
5274
5275 server.abort();
5276 let _ = server.await;
5277 }
5278
5279 #[tokio::test]
5280 async fn interrupting_an_idle_thread_is_not_an_error() {
5281 let (state, _tmp) = capability_test_state();
5282 let response = dispatch_stdio_request(
5283 &state,
5284 "thread/interrupt",
5285 json!({ "thread_id": "thr_nothing_running" }),
5286 )
5287 .await
5288 .expect("interrupt dispatch");
5289 assert_eq!(response.result["interrupted"], json!(false));
5290 }
5291
5292 #[tokio::test]
5293 async fn output_cap_bridge_checks_support_before_creation_and_forwards_each_surface() {
5294 use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
5295 #[derive(Clone)]
5296 struct Fixture {
5297 supported: Arc<AtomicBool>,
5298 created: Arc<AtomicUsize>,
5299 requests: Arc<Mutex<Vec<Value>>>,
5300 }
5301 async fn info(State(f): State<Fixture>, headers: axum::http::HeaderMap) -> Json<Value> {
5302 assert_eq!(
5303 headers.get(header::AUTHORIZATION).unwrap(),
5304 "Bearer fixture-output-cap"
5305 );
5306 Json(
5307 json!({"capabilities":{"turn_output_token_limit":f.supported.load(Ordering::SeqCst)}}),
5308 )
5309 }
5310 async fn providers() -> Json<Value> {
5311 Json(
5312 json!({"current":"custom","providers":[{"id":"custom","default_model":"fixture-model"}]}),
5313 )
5314 }
5315 async fn models() -> Json<Value> {
5316 Json(
5317 json!({"models":[{"id":"fixture-model","output_token_limit":"supported"},{"id":"uncapped-transport","output_token_limit":"unsupported"}]}),
5318 )
5319 }
5320 async fn create_thread(State(f): State<Fixture>) -> Json<Value> {
5321 let n = f.created.fetch_add(1, Ordering::SeqCst);
5322 Json(json!({"id":format!("thr_cap_{n}")}))
5323 }
5324 async fn create_turn(State(f): State<Fixture>, Json(body): Json<Value>) -> Json<Value> {
5325 f.requests.lock().await.push(body);
5326 Json(json!({"turn":{"id":"turn_cap"}}))
5327 }
5328 async fn events(State(f): State<Fixture>) -> impl IntoResponse {
5329 let seq = f.requests.lock().await.len();
5330 (
5331 [(header::CONTENT_TYPE, "text/event-stream")],
5332 sse_frame(
5333 "turn.completed",
5334 json!({
5335 "seq":seq,"turn_id":"turn_cap","payload":{"turn":{"status":"completed"}}
5336 }),
5337 ),
5338 )
5339 }
5340 let fixture = Fixture {
5341 supported: Arc::new(AtomicBool::new(false)),
5342 created: Arc::new(AtomicUsize::new(0)),
5343 requests: Arc::new(Mutex::new(Vec::new())),
5344 };
5345 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
5346 let addr = listener.local_addr().unwrap();
5347 let (state, _tmp, owner_router) = thread_control::compatibility_router();
5348 let router = Router::new()
5349 .route("/v1/runtime/info", get(info))
5350 .route("/v1/providers", get(providers))
5351 .route("/v1/providers/custom/models", get(models))
5352 .route("/v1/threads", post(create_thread))
5353 .route("/v1/threads/{id}/turns", post(create_turn))
5354 .route("/v1/threads/{id}/events", get(events))
5355 .with_state(fixture.clone())
5356 .fallback_service(owner_router);
5357 let server = tokio::spawn(async move { axum::serve(listener, router).await.unwrap() });
5358 let mut bridge = RuntimeBridge::from_base_url_for_test(format!("http://{addr}"));
5359 bridge.auth_token = Some("fixture-output-cap".into());
5360 *state.runtime_bridge.lock().await = Some(Arc::new(Mutex::new(bridge)));
5361 let result = dispatch_stdio_request(
5362 &state,
5363 "prompt/run",
5364 json!({"prompt":"review","maxOutputTokens":1500}),
5365 )
5366 .await;
5367 assert!(result.unwrap_err().message.contains("does not support"));
5368 assert_eq!(fixture.created.load(Ordering::SeqCst), 0);
5369 assert!(fixture.requests.lock().await.is_empty());
5370 fixture.supported.store(true, Ordering::SeqCst);
5371 for model in ["auto", "uncapped-transport", "unknown-model"] {
5372 assert!(
5373 dispatch_stdio_request(
5374 &state,
5375 "prompt/run",
5376 json!({"prompt":"review","model":model,"maxOutputTokens":1500})
5377 )
5378 .await
5379 .is_err()
5380 );
5381 }
5382 assert_eq!(fixture.created.load(Ordering::SeqCst), 0);
5383 for thread_id in ["stdio-cap", "request-cap", "http-cap"] {
5384 seed_client_thread(&state, thread_id).await;
5385 }
5386 for (method, params) in [
5387 (
5388 "prompt/run",
5389 json!({"prompt":"review","maxOutputTokens":1500}),
5390 ),
5391 (
5392 "thread/message",
5393 json!({"thread_id":"stdio-cap","input":"review","maxOutputTokens":1500}),
5394 ),
5395 (
5396 "thread/request",
5397 json!({"kind":"message","thread_id":"request-cap","input":"review","maxOutputTokens":1500}),
5398 ),
5399 ] {
5400 dispatch_stdio_request(&state, method, params)
5401 .await
5402 .expect("existing app-server caller forwards allowance");
5403 }
5404 run_http_thread_message(
5405 &state,
5406 "http-cap".into(),
5407 "review".into(),
5408 Vec::new(),
5409 std::num::NonZeroU32::new(1500),
5410 )
5411 .await
5412 .unwrap();
5413 let requests = fixture.requests.lock().await.clone();
5414 assert_eq!(requests.len(), 4);
5415 assert!(
5416 requests
5417 .iter()
5418 .all(|request| request["maxOutputTokens"] == 1500)
5419 );
5420 let count = fixture.created.load(Ordering::SeqCst);
5421 for invalid in [
5422 json!(0),
5423 json!(-1),
5424 json!(1.5),
5425 json!("1500"),
5426 json!(4_294_967_296u64),
5427 ] {
5428 assert!(
5429 dispatch_stdio_request(
5430 &state,
5431 "prompt/run",
5432 json!({"prompt":"review","maxOutputTokens":invalid})
5433 )
5434 .await
5435 .is_err()
5436 );
5437 }
5438 assert_eq!(fixture.created.load(Ordering::SeqCst), count);
5439 assert_eq!(fixture.requests.lock().await.len(), 4);
5440 server.abort();
5441 let _ = server.await;
5442 }
5443
5444 #[tokio::test]
5445 async fn stdio_runtime_bridge_streams_response_delta_events() {
5446 async fn create_turn(AxumPath(thread_id): AxumPath<String>) -> Json<Value> {
5447 Json(json!({
5448 "thread": { "id": thread_id },
5449 "turn": { "id": "turn_test" },
5450 }))
5451 }
5452
5453 async fn thread_events(
5454 AxumPath(thread_id): AxumPath<String>,
5455 Query(query): Query<HashMap<String, String>>,
5456 ) -> ([(header::HeaderName, &'static str); 1], String) {
5457 assert_eq!(thread_id, "thr_test");
5458 assert_eq!(query.get("since_seq").map(String::as_str), Some("0"));
5459
5460 let body = [
5461 sse_frame(
5462 "item.delta",
5463 json!({
5464 "seq": 1,
5465 "turn_id": "turn_test",
5466 "payload": {
5467 "kind": "agent_message",
5468 "delta": "hello"
5469 }
5470 }),
5471 ),
5472 sse_frame(
5473 "turn.completed",
5474 json!({
5475 "seq": 2,
5476 "turn_id": "turn_test",
5477 "payload": {
5478 "turn": {
5479 "status": "completed"
5480 }
5481 }
5482 }),
5483 ),
5484 ]
5485 .concat();
5486
5487 ([(header::CONTENT_TYPE, "text/event-stream")], body)
5488 }
5489
5490 let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
5491 .await
5492 .expect("bind test listener");
5493 let addr = listener.local_addr().expect("listener addr");
5494 let app = Router::new()
5495 .route("/v1/threads/{thread_id}/turns", post(create_turn))
5496 .route("/v1/threads/{thread_id}/events", get(thread_events));
5497
5498 let server = tokio::spawn(async move {
5499 axum::serve(listener, app)
5500 .await
5501 .expect("serve test runtime");
5502 });
5503
5504 let mut bridge = RuntimeBridge::from_base_url_for_test(format!("http://{addr}"));
5505 let (mut reader, mut writer) = tokio::io::duplex(4096);
5506
5507 let result = bridge
5508 .message_thread(
5509 "thr_test",
5510 RuntimeTurnInput {
5511 input: "hello",
5512 images: &[],
5513 max_output_tokens: None,
5514 expected_workspace: None,
5515 },
5516 &mut writer,
5517 None,
5518 None,
5519 )
5520 .await
5521 .expect("message_thread should succeed");
5522 drop(writer);
5523
5524 let mut stdout = Vec::new();
5525 reader
5526 .read_to_end(&mut stdout)
5527 .await
5528 .expect("read stdio output");
5529 server.abort();
5530 let _ = server.await;
5531
5532 let lines: Vec<Value> = String::from_utf8(stdout)
5533 .expect("utf8 output")
5534 .lines()
5535 .map(|line| serde_json::from_str(line).expect("json line"))
5536 .collect();
5537
5538 assert_eq!(
5539 result.get("status").and_then(Value::as_str),
5540 Some("accepted")
5541 );
5542 assert_eq!(
5543 result.pointer("/data/turn_id").and_then(Value::as_str),
5544 Some("turn_test")
5545 );
5546 assert_eq!(bridge.last_seq_by_thread.get("thr_test"), Some(&2));
5547
5548 let event_types: Vec<&str> = lines
5549 .iter()
5550 .map(|line| {
5551 line.get("type")
5552 .and_then(Value::as_str)
5553 .expect("event type")
5554 })
5555 .collect();
5556 assert_eq!(
5557 event_types,
5558 vec!["response_start", "response_delta", "response_end"]
5559 );
5560 assert_eq!(lines[1]["delta"], "hello");
5561 }
5562
5563 /// Audit R03-05: a stream that breaks before `turn.completed` must not
5564 /// report `response_end` ahead of the error, and must not leave the
5565 /// runtime turn running with nothing able to interrupt it.
5566 #[tokio::test]
5567 async fn stdio_runtime_bridge_interrupts_a_turn_whose_stream_breaks() {
5568 static INTERRUPTED: std::sync::atomic::AtomicBool =
5569 std::sync::atomic::AtomicBool::new(false);
5570
5571 async fn create_turn(AxumPath(thread_id): AxumPath<String>) -> Json<Value> {
5572 Json(json!({
5573 "thread": { "id": thread_id },
5574 "turn": { "id": "turn_broken" },
5575 }))
5576 }
5577
5578 async fn thread_events() -> ([(header::HeaderName, &'static str); 1], String) {
5579 // One delta, then the stream ends without `turn.completed`.
5580 let body = sse_frame(
5581 "item.delta",
5582 json!({
5583 "seq": 1,
5584 "turn_id": "turn_broken",
5585 "payload": { "kind": "agent_message", "delta": "partial" }
5586 }),
5587 );
5588 ([(header::CONTENT_TYPE, "text/event-stream")], body)
5589 }
5590
5591 async fn interrupt(
5592 AxumPath((thread_id, turn_id)): AxumPath<(String, String)>,
5593 ) -> Json<Value> {
5594 assert_eq!(thread_id, "thr_broken");
5595 assert_eq!(turn_id, "turn_broken");
5596 INTERRUPTED.store(true, std::sync::atomic::Ordering::SeqCst);
5597 Json(json!({ "interrupted": true }))
5598 }
5599
5600 let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
5601 .await
5602 .expect("bind test listener");
5603 let addr = listener.local_addr().expect("listener addr");
5604 let app = Router::new()
5605 .route("/v1/threads/{thread_id}/turns", post(create_turn))
5606 .route("/v1/threads/{thread_id}/events", get(thread_events))
5607 .route(
5608 "/v1/threads/{thread_id}/turns/{turn_id}/interrupt",
5609 post(interrupt),
5610 );
5611 let server = tokio::spawn(async move {
5612 axum::serve(listener, app)
5613 .await
5614 .expect("serve test runtime");
5615 });
5616
5617 let mut bridge = RuntimeBridge::from_base_url_for_test(format!("http://{addr}"));
5618 let (mut reader, mut writer) = tokio::io::duplex(4096);
5619 let result = bridge
5620 .message_thread(
5621 "thr_broken",
5622 RuntimeTurnInput {
5623 input: "hello",
5624 images: &[],
5625 max_output_tokens: None,
5626 expected_workspace: None,
5627 },
5628 &mut writer,
5629 None,
5630 None,
5631 )
5632 .await;
5633 drop(writer);
5634 let mut stdout = Vec::new();
5635 reader
5636 .read_to_end(&mut stdout)
5637 .await
5638 .expect("read stdio output");
5639 server.abort();
5640 let _ = server.await;
5641
5642 assert!(result.is_err(), "a broken stream is an error");
5643 let event_types: Vec<String> = String::from_utf8(stdout)
5644 .expect("utf8 output")
5645 .lines()
5646 .map(|line| {
5647 serde_json::from_str::<Value>(line).expect("json line")["type"]
5648 .as_str()
5649 .expect("event type")
5650 .to_string()
5651 })
5652 .collect();
5653 assert_eq!(event_types, vec!["response_start", "response_delta"]);
5654 assert!(
5655 INTERRUPTED.load(std::sync::atomic::Ordering::SeqCst),
5656 "the orphaned runtime turn was not interrupted"
5657 );
5658 }
5659
5660 #[tokio::test]
5661 async fn stdio_runtime_bridge_applies_thread_start_hints() {
5662 async fn create_thread(Json(body): Json<Value>) -> Json<Value> {
5663 assert_eq!(body["model"], "deepseek-v4");
5664 assert_eq!(body["workspace"], "/tmp/codewhale-stdio");
5665 Json(json!({
5666 "id": "thr_runtime",
5667 "model": body["model"].clone(),
5668 "workspace": body["workspace"].clone(),
5669 }))
5670 }
5671
5672 let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
5673 .await
5674 .expect("bind test listener");
5675 let addr = listener.local_addr().expect("listener addr");
5676 let app = Router::new().route("/v1/threads", post(create_thread));
5677
5678 let server = tokio::spawn(async move {
5679 axum::serve(listener, app)
5680 .await
5681 .expect("serve test runtime");
5682 });
5683
5684 let mut bridge = RuntimeBridge::from_base_url_for_test(format!("http://{addr}"));
5685 let mut thread_map = HashMap::new();
5686 let runtime_id = bridge
5687 .ensure_runtime_thread(
5688 &mut thread_map,
5689 "legacy_thread",
5690 Some(RuntimeThreadHint {
5691 model: Some("deepseek-v4".to_string()),
5692 workspace: Some(PathBuf::from("/tmp/codewhale-stdio")),
5693 }),
5694 )
5695 .await
5696 .expect("runtime thread");
5697 server.abort();
5698 let _ = server.await;
5699
5700 assert_eq!(runtime_id, "thr_runtime");
5701 assert_eq!(
5702 thread_map.get("legacy_thread").map(String::as_str),
5703 Some("thr_runtime")
5704 );
5705 }
5706
5707 // ── prompt routing runs a real turn ────────────────────────────────
5708 //
5709 // `/prompt`, `prompt/request` and `prompt/run` used to return HTTP 200
5710 // with a stringified echo of the caller's own routing metadata, having
5711 // called no model at all. These stand up the in-crate stub runtime and
5712 // assert the response is what the model streamed — not an echo — and
5713 // that an unreachable runtime is an explicit typed failure.
5714
5715 /// Prompts the stub runtime was actually asked to run.
5716 type StubPrompts = Arc<Mutex<Vec<String>>>;
5717
5718 /// A minimal but honest runtime: it creates threads, starts turns, and
5719 /// streams `agent_message` deltas followed by `turn.completed`.
5720 async fn spawn_stub_runtime(
5721 owner_router: Option<Router>,
5722 ) -> (String, StubPrompts, tokio::task::JoinHandle<()>) {
5723 async fn create_thread(Json(body): Json<Value>) -> Json<Value> {
5724 Json(json!({
5725 "id": "thr_stub",
5726 "model": body["model"].as_str().unwrap_or("stub-model-v1"),
5727 }))
5728 }
5729
5730 async fn create_turn(
5731 State(prompts): State<StubPrompts>,
5732 AxumPath(thread_id): AxumPath<String>,
5733 Json(body): Json<Value>,
5734 ) -> Json<Value> {
5735 prompts
5736 .lock()
5737 .await
5738 .push(body["prompt"].as_str().unwrap_or_default().to_string());
5739 Json(json!({
5740 "thread": { "id": thread_id, "model": "stub-model-v1" },
5741 "turn": { "id": "turn_stub" },
5742 }))
5743 }
5744
5745 async fn thread_events(
5746 AxumPath(_thread_id): AxumPath<String>,
5747 ) -> ([(header::HeaderName, &'static str); 1], String) {
5748 let body = [
5749 sse_frame(
5750 "item.delta",
5751 json!({
5752 "seq": 1,
5753 "turn_id": "turn_stub",
5754 "payload": { "kind": "agent_message", "delta": "the answer" }
5755 }),
5756 ),
5757 sse_frame(
5758 "item.delta",
5759 json!({
5760 "seq": 2,
5761 "turn_id": "turn_stub",
5762 "payload": { "kind": "agent_message", "delta": " is 4" }
5763 }),
5764 ),
5765 sse_frame(
5766 "turn.completed",
5767 json!({
5768 "seq": 3,
5769 "turn_id": "turn_stub",
5770 "payload": { "turn": { "status": "completed" } }
5771 }),
5772 ),
5773 ]
5774 .concat();
5775 ([(header::CONTENT_TYPE, "text/event-stream")], body)
5776 }
5777
5778 let prompts: StubPrompts = Arc::new(Mutex::new(Vec::new()));
5779 let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
5780 .await
5781 .expect("bind stub runtime");
5782 let addr = listener.local_addr().expect("listener addr");
5783 let app = Router::new()
5784 .route("/v1/threads", post(create_thread))
5785 .route("/v1/threads/{thread_id}/turns", post(create_turn))
5786 .route("/v1/threads/{thread_id}/events", get(thread_events))
5787 .with_state(prompts.clone());
5788 let app = match owner_router {
5789 Some(owner) => app.fallback_service(owner),
5790 None => app,
5791 };
5792 let server = tokio::spawn(async move {
5793 let _ = axum::serve(listener, app).await;
5794 });
5795 (format!("http://{addr}"), prompts, server)
5796 }
5797
5798 async fn seed_bridge_at(state: &AppState, base_url: String) -> SharedRuntimeBridge {
5799 let bridge = Arc::new(Mutex::new(RuntimeBridge::from_base_url_for_test(base_url)));
5800 *state.runtime_bridge.lock().await = Some(bridge.clone());
5801 bridge
5802 }
5803
5804 #[tokio::test]
5805 async fn prompt_request_executes_a_genuine_model_turn() {
5806 let (state, _tmp) = capability_test_state();
5807 let (base_url, prompts, server) = spawn_stub_runtime(None).await;
5808 seed_bridge_at(&state, base_url).await;
5809
5810 let (mut reader, mut writer) = tokio::io::duplex(4096);
5811 let dispatched = dispatch_stdio_request_with_writer(
5812 &state,
5813 &mut writer,
5814 "prompt/request",
5815 json!({ "prompt": "what is 2+2" }),
5816 AppTransport::Stdio,
5817 )
5818 .await
5819 .expect("prompt/request dispatch");
5820 drop(writer);
5821
5822 let response: PromptResponse =
5823 serde_json::from_value(dispatched.result).expect("prompt response");
5824
5825 // The model's words, not a restatement of the request.
5826 assert_eq!(response.output, "the answer is 4");
5827 assert!(
5828 !response.output.contains("what is 2+2"),
5829 "prompt echo leaked into the output: {}",
5830 response.output
5831 );
5832 assert_eq!(response.model, "stub-model-v1");
5833 assert_eq!(
5834 prompts.lock().await.as_slice(),
5835 ["what is 2+2".to_string()],
5836 "the prompt must reach the runtime's turn endpoint"
5837 );
5838
5839 // Real streaming frames, not three canned ones.
5840 let deltas: Vec<String> = response
5841 .events
5842 .iter()
5843 .filter_map(|event| match event {
5844 EventFrame::ResponseDelta { delta, .. } => Some(delta.clone()),
5845 _ => None,
5846 })
5847 .collect();
5848 assert_eq!(deltas, vec!["the answer".to_string(), " is 4".to_string()]);
5849 assert!(matches!(
5850 response.events.first(),
5851 Some(EventFrame::ResponseStart { .. })
5852 ));
5853 assert!(matches!(
5854 response.events.last(),
5855 Some(EventFrame::ResponseEnd { .. })
5856 ));
5857
5858 // The stdio transport sees the same turn stream `thread/message` emits.
5859 let mut stdout = Vec::new();
5860 reader.read_to_end(&mut stdout).await.expect("read stdout");
5861 let stdout = String::from_utf8(stdout).expect("utf8 stdout");
5862 assert!(
5863 stdout.contains("\"type\":\"response_delta\"") && stdout.contains("the answer"),
5864 "stdio prompt turn must stream its deltas, got: {stdout}"
5865 );
5866
5867 // A prompt without a thread_id must not leave a mapping behind.
5868 assert!(
5869 state.runtime_thread_map.lock().await.is_empty(),
5870 "one-shot prompt threads must not accumulate in the map"
5871 );
5872
5873 server.abort();
5874 let _ = server.await;
5875 }
5876
5877 #[tokio::test]
5878 async fn prompt_without_a_reachable_runtime_fails_explicitly() {
5879 let (state, _tmp) = capability_test_state();
5880 // Port 9 (discard) refuses immediately: no runtime is listening.
5881 seed_bridge_at(&state, "http://127.0.0.1:9".to_string()).await;
5882
5883 let err = dispatch_stdio_request(&state, "prompt/run", json!({ "prompt": "hello" }))
5884 .await
5885 .expect_err("a prompt with no reachable runtime must fail, not echo");
5886 assert_eq!(err.code, RUNTIME_UNAVAILABLE_CODE);
5887
5888 let (status, Json(body)) = http_error_from_jsonrpc(err);
5889 assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE);
5890 assert_eq!(body["error"]["code"], "runtime_unavailable");
5891 assert!(
5892 body.get("output").is_none(),
5893 "a failure must not be shaped like a PromptResponse: {body}"
5894 );
5895 }
5896
5897 #[tokio::test]
5898 async fn empty_prompt_is_rejected_before_any_runtime_work() {
5899 let (state, _tmp) = capability_test_state();
5900 let err = dispatch_stdio_request(&state, "prompt/request", json!({ "prompt": " " }))
5901 .await
5902 .expect_err("an empty prompt must be rejected");
5903 assert_eq!(err.code, -32602);
5904 assert!(
5905 state.runtime_bridge.lock().await.is_none(),
5906 "a rejected prompt must not start a runtime"
5907 );
5908 }
5909
5910 #[tokio::test]
5911 async fn http_thread_message_runs_the_turn_instead_of_queueing_it() {
5912 let (state, _tmp, owner_router) = thread_control::compatibility_router();
5913 seed_client_thread(&state, "thr_http").await;
5914 let (base_url, prompts, server) = spawn_stub_runtime(Some(owner_router)).await;
5915 seed_bridge_at(&state, base_url).await;
5916
5917 let response = run_http_thread_message(
5918 &state,
5919 "thr_http".to_string(),
5920 "go".to_string(),
5921 Vec::new(),
5922 None,
5923 )
5924 .await
5925 .expect("http thread message");
5926
5927 assert_eq!(response.status, "completed");
5928 assert_eq!(response.thread_id, "thr_http");
5929 assert_eq!(response.data["turn_id"], "turn_stub");
5930 assert_eq!(prompts.lock().await.as_slice(), ["go".to_string()]);
5931 assert!(
5932 response
5933 .events
5934 .iter()
5935 .any(|event| matches!(event, EventFrame::ResponseDelta { .. })),
5936 "a completed turn must carry the deltas it streamed"
5937 );
5938
5939 server.abort();
5940 let _ = server.await;
5941 }
5942
5943 #[tokio::test]
5944 async fn http_thread_message_without_a_runtime_is_a_typed_error() {
5945 let (state, _tmp, _owner_router) = thread_control::compatibility_router();
5946 seed_client_thread(&state, "thr_http").await;
5947 seed_bridge_at(&state, "http://127.0.0.1:9".to_string()).await;
5948
5949 let err = run_http_thread_message(
5950 &state,
5951 "thr_http".to_string(),
5952 "go".to_string(),
5953 Vec::new(),
5954 None,
5955 )
5956 .await
5957 .expect_err("no runtime means no turn");
5958 assert_eq!(err.code, RUNTIME_UNAVAILABLE_CODE);
5959 }
5960
5961 #[tokio::test]
5962 async fn submit_user_input_refuses_instead_of_claiming_resolution() {
5963 let (state, _tmp) = capability_test_state();
5964 let response = process_app_request(
5965 &state,
5966 AppRequest::SubmitUserInput {
5967 request_id: "user-input-1".to_string(),
5968 answers: Vec::new(),
5969 },
5970 AppTransport::Stdio,
5971 )
5972 .await;
5973
5974 assert!(!response.ok, "this transport cannot deliver the answer");
5975 assert_eq!(response.data["error"], "user_input_reply_unsupported");
5976 assert!(
5977 response.data.get("resolved").is_none(),
5978 "nothing was resolved: {}",
5979 response.data
5980 );
5981 assert!(
5982 response.data["message"]
5983 .as_str()
5984 .expect("message")
5985 .contains("/v1/user-input/"),
5986 "the refusal must name the transport that can accept the answer"
5987 );
5988 assert!(
5989 !response.data["message"]
5990 .as_str()
5991 .expect("message")
5992 .contains(" "),
5993 "the refusal must not expose source-formatting whitespace"
5994 );
5995 }
5996
5997 // ── capability drift guard ─────────────────────────────────────────
5998 //
5999 // The stdio `capabilities` method is the benchmark/SDK contract: external
6000 // harnesses probe it (without spending model tokens) to learn what the
6001 // app-server can do. Pin the advertised method set so any change forces a
6002 // deliberate update here, in the dispatcher, and in docs/RUNTIME_API.md.
6003
6004 /// Methods advertised by the top-level `capabilities` probe, in order.
6005 const EXPECTED_CAPABILITY_METHODS: &[&str] = &[
6006 "healthz",
6007 "thread/capabilities",
6008 "thread/request",
6009 "thread/create",
6010 "thread/start",
6011 "thread/resume",
6012 "thread/fork",
6013 "thread/list",
6014 "thread/read",
6015 "thread/set_name",
6016 "thread/goal/set",
6017 "thread/goal/get",
6018 "thread/goal/clear",
6019 "thread/archive",
6020 "thread/unarchive",
6021 "thread/message",
6022 "thread/interrupt",
6023 "app/capabilities",
6024 "app/request",
6025 "app/config/get",
6026 "app/config/set",
6027 "app/config/unset",
6028 "app/config/list",
6029 "app/config/reload",
6030 "app/models",
6031 "app/thread_loaded_list",
6032 "prompt/capabilities",
6033 "prompt/request",
6034 "prompt/run",
6035 "shutdown",
6036 ];
6037
6038 fn capability_test_state() -> (AppState, tempfile::TempDir) {
6039 let tmp = tempfile::tempdir().expect("tempdir");
6040 let config_path = tmp.path().join("config.toml");
6041 fs::write(&config_path, "").expect("write config");
6042 let state = build_state(Some(config_path), None).expect("state");
6043 (state, tmp)
6044 }
6045
6046 /// Persist a client thread under a fixed id, as `thread/create` would,
6047 /// so `thread/message` accepts it.
6048 fn test_client_metadata(thread_id: &str) -> codewhale_state::ThreadMetadata {
6049 codewhale_state::ThreadMetadata {
6050 id: thread_id.to_string(),
6051 rollout_path: None,
6052 preview: String::new(),
6053 ephemeral: false,
6054 model_provider: "deepseek".to_string(),
6055 created_at: 1,
6056 updated_at: 1,
6057 status: codewhale_state::ThreadStatus::Idle,
6058 path: None,
6059 cwd: PathBuf::from("/tmp/codewhale"),
6060 cli_version: "0.0.0-test".to_string(),
6061 source: codewhale_state::SessionSource::Api,
6062 name: None,
6063 sandbox_policy: None,
6064 approval_mode: None,
6065 archived: false,
6066 archived_at: None,
6067 git_sha: None,
6068 git_branch: None,
6069 git_origin_url: None,
6070 memory_mode: None,
6071 current_leaf_id: None,
6072 }
6073 }
6074
6075 fn restore_test_thread_archive(store: &StateStore, metadata: &codewhale_state::ThreadMetadata) {
6076 store
6077 .restore_legacy_thread_archive(&codewhale_state::LegacyThreadArchive {
6078 thread: metadata.clone(),
6079 messages: Vec::new(),
6080 goal: None,
6081 checkpoints: Vec::new(),
6082 })
6083 .expect("restore absent historical fixture");
6084 }
6085
6086 async fn seed_client_thread(state: &AppState, thread_id: &str) {
6087 let mut metadata = test_client_metadata(thread_id);
6088 if let Some(workspace) = state.frontend_workspace.as_ref() {
6089 metadata.cwd = workspace.clone();
6090 }
6091 let store = state.runtime.read().await.state_store().clone();
6092 restore_test_thread_archive(&store, &metadata);
6093 }
6094
6095 #[tokio::test]
6096 async fn capabilities_method_set_is_stable() {
6097 let (state, _tmp) = capability_test_state();
6098 let caps = dispatch_stdio_request(&state, "capabilities", json!({}))
6099 .await
6100 .expect("capabilities dispatch");
6101 let methods: Vec<String> = caps.result["methods"]
6102 .as_array()
6103 .expect("methods array")
6104 .iter()
6105 .map(|m| m.as_str().expect("method string").to_string())
6106 .collect();
6107 assert_eq!(
6108 methods, EXPECTED_CAPABILITY_METHODS,
6109 "app-server stdio capability set drifted; update the dispatcher, this \
6110 snapshot, and docs/RUNTIME_API.md together"
6111 );
6112 }
6113
6114 /// The socket transport advertises the `daemon/attach` handshake right
6115 /// after `healthz`; the stdio pin above must stay untouched by it.
6116 #[tokio::test]
6117 async fn socket_transport_advertises_daemon_attach() {
6118 let (state, _tmp) = capability_test_state();
6119 let mut sink = tokio::io::sink();
6120 let caps = dispatch_stdio_request_with_writer(
6121 &state,
6122 &mut sink,
6123 "capabilities",
6124 json!({}),
6125 AppTransport::Socket,
6126 )
6127 .await
6128 .expect("capabilities dispatch");
6129 let expected_transport = if cfg!(windows) {
6130 "named-pipe"
6131 } else {
6132 "unix-socket"
6133 };
6134 assert_eq!(caps.result["transport"], json!(expected_transport));
6135 let methods: Vec<String> = caps.result["methods"]
6136 .as_array()
6137 .expect("methods array")
6138 .iter()
6139 .map(|m| m.as_str().expect("method string").to_string())
6140 .collect();
6141 let mut expected: Vec<String> = EXPECTED_CAPABILITY_METHODS
6142 .iter()
6143 .map(|m| m.to_string())
6144 .collect();
6145 expected.insert(1, daemon_socket::ATTACH_METHOD.to_string());
6146 assert_eq!(methods, expected);
6147 }
6148
6149 #[tokio::test]
6150 async fn every_advertised_capability_is_dispatchable() {
6151 let (state, _tmp) = capability_test_state();
6152 // Empty params: methods may fail validation (-32602), but none may report
6153 // method-not-found (-32601). Required fields (e.g. PromptRequest.prompt)
6154 // make the prompt routes fail at parse time, so no model tokens are spent.
6155 for method in EXPECTED_CAPABILITY_METHODS {
6156 if let Err(err) = dispatch_stdio_request(&state, method, json!({})).await {
6157 assert_ne!(
6158 err.code,
6159 JsonRpcError::method_not_found(method).code,
6160 "advertised capability `{method}` is not dispatchable"
6161 );
6162 }
6163 }
6164 }
6165
6166 // ── resolve_auth_token ─────────────────────────────────────────────
6167
6168 #[test]
6169 fn auth_token_empty_string_fails() {
6170 let options = AppServerOptions {
6171 listen: "127.0.0.1:0".parse().expect("addr"),
6172 config_path: None,
6173 auth_token: Some(" ".to_string()),
6174 insecure_no_auth: false,
6175 cors_origins: Vec::new(),
6176 };
6177 let err = resolve_auth_token(&options).expect_err("empty token should fail");
6178 assert!(err.to_string().contains("cannot be empty"));
6179 }
6180
6181 #[test]
6182 fn auth_token_generated_when_none_provided() {
6183 let options = AppServerOptions {
6184 listen: "127.0.0.1:0".parse().expect("addr"),
6185 config_path: None,
6186 auth_token: None,
6187 insecure_no_auth: false,
6188 cors_origins: Vec::new(),
6189 };
6190 let token = resolve_auth_token(&options).unwrap();
6191 assert!(token.is_some());
6192 assert!(token.unwrap().starts_with("cwapp_"));
6193 }
6194
6195 #[test]
6196 fn runtime_child_endpoint_must_be_a_reported_loopback_port() {
6197 let ok = parse_runtime_endpoint("Runtime API listening on http://127.0.0.1:49152\r\n")
6198 .expect("loopback endpoint");
6199 assert_eq!(ok.port(), 49152);
6200 for bad in [
6201 "Runtime API listening on http://127.0.0.1:0",
6202 "Runtime API listening on http://10.0.0.5:7878",
6203 "Runtime API listening on http://[::1]:7878",
6204 "Runtime API listening on http://example.com:80",
6205 "Runtime API listening on http://127.0.0.1:80/redirect",
6206 "listening on http://127.0.0.1:7878",
6207 "",
6208 ] {
6209 assert!(parse_runtime_endpoint(bad).is_err(), "{bad:?}");
6210 }
6211
6212 let mut first_line_only: &[u8] =
6213 b"Runtime API listening on http://127.0.0.1:5000\nRuntime API listening on http://127.0.0.1:6000\n";
6214 assert_eq!(
6215 read_runtime_ready_line(&mut first_line_only).unwrap(),
6216 "Runtime API listening on http://127.0.0.1:5000"
6217 );
6218 let mut closed: &[u8] = b"Runtime API listening on http://127.0.0.1:5000";
6219 assert!(read_runtime_ready_line(&mut closed).is_err(), "no newline");
6220 let oversized = vec![b'a'; RUNTIME_READY_MAX_BYTES + 1];
6221 assert!(read_runtime_ready_line(&mut oversized.as_slice()).is_err());
6222 }
6223
6224 #[test]
6225 fn runtime_bridge_command_keeps_auth_token_out_of_argv() {
6226 // FR001-C001: runtime auth token must not appear on the child argv
6227 // (visible via local `ps`); pass it via env instead.
6228 let token = "cwrt_unit_test_secret_token_not_for_argv";
6229 let cmd = RuntimeBridge::runtime_command(None, token).expect("command");
6230 let argv: Vec<String> = cmd
6231 .get_args()
6232 .map(|a| a.to_string_lossy().into_owned())
6233 .collect();
6234 assert!(
6235 argv.windows(2).any(|pair| pair == ["--port", "0"]),
6236 "the child picks and reports its own port: {argv:?}"
6237 );
6238 assert!(
6239 !argv
6240 .iter()
6241 .any(|a| a.contains(token) || a == "--auth-token"),
6242 "auth token must not be present in child argv: {argv:?}"
6243 );
6244 let envs: Vec<(String, String)> = cmd
6245 .get_envs()
6246 .filter_map(|(k, v)| {
6247 Some((
6248 k.to_string_lossy().into_owned(),
6249 v?.to_string_lossy().into_owned(),
6250 ))
6251 })
6252 .collect();
6253 assert!(
6254 envs.iter()
6255 .any(|(k, v)| k == "CODEWHALE_RUNTIME_TOKEN" && v == token),
6256 "token must be carried via CODEWHALE_RUNTIME_TOKEN: {envs:?}"
6257 );
6258 assert!(
6259 envs.iter()
6260 .any(|(k, v)| k == "DEEPSEEK_RUNTIME_TOKEN" && v == token),
6261 "legacy alias DEEPSEEK_RUNTIME_TOKEN must also carry the token: {envs:?}"
6262 );
6263 }
6264
6265 #[test]
6266 fn generated_auth_status_does_not_render_token() {
6267 let rendered = app_server_auth_status_lines(false).join("\n");
6268
6269 assert!(!rendered.contains("Authorization: Bearer"));
6270 assert!(rendered.contains("not printed"));
6271 assert!(rendered.contains("CODEWHALE_APP_SERVER_TOKEN"));
6272 }
6273
6274 #[test]
6275 fn auth_token_explicit_is_preserved() {
6276 let options = AppServerOptions {
6277 listen: "127.0.0.1:0".parse().expect("addr"),
6278 config_path: None,
6279 auth_token: Some("my-secret".to_string()),
6280 insecure_no_auth: false,
6281 cors_origins: Vec::new(),
6282 };
6283 let token = resolve_auth_token(&options).unwrap();
6284 assert_eq!(token.as_deref(), Some("my-secret"));
6285 }
6286
6287 #[test]
6288 fn auth_token_explicit_allows_non_loopback_bind() {
6289 let options = AppServerOptions {
6290 listen: "0.0.0.0:8787".parse().expect("socket addr"),
6291 config_path: None,
6292 auth_token: Some("my-secret".to_string()),
6293 insecure_no_auth: false,
6294 cors_origins: Vec::new(),
6295 };
6296 let token = resolve_auth_token(&options).unwrap();
6297 assert_eq!(token.as_deref(), Some("my-secret"));
6298 }
6299
6300 #[test]
6301 fn insecure_no_auth_on_loopback_returns_none() {
6302 let options = AppServerOptions {
6303 listen: "127.0.0.1:0".parse().expect("addr"),
6304 config_path: None,
6305 auth_token: None,
6306 insecure_no_auth: true,
6307 cors_origins: Vec::new(),
6308 };
6309 let token = resolve_auth_token(&options).unwrap();
6310 assert!(token.is_none());
6311 }
6312
6313 #[test]
6314 fn insecure_no_auth_on_non_loopback_fails_fast() {
6315 let options = AppServerOptions {
6316 listen: "0.0.0.0:8787".parse().expect("socket addr"),
6317 config_path: None,
6318 auth_token: None,
6319 insecure_no_auth: true,
6320 cors_origins: Vec::new(),
6321 };
6322
6323 let err = resolve_auth_token(&options).expect_err("non-loopback unauth should fail");
6324 assert!(
6325 err.to_string()
6326 .contains("refusing unauthenticated app-server bind")
6327 );
6328 }
6329
6330 // ── cors_layer ─────────────────────────────────────────────────────
6331
6332 #[test]
6333 fn cors_layer_includes_default_origins() {
6334 let layer = cors_layer(&[]);
6335 // Just verify it doesn't panic and creates successfully
6336 let _ = layer;
6337 }
6338
6339 #[test]
6340 fn cors_layer_adds_extra_origins() {
6341 let extras = vec!["https://example.com".to_string()];
6342 let layer = cors_layer(&extras);
6343 let _ = layer;
6344 }
6345
6346 #[test]
6347 fn cors_layer_skips_empty_origins() {
6348 let extras = vec!["".to_string(), " ".to_string()];
6349 let layer = cors_layer(&extras);
6350 let _ = layer;
6351 }
6352
6353 // ── JsonRpc helpers ────────────────────────────────────────────────
6354
6355 #[test]
6356 fn params_or_object_returns_object_for_null() {
6357 let result = params_or_object(Value::Null);
6358 assert_eq!(result, json!({}));
6359 }
6360
6361 #[test]
6362 fn params_or_object_passthrough_for_non_null() {
6363 let input = json!({"key": "value"});
6364 let result = params_or_object(input.clone());
6365 assert_eq!(result, input);
6366 }
6367
6368 #[test]
6369 fn jsonrpc_result_format() {
6370 let result = jsonrpc_result(Some(json!(1)), json!({"ok": true}));
6371 assert_eq!(result["jsonrpc"], "2.0");
6372 assert_eq!(result["id"], 1);
6373 assert_eq!(result["result"]["ok"], true);
6374 }
6375
6376 #[test]
6377 fn jsonrpc_result_null_id() {
6378 let result = jsonrpc_result(None, json!(null));
6379 assert_eq!(result["id"], Value::Null);
6380 }
6381
6382 #[test]
6383 fn jsonrpc_error_format() {
6384 let err = jsonrpc_error(Some(json!(2)), JsonRpcError::internal("oops"));
6385 assert_eq!(err["jsonrpc"], "2.0");
6386 assert_eq!(err["id"], 2);
6387 assert_eq!(err["error"]["code"], -32603);
6388 assert_eq!(err["error"]["message"], "oops");
6389 }
6390
6391 #[test]
6392 fn jsonrpc_error_codes() {
6393 assert_eq!(JsonRpcError::parse_error("").code, -32700);
6394 assert_eq!(JsonRpcError::invalid_request("").code, -32600);
6395 assert_eq!(JsonRpcError::method_not_found("x").code, -32601);
6396 assert_eq!(JsonRpcError::invalid_params("").code, -32602);
6397 assert_eq!(JsonRpcError::internal("").code, -32603);
6398 }
6399
6400 // ── AppServerOptions ───────────────────────────────────────────────
6401
6402 #[test]
6403 fn app_server_options_debug_does_not_leak_token() {
6404 let options = AppServerOptions {
6405 listen: "127.0.0.1:8080".parse().expect("addr"),
6406 config_path: None,
6407 auth_token: Some("secret-token".to_string()),
6408 insecure_no_auth: false,
6409 cors_origins: vec!["https://example.com".to_string()],
6410 };
6411 let debug = format!("{options:?}");
6412 assert!(!debug.contains("secret-token"));
6413 assert!(debug.contains("<redacted>"));
6414 assert!(debug.contains("8080"));
6415 }
6416
6417 // ── Default CORS origins ──────────────────────────────────────────
6418
6419 #[test]
6420 fn default_cors_origins_include_common_dev_ports() {
6421 assert!(DEFAULT_CORS_ORIGINS.contains(&"http://localhost:3000"));
6422 assert!(DEFAULT_CORS_ORIGINS.contains(&"http://localhost:5173"));
6423 assert!(DEFAULT_CORS_ORIGINS.contains(&"tauri://localhost"));
6424 }
6425 #[tokio::test]
6426 async fn runtime_image_daemon_bridge_checks_transport_and_forwards_exact_wire() {
6427 async fn capture(
6428 State(seen): State<Arc<Mutex<Vec<Value>>>>,
6429 Json(body): Json<Value>,
6430 ) -> (StatusCode, Json<Value>) {
6431 seen.lock().await.push(body);
6432 (
6433 StatusCode::BAD_REQUEST,
6434 Json(json!({"error":"fixture stops before an Engine"})),
6435 )
6436 }
6437 for supported in [false, true] {
6438 let seen = Arc::new(Mutex::new(Vec::new()));
6439 let app = Router::new()
6440 .route(
6441 "/v1/runtime/info",
6442 get(move || async move {
6443 Json(json!({"capabilities":{"turn_image_inputs":supported}}))
6444 }),
6445 )
6446 .route("/v1/threads/{id}/turns", post(capture))
6447 .with_state(seen.clone());
6448 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
6449 let addr = listener.local_addr().unwrap();
6450 let server = tokio::spawn(async move {
6451 axum::serve(listener, app).await.unwrap();
6452 });
6453 let mut bridge = RuntimeBridge::from_base_url_for_test(format!("http://{addr}"));
6454 let images = vec![RuntimeImageInput {
6455 mime: "image/png".into(),
6456 data_base64: "fixture-bytes-validated-by-Core".into(),
6457 }];
6458 let mut writer = tokio::io::sink();
6459 assert!(
6460 bridge
6461 .message_thread(
6462 "thr_fixture",
6463 RuntimeTurnInput {
6464 input: "look",
6465 images: &images,
6466 max_output_tokens: None,
6467 expected_workspace: None
6468 },
6469 &mut writer,
6470 None,
6471 None
6472 )
6473 .await
6474 .is_err()
6475 );
6476 let requests = seen.lock().await;
6477 assert_eq!(requests.len(), usize::from(supported));
6478 if supported {
6479 assert_eq!(requests[0], json!({"prompt":"look","images":images}));
6480 }
6481 server.abort();
6482 }
6483 }
6484
6485 #[test]
6486 fn runtime_image_daemon_all_input_families_preserve_images() {
6487 let image = json!({"mime":"image/png","dataBase64":"AQ=="});
6488 let thread: ThreadMessageParams = serde_json::from_value(
6489 json!({"thread_id":"thr_fixture","input":"look","images":[image.clone()]}),
6490 )
6491 .unwrap();
6492 let prompt: PromptRequest =
6493 serde_json::from_value(json!({"prompt":"look","images":[image.clone()]})).unwrap();
6494 let generic: ThreadRequest = serde_json::from_value(
6495 json!({"kind":"message","thread_id":"thr_fixture","input":"look","images":[image]}),
6496 )
6497 .unwrap();
6498 assert_eq!(thread.images, prompt.images);
6499 let ThreadRequest::Message { images, .. } = generic else {
6500 panic!("message");
6501 };
6502 assert_eq!(thread.images, images);
6503 assert!(matches!(
6504 parse_stdio_line(&" ".repeat(MAX_RUNTIME_IMAGE_BODY_BYTES + 1)),
6505 ParsedStdioLine::Rejected(_)
6506 ));
6507 }
6508 }
6509
6509 lines RUST