Skip to main content

agentos_client/
agent_os.rs

1//! The `AgentOs` struct (all fields from ADR-001 ยง3), the `create` builder, and the `shutdown`
2//! (dispose) teardown.
3//!
4//! `AgentOs` is `Arc`-cloneable; all interior state lives behind concurrent maps / atomics /
5//! channels so `&self` methods never need an outer lock. Module files add only `impl AgentOs` blocks
6//! and never introduce new struct fields.
7
8use 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
48// ---------------------------------------------------------------------------
49// Registry entries
50// ---------------------------------------------------------------------------
51
52/// An SDK-spawned process (TS `_processes` value). Keyed by user-facing pid.
53pub(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    /// Seeded `None`; the already-exited branch fires immediately once it holds `Some(code)`.
59    pub exit_tx: watch::Sender<Option<i32>>,
60    /// The sidecar-side process id used on the wire.
61    pub process_id: String,
62    /// The kernel pid returned by the `Execute` response, seeded once the spawn lands. The TS native
63    /// path builds `displayPidByKernelPid` from this so `all_processes`/`process_tree` report the
64    /// public spawn pid (the map key) for the spawned root, not the raw kernel pid.
65    pub kernel_pid: watch::Sender<Option<u32>>,
66}
67
68/// A PTY-backed shell (TS `_shells` value). Keyed by synthetic `shell-N` id.
69///
70/// `data_tx` carries stdout only, matching TS where the kernel handle's `onData` is fed exclusively
71/// by `stdoutHandlers`. `stderr_tx` is the dedicated stderr channel that backs the `on_stderr` option
72/// and `on_shell_stderr`, matching TS where stderr reaches the host only through `stderrHandlers`.
73pub(crate) struct ShellEntry {
74    pub pid: u32,
75    pub data_tx: broadcast::Sender<Vec<u8>>,
76    pub stderr_tx: broadcast::Sender<Vec<u8>>,
77    /// The sidecar-side process id used on the wire.
78    pub process_id: String,
79    /// Spawn-readiness gate. Seeded `false`; flips to `true` once the background `Execute` request is
80    /// acked. TS `openShell` is fully synchronous so `writeShell` always addresses a live spawn; the
81    /// Rust wire spawn is async, so `write_shell`/`close_shell` await this gate before issuing their
82    /// wire request to preserve the deterministic ordering and avoid dropping early input.
83    pub spawned_tx: watch::Sender<bool>,
84}
85
86/// A connected ACP terminal process and its output fan-out task.
87pub(crate) struct AcpTerminalEntry {
88    pub exit_task: JoinHandle<()>,
89}
90
91/// Mutable output state of a host-request ACP terminal (mirrors the TS `AcpTerminalEntry`
92/// `output` / `truncated` accumulation behavior).
93pub(crate) struct HostAcpTerminalOutput {
94    /// Accumulated UTF-8 terminal output (stdout + stderr interleaved, like the TS handle).
95    pub buffer: String,
96    pub truncated: bool,
97    /// Byte limit; `output` is trimmed from the front once it exceeds this. Mirrors the TS
98    /// `outputByteLimit` (default 1 MiB).
99    pub output_byte_limit: usize,
100}
101
102/// A host-request ACP terminal created via `terminal/create` (mirrors the TS `_acpTerminals`
103/// value). Backed by a real PTY shell (`open_shell`); the background fan-out task accumulates
104/// output and records the exit code.
105pub(crate) struct HostAcpTerminal {
106    /// The backing shell id (`shell-N`) used for `terminal/write` / `terminal/resize` /
107    /// `terminal/kill`.
108    pub shell_id: String,
109    /// Shared output buffer updated by the fan-out task and read by `terminal/output`.
110    pub output: Arc<parking_lot::Mutex<HostAcpTerminalOutput>>,
111    /// Exit code once the process has exited (`None` while running). Mirrors `exitCode`.
112    pub exit_rx: watch::Receiver<Option<i32>>,
113}
114
115/// An ACP session (TS `_sessions` value). Keyed by ACP session id.
116pub(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    /// Pending prompt resolvers, for cancel prompt-fallback + abort-on-close.
128    ///
129    /// The resolver carries the intended [`JsonRpcResponse`], mirroring the TS resolver shape
130    /// `{ method, resolve: (response) => void }`. The cause (close vs cancel) decides the payload at
131    /// the abort/cancel site: abort-on-close resolves with the `-32000` `Session closed: <id>` error,
132    /// while prompt-cancel resolves with `{ result: { stopReason: "cancelled" } }`. The shape is NOT
133    /// re-derived from the method downstream.
134    pub pending_prompt_resolvers:
135        SccHashMap<i64, oneshot::Sender<crate::json_rpc::JsonRpcResponse>>,
136}
137
138// ---------------------------------------------------------------------------
139// AgentOs
140// ---------------------------------------------------------------------------
141
142/// The high-level client. Cheaply cloneable via `Arc`.
143#[derive(Clone)]
144pub struct AgentOs {
145    inner: Arc<AgentOsInner>,
146}
147
148pub(crate) struct AgentOsInner {
149    // Transport / connection / VM handle.
150    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    // Process registries.
157    pub(crate) process_registry_lock: parking_lot::Mutex<()>,
158    pub(crate) processes: SccHashMap<u32, ProcessEntry>,
159    /// Wire `process_id` allocator for `exec` (the kernel-process view). Distinct from the
160    /// spawn synthetic-pid space so an `exec` call never perturbs the observable `spawn` pid sequence
161    /// (TS `nextSyntheticPid` is advanced only by `spawn`, never by `exec`).
162    pub(crate) process_counter: AtomicU64,
163    /// Synthetic display-pid allocator for `spawn` (TS `nextSyntheticPid`, seeded at
164    /// [`crate::process::SYNTHETIC_PID_BASE`]). The first spawned process gets `SYNTHETIC_PID_BASE`.
165    pub(crate) synthetic_pid_counter: AtomicU64,
166    pub(crate) observed_process_time_lock: parking_lot::Mutex<()>,
167    /// First-observed start time (epoch ms) per `"<process_id>:<kernel_pid>"`, mirroring TS
168    /// `observedProcessStartTimes`. A process keeps the timestamp first seen in `all_processes` across
169    /// later calls instead of advancing on every snapshot.
170    pub(crate) observed_process_start_times: SccHashMap<String, f64>,
171    /// First-observed exit time (epoch ms) per SDK-spawned wire `process_id`, mirroring TS
172    /// `tracked.exitTime` (set once when the process is first seen exited).
173    pub(crate) observed_process_exit_times: SccHashMap<String, f64>,
174
175    // Shell registries.
176    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    /// Host-request ACP terminals created via `terminal/create` (TS `_acpTerminals`). Keyed by the
183    /// `acp-terminal-N` id the agent uses in subsequent `terminal/*` calls.
184    pub(crate) host_acp_terminals: SccHashMap<String, HostAcpTerminal>,
185    /// Monotonic counter for the `acp-terminal-N` ids (TS `_acpTerminalCounter`).
186    pub(crate) host_acp_terminal_counter: AtomicU64,
187
188    // Session registries.
189    pub(crate) sessions: SccHashMap<String, SessionEntry>,
190    /// Bounded ordered set (cap [`crate::CLOSED_SESSION_ID_RETENTION_LIMIT`]) for close idempotence.
191    pub(crate) closed_session_ids: parking_lot::Mutex<VecDeque<String>>,
192    /// Session ids with an in-flight close in progress. Mirrors TS `_sessionClosePromises`: because
193    /// `close_session` runs the actual close on a detached task, this set keeps the id "known" during
194    /// the window between removal from `sessions` and insertion into `closed_session_ids`, so a second
195    /// `close_session` (or close-after-destroy) does not spuriously throw `SessionNotFound`.
196    pub(crate) closing_session_ids: SccHashSet<String>,
197
198    // Cron.
199    pub(crate) cron: Arc<CronManager>,
200
201    // Config / lifecycle.
202    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    /// The sole public VM entry point. Processes software, spawns/authenticates the sidecar, creates
211    /// the VM, waits for ready (10s), configures it, takes a lease, and constructs the cron manager
212    /// (default [`crate::config::TimerScheduleDriver`]).
213    pub async fn create(options: AgentOsConfig) -> Result<AgentOs, ClientError> {
214        let config = Arc::new(options);
215
216        // 1. Resolve the sidecar handle (shared "default" pool unless configured otherwise) and
217        //    establish/reuse its shared process + authenticated connection. A shared sidecar hosts
218        //    multiple VMs in one process, each opening its own session + VM below.
219        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        // 2. Open a session for this VM (connection scope) on the shared connection.
230        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        // 3. Subscribe to events BEFORE CreateVm so the `ready` lifecycle event cannot be missed.
280        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        // 4. Create the VM (session scope).
290        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        // 5. Wait for the VM to reach `ready` (bounded by VM_READY_TIMEOUT_MS).
344        wait_for_vm_ready(&mut events, &vm_id, crate::VM_READY_TIMEOUT_MS).await?;
345
346        // Resolve software packages to host roots (port of TS `processSoftware` for the
347        // ConfigureVm descriptors). Each `package` is resolved under `module_access_cwd/node_modules`;
348        // an unresolvable package is an explicit error rather than a silent no-op. Wasm command
349        // packages additionally become `/__secure_exec/commands/{index}/` mounts so the sidecar can
350        // discover and resolve guest commands.
351        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        // Native plugin mounts configured on the client, combined with the wasm command-dir mounts.
359        let mut mounts = serialize_mounts(&config)?;
360        mounts.extend(command_mounts);
361
362        // 6. Configure the VM (vm scope).
363        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        // 6b. Register host tool kits (if any): forward each tool definition via `register_host_callbacks`,
418        //     record the host execute callbacks in the per-VM registry, and install the shared
419        //     host-callback that routes guest tool calls back to the host by VM.
420        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        // 7. Lease this VM on the (possibly shared) sidecar, build cron, and assemble the client.
504        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        // Register the permission router and callback unconditionally (unlike `host_callback`,
551        // which is gated on configured tool kits): any agent session can raise a permission
552        // request. Re-registering on a shared transport replaces an identical stateless callback,
553        // same as the `host_callback` pattern.
554        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    /// Dispose the VM (= TS `dispose`). Teardown order:
565    /// 1. cron dispose
566    /// 2. close all sessions (swallow errors)
567    /// 3. kill all shells + snapshot pending exits
568    /// 4. kill all ACP terminals
569    /// 5. drain tracked shell-exit tasks (two-phase, bounded by
570    ///    [`crate::SHELL_DISPOSE_TIMEOUT_MS`])
571    /// 6. unregister the sidecar event listener
572    /// 7. release the lease (or tear down the transport)
573    ///
574    /// Idempotent (guarded by `disposed`).
575    pub async fn shutdown(&self) -> Result<(), ClientError> {
576        // Idempotent: only the first caller runs teardown.
577        if self.inner.disposed.swap(true, Ordering::SeqCst) {
578            return Ok(());
579        }
580
581        // 1. Cron dispose (cancel armed timers + tear down the driver).
582        self.inner.cron.dispose();
583
584        // 2-5. Best-effort drain tracked shell and terminal tasks before the VM is disposed, bounded
585        //      by SHELL_DISPOSE_TIMEOUT_MS so late output cannot race a closed transport.
586        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        // Tear down host-request ACP terminals (`terminal/create`). Close the backing shell, which
629        // sends SIGTERM, removes the shell entry, and ends the fan-out/exit task; the task itself is
630        // tracked in `pending_shell_exits` above and drained with the other shell exit tasks.
631        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        // 6-7. Release this VM (DisposeVm best-effort) and its lease. The transport is shared across
656        //      VMs on the same sidecar, so it is only torn down when this was the last VM (matching
657        //      the TS lease/shared-sidecar lifecycle); otherwise sibling VMs keep using it.
658        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    // --- internal accessors used by sibling impl blocks ---
691
692    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    /// The (possibly shared) sidecar handle backing this VM. Public for parity with TS
721    /// `AgentOs.sidecar` (e.g. `describe()` reports `active_vm_count` across VMs sharing a pool).
722    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
811/// Convert a sidecar's client-side placement into the wire `SidecarPlacement` for OpenSession.
812fn 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        // 0.3: the Node builtin allow-list moved from ConfigureVmRequest to
864        // VM creation. `None` => engine default allow-list; `Some([..])` =>
865        // exactly those (`Some([])` denies all). Platform/module-resolution
866        // keep their engine defaults (full Node emulation), matching prior
867        // behavior where Agent OS only ever constrained the builtin allow-list.
868        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                // Agent SDK snapshotting is driven by the TypeScript client
874                // (`packages/core` resolves the per-agent `dist/sdk-snapshot.js`
875                // bundle). The Rust client does not resolve npm package bundles, so
876                // it forwards no snapshot. TODO: expose a snapshot bundle input on
877                // the Rust client config for parity if a Rust consumer needs it.
878                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
1020/// Hosts the VM may reach by default (egress). The default network policy is an
1021/// allowlist of the common hosted LLM provider API endpoints so the standard
1022/// agent quickstart works with zero network configuration, while still matching
1023/// the Workers-style default-deny egress model: every other host is denied
1024/// unless the client widens the `network` permission. Clients opt out by
1025/// configuring `network` explicitly (e.g. `{ network: "allow" }`).
1026const DEFAULT_EGRESS_HOSTS: &[&str] = &[
1027    "api.anthropic.com",
1028    "api.openai.com",
1029    "generativelanguage.googleapis.com",
1030    "openrouter.ai",
1031];
1032
1033/// Resource patterns for the default egress allowlist. Network permission
1034/// resources are `dns://<host>` for name resolution and `tcp://<host>:<port>`
1035/// for the connection itself, so each allowed host needs both forms.
1036fn 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
1043/// vm_config variant of the default egress allowlist (deny-by-default rule set).
1044fn 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
1055/// Wire variant of the default egress allowlist (deny-by-default rule set).
1056fn 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
1128/// Default permission policy when the client supplies no `permissions`:
1129/// allow-all for fs/childProcess/process/env/binding (the VM is itself the
1130/// isolation boundary), with network egress restricted to the default LLM
1131/// allowlist (see [`default_network_egress_scope_config`]).
1132fn 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
1210/// Await the `ready` VM lifecycle event for `vm_id`, bounded by `timeout_ms`.
1211async 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
1248/// Process-global per-VM host-tool registry. The shared transport's single host-callback routes to
1249/// the right VM's toolkits by frame ownership.
1250static 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
1263/// Process-global map of vm id -> client inner, so the shared `permission_request` transport
1264/// callback can route a sidecar permission request to the owning client. `Weak` so the registry
1265/// never extends a client's lifetime; entries are removed in `shutdown`.
1266static 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
1272/// Process-global map of sidecar session -> Rust-host js_bridge callback.
1273///
1274/// Native root plugins can issue callbacks while `CreateVm` is still in flight, before the client
1275/// knows the generated VM id. Session ownership is already known by then and stays stable for the VM.
1276static 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
1395/// The transport callback that answers sidecar permission requests by routing them to the owning
1396/// client's `on_permission_request` subscribers. Mirrors TS `_handlePermissionSidecarRequest`.
1397fn 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
1484// ---------------------------------------------------------------------------
1485// ACP host-request dispatch (mirrors TS `_dispatchAcpSidecarRequest` ->
1486// `_handleSupportedAcpSidecarRequest`)
1487// ---------------------------------------------------------------------------
1488
1489/// The default `terminal/create` output cap (1 MiB), matching the TS reference.
1490const ACP_TERMINAL_DEFAULT_OUTPUT_BYTE_LIMIT: usize = 1_048_576;
1491
1492/// A JSON-RPC error raised while handling an ACP host request. Mirrors the TS `AcpDispatchError`.
1493struct 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            // Preserve the kernel errno code where one exists (e.g. ENOENT), surfaced through the
1521            // JSON-RPC `data.code`, while keeping a JSON-RPC internal-error envelope.
1522            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        // The filesystem methods return `anyhow::Result`; downcast to recover the kernel errno where
1533        // the underlying cause is a `ClientError::Kernel` (so e.g. ENOENT survives into `data.code`).
1534        match error.downcast::<ClientError>() {
1535            Ok(client_error) => client_error.into(),
1536            Err(error) => AcpDispatchError::new(-32603, error.to_string()),
1537        }
1538    }
1539}
1540
1541/// Decode the inbound JSON-RPC request, dispatch it to the matching VM operation, and serialize the
1542/// JSON-RPC response (success or error). Always returns a valid JSON-RPC response string; the
1543/// `id`/`error` shape mirrors `_dispatchAcpSidecarRequest`.
1544async 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
1595/// Resolve the `AgentOs` that owns the VM named in `ownership`, mirroring `route_permission_request`.
1596fn 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
1606/// Mirror of TS `_handleSupportedAcpSidecarRequest`: dispatch the JSON-RPC method to the matching VM
1607/// operation. Returns the JSON-RPC `result` value on success.
1608async 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, &params).await
1617        }
1618        "fs/read" | "fs/read_text_file" => {
1619            let agent = resolve_acp_agent(ownership)?;
1620            handle_acp_read_file(&agent, &params).await
1621        }
1622        "fs/write" | "fs/write_text_file" => {
1623            let agent = resolve_acp_agent(ownership)?;
1624            handle_acp_write_file(&agent, &params).await
1625        }
1626        "fs/readDir" | "fs/read_dir" => {
1627            let agent = resolve_acp_agent(ownership)?;
1628            handle_acp_read_dir(&agent, &params).await
1629        }
1630        "terminal/create" => {
1631            let agent = resolve_acp_agent(ownership)?;
1632            handle_acp_create_terminal(&agent, &params)
1633        }
1634        "terminal/write" => {
1635            let agent = resolve_acp_agent(ownership)?;
1636            handle_acp_write_terminal(&agent, &params)
1637        }
1638        "terminal/output" | "terminal/read" => {
1639            let agent = resolve_acp_agent(ownership)?;
1640            handle_acp_read_terminal(&agent, &params)
1641        }
1642        "terminal/wait_for_exit" | "terminal/waitForExit" => {
1643            let agent = resolve_acp_agent(ownership)?;
1644            handle_acp_wait_for_terminal_exit(&agent, &params).await
1645        }
1646        "terminal/kill" => {
1647            let agent = resolve_acp_agent(ownership)?;
1648            handle_acp_kill_terminal(&agent, &params)
1649        }
1650        "terminal/release" | "terminal/close" => {
1651            let agent = resolve_acp_agent(ownership)?;
1652            handle_acp_release_terminal(&agent, &params)
1653        }
1654        "terminal/resize" => {
1655            let agent = resolve_acp_agent(ownership)?;
1656            handle_acp_resize_terminal(&agent, &params)
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
1666// --- ACP host-request param helpers (mirror TS `_acpParams` / `_require*` / `_optional*`) ---
1667
1668fn 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
1759/// Parse the ACP `env` param, accepting either an object map or a `[{ name, value }]` array, matching
1760/// the TS `_optionalAcpEnvParam`.
1761fn 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
1820// --- fs/* handlers ---
1821
1822async 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
1910// --- session/request_permission handler ---
1911
1912async 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            // The host-request id is not available here as the permission key; use a generated key
1924            // scoped to the session so concurrent permission requests do not collide.
1925            permission_id: format!("acp-permission-{}", uuid::Uuid::new_v4()),
1926            params: Value::Object(params.clone()),
1927        },
1928    )
1929    .await;
1930
1931    // `reply: None` means the session/VM is gone or the request timed out -> cancelled outcome.
1932    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
1947/// Mirror of TS `_normalizeAcpPermissionOptionId`: pick the matching option id from the request's
1948/// `options`, falling back to the canonical id for the decision.
1949fn 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
1981/// Mirror of TS `_buildAcpPermissionResult`: produce `{ outcome: { outcome: "selected", optionId } }`.
1982fn 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
1993// --- terminal/* handlers ---
1994
1995fn 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    // Build the PTY shell. Both stdout and stderr are appended to the same output buffer, mirroring
2032    // the TS handle where `onData` and `onStderr` both append to `terminal.output`.
2033    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    // Both stdout and stderr are appended to the single combined output buffer inside
2051    // `acp_open_terminal`'s fan-out task (mirroring the TS handle's `onData`/`onStderr`).
2052    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        // Trim from the front to the limit, on a char boundary, matching the TS slice-to-limit
2093        // behavior (which trims to the last `limit` UTF-16 code units; bytes are an acceptable port).
2094        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            // Sender dropped (terminal released / VM disposed) without a recorded
2184            // exit code. Surface that as an abnormal exit instead of pretending
2185            // the terminal completed cleanly with exit 0.
2186            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    // The native shell API only exposes SIGTERM teardown via `close_shell`'s kill; the explicit
2200    // `signal` param is accepted for parity but the underlying kill is fixed to SIGTERM. The terminal
2201    // entry is retained (matching TS `kill`, which does not delete the terminal) so `terminal/output`
2202    // and `terminal/wait_for_exit` still work afterward.
2203    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 the process has not exited yet, kill it (TS releases by killing when `exitCode === null`).
2222    if terminal.exit_rx.borrow().is_none() {
2223        let _ = agent.acp_kill_terminal_shell(&terminal.shell_id);
2224    }
2225    // Closing the shell removes the registry entry and ends the fan-out/exit task naturally.
2226    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
2251/// Look up the backing shell id for a host-request terminal, or a JSON-RPC -32602 error.
2252fn 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
2262/// The transport callback that answers guest tool invocations by running the matching host tool.
2263fn 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
2293/// Run a single tool invocation against the per-VM host-tool registry, honoring the timeout. Mirrors
2294/// TS `handleHostCallback` (unknown-tool + timeout + error shapes).
2295async 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(&registry.tool_kits, toolkit_name),
2443            None => Ok(list_toolkits_payload(&registry.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(&registry.tool_kits)
2455        ));
2456    };
2457
2458    let Some(tool_name) = command.args.get(1) else {
2459        return describe_toolkit_payload(&registry.tool_kits, subcommand);
2460    };
2461    if is_help_flag(tool_name) {
2462        return describe_toolkit_payload(&registry.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(&registry.tool_kits, &toolkit.name);
2486    };
2487    if is_help_flag(tool_name) {
2488        return describe_toolkit_payload(&registry.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
3412/// A software package resolved to its host root, paired with the kind that decides how it is mounted.
3413struct ResolvedSoftware {
3414    descriptor: wire::SoftwareDescriptor,
3415    kind: SoftwareKind,
3416}
3417
3418/// Resolve `config.software` package inputs to host roots, each rooted at its host `node_modules`
3419/// directory under `module_access_cwd` (default `.`). An absolute `package` path bypasses the
3420/// `node_modules` prefix (via `Path::join` semantics), which is how wasm command directories are
3421/// passed directly. Mirrors the TS `processSoftware` mapping. An unresolvable package is an explicit
3422/// error, not a silent no-op.
3423fn 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
3454/// Build the `host_dir` mount descriptors that expose each wasm command directory at
3455/// `/__secure_exec/commands/{index}/` in the guest, so the sidecar's `discover_command_guest_paths` can
3456/// resolve guest commands. Indices are zero-padded so the sidecar's lexical sort preserves numeric
3457/// resolution priority past nine packages. Agent/tool packages are skipped here (they are not
3458/// command directories). Mirrors the TS `commandDirs` mount loop in `agent-os.ts`.
3459fn 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
3580/// Default permission policy (wire form) when the client supplies no
3581/// `permissions`: allow-all for fs/childProcess/process/env/binding, with network
3582/// egress restricted to the default LLM allowlist
3583/// (see [`default_network_egress_scope`]).
3584fn 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
3671/// Extract the `vm_id` from a generated ownership scope, if it is VM-scoped.
3672fn 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
3680/// Map a `Rejected` response into a [`ClientError::Kernel`] so the errno `code` survives.
3681fn 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        // fs/childProcess/process/env stay allow-all (the VM is the boundary).
3725        assert_eq!(
3726            policy.child_process,
3727            Some(PatternPermissionScope::PermissionMode(
3728                WirePermissionMode::Allow
3729            ))
3730        );
3731
3732        // Network egress is a deny-by-default allowlist of LLM provider hosts,
3733        // covering both DNS resolution and the TCP connection for each host.
3734        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}