fix(devtools): split Worker inspector execution paths

This commit is contained in:
ldm0
2026-08-18 17:49:10 +08:00
parent 0ae3508a87
commit 7f62357cd2
11 changed files with 1339 additions and 332 deletions
@@ -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<HashMap<u64, tokio::sync::mpsc::UnboundedSender<crate::worker::WorkerMessage>>>,
dedicated_worker_devtools_handles: Mutex<HashMap<u64, crate::worker::WorkerDevToolsHandle>>,
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,
@@ -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<WorkerMessage>,
) {
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<mpsc::UnboundedSender<WorkerMessage>> {
fn dedicated_worker_devtools_handle(&self, instance_id: u64) -> Option<WorkerDevToolsHandle> {
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<RendererRuntimeInspectorResponseSender>,
) -> Result<Vec<RendererRuntimeInspectorMessage>, 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<String>,
) -> 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<String>,
) -> 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())
}
}
@@ -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<WorkerDevToolsHandle> {
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;
@@ -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<RendererRuntimeInspectorResponseSender>,
) -> Result<Vec<RendererRuntimeInspectorMessage>, 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<String>,
) -> 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))
}
}
+115 -33
View File
@@ -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<WorkerToParentMessage>,
},
/// Dispatch a CDP Runtime protocol message inside this worker's V8 inspector session.
DispatchRuntimeProtocolMessage {
inspector_session_id: Option<String>,
raw_json: String,
deferred_response: Option<RendererRuntimeInspectorResponseSender>,
response_tx: oneshot::Sender<Result<Vec<RendererRuntimeInspectorMessage>, String>>,
},
/// Attach one renderer-side V8 inspector session before its first command.
AttachRuntimeInspectorSession {
inspector_session_id: Option<String>,
},
/// 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<String>,
},
/// 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<WorkerMessage>,
inspector_tasks: WorkerInspectorTaskRunner,
}
impl WorkerDevToolsHandle {
pub(crate) fn new(
wake_tx: mpsc::UnboundedSender<WorkerMessage>,
isolate_handle: Arc<Mutex<Option<v8::IsolateHandle>>>,
) -> 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<String>,
raw_json: String,
deferred_response: Option<RendererRuntimeInspectorResponseSender>,
response_tx: oneshot::Sender<Result<Vec<RendererRuntimeInspectorMessage>, 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<String>,
) -> bool {
self.inspector_tasks.append_attach(inspector_session_id)
}
pub(crate) fn detach_runtime_inspector_session(
&self,
inspector_session_id: Option<String>,
) -> 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<WorkerMessage>,
@@ -700,6 +752,7 @@ pub(crate) struct WorkerHandle {
join_handle: Option<std::thread::JoinHandle<()>>,
isolate_handle: Arc<Mutex<Option<v8::IsolateHandle>>>,
termination_requested: Arc<AtomicBool>,
devtools: WorkerDevToolsHandle,
}
impl WorkerHandle {
@@ -719,12 +772,32 @@ impl WorkerHandle {
)
}
#[cfg(test)]
pub(crate) fn new_with_termination_requested(
tx: mpsc::UnboundedSender<WorkerMessage>,
rx: mpsc::UnboundedReceiver<WorkerToParentMessage>,
join_handle: std::thread::JoinHandle<()>,
isolate_handle: Arc<Mutex<Option<v8::IsolateHandle>>>,
termination_requested: Arc<AtomicBool>,
) -> 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<WorkerMessage>,
rx: mpsc::UnboundedReceiver<WorkerToParentMessage>,
join_handle: std::thread::JoinHandle<()>,
isolate_handle: Arc<Mutex<Option<v8::IsolateHandle>>>,
termination_requested: Arc<AtomicBool>,
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<RendererRuntimeInspectorResponseSender>,
response_tx: oneshot::Sender<Result<Vec<RendererRuntimeInspectorMessage>, 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<String>,
) -> bool {
self.devtools
.attach_runtime_inspector_session(inspector_session_id)
}
pub(crate) fn detach_runtime_inspector_session(
&self,
inspector_session_id: Option<String>,
) -> 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)]) {
@@ -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<HashMap<u64, Weak<dyn WorkerInspectorInterruptExecutor>>> =
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<dyn WorkerInspectorInterruptExecutor>,
) {
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<Rc<dyn WorkerInspectorInterruptExecutor>> {
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::<WorkerInspectorInterruptTarget>()) };
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<String>,
raw_json: String,
deferred_response: Option<RendererRuntimeInspectorResponseSender>,
response_tx: oneshot::Sender<Result<Vec<RendererRuntimeInspectorMessage>, String>>,
},
AttachSession {
inspector_session_id: Option<String>,
},
DetachSession {
inspector_session_id: Option<String>,
},
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<WorkerInspectorTaskEntry>,
non_interrupting_tasks: VecDeque<WorkerInspectorTaskEntry>,
next_sequence: u64,
disposed: bool,
isolate_ready: bool,
pause_loop_active: bool,
quit_pause_loop: bool,
}
struct WorkerInspectorTaskRunnerShared {
state: Mutex<WorkerInspectorTaskRunnerState>,
pause_work: Condvar,
wake_tx: mpsc::UnboundedSender<WorkerMessage>,
isolate_handle: Arc<Mutex<Option<v8::IsolateHandle>>>,
interrupt_target: Arc<WorkerInspectorInterruptTarget>,
interrupt_armed: AtomicBool,
resume_requested: AtomicBool,
}
#[derive(Clone)]
pub(crate) struct WorkerInspectorTaskRunner {
shared: Arc<WorkerInspectorTaskRunnerShared>,
}
impl WorkerInspectorTaskRunner {
pub(crate) fn new(
wake_tx: mpsc::UnboundedSender<WorkerMessage>,
isolate_handle: Arc<Mutex<Option<v8::IsolateHandle>>>,
) -> 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<String>,
raw_json: String,
deferred_response: Option<RendererRuntimeInspectorResponseSender>,
response_tx: oneshot::Sender<Result<Vec<RendererRuntimeInspectorMessage>, String>>,
) -> bool {
let mode = serde_json::from_str::<serde_json::Value>(&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<String>) -> bool {
self.append(
WorkerInspectorTaskMode::Interrupt,
WorkerInspectorTask::AttachSession {
inspector_session_id,
},
)
}
pub(crate) fn append_detach(&self, inspector_session_id: Option<String>) -> 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<WorkerInspectorTask> {
self.shared
.state
.lock()
.interrupting_tasks
.pop_front()
.map(|entry| entry.task)
}
pub(crate) fn claim_non_interrupting_task(&self) -> Option<WorkerInspectorTask> {
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::<c_void>();
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<WorkerInspectorTask> {
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::<Vec<_>>();
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<super::WorkerMessage>,
) {
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<Result<Vec<crate::runtime::RendererRuntimeInspectorMessage>, 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"
);
}
}
+2 -1
View File
@@ -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,
+25 -8
View File
@@ -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<WorkerRuntimeInspector>,
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<WorkerToParentMessage>,
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<WorkerRuntimeInspector>) {
(&mut self.isolate, Rc::clone(&self.runtime_inspector))
}
pub(super) fn unregister_worker_isolate_platform(&mut self) {
+91 -152
View File
@@ -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<v8::Context>,
inspector: &mut WorkerRuntimeInspector,
inspector_session_id: Option<&str>,
raw_json: &str,
mut deferred_response: Option<crate::runtime::RendererRuntimeInspectorResponseSender>,
task: WorkerInspectorTask,
state: &Rc<RefCell<WorkerGlobalState>>,
module_graph_fetch_tx: &mpsc::UnboundedSender<WorkerModuleGraphFetchCompletion>,
parent_tx: &mpsc::UnboundedSender<WorkerToParentMessage>,
) -> Result<Vec<RendererRuntimeInspectorMessage>, 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<v8::Context>,
state: &Rc<RefCell<WorkerGlobalState>>,
module_graph_fetch_tx: &mpsc::UnboundedSender<WorkerModuleGraphFetchCompletion>,
parent_tx: &mpsc::UnboundedSender<WorkerToParentMessage>,
script_url: &str,
rx: &mut mpsc::UnboundedReceiver<WorkerMessage>,
pending_bootstrap_messages: &mut VecDeque<WorkerMessage>,
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::<serde_json::Value>(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<WorkerRuntimeInspectorMessageBatch> {
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<WorkerToParentMessage>,
) {
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<WorkerToParentMessage>,
) {
// 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<WorkerToParentMessage>,
isolate_handle: Arc<Mutex<Option<v8::IsolateHandle>>>,
termination_requested: Arc<AtomicBool>,
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);
@@ -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<WorkerInspectorPendingMessageBatch> {
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<v8::UnsafeRawIsolatePtr>,
default_context: Rc<RefCell<Option<v8::Global<v8::Context>>>>,
executor: Rc<WorkerInspectorExecutor>,
}
impl WorkerInspectorClient {
fn new(
isolate: v8::UnsafeRawIsolatePtr,
default_context: Rc<RefCell<Option<v8::Global<v8::Context>>>>,
executor: Rc<WorkerInspectorExecutor>,
) -> 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<String, WorkerInspectorSessionState>,
sessions: RefCell<HashMap<String, Rc<v8::inspector::V8InspectorSession>>>,
detached_sessions: RefCell<HashSet<String>>,
inspector: v8::inspector::V8Inspector,
outbound: WorkerInspectorOutbound,
default_context: Rc<RefCell<Option<v8::Global<v8::Context>>>>,
default_execution_context_id: Option<i64>,
default_execution_context_id: Cell<Option<i64>>,
task_runner: WorkerInspectorTaskRunner,
parent_tx: mpsc::UnboundedSender<WorkerToParentMessage>,
shared_worker: bool,
}
struct WorkerInspectorExecutor {
isolate: UnsafeCell<v8::UnsafeRawIsolatePtr>,
inspector: Weak<WorkerRuntimeInspector>,
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<WorkerToParentMessage>,
shared_worker: bool,
) -> Rc<Self> {
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<dyn WorkerInspectorInterruptExecutor> = 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<Self> {
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<v8::Context>,
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<Vec<RendererRuntimeInspectorMessage>, 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<RendererRuntimeInspectorResponseSender>,
) -> Result<Vec<RendererRuntimeInspectorMessage>, 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<RendererRuntimeInspectorResponseSender>,
) -> Result<Vec<RendererRuntimeInspectorMessage>, 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<WorkerRuntimeInspectorMessageBatch> {
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<WorkerRuntimeInspectorMessageBatch> {
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<WorkerRuntimeInspectorMessageBatch> {
let mut session_keys = self.sessions.keys().collect::<Vec<_>>();
let sessions = self.sessions.borrow();
let mut session_keys = sessions.keys().collect::<Vec<_>>();
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<v8::inspector::V8InspectorSession> {
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::<Vec<_>>();
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::<Value>(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"
);
}
}
@@ -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<serde_json::Value>,
) -> oneshot::Receiver<Result<Vec<crate::runtime::RendererRuntimeInspectorMessage>, 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();