返回 CodeWhale
acp_server.rs
根目录 / crates / tui / src / acp_server.rs
1 //! ACP stdio projection of the existing Runtime/Engine owner.
2 //! History, provider requests, tools, approval and cancellation are Core facts.
3 //! This module keeps only connection configuration, canonical bindings and
4 //! JSON-RPC correlation. Secure cross-process attachment is not yet qualified.
5 use crate::config::{Config, ProviderKind};
6 #[cfg(test)]
7 use crate::runtime_threads::RuntimeThreadManagerConfig;
8 use crate::runtime_threads::{CreateThreadRequest, RuntimeThreadManager, UpdateThreadRequest};
9 use anyhow::{Result, anyhow};
10 use codewhale_config::AppMode;
11 use codewhale_execpolicy::ApprovalMode;
12 use serde_json::{Value, json};
13 use std::collections::{HashMap, VecDeque};
14 use std::path::PathBuf;
15 use std::sync::Arc;
16 use std::sync::atomic::{AtomicU64, Ordering};
17 #[cfg(test)]
18 use tokio::io::BufReader;
19 use tokio::io::{AsyncBufRead, AsyncWrite, AsyncWriteExt};
20 mod runtime;
21 const ACP_PROTOCOL_VERSION: u64 = 1;
22 const MAX_ACP_SESSIONS: usize = 64;
23 const TOOL_CALL_CONTENT_PREVIEW_CHARS: usize = 4_000;
24 static NEXT_ACP_PERMISSION_REQUEST_ID: AtomicU64 = AtomicU64::new(1);
25
26 pub async fn run_acp_server(
27 config: Config,
28 model: String,
29 default_cwd: PathBuf,
30 plugins: Arc<crate::plugins::PluginDiscoveryContext>,
31 config_path: Option<PathBuf>,
32 config_profile: Option<String>,
33 ) -> Result<()> {
34 crate::runtime_api::run_http_server(
35 config,
36 default_cwd,
37 plugins,
38 crate::runtime_api::RuntimeApiOptions {
39 port: 0,
40 config_path,
41 config_profile,
42 control_frontend: Some(codewhale_app_server::RuntimeControlFrontend::Acp { model }),
43 ..Default::default()
44 },
45 )
46 .await
47 }
48
49 pub(crate) struct CapturedAcpFrontend {
50 config: Config,
51 model: String,
52 cwd: PathBuf,
53 manager: Arc<RuntimeThreadManager>,
54 sessions_dir: PathBuf,
55 config_path: Option<PathBuf>,
56 config_profile: Option<String>,
57 }
58 pub(crate) fn capture_frontend(
59 config: Config,
60 model: String,
61 cwd: PathBuf,
62 manager: Arc<RuntimeThreadManager>,
63 sessions_dir: PathBuf,
64 config_path: Option<PathBuf>,
65 config_profile: Option<String>,
66 ) -> Result<Arc<CapturedAcpFrontend>> {
67 Ok(Arc::new(CapturedAcpFrontend {
68 config,
69 model,
70 cwd,
71 manager,
72 sessions_dir,
73 config_path,
74 config_profile,
75 }))
76 }
77 impl CapturedAcpFrontend {
78 pub(crate) fn serve(
79 &self,
80 input: Box<dyn AsyncBufRead + Send + Unpin>,
81 mut output: Box<dyn AsyncWrite + Send + Unpin>,
82 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = Result<()>> + Send + '_>> {
83 Box::pin(async move {
84 let mut server = AcpServer::new(
85 self.config.clone(),
86 self.model.clone(),
87 self.cwd.clone(),
88 self.manager.clone(),
89 self.sessions_dir.clone(),
90 self.config_path.clone(),
91 self.config_profile.clone(),
92 );
93 let mut input = codewhale_app_server::BoundedLines::new(input);
94 server.serve(&mut input, &mut output).await
95 })
96 }
97 }
98
99 #[cfg(test)]
100 impl codewhale_app_server::RuntimeOwnerFrontend for CapturedAcpFrontend {
101 fn validate_selection<'a>(
102 &'a self,
103 selection: &'a codewhale_app_server::RuntimeOwnerFrontendSelection,
104 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = Result<()>> + Send + 'a>> {
105 Box::pin(async move {
106 anyhow::ensure!(
107 matches!(
108 selection,
109 codewhale_app_server::RuntimeOwnerFrontendSelection::Acp { .. }
110 ),
111 "ACP fixture does not capture HTTP services"
112 );
113 Ok(())
114 })
115 }
116 fn serve(
117 &self,
118 selection: codewhale_app_server::RuntimeOwnerFrontendSelection,
119 _compatibility: codewhale_app_server::AppState,
120 input: Box<dyn AsyncBufRead + Send + Unpin>,
121 output: Box<dyn AsyncWrite + Send + Unpin>,
122 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = Result<()>> + Send + '_>> {
123 Box::pin(async move {
124 anyhow::ensure!(
125 matches!(
126 selection,
127 codewhale_app_server::RuntimeOwnerFrontendSelection::Acp { .. }
128 ),
129 "ACP fixture does not capture HTTP services"
130 );
131 CapturedAcpFrontend::serve(self, input, output).await
132 })
133 }
134 }
135
136 struct AcpServer {
137 config: Config,
138 model: String,
139 default_cwd: PathBuf,
140 runtime: Arc<RuntimeThreadManager>,
141 sessions_dir: PathBuf,
142 config_path: Option<PathBuf>,
143 config_profile: Option<String>,
144 sessions: HashMap<String, AcpSession>,
145 insertion_order: VecDeque<String>,
146 client_supports_terminal: bool,
147 response_id_policy: JsonRpcResponseIdPolicy,
148 }
149 /// A transport binding, never another conversation or executable registry.
150 struct AcpSession {
151 thread_id: String,
152 cursor: u64,
153 }
154 struct PendingToolCall {
155 execution_id: String,
156 name: String,
157 input: Value,
158 }
159 enum AcpDispatch {
160 Response(Value),
161 Shutdown,
162 }
163 #[derive(Debug)]
164 struct AcpError {
165 code: i32,
166 message: String,
167 }
168 impl AcpServer {
169 fn new(
170 config: Config,
171 model: String,
172 default_cwd: PathBuf,
173 runtime: Arc<RuntimeThreadManager>,
174 sessions_dir: PathBuf,
175 config_path: Option<PathBuf>,
176 config_profile: Option<String>,
177 ) -> Self {
178 Self {
179 config,
180 model,
181 default_cwd,
182 runtime,
183 sessions_dir,
184 config_path,
185 config_profile,
186 sessions: HashMap::new(),
187 insertion_order: VecDeque::new(),
188 client_supports_terminal: false,
189 response_id_policy: JsonRpcResponseIdPolicy::Preserve,
190 }
191 }
192 async fn serve<R: AsyncBufRead + Unpin, W: AsyncWrite + Unpin>(
193 &mut self,
194 reader: &mut codewhale_app_server::BoundedLines<R>,
195 writer: &mut W,
196 ) -> Result<()> {
197 while let Some(line) = reader.next_line().await? {
198 if line.trim().is_empty() {
199 continue;
200 }
201 let message: Value = match serde_json::from_str(&line) {
202 Ok(value) => value,
203 Err(error) => {
204 write_jsonrpc_error(writer, None, -32700, format!("invalid json: {error}"))
205 .await?;
206 continue;
207 }
208 };
209 if codewhale_app_server::is_control_input_closed(&message) {
210 break;
211 }
212 let response_id = message
213 .get("id")
214 .cloned()
215 .map(|id| self.response_id_policy.response_id(id));
216 if message.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
217 write_jsonrpc_error(writer, response_id, -32600, "jsonrpc version must be 2.0")
218 .await?;
219 continue;
220 }
221 let Some(method) = message.get("method").and_then(Value::as_str) else {
222 if is_jsonrpc_response(&message) {
223 continue;
224 }
225 write_jsonrpc_error(writer, response_id, -32600, "missing method").await?;
226 continue;
227 };
228 let params = message.get("params").cloned().unwrap_or_else(|| json!({}));
229 if method == "session/prompt" {
230 if let Err(error) = self.validate_prompt(&params) {
231 write_jsonrpc_error(writer, response_id, error.code, error.message).await?;
232 continue;
233 }
234 let result = self.drive_prompt(params, reader, writer).await;
235 match result {
236 Ok(reason) => {
237 if let Some(id) = response_id {
238 write_jsonrpc_result(writer, id, json!({"stopReason":reason})).await?;
239 }
240 }
241 Err(error) => {
242 write_jsonrpc_error(
243 writer,
244 response_id,
245 -32603,
246 crate::client::redact_model_bound_text(&error.to_string(), &[]),
247 )
248 .await?
249 }
250 }
251 continue;
252 }
253 match self.handle_request(method, params).await {
254 Ok(AcpDispatch::Response(result)) => {
255 if let Some(id) = response_id {
256 write_jsonrpc_result(writer, id, result).await?;
257 }
258 }
259 Ok(AcpDispatch::Shutdown) => {
260 if let Some(id) = response_id {
261 write_jsonrpc_result(writer, id, json!(null)).await?;
262 }
263 break;
264 }
265 Err(error) => {
266 write_jsonrpc_error(writer, response_id, error.code, error.message).await?
267 }
268 }
269 }
270 Ok(())
271 }
272 // `session/prompt` is handled in the main loop (it needs to run concurrently
273 // with the reader for cancellation); every other method is request/response.
274 async fn handle_request(
275 &mut self,
276 method: &str,
277 params: Value,
278 ) -> std::result::Result<AcpDispatch, AcpError> {
279 match method {
280 "initialize" => {
281 if let Some(terminal) = params
282 .pointer("/clientCapabilities/terminal")
283 .and_then(Value::as_bool)
284 {
285 self.client_supports_terminal = terminal;
286 }
287 self.response_id_policy = JsonRpcResponseIdPolicy::from_initialize_params(&params);
288 Ok(AcpDispatch::Response(initialize_result(
289 params.get("protocolVersion").and_then(Value::as_u64),
290 &self.config,
291 )))
292 }
293 "session/new" => Ok(AcpDispatch::Response(self.new_session(params).await?)),
294 "session/list" => Ok(AcpDispatch::Response(self.list_sessions(params)?)),
295 "session/load" => Ok(AcpDispatch::Response(self.load_session(params).await?)),
296 "session/listProviders" => Ok(AcpDispatch::Response(self.list_providers())),
297 "session/currentModel" => Ok(AcpDispatch::Response(self.current_model())),
298 "session/selectModel" => Ok(AcpDispatch::Response(self.select_model(params)?)),
299 "session/set_config_option" => Ok(AcpDispatch::Response(
300 self.set_session_config(params).await?,
301 )),
302 "session/set_mode" | "session/set_model" => {
303 let (config_id, field) = if method == "session/set_mode" {
304 ("mode", "modeId")
305 } else {
306 ("model", "modelId")
307 };
308 self.set_session_config(json!({
309 "sessionId": params.get("sessionId"),
310 "configId": config_id,
311 "value": params.get(field),
312 }))
313 .await?;
314 Ok(AcpDispatch::Response(json!({})))
315 }
316 // A cancel that arrives with no prompt in flight is an idempotent
317 // no-op (the in-flight case is handled by the prompt driver).
318 "session/cancel" => Ok(AcpDispatch::Response(json!(null))),
319 "shutdown" => Ok(AcpDispatch::Shutdown),
320 _ => Err(AcpError::method_not_found(method)),
321 }
322 }
323
324 fn validate_prompt(&self, params: &Value) -> std::result::Result<(String, String), AcpError> {
325 let id = params
326 .get("sessionId")
327 .and_then(Value::as_str)
328 .ok_or_else(|| AcpError::invalid_params("sessionId is required"))?;
329 if !self.sessions.contains_key(id) {
330 return Err(AcpError::invalid_params("unknown sessionId"));
331 }
332 let prompt = extract_prompt_text(params.get("prompt"))
333 .filter(|text| !text.trim().is_empty())
334 .ok_or_else(|| AcpError::invalid_params("prompt must include text content"))?;
335 Ok((id.to_string(), prompt))
336 }
337 fn shell_allowed(&self) -> bool {
338 self.client_supports_terminal && self.config.allow_shell()
339 }
340 fn remember(&mut self, id: String, thread_id: String, cursor: u64) {
341 if self.sessions.contains_key(&id) {
342 return;
343 }
344 if self.sessions.len() >= MAX_ACP_SESSIONS
345 && let Some(old) = self.insertion_order.pop_front()
346 {
347 self.sessions.remove(&old);
348 }
349 self.insertion_order.push_back(id.clone());
350 self.sessions.insert(id, AcpSession { thread_id, cursor });
351 }
352 async fn new_session(&mut self, params: Value) -> std::result::Result<Value, AcpError> {
353 let cwd = params
354 .get("cwd")
355 .and_then(Value::as_str)
356 .map(PathBuf::from)
357 .unwrap_or_else(|| self.default_cwd.clone());
358 let identity = self
359 .config
360 .active_provider_identity()
361 .map_err(|error| AcpError::internal(error.to_string()))?;
362 let thread = self
363 .runtime
364 .create_thread_with_shell_policy(
365 CreateThreadRequest {
366 model: Some(self.model.clone()),
367 model_provider: Some(identity.persisted_kind().to_string()),
368 model_provider_id: identity.persisted_id().map(str::to_string),
369 workspace: Some(cwd),
370 mode: Some(
371 if acp_mode(&self.config) == AppMode::Plan {
372 "plan"
373 } else {
374 "agent"
375 }
376 .to_string(),
377 ),
378 permission_posture: Some(self.permission_value().to_string()),
379 allow_shell: Some(self.shell_allowed()),
380 ..Default::default()
381 },
382 self.config_path.as_deref(),
383 self.config_profile.as_deref(),
384 )
385 .await
386 .map_err(|error| AcpError::internal(error.to_string()))?;
387 let id = uuid::Uuid::new_v4().to_string();
388 crate::runtime_api::sessions::initialize_empty_session(
389 &self.runtime,
390 &self.sessions_dir,
391 &thread.id,
392 &id,
393 )
394 .await
395 .map_err(|error| AcpError::internal(error.message))?;
396 let cursor = self
397 .runtime
398 .get_thread_detail(&thread.id)
399 .await
400 .map_err(|error| AcpError::internal(error.to_string()))?
401 .latest_seq;
402 self.remember(id.clone(), thread.id, cursor);
403 self.session_configuration(&id).await
404 }
405 /// Durable Codewhale sessions an ACP client can resume (#5864).
406 ///
407 /// ACP transport bindings are in-memory and capped; Core sessions are the
408 /// durable record, and an IDE that offers "resume" means those. A store
409 /// that cannot be read is an empty list, not a failed request: enumeration
410 /// is discovery, and a client asking what exists should not be broken by a
411 /// missing sessions directory.
412 fn list_sessions(&self, params: Value) -> std::result::Result<Value, AcpError> {
413 let cwd = match params.get("cwd") {
414 None | Some(Value::Null) => None,
415 Some(Value::String(path)) if std::path::Path::new(path).is_absolute() => {
416 Some(PathBuf::from(path))
417 }
418 Some(_) => {
419 return Err(AcpError::invalid_params(
420 "session/list cwd must be an absolute path",
421 ));
422 }
423 };
424 let sessions = crate::session_manager::SessionManager::new(self.sessions_dir.clone())
425 .ok()
426 .and_then(|manager| manager.list_sessions().ok())
427 .unwrap_or_default();
428 let sessions: Vec<Value> = sessions
429 .into_iter()
430 .filter(|meta| {
431 cwd.as_ref().is_none_or(|cwd| {
432 crate::session_manager::paths_equivalent(&meta.workspace, cwd)
433 })
434 })
435 .map(|meta| {
436 json!({
437 "sessionId": meta.id,
438 "title": meta.title,
439 "cwd": meta.workspace.to_string_lossy(),
440 "createdAt": meta.created_at.to_rfc3339(),
441 "updatedAt": meta.updated_at.to_rfc3339(),
442 "messageCount": meta.message_count,
443 })
444 })
445 .collect();
446 Ok(json!({ "sessions": sessions }))
447 }
448
449 async fn load_session(&mut self, params: Value) -> std::result::Result<Value, AcpError> {
450 let requested = params
451 .get("sessionId")
452 .and_then(Value::as_str)
453 .ok_or_else(|| AcpError::invalid_params("session/load requires sessionId"))?;
454 if self.sessions.contains_key(requested) {
455 return self.session_configuration(requested).await;
456 }
457 let store = crate::session_manager::SessionManager::new(self.sessions_dir.clone())
458 .map_err(|error| AcpError::internal(error.to_string()))?;
459 let saved = store
460 .resume_session_by_prefix(requested)
461 .map_err(|error| {
462 AcpError::invalid_params(format!("could not load session {requested}: {error}"))
463 })?
464 .session;
465 let id = saved.metadata.id;
466 if self.sessions.contains_key(&id) {
467 return self.session_configuration(&id).await;
468 }
469 let (_, axum::Json(resumed)) = crate::runtime_api::sessions::resume_session_in_runtime(
470 &self.runtime,
471 &self.sessions_dir,
472 &id,
473 crate::runtime_api::sessions::ResumeSessionRequest {
474 model: None,
475 mode: Some(
476 if acp_mode(&self.config) == AppMode::Plan {
477 "plan"
478 } else {
479 saved.metadata.mode.as_deref().unwrap_or("agent")
480 }
481 .to_string(),
482 ),
483 },
484 (self.config_path.as_deref(), self.config_profile.as_deref()),
485 Some(self.shell_allowed()),
486 )
487 .await
488 .map_err(|error| AcpError::internal(error.message))?;
489 // A saved posture cannot widen the server, and connection terminal
490 // negotiation only narrows the canonical thread's shell ceiling.
491 self.runtime
492 .update_thread_with_shell_policy(
493 &resumed.thread_id,
494 UpdateThreadRequest {
495 permission_posture: Some(self.permission_value().to_string()),
496 allow_shell: Some(self.shell_allowed()),
497 ..Default::default()
498 },
499 self.config_path.as_deref(),
500 self.config_profile.as_deref(),
501 )
502 .await
503 .map_err(|error| AcpError::internal(error.to_string()))?;
504 let cursor = self
505 .runtime
506 .get_thread_detail(&resumed.thread_id)
507 .await
508 .map_err(|error| AcpError::internal(error.to_string()))?
509 .latest_seq;
510 self.remember(id.clone(), resumed.thread_id, cursor);
511 self.session_configuration(&id).await
512 }
513 fn permission_value(&self) -> &'static str {
514 match acp_approval_mode(&self.config) {
515 ApprovalMode::Bypass => "full-access",
516 ApprovalMode::Auto => "auto-review",
517 ApprovalMode::Never => "never",
518 ApprovalMode::Suggest => "ask",
519 }
520 }
521 async fn session_configuration(
522 &self,
523 session_id: &str,
524 ) -> std::result::Result<Value, AcpError> {
525 use codewhale_localization::{MessageId, resolve_locale, tr};
526 let settings = crate::settings::Settings::load().unwrap_or_default();
527 let locale = resolve_locale(&settings.locale);
528 let binding = self
529 .sessions
530 .get(session_id)
531 .ok_or_else(|| AcpError::invalid_params("unknown sessionId"))?;
532 let thread = self
533 .runtime
534 .get_thread(&binding.thread_id)
535 .await
536 .map_err(|error| AcpError::internal(error.to_string()))?;
537 let identity = self
538 .config
539 .resolve_persisted_provider_identity(
540 thread.model_provider.as_deref(),
541 thread.model_provider_id.as_deref(),
542 )
543 .map_err(AcpError::invalid_params)?;
544 let mut models = crate::provider_lake::models_for_provider(&self.config, &identity);
545 if !models.contains(&thread.model) {
546 models.push(thread.model.clone());
547 }
548 let mut modes = vec![
549 json!({"id": "plan", "name": tr(locale, MessageId::AppModePlan), "description": tr(locale, MessageId::AppModePlanHint)}),
550 ];
551 // #6310: the permission posture is server-owned (a client can never
552 // relax it), but it must be discoverable. Work under Full Access must
553 // not claim that edits ask for approval, and the posture is surfaced
554 // below as a read-only select that names how Full Access is enabled.
555 let posture = acp_approval_mode(&self.config);
556 let agent_hint = if posture == ApprovalMode::Bypass {
557 tr(locale, MessageId::HomeYoloModeTip)
558 } else {
559 tr(locale, MessageId::AppModeAgentHint)
560 };
561 if acp_mode(&self.config) != AppMode::Plan {
562 modes.insert(0, json!({"id": "agent", "name": tr(locale, MessageId::AppModeAgent), "description": agent_hint}));
563 }
564 let current_mode = if thread.mode == "plan" {
565 "plan"
566 } else {
567 "agent"
568 };
569 let (posture_value, posture_name, posture_description) = match posture {
570 ApprovalMode::Bypass => (
571 "full-access",
572 MessageId::ConfigChoiceFullAccess,
573 MessageId::PermissionsPostureBypass,
574 ),
575 ApprovalMode::Auto => (
576 "auto-review",
577 MessageId::ConfigChoiceAutoReview,
578 MessageId::PermissionsPostureAuto,
579 ),
580 ApprovalMode::Never => (
581 "never",
582 MessageId::ConfigChoiceNever,
583 MessageId::PermissionsPostureNever,
584 ),
585 ApprovalMode::Suggest => (
586 "ask",
587 MessageId::ConfigChoiceAsk,
588 MessageId::PermissionsPostureAsk,
589 ),
590 };
591 Ok(json!({
592 "sessionId": session_id,
593 "modes": {"currentModeId": current_mode, "availableModes": modes},
594 "models": {
595 "currentModelId": thread.model,
596 "availableModels": models.iter().map(|model| json!({"modelId": model, "name": model})).collect::<Vec<_>>()
597 },
598 "configOptions": [
599 {"id": "mode", "name": tr(locale, MessageId::SettingSubjectMode), "category": "mode", "type": "select", "currentValue": current_mode,
600 "options": modes.iter().map(|mode| json!({"value": mode["id"], "name": mode["name"], "description": mode["description"]})).collect::<Vec<_>>()},
601 {"id": "model", "name": tr(locale, MessageId::SettingSubjectModel), "category": "model", "type": "select", "currentValue": thread.model,
602 "options": models.iter().map(|model| json!({"value": model, "name": model})).collect::<Vec<_>>()},
603 // Exactly one option: the posture the server was started
604 // with. Offering a looser value here would let a client relax
605 // the operator's floor.
606 {"id": "permission", "name": tr(locale, MessageId::SettingSubjectPermissions), "category": "_permission", "type": "select", "currentValue": posture_value,
607 "options": [{"value": posture_value, "name": tr(locale, posture_name), "description": tr(locale, posture_description)}],
608 "_meta": {"codewhale": {
609 "readOnly": true,
610 "fullAccess": posture == ApprovalMode::Bypass,
611 "enableFullAccess": ACP_FULL_ACCESS_HINT,
612 }}}
613 ]
614 }))
615 }
616
617 async fn set_session_config(&mut self, params: Value) -> std::result::Result<Value, AcpError> {
618 let session_id = params
619 .get("sessionId")
620 .and_then(Value::as_str)
621 .ok_or_else(|| AcpError::invalid_params("sessionId is required"))?;
622 let config_id = params
623 .get("configId")
624 .and_then(Value::as_str)
625 .ok_or_else(|| AcpError::invalid_params("configId is required"))?;
626 let value = params
627 .get("value")
628 .and_then(Value::as_str)
629 .ok_or_else(|| AcpError::invalid_params("value must be an offered string option"))?;
630 if !self.sessions.contains_key(session_id) {
631 return Err(AcpError::invalid_params("unknown sessionId"));
632 }
633 let state = self.session_configuration(session_id).await?;
634 let offered = state["configOptions"]
635 .as_array()
636 .unwrap()
637 .iter()
638 .any(|option| {
639 option["id"] == config_id
640 && option["options"]
641 .as_array()
642 .unwrap()
643 .iter()
644 .any(|choice| choice["value"] == value)
645 });
646 if !offered {
647 return Err(AcpError::invalid_params(
648 "unknown configuration option or value",
649 ));
650 }
651 let binding = &self.sessions[session_id];
652 let request = match config_id {
653 "model" => UpdateThreadRequest {
654 model: Some(value.to_string()),
655 ..Default::default()
656 },
657 "mode" => UpdateThreadRequest {
658 mode: Some(value.to_string()),
659 allow_shell: Some(value != "plan" && self.shell_allowed()),
660 ..Default::default()
661 },
662 "permission" => return Ok(json!({"configOptions": state["configOptions"]})),
663 _ => unreachable!("validated offered option"),
664 };
665 self.runtime
666 .update_thread_with_shell_policy(
667 &binding.thread_id,
668 request,
669 self.config_path.as_deref(),
670 self.config_profile.as_deref(),
671 )
672 .await
673 .map_err(|error| AcpError::internal(error.to_string()))?;
674 Ok(json!({"configOptions": self.session_configuration(session_id).await?["configOptions"]}))
675 }
676
677 fn list_providers(&self) -> Value {
678 let mut providers = self.config.provider_identities().into_iter().filter(|identity| identity.provider != ProviderKind::Antigravity).map(|identity| {
679 json!({"id": identity.key, "displayName": identity.compatibility().map(|row| row.label).unwrap_or(identity.key.as_str()), "defaultModel": crate::model_inventory::provider_default_model(&self.config, &identity)})
680 }).collect::<Vec<_>>();
681 providers.extend(self.config.unadmitted_provider_keys().into_iter().map(|key| json!({"id": key, "displayName": format!("{key} (unavailable)"), "defaultModel": "", "available": false, "reason": "provider_identity_unavailable"})));
682 json!({"providers": providers})
683 }
684
685 fn current_model(&self) -> Value {
686 // Prefer the raw configured provider key so a custom `[providers.<name>]`
687 // entry round-trips through ACP instead of canonicalizing to "custom".
688 let provider = self
689 .config
690 .active_provider_identity()
691 .ok()
692 .map(|identity| identity.key.to_string())
693 .or_else(|| self.config.provider.clone())
694 .unwrap_or_else(|| "unavailable".to_string());
695 json!({
696 "provider": provider,
697 "model": self.model.as_str()
698 })
699 }
700
701 fn select_model(&mut self, params: Value) -> std::result::Result<Value, AcpError> {
702 let model = params
703 .get("model")
704 .and_then(Value::as_str)
705 .ok_or_else(|| AcpError::invalid_params("model is required"))?
706 .to_string();
707
708 if let Some(provider_value) = params.get("provider") {
709 let provider_name = provider_value
710 .as_str()
711 .ok_or_else(|| AcpError::invalid_params("provider must be a string"))?;
712 let identity = self
713 .config
714 .resolve_provider_selection_identity(provider_name)
715 .map_err(AcpError::invalid_params)?;
716 self.config
717 .scope_to_provider_identity(&identity)
718 .map_err(AcpError::invalid_params)?;
719 }
720
721 self.model = model;
722 Ok(self.current_model())
723 }
724 }
725 const ACP_FULL_ACCESS_HINT: &str = "Start the server with `codewhale --approval-policy full-access serve --acp`, or set approval_policy = \"full-access\" in config.toml. Full Access also turns off Codewhale's own sandbox unless sandbox_mode tightens it; Plan stays read-only.";
726
727 fn acp_mode(config: &Config) -> AppMode {
728 if config.sandbox_mode.as_deref() == Some("read-only") {
729 AppMode::Plan
730 } else {
731 AppMode::Agent
732 }
733 }
734
735 /// Approval posture for ACP turns, derived from server config instead of
736 /// hardcoded: `--yolo` resolves to Bypass so an unattended headless session
737 /// actually executes tools (#6337); otherwise the configured approval policy,
738 /// else the Suggest default. Plan mode still pins read-only downstream
739 /// regardless of posture.
740 fn acp_approval_mode(config: &Config) -> ApprovalMode {
741 if config.yolo.unwrap_or(false) {
742 ApprovalMode::Bypass
743 } else {
744 config
745 .approval_policy
746 .as_deref()
747 .and_then(ApprovalMode::from_config_value)
748 .unwrap_or_default()
749 }
750 }
751
752 /// ACP `kind` hint for a tool call, used by the client to pick an icon/label.
753 /// Falls back to `"other"` for tools without an obvious category.
754 ///
755 /// `File` is a single canonical tool covering read/list/search/write/edit/
756 /// patch (#4625), so its kind depends on the `action` argument rather than
757 /// the tool name alone.
758 fn tool_call_kind(call: &PendingToolCall) -> &'static str {
759 match call.name.as_str() {
760 "File" => match call.input.get("action").and_then(Value::as_str) {
761 Some("write" | "edit" | "patch") => "edit",
762 _ => "read",
763 },
764 "apply_patch" => "edit",
765 "Git" => "read",
766 "bash" | "Bash" | "terminal/run" | "terminal/send" | "terminal/wait"
767 | "terminal/cancel" | "terminal/reset" => "execute",
768 _ => "other",
769 }
770 }
771
772 /// Human-readable title for a tool call: the tool name plus its primary
773 /// argument (path/command/pattern) when present, so the client's tool-call
774 /// card is legible without expanding raw input.
775 fn tool_call_title(call: &PendingToolCall) -> String {
776 let detail = call
777 .input
778 .get("path")
779 .or_else(|| call.input.get("command"))
780 .or_else(|| call.input.get("pattern"))
781 .or_else(|| call.input.get("task_id"))
782 .and_then(Value::as_str);
783 match detail {
784 Some(detail) => format!("{}: {}", call.name, detail),
785 None => call.name.clone(),
786 }
787 }
788
789 fn truncate_for_acp(content: &str) -> String {
790 if content.chars().count() <= TOOL_CALL_CONTENT_PREVIEW_CHARS {
791 return content.to_string();
792 }
793 let truncated: String = content
794 .chars()
795 .take(TOOL_CALL_CONTENT_PREVIEW_CHARS)
796 .collect();
797 format!("{truncated}\n… [truncated for display; the full result was sent to the model]")
798 }
799
800 async fn write_tool_call_start<W>(
801 writer: &mut W,
802 session_id: &str,
803 call: &PendingToolCall,
804 ) -> Result<()>
805 where
806 W: AsyncWrite + Unpin,
807 {
808 let notification = json!({
809 "jsonrpc": "2.0",
810 "method": "session/update",
811 "params": {
812 "sessionId": session_id,
813 "update": {
814 "sessionUpdate": "tool_call",
815 "toolCallId": call.execution_id,
816 "title": tool_call_title(call),
817 "kind": tool_call_kind(call),
818 "status": "pending",
819 "rawInput": call.input,
820 }
821 }
822 });
823 write_json_line(writer, notification).await
824 }
825
826 async fn write_tool_call_update<W>(
827 writer: &mut W,
828 session_id: &str,
829 call: &PendingToolCall,
830 status: &str,
831 content: Option<&str>,
832 ) -> Result<()>
833 where
834 W: AsyncWrite + Unpin,
835 {
836 write_tool_call_update_with_blocks(writer, session_id, call, status, content, &[]).await
837 }
838
839 async fn write_tool_call_update_with_blocks<W>(
840 writer: &mut W,
841 session_id: &str,
842 call: &PendingToolCall,
843 status: &str,
844 content: Option<&str>,
845 rich_blocks: &[codewhale_tools::ToolResultContentBlock],
846 ) -> Result<()>
847 where
848 W: AsyncWrite + Unpin,
849 {
850 let mut update = json!({
851 "sessionUpdate": "tool_call_update",
852 "toolCallId": call.execution_id,
853 "status": status,
854 });
855 if content.is_some() || !rich_blocks.is_empty() {
856 let mut blocks = Vec::with_capacity(rich_blocks.len() + usize::from(content.is_some()));
857 if let Some(content) = content {
858 blocks.push(json!({
859 "type": "content",
860 "content": { "type": "text", "text": truncate_for_acp(content) }
861 }));
862 }
863 blocks.extend(rich_blocks.iter().map(|block| match block {
864 codewhale_tools::ToolResultContentBlock::Image { mime_type, data } => json!({
865 "type": "content",
866 "content": { "type": "image", "data": data, "mimeType": mime_type }
867 }),
868 }));
869 update["content"] = json!(blocks);
870 }
871 let notification = json!({
872 "jsonrpc": "2.0",
873 "method": "session/update",
874 "params": {
875 "sessionId": session_id,
876 "update": update
877 }
878 });
879 write_json_line(writer, notification).await
880 }
881
882 impl AcpError {
883 fn invalid_params(message: impl Into<String>) -> Self {
884 Self {
885 code: -32602,
886 message: message.into(),
887 }
888 }
889
890 /// JSON-RPC internal error: the request was well-formed and the agent
891 /// could not serve it.
892 fn internal(message: impl Into<String>) -> Self {
893 Self {
894 code: -32603,
895 message: message.into(),
896 }
897 }
898
899 fn method_not_found(method: &str) -> Self {
900 Self {
901 code: -32601,
902 message: format!("method not found: {method}"),
903 }
904 }
905 }
906
907 fn initialize_result(client_protocol_version: Option<u64>, config: &Config) -> Value {
908 json!({
909 "protocolVersion": client_protocol_version
910 .map(|version| version.min(ACP_PROTOCOL_VERSION))
911 .unwrap_or(ACP_PROTOCOL_VERSION),
912 "agentCapabilities": {
913 "loadSession": true,
914 "modelSelection": true,
915 "promptCapabilities": {
916 "image": false,
917 "audio": false,
918 "embeddedContext": true
919 },
920 "mcpCapabilities": {
921 "http": false,
922 "sse": false
923 },
924 // ACP `SessionCapabilities` fields are objects, never booleans:
925 // `{}` means "supported", absent/null means "not supported"
926 // (#5969 — a boolean here made JetBrains' strictly-typed client
927 // fail the handshake and kill the agent). `session/load` support
928 // is advertised by the top-level `loadSession` above; it is not a
929 // field of `sessionCapabilities`. We only claim `list` because
930 // `session/list` is the only one of the optional session methods
931 // this server dispatches.
932 "sessionCapabilities": {
933 "list": {}
934 }
935 },
936 "agentInfo": {
937 "name": "codewhale",
938 "title": "codewhale",
939 "version": env!("CARGO_PKG_VERSION")
940 },
941 "authMethods": acp_auth_methods(config)
942 })
943 }
944
945 fn acp_auth_methods(config: &Config) -> Value {
946 let Some(identity) = config.active_provider_identity().ok().filter(|identity| {
947 identity.provider != ProviderKind::Custom && identity.provider != ProviderKind::Antigravity
948 }) else {
949 return json!([]);
950 };
951 let provider = identity.provider.as_str();
952 json!([
953 {
954 "id": "codewhale-terminal-auth",
955 "name": "Set Codewhale API key",
956 "description": format!("Run Codewhale's terminal credential setup for the {provider} provider."),
957 "type": "terminal",
958 "args": ["auth", "set", "--provider", provider],
959 "env": {}
960 }
961 ])
962 }
963
964 fn extract_prompt_text(prompt: Option<&Value>) -> Option<String> {
965 match prompt? {
966 Value::String(text) => Some(text.clone()),
967 Value::Array(blocks) => {
968 let parts = blocks
969 .iter()
970 .filter_map(content_block_text)
971 .collect::<Vec<_>>();
972 (!parts.is_empty()).then(|| parts.join("\n\n"))
973 }
974 _ => None,
975 }
976 }
977
978 fn content_block_text(block: &Value) -> Option<String> {
979 match block.get("type").and_then(Value::as_str)? {
980 "text" => block
981 .get("text")
982 .and_then(Value::as_str)
983 .map(str::to_string),
984 "resource" => resource_text(block),
985 "resource_link" | "resourceLink" => resource_link_text(block),
986 _ => None,
987 }
988 }
989
990 fn resource_text(block: &Value) -> Option<String> {
991 let resource = block.get("resource").unwrap_or(block);
992 if let Some(text) = resource.get("text").and_then(Value::as_str) {
993 return Some(text.to_string());
994 }
995 resource_link_text(resource)
996 }
997
998 fn resource_link_text(block: &Value) -> Option<String> {
999 let uri = block
1000 .get("uri")
1001 .or_else(|| block.pointer("/resource/uri"))
1002 .and_then(Value::as_str)?;
1003 Some(format!("@{uri}"))
1004 }
1005
1006 async fn write_session_update<W>(writer: &mut W, session_id: &str, text: String) -> Result<()>
1007 where
1008 W: AsyncWrite + Unpin,
1009 {
1010 let notification = json!({
1011 "jsonrpc": "2.0",
1012 "method": "session/update",
1013 "params": {
1014 "sessionId": session_id,
1015 "update": {
1016 "sessionUpdate": "agent_message_chunk",
1017 "content": {
1018 "type": "text",
1019 "text": text
1020 }
1021 }
1022 }
1023 });
1024 write_json_line(writer, notification).await
1025 }
1026
1027 async fn write_jsonrpc_result<W>(writer: &mut W, id: Value, result: Value) -> Result<()>
1028 where
1029 W: AsyncWrite + Unpin,
1030 {
1031 write_json_line(
1032 writer,
1033 json!({
1034 "jsonrpc": "2.0",
1035 "id": id,
1036 "result": result
1037 }),
1038 )
1039 .await
1040 }
1041
1042 async fn write_jsonrpc_error<W>(
1043 writer: &mut W,
1044 id: Option<Value>,
1045 code: i32,
1046 message: impl Into<String>,
1047 ) -> Result<()>
1048 where
1049 W: AsyncWrite + Unpin,
1050 {
1051 write_json_line(
1052 writer,
1053 json!({
1054 "jsonrpc": "2.0",
1055 "id": id,
1056 "error": {
1057 "code": code,
1058 "message": message.into()
1059 }
1060 }),
1061 )
1062 .await
1063 }
1064
1065 async fn write_json_line<W>(writer: &mut W, value: Value) -> Result<()>
1066 where
1067 W: AsyncWrite + Unpin,
1068 {
1069 tokio::time::timeout(std::time::Duration::from_secs(30), async {
1070 writer.write_all(value.to_string().as_bytes()).await?;
1071 writer.write_all(b"\n").await?;
1072 writer.flush().await
1073 })
1074 .await
1075 .map_err(|_| anyhow!("ACP transport write did not settle within 30 seconds"))??;
1076 Ok(())
1077 }
1078
1079 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
1080 enum JsonRpcResponseIdPolicy {
1081 /// JSON-RPC's normal contract: echo the request id without changing type.
1082 Preserve,
1083 /// Zed's ACP client currently decodes response ids as strings even when it
1084 /// sent a number. Keep this narrow compatibility mode client-identified.
1085 StringifyNumeric,
1086 }
1087
1088 impl JsonRpcResponseIdPolicy {
1089 fn from_initialize_params(params: &Value) -> Self {
1090 let client_name = params
1091 .pointer("/clientInfo/name")
1092 .and_then(Value::as_str)
1093 .unwrap_or_default();
1094 if client_name.eq_ignore_ascii_case("zed") {
1095 Self::StringifyNumeric
1096 } else {
1097 Self::Preserve
1098 }
1099 }
1100
1101 fn response_id(self, id: Value) -> Value {
1102 match (self, id) {
1103 (Self::StringifyNumeric, Value::Number(number)) => Value::String(number.to_string()),
1104 (_, id) => id,
1105 }
1106 }
1107 }
1108
1109 fn is_jsonrpc_response(message: &Value) -> bool {
1110 message.get("id").is_some()
1111 && (message.get("result").is_some() || message.get("error").is_some())
1112 }
1113
1114 #[cfg(test)]
1115 mod tests {
1116 use super::*;
1117 use crate::llm_client::mock::{MockLlmClient, canned};
1118 use codewhale_models::{ContentBlock, Role};
1119 use std::time::Duration;
1120
1121 fn parse_lines(output: Vec<u8>) -> Vec<Value> {
1122 String::from_utf8(output)
1123 .unwrap()
1124 .lines()
1125 .map(|line| serde_json::from_str(line).unwrap())
1126 .collect()
1127 }
1128 pub(super) fn fixture_config() -> Config {
1129 let mut config = Config {
1130 provider: Some("acp-fixture".into()),
1131 providers: Some(crate::config::ProvidersConfig {
1132 custom: HashMap::from([(
1133 "acp-fixture".into(),
1134 crate::config::ProviderConfig {
1135 kind: Some("openai-compatible".into()),
1136 base_url: Some("http://127.0.0.1:18181/v1".into()),
1137 model: Some("fixture-model".into()),
1138 api_key: Some("local-test-key".into()),
1139 ..Default::default()
1140 },
1141 )]),
1142 ..Default::default()
1143 }),
1144 snapshots: Some(crate::config::SnapshotsConfig {
1145 enabled: false,
1146 ..Default::default()
1147 }),
1148 ..Default::default()
1149 };
1150 config.set_feature("mcp", false).unwrap();
1151 config.set_feature("subagents", false).unwrap();
1152 config
1153 }
1154 pub(super) struct Rig {
1155 pub(super) server: AcpServer,
1156 mock: Arc<MockLlmClient>,
1157 pub(super) workspace: PathBuf,
1158 _home: crate::test_support::SealedHome,
1159 _dir: tempfile::TempDir,
1160 }
1161 impl Rig {
1162 pub(super) fn new(
1163 config: Config,
1164 turns: Vec<crate::llm_client::mock::CannedTurn>,
1165 ) -> Result<Self> {
1166 Self::with_base(config, turns, true)
1167 }
1168 fn with_base(
1169 config: Config,
1170 turns: Vec<crate::llm_client::mock::CannedTurn>,
1171 acp: bool,
1172 ) -> Result<Self> {
1173 let dir = tempfile::tempdir()?;
1174 let home = crate::test_support::SealedHome::at(dir.path());
1175 let workspace = dir.path().join("workspace");
1176 std::fs::create_dir(&workspace)?;
1177 let runtime_dir = dir.path().join("runtime");
1178 let manager_config = RuntimeThreadManagerConfig {
1179 data_dir: runtime_dir.clone(),
1180 task_data_dir: runtime_dir,
1181 sessions_dir: None,
1182 max_active_threads: 4,
1183 };
1184 let plugins = Arc::new(crate::plugins::PluginRegistry::empty(&workspace));
1185 let manager = Arc::new(if acp {
1186 RuntimeThreadManager::open_acp(
1187 config.clone(),
1188 workspace.clone(),
1189 manager_config,
1190 plugins,
1191 )?
1192 } else {
1193 RuntimeThreadManager::open_with_plugin_registry(
1194 config.clone(),
1195 workspace.clone(),
1196 manager_config,
1197 plugins,
1198 )?
1199 });
1200 let mock = Arc::new(MockLlmClient::new(turns));
1201 manager.set_test_model_client(mock.clone());
1202 let sessions = crate::session_manager::default_sessions_dir()?;
1203 Ok(Self {
1204 server: AcpServer::new(
1205 config,
1206 "fixture-model".into(),
1207 workspace.clone(),
1208 manager,
1209 sessions,
1210 None,
1211 None,
1212 ),
1213 mock,
1214 workspace,
1215 _home: home,
1216 _dir: dir,
1217 })
1218 }
1219 pub(super) async fn new_session(&mut self) -> String {
1220 self.server.new_session(json!({})).await.unwrap()["sessionId"]
1221 .as_str()
1222 .unwrap()
1223 .to_string()
1224 }
1225 async fn prompt(
1226 &mut self,
1227 session: &str,
1228 text: &str,
1229 ) -> Result<(&'static str, Vec<Value>)> {
1230 // Keep the transport input open: an empty reader is a real EOF/cancel.
1231 let (client, input) = tokio::io::duplex(4096);
1232 let mut reader = codewhale_app_server::BoundedLines::new(BufReader::new(input));
1233 let mut output = Vec::new();
1234 let result = tokio::time::timeout(
1235 Duration::from_secs(20),
1236 self.server.drive_prompt(
1237 json!({"sessionId":session,"prompt":text}),
1238 &mut reader,
1239 &mut output,
1240 ),
1241 )
1242 .await??;
1243 drop(client);
1244 Ok((result, parse_lines(output)))
1245 }
1246 async fn history(&self, session: &str) -> Vec<codewhale_models::Message> {
1247 let thread = &self.server.sessions[session].thread_id;
1248 self.server
1249 .runtime
1250 .get_engine(thread)
1251 .await
1252 .unwrap()
1253 .get_session_snapshot()
1254 .await
1255 .unwrap()
1256 .messages
1257 }
1258 pub(super) async fn close(&self) {
1259 self.server.runtime.shutdown_and_wait().await.unwrap();
1260 }
1261 }
1262 #[cfg(any(unix, windows))]
1263 #[tokio::test(flavor = "current_thread")]
1264 async fn authenticated_acp_owner_projects_actual_engine_and_guest_eof_preserves_lease()
1265 -> Result<()> {
1266 authenticated_owner_profile_case(true).await
1267 }
1268 #[cfg(any(unix, windows))]
1269 #[tokio::test(flavor = "current_thread")]
1270 async fn normal_owner_admits_authenticated_acp_then_ordinary_successor_without_profile_leak()
1271 -> Result<()> {
1272 authenticated_owner_profile_case(false).await
1273 }
1274 #[cfg(any(unix, windows))]
1275 async fn authenticated_owner_profile_case(base_acp: bool) -> Result<()> {
1276 struct AbortOnDrop(tokio::task::AbortHandle);
1277 impl Drop for AbortOnDrop {
1278 fn drop(&mut self) {
1279 self.0.abort();
1280 }
1281 }
1282 async fn correlated(
1283 guest: &mut codewhale_app_server::daemon_client::OwnerClient,
1284 id: Value,
1285 ) -> Result<Value> {
1286 tokio::time::timeout(Duration::from_secs(30), async {
1287 loop {
1288 let frame = guest
1289 .recv()
1290 .await?
1291 .ok_or_else(|| anyhow!("ACP owner closed before reply"))?;
1292 if frame["id"] == id {
1293 return Ok(frame);
1294 }
1295 }
1296 })
1297 .await?
1298 }
1299 let config = fixture_config();
1300 let rig = Rig::with_base(
1301 config.clone(),
1302 vec![canned::simple_text_turn("same owner answer")],
1303 base_acp,
1304 )?;
1305 let manager = rig.server.runtime.clone();
1306 let capture = manager.clone();
1307 let (binding, generation) = codewhale_app_server::daemon_socket::owner_work(move || {
1308 capture.capture_control_owner()
1309 })
1310 .await?;
1311 #[cfg(unix)]
1312 let socket_directory = rig._dir.path().canonicalize()?.join("endpoint");
1313 #[cfg(unix)]
1314 let _socket_directory =
1315 codewhale_config::private_directory::PrivateDirectory::open(&socket_directory)?;
1316 #[cfg(unix)]
1317 let socket_path = socket_directory.join("acp-owner.sock");
1318 #[cfg(windows)]
1319 let socket_path = PathBuf::from(format!(
1320 r"\\.\pipe\codewhale-acp-test-{}",
1321 uuid::Uuid::new_v4()
1322 ));
1323 #[cfg(unix)]
1324 let principal =
1325 codewhale_config::private_directory::PrivateDirectory::current_user_id().to_string();
1326 #[cfg(windows)]
1327 let principal = codewhale_app_server::daemon_socket::owner_work(|| {
1328 codewhale_config::windows_identity::CurrentWindowsUser::open()?.sid_string()
1329 })
1330 .await?;
1331 let receipt = codewhale_protocol::RuntimeOwnerReceipt {
1332 version: 1,
1333 data_dir: binding.data_dir,
1334 execution_scope: binding.execution_scope,
1335 lease_generation: generation,
1336 pid: std::process::id(),
1337 process_start: codewhale_app_server::daemon_socket::capture_process_start(
1338 std::process::id(),
1339 )
1340 .await?,
1341 principal,
1342 socket_path: socket_path.clone(),
1343 config_path: None,
1344 };
1345 let frontend = capture_frontend(
1346 config,
1347 "fixture-model".into(),
1348 rig.workspace.clone(),
1349 manager.clone(),
1350 rig.server.sessions_dir.clone(),
1351 None,
1352 None,
1353 )?;
1354 // ACP calls the held manager directly through its captured projection;
1355 // this closed endpoint has no fake HTTP service and is never used.
1356 let (daemon, _state) = codewhale_app_server::bind_runtime_frontends(
1357 None,
1358 None,
1359 receipt.clone(),
1360 codewhale_app_server::RuntimeOwnerRouting {
1361 workers: None,
1362 workspace: Some(rig.workspace.clone()),
1363 endpoint: "127.0.0.1:1".parse()?,
1364 mobile: false,
1365 web: false,
1366 acp: true,
1367 acp_only: base_acp,
1368 },
1369 Some(frontend),
1370 )
1371 .await?;
1372 let shutdown = daemon.shutdown_handle();
1373 let owner = tokio::spawn(daemon.serve());
1374 let _abort = AbortOnDrop(owner.abort_handle());
1375 let mut guest = codewhale_app_server::daemon_client::connect_acp_if_published(
1376 None,
1377 Some(socket_path.clone()),
1378 )
1379 .await?
1380 .ok_or_else(|| anyhow!("published ACP owner unavailable"))?;
1381 assert_eq!(guest.receipt(), &receipt);
1382 assert!(guest.routing().is_some_and(|routing| routing.acp));
1383 guest
1384 .send(
1385 json!(1),
1386 "initialize",
1387 json!({"protocolVersion":1,"clientCapabilities":{}}),
1388 )
1389 .await?;
1390 assert!(correlated(&mut guest, json!(1)).await?["error"].is_null());
1391 guest
1392 .send(json!(2), "session/new", json!({"cwd":rig.workspace}))
1393 .await?;
1394 let created = correlated(&mut guest, json!(2)).await?;
1395 let session = created["result"]["sessionId"]
1396 .as_str()
1397 .ok_or_else(|| anyhow!("{created}"))?
1398 .to_string();
1399 guest
1400 .send(
1401 json!(3),
1402 "session/prompt",
1403 json!({"sessionId":session,"prompt":"one actual turn"}),
1404 )
1405 .await?;
1406 let completed = correlated(&mut guest, json!(3)).await?;
1407 assert_eq!(completed["result"]["stopReason"], "end_turn", "{completed}");
1408 tokio::time::timeout(
1409 Duration::from_secs(5),
1410 guest.forward(tokio::io::empty(), tokio::io::sink()),
1411 )
1412 .await??;
1413 let mut shutdown_guest = codewhale_app_server::daemon_client::connect_acp_if_published(
1414 None,
1415 Some(socket_path.clone()),
1416 )
1417 .await?
1418 .ok_or_else(|| anyhow!("owner ended after guest EOF"))?;
1419 shutdown_guest.send(json!(4), "shutdown", json!({})).await?;
1420 let closed = correlated(&mut shutdown_guest, json!(4)).await?;
1421 assert!(closed["error"].is_null());
1422 assert!(
1423 tokio::time::timeout(Duration::from_secs(5), shutdown_guest.recv())
1424 .await??
1425 .is_none(),
1426 "ACP shutdown closes only this connection"
1427 );
1428 let rows = manager
1429 .list_threads(
1430 crate::runtime_threads::ThreadListFilter::IncludeArchived,
1431 None,
1432 )
1433 .await?;
1434 assert_eq!(rows.len(), 1);
1435 let detail = manager.get_thread_detail(&rows[0].id).await?;
1436 assert_eq!(detail.turns.len(), 1);
1437 assert_eq!(rig.mock.call_count(), 1);
1438 if !base_acp {
1439 rig.mock
1440 .push_turn(canned::simple_text_turn("ordinary successor"));
1441 let ordinary = manager
1442 .start_turn(
1443 &rows[0].id,
1444 crate::runtime_threads::StartTurnRequest {
1445 prompt: "ordinary successor".into(),
1446 ..Default::default()
1447 },
1448 )
1449 .await?;
1450 tokio::time::timeout(Duration::from_secs(30), async {
1451 loop {
1452 let detail = manager.get_thread_detail(&rows[0].id).await?;
1453 if detail.turns.iter().any(|t| {
1454 t.id == ordinary.id
1455 && t.status == crate::runtime_threads::RuntimeTurnStatus::Completed
1456 }) {
1457 return Ok::<_, anyhow::Error>(());
1458 }
1459 tokio::time::sleep(Duration::from_millis(20)).await;
1460 }
1461 })
1462 .await??;
1463 let requests = rig.mock.captured_requests();
1464 assert_eq!(requests.len(), 2);
1465 let acp = requests[0]
1466 .tools
1467 .as_ref()
1468 .ok_or_else(|| anyhow!("ACP tools missing"))?;
1469 assert!(acp.iter().all(|tool| !matches!(
1470 tool.name.as_str(),
1471 "update_goal"
1472 | "create_goal"
1473 | "request_user_input"
1474 | "tool_search"
1475 | "spawn_agent"
1476 )));
1477 let normal = requests[1]
1478 .tools
1479 .as_ref()
1480 .ok_or_else(|| anyhow!("ordinary tools missing"))?;
1481 assert_ne!(
1482 acp.iter().map(|t| &t.name).collect::<Vec<_>>(),
1483 normal.iter().map(|t| &t.name).collect::<Vec<_>>(),
1484 "ordinary successor must rebuild its ordinary catalog"
1485 );
1486 assert_eq!(manager.get_thread_detail(&rows[0].id).await?.turns.len(), 2);
1487 }
1488 let capture = manager.clone();
1489 let (after, generation) = codewhale_app_server::daemon_socket::owner_work(move || {
1490 capture.capture_control_owner()
1491 })
1492 .await?;
1493 assert_eq!(after.data_dir, receipt.data_dir);
1494 assert_eq!(
1495 generation, receipt.lease_generation,
1496 "guest EOF cannot release the host lease"
1497 );
1498 shutdown.trigger();
1499 tokio::time::timeout(Duration::from_secs(5), owner).await???;
1500 assert!(
1501 codewhale_app_server::daemon_client::connect_acp_if_published(None, Some(socket_path))
1502 .await?
1503 .is_none(),
1504 "shutdown withdraws the authenticated owner receipt instead of advertising a guest"
1505 );
1506 rig.close().await;
1507 Ok(())
1508 }
1509
1510 #[test]
1511 fn acp_approval_mode_derives_from_server_config() {
1512 // #6337: `--yolo --danger-full-access` must not silently run as Ask.
1513 let yolo = Config {
1514 yolo: Some(true),
1515 ..Config::default()
1516 };
1517 assert_eq!(acp_approval_mode(&yolo), ApprovalMode::Bypass);
1518 let policy = Config {
1519 approval_policy: Some("never".into()),
1520 ..Config::default()
1521 };
1522 assert_eq!(acp_approval_mode(&policy), ApprovalMode::Never);
1523 assert_eq!(acp_approval_mode(&Config::default()), ApprovalMode::Suggest);
1524 }
1525
1526 #[test]
1527 fn initialize_advertises_baseline_acp_agent() {
1528 let result = initialize_result(Some(1), &Config::default());
1529
1530 assert_eq!(result["protocolVersion"], 1);
1531 assert_eq!(result["agentInfo"]["name"], "codewhale");
1532 // #5864: enumerating and resuming durable Codewhale sessions is now
1533 // served, so the capability says so rather than declining it.
1534 // #5969: this test used to assert `list == true` and `load == true`,
1535 // certifying the wire format that broke every strictly-typed client.
1536 // `loadSession` is the top-level boolean that advertises
1537 // `session/load`; inside `sessionCapabilities` every field is a
1538 // capability *object*, and `load` is not a field at all. Comparing the
1539 // whole object pins both halves: a boolean `list` or a resurrected
1540 // `load` key fails here.
1541 assert_eq!(result["agentCapabilities"]["loadSession"], true);
1542 assert_eq!(
1543 result["agentCapabilities"]["sessionCapabilities"],
1544 json!({"list": {}})
1545 );
1546 assert!(
1547 result["agentCapabilities"]["sessionCapabilities"]["list"].is_object(),
1548 "sessionCapabilities.list must be a SessionListCapabilities object, got {}",
1549 result["agentCapabilities"]["sessionCapabilities"]["list"]
1550 );
1551 assert_eq!(
1552 result["agentCapabilities"]["promptCapabilities"]["embeddedContext"],
1553 true
1554 );
1555 assert_eq!(result["authMethods"][0]["type"], "terminal");
1556 assert_eq!(
1557 result["authMethods"][0]["args"],
1558 json!(["auth", "set", "--provider", "deepseek"])
1559 );
1560 }
1561
1562 #[test]
1563 fn initialize_advertises_model_selection_capability() {
1564 let result = initialize_result(Some(1), &Config::default());
1565
1566 assert_eq!(result["agentCapabilities"]["modelSelection"], true);
1567 }
1568
1569 #[test]
1570 fn extract_prompt_text_accepts_text_and_resource_blocks() {
1571 let prompt = json!([
1572 { "type": "text", "text": "Review this file" },
1573 {
1574 "type": "resource",
1575 "resource": {
1576 "uri": "file:///tmp/app.rs",
1577 "mimeType": "text/rust",
1578 "text": "fn main() {}"
1579 }
1580 },
1581 { "type": "resource_link", "uri": "file:///tmp/lib.rs" }
1582 ]);
1583
1584 let text = extract_prompt_text(Some(&prompt)).expect("prompt text");
1585
1586 assert!(text.contains("Review this file"));
1587 assert!(text.contains("fn main() {}"));
1588 assert!(text.contains("@file:///tmp/lib.rs"));
1589 }
1590
1591 #[tokio::test]
1592 async fn session_update_is_protocol_clean_single_line_json() {
1593 let mut out = Vec::new();
1594
1595 write_session_update(&mut out, "sess_1", "hello\nworld".to_string())
1596 .await
1597 .expect("write update");
1598
1599 let line = String::from_utf8(out).expect("utf8");
1600 assert_eq!(line.lines().count(), 1);
1601 let value: Value = serde_json::from_str(line.trim()).expect("json");
1602 assert_eq!(value["method"], "session/update");
1603 assert_eq!(value["params"]["sessionId"], "sess_1");
1604 assert_eq!(value["params"]["update"]["content"]["text"], "hello\nworld");
1605 }
1606
1607 #[tokio::test]
1608 async fn jsonrpc_result_preserves_numeric_ids_for_avante_acp() {
1609 let mut out = Vec::new();
1610
1611 let params = json!({
1612 "protocolVersion": 1,
1613 "clientCapabilities": {}
1614 });
1615 let id = JsonRpcResponseIdPolicy::from_initialize_params(&params).response_id(json!(1));
1616 write_jsonrpc_result(&mut out, id, json!({"ok": true}))
1617 .await
1618 .expect("write result");
1619
1620 let line = String::from_utf8(out).expect("utf8");
1621 let value: Value = serde_json::from_str(line.trim()).expect("json");
1622 // Numeric ID must stay numeric — avante.nvim's Lua client uses
1623 // strict table keys (callbacks[1] ≠ callbacks["1"]).
1624 assert!(
1625 value["id"].is_number(),
1626 "numeric id must stay numeric, got {:?}",
1627 value["id"]
1628 );
1629 assert_eq!(value["result"], json!({"ok": true}));
1630 }
1631
1632 #[tokio::test]
1633 async fn jsonrpc_result_stringifies_numeric_ids_for_zed_acp() {
1634 let mut out = Vec::new();
1635
1636 let params = json!({
1637 "protocolVersion": 1,
1638 "clientCapabilities": {},
1639 "clientInfo": {
1640 "name": "zed",
1641 "version": "1.2.6"
1642 }
1643 });
1644 let id = JsonRpcResponseIdPolicy::from_initialize_params(&params).response_id(json!(1));
1645 write_jsonrpc_result(&mut out, id, json!({"ok": true}))
1646 .await
1647 .expect("write result");
1648
1649 let line = String::from_utf8(out).expect("utf8");
1650 let value: Value = serde_json::from_str(line.trim()).expect("json");
1651 assert_eq!(value["id"], "1");
1652 assert_eq!(value["result"], json!({"ok": true}));
1653 }
1654
1655 #[tokio::test]
1656 async fn jsonrpc_error_keeps_absent_id_null() {
1657 let mut out = Vec::new();
1658
1659 write_jsonrpc_error(&mut out, None, -32700, "invalid json")
1660 .await
1661 .expect("write error");
1662
1663 let line = String::from_utf8(out).expect("utf8");
1664 let value: Value = serde_json::from_str(line.trim()).expect("json");
1665 assert_eq!(value["id"], Value::Null);
1666 assert_eq!(value["error"]["code"], -32700);
1667 }
1668
1669 #[tokio::test]
1670 async fn tool_update_emits_typed_acp_image_content() {
1671 let mut output = Vec::new();
1672 let call = PendingToolCall {
1673 execution_id: uuid::Uuid::new_v4().to_string(),
1674 name: "read".to_string(),
1675 input: json!({"path": "shot.png"}),
1676 };
1677 write_tool_call_update_with_blocks(
1678 &mut output,
1679 "session_1",
1680 &call,
1681 "completed",
1682 Some("screenshot captured"),
1683 &[codewhale_tools::ToolResultContentBlock::Image {
1684 mime_type: "image/png".to_string(),
1685 data: "QUJD".to_string(),
1686 }],
1687 )
1688 .await
1689 .expect("ACP update");
1690
1691 let lines = parse_lines(output);
1692 let content = lines[0]["params"]["update"]["content"]
1693 .as_array()
1694 .expect("ACP content blocks");
1695 assert_eq!(content[0]["content"]["type"], "text");
1696 assert_eq!(content[1]["content"]["type"], "image");
1697 assert_eq!(content[1]["content"]["mimeType"], "image/png");
1698 assert_eq!(content[1]["content"]["data"], "QUJD");
1699 }
1700
1701 #[tokio::test(flavor = "current_thread")]
1702 async fn new_session_returns_durable_bare_uuid_without_provider_or_history() -> Result<()> {
1703 let mut rig = Rig::new(fixture_config(), vec![])?;
1704 let id = rig.new_session().await;
1705 uuid::Uuid::parse_str(&id)?;
1706 let store = crate::session_manager::SessionManager::new(rig.server.sessions_dir.clone())?;
1707 let saved = store.load_session(&id)?;
1708 assert!(saved.messages.is_empty());
1709 assert_eq!(
1710 saved.metadata.runtime_store.as_ref(),
1711 Some(&rig.server.runtime.session_store_binding())
1712 );
1713 let thread = rig.server.sessions[&id].thread_id.clone();
1714 assert_eq!(
1715 rig.server
1716 .runtime
1717 .get_thread(&thread)
1718 .await?
1719 .session_id
1720 .as_deref(),
1721 Some(id.as_str())
1722 );
1723 assert_eq!(rig.mock.call_count(), 0);
1724 let loaded = rig
1725 .server
1726 .load_session(json!({"sessionId":&id[..8]}))
1727 .await
1728 .unwrap();
1729 assert_eq!(loaded["sessionId"], id);
1730 assert_eq!(rig.server.sessions.len(), 1);
1731 assert_eq!(rig.server.insertion_order.len(), 1);
1732 assert!(
1733 rig.server
1734 .load_session(json!({"sessionId":uuid::Uuid::new_v4().to_string()}))
1735 .await
1736 .is_err()
1737 );
1738 rig.close().await;
1739 Ok(())
1740 }
1741 #[tokio::test(flavor = "current_thread")]
1742 async fn canonical_sessions_keep_independent_history_and_scoped_configuration() -> Result<()> {
1743 let mut rig = Rig::new(
1744 fixture_config(),
1745 vec![
1746 canned::simple_text_turn("A answer"),
1747 canned::simple_text_turn("B answer"),
1748 ],
1749 )?;
1750 let a = rig.new_session().await;
1751 let b = rig.new_session().await;
1752 rig.server
1753 .set_session_config(json!({"sessionId":a,"configId":"mode","value":"plan"}))
1754 .await
1755 .unwrap();
1756 assert_eq!(
1757 rig.server.session_configuration(&a).await.unwrap()["modes"]["currentModeId"],
1758 "plan"
1759 );
1760 assert_eq!(
1761 rig.server.session_configuration(&b).await.unwrap()["modes"]["currentModeId"],
1762 "agent"
1763 );
1764 assert!(
1765 rig.server
1766 .set_session_config(
1767 json!({"sessionId":a,"configId":"permission","value":"full-access"})
1768 )
1769 .await
1770 .is_err()
1771 );
1772 assert!(
1773 rig.server
1774 .set_session_config(json!({"sessionId":a,"configId":"mode","value":"fabricated"}))
1775 .await
1776 .is_err()
1777 );
1778 assert_eq!(rig.prompt(&a, "A prompt").await?.0, "end_turn");
1779 assert_eq!(rig.prompt(&b, "B prompt").await?.0, "end_turn");
1780 let ah = serde_json::to_string(&rig.history(&a).await)?;
1781 let bh = serde_json::to_string(&rig.history(&b).await)?;
1782 assert!(ah.contains("A prompt") && ah.contains("A answer") && !ah.contains("B prompt"));
1783 assert!(bh.contains("B prompt") && bh.contains("B answer") && !bh.contains("A prompt"));
1784 let listed = rig
1785 .server
1786 .list_sessions(json!({"cwd":rig.workspace}))
1787 .unwrap();
1788 assert_eq!(listed["sessions"].as_array().unwrap().len(), 2);
1789 rig.close().await;
1790 Ok(())
1791 }
1792 #[tokio::test(flavor = "current_thread")]
1793 async fn agentic_turn_chains_real_write_read_and_retains_full_core_pairs() -> Result<()> {
1794 let mut config = fixture_config();
1795 config.yolo = Some(true);
1796 let mut rig = Rig::new(
1797 config,
1798 vec![
1799 canned::tool_call_turn(
1800 "provider-write",
1801 "write",
1802 r#"{"path":"receipt.txt","content":"full receipt"}"#,
1803 ),
1804 canned::tool_call_turn("provider-read", "read", r#"{"path":"receipt.txt"}"#),
1805 canned::simple_text_turn("done"),
1806 ],
1807 )?;
1808 let id = rig.new_session().await;
1809 let (reason, wire) = rig.prompt(&id, "write then read").await?;
1810 assert_eq!(reason, "end_turn");
1811 assert_eq!(
1812 std::fs::read_to_string(rig.workspace.join("receipt.txt"))?,
1813 "full receipt"
1814 );
1815 assert!(
1816 !wire
1817 .iter()
1818 .any(|v| v["method"] == "session/request_permission")
1819 );
1820 let statuses: Vec<_> = wire
1821 .iter()
1822 .filter_map(|v| v.pointer("/params/update/status").and_then(Value::as_str))
1823 .collect();
1824 assert_eq!(
1825 statuses,
1826 vec![
1827 "pending",
1828 "in_progress",
1829 "completed",
1830 "pending",
1831 "in_progress",
1832 "completed"
1833 ]
1834 );
1835 let history = rig.history(&id).await;
1836 let uses = history
1837 .iter()
1838 .flat_map(|m| &m.content)
1839 .filter(|b| matches!(b, ContentBlock::ToolUse { .. }))
1840 .count();
1841 let results = history
1842 .iter()
1843 .flat_map(|m| &m.content)
1844 .filter(|b| matches!(b, ContentBlock::ToolResult { .. }))
1845 .count();
1846 assert_eq!(uses, 2);
1847 assert_eq!(results, 2);
1848 let saved = crate::session_manager::SessionManager::new(rig.server.sessions_dir.clone())?
1849 .load_session(&id)?;
1850 assert_eq!(
1851 saved.messages, history,
1852 "checkpoint is the full Engine snapshot, not display chunks"
1853 );
1854 let request = rig.mock.captured_requests();
1855 assert_eq!(request.len(), 3);
1856 assert!(
1857 request.iter().all(
1858 |r| r.tools.as_ref().is_some_and(|t| t.iter().all(|t| !matches!(
1859 t.name.as_str(),
1860 "code_execution"
1861 | "js_execution"
1862 | "execute_tools"
1863 | "tool_search"
1864 | "task"
1865 | "subagent"
1866 )))
1867 )
1868 );
1869 rig.close().await;
1870 Ok(())
1871 }
1872 #[tokio::test(flavor = "current_thread")]
1873 async fn agentic_turn_preserves_core_partial_write_when_later_provider_fails() -> Result<()> {
1874 let mut config = fixture_config();
1875 config.yolo = Some(true);
1876 let mut rig = Rig::new(
1877 config,
1878 vec![canned::tool_call_turn(
1879 "partial",
1880 "write",
1881 r#"{"path":"partial.txt","content":"completed effect"}"#,
1882 )],
1883 )?;
1884 rig.mock.push_error("fixture terminal provider refusal");
1885 let id = rig.new_session().await;
1886 let error = rig.prompt(&id, "write then fail").await.unwrap_err();
1887 assert!(!error.to_string().is_empty());
1888 assert_eq!(
1889 std::fs::read_to_string(rig.workspace.join("partial.txt"))?,
1890 "completed effect"
1891 );
1892 let history = rig.history(&id).await;
1893 assert!(
1894 history
1895 .iter()
1896 .flat_map(|m| &m.content)
1897 .any(|b| matches!(b, ContentBlock::ToolResult { .. }))
1898 );
1899 let saved = crate::session_manager::SessionManager::new(rig.server.sessions_dir.clone())?
1900 .load_session(&id)?;
1901 assert_eq!(saved.messages, history);
1902 let detail = rig
1903 .server
1904 .runtime
1905 .get_thread_detail(&rig.server.sessions[&id].thread_id)
1906 .await?;
1907 assert_eq!(
1908 detail.turns.last().unwrap().status,
1909 crate::runtime_threads::RuntimeTurnStatus::Failed
1910 );
1911 rig.close().await;
1912 Ok(())
1913 }
1914 #[tokio::test(flavor = "current_thread")]
1915 async fn core_prompt_is_stable_after_instruction_write_and_streams_each_text_delta()
1916 -> Result<()> {
1917 let mut config = fixture_config();
1918 config.yolo = Some(true);
1919 let final_turn = vec![
1920 canned::message_start("final"),
1921 canned::text_block_start(0),
1922 canned::text_delta(0, "hello"),
1923 canned::text_delta(0, " world"),
1924 canned::block_stop(0),
1925 canned::message_delta("end_turn", None),
1926 canned::message_stop(),
1927 ];
1928 let mut rig = Rig::new(
1929 config,
1930 vec![
1931 canned::tool_call_turn(
1932 "law-write",
1933 "write",
1934 r#"{"path":"AGENTS.md","content":"new-self-authored-law"}"#,
1935 ),
1936 final_turn,
1937 ],
1938 )?;
1939 std::fs::write(rig.workspace.join("AGENTS.md"), "initial-project-law")?;
1940 let id = rig.new_session().await;
1941 let (reason, wire) = rig.prompt(&id, "update instructions then answer").await?;
1942 assert_eq!(reason, "end_turn");
1943 assert_eq!(
1944 std::fs::read_to_string(rig.workspace.join("AGENTS.md"))?,
1945 "new-self-authored-law"
1946 );
1947 let requests = rig.mock.captured_requests();
1948 assert_eq!(requests.len(), 2);
1949 assert_eq!(
1950 serde_json::to_vec(&requests[0].system)?,
1951 serde_json::to_vec(&requests[1].system)?
1952 );
1953 let prompt = crate::prompts::system_prompt_flat_text(requests[0].system.as_ref().unwrap());
1954 assert!(prompt.contains(crate::prompts::text::BASE_PROMPT.trim()));
1955 assert!(prompt.contains("initial-project-law"));
1956 assert!(!prompt.contains("new-self-authored-law"));
1957 let chunks: Vec<_> = wire
1958 .iter()
1959 .filter(|v| {
1960 v.pointer("/params/update/sessionUpdate") == Some(&json!("agent_message_chunk"))
1961 })
1962 .filter_map(|v| {
1963 v.pointer("/params/update/content/text")
1964 .and_then(Value::as_str)
1965 })
1966 .collect();
1967 // Runtime may coalesce provider frames. ACP must project its durable
1968 // deltas once each, rather than inventing a second streaming source.
1969 let events = rig
1970 .server
1971 .runtime
1972 .events_since(&rig.server.sessions[&id].thread_id, None)?;
1973 let deltas: Vec<_> = events
1974 .iter()
1975 .filter(|event| {
1976 event.event == "item.delta"
1977 && event.payload["kind"].as_str() == Some("agent_message")
1978 })
1979 .collect();
1980 let canonical_chunks: Vec<_> = deltas
1981 .iter()
1982 .map(|event| event.payload["delta"].as_str().unwrap())
1983 .collect();
1984 assert!(!chunks.is_empty());
1985 assert_eq!(chunks, canonical_chunks);
1986 assert_eq!(chunks.concat(), "hello world");
1987 let completed = events
1988 .iter()
1989 .find(|event| event.event == "turn.completed")
1990 .unwrap();
1991 assert!(deltas.iter().all(|event| event.seq < completed.seq));
1992 let history = serde_json::to_string(&rig.history(&id).await)?;
1993 assert!(history.contains("hello world"));
1994 rig.close().await;
1995 Ok(())
1996 }
1997 #[tokio::test(flavor = "current_thread")]
1998 async fn core_file_batch_orders_real_results_and_returns_missing_path_to_model() -> Result<()> {
1999 let mut events = vec![canned::message_start("file-batch")];
2000 for (index, id, name, input) in [
2001 (0, "list", "File", r#"{"action":"list","path":"."}"#),
2002 (1, "read", "read", r#"{"path":"present.txt"}"#),
2003 (2, "missing", "read", r#"{"path":"missing.txt"}"#),
2004 ] {
2005 events.extend([
2006 canned::tool_use_block_start(index, id, name),
2007 canned::tool_input_delta(index, input),
2008 canned::block_stop(index),
2009 ]);
2010 }
2011 events.extend([
2012 canned::message_delta("tool_use", None),
2013 canned::message_stop(),
2014 ]);
2015 let mut rig = Rig::new(
2016 fixture_config(),
2017 vec![events, canned::simple_text_turn("handled file failure")],
2018 )?;
2019 std::fs::write(rig.workspace.join("present.txt"), "real-file-marker")?;
2020 let id = rig.new_session().await;
2021 let (reason, wire) = rig
2022 .prompt(&id, "list then read present and missing")
2023 .await?;
2024 assert_eq!(reason, "end_turn");
2025 let statuses: Vec<_> = wire
2026 .iter()
2027 .filter_map(|v| v.pointer("/params/update/status").and_then(Value::as_str))
2028 .collect();
2029 assert_eq!(
2030 statuses,
2031 vec![
2032 "pending",
2033 "pending",
2034 "pending",
2035 "in_progress",
2036 "completed",
2037 "in_progress",
2038 "completed",
2039 "in_progress",
2040 "failed"
2041 ]
2042 );
2043 let requests = rig.mock.captured_requests();
2044 assert_eq!(requests.len(), 2);
2045 let feedback = serde_json::to_string(&requests[1].messages)?;
2046 assert!(
2047 feedback.contains("present.txt")
2048 && feedback.contains("real-file-marker")
2049 && feedback.contains("missing.txt")
2050 );
2051 assert!(
2052 requests[1]
2053 .messages
2054 .iter()
2055 .flat_map(|m| &m.content)
2056 .any(|b| matches!(
2057 b,
2058 ContentBlock::ToolResult {
2059 is_error: Some(true),
2060 ..
2061 }
2062 ))
2063 );
2064 rig.close().await;
2065 Ok(())
2066 }
2067 #[derive(Clone, Copy)]
2068 enum PermissionReply {
2069 AllowAfterWrong,
2070 ForeignCancelAndBusy,
2071 Reject,
2072 CancelLateAllow,
2073 Eof,
2074 }
2075 async fn permission_case(reply: PermissionReply) -> Result<()> {
2076 permission_case_with_base(reply, true).await
2077 }
2078 async fn permission_case_with_base(reply: PermissionReply, base_acp: bool) -> Result<()> {
2079 let mut rig = Rig::with_base(
2080 fixture_config(),
2081 vec![
2082 canned::tool_call_turn(
2083 "provider-id-not-authority",
2084 "write",
2085 r#"{"path":"ask.txt","content":"allowed"}"#,
2086 ),
2087 canned::simple_text_turn("settled"),
2088 ],
2089 base_acp,
2090 )?;
2091 let id = rig.new_session().await;
2092 let thread = rig.server.sessions[&id].thread_id.clone();
2093 let manager = Arc::clone(&rig.server.runtime);
2094 let target = rig.workspace.join("ask.txt");
2095 let (client, transport) = tokio::io::duplex(65536);
2096 let (read, mut write) = tokio::io::split(transport);
2097 let mut reader = codewhale_app_server::BoundedLines::new(BufReader::new(read));
2098 let (cr, mut cw) = tokio::io::split(client);
2099 let mut cr = codewhale_app_server::BoundedLines::new(BufReader::new(cr));
2100 let sid = id.clone();
2101 let target_in = target.clone();
2102 let client = async move {
2103 let mut seen = Vec::new();
2104 loop {
2105 let line = cr
2106 .next_line()
2107 .await?
2108 .ok_or_else(|| anyhow!("permission request missing"))?;
2109 let v: Value = serde_json::from_str(&line)?;
2110 seen.push(v.clone());
2111 if v["method"] != "session/request_permission" {
2112 continue;
2113 }
2114 let detail = manager.get_thread_detail(&thread).await?;
2115 let pending = detail
2116 .pending_approvals
2117 .first()
2118 .ok_or_else(|| anyhow!("actual pending approval missing"))?;
2119 assert_eq!(
2120 pending.tool_call_id.as_deref(),
2121 v.pointer("/params/toolCall/toolCallId")
2122 .and_then(Value::as_str)
2123 );
2124 assert!(!target_in.exists());
2125 assert!(
2126 !seen
2127 .iter()
2128 .any(|v| v.pointer("/params/update/status") == Some(&json!("in_progress"))),
2129 "no dispatch before the real decision"
2130 );
2131 match reply {
2132 PermissionReply::AllowAfterWrong=>{
2133 write_json_line(&mut cw,json!({"jsonrpc":"2.0","id":"wrong-id","result":{"outcome":{"outcome":"selected","optionId":"allow-once"}}})).await?;
2134 tokio::time::sleep(Duration::from_millis(30)).await;
2135 assert!(!target_in.exists());assert_eq!(manager.get_thread_detail(&thread).await?.pending_approvals.len(),1);
2136 write_json_line(&mut cw,json!({"jsonrpc":"2.0","id":v["id"],"result":{"outcome":{"outcome":"selected","optionId":"allow-once"}}})).await?;
2137 },
2138 PermissionReply::ForeignCancelAndBusy=>{
2139 write_json_line(&mut cw,json!({"jsonrpc":"2.0","id":41,"method":"session/cancel","params":{"sessionId":"foreign-session"}})).await?;
2140 write_json_line(&mut cw,json!({"jsonrpc":"2.0","id":42,"method":"session/prompt","params":{"sessionId":sid,"prompt":"must not run"}})).await?;
2141 for expected in [41,42] {
2142 let response:Value=serde_json::from_str(&cr.next_line().await?.ok_or_else(||anyhow!("control response missing"))?)?;
2143 assert_eq!(response["id"],expected);
2144 if expected==41 {assert_eq!(response["result"],Value::Null);} else {assert_eq!(response["error"]["code"],-32603);}
2145 seen.push(response);
2146 }
2147 assert!(!target_in.exists());
2148 assert_eq!(manager.get_thread_detail(&thread).await?.pending_approvals.len(),1);
2149 write_json_line(&mut cw,json!({"jsonrpc":"2.0","id":v["id"],"result":{"outcome":{"outcome":"selected","optionId":"allow-once"}}})).await?;
2150 },
2151 PermissionReply::Reject=>write_json_line(&mut cw,json!({"jsonrpc":"2.0","id":v["id"],"result":{"outcome":{"outcome":"selected","optionId":"fabricated"}}})).await?,
2152 PermissionReply::Eof=>cw.shutdown().await?,
2153 PermissionReply::CancelLateAllow=>{
2154 write_json_line(&mut cw,json!({"jsonrpc":"2.0","method":"session/cancel","params":{"sessionId":sid}})).await?;
2155 write_json_line(&mut cw,json!({"jsonrpc":"2.0","id":v["id"],"result":{"outcome":{"outcome":"selected","optionId":"allow-once"}}})).await?;
2156 },
2157 }
2158 return Ok::<_, anyhow::Error>((cr, cw, seen));
2159 }
2160 };
2161 let driver = rig.server.drive_prompt(
2162 json!({"sessionId":id,"prompt":"write under Ask"}),
2163 &mut reader,
2164 &mut write,
2165 );
2166 let (result, client) = tokio::time::timeout(Duration::from_secs(20), async {
2167 tokio::join!(driver, client)
2168 })
2169 .await?;
2170 let (mut cr, _cw, mut wire) = client?;
2171 write.shutdown().await?;
2172 while let Some(line) = cr.next_line().await? {
2173 wire.push(serde_json::from_str(&line)?);
2174 }
2175 match reply {
2176 PermissionReply::AllowAfterWrong | PermissionReply::ForeignCancelAndBusy => {
2177 assert_eq!(result?, "end_turn");
2178 assert_eq!(std::fs::read_to_string(&target)?, "allowed");
2179 assert_eq!(
2180 wire.iter()
2181 .filter(
2182 |v| v.pointer("/params/update/status") == Some(&json!("in_progress"))
2183 )
2184 .count(),
2185 1
2186 );
2187 }
2188 PermissionReply::Reject => {
2189 assert_eq!(result?, "end_turn");
2190 assert!(!target.exists());
2191 }
2192 PermissionReply::CancelLateAllow | PermissionReply::Eof => {
2193 assert_eq!(result?, "cancelled");
2194 assert!(!target.exists());
2195 }
2196 }
2197 let saved = crate::session_manager::SessionManager::new(rig.server.sessions_dir.clone())?
2198 .load_session(&id)?;
2199 assert!(saved.messages.iter().any(|m| m.role == Role::User));
2200 if !base_acp {
2201 let thread = rig.server.sessions[&id].thread_id.clone();
2202 let successor = rig
2203 .server
2204 .runtime
2205 .start_turn(
2206 &thread,
2207 crate::runtime_threads::StartTurnRequest {
2208 prompt: "ordinary after cancelled ACP".into(),
2209 ..Default::default()
2210 },
2211 )
2212 .await?;
2213 wait_real_turn(&rig.server.runtime, &thread, &successor.id).await?;
2214 let requests = rig.mock.captured_requests();
2215 assert_eq!(
2216 requests.len(),
2217 2,
2218 "cancelled approval cannot replay the original model request"
2219 );
2220 assert_ne!(
2221 requests[0].tools, requests[1].tools,
2222 "cancelled ACP narrowing must not leak into its ordinary successor"
2223 );
2224 assert!(
2225 !target.exists(),
2226 "late allow cannot execute after cancellation"
2227 );
2228 }
2229 rig.close().await;
2230 Ok(())
2231 }
2232 #[tokio::test(flavor = "current_thread")]
2233 async fn normal_owner_acp_cancel_then_ordinary_successor_keeps_effects_and_profile_scoped()
2234 -> Result<()> {
2235 permission_case_with_base(PermissionReply::CancelLateAllow, false).await
2236 }
2237
2238 async fn wait_real_turn(
2239 manager: &RuntimeThreadManager,
2240 thread: &str,
2241 turn: &str,
2242 ) -> Result<crate::runtime_threads::TurnRecord> {
2243 tokio::time::timeout(Duration::from_secs(30), async {
2244 loop {
2245 let detail = manager.get_thread_detail(thread).await?;
2246 if let Some(record) = detail.turns.iter().find(|record| {
2247 record.id == turn
2248 && !matches!(
2249 record.status,
2250 crate::runtime_threads::RuntimeTurnStatus::InProgress
2251 | crate::runtime_threads::RuntimeTurnStatus::Queued
2252 )
2253 }) && manager
2254 .events_since_async(thread, None)
2255 .await?
2256 .iter()
2257 .any(|event| {
2258 event.event == "turn.completed" && event.turn_id.as_deref() == Some(turn)
2259 })
2260 {
2261 return Ok::<_, anyhow::Error>(record.clone());
2262 }
2263 tokio::time::sleep(Duration::from_millis(20)).await;
2264 }
2265 })
2266 .await?
2267 }
2268
2269 #[tokio::test(flavor = "current_thread")]
2270 async fn real_normal_engine_acp_operation_replay_is_profile_bound_and_successor_is_ordinary()
2271 -> Result<()> {
2272 let mut rig = Rig::with_base(
2273 fixture_config(),
2274 vec![canned::simple_text_turn("ACP initial")],
2275 false,
2276 )?;
2277 let session = rig.new_session().await;
2278 let thread = rig.server.sessions[&session].thread_id.clone();
2279 let request = crate::runtime_threads::StartTurnRequest {
2280 prompt: "same semantic request".into(),
2281 operation_key: Some("acp-operation".into()),
2282 ..Default::default()
2283 };
2284 let first = rig
2285 .server
2286 .runtime
2287 .start_acp_turn(&thread, request.clone())
2288 .await?;
2289 let replay = rig
2290 .server
2291 .runtime
2292 .start_acp_turn(&thread, request.clone())
2293 .await?;
2294 assert_eq!(replay.id, first.id);
2295 assert!(
2296 rig.server
2297 .runtime
2298 .start_turn(&thread, request.clone())
2299 .await
2300 .is_err(),
2301 "ordinary admission cannot redeem an ACP operation binding"
2302 );
2303 assert_eq!(
2304 wait_real_turn(&rig.server.runtime, &thread, &first.id)
2305 .await?
2306 .status,
2307 crate::runtime_threads::RuntimeTurnStatus::Completed
2308 );
2309 assert_eq!(
2310 rig.mock.call_count(),
2311 1,
2312 "exact replay executes the Engine once"
2313 );
2314 rig.mock
2315 .push_turn(canned::simple_text_turn("ordinary initial"));
2316 let ordinary_request = crate::runtime_threads::StartTurnRequest {
2317 operation_key: Some("ordinary-operation".into()),
2318 ..request
2319 };
2320 let ordinary = rig
2321 .server
2322 .runtime
2323 .start_turn(&thread, ordinary_request.clone())
2324 .await?;
2325 assert!(
2326 rig.server
2327 .runtime
2328 .start_acp_turn(&thread, ordinary_request.clone())
2329 .await
2330 .is_err(),
2331 "ACP admission cannot redeem an ordinary historical binding"
2332 );
2333 assert_eq!(
2334 wait_real_turn(&rig.server.runtime, &thread, &ordinary.id)
2335 .await?
2336 .status,
2337 crate::runtime_threads::RuntimeTurnStatus::Completed
2338 );
2339 assert_eq!(
2340 rig.server
2341 .runtime
2342 .start_turn(&thread, ordinary_request)
2343 .await?
2344 .id,
2345 ordinary.id
2346 );
2347 assert_eq!(rig.mock.call_count(), 2);
2348 assert_eq!(
2349 rig.server
2350 .runtime
2351 .get_thread_detail(&thread)
2352 .await?
2353 .turns
2354 .len(),
2355 2
2356 );
2357 let requests = rig.mock.captured_requests();
2358 assert_ne!(
2359 requests[0].tools, requests[1].tools,
2360 "ordinary successor must retain its own catalog"
2361 );
2362 rig.close().await;
2363 Ok(())
2364 }
2365
2366 #[tokio::test(flavor = "current_thread")]
2367 async fn acp_permission_allow_after_wrong_id_executes_once() -> Result<()> {
2368 permission_case(PermissionReply::AllowAfterWrong).await
2369 }
2370 #[tokio::test(flavor = "current_thread")]
2371 async fn foreign_cancel_and_busy_request_leave_the_same_core_approval_pending() -> Result<()> {
2372 permission_case(PermissionReply::ForeignCancelAndBusy).await
2373 }
2374 #[tokio::test(flavor = "current_thread")]
2375 async fn acp_permission_invalid_option_refuses_without_effect() -> Result<()> {
2376 permission_case(PermissionReply::Reject).await
2377 }
2378 #[tokio::test(flavor = "current_thread")]
2379 async fn acp_permission_cancel_ignores_late_allow_and_never_runs_tool() -> Result<()> {
2380 permission_case(PermissionReply::CancelLateAllow).await
2381 }
2382 #[tokio::test(flavor = "current_thread")]
2383 async fn plan_and_operator_floor_cannot_be_relaxed_by_full_access_or_client() -> Result<()> {
2384 let mut config = fixture_config();
2385 config.yolo = Some(true);
2386 config.sandbox_mode = Some("read-only".into());
2387 let mut rig = Rig::new(
2388 config,
2389 vec![
2390 canned::tool_call_turn(
2391 "write-denied",
2392 "write",
2393 r#"{"path":"denied.txt","content":"forbidden"}"#,
2394 ),
2395 canned::simple_text_turn("read only"),
2396 ],
2397 )?;
2398 let id = rig.new_session().await;
2399 assert!(
2400 rig.server
2401 .set_session_config(json!({"sessionId":id,"configId":"mode","value":"agent"}))
2402 .await
2403 .is_err()
2404 );
2405 let view = rig.server.session_configuration(&id).await.unwrap();
2406 assert_eq!(view["modes"]["currentModeId"], "plan");
2407 assert_eq!(
2408 view["configOptions"][2]["options"]
2409 .as_array()
2410 .unwrap()
2411 .len(),
2412 1
2413 );
2414 let (_, wire) = rig.prompt(&id, "try write").await?;
2415 assert!(!rig.workspace.join("denied.txt").exists());
2416 assert!(
2417 !wire
2418 .iter()
2419 .any(|v| v["method"] == "session/request_permission")
2420 );
2421 rig.close().await;
2422 Ok(())
2423 }
2424 #[tokio::test(flavor = "current_thread")]
2425 async fn foreign_canonical_store_refuses_before_allocating_or_copying_history() -> Result<()> {
2426 let mut rig = Rig::new(fixture_config(), vec![])?;
2427 let id = rig.new_session().await;
2428 let store = crate::session_manager::SessionManager::new(rig.server.sessions_dir.clone())?;
2429 let mut saved = store.load_session(&id)?;
2430 let path = store.save_session(&saved)?;
2431 let original = std::fs::read(&path)?;
2432 saved.metadata.runtime_store.as_mut().unwrap().data_dir =
2433 rig.workspace.join("foreign-owner");
2434 let write_error = store.save_session(&saved).unwrap_err();
2435 assert_eq!(write_error.kind(), std::io::ErrorKind::PermissionDenied);
2436 assert_eq!(std::fs::read(&path)?, original);
2437 // A normal writer rejects this binding. Inject a synthetic stale
2438 // document only into this private fixture to exercise load refusal.
2439 std::fs::write(&path, serde_json::to_vec_pretty(&saved)?)?;
2440 rig.server.sessions.clear();
2441 rig.server.insertion_order.clear();
2442 let before = rig
2443 .server
2444 .runtime
2445 .list_threads(
2446 crate::runtime_threads::ThreadListFilter::IncludeArchived,
2447 None,
2448 )
2449 .await?
2450 .len();
2451 let error = rig
2452 .server
2453 .load_session(json!({"sessionId":id}))
2454 .await
2455 .unwrap_err();
2456 assert!(error.message.contains("secure owner attachment"));
2457 assert_eq!(
2458 rig.server
2459 .runtime
2460 .list_threads(
2461 crate::runtime_threads::ThreadListFilter::IncludeArchived,
2462 None
2463 )
2464 .await?
2465 .len(),
2466 before
2467 );
2468 assert_eq!(rig.mock.call_count(), 0);
2469 rig.close().await;
2470 Ok(())
2471 }
2472 #[tokio::test(flavor = "current_thread")]
2473 async fn transport_eof_cancels_actual_turn_and_retains_checkpoint() -> Result<()> {
2474 permission_case(PermissionReply::Eof).await
2475 }
2476 #[tokio::test(flavor = "current_thread")]
2477 async fn real_image_tool_result_projects_typed_blocks_without_losing_core_history() -> Result<()>
2478 {
2479 use base64::Engine as _;
2480 let mut rig = Rig::new(
2481 fixture_config(),
2482 vec![
2483 canned::tool_call_turn("image", "read", r#"{"path":"shot.png"}"#),
2484 canned::simple_text_turn("image read"),
2485 ],
2486 )?;
2487 let image = crate::image_attach::tests::runtime_image_fixture(21);
2488 std::fs::write(
2489 rig.workspace.join("shot.png"),
2490 base64::engine::general_purpose::STANDARD.decode(image.data_base64)?,
2491 )?;
2492 let id = rig.new_session().await;
2493 let (_, wire) = rig.prompt(&id, "read image").await?;
2494 assert!(
2495 wire.iter()
2496 .filter_map(|v| v
2497 .pointer("/params/update/content")
2498 .and_then(Value::as_array))
2499 .flatten()
2500 .any(|block| block.pointer("/content/type") == Some(&json!("image")))
2501 );
2502 assert!(rig.history(&id).await.iter().flat_map(|m|&m.content).any(|b|matches!(b,ContentBlock::ToolResult{content_blocks:Some(blocks),..} if !blocks.is_empty())));
2503 rig.close().await;
2504 Ok(())
2505 }
2506 #[tokio::test(flavor = "current_thread")]
2507 async fn finite_engine_step_budget_is_a_typed_acp_stop_reason() -> Result<()> {
2508 finite_engine_step_budget_case(true).await
2509 }
2510 #[tokio::test(flavor = "current_thread")]
2511 async fn normal_owner_acp_step_stop_preserves_ordinary_successor() -> Result<()> {
2512 finite_engine_step_budget_case(false).await
2513 }
2514 async fn finite_engine_step_budget_case(base_acp: bool) -> Result<()> {
2515 let mut config = fixture_config();
2516 config.tui = Some(crate::config::TuiConfig {
2517 max_model_steps: Some(1),
2518 ..Default::default()
2519 });
2520 let mut rig = Rig::with_base(
2521 config,
2522 vec![
2523 canned::tool_call_turn("budget", "read", r#"{"path":"one.txt"}"#),
2524 canned::simple_text_turn("bounded final report"),
2525 ],
2526 base_acp,
2527 )?;
2528 std::fs::write(rig.workspace.join("one.txt"), "one")?;
2529 let id = rig.new_session().await;
2530 assert_eq!(
2531 rig.prompt(&id, "read until budget").await?.0,
2532 "max_turn_requests"
2533 );
2534 let detail = rig
2535 .server
2536 .runtime
2537 .get_thread_detail(&rig.server.sessions[&id].thread_id)
2538 .await?;
2539 assert_eq!(
2540 detail
2541 .turns
2542 .last()
2543 .unwrap()
2544 .model_request_diagnostics
2545 .as_ref()
2546 .unwrap()
2547 .stop_reason,
2548 Some(crate::runtime_threads::RuntimeTurnStopReason::StepBudgetExhausted)
2549 );
2550 assert_eq!(
2551 rig.mock.call_count(),
2552 2,
2553 "one admitted model step plus existing Core final-report step"
2554 );
2555 if !base_acp {
2556 rig.mock
2557 .push_turn(canned::simple_text_turn("ordinary after ACP step stop"));
2558 let thread = rig.server.sessions[&id].thread_id.clone();
2559 let successor = rig
2560 .server
2561 .runtime
2562 .start_turn(
2563 &thread,
2564 crate::runtime_threads::StartTurnRequest {
2565 prompt: "ordinary after step stop".into(),
2566 ..Default::default()
2567 },
2568 )
2569 .await?;
2570 assert_eq!(
2571 wait_real_turn(&rig.server.runtime, &thread, &successor.id)
2572 .await?
2573 .status,
2574 crate::runtime_threads::RuntimeTurnStatus::Completed
2575 );
2576 let requests = rig.mock.captured_requests();
2577 assert_eq!(requests.len(), 3);
2578 assert_ne!(requests[0].tools, requests[2].tools);
2579 }
2580 rig.close().await;
2581 Ok(())
2582 }
2583 #[tokio::test(flavor = "current_thread")]
2584 async fn invalid_prompt_and_rpc_shapes_refuse_before_core_admission() -> Result<()> {
2585 let mut rig = Rig::new(fixture_config(), vec![])?;
2586 let id = rig.new_session().await;
2587 let input = format!(
2588 "{}\n{}\n{}\n{}\n",
2589 json!({"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"sessionId":id,"prompt":[]}}),
2590 json!({"jsonrpc":"2.0","id":4,"method":"session/prompt","params":{"sessionId":"unknown","prompt":"x"}}),
2591 json!({"jsonrpc":"2.0","id":5,"method":"fabricated"}),
2592 json!({"jsonrpc":"2.0","id":6,"method":"shutdown"})
2593 );
2594 let mut reader = codewhale_app_server::BoundedLines::new(BufReader::new(input.as_bytes()));
2595 let mut output = Vec::new();
2596 rig.server.serve(&mut reader, &mut output).await?;
2597 let wire = parse_lines(output);
2598 assert_eq!(wire[0]["error"]["code"], -32602);
2599 assert_eq!(wire[1]["error"]["code"], -32602);
2600 assert_eq!(wire[2]["error"]["code"], -32601);
2601 assert_eq!(wire[0]["id"], 3);
2602 assert_eq!(rig.mock.call_count(), 0);
2603 assert!(
2604 rig.server
2605 .runtime
2606 .get_thread_detail(&rig.server.sessions[&id].thread_id)
2607 .await?
2608 .turns
2609 .is_empty()
2610 );
2611 rig.close().await;
2612 Ok(())
2613 }
2614 #[tokio::test(flavor = "current_thread")]
2615 async fn canonical_thread_shell_ceiling_requires_operator_and_client_opt_in() -> Result<()> {
2616 for (operator, client, expected) in [
2617 (true, false, false),
2618 (false, true, false),
2619 (true, true, true),
2620 ] {
2621 let mut config = fixture_config();
2622 config.allow_shell = Some(operator);
2623 let operator_config = format!("allow_shell = {operator}\n");
2624 let mut rig = Rig::new(config, vec![])?;
2625 let config_path = rig._dir.path().join("operator.toml");
2626 std::fs::write(&config_path, operator_config)?;
2627 rig.server.config_path = Some(config_path);
2628 rig.server.client_supports_terminal = client;
2629 let id = rig.new_session().await;
2630 let thread = rig
2631 .server
2632 .runtime
2633 .get_thread(&rig.server.sessions[&id].thread_id)
2634 .await?;
2635 assert_eq!(thread.allow_shell, expected);
2636 assert_eq!(rig.mock.call_count(), 0);
2637 rig.close().await;
2638 }
2639 let mut config = fixture_config();
2640 config.allow_shell = Some(true);
2641 let mut rig = Rig::new(config, vec![])?;
2642 rig.server.client_supports_terminal = true;
2643 let error = rig.server.new_session(json!({})).await.unwrap_err();
2644 assert!(error.message.contains("unresolved configuration source"));
2645 assert!(rig.server.sessions.is_empty());
2646 assert!(
2647 rig.server
2648 .runtime
2649 .list_threads(
2650 crate::runtime_threads::ThreadListFilter::IncludeArchived,
2651 None,
2652 )
2653 .await?
2654 .is_empty()
2655 );
2656 assert_eq!(rig.mock.call_count(), 0);
2657 rig.close().await;
2658 Ok(())
2659 }
2660 #[tokio::test(flavor = "current_thread")]
2661 async fn model_discovery_and_selection_keep_custom_provider_identity_and_refuse_invalid_values()
2662 -> Result<()> {
2663 let mut rig = Rig::new(fixture_config(), vec![])?;
2664 assert_eq!(rig.server.current_model()["provider"], "acp-fixture");
2665 assert_eq!(rig.server.current_model()["model"], "fixture-model");
2666 assert!(
2667 rig.server.list_providers()["providers"]
2668 .as_array()
2669 .unwrap()
2670 .iter()
2671 .any(|p| p["id"] == "acp-fixture")
2672 );
2673 assert!(
2674 rig.server
2675 .select_model(json!({"provider":"not-a-provider","model":"x"}))
2676 .is_err()
2677 );
2678 assert!(
2679 rig.server
2680 .select_model(json!({"provider":"acp-fixture"}))
2681 .is_err()
2682 );
2683 let selected = rig
2684 .server
2685 .select_model(json!({"provider":"acp-fixture","model":"fixture-model"}))
2686 .unwrap();
2687 assert_eq!(selected["provider"], "acp-fixture");
2688 assert_eq!(rig.mock.call_count(), 0);
2689 rig.close().await;
2690 Ok(())
2691 }
2692 }
2693
2693 lines RUST