From 7f62357cd209cdd465b32f7203bbbc26e1ffbc3a Mon Sep 17 00:00:00 2001 From: ldm0 Date: Tue, 18 Aug 2026 17:49:10 +0800 Subject: [PATCH] fix(devtools): split Worker inspector execution paths --- .../src/runtime/browser_context_runtime.rs | 13 +- .../dedicated_workers.rs | 75 +-- .../src/shared_worker_runtime/host_worker.rs | 16 +- .../protocol_commands.rs | 16 +- moli-renderer-v8/src/worker/handle.rs | 148 ++++- .../src/worker/inspector_task_runner.rs | 618 ++++++++++++++++++ moli-renderer-v8/src/worker/mod.rs | 3 +- moli-renderer-v8/src/worker/thread/isolate.rs | 33 +- moli-renderer-v8/src/worker/thread/mod.rs | 243 +++---- .../src/worker/thread/runtime_inspector.rs | 363 ++++++++-- .../src/worker/thread/tests/lifecycle.rs | 143 +++- 11 files changed, 1339 insertions(+), 332 deletions(-) create mode 100644 moli-renderer-v8/src/worker/inspector_task_runner.rs diff --git a/moli-renderer-v8/src/runtime/browser_context_runtime.rs b/moli-renderer-v8/src/runtime/browser_context_runtime.rs index dfc5138ba5..d3d372992d 100644 --- a/moli-renderer-v8/src/runtime/browser_context_runtime.rs +++ b/moli-renderer-v8/src/runtime/browser_context_runtime.rs @@ -181,8 +181,7 @@ struct RendererBrowserContextRuntimeInner { next_child_document_loader_id: AtomicU64, next_detached_parser_script_fetch_id: AtomicU64, next_dedicated_worker_instance_id: AtomicU64, - dedicated_worker_devtools_senders: - Mutex>>, + dedicated_worker_devtools_handles: Mutex>, dedicated_worker_pause_on_start_for_devtools: AtomicBool, javascript_dialog_handler_enabled: AtomicBool, renderer_output_transport_tx: RendererOutputTransportSenderSlot, @@ -310,10 +309,10 @@ impl Drop for RendererBrowserContextRuntimeInner { } fn terminate_browser_context_resource_producers(inner: &RendererBrowserContextRuntimeInner) { - let dedicated_worker_senders = - std::mem::take(&mut *inner.dedicated_worker_devtools_senders.lock()); - for sender in dedicated_worker_senders.into_values() { - let _ = sender.send(crate::worker::WorkerMessage::Terminate); + let dedicated_worker_handles = + std::mem::take(&mut *inner.dedicated_worker_devtools_handles.lock()); + for handle in dedicated_worker_handles.into_values() { + let _ = handle.terminate_for_devtools(); } inner .shared_worker_runtime @@ -521,7 +520,7 @@ impl RendererBrowserContextRuntime { next_child_document_loader_id: AtomicU64::default(), next_detached_parser_script_fetch_id: AtomicU64::default(), next_dedicated_worker_instance_id: AtomicU64::default(), - dedicated_worker_devtools_senders: Mutex::new(HashMap::new()), + dedicated_worker_devtools_handles: Mutex::new(HashMap::new()), dedicated_worker_pause_on_start_for_devtools: AtomicBool::new(false), javascript_dialog_handler_enabled: AtomicBool::new(false), renderer_output_transport_tx, diff --git a/moli-renderer-v8/src/runtime/browser_context_runtime/dedicated_workers.rs b/moli-renderer-v8/src/runtime/browser_context_runtime/dedicated_workers.rs index dd64b85011..a4d0b5d166 100644 --- a/moli-renderer-v8/src/runtime/browser_context_runtime/dedicated_workers.rs +++ b/moli-renderer-v8/src/runtime/browser_context_runtime/dedicated_workers.rs @@ -1,7 +1,7 @@ -use tokio::sync::{mpsc, oneshot}; +use tokio::sync::oneshot; use crate::runtime::{RendererRuntimeInspectorMessage, RendererRuntimeInspectorResponseSender}; -use crate::worker::{WorkerHandle, WorkerMessage}; +use crate::worker::{WorkerDevToolsHandle, WorkerHandle}; use super::RendererBrowserContextRuntime; @@ -20,24 +20,16 @@ impl RendererBrowserContextRuntime { &self, instance_id: u64, handle: &WorkerHandle, - ) { - self.attach_dedicated_worker_devtools_sender(instance_id, handle.tx.clone()); - } - - fn attach_dedicated_worker_devtools_sender( - &self, - instance_id: u64, - sender: mpsc::UnboundedSender, ) { self.inner - .dedicated_worker_devtools_senders + .dedicated_worker_devtools_handles .lock() - .insert(instance_id, sender); + .insert(instance_id, handle.devtools_handle()); } pub(crate) fn unregister_dedicated_worker_devtools_handle(&self, instance_id: u64) { self.inner - .dedicated_worker_devtools_senders + .dedicated_worker_devtools_handles .lock() .remove(&instance_id); } @@ -54,12 +46,9 @@ impl RendererBrowserContextRuntime { .load(std::sync::atomic::Ordering::Relaxed) } - fn dedicated_worker_devtools_sender( - &self, - instance_id: u64, - ) -> Option> { + fn dedicated_worker_devtools_handle(&self, instance_id: u64) -> Option { self.inner - .dedicated_worker_devtools_senders + .dedicated_worker_devtools_handles .lock() .get(&instance_id) .cloned() @@ -103,18 +92,18 @@ impl RendererBrowserContextRuntime { raw_json: String, deferred_response: Option, ) -> Result, String> { - let Some(sender) = self.dedicated_worker_devtools_sender(instance_id) else { + let Some(handle) = self.dedicated_worker_devtools_handle(instance_id) else { return Err("DedicatedWorkerRuntimeUnavailable".to_owned()); }; let (response_tx, response_rx) = oneshot::channel(); - sender - .send(WorkerMessage::DispatchRuntimeProtocolMessage { - inspector_session_id, - raw_json, - deferred_response, - response_tx, - }) - .map_err(|_| "DedicatedWorkerRuntimeUnavailable".to_owned())?; + if !handle.dispatch_runtime_protocol_message( + inspector_session_id, + raw_json, + deferred_response, + response_tx, + ) { + return Err("DedicatedWorkerRuntimeUnavailable".to_owned()); + } response_rx .await .map_err(|_| "DedicatedWorkerRuntimeUnavailable".to_owned())? @@ -125,14 +114,8 @@ impl RendererBrowserContextRuntime { instance_id: u64, inspector_session_id: Option, ) -> bool { - self.dedicated_worker_devtools_sender(instance_id) - .is_some_and(|sender| { - sender - .send(WorkerMessage::AttachRuntimeInspectorSession { - inspector_session_id, - }) - .is_ok() - }) + self.dedicated_worker_devtools_handle(instance_id) + .is_some_and(|handle| handle.attach_runtime_inspector_session(inspector_session_id)) } pub fn detach_dedicated_worker_runtime_inspector_session( @@ -140,30 +123,20 @@ impl RendererBrowserContextRuntime { instance_id: u64, inspector_session_id: Option, ) -> bool { - self.dedicated_worker_devtools_sender(instance_id) - .is_some_and(|sender| { - sender - .send(WorkerMessage::DetachRuntimeInspectorSession { - inspector_session_id, - }) - .is_ok() - }) + self.dedicated_worker_devtools_handle(instance_id) + .is_some_and(|handle| handle.detach_runtime_inspector_session(inspector_session_id)) } pub fn run_dedicated_worker_if_waiting_for_debugger_for_devtools( &self, instance_id: u64, ) -> bool { - self.dedicated_worker_devtools_sender(instance_id) - .is_some_and(|sender| { - sender - .send(WorkerMessage::RunIfWaitingForDebuggerForDevtools) - .is_ok() - }) + self.dedicated_worker_devtools_handle(instance_id) + .is_some_and(|handle| handle.run_if_waiting_for_debugger()) } pub fn close_dedicated_worker_for_devtools(&self, instance_id: u64) -> bool { - self.dedicated_worker_devtools_sender(instance_id) - .is_some_and(|sender| sender.send(WorkerMessage::Terminate).is_ok()) + self.dedicated_worker_devtools_handle(instance_id) + .is_some_and(|handle| handle.terminate_for_devtools()) } } diff --git a/moli-renderer-v8/src/shared_worker_runtime/host_worker.rs b/moli-renderer-v8/src/shared_worker_runtime/host_worker.rs index b65960fc64..092fc6916d 100644 --- a/moli-renderer-v8/src/shared_worker_runtime/host_worker.rs +++ b/moli-renderer-v8/src/shared_worker_runtime/host_worker.rs @@ -3,7 +3,8 @@ use std::sync::Arc; use tokio::sync::mpsc; use crate::worker::{ - WorkerGlobalKind, WorkerMessage, WorkerSpawnOptions, spawn_worker_with_options, + WorkerDevToolsHandle, WorkerGlobalKind, WorkerMessage, WorkerSpawnOptions, + spawn_worker_with_options, }; use super::{ @@ -113,6 +114,19 @@ impl RendererSharedWorkerHost { } } + pub(super) fn running_devtools_handle(&self) -> Option { + let state = self.state.lock(); + match &*state { + RendererSharedWorkerHostState::Running { + handle: Some(handle), + .. + } => Some(handle.devtools_handle()), + RendererSharedWorkerHostState::Loading { .. } + | RendererSharedWorkerHostState::Running { handle: None, .. } + | RendererSharedWorkerHostState::Closed => None, + } + } + pub(super) fn send_worker_message(&self, message: WorkerMessage) -> bool { let Some(tx) = self.running_tx() else { return false; diff --git a/moli-renderer-v8/src/shared_worker_runtime/protocol_commands.rs b/moli-renderer-v8/src/shared_worker_runtime/protocol_commands.rs index bb12d57e39..c1c8a89ad3 100644 --- a/moli-renderer-v8/src/shared_worker_runtime/protocol_commands.rs +++ b/moli-renderer-v8/src/shared_worker_runtime/protocol_commands.rs @@ -1,9 +1,7 @@ use tokio::sync::oneshot; -use crate::runtime::{RendererRuntimeInspectorMessage, RendererRuntimeInspectorResponseSender}; -use crate::worker::WorkerMessage; - use super::host::RendererSharedWorkerHost; +use crate::runtime::{RendererRuntimeInspectorMessage, RendererRuntimeInspectorResponseSender}; impl RendererSharedWorkerHost { pub(super) async fn dispatch_runtime_protocol_message( @@ -39,13 +37,16 @@ impl RendererSharedWorkerHost { raw_json: String, deferred_response: Option, ) -> Result, String> { + let Some(handle) = self.running_devtools_handle() else { + return Err("SharedWorkerRuntimeUnavailable".to_owned()); + }; let (response_tx, response_rx) = oneshot::channel(); - if !self.send_worker_message(WorkerMessage::DispatchRuntimeProtocolMessage { + if !handle.dispatch_runtime_protocol_message( inspector_session_id, raw_json, deferred_response, response_tx, - }) { + ) { return Err("SharedWorkerRuntimeUnavailable".to_owned()); } response_rx @@ -57,8 +58,7 @@ impl RendererSharedWorkerHost { &self, inspector_session_id: Option, ) -> bool { - self.send_worker_message(WorkerMessage::DetachRuntimeInspectorSession { - inspector_session_id, - }) + self.running_devtools_handle() + .is_some_and(|handle| handle.detach_runtime_inspector_session(inspector_session_id)) } } diff --git a/moli-renderer-v8/src/worker/handle.rs b/moli-renderer-v8/src/worker/handle.rs index 6a8db5d4ea..4873b1048d 100644 --- a/moli-renderer-v8/src/worker/handle.rs +++ b/moli-renderer-v8/src/worker/handle.rs @@ -33,6 +33,7 @@ use crate::runtime::{ }; use crate::structured_clone::V8StructuredClonePayload; use crate::types::{BroadcastChannelId, DedicatedWorkerId, MessagePortId, NetworkBodySourceId}; +use crate::worker::inspector_task_runner::WorkerInspectorTaskRunner; use moli_crypto::sha256_hex; use moli_fetch::{RequestCredentialsMode, ResponseHead}; use moli_shared_worker::SharedWorkerInstanceId; @@ -143,23 +144,10 @@ pub(crate) enum WorkerMessage { worker_id: DedicatedWorkerId, message: Box, }, - /// Dispatch a CDP Runtime protocol message inside this worker's V8 inspector session. - DispatchRuntimeProtocolMessage { - inspector_session_id: Option, - raw_json: String, - deferred_response: Option, - response_tx: oneshot::Sender, String>>, - }, - /// Attach one renderer-side V8 inspector session before its first command. - AttachRuntimeInspectorSession { - inspector_session_id: Option, - }, - /// Release a pre-bootstrap debugger pause already acknowledged by the owner. - RunIfWaitingForDebuggerForDevtools, - /// Detach one renderer-side V8 inspector session from this worker. - DetachRuntimeInspectorSession { - inspector_session_id: Option, - }, + /// Owner-thread fallback for one queued interrupting Inspector task. + RunInterruptingInspectorTask, + /// Owner-thread dispatch for one Inspector task that may run JavaScript. + RunInspectorTaskDontInterrupt, #[cfg(test)] /// Inspect worker resource-owner V8 slots from inside the worker thread. ResourceOwnerSlotDiagnostics { @@ -691,6 +679,70 @@ impl WorkerRuntimeEvent { /// Handle held by the parent (main-frame) context to communicate with a /// running worker. +#[derive(Clone, Debug)] +pub(crate) struct WorkerDevToolsHandle { + worker_tx: mpsc::UnboundedSender, + inspector_tasks: WorkerInspectorTaskRunner, +} + +impl WorkerDevToolsHandle { + pub(crate) fn new( + wake_tx: mpsc::UnboundedSender, + isolate_handle: Arc>>, + ) -> Self { + Self { + inspector_tasks: WorkerInspectorTaskRunner::new(wake_tx.clone(), isolate_handle), + worker_tx: wake_tx, + } + } + + pub(crate) fn inspector_tasks(&self) -> &WorkerInspectorTaskRunner { + &self.inspector_tasks + } + + pub(crate) fn dispatch_runtime_protocol_message( + &self, + inspector_session_id: Option, + raw_json: String, + deferred_response: Option, + response_tx: oneshot::Sender, String>>, + ) -> bool { + self.inspector_tasks.append_protocol_message( + inspector_session_id, + raw_json, + deferred_response, + response_tx, + ) + } + + pub(crate) fn attach_runtime_inspector_session( + &self, + inspector_session_id: Option, + ) -> bool { + self.inspector_tasks.append_attach(inspector_session_id) + } + + pub(crate) fn detach_runtime_inspector_session( + &self, + inspector_session_id: Option, + ) -> bool { + self.inspector_tasks.append_detach(inspector_session_id) + } + + pub(crate) fn run_if_waiting_for_debugger(&self) -> bool { + self.inspector_tasks.append_run_if_waiting_for_debugger() + } + + pub(crate) fn dispose(&self, message: &str) { + self.inspector_tasks.dispose(message); + } + + pub(crate) fn terminate_for_devtools(&self) -> bool { + self.dispose("Worker closed before Inspector task dispatch"); + self.worker_tx.send(WorkerMessage::Terminate).is_ok() + } +} + pub(crate) struct WorkerHandle { /// Send messages *to* the worker. pub(crate) tx: mpsc::UnboundedSender, @@ -700,6 +752,7 @@ pub(crate) struct WorkerHandle { join_handle: Option>, isolate_handle: Arc>>, termination_requested: Arc, + devtools: WorkerDevToolsHandle, } impl WorkerHandle { @@ -719,12 +772,32 @@ impl WorkerHandle { ) } + #[cfg(test)] pub(crate) fn new_with_termination_requested( tx: mpsc::UnboundedSender, rx: mpsc::UnboundedReceiver, join_handle: std::thread::JoinHandle<()>, isolate_handle: Arc>>, termination_requested: Arc, + ) -> Self { + let devtools = WorkerDevToolsHandle::new(tx.clone(), Arc::clone(&isolate_handle)); + Self::new_with_termination_requested_and_devtools( + tx, + rx, + join_handle, + isolate_handle, + termination_requested, + devtools, + ) + } + + pub(crate) fn new_with_termination_requested_and_devtools( + tx: mpsc::UnboundedSender, + rx: mpsc::UnboundedReceiver, + join_handle: std::thread::JoinHandle<()>, + isolate_handle: Arc>>, + termination_requested: Arc, + devtools: WorkerDevToolsHandle, ) -> Self { Self { tx, @@ -732,6 +805,7 @@ impl WorkerHandle { join_handle: Some(join_handle), isolate_handle, termination_requested, + devtools, } } @@ -745,6 +819,8 @@ impl WorkerHandle { // Publish the lifecycle transition before interrupting V8. The worker // event loop can then reject an already-selected task without relying // on ordering between cloned mpsc senders. + self.devtools + .dispose("Worker terminated before Inspector task dispatch"); self.termination_requested.store(true, Ordering::Release); self.terminate_execution_if_ready(); let _ = self.tx.send(WorkerMessage::Terminate); @@ -776,31 +852,37 @@ impl WorkerHandle { deferred_response: Option, response_tx: oneshot::Sender, String>>, ) -> bool { - self.tx - .send(WorkerMessage::DispatchRuntimeProtocolMessage { - inspector_session_id, - raw_json, - deferred_response, - response_tx, - }) - .is_ok() + self.devtools.dispatch_runtime_protocol_message( + inspector_session_id, + raw_json, + deferred_response, + response_tx, + ) + } + + #[cfg(test)] + pub(crate) fn attach_runtime_inspector_session( + &self, + inspector_session_id: Option, + ) -> bool { + self.devtools + .attach_runtime_inspector_session(inspector_session_id) } pub(crate) fn detach_runtime_inspector_session( &self, inspector_session_id: Option, ) -> bool { - self.tx - .send(WorkerMessage::DetachRuntimeInspectorSession { - inspector_session_id, - }) - .is_ok() + self.devtools + .detach_runtime_inspector_session(inspector_session_id) } pub(crate) fn run_if_waiting_for_debugger_for_devtools(&self) -> bool { - self.tx - .send(WorkerMessage::RunIfWaitingForDebuggerForDevtools) - .is_ok() + self.devtools.run_if_waiting_for_debugger() + } + + pub(crate) fn devtools_handle(&self) -> WorkerDevToolsHandle { + self.devtools.clone() } pub(crate) fn set_extra_http_headers(&self, headers: &[(String, String)]) { diff --git a/moli-renderer-v8/src/worker/inspector_task_runner.rs b/moli-renderer-v8/src/worker/inspector_task_runner.rs new file mode 100644 index 0000000000..9a248cb217 --- /dev/null +++ b/moli-renderer-v8/src/worker/inspector_task_runner.rs @@ -0,0 +1,618 @@ +use std::{ + cell::RefCell, + collections::{HashMap, VecDeque}, + ffi::c_void, + rc::{Rc, Weak}, + sync::{ + Arc, + atomic::{AtomicBool, AtomicU64, Ordering}, + }, +}; + +use parking_lot::{Condvar, Mutex}; +use serde_json::json; +use tokio::sync::{mpsc, oneshot}; + +use crate::runtime::{RendererRuntimeInspectorMessage, RendererRuntimeInspectorResponseSender}; + +use super::handle::WorkerMessage; + +static NEXT_WORKER_INSPECTOR_ROUTE_ID: AtomicU64 = AtomicU64::new(1); + +thread_local! { + static WORKER_INSPECTOR_EXECUTORS: RefCell>> = + RefCell::new(HashMap::new()); +} + +pub(crate) trait WorkerInspectorInterruptExecutor { + fn dispatch_interrupt(&self, isolate: v8::UnsafeRawIsolatePtr); +} + +pub(crate) fn register_worker_inspector_executor( + route_id: u64, + executor: &Rc, +) { + let previous = WORKER_INSPECTOR_EXECUTORS.with(|executors| { + executors + .borrow_mut() + .insert(route_id, Rc::downgrade(executor)) + }); + assert!( + previous.is_none(), + "worker Inspector executor route IDs must be unique" + ); +} + +pub(crate) fn unregister_worker_inspector_executor(route_id: u64) { + let _ = WORKER_INSPECTOR_EXECUTORS.try_with(|executors| { + executors.borrow_mut().remove(&route_id); + }); +} + +fn worker_inspector_executor(route_id: u64) -> Option> { + WORKER_INSPECTOR_EXECUTORS + .try_with(|executors| executors.borrow().get(&route_id).and_then(Weak::upgrade)) + .ok() + .flatten() +} + +struct WorkerInspectorInterruptTarget { + route_id: u64, +} + +unsafe extern "C" fn dispatch_worker_inspector_interrupt( + isolate: v8::UnsafeRawIsolatePtr, + data: *mut c_void, +) { + // SAFETY: every accepted V8 interrupt owns exactly one strong reference + // created with `Arc::into_raw`. V8 invokes an accepted callback at most + // once, so this callback is the unique consumer of that reference. + let callback_target = unsafe { Arc::from_raw(data.cast::()) }; + let Some(executor) = worker_inspector_executor(callback_target.route_id) else { + return; + }; + executor.dispatch_interrupt(isolate); +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub(crate) enum WorkerInspectorTaskMode { + Interrupt, + DontInterrupt, +} + +pub(crate) fn worker_inspector_task_mode(method: &str) -> WorkerInspectorTaskMode { + match method { + "Debugger.evaluateOnCallFrame" + | "Runtime.evaluate" + | "Runtime.callFunctionOn" + | "Runtime.getProperties" + | "Runtime.runScript" => WorkerInspectorTaskMode::DontInterrupt, + _ => WorkerInspectorTaskMode::Interrupt, + } +} + +pub(crate) enum WorkerInspectorTask { + DispatchProtocolMessage { + inspector_session_id: Option, + raw_json: String, + deferred_response: Option, + response_tx: oneshot::Sender, String>>, + }, + AttachSession { + inspector_session_id: Option, + }, + DetachSession { + inspector_session_id: Option, + }, + RunIfWaitingForDebugger, +} + +impl WorkerInspectorTask { + fn fail(self, message: &str) { + let Self::DispatchProtocolMessage { + deferred_response, + response_tx, + .. + } = self + else { + return; + }; + if let Some(response) = deferred_response { + let call_id = response.call_id(); + let _ = response.send(json!({ + "id": call_id, + "error": { + "code": -32000, + "message": message, + }, + })); + } + let _ = response_tx.send(Err(message.to_owned())); + } +} + +struct WorkerInspectorTaskEntry { + sequence: u64, + task: WorkerInspectorTask, +} + +struct WorkerInspectorTaskRunnerState { + interrupting_tasks: VecDeque, + non_interrupting_tasks: VecDeque, + next_sequence: u64, + disposed: bool, + isolate_ready: bool, + pause_loop_active: bool, + quit_pause_loop: bool, +} + +struct WorkerInspectorTaskRunnerShared { + state: Mutex, + pause_work: Condvar, + wake_tx: mpsc::UnboundedSender, + isolate_handle: Arc>>, + interrupt_target: Arc, + interrupt_armed: AtomicBool, + resume_requested: AtomicBool, +} + +#[derive(Clone)] +pub(crate) struct WorkerInspectorTaskRunner { + shared: Arc, +} + +impl WorkerInspectorTaskRunner { + pub(crate) fn new( + wake_tx: mpsc::UnboundedSender, + isolate_handle: Arc>>, + ) -> Self { + let route_id = NEXT_WORKER_INSPECTOR_ROUTE_ID + .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| { + current.checked_add(1) + }) + .expect("worker Inspector route ID exhausted"); + Self { + shared: Arc::new(WorkerInspectorTaskRunnerShared { + state: Mutex::new(WorkerInspectorTaskRunnerState { + interrupting_tasks: VecDeque::new(), + non_interrupting_tasks: VecDeque::new(), + next_sequence: 1, + disposed: false, + isolate_ready: false, + pause_loop_active: false, + quit_pause_loop: false, + }), + pause_work: Condvar::new(), + wake_tx, + isolate_handle, + interrupt_target: Arc::new(WorkerInspectorInterruptTarget { route_id }), + interrupt_armed: AtomicBool::new(false), + resume_requested: AtomicBool::new(false), + }), + } + } + + pub(crate) fn route_id(&self) -> u64 { + self.shared.interrupt_target.route_id + } + + pub(crate) fn append_protocol_message( + &self, + inspector_session_id: Option, + raw_json: String, + deferred_response: Option, + response_tx: oneshot::Sender, String>>, + ) -> bool { + let mode = serde_json::from_str::(&raw_json) + .ok() + .and_then(|value| { + value + .get("method") + .and_then(serde_json::Value::as_str) + .map(worker_inspector_task_mode) + }) + .unwrap_or(WorkerInspectorTaskMode::Interrupt); + self.append( + mode, + WorkerInspectorTask::DispatchProtocolMessage { + inspector_session_id, + raw_json, + deferred_response, + response_tx, + }, + ) + } + + pub(crate) fn append_attach(&self, inspector_session_id: Option) -> bool { + self.append( + WorkerInspectorTaskMode::Interrupt, + WorkerInspectorTask::AttachSession { + inspector_session_id, + }, + ) + } + + pub(crate) fn append_detach(&self, inspector_session_id: Option) -> bool { + self.append( + WorkerInspectorTaskMode::Interrupt, + WorkerInspectorTask::DetachSession { + inspector_session_id, + }, + ) + } + + pub(crate) fn append_run_if_waiting_for_debugger(&self) -> bool { + self.append( + WorkerInspectorTaskMode::Interrupt, + WorkerInspectorTask::RunIfWaitingForDebugger, + ) + } + + fn append(&self, mode: WorkerInspectorTaskMode, task: WorkerInspectorTask) -> bool { + let mut state = self.shared.state.lock(); + if state.disposed { + drop(state); + task.fail("Worker Inspector task runner is disposed"); + return false; + } + let sequence = state.next_sequence; + state.next_sequence = state + .next_sequence + .checked_add(1) + .expect("worker Inspector task sequence exhausted"); + let entry = WorkerInspectorTaskEntry { sequence, task }; + match mode { + WorkerInspectorTaskMode::Interrupt => state.interrupting_tasks.push_back(entry), + WorkerInspectorTaskMode::DontInterrupt => { + state.non_interrupting_tasks.push_back(entry); + } + } + drop(state); + + let wake = match mode { + WorkerInspectorTaskMode::Interrupt => WorkerMessage::RunInterruptingInspectorTask, + WorkerInspectorTaskMode::DontInterrupt => WorkerMessage::RunInspectorTaskDontInterrupt, + }; + if self.shared.wake_tx.send(wake).is_err() { + self.dispose("Worker Inspector owner is unavailable"); + return false; + } + if mode == WorkerInspectorTaskMode::Interrupt { + self.request_interrupt_if_needed(); + } + self.shared.pause_work.notify_one(); + true + } + + pub(crate) fn activate_isolate(&self) { + let mut state = self.shared.state.lock(); + if state.disposed { + return; + } + state.isolate_ready = true; + drop(state); + self.request_interrupt_if_needed(); + } + + pub(crate) fn claim_interrupting_task(&self) -> Option { + self.shared + .state + .lock() + .interrupting_tasks + .pop_front() + .map(|entry| entry.task) + } + + pub(crate) fn claim_non_interrupting_task(&self) -> Option { + self.shared + .state + .lock() + .non_interrupting_tasks + .pop_front() + .map(|entry| entry.task) + } + + pub(crate) fn interrupt_callback_started(&self) { + self.shared.interrupt_armed.store(false, Ordering::Release); + } + + pub(crate) fn request_interrupt_if_needed(&self) { + let should_request = { + let state = self.shared.state.lock(); + !state.disposed && state.isolate_ready && !state.interrupting_tasks.is_empty() + }; + if !should_request + || self + .shared + .interrupt_armed + .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire) + .is_err() + { + return; + } + + let callback_target = Arc::into_raw(Arc::clone(&self.shared.interrupt_target)); + let callback_data = callback_target.cast_mut().cast::(); + let accepted = self + .shared + .isolate_handle + .lock() + .as_ref() + .is_some_and(|handle| { + handle.request_interrupt(dispatch_worker_inspector_interrupt, callback_data) + }); + if !accepted { + // SAFETY: V8 rejected the request, so no callback can consume the + // strong reference created immediately above. + unsafe { drop(Arc::from_raw(callback_target)) }; + self.shared.interrupt_armed.store(false, Ordering::Release); + } + } + + pub(crate) fn begin_pause_loop(&self) -> bool { + let mut state = self.shared.state.lock(); + if state.disposed || state.pause_loop_active { + return false; + } + state.pause_loop_active = true; + state.quit_pause_loop = false; + true + } + + pub(crate) fn wait_for_pause_task(&self) -> Option { + let mut state = self.shared.state.lock(); + loop { + if state.disposed || state.quit_pause_loop { + return None; + } + let interrupt_sequence = state.interrupting_tasks.front().map(|entry| entry.sequence); + let non_interrupt_sequence = state + .non_interrupting_tasks + .front() + .map(|entry| entry.sequence); + match (interrupt_sequence, non_interrupt_sequence) { + (Some(interrupt), Some(non_interrupt)) if interrupt <= non_interrupt => { + let task = state + .interrupting_tasks + .pop_front() + .expect("front worker Inspector interrupt task disappeared") + .task; + return Some(task); + } + (Some(_), Some(_)) | (None, Some(_)) => { + let task = state + .non_interrupting_tasks + .pop_front() + .expect("front worker Inspector non-interrupt task disappeared") + .task; + return Some(task); + } + (Some(_), None) => { + let task = state + .interrupting_tasks + .pop_front() + .expect("front worker Inspector interrupt task disappeared") + .task; + return Some(task); + } + (None, None) => self.shared.pause_work.wait(&mut state), + } + } + } + + pub(crate) fn finish_pause_loop(&self) { + let mut state = self.shared.state.lock(); + state.pause_loop_active = false; + state.quit_pause_loop = false; + } + + pub(crate) fn request_quit_pause_loop(&self) { + self.shared.state.lock().quit_pause_loop = true; + self.shared.pause_work.notify_all(); + } + + pub(crate) fn request_resume(&self) { + self.shared.resume_requested.store(true, Ordering::Release); + } + + pub(crate) fn take_resume_requested(&self) -> bool { + self.shared.resume_requested.swap(false, Ordering::AcqRel) + } + + pub(crate) fn dispose(&self, message: &str) { + let tasks = { + let mut state = self.shared.state.lock(); + if state.disposed { + return; + } + state.disposed = true; + state.isolate_ready = false; + state.quit_pause_loop = true; + let mut tasks = state + .interrupting_tasks + .drain(..) + .map(|entry| entry.task) + .collect::>(); + tasks.extend( + state + .non_interrupting_tasks + .drain(..) + .map(|entry| entry.task), + ); + tasks + }; + self.shared.pause_work.notify_all(); + for task in tasks { + task.fail(message); + } + } +} + +impl std::fmt::Debug for WorkerInspectorTaskRunner { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + let state = self.shared.state.lock(); + formatter + .debug_struct("WorkerInspectorTaskRunner") + .field("route_id", &self.route_id()) + .field("interrupting_tasks", &state.interrupting_tasks.len()) + .field( + "non_interrupting_tasks", + &state.non_interrupting_tasks.len(), + ) + .field("disposed", &state.disposed) + .field("isolate_ready", &state.isolate_ready) + .finish() + } +} + +#[cfg(test)] +mod tests { + use std::sync::Arc; + + use parking_lot::Mutex; + use tokio::sync::{mpsc, oneshot}; + + use super::{ + WorkerInspectorTask, WorkerInspectorTaskMode, WorkerInspectorTaskRunner, + worker_inspector_task_mode, + }; + + fn runner() -> ( + WorkerInspectorTaskRunner, + mpsc::UnboundedReceiver, + ) { + let (wake_tx, wake_rx) = mpsc::unbounded_channel(); + ( + WorkerInspectorTaskRunner::new(wake_tx, Arc::new(Mutex::new(None))), + wake_rx, + ) + } + + fn append_protocol( + runner: &WorkerInspectorTaskRunner, + id: i64, + method: &str, + ) -> oneshot::Receiver, String>> + { + let (response_tx, response_rx) = oneshot::channel(); + assert!(runner.append_protocol_message( + Some("SID-worker-runner".to_owned()), + serde_json::json!({ "id": id, "method": method }).to_string(), + None, + response_tx, + )); + response_rx + } + + #[test] + fn chromium_worker_inspector_dont_interrupt_catalog_is_exact() { + for method in [ + "Debugger.evaluateOnCallFrame", + "Runtime.evaluate", + "Runtime.callFunctionOn", + "Runtime.getProperties", + "Runtime.runScript", + ] { + assert_eq!( + worker_inspector_task_mode(method), + WorkerInspectorTaskMode::DontInterrupt, + "{method} should use Chromium's non-interrupting worker path" + ); + } + for method in [ + "Debugger.pause", + "Debugger.resume", + "Debugger.stepInto", + "Runtime.enable", + "Runtime.terminateExecution", + "Inspector.disable", + ] { + assert_eq!( + worker_inspector_task_mode(method), + WorkerInspectorTaskMode::Interrupt, + "{method} should interrupt a busy worker" + ); + } + } + + #[test] + fn interrupt_task_overtakes_earlier_non_interrupting_owner_task() { + let (runner, mut wake_rx) = runner(); + let _evaluate = append_protocol(&runner, 1, "Runtime.evaluate"); + let _terminate = append_protocol(&runner, 2, "Runtime.terminateExecution"); + + assert!(matches!( + wake_rx.try_recv(), + Ok(super::WorkerMessage::RunInspectorTaskDontInterrupt) + )); + assert!(matches!( + wake_rx.try_recv(), + Ok(super::WorkerMessage::RunInterruptingInspectorTask) + )); + assert!(matches!( + runner.claim_interrupting_task(), + Some(WorkerInspectorTask::DispatchProtocolMessage { raw_json, .. }) + if raw_json.contains("Runtime.terminateExecution") + )); + assert!(matches!( + runner.claim_non_interrupting_task(), + Some(WorkerInspectorTask::DispatchProtocolMessage { raw_json, .. }) + if raw_json.contains("Runtime.evaluate") + )); + assert!(runner.claim_interrupting_task().is_none()); + assert!(runner.claim_non_interrupting_task().is_none()); + } + + #[test] + fn owner_fallback_and_interrupt_claim_one_fifo_exactly_once() { + let (runner, _wake_rx) = runner(); + let _first = append_protocol(&runner, 1, "Debugger.pause"); + let _second = append_protocol(&runner, 2, "Debugger.resume"); + + let first = runner + .claim_interrupting_task() + .expect("owner fallback should claim first task"); + let second = runner + .claim_interrupting_task() + .expect("interrupt callback should claim second task"); + assert!(runner.claim_interrupting_task().is_none()); + assert!(matches!( + first, + WorkerInspectorTask::DispatchProtocolMessage { raw_json, .. } + if raw_json.contains("Debugger.pause") + )); + assert!(matches!( + second, + WorkerInspectorTask::DispatchProtocolMessage { raw_json, .. } + if raw_json.contains("Debugger.resume") + )); + } + + #[test] + fn dispose_rejects_queued_and_late_protocol_tasks() { + let (runner, _wake_rx) = runner(); + let mut queued = append_protocol(&runner, 1, "Runtime.evaluate"); + runner.dispose("worker teardown"); + assert_eq!( + queued + .try_recv() + .expect("queued task should be canceled") + .expect_err("queued task must not dispatch"), + "worker teardown" + ); + + let (response_tx, mut late) = oneshot::channel(); + assert!(!runner.append_protocol_message( + None, + serde_json::json!({ "id": 2, "method": "Debugger.pause" }).to_string(), + None, + response_tx, + )); + assert_eq!( + late.try_recv() + .expect("late task should be rejected") + .expect_err("late task must not dispatch"), + "Worker Inspector task runner is disposed" + ); + } +} diff --git a/moli-renderer-v8/src/worker/mod.rs b/moli-renderer-v8/src/worker/mod.rs index a97fe651e4..cf707a769a 100644 --- a/moli-renderer-v8/src/worker/mod.rs +++ b/moli-renderer-v8/src/worker/mod.rs @@ -14,6 +14,7 @@ pub(crate) mod abort; mod data_url; mod global_scope; mod handle; +mod inspector_task_runner; mod module_mime; mod module_runtime; mod script_loading; @@ -56,7 +57,7 @@ pub(crate) use handle::{ WorkerScriptResource, WorkerScriptResourceKind, WorkerToParentMessage, WorkerWebSocketFrameEvent, WorkerWebSocketLifecycleEvent, worker_secure_context_for_script_url, }; -pub(crate) use handle::{WorkerHandle, WorkerNetworkPolicy}; +pub(crate) use handle::{WorkerDevToolsHandle, WorkerHandle, WorkerNetworkPolicy}; pub(crate) use module_mime::{ ensure_worker_css_module_mime, ensure_worker_json_module_mime, ensure_worker_text_module_mime, ensure_worker_wasm_module_mime, diff --git a/moli-renderer-v8/src/worker/thread/isolate.rs b/moli-renderer-v8/src/worker/thread/isolate.rs index f0d518d1dc..80684b9ddc 100644 --- a/moli-renderer-v8/src/worker/thread/isolate.rs +++ b/moli-renderer-v8/src/worker/thread/isolate.rs @@ -1,4 +1,11 @@ +use std::rc::Rc; + +use tokio::sync::mpsc; + use crate::v8_platform::{V8ForegroundTaskWake, V8PlatformIsolateRegistration}; +use crate::worker::{ + handle::WorkerToParentMessage, inspector_task_runner::WorkerInspectorTaskRunner, +}; use super::{ super::module_runtime::{ @@ -12,13 +19,18 @@ pub(super) struct WorkerIsolateState { // Inspector teardown touches V8-owned state, and platform registration must // be gone before the isolate is destroyed. Keep both fields before // `isolate` so Rust drops them first on normal worker teardown. - runtime_inspector: WorkerRuntimeInspector, + runtime_inspector: Rc, platform_registration: V8PlatformIsolateRegistration, isolate: v8::OwnedIsolate, } impl WorkerIsolateState { - pub(super) fn new(platform_wake: V8ForegroundTaskWake) -> Self { + pub(super) fn new( + platform_wake: V8ForegroundTaskWake, + inspector_task_runner: WorkerInspectorTaskRunner, + parent_tx: mpsc::UnboundedSender, + shared_worker: bool, + ) -> Self { let mut isolate = v8::Isolate::new(Default::default()); crate::context_bootstrap::install_agent_microtask_checkpoint_tasks(&mut isolate); isolate.set_microtasks_policy(v8::MicrotasksPolicy::Explicit); @@ -33,7 +45,12 @@ impl WorkerIsolateState { isolate.set_modify_code_generation_from_strings_callback( crate::context_bootstrap::trusted_types_code_generation_check_callback, ); - let runtime_inspector = WorkerRuntimeInspector::new(&mut isolate); + let runtime_inspector = WorkerRuntimeInspector::new( + &mut isolate, + inspector_task_runner, + parent_tx, + shared_worker, + ); let platform_registration = V8PlatformIsolateRegistration::register( &mut isolate, platform_wake.into_platform_wake(), @@ -50,14 +67,14 @@ impl WorkerIsolateState { &mut self.isolate } - pub(super) fn worker_runtime_inspector_mut(&mut self) -> &mut WorkerRuntimeInspector { - &mut self.runtime_inspector + pub(super) fn worker_runtime_inspector(&self) -> &WorkerRuntimeInspector { + &self.runtime_inspector } - pub(super) fn worker_isolate_and_runtime_inspector_mut( + pub(super) fn worker_isolate_and_runtime_inspector( &mut self, - ) -> (&mut v8::OwnedIsolate, &mut WorkerRuntimeInspector) { - (&mut self.isolate, &mut self.runtime_inspector) + ) -> (&mut v8::OwnedIsolate, Rc) { + (&mut self.isolate, Rc::clone(&self.runtime_inspector)) } pub(super) fn unregister_worker_isolate_platform(&mut self) { diff --git a/moli-renderer-v8/src/worker/thread/mod.rs b/moli-renderer-v8/src/worker/thread/mod.rs index 4f25bd2f24..79868e8152 100644 --- a/moli-renderer-v8/src/worker/thread/mod.rs +++ b/moli-renderer-v8/src/worker/thread/mod.rs @@ -30,7 +30,7 @@ use crate::network::{ context::{WorkerResourceLoader, WorkerResourceOwner}, loads::{ResourceLoadDisposition, ResourceLoadKind}, }; -use crate::runtime::{RendererRuntimeInspectorMessage, RendererWorkerContextRuntime}; +use crate::runtime::RendererWorkerContextRuntime; use crate::service_worker_runtime::{ ServiceWorkerClientId, ServiceWorkerClientType, ServiceWorkerRuntimeService, }; @@ -92,11 +92,13 @@ use super::global_scope::{ service_worker_fetch_handler_type, }; use super::handle::{ - WorkerBootstrapCompletion, WorkerBootstrapFailure, WorkerBootstrapSuccess, WorkerErrorPhase, - WorkerErrorSource, WorkerFetchHandlerType, WorkerHandle, WorkerMessage, WorkerNetworkPolicy, - WorkerParentErrorEventKind, WorkerRuntimeInspectorMessageBatch, WorkerScriptResource, - WorkerScriptResourceKind, WorkerToParentMessage, + WorkerBootstrapCompletion, WorkerBootstrapFailure, WorkerBootstrapSuccess, + WorkerDevToolsHandle, WorkerErrorPhase, WorkerErrorSource, WorkerFetchHandlerType, + WorkerHandle, WorkerMessage, WorkerNetworkPolicy, WorkerParentErrorEventKind, + WorkerRuntimeInspectorMessageBatch, WorkerScriptResource, WorkerScriptResourceKind, + WorkerToParentMessage, }; +use super::inspector_task_runner::{WorkerInspectorTask, WorkerInspectorTaskRunner}; use super::module_runtime::{ WorkerBootstrapError, WorkerDynamicModuleImportAdvance, WorkerModuleBootstrapResume, WorkerModuleBootstrapStart, WorkerModuleEvaluationCompletion, WorkerModuleFetchedSource, @@ -975,36 +977,21 @@ fn drain_worker_dynamic_module_imports_for_context( drain_worker_dynamic_module_imports(scope, state, module_graph_fetch_tx); } -fn dispatch_worker_runtime_protocol_message( - isolate: &mut v8::OwnedIsolate, +fn dispatch_worker_inspector_task( + worker_isolate: &mut WorkerIsolateState, context: &v8::Global, - inspector: &mut WorkerRuntimeInspector, - inspector_session_id: Option<&str>, - raw_json: &str, - mut deferred_response: Option, + task: WorkerInspectorTask, state: &Rc>, module_graph_fetch_tx: &mpsc::UnboundedSender, - parent_tx: &mpsc::UnboundedSender, -) -> Result, String> { - if matches!(state.borrow().global_kind, WorkerGlobalKind::Shared { .. }) { - deferred_response = deferred_response - .map(|response| response.defer_publication_to_shared_worker_parent(parent_tx.clone())); - } +) { + let (isolate, inspector) = worker_isolate.worker_isolate_and_runtime_inspector(); + inspector.execute_task(isolate, task); let scope = pin!(v8::HandleScope::new(isolate)); let scope = &mut scope.init(); let ctx = v8::Local::new(scope, context); let scope = &mut v8::ContextScope::new(scope, ctx); - let messages = match deferred_response { - Some(callback) => inspector.dispatch_protocol_message_with_deferred_response( - scope, - inspector_session_id, - raw_json, - callback, - )?, - None => inspector.dispatch_protocol_message(scope, inspector_session_id, raw_json)?, - }; + perform_worker_microtask_checkpoint_and_report_pending_promise_rejections(scope); drain_worker_dynamic_module_imports(scope, state, module_graph_fetch_tx); - Ok(messages) } async fn run_worker_pre_bootstrap_debugger_pause( @@ -1012,64 +999,44 @@ async fn run_worker_pre_bootstrap_debugger_pause( context: &v8::Global, state: &Rc>, module_graph_fetch_tx: &mpsc::UnboundedSender, - parent_tx: &mpsc::UnboundedSender, script_url: &str, rx: &mut mpsc::UnboundedReceiver, pending_bootstrap_messages: &mut VecDeque, + inspector_task_runner: &WorkerInspectorTaskRunner, ) -> bool { loop { + if inspector_task_runner.take_resume_requested() { + return true; + } match rx.recv().await { - Some(WorkerMessage::DispatchRuntimeProtocolMessage { - inspector_session_id, - raw_json, - deferred_response, - response_tx, - }) => { - let should_resume = - worker_runtime_protocol_message_is_run_if_waiting_for_debugger(&raw_json); - let (isolate, runtime_inspector) = - worker_isolate.worker_isolate_and_runtime_inspector_mut(); - let result = dispatch_worker_runtime_protocol_message( - isolate, - context, - runtime_inspector, - inspector_session_id.as_deref(), - &raw_json, - deferred_response, - state, - module_graph_fetch_tx, - parent_tx, - ); - let did_resume = should_resume && result.is_ok(); - let _ = response_tx.send(result); - forward_pending_worker_runtime_protocol_messages( - worker_isolate.worker_runtime_inspector_mut(), - parent_tx, - ); - if did_resume { + Some(WorkerMessage::RunInterruptingInspectorTask) => { + if let Some(task) = inspector_task_runner.claim_interrupting_task() { + dispatch_worker_inspector_task( + worker_isolate, + context, + task, + state, + module_graph_fetch_tx, + ); + } + inspector_task_runner.request_interrupt_if_needed(); + if inspector_task_runner.take_resume_requested() { return true; } } - Some(WorkerMessage::AttachRuntimeInspectorSession { - inspector_session_id, - }) => { - worker_isolate - .worker_runtime_inspector_mut() - .attach_session(inspector_session_id.as_deref()); - } - Some(WorkerMessage::RunIfWaitingForDebuggerForDevtools) => { - return true; - } - Some(WorkerMessage::DetachRuntimeInspectorSession { - inspector_session_id, - }) => { - worker_isolate - .worker_runtime_inspector_mut() - .detach_session(inspector_session_id.as_deref()); - forward_pending_worker_runtime_protocol_messages( - worker_isolate.worker_runtime_inspector_mut(), - parent_tx, - ); + Some(WorkerMessage::RunInspectorTaskDontInterrupt) => { + if let Some(task) = inspector_task_runner.claim_non_interrupting_task() { + dispatch_worker_inspector_task( + worker_isolate, + context, + task, + state, + module_graph_fetch_tx, + ); + } + if inspector_task_runner.take_resume_requested() { + return true; + } } Some(WorkerMessage::SetExtraHttpHeaders(headers)) => { let (loader, headers_for_loader) = { @@ -1122,17 +1089,8 @@ async fn run_worker_pre_bootstrap_debugger_pause( } } -fn worker_runtime_protocol_message_is_run_if_waiting_for_debugger(raw_json: &str) -> bool { - serde_json::from_str::(raw_json).is_ok_and(|value| { - value - .get("method") - .and_then(serde_json::Value::as_str) - .is_some_and(|method| method == "Runtime.runIfWaitingForDebugger") - }) -} - fn drain_worker_runtime_protocol_messages( - inspector: &mut WorkerRuntimeInspector, + inspector: &WorkerRuntimeInspector, ) -> Vec { inspector.take_pending_messages() } @@ -1176,7 +1134,7 @@ fn worker_resource_owner_slot_diagnostics( } fn forward_pending_worker_runtime_protocol_messages( - inspector: &mut WorkerRuntimeInspector, + inspector: &WorkerRuntimeInspector, parent_tx: &mpsc::UnboundedSender, ) { let messages = drain_worker_runtime_protocol_messages(inspector); @@ -1186,7 +1144,7 @@ fn forward_pending_worker_runtime_protocol_messages( } fn forward_worker_script_loaded( - inspector: &mut WorkerRuntimeInspector, + inspector: &WorkerRuntimeInspector, parent_tx: &mpsc::UnboundedSender, ) { // Blink notifies each attached worker Inspector agent only after top-level @@ -1410,6 +1368,9 @@ pub(crate) fn spawn_worker_with_options(options: WorkerSpawnOptions) -> WorkerHa let worker_wake_tx = parent_to_worker_tx.clone(); let isolate_handle = Arc::new(Mutex::new(None)); let worker_isolate_handle = Arc::clone(&isolate_handle); + let devtools = + WorkerDevToolsHandle::new(parent_to_worker_tx.clone(), Arc::clone(&isolate_handle)); + let worker_inspector_tasks = devtools.inspector_tasks().clone(); let termination_requested = Arc::new(AtomicBool::new(false)); let worker_termination_requested = Arc::clone(&termination_requested); @@ -1454,16 +1415,18 @@ pub(crate) fn spawn_worker_with_options(options: WorkerSpawnOptions) -> WorkerHa worker_to_parent_tx, worker_isolate_handle, worker_termination_requested, + worker_inspector_tasks, )); }) .expect("failed to spawn worker thread"); - WorkerHandle::new_with_termination_requested( + WorkerHandle::new_with_termination_requested_and_devtools( parent_to_worker_tx, worker_to_parent_rx, join_handle, isolate_handle, termination_requested, + devtools, ) } @@ -1583,6 +1546,7 @@ async fn worker_main( parent_tx: mpsc::UnboundedSender, isolate_handle: Arc>>, termination_requested: Arc, + inspector_task_runner: WorkerInspectorTaskRunner, ) { debug!(url = %script_url, "worker started"); let mut bootstrap_completion = WorkerBootstrapCompletionReporter::new(bootstrap_completion_tx); @@ -1613,6 +1577,9 @@ async fn worker_main( let resource_owner_id = crate::resource_owner::ResourceOwnerId::new(); let mut worker_isolate = WorkerIsolateState::new( crate::v8_platform::V8ForegroundTaskWake::worker(worker_runtime_wake_tx.clone()), + inspector_task_runner.clone(), + parent_tx.clone(), + matches!(global_kind, WorkerGlobalKind::Shared { .. }), ); install_worker_promise_rejection_dispatch( worker_isolate.worker_isolate_mut(), @@ -1783,8 +1750,7 @@ async fn worker_main( let mut bootstrap_failed = false; let mut install_global_failed = false; { - let (isolate, runtime_inspector) = - worker_isolate.worker_isolate_and_runtime_inspector_mut(); + let (isolate, runtime_inspector) = worker_isolate.worker_isolate_and_runtime_inspector(); let scope = pin!(v8::HandleScope::new(isolate)); let scope = &mut scope.init(); *isolate_handle.lock() = Some(scope.thread_safe_handle()); @@ -1821,9 +1787,14 @@ async fn worker_main( } } if install_global_failed { + inspector_task_runner.dispose("Worker global installation failed"); + *isolate_handle.lock() = None; worker_isolate.unregister_worker_isolate_platform(); return; } + // Commands can queue before the isolate exists. Arm V8 interrupts only + // after the worker context and its Inspector routing are fully installed. + inspector_task_runner.activate_isolate(); let mut terminated_before_bootstrap = false; if pause_evaluation_until_debugger @@ -1832,10 +1803,10 @@ async fn worker_main( &context, &state, &module_graph_fetch_tx, - &parent_tx, &script_url, &mut rx, &mut pending_bootstrap_messages, + &inspector_task_runner, ) .await { @@ -1844,7 +1815,7 @@ async fn worker_main( } if !terminated_before_bootstrap { - let (isolate, _) = worker_isolate.worker_isolate_and_runtime_inspector_mut(); + let (isolate, _) = worker_isolate.worker_isolate_and_runtime_inspector(); let scope = pin!(v8::HandleScope::new(isolate)); let scope = &mut scope.init(); let ctx = v8::Local::new(scope, &context); @@ -1909,7 +1880,7 @@ async fn worker_main( drain_worker_dynamic_module_imports(scope, &state, &module_graph_fetch_tx); } if !terminated_before_bootstrap && pending_module_bootstrap.is_none() { - forward_worker_script_loaded(worker_isolate.worker_runtime_inspector_mut(), &parent_tx); + forward_worker_script_loaded(worker_isolate.worker_runtime_inspector(), &parent_tx); } if bootstrap_failed { state.borrow_mut().closed = true; @@ -2704,60 +2675,28 @@ async fn worker_main( drain_worker_dynamic_module_imports(scope, &state, &module_graph_fetch_tx); } } - WorkerLoopWake::Message(Some(WorkerMessage::DispatchRuntimeProtocolMessage { - inspector_session_id, - raw_json, - deferred_response, - response_tx, - })) => { - if pending_module_bootstrap.is_some() { - pending_bootstrap_messages.push_back( - WorkerMessage::DispatchRuntimeProtocolMessage { - inspector_session_id, - raw_json, - deferred_response, - response_tx, - }, + WorkerLoopWake::Message(Some(WorkerMessage::RunInterruptingInspectorTask)) => { + if let Some(task) = inspector_task_runner.claim_interrupting_task() { + dispatch_worker_inspector_task( + &mut worker_isolate, + &context, + task, + &state, + &module_graph_fetch_tx, ); - continue; } - let (isolate, runtime_inspector) = - worker_isolate.worker_isolate_and_runtime_inspector_mut(); - let result = dispatch_worker_runtime_protocol_message( - isolate, - &context, - runtime_inspector, - inspector_session_id.as_deref(), - &raw_json, - deferred_response, - &state, - &module_graph_fetch_tx, - &parent_tx, - ); - let _ = response_tx.send(result); + inspector_task_runner.request_interrupt_if_needed(); } - WorkerLoopWake::Message(Some(WorkerMessage::AttachRuntimeInspectorSession { - inspector_session_id, - })) => { - worker_isolate - .worker_runtime_inspector_mut() - .attach_session(inspector_session_id.as_deref()); - } - WorkerLoopWake::Message(Some(WorkerMessage::RunIfWaitingForDebuggerForDevtools)) => {} - WorkerLoopWake::Message(Some(WorkerMessage::DetachRuntimeInspectorSession { - inspector_session_id, - })) => { - if pending_module_bootstrap.is_some() { - pending_bootstrap_messages.push_back( - WorkerMessage::DetachRuntimeInspectorSession { - inspector_session_id, - }, + WorkerLoopWake::Message(Some(WorkerMessage::RunInspectorTaskDontInterrupt)) => { + if let Some(task) = inspector_task_runner.claim_non_interrupting_task() { + dispatch_worker_inspector_task( + &mut worker_isolate, + &context, + task, + &state, + &module_graph_fetch_tx, ); - continue; } - worker_isolate - .worker_runtime_inspector_mut() - .detach_session(inspector_session_id.as_deref()); } #[cfg(test)] WorkerLoopWake::Message(Some(WorkerMessage::ResourceOwnerSlotDiagnostics { @@ -3097,7 +3036,7 @@ async fn worker_main( } WorkerLoopWake::ModuleGraphFetch(Some(completion)) => { let (isolate, runtime_inspector) = - worker_isolate.worker_isolate_and_runtime_inspector_mut(); + worker_isolate.worker_isolate_and_runtime_inspector(); let scope = pin!(v8::HandleScope::new(isolate)); let scope = &mut scope.init(); let ctx = v8::Local::new(scope, &context); @@ -3159,7 +3098,7 @@ async fn worker_main( &state, &module_graph_fetch_tx, ); - forward_worker_script_loaded(runtime_inspector, &parent_tx); + forward_worker_script_loaded(&runtime_inspector, &parent_tx); } WorkerModuleBootstrapResume::NeedFetches(requests) => { start_worker_module_graph_fetch_batch( @@ -3192,7 +3131,7 @@ async fn worker_main( &parent_tx, &script_url, ); - forward_worker_script_loaded(runtime_inspector, &parent_tx); + forward_worker_script_loaded(&runtime_inspector, &parent_tx); break; } } @@ -3220,7 +3159,7 @@ async fn worker_main( } WorkerLoopWake::ModuleEvaluation(Some(completion)) => { let (isolate, runtime_inspector) = - worker_isolate.worker_isolate_and_runtime_inspector_mut(); + worker_isolate.worker_isolate_and_runtime_inspector(); let scope = pin!(v8::HandleScope::new(isolate)); let scope = &mut scope.init(); let ctx = v8::Local::new(scope, &context); @@ -3244,7 +3183,7 @@ async fn worker_main( &state, &module_graph_fetch_tx, ); - forward_worker_script_loaded(runtime_inspector, &parent_tx); + forward_worker_script_loaded(&runtime_inspector, &parent_tx); } WorkerModuleBootstrapResume::NeedFetches(requests) => { start_worker_module_graph_fetch_batch( @@ -3277,7 +3216,7 @@ async fn worker_main( &parent_tx, &script_url, ); - forward_worker_script_loaded(runtime_inspector, &parent_tx); + forward_worker_script_loaded(&runtime_inspector, &parent_tx); break; } } @@ -3355,7 +3294,7 @@ async fn worker_main( } forward_pending_worker_runtime_protocol_messages( - worker_isolate.worker_runtime_inspector_mut(), + worker_isolate.worker_runtime_inspector(), &parent_tx, ); @@ -3370,14 +3309,14 @@ async fn worker_main( // routes are destroyed. Ordinary transports are cancelled here; // explicitly keepalive loads are reduced to browser-runtime network-only // records and therefore cannot retain this WorkerGlobalScope. + inspector_task_runner.dispose("Worker exited before Inspector task dispatch"); let resource_loader = state.borrow().loader.clone(); resource_loader.begin_detach(); if matches!(state.borrow().global_kind, WorkerGlobalKind::Shared { .. }) { let _ = parent_tx.send(WorkerToParentMessage::SharedWorkerClosed); } { - let (isolate, runtime_inspector) = - worker_isolate.worker_isolate_and_runtime_inspector_mut(); + let (isolate, runtime_inspector) = worker_isolate.worker_isolate_and_runtime_inspector(); let scope = pin!(v8::HandleScope::new(isolate)); let scope = &mut scope.init(); let ctx = v8::Local::new(scope, &context); diff --git a/moli-renderer-v8/src/worker/thread/runtime_inspector.rs b/moli-renderer-v8/src/worker/thread/runtime_inspector.rs index 854952c69d..c097b7ae78 100644 --- a/moli-renderer-v8/src/worker/thread/runtime_inspector.rs +++ b/moli-renderer-v8/src/worker/thread/runtime_inspector.rs @@ -1,15 +1,23 @@ use std::{ - cell::{RefCell, UnsafeCell}, - collections::{HashMap, VecDeque}, - rc::Rc, + cell::{Cell, RefCell, UnsafeCell}, + collections::{HashMap, HashSet, VecDeque}, + rc::{Rc, Weak}, }; use serde_json::{Value, json}; use crate::inspector_microtasks::with_scoped_inspector_microtasks; use crate::runtime::{RendererRuntimeInspectorMessage, RendererRuntimeInspectorResponseSender}; -use crate::worker::handle::WorkerRuntimeInspectorMessageBatch; +use crate::worker::{ + handle::{WorkerRuntimeInspectorMessageBatch, WorkerToParentMessage}, + inspector_task_runner::{ + WorkerInspectorInterruptExecutor, WorkerInspectorTask, WorkerInspectorTaskRunner, + register_worker_inspector_executor, unregister_worker_inspector_executor, + }, +}; +use tokio::sync::mpsc; +#[cfg(test)] use super::dispatch::perform_worker_microtask_checkpoint_and_report_pending_promise_rejections; const WORKER_INSPECTOR_CONTEXT_GROUP_ID: i32 = 1; @@ -216,6 +224,26 @@ impl WorkerInspectorOutbound { ) } + fn take_active_dispatch_notifications(&self) -> Vec { + let mut state = self.0.borrow_mut(); + let mut messages = Vec::new(); + for scope in &mut state.active_dispatch_scopes { + let mut responses = Vec::new(); + for message in scope.messages.drain(..) { + if message.has_v8_inspector_method() { + messages.push(WorkerInspectorPendingMessage { + session_key: scope.session_key.clone(), + message, + }); + } else { + responses.push(message); + } + } + scope.messages = responses; + } + coalesce_worker_inspector_messages(messages) + } + fn push_dispatch_scope(&self, session_key: &str) -> WorkerInspectorDispatchScopeGuard { self.0 .borrow_mut() @@ -252,21 +280,32 @@ impl v8::inspector::ChannelImpl for WorkerInspectorChannel { struct WorkerInspectorClient { isolate: UnsafeCell, default_context: Rc>>>, + executor: Rc, } impl WorkerInspectorClient { fn new( isolate: v8::UnsafeRawIsolatePtr, default_context: Rc>>>, + executor: Rc, ) -> Self { Self { isolate: UnsafeCell::new(isolate), default_context, + executor, } } } impl v8::inspector::V8InspectorClientImpl for WorkerInspectorClient { + fn run_message_loop_on_pause(&self, _context_group_id: i32) { + self.executor.run_pause_loop(); + } + + fn quit_message_loop_on_pause(&self) { + self.executor.task_runner.request_quit_pause_loop(); + } + fn ensure_default_context_in_group( &self, context_group_id: i32, @@ -283,36 +322,113 @@ impl v8::inspector::V8InspectorClientImpl for WorkerInspectorClient { } } -struct WorkerInspectorSessionState { - session: v8::inspector::V8InspectorSession, -} - pub(super) struct WorkerRuntimeInspector { - sessions: HashMap, + sessions: RefCell>>, + detached_sessions: RefCell>, inspector: v8::inspector::V8Inspector, outbound: WorkerInspectorOutbound, default_context: Rc>>>, - default_execution_context_id: Option, + default_execution_context_id: Cell>, + task_runner: WorkerInspectorTaskRunner, + parent_tx: mpsc::UnboundedSender, + shared_worker: bool, +} + +struct WorkerInspectorExecutor { + isolate: UnsafeCell, + inspector: Weak, + task_runner: WorkerInspectorTaskRunner, +} + +impl Drop for WorkerInspectorExecutor { + fn drop(&mut self) { + unregister_worker_inspector_executor(self.task_runner.route_id()); + } +} + +impl WorkerInspectorExecutor { + fn run_pause_loop(&self) { + let Some(inspector) = self.inspector.upgrade() else { + return; + }; + inspector.forward_active_dispatch_messages(); + if !self.task_runner.begin_pause_loop() { + return; + } + while let Some(task) = self.task_runner.wait_for_pause_task() { + let isolate = unsafe { &mut *self.isolate.get() }; + let isolate = unsafe { v8::Isolate::ref_from_raw_isolate_ptr_mut(isolate) }; + inspector.execute_task(isolate, task); + } + self.task_runner.finish_pause_loop(); + } +} + +impl WorkerInspectorInterruptExecutor for WorkerInspectorExecutor { + fn dispatch_interrupt(&self, isolate: v8::UnsafeRawIsolatePtr) { + self.task_runner.interrupt_callback_started(); + let Some(task) = self.task_runner.claim_interrupting_task() else { + self.task_runner.request_interrupt_if_needed(); + return; + }; + if let Some(inspector) = self.inspector.upgrade() { + let mut isolate_ptr = isolate; + let isolate = unsafe { v8::Isolate::ref_from_raw_isolate_ptr_mut(&mut isolate_ptr) }; + inspector.execute_task(isolate, task); + } + self.task_runner.request_interrupt_if_needed(); + } } impl WorkerRuntimeInspector { - pub(super) fn new(isolate: &mut v8::Isolate) -> Self { + pub(super) fn new( + isolate: &mut v8::Isolate, + task_runner: WorkerInspectorTaskRunner, + parent_tx: mpsc::UnboundedSender, + shared_worker: bool, + ) -> Rc { let isolate_ptr = unsafe { isolate.as_raw_isolate_ptr() }; let default_context = Rc::new(RefCell::new(None)); - let inspector_client = v8::inspector::V8InspectorClient::new(Box::new( - WorkerInspectorClient::new(isolate_ptr, Rc::clone(&default_context)), - )); - Self { - inspector: v8::inspector::V8Inspector::create(isolate, inspector_client), - sessions: HashMap::new(), - outbound: WorkerInspectorOutbound::default(), - default_context, - default_execution_context_id: None, - } + Rc::new_cyclic(|weak_inspector| { + let executor = Rc::new(WorkerInspectorExecutor { + isolate: UnsafeCell::new(isolate_ptr), + inspector: weak_inspector.clone(), + task_runner: task_runner.clone(), + }); + let interrupt_executor: Rc = executor.clone(); + register_worker_inspector_executor(task_runner.route_id(), &interrupt_executor); + let inspector_client = + v8::inspector::V8InspectorClient::new(Box::new(WorkerInspectorClient::new( + isolate_ptr, + Rc::clone(&default_context), + executor.clone(), + ))); + Self { + inspector: v8::inspector::V8Inspector::create(isolate, inspector_client), + sessions: RefCell::new(HashMap::new()), + detached_sessions: RefCell::new(HashSet::new()), + outbound: WorkerInspectorOutbound::default(), + default_context, + default_execution_context_id: Cell::new(None), + task_runner, + parent_tx, + shared_worker, + } + }) + } + + #[cfg(test)] + fn new_for_test(isolate: &mut v8::Isolate) -> Rc { + let (wake_tx, _wake_rx) = mpsc::unbounded_channel(); + let (parent_tx, _parent_rx) = mpsc::unbounded_channel(); + let isolate_handle = + std::sync::Arc::new(parking_lot::Mutex::new(Some(isolate.thread_safe_handle()))); + let task_runner = WorkerInspectorTaskRunner::new(wake_tx, isolate_handle); + Self::new(isolate, task_runner, parent_tx, false) } pub(super) fn attach_context<'s>( - &mut self, + &self, context: v8::Local<'s, v8::Context>, default_context: v8::Global, script_url: &str, @@ -325,27 +441,30 @@ impl WorkerRuntimeInspector { v8::inspector::StringView::from(script_url.as_bytes()), v8::inspector::StringView::from(&br#"{"isDefault":true,"type":"worker"}"#[..]), ); - self.default_execution_context_id = Some(i64::from( + self.default_execution_context_id.set(Some(i64::from( v8::inspector::V8Inspector::execution_context_id(context), - )); + ))); } pub(super) fn context_destroyed<'s>(&self, context: v8::Local<'s, v8::Context>) { self.inspector.context_destroyed(context); } - pub(super) fn detach_session(&mut self, inspector_session_id: Option<&str>) { + pub(super) fn detach_session(&self, inspector_session_id: Option<&str>) { let session_key = worker_inspector_session_key(inspector_session_id); - self.sessions.remove(&session_key); + self.sessions.borrow_mut().remove(&session_key); + self.detached_sessions.borrow_mut().insert(session_key); } - pub(super) fn attach_session(&mut self, inspector_session_id: Option<&str>) { + pub(super) fn attach_session(&self, inspector_session_id: Option<&str>) { let session_key = worker_inspector_session_key(inspector_session_id); + self.detached_sessions.borrow_mut().remove(&session_key); let _ = self.ensure_session(&session_key); } + #[cfg(test)] pub(super) fn dispatch_protocol_message( - &mut self, + &self, scope: &mut v8::PinScope<'_, '_>, inspector_session_id: Option<&str>, raw_json: &str, @@ -358,45 +477,86 @@ impl WorkerRuntimeInspector { ) } - pub(super) fn dispatch_protocol_message_with_deferred_response( - &mut self, - scope: &mut v8::PinScope<'_, '_>, - inspector_session_id: Option<&str>, - raw_json: &str, - deferred_response: RendererRuntimeInspectorResponseSender, - ) -> Result, String> { - self.dispatch_protocol_message_with_optional_deferred_response( - scope, - inspector_session_id, - raw_json, - Some(deferred_response), - ) - } - + #[cfg(test)] fn dispatch_protocol_message_with_optional_deferred_response( - &mut self, + &self, scope: &mut v8::PinScope<'_, '_>, inspector_session_id: Option<&str>, raw_json: &str, deferred_response: Option, + ) -> Result, String> { + let messages = with_scoped_inspector_microtasks(scope, || { + self.dispatch_protocol_message_scoped(inspector_session_id, raw_json, deferred_response) + })?; + perform_worker_microtask_checkpoint_and_report_pending_promise_rejections(scope); + Ok(messages) + } + + fn dispatch_protocol_message_scoped( + &self, + inspector_session_id: Option<&str>, + raw_json: &str, + deferred_response: Option, ) -> Result, String> { let session_key = worker_inspector_session_key(inspector_session_id); + if self.detached_sessions.borrow().contains(&session_key) { + return Err("Worker Inspector session has been detached".to_owned()); + } if let Some(callback) = deferred_response { self.outbound .register_response_callback(&session_key, callback); } let dispatch_scope = self.outbound.push_dispatch_scope(&session_key); let session = self.ensure_session(&session_key); - with_scoped_inspector_microtasks(scope, || { - session.dispatch_protocol_message(v8::inspector::StringView::from(raw_json.as_bytes())); - }); - perform_worker_microtask_checkpoint_and_report_pending_promise_rejections(scope); + session.dispatch_protocol_message(v8::inspector::StringView::from(raw_json.as_bytes())); let messages = dispatch_scope.finish(); self.record_execution_context_state(&messages); Ok(messages) } - pub(super) fn take_pending_messages(&mut self) -> Vec { + pub(super) fn execute_task(&self, isolate: &mut v8::Isolate, task: WorkerInspectorTask) { + match task { + WorkerInspectorTask::DispatchProtocolMessage { + inspector_session_id, + raw_json, + deferred_response, + response_tx, + } => { + let deferred_response = if self.shared_worker { + deferred_response.map(|response| { + response.defer_publication_to_shared_worker_parent(self.parent_tx.clone()) + }) + } else { + deferred_response + }; + let result = with_scoped_inspector_microtasks(isolate, || { + self.dispatch_protocol_message_scoped( + inspector_session_id.as_deref(), + &raw_json, + deferred_response, + ) + }); + if result.is_ok() + && worker_runtime_protocol_message_is_run_if_waiting_for_debugger(&raw_json) + { + self.task_runner.request_resume(); + } + let _ = response_tx.send(result); + } + WorkerInspectorTask::AttachSession { + inspector_session_id, + } => self.attach_session(inspector_session_id.as_deref()), + WorkerInspectorTask::DetachSession { + inspector_session_id, + } => self.detach_session(inspector_session_id.as_deref()), + WorkerInspectorTask::RunIfWaitingForDebugger => { + self.task_runner.request_resume(); + } + } + self.forward_pending_messages(); + } + + pub(super) fn take_pending_messages(&self) -> Vec { let batches = self.outbound.take_all(); for batch in &batches { self.record_execution_context_state(&batch.messages); @@ -408,7 +568,8 @@ impl WorkerRuntimeInspector { } pub(super) fn worker_script_loaded_messages(&self) -> Vec { - let mut session_keys = self.sessions.keys().collect::>(); + let sessions = self.sessions.borrow(); + let mut session_keys = sessions.keys().collect::>(); session_keys.sort_unstable(); session_keys .into_iter() @@ -424,12 +585,12 @@ impl WorkerRuntimeInspector { .collect() } - fn ensure_session(&mut self, session_key: &str) -> &v8::inspector::V8InspectorSession { - &self - .sessions + fn ensure_session(&self, session_key: &str) -> Rc { + self.sessions + .borrow_mut() .entry(session_key.to_owned()) - .or_insert_with(|| WorkerInspectorSessionState { - session: self.inspector.connect( + .or_insert_with(|| { + Rc::new(self.inspector.connect( WORKER_INSPECTOR_CONTEXT_GROUP_ID, v8::inspector::Channel::new(Box::new(WorkerInspectorChannel { outbound: self.outbound.clone(), @@ -437,12 +598,12 @@ impl WorkerRuntimeInspector { })), v8::inspector::StringView::from(&b"{}"[..]), v8::inspector::V8InspectorClientTrustLevel::FullyTrusted, - ), + )) }) - .session + .clone() } - fn record_execution_context_state(&mut self, messages: &[RendererRuntimeInspectorMessage]) { + fn record_execution_context_state(&self, messages: &[RendererRuntimeInspectorMessage]) { for message in messages { match message { RendererRuntimeInspectorMessage::RuntimeContext( @@ -451,25 +612,59 @@ impl WorkerRuntimeInspector { if event.context_type.as_deref() == Some("worker") && let Some(id) = event.context_id { - self.default_execution_context_id = Some(id); + self.default_execution_context_id.set(Some(id)); } } RendererRuntimeInspectorMessage::RuntimeContext( crate::protocol_types::RuntimeContextRestoreEvent::Destroyed(event), ) => { - if event.context_id == self.default_execution_context_id { - self.default_execution_context_id = None; + if event.context_id == self.default_execution_context_id.get() { + self.default_execution_context_id.set(None); } } RendererRuntimeInspectorMessage::RuntimeContext( crate::protocol_types::RuntimeContextRestoreEvent::Cleared(_), ) => { - self.default_execution_context_id = None; + self.default_execution_context_id.set(None); } _ => {} } } } + + fn forward_pending_messages(&self) { + let messages = self.take_pending_messages(); + if !messages.is_empty() { + let _ = self + .parent_tx + .send(WorkerToParentMessage::RuntimeInspectorMessages(messages)); + } + } + + fn forward_active_dispatch_messages(&self) { + let batches = self.outbound.take_active_dispatch_notifications(); + for batch in &batches { + self.record_execution_context_state(&batch.messages); + } + let messages = batches + .into_iter() + .map(worker_runtime_inspector_message_batch) + .collect::>(); + if !messages.is_empty() { + let _ = self + .parent_tx + .send(WorkerToParentMessage::RuntimeInspectorMessages(messages)); + } + } +} + +fn worker_runtime_protocol_message_is_run_if_waiting_for_debugger(raw_json: &str) -> bool { + serde_json::from_str::(raw_json).is_ok_and(|value| { + value + .get("method") + .and_then(Value::as_str) + .is_some_and(|method| method == "Runtime.runIfWaitingForDebugger") + }) } fn worker_inspector_session_key(inspector_session_id: Option<&str>) -> String { @@ -558,7 +753,7 @@ mod tests { crate::ensure_v8_for_test(); let mut isolate = v8::Isolate::new(Default::default()); isolate.set_microtasks_policy(v8::MicrotasksPolicy::Explicit); - let mut inspector = WorkerRuntimeInspector::new(&mut isolate); + let inspector = WorkerRuntimeInspector::new_for_test(&mut isolate); let context = { let scope = pin!(v8::HandleScope::new(&mut isolate)); let scope = &mut scope.init(); @@ -772,4 +967,46 @@ mod tests { }] ); } + + #[test] + fn worker_pause_flushes_active_notifications_but_retains_command_response() { + let outbound = WorkerInspectorOutbound::default(); + let dispatch_scope = outbound.push_dispatch_scope("SID-1"); + outbound.push_value("SID-1", json!({"method": "Debugger.paused", "params": {}})); + outbound.push_response_value("SID-1", 7, json!({"id": 7, "result": {}})); + + assert_eq!( + outbound.take_active_dispatch_notifications(), + vec![WorkerInspectorPendingMessageBatch { + inspector_session_id: Some("SID-1".to_owned()), + messages: vec![inspector_message( + json!({"method": "Debugger.paused", "params": {}}) + )], + }] + ); + assert_eq!( + dispatch_scope.finish(), + vec![inspector_message(json!({"id": 7, "result": {}}))], + ); + } + + #[test] + fn detached_worker_inspector_session_is_not_lazily_recreated() { + crate::ensure_v8_for_test(); + let mut isolate = v8::Isolate::new(Default::default()); + let inspector = WorkerRuntimeInspector::new_for_test(&mut isolate); + inspector.attach_session(Some("SID-detached")); + inspector.detach_session(Some("SID-detached")); + + assert_eq!( + inspector + .dispatch_protocol_message_scoped( + Some("SID-detached"), + &json!({"id": 1, "method": "Runtime.enable"}).to_string(), + None, + ) + .expect_err("detached session must reject a queued late command"), + "Worker Inspector session has been detached" + ); + } } diff --git a/moli-renderer-v8/src/worker/thread/tests/lifecycle.rs b/moli-renderer-v8/src/worker/thread/tests/lifecycle.rs index a090037c85..fc59d50dc1 100644 --- a/moli-renderer-v8/src/worker/thread/tests/lifecycle.rs +++ b/moli-renderer-v8/src/worker/thread/tests/lifecycle.rs @@ -363,14 +363,10 @@ async fn worker_attach_before_runtime_enable_preserves_script_loaded_after_resum .with_pause_evaluation_until_debugger(true), ); - handle - .tx - .send( - crate::worker::WorkerMessage::AttachRuntimeInspectorSession { - inspector_session_id: Some("SID-worker-attached".to_owned()), - }, - ) - .expect("worker should accept its Inspector session before the first command"); + assert!( + handle.attach_runtime_inspector_session(Some("SID-worker-attached".to_owned())), + "worker should accept its Inspector session before the first command" + ); assert!( handle.run_if_waiting_for_debugger_for_devtools(), "worker should accept debugger resume after the Inspector session is attached" @@ -414,6 +410,137 @@ async fn worker_attach_before_runtime_enable_preserves_script_loaded_after_resum handle.terminate_and_join(); } +#[tokio::test] +async fn worker_inspector_interrupt_overtakes_js_running_command_during_active_javascript() { + ensure_v8(); + let mut handle = spawn_worker( + r#" + let deliveries = 0; + self.onmessage = () => { + deliveries += 1; + postMessage(deliveries === 1 ? "entered" : "recovered"); + if (deliveries === 1) { + while (true) {} + } + }; + postMessage("ready"); + "# + .to_owned(), + "https://example.test/app/interruptible-worker.js".to_owned(), + ); + + fn dispatch_runtime( + handle: &WorkerHandle, + id: i64, + method: &str, + params: Option, + ) -> oneshot::Receiver, String>> + { + let (response_tx, response_rx) = oneshot::channel(); + let mut message = serde_json::json!({ + "id": id, + "method": method, + }); + if let Some(params) = params { + message["params"] = params; + } + assert!( + handle.dispatch_runtime_protocol_message( + Some("SID-worker-interrupt".to_owned()), + message.to_string(), + None, + response_tx, + ), + "worker should accept {method}" + ); + response_rx + } + + let ready = timeout(TIMEOUT, handle.recv()) + .await + .expect("timed out waiting for worker readiness") + .expect("worker closed before readiness"); + assert!(matches!( + ready, + WorkerToParentMessage::Post(ref payload) if stringify_payload(payload) == r#""ready""# + )); + + let enable = dispatch_runtime(&handle, 1, "Runtime.enable", None); + timeout(TIMEOUT, enable) + .await + .expect("timed out enabling worker Runtime") + .expect("worker Runtime enable response channel closed") + .expect("worker Runtime.enable failed"); + + handle.post_message(serialize_test_string("start")); + let entered = timeout(TIMEOUT, handle.recv()) + .await + .expect("timed out waiting for active worker JavaScript") + .expect("worker closed before entering active JavaScript"); + assert!(matches!( + entered, + WorkerToParentMessage::Post(ref payload) if stringify_payload(payload) == r#""entered""# + )); + + let mut evaluate = dispatch_runtime( + &handle, + 2, + "Runtime.evaluate", + Some(serde_json::json!({ + "expression": "40 + 2", + "returnByValue": true, + })), + ); + assert!( + timeout(Duration::from_millis(100), &mut evaluate) + .await + .is_err(), + "Runtime.evaluate must not interrupt active worker JavaScript" + ); + + let terminate = dispatch_runtime(&handle, 3, "Runtime.terminateExecution", None); + timeout(TIMEOUT, terminate) + .await + .expect("Runtime.terminateExecution did not interrupt active worker JavaScript") + .expect("worker Runtime termination response channel closed") + .expect("worker Runtime.terminateExecution failed"); + + let evaluate_messages = timeout(TIMEOUT, &mut evaluate) + .await + .expect("queued Runtime.evaluate did not run after termination") + .expect("queued Runtime.evaluate response channel closed") + .expect("queued Runtime.evaluate failed after termination"); + let evaluate_response = evaluate_messages + .into_iter() + .map(crate::runtime::RendererRuntimeInspectorMessage::into_v8_inspector_message) + .find(|message| message["id"] == 2) + .expect("queued Runtime.evaluate response"); + assert_eq!(evaluate_response["result"]["result"]["value"], 42); + + handle.post_message(serialize_test_string("again")); + loop { + let recovered = timeout(TIMEOUT, handle.recv()) + .await + .expect("timed out waiting for worker recovery") + .expect("worker closed instead of recovering"); + match recovered { + WorkerToParentMessage::Post(payload) => { + assert_eq!(stringify_payload(&payload), r#""recovered""#); + break; + } + WorkerToParentMessage::RuntimeInspectorMessages(_) => {} + WorkerToParentMessage::Error { message, .. } => { + assert_eq!( + message, "null", + "worker emitted an unexpected error while recovering" + ); + } + other => panic!("unexpected worker output while checking recovery: {other:?}"), + } + } + handle.terminate_and_join(); +} + #[tokio::test] async fn real_workers_register_service_worker_clients_until_thread_exit() { ensure_v8();