From ac654b00951dae3895255808b7ce04bb6c9d56ab Mon Sep 17 00:00:00 2001 From: Sergio Esteves Date: Fri, 25 Sep 2026 14:05:48 +0100 Subject: [PATCH] test: reproduce unexecuted task metrics hang --- tests/metrics_collection.rs | 38 +++++++++++++++++++++++++++++++++++++ 1 file changed, 38 insertions(+) diff --git a/tests/metrics_collection.rs b/tests/metrics_collection.rs index c6a421da7..911568eac 100644 --- a/tests/metrics_collection.rs +++ b/tests/metrics_collection.rs @@ -22,6 +22,7 @@ mod tests { }; use futures::TryStreamExt; use std::sync::Arc; + use std::time::Duration; use test_case::test_case; #[test_case(DistributedMetricsFormat::Aggregated ; "aggregated_metrics")] @@ -379,6 +380,43 @@ mod tests { Ok(()) } + /// Regression for #739: an empty build side can leave sampled probe tasks unexecuted. + #[tokio::test] + #[ignore = "metrics rewrite hangs on planned but unexecuted tasks"] + async fn metrics_rewrite_after_unexecuted_aqe_tasks() -> Result<(), Box> + { + let (mut ctx, _guard, _) = start_localhost_context(3, DefaultSessionBuilder).await; + ctx.set_distributed_dynamic_task_count(true)?; + register_parquet_tables(&ctx).await?; + { + let state = ctx.state_ref(); + let mut state = state.write(); + let options = state.config_mut().options_mut(); + options.optimizer.hash_join_single_partition_threshold = 0; + options.optimizer.hash_join_single_partition_threshold_rows = 0; + } + + let plan = ctx + .sql( + r#"SELECT a."MinTemp" FROM weather a JOIN weather b + ON a."RainToday" = b."RainToday" WHERE a."MinTemp" > 1000000"#, + ) + .await? + .create_physical_plan() + .await?; + let batches = execute_stream(plan.clone(), ctx.task_ctx())? + .try_collect::>() + .await?; + assert_eq!(batches.iter().map(|b| b.num_rows()).sum::(), 0); + + tokio::time::timeout( + Duration::from_secs(10), + rewrite_distributed_plan_with_metrics(plan, DistributedMetricsFormat::PerTask), + ) + .await??; + Ok(()) + } + /// Looks for an [ExecutionPlan] that matches the provided type parameter `T1` in /// the left node and `T2` in the right node and compares its metrics. /// There might be more than one, so `index` determines which one is compared.