mirror of
https://github.com/lexmount/moli.git
synced 2026-10-05 16:00:54 +00:00
fix(scheduler): let replacement navigation pass load-blocked work
Centralize navigation dependency decisions for readiness probing and snapshot draining. Replacement actions can reach the owner that cancels the pending request even when earlier observations wait for that load, while exact owner authorization remains mandatory. Keep retirement admission policy beside publication ownership. Extend coverage for stale owners, cancellation, reused frontend ids, duplicate terminals, and late results without reviving retired notifications or remote objects. The new queue regression fails before the scheduler fix. Validation: cargo fmt --all; cargo clippy --workspace --all-targets --all-features -- -D warnings; cargo nextest run --no-fail-fast (18833 passed, 16 skipped on the final change set).
This commit is contained in:
@@ -1883,12 +1883,11 @@ impl CdpScheduler {
|
||||
return out;
|
||||
}
|
||||
let target_ids = self.protocol_residence_navigation_gate_target_ids(&residence);
|
||||
let blocked_by_prior_residence = target_ids
|
||||
.iter()
|
||||
.any(|target_id| blocked_target_ids.contains(target_id));
|
||||
let blocked_by_navigation = !residence.bypasses_inflight_navigation_gate()
|
||||
&& self.protocol_targets_have_inflight_background_navigation(&target_ids);
|
||||
if blocked_by_prior_residence || blocked_by_navigation {
|
||||
if self.protocol_residence_waits_for_navigation(
|
||||
&residence,
|
||||
&target_ids,
|
||||
&blocked_target_ids,
|
||||
) {
|
||||
if target_ids.is_empty() {
|
||||
retained.push_back(residence);
|
||||
retained.append(&mut snapshot);
|
||||
@@ -1946,6 +1945,30 @@ impl CdpScheduler {
|
||||
})
|
||||
}
|
||||
|
||||
/// Share this decision between readiness probing and snapshot draining.
|
||||
/// Otherwise a load-blocked observation in front can indirectly reapply
|
||||
/// the very gate that a replacement navigation is allowed to bypass.
|
||||
fn protocol_residence_waits_for_navigation(
|
||||
&self,
|
||||
residence: &ProtocolSchedulerResidence,
|
||||
target_ids: &[String],
|
||||
blocked_target_ids: &[String],
|
||||
) -> bool {
|
||||
use moli_protocol::ProtocolNavigationDependency;
|
||||
let pending_load = self.protocol_targets_have_inflight_background_navigation(target_ids);
|
||||
let dependency = residence.navigation_dependency();
|
||||
if dependency == ProtocolNavigationDependency::ReplacesPendingLoad
|
||||
&& pending_load
|
||||
&& !target_ids.is_empty()
|
||||
{
|
||||
return false;
|
||||
}
|
||||
target_ids
|
||||
.iter()
|
||||
.any(|target| blocked_target_ids.contains(target))
|
||||
|| (dependency == ProtocolNavigationDependency::AfterLoad && pending_load)
|
||||
}
|
||||
|
||||
fn next_ungated_protocol_residence_index(&self) -> Option<usize> {
|
||||
// A target-local navigation is an ordering barrier only for later
|
||||
// work from the same target. Keep those lanes ordered while allowing
|
||||
@@ -1954,15 +1977,11 @@ impl CdpScheduler {
|
||||
let mut blocked_target_ids = Vec::new();
|
||||
for (index, residence) in self.queues.protocol_residences.iter().enumerate() {
|
||||
let target_ids = self.protocol_residence_navigation_gate_target_ids(residence);
|
||||
if target_ids
|
||||
.iter()
|
||||
.any(|target_id| blocked_target_ids.contains(target_id))
|
||||
{
|
||||
continue;
|
||||
}
|
||||
if !residence.bypasses_inflight_navigation_gate()
|
||||
&& self.protocol_targets_have_inflight_background_navigation(&target_ids)
|
||||
{
|
||||
if self.protocol_residence_waits_for_navigation(
|
||||
residence,
|
||||
&target_ids,
|
||||
&blocked_target_ids,
|
||||
) {
|
||||
if target_ids.is_empty() {
|
||||
return None;
|
||||
}
|
||||
|
||||
@@ -3,7 +3,7 @@ use std::collections::VecDeque;
|
||||
use moli_core::RendererOutputCursor;
|
||||
use moli_protocol::{
|
||||
DeferredMainDocumentLoadObservationId, DeferredMainDocumentLoadPredecessorCandidate,
|
||||
ProtocolSchedulerWork, ProtocolSchedulerWorkKind,
|
||||
ProtocolNavigationDependency, ProtocolSchedulerWork, ProtocolSchedulerWorkKind,
|
||||
};
|
||||
|
||||
use super::ProtocolOutputSequence;
|
||||
@@ -108,11 +108,11 @@ pub(super) enum ClientTurnPredecessor {
|
||||
}
|
||||
|
||||
impl ProtocolSchedulerResidence {
|
||||
pub(super) fn bypasses_inflight_navigation_gate(&self) -> bool {
|
||||
matches!(
|
||||
self,
|
||||
Self::ProtocolWork { work, .. } if work.bypasses_inflight_navigation_gate()
|
||||
)
|
||||
pub(super) fn navigation_dependency(&self) -> ProtocolNavigationDependency {
|
||||
match self {
|
||||
Self::ProtocolWork { work, .. } => work.navigation_dependency(),
|
||||
Self::RendererOutputPublication(_) => ProtocolNavigationDependency::AfterLoad,
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn should_yield_to_client_turn(&self) -> bool {
|
||||
|
||||
@@ -1252,6 +1252,74 @@ fn replacement_navigation_cancels_and_exactly_settles_the_target_owned_request()
|
||||
assert!(!conn.has_inflight_background_navigation());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn replacement_navigation_passes_load_blocked_observations_but_checks_its_owner() {
|
||||
for snapshot_drain in [false, true] {
|
||||
let mut conn = CdpConnection::new();
|
||||
let navigation = arm_background_navigation_request(&mut conn, "LOADER-pending");
|
||||
let target = navigation.target_id().to_owned();
|
||||
let context = conn.default_browser_context_id().to_owned();
|
||||
let mut scheduler = CdpScheduler::new(conn);
|
||||
scheduler.apply_scheduler_events(vec![
|
||||
CdpSchedulerEvent::ProtocolWorkPublished {
|
||||
work: root_frame_stopped_loading_work_for_target(
|
||||
1,
|
||||
vec![None],
|
||||
context.clone(),
|
||||
target.clone(),
|
||||
"FRAME-pending".to_owned(),
|
||||
"LOADER-pending".to_owned(),
|
||||
),
|
||||
},
|
||||
CdpSchedulerEvent::ProtocolWorkPublished {
|
||||
work: moli_protocol::test_support::retired_location_navigation_work_for_target(
|
||||
2,
|
||||
context,
|
||||
target,
|
||||
root_document_lifecycle_identity(PageId::new_for_testing(99), 1),
|
||||
),
|
||||
},
|
||||
]);
|
||||
assert_eq!(
|
||||
scheduler.next_ungated_protocol_residence_index(),
|
||||
Some(1),
|
||||
"replacement must reach its owner even behind a load-blocked observation"
|
||||
);
|
||||
if snapshot_drain {
|
||||
let snapshot = std::mem::take(&mut scheduler.queues.protocol_residences);
|
||||
assert!(
|
||||
scheduler
|
||||
.complete_protocol_residence_snapshot(snapshot)
|
||||
.await
|
||||
.is_empty()
|
||||
);
|
||||
} else {
|
||||
scheduler.satisfy_front_protocol_residence_client_turn_predecessor();
|
||||
assert!(
|
||||
scheduler
|
||||
.complete_next_protocol_residence()
|
||||
.await
|
||||
.is_empty()
|
||||
);
|
||||
}
|
||||
assert_eq!(
|
||||
scheduler.queues.protocol_residences.len(),
|
||||
1,
|
||||
"the held observation remains, while the replacement action is consumed"
|
||||
);
|
||||
assert!(
|
||||
!navigation.is_cancelled(),
|
||||
"an action from a retired owner has no execution authority"
|
||||
);
|
||||
assert!(scheduler.conn.has_inflight_background_navigation());
|
||||
assert!(settle_background_navigation_request(
|
||||
&mut scheduler.conn,
|
||||
&navigation
|
||||
));
|
||||
assert_eq!(scheduler.next_ungated_protocol_residence_index(), Some(0));
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn scheduler_defers_subresource_network_events_until_background_navigation_gate_clears() {
|
||||
let mut conn = CdpConnection::new();
|
||||
|
||||
+11
@@ -577,6 +577,17 @@ async fn websocket_cdp_renderer_navigation_replaces_pending_response_without_cli
|
||||
let replacement_url = format!("http://{fixture_addr}/replacement");
|
||||
cdp_navigate_and_wait_for_load(&mut socket, 5, &session.session_id, &source_url).await;
|
||||
|
||||
// Emit a real Runtime observation before the replacement owner action.
|
||||
// With Runtime disabled the console call below creates no blocked residence.
|
||||
let _ = send_cdp_command(
|
||||
&mut socket,
|
||||
60,
|
||||
"Runtime.enable",
|
||||
Some(&session.session_id),
|
||||
json!({}),
|
||||
)
|
||||
.await;
|
||||
|
||||
// The response to /ready proves the first navigation reached the server.
|
||||
// Both navigations run outside the Runtime command, and no further client
|
||||
// command may be needed to publish or execute the replacement action.
|
||||
|
||||
@@ -32,7 +32,8 @@ pub use runtime_command_barrier::{
|
||||
};
|
||||
pub(crate) use scheduler_work::ReadyProtocolSchedulerWork;
|
||||
pub use scheduler_work::{
|
||||
ProtocolSchedulerWork, ProtocolSchedulerWorkKind, ProtocolWorkPublishSequence,
|
||||
ProtocolNavigationDependency, ProtocolSchedulerWork, ProtocolSchedulerWorkKind,
|
||||
ProtocolWorkPublishSequence,
|
||||
};
|
||||
pub(in crate::domains) use subresource::{
|
||||
PreparedSubresourceContinueAction,
|
||||
|
||||
@@ -172,30 +172,8 @@ async fn project_renderer_output_records_for_owner(
|
||||
) {
|
||||
for record in records {
|
||||
let (renderer_cause, mut item) = record.into_parts();
|
||||
if projection == RendererPublicationProjection::RetiringNetworkAndResponses {
|
||||
match &mut item {
|
||||
RendererOutputItem::Observation(
|
||||
moli_core::RendererProtocolObservation::Network { .. },
|
||||
) => {}
|
||||
RendererOutputItem::Observation(
|
||||
moli_core::RendererProtocolObservation::RuntimeInspector(batch),
|
||||
) => {
|
||||
// A completed command can reach ingress after its Page was
|
||||
// replaced. Let the session's exact call/attachment correlation
|
||||
// authorize that response, without reviving old notifications.
|
||||
batch.messages.retain(|message| {
|
||||
matches!(
|
||||
message,
|
||||
moli_core::page::RendererRuntimeInspectorMessage::Protocol(message)
|
||||
if message.renderer_call_id().is_some()
|
||||
)
|
||||
});
|
||||
if batch.messages.is_empty() {
|
||||
continue;
|
||||
}
|
||||
}
|
||||
_ => continue,
|
||||
}
|
||||
if !projection.admit_record(&mut item) {
|
||||
continue;
|
||||
}
|
||||
match item {
|
||||
RendererOutputItem::OwnerAction(action) => {
|
||||
@@ -433,6 +411,81 @@ mod tests {
|
||||
command_context.take_protocol_events().is_empty(),
|
||||
"a duplicate terminal response must not be delivered twice"
|
||||
);
|
||||
// Reuse a frontend id, then cancel the new pending call as session
|
||||
// disposal does. Neither the old success nor the cancelled call's
|
||||
// late success may deliver a second terminal or consume a future call.
|
||||
let attachment = RendererAgentAttachmentId::allocate();
|
||||
let prepared = conn
|
||||
.try_register_renderer_call_for_session_owner(
|
||||
Some("SID-retiring-output"),
|
||||
901_002,
|
||||
Some(attachment),
|
||||
RendererCommandDescriptor::from_frontend_policy(
|
||||
frontend.json().to_owned(),
|
||||
frontend.renderer_policy(),
|
||||
RendererInspectorResponseDelivery::SessionSink,
|
||||
),
|
||||
)
|
||||
.expect("frontend id can be reused after its terminal");
|
||||
let cancelled_call_id = prepared.correlation().renderer_call_id().get();
|
||||
drop(prepared);
|
||||
project_renderer_output_records_for_owner(
|
||||
&mut conn,
|
||||
&owner,
|
||||
vec![response(retired_attachment, renderer_call_id)],
|
||||
RendererOutputCursor::new_for_test(stream, 4),
|
||||
RendererPublicationProjection::RetiringNetworkAndResponses,
|
||||
&mut barriers,
|
||||
&mut command_context,
|
||||
)
|
||||
.await;
|
||||
assert!(command_context.take_protocol_events().is_empty());
|
||||
assert!(
|
||||
conn.renderer_runtime_command_cause_for_frontend(Some("SID-retiring-output"), 901_002)
|
||||
.is_some(),
|
||||
"the old terminal must not consume a new call reusing the frontend id"
|
||||
);
|
||||
let mut cancellations = Vec::new();
|
||||
let mut claimed = Vec::new();
|
||||
conn.fail_pending_inspector_awaits_for_owner_background_events_into(
|
||||
&mut cancellations,
|
||||
&mut claimed,
|
||||
&owner,
|
||||
"Inspector detached",
|
||||
);
|
||||
let cancellations = cancellations
|
||||
.into_iter()
|
||||
.chain(claimed)
|
||||
.map(crate::conn::BackgroundProtocolEvent::into_protocol_message)
|
||||
.collect::<Vec<_>>();
|
||||
assert_eq!(
|
||||
cancellations.len(),
|
||||
1,
|
||||
"cancellation owes exactly one terminal"
|
||||
);
|
||||
assert_eq!(cancellations[0]["id"], 901_002);
|
||||
assert!(cancellations[0].get("error").is_some());
|
||||
project_renderer_output_records_for_owner(
|
||||
&mut conn,
|
||||
&owner,
|
||||
vec![
|
||||
response(retired_attachment, renderer_call_id),
|
||||
response(attachment, cancelled_call_id),
|
||||
],
|
||||
RendererOutputCursor::new_for_test(stream, 5),
|
||||
RendererPublicationProjection::RetiringNetworkAndResponses,
|
||||
&mut barriers,
|
||||
&mut command_context,
|
||||
)
|
||||
.await;
|
||||
assert!(
|
||||
command_context.take_protocol_events().is_empty(),
|
||||
"late results after cancellation must be discarded"
|
||||
);
|
||||
assert!(
|
||||
conn.renderer_runtime_command_cause_for_frontend(Some("SID-retiring-output"), 901_002)
|
||||
.is_none()
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
@@ -51,6 +51,34 @@ pub(crate) enum RendererPublicationProjection {
|
||||
RetiringNetworkAndResponses,
|
||||
}
|
||||
|
||||
impl RendererPublicationProjection {
|
||||
/// Retirement revokes new document actions and notifications, not the
|
||||
/// obligation to deliver a terminal response. The pending-call registry
|
||||
/// still decides whether its exact attachment/call is owed, and consumes
|
||||
/// that correlation once. Network records retain their own document check.
|
||||
pub(super) fn admit_record(self, item: &mut moli_core::RendererOutputItem) -> bool {
|
||||
use moli_core::{RendererOutputItem, RendererProtocolObservation};
|
||||
if self == Self::CurrentOwner {
|
||||
return true;
|
||||
}
|
||||
match item {
|
||||
RendererOutputItem::Observation(RendererProtocolObservation::Network { .. }) => true,
|
||||
RendererOutputItem::Observation(RendererProtocolObservation::RuntimeInspector(
|
||||
batch,
|
||||
)) => {
|
||||
batch.messages.retain(|message| {
|
||||
matches!(message,
|
||||
moli_core::page::RendererRuntimeInspectorMessage::Protocol(message)
|
||||
if message.renderer_call_id().is_some()
|
||||
)
|
||||
});
|
||||
!batch.messages.is_empty()
|
||||
}
|
||||
_ => false,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl RendererPublicationRoute {
|
||||
fn for_target(
|
||||
browser_context_id: String,
|
||||
|
||||
@@ -53,6 +53,16 @@ pub enum ProtocolSchedulerWorkKind {
|
||||
PageTargetTerminationOwnerAction,
|
||||
}
|
||||
|
||||
/// Relationship to an outstanding load on this work's exact target. This
|
||||
/// controls selection only: owner identity and terminal response correlation
|
||||
/// remain with their existing owners.
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
pub enum ProtocolNavigationDependency {
|
||||
AfterLoad,
|
||||
Independent,
|
||||
ReplacesPendingLoad,
|
||||
}
|
||||
|
||||
/// Durable protocol-owned work with concrete payload, exact route and one
|
||||
/// connection-local publication sequence.
|
||||
///
|
||||
@@ -304,23 +314,19 @@ impl ProtocolSchedulerWork {
|
||||
}
|
||||
}
|
||||
|
||||
/// Reports work whose transition is independent of an in-flight document
|
||||
/// load on the same target.
|
||||
///
|
||||
/// Foreground selection happens when Chromium accepts a user activation,
|
||||
/// not when the selected target finishes loading. This also lets a target
|
||||
/// paused by `waitForDebuggerOnStart` become active before its initial
|
||||
/// navigation is released. The exact target id above still preserves
|
||||
/// target-local ordering against earlier scheduler residences.
|
||||
/// A navigation requested by the current Document must also reach the
|
||||
/// browser while a previous response is pending, so it can supersede that
|
||||
/// request instead of waiting for the request it needs to cancel.
|
||||
pub fn bypasses_inflight_navigation_gate(&self) -> bool {
|
||||
matches!(
|
||||
&self.payload,
|
||||
ProtocolSchedulerWorkPayload::PopupTargetActivationAction(_)
|
||||
| ProtocolSchedulerWorkPayload::TopLevelLocationNavigationOwnerAction(_)
|
||||
)
|
||||
/// A current-document navigation must reach the owner that cancels the
|
||||
/// old request, including when earlier observations are waiting on it.
|
||||
/// Independent activation still preserves ordering with earlier work.
|
||||
pub fn navigation_dependency(&self) -> ProtocolNavigationDependency {
|
||||
match &self.payload {
|
||||
ProtocolSchedulerWorkPayload::TopLevelLocationNavigationOwnerAction(_) => {
|
||||
ProtocolNavigationDependency::ReplacesPendingLoad
|
||||
}
|
||||
ProtocolSchedulerWorkPayload::PopupTargetActivationAction(_) => {
|
||||
ProtocolNavigationDependency::Independent
|
||||
}
|
||||
_ => ProtocolNavigationDependency::AfterLoad,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn is_top_level_location_navigation_owner_action(&self) -> bool {
|
||||
|
||||
@@ -34,9 +34,10 @@ pub use conn::{
|
||||
PendingRuntimeProtocolMessageDispatch,
|
||||
};
|
||||
pub use domains::activity::{
|
||||
ProtocolSchedulerWork, ProtocolSchedulerWorkKind, ProtocolWorkPublishSequence,
|
||||
RuntimeCommandOutputBarrierCompletion, RuntimeCommandOutputBarrierPermit,
|
||||
RuntimeCommandOutputBarrierTerminal, RuntimeCommandOutputBarriers,
|
||||
ProtocolNavigationDependency, ProtocolSchedulerWork, ProtocolSchedulerWorkKind,
|
||||
ProtocolWorkPublishSequence, RuntimeCommandOutputBarrierCompletion,
|
||||
RuntimeCommandOutputBarrierPermit, RuntimeCommandOutputBarrierTerminal,
|
||||
RuntimeCommandOutputBarriers,
|
||||
};
|
||||
pub use domains::page::{
|
||||
BackgroundNavigationCompletion, CompletedPageScreencastCapture,
|
||||
|
||||
@@ -128,3 +128,35 @@ pub fn root_frame_stopped_loading_work_for_target(
|
||||
loader_id,
|
||||
)
|
||||
}
|
||||
|
||||
/// A replacement action with a deliberately retired Page owner. Scheduler
|
||||
/// tests can check selection independently of the later owner authorization.
|
||||
pub fn retired_location_navigation_work_for_target(
|
||||
publish_sequence: u64,
|
||||
browser_context_id: String,
|
||||
target_id: String,
|
||||
source_document: RendererDocumentLifecycleIdentity,
|
||||
) -> ProtocolSchedulerWork {
|
||||
let owner =
|
||||
crate::conn::CommandOwnerScope::for_route(crate::conn::CdpSessionRoute::PageTarget {
|
||||
browser_context_id: browser_context_id.clone(),
|
||||
target_id: target_id.clone(),
|
||||
session_key: moli_page_types::DevToolsSessionKey::Primary,
|
||||
});
|
||||
let page_owner = crate::conn::TargetPageResidenceIdentity::new(
|
||||
browser_context_id,
|
||||
Some(target_id),
|
||||
crate::conn::TargetPageAttachmentId::allocate(),
|
||||
);
|
||||
ProtocolSchedulerWork::top_level_location_navigation_owner_action(
|
||||
crate::ProtocolWorkPublishSequence::new(publish_sequence),
|
||||
crate::conn::TopLevelLocationNavigationOwnerAction::from_prepared(
|
||||
owner,
|
||||
page_owner,
|
||||
moli_core::page::RendererDocumentSourcedTopLevelLocationNavigation::new(
|
||||
source_document,
|
||||
"data:text/html,replacement".to_owned(),
|
||||
),
|
||||
),
|
||||
)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user