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};
55use 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
70const 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
104pub use crate::state::{DispatchResult, NativeSidecarConfig, SidecarError};
106
107#[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 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
791impl 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
843pub 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 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 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 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 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 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 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 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 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 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 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 pub(crate) fn take_disposed_sessions(&mut self) -> Vec<(String, String)> {
1925 std::mem::take(&mut self.disposed_sessions)
1926 }
1927
1928 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 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 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 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 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
3194pub(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
3234pub(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#[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#[allow(dead_code)]
3355fn symlinked_node_modules_hint(stderr: &str) -> Option<&'static str> {
3356 const STORE_MARKERS: &[&str] = &[
3361 "node_modules/.pnpm/",
3362 "node_modules/.bun/",
3363 "node_modules/.store/",
3364 "/__virtual__/",
3365 ];
3366 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 #[test]
3397 fn matches_pnpm_store_enoent() {
3398 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 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 #[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 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 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 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 #[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 #[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 #[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 #[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 #[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}