Skip to main content

micromegas/servers/
flight_sql_service_impl.rs

1use super::query_audit::{QueryAuditRecord, ScanMetrics, aggregate_scan_metrics};
2use super::sqlinfo::{
3    SQL_INFO_DATE_TIME_FUNCTIONS, SQL_INFO_NUMERIC_FUNCTIONS, SQL_INFO_SQL_KEYWORDS,
4    SQL_INFO_STRING_FUNCTIONS, SQL_INFO_SYSTEM_FUNCTIONS,
5};
6use anyhow::Result;
7use arrow_flight::decode::FlightRecordBatchStream;
8use arrow_flight::encode::{DictionaryHandling, FlightDataEncoderBuilder};
9use arrow_flight::error::FlightError;
10use arrow_flight::sql::DoPutPreparedStatementResult;
11use arrow_flight::sql::metadata::{SqlInfoData, SqlInfoDataBuilder};
12use arrow_flight::sql::server::PeekableFlightDataStream;
13use arrow_flight::sql::{
14    ActionBeginSavepointRequest, ActionBeginSavepointResult, ActionBeginTransactionRequest,
15    ActionBeginTransactionResult, ActionCancelQueryRequest, ActionCancelQueryResult,
16    ActionClosePreparedStatementRequest, ActionCreatePreparedStatementRequest,
17    ActionCreatePreparedStatementResult, ActionCreatePreparedSubstraitPlanRequest,
18    ActionEndSavepointRequest, ActionEndTransactionRequest, Any, CommandGetCatalogs,
19    CommandGetCrossReference, CommandGetDbSchemas, CommandGetExportedKeys, CommandGetImportedKeys,
20    CommandGetPrimaryKeys, CommandGetSqlInfo, CommandGetTableTypes, CommandGetTables,
21    CommandGetXdbcTypeInfo, CommandPreparedStatementQuery, CommandPreparedStatementUpdate,
22    CommandStatementIngest, CommandStatementQuery, CommandStatementSubstraitPlan,
23    CommandStatementUpdate, ProstMessageExt, SqlInfo, TicketStatementQuery,
24    server::FlightSqlService,
25};
26use arrow_flight::{
27    Action, FlightDescriptor, FlightEndpoint, FlightInfo, HandshakeRequest, HandshakeResponse,
28    Ticket, flight_service_server::FlightService,
29};
30use core::str;
31use datafusion::arrow::datatypes::Schema;
32use datafusion::arrow::ipc::writer::StreamWriter;
33use datafusion::physical_plan::{ExecutionPlan, execute_stream};
34use futures::StreamExt;
35use futures::{Stream, TryStreamExt};
36use micromegas_analytics::lakehouse::lakehouse_context::LakehouseContext;
37use micromegas_analytics::lakehouse::partition_cache::QueryPartitionProvider;
38use micromegas_analytics::lakehouse::query::make_session_context;
39use micromegas_analytics::lakehouse::session_configurator::SessionConfigurator;
40use micromegas_analytics::lakehouse::view_factory::ViewFactory;
41use micromegas_analytics::replication::bulk_ingest;
42use micromegas_analytics::time::TimeRange;
43use micromegas_auth::user_attribution::{is_admin, validate_and_resolve_user_attribution_grpc};
44use micromegas_tracing::prelude::*;
45use once_cell::sync::Lazy;
46use prost::Message;
47use std::pin::Pin;
48use std::str::FromStr;
49use std::sync::Arc;
50use std::task::{Context, Poll};
51use std::time::Instant;
52use tonic::metadata::MetadataMap;
53use tonic::{Request, Response, Status, Streaming};
54
55type FlightDataStream =
56    Pin<Box<dyn Stream<Item = Result<arrow_flight::FlightData, Status>> + Send>>;
57
58macro_rules! status {
59    ($desc:expr, $err:expr) => {
60        Status::internal(format!("{}: {} at {}:{}", $desc, $err, file!(), line!()))
61    };
62}
63
64macro_rules! api_entry_not_implemented {
65    () => {{
66        let function_name = micromegas_tracing::__function_name!();
67        error!("not implemented: {function_name}");
68        Err(Status::unimplemented(format!(
69            "{}:{} not implemented: {function_name}",
70            file!(),
71            line!()
72        )))
73    }};
74}
75
76/// Attribution, per-stage timing, and (once created) the physical plan for
77/// one query. Built as soon as attribution is resolved, updated as each
78/// setup stage completes, and emitted exactly once as a [`QueryAuditRecord`]
79/// under the `flightsql_query_audit` log target — either on an early setup
80/// failure, at stream completion/error, or (if the stream is dropped mid-drain)
81/// from `CompletionTrackedStream`'s `Drop` impl.
82struct QueryAuditState {
83    client: String,
84    user: String,
85    email: String,
86    name: Option<String>,
87    service_account: bool,
88    service_account_name: Option<String>,
89    sql: String,
90    range_begin: Option<String>,
91    range_end: Option<String>,
92    limit: Option<u64>,
93    context_init_ms: f64,
94    planning_ms: f64,
95    execution_ms: f64,
96    setup_ms: f64,
97    request_start: Instant,
98    /// `None` until the physical plan is created; set as soon as it is, and
99    /// still `None` for a record emitted on a setup failure that happens
100    /// before that point.
101    plan: Option<Arc<dyn ExecutionPlan>>,
102}
103
104impl QueryAuditState {
105    /// Aggregate the plan's DataFusion metrics (if a physical plan was
106    /// created before the failure/completion), assemble the audit record,
107    /// and emit it as a single JSON log line.
108    ///
109    /// Takes `&self` (rather than consuming) so it can be called from a
110    /// setup-error `map_err` closure while leaving the state available for
111    /// further updates/completion on the success path, and so `Drop` can
112    /// call it on an abandoned/cancelled stream without needing to
113    /// reconstruct anything.
114    fn emit(&self, status: &'static str, error: Option<String>) {
115        let scan = match &self.plan {
116            Some(plan) => aggregate_scan_metrics(plan.as_ref()),
117            None => ScanMetrics {
118                output_rows: None,
119                bytes_scanned: 0,
120            },
121        };
122        let total_ms = self.request_start.elapsed().as_secs_f64() * 1000.0;
123        let record = QueryAuditRecord {
124            client: self.client.clone(),
125            user: self.user.clone(),
126            email: self.email.clone(),
127            name: self.name.clone(),
128            service_account: self.service_account,
129            service_account_name: self.service_account_name.clone(),
130            sql: self.sql.clone(),
131            range_begin: self.range_begin.clone(),
132            range_end: self.range_end.clone(),
133            limit: self.limit,
134            context_init_ms: self.context_init_ms,
135            planning_ms: self.planning_ms,
136            execution_ms: self.execution_ms,
137            setup_ms: self.setup_ms,
138            total_ms,
139            status,
140            error,
141            output_rows: scan.output_rows,
142            bytes_scanned: scan.bytes_scanned,
143        };
144        match serde_json::to_string(&record) {
145            Ok(json) => info!(target: "flightsql_query_audit", "{json}"),
146            Err(e) => warn!("failed to serialize query audit record: {e}"),
147        }
148    }
149}
150
151/// Stream wrapper that tracks when the stream is fully consumed
152struct CompletionTrackedStream<S> {
153    inner: S,
154    start_time: i64,
155    completed: bool,
156    audit: Option<QueryAuditState>,
157}
158
159impl<S> CompletionTrackedStream<S> {
160    fn new(inner: S, start_time: i64, audit: QueryAuditState) -> Self {
161        Self {
162            inner,
163            start_time,
164            completed: false,
165            audit: Some(audit),
166        }
167    }
168}
169
170impl<S> Drop for CompletionTrackedStream<S> {
171    /// If the stream is dropped before yielding `None` or an `Err` (client
172    /// disconnect/cancel mid-drain), `poll_next` never ran its completion
173    /// arms, so `self.audit` is still `Some(...)`. Emit it here with a
174    /// terminal "incomplete" status so cancelled/abandoned queries are still
175    /// audited instead of silently vanishing. A no-op for streams that
176    /// completed or errored normally, since those arms already `take()` the
177    /// audit state.
178    fn drop(&mut self) {
179        if let Some(state) = self.audit.take() {
180            state.emit("incomplete", None);
181        }
182    }
183}
184
185impl<S> Stream for CompletionTrackedStream<S>
186where
187    S: Stream<Item = Result<arrow_flight::FlightData, Status>> + Unpin + Send,
188{
189    type Item = Result<arrow_flight::FlightData, Status>;
190
191    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
192        match Pin::new(&mut self.inner).poll_next(cx) {
193            Poll::Ready(Some(result)) => {
194                // Check if this is an error result and log it
195                if let Err(ref err) = result {
196                    let sql = self.audit.as_ref().map(|state| state.sql.as_str());
197                    error!("stream error occurred: {err:?} sql={sql:?}");
198                    if !self.completed {
199                        let total_duration = now() - self.start_time;
200                        imetric!("query_duration_with_error", "ticks", total_duration as u64);
201                        imetric!("query_failed", "count", 1);
202                        self.completed = true;
203                        if let Some(state) = self.audit.take() {
204                            state.emit("error", Some(err.to_string()));
205                        }
206                    }
207                }
208                Poll::Ready(Some(result))
209            }
210            Poll::Ready(None) => {
211                // Stream completed successfully
212                if !self.completed {
213                    let total_duration = now() - self.start_time;
214                    imetric!("query_duration_total", "ticks", total_duration as u64);
215                    imetric!("query_completed_successfully", "count", 1);
216                    self.completed = true;
217                    if let Some(state) = self.audit.take() {
218                        state.emit("ok", None);
219                    }
220                }
221                Poll::Ready(None)
222            }
223            Poll::Pending => Poll::Pending,
224        }
225    }
226}
227
228static INSTANCE_SQL_DATA: Lazy<SqlInfoData> = Lazy::new(|| {
229    let mut builder = SqlInfoDataBuilder::new();
230    // Server information
231    builder.append(SqlInfo::FlightSqlServerName, "Micromegas Flight SQL Server");
232    builder.append(SqlInfo::FlightSqlServerVersion, "1");
233    // 1.3 comes from https://github.com/apache/arrow/blob/f9324b79bf4fc1ec7e97b32e3cce16e75ef0f5e3/format/Schema.fbs#L24
234    builder.append(SqlInfo::FlightSqlServerArrowVersion, "1.3");
235    builder.append(SqlInfo::SqlKeywords, SQL_INFO_SQL_KEYWORDS);
236    builder.append(SqlInfo::SqlNumericFunctions, SQL_INFO_NUMERIC_FUNCTIONS);
237    builder.append(SqlInfo::SqlStringFunctions, SQL_INFO_STRING_FUNCTIONS);
238    builder.append(SqlInfo::SqlSystemFunctions, SQL_INFO_SYSTEM_FUNCTIONS);
239    builder.append(SqlInfo::SqlDatetimeFunctions, SQL_INFO_DATE_TIME_FUNCTIONS);
240    builder.build().unwrap()
241});
242
243/// Implementation of the Flight SQL service.
244#[derive(Clone)]
245pub struct FlightSqlServiceImpl {
246    lakehouse: Arc<LakehouseContext>,
247    part_provider: Arc<dyn QueryPartitionProvider>,
248    view_factory: Arc<ViewFactory>,
249    session_configurator: Arc<dyn SessionConfigurator>,
250}
251
252impl FlightSqlServiceImpl {
253    pub fn new(
254        lakehouse: Arc<LakehouseContext>,
255        part_provider: Arc<dyn QueryPartitionProvider>,
256        view_factory: Arc<ViewFactory>,
257        session_configurator: Arc<dyn SessionConfigurator>,
258    ) -> Self {
259        Self {
260            lakehouse,
261            part_provider,
262            view_factory,
263            session_configurator,
264        }
265    }
266
267    fn should_preserve_dictionary(metadata: &MetadataMap) -> bool {
268        metadata
269            .get("preserve_dictionary")
270            .and_then(|v| v.to_str().ok())
271            .map(|s| s.eq_ignore_ascii_case("true"))
272            .unwrap_or(false)
273    }
274
275    #[span_fn]
276    async fn execute_query(
277        &self,
278        ticket_stmt: TicketStatementQuery,
279        metadata: &MetadataMap,
280    ) -> Result<Response<FlightDataStream>, Status> {
281        let begin_request = now();
282        let request_start = Instant::now();
283        let sql = std::str::from_utf8(&ticket_stmt.statement_handle)
284            .map_err(|e| status!("Unable to parse query", e))?;
285
286        let mut begin = metadata.get("query_range_begin");
287        if let Some(s) = &begin
288            && s.is_empty()
289        {
290            begin = None;
291        }
292        let mut end = metadata.get("query_range_end");
293        if let Some(s) = &end
294            && s.is_empty()
295        {
296            end = None;
297        }
298        let query_range = if begin.is_some() && end.is_some() {
299            let begin_datetime = chrono::DateTime::parse_from_rfc3339(
300                begin
301                    .unwrap()
302                    .to_str()
303                    .map_err(|e| status!("Unable to convert query_range_begin to string", e))?,
304            )
305            .map_err(|e| status!("Unable to parse query_range_begin as a rfc3339 datetime", e))?;
306            let end_datetime = chrono::DateTime::parse_from_rfc3339(
307                end.unwrap()
308                    .to_str()
309                    .map_err(|e| status!("Unable to convert query_range_end to string", e))?,
310            )
311            .map_err(|e| status!("Unable to parse query_range_end as a rfc3339 datetime", e))?;
312            Some(TimeRange::new(begin_datetime.into(), end_datetime.into()))
313        } else {
314            None
315        };
316
317        // Validate and resolve user attribution
318        let attr = validate_and_resolve_user_attribution_grpc(metadata).map_err(|e| *e)?;
319
320        let client_type = metadata
321            .get("x-client-type")
322            .and_then(|v| v.to_str().ok())
323            .unwrap_or("unknown");
324
325        let user_name_display = attr.user_name.as_deref().unwrap_or("");
326
327        // Log query with full attribution
328        if let Some(service_account_name) = &attr.service_account {
329            info!(
330                "execute_query range={query_range:?} sql={sql:?} limit={:?} user={} email={} name={user_name_display:?} service_account={service_account_name} client={client_type}",
331                metadata.get("limit"),
332                attr.user_id,
333                attr.user_email
334            );
335        } else {
336            info!(
337                "execute_query range={query_range:?} sql={sql:?} limit={:?} user={} email={} name={user_name_display:?} client={client_type}",
338                metadata.get("limit"),
339                attr.user_id,
340                attr.user_email
341            );
342        }
343
344        // Attribution is resolved from here on, so build the audit state now
345        // (durations/limit/plan filled in as they become known) instead of
346        // only after the physical plan exists. This lets every subsequent
347        // setup failure (session context, planning, limit, physical plan,
348        // stream construction) still emit an "error" audit record instead of
349        // silently disappearing on an early `?` return.
350        let mut audit_state = QueryAuditState {
351            client: client_type.to_string(),
352            user: attr.user_id.clone(),
353            email: attr.user_email.clone(),
354            name: attr.user_name.clone(),
355            service_account: attr.service_account.is_some(),
356            service_account_name: attr.service_account.clone(),
357            sql: sql.to_string(),
358            range_begin: query_range.as_ref().map(|r| r.begin.to_rfc3339()),
359            range_end: query_range.as_ref().map(|r| r.end.to_rfc3339()),
360            limit: None,
361            context_init_ms: 0.0,
362            planning_ms: 0.0,
363            execution_ms: 0.0,
364            setup_ms: 0.0,
365            request_start,
366            plan: None,
367        };
368
369        // Session context creation phase
370        let session_begin = now();
371        let session_begin_instant = Instant::now();
372        let ctx = make_session_context(
373            self.lakehouse.clone(),
374            self.part_provider.clone(),
375            query_range,
376            self.view_factory.clone(),
377            self.session_configurator.clone(),
378            is_admin(metadata),
379        )
380        .await
381        .map_err(|e| {
382            audit_state.emit("error", Some(format!("error in make_session_context: {e}")));
383            status!("error in make_session_context", e)
384        })?;
385        let context_init_duration = now() - session_begin;
386        audit_state.context_init_ms = session_begin_instant.elapsed().as_secs_f64() * 1000.0;
387
388        // Query planning phase
389        let planning_begin = now();
390        let planning_begin_instant = Instant::now();
391        let mut df = ctx.sql(sql).await.map_err(|e| {
392            audit_state.emit("error", Some(format!("error building dataframe: {e}")));
393            status!("error building dataframe", e)
394        })?;
395        let planning_duration = now() - planning_begin;
396        audit_state.planning_ms = planning_begin_instant.elapsed().as_secs_f64() * 1000.0;
397
398        if let Some(limit_str) = metadata.get("limit") {
399            let parsed_limit: usize = usize::from_str(limit_str.to_str().map_err(|e| {
400                audit_state.emit("error", Some(format!("error converting limit to str: {e}")));
401                status!("error converting limit to str", e)
402            })?)
403            .map_err(|e| {
404                audit_state.emit("error", Some(format!("error parsing limit: {e}")));
405                status!("error parsing limit", e)
406            })?;
407            audit_state.limit = Some(parsed_limit as u64);
408            df = df.limit(0, Some(parsed_limit)).map_err(|e| {
409                audit_state.emit(
410                    "error",
411                    Some(format!("error building dataframe with limit: {e}")),
412                );
413                status!("error building dataframe with limit", e)
414            })?;
415        }
416
417        // Query execution phase: build the physical plan (kept for post-drain metrics)
418        // and run it via the free-function `execute_stream`, which is what
419        // `DataFrame::execute_stream` does internally, minus dropping the plan.
420        let execution_begin = now();
421        let execution_begin_instant = Instant::now();
422        let schema = Arc::new(df.schema().as_arrow().clone());
423        let task_ctx = Arc::new(df.task_ctx());
424        let plan = df.create_physical_plan().await.map_err(|e| {
425            audit_state.emit("error", Some(format!("error creating physical plan: {e}")));
426            status!("error creating physical plan", e)
427        })?;
428        audit_state.plan = Some(plan.clone());
429        let stream = execute_stream(plan, task_ctx)
430            .map_err(|e| {
431                audit_state.emit("error", Some(format!("Error executing plan: {e:?}")));
432                Status::internal(format!("Error executing plan: {e:?}"))
433            })?
434            .map_err(|e| FlightError::ExternalError(Box::new(e)));
435        let builder = if Self::should_preserve_dictionary(metadata) {
436            FlightDataEncoderBuilder::new()
437                .with_schema(schema.clone())
438                .with_dictionary_handling(DictionaryHandling::Resend)
439        } else {
440            FlightDataEncoderBuilder::new().with_schema(schema.clone())
441        };
442        let flight_data_stream = builder.build(stream);
443        let execution_duration = now() - execution_begin;
444        audit_state.execution_ms = execution_begin_instant.elapsed().as_secs_f64() * 1000.0;
445
446        // Calculate total setup time and record detailed metrics
447        let total_setup_duration = now() - begin_request;
448        audit_state.setup_ms = request_start.elapsed().as_secs_f64() * 1000.0;
449
450        // Record detailed timing metrics
451        imetric!(
452            "context_init_duration",
453            "ticks",
454            context_init_duration as u64
455        );
456        imetric!("query_planning_duration", "ticks", planning_duration as u64);
457        imetric!(
458            "query_execution_duration",
459            "ticks",
460            execution_duration as u64
461        );
462        imetric!("query_setup_duration", "ticks", total_setup_duration as u64);
463
464        // Create instrumented stream that tracks completion
465        let instrumented_stream =
466            flight_data_stream.map_err(|e| status!("error building data stream", e));
467        let completion_tracked_stream =
468            CompletionTrackedStream::new(instrumented_stream.boxed(), begin_request, audit_state);
469        Ok(Response::new(
470            Box::pin(completion_tracked_stream) as FlightDataStream
471        ))
472    }
473}
474
475#[tonic::async_trait]
476impl FlightSqlService for FlightSqlServiceImpl {
477    type FlightService = FlightSqlServiceImpl;
478
479    async fn do_handshake(
480        &self,
481        _request: Request<Streaming<HandshakeRequest>>,
482    ) -> Result<
483        Response<Pin<Box<dyn Stream<Item = Result<HandshakeResponse, Status>> + Send>>>,
484        Status,
485    > {
486        api_entry_not_implemented!()
487    }
488
489    #[span_fn]
490    async fn do_get_fallback(
491        &self,
492        request: Request<Ticket>,
493        _message: Any,
494    ) -> Result<Response<<Self as FlightService>::DoGetStream>, Status> {
495        let ticket_stmt = TicketStatementQuery::decode(request.get_ref().ticket.clone())
496            .map_err(|e| status!("Could not read ticket", e))?;
497        self.execute_query(ticket_stmt, request.metadata()).await
498    }
499
500    #[span_fn]
501    async fn get_flight_info_statement(
502        &self,
503        query: CommandStatementQuery,
504        _request: Request<FlightDescriptor>,
505    ) -> Result<Response<FlightInfo>, Status> {
506        let begin_request = now();
507        info!("get_flight_info_statement {query:?} ");
508        let CommandStatementQuery { query, .. } = query;
509        let schema = Schema::empty();
510        let ticket = TicketStatementQuery {
511            statement_handle: query.into(),
512        };
513        let mut bytes: Vec<u8> = Vec::new();
514        if ticket.encode(&mut bytes).is_ok() {
515            let info = FlightInfo::new()
516                .try_with_schema(&schema)
517                .unwrap()
518                .with_endpoint(FlightEndpoint::new().with_ticket(Ticket::new(bytes)));
519            let duration = now() - begin_request;
520            imetric!("request_duration", "ticks", duration as u64);
521            Ok(Response::new(info))
522        } else {
523            error!("Error encoding ticket");
524            Err(Status::internal("Error encoding ticket"))
525        }
526    }
527
528    async fn get_flight_info_substrait_plan(
529        &self,
530        _query: CommandStatementSubstraitPlan,
531        _request: Request<FlightDescriptor>,
532    ) -> Result<Response<FlightInfo>, Status> {
533        api_entry_not_implemented!()
534    }
535
536    async fn get_flight_info_prepared_statement(
537        &self,
538        _cmd: CommandPreparedStatementQuery,
539        _request: Request<FlightDescriptor>,
540    ) -> Result<Response<FlightInfo>, Status> {
541        api_entry_not_implemented!()
542    }
543
544    async fn get_flight_info_catalogs(
545        &self,
546        _query: CommandGetCatalogs,
547        _request: Request<FlightDescriptor>,
548    ) -> Result<Response<FlightInfo>, Status> {
549        api_entry_not_implemented!()
550    }
551
552    async fn get_flight_info_schemas(
553        &self,
554        _query: CommandGetDbSchemas,
555        _request: Request<FlightDescriptor>,
556    ) -> Result<Response<FlightInfo>, Status> {
557        api_entry_not_implemented!()
558    }
559
560    #[span_fn]
561    async fn get_flight_info_tables(
562        &self,
563        query: CommandGetTables,
564        request: Request<FlightDescriptor>,
565    ) -> Result<Response<FlightInfo>, Status> {
566        let begin_request = now();
567        info!("get_flight_info_tables");
568        let flight_descriptor = request.into_inner();
569        let ticket = Ticket {
570            ticket: query.as_any().encode_to_vec().into(),
571        };
572        let endpoint = FlightEndpoint::new().with_ticket(ticket);
573        let flight_info = FlightInfo::new()
574            .try_with_schema(&query.into_builder().schema())
575            .map_err(|e| status!("Unable to encode schema", e))?
576            .with_endpoint(endpoint)
577            .with_descriptor(flight_descriptor);
578        let duration = now() - begin_request;
579        imetric!("request_duration", "ticks", duration as u64);
580        Ok(tonic::Response::new(flight_info))
581    }
582
583    async fn get_flight_info_table_types(
584        &self,
585        _query: CommandGetTableTypes,
586        _request: Request<FlightDescriptor>,
587    ) -> Result<Response<FlightInfo>, Status> {
588        api_entry_not_implemented!()
589    }
590
591    #[span_fn]
592    async fn get_flight_info_sql_info(
593        &self,
594        query: CommandGetSqlInfo,
595        request: Request<FlightDescriptor>,
596    ) -> Result<Response<FlightInfo>, Status> {
597        let begin_request = now();
598        info!("get_flight_info_sql_info");
599        let flight_descriptor = request.into_inner();
600        let ticket = Ticket::new(query.as_any().encode_to_vec());
601        let endpoint = FlightEndpoint::new().with_ticket(ticket);
602        let flight_info = FlightInfo::new()
603            .try_with_schema(query.into_builder(&INSTANCE_SQL_DATA).schema().as_ref())
604            .map_err(|e| status!("Unable to encode schema", e))?
605            .with_endpoint(endpoint)
606            .with_descriptor(flight_descriptor);
607        let duration = now() - begin_request;
608        imetric!("request_duration", "ticks", duration as u64);
609        Ok(tonic::Response::new(flight_info))
610    }
611
612    async fn get_flight_info_primary_keys(
613        &self,
614        _query: CommandGetPrimaryKeys,
615        _request: Request<FlightDescriptor>,
616    ) -> Result<Response<FlightInfo>, Status> {
617        api_entry_not_implemented!()
618    }
619
620    async fn get_flight_info_exported_keys(
621        &self,
622        _query: CommandGetExportedKeys,
623        _request: Request<FlightDescriptor>,
624    ) -> Result<Response<FlightInfo>, Status> {
625        api_entry_not_implemented!()
626    }
627
628    async fn get_flight_info_imported_keys(
629        &self,
630        _query: CommandGetImportedKeys,
631        _request: Request<FlightDescriptor>,
632    ) -> Result<Response<FlightInfo>, Status> {
633        api_entry_not_implemented!()
634    }
635
636    async fn get_flight_info_cross_reference(
637        &self,
638        _query: CommandGetCrossReference,
639        _request: Request<FlightDescriptor>,
640    ) -> Result<Response<FlightInfo>, Status> {
641        api_entry_not_implemented!()
642    }
643
644    async fn get_flight_info_xdbc_type_info(
645        &self,
646        _query: CommandGetXdbcTypeInfo,
647        _request: Request<FlightDescriptor>,
648    ) -> Result<Response<FlightInfo>, Status> {
649        api_entry_not_implemented!()
650    }
651
652    #[span_fn]
653    async fn do_get_statement(
654        &self,
655        ticket: TicketStatementQuery,
656        request: Request<Ticket>,
657    ) -> Result<Response<<Self as FlightService>::DoGetStream>, Status> {
658        self.execute_query(ticket, request.metadata()).await
659    }
660
661    async fn do_get_prepared_statement(
662        &self,
663        _query: CommandPreparedStatementQuery,
664        _request: Request<Ticket>,
665    ) -> Result<Response<<Self as FlightService>::DoGetStream>, Status> {
666        api_entry_not_implemented!()
667    }
668
669    async fn do_get_catalogs(
670        &self,
671        _query: CommandGetCatalogs,
672        _request: Request<Ticket>,
673    ) -> Result<Response<<Self as FlightService>::DoGetStream>, Status> {
674        api_entry_not_implemented!()
675    }
676
677    async fn do_get_schemas(
678        &self,
679        _query: CommandGetDbSchemas,
680        _request: Request<Ticket>,
681    ) -> Result<Response<<Self as FlightService>::DoGetStream>, Status> {
682        api_entry_not_implemented!()
683    }
684
685    #[span_fn]
686    async fn do_get_tables(
687        &self,
688        query: CommandGetTables,
689        _request: Request<Ticket>,
690    ) -> Result<Response<<Self as FlightService>::DoGetStream>, Status> {
691        let begin_request = now();
692        info!("do_get_tables {query:?}");
693        let mut builder = query.into_builder();
694        for view in self.view_factory.get_global_views() {
695            let catalog_name = "";
696            let schema_name = "";
697            builder
698                .append(
699                    catalog_name,
700                    schema_name,
701                    &*view.get_view_set_name(),
702                    "table",
703                    &view.get_file_schema(),
704                )
705                .map_err(Status::from)?;
706        }
707        let schema = builder.schema();
708        let batch = builder.build();
709        let stream = FlightDataEncoderBuilder::new()
710            .with_schema(schema)
711            .build(futures::stream::once(async { batch }))
712            .map_err(Status::from);
713        let duration = now() - begin_request;
714        imetric!("request_duration", "ticks", duration as u64);
715        Ok(Response::new(Box::pin(stream)))
716    }
717
718    async fn do_get_table_types(
719        &self,
720        _query: CommandGetTableTypes,
721        _request: Request<Ticket>,
722    ) -> Result<Response<<Self as FlightService>::DoGetStream>, Status> {
723        api_entry_not_implemented!()
724    }
725
726    #[span_fn]
727    async fn do_get_sql_info(
728        &self,
729        query: CommandGetSqlInfo,
730        _request: Request<Ticket>,
731    ) -> Result<Response<<Self as FlightService>::DoGetStream>, Status> {
732        info!("do_get_sql_info");
733        let builder = query.into_builder(&INSTANCE_SQL_DATA);
734        let schema = builder.schema();
735        let batch = builder.build();
736        let stream = FlightDataEncoderBuilder::new()
737            .with_schema(schema)
738            .build(futures::stream::once(async { batch }))
739            .map_err(Status::from);
740        Ok(Response::new(Box::pin(stream)))
741    }
742
743    async fn do_get_primary_keys(
744        &self,
745        _query: CommandGetPrimaryKeys,
746        _request: Request<Ticket>,
747    ) -> Result<Response<<Self as FlightService>::DoGetStream>, Status> {
748        api_entry_not_implemented!()
749    }
750
751    async fn do_get_exported_keys(
752        &self,
753        _query: CommandGetExportedKeys,
754        _request: Request<Ticket>,
755    ) -> Result<Response<<Self as FlightService>::DoGetStream>, Status> {
756        api_entry_not_implemented!()
757    }
758
759    async fn do_get_imported_keys(
760        &self,
761        _query: CommandGetImportedKeys,
762        _request: Request<Ticket>,
763    ) -> Result<Response<<Self as FlightService>::DoGetStream>, Status> {
764        api_entry_not_implemented!()
765    }
766
767    async fn do_get_cross_reference(
768        &self,
769        _query: CommandGetCrossReference,
770        _request: Request<Ticket>,
771    ) -> Result<Response<<Self as FlightService>::DoGetStream>, Status> {
772        api_entry_not_implemented!()
773    }
774
775    async fn do_get_xdbc_type_info(
776        &self,
777        _query: CommandGetXdbcTypeInfo,
778        _request: Request<Ticket>,
779    ) -> Result<Response<<Self as FlightService>::DoGetStream>, Status> {
780        api_entry_not_implemented!()
781    }
782
783    async fn do_put_statement_update(
784        &self,
785        _ticket: CommandStatementUpdate,
786        _request: Request<PeekableFlightDataStream>,
787    ) -> Result<i64, Status> {
788        api_entry_not_implemented!()
789    }
790
791    #[span_fn]
792    async fn do_put_statement_ingest(
793        &self,
794        command: CommandStatementIngest,
795        request: Request<PeekableFlightDataStream>,
796    ) -> Result<i64, Status> {
797        let table_name = command.table;
798        info!("do_put_statement_ingest table_name={table_name}");
799        let stream = FlightRecordBatchStream::new_from_flight_data(
800            request.into_inner().map_err(|e| e.into()),
801        );
802        bulk_ingest(self.lakehouse.lake().clone(), &table_name, stream)
803            .await
804            .map_err(|e| {
805                let msg = format!("error ingesting into {table_name}: {e:?}");
806                error!("{msg}");
807                status!(msg, e)
808            })
809    }
810
811    async fn do_put_substrait_plan(
812        &self,
813        _ticket: CommandStatementSubstraitPlan,
814        _request: Request<PeekableFlightDataStream>,
815    ) -> Result<i64, Status> {
816        api_entry_not_implemented!()
817    }
818
819    async fn do_put_prepared_statement_query(
820        &self,
821        _query: CommandPreparedStatementQuery,
822        _request: Request<PeekableFlightDataStream>,
823    ) -> Result<DoPutPreparedStatementResult, Status> {
824        api_entry_not_implemented!()
825    }
826
827    async fn do_put_prepared_statement_update(
828        &self,
829        _query: CommandPreparedStatementUpdate,
830        _request: Request<PeekableFlightDataStream>,
831    ) -> Result<i64, Status> {
832        api_entry_not_implemented!()
833    }
834
835    #[span_fn]
836    async fn do_action_create_prepared_statement(
837        &self,
838        query: ActionCreatePreparedStatementRequest,
839        request: Request<Action>,
840    ) -> Result<ActionCreatePreparedStatementResult, Status> {
841        info!("do_action_create_prepared_statement query={}", &query.query);
842
843        let ctx = make_session_context(
844            self.lakehouse.clone(),
845            self.part_provider.clone(),
846            None,
847            self.view_factory.clone(),
848            self.session_configurator.clone(),
849            is_admin(request.metadata()),
850        )
851        .await
852        .map_err(|e| status!("error in make_session_context", e))?;
853
854        let df = ctx
855            .sql(&query.query)
856            .await
857            .map_err(|e| status!("error building dataframe", e))?;
858        let schema = df.schema().as_arrow();
859        let mut schema_buffer = Vec::new();
860        let mut writer = StreamWriter::try_new(&mut schema_buffer, schema)
861            .map_err(|e| status!("error writing schema to in-memory buffer", e))?;
862        writer
863            .finish()
864            .map_err(|e| status!("error closing arrow ipc stream writer", e))?;
865        // here we could serialize the logical plan and return that as the prepared statement, but we would
866        // need to register LogicalExtensionCodec for user-defined functions
867        // instead, we are sending back the sql as we received it
868        let result = ActionCreatePreparedStatementResult {
869            prepared_statement_handle: query.query.into(),
870            dataset_schema: schema_buffer.into(),
871            parameter_schema: "".into(),
872        };
873        Ok(result)
874    }
875
876    async fn do_action_close_prepared_statement(
877        &self,
878        _query: ActionClosePreparedStatementRequest,
879        _request: Request<Action>,
880    ) -> Result<(), Status> {
881        info!("do_action_close_prepared_statement");
882        Ok(())
883    }
884
885    async fn do_action_create_prepared_substrait_plan(
886        &self,
887        _query: ActionCreatePreparedSubstraitPlanRequest,
888        _request: Request<Action>,
889    ) -> Result<ActionCreatePreparedStatementResult, Status> {
890        api_entry_not_implemented!()
891    }
892
893    async fn do_action_begin_transaction(
894        &self,
895        _query: ActionBeginTransactionRequest,
896        _request: Request<Action>,
897    ) -> Result<ActionBeginTransactionResult, Status> {
898        api_entry_not_implemented!()
899    }
900
901    async fn do_action_end_transaction(
902        &self,
903        _query: ActionEndTransactionRequest,
904        _request: Request<Action>,
905    ) -> Result<(), Status> {
906        api_entry_not_implemented!()
907    }
908
909    async fn do_action_begin_savepoint(
910        &self,
911        _query: ActionBeginSavepointRequest,
912        _request: Request<Action>,
913    ) -> Result<ActionBeginSavepointResult, Status> {
914        api_entry_not_implemented!()
915    }
916
917    async fn do_action_end_savepoint(
918        &self,
919        _query: ActionEndSavepointRequest,
920        _request: Request<Action>,
921    ) -> Result<(), Status> {
922        api_entry_not_implemented!()
923    }
924
925    async fn do_action_cancel_query(
926        &self,
927        _query: ActionCancelQueryRequest,
928        _request: Request<Action>,
929    ) -> Result<ActionCancelQueryResult, Status> {
930        api_entry_not_implemented!()
931    }
932
933    async fn register_sql_info(&self, _id: i32, _result: &SqlInfo) {
934        info!("register_sql_info");
935    }
936}