andygrove commented on code in PR #2416:
URL: 
https://github.com/apache/datafusion-ballista/pull/2416#discussion_r4112938218


##########
ballista/flight-sql/src/service.rs:
##########
@@ -0,0 +1,852 @@
+// 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.
+
+//! The Arrow Flight SQL frontend itself.
+
+use std::pin::Pin;
+use std::sync::Arc;
+use std::time::Duration;
+
+use arrow::array::RecordBatch;
+use arrow::datatypes::{Schema, SchemaRef};
+use arrow::ipc::writer::IpcWriteOptions;
+use arrow_flight::encode::FlightDataEncoderBuilder;
+use arrow_flight::flight_service_server::FlightService;
+use arrow_flight::sql::server::FlightSqlService;
+use arrow_flight::sql::{
+    ActionCancelQueryRequest, ActionCancelQueryResult,
+    ActionClosePreparedStatementRequest, ActionCreatePreparedStatementRequest,
+    ActionCreatePreparedStatementResult, Any, CommandGetCatalogs, 
CommandGetDbSchemas,
+    CommandGetSqlInfo, CommandGetTableTypes, CommandGetTables, 
CommandGetXdbcTypeInfo,
+    CommandPreparedStatementQuery, CommandStatementQuery, 
CommandStatementUpdate,
+    ProstMessageExt, SqlInfo, TicketStatementQuery,
+    metadata::{SqlInfoData, XdbcTypeInfoData},
+    server::PeekableFlightDataStream,
+};
+use arrow_flight::{
+    Action, FlightData, FlightDescriptor, FlightEndpoint, FlightInfo, 
HandshakeRequest,
+    HandshakeResponse, IpcMessage, SchemaAsIpc, Ticket,
+};
+use ballista_core::error::BallistaError;
+use ballista_core::flight_proxy_service::BallistaFlightProxyService;
+use ballista_core::planner::scans_only_local_tables;
+use ballista_core::serde::protobuf::PartitionLocation;
+use ballista_core::serde::scheduler::{Action as BallistaAction, 
ShuffleFileKind};
+use ballista_core::serde::{decode_protobuf, protobuf};
+use datafusion::logical_expr::{DdlStatement, LogicalPlan};
+use datafusion::prelude::SessionContext;
+use futures::{Stream, TryStreamExt};
+use prost::Message;
+use tonic::metadata::MetadataMap;
+use tonic::{Request, Response, Status, Streaming};
+use uuid::Uuid;
+
+use crate::auth::{AnonymousAuthenticator, Authenticator};
+use crate::backend::QueryBackend;
+use crate::metadata;
+use crate::session::{LocalResult, Prepared, SessionStore};
+use crate::ticket::StatementHandle;
+
+/// Session shared by every client that connects without authenticating.
+///
+/// Anonymous clients cannot be told apart, so they necessarily share catalog
+/// state. Configure an [`Authenticator`] to get a session per connection.
+pub const ANONYMOUS_SESSION: &str = "flight-sql-anonymous";
+
+/// How long a session, prepared statement, or unredeemed local result may sit
+/// idle before it is discarded.
+const DEFAULT_TTL: Duration = Duration::from_secs(30 * 60);
+
+/// How often expired handles are swept.
+const REAP_INTERVAL: Duration = Duration::from_secs(60);
+
+type DoGetStream =
+    Pin<Box<dyn Stream<Item = Result<FlightData, Status>> + Send + 'static>>;
+
+/// Serves Arrow Flight SQL on behalf of a Ballista cluster.
+///
+/// Clients send SQL text; the frontend plans it against the session's catalog,
+/// submits the plan through a [`QueryBackend`], and hands back one
+/// `FlightEndpoint` per output partition. `DoGet` on those tickets is proxied
+/// to the executor holding the partition, so clients never need to reach
+/// executors themselves — the failure mode that made the pre-46.0.0
+/// implementation unusable behind NAT, Docker, and Kubernetes.
+pub struct BallistaFlightSqlService<B: QueryBackend> {
+    backend: Arc<B>,
+    proxy: BallistaFlightProxyService,
+    auth: Arc<dyn Authenticator>,
+    store: Arc<SessionStore>,
+    sql_info: SqlInfoData,
+    xdbc_info: XdbcTypeInfoData,
+}
+
+impl<B: QueryBackend> BallistaFlightSqlService<B> {
+    /// Builds a frontend over `backend`, using `proxy` to stream partition
+    /// data back from executors.
+    ///
+    /// The service authenticates nobody until an [`Authenticator`] is supplied
+    /// via [`with_authenticator`](Self::with_authenticator).
+    pub fn new(backend: Arc<B>, proxy: BallistaFlightProxyService) -> Self {
+        let store = Arc::new(SessionStore::new(DEFAULT_TTL));
+
+        let backend_for_reaper = backend.clone();
+        store.spawn_reaper(REAP_INTERVAL, move |session_id| {
+            let backend = backend_for_reaper.clone();
+            async move {
+                if let Err(e) = backend.close_session(&session_id).await {
+                    log::warn!("flight-sql: failed to close session 
{session_id}: {e}");
+                }
+            }
+        });
+
+        Self {
+            backend,
+            proxy,
+            auth: Arc::new(AnonymousAuthenticator),
+            store,
+            sql_info: metadata::sql_info(),
+            xdbc_info: metadata::xdbc_type_info(),
+        }
+    }
+
+    /// Installs an authenticator. Without one, every handshake is accepted and
+    /// unauthenticated clients share a single session.
+    pub fn with_authenticator(mut self, auth: Arc<dyn Authenticator>) -> Self {
+        self.auth = auth;
+        self
+    }
+
+    /// True when the service will accept unauthenticated clients, which the
+    /// scheduler logs at startup.
+    pub fn allows_anonymous(&self) -> bool {
+        self.auth.allows_anonymous()
+    }
+
+    /// Resolves the Ballista session for a request from its bearer token.
+    fn session_id(&self, metadata: &MetadataMap) -> Result<String, Status> {
+        match bearer_token(metadata) {
+            Some(token) => self.store.session(&token).ok_or_else(|| {
+                Status::unauthenticated(
+                    "unknown or expired session token; re-run the Flight 
handshake",
+                )
+            }),
+            None if self.auth.allows_anonymous() => 
Ok(ANONYMOUS_SESSION.to_string()),
+            None => Err(Status::unauthenticated(
+                "missing bearer token; authenticate with the Flight handshake 
first",
+            )),
+        }
+    }
+
+    /// Returns the context for `session_id`, building it on first use.
+    ///
+    /// The cache is what makes a session a session: [`QueryBackend::session`]
+    /// builds a fresh `SessionContext` every call, so without it a table
+    /// created by one request would be invisible to the next — and every
+    /// request would pay for a full DataFusion session to be constructed.
+    async fn open_session(
+        &self,
+        session_id: &str,
+    ) -> Result<Arc<SessionContext>, Status> {
+        if let Some(ctx) = self.store.context(session_id) {
+            return Ok(ctx);
+        }
+
+        let ctx = self
+            .backend
+            .session(session_id)
+            .await
+            .map_err(|e| Status::internal(format!("failed to open session: 
{e}")))?;
+
+        Ok(self.store.insert_context(session_id.to_string(), ctx))
+    }
+
+    /// Resolves a request to the context it should be planned against.
+    async fn context(
+        &self,
+        metadata: &MetadataMap,
+    ) -> Result<Arc<SessionContext>, Status> {
+        self.open_session(&self.session_id(metadata)?).await
+    }
+
+    /// Plans `sql` against the session's catalog.
+    async fn plan(ctx: &SessionContext, sql: &str) -> Result<LogicalPlan, 
Status> {
+        ctx.state()
+            .create_logical_plan(sql)
+            .await
+            .map_err(|e| Status::invalid_argument(format!("failed to plan 
query: {e}")))
+    }
+
+    /// Runs a planned statement and describes where to collect its results.
+    async fn flight_info_for(
+        &self,
+        ctx: Arc<SessionContext>,
+        plan: LogicalPlan,
+        descriptor: FlightDescriptor,
+        job_name: &str,
+    ) -> Result<FlightInfo, Status> {
+        if let Disposition::Unsupported(reason) = disposition(&plan) {
+            return Err(Status::unimplemented(reason));
+        }
+
+        if disposition(&plan) == Disposition::RunOnScheduler {
+            let (schema, batches) = execute_locally(&ctx, plan).await?;
+            let handle = Uuid::new_v4().to_string();
+            self.store.insert_result(
+                handle.clone(),
+                LocalResult {
+                    schema: schema.clone(),
+                    batches,
+                },
+            );
+
+            let endpoint = Self::endpoint(StatementHandle::Local(handle));
+            return build_flight_info(&schema, vec![endpoint], descriptor);
+        }
+
+        let result = self
+            .backend
+            .execute(job_name, ctx, plan)
+            .await
+            .map_err(query_failed)?;
+
+        let endpoints = result
+            .partitions
+            .into_iter()
+            .map(|location| partition_handle(location).map(Self::endpoint))
+            .collect::<Result<Vec<_>, _>>()
+            .map_err(|e| Status::internal(format!("invalid partition location: 
{e}")))?;
+
+        log::debug!(
+            "flight-sql: job {} produced {} endpoint(s)",
+            result.job_id,
+            endpoints.len()
+        );
+
+        build_flight_info(&result.schema, endpoints, descriptor)
+    }
+
+    /// Builds the endpoint a client redeems for one slice of the result.
+    ///
+    /// Endpoints carry no location, which Flight defines as "fetch from the
+    /// server that gave you this FlightInfo". That keeps every 
cluster-internal
+    /// address off the wire and lets the frontend work unchanged behind NAT, a
+    /// load balancer, or an ingress.
+    fn endpoint(handle: StatementHandle) -> FlightEndpoint {
+        let ticket = TicketStatementQuery {
+            statement_handle: handle.encode().into(),
+        };
+        FlightEndpoint::new().with_ticket(Ticket {
+            ticket: ticket.as_any().encode_to_vec().into(),
+        })
+    }
+
+    /// Serves a metadata command by round-tripping the command itself as the
+    /// ticket, so `DoGet` lands back on the matching handler.
+    fn metadata_info<C: ProstMessageExt>(
+        command: C,
+        schema: SchemaRef,
+        descriptor: FlightDescriptor,
+    ) -> Result<Response<FlightInfo>, Status> {
+        let endpoint = FlightEndpoint::new().with_ticket(Ticket {
+            ticket: command.as_any().encode_to_vec().into(),
+        });
+        build_flight_info(&schema, vec![endpoint], 
descriptor).map(Response::new)
+    }
+}
+
+/// What the frontend should do with a planned statement.
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+enum Disposition {
+    /// Execute on the scheduler: it mutates session state and produces no data
+    /// worth distributing.
+    RunOnScheduler,
+    /// Submit to the cluster.
+    Distribute,
+    /// Refuse, with an explanation for the client.
+    Unsupported(&'static str),
+}
+
+/// Classifies a plan.
+///
+/// One function rather than two predicates, because the interesting cases are
+/// the ones where the answers overlap: `CREATE TABLE AS SELECT` is DDL, and
+/// DDL runs on the scheduler, so a caller that asked "is this DDL?" before
+/// asking "is this supported?" would silently execute its query on one node.
+/// Returning a single verdict makes that ordering impossible to get wrong.
+fn disposition(plan: &LogicalPlan) -> Disposition {
+    match plan {
+        LogicalPlan::Dml(_) => Disposition::Unsupported(
+            "Ballista Flight SQL does not support INSERT/UPDATE/DELETE; \
+             the distributed write path is not implemented",
+        ),
+        LogicalPlan::Copy(_) => Disposition::Unsupported(
+            "Ballista Flight SQL does not support COPY; \
+             the distributed write path is not implemented",
+        ),
+        LogicalPlan::Ddl(DdlStatement::CreateMemoryTable(_)) => 
Disposition::Unsupported(
+            "Ballista Flight SQL does not support CREATE TABLE AS SELECT, \
+             because it would execute on the scheduler rather than the 
cluster; \
+             use CREATE EXTERNAL TABLE over data the executors can read",
+        ),
+        // Other DDL only edits the session catalog, and `SET`-style statements
+        // only edit session config; neither has anything to distribute.
+        LogicalPlan::Ddl(_) | LogicalPlan::Statement(_) => 
Disposition::RunOnScheduler,
+        // `SHOW ...` and anything else reading `information_schema` describes
+        // the catalog held here. There is nothing to distribute, and the
+        // physical form of those scans cannot be serialized for an executor,
+        // so distributing one would hand the client a job that never runs.
+        _ if scans_only_local_tables(plan) => Disposition::RunOnScheduler,
+        _ => Disposition::Distribute,
+    }
+}
+
+async fn execute_locally(
+    ctx: &SessionContext,
+    plan: LogicalPlan,
+) -> Result<(SchemaRef, Vec<RecordBatch>), Status> {
+    let df = ctx
+        .execute_logical_plan(plan)
+        .await
+        .map_err(|e| Status::internal(format!("failed to execute statement: 
{e}")))?;
+    let planned_schema: SchemaRef = Arc::new(df.schema().as_arrow().clone());
+    let batches = df
+        .collect()
+        .await
+        .map_err(|e| Status::internal(format!("failed to execute statement: 
{e}")))?;
+
+    // Prefer the schema the data actually carries; DataFusion's DDL results
+    // are empty and their frame schema is not always the same object.
+    let schema = batches
+        .first()
+        .map(|batch| batch.schema())
+        .unwrap_or(planned_schema);
+
+    Ok((schema, batches))
+}
+
+/// Turns a shuffle partition into the ticket payload the Flight proxy already
+/// knows how to redeem.
+fn partition_handle(
+    location: PartitionLocation,
+) -> Result<StatementHandle, BallistaError> {
+    let layout = location.layout();
+    let partition_id = location.partition_id.ok_or_else(|| {
+        BallistaError::Internal("partition location has no partition 
id".to_string())
+    })?;
+    let executor = location.executor_meta.ok_or_else(|| {
+        BallistaError::Internal("partition location has no executor 
metadata".to_string())
+    })?;
+
+    let action = BallistaAction::FetchPartition {
+        job_id: partition_id.job_id.into(),
+        stage_id: partition_id.stage_id as usize,
+        partition_id: partition_id.partition_id as usize,
+        host: executor.host,
+        port: executor.port as u16,
+        file_id: location.file_id,
+        layout,
+        // A Flight client wants the whole partition: the data file, in full.
+        file_kind: ShuffleFileKind::Data,
+        byte_ranges: vec![],
+    };
+
+    let encoded: protobuf::Action = action.try_into()?;
+    Ok(StatementHandle::Partition(encoded.encode_to_vec()))
+}
+
+fn build_flight_info(
+    schema: &Schema,
+    endpoints: Vec<FlightEndpoint>,
+    descriptor: FlightDescriptor,
+) -> Result<FlightInfo, Status> {
+    FlightInfo::new()
+        .try_with_schema(schema)
+        .map_err(|e| Status::internal(format!("failed to encode result schema: 
{e}")))
+        .map(|info| info.with_descriptor(descriptor).with_endpoints(endpoints))
+}
+
+/// Streams a single metadata batch back to the client.
+fn one_batch_response(batch: RecordBatch) -> Response<DoGetStream> {
+    batch_response(batch.schema(), vec![batch])
+}
+
+fn batch_response(schema: SchemaRef, batches: Vec<RecordBatch>) -> 
Response<DoGetStream> {
+    let stream = FlightDataEncoderBuilder::new()
+        .with_schema(schema)
+        .build(futures::stream::iter(batches.into_iter().map(Ok)))
+        .map_err(|e| Status::internal(format!("failed to encode results: 
{e}")));
+
+    Response::new(Box::pin(stream) as DoGetStream)
+}
+
+fn bearer_token(metadata: &MetadataMap) -> Option<String> {
+    let value = metadata.get("authorization")?.to_str().ok()?;
+    value
+        .strip_prefix("Bearer ")
+        .or_else(|| value.strip_prefix("bearer "))
+        .map(str::to_string)
+}
+
+fn invalid_ticket(e: impl std::fmt::Display) -> Status {
+    Status::invalid_argument(format!("invalid ticket: {e}"))
+}
+
+fn query_failed(e: BallistaError) -> Status {
+    // A failed job is the client's problem to see, not an opaque 500.
+    Status::internal(format!("query execution failed: {e}"))
+}
+
+fn encode_schema(schema: &Schema) -> Result<Vec<u8>, Status> {
+    let message: IpcMessage = SchemaAsIpc::new(schema, 
&IpcWriteOptions::default())
+        .try_into()
+        .map_err(|e| Status::internal(format!("failed to encode schema: 
{e}")))?;
+    Ok(message.0.to_vec())
+}
+
+#[tonic::async_trait]
+impl<B: QueryBackend> FlightSqlService for BallistaFlightSqlService<B> {
+    type FlightService = BallistaFlightSqlService<B>;
+
+    async fn do_handshake(
+        &self,
+        request: Request<Streaming<HandshakeRequest>>,
+    ) -> Result<
+        Response<Pin<Box<dyn Stream<Item = Result<HandshakeResponse, Status>> 
+ Send>>>,
+        Status,
+    > {
+        let identity = self.auth.authenticate(request.metadata()).await?;
+
+        let token = Uuid::new_v4().to_string();
+        let session_id = format!("flight-sql-{}", Uuid::new_v4());
+
+        // Build the session eagerly so a failure surfaces at handshake time
+        // rather than on the client's first query, and so the first query does
+        // not pay for it.
+        self.open_session(&session_id).await?;
+        self.store.insert_session(token.clone(), session_id.clone());
+
+        log::debug!(
+            "flight-sql: handshake for {:?} bound to session {session_id}",
+            identity.user
+        );
+
+        let result = HandshakeResponse {
+            protocol_version: 0,
+            payload: token.clone().into(),
+        };
+        let stream = futures::stream::once(async move { Ok(result) });
+
+        let mut response = Response::new(Box::pin(stream) as _);
+        response.metadata_mut().insert(
+            "authorization",
+            format!("Bearer {token}")
+                .parse()
+                .map_err(|_| Status::internal("failed to encode session 
token"))?,
+        );
+        Ok(response)
+    }
+
+    /// Redeems tickets that are not Flight SQL commands.
+    ///
+    /// Ballista's own Rust client fetches shuffle output with a
+    /// `ballista.protobuf.Action` ticket. When the Flight SQL frontend is
+    /// mounted it replaces the standalone proxy on the scheduler's port, so it
+    /// has to keep serving those tickets.
+    async fn do_get_fallback(
+        &self,
+        request: Request<Ticket>,
+        _message: Any,
+    ) -> Result<Response<DoGetStream>, Status> {
+        decode_protobuf(&request.get_ref().ticket).map_err(invalid_ticket)?;
+        self.proxy.do_get(request).await
+    }
+
+    async fn get_flight_info_statement(
+        &self,
+        query: CommandStatementQuery,
+        request: Request<FlightDescriptor>,
+    ) -> Result<Response<FlightInfo>, Status> {
+        let ctx = self.context(request.metadata()).await?;
+        let plan = Self::plan(&ctx, &query.query).await?;
+        let descriptor = request.into_inner();
+
+        self.flight_info_for(ctx, plan, descriptor, &job_name(&query.query))
+            .await
+            .map(Response::new)
+    }
+
+    async fn get_flight_info_prepared_statement(
+        &self,
+        query: CommandPreparedStatementQuery,
+        request: Request<FlightDescriptor>,
+    ) -> Result<Response<FlightInfo>, Status> {
+        let handle = prepared_handle(&query.prepared_statement_handle)?;
+        let prepared = self.store.prepared(&handle).ok_or_else(|| {
+            Status::not_found("unknown or expired prepared statement handle")
+        })?;
+
+        let ctx = self.open_session(&prepared.session_id).await?;
+        let descriptor = request.into_inner();
+
+        self.flight_info_for(ctx, prepared.plan, descriptor, "flight-sql 
prepared")
+            .await
+            .map(Response::new)
+    }
+
+    async fn get_flight_info_catalogs(
+        &self,
+        query: CommandGetCatalogs,
+        request: Request<FlightDescriptor>,
+    ) -> Result<Response<FlightInfo>, Status> {
+        let schema = query.into_builder().schema();
+        Self::metadata_info(query, schema, request.into_inner())
+    }
+
+    async fn get_flight_info_schemas(
+        &self,
+        query: CommandGetDbSchemas,
+        request: Request<FlightDescriptor>,
+    ) -> Result<Response<FlightInfo>, Status> {
+        let schema = query.clone().into_builder().schema();
+        Self::metadata_info(query, schema, request.into_inner())
+    }
+
+    async fn get_flight_info_tables(
+        &self,
+        query: CommandGetTables,
+        request: Request<FlightDescriptor>,
+    ) -> Result<Response<FlightInfo>, Status> {
+        let schema = query.clone().into_builder().schema();
+        Self::metadata_info(query, schema, request.into_inner())
+    }
+
+    async fn get_flight_info_table_types(
+        &self,
+        query: CommandGetTableTypes,
+        request: Request<FlightDescriptor>,
+    ) -> Result<Response<FlightInfo>, Status> {
+        let schema = query.into_builder().schema();
+        Self::metadata_info(query, schema, request.into_inner())
+    }
+
+    async fn get_flight_info_sql_info(
+        &self,
+        query: CommandGetSqlInfo,
+        request: Request<FlightDescriptor>,
+    ) -> Result<Response<FlightInfo>, Status> {
+        let schema = query.clone().into_builder(&self.sql_info).schema();
+        Self::metadata_info(query, schema, request.into_inner())
+    }
+
+    async fn get_flight_info_xdbc_type_info(
+        &self,
+        query: CommandGetXdbcTypeInfo,
+        request: Request<FlightDescriptor>,
+    ) -> Result<Response<FlightInfo>, Status> {
+        let schema = query.into_builder(&self.xdbc_info).schema();
+        Self::metadata_info(query, schema, request.into_inner())
+    }
+
+    async fn do_get_statement(

Review Comment:
   Every Flight SQL handler now checks the token, and prepared statements and 
local results are bound to the session that created them.
   
   I didn't sign the partition tickets though. `do_get_fallback` has to serve 
raw `FetchPartition` tickets to native clients with no token, so signed Flight 
SQL tickets wouldn't stop anyone. I added a note in the code, and the docs now 
say the authenticator isolates sessions but doesn't secure the port.
   
   I kept `Identity` for now. Handles are bound to the session rather than the 
identity, which also works for anonymous clients.



##########
ballista/scheduler/tests/flight_sql.rs:
##########
@@ -0,0 +1,339 @@
+// 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.
+
+//! End-to-end coverage for the Arrow Flight SQL frontend against a real
+//! (in-process) cluster: scheduler plus one executor, driven entirely through
+//! the Flight SQL protocol.
+//!
+//! The pre-46.0.0 implementation shipped with no tests at all, and the bugs
+//! that got it removed (#1012, #941, #839, #756) were exactly the ones an
+//! end-to-end test catches: endpoints pointing somewhere the client cannot
+//! reach, and results that fail to decode.
+
+#![cfg(feature = "flight-sql")]
+
+use std::sync::Arc;
+use std::time::Duration;
+
+use arrow_flight::sql::client::FlightSqlServiceClient;
+use arrow_flight::sql::{CommandGetTables, SqlInfo};
+use ballista_core::config::TaskSchedulingPolicy;
+use ballista_core::error::BallistaError;
+use ballista_core::serde::BallistaCodec;
+use ballista_core::serde::protobuf::scheduler_grpc_client::SchedulerGrpcClient;
+use ballista_core::utils::{
+    GrpcClientConfig, create_grpc_client_connection, default_config_producer,
+    default_session_builder,
+};
+use ballista_scheduler::cluster::BallistaCluster;
+use ballista_scheduler::config::SchedulerConfig;
+use ballista_scheduler::metrics::default_metrics_collector;
+use ballista_scheduler::scheduler_process::start_grpc_service_with_listener;
+use ballista_scheduler::scheduler_server::SchedulerServer;
+use datafusion::arrow::array::RecordBatch;
+use datafusion_proto::protobuf::{LogicalPlanNode, PhysicalPlanNode};
+use futures::TryStreamExt;
+use tonic::transport::Channel;
+
+/// Test errors are only ever printed, and the client, transport, and Ballista
+/// error types all implement `Error` — so one boxed type spares every call 
site
+/// a `map_err`.
+type Result<T = ()> = std::result::Result<T, Box<dyn std::error::Error>>;
+
+/// Boots a scheduler serving Flight SQL plus one executor, and returns the
+/// scheduler's URL.
+async fn start_cluster() -> Result<String> {

Review Comment:
   For the push path I went with a unit test instead. 
`launch_tasks_fails_jobs_whose_tasks_cannot_be_prepared` in `state/mod.rs` 
checks that the job is reported failed and its slots are freed, and it fails if 
the fix is reverted. After merging main the fix is also much smaller, since 
#2477 and #2016 already handle rejected jobs and refunds.
   
   The test now uses `with_flight_sql(true)`. I couldn't reuse 
`new_standalone_scheduler_with_builder` because it hardcodes its 
`SchedulerConfig` and only mounts `SchedulerGrpc`, so there's no way to turn 
Flight SQL on through it.



-- 
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]

Reply via email to