1use std::collections::{BTreeMap, HashMap, VecDeque};
9use std::io::Write;
10use std::sync::atomic::{AtomicBool, AtomicI64, AtomicU64, AtomicUsize, Ordering};
11use std::sync::{Arc, Weak};
12use std::time::Duration;
13
14use scc::{HashMap as SccHashMap, HashSet as SccHashSet};
15use serde::Deserialize;
16use serde_json::{Map, Value};
17use tokio::sync::{broadcast, oneshot, watch};
18use tokio::task::JoinHandle;
19
20use agentos_protocol::generated::v1::{
21 AcpCallback, AcpCallbackResponse, AcpEvent, AcpHostRequestCallbackResponse,
22 AcpPermissionCallbackResponse,
23};
24use agentos_protocol::ACP_EXTENSION_NAMESPACE;
25use agentos_sidecar_client::wire;
26use agentos_vm_config as vm_config;
27
28use crate::config::{
29 AgentOsConfig, AgentOsLimits, HostTool, MountConfig, PermissionMode, Permissions,
30 RootFilesystemConfig, RootFilesystemKind, RootFilesystemMode as ConfigRootFilesystemMode,
31 RootLowerInput, SidecarJsBridgeCall, SidecarJsBridgeCallback, TimerScheduleDriver, ToolKit,
32};
33use crate::cron::CronManager;
34use crate::error::ClientError;
35use crate::json_rpc::JsonRpcNotification;
36use crate::process::SYNTHETIC_PID_BASE;
37use crate::session::{
38 record_live_session_event, AgentCapabilities, AgentExitEvent, AgentInfo, PermissionReply,
39 PermissionRequest, PermissionRouteRequest, PermissionRouteResult, SessionConfigOption,
40 SessionModeState,
41};
42use crate::sidecar::{AgentOsSidecar, AgentOsSidecarPlacement, AgentOsSidecarVmLease};
43use crate::transport::{SidecarProcess, WireSidecarCallback};
44use agentos_sidecar_client::TransportError;
45
46use once_cell::sync::OnceCell;
47
48pub(crate) struct ProcessEntry {
54 pub command: String,
55 pub args: Vec<String>,
56 pub stdout_tx: broadcast::Sender<Vec<u8>>,
57 pub stderr_tx: broadcast::Sender<Vec<u8>>,
58 pub exit_tx: watch::Sender<Option<i32>>,
60 pub process_id: String,
62 pub kernel_pid: watch::Sender<Option<u32>>,
66 pub output_tasks: Vec<JoinHandle<()>>,
70 pub started_at: i64,
72}
73
74pub(crate) struct ShellEntry {
80 pub pid: u32,
81 pub data_tx: broadcast::Sender<Vec<u8>>,
82 pub stderr_tx: broadcast::Sender<Vec<u8>>,
83 pub process_id: String,
85 pub spawned_tx: watch::Sender<bool>,
90 pub exit_tx: watch::Sender<Option<i32>>,
93}
94
95pub(crate) struct AcpTerminalEntry {
97 pub exit_task: JoinHandle<()>,
98}
99
100pub(crate) struct HostAcpTerminalOutput {
103 pub buffer: String,
105 pub truncated: bool,
106 pub output_byte_limit: usize,
109}
110
111pub(crate) struct HostAcpTerminal {
115 pub shell_id: String,
118 pub output: Arc<parking_lot::Mutex<HostAcpTerminalOutput>>,
120 pub exit_rx: watch::Receiver<Option<i32>>,
122}
123
124pub(crate) struct SessionEntry {
126 pub agent_type: String,
127 pub modes: parking_lot::Mutex<Option<SessionModeState>>,
128 pub config_options: parking_lot::Mutex<Vec<SessionConfigOption>>,
129 pub capabilities: parking_lot::Mutex<Option<AgentCapabilities>>,
130 pub agent_info: parking_lot::Mutex<Option<AgentInfo>>,
131 pub config_overrides: parking_lot::Mutex<std::collections::BTreeMap<String, String>>,
132 pub event_tx: broadcast::Sender<JsonRpcNotification>,
133 pub permission_tx: broadcast::Sender<PermissionRequest>,
134 pub agent_exit_tx: broadcast::Sender<AgentExitEvent>,
135 pub pending_permission_replies: SccHashMap<String, oneshot::Sender<PermissionReply>>,
136 pub pending_session_request_lock: parking_lot::Mutex<()>,
137 pub pending_prompt_resolvers:
145 SccHashMap<i64, oneshot::Sender<crate::json_rpc::JsonRpcResponse>>,
146}
147
148#[derive(Debug, Clone)]
158pub struct PackageDescriptor {
159 pub path: String,
160}
161
162#[derive(Debug, Clone, PartialEq, Eq)]
163pub struct ProjectedAgent {
164 pub id: String,
165 pub acp_entrypoint: String,
166 pub adapter_entrypoint: String,
167}
168
169#[derive(Clone)]
171pub struct AgentOs {
172 inner: Arc<AgentOsInner>,
173}
174
175pub(crate) struct AgentOsInner {
176 pub(crate) transport: Arc<SidecarProcess>,
178 pub(crate) connection_id: String,
179 pub(crate) session_id: String,
180 pub(crate) vm_id: String,
181 pub(crate) request_counter: AtomicI64,
182 pub(crate) projected_commands: parking_lot::Mutex<BTreeMap<String, String>>,
184 pub(crate) projected_agents: parking_lot::Mutex<Vec<ProjectedAgent>>,
186
187 pub(crate) process_registry_lock: parking_lot::Mutex<()>,
189 pub(crate) processes: SccHashMap<u32, ProcessEntry>,
190 pub(crate) process_counter: AtomicU64,
194 pub(crate) synthetic_pid_counter: AtomicU64,
197 pub(crate) observed_process_time_lock: parking_lot::Mutex<()>,
198 pub(crate) observed_process_start_times: SccHashMap<String, f64>,
202 pub(crate) observed_process_exit_times: SccHashMap<String, f64>,
205
206 pub(crate) shells: SccHashMap<String, ShellEntry>,
208 pub(crate) shell_counter: AtomicU64,
209 pub(crate) pending_shell_exits: SccHashMap<u64, JoinHandle<()>>,
210 pub(crate) closed_shell_exit_codes: parking_lot::Mutex<VecDeque<(String, i32)>>,
214 pub(crate) acp_terminals: SccHashMap<String, AcpTerminalEntry>,
215 pub(crate) acp_terminal_count: AtomicUsize,
216 pub(crate) acp_terminal_lifecycle_lock: tokio::sync::Mutex<()>,
217 pub(crate) host_acp_terminals: SccHashMap<String, HostAcpTerminal>,
220 pub(crate) host_acp_terminal_counter: AtomicU64,
222
223 pub(crate) sessions: SccHashMap<String, SessionEntry>,
225 pub(crate) closed_session_ids: parking_lot::Mutex<VecDeque<String>>,
227 pub(crate) closing_session_ids: SccHashSet<String>,
232
233 pub(crate) cron: Arc<CronManager>,
235
236 pub(crate) config: Arc<AgentOsConfig>,
238 pub(crate) sidecar: Arc<AgentOsSidecar>,
239 pub(crate) sidecar_lease: parking_lot::Mutex<Option<AgentOsSidecarVmLease>>,
240 pub(crate) in_process_mounts: SccHashMap<String, crate::fs::MountedFs>,
241 pub(crate) disposed: AtomicBool,
242 pub(crate) acp_event_pump: parking_lot::Mutex<Option<JoinHandle<()>>>,
246}
247
248impl AgentOs {
249 pub async fn create(options: AgentOsConfig) -> Result<AgentOs, ClientError> {
253 let config = Arc::new(options);
254
255 let sidecar = match &config.sidecar {
259 Some(crate::config::AgentOsSidecarConfig::Explicit { handle }) => handle.clone(),
260 Some(crate::config::AgentOsSidecarConfig::Shared { pool }) => {
261 AgentOs::get_shared_sidecar(pool.clone(), config.sidecar_binary_path.clone())
262 .await?
263 }
264 None => AgentOs::get_shared_sidecar(None, config.sidecar_binary_path.clone()).await?,
265 };
266 let (transport, connection_id, _) = sidecar.ensure_connection().await?;
267
268 let session = match transport
270 .request_wire(
271 wire_connection_ownership(&connection_id),
272 wire::RequestPayload::OpenSessionRequest(wire::OpenSessionRequest {
273 placement: sidecar_wire_placement(&sidecar),
274 metadata: HashMap::new(),
275 }),
276 )
277 .await?
278 {
279 wire::ResponsePayload::SessionOpenedResponse(opened) => opened,
280 wire::ResponsePayload::RejectedResponse(rejected) => {
281 return Err(rejected_to_error(rejected));
282 }
283 wire::ResponsePayload::AuthenticatedResponse(_)
284 | wire::ResponsePayload::VmCreatedResponse(_)
285 | wire::ResponsePayload::VmDisposedResponse(_)
286 | wire::ResponsePayload::RootFilesystemBootstrappedResponse(_)
287 | wire::ResponsePayload::VmConfiguredResponse(_)
288 | wire::ResponsePayload::HostCallbacksRegisteredResponse(_)
289 | wire::ResponsePayload::LayerCreatedResponse(_)
290 | wire::ResponsePayload::LayerSealedResponse(_)
291 | wire::ResponsePayload::SnapshotImportedResponse(_)
292 | wire::ResponsePayload::SnapshotExportedResponse(_)
293 | wire::ResponsePayload::OverlayCreatedResponse(_)
294 | wire::ResponsePayload::GuestFilesystemResultResponse(_)
295 | wire::ResponsePayload::RootFilesystemSnapshotResponse(_)
296 | wire::ResponsePayload::ProcessStartedResponse(_)
297 | wire::ResponsePayload::StdinWrittenResponse(_)
298 | wire::ResponsePayload::PtyResizedResponse(_)
299 | wire::ResponsePayload::StdinClosedResponse(_)
300 | wire::ResponsePayload::ProcessKilledResponse(_)
301 | wire::ResponsePayload::ProcessSnapshotResponse(_)
302 | wire::ResponsePayload::ListenerSnapshotResponse(_)
303 | wire::ResponsePayload::BoundUdpSnapshotResponse(_)
304 | wire::ResponsePayload::SignalStateResponse(_)
305 | wire::ResponsePayload::ZombieTimerCountResponse(_)
306 | wire::ResponsePayload::FilesystemResultResponse(_)
307 | wire::ResponsePayload::PermissionDecisionResponse(_)
308 | wire::ResponsePayload::PersistenceStateResponse(_)
309 | wire::ResponsePayload::PersistenceFlushedResponse(_)
310 | wire::ResponsePayload::VmFetchResponse(_)
311 | wire::ResponsePayload::ExtEnvelope(_)
312 | wire::ResponsePayload::GuestKernelResultResponse(_)
313 | wire::ResponsePayload::ResourceSnapshotResponse(_)
314 | wire::ResponsePayload::PackageLinkedResponse(_)
315 | wire::ResponsePayload::ProvidedCommandsResponse(_) => {
316 return Err(ClientError::Sidecar(
317 "unexpected open_session response".to_string(),
318 ));
319 }
320 };
321 let session_id = session.session_id;
322
323 let mut events = transport.subscribe_wire_events();
325 let permissions = permissions_policy(&config);
326 let create_vm_config = serialize_create_vm_config_for_sidecar(&config)?;
327 if let Some(callback) = config.sidecar_js_bridge_callback.clone() {
328 let _ = session_js_bridge_callbacks()
329 .insert(sidecar_session_key(&connection_id, &session_id), callback);
330 transport.register_wire_callback("js_bridge_call", js_bridge_call_callback());
331 }
332
333 let vm = match transport
335 .request_wire(
336 wire_session_ownership(&connection_id, &session_id),
337 wire::RequestPayload::CreateVmRequest(wire::CreateVmRequest {
338 runtime: wire::GuestRuntimeKind::JavaScript,
339 config: serde_json::to_string(&create_vm_config).map_err(|error| {
340 ClientError::Sidecar(format!(
341 "failed to serialize create VM config: {error}"
342 ))
343 })?,
344 }),
345 )
346 .await?
347 {
348 wire::ResponsePayload::VmCreatedResponse(created) => created,
349 wire::ResponsePayload::RejectedResponse(rejected) => {
350 return Err(rejected_to_error(rejected));
351 }
352 wire::ResponsePayload::AuthenticatedResponse(_)
353 | wire::ResponsePayload::SessionOpenedResponse(_)
354 | wire::ResponsePayload::VmDisposedResponse(_)
355 | wire::ResponsePayload::RootFilesystemBootstrappedResponse(_)
356 | wire::ResponsePayload::VmConfiguredResponse(_)
357 | wire::ResponsePayload::HostCallbacksRegisteredResponse(_)
358 | wire::ResponsePayload::LayerCreatedResponse(_)
359 | wire::ResponsePayload::LayerSealedResponse(_)
360 | wire::ResponsePayload::SnapshotImportedResponse(_)
361 | wire::ResponsePayload::SnapshotExportedResponse(_)
362 | wire::ResponsePayload::OverlayCreatedResponse(_)
363 | wire::ResponsePayload::GuestFilesystemResultResponse(_)
364 | wire::ResponsePayload::RootFilesystemSnapshotResponse(_)
365 | wire::ResponsePayload::ProcessStartedResponse(_)
366 | wire::ResponsePayload::StdinWrittenResponse(_)
367 | wire::ResponsePayload::PtyResizedResponse(_)
368 | wire::ResponsePayload::StdinClosedResponse(_)
369 | wire::ResponsePayload::ProcessKilledResponse(_)
370 | wire::ResponsePayload::ProcessSnapshotResponse(_)
371 | wire::ResponsePayload::ListenerSnapshotResponse(_)
372 | wire::ResponsePayload::BoundUdpSnapshotResponse(_)
373 | wire::ResponsePayload::SignalStateResponse(_)
374 | wire::ResponsePayload::ZombieTimerCountResponse(_)
375 | wire::ResponsePayload::FilesystemResultResponse(_)
376 | wire::ResponsePayload::PermissionDecisionResponse(_)
377 | wire::ResponsePayload::PersistenceStateResponse(_)
378 | wire::ResponsePayload::PersistenceFlushedResponse(_)
379 | wire::ResponsePayload::VmFetchResponse(_)
380 | wire::ResponsePayload::ExtEnvelope(_)
381 | wire::ResponsePayload::GuestKernelResultResponse(_)
382 | wire::ResponsePayload::ResourceSnapshotResponse(_)
383 | wire::ResponsePayload::PackageLinkedResponse(_)
384 | wire::ResponsePayload::ProvidedCommandsResponse(_) => {
385 return Err(ClientError::Sidecar(
386 "unexpected create_vm response".to_string(),
387 ));
388 }
389 };
390 let vm_id = vm.vm_id;
391
392 wait_for_vm_ready(&mut events, &vm_id, crate::VM_READY_TIMEOUT_MS).await?;
394
395 let packages = build_package_descriptors(&config);
398
399 let mounts = serialize_mounts(&config)?;
401
402 let (projected_commands, projected_agents) = match transport
406 .request_wire(
407 wire_vm_ownership(&connection_id, &session_id, &vm_id),
408 wire::RequestPayload::ConfigureVmRequest(wire::ConfigureVmRequest {
409 mounts,
410 software: Vec::new(),
413 permissions: Some(permissions),
414 module_access_cwd: None,
418 instructions: config.additional_instructions.clone().into_iter().collect(),
419 projected_modules: Vec::new(),
420 command_permissions: HashMap::new(),
421 loopback_exempt_ports: config.loopback_exempt_ports.clone(),
422 packages,
423 packages_mount_at: config.packages_mount_at.clone().unwrap_or_default(),
424 bootstrap_commands: Vec::new(),
425 tool_shim_commands: Vec::new(),
426 }),
427 )
428 .await?
429 {
430 wire::ResponsePayload::VmConfiguredResponse(configured) => (
431 configured
432 .projected_commands
433 .into_iter()
434 .map(|command| (command.name, command.guest_path))
435 .collect(),
436 projected_agents_from_wire(configured.agents),
437 ),
438 wire::ResponsePayload::RejectedResponse(rejected) => {
439 return Err(rejected_to_error(rejected));
440 }
441 wire::ResponsePayload::AuthenticatedResponse(_)
442 | wire::ResponsePayload::SessionOpenedResponse(_)
443 | wire::ResponsePayload::VmCreatedResponse(_)
444 | wire::ResponsePayload::VmDisposedResponse(_)
445 | wire::ResponsePayload::RootFilesystemBootstrappedResponse(_)
446 | wire::ResponsePayload::HostCallbacksRegisteredResponse(_)
447 | wire::ResponsePayload::LayerCreatedResponse(_)
448 | wire::ResponsePayload::LayerSealedResponse(_)
449 | wire::ResponsePayload::SnapshotImportedResponse(_)
450 | wire::ResponsePayload::SnapshotExportedResponse(_)
451 | wire::ResponsePayload::OverlayCreatedResponse(_)
452 | wire::ResponsePayload::GuestFilesystemResultResponse(_)
453 | wire::ResponsePayload::RootFilesystemSnapshotResponse(_)
454 | wire::ResponsePayload::ProcessStartedResponse(_)
455 | wire::ResponsePayload::StdinWrittenResponse(_)
456 | wire::ResponsePayload::PtyResizedResponse(_)
457 | wire::ResponsePayload::StdinClosedResponse(_)
458 | wire::ResponsePayload::ProcessKilledResponse(_)
459 | wire::ResponsePayload::ProcessSnapshotResponse(_)
460 | wire::ResponsePayload::ListenerSnapshotResponse(_)
461 | wire::ResponsePayload::BoundUdpSnapshotResponse(_)
462 | wire::ResponsePayload::SignalStateResponse(_)
463 | wire::ResponsePayload::ZombieTimerCountResponse(_)
464 | wire::ResponsePayload::FilesystemResultResponse(_)
465 | wire::ResponsePayload::PermissionDecisionResponse(_)
466 | wire::ResponsePayload::PersistenceStateResponse(_)
467 | wire::ResponsePayload::PersistenceFlushedResponse(_)
468 | wire::ResponsePayload::VmFetchResponse(_)
469 | wire::ResponsePayload::ExtEnvelope(_)
470 | wire::ResponsePayload::GuestKernelResultResponse(_)
471 | wire::ResponsePayload::ResourceSnapshotResponse(_)
472 | wire::ResponsePayload::PackageLinkedResponse(_)
473 | wire::ResponsePayload::ProvidedCommandsResponse(_) => {
474 return Err(ClientError::Sidecar(
475 "unexpected configure_vm response".to_string(),
476 ));
477 }
478 };
479
480 if !config.tool_kits.is_empty() {
484 let mut tool_map: HashMap<String, HostTool> = HashMap::new();
485 for kit in &config.tool_kits {
486 let mut tools = HashMap::new();
487 for tool in &kit.tools {
488 tools.insert(
489 tool.name.clone(),
490 wire::RegisteredHostCallbackDefinition {
491 description: tool.description.clone(),
492 input_schema: json_utf8(
493 &tool.input_schema,
494 "host callback input schema",
495 )?,
496 timeout_ms: tool.timeout_ms,
497 examples: Vec::new(),
498 },
499 );
500 tool_map.insert(format!("{}:{}", kit.name, tool.name), tool.clone());
501 }
502 match transport
503 .request_wire(
504 wire_vm_ownership(&connection_id, &session_id, &vm_id),
505 wire::RequestPayload::RegisterHostCallbacksRequest(
506 wire::RegisterHostCallbacksRequest {
507 name: kit.name.clone(),
508 description: kit.description.clone(),
509 command_aliases: vec![format!("agentos-{}", kit.name)],
510 registry_command_aliases: vec![String::from("agentos")],
511 callbacks: tools,
512 },
513 ),
514 )
515 .await?
516 {
517 wire::ResponsePayload::HostCallbacksRegisteredResponse(_) => {}
518 wire::ResponsePayload::RejectedResponse(rejected) => {
519 return Err(rejected_to_error(rejected));
520 }
521 wire::ResponsePayload::AuthenticatedResponse(_)
522 | wire::ResponsePayload::SessionOpenedResponse(_)
523 | wire::ResponsePayload::VmCreatedResponse(_)
524 | wire::ResponsePayload::VmDisposedResponse(_)
525 | wire::ResponsePayload::RootFilesystemBootstrappedResponse(_)
526 | wire::ResponsePayload::VmConfiguredResponse(_)
527 | wire::ResponsePayload::LayerCreatedResponse(_)
528 | wire::ResponsePayload::LayerSealedResponse(_)
529 | wire::ResponsePayload::SnapshotImportedResponse(_)
530 | wire::ResponsePayload::SnapshotExportedResponse(_)
531 | wire::ResponsePayload::OverlayCreatedResponse(_)
532 | wire::ResponsePayload::GuestFilesystemResultResponse(_)
533 | wire::ResponsePayload::RootFilesystemSnapshotResponse(_)
534 | wire::ResponsePayload::ProcessStartedResponse(_)
535 | wire::ResponsePayload::StdinWrittenResponse(_)
536 | wire::ResponsePayload::PtyResizedResponse(_)
537 | wire::ResponsePayload::StdinClosedResponse(_)
538 | wire::ResponsePayload::ProcessKilledResponse(_)
539 | wire::ResponsePayload::ProcessSnapshotResponse(_)
540 | wire::ResponsePayload::ListenerSnapshotResponse(_)
541 | wire::ResponsePayload::BoundUdpSnapshotResponse(_)
542 | wire::ResponsePayload::SignalStateResponse(_)
543 | wire::ResponsePayload::ZombieTimerCountResponse(_)
544 | wire::ResponsePayload::FilesystemResultResponse(_)
545 | wire::ResponsePayload::PermissionDecisionResponse(_)
546 | wire::ResponsePayload::PersistenceStateResponse(_)
547 | wire::ResponsePayload::PersistenceFlushedResponse(_)
548 | wire::ResponsePayload::VmFetchResponse(_)
549 | wire::ResponsePayload::ExtEnvelope(_)
550 | wire::ResponsePayload::GuestKernelResultResponse(_)
551 | wire::ResponsePayload::ResourceSnapshotResponse(_)
552 | wire::ResponsePayload::PackageLinkedResponse(_)
553 | wire::ResponsePayload::ProvidedCommandsResponse(_) => {
554 return Err(ClientError::Sidecar(
555 "unexpected register_host_callbacks response".to_string(),
556 ));
557 }
558 }
559 }
560 let _ = vm_tools().insert(
561 vm_id.clone(),
562 Arc::new(VmHostToolRegistry {
563 tool_kits: config.tool_kits.clone(),
564 tool_map,
565 permissions: config.permissions.clone(),
566 }),
567 );
568 transport.register_wire_callback("host_callback", host_callback_callback());
569 }
570
571 sidecar.active_vm_count.fetch_add(1, Ordering::SeqCst);
573 let lease = AgentOsSidecarVmLease {
574 sidecar: sidecar.clone(),
575 };
576
577 let driver = config
578 .schedule_driver
579 .clone()
580 .unwrap_or_else(|| Arc::new(TimerScheduleDriver::new()));
581 let cron = Arc::new(CronManager::new(driver));
582
583 let inner = AgentOsInner {
584 transport,
585 connection_id,
586 session_id,
587 vm_id,
588 request_counter: AtomicI64::new(1),
589 projected_commands: parking_lot::Mutex::new(projected_commands),
590 projected_agents: parking_lot::Mutex::new(projected_agents),
591 process_registry_lock: parking_lot::Mutex::new(()),
592 processes: SccHashMap::new(),
593 process_counter: AtomicU64::new(1),
594 synthetic_pid_counter: AtomicU64::new(SYNTHETIC_PID_BASE),
595 observed_process_time_lock: parking_lot::Mutex::new(()),
596 observed_process_start_times: SccHashMap::new(),
597 observed_process_exit_times: SccHashMap::new(),
598 shells: SccHashMap::new(),
599 shell_counter: AtomicU64::new(0),
600 pending_shell_exits: SccHashMap::new(),
601 closed_shell_exit_codes: parking_lot::Mutex::new(VecDeque::new()),
602 acp_terminals: SccHashMap::new(),
603 acp_terminal_count: AtomicUsize::new(0),
604 acp_terminal_lifecycle_lock: tokio::sync::Mutex::new(()),
605 host_acp_terminals: SccHashMap::new(),
606 host_acp_terminal_counter: AtomicU64::new(0),
607 sessions: SccHashMap::new(),
608 closed_session_ids: parking_lot::Mutex::new(VecDeque::new()),
609 closing_session_ids: SccHashSet::new(),
610 cron,
611 config,
612 sidecar,
613 sidecar_lease: parking_lot::Mutex::new(Some(lease)),
614 in_process_mounts: SccHashMap::new(),
615 disposed: AtomicBool::new(false),
616 acp_event_pump: parking_lot::Mutex::new(None),
617 };
618
619 let client = AgentOs {
620 inner: Arc::new(inner),
621 };
622 let _ = vm_permission_routers()
627 .insert(client.inner.vm_id.clone(), Arc::downgrade(&client.inner));
628 client
629 .inner
630 .transport
631 .register_wire_callback("ext", permission_request_callback());
632 spawn_acp_event_pump(&client);
633 Ok(client)
634 }
635
636 pub async fn link_software(&self, descriptor: PackageDescriptor) -> Result<(), ClientError> {
653 let inner = self.inner();
654 let response = self
655 .transport()
656 .request_wire(
657 wire_vm_ownership(&inner.connection_id, &inner.session_id, &inner.vm_id),
658 wire::RequestPayload::LinkPackageRequest(wire::LinkPackageRequest {
659 package: wire::PackageDescriptor {
662 path: descriptor.path,
663 },
664 }),
665 )
666 .await?;
667 match response {
668 wire::ResponsePayload::PackageLinkedResponse(linked) => {
669 let mut guard = inner.projected_commands.lock();
670 for command in linked.projected_commands {
671 guard.insert(command.name, command.guest_path);
672 }
673 register_projected_agents(
674 &inner.projected_agents,
675 projected_agents_from_wire(linked.agents),
676 );
677 Ok(())
678 }
679 wire::ResponsePayload::RejectedResponse(rejected) => Err(rejected_to_error(rejected)),
680 other => Err(ClientError::Sidecar(format!(
681 "unexpected link_package response: {other:?}"
682 ))),
683 }
684 }
685
686 pub async fn provided_commands(&self) -> Result<BTreeMap<String, Vec<String>>, ClientError> {
687 let inner = self.inner();
688 let response = self
689 .transport()
690 .request_wire(
691 wire_vm_ownership(&inner.connection_id, &inner.session_id, &inner.vm_id),
692 wire::RequestPayload::ProvidedCommandsRequest,
693 )
694 .await?;
695 match response {
696 wire::ResponsePayload::ProvidedCommandsResponse(provided) => Ok(provided
697 .packages
698 .into_iter()
699 .map(|package| (package.package_name, package.commands))
700 .collect()),
701 wire::ResponsePayload::RejectedResponse(rejected) => Err(rejected_to_error(rejected)),
702 other => Err(ClientError::Sidecar(format!(
703 "unexpected provided_commands response: {other:?}"
704 ))),
705 }
706 }
707
708 pub async fn shutdown(&self) -> Result<(), ClientError> {
709 if self.inner.disposed.swap(true, Ordering::SeqCst) {
711 return Ok(());
712 }
713
714 self.inner.cron.dispose();
719
720 abort_tracked_task(&self.inner.acp_event_pump);
726 crate::process::drain_process_output_tasks(&self.inner.processes);
727
728 let mut exit_tasks = Vec::new();
731 self.inner.pending_shell_exits.retain(|_, task| {
732 exit_tasks.push(std::mem::replace(task, tokio::spawn(async {})));
733 false
734 });
735
736 {
737 let _terminal_lifecycle_guard = self.inner.acp_terminal_lifecycle_lock.lock().await;
738 let mut terminal_entries = Vec::new();
739 self.inner.acp_terminals.retain(|process_id, entry| {
740 terminal_entries.push((
741 process_id.clone(),
742 std::mem::replace(&mut entry.exit_task, tokio::spawn(async {})),
743 ));
744 false
745 });
746 self.inner.acp_terminal_count.store(0, Ordering::SeqCst);
747 for (process_id, _) in &terminal_entries {
748 let transport = self.transport().clone();
749 let ownership = wire::OwnershipScope::VmOwnership(wire::VmOwnership {
750 connection_id: self.inner.connection_id.clone(),
751 session_id: self.inner.session_id.clone(),
752 vm_id: self.inner.vm_id.clone(),
753 });
754 let process_id = process_id.clone();
755 exit_tasks.push(tokio::spawn(async move {
756 let _ = transport
757 .request_wire(
758 ownership,
759 wire::RequestPayload::KillProcessRequest(wire::KillProcessRequest {
760 process_id,
761 signal: String::from("SIGTERM"),
762 }),
763 )
764 .await;
765 }));
766 }
767 for (_, task) in terminal_entries {
768 exit_tasks.push(task);
769 }
770 }
771
772 let mut host_terminal_shells = Vec::new();
776 self.inner.host_acp_terminals.retain(|_, terminal| {
777 host_terminal_shells.push(terminal.shell_id.clone());
778 false
779 });
780 for shell_id in host_terminal_shells {
781 let _ = self.close_shell(&shell_id);
782 }
783
784 if !exit_tasks.is_empty() {
785 let mut drain_tasks = exit_tasks;
786 if tokio::time::timeout(
787 Duration::from_millis(crate::SHELL_DISPOSE_TIMEOUT_MS),
788 futures::future::join_all(drain_tasks.iter_mut()),
789 )
790 .await
791 .is_err()
792 {
793 for task in drain_tasks {
794 task.abort();
795 }
796 }
797 }
798
799 let lease = self.inner.sidecar_lease.lock().take();
803 let _ = self
804 .transport()
805 .request_wire(
806 wire::OwnershipScope::VmOwnership(wire::VmOwnership {
807 connection_id: self.inner.connection_id.clone(),
808 session_id: self.inner.session_id.clone(),
809 vm_id: self.inner.vm_id.clone(),
810 }),
811 wire::RequestPayload::DisposeVmRequest(wire::DisposeVmRequest {
812 reason: wire::DisposeReason::Requested,
813 }),
814 )
815 .await;
816 let _ = vm_tools().remove(&self.inner.vm_id);
817 let _ = vm_permission_routers().remove(&self.inner.vm_id);
818 let _ = session_js_bridge_callbacks().remove(&sidecar_session_key(
819 &self.inner.connection_id,
820 &self.inner.session_id,
821 ));
822 let sidecar = self.inner.sidecar.clone();
823 if let Some(lease) = lease {
824 lease.dispose().await?;
825 }
826 if sidecar.active_vm_count.load(Ordering::SeqCst) == 0 {
827 sidecar.kill_connection().await;
828 let _ = sidecar.dispose().await;
829 }
830
831 Ok(())
832 }
833
834 pub(crate) fn inner(&self) -> &AgentOsInner {
837 &self.inner
838 }
839
840 pub(crate) fn transport(&self) -> &Arc<SidecarProcess> {
841 &self.inner.transport
842 }
843
844 pub(crate) fn connection_id(&self) -> &str {
845 &self.inner.connection_id
846 }
847
848 pub(crate) fn wire_session_id(&self) -> &str {
849 &self.inner.session_id
850 }
851
852 pub(crate) fn vm_id(&self) -> &str {
853 &self.inner.vm_id
854 }
855
856 pub(crate) fn config(&self) -> &Arc<AgentOsConfig> {
857 &self.inner.config
858 }
859
860 pub(crate) fn cron(&self) -> &Arc<CronManager> {
861 &self.inner.cron
862 }
863
864 pub fn sidecar(&self) -> Arc<AgentOsSidecar> {
867 self.inner.sidecar.clone()
868 }
869
870 pub fn projected_agents(&self) -> Vec<ProjectedAgent> {
871 self.inner.projected_agents.lock().clone()
872 }
873}
874
875fn abort_tracked_task(slot: &parking_lot::Mutex<Option<JoinHandle<()>>>) {
878 if let Some(handle) = slot.lock().take() {
879 handle.abort();
880 }
881}
882
883fn spawn_acp_event_pump(client: &AgentOs) {
884 let mut events = client.transport().subscribe_wire_events();
885 let inner = Arc::downgrade(&client.inner);
886 let handle = tokio::spawn(async move {
887 loop {
888 match events.recv().await {
889 Ok((ownership, wire::EventPayload::ExtEnvelope(envelope))) => {
890 let Some(inner) = inner.upgrade() else {
891 break;
892 };
893 if inner.disposed.load(Ordering::SeqCst) {
894 break;
895 }
896 if wire_ownership_vm_id(&ownership) != Some(inner.vm_id.as_str()) {
897 continue;
898 }
899 if let Err(error) = deliver_acp_ext_event(&inner, envelope) {
900 tracing::warn!(?error, "failed to deliver acp extension event");
901 }
902 }
903 Ok((
904 _,
905 wire::EventPayload::VmLifecycleEvent(_)
906 | wire::EventPayload::ProcessOutputEvent(_)
907 | wire::EventPayload::ProcessExitedEvent(_)
908 | wire::EventPayload::StructuredEvent(_),
909 )) => {}
910 Err(broadcast::error::RecvError::Lagged(_)) => {}
911 Err(broadcast::error::RecvError::Closed) => break,
912 }
913 }
914 });
915 *client.inner.acp_event_pump.lock() = Some(handle);
916}
917
918fn deliver_acp_ext_event(
919 inner: &AgentOsInner,
920 envelope: wire::ExtEnvelope,
921) -> Result<(), ClientError> {
922 if envelope.namespace != ACP_EXTENSION_NAMESPACE {
923 return Ok(());
924 }
925 let event: AcpEvent = serde_bare::from_slice(&envelope.payload)
926 .map_err(|error| ClientError::Sidecar(format!("invalid ACP event: {error}")))?;
927 match event {
928 AcpEvent::AcpSessionEvent(event) => {
929 let notification: JsonRpcNotification = serde_json::from_str(&event.notification)
930 .map_err(|error| {
931 ClientError::Sidecar(format!("invalid ACP session notification: {error}"))
932 })?;
933 let delivered = inner
934 .sessions
935 .read(&event.session_id, |_, entry| {
936 record_live_session_event(entry, notification.clone());
937 })
938 .is_some();
939 if !delivered {
940 tracing::warn!(
941 session_id = event.session_id,
942 "received acp event for unknown session"
943 );
944 }
945 Ok(())
946 }
947 AcpEvent::AcpAgentStderrEvent(event) => {
948 if !event.session_id.is_empty()
949 && inner.sessions.read(&event.session_id, |_, _| ()).is_none()
950 {
951 tracing::warn!(
952 session_id = event.session_id,
953 agent_type = event.agent_type,
954 process_id = event.process_id,
955 "received acp stderr event for unknown session"
956 );
957 }
958
959 let mut stderr = std::io::stderr().lock();
960 if let Err(error) = stderr.write_all(&event.chunk).and_then(|_| stderr.flush()) {
961 tracing::warn!(?error, "failed to write acp stderr event");
962 }
963 Ok(())
964 }
965 AcpEvent::AcpAgentExitedEvent(event) => {
966 tracing::warn!(
967 session_id = event.session_id,
968 agent_type = event.agent_type,
969 process_id = event.process_id,
970 exit_code = ?event.exit_code,
971 restart = event.restart,
972 restart_count = event.restart_count,
973 max_restarts = event.max_restarts,
974 "acp agent adapter exited unexpectedly"
975 );
976 let delivered = inner
977 .sessions
978 .read(&event.session_id, |_, entry| {
979 let _ = entry.agent_exit_tx.send(AgentExitEvent {
980 session_id: event.session_id.clone(),
981 agent_type: event.agent_type.clone(),
982 process_id: event.process_id.clone(),
983 exit_code: event.exit_code,
984 restart: event.restart.clone(),
985 restart_count: event.restart_count,
986 max_restarts: event.max_restarts,
987 });
988 })
989 .is_some();
990 if !delivered {
991 tracing::warn!(
992 session_id = event.session_id,
993 "received acp agent exit event for unknown session"
994 );
995 }
996 Ok(())
997 }
998 }
999}
1000
1001fn sidecar_wire_placement(sidecar: &AgentOsSidecar) -> wire::SidecarPlacement {
1003 match &sidecar.placement {
1004 AgentOsSidecarPlacement::Shared { pool } => {
1005 wire::SidecarPlacement::SidecarPlacementShared(wire::SidecarPlacementShared {
1006 pool: pool.clone(),
1007 })
1008 }
1009 AgentOsSidecarPlacement::Explicit { sidecar_id } => {
1010 wire::SidecarPlacement::SidecarPlacementExplicit(wire::SidecarPlacementExplicit {
1011 sidecar_id: sidecar_id.clone(),
1012 })
1013 }
1014 }
1015}
1016
1017fn wire_connection_ownership(connection_id: &str) -> wire::OwnershipScope {
1018 wire::OwnershipScope::ConnectionOwnership(wire::ConnectionOwnership {
1019 connection_id: connection_id.to_string(),
1020 })
1021}
1022
1023fn wire_session_ownership(connection_id: &str, session_id: &str) -> wire::OwnershipScope {
1024 wire::OwnershipScope::SessionOwnership(wire::SessionOwnership {
1025 connection_id: connection_id.to_string(),
1026 session_id: session_id.to_string(),
1027 })
1028}
1029
1030fn wire_vm_ownership(connection_id: &str, session_id: &str, vm_id: &str) -> wire::OwnershipScope {
1031 wire::OwnershipScope::VmOwnership(wire::VmOwnership {
1032 connection_id: connection_id.to_string(),
1033 session_id: session_id.to_string(),
1034 vm_id: vm_id.to_string(),
1035 })
1036}
1037
1038fn serialize_create_vm_config_for_sidecar(
1039 config: &AgentOsConfig,
1040) -> Result<vm_config::CreateVmConfig, ClientError> {
1041 let (root_filesystem, native_root) =
1042 serialize_root_filesystem_config_for_sidecar(&config.root_filesystem)?;
1043 Ok(vm_config::CreateVmConfig {
1044 cwd: None,
1045 env: BTreeMap::new(),
1046 root_filesystem,
1047 permissions: Some(permissions_policy_config(config)),
1048 limits: serialize_limits_config_for_sidecar(config.limits.as_ref())?,
1049 dns: None,
1050 native_root,
1051 listen: None,
1052 loopback_exempt_ports: config.loopback_exempt_ports.clone(),
1053 js_runtime: config.allowed_node_builtins.as_ref().map(|allowed| {
1059 vm_config::JsRuntimeConfig {
1060 platform: vm_config::JsRuntimePlatform::default(),
1061 module_resolution: vm_config::JsModuleResolution::default(),
1062 allowed_builtins: Some(allowed.clone()),
1063 high_resolution_time: None,
1064 }
1065 }),
1066 bootstrap_commands: Some(vec![
1067 String::from("node"),
1068 String::from("npm"),
1069 String::from("npx"),
1070 String::from("python"),
1071 String::from("python3"),
1072 ]),
1073 })
1074}
1075
1076fn serialize_root_filesystem_config_for_sidecar(
1077 config: &RootFilesystemConfig,
1078) -> Result<
1079 (
1080 vm_config::RootFilesystemConfig,
1081 Option<vm_config::NativeRootFilesystemConfig>,
1082 ),
1083 ClientError,
1084> {
1085 let mode = match config.mode.unwrap_or(ConfigRootFilesystemMode::Ephemeral) {
1086 ConfigRootFilesystemMode::Ephemeral => vm_config::RootFilesystemMode::Ephemeral,
1087 ConfigRootFilesystemMode::ReadOnly => vm_config::RootFilesystemMode::ReadOnly,
1088 };
1089 match config.kind {
1090 RootFilesystemKind::Overlay => {
1091 if config.native_plugin.is_some() {
1092 return Err(ClientError::Sidecar(
1093 "rootFilesystem.nativePlugin requires type \"native\"".to_string(),
1094 ));
1095 }
1096 let lowers = config
1097 .lowers
1098 .iter()
1099 .map(serialize_root_lower_config_for_sidecar)
1100 .collect::<Result<Vec<_>, _>>()?;
1101 Ok((
1102 vm_config::RootFilesystemConfig {
1103 mode,
1104 disable_default_base_layer: config.disable_default_base_layer,
1105 lowers,
1106 bootstrap_entries: Vec::new(),
1107 },
1108 None,
1109 ))
1110 }
1111 RootFilesystemKind::Native => {
1112 if !config.lowers.is_empty() {
1113 return Err(ClientError::Sidecar(
1114 "native root filesystems do not support rootFilesystem.lowers".to_string(),
1115 ));
1116 }
1117 let plugin = config.native_plugin.as_ref().ok_or_else(|| {
1118 ClientError::Sidecar(
1119 "rootFilesystem.nativePlugin is required for type \"native\"".to_string(),
1120 )
1121 })?;
1122 Ok((
1123 vm_config::RootFilesystemConfig {
1124 mode,
1125 disable_default_base_layer: config.disable_default_base_layer,
1126 lowers: Vec::new(),
1127 bootstrap_entries: Vec::new(),
1128 },
1129 Some(vm_config::NativeRootFilesystemConfig {
1130 plugin: vm_config::MountPluginDescriptor {
1131 id: plugin.id.clone(),
1132 config: plugin
1133 .config
1134 .clone()
1135 .unwrap_or_else(|| serde_json::Value::Object(serde_json::Map::new())),
1136 },
1137 read_only: config.mode == Some(ConfigRootFilesystemMode::ReadOnly),
1138 }),
1139 ))
1140 }
1141 }
1142}
1143
1144fn serialize_root_lower_config_for_sidecar(
1145 lower: &RootLowerInput,
1146) -> Result<vm_config::RootFilesystemLowerDescriptor, ClientError> {
1147 match lower {
1148 RootLowerInput::BundledBaseFilesystem => {
1149 Ok(vm_config::RootFilesystemLowerDescriptor::BundledBaseFilesystem)
1150 }
1151 RootLowerInput::SnapshotExport(snapshot) => {
1152 let entries = snapshot
1153 .source
1154 .filesystem
1155 .entries
1156 .iter()
1157 .map(serialize_filesystem_entry_config_for_sidecar)
1158 .collect::<Result<Vec<_>, _>>()?;
1159 Ok(vm_config::RootFilesystemLowerDescriptor::Snapshot { entries })
1160 }
1161 }
1162}
1163
1164fn serialize_filesystem_entry_config_for_sidecar(
1165 entry: &crate::fs::FilesystemEntry,
1166) -> Result<vm_config::RootFilesystemEntry, ClientError> {
1167 let mode = u32::from_str_radix(entry.mode.trim_start_matches("0o"), 8).map_err(|error| {
1168 ClientError::Sidecar(format!(
1169 "invalid root filesystem mode {} for {}: {error}",
1170 entry.mode, entry.path
1171 ))
1172 })?;
1173 let kind = match entry.entry_type {
1174 crate::fs::DirEntryType::File => vm_config::RootFilesystemEntryKind::File,
1175 crate::fs::DirEntryType::Directory => vm_config::RootFilesystemEntryKind::Directory,
1176 crate::fs::DirEntryType::Symlink => vm_config::RootFilesystemEntryKind::Symlink,
1177 };
1178 let encoding = entry.encoding.map(|encoding| match encoding {
1179 crate::fs::FilesystemEntryEncoding::Utf8 => vm_config::RootFilesystemEntryEncoding::Utf8,
1180 crate::fs::FilesystemEntryEncoding::Base64 => {
1181 vm_config::RootFilesystemEntryEncoding::Base64
1182 }
1183 });
1184
1185 Ok(vm_config::RootFilesystemEntry {
1186 path: entry.path.clone(),
1187 kind,
1188 mode: Some(mode),
1189 uid: Some(entry.uid),
1190 gid: Some(entry.gid),
1191 content: entry.content.clone(),
1192 encoding,
1193 target: entry.target.clone(),
1194 executable: entry.entry_type == crate::fs::DirEntryType::File && (mode & 0o111) != 0,
1195 })
1196}
1197
1198fn serialize_limits_config_for_sidecar(
1199 limits: Option<&AgentOsLimits>,
1200) -> Result<Option<vm_config::VmLimitsConfig>, ClientError> {
1201 let Some(limits) = limits else {
1202 return Ok(None);
1203 };
1204 let value = serde_json::to_value(limits).map_err(|error| {
1205 ClientError::Sidecar(format!("failed to serialize VM limits config: {error}"))
1206 })?;
1207 serde_json::from_value(value).map(Some).map_err(|error| {
1208 ClientError::Sidecar(format!("failed to encode VM limits config: {error}"))
1209 })
1210}
1211
1212const DEFAULT_EGRESS_HOSTS: &[&str] = &[
1219 "api.anthropic.com",
1220 "api.openai.com",
1221 "generativelanguage.googleapis.com",
1222 "openrouter.ai",
1223];
1224
1225fn default_egress_patterns() -> Vec<String> {
1229 DEFAULT_EGRESS_HOSTS
1230 .iter()
1231 .flat_map(|host| [format!("dns://{host}"), format!("tcp://{host}:*")])
1232 .collect()
1233}
1234
1235fn default_network_egress_scope_config() -> vm_config::PatternPermissionScope {
1237 vm_config::PatternPermissionScope::Rules(vm_config::PatternPermissionRuleSet {
1238 default: Some(vm_config::PermissionMode::Deny),
1239 rules: vec![vm_config::PatternPermissionRule {
1240 mode: vm_config::PermissionMode::Allow,
1241 operations: vec!["*".to_string()],
1242 patterns: default_egress_patterns(),
1243 }],
1244 })
1245}
1246
1247fn default_network_egress_scope() -> wire::PatternPermissionScope {
1249 wire::PatternPermissionScope::PatternPermissionRuleSet(wire::PatternPermissionRuleSet {
1250 default: Some(wire::PermissionMode::Deny),
1251 rules: vec![wire::PatternPermissionRule {
1252 mode: wire::PermissionMode::Allow,
1253 operations: vec!["*".to_string()],
1254 patterns: default_egress_patterns(),
1255 }],
1256 })
1257}
1258
1259fn permissions_policy_config(config: &AgentOsConfig) -> vm_config::PermissionsPolicy {
1260 let Some(permissions) = config.permissions.as_ref() else {
1261 return default_permissions_policy_config();
1262 };
1263
1264 vm_config::PermissionsPolicy {
1265 fs: Some(
1266 permissions
1267 .fs
1268 .as_ref()
1269 .map(serialize_fs_permissions_config)
1270 .unwrap_or(vm_config::FsPermissionScope::Mode(
1271 vm_config::PermissionMode::Allow,
1272 )),
1273 ),
1274 network: Some(
1275 permissions
1276 .network
1277 .as_ref()
1278 .map(serialize_pattern_permissions_config)
1279 .unwrap_or_else(default_network_egress_scope_config),
1280 ),
1281 child_process: Some(
1282 permissions
1283 .child_process
1284 .as_ref()
1285 .map(serialize_pattern_permissions_config)
1286 .unwrap_or(vm_config::PatternPermissionScope::Mode(
1287 vm_config::PermissionMode::Allow,
1288 )),
1289 ),
1290 process: Some(
1291 permissions
1292 .process
1293 .as_ref()
1294 .map(serialize_pattern_permissions_config)
1295 .unwrap_or(vm_config::PatternPermissionScope::Mode(
1296 vm_config::PermissionMode::Allow,
1297 )),
1298 ),
1299 env: Some(
1300 permissions
1301 .env
1302 .as_ref()
1303 .map(serialize_pattern_permissions_config)
1304 .unwrap_or(vm_config::PatternPermissionScope::Mode(
1305 vm_config::PermissionMode::Allow,
1306 )),
1307 ),
1308 binding: Some(
1309 permissions
1310 .binding
1311 .as_ref()
1312 .map(serialize_pattern_permissions_config)
1313 .unwrap_or(vm_config::PatternPermissionScope::Mode(
1314 vm_config::PermissionMode::Allow,
1315 )),
1316 ),
1317 }
1318}
1319
1320fn default_permissions_policy_config() -> vm_config::PermissionsPolicy {
1325 vm_config::PermissionsPolicy {
1326 fs: Some(vm_config::FsPermissionScope::Mode(
1327 vm_config::PermissionMode::Allow,
1328 )),
1329 network: Some(default_network_egress_scope_config()),
1330 child_process: Some(vm_config::PatternPermissionScope::Mode(
1331 vm_config::PermissionMode::Allow,
1332 )),
1333 process: Some(vm_config::PatternPermissionScope::Mode(
1334 vm_config::PermissionMode::Allow,
1335 )),
1336 env: Some(vm_config::PatternPermissionScope::Mode(
1337 vm_config::PermissionMode::Allow,
1338 )),
1339 binding: Some(vm_config::PatternPermissionScope::Mode(
1340 vm_config::PermissionMode::Allow,
1341 )),
1342 }
1343}
1344
1345fn serialize_fs_permissions_config(
1346 permissions: &crate::config::FsPermissions,
1347) -> vm_config::FsPermissionScope {
1348 match permissions {
1349 crate::config::FsPermissions::Mode(mode) => {
1350 vm_config::FsPermissionScope::Mode(serialize_permission_mode_config(*mode))
1351 }
1352 crate::config::FsPermissions::Rules(rules) => {
1353 vm_config::FsPermissionScope::Rules(vm_config::FsPermissionRuleSet {
1354 default: rules.default.map(serialize_permission_mode_config),
1355 rules: rules
1356 .rules
1357 .iter()
1358 .map(|rule| vm_config::FsPermissionRule {
1359 mode: serialize_permission_mode_config(rule.mode),
1360 operations: operation_wildcard_if_omitted(&rule.operations),
1361 paths: resource_wildcard_if_omitted(&rule.paths),
1362 })
1363 .collect(),
1364 })
1365 }
1366 }
1367}
1368
1369fn serialize_pattern_permissions_config(
1370 permissions: &crate::config::PatternPermissions,
1371) -> vm_config::PatternPermissionScope {
1372 match permissions {
1373 crate::config::PatternPermissions::Mode(mode) => {
1374 vm_config::PatternPermissionScope::Mode(serialize_permission_mode_config(*mode))
1375 }
1376 crate::config::PatternPermissions::Rules(rules) => {
1377 vm_config::PatternPermissionScope::Rules(vm_config::PatternPermissionRuleSet {
1378 default: rules.default.map(serialize_permission_mode_config),
1379 rules: rules
1380 .rules
1381 .iter()
1382 .map(|rule| vm_config::PatternPermissionRule {
1383 mode: serialize_permission_mode_config(rule.mode),
1384 operations: operation_wildcard_if_omitted(&rule.operations),
1385 patterns: resource_wildcard_if_omitted(&rule.patterns),
1386 })
1387 .collect(),
1388 })
1389 }
1390 }
1391}
1392
1393fn serialize_permission_mode_config(
1394 mode: crate::config::PermissionMode,
1395) -> vm_config::PermissionMode {
1396 match mode {
1397 crate::config::PermissionMode::Allow => vm_config::PermissionMode::Allow,
1398 crate::config::PermissionMode::Deny => vm_config::PermissionMode::Deny,
1399 }
1400}
1401
1402async fn wait_for_vm_ready(
1404 events: &mut broadcast::Receiver<(wire::OwnershipScope, wire::EventPayload)>,
1405 vm_id: &str,
1406 timeout_ms: u64,
1407) -> Result<(), ClientError> {
1408 let wait = async {
1409 loop {
1410 match events.recv().await {
1411 Ok((ownership, payload)) => match payload {
1412 wire::EventPayload::VmLifecycleEvent(event) => {
1413 if matches!(event.state, wire::VmLifecycleState::Ready)
1414 && wire_ownership_vm_id(&ownership) == Some(vm_id)
1415 {
1416 return Ok(());
1417 }
1418 }
1419 wire::EventPayload::ProcessOutputEvent(_)
1420 | wire::EventPayload::ProcessExitedEvent(_)
1421 | wire::EventPayload::StructuredEvent(_)
1422 | wire::EventPayload::ExtEnvelope(_) => {}
1423 },
1424 Err(broadcast::error::RecvError::Lagged(_)) => {}
1425 Err(broadcast::error::RecvError::Closed) => {
1426 return Err(ClientError::Sidecar(
1427 "sidecar transport closed before the VM became ready".to_string(),
1428 ));
1429 }
1430 }
1431 }
1432 };
1433 tokio::time::timeout(Duration::from_millis(timeout_ms), wait)
1434 .await
1435 .map_err(|_| {
1436 ClientError::Sidecar("timed out waiting for the VM to become ready".to_string())
1437 })?
1438}
1439
1440static VM_TOOLS: OnceCell<SccHashMap<String, Arc<VmHostToolRegistry>>> = OnceCell::new();
1443
1444#[derive(Clone)]
1445struct VmHostToolRegistry {
1446 tool_kits: Vec<ToolKit>,
1447 tool_map: HashMap<String, HostTool>,
1448 permissions: Option<Permissions>,
1449}
1450
1451fn vm_tools() -> &'static SccHashMap<String, Arc<VmHostToolRegistry>> {
1452 VM_TOOLS.get_or_init(SccHashMap::new)
1453}
1454
1455static VM_PERMISSION_ROUTERS: OnceCell<SccHashMap<String, Weak<AgentOsInner>>> = OnceCell::new();
1459
1460fn vm_permission_routers() -> &'static SccHashMap<String, Weak<AgentOsInner>> {
1461 VM_PERMISSION_ROUTERS.get_or_init(SccHashMap::new)
1462}
1463
1464static SESSION_JS_BRIDGE_CALLBACKS: OnceCell<SccHashMap<String, SidecarJsBridgeCallback>> =
1469 OnceCell::new();
1470
1471fn session_js_bridge_callbacks() -> &'static SccHashMap<String, SidecarJsBridgeCallback> {
1472 SESSION_JS_BRIDGE_CALLBACKS.get_or_init(SccHashMap::new)
1473}
1474
1475fn sidecar_session_key(connection_id: &str, session_id: &str) -> String {
1476 format!("{connection_id}\0{session_id}")
1477}
1478
1479fn wire_ownership_session_key(ownership: &wire::OwnershipScope) -> Option<String> {
1480 match ownership {
1481 wire::OwnershipScope::SessionOwnership(ownership) => Some(sidecar_session_key(
1482 &ownership.connection_id,
1483 &ownership.session_id,
1484 )),
1485 wire::OwnershipScope::VmOwnership(ownership) => Some(sidecar_session_key(
1486 &ownership.connection_id,
1487 &ownership.session_id,
1488 )),
1489 wire::OwnershipScope::ConnectionOwnership(_) => None,
1490 }
1491}
1492
1493fn js_bridge_call_callback() -> WireSidecarCallback {
1494 Arc::new(|payload, ownership| {
1495 Box::pin(async move {
1496 let request = match payload {
1497 wire::SidecarRequestPayload::JsBridgeCallRequest(request) => request,
1498 wire::SidecarRequestPayload::HostCallbackRequest(_) => {
1499 return Ok(wire::SidecarResponsePayload::JsBridgeResultResponse(
1500 wire::JsBridgeResultResponse {
1501 call_id: "unknown".to_string(),
1502 result: None,
1503 error: Some(
1504 "js-bridge callback received a host callback request".to_string(),
1505 ),
1506 },
1507 ));
1508 }
1509 wire::SidecarRequestPayload::ExtEnvelope(_) => {
1510 return Ok(wire::SidecarResponsePayload::JsBridgeResultResponse(
1511 wire::JsBridgeResultResponse {
1512 call_id: "unknown".to_string(),
1513 result: None,
1514 error: Some(
1515 "js-bridge callback received an extension request".to_string(),
1516 ),
1517 },
1518 ));
1519 }
1520 };
1521 Ok(wire::SidecarResponsePayload::JsBridgeResultResponse(
1522 run_js_bridge_callback(&ownership, request).await,
1523 ))
1524 })
1525 })
1526}
1527
1528async fn run_js_bridge_callback(
1529 ownership: &wire::OwnershipScope,
1530 request: wire::JsBridgeCallRequest,
1531) -> wire::JsBridgeResultResponse {
1532 let call_id = request.call_id;
1533 let args = match serde_json::from_str::<Value>(&request.args) {
1534 Ok(args) => args,
1535 Err(error) => {
1536 return wire::JsBridgeResultResponse {
1537 call_id,
1538 result: None,
1539 error: Some(format!("Invalid js_bridge args: {error}")),
1540 };
1541 }
1542 };
1543 let callback = wire_ownership_session_key(ownership)
1544 .and_then(|key| session_js_bridge_callbacks().read(&key, |_, callback| callback.clone()));
1545 let Some(callback) = callback else {
1546 return wire::JsBridgeResultResponse {
1547 call_id,
1548 result: None,
1549 error: Some("No js_bridge callback registered for sidecar session".to_string()),
1550 };
1551 };
1552
1553 let call = SidecarJsBridgeCall {
1554 call_id: call_id.clone(),
1555 mount_id: request.mount_id,
1556 operation: request.operation,
1557 args,
1558 };
1559 match callback(call).await {
1560 Ok(result) => match result {
1561 Some(value) => match serde_json::to_string(&value) {
1562 Ok(result) => wire::JsBridgeResultResponse {
1563 call_id,
1564 result: Some(result),
1565 error: None,
1566 },
1567 Err(error) => wire::JsBridgeResultResponse {
1568 call_id,
1569 result: None,
1570 error: Some(format!("Invalid js_bridge result: {error}")),
1571 },
1572 },
1573 None => wire::JsBridgeResultResponse {
1574 call_id,
1575 result: None,
1576 error: None,
1577 },
1578 },
1579 Err(error) => wire::JsBridgeResultResponse {
1580 call_id,
1581 result: None,
1582 error: Some(error),
1583 },
1584 }
1585}
1586
1587fn permission_request_callback() -> WireSidecarCallback {
1590 Arc::new(|payload, ownership| {
1591 Box::pin(async move {
1592 match payload {
1593 wire::SidecarRequestPayload::ExtEnvelope(envelope) => {
1594 handle_acp_ext_callback(envelope, &ownership)
1595 .await
1596 .map_err(|error| TransportError::Sidecar(error.to_string()))
1597 }
1598 wire::SidecarRequestPayload::HostCallbackRequest(_)
1599 | wire::SidecarRequestPayload::JsBridgeCallRequest(_) => Ok(
1600 wire::SidecarResponsePayload::ExtEnvelope(wire::ExtEnvelope {
1601 namespace: ACP_EXTENSION_NAMESPACE.to_string(),
1602 payload: b"permission callback received a non-extension request".to_vec(),
1603 }),
1604 ),
1605 }
1606 })
1607 })
1608}
1609
1610async fn handle_acp_ext_callback(
1611 envelope: wire::ExtEnvelope,
1612 ownership: &wire::OwnershipScope,
1613) -> Result<wire::SidecarResponsePayload, ClientError> {
1614 if envelope.namespace != ACP_EXTENSION_NAMESPACE {
1615 return Ok(wire::SidecarResponsePayload::ExtEnvelope(
1616 wire::ExtEnvelope {
1617 namespace: envelope.namespace,
1618 payload: b"unknown extension namespace".to_vec(),
1619 },
1620 ));
1621 }
1622 let callback: AcpCallback = serde_bare::from_slice(&envelope.payload)
1623 .map_err(|error| ClientError::Sidecar(format!("invalid ACP callback: {error}")))?;
1624 let response = match callback {
1625 AcpCallback::AcpPermissionCallback(callback) => {
1626 let params =
1627 serde_json::from_str(&callback.params).unwrap_or_else(|_| serde_json::json!({}));
1628 let result = route_permission_request(
1629 ownership,
1630 PermissionRouteRequest {
1631 session_id: callback.session_id,
1632 permission_id: callback.permission_id.clone(),
1633 params,
1634 },
1635 )
1636 .await;
1637 let reply = result.reply.unwrap_or_else(|| String::from("reject"));
1638 AcpCallbackResponse::AcpPermissionCallbackResponse(AcpPermissionCallbackResponse {
1639 permission_id: callback.permission_id,
1640 reply,
1641 })
1642 }
1643 AcpCallback::AcpHostRequestCallback(callback) => {
1644 let response = dispatch_acp_host_request(ownership, &callback.request).await;
1645 AcpCallbackResponse::AcpHostRequestCallbackResponse(AcpHostRequestCallbackResponse {
1646 response: Some(response),
1647 })
1648 }
1649 };
1650 let payload = serde_bare::to_vec(&response).map_err(|error| {
1651 ClientError::Sidecar(format!("failed to encode ACP callback response: {error}"))
1652 })?;
1653 Ok(wire::SidecarResponsePayload::ExtEnvelope(
1654 wire::ExtEnvelope {
1655 namespace: ACP_EXTENSION_NAMESPACE.to_string(),
1656 payload,
1657 },
1658 ))
1659}
1660
1661async fn route_permission_request(
1662 ownership: &wire::OwnershipScope,
1663 request: PermissionRouteRequest,
1664) -> PermissionRouteResult {
1665 let vm_id = wire_ownership_vm_id(ownership).unwrap_or("");
1666 let inner = vm_permission_routers()
1667 .read(vm_id, |_, weak| weak.clone())
1668 .and_then(|weak| weak.upgrade());
1669 let Some(inner) = inner else {
1670 return PermissionRouteResult { reply: None };
1671 };
1672 let client = AgentOs { inner };
1673 client.deliver_sidecar_permission_request(request).await
1674}
1675
1676const ACP_TERMINAL_DEFAULT_OUTPUT_BYTE_LIMIT: usize = 1_048_576;
1683
1684struct AcpDispatchError {
1686 code: i64,
1687 message: String,
1688 data: Option<Value>,
1689}
1690
1691impl AcpDispatchError {
1692 fn new(code: i64, message: impl Into<String>) -> Self {
1693 Self {
1694 code,
1695 message: message.into(),
1696 data: None,
1697 }
1698 }
1699
1700 fn with_data(code: i64, message: impl Into<String>, data: Value) -> Self {
1701 Self {
1702 code,
1703 message: message.into(),
1704 data: Some(data),
1705 }
1706 }
1707}
1708
1709impl From<ClientError> for AcpDispatchError {
1710 fn from(error: ClientError) -> Self {
1711 match error {
1712 ClientError::Kernel { code, message } => {
1715 AcpDispatchError::with_data(-32603, message, serde_json::json!({ "code": code }))
1716 }
1717 other => AcpDispatchError::new(-32603, other.to_string()),
1718 }
1719 }
1720}
1721
1722impl From<anyhow::Error> for AcpDispatchError {
1723 fn from(error: anyhow::Error) -> Self {
1724 match error.downcast::<ClientError>() {
1727 Ok(client_error) => client_error.into(),
1728 Err(error) => AcpDispatchError::new(-32603, error.to_string()),
1729 }
1730 }
1731}
1732
1733async fn dispatch_acp_host_request(ownership: &wire::OwnershipScope, request: &str) -> String {
1737 let parsed = serde_json::from_str::<Value>(request);
1738 let (id, method, params_value) = match parsed {
1739 Ok(value) => {
1740 let id = value.get("id").cloned().unwrap_or(Value::Null);
1741 let method = value
1742 .get("method")
1743 .and_then(Value::as_str)
1744 .map(str::to_string);
1745 (id, method, value.get("params").cloned())
1746 }
1747 Err(error) => {
1748 return acp_error_response(Value::Null, -32700, &format!("Parse error: {error}"), None);
1749 }
1750 };
1751
1752 let Some(method) = method else {
1753 return acp_error_response(id, -32600, "Invalid Request: missing method", None);
1754 };
1755
1756 match handle_acp_host_request(ownership, &method, params_value).await {
1757 Ok(result) => serde_json::to_string(&serde_json::json!({
1758 "jsonrpc": "2.0",
1759 "id": id,
1760 "result": result,
1761 }))
1762 .unwrap_or_else(|error| acp_error_response(Value::Null, -32603, &error.to_string(), None)),
1763 Err(error) => acp_error_response(id, error.code, &error.message, error.data),
1764 }
1765}
1766
1767fn acp_error_response(id: Value, code: i64, message: &str, data: Option<Value>) -> String {
1768 let mut error = serde_json::json!({
1769 "code": code,
1770 "message": message,
1771 });
1772 if let Some(data) = data {
1773 if let Some(map) = error.as_object_mut() {
1774 map.insert("data".to_string(), data);
1775 }
1776 }
1777 serde_json::to_string(&serde_json::json!({
1778 "jsonrpc": "2.0",
1779 "id": id,
1780 "error": error,
1781 }))
1782 .unwrap_or_else(|_| {
1783 String::from(r#"{"jsonrpc":"2.0","id":null,"error":{"code":-32603,"message":"failed to encode error response"}}"#)
1784 })
1785}
1786
1787fn resolve_acp_agent(ownership: &wire::OwnershipScope) -> Result<AgentOs, AcpDispatchError> {
1789 let vm_id = wire_ownership_vm_id(ownership).unwrap_or("");
1790 let inner = vm_permission_routers()
1791 .read(vm_id, |_, weak| weak.clone())
1792 .and_then(|weak| weak.upgrade());
1793 inner
1794 .map(|inner| AgentOs { inner })
1795 .ok_or_else(|| AcpDispatchError::new(-32603, "VM is no longer available"))
1796}
1797
1798async fn handle_acp_host_request(
1801 ownership: &wire::OwnershipScope,
1802 method: &str,
1803 params_value: Option<Value>,
1804) -> Result<Value, AcpDispatchError> {
1805 let params = acp_params(method, params_value)?;
1806 match method {
1807 crate::session::ACP_PERMISSION_METHOD => {
1808 handle_acp_permission_request(ownership, method, ¶ms).await
1809 }
1810 "fs/read" | "fs/read_text_file" => {
1811 let agent = resolve_acp_agent(ownership)?;
1812 handle_acp_read_file(&agent, ¶ms).await
1813 }
1814 "fs/write" | "fs/write_text_file" => {
1815 let agent = resolve_acp_agent(ownership)?;
1816 handle_acp_write_file(&agent, ¶ms).await
1817 }
1818 "fs/readDir" | "fs/read_dir" => {
1819 let agent = resolve_acp_agent(ownership)?;
1820 handle_acp_read_dir(&agent, ¶ms).await
1821 }
1822 "terminal/create" => {
1823 let agent = resolve_acp_agent(ownership)?;
1824 handle_acp_create_terminal(&agent, ¶ms)
1825 }
1826 "terminal/write" => {
1827 let agent = resolve_acp_agent(ownership)?;
1828 handle_acp_write_terminal(&agent, ¶ms)
1829 }
1830 "terminal/output" | "terminal/read" => {
1831 let agent = resolve_acp_agent(ownership)?;
1832 handle_acp_read_terminal(&agent, ¶ms)
1833 }
1834 "terminal/wait_for_exit" | "terminal/waitForExit" => {
1835 let agent = resolve_acp_agent(ownership)?;
1836 handle_acp_wait_for_terminal_exit(&agent, ¶ms).await
1837 }
1838 "terminal/kill" => {
1839 let agent = resolve_acp_agent(ownership)?;
1840 handle_acp_kill_terminal(&agent, ¶ms)
1841 }
1842 "terminal/release" | "terminal/close" => {
1843 let agent = resolve_acp_agent(ownership)?;
1844 handle_acp_release_terminal(&agent, ¶ms)
1845 }
1846 "terminal/resize" => {
1847 let agent = resolve_acp_agent(ownership)?;
1848 handle_acp_resize_terminal(&agent, ¶ms)
1849 }
1850 other => Err(AcpDispatchError::with_data(
1851 -32601,
1852 format!("Method not found: {other}"),
1853 serde_json::json!({ "method": other }),
1854 )),
1855 }
1856}
1857
1858fn acp_params(
1861 method: &str,
1862 params_value: Option<Value>,
1863) -> Result<Map<String, Value>, AcpDispatchError> {
1864 match params_value {
1865 None | Some(Value::Null) => Ok(Map::new()),
1866 Some(Value::Object(map)) => Ok(map),
1867 Some(_) => Err(AcpDispatchError::new(
1868 -32602,
1869 format!("{method} requires object params"),
1870 )),
1871 }
1872}
1873
1874fn require_acp_string(
1875 params: &Map<String, Value>,
1876 name: &str,
1877 method: &str,
1878) -> Result<String, AcpDispatchError> {
1879 match params.get(name).and_then(Value::as_str) {
1880 Some(value) => Ok(value.to_string()),
1881 None => Err(AcpDispatchError::new(
1882 -32602,
1883 format!("{method} requires a string {name}"),
1884 )),
1885 }
1886}
1887
1888fn optional_acp_string(
1889 params: &Map<String, Value>,
1890 name: &str,
1891 method: &str,
1892) -> Result<Option<String>, AcpDispatchError> {
1893 match params.get(name) {
1894 None | Some(Value::Null) => Ok(None),
1895 Some(Value::String(value)) => Ok(Some(value.clone())),
1896 Some(_) => Err(AcpDispatchError::new(
1897 -32602,
1898 format!("{method} requires {name} to be a string when provided"),
1899 )),
1900 }
1901}
1902
1903fn optional_acp_number(
1904 params: &Map<String, Value>,
1905 name: &str,
1906 method: &str,
1907) -> Result<Option<f64>, AcpDispatchError> {
1908 match params.get(name) {
1909 None | Some(Value::Null) => Ok(None),
1910 Some(value) => match value.as_f64() {
1911 Some(number) if number.is_finite() => Ok(Some(number)),
1912 _ => Err(AcpDispatchError::new(
1913 -32602,
1914 format!("{method} requires {name} to be a number when provided"),
1915 )),
1916 },
1917 }
1918}
1919
1920fn optional_acp_string_array(
1921 params: &Map<String, Value>,
1922 name: &str,
1923 method: &str,
1924) -> Result<Option<Vec<String>>, AcpDispatchError> {
1925 match params.get(name) {
1926 None | Some(Value::Null) => Ok(None),
1927 Some(Value::Array(items)) => {
1928 let mut out = Vec::with_capacity(items.len());
1929 for item in items {
1930 match item.as_str() {
1931 Some(value) => out.push(value.to_string()),
1932 None => {
1933 return Err(AcpDispatchError::new(
1934 -32602,
1935 format!(
1936 "{method} requires {name} to be an array of strings when provided"
1937 ),
1938 ))
1939 }
1940 }
1941 }
1942 Ok(Some(out))
1943 }
1944 Some(_) => Err(AcpDispatchError::new(
1945 -32602,
1946 format!("{method} requires {name} to be an array of strings when provided"),
1947 )),
1948 }
1949}
1950
1951fn optional_acp_env(
1954 params: &Map<String, Value>,
1955 name: &str,
1956 method: &str,
1957) -> Result<Option<BTreeMap<String, String>>, AcpDispatchError> {
1958 match params.get(name) {
1959 None | Some(Value::Null) => Ok(None),
1960 Some(Value::Array(items)) => {
1961 let mut env = BTreeMap::new();
1962 for entry in items {
1963 let Some(record) = entry.as_object() else {
1964 return Err(AcpDispatchError::new(
1965 -32602,
1966 format!("{method} requires {name} entries to be {{ name, value }} objects"),
1967 ));
1968 };
1969 match (
1970 record.get("name").and_then(Value::as_str),
1971 record.get("value").and_then(Value::as_str),
1972 ) {
1973 (Some(key), Some(value)) => {
1974 env.insert(key.to_string(), value.to_string());
1975 }
1976 _ => {
1977 return Err(AcpDispatchError::new(
1978 -32602,
1979 format!(
1980 "{method} requires {name} entries to be {{ name, value }} objects"
1981 ),
1982 ))
1983 }
1984 }
1985 }
1986 Ok(Some(env))
1987 }
1988 Some(Value::Object(map)) => {
1989 let mut env = BTreeMap::new();
1990 for (key, value) in map {
1991 match value.as_str() {
1992 Some(value) => {
1993 env.insert(key.clone(), value.to_string());
1994 }
1995 None => {
1996 return Err(AcpDispatchError::new(
1997 -32602,
1998 format!("{method} requires {name} values to be strings"),
1999 ))
2000 }
2001 }
2002 }
2003 Ok(Some(env))
2004 }
2005 Some(_) => Err(AcpDispatchError::new(
2006 -32602,
2007 format!("{method} requires {name} to be an object or name/value array"),
2008 )),
2009 }
2010}
2011
2012async fn handle_acp_read_file(
2015 agent: &AgentOs,
2016 params: &Map<String, Value>,
2017) -> Result<Value, AcpDispatchError> {
2018 let method = "fs/read";
2019 let path = require_acp_string(params, "path", method)?;
2020 let line = optional_acp_number(params, "line", method)?;
2021 let limit = optional_acp_number(params, "limit", method)?;
2022 let encoding = optional_acp_string(params, "encoding", method)?;
2023 let bytes = agent.read_file(&path).await?;
2024 if encoding.as_deref() == Some("base64") {
2025 use base64::engine::general_purpose::STANDARD as BASE64;
2026 use base64::Engine as _;
2027 return Ok(serde_json::json!({ "content": BASE64.encode(&bytes) }));
2028 }
2029 let text = String::from_utf8_lossy(&bytes).into_owned();
2030 if line.is_none() && limit.is_none() {
2031 return Ok(serde_json::json!({ "content": text }));
2032 }
2033 let start_line = line.map(|n| n.trunc() as i64).unwrap_or(1).max(1);
2034 let lines: Vec<&str> = text.split('\n').collect();
2035 let start_index = (start_line - 1).max(0) as usize;
2036 let selected: Vec<&str> = match limit {
2037 None => lines.into_iter().skip(start_index).collect(),
2038 Some(limit) => {
2039 let limit = limit.trunc().max(0.0) as usize;
2040 lines.into_iter().skip(start_index).take(limit).collect()
2041 }
2042 };
2043 Ok(serde_json::json!({ "content": selected.join("\n") }))
2044}
2045
2046async fn handle_acp_write_file(
2047 agent: &AgentOs,
2048 params: &Map<String, Value>,
2049) -> Result<Value, AcpDispatchError> {
2050 let method = "fs/write";
2051 let path = require_acp_string(params, "path", method)?;
2052 let content = require_acp_string(params, "content", method)?;
2053 let encoding = optional_acp_string(params, "encoding", method)?;
2054 if encoding.as_deref() == Some("base64") {
2055 use base64::engine::general_purpose::STANDARD as BASE64;
2056 use base64::Engine as _;
2057 let decoded = BASE64.decode(content.as_bytes()).map_err(|error| {
2058 AcpDispatchError::new(
2059 -32602,
2060 format!("{method} content is not valid base64: {error}"),
2061 )
2062 })?;
2063 agent.write_file(&path, decoded).await?;
2064 } else {
2065 agent.write_file(&path, content).await?;
2066 }
2067 Ok(Value::Null)
2068}
2069
2070async fn handle_acp_read_dir(
2071 agent: &AgentOs,
2072 params: &Map<String, Value>,
2073) -> Result<Value, AcpDispatchError> {
2074 let method = "fs/readDir";
2075 let path = require_acp_string(params, "path", method)?;
2076 let entries = agent.acp_read_dir_with_types(&path).await?;
2077 let mapped: Vec<Value> = entries
2078 .into_iter()
2079 .map(|entry| {
2080 let child_path = if path == "/" {
2081 format!("/{}", entry.name)
2082 } else {
2083 format!("{path}/{}", entry.name)
2084 };
2085 let entry_type = if entry.is_symbolic_link {
2086 "symlink"
2087 } else if entry.is_directory {
2088 "directory"
2089 } else {
2090 "file"
2091 };
2092 serde_json::json!({
2093 "name": entry.name,
2094 "path": child_path,
2095 "type": entry_type,
2096 })
2097 })
2098 .collect();
2099 Ok(serde_json::json!({ "entries": mapped }))
2100}
2101
2102async fn handle_acp_permission_request(
2105 ownership: &wire::OwnershipScope,
2106 method: &str,
2107 params: &Map<String, Value>,
2108) -> Result<Value, AcpDispatchError> {
2109 let session_id = require_acp_string(params, "sessionId", method)?;
2110
2111 let result = route_permission_request(
2112 ownership,
2113 PermissionRouteRequest {
2114 session_id: session_id.clone(),
2115 permission_id: format!("acp-permission-{}", uuid::Uuid::new_v4()),
2118 params: Value::Object(params.clone()),
2119 },
2120 )
2121 .await;
2122
2123 let reply = match result.reply.as_deref() {
2125 Some("always") => PermissionDecision::Always,
2126 Some("once") => PermissionDecision::Once,
2127 _ => PermissionDecision::Reject,
2128 };
2129 Ok(build_acp_permission_result(reply, params))
2130}
2131
2132#[derive(Clone, Copy)]
2133enum PermissionDecision {
2134 Always,
2135 Once,
2136 Reject,
2137}
2138
2139fn normalize_acp_permission_option_id(
2142 options: Option<&Vec<Value>>,
2143 decision: PermissionDecision,
2144) -> String {
2145 let (option_ids, kinds, fallback): (&[&str], &[&str], &str) = match decision {
2146 PermissionDecision::Always => (
2147 &["always", "allow_always"],
2148 &["allow_always"],
2149 "allow_always",
2150 ),
2151 PermissionDecision::Once => (&["once", "allow_once"], &["allow_once"], "allow_once"),
2152 PermissionDecision::Reject => (&["reject", "reject_once"], &["reject_once"], "reject_once"),
2153 };
2154 if let Some(options) = options {
2155 for option in options {
2156 let Some(record) = option.as_object() else {
2157 continue;
2158 };
2159 let option_id = record.get("optionId").and_then(Value::as_str);
2160 let kind = record.get("kind").and_then(Value::as_str);
2161 let matches = option_id.is_some_and(|id| option_ids.contains(&id))
2162 || kind.is_some_and(|k| kinds.contains(&k));
2163 if matches {
2164 if let Some(id) = option_id {
2165 return id.to_string();
2166 }
2167 }
2168 }
2169 }
2170 fallback.to_string()
2171}
2172
2173fn build_acp_permission_result(decision: PermissionDecision, params: &Map<String, Value>) -> Value {
2175 let options = params.get("options").and_then(Value::as_array);
2176 let option_id = normalize_acp_permission_option_id(options, decision);
2177 serde_json::json!({
2178 "outcome": {
2179 "outcome": "selected",
2180 "optionId": option_id,
2181 }
2182 })
2183}
2184
2185fn require_acp_terminal_id(
2188 params: &Map<String, Value>,
2189 method: &str,
2190) -> Result<String, AcpDispatchError> {
2191 require_acp_string(params, "terminalId", method)
2192}
2193
2194fn handle_acp_create_terminal(
2195 agent: &AgentOs,
2196 params: &Map<String, Value>,
2197) -> Result<Value, AcpDispatchError> {
2198 let method = "terminal/create";
2199 let command = require_acp_string(params, "command", method)?;
2200 let args = optional_acp_string_array(params, "args", method)?;
2201 let env = optional_acp_env(params, "env", method)?;
2202 let cwd = optional_acp_string(params, "cwd", method)?;
2203 let cols = optional_acp_number(params, "cols", method)?;
2204 let rows = optional_acp_number(params, "rows", method)?;
2205 let output_byte_limit = optional_acp_number(params, "outputByteLimit", method)?
2206 .map(|n| n.trunc().max(0.0) as usize)
2207 .unwrap_or(ACP_TERMINAL_DEFAULT_OUTPUT_BYTE_LIMIT);
2208
2209 let counter = agent
2210 .inner()
2211 .host_acp_terminal_counter
2212 .fetch_add(1, Ordering::SeqCst)
2213 + 1;
2214 let terminal_id = format!("acp-terminal-{counter}");
2215
2216 let output = Arc::new(parking_lot::Mutex::new(HostAcpTerminalOutput {
2217 buffer: String::new(),
2218 truncated: false,
2219 output_byte_limit,
2220 }));
2221 let (exit_tx, exit_rx) = watch::channel::<Option<i32>>(None);
2222
2223 let mut shell_options = crate::shell::OpenShellOptions {
2226 command: Some(command),
2227 cwd,
2228 ..Default::default()
2229 };
2230 if let Some(args) = args {
2231 shell_options.args = args;
2232 }
2233 if let Some(env) = env {
2234 shell_options.env = env;
2235 }
2236 if let Some(cols) = cols {
2237 shell_options.cols = Some(cols.trunc() as u16);
2238 }
2239 if let Some(rows) = rows {
2240 shell_options.rows = Some(rows.trunc() as u16);
2241 }
2242 let buffer_sink = output.clone();
2245 let handle = agent
2246 .acp_open_terminal(shell_options, exit_tx, move |data: &[u8]| {
2247 append_acp_terminal_output(&buffer_sink, data);
2248 })
2249 .map_err(|error| AcpDispatchError::new(-32603, error.to_string()))?;
2250 let shell_id = handle.shell_id.clone();
2251
2252 let entry = HostAcpTerminal {
2253 shell_id,
2254 output,
2255 exit_rx,
2256 };
2257 if agent
2258 .inner()
2259 .host_acp_terminals
2260 .insert(terminal_id.clone(), entry)
2261 .is_err()
2262 {
2263 return Err(AcpDispatchError::new(
2264 -32603,
2265 format!("ACP terminal id collision: {terminal_id}"),
2266 ));
2267 }
2268
2269 Ok(serde_json::json!({ "terminalId": terminal_id }))
2270}
2271
2272fn append_acp_terminal_output(
2273 output: &Arc<parking_lot::Mutex<HostAcpTerminalOutput>>,
2274 data: &[u8],
2275) {
2276 let chunk = String::from_utf8_lossy(data);
2277 if chunk.is_empty() {
2278 return;
2279 }
2280 let mut state = output.lock();
2281 state.buffer.push_str(&chunk);
2282 let limit = state.output_byte_limit;
2283 if state.buffer.len() > limit {
2284 let overflow = state.buffer.len() - limit;
2287 let mut cut = overflow;
2288 while cut < state.buffer.len() && !state.buffer.is_char_boundary(cut) {
2289 cut += 1;
2290 }
2291 state.buffer = state.buffer.split_off(cut);
2292 state.truncated = true;
2293 }
2294}
2295
2296fn handle_acp_write_terminal(
2297 agent: &AgentOs,
2298 params: &Map<String, Value>,
2299) -> Result<Value, AcpDispatchError> {
2300 let method = "terminal/write";
2301 let terminal_id = require_acp_terminal_id(params, method)?;
2302 let shell_id = acp_terminal_shell_id(agent, &terminal_id)?;
2303 let data = require_acp_string(params, "data", method)?;
2304 let encoding = optional_acp_string(params, "encoding", method)?;
2305 let input = if encoding.as_deref() == Some("base64") {
2306 use base64::engine::general_purpose::STANDARD as BASE64;
2307 use base64::Engine as _;
2308 let decoded = BASE64.decode(data.as_bytes()).map_err(|error| {
2309 AcpDispatchError::new(
2310 -32602,
2311 format!("{method} data is not valid base64: {error}"),
2312 )
2313 })?;
2314 crate::process::StdinInput::Bytes(decoded)
2315 } else {
2316 crate::process::StdinInput::Text(data)
2317 };
2318 agent
2319 .write_shell(&shell_id, input)
2320 .map_err(|error| AcpDispatchError::new(-32603, error.to_string()))?;
2321 Ok(Value::Null)
2322}
2323
2324fn handle_acp_read_terminal(
2325 agent: &AgentOs,
2326 params: &Map<String, Value>,
2327) -> Result<Value, AcpDispatchError> {
2328 let method = "terminal/output";
2329 let terminal_id = require_acp_terminal_id(params, method)?;
2330 agent
2331 .inner()
2332 .host_acp_terminals
2333 .read(&terminal_id, |_, terminal| {
2334 let (output, truncated) = {
2335 let state = terminal.output.lock();
2336 (state.buffer.clone(), state.truncated)
2337 };
2338 let mut result = serde_json::json!({
2339 "output": output,
2340 "truncated": truncated,
2341 });
2342 if let Some(exit_code) = *terminal.exit_rx.borrow() {
2343 if let Some(map) = result.as_object_mut() {
2344 map.insert(
2345 "exitStatus".to_string(),
2346 serde_json::json!({ "exitCode": exit_code, "signal": Value::Null }),
2347 );
2348 }
2349 }
2350 result
2351 })
2352 .ok_or_else(|| {
2353 AcpDispatchError::new(-32602, format!("ACP terminal not found: {terminal_id}"))
2354 })
2355}
2356
2357async fn handle_acp_wait_for_terminal_exit(
2358 agent: &AgentOs,
2359 params: &Map<String, Value>,
2360) -> Result<Value, AcpDispatchError> {
2361 let method = "terminal/wait_for_exit";
2362 let terminal_id = require_acp_terminal_id(params, method)?;
2363 let mut exit_rx = agent
2364 .inner()
2365 .host_acp_terminals
2366 .read(&terminal_id, |_, terminal| terminal.exit_rx.clone())
2367 .ok_or_else(|| {
2368 AcpDispatchError::new(-32602, format!("ACP terminal not found: {terminal_id}"))
2369 })?;
2370 let exit_code = loop {
2371 if let Some(code) = *exit_rx.borrow() {
2372 break code;
2373 }
2374 if exit_rx.changed().await.is_err() {
2375 break exit_rx.borrow().unwrap_or(1);
2379 }
2380 };
2381 Ok(serde_json::json!({ "exitCode": exit_code, "signal": Value::Null }))
2382}
2383
2384fn handle_acp_kill_terminal(
2385 agent: &AgentOs,
2386 params: &Map<String, Value>,
2387) -> Result<Value, AcpDispatchError> {
2388 let method = "terminal/kill";
2389 let terminal_id = require_acp_terminal_id(params, method)?;
2390 let shell_id = acp_terminal_shell_id(agent, &terminal_id)?;
2391 agent
2396 .acp_kill_terminal_shell(&shell_id)
2397 .map_err(|error| AcpDispatchError::new(-32603, error.to_string()))?;
2398 Ok(Value::Null)
2399}
2400
2401fn handle_acp_release_terminal(
2402 agent: &AgentOs,
2403 params: &Map<String, Value>,
2404) -> Result<Value, AcpDispatchError> {
2405 let method = "terminal/release";
2406 let terminal_id = require_acp_terminal_id(params, method)?;
2407 let Some((_, terminal)) = agent.inner().host_acp_terminals.remove(&terminal_id) else {
2408 return Err(AcpDispatchError::new(
2409 -32602,
2410 format!("ACP terminal not found: {terminal_id}"),
2411 ));
2412 };
2413 if terminal.exit_rx.borrow().is_none() {
2415 let _ = agent.acp_kill_terminal_shell(&terminal.shell_id);
2416 }
2417 let _ = agent.close_shell(&terminal.shell_id);
2419 Ok(Value::Null)
2420}
2421
2422fn handle_acp_resize_terminal(
2423 agent: &AgentOs,
2424 params: &Map<String, Value>,
2425) -> Result<Value, AcpDispatchError> {
2426 let method = "terminal/resize";
2427 let terminal_id = require_acp_terminal_id(params, method)?;
2428 let shell_id = acp_terminal_shell_id(agent, &terminal_id)?;
2429 let cols = optional_acp_number(params, "cols", method)?;
2430 let rows = optional_acp_number(params, "rows", method)?;
2431 let (Some(cols), Some(rows)) = (cols, rows) else {
2432 return Err(AcpDispatchError::new(
2433 -32602,
2434 format!("{method} requires numeric cols and rows"),
2435 ));
2436 };
2437 agent
2438 .resize_shell(&shell_id, cols.trunc() as u16, rows.trunc() as u16)
2439 .map_err(|error| AcpDispatchError::new(-32603, error.to_string()))?;
2440 Ok(Value::Null)
2441}
2442
2443fn acp_terminal_shell_id(agent: &AgentOs, terminal_id: &str) -> Result<String, AcpDispatchError> {
2445 agent
2446 .inner()
2447 .host_acp_terminals
2448 .read(terminal_id, |_, terminal| terminal.shell_id.clone())
2449 .ok_or_else(|| {
2450 AcpDispatchError::new(-32602, format!("ACP terminal not found: {terminal_id}"))
2451 })
2452}
2453
2454fn host_callback_callback() -> WireSidecarCallback {
2456 Arc::new(|payload, ownership| {
2457 Box::pin(async move {
2458 let request = match payload {
2459 wire::SidecarRequestPayload::HostCallbackRequest(request) => request,
2460 wire::SidecarRequestPayload::JsBridgeCallRequest(_) => {
2461 return Ok(wire::SidecarResponsePayload::HostCallbackResultResponse(
2462 wire::HostCallbackResultResponse {
2463 invocation_id: "unknown".to_string(),
2464 result: None,
2465 error: Some("host-callback received a non-tool request".to_string()),
2466 },
2467 ));
2468 }
2469 wire::SidecarRequestPayload::ExtEnvelope(envelope) => {
2470 return Ok(wire::SidecarResponsePayload::ExtEnvelope(
2471 wire::ExtEnvelope {
2472 namespace: envelope.namespace,
2473 payload: b"host-callback received an extension request".to_vec(),
2474 },
2475 ));
2476 }
2477 };
2478 Ok(wire::SidecarResponsePayload::HostCallbackResultResponse(
2479 run_host_callback(&ownership, request).await,
2480 ))
2481 })
2482 })
2483}
2484
2485async fn run_host_callback(
2488 ownership: &wire::OwnershipScope,
2489 request: wire::HostCallbackRequest,
2490) -> wire::HostCallbackResultResponse {
2491 let input = match serde_json::from_str::<Value>(&request.input) {
2492 Ok(input) => input,
2493 Err(error) => {
2494 return wire::HostCallbackResultResponse {
2495 invocation_id: request.invocation_id,
2496 result: None,
2497 error: Some(format!("Invalid host callback input: {error}")),
2498 };
2499 }
2500 };
2501 let vm_id = wire_ownership_vm_id(ownership).unwrap_or("");
2502 let registry = vm_tools().read(vm_id, |_, registry| registry.clone());
2503 let Some(registry) = registry else {
2504 return wire::HostCallbackResultResponse {
2505 invocation_id: request.invocation_id,
2506 result: None,
2507 error: Some(format!("Unknown tool \"{}\"", request.callback_key)),
2508 };
2509 };
2510
2511 if let Some(command) = parse_host_command_callback_input(&input) {
2512 return match run_host_command_callback(ownership, registry.as_ref(), command).await {
2513 Ok(value) => match host_callback_json_result(value) {
2514 Ok(result) => wire::HostCallbackResultResponse {
2515 invocation_id: request.invocation_id,
2516 result: Some(result),
2517 error: None,
2518 },
2519 Err(error) => wire::HostCallbackResultResponse {
2520 invocation_id: request.invocation_id,
2521 result: None,
2522 error: Some(error),
2523 },
2524 },
2525 Err(error) => wire::HostCallbackResultResponse {
2526 invocation_id: request.invocation_id,
2527 result: None,
2528 error: Some(error),
2529 },
2530 };
2531 }
2532
2533 let tool = registry.tool_map.get(&request.callback_key).cloned();
2534 let Some(tool) = tool else {
2535 return wire::HostCallbackResultResponse {
2536 invocation_id: request.invocation_id,
2537 result: None,
2538 error: Some(format!("Unknown tool \"{}\"", request.callback_key)),
2539 };
2540 };
2541 let timeout = Duration::from_millis(request.timeout_ms.max(1));
2542 match tokio::time::timeout(timeout, (tool.execute)(input)).await {
2543 Ok(Ok(value)) => match host_callback_json_result(value) {
2544 Ok(result) => wire::HostCallbackResultResponse {
2545 invocation_id: request.invocation_id,
2546 result: Some(result),
2547 error: None,
2548 },
2549 Err(error) => wire::HostCallbackResultResponse {
2550 invocation_id: request.invocation_id,
2551 result: None,
2552 error: Some(error),
2553 },
2554 },
2555 Ok(Err(error)) => wire::HostCallbackResultResponse {
2556 invocation_id: request.invocation_id,
2557 result: None,
2558 error: Some(error),
2559 },
2560 Err(_) => wire::HostCallbackResultResponse {
2561 invocation_id: request.invocation_id,
2562 result: None,
2563 error: Some(format!(
2564 "Tool \"{}\" timed out after {}ms",
2565 request.callback_key, request.timeout_ms
2566 )),
2567 },
2568 }
2569}
2570
2571#[derive(Debug, Deserialize)]
2572struct HostCommandCallbackInput {
2573 #[serde(rename = "type")]
2574 kind: String,
2575 command: String,
2576 #[serde(default)]
2577 args: Vec<String>,
2578 cwd: String,
2579}
2580
2581fn parse_host_command_callback_input(input: &Value) -> Option<HostCommandCallbackInput> {
2582 let command = serde_json::from_value::<HostCommandCallbackInput>(input.clone()).ok()?;
2583 if command.kind == "command" {
2584 Some(command)
2585 } else {
2586 None
2587 }
2588}
2589
2590async fn run_host_command_callback(
2591 ownership: &wire::OwnershipScope,
2592 registry: &VmHostToolRegistry,
2593 command: HostCommandCallbackInput,
2594) -> Result<Value, String> {
2595 if command.command == "agentos" {
2596 return handle_agentos_registry_command(ownership, registry, &command).await;
2597 }
2598 let Some(toolkit) = registry
2599 .tool_kits
2600 .iter()
2601 .find(|toolkit| format!("agentos-{}", toolkit.name) == command.command)
2602 else {
2603 return Err(format!(
2604 "Unknown host callback command \"{}\"",
2605 command.command
2606 ));
2607 };
2608 handle_agentos_toolkit_command(ownership, registry, &command, toolkit).await
2609}
2610
2611async fn handle_agentos_registry_command(
2612 ownership: &wire::OwnershipScope,
2613 registry: &VmHostToolRegistry,
2614 command: &HostCommandCallbackInput,
2615) -> Result<Value, String> {
2616 let Some(subcommand) = command.args.first() else {
2617 return Ok(json_object([(
2618 "usage",
2619 Value::String(String::from(
2620 "agentos <command>: list-tools [toolkit], <toolkit> --help, or <toolkit> <tool> ...",
2621 )),
2622 )]));
2623 };
2624 if is_help_flag(subcommand) {
2625 return Ok(json_object([(
2626 "usage",
2627 Value::String(String::from(
2628 "agentos <command>: list-tools [toolkit], <toolkit> --help, or <toolkit> <tool> ...",
2629 )),
2630 )]));
2631 }
2632 if subcommand == "list-tools" {
2633 return match command.args.get(1) {
2634 Some(toolkit_name) => describe_toolkit_payload(®istry.tool_kits, toolkit_name),
2635 None => Ok(list_toolkits_payload(®istry.tool_kits)),
2636 };
2637 }
2638
2639 let Some(toolkit) = registry
2640 .tool_kits
2641 .iter()
2642 .find(|toolkit| toolkit.name == *subcommand)
2643 else {
2644 return Err(format!(
2645 "No toolkit \"{subcommand}\". Available: {}",
2646 toolkit_names(®istry.tool_kits)
2647 ));
2648 };
2649
2650 let Some(tool_name) = command.args.get(1) else {
2651 return describe_toolkit_payload(®istry.tool_kits, subcommand);
2652 };
2653 if is_help_flag(tool_name) {
2654 return describe_toolkit_payload(®istry.tool_kits, subcommand);
2655 }
2656 if command.args.get(2).is_some_and(|value| is_help_flag(value)) {
2657 return describe_tool_payload(toolkit, tool_name);
2658 }
2659 invoke_host_tool(
2660 ownership,
2661 registry,
2662 toolkit,
2663 tool_name,
2664 command.args.get(2..).unwrap_or_default(),
2665 &command.cwd,
2666 )
2667 .await
2668}
2669
2670async fn handle_agentos_toolkit_command(
2671 ownership: &wire::OwnershipScope,
2672 registry: &VmHostToolRegistry,
2673 command: &HostCommandCallbackInput,
2674 toolkit: &ToolKit,
2675) -> Result<Value, String> {
2676 let Some(tool_name) = command.args.first() else {
2677 return describe_toolkit_payload(®istry.tool_kits, &toolkit.name);
2678 };
2679 if is_help_flag(tool_name) {
2680 return describe_toolkit_payload(®istry.tool_kits, &toolkit.name);
2681 }
2682 if command.args.get(1).is_some_and(|value| is_help_flag(value)) {
2683 return describe_tool_payload(toolkit, tool_name);
2684 }
2685 invoke_host_tool(
2686 ownership,
2687 registry,
2688 toolkit,
2689 tool_name,
2690 command.args.get(1..).unwrap_or_default(),
2691 &command.cwd,
2692 )
2693 .await
2694}
2695
2696async fn invoke_host_tool(
2697 ownership: &wire::OwnershipScope,
2698 registry: &VmHostToolRegistry,
2699 toolkit: &ToolKit,
2700 tool_name: &str,
2701 args: &[String],
2702 cwd: &str,
2703) -> Result<Value, String> {
2704 let callback_key = format!("{}:{tool_name}", toolkit.name);
2705 let Some(tool) = registry.tool_map.get(&callback_key).cloned() else {
2706 return Err(format!(
2707 "No tool \"{tool_name}\" in toolkit \"{}\". Available: {}",
2708 toolkit.name,
2709 tool_names(toolkit)
2710 ));
2711 };
2712
2713 if tool_permission_mode(registry.permissions.as_ref(), &callback_key) != PermissionMode::Allow {
2714 return Err(format!(
2715 "EACCES: blocked by binding.invoke policy for {callback_key}"
2716 ));
2717 }
2718
2719 let input = parse_host_tool_input(ownership, &tool, args, cwd).await?;
2720 validate_tool_input(&tool.input_schema, &input).map_err(|error| error.to_string())?;
2721
2722 let timeout = Duration::from_millis(tool.timeout_ms.unwrap_or(30_000).max(1));
2723 match tokio::time::timeout(timeout, (tool.execute)(input)).await {
2724 Ok(Ok(value)) => Ok(value),
2725 Ok(Err(error)) => Err(error),
2726 Err(_) => Err(format!(
2727 "Tool \"{callback_key}\" timed out after {}ms",
2728 tool.timeout_ms.unwrap_or(30_000)
2729 )),
2730 }
2731}
2732
2733async fn parse_host_tool_input(
2734 ownership: &wire::OwnershipScope,
2735 tool: &HostTool,
2736 args: &[String],
2737 cwd: &str,
2738) -> Result<Value, String> {
2739 if args.first().is_some_and(|arg| arg == "--json") {
2740 let value = args
2741 .get(1)
2742 .ok_or_else(|| String::from("Flag --json requires a value"))?;
2743 return serde_json::from_str(value)
2744 .map_err(|error| format!("Invalid JSON for --json: {error}"));
2745 }
2746
2747 if args.first().is_some_and(|arg| arg == "--json-file") {
2748 let path = args
2749 .get(1)
2750 .ok_or_else(|| String::from("Flag --json-file requires a value"))?;
2751 let guest_path = normalize_guest_path(if path.starts_with('/') {
2752 path.clone()
2753 } else {
2754 format!("{cwd}/{path}")
2755 });
2756 let vm_id = wire_ownership_vm_id(ownership).unwrap_or("");
2757 let inner = vm_permission_routers()
2758 .read(vm_id, |_, weak| weak.clone())
2759 .and_then(|weak| weak.upgrade())
2760 .ok_or_else(|| String::from("Invalid JSON file: VM is no longer available"))?;
2761 let bytes = AgentOs { inner }
2762 .read_file(&guest_path)
2763 .await
2764 .map_err(|error| format!("Invalid JSON file: {error}"))?;
2765 let text =
2766 String::from_utf8(bytes).map_err(|error| format!("Invalid JSON file: {error}"))?;
2767 return serde_json::from_str(&text).map_err(|error| format!("Invalid JSON file: {error}"));
2768 }
2769
2770 parse_tool_argv(&tool.input_schema, args)
2771}
2772
2773fn host_callback_json_result(value: Value) -> Result<String, String> {
2774 serde_json::to_string(&value).map_err(|error| format!("Invalid host callback result: {error}"))
2775}
2776
2777fn parse_tool_argv(schema: &Value, argv: &[String]) -> Result<Value, String> {
2778 let properties = schema
2779 .get("properties")
2780 .and_then(Value::as_object)
2781 .cloned()
2782 .unwrap_or_default();
2783 let required = schema
2784 .get("required")
2785 .and_then(Value::as_array)
2786 .map(|items| {
2787 items
2788 .iter()
2789 .filter_map(Value::as_str)
2790 .map(str::to_owned)
2791 .collect::<std::collections::BTreeSet<_>>()
2792 })
2793 .unwrap_or_default();
2794
2795 let mut flag_to_field = BTreeMap::new();
2796 for (field_name, field_schema) in &properties {
2797 flag_to_field.insert(
2798 camel_to_kebab(field_name),
2799 (field_name.clone(), field_schema.clone()),
2800 );
2801 }
2802
2803 let mut input = Map::new();
2804 let mut index = 0;
2805 while index < argv.len() {
2806 let arg = &argv[index];
2807 if !arg.starts_with("--") {
2808 return Err(format!("Unexpected positional argument: \"{arg}\""));
2809 }
2810
2811 let raw_flag = &arg[2..];
2812 let (flag_name, negated) = raw_flag
2813 .strip_prefix("no-")
2814 .map(|name| (name, true))
2815 .unwrap_or((raw_flag, false));
2816 let Some((field_name, field_schema)) = flag_to_field.get(flag_name) else {
2817 return Err(format!("Unknown flag: --{raw_flag}"));
2818 };
2819 let field_type = json_schema_type(field_schema);
2820
2821 if negated {
2822 if field_type != Some("boolean") {
2823 return Err(format!("Unknown flag: --{raw_flag}"));
2824 }
2825 input.insert(field_name.clone(), Value::Bool(false));
2826 index += 1;
2827 continue;
2828 }
2829
2830 match field_type {
2831 Some("boolean") => {
2832 input.insert(field_name.clone(), Value::Bool(true));
2833 index += 1;
2834 }
2835 Some("number") | Some("integer") => {
2836 let value = argv
2837 .get(index + 1)
2838 .ok_or_else(|| format!("Flag --{raw_flag} requires a value"))?;
2839 let number = value
2840 .parse::<f64>()
2841 .map_err(|_| format!("Flag --{raw_flag} expects a number, got \"{value}\""))?;
2842 let number = serde_json::Number::from_f64(number).ok_or_else(|| {
2843 format!("Flag --{raw_flag} expects a finite number, got \"{value}\"")
2844 })?;
2845 input.insert(field_name.clone(), Value::Number(number));
2846 index += 2;
2847 }
2848 Some("array") => {
2849 let value = argv
2850 .get(index + 1)
2851 .ok_or_else(|| format!("Flag --{raw_flag} requires a value"))?;
2852 let item_type = field_schema.get("items").and_then(json_schema_type);
2853 let parsed_value = match item_type {
2854 Some("number") | Some("integer") => {
2855 let number = value.parse::<f64>().map_err(|_| {
2856 format!("Flag --{raw_flag} expects a number value, got \"{value}\"")
2857 })?;
2858 let number = serde_json::Number::from_f64(number).ok_or_else(|| {
2859 format!(
2860 "Flag --{raw_flag} expects a finite number value, got \"{value}\""
2861 )
2862 })?;
2863 Value::Number(number)
2864 }
2865 Some("boolean") => {
2866 let boolean = value.parse::<bool>().map_err(|_| {
2867 format!("Flag --{raw_flag} expects a boolean value, got \"{value}\"")
2868 })?;
2869 Value::Bool(boolean)
2870 }
2871 _ => Value::String(value.clone()),
2872 };
2873 input
2874 .entry(field_name.clone())
2875 .or_insert_with(|| Value::Array(Vec::new()))
2876 .as_array_mut()
2877 .expect("array field should always contain an array")
2878 .push(parsed_value);
2879 index += 2;
2880 }
2881 _ => {
2882 let value = argv
2883 .get(index + 1)
2884 .ok_or_else(|| format!("Flag --{raw_flag} requires a value"))?;
2885 input.insert(field_name.clone(), Value::String(value.clone()));
2886 index += 2;
2887 }
2888 }
2889 }
2890
2891 for field_name in required {
2892 if !input.contains_key(&field_name) {
2893 return Err(format!(
2894 "Missing required flag: --{}",
2895 camel_to_kebab(&field_name)
2896 ));
2897 }
2898 }
2899
2900 Ok(Value::Object(input))
2901}
2902
2903#[derive(Debug, Clone, PartialEq, Eq)]
2904struct ToolInputSchemaViolation {
2905 path: String,
2906 expected: String,
2907 actual: String,
2908}
2909
2910impl ToolInputSchemaViolation {
2911 fn new(
2912 path: impl Into<String>,
2913 expected: impl Into<String>,
2914 actual: impl Into<String>,
2915 ) -> Self {
2916 Self {
2917 path: path.into(),
2918 expected: expected.into(),
2919 actual: actual.into(),
2920 }
2921 }
2922}
2923
2924impl std::fmt::Display for ToolInputSchemaViolation {
2925 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2926 write!(
2927 f,
2928 "ToolInputSchemaViolation at {}: expected {}, got {}",
2929 self.path, self.expected, self.actual
2930 )
2931 }
2932}
2933
2934fn validate_tool_input(schema: &Value, input: &Value) -> Result<(), ToolInputSchemaViolation> {
2935 validate_tool_input_at_path(schema, input, "$")
2936}
2937
2938fn validate_tool_input_at_path(
2939 schema: &Value,
2940 input: &Value,
2941 path: &str,
2942) -> Result<(), ToolInputSchemaViolation> {
2943 if schema.is_null() || schema.as_object().is_some_and(|object| object.is_empty()) {
2944 return Ok(());
2945 }
2946 if let Some(branches) = schema.get("anyOf").and_then(Value::as_array) {
2947 return validate_schema_branches(branches, input, path, "anyOf");
2948 }
2949 if let Some(branches) = schema.get("oneOf").and_then(Value::as_array) {
2950 return validate_schema_branches(branches, input, path, "oneOf");
2951 }
2952 if let Some(enum_values) = schema.get("enum").and_then(Value::as_array) {
2953 if enum_values.iter().any(|candidate| candidate == input) {
2954 return Ok(());
2955 }
2956 return Err(ToolInputSchemaViolation::new(
2957 path,
2958 format!(
2959 "one of {}",
2960 enum_values
2961 .iter()
2962 .map(compact_json)
2963 .collect::<Vec<_>>()
2964 .join(", ")
2965 ),
2966 describe_value(input),
2967 ));
2968 }
2969 if let Some(expected) = schema.get("const") {
2970 if expected == input {
2971 return Ok(());
2972 }
2973 return Err(ToolInputSchemaViolation::new(
2974 path,
2975 format!("constant {}", compact_json(expected)),
2976 describe_value(input),
2977 ));
2978 }
2979
2980 match schema.get("type") {
2981 Some(Value::String(expected_type)) => {
2982 validate_typed_tool_input(schema, input, path, expected_type)
2983 }
2984 Some(Value::Array(expected_types)) => {
2985 let mut first_error = None;
2986 for expected_type in expected_types.iter().filter_map(Value::as_str) {
2987 match validate_typed_tool_input(schema, input, path, expected_type) {
2988 Ok(()) => return Ok(()),
2989 Err(error) if first_error.is_none() => first_error = Some(error),
2990 Err(_) => {}
2991 }
2992 }
2993 Err(first_error.unwrap_or_else(|| {
2994 ToolInputSchemaViolation::new(
2995 path,
2996 describe_expected(schema),
2997 describe_value(input),
2998 )
2999 }))
3000 }
3001 Some(_) => Ok(()),
3002 None if has_object_keywords(schema) => {
3003 validate_typed_tool_input(schema, input, path, "object")
3004 }
3005 None => Ok(()),
3006 }
3007}
3008
3009fn validate_schema_branches(
3010 branches: &[Value],
3011 input: &Value,
3012 path: &str,
3013 keyword: &str,
3014) -> Result<(), ToolInputSchemaViolation> {
3015 let mut first_error = None;
3016 for branch in branches {
3017 match validate_tool_input_at_path(branch, input, path) {
3018 Ok(()) => return Ok(()),
3019 Err(error) if first_error.is_none() => first_error = Some(error),
3020 Err(_) => {}
3021 }
3022 }
3023 Err(first_error.unwrap_or_else(|| {
3024 ToolInputSchemaViolation::new(
3025 path,
3026 format!(
3027 "{keyword} branch ({})",
3028 branches
3029 .iter()
3030 .map(describe_expected)
3031 .collect::<Vec<_>>()
3032 .join(" | ")
3033 ),
3034 describe_value(input),
3035 )
3036 }))
3037}
3038
3039fn validate_typed_tool_input(
3040 schema: &Value,
3041 input: &Value,
3042 path: &str,
3043 expected_type: &str,
3044) -> Result<(), ToolInputSchemaViolation> {
3045 match expected_type {
3046 "null" if input.is_null() => Ok(()),
3047 "null" => Err(type_violation(path, expected_type, input)),
3048 "boolean" if input.is_boolean() => Ok(()),
3049 "boolean" => Err(type_violation(path, expected_type, input)),
3050 "string" => validate_string_tool_input(schema, input, path),
3051 "number" => validate_number_tool_input(schema, input, path, false),
3052 "integer" => validate_number_tool_input(schema, input, path, true),
3053 "array" => validate_array_tool_input(schema, input, path),
3054 "object" => validate_object_tool_input(schema, input, path),
3055 _ => Ok(()),
3056 }
3057}
3058
3059fn validate_string_tool_input(
3060 schema: &Value,
3061 input: &Value,
3062 path: &str,
3063) -> Result<(), ToolInputSchemaViolation> {
3064 let Some(value) = input.as_str() else {
3065 return Err(type_violation(path, "string", input));
3066 };
3067 if let Some(min_length) = schema.get("minLength").and_then(Value::as_u64) {
3068 if value.chars().count() < min_length as usize {
3069 return Err(ToolInputSchemaViolation::new(
3070 path,
3071 format!("string with minLength {min_length}"),
3072 format!("string length {}", value.chars().count()),
3073 ));
3074 }
3075 }
3076 if let Some(max_length) = schema.get("maxLength").and_then(Value::as_u64) {
3077 if value.chars().count() > max_length as usize {
3078 return Err(ToolInputSchemaViolation::new(
3079 path,
3080 format!("string with maxLength {max_length}"),
3081 format!("string length {}", value.chars().count()),
3082 ));
3083 }
3084 }
3085 Ok(())
3086}
3087
3088fn validate_number_tool_input(
3089 schema: &Value,
3090 input: &Value,
3091 path: &str,
3092 expect_integer: bool,
3093) -> Result<(), ToolInputSchemaViolation> {
3094 let Some(number) = input.as_f64() else {
3095 return Err(type_violation(
3096 path,
3097 if expect_integer { "integer" } else { "number" },
3098 input,
3099 ));
3100 };
3101 if expect_integer && number.fract() != 0.0 {
3102 return Err(type_violation(path, "integer", input));
3103 }
3104 if let Some(minimum) = schema.get("minimum").and_then(Value::as_f64) {
3105 if number < minimum {
3106 return Err(ToolInputSchemaViolation::new(
3107 path,
3108 format!(
3109 "{} >= {}",
3110 if expect_integer { "integer" } else { "number" },
3111 minimum
3112 ),
3113 compact_json(input),
3114 ));
3115 }
3116 }
3117 if let Some(minimum) = schema.get("exclusiveMinimum").and_then(Value::as_f64) {
3118 if number <= minimum {
3119 return Err(ToolInputSchemaViolation::new(
3120 path,
3121 format!(
3122 "{} > {}",
3123 if expect_integer { "integer" } else { "number" },
3124 minimum
3125 ),
3126 compact_json(input),
3127 ));
3128 }
3129 }
3130 if let Some(maximum) = schema.get("maximum").and_then(Value::as_f64) {
3131 if number > maximum {
3132 return Err(ToolInputSchemaViolation::new(
3133 path,
3134 format!(
3135 "{} <= {}",
3136 if expect_integer { "integer" } else { "number" },
3137 maximum
3138 ),
3139 compact_json(input),
3140 ));
3141 }
3142 }
3143 if let Some(maximum) = schema.get("exclusiveMaximum").and_then(Value::as_f64) {
3144 if number >= maximum {
3145 return Err(ToolInputSchemaViolation::new(
3146 path,
3147 format!(
3148 "{} < {}",
3149 if expect_integer { "integer" } else { "number" },
3150 maximum
3151 ),
3152 compact_json(input),
3153 ));
3154 }
3155 }
3156 Ok(())
3157}
3158
3159fn validate_array_tool_input(
3160 schema: &Value,
3161 input: &Value,
3162 path: &str,
3163) -> Result<(), ToolInputSchemaViolation> {
3164 let Some(items) = input.as_array() else {
3165 return Err(type_violation(path, "array", input));
3166 };
3167 if let Some(min_items) = schema.get("minItems").and_then(Value::as_u64) {
3168 if items.len() < min_items as usize {
3169 return Err(ToolInputSchemaViolation::new(
3170 path,
3171 format!("array with minItems {min_items}"),
3172 format!("array length {}", items.len()),
3173 ));
3174 }
3175 }
3176 if let Some(max_items) = schema.get("maxItems").and_then(Value::as_u64) {
3177 if items.len() > max_items as usize {
3178 return Err(ToolInputSchemaViolation::new(
3179 path,
3180 format!("array with maxItems {max_items}"),
3181 format!("array length {}", items.len()),
3182 ));
3183 }
3184 }
3185 if let Some(item_schema) = schema.get("items") {
3186 for (index, item) in items.iter().enumerate() {
3187 validate_tool_input_at_path(item_schema, item, &format!("{path}[{index}]"))?;
3188 }
3189 }
3190 Ok(())
3191}
3192
3193fn validate_object_tool_input(
3194 schema: &Value,
3195 input: &Value,
3196 path: &str,
3197) -> Result<(), ToolInputSchemaViolation> {
3198 let Some(object) = input.as_object() else {
3199 return Err(type_violation(path, "object", input));
3200 };
3201 let properties = schema
3202 .get("properties")
3203 .and_then(Value::as_object)
3204 .cloned()
3205 .unwrap_or_default();
3206 let required = schema
3207 .get("required")
3208 .and_then(Value::as_array)
3209 .cloned()
3210 .unwrap_or_default();
3211 for field in required.iter().filter_map(Value::as_str) {
3212 if !object.contains_key(field) {
3213 let field_path = format!("{path}.{field}");
3214 let expected = properties
3215 .get(field)
3216 .map(describe_expected)
3217 .unwrap_or_else(|| String::from("required value"));
3218 return Err(ToolInputSchemaViolation::new(
3219 field_path,
3220 expected,
3221 "missing value",
3222 ));
3223 }
3224 }
3225 for (field, value) in object {
3226 let field_path = format!("{path}.{field}");
3227 if let Some(field_schema) = properties.get(field) {
3228 validate_tool_input_at_path(field_schema, value, &field_path)?;
3229 continue;
3230 }
3231 match schema.get("additionalProperties") {
3232 Some(Value::Bool(false)) => {
3233 return Err(ToolInputSchemaViolation::new(
3234 field_path,
3235 "no additional properties",
3236 describe_value(value),
3237 ));
3238 }
3239 Some(additional_schema) => {
3240 validate_tool_input_at_path(additional_schema, value, &field_path)?;
3241 }
3242 None => {}
3243 }
3244 }
3245 Ok(())
3246}
3247
3248fn has_object_keywords(schema: &Value) -> bool {
3249 schema.get("properties").is_some()
3250 || schema.get("required").is_some()
3251 || schema.get("additionalProperties").is_some()
3252}
3253
3254fn type_violation(path: &str, expected: &str, input: &Value) -> ToolInputSchemaViolation {
3255 ToolInputSchemaViolation::new(path, expected, describe_value(input))
3256}
3257
3258fn describe_expected(schema: &Value) -> String {
3259 if let Some(enum_values) = schema.get("enum").and_then(Value::as_array) {
3260 return format!(
3261 "one of {}",
3262 enum_values
3263 .iter()
3264 .map(compact_json)
3265 .collect::<Vec<_>>()
3266 .join(", ")
3267 );
3268 }
3269 if let Some(expected) = schema.get("const") {
3270 return format!("constant {}", compact_json(expected));
3271 }
3272 match schema.get("type") {
3273 Some(Value::String(expected_type)) => expected_type.clone(),
3274 Some(Value::Array(expected_types)) => expected_types
3275 .iter()
3276 .filter_map(Value::as_str)
3277 .collect::<Vec<_>>()
3278 .join(" | "),
3279 _ if has_object_keywords(schema) => String::from("object"),
3280 _ => String::from("value"),
3281 }
3282}
3283
3284fn describe_value(value: &Value) -> String {
3285 match value {
3286 Value::Null => String::from("null"),
3287 Value::Bool(_) => String::from("boolean"),
3288 Value::Number(number) => {
3289 let is_integer = number.as_i64().is_some()
3290 || number.as_u64().is_some()
3291 || number.as_f64().is_some_and(|float| float.fract() == 0.0);
3292 if is_integer {
3293 String::from("integer")
3294 } else {
3295 String::from("number")
3296 }
3297 }
3298 Value::String(_) => String::from("string"),
3299 Value::Array(_) => String::from("array"),
3300 Value::Object(_) => String::from("object"),
3301 }
3302}
3303
3304fn compact_json(value: &Value) -> String {
3305 serde_json::to_string(value).unwrap_or_else(|_| String::from("<invalid json>"))
3306}
3307
3308fn list_toolkits_payload(tool_kits: &[ToolKit]) -> Value {
3309 Value::Object(Map::from_iter([(
3310 String::from("toolkits"),
3311 Value::Array(
3312 tool_kits
3313 .iter()
3314 .map(|toolkit| {
3315 json_object([
3316 ("name", Value::String(toolkit.name.clone())),
3317 ("description", Value::String(toolkit.description.clone())),
3318 (
3319 "tools",
3320 Value::Array(
3321 toolkit
3322 .tools
3323 .iter()
3324 .map(|tool| Value::String(tool.name.clone()))
3325 .collect(),
3326 ),
3327 ),
3328 ])
3329 })
3330 .collect(),
3331 ),
3332 )]))
3333}
3334
3335fn describe_toolkit_payload(tool_kits: &[ToolKit], toolkit_name: &str) -> Result<Value, String> {
3336 let Some(toolkit) = tool_kits
3337 .iter()
3338 .find(|toolkit| toolkit.name == toolkit_name)
3339 else {
3340 return Err(format!(
3341 "No toolkit \"{toolkit_name}\". Available: {}",
3342 toolkit_names(tool_kits)
3343 ));
3344 };
3345 Ok(json_object([
3346 ("name", Value::String(toolkit.name.clone())),
3347 ("description", Value::String(toolkit.description.clone())),
3348 (
3349 "tools",
3350 Value::Object(Map::from_iter(toolkit.tools.iter().map(|tool| {
3351 (
3352 tool.name.clone(),
3353 json_object([
3354 ("description", Value::String(tool.description.clone())),
3355 (
3356 "flags",
3357 Value::Array(describe_tool_flags(&tool.input_schema)),
3358 ),
3359 ]),
3360 )
3361 }))),
3362 ),
3363 ]))
3364}
3365
3366fn describe_tool_payload(toolkit: &ToolKit, tool_name: &str) -> Result<Value, String> {
3367 let Some(tool) = toolkit.tools.iter().find(|tool| tool.name == tool_name) else {
3368 return Err(format!(
3369 "No tool \"{tool_name}\" in toolkit \"{}\". Available: {}",
3370 toolkit.name,
3371 tool_names(toolkit)
3372 ));
3373 };
3374 Ok(json_object([
3375 ("toolkit", Value::String(toolkit.name.clone())),
3376 ("tool", Value::String(tool_name.to_string())),
3377 ("description", Value::String(tool.description.clone())),
3378 (
3379 "flags",
3380 Value::Array(describe_tool_flags(&tool.input_schema)),
3381 ),
3382 ("examples", Value::Array(Vec::new())),
3383 ]))
3384}
3385
3386fn describe_tool_flags(schema: &Value) -> Vec<Value> {
3387 let properties = schema
3388 .get("properties")
3389 .and_then(Value::as_object)
3390 .cloned()
3391 .unwrap_or_default();
3392 let required = schema
3393 .get("required")
3394 .and_then(Value::as_array)
3395 .map(|items| {
3396 items
3397 .iter()
3398 .filter_map(Value::as_str)
3399 .map(str::to_owned)
3400 .collect::<std::collections::BTreeSet<_>>()
3401 })
3402 .unwrap_or_default();
3403 properties
3404 .into_iter()
3405 .map(|(field_name, field_schema)| {
3406 json_object([
3407 (
3408 "name",
3409 Value::String(format!("--{}", camel_to_kebab(&field_name))),
3410 ),
3411 (
3412 "type",
3413 Value::String(describe_tool_flag_type(&field_schema)),
3414 ),
3415 ("required", Value::Bool(required.contains(&field_name))),
3416 ])
3417 })
3418 .collect()
3419}
3420
3421fn describe_tool_flag_type(schema: &Value) -> String {
3422 match json_schema_type(schema) {
3423 Some("array") => {
3424 let item_type = schema
3425 .get("items")
3426 .and_then(json_schema_type)
3427 .unwrap_or("string");
3428 format!("{item_type}[]")
3429 }
3430 Some("string") => schema
3431 .get("enum")
3432 .and_then(Value::as_array)
3433 .map(|values| values.iter().filter_map(Value::as_str).collect::<Vec<_>>())
3434 .filter(|values| !values.is_empty())
3435 .map(|values| values.join("|"))
3436 .unwrap_or_else(|| String::from("string")),
3437 Some(other) => other.to_string(),
3438 None => String::from("string"),
3439 }
3440}
3441
3442fn tool_permission_mode(permissions: Option<&Permissions>, callback_key: &str) -> PermissionMode {
3443 let Some(permissions) = permissions else {
3444 return PermissionMode::Allow;
3445 };
3446 let Some(scope) = permissions.binding.as_ref() else {
3447 return PermissionMode::Allow;
3448 };
3449 match scope {
3450 crate::config::PatternPermissions::Mode(mode) => *mode,
3451 crate::config::PatternPermissions::Rules(rules) => {
3452 let mut mode = rules.default.unwrap_or(PermissionMode::Deny);
3453 for rule in &rules.rules {
3454 let operations_match = rule
3455 .operations
3456 .as_ref()
3457 .map(|operations| {
3458 operations
3459 .iter()
3460 .any(|operation| operation == "*" || operation == "invoke")
3461 })
3462 .unwrap_or(true);
3463 let patterns_match = rule
3464 .patterns
3465 .as_ref()
3466 .map(|patterns| {
3467 patterns
3468 .iter()
3469 .any(|pattern| permission_pattern_matches(pattern, callback_key))
3470 })
3471 .unwrap_or(true);
3472 if operations_match && patterns_match {
3473 mode = rule.mode;
3474 }
3475 }
3476 mode
3477 }
3478 }
3479}
3480
3481fn permission_pattern_matches(pattern: &str, value: &str) -> bool {
3482 if pattern == "*" || pattern == "**" || pattern == value {
3483 return true;
3484 }
3485 let mut pattern_index = 0;
3486 let mut value_index = 0;
3487 let pattern_bytes = pattern.as_bytes();
3488 let value_bytes = value.as_bytes();
3489 let mut star_index = None;
3490 let mut match_index = 0;
3491 while value_index < value_bytes.len() {
3492 if pattern_index < pattern_bytes.len()
3493 && pattern_bytes[pattern_index] == b'*'
3494 && pattern_index + 1 < pattern_bytes.len()
3495 && pattern_bytes[pattern_index + 1] == b'*'
3496 {
3497 star_index = Some(pattern_index);
3498 match_index = value_index;
3499 pattern_index += 2;
3500 } else if pattern_index < pattern_bytes.len() && pattern_bytes[pattern_index] == b'*' {
3501 star_index = Some(pattern_index);
3502 match_index = value_index;
3503 pattern_index += 1;
3504 } else if pattern_index < pattern_bytes.len()
3505 && pattern_bytes[pattern_index] == value_bytes[value_index]
3506 {
3507 pattern_index += 1;
3508 value_index += 1;
3509 } else if let Some(star) = star_index {
3510 if pattern_bytes[star] == b'*'
3511 && star + 1 < pattern_bytes.len()
3512 && pattern_bytes[star + 1] != b'*'
3513 && value_bytes.get(match_index) == Some(&b':')
3514 {
3515 return false;
3516 }
3517 pattern_index = if star + 1 < pattern_bytes.len() && pattern_bytes[star + 1] == b'*' {
3518 star + 2
3519 } else {
3520 star + 1
3521 };
3522 match_index += 1;
3523 value_index = match_index;
3524 } else {
3525 return false;
3526 }
3527 }
3528 while pattern_index < pattern_bytes.len() && pattern_bytes[pattern_index] == b'*' {
3529 pattern_index += if pattern_index + 1 < pattern_bytes.len()
3530 && pattern_bytes[pattern_index + 1] == b'*'
3531 {
3532 2
3533 } else {
3534 1
3535 };
3536 }
3537 pattern_index == pattern_bytes.len()
3538}
3539
3540fn toolkit_names(tool_kits: &[ToolKit]) -> String {
3541 tool_kits
3542 .iter()
3543 .map(|toolkit| toolkit.name.clone())
3544 .collect::<Vec<_>>()
3545 .join(", ")
3546}
3547
3548fn tool_names(toolkit: &ToolKit) -> String {
3549 toolkit
3550 .tools
3551 .iter()
3552 .map(|tool| tool.name.clone())
3553 .collect::<Vec<_>>()
3554 .join(", ")
3555}
3556
3557fn is_help_flag(value: &str) -> bool {
3558 matches!(value, "--help" | "-h")
3559}
3560
3561fn json_schema_type(schema: &Value) -> Option<&str> {
3562 schema.get("type").and_then(Value::as_str)
3563}
3564
3565fn camel_to_kebab(value: &str) -> String {
3566 let mut output = String::new();
3567 for (index, ch) in value.chars().enumerate() {
3568 if ch.is_ascii_uppercase() && index > 0 {
3569 output.push('-');
3570 }
3571 output.push(ch.to_ascii_lowercase());
3572 }
3573 output
3574}
3575
3576fn normalize_guest_path(path: String) -> String {
3577 let absolute = path.starts_with('/');
3578 let mut parts = Vec::new();
3579 for part in path.split('/') {
3580 match part {
3581 "" | "." => {}
3582 ".." => {
3583 parts.pop();
3584 }
3585 _ => parts.push(part),
3586 }
3587 }
3588 let normalized = parts.join("/");
3589 if absolute {
3590 format!("/{normalized}")
3591 } else {
3592 normalized
3593 }
3594}
3595
3596fn json_object<const N: usize>(entries: [(&str, Value); N]) -> Value {
3597 Value::Object(Map::from_iter(
3598 entries
3599 .into_iter()
3600 .map(|(key, value)| (key.to_string(), value)),
3601 ))
3602}
3603
3604fn build_package_descriptors(config: &AgentOsConfig) -> Vec<wire::PackageDescriptor> {
3607 config
3608 .packages
3609 .iter()
3610 .map(|package| wire::PackageDescriptor {
3611 path: package.path.clone(),
3612 })
3613 .collect()
3614}
3615
3616fn projected_agents_from_wire(agents: Vec<wire::AgentosProjectedAgent>) -> Vec<ProjectedAgent> {
3617 agents
3618 .into_iter()
3619 .map(|agent| ProjectedAgent {
3620 id: agent.id,
3621 acp_entrypoint: agent.acp_entrypoint,
3622 adapter_entrypoint: agent.adapter_entrypoint,
3623 })
3624 .collect()
3625}
3626
3627fn register_projected_agents(
3628 projected_agents: &parking_lot::Mutex<Vec<ProjectedAgent>>,
3629 agents: Vec<ProjectedAgent>,
3630) {
3631 let mut guard = projected_agents.lock();
3632 for agent in agents {
3633 guard.retain(|existing| existing.id != agent.id);
3634 guard.push(agent);
3635 }
3636}
3637
3638fn serialize_mounts(config: &AgentOsConfig) -> Result<Vec<wire::MountDescriptor>, ClientError> {
3639 config
3640 .mounts
3641 .iter()
3642 .map(|mount| match mount {
3643 MountConfig::Native {
3644 path,
3645 plugin,
3646 read_only,
3647 } => {
3648 let plugin_config = plugin
3649 .config
3650 .clone()
3651 .unwrap_or_else(|| serde_json::Value::Object(Default::default()));
3652 Ok(wire::MountDescriptor {
3653 guest_path: path.clone(),
3654 read_only: *read_only,
3655 plugin: wire::MountPluginDescriptor {
3656 id: plugin.id.clone(),
3657 config: json_utf8(&plugin_config, "native mount plugin config")?,
3658 },
3659 })
3660 }
3661 MountConfig::Plain { .. } => Err(ClientError::Sidecar(
3662 "plain mounts cannot be configured during Rust client VM creation".to_string(),
3663 )),
3664 MountConfig::Overlay { .. } => Err(ClientError::Sidecar(
3665 "overlay mounts cannot be configured during Rust client VM creation".to_string(),
3666 )),
3667 })
3668 .collect()
3669}
3670
3671fn permissions_policy(config: &AgentOsConfig) -> wire::PermissionsPolicy {
3672 let Some(permissions) = config.permissions.as_ref() else {
3673 return default_permissions_policy();
3674 };
3675
3676 wire::PermissionsPolicy {
3677 fs: Some(
3678 permissions
3679 .fs
3680 .as_ref()
3681 .map(serialize_fs_permissions)
3682 .unwrap_or(wire::FsPermissionScope::PermissionMode(
3683 wire::PermissionMode::Allow,
3684 )),
3685 ),
3686 network: Some(
3687 permissions
3688 .network
3689 .as_ref()
3690 .map(serialize_pattern_permissions)
3691 .unwrap_or_else(default_network_egress_scope),
3692 ),
3693 child_process: Some(
3694 permissions
3695 .child_process
3696 .as_ref()
3697 .map(serialize_pattern_permissions)
3698 .unwrap_or(wire::PatternPermissionScope::PermissionMode(
3699 wire::PermissionMode::Allow,
3700 )),
3701 ),
3702 process: Some(
3703 permissions
3704 .process
3705 .as_ref()
3706 .map(serialize_pattern_permissions)
3707 .unwrap_or(wire::PatternPermissionScope::PermissionMode(
3708 wire::PermissionMode::Allow,
3709 )),
3710 ),
3711 env: Some(
3712 permissions
3713 .env
3714 .as_ref()
3715 .map(serialize_pattern_permissions)
3716 .unwrap_or(wire::PatternPermissionScope::PermissionMode(
3717 wire::PermissionMode::Allow,
3718 )),
3719 ),
3720 binding: Some(
3721 permissions
3722 .binding
3723 .as_ref()
3724 .map(serialize_pattern_permissions)
3725 .unwrap_or(wire::PatternPermissionScope::PermissionMode(
3726 wire::PermissionMode::Allow,
3727 )),
3728 ),
3729 }
3730}
3731
3732fn default_permissions_policy() -> wire::PermissionsPolicy {
3737 wire::PermissionsPolicy {
3738 fs: Some(wire::FsPermissionScope::PermissionMode(
3739 wire::PermissionMode::Allow,
3740 )),
3741 network: Some(default_network_egress_scope()),
3742 child_process: Some(wire::PatternPermissionScope::PermissionMode(
3743 wire::PermissionMode::Allow,
3744 )),
3745 process: Some(wire::PatternPermissionScope::PermissionMode(
3746 wire::PermissionMode::Allow,
3747 )),
3748 env: Some(wire::PatternPermissionScope::PermissionMode(
3749 wire::PermissionMode::Allow,
3750 )),
3751 binding: Some(wire::PatternPermissionScope::PermissionMode(
3752 wire::PermissionMode::Allow,
3753 )),
3754 }
3755}
3756
3757fn serialize_fs_permissions(permissions: &crate::config::FsPermissions) -> wire::FsPermissionScope {
3758 match permissions {
3759 crate::config::FsPermissions::Mode(mode) => {
3760 wire::FsPermissionScope::PermissionMode(serialize_permission_mode(*mode))
3761 }
3762 crate::config::FsPermissions::Rules(rules) => {
3763 wire::FsPermissionScope::FsPermissionRuleSet(wire::FsPermissionRuleSet {
3764 default: rules.default.map(serialize_permission_mode),
3765 rules: rules
3766 .rules
3767 .iter()
3768 .map(|rule| wire::FsPermissionRule {
3769 mode: serialize_permission_mode(rule.mode),
3770 operations: operation_wildcard_if_omitted(&rule.operations),
3771 paths: resource_wildcard_if_omitted(&rule.paths),
3772 })
3773 .collect(),
3774 })
3775 }
3776 }
3777}
3778
3779fn serialize_pattern_permissions(
3780 permissions: &crate::config::PatternPermissions,
3781) -> wire::PatternPermissionScope {
3782 match permissions {
3783 crate::config::PatternPermissions::Mode(mode) => {
3784 wire::PatternPermissionScope::PermissionMode(serialize_permission_mode(*mode))
3785 }
3786 crate::config::PatternPermissions::Rules(rules) => {
3787 wire::PatternPermissionScope::PatternPermissionRuleSet(wire::PatternPermissionRuleSet {
3788 default: rules.default.map(serialize_permission_mode),
3789 rules: rules
3790 .rules
3791 .iter()
3792 .map(|rule| wire::PatternPermissionRule {
3793 mode: serialize_permission_mode(rule.mode),
3794 operations: operation_wildcard_if_omitted(&rule.operations),
3795 patterns: resource_wildcard_if_omitted(&rule.patterns),
3796 })
3797 .collect(),
3798 })
3799 }
3800 }
3801}
3802
3803fn serialize_permission_mode(mode: crate::config::PermissionMode) -> wire::PermissionMode {
3804 match mode {
3805 crate::config::PermissionMode::Allow => wire::PermissionMode::Allow,
3806 crate::config::PermissionMode::Deny => wire::PermissionMode::Deny,
3807 }
3808}
3809
3810fn json_utf8(value: &serde_json::Value, context: &str) -> Result<String, ClientError> {
3811 serde_json::to_string(value)
3812 .map_err(|error| ClientError::Sidecar(format!("failed to serialize {context}: {error}")))
3813}
3814
3815fn operation_wildcard_if_omitted(values: &Option<Vec<String>>) -> Vec<String> {
3816 values.clone().unwrap_or_else(|| vec!["*".to_string()])
3817}
3818
3819fn resource_wildcard_if_omitted(values: &Option<Vec<String>>) -> Vec<String> {
3820 values.clone().unwrap_or_else(|| vec!["**".to_string()])
3821}
3822
3823fn wire_ownership_vm_id(ownership: &wire::OwnershipScope) -> Option<&str> {
3825 match ownership {
3826 wire::OwnershipScope::VmOwnership(ownership) => Some(ownership.vm_id.as_str()),
3827 wire::OwnershipScope::ConnectionOwnership(_)
3828 | wire::OwnershipScope::SessionOwnership(_) => None,
3829 }
3830}
3831
3832fn rejected_to_error(rejected: wire::RejectedResponse) -> ClientError {
3834 ClientError::Kernel {
3835 code: rejected.code,
3836 message: rejected.message,
3837 }
3838}
3839
3840#[cfg(test)]
3841mod tests {
3842 use super::{
3843 abort_tracked_task, default_permissions_policy, permissions_policy,
3844 serialize_create_vm_config_for_sidecar, serialize_root_filesystem_config_for_sidecar,
3845 JoinHandle,
3846 };
3847 use crate::config::{
3848 AgentOsConfig, AgentOsLimits, FsPermissionRule, FsPermissions, HttpLimits, JsRuntimeLimits,
3849 MountPlugin, PatternPermissions, PermissionMode, Permissions, PythonLimits, ResourceLimits,
3850 RootFilesystemConfig, RootFilesystemKind, RootFilesystemMode, RootLowerInput,
3851 RulePermissions, ToolLimits, WasmLimits,
3852 };
3853 use crate::fs::{
3854 DirEntryType, FilesystemEntry, FilesystemEntryEncoding, FilesystemSnapshotEntries,
3855 FilesystemSnapshotExport, RootSnapshotExport, SnapshotExportKind,
3856 };
3857 use agentos_sidecar_client::wire::{
3858 FsPermissionScope, PatternPermissionScope, PermissionMode as WirePermissionMode,
3859 };
3860 use agentos_vm_config::{
3861 RootFilesystemEntryKind, RootFilesystemLowerDescriptor,
3862 RootFilesystemMode as ConfigRootFilesystemMode,
3863 };
3864
3865 #[tokio::test]
3877 async fn abort_tracked_task_aborts_and_clears_the_handle() {
3878 let slot: parking_lot::Mutex<Option<JoinHandle<()>>> = parking_lot::Mutex::new(None);
3880 assert!(
3881 slot.lock().is_none(),
3882 "pump slot starts empty like AgentOsInner"
3883 );
3884
3885 let task = tokio::spawn(async {
3886 loop {
3887 tokio::time::sleep(std::time::Duration::from_secs(3600)).await;
3888 }
3889 });
3890 let abort_handle = task.abort_handle();
3891 *slot.lock() = Some(task);
3893 assert!(
3894 slot.lock().is_some(),
3895 "spawning the pump must populate the tracked handle"
3896 );
3897
3898 assert!(!abort_handle.is_finished(), "pump task should start alive");
3899
3900 abort_tracked_task(&slot);
3901
3902 assert!(
3903 slot.lock().is_none(),
3904 "tracked handle must be taken on abort"
3905 );
3906
3907 for _ in 0..100 {
3909 if abort_handle.is_finished() {
3910 break;
3911 }
3912 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
3913 }
3914 assert!(
3915 abort_handle.is_finished(),
3916 "pump task must be aborted on shutdown"
3917 );
3918 }
3919
3920 #[test]
3921 fn permissions_policy_defaults_to_default_policy_when_unset() {
3922 assert_eq!(
3923 permissions_policy(&AgentOsConfig::default()),
3924 default_permissions_policy()
3925 );
3926 }
3927
3928 #[test]
3929 fn default_network_egress_is_llm_allowlist_not_allow_all() {
3930 let policy = permissions_policy(&AgentOsConfig::default());
3931
3932 assert_eq!(
3934 policy.child_process,
3935 Some(PatternPermissionScope::PermissionMode(
3936 WirePermissionMode::Allow
3937 ))
3938 );
3939
3940 let Some(PatternPermissionScope::PatternPermissionRuleSet(rules)) = policy.network else {
3943 panic!("expected default network egress to be a rule set, not allow-all");
3944 };
3945 assert_eq!(rules.default, Some(WirePermissionMode::Deny));
3946 assert_eq!(rules.rules.len(), 1);
3947 assert_eq!(rules.rules[0].mode, WirePermissionMode::Allow);
3948 let patterns = &rules.rules[0].patterns;
3949 assert!(patterns.contains(&"dns://api.anthropic.com".to_string()));
3950 assert!(patterns.contains(&"tcp://api.anthropic.com:*".to_string()));
3951 assert!(patterns.contains(&"dns://api.openai.com".to_string()));
3952 assert!(patterns.contains(&"dns://generativelanguage.googleapis.com".to_string()));
3953 assert!(patterns.contains(&"dns://openrouter.ai".to_string()));
3954 }
3955
3956 #[test]
3957 fn permissions_policy_preserves_configured_denies_and_allows_omitted_domains() {
3958 let policy = permissions_policy(&AgentOsConfig {
3959 permissions: Some(Permissions {
3960 network: Some(PatternPermissions::Mode(PermissionMode::Deny)),
3961 ..Default::default()
3962 }),
3963 ..Default::default()
3964 });
3965
3966 assert_eq!(
3967 policy.network,
3968 Some(PatternPermissionScope::PermissionMode(
3969 WirePermissionMode::Deny
3970 ))
3971 );
3972 assert_eq!(
3973 policy.child_process,
3974 Some(PatternPermissionScope::PermissionMode(
3975 WirePermissionMode::Allow
3976 ))
3977 );
3978 }
3979
3980 #[test]
3981 fn permissions_policy_expands_omitted_rule_fields_to_domain_wildcards() {
3982 let policy = permissions_policy(&AgentOsConfig {
3983 permissions: Some(Permissions {
3984 fs: Some(FsPermissions::Rules(RulePermissions {
3985 default: Some(PermissionMode::Deny),
3986 rules: vec![FsPermissionRule {
3987 mode: PermissionMode::Allow,
3988 operations: None,
3989 paths: Some(vec!["/workspace/**".to_string()]),
3990 }],
3991 })),
3992 ..Default::default()
3993 }),
3994 ..Default::default()
3995 });
3996
3997 let Some(FsPermissionScope::FsPermissionRuleSet(rules)) = policy.fs else {
3998 panic!("expected fs rule set");
3999 };
4000 assert_eq!(rules.default, Some(WirePermissionMode::Deny));
4001 assert_eq!(rules.rules[0].operations, vec!["*"]);
4002 assert_eq!(rules.rules[0].paths, vec!["/workspace/**"]);
4003
4004 let policy = permissions_policy(&AgentOsConfig {
4005 permissions: Some(Permissions {
4006 network: Some(PatternPermissions::Rules(RulePermissions {
4007 default: Some(PermissionMode::Allow),
4008 rules: vec![crate::config::PatternPermissionRule {
4009 mode: PermissionMode::Deny,
4010 operations: None,
4011 patterns: None,
4012 }],
4013 })),
4014 ..Default::default()
4015 }),
4016 ..Default::default()
4017 });
4018
4019 let Some(PatternPermissionScope::PatternPermissionRuleSet(rules)) = policy.network else {
4020 panic!("expected network rule set");
4021 };
4022 assert_eq!(rules.default, Some(WirePermissionMode::Allow));
4023 assert_eq!(rules.rules[0].operations, vec!["*"]);
4024 assert_eq!(rules.rules[0].patterns, vec!["**"]);
4025 }
4026
4027 #[test]
4028 fn root_filesystem_serializer_preserves_configured_descriptor() {
4029 let (descriptor, native_root) =
4030 serialize_root_filesystem_config_for_sidecar(&RootFilesystemConfig {
4031 mode: Some(RootFilesystemMode::ReadOnly),
4032 disable_default_base_layer: true,
4033 lowers: vec![
4034 RootLowerInput::BundledBaseFilesystem,
4035 RootLowerInput::SnapshotExport(RootSnapshotExport {
4036 kind: SnapshotExportKind::SnapshotExport,
4037 source: FilesystemSnapshotExport {
4038 format: "agentos-filesystem-snapshot-v1".to_string(),
4039 filesystem: FilesystemSnapshotEntries {
4040 entries: vec![
4041 FilesystemEntry {
4042 path: "/bin/run".to_string(),
4043 entry_type: DirEntryType::File,
4044 mode: "0755".to_string(),
4045 uid: 1000,
4046 gid: 1000,
4047 content: Some("#!/bin/sh".to_string()),
4048 encoding: Some(FilesystemEntryEncoding::Utf8),
4049 target: None,
4050 },
4051 FilesystemEntry {
4052 path: "/link".to_string(),
4053 entry_type: DirEntryType::Symlink,
4054 mode: "0777".to_string(),
4055 uid: 0,
4056 gid: 0,
4057 content: None,
4058 encoding: None,
4059 target: Some("/bin/run".to_string()),
4060 },
4061 ],
4062 },
4063 },
4064 }),
4065 ],
4066 ..Default::default()
4067 })
4068 .expect("serialize root filesystem");
4069
4070 assert!(native_root.is_none());
4071 assert_eq!(descriptor.mode, ConfigRootFilesystemMode::ReadOnly);
4072 assert!(descriptor.disable_default_base_layer);
4073 assert_eq!(descriptor.bootstrap_entries, Vec::new());
4074 assert!(matches!(
4075 descriptor.lowers[0],
4076 RootFilesystemLowerDescriptor::BundledBaseFilesystem
4077 ));
4078
4079 let RootFilesystemLowerDescriptor::Snapshot { entries } = &descriptor.lowers[1] else {
4080 panic!("expected snapshot lower");
4081 };
4082 assert_eq!(entries[0].path, "/bin/run");
4083 assert_eq!(entries[0].kind, RootFilesystemEntryKind::File);
4084 assert_eq!(entries[0].mode, Some(0o755));
4085 assert!(entries[0].executable);
4086 assert_eq!(entries[1].kind, RootFilesystemEntryKind::Symlink);
4087 assert_eq!(entries[1].target.as_deref(), Some("/bin/run"));
4088 }
4089
4090 #[test]
4091 fn create_vm_config_preserves_native_root_config() {
4092 let config = serialize_create_vm_config_for_sidecar(&AgentOsConfig {
4093 root_filesystem: RootFilesystemConfig {
4094 kind: RootFilesystemKind::Native,
4095 mode: Some(RootFilesystemMode::ReadOnly),
4096 native_plugin: Some(MountPlugin {
4097 id: "sqlite_vfs".to_string(),
4098 config: Some(serde_json::json!({
4099 "databasePath": "/tmp/agentos-root.sqlite"
4100 })),
4101 }),
4102 ..Default::default()
4103 },
4104 ..Default::default()
4105 })
4106 .expect("serialize create VM config");
4107 let native_root = config.native_root.expect("native root config");
4108
4109 assert_eq!(native_root.plugin.id, "sqlite_vfs");
4110 assert_eq!(
4111 native_root.plugin.config,
4112 serde_json::json!({ "databasePath": "/tmp/agentos-root.sqlite" })
4113 );
4114 assert!(native_root.read_only);
4115 }
4116
4117 #[test]
4118 fn create_vm_config_preserves_typed_limits() {
4119 let config = serialize_create_vm_config_for_sidecar(&AgentOsConfig {
4120 limits: Some(AgentOsLimits {
4121 resources: Some(ResourceLimits {
4122 max_processes: Some(7),
4123 max_filesystem_bytes: Some(4096),
4124 ..Default::default()
4125 }),
4126 http: Some(HttpLimits {
4127 max_fetch_response_bytes: Some(1024),
4128 }),
4129 tools: Some(ToolLimits {
4130 default_tool_timeout_ms: Some(500),
4131 max_registered_tools_per_vm: Some(12),
4132 ..Default::default()
4133 }),
4134 js_runtime: Some(JsRuntimeLimits {
4135 v8_heap_limit_mb: Some(64),
4136 sync_rpc_wait_timeout_ms: Some(2_000),
4137 cpu_time_limit_ms: Some(30_000),
4138 wall_clock_limit_ms: Some(0),
4139 import_cache_materialize_timeout_ms: Some(30_000),
4140 ..Default::default()
4141 }),
4142 python: Some(PythonLimits {
4143 max_old_space_mb: Some(256),
4144 ..Default::default()
4145 }),
4146 wasm: Some(WasmLimits {
4147 prewarm_timeout_ms: Some(30_000),
4148 runner_heap_limit_mb: Some(2_048),
4149 ..Default::default()
4150 }),
4151 ..Default::default()
4152 }),
4153 ..Default::default()
4154 })
4155 .expect("serialize create VM config");
4156 let limits = config.limits.expect("limits config");
4157
4158 let resources = limits.resources.expect("resource limits");
4159 assert_eq!(resources.max_processes, Some(7));
4160 assert_eq!(resources.max_filesystem_bytes, Some(4096));
4161 assert_eq!(
4162 limits.http.expect("http limits").max_fetch_response_bytes,
4163 Some(1024)
4164 );
4165 assert_eq!(
4166 limits
4167 .tools
4168 .as_ref()
4169 .expect("tool limits")
4170 .default_tool_timeout_ms,
4171 Some(500)
4172 );
4173 assert_eq!(
4174 limits
4175 .tools
4176 .expect("tool limits")
4177 .max_registered_tools_per_vm,
4178 Some(12)
4179 );
4180 assert_eq!(
4181 limits
4182 .js_runtime
4183 .as_ref()
4184 .expect("js runtime limits")
4185 .v8_heap_limit_mb,
4186 Some(64)
4187 );
4188 let js_runtime = limits.js_runtime.expect("js runtime limits");
4189 assert_eq!(js_runtime.sync_rpc_wait_timeout_ms, Some(2_000));
4190 assert_eq!(js_runtime.cpu_time_limit_ms, Some(30_000));
4191 assert_eq!(js_runtime.wall_clock_limit_ms, Some(0));
4192 assert_eq!(js_runtime.import_cache_materialize_timeout_ms, Some(30_000));
4193 assert_eq!(
4194 limits.python.expect("python limits").max_old_space_mb,
4195 Some(256)
4196 );
4197 let wasm = limits.wasm.expect("wasm limits");
4198 assert_eq!(wasm.prewarm_timeout_ms, Some(30_000));
4199 assert_eq!(wasm.runner_heap_limit_mb, Some(2_048));
4200 }
4201}