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