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