Skip to main content

secure_exec_sidecar/
service.rs

1use crate::bridge::{build_mount_plugin_registry, MountPluginContext};
2pub(crate) use crate::execution::{
3    build_javascript_socket_path_context, canonical_signal_name, dispatch_loopback_http_request,
4    error_code, format_tcp_resource, ignore_stale_javascript_sync_rpc_response,
5    javascript_sync_rpc_arg_i32, javascript_sync_rpc_arg_str, javascript_sync_rpc_arg_u32,
6    javascript_sync_rpc_arg_u32_optional, javascript_sync_rpc_arg_u64,
7    javascript_sync_rpc_arg_u64_optional, javascript_sync_rpc_bytes_arg,
8    javascript_sync_rpc_bytes_value, javascript_sync_rpc_encoding, javascript_sync_rpc_error_code,
9    javascript_sync_rpc_option_bool, javascript_sync_rpc_option_u32, parse_signal,
10    sanitize_javascript_child_process_internal_bootstrap_env, service_javascript_sync_rpc,
11    vm_network_resource_counts, JavascriptSyncRpcServiceRequest, LoopbackHttpDispatchRequest,
12};
13use crate::extension::{
14    Extension, ExtensionBufferedProcessOutput, ExtensionContext, ExtensionFuture, ExtensionHost,
15    ExtensionSnapshot,
16};
17use crate::filesystem::guest_filesystem_call as filesystem_guest_filesystem_call;
18use crate::limits::DEFAULT_ACP_STDOUT_BUFFER_BYTE_LIMIT;
19use crate::protocol::{
20    AuthenticatedResponse, CloseStdinRequest, DisposeReason, EventFrame, EventPayload,
21    ExecuteRequest, ExtEnvelope, FsPermissionScope, GuestFilesystemCallRequest,
22    GuestFilesystemResultResponse, JavascriptChildProcessSpawnOptions,
23    JavascriptChildProcessSpawnRequest, KillProcessRequest, OpenSessionRequest, OwnershipScope,
24    PatternPermissionRule, PatternPermissionScope, PermissionMode, PermissionsPolicy,
25    ProcessKilledResponse, ProcessStartedResponse, ProtocolSchema, RejectedResponse, RequestFrame,
26    RequestId, RequestPayload, ResponseFrame, ResponsePayload, SessionOpenedResponse,
27    SidecarRequestFrame, SidecarRequestPayload, SidecarResponseFrame, SidecarResponsePayload,
28    SidecarResponseTracker, SidecarResponseTrackerError, SignalDispositionAction,
29    SignalHandlerRegistration, StdinClosedResponse, StdinWrittenResponse, VmLifecycleEvent,
30    VmLifecycleState, WriteStdinRequest,
31};
32use crate::state::{
33    ActiveExecutionEvent, BridgeError, ConnectionState, EventSinkTransport, JavascriptSocketFamily,
34    JavascriptSocketPathContext, ProcessEventEnvelope, SessionState, SharedBridge, SharedEventSink,
35    SharedSidecarRequestClient, SidecarRequestTransport, VmState, EXECUTION_DRIVER_NAME,
36};
37use crate::tools::register_host_callbacks;
38use crate::NativeSidecarBridge;
39use secure_exec_bridge::queue_tracker::{register_queue, QueueGauge, TrackedLimit};
40use secure_exec_bridge::{
41    CommandPermissionRequest, EnvironmentAccess, EnvironmentPermissionRequest, FilesystemAccess,
42    FilesystemPermissionRequest, LifecycleEventRecord, LifecycleState, LogLevel, LogRecord,
43    NetworkAccess, NetworkPermissionRequest, StructuredEventRecord,
44};
45use secure_exec_execution::{
46    JavascriptExecutionEngine, JavascriptExecutionError, JavascriptSyncRpcRequest,
47    PythonExecutionEngine, PythonExecutionError, WasmExecutionEngine, WasmExecutionError,
48};
49use secure_exec_kernel::kernel::KernelError;
50use secure_exec_kernel::mount_plugin::{FileSystemPluginRegistry, PluginError};
51use secure_exec_kernel::permissions::{
52    permission_glob_matches, CommandAccessRequest, EnvAccessRequest, EnvironmentOperation,
53    NetworkAccessRequest, NetworkOperation, PermissionDecision,
54};
55// root_fs types moved to crate::vm
56use secure_exec_kernel::vfs::VfsError;
57use serde::Deserialize;
58use serde_json::{json, Value};
59use std::collections::{BTreeMap, BTreeSet, VecDeque};
60use std::fmt;
61use std::fs;
62use std::os::unix::fs::PermissionsExt;
63use std::path::{Component, Path, PathBuf};
64use std::sync::{Arc, Mutex};
65use std::task::{Context, Poll, Waker};
66use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
67use tokio::sync::mpsc::{channel, Receiver, Sender};
68use tokio::time;
69
70// Constants and type aliases moved to crate::state
71
72const INTERNAL_JAVASCRIPT_ENTRYPOINT_ENV_KEYS: &[&str] =
73    &["AGENTOS_ENTRYPOINT", "AGENTOS_BOOTSTRAP_MODULE"];
74const INTERNAL_WASM_ENTRYPOINT_ENV_KEYS: &[&str] =
75    &["AGENTOS_WASM_MODULE_PATH", "AGENTOS_WASM_MODULE_BASE64"];
76const INTERNAL_PYTHON_ENTRYPOINT_ENV_PREFIXES: &[&str] = &["AGENTOS_PYTHON_"];
77pub(crate) const MAX_PROCESS_EVENT_QUEUE: usize = 10_000;
78pub(crate) const MAX_PENDING_SIDECAR_RESPONSES: usize = 10_000;
79pub(crate) const MAX_OUTBOUND_SIDECAR_REQUESTS: usize = 10_000;
80pub(crate) const MAX_COMPLETED_SIDECAR_RESPONSES: usize = 10_000;
81
82pub(crate) fn process_event_queue_overflow_error() -> SidecarError {
83    SidecarError::InvalidState(format!(
84        "process event queue exceeded {MAX_PROCESS_EVENT_QUEUE} pending events"
85    ))
86}
87
88fn sidecar_response_pending_overflow_error() -> SidecarError {
89    SidecarError::InvalidState(format!(
90        "sidecar response tracker exceeded {MAX_PENDING_SIDECAR_RESPONSES} pending responses"
91    ))
92}
93
94fn outbound_sidecar_request_queue_overflow_error() -> SidecarError {
95    SidecarError::InvalidState(format!(
96        "outbound sidecar request queue exceeded {MAX_OUTBOUND_SIDECAR_REQUESTS} pending requests"
97    ))
98}
99
100fn wire_protocol_error(error: crate::wire::ProtocolCodecError) -> SidecarError {
101    SidecarError::InvalidState(format!("invalid generated wire protocol frame: {error}"))
102}
103
104// NativeSidecarConfig, DispatchResult, SidecarError moved to crate::state
105pub use crate::state::{DispatchResult, NativeSidecarConfig, SidecarError};
106
107// SharedBridge struct and Clone impl moved to crate::state
108
109#[derive(Debug, Default, Deserialize)]
110struct LegacyJavascriptChildProcessSpawnOptions {
111    #[serde(default)]
112    cwd: Option<String>,
113    #[serde(default)]
114    env: BTreeMap<String, String>,
115    #[serde(default)]
116    input: Option<Value>,
117    #[serde(default)]
118    shell: bool,
119    #[serde(default)]
120    detached: bool,
121    #[serde(default)]
122    stdio: Vec<String>,
123    #[serde(default, rename = "maxBuffer")]
124    max_buffer: Option<usize>,
125    #[serde(default)]
126    timeout: Option<u64>,
127    #[serde(default, rename = "killSignal")]
128    kill_signal: Option<String>,
129}
130
131#[derive(Debug, Deserialize)]
132#[serde(rename_all = "camelCase")]
133struct JavascriptHttpLoopbackRequest {
134    process_id: String,
135    server_id: u64,
136    host: String,
137    port: u16,
138    request: String,
139}
140
141fn is_javascript_loopback_host(host: &str) -> bool {
142    host == "127.0.0.1" || host == "::1" || host.eq_ignore_ascii_case("localhost")
143}
144
145pub(crate) fn parse_javascript_child_process_spawn_request(
146    vm: &VmState,
147    args: &[Value],
148) -> Result<(JavascriptChildProcessSpawnRequest, Option<usize>), SidecarError> {
149    if let Some(value) = args.first().cloned() {
150        if let Ok(request) = serde_json::from_value::<JavascriptChildProcessSpawnRequest>(value) {
151            return Ok((request, None));
152        }
153    }
154
155    let command = javascript_sync_rpc_arg_str(args, 0, "child_process.spawn command")?.to_owned();
156    let raw_args = javascript_sync_rpc_arg_str(args, 1, "child_process.spawn args")?;
157    let raw_options = javascript_sync_rpc_arg_str(args, 2, "child_process.spawn options")?;
158
159    let parsed_args = serde_json::from_str::<Vec<String>>(raw_args).map_err(|error| {
160        SidecarError::InvalidState(format!("invalid child_process.spawn args payload: {error}"))
161    })?;
162    let parsed_options = serde_json::from_str::<LegacyJavascriptChildProcessSpawnOptions>(
163        raw_options,
164    )
165    .map_err(|error| {
166        SidecarError::InvalidState(format!(
167            "invalid child_process.spawn options payload: {error}"
168        ))
169    })?;
170
171    Ok((
172        JavascriptChildProcessSpawnRequest {
173            command,
174            args: parsed_args,
175            options: JavascriptChildProcessSpawnOptions {
176                cwd: parsed_options.cwd,
177                env: parsed_options.env,
178                internal_bootstrap_env: sanitize_javascript_child_process_internal_bootstrap_env(
179                    &vm.guest_env,
180                ),
181                input: parsed_options.input,
182                shell: parsed_options.shell,
183                detached: parsed_options.detached,
184                stdio: parsed_options.stdio,
185                timeout: parsed_options.timeout,
186                kill_signal: parsed_options.kill_signal,
187            },
188        },
189        parsed_options.max_buffer,
190    ))
191}
192
193impl<B> SharedBridge<B> {
194    fn new(bridge: B) -> Self {
195        Self {
196            inner: Arc::new(Mutex::new(bridge)),
197            permissions: Arc::new(Mutex::new(BTreeMap::new())),
198            #[cfg(test)]
199            set_vm_permissions_outcomes: Arc::new(Mutex::new(VecDeque::new())),
200        }
201    }
202}
203
204impl<B> SharedBridge<B>
205where
206    B: NativeSidecarBridge + Send + 'static,
207    BridgeError<B>: fmt::Debug + Send + Sync + 'static,
208{
209    pub(crate) fn with_mut<T>(
210        &self,
211        operation: impl FnOnce(&mut B) -> Result<T, BridgeError<B>>,
212    ) -> Result<T, SidecarError> {
213        let mut bridge = self.inner.lock().map_err(|_| {
214            SidecarError::Bridge(String::from("native sidecar bridge lock poisoned"))
215        })?;
216        operation(&mut bridge).map_err(|error| SidecarError::Bridge(format!("{error:?}")))
217    }
218
219    fn inspect<T>(&self, operation: impl FnOnce(&mut B) -> T) -> Result<T, SidecarError> {
220        let mut bridge = self.inner.lock().map_err(|_| {
221            SidecarError::Bridge(String::from("native sidecar bridge lock poisoned"))
222        })?;
223        Ok(operation(&mut bridge))
224    }
225
226    #[cfg(test)]
227    #[allow(dead_code)]
228    pub(crate) fn queue_set_vm_permissions_result(
229        &self,
230        result: Result<(), SidecarError>,
231    ) -> Result<(), SidecarError> {
232        let mut outcomes = self.set_vm_permissions_outcomes.lock().map_err(|_| {
233            SidecarError::Bridge(String::from(
234                "native sidecar test set_vm_permissions outcome lock poisoned",
235            ))
236        })?;
237        outcomes.push_back(result.err());
238        Ok(())
239    }
240
241    pub(crate) fn emit_lifecycle(
242        &self,
243        vm_id: &str,
244        state: LifecycleState,
245    ) -> Result<(), SidecarError> {
246        self.with_mut(|bridge| {
247            bridge.emit_lifecycle(LifecycleEventRecord {
248                vm_id: vm_id.to_owned(),
249                state,
250                detail: None,
251            })
252        })
253    }
254
255    pub(crate) fn emit_log(
256        &self,
257        vm_id: &str,
258        message: impl Into<String>,
259    ) -> Result<(), SidecarError> {
260        self.with_mut(|bridge| {
261            bridge.emit_log(LogRecord {
262                vm_id: vm_id.to_owned(),
263                level: LogLevel::Info,
264                message: message.into(),
265            })
266        })
267    }
268
269    pub(crate) fn filesystem_decision(
270        &self,
271        vm_id: &str,
272        path: &str,
273        access: FilesystemAccess,
274    ) -> PermissionDecision {
275        if let Some(decision) = self.static_permission_decision(
276            vm_id,
277            filesystem_permission_capability(access),
278            "fs",
279            Some(path),
280        ) {
281            return decision;
282        }
283        match self.with_mut(|bridge| {
284            bridge.check_filesystem_access(FilesystemPermissionRequest {
285                vm_id: vm_id.to_owned(),
286                path: path.to_owned(),
287                access,
288            })
289        }) {
290            Ok(decision) => map_bridge_permission(decision),
291            Err(error) => PermissionDecision::deny(error.to_string()),
292        }
293    }
294
295    pub(crate) fn command_decision(
296        &self,
297        vm_id: &str,
298        request: &CommandAccessRequest,
299    ) -> PermissionDecision {
300        if is_internal_runtime_command_request(request) {
301            return PermissionDecision::allow();
302        }
303        if let Some(decision) = self.static_permission_decision(
304            vm_id,
305            "child_process.spawn",
306            "child_process",
307            Some(&request.command),
308        ) {
309            return decision;
310        }
311        match self.with_mut(|bridge| {
312            bridge.check_command_execution(CommandPermissionRequest {
313                vm_id: vm_id.to_owned(),
314                command: request.command.clone(),
315                args: request.args.clone(),
316                cwd: request.cwd.clone(),
317                env: request.env.clone(),
318            })
319        }) {
320            Ok(decision) => map_bridge_permission(decision),
321            Err(error) => PermissionDecision::deny(error.to_string()),
322        }
323    }
324
325    pub(crate) fn environment_decision(
326        &self,
327        vm_id: &str,
328        request: &EnvAccessRequest,
329    ) -> PermissionDecision {
330        if let Some(decision) = self.static_permission_decision(
331            vm_id,
332            environment_permission_capability(request.op),
333            "env",
334            Some(&request.key),
335        ) {
336            return decision;
337        }
338        match self.with_mut(|bridge| {
339            bridge.check_environment_access(EnvironmentPermissionRequest {
340                vm_id: vm_id.to_owned(),
341                access: match request.op {
342                    EnvironmentOperation::Read => EnvironmentAccess::Read,
343                    EnvironmentOperation::Write => EnvironmentAccess::Write,
344                },
345                key: request.key.clone(),
346                value: request.value.clone(),
347            })
348        }) {
349            Ok(decision) => map_bridge_permission(decision),
350            Err(error) => PermissionDecision::deny(error.to_string()),
351        }
352    }
353
354    pub(crate) fn network_decision(
355        &self,
356        vm_id: &str,
357        request: &NetworkAccessRequest,
358    ) -> PermissionDecision {
359        if let Some(decision) = self.static_permission_decision(
360            vm_id,
361            network_permission_capability(request.op),
362            "network",
363            Some(&request.resource),
364        ) {
365            return decision;
366        }
367        match self.with_mut(|bridge| {
368            bridge.check_network_access(NetworkPermissionRequest {
369                vm_id: vm_id.to_owned(),
370                access: match request.op {
371                    NetworkOperation::Fetch => NetworkAccess::Fetch,
372                    NetworkOperation::Http => NetworkAccess::Http,
373                    NetworkOperation::Dns => NetworkAccess::Dns,
374                    NetworkOperation::Listen => NetworkAccess::Listen,
375                },
376                resource: request.resource.clone(),
377            })
378        }) {
379            Ok(decision) => map_bridge_permission(decision),
380            Err(error) => PermissionDecision::deny(error.to_string()),
381        }
382    }
383
384    pub(crate) fn require_network_access(
385        &self,
386        vm_id: &str,
387        op: NetworkOperation,
388        resource: impl Into<String>,
389    ) -> Result<(), SidecarError> {
390        let resource = resource.into();
391        let decision = self.network_decision(
392            vm_id,
393            &NetworkAccessRequest {
394                vm_id: vm_id.to_owned(),
395                op,
396                resource: resource.clone(),
397            },
398        );
399        if decision.allow {
400            return Ok(());
401        }
402
403        let message = match decision.reason.as_deref() {
404            Some(reason) => format!("EACCES: permission denied, {resource}: {reason}"),
405            None => format!("EACCES: permission denied, {resource}"),
406        };
407        Err(SidecarError::Execution(message))
408    }
409
410    pub(crate) fn set_vm_permissions(
411        &self,
412        vm_id: &str,
413        permissions: &PermissionsPolicy,
414    ) -> Result<(), SidecarError> {
415        #[cfg(test)]
416        {
417            let mut outcomes = self.set_vm_permissions_outcomes.lock().map_err(|_| {
418                SidecarError::Bridge(String::from(
419                    "native sidecar test set_vm_permissions outcome lock poisoned",
420                ))
421            })?;
422            if let Some(Some(error)) = outcomes.pop_front() {
423                return Err(error);
424            }
425        }
426
427        let mut stored = self.permissions.lock().map_err(|_| {
428            SidecarError::Bridge(String::from(
429                "native sidecar permission policy lock poisoned",
430            ))
431        })?;
432        stored.insert(vm_id.to_owned(), permissions.clone());
433        Ok(())
434    }
435
436    pub(crate) fn restore_vm_permissions_fail_closed(
437        &self,
438        vm_id: &str,
439        original_permissions: &PermissionsPolicy,
440        context: &str,
441        operation_error: &SidecarError,
442    ) -> Result<(), SidecarError> {
443        match self.set_vm_permissions(vm_id, original_permissions) {
444            Ok(()) => Ok(()),
445            Err(restore_error) => {
446                let deny_all = PermissionsPolicy::deny_all();
447                match self.set_vm_permissions(vm_id, &deny_all) {
448                    Ok(()) => Err(SidecarError::InvalidState(format!(
449                        "{context} failed: {operation_error}; restoring original permissions failed: {restore_error}; applied deny-all fallback"
450                    ))),
451                    Err(deny_all_error) => panic!(
452                        "{context} failed: {operation_error}; restoring original permissions failed: {restore_error}; deny-all fallback failed: {deny_all_error}"
453                    ),
454                }
455            }
456        }
457    }
458
459    pub(crate) fn clear_vm_permissions(&self, vm_id: &str) -> Result<(), SidecarError> {
460        let mut stored = self.permissions.lock().map_err(|_| {
461            SidecarError::Bridge(String::from(
462                "native sidecar permission policy lock poisoned",
463            ))
464        })?;
465        stored.remove(vm_id);
466        Ok(())
467    }
468
469    pub(crate) fn static_permission_decision(
470        &self,
471        vm_id: &str,
472        capability: &str,
473        domain: &str,
474        resource: Option<&str>,
475    ) -> Option<PermissionDecision> {
476        let stored = self.permissions.lock().ok()?;
477        let permissions = stored.get(vm_id)?;
478        let mode = evaluate_permissions_policy(permissions, domain, capability, resource);
479        Some(permission_mode_to_kernel_decision(mode, capability))
480    }
481}
482
483pub(crate) fn evaluate_permissions_policy(
484    permissions: &PermissionsPolicy,
485    domain: &str,
486    capability: &str,
487    resource: Option<&str>,
488) -> PermissionMode {
489    match domain {
490        "fs" => evaluate_fs_permission_scope(
491            permissions.fs.as_ref(),
492            capability_operation(capability, domain),
493            resource,
494        ),
495        "network" => evaluate_pattern_permission_scope(
496            permissions.network.as_ref(),
497            capability_operation(capability, domain),
498            resource,
499        ),
500        "child_process" => evaluate_pattern_permission_scope(
501            permissions.child_process.as_ref(),
502            capability_operation(capability, domain),
503            resource,
504        ),
505        "process" => evaluate_pattern_permission_scope(
506            permissions.process.as_ref(),
507            capability_operation(capability, domain),
508            resource,
509        ),
510        "env" => evaluate_pattern_permission_scope(
511            permissions.env.as_ref(),
512            capability_operation(capability, domain),
513            resource,
514        ),
515        "binding" => evaluate_pattern_permission_scope(
516            permissions.binding.as_ref(),
517            capability_operation(capability, domain),
518            resource,
519        ),
520        _ => PermissionMode::Deny,
521    }
522}
523
524fn evaluate_fs_permission_scope(
525    scope: Option<&FsPermissionScope>,
526    operation: &str,
527    resource: Option<&str>,
528) -> PermissionMode {
529    match scope {
530        Some(FsPermissionScope::PermissionMode(mode)) => mode.clone(),
531        Some(FsPermissionScope::FsPermissionRuleSet(rules)) => {
532            let mut mode = rules.default.clone().unwrap_or(PermissionMode::Deny);
533            for rule in &rules.rules {
534                if fs_rule_matches(rule, operation, resource) {
535                    mode = rule.mode.clone();
536                }
537            }
538            mode
539        }
540        None => PermissionMode::Deny,
541    }
542}
543
544fn evaluate_pattern_permission_scope(
545    scope: Option<&PatternPermissionScope>,
546    operation: &str,
547    resource: Option<&str>,
548) -> PermissionMode {
549    match scope {
550        Some(PatternPermissionScope::PermissionMode(mode)) => mode.clone(),
551        Some(PatternPermissionScope::PatternPermissionRuleSet(rules)) => {
552            let mut mode = rules.default.clone().unwrap_or(PermissionMode::Deny);
553            for rule in &rules.rules {
554                if pattern_rule_matches(rule, operation, resource) {
555                    mode = rule.mode.clone();
556                }
557            }
558            mode
559        }
560        None => PermissionMode::Deny,
561    }
562}
563
564fn fs_rule_matches(
565    rule: &crate::protocol::FsPermissionRule,
566    operation: &str,
567    resource: Option<&str>,
568) -> bool {
569    let operations_match = permission_operation_matches(&rule.operations, operation);
570    let paths_match = permission_resource_matches(&rule.paths, resource);
571    operations_match && paths_match
572}
573
574fn pattern_rule_matches(
575    rule: &PatternPermissionRule,
576    operation: &str,
577    resource: Option<&str>,
578) -> bool {
579    let operations_match = permission_operation_matches(&rule.operations, operation);
580    let patterns_match = permission_resource_matches(&rule.patterns, resource);
581    operations_match && patterns_match
582}
583
584fn permission_operation_matches(candidates: &[String], operation: &str) -> bool {
585    candidates
586        .iter()
587        .any(|candidate| candidate == "*" || candidate == operation)
588}
589
590fn permission_resource_matches(patterns: &[String], resource: Option<&str>) -> bool {
591    resource.is_some_and(|value| {
592        patterns
593            .iter()
594            .any(|pattern| permission_glob_matches(pattern, value))
595    })
596}
597
598pub(crate) fn validate_permissions_policy(
599    permissions: &PermissionsPolicy,
600) -> Result<(), SidecarError> {
601    if let Some(scope) = permissions.fs.as_ref() {
602        validate_fs_permission_scope("fs", scope)?;
603    }
604    if let Some(scope) = permissions.network.as_ref() {
605        validate_pattern_permission_scope("network", scope)?;
606    }
607    if let Some(scope) = permissions.child_process.as_ref() {
608        validate_pattern_permission_scope("child_process", scope)?;
609    }
610    if let Some(scope) = permissions.process.as_ref() {
611        validate_pattern_permission_scope("process", scope)?;
612    }
613    if let Some(scope) = permissions.env.as_ref() {
614        validate_pattern_permission_scope("env", scope)?;
615    }
616    if let Some(scope) = permissions.binding.as_ref() {
617        validate_pattern_permission_scope("binding", scope)?;
618    }
619    Ok(())
620}
621
622fn validate_fs_permission_scope(
623    domain: &str,
624    scope: &FsPermissionScope,
625) -> Result<(), SidecarError> {
626    let FsPermissionScope::FsPermissionRuleSet(rule_set) = scope else {
627        return Ok(());
628    };
629
630    for (index, rule) in rule_set.rules.iter().enumerate() {
631        validate_permission_rule_field(
632            &rule.operations,
633            &format!("{domain}.rules[{index}].operations"),
634        )?;
635        validate_permission_rule_field(&rule.paths, &format!("{domain}.rules[{index}].paths"))?;
636    }
637
638    Ok(())
639}
640
641fn validate_pattern_permission_scope(
642    domain: &str,
643    scope: &PatternPermissionScope,
644) -> Result<(), SidecarError> {
645    let PatternPermissionScope::PatternPermissionRuleSet(rule_set) = scope else {
646        return Ok(());
647    };
648
649    for (index, rule) in rule_set.rules.iter().enumerate() {
650        validate_permission_rule_field(
651            &rule.operations,
652            &format!("{domain}.rules[{index}].operations"),
653        )?;
654        validate_permission_rule_field(
655            &rule.patterns,
656            &format!("{domain}.rules[{index}].patterns"),
657        )?;
658    }
659
660    Ok(())
661}
662
663fn validate_permission_rule_field(values: &[String], field: &str) -> Result<(), SidecarError> {
664    if values.is_empty() {
665        return Err(SidecarError::InvalidState(format!(
666            "invalid permissions policy: {field} must not be empty; use [\"*\"] for wildcard"
667        )));
668    }
669    Ok(())
670}
671
672fn capability_operation<'a>(capability: &'a str, domain: &str) -> &'a str {
673    capability
674        .strip_prefix(domain)
675        .and_then(|value| value.strip_prefix('.'))
676        .unwrap_or("")
677}
678
679fn permission_mode_to_kernel_decision(
680    mode: PermissionMode,
681    capability: &str,
682) -> PermissionDecision {
683    match mode {
684        PermissionMode::Allow => PermissionDecision::allow(),
685        PermissionMode::Ask => {
686            PermissionDecision::deny(format!("permission prompt required for {capability}"))
687        }
688        PermissionMode::Deny => PermissionDecision::deny(format!("blocked by {capability} policy")),
689    }
690}
691
692pub(crate) fn filesystem_permission_capability(access: FilesystemAccess) -> &'static str {
693    match access {
694        FilesystemAccess::Read => "fs.read",
695        FilesystemAccess::Write => "fs.write",
696        FilesystemAccess::Stat => "fs.stat",
697        FilesystemAccess::ReadDir => "fs.readdir",
698        FilesystemAccess::CreateDir => "fs.create_dir",
699        FilesystemAccess::Remove => "fs.rm",
700        FilesystemAccess::Rename => "fs.rename",
701        FilesystemAccess::Symlink => "fs.symlink",
702        FilesystemAccess::ReadLink => "fs.readlink",
703        FilesystemAccess::Chmod => "fs.chmod",
704        FilesystemAccess::Truncate => "fs.truncate",
705    }
706}
707
708fn network_permission_capability(operation: NetworkOperation) -> &'static str {
709    match operation {
710        NetworkOperation::Fetch => "network.fetch",
711        NetworkOperation::Http => "network.http",
712        NetworkOperation::Dns => "network.dns",
713        NetworkOperation::Listen => "network.listen",
714    }
715}
716
717fn environment_permission_capability(operation: EnvironmentOperation) -> &'static str {
718    match operation {
719        EnvironmentOperation::Read => "env.read",
720        EnvironmentOperation::Write => "env.write",
721    }
722}
723
724fn is_internal_runtime_command_request(request: &CommandAccessRequest) -> bool {
725    match request.command.as_str() {
726        "node" => request
727            .env
728            .keys()
729            .any(|key| INTERNAL_JAVASCRIPT_ENTRYPOINT_ENV_KEYS.contains(&key.as_str())),
730        "wasm" => request
731            .env
732            .keys()
733            .any(|key| INTERNAL_WASM_ENTRYPOINT_ENV_KEYS.contains(&key.as_str())),
734        "python" => request.env.keys().any(|key| {
735            INTERNAL_PYTHON_ENTRYPOINT_ENV_PREFIXES
736                .iter()
737                .any(|prefix| key.starts_with(prefix))
738        }),
739        _ => false,
740    }
741}
742
743fn ownership_matches_process_event(
744    ownership: &OwnershipScope,
745    event: &ProcessEventEnvelope,
746) -> bool {
747    match ownership {
748        OwnershipScope::ConnectionOwnership(inner) => inner.connection_id == event.connection_id,
749        OwnershipScope::SessionOwnership(inner) => {
750            inner.connection_id == event.connection_id && inner.session_id == event.session_id
751        }
752        OwnershipScope::VmOwnership(inner) => {
753            inner.connection_id == event.connection_id
754                && inner.session_id == event.session_id
755                && inner.vm_id == event.vm_id
756        }
757    }
758}
759
760fn public_process_event_matches_ownership<B>(
761    sidecar: &NativeSidecar<B>,
762    ownership: &OwnershipScope,
763    event: &ProcessEventEnvelope,
764) -> bool
765where
766    B: NativeSidecarBridge + Send + 'static,
767    BridgeError<B>: fmt::Debug + Send + Sync + 'static,
768{
769    if !ownership_matches_process_event(ownership, event) {
770        return false;
771    }
772
773    if event.process_id.contains('/') {
774        return false;
775    }
776
777    // Stale queued events must still be drained through handle_process_event_envelope()
778    // so the sidecar can emit the expected fail-closed log when teardown wins the race.
779    let _ = sidecar;
780    true
781}
782
783fn poll_future_once<F: std::future::Future>(future: std::pin::Pin<&mut F>) -> Option<F::Output> {
784    let mut context = Context::from_waker(Waker::noop());
785    match future.poll(&mut context) {
786        Poll::Ready(output) => Some(output),
787        Poll::Pending => None,
788    }
789}
790
791// ConnectionState, SessionState, VmConfiguration, VmState moved to crate::state
792
793// JavascriptSocketPathContext, JavascriptSocketFamily, VmListenPolicy moved to crate::state
794
795impl JavascriptSocketPathContext {
796    pub(crate) fn loopback_port_allowed(&self, port: u16) -> bool {
797        self.loopback_exempt_ports.contains(&port)
798            || self
799                .tcp_loopback_guest_to_host_ports
800                .keys()
801                .any(|(_, guest_port)| *guest_port == port)
802    }
803
804    pub(crate) fn translate_tcp_loopback_port(
805        &self,
806        family: JavascriptSocketFamily,
807        port: u16,
808    ) -> Option<u16> {
809        self.tcp_loopback_guest_to_host_ports
810            .get(&(family, port))
811            .copied()
812    }
813
814    pub(crate) fn http_loopback_target(
815        &self,
816        family: JavascriptSocketFamily,
817        port: u16,
818    ) -> Option<&crate::state::JavascriptHttpLoopbackTarget> {
819        self.http_loopback_targets.get(&(family, port))
820    }
821
822    pub(crate) fn translate_udp_loopback_port(
823        &self,
824        family: JavascriptSocketFamily,
825        port: u16,
826    ) -> Option<u16> {
827        self.udp_loopback_guest_to_host_ports
828            .get(&(family, port))
829            .copied()
830    }
831
832    pub(crate) fn guest_udp_port_for_host_port(
833        &self,
834        family: JavascriptSocketFamily,
835        port: u16,
836    ) -> Option<u16> {
837        self.udp_loopback_host_to_guest_ports
838            .get(&(family, port))
839            .copied()
840    }
841}
842
843// ActiveProcess, NetworkResourceCounts moved to crate::state
844
845pub struct NativeSidecar<B> {
846    pub(crate) config: NativeSidecarConfig,
847    pub(crate) bridge: SharedBridge<B>,
848    pub(crate) mount_plugins: FileSystemPluginRegistry<MountPluginContext<B>>,
849    pub(crate) cache_root: PathBuf,
850    pub(crate) javascript_engine: JavascriptExecutionEngine,
851    pub(crate) python_engine: PythonExecutionEngine,
852    pub(crate) wasm_engine: WasmExecutionEngine,
853    pub(crate) next_connection_id: usize,
854    pub(crate) next_session_id: usize,
855    pub(crate) next_vm_id: usize,
856    pub(crate) next_sidecar_request_id: RequestId,
857    pub(crate) connections: BTreeMap<String, ConnectionState>,
858    pub(crate) sessions: BTreeMap<String, SessionState>,
859    pub(crate) vms: BTreeMap<String, VmState>,
860    #[allow(dead_code)]
861    pub(crate) process_event_sender: Sender<ProcessEventEnvelope>,
862    pub(crate) process_event_receiver: Option<Receiver<ProcessEventEnvelope>>,
863    pub(crate) pending_process_events: VecDeque<ProcessEventEnvelope>,
864    pub(crate) pending_sidecar_responses: SidecarResponseTracker,
865    pub(crate) outbound_sidecar_requests: VecDeque<SidecarRequestFrame>,
866    pub(crate) completed_sidecar_responses: BTreeMap<RequestId, SidecarResponseFrame>,
867    pub(crate) completed_sidecar_response_order: VecDeque<RequestId>,
868    pub(crate) completed_sidecar_responses_gauge: Arc<QueueGauge>,
869    pub(crate) pending_process_events_gauge: Arc<QueueGauge>,
870    pub(crate) pending_sidecar_responses_gauge: Arc<QueueGauge>,
871    pub(crate) outbound_sidecar_requests_gauge: Arc<QueueGauge>,
872    pub(crate) sidecar_requests: SharedSidecarRequestClient,
873    pub(crate) event_sink: SharedEventSink,
874    pub(crate) extensions: BTreeMap<String, Arc<dyn Extension>>,
875    pub(crate) extension_sessions: BTreeMap<(String, String), ExtensionSessionResources>,
876    pub(crate) extension_process_output_buffers:
877        BTreeMap<(String, String), ExtensionBufferedProcessOutput>,
878    /// Session scopes (connection_id, session_id) disposed since the stdio
879    /// transport last drained them. Lets the transport remove dead sessions from
880    /// its active-session set instead of iterating them forever (M5).
881    pub(crate) disposed_sessions: Vec<(String, String)>,
882}
883
884#[derive(Debug)]
885pub(crate) struct ExtensionSessionResources {
886    pub(crate) ownership: OwnershipScope,
887    pub(crate) process_ids: BTreeSet<String>,
888    pub(crate) vm_ids: BTreeSet<String>,
889}
890
891impl<B> fmt::Debug for NativeSidecar<B> {
892    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
893        f.debug_struct("NativeSidecar")
894            .field("config", &self.config)
895            .field("cache_root", &self.cache_root)
896            .field("next_connection_id", &self.next_connection_id)
897            .field("next_session_id", &self.next_session_id)
898            .field("next_vm_id", &self.next_vm_id)
899            .field("connection_count", &self.connections.len())
900            .field("session_count", &self.sessions.len())
901            .field("vm_count", &self.vms.len())
902            .field("extension_session_count", &self.extension_sessions.len())
903            .field(
904                "extension_process_output_buffer_count",
905                &self.extension_process_output_buffers.len(),
906            )
907            .finish()
908    }
909}
910
911impl<B> NativeSidecar<B>
912where
913    B: NativeSidecarBridge + Send + 'static,
914    BridgeError<B>: fmt::Debug + Send + Sync + 'static,
915{
916    pub fn new(bridge: B) -> Result<Self, SidecarError> {
917        Self::with_config(bridge, NativeSidecarConfig::default())
918    }
919
920    pub fn with_config(bridge: B, config: NativeSidecarConfig) -> Result<Self, SidecarError> {
921        if matches!(config.expected_auth_token.as_deref(), Some("")) {
922            return Err(SidecarError::InvalidState(String::from(
923                "native sidecar expected_auth_token must not be empty",
924            )));
925        }
926
927        let cache_root = config.compile_cache_root.clone().unwrap_or_else(|| {
928            std::env::temp_dir().join(format!(
929                "{}-{}",
930                config.sidecar_id,
931                SystemTime::now()
932                    .duration_since(UNIX_EPOCH)
933                    .expect("system time before unix epoch")
934                    .as_nanos()
935            ))
936        });
937        fs::create_dir_all(&cache_root).map_err(|error| {
938            SidecarError::Io(format!("failed to prepare sidecar cache root: {error}"))
939        })?;
940
941        let bridge = SharedBridge::new(bridge);
942        let mount_plugins = build_mount_plugin_registry::<B>()?;
943        let (process_event_sender, process_event_receiver) = channel(MAX_PROCESS_EVENT_QUEUE);
944
945        Ok(Self {
946            config,
947            bridge,
948            mount_plugins,
949            cache_root,
950            javascript_engine: JavascriptExecutionEngine::default(),
951            python_engine: PythonExecutionEngine::default(),
952            wasm_engine: WasmExecutionEngine::default(),
953            next_connection_id: 0,
954            next_session_id: 0,
955            next_vm_id: 0,
956            next_sidecar_request_id: -1,
957            connections: BTreeMap::new(),
958            sessions: BTreeMap::new(),
959            vms: BTreeMap::new(),
960            process_event_sender,
961            process_event_receiver: Some(process_event_receiver),
962            pending_process_events: VecDeque::new(),
963            pending_sidecar_responses: SidecarResponseTracker::default(),
964            outbound_sidecar_requests: VecDeque::new(),
965            completed_sidecar_responses: BTreeMap::new(),
966            completed_sidecar_response_order: VecDeque::new(),
967            completed_sidecar_responses_gauge: register_queue(
968                TrackedLimit::CompletedSidecarResponses,
969                MAX_COMPLETED_SIDECAR_RESPONSES,
970            ),
971            pending_process_events_gauge: register_queue(
972                TrackedLimit::PendingProcessEvents,
973                MAX_PROCESS_EVENT_QUEUE,
974            ),
975            pending_sidecar_responses_gauge: register_queue(
976                TrackedLimit::PendingSidecarResponses,
977                MAX_PENDING_SIDECAR_RESPONSES,
978            ),
979            outbound_sidecar_requests_gauge: register_queue(
980                TrackedLimit::OutboundSidecarRequests,
981                MAX_OUTBOUND_SIDECAR_REQUESTS,
982            ),
983            sidecar_requests: SharedSidecarRequestClient::default(),
984            event_sink: SharedEventSink::default(),
985            extensions: BTreeMap::new(),
986            extension_sessions: BTreeMap::new(),
987            extension_process_output_buffers: BTreeMap::new(),
988            disposed_sessions: Vec::new(),
989        })
990    }
991
992    pub fn with_config_and_extensions(
993        bridge: B,
994        config: NativeSidecarConfig,
995        extensions: Vec<Box<dyn Extension>>,
996    ) -> Result<Self, SidecarError> {
997        let mut sidecar = Self::with_config(bridge, config)?;
998        for extension in extensions {
999            sidecar.register_extension(extension)?;
1000        }
1001        Ok(sidecar)
1002    }
1003
1004    pub(crate) fn prune_extension_process_resource(&mut self, process_id: &str) {
1005        self.extension_sessions.retain(|_, resources| {
1006            resources.process_ids.remove(process_id);
1007            !resources.process_ids.is_empty() || !resources.vm_ids.is_empty()
1008        });
1009    }
1010
1011    pub(crate) fn prune_extension_vm_resource(&mut self, vm_id: &str) {
1012        self.extension_sessions.retain(|_, resources| {
1013            if matches!(
1014                &resources.ownership,
1015                OwnershipScope::VmOwnership(inner) if inner.vm_id == vm_id
1016            ) {
1017                resources.process_ids.clear();
1018            }
1019            resources.vm_ids.remove(vm_id);
1020            !resources.process_ids.is_empty() || !resources.vm_ids.is_empty()
1021        });
1022    }
1023
1024    /// Reclaim every per-VM tracking entry owned by the sidecar for `vm_id`.
1025    ///
1026    /// Called unconditionally from `dispose_vm_internal` so that a fallible
1027    /// teardown step (root-filesystem snapshot/flush, kernel dispose, permission
1028    /// reset) erroring out with `?` can never strand these maps for the rest of
1029    /// the process lifetime (H1). This also reclaims the ACP output-buffer map,
1030    /// which was previously removed only on a successful handoff and leaked on VM
1031    /// or session disposal (M6).
1032    pub(crate) fn reclaim_vm_tracking(&mut self, session_id: &str, vm_id: &str) {
1033        self.javascript_engine.dispose_vm(vm_id);
1034        self.python_engine.dispose_vm(vm_id);
1035        self.wasm_engine.dispose_vm(vm_id);
1036        self.prune_extension_vm_resource(vm_id);
1037        self.extension_process_output_buffers
1038            .retain(|(buffer_vm_id, _process_id), _| buffer_vm_id != vm_id);
1039        if let Some(session) = self.sessions.get_mut(session_id) {
1040            session.vm_ids.remove(vm_id);
1041        }
1042    }
1043
1044    pub(crate) fn capture_extension_process_output_event(
1045        &mut self,
1046        vm_id: &str,
1047        process_id: &str,
1048        event: &ActiveExecutionEvent,
1049    ) -> bool {
1050        let Some(buffer) = self
1051            .extension_process_output_buffers
1052            .get_mut(&(vm_id.to_string(), process_id.to_string()))
1053        else {
1054            return false;
1055        };
1056        match event {
1057            ActiveExecutionEvent::Stdout(chunk) => {
1058                buffer.append_stdout(chunk, DEFAULT_ACP_STDOUT_BUFFER_BYTE_LIMIT);
1059                true
1060            }
1061            ActiveExecutionEvent::Stderr(chunk) => {
1062                buffer.append_stderr(chunk, DEFAULT_ACP_STDOUT_BUFFER_BYTE_LIMIT);
1063                true
1064            }
1065            ActiveExecutionEvent::JavascriptSyncRpcRequest(_)
1066            | ActiveExecutionEvent::PythonVfsRpcRequest(_)
1067            | ActiveExecutionEvent::SignalState { .. }
1068            | ActiveExecutionEvent::Exited(_) => false,
1069        }
1070    }
1071
1072    fn bind_extension_process_resource(
1073        &mut self,
1074        ownership: OwnershipScope,
1075        namespace: String,
1076        ext_session_id: String,
1077        process_id: String,
1078    ) -> Result<(), SidecarError> {
1079        if ext_session_id.is_empty() {
1080            return Err(SidecarError::InvalidState(String::from(
1081                "extension session id must not be empty",
1082            )));
1083        }
1084        let (connection_id, session_id, vm_id) = self.vm_scope_for(&ownership)?;
1085        self.require_owned_vm(&connection_id, &session_id, &vm_id)?;
1086        let process_exists = self
1087            .vms
1088            .get(&vm_id)
1089            .is_some_and(|vm| vm.active_processes.contains_key(&process_id));
1090        if !process_exists {
1091            return Err(SidecarError::InvalidState(format!(
1092                "VM {vm_id} has no active process {process_id}"
1093            )));
1094        }
1095
1096        let key = (namespace, ext_session_id);
1097        if let Some(resources) = self.extension_sessions.get_mut(&key) {
1098            if resources.ownership != ownership {
1099                return Err(SidecarError::InvalidState(String::from(
1100                    "extension session ownership did not match existing resources",
1101                )));
1102            }
1103            resources.process_ids.insert(process_id);
1104        } else {
1105            self.extension_sessions.insert(
1106                key,
1107                ExtensionSessionResources {
1108                    ownership,
1109                    process_ids: BTreeSet::from([process_id]),
1110                    vm_ids: BTreeSet::new(),
1111                },
1112            );
1113        }
1114        Ok(())
1115    }
1116
1117    fn bind_extension_vm_resource(
1118        &mut self,
1119        ownership: OwnershipScope,
1120        namespace: String,
1121        ext_session_id: String,
1122    ) -> Result<(), SidecarError> {
1123        if ext_session_id.is_empty() {
1124            return Err(SidecarError::InvalidState(String::from(
1125                "extension session id must not be empty",
1126            )));
1127        }
1128        let (connection_id, session_id, vm_id) = self.vm_scope_for(&ownership)?;
1129        self.require_owned_vm(&connection_id, &session_id, &vm_id)?;
1130
1131        let key = (namespace, ext_session_id);
1132        if let Some(resources) = self.extension_sessions.get_mut(&key) {
1133            if resources.ownership != ownership {
1134                return Err(SidecarError::InvalidState(String::from(
1135                    "extension session ownership did not match existing resources",
1136                )));
1137            }
1138            resources.vm_ids.insert(vm_id);
1139        } else {
1140            self.extension_sessions.insert(
1141                key,
1142                ExtensionSessionResources {
1143                    ownership,
1144                    process_ids: BTreeSet::new(),
1145                    vm_ids: BTreeSet::from([vm_id]),
1146                },
1147            );
1148        }
1149        Ok(())
1150    }
1151
1152    pub fn sidecar_id(&self) -> &str {
1153        &self.config.sidecar_id
1154    }
1155
1156    pub fn with_bridge_mut<T>(
1157        &self,
1158        operation: impl FnOnce(&mut B) -> T,
1159    ) -> Result<T, SidecarError> {
1160        self.bridge.inspect(operation)
1161    }
1162
1163    pub fn set_sidecar_request_transport(&mut self, transport: Arc<dyn SidecarRequestTransport>) {
1164        self.sidecar_requests.set_transport(transport);
1165    }
1166
1167    pub fn set_event_transport(&mut self, transport: Arc<dyn EventSinkTransport>) {
1168        self.event_sink.set_transport(transport);
1169    }
1170
1171    pub fn register_extension(
1172        &mut self,
1173        extension: Box<dyn Extension>,
1174    ) -> Result<(), SidecarError> {
1175        let namespace = extension.namespace().to_owned();
1176        if namespace.is_empty() {
1177            return Err(SidecarError::InvalidState(String::from(
1178                "extension namespace must not be empty",
1179            )));
1180        }
1181        if self.extensions.contains_key(&namespace) {
1182            return Err(SidecarError::Conflict(format!(
1183                "extension namespace {namespace} is already registered",
1184            )));
1185        }
1186        self.extensions.insert(namespace, Arc::from(extension));
1187        Ok(())
1188    }
1189
1190    pub fn set_sidecar_request_handler<F>(&mut self, handler: F)
1191    where
1192        F: Fn(SidecarRequestFrame) -> Result<SidecarResponsePayload, SidecarError>
1193            + Send
1194            + Sync
1195            + 'static,
1196    {
1197        struct HandlerTransport<F>(F);
1198
1199        impl<F> SidecarRequestTransport for HandlerTransport<F>
1200        where
1201            F: Fn(SidecarRequestFrame) -> Result<SidecarResponsePayload, SidecarError>
1202                + Send
1203                + Sync
1204                + 'static,
1205        {
1206            fn send_request(
1207                &self,
1208                request: SidecarRequestFrame,
1209                _timeout: Duration,
1210            ) -> Result<SidecarResponseFrame, SidecarError> {
1211                let payload = (self.0)(request.clone())?;
1212                Ok(SidecarResponseFrame::new(
1213                    request.request_id,
1214                    request.ownership,
1215                    payload,
1216                ))
1217            }
1218        }
1219
1220        self.set_sidecar_request_transport(Arc::new(HandlerTransport(handler)));
1221    }
1222
1223    pub fn set_wire_sidecar_request_handler<F>(&mut self, handler: F)
1224    where
1225        F: Fn(
1226                crate::wire::SidecarRequestFrame,
1227            ) -> Result<crate::wire::SidecarResponseFrame, SidecarError>
1228            + Send
1229            + Sync
1230            + 'static,
1231    {
1232        self.set_sidecar_request_handler(move |request| {
1233            let request = crate::wire::sidecar_request_frame_from_compat(request)
1234                .map_err(wire_protocol_error)?;
1235            let response = handler(request)?;
1236            let response = crate::wire::sidecar_response_frame_to_compat(response)
1237                .map_err(wire_protocol_error)?;
1238            Ok(response.payload)
1239        });
1240    }
1241
1242    pub(crate) fn queue_pending_process_event(
1243        &mut self,
1244        envelope: ProcessEventEnvelope,
1245    ) -> Result<(), SidecarError> {
1246        if self.pending_process_events.len() >= MAX_PROCESS_EVENT_QUEUE {
1247            return Err(process_event_queue_overflow_error());
1248        }
1249        self.pending_process_events.push_back(envelope);
1250        self.pending_process_events_gauge
1251            .observe_depth(self.pending_process_events.len());
1252        Ok(())
1253    }
1254
1255    pub(crate) fn queue_front_pending_process_event(
1256        &mut self,
1257        envelope: ProcessEventEnvelope,
1258    ) -> Result<(), SidecarError> {
1259        if self.pending_process_events.len() >= MAX_PROCESS_EVENT_QUEUE {
1260            return Err(process_event_queue_overflow_error());
1261        }
1262        self.pending_process_events.push_front(envelope);
1263        self.pending_process_events_gauge
1264            .observe_depth(self.pending_process_events.len());
1265        Ok(())
1266    }
1267
1268    pub(crate) fn pending_process_event_capacity(&self) -> usize {
1269        MAX_PROCESS_EVENT_QUEUE.saturating_sub(self.pending_process_events.len())
1270    }
1271
1272    pub fn dispatch_blocking(
1273        &mut self,
1274        request: RequestFrame,
1275    ) -> Result<DispatchResult, SidecarError> {
1276        let inside_runtime = tokio::runtime::Handle::try_current().is_ok();
1277        if matches!(
1278            request.payload,
1279            RequestPayload::DisposeVm(_) | RequestPayload::Ext(_)
1280        ) && !inside_runtime
1281        {
1282            return tokio::runtime::Builder::new_current_thread()
1283                .enable_all()
1284                .build()
1285                .expect("sidecar dispatch runtime")
1286                .block_on(self.dispatch(request));
1287        }
1288
1289        let mut future = std::pin::pin!(self.dispatch(request));
1290        match poll_future_once(future.as_mut()) {
1291            Some(result) => result,
1292            None if inside_runtime => Err(SidecarError::InvalidState(String::from(
1293                "dispatch_blocking cannot wait for an async sidecar request inside a Tokio runtime; use dispatch().await",
1294            ))),
1295            None => tokio::runtime::Builder::new_current_thread()
1296                .enable_all()
1297                .build()
1298                .expect("sidecar dispatch runtime")
1299                .block_on(future),
1300        }
1301    }
1302
1303    pub fn dispatch_wire_blocking(
1304        &mut self,
1305        request: crate::wire::RequestFrame,
1306    ) -> Result<crate::wire::WireDispatchResult, SidecarError> {
1307        let request = crate::wire::request_frame_to_compat(request).map_err(wire_protocol_error)?;
1308        let result = self.dispatch_blocking(request)?;
1309        crate::wire::dispatch_result_from_compat(result).map_err(wire_protocol_error)
1310    }
1311
1312    pub fn poll_event_blocking(
1313        &mut self,
1314        ownership: &OwnershipScope,
1315        timeout: Duration,
1316    ) -> Result<Option<EventFrame>, SidecarError> {
1317        tokio::runtime::Builder::new_current_thread()
1318            .enable_all()
1319            .build()
1320            .expect("sidecar poll runtime")
1321            .block_on(self.poll_event(ownership, timeout))
1322    }
1323
1324    pub fn poll_event_wire_blocking(
1325        &mut self,
1326        ownership: &crate::wire::OwnershipScope,
1327        timeout: Duration,
1328    ) -> Result<Option<crate::wire::EventFrame>, SidecarError> {
1329        let ownership = crate::wire::ownership_scope_to_compat(ownership.clone());
1330        self.poll_event_blocking(&ownership, timeout)?
1331            .map(crate::wire::event_frame_from_compat)
1332            .transpose()
1333            .map_err(wire_protocol_error)
1334    }
1335
1336    pub fn close_session_blocking(
1337        &mut self,
1338        connection_id: &str,
1339        session_id: &str,
1340    ) -> Result<Vec<EventFrame>, SidecarError> {
1341        tokio::runtime::Builder::new_current_thread()
1342            .enable_all()
1343            .build()
1344            .expect("sidecar close-session runtime")
1345            .block_on(self.close_session(connection_id, session_id))
1346    }
1347
1348    pub fn remove_connection_blocking(
1349        &mut self,
1350        connection_id: &str,
1351    ) -> Result<Vec<EventFrame>, SidecarError> {
1352        tokio::runtime::Builder::new_current_thread()
1353            .enable_all()
1354            .build()
1355            .expect("sidecar remove-connection runtime")
1356            .block_on(self.remove_connection(connection_id))
1357    }
1358
1359    pub fn dispose_vm_internal_blocking(
1360        &mut self,
1361        connection_id: &str,
1362        session_id: &str,
1363        vm_id: &str,
1364        reason: DisposeReason,
1365    ) -> Result<Vec<EventFrame>, SidecarError> {
1366        tokio::runtime::Builder::new_current_thread()
1367            .enable_all()
1368            .build()
1369            .expect("sidecar dispose-vm runtime")
1370            .block_on(self.dispose_vm_internal(connection_id, session_id, vm_id, reason))
1371    }
1372
1373    pub async fn dispatch(
1374        &mut self,
1375        request: RequestFrame,
1376    ) -> Result<DispatchResult, SidecarError> {
1377        if let Err(error) = self.ensure_request_within_frame_limit(&request) {
1378            return Ok(DispatchResult {
1379                response: self.reject(&request, error_code(&error), &error.to_string()),
1380                events: Vec::new(),
1381            });
1382        }
1383
1384        let result = match request.payload.clone() {
1385            RequestPayload::Authenticate(payload) => {
1386                self.authenticate_connection(&request, payload).await
1387            }
1388            RequestPayload::OpenSession(payload) => self.open_session(&request, payload).await,
1389            RequestPayload::CreateVm(payload) => self.create_vm(&request, payload).await,
1390            RequestPayload::DisposeVm(payload) => self.dispose_vm(&request, payload).await,
1391            RequestPayload::BootstrapRootFilesystem(payload) => {
1392                self.bootstrap_root_filesystem(&request, payload.entries)
1393                    .await
1394            }
1395            RequestPayload::ConfigureVm(payload) => self.configure_vm(&request, payload).await,
1396            RequestPayload::RegisterHostCallbacks(payload) => {
1397                register_host_callbacks(self, &request, payload)
1398            }
1399            RequestPayload::CreateLayer(payload) => self.create_layer(&request, payload).await,
1400            RequestPayload::SealLayer(payload) => self.seal_layer(&request, payload).await,
1401            RequestPayload::ImportSnapshot(payload) => {
1402                self.import_snapshot(&request, payload).await
1403            }
1404            RequestPayload::ExportSnapshot(payload) => {
1405                self.export_snapshot(&request, payload).await
1406            }
1407            RequestPayload::CreateOverlay(payload) => self.create_overlay(&request, payload).await,
1408            RequestPayload::GuestFilesystemCall(payload) => {
1409                self.guest_filesystem_call(&request, payload).await
1410            }
1411            RequestPayload::SnapshotRootFilesystem(payload) => {
1412                self.snapshot_root_filesystem(&request, payload).await
1413            }
1414            RequestPayload::Execute(payload) => self.execute(&request, payload).await,
1415            RequestPayload::WriteStdin(payload) => self.write_stdin(&request, payload).await,
1416            RequestPayload::CloseStdin(payload) => self.close_stdin(&request, payload).await,
1417            RequestPayload::KillProcess(payload) => self.kill_process(&request, payload).await,
1418            RequestPayload::GetProcessSnapshot(payload) => {
1419                self.get_process_snapshot(&request, payload).await
1420            }
1421            RequestPayload::FindListener(payload) => self.find_listener(&request, payload).await,
1422            RequestPayload::FindBoundUdp(payload) => self.find_bound_udp(&request, payload).await,
1423            RequestPayload::VmFetch(payload) => self.vm_fetch(&request, payload).await,
1424            RequestPayload::GetSignalState(payload) => {
1425                self.get_signal_state(&request, payload).await
1426            }
1427            RequestPayload::GetZombieTimerCount(payload) => {
1428                self.get_zombie_timer_count(&request, payload).await
1429            }
1430            RequestPayload::HostFilesystemCall(_)
1431            | RequestPayload::PersistenceLoad(_)
1432            | RequestPayload::PersistenceFlush(_) => Ok(DispatchResult {
1433                response: self.reject(
1434                    &request,
1435                    "unsupported_direction",
1436                    "host callback request categories are sidecar-to-host only in this scaffold",
1437                ),
1438                events: Vec::new(),
1439            }),
1440            RequestPayload::Ext(payload) => {
1441                self.dispatch_extension_request(&request, payload).await
1442            }
1443        };
1444
1445        match result {
1446            Ok(dispatch) => Ok(dispatch),
1447            Err(error @ SidecarError::Io(_)) => Err(error),
1448            Err(error) => Ok(DispatchResult {
1449                response: self.reject(&request, error_code(&error), &error.to_string()),
1450                events: Vec::new(),
1451            }),
1452        }
1453    }
1454
1455    pub async fn dispatch_wire(
1456        &mut self,
1457        request: crate::wire::RequestFrame,
1458    ) -> Result<crate::wire::WireDispatchResult, SidecarError> {
1459        let request = crate::wire::request_frame_to_compat(request).map_err(wire_protocol_error)?;
1460        let result = self.dispatch(request).await?;
1461        crate::wire::dispatch_result_from_compat(result).map_err(wire_protocol_error)
1462    }
1463
1464    pub async fn poll_event_wire(
1465        &mut self,
1466        ownership: &crate::wire::OwnershipScope,
1467        timeout: Duration,
1468    ) -> Result<Option<crate::wire::EventFrame>, SidecarError> {
1469        let ownership = crate::wire::ownership_scope_to_compat(ownership.clone());
1470        self.poll_event(&ownership, timeout)
1471            .await?
1472            .map(crate::wire::event_frame_from_compat)
1473            .transpose()
1474            .map_err(wire_protocol_error)
1475    }
1476
1477    async fn dispatch_extension_request(
1478        &mut self,
1479        request: &RequestFrame,
1480        envelope: ExtEnvelope,
1481    ) -> Result<DispatchResult, SidecarError> {
1482        let namespace = envelope.namespace;
1483        let Some(extension) = self.extensions.get(&namespace).cloned() else {
1484            return Ok(DispatchResult {
1485                response: self.reject(
1486                    request,
1487                    "unknown_extension",
1488                    &format!("no extension registered for namespace {namespace}"),
1489                ),
1490                events: Vec::new(),
1491            });
1492        };
1493        let snapshot = ExtensionSnapshot::new(
1494            namespace.clone(),
1495            request.ownership.clone(),
1496            self.sidecar_requests.clone(),
1497            self.event_sink.clone(),
1498        );
1499        let ctx = ExtensionContext::new(snapshot, self);
1500        let response = extension.handle_request(ctx, envelope.payload).await?;
1501        Ok(DispatchResult {
1502            response: self.respond(
1503                request,
1504                ResponsePayload::ExtResult(ExtEnvelope {
1505                    namespace,
1506                    payload: response.payload,
1507                }),
1508            ),
1509            events: response.events,
1510        })
1511    }
1512
1513    pub async fn poll_event(
1514        &mut self,
1515        ownership: &OwnershipScope,
1516        timeout: Duration,
1517    ) -> Result<Option<EventFrame>, SidecarError> {
1518        let deadline = Instant::now() + timeout;
1519        loop {
1520            if let Some(index) = self
1521                .pending_process_events
1522                .iter()
1523                .position(|event| public_process_event_matches_ownership(self, ownership, event))
1524            {
1525                let Some(envelope) = self.pending_process_events.remove(index) else {
1526                    continue;
1527                };
1528                if let Some(frame) = self.handle_process_event_envelope(envelope)? {
1529                    return Ok(Some(frame));
1530                }
1531                continue;
1532            }
1533
1534            if !timeout.is_zero() {
1535                let _ = self.pump_process_events(ownership).await?;
1536            }
1537
1538            let queued_envelopes = {
1539                let pending_capacity = self.pending_process_event_capacity();
1540                let receiver = self.process_event_receiver.as_mut().ok_or_else(|| {
1541                    SidecarError::InvalidState(String::from("process event receiver unavailable"))
1542                })?;
1543                let mut queued = Vec::new();
1544                loop {
1545                    if queued.len() >= pending_capacity {
1546                        if receiver.is_empty() {
1547                            break;
1548                        }
1549                        return Err(process_event_queue_overflow_error());
1550                    }
1551                    match receiver.try_recv() {
1552                        Ok(envelope) => queued.push(envelope),
1553                        Err(tokio::sync::mpsc::error::TryRecvError::Empty) => break,
1554                        Err(tokio::sync::mpsc::error::TryRecvError::Disconnected) => break,
1555                    }
1556                }
1557                queued
1558            };
1559
1560            let mut matching_envelope = None;
1561            for envelope in queued_envelopes {
1562                if matching_envelope.is_none()
1563                    && public_process_event_matches_ownership(self, ownership, &envelope)
1564                {
1565                    matching_envelope = Some(envelope);
1566                } else {
1567                    self.queue_pending_process_event(envelope)?;
1568                }
1569            }
1570
1571            if let Some(envelope) = matching_envelope {
1572                if let Some(frame) = self.handle_process_event_envelope(envelope)? {
1573                    return Ok(Some(frame));
1574                }
1575                continue;
1576            }
1577
1578            if Instant::now() >= deadline {
1579                return Ok(None);
1580            }
1581
1582            let remaining = deadline.saturating_duration_since(Instant::now());
1583            time::sleep(remaining.min(Duration::from_millis(10))).await;
1584        }
1585    }
1586
1587    pub(crate) fn handle_process_event_envelope(
1588        &mut self,
1589        envelope: ProcessEventEnvelope,
1590    ) -> Result<Option<EventFrame>, SidecarError> {
1591        let ProcessEventEnvelope {
1592            connection_id,
1593            session_id,
1594            vm_id,
1595            process_id,
1596            event,
1597        } = envelope;
1598
1599        if matches!(event, ActiveExecutionEvent::Exited(_)) {
1600            let mut trailing = Vec::new();
1601            let mut deferred = VecDeque::new();
1602            while let Some(pending) = self.pending_process_events.pop_front() {
1603                if pending.vm_id == vm_id
1604                    && pending.process_id == process_id
1605                    && !matches!(pending.event, ActiveExecutionEvent::Exited(_))
1606                {
1607                    trailing.push(pending.event);
1608                } else {
1609                    deferred.push_back(pending);
1610                }
1611            }
1612            self.pending_process_events = deferred;
1613            let drain_limit = self
1614                .pending_process_event_capacity()
1615                .saturating_sub(trailing.len().saturating_add(1));
1616            trailing.extend(
1617                self.drain_process_events_blocking_with_limit(&vm_id, &process_id, drain_limit)?
1618                    .into_iter()
1619                    .filter(|event| !matches!(event, ActiveExecutionEvent::Exited(_))),
1620            );
1621
1622            if !trailing.is_empty() {
1623                if self.pending_process_event_capacity() < trailing.len() {
1624                    return Err(process_event_queue_overflow_error());
1625                }
1626                let emit_now = if self.pending_process_event_capacity() == trailing.len() {
1627                    Some(trailing.remove(0))
1628                } else {
1629                    None
1630                };
1631                self.queue_front_pending_process_event(ProcessEventEnvelope {
1632                    connection_id: connection_id.clone(),
1633                    session_id: session_id.clone(),
1634                    vm_id: vm_id.clone(),
1635                    process_id: process_id.clone(),
1636                    event,
1637                })?;
1638                for event in trailing.into_iter().rev() {
1639                    self.queue_front_pending_process_event(ProcessEventEnvelope {
1640                        connection_id: connection_id.clone(),
1641                        session_id: session_id.clone(),
1642                        vm_id: vm_id.clone(),
1643                        process_id: process_id.clone(),
1644                        event,
1645                    })?;
1646                }
1647                if let Some(event) = emit_now {
1648                    return self.handle_execution_event(&vm_id, &process_id, event);
1649                }
1650                return Ok(None);
1651            }
1652        }
1653
1654        self.handle_execution_event(&vm_id, &process_id, event)
1655    }
1656
1657    // try_poll_event moved to crate::execution
1658
1659    pub async fn close_session(
1660        &mut self,
1661        connection_id: &str,
1662        session_id: &str,
1663    ) -> Result<Vec<EventFrame>, SidecarError> {
1664        self.dispose_session(connection_id, session_id, DisposeReason::Requested)
1665            .await
1666    }
1667
1668    pub async fn remove_connection(
1669        &mut self,
1670        connection_id: &str,
1671    ) -> Result<Vec<EventFrame>, SidecarError> {
1672        self.require_authenticated_connection(connection_id)?;
1673
1674        let session_ids = self
1675            .connections
1676            .get(connection_id)
1677            .expect("authenticated connection should exist")
1678            .sessions
1679            .iter()
1680            .cloned()
1681            .collect::<Vec<_>>();
1682
1683        let mut events = Vec::new();
1684        let mut first_error: Option<SidecarError> = None;
1685        for session_id in session_ids {
1686            // Attempt EVERY session; aggregate errors instead of `?`-ing out on
1687            // the first so one wedged session cannot abandon the rest (H1).
1688            match self
1689                .dispose_session(connection_id, &session_id, DisposeReason::ConnectionClosed)
1690                .await
1691            {
1692                Ok(session_events) => events.extend(session_events),
1693                Err(error) => {
1694                    if first_error.is_none() {
1695                        first_error = Some(error);
1696                    }
1697                }
1698            }
1699        }
1700
1701        self.connections.remove(connection_id);
1702        if let Some(error) = first_error {
1703            return Err(error);
1704        }
1705        Ok(events)
1706    }
1707
1708    async fn authenticate_connection(
1709        &mut self,
1710        request: &RequestFrame,
1711        payload: crate::protocol::AuthenticateRequest,
1712    ) -> Result<DispatchResult, SidecarError> {
1713        let _ = self.connection_id_for(&request.ownership)?;
1714        if let Err(error) = self.validate_auth_token(&payload.auth_token) {
1715            let mut fields = audit_fields([
1716                (String::from("source"), payload.client_name.clone()),
1717                (String::from("reason"), error.to_string()),
1718            ]);
1719            if let OwnershipScope::ConnectionOwnership(inner) = &request.ownership {
1720                fields.insert(String::from("connection_id"), inner.connection_id.clone());
1721            }
1722            emit_security_audit_event(
1723                &self.bridge,
1724                &self.config.sidecar_id,
1725                "security.auth.failed",
1726                fields,
1727            );
1728            return Err(error);
1729        }
1730
1731        if payload.protocol_version != crate::wire::PROTOCOL_VERSION {
1732            return Err(SidecarError::ProtocolVersionMismatch(format!(
1733                "sidecar protocol version mismatch: expected {}, got {}",
1734                crate::wire::PROTOCOL_VERSION,
1735                payload.protocol_version
1736            )));
1737        }
1738
1739        let expected_bridge_version = secure_exec_bridge::bridge_contract().version;
1740        if payload.bridge_version != expected_bridge_version {
1741            return Err(SidecarError::BridgeVersionMismatch(format!(
1742                "bridge contract version mismatch: expected {expected_bridge_version}, got {}",
1743                payload.bridge_version
1744            )));
1745        }
1746
1747        let connection_id = self.allocate_connection_id();
1748        self.connections.insert(
1749            connection_id.clone(),
1750            ConnectionState {
1751                auth_token: payload.auth_token,
1752                sessions: BTreeSet::new(),
1753            },
1754        );
1755
1756        let response = self.response_with_ownership(
1757            request.request_id,
1758            OwnershipScope::connection(&connection_id),
1759            ResponsePayload::Authenticated(AuthenticatedResponse {
1760                sidecar_id: self.config.sidecar_id.clone(),
1761                connection_id,
1762                max_frame_bytes: self.config.max_frame_bytes as u32,
1763            }),
1764        );
1765        Ok(DispatchResult {
1766            response,
1767            events: Vec::new(),
1768        })
1769    }
1770
1771    async fn open_session(
1772        &mut self,
1773        request: &RequestFrame,
1774        payload: OpenSessionRequest,
1775    ) -> Result<DispatchResult, SidecarError> {
1776        let connection_id = self.connection_id_for(&request.ownership)?;
1777        self.require_authenticated_connection(&connection_id)?;
1778
1779        self.next_session_id += 1;
1780        let session_id = format!("session-{}", self.next_session_id);
1781        self.sessions.insert(
1782            session_id.clone(),
1783            SessionState {
1784                connection_id: connection_id.clone(),
1785                placement: payload.placement,
1786                metadata: payload.metadata.into_iter().collect(),
1787                vm_ids: BTreeSet::new(),
1788            },
1789        );
1790        self.connections
1791            .get_mut(&connection_id)
1792            .expect("authenticated connection should exist")
1793            .sessions
1794            .insert(session_id.clone());
1795
1796        Ok(DispatchResult {
1797            response: self.respond(
1798                request,
1799                ResponsePayload::SessionOpened(SessionOpenedResponse {
1800                    session_id,
1801                    owner_connection_id: connection_id,
1802                }),
1803            ),
1804            events: Vec::new(),
1805        })
1806    }
1807
1808    // create_vm, dispose_vm, bootstrap_root_filesystem, configure_vm moved to crate::vm
1809
1810    async fn guest_filesystem_call(
1811        &mut self,
1812        request: &RequestFrame,
1813        payload: GuestFilesystemCallRequest,
1814    ) -> Result<DispatchResult, SidecarError> {
1815        filesystem_guest_filesystem_call(self, request, payload).await
1816    }
1817
1818    // snapshot_root_filesystem moved to crate::vm
1819
1820    // execute, write_stdin, close_stdin, kill_process, find_listener, find_bound_udp,
1821    // get_signal_state, get_zombie_timer_count moved to crate::execution
1822
1823    async fn dispose_session(
1824        &mut self,
1825        connection_id: &str,
1826        session_id: &str,
1827        reason: DisposeReason,
1828    ) -> Result<Vec<EventFrame>, SidecarError> {
1829        self.require_owned_session(connection_id, session_id)?;
1830
1831        let vm_ids = self
1832            .sessions
1833            .get(session_id)
1834            .expect("owned session should exist")
1835            .vm_ids
1836            .iter()
1837            .cloned()
1838            .collect::<Vec<_>>();
1839
1840        let mut events = Vec::new();
1841        let mut first_error: Option<SidecarError> = None;
1842        for vm_id in vm_ids {
1843            // Attempt EVERY VM; aggregate errors instead of `?`-ing out on the
1844            // first so one stuck VM cannot strand the remaining VMs' teardown and
1845            // leave the session permanently un-reclaimed (H1).
1846            match self
1847                .dispose_vm_internal(connection_id, session_id, &vm_id, reason.clone())
1848                .await
1849            {
1850                Ok(vm_events) => events.extend(vm_events),
1851                Err(error) => {
1852                    if first_error.is_none() {
1853                        first_error = Some(error);
1854                    }
1855                }
1856            }
1857        }
1858
1859        // On client disconnect, give every registered extension a chance to free
1860        // the per-session state it tracks (H4): the host owns the only signal an
1861        // extension gets that a session has gone away.
1862        if matches!(reason, DisposeReason::ConnectionClosed) {
1863            if let Err(error) = self
1864                .dispose_extension_session_state(connection_id, session_id)
1865                .await
1866            {
1867                if first_error.is_none() {
1868                    first_error = Some(error);
1869                }
1870            }
1871        }
1872
1873        self.sessions.remove(session_id);
1874        if let Some(connection) = self.connections.get_mut(connection_id) {
1875            connection.sessions.remove(session_id);
1876        }
1877        // Tell the stdio transport this session is gone so it stops iterating a
1878        // dead entry every event-pump tick and the set stops growing (M5).
1879        self.disposed_sessions
1880            .push((connection_id.to_owned(), session_id.to_owned()));
1881
1882        if let Some(error) = first_error {
1883            return Err(error);
1884        }
1885        Ok(events)
1886    }
1887
1888    /// Invoke each registered extension's per-session teardown hook so it can
1889    /// release the state it keyed on this host session. Errors are aggregated so
1890    /// one misbehaving extension cannot prevent the others from cleaning up.
1891    async fn dispose_extension_session_state(
1892        &mut self,
1893        connection_id: &str,
1894        session_id: &str,
1895    ) -> Result<(), SidecarError> {
1896        let ownership = OwnershipScope::session(connection_id, session_id);
1897        let extensions = self
1898            .extensions
1899            .values()
1900            .cloned()
1901            .collect::<Vec<Arc<dyn Extension>>>();
1902        let mut first_error: Option<SidecarError> = None;
1903        for extension in extensions {
1904            let snapshot = ExtensionSnapshot::new(
1905                extension.namespace().to_owned(),
1906                ownership.clone(),
1907                self.sidecar_requests.clone(),
1908                self.event_sink.clone(),
1909            );
1910            if let Err(error) = extension.on_session_disposed(snapshot).await {
1911                if first_error.is_none() {
1912                    first_error = Some(error);
1913                }
1914            }
1915        }
1916        match first_error {
1917            Some(error) => Err(error),
1918            None => Ok(()),
1919        }
1920    }
1921
1922    /// Drain the session scopes disposed since the last call so the stdio
1923    /// transport can untrack them from its active-session set (M5).
1924    pub(crate) fn take_disposed_sessions(&mut self) -> Vec<(String, String)> {
1925        std::mem::take(&mut self.disposed_sessions)
1926    }
1927
1928    // dispose_vm_internal, terminate_vm_processes, wait_for_vm_processes_to_exit moved to crate::vm
1929
1930    // kill_process_internal, handle_execution_event, handle_python_vfs_rpc_request,
1931    // resolve_javascript_child_process_execution, spawn_javascript_child_process,
1932    // poll_javascript_child_process, write_javascript_child_process_stdin,
1933    // close_javascript_child_process_stdin, kill_javascript_child_process moved to crate::execution
1934
1935    pub(crate) fn handle_javascript_sync_rpc_request(
1936        &mut self,
1937        vm_id: &str,
1938        process_id: &str,
1939        request: JavascriptSyncRpcRequest,
1940    ) -> Result<(), SidecarError> {
1941        let Some(vm) = self.vms.get(vm_id) else {
1942            log_stale_process_event(&self.bridge, vm_id, process_id, "javascript sync RPC");
1943            return Ok(());
1944        };
1945        if !vm.active_processes.contains_key(process_id) {
1946            log_stale_process_event(&self.bridge, vm_id, process_id, "javascript sync RPC");
1947            return Ok(());
1948        }
1949
1950        let response: Result<Value, SidecarError> = match request.method.as_str() {
1951            "child_process.spawn" => {
1952                let Some(vm) = self.vms.get(vm_id) else {
1953                    log_stale_process_event(
1954                        &self.bridge,
1955                        vm_id,
1956                        process_id,
1957                        "javascript sync RPC child_process.spawn",
1958                    );
1959                    return Ok(());
1960                };
1961                let (payload, _) = parse_javascript_child_process_spawn_request(vm, &request.args)?;
1962                self.spawn_javascript_child_process(vm_id, process_id, payload)
1963            }
1964            "child_process.spawn_sync" => {
1965                let Some(vm) = self.vms.get(vm_id) else {
1966                    log_stale_process_event(
1967                        &self.bridge,
1968                        vm_id,
1969                        process_id,
1970                        "javascript sync RPC child_process.spawn_sync",
1971                    );
1972                    return Ok(());
1973                };
1974                let (payload, max_buffer) =
1975                    parse_javascript_child_process_spawn_request(vm, &request.args)?;
1976                self.spawn_javascript_child_process_sync(vm_id, process_id, payload, max_buffer)
1977            }
1978            "child_process.poll" => {
1979                let child_process_id =
1980                    javascript_sync_rpc_arg_str(&request.args, 0, "child_process.poll child id")?;
1981                let wait_ms = javascript_sync_rpc_arg_u64_optional(
1982                    &request.args,
1983                    1,
1984                    "child_process.poll wait ms",
1985                )?
1986                .unwrap_or_default();
1987                self.poll_javascript_child_process(vm_id, process_id, child_process_id, wait_ms)
1988            }
1989            "child_process.write_stdin" => {
1990                let child_process_id = javascript_sync_rpc_arg_str(
1991                    &request.args,
1992                    0,
1993                    "child_process.write_stdin child id",
1994                )?;
1995                let chunk = javascript_sync_rpc_bytes_arg(
1996                    &request.args,
1997                    1,
1998                    "child_process.write_stdin chunk",
1999                )?;
2000                self.write_javascript_child_process_stdin(
2001                    vm_id,
2002                    process_id,
2003                    child_process_id,
2004                    &chunk,
2005                )?;
2006                Ok(Value::Null)
2007            }
2008            "child_process.close_stdin" => {
2009                let child_process_id = javascript_sync_rpc_arg_str(
2010                    &request.args,
2011                    0,
2012                    "child_process.close_stdin child id",
2013                )?;
2014                self.close_javascript_child_process_stdin(vm_id, process_id, child_process_id)?;
2015                Ok(Value::Null)
2016            }
2017            "child_process.kill" => {
2018                let child_process_id =
2019                    javascript_sync_rpc_arg_str(&request.args, 0, "child_process.kill child id")?;
2020                let signal =
2021                    javascript_sync_rpc_arg_str(&request.args, 1, "child_process.kill signal")?;
2022                self.kill_javascript_child_process(vm_id, process_id, child_process_id, signal)?;
2023                Ok(Value::Null)
2024            }
2025            "process.kill" => {
2026                let target_pid =
2027                    javascript_sync_rpc_arg_i32(&request.args, 0, "process.kill target pid")?;
2028                let signal = javascript_sync_rpc_arg_str(&request.args, 1, "process.kill signal")?;
2029                let parsed_signal = parse_signal(signal)?;
2030                if parsed_signal == 0 {
2031                    let Some(vm) = self.vms.get(vm_id) else {
2032                        log_stale_process_event(
2033                            &self.bridge,
2034                            vm_id,
2035                            process_id,
2036                            "javascript sync RPC process.kill",
2037                        );
2038                        return Ok(());
2039                    };
2040                    if !vm.active_processes.contains_key(process_id) {
2041                        log_stale_process_event(
2042                            &self.bridge,
2043                            vm_id,
2044                            process_id,
2045                            "javascript sync RPC process.kill",
2046                        );
2047                        return Ok(());
2048                    }
2049                    vm.kernel
2050                        .signal_process(EXECUTION_DRIVER_NAME, target_pid, parsed_signal)
2051                        .map(|()| Value::Null)
2052                        .map_err(kernel_error)
2053                } else if target_pid < 0 {
2054                    let caller_kernel_pid = {
2055                        let Some(vm) = self.vms.get(vm_id) else {
2056                            log_stale_process_event(
2057                                &self.bridge,
2058                                vm_id,
2059                                process_id,
2060                                "javascript sync RPC process.kill",
2061                            );
2062                            return Ok(());
2063                        };
2064                        let Some(caller) = vm.active_processes.get(process_id) else {
2065                            log_stale_process_event(
2066                                &self.bridge,
2067                                vm_id,
2068                                process_id,
2069                                "javascript sync RPC process.kill",
2070                            );
2071                            return Ok(());
2072                        };
2073                        caller.kernel_pid
2074                    };
2075                    let pgid = target_pid.unsigned_abs();
2076                    match self.signal_vm_process_group(vm_id, caller_kernel_pid, pgid, signal) {
2077                        Ok(true) => {
2078                            Ok(self.apply_self_process_kill(vm_id, process_id, parsed_signal))
2079                        }
2080                        Ok(false) => Ok(Value::Null),
2081                        Err(error) => Err(error),
2082                    }
2083                } else {
2084                    enum ProcessKillTarget {
2085                        SelfProcess,
2086                        Child(String),
2087                        TopLevel(String),
2088                        KernelPid(u32),
2089                    }
2090                    let target = {
2091                        let Some(vm) = self.vms.get(vm_id) else {
2092                            log_stale_process_event(
2093                                &self.bridge,
2094                                vm_id,
2095                                process_id,
2096                                "javascript sync RPC process.kill",
2097                            );
2098                            return Ok(());
2099                        };
2100                        let Some(caller) = vm.active_processes.get(process_id) else {
2101                            log_stale_process_event(
2102                                &self.bridge,
2103                                vm_id,
2104                                process_id,
2105                                "javascript sync RPC process.kill",
2106                            );
2107                            return Ok(());
2108                        };
2109                        let caller_pid = i32::try_from(caller.kernel_pid).map_err(|_| {
2110                            SidecarError::InvalidState("caller pid exceeds i32".into())
2111                        })?;
2112                        if caller_pid == target_pid {
2113                            ProcessKillTarget::SelfProcess
2114                        } else if let Some((child_process_id, _)) = caller
2115                            .child_processes
2116                            .iter()
2117                            .find(|(_, child)| i32::try_from(child.kernel_pid) == Ok(target_pid))
2118                        {
2119                            ProcessKillTarget::Child(child_process_id.clone())
2120                        } else if let Some((target_process_id, _)) =
2121                            vm.active_processes.iter().find(|(_, process)| {
2122                                i32::try_from(process.kernel_pid) == Ok(target_pid)
2123                            })
2124                        {
2125                            ProcessKillTarget::TopLevel(target_process_id.clone())
2126                        } else {
2127                            let target_kernel_pid = u32::try_from(target_pid).map_err(|_| {
2128                                SidecarError::InvalidState(format!(
2129                                    "EINVAL: invalid process pid {target_pid}"
2130                                ))
2131                            })?;
2132                            ProcessKillTarget::KernelPid(target_kernel_pid)
2133                        }
2134                    };
2135                    match target {
2136                        ProcessKillTarget::SelfProcess => {
2137                            Ok(self.apply_self_process_kill(vm_id, process_id, parsed_signal))
2138                        }
2139                        ProcessKillTarget::Child(child_process_id) => {
2140                            self.kill_javascript_child_process(
2141                                vm_id,
2142                                process_id,
2143                                &child_process_id,
2144                                signal,
2145                            )?;
2146                            Ok(Value::Null)
2147                        }
2148                        ProcessKillTarget::TopLevel(target_process_id) => {
2149                            self.kill_process_internal(vm_id, &target_process_id, signal)?;
2150                            Ok(Value::Null)
2151                        }
2152                        ProcessKillTarget::KernelPid(target_kernel_pid) => {
2153                            // Grandchildren and untracked kernel processes are
2154                            // resolved VM-wide instead of failing with an
2155                            // unknown-pid error.
2156                            self.signal_vm_kernel_pid(vm_id, target_kernel_pid, signal)
2157                                .map(|()| Value::Null)
2158                        }
2159                    }
2160                }
2161            }
2162            "process.signal_state" => {
2163                let signal =
2164                    javascript_sync_rpc_arg_u32(&request.args, 0, "process.signal_state signal")?;
2165                let action =
2166                    javascript_sync_rpc_arg_str(&request.args, 1, "process.signal_state action")?;
2167                let mask_json =
2168                    javascript_sync_rpc_arg_str(&request.args, 2, "process.signal_state mask")?;
2169                let flags =
2170                    javascript_sync_rpc_arg_u32(&request.args, 3, "process.signal_state flags")?;
2171                let mask: Vec<u32> = serde_json::from_str(mask_json).map_err(|error| {
2172                    SidecarError::InvalidState(format!(
2173                        "process.signal_state mask must be valid JSON: {error}"
2174                    ))
2175                })?;
2176                let action = match action.trim().to_ascii_lowercase().as_str() {
2177                    "default" => SignalDispositionAction::Default,
2178                    "ignore" => SignalDispositionAction::Ignore,
2179                    "user" => SignalDispositionAction::User,
2180                    other => {
2181                        return Err(SidecarError::InvalidState(format!(
2182                            "unsupported process.signal_state action {other}"
2183                        )));
2184                    }
2185                };
2186                let Some(vm) = self.vms.get_mut(vm_id) else {
2187                    log_stale_process_event(
2188                        &self.bridge,
2189                        vm_id,
2190                        process_id,
2191                        "javascript sync RPC process.signal_state",
2192                    );
2193                    return Ok(());
2194                };
2195                if action == SignalDispositionAction::Default && mask.is_empty() && flags == 0 {
2196                    let remove_process_entry = vm
2197                        .signal_states
2198                        .get_mut(process_id)
2199                        .map(|handlers| {
2200                            handlers.remove(&signal);
2201                            handlers.is_empty()
2202                        })
2203                        .unwrap_or(false);
2204                    if remove_process_entry {
2205                        vm.signal_states.remove(process_id);
2206                    }
2207                } else {
2208                    vm.signal_states
2209                        .entry(process_id.to_owned())
2210                        .or_default()
2211                        .insert(
2212                            signal,
2213                            SignalHandlerRegistration {
2214                                action,
2215                                mask,
2216                                flags,
2217                            },
2218                        );
2219                }
2220                Ok(Value::Null)
2221            }
2222            "net.http_request" => {
2223                let payload = request
2224                    .args
2225                    .first()
2226                    .cloned()
2227                    .ok_or_else(|| {
2228                        SidecarError::InvalidState(String::from(
2229                            "net.http_request requires a request payload",
2230                        ))
2231                    })
2232                    .and_then(|value| {
2233                        serde_json::from_value::<JavascriptHttpLoopbackRequest>(value).map_err(
2234                            |error| {
2235                                SidecarError::InvalidState(format!(
2236                                    "invalid net.http_request payload: {error}"
2237                                ))
2238                            },
2239                        )
2240                    })?;
2241                if !is_javascript_loopback_host(&payload.host) {
2242                    return Err(SidecarError::Execution(format!(
2243                        "EACCES: HTTP loopback request requires a loopback host, got {}",
2244                        payload.host
2245                    )));
2246                }
2247                self.bridge.require_network_access(
2248                    vm_id,
2249                    NetworkOperation::Http,
2250                    format_tcp_resource(&payload.host, payload.port),
2251                )?;
2252                let Some(vm) = self.vms.get_mut(vm_id) else {
2253                    log_stale_process_event(
2254                        &self.bridge,
2255                        vm_id,
2256                        process_id,
2257                        "javascript sync RPC net.http_request",
2258                    );
2259                    return Ok(());
2260                };
2261                let resource_limits = vm.kernel.resource_limits().clone();
2262                let socket_paths = build_javascript_socket_path_context(vm)?;
2263                let target_is_current =
2264                    [JavascriptSocketFamily::Ipv4, JavascriptSocketFamily::Ipv6]
2265                        .iter()
2266                        .any(|family| {
2267                            socket_paths
2268                                .http_loopback_target(*family, payload.port)
2269                                .is_some_and(|target| {
2270                                    target.process_id == payload.process_id
2271                                        && target.server_id == payload.server_id
2272                                })
2273                        });
2274                if !target_is_current {
2275                    return Err(SidecarError::InvalidState(format!(
2276                        "unknown HTTP loopback target {}:{} for server {} in process {}",
2277                        payload.host, payload.port, payload.server_id, payload.process_id
2278                    )));
2279                }
2280                let Some(target_process) = vm.active_processes.get_mut(&payload.process_id) else {
2281                    return Err(SidecarError::InvalidState(format!(
2282                        "unknown HTTP loopback process {}",
2283                        payload.process_id
2284                    )));
2285                };
2286                dispatch_loopback_http_request(LoopbackHttpDispatchRequest {
2287                    bridge: &self.bridge,
2288                    vm_id,
2289                    dns: &vm.dns,
2290                    socket_paths: &socket_paths,
2291                    kernel: &mut vm.kernel,
2292                    process: target_process,
2293                    resource_limits: &resource_limits,
2294                    server_id: payload.server_id,
2295                    request_json: &payload.request,
2296                })
2297                .map(Value::String)
2298            }
2299            _ => {
2300                let Some(vm) = self.vms.get_mut(vm_id) else {
2301                    log_stale_process_event(
2302                        &self.bridge,
2303                        vm_id,
2304                        process_id,
2305                        "javascript sync RPC bridge dispatch",
2306                    );
2307                    return Ok(());
2308                };
2309                let resource_limits = vm.kernel.resource_limits().clone();
2310                let network_counts = vm_network_resource_counts(vm);
2311                let socket_paths = build_javascript_socket_path_context(vm)?;
2312                let Some(process) = vm.active_processes.get_mut(process_id) else {
2313                    log_stale_process_event(
2314                        &self.bridge,
2315                        vm_id,
2316                        process_id,
2317                        "javascript sync RPC bridge dispatch",
2318                    );
2319                    return Ok(());
2320                };
2321                service_javascript_sync_rpc(JavascriptSyncRpcServiceRequest {
2322                    bridge: &self.bridge,
2323                    vm_id,
2324                    dns: &vm.dns,
2325                    socket_paths: &socket_paths,
2326                    kernel: &mut vm.kernel,
2327                    process,
2328                    sync_request: &request,
2329                    resource_limits: &resource_limits,
2330                    network_counts,
2331                })
2332            }
2333        };
2334
2335        let Some(vm) = self.vms.get_mut(vm_id) else {
2336            log_stale_process_event(
2337                &self.bridge,
2338                vm_id,
2339                process_id,
2340                "javascript sync RPC response delivery",
2341            );
2342            return Ok(());
2343        };
2344        let shadow_root = vm.cwd.clone();
2345        let Some(process) = vm.active_processes.get_mut(process_id) else {
2346            log_stale_process_event(
2347                &self.bridge,
2348                vm_id,
2349                process_id,
2350                "javascript sync RPC response delivery",
2351            );
2352            return Ok(());
2353        };
2354
2355        if response.is_ok()
2356            && matches!(
2357                request.method.as_str(),
2358                "fs.chmodSync" | "fs.promises.chmod"
2359            )
2360        {
2361            let guest_path =
2362                javascript_sync_rpc_arg_str(&request.args, 0, "filesystem chmod path")?;
2363            let mode =
2364                javascript_sync_rpc_arg_u32(&request.args, 1, "filesystem chmod mode")? & 0o7777;
2365            let host_path =
2366                shadow_host_path_for_process(&shadow_root, &process.guest_cwd, guest_path);
2367            if host_path.exists() {
2368                fs::set_permissions(&host_path, fs::Permissions::from_mode(mode)).map_err(
2369                    |error| {
2370                        SidecarError::Io(format!(
2371                            "failed to mirror chmod to shadow path {}: {error}",
2372                            host_path.display()
2373                        ))
2374                    },
2375                )?;
2376            }
2377        }
2378
2379        match response {
2380            Ok(result) => process
2381                .execution
2382                .respond_javascript_sync_rpc_success(request.id, result)
2383                .or_else(ignore_stale_javascript_sync_rpc_response),
2384            Err(error) => process
2385                .execution
2386                .respond_javascript_sync_rpc_error(
2387                    request.id,
2388                    javascript_sync_rpc_error_code(&error),
2389                    error.to_string(),
2390                )
2391                .or_else(ignore_stale_javascript_sync_rpc_response),
2392        }
2393    }
2394
2395    /// Applies a `process.kill` aimed at the calling process itself and
2396    /// returns the self-delivery action payload for the bridge.
2397    fn apply_self_process_kill(
2398        &mut self,
2399        vm_id: &str,
2400        process_id: &str,
2401        parsed_signal: i32,
2402    ) -> Value {
2403        let action = self
2404            .vms
2405            .get(vm_id)
2406            .and_then(|vm| vm.signal_states.get(process_id))
2407            .and_then(|handlers| handlers.get(&(parsed_signal as u32)))
2408            .map(|registration| registration.action.clone())
2409            .unwrap_or(SignalDispositionAction::Default);
2410        if action == SignalDispositionAction::Default
2411            && parsed_signal != 0
2412            && !matches!(
2413                canonical_signal_name(parsed_signal),
2414                Some("SIGWINCH" | "SIGCHLD" | "SIGCONT" | "SIGURG")
2415            )
2416        {
2417            if let Some(vm) = self.vms.get_mut(vm_id) {
2418                if let Some(process) = vm.active_processes.get_mut(process_id) {
2419                    process.pending_self_signal_exit = Some(parsed_signal);
2420                }
2421            }
2422        }
2423        json!({
2424            "self": true,
2425            "action": match action {
2426                SignalDispositionAction::Default => "default",
2427                SignalDispositionAction::Ignore => "ignore",
2428                SignalDispositionAction::User => "user",
2429            },
2430        })
2431    }
2432
2433    pub(crate) fn vm_ids_for_scope(
2434        &self,
2435        ownership: &OwnershipScope,
2436    ) -> Result<Vec<String>, SidecarError> {
2437        match ownership {
2438            OwnershipScope::SessionOwnership(inner) => {
2439                self.require_owned_session(&inner.connection_id, &inner.session_id)?;
2440                Ok(self
2441                    .sessions
2442                    .get(&inner.session_id)
2443                    .expect("owned session should exist")
2444                    .vm_ids
2445                    .iter()
2446                    .cloned()
2447                    .collect())
2448            }
2449            OwnershipScope::VmOwnership(inner) => {
2450                self.require_owned_vm(&inner.connection_id, &inner.session_id, &inner.vm_id)?;
2451                Ok(vec![inner.vm_id.clone()])
2452            }
2453            OwnershipScope::ConnectionOwnership(..) => Err(SidecarError::InvalidState(
2454                String::from("event polling requires session or VM ownership scope"),
2455            )),
2456        }
2457    }
2458
2459    pub(crate) fn vm_ownership(&self, vm_id: &str) -> Result<OwnershipScope, SidecarError> {
2460        let vm = self
2461            .vms
2462            .get(vm_id)
2463            .ok_or_else(|| SidecarError::InvalidState(format!("unknown sidecar VM {vm_id}")))?;
2464        Ok(OwnershipScope::vm(&vm.connection_id, &vm.session_id, vm_id))
2465    }
2466
2467    pub(crate) fn vm_has_active_processes(&self, vm_id: &str) -> bool {
2468        self.vms
2469            .get(vm_id)
2470            .is_some_and(|vm| !vm.active_processes.is_empty())
2471    }
2472
2473    fn require_authenticated_connection(&self, connection_id: &str) -> Result<(), SidecarError> {
2474        if self.connections.contains_key(connection_id) {
2475            Ok(())
2476        } else {
2477            Err(SidecarError::InvalidState(format!(
2478                "connection {connection_id} has not authenticated"
2479            )))
2480        }
2481    }
2482
2483    pub(crate) fn require_owned_session(
2484        &self,
2485        connection_id: &str,
2486        session_id: &str,
2487    ) -> Result<(), SidecarError> {
2488        self.require_authenticated_connection(connection_id)?;
2489        let session = self.sessions.get(session_id).ok_or_else(|| {
2490            SidecarError::InvalidState(format!("unknown sidecar session {session_id}"))
2491        })?;
2492        if session.connection_id == connection_id {
2493            Ok(())
2494        } else {
2495            Err(SidecarError::InvalidState(format!(
2496                "session {session_id} is not owned by connection {connection_id}"
2497            )))
2498        }
2499    }
2500
2501    pub(crate) fn require_owned_vm(
2502        &self,
2503        connection_id: &str,
2504        session_id: &str,
2505        vm_id: &str,
2506    ) -> Result<(), SidecarError> {
2507        self.require_owned_session(connection_id, session_id)?;
2508        let vm = self
2509            .vms
2510            .get(vm_id)
2511            .ok_or_else(|| SidecarError::InvalidState(format!("unknown sidecar VM {vm_id}")))?;
2512        if vm.connection_id != connection_id || vm.session_id != session_id {
2513            return Err(SidecarError::InvalidState(format!(
2514                "VM {vm_id} is not owned by {connection_id}/{session_id}"
2515            )));
2516        }
2517        Ok(())
2518    }
2519
2520    fn connection_id_for(&self, ownership: &OwnershipScope) -> Result<String, SidecarError> {
2521        match ownership {
2522            OwnershipScope::ConnectionOwnership(inner) => Ok(inner.connection_id.clone()),
2523            OwnershipScope::SessionOwnership(..) | OwnershipScope::VmOwnership(..) => {
2524                Err(SidecarError::InvalidState(String::from(
2525                    "request requires connection ownership scope",
2526                )))
2527            }
2528        }
2529    }
2530
2531    fn validate_auth_token(&self, auth_token: &str) -> Result<(), SidecarError> {
2532        let Some(expected_auth_token) = self.config.expected_auth_token.as_deref() else {
2533            return Ok(());
2534        };
2535
2536        if auth_token == expected_auth_token {
2537            Ok(())
2538        } else {
2539            Err(SidecarError::Unauthorized(String::from(
2540                "authenticate request provided an invalid auth token",
2541            )))
2542        }
2543    }
2544
2545    fn allocate_connection_id(&mut self) -> String {
2546        self.next_connection_id += 1;
2547        format!("conn-{}", self.next_connection_id)
2548    }
2549
2550    fn take_matching_process_event_envelope(
2551        &mut self,
2552        vm_id: &str,
2553        process_id: &str,
2554    ) -> Result<Option<ProcessEventEnvelope>, SidecarError> {
2555        if let Some(index) = self
2556            .pending_process_events
2557            .iter()
2558            .position(|event| event.vm_id == vm_id && event.process_id == process_id)
2559        {
2560            return Ok(self.pending_process_events.remove(index));
2561        }
2562
2563        let mut matching_envelope = None;
2564        let mut deferred = Vec::new();
2565        {
2566            let pending_capacity = self.pending_process_event_capacity();
2567            let receiver = self.process_event_receiver.as_mut().ok_or_else(|| {
2568                SidecarError::InvalidState(String::from("process event receiver unavailable"))
2569            })?;
2570            loop {
2571                if deferred.len() >= pending_capacity {
2572                    if receiver.is_empty() {
2573                        break;
2574                    }
2575                    return Err(process_event_queue_overflow_error());
2576                }
2577                let envelope = match receiver.try_recv() {
2578                    Ok(envelope) => envelope,
2579                    Err(tokio::sync::mpsc::error::TryRecvError::Empty) => break,
2580                    Err(tokio::sync::mpsc::error::TryRecvError::Disconnected) => break,
2581                };
2582                if matching_envelope.is_none()
2583                    && envelope.vm_id == vm_id
2584                    && envelope.process_id == process_id
2585                {
2586                    matching_envelope = Some(envelope);
2587                    break;
2588                }
2589                deferred.push(envelope);
2590            }
2591        }
2592        for envelope in deferred {
2593            self.queue_pending_process_event(envelope)?;
2594        }
2595
2596        Ok(matching_envelope)
2597    }
2598
2599    fn allocate_sidecar_request_id(&mut self) -> RequestId {
2600        let request_id = self.next_sidecar_request_id;
2601        self.next_sidecar_request_id -= 1;
2602        request_id
2603    }
2604
2605    pub(crate) fn session_scope_for(
2606        &self,
2607        ownership: &OwnershipScope,
2608    ) -> Result<(String, String), SidecarError> {
2609        match ownership {
2610            OwnershipScope::SessionOwnership(inner) => {
2611                Ok((inner.connection_id.clone(), inner.session_id.clone()))
2612            }
2613            OwnershipScope::ConnectionOwnership(..) | OwnershipScope::VmOwnership(..) => {
2614                Err(SidecarError::InvalidState(String::from(
2615                    "request requires session ownership scope",
2616                )))
2617            }
2618        }
2619    }
2620
2621    pub(crate) fn vm_scope_for(
2622        &self,
2623        ownership: &OwnershipScope,
2624    ) -> Result<(String, String, String), SidecarError> {
2625        match ownership {
2626            OwnershipScope::VmOwnership(inner) => Ok((
2627                inner.connection_id.clone(),
2628                inner.session_id.clone(),
2629                inner.vm_id.clone(),
2630            )),
2631            OwnershipScope::ConnectionOwnership(..) | OwnershipScope::SessionOwnership(..) => Err(
2632                SidecarError::InvalidState(String::from("request requires VM ownership scope")),
2633            ),
2634        }
2635    }
2636
2637    fn response_with_ownership(
2638        &self,
2639        request_id: RequestId,
2640        ownership: OwnershipScope,
2641        payload: ResponsePayload,
2642    ) -> ResponseFrame {
2643        ResponseFrame {
2644            schema: ProtocolSchema::current(),
2645            request_id,
2646            ownership,
2647            payload,
2648        }
2649    }
2650
2651    pub(crate) fn respond(
2652        &self,
2653        request: &RequestFrame,
2654        payload: ResponsePayload,
2655    ) -> ResponseFrame {
2656        self.response_with_ownership(request.request_id, request.ownership.clone(), payload)
2657    }
2658
2659    fn reject(&self, request: &RequestFrame, code: &str, message: &str) -> ResponseFrame {
2660        self.respond(
2661            request,
2662            ResponsePayload::Rejected(RejectedResponse {
2663                code: code.to_owned(),
2664                message: message.to_owned(),
2665            }),
2666        )
2667    }
2668
2669    pub fn queue_sidecar_request(
2670        &mut self,
2671        ownership: OwnershipScope,
2672        payload: SidecarRequestPayload,
2673    ) -> Result<RequestId, SidecarError> {
2674        if self.outbound_sidecar_requests.len() >= MAX_OUTBOUND_SIDECAR_REQUESTS {
2675            return Err(outbound_sidecar_request_queue_overflow_error());
2676        }
2677        if self.pending_sidecar_responses.pending_count() >= MAX_PENDING_SIDECAR_RESPONSES {
2678            return Err(sidecar_response_pending_overflow_error());
2679        }
2680        let request_id = self.allocate_sidecar_request_id();
2681        let request = SidecarRequestFrame::new(request_id, ownership, payload);
2682        self.pending_sidecar_responses
2683            .register_request(&request)
2684            .map_err(sidecar_response_tracker_error)?;
2685        self.outbound_sidecar_requests.push_back(request);
2686        self.outbound_sidecar_requests_gauge
2687            .observe_depth(self.outbound_sidecar_requests.len());
2688        self.pending_sidecar_responses_gauge
2689            .observe_depth(self.pending_sidecar_responses.pending_count());
2690        Ok(request_id)
2691    }
2692
2693    pub fn queue_wire_sidecar_request(
2694        &mut self,
2695        ownership: crate::wire::OwnershipScope,
2696        payload: crate::wire::SidecarRequestPayload,
2697    ) -> Result<crate::wire::RequestId, SidecarError> {
2698        let ownership = crate::wire::ownership_scope_to_compat(ownership);
2699        let payload = crate::wire::sidecar_request_payload_to_compat(&ownership, payload)
2700            .map_err(wire_protocol_error)?;
2701        self.queue_sidecar_request(ownership, payload)
2702    }
2703
2704    pub fn pop_sidecar_request(&mut self) -> Option<SidecarRequestFrame> {
2705        let request = self.outbound_sidecar_requests.pop_front();
2706        self.outbound_sidecar_requests_gauge
2707            .observe_depth(self.outbound_sidecar_requests.len());
2708        request
2709    }
2710
2711    pub fn pop_wire_sidecar_request(
2712        &mut self,
2713    ) -> Result<Option<crate::wire::SidecarRequestFrame>, SidecarError> {
2714        self.pop_sidecar_request()
2715            .map(crate::wire::sidecar_request_frame_from_compat)
2716            .transpose()
2717            .map_err(wire_protocol_error)
2718    }
2719
2720    pub fn accept_sidecar_response(
2721        &mut self,
2722        response: SidecarResponseFrame,
2723    ) -> Result<(), SidecarError> {
2724        match self.pending_sidecar_responses.accept_response(&response) {
2725            Ok(()) => {}
2726            // A response for a request that is no longer pending (its owning VM
2727            // was disposed, abandoning the in-flight callback) or already
2728            // completed is a benign late/stale reply on the shared sidecar — a
2729            // per-VM `sidecar_request` can be answered by the host after that VM
2730            // has been torn down (multiple VMs share one sidecar process). Drop
2731            // it instead of failing the whole sidecar over a harmless straggler.
2732            Err(
2733                error @ (SidecarResponseTrackerError::UnmatchedResponse { .. }
2734                | SidecarResponseTrackerError::DuplicateResponse { .. }),
2735            ) => {
2736                tracing::warn!(
2737                    request_id = response.request_id,
2738                    "dropping stale sidecar response with no matching pending request: {error}"
2739                );
2740                return Ok(());
2741            }
2742            Err(error) => return Err(sidecar_response_tracker_error(error)),
2743        }
2744        self.pending_sidecar_responses_gauge
2745            .observe_depth(self.pending_sidecar_responses.pending_count());
2746        self.completed_sidecar_response_order
2747            .push_back(response.request_id);
2748        self.completed_sidecar_responses
2749            .insert(response.request_id, response);
2750        self.completed_sidecar_responses_gauge
2751            .observe_depth(self.completed_sidecar_responses.len());
2752        while self.completed_sidecar_responses.len() > MAX_COMPLETED_SIDECAR_RESPONSES {
2753            match self.completed_sidecar_response_order.pop_front() {
2754                // Only a response that was never retrieved is a real loss; an id
2755                // already taken via take_sidecar_response leaves a stale order
2756                // entry that removes to None and is not a dropped response.
2757                Some(evicted) => {
2758                    if self.completed_sidecar_responses.remove(&evicted).is_some() {
2759                        tracing::warn!(
2760                            queue = "completed_sidecar_responses",
2761                            evicted_request_id = evicted,
2762                            capacity = MAX_COMPLETED_SIDECAR_RESPONSES,
2763                            "dropping an unretrieved completed sidecar response to stay within cap; the host can no longer fetch it (response lost)"
2764                        );
2765                        self.completed_sidecar_responses_gauge
2766                            .observe_depth(self.completed_sidecar_responses.len());
2767                    }
2768                }
2769                None => break,
2770            }
2771        }
2772        Ok(())
2773    }
2774
2775    pub fn accept_wire_sidecar_response(
2776        &mut self,
2777        response: crate::wire::SidecarResponseFrame,
2778    ) -> Result<(), SidecarError> {
2779        let response =
2780            crate::wire::sidecar_response_frame_to_compat(response).map_err(wire_protocol_error)?;
2781        self.accept_sidecar_response(response)
2782    }
2783
2784    pub fn take_sidecar_response(&mut self, request_id: RequestId) -> Option<SidecarResponseFrame> {
2785        let response = self.completed_sidecar_responses.remove(&request_id);
2786        if response.is_some() {
2787            self.completed_sidecar_response_order
2788                .retain(|completed_id| completed_id != &request_id);
2789            self.completed_sidecar_responses_gauge
2790                .observe_depth(self.completed_sidecar_responses.len());
2791        }
2792        response
2793    }
2794
2795    pub fn take_wire_sidecar_response(
2796        &mut self,
2797        request_id: crate::wire::RequestId,
2798    ) -> Result<Option<crate::wire::SidecarResponseFrame>, SidecarError> {
2799        self.take_sidecar_response(request_id)
2800            .map(|response| {
2801                crate::wire::sidecar_response_frame_from_compat(response)
2802                    .map_err(wire_protocol_error)
2803            })
2804            .transpose()
2805    }
2806
2807    pub(crate) fn vm_lifecycle_event(
2808        &self,
2809        connection_id: &str,
2810        session_id: &str,
2811        vm_id: &str,
2812        state: VmLifecycleState,
2813    ) -> EventFrame {
2814        EventFrame::new(
2815            OwnershipScope::vm(connection_id, session_id, vm_id),
2816            EventPayload::VmLifecycle(VmLifecycleEvent { state }),
2817        )
2818    }
2819
2820    fn ensure_request_within_frame_limit(
2821        &self,
2822        request: &RequestFrame,
2823    ) -> Result<(), SidecarError> {
2824        let frame = crate::protocol::to_generated_protocol_frame(
2825            &crate::protocol::ProtocolFrame::Request(request.clone()),
2826        )
2827        .map_err(|error| {
2828            SidecarError::InvalidState(format!("failed to convert request frame: {error}"))
2829        })?;
2830        let crate::wire::ProtocolFrame::RequestFrame(_) = &frame else {
2831            return Err(SidecarError::InvalidState(String::from(
2832                "request converted to non-request wire frame",
2833            )));
2834        };
2835
2836        crate::wire::WireFrameCodec::new(self.config.max_frame_bytes)
2837            .encode(&frame)
2838            .map(|_| ())
2839            .map_err(|error| SidecarError::FrameTooLarge(error.to_string()))
2840    }
2841}
2842
2843impl<B> ExtensionHost for NativeSidecar<B>
2844where
2845    B: NativeSidecarBridge + Send + 'static,
2846    BridgeError<B>: fmt::Debug + Send + Sync + 'static,
2847{
2848    fn spawn_process<'a>(
2849        &'a mut self,
2850        ownership: OwnershipScope,
2851        payload: ExecuteRequest,
2852    ) -> ExtensionFuture<'a, ProcessStartedResponse> {
2853        Box::pin(async move {
2854            let request = RequestFrame::new(0, ownership, RequestPayload::Execute(payload.clone()));
2855            let dispatch = NativeSidecar::execute(self, &request, payload).await?;
2856            match dispatch.response.payload {
2857                ResponsePayload::ProcessStarted(response) => Ok(response),
2858                other => Err(unexpected_extension_host_response("execute", other)),
2859            }
2860        })
2861    }
2862
2863    fn write_stdin<'a>(
2864        &'a mut self,
2865        ownership: OwnershipScope,
2866        payload: WriteStdinRequest,
2867    ) -> ExtensionFuture<'a, StdinWrittenResponse> {
2868        Box::pin(async move {
2869            let request =
2870                RequestFrame::new(0, ownership, RequestPayload::WriteStdin(payload.clone()));
2871            let dispatch = NativeSidecar::write_stdin(self, &request, payload).await?;
2872            match dispatch.response.payload {
2873                ResponsePayload::StdinWritten(response) => Ok(response),
2874                other => Err(unexpected_extension_host_response("write_stdin", other)),
2875            }
2876        })
2877    }
2878
2879    fn close_stdin<'a>(
2880        &'a mut self,
2881        ownership: OwnershipScope,
2882        payload: CloseStdinRequest,
2883    ) -> ExtensionFuture<'a, StdinClosedResponse> {
2884        Box::pin(async move {
2885            let request =
2886                RequestFrame::new(0, ownership, RequestPayload::CloseStdin(payload.clone()));
2887            let dispatch = NativeSidecar::close_stdin(self, &request, payload).await?;
2888            match dispatch.response.payload {
2889                ResponsePayload::StdinClosed(response) => Ok(response),
2890                other => Err(unexpected_extension_host_response("close_stdin", other)),
2891            }
2892        })
2893    }
2894
2895    fn kill_process<'a>(
2896        &'a mut self,
2897        ownership: OwnershipScope,
2898        payload: KillProcessRequest,
2899    ) -> ExtensionFuture<'a, ProcessKilledResponse> {
2900        Box::pin(async move {
2901            let request =
2902                RequestFrame::new(0, ownership, RequestPayload::KillProcess(payload.clone()));
2903            let dispatch = NativeSidecar::kill_process(self, &request, payload).await?;
2904            match dispatch.response.payload {
2905                ResponsePayload::ProcessKilled(response) => Ok(response),
2906                other => Err(unexpected_extension_host_response("kill_process", other)),
2907            }
2908        })
2909    }
2910
2911    fn poll_event<'a>(
2912        &'a mut self,
2913        ownership: OwnershipScope,
2914        timeout: Duration,
2915    ) -> ExtensionFuture<'a, Option<EventFrame>> {
2916        Box::pin(async move { NativeSidecar::poll_event(self, &ownership, timeout).await })
2917    }
2918
2919    fn guest_filesystem_call<'a>(
2920        &'a mut self,
2921        ownership: OwnershipScope,
2922        payload: GuestFilesystemCallRequest,
2923    ) -> ExtensionFuture<'a, GuestFilesystemResultResponse> {
2924        Box::pin(async move {
2925            let request = RequestFrame::new(
2926                0,
2927                ownership,
2928                RequestPayload::GuestFilesystemCall(payload.clone()),
2929            );
2930            let dispatch = NativeSidecar::guest_filesystem_call(self, &request, payload).await?;
2931            match dispatch.response.payload {
2932                ResponsePayload::GuestFilesystemResult(response) => Ok(response),
2933                other => Err(unexpected_extension_host_response(
2934                    "guest_filesystem_call",
2935                    other,
2936                )),
2937            }
2938        })
2939    }
2940
2941    fn bind_process_to_session<'a>(
2942        &'a mut self,
2943        ownership: OwnershipScope,
2944        namespace: String,
2945        ext_session_id: String,
2946        process_id: String,
2947    ) -> ExtensionFuture<'a, ()> {
2948        Box::pin(async move {
2949            self.bind_extension_process_resource(ownership, namespace, ext_session_id, process_id)
2950        })
2951    }
2952
2953    fn bind_vm_to_session<'a>(
2954        &'a mut self,
2955        ownership: OwnershipScope,
2956        namespace: String,
2957        ext_session_id: String,
2958    ) -> ExtensionFuture<'a, ()> {
2959        Box::pin(
2960            async move { self.bind_extension_vm_resource(ownership, namespace, ext_session_id) },
2961        )
2962    }
2963
2964    fn dispose_session_resources<'a>(
2965        &'a mut self,
2966        ownership: OwnershipScope,
2967        namespace: String,
2968        ext_session_id: String,
2969    ) -> ExtensionFuture<'a, Vec<EventFrame>> {
2970        Box::pin(async move {
2971            let key = (namespace, ext_session_id);
2972            let Some(resources) = self.extension_sessions.get(&key) else {
2973                return Ok(Vec::new());
2974            };
2975            if resources.ownership != ownership {
2976                return Err(SidecarError::InvalidState(String::from(
2977                    "extension session ownership did not match dispose request",
2978                )));
2979            }
2980            let resources = self
2981                .extension_sessions
2982                .remove(&key)
2983                .expect("extension resources existed before removal");
2984            let (connection_id, session_id, vm_id) = self.vm_scope_for(&ownership)?;
2985            for process_id in resources.process_ids {
2986                if self
2987                    .vms
2988                    .get(&vm_id)
2989                    .is_some_and(|vm| vm.active_processes.contains_key(&process_id))
2990                {
2991                    self.kill_process_internal(&vm_id, &process_id, "SIGTERM")?;
2992                }
2993            }
2994            let mut events = Vec::new();
2995            for resource_vm_id in resources.vm_ids {
2996                if self.vms.contains_key(&resource_vm_id) {
2997                    events.extend(
2998                        self.dispose_vm_internal(
2999                            &connection_id,
3000                            &session_id,
3001                            &resource_vm_id,
3002                            DisposeReason::Requested,
3003                        )
3004                        .await?,
3005                    );
3006                }
3007            }
3008            Ok(events)
3009        })
3010    }
3011
3012    fn start_buffering_process_output<'a>(
3013        &'a mut self,
3014        ownership: OwnershipScope,
3015        process_id: String,
3016    ) -> ExtensionFuture<'a, ()> {
3017        Box::pin(async move {
3018            let (connection_id, session_id, vm_id) = self.vm_scope_for(&ownership)?;
3019            self.require_owned_vm(&connection_id, &session_id, &vm_id)?;
3020            let key = (vm_id, process_id);
3021            if self.extension_process_output_buffers.contains_key(&key) {
3022                return Err(SidecarError::Conflict(String::from(
3023                    "extension process output buffering already started",
3024                )));
3025            }
3026            self.extension_process_output_buffers
3027                .insert(key, ExtensionBufferedProcessOutput::default());
3028            Ok(())
3029        })
3030    }
3031
3032    fn handoff_buffered_process_output<'a>(
3033        &'a mut self,
3034        ownership: OwnershipScope,
3035        namespace: String,
3036        ext_session_id: String,
3037        process_id: String,
3038        timeout: Duration,
3039    ) -> ExtensionFuture<'a, ExtensionBufferedProcessOutput> {
3040        Box::pin(async move {
3041            let (connection_id, session_id, vm_id) = self.vm_scope_for(&ownership)?;
3042            self.require_owned_vm(&connection_id, &session_id, &vm_id)?;
3043            let key = (vm_id.clone(), process_id.clone());
3044            let deadline = Instant::now() + timeout;
3045            loop {
3046                self.pump_process_events(&ownership).await?;
3047                while let Some(envelope) =
3048                    self.take_matching_process_event_envelope(&vm_id, &process_id)?
3049                {
3050                    if self.capture_extension_process_output_event(
3051                        &vm_id,
3052                        &process_id,
3053                        &envelope.event,
3054                    ) {
3055                        continue;
3056                    }
3057                    self.queue_pending_process_event(envelope)?;
3058                    break;
3059                }
3060                let buffered = self
3061                    .extension_process_output_buffers
3062                    .get(&key)
3063                    .is_some_and(|buffer| !buffer.stdout.is_empty() || !buffer.stderr.is_empty());
3064                if buffered || timeout.is_zero() || Instant::now() >= deadline {
3065                    break;
3066                }
3067                let remaining = deadline.saturating_duration_since(Instant::now());
3068                time::sleep(remaining.min(Duration::from_millis(10))).await;
3069            }
3070            self.bind_extension_process_resource(
3071                ownership,
3072                namespace,
3073                ext_session_id,
3074                process_id.clone(),
3075            )?;
3076            self.extension_process_output_buffers
3077                .remove(&key)
3078                .ok_or_else(|| {
3079                    SidecarError::InvalidState(String::from(
3080                        "extension process output buffering was not started",
3081                    ))
3082                })
3083        })
3084    }
3085}
3086
3087fn unexpected_extension_host_response(operation: &str, payload: ResponsePayload) -> SidecarError {
3088    match payload {
3089        ResponsePayload::Rejected(response) => SidecarError::InvalidState(format!(
3090            "extension {operation} rejected with {}: {}",
3091            response.code, response.message
3092        )),
3093        other => SidecarError::InvalidState(format!(
3094            "extension {operation} returned unexpected response: {other:?}"
3095        )),
3096    }
3097}
3098
3099fn shadow_host_path_for_process(
3100    shadow_root: &Path,
3101    process_guest_cwd: &str,
3102    guest_path: &str,
3103) -> PathBuf {
3104    let normalized_guest_path = if guest_path.starts_with('/') {
3105        normalize_path(guest_path)
3106    } else {
3107        normalize_path(&format!(
3108            "{}/{}",
3109            process_guest_cwd.trim_end_matches('/'),
3110            guest_path
3111        ))
3112    };
3113    if normalized_guest_path == "/" {
3114        shadow_root.to_path_buf()
3115    } else {
3116        shadow_root.join(normalized_guest_path.trim_start_matches('/'))
3117    }
3118}
3119
3120fn sidecar_response_tracker_error(error: SidecarResponseTrackerError) -> SidecarError {
3121    SidecarError::InvalidState(format!(
3122        "invalid sidecar response correlation state: {error}"
3123    ))
3124}
3125
3126fn map_bridge_permission(decision: secure_exec_bridge::PermissionDecision) -> PermissionDecision {
3127    match decision.verdict {
3128        secure_exec_bridge::PermissionVerdict::Allow => PermissionDecision::allow(),
3129        secure_exec_bridge::PermissionVerdict::Deny => PermissionDecision::deny(
3130            decision
3131                .reason
3132                .unwrap_or_else(|| String::from("denied by host")),
3133        ),
3134        secure_exec_bridge::PermissionVerdict::Prompt => PermissionDecision::deny(
3135            decision
3136                .reason
3137                .unwrap_or_else(|| String::from("permission prompt required")),
3138        ),
3139    }
3140}
3141
3142fn audit_timestamp() -> String {
3143    SystemTime::now()
3144        .duration_since(UNIX_EPOCH)
3145        .expect("system time before unix epoch")
3146        .as_millis()
3147        .to_string()
3148}
3149
3150pub(crate) fn audit_fields<I, K, V>(fields: I) -> BTreeMap<String, String>
3151where
3152    I: IntoIterator<Item = (K, V)>,
3153    K: Into<String>,
3154    V: Into<String>,
3155{
3156    let mut mapped = BTreeMap::from([(String::from("timestamp"), audit_timestamp())]);
3157    for (key, value) in fields {
3158        mapped.insert(key.into(), value.into());
3159    }
3160    mapped
3161}
3162
3163pub(crate) fn emit_structured_event<B>(
3164    bridge: &SharedBridge<B>,
3165    vm_id: &str,
3166    name: &str,
3167    fields: BTreeMap<String, String>,
3168) -> Result<(), SidecarError>
3169where
3170    B: NativeSidecarBridge + Send + 'static,
3171    BridgeError<B>: fmt::Debug + Send + Sync + 'static,
3172{
3173    bridge.with_mut(|bridge| {
3174        bridge.emit_structured_event(StructuredEventRecord {
3175            vm_id: vm_id.to_owned(),
3176            name: name.to_owned(),
3177            fields,
3178        })
3179    })
3180}
3181
3182pub(crate) fn emit_security_audit_event<B>(
3183    bridge: &SharedBridge<B>,
3184    vm_id: &str,
3185    name: &str,
3186    fields: BTreeMap<String, String>,
3187) where
3188    B: NativeSidecarBridge + Send + 'static,
3189    BridgeError<B>: fmt::Debug + Send + Sync + 'static,
3190{
3191    let _ = emit_structured_event(bridge, vm_id, name, fields);
3192}
3193
3194/// Build a wire `EventFrame` carrying a `StructuredEvent` (name + string-map
3195/// detail) scoped to a connection. Used to forward limit-registry warnings to the
3196/// host as `{type:"structured", name:"limit_warning", detail}` events without a
3197/// protocol schema change. Emitted directly to the host (not via the polled,
3198/// per-session bridge queue, which is a no-op in the stdio sidecar), so a
3199/// process-global signal is delivered against the active connection.
3200pub(crate) fn structured_event_frame(
3201    connection_id: &str,
3202    name: &str,
3203    detail: std::collections::HashMap<String, String>,
3204) -> Result<crate::wire::EventFrame, SidecarError> {
3205    let event = EventFrame::new(
3206        OwnershipScope::connection(connection_id),
3207        EventPayload::Structured(crate::protocol::StructuredEvent {
3208            name: name.to_owned(),
3209            detail,
3210        }),
3211    );
3212    crate::wire::event_frame_from_compat(event).map_err(|error| {
3213        SidecarError::InvalidState(format!("invalid structured event frame: {error}"))
3214    })
3215}
3216
3217pub(crate) fn log_stale_process_event<B>(
3218    bridge: &SharedBridge<B>,
3219    vm_id: &str,
3220    process_id: &str,
3221    context: &str,
3222) where
3223    B: NativeSidecarBridge + Send + 'static,
3224    BridgeError<B>: fmt::Debug + Send + Sync + 'static,
3225{
3226    let _ = bridge.emit_log(
3227        vm_id,
3228        format!(
3229            "Ignoring stale process event during {context}: VM {vm_id} process {process_id} was already reaped"
3230        ),
3231    );
3232}
3233
3234// filesystem_operation_label moved to crate::vm
3235
3236pub(crate) fn root_filesystem_error(error: impl std::fmt::Display) -> SidecarError {
3237    SidecarError::InvalidState(format!("root filesystem: {error}"))
3238}
3239
3240pub(crate) fn normalize_path(path: &str) -> String {
3241    let mut segments = Vec::new();
3242    for component in Path::new(path).components() {
3243        match component {
3244            Component::RootDir => segments.clear(),
3245            Component::ParentDir => {
3246                segments.pop();
3247            }
3248            Component::CurDir => {}
3249            Component::Normal(value) => segments.push(value.to_string_lossy().into_owned()),
3250            Component::Prefix(prefix) => {
3251                segments.push(prefix.as_os_str().to_string_lossy().into_owned());
3252            }
3253        }
3254    }
3255
3256    let normalized = format!("/{}", segments.join("/"));
3257    if normalized.is_empty() {
3258        String::from("/")
3259    } else {
3260        normalized
3261    }
3262}
3263
3264pub(crate) fn normalize_host_path(path: &Path) -> PathBuf {
3265    let mut normalized = PathBuf::new();
3266
3267    for component in path.components() {
3268        match component {
3269            Component::Prefix(prefix) => normalized.push(prefix.as_os_str()),
3270            Component::RootDir => normalized.push(Path::new("/")),
3271            Component::CurDir => {}
3272            Component::ParentDir => {
3273                if normalized != Path::new("/") {
3274                    normalized.pop();
3275                }
3276            }
3277            Component::Normal(part) => normalized.push(part),
3278        }
3279    }
3280
3281    if normalized.as_os_str().is_empty() {
3282        if path.is_absolute() {
3283            PathBuf::from("/")
3284        } else {
3285            PathBuf::from(".")
3286        }
3287    } else {
3288        normalized
3289    }
3290}
3291
3292pub(crate) fn path_is_within_root(path: &Path, root: &Path) -> bool {
3293    path == root || path.starts_with(root)
3294}
3295
3296pub(crate) fn dirname(path: &str) -> String {
3297    let normalized = normalize_path(path);
3298    let parent = Path::new(&normalized)
3299        .parent()
3300        .unwrap_or_else(|| Path::new("/"));
3301    let value = parent.to_string_lossy();
3302    if value.is_empty() {
3303        String::from("/")
3304    } else {
3305        value.into_owned()
3306    }
3307}
3308
3309pub(crate) fn kernel_error(error: KernelError) -> SidecarError {
3310    SidecarError::Kernel(error.to_string())
3311}
3312
3313pub(crate) fn plugin_error(error: PluginError) -> SidecarError {
3314    SidecarError::Plugin(error.to_string())
3315}
3316
3317pub(crate) fn javascript_error(error: JavascriptExecutionError) -> SidecarError {
3318    SidecarError::Execution(error.to_string())
3319}
3320
3321pub(crate) fn wasm_error(error: WasmExecutionError) -> SidecarError {
3322    SidecarError::Execution(error.to_string())
3323}
3324
3325pub(crate) fn python_error(error: PythonExecutionError) -> SidecarError {
3326    SidecarError::Execution(error.to_string())
3327}
3328
3329pub(crate) fn vfs_error(error: VfsError) -> SidecarError {
3330    SidecarError::Kernel(error.to_string())
3331}
3332
3333/// Actionable guidance shown when guest package resolution fails because the packages live in a
3334/// non-flat `node_modules` whose package store is not visible in the VM. Mounting host `node_modules`
3335/// is a bind mount, so symlinked/store layouts
3336/// do not resolve inside the VM: Node canonicalizes a module to its store
3337/// realpath (e.g. `node_modules/.pnpm/...`, `.bun/...`, `.store/...`) which lives
3338/// above the mounted directory and the guest `fs` cannot read. Plug'n'Play
3339/// (yarn-berry default) has no `node_modules` at all. A flat (hoisted) layout is
3340/// required. The empirically-supported package managers are captured in
3341/// `crates/sidecar/tests/module_layout_e2e.rs`.
3342#[allow(dead_code)]
3343const HOISTED_NODE_MODULES_GUIDANCE: &str = "secure-exec can't load mounted node_modules: the directory uses a non-flat layout (pnpm / bun / yarn workspaces store, or yarn Plug'n'Play) whose package store isn't visible inside the VM. A flat (hoisted) node_modules is required.\n  - pnpm        -> add `node-linker=hoisted` to .npmrc, then reinstall\n  - yarn berry  -> set `nodeLinker: node-modules` in .yarnrc.yml (not pnp/pnpm)\n  - bun         -> install dependencies outside a workspace (workspaces use a .bun store)\n  - npm / yarn classic -> already flat, no change needed";
3344
3345/// Detect, from an adapter's captured stderr, a non-flat-`node_modules` failure
3346/// signature. Returns the actionable guidance to fold into the surfaced error,
3347/// or `None` when the failure is unrelated.
3348///
3349/// Two signatures, both kept specific so they never fire on unrelated crashes:
3350/// - a missing-file / cannot-resolve error referencing a package STORE path that
3351///   lives above the mounted project (`.pnpm`, `.bun`, `.store`, PnP `__virtual__`),
3352/// - a yarn Plug'n'Play fingerprint (`.pnp.cjs`, the zip cache, or PnP's
3353///   "isn't declared in your dependencies" resolver error).
3354#[allow(dead_code)]
3355fn symlinked_node_modules_hint(stderr: &str) -> Option<&'static str> {
3356    // Package stores that only appear in a path when a non-flat layout is used.
3357    // pnpm (isolated), bun (workspace), yarn-berry (nodeLinker: pnpm), and PnP
3358    // virtual instances all keep real package files under these store dirs, which
3359    // sit above the mounted project node_modules and so are not guest-visible.
3360    const STORE_MARKERS: &[&str] = &[
3361        "node_modules/.pnpm/",
3362        "node_modules/.bun/",
3363        "node_modules/.store/",
3364        "/__virtual__/",
3365    ];
3366    // Yarn Plug'n'Play has no node_modules at all; resolution fails against the
3367    // .pnp runtime / zip cache. "isn't declared in your dependencies" is PnP's
3368    // distinctive resolver error and is specific enough to fire on its own.
3369    const PNP_STRICT_MARKERS: &[&str] = &["isn't declared in your dependencies"];
3370    const PNP_PATH_MARKERS: &[&str] = &[".pnp.cjs", ".pnp.loader.mjs", "/.yarn/cache/"];
3371
3372    if PNP_STRICT_MARKERS.iter().any(|m| stderr.contains(m)) {
3373        return Some(HOISTED_NODE_MODULES_GUIDANCE);
3374    }
3375
3376    let missing = stderr.contains("ENOENT")
3377        || stderr.contains("no such file or directory")
3378        || stderr.contains("Cannot find module")
3379        || stderr.contains("MODULE_NOT_FOUND");
3380    if !missing {
3381        return None;
3382    }
3383    if STORE_MARKERS.iter().any(|m| stderr.contains(m))
3384        || PNP_PATH_MARKERS.iter().any(|m| stderr.contains(m))
3385    {
3386        return Some(HOISTED_NODE_MODULES_GUIDANCE);
3387    }
3388    None
3389}
3390
3391#[cfg(test)]
3392mod symlinked_node_modules_hint_tests {
3393    use super::symlinked_node_modules_hint;
3394
3395    // Positive cases: each non-flat package manager's store/PnP signature.
3396    #[test]
3397    fn matches_pnpm_store_enoent() {
3398        // Real pi-coding-agent failure: getPackageDir() falls back to a
3399        // dist/package.json inside the unreachable .pnpm store.
3400        let stderr = "Error: ENOENT: no such file or directory, open '/root/node_modules/.pnpm/@mariozechner+pi-coding-agent@0.60.0_x/node_modules/@mariozechner/pi-coding-agent/dist/package.json'";
3401        let hint = symlinked_node_modules_hint(stderr).expect("expected hoisted guidance");
3402        assert!(hint.contains("secure-exec can't load mounted node_modules"));
3403        assert!(!hint.contains("agentos"));
3404    }
3405
3406    #[test]
3407    fn matches_bun_store_enoent() {
3408        let stderr = "Error: ENOENT: no such file or directory, open '/root/node_modules/.bun/is-odd@3.0.1/node_modules/is-odd/package.json'";
3409        assert!(symlinked_node_modules_hint(stderr).is_some());
3410    }
3411
3412    #[test]
3413    fn matches_yarn_pnpm_store_enoent() {
3414        let stderr = "Error: ENOENT: no such file or directory, open '/root/node_modules/.store/is-odd-npm-3.0.1-93c3c3f41b/package/package.json'";
3415        assert!(symlinked_node_modules_hint(stderr).is_some());
3416    }
3417
3418    #[test]
3419    fn matches_pnp_declared_error() {
3420        // Yarn PnP's distinctive resolver error (no node_modules at all).
3421        let stderr = "Error: Your application tried to access is-number, but it isn't declared in your dependencies; this makes the require call ambiguous and unsound.";
3422        assert!(symlinked_node_modules_hint(stderr).is_some());
3423    }
3424
3425    #[test]
3426    fn matches_pnp_cjs_module_not_found() {
3427        let stderr = "Error: Cannot find module 'is-odd'\n    at /root/.pnp.cjs:12345:18\n    code: 'MODULE_NOT_FOUND'";
3428        assert!(symlinked_node_modules_hint(stderr).is_some());
3429    }
3430
3431    #[test]
3432    fn matches_virtual_instance() {
3433        let stderr = "Error: ENOENT: no such file or directory, open '/root/.yarn/__virtual__/is-odd-abc/1/node_modules/is-odd/package.json'";
3434        assert!(symlinked_node_modules_hint(stderr).is_some());
3435    }
3436
3437    // Negative cases: must not fire.
3438    #[test]
3439    fn ignores_enoent_outside_a_store() {
3440        let stderr = "Error: ENOENT: no such file or directory, open '/tmp/scratch/config.json'";
3441        assert!(symlinked_node_modules_hint(stderr).is_none());
3442    }
3443
3444    #[test]
3445    fn ignores_store_path_without_missing_file() {
3446        let stderr =
3447            "loaded /root/node_modules/.pnpm/some-pkg@1.0.0/node_modules/some-pkg/index.js";
3448        assert!(symlinked_node_modules_hint(stderr).is_none());
3449    }
3450
3451    #[test]
3452    fn ignores_flat_node_modules_enoent() {
3453        // npm / yarn-nm / pnpm-hoisted: flat, no store dir in the path.
3454        let stderr = "Error: ENOENT: no such file or directory, open '/root/node_modules/is-odd/missing-asset.json'";
3455        assert!(symlinked_node_modules_hint(stderr).is_none());
3456    }
3457
3458    #[test]
3459    fn ignores_unrelated_failure() {
3460        let stderr = "Error: connect ECONNREFUSED 127.0.0.1:443";
3461        assert!(symlinked_node_modules_hint(stderr).is_none());
3462    }
3463}
3464
3465#[cfg(test)]
3466mod structured_event_frame_tests {
3467    use super::*;
3468
3469    #[test]
3470    fn structured_event_frame_round_trips_limit_warning() {
3471        let mut detail = std::collections::HashMap::new();
3472        // Pin a real emitted limit name rather than a fictional string.
3473        let limit_name = TrackedLimit::JavascriptEventChannel.as_str();
3474        detail.insert(String::from("limit"), String::from(limit_name));
3475        detail.insert(String::from("fillPercent"), String::from("82"));
3476
3477        let wire = structured_event_frame("conn-1", "limit_warning", detail)
3478            .expect("build structured event frame");
3479        let compat = crate::wire::event_frame_to_compat(wire).expect("convert to compat");
3480
3481        match compat.payload {
3482            EventPayload::Structured(event) => {
3483                assert_eq!(event.name, "limit_warning");
3484                assert_eq!(
3485                    event.detail.get("limit").map(String::as_str),
3486                    Some(limit_name)
3487                );
3488                assert_eq!(
3489                    event.detail.get("fillPercent").map(String::as_str),
3490                    Some("82")
3491                );
3492            }
3493            other => panic!("expected structured payload, got {other:?}"),
3494        }
3495        match compat.ownership {
3496            OwnershipScope::ConnectionOwnership(inner) => {
3497                assert_eq!(inner.connection_id, "conn-1");
3498            }
3499            other => panic!("expected connection ownership, got {other:?}"),
3500        }
3501    }
3502}
3503
3504#[cfg(test)]
3505mod dispose_lifecycle_tests {
3506    use super::*;
3507    use crate::extension::ExtensionResponse;
3508    use crate::stdio::LocalBridge;
3509    use std::sync::atomic::{AtomicUsize, Ordering};
3510
3511    fn block_on<F: std::future::Future>(future: F) -> F::Output {
3512        tokio::runtime::Builder::new_current_thread()
3513            .enable_all()
3514            .build()
3515            .expect("dispose lifecycle test runtime")
3516            .block_on(future)
3517    }
3518
3519    fn test_sidecar() -> NativeSidecar<LocalBridge> {
3520        NativeSidecar::new(LocalBridge::default()).expect("build test sidecar")
3521    }
3522
3523    // Register a connection + session directly so the dispose paths can be
3524    // exercised without spinning up a V8-backed VM.
3525    fn insert_session(
3526        sidecar: &mut NativeSidecar<LocalBridge>,
3527        connection_id: &str,
3528        session_id: &str,
3529        vm_ids: BTreeSet<String>,
3530    ) {
3531        sidecar.connections.insert(
3532            connection_id.to_string(),
3533            ConnectionState {
3534                auth_token: String::new(),
3535                sessions: BTreeSet::from([session_id.to_string()]),
3536            },
3537        );
3538        sidecar.sessions.insert(
3539            session_id.to_string(),
3540            SessionState {
3541                connection_id: connection_id.to_string(),
3542                placement: crate::protocol::SidecarPlacement::SidecarPlacementShared(
3543                    crate::protocol::SidecarPlacementShared { pool: None },
3544                ),
3545                metadata: BTreeMap::new(),
3546                vm_ids,
3547            },
3548        );
3549    }
3550
3551    struct RecordingExtension {
3552        namespace: String,
3553        session_disposed: Arc<AtomicUsize>,
3554    }
3555
3556    impl Extension for RecordingExtension {
3557        fn namespace(&self) -> &str {
3558            &self.namespace
3559        }
3560
3561        fn handle_request<'a>(
3562            &'a self,
3563            _ctx: ExtensionContext<'a>,
3564            _payload: Vec<u8>,
3565        ) -> ExtensionFuture<'a, ExtensionResponse> {
3566            Box::pin(async { Ok(ExtensionResponse::new(Vec::new())) })
3567        }
3568
3569        fn on_session_disposed<'a>(&'a self, _ctx: ExtensionSnapshot) -> ExtensionFuture<'a, ()> {
3570            let counter = self.session_disposed.clone();
3571            Box::pin(async move {
3572                counter.fetch_add(1, Ordering::SeqCst);
3573                Ok(())
3574            })
3575        }
3576    }
3577
3578    fn register_recording_extension(sidecar: &mut NativeSidecar<LocalBridge>) -> Arc<AtomicUsize> {
3579        let counter = Arc::new(AtomicUsize::new(0));
3580        sidecar
3581            .register_extension(Box::new(RecordingExtension {
3582                namespace: String::from("dev.test.dispose"),
3583                session_disposed: counter.clone(),
3584            }))
3585            .expect("register recording extension");
3586        counter
3587    }
3588
3589    // H4: the extension per-session teardown hook fires on ConnectionClosed so an
3590    // ACP-style extension can release per-session state on client disconnect.
3591    #[test]
3592    fn connection_closed_dispose_invokes_extension_session_teardown() {
3593        let mut sidecar = test_sidecar();
3594        let counter = register_recording_extension(&mut sidecar);
3595        insert_session(&mut sidecar, "conn-1", "session-1", BTreeSet::new());
3596
3597        block_on(sidecar.dispose_session("conn-1", "session-1", DisposeReason::ConnectionClosed))
3598            .expect("dispose session on connection close");
3599
3600        assert_eq!(
3601            counter.load(Ordering::SeqCst),
3602            1,
3603            "extension session-teardown hook must fire on ConnectionClosed"
3604        );
3605        assert!(
3606            !sidecar.sessions.contains_key("session-1"),
3607            "the disposed session must be reclaimed"
3608        );
3609    }
3610
3611    // H4 (negative): a client-requested dispose is not a disconnect, so the
3612    // teardown hook must not fire.
3613    #[test]
3614    fn requested_dispose_does_not_invoke_extension_session_teardown() {
3615        let mut sidecar = test_sidecar();
3616        let counter = register_recording_extension(&mut sidecar);
3617        insert_session(&mut sidecar, "conn-1", "session-1", BTreeSet::new());
3618
3619        block_on(sidecar.dispose_session("conn-1", "session-1", DisposeReason::Requested))
3620            .expect("dispose session on request");
3621
3622        assert_eq!(
3623            counter.load(Ordering::SeqCst),
3624            0,
3625            "the teardown hook is reserved for client disconnect"
3626        );
3627    }
3628
3629    // M5: disposing a session records its scope for the stdio transport to drain.
3630    #[test]
3631    fn dispose_session_records_disposed_scope() {
3632        let mut sidecar = test_sidecar();
3633        insert_session(&mut sidecar, "conn-1", "session-1", BTreeSet::new());
3634
3635        block_on(sidecar.dispose_session("conn-1", "session-1", DisposeReason::Requested))
3636            .expect("dispose session");
3637
3638        assert_eq!(
3639            sidecar.take_disposed_sessions(),
3640            vec![(String::from("conn-1"), String::from("session-1"))],
3641            "dispose must publish the session scope so stdio can untrack it"
3642        );
3643    }
3644
3645    // H1 + M6: every per-VM tracking map is reclaimed for a disposed VM. The
3646    // output-buffer map (M6) was previously only removed on a successful handoff,
3647    // and the engine/extension maps (H1) were only reclaimed after the fallible
3648    // teardown steps' `?`, so any failure stranded them.
3649    #[test]
3650    fn reclaim_vm_tracking_clears_every_per_vm_map() {
3651        let mut sidecar = test_sidecar();
3652        insert_session(
3653            &mut sidecar,
3654            "conn-1",
3655            "session-1",
3656            BTreeSet::from([String::from("vm-1")]),
3657        );
3658        sidecar.extension_process_output_buffers.insert(
3659            (String::from("vm-1"), String::from("proc-1")),
3660            ExtensionBufferedProcessOutput::default(),
3661        );
3662        sidecar.extension_sessions.insert(
3663            (String::from("ns"), String::from("ext-sess-1")),
3664            ExtensionSessionResources {
3665                ownership: OwnershipScope::vm("conn-1", "session-1", "vm-1"),
3666                process_ids: BTreeSet::new(),
3667                vm_ids: BTreeSet::from([String::from("vm-1")]),
3668            },
3669        );
3670
3671        sidecar.reclaim_vm_tracking("session-1", "vm-1");
3672
3673        assert!(
3674            sidecar.extension_process_output_buffers.is_empty(),
3675            "M6: the output-buffer map must be reclaimed on VM disposal"
3676        );
3677        assert!(
3678            sidecar.extension_sessions.is_empty(),
3679            "H1: an extension session bound only to the VM must be reclaimed"
3680        );
3681        assert!(
3682            !sidecar
3683                .sessions
3684                .get("session-1")
3685                .expect("session present")
3686                .vm_ids
3687                .contains("vm-1"),
3688            "the VM id must be removed from its session"
3689        );
3690    }
3691
3692    // H1: a failing VM dispose inside the loop must not abandon the session. With
3693    // unregistered VM ids, `dispose_vm_internal` fails on `require_owned_vm`;
3694    // pre-fix the loop `?`-ed out and left the session in `self.sessions`.
3695    #[test]
3696    fn dispose_session_reclaims_session_even_when_a_vm_dispose_fails() {
3697        let mut sidecar = test_sidecar();
3698        insert_session(
3699            &mut sidecar,
3700            "conn-1",
3701            "session-1",
3702            BTreeSet::from([String::from("vm-a"), String::from("vm-b")]),
3703        );
3704
3705        let result =
3706            block_on(sidecar.dispose_session("conn-1", "session-1", DisposeReason::Requested));
3707
3708        assert!(
3709            result.is_err(),
3710            "a failing VM dispose must still surface an error"
3711        );
3712        assert!(
3713            !sidecar.sessions.contains_key("session-1"),
3714            "the session must be reclaimed even though VM dispose failed"
3715        );
3716    }
3717}