refactor(renderer): remove unused worker task generations

This commit is contained in:
ldm0
2026-08-18 02:26:17 +08:00
committed by Donough Liu
parent 5785cbeb8b
commit 51e5d208e4
4 changed files with 34 additions and 87 deletions
@@ -11,7 +11,6 @@ pub(crate) enum WebCryptoCompletionSink {
Page(crate::page_task_queue::RendererPageWebCryptoTaskProducer),
Worker {
task_id: u64,
generation: u64,
tx: tokio::sync::mpsc::UnboundedSender<crate::worker::WorkerWebCryptoCompletion>,
},
}
@@ -23,16 +22,8 @@ impl WebCryptoCompletionSink {
WebCryptoCompletionSink::Page(producer) => {
let _ = producer.send(result);
}
WebCryptoCompletionSink::Worker {
task_id,
generation,
tx,
} => {
let _ = tx.send(crate::worker::WorkerWebCryptoCompletion {
task_id,
generation,
result,
});
WebCryptoCompletionSink::Worker { task_id, tx } => {
let _ = tx.send(crate::worker::WorkerWebCryptoCompletion { task_id, result });
}
}
}
@@ -57,13 +48,12 @@ pub(crate) fn register_webcrypto_task<'s>(
unsafe { &mut *host_ptr }.register_pending_webcrypto_task(scope, promise.resolver())?;
return Some((handle, WebCryptoCompletionSink::Page(producer)));
}
let (task_id, generation, completion_tx) =
let (task_id, completion_tx) =
crate::worker::register_worker_webcrypto_task(scope, promise.resolver())?;
Some((
handle,
WebCryptoCompletionSink::Worker {
task_id,
generation,
tx: completion_tx,
},
))
+11 -25
View File
@@ -350,7 +350,6 @@ enum OpfsCompletionSink {
Page(crate::page_task_queue::RendererPageOpfsTaskProducer),
Worker {
task_id: u64,
generation: u64,
sender: tokio::sync::mpsc::UnboundedSender<crate::worker::WorkerOpfsCompletion>,
},
}
@@ -359,23 +358,15 @@ impl OpfsCompletionSink {
/// Consume the exact terminal capability.
///
/// A Page capability already contains its task identity. Worker identity
/// stays local to the Worker transport because that queue still performs
/// its own generation authorization.
/// stays local to the Worker transport and is valid for that worker state's
/// lifetime.
fn send(self, result: OpfsTaskResult) {
match self {
Self::Page(sender) => {
let _ = sender.send(result);
}
Self::Worker {
task_id,
generation,
sender,
} => {
let _ = sender.send(crate::worker::WorkerOpfsCompletion {
task_id,
generation,
result,
});
Self::Worker { task_id, sender } => {
let _ = sender.send(crate::worker::WorkerOpfsCompletion { task_id, result });
}
}
}
@@ -407,16 +398,11 @@ impl RegisteredOpfsTask {
fn worker(
task_id: u64,
generation: u64,
sender: tokio::sync::mpsc::UnboundedSender<crate::worker::WorkerOpfsCompletion>,
) -> Self {
Self {
cancellation: RegisteredOpfsTaskCancellation::Worker { task_id },
completion: OpfsCompletionSink::Worker {
task_id,
generation,
sender,
},
completion: OpfsCompletionSink::Worker { task_id, sender },
}
}
@@ -443,9 +429,9 @@ fn register_opfs_task<'s>(
)?;
return Some(RegisteredOpfsTask::page(task_id, sender));
}
let (task_id, generation, sender) =
let (task_id, sender) =
crate::worker::register_worker_opfs_task(scope, resolver, locator, handle_access)?;
Some(RegisteredOpfsTask::worker(task_id, generation, sender))
Some(RegisteredOpfsTask::worker(task_id, sender))
}
fn register_opfs_iterator_task<'s>(
@@ -470,7 +456,7 @@ fn register_opfs_iterator_task<'s>(
)?;
return Some(RegisteredOpfsTask::page(task_id, sender));
}
let (task_id, generation, sender) = crate::worker::register_worker_opfs_iterator_task(
let (task_id, sender) = crate::worker::register_worker_opfs_iterator_task(
scope,
locator,
registry,
@@ -478,7 +464,7 @@ fn register_opfs_iterator_task<'s>(
v8::Global::new(scope, iterator),
handle_access,
)?;
Some(RegisteredOpfsTask::worker(task_id, generation, sender))
Some(RegisteredOpfsTask::worker(task_id, sender))
}
fn register_opfs_move_task<'s>(
@@ -504,7 +490,7 @@ fn register_opfs_move_task<'s>(
)?;
return Some(RegisteredOpfsTask::page(task_id, sender));
}
let (task_id, generation, sender) = crate::worker::register_worker_opfs_move_task(
let (task_id, sender) = crate::worker::register_worker_opfs_move_task(
scope,
resolver,
handle,
@@ -512,7 +498,7 @@ fn register_opfs_move_task<'s>(
locator,
handle_access,
)?;
Some(RegisteredOpfsTask::worker(task_id, generation, sender))
Some(RegisteredOpfsTask::worker(task_id, sender))
}
fn current_opfs_handle_registry(scope: &mut v8::PinScope<'_, '_>) -> Option<OpfsHandleRegistry> {
+20 -47
View File
@@ -1037,10 +1037,9 @@ pub(super) enum PendingServiceWorkerClientQueryType {
}
/// A WebCrypto primitive dispatched off the worker event loop to the blocking
/// pool. The worker owns the resolver; completion is matched back by `task_id`
/// and dropped if the worker generation advanced (navigation/`close()` reset).
/// pool. The worker owns the resolver and matches the completion back by the
/// task id allocated for this worker lifetime.
pub(super) struct PendingWorkerWebCryptoTask {
pub(super) generation: u64,
pub(super) resolver: v8::Global<v8::PromiseResolver>,
}
@@ -1049,12 +1048,10 @@ pub(super) struct PendingWorkerWebCryptoTask {
/// page-side typed producer captures its exact Page/Window owner.
pub(crate) struct WorkerWebCryptoCompletion {
pub(crate) task_id: u64,
pub(crate) generation: u64,
pub(crate) result: Result<WebCryptoTaskResult, WebCryptoRejection>,
}
pub(super) struct PendingWorkerOpfsTask {
pub(super) generation: u64,
pub(super) locator: moli_storage_service::StorageBucketLocator,
pub(super) handle_access: Option<crate::opfs_owner_tasks::OpfsHandleAccessContext>,
pub(super) settlement: crate::opfs_owner_tasks::OpfsTaskSettlement,
@@ -1086,7 +1083,6 @@ impl WorkerOpfsOwnerState {
pub(crate) struct WorkerOpfsCompletion {
pub(crate) task_id: u64,
pub(crate) generation: u64,
pub(crate) result: OpfsTaskResult,
}
@@ -1597,15 +1593,8 @@ pub(crate) struct WorkerGlobalState {
pub(super) pending_webcrypto: HashMap<u64, PendingWorkerWebCryptoTask>,
/// Worker WebCrypto task id counter (never zero so 0 can mean "unset").
pub(super) next_webcrypto_task_id: u64,
/// Worker WebCrypto generation; bumping it invalidates in-flight completions.
/// Workers have no in-place context reset (terminate/`close()` drops the
/// whole state and tears down the isolate), so this is currently constant.
/// It mirrors the page-side generation guard and would let a future in-place
/// worker reset drop stale blocking-task completions without a code change.
pub(super) webcrypto_generation: u64,
/// OPFS completions routed from the partition-owned storage IO sequence.
pub(super) opfs_completion_tx: mpsc::UnboundedSender<WorkerOpfsCompletion>,
pub(super) opfs_generation: u64,
pub(super) opfs_owner_state: Option<WorkerOpfsOwnerState>,
/// In-flight Service Worker `periodicsync` events keyed by runtime event id.
pub(super) pending_service_worker_periodic_sync_events:
@@ -1715,18 +1704,14 @@ impl WorkerGlobalState {
pub(super) fn register_pending_webcrypto_task(
&mut self,
resolver: v8::Global<v8::PromiseResolver>,
) -> (u64, u64, mpsc::UnboundedSender<WorkerWebCryptoCompletion>) {
) -> (u64, mpsc::UnboundedSender<WorkerWebCryptoCompletion>) {
let task_id = self.next_webcrypto_task_id;
self.next_webcrypto_task_id = self.next_webcrypto_task_id.wrapping_add(1).max(1);
let generation = self.webcrypto_generation;
self.pending_webcrypto.insert(
task_id,
PendingWorkerWebCryptoTask {
generation,
resolver,
},
);
(task_id, generation, self.webcrypto_completion_tx.clone())
self.next_webcrypto_task_id = task_id
.checked_add(1)
.expect("worker WebCrypto task id exhausted");
self.pending_webcrypto
.insert(task_id, PendingWorkerWebCryptoTask { resolver });
(task_id, self.webcrypto_completion_tx.clone())
}
pub(super) fn take_pending_webcrypto_task(
@@ -1741,23 +1726,23 @@ impl WorkerGlobalState {
locator: moli_storage_service::StorageBucketLocator,
handle_access: Option<crate::opfs_owner_tasks::OpfsHandleAccessContext>,
settlement: crate::opfs_owner_tasks::OpfsTaskSettlement,
) -> (u64, u64, mpsc::UnboundedSender<WorkerOpfsCompletion>) {
) -> (u64, mpsc::UnboundedSender<WorkerOpfsCompletion>) {
let state = self
.opfs_owner_state
.get_or_insert_with(WorkerOpfsOwnerState::default);
let task_id = state.next_task_id;
state.next_task_id = state.next_task_id.wrapping_add(1).max(1);
let generation = self.opfs_generation;
state.next_task_id = task_id
.checked_add(1)
.expect("worker OPFS task id exhausted");
state.pending_tasks.insert(
task_id,
PendingWorkerOpfsTask {
generation,
locator,
handle_access,
settlement,
},
);
(task_id, generation, self.opfs_completion_tx.clone())
(task_id, self.opfs_completion_tx.clone())
}
pub(super) fn take_pending_opfs_task(&mut self, task_id: u64) -> Option<PendingWorkerOpfsTask> {
@@ -2138,7 +2123,7 @@ impl WorkerGlobalState {
pub(crate) fn register_worker_webcrypto_task(
scope: &mut v8::PinScope<'_, '_>,
resolver: v8::Local<'_, v8::PromiseResolver>,
) -> Option<(u64, u64, mpsc::UnboundedSender<WorkerWebCryptoCompletion>)> {
) -> Option<(u64, mpsc::UnboundedSender<WorkerWebCryptoCompletion>)> {
let state = get_worker_state(scope)?;
let resolver = v8::Global::new(scope, resolver);
let mut state = state.borrow_mut();
@@ -2150,7 +2135,7 @@ pub(crate) fn register_worker_opfs_task(
resolver: v8::Local<'_, v8::PromiseResolver>,
locator: moli_storage_service::StorageBucketLocator,
handle_access: Option<crate::opfs_owner_tasks::OpfsHandleAccessContext>,
) -> Option<(u64, u64, mpsc::UnboundedSender<WorkerOpfsCompletion>)> {
) -> Option<(u64, mpsc::UnboundedSender<WorkerOpfsCompletion>)> {
let state = get_worker_state(scope)?;
let resolver = v8::Global::new(scope, resolver);
let mut state = state.borrow_mut();
@@ -2168,7 +2153,7 @@ pub(crate) fn register_worker_opfs_iterator_task(
iterator_id: u32,
keep_alive: v8::Global<v8::Object>,
handle_access: Option<crate::opfs_owner_tasks::OpfsHandleAccessContext>,
) -> Option<(u64, u64, mpsc::UnboundedSender<WorkerOpfsCompletion>)> {
) -> Option<(u64, mpsc::UnboundedSender<WorkerOpfsCompletion>)> {
let state = get_worker_state(scope)?;
let mut state = state.borrow_mut();
Some(state.register_pending_opfs_task(
@@ -2189,7 +2174,7 @@ pub(crate) fn register_worker_opfs_move_task(
mutation: crate::opfs_owner_tasks::OpfsHandleMutationGuard,
locator: moli_storage_service::StorageBucketLocator,
handle_access: Option<crate::opfs_owner_tasks::OpfsHandleAccessContext>,
) -> Option<(u64, u64, mpsc::UnboundedSender<WorkerOpfsCompletion>)> {
) -> Option<(u64, mpsc::UnboundedSender<WorkerOpfsCompletion>)> {
let state = get_worker_state(scope)?;
let mut state = state.borrow_mut();
Some(state.register_pending_opfs_task(
@@ -2262,8 +2247,8 @@ pub(crate) fn cancel_worker_opfs_task(scope: &mut v8::PinScope<'_, '_>, task_id:
}
/// Settle a worker WebCrypto promise once its blocking task reports back on the
/// worker event loop. Stale completions (generation advanced by a worker reset)
/// are dropped without touching the resolver.
/// worker event loop. Terminating a worker drops this state and its completion
/// receiver together, so a separate reset generation is unnecessary.
pub(super) fn drain_worker_webcrypto_completion(
scope: &mut v8::PinScope<'_, '_>,
state: &Rc<RefCell<WorkerGlobalState>>,
@@ -2271,15 +2256,9 @@ pub(super) fn drain_worker_webcrypto_completion(
) {
let pending = {
let mut state = state.borrow_mut();
let current_generation = state.webcrypto_generation;
let Some(pending) = state.take_pending_webcrypto_task(completion.task_id) else {
return;
};
if pending.generation != completion.generation
|| completion.generation != current_generation
{
return;
}
pending
};
let resolver = v8::Local::new(scope, &pending.resolver);
@@ -2297,15 +2276,9 @@ pub(super) fn drain_worker_opfs_completion(
) {
let pending = {
let mut state = state.borrow_mut();
let current_generation = state.opfs_generation;
let Some(pending) = state.take_pending_opfs_task(completion.task_id) else {
return;
};
if pending.generation != completion.generation
|| completion.generation != current_generation
{
return;
}
pending
};
let handle_access = pending.handle_access;
@@ -1733,9 +1733,7 @@ async fn worker_main(
webcrypto_completion_tx,
pending_webcrypto: std::collections::HashMap::new(),
next_webcrypto_task_id: 1,
webcrypto_generation: 0,
opfs_completion_tx,
opfs_generation: 0,
opfs_owner_state: None,
pending_service_worker_periodic_sync_events: std::collections::HashMap::new(),
pending_service_worker_client_queries: std::collections::HashMap::new(),