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