1use std::collections::HashSet;
6use std::fs;
7use std::io::{self, Read, Write};
8
9use alopex_embedded::Database;
10
11use crate::batch::{BatchMode, DistributedReadOutcome};
12use crate::cli::{OutputFormat, RoutingReportFormat, SqlCommand, SqlReadMode};
13use crate::client::http::{ClientError, HttpClient};
14use crate::error::{CliError, Result};
15use crate::models::{Column, DataType, Row, Value};
16use crate::output::formatter::{create_formatter, Formatter};
17use crate::profile::config::ResolvedSqlReadMode;
18use crate::streaming::timeout::parse_deadline;
19use crate::streaming::{
20 write_distributed_read_routing_report, CancelSignal, Deadline, DistributedReadRoutingReport,
21 StreamingWriter, WriteStatus,
22};
23use crate::tui::{is_tty, TuiApp};
24use crate::ui::mode::UiMode;
25use futures_util::StreamExt;
26use reqwest::Response;
27
28#[doc(hidden)]
29pub struct SqlExecutionOptions<'a> {
30 pub limit: Option<usize>,
31 pub quiet: bool,
32 pub cancel: &'a CancelSignal,
33 pub deadline: &'a Deadline,
34 pub admin_launcher: Option<Box<dyn FnMut() -> Result<()> + 'a>>,
35}
36
37enum SqlOutput<'a> {
42 Format(OutputFormat),
45 Custom(&'a mut dyn FnMut() -> Box<dyn Formatter>),
48}
49
50impl SqlOutput<'_> {
51 fn create_block_formatter(&mut self) -> Box<dyn Formatter> {
53 match self {
54 SqlOutput::Format(format) => create_formatter(*format),
55 SqlOutput::Custom(factory) => factory(),
56 }
57 }
58
59 fn statement_array(&self) -> bool {
61 matches!(self, SqlOutput::Format(OutputFormat::Json))
62 }
63
64 fn supports_streaming(&mut self) -> bool {
66 match self {
67 SqlOutput::Format(format) => format.supports_streaming(),
68 SqlOutput::Custom(factory) => factory().supports_streaming(),
69 }
70 }
71}
72
73struct CountingWriter<W> {
78 inner: W,
79 bytes: u64,
80}
81
82impl<W: Write> CountingWriter<W> {
83 fn new(inner: W) -> Self {
84 Self { inner, bytes: 0 }
85 }
86
87 fn bytes_written(&self) -> u64 {
88 self.bytes
89 }
90}
91
92impl<W: Write> Write for CountingWriter<W> {
93 fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
94 let written = self.inner.write(buf)?;
95 self.bytes += written as u64;
96 Ok(written)
97 }
98
99 fn flush(&mut self) -> io::Result<()> {
100 self.inner.flush()
101 }
102}
103
104#[allow(clippy::too_many_arguments)]
126pub fn execute_with_formatter<'a, W: Write>(
127 db: &Database,
128 cmd: SqlCommand,
129 batch_mode: &BatchMode,
130 ui_mode: UiMode,
131 writer: &mut W,
132 output_format: OutputFormat,
133 admin_launcher: Option<Box<dyn FnMut() -> Result<()> + 'a>>,
134 limit: Option<usize>,
135 quiet: bool,
136) -> Result<()> {
137 let deadline = Deadline::new(parse_deadline(cmd.deadline.as_deref())?);
138 let cancel = CancelSignal::new();
139
140 execute_with_formatter_control(
141 db,
142 cmd,
143 batch_mode,
144 ui_mode,
145 writer,
146 output_format,
147 SqlExecutionOptions {
148 limit,
149 quiet,
150 cancel: &cancel,
151 deadline: &deadline,
152 admin_launcher,
153 },
154 )
155}
156
157#[allow(clippy::too_many_arguments)]
164pub fn execute_with_formatter_factory<'a, W: Write>(
165 db: &Database,
166 cmd: SqlCommand,
167 batch_mode: &BatchMode,
168 ui_mode: UiMode,
169 writer: &mut W,
170 make_formatter: &mut dyn FnMut() -> Box<dyn Formatter>,
171 admin_launcher: Option<Box<dyn FnMut() -> Result<()> + 'a>>,
172 limit: Option<usize>,
173 quiet: bool,
174) -> Result<()> {
175 let deadline = Deadline::new(parse_deadline(cmd.deadline.as_deref())?);
176 let cancel = CancelSignal::new();
177 let mut output = SqlOutput::Custom(make_formatter);
178
179 execute_with_output_control(
180 db,
181 cmd,
182 batch_mode,
183 ui_mode,
184 writer,
185 &mut output,
186 SqlExecutionOptions {
187 limit,
188 quiet,
189 cancel: &cancel,
190 deadline: &deadline,
191 admin_launcher,
192 },
193 )
194}
195
196#[doc(hidden)]
197pub fn execute_with_formatter_control<W: Write>(
198 db: &Database,
199 cmd: SqlCommand,
200 batch_mode: &BatchMode,
201 ui_mode: UiMode,
202 writer: &mut W,
203 output_format: OutputFormat,
204 options: SqlExecutionOptions<'_>,
205) -> Result<()> {
206 let mut output = SqlOutput::Format(output_format);
207 execute_with_output_control(db, cmd, batch_mode, ui_mode, writer, &mut output, options)
208}
209
210fn execute_with_output_control<W: Write>(
211 db: &Database,
212 cmd: SqlCommand,
213 batch_mode: &BatchMode,
214 ui_mode: UiMode,
215 writer: &mut W,
216 output: &mut SqlOutput<'_>,
217 mut options: SqlExecutionOptions<'_>,
218) -> Result<()> {
219 let sql = cmd.resolve_query(batch_mode)?;
220 let effective_limit = merge_limit(options.limit, cmd.max_rows);
221 options.limit = effective_limit;
222
223 if ui_mode == UiMode::Tui {
224 return execute_tui_local_or_fallback(db, &sql, writer, output, options);
225 }
226
227 execute_sql_with_formatter(db, &sql, writer, output, &options)
228}
229
230fn is_select_query(sql: &str) -> Result<bool> {
231 use alopex_sql::{AlopexDialect, Parser};
232
233 let dialect = AlopexDialect;
234 let stmts = Parser::parse_sql(&dialect, sql).map_err(|e| CliError::Parse(format!("{}", e)))?;
235 Ok(stmts.len() == 1
236 && stmts
237 .first()
238 .is_some_and(|statement| statement.kind.is_query()))
239}
240
241#[allow(clippy::too_many_arguments)]
246pub async fn execute_remote_with_formatter<'a, W: Write>(
247 client: &HttpClient,
248 cmd: &SqlCommand,
249 batch_mode: &BatchMode,
250 ui_mode: UiMode,
251 writer: &mut W,
252 output_format: OutputFormat,
253 admin_launcher: Option<Box<dyn FnMut() -> Result<()> + 'a>>,
254 limit: Option<usize>,
255 quiet: bool,
256) -> Result<()> {
257 let effective_limit = merge_limit(limit, cmd.max_rows);
258 let deadline = Deadline::new(parse_deadline(cmd.deadline.as_deref())?);
259 let cancel = CancelSignal::new();
260 let options = SqlExecutionOptions {
261 limit: effective_limit,
262 quiet,
263 cancel: &cancel,
264 deadline: &deadline,
265 admin_launcher,
266 };
267
268 execute_remote_with_formatter_control(
269 client,
270 cmd,
271 batch_mode,
272 ui_mode,
273 writer,
274 output_format,
275 options,
276 )
277 .await
278}
279
280#[allow(clippy::too_many_arguments)]
287pub async fn execute_remote_with_routing<'a, W: Write>(
288 client: &HttpClient,
289 cmd: &SqlCommand,
290 read_mode: ResolvedSqlReadMode,
291 batch_mode: &BatchMode,
292 ui_mode: UiMode,
293 writer: &mut W,
294 output_format: OutputFormat,
295 admin_launcher: Option<Box<dyn FnMut() -> Result<()> + 'a>>,
296 limit: Option<usize>,
297 quiet: bool,
298) -> Result<()> {
299 let requested_mode = requested_read_mode_name(cmd.read_mode, read_mode);
300 let (effective_mode, decision) = match read_mode {
301 ResolvedSqlReadMode::Local => (Some("local".to_string()), "local"),
302 ResolvedSqlReadMode::Cluster(mode) => (Some(read_mode_name(mode).to_string()), "cluster"),
303 };
304 let result = match read_mode {
305 ResolvedSqlReadMode::Local => {
306 execute_remote_with_formatter(
307 client,
308 cmd,
309 batch_mode,
310 ui_mode,
311 writer,
312 output_format,
313 admin_launcher,
314 limit,
315 quiet,
316 )
317 .await
318 }
319 ResolvedSqlReadMode::Cluster(mode) => {
320 execute_remote_cluster_read(client, cmd, batch_mode, ui_mode, mode).await
321 }
322 };
323 if let Some(format) = cmd.routing_report {
324 emit_routing_report(format, requested_mode, effective_mode, decision, &result)?;
325 }
326 result
327}
328
329async fn execute_remote_cluster_read(
330 client: &HttpClient,
331 cmd: &SqlCommand,
332 batch_mode: &BatchMode,
333 ui_mode: UiMode,
334 mode: SqlReadMode,
335) -> Result<()> {
336 if ui_mode == UiMode::Tui {
337 return Err(distributed_read_error(
338 DistributedReadOutcome::Unsupported,
339 "cluster SQL read does not support TUI preview",
340 ));
341 }
342 let sql = cmd.resolve_query(batch_mode)?;
343 if !is_select_query(&sql)? {
344 return Err(distributed_read_error(
345 DistributedReadOutcome::Unsupported,
346 "remote DDL, DML, and transactions are not supported for cluster reads",
347 ));
348 }
349 let request = RemoteDistributedReadRequest {
350 sql,
351 read_mode: read_mode_name(mode),
352 };
353 let _: serde_json::Value = client
358 .post_json("v1/sql/reads", &request)
359 .await
360 .map_err(map_distributed_client_error)?;
361 Err(distributed_read_error(
362 DistributedReadOutcome::TerminalFailure,
363 "distributed-read server response requires a prepared-result stream adapter",
364 ))
365}
366
367fn read_mode_name(mode: SqlReadMode) -> &'static str {
368 match mode {
369 SqlReadMode::Local => "local",
370 SqlReadMode::Inherit => "inherit",
371 SqlReadMode::Strong => "strong",
372 SqlReadMode::Stale => "stale",
373 }
374}
375
376fn requested_read_mode_name(
377 command_mode: Option<SqlReadMode>,
378 resolved_mode: ResolvedSqlReadMode,
379) -> &'static str {
380 let default = match resolved_mode {
381 ResolvedSqlReadMode::Local => SqlReadMode::Local,
382 ResolvedSqlReadMode::Cluster(_) => SqlReadMode::Inherit,
383 };
384 read_mode_name(command_mode.unwrap_or(default))
385}
386
387fn emit_routing_report(
388 format: RoutingReportFormat,
389 requested_mode: &str,
390 effective_mode: Option<String>,
391 decision: &str,
392 result: &Result<()>,
393) -> Result<()> {
394 let (outcome, reason) = match result {
395 Ok(()) => (DistributedReadOutcome::Success.as_str().to_string(), None),
396 Err(CliError::DistributedReadOutcome {
397 outcome, reason, ..
398 }) => (outcome.clone(), Some(reason.clone())),
399 Err(CliError::Cancelled) => (
400 DistributedReadOutcome::Cancelled.as_str().to_string(),
401 Some("client cancellation".into()),
402 ),
403 Err(CliError::Timeout(reason)) => (
404 DistributedReadOutcome::RetryableFailure
405 .as_str()
406 .to_string(),
407 Some(reason.clone()),
408 ),
409 Err(CliError::ServerUnsupported(reason)) => (
410 DistributedReadOutcome::Unsupported.as_str().to_string(),
411 Some(reason.clone()),
412 ),
413 Err(error) => (
414 DistributedReadOutcome::TerminalFailure.as_str().to_string(),
415 Some(error.to_string()),
416 ),
417 };
418 let report = DistributedReadRoutingReport::new(
419 requested_mode,
420 effective_mode,
421 decision,
422 outcome,
423 reason,
424 );
425 let stderr = std::io::stderr();
426 let mut stderr = stderr.lock();
427 write_distributed_read_routing_report(&mut stderr, format, &report)
428}
429
430fn distributed_read_error(outcome: DistributedReadOutcome, reason: impl Into<String>) -> CliError {
431 let reason = reason.into();
432 CliError::DistributedReadOutcome {
433 outcome: outcome.as_str().to_string(),
434 reason,
435 exit_code: outcome.exit_code(),
436 }
437}
438
439fn map_distributed_client_error(error: ClientError) -> CliError {
440 match error {
441 ClientError::Request { source, .. } => distributed_read_error(
442 DistributedReadOutcome::RetryableFailure,
443 format!("request failed: {source}"),
444 ),
445 ClientError::Auth(error) => distributed_read_error(
446 DistributedReadOutcome::AuthorizationFailure,
447 error.to_string(),
448 ),
449 ClientError::HttpStatus { status, body } => {
450 let code = serde_json::from_str::<serde_json::Value>(&body)
451 .ok()
452 .and_then(|value| {
453 value
454 .get("error")
455 .and_then(|error| error.get("code"))
456 .and_then(|code| code.as_str())
457 .map(str::to_owned)
458 });
459 let outcome = match code.as_deref() {
460 Some("UNAUTHORIZED") | Some("AUTHORIZATION_FAILURE") => {
461 DistributedReadOutcome::AuthorizationFailure
462 }
463 Some("UNSUPPORTED_REMOTE")
464 | Some("NOT_IMPLEMENTED")
465 | Some("FUTURE_DISTRIBUTED_EXECUTION_REQUIRED") => {
466 DistributedReadOutcome::Unsupported
467 }
468 Some("CAPABILITY_UNAVAILABLE") => DistributedReadOutcome::TerminalFailure,
469 _ if status.is_server_error() || status.as_u16() == 408 => {
470 DistributedReadOutcome::RetryableFailure
471 }
472 _ => DistributedReadOutcome::TerminalFailure,
473 };
474 distributed_read_error(outcome, format!("HTTP {}: {body}", status.as_u16()))
475 }
476 ClientError::InvalidUrl(message) | ClientError::Build(message) => {
477 distributed_read_error(DistributedReadOutcome::TerminalFailure, message)
478 }
479 }
480}
481
482fn sql_context_message(sql: &str) -> String {
483 let condensed = sql.split_whitespace().collect::<Vec<_>>().join(" ");
484 let max_len = 200;
485 if condensed.chars().count() > max_len {
486 let truncated: String = condensed.chars().take(max_len).collect();
487 format!("SQL: {truncated}...")
488 } else {
489 format!("SQL: {condensed}")
490 }
491}
492
493#[allow(clippy::too_many_arguments)]
497pub async fn execute_remote_with_formatter_factory<'a, W: Write>(
498 client: &HttpClient,
499 cmd: &SqlCommand,
500 batch_mode: &BatchMode,
501 ui_mode: UiMode,
502 writer: &mut W,
503 make_formatter: &mut dyn FnMut() -> Box<dyn Formatter>,
504 admin_launcher: Option<Box<dyn FnMut() -> Result<()> + 'a>>,
505 limit: Option<usize>,
506 quiet: bool,
507) -> Result<()> {
508 let effective_limit = merge_limit(limit, cmd.max_rows);
509 let deadline = Deadline::new(parse_deadline(cmd.deadline.as_deref())?);
510 let cancel = CancelSignal::new();
511 let options = SqlExecutionOptions {
512 limit: effective_limit,
513 quiet,
514 cancel: &cancel,
515 deadline: &deadline,
516 admin_launcher,
517 };
518 let mut output = SqlOutput::Custom(make_formatter);
519 execute_remote_with_output_control(
520 client,
521 cmd,
522 batch_mode,
523 ui_mode,
524 writer,
525 &mut output,
526 options,
527 )
528 .await
529}
530
531#[doc(hidden)]
532pub async fn execute_remote_with_formatter_control<W: Write>(
533 client: &HttpClient,
534 cmd: &SqlCommand,
535 batch_mode: &BatchMode,
536 ui_mode: UiMode,
537 writer: &mut W,
538 output_format: OutputFormat,
539 options: SqlExecutionOptions<'_>,
540) -> Result<()> {
541 let mut output = SqlOutput::Format(output_format);
542 execute_remote_with_output_control(
543 client,
544 cmd,
545 batch_mode,
546 ui_mode,
547 writer,
548 &mut output,
549 options,
550 )
551 .await
552}
553
554async fn execute_remote_with_output_control<W: Write>(
555 client: &HttpClient,
556 cmd: &SqlCommand,
557 batch_mode: &BatchMode,
558 ui_mode: UiMode,
559 writer: &mut W,
560 output: &mut SqlOutput<'_>,
561 options: SqlExecutionOptions<'_>,
562) -> Result<()> {
563 let sql = cmd.resolve_query(batch_mode)?;
564 if ui_mode == UiMode::Tui {
565 return execute_tui_remote_or_fallback(client, &sql, cmd, writer, output, options).await;
566 }
567 execute_remote_with_formatter_impl(client, &sql, cmd, writer, output, &options).await
568}
569
570async fn execute_remote_with_formatter_impl<W: Write>(
571 client: &HttpClient,
572 sql: &str,
573 cmd: &SqlCommand,
574 writer: &mut W,
575 output: &mut SqlOutput<'_>,
576 options: &SqlExecutionOptions<'_>,
577) -> Result<()> {
578 if is_select_query(sql)? && output.supports_streaming() {
579 return execute_remote_streaming(
580 client,
581 sql,
582 writer,
583 output,
584 options,
585 cmd.fetch_size,
586 cmd.max_rows,
587 )
588 .await;
589 }
590
591 let request = RemoteSqlRequest {
592 sql: sql.to_string(),
593 streaming: false,
594 fetch_size: cmd.fetch_size,
595 max_rows: cmd.max_rows,
596 };
597 let response: RemoteSqlResponse = tokio::select! {
598 result = tokio::time::timeout(options.deadline.remaining(), client.post_json("api/sql/query", &request)) => {
599 match result {
600 Ok(value) => value.map_err(map_client_error)?,
601 Err(_) => {
602 let _ = send_cancel_request(client).await;
603 return Err(CliError::Timeout(format!(
604 "deadline exceeded after {}",
605 humantime::format_duration(options.deadline.duration())
606 )));
607 }
608 }
609 }
610 _ = options.cancel.wait() => {
611 let _ = send_cancel_request(client).await;
612 return Err(CliError::Cancelled);
613 }
614 };
615
616 let json_array = output.statement_array();
618 if json_array {
619 writeln!(writer, "[")?;
620 }
621 let mut emitted = false;
622 for response in response.into_results() {
623 if response.columns.is_empty() && options.quiet {
624 continue;
625 }
626 if json_array && emitted {
627 writeln!(writer, ",")?;
628 }
629 emitted = true;
630
631 if response.columns.is_empty() {
632 let message = match response.affected_rows {
633 Some(count) => format!("{count} row(s) affected"),
634 None => "Operation completed successfully".to_string(),
635 };
636 if json_array {
637 writeln!(writer, "[")?;
638 }
639 let columns = sql_status_columns();
640 {
641 let mut streaming_writer = StreamingWriter::new(
642 &mut *writer,
643 output.create_block_formatter(),
644 columns,
645 options.limit,
646 )
647 .with_quiet(options.quiet);
648 streaming_writer.prepare(Some(1))?;
649 let row = Row::new(vec![Value::Text("OK".to_string()), Value::Text(message)]);
650 streaming_writer.write_row(row)?;
651 streaming_writer.finish()?;
652 }
653 if json_array {
654 writeln!(writer, "]")?;
655 }
656 continue;
657 }
658
659 let columns: Vec<Column> = response
660 .columns
661 .iter()
662 .map(|col| Column::new(&col.name, data_type_from_string(&col.data_type)))
663 .collect();
664 if json_array {
665 writeln!(writer, "[")?;
666 }
667 {
668 let mut streaming_writer = StreamingWriter::new(
669 &mut *writer,
670 output.create_block_formatter(),
671 columns,
672 options.limit,
673 )
674 .with_quiet(options.quiet);
675 streaming_writer.prepare(Some(response.rows.len()))?;
676 for row in response.rows {
677 if options.cancel.is_cancelled() {
678 let _ = send_cancel_request(client).await;
679 return Err(CliError::Cancelled);
680 }
681 options.deadline.check()?;
682 let values = row.into_iter().map(remote_value_to_value).collect();
683 match streaming_writer.write_row(Row::new(values))? {
684 WriteStatus::LimitReached => break,
685 WriteStatus::Continue => {}
686 }
687 }
688 streaming_writer.finish()?;
689 }
690 if json_array {
691 writeln!(writer, "]")?;
692 }
693 }
694 if json_array {
695 writeln!(writer, "]")?;
696 }
697 Ok(())
698}
699
700fn execute_tui_local_or_fallback<'a, W: Write>(
701 db: &Database,
702 sql: &str,
703 writer: &mut W,
704 output: &mut SqlOutput<'_>,
705 mut options: SqlExecutionOptions<'a>,
706) -> Result<()> {
707 if !is_tty() {
708 if !options.quiet {
709 eprintln!("Warning: --tui requires a TTY, falling back to batch output.");
710 }
711 return execute_sql_with_formatter(db, sql, writer, output, &options);
712 }
713
714 let admin_launcher = options.admin_launcher.take();
715 match execute_tui_local(db, sql, &options, admin_launcher) {
716 Ok(()) => Ok(()),
717 Err(err) => {
718 if !options.quiet {
719 eprintln!("Warning: TUI failed ({err}); falling back to batch output.");
720 }
721 execute_sql_with_formatter(db, sql, writer, output, &options)
722 }
723 }
724}
725
726fn execute_tui_local<'a>(
727 db: &Database,
728 sql: &str,
729 options: &SqlExecutionOptions<'a>,
730 admin_launcher: Option<Box<dyn FnMut() -> Result<()> + 'a>>,
731) -> Result<()> {
732 use alopex_sql::ExecutionResult;
733
734 options.deadline.check()?;
735 let result = db.execute_sql(sql)?;
736 options.deadline.check()?;
737
738 let (columns, rows) = match result {
739 ExecutionResult::Success => {
740 let columns = sql_status_columns();
741 let row = Row::new(vec![
742 Value::Text("OK".to_string()),
743 Value::Text("Operation completed successfully".to_string()),
744 ]);
745 (columns, vec![row])
746 }
747 ExecutionResult::RowsAffected(count) => {
748 let columns = sql_status_columns();
749 let row = Row::new(vec![
750 Value::Text("OK".to_string()),
751 Value::Text(format!("{count} row(s) affected")),
752 ]);
753 (columns, vec![row])
754 }
755 ExecutionResult::Query(query_result) => {
756 let columns = query_result
757 .columns
758 .iter()
759 .map(|col| Column::new(&col.name, DataType::Text))
760 .collect::<Vec<_>>();
761 let mut rows = Vec::with_capacity(query_result.rows.len());
762 for sql_row in query_result.rows {
763 let values = sql_row.into_iter().map(sql_value_to_value).collect();
764 rows.push(Row::new(values));
765 }
766 if let Some(limit) = options.limit {
767 rows.truncate(limit);
768 }
769 (columns, rows)
770 }
771 };
772
773 let app = TuiApp::new(columns, rows, "local", false)
774 .with_context_message(Some(sql_context_message(sql)))
775 .with_admin_launcher(admin_launcher);
776 app.run()
777}
778
779async fn execute_tui_remote_or_fallback<'a, W: Write>(
780 client: &HttpClient,
781 sql: &str,
782 cmd: &SqlCommand,
783 writer: &mut W,
784 output: &mut SqlOutput<'_>,
785 mut options: SqlExecutionOptions<'a>,
786) -> Result<()> {
787 if !is_tty() {
788 if !options.quiet {
789 eprintln!("Warning: --tui requires a TTY, falling back to batch output.");
790 }
791 return execute_remote_with_formatter_impl(client, sql, cmd, writer, output, &options)
792 .await;
793 }
794
795 let admin_launcher = options.admin_launcher.take();
796 match execute_tui_remote(client, sql, cmd, &options, admin_launcher).await {
797 Ok(()) => Ok(()),
798 Err(err) => {
799 if !options.quiet {
800 eprintln!("Warning: TUI failed ({err}); falling back to batch output.");
801 }
802 execute_remote_with_formatter_impl(client, sql, cmd, writer, output, &options).await
803 }
804 }
805}
806
807async fn execute_tui_remote<'a>(
808 client: &HttpClient,
809 sql: &str,
810 cmd: &SqlCommand,
811 options: &SqlExecutionOptions<'a>,
812 admin_launcher: Option<Box<dyn FnMut() -> Result<()> + 'a>>,
813) -> Result<()> {
814 if is_select_query(sql)? {
815 let (columns, rows) =
816 collect_remote_streaming_rows(client, sql, options, cmd.fetch_size, cmd.max_rows)
817 .await?;
818 let app = TuiApp::new(columns, rows, "server", false)
819 .with_context_message(Some(sql_context_message(sql)))
820 .with_admin_launcher(admin_launcher);
821 return app.run();
822 }
823
824 let request = RemoteSqlRequest {
825 sql: sql.to_string(),
826 streaming: false,
827 fetch_size: cmd.fetch_size,
828 max_rows: cmd.max_rows,
829 };
830
831 let response: RemoteSqlResponse = tokio::select! {
832 result = tokio::time::timeout(options.deadline.remaining(), client.post_json("api/sql/query", &request)) => {
833 match result {
834 Ok(value) => value.map_err(map_client_error)?,
835 Err(_) => {
836 let _ = send_cancel_request(client).await;
837 return Err(CliError::Timeout(format!(
838 "deadline exceeded after {}",
839 humantime::format_duration(options.deadline.duration())
840 )));
841 }
842 }
843 }
844 _ = options.cancel.wait() => {
845 let _ = send_cancel_request(client).await;
846 return Err(CliError::Cancelled);
847 }
848 };
849
850 let response = response
851 .into_results()
852 .into_iter()
853 .last()
854 .unwrap_or_default();
855 let (columns, rows) = if response.columns.is_empty() {
856 let columns = sql_status_columns();
857 let message = match response.affected_rows {
858 Some(count) => format!("{count} row(s) affected"),
859 None => "Operation completed successfully".to_string(),
860 };
861 let row = Row::new(vec![Value::Text("OK".to_string()), Value::Text(message)]);
862 (columns, vec![row])
863 } else {
864 let columns: Vec<Column> = response
865 .columns
866 .iter()
867 .map(|col| Column::new(&col.name, data_type_from_string(&col.data_type)))
868 .collect();
869 let mut rows = response
870 .rows
871 .into_iter()
872 .map(|row| Row::new(row.into_iter().map(remote_value_to_value).collect()))
873 .collect::<Vec<_>>();
874 if let Some(limit) = options.limit {
875 rows.truncate(limit);
876 }
877 (columns, rows)
878 };
879
880 let app = TuiApp::new(columns, rows, "server", false)
881 .with_context_message(Some(sql_context_message(sql)))
882 .with_admin_launcher(admin_launcher);
883 app.run()
884}
885
886async fn collect_remote_streaming_rows(
887 client: &HttpClient,
888 sql: &str,
889 options: &SqlExecutionOptions<'_>,
890 fetch_size: Option<usize>,
891 max_rows: Option<usize>,
892) -> Result<(Vec<Column>, Vec<Row>)> {
893 let request = RemoteSqlRequest {
894 sql: sql.to_string(),
895 streaming: true,
896 fetch_size,
897 max_rows,
898 };
899
900 let response = tokio::select! {
901 result = tokio::time::timeout(options.deadline.remaining(), client.post_json_stream("api/sql/query", &request)) => {
902 match result {
903 Ok(value) => value.map_err(map_client_error)?,
904 Err(_) => {
905 let _ = send_cancel_request(client).await;
906 return Err(CliError::Timeout(format!(
907 "deadline exceeded after {}",
908 humantime::format_duration(options.deadline.duration())
909 )));
910 }
911 }
912 }
913 _ = options.cancel.wait() => {
914 let _ = send_cancel_request(client).await;
915 return Err(CliError::Cancelled);
916 }
917 };
918
919 if let Some(content_type) = response
920 .headers()
921 .get(reqwest::header::CONTENT_TYPE)
922 .and_then(|value| value.to_str().ok())
923 {
924 if content_type.starts_with("application/jsonl") {
925 return collect_remote_jsonl_rows(client, response, options).await;
926 }
927 }
928
929 let mut stream = response.bytes_stream();
930 let mut buffer: Vec<u8> = Vec::new();
931 let mut pos: usize = 0;
932 let mut done = false;
933 let mut saw_array_start = false;
934 let mut columns: Option<Vec<String>> = None;
935 let mut column_set: Option<HashSet<String>> = None;
936 let mut rows: Vec<Row> = Vec::new();
937
938 while !done {
939 if options.cancel.is_cancelled() {
940 let _ = send_cancel_request(client).await;
941 return Err(CliError::Cancelled);
942 }
943 if let Err(err) = options.deadline.check() {
944 let _ = send_cancel_request(client).await;
945 return Err(err);
946 }
947
948 let next = tokio::select! {
949 _ = options.cancel.wait() => {
950 let _ = send_cancel_request(client).await;
951 return Err(CliError::Cancelled);
952 }
953 result = tokio::time::timeout(options.deadline.remaining(), stream.next()) => {
954 match result {
955 Ok(value) => value,
956 Err(_) => {
957 let _ = send_cancel_request(client).await;
958 return Err(CliError::Timeout(format!(
959 "deadline exceeded after {}",
960 humantime::format_duration(options.deadline.duration())
961 )));
962 }
963 }
964 }
965 };
966
967 let chunk = match next {
968 Some(chunk) => chunk,
969 None => break,
970 };
971
972 let bytes = match chunk {
973 Ok(bytes) => bytes,
974 Err(err) => return Err(CliError::ServerConnection(format!("request failed: {err}"))),
975 };
976
977 buffer.extend_from_slice(&bytes);
978
979 loop {
980 skip_whitespace(&buffer, &mut pos);
981 if pos >= buffer.len() {
982 break;
983 }
984
985 if !saw_array_start {
986 if buffer[pos] != b'[' {
987 return Err(CliError::InvalidArgument(
988 "Invalid streaming response: expected JSON array".into(),
989 ));
990 }
991 pos += 1;
992 saw_array_start = true;
993 continue;
994 }
995
996 skip_whitespace(&buffer, &mut pos);
997 if pos >= buffer.len() {
998 break;
999 }
1000
1001 if buffer[pos] == b']' {
1002 pos += 1;
1003 done = true;
1004 break;
1005 }
1006
1007 let slice = &buffer[pos..];
1008 let mut stream =
1009 serde_json::Deserializer::from_slice(slice).into_iter::<serde_json::Value>();
1010 let value = match stream.next() {
1011 Some(Ok(value)) => value,
1012 Some(Err(err)) if err.is_eof() => break,
1013 Some(Err(err)) => return Err(CliError::Json(err)),
1014 None => break,
1015 };
1016 pos = pos.saturating_add(stream.byte_offset());
1017
1018 let object = value.as_object().ok_or_else(|| {
1019 CliError::InvalidArgument("Invalid streaming row: expected JSON object".into())
1020 })?;
1021
1022 if columns.is_none() {
1023 let names: Vec<String> = object.keys().cloned().collect();
1024 let set: HashSet<String> = names.iter().cloned().collect();
1025 if names.is_empty() {
1026 return Err(CliError::InvalidArgument(
1027 "Invalid streaming row: empty object".into(),
1028 ));
1029 }
1030 columns = Some(names);
1031 column_set = Some(set);
1032 }
1033
1034 let names = columns
1035 .as_ref()
1036 .ok_or_else(|| CliError::InvalidArgument("Missing columns".into()))?;
1037 let set = column_set
1038 .as_ref()
1039 .ok_or_else(|| CliError::InvalidArgument("Missing column set".into()))?;
1040
1041 if object.len() != names.len() || !object.keys().all(|key| set.contains(key)) {
1042 return Err(CliError::InvalidArgument(
1043 "Invalid streaming row: column mismatch".into(),
1044 ));
1045 }
1046
1047 let values = names
1048 .iter()
1049 .map(|name| {
1050 object.get(name).ok_or_else(|| {
1051 CliError::InvalidArgument(format!(
1052 "Invalid streaming row: missing column '{name}'"
1053 ))
1054 })
1055 })
1056 .map(|value| value.and_then(json_value_to_value))
1057 .collect::<Result<Vec<_>>>()?;
1058 rows.push(Row::new(values));
1059
1060 if let Some(limit) = options.limit {
1061 if rows.len() >= limit {
1062 let _ = send_cancel_request(client).await;
1063 done = true;
1064 break;
1065 }
1066 }
1067
1068 skip_whitespace(&buffer, &mut pos);
1069 if pos >= buffer.len() {
1070 break;
1071 }
1072 match buffer[pos] {
1073 b',' => {
1074 pos += 1;
1075 }
1076 b']' => {
1077 pos += 1;
1078 done = true;
1079 break;
1080 }
1081 _ => {
1082 return Err(CliError::InvalidArgument(
1083 "Invalid streaming response: expected ',' or ']'".into(),
1084 ));
1085 }
1086 }
1087 }
1088
1089 if pos > 0 {
1090 buffer.drain(..pos);
1091 pos = 0;
1092 }
1093 }
1094
1095 if done {
1096 if has_non_whitespace(&buffer) {
1097 return Err(CliError::InvalidArgument(
1098 "Invalid streaming response: unexpected trailing data".into(),
1099 ));
1100 }
1101 buffer.clear();
1102 loop {
1103 let next = tokio::select! {
1104 _ = options.cancel.wait() => {
1105 let _ = send_cancel_request(client).await;
1106 return Err(CliError::Cancelled);
1107 }
1108 result = tokio::time::timeout(options.deadline.remaining(), stream.next()) => {
1109 match result {
1110 Ok(value) => value,
1111 Err(_) => {
1112 let _ = send_cancel_request(client).await;
1113 return Err(CliError::Timeout(format!(
1114 "deadline exceeded after {}",
1115 humantime::format_duration(options.deadline.duration())
1116 )));
1117 }
1118 }
1119 }
1120 };
1121
1122 let chunk = match next {
1123 Some(chunk) => chunk,
1124 None => break,
1125 };
1126
1127 let bytes = match chunk {
1128 Ok(bytes) => bytes,
1129 Err(err) => {
1130 return Err(CliError::ServerConnection(format!("request failed: {err}")));
1131 }
1132 };
1133
1134 if has_non_whitespace(&bytes) {
1135 return Err(CliError::InvalidArgument(
1136 "Invalid streaming response: unexpected trailing data".into(),
1137 ));
1138 }
1139 }
1140 } else {
1141 skip_whitespace(&buffer, &mut pos);
1142 if pos < buffer.len() {
1143 return Err(CliError::InvalidArgument(
1144 "Invalid streaming response: unexpected trailing data".into(),
1145 ));
1146 }
1147 return Err(CliError::InvalidArgument(
1148 "Invalid streaming response: unexpected end of stream".into(),
1149 ));
1150 }
1151
1152 let columns = columns
1153 .unwrap_or_default()
1154 .into_iter()
1155 .map(|name| Column::new(name, DataType::Text))
1156 .collect();
1157 Ok((columns, rows))
1158}
1159
1160#[allow(dead_code)]
1161async fn collect_remote_non_streaming_rows(
1162 client: &HttpClient,
1163 sql: &str,
1164 options: &SqlExecutionOptions<'_>,
1165 fetch_size: Option<usize>,
1166 max_rows: Option<usize>,
1167) -> Result<(Vec<Column>, Vec<Row>)> {
1168 let request = RemoteSqlRequest {
1169 sql: sql.to_string(),
1170 streaming: false,
1171 fetch_size,
1172 max_rows,
1173 };
1174 let response: RemoteSqlResponse = tokio::select! {
1175 result = tokio::time::timeout(options.deadline.remaining(), client.post_json("api/sql/query", &request)) => {
1176 match result {
1177 Ok(value) => value.map_err(map_client_error)?,
1178 Err(_) => {
1179 let _ = send_cancel_request(client).await;
1180 return Err(CliError::Timeout(format!(
1181 "deadline exceeded after {}",
1182 humantime::format_duration(options.deadline.duration())
1183 )));
1184 }
1185 }
1186 }
1187 _ = options.cancel.wait() => {
1188 let _ = send_cancel_request(client).await;
1189 return Err(CliError::Cancelled);
1190 }
1191 };
1192
1193 let response = response
1194 .into_results()
1195 .into_iter()
1196 .last()
1197 .unwrap_or_default();
1198 if response.columns.is_empty() {
1199 let message = match response.affected_rows {
1200 Some(count) => format!("{count} row(s) affected"),
1201 None => "Operation completed successfully".to_string(),
1202 };
1203 let columns = sql_status_columns();
1204 let row = Row::new(vec![Value::Text("OK".to_string()), Value::Text(message)]);
1205 return Ok((columns, vec![row]));
1206 }
1207
1208 let columns = response
1209 .columns
1210 .iter()
1211 .map(|col| Column::new(&col.name, data_type_from_string(&col.data_type)))
1212 .collect::<Vec<_>>();
1213 let mut rows = response
1214 .rows
1215 .into_iter()
1216 .map(|row| Row::new(row.into_iter().map(remote_value_to_value).collect()))
1217 .collect::<Vec<_>>();
1218 if let Some(limit) = options.limit {
1219 rows.truncate(limit);
1220 }
1221 Ok((columns, rows))
1222}
1223
1224async fn execute_remote_streaming<W: Write>(
1225 client: &HttpClient,
1226 sql: &str,
1227 writer: &mut W,
1228 output: &mut SqlOutput<'_>,
1229 options: &SqlExecutionOptions<'_>,
1230 fetch_size: Option<usize>,
1231 max_rows: Option<usize>,
1232) -> Result<()> {
1233 let request = RemoteSqlRequest {
1234 sql: sql.to_string(),
1235 streaming: true,
1236 fetch_size,
1237 max_rows,
1238 };
1239
1240 let response = tokio::select! {
1241 result = tokio::time::timeout(options.deadline.remaining(), client.post_json_stream("api/sql/query", &request)) => {
1242 match result {
1243 Ok(value) => value.map_err(map_client_error)?,
1244 Err(_) => {
1245 let _ = send_cancel_request(client).await;
1246 return Err(CliError::Timeout(format!(
1247 "deadline exceeded after {}",
1248 humantime::format_duration(options.deadline.duration())
1249 )));
1250 }
1251 }
1252 }
1253 _ = options.cancel.wait() => {
1254 let _ = send_cancel_request(client).await;
1255 return Err(CliError::Cancelled);
1256 }
1257 };
1258
1259 if !output.statement_array() {
1262 return stream_remote_result_set(
1263 client,
1264 response,
1265 writer,
1266 output.create_block_formatter(),
1267 options,
1268 )
1269 .await;
1270 }
1271
1272 writeln!(writer, "[")?;
1273 let mut result_set_writer = CountingWriter::new(&mut *writer);
1274 let result = stream_remote_result_set(
1275 client,
1276 response,
1277 &mut result_set_writer,
1278 output.create_block_formatter(),
1279 options,
1280 )
1281 .await;
1282 let result_set_started = result_set_writer.bytes_written() > 0;
1283 match result {
1284 Ok(()) => {
1285 writeln!(writer, "]")?;
1286 Ok(())
1287 }
1288 Err(err) => {
1289 if result_set_started {
1294 let _ = writeln!(writer, "]");
1295 }
1296 let _ = writeln!(writer, "]");
1297 Err(err)
1298 }
1299 }
1300}
1301
1302async fn stream_remote_result_set<W: Write>(
1304 client: &HttpClient,
1305 response: Response,
1306 writer: &mut W,
1307 formatter: Box<dyn Formatter>,
1308 options: &SqlExecutionOptions<'_>,
1309) -> Result<()> {
1310 if let Some(content_type) = response
1311 .headers()
1312 .get(reqwest::header::CONTENT_TYPE)
1313 .and_then(|value| value.to_str().ok())
1314 {
1315 if content_type.starts_with("application/jsonl") {
1316 return execute_remote_jsonl_streaming(client, response, writer, formatter, options)
1317 .await;
1318 }
1319 }
1320
1321 let mut stream = response.bytes_stream();
1322 let mut buffer: Vec<u8> = Vec::new();
1323 let mut pos: usize = 0;
1324 let mut streaming_writer: Option<StreamingWriter<&mut W>> = None;
1325 let mut formatter = Some(formatter);
1326 let mut columns: Option<Vec<String>> = None;
1327 let mut column_set: Option<HashSet<String>> = None;
1328 let mut done = false;
1329 let mut saw_array_start = false;
1330
1331 while !done {
1332 if options.cancel.is_cancelled() {
1333 let _ = send_cancel_request(client).await;
1334 return Err(CliError::Cancelled);
1335 }
1336 if let Err(err) = options.deadline.check() {
1337 let _ = send_cancel_request(client).await;
1338 return Err(err);
1339 }
1340
1341 let next = tokio::select! {
1342 _ = options.cancel.wait() => {
1343 let _ = send_cancel_request(client).await;
1344 return Err(CliError::Cancelled);
1345 }
1346 result = tokio::time::timeout(options.deadline.remaining(), stream.next()) => {
1347 match result {
1348 Ok(value) => value,
1349 Err(_) => {
1350 let _ = send_cancel_request(client).await;
1351 return Err(CliError::Timeout(format!(
1352 "deadline exceeded after {}",
1353 humantime::format_duration(options.deadline.duration())
1354 )));
1355 }
1356 }
1357 }
1358 };
1359
1360 let chunk = match next {
1361 Some(chunk) => chunk,
1362 None => break,
1363 };
1364
1365 let bytes = match chunk {
1366 Ok(bytes) => bytes,
1367 Err(err) => return Err(CliError::ServerConnection(format!("request failed: {err}"))),
1368 };
1369
1370 buffer.extend_from_slice(&bytes);
1371
1372 loop {
1373 skip_whitespace(&buffer, &mut pos);
1374 if pos >= buffer.len() {
1375 break;
1376 }
1377
1378 if !saw_array_start {
1379 if buffer[pos] != b'[' {
1380 return Err(CliError::InvalidArgument(
1381 "Invalid streaming response: expected JSON array".into(),
1382 ));
1383 }
1384 pos += 1;
1385 saw_array_start = true;
1386 continue;
1387 }
1388
1389 skip_whitespace(&buffer, &mut pos);
1390 if pos >= buffer.len() {
1391 break;
1392 }
1393
1394 if buffer[pos] == b']' {
1395 pos += 1;
1396 done = true;
1397 break;
1398 }
1399
1400 let slice = &buffer[pos..];
1401 let mut stream =
1402 serde_json::Deserializer::from_slice(slice).into_iter::<serde_json::Value>();
1403 let value = match stream.next() {
1404 Some(Ok(value)) => value,
1405 Some(Err(err)) if err.is_eof() => break,
1406 Some(Err(err)) => return Err(CliError::Json(err)),
1407 None => break,
1408 };
1409 pos = pos.saturating_add(stream.byte_offset());
1410
1411 let object = value.as_object().ok_or_else(|| {
1412 CliError::InvalidArgument("Invalid streaming row: expected JSON object".into())
1413 })?;
1414
1415 if columns.is_none() {
1416 let names: Vec<String> = object.keys().cloned().collect();
1417 let set: HashSet<String> = names.iter().cloned().collect();
1418 if names.is_empty() {
1419 return Err(CliError::InvalidArgument(
1420 "Invalid streaming row: empty object".into(),
1421 ));
1422 }
1423 let cols = names
1424 .iter()
1425 .map(|name| Column::new(name, DataType::Text))
1426 .collect::<Vec<_>>();
1427 let formatter = formatter
1428 .take()
1429 .ok_or_else(|| CliError::InvalidArgument("Missing formatter".into()))?;
1430 let mut writer = StreamingWriter::new(&mut *writer, formatter, cols, options.limit)
1431 .with_quiet(options.quiet);
1432 writer.prepare(None)?;
1433 streaming_writer = Some(writer);
1434 columns = Some(names);
1435 column_set = Some(set);
1436 }
1437
1438 let names = columns
1439 .as_ref()
1440 .ok_or_else(|| CliError::InvalidArgument("Missing columns".into()))?;
1441 let set = column_set
1442 .as_ref()
1443 .ok_or_else(|| CliError::InvalidArgument("Missing column set".into()))?;
1444
1445 if object.len() != names.len() || !object.keys().all(|key| set.contains(key)) {
1446 return Err(CliError::InvalidArgument(
1447 "Invalid streaming row: column mismatch".into(),
1448 ));
1449 }
1450
1451 let values = names
1452 .iter()
1453 .map(|name| {
1454 object.get(name).ok_or_else(|| {
1455 CliError::InvalidArgument(format!(
1456 "Invalid streaming row: missing column '{name}'"
1457 ))
1458 })
1459 })
1460 .map(|value| value.and_then(json_value_to_value))
1461 .collect::<Result<Vec<_>>>()?;
1462
1463 if let Some(writer) = streaming_writer.as_mut() {
1464 match writer.write_row(Row::new(values))? {
1465 WriteStatus::LimitReached => {
1466 let _ = send_cancel_request(client).await;
1467 return writer.finish();
1468 }
1469 WriteStatus::Continue => {}
1470 }
1471 }
1472
1473 skip_whitespace(&buffer, &mut pos);
1474 if pos >= buffer.len() {
1475 break;
1476 }
1477 match buffer[pos] {
1478 b',' => {
1479 pos += 1;
1480 }
1481 b']' => {
1482 pos += 1;
1483 done = true;
1484 break;
1485 }
1486 _ => {
1487 return Err(CliError::InvalidArgument(
1488 "Invalid streaming response: expected ',' or ']'".into(),
1489 ));
1490 }
1491 }
1492 }
1493
1494 if pos > 0 {
1495 buffer.drain(..pos);
1496 pos = 0;
1497 }
1498 }
1499
1500 if done {
1501 if has_non_whitespace(&buffer) {
1502 return Err(CliError::InvalidArgument(
1503 "Invalid streaming response: unexpected trailing data".into(),
1504 ));
1505 }
1506 buffer.clear();
1507 loop {
1508 let next = tokio::select! {
1509 _ = options.cancel.wait() => {
1510 let _ = send_cancel_request(client).await;
1511 return Err(CliError::Cancelled);
1512 }
1513 result = tokio::time::timeout(options.deadline.remaining(), stream.next()) => {
1514 match result {
1515 Ok(value) => value,
1516 Err(_) => {
1517 let _ = send_cancel_request(client).await;
1518 return Err(CliError::Timeout(format!(
1519 "deadline exceeded after {}",
1520 humantime::format_duration(options.deadline.duration())
1521 )));
1522 }
1523 }
1524 }
1525 };
1526
1527 let chunk = match next {
1528 Some(chunk) => chunk,
1529 None => break,
1530 };
1531
1532 let bytes = match chunk {
1533 Ok(bytes) => bytes,
1534 Err(err) => {
1535 return Err(CliError::ServerConnection(format!("request failed: {err}")));
1536 }
1537 };
1538
1539 if has_non_whitespace(&bytes) {
1540 return Err(CliError::InvalidArgument(
1541 "Invalid streaming response: unexpected trailing data".into(),
1542 ));
1543 }
1544 }
1545 } else {
1546 skip_whitespace(&buffer, &mut pos);
1547 if pos < buffer.len() {
1548 return Err(CliError::InvalidArgument(
1549 "Invalid streaming response: unexpected trailing data".into(),
1550 ));
1551 }
1552 return Err(CliError::InvalidArgument(
1553 "Invalid streaming response: unexpected end of stream".into(),
1554 ));
1555 }
1556
1557 if let Some(mut writer) = streaming_writer {
1558 return writer.finish();
1559 }
1560
1561 if done && saw_array_start {
1562 if let Some(formatter) = formatter.take() {
1563 let mut writer =
1564 StreamingWriter::new(&mut *writer, formatter, Vec::new(), options.limit)
1565 .with_quiet(options.quiet);
1566 writer.prepare(None)?;
1567 return writer.finish();
1568 }
1569 }
1570
1571 Ok(())
1572}
1573
1574async fn collect_remote_jsonl_rows(
1575 client: &HttpClient,
1576 response: Response,
1577 options: &SqlExecutionOptions<'_>,
1578) -> Result<(Vec<Column>, Vec<Row>)> {
1579 let mut stream = response.bytes_stream();
1580 let mut buffer: Vec<u8> = Vec::new();
1581 let mut columns: Option<Vec<Column>> = None;
1582 let mut rows: Vec<Row> = Vec::new();
1583 let mut done = false;
1584
1585 while !done {
1586 if options.cancel.is_cancelled() {
1587 let _ = send_cancel_request(client).await;
1588 return Err(CliError::Cancelled);
1589 }
1590 if let Err(err) = options.deadline.check() {
1591 let _ = send_cancel_request(client).await;
1592 return Err(err);
1593 }
1594
1595 let next = tokio::select! {
1596 _ = options.cancel.wait() => {
1597 let _ = send_cancel_request(client).await;
1598 return Err(CliError::Cancelled);
1599 }
1600 result = tokio::time::timeout(options.deadline.remaining(), stream.next()) => {
1601 match result {
1602 Ok(value) => value,
1603 Err(_) => {
1604 let _ = send_cancel_request(client).await;
1605 return Err(CliError::Timeout(format!(
1606 "deadline exceeded after {}",
1607 humantime::format_duration(options.deadline.duration())
1608 )));
1609 }
1610 }
1611 }
1612 };
1613
1614 let chunk = match next {
1615 Some(chunk) => chunk,
1616 None => break,
1617 };
1618
1619 let bytes = match chunk {
1620 Ok(bytes) => bytes,
1621 Err(err) => return Err(CliError::ServerConnection(format!("request failed: {err}"))),
1622 };
1623
1624 buffer.extend_from_slice(&bytes);
1625
1626 while let Some(newline) = buffer.iter().position(|&b| b == b'\n') {
1627 let line = buffer.drain(..=newline).collect::<Vec<u8>>();
1628 if let Some(item) = parse_jsonl_line(&line)? {
1629 if let Some(error) = item.error {
1630 return Err(CliError::InvalidArgument(format!(
1631 "Server error: {}",
1632 error.message
1633 )));
1634 }
1635 if item.done {
1636 done = true;
1637 break;
1638 }
1639 let row = item.row.ok_or_else(|| {
1640 CliError::InvalidArgument("Invalid streaming response: missing row".into())
1641 })?;
1642 if columns.is_none() {
1643 columns = Some(default_stream_columns(row.len()));
1644 }
1645 rows.push(Row::new(
1646 row.into_iter().map(remote_value_to_value).collect(),
1647 ));
1648 if let Some(limit) = options.limit {
1649 if rows.len() >= limit {
1650 let _ = send_cancel_request(client).await;
1651 done = true;
1652 break;
1653 }
1654 }
1655 }
1656 }
1657 }
1658
1659 if !done {
1660 if let Some(item) = parse_jsonl_line(&buffer)? {
1661 if let Some(error) = item.error {
1662 return Err(CliError::InvalidArgument(format!(
1663 "Server error: {}",
1664 error.message
1665 )));
1666 }
1667 if item.done {
1668 done = true;
1669 } else if let Some(row) = item.row {
1670 if columns.is_none() {
1671 columns = Some(default_stream_columns(row.len()));
1672 }
1673 rows.push(Row::new(
1674 row.into_iter().map(remote_value_to_value).collect(),
1675 ));
1676 if let Some(limit) = options.limit {
1677 if rows.len() >= limit {
1678 let _ = send_cancel_request(client).await;
1679 done = true;
1680 }
1681 }
1682 } else {
1683 return Err(CliError::InvalidArgument(
1684 "Invalid streaming response: missing row".into(),
1685 ));
1686 }
1687 }
1688 }
1689
1690 if !done {
1691 return Err(CliError::InvalidArgument(
1692 "Invalid streaming response: unexpected end of stream".into(),
1693 ));
1694 }
1695
1696 Ok((columns.unwrap_or_default(), rows))
1697}
1698
1699async fn execute_remote_jsonl_streaming<W: Write>(
1700 client: &HttpClient,
1701 response: Response,
1702 writer: &mut W,
1703 formatter: Box<dyn Formatter>,
1704 options: &SqlExecutionOptions<'_>,
1705) -> Result<()> {
1706 let mut stream = response.bytes_stream();
1707 let mut buffer: Vec<u8> = Vec::new();
1708 let mut streaming_writer: Option<StreamingWriter<&mut W>> = None;
1709 let mut formatter = Some(formatter);
1710 let mut done = false;
1711
1712 while !done {
1713 if options.cancel.is_cancelled() {
1714 let _ = send_cancel_request(client).await;
1715 return Err(CliError::Cancelled);
1716 }
1717 if let Err(err) = options.deadline.check() {
1718 let _ = send_cancel_request(client).await;
1719 return Err(err);
1720 }
1721
1722 let next = tokio::select! {
1723 _ = options.cancel.wait() => {
1724 let _ = send_cancel_request(client).await;
1725 return Err(CliError::Cancelled);
1726 }
1727 result = tokio::time::timeout(options.deadline.remaining(), stream.next()) => {
1728 match result {
1729 Ok(value) => value,
1730 Err(_) => {
1731 let _ = send_cancel_request(client).await;
1732 return Err(CliError::Timeout(format!(
1733 "deadline exceeded after {}",
1734 humantime::format_duration(options.deadline.duration())
1735 )));
1736 }
1737 }
1738 }
1739 };
1740
1741 let chunk = match next {
1742 Some(chunk) => chunk,
1743 None => break,
1744 };
1745
1746 let bytes = match chunk {
1747 Ok(bytes) => bytes,
1748 Err(err) => return Err(CliError::ServerConnection(format!("request failed: {err}"))),
1749 };
1750
1751 buffer.extend_from_slice(&bytes);
1752
1753 while let Some(newline) = buffer.iter().position(|&b| b == b'\n') {
1754 let line = buffer.drain(..=newline).collect::<Vec<u8>>();
1755 if let Some(item) = parse_jsonl_line(&line)? {
1756 if let Some(error) = item.error {
1757 return Err(CliError::InvalidArgument(format!(
1758 "Server error: {}",
1759 error.message
1760 )));
1761 }
1762 if item.done {
1763 done = true;
1764 break;
1765 }
1766 let row = item.row.ok_or_else(|| {
1767 CliError::InvalidArgument("Invalid streaming response: missing row".into())
1768 })?;
1769 if streaming_writer.is_none() {
1770 let columns = default_stream_columns(row.len());
1771 let formatter = formatter
1772 .take()
1773 .ok_or_else(|| CliError::InvalidArgument("Missing formatter".into()))?;
1774 let mut writer =
1775 StreamingWriter::new(&mut *writer, formatter, columns, options.limit)
1776 .with_quiet(options.quiet);
1777 writer.prepare(None)?;
1778 streaming_writer = Some(writer);
1779 }
1780 let values = row.into_iter().map(remote_value_to_value).collect();
1781 if let Some(writer) = streaming_writer.as_mut() {
1782 match writer.write_row(Row::new(values))? {
1783 WriteStatus::LimitReached => {
1784 let _ = send_cancel_request(client).await;
1785 done = true;
1786 break;
1787 }
1788 WriteStatus::Continue => {}
1789 }
1790 }
1791 }
1792 }
1793 }
1794
1795 if !done {
1796 if let Some(item) = parse_jsonl_line(&buffer)? {
1797 if let Some(error) = item.error {
1798 return Err(CliError::InvalidArgument(format!(
1799 "Server error: {}",
1800 error.message
1801 )));
1802 }
1803 if item.done {
1804 done = true;
1805 } else if let Some(row) = item.row {
1806 if streaming_writer.is_none() {
1807 let columns = default_stream_columns(row.len());
1808 let formatter = formatter
1809 .take()
1810 .ok_or_else(|| CliError::InvalidArgument("Missing formatter".into()))?;
1811 let mut writer =
1812 StreamingWriter::new(&mut *writer, formatter, columns, options.limit)
1813 .with_quiet(options.quiet);
1814 writer.prepare(None)?;
1815 streaming_writer = Some(writer);
1816 }
1817 let values = row.into_iter().map(remote_value_to_value).collect();
1818 if let Some(writer) = streaming_writer.as_mut() {
1819 match writer.write_row(Row::new(values))? {
1820 WriteStatus::LimitReached => {
1821 let _ = send_cancel_request(client).await;
1822 done = true;
1823 }
1824 WriteStatus::Continue => {}
1825 }
1826 }
1827 } else {
1828 return Err(CliError::InvalidArgument(
1829 "Invalid streaming response: missing row".into(),
1830 ));
1831 }
1832 }
1833 }
1834
1835 if !done {
1836 return Err(CliError::InvalidArgument(
1837 "Invalid streaming response: unexpected end of stream".into(),
1838 ));
1839 }
1840
1841 if let Some(mut writer) = streaming_writer {
1842 return writer.finish();
1843 }
1844
1845 if let Some(formatter) = formatter.take() {
1846 let mut writer = StreamingWriter::new(&mut *writer, formatter, Vec::new(), options.limit)
1847 .with_quiet(options.quiet);
1848 writer.prepare(None)?;
1849 return writer.finish();
1850 }
1851
1852 Ok(())
1853}
1854
1855fn default_stream_columns(count: usize) -> Vec<Column> {
1856 (0..count)
1857 .map(|idx| Column::new(format!("col{}", idx + 1), DataType::Text))
1858 .collect()
1859}
1860
1861fn parse_jsonl_line(line: &[u8]) -> Result<Option<RemoteStreamItem>> {
1862 if line
1863 .iter()
1864 .all(|b| matches!(b, b' ' | b'\n' | b'\r' | b'\t'))
1865 {
1866 return Ok(None);
1867 }
1868 let text = std::str::from_utf8(line)
1869 .map_err(|err| CliError::InvalidArgument(format!("Invalid UTF-8: {err}")))?;
1870 let trimmed = text.trim();
1871 if trimmed.is_empty() {
1872 return Ok(None);
1873 }
1874 let item = serde_json::from_str::<RemoteStreamItem>(trimmed)?;
1875 Ok(Some(item))
1876}
1877
1878fn skip_whitespace(buffer: &[u8], pos: &mut usize) {
1879 while *pos < buffer.len() {
1880 match buffer[*pos] {
1881 b' ' | b'\n' | b'\r' | b'\t' => *pos += 1,
1882 _ => break,
1883 }
1884 }
1885}
1886
1887fn has_non_whitespace(buffer: &[u8]) -> bool {
1888 buffer
1889 .iter()
1890 .any(|byte| !matches!(byte, b' ' | b'\n' | b'\r' | b'\t'))
1891}
1892
1893fn json_value_to_value(value: &serde_json::Value) -> Result<Value> {
1894 match value {
1895 serde_json::Value::Null => Ok(Value::Null),
1896 serde_json::Value::Bool(value) => Ok(Value::Bool(*value)),
1897 serde_json::Value::Number(value) => {
1898 if let Some(value) = value.as_i64() {
1899 Ok(Value::Int(value))
1900 } else if let Some(value) = value.as_f64() {
1901 Ok(Value::Float(value))
1902 } else {
1903 Err(CliError::InvalidArgument(
1904 "Invalid numeric value in streaming row".into(),
1905 ))
1906 }
1907 }
1908 serde_json::Value::String(value) => Ok(Value::Text(value.clone())),
1909 serde_json::Value::Array(values) => {
1910 let mut vector = Vec::with_capacity(values.len());
1911 for entry in values {
1912 let number = entry.as_f64().ok_or_else(|| {
1913 CliError::InvalidArgument("Invalid vector value in streaming row".into())
1914 })?;
1915 vector.push(number as f32);
1916 }
1917 Ok(Value::Vector(vector))
1918 }
1919 serde_json::Value::Object(_) => Err(CliError::InvalidArgument(
1920 "Invalid streaming row: nested objects are not supported".into(),
1921 )),
1922 }
1923}
1924
1925#[allow(dead_code)]
1927pub fn execute<W: Write>(
1928 db: &Database,
1929 cmd: SqlCommand,
1930 batch_mode: &BatchMode,
1931 writer: &mut StreamingWriter<W>,
1932) -> Result<()> {
1933 let sql = cmd.resolve_query(batch_mode)?;
1934
1935 execute_sql(db, &sql, writer)
1936}
1937
1938impl SqlCommand {
1939 pub fn resolve_query(&self, batch_mode: &BatchMode) -> Result<String> {
1941 match (&self.query, &self.file) {
1942 (Some(query), None) => Ok(query.clone()),
1943 (None, Some(file)) => fs::read_to_string(file).map_err(|e| {
1944 CliError::InvalidArgument(format!("Failed to read SQL file '{}': {}", file, e))
1945 }),
1946 (None, None) if !batch_mode.is_tty => {
1947 let mut buf = String::new();
1948 io::stdin().read_to_string(&mut buf)?;
1949 Ok(buf)
1950 }
1951 (None, None) => Err(CliError::NoQueryProvided),
1952 (Some(_), Some(_)) => Err(CliError::InvalidArgument(
1953 "Cannot specify both query and file".to_string(),
1954 )),
1955 }
1956 }
1957}
1958
1959fn execute_sql<W: Write>(db: &Database, sql: &str, writer: &mut StreamingWriter<W>) -> Result<()> {
1961 use alopex_embedded::SqlResult;
1962
1963 let result = db.execute_sql(sql)?;
1964
1965 match result {
1966 SqlResult::Success => {
1967 writer.prepare(Some(1))?;
1969 let row = Row::new(vec![
1970 Value::Text("OK".to_string()),
1971 Value::Text("Operation completed successfully".to_string()),
1972 ]);
1973 writer.write_row(row)?;
1974 writer.finish()?;
1975 }
1976 SqlResult::RowsAffected(count) => {
1977 writer.prepare(Some(1))?;
1979 let row = Row::new(vec![
1980 Value::Text("OK".to_string()),
1981 Value::Text(format!("{} row(s) affected", count)),
1982 ]);
1983 writer.write_row(row)?;
1984 writer.finish()?;
1985 }
1986 SqlResult::Query(query_result) => {
1987 let row_count = query_result.rows.len();
1989 writer.prepare(Some(row_count))?;
1990
1991 for sql_row in query_result.rows {
1992 let values: Vec<Value> = sql_row.into_iter().map(sql_value_to_value).collect();
1993 let row = Row::new(values);
1994
1995 match writer.write_row(row)? {
1996 WriteStatus::LimitReached => break,
1997 WriteStatus::Continue => {}
1998 }
1999 }
2000
2001 writer.finish()?;
2002 }
2003 }
2004
2005 Ok(())
2006}
2007
2008fn execute_sql_with_formatter<W: Write>(
2024 db: &Database,
2025 sql: &str,
2026 writer: &mut W,
2027 output: &mut SqlOutput<'_>,
2028 options: &SqlExecutionOptions<'_>,
2029) -> Result<()> {
2030 use alopex_sql::{AlopexDialect, Parser};
2031
2032 let dialect = AlopexDialect;
2035 let stmts = Parser::parse_sql(&dialect, sql).map_err(|e| CliError::Parse(format!("{}", e)))?;
2036
2037 let is_single_query = stmts.len() == 1
2038 && stmts
2039 .first()
2040 .is_some_and(|statement| statement.kind.is_query());
2041
2042 if is_single_query {
2043 let formatter = output.create_block_formatter();
2045 if !output.statement_array() {
2046 return execute_sql_select_streaming(db, sql, writer, formatter, options, false);
2047 }
2048 let mut result_set_writer = CountingWriter::new(&mut *writer);
2049 let result =
2050 execute_sql_select_streaming(db, sql, &mut result_set_writer, formatter, options, true);
2051 let output_started = result_set_writer.bytes_written() > 0;
2052 match result {
2053 Ok(()) => {
2054 writeln!(writer, "]")?;
2056 Ok(())
2057 }
2058 Err(err) => {
2059 if output_started {
2066 let _ = writeln!(writer, "]");
2067 let _ = writeln!(writer, "]");
2068 }
2069 Err(err)
2070 }
2071 }
2072 } else {
2073 execute_sql_statements(db, sql, writer, output, options)
2075 }
2076}
2077
2078fn execute_sql_select_streaming<W: Write>(
2088 db: &Database,
2089 sql: &str,
2090 writer: &mut W,
2091 formatter: Box<dyn Formatter>,
2092 options: &SqlExecutionOptions<'_>,
2093 open_statement_array: bool,
2094) -> Result<()> {
2095 use alopex_embedded::StreamingQueryResult;
2096 use std::sync::atomic::{AtomicBool, Ordering};
2097 use std::sync::Arc;
2098
2099 fn cli_err_to_embedded(e: crate::error::CliError) -> alopex_embedded::Error {
2101 alopex_embedded::Error::Sql(alopex_sql::SqlError::Execution {
2102 message: e.to_string(),
2103 code: "ALOPEX-C001",
2104 })
2105 }
2106
2107 let cancelled = Arc::new(AtomicBool::new(false));
2108 let timed_out = Arc::new(AtomicBool::new(false));
2109 let cancel_flag = cancelled.clone();
2110 let timeout_flag = timed_out.clone();
2111
2112 let result = db.execute_sql_with_rows(sql, |mut rows| {
2113 let columns = columns_from_streaming_rows(&rows);
2115 if open_statement_array {
2118 writeln!(writer, "[").map_err(|e| cli_err_to_embedded(CliError::Io(e)))?;
2119 }
2120 let mut streaming_writer = StreamingWriter::new(writer, formatter, columns, options.limit)
2121 .with_quiet(options.quiet);
2122
2123 streaming_writer
2125 .prepare(None)
2126 .map_err(cli_err_to_embedded)?;
2127
2128 if let Err(err) = options.deadline.check() {
2129 timeout_flag.store(true, Ordering::SeqCst);
2130 return Err(cli_err_to_embedded(err));
2131 }
2132
2133 while let Some(sql_row) = rows.next_row()? {
2137 if options.cancel.is_cancelled() {
2138 cancel_flag.store(true, Ordering::SeqCst);
2139 return Err(cli_err_to_embedded(CliError::Cancelled));
2140 }
2141 if let Err(err) = options.deadline.check() {
2142 timeout_flag.store(true, Ordering::SeqCst);
2143 return Err(cli_err_to_embedded(err));
2144 }
2145 let values: Vec<Value> = sql_row.into_iter().map(sql_value_to_value).collect();
2146 let row = Row::new(values);
2147
2148 match streaming_writer
2149 .write_row(row)
2150 .map_err(cli_err_to_embedded)?
2151 {
2152 WriteStatus::LimitReached => break,
2153 WriteStatus::Continue => {}
2154 }
2155 }
2156
2157 streaming_writer.finish().map_err(cli_err_to_embedded)?;
2158 Ok(())
2159 });
2160
2161 let result = match result {
2162 Ok(value) => value,
2163 Err(err) => {
2164 if cancelled.load(Ordering::SeqCst) {
2165 return Err(CliError::Cancelled);
2166 }
2167 if timed_out.load(Ordering::SeqCst) {
2168 return Err(CliError::Timeout(format!(
2169 "deadline exceeded after {}",
2170 humantime::format_duration(options.deadline.duration())
2171 )));
2172 }
2173 return Err(CliError::Database(err));
2174 }
2175 };
2176
2177 match result {
2178 StreamingQueryResult::QueryProcessed(()) => Ok(()),
2179 StreamingQueryResult::Success | StreamingQueryResult::RowsAffected(_) => {
2180 Ok(())
2182 }
2183 }
2184}
2185
2186#[derive(serde::Serialize)]
2187struct RemoteSqlRequest {
2188 sql: String,
2189 #[serde(default)]
2190 streaming: bool,
2191 #[serde(skip_serializing_if = "Option::is_none")]
2192 fetch_size: Option<usize>,
2193 #[serde(skip_serializing_if = "Option::is_none")]
2194 max_rows: Option<usize>,
2195}
2196
2197#[derive(serde::Serialize)]
2198struct RemoteDistributedReadRequest {
2199 sql: String,
2200 read_mode: &'static str,
2201}
2202
2203#[derive(serde::Deserialize)]
2204struct RemoteColumnInfo {
2205 name: String,
2206 data_type: String,
2207}
2208
2209#[derive(serde::Deserialize, Default)]
2210struct RemoteSqlResult {
2211 #[serde(default)]
2212 columns: Vec<RemoteColumnInfo>,
2213 #[serde(default)]
2214 rows: Vec<Vec<alopex_sql::storage::SqlValue>>,
2215 #[serde(default)]
2216 affected_rows: Option<u64>,
2217}
2218
2219#[derive(serde::Deserialize)]
2220struct RemoteSqlResponse {
2221 #[serde(default)]
2222 results: Vec<RemoteSqlResult>,
2223 #[serde(flatten)]
2224 legacy: RemoteSqlResult,
2225}
2226
2227impl RemoteSqlResponse {
2228 fn into_results(self) -> Vec<RemoteSqlResult> {
2229 if self.results.is_empty() {
2230 vec![self.legacy]
2231 } else {
2232 self.results
2233 }
2234 }
2235}
2236
2237#[derive(serde::Deserialize)]
2238struct RemoteStreamItem {
2239 row: Option<Vec<alopex_sql::storage::SqlValue>>,
2240 error: Option<RemoteStreamError>,
2241 #[serde(default)]
2242 done: bool,
2243}
2244
2245#[derive(serde::Deserialize)]
2246struct RemoteStreamError {
2247 message: String,
2248}
2249
2250fn map_client_error(err: ClientError) -> CliError {
2251 match err {
2252 ClientError::Request { source, .. } => {
2253 CliError::ServerConnection(format!("request failed: {source}"))
2254 }
2255 ClientError::InvalidUrl(message) => CliError::InvalidArgument(message),
2256 ClientError::Build(message) => CliError::InvalidArgument(message),
2257 ClientError::Auth(err) => CliError::InvalidArgument(err.to_string()),
2258 ClientError::HttpStatus { status, body } => {
2259 CliError::InvalidArgument(format!("Server error: HTTP {} - {}", status.as_u16(), body))
2260 }
2261 }
2262}
2263
2264async fn send_cancel_request(client: &HttpClient) -> Result<()> {
2265 #[derive(serde::Serialize)]
2266 struct CancelRequest {}
2267
2268 let request = CancelRequest {};
2269 let _: serde_json::Value = client
2270 .post_json("api/sql/cancel", &request)
2271 .await
2272 .map_err(map_client_error)?;
2273 Ok(())
2274}
2275
2276fn merge_limit(limit: Option<usize>, max_rows: Option<usize>) -> Option<usize> {
2277 match (limit, max_rows) {
2278 (Some(a), Some(b)) => Some(a.min(b)),
2279 (Some(value), None) | (None, Some(value)) => Some(value),
2280 (None, None) => None,
2281 }
2282}
2283
2284fn execute_sql_statements<W: Write>(
2297 db: &Database,
2298 sql: &str,
2299 writer: &mut W,
2300 output: &mut SqlOutput<'_>,
2301 options: &SqlExecutionOptions<'_>,
2302) -> Result<()> {
2303 use alopex_sql::ExecutionResult;
2304
2305 options.deadline.check()?;
2306 let results = db.execute_sql_multi(sql)?;
2307 options.deadline.check()?;
2308
2309 let json_array = output.statement_array();
2310 if json_array {
2311 writeln!(writer, "[")?;
2312 }
2313 let mut first = true;
2314 for result in results {
2315 let is_status = !matches!(result, ExecutionResult::Query(_));
2316 if is_status && options.quiet {
2317 continue;
2318 }
2319 if json_array && !first {
2320 writeln!(writer, ",")?;
2321 }
2322 first = false;
2323 emit_execution_result(result, writer, output, options)?;
2324 }
2325 if json_array {
2326 writeln!(writer, "]")?;
2327 }
2328 Ok(())
2329}
2330
2331fn emit_execution_result<W: Write>(
2333 result: alopex_sql::ExecutionResult,
2334 writer: &mut W,
2335 output: &mut SqlOutput<'_>,
2336 options: &SqlExecutionOptions<'_>,
2337) -> Result<()> {
2338 use alopex_sql::ExecutionResult;
2339
2340 match result {
2341 ExecutionResult::Success => emit_status_result(
2342 writer,
2343 output,
2344 options,
2345 "Operation completed successfully".to_string(),
2346 ),
2347 ExecutionResult::RowsAffected(count) => {
2348 emit_status_result(writer, output, options, format!("{count} row(s) affected"))
2349 }
2350 ExecutionResult::Query(query_result) => {
2351 let columns = columns_from_query_result(&query_result);
2352 let mut streaming_writer = StreamingWriter::new(
2353 &mut *writer,
2354 output.create_block_formatter(),
2355 columns,
2356 options.limit,
2357 )
2358 .with_quiet(options.quiet);
2359 streaming_writer.prepare(Some(query_result.rows.len()))?;
2360 for sql_row in query_result.rows {
2361 let values: Vec<Value> = sql_row.into_iter().map(sql_value_to_value).collect();
2362 match streaming_writer.write_row(Row::new(values))? {
2363 WriteStatus::LimitReached => break,
2364 WriteStatus::Continue => {}
2365 }
2366 }
2367 streaming_writer.finish()
2368 }
2369 }
2370}
2371
2372fn emit_status_result<W: Write>(
2374 writer: &mut W,
2375 output: &mut SqlOutput<'_>,
2376 options: &SqlExecutionOptions<'_>,
2377 message: String,
2378) -> Result<()> {
2379 let columns = sql_status_columns();
2380 let mut streaming_writer = StreamingWriter::new(
2381 &mut *writer,
2382 output.create_block_formatter(),
2383 columns,
2384 options.limit,
2385 )
2386 .with_quiet(options.quiet);
2387 streaming_writer.prepare(Some(1))?;
2388 streaming_writer.write_row(Row::new(vec![
2389 Value::Text("OK".to_string()),
2390 Value::Text(message),
2391 ]))?;
2392 streaming_writer.finish()
2393}
2394
2395fn sql_value_to_value(sql_value: alopex_sql::SqlValue) -> Value {
2397 use alopex_sql::SqlValue;
2398
2399 match sql_value {
2400 SqlValue::Null => Value::Null,
2401 SqlValue::Integer(i) => Value::Int(i as i64),
2402 SqlValue::BigInt(i) => Value::Int(i),
2403 SqlValue::Float(f) => Value::Float(f as f64),
2404 SqlValue::Double(f) => Value::Float(f),
2405 SqlValue::Text(s) => Value::Text(s),
2406 SqlValue::Blob(b) => Value::Bytes(b),
2407 SqlValue::Boolean(b) => Value::Bool(b),
2408 SqlValue::Timestamp(ts) => {
2409 Value::Text(format!("{}", ts))
2411 }
2412 SqlValue::Vector(v) => Value::Vector(v),
2413 value @ (SqlValue::Date(_) | SqlValue::Time(_) | SqlValue::Interval { .. }) => {
2414 Value::Text(value.temporal_text().expect("valid stored temporal value"))
2415 }
2416 SqlValue::Decimal(value) => Value::Text(value.to_string()),
2417 SqlValue::Json(value) => Value::Text(value.to_string()),
2418 value @ (SqlValue::Array(_) | SqlValue::Map(_) | SqlValue::Struct(_)) => Value::Text(
2419 value
2420 .nested_json_text()
2421 .expect("nested SQL value has a JSON mapping"),
2422 ),
2423 }
2424}
2425
2426fn remote_value_to_value(sql_value: alopex_sql::storage::SqlValue) -> Value {
2427 use alopex_sql::storage::SqlValue;
2428
2429 match sql_value {
2430 SqlValue::Null => Value::Null,
2431 SqlValue::Integer(i) => Value::Int(i as i64),
2432 SqlValue::BigInt(i) => Value::Int(i),
2433 SqlValue::Float(f) => Value::Float(f as f64),
2434 SqlValue::Double(f) => Value::Float(f),
2435 SqlValue::Text(s) => Value::Text(s),
2436 SqlValue::Blob(b) => Value::Bytes(b),
2437 SqlValue::Boolean(b) => Value::Bool(b),
2438 SqlValue::Timestamp(ts) => Value::Text(ts.to_string()),
2439 SqlValue::Vector(v) => Value::Vector(v),
2440 value @ (SqlValue::Date(_) | SqlValue::Time(_) | SqlValue::Interval { .. }) => {
2441 Value::Text(value.temporal_text().expect("valid stored temporal value"))
2442 }
2443 SqlValue::Decimal(value) => Value::Text(value.to_string()),
2444 SqlValue::Json(value) => Value::Text(value.to_string()),
2445 value @ (SqlValue::Array(_) | SqlValue::Map(_) | SqlValue::Struct(_)) => Value::Text(
2446 value
2447 .nested_json_text()
2448 .expect("nested SQL value has a JSON mapping"),
2449 ),
2450 }
2451}
2452
2453fn data_type_from_string(value: &str) -> DataType {
2454 let upper = value.to_ascii_uppercase();
2455 if upper.starts_with("INT") || upper.starts_with("BIGINT") {
2456 DataType::Int
2457 } else if upper.starts_with("FLOAT") || upper.starts_with("DOUBLE") {
2458 DataType::Float
2459 } else if upper.starts_with("BLOB") {
2460 DataType::Bytes
2461 } else if upper.starts_with("BOOLEAN") {
2462 DataType::Bool
2463 } else if upper.starts_with("VECTOR") {
2464 DataType::Vector
2465 } else {
2466 DataType::Text
2467 }
2468}
2469
2470fn sql_column_to_column(col: &alopex_sql::executor::ColumnInfo) -> Column {
2472 use alopex_sql::planner::ResolvedType;
2473
2474 let data_type = match &col.data_type {
2475 ResolvedType::Integer | ResolvedType::BigInt => DataType::Int,
2476 ResolvedType::Float | ResolvedType::Double => DataType::Float,
2477 ResolvedType::Text => DataType::Text,
2478 ResolvedType::Blob => DataType::Bytes,
2479 ResolvedType::Boolean => DataType::Bool,
2480 ResolvedType::Timestamp => DataType::Text, ResolvedType::Date | ResolvedType::Time | ResolvedType::Interval => DataType::Text,
2482 ResolvedType::Decimal { .. } => DataType::Text,
2483 ResolvedType::Json => DataType::Text,
2484 ResolvedType::Array(_) | ResolvedType::Map { .. } | ResolvedType::Struct(_) => {
2485 DataType::Text
2486 }
2487 ResolvedType::Vector { .. } => DataType::Vector,
2488 ResolvedType::Null => DataType::Text, };
2490
2491 Column::new(&col.name, data_type)
2492}
2493
2494fn columns_from_query_result(query_result: &alopex_sql::executor::QueryResult) -> Vec<Column> {
2496 query_result
2497 .columns
2498 .iter()
2499 .map(sql_column_to_column)
2500 .collect()
2501}
2502
2503#[allow(dead_code)] fn columns_from_streaming_result(
2506 query_iter: &alopex_embedded::QueryRowIterator<'_>,
2507) -> Vec<Column> {
2508 query_iter
2509 .columns()
2510 .iter()
2511 .map(sql_column_to_column)
2512 .collect()
2513}
2514
2515fn columns_from_streaming_rows(rows: &alopex_embedded::StreamingRows<'_>) -> Vec<Column> {
2517 rows.columns().iter().map(sql_column_to_column).collect()
2518}
2519
2520pub fn sql_status_columns() -> Vec<Column> {
2522 vec![
2523 Column::new("status", DataType::Text),
2524 Column::new("message", DataType::Text),
2525 ]
2526}
2527
2528#[cfg(test)]
2529mod tests {
2530 use super::*;
2531 use crate::batch::BatchModeSource;
2532 use crate::output::jsonl::JsonlFormatter;
2533
2534 fn create_test_db() -> Database {
2535 Database::open_in_memory().unwrap()
2536 }
2537
2538 fn create_status_writer(output: &mut Vec<u8>) -> StreamingWriter<&mut Vec<u8>> {
2539 let formatter = Box::new(JsonlFormatter::new());
2540 let columns = sql_status_columns();
2541 StreamingWriter::new(output, formatter, columns, None)
2542 }
2543
2544 fn create_query_writer(
2545 output: &mut Vec<u8>,
2546 columns: Vec<Column>,
2547 ) -> StreamingWriter<&mut Vec<u8>> {
2548 let formatter = Box::new(JsonlFormatter::new());
2549 StreamingWriter::new(output, formatter, columns, None)
2550 }
2551
2552 fn default_batch_mode() -> BatchMode {
2553 BatchMode {
2554 is_batch: false,
2555 is_tty: true,
2556 source: BatchModeSource::Default,
2557 }
2558 }
2559
2560 #[test]
2561 fn test_create_table() {
2562 let db = create_test_db();
2563
2564 let mut output = Vec::new();
2565 {
2566 let mut writer = create_status_writer(&mut output);
2567 execute_sql(
2568 &db,
2569 "CREATE TABLE users (id INTEGER PRIMARY KEY, name TEXT);",
2570 &mut writer,
2571 )
2572 .unwrap();
2573 }
2574
2575 let result = String::from_utf8(output).unwrap();
2576 assert!(result.contains("OK"));
2577 }
2578
2579 #[test]
2580 fn test_insert_and_select() {
2581 let db = create_test_db();
2582
2583 {
2585 let mut output = Vec::new();
2586 let mut writer = create_status_writer(&mut output);
2587 execute_sql(
2588 &db,
2589 "CREATE TABLE users (id INTEGER PRIMARY KEY, name TEXT);",
2590 &mut writer,
2591 )
2592 .unwrap();
2593 }
2594
2595 {
2597 let mut output = Vec::new();
2598 let mut writer = create_status_writer(&mut output);
2599 execute_sql(
2600 &db,
2601 "INSERT INTO users (id, name) VALUES (1, 'Alice');",
2602 &mut writer,
2603 )
2604 .unwrap();
2605 let result = String::from_utf8(output).unwrap();
2606 assert!(result.contains("row(s) affected"));
2607 }
2608
2609 {
2611 let mut output = Vec::new();
2612 let columns = vec![
2613 Column::new("id", DataType::Int),
2614 Column::new("name", DataType::Text),
2615 ];
2616 let mut writer = create_query_writer(&mut output, columns);
2617 execute_sql(&db, "SELECT id, name FROM users;", &mut writer).unwrap();
2618 let result = String::from_utf8(output).unwrap();
2619 assert!(result.contains("Alice"));
2620 }
2621 }
2622
2623 #[test]
2624 fn test_syntax_error() {
2625 let db = create_test_db();
2626
2627 let mut output = Vec::new();
2628 let mut writer = create_status_writer(&mut output);
2629 let result = execute_sql(&db, "CREATE TABEL invalid_syntax;", &mut writer);
2630 assert!(result.is_err());
2631 }
2632
2633 #[test]
2634 fn test_multiple_statements() {
2635 let db = create_test_db();
2636
2637 let mut output = Vec::new();
2638 {
2639 let mut writer = create_status_writer(&mut output);
2640 execute_sql(
2641 &db,
2642 "CREATE TABLE t (id INTEGER PRIMARY KEY); INSERT INTO t (id) VALUES (1);",
2643 &mut writer,
2644 )
2645 .unwrap();
2646 }
2647
2648 {
2650 let mut output = Vec::new();
2651 let columns = vec![Column::new("id", DataType::Int)];
2652 let mut writer = create_query_writer(&mut output, columns);
2653 execute_sql(&db, "SELECT id FROM t;", &mut writer).unwrap();
2654 let result = String::from_utf8(output).unwrap();
2655 assert!(result.contains("1"));
2656 }
2657 }
2658
2659 #[test]
2664 fn streaming_select_row_error_propagates() {
2665 let db = create_test_db();
2666 db.execute_sql("CREATE TABLE t (id INTEGER PRIMARY KEY);")
2667 .unwrap();
2668 db.execute_sql("INSERT INTO t (id) VALUES (1);").unwrap();
2669
2670 let deadline = Deadline::new(parse_deadline(None).unwrap());
2671 let cancel = CancelSignal::new();
2672 let options = SqlExecutionOptions {
2673 limit: None,
2674 quiet: false,
2675 cancel: &cancel,
2676 deadline: &deadline,
2677 admin_launcher: None,
2678 };
2679
2680 let mut output = Vec::new();
2681 let formatter = Box::new(JsonlFormatter::new());
2682 let result = execute_sql_select_streaming(
2685 &db,
2686 "SELECT id / 0 AS x FROM t;",
2687 &mut output,
2688 formatter,
2689 &options,
2690 false,
2691 );
2692 assert!(
2693 result.is_err(),
2694 "row evaluation errors must propagate instead of yielding an empty result"
2695 );
2696 }
2697
2698 #[test]
2699 fn test_sql_value_conversion() {
2700 use alopex_sql::SqlValue;
2701
2702 assert!(matches!(sql_value_to_value(SqlValue::Null), Value::Null));
2703 assert!(matches!(
2704 sql_value_to_value(SqlValue::Integer(42)),
2705 Value::Int(42)
2706 ));
2707 assert!(matches!(
2708 sql_value_to_value(SqlValue::BigInt(100)),
2709 Value::Int(100)
2710 ));
2711 assert!(matches!(
2712 sql_value_to_value(SqlValue::Boolean(true)),
2713 Value::Bool(true)
2714 ));
2715 assert!(
2716 matches!(sql_value_to_value(SqlValue::Text("hello".to_string())), Value::Text(s) if s == "hello")
2717 );
2718 }
2719
2720 #[test]
2721 fn resolve_query_from_argument() {
2722 let cmd = SqlCommand {
2723 query: Some("SELECT 1".to_string()),
2724 file: None,
2725 fetch_size: None,
2726 max_rows: None,
2727 deadline: None,
2728 read_mode: None,
2729 routing_report: None,
2730 tui: false,
2731 };
2732
2733 let sql = cmd.resolve_query(&default_batch_mode()).unwrap();
2734 assert_eq!(sql, "SELECT 1");
2735 }
2736
2737 #[test]
2738 fn resolve_query_from_file() {
2739 let mut file = tempfile::NamedTempFile::new().unwrap();
2740 writeln!(file, "SELECT * FROM users").unwrap();
2741
2742 let cmd = SqlCommand {
2743 query: None,
2744 file: Some(file.path().display().to_string()),
2745 fetch_size: None,
2746 max_rows: None,
2747 deadline: None,
2748 read_mode: None,
2749 routing_report: None,
2750 tui: false,
2751 };
2752
2753 let sql = cmd.resolve_query(&default_batch_mode()).unwrap();
2754 assert_eq!(sql, "SELECT * FROM users\n");
2755 }
2756
2757 #[test]
2758 fn resolve_query_returns_no_query_error() {
2759 let cmd = SqlCommand {
2760 query: None,
2761 file: None,
2762 fetch_size: None,
2763 max_rows: None,
2764 deadline: None,
2765 read_mode: None,
2766 routing_report: None,
2767 tui: false,
2768 };
2769
2770 let err = cmd.resolve_query(&default_batch_mode()).unwrap_err();
2771 assert!(matches!(err, CliError::NoQueryProvided));
2772 }
2773
2774 #[test]
2775 fn resolve_query_rejects_query_and_file() {
2776 let cmd = SqlCommand {
2777 query: Some("SELECT 1".to_string()),
2778 file: Some("query.sql".into()),
2779 fetch_size: None,
2780 max_rows: None,
2781 deadline: None,
2782 read_mode: None,
2783 routing_report: None,
2784 tui: false,
2785 };
2786
2787 let err = cmd.resolve_query(&default_batch_mode()).unwrap_err();
2788 assert!(matches!(
2789 err,
2790 CliError::InvalidArgument(msg) if msg == "Cannot specify both query and file"
2791 ));
2792 }
2793}