From fe572867346fa73e4d036c4da8136e8770e58d9b Mon Sep 17 00:00:00 2001 From: rjhallsted Date: Sat, 26 Sep 2026 12:46:10 -0600 Subject: [PATCH 1/2] Add config flag --- src/distributed_planner/distributed_config.rs | 3 +++ 1 file changed, 3 insertions(+) diff --git a/src/distributed_planner/distributed_config.rs b/src/distributed_planner/distributed_config.rs index 4b8fd845e..ee2a9ff01 100644 --- a/src/distributed_planner/distributed_config.rs +++ b/src/distributed_planner/distributed_config.rs @@ -33,6 +33,9 @@ extensions_options! { /// the distributed plan. This does not control whether dynamic filtering is used during /// query execution. pub collect_dynamic_filters: bool, default = true + /// Control whether dynamic filters are distributed accross network boundaries. This does + /// not control intra-stage filter behavior. + pub remote_dynamic_filters: bool, default = true /// Enable broadcast joins for CollectLeft hash joins. When enabled, the build side of /// a CollectLeft join is broadcast to all consumer tasks. pub broadcast_joins: bool, default = true From 0f63462c7715a9c2ed173021c44170de25ed1f6f Mon Sep 17 00:00:00 2001 From: rjhallsted Date: Sat, 26 Sep 2026 13:44:37 -0600 Subject: [PATCH 2/2] Add gates and tests --- src/coordinator/distributed.rs | 5 +- src/coordinator/prepare_dynamic_plan.rs | 11 +++- src/coordinator/prepare_static_plan.rs | 11 +++- src/coordinator/query_coordinator.rs | 23 +++++--- src/distributed_ext.rs | 39 +++++++++++++ src/distributed_planner/distributed_config.rs | 2 +- src/dynamic_filtering/discovery.rs | 55 ++++++++++++++++++- src/dynamic_filtering/mod.rs | 10 +++- tests/dynamic_filtering/common.rs | 45 +++++++++++---- tests/dynamic_filtering/config.rs | 50 +++++++++++++++++ 10 files changed, 222 insertions(+), 29 deletions(-) diff --git a/src/coordinator/distributed.rs b/src/coordinator/distributed.rs index e2b2508b7..9b29f37e9 100644 --- a/src/coordinator/distributed.rs +++ b/src/coordinator/distributed.rs @@ -4,7 +4,7 @@ use crate::coordinator::prepare_static_plan::prepare_static_plan; use crate::coordinator::query_coordinator::QueryCoordinator; use crate::coordinator::store::{Store, task_keys_for_plan}; use crate::dynamic_filtering::{ - is_dynamic_filtering_enabled, sever_dynamic_filter_relationships_in_plan_for_display, + is_local_dynamic_filtering_enabled, sever_dynamic_filter_relationships_in_plan_for_display, }; use crate::{DistributedConfig, TaskCompletedDynamicFilters, TaskKey, TaskMetrics}; use datafusion::common::internal_datafusion_err; @@ -251,7 +251,8 @@ impl ExecutionPlan for DistributedExec { false => prepare_static_plan(&query_coordinator, &base_plan).await?, }; - let dynamic_filtering_enabled = is_dynamic_filtering_enabled(context.session_config()); + let dynamic_filtering_enabled = + is_local_dynamic_filtering_enabled(context.session_config()); prepared.plan_for_viz = match dynamic_filtering_enabled && collect_dynamic_filters { true => sever_dynamic_filter_relationships_in_plan_for_display( prepared.plan_for_viz, diff --git a/src/coordinator/prepare_dynamic_plan.rs b/src/coordinator/prepare_dynamic_plan.rs index 6b7f95857..9b92804b7 100644 --- a/src/coordinator/prepare_dynamic_plan.rs +++ b/src/coordinator/prepare_dynamic_plan.rs @@ -5,7 +5,9 @@ use crate::distributed_planner::{ InjectNetworkBoundaryContext, NetworkBoundaryBuilderResult, ProducerHead, calculate_cost, inject_network_boundaries, }; -use crate::dynamic_filtering::orphan_dynamic_filter_consumers; +use crate::dynamic_filtering::{ + is_remote_dynamic_filtering_enabled, orphan_dynamic_filter_consumers, +}; use crate::events::TaskCountAnnotation::{Desired, Maximum}; use crate::execution_plans::SamplerExec; use crate::stage::{LocalStage, RemoteStage}; @@ -82,7 +84,12 @@ pub(super) async fn prepare_dynamic_plan( // In order to infer the compute the cost of the stage above this one, here a sampler // is injected to gather runtime statistics. input_stage.plan = ProducerHead::insert_sampler(input_stage.plan)?; - let dynamic_filter_anchors = orphan_dynamic_filter_consumers(&input_stage.plan)?; + let dynamic_filter_anchors = + if is_remote_dynamic_filtering_enabled(query_coordinator.session_config())? { + orphan_dynamic_filter_consumers(&input_stage.plan)? + } else { + vec![] + }; let mut load_info_rxs = Vec::with_capacity(input_stage.tasks); diff --git a/src/coordinator/prepare_static_plan.rs b/src/coordinator/prepare_static_plan.rs index f83838d5c..51b0da064 100644 --- a/src/coordinator/prepare_static_plan.rs +++ b/src/coordinator/prepare_static_plan.rs @@ -1,7 +1,9 @@ use crate::common::TreeNodeExt; use crate::coordinator::distributed::PreparedPlan; use crate::coordinator::query_coordinator::QueryCoordinator; -use crate::dynamic_filtering::orphan_dynamic_filter_consumers; +use crate::dynamic_filtering::{ + is_remote_dynamic_filtering_enabled, orphan_dynamic_filter_consumers, +}; use crate::stage::RemoteStage; use crate::{NetworkBoundaryExt, Stage}; use datafusion::common::tree_node::Transformed; @@ -32,7 +34,12 @@ pub(super) async fn prepare_static_plan( let Stage::Local(stage) = plan.input_stage() else { return exec_err!("Input stage from network boundary was not in Local state"); }; - let dynamic_filter_anchors = orphan_dynamic_filter_consumers(&stage.plan)?; + let dynamic_filter_anchors = + if is_remote_dynamic_filtering_enabled(query_coordinator.session_config())? { + orphan_dynamic_filter_consumers(&stage.plan)? + } else { + vec![] + }; let mut stage_coordinator = query_coordinator.stage_coordinator(stage); let mut futures = Vec::with_capacity(stage.tasks); diff --git a/src/coordinator/query_coordinator.rs b/src/coordinator/query_coordinator.rs index 4889f888f..93e5d943f 100644 --- a/src/coordinator/query_coordinator.rs +++ b/src/coordinator/query_coordinator.rs @@ -5,7 +5,8 @@ use crate::coordinator::DynamicFilterRegistry; use crate::coordinator::Store; use crate::coordinator::latency_metric::LatencyMetric; use crate::dynamic_filtering::{ - dynamic_filter_remote_producer_ids, is_dynamic_filtering_enabled, + dynamic_filter_remote_producer_ids, is_local_dynamic_filtering_enabled, + is_remote_dynamic_filtering_enabled, maybe_roundtrip_plan_to_sever_in_memory_dynamic_filter_relationships, }; use crate::events::{ @@ -176,8 +177,10 @@ impl<'a> StageCoordinator<'a> { task_number: task_i, }; - self.dynamic_filter_registry - .register_task(&plan, task_key)?; + if is_remote_dynamic_filtering_enabled(session_config)? { + self.dynamic_filter_registry + .register_task(&plan, task_key)?; + } let mut headers = get_config_extension_propagation_headers(session_config)?; headers.extend(get_passthrough_headers(session_config)); @@ -433,7 +436,8 @@ impl<'a> StageCoordinator<'a> { let wuf_registry = session_config .get_extension::() .unwrap_or_default(); - let dynamic_filtering_enabled = is_dynamic_filtering_enabled(session_config); + let dynamic_filtering_enabled = is_local_dynamic_filtering_enabled(session_config); + let remote_dynamic_filtering_enabled = is_remote_dynamic_filtering_enabled(session_config)?; let mut work_unit_feed_declarations = vec![]; let d_ctx = DistributedTaskContext { @@ -491,11 +495,12 @@ impl<'a> StageCoordinator<'a> { } else { transformed.data }; - let dynamic_filter_remote_producer_ids = if dynamic_filtering_enabled { - dynamic_filter_remote_producer_ids(&plan)? - } else { - vec![] - }; + let dynamic_filter_remote_producer_ids = + if dynamic_filtering_enabled && remote_dynamic_filtering_enabled { + dynamic_filter_remote_producer_ids(&plan)? + } else { + vec![] + }; Ok(TaskSpecializedPlan { plan, work_unit_feed_declarations, diff --git a/src/distributed_ext.rs b/src/distributed_ext.rs index 287377607..c8ccf0426 100644 --- a/src/distributed_ext.rs +++ b/src/distributed_ext.rs @@ -384,6 +384,17 @@ pub trait DistributedExt: Sized { enabled: bool, ) -> Result<(), DataFusionError>; + /// Enables or disables the use of distributed dynamic filters across network boundaries. This + /// does not affect dynamic filter pushdown intra-stage. + fn with_distributed_dynamic_filters_used(self, enabled: bool) -> Result; + + /// Same as [`DistributedExt::with_distributed_dynamic_filters_used`] but with an in-place + /// mutation. + fn set_distributed_dynamic_filters_used( + &mut self, + enabled: bool, + ) -> Result<(), DataFusionError>; + /// Enables children isolator unions for distributing UNION operations across as many tasks as /// the sum of all the tasks required for each child. /// @@ -841,6 +852,15 @@ impl DistributedExt for SessionConfig { Ok(()) } + fn set_distributed_dynamic_filters_used( + &mut self, + enabled: bool, + ) -> Result<(), DataFusionError> { + let d_cfg = DistributedConfig::from_config_options_mut(self.options_mut())?; + d_cfg.remote_dynamic_filters = enabled; + Ok(()) + } + fn set_distributed_children_isolator_unions( &mut self, enabled: bool, @@ -1008,6 +1028,10 @@ impl DistributedExt for SessionConfig { #[expr($?;Ok(self))] fn with_distributed_dynamic_filter_collection(mut self, enabled: bool) -> Result; + #[call(set_distributed_dynamic_filters_used)] + #[expr($?;Ok(self))] + fn with_distributed_dynamic_filters_used(mut self, enabled: bool) -> Result; + #[call(set_distributed_children_isolator_unions)] #[expr($?;Ok(self))] fn with_distributed_children_isolator_unions(mut self, enabled: bool) -> Result; @@ -1143,6 +1167,11 @@ impl DistributedExt for SessionStateBuilder { #[expr($?;Ok(self))] fn with_distributed_dynamic_filter_collection(mut self, enabled: bool) -> Result; + fn set_distributed_dynamic_filters_used(&mut self, enabled: bool) -> Result<(), DataFusionError>; + #[call(set_distributed_dynamic_filters_used)] + #[expr($?;Ok(self))] + fn with_distributed_dynamic_filters_used(mut self, enabled: bool) -> Result; + fn set_distributed_children_isolator_unions(&mut self, enabled: bool) -> Result<(), DataFusionError>; #[call(set_distributed_children_isolator_unions)] #[expr($?;Ok(self))] @@ -1303,6 +1332,11 @@ impl DistributedExt for SessionState { #[expr($?;Ok(self))] fn with_distributed_dynamic_filter_collection(mut self, enabled: bool) -> Result; + fn set_distributed_dynamic_filters_used(&mut self, enabled: bool) -> Result<(), DataFusionError>; + #[call(set_distributed_dynamic_filters_used)] + #[expr($?;Ok(self))] + fn with_distributed_dynamic_filters_used(mut self, enabled: bool) -> Result; + fn set_distributed_children_isolator_unions(&mut self, enabled: bool) -> Result<(), DataFusionError>; #[call(set_distributed_children_isolator_unions)] #[expr($?;Ok(self))] @@ -1456,6 +1490,11 @@ impl DistributedExt for SessionContext { #[expr($?;Ok(self))] fn with_distributed_dynamic_filter_collection(self, enabled: bool) -> Result; + fn set_distributed_dynamic_filters_used(&mut self, enabled: bool) -> Result<(), DataFusionError>; + #[call(set_distributed_dynamic_filters_used)] + #[expr($?;Ok(self))] + fn with_distributed_dynamic_filters_used(self, enabled: bool) -> Result; + fn set_distributed_children_isolator_unions(&mut self, enabled: bool) -> Result<(), DataFusionError>; #[call(set_distributed_children_isolator_unions)] #[expr($?;Ok(self))] diff --git a/src/distributed_planner/distributed_config.rs b/src/distributed_planner/distributed_config.rs index ee2a9ff01..19836d84e 100644 --- a/src/distributed_planner/distributed_config.rs +++ b/src/distributed_planner/distributed_config.rs @@ -33,7 +33,7 @@ extensions_options! { /// the distributed plan. This does not control whether dynamic filtering is used during /// query execution. pub collect_dynamic_filters: bool, default = true - /// Control whether dynamic filters are distributed accross network boundaries. This does + /// Control whether dynamic filters are distributed across network boundaries. This does /// not control intra-stage filter behavior. pub remote_dynamic_filters: bool, default = true /// Enable broadcast joins for CollectLeft hash joins. When enabled, the build side of diff --git a/src/dynamic_filtering/discovery.rs b/src/dynamic_filtering/discovery.rs index 9b833a01a..ccddd5cf4 100644 --- a/src/dynamic_filtering/discovery.rs +++ b/src/dynamic_filtering/discovery.rs @@ -235,6 +235,7 @@ mod tests { ) other_build ON other_build.key = probe."RainTomorrow" WHERE probe."MinTemp" > 0 "#, + true, ) .await?; assert_snapshot!(display, @r" @@ -285,6 +286,7 @@ mod tests { GROUP BY "RainTomorrow" ) probe ON build.key = probe.key "#, + true, ) .await?; assert_snapshot!(display, @r" @@ -314,11 +316,62 @@ mod tests { Ok(()) } - async fn display_query(sql: &str) -> Result { + /// Disabling remote dynamic filters removes the anchors and remote producers, but the + /// task-local producers and consumers are left intact. + #[tokio::test] + async fn skips_anchors_when_remote_dynamic_filters_are_disabled() -> Result<()> { + let display = display_query( + r#" + SELECT COUNT(*) + FROM ( + SELECT DISTINCT "RainToday" AS key + FROM weather + ) build + JOIN weather probe ON build.key = probe."RainToday" + JOIN ( + SELECT DISTINCT "RainTomorrow" AS key + FROM weather + ) other_build ON other_build.key = probe."RainTomorrow" + WHERE probe."MinTemp" > 0 + "#, + false, + ) + .await?; + assert_snapshot!(display, @r" + Stage 5 + AggregateExec + HashJoinExec producers=[1] + NetworkShuffleExec + AggregateExec + NetworkShuffleExec + Stage 4 + RepartitionExec + AggregateExec + DataSourceExec consumers=[1] + Stage 3 + RepartitionExec + HashJoinExec producers=[2] + NetworkShuffleExec + AggregateExec + NetworkShuffleExec + Stage 2 + RepartitionExec + AggregateExec + DataSourceExec consumers=[2] + Stage 1 + RepartitionExec + FilterExec + DataSourceExec + "); + Ok(()) + } + + async fn display_query(sql: &str, remote_dynamic_filters: bool) -> Result { let captured_plans = CapturePlans::default(); let (ctx, _guard, _) = start_localhost_context(2, DefaultSessionBuilder).await; let ctx = ctx .with_distributed_broadcast_joins(false)? + .with_distributed_dynamic_filters_used(remote_dynamic_filters)? .with_distributed_route_task_handler(captured_plans.clone()); { let state = ctx.state_ref(); diff --git a/src/dynamic_filtering/mod.rs b/src/dynamic_filtering/mod.rs index 3752ba82b..3eadda4d9 100644 --- a/src/dynamic_filtering/mod.rs +++ b/src/dynamic_filtering/mod.rs @@ -1,6 +1,7 @@ mod discovery; mod display; +use crate::DistributedConfig; use crate::codec::roundtrip_pb; use datafusion::common::Result; use datafusion::execution::TaskContext; @@ -12,13 +13,20 @@ pub(crate) use discovery::*; pub use display::rewrite_distributed_plan_with_dynamic_filters; pub(crate) use display::sever_dynamic_filter_relationships_in_plan_for_display; -pub(crate) fn is_dynamic_filtering_enabled(session_config: &SessionConfig) -> bool { +pub(crate) fn is_local_dynamic_filtering_enabled(session_config: &SessionConfig) -> bool { session_config .options() .optimizer .enable_dynamic_filter_pushdown } +pub(crate) fn is_remote_dynamic_filtering_enabled(session_config: &SessionConfig) -> Result { + let remote_enabled = + DistributedConfig::from_session_config(session_config)?.remote_dynamic_filters; + let local_enabled = is_local_dynamic_filtering_enabled(session_config); + Ok(remote_enabled && local_enabled) +} + /// Deepcopies the plan if it contains any dynamic filter producers or consumers. This isolates /// any dynamic filters in this plan from dynamic filters *outside* the plan. Plan nodes /// *within* this plan will share in-memory dynamic filter state with eachother, even after copying. diff --git a/tests/dynamic_filtering/common.rs b/tests/dynamic_filtering/common.rs index e468ecf8d..c0c4a6d5d 100644 --- a/tests/dynamic_filtering/common.rs +++ b/tests/dynamic_filtering/common.rs @@ -24,7 +24,8 @@ pub(crate) struct TestQuery<'a> { broadcast_joins: bool, one_task_per_leaf: bool, collect_dynamic_filters: bool, - expect_dynamic_filter_updates: bool, + remote_dynamic_filters: bool, + expect_dynamic_filter_updates: Option, } impl<'a> TestQuery<'a> { @@ -35,7 +36,8 @@ impl<'a> TestQuery<'a> { broadcast_joins: false, one_task_per_leaf: false, collect_dynamic_filters: true, - expect_dynamic_filter_updates: false, + remote_dynamic_filters: true, + expect_dynamic_filter_updates: None, } } @@ -63,8 +65,21 @@ impl<'a> TestQuery<'a> { self } + /// Disables distributing dynamic filters across network boundaries. + pub(crate) fn without_remote_dynamic_filters(mut self) -> Self { + self.remote_dynamic_filters = false; + self + } + + /// Asserts that the coordinator received at least one dynamic filter update. pub(crate) fn expect_dynamic_filter_updates(mut self) -> Self { - self.expect_dynamic_filter_updates = true; + self.expect_dynamic_filter_updates = Some(true); + self + } + + /// Asserts that the coordinator received no dynamic filter updates. + pub(crate) fn expect_no_dynamic_filter_updates(mut self) -> Self { + self.expect_dynamic_filter_updates = Some(false); self } @@ -72,7 +87,8 @@ impl<'a> TestQuery<'a> { let (ctx, _guard, _) = start_localhost_context(2, DefaultSessionBuilder).await; let mut ctx = ctx .with_distributed_broadcast_joins(self.broadcast_joins)? - .with_distributed_dynamic_filter_collection(self.collect_dynamic_filters)?; + .with_distributed_dynamic_filter_collection(self.collect_dynamic_filters)? + .with_distributed_dynamic_filters_used(self.remote_dynamic_filters)?; if self.one_task_per_leaf { ctx = ctx.with_distributed_desired_task_count_handler(1usize); } @@ -124,7 +140,7 @@ pub(crate) async fn execute_range_partitioned_query( register_range_partitioned_table(&ctx, "dim", "testdata/join/parquet/dim", "d_dkey").await?; register_range_partitioned_table(&ctx, "fact", "testdata/join/parquet/fact", "f_dkey").await?; - execute_query_and_display(&ctx, sql, expected_rows, true, false).await + execute_query_and_display(&ctx, sql, expected_rows, true, None).await } async fn register_range_partitioned_table( @@ -156,7 +172,7 @@ async fn execute_query_and_display( sql: &str, expected_rows: usize, collect_dynamic_filters: bool, - expect_dynamic_filter_updates: bool, + expect_dynamic_filter_updates: Option, ) -> Result { let plan = ctx.sql(sql).await?.create_physical_plan().await?; let task_ctx = ctx.task_ctx(); @@ -176,16 +192,23 @@ async fn execute_query_and_display( ); assert_eq!(display_plan_ascii(plan.as_ref(), false), original_display); - if expect_dynamic_filter_updates { + if let Some(expect_updates) = expect_dynamic_filter_updates { let updates = plan .metrics() .expect("DistributedExec has metrics") .sum(|metric| metric.value().name() == "dynamic_filter_updates_received") .map_or(0, |metric| metric.as_usize()); - assert!( - updates > 0, - "expected dynamic_filter_updates_received > 0, got {updates}" - ); + if expect_updates { + assert!( + updates > 0, + "expected dynamic_filter_updates_received > 0, got {updates}" + ); + } else { + assert_eq!( + updates, 0, + "expected dynamic_filter_updates_received == 0, got {updates}" + ); + } } Ok(display_plan_ascii( diff --git a/tests/dynamic_filtering/config.rs b/tests/dynamic_filtering/config.rs index 0a34337d0..101f609c0 100644 --- a/tests/dynamic_filtering/config.rs +++ b/tests/dynamic_filtering/config.rs @@ -2,6 +2,7 @@ mod tests { use crate::common::TestQuery; use datafusion::common::Result; + use datafusion_distributed::assert_snapshot; #[tokio::test] async fn completed_filter_collection_can_be_disabled() -> Result<()> { @@ -22,4 +23,53 @@ mod tests { .await?; Ok(()) } + + /// Disabling remote dynamic filters leaves the plan and its task-local dynamic filters + /// intact, but no producer updates are streamed to the coordinator. + #[tokio::test] + async fn remote_dynamic_filters_can_be_disabled() -> Result<()> { + let display = TestQuery::new( + r#" + SELECT COUNT(*) + FROM ( + SELECT DISTINCT "RainToday" AS key + FROM weather + ) build + JOIN weather probe ON build.key = probe."RainToday" + "#, + ) + .without_remote_dynamic_filters() + .expect_no_dynamic_filter_updates() + .execute() + .await?; + assert_snapshot!(display, @" + ┌───── DistributedExec + │ ProjectionExec: expr=[count(Int64(1))@0 as count(*)] + │ AggregateExec: mode=Final, gby=[], aggr=[count(Int64(1))] + │ CoalescePartitionsExec + │ [Stage 3] => NetworkCoalesceExec: output_partitions=6, input_tasks=2 + └────────────────────────────────────────────────── + ┌───── Stage 3 ── tasks=2, partitions=3 + │ AggregateExec: mode=Partial, gby=[], aggr=[count(Int64(1))] + │ HashJoinExec: mode=Partitioned, join_type=RightSemi, on=[(key@0, RainToday@0)], projection=[] + │ AggregateExec: mode=FinalPartitioned, gby=[key@0 as key], aggr=[] + │ [Stage 1] => NetworkShuffleExec: output_partitions=3, input_tasks=2 + │ [Stage 2] => NetworkShuffleExec: output_partitions=3, input_tasks=2 + └────────────────────────────────────────────────── + ┌───── Stage 1 ── tasks=2, partitions=6 + │ RepartitionExec: partitioning=Hash([key@0], 6), input_partitions=3 + │ AggregateExec: mode=Partial, gby=[key@0 as key], aggr=[] + │ DistributedLeafExec: + │ t0: DataSourceExec: file_groups={3 groups: [[/testdata/weather/result-000000.parquet:..], [/testdata/weather/result-000000.parquet:.., /testdata/weather/result-000001.parquet:..], [/testdata/weather/result-000002.parquet:..]]}, projection=[RainToday@19 as key], file_type=parquet + │ t1: DataSourceExec: file_groups={3 groups: [[/testdata/weather/result-000000.parquet:..], [/testdata/weather/result-000001.parquet:.., /testdata/weather/result-000002.parquet:..], [/testdata/weather/result-000002.parquet:..]]}, projection=[RainToday@19 as key], file_type=parquet + └────────────────────────────────────────────────── + ┌───── Stage 2 ── tasks=2, partitions=6 + │ RepartitionExec: partitioning=Hash([RainToday@0], 6), input_partitions=3 + │ DistributedLeafExec: + │ t0: DataSourceExec: file_groups={3 groups: [[/testdata/weather/result-000000.parquet:..], [/testdata/weather/result-000000.parquet:.., /testdata/weather/result-000001.parquet:..], [/testdata/weather/result-000002.parquet:..]]}, projection=[RainToday], file_type=parquet, predicate=DynamicFilter [ empty ], dynamic_rg_pruning=eligible + │ t1: DataSourceExec: file_groups={3 groups: [[/testdata/weather/result-000000.parquet:..], [/testdata/weather/result-000001.parquet:.., /testdata/weather/result-000002.parquet:..], [/testdata/weather/result-000002.parquet:..]]}, projection=[RainToday], file_type=parquet, predicate=DynamicFilter [ empty ], dynamic_rg_pruning=eligible + └────────────────────────────────────────────────── + "); + Ok(()) + } }