andygrove commented on code in PR #2425: URL: https://github.com/apache/datafusion-ballista/pull/2425#discussion_r4112254637
########## chaos-testing/src/rest.rs: ########## @@ -0,0 +1,192 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! Scheduler REST-API polling shared by both cluster backends. +//! +//! [`crate::cluster::TestCluster`] (local processes) and [`crate::k8s`] +//! (`kind` pods) both drive scenarios by polling the scheduler's REST API — +//! the same endpoints (`/api/executors`, `/api/jobs`, `/api/job/{id}/stages`) +//! and the same JSON shape, differing only in the base URL (loopback vs. the +//! port-forward). These free functions take that base URL so the polling logic +//! lives in one place; each backend wraps them in thin methods. + +use serde_json::Value; +use std::sync::OnceLock; +use std::time::{Duration, Instant}; + +/// Per-request timeout. Short so a stalled connection (e.g. a k8s port-forward +/// that dropped) surfaces as a retryable error inside a polling loop rather than +/// hanging the whole wait; harmless for the loopback (process) backend. +const REQUEST_TIMEOUT: Duration = Duration::from_secs(5); + +/// One shared client (and connection pool) for all polling. The tight polling +/// loops below would otherwise build a fresh `Client` per request — ~1200 new +/// connections through the k8s port-forward over a single 60s wait. +fn client() -> &'static reqwest::Client { Review Comment: One thing I noticed here. This client lives in a process-wide `OnceLock`, but each `#[tokio::test]` gets its own runtime. reqwest can hand back a pooled connection whose runtime has already been dropped, and that fails with "dispatch task is gone". Fresh ephemeral ports make it unlikely, and the polling loops retry, so at worst it's a rare flake. Still, keeping a client on each `TestCluster` / `K8sCluster` would sidestep it. Also, if the builder fails, `unwrap_or_else(|_| reqwest::Client::new())` quietly drops the timeout, so an `expect` might be clearer. ########## chaos-testing/src/rest.rs: ########## @@ -0,0 +1,192 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! Scheduler REST-API polling shared by both cluster backends. +//! +//! [`crate::cluster::TestCluster`] (local processes) and [`crate::k8s`] +//! (`kind` pods) both drive scenarios by polling the scheduler's REST API — +//! the same endpoints (`/api/executors`, `/api/jobs`, `/api/job/{id}/stages`) +//! and the same JSON shape, differing only in the base URL (loopback vs. the +//! port-forward). These free functions take that base URL so the polling logic +//! lives in one place; each backend wraps them in thin methods. + +use serde_json::Value; +use std::sync::OnceLock; +use std::time::{Duration, Instant}; + +/// Per-request timeout. Short so a stalled connection (e.g. a k8s port-forward +/// that dropped) surfaces as a retryable error inside a polling loop rather than +/// hanging the whole wait; harmless for the loopback (process) backend. +const REQUEST_TIMEOUT: Duration = Duration::from_secs(5); + +/// One shared client (and connection pool) for all polling. The tight polling +/// loops below would otherwise build a fresh `Client` per request — ~1200 new +/// connections through the k8s port-forward over a single 60s wait. +fn client() -> &'static reqwest::Client { + static CLIENT: OnceLock<reqwest::Client> = OnceLock::new(); + CLIENT.get_or_init(|| { + reqwest::Client::builder() + .timeout(REQUEST_TIMEOUT) + .build() + .unwrap_or_else(|_| reqwest::Client::new()) + }) +} + +async fn get_json(url: String) -> Result<Value, String> { + client() + .get(url) + .send() + .await + .map_err(|e| e.to_string())? + .json() + .await + .map_err(|e| e.to_string()) +} + +/// How many executors the scheduler currently considers registered. +pub(crate) async fn registered_executors(rest_url: &str) -> Result<usize, String> { + let body = get_json(format!("{rest_url}/api/executors")).await?; + Ok(body.as_array().map(|a| a.len()).unwrap_or(0)) +} + +/// The ids of every executor the scheduler currently lists. Lets a scenario +/// prove a *new* executor (a fresh id) replaced a killed one after rescheduling. +/// Only the k8s backend reschedules automatically, so this is k8s-only. +#[cfg(feature = "k8s")] +pub(crate) async fn executor_ids(rest_url: &str) -> Result<Vec<String>, String> { + let body = get_json(format!("{rest_url}/api/executors")).await?; + Ok(body + .as_array() + .into_iter() + .flatten() + .filter_map(|e| e.get("id").and_then(|v| v.as_str()).map(String::from)) + .collect()) +} + +/// The id of the single job the scheduler currently knows about. +/// +/// The harness runs one query at a time, so "the running job" is unambiguous. +pub(crate) async fn running_job_id(rest_url: &str) -> Result<String, String> { Review Comment: Just flagging that this slightly changes behavior for the process backend too. Before, `TestCluster::running_job_id` and `await_any_stage_running` failed as soon as a request errored, and there was no request timeout. Now they retry until the deadline. I think that's fine and probably better, but if the scheduler dies in an `ha.rs` scenario it will show up as a 30 to 60s timeout rather than an immediate error. It might be worth a line in the PR description. ########## chaos-testing/tests/k8s.rs: ########## @@ -27,87 +27,336 @@ //! kind load docker-image ballista-chaos:test //! CHAOS_BACKEND=kind cargo test -p ballista-chaos --features k8s --test k8s -- --test-threads=1 //! ``` +//! +//! These are the scenarios that genuinely need a cluster — real pod lifecycle, +//! rescheduling, and the port-forward/flight-proxy path — rather than the +//! fault-injection scenarios in `ha.rs`, which are backend-agnostic and stay on +//! the fast process harness. Which planner each scenario runs under mirrors +//! `ha.rs`: both AQE settings only where the planner changes the code path. #![cfg(feature = "k8s")] +use std::time::{Duration, Instant}; + use ballista::prelude::{SessionConfigExt, SessionContextExt}; +use ballista_core::config::BALLISTA_ADAPTIVE_PLANNER_ENABLED; use chaos_testing::fixture::Fixture; -use chaos_testing::k8s::K8sCluster; +use chaos_testing::k8s::{K8sCluster, KillMode}; use datafusion::arrow::util::pretty::pretty_format_batches; use datafusion::execution::session_state::SessionStateBuilder; use datafusion::prelude::{SessionConfig, SessionContext}; +use rstest::rstest; -/// The k8s scenarios need a running kind cluster with the chaos images loaded; +/// The k8s scenarios need a running kind cluster with the chaos image loaded; /// they are opt-in via `CHAOS_BACKEND=kind` so a plain `cargo test` skips them. fn kind_backend_selected() -> bool { if std::env::var("CHAOS_BACKEND").as_deref() == Ok("kind") { true } else { eprintln!( "skipping k8s scenario: set CHAOS_BACKEND=kind and provide a kind cluster \ - with the chaos images loaded (see the crate README runbook)" + with the chaos image loaded (see the crate README runbook)" ); false } } -/// The chaos-free baseline query, run on a fresh local DataFusion context. This -/// is the reference the cluster must reproduce exactly. -async fn local_baseline(fixture: &Fixture) -> String { - let ctx = SessionContext::new(); - for stmt in fixture.register_sql() { - ctx.sql(&stmt).await.unwrap().collect().await.unwrap(); - } - let batches = ctx - .sql(Fixture::baseline_query()) - .await - .unwrap() - .collect() +/// One kind cluster plus its fixture and a connected client, wired for a single +/// scenario. The k8s counterpart of `ha.rs`'s `ChaosRun`: it centralises the +/// fixture write, the client connect, and the UDF-after-upgrade registration so +/// each scenario reads as just its fault and its assertions. +struct K8sRun { + cluster: K8sCluster, + fixture: Fixture, + ctx: SessionContext, +} + +impl K8sRun { + /// Deploy a cluster of `executors`, write the fixture into the shared mount, + /// and connect a client with AQE set to `aqe`. Executor pods use the default + /// (realistic) graceful-shutdown window. + async fn start(aqe: bool, executors: usize) -> Self { + Self::start_with_executor_grace( + aqe, + executors, + chaos_testing::k8s::DEFAULT_EXECUTOR_GRACE_SECONDS, + ) .await - .unwrap(); - pretty_format_batches(&batches).unwrap().to_string() + } + + /// Like [`Self::start`], but with an explicit executor pod + /// `terminationGracePeriodSeconds`. Scenario G uses a short grace so the kill + /// is abrupt (the kubelet SIGKILLs the executors before they can drain). + async fn start_with_executor_grace( + aqe: bool, + executors: usize, + executor_grace_seconds: u64, + ) -> Self { + let cluster = + K8sCluster::start_with_executor_grace(executors, executor_grace_seconds) + .await + .expect("kind cluster must start"); + + // Written into the shared mount, so the scheduler and executor pods see it. + let fixture = Fixture::write(cluster.shared_dir()) + .await + .expect("fixture must be written to the shared mount"); + + let config = SessionConfig::new_with_ballista() + .set_bool(BALLISTA_ADAPTIVE_PLANNER_ENABLED, aqe); + let state = SessionStateBuilder::new() + .with_config(config) + .with_default_features() + .build(); + let ctx = SessionContext::remote_with_state(&cluster.scheduler_url(), state) + .await + .expect("client must connect to the scheduler"); + + // Registered *after* `remote_with_state`: `upgrade_for_ballista` rebuilds + // the state with `with_scalar_functions(...)`, which replaces rather than + // merges the scalar-function map and would drop a UDF registered before. + // (Same subtlety documented in `ha.rs`'s `ChaosRun`.) + ctx.register_udf(chaos_testing::udf::chaos_fail_udf().as_ref().clone()); + ctx.register_udf(chaos_testing::udf::chaos_delay_udf().as_ref().clone()); + + for stmt in fixture.register_sql() { + ctx.sql(&stmt).await.unwrap().collect().await.unwrap(); + } + + Self { + cluster, + fixture, + ctx, + } + } + + /// A clone of the session context, for running a query concurrently with a + /// fault (the query is spawned while the main task kills executors). + fn clone_ctx(&self) -> SessionContext { + self.ctx.clone() + } + + /// Run a query on the cluster, returning the formatted result. + async fn sql(&self, sql: &str) -> Result<String, String> { + let df = self.ctx.sql(sql).await.map_err(|e| e.to_string())?; + let batches = df.collect().await.map_err(|e| e.to_string())?; + Ok(pretty_format_batches(&batches) + .map_err(|e| e.to_string())? + .to_string()) + } + + /// The expected answer, computed by plain local DataFusion over the same + /// fixture — the reference the cluster must reproduce exactly. + async fn local_baseline(&self) -> String { + let ctx = SessionContext::new(); + for stmt in self.fixture.register_sql() { + ctx.sql(&stmt).await.unwrap().collect().await.unwrap(); + } + let batches = ctx + .sql(Fixture::baseline_query()) + .await + .unwrap() + .collect() + .await + .unwrap(); + pretty_format_batches(&batches).unwrap().to_string() + } } /// Smoke test: a real query on a real kind cluster returns the same result as /// plain local DataFusion. Exercises the whole path — client → scheduler → pods /// → shuffle → result — with the fixture shared through the `hostPath` mount. +/// +/// Single planner: the wiring this smoke-tests (mount, port-forward, flight +/// proxy, shuffle) does not vary by planner, so there is nothing to gain from +/// running it under both. #[tokio::test] async fn baseline_matches_local_datafusion_on_k8s() { if !kind_backend_selected() { return; } - let cluster = K8sCluster::start(2).await.expect("kind cluster must start"); - - // Written into the shared mount, so the scheduler and executor pods see it. - let fixture = Fixture::write(cluster.shared_dir()) + let run = K8sRun::start(false, 2).await; + let expected = run.local_baseline().await; + let actual = run + .sql(Fixture::baseline_query()) .await - .expect("fixture must be written to the shared mount"); + .expect("cluster must serve the baseline query"); + + assert_eq!( + actual, expected, + "cluster result must match plain local DataFusion" + ); +} - let expected = local_baseline(&fixture).await; +/// Scenario G (#2029) on k8s: every executor is lost mid-query. +/// +/// The kill scales the executor Deployment to 0 (dropping the ReplicaSet's desired +/// count so it will not reschedule — unlike `--cascade=orphan`, which leaves it +/// recreating pods and the job recovers) then force-deletes the pods. k8s always +/// SIGTERMs before SIGKILL; the abruptness comes from the pods' short grace +/// (`ABRUPT_EXECUTOR_GRACE_SECONDS`), which SIGKILLs the executor before it can +/// drain the in-flight query. With no executor left and none rescheduled, the job +/// fails after the bounded no-executors grace rather than hanging. +/// +/// The per-batch `chaos_delay` just has to outlast that grace so a task is still +/// running when the SIGKILL lands (otherwise the executor drains and the query +/// succeeds — how an earlier version, with the default 30s grace, spuriously +/// passed). +/// +/// Both AQE settings: #2029 had an AQE-on-only second hang path (a task-launch +/// failure that removed the last executor without arming the grace timer), so +/// the planner genuinely changes the code path here. +#[rstest] +#[case::aqe_off(false)] +#[case::aqe_on(true)] +#[tokio::test] +async fn killing_every_executor_terminates_the_job_on_k8s(#[case] aqe: bool) { + if !kind_backend_selected() { + return; + } + + // Short executor grace so the kill is abrupt: the kubelet SIGKILLs the + // executors before their graceful path can drain the in-flight query. + let run = K8sRun::start_with_executor_grace( + aqe, + 2, + chaos_testing::k8s::ABRUPT_EXECUTOR_GRACE_SECONDS, + ) + .await; + // Per-batch delay long enough that a task is still running when the SIGKILL + // lands (> the pods' terminationGracePeriodSeconds); see the doc comment. + let sql = Fixture::chaos_query("chaos_delay(f.key >= 0, 10000)"); Review Comment: Small thing. This works as long as the scheduler's no-executor grace stays well under the 120s query timeout below. A quick comment linking the two numbers would help whoever changes one of them later. ########## chaos-testing/tests/k8s.rs: ########## @@ -27,87 +27,336 @@ //! kind load docker-image ballista-chaos:test //! CHAOS_BACKEND=kind cargo test -p ballista-chaos --features k8s --test k8s -- --test-threads=1 //! ``` +//! +//! These are the scenarios that genuinely need a cluster — real pod lifecycle, +//! rescheduling, and the port-forward/flight-proxy path — rather than the +//! fault-injection scenarios in `ha.rs`, which are backend-agnostic and stay on +//! the fast process harness. Which planner each scenario runs under mirrors +//! `ha.rs`: both AQE settings only where the planner changes the code path. #![cfg(feature = "k8s")] +use std::time::{Duration, Instant}; + use ballista::prelude::{SessionConfigExt, SessionContextExt}; +use ballista_core::config::BALLISTA_ADAPTIVE_PLANNER_ENABLED; use chaos_testing::fixture::Fixture; -use chaos_testing::k8s::K8sCluster; +use chaos_testing::k8s::{K8sCluster, KillMode}; use datafusion::arrow::util::pretty::pretty_format_batches; use datafusion::execution::session_state::SessionStateBuilder; use datafusion::prelude::{SessionConfig, SessionContext}; +use rstest::rstest; -/// The k8s scenarios need a running kind cluster with the chaos images loaded; +/// The k8s scenarios need a running kind cluster with the chaos image loaded; /// they are opt-in via `CHAOS_BACKEND=kind` so a plain `cargo test` skips them. fn kind_backend_selected() -> bool { if std::env::var("CHAOS_BACKEND").as_deref() == Ok("kind") { true } else { eprintln!( "skipping k8s scenario: set CHAOS_BACKEND=kind and provide a kind cluster \ - with the chaos images loaded (see the crate README runbook)" + with the chaos image loaded (see the crate README runbook)" ); false } } -/// The chaos-free baseline query, run on a fresh local DataFusion context. This -/// is the reference the cluster must reproduce exactly. -async fn local_baseline(fixture: &Fixture) -> String { - let ctx = SessionContext::new(); - for stmt in fixture.register_sql() { - ctx.sql(&stmt).await.unwrap().collect().await.unwrap(); - } - let batches = ctx - .sql(Fixture::baseline_query()) - .await - .unwrap() - .collect() +/// One kind cluster plus its fixture and a connected client, wired for a single +/// scenario. The k8s counterpart of `ha.rs`'s `ChaosRun`: it centralises the +/// fixture write, the client connect, and the UDF-after-upgrade registration so +/// each scenario reads as just its fault and its assertions. +struct K8sRun { + cluster: K8sCluster, + fixture: Fixture, + ctx: SessionContext, +} + +impl K8sRun { + /// Deploy a cluster of `executors`, write the fixture into the shared mount, + /// and connect a client with AQE set to `aqe`. Executor pods use the default + /// (realistic) graceful-shutdown window. + async fn start(aqe: bool, executors: usize) -> Self { + Self::start_with_executor_grace( + aqe, + executors, + chaos_testing::k8s::DEFAULT_EXECUTOR_GRACE_SECONDS, + ) .await - .unwrap(); - pretty_format_batches(&batches).unwrap().to_string() + } + + /// Like [`Self::start`], but with an explicit executor pod + /// `terminationGracePeriodSeconds`. Scenario G uses a short grace so the kill + /// is abrupt (the kubelet SIGKILLs the executors before they can drain). + async fn start_with_executor_grace( + aqe: bool, + executors: usize, + executor_grace_seconds: u64, + ) -> Self { + let cluster = + K8sCluster::start_with_executor_grace(executors, executor_grace_seconds) + .await + .expect("kind cluster must start"); + + // Written into the shared mount, so the scheduler and executor pods see it. + let fixture = Fixture::write(cluster.shared_dir()) + .await + .expect("fixture must be written to the shared mount"); + + let config = SessionConfig::new_with_ballista() + .set_bool(BALLISTA_ADAPTIVE_PLANNER_ENABLED, aqe); + let state = SessionStateBuilder::new() + .with_config(config) + .with_default_features() + .build(); + let ctx = SessionContext::remote_with_state(&cluster.scheduler_url(), state) + .await + .expect("client must connect to the scheduler"); + + // Registered *after* `remote_with_state`: `upgrade_for_ballista` rebuilds + // the state with `with_scalar_functions(...)`, which replaces rather than + // merges the scalar-function map and would drop a UDF registered before. + // (Same subtlety documented in `ha.rs`'s `ChaosRun`.) + ctx.register_udf(chaos_testing::udf::chaos_fail_udf().as_ref().clone()); + ctx.register_udf(chaos_testing::udf::chaos_delay_udf().as_ref().clone()); + + for stmt in fixture.register_sql() { + ctx.sql(&stmt).await.unwrap().collect().await.unwrap(); + } + + Self { + cluster, + fixture, + ctx, + } + } + + /// A clone of the session context, for running a query concurrently with a + /// fault (the query is spawned while the main task kills executors). + fn clone_ctx(&self) -> SessionContext { + self.ctx.clone() + } + + /// Run a query on the cluster, returning the formatted result. + async fn sql(&self, sql: &str) -> Result<String, String> { + let df = self.ctx.sql(sql).await.map_err(|e| e.to_string())?; + let batches = df.collect().await.map_err(|e| e.to_string())?; + Ok(pretty_format_batches(&batches) + .map_err(|e| e.to_string())? + .to_string()) + } + + /// The expected answer, computed by plain local DataFusion over the same + /// fixture — the reference the cluster must reproduce exactly. + async fn local_baseline(&self) -> String { + let ctx = SessionContext::new(); + for stmt in self.fixture.register_sql() { + ctx.sql(&stmt).await.unwrap().collect().await.unwrap(); + } + let batches = ctx + .sql(Fixture::baseline_query()) + .await + .unwrap() + .collect() + .await + .unwrap(); + pretty_format_batches(&batches).unwrap().to_string() + } } /// Smoke test: a real query on a real kind cluster returns the same result as /// plain local DataFusion. Exercises the whole path — client → scheduler → pods /// → shuffle → result — with the fixture shared through the `hostPath` mount. +/// +/// Single planner: the wiring this smoke-tests (mount, port-forward, flight +/// proxy, shuffle) does not vary by planner, so there is nothing to gain from +/// running it under both. #[tokio::test] async fn baseline_matches_local_datafusion_on_k8s() { if !kind_backend_selected() { return; } - let cluster = K8sCluster::start(2).await.expect("kind cluster must start"); - - // Written into the shared mount, so the scheduler and executor pods see it. - let fixture = Fixture::write(cluster.shared_dir()) + let run = K8sRun::start(false, 2).await; + let expected = run.local_baseline().await; + let actual = run + .sql(Fixture::baseline_query()) .await - .expect("fixture must be written to the shared mount"); + .expect("cluster must serve the baseline query"); + + assert_eq!( + actual, expected, + "cluster result must match plain local DataFusion" + ); +} - let expected = local_baseline(&fixture).await; +/// Scenario G (#2029) on k8s: every executor is lost mid-query. +/// +/// The kill scales the executor Deployment to 0 (dropping the ReplicaSet's desired +/// count so it will not reschedule — unlike `--cascade=orphan`, which leaves it +/// recreating pods and the job recovers) then force-deletes the pods. k8s always +/// SIGTERMs before SIGKILL; the abruptness comes from the pods' short grace +/// (`ABRUPT_EXECUTOR_GRACE_SECONDS`), which SIGKILLs the executor before it can +/// drain the in-flight query. With no executor left and none rescheduled, the job +/// fails after the bounded no-executors grace rather than hanging. +/// +/// The per-batch `chaos_delay` just has to outlast that grace so a task is still +/// running when the SIGKILL lands (otherwise the executor drains and the query +/// succeeds — how an earlier version, with the default 30s grace, spuriously +/// passed). +/// +/// Both AQE settings: #2029 had an AQE-on-only second hang path (a task-launch +/// failure that removed the last executor without arming the grace timer), so +/// the planner genuinely changes the code path here. +#[rstest] +#[case::aqe_off(false)] +#[case::aqe_on(true)] +#[tokio::test] +async fn killing_every_executor_terminates_the_job_on_k8s(#[case] aqe: bool) { + if !kind_backend_selected() { + return; + } + + // Short executor grace so the kill is abrupt: the kubelet SIGKILLs the + // executors before their graceful path can drain the in-flight query. + let run = K8sRun::start_with_executor_grace( + aqe, + 2, + chaos_testing::k8s::ABRUPT_EXECUTOR_GRACE_SECONDS, + ) + .await; + // Per-batch delay long enough that a task is still running when the SIGKILL + // lands (> the pods' terminationGracePeriodSeconds); see the doc comment. + let sql = Fixture::chaos_query("chaos_delay(f.key >= 0, 10000)"); + + let query = tokio::spawn({ + let ctx = run.clone_ctx(); + async move { ctx.sql(&sql).await?.collect().await } + }); + + let job_id = match run.cluster.running_job_id().await { + Ok(id) => id, + Err(e) => { + fail_with_diagnostics(&run.cluster, format!("job must appear: {e}")).await + } + }; + if let Err(e) = run.cluster.await_any_stage_running(&job_id).await { + fail_with_diagnostics( + &run.cluster, + format!( + "the job must start running a task before we remove its executors: {e}" + ), + ) + .await; + } - let config = SessionConfig::new_with_ballista(); - let state = SessionStateBuilder::new() - .with_config(config) - .with_default_features() - .build(); - let ctx = SessionContext::remote_with_state(&cluster.scheduler_url(), state) + // Total loss that stays lost and is not drained: scale to 0 so no replacement + // is scheduled, then force-kill the pods. + run.cluster + .kill_all_executors_hard() .await - .expect("client must connect to the scheduler"); + .expect("force-kill every executor"); - for stmt in fixture.register_sql() { - ctx.sql(&stmt).await.unwrap().collect().await.unwrap(); + // Dump diagnostics on timeout — this is the assertion the scenario exists for, + // and the namespace is dropped immediately after, so a CI hang leaves nothing. + let result = match tokio::time::timeout(Duration::from_secs(120), query).await { + Ok(joined) => joined.expect("query task should not panic"), + Err(_) => { + fail_with_diagnostics( + &run.cluster, + "job must terminate, not hang, after every executor is lost".to_string(), + ) + .await + } + }; + let err = result.expect_err("query must fail once every executor is lost"); + let msg = err.to_string().to_lowercase(); + assert!( + msg.contains("executor"), Review Comment: This check is pretty loose, since almost any Ballista error mentions an executor. The timeout already covers the "doesn't hang" part, so this is just about precision. If the scheduler's no-executors error has a stable phrase, could we match on that instead? ########## chaos-testing/tests/k8s.rs: ########## @@ -27,87 +27,336 @@ //! kind load docker-image ballista-chaos:test //! CHAOS_BACKEND=kind cargo test -p ballista-chaos --features k8s --test k8s -- --test-threads=1 //! ``` +//! +//! These are the scenarios that genuinely need a cluster — real pod lifecycle, +//! rescheduling, and the port-forward/flight-proxy path — rather than the +//! fault-injection scenarios in `ha.rs`, which are backend-agnostic and stay on +//! the fast process harness. Which planner each scenario runs under mirrors +//! `ha.rs`: both AQE settings only where the planner changes the code path. #![cfg(feature = "k8s")] +use std::time::{Duration, Instant}; + use ballista::prelude::{SessionConfigExt, SessionContextExt}; +use ballista_core::config::BALLISTA_ADAPTIVE_PLANNER_ENABLED; use chaos_testing::fixture::Fixture; -use chaos_testing::k8s::K8sCluster; +use chaos_testing::k8s::{K8sCluster, KillMode}; use datafusion::arrow::util::pretty::pretty_format_batches; use datafusion::execution::session_state::SessionStateBuilder; use datafusion::prelude::{SessionConfig, SessionContext}; +use rstest::rstest; -/// The k8s scenarios need a running kind cluster with the chaos images loaded; +/// The k8s scenarios need a running kind cluster with the chaos image loaded; /// they are opt-in via `CHAOS_BACKEND=kind` so a plain `cargo test` skips them. fn kind_backend_selected() -> bool { if std::env::var("CHAOS_BACKEND").as_deref() == Ok("kind") { true } else { eprintln!( "skipping k8s scenario: set CHAOS_BACKEND=kind and provide a kind cluster \ - with the chaos images loaded (see the crate README runbook)" + with the chaos image loaded (see the crate README runbook)" ); false } } -/// The chaos-free baseline query, run on a fresh local DataFusion context. This -/// is the reference the cluster must reproduce exactly. -async fn local_baseline(fixture: &Fixture) -> String { - let ctx = SessionContext::new(); - for stmt in fixture.register_sql() { - ctx.sql(&stmt).await.unwrap().collect().await.unwrap(); - } - let batches = ctx - .sql(Fixture::baseline_query()) - .await - .unwrap() - .collect() +/// One kind cluster plus its fixture and a connected client, wired for a single +/// scenario. The k8s counterpart of `ha.rs`'s `ChaosRun`: it centralises the +/// fixture write, the client connect, and the UDF-after-upgrade registration so +/// each scenario reads as just its fault and its assertions. +struct K8sRun { + cluster: K8sCluster, + fixture: Fixture, + ctx: SessionContext, +} + +impl K8sRun { + /// Deploy a cluster of `executors`, write the fixture into the shared mount, + /// and connect a client with AQE set to `aqe`. Executor pods use the default + /// (realistic) graceful-shutdown window. + async fn start(aqe: bool, executors: usize) -> Self { + Self::start_with_executor_grace( + aqe, + executors, + chaos_testing::k8s::DEFAULT_EXECUTOR_GRACE_SECONDS, + ) .await - .unwrap(); - pretty_format_batches(&batches).unwrap().to_string() + } + + /// Like [`Self::start`], but with an explicit executor pod + /// `terminationGracePeriodSeconds`. Scenario G uses a short grace so the kill + /// is abrupt (the kubelet SIGKILLs the executors before they can drain). + async fn start_with_executor_grace( + aqe: bool, + executors: usize, + executor_grace_seconds: u64, + ) -> Self { + let cluster = + K8sCluster::start_with_executor_grace(executors, executor_grace_seconds) + .await + .expect("kind cluster must start"); + + // Written into the shared mount, so the scheduler and executor pods see it. + let fixture = Fixture::write(cluster.shared_dir()) + .await + .expect("fixture must be written to the shared mount"); + + let config = SessionConfig::new_with_ballista() + .set_bool(BALLISTA_ADAPTIVE_PLANNER_ENABLED, aqe); + let state = SessionStateBuilder::new() + .with_config(config) + .with_default_features() + .build(); + let ctx = SessionContext::remote_with_state(&cluster.scheduler_url(), state) + .await + .expect("client must connect to the scheduler"); + + // Registered *after* `remote_with_state`: `upgrade_for_ballista` rebuilds + // the state with `with_scalar_functions(...)`, which replaces rather than + // merges the scalar-function map and would drop a UDF registered before. + // (Same subtlety documented in `ha.rs`'s `ChaosRun`.) + ctx.register_udf(chaos_testing::udf::chaos_fail_udf().as_ref().clone()); + ctx.register_udf(chaos_testing::udf::chaos_delay_udf().as_ref().clone()); + + for stmt in fixture.register_sql() { + ctx.sql(&stmt).await.unwrap().collect().await.unwrap(); + } + + Self { + cluster, + fixture, + ctx, + } + } + + /// A clone of the session context, for running a query concurrently with a + /// fault (the query is spawned while the main task kills executors). + fn clone_ctx(&self) -> SessionContext { + self.ctx.clone() + } + + /// Run a query on the cluster, returning the formatted result. + async fn sql(&self, sql: &str) -> Result<String, String> { + let df = self.ctx.sql(sql).await.map_err(|e| e.to_string())?; + let batches = df.collect().await.map_err(|e| e.to_string())?; + Ok(pretty_format_batches(&batches) + .map_err(|e| e.to_string())? + .to_string()) + } + + /// The expected answer, computed by plain local DataFusion over the same + /// fixture — the reference the cluster must reproduce exactly. + async fn local_baseline(&self) -> String { + let ctx = SessionContext::new(); + for stmt in self.fixture.register_sql() { + ctx.sql(&stmt).await.unwrap().collect().await.unwrap(); + } + let batches = ctx + .sql(Fixture::baseline_query()) + .await + .unwrap() + .collect() + .await + .unwrap(); + pretty_format_batches(&batches).unwrap().to_string() + } } /// Smoke test: a real query on a real kind cluster returns the same result as /// plain local DataFusion. Exercises the whole path — client → scheduler → pods /// → shuffle → result — with the fixture shared through the `hostPath` mount. +/// +/// Single planner: the wiring this smoke-tests (mount, port-forward, flight +/// proxy, shuffle) does not vary by planner, so there is nothing to gain from +/// running it under both. #[tokio::test] async fn baseline_matches_local_datafusion_on_k8s() { if !kind_backend_selected() { return; } - let cluster = K8sCluster::start(2).await.expect("kind cluster must start"); - - // Written into the shared mount, so the scheduler and executor pods see it. - let fixture = Fixture::write(cluster.shared_dir()) + let run = K8sRun::start(false, 2).await; + let expected = run.local_baseline().await; + let actual = run + .sql(Fixture::baseline_query()) .await - .expect("fixture must be written to the shared mount"); + .expect("cluster must serve the baseline query"); + + assert_eq!( + actual, expected, + "cluster result must match plain local DataFusion" + ); +} - let expected = local_baseline(&fixture).await; +/// Scenario G (#2029) on k8s: every executor is lost mid-query. +/// +/// The kill scales the executor Deployment to 0 (dropping the ReplicaSet's desired +/// count so it will not reschedule — unlike `--cascade=orphan`, which leaves it +/// recreating pods and the job recovers) then force-deletes the pods. k8s always +/// SIGTERMs before SIGKILL; the abruptness comes from the pods' short grace +/// (`ABRUPT_EXECUTOR_GRACE_SECONDS`), which SIGKILLs the executor before it can +/// drain the in-flight query. With no executor left and none rescheduled, the job +/// fails after the bounded no-executors grace rather than hanging. +/// +/// The per-batch `chaos_delay` just has to outlast that grace so a task is still +/// running when the SIGKILL lands (otherwise the executor drains and the query +/// succeeds — how an earlier version, with the default 30s grace, spuriously +/// passed). +/// +/// Both AQE settings: #2029 had an AQE-on-only second hang path (a task-launch +/// failure that removed the last executor without arming the grace timer), so +/// the planner genuinely changes the code path here. +#[rstest] +#[case::aqe_off(false)] +#[case::aqe_on(true)] +#[tokio::test] +async fn killing_every_executor_terminates_the_job_on_k8s(#[case] aqe: bool) { + if !kind_backend_selected() { + return; + } + + // Short executor grace so the kill is abrupt: the kubelet SIGKILLs the + // executors before their graceful path can drain the in-flight query. + let run = K8sRun::start_with_executor_grace( + aqe, + 2, + chaos_testing::k8s::ABRUPT_EXECUTOR_GRACE_SECONDS, + ) + .await; + // Per-batch delay long enough that a task is still running when the SIGKILL + // lands (> the pods' terminationGracePeriodSeconds); see the doc comment. + let sql = Fixture::chaos_query("chaos_delay(f.key >= 0, 10000)"); + + let query = tokio::spawn({ + let ctx = run.clone_ctx(); + async move { ctx.sql(&sql).await?.collect().await } + }); + + let job_id = match run.cluster.running_job_id().await { + Ok(id) => id, + Err(e) => { + fail_with_diagnostics(&run.cluster, format!("job must appear: {e}")).await + } + }; + if let Err(e) = run.cluster.await_any_stage_running(&job_id).await { + fail_with_diagnostics( + &run.cluster, + format!( + "the job must start running a task before we remove its executors: {e}" + ), + ) + .await; + } - let config = SessionConfig::new_with_ballista(); - let state = SessionStateBuilder::new() - .with_config(config) - .with_default_features() - .build(); - let ctx = SessionContext::remote_with_state(&cluster.scheduler_url(), state) + // Total loss that stays lost and is not drained: scale to 0 so no replacement + // is scheduled, then force-kill the pods. + run.cluster + .kill_all_executors_hard() .await - .expect("client must connect to the scheduler"); + .expect("force-kill every executor"); - for stmt in fixture.register_sql() { - ctx.sql(&stmt).await.unwrap().collect().await.unwrap(); + // Dump diagnostics on timeout — this is the assertion the scenario exists for, + // and the namespace is dropped immediately after, so a CI hang leaves nothing. + let result = match tokio::time::timeout(Duration::from_secs(120), query).await { + Ok(joined) => joined.expect("query task should not panic"), + Err(_) => { + fail_with_diagnostics( + &run.cluster, + "job must terminate, not hang, after every executor is lost".to_string(), + ) + .await + } + }; + let err = result.expect_err("query must fail once every executor is lost"); + let msg = err.to_string().to_lowercase(); + assert!( + msg.contains("executor"), + "failure should name the executor loss, got: {err}" + ); +} + +/// Dump cluster diagnostics, then panic. Used at scenario G's failure points so a +/// CI failure leaves the pod state and logs behind before the namespace is torn +/// down on drop. +async fn fail_with_diagnostics(cluster: &K8sCluster, msg: String) -> ! { + cluster.dump_diagnostics().await; + panic!("{msg}"); +} + +/// Scenario F on k8s: an executor pod is killed and the cluster reabsorbs its +/// replacement. +/// +/// Unlike the process harness — where the test spawns a fresh executor itself — +/// deleting a pod lets the Deployment reschedule a replacement automatically, +/// with a brand-new executor id. That is the k8s-unique behaviour this asserts: +/// after a forced pod delete, the scheduler settles back to two executors, one +/// of which is genuinely new (an id not present before), and the cluster still +/// serves queries. +/// +/// Single planner (aqe off): rescheduling and re-registration are +/// planner-independent, and the post-restart query correctness is already +/// covered by the baseline scenario. +#[tokio::test] +async fn restarted_executor_rejoins_and_serves_queries_on_k8s() { + if !kind_backend_selected() { + return; } - let batches = ctx - .sql(Fixture::baseline_query()) + + let run = K8sRun::start(false, 2).await; + let expected = run.local_baseline().await; + + let before = run + .cluster + .executor_ids() .await - .unwrap() - .collect() + .expect("must list executors before the kill"); + assert_eq!(before.len(), 2, "expected two executors to start"); + + run.cluster + .kill_one_executor(KillMode::Forced) Review Comment: Your work on G shows that a force-delete still honors the pod's `terminationGracePeriodSeconds`. If I'm reading that right, F's forced kill with the default 30s grace is really a graceful shutdown, so the executor might deregister cleanly instead of being lost. F still proves the rescheduled pod gets a new id and rejoins, which is great. But the `KillMode::Forced` docs in `k8s.rs` still describe it as an abrupt, SIGKILL-like loss. Would you be up for either running F with `ABRUPT_EXECUTOR_GRACE_SECONDS` too, or updating the `KillMode::Forced` docs and F's comment to match what we now know? -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
