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
76struct 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 plan: Option<Arc<dyn ExecutionPlan>>,
102}
103
104impl QueryAuditState {
105 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
151struct 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 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 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 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 builder.append(SqlInfo::FlightSqlServerName, "Micromegas Flight SQL Server");
232 builder.append(SqlInfo::FlightSqlServerVersion, "1");
233 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#[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 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 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 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 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 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 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 let total_setup_duration = now() - begin_request;
448 audit_state.setup_ms = request_start.elapsed().as_secs_f64() * 1000.0;
449
450 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 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 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}